§
    šŠtj/Ÿ  ã                  ó^  — d dl mZ d dlZd dlmZmZmZmZ d dlm	Z	 d dl
mZmZ d dl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mZ d d
lmZmZ d dlmZmZ d dl m!Z! erd dl"m#Z#m$Z$ d dl%m&Z&  ej'        e(¦  «        Z) G d„ de¦  «        Z* G d„ de¦  «        Z+ G d„ de¦  «        Z, G d„ de¦  «        Z-ed         Z.d4d„Z/ G d„ ded¬ ¦  «        Z0 G d!„ d"e¦  «        Z1d5d&„Z2d6d)„Z3 G d*„ d+e1¦  «        Z4 G d,„ d-e1¦  «        Z5 G d.„ d/e¦  «        Z6 G d0„ d1e¦  «        Z7 G d2„ d3e¦  «        Z8dS )7é    )ÚannotationsN)ÚTYPE_CHECKINGÚAnyÚLiteralÚcast)Úmessage_to_events)ÚAsyncChatModelStreamÚChatModelStream)ÚAIMessageChunkÚBaseMessageÚToolMessage)ÚLifecycleCauseÚMessagesData)ÚNotRequiredÚ	TypedDict)ÚGraphDrainedÚGraphInterrupt)ÚProtocolEventÚStreamTransformer)ÚAsyncSubgraphRunStreamÚSubgraphRunStream)ÚStreamChannel)Ú	AwaitableÚCallable)Ú	StreamMuxc                  óV   ‡ — e Zd ZdZdZdZddˆ fd	„Zdd„Zedd„¦   «         Z	dd„Z
ˆ xZS )ÚValuesTransformeru  Capture values events as a drainable stream of state snapshots.

    Provides the `run.values` projection. `run.output`,
    `run.interrupted` and `run.interrupts` are tracked directly
    by the run stream and do not depend on this transformer.

    Native transformer â€” projection keys are exposed as direct
    attributes on the run stream (e.g. `run.values`).

    Only values events at the run's own level are captured; snapshots
    from deeper subgraphs are left in the main event log but excluded
    from the projection. "Own level" is defined by `scope`, which
    `stream_events(version="v3")` / `astream_events(version="v3")` populate from the caller's
    checkpoint namespace so that a nested `stream_events(version="v3")` call still
    sees its own root snapshots.
    T)Úvalues© Úscopeútuple[str, ...]ÚreturnÚNonec                óÂ   •— t          ¦   «                              |¦  «         t          ¦   «         | _        d | _        d| _        g | _        t          |¦  «        | _        d S )NF)	ÚsuperÚ__init__r   Ú_logÚ_latestÚ_interruptedÚ_interruptsÚlistÚ_scope_list©Úselfr    Ú	__class__s     €ú[/var/www/html/CA-Chatbot/venv/lib/python3.11/site-packages/langgraph/stream/transformers.pyr&   zValuesTransformer.__init__1   sS   ø€ Ý‰Œ×Ò˜ÑÔÐÝ3@±?´?ˆŒ	Ø.2ˆŒØ!ˆÔØ&(ˆÔõ '+¨5¡k¤kˆÔÐÐó    údict[str, Any]c                ó   — d| j         iS )Nr   ©r'   ©r.   s    r0   ÚinitzValuesTransformer.init;   ó   € Ø˜$œ)Ð$Ð$r1   úBaseException | Nonec                ó   — | j         j        S )z€The error that ended the run, or `None` if it succeeded.

        Set by the mux when it auto-fails the projection log.
        )r'   Ú_errorr5   s    r0   ÚerrorzValuesTransformer.error>   s   € ð ŒyÔÐr1   Úeventr   Úboolc                ó$  — |d         dk    rdS |d         }|d         | j         k    rdS |d         | _        |                     dd¦  «        }|r!d| _        | j                             |¦  «         | j                             |d         ¦  «         dS )	NÚmethodr   TÚparamsÚ	namespaceÚdataÚ
interruptsr   )r,   r(   Úgetr)   r*   Úextendr'   Úpush)r.   r<   r@   rC   s       r0   ÚprocesszValuesTransformer.processF   s™   € Ø�Œ?˜hÒ&Ð&Ø�4Ø�x”ˆØ�+Ô $Ô"2Ò2Ð2Ø�4Ø˜f”~ˆŒØ—Z’Z ¨bÑ1Ô1ˆ
Øð 	0Ø $ˆDÔØÔ×#Ò# JÑ/Ô/Ð/ØŒ	�Š�v˜f”~Ñ&Ô&Ð&Øˆtr1   ©r   ©r    r!   r"   r#   ©r"   r2   ©r"   r8   ©r<   r   r"   r=   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__Ú_nativeÚrequired_stream_modesr&   r6   Úpropertyr;   rG   Ú__classcell__©r/   s   @r0   r   r      s¡   ø€ € € € € ðð ð" €GØ'Ðð2ð 2ð 2ð 2ð 2ð 2ð 2ð%ð %ð %ð %ð ð ð  ð  ñ „Xð ðð ð ð ð ð ð ð r1   r   c                  ó>   ‡ — e Zd ZdZdZdZddˆ fd	„Zdd„Zdd„Zˆ xZ	S )ÚCustomTransformeruå  Capture custom events as a drainable stream of arbitrary payloads.

    Nodes emit custom data via `get_stream_writer()`. This transformer
    surfaces those events on `run.custom` as a `StreamChannel[Any]`,
    preserving payloads in arrival order.

    Only events at the run's own scope are captured; custom data from
    deeper subgraphs is available on the respective subgraph handle's
    `.custom` projection.

    Native transformer â€” `run.custom` is a direct attribute.
    T)Úcustomr   r    r!   r"   r#   c                ó˜   •— t          ¦   «                              |¦  «         t          ¦   «         | _        t	          |¦  «        | _        d S ©N©r%   r&   r   r'   r+   r,   r-   s     €r0   r&   zCustomTransformer.__init__f   s:   ø€ Ý‰Œ×Ò˜ÑÔÐÝ(5©¬ˆŒ	Ý&*¨5¡k¤kˆÔÐÐr1   r2   c                ó   — d| j         iS )NrX   r4   r5   s    r0   r6   zCustomTransformer.initk   r7   r1   r<   r   r=   c                ó˜   — |d         dk    rdS |d         }|d         | j         k    rdS | j                             |d         ¦  «         dS )Nr?   rX   Tr@   rA   rB   ©r,   r'   rF   ©r.   r<   r@   s      r0   rG   zCustomTransformer.processn   sT   € Ø�Œ?˜hÒ&Ð&Ø�4Ø�x”ˆØ�+Ô $Ô"2Ò2Ð2Ø�4ØŒ	�Š�v˜f”~Ñ&Ô&Ð&Øˆtr1   rH   rI   rJ   rL   ©
