o
    í6Wjx&  ã                   @  s’  U d dl m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m	Z	 d dl
mZ d dlmZ d dlmZmZ dd	lmZmZmZ dd
lmZ ddlmZ ddlmZ ddlmZmZ ddlmZm Z m!Z! ddl"m#Z#m$Z$ ddl%m&Z& ej'dkr�d dlm(Z(m)Z) nd dl*m(Z(m)Z) dZ+edƒZ,e(dƒZ-e#dƒZ.de/d< e#dƒZ0de/d< e#dƒZ1de/d< dddœd0d'd(„Z2d1d*d+„Z3d2d-d.„Z4e5d/krÇe4ƒ  dS dS )3é    )Úannotations)Úcurrent_default_process_limiterÚprocess_workerÚrun_syncN)Údeque)ÚCallable)Ú
ModuleType)ÚTypeVarÚcasté   )Úcurrent_timeÚget_async_backendÚget_cancelled_exc_class)ÚBrokenWorkerProcess)Úopen_process)ÚCapacityLimiter)ÚCancelScopeÚ
fail_after)ÚByteReceiveStreamÚByteSendStreamÚProcess)ÚRunVarÚcheckpoint_if_cancelled)ÚBufferedByteReceiveStream)é   é   )ÚTypeVarTupleÚUnpacki,  ÚT_RetvalÚPosArgsTÚ_process_pool_workerszRunVar[set[Process]]Ú_process_pool_idle_workersz$RunVar[deque[tuple[Process, float]]]Ú_default_process_limiterzRunVar[CapacityLimiter]F)ÚcancellableÚlimiterÚfuncú&Callable[[Unpack[PosArgsT]], T_Retval]ÚargsúUnpack[PosArgsT]r#   Úboolr$   úCapacityLimiter | NoneÚreturnc                ‡  sZ  �d‡ ‡‡‡fdd„}t ƒ I dH  tjd| |ftjd	�}z
t ¡ ‰t ¡ }W n tyE   tƒ ‰t	ƒ }t ˆ¡ t |¡ t
ƒ  ˆ¡ Y nw |pJtƒ 4 I dH �šO |r½| ¡ \‰}ˆjdu r¶ttˆjƒ‰tttˆjƒƒ‰ tƒ }g }	|r”||d
 d  tk r~n| ¡ \}
}|
 ¡  ˆ |
¡ |	 |
¡ |sstdd�� |	D ]	}| ¡ I dH  qœW d  ƒ n1 s°w   Y  n•ˆ ˆ¡ |sStjddtg}t |t!j"t!j"d�I dH ‰zTttˆjƒ‰tttˆjƒƒ‰ t#dƒ� ˆ  $d¡I dH }W d  ƒ n1 söw   Y  |dk�rt%d|›�ƒ‚t&tj'd ddƒ}tjdtj(|ftjd	�}||ƒI dH  W n! t%t)ƒ f�y0   ‚  t*�yE } z	ˆ ¡  t%dƒ|‚d}~ww ˆ +ˆ¡ t| d��: z)tt,||ƒI dH ƒW ˆˆv �rj| ˆtƒ f¡ W  d  ƒ W  d  ƒI dH  S ˆˆv �rŠ| ˆtƒ f¡ w w 1 �s�w   Y  W d  ƒI dH  dS 1 I dH �s¦w   Y  dS )a'  
    Call the given function with the given arguments in a worker process.

    If the ``cancellable`` option is enabled and the task waiting for its completion is
    cancelled, the worker process running it will be abruptly terminated using SIGKILL
    (or ``terminateProcess()`` on Windows).

    :param func: a callable
    :param args: positional arguments for the callable
    :param cancellable: ``True`` to allow cancellation of the operation while it's
        running
    :param limiter: capacity limiter to use to limit the total amount of processes
        running (if omitted, the default limiter is used)
    :raises NoEventLoopError: if no supported asynchronous event loop is running in the
        current thread
    :return: an awaitable that yields the return value of the function.

    Úpickled_cmdÚbytesr+   Úobjectc                 “  s  �z/ˆ  | ¡I d H  ˆ  dd¡I d H }| d¡\}}|dvr%td|›�ƒ‚ˆ  t|ƒ¡I d H }W nG tyw } z;ˆ ˆ¡ z"ˆ ¡  t	dd�� ˆ 
