ó
    ýÞ jå5  ã                  ó*  • S r SSKJr  SSKJrJrJr  SSKJrJ	r	J
r
  Sr\" 1 Sk5      rSS jrSS jrSS	 jrSS S jjr\	S   rS!S jr    S"S jrS#S jr " S S\
5      r " S S5      r " S S5      r " S S5      r " S S5      r " S S5      rg
)$u;  Per-channel event â†’ items state machines.

Used both by the projection iterators (`_ValuesProjection`,
`_MessagesProjection`, `_ToolCallsProjection`, `_SubgraphsProjection`) on
`AsyncThreadStream` / `SyncThreadStream`, and by `interleave_projections`,
which drives multiple decoders from one shared subscription.
é    )Úannotations)ÚCallableÚIterableÚMapping)ÚAnyÚLiteralÚProtocol)ÚvaluesÚmessagesÚ
tool_callsÚ	subgraphsÚupdatesÚcheckpointsÚtasks>   ÚinputÚtoolsÚ	lifecyclec           
     óŽ   • U  H?  nU[         ;   d  M  US:X  a  SOSn[        U< SU SSR                  [        5       S35      e   g)	a9  Reject reserved protocol channel names before they hit the fallback.

Genuine extension names pass through untouched; only names that
``infer_channel`` treats as built-in methods without an interleave decoder
are rejected, so a typo'd or unsupported protocol channel surfaces an error
instead of an empty stream.
r   z (use "tool_calls")Ú z. is not a valid interleave_projections channelz. Supported channels: z, z, or an extension name.N)ÚRESERVED_INTERLEAVE_CHANNELSÚ
ValueErrorÚjoinÚSUPPORTED_INTERLEAVE_CHANNELS)ÚchannelsÚchÚhints      ÚW/var/www/html/gaurav/venv/lib/python3.13/site-packages/langgraph_sdk/stream/decoders.pyÚvalidate_interleave_channelsr   "   s_   € ó ˆØÔ-Õ-Ø,.°'«MÑ(¸rˆDÜØ‘&ÐFÀtÀfð M'Ø'+§y¡yÔ1NÓ'OÐ&Pð Q(ð(óð ò ó    c                ó¨   • [        U [        5      (       d  / $ U R                  S5      =(       d    / n[        U[        5      (       a  [        U5      $ / $ )NÚ	namespace)Ú
isinstanceÚdictÚgetÚlist)Úparams_fieldr!   s     r   Ú_event_namespacer'   4   sD   € Ü�l¤D×)Ñ)Øˆ	Ø× Ñ  Ó-×3°€IÜ(¨´D×9Ñ9Œ4�	‹?ÐA¸rÐAr   c                ót   • U R                  S5      =(       d    U R                  S5      nUb  [        U5      $ S $ )NÚidÚ
message_id©r$   Ústr)Údatar*   s     r   Ú_message_event_idr.   ;   s1   € Ø—‘˜$“×9 4§8¡8¨LÓ#9€JØ(Ñ4Œ3ˆz‹?Ð>¸$Ð>r   Nc                ó:   • [        U 5      nUb  SU 3$ Ub  SU 3$ g)a  Return the routing key for a message-channel event in `active`.

Keys on `message_id` when available so concurrent messages that share the
same `run_id` (two AI turns in one agent step) route to independent streams
rather than colliding on a shared `run:<run_id>` slot.
zmessage:Ú
__single__)r.   )r-   Úfallbackr*   s      r   Ú_message_route_keyr2   @   s7   € ô # 4Ó(€JØÑØ˜*˜Ð&Ð&ØÑØ˜(˜Ð$Ð$Ør   )ÚstartedÚ	completedÚfailedÚinterruptedc                óD   • U R                  S5      u  pnX(       a  U4$ S 4$ )NÚ:)Ú	partition)ÚsegmentÚnameÚsepÚtask_ids       r   Ú_parse_namespace_segmentr>   R   s,   € Ø ×*Ñ*¨3Ó/Ñ€DˆwØ�C�Ð)Ð) TÐ)Ð)r   c                ó|   • U R                  S5      (       a  gU R                  S5      nU(       a  S[        U5      4$ g)NÚ
interrupts)r6   NÚerrorr5   )r4   Nr+   )r-   rA   s     r   Ú_terminal_from_tasks_resultrB   W   s9   € ð ‡x�x�×ÑØ"Ø�H‰H�WÓ€EÞØœ˜U›Ð#Ð#Ør   c                óx   • [        U 5      [        U5      S-   :H  =(       a    [        U S [        U5       5      U:H  $ )Né   )ÚlenÚtuple)r!   Úscopes     r   Ú_is_direct_childrH   b   s4   € Üˆy‹>œS ›Z¨!™^Ñ+×W´°iÀÄ#ÀeÃ*Ð6MÓ0NÐRWÑ0WÐWr   c                  ó   • \ rS rSrSS jrSrg)ÚDecoderéf   c                ó   • g ©N© )ÚselfÚevents     r   ÚfeedÚDecoder.feedg   s   € ¸sr   rN   N©rP   zMapping[str, Any]ÚreturnúIterable[Any])Ú__name__Ú
__module__Ú__qualname__Ú__firstlineno__rQ   Ú__static_attributes__rN   r   r   rJ   rJ   f   s   † ßBr   rJ   c                  ó0   • \ rS rSrSrSSS jjrS	S jrSrg)
ÚDataDecoderéj   u#  Yields `params.data` from events of a single `method`.

