§
    šŠtjçS  ã                  ó�   — d dl mZ d dlZd dlZd dlmZmZ d dlmZ d dl	m
Z
mZmZ d dlmZ edgef         Z	  G d„ d	¦  «        ZdS )
é    )ÚannotationsN)Ú	AwaitableÚCallable)ÚAny)ÚProtocolEventÚStreamTransformerÚtransformer_requires_async)ÚStreamChannelútuple[str, ...]c                  ó¨   — e Zd ZdZ	 d7dddddœd8d„Zd9d„Zd:d„Zd;d„Zd<d„Zd=d„Z	d>d!„Z
d?d$„Zd@d%„ZdAd(„Zd?d)„Zd@d*„ZdAd+„ZdBd-„Zdd.œdCd2„ZdDd6„ZdS )EÚ	StreamMuxuu  Central event dispatcher for the streaming infrastructure.

    Owns the main event log and routes events through a transformer
    pipeline. StreamChannels with a name discovered in transformer
    projections are auto-wired so that every `push()` also injects a
    `ProtocolEvent` into the main log. StreamChannels without a name
    are local-only.

    Pass `is_async=True` when the mux will be consumed via async
    iteration (`handler.astream()`). All StreamChannel instances
    discovered during registration are automatically bound to the
    matching mode.

    Attributes:
        extensions: Merged projection dict across all registered
            transformers. Treat as read-only â€” mutations won't be
            reflected back in individual transformers' state.
        native_keys: Projection keys contributed by transformers with
            `_native = True`.
    NF© T)Úis_asyncÚ	factoriesÚscopeÚ_assign_seqÚtransformersúlist[StreamTransformer] | Noner   Úboolr   úlist[TransformerFactory] | Noner   r   r   ÚreturnÚNonec               ó:  — || _         || _        || _        t          ¦   «         | _        | j                             |¬¦  «         | j                             | ¦  «         g | _        g | _        d| _	        d| _
        i | _        t          ¦   «         | _        i | _        i | _        |�t!          |¦  «        nd| _        d| _        d| _        g }g }|�|D ]5} ||¦  «        }	t)          |	dd¦  «        r|n|                     |	¦  «         Œ6g |¢|¢R D ]}	|                      |	¦  «         Œ|                     ¦   «          |                     ¦   «          |pdD ]*}	t)          |	dd¦  «        r|n|                     |	¦  «         Œ+g |¢|¢R D ]}	|                      |	¦  «         ŒdS )uy  Initialize the mux and register transformers in order.

        Callers pass either `transformers` (pre-built instances) or
        `factories` (callables producing fresh instances per mux). Each
        transformer's `init()` is called, projections are merged into
        `extensions`, `_native` keys are recorded in `native_keys`, and
        any StreamChannel instances are bound and (if named) wired.

        Transformers with `StreamTransformer.before_builtins = True` are
        registered ahead of the rest, preserving relative order within
        each lane. This lets content-mutating transformers (PII
        redaction, content filters, etc.) run before built-ins like
        `MessagesTransformer` that eagerly snapshot text fields into
        their projections. See `StreamTransformer.before_builtins` for
        the contract and foot-guns.

        Args:
            transformers: Already-built transformer instances. Registered
                only on this mux â€” they are NOT cloned into child
                mini-muxes built by `_make_child`. Use `factories` for
                transformers that should propagate to nested scopes.
            is_async: True for async dispatch (`apush` / `aclose` /
                `afail`), False for the sync path.
            factories: One-argument callables `(scope) -> StreamTransformer`.
                Called once with this mux's `scope` here, and cloned
                again per child scope by `_make_child` so each
                sub-mux gets fresh instances.
            scope: The namespace the mux operates within. The root mux
                is `()`.
            _assign_seq: Internal flag for child muxes. Root muxes assign
                monotonic `seq` numbers when appending to their main event
                log; child muxes share forwarded event objects and must not
                mutate their envelopes.

        Raises:
            RuntimeError: If any transformer requires an async run but
                the mux is in sync mode.
            TypeError: If a transformer's `init()` doesn't return a dict.
            ValueError: If transformers' projection keys collide.
        ©r   r   NÚbefore_builtinsFr   )r   r   r   r
   Ú_eventsÚ_bindÚ	_bind_muxÚ_transformersÚ	_channelsÚ_seqÚ	_push_seqÚ
