Ë
    ÷Q(ho*  ã                   ó¸   — 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j                  «       Zd adadad„ Z	 	 	 	 	 	 	 	 	 	 dd„Z G d	„ d
e«      Zy)é    Né   )ÚProcessPoolExecutorÚEXTRA_QUEUED_CALLS)Ú	cpu_count)Úget_contextÚget_reusable_executorc                  ó\   — t         5  t        } t        dz  a| cddd«       S # 1 sw Y   yxY w)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_ids    úe/var/www/skyplay_api_hub/venv/lib/python3.12/site-packages/joblib/externals/loky/reusable_executor.pyÚ_get_next_executor_idr      s/   € ô 
ñ Ü'ˆÜ˜QÑÐØ÷÷ ò ús   ‡"¢+c
                 óJ   — t         j                  | |||||||||	¬«
      \  }
}|
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Ú_s               r   r   r   %   sD   € ôh )×>Ñ>ØØØØ!ØØ!Ø'ØØØð ?ó �L€Iˆqð Ðó    c                   óx   ‡ — e Zd Z	 	 	 	 	 	 	 	 	 dˆ fd„	Ze	 	 	 	 	 	 	 	 	 	 dd„«       Zˆ fd„Zd„ Zd„ Zˆ fd„Z	ˆ xZ
S )	r   c           
      óP   •— t         ‰| �  |||||||	|
¬«       || _        || _        y )N)r   r   r   r   r   r   r   r   )ÚsuperÚ__init__r   Ú_submit_resize_lock)ÚselfÚsubmit_resize_lockr   r   r   r   r   r   r   r   r   Ú	__class__s              €r   r!   z_ReusablePoolExecutor.__init__i   sA   ø€ ô 	‰ÑØ#ØØØ%Ø+Ø#ØØð 	ô 		
ð 'ˆÔØ#5ˆÕ r   c           
      ó<  — t         5  t        }|€|du r|�|j                  }nt        «       }n|dk  rt	        d|› d�«      ‚t        |t        «      rt        |«      }|�|j                  «       dk(  rt	        d«      ‚t        ||||||	|
¬«      }|€Ed}t        j                  j                  d	|› d�«       t        «       }|a | t         f||d
œ|¤Žxa}�n-|dk(  r	|t        k(  }|j                  j                   s'|j                  j"                  s|r|j$                  |k  r¢|j                  j                   rd}n-|j                  j"                  rd}n|j$                  |k  rd}nd}t        j                  j                  d|› d|› d�«       |j#                  d|¬«       d xax}a | j&                  dd|i|¤Žcd d d «       S t        j                  j                  d|j                  › d�«       d}|j)                  |«       d d d «       ||fS # 1 sw Y   fS xY w)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   ÚautoÚbrokenÚshutdownzqueue size is too smallz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+   Ú
queue_sizer   Ú_resize)Úclsr   r   r   r   r   r   r   r   r   r   ÚexecutorÚkwargsÚ	is_reusedr   Úreasons                   r   r   z+_ReusablePoolExecutor.get_reusable_executorƒ   sz  € ô ñ O	2ä ˆHàÐ"Ø˜D‘= XÐ%9Ø"*×"7Ñ"7‘Kä"+£+‘KØ Ò!Ü Ø>¸{¸mÈ1ÐMóð ô ˜'¤3Ô'Ü% gÓ.�ØÐ" w×'?Ñ'?Ó'AÀVÒ'KÜ ØJóð ô ØØØ)Ø /Ø'Ø!ØôˆFð ÐØ!�	Ü—‘—‘Ø9¸+¸ÀaÐHôô 4Ó5�Ø#)Ð Ù'*Ü"ð(à +Ø +ñ(ð ñ	(ð �	šHð ˜F’?Ø"Ô&6Ñ6�Eà—O‘O×*Ò*Ø—‘×/Ò/Ù Ø×*Ñ*¨[Ò8à—‘×-Ò-Ø!)™Ø!Ÿ™×1Ò1Ø!+™Ø!×,Ñ,¨{Ò:ð ";™à!9˜Ü—G‘G—M‘MØCØ&˜-ð (#Ø#) (¨"ð.ôð
 ×%Ñ%¨4¸lÐ%ÔKØ>BÐB�IÐB Ð+;à4˜3×4Ñ4ñ Ø$/ðØ39ñ÷MO	2ñ O	2ôT —G‘G—M‘Mð'Ø'/×'<Ñ'<Ð&=¸Qð@ôð !%�IØ×$Ñ$ [Ô1÷_O	2ðb ˜Ð"Ð"÷cO	2ðb ˜Ð"Ð"ús   ‡F2HÇA HÈHc                 ón   •— | j                   5  t        ‰| �  |g|¢­i |¤Žcd d d «       S # 1 sw Y   y xY w©N)r"   r    Úsubmit)r#   ÚfnÚargsr=   r%   s       €r   rB   z_ReusablePoolExecutor.submitä   s7   ø€ Ø×%Ñ%ñ 	7Ü‘7‘> "Ð6 tÒ6¨vÑ6÷	7÷ 	7ò 	7ús   Ž+«4c                 ó¾  — | j                   5  |€t        d«      ‚|| j                  k(  r
	 d d d «       y | j                  €|| _        	 d d d «       y | j	                  «        | j
                  5  t        | j                  j                  «       «      }t        d„ |D «       «      }|| _        t        ||«      D ]  }| j                  j                  d «       Œ 	 d d d «       t        | j                  «      |kD  rZ| j                  j                  sDt!        j"                  d«       t        | j                  «      |kD  r| j                  j                  sŒD| j%                  «        t        | j                  j                  «       «      }t'        d„ |D «       «      s(t!        j"                  d«       t'        d„ |D «       «      sŒ(d d d «       y # 1 sw Y   ŒñxY w# 1 sw Y   y xY w)Nz&Trying to resize with max_workers=Nonec              3   ó<   K  — | ]  }|j                  «       –— Œ y ­wrA   ©Úis_alive©Ú.0Úps     r   ú	<genexpr>z0_ReusablePoolExecutor._resize.<locals>.<genexpr>ý   s   è ø€ Ò'H¸¨¯