Covers the channels whose projection is just "emit the payload": `values`,
`updates`, `checkpoints`, `tasks` â€” the SDK analog of local's
`Values`/`Updates`/`Checkpoints`/`TasksTransformer`, all of which push
`params["data"]` unchanged. The REST-state seeding for `values` stays at
the projection layer; it is a one-shot pre-stream fetch, not part of the
event state machine.

Args:
    method: The protocol `method` this decoder consumes.
    namespace: When not `None`, events whose namespace differs are ignored
        (scope filter, mirroring the local transformers' `namespace != scope`
        check). `None` consumes every namespace â€” the historical `values`
        projection behavior, where subscription scoping is handled upstream.
Nc                óF   • Xl         Ub  [        U5      U l        g S U l        g rM   )Ú_methodr%   Ú
_namespace)rO   Úmethodr!   s      r   Ú__init__ÚDataDecoder.__init__|   s   € ØŒØ-6Ñ-Bœ$˜y›/ˆ�Èˆ�r   c              #  ó   #   • UR                  S5      U R                  :w  a  g UR                  S5      =(       d    0 nU R                  b  [        U5      U R                  :w  a  g UR                  S5      nUb  Uv •  g g 7f)Nra   Úparamsr-   )r$   r_   r`   r'   ©rO   rP   re   r-   s       r   rQ   ÚDataDecoder.feed€   sl   é € Ø�9‰9�XÓ $§,¡,Ó.ØØ—‘˜8Ó$×*¨ˆØ�?‰?Ñ&Ô+;¸FÓ+CÀtÇÁÓ+VØØ�z‰z˜&Ó!ˆØÑØ‹Jð ùs   ‚A<A>)r_   r`   rM   )ra   r,   r!   zlist[str] | NonerS   ©rV   rW   rX   rY   Ú__doc__rb   rQ   rZ   rN   r   r   r\   r\   j   s   † ñö"M÷r   r\   c                  ó4   • \ rS rSrSr    SS jrSS jrSrg)	ÚMessagesDecoderé‹   a3  Yields one chat-model stream per `message-start` event.

Subsequent events route to the matching stream via `stream.dispatch(data)`.
Mirrors the per-event body of `_MessagesProjection._messages_iter`
(`_async/stream.py:404-458`). The subscription open/close and the
`_root_messages_inbox` drain branch stay at the projection layer.

Args:
    namespace: Events whose namespace differs are ignored (scope filter).
    stream_factory: Keyword-only `(namespace, node, message_id) -> stream`.
        Sync binds `ChatModelStream`; async binds `AsyncChatModelStream`.
c                ó>   • [        U5      U l        X l        0 U l        g rM   )r%   r`   Ú_stream_factoryÚ_active)rO   r!   Ústream_factorys      r   rb   ÚMessagesDecoder.__init__™   s   € ô
 ˜y›/ˆŒØ-ÔØ')ˆ�r   c              #  ó*  #   • UR                  S5      S:w  a  g UR                  S5      =(       d    0 n[        U5      U R                  :w  a  g UR                  S5      n[        U[        5      (       d  g UR                  S5      S:X  a«  [        U5      n[        X4S9n[        UR                  S5      [        5      (       a  UR                  S5      O0 nU R                  [        U R                  5      U(       a  UR                  S	5      OS US
9nXpR                  U'   UR                  U5        Uv •  g [        U5      nU R                  R                  U5      nUcK  US:X  aE  [        U R                  5      S:X  a,  [        [        U R                  R                  5       5      5      nUc  g UR                  U5        UR                  S5      S;   a@  [        U R                  R                  5       5       H  u  p‰X—L d  M  U R                  U	 M     g g 7f)Nra   r   re   r-   rP   zmessage-start)r1   ÚmetadataÚlanggraph_node)r!   Únoder*   r0   rD   )zmessage-finishrA   )r$   r'   r`   r"   r#   r.   r2   rn   r%   ro   ÚdispatchrE   ÚnextÚiterr
   Úitems)