extensionsÚsetÚnative_keysÚ_projection_ownersÚ_transformer_by_keyÚlistÚ
_factoriesÚ_pump_fnÚ	_apump_fnÚgetattrÚappendÚ	_registerÚclear)
Úselfr   r   r   r   r   ÚpreÚrestÚfactoryÚtransformers
             úS/var/www/html/CA-Chatbot/venv/lib/python3.11/site-packages/langgraph/stream/_mux.pyÚ__init__zStreamMux.__init__0   sæ  € ðb !ˆŒØ&+ˆŒ
Ø&ˆÔÝ5B±_´_ˆŒØŒ×Ò HÐÑ-Ô-Ð-ØŒ×Ò˜tÑ$Ô$Ð$Ø68ˆÔØ35ˆŒØˆŒ	ØˆŒà*,ˆŒÝ%(¡U¤UˆÔØ24ˆÔØACˆÔ ð  )Ð4�D�‰OŒOˆO¸$ð 	Œð 48ˆŒØ?CˆŒð (*ˆØ(*ˆØÐ Ø$ð &ð &�Ø%˜g e™nœn�å" ;Ð0AÀ5ÑIÔIÐS�C�CÈtß’&˜Ñ%Ô%Ð%Ð%Ø, ˜} t˜}˜}ð ,ð ,�Ø—’˜{Ñ+Ô+Ð+Ð+Ø�IŠI‰KŒKˆKØ�JŠJ‰LŒLˆLØ'Ð-¨2ð 	ð 	ˆKÝ˜KÐ):¸EÑBÔBÐLˆSˆSÈ×TÒTØñô ð ð ð )˜S˜= 4˜=˜=ð 	(ð 	(ˆKØ�NŠN˜;Ñ'Ô'Ð'Ð'ð	(ð 	(ó    ÚkeyÚstrúStreamTransformer | Nonec                ó6   — | j                              |¦  «        S )z@Return the transformer that contributed `key` to the projection.)r'   Úget)r0   r8   s     r5   Útransformer_by_keyzStreamMux.transformer_by_key–   s   € àÔ'×+Ò+¨CÑ0Ô0Ð0r7   Úintc                ó0   — | xj         dz  c_         | j         S )Né   )r"   ©r0   s    r5   Ú_next_push_seqzStreamMux._next_push_seqš   s   € ØˆŒ˜!ÑˆŒØŒ~Ðr7   ÚfnúCallable[[], bool]c                óž   — || _         || j        _        | j        D ]	}||_        Œ
| j        D ] }t          |dd¦  «        }|� ||¦  «         Œ!dS )aÝ  Wire the sync pull callback onto every projection in this mux.

        Records the pump on the mux so child mini-muxes built by
        `_make_child` can inherit it. Propagates to:
        - the main event log (`self._events`)
        - every projection StreamChannel in `extensions`
        - any registered transformer that exposes `_bind_pump` (e.g.
          `MessagesTransformer` so `ChatModelStream` instances drive the
          shared pump from their cursors)
        Ú
