o
    í6Wj¼  ã                   @  sÖ  U d dl mZ d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mZmZ ddlmZmZ dd	lmZ dd
lmZ ddlmZ ejdkrSd dl
mZmZ nd dlmZmZ ejdkrud dlmZmZ d=dd„ZG dd„ dƒZnDejdkr²d dl Z d dl!Z!dZ"de#d< d Z$de#d< dZ%de#d< e%e"fZ&de#d < e$e"fZ'de#d!< e(d"d#d$ƒZ)G d%d„ dƒZnG d&d„ dƒZd'Z*de#d(< d)Z+ed*ƒZ,ed+ƒZ-eee  d,ƒZ.ee d-ƒZ/d>d1d2„Z0dd3œd?d8d9„Z1d@d;d<„Z2dS )Aé    )Úannotations)Úrun_syncÚ#current_default_interpreter_limiterN)Údeque)ÚCallable)ÚAnyÚFinalÚTypeVaré   )Úcurrent_timeÚ	to_thread)ÚBrokenWorkerInterpreter)ÚCapacityLimiter)ÚRunVar)é   é   )ÚTypeVarTupleÚUnpack)r   é   )ÚExecutionFailedÚcreateÚfuncúCallable[..., Any]Úargsútuple[Any, ...]Úreturnútuple[Any, bool]c              
   C  s@   z	| |Ž }W |dfS  t y } z
|dfW  Y d }~S d }~ww )NTF)ÚBaseException)r   r   ÚretvalÚexc© r    úc/home/esfera/Documents/content_generation/venv/lib/python3.10/site-packages/anyio/to_interpreter.pyÚ_interp_call   s   
ý€ÿr"   c                   @  ó8   e Zd ZU dZded< ddd„Zddd	„Zddd„ZdS )Ú_Workerr   ÚfloatÚ	last_usedr   ÚNonec                 C  s   t ƒ | _d S ©N)r   Ú_interpreter©Úselfr    r    r!   Ú__init__)   s   ú_Worker.__init__c                 C  s   | j  ¡  d S r(   )r)   Úcloser*   r    r    r!   Údestroy,   s   ú_Worker.destroyr   úCallable[..., T_Retval]r   r   ÚT_Retvalc              
   C  sJ   z| j  t||¡\}}W n ty } zt|jƒ|‚d }~ww |r#|‚|S r(   )r)   Úcallr"   r   r   Úexcinfo)r+   r   r   ÚresÚis_exceptionr   r    r    r!   r3   /   s   €ÿú_Worker.callN©r   r'   ©r   r1   r   r   r   r2   ©Ú__name__Ú
__module__Ú__qualname__r&   Ú__annotations__r,   r/   r3   r    r    r    r!   r$   &   s
   
 

r$   )r   é   é   r   ÚUNBOUNDÚFMT_UNPICKLEDÚFMT_PICKLEDÚQUEUE_PICKLE_ARGSÚQUEUE_UNPICKLE_ARGSa_  
import _interpqueues
from _interpreters import NotShareableError
from pickle import loads, dumps, HIGHEST_PROTOCOL

QUEUE_PICKLE_ARGS = (1, 2)
QUEUE_UNPICKLE_ARGS = (0, 2)

item = _interpqueues.get(queue_id)[0]
try:
    func, args = loads(item)
    retval = func(*args)
except BaseException as exc:
    is_exception = True
    retval = exc
else:
    is_exception = False

try:
    _interpqueues.put(queue_id, (retval, is_exception), *QUEUE_UNPICKLE_ARGS)
except NotShareableError:
    retval = dumps(retval, HIGHEST_PROTOCOL)
    _interpqueues.put(queue_id, (retval, is_exception), *QUEUE_PICKLE_ARGS)
    z<string>Úexecc                   @  r#   )r$   r   r%   r&   r   r'   c                 C  s6   t  ¡ | _tjdgt¢R Ž | _t  | jd| ji¡ d S )Nr
   Úqueue_id)Ú_interpretersr   Ú_interpreter_idÚ_interpqueuesrE   Ú	_queue_idÚset___main___attrsr*   r    r    r!   r,   g   s
   
