U
    ¬º|eA(  ã                
   @   sˆ   d dl Z d dlZd dlZd dlZddlmZmZ ddlm	Z	 ddl
mZ dgZe ¡ Zd adadadd„ Zddd„ZG dd„ deƒZdS )é    Né   )ÚProcessPoolExecutorÚEXTRA_QUEUED_CALLS)Ú	cpu_count)Úget_contextÚget_reusable_executorc               
   C   s,   t � t} td7 a| W  5 Q R £ S Q R X dS )z¯Ensure that each successive executor instance has a unique, monotonic id.

    The purpose of this monotonic id is to help debug and test automated
    instance creation.
    r   N)Ú_executor_lockÚ_next_executor_id)Úexecutor_id© r   úd/var/www/website-v5/atlas_env/lib/python3.8/site-packages/joblib/externals/loky/reusable_executor.pyÚ_get_next_executor_id   s    r   é
   FÚautor   c
                 C   s&   t j| |||||||||	d�
\}
}|
S )a¬  Return the current ReusableExectutor instance.

    Start a new instance if it has not been started already or if the previous
    instance was left in a broken state.

    If the previous instance does not have the requested number of workers, the
    executor is dynamically resized to adjust the number of workers prior to
    returning.

    Reusing a singleton instance spares the overhead of starting new worker
    processes and importing common python packages each time.

    ``max_workers`` controls the maximum number of tasks that can be running in
    parallel in worker processes. By default this is set to the number of
    CPUs on the host.

    Setting ``timeout`` (in seconds) makes idle workers automatically shutdown
    so as to release system resources. New workers are respawn upon submission
    of new tasks so that ``max_workers`` are available to accept the newly
    submitted tasks. Setting ``timeout`` to around 100 times the time required
    to spawn new processes and import packages in them (on the order of 100ms)
    ensures that the overhead of spawning workers is negligible.

    Setting ``kill_workers=True`` makes it possible to forcibly interrupt
    previously spawned jobs to get a new instance of the reusable executor
    with new constructor argument values.

    The ``job_reducers`` and ``result_reducers`` are used to customize the
    pickling of tasks and results send to the executor.

    When provided, the ``initializer`` is run first in newly spawned
    processes with argument ``initargs``.

    The environment variable in the child process are a copy of the values in
    the main process. One can provide a dict ``{ENV: VAL}`` where ``ENV`` and
    ``VAL`` are string literals to overwrite the environment variable ``ENV``
    in the child processes to value ``VAL``. The environment variables are set
    in the children before any module is loaded. This only works with the
    ``loky`` context.
    )
Úmax_workersÚcontextÚtimeoutÚkill_workersÚreuseÚjob_reducersÚresult_reducersÚinitializerÚinitargsÚenv)Ú_ReusablePoolExecutorr   )r   r   r   r   r   r   r   r   r   r   Ú	_executorÚ_r   r   r   r   %   s    4ö
c                       sT   e Zd Zd‡ fdd„	Zedd	d
„ƒZ‡ fdd„Zdd„ Zdd„ Z‡ fdd„Z	‡  Z
S )r   Nr   r   c              
      s,   t ƒ j|||||||	|
d� || _|| _d S )N)r   r   r   r   r   r   r   r   )ÚsuperÚ__init__r
   Ú_submit_resize_lock)ÚselfZsubmit_resize_lockr   r   r   r
   r   r   r   r   r   ©Ú	__class__r   r   r   i   s    ø
z_ReusablePoolExecutor.__init__r   Fr   c              
   C   sª  t ��– t}|d kr4|dkr,|d k	r,|j}qLtƒ }n|dkrLtd|› d�ƒ‚t|tƒr^t|ƒ}|d k	rz| ¡ dkrztdƒ‚t	||||||	|
