o
    l¨Êh—²  ã                   @   sØ  d Z dZddlZddlZddlZddlZddlZddlZddlZddl	Z	ddl
Z
ddlmZ ddlZddlmZ ddlmZ ddlmZ dd	lmZ dd
lmZ ddlmZ ddlmZ ddlmZ ddlmZmZmZ ddl m!Z!m"Z" ddl#m$Z$m%Z% zddl&m'Z( W n e)y�   e*Z(Y nw ej+d dkr™e,Z-e .¡ Z/da0e1ej2 3dd¡ƒZ4da5dZ6e1dƒZ7zddl8m9Z9 dZ:dLdd„Z;W n e)yÊ   dZ:Y nw G dd„ dƒZ<G dd„ de=ƒZ>d d!„ Z?da@dZAG d"d#„ d#eBƒZCG d$d%„ d%eDƒZEd&d'„ ZFG d(d)„ d)e=ƒZGG d*d+„ d+e=ƒZHG d,d-„ d-e=ƒZIG d.d/„ d/eƒZJd0d1„ ZKd2d3„ ZLdMd4d5„ZMd6d7„ ZNd8d9„ ZOd:d;„ ZPdaQdaRd<d=„ ZSd>d?„ ZTd@dA„ ZUG dBdC„ dCe*ƒZVG dDdE„ dEe(ƒZ'G dFdG„ dGe'ƒZWe'ZXG dHdI„ dIe*ƒZYG dJdK„ dKejZƒZ[dS )Na*	  Implements ProcessPoolExecutor.

The follow diagram and text describe the data-flow through the system:

|======================= In-process =====================|== Out-of-process ==|

+----------+     +----------+       +--------+     +-----------+    +---------+
|          |  => | Work Ids |       |        |     | Call Q    |    | Process |
|          |     +----------+       |        |     +-----------+    |  Pool   |
|          |     | ...      |       |        |     | ...       |    +---------+
|          |     | 6        |    => |        |  => | 5, call() | => |         |
|          |     | 7        |       |        |     | ...       |    |         |
| Process  |     | ...      |       | Local  |     +-----------+    | Process |
|  Pool    |     +----------+       | Worker |                      |  #1..n  |
| Executor |                        | Thread |                      |         |
|          |     +----------- +     |        |     +-----------+    |         |
|          | <=> | Work Items | <=> |        | <=  | Result Q  | <= |         |
|          |     +------------+     |        |     +-----------+    |         |
|          |     | 6: call()  |     |        |     | ...       |    |         |
|          |     |    future  |     +--------+     | 4, result |    |         |
|          |     | ...        |                    | 3, except |    |         |
+----------+     +------------+                    +-----------+    +---------+

Executor.submit() called:
- creates a uniquely numbered _WorkItem and adds it to the "Work Items" dict
- adds the id of the _WorkItem to the "Work Ids" queue

Local worker thread:
- reads work ids from the "Work Ids" queue and looks up the corresponding
  WorkItem from the "Work Items" dict: if the work item has been cancelled then
  it is simply removed from the dict, otherwise it is repackaged as a
  _CallItem and put in the "Call Q". New _CallItems are put in the "Call Q"
  until "Call Q" is full. NOTE: the size of the "Call Q" is kept small because
  calls placed in the "Call Q" can no longer be cancelled with Future.cancel().
- reads _ResultItems from "Result Q", updates the future stored in the
  "Work Items" dict and deletes the dict entry

Process #1..n:
- reads _CallItems from "Call Q", executes the calls, and puts the resulting
  _ResultItems in "Result Q"
z,Thomas Moreau (thomas.moreau.2010@gmail.com)é    N)Útime)Úpartial)ÚPicklingErroré   )Ú_base)Úget_context)Úqueue)Úwait)Ú	set_cause)Ú	cpu_count)ÚQueueÚSimpleQueueÚFull)Úset_loky_picklerÚget_loky_pickler_name)Úrecursive_terminateÚget_exitcodes_terminated_worker)ÚBrokenProcessPoolé   FÚLOKY_MAX_DEPTHé
   g      ð?g    „×—A)ÚProcessTc                 C   s   |rt  ¡  t| ƒ ¡ jS ©N)ÚgcÚcollectr   Úmemory_infoÚrss)ÚpidÚforce_gc© r   úX/var/www/html/env/lib/python3.10/site-packages/joblib/externals/loky/process_executor.pyÚ_get_memory_usageƒ   s   r!   c                   @   s,   e Zd Zdd„ Zdd„ Zdd„ Zdd„ Zd	S )
