ó
    ýÞ j¦c  ã                  ób  • S SK Jr  S SKrS SKJrJrJrJrJr  S SK	J
r
Jr  S SKJrJr  S SKJr  S SKJr  S SKJr  S S	KJr  \(       a  S S
KJrJr  S SKJr  S SKJrJr  SS jrSS jr \" SS9 " S S5      5       r!\" SS9 " S S5      5       r" " S S5      r# " S S\!\#5      r$ " S S\"\#5      r%g)é    )ÚannotationsN)ÚAsyncIteratorÚ	AwaitableÚCallableÚIteratorÚMapping)ÚMappingProxyTypeÚTracebackType)ÚTYPE_CHECKINGÚAny)Úbeta)Úconvert_to_protocol_event)Ú	StreamMux)ÚProtocolEvent)ÚAsyncChatModelStreamÚChatModelStream©ÚStreamChannel)ÚLifecyclePayloadÚSubgraphStatusc                ó<   • U " 5       (       a   U " 5       (       a  M  gg)z*Call the sync pump until it returns False.N© ©Úpumps    ÚU/var/www/html/gaurav/venv/lib/python3.13/site-packages/langgraph/stream/run_stream.pyÚ_drive_until_doner      s   € á
�&‰&Øñ �&�&ó    c              ƒ  ól   #   • U " 5       I Sh  v•N (       a   U " 5       I Sh  v•N (       a  M  gg N" N7f)z+Call the async pump until it returns False.Nr   r   s    r   Ú_adrive_until_doner      s   é € á“�,�,Øñ “�,�,�,ùs   ‚4�0Ž4¢2£	4®4²4z4The v3 streaming protocol on Pregel is experimental.)Úmessagec                  óþ   • \ rS rSr% SrS\S'   S\S'   S\S'   S	\S
'   SS.       SS jjrSS jrSS jrSS jr	SS jr
S S jr        S!S jr\S"S j5       r\SS j5       r\S#S j5       rS$S jrS%S jrSrg)&ÚGraphRunStreamé$   u�  Sync run stream with caller-driven pumping.

The caller's iteration on any projection (`values`, `messages`,
raw events, or `output`) drives the graph forward. No background
thread is used â€” the caller's `for` loop is the pump.

Projections are single-consumer â€” iterating `run.values` twice
raises. Use `projection.tee(n)` if you genuinely need fan-out.

All transformer projections live in `extensions`. Native transformer
projections (those with `_native = True`) are also set as direct
attributes on this instance (e.g. `run.values`, `run.messages`).

!!! warning

    Returned by `Pregel.stream_events(version="v3")`, which is
    experimental and may change.