rM   rN   rO   rP   rQ   rR   r&   r6   rG   rT   rU   s   @r0   rW   rW   U   s�   ø€ € € € € ðð ð €GØ'Ðð2ð 2ð 2ð 2ð 2ð 2ð 2ð
%ð %ð %ð %ðð ð ð ð ð ð ð r1   rW   c                  ó>   ‡ — e Zd ZdZdZdZddˆ fd	„Zdd„Zdd„Zˆ xZ	S )ÚUpdatesTransformeruì  Capture updates events as a drainable stream of node outputs.

    Surfaces `stream_mode="updates"` data on `run.updates` as a
    `StreamChannel[dict[str, Any]]`. Each item is a dict mapping a node
    (or task) name to the update it returned after a step.

    Only events at the run's own scope are captured; updates from deeper
    subgraphs are available on the respective subgraph handle's
    `.updates` projection.

    Native transformer â€” `run.updates` is a direct attribute.
    T)Úupdatesr   r    r!   r"   r#   c                ó˜   •— t          ¦   «                              |¦  «         t          ¦   «         | _        t	          |¦  «        | _        d S rZ   r[   r-   s     €r0   r&   zUpdatesTransformer.__init__‰   ó:   ø€ Ý‰Œ×Ò˜ÑÔÐÝ3@±?´?ˆŒ	Ý&*¨5¡k¤kˆÔÐÐr1   r2   c                ó   — d| j         iS )Nrc   r4   r5   s    r0   r6   zUpdatesTransformer.initŽ   s   € Ø˜4œ9Ð%Ð%r1   r<   r   r=   c                ó˜   — |d         dk    rdS |d         }|d         | j         k    rdS | j                             |d         ¦  «         dS )Nr?   rc   Tr@   rA   rB   r^   r_   s      r0   rG   zUpdatesTransformer.process‘   sT   € Ø�Œ?˜iÒ'Ð'Ø�4Ø�x”ˆØ�+Ô $Ô"2Ò2Ð2Ø�4ØŒ	�Š�v˜f”~Ñ&Ô&Ð&Øˆtr1   rH   rI   rJ   rL   r`   rU   s   @r0   rb   rb   x   s�   ø€ € € € € ðð ð €GØ(Ðð2ð 2ð 2ð 2ð 2ð 2ð 2ð
&ð &ð &ð &ðð ð ð ð ð ð ð r1   rb   c                  óv   ‡ — e Zd ZdZdZdZd'd(ˆ fd	„Zd)d„Zd*d„Zd+d„Z	d,d„Z
d-d„Zd.d„Zd/d"„Zd0d#„Zd1d&„Zˆ xZS )2ÚMessagesTransformeru>  Capture messages events as ChatModelStream objects.

    The messages projection yields one `ChatModelStream` (or
    `AsyncChatModelStream`) per LLM call. Consumers iterate
    `run.messages` to get stream handles, then use each handle's typed
    projections (`.text`, `.reasoning`, `.tool_calls`, `.usage`,
    `.output`) for per-message content.

    Two input shapes are handled (via `params["data"] = (payload,
    metadata)` from `StreamMessagesHandler`):

    1. Protocol event (dict with `"event"` key) â€” emitted by
       `stream_events(version="v3")` / `astream_events(version="v3")` via the `on_stream_event`
       callback. Routed to an existing `ChatModelStream` by
       `metadata["run_id"]`. A `message-start` event creates a new
       stream; `message-finish` closes it.
    2. Whole `AIMessage` â€” emitted from `on_chain_end` when a node
       returns a finalized message. Replayed as a synthetic protocol
       event lifecycle via `message_to_events`, then the
       already-complete stream is pushed to the log.

    V1 `AIMessageChunk` tuples (from `on_llm_new_token`) are not
    streamed into this projection: chat models that want to populate
    `run.messages` with content-block streaming must use
    `stream_events(version="v3")` / `astream_events(version="v3")`. Models called via the legacy
    `stream()` method still surface their final `AIMessage` via
    `on_chain_end` when a node returns it as state.

    Only events at the run's own level are projected; tokens from
    deeper subgraphs are left in the main event log but excluded from
    `.messages`. "Own level" is defined by `scope`, which
    `stream_events(version="v3")` / `astream_events(version="v3")` populate from the caller's checkpoint
    namespace so that a `stream_events(version="v3")` call inside a node still sees its
    own root chat model streams on `.messages`. Consumers that need
    subgraph tokens should iterate the raw event stream or register a
    custom transformer.

    Native transformer â€” the `messages` projection is exposed as a
    direct attribute on the run stream.
    T)Úmessagesr   r    r!   r"   r#   c                óè   •— t          ¦   «                              |¦  «         t          ¦   «         | _        i | _        t          ¦   «         | _        d | _        d | _        t          |¦  «        | _
        d S rZ   )r%   r&   r   r'   Ú_by_runÚsetÚ_ignored_runsÚ_pump_fnÚ	_apump_fnr+   r,   r-   s     €r0   r&   zMessagesTransformer.__init__È   s_   ø€ Ý‰Œ×Ò˜ÑÔÐÝ4A±O´OˆŒ	ð 46ˆŒÝ'*¡u¤uˆÔØ37ˆŒØ?CˆŒõ '+¨5¡k¤kˆÔÐÐr1   r2   c                ó   — d| j         iS )Nrj   r4   r5   s    r0   r6   zMessagesTransformer.initÕ   s   € Ø˜DœIÐ&Ð&r1   ÚfnúCallable[[], bool]c                ó   — || _         dS )zIWire the sync pull callback. Called by GraphRunStream._wire_request_more.N)ro   ©r.   rr   s     r0   Ú
_bind_pumpzMessagesTransformer._bind_pumpØ   s   € àˆŒˆˆr1   úCallable[[], Awaitable[bool]]c                ó   — || _         dS )zèWire the async pull callback.

        Called by `AsyncGraphRunStream._wire_arequest_more` so each
        `AsyncChatModelStream` this transformer creates can drive the
        shared graph pump from its projection cursors.
        N)rp   ru   s     r0   Ú_bind_apumpzMessagesTransformer._bind_apumpÜ   s   € ð ˆŒˆˆr1   rA   ú	list[str]Únodeú
str | NoneÚ
message_idr
   c               óú   — | j         �.t          |||¬¦  «        }|                     | j         ¦  «         |S | j        �.t	          |||¬¦  «        }|                     | j        ¦  «         |S t          |||¬¦  «        S )a^  Create a ChatModelStream (sync) or AsyncChatModelStream (async).

        Wires whichever pump is bound. Prefers the async pump so nested
        iteration under `AsyncGraphRunStream` drives the graph forward
        without a background task. The unwired fallback (no pump bound)
        is used by unit tests that dispatch events manually.
        N©rA   r{   r}   )rp   r	   Úset_arequest_morero   r
   Úset_request_more)r.   rA   r{   r}   ÚastreamÚstreams         r0   Ú_make_streamz MessagesTransformer._make_streamå   s¨   € ð Œ>Ð%Ý*Ø#ØØ%ðñ ô ˆGð
 ×%Ò% d¤nÑ5Ô5Ð5ØˆNØŒ=Ð$Ý&5Ø#ØØ%ð'ñ 'ô 'ˆFð
 ×#Ò# D¤MÑ2Ô2Ð2ØˆMÝ#ØØØ!ð
ñ 
ô 
ð 	
r1   r<   r   r=   c                ó  — |d         dk    rdS |d         }|d         | j         k    rdS |d         \  }}|                     d¦  «        }|r#t          |                     dd	¦  «        ¦  «        nd	}t          |t          ¦  «        r+d
|v r'|                      t          d|¦  «        ||¬¦  «         nVt          |t          ¦  «        rAt          |t          ¦  «        s,t          |t          ¦  «        s|  
                    ||¬¦  «         dS )Nr?   rj   Tr@   rA   rB   Úlanggraph_nodeÚrun_idÚ r<   r   )r‡   r{   )r{   )r,   rD   ÚstrÚ
isinstanceÚdictÚ_route_protocol_eventr   r   r   r   Ú_route_whole_message)r.   r<   r@   ÚpayloadÚmetadatar{   r‡   s          r0   rG   zMessagesTransformer.process	  s'  € Ø�Œ?˜jÒ(Ð(Ø�4Ø�x”ˆØ�+Ô $Ô"2Ò2Ð2Ø�4à" 6œNÑˆ�Ø#Ÿ<š<Ð(8Ñ9Ô9ˆØ4<ÐD•�X—\’\ (¨BÑ/Ô/Ñ0Ô0Ð0À"ˆå�g�tÑ$Ô$ð 		:¨°GÐ);Ð);Ø×&Ò&Ý�^ WÑ-Ô-°fÀ4ð 'ñ ô ð ð õ �w¥Ñ,Ô,ð	:å˜w­Ñ7Ô7ð	:õ ˜w­Ñ4Ô4ð	:ð
 ×%Ò% g°DÐ%Ñ9Ô9Ð9ð
 ˆtr1   r   r‡   r‰   c               ól  — |                      d¦  «        }|dk    r®|                      d¦  «        dk    r| j                             |¦  «         d S |                      d¦  «        }|                      g ||�t	          |¦  «        nd ¬¦  «        }|| j        |<   | j                             |¦  «         |                     |¦  «         d S || j        v r$|dk    r| j         	                    |¦  «         d S d S || j        v r2| j        |         }|                     |¦  «         |dk    r| j        |= d S d S d S )Nr<   zmessage-startÚroleÚtoolr}   r   zmessage-finish)
rD   rn   Úaddr„   r‰   rl   r'   rF   ÚdispatchÚdiscard)r.   r<   r‡   r{   Ú
event_typer}   rƒ   s          r0   rŒ   z)MessagesTransformer._route_protocol_event$  sh  € ð —Y’Y˜wÑ'Ô'ˆ
Ø˜Ò(Ð(ð �yŠy˜Ñ Ô  FÒ*Ð*ØÔ"×&Ò& vÑ.Ô.Ð.Ø�ØŸš <Ñ0Ô0ˆJØ×&Ò&ØØØ.8Ð.D�3˜z™?œ?˜?È$ð 'ñ ô ˆFð
 $*ˆDŒL˜Ñ ØŒI�NŠN˜6Ñ"Ô"Ð"Ø�OŠO˜EÑ"Ô"Ð"Ð"Ð"Ø�tÔ)Ð)Ð)ØÐ-Ò-Ð-ØÔ"×*Ò*¨6Ñ2Ô2Ð2Ð2Ð2ð .Ð-à�t”|Ð#Ð#Ø”\ &Ô)ˆFØ�OŠO˜EÑ"Ô"Ð"ØÐ-Ò-Ð-Ø”L Ð(Ð(Ð(ð	 $Ð#ð .Ð-r1   Úmessager   c               óÐ   — |                       g ||j        ¬¦  «        }t          ||j        ¬¦  «        D ]}|                     |¦  «         Œ| j                             |¦  «         d S )Nr   )r}   )r„   Úidr   r”   r'   rF   )r.   r—   r{   rƒ   Úevts        r0   r�   z(MessagesTransformer._route_whole_messageD  sk   € Ø×"Ò"¨R°dÀwÄzÐ"ÑRÔRˆÝ$ W¸¼ÐDÑDÔDð 	!ð 	!ˆCØ�OŠO˜CÑ Ô Ð Ð ØŒ	�Š�vÑÔÐÐÐr1   c                ój   — | j                              ¦   «          | j                             ¦   «          dS )uJ   Clear any routing state â€” streams close themselves via `message-finish`.N)rl   Úclearrn   r5   s    r0   ÚfinalizezMessagesTransformer.finalizeJ  s1   € àŒ×ÒÑÔÐØÔ× Ò Ñ"Ô"Ð"Ð"Ð"r1   ÚerrÚBaseExceptionc                óæ   — t          | j                             ¦   «         ¦  «        D ]}|                     |¦  «         Œ| j                             ¦   «          | j                             ¦   «          dS )zCPropagate run error to any streams still open when the graph fails.N)r+   rl   r   Úfailrœ   rn   )r.   rž   rƒ   s      r0   r¡   zMessagesTransformer.failO  sk   € å˜4œ<×.Ò.Ñ0Ô0Ñ1Ô1ð 	ð 	ˆFØ�KŠK˜ÑÔÐÐØŒ×ÒÑÔÐØÔ× Ò Ñ"Ô"Ð"Ð"Ð"r1   rH   rI   rJ   )rr   rs   r"   r#   )rr   rw   r"   r#   )rA   rz   r{   r|   r}   r|   r"   r
   rL   )r<   r   r‡   r‰   r{   r|   r"   r#   )r—   r   r{   r|   r"   r#   ©r"   r#   ©rž   rŸ   r"   r#   )rM   rN   rO   rP   rQ   rR   r&   r6   rv   ry   r„   rG   rŒ   r�   r�   r¡   rT   rU   s   @r0   ri   ri   ›   s  ø€ € € € € ð'ð 'ðR €GØ)Ðð2ð 2ð 2ð 2ð 2ð 2ð 2ð'ð 'ð 'ð 'ðð ð ð ðð ð ð ð"
ð "
ð "
ð "
ðHð ð ð ð6)ð )ð )ð )ð@ð ð ð ð#ð #ð #ð #ð
#ð #ð #ð #ð #ð #ð #ð #r1   ri   )ÚstartedÚ	completedÚfailedÚinterruptedÚdrainedÚsegmentr‰   r"   útuple[str, str | None]c                óD   — |                       d¦  «        \  }}}||r|ndfS )zÁSplit a namespace segment into `(graph_name, trigger_call_id)`.

    Segments are formatted `node_name:task_id` by `prepare_next_tasks`.
    Returns `(segment, None)` if no `:` is present.
    ú:N)Ú	partition)r©   ÚnameÚsepÚtask_ids       r0   Ú_parse_ns_segmentr±   Z  s2   € ð !×*Ò*¨3Ñ/Ô/Ñ€Dˆ#ˆwØ˜CÐ)�� TÐ)Ð)r1   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 )ÚLifecyclePayloada,  Payload of a lifecycle event surfaced on the `lifecycle` channel.

    Auto-forwarded as `lifecycle` protocol events (no `custom:` prefix
    because `LifecycleTransformer` is a native transformer) so remote
    SDK clients receive the same data in-process consumers see via
    `run.lifecycle`.
    ÚSubgraphStatusr<   rz   rA   zNotRequired[str]Ú
graph_nameÚtrigger_call_idzNotRequired[LifecycleCause]Úcauser;   N)rM   rN   rO   rP   Ú__annotations__r   r1   r0   r³   r³   d  sf   € € € € € € ðð ð ÐÐÑØÐÐÑØ Ð Ð Ñ Ø%Ð%Ð%Ñ%Ø&Ð&Ð&Ñ&ØÐÐÑÐÐr1   r³   F)Útotalc                  ó‚   ‡ — e Zd ZdZdZd#d$ˆ fd„Zd%d„Zd&d„Zd'd„Zd(d„Z	d)d„Z
d*d„Zd)d„Zd+d„Zd)d„Zd,d„Zd-d"„Zˆ xZS ).Ú_TasksLifecycleBaseu  Shared bookkeeping for `tasks`-event-driven lifecycle inference.

    Both `LifecycleTransformer` (wire-serializable channel) and
    `SubgraphTransformer` (in-process navigation handles) discover
    subgraphs by watching the same `tasks` stream â€” `started` on the
    first event at a tracked namespace, terminal status when the
    parent's `TaskResultPayload` arrives. Centralizing the dispatch
    + open-set bookkeeping here keeps the inference rules from
    drifting between the two surfaces.

    Subclasses provide three template-method hooks:

    - `_should_track(ns)` â€” scope filter (e.g. multi-depth vs
      direct-children-only).
    - `_on_started(ns, graph_name, trigger_call_id)` â€” first sighting
      action (push payload / build handle / etc.). Called once per
      discovered namespace.
    - `_on_terminal(ns, status, error)` â€” terminal action (push
      terminal payload / mark handle status). Called once per
      tracked namespace at result time, or via `finalize` / `fail`
      sweeps if no parent result arrived.

    Tasks events are suppressed from the main event log (`process`
    returns False) â€” they're folded into whichever projection the
    subclass populates; consumers iterating the raw protocol stream
    see the higher-level view.
    ©Útasksr   r    r!   r"   r#   c                ó¨   •— t          ¦   «                              |¦  «         t          ¦   «         | _        i | _        i | _        i | _        d | _        d S rZ   )r%   r&   rm   Ú_seenÚ_openÚ	_lc_by_nsÚ_pending_tool_callsÚ_pending_causer-   s     €r0   r&   z_TasksLifecycleBase.__init__”  sR   ø€ Ý‰Œ×Ò˜ÑÔÐÝ+.©5¬5ˆŒ
ð 24ˆŒ
ð =?ˆŒð 46ˆÔ ð
 6:ˆÔÐÐr1   Únsr=   c                ó   — t           ‚)uF   Scope filter â€” return True iff `ns` is in this transformer's region.©ÚNotImplementedError©r.   rÄ   s     r0   Ú_should_trackz!_TasksLifecycleBase._should_track®  s   € å!Ð!r1   rµ   r|   r¶   c                ó   — t           ‚)z™Fired once per discovered namespace (first observed task event).

        The triggering `cause` (if any) is available as `self._pending_cause`.
        rÆ   )r.   rÄ   rµ   r¶   s       r0   Ú_on_startedz_TasksLifecycleBase._on_started²  s
   € õ "Ð!r1   Ústatusr´   r;   c                ó   — t           ‚)z{Fired once per tracked namespace when its parent's result arrives,
        or via finalize/fail safety-net sweeps.
        rÆ   )r.   rÄ   rÌ   r;   s       r0   Ú_on_terminalz _TasksLifecycleBase._on_terminal¾  s
   € õ "Ð!r1   r<   r   c                ó,  — |d         dk    rdS t          |d         d         ¦  «        }|d         d         }d|v r|                      ||¦  «         nA|                      ||¦  «         |                      |¦  «         |                      ||¦  «         dS )	Nr?   r½   Tr@   rA   rB   ÚresultF)ÚtupleÚ_handle_task_resultÚ_record_identityÚ_record_pending_tool_callsÚ_handle_task_start)r.   r<   rÄ   rB   s       r0   rG   z_TasksLifecycleBase.processË  s¥   € Ø�Œ?˜gÒ%Ð%Ø�4Ý�5˜”? ;Ô/Ñ0Ô0ˆØ�XŒ˜vÔ&ˆØ�tÐÐØ×$Ò$ R¨Ñ.Ô.Ð.Ð.à×!Ò! " dÑ+Ô+Ð+Ø×+Ò+¨DÑ1Ô1Ð1Ø×#Ò# B¨Ñ-Ô-Ð-ð ˆur1   rB   r2   c                ó„   — || j         v rdS |                     d¦  «        pi }|                     d¦  «        | j         |<   dS )ay  Record this namespace's `lc_agent_name` (first task event wins).

        Runs for every task-start event, including `ns == self.scope` and
        tracked children. Pregel emits parent-namespace tasks before
        child-namespace tasks, so under that ordering the parent's identity is
        recorded by the time a child event is evaluated in `_handle_task_start`.
        Nr�   Úlc_agent_name)rÁ   rD   )r.   rÄ   rB   r�   s       r0   rÓ   z$_TasksLifecycleBase._record_identityÛ  sJ   € ð �”ÐÐØˆFØ—8’8˜JÑ'Ô'Ð-¨2ˆØ%Ÿ\š\¨/Ñ:Ô:ˆŒ�rÑÐÐr1   c                óJ  — |                      d¦  «        }t          |t          ¦  «        sdS |                      d¦  «        }d}t          |t          ¦  «        r[t          |                      d¦  «        t          ¦  «        r3|d                               d¦  «        }t          |t          ¦  «        r|}nat          |t          ¦  «        rL|D ]I}t          |t          ¦  «        r2t          |                      d¦  «        t          ¦  «        r
|d         } nŒJ|�|| j        |<   dS dS )a  Harvest a task's triggering tool_call_id keyed by its task id.

        A tool-dispatch task seeds `task_id -> tool_call_id`; the spawned
        subgraph's namespace segment `node:<task_id>` shares that id, letting
        a subagent recover the tool call that caused it across payloads. Two
        input shapes are handled: the current Pregel push model schedules each
        tool call as its own task whose `input` is a `tool_call_with_context`
        dict, while a legacy / batched model passes a list of tool-call dicts.
        r™   NÚinputÚ	tool_call)rD   rŠ   r‰   r‹   r+   rÂ   )r.   rB   r°   rŽ   Útool_call_idÚ	candidateÚtcs          r0   rÔ   z._TasksLifecycleBase._record_pending_tool_callsè  s  € ð —(’(˜4‘.”.ˆÝ˜'¥3Ñ'Ô'ð 	ØˆFØ—(’(˜7Ñ#Ô#ˆØ#'ˆõ �g�tÑ$Ô$ð 		­°G·K²KÀÑ4LÔ4LÍdÑ)SÔ)Sð 		Ø Ô,×0Ò0°Ñ6Ô6ˆIÝ˜)¥SÑ)Ô)ð )Ø(�øå˜¥Ñ&Ô&ð 	Øð ð �Ý˜b¥$Ñ'Ô'ð ­J°r·v²v¸d±|´|ÅSÑ,IÔ,Ið Ø#% d¤8�LØ�EøØÐ#Ø0<ˆDÔ$ WÑ-Ð-Ð-ð $Ð#r1   c                óÚ  — |                       |¦  «        r	|| j        v rd S | j                             |¦  «         t          |d         ¦  «        \  }}|                     d¦  «        pi }|                     d¦  «        }|d u}|r|n|pd }d }	|r0|�.| j                             |¦  «        }
|
rdt          |
¦  «        dœ}	|	| _        |                      |||¦  «         |�|| j	        |<   d S d S )Néÿÿÿÿr�   r×   ÚtoolCall)ÚtyperÛ   )
rÉ   r¿   r“   r±   rD   rÂ   r‰   rÃ   rË   rÀ   )r.   rÄ   rB   Úparsed_namer¶   r�   Úchild_lcÚis_subagentrµ   r·   rÛ   s              r0   rÕ   z&_TasksLifecycleBase._handle_task_start  s  € Ø×!Ò! "Ñ%Ô%ð 	¨¨t¬zÐ)9Ð)9ØˆFØŒ
�Š�rÑÔÐÝ'8¸¸B¼Ñ'@Ô'@Ñ$ˆ�_Ø—8’8˜JÑ'Ô'Ð-¨2ˆØ—<’< Ñ0Ô0ˆð  dÐ*ˆØ!,ÐG�X�X°;Ð3FÀ$ˆ
Ø'+ˆØð 	P˜?Ð6ØÔ3×7Ò7¸ÑHÔHˆLØð PØ!+½SÀÑ=NÔ=NÐOÐO�ð $ˆÔØ×Ò˜˜Z¨Ñ9Ô9Ð9ØÐ&Ø,ˆDŒJ�r‰NˆNˆNð 'Ð&r1   ú8list[tuple[tuple[str, ...], SubgraphStatus, str | None]]c                ó"  — |                      d¦  «        }|sg S g }t          | j                             ¦   «         ¦  «        D ]L\  }}|dd…         |k    s||k    rŒt	          |¦  «        \  }}|                     |||f¦  «         | j        |= ŒM|S )z>Return and remove tracked children closed by this task result.r™   Nrß   )rD   r+   rÀ   ÚitemsÚ_terminal_from_resultÚappend)	r.   rÄ   rB   Ú	result_idÚtransitionsÚchild_nsÚparent_task_idrÌ   r;   s	            r0   Ú_pop_terminal_transitionsz-_TasksLifecycleBase._pop_terminal_transitions$  s®   € ð —H’H˜T‘N”Nˆ	Øð 	ØˆIØPRˆÝ(,¨T¬Z×-=Ò-=Ñ-?Ô-?Ñ(@Ô(@ð 	%ð 	%Ñ$ˆH�nØ˜˜˜Œ} Ò"Ð" n¸	Ò&AÐ&AØÝ1°$Ñ7Ô7‰MˆF�EØ×Ò ¨&°%Ð8Ñ9Ô9Ð9Ø”
˜8Ð$Ð$ØÐr1   c                ón   — |                       ||¦  «        D ]\  }}}|                      |||¦  «         Œd S rZ   )rî   rÎ   )r.   rÄ   rB   rì   rÌ   r;   s         r0   rÒ   z'_TasksLifecycleBase._handle_task_result4  sM   € Ø'+×'EÒ'EÀbÈ$Ñ'OÔ'Oð 	7ð 	7Ñ#ˆH�f˜eØ×Ò˜h¨°Ñ6Ô6Ð6Ð6ð	7ð 	7r1   c                ó”   — t          | j        ¦  «        D ]}|                      |dd¦  «         Œ| j                             ¦   «          dS )zAEmit `completed` for any tracked namespace still open at run end.r¥   N)r+   rÀ   rÎ   rœ   rÈ   s     r0   r�   z_TasksLifecycleBase.finalize8  sP   € å�t”zÑ"Ô"ð 	5ð 	5ˆBØ×Ò˜b +¨tÑ4Ô4Ð4Ð4ØŒ
×ÒÑÔÐÐÐr1   rž   rŸ   c                ó¸   — t          |¦  «        \  }}t          | j        ¦  «        D ]}|                      |||¦  «         Œ| j                             ¦   «          dS )z:Emit terminal status for any tracked namespace still open.N)Ú_status_from_exceptionr+   rÀ   rÎ   rœ   )r.   rž   rÌ   Ú	error_strrÄ   s        r0   r¡   z_TasksLifecycleBase.fail>  sb   € å2°3Ñ7Ô7Ñˆ�	Ý�t”zÑ"Ô"ð 	5ð 	5ˆBØ×Ò˜b &¨)Ñ4Ô4Ð4Ð4ØŒ
×ÒÑÔÐÐÐr1   rH   rI   ©rÄ   r!   r"   r=   ©rÄ   r!   rµ   r|   r¶   r|   r"   r#   ©rÄ   r!   rÌ   r´   r;   r|   r"   r#   rL   )rÄ   r!   rB   r2   r"   r#   )rB   r2   r"   r#   )rÄ   r!   rB   r2   r"   rå   r¢   r£   )rM   rN   rO   rP   rR   r&   rÉ   rË   rÎ   rG   rÓ   rÔ   rÕ   rî   rÒ   r�   r¡   rT   rU   s   @r0   r»   r»   u  s0  ø€ € € € € ðð ð8 'Ðð:ð :ð :ð :ð :ð :ð :ð4"ð "ð "ð "ð
"ð 
"ð 
"ð 
"ð	"ð 	"ð 	"ð 	"ðð ð ð ð ;ð ;ð ;ð ;ð=ð =ð =ð =ð<-ð -ð -ð -ð<ð ð ð ð 7ð 7ð 7ð 7ðð ð ð ðð ð ð ð ð ð ð r1   r»   rž   rŸ   ú!tuple[SubgraphStatus, str | None]c                ó€   — t          | t          ¦  «        rdS t          | t          ¦  «        rdS dt          | ¦  «        fS )zCMap a run exception to a subgraph terminal status and error string.)r¨   N©r§   Nr¦   )rŠ   r   r   r‰   )rž   s    r0   rò   rò   F  sB   € å�#•|Ñ$Ô$ð ØˆÝ�#•~Ñ&Ô&ð #Ø"Ð"Ø•S˜‘X”XÐÐr1   rŽ   r2   c                ó„   — |                       d¦  «        rdS |                       d¦  «        }|rdt          |¦  «        fS dS )züMap a `TaskResultPayload` to a `(status, error)` pair.

    Order matters: a result with both `error` and `interrupts` prefers
    the interrupt classification, since `GraphInterrupt` manifests as
    a populated `interrupts` list, not as `error`.
    rC   rù   r;   r¦   )r¥   N)rD   r‰   )rŽ   r;   s     r0   rè   rè   O  sN   € ð ‡{‚{�<Ñ Ô ð #Ø"Ð"Ø�KŠK˜Ñ Ô €EØð $Ø�˜U™œÐ#Ð#ØÐr1   c                  óJ   ‡ — e Zd ZdZdZddˆ fd„Zdd
„Zdd„Zdd„Zdd„Z	ˆ xZ
S )ÚLifecycleTransformeru®  Surface subgraph lifecycle as `lifecycle` protocol events.

    Pushes `LifecyclePayload` to a `StreamChannel` named `lifecycle`.
    The channel is auto-forwarded by the mux so payloads land in the
    main event log under `method = "lifecycle"` (native transformer â€”
    no `custom:` prefix) â€” visible to remote SDK clients over the
    wire and to in-process consumers via `run.lifecycle`.

    Tracks subgraphs at every depth strictly below the transformer's
    scope, so a graph â†’ subgraph â†’ subgraph chain produces lifecycle
    events for both nested levels in a flat stream.

    Native transformer â€” projection key `lifecycle` is exposed as
    `run.lifecycle`.
    Tr   r    r!   r"   r#   c                ór   •— t          ¦   «                              |¦  «         t          d¦  «        | _        d S ©NÚ	lifecycle)r%   r&   r   Ú_channelr-   s     €r0   r&   zLifecycleTransformer.__init__s  s.   ø€ Ý‰Œ×Ò˜ÑÔÐÝ9FÀ{Ñ9SÔ9SˆŒˆˆr1   r2   c                ó   — d| j         iS rþ   )r   r5   s    r0   r6   zLifecycleTransformer.initw  s   € Ø˜Tœ]Ð+Ð+r1   rÄ   r=   c                óv   — t          | j        ¦  «        }t          |¦  «        |k    o|d |…         | j        k    S rZ   ©Úlenr    ©r.   rÄ   Údepths      r0   rÉ   z"LifecycleTransformer._should_trackz  s3   € Ý�D”J‘”ˆÝ�2‰wŒw˜ŠÐ; 2 f u f¤:°´Ò#;Ð;r1   rµ   r|   r¶   c                óš   — |€d S dt          |¦  «        dœ}|r||d<   ||d<   | j        }|�||d<   | j                             |¦  «         d S )Nr¤   ©r<   rA   rµ   r¶   r·   )r+   rÃ   r   rF   )r.   rÄ   rµ   r¶   rŽ   r·   s         r0   rË   z LifecycleTransformer._on_started~  su   € ð Ð"ð ˆFØ.7ÅdÈ2ÁhÄhÐ$OÐ$OˆØð 	/Ø$.ˆG�LÑ!Ø%4ˆÐ!Ñ"ØÔ#ˆØÐØ$ˆG�GÑØŒ×Ò˜7Ñ#Ô#Ð#Ð#Ð#r1   rÌ   r´   r;   c                ól   — |t          |¦  «        dœ}|�||d<   | j                             |¦  «         d S )Nr  r;   )r+   r   rF   )r.   rÄ   rÌ   r;   rŽ   s        r0   rÎ   z!LifecycleTransformer._on_terminal’  sC   € ð /5Å4ÈÁ8Ä8Ð$LÐ$LˆØÐØ$ˆG�GÑØŒ×Ò˜7Ñ#Ô#Ð#Ð#Ð#r1   rH   rI   rJ   rô   rõ   rö   )rM   rN   rO   rP   rQ   r&   r6   rÉ   rË   rÎ   rT   rU   s   @r0   rü   rü   `  sª   ø€ € € € € ðð ð  €GðTð Tð Tð Tð Tð Tð Tð,ð ,ð ,ð ,ð<ð <ð <ð <ð$ð $ð $ð $ð(	$ð 	$ð 	$ð 	$ð 	$ð 	$ð 	$ð 	$r1   rü   c                  óÂ   ‡ — e Zd ZdZdZdZd.d/ˆ fd„Zd0d
„Zd1d„Zd2d„Z	d3d„Z
d4d„Zd4d„Zd5d„Zd6d„Zd6d„Zd7d"„Zd8ˆ fd#„Zd8d$„Zd9d&„Zd9d'„Zd:d(„Zd:d)„Zd;d,„Zd;d-„Zˆ xZS )<ÚSubgraphTransformeru#  Discover subgraph invocations as in-process navigation handles.

    Per discovered direct-child subgraph, builds a `SubgraphRunStream`
    (or `AsyncSubgraphRunStream`) wrapping a child mini-mux scoped to
    the subgraph's namespace. Consumers iterate `run.subgraphs` to
    receive handles, then drill into `handle.values` / `handle.messages`
    / `handle.subgraphs` (recursive grandchildren) / `handle.lifecycle`.

    Each mini-mux owns its own scope and uses its own
    `SubgraphTransformer` to discover its direct children, so
    grandchildren live on the child handle â€” never on the root's
    `subgraphs` log. Forwarding events into the matching child mini-mux
    is what keeps the child's projections populated.

    Native transformer â€” `subgraphs` is exposed as `run.subgraphs`.
    Tr   r    r!   r"   r#   c                óŒ   •— t          ¦   «                              |¦  «         t          ¦   «         | _        i | _        d | _        d S rZ   )r%   r&   r   r'   Ú_handlesÚ_muxr-   s     €r0   r&   zSubgraphTransformer.__init__³  sB   ø€ Ý‰Œ×Ò˜ÑÔÐå‰OŒOð 	Œ	ð
 ð 	Œð '+ˆŒ	ˆ	ˆ	r1   r2   c                ó   — d| j         iS )NÚ	subgraphsr4   r5   s    r0   r6   zSubgraphTransformer.init½  s   € Ø˜TœYÐ'Ð'r1   Úmuxr   c                ó   — || _         d S rZ   )r  )r.   r  s     r0   Ú_on_registerz SubgraphTransformer._on_registerÀ  s   € ØˆŒ	ˆ	ˆ	r1   rÄ   r=   c                ó|   — t          | j        ¦  «        }t          |¦  «        |dz   k    o|d |…         | j        k    S )Né   r  r  s      r0   rÉ   z!SubgraphTransformer._should_trackÃ  s:   € õ �D”J‘”ˆÝ�2‰wŒw˜% !™)Ò#Ð@¨¨6¨E¨6¬
°d´jÒ(@Ð@r1   rµ   r|   r¶   c                ó  — | j         €d S 	 | j                              |¦  «        }n# t          $ r Y d S w xY w|j        rt          nt
          } |||||¬¦  «        }|| j        |<   | j                             |¦  «         d S )N)r  Úpathrµ   r¶   )	r  Ú_make_childÚRuntimeErrorÚis_asyncr   r   r  r'   rF   )r.   rÄ   rµ   r¶   Ú	child_muxÚ