_bind_pumpN)r*   r   Ú_request_morer    r   r,   )r0   rC   Úchr4   Úbinds        r5   Ú	bind_pumpzStreamMux.bind_pump¢   sr   € ð ˆŒØ%'ˆŒÔ"Ø”.ð 	"ð 	"ˆBØ!ˆBÔÐØÔ-ð 	ð 	ˆKÝ˜;¨°dÑ;Ô;ˆDØÐØ��R‘”�øð	ð 	r7   úCallable[[], Awaitable[bool]]c                óž   — || _         || j        _        | j        D ]	}||_        Œ
| j        D ] }t          |dd¦  «        }|� ||¦  «         Œ!dS )z!Async counterpart to `bind_pump`.Ú_bind_apumpN)r+   r   Ú_arequest_morer    r   r,   )r0   rC   rH   r4   Úabinds        r5   Ú
bind_apumpzStreamMux.bind_apump¶   sp   € àˆŒØ&(ˆŒÔ#Ø”.ð 	#ð 	#ˆBØ "ˆBÔÐØÔ-ð 	ð 	ˆKÝ˜K¨¸Ñ=Ô=ˆEØÐ Ø��b‘	”	�	øð	ð 	r7   c                óð   — | j         €t          d¦  «        ‚t          | j         | j        |d¬¦  «        }| j        �|                     | j        ¦  «         | j        �|                     | j        ¦  «         |S )aÊ  Build a mini-mux with the same factories scoped to `scope`.

        Used by `SubgraphTransformer` to attach a fresh transformer
        pipeline to each discovered subgraph handle. The child mux
        inherits the current pump bindings (so cursors on its
        projection logs drive the root pump), carries the same factory
        list forward to any grandchild subgraphs, and does not assign
        `seq` numbers so forwarded events can be shared without
        mutating their envelope.

        Raises:
            RuntimeError: If the mux was not constructed with
                `factories=`. Mini-muxes require factories so each scope
                gets its own fresh transformer instances.
        Nz‚StreamMux._make_child requires the mux to be constructed with `factories=`; pre-built transformers can't be cloned to a new scope.F)r   r   r   r   )r)   ÚRuntimeErrorr   r   r*   rJ   r+   rP   )r0   r   Úchilds      r5   Ú_make_childzStreamMux._make_childÁ   sŠ   € ð  Œ?Ð"Ýð)ñô ð õ
 Ø”oØ”]ØØð	
ñ 
ô 
ˆð Œ=Ð$Ø�OŠO˜DœMÑ*Ô*Ð*ØŒ>Ð%Ø×Ò˜Tœ^Ñ,Ô,Ð,Øˆr7   r4   r   c                ó¾  ‡ — t          |¦  «        r+‰ j        s$t          t          |¦  «        j        › d�¦  «        ‚|                     ¦   «         }t          |t          ¦  «        s$t          dt          |¦  «        j        › �¦  «        ‚t          |¦  «        t          ‰ j
        ¦  «        z  }|rUd                     ˆ fd„t          |¦  «        D ¦   «         ¦  «        }t          dt          |¦  «        j        › d|› �¦  «        ‚t          t          |dd¦  «        ¦  «        }‰ j                             |¦  «         ‰                      ||¬	¦  «         ‰ j
                             |¦  «         t          |¦  «        j        }|D ]}|‰ j        |<   |‰ j        |<   Œ|r,‰ j                             |                     ¦   «         ¦  «         |                     ‰ ¦  «         d
S )zëRegister a single transformer.

        Calls `transformer.init()`, stores the transformer for event
        processing, binds any StreamChannel instances in the projection,
        and merges the projection into `extensions`.
        uz    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, c              3  ó>   •K  — | ]}|›d ‰j         |         › d�V — ŒdS )z (owned by ú)N)r&   )Ú.0r8   r0   s     €r5   ú	<genexpr>z&StreamMux._register.<locals>.<genexpr>ø   sP   øè è € ð %ð %àð ÐDÐD TÔ%<¸SÔ%AÐDÐDÐDð%ð %ð %ð %ð %ð %r7   zTransformer zF returned projection keys that conflict with already-registered keys: Ú_nativeF©ÚnativeN)r	   r   rR   ÚtypeÚ__name__ÚinitÚ