úStreamChannel[dict[str, Any]]ÚvalueszStreamChannel[ChatModelStream]ÚmessagesúStreamChannel[LifecyclePayload]Ú	lifecyclez StreamChannel[SubgraphRunStream]Ú	subgraphsT©Ú	wire_pumpc               óF  • Xl         X l        [        UR                  5      U l        SU l        SU l        SU l        / U l        [        UR                  5      U l
        UR                   H  n[        XUR                  U   5        M     U(       a  U R                  U5        gg)aë  Initialize the run stream.

Args:
    graph_iter: Pull-based iterator over the graph's stream,
        or `None` for nested run streams whose pump is driven
        by an outer run (e.g. `SubgraphRunStream`).
    mux: The StreamMux owning projections and the main log.
    wire_pump: When True (default), bind `_pump_next` as the
        mux's pump callable. Subclasses that inherit a parent
        pump via `StreamMux._make_child` should pass False to
        preserve the parent binding.
FN)Ú_graph_iterÚ_muxr	   Ú
extensionsÚ
_exhaustedÚ_latestÚ_interruptedÚ_interruptsÚlistÚscopeÚ_scope_listÚnative_keysÚsetattrÚ_wire_request_more)ÚselfÚ
graph_iterÚmuxr+   Úkeys        r   Ú__init__ÚGraphRunStream.__init__C   s„   € ð& &ÔØŒ	Ü-=¸c¿n¹nÓ-MˆŒØˆŒØ.2ˆŒØ!ˆÔØ&(ˆÔÜ&*¨3¯9©9£oˆÔØ—?”?ˆCÜ�D˜sŸ~™~¨cÑ2Ö3ñ #æØ×#Ñ# CÕ(ð r   c                ó:   • UR                  U R                  5        g)a<  Wire the sync pull callback through the mux.

Routing through `mux.bind_pump` (rather than walking
projections directly here) lets child mini-muxes built by
`mux._make_child(...)` inherit the same pump callable, so
cursors on a subgraph handle's projections drive the root
pump just like cursors on `run.values` do.
N)Ú	bind_pumpÚ
_pump_next©r:   r<   s     r   r9   Ú!GraphRunStream._wire_request_morec   s   € ð 	�‰�d—o‘oÕ&r   c                óÖ   • US   S:w  a  gUS   nUS   U R                   :w  a  gUS   U l        UR                  SS5      nU(       a#  S	U l        U R                  R                  U5        gg©
z;Track values-event state for output/interrupted/interrupts.Úmethodr%   NÚparamsÚ	namespaceÚdataÚ
interruptsr   T©r6   r1   Úgetr2   r3   Úextend©r:   ÚeventrH   rK   s       r   Ú_observe_eventÚGraphRunStream._observe_eventn   óo   € à�‰?˜hÓ&ØØ�x‘ˆØ�+Ñ $×"2Ñ"2Ó2ØØ˜f‘~ˆŒØ—Z‘Z ¨bÓ1ˆ
ÞØ $ˆDÔØ×Ñ×#Ñ# JÕ/ð r   c                ó¦  • U R                   (       d  U R                  c  g [        U R                  5      n[        U5      nU R	                  U5        U R
                  R                  U5        g! [         a$    U R
                  R                  5         SU l          g[         a,  nU R
                  R                  U5        SU l          SnAgSnAff = f)zøPull one event from the graph and push it through the mux.

Returns:
    True if an event was pulled, False if the graph is exhausted
    or has raised. Always False when constructed with
    `graph_iter=None` (the run is driven by an outer pump).
NFT)r0   r-   Únextr   rQ   r.   ÚpushÚStopIterationÚcloseÚ	ExceptionÚfail©r:   ÚpartrP   Úes       r   rB   ÚGraphRunStream._pump_next{   s£   € ð �?�?˜d×.Ñ.Ñ6Øð	Ü˜×(Ñ(Ó)ˆDÜ-¨dÓ3ˆEØ×Ñ Ô&Ø�I‰I�N‰N˜5Ô!ØøÜó 	Ø�I‰I�O‰OÔØ"ˆDŒOÙÜó 	Ø�I‰I�N‰N˜1ÔØ"ˆDŒOÜûð	ús   ¡AA. Á.+CÂ	CÂ$"CÃCc                ó  • U R                   (       a  gSU l         U R                  nSU l        Ub  [        USS5      =nb   U" 5          U R                  R                  5         g! [         a     N(f = f! [         a     gf = f)zÄStop the run early.

Closes the underlying graph iterator (propagating `GeneratorExit`
so in-flight nodes and subgraphs are cancelled), closes the mux,
and marks the stream exhausted. Idempotent.
NTrX   )r0   r-   ÚgetattrrY   r.   rX   )r:   r;   rX   s      r   ÚabortÚGraphRunStream.abort”   s†   € ð �?�?ØØˆŒØ×%Ñ%ˆ
ØˆÔàÑ"Ü! *¨g°tÓ<Ð<�ÑIðÙ”ð	Ø�I‰I�O‰OÕøô ó Ùðûô ó 	Ùð	ús$   ÁA$ Á	A4 Á$
A1Á0A1Á4
BÂ Bc                ó   • U $ ©Nr   ©r:   s    r   Ú	__enter__ÚGraphRunStream.__enter__­   s   € Øˆr   c                ó$   • U R                  5         g rd   ©ra   ©r:   Úexc_typeÚexcÚtbs       r   Ú__exit__ÚGraphRunStream.__exit__°   s   € ð 	�
‰
�r   c                óŽ   • [        U R                  5        U R                  R                  R                  =nb  UeU R
                  $ )z7Drive the run to completion and return the final state.)r   rB   r.   Ú_eventsÚ_errorr1   ©r:   Úerrs     r   ÚoutputÚGraphRunStream.output¸   s:   € ô 	˜$Ÿ/™/Ô*Ø—9‘9×$Ñ$×+Ñ+Ð+ˆCÑ8ØˆIØ�|‰|Ðr   c                óŽ   • [        U R                  5        U R                  R                  R                  =nb  UeU R
                  $ )z�Drive the run to completion, then return whether it was
interrupted.

Raises:
    BaseException: If the run ended with an error.
)r   rB   r.   rq   rr   r2   rs   s     r   ÚinterruptedÚGraphRunStream.interruptedÀ   s<   € ô 	˜$Ÿ/™/Ô*Ø—9‘9×$Ñ$×+Ñ+Ð+ˆCÑ8ØˆIØ× Ñ Ð r   c                óŽ   • [        U R                  5        U R                  R                  R                  =nb  UeU R
                  $ )zyDrive the run to completion, then return interrupt payloads.

Raises:
    BaseException: If the run ended with an error.
)r   rB   r.   rq   rr   r3   rs   s     r   rK   ÚGraphRunStream.interruptsÍ   s<   € ô 	˜$Ÿ/™/Ô*Ø—9‘9×$Ñ$×+Ñ+Ð+ˆCÑ8ØˆIØ×ÑÐr   c                ó@   • [        U R                  R                  5      $ ©z<Subscribe to the main event log and iterate protocol events.)Úiterr.   rq   re   s    r   Ú__iter__ÚGraphRunStream.__iter__Ù   s   € ä�D—I‘I×%Ñ%Ó&Ð&r   c              '  ó�  #   • SSK Jn  0 n U H±  nU R                  U   n[        XR5      (       d%  [	        S[        U5      R                   SU< 35      eUR                  c  [	        SU< S35      eUR                  (       a  [	        SU< S35      eUR                  (       a  [        SU< S	35      eS
Ul        XSU'   M³     [        5       n[        U5      [        U5      :  Gaž  SnUR                  5        H�  u  pEXF;   a  M  UR                  (       a=  UR                  (       d,  UR                  b  UR                  eUR!                  U5        MZ  UR                  (       d  Mm  UR                  S   S   nUb
  X‡S   :  d  MŒ  X„4nM‘     Ub+  X7S      R                  R#                  5       u  pšUS   U
4v •  O°U R$                  R&                  nUb  U" 5       (       d‹  [        U5      nUR                  5        H\  u  pEXF;  d  M  UR                  (       a  M  UR                  (       d  M2  UR                  b  UR                  eUR!                  U5        M^     [        U5      U:X  a  O[        U5      [        U5      :  a  GMž  UR)                  5        H
  nSUl        M     g! UR)                  5        H
  nSUl        M     f = f7f)u—  Iterate multiple projections in arrival order, yielding ``(name, item)``.

