§
    šŠtj¦c  ã                  ó˜  — d dl mZ d dlZd dlmZmZmZmZmZ d dl	m
Z
mZ d dlmZmZ d dlmZ d dlmZ d dlmZ d d	lmZ erd d
lmZmZ d dlmZ d dlmZmZ d d„Zd!d„Z  ed¬¦  «         G d„ d¦  «        ¦   «         Z! ed¬¦  «         G d„ d¦  «        ¦   «         Z" G d„ d¦  «        Z# G d„ de!e#¦  «        Z$ G d„ de"e#¦  «        Z%dS )"é    )ÚannotationsN)ÚAsyncIteratorÚ	AwaitableÚCallableÚIteratorÚMapping)ÚMappingProxyTypeÚTracebackType)ÚTYPE_CHECKINGÚAny)Úbeta)Úconvert_to_protocol_event)Ú	StreamMux)ÚProtocolEvent)ÚAsyncChatModelStreamÚChatModelStream©ÚStreamChannel)ÚLifecyclePayloadÚSubgraphStatusÚpumpúCallable[[], bool]ÚreturnÚNonec                ó4   —  | ¦   «         r	  | ¦   «         °dS dS )z*Call the sync pump until it returns False.N© ©r   s    úY/var/www/html/CA-Chatbot/venv/lib/python3.11/site-packages/langgraph/stream/run_stream.pyÚ_drive_until_doner      s7   € à
ˆ$‰&Œ&ð Øð ˆ$‰&Œ&ð ð ð ð ð ó    úCallable[[], Awaitable[bool]]c              ƒ  óP   K  —  | ¦   «         ƒ d{V —†r	  | ¦   «         ƒ d{V —†°dS dS )z+Call the async pump until it returns False.Nr   r   s    r   Ú_adrive_until_doner#      sS   è è € à�‘”ˆ,ˆ,ˆ,ˆ,ˆ,ˆ,ð Øð �‘”ˆ,ˆ,ˆ,ˆ,ˆ,ˆ,ð ð ð ð ð r    z4The v3 streaming protocol on Pregel is experimental.)Úmessagec                  óÒ   — e Zd ZU dZded<   ded<   ded<   ded	<   d
dœd/d„Zd0d„Zd1d„Zd2d„Zd3d„Z	d4d„Z
d5d"„Zed6d$„¦   «         Zed2d%„¦   «         Zed7d'„¦   «         Zd8d)„Zd9d-„Zd.S ):ÚGraphRunStreamuÍ  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_pumpÚ
graph_iterúIterator[Any] | NoneÚmuxr   r.   Úboolr   r   c               ó<  — || _         || _        t          |j        ¦  «        | _        d| _        d| _        d| _        g | _        t          |j	        ¦  «        | _
        |j        D ]}t          | ||j        |         ¦  «         Œ|r|                      |¦  «         dS dS )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)Úselfr/   r1   r.   Úkeys        r   Ú__init__zGraphRunStream.__init__C   s«   € ð& &ˆÔØˆŒ	Ý-=¸c¼nÑ-MÔ-MˆŒØˆŒØ.2ˆŒØ!ˆÔØ&(ˆÔÝ&*¨3¬9¡o¤oˆÔØ”?ð 	4ð 	4ˆCÝ�D˜#˜sœ~¨cÔ2Ñ3Ô3Ð3Ð3Øð 	)Ø×#Ò# CÑ(Ô(Ð(Ð(Ð(ð	)ð 	)r    c                ó:   — |                      | j        ¦  «         dS )al  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©rA   r1   s     r   r@   z!GraphRunStream._wire_request_morec   s   € ð 	�Š�d”oÑ&Ô&Ð&Ð&Ð&r    Úeventr   c                óè   — |d         dk    rdS |d         }|d         | j         k    rdS |d         | _        |                     dd¦  «        }|r#d	| _        | j                             |¦  «         dS dS ©
z;Track values-event state for output/interrupted/interrupts.Úmethodr(   NÚparamsÚ	namespaceÚdataÚ
interruptsr   T©r=   r8   Úgetr9   r:   Úextend©rA   rH   rL   rO   s       r   Ú_observe_eventzGraphRunStream._observe_eventn   óŒ   € à�Œ?˜hÒ&Ð&ØˆFØ�x”ˆØ�+Ô $Ô"2Ò2Ð2ØˆFØ˜f”~ˆŒØ—Z’Z ¨bÑ1Ô1ˆ
Øð 	0Ø $ˆDÔØÔ×#Ò# JÑ/Ô/Ð/Ð/Ð/ð	0ð 	0r    c                ó–  — | j         s| j        €dS 	 t          | j        ¦  «        }t          |¦  «        }|                      |¦  «         | j                             |¦  «         dS # t          $ r$ | j                             ¦   «          d| _         Y dS t          $ r,}| j         
                    |¦  «         d| _         Y d}~dS d}~ww xY w)a   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)r7   r4   Únextr   rT   r5   ÚpushÚStopIterationÚcloseÚ	ExceptionÚfail©rA   ÚpartrH   Úes       r   rF   zGraphRunStream._pump_next{   sã   € ð Œ?ð 	˜dÔ.Ð6Ø�5ð	Ý˜Ô(Ñ)Ô)ˆDÝ-¨dÑ3Ô3ˆEØ×Ò Ñ&Ô&Ð&ØŒI�NŠN˜5Ñ!Ô!Ð!Ø�4øÝð 	ð 	ð 	ØŒI�OŠOÑÔÐØ"ˆDŒOØ�5�5Ýð 	ð 	ð 	ØŒI�NŠN˜1ÑÔÐØ"ˆDŒOØ�5�5�5�5�5øøøøð	øøøs   ’AA& Á&*CÂ	CÂ!CÃCc                óú   — | j         rdS d| _         | j        }d| _        |�/t          |dd¦  «        x}�	  |¦   «          n# t          $ r Y nw xY w	 | j                             ¦   «          dS # t          $ r Y dS w xY w)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.
        NTrZ   )r7   r4   Úgetattrr[   r5   rZ   )rA   r/   rZ   s      r   ÚabortzGraphRunStream.abort”   s¸   € ð Œ?ð 	ØˆFØˆŒØÔ%ˆ