ÿr-   c                 C  s   t  | j¡ t | j¡ d S r(   )rJ   r/   rK   rH   rI   r*   r    r    r!   r/   n   s   r0   r   r1   r   r   r2   c           	      C  sˆ   dd l }| ||f|j¡}tj| j|gt¢R Ž  t | j	t
¡}|r%t|ƒ‚t | j¡}|d d… \\}}}|tkr>| |¡}|rB|‚|S )Nr   r@   )ÚpickleÚdumpsÚHIGHEST_PROTOCOLrJ   ÚputrK   rD   rH   rF   rI   Ú	_run_funcr   ÚgetrC   Úloads)	r+   r   r   rM   ÚitemÚexc_infor5   r6   Úfmtr    r    r!   r3   r   s   
r7   Nr8   r9   r:   r    r    r    r!   r$   d   s
   
 

c                   @  s8   e Zd ZU dZded< ddd„Zddd„Zddd„ZdS )r$   r   r%   r&   r   r'   c                 C  s   t dƒ‚)Nz,subinterpreters require at least Python 3.13)ÚRuntimeErrorr*   r    r    r!   r,   �   s   r-   r   r1   r   r   r2   c                 C  s   t ‚r(   )ÚNotImplementedError)r+   r   r   r    r    r!   r3   �   s   r7   c                 C  s   d S r(   r    r*   r    r    r!   r/   —   s   r0   Nr8   r9   )r;   r<   r=   r&   r>   r,   r3   r/   r    r    r    r!   r$   Š   s
   
 

é   ÚDEFAULT_CPU_COUNTé   r2   ÚPosArgsTÚ_available_workersÚ_default_interpreter_limiterÚworkersúdeque[_Worker]r'   c                 C  s   | D ]}|  ¡  q|  ¡  d S r(   )r/   Úclear)r_   Úworkerr    r    r!   Ú_stop_workers§   s   
rc   ©Úlimiterú&Callable[[Unpack[PosArgsT]], T_Retval]úUnpack[PosArgsT]re   úCapacityLimiter | Nonec             
   Ç  sf  �|du rt ƒ }zt ¡ }W n ty%   tƒ }t |¡ t t|¡ Y nw |4 I dH š z| 	¡ }W n t
y?   tƒ }Y nw W d  ƒI dH  n1 I dH sPw   Y  z5tj|j| ||d�I dH W tƒ }|r�||d j tkrrntj| ¡ j|d�I dH  |shtƒ |_| |¡ S tƒ }|r©||d j tkršntj| ¡ j|d�I dH  |s�tƒ |_| |¡ w )aç  
    Call the given function with the given arguments in a subinterpreter.

    .. warning:: On Python 3.13, the :mod:`concurrent.interpreters` module was not yet
        available, so the code path for that Python version relies on an undocumented,
        private API. As such, it is recommended to not rely on this function for anything
        mission-critical on Python 3.13.

    :param func: a callable
    :param args: the positional arguments for the callable
    :param limiter: capacity limiter to use to limit the total number of subinterpreters
        running (if omitted, the default limiter is used)
    :return: the result of the call
    :raises BrokenWorkerInterpreter: if there's an internal error in a subinterpreter

    Nrd   r   )r   Ú_idle_workersrR   ÚLookupErrorr   ÚsetÚatexitÚregisterrc   ÚpopÚ
IndexErrorr$   r   r   r3   r   r&   ÚMAX_WORKER_IDLE_TIMEÚpopleftr/   Úappend)r   re   r   Úidle_workersrb   Únowr    r    r!   r   ®   sR   €
ý
ÿ€(ýüüøür   r   c                  C  s<   zt  ¡ W S  ty   tt ¡ ptƒ} t  | ¡ |  Y S w )zÉ
    Return the capacity limiter used by default to limit the number of concurrently
    running subinterpreters.

    Defaults to the number of CPU cores.

    :return: a capacity limiter object

    )r^   rR   rj   r   ÚosÚ	cpu_countrZ   rk   rd   r    r    r!   r   ç   s   


ýr   )r   r   r   r   r   r   )r_   r`   r   r'   )r   rf   r   rg   re   rh   r   r2   )r   r   )3Ú
__future__r   Ú__all__rl   ru   ÚsysÚcollectionsr   Úcollections.abcr   Útypingr   r   r	   Ú r   r   Ú_core._exceptionsr   Ú_core._synchronizationr   Úlowlevelr   Úversion_infor   r   Útyping_extensionsÚconcurrent.interpretersr   r   r"   r$   rJ   rH   rA   r>   rB   rC   rD   rE   ÚcompilerQ   rZ   rp   r2   r\   ri   r^   rc   r   r   r    r    r    r!   Ú<module>   sZ    




æ&ÿ

ý9