Items are ordered by a monotonic push stamp assigned when each
transformer pushes into its `StreamChannel`. This gives strict
arrival ordering across projections, unlike round-robin.

Args:
    *names: Projection keys to interleave. Must match keys in
        ``extensions``.

Yields:
    ``(name, item)`` tuples in arrival order across the named
    projections.

Each named channel is locked for the duration of iteration and
released when the generator completes, is closed, or raises.
Channels cannot be subscribed concurrently â€” use `.tee(n)` if
you need fan-out.

Raises:
    KeyError: If a name doesn't match a registered projection.

Example:
    ```python
    for name, item in run.interleave("messages", "values"):
        if name == "messages":
            print("msg:", item)
        else:
            print("val:", item)
    ```
r   r   z5interleave() requires StreamChannel projections, got z for NzStreamChannel zI has not been bound yet. Register the transformer with a StreamMux first.uL    is bound to async mode â€” sync interleave() cannot consume async channels.z3 already has a subscriber; use .tee(n) for fan-out.Té   F)Úlanggraph.stream.stream_channelr   r/   Ú
isinstanceÚ	TypeErrorÚtypeÚ__name__Ú	_is_asyncÚ_subscribedÚRuntimeErrorÚsetÚlenÚitemsÚ_closedÚ_itemsrr   ÚaddÚpopleftr.   Ú_pump_fnr%   )r:   Únamesr   ÚchannelsÚnameÚchÚdoneÚbestÚstampÚ_stampÚitemr   Úbefores                r   Ú
interleaveÚGraphRunStream.interleaveÝ   s{  é € õ@ 	Bà24ˆð<	'Û�Ø—_‘_ TÑ*�Ü! "×4Ñ4Ü#ðÜ# B›x×0Ñ0Ð1°°t±hð@óð ð —<‘<Ñ'Ü#Ø(¨©ð 1Kð Kóð ð —<—<Ü#Ø(¨©ð 1Kð Kóð ð —>—>Ü&Ø(¨©ð 13ð 3óð ð "&�”Ø!#˜“ñ/ ô2 !›UˆDä�d“)œc (›mÔ+Ø/3�Ø (§¡Ö 0‘H�DØ“|Ù Ø—z—z¨"¯)¯)ØŸ9™9Ñ0Ø"$§)¡)˜OØŸ™ œÙ Ø—y—y‘yØ "§	¡	¨!¡¨Q¡˜Ø™<¨5¸±7­?Ø$) =šDñ !1ð Ñ#Ø#+°©GÑ#4×#;Ñ#;×#CÑ#CÓ#E‘L�FØ ™7 D˜/Ó)àŸ9™9×-Ñ-�DØ‘|©4¯6©6Ü!$ T£˜Ø(0¯©Ö(8™H˜DØ#Õ/¸¿	¿	¹	Ø#%§:§:¡:Ø')§y¡yÑ'<Ø.0¯i©i¨Ø$(§H¡H¨T¦Nñ )9ô ˜t›9¨Ó.Ø!ô; �d“)œc (›mÖ+ð> —o‘oÖ'�Ø!&�–ò (ø�h—o‘oÖ'�Ø!&�–ò (üs=   ‚	KŒEJ# Å%J# ÆA=J# ÈJ# ÈJ# È+AJ# ÊKÊ# KËK)r0   r-   r2   r3   r1   r.   r6   r/   N)r;   zIterator[Any] | Noner<   r   r+   ÚboolÚreturnÚNone©r<   r   r    r¡   ©rP   r   r    r¡   ©r    rŸ   ©r    r¡   )r    r"   ©rk   ztype[BaseException] | Nonerl   zBaseException | Nonerm   zTracebackType | Noner    r¡   ©r    zdict[str, Any] | None©r    z	list[Any])r    zIterator[ProtocolEvent])r“   Ústrr    zIterator[tuple[str, Any]])r‡   Ú
__module__Ú__qualname__Ú__firstlineno__Ú__doc__Ú__annotations__r>   r9   rQ   rB   ra   rf   rn   Úpropertyru   rx   rK   r   r�   Ú__static_attributes__r   r   r   r"   r"   $   sÝ   ‡ ñð0 *Ó)Ø,Ó,Ø.Ó.Ø/Ó/ð ñ)à(ð)ð ð)ð
 ð)ð 