ØˆÔàÐ"Ý! *¨g°tÑ<Ô<Ð<�ÐIðØ�‘”��øÝð ð ð Ø�ðøøøð	ØŒI�OŠOÑÔÐÐÐøÝð 	ð 	ð 	ØˆDˆDð	øøøs#   µ
A  Á 
AÁAÁA, Á,
A:Á9A:c                ó   — | S ©Nr   ©rA   s    r   Ú	__enter__zGraphRunStream.__enter__­   s   € Øˆr    Úexc_typeútype[BaseException] | NoneÚexcúBaseException | NoneÚtbúTracebackType | Nonec                ó.   — |                       ¦   «          d S rd   ©rb   ©rA   rg   ri   rk   s       r   Ú__exit__zGraphRunStream.__exit__°   s   € ð 	�
Š
‰Œˆˆˆr    údict[str, Any] | Nonec                ób   — t          | j        ¦  «         | j        j        j        x}�|‚| j        S )z7Drive the run to completion and return the final state.)r   rF   r5   Ú_eventsÚ_errorr8   ©rA   Úerrs     r   ÚoutputzGraphRunStream.output¸   s4   € õ 	˜$œ/Ñ*Ô*Ð*Ø”9Ô$Ô+Ð+ˆCÐ8ØˆIØŒ|Ðr    c                ób   — t          | j        ¦  «         | j        j        j        x}�|‚| j        S )z¡Drive the run to completion, then return whether it was
        interrupted.

        Raises:
            BaseException: If the run ended with an error.
        )r   rF   r5   rs   rt   r9   ru   s     r   ÚinterruptedzGraphRunStream.interruptedÀ   s5   € õ 	˜$œ/Ñ*Ô*Ð*Ø”9Ô$Ô+Ð+ˆCÐ8ØˆIØÔ Ð r    ú	list[Any]c                ób   — t          | j        ¦  «         | j        j        j        x}�|‚| j        S )z‘Drive the run to completion, then return interrupt payloads.

        Raises:
            BaseException: If the run ended with an error.
        )r   rF   r5   rs   rt   r:   ru   s     r   rO   zGraphRunStream.interruptsÍ   s5   € õ 	˜$œ/Ñ*Ô*Ð*Ø”9Ô$Ô+Ð+ˆCÐ8ØˆIØÔÐr    úIterator[ProtocolEvent]c                ó4   — t          | j        j        ¦  «        S ©z<Subscribe to the main event log and iterate protocol events.)Úiterr5   rs   re   s    r   Ú__iter__zGraphRunStream.__iter__Ù   s   € å�D”IÔ%Ñ&Ô&Ð&r    ÚnamesÚstrúIterator[tuple[str, Any]]c              '  ó  K  — ddl m} i }	 |D ] }| j        |         }t          ||¦  «        s't	          dt          |¦  «        j        › d|›�¦  «        ‚|j        €t	          d|›d�¦  «        ‚|j        rt	          d|›d�¦  «        ‚|j        rt          d|›d	�¦  «        ‚d
