ó
    ýÞ jÀM ã                  ó0  • % S r SSKJr  SSKrSSKrSSKrSSKJrJrJ	r	J
r
  SSKJrJ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  SS
KJrJr  SSKJr  SSKJ r J!r!J"r"J#r#J$r$J%r%J&r&  SSK'J(r(J)r)  SSK*J+r+J,r,J-r-J.r.   " S S\5      r/\ " S S5      5       r0\ " S S5      5       r1/ SQr2S\3S'         SES jr4SFS jr5\6" SS15      r7SGS jr8 " S S5      r9 " S S 5      r:S!S".SHS# jjr; " S$ S%5      r< " S& S'5      r= " S( S)5      r>SIS* jr?SJSKS+ jjr@\S,   rASLS- jrB    SMS. jrCSNS/ jrDSOS0 jrE " S1 S25      rF " S3 S45      rG " S5 S65      rH " S7 S85      rI " S9 S:5      rJ " S; S<5      rK " S= S>5      rL " S? S@5      rM " SA SB5      rN " SC SD5      rOg)Pa%  Async thread-centric streaming surface for the v3 protocol.

`AsyncThreadStream` is an async context manager that owns a
`ProtocolSseTransport` for one thread, dispatches commands (`run.start`,
`run.respond`), exposes typed subscriptions over a single shared SSE
(`subscribe`, `events`), surfaces lifecycle state (`interrupted`,
`interrupts`) via an always-on lifecycle watcher SSE, and provides typed
projections (`thread.values`, `thread.messages`, `thread.tool_calls`,
`thread.extensions`).

Direct port of `libs/sdk/src/client/stream/index.ts`.
é    )ÚannotationsN)ÚAsyncGeneratorÚAsyncIteratorÚ	GeneratorÚMapping)Ú	dataclassÚfield)ÚAnyÚLiteralÚ	TypedDictÚcast©ÚAsyncChatModelStream)ÚEventÚSubscribeParams)Ú
HttpClient)ÚLangSmithTracingÚQueryParamTypes)Ú_SeenEventIds)ÚDataDecoderÚDecoderÚExtensionsDecoderÚMessagesDecoderÚSubgraphsDecoderÚToolCallsDecoderÚvalidate_interleave_channels)Úcompute_union_filterÚinfer_channel)ÚAsyncProtocolTransportÚEventStreamHandleÚProtocolSseTransportÚProtocolWebSocketTransportc                  ó8   • \ rS rSr% SrS\S'   S\S'   S\S'   S	rg
)ÚInterruptPayloadé/   zCPayload surfaced when the server requests human input for a thread.ÚstrÚinterrupt_idr
   Úvalueú	list[str]Ú	namespace© N)Ú__name__Ú
__module__Ú__qualname__Ú__firstlineno__Ú__doc__Ú__annotations__Ú__static_attributes__r+   ó    ÚU/var/www/html/gaurav/venv/lib/python3.13/site-packages/langgraph_sdk/_async/stream.pyr$   r$   /   s   ‡ ÙMàÓØƒJØÖr3   r$   c                  ó2   • \ rS rSr% SrS\S'   SrS\S'   Srg)	Ú_RunTerminalé7   zHTerminal state record resolved into `_run_done` on lifecycle completion.zLiteral['completed', 'errored']ÚstatusNzBaseException | NoneÚerrorr+   )r,   r-   r.   r/   r0   r1   r9   r2   r+   r3   r4   r6   r6   7   s   ‡ áRà+Ó+Ø"&€EÐÖ&r3   r6   c                  óX   • \ rS rSr% SrS\S'   S\S'   \" \R                  S9r	S\S	'   S
r
g)Ú_Subscriptioné?   zFInternal record for one active subscription on an `AsyncThreadStream`.ÚintÚidr   Úparams)Údefault_factoryzasyncio.QueueÚqueuer+   N)r,   r-   r.   r/   r0   r1   r	   ÚasyncioÚQueuerA   r2   r+   r3   r4   r;   r;   ?   s#   ‡ áPàƒGØÓÙ °·±Ñ?€Eˆ=Ö?r3   r;   )	ÚvaluesÚupdatesÚmessagesÚtoolsÚ	lifecycleÚinputÚcheckpointsÚtasksÚcustomr)   Ú_ALL_CHANNELSc                ó"   • U [        U5      /SS.$ )Nr   )ÚchannelsÚ
namespacesÚdepth©Úlist)rO   r*   s     r4   Ú_exact_namespace_paramsrT   Y   s   € ð
 Ü˜I“Ð'Øñð r3   c                ó¨   • [        U [        5      (       d  / $ U R                  S5      =(       d    / n[        U[        5      (       a  [        U5      $ / $ )Nr*   )Ú
isinstanceÚdictÚgetrS   )Úparams_fieldr*   s     r4   Ú_event_namespacerZ   d   sD   € Ü�l¤D×)Ñ)Øˆ	Ø× Ñ  Ó-×3°€IÜ(¨´D×9Ñ9Œ4�	‹?ÐA¸rÐAr3   Ú	completedÚfailedc                ó†  • [        U [        5      (       d  gU R                  S5      S:w  a  gU R                  S5      =(       d    0 n[        U[        5      (       d  gUR                  S5      (       d  / (       a  gUR                  S5      =(       d    0 n[        U[        5      (       d  gUR                  S5      [        ;   $ )a  Return True for a root-namespace lifecycle event marking run end.

Matches the wire shape ``{method: "lifecycle", params: {namespace: [],
data: {event: "completed" | "failed"}}}``. Subgraph lifecycle events
(non-empty namespace) do not terminate the parent run.
FÚmethodrH   r?   r*   ÚdataÚevent)rV   rW   rX   Ú_ROOT_TERMINAL_LIFECYCLE_EVENTS)r`   r?   r_   s      r4   Ú_is_root_terminal_lifecyclerb   n   s—   € ô �eœT×"Ñ"ØØ‡y�y�Ó˜kÓ)ØØ�Y‰Y�xÓ ×& B€FÜ�fœd×#Ñ#ØØ‡z�z�+×Ñ¦"ØØ�:‰:�fÓ×# €DÜ�dœD×!Ñ!ØØ�8‰8�GÓÔ ?Ñ?Ð?r3   c                  óF   • \ rS rSrSrS	S jrSSSS.       S
S jjrSrg)Ú_AgentModuleé„   z4Assistant graph helpers scoped to one thread stream.c                ó   • Xl         g ©N©Ú_owner©ÚselfÚowners     r4   Ú__init__Ú_AgentModule.__init__‡   ó   € Ø�r3   FN)ÚxrayÚheadersr?   c             ƒ  óª  #   • U R                   R                  (       a  [        S5      e0 nU(       a  XS'   U(       a  UR                  [	        U5      5        0 U R                   R
                  E[	        U=(       d    0 5      EnU R                   R                  R                  SU R                   R                   S3UU=(       d    S S9I S h  v•N $  N7f)NzAsyncThreadStream is closed.rp   z/assistants/z/graph)r?   rq   )	ri   Ú_closedÚRuntimeErrorÚupdaterW   Ú_headersÚ_httprX   Úassistant_id)rk   rp   rq   r?   Úquery_paramsÚrequest_headerss         r4   Úget_treeÚ_AgentModule.get_treeŠ   sµ   é € ð �;‰;××ÜÐ=Ó>Ð>Ø')ˆÞØ#'˜Ñ ÞØ×Ñ¤ V£Ô-ØI˜TŸ[™[×1Ñ1ÐI´T¸'¿-ÀRÓ5HÐIˆØ—[‘[×&Ñ&×*Ñ*Ø˜4Ÿ;™;×3Ñ3Ð4°FÐ;ØØ#×+ tð +ð 
÷ 
ð 	
ñ 
ùs   ‚C
CÃCÃCrh   ©rl   ÚAsyncThreadStreamÚreturnÚNone)rp   z
int | boolrq   úMapping[str, str] | Noner?   zQueryParamTypes | Noner   zdict[str, list[dict[str, Any]]])r,   r-   r.   r/   r0   rm   r{   r2   r+   r3   r4   rd   rd   „   sG   † Ù>ôð !Ø,0Ø)-ñ
ð ð
ð *ð	
ð
 'ð
ð 
)÷
ð 
r3   rd   c                  óh   • \ rS rSrSrS
S jrSSSSS.         SS jjrSS.     SS jjrS	rg)Ú	RunModuleé    zpCommand dispatcher for `run.start`.

Bound to one `AsyncThreadStream`; accesses its transport and id allocator.
c                ó   • Xl         g rg   rh   rj   s     r4   rm   ÚRunModule.__init__¦   ro   r3   N)rI   ÚconfigÚmetadataÚlangsmith_tracingc             ƒ  ó°  #   • SU R                   R                  0nUb  XS'   Ub  X%S'   Ub  X5S'   Ub  XES'   [        R                  " 5       nUR	                  5       nXpR                   l         U R                   R                  SU5      I Sh  v•N nUR                  5       (       d  UR                  S5        SU R                   l	        UU R                   R
                  UL a  SU R                   l        UR                  5       (       a'  UR                  5       (       d  UR                  5         $ $ $  N£! [         a,  n	UR                  5       (       d  UR                  U	5        e Sn	A	ff = f! U R                   R
                  UL a  SU R                   l        UR                  5       (       a'  UR                  5       (       d  UR                  5         f f f = f7f)	zGSend `run.start` to the server. Returns the result (`{"run_id": ...}`).rx   NrI   r‡   rˆ   Úlangsmith_tracerz	run.startT)ri   rx   rB   Úget_running_loopÚcreate_futureÚ_run_start_readyÚ_send_commandÚdoneÚ
set_resultÚ	_run_seenÚ	cancelledÚ	exceptionÚBaseExceptionÚset_exception)
rk   rI   r‡   rˆ   r‰   r?   ÚloopÚgateÚresultÚerrs
             r4   ÚstartÚRunModule.start©   s‡  é € ð #1°$·+±+×2JÑ2JÐ!KˆØÑØ#�7‰OØÑØ%�8ÑØÑØ!)�:ÑØÑ(Ø):Ð%Ñ&Ü×'Ò'Ó)ˆØ%)×%7Ñ%7Ó%9ˆØ'+�‰Ô$ð	!ØŸ;™;×4Ñ4°[À&ÓI×IˆFØ—9‘9—;‘;Ø—‘ Ô%Ø$(ˆD�K‰KÔ!Øð �{‰{×+Ñ+¨tÒ3Ø/3�—‘Ô,ð �y‰y�{‰{ 4§>¡>×#3Ñ#3Ø—‘Õ ð $4ˆ{ñ) Jøô
 ó 	ð —9‘9—;‘;Ø×"Ñ" 3Ô'Øûð	ûð �{‰{×+Ñ+¨tÒ3Ø/3�—‘Ô,ð �y‰y�{‰{ 4§>¡>×#3Ñ#3Ø—‘Õ ð $4ˆ{üsJ   ‚A*GÁ-D2 ÂD0Â<D2 Ã	A'GÄ0D2 Ä2
E(Ä<'E#Å#E(Å(E+ Å+A(GÇG)r'   c             ƒ  óh  ^#   • U R                   R                   ISh  v•N   [        U R                   R                  5      nTc^  [	        U5      S:X  a  [        S5      e[	        U5      S:”  a/  U Vs/ sH  oDS   PM	     nn[        S[	        U5       SU< S35      eUS   nO)[        U4S	 jU 5       S5      nUc  [        S
T< S35      eUS   US   US.nU R                   R                  SU5      I Sh  v•N sSSS5      ISh  v•N   $  Nîs  snf  N N! , ISh  v•N  (       d  f       g= f7f)aï  Reply to a server-side interrupt and resume the run.

Args:
    response: the response value forwarded as `params.response` on the
        wire (protocol field name).
    interrupt_id: optional explicit id. When omitted, requires exactly
        one outstanding interrupt and uses its id.

Raises:
    RuntimeError: no outstanding interrupts; `interrupt_id` is None but
        multiple interrupts are outstanding; or the explicit
        `interrupt_id` doesn't match any outstanding interrupt.
Nr   zrthread.run.respond: no outstanding interrupt. Provide an explicit `interrupt_id` or wait for `thread.interrupted`.é   r'   u"   thread.run.respond: ambiguous â€” z outstanding interrupts (z&). Provide an explicit `interrupt_id`.c              3  ó:   >#   • U H  oS    T:X  d  M  Uv •  M     g7f)r'   Nr+   )Ú.0Úpr'   s     €r4   Ú	<genexpr>Ú$RunModule.respond.<locals>.<genexpr>ÿ   s   øé € ÐQ¡˜1°Ñ/@ÀLÑ/P—Q‘Q¢ùs   ƒ’	z!thread.run.respond: interrupt_id zA does not match any outstanding interrupt in `thread.interrupts`.r*   )r'   r*   Úresponsezinput.respond)ri   Ú_interrupts_lockrS   Ú
interruptsÚlenrt   Únextr�   )rk   r¤   r'   Úoutstandingr¡   ÚidsÚmatchr?   s     `     r4   ÚrespondÚRunModule.respondÖ   sI  øé € ð, —;‘;×/×/Ó/Ü˜tŸ{™{×5Ñ5Ó6ˆKØÑ#Ü�{Ó# qÓ(Ü&ð0óð ô
 �{Ó# aÓ'Ù6AÓB±k°˜^Ô,±k�CÐBÜ&Ø<¼SÀÓ=MÐ<Nð O3Ø36±'ð :*ð*óð ð
 $ A™‘äÜQ¡ÓQØó�ð ‘=Ü&Ø;¸LÑ;Kð L/ð /óð ð !& nÑ 5Ø" ;Ñ/Ø$ñˆFð
 Ÿ™×2Ñ2°?ÀFÓK×K÷C 0×/Ó/ùò Cñ. L÷C 0×/×/Ð/üsd   ƒD2žDŸD2¢ADÁ2DÁ?A8DÃ7DÃ8DÃ;D2ÄDÄD2ÄDÄD2ÄD/ÄD!ÄD/Ä+D2rh   r}   )
rI   r
   r‡   údict[str, Any] | Nonerˆ   r®   r‰   zLangSmithTracing | Noner   údict[str, Any])r¤   r
   r'   ú
str | Noner   r¯   )	r,   r-   r.   r/   r0   rm   r›   r¬   r2   r+   r3   r4   rƒ   rƒ       s‚   † ñô
ð Ø(,Ø*.Ø59ñ+!ð ð+!ð &ð	+!ð
 (ð+!ð 3ð+!ð 
õ+!ðb $(ñ	7Làð7Lð !ð	7Lð
 
÷7Lð 7Lr3   rƒ   g        )Údelayc             ƒ  óŽ   #   • U(       a  [         R                  " U5      I Sh  v•N   U R                  5       I Sh  v•N   g N N7f)z¹Close a handle, optionally after a brief delay. Used to detach
closing the old stream from the synchronous rotation step so the new
stream can absorb server-side replayed events first.
N)rB   ÚsleepÚclose)Úhandler±   s     r4   Ú_close_afterr¶     s3   é € ö
 Ü�mŠm˜EÓ"×"Ð"Ø
�,‰,‹.×Ññ 	#Ùùs   ‚!A£A¤A»A¼AÁAc                  óF   • \ rS rSrSrS
S jrS rSS jrSS jrSS jr	Sr
g	)Ú_OutputAwaitablei  zãAwaitable that waits for lifecycle completion then fetches durable thread state.

Multiple awaiters share one underlying task (idempotent task caching).
Call `with_timeout(seconds)` to bound the wait on the lifecycle terminal.
c                ó,   • Xl         S U l        S U l        g rg   )Ú_threadÚ_taskÚ_timeout©rk   Úthreads     r4   rm   Ú_OutputAwaitable.__init__!  s   € ØŒØ/3ˆŒ
Ø&*ˆ�r3   c                ó>   • U R                  5       R                  5       $ rg   )Ú	_get_taskÚ	__await__©rk   s    r4   rÂ   Ú_OutputAwaitable.__await__&  s   € Ø�~‰~Ó×)Ñ)Ó+Ð+r3   c                ó<   • [        U R                  5      nXl        U$ )a"  Return a new awaitable that raises `asyncio.TimeoutError` after `timeout` seconds.

Bounds the wait for the lifecycle terminal (and only that wait); the
subsequent REST GET for terminal state is not bounded. Returns a
fresh `_OutputAwaitable` so the original `thread.output` is unaffected.
)r¸   rº   r¼   )rk   ÚtimeoutÚboundeds      r4   Úwith_timeoutÚ_OutputAwaitable.with_timeout)  s   € ô # 4§<¡<Ó0ˆØ"ÔØˆr3   c                ó†   • U R                   c)  [        R                  " U R                  5       5      U l         U R                   $ )a|  Return the shared fetch task, creating it on first call.

A cancelled task is intentionally NOT respawned: subsequent awaiters
receive `asyncio.CancelledError` from the shared task instead of
triggering a fresh REST GET. This preserves "one fetch per
`thread.output`" semantics even when callers wrap awaits with
`asyncio.wait_for` (which cancels the underlying task on timeout).
)r»   rB   Úcreate_taskÚ_fetchrÃ   s    r4   rÁ   Ú_OutputAwaitable._get_task4  s0   € ð �:‰:ÑÜ ×,Ò,¨T¯[©[«]Ó;ˆDŒJØ�z‰zÐr3   c              ƒ  óD  #   • U R                   R                  5       (       aG  U R                   R                  5       I Sh  v•N nU R                   R                  U5      (       a  US   $ U R                  b@  [
        R                  " U R                   R                  5       U R                  S9I Sh  v•N nO"U R                   R                  5       I Sh  v•N nUR                  b  UR                  eU R                   R                  5       I Sh  v•N nUS   $  NØ Ni NH N7f)zAFetch terminal thread state, waiting for the lifecycle if needed.NrD   ©rÆ   )	rº   Ú&_can_return_existing_state_immediatelyÚ_fetch_stateÚ_state_is_terminalr¼   rB   Úwait_forÚ_wait_for_run_doner9   )rk   ÚstateÚterminals      r4   rÌ   Ú_OutputAwaitable._fetchA  sã   é € ð �<‰<×>Ñ>×@Ñ@ØŸ,™,×3Ñ3Ó5×5ˆEØ�|‰|×.Ñ.¨u×5Ñ5Ø˜X‘Ð&ð �=‰=Ñ$Ü$×-Ò-Ø—‘×/Ñ/Ó1¸4¿=¹=ñ÷ ‰Hð "Ÿ\™\×<Ñ<Ó>×>ˆHØ�>‰>Ñ%Ø—.‘.Ð Ø—l‘l×/Ñ/Ó1×1ˆØ�X‰Ðñ 6ññ ?ñ 2ùsG   ‚=D ¿DÁ A0D Â0DÂ1"D ÃDÃ:D ÄDÄ
