import time,random, copy
import multiprocessing
import threading

def map_single(fn):
    def funk(lst):
        nulst = [   fn(n)
                    for n in lst    ]
        return nulst
    return funk

def map_singleB(fn):
    def funk(lst):
        nulst = copy.deepcopy(lst)
        for i in range(len(nulst)):
            nulst[i] = fn(nulst[i])
        return nulst
    return funk

def map_multi(n,fn):
    #enqueue = multiprocessing.Queue(999999999)
    #dequeue = multiprocessing.Queue(999999999)
    #toprocess, frommain = multiprocessing.Pipe()
    #fromprocess, tomain = multiprocessing.Pipe()
    def process(n,enqueue,dequeue,loop=True):
        score = 0
        while loop:
            data,loop = enqueue.get()
            #data,loop = frommain.recv()
            if loop:
                score = score + 1
                nudata = fn(data)
                dequeue.put(nudata)
                #tomain.send(nudata)
        print "process",n,"score",score
    def generateProcess(i,process):
        enqueue = multiprocessing.Queue()
        dequeue = multiprocessing.Queue()
        process = multiprocessing.Process(target=process, args=(i,enqueue,dequeue))
        return (process,enqueue,dequeue)
    prlst = [   generateProcess(i,process) for i in range(n)  ]
    def funk(lst):
        def threadA():
            start_all_process = [   pr.start() for pr,_,_ in prlst   ]
            wait_all_process = [    pr.join() for pr,_,_ in prlst   ]
            print "threadA done"
        def threadB():
            prid = 0
            for data in lst:
                process,enqueue,dequeue = prlst[prid]
                enqueue.put((data,True))
                prid = (prid + 1) % n
            
            send_to_process = [ enqueue.put((data,True)) for data in lst   ]
            shut_all_process = [enqueue.put((None,None)) for _,enqueue,_ in prlst  ]
            #send_to_process = [ fromprocess.send((data,True)) for data in lst   ]
            #shut_all_process = [toprocess.send((None,None)) for pr in prlst  ]
            print "threadB done"
        thrlst = [ threading.Thread(target=thread) for thread in [threadA, threadB ]]
        start_all_thread = [ t.start() for t in thrlst  ]
        nulst = zip([data for data in lst],(dequeue.get() for _,_,dequeue in prlst))
        #nulst = [   fromprocess.recv() for i in range(len(lst)) ]
        wait_all_thread = [ t.join() for t in thrlst    ]
        return nulst
    return funk




def bench(fnmap,lst):
    start = time.time()
    listB = fnmap(lst)
    end = time.time()
    return end - start

if __name__ == '__main__':
    print "generate random list"
    listA = list(range(999999))
    random.shuffle(listA)

    def inc(n):
        return n+1

    for _name,_map in { 'map single thread' : map_single(inc), 
                        'map single thread B' : map_singleB(inc),
                        'map single process queue' : map_multi(1,inc),
                        'map dual process queue' : map_multi(2,inc),
                        'map triple process queue' : map_multi(3,inc),
                        'map quad process queue' : map_multi(4,inc),
                        'map octo process queue' : map_multi(8,inc)
                        }.items():
        print "__________________________________________"
        print "begin",_name
        tim = bench(_map,listA)
        print _name,"took",tim,"seconds to complete"

    