|_        |||<   Œ¡t          ¦   «         }t          |¦  «        t          |¦  «        k     �rad}|                     ¦   «         D ]h\  }}||v rŒ
|j        r+|j        s$|j        �|j        ‚|                     |¦  «         Œ<|j        r%|j        d         d         }|�||d         k     r||f}Œi|�5||d                  j                             ¦   «         \  }	}
|d         |
fV — nŠ| j        j        }|�
 |¦   «         srt          |¦  «        }|                     ¦   «         D ]:\  }}||vr1|j        s*|j        r#|j        �|j        ‚|                     |¦  «         Œ;t          |¦  «        |k    rn!t          |¦  «        t          |¦  «        k     �°a|                     ¦   «         D ]	}d|_        Œ
dS # |                     ¦   «         D ]	}d|_        Œ
w xY w)uW  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   r6   Ú
isinstanceÚ	TypeErrorÚtypeÚ__name__Ú	_is_asyncÚ_subscribedÚRuntimeErrorÚsetÚlenÚitemsÚ_closedÚ_itemsrt   ÚaddÚpopleftr5   Ú_pump_fnr(   )rA   r�   r   ÚchannelsÚnameÚchÚdoneÚbestÚstampÚ_stampÚitemr   Úbefores                r   Ú
interleavezGraphRunStream.interleaveÝ   s=  è è € ð@ 	BÐAÐAÐAÐAÐAà24ˆð<	'Øð $ð $�Ø”_ TÔ*�Ý! " mÑ4Ô4ð Ý#ð@Ý# B™xœxÔ0ð@ð @Ø7;ð@ð @ñô ð ð ”<Ð'Ý#ðK¨ð Kð Kð Kñô ð ð ”<ð Ý#ðK¨ð Kð Kð Kñô ð ð ”>ð Ý&ð3¨ð 3ð 3ð 3ñô ð ð "&�”Ø!#�˜‘�å ™UœUˆDå�d‘)”)�c (™mœmÒ+Ñ+Ø/3�Ø (§¢Ñ 0Ô 0ð 1ð 1‘H�D˜"Ø˜t�|�|Ø Ø”zð !¨"¬)ð !Øœ9Ð0Ø"$¤)˜OØŸš ™œ˜Ø Ø”yð 1Ø "¤	¨!¤¨Q¤˜Ø˜<¨5°4¸´7ª?¨?Ø$)¨4 =˜DøàÐ#Ø#+¨D°¬GÔ#4Ô#;×#CÒ#CÑ#EÔ#E‘L�F˜DØ œ7 D˜/Ð)Ð)Ð)Ð)àœ9Ô-�DØ�|¨4¨4©6¬6�|Ý!$ T¡¤˜Ø(0¯ªÑ(8Ô(8ð 3ð 3™H˜D "Ø#¨4Ð/Ð/¸¼	Ð/Ø#%¤:ð !3Ø')¤yÐ'<Ø.0¬i¨Ø$(§H¢H¨T¡N¤N NøÝ˜t™9œ9¨Ò.Ð.Ø!õ; �d‘)”)�c (™mœmÒ+Ñ+ð> —o’oÑ'Ô'ð 'ð '�Ø!&�”�ð'ð 'ø�h—o’oÑ'Ô'ð 'ð '�Ø!&�”�ð'øøøs   ŒH3I É I?N)r/   r0   r1   r   r.   r2   r   r   ©r1   r   r   r   ©rH   r   r   r   ©r   r2   ©r   r   )r   r&   ©rg   rh   ri   rj   rk   rl   r   r   ©r   rq   ©r   rz   )r   r|   )r�   r‚   r   rƒ   )rŠ   Ú
__module__Ú__qualname__Ú__doc__Ú__annotations__rC   r@   rT   rF   rb   rf   rp   Úpropertyrw   ry   rO   r€   rŸ   r   r    r   r&   r&   $   s|  € € € € € € ðð ð0 *Ð)Ð)Ñ)Ø,Ð,Ð,Ñ,Ø.Ð.Ð.Ñ.Ø/Ð/Ð/Ñ/ð ð)ð )ð )ð )ð )ð )ð@	'ð 	'ð 	'ð 	'ð0ð 0ð 0ð 0ðð ð ð ð2ð ð ð ð2ð ð ð ðð ð ð ð ðð ð ñ „Xðð ð
!ð 
!ð 
!ñ „Xð
!ð ð	 ð 	 ð 	 ñ „Xð	 ð'ð 'ð 'ð 'ð_'ð _'ð _'ð _'ð _'ð _'r    r&   c                  óš   — e Zd ZU dZded<   ded<   ded<   ded	<   d
dœd+d„Zd,d„Zd-d„Zd.d„Zd/d„Z	d0d„Z
d1d"„Zd2d$„Zd.d%„Zd3d'„Zd4d)„Zd*S )5ÚAsyncGraphRunStreamuŸ  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-   Úgraph_aiterúAsyncIterator[Any] | Noner1   r   r.   r2   r   r   c               ó–  — || _         || _        t          |j        ¦  «        | _        d| _        d| _        d| _        g | _        t          |j	        ¦  «        | _
        t          j        ¦   «         | _        d| _        d| _        d| _        |j        D ]}t%          | ||j        |         ¦  «         Œ|r|                      |¦  «         dS dS )aI  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_aiterr5   r	   r6   r7   r8   r9   r:   r;   r<   r=   ÚasyncioÚ	ConditionÚ