D ÄD ÄD ÄD )r»   rº   r¼   N©r¾   r~   r   r€   )rÆ   Úfloatr   r¸   )r   zasyncio.Task[Any])r   r
   )r,   r-   r.   r/   r0   rm   rÂ   rÈ   rÁ   rÌ   r2   r+   r3   r4   r¸   r¸     s    † ñô+ò
,ô	ô÷r3   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)Ú_ValuesProjectioniW  uÕ   Typed projection for `thread.values` â€” yields state snapshots as they arrive.

Supports both `async for` (live stream of state snapshots) and `await`
(delegates to `thread.output` for the terminal state value).
c                ó   • Xl         g rg   ©rº   r½   s     r4   rm   Ú_ValuesProjection.__init__^  ó   € Ø�r3   c                óJ   • U R                   R                  R                  5       $ rg   )rº   ÚoutputrÂ   rÃ   s    r4   rÂ   Ú_ValuesProjection.__await__a  s   € Ø�|‰|×"Ñ"×,Ñ,Ó.Ð.r3   c                ó"   • U R                  5       $ rg   )Ú_values_iterrÃ   s    r4   Ú	__aiter__Ú_ValuesProjection.__aiter__d  s   € Ø× Ñ Ó"Ð"r3   c               ó¬  #   • U R                   R                  c  [        S5      eSS/0nU R                   R                  U5      n[	        S5      n U R                   R                  U5      I S h  v•N   U R                   R                  5         U R                   R                  5       I S h  v•N nUS   7v •   UR                  R                  5       I S h  v•N nUc'   U R                   R                  UR                  5        g UR                  U5       H  nU7v •  M
     Mk   N´ Nz NQ! U R                   R                  UR                  5        f = f7f)Nõ3   AsyncThreadStream not entered â€” use `async with`.rO   rD   )rº   Ú
_transportrt   Ú_register_subscriptionr   Ú_reconcile_streamÚ_ensure_fanout_runningrÑ   rA   rX   Ú_unregister_subscriptionr>   Úfeed)rk   r?   ÚsubÚdecoderrÕ   ÚitemÚouts          r4   rä   Ú_ValuesProjection._values_iterg  s  é € Ø�<‰<×"Ñ"Ñ*ÜÐTÓUÐUØ#-°¨zÐ":ˆØ�l‰l×1Ñ1°&Ó9ˆÜ˜hÓ'ˆð	:Ø—,‘,×0Ñ0°Ó8×8Ð8Ø�L‰L×/Ñ/Ô1ØŸ,™,×3Ñ3Ó5×5ˆEØ˜‘/Ó!ØØ ŸY™YŸ]™]›_×,�Ø‘<Øð �L‰L×1Ñ1°#·&±&Õ9ð #Ÿ<™<¨Ö-�CØ•Iñ .ñ	 ñ	 9á5ñ -øð �L‰L×1Ñ1°#·&±&Õ9üsT   ‚AEÁD* Á/D$Á0;D* Â+D&Â,*D* ÃD(ÃD* Ã&EÄ D* Ä&D* Ä(D* Ä*'EÅErÝ   NrØ   )r   zGenerator[Any, None, Any]©r   zAsyncIterator[Any]©r   zAsyncGenerator[Any, None])
r,   r-   r.   r/   r0   rm   rÂ   rå   rä   r2   r+   r3   r4   rÛ   rÛ   W  s   † ñôô/ô#÷:r3   rÛ   c                  óX   • \ rS rSrSr S	     S
S jjrSS jrSS jr    SS jrSr	g)Ú_MessagesProjectioni|  zÜTyped projection for root-scope `thread.messages`.

Iterating yields one `AsyncChatModelStream` per message-start event.
Each iterator owns its own `messages` subscription and routes events
from the root namespace only.
Nc                óB   • Xl         [        U=(       d    / 5      U l        g rg   ©rº   rS   Ú
_namespace©rk   r¾   r*   s      r4   rm   Ú_MessagesProjection.__init__„  ó   € ð ŒÜ˜yŸ¨BÓ/ˆ�r3   c                ó"   • U R                  5       $ rg   ©Ú_messages_iterrÃ   s    r4   rå   Ú_MessagesProjection.__aiter__Š  ó   € Ø×"Ñ"Ó$Ð$r3   c               ó  #   • U R                   R                  c  [        S5      eU R                  (       d  U R                   R                  OS nUb!  U R                  U5        S h  v•N nU7v •  M  [        S/U R                  5      nU R                   R                  U5      n[        U R                  S S9n/ n U R                   R                  U5      I S h  v•N   U R                   R                  5          UR                  R                  5       I S h  v•N nUcK   U H  nU R                   R                  U5        M      U R                   R                  UR                  5        g UR!                  U5       H4  nU R                   R#                  U5        UR%                  U5        U7v •  M6     M»   GNO
 g  Nß N¤! U H  nU R                   R                  U5        M      U R                   R                  UR                  5        f = f7f)Nrè   rF   c                ó   • [        XUS9$ ©N©r*   ÚnodeÚ
message_idr   r  s      r4   Ú<lambda>Ú4_MessagesProjection._messages_iter.<locals>.<lambda>œ  s   € ÔBVØ#¸:òCr3   ©r*   Ústream_factory)rº   ré   rt   rú   Ú_root_messages_inboxÚ_drain_inboxrT   rê   r   rë   rì   rA   rX   Ú!_unregister_active_message_streamrí   r>   rî   Ú_register_active_message_streamÚappend)rk   Ú
root_inboxÚstreamr?   rï   rð   Ú
registeredrñ   s           r4   r   Ú"_MessagesProjection._messages_iter�  s¦  é € Ø�<‰<×"Ñ"Ñ*ÜÐTÓUÐUð ?C¿o¿o�T—\‘\×6Ò6ÐSWˆ
ØÑ!Ø $× 1Ñ 1°*Ô =÷ �fØ•ä(¨*¨°t·±ÓGˆØ�l‰l×1Ñ1°&Ó9ˆÜ!Ø—o‘oññ
ˆð 24ˆ
ð	:Ø—,‘,×0Ñ0°Ó8×8Ð8Ø�L‰L×/Ñ/Ô1ØØ ŸY™YŸ]™]›_×,�Ø‘<Øó %�Ø—‘×>Ñ>¸vÖFñ %à�L‰L×1Ñ1°#·&±&Õ9ð &Ÿl™l¨4Ö0�FØ—L‘L×@Ñ@ÀÔHØ×%Ñ% fÔ-Ø •Lñ 1ñ	 òÐ =àñ 9ñ -øó %�Ø—‘×>Ñ>¸vÖFñ %à�L‰L×1Ñ1°#·&±&Õ9üsp   ‚A HÁ"F8Á&F5Á'F8Á*AHÂ<F> ÃF:Ã<F> ÄF<ÄF> Ä A
HÅ*AF> Æ5F8Æ8HÆ:F> Æ<F> Æ>AH	È	Hc               óâ  #   • SSK Jn  0 n  UR                  5       I Sh  v•N nUc4   UR                  5        H  nU R                  R                  U5        M      gUR                  S5      =(       d    0 n[        U[        5      (       a  UR                  S5      OSn[        U[        5      (       d  M©  UR                  S5      nUS:X  a´  [        U5      n	[        XyS9n
[        UR                  S	5      [        5      (       a  UR                  S	5      O0 nU" [        U R                  5      U(       a  UR                  S
5      OSU	S9nXSU
'   U R                  R                  U5        UR                  U5        U7v •  O²[        U5      n
UR                  U
5      nUc1  [        U5      S:X  a"  [        [!        UR                  5       5      5      nUc  GMÊ  UR                  U5        US;   aE  U R                  R                  U5        [        UR#                  5       5       H  u  pÍXÕL d  M  X<	 M     GM(   GN! UR                  5        H  nU R                  R                  U5        M      f = f7f)zMDrain a pre-filled inbox of messages events, yielding one stream per message.r   r   Nr?   r_   r`   úmessage-start©Úfallbackrˆ   Úlanggraph_noder  rž   ©zmessage-finishr9   )Ú0langchain_core.language_models.chat_model_streamr   rX   rD   rº   r  rV   rW   Ú_message_event_idÚ_message_route_keyrS   rú   r  Údispatchr§   r¨   ÚiterÚitems)rk   Úinboxr   Úactiverñ   r  rY   r_   Ú
event_typer  Úkeyrˆ   Ú	route_keyÚ	candidates                 r4   r  Ú _MessagesProjection._drain_inbox±  s  é € õ	
ð 35ˆð,	GØØ"ŸY™Y›[×(�Ø‘<ØðN !Ÿ-™-ž/�Ø—‘×>Ñ>¸vÖFò *ðM  $Ÿx™x¨Ó1×7°R�ä0:¸<Ì×0NÑ0N�L×$Ñ$ VÔ,ÐTXð ô " $¬×-Ñ-ÙØ!ŸX™X gÓ.�
Ø Ó0Ü!2°4Ó!8�JÜ,¨TÑG�Cô & d§h¡h¨zÓ&:¼D×AÑAð Ÿ™ Ô,àð ñ
 2Ü"& t§¡Ó"7Þ?G˜XŸ\™\Ð*:Ô;ÈTØ#-ñ�Fð
 #)˜3‘KØ—L‘L×@Ñ@ÀÔHØ—O‘O DÔ)Ø ”Lä,¨TÓ2�CØ#ŸZ™Z¨›_�FØ‘~¬#¨f«+¸Ó*:Ü!%¤d¨6¯=©=«?Ó&;Ó!<˜Ø‘~Ú Ø—O‘O DÔ)Ø!Ð%@Ó@ØŸ™×FÑFÀvÔNÜ48¸¿¹»Ö4HÑ0˜IØ(Ô2Ø$*Ò$5ñ 5IòM Ú(øðR !Ÿ-™-ž/�Ø—‘×>Ñ>¸vÖFò *üs2   ‚	I/ŒH8  H5¡H8 ©3I/ÁGH8 È,
H8 È84I,É,I/©rú   rº   rg   ©r¾   r~   r*   úlist[str] | Noner   r€   )r   z#AsyncIterator[AsyncChatModelStream])r   ú*AsyncGenerator[AsyncChatModelStream, None])r"  úasyncio.Queue[Event | None]r   r,  )
r,   r-   r.   r/   r0   rm   rå   r   r  r2   r+   r3   r4   r÷   r÷   |  sN   † ñð HLð0Ø'ð0Ø4Dð0à	õ0ô%ô":ðH5GØ0ð5Gà	3÷5Gr3   r÷   c                ót   • U R                  S5      =(       d    U R                  S5      nUb  [        U5      $ S $ )Nr>   r  ©rX   r&   )r_   r  s     r4   r  r  é  s1   € Ø—‘˜$“×9 4§8¡8¨LÓ#9€JØ(Ñ4Œ3ˆz‹?Ð>¸$Ð>r3   c                ó:   • [        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_   r  r  s      r4   r  r  î  s7   € ô # 4Ó(€JØÑØ˜*˜Ð&Ð&ØÑØ˜(˜Ð$Ð$Ør3   )Ústartedr[   r\   Úinterruptedc                óD   • U R                  S5      u  pnX(       a  U4$ S 4$ )NÚ:)Ú	partition)ÚsegmentÚnameÚsepÚtask_ids       r4   Ú_parse_namespace_segmentr;     s,   € Ø ×*Ñ*¨3Ó/Ñ€DˆwØ�C�Ð)Ð) TÐ)Ð)r3   c                ó|   • U R                  S5      (       a  gU R                  S5      nU(       a  S[        U5      4$ g)Nr¦   )r3  Nr9   r\   )r[   Nr/  )r_   r9   s     r4   Ú_terminal_from_tasks_resultr=    s9   € ð ‡x�x�×ÑØ"Ø�H‰H�WÓ€EÞØœ˜U›Ð#Ð#Ør3   c                óx   • [        U 5      [        U5      S-   :H  =(       a    [        U S [        U5       5      U:H  $ )Nrž   )r§   Útuple)r*   Úscopes     r4   Ú_is_direct_childrA    s4   € Üˆy‹>œS ›Z¨!™^Ñ+×W´°iÀÄ#ÀeÃ*Ð6MÓ0NÐRWÑ0WÐWr3   c                ó$   • / SQ[        U 5      /S.$ )N)rF   rK   rG   rH   )rO   rP   rR   ©r@  s    r4   Ú_subgraph_subscription_paramsrD    s   € ò @Ü˜E“{�mñð r3   c                  ó€   • \ rS rSrSrSS.           SS jjrSS jrSS jrSS jrSS	 jr	SS
 jr
SSS jjrSrg)ÚScopedStreamHandlei#  z<Scoped streaming handle for one discovered child invocation.r   )Úmax_queue_sizec               óð  • Xl         X l        [        U5      U l        X0l        X@l        SU l        S U l        XPl        [        R                  " US9U l        [        R                  " US9U l        [        R                  " US9U l        0 U l        [        5       U l        [#        U 5      U l        ['        U 5      U l        [+        U 5      U l        U R,                  U l        [1        U[        U5      S9U l        g )Nr2  ©Úmaxsize©r*   )rº   ÚpathrS   r*   Ú
graph_nameÚtrigger_call_idr8   r9   Ú_max_queue_sizerB   rC   Ú_messages_inboxÚ_tools_inboxÚ_tasks_inboxÚ_descendant_handlesÚsetÚ_iterated_inboxesÚ_HandleMessagesProjectionrF   Ú_HandleToolCallsProjectionÚ
tool_callsÚ_HandleSubgraphsProjectionÚ	subgraphsÚ	subagentsÚ_ExtensionsProjectionÚ
extensions)rk   r¾   rL  rM  rN  rG  s         r4   rm   ÚScopedStreamHandle.__init__&  sÖ   € ð ŒØŒ	Ü˜d›ˆŒØ$ŒØ.ÔØ&/ˆŒØ!%ˆŒ
Ø-Ôô =D¿MºMØ"ñ=
ˆÔô :A¿ºØ"ñ:
ˆÔô :A¿ºØ"ñ:
ˆÔð OQˆÔ ô ,/«5ˆÔÜ1°$Ó7ˆŒÜ4°TÓ:ˆŒÜ3°DÓ9ˆŒØŸ™ˆŒÜ/°Ä$ÀtÃ*ÑMˆ�r3   c                ó  • UR                  S5      nUS:X  a  U R                  R                  U5        OCUS:X  a  U R                  R                  U5        O!US:X  a  U R                  R                  U5        US;   aˆ  [        [        UR                  S5      =(       d    0 5      5      nU R                  R                  5        H=  u  pE[        U5      n[        U5      U:¼  d  M!  USU U:X  d  M,  UR                  U5        M?     gg)zýRoute a descendant event into the appropriate channel inbox.

Also fans out to any registered descendant handles whose path is a
prefix of the event namespace, so grandchild events are delivered at
push time rather than via a post-hoc drain-and-replay.
r^   rF   rG   rK   ©rF   rG   rK   r?   N)rX   rP  Ú
put_nowaitrQ  rR  r?  rZ   rS  r!  r§   Ú_push_event)rk   r`   r^   Úns_tupleÚ	desc_pathÚdesc_handleÚdesc_lens          r4   rb  ÚScopedStreamHandle._push_eventR  sâ   € ð —‘˜8Ó$ˆØ�ZÓØ× Ñ ×+Ñ+¨EÕ2Ø�wÓØ×Ñ×(Ñ(¨Õ/Ø�wÓØ×Ñ×(Ñ(¨Ô/ð Ð3Ó3ÜÔ-¨e¯i©i¸Ó.A×.GÀRÓHÓIˆHØ*.×*BÑ*B×*HÑ*HÖ*JÑ&�	Ü˜y›>�Ü�x“= HÕ,°¸)¸8Ð1DÈ	Õ1QØ×+Ñ+¨EÖ2ò +Kð 4r3   c           	     ó6  • XR                   UR                  '   [        UR                  5      nS Hæ  n[        X5      n/ nUR	                  5       (       d6  UR                  UR                  5       5        UR	                  5       (       d  M6  U H…  nUR                  U5        Uc  M  [        [        UR                  S5      =(       d    0 5      5      n[        U5      U:¼  d  MV  USU UR                  :X  d  Mk  [        X5      R                  U5        M‡     Mè     g)a  Register a newly-discovered grandchild so future events are fanned out.

Also drains any events already buffered in this handle's inboxes whose
namespace matches the grandchild, so events that arrived before the
grandchild was discovered are forwarded in arrival order.
)rP  rQ  rR  Nr?   )rS  rL  r§   ÚgetattrÚemptyr  Ú
get_nowaitra  r?  rZ   rX   )rk   rµ   rf  Ú
inbox_attrr"  Ústagingr`   rc  s           r4   Ú_register_descendantÚ'ScopedStreamHandle._register_descendanti  sÝ   € ð 17× Ñ  §¡Ñ-Ü�v—{‘{Ó#ˆó
ˆJô
 29¸Ó1JˆEØ*,ˆGØ—k‘k—m‘mØ—‘˜u×/Ñ/Ó1Ô2ð —k‘k—m“mã �Ø× Ñ  Ô'Ø‘=ÙÜ Ô!1°%·)±)¸HÓ2E×2KÈÓ!LÓM�Ü�x“= HÕ,°¸)¸8Ð1DÈÏÉÕ1SÜ˜FÓ/×:Ñ:¸5ÖAó !ò
r3   c                ó<   • U R                   R                  US5        g)z6Remove a grandchild after it reaches a terminal state.N)rS  Úpop)rk   rL  s     r4   Ú_unregister_descendantÚ)ScopedStreamHandle._unregister_descendantƒ  s   € à× Ñ ×$Ñ$ T¨4Õ0r3   c                óš   • U R                   R                  U5        U R                  S:w  a   [        U SU S35      R	                  S5        gg)a:  Record that an inbox has an active consumer.