isinstanceÚdictÚ	TypeErrorr$   r#   ÚjoinÚsortedÚ
ValueErrorr   r,   r   r-   Ú_bind_and_wireÚupdater&   r'   r%   ÚkeysÚ_on_register)r0   r4   Ú
projectionÚ	conflictsÚattributionsÚ	is_nativeÚ
owner_namer8   s   `       r5   r.   zStreamMux._registerã   s&  ø€ õ & kÑ2Ô2ð 	¸4¼=ð 	ÝÝ˜Ñ$Ô$Ô-ð Dð Dð Dñô ð ð
 !×%Ò%Ñ'Ô'ˆ
Ý˜*¥dÑ+Ô+ð 	Ýð3Ý˜JÑ'Ô'Ô0ð3ð 3ñô ð õ ˜
‘O”O¥c¨$¬/Ñ&:Ô&:Ñ:ˆ	Øð 		ØŸ9š9ð %ð %ð %ð %å! )Ñ,Ô,ð%ñ %ô %ñ ô ˆLõ ð(�t KÑ0Ô0Ô9ð (ð (à%ð(ð (ñô ð õ
 � ¨i¸Ñ?Ô?Ñ@Ô@ˆ	ØÔ×!Ò! +Ñ.Ô.Ð.Ø×Ò˜J¨yÐÑ9Ô9Ð9ØŒ×Ò˜zÑ*Ô*Ð*Ý˜+Ñ&Ô&Ô/ˆ
Øð 	8ð 	8ˆCØ+5ˆDÔ# CÑ(Ø,7ˆDÔ$ SÑ)Ð)Øð 	7ØÔ×#Ò# J§O¢OÑ$5Ô$5Ñ6Ô6Ð6Ø× Ò  Ñ&Ô&Ð&Ð&Ð&r7   Úeventr   c                óÊ   — d}| j         D ]}|                     |¦  «        sd}Œ|r=| j        r| xj        dz  c_        | j        |d<   | j                             |¦  «         dS dS )aN  Route an event through all transformers, then append to the main log.

        Each transformer's `process()` is called in registration order.
        If any transformer returns False, the event is suppressed from
        the main log, but transformers that already saw it keep their
        side effects.

        On the root mux, `seq` is assigned right before an event enters
        the main log, not before the transformer pipeline runs. This
        ensures that events auto-forwarded from StreamChannels during
        `process()` get earlier seq numbers than the original event,
        preserving monotonic ordering in the root log. Child muxes do
        not assign `seq`, so subgraph forwarding can share event objects
        without mutating their envelopes.

        Args:
            event: The protocol event to dispatch.
        TFr@   ÚseqN)r   Úprocessr   r!   r   Úpush©r0   ro   Úkeepr4   s       r5   rs   zStreamMux.push  sŠ   € ð& ˆØÔ-ð 	ð 	ˆKØ×&Ò& uÑ-Ô-ð Ø�øØð 	%ØÔð )Ø�	”	˜Q‘�	”	Ø#œy��e‘ØŒL×Ò˜eÑ$Ô$Ð$Ð$Ð$ð		%ð 	%r7   c                ó  — d}| j         D ]2}	 |                     ¦   «          Œ# t          $ r}|€|}Y d}~Œ+d}~ww xY w| j        D ]}|j        s|                     ¦   «          Œ| j                             ¦   «          |�|‚dS )uO  Finalize all transformers, close all projections and the main log.

        StreamChannels discovered in transformer projections are
        auto-closed after `finalize()` runs â€” transformers don't need
        to close them manually. If any transformer's `finalize()` raises,
        the remaining transformers, projections, and the main log are
        still closed; the first error is re-raised after cleanup
        completes.

        Raises:
            BaseException: The first error raised by a transformer's
                `finalize()`, re-raised after cleanup finishes.
        N)r   ÚfinalizeÚBaseExceptionr    Ú_closedÚcloser   )r0   Úfirst_errorr4   ÚerH   s        r5   rz   zStreamMux.close*  s·   € ð -1ˆØÔ-ð 	$ð 	$ˆKð$Ø×$Ò$Ñ&Ô&Ð&Ð&øÝ ð $ð $ð $ØÐ&Ø"#�Køøøøøøøøøð$øøøð ”.ð 	ð 	ˆBØ”:ð Ø—’‘
”
�
øØŒ×ÒÑÔÐØÐ"ØÐð #Ð"s   �"¢
:¬5µ:Úerrrx   c                óæ   — | j         D ](}	 |                     |¦  «         Œ# t          $ r Y Œ%w xY w| j        D ]}|j        s|                     |¦  «         Œ| j                             |¦  «         dS )u‹  Fail all transformers, projections, and the main log.

        StreamChannels discovered in transformer projections are
        auto-failed â€” transformers don't need to fail them manually.
        If any transformer's `fail()` raises, the remaining
        transformers, projections, and the main log are still failed.

        Args:
            err: The exception that ended the run.
        N)r   Úfailrx   r    ry   r   )r0   r}   r4   rH   s       r5   r   zStreamMux.failF  s™   € ð  Ô-ð 	ð 	ˆKðØ× Ò  Ñ%Ô%Ð%Ð%øÝ ð ð ð Ø�ðøøøà”.ð 	ð 	ˆBØ”:ð Ø—’˜‘”�øØŒ×Ò˜#ÑÔÐÐÐs   ‹!¡
.­.c              ƒ  óÚ   K  — d}| j         D ]}|                     |¦  «        ƒ d{V —†sd}Œ |r=| j        r| xj        dz  c_        | j        |d<   | j                             |¦  «         dS dS )uM  Dispatch an event on the async lane.

        Awaits each transformer's `aprocess` in registration order
        before appending to the main log. A slow `aprocess` serializes
        the pipeline by design â€” that's the guarantee that lets a later
        transformer (or a synchronous consumer) see the result of the
        async work. For decoupled work, use `schedule()` from inside
        `process` / `aprocess` instead.

        The main log append is a non-blocking `push` â€” matching v1's
        `put_nowait` shape. The root mux assigns `seq`; child muxes do
        not, so forwarded subgraph events can be shared without copying.
        Memory is bounded by caller pace via the caller-driven pump; see
        `StreamChannel` for the full tradeoff story.

        Args:
            event: The protocol event to dispatch.
        TNFr@   rq   )r   Úaprocessr   r!   r   rs   rt   s       r5   ÚapushzStreamMux.apush_  s    è è € ð& ˆØÔ-ð 	ð 	ˆKØ$×-Ò-¨eÑ4Ô4Ð4Ð4Ð4Ð4Ð4Ð4ð Ø�øØð 	%ØÔð )Ø�	”	˜Q‘�	”	Ø#œy��e‘ØŒL×Ò˜eÑ$Ô$Ð$Ð$Ð$ð		%ð 	%r7   c              ƒ  ó¨  K  — |                       ¦   «         }|r5t          j        |ddiŽƒ d{V —†}t          d„ |D ¦   «         d¦  «        }|�|‚d}| j        D ]8}	 |                     ¦   «         ƒ d{V —† Œ# t          $ r}|€|}Y d}~Œ1d}~ww xY w| j        D ]}|j        s| 	                    ¦   «          Œ| j
         	                    ¦   «          |�|‚dS )a8  Finalize on the async lane.

        Awaits every task started via `StreamTransformer.schedule()`
        across all transformers, then calls `afinalize()` on each,
        then auto-closes channels and the main event log.

        If any scheduled task raised under `on_error="raise"`, or any
        transformer's `afinalize` raises, the exception propagates.
        The caller (the pump) handles it by routing into `afail`.

        Raises:
            BaseException: The first scheduled-task or `afinalize`
                error, re-raised after cleanup.
        Úreturn_exceptionsTNc              3  óx   K  — | ]5}t          |t          ¦  «        ¯t          |t          j        ¦  «        °1|V — Œ6d S ©N)r`   rx   ÚasyncioÚCancelledError)rX   Úrs     r5   rY   z#StreamMux.aclose.<locals>.<genexpr>�  s]   è è € ð ð àÝ! !¥]Ñ3Ô3ðõ ' q­'Ô*@ÑAÔAð	Øðð ð ð ð ð r7   )Ú_collect_scheduled_tasksr‡   ÚgatherÚnextr   Ú	afinalizerx   r    ry   rz   r   )r0   ÚpendingÚresultsÚ	first_errr{   r4   r|   rH   s           r5   ÚaclosezStreamMux.aclose|  sI  è è € ð ×/Ò/Ñ1Ô1ˆØð 	 Ý#œN¨GÐLÀtÐLÐLÐLÐLÐLÐLÐLÐLˆGÝðð à$ðñ ô ð ñô ˆIð Ð$Ø�à,0ˆØÔ-ð 	$ð 	$ˆKð$Ø!×+Ò+Ñ-Ô-Ð-Ð-Ð-Ð-Ð-Ð-Ð-Ð-øÝ ð $ð $ð $ØÐ&Ø"#�Køøøøøøøøøð$øøøð ”.ð 	ð 	ˆBØ”:ð Ø—’‘