handle_clsÚhandles          r0   rË   zSubgraphTransformer._on_startedÉ  s¯   € ð Œ9ÐØˆFð	Øœ	×-Ò-¨bÑ1Ô1ˆIˆIøÝð 	ð 	ð 	ØˆFˆFð	øøøà/8Ô/AÐXÕ+Ð+ÕGXˆ
ð �ØØØ!Ø+ð	
ñ 
ô 
ˆð #ˆŒ�bÑØŒ	�Š�vÑÔÐÐÐs   ‹& ¦
4³4rÌ   r´   r;   c                óž   — | j                              |¦  «        }|�|                      |||¦  «        sd S |                      |||¦  «         d S rZ   )r  rD   Ú_mark_terminalÚ_close_or_fail_handle©r.   rÄ   rÌ   r;   r  s        r0   rÎ   z SubgraphTransformer._on_terminalâ  sW   € ð ”×"Ò" 2Ñ&Ô&ˆØˆ> ×!4Ò!4°V¸VÀUÑ!KÔ!Kˆ>ØˆFØ×"Ò" 6¨6°5Ñ9Ô9Ð9Ð9Ð9r1   c              ƒ  ó®   K  — | j                              |¦  «        }|�|                      |||¦  «        sd S |                      |||¦  «        ƒ d {V —† d S rZ   )r  rD   r  Ú_aclose_or_fail_handler!  s        r0   Ú_aon_terminalz!SubgraphTransformer._aon_terminalí  sm   è è € ð ”×"Ò" 2Ñ&Ô&ˆØˆ> ×!4Ò!4°V¸VÀUÑ!KÔ!Kˆ>ØˆFØ×)Ò)¨&°&¸%Ñ@Ô@Ð@Ð@Ð@Ð@Ð@Ð@Ð@Ð@Ð@r1   r  ú*SubgraphRunStream | AsyncSubgraphRunStreamc                óT   — |j         rdS ||_        |�|j        €||_        d|_         dS )z>Mark a handle terminal once. Returns True on first transition.FNT)Ú_seen_terminalrÌ   r;   ©r.   r  rÌ   r;   s       r0   r  z"SubgraphTransformer._mark_terminalø  s<   € ð Ô ð 	Ø�5ØˆŒØÐ ¤Ð!5Ø ˆFŒLØ $ˆÔØˆtr1   c                óÎ   — |j         �|j         j        j        rd S |dk    r+|j                              t	          |pd¦  «        ¦  «         d S |j                              ¦   «          d S ©Nr¦   zSubgraph failed)r  Ú_eventsÚ_closedr¡   r  Úcloser(  s       r0   r   z)SubgraphTransformer._close_or_fail_handle  sm   € ð Œ;Ð &¤+Ô"5Ô"=ÐØˆFØ�XÒÐØŒK×Ò�\¨%Ð*DÐ3DÑEÔEÑFÔFÐFÐFÐFàŒK×ÒÑÔÐÐÐr1   c              ƒ  óê   K  — |j         �|j         j        j        rd S |dk    r1|j                              t	          |pd¦  «        ¦  «        ƒ d {V —† d S |j                              ¦   «         ƒ d {V —† d S r*  )r  r+  r,  Úafailr  Úacloser(  s       r0   r#  z*SubgraphTransformer._aclose_or_fail_handle  s•   è è € ð Œ;Ð &¤+Ô"5Ô"=ÐØˆFØ�XÒÐØ”+×#Ò#¥L°Ð1KÐ:KÑ$LÔ$LÑMÔMÐMÐMÐMÐMÐMÐMÐMÐMÐMà”+×$Ò$Ñ&Ô&Ð&Ð&Ð&Ð&Ð&Ð&Ð&Ð&Ð&r1   r<   r   ú1SubgraphRunStream | AsyncSubgraphRunStream | Nonec                ó  — t          |d         d         ¦  «        }t          | j        ¦  «        }t          |¦  «        |dz   k     rd S | j                             |d |dz   …         ¦  «        }|�|j        �|j        j        j        rd S |S )Nr@   rA   r  )rÑ   r  r    r  rD   r  r+  r,  )r.   r<   rÄ   r  r  s        r0   Ú_handle_for_eventz%SubgraphTransformer._handle_for_event!  s…   € õ �5˜”? ;Ô/Ñ0Ô0ˆÝ�D”J‘”ˆÝˆr‰7Œ7�U˜Q‘YÒÐØ�4Ø”×"Ò" 2 k¨°©	 k¤?Ñ3Ô3ˆØˆ>˜Vœ[Ð0°F´KÔ4GÔ4OÐ0Ø�4Øˆr1   c                óÖ   •— t          ¦   «                              |¦  «        }|                      |¦  «        }|�/|                     |¦  «         |j                             |¦  «         |S rZ   )r%   rG   r3  Ú_observe_eventr  rF   )r.   r<   Úkeepr  r/   s       €r0   rG   zSubgraphTransformer.process-  sb   ø€ õ ‰wŒw�Š˜uÑ%Ô%ˆØ×'Ò'¨Ñ.Ô.ˆØÐØ×!Ò! %Ñ(Ô(Ð(ØŒK×Ò˜UÑ#Ô#Ð#Øˆr1   c              ƒ  ó  K  — |d         dk    r¬t          |d         d         ¦  «        }|d         d         }d|v r;|                      ||¦  «        D ]#\  }}}|                      |||¦  «        ƒ d {V —† Œ$nA|                      ||¦  «         |                      |¦  «         |                      ||¦  «         d}nd}|                      |¦  «        }|�5|                     |¦  «         |j         	                    |¦  «        ƒ d {V —† |S )	Nr?   r½   r@   rA   rB   rÐ   FT)
rÑ   rî   r$  rÓ   rÔ   rÕ   r3  r5  r  Úapush)	r.   r<   rÄ   rB   rì   rÌ   r;   r6  r  s	            r0   ÚaprocesszSubgraphTransformer.aprocess7  sJ  è è € ð �Œ?˜gÒ%Ð%Ý�u˜X” {Ô3Ñ4Ô4ˆBØ˜”? 6Ô*ˆDØ˜4ÐÐØ/3×/MÒ/MÈbÐRVÑ/WÔ/Wð Fð FÑ+�H˜f eØ×,Ò,¨X°v¸uÑEÔEÐEÐEÐEÐEÐEÐEÐEÐEðFð ×%Ò% b¨$Ñ/Ô/Ð/Ø×/Ò/°Ñ5Ô5Ð5Ø×'Ò'¨¨DÑ1Ô1Ð1ØˆDˆDàˆDØ×'Ò'¨Ñ.Ô.ˆØÐØ×!Ò! %Ñ(Ô(Ð(Ø”+×#Ò# EÑ*Ô*Ð*Ð*Ð*Ð*Ð*Ð*Ð*Øˆr1   r8   c                óž  — d }t          | j        ¦  «        D ]5}	 |                      |dd ¦  «         Œ# t          $ r}|€|}Y d }~Œ.d }~ww xY w| j                             ¦   «          | j                             ¦   «         D ]M}|                      |dd ¦  «        r4	 |                      |dd ¦  «         Œ2# t          $ r}|€|}Y d }~ŒEd }~ww xY wŒN|S ©Nr¥   )	r+   rÀ   rÎ   rŸ   rœ   r  r   r  r   ©r.   Úfirst_errorrÄ   Úer  s        r0   Ú_complete_open_handlesz*SubgraphTransformer._complete_open_handlesP  s  € Ø,0ˆÝ�t”zÑ"Ô"ð 	$ð 	$ˆBð$Ø×!Ò! " k°4Ñ8Ô8Ð8Ð8øÝ ð $ð $ð $ØÐ&Ø"#�Køøøøøøøøøð$øøøð 	Œ
×ÒÑÔÐØ”m×*Ò*Ñ,Ô,ð 	(ð 	(ˆFØ×"Ò" 6¨;¸Ñ=Ô=ð (ð(Ø×.Ò.¨v°{ÀDÑIÔIÐIÐIøÝ$ð (ð (ð (Ø"Ð*Ø&'˜øøøøøøøøøð(øøøð(ð Ðs,   š2²
A
¼AÁA
ÂB1Â1
C	Â;CÃC	c              ƒ  óº  K  — d }t          | j        ¦  «        D ];}	 |                      |dd ¦  «        ƒ d {V —† Œ!# t          $ r}|€|}Y d }~Œ4d }~ww xY w| j                             ¦   «          | j                             ¦   «         D ]S}|                      |dd ¦  «        r:	 |                      |dd ¦  «        ƒ d {V —† Œ8# t          $ r}|€|}Y d }~ŒKd }~ww xY wŒT|S r;  )	r+   rÀ   r$  rŸ   rœ   r  r   r  r#  r<  s        r0   Ú_acomplete_open_handlesz+SubgraphTransformer._acomplete_open_handlesb  s@  è è € Ø,0ˆÝ�t”zÑ"Ô"ð 	$ð 	$ˆBð$Ø×(Ò(¨¨[¸$Ñ?Ô?Ð?Ð?Ð?Ð?Ð?Ð?Ð?Ð?øÝ ð $ð $ð $ØÐ&Ø"#�Køøøøøøøøøð$øøøð 	Œ
×ÒÑÔÐØ”m×*Ò*Ñ,Ô,ð 	(ð 	(ˆFØ×"Ò" 6¨;¸Ñ=Ô=ð (ð(Ø×5Ò5°f¸kÈ4ÑPÔPÐPÐPÐPÐPÐPÐPÐPÐPøÝ$ð (ð (ð (Ø"Ð*Ø&'˜øøøøøøøøøð(øøøð(ð Ðs-   œ:º
AÁAÁAÂ!B?Â?
CÃ	CÃCc                ó6   — |                       ¦   «         }|�|‚d S rZ   )r?  ©r.   r=  s     r0   r�   zSubgraphTransformer.finalizet  s'   € Ø×1Ò1Ñ3Ô3ˆØÐ"ØÐð #Ð"r1   c              ƒ  óF   K  — |                       ¦   «         ƒ d {V —†}|�|‚d S rZ   )rA  rC  s     r0   Ú	afinalizezSubgraphTransformer.afinalizey  s=   è è € Ø ×8Ò8Ñ:Ô:Ð:Ð:Ð:Ð:Ð:Ð:ˆØÐ"ØÐð #Ð"r1   rž   rŸ   c                óŽ  — t          |¦  «        \  }}| j                             ¦   «          | j                             ¦   «         D ]}|                      |||¦  «         |j        �_|j        j        j        sN	 |j         	                    |¦  «         ŒM# t          $ r% t                               d|j        d¬¦  «         Y Œ{w xY wŒ€d S ©NzRError failing subgraph mini-mux at %s; subscribers may not see the terminal error.T)Úexc_info)rò   rÀ   rœ   r  r   r  r  r+  r,  r¡   Ú	ExceptionÚ_loggerÚwarningr  ©r.   rž   rÌ   ró   r  s        r0   r¡   zSubgraphTransformer.fail~  sé   € Ý2°3Ñ7Ô7Ñˆ�	ØŒ
×ÒÑÔÐØ”m×*Ò*Ñ,Ô,ð 	ð 	ˆFØ×Ò ¨°	Ñ:Ô:Ð:ØŒ{Ð&¨v¬{Ô/BÔ/JÐ&ðØ”K×$Ò$ SÑ)Ô)Ð)Ð)øÝ ð ð ð Ý—O’OðFàœØ!%ð	 $ñ ô ð ð ð ðøøøøð	ð 	s   Á7BÂ,CÃ Cc              ƒ  óž  K  — t          |¦  «        \  }}| j                             ¦   «          | j                             ¦   «         D ]…}|                      |||¦  «         |j        �e|j        j        j        sT	 |j         	                    |¦  «        ƒ d {V —† ŒS# t          $ r% t                               d|j        d¬¦  «         Y Œ�w xY wŒ†d S rG  )rò   rÀ   rœ   r  r   r  r  r+  r,  r/  rI  rJ  rK  r  rL  s        r0   r/  zSubgraphTransformer.afailŽ  sÿ   è è € Ý2°3Ñ7Ô7Ñˆ�	ØŒ
×ÒÑÔÐØ”m×*Ò*Ñ,Ô,ð 	ð 	ˆFØ×Ò ¨°	Ñ:Ô:Ð:ØŒ{Ð&¨v¬{Ô/BÔ/JÐ&ðØ œ+×+Ò+¨CÑ0Ô0Ð0Ð0Ð0Ð0Ð0Ð0Ð0Ð0øÝ ð ð ð Ý—O’OðFàœØ!%ð	 $ñ ô ð ð ð ðøøøøð	ð 	s   Á9 BÂ,C	ÃC	rH   rI   rJ   )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   r"   r1  rL   rK   r¢   r£   )rM   rN   rO   rP   rQ   Úsupports_syncr&   r6   r  rÉ   rË   rÎ   r$  r  r   r#  r3  rG   r9  r?  rA  r�   rE  r¡   r/  rT   rU   s   @r0   r  r  ž  sÒ  ø€ € € € € ðð ð" €GØ€Mð+ð +ð +ð +ð +ð +ð +ð(ð (ð (ð (ðð ð ð ðAð Að Að Aðð ð ð ð2	:ð 	:ð 	:ð 	:ð	Að 	Að 	Að 	Aðð ð ð ð ð  ð  ð  ð'ð 'ð 'ð 'ð
ð 
ð 
ð 
ðð ð ð ð ð ðð ð ð ð2ð ð ð ð$ð ð ð ð$ð ð ð ð
ð ð ð ð
ð ð ð ð ð ð ð ð ð ð ð r1   r  c                  ó>   ‡ — e Zd ZdZdZdZddˆ fd	„Zdd„Zdd„Zˆ xZ	S )ÚCheckpointsTransformeru‰  Capture checkpoint events as a drainable stream.

    Surfaces `stream_mode="checkpoints"` data on `run.checkpoints` as
    a `StreamChannel[dict[str, Any]]`. Each item is in the same format
    as returned by `get_state()`.

    Checkpoint events are only emitted when a checkpointer is configured
    on the graph. When no checkpointer is present, the projection exists
    but receives no events.

    Only events at the run's own scope are captured; checkpoint data from
    deeper subgraphs is available on the respective subgraph handle's
    `.checkpoints` projection.

    Native transformer â€” `run.checkpoints` is a direct attribute.
    T)Úcheckpointsr   r    r!   r"   r#   c                ó˜   •— t          ¦   «                              |¦  «         t          ¦   «         | _        t	          |¦  «        | _        d S rZ   r[   r-   s     €r0   r&   zCheckpointsTransformer.__init__´  re   r1   r2   c                ó   — d| j         iS )NrQ  r4   r5   s    r0   r6   zCheckpointsTransformer.init¹  s   € Ø˜tœyÐ)Ð)r1   r<   r   r=   c                ó˜   — |d         dk    rdS |d         }|d         | j         k    rdS | j                             |d         ¦  «         dS )Nr?   rQ  Tr@   rA   rB   r^   r_   s      r0   rG   zCheckpointsTransformer.process¼  sT   € Ø�Œ?˜mÒ+Ð+Ø�4Ø�x”ˆØ�+Ô $Ô"2Ò2Ð2Ø�4ØŒ	�Š�v˜f”~Ñ&Ô&Ð&Øˆtr1   rH   rI   rJ   rL   r`   rU   s   @r0   rP  rP  Ÿ  s�   ø€ € € € € ðð ð" €GØ,Ðð2ð 2ð 2ð 2ð 2ð 2ð 2ð
*ð *ð *ð *ðð ð ð ð ð ð ð r1   rP  c                  ó>   ‡ — e Zd ZdZdZdZddˆ fd	„Zdd„Zdd„Zˆ xZ	S )ÚDebugTransformeru  Capture debug events as a drainable stream.

    Surfaces `stream_mode="debug"` data on `run.debug` as a
    `StreamChannel[dict[str, Any]]`. Each item is a debug event with
    step-level detail (checkpoint snapshots, task payloads, and
    task results wrapped with step number and timestamp).

    Only events at the run's own scope are captured; debug data from
    deeper subgraphs is available on the respective subgraph handle's
    `.debug` projection.

    Native transformer â€” `run.debug` is a direct attribute.
    T)Údebugr   r    r!   r"   r#   c                ó˜   •— t          ¦   «                              |¦  «         t          ¦   «         | _        t	          |¦  «        | _        d S rZ   r[   r-   s     €r0   r&   zDebugTransformer.__init__Ø  re   r1   r2   c                ó   — d| j         iS )NrW  r4   r5   s    r0   r6   zDebugTransformer.initÝ  ó   € Ø˜œÐ#Ð#r1   r<   r   r=   c                ó˜   — |d         dk    rdS |d         }|d         | j         k    rdS | j                             |d         ¦  «         dS )Nr?   rW  Tr@   rA   rB   r^   r_   s      r0   rG   zDebugTransformer.processà  óT   € Ø�Œ?˜gÒ%Ð%Ø�4Ø�x”ˆØ�+Ô $Ô"2Ò2Ð2Ø�4ØŒ	�Š�v˜f”~Ñ&Ô&Ð&Øˆtr1   rH   rI   rJ   rL   r`   rU   s   @r0   rV  rV  Æ  s�   ø€ € € € € ðð ð €GØ&Ðð2ð 2ð 2ð 2ð 2ð 2ð 2ð
$ð $ð $ð $ðð ð ð ð ð ð ð r1   rV  c                  ó>   ‡ — e Zd ZdZdZdZddˆ fd	„Zdd„Zdd„Zˆ xZ	S )ÚTasksTransformeru›  Capture raw task events as a drainable stream.

    Surfaces `stream_mode="tasks"` data on `run.tasks` as a
    `StreamChannel[dict[str, Any]]`. Each item is a task payload
    (start or result).

    `LifecycleTransformer` and `SubgraphTransformer` also consume
    `tasks` events for subgraph discovery and lifecycle tracking.
    This transformer captures the raw payloads independently for
    consumers who need task-level detail.

    Only events at the run's own scope are captured; task data from
    deeper subgraphs is available on the respective subgraph handle's
    `.tasks` projection.

    Native transformer â€” `run.tasks` is a direct attribute.
    Tr¼   r   r    r!   r"   r#   c                ó˜   •— t          ¦   «                              |¦  «         t          ¦   «         | _        t	          |¦  «        | _        d S rZ   r[   r-   s     €r0   r&   zTasksTransformer.__init__   re   r1   r2   c                ó   — d| j         iS )Nr½   r4   r5   s    r0   r6   zTasksTransformer.init  rZ  r1   r<   r   r=   c                ó˜   — |d         dk    rdS |d         }|d         | j         k    rdS | j                             |d         ¦  «         dS )Nr?   r½   Tr@   rA   rB   r^   r_   s      r0   rG   zTasksTransformer.process  r\  r1   rH   rI   rJ   rL   r`   rU   s   @r0   r^  r^  ê  s�   ø€ € € € € ðð ð$ €GØ&Ðð2ð 2ð 2ð 2ð 2ð 2ð 2ð
$ð $ð $ð $ðð ð ð ð ð ð ð r1   r^  )r©   r‰   r"   rª   )rž   rŸ   r"   r÷   )rŽ   r2   r"   r÷   )9Ú
__future__r   ÚloggingÚtypingr   r   r   r   Ú-langchain_core.language_models._compat_bridger   Ú0langchain_core.language_models.chat_model_streamr	   r
   Úlangchain_core.messagesr   r   r   Úlangchain_protocol.protocolr   r   Útyping_extensionsr   r   Úlanggraph.errorsr   r   Úlanggraph.stream._typesr   r   Úlanggraph.stream.run_streamr   r   Úlanggraph.stream.stream_channelr   Úcollections.abcr   r   Úlanggraph.stream._muxr   Ú	getLoggerrM   rJ  r   rW   rb   ri   r´   r±   r³   r»   rò   rè   rü   r  rP  rV  r^  r   r1   r0   ú<module>rq     sÀ  ðØ "Ð "Ð "Ð "Ð "Ð "à €€€Ø 4Ð 4Ð 4Ð 4Ð 4Ð 4Ð 4Ð 4Ð 4Ð 4Ð 4Ð 4à KÐ KÐ KÐ KÐ KÐ Kðð ð ð ð ð ð ð ð MÐ LÐ LÐ LÐ LÐ LÐ LÐ LÐ LÐ LØ DÐ DÐ DÐ DÐ DÐ DÐ DÐ DØ 4Ð 4Ð 4Ð 4Ð 4Ð 4Ð 4Ð 4à 9Ð 9Ð 9Ð 9Ð 9Ð 9Ð 9Ð 9Ø DÐ DÐ DÐ DÐ DÐ DÐ DÐ DØ QÐ QÐ QÐ QÐ QÐ QÐ QÐ QØ 9Ð 9Ð 9Ð 9Ð 9Ð 9àð 0Ø3Ð3Ð3Ð3Ð3Ð3Ð3Ð3à/Ð/Ð/Ð/Ð/Ð/à
ˆ'Ô
˜HÑ
%Ô
%€ð6ð 6ð 6ð 6ð 6Ð)ñ 6ô 6ð 6ðr ð  ð  ð  ð  Ð)ñ  ô  ð  ðF ð  ð  ð  ð  Ð*ñ  ô  ð  ðFy#ð y#ð y#ð y#ð y#Ð+ñ y#ô y#ð y#ðx ÐSÔT€ð*ð *ð *ð *ðð ð ð ð �y¨ð ñ ô ð ð"Nð Nð Nð Nð NÐ+ñ Nô Nð Nðbð ð ð ðð ð ð ð";$ð ;$ð ;$ð ;$ð ;$Ð.ñ ;$ô ;$ð ;$ð|~ð ~ð ~ð ~ð ~Ð-ñ ~ô ~ð ~ðB$ð $ð $ð $ð $Ð.ñ $ô $ð $ðN!ð !ð !ð !ð !Ð(ñ !ô !ð !ðH%ð %ð %ð %ð %Ð(ñ %ô %ð %ð %ð %r1   