õ)ô@	'ô0ôô2ô2ðà,ðð "ðð !ð	ð
 
ôð óó ðð ó
!ó ð
!ð ó	 ó ð	 ô'÷_'r   r"   c                  óÖ   • \ rS rSr% SrS\S'   S\S'   S\S'   S	\S
'   SS.       SS jjrSS jrSS jrSS jr	SS jr
SS jr        S S jrS!S jrSS jrS"S jrS#S jrSrg)$ÚAsyncGraphRunStreami?  uS  Async run stream with caller-driven pumping.

Async iteration on any projection drives the graph forward â€” there
is no background task. Concurrent consumers share a single-flight
pump via an `asyncio.Lock`, so each awaiting cursor contributes one
event per acquisition. Backpressure comes from the logs: when a
subscribed log's buffer reaches `maxlen`, `apush` awaits the
subscriber to drain, which holds back the pump and paces the graph.

Projections are single-consumer â€” a second `aiter(run.values)`
raises. Use `projection.tee(n)` for fan-out.

Use as an async context manager to guarantee clean shutdown on
early exit:

```python
async with await handler.astream(input) as run:
    async for msg in run.messages:
        ...
```

!!! warning

    Awaited from `Pregel.astream_events(version="v3")`, which is
    experimental and may change.
r$   r%   z#StreamChannel[AsyncChatModelStream]r&   r'   r(   z%StreamChannel[AsyncSubgraphRunStream]r)   Tr*   c               ó¤  • Xl         X l        [        UR                  5      U l        SU l        SU l        SU l        / U l        [        UR                  5      U l
        [        R                  " 5       U l        SU l        SU l        SU l        UR"                   H  n[%        XUR                  U   5        M     U(       a  U R'                  U5        gg)aù  Initialize the async run stream.

Args:
    graph_aiter: Async iterator over the graph's stream, or
        `None` for nested run streams whose pump is driven by
        an outer run (e.g. `AsyncSubgraphRunStream`).
    mux: The StreamMux owning projections and the main log.
    wire_pump: When True (default), bind `_apump_next` as the
        mux's async pump callable. Subclasses that inherit a
        parent pump via `StreamMux._make_child` should pass
        False to preserve the parent binding.
FN)Ú_graph_aiterr.   r	   r/   r0   r1   r2   r3   r4   r5   r6   ÚasyncioÚ	ConditionÚ
_pump_condÚ_pumpingÚ_anext_taskÚ	_abortingr7   r8   Ú_wire_arequest_more)r:   Úgraph_aiterr<   r+   r=   s        r   r>   ÚAsyncGraphRunStream.__init__f  sª   € ð& (ÔØŒ	Ü-=¸c¿n¹nÓ-MˆŒØˆŒØ.2ˆŒØ!ˆÔØ&(ˆÔÜ&*¨3¯9©9£oˆÔÜ!×+Ò+Ó-ˆŒØˆŒØ7;ˆÔØˆŒØ—?”?ˆCÜ�D˜sŸ~™~¨cÑ2Ö3ñ #æØ×$Ñ$ SÕ)ð r   c                óÖ   • US   S:w  a  gUS   nUS   U R                   :w  a  gUS   U l        UR                  SS5      nU(       a#  S	U l        U R                  R                  U5        ggrF   rL   rO   s       r   rQ   Ú"AsyncGraphRunStream._observe_eventŠ  rS   r   c                ó:   • UR                  U R                  5        g)zÒWire the async pull callback through the mux.

Mirrors `_wire_request_more`: routing through
`mux.bind_apump` lets child mini-muxes inherit the pump
callable so cursors on subgraph handles drive the root
pump.
N)Ú
bind_apumpÚ_apump_nextrC   s     r   r»   Ú'AsyncGraphRunStream._wire_arequest_more—  s   € ð 	�‰�t×'Ñ'Õ(r   c           	   ƒ  ó|  #   • U R                    ISh  v•N   U R                  (       d  U R                  c   SSS5      ISh  v•N   gU R                  (       aD  U R                   R	                  5       I Sh  v•N   U R                  (       + sSSS5      ISh  v•N   $ SU l        SSS5      ISh  v•N     [
        R                  " U R                  R                  5       5      U l         U R                  I Sh  v•N n SU l        [        U5      nU R                  U5        U R                  R!                  U5      I Sh  v•N    U R                    ISh  v•N   SU l        U R                   R                  5         SSS5      ISh  v•N   g GNz GNN GN Nþ Nè! , ISh  v•N  (       d  f       Ný= f Nº! [
        R                   aŸ    U R                  (       ar  SU l         SU l        U R                    ISh  v•N    SU l        U R                   R                  5         SSS5      ISh  v•N    g! , ISh  v•N  (       d  f       g= fU R                  R                  5         e f = f! SU l        f = f GN8 GN% Nö! , ISh  v•N  (       d  f       g= f! ["         a�    SU l        U R                  R%                  5       I Sh  v•N     U R                    ISh  v•N    SU l        U R                   R                  5         SSS5      ISh  v•N    g! , ISh  v•N  (       d  f       g= f[&         a—  nSU l        U R                  R)                  U5      I Sh  v•N     SnAU R                    ISh  v•N    SU l        U R                   R                  5         SSS5      ISh  v•N    g! , ISh  v•N  (       d  f       g= fSnAff = f! U R                    ISh  v•N    SU l        U R                   R                  5         SSS5      ISh  v•N    f ! , ISh  v•N  (       d  f       f = f= f7f)u  Drive one pump step, or wait for the active pumper to drive one.

"Take-a-number" semantics: at most one task at a time calls
`graph_aiter.__anext__()` (asyncio iterators can't be advanced
concurrently). Other callers wait on a Condition that the
active pumper notifies after each step. This lets a "passive"
consumer â€” one whose projection's buffer is being filled by the
active pumper's push â€” wake up as soon as its data lands,
instead of queueing on the pump and only observing its data one
graph event late.

`except Exception` is intentional â€” `CancelledError` and other
`BaseException` subclasses propagate, matching asyncio's
cancellation contract.

Returns:
    True if a pump step completed (by this task or another),
    False if the graph is exhausted.
NFT)r·   r0   r´   r¸   Úwaitrµ   Úensure_futureÚ	__anext__r¹   ÚCancelledErrorrº   Ú
notify_allÚcancelr   rQ   r.   ÚapushÚStopAsyncIterationÚacloserY   Úafailr[   s       r   rÂ   ÚAsyncGraphRunStream._apump_next¡  sµ  é € ð( —?—?“?Ø�� $×"3Ñ"3Ñ";Ø÷ #—?�?ð �}�}à—o‘o×*Ñ*Ó,×,Ð,ØŸ?™?Ô*÷ #—?‘?ð !ˆDŒM÷ #—?ð$	-ðô $+×#8Ò#8¸×9JÑ9J×9TÑ9TÓ9VÓ#W�Ô ð,Ø!%×!1Ñ!1×1‘Dð (,�DÔ$Ü1°$Ó7�Ø×#Ñ# EÔ*Ø—i‘i—o‘o eÓ,×,Ð,Øð ——“Ø %�”Ø—‘×*Ñ*Ô,÷ '—�õW #ò
 -÷ #—?—?’?úñ$ 2øÜ×-Ñ-ó Ø—~—~à*.˜œØ$ð (,�DÔ$ð ——”Ø %�”Ø—‘×*Ñ*Ô,÷ '————�úð# ×$Ñ$×+Ñ+Ô-Øðûð (,�DÕ$úò -÷ '——“ûô &ó Ø"&�”Ø—i‘i×&Ñ&Ó(×(Ñ(Øð ——”Ø %�”Ø—‘×*Ñ*Ô,÷ '————�úô ó Ø"&�”Ø—i‘i—o‘o aÓ(×(Ñ(Ûà——”Ø %�”Ø—‘×*Ñ*Ô,÷ '————�þðûð
 ——”Ø %�”Ø—‘×*Ñ*Ô,÷ '————�ÿsŠ  ‚P<“F”P<— F·P<ÁFÁP<Á.FÁ6FÁ7FÂP<ÂFÂP<ÂFÂ$P<Â/FÂ0P<Â63J Ã*F6 Ã9F4Ã:F6 Ã?AJ Å I8ÅJ ÅP<ÅI;ÅP<Å"J Å<P<ÆI>ÆP<ÆP<ÆFÆP<ÆP<ÆF1Æ F#Æ!F1Æ-P<Æ4F6 Æ6,I)Ç"I, Ç#J Ç*P<Ç:G=Ç;P<Ç?"H3È!P<È,H/È-P<È3I
È9H<È:I
ÉP<ÉI)É)I, É,	I5É5J É;P<É>P<Ê JÊJ	ÊJÊP<Ê.OËKË	OËO ËP<ËK"Ë P<Ë$"LÌP<ÌLÌP<ÌL/ÌL!ÌL/Ì+P<Ì2	OÌ;%OÍ M#Í!OÍ&O Í*P<Í:M=Í;P<Í?"N3Î!P<Î,N/Î-P<Î3O
Î9N<Î:O
ÏP<ÏOÏO ÏP9Ï&O)
Ï'P9Ï+"PÐP9ÐPÐP9ÐP6Ð%P(Ð&P6Ð2P9Ð9P<c              ƒ  óè  #   • U R                    ISh  v•N   U R                  (       a   SSS5      ISh  v•N   gSU l        SU l        U R                  nSU l        U R                  nU R                   R                  5         SSS5      ISh  v•N   Wb0  UR                  5       (       d  UR                  5          UI Sh  v•N   Wb   [        USS5      =nb   U" 5       I Sh  v•N    U R                  R                  5       I Sh  v•N   g Nø NØ N‚! , ISh  v•N  (       d  f       N—= f Nk! [        R                  [        4 a     N…f = f Nj! [         a     Ntf = f NY! [         a     gf = f7f)a`  Stop the run early.

Marks the stream exhausted and wakes any pump-waiters. Cancels an
in-flight pull if one is running, then closes the underlying graph
iterator, so running nodes and nested subgraphs are cancelled
whether or not a pump is mid-pull. Closes the mux; any `apush`
blocked on backpressure wakes and returns without appending.
Idempotent.
NTrÍ   )r·   r0   rº   r´   r¹   rÉ   r—   rÊ   rµ   rÈ   rY   r`   r.   rÍ   )r:   r¼   Ú
anext_taskrÍ   s       r   ra   ÚAsyncGraphRunStream.abortä  s6  é € ð —?—?“?Ø��Ø÷ #—?�?ð #ˆDŒOØ!ˆDŒNØ×+Ñ+ˆKØ $ˆDÔØ×)Ñ)ˆJØ�O‰O×&Ñ&Ô(÷ #—?ð Ñ!¨*¯/©/×*;Ñ*;Ø×ÑÔðØ × Ð ð Ñ#Ü" ;°¸$Ó?Ð?�ÑLðÙ“h—�ð	Ø—)‘)×"Ñ"Ó$×$Ñ$÷9 #—?—?”?úñ  !øÜ×*Ñ*¬IÐ6ó Ùðúñ øÜó Ùðúñ %øÜó 	Ùð	üsû   ‚E2“D”E2—DªE2µD¶E2»ADÂE2ÂDÂ,E2Â;D. Ã D,ÃD. ÃE2Ã
E Ã#EÃ$E Ã)E" ÄE ÄE" ÄE2ÄE2ÄE2ÄD)ÄDÄD)Ä%E2Ä,D. Ä.EÅE2Å
EÅE2ÅE Å
EÅE2ÅEÅE2Å E" Å"
E/Å,E2Å.E/Å/E2c              ƒ  ó   #   • U $ 7frd   r   re   s    r   Ú
__aenter__ÚAsyncGraphRunStream.__aenter__  s
   é € Øˆùs   ‚c              ƒ  ó@   #   • U R                  5       I S h  v•N   g  N7frd   ri   rj   s       r   Ú	__aexit__ÚAsyncGraphRunStream.__aexit__  s   é € ð �j‰j‹l×Óùs   ‚–—c              ƒ  óª   #   • [        U R                  5      I Sh  v•N   U R                  R                  R                  =nb  UeU R
                  $  N57f)aK  Drive the run to completion and return the final state.

Methods (not properties) on the async lane so `run.output`
without `await` raises at type-check time instead of silently
yielding a coroutine object.

Example:
    ```python
    output = await run.output()
    ```

Raises:
    BaseException: If the run ended with an error.
N)r   rÂ   r.   rq   rr   r1   rs   s     r   ru   ÚAsyncGraphRunStream.output  sJ   é € ô ! ×!1Ñ!1Ó2×2Ð2Ø—9‘9×$Ñ$×+Ñ+Ð+ˆCÑ8ØˆIØ�|‰|Ðñ 	3ùó   ‚A›Aœ6Ac              ƒ  óª   #   • [        U R                  5      I Sh  v•N   U R                  R                  R                  =nb  UeU R
                  $  N57f)zDrive the run to completion and return whether it was
interrupted.

Raises:
    BaseException: If the run ended with an error.
N)r   rÂ   r.   rq   rr   r2   rs   s     r   rx   ÚAsyncGraphRunStream.interrupted-  sL   é € ô ! ×!1Ñ!1Ó2×2Ð2Ø—9‘9×$Ñ$×+Ñ+Ð+ˆCÑ8ØˆIØ× Ñ Ð ñ 	3ùrÛ   c              ƒ  óª   #   • [        U R                  5      I Sh  v•N   U R                  R                  R                  =nb  UeU R
                  $  N57f)zwDrive the run to completion and return interrupt payloads.

Raises:
    BaseException: If the run ended with an error.
N)r   rÂ   r.   rq   rr   r3   rs   s     r   rK   ÚAsyncGraphRunStream.interrupts9  sL   é € ô ! ×!1Ñ!1Ó2×2Ð2Ø—9‘9×$Ñ$×+Ñ+Ð+ˆCÑ8ØˆIØ×ÑÐñ 	3ùrÛ   c                óJ   • U R                   R                  R                  5       $ r}   )r.   rq   Ú	__aiter__re   s    r   rá   ÚAsyncGraphRunStream.__aiter__D  s   € à�y‰y× Ñ ×*Ñ*Ó,Ð,r   )rº   r¹   r0   r´   r2   r3   r1   r.   r·   r¸   r6   r/   N)r¼   zAsyncIterator[Any] | Noner<   r   r+   rŸ   r    r¡   r£   r¢   r¤   r¥   )r    r²   r¦   r§   r¨   )r    zAsyncIterator[ProtocolEvent])r‡   rª   r«   r¬   r­   r®   r>   rQ   r»   rÂ   ra   rÔ   r×   ru   rx   rK   rá   r°   r   r   r   r²   r²   ?  s®   ‡ ñð@ *Ó)Ø1Ó1Ø.Ó.Ø4Ó4ð ñ"*à.ð"*ð ð"*ð
 ð"*ð 
õ"*ôH0ô)ôA-ôF(ôTðà,ðð "ðð !ð	ð
 
ôôô(
!ô	 ÷-r   r²   c                  óV   • \ rS rSr% SrS\S'   S\S'   S\S'   S\S	'   S\S
'   S\S'   Srg)Ú_SubgraphRunStreamMixiniI  u5  Subgraph metadata + parent-pump delegation shared by both lanes.

Inherits from `GraphRunStream` (or `AsyncGraphRunStream`) with
`graph_iter=None` + `wire_pump=False` â€” the mini-mux is driven
by the parent's pump (inherited via `StreamMux._make_child`), and
the handle never pulls upstream itself. Pump-driving methods
delegate to the parent pump so `handle.output` and friends drive
the root run.

Subclasses set the parent pump function captured at construction
(`_parent_pump_fn` / `_parent_apump_fn`) and override
`_pump_next` / `_apump_next` to delegate to it.

Status is updated in place by `SubgraphTransformer`. Iterate
`run.subgraphs` to receive handles as subgraphs spawn, then
drill into projections inside the loop body **before** the next
pump cycle â€” same lazy-subscribe constraint as root projections.
útuple[str, ...]Úpathú
str | NoneÚ
graph_nameÚtrigger_call_idr   ÚstatusÚerrorrŸ   Ú_seen_terminalr   N)r‡   rª   r«   r¬   r­   r®   r°   r   r   r   rä   rä   I  s-   ‡ ñð& ÓØÓØÓØÓØÓØÖr   rä   c                  óV   ^ • \ rS rSrSrSSS.         SU 4S jjjrS	S jrSrU =r$ )
ÚSubgraphRunStreamie  zASync handle for a discovered subgraph (extends `GraphRunStream`).N©rè   ré   c               ó”   >• UR                   U l        [        TU ]  S USS9  X l        X0l        X@l        SU l        S U l        SU l	        g )NF)r;   r<   r+   Ústarted)
r’   Ú_parent_pump_fnÚsuperr>   ræ   rè   ré   rê   rë   rì   ©r:   r<   ræ   rè   ré   Ú	__class__s        €r   r>   ÚSubgraphRunStream.__init__h  sT   ø€ ð ;>¿,¹,ˆÔÜ‰ÑØØØð 	ñ 	
ð
 Œ	Ø$ŒØ.ÔØˆŒØˆŒ
Ø#ˆÕr   c                óÌ   • U R                   (       dC  U R                  (       d2  U R                  R                  R                  (       d  U R
                  c  gU R                  5       $ )zÂDelegate to the parent's pump.

Cursors on this handle's projections call here when their
buffers empty. Driving the parent fans events into our
mini-mux, transparently advancing the whole run.
F)r0   rì   r.   rq   rŽ   rò   re   s    r   rB   ÚSubgraphRunStream._pump_next  sE   € ð �O�OØ×"×"Ø�y‰y× Ñ ×(×(Ø×#Ñ#Ñ+àØ×#Ñ#Ó%Ð%r   )rò   rì   rë   rè   ræ   rê   ré   ©
r<   r   ræ   rå   rè   rç   ré   rç   r    r¡   r¤   )	r‡   rª   r«   r¬   r­   r>   rB   r°   Ú__classcell__©rõ   s   @r   rî   rî   e  sR   ø† ÙKð "&Ø&*ñ$àð$ð ð	$ð
 ð$ð $ð$ð 
÷$ð $÷.&ò &r   rî   c                  óV   ^ • \ rS rSrSrSSS.         SU 4S jjjrS	S jrSrU =r$ )
ÚAsyncSubgraphRunStreami�  zGAsync handle for a discovered subgraph (extends `AsyncGraphRunStream`).Nrï   c               ó”   >• UR                   U l        [        TU ]  S USS9  X l        X0l        X@l        SU l        S U l        SU l	        g )NF)r¼   r<   r+   rñ   )
Ú	_apump_fnÚ_parent_apump_fnró   r>   ræ   rè   ré   rê   rë   rì   rô   s        €r   r>   ÚAsyncSubgraphRunStream.__init__“  sV   ø€ ð GJÇmÁmˆÔÜ‰ÑØØØð 	ñ 	
ð
 Œ	Ø$ŒØ.ÔØˆŒØˆŒ
Ø#ˆÕr   c              ƒ  óè   #   • U R                   (       dC  U R                  (       d2  U R                  R                  R                  (       d  U R
                  c  gU R                  5       I Sh  v•N $  N7f)z$Delegate to the parent's async pump.NF)r0   rì   r.   rq   rŽ   r   re   s    r   rÂ   Ú"AsyncSubgraphRunStream._apump_next¨  sN   é € ð �O�OØ×"×"Ø�y‰y× Ñ ×(×(Ø×$Ñ$Ñ,àØ×*Ñ*Ó,×,Ð,Ñ,ùs   ‚A)A2Á+A0Á,A2)r   rì   rë   rè   ræ   rê   ré   rù   r¤   )	r‡   rª   r«   r¬   r­   r>   rÂ   r°   rú   rû   s   @r   rý   rý   �  sR   ø† ÙQð "&Ø&*ñ$àð$ð ð	$ð
 ð$ð $ð$ð 
÷$ð $÷*	-ò 	-r   rý   )r   zCallable[[], bool]r    r¡   )r   zCallable[[], Awaitable[bool]]r    r¡   )&Ú
__future__r   rµ   Úcollections.abcr   r   r   r   r   Útypesr	   r
   Útypingr   r   Úlangchain_core._apir   Úlanggraph.stream._convertr   Úlanggraph.stream._muxr   Úlanggraph.stream._typesr   Ú0langchain_core.language_models.chat_model_streamr   r   rƒ   r   Úlanggraph.stream.transformersr   r   r   r   r"   r²   rä   rî   rý   r   r   r   Ú<module>r     s¯   ðÝ "ã ß QÕ Qß 1ß %å $å ?Ý +Ý 1æ÷õ
 >ßNôôñ ÐDÑE÷W'ð W'ó FðW'ñt ÐDÑE÷F-ð F-ó FðF-÷Rñ ô8(&˜Ð(?ô (&ôV!-Ð0Ð2Iõ !-r   