If the handle is already closed (status != 'started'), immediately
enqueue a sentinel so the consumer's `await get()` terminates. This
handles sequential consumption (iterate after the handle is finished).

Must be called by each projection at the start of iteration.
r2  Ú_Ú_inboxN)rU  Úaddr8   ri  ra  ©rk   Úkinds     r4   Ú_mark_iteratedÚ!ScopedStreamHandle._mark_iterated‡  sI   € ð 	×Ñ×"Ñ" 4Ô(Ø�;‰;˜)Ó#ô �D˜A˜d˜V 6Ð*Ó+×6Ñ6°tÕ<ð $r3   c                óv   • S H3  nXR                   ;   d  M  [        U SU S35      R                  S5        M5     g)a  Signal EOF only on channel inboxes that have an active consumer.

Inboxes without a consumer would accumulate a leaked None sentinel
that is never drained, so we skip them. For inboxes whose consumer
starts after this call, `_mark_iterated` sends the sentinel lazily.
r`  ru  rv  N)rU  ri  ra  rx  s     r4   Ú_close_inboxesÚ!ScopedStreamHandle._close_inboxes–  s8   € ó 3ˆDØ×-Ñ-Õ-Ü˜  $  vÐ.Ó/×:Ñ:¸4Ö@ò 3r3   Nc                ó^   • U R                   S:w  a  g Xl         X l        U R                  5         g )Nr2  )r8   r9   r}  )rk   r8   r9   s      r4   Ú_finishÚScopedStreamHandle._finish¡  s'   € Ø�;‰;˜)Ó#ØØŒØŒ
Ø×ÑÕr3   )rS  rU  rO  rP  rR  rº   rQ  r9   r]  rM  rF   r*   rL  r8   r[  rZ  rX  rN  )r¾   r~   rL  útuple[str, ...]rM  r°   rN  r°   rG  r=   r   r€   ©r`   r   r   r€   ©rµ   rF  r   r€   )rL  r‚  r   r€   )ry  r&   r   r€   ©r   r€   rg   )r8   ÚSubgraphStatusr9   r°   r   r€   )r,   r-   r.   r/   r0   rm   rb  rn  rr  rz  r}  r€  r2   r+   r3   r4   rF  rF  #  sx   † ÙFð  ñ*Nð "ð*Nð ð	*Nð
 ð*Nð $ð*Nð ð*Nð 
