o
    éT•jC"  ã                   @   sÊ   d Z dZddlm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	e 
¡ Zdae ¡ Zdd„ Ze e¡ ee	dƒrEe	jejejejd	� G d
d„ deƒZdd„ ZG dd„ dejƒZG dd„ dejƒZdS )zImplements ThreadPoolExecutor.z"Brian Quinlan (brian@sweetapp.com)é    )Ú_baseNFc                  C   sf   t �
 daW d   ƒ n1 sw   Y  tt ¡ ƒ} | D ]	\}}| d ¡ q| D ]\}}| ¡  q(d S ©NT)Ú_global_shutdown_lockÚ	_shutdownÚlistÚ_threads_queuesÚitemsÚputÚjoin)r   ÚtÚq© r   ú0/usr/lib/python3.10/concurrent/futures/thread.pyÚ_python_exit   s   ÿ
ÿr   Úregister_at_fork)ÚbeforeÚafter_in_childÚafter_in_parentc                   @   s&   e Zd Zdd„ Zdd„ ZeejƒZdS )Ú	_WorkItemc                 C   s   || _ || _|| _|| _d S ©N)ÚfutureÚfnÚargsÚkwargs)Úselfr   r   r   r   r   r   r   Ú__init__/   s   
z_WorkItem.__init__c              
   C   sn   | j  ¡ sd S z| j| ji | j¤Ž}W n ty. } z| j  |¡ d } W Y d }~d S d }~ww | j  |¡ d S r   )r   Úset_running_or_notify_cancelr   r   r   ÚBaseExceptionÚset_exceptionÚ
set_result)r   ÚresultÚexcr   r   r   Úrun5   s   
€ýz_WorkItem.runN)	Ú__name__Ú
__module__Ú__qualname__r   r"   ÚclassmethodÚtypesÚGenericAliasÚ__class_getitem__r   r   r   r   r   .   s    r   c                 C   sì   |d ur(z||Ž  W n t y'   tjjddd� | ƒ }|d ur$| ¡  Y d S w z;	 |jdd�}|d urG| ¡  ~| ƒ }|d urE|j ¡  ~q)| ƒ }t	sS|d u sS|j	rb|d urZd|_	| 
d ¡ W d S ~q* t yu   tjjddd� Y d S w )NzException in initializer:T)Úexc_info)ÚblockzException in worker)r   r   ÚLOGGERÚcriticalÚ_initializer_failedÚgetr"   Ú_idle_semaphoreÚreleaser   r	   )Úexecutor_referenceÚ
work_queueÚinitializerÚinitargsÚexecutorÚ	work_itemr   r   r   Ú_workerE   s@   û

åÿr8   c                   @   s   e Zd ZdZdS )ÚBrokenThreadPoolzR
    Raised when a worker thread in a ThreadPoolExecutor failed initializing.
    N)r#   r$   r%   Ú__doc__r   r   r   r   r9   p   s    r9   c                   @   sd   e Zd Ze ¡ jZ		ddd„Zdd„ Ze	j
jje_dd	„ Zd
d„ Zdddœdd„Ze	j
jje_dS )ÚThreadPoolExecutorNÚ r   c                 C   s¢   |du rt dt ¡ pdd ƒ}|dkrtdƒ‚|dur#t|ƒs#tdƒ‚|| _t ¡ | _	t
 d¡| _tƒ | _d| _d| _t
 ¡ | _|pGd	|  ¡  | _|| _|| _dS )
a•  Initializes a new ThreadPoolExecutor instance.

        Args:
            max_workers: The maximum number of threads that can be used to
                execute the given calls.
            thread_name_prefix: An optional name prefix to give our threads.
            initializer: A callable used to initialize worker threads.
            initargs: A tuple of arguments to pass to the initializer.
        Né    é   é   r   z"max_workers must be greater than 0zinitializer must be a callableFzThreadPoolExecutor-%d)ÚminÚosÚ	cpu_countÚ
ValueErrorÚcallableÚ	TypeErrorÚ_max_workersÚqueueÚSimpleQueueÚ_work_queueÚ	threadingÚ	Semaphorer0   ÚsetÚ_threadsÚ_brokenr   ÚLockÚ_shutdown_lockÚ_counterÚ_thread_name_prefixÚ_initializerÚ	_initargs)r   Úmax_workersÚthread_name_prefixr4   r5   r   r   r   r   {   s$   


ÿ
zThreadPoolExecutor.__init__c             	   O   s¶   | j �N t�; | jrt| jƒ‚| jrtdƒ‚trtdƒ‚t ¡ }t||||ƒ}| j	 
|¡ |  ¡  |W  d   ƒ W  d   ƒ S 1 sDw   Y  W d   ƒ d S 1 sTw   Y  d S )Nz*cannot schedule new futures after shutdownz6cannot schedule new futures after interpreter shutdown)rP   r   rN   r9   r   ÚRuntimeErrorr   ÚFuturer   rI   r	   Ú_adjust_thread_count)r   r   r   r   ÚfÚwr   r   r   Úsubmit¡   s   
RñzThreadPoolExecutor.submitc                 C   s’   | j jdd�r	d S | jfdd„}t| jƒ}|| jk rGd| jp| |f }tj|t	t
 | |¡| j| j| jfd�}| ¡  | j |¡ | jt|< d S d S )Nr   )Útimeoutc                 S   s   |  d ¡ d S r   )r	   )Ú_r   r   r   r   Ú
weakref_cb»   s   z;ThreadPoolExecutor._adjust_thread_count.<locals>.weakref_cbz%s_%d)ÚnameÚtargetr   )r0   ÚacquirerI   ÚlenrM   rF   rR   rJ   ÚThreadr8   ÚweakrefÚrefrS   rT   ÚstartÚaddr   )r   r_   Únum_threadsÚthread_namer   r   r   r   rY   ´   s&   


ÿ
ýÿöz'ThreadPoolExecutor._adjust_thread_countc              	   C   st   | j �- d| _	 z| j ¡ }W n
 tjy   Y nw |d ur'|j t| jƒ¡ qW d   ƒ d S 1 s3w   Y  d S )NzBA thread initializer failed, the thread pool is not usable anymore)	rP   rN   rI   Ú
get_nowaitrG   ÚEmptyr   r   r9   )r   r7   r   r   r   r.   Ë   s   ÿú"øz&ThreadPoolExecutor._initializer_failedTF)Úcancel_futuresc             	   C   s–   | j �0 d| _|r&	 z| j ¡ }W n
 tjy   Y nw |d ur%|j ¡  q
| j d ¡ W d   ƒ n1 s6w   Y  |rG| j	D ]}| 
¡  q@d S d S r   )rP   r   rI   rk   rG   rl   r   Úcancelr	   rM   r
   )r   Úwaitrm   r7   r   r   r   r   ÚshutdownØ   s&   ÿ
ú
ñ

þzThreadPoolExecutor.shutdown)Nr<   Nr   )T)r#   r$   r%   Ú	itertoolsÚcountÚ__next__rQ   r   r\   r   ÚExecutorr:   rY   r.   rp   r   r   r   r   r;   v   s    

ÿ&r;   )r:   Ú
__author__Úconcurrent.futuresr   rq   rG   rJ   r'   re   rA   ÚWeakKeyDictionaryr   r   rO   r   r   Ú_register_atexitÚhasattrr   rb   Ú_at_fork_reinitr1   Úobjectr   r8   ÚBrokenExecutorr9   rt   r;   r   r   r   r   Ú<module>   s.   

þ+