Ú_ThreadWakeupc                 C   s   t jdd�\| _| _d S )NF)Úduplex)ÚmpÚPipeÚ_readerÚ_writer©Úselfr   r   r    Ú__init__Ž   s   z_ThreadWakeup.__init__c                 C   s   | j  ¡  | j ¡  d S r   )r'   Úcloser&   r(   r   r   r    r+   ‘   s   
z_ThreadWakeup.closec                 C   s<   t jdkrt jd d… dk r| j d¡ d S | j d¡ d S )NÚwin32r   )é   é   ó   0ó    )ÚsysÚplatformÚversion_infor'   Ú
send_bytesr(   r   r   r    Úwakeup•   s   z_ThreadWakeup.wakeupc                 C   s&   | j  ¡ r| j  ¡  | j  ¡ sd S d S r   )r&   ÚpollÚ
recv_bytesr(   r   r   r    Úclear�   s   

ÿz_ThreadWakeup.clearN)Ú__name__Ú
__module__Ú__qualname__r*   r+   r5   r8   r   r   r   r    r"   �   s
    r"   c                   @   s*   e Zd ZdZdd„ Zd
dd„Zdd„ Zd	S )Ú_ExecutorFlagsa  necessary references to maintain executor states without preventing gc

    It permits to keep the information needed by queue_management_thread
    and crash_detection_thread to maintain the pool without preventing the
    garbage collection of unreferenced executors.
    c                 C   s    d| _ d | _d| _t ¡ | _d S )NF)ÚshutdownÚbrokenÚkill_workersÚ	threadingÚLockÚshutdown_lockr(   r   r   r    r*   ©   s   z_ExecutorFlags.__init__Fc                 C   ó8   | j � d| _|| _W d   ƒ d S 1 sw   Y  d S ©NT)rB   r=   r?   )r)   r?   r   r   r    Úflag_as_shutting_down°   ó   "þz$_ExecutorFlags.flag_as_shutting_downc                 C   rC   rD   )rB   r=   r>   )r)   r>   r   r   r    Úflag_as_brokenµ   rF   z_ExecutorFlags.flag_as_brokenN©F)r9   r:   r;   Ú__doc__r*   rE   rG   r   r   r   r    r<   ¢   s
    
r<   c                  C   sZ   da tt ¡ ƒ} tj d | ¡¡ | D ]\}}| ¡ r| 	¡  q| D ]\}}| 
¡  q"d S )NTz=Interpreter shutting down. Waking up queue_manager_threads {})Ú_global_shutdownÚlistÚ_threads_wakeupsÚitemsr$   ÚutilÚdebugÚformatÚis_aliver5   Újoin)rM   ÚthreadÚthread_wakeupÚ_r   r   r    Ú_python_exit»   s   ÿ€
ÿrV   c                   @   s"   e Zd ZdZddd„Zdd„ ZdS )Ú_RemoteTracebackzAEmbed stringification of remote traceback in local traceback
    Nc                 C   s
   || _ d S r   ©Útb)r)   rY   r   r   r    r*   Õ   s   
z_RemoteTraceback.__init__c                 C   s   | j S r   rX   r(   r   r   r    Ú__str__Ø   s   z_RemoteTraceback.__str__r   )r9   r:   r;   rI   r*   rZ   r   r   r   r    rW   Ò   s    
rW   c                   @   s   e Zd Zdd„ Zdd„ ZdS )Ú_ExceptionWithTracebackc                 C   sR   t |dd ƒ}|d u rt ¡ \}}}t t|ƒ||¡}d |¡}|| _d| | _d S )NÚ__traceback__Ú z

"""
%s""")	Úgetattrr1   Úexc_infoÚ	tracebackÚformat_exceptionÚtyperR   ÚexcrY   )r)   rc   rY   rU   r   r   r    r*   Þ   s   
z _ExceptionWithTraceback.__init__c                 C   s   t | j| jffS r   )Ú_rebuild_excrc   rY   r(   r   r   r    Ú
__reduce__ç   s   z"_ExceptionWithTraceback.__reduce__N)r9   r:   r;   r*   re   r   r   r   r    r[   Ü   s    	r[   c                 C   s   t | t|ƒƒ} | S r   )r
   rW   )rc   rY   r   r   r    rd   ë   s   rd   c                   @   s   e Zd Zg d¢Zdd„ ZdS )Ú	_WorkItem©ÚfutureÚfnÚargsÚkwargsc                 C   s   || _ || _|| _|| _d S r   rg   )r)   rh   ri   rj   rk   r   r   r    r*   ô   s   