õ*NôX3ô.Bô41ô=ô	A÷ñ r3   rF  c                  ó6   • \ rS rSrSrSS jrS	S jrS
S jrSrg)rV  i©  zHMessages projection that drains a `ScopedStreamHandle`'s messages inbox.c                ó   • Xl         g rg   ©Ú_handle©rk   rµ   s     r4   rm   Ú"_HandleMessagesProjection.__init__¬  rß   r3   c                ó"   • U R                  5       $ rg   rÿ   rÃ   s    r4   rå   Ú#_HandleMessagesProjection.__aiter__¯  r  r3   c               ó^  #   • SSK Jn  U R                  R                  S5        0 n U R                  R                  R                  5       I S h  v•N nUc  g UR                  S5      =(       d    0 n[        U5      nXPR                  R                  :w  a  Mq  [        U[        5      (       a  UR                  S5      OS n[        U[        5      (       d  M°  UR                  S5      nUS:X  a£  [        U5      n[        XhS9n	[        UR                  S	5      [        5      (       a  UR                  S	5      O0 n
U" [        U R                  R                  5      U
(       a  U
R                  S
5      OS US9nX²U	'   UR                  U5        U7v •  O—[        U5      n	UR                  U	5      nUc1  [        U5      S:X  a"  [        [!        UR#                  5       5      5      nUc  GMÀ  UR                  U5        US;   a*  [        UR%                  5       5       H  u  pÍXÛL d  M  X,	 M     GM   GNÜ7f)Nr   r   rF   r?   r_   r`   r  r  rˆ   r  r  rž   r  )r  r   rŠ  rz  rP  rX   rZ   r*   rV   rW   r  r  rS   r  r§   r¨   r   rD   r!  )rk   r   r#  rñ   rY   Únsr_   r$  r  r%  rˆ   r  r&  r'  s                 r4   r   Ú(_HandleMessagesProjection._messages_iter²  sÓ  é € õ	
ð 	�‰×#Ñ# JÔ/Ø24ˆØØŸ™×5Ñ5×9Ñ9Ó;×;ˆDØ‰|ØØŸ8™8 HÓ-×3°ˆLÜ! ,Ó/ˆBØ—\‘\×+Ñ+Ó+ÙÜ/9¸,Ì×/MÑ/M�<×#Ñ# FÔ+ÐSWˆDÜ˜d¤D×)Ñ)ÙØŸ™ 'Ó*ˆJØ˜_Ó,Ü.¨tÓ4�
Ü(¨ÑC�ô " $§(¡(¨:Ó"6¼×=Ñ=ð —H‘H˜ZÔ(àð ñ
 .Ü" 4§<¡<×#9Ñ#9Ó:Þ;C˜Ÿ™Ð&6Ô7ÈØ)ñ�ð
 %�s‘Ø—‘ Ô%Ø”ä(¨Ó.�ØŸ™ C›�Ø‘>¤c¨&£k°QÓ&6Ü!¤$ v§}¡}£Ó"7Ó8�FØ‘>ÚØ—‘ Ô%ØÐ!<Ó<Ü04°V·\±\³^Ö0DÑ,˜	Ø$Ô.Ø &Ò 1ñ 1EòK Ú;ùs   ‚AH-ÁH*ÁGH-È!
H-r‰  Nr„  rô   rõ   )	r,   r-   r.   r/   r0   rm   rå   r   r2   r+   r3   r4   rV  rV  ©  s   † ÙRôô%÷.2r3   rV  c                  ó6   • \ rS rSrSrSS jrS	S jrS
S jrSrg)rW  iã  zGTool calls projection that drains a `ScopedStreamHandle`'s tools inbox.c                ó   • Xl         g rg   r‰  r‹  s     r4   rm   Ú#_HandleToolCallsProjection.__init__æ  rß   r3   c                ó"   • U R                  5       $ rg   ©Ú_tool_calls_iterrÃ   s    r4   rå   Ú$_HandleToolCallsProjection.__aiter__é  ó   € Ø×$Ñ$Ó&Ð&r3   c               ó6  #   • U R                   R                  S5        0 n U R                   R                  R                  5       I S h  v•N nUc4  [	        S5      nUR                  5        H  nUR                  U5        M     g UR                  S5      =(       d    0 n[        U5      nX`R                   R                  :w  a  M¤  [        U[        5      (       a  UR                  S5      OS n[        U[        5      (       d  Mã  UR                  S5      nUR                  S5      n	[        U	[        5      (       d  GM  US:X  aj  UR                  S5      n
[        U
[        5      (       d  S	n
[        U	U
UR                  S
5      [        U R                   R                  5      S9nXAU	'   U7v •  OæUS:X  aL  UR                  U	5      nUR                  S5      nUb&  [        U[        5      (       a  UR                  U5        O”US:X  a6  UR                  U	S 5      nUb   UR!                  UR                  S5      5        OXUS:X  aR  UR                  U	S 5      nUb=  UR                  S5      nUR                  [	        U(       a  [        U5      OS5      5        GMu   GNN7f)NrG   ú3Tool call stream closed before terminal tool event.r?   r_   r`   Útool_call_idztool-startedÚ	tool_nameÚ rI   ©rœ  r8  rI   r*   ztool-output-deltaÚdeltaztool-finishedrá   z
tool-errorÚmessagezTool call errored)rŠ  rz  rQ  rX   rt   rD   Ú_failrZ   r*   rV   rW   r&   ÚToolCallHandlerS   Ú_push_deltarq  r€  )rk   r#  rñ   rš   rµ   rY   r�  r_   r$  rœ  r�  r   r¡  s                r4   r—  Ú+_HandleToolCallsProjection._tool_calls_iterì  s)  é € Ø�‰×#Ñ# GÔ,Ø,.ˆØØŸ™×2Ñ2×6Ñ6Ó8×8ˆDØ‰|Ü"ØIó�ð %Ÿm™mžo�FØ—L‘L Ö%ñ .àØŸ8™8 HÓ-×3°ˆLÜ! ,Ó/ˆBØ—\‘\×+Ñ+Ó+ÙÜ/9¸,Ì×/MÑ/M�<×#Ñ# FÔ+ÐSWˆDÜ˜d¤D×)Ñ)ÙØŸ™ 'Ó*ˆJØŸ8™8 NÓ3ˆLÜ˜l¬C×0Ñ0ÚØ˜^Ó+Ø ŸH™H [Ó1�	Ü! )¬S×1Ñ1Ø "�IÜ'Ø!-Ø"ØŸ(™( 7Ó+Ü" 4§<¡<×#9Ñ#9Ó:ñ	�ð (.�|Ñ$Ø”ØÐ2Ó2ØŸ™ LÓ1�ØŸ™ Ó)�ØÑ%¬*°U¼C×*@Ñ*@Ø×&Ñ& uÔ-øØ˜Ó.ØŸ™ L°$Ó7�ØÑ%Ø—N‘N 4§8¡8¨HÓ#5Ô6øØ˜|Ó+ØŸ™ L°$Ó7�ØÑ%Ø"Ÿh™h yÓ1�GØ—L‘LÜ$¶W¤S¨¤\ÐBUÓVôò[ Ú8ùs   ‚AJÁJÁ	IJr‰  Nr„  rô   rõ   ©	r,   r-   r.   r/   r0   rm   rå   r—  r2   r+   r3   r4   rW  rW  ã  s   † ÙQôô'÷2r3   rW  c                  óH   • \ 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)rY  i!  zFSubgraphs projection that drains a `ScopedStreamHandle`'s tasks inbox.c                ó   • Xl         g rg   r‰  r‹  s     r4   rm   Ú#_HandleSubgraphsProjection.__init__$  rß   r3   c                ó"   • U R                  5       $ rg   ©Ú_subgraphs_iterrÃ   s    r4   rå   Ú$_HandleSubgraphsProjection.__aiter__'  ó   € Ø×#Ñ#Ó%Ð%r3   c                óœ  • U R                   nS GH9  u  p4[        X#5      n/ nUR                  5       (       d6  UR                  UR	                  5       5        UR                  5       (       d  M6  U HÕ  nUc  UR                  S5        M  UR                  S5      =(       d    0 n[        [        U5      5      n	Sn
UR                  5        H^  u  p¼[        UR                  5      n[        U	5      U:¼  d  M+  U	SU UR                  :X  d  M@  [        XÄ5      nUR                  U5        Sn
  O   U
(       a  MÄ  UR                  U5        M×     GM<     g)zýDrain non-blocking events from parent's messages/tools inboxes to grandchildren.

Called after each tasks event so grandchild handles receive events that
were enqueued in the parent handle's inboxes before (or just after) the
grandchild was discovered.
))rP  rP  )rQ  rQ  Nr?   FT)rŠ  ri  rj  r  rk  ra  rX   r?  rZ   r!  r§   rL  )rk   r#  Úparentrl  Úgrandchild_attrr"  rm  r`   Úevent_paramsrc  ÚroutedÚ_child_pathÚ
grandchildÚgrandchild_lenÚgc_inboxs                  r4   Ú'_route_sibling_inboxes_to_grandchildrenÚB_HandleSubgraphsProjection._route_sibling_inboxes_to_grandchildren*  s  € ð —‘ˆô,
Ñ'ˆJô 29¸Ó1LˆEØ*,ˆGà—k‘k—m‘mØ—‘˜u×/Ñ/Ó1Ô2ð —k‘k—m“mã �Ø‘=à×$Ñ$ TÔ*ÙØ$Ÿy™y¨Ó2×8°b�Ü Ô!1°,Ó!?Ó@�Ø�Ø/5¯|©|®~Ñ+�KÜ%(¨¯©Ó%9�Nä˜H›¨Õ7Ø$ _ nÐ5¸¿¹ÕHä@GØ&óA˜ð !×+Ñ+¨EÔ2Ø!%˜Ùñ 0>÷ �và×$Ñ$ UÖ+ô- !ò,
r3   c               óÒ  #   • U R                   R                  S5        [        5       n0 nU R                   R                  n U R                   R                  R                  5       I S h  v•N nUc;  UR                  5        H&  nUR                  S:X  d  M  UR                  S5        M(     g UR                  S5      =(       d    0 n[        U5      n[        U[        5      (       a  UR                  S5      OS n[        U[        5      (       d  MÏ  SU;   a¢  UR                  S5      n	U	(       d  Mï  [        U5      n
[        UR                  5       5       H]  u  p¼US S U
:w  a  M  UR                  U	:w  a  M"  [!        U5      u  pÞUR                  XÞ5        X+	 U R                   R#                  U5        M_     GMw  [%        Xs5      (       d  GMŠ  [        U5      nXñ;   a  GM�  UR'                  U5        [)        US   5      u  nn[+        U R                   R,                  UU=(       d    S UU R                   R.                  S	9nXÂU'   U R                   R1                  U5        U7v •  GM#   GNü7f)
NrK   r2  r[   r?   r_   r™   r>   éÿÿÿÿ)r¾   rL  rM  rN  rG  )rŠ  rz  rT  rL  rR  rX   rD   r8   r€  rZ   rV   rW   r?  rS   r!  rN  r=  rr  rA  rw  r;  rF  rº   rO  rn  )rk   Úseenr#  r@  rñ   ÚchildrY   r*   r_   Ú	result_idÚparent_pathÚ
child_pathÚchild_handler8   r9   rL  rM  rN  s                     r4   r¬  Ú*_HandleSubgraphsProjection._subgraphs_iterV  s  é € Ø�‰×#Ñ# GÔ,Ü%(£UˆØ<>ˆØ—‘×!Ñ!ˆØØŸ™×2Ñ2×6Ñ6Ó8×8ˆDØ‰|Ø#Ÿ]™]ž_�EØ—|‘| yÕ0ØŸ™ kÖ2ñ -ð ØŸ8™8 HÓ-×3°ˆLÜ(¨Ó6ˆIÜ/9¸,Ì×/MÑ/M�<×#Ñ# FÔ+ÐSWˆDÜ˜d¤D×)Ñ)ÙØ˜4ÓØ ŸH™H T›N�	Þ ÙÜ# IÓ.�Ü04°V·\±\³^Ö0DÑ,�JØ! # 2�¨+Ó5Ù Ø#×3Ñ3°yÓ@Ù Ü$?ÀÓ$E‘M�FØ ×(Ñ(¨Ô7ØÐ*Ø—L‘L×7Ñ7¸
ÖCñ 1Eò Ü# I×5Ñ5ÚÜ˜Ó#ˆDØ‹|ÚØ�H‰H�TŒNÜ*BÀ4ÈÁ8Ó*LÑ'ˆJ˜Ü-Ø—|‘|×+Ñ+ØØ%×-¨Ø /Ø#Ÿ|™|×;Ñ;ñˆLð (�4‰Lð �L‰L×-Ñ-¨lÔ;ØÓò[ Ú8ùs   ‚A&I'Á(I$Á)(I'ÂGI'r‰  Nr„  ©r   z!AsyncIterator[ScopedStreamHandle])r#  z)dict[tuple[str, ...], ScopedStreamHandle]r   r€   ©r   z(AsyncGenerator[ScopedStreamHandle, None])
r,   r-   r.   r/   r0   rm   rå   r¸  r¬  r2   r+   r3   r4   rY  rY  !  s,   † ÙPôô&ð*,à9ð*,ð 
ô*,÷X2r3   rY  c                  ó:   • \ rS rSrSrSS	S jjrS
S jrSS jrSrg)Ú_SubgraphsProjectioni‹  z8Discover direct child invocations for a namespace scope.c                ó   • Xl         X l        g rg   )rº   Ú_scope)rk   r¾   r@  s      r4   rm   Ú_SubgraphsProjection.__init__Ž  s   € ØŒØ�r3   c                ó"   • U R                  5       $ rg   r«  rÃ   s    r4   rå   Ú_SubgraphsProjection.__aiter__’  r®  r3   c               ó  ^ #   • T R                   R                  c  [        S5      e[        T R                  5      nT R                   R                  U5      n[        T R                  U 4S jS9nT R                  (       d  T R                   R                  5       OS n T R                   R                  U5      I S h  v•N   T R                   R                  5          UR                  R                  5       I S h  v•N nUcü   SnT R                   R                  nUba  UR                  5       (       aL  UR                  5       (       d7  UR                  5       n[!        U["        5      (       a  UR$                  S:X  a  SnUR&                  R)                  5        H&  n	U	R$                  S:X  d  M  U	R+                  U5        M(     T R                   R-                  UR.                  5        Ub  UR1                  S 5        g g UR                  S5      =(       d    0 n
UbH  UR                  S	5      S
:X  a3  [3        [5        U
5      5      T R                  :X  a  UR1                  U5        UR7                  U5       H  n	U	7v •  M
     GM¦   GNÆ GNŒ! SnT R                   R                  nUba  UR                  5       (       aL  UR                  5       (       d7  UR                  5       n[!        U["        5      (       a  UR$                  S:X  a  SnUR&                  R)                  5        H&  n	U	R$                  S:X  d  M  U	R+                  U5        M(     T R                   R-                  UR.                  5        Ub  UR1                  S 5        f f = f7f)Nú1AsyncThreadStream not entered - use `async with`.c                ó0   >• [        TR                  U UUS9$ ©N)r¾   rL  rM  rN  )rF  rº   ©rL  rM  rN  rk   s      €r4   r	  Ú6_SubgraphsProjection._subgraphs_iter.<locals>.<lambda>œ  s   ø€ Ü"ØŸ<™<ØØ)Ø$3ò	r3   ©r@  Úhandle_factoryr[   Úerroredr\   r2  r?   r^   rF   )rº   ré   rt   rD  rÈ  rê   r   Ú_activate_root_messages_inboxrë   rì   rA   rX   Ú	_run_doner�   r“   r™   rV   r6   r8   Ú_activerD   r€  rí   r>   ra  r?  rZ   rî   )rk   r?   rï   rð   r  rñ   Úterminal_statusÚrun_doner™   rµ   rY   s   `          r4   r¬  Ú$_SubgraphsProjection._subgraphs_iter•  s¹  øé € Ø�<‰<×"Ñ"Ñ*ÜÐRÓSÐSÜ.¨t¯{©{Ó;ˆØ�l‰l×1Ñ1°&Ó9ˆÜ"Ø—+‘+ôñ

ˆð AEÇÇˆD�L‰L×6Ñ6Ô8ÐQUð 	ð	,Ø—,‘,×0Ñ0°Ó8×8Ð8Ø�L‰L×/Ñ/Ô1ØØ ŸY™YŸ]™]›_×,�Ø‘<Øð /:ˆOØ—|‘|×-Ñ-ˆHØÑ#¨¯©¯©À×@RÑ@R×@TÑ@TØ!Ÿ™Ó*�Ü˜f¤l×3Ñ3¸¿¹ÈÓ8RØ&.�OØ!Ÿ/™/×0Ñ0Ö2�Ø—=‘= IÕ-Ø—N‘N ?Ö3ñ 3ð �L‰L×1Ñ1°#·&±&Ô9ØÑ%Ø×%Ñ% dÕ+ð &ð/  $Ÿx™x¨Ó1×7°R�àÑ*ØŸ™ Ó*¨jÓ8ÜÔ.¨|Ó<Ó=ÀÇÁÓLà×)Ñ)¨$Ô/Ø%Ÿl™l¨4Ö0�FØ •Lñ 1ò ò 9ò -øð  /:ˆOØ—|‘|×-Ñ-ˆHØÑ#¨¯©¯©À×@RÑ@R×@TÑ@TØ!Ÿ™Ó*�Ü˜f¤l×3Ñ3¸¿¹ÈÓ8RØ&.�OØ!Ÿ/™/×0Ñ0Ö2�Ø—=‘= IÕ-Ø—N‘N ?Ö3ñ 3ð �L‰L×1Ñ1°#·&±&Ô9ØÑ%Ø×%Ñ% dÕ+ð &üsS   ƒBNÂJ Â:J Â;<J Ã7JÃ8J Ä B'NÆ+ANÇ;BJ ÊJ ÊB(NÌ2ANÎN)rÈ  rº   N)r+   )r¾   r~   r@  r‚  r   r€   rÃ  rÄ  )	r,   r-   r.   r/   r0   rm   rå   r¬  r2   r+   r3   r4   rÆ  rÆ  ‹  s   † ÙBöô&÷2,r3   rÆ  c                  ó€   • \ rS rSrSrSSSS.           SS jjr\SS j5       rSS jrSS	 jr	SS
 jr
SS jrSrg)r£  iÊ  z*Async handle for one root-scope tool call.Né   )rI   r*   rG  c               ó
  • Xl         X l        X0l        [        U=(       d    / 5      U l        SU l        S U l        [        R                  " 5       nUR                  5       U l
        [        R                  " US9U l        SU l        g )NFrI  )rœ  r8  rI   rS   r*   r�   r9   rB   rŒ   r�   rá   rC   Ú_deltasÚ_deltas_consumed)rk   rœ  r8  rI   r*   rG  r—   s          r4   rm   ÚToolCallHandle.__init__Í  sh   € ð )ÔØŒ	ØŒ
Ü˜iŸo¨2Ó.ˆŒØˆŒ	Ø+/ˆŒ
Ü×'Ò'Ó)ˆØ+/×+=Ñ+=Ó+?ˆŒÜ29·-²-ÈÑ2WˆŒØ %ˆÕr3   c                óh   • U R                   (       a  [        S5      eSU l         U R                  5       $ )uÆ   Stream tool output deltas emitted before the terminal event.

Raises:
    RuntimeError: if called more than once â€” the underlying queue is
        single-consumer and cannot be fanned out safely.
z@ToolCallHandle.deltas can only be iterated by a single consumer.T)rß  rt   Ú_delta_iterrÃ   s    r4   ÚdeltasÚToolCallHandle.deltasá  s6   € ð × × ÜØRóð ð !%ˆÔØ×ÑÓ!Ð!r3   c               ój   #   •  U R                   R                  5       I S h  v•N nUc  g U7v •  M-   N7frg   )rÞ  rX   )rk   rñ   s     r4   râ  ÚToolCallHandle._delta_iterð  s2   é € ØØŸ™×)Ñ)Ó+×+ˆDØ‰|ØØ‹Jñ	 Ù+ùs   ‚3¡1¢3c                ó^   • U R                   (       a  g U R                  R                  U5        g rg   )r�   rÞ  ra  )rk   r   s     r4   r¤  ÚToolCallHandle._push_delta÷  s   € Ø�9�9ØØ�‰×Ñ Õ&r3   c                óà   • U R                   (       a  g SU l         U R                  R                  5       (       d  U R                  R                  U5        U R                  R	                  S 5        g ©NT)r�   rá   r‘   rÞ  ra  )rk   rá   s     r4   r€  ÚToolCallHandle._finishü  sJ   € Ø�9�9ØØˆŒ	Ø�{‰{×Ñ×!Ñ!Ø�K‰K×"Ñ" 6Ô*Ø�‰×Ñ Õ%r3   c                óì   • U R                   (       a  g SU l         Xl        U R                  R                  5       (       d  U R                  R                  U5        U R                  R                  S 5        g rê  )r�   r9   rá   r–   rÞ  ra  )rk   rš   s     r4   r¢  ÚToolCallHandle._fail  sO   € Ø�9�9ØØˆŒ	ØŒ
Ø�{‰{×Ñ×!Ñ!Ø�K‰K×%Ñ% cÔ*Ø�‰×Ñ Õ%r3   )	rÞ  rß  r�   r9   rI   r8  r*   rá   rœ  )rœ  r&   r8  r&   rI   r
   r*   r+  rG  r=   r   r€   )r   zAsyncIterator[str])r   zAsyncGenerator[str, None])r   r&   r   r€   )rá   r
   r   r€   ©rš   r•   r   r€   )r,   r-   r.   r/   r0   rm   Úpropertyrã  râ  r¤  r€  r¢  r2   r+   r3   r4   r£  r£  Ê  sy   † Ù4ð Ø&*Ø"ñ&ð ð&ð ð	&ð
 ð&ð $ð&ð ð&ð 
õ&ð( ó"ó ð"ôô'ô
&÷&r3   r£  c                  óF   • \ rS rSrSr S     S	S jjrS
S jrSS jrSrg)Ú_ToolCallsProjectioni  z4Typed projection for root-scope `thread.tool_calls`.Nc                óB   • Xl         [        U=(       d    / 5      U l        g rg   rù   rû   s      r4   rm   Ú_ToolCallsProjection.__init__  rý   r3   c                ó"   • U R                  5       $ rg   r–  rÃ   s    r4   rå   Ú_ToolCallsProjection.__aiter__  r™  r3   c               ó.  #   • U R                   R                  c  [        S5      e[        S/U R                  5      nU R                   R                  U5      n[        U R                  S S9n/ n U R                   R                  U5      I S h  v•N   U R                   R                  5          UR                  R                  5       I S h  v•N nUc÷   U R                   R                  nS nUbF  UR                  5       (       a1  UR                  5       (       d  UR                  5       nUR                  nUb  UO
[        S5      n	[!        UR"                  R%                  5       5       H  n
U
R'                  U	5        M     U H  n
U R                   R)                  U
5        M      U R                   R+                  UR,                  5        g UR/                  U5       H4  n
U R                   R1                  U
5        UR3                  U
5        U
7v •  M6     GMh   GNˆ GNN! U R                   R                  nS nUbF  UR                  5       (       a1  UR                  5       (       d  UR                  5       nUR                  nUb  UO
[        S5      n	[!        UR"                  R%                  5       5       H  n
U
R'                  U	5        M     U H  n
U R                   R)                  U
5        M      U R                   R+                  UR,                  5        f = f7f)NrÍ  rG   c                ó   • [        U UUUS9$ ©NrŸ  ©r£  rŸ  s       r4   r	  Ú7_ToolCallsProjection._tool_calls_iter.<locals>.<lambda>!  s   € ÜØ!-ØØØ'ò	r3   ©r*   rÓ  r›  )rº   ré   rt   rT   rú   rê   r   rë   rì   rA   rX   rÖ  r�   r“   r™   r9   rS   r×  rD   r¢  Ú_unregister_active_tool_callrí   r>   rî   Ú_register_active_tool_callr  )rk   r?   rï   rð   r  rñ   rÙ  Úterminal_errrÖ   rš   rµ   s              r4   r—  Ú%_ToolCallsProjection._tool_calls_iter  sx  é € Ø�<‰<×"Ñ"Ñ*ÜÐRÓSÐSÜ(¨'¨°D·O±OÓDˆØ�l‰l×1Ñ1°&Ó9ˆÜ"Ø—o‘oññ

ˆð ,.ˆ
ð	:Ø—,‘,×0Ñ0°Ó8×8Ð8Ø�L‰L×/Ñ/Ô1ØØ ŸY™YŸ]™]›_×,�Ø‘<Øð —|‘|×-Ñ-ˆHØ15ˆLØÑ#¨¯©¯©À×@RÑ@R×@TÑ@TØ#Ÿ?™?Ó,�Ø'Ÿ~™~�ð  Ñ+ñ ä!Ð"WÓXð ô
 ˜wŸ™×5Ñ5Ó7Ö8�Ø—‘˜SÖ!ñ 9ã$�Ø—‘×9Ñ9¸&ÖAñ %à�L‰L×1Ñ1°#·&±&Õ9ð1 &Ÿl™l¨4Ö0�FØ—L‘L×;Ñ;¸FÔCØ×%Ñ% fÔ-Ø •Lñ 1ò	 ò 9ò -øð —|‘|×-Ñ-ˆHØ15ˆLØÑ#¨¯©¯©À×@RÑ@R×@TÑ@TØ#Ÿ?™?Ó,�Ø'Ÿ~™~�ð  Ñ+ñ ä!Ð"WÓXð ô
 ˜wŸ™×5Ñ5Ó7Ö8�Ø—‘˜SÖ!ñ 9ã$�Ø—‘×9Ñ9¸&ÖAñ %à�L‰L×1Ñ1°#·&±&Õ9üsE   ‚A,LÁ/H ÂHÂ<H Ã
HÃH ÃC6LÇ	AH ÈH ÈC7LÌLr)  rg   r*  )r   zAsyncIterator[ToolCallHandle])r   z$AsyncGenerator[ToolCallHandle, None]r¦  r+   r3   r4   rñ  rñ    s3   † Ù>ð HLð0Ø'ð0Ø4Dð0à	õ0ô'÷0:r3   rñ  c                  ó,   • \ rS rSrSrSS jrSS jrSrg)	r\  iM  a  Mapping from extension name to custom event payload stream.

Repeated access for the same `name` returns the cached projection so that
callers receive the same subscription handle across multiple references to
`thread.extensions["foo"]` within one session.
c                ó*   • Xl         X l        0 U l        g rg   )rº   rú   Ú_cacherû   s      r4   rm   Ú_ExtensionsProjection.__init__U  s   € ØŒØ#ŒØ79ˆ�r3   c                ó¸   • U(       d  [        S5      eXR                  ;  a+  [        U R                  XR                  S9U R                  U'   U R                  U   $ )Nz!extension name must be non-empty.)r8  r*   )Ú
ValueErrorr  Ú_ExtensionProjectionrº   rú   )rk   r8  s     r4   Ú__getitem__Ú!_ExtensionsProjection.__getitem__Z  sL   € ÞÜÐ@ÓAÐAØ—{‘{Ó"Ü 4Ø—‘ 4·?±?ñ!ˆD�K‰K˜Ñð �{‰{˜4Ñ Ð r3   )r  rú   rº   N)r¾   r~   r*   r)   r   r€   )r8  r&   r   r  )r,   r-   r.   r/   r0   rm   r  r2   r+   r3   r4   r\  r\  M  s   † ñô:÷
!r3   r\  c                  óB   • \ rS rSr        SS jrSS jrS	S jrSrg)
r  id  c               ó(   • Xl         X l        X0l        g rg   )rº   Ú_namerú   )rk   r¾   r8  r*   s       r4   rm   Ú_ExtensionProjection.__init__e  s   € ð ŒØŒ
Ø#�r3   c                ó"   • U R                  5       $ rg   )Ú_iterrÃ   s    r4   rå   Ú_ExtensionProjection.__aiter__p  s   € Ø�z‰z‹|Ðr3   c               ó   #   • SSU R                    3/0nU R                  (       a  U R                  /US'   U R                  R                  U5      n[	        U R                   S9n U R                  R
                  (       a'   U R                  R                  UR                  5        g U R                  R                  U5      I S h  v•N   U R                  R                  5          UR                  R                  5       I S h  v•N nUc'   U R                  R                  UR                  5        g UR                  U5       H  nU7v •  M
     Mk   NŠ NO! U R                  R                  UR                  5        f = f7f)NrO   úcustom:rP   ©r8  )r  rú   rº   rê   r   rs   rí   r>   rë   rì   rA   rX   rî   )rk   r?   rï   rð   rñ   rò   s         r4   r  Ú_ExtensionProjection._iters  s%  é € Ø#-°'¸$¿*¹*¸Ð0FÐ/GÐ"HˆØ�?�?Ø$(§O¡OÐ#4ˆF�<Ñ Ø�l‰l×1Ñ1°&Ó9ˆÜ#¨¯©Ñ4ˆð	:Ø�|‰|×#×#Øð �L‰L×1Ñ1°#·&±&Õ9ð —,‘,×0Ñ0°Ó8×8Ð8Ø�L‰L×/Ñ/Ô1ØØ ŸY™YŸ]™]›_×,�Ø‘<Øð �L‰L×1Ñ1°#·&±&Õ9ð #Ÿ<™<¨Ö-�CØ•Iñ .ñ	 ñ 9ñ -øð �L‰L×1Ñ1°#·&±&Õ9üsN   ‚A"E>Á%E Â&E>Â'E ÃEÃ<E ÄEÄE Ä&E>Ä1 E ÅE Å'E;Å;E>)r  rú   rº   N)r¾   r~   r8  r&   r*   r)   r   r€   )r   zAsyncIterator[dict[str, Any]])r   z$AsyncGenerator[dict[str, Any], None])r,   r-   r.   r/   rm   rå   r  r2   r+   r3   r4   r  r  d  s7   † ð	$à!ð	$ð ð		$ð
 ð	$ð 
ô	$ô÷:r3   r  c                  ó^  • \ rS rSrSrSSSSSS.                 S4S jjr\S5S	 j5       rS5S
 jrS6S jr	\S7S j5       r
S8S jrS9S 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S8S jrS?S jrS@S jrSSS.       SAS jjr    SBS jr\SCS j5       r          SDS jr    SES jrS8S  jrS8S! jrSFS" jr SGS# jr!SHS$ jr" SI   SJS% jjr#SKS& jr$      SLS' jr%SS(.SMS) jjr&S8S* jr'S@S+ jr(SNS, jr)S8S- jr*SNS. jr+SOS/ jr,SGS0 jr-SPS1 jr.S@S2 jr/S3r0g)Qr~   iˆ  z®Async context manager for one thread's v3 streaming session.

Construct via `client.threads.stream(thread_id=None, *, assistant_id, ...)`
rather than instantiating directly.
NrÜ  FÚsse)rq   rG  Úrun_start_timeoutÚexplicit_thread_idÚtransport_kindc               ót  • Xl         [        U=(       d    0 5      U l        X l        X0l        XPl        X`l        Xpl        X€l        SU l	        S U l
        / U l        SU l        SU l        0 U l        [        5       U l        S U l        S U l        S U l        SU l        / U l        [,        R.                  " 5       U l        S U l        S U l        S U l        SU l        SU l        SU l        SU l        S U l         SU l!        S U l"        S U l#        [I        5       U l%        [I        5       U l&        S U l'        [Q        U 5      U l)        [U        U 5      U l+        [Y        U 5      U l-        []        U 5      U l/        [a        U / S9U l1        [e        U / S9U l3        [i        U SS9U l5        U Rj                  U l6        [o        U / S9U l8        g )	NFrž   é   gš™™™™™¹?g       @rK  r+   rC  )9rw   rW   rv   Ú	thread_idrx   rO  Ú_run_start_timeoutÚ_explicit_thread_idÚ_transport_kindrs   ré   Ú_open_handlesÚ_next_command_idÚ_next_subscription_idÚ_subscriptionsr   Ú_seen_event_idsÚ_shared_streamÚ_shared_stream_filterÚ_fanout_taskr3  r¦   rB   ÚLockr¥   Ú_lifecycle_watcher_taskÚ_lifecycle_watcher_handleÚ_lifecycle_cursorÚ!_lifecycle_max_reconnect_attemptsÚ_shared_max_reconnect_attemptsÚ_shared_reconnect_backoff_baseÚ_shared_reconnect_backoff_caprŽ   r’   rÖ  Ú_cursorrT  Ú_active_message_streamsÚ_active_tool_callsr  rƒ   Úrunrd   Úagentr¸   rá   rÛ   rD   r÷   rF   rñ  rX  rÆ  rZ  r[  r\  r]  )	rk   Úhttpr  rx   rq   rG  r  r  r  s	            r4   rm   ÚAsyncThreadStream.__init__�  s”  € ð Œ
Ü˜WŸ]¨Ó+ˆŒØ"ŒØ(ÔØ-ÔØ"3ÔØ#5Ô Ø-ÔØˆŒØ9=ˆŒØ68ˆÔØ !ˆÔØ%&ˆÔ"Ø8:ˆÔÜ,›ˆÔØ8<ˆÔØ<@ˆÔ"Ø7;ˆÔØ!&ˆÔØ24ˆŒô
 !(§¢£ˆÔØBFˆÔ$ØCGˆÔ&Ø-1ˆÔØ12ˆÔ.ð /0ˆÔ+Ø.1ˆÔ+Ø-0ˆÔ*Ø=AˆÔØ$ˆŒØ>BˆŒØ#'ˆŒÜBEÃ%ˆÔ$Ü7:³uˆÔð IMˆÔ!Ü˜T“?ˆŒÜ! $Ó'ˆŒ
Ü& tÓ,ˆŒÜ'¨Ó-ˆŒÜ+¨D¸BÑ?ˆŒÜ.¨t¸rÑBˆŒÜ-¨d¸"Ñ=ˆŒØŸ™ˆŒÜ/°ÀÑCˆ�r3   c                ó   • U $ )zóReturn self as the subscription controller (duck-type compatible with StreamController).

Exposes `_subscriptions` so tests can verify subscription counts via
`thread._controller._subscriptions` without requiring a separate controller object.
r+   rÃ   s    r4   Ú_controllerÚAsyncThreadStream._controllerÕ  s	   € ð ˆr3   c              ƒ  ón  #   • U R                   (       a  [        S5      eU R                  S:X  a  [        O[        nU" U R
                  R                  U R                  U R                  U R                  S9U l
        [        R                  " 5       R                  5       U l        U R                  5         U $ 7f)Nz5AsyncThreadStream is closed and cannot be re-entered.Ú	websocket)Úclientr  rq   rG  )rs   rt   r  r"   r!   rw   r;  r  rv   rO  ré   rB   rŒ   r�   rÖ  Ú!_ensure_lifecycle_watcher_running)rk   Útransport_clss     r4   Ú
__aenter__ÚAsyncThreadStream.__aenter__Þ  s•   é € Ø�<�<ÜÐVÓWÐWð ×#Ñ# {Ó2õ 'ä%ð 	ñ
 (Ø—:‘:×$Ñ$Ø—n‘nØ—M‘MØ×/Ñ/ñ	
ˆŒô !×1Ò1Ó3×AÑAÓCˆŒð 	×.Ñ.Ô0Øˆùs   ‚B3B5c              ƒ  ó„   #   •  U R                  5       I S h  v•N   g  N! [         a  nUc  e X$l         S nAg S nAff = f7frg   )r´   r•   Ú__context__)rk   Úexc_typeÚexcÚtbÚ	close_errs        r4   Ú	__aexit__ÚAsyncThreadStream.__aexit__ó  s9   é € ð	(Ø—*‘*“,×ÓøÜó 	(Ø‰{Øà$'×!Ñ!ûð		(üs0   ‚A „ —˜ œA � Ÿ
=©
8³A ¸=½A c                óÂ   • U R                   c  [        S5      eU R                   R                  S[        05      nU R                  R                  U5        UR                  $ )a  Return a fresh subscription to ALL channels.

Each property access opens a new subscription; callers iterating twice
will see two independent streams (both filtered by the same channel union).
Terminates when the stream closes (server hangup, `__aexit__`, or
transport-level close).
rè   rO   )ré   rt   Úopen_event_streamrM   r  r  Úeventsr‹  s     r4   rJ  ÚAsyncThreadStream.eventsü  sQ   € ð �?‰?Ñ"ÜÐTÓUÐUØ—‘×2Ñ2°JÄÐ3NÓOˆØ×Ñ×!Ñ! &Ô)Ø�}‰}Ðr3   c              ƒ  óÀ  #   • U R                   (       a  gSU l         U R                   H  nUR                  5       I Sh  v•N   M     U R                  nUb%  UR	                  5       (       d  UR                  5         U R                  b`  U R                  R                  5         [        R                  " [        [        R                  5         U R                  I Sh  v•N   SSS5        U R                  b"  U R                  R                  5       I Sh  v•N   U R                  [        R                  " 5       5        U R                  [        R                  " 5       5        U R                  b`  U R                  R                  5         [        R                  " [        [        R                  5         U R                  I Sh  v•N   SSS5        U R                   b"  U R                   R                  5       I Sh  v•N   U R"                  b#  U R"                  R                  5       I Sh  v•N   gg GNî GNT! , (       d  f       GNY= f GN2 N‡! , (       d  f       N‹= f Nc N67f)z(Tear down the thread stream. Idempotent.NT)rs   r  r´   rÖ  r�   Úcancelr(  Ú
contextlibÚsuppressÚ	ExceptionrB   ÚCancelledErrorr)  Ú_fail_active_message_streamsÚ_fail_active_tool_callsr&  r$  ré   )rk   rµ   rÙ  s      r4   r´   ÚAsyncThreadStream.close  s·  é € à�<�<ØØˆŒØ×(Ô(ˆFØ—,‘,“.× Ò ñ )ð —>‘>ˆØÑ¨¯©¯©Ø�O‰OÔØ×'Ñ'Ñ3Ø×(Ñ(×/Ñ/Ô1Ü×$Ò$¤Y´×0FÑ0FÕGØ×2Ñ2×2Ð2÷ Hà×)Ñ)Ñ5Ø×0Ñ0×6Ñ6Ó8×8Ð8Ø×)Ñ)¬'×*@Ò*@Ó*BÔCØ×$Ñ$¤W×%;Ò%;Ó%=Ô>Ø×ÑÑ(Ø×Ñ×$Ñ$Ô&Ü×$Ò$¤Y´×0FÑ0FÕGØ×'Ñ'×'Ð'÷ Hà×ÑÑ*Ø×%Ñ%×+Ñ+Ó-×-Ð-Ø�?‰?Ñ&Ø—/‘/×'Ñ'Ó)×)Ñ)ð 'ò' !ò 3÷ HÖGúò 9ñ (÷ HÕGúñ .á)ùs�   ‚<I¾H,¿BIÃH2ÃH/ÃH2Ã 2IÄIÄBIÆ/I	Æ?IÇ I	Ç2IÇ6IÇ7.IÈ%IÈ&IÈ/H2È2
IÈ<	IÉI	É	
IÉIÉIc                óÂ   • [        U R                  U[        R                  " U R                  S9S9nU =R                  S-  sl        X R
                  UR                  '   U$ )z6Allocate a subscription id and add it to the registry.rI  )r>   r?   rA   rž   )r;   r!  rB   rC   rO  r"  r>   )rk   r?   rï   s      r4   rê   Ú(AsyncThreadStream._register_subscription'  sT   € äØ×)Ñ)ØÜ—-’-¨×(<Ñ(<Ñ=ñ
ˆð
 	×"Ò" aÑ'Õ"Ø&)×Ñ˜CŸF™FÑ#Øˆ
r3   c                ó<   • U R                   R                  US5        g)zARemove a subscription from the registry. No-op if already absent.N)r"  rq  )rk   Úsubscription_ids     r4   rí   Ú*AsyncThreadStream._unregister_subscription2  s   € à×Ñ×Ñ °Õ6r3   c                óh   • U R                   c  [        R                  " 5       U l         U R                   $ )zÜCreate the root-scope messages inbox if not already active and return it.

Called by `_SubgraphsProjection` at scope `()` to capture messages events
that arrive at namespace `[]` before `thread.messages` has subscribed.
)r  rB   rC   rÃ   s    r4   rÕ  Ú/AsyncThreadStream._activate_root_messages_inbox6  s*   € ð ×$Ñ$Ñ,Ü(/¯ª«ˆDÔ%Ø×(Ñ(Ð(r3   c                ó:   • U R                   R                  U5        g rg   )r0  rw  ©rk   r  s     r4   r  Ú1AsyncThreadStream._register_active_message_stream@  s   € Ø×$Ñ$×(Ñ(¨Õ0r3   c                ó:   • U R                   R                  U5        g rg   )r0  Údiscardr]  s     r4   r  Ú3AsyncThreadStream._unregister_active_message_streamC  s   € Ø×$Ñ$×,Ñ,¨VÕ4r3   c                ó’   • [        U R                  5       H  nUR                  U5        M     U R                  R                  5         g rg   )rS   r0  ÚfailÚclear)rk   rš   r  s      r4   rR  Ú.AsyncThreadStream._fail_active_message_streamsF  s5   € Ü˜4×7Ñ7Ö8ˆFØ�K‰K˜Öñ 9à×$Ñ$×*Ñ*Õ,r3   c                ó:   • U R                   R                  U5        g rg   )r1  rw  r‹  s     r4   rý  Ú,AsyncThreadStream._register_active_tool_callK  s   € Ø×Ñ×#Ñ# FÕ+r3   c                ó:   • U R                   R                  U5        g rg   )r1  r`  r‹  s     r4   rü  Ú.AsyncThreadStream._unregister_active_tool_callN  s   € Ø×Ñ×'Ñ'¨Õ/r3   c                ó’   • [        U R                  5       H  nUR                  U5        M     U R                  R                  5         g rg   )rS   r1  r¢  rd  )rk   rš   rµ   s      r4   rS  Ú)AsyncThreadStream._fail_active_tool_callsQ  s5   € Ü˜4×2Ñ2Ö3ˆFØ�L‰L˜Öñ 4à×Ñ×%Ñ%Õ'r3   c                ó  • [        U R                  R                  5       5       HK  n[        R                  " [
        R                  5         UR                  R                  S5        SSS5        MM     g! , (       d  f       M_  = f)aa  Wake every active projection iterator on interrupt / run end.

Pushes the terminal sentinel (`None`) into every subscription
queue. Iterators see `None` and return; the shared SSE keeps
running so re-iteration after `run.respond(...)` registers a
fresh subscription and resumes.

`root_messages_inbox` is intentionally NOT signaled here: the
subgraphs projection that populates it is responsible for
pushing the terminal `None` in its own `finally` block, so any
message events it redirected to the inbox land before the
sentinel. Signaling root_inbox here would race the redirection
and could drop messages.
N)	rS   r"  rD   rN  rO  rB   Ú	QueueFullrA   ra  )rk   rï   s     r4   Ú_signal_pausedÚ AsyncThreadStream._signal_pausedV  sW   € ô" ˜×+Ñ+×2Ñ2Ó4Ö5ˆCÜ×$Ò$¤W×%6Ñ%6Õ7Ø—	‘	×$Ñ$ TÔ*÷ 8Ñ7ò 6ß7Ö7ús   ÁA4Á4
B	c                óv   • [        U[        5      (       a$  U R                  b  XR                  :”  a  Xl        ggg)zCAdvance the reconnect cursor from a command response meta sequence.N)rV   r=   r/  )rk   Úseqs     r4   Úobserve_applied_through_seqÚ-AsyncThreadStream.observe_applied_through_seqk  s/   € ä�cœ3×Ñ T§\¡\Ñ%9¸SÇ<Á<Ó=OØ�Lð >PÐr3   c                ó˜   • UR                  S5      n[        U[        5      (       a$  U R                  b  X R                  :”  a  X l        g g g ©Nrq  )rX   rV   r=   r/  ©rk   r`   rq  s      r4   Ú_observe_eventÚ AsyncThreadStream._observe_eventp  s=   € Ø�i‰i˜ÓˆÜ�cœ3×Ñ T§\¡\Ñ%9¸SÇ<Á<Ó=OØ�Lð >PÐr3   )rP   rQ   c               óŠ   • U R                   c  [        S5      eS[        U5      0nUb  X$S'   Ub  X4S'   U R                  U5      $ )züOpen a typed subscription against the shared SSE.

Returns an async iterator that yields raw `Event` dicts matching the
given filter. Multiple concurrent subscribes share one HTTP connection
whose union expands or rotates as subscriptions come and go.
rè   rO   rP   rQ   )ré   rt   rS   Ú_subscription_iter)rk   rO   rP   rQ   r?   s        r4   Ú	subscribeÚAsyncThreadStream.subscribeu  sT   € ð �?‰?Ñ"ÜÐTÓUÐUØ#-¬t°H«~Ð">ˆØÑ!Ø#-�<Ñ ØÑØ#�7‰OØ×&Ñ& vÓ.Ð.r3   c               ó¶  ^ #   • [        U5        T R                  c  [        S5      e0 n/ nU GH  nUS:X  a#  [        S5      X$'   UR	                  SS/05        M-  US;   a*  [        U/ S9X$'   UR	                  [        U// 5      5        M]  US:X  a+  [        / S S	9X$'   UR	                  [        S// 5      5        MŽ  US
:X  a+  [        / S S9X$'   UR	                  [        S// 5      5        M¿  US:X  a,  [        SU 4S jS9X$'   UR	                  [        S5      5        Mñ  [        US9X$'   UR	                  SSU 3/05        GM     U(       d  g[        [        [        [        [        [        [         ["        4      U5      5      5      nUR%                  S5      n/ n/ n T R'                  U5        Sh  v•N n	Ub  UR)                  U	5       H
  n
SU
47v •  M     [+        U	5      nT R-                  U5      nUc  ML  US:w  d  MT  UR%                  U5      nUc  Mj  UR)                  U	5       HZ  n
US
:X  a#  T R/                  U
5        UR	                  U
5        O(US:X  a"  T R1                  U
5        UR	                  U
5        XÊ47v •  M\     MÛ   NÖ
 T R3                  UR%                  S
5      UUU5        g! T R3                  UR%                  S
5      UUU5        f = f7f)a:  Yield `(channel_name, item)` tuples across multiple projections.

One shared subscription drives all per-channel decoders; items arrive
in server-emit order (the SDK analog of `GraphRunStream.interleave`).

Args:
    channels: Flat list of `"values"`, `"messages"`, `"tool_calls"`,
        `"subgraphs"`, and/or extension names. Built-ins yield their
        typed item (snapshot dict / `AsyncChatModelStream` /
        `ToolCallHandle` / `ScopedStreamHandle`); an extension yields
        its payload dict, keyed by the bare extension name.

Note:
    Handles and streams are yielded eagerly (before their sub-stream
    completes), so items arrive interleaved in real time. To receive a
    fully-resolved handle (output already populated), use the dedicated
    `thread.tool_calls` / `thread.messages` projections instead.
Nrè   rD   rO   )rE   rJ   rK   rK  rF   c                ó   • [        XUS9$ r  r   r  s      r4   r	  Ú:AsyncThreadStream.interleave_projections.<locals>.<lambda>³  s   € Ü,Ø&/Àzòr3   r  rX  c                ó   • [        U UUUS9$ rø  rù  rŸ  s       r4   r	  r  ½  s   € Ü&Ø)5Ø!%Ø"'Ø&/ò	r3   rû  rG   rZ  r+   c                ó   >• [        TU UUS9$ rÏ  )rF  rÐ  s      €r4   r	  r  Ê  s   ø€ Ü*Ø#'Ø!%Ø'1Ø,;ò	r3   rÒ  r  r  )r   ré   rt   r   r  rT   r   r   r   rD  r   r   r   r   rS   rW   r&   r
   rX   rz  rî   r   Ú_interleave_public_namerý  r  Ú_finalize_interleave_decoders)rk   rO   ÚdecodersÚ
sub_paramsÚchÚmergedrZ  Úregistered_tool_callsÚregistered_message_streamsr`   rñ   ÚwireÚpublicrð   s   `             r4   Úinterleave_projectionsÚ(AsyncThreadStream.interleave_projections‹  sÖ  øé € ô* 	% XÔ.Ø�?‰?Ñ"ÜÐTÓUÐUØ')ˆØ,.ˆ
ÜˆBØ�X‹~Ü*¨8Ó4�‘Ø×!Ñ! :°¨zÐ":Ö;ØÐ:Ó:ô
  +¨2¸Ñ<�‘Ø×!Ñ!Ô"9¸2¸$ÀÓ"CÖDØ�zÓ!Ü.Ø ñ$ñ �‘ð ×!Ñ!Ô"9¸:¸,ÈÓ"KÖLØ�|Ó#Ü/Ø ñ$ñ
 �‘ð ×!Ñ!Ô"9¸7¸)ÀRÓ"HÖIØ�{Ó"Ü/Øô$ñ
 �‘ð ×!Ñ!Ô"?ÀÓ"CÖDä0°bÑ9�‘Ø×!Ñ! :°'¸"¸°Ð/?Ð"@×Añc öd ØÜÜÜ ¤¤d¬4´´S°©>Ñ&:¸JÓ!GÓHó
ˆð —L‘L Ó-ˆ	ð 79ÐØACÐ"ð	Ø#×6Ñ6°vÔ>÷ 1�eØÑ(Ø )§¡¨uÖ 5˜Ø*¨DÐ1Õ1ñ !6ä$ UÓ+�Ø×5Ñ5°dÓ;�àÓ%¨&°KÕ*?Ø&Ÿl™l¨6Ó2�GØÓ*Ø$+§L¡L°Ö$7˜DØ%¨Ó5Ø $× ?Ñ ?ÀÔ EØ 5× <Ñ <¸TÕ BØ!'¨:Ó!5Ø $× DÑ DÀTÔ JØ :× AÑ AÀ$Ô GØ#) .Õ0ó %8ñ1Ð>ð& ×.Ñ.Ø—‘˜\Ó*ØØ%Ø*õ	øˆD×.Ñ.Ø—‘˜\Ó*ØØ%Ø*õ	üsV   ƒFKÆJ1 Æ/JÆ3J
Æ4JÆ7A J1 Ç;J1 ÈJ1 ÈA1J1 Ê
JÊJ1 Ê$KÊ1%KËKc                ód   • U c  gU S:X  a  gU R                  S5      (       a  U [        S5      S $ U $ )zMMap a wire channel name to the public channel name used in interleave tuples.NrG   rX  r  )Ú
startswithr§   )rŠ  s    r4   r‚  Ú)AsyncThreadStream._interleave_public_nameþ  s<   € ð ‰<ØØ�7‹?ØØ�?‰?˜9×%Ñ%Øœ˜I›Ð(Ð)Ð)Øˆr3   c                óö  • U R                   nUb:  UR                  5       (       a%  UR                  5       (       d  UR                  5       OSn[	        U[
        5      (       ab  Ub  UR                  b  UR                  O
[        S5      n[        UR                  R                  5       5       H  nUR                  U5        M     U H  nU R                  U5        M     U H  n	U R                  U	5        M     [	        U[        5      (       an  [	        U[        5      (       a  UR                   S:X  a  SOSn
UR                  R                  5        H&  nUR                   S:X  d  M  UR#                  U
5        M(     gg)aD  Finalize in-flight handles when `interleave_projections` tears down.

Mirrors the terminal handling of the dedicated `_ToolCallsProjection` /
`_SubgraphsProjection`: in-flight tool calls are failed (so awaiting
`handle.output` can't hang) and discovered subgraph children are
force-completed with the run's terminal status.
Nr›  rÔ  r\   r[   r2  )rÖ  r�   r“   r™   rV   r   r9   rt   rS   r×  rD   r¢  rü  r  r   r6   r8   r€  )rk   rX  rZ  rˆ  r‰  rÙ  Úresolvedrš   rµ   r  rØ  r½  s               r4   rƒ  Ú/AsyncThreadStream._finalize_interleave_decoders	  sB  € ð —>‘>ˆð Ñ#¨¯©¯©À×@RÑ@R×@TÑ@Tð �O‰OÔàð 	ô
 �jÔ"2×3Ñ3ð Ñ'¨H¯N©NÑ,Fð —’ä!Ð"WÓXð ô
 ˜z×1Ñ1×8Ñ8Ó:Ö;�Ø—‘˜SÖ!ñ <ã+ˆFØ×-Ñ-¨fÖ5ñ ,ã0ˆFØ×2Ñ2°6Ö:ñ 1ä�iÔ!1×2Ñ2ô ˜h¬×5Ñ5¸(¿/¹/ÈYÓ:Vñ à ð ð
 #×*Ñ*×1Ñ1Ö3�Ø—<‘< 9Õ,Ø—M‘M /Ö2ò 4ð 3r3   c               ó¸  #   • U R                  U5      n U R                  (       a   U R                  UR                  5        g U R	                  U5      I S h  v•N   U R                  5          UR                  R                  5       I S h  v•N nUc   U R                  UR                  5        g U7v •  MI   N^ N-! U R                  UR                  5        f = f7frg   )rê   rs   rí   r>   rë   rì   rA   rX   )rk   r?   rï   rñ   s       r4   rz  Ú$AsyncThreadStream._subscription_iter3  s¸   é € ð ×)Ñ)¨&Ó1ˆð	2Ø�|�|Øð ×)Ñ)¨#¯&©&Õ1ð ×(Ñ(¨Ó0×0Ð0Ø×'Ñ'Ô)ØØ ŸY™YŸ]™]›_×,�Ø‘<Øð ×)Ñ)¨#¯&©&Õ1ð “
ñ	 ñ 1ñ -øð
 ×)Ñ)¨#¯&©&Õ1üsK   ‚C•B: §CÁB: ÁB6Á2B: Â
B8ÂB: ÂCÂ/B: Â8B: Â:CÃCc                ó°   • U R                   b  U R                   R                  5       (       a*  [        R                  " U R	                  5       5      U l         g g rg   )r&  r�   rB   rË   Ú_fanoutrÃ   s    r4   rì   Ú(AsyncThreadStream._ensure_fanout_runningD  sA   € Ø×ÑÑ$¨×(9Ñ(9×(>Ñ(>×(@Ñ(@Ü '× 3Ò 3°D·L±L³NÓ CˆDÕð )Ar3   c              ƒ  óR  #   • SSK Jn  U R                  (       Gd¡  U R                  nUc  g U R	                  UR
                  5        Sh  v•N nU R                  (       a    O“U R                  U5        [        U R                  R                  5       5       H7  nU" X4R                  5      (       d  M  UR                  R                  U5        M9     [        U5      (       d  M�  U R                  5         M¯  U R                  UL a£  UR                   I Sh  v•N nUb‹  [#        U[$        R&                  5      (       dl  U R                  (       d[  [(        R*                  " [        5         UR-                  5       I Sh  v•N   SSS5        U R/                  5       I Sh  v•N (       a  GMž  OU R                  (       d  GM¡  U R                  R                  5        H  nUR                  R                  S5        M      g GN®
 GN	! [         a     GNf = f Nû N—! , (       d  f       N›= f NŠ7f)a}  Single consumer of the shared SSE; routes events to subscriptions.

Why: rotation in `_reconcile_stream` replaces `_shared_stream` mid-loop.
Re-read `self._shared_stream` on each outer iteration so we always
consume from the current handle. The old handle's iterator exhausts
naturally after `_close_after` closes it.

On a post-ready transport drop (non-cancelled error in `shared.done`),
attempts to reconnect up to `_shared_max_reconnect_attempts` times so
scoped projections (subgraph child handles, message streams) survive
without losing buffered events. The reconnect replays `since=<cursor>`
and `_dedup_iter` drops any overlap.
r   )Úmatches_subscriptionN)Ú!langgraph_sdk.stream.subscriptionrš  rs   r$  Ú_dedup_iterrJ  rw  rS   r"  rD   r?   rA   ra  rb   rn  rP  r�   rV   rB   rQ  rN  rO  r´   Ú_reconnect_shared_stream)rk   rš  Úsharedr`   rï   rš   s         r4   r—  ÚAsyncThreadStream._fanoutH  sš  é € õ 	Kà—,—,�,Ø×(Ñ(ˆFØ‰~ØðØ#'×#3Ñ#3°F·M±MÔ#B÷ .˜%Ø—|—|ÙØ×'Ñ'¨Ô.Ü# D×$7Ñ$7×$>Ñ$>Ó$@ÖA˜Ù/°·z±z×BÓBØŸI™I×0Ñ0°Ö7ñ  Bô 3°5×9Ó9Ø×+Ñ+Ö-ð ×"Ñ" fÒ,ð
 #ŸK™K×'�à‘OÜ& s¬G×,BÑ,B×CÑCØ ŸLŸLä#×,Ò,¬YÕ7Ø$Ÿl™l›n×,Ð,÷ 8à!×:Ñ:Ó<×<Õ<Ú ØðO —,—,’,ðV ×&Ñ&×-Ñ-Ö/ˆCØ�I‰I× Ñ  Ö&ò 0òM.Ò#Bøô  ó âðúñ (ñ -÷ 8Õ7úá<ùs´   ‚)H'¬G? ÁG<ÁG9ÁG<ÁG? Á"H'Á#AG? Â3-G? Ã$G? Ã6H'ÄHÄAH'Å&HÅ:HÅ;HÅ?H'ÆH%ÆH'Æ<=H'Ç9G<Ç<G? Ç=H'Ç?
HÈ	H'ÈHÈH'ÈHÈ
H"ÈH'c              ƒ  óÖ   #   • U R                   nU R                  n[        X2SU-  -  5      n[        R                  " SUS-  5      n[
        R                  " XE-   5      I Sh  v•N   g N7f)zJSleep with exponential backoff and jitter for reconnect attempt `attempt`.é   r   g      Ð?N)r-  r.  ÚminÚrandomÚuniformrB   r³   )rk   ÚattemptÚbaseÚcapr±   Újitters         r4   Ú_reconnect_sleepÚ"AsyncThreadStream._reconnect_sleep†  sV   é € à×2Ñ2ˆØ×0Ñ0ˆÜ�C  G¡Ñ,Ó-ˆÜ—’  5¨4¡<Ó0ˆÜ�mŠm˜E™NÓ+×+Ó+ùs   ‚AA)Á!A'Á"A)c              ƒ  óâ  #   • U R                   c  gU R                  nUc  g[        U R                  5       Hs  nU R                  (       a    g[        U5      nU R                  b  U R                  US'    U R                   R                  U5      nUR                  I Sh  v•N   X@l          g   g N! [        R                   a    e [         a    U R                  U5      I Sh  v•N     Mµ  f = f7f)zýAttempt to reopen the shared stream after a post-ready transport drop.

Returns:
    `True` if a new stream was opened (caller should resume fanout),
    `False` if all reconnect attempts were exhausted or the controller
    was closed in the meantime.
NFÚsinceT)ré   r%  Úranger,  rs   rW   r/  rI  ÚreadyrB   rQ  rP  r©  r$  )rk   Úbase_filterr¥  Ústream_paramsÚ
new_streams        r4   r�  Ú*AsyncThreadStream._reconnect_shared_streamŽ  sß   é € ð �?‰?Ñ"Øð ×0Ñ0ˆØÑØÜ˜T×@Ñ@ÖAˆGØ�|�|ÙÜ,0°Ó,=ˆMØ�|‰|Ñ'Ø)-¯©�˜gÑ&ðØ!Ÿ_™_×>Ñ>¸}ÓM�
Ø ×&Ñ&×&Ð&ð #-ÔÙñ Bð  ñ 'øÜ×)Ñ)ó ØÜó Ø×+Ñ+¨GÓ4×4Ñ4ÚðüsH   ‚A1C/Á4*B0ÂB.ÂB0Â#C/Â.B0Â02C,Ã"C%Ã#C,Ã(C/Ã+C,Ã,C/c              ƒ  óJ  #   • U R                  U R                  S9I Sh  v•N   SSKJn  U R                  c  [        S5      eU R                  b/  U R                  b"  U" U R                  [        U5      5      (       a  gU R                  US9n[        U5      nU R                  b  U R                  US'   U R                  R                  U5      nU R                  nXPl        X0l        UR                  I Sh  v•N   Ub   [        R                  " [        U5      5        gg Nÿ N*7f)aÕ  Ensure the shared SSE covers `candidate_filter`. Rotate if not.

Open-new-before-close-old: any events buffered server-side between
the two opens are replayed on the new SSE, and the per-thread
`_seen_event_ids` set dedupes the overlap. Awaits `new_stream.ready`
so the HTTP connection is established before returning, guaranteeing
that both old and new streams are simultaneously connected during
rotation (enabling correct peak-count tracking and dedup correctness).
rÏ   Nr   )Úfilter_coversrè   )Úextrar¬  )Ú_await_run_start_gater  r›  r´  ré   rt   r$  r%  rW   Ú_compute_current_unionr/  rI  r®  rB   rË   r¶   )rk   Úcandidate_filterr´  Ú
new_filterr°  r±  Ú
old_streams          r4   rë   Ú#AsyncThreadStream._reconcile_stream¯  s  é € ð ×(Ñ(°×1HÑ1HÐ(ÐI×IÐIÝCà�?‰?Ñ"ÜÐTÓUÐUð ×ÑÑ+Ø×*Ñ*Ñ6Ù˜d×8Ñ8¼$Ð?OÓ:P×QÑQàà×0Ñ0Ð7GÐ0ÐHˆ
Ü(,¨ZÓ(8ˆØ�<‰<Ñ#Ø%)§\¡\ˆM˜'Ñ"Ø—_‘_×6Ñ6°}ÓEˆ
Ø×(Ñ(ˆ
Ø(ÔØ%/Ô"ð ×Ñ×ÐØÑ!ô ×Ò¤¨ZÓ 8Õ9ð "ñ3 	Jñ0 	ùs"   ‚D#ŸD CD#Ã6D!Ã7)D#Ä!D#c                ó  • SSK Jn  U R                  R                  5        Vs/ sH  n[	        UR
                  5      PM     nnUb  UR                  [	        U5      5        UR                  SS/05        U" U5      $ s  snf )Nr   )r   rO   rH   )r›  r   r"  rD   rW   r?   r  )rk   rµ  r   rï   Úfilterss        r4   r·  Ú(AsyncThreadStream._compute_current_union×  s|   € õ 	Kð )-×(;Ñ(;×(BÑ(BÔ(Dó)
Ù(D ŒD�—‘ÖÑ(Dð 	ð )
ð ÑØ�N‰Nœ4 ›;Ô'ð 	�‰˜
 [ MÐ2Ô3Ù# GÓ,Ð,ùò)
s   £A<c               ó¶   #   • U  S h  v•N nUR                  S5      nUb,  X0R                  ;   a  M.  U R                  R                  U5        U7v •  MP   NK
 g 7f)NÚevent_id)rX   r#  rw  )rk   Úsourcer`   rÀ  s       r4   rœ  ÚAsyncThreadStream._dedup_iterí  sR   é € Ù!÷ 	�%Ø—y‘y Ó,ˆHØÑ#Ø×3Ñ3Ó3ÙØ×$Ñ$×(Ñ(¨Ô2Ø�Kñ	™6ùs&   ‚A…A‰AŠA�AAÁAÁAc              ƒ  óL  #   • U R                   c  [        S5      eU R                  nU =R                  S-  sl        U R                   R                  X1US.5      I Sh  v•N nUc  0 $ UR	                  S5      S:X  a5  UR	                  SS5      nUR	                  SS	5      n[        S
U SU 35      eUR	                  S5      n[        U[        5      (       a9  UR	                  S5      nU R                  b  U R                  R                  U5        UR	                  S0 5      $  NÄ7f)zãSend a protocol command and return the `result` payload.

Returns `{}` for 202/204 responses (no body). Raises `RuntimeError`
with the protocol code/message when the server returns an error
envelope (`{"type": "error", ...}`).
Nrè   rž   )r>   r^   r?   Útyper9   Úunknownr¡  rž  zProtocol error [z]: ÚmetaÚapplied_through_seqr™   )	ré   rt   r   Úsend_commandrX   rV   rW   r7  rr  )	rk   r^   r?   Ú
command_idr¤   Úcoder¡  rÆ  rÇ  s	            r4   r�   ÚAsyncThreadStream._send_commandö  s  é € ð �?‰?Ñ"ÜÐTÓUÐUØ×*Ñ*ˆ
Ø×Ò Ñ"ÕØŸ™×5Ñ5Ø¸6ÑBó
÷ 
ˆð ÑàˆIØ�<‰<˜Ó 7Ó*Ø—<‘< ¨Ó3ˆDØ—l‘l 9¨bÓ1ˆGÜÐ!1°$°°s¸7¸)ÐDÓEÐEØ�|‰|˜FÓ#ˆÜ�dœD×!Ñ!Ø"&§(¡(Ð+@Ó"AÐØ×ÑÑ+Ø× Ñ ×<Ñ<Ð=PÔQØ�|‰|˜H bÓ)Ð)ñ
ùs   ‚AD$ÁD"ÁCD$rÏ   c             ƒ  óÜ   #   • U R                   nUb  UR                  5       (       a  gUc  UI Sh  v•N   g[        R                  " [        R                  " U5      US9I Sh  v•N   g N7 N7f)a  Wait for the current run.start to commit the thread server-side.

No-op when no run.start is in flight. Re-raises if run.start failed.
Raises `asyncio.TimeoutError` if `timeout` is set and the gate does
not resolve in time; the gate itself is left intact for later callers.
NrÏ   )rŽ   r�   rB   rÓ   Úshield)rk   rÆ   r˜   s      r4   r¶  Ú'AsyncThreadStream._await_run_start_gate  sW   é € ð ×$Ñ$ˆØ‰<˜4Ÿ9™9Ÿ;™;ØØ‰?Ø�J‰Jä×"Ò"¤7§>¢>°$Ó#7ÀÑI×IÑIñ áIùs!   ‚.A,°A(±1A,Á"A*Á#A,Á*A,c                ór   • U R                   b  g [        R                  " U R                  5       5      U l         g rg   )r(  rB   rË   Ú_run_lifecycle_watcherrÃ   s    r4   r<  Ú3AsyncThreadStream._ensure_lifecycle_watcher_running#  s0   € Ø×'Ñ'Ñ3ØÜ'.×':Ò':Ø×'Ñ'Ó)ó(
ˆÕ$r3   c                ó˜   • UR                  S5      n[        U[        5      (       a$  U R                  b  X R                  :”  a  X l        g g g ru  )rX   rV   r=   r*  rv  s      r4   Ú_observe_lifecycle_eventÚ*AsyncThreadStream._observe_lifecycle_event*  sE   € Ø�i‰i˜ÓˆÜ�cœ3×ÑØ×"Ñ"Ñ*¨c×4JÑ4JÓ.Jà%(Õ"ð /Kð  r3   c                óJ   • SSS/0nU R                   b  U R                   US'   U$ )NrO   rH   rI   r¬  )r*  )rk   r?   s     r4   Ú_lifecycle_stream_paramsÚ*AsyncThreadStream._lifecycle_stream_params1  s1   € Ø",¨{¸GÐ.DÐ!EˆØ×!Ñ!Ñ-Ø"×4Ñ4ˆF�7‰OØˆr3   c           
   ƒ  ó¦  #   • U R                   c  gSnU R                  (       d«   U R                   R                  U R                  5       5      nX l        [
        R                  " UR                  SS9I Sh  v•N   UR                    Sh  v•N nU R                  (       a    gU R                  U5        U R                  U5      I Sh  v•N   MH  g NY NF N
 UR                  I Sh  v•N  nUb  [        U[
        R                  5      (       aJ  UcF  U R                  nUb7  UR                  5       (       d"  UR                  [!        S[#        S5      S95        gUS-  nXR$                  :”  a  Ue[
        R&                  " S	5      I Sh  v•N    O®! [
        R                   a    e [(         a�  nUS-  nXR$                  ::  a&  [
        R&                  " S	5      I Sh  v•N     SnAGMÓ  U R                  nUb:  UR                  5       (       d%  UR                  [!        S[#        S
U 35      S95         SnAgSnAff = fU R                  (       d  GM,  g7f)z3Always-on SSE consuming lifecycle + input channels.Nr   g      @rÏ   rÔ  z,lifecycle stream ended before terminal event©r8   r9   rž   gš™™™™™©?zLifecycle transport failed: )ré   rs   rI  rÖ  r)  rB   rÓ   r®  rJ  rÓ  Ú_apply_lifecycle_eventr�   rV   rQ  rÖ  r‘   r6   rt   r+  r³   rP  )rk   Úreconnect_attemptsrµ   r`   rš   rÙ  rC  s          r4   rÐ  Ú(AsyncThreadStream._run_lifecycle_watcher7  sù  é € à�?‰?Ñ"ØØÐØ—,—,ð/ØŸ™×:Ñ:Ø×1Ñ1Ó3ó�ð 28Ô.Ü×&Ò& v§|¡|¸SÑA×AÐAØ#)§=¢=÷ =˜%Ø—|—|ÙØ×1Ñ1°%Ô8Ø×5Ñ5°eÓ<×<Ò<ð ñ Bñ=ñ =ð	 $1ð
 #ŸK™K×'Ð'�Ø‘;¤*¨S´'×2HÑ2H×"IÑ"Ið ‘{Ø#'§>¡>˜Ø#Ñ/¸¿¹¿¹Ø$×/Ñ/Ü ,Ø+4Ü*6Ø(Vó+&ñ!"ôð Ø" aÑ'Ð"Ø%×(NÑ(NÓNØ�IÜ—m’m DÓ)×)Ò)øÜ×)Ñ)ó ØÜó Ø" aÑ'Ð"Ø%×)OÑ)OÓOÜ!Ÿ-š-¨Ó-×-Ñ-ÞØŸ>™>�ØÑ'°·±·±Ø×'Ñ'Ü$Ø#,Ü".Ð1MÈcÈUÐ/SÓ"Tñôô ûðúðG —,—,“,ùsÍ   ‚"I¥AF Á6CÁ7F ÂCÂCÂCÂF Â"IÂ#%F ÃCÃ	F ÃIÃF ÃCÃF ÃF Ã&C)Ã'A0F ÅIÅ/F ÆF
ÆF ÆIÆH9Æ,-H4ÇGÇH4ÇIÇ&A	H4È/IÈ4H9È9IÉIc              ƒ  ó”   #   • U R                   R                  SU R                   S3U R                  =(       d    SS9I Sh  v•N $  N7f)z6Fetch the current thread state from the REST endpoint.z	/threads/z/stateN)rq   )rw   rX   r  rv   rÃ   s    r4   rÑ   ÚAsyncThreadStream._fetch_staten  sF   é € à—Z‘Z—^‘^Ø˜Ÿ™Ð' vÐ.Ø—M‘M×) Tð $ð 
÷ 
ð 	
ñ 
ùs   ‚?AÁAÁAc                óh   • UR                  S5      (       + =(       a    UR                  S5      (       + $ )zEReturn `True` if the thread state has no pending tasks or next nodes.r¨   rK   )rX   )rk   rÕ   s     r4   rÒ   Ú$AsyncThreadStream._state_is_terminalu  s%   € à—9‘9˜VÓ$Ô$×?¨U¯Y©Y°wÓ-?Ô)?Ð?r3   c                óJ   • U R                   =(       a    U R                  (       + $ )z÷Return `True` if we can try the REST state before waiting on the lifecycle.

True only when the caller passed an explicit `thread_id` (not a minted
UUID) and no run has been seen yet, indicating a potential reattach to
an already-terminal thread.
)r  r’   rÃ   s    r4   rÐ   Ú8AsyncThreadStream._can_return_existing_state_immediatelyy  s   € ð ×'Ñ'×>°·±Ô,>Ð>r3   c              ƒ  óÀ   #   • U R                   c  [        S5      eU R                  (       d  U R                  (       d  [        S5      eU R                   I Sh  v•N $  N7f)z¾Await `_run_done`, raising if the stream was never entered or no run exists.

Raises:
    RuntimeError: stream not entered, or no run started and no explicit
        thread_id was provided.
Nu0   AsyncThreadStream not entered â€” use async withzmthread.output: no run has been started and no explicit thread_id was provided. Call thread.run.start() first.)rÖ  rt   r’   r  rÃ   s    r4   rÔ   Ú$AsyncThreadStream._wait_for_run_done‚  sP   é € ð �>‰>Ñ!ÜÐQÓRÐRØ�~�~ d×&>×&>Üð?óð ð —^‘^×#Ð#Ñ#ùs   ‚AAÁAÁAc              ƒ  ó>  #   • UR                  S5      nUS:X  GaG  UR                  S5      =(       d    0 n[        U[        5      (       a  UR                  S5      OSn[        U[        5      (       a  UR                  S5      OSn[        U[        5      (       aÇ  U[        U[        5      (       a  UR                  S5      OS[        U[        5      (       a  UR                  S5      =(       d    / O/ S	.nU R                   ISh  v•N   U R
                  nU R                  R                  U5        S
U l        SSS5      ISh  v•N   W(       d  U R                  5         gggUS:X  Gar  UR                  S5      =(       d    0 n[        U[        5      (       a  UR                  S5      OSn[        U[        5      (       a  UR                  S5      OSnUS;   a  S
U l	        gUS;   aó  U R                   ISh  v•N   SU l        / U l        SSS5      ISh  v•N   U R                  n	U	b°  U	R                  5       (       dš  US:X  a{  [        U[        5      (       a  UR                  S5      OSn
[        U
(       a  SU
 3OS5      nU R                  U5        U R                  U5        U	R                  [!        SUS95        gU	R                  [!        SS95        ggggg GN× GNœ! , ISh  v•N  (       d  f       GN²= f GN Nê! , ISh  v•N  (       d  f       GN = f7f)zRUpdate `interrupted` / `interrupts` / `_run_done` from a lifecycle or input event.r^   zinput.requestedr?   r_   Nr'   r(   r*   )r'   r(   r*   TrH   r`   )r2  Úrunning)r[   r\   Fr\   r9   zRun errored: zRun erroredrÔ  rÙ  r[   )r8   )rX   rV   rW   r&   r¥   r3  r¦   r  rn  r’   rÖ  r�   rt   rR  rS  r‘   r6   )rk   r`   r^   r?   r_   r'   ÚpayloadÚwas_interruptedÚphaserÙ  Ú	error_msgr9   s               r4   rÚ  Ú(AsyncThreadStream._apply_lifecycle_event’  s{  é € à—‘˜8Ó$ˆØÐ&Ô&Ø—Y‘Y˜xÓ(×.¨BˆFÜ)3°F¼D×)AÑ)A�6—:‘:˜fÔ%ÀtˆDÜ7AÀ$Ì×7MÑ7M˜4Ÿ8™8 NÔ3ÐSWˆLÜ˜,¬×,Ñ,à$0Ü2<¸TÄ4×2HÑ2H˜TŸX™X gÔ.Èdä! &¬$×/Ñ/ð "(§¡¨KÓ!8×!>¸Bøàñ-�ð  ×0×0Ó0Ø&*×&6Ñ&6�OØ—O‘O×*Ñ*¨7Ô3Ø'+�DÔ$÷ 1×0ö 'Ø×'Ñ'Õ)ð 'ð' -ð* �{Ô"Ø—Y‘Y˜xÓ(×.¨BˆFÜ)3°F¼D×)AÑ)A�6—:‘:˜fÔ%ÀtˆDÜ)3°D¼$×)?Ñ)?�D—H‘H˜WÔ%ÀTˆEØÐ.Ó.ð "&�•ØÐ1Ó1ð  ×0×0Ó0Ø',�DÔ$Ø&(�D”O÷ 1×0ð  Ÿ>™>�ØÑ'°·±·±Ø Ó(ä1;¸DÄ$×1GÑ1G˜DŸH™H WÔ-ÈTð "ô !-Þ;D˜m¨I¨;Ñ7È-ó!˜ð ×9Ñ9¸%Ô@Ø×4Ñ4°UÔ;Ø ×+Ñ+¬LÀ	ÐQVÑ,WÕXà ×+Ñ+¬LÀÑ,LÕMð 1@Ð'ð 2ð #÷ 1×0×0Õ0ú÷6 1×0×0Ô0üsŒ   ‚DLÄKÄLÄ	/K"Ä8LÅKÅB3LÇ7K=Ç8LÇ;LÈ
LÈL ÈCLËLË"K:Ë(K+Ë)K:Ë5	LÌ LÌLÌLÌ	LÌL),r0  r1  rs   r/  r  r&  rv   rw   r¥   r*  r+  r)  r(  rO  r   r!  r  r  rÖ  r’   rŽ   r  r#  r,  r-  r.  r$  r%  r"  ré   r  r3  rx   r]  r3  r¦   rF   rá   r2  r[  rZ  r  rX  rD   )r4  r   r  r&   rx   r&   rq   r�   rG  r=   r  úfloat | Noner  Úboolr  zLiteral['sse', 'websocket']r   r€   )r   r~   )rB  r
   rC  r
   rD  r
   r   r€   )r   úAsyncIterator[Event]r…  )r?   r   r   r;   )rX  r=   r   r€   )r   r-  )r  r   r   r€   rî  )rµ   r£  r   r€   )rq  r
   r   r€   rƒ  )rO   r)   rP   zlist[list[str]] | NonerQ   z
int | Noner   rî  )rO   r)   r   zAsyncIterator[tuple[str, Any]])rŠ  r°   r   r°   )
rX  úDecoder | NonerZ  rï  rˆ  zlist[ToolCallHandle]r‰  zlist[AsyncChatModelStream]r   r€   )r?   r   r   zAsyncGenerator[Event, None])r¥  r=   r   r€   )r   rí  )r¸  r   r   r€   rg   )rµ  zSubscribeParams | Noner   r¯   )rÁ  rî  r   rî  )r^   r&   r?   r¯   r   r¯   )rÆ   rì  r   r€   )r   r¯   )rÕ   r¯   r   rí  )r   r6   )1r,   r-   r.   r/   r0   rm   rï  r7  r>  rF  rJ  r´   rê   rí   rÕ  r  r  rR  rý  rü  rS  rn  rr  rw  r{  rŒ  Ústaticmethodr‚  rƒ  rz  rì   r—  r©  r�  rë   r·  rœ  r�   r¶  r<  rÓ  rÖ  rÐ  rÑ   rÒ   rÐ   rÔ   rÚ  r2   r+   r3   r4   r~   r~   ˆ  s@  † ñð -1Ø"Ø*.Ø#(Ø6;ñDDð ðDDð ð	DDð
 ðDDð *ðDDð ðDDð (ðDDð !ðDDð 4ðDDð 
õDDðL óó ðôô*(ð óó ðô*ô8	ô7ô)ô1ô5ô-ô
,ô0ô(ô
+ô*ô
ð .2Ø ñ/àð/ð +ð	/ð
 ð/ð 
õ/ð,qØ!ðqà	'ôqðf óó ðð(3à"ð(3ð "ð(3ð  4ð	(3ð
 %?ð(3ð 
ô(3ðT2Ø%ð2à	$ô2ô"Dô<'ô|,ôôB&:ðR /3ð-Ø+ð-à	õ-ô,ð*Øð*Ø#1ð*à	ô*ð< FJ÷ Jô
ô)ôô5ôn
ô@ô?ô$÷ :Nr3   r~   )rO   r)   r*   r)   r   r   )rY   r
   r   r)   )r`   r
   r   rí  )rµ   r    r±   rÙ   r   r€   )r_   r¯   r   r°   rg   )r_   r¯   r  r°   r   r&   )r7  r&   r   ztuple[str, str | None])r_   r¯   r   z!tuple[SubgraphStatus, str | None])r*   r)   r@  r‚  r   rí  )r@  r‚  r   r   )Pr0   Ú
__future__r   rB   rN  r£  Úcollections.abcr   r   r   r   Údataclassesr   r	   Útypingr
   r   r   r   r  r   Úlangchain_protocolr   r   Úlanggraph_sdk._async.httpr   Úlanggraph_sdk.schemar   r   Úlanggraph_sdk.stream.controllerr   Úlanggraph_sdk.stream.decodersr   r   r   r   r   r   r   r›  r   r   Úlanggraph_sdk.stream.transportr   r    r!   r"   r$   r6   r;   rM   r1   rT   rZ   Ú	frozensetra   rb   rd   rƒ   r¶   r¸   rÛ   r÷   r  r  r†  r;  r=  rA  rD  rF  rV  rW  rY  rÆ  r£  rñ  r\  r  r~   r+   r3   r4   Ú<module>rü     sÝ  ðòõ #ã Û Û ß MÓ Mß (ß 0Ó 0å Qß 5å 0ß BÝ 9÷÷ ñ ÷ R÷ó ô�yô ð ÷'ð 'ó ð'ð ÷@ð @ó ð@ò
€ˆyó 
ðØðàðð ôôBñ #,¨[¸(Ð,CÓ"DÐ ô@÷,
ñ 
÷8mLñ mLð` EH÷ ÷:ñ :÷z":ñ ":÷JjGñ jGôZ?ö
ð ÐHÑI€ô*ð
Ø
ðà&ôôXô÷Cñ C÷L72ñ 72÷t;ñ ;÷|gñ g÷T<,ñ <,÷~A&ñ A&÷H<:ñ <:÷~!ñ !÷.!:ñ !:÷HDNò DNr3   