”
�
øØŒ×ÒÑÔÐØÐ"ØÐð #Ð"s   ÁA5Á5
BÁ?BÂBc              ƒ  óž  K  — |                       ¦   «         }|D ]}|                     ¦   «          Œ|rt          j        |ddiŽƒ d{V —† | j        D ].}	 |                     |¦  «        ƒ d{V —† Œ# t          $ r Y Œ+w xY w| j        D ]}|j        s| 	                    |¦  «         Œ| j
        j        s| j
         	                    |¦  «         dS dS )a&  Fail on the async lane.

        Cancels every scheduled task across all transformers, awaits
        them to completion, then runs each transformer's `afail` hook
        and auto-fails channels and the main event log.

        Args:
            err: The exception that ended the run.
        r„   TN)rŠ   Úcancelr‡   r‹   r   Úafailrx   r    ry   r   r   )r0   r}   rŽ   Útaskr4   rH   s         r5   r”   zStreamMux.afail¨  s%  è è € ð ×/Ò/Ñ1Ô1ˆØð 	ð 	ˆDØ�KŠK‰MŒMˆMˆMØð 	CÝ”. 'ÐB¸TÐBÐBÐBÐBÐBÐBÐBÐBÐBàÔ-ð 	ð 	ˆKðØ!×'Ò'¨Ñ,Ô,Ð,Ð,Ð,Ð,Ð,Ð,Ð,Ð,øÝ ð ð ð Ø�ðøøøà”.ð 	ð 	ˆBØ”:ð Ø—’˜‘”�øØŒ|Ô#ð 	#ØŒL×Ò˜cÑ"Ô"Ð"Ð"Ð"ð	#ð 	#s   ÁA/Á/