_pump_condÚ_pumpingÚ_anext_taskÚ	_abortingr>   r?   Ú_wire_arequest_more)rA   r®   r1   r.   rB   s        r   rC   zAsyncGraphRunStream.__init__f  sÑ   € ð& (ˆÔØˆŒ	Ý-=¸c¼nÑ-MÔ-MˆŒØˆŒØ.2ˆŒØ!ˆÔØ&(ˆÔÝ&*¨3¬9¡o¤oˆÔÝ!Ô+Ñ-Ô-ˆŒØˆŒØ7;ˆÔØˆŒØ”?ð 	4ð 	4ˆCÝ�D˜#˜sœ~¨cÔ2Ñ3Ô3Ð3Ð3Øð 	*Ø×$Ò$ SÑ)Ô)Ð)Ð)Ð)ð	*ð 	*r    rH   r   c                óè   — |d         dk    rdS |d         }|d         | j         k    rdS |d         | _        |                     dd¦  «        }|r#d	| _        | j                             |¦  «         dS dS rJ   rP   rS   s       r   rT   z"AsyncGraphRunStream._observe_eventŠ  rU   r    c                ó:   — |                      | j        ¦  «         dS )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_nextrG   s     r   r¸   z'AsyncGraphRunStream._wire_arequest_more—  s   € ð 	�Š�tÔ'Ñ(Ô(Ð(Ð(Ð(r    c              ƒ  ó  K  — | j         4 ƒd{V —† | j        s| j        €	 ddd¦  «        ƒd{V —† dS | j        r9| j                              ¦   «         ƒ d{V —† | j         cddd¦  «        ƒd{V —† S d| _        ddd¦  «        ƒd{V —† n# 1 ƒd{V —†swxY w Y   	 	 t          j        | j                             ¦   «         ¦  «        | _        	 | j        ƒ d{V —†}n—# t
          j	        $ r… | j
        rcd| _        Y d| _        | j         4 ƒd{V —† d| _        | j                              ¦   «          ddd¦  «        ƒd{V —† dS # 1 ƒd{V —†swxY w Y   dS | j                             ¦   «          ‚ w xY w	 d| _        n# d| _        w xY wt          |¦  «        }|                      |¦  «         | j                             |¦  «        ƒ d{V —† 	 | j         4 ƒd{V —† d| _        | j                              ¦   «          ddd¦  «        ƒd{V —† dS # 1 ƒd{V —†swxY w Y   dS # t"          $ r| d| _        | j                             ¦   «         ƒ d{V —† Y | j         4 ƒd{V —† d| _        | j                              ¦   «          ddd¦  «        ƒd{V —† dS # 1 ƒd{V —†swxY w Y   dS t&          $ r„}d| _        | j                             |¦  «        ƒ d{V —† Y d}~| j         4 ƒd{V —† d| _        | j                              ¦   «          ddd¦  «        ƒd{V —† dS # 1 ƒd{V —†swxY w Y   dS d}~ww xY w# | j         4 ƒd{V —† d| _        | j                              ¦   «          ddd¦  «        ƒd{V —† w # 1 ƒd{V —†swxY w Y   w xY w)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´   r7   r±   rµ   Úwaitr²   Úensure_futureÚ	__anext__r¶   ÚCancelledErrorr·   Ú
notify_allÚcancelr   rT   r5   ÚapushÚStopAsyncIterationÚacloser[   Úafailr]   s       r   r¼   zAsyncGraphRunStream._apump_next¡  sº  è è € ð( ”?ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!ØŒð  $Ô"3Ð";Øð	!ð 	!ð 	!ñ 	!ô 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð Œ}ð +à”o×*Ò*Ñ,Ô,Ð,Ð,Ð,Ð,Ð,Ð,Ð,Øœ?Ð*ð	!ð 	!ð 	!ð 	!ñ 	!ô 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð !ˆDŒMð	!ð 	!ð 	!ñ 	!ô 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!øøøð 	!ð 	!ð 	!ð 	!ð$	-ðõ $+Ô#8¸Ô9J×9TÒ9TÑ9VÔ9VÑ#WÔ#W�Ô ð,Ø!%Ô!1Ð1Ð1Ð1Ð1Ð1Ð1�D�DøÝÔ-ð ð ð Ø”~ð %à*.˜œØ$ð (,�DÔ$ð ”ð -ð -ð -ð -ð -ð -ð -ð -Ø %�”Ø”×*Ò*Ñ,Ô,Ð,ð-ð -ð -ñ -ô -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -øøøð -ð -ð -ð -ð -ð -ð# Ô$×+Ò+Ñ-Ô-Ð-Øðøøøð ð (,�DÔ$Ð$ø t�DÔ$Ð+Ð+Ð+Ð+Ý1°$Ñ7Ô7�Ø×#Ò# EÑ*Ô*Ð*Ø”i—o’o eÑ,Ô,Ð,Ð,Ð,Ð,Ð,Ð,Ð,Øð ”ð -ð -ð -ð -ð -ð -ð -ð -Ø %�”Ø”×*Ò*Ñ,Ô,Ð,ð-ð -ð -ñ -ô -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -øøøð -ð -ð -ð -ð -ð -øõ &ð ð ð Ø"&�”Ø”i×&Ò&Ñ(Ô(Ð(Ð(Ð(Ð(Ð(Ð(Ð(Øð ”ð -ð -ð -ð -ð -ð -ð -ð -Ø %�”Ø”×*Ò*Ñ,Ô,Ð,ð-ð -ð -ñ -ô -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -øøøð -ð -ð -ð -ð -ð -õ ð ð ð Ø"&�”Ø”i—o’o aÑ(Ô(Ð(Ð(Ð(Ð(Ð(Ð(Ð(Ø�u�u�uà”ð -ð -ð -ð -ð -ð -ð -ð -Ø %�”Ø”×*Ò*Ñ,Ô,Ð,ð-ð -ð -ñ -ô -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -øøøð -ð -ð -ð -ð -ð -øøøøðøøøøð
 ”ð -ð -ð -ð -ð -ð -ð -ð -Ø %�”Ø”×*Ò*Ñ,Ô,Ð,ð-ð -ð -ñ -ô -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -ð -øøøð -ð -ð -ð -ð -øøøsû   �B³-BÁ3BÂ
BÂBÂ0H# ÃC ÃE> ÃE2Ã;E> Ã<H# Ä!EÅ
EÅEÅE2Å2E> Å6H# Å>	FÆAH# Ç!HÈ
HÈHÈ#0L5ÉL8 É!!JÊ
JÊ"JÊ(	L5Ê1'L0ËL8 Ë)!LÌ
L'Ì*L'Ì0L5Ì5L8 Ì8NÍ!M9Í'NÍ9
NÎNÎNÎNc              ƒ  óŠ  K  — | j         4 ƒd{V —† | j        r	 ddd¦  «        ƒd{V —† dS d| _        d| _        | j        }d| _        | j        }| j                              ¦   «          ddd¦  «        ƒd{V —† n# 1 ƒd{V —†swxY w Y   |�N|                     ¦   «         s:|                     ¦   «          	 |ƒ d{V —† n# t          j	        t          f$ r Y nw xY w|�5t          |dd¦  «        x}�"	  |¦   «         ƒ d{V —† n# t          $ r Y nw xY w	 | j                             ¦   «         ƒ d{V —† dS # t          $ r Y dS w xY w)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´   r7   r·   r±   r¶   rÂ   r™   rÃ   r²   rÁ   r[   ra   r5   rÆ   )rA   r®   Ú
anext_taskrÆ   s       r   rb   zAsyncGraphRunStream.abortä  so  è è € ð ”?ð 	)ð 	)ð 	)ð 	)ð 	)ð 	)ð 	)ð 	)ØŒð Øð	)ð 	)ð 	)ñ 	)ô 	)ð 	)ð 	)ð 	)ð 	)ð 	)ð 	)ð 	)ð 	)ð 	)ð #ˆDŒOØ!ˆDŒNØÔ+ˆKØ $ˆDÔØÔ)ˆJØŒO×&Ò&Ñ(Ô(Ð(ð	)ð 	)ð 	)ñ 	)ô 	)ð 	)ð 	)ð 	)ð 	)ð 	)ð 	)ð 	)ð 	)ð 	)ð 	)ð 	)ð 	)ð 	)ð 	)ð 	)ð 	)ð 	)ð 	)øøøð 	)ð 	)ð 	)ð 	)ð Ð!¨*¯/ª/Ñ*;Ô*;Ð!Ø×ÒÑÔÐðØ Ð Ð Ð Ð Ð Ð Ð Ð øÝÔ*­IÐ6ð ð ð Ø�ðøøøð Ð#Ý" ;°¸$Ñ?Ô?Ð?�ÐLðØ�f‘h”h��������øÝð ð ð Ø�ðøøøð	Ø”)×"Ò"Ñ$Ô$Ð$Ð$Ð$Ð$Ð$Ð$Ð$Ð$Ð$øÝð 	ð 	ð 	ØˆDˆDð	øøøsL   �	A:¬<A:Á:
BÂBÂ6B? Â?CÃCÃ1D Ä
DÄDÄD4 Ä4
EÅEc              ƒ  ó
   K  — | S rd   r   re   s    r   Ú