d�}|d krÖd}t
j d	|› d�¡ tƒ }|a| t f||d
œ|—Ž a}nÂ|dkræ|tk}|jjsü|jjsü|�st|jj�rd}n|jj�rd}nd}t
j d|› d|› d�¡ |jd|d� d  a }a| jf d|i|—ŽW  5 Q R £ S t
j d|j› d�¡ d}| |¡ W 5 Q R X ||fS )NTr   z(max_workers must be greater than 0, got Ú.Úforkz4Cannot use reusable executor with the 'fork' context)r   r   r   r   r   r   r   Fz#Create a executor with max_workers=)r   r
   r   ÚbrokenÚshutdownzarguments have changedz)Creating a new executor with max_workers=z, as the previous instance cannot be reused (z).)Úwaitr   r   z+Reusing existing executor with max_workers=)r   r   Ú_max_workersr   Ú
ValueErrorÚ
isinstanceÚstrr   Úget_start_methodÚdictÚmpÚutilÚdebugr   Ú_executor_kwargsÚ_flagsr%   r&   r   Ú_resize)Úclsr   r   r   r   r   r   r   r   r   r   ÚexecutorÚkwargsZ	is_reusedr
   Úreasonr   r   r   r   ƒ   sŠ    
ÿ
ÿù	
ÿÿýüÿþý

ÿÿÿÿz+_ReusablePoolExecutor.get_reusable_executorc              
      s2   | j �" tƒ j|f|ž|ŽW  5 Q R £ S Q R X d S ©N)r   r   Úsubmit)r    ÚfnÚargsr6   r!   r   r   r9   ß   s    z_ReusablePoolExecutor.submitc              
   C   s  | j �� |d krtdƒ‚n|| jkr4W 5 Q R £ d S | jd krR|| _W 5 Q R £ d S |  ¡  | j�H t| j ¡ ƒ}t	dd„ |D ƒƒ}|| _t
||ƒD ]}| j d ¡ q’W 5 Q R X t| jƒ|krÐ| jjsÐt d¡ q®|  ¡  t| j ¡ ƒ}tdd„ |D ƒƒ�st d¡ qæW 5 Q R X d S )Nz&Trying to resize with max_workers=Nonec                 s   s   | ]}|  ¡ V  qd S r8   ©Úis_alive©Ú.0Úpr   r   r   Ú	<genexpr>ø   s     z0_ReusablePoolExecutor._resize.<locals>.<genexpr>çü©ñÒMbP?c                 s   s   | ]}|  ¡ V  qd S r8   r<   r>   r   r   r   rA     s     )r   r)   r(   Z_executor_manager_threadÚ_wait_job_completionZ_processes_management_lockÚlistÚ
_processesÚvaluesÚsumÚrangeÚ_call_queueÚputÚlenr2   r%   ÚtimeÚsleepÚ_adjust_process_countÚall)r    r   Ú	processesZnb_children_aliver   r   r   r   r3   ã   s0    



ÿÿz_ReusablePoolExecutor._resizec                 C   s>   | j r(t dt¡ tj d| j› d�¡ | j r:t 	d¡ q(dS )z8Wait for the cache to be empty before resizing the pool.z\Trying to resize an executor with running jobs: waiting for jobs completion before resizing.z	Executor z, waiting for jobs completion before resizingrB   N)
Ú_pending_work_itemsÚwarningsÚwarnÚUserWarningr.   r/   r0   r
   rL   rM   )r    r   r   r   rC     s    ýÿz*_ReusablePoolExecutor._wait_job_completionc                    s$   dt ƒ  t }tƒ j|||d� d S )Né   )Ú
queue_size)r   r   r   Ú_setup_queues)r    r   r   rV   r!   r   r   rW     s      ÿz#_ReusablePoolExecutor._setup_queues)	NNNr   NNNr   N)
NNr   Fr   NNNr   N)Ú__name__Ú
__module__Ú__qualname__r   Úclassmethodr   r9   r3   rC   rW   Ú__classcell__r   r   r!   r   r   h   s4            õ          õ[#r   )
NNr   Fr   NNNr   N)rL   rR   Ú	threadingÚmultiprocessingr.   Úprocess_executorr   r   Úbackend.contextr   Úbackendr   Ú__all__ÚRLockr   r	   r   r1   r   r   r   r   r   r   r   Ú<module>   s0             ö
C