¡ I d H  W d   ƒ n1 sYw   Y  W n	 tyh   Y nw t|tƒ ƒrp‚ t|‚d }~ww t |¡}|dkrŠt|tƒsˆJ ‚|‚|S )	Nó   
é2   ó    )ó   RETURNó	   EXCEPTIONú-Worker process returned unexpected response: T©Úshieldr3   )ÚsendÚreceive_untilÚsplitÚRuntimeErrorÚreceive_exactlyÚintÚBaseExceptionÚdiscardÚkillr   ÚacloseÚProcessLookupErrorÚ
isinstancer   r   ÚpickleÚloads)r,   ÚresponseÚstatusÚlengthÚpickled_responseÚexcÚretval©ÚbufferedÚprocessÚstdinÚworkers© ú_/home/esfera/Documents/content_generation/venv/lib/python3.10/site-packages/anyio/to_process.pyÚsend_raw_commandF   s>   €ÿ
ÿ€ÿ€ô
z"run_sync.<locals>.send_raw_commandNÚrun)Úprotocolr   r   Tr5   z-uz-m)rN   Ústdouté   é   ó   READY
r4   Ú__main__Ú__file__Úinitz*Error during worker process initialization)r,   r-   r+   r.   )-r   rC   ÚdumpsÚHIGHEST_PROTOCOLr    Úgetr!   ÚLookupErrorÚsetr   r   Ú#setup_process_pool_exit_at_shutdownr   ÚpopÚ
returncoder
   r   rN   r   r   rU   r   ÚWORKER_MAX_IDLE_TIMEÚpopleftr?   ÚremoveÚappendr   r@   ÚsysÚ
executableÚ__name__r   Ú
subprocessÚPIPEr   Úreceiver   ÚgetattrÚmodulesÚpathr   r=   Úaddr   )r%   r#   r$   r'   rR   ÚrequestÚidle_workersÚ
idle_sinceÚnowÚkilled_processesÚprocess_to_killÚkilled_processÚcommandÚmessageÚmain_module_pathÚpickledrI   rP   rK   rQ   r   -   s¬   €!

û

ÿ

ù	ÿÿ
å
ÿ
ÿ
ÿ
ÿ
þÿþ€þ

û¾
Fÿü0¾r   r   c                  C  s<   zt  ¡ W S  ty   tt ¡ pdƒ} t  | ¡ |  Y S w )z“
    Return the capacity limiter that is used by default to limit the number of worker
    processes.

    :return: a capacity limiter object

    é   )r"   r^   r_   r   ÚosÚ	cpu_countr`   )r$   rP   rP   rQ   r   ¿   s   

ýr   ÚNonec               
   C  s  t j} t j}ttjƒt _ttjdƒt _ttjdƒt _|j d¡ 	 d  }}z
t	 
| j¡^}}W n ty9   Y d S  tyL } z|}W Y d }~nod }~ww |dkrp|\}}z||Ž }W n[ tyo } z|}W Y d }~nLd }~ww |dkr·|\t _}t jd= |r·tj |¡r·ztdƒ}	tj|dd�}
|	j |
¡ |	 t jd< t jd< W n ty¶ } z|}W Y d }~nd }~ww z|d urÆd	}t	 |t	j¡}n	d
}t	 |t	j¡}W n tyí } z|}d	}t	 |t	j¡}W Y d }~nd }~ww |j d|t|ƒf ¡ |j |¡ t|tƒ�r|‚q!)NÚwrX   TrS   r[   rY   Ú__mp_main__)Úrun_namer3   r2   s   %s %d
)rh   rN   rU   Úopenr~   ÚdevnullÚstderrÚbufferÚwriterC   ÚloadÚEOFErrorr=   rp   ro   Úisfiler   ÚrunpyÚrun_pathÚ__dict__Úupdater\   r]   ÚlenrB   Ú
SystemExit)rN   rU   rJ   Ú	exceptionry   r'   rI   r%   r{   ÚmainÚmain_contentrF   r|   rP   rP   rQ   r   Ï   sr   €ÿ€ÿ
ÿ€ÿ€€ýÐr   rY   )
r%   r&   r'   r(   r#   r)   r$   r*   r+   r   )r+   r   )r+   r€   )6Ú
__future__r   Ú__all__r~   rC   rŒ   rk   rh   Úcollectionsr   Úcollections.abcr   Útypesr   Útypingr	   r
   Ú_core._eventloopr   r   r   Ú_core._exceptionsr   Ú_core._subprocessesr   Ú_core._synchronizationr   Ú_core._tasksr   r   Úabcr   r   r   Úlowlevelr   r   Ústreams.bufferedr   Úversion_infor   r   Útyping_extensionsrd   r   r   r    Ú__annotations__r!   r"   r   r   r   rj   rP   rP   rP   rQ   Ú<module>   sN    
ÿü 

=
ÿ