__aenter__zAsyncGraphRunStream.__aenter__  s   è è € Øˆr    rg   rh   ri   rj   rk   rl   c              ƒ  ó>   K  — |                       ¦   «         ƒ d {V —† d S rd   rn   ro   s       r   Ú	__aexit__zAsyncGraphRunStream.__aexit__  s.   è è € ð �jŠj‰lŒlÐÐÐÐÐÐÐÐÐr    rq   c              ƒ  ór   K  — t          | j        ¦  «        ƒ d{V —† | j        j        j        x}�|‚| j        S )a›  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¼   r5   rs   rt   r8   ru   s     r   rw   zAsyncGraphRunStream.output  sK   è è € õ ! Ô!1Ñ2Ô2Ð2Ð2Ð2Ð2Ð2Ð2Ð2Ø”9Ô$Ô+Ð+ˆCÐ8ØˆIØŒ|Ðr    c              ƒ  ór   K  — t          | j        ¦  «        ƒ d{V —† | j        j        j        x}�|‚| j        S )zŸDrive the run to completion and return whether it was
        interrupted.

        Raises:
            BaseException: If the run ended with an error.
        N)r#   r¼   r5   rs   rt   r9   ru   s     r   ry   zAsyncGraphRunStream.interrupted-  sL   è è € õ ! Ô!1Ñ2Ô2Ð2Ð2Ð2Ð2Ð2Ð2Ð2Ø”9Ô$Ô+Ð+ˆCÐ8ØˆIØÔ Ð r    rz   c              ƒ  ór   K  — t          | j        ¦  «        ƒ d{V —† | j        j        j        x}�|‚| j        S )z�Drive the run to completion and return interrupt payloads.

        Raises:
            BaseException: If the run ended with an error.
        N)r#   r¼   r5   rs   rt   r:   ru   s     r   rO   zAsyncGraphRunStream.interrupts9  sL   è è € õ ! Ô!1Ñ2Ô2Ð2Ð2Ð2Ð2Ð2Ð2Ð2Ø”9Ô$Ô+Ð+ˆCÐ8ØˆIØÔÐr    úAsyncIterator[ProtocolEvent]c                ó>   — | j         j                             ¦   «         S r~   )r5   rs   Ú	__aiter__re   s    r   rÓ   zAsyncGraphRunStream.__aiter__D  s   € àŒyÔ ×*Ò*Ñ,Ô,Ð,r    N)r®   r¯   r1   r   r.   r2   r   r   r¡   r    r¢   r£   )r   r­   r¤   r¥   r¦   )r   rÑ   )rŠ   r§   r¨   r©   rª   rC   rT   r¸   r¼   rb   rË   rÍ   rw   ry   rO   rÓ   r   r    r   r­   r­   ?  sE  € € € € € € ðð ð@ *Ð)Ð)Ñ)Ø1Ð1Ð1Ñ1Ø.Ð.Ð.Ñ.Ø4Ð4Ð4Ñ4ð ð"*ð "*ð "*ð "*ð "*ð "*ðH0ð 0ð 0ð 0ð)ð )ð )ð )ðA-ð A-ð A-ð A-ðF(ð (ð (ð (ðTð ð ð ðð ð ð ðð ð ð ð(
