o
    l¨Êh4$  ã                   @   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dƒ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_executorÚ c                  C   s8   t � t} td7 a| W  d  ƒ S 1 sw   Y  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   úY/var/www/html/env/lib/python3.10/site-packages/joblib/externals/loky/reusable_executor.pyÚ_get_next_executor_id   s
   $ýr   é
   FÚautor   c
              
   C   s°  t �Ì t}
| du r|du r|
dur|
j} ntƒ } n| dkr$td | ¡ƒ‚t|tƒr-t|ƒ}|dur;| 	¡ dkr;tdƒ‚t
|||||||	d�}|
du rftj d | ¡¡ tƒ }|att f| |d	œ|¤Ž a}
n`|d
krn|tk}|
jjsx|
jjsx|s¯|
jjrd}n	|
jjr†d}nd}tj d | |¡¡ |
jd|d� d a }
atdd| i|¤ŽW  d  ƒ S tj d |
j¡¡ |
 | ¡ W d  ƒ |
S W d  ƒ |
S 1 sÑw   Y  |
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
    ``VAR`` 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 with the
    ``loky`` context and it is unreliable on Windows with Python < 3.6.
    NTr   z+max_workers must be greater than 0, got {}.Úforkz4Cannot use reusable executor with the 'fork' context)ÚcontextÚtimeoutÚjob_reducersÚresult_reducersÚinitializerÚinitargsÚenvz&Create a executor with max_workers={}.)Úmax_workersr   r   ÚbrokenÚshutdownzarguments have changedz[Creating a new executor with max_workers={} as the previous instance cannot be reused ({}).)ÚwaitÚkill_workersr   z.Reusing existing executor with max_workers={}.r   )r	   Ú	_executorÚ_max_workersr   Ú
ValueErrorÚformatÚ
isinstanceÚSTRING_TYPEr   Úget_start_methodÚdictÚmpÚutilÚdebugr   Ú_executor_kwargsÚ_ReusablePoolExecutorÚ_flagsr   r   r   Ú_resize)r   r   r   r   Úreuser   r   r   r   r   ÚexecutorÚkwargsr   Úreasonr   r   r   r   (   s„   ,þ
üÿÿþþÿý
ÿÍ6ÿ
È:ä
â:Æ:c                       sN   e Zd Z				d‡ fdd„	Z‡ fdd„Zdd	„ Zd
d„ Z‡ fdd„Z‡  ZS )r*   Nr   r   c              
      s0   t t| ƒj|||||||	|
d� || _|| _d S )N)r   r   r   r   r   r   r   r   )Úsuperr*   Ú__init__r   Ú_submit_resize_lock)ÚselfÚsubmit_resize_lockr   r   r   r   r   r   r   r   r   ©Ú	__class__r   r   r2   ’   s   
ý
z_ReusablePoolExecutor.__init__c                    sH   | j � tt| ƒj|g|¢R i |¤ŽW  d   ƒ S 1 sw   Y  d S ©N)r3   r1   r*   Úsubmit)r4   ÚfnÚargsr/   r6   r   r   r9   �   s   
ÿÿÿ$ÿz_ReusablePoolExecutor.submitc              	   C   st  | j �­ |d u rtdƒ‚|| jkr	 W d   ƒ d S | jd u r+|| _	 W d   ƒ d S |  ¡  | j�) t| j ¡ ƒ}t	dd„ |D ƒƒ}|| _t
||ƒD ]}| j d ¡ qKW d   ƒ n1 s^w   Y  t| jƒ|kr~| jjs~t d¡ t| jƒ|kr~| jjrn|  ¡  t| j ¡ ƒ}tdd„ |D ƒƒs¨t d¡ tdd„ |D ƒƒr’W d   ƒ d S W d   ƒ d S 1 s³w   Y  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   g | ]}|  ¡ ‘qS r   r<   r>   r   r   r   Ú
<listcomp>Á   s    z1_ReusablePoolExecutor._resize.<locals>.<listcomp>)r3   r    r   Ú_queue_management_threadÚ_wait_job_completionÚ_processes_management_lockÚlistÚ
_processesÚvaluesÚsumÚrangeÚ_call_queueÚputÚlenr+   r   ÚtimeÚsleepÚ_adjust_process_countÚall)r4   r   Ú	processesÚnb_children_aliveÚ_r   r   r   r,   ¢   sD   
ü
õÿüÿ
þÿ
ÿâ"âz_ReusablePoolExecutor._resizec                 C   s\   t | jƒdkrt dt¡ tj d | j	¡¡ t | jƒdkr,t
 d¡ t | jƒdksdS dS )z8Wait for the cache to be empty before resizing the pool.r   z\Trying to resize an executor with running jobs: waiting for jobs completion before resizing.z7Executor {} waiting for jobs completion before resizingrB   N)rN   Ú_pending_work_itemsÚwarningsÚwarnÚUserWarningr&   r'   r(   r!   r   rO   rP   )r4   r   r   r   rE   Ä   s   þÿ
ÿz*_ReusablePoolExecutor._wait_job_completionc                    s(   dt ƒ  t }tt| ƒj|||d� d S )Né   )Ú
queue_size)r   r   r1   r*   Ú_setup_queues)r4   r   r   r[   r6   r   r   r\   Ñ   s   

ÿz#_ReusablePoolExecutor._setup_queues)	NNNr   NNNr   N)	Ú__name__Ú
__module__Ú__qualname__r2   r9   r,   rE   r\   Ú__classcell__r   r   r6   r   r*   ‘   s    ý"r*   )
NNr   Fr   NNNr   N)rO   rW   Ú	threadingÚmultiprocessingr&   Úprocess_executorr   r   Úbackend.contextr   Úbackendr   Ú__all__Útyper#   ÚRLockr	   r
   r   r)   r   r   r*   r   r   r   r   Ú<module>   s(   
ýi