ó
    ýÞ j/Ÿ  ã                  ó*  • S SK Jr  S SKrS SKJrJrJrJr  S SKJ	r	  S SK
JrJr  S SK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Jr  S S
KJrJr  S SKJrJr  S SK J!r!  \(       a  S SK"J#r#J$r$  S SK%J&r&  \RN                  " \(5      r) " S S\5      r* " S S\5      r+ " S S\5      r, " S S\5      r-\S   r.S+S jr/ " S S\SS9r0 " S S\5      r1S,S jr2    S-S  jr3 " S! S"\15      r4 " S# S$\15      r5 " S% S&\5      r6 " S' S(\5      r7 " S) S*\5      r8g).é    )Ú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                  ód   ^ • \ rS rSrSrSrSrS
SU 4S jjjrSS jr\	SS j5       r
SS jrS	rU =r$ )ÚValuesTransformeré   uã  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)Úvaluesc                óŒ   >• [         TU ]  U5        [        5       U l        S U l        SU l        / U l        [        U5      U l        g )NF)	ÚsuperÚ__init__r   Ú_logÚ_latestÚ_interruptedÚ_interruptsÚlistÚ_scope_list©ÚselfÚscopeÚ	__class__s     €ÚW/var/www/html/gaurav/venv/lib/python3.13/site-packages/langgraph/stream/transformers.pyr"   ÚValuesTransformer.__init__1   s>   ø€ Ü‰Ñ˜ÔÜ3@³?ˆŒ	Ø.2ˆŒØ!ˆÔØ&(ˆÔô '+¨5£kˆÕó    c                ó   • SU R                   0$ )Nr   ©r#   ©r*   s    r-   ÚinitÚValuesTransformer.init;   ó   € Ø˜$Ÿ)™)Ð$Ð$r/   c                ó.   • U R                   R                  $ )zpThe error that ended the run, or `None` if it succeeded.

Set by the mux when it auto-fails the projection log.
)r#   Ú_errorr2   s    r-   ÚerrorÚValuesTransformer.error>   s   € ð �y‰y×ÑÐ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        U R                  R                  US   5        g)	NÚmethodr   TÚparamsÚ	namespaceÚdataÚ
interrupts© )r(   r$   Úgetr%   r&   Úextendr#   Úpush)r*   Úeventr<   r?   s       r-   ÚprocessÚValuesTransformer.processF   s�   € Ø�‰?˜hÓ&ØØ�x‘ˆØ�+Ñ $×"2Ñ"2Ó2ØØ˜f‘~ˆŒØ—Z‘Z ¨bÓ1ˆ
ÞØ $ˆDÔØ×Ñ×#Ñ# JÔ/Ø�	‰	�‰�v˜f‘~Ô&Ør/   )r%   r&   r$   r#   r(   ©r@   ©r+   útuple[str, ...]ÚreturnÚNone©rJ   údict[str, Any]©rJ   zBaseException | None©rD   r   rJ   Úbool)Ú__name__Ú
__module__Ú__qualname__Ú__firstlineno__Ú__doc__Ú_nativeÚrequired_stream_modesr"   r3   Úpropertyr8   rE   Ú__static_attributes__Ú__classcell__©r,   s   @r-   r   r      sB   ø† ñð" €GØ'Ð÷2ñ 2ô%ð ó ó ð ÷ò r/   r   c                  óP   ^ • \ rS rSrSrSrSrS	S
U 4S jjjrSS jrSS jr	Sr
U =r$ )ÚCustomTransformeréU   uÅ  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)Úcustomc                ób   >• [         TU ]  U5        [        5       U l        [	        U5      U l        g ©N©r!   r"   r   r#   r'   r(   r)   s     €r-   r"   ÚCustomTransformer.__init__f   s%   ø€ Ü‰Ñ˜ÔÜ(5«ˆŒ	Ü&*¨5£kˆÕr/   c                ó   • SU R                   0$ )Nr_   r1   r2   s    r-   r3   ÚCustomTransformer.initk   r5   r/   c                ó†   • US   S:w  a  gUS   nUS   U R                   :w  a  gU R                  R                  US   5        g)Nr;   r_   Tr<   r=   r>   ©r(   r#   rC   ©r*   rD   r<   s      r-   rE   ÚCustomTransformer.processn   sG   € Ø�‰?˜hÓ&ØØ�x‘ˆØ�+Ñ $×"2Ñ"2Ó2ØØ�	‰	�‰�v˜f‘~Ô&Ør/   ©r#   r(   rG   rH   rL   rO   ©rQ   rR   rS   rT   rU   rV   rW   r"   r3   rE   rY   rZ   r[   s   @r-   r]   r]   U   s.   ø† ñð €GØ'Ð÷2ñ 2ô
%÷ò r/   r]   c                  óP   ^ • \ rS rSrSrSrSrS	S
U 4S jjjrSS jrSS jr	Sr
U =r$ )ÚUpdatesTransformeréx   uÌ  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)Úupdatesc                ób   >• [         TU ]  U5        [        5       U l        [	        U5      U l        g ra   rb   r)   s     €r-   r"   ÚUpdatesTransformer.__init__‰   ó%   ø€ Ü‰Ñ˜ÔÜ3@³?ˆŒ	Ü&*¨5£kˆÕr/   c                ó   • SU R                   0$ )Nro   r1   r2   s    r-   r3   ÚUpdatesTransformer.initŽ   s   € Ø˜4Ÿ9™9Ð%Ð%r/   c                ó†   • US   S:w  a  gUS   nUS   U R                   :w  a  gU R                  R                  US   5        g)Nr;   ro   Tr<   r=   r>   rg   rh   s      r-   rE   ÚUpdatesTransformer.process‘   sG   € Ø�‰?˜iÓ'ØØ�x‘ˆØ�+Ñ $×"2Ñ"2Ó2ØØ�	‰	�‰�v˜f‘~Ô&Ør/   rj   rG   rH   rL   rO   rk   r[   s   @r-   rm   rm   x   s.   ø† ñð €GØ(Ð÷2ñ 2ô
&÷ò r/   rm   c                  ó¶   ^ • \ rS rSrSrSrSrSSU 4S j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rU =r$ )ÚMessagesTransformeré›   uº  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)Úmessagesc                óª   >• [         TU ]  U5        [        5       U l        0 U l        [        5       U l        S U l        S U l        [        U5      U l
        g ra   )r!   r"   r   r#   Ú_by_runÚsetÚ_ignored_runsÚ_pump_fnÚ	_apump_fnr'   r(   r)   s     €r-   r"   ÚMessagesTransformer.__init__È   sH   ø€ Ü‰Ñ˜ÔÜ4A³OˆŒ	ð 46ˆŒÜ'*£uˆÔØ37ˆŒØ?CˆŒô '+¨5£kˆÕr/   c                ó   • SU R                   0$ )Nrz   r1   r2   s    r-   r3   ÚMessagesTransformer.initÕ   s   € Ø˜DŸI™IÐ&Ð&r/   c                ó   • Xl         g)zIWire the sync pull callback. Called by GraphRunStream._wire_request_more.N)r   ©r*   Úfns     r-   Ú