z_WorkItem.__init__N)r9   r:   r;   Ú	__slots__r*   r   r   r   r    rf   ð   s    rf   c                   @   s   e Zd Zddd„ZdS )Ú_ResultItemNc                 C   s   || _ || _|| _d S r   )Úwork_idÚ	exceptionÚresult)r)   rn   ro   rp   r   r   r    r*   ý   s   
z_ResultItem.__init__©NN)r9   r:   r;   r*   r   r   r   r    rm   û   s    rm   c                   @   s$   e Zd Zdd„ Zdd„ Zdd„ ZdS )Ú	_CallItemc                 C   s$   || _ || _|| _|| _tƒ | _d S r   )rn   ri   rj   rk   r   Úloky_pickler)r)   rn   ri   rj   rk   r   r   r    r*     s
   z_CallItem.__init__c                 C   s   t | jƒ | j| ji | j¤ŽS r   )r   rs   ri   rj   rk   r(   r   r   r    Ú__call__  s   
z_CallItem.__call__c                 C   s   d  | j| j| j| j¡S )NzCallItem({}, {}, {}, {}))rP   rn   ri   rj   rk   r(   r   r   r    Ú__repr__  s   ÿz_CallItem.__repr__N)r9   r:   r;   r*   rt   ru   r   r   r   r    rr     s    	rr   c                       s2   e Zd ZdZ		d‡ fdd„	Z‡ fdd„Z‡  ZS )	Ú
_SafeQueuez=Safe Queue set exception to the future object linked to a jobr   Nc                    s,   || _ || _|| _tt| ƒj|||d� d S )N©ÚreducersÚctx)rT   Úpending_work_itemsÚrunning_work_itemsÚsuperrv   r*   )r)   Úmax_sizery   rz   r{   rT   rx   ©Ú	__class__r   r    r*     s   z_SafeQueue.__init__c                    s´   t |tƒrOt |tjƒrtdƒ}ntdƒ}t t|ƒ|t	|dd ƒ¡}t
|td d |¡¡ƒƒ}| j |jd ¡}| j |j¡ |d urH|j |¡ ~| j ¡  d S tt| ƒ ||¡ d S )NzNThe task could not be sent to the workers as it is too large for `send_bytes`.z4Could not pickle the task to send it to the workers.r\   z

"""
{}"""r]   )Ú
isinstancerr   ÚstructÚerrorÚRuntimeErrorr   r`   ra   rb   r^   r
   rW   rP   rR   rz   Úpoprn   r{   Úremoverh   Úset_exceptionrT   r5   r|   rv   Ú_on_queue_feeder_error)r)   ÚeÚobjÚraised_errorrY   Ú	work_itemr~   r   r    r‡      s*   
ÿÿÿÿz!_SafeQueue._on_queue_feeder_error)r   NNNNN)r9   r:   r;   rI   r*   r‡   Ú__classcell__r   r   r~   r    rv     s    ÿrv   c                 g   sB   � t jdk rtj|Ž }nt|Ž }	 tt || ¡ƒ}|sdS |V  q)z+Iterates over zip()ed iterables in chunks. )r-   r-   TN)r1   r3   Ú	itertoolsÚizipÚzipÚtupleÚislice)Ú	chunksizeÚ	iterablesÚitÚchunkr   r   r    Ú_get_chunks;  s   €
ür–   c                    s   ‡ fdd„|D ƒS )z»Processes a chunk of an iterable passed to map.

    Runs the function passed to map() on a chunk of the
    iterable passed to map.

    This function is run in a separate process.

    c                    s   g | ]}ˆ |Ž ‘qS r   r   )Ú.0rj   ©ri   r   r    Ú
<listcomp>Q  s    z"_process_chunk.<locals>.<listcomp>r   )ri   r•   r   r˜   r    Ú_process_chunkH  s   	rš   c              
   C   s\   z|   t|||d�¡ W dS  ty- } zt|ƒ}|   t||d�¡ W Y d}~dS d}~ww )z.Safely send back the given result or exception)rp   ro   ©ro   N)Úputrm   ÚBaseExceptionr[   )Úresult_queuern   rp   ro   rˆ   rc   r   r   r    Ú_sendback_resultT  s   
ÿ €þrŸ   c                 C   sž  |durz||Ž  W n t y   tjjddd� Y dS w |ad}d}	t ¡ }
tj 	d| ¡ 	 z| j
d|d�}|du rBtj d¡ W nO tjyj   tj d| ¡ |jd	d
�r`| ¡  d}ntj d¡ Y q/Y n) t y’   t ¡ }z	| t|ƒ¡ W n t yŠ   t|ƒ Y nw t d¡ Y nw |du r°| |
¡ |�
 	 W d  ƒ dS 1 s«w   Y  z|ƒ }W n  t yÕ } zt|ƒ}| t|j|d�¡ W Y d}~nd}~ww t||j|d� ~~t�r:|du rñt|
