XjpUddlmZdZddlZddlZddlZddlmZddlm Z ddl m Z m Z m Z ddlmZmZdd lmZdd lmZdd lmZej.d k\r dd l mZmZndd lmZmZej.dk\rddlmZmZ d+dZGddZn]ej.dk\rDddl Z ddl!Z!dZ"de#d<dZ$de#d<dZ%de#d<e%e"fZ&de#d<e$e"fZ'de#d<e(dddZ)GddZn GddZd Z*de#d!<d"Z+e d#Z,ed$Z-eeed%Z.eed&Z/d,d'Z0dd( d-d)Z1d.d*Z2y)/) annotations)run_sync#current_default_interpreter_limiterN)deque)Callable)AnyFinalTypeVar) current_time to_thread)BrokenWorkerInterpreter)CapacityLimiter)RunVar) ) TypeVarTupleUnpack)r)ExecutionFailedcreatecJ ||}|dfS#t$r}|dfcYd}~Sd}~wwxYw)NFT) BaseException)funcargsretvalexcs h/mnt/ssd/data/Dropbox/adrian/vault-secondbrain/venv/lib/python3.12/site-packages/anyio/to_interpreter.py _interp_callrs8 !4[F5=  9  s """c@eZdZUdZded<ddZddZ d dZy) _Workerrfloat last_usedc"t|_yN)r _interpreterselfs r__init__z_Worker.__init__)s &D c8|jjyr%)r&closer's rdestroyz_Worker.destroy,s    # # %r*c |jjt||\}}|r||S#t$r}t |j |d}~wwxYwr%)r&callrrrexcinfo)r(rrres is_exceptionrs rr/z _Worker.call/s[  D$($5$5$:$:<t$T!\ J # D-ckk:C Ds$, AA  ANreturnNonerzCallable[..., T_Retval]rtuple[Any, ...]r4T_Retval__name__ __module__ __qualname__r#__annotations__r)r-r/r*rr!r!&s7 5 ) & ) "  r*r!)r r UNBOUND FMT_UNPICKLED FMT_PICKLEDQUEUE_PICKLE_ARGSQUEUE_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) zexecc@eZdZUdZded<ddZddZ d dZy) r!rr"r#ctj|_tjdgt|_tj |jd|j iy)Nr queue_id) _interpretersr_interpreter_id _interpqueuesrE _queue_idset___main___attrsr's rr)z_Worker.__init__gsM#0#7#7#9D *11!J6IJDN  , ,$$z4>>&B r*ctj|jtj|jyr%)rLr-rMrJrKr's rr-z_Worker.destroyns(  ! !$.. 1  ! !$"6"6 7r*cddl}|j||f|j}tj|j |gt tj|jt}|r t|tj|j }|dd\\}}}|tk(r|j|}|r||S)Nrr@)pickledumpsHIGHEST_PROTOCOLrLputrMrDrJrFrK _run_funcrgetrCloads) r(rrrQitemexc_infor1r2fmts rr/z _Worker.callrs <<t f.E.EFD   dnnd G5F G$))$*>*> JH-h77##DNN3C'*2Aw $ S,k!ll3' Jr*Nr3r6r9r>r*rr!r!ds7 5  8 ) "   r*c@eZdZUdZded<ddZ d dZddZy) r!rr"r#ctd)Nz,subinterpreters require at least Python 3.13) RuntimeErrorr's rr)z_Worker.__init__sMN Nr*ctr%)NotImplementedError)r(rrs rr/z _Worker.calls & %r*cyr%r>r's rr-z_Worker.destroys r*Nr3r6)r:r;r<r#r=r)r/r-r>r*rr!r!s8 5 O &) &" &  & r*DEFAULT_CPU_COUNTr8PosArgsT_available_workers_default_interpreter_limitercR|D]}|j|jyr%)r-clear)workersworkers r _stop_workersrks& MMOr*limitercK| t} tj}|4d{ |j}dddd{ tjj|||d{t}|rT||dj z t"krn:tj|j%j&|d{|rTt|_|j)|S#t$r=t }tj |t jt|YwxYw7#t$rt}YwxYw7#1d{7swY(xYw77#t}|rU||dj z t"krn;tj|j%j&|d{7|rUt_|j)|wxYww)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 Nrlr)r _idle_workersrV LookupErrorrsetatexitregisterrkpop IndexErrorr!r rr/r r#MAX_WORKER_IDLE_TIMEpopleftr-append)rrmr idle_workersrjnows rrrs*575$((*  !%%'F $'' KK     n\!_...2FF$$\%9%9%;%C%CWU U U  (>F#9 5w ,' |45  YF   V n\!_...2FF$$\%9%9%;%C%CWU U U  (>F#s HC<HEHE%E H E" H&F7E;8F;AHE>H!HHEHEE%EE%"H%E8+E. ,E83H;F>HAHGH"!HHc tjS#t$r?tt j xst }tj||cYSwxYw)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 )rfrVrpros cpu_countrbrqrls rrrsN+//11 !",,."E4EF$((1sAAA)rzCallable[..., Any]rr7r4ztuple[Any, bool])rizdeque[_Worker]r4r5)rz&Callable[[Unpack[PosArgsT]], T_Retval]rzUnpack[PosArgsT]rmzCapacityLimiter | Noner4r8)r4r)3 __future__r__all__rrr|sys collectionsrcollections.abcrtypingrr r r r _core._exceptionsr_core._synchronizationrlowlevelr version_inforrtyping_extensionsconcurrent.interpretersrrrr!rLrJrAr=rBrCrDrEcompilerUrbrvr8rdrorfrkrrr>r*rrs"   $&&%63w++6w?! !(7! !.  GUM5K +W5u5"/!99 0 5I:##L"5 :   #&uW~&';< 6vo67UV'+6$ 06$ 6$$6$ 6$rr*