!ð 
!ð 
!ð 
!ð	 ð 	 ð 	 ð 	 ð-ð -ð -ð -ð -ð -r    r­   c                  óP   — e Zd ZU dZded<   ded<   ded<   ded<   ded	<   d
ed<   dS )Ú_SubgraphRunStreamMixinum  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Úerrorr2   Ú_seen_terminalN)rŠ   r§   r¨   r©   rª   r   r    r   rÕ   rÕ   I  sf   € € € € € € ðð ð& ÐÐÑØÐÐÑØÐÐÑØÐÐÑØÐÐÑØÐÐÑÐÐr    rÕ   c                  ó4   ‡ — e Zd ZdZdddœdˆ fd„Zdd„Zˆ xZS )ÚSubgraphRunStreamzASync handle for a discovered subgraph (extends `GraphRunStream`).N©rÙ   rÚ   r1   r   r×   rÖ   rÙ   rØ   rÚ   r   r   c               ó¼   •— |j         | _        t          ¦   «                              d |d¬¦  «         || _        || _        || _        d| _        d | _        d| _	        d S )NF)r/   r1   r.   Ústarted)
r•   Ú_parent_pump_fnÚsuperrC   r×   rÙ   rÚ   rÛ   rÜ   rÝ   ©rA   r1   r×   rÙ   rÚ   Ú	__class__s        €r   rC   zSubgraphRunStream.__init__h  sm   ø€ ð ;>¼,ˆÔÝ‰Œ×ÒØØØð 	ñ 	
ô 	
ð 	
ð
 ˆŒ	Ø$ˆŒØ.ˆÔØˆŒØˆŒ