dd�}tƒ }	q/tƒ |	 tk�r9t|
ƒ}tƒ }	|| tk �rq/t|
dd�}tƒ }	|| tk �rq/tj d¡ | |
¡ |�
 	 W d  ƒ dS 1 �s4w   Y  n|	du �sGtƒ |	 tk�rNt  !¡  tƒ }	q0)an  Evaluates calls from call_queue and places the results in result_queue.

    This worker is run in a separate process.

    Args:
        call_queue: A ctx.Queue of _CallItems that will be read and
            evaluated by the worker.
        result_queue: A ctx.Queue of _ResultItems that will written
            to by the worker.
        initializer: A callable initializer, or None
        initargs: A tuple of args for the initializer
        process_management_lock: A ctx.Lock avoiding worker timeout while some
            workers are being spawned.
        timeout: maximum time to wait for a new item in the call_queue. If that
            time is expired, the worker will shutdown.
        worker_exit_lock: Lock to avoid flagging the executor as broken on
            workers timeout.
        current_depth: Nested parallelism level, to avoid infinite spawning.
    NzException in initializer:T)r_   zWorker started with timeout=%s)ÚblockÚtimeoutz Shutting down worker on sentinelz)Shutting down worker after timeout %0.3fsF©r    z+Could not acquire processes_management_lockr   r›   )rp   )r   z*Memory leak detected: shutting down worker)"r�   r   ÚLOGGERÚcriticalÚ_CURRENT_DEPTHÚosÚgetpidr$   rN   rO   ÚgetÚinfor   ÚEmptyÚacquireÚreleaser`   Ú
format_excrœ   rW   Úprintr1   Úexitr[   rm   rn   rŸ   Ú_USE_PSUTILr!   r   Ú_MEMORY_LEAK_CHECK_DELAYÚ_MAX_MEMORY_LEAK_SIZEr   r   )Ú
call_queuerž   ÚinitializerÚinitargsÚprocesses_management_lockr¡   Úworker_exit_lockÚcurrent_depthÚ_process_reference_sizeÚ_last_memory_leak_checkr   Ú	call_itemÚprevious_tbÚrrˆ   rc   Ú	mem_usager   r   r    Ú_process_worker^  sž   ü€ÿýýø	
 ÿ
 €þ
"ÿ€
ÿ´r¿   c                 C   s|   	 |  ¡ rdS z|jdd�}W n tjy   Y dS w | | }|j ¡ r9||g7 }|jt||j|j	|j
ƒdd� n| |= q q)aA  Fills call_queue with _WorkItems from pending_work_items.

    This function never blocks.

    Args:
        pending_work_items: A dict mapping work ids to _WorkItems e.g.
            {5: <_WorkItem...>, 6: <_WorkItem...>, ...}
        work_ids: A queue.Queue of work ids e.g. Queue([5, 6, ...]). Work ids
            are consumed and the corresponding _WorkItems from
            pending_work_items are transformed into _CallItems and put in
            call_queue.
        call_queue: A ctx.Queue that will be filled with _CallItems
            derived from _WorkItems.
    TNFr¢   )Úfullr¨   r   rª   rh   Úset_running_or_notify_cancelrœ   rr   ri   rj   rk   )rz   r{   Úwork_idsr³   rn   r‹   r   r   r    Ú_add_call_item_to_queueÔ  s*   ÿ

ýüírÃ   c
              
      sØ  d‰‡‡fdd„}
