Team Ai
Datasetpublic

codekingpro/portable-devtools

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

2��j�N���SSKJr SSKrSSKrSSKJrJr SSKJr SSK	J3r4JrJr SSK
Jr \S/\4r"SS	5rg)5�)�annotationsN)�	Awaitable�Callable)�Any)�
ProtocolEvent�StreamTransformer�transformer_requires_async)�
StreamChannel�tuple[str, ...]c���\rSrSrSrSSSSSS.SSjjjrSS	jrSS6jrSSjrSSjr	S S
jr7S!SjrS"SjrS#Sjr
S$SjrS"SjrS#SjrS$SjrS%SjrSS.S&SjjrS'SjrSrg)(�	StreamMux�u5Central event dispatcher for the streaming infrastructure.8 9Owns the main event log and routes events through a transformer10pipeline. StreamChannels with a name discovered in transformer11projections are auto-wired so that every `push()` also injects a12`ProtocolEvent` into the main log. StreamChannels without a name13are local-only.14 15Pass `is_async=True` when the mux will be consumed via async16iteration (`handler.astream()`). All StreamChannel instances17discovered during registration are automatically bound to the18matching mode.19 20Attributes:21    extensions: Merged projection dict across all registered22        transformers. Treat as read-only — mutations won't be23        reflected back in individual transformers' state.24    native_keys: Projection keys contributed by transformers with25        `_native = True`.26NF�T)�is_async�	factories�scope�_assign_seqc��X lX@lXPl[5UlURRUS9 URR
U5 /Ul/UlSUl	SUl270Ul[5Ul
0Ul0UlUb[!U5OSUlSUlSUlUb UHnUR)U"U55 M U=(d SHnUR)U5 M g)u�Initialize the mux and register transformers in order.28 29Callers pass either `transformers` (pre-built instances) or30`factories` (callables producing fresh instances per mux). Each31transformer's `init()` is called, projections are merged into32`extensions`, `_native` keys are recorded in `native_keys`, and33any StreamChannel instances are bound and (if named) wired.34 35Args:36    transformers: Already-built transformer instances. Registered37        only on this mux — they are NOT cloned into child38        mini-muxes built by `_make_child`. Use `factories` for39        transformers that should propagate to nested scopes.40    is_async: True for async dispatch (`apush` / `aclose` /41        `afail`), False for the sync path.42    factories: One-argument callables `(scope) -> StreamTransformer`.43        Called once with this mux's `scope` here, and cloned44        again per child scope by `_make_child` so each45        sub-mux gets fresh instances.46    scope: The namespace the mux operates within. The root mux47        is `()`.48    _assign_seq: Internal flag for child muxes. Root muxes assign49        monotonic `seq` numbers when appending to their main event50        log; child muxes share forwarded event objects and must not51        mutate their envelopes.52 53Raises:54    RuntimeError: If any transformer requires an async run but55        the mux is in sync mode.56    TypeError: If a transformer's `init()` doesn't return a dict.57    ValueError: If transformers' projection keys collide.58�rrNr)rrrr59�_events�_bind�	_bind_mux�
_transformers�	_channels�_seq�	_push_seq�60extensions�set�native_keys�_projection_owners�_transformer_by_key�list�61_factories�_pump_fn�	_apump_fn�	_register)�self�transformersrrrr�factory�transformers        �[D:\code\apps\devtools\python\user_packages\Python313\site-packages\langgraph/stream/_mux.py�__init__�StreamMux.__init__0s���R!�
�&+�62�&��5B�_��������H��-������t�$�68���35�����	����*,���%(�U���24���AC�� � )�4�D��O�$�	
��48��
�?C���� �$�����w�u�~�.�%�'�-�2�-�K��N�N�;�'�.�c�8�URRU5$)z@Return the transformer that contributed `key` to the projection.)r!�get)r'�keys  r+�transformer_by_key�StreamMux.transformer_by_key}s���'�'�+�+�C�0�0r.c�D�U=RS-
slUR$)N�)r)r's r+�_next_push_seq�StreamMux._next_push_seq�s�����!����~�~�r.c��XlXRlURH	nXlM URHn[USS5nUcMU"U5 M g)a�Wire the sync pull callback onto every projection in this mux.63 64Records the pump on the mux so child mini-muxes built by65`_make_child` can inherit it. Propagates to:66- the main event log (`self._events`)67- every projection StreamChannel in `extensions`68- any registered transformer that exposes `_bind_pump` (e.g.69  `MessagesTransformer` so `ChatModelStream` instances drive the70  shared pump from their cursors)71�72_bind_pumpN)r$r�
_request_morerr�getattr)r'�fn�chr*�binds     r+�	bind_pump�StreamMux.bind_pump�sR���
�%'���"��.�.�B�!��!��-�-�K��;��d�;�D����R��.r.c��XlXRlURH	nXlM URHn[USS5nUcMU"U5 M g)z!Async counterpart to `bind_pump`.�_bind_apumpN)r%r�_arequest_morerrr;)r'r<r=r*�abinds     r+�73bind_apump�StreamMux.bind_apump�sP����&(���#��.�.�B� "��!��-�-�K��K���=�E�� ��b�	�.r.c��URc[S5e[URURUSS9nURbURUR5 URbURUR5 U$)ajBuild a mini-mux with the same factories scoped to `scope`.74 75Used by `SubgraphTransformer` to attach a fresh transformer76pipeline to each discovered subgraph handle. The child mux77inherits the current pump bindings (so cursors on its78projection logs drive the root pump), carries the same factory79list forward to any grandchild subgraphs, and does not assign80`seq` numbers so forwarded events can be shared without81mutating their envelope.82 83Raises:84    RuntimeError: If the mux was not constructed with85        `factories=`. Mini-muxes require factories so each scope86        gets its own fresh transformer instances.87z�StreamMux._make_child requires the mux to be constructed with `factories=`; pre-built transformers can't be cloned to a new scope.F)rrrr)r#�RuntimeErrorr
rr$r?r%rE)r'r�childs   r+�_make_child�StreamMux._make_child�s}�� �?�?�"��)��
�88��o�o��]�]���	89���=�=�$��O�O�D�M�M�*��>�>�%����T�^�^�,��r.c�^�[U5(a2TR(d![[U5RS35eUR5n[
U[5(d![S[U5R35e[U5[TR5-nU(aHSRU4Sj[U555n[S[U5RSU35e[[USS55nTR R#U5 TR%X%S	9 TRR'U5 [U5RnUH!nUTR(U'UTR*U'M# U(a)TR,R'UR/55 UR1T5 g90)z�Register a single transformer.91 92Calls `transformer.init()`, stores the transformer for event93processing, binds any StreamChannel instances in the projection,94and merges the projection into `extensions`.95uz requires an async run — it overrides aprocess/afinalize/afail or sets requires_async=True. Use astream(), not stream().z1StreamTransformer.init() must return a dict, got z, c3�P># �UHnU<STRUS3v� M g7f)z (owned by �)N)r )�.0r1r's  �r+�	<genexpr>�&StreamMux._register.<locals>.<genexpr>�s1����%�,�C��'��T�%<�%<�S�%A�$B�!�D�,�s�#&zTransformer zF returned projection keys that conflict with already-registered keys: �_nativeF��nativeN)r	rrH�type�__name__�init�96isinstance�dict�	TypeErrorrr�join�sorted�97ValueError�boolr;r�append�_bind_and_wire�updater r!r�keys�_on_register)r'r*�98projection�	conflicts�attributions�	is_native�99owner_namer1s`       r+r&�StreamMux._register�s����&�k�2�2�4�=�=����$�-�-�.�/D�D��
�100!�%�%�'�101��*�d�+�+����J�'�0�0�1�3��
��102�O�c�$�/�/�&:�:�	���9�9�%�!�)�,�%��L���t�K�0�9�9�:�;�%��(��
�103���i��?�@�	����!�!�+�.����J��9������z�*��+�&�/�/�104��C�+5�D�#�#�C�(�,7�D�$�$�S�)������#�#�J�O�O�$5�6�� � ��&r.c��SnURHnURU5(aMSnM U(aQUR(a$U=RS-
slURUS'URRU5 gg)a�Route an event through all transformers, then append to the main log.105 106Each transformer's `process()` is called in registration order.107If any transformer returns False, the event is suppressed from108the main log, but transformers that already saw it keep their109side effects.110 111On the root mux, `seq` is assigned right before an event enters112the main log, not before the transformer pipeline runs. This113ensures that events auto-forwarded from StreamChannels during114`process()` get earlier seq numbers than the original event,115preserving monotonic ordering in the root log. Child muxes do116not assign `seq`, so subgraph forwarding can share event objects117without mutating their envelopes.118 119Args:120    event: The protocol event to dispatch.121TFr5�seqN)r�processrrr�push�r'�event�keepr*s    r+rm�StreamMux.push�sn��&���-�-�K��&�&�u�-�-���.������	�	�Q��	�#�y�y��e���L�L���e�$�	r.c�@�SnURHnUR5 M URH&nUR(aMUR5 M( URR5 UbUeg![anUcUnSnAM�SnAM�SnAff=f)u�Finalize all transformers, close all projections and the main log.122 123StreamChannels discovered in transformer projections are124auto-closed after `finalize()` runs — transformers don't need125to close them manually. If any transformer's `finalize()` raises,126the remaining transformers, projections, and the main log are127still closed; the first error is re-raised after cleanup128completes.129 130Raises:131    BaseException: The first error raised by a transformer's132        `finalize()`, re-raised after cleanup finishes.133N)r�finalize�
BaseExceptionr�_closed�closer)r'�first_errorr*�er=s     r+rv�StreamMux.closes���-1���-�-�K�
$��$�$�&�.��.�.�B��:�:�:����134�!�	
�������"���#��!�
$��&�"#�K�'��
$�s�A=�=135B�B�Bc��URHnURU5 M URH'nUR(aMURU5 M) UR136RU5 g![a Myf=f)uSFail all transformers, projections, and the main log.137 138StreamChannels discovered in transformer projections are139auto-failed — transformers don't need to fail them manually.140If any transformer's `fail()` raises, the remaining141transformers, projections, and the main log are still failed.142 143Args:144    err: The exception that ended the run.145N)r�failrtrrur)r'�errr*r=s    r+r{�StreamMux.fail-ss�� �-�-�K�
�� � ��%�.�146�.�.�B��:�:�:������!�	
�����#���!�
��
�s�A9�9147B�Bc��.# �SnURH%nURU5IShv�N(aM#SnM' U(aQUR(a$U=RS-
slURUS'URRU5 ggNj7f)u�Dispatch an event on the async lane.148 149Awaits each transformer's `aprocess` in registration order150before appending to the main log. A slow `aprocess` serializes151the pipeline by design — that's the guarantee that lets a later152transformer (or a synchronous consumer) see the result of the153async work. For decoupled work, use `schedule()` from inside154`process` / `aprocess` instead.155 156The main log append is a non-blocking `push` — matching v1's157`put_nowait` shape. The root mux assigns `seq`; child muxes do158not, so forwarded subgraph events can be shared without copying.159Memory is bounded by caller pace via the caller-driven pump; see160`StreamChannel` for the full tradeoff story.161 162Args:163    event: The protocol event to dispatch.164TNFr5rk)r�aprocessrrrrmrns    r+�apush�StreamMux.apushFsy���&���-�-�K�$�-�-�e�4�4�4���.������	�	�Q��	�#�y�y��e���L�L���e�$�	�5�s�&B�B�B�A Bc���# �UR5nU(a6[R"USS06IShv�Nn[SU5S5nUbUeSnURHnUR5IShv�N M URH&nUR(aMUR5 M( URR5 UbUegN�N`![anUcUnSnAM�SnAM�SnAff=f7f)a�Finalize on the async lane.165 166Awaits every task started via `StreamTransformer.schedule()`167across all transformers, then calls `afinalize()` on each,168then auto-closes channels and the main event log.169 170If any scheduled task raised under `on_error="raise"`, or any171transformer's `afinalize` raises, the exception propagates.172The caller (the pump) handles it by routing into `afail`.173 174Raises:175    BaseException: The first scheduled-task or `afinalize`176        error, re-raised after cleanup.177�return_exceptionsTNc3�# �UH?n[U[5(dM[U[R5(aM;Uv� MA g7f�N)rXrt�asyncio�CancelledError)rO�rs  r+rP�#StreamMux.aclose.<locals>.<genexpr>vs:����$��!�!�]�3��'�q�'�*@�*@�A��A�$�s�A	�A	�	A	)�_collect_scheduled_tasksr��gather�nextr�	afinalizertrrurvr)r'�pending�results�	first_errrwr*rxr=s        r+�aclose�StreamMux.aclosecs�����/�/�1���#�N�N�G�L�t�L�L�G���$����I��$���,0���-�-�K�
$�!�+�+�-�-�-�.��.�.�B��:�:�:����178�!�	
�������"���#�1M� .�� �
$��&�"#�K�'��
$�sQ�1C;�C�-C;�"C�5C�6C�:"C;� 5C;�C�179C8�"C3�'C;�3C8�8C;c��# �UR5nUHnUR5 M U(a[R"USS06IShv�N URHnURU5IShv�N M URH'nUR(aMURU5 M) URR(dURRU5 ggN�Ny![a M�f=f7f)z�Fail on the async lane.180 181Cancels every scheduled task across all transformers, awaits182them to completion, then runs each transformer's `afail` hook183and auto-fails channels and the main event log.184 185Args:186    err: The exception that ended the run.187r�TN)r��cancelr�r�r�afailrtrrur{r)r'r|r��taskr*r=s      r+r��StreamMux.afail�s�����/�/�1���D��K�K�M����.�.�'�B�T�B�B�B��-�-�K�
�!�'�'��,�,�,�.�188�.�.�B��:�:�:������!��|�|�#�#��L�L���c�"�$�
C�-�� �
��
�sO�A189D�C-�
D�!C1�5C/�6C1�:"D� AD�/C1�1190C?�;D�>C?�?Dc	��URVVs/sH1n[USS5HnUR5(aMUPM M3 snn$s snnf)z@Return a snapshot of in-flight tasks scheduled via transformers.�_stream_scheduled_tasksr)rr;�done)r'r*r�s   r+r��"StreamMux._collect_scheduled_tasks�sO�� $�1�1�191�1����-F��K���9�9�;�
�K�
�1�192�	193��194s195�*A�196ArSc�^�UR5H�n[U[5(dMURTRS9 URT5 TRRU5 URcMnU(aUROSUR3nSU4SjjnURU"U55 M� g)a�Bind and optionally wire StreamChannel instances in a projection.197 198All StreamChannels are bound and tracked. Channels with a name199are additionally wired for protocol auto-forwarding.200 201Args:202    projection: The projection dict returned by a transformer's203        `init()`.204    native: True when the owning transformer is `_native`.205        Named channels owned by a native transformer use the206        channel name directly as the protocol method;207        user-defined channels are prefixed with `custom:`.208rNzcustom:c�>^�SUU4SjjnU$)Nc�*>�TRTU5 gr�)�_forward)�item�method_namer's ��r+r��AStreamMux._bind_and_wire.<locals>._make_forward.<locals>._forward�s��� �M�M�+�t�<r.)r�r�return�Noner)r�r�r's` �r+�
_make_forward�/StreamMux._bind_and_wire.<locals>._make_forward�s���=�=� (�r.)r��strr�zCallable[[Any], None])209�valuesrXr210rrrrr_�name�_wire)r'rdrT�value�methodr�s`     r+r`�StreamMux._bind_and_wire�s����  �&�&�(�E��%��/�/����T�]�]��3�����%����%�%�e�,��:�:�)�+1�U�Z�Z������7M�F�(��K�K�
�f� 5�6�)r.c���SU/[[R"5S-5US.S.nUR(a$U=RS-
slURUS'URRU5 g)a6Inject a ProtocolEvent for a StreamChannel push.211 212Forwarded events bypass the transformer pipeline to avoid213infinite recursion (a transformer that pushes to a channel214during `process()` would re-trigger itself). These events are215visible in this mux's main event log but are not passed through216transformers' `process()` methods. Only the root mux assigns217`seq` to forwarded channel events.218 219Args:220    method: The full protocol method (already with or without221        the `custom:` prefix; resolved by `_bind_and_wire`).222    item: The payload pushed onto the channel.223roi�)�	namespace�	timestamp�data)rUr��paramsr5rkN)�int�timerrrrm)r'r�r�ros    r+r��StreamMux._forward�sf�� ��� �����t�!3�4��� 224������I�I��N�I��9�9�E�%�L������%� r.)r%rrrr#r r$rrr!rrrrrr�)r(zlist[StreamTransformer] | Nonerr^rzlist[TransformerFactory] | Nonerrrr^r�r�)r1r�r�zStreamTransformer | None)r�r�)r<zCallable[[], bool]r�r�)r<zCallable[[], Awaitable[bool]]r�r�)rrr�r
)r*rr�r�)rorr�r�)r�r�)r|rtr�r�)r�zlist[asyncio.Task[Any]])rdzdict[str, Any]rTr^r�r�)r�r�r�rr�r�)rV�225__module__�__qualname__�__firstlineno__�__doc__r,r2r6r?rErJr&rmrvr{r�r�r�r�r`r��__static_attributes__rr.r+r
r
s����.8<�K(��59�!#� �K(�4�K(��	K(�2263�K(��
K(��K(�227�K(�Z1���(	� �D('�T%�:�8�2%�:*�X#�6228�=B�7�(�7�59�7�	
�7�@!r.r
)�229__future__rr�r��collections.abcrr�typingr�langgraph.stream._typesrrr	�langgraph.stream.stream_channelr230�TransformerFactoryr
rr.r+�<module>r�sI��"���/����231:��0�1�3D�D�E���X!�X!r.
codekingpro/portable-devtools · Team Ai