Ø#ˆÔÐÐr    r2   c                óz   — | j         s| j        s| j        j        j        s| j        €dS |                      ¦   «         S )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.
        NF)r7   rÝ   r5   rs   r‘   rã   re   s    r   rF   zSubgraphRunStream._pump_next  sM   € ð ŒOð	àÔ"ð	ð ŒyÔ Ô(ð	ð Ô#Ð+à�5Ø×#Ò#Ñ%Ô%Ð%r    ©
r1   r   r×   rÖ   rÙ   rØ   rÚ   rØ   r   r   r¢   )rŠ   r§   r¨   r©   rC   rF   Ú__classcell__©ræ   s   @r   rß   rß   e  si   ø€ € € € € ØKÐKð "&Ø&*ð$ð $ð $ð $ð $ð $ð $ð $ð.&ð &ð &ð &ð &ð &ð &ð &r    rß   c                  ó4   ‡ — e Zd ZdZdddœdˆ fd„Zdd„Zˆ xZS )ÚAsyncSubgraphRunStreamzGAsync handle for a discovered subgraph (extends `AsyncGraphRunStream`).Nrà   r1   r   r×   rÖ   rÙ   rØ   rÚ   r   r   c               ó¼   •— |j         | _        t          ¦   «                              d |d¬¦  «         || _        || _        || _        d| _        d | _        d| _	        d S )NF)r®   r1   r.   râ   )
Ú	_apump_fnÚ_parent_apump_fnrä   rC   r×   rÙ   rÚ   rÛ   rÜ   rÝ   rå   s        €r   rC   zAsyncSubgraphRunStream.__init__“  so   ø€ ð GJÄmˆÔÝ‰Œ×ÒØØØð 	ñ 	
ô 	
ð 	
ð
 ˆŒ	Ø$ˆŒØ.ˆÔØˆŒØˆŒ