rO   rP   re   r-   r*   Úkeyrs   ÚstreamÚ	route_keyÚ	candidates
             r   rQ   ÚMessagesDecoder.feed¢   s¶  é € Ø�9‰9�XÓ *Ó,ØØ—‘˜8Ó$×*¨ˆÜ˜FÓ# t§¡Ó6ØØ�z‰z˜&Ó!ˆÜ˜$¤×%Ñ%ØØ�8‰8�GÓ Ó/Ü*¨4Ó0ˆJÜ$ TÑ?ˆCä(2°4·8±8¸JÓ3GÌ×(NÑ(N�—‘˜Ô$ÐTVð ð ×)Ñ)Ü˜tŸ™Ó/Þ7?�X—\‘\Ð"2Ô3ÀTØ%ð *ð ˆFð
 !'�L‰L˜ÑØ�O‰O˜DÔ!Ø‹Lä$ TÓ*ˆCØ—\‘\×%Ñ% cÓ*ˆFØ‰~ #¨Ó"5¼#¸d¿l¹lÓ:KÈqÓ:PÜœd 4§<¡<×#6Ñ#6Ó#8Ó9Ó:�Ø‰~ØØ�O‰O˜DÔ!Ø�x‰x˜Ó Ð$?Ó?Ü,0°·±×1CÑ1CÓ1EÖ,FÑ(�IØ Ô*Ø ŸL™L¨Ò3ò -Gð @ùs   ‚G:HÈ H)ro   r`   rn   N)r!   ú	list[str]rp   úCallable[..., Any]rS   rh   rN   r   r   rk   rk   ‹   s#   † ñð*àð*ð +ô*÷"4r   rk   c                  ó,   • \ rS rSrSrSS jrSS jrSrg)	ÚToolCallsDecoderéÇ   a£  Yields one tool-call handle per `tool-started` event.

Mirrors the per-event body of `_ToolCallsProjection._tool_calls_iter`
(`_async/stream.py:1168-1217`). The thread register/unregister and the
terminal-error-on-close finally stay at the projection / wrapper layer.

Args:
    namespace: Events whose namespace differs are ignored.
    handle_factory: Keyword-only `(tool_call_id, name, input, namespace) -> handle`.
c                ó>   • [        U5      U l        X l        0 U l        g rM   )r%   r`   Ú_handle_factoryro   )rO   r!   Úhandle_factorys      r   rb   ÚToolCallsDecoder.__init__Ó   s   € Ü˜y›/ˆŒØ-ÔØ')ˆ�r   c              #  óZ  #   • UR                  S5      S:w  a  g UR                  S5      =(       d    0 n[        U5      U R                  :w  a  g UR                  S5      n[        U[        5      (       d  g UR                  S5      n[        U[
        5      (       d  g UR                  S5      nUS:X  ao  UR                  S5      nU R                  U[        U[
        5      (       a  UOS	UR                  S
5      [        U R                  5      S9nXpR                  U'   Uv •  g US:X  aX  U R                  R                  U5      nUR                  S5      nUb(  [        U[
        5      (       a  UR                  U5        g g g US:X  aA  U R                  R                  US 5      nUb!  UR                  UR                  S5      5        g g US:X  a^  U R                  R                  US 5      nUb>  UR                  S5      n	UR                  [        U	(       a  [        U	5      OS5      5        g g g 7f)Nra   r   re   r-   Útool_call_idrP   ztool-startedÚ	tool_namer   r   )r‰   r;   r   r!   ztool-output-deltaÚdeltaztool-finishedÚoutputz
tool-errorÚmessagezTool call errored)r$   r'   r`   r"   r#   r,   r…   r%   ro   Ú_push_deltaÚpopÚ_finishÚ_failÚRuntimeError)
rO   rP   re   r-   r‰   Ú
event_typer;   Úhandler‹   r�   s
             r   rQ   ÚToolCallsDecoder.feedØ   sÞ  é € Ø�9‰9�XÓ 'Ó)ØØ—‘˜8Ó$×*¨ˆÜ˜FÓ# t§¡Ó6ØØ�z‰z˜&Ó!ˆÜ˜$¤×%Ñ%ØØ—x‘x Ó/ˆÜ˜,¬×,Ñ,ØØ—X‘X˜gÓ&ˆ
