Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
_executor.cpython-313.pyc75 linesDownload Raw Back to __pycache__
1�

2��j���<�SSKJr SSKrSSKrSSKrSSKJrJrJ	r	 SSK3JrJrJ
r
 SSKJr SSKJr SSKJrJrJr SSKJr SS	KJr SS4KJr SSKJrJr SSKJ r  \"S
5r!\"S5r""SS\\!\"45r#"SS\5r$"SS\5r%SSjr&SSjr'g)�)�annotationsN)�	Awaitable�Callable�	Coroutine)�AbstractAsyncContextManager�AbstractContextManager�	ExitStack)�copy_context)�
TracebackType)�Protocol�TypeVar�cast)�RunnableConfig)�get_executor_for_config)�	ParamSpec)�CONTEXT_NOT_SUPPORTED�run_coroutine_threadsafe)�
GraphBubbleUp�P�Tc�J�\rSrSrSSSSS.SSjjrSrg)	�Submit�NFT��__name__�__cancel_on_exit__�__reraise_on_exit__�
__next_tick__c��g�N�)�self�fnrrrr�args�kwargss        �`D:\code\apps\devtools\python\user_packages\Python313\site-packages\langgraph/pregel/_executor.py�__call__�Submit.__call__s��(+�r!�r#�Callable[P, T]r$�P.argsr�5str | Noner�boolrr.rr.r%�P.kwargs�returnzconcurrent.futures.Future[T])r�6__module__�__qualname__�__firstlineno__r'�__static_attributes__r!r)r&rrsh��7 $�#(�$(�#�	+��	+��	+��		+�8!�	+�"�
	+��	+��	+�9&�	+�	+r)rc��\rSrSrSrS
SjrSSSSS.SSjjrSS	jrSS10jrSSjr	Sr11g)�BackgroundExecutor�(a5A context manager that runs sync tasks in the background.12Uses a thread pool executor to delegate tasks to separate threads.13On exit,14- cancels any (not yet started) tasks with `__cancel_on_exit__=True`15- waits for all tasks to finish16- re-raises the first exception from tasks with `__reraise_on_exit__=True`c��[5UlURR[U55Ul0Ulgr )r	�stack�
enter_contextr�executor�tasks)r"�configs  r&�__init__�BackgroundExecutor.__init__0s.���[��17��18�19�0�0�1H��1P�Q��
�IK��20r)NFTrc��[5nU(aZ[[RR[21URR"[URU/UQ70UD65n	O+URR"URU/UQ70UD6n	X44URU	'U	RUR5 U	$r )
r22r�23concurrent�futures�Futurerr;�submit�	next_tick�runr<�add_done_callback�done)24r"r#rrrrr$r%�ctx�tasks25          r&rD�BackgroundExecutor.submit6s����n�����"�"�)�)�!�,��
�
�$�$�Y�����M�d�M�f�M��D�26�=�=�'�'�����E�d�E�f�E�D�.�D��27�28�4�����t�y�y�)��r)c���UR5 URRU5 g![a URRU5 g[a gf=f)z3Remove the task from the tasks dict when it's done.N)�resultr<�popr�
BaseException)r"rJs  r&rH�BackgroundExecutor.doneMsP��		!��K�K�M�
�J�J�N�N�4� ���	!�
�J�J�N�N�4� ��	��	�s�.�%A!�	A!� A!c��UR$r �rD�r"s r&�	__enter__�BackgroundExecutor.__enter__Zs���{�{�r)c�@�URR5nUR5H!unupgU(dMUR5 M# UVs1sHo�R	5(aMUiM sn=n	(a[29RRU	5 URRXU5 Uc7UR5H"unupzU30(dMUR5 M$ ggs snf![31RRa MNf=fr )r<�copy�items�cancelrHrArB�waitr9�__exit__rM�CancelledError)r"�exc_type�	exc_value�	tracebackr<rJrY�_�t�pending�reraises           r&r[�BackgroundExecutor.__exit__]s����32�33���!��!&�����D�+�6��v����
�"/�#(�8�%�Q�v�v�x�q�%�8�8�7�8����#�#�G�,��34�35���H��;���&+�k�k�m�"��l�q�����K�K�M�	'4���9��"�)�)�8�8����s�C6�0C6� C;�;D�D)r;r9r<�r=rr0�Noner*)rJzconcurrent.futures.Futurer0rf�r0r)r]�type[BaseException] | Noner^�BaseException | Noner_�TracebackType | Noner0zbool | None)rr1r2r3�__doc__r>rDrHrTr[r4r!r)r&r6r6(s���R�L� $�#(�$(�#�������	�36!��"�
�����37&��.!���,��(��(�	�3839�r)r6c��\rSrSrSrS
SjrSSSSS.SSjjrSS	jrSS40jrSSjr	Sr41g)�AsyncBackgroundExecutor�za;A context manager that runs async tasks in the background.42Uses the current event loop to delegate tasks to asyncio tasks.43On exit,44- cancels any tasks with `__cancel_on_exit__=True`45- waits for all tasks to finish46- re-raises the first exception from tasks with `__reraise_on_exit__=True`47  ignoring CancelledErrorc���0Ul[5Ul[R"5UlUR
S5=n(a[R"U5UlgSUlg)N�max_concurrency)	r<�object�sentinel�asyncio�get_running_loop�loop�get�	Semaphore�	semaphore)r"r=rps   r&r>� AsyncBackgroundExecutor.__init__�sU��>@��48����
��,�,�.��	�$�j�j�):�;�;�?�;�7>�7H�7H��8�D�N�"�D�Nr)NFTrc�h�[[SS[4U"U0UD65nUR(a[	URU5n[49(a[
X�RX%S9n	O[
UURU[5US9n	X44URU	'U	RUR5 U	$)N)�name�lazy)r{�contextr|)rrrrx�gatedrrrur50r<rGrH)51r"r#rrrrr$r%�cororJs52          r&rD�AsyncBackgroundExecutor.submit�s����I�d�D�!�m�,�b�$�.A�&�.A�B���>�>������.�D� � �+��i�i�h��D�,���	�	��$��"��D�/�D��53�54�4�����t�y�y�)��r)c�8�UR5=n(a2[U[5(aURR	U5 ggURR	U5 g![55Ra URR	U5 gf=fr )�	exception�56isinstancerr<rNrsr\)r"rJ�excs   r&rH�AsyncBackgroundExecutor.done�st��		!��n�n�&�&�s�&��c�=�1�1��J�J�N�N�4�(�2��57�58���t�$���%�%�	!��J�J�N�N�4� �	!�s�AA'�A'�'/B�Bc��"# �UR$7fr rRrSs r&�59__aenter__�"AsyncBackgroundExecutor.__aenter__�s����{�{��s�
c���# �URR5nUR5H,unupgU(dMURUR5 M. U(a[60R"U5IShv�N Uc@UR5H+unupxU(dMUR5=n	(aU	eM- ggNH![61Ra MJf=f7fr )	r<rWrXrYrrrsrZr�r\)62r"r]r^r_r<rJrYr`rcr�s63          r&�	__aexit__�!AsyncBackgroundExecutor.__aexit__�s�����64�65���!��!&�����D�+�6��v����D�M�M�*�"/���,�,�u�%�%�%���&+�k�k�m�"��l�q����"�n�n�.�.�s�.�!�	�/�	'4��
&���-�-����s:�8C#�?C#�=C�>(C#�'C�C#�C �C#�C � C#)rurxrrr<re)r#zCallable[P, Awaitable[T]]r$r,rr-rr.rr.rr.r%r/r0zasyncio.Future[T])rJzasyncio.Futurer0rfrg)r]rhr^rir_rjr0rf)rr1r2r3rkr>rDrHr�r�r4r!r)r&rmrmzs���!�	"� $�#(�$(�#��%�����	�66!��"�
�����67��:68!���,��(��(�	�6970�r)rmc��# �UIShv�N UIShv�NsSSS5IShv�N $NNN	!,IShv�N(df   g=f7f)zHA coroutine that waits for a semaphore before running another coroutine.Nr!)rxrs  r&r~r~�s&����y��z��y�y���y�y�y�sE�A	�)�A	�/�+�/�A	�-�A	�/�A	�A�8�A�A	c�>�[R"S5 U"U0UD6$)zPA function that yields control to other threads before running another function.r)�time�sleep)r#r$r%s   r&rErE�s���J�J�q�M�
�t��v��r))rxzasyncio.SemaphorerzCoroutine[None, None, T]r0r)r#r+r$r,r%r/r0r)(�71__future__rrs�concurrent.futuresrAr��collections.abcrrr�72contextlibrrr	�contextvarsr73�typesr�typingrr
r�langchain_core.runnablesr�langchain_core.runnables.configr�typing_extensionsr�langgraph._internal._futurerr�langgraph.errorsrrrrr6rmr~rEr!r)r&�<module>r�s���"����:�:�U�U�$����4�C�'�W�*�
�c�N���C�L��74+�X�a��d�^�75+�O�/�O�dY�9�Y�x�r)
codekingpro/portable-devtools · Team Ai