A<Á;A<úlist[asyncio.Task[Any]]c                ó$   — d„ | j         D ¦   «         S )z@Return a snapshot of in-flight tasks scheduled via transformers.c                ób   — g | ],}t          |d d¦  «        D ]}|                     ¦   «         °|‘ŒŒ-S )Ú_stream_scheduled_tasksr   )r,   Údone)rX   r4   r•   s      r5   ú
<listcomp>z6StreamMux._collect_scheduled_tasks.<locals>.<listcomp>Å  sZ   € ð 
ð 
ð 
àÝ Ð-FÈÑKÔKð
ð 
ð Ø—9’9‘;”;ð	
Øð
ð 
ð 
ð 
r7   )r   rA   s    r5   rŠ   z"StreamMux._collect_scheduled_tasksÃ  s&   € ð
ð 
à#Ô1ð
ñ 
ô 
ð 	
r7   r[   rj   údict[str, Any]r\   c               óp  ‡ — |                      ¦   «         D ]Ÿ}t          |t          ¦  «        rˆ|                     ‰ j        ¬¦  «         |                     ‰ ¦  «         ‰ j                             |¦  «         |j        �7|r|j        n	d|j        › �}d	ˆ fd„}| 	                     ||¦  «        ¦  «         Œ dS )
aI  Bind and optionally wire StreamChannel instances in a projection.

        All StreamChannels are bound and tracked. Channels with a name
        are additionally wired for protocol auto-forwarding.

        Args:
            projection: The projection dict returned by a transformer's
                `init()`.
            native: True when the owning transformer is `_native`.
                Named channels owned by a native transformer use the
                channel name directly as the protocol method;
                user-defined channels are prefixed with `custom:`.
        r   Nzcustom:Úmethod_namer9   r   úCallable[[Any], None]c                ó   •‡ — dˆ ˆfd„}|S )NÚitemr   r   r   c                ó4   •— ‰                      ‰| ¦  «         d S r†   )Ú_forward)r¡   rž   r0   s    €€r5   r£   zAStreamMux._bind_and_wire.<locals>._make_forward.<locals>._forwardé  s   ø€ Ø ŸMšM¨+°tÑ<Ô<Ð<Ð<Ð<r7   )r¡   r   r   r   r   )rž   r£   r0   s   ` €r5   Ú_make_forwardz/StreamMux._bind_and_wire.<locals>._make_forwardè  s.   øø€ ð=ð =ð =ð =ð =ð =ð =ð  (˜r7   )rž   r9   r   rŸ   )
Úvaluesr`   r
   r   r   r   r    r-   ÚnameÚ_wire)r0   rj   r\   ÚvalueÚmethodr¤   s   `     r5   rf   zStreamMux._bind_and_wireÐ  sÚ   ø€ ð   ×&Ò&Ñ(Ô(ð 	7ð 	7ˆEÝ˜%¥Ñ/Ô/ð 7Ø—’ T¤]�Ñ3Ô3Ð3Ø—’ Ñ%Ô%Ð%Ø”×%Ò% eÑ,Ô,Ð,Ø”:Ð)Ø+1ÐM˜UœZ˜ZÐ7MÀÄÐ7MÐ7M�Fð(ð (ð (ð (ð (ð (ð —K’K  ¨fÑ 5Ô 5Ñ6Ô6Ð6øð	7ð 	7r7   r©   r¡   r   c                óÒ   — d|g t          t          j        ¦   «         dz  ¦  «        |dœdœ}| j        r| xj        dz  c_        | j        |d<   | j                             |¦  «         dS )aŽ  Inject a ProtocolEvent for a StreamChannel push.

        Forwarded events bypass the transformer pipeline to avoid
        infinite recursion (a transformer that pushes to a channel
        during `process()` would re-trigger itself). These events are
        visible in this mux's main event log but are not passed through
        transformers' `process()` methods. Only the root mux assigns
        `seq` to forwarded channel events.

        Args:
            method: The full protocol method (already with or without
                the `custom:` prefix; resolved by `_bind_and_wire`).
            item: The payload pushed onto the channel.
        ro   iè  )Ú	namespaceÚ	timestampÚdata)r]   r©   Úparamsr@   rq   N)r>   Útimer   r!   r   rs   )r0   r©   r¡   ro   s       r5   r£   zStreamMux._forwardð  s   € ð  ØàÝ ¥¤¡¤¨tÑ!3Ñ4Ô4Øðð ð 
ð  
ˆð Ôð 	%ØˆIŒI˜‰NˆIŒIØœ9ˆE�%‰LØŒ×Ò˜%Ñ Ô Ð Ð Ð r7   r†   )r   r   r   r   r   r   r   r   r   r   r   r   )r8   r9   r   r:   )r   r>   )rC   rD   r   r   )rC   rK   r   r   )r   r   r   r   )r4   r   r   r   )ro   r   r   r   )r   r   )r}   rx   r   r   )r   r–   )rj   rœ   r\   r   r   r   )r©   r9   r¡   r   r   r   )r^   Ú
__module__Ú__qualname__Ú__doc__r6   r=   rB   rJ   rP   rT   r.   rs   rz   r   r‚   r‘   r”   rŠ   rf   r£   r   r7   r5   r   r      sœ  € € € € € ðð ð. 8<ðd(ð Ø59Ø!#Ø ðd(ð d(ð d(ð d(ð d(ð d(ðL1ð 1ð 1ð 1ðð ð ð ðð ð ð ð(	ð 	ð 	ð 	ð ð  ð  ð  ðD('ð ('ð ('ð ('ðT%ð %ð %ð %ð:ð ð ð ð8ð ð ð ð2%ð %ð %ð %ð:*ð *ð *ð *ðX#ð #ð #ð #ð6
ð 
ð 
ð 
ð =Bð7ð 7ð 7ð 7ð 7ð 7ð@!ð !ð !ð !ð !ð !r7   r   )Ú
__future__r   r‡   r¯   Úcollections.abcr   r   Útypingr   Úlanggraph.stream._typesr   r   r	   Úlanggraph.stream.stream_channelr
   ÚTransformerFactoryr   r   r7   r5   ú<module>r¹      sï   ðØ "Ð "Ð "Ð "Ð "Ð "à €€€Ø €€€Ø /Ð /Ð /Ð /Ð /Ð /Ð /Ð /Ø Ð Ð Ð Ð Ð ðð ð ð ð ð ð ð ð ð ð
 :Ð 9Ð 9Ð 9Ð 9Ð 9àÐ0Ð1Ð3DÐDÔEÐ ððq!ð q!ð q!ð q!ð q!ñ q!ô q!ð q!ð q!ð q!r7   