Ø˜Ó'Ø—8‘8˜KÓ(ˆDØ×)Ñ)Ø)Ü'¨¬c×2Ñ2‘T¸Ø—h‘h˜wÓ'Ü˜tŸ™Ó/ð	 *ð ˆFð *0�L‰L˜Ñ&Ø‹LØÐ.Ó.Ø—\‘\×%Ñ% lÓ3ˆFØ—H‘H˜WÓ%ˆEØÑ!¤j°¼×&<Ñ&<Ø×"Ñ" 5Õ)ð '=Ð!à˜?Ó*Ø—\‘\×%Ñ% l°DÓ9ˆFØÑ!Ø—‘˜tŸx™x¨Ó1Õ2ð "à˜<Ó'Ø—\‘\×%Ñ% l°DÓ9ˆFØÑ!ØŸ(™( 9Ó-�Ø—‘Ü ¶¤ W¤Ð>QÓRõð "ð (ùs   ‚H)H+)ro   r…   r`   N)r!   r   r†   r€   rS   rh   rN   r   r   r‚   r‚   Ç   s   † ñ	ô*÷
&r   r‚   c                  ó@   • \ rS rSrSrS	S jrS
S jrSS jrSS jrSr	g)ÚSubgraphsDecoderi  aÂ  Discovers child subgraph handles and fans out events to active ones.

Mirrors the per-event body of `_SubgraphsProjection._subgraphs_iter`
(`_async/stream.py:963-1041`) plus `_apply_tasks_result`. Root-inbox
forwarding and terminal-status-on-close stay at the projection / wrapper
layer.

Args:
    scope: Tuple-form namespace of this decoder's parent. `()` for root.
    handle_factory: Keyword-only `(path, graph_name, trigger_call_id) -> handle`.
c                óH   • Xl         X l        0 U l        [        5       U l        g rM   )Ú_scoper…   ro   ÚsetÚ_seen)rO   rG   r†   s      r   rb   ÚSubgraphsDecoder.__init__  s   € ØŒØ-ÔØ35ˆŒÜ+.«5ˆ�
r   c              #  óâ  #   • UR                  S5      =(       d    0 n[        U5      nUR                  S5      n[        U[        5      (       d  g UR                  S5      n[	        U5      nU R
                  R                  5        H=  u  px[        U5      n	[        U5      U	:¼  d  M!  US U	 U:X  d  M,  UR                  U5          O   US:X  aM  SU;   a  U R                  X45        g [        X0R                  5      (       a  U R                  U5       S h  v•N   g g US:X  aK  UR                  S5      S:X  a5  [        X0R                  5      (       a  U R                  U5       S h  v•N   g g g g  NX N
7f)	Nre   r-   ra   r   Úresultr   rP   r3   )r$   r'   r"   r#   rF   ro   ry   rE   Ú_push_eventÚ_apply_tasks_resultrH   r™   Ú	_discover)
rO   rP   re   r!   r-   ra   Úns_tupleÚ
child_pathÚchild_handleÚ	child_lens
             r   rQ   ÚSubgraphsDecoder.feed  sC  é € Ø—‘˜8Ó$×*¨ˆÜ$ VÓ,ˆ	Ø�z‰z˜&Ó!ˆÜ˜$¤×%Ñ%ØØ—‘˜8Ó$ˆô ˜Ó#ˆØ(,¯©×(:Ñ(:Ö(<Ñ$ˆJÜ˜J›ˆIÜ�8‹} 	Õ)¨h°z¸	Ð.BÀjÕ.PØ×(Ñ(¨Ô/Ùñ	 )=ð �WÓØ˜4ÓØ×(Ñ(¨Õ9Ü! )¯[©[×9Ñ9ØŸ>™>¨)Ó4×4Ñ4ð :ð �kÓ!Ø—‘˜Ó! YÓ.Ü  ¯K©K×8Ñ8à—~‘~ iÓ0×0Ñ0ð 9ð /ð "ñ 5ñ 1ùs2   ‚B E/Â&E/Â1A!E/ÄE+ÄAE/Å"E-Å#	E/Å-E/c              #  óð   #   • [        U5      nX R                  ;   a  g U R                  R                  U5        [        US   5      u  p4U R	                  UU=(       d    S US9nXPR
                  U'   Uv •  g 7f)Néÿÿÿÿ)ÚpathÚ
graph_nameÚtrigger_call_id)rF   r›   Úaddr>   r…   ro   )rO   r!   r©   rª   r«   r”   s         r   r¡   ÚSubgraphsDecoder._discover1  ss   é € Ü�YÓˆØ—:‘:ÓØØ�
‰
�‰�tÔÜ&>¸tÀB¹xÓ&HÑ#ˆ
Ø×%Ñ%ØØ!×) TØ+ð &ð 
ˆð
 $�‰�TÑØ‹ùs   ‚A4A6c                ó4  • UR                  S5      nU(       d  g [        U5      n[        U R                  R	                  5       5       HM  u  pVUS S U:w  a  M  UR
                  U:w  a  M"  [        U5      u  pxUR                  Xx5        U R                  U	 MO     g )Nr)   r¨   )r$   rF   r%   ro   ry   r«   rB   r�   )	rO   r!   r-   Ú	result_idÚparent_pathr£   r”   ÚstatusrA   s	            r   r    Ú$SubgraphsDecoder._apply_tasks_result?  s„   € Ø—H‘H˜T“Nˆ	ÞØÜ˜IÓ&ˆÜ"& t§|¡|×'9Ñ'9Ó';Ö"<ÑˆJØ˜#˜2ˆ +Ó-ÙØ×%Ñ%¨Ó2ÙÜ7¸Ó=‰MˆFØ�N‰N˜6Ô)Ø—‘˜ZÒ(ò #=r   )ro   r…   r™   r›   N)rG   útuple[str, ...]r†   r€   rS   )r!   r   rT   rU   )r!   r   r-   údict[str, Any]rT   ÚNone)
rV   rW   rX   rY   ri   rb   rQ   r¡   r    rZ   rN   r   r   r—   r—     s   † ñ
ô1ô1ô:÷)r   r—   c                  ó,   • \ rS rSrSrSS jrSS jrSrg)	ÚExtensionsDecoderiN  a1  Yields `params.data` from one named custom channel.

Mirrors `_ExtensionProjection._iter` (`_async/stream.py:1278-1299`), with
an added name filter so it can share one subscription in interleave.

Args:
    name: The extension name. Only `custom` events whose `data["name"]`
        matches are consumed.
c                ó4   • U(       d  [        S5      eXl        g )Nz!extension name must be non-empty.)r   Ú_name)rO   r;   s     r   rb   ÚExtensionsDecoder.__init__Y  s   € ÞÜÐ@ÓAÐAØ�
r   c              #  ó  #   • UR                  S5      S:w  a  g UR                  S5      =(       d    0 nUR                  S5      n[        U[        5      (       d  g UR                  S5      U R                  :w  a  g Uv •  g 7f)Nra   Úcustomre   r-   r;   )r$   r"   r#   r¹   rf   s       r   rQ   ÚExtensionsDecoder.feed^  sg   é € Ø�9‰9�XÓ (Ó*ØØ—‘˜8Ó$×*¨ˆØ�z‰z˜&Ó!ˆÜ˜$¤×%Ñ%ØØ�8‰8�FÓ˜tŸz™zÓ)ØØ‹
ùs   ‚A=A?)r¹   N)r;   r,   rS   rh   rN   r   r   r·   r·   N  s   † ñô÷
	r   r·   )r   r   rT   rµ   )r&   r   rT   r   )r-   r´   rT   ú
str | NonerM   )r-   r´   r1   r¾   rT   r,   )r:   r,   rT   ztuple[str, str | None])r-   r´   rT   z!tuple[SubgraphStatus, str | None])r!   r   rG   r³   rT   Úbool)ri   Ú
__future__r   Úcollections.abcr   r   r   Útypingr   r   r	   r   Ú	frozensetr   r   r'   r.   r2   ÚSubgraphStatusr>   rB   rH   rJ   r\   rk   r‚   r—   r·   rN   r   r   Ú<module>rÅ      s¹   ðñõ #ç 7Ñ 7ß )Ñ )ð!Ð ñ   )Ò)HÓIÐ ôô$Bô?ö
ð ÐHÑI€ô*ð
Ø
ðà&ôôXôCˆhô C÷ñ ÷B94ñ 94÷x7ñ 7÷tJ)ñ J)÷Zò r   