‡ ‡‡‡fdd„}|j }|j }||g}	 t|||ˆ ƒ dd„ tˆ ¡ ƒD ƒ}t|| ƒ}d	dtf}||v r�z| ¡ }d}t|tƒrPd
|j	t
f}W n7 ty€ } z#t|ddƒ}|du rjt ¡ \}}}dt t|ƒ||¡t
f}W Y d}~nd}~ww ||v r‰d}d}| ¡  |du�r|\}}}t|tƒrªtjdkrª|d tˆƒ¡7 }||ƒ}|dur¿t|td d |¡¡ƒƒ}ˆ |¡ | ¡ D ]\}}|j |¡ ~qÈ| ¡  ˆrüˆ ¡ \}}tj  d |j!¡¡ zt"|ƒ W n	 t#yù   Y nw ˆsÚ|ƒ  dS t|t$ƒ�rbˆ� ˆ %|d¡}W d  ƒ n	1 �sw   Y  |du�r/|j& '¡  | ¡  ~t(|ƒ}t(|ƒ}|| dk�sE|t(ˆƒk�ra| ƒ ‰ˆdu�rat(ˆƒˆj)k �rat* +dt,¡ ˆ -¡  d‰n,|du�rŽ| %|j.d¡}|du�r�|j/�r|j |j/¡ n|j 0|j1¡ ~| 2|j.¡ ~| ƒ ‰|
ƒ �rãˆj3� dˆ_4W d  ƒ n	1 �s§w   Y  ˆj5�rÚ|�rÅ| ¡ \}}|j t6dƒ¡ ~|�s³ˆ�rÕˆ ¡ \}}t"|ƒ ˆ�sÈ|ƒ  dS |�sâ|ƒ  dS nˆj7�rédS d‰q)aê  Manages the communication between this process and the worker processes.

    This function is run in a local thread.

    Args:
        executor_reference: A weakref.ref to the ProcessPoolExecutor that owns
            this thread. Used to determine if the ProcessPoolExecutor has been
            garbage collected and that this function can exit.
        executor_flags: A ExecutorFlags holding internal states of the
            ProcessPoolExecutor. It permits to know if the executor is broken
            even the object has been gc.
        process: A list of the ctx.Process instances used as
            workers.
        pending_work_items: A dict mapping work ids to _WorkItems e.g.
            {5: <_WorkItem...>, 6: <_WorkItem...>, ...}
        work_ids_queue: A queue.Queue of work ids e.g. Queue([5, 6, ...]).
        call_queue: A ctx.Queue that will be filled with _CallItems
            derived from _WorkItems for processing by the process workers.
        result_queue: A ctx.SimpleQueue of _ResultItems generated by the
            process workers.
        thread_wakeup: A _ThreadWakeup to allow waking up the
            queue_manager_thread from the main Thread and avoid deadlocks
            caused by permanently locked queues.
    Nc                      s   t pˆ d u s	ˆjoˆj S r   )rJ   r=   r>   r   )ÚexecutorÚexecutor_flagsr   r    Úis_shutting_down   s   þz2_queue_management_worker.<locals>.is_shutting_downc               	      sX  t j d¡ ˆ ¡  ˆ� d} tˆ ¡ ƒD ]}|j ¡  | d7 } qW d   ƒ n1 s+w   Y  | }d}||k r�| dkr�t|| ƒD ]}zˆ  	d ¡ |d7 }W qB t
yY   Y  nw ˆ� tdd„ tˆ ¡ ƒD ƒƒ} W d   ƒ n1 stw   Y  ||k r�| dks<t j d¡ ˆ  ¡  t j d¡ ˆrŸˆ ¡ \}}| ¡  ˆs“t j d tˆƒ¡¡ d S )	Nz%queue management thread shutting downr   r   c                 s   s   � | ]}|  ¡ V  qd S r   )rQ   ©r—   Úpr   r   r    Ú	<genexpr>C  s   € 
ÿzI_queue_management_worker.<locals>.shutdown_all_workers.<locals>.<genexpr>zclosing call_queuezjoining processesz>queue management thread clean shutdown of worker processes: {})r$   rN   rO   rE   rK   ÚvaluesÚ_worker_exit_lockr¬   ÚrangeÚ
put_nowaitr   Úsumr+   ÚpopitemrR   rP   )Ún_children_aliverÈ   Ún_children_to_stopÚn_sentinels_sentÚirU   )r³   rÅ   Ú	processesr¶   r   r    Úshutdown_all_workers+  sF   

þþ
ÿ