©
¯Ñ'Hùó   ‚çü©ñÒMbP?c              3   ó<   K  — | ]  }|j                  «       –— Œ y ­wrA   rG   rI   s     r   rL   z0_ReusablePoolExecutor._resize.<locals>.<genexpr>  s   è ø€ Ò:¨1˜!Ÿ*™*Ÿ,Ñ:ùrM   )r"   r/   r.   Ú_executor_manager_threadÚ_wait_job_completionÚ_processes_management_lockÚlistÚ
_processesÚvaluesÚsumÚrangeÚ_call_queueÚputÚlenr8   r*   ÚtimeÚsleepÚ_adjust_process_countÚall)r#   r   Ú	processesÚnb_children_aliver   s        r   r:   z_ReusablePoolExecutor._resizeè   s™  € Ø×%Ñ%ñ  	!ØÐ"Ü Ð!IÓJÐJØ × 1Ñ 1Ò1Ø÷	 	!ð  	!ð ×,Ñ,Ð4ð %0�Ô!Ø÷ 	!ð  	!ð ×%Ñ%Ô'ð
 ×0Ñ0ñ /Ü  §¡×!7Ñ!7Ó!9Ó:�	Ü$'Ñ'H¸iÔ'HÓ$HÐ!Ø$/�Ô!Ü˜{Ð,=Ó>ò /�AØ×$Ñ$×(Ñ(¨Õ.ñ/÷	/ô �D—O‘OÓ$ {Ò2¸4¿;¹;×;MÒ;Mä—
‘
˜4Ô ô �D—O‘OÓ$ {Ò2¸4¿;¹;×;MÓ;Mð ×&Ñ&Ô(Ü˜TŸ_™_×3Ñ3Ó5Ó6ˆIÜÑ:°	Ô:Ô:Ü—
‘
˜4Ô ô Ñ:°	Ô:Õ:÷? 	!ð  	!÷$/ð /ú÷% 	!ð  	!ús7   �G´GÁGÁ-A)GÃA9GÅA,GÇG	ÇGÇGc                 ó  — | j                   rGt        j                  dt        «       t        j
                  j                  d| j                  › d�«       | j                   r#t        j                  d«       | j                   rŒ"yy)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 resizingrN   N)
Ú_pending_work_itemsÚwarningsÚwarnÚUserWarningr4   r5   r6   r   r[   r\   )r#   s    r   rQ   z*_ReusablePoolExecutor._wait_job_completion  sm   € ð ×#Ò#Ü�M‰Mð?äôô
 �G‰G�M‰MØ˜D×,Ñ,Ð-ð ."ð "ôð
 ×&Ò&Ü�J‰J�tÔð ×&Õ&r   c                 óœ   •— t        t        «       | j                  «      }d|z  t        z   | _        t
        ‰| �  ||| j                  ¬«       y )Né   )r9   )Úmaxr   r.   r   r9   r    Ú_setup_queues)r#   r   r   Úmin_queue_sizer%   s       €r   ri   z#_ReusablePoolExecutor._setup_queues  sH   ø€ ô œY›[¨$×*;Ñ*;Ó<ˆØ˜nÑ,Ô/AÑAˆŒÜ‰ÑØ˜/°d·o±oð 	õ 	
r   )	NNNr   NNNr-   N©
NNé
   Fr)   NNNr-   N)Ú__name__Ú
__module__Ú__qualname__r!   Úclassmethodr   rB   r:   rQ   ri   Ú__classcell__)r%   s   @r   r   r   h   sv   ø„ ð ØØØØØØØØõ6ð4 ð ØØØØØØØØØò^#ó ð^#ô@7ò!!òF÷"

ð 

r   r   rk   )r[   rc   Ú	threadingÚmultiprocessingr4   Úprocess_executorr   r   Úbackend.contextr   Úbackendr   Ú__all__ÚRLockr
   r   r   r7   r   r   r   r-   r   r   ú<module>ry      s€   ðó Û Û Û ç EÝ &Ý  à"Ð
#€ð !�—‘Ó"€ØÐ Ø€	ØÐ ò
ð ØØØØ
ØØØØØó@ôF~
Ð/õ ~
r   