Ø#ˆÔÐÐr    r2   c              ƒ  óŠ   K  — | j         s| j        s| j        j        j        s| j        €dS |                      ¦   «         ƒ d{V —†S )z$Delegate to the parent's async pump.NF)r7   rÝ   r5   rs   r‘   rï   re   s    r   r¼   z"AsyncSubgraphRunStream._apump_next¨  sc   è è € ð ŒOð	àÔ"ð	ð ŒyÔ Ô(ð	ð Ô$Ð,à�5Ø×*Ò*Ñ,Ô,Ð,Ð,Ð,Ð,Ð,Ð,Ð,r    rè   r¢   )rŠ   r§   r¨   r©   rC   r¼   ré   rê   s   @r   rì   rì   �  si   ø€ € € € € ØQÐQð "&Ø&*ð$ð $ð $ð $ð $ð $ð $ð $ð*	-ð 	-ð 	-ð 	-ð 	-ð 	-ð 	-ð 	-r    rì   )r   r   r   r   )r   r!   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û      so  ðØ "Ð "Ð "Ð "Ð "Ð "à €€€Ø QÐ QÐ QÐ QÐ QÐ QÐ QÐ QÐ QÐ QÐ QÐ QÐ QÐ QØ 1Ð 1Ð 1Ð 1Ð 1Ð 1Ð 1Ð 1Ø %Ð %Ð %Ð %Ð %Ð %Ð %Ð %à $Ð $Ð $Ð $Ð $Ð $à ?Ð ?Ð ?Ð ?Ð ?Ð ?Ø +Ð +Ð +Ð +Ð +Ð +Ø 1Ð 1Ð 1Ð 1Ð 1Ð 1àð Oðð ð ð ð ð ð ð ð
 >Ð=Ð=Ð=Ð=Ð=ØNÐNÐNÐNÐNÐNÐNÐNðð ð ð ðð ð ð ð €ÐDÐEÑEÔEðW'ð W'ð W'ð W'ð W'ñ W'ô W'ñ FÔEðW'ðt €ÐDÐEÑEÔEðF-ð F-ð F-ð F-ð F-ñ F-ô F-ñ FÔEðF-ðRð ð ð ð ñ ô ð ð8(&ð (&ð (&ð (&ð (&˜Ð(?ñ (&ô (&ð (&ðV!-ð !-ð !-ð !-ð !-Ð0Ð2Iñ !-ô !-ð !-ð !-ð !-r    