ÿÿùþ
ÿz6_queue_management_worker.<locals>.shutdown_all_workersTc                 S   s   g | ]}|j ‘qS r   )ÚsentinelrÇ   r   r   r    r™   f  s    z,_queue_management_worker.<locals>.<listcomp>zÞA worker process managed by the executor was unexpectedly terminated. This could be caused by a segmentation fault while calling the function or by an excessive memory usage causing the Operating System to kill the worker.zfA task has failed to un-serialize. Please ensure that the arguments of the function are all picklable.r\   zrA result has failed to un-serialize. Please ensure that the objects returned by the function are always picklable.r,   z% The exit codes of the workers are {}z

'''
{}'''r]   zterminate process {}r   z‚A worker stopped while some jobs were given to the executor. This can be caused by a too short worker timeout or by a memory leak.z9The Executor was shutdown before this job could complete.)8r&   rÃ   rK   rÊ   r	   ÚTerminatedWorkerErrorÚrecvr€   rW   rY   r   r�   r^   r1   r_   r`   ra   rb   r8   Ú
issubclassr2   rP   r   r
   rR   rG   rM   rh   r†   rÏ   r$   rN   rO   Únamer   ÚProcessLookupErrorÚintr„   rË   r¬   ÚlenÚ_max_workersÚwarningsÚwarnÚUserWarningÚ_adjust_process_countrn   ro   Ú
set_resultrp   r…   rB   r=   r?   ÚShutdownExecutorErrorr>   )Úexecutor_referencerÅ   rÔ   rz   r{   Úwork_ids_queuer³   rž   rT   r¶   rÆ   rÕ   Úresult_readerÚwakeup_readerÚreadersÚworker_sentinelsÚreadyr>   Úresult_itemrˆ   rY   rU   ÚmsgÚcause_tbÚexc_typeÚbpern   r‹   rÈ   Ú	n_pendingÚ	n_runningr   )r³   rÄ   rÅ   rÔ   r¶   r    Ú_queue_management_workerü  s   "-ý	ü
ý€ü€ü	



ÿÿ
ÿûÿ


ý€

ÿÿûþþ �æró   c               	   C   sd   t rtrttƒ‚da zt d¡} W n ttfy   Y d S w | dkr$d S | dkr*d S d|  attƒ‚)NTÚSC_SEM_NSEMS_MAXéÿÿÿÿé   z@system provides too few semaphores (%d available, 256 necessary))Ú_system_limits_checkedÚ_system_limitedÚNotImplementedErrorr¦   ÚsysconfÚAttributeErrorÚ
ValueError)Ú	nsems_maxr   r   r    Ú_check_system_limitsü  s"   þÿrþ   c                 c   s*   � | D ]}|  ¡  |r| ¡ V  |sqdS )z½
    Specialized implementation of itertools.chain.from_iterable.
    Each item in *iterable* should be a list.  This function is
    careful not to keep references to yielded objects.
    N)Úreverser„   )ÚiterableÚelementr   r   r    Ú_chain_from_iterable_of_lists  s   €
ÿ€þr  c                 C   sF   |   ¡ dkrtdkrtdƒ‚dtk rtd tkr!td t¡ƒ‚d S d S )NÚforkr   z–Could not spawn extra nested processes at depth superior to MAX_DEPTH=1. It is not possible to increase this limit when using the 'fork' start method.r   z§Could not spawn extra nested processes at depth superior to MAX_DEPTH={}. If this is intendend, you can change this limit with the LOKY_MAX_DEPTH environment variable.)Úget_start_methodr¥   ÚLokyRecursionErrorÚ	MAX_DEPTHrP   )Úcontextr   r   r    Ú_check_max_depth   s   ÿýÿr  c                   @   ó   e Zd ZdZdS )r  zLRaised when a process try to spawn too many levels of nested processes.
    N©r9   r:   r;   rI   r   r   r   r    r  0  ó    r  c                   @   r	  )r   a2  
    Raised when the executor is broken while a future was in the running state.
    The cause can an error raised when unpickling the task in the worker
    process or when unpickling the result value in the parent process. It can
    also be caused by a worker process being terminated unexpectedly.
    Nr
  r   r   r   r    r   5  r  r   c                   @   r	  )r×   zy
    Raised when a process in a ProcessPoolExecutor terminated abruptly
    while a future was in the running state.
    Nr
  r   r   r   r    r×   >  r  r×   c                   @   r	  )rä   zo
    Raised when a ProcessPoolExecutor is shutdown while a future was in the
    running or pending state.
    Nr
  r   r   r   r    rä   J  s    rä   c                       s|   e Zd ZdZ			ddd„Zddd„Zdd„ Zd	d
„ Zdd„ Zdd„ Z	e
jj	je	_‡ fdd„Zddd„Ze
jjje_‡  ZS )ÚProcessPoolExecutorNr   c	           	      C   sè   t ƒ  |du rtƒ | _n|dkrtdƒ‚|| _|du rtƒ }|| _|| _|dur0t|ƒs0tdƒ‚|| _	|| _
t| jƒ |du rA|}|| _i | _d| _i | _g | _t ¡ | _| j ¡ | _d| _tƒ | _tƒ | _|  ||¡ tj d¡ dS )aõ  Initializes a new ProcessPoolExecutor instance.

        Args:
            max_workers: int, optional (default: cpu_count())
                The maximum number of processes that can be used to execute the
                given calls. If None or not given then as many worker processes
                will be created as the number of CPUs the current process
                can use.
            job_reducers, result_reducers: dict(type: reducer_func)
                Custom reducer for pickling the jobs and the results from the
                Executor. If only `job_reducers` is provided, `result_reducer`
                will use the same reducers
            timeout: int, optional (default: None)
                Idle workers exit after timeout seconds. If a new job is
                submitted after the timeout, the executor will start enough
                new Python processes to make sure the pool of workers is full.
            context: A multiprocessing context to launch the workers. This
                object should provide SimpleQueue, Queue and Process.
            initializer: An callable used to initialize worker processes.
            initargs: A tuple of arguments to pass to the initializer.
            env: A dict of environment variable to overwrite in the child
                process. The environment variables are set before any module is
                loaded. Note that this only works with the loky context and it
                is unreliable under windows with Python < 3.6.
        Nr   z"max_workers must be greater than 0zinitializer must be a callablezProcessPoolExecutor is setup)rþ   r   rÞ   rü   r   Ú_contextÚ_envÚcallableÚ	TypeErrorÚ_initializerÚ	_initargsr  Ú_timeoutÚ
_processesÚ_queue_countÚ_pending_work_itemsÚ_running_work_itemsr   r   Ú	_work_idsrA   Ú_processes_management_lockÚ_queue_management_threadr"   Ú_queue_management_thread_wakeupr<   Ú_flagsÚ_setup_queuesr$   rN   rO   )	r)   Úmax_workersÚjob_reducersÚresult_reducersr¡   r  r´   rµ   Úenvr   r   r    r*   V  s:   


zProcessPoolExecutor.__init__c                 C   sP   |d u rd| j  t }t|| j| j| j|| jd�| _d| j_t	|| jd�| _
d S )Nr   )r}   rz   r{   rT   rx   ry   Trw   )rÞ   ÚEXTRA_QUEUED_CALLSrv   r  r  r  r  Ú_call_queueÚ_ignore_epiper   Ú_result_queue)r)   r  r   Ú
queue_sizer   r   r    r  §  s   üÿz!ProcessPoolExecutor._setup_queuesc                 C   s¨   | j d u rPtj d¡ | jfdd„}tjtt 	| |¡| j
| j| j| j| j| j| j| j| jf
dd�| _ d| j _| j  ¡  | jt| j < td u rRtjjd tdd�ad S d S d S )	Nz%_start_queue_management_thread calledc                 S   s   t j d¡ | ¡  d S )Nz?Executor collected: triggering callback for QueueManager wakeup)r$   rN   rO   r5   )rU   rT   r   r   r    Ú
weakref_cbÁ  s   zFProcessPoolExecutor._start_queue_management_thread.<locals>.weakref_cbÚQueueManagerThread)Útargetrj   rÚ   Té   )Úexitpriority)r  r$   rN   rO   r  r@   ÚThreadró   ÚweakrefÚrefr  r  r  r  r  r#  r%  r  ÚdaemonÚstartrL   Úprocess_pool_executor_at_exitÚFinalizerV   )r)   r'  r   r   r    Ú_start_queue_management_threadº  s:   

ÿ
÷
ô
ÿ
ÿÙ#z2ProcessPoolExecutor._start_queue_management_threadc              
   C   s¾   t t| jƒ| jƒD ]I}| j d¡}| j| j| j| j	| j
| j|td f}| ¡  z| jjt|| jd�}W n tyD   | jjt|d�}Y nw ||_| ¡  || j|j< q	tj d | j¡¡ d S )Nr   )r)  rj   r!  )r)  rj   zAdjust process count : {})rÌ   rÝ   r  rÞ   r  ÚBoundedSemaphorer#  r%  r  r  r  r  r¥   r«   r   r¿   r  r  rË   r0  r   r$   rN   rO   rP   )r)   rU   r·   rj   rÈ   r   r   r    râ   å  s$   þ

ÿÿz)ProcessPoolExecutor._adjust_process_countc                 C   sL   | j � t| jƒ| jkr|  ¡  |  ¡  W d  ƒ dS 1 sw   Y  dS )z>ensures all workers and management thread are running
        N)r  rÝ   r  rÞ   râ   r3  r(   r   r   r    Ú_ensure_executor_runningø  s
   
"ýz,ProcessPoolExecutor._ensure_executor_runningc                 O   s°   | j j�J | j jd ur| j j‚| j jrtdƒ‚trtdƒ‚t ¡ }t	||||ƒ}|| j
| j< | j | j¡ |  jd7  _| j ¡  |  ¡  |W  d   ƒ S 1 sQw   Y  d S )Nz*cannot schedule new futures after shutdownz6cannot schedule new futures after interpreter shutdownr   )r  rB   r>   r=   rä   rJ   rƒ   r   ÚFuturerf   r  r  r  rœ   r  r5   r5  )r)   ri   rj   rk   ÚfÚwr   r   r    Úsubmit   s$   
ÿ
$ézProcessPoolExecutor.submitc                    sX   |  dd¡}|  dd¡}|dk rtdƒ‚tt| ƒjtt|ƒt|g|¢R Ž |d�}t|ƒS )az  Returns an iterator equivalent to map(fn, iter).

        Args:
            fn: A callable that will take as many arguments as there are
                passed iterables.
            timeout: The maximum number of seconds to wait. If None, then there
                is no limit on the wait time.
            chunksize: If greater than one, the iterables will be chopped into
                chunks of size chunksize and submitted to the process pool.
                If set to one, the items in the list will be sent one at a
                time.

        Returns:
            An iterator equivalent to: map(func, *iterables) but the calls may
            be evaluated out-of-order.

        Raises:
            TimeoutError: If the entire result iterator could not be generated
                before the given timeout.
            Exception: If fn(*args) raises for any values.
        r¡   Nr’   r   zchunksize must be >= 1.)r¡   )	r¨   rü   r|   r  Úmapr   rš   r–   r  )r)   ri   r“   rk   r¡   r’   Úresultsr~   r   r    r:    s   
þzProcessPoolExecutor.mapTFc                 C   sÌ   t j d|  ¡ | j |¡ | j}| j}|r8d | _|rd | _|d ur2z| ¡  W n	 ty1   Y nw |r8| 	¡  | j
}|rJd | _
| ¡  |rJ| ¡  d | _d | _|rdz| ¡  W d S  tyc   Y d S w d S )Nzshutting down executor %s)r$   rN   rO   r  rE   r  r  r5   ÚOSErrorrR   r#  r+   Újoin_threadr%  r  )r)   r	   r?   ÚqmtÚqmtwÚcqr   r   r    r=   ;  s>   þþýzProcessPoolExecutor.shutdown)NNNNNNr   Nr   )TF)r9   r:   r;   Ú_at_exitr*   r  r3  râ   r5  r9  r   ÚExecutorrI   r:  r=   rŒ   r   r   r~   r    r  R  s    
þ
Q+
 #r  rH   rq   )\rI   Ú
__author__r¦   r   r1   r�   r-  rß   r�   r`   r@   r   Úmultiprocessingr$   Ú	functoolsr   Úpickler   r]   r   Úbackendr   Úbackend.compatr   r	   r
   Úbackend.contextr   Úbackend.queuesr   r   r   Úbackend.reductionr   r   Úbackend.utilsr   r   Úconcurrent.futures.processr   Ú_BPPExceptionÚImportErrorrƒ   r3   r<  rÛ   ÚWeakKeyDictionaryrL   rJ   rÜ   Úenvironr¨   r  r¥   r±   r²   Úpsutilr   r°   r!   r"   Úobjectr<   rV   r1  r"  Ú	ExceptionrW   r�   r[   rd   rf   rm   rr   rv   r–   rš   rŸ   r¿   rÃ   ró   r÷   rø   rþ   r  r  r  r×   ÚBrokenExecutorrä   rB  r  r   r   r   r    Ú<module>   s”   +ÿÿ
$

v( }		