_bind_pumpÚMessagesTransformer._bind_pumpØ   s   € à�r/   c                ó   • Xl         g)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)r€   r…   s     r-   Ú_bind_apumpÚMessagesTransformer._bind_apumpÜ   s	   € ð �r/   c               óì   • U R                   b(  [        UUUS9nUR                  U R                   5        U$ U R                  b(  [	        UUUS9nUR                  U R                  5        U$ [        UUUS9$ )a6  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.
©r=   ÚnodeÚ
message_id)r€   r	   Úset_arequest_morer   r
   Úset_request_more)r*   r=   rŽ   r�   ÚastreamÚstreams         r-   Ú_make_streamÚ MessagesTransformer._make_streamå   s†   € ð �>‰>Ñ%Ü*Ø#ØØ%ñˆGð
 ×%Ñ% d§n¡nÔ5ØˆNØ�=‰=Ñ$Ü&5Ø#ØØ%ñ'ˆFð
 ×#Ñ# D§M¡MÔ2ØˆMÜ#ØØØ!ñ
ð 	
r/   c                óÊ  • US   S:w  a  gUS   nUS   U R                   :w  a  gUS   u  p4UR                  S5      nU(       a  [        UR                  SS	5      5      OS	n[        U[        5      (       a!  S
U;   a  U R                  [        SU5      XeS9  g[        U[        5      (       a9  [        U[        5      (       d$  [        U[        5      (       d  U R                  X5S9  g)Nr;   rz   Tr<   r=   r>   Úlanggraph_nodeÚrun_idÚ rD   r   )r˜   rŽ   )rŽ   )r(   rA   ÚstrÚ
isinstanceÚdictÚ_route_protocol_eventr   r   r   r   Ú_route_whole_message)r*   rD   r<   ÚpayloadÚmetadatarŽ   r˜   s          r-   rE   ÚMessagesTransformer.process	  sÛ   € Ø�‰?˜jÓ(ØØ�x‘ˆØ�+Ñ $×"2Ñ"2Ó2Øà" 6™NÑˆØ#Ÿ<™<Ð(8Ó9ˆÞ4<”�X—\‘\ (¨BÓ/Ô0À"ˆä�gœt×$Ñ$¨°GÓ);Ø×&Ñ&Ü�^ WÓ-°fð 'ñ ð ô �w¤×,Ñ,Ü˜w¬×7Ñ7Ü˜w¬×4Ñ4à×%Ñ% gÐ%Ñ9ð
 r/   c               óV  • UR                  S5      nUS:X  aœ  UR                  S5      S:X  a  U R                  R                  U5        g UR                  S5      nU R                  / UUb  [	        U5      OS S9nX`R
                  U'   U R                  R                  U5        UR                  U5        g X R                  ;   a#  US:X  a  U R                  R                  U5        g g X R
                  ;   a5  U R
                  U   nUR                  U5        US:X  a  U R
                  U	 g g g )NrD   zmessage-startÚroleÚtoolr�   r�   zmessage-finish)
rA   r~   Úaddr”   rš   r|   r#   rC   ÚdispatchÚdiscard)r*   rD   r˜   rŽ   Ú
event_typer�   r“   s          r-   r�   Ú)MessagesTransformer._route_protocol_event$  s  € ð —Y‘Y˜wÓ'ˆ
Ø˜Ó(ð �y‰y˜Ó  FÓ*Ø×"Ñ"×&Ñ& vÔ.ØØŸ™ <Ó0ˆJØ×&Ñ&ØØØ.8Ñ.Dœ3˜zœ?È$ð 'ð ˆFð
 $*�L‰L˜Ñ Ø�I‰I�N‰N˜6Ô"Ø�O‰O˜EÕ"Ø×)Ñ)Ó)ØÐ-Ó-Ø×"Ñ"×*Ñ*¨6Õ2ð .à—|‘|Ó#Ø—\‘\ &Ñ)ˆFØ�O‰O˜EÔ"ØÐ-Ó-Ø—L‘L Ñ(ð .ð $r/   c               óÄ   • U R                  / X!R                  S9n[        XR                  S9 H  nUR                  U5        M     U R                  R                  U5        g )Nr�   )r�   )r”   Úidr   r¦   r#   rC   )r*   ÚmessagerŽ   r“   Úevts        r-   rž   Ú(MessagesTransformer._route_whole_messageD  sK   € Ø×"Ñ"¨R°dÇzÁzÐ"ÐRˆÜ$ W¿¹ÔDˆCØ�O‰O˜CÖ ñ Eà�	‰	�‰�vÕr/   c                ól   • U R                   R                  5         U R                  R                  5         g)uJ   Clear any routing state â€” streams close themselves via `message-finish`.N)r|   Úclearr~   r2   s    r-   ÚfinalizeÚMessagesTransformer.finalizeJ  s$   € à�‰×ÑÔØ×Ñ× Ñ Õ"r/   c                óâ   • [        U R                  R                  5       5       H  nUR                  U5        M     U R                  R	                  5         U R
                  R	                  5         g)zCPropagate run error to any streams still open when the graph fails.N)r'   r|   r   Úfailr°   r~   )r*   Úerrr“   s      r-   r´   ÚMessagesTransformer.failO  sL   € ä˜4Ÿ<™<×.Ñ.Ó0Ö1ˆFØ�K‰K˜Öñ 2à�‰×ÑÔØ×Ñ× Ñ Õ"r/   )r€   r|   r~   r#   r   r(   rG   rH   rL   )r†   zCallable[[], bool]rJ   rK   )r†   zCallable[[], Awaitable[bool]]rJ   rK   )r=   ú	list[str]rŽ   ú
str | Noner�   r¸   rJ   r
   rO   )rD   r   r˜   rš   rŽ   r¸   rJ   rK   )r¬   r   rŽ   r¸   rJ   rK   ©rJ   rK   ©rµ   ÚBaseExceptionrJ   rK   )rQ   rR   rS   rT   rU   rV   rW   r"   r3   r‡   rŠ   r”   rE   r�   rž   r±   r´   rY   rZ   r[   s   @r-   rx   rx   ›   s¢   ø† ñ'ðR €GØ)Ð÷2ñ 2ô'ôôð"
ð ð"
ð ð	"
ð
 ð"
ð 
ô"
ôHð6)àð)ð ð	)ð
 ð)ð 
ô)ô@ô#÷
#ò #r/   rx   )ÚstartedÚ	completedÚfailedÚinterruptedÚdrainedc                óD   • U R                  S5      u  pnX(       a  U4$ S4$ )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)ÚsegmentÚnameÚsepÚtask_ids       r-   Ú_parse_ns_segmentrÈ   Z  s.   € ð !×*Ñ*¨3Ó/Ñ€DˆwØ�C�Ð)Ð) TÐ)Ð)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)ÚLifecyclePayloadid  a  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`.
ÚSubgraphStatusrD   r·   r=   zNotRequired[str]Ú
graph_nameÚtrigger_call_idzNotRequired[LifecycleCause]Úcauser8   r@   N)rQ   rR   rS   rT   rU   Ú__annotations__rY   r@   r/   r-   rÊ   rÊ   d  s-   ‡ ñð ÓØÓØ Ó Ø%Ó%Ø&Ó&ØÖr/   rÊ   F)Útotalc                  óÒ   ^ • \ rS rSrSrSrSSU 4S j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S jrSrU =r$ )Ú_TasksLifecycleBaseiu  uµ  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.
©Útasksc                óz   >• [         TU ]  U5        [        5       U l        0 U l        0 U l        0 U l        S U l        g ra   )r!   r"   r}   Ú_seenÚ_openÚ	_lc_by_nsÚ_pending_tool_callsÚ_pending_causer)   s     €r-   r"   Ú_TasksLifecycleBase.__init__”  s?   ø€ Ü‰Ñ˜ÔÜ+.«5ˆŒ
ð 24ˆŒ
ð =?ˆŒð 46ˆÔ ð
 6:ˆÕr/   c                ó   • [         e)uF   Scope filter â€” return True iff `ns` is in this transformer's region.©ÚNotImplementedError©r*   Únss     r-   Ú_should_trackÚ!_TasksLifecycleBase._should_track®  s   € ä!Ð!r/   c                ó   • [         e)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       r-   Ú_on_startedÚ_TasksLifecycleBase._on_started²  s
   € ô "Ð!r/   c                ó   • [         e)zkFired once per tracked namespace when its parent's result arrives,
or via finalize/fail safety-net sweeps.
rÝ   )r*   rà   Ústatusr8   s       r-   Ú_on_terminalÚ _TasksLifecycleBase._on_terminal¾  s
   € ô "Ð!r/   c                óà   • US   S:w  a  g[        US   S   5      nUS   S   nSU;   a  U R                  X#5        gU R                  X#5        U R                  U5        U R	                  X#5        g)	Nr;   rÔ   Tr<   r=   r>   ÚresultF)ÚtupleÚ_handle_task_resultÚ_record_identityÚ_record_pending_tool_callsÚ_handle_task_start)r*   rD   rà   r>   s       r-   rE   Ú_TasksLifecycleBase.processË  s~   € Ø�‰?˜gÓ%ØÜ�5˜‘? ;Ñ/Ó0ˆØ�X‰˜vÑ&ˆØ�tÓØ×$Ñ$ RÔ.ð ð ×!Ñ! "Ô+Ø×+Ñ+¨DÔ1Ø×#Ñ# BÔ-ð r/   c                ó”   • XR                   ;   a  gUR                  S5      =(       d    0 nUR                  S5      U R                   U'   g)aQ  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Ø   rA   )r*   rà   r>   r    s       r-   rî   Ú$_TasksLifecycleBase._record_identityÛ  s;   € ð —‘ÓØØ—8‘8˜JÓ'×-¨2ˆØ%Ÿ\™\¨/Ó:ˆ�‰�rÒr/   c                ó&  • UR                  S5      n[        U[        5      (       d  gUR                  S5      nSn[        U[        5      (       aP  [        UR                  S5      [        5      (       a,  US   R                  S5      n[        U[        5      (       a  UnO`[        U[        5      (       aK  U HE  n[        U[        5      (       d  M  [        UR                  S5      [        5      (       d  M@  US   n  O   Ub  X@R
                  U'   gg)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)rA   r›   rš   rœ   r'   rÙ   )r*   r>   rÇ   rŸ   Útool_call_idÚ	candidateÚtcs          r-   rï   Ú._TasksLifecycleBase._record_pending_tool_callsè  sß   € ð —(‘(˜4“.ˆÜ˜'¤3×'Ñ'ØØ—(‘(˜7Ó#ˆØ#'ˆô �gœt×$Ñ$¬°G·K±KÀÓ4LÌd×)SÑ)SØ Ñ,×0Ñ0°Ó6ˆIÜ˜)¤S×)Ñ)Ø(�øä˜¤×&Ñ&Û�Ü˜b¤$×'Ó'¬J°r·v±v¸d³|ÄS×,IÓ,IØ#% d¡8�LÙñ ð Ñ#Ø0<×$Ñ$ WÒ-ð $r/   c                óø  • U R                  U5      (       a  XR                  ;   a  g U R                  R                  U5        [        US   5      u  p4UR	                  S5      =(       d    0 nUR	                  S5      nUS LnU(       a  UO
U=(       d    S nS n	U(       a3  Ub0  U R
                  R	                  U5      n
U
(       a  S[        U
5      S.n	X�l        U R                  XU5        Ub  X@R                  U'   g g )Néÿÿÿÿr    ró   ÚtoolCall)Útyperø   )
rá   rÖ   r¥   rÈ   rA   rÙ   rš   rÚ   rä   r×   )r*   rà   r>   Úparsed_namerÍ   r    Úchild_lcÚis_subagentrÌ   rÎ   rø   s              r-   rð   Ú&_TasksLifecycleBase._handle_task_start  sÞ   € Ø×!Ñ! "×%Ñ%¨¯z©zÓ)9ØØ�
‰
�‰�rÔÜ'8¸¸B¹Ó'@Ñ$ˆØ—8‘8˜JÓ'×-¨2ˆØ—<‘< Ó0ˆð  dÐ*ˆÞ!,‘X°;×3FÀ$ˆ
Ø'+ˆÞ˜?Ñ6Ø×3Ñ3×7Ñ7¸ÓHˆLÞØ!+¼SÀÓ=NÑO�ð $ÔØ×Ñ˜¨Ô9ØÑ&Ø,�J‰J�rŠNð 'r/   c                ó  • UR                  S5      nU(       d  / $ / n[        U R                  R                  5       5       HB  u  pVUSS U:w  d  Xc:w  a  M  [	        U5      u  pxUR                  XWU45        U R                  U	 MD     U$ )z>Return and remove tracked children closed by this task result.r«   Nrý   )rA   r'   r×   ÚitemsÚ_terminal_from_resultÚappend)	r*   rà   r>   Ú	result_idÚtransitionsÚchild_nsÚparent_task_idrç   r8   s	            r-   Ú_pop_terminal_transitionsÚ-_TasksLifecycleBase._pop_terminal_transitions$  s‡   € ð —H‘H˜T“Nˆ	ÞØˆIØPRˆÜ(,¨T¯Z©Z×-=Ñ-=Ó-?Ö(@Ñ$ˆHØ˜˜ˆ} Ó" nÓ&AÙÜ1°$Ó7‰MˆFØ×Ñ °%Ð8Ô9Ø—
‘
˜8Ò$ñ )Að Ðr/   c                ó^   • U R                  X5       H  u  p4nU R                  X4U5        M     g ra   )r  rè   )r*   rà   r>   r
  rç   r8   s         r-   rí   Ú'_TasksLifecycleBase._handle_task_result4  s-   € Ø'+×'EÑ'EÀbÖ'OÑ#ˆH˜eØ×Ñ˜h°Ö6ò (Pr/   c                ó–   • [        U R                  5       H  nU R                  USS5        M     U R                  R                  5         g)zAEmit `completed` for any tracked namespace still open at run end.r½   N)r'   r×   rè   r°   rß   s     r-   r±   Ú_TasksLifecycleBase.finalize8  s7   € ä�t—z‘zÖ"ˆBØ×Ñ˜b +¨tÖ4ñ #à�
‰
×ÑÕr/   c                ó®   • [        U5      u  p#[        U R                  5       H  nU R                  XBU5        M     U R                  R	                  5         g)z:Emit terminal status for any tracked namespace still open.N)Ú_status_from_exceptionr'   r×   rè   r°   )r*   rµ   rç   Ú	error_strrà   s        r-   r´   Ú_TasksLifecycleBase.fail>  sB   € ä2°3Ó7ÑˆÜ�t—z‘zÖ"ˆBØ×Ñ˜b¨)Ö4ñ #à�
‰
×ÑÕr/   )rØ   r×   rÚ   rÙ   rÖ   rG   rH   ©rà   rI   rJ   rP   ©rà   rI   rÌ   r¸   rÍ   r¸   rJ   rK   ©rà   rI   rç   rË   r8   r¸   rJ   rK   rO   )rà   rI   r>   rM   rJ   rK   )r>   rM   rJ   rK   )rà   rI   r>   rM   rJ   z8list[tuple[tuple[str, ...], SubgraphStatus, str | None]]r¹   rº   )rQ   rR   rS   rT   rU   rW   r"   rá   rä   rè   rE   rî   rï   rð   r  rí   r±   r´   rY   rZ   r[   s   @r-   rÒ   rÒ   u  sº   ø† ñð8 'Ð÷:ñ :ô4"ð
"àð
"ð ð
"ð $ð	
"ð
 
ô
"ð	"àð	"ð ð	"ð ð		"ð
 
ô	"ôô ;ô=ô<-ð<Ø!ðØ)7ðà	Aôô 7ô÷ò r/   rÒ   c                ót   • [        U [        5      (       a  g[        U [        5      (       a  gS[        U 5      4$ )zCMap a run exception to a subgraph terminal status and error string.)rÀ   N©r¿   Nr¾   )r›   r   r   rš   )rµ   s    r-   r  r  F  s1   € ä�#”|×$Ñ$ØÜ�#”~×&Ñ&Ø"Ø”S˜“XÐÐr/   c                ó|   • U R                  S5      (       a  gU R                  S5      nU(       a  S[        U5      4$ g)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`.
r?   r  r8   r¾   )r½   N)rA   rš   )rŸ   r8   s     r-   r  r  O  s9   € ð ‡{�{�<× Ñ Ø"Ø�K‰K˜Ó €EÞØœ˜U›Ð#Ð#Ør/   c                  ó€   ^ • \ rS rSrSrSrS
SU 4S jjjrSS jrSS jr        SS jr	        SS jr
S	rU =r$ )ÚLifecycleTransformeri`  u‚  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`.
Tc                óD   >• [         TU ]  U5        [        S5      U l        g ©NÚ	lifecycle)r!   r"   r   Ú_channelr)   s     €r-   r"   ÚLifecycleTransformer.__init__s  s   ø€ Ü‰Ñ˜ÔÜ9FÀ{Ó9Sˆ�r/   c                ó   • SU R                   0$ r  ©r!  r2   s    r-   r3   ÚLifecycleTransformer.initw  s   € Ø˜TŸ]™]Ð+Ð+r/   c                óz   • [        U R                  5      n[        U5      U:„  =(       a    US U U R                  :H  $ ra   ©Úlenr+   ©r*   rà   Údepths      r-   rá   Ú"LifecycleTransformer._should_trackz  s1   € Ü�D—J‘J“ˆÜ�2‹w˜‰×; 2 f u :°·±Ñ#;Ð;r/   c                ó¢   • Uc  g S[        U5      S.nU(       a  X$S'   X4S'   U R                  nUb  XTS'   U R                  R                  U5        g )Nr¼   ©rD   r=   rÌ   rÍ   rÎ   )r'   rÚ   r!  rC   )r*   rà   rÌ   rÍ   rŸ   rÎ   s         r-   rä   Ú LifecycleTransformer._on_started~  s\   € ð Ñ"ð Ø.7ÄdÈ2ÃhÑ$OˆÞØ$.�LÑ!Ø%4Ð!Ñ"Ø×#Ñ#ˆØÑØ$�GÑØ�‰×Ñ˜7Õ#r/   c                ód   • U[        U5      S.nUb  X4S'   U R                  R                  U5        g )Nr-  r8   )r'   r!  rC   )r*   rà   rç   r8   rŸ   s        r-   rè   Ú!LifecycleTransformer._on_terminal’  s2   € ð /5Ä4ÈÃ8Ñ$LˆØÑØ$�GÑØ�‰×Ñ˜7Õ#r/   r$  rG   rH   rL   r  r  r  )rQ   rR   rS   rT   rU   rV   r"   r3   rá   rä   rè   rY   rZ   r[   s   @r-   r  r  `  s€   ø† ñð  €G÷Tñ Tô,ô<ð$àð$ð ð$ð $ð	$ð
 
ô$ð(	$àð	$ð ð	$ð ð		$ð
 
÷	$ò 	$r/   r  c                  ó^  ^ • \ rS rSrSrSrSrSSU 4S j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"U 4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rU =r$ )&ÚSubgraphTransformeriž  uó  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`.
Tc                ó^   >• [         TU ]  U5        [        5       U l        0 U l        S U l        g ra   )r!   r"   r   r#   Ú_handlesÚ_muxr)   s     €r-   r"   ÚSubgraphTransformer.__init__³  s1   ø€ Ü‰Ñ˜Ôä‹Oð 	Œ	ð
 ð 	Œð '+ˆ�	r/   c                ó   • SU R                   0$ )NÚ	subgraphsr1   r2   s    r-   r3   ÚSubgraphTransformer.init½  s   € Ø˜TŸY™YÐ'Ð'r/   c                ó   • Xl         g ra   )r5  )r*   Úmuxs     r-   Ú_on_registerÚ SubgraphTransformer._on_registerÀ  s   € Ø�	r/   c                ó€   • [        U R                  5      n[        U5      US-   :H  =(       a    US U U R                  :H  $ )Né   r'  r)  s      r-   rá   Ú!SubgraphTransformer._should_trackÃ  s8   € ô �D—J‘J“ˆÜ�2‹w˜% !™)Ñ#×@¨¨6¨E¨
°d·j±jÑ(@Ð@r/   c                ó  • U R                   c  g  U R                   R                  U5      nUR                  (       a  [        O[
        nU" UUUUS9nX`R                  U'   U R                  R                  U5        g ! [         a     g f = f)N)r;  ÚpathrÌ   rÍ   )	r5  Ú_make_childÚRuntimeErrorÚis_asyncr   r   r4  r#   rC   )r*   rà   rÌ   rÍ   Ú	child_muxÚ
handle_clsÚhandles          r-   rä   ÚSubgraphTransformer._on_startedÉ  s‡   € ð �9‰9ÑØð	ØŸ	™	×-Ñ-¨bÓ1ˆIð 09×/A×/AÕ+ÔGXˆ
ñ ØØØ!Ø+ñ	
ˆð #�‰�bÑØ�	‰	�‰�vÕøô ó 	Ùð	ús   �A; Á;
BÂBc                ó”   • U R                   R                  U5      nUb  U R                  XBU5      (       d  g U R                  XBU5        g ra   )r4  rA   Ú_mark_terminalÚ_close_or_fail_handle©r*   rà   rç   r8   rH  s        r-   rè   Ú SubgraphTransformer._on_terminalâ  sB   € ð —‘×"Ñ" 2Ó&ˆØ‰> ×!4Ñ!4°VÀU×!KÑ!KØØ×"Ñ" 6°5Õ9r/   c              ƒ  ó°   #   • U R                   R                  U5      nUb  U R                  XBU5      (       d  g U R                  XBU5      I S h  v•N   g  N7fra   )r4  rA   rK  Ú_aclose_or_fail_handlerM  s        r-   Ú_aon_terminalÚ!SubgraphTransformer._aon_terminalí  sK   é € ð —‘×"Ñ" 2Ó&ˆØ‰> ×!4Ñ!4°VÀU×!KÑ!KØØ×)Ñ)¨&¸%Ó@×@Ó@ùs   ‚AAÁAÁAc                ón   • UR                   (       a  gX!l        Ub  UR                  c  X1l        SUl         g)z>Mark a handle terminal once. Returns True on first transition.FT)Ú_seen_terminalrç   r8   ©r*   rH  rç   r8   s       r-   rK  Ú"SubgraphTransformer._mark_terminalø  s4   € ð × × ØØŒØÑ §¡Ñ!5Ø ŒLØ $ˆÔØr/   c                ó  • UR                   b%  UR                   R                  R                  (       a  g US:X  a.  UR                   R                  [	        U=(       d    S5      5        g UR                   R                  5         g ©Nr¾   zSubgraph failed)r5  Ú_eventsÚ_closedr´   rD  ÚcloserU  s       r-   rL  Ú)SubgraphTransformer._close_or_fail_handle  sX   € ð �;‰;Ñ &§+¡+×"5Ñ"5×"=×"=ØØ�XÓØ�K‰K×Ñœ\¨%×*DÐ3DÓEÕFà�K‰K×ÑÕr/   c              ƒ  ó6  #   • UR                   b%  UR                   R                  R                  (       a  g US:X  a6  UR                   R                  [	        U=(       d    S5      5      I S h  v•N   g UR                   R                  5       I S h  v•N   g  N( N7frX  )r5  rY  rZ  ÚafailrD  ÚacloserU  s       r-   rP  Ú*SubgraphTransformer._aclose_or_fail_handle  sp   é € ð �;‰;Ñ &§+¡+×"5Ñ"5×"=×"=ØØ�XÓØ—+‘+×#Ñ#¤L°×1KÐ:KÓ$LÓM×MÑMà—+‘+×$Ñ$Ó&×&Ñ&ñ Ná&ùs$   ‚A*BÁ,BÁ-"BÂBÂBÂBc                ó&  • [        US   S   5      n[        U R                  5      n[        U5      US-   :  a  g U R                  R	                  US US-    5      nUb2  UR
                  b%  UR
                  R                  R                  (       a  g U$ )Nr<   r=   r?  )rì   r(  r+   r4  rA   r5  rY  rZ  )r*   rD   rà   r*  rH  s        r-   Ú_handle_for_eventÚ%SubgraphTransformer._handle_for_event!  s}   € ô �5˜‘? ;Ñ/Ó0ˆÜ�D—J‘J“ˆÜˆr‹7�U˜Q‘YÓØØ—‘×"Ñ" 2 k¨°©	 ?Ó3ˆØ‰>˜VŸ[™[Ñ0°F·K±K×4GÑ4G×4O×4OØØˆr/   c                ó¦   >• [         TU ]  U5      nU R                  U5      nUb,  UR                  U5        UR                  R                  U5        U$ ra   )r!   rE   rb  Ú_observe_eventr5  rC   )r*   rD   ÚkeeprH  r,   s       €r-   rE   ÚSubgraphTransformer.process-  sN   ø€ ô ‰w‰˜uÓ%ˆØ×'Ñ'¨Ó.ˆØÑØ×!Ñ! %Ô(Ø�K‰K×Ñ˜UÔ#Øˆr/   c              ƒ  óÒ  #   • US   S:X  a‹  [        US   S   5      nUS   S   nSU;   a6  U R                  X#5       H   u  pEnU R                  XEU5      I S h  v•N   M"     O3U R                  X#5        U R	                  U5        U R                  X#5        SnOSnU R                  U5      nUb4  UR                  U5        UR                  R                  U5      I S h  v•N   U$  N‹ N7f)	Nr;   rÔ   r<   r=   r>   rë   FT)
rì   r  rQ  rî   rï   rð   rb  re  r5  Úapush)	r*   rD   rà   r>   r
  rç   r8   rf  rH  s	            r-   ÚaprocessÚSubgraphTransformer.aprocess7  sï   é € ð �‰?˜gÓ%Ü�u˜X‘ {Ñ3Ó4ˆBØ˜‘? 6Ñ*ˆDØ˜4ÓØ/3×/MÑ/MÈbÖ/WÑ+�H eØ×,Ñ,¨X¸uÓE×EÒEò 0Xð ×%Ñ% bÔ/Ø×/Ñ/°Ô5Ø×'Ñ'¨Ô1Ø‰DàˆDØ×'Ñ'¨Ó.ˆØÑØ×!Ñ! %Ô(Ø—+‘+×#Ñ# EÓ*×*Ð*Øˆñ Fñ +ùs%   ‚AC'ÁC#ÁBC'ÃC%ÃC'Ã%C'c                óÈ  • S n[        U R                  5       H  n U R                  USS 5        M     U R                  R	                  5         U R
                  R                  5        H1  nU R                  USS 5      (       d  M   U R                  USS 5        M3     U$ ! [         a  nUc  Un S nAM›   S nAM¡  S nAff = f! [         a  nUc  Un S nAMo   S nAMu  S nAff = f©Nr½   )	r'   r×   rè   r»   r°   r4  r   rK  rL  ©r*   Úfirst_errorrà   ÚerH  s        r-   Ú_complete_open_handlesÚ*SubgraphTransformer._complete_open_handlesP  sÑ   € Ø,0ˆÜ�t—z‘zÖ"ˆBð$Ø×!Ñ! " k°4Ö8ñ #ð 	�
‰
×ÑÔØ—m‘m×*Ñ*Ö,ˆFØ×"Ñ" 6¨;¸×=Ó=ð(Ø×.Ñ.¨v°{ÀDÖIñ -ð Ðøô !ó $ØÑ&Ø"#–Kõ 'ûð$ûô %ó (Ø"Ñ*Ø&'žõ +ûð(ús/   œBÂCÂ
B>Â(B9Â9B>Ã
C!ÃCÃC!c              ƒ  óø  #   • S n[        U R                  5       H  n U R                  USS 5      I S h  v•N   M!     U R                  R	                  5         U R
                  R                  5        H9  nU R                  USS 5      (       d  M   U R                  USS 5      I S h  v•N   M;     U$  N{! [         a  nUc  Un S nAM­   S nAM³  S nAff = f N/! [         a  nUc  Un S nAM{   S nAM�  S nAff = f7frm  )	r'   r×   rQ  r»   r°   r4  r   rK  rP  rn  s        r-   Ú_acomplete_open_handlesÚ+SubgraphTransformer._acomplete_open_handlesb  sè   é € Ø,0ˆÜ�t—z‘zÖ"ˆBð$Ø×(Ñ(¨¨[¸$Ó?×?Ò?ñ #ð 	�
‰
×ÑÔØ—m‘m×*Ñ*Ö,ˆFØ×"Ñ" 6¨;¸×=Ó=ð(Ø×5Ñ5°f¸kÈ4ÓP×PÒPñ -ð Ðñ @øÜ ó $ØÑ&Ø"#–Kõ 'ûð$úñ QøÜ$ó (Ø"Ñ*Ø&'žõ +ûð(üsz   ‚C:žB2´B0µB2¹AC:ÂCÂ%CÂ&CÂ*C:Â0B2Â2
CÂ<CÃC:ÃCÃC:ÃCÃ
C7Ã!C2Ã&C:Ã2C7Ã7C:c                ó.   • U R                  5       nUb  Ueg ra   )rq  ©r*   ro  s     r-   r±   ÚSubgraphTransformer.finalizet  s!   € Ø×1Ñ1Ó3ˆØÑ"ØÐð #r/   c              ƒ  óJ   #   • U R                  5       I S h  v•N nUb  Ueg  N
7fra   )rt  rw  s     r-   Ú	afinalizeÚSubgraphTransformer.afinalizey  s,   é € Ø ×8Ñ8Ó:×:ˆØÑ"ØÐð #ñ ;ùs   ‚#–!—#c                ó¼  • [        U5      u  p#U R                  R                  5         U R                  R	                  5        Hg  nU R                  XBU5        UR                  c  M$  UR                  R                  R                  (       a  MK   UR                  R                  U5        Mi     g ! [         a#    [        R                  SUR                  SS9   M˜  f = f©NzRError failing subgraph mini-mux at %s; subscribers may not see the terminal error.T)Úexc_info)r  r×   r°   r4  r   rK  r5  rY  rZ  r´   Ú	ExceptionÚ_loggerÚwarningrB  ©r*   rµ   rç   r  rH  s        r-   r´   ÚSubgraphTransformer.fail~  s®   € Ü2°3Ó7ÑˆØ�
‰
×ÑÔØ—m‘m×*Ñ*Ö,ˆFØ×Ñ °	Ô:Ø�{‰{Ó&¨v¯{©{×/BÑ/B×/J×/JÑ/JðØ—K‘K×$Ñ$ SÖ)ò	 -øô
 !ó Ü—O‘OðFàŸ™Ø!%ð	 $ô ðús   ÂB.Â.)CÃCc              ƒ  óØ  #   • [        U5      u  p#U R                  R                  5         U R                  R	                  5        Ho  nU R                  XBU5        UR                  c  M$  UR                  R                  R                  (       a  MK   UR                  R                  U5      I S h  v•N   Mq     g  N	! [         a#    [        R                  SUR                  SS9   M¢  f = f7fr}  )r  r×   r°   r4  r   rK  r5  rY  rZ  r^  r  r€  r�  rB  r‚  s        r-   r^  ÚSubgraphTransformer.afailŽ  s¹   é € Ü2°3Ó7ÑˆØ�
‰
×ÑÔØ—m‘m×*Ñ*Ö,ˆFØ×Ñ °	Ô:Ø�{‰{Ó&¨v¯{©{×/BÑ/B×/J×/JÑ/JðØ Ÿ+™+×+Ñ+¨CÓ0×0Ò0ò	 -ñ 1øÜ ó Ü—O‘OðFàŸ™Ø!%ð	 $ô ðüsB   ‚A"C*Á(#C*ÂB:Â.B8Â/B:Â3C*Â8B:Â:)C'Ã#C*Ã&C'Ã'C*)r4  r#   r5  rG   rH   rL   )r;  r   rJ   rK   r  r  r  )rH  ú*SubgraphRunStream | AsyncSubgraphRunStreamrç   rË   r8   r¸   rJ   rP   )rH  r†  rç   rË   r8   r¸   rJ   rK   )rD   r   rJ   z1SubgraphRunStream | AsyncSubgraphRunStream | NonerO   rN   r¹   rº   )rQ   rR   rS   rT   rU   rV   Úsupports_syncr"   r3   r<  rá   rä   rè   rQ  rK  rL  rP  rb  rE   rj  rq  rt  r±   rz  r´   r^  rY   rZ   r[   s   @r-   r2  r2  ž  sw  ø† ñð" €GØ€M÷+ñ +ô(ôôAðàðð ðð $ð	ð
 
ôð2	:àð	:ð ð	:ð ð		:ð
 
ô	:ð	Aàð	Að ð	Að ð		Að
 
ô	Aðà:ðð ðð ð	ð
 
ôð à:ð ð ð ð ð	 ð
 
ô ð'à:ð'ð ð'ð ð	'ð
 
ô'ð
Ø"ð
à	:ô
÷ôô2ô$ô$ô
ô
÷ ò r/   r2  c                  óP   ^ • \ rS rSrSrSrSrS	S
U 4S jjjrSS jrSS jr	Sr
U =r$ )ÚCheckpointsTransformeriŸ  u]  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)Úcheckpointsc                ób   >• [         TU ]  U5        [        5       U l        [	        U5      U l        g ra   rb   r)   s     €r-   r"   ÚCheckpointsTransformer.__init__´  rr   r/   c                ó   • SU R                   0$ )NrŠ  r1   r2   s    r-   r3   ÚCheckpointsTransformer.init¹  s   € Ø˜tŸy™yÐ)Ð)r/   c                ó†   • US   S:w  a  gUS   nUS   U R                   :w  a  gU R                  R                  US   5        g)Nr;   rŠ  Tr<   r=   r>   rg   rh   s      r-   rE   ÚCheckpointsTransformer.process¼  sG   € Ø�‰?˜mÓ+ØØ�x‘ˆØ�+Ñ $×"2Ñ"2Ó2ØØ�	‰	�‰�v˜f‘~Ô&Ør/   rj   rG   rH   rL   rO   rk   r[   s   @r-   r‰  r‰  Ÿ  s.   ø† ñð" €GØ,Ð÷2ñ 2ô
*÷ò r/   r‰  c                  óP   ^ • \ rS rSrSrSrSrS	S
U 4S jjjrSS jrSS jr	Sr
U =r$ )ÚDebugTransformeriÆ  uì  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)Údebugc                ób   >• [         TU ]  U5        [        5       U l        [	        U5      U l        g ra   rb   r)   s     €r-   r"   ÚDebugTransformer.__init__Ø  rr   r/   c                ó   • SU R                   0$ )Nr“  r1   r2   s    r-   r3   ÚDebugTransformer.initÝ  ó   € Ø˜Ÿ™Ð#Ð#r/   c                ó†   • US   S:w  a  gUS   nUS   U R                   :w  a  gU R                  R                  US   5        g)Nr;   r“  Tr<   r=   r>   rg   rh   s      r-   rE   ÚDebugTransformer.processà  óG   € Ø�‰?˜gÓ%ØØ�x‘ˆØ�+Ñ $×"2Ñ"2Ó2ØØ�	‰	�‰�v˜f‘~Ô&Ør/   rj   rG   rH   rL   rO   rk   r[   s   @r-   r’  r’  Æ  s.   ø† ñð €GØ&Ð÷2ñ 2ô
$÷ò r/   r’  c                  óP   ^ • \ rS rSrSrSrSrS	S
U 4S jjjrSS jrSS jr	Sr
U =r$ )ÚTasksTransformeriê  uk  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Ó   c                ób   >• [         TU ]  U5        [        5       U l        [	        U5      U l        g ra   rb   r)   s     €r-   r"   ÚTasksTransformer.__init__   rr   r/   c                ó   • SU R                   0$ )NrÔ   r1   r2   s    r-   r3   ÚTasksTransformer.init  r˜  r/   c                ó†   • US   S:w  a  gUS   nUS   U R                   :w  a  gU R                  R                  US   5        g)Nr;   rÔ   Tr<   r=   r>   rg   rh   s      r-   rE   ÚTasksTransformer.process  r›  r/   rj   rG   rH   rL   rO   rk   r[   s   @r-   r�  r�  ê  s.   ø† ñð$ €GØ&Ð÷2ñ 2ô
$÷ò r/   r�  )rÄ   rš   rJ   ztuple[str, str | None])rµ   r»   rJ   ú!tuple[SubgraphStatus, str | None])rŸ   rM   rJ   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   Ú	getLoggerrQ   r€  r   r]   rm   rx   rË   rÈ   rÊ   rÒ   r  r  r  r2  r‰  r’  r�  r@   r/   r-   Ú<module>r´     s  ðÝ "ã ß 4Ó 4å K÷÷ MÑ Lß Dß 4ç 9ß Dß QÝ 9æß3å/à
×
Ò
˜HÓ
%€ô6Ð)ô 6ôr Ð)ô  ôF Ð*ô  ôFy#Ð+ô y#ðx ÐSÑT€ô*ô�y¨ò ô"NÐ+ô NôbðØðà&ôô";$Ð.ô ;$ô|~Ð-ô ~ôB$Ð.ô $ôN!Ð(ô !ôH%Ð(õ %r/   