ó
    ýÞ jV6  ã                  ó&  • S SK Jr  S SKrS SKrS SK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  S SKJr  S S	KJrJr  S S
KJr  S SKJr  \R4                  " \5      rSS jr " S S5      r " S S\\\4   5      r  " S S5      r! " S S5      r"g)é    )ÚannotationsN)ÚAsyncIteratorÚIteratorÚMapping)ÚTracebackType)ÚAnyÚcast)ÚRunnableConfig)ÚAsyncThreadStream)ÚSyncThreadStream)ÚLangGraphClientÚSyncLangGraphClient)ÚDataDecoder)ÚCommandc                ó¢   • [        U [        5      (       a9  U R                  (       d  U R                  (       a  [	        S5      eU R
                  $ U $ )ac  Translate a local `Command` into the v3 wire `input`, else passthrough.

The v3 server decides start-vs-resume from thread state (an interrupted
run or pending interrupts) and, on resume, wraps the whole `input` as
`{"resume": input}` itself. So a resume `Command` must surface its raw
`resume` value as the wire `input` (not the serialized dataclass, which
the server would double-wrap). The v3 `run.start` path has no `goto` /
`update` channel, so those are rejected.

`langgraph_sdk` is upstream of `langgraph`, so this `Command`-aware
marshalling lives here on the adapter (langgraph) side of the boundary.
z}RemoteGraph v3 streaming supports `Command(resume=...)` only; `goto` / `update` are not supported by the v3 `run.start` path.)Ú
isinstancer   ÚgotoÚupdateÚNotImplementedErrorÚresume)Úinputs    Ú]/var/www/html/gaurav/venv/lib/python3.13/site-packages/langgraph/pregel/_remote_run_stream.pyÚ_translate_command_inputr      sB   € ô �%œ×!Ñ!Ø�:�:˜ŸŸÜ%ðRóð ð �|‰|ÐØ€Ló    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)Ú_ChannelProjectioné,   u§  Decoded projection for a wire channel the SDK doesn't type natively.

Subscribes to `channel` and decodes each event's `params["data"]` through the
SDK's `DataDecoder` â€” the same decoder the SDK's own plain-payload projections
(`values` / `updates` / `checkpoints` / `tasks`) use, which yields the item
shape that local's `UpdatesTransformer` / `CheckpointsTransformer` /
`TasksTransformer` / `CustomTransformer` push, so iterating this matches the
corresponding local projection. Iterate with `for` against a sync stream and
`async for` against an async stream (matching the underlying SDK). Opening
the subscription requires the stream to be entered (`with` / `async with`).
c                ó   • Xl         X l        g ©N)Ú_sdkÚ_channel)ÚselfÚsdkÚchannels      r   Ú__init__Ú_ChannelProjection.__init__9   s   € ØŒ	Ø�r   c              #  óô   #   • [        U R                  5      n[        [        [           U R
                  R                  U R                  /5      5      nU H  nUR                  U5       S h  v•N   M     g  N	7fr   )r   r!   r	   r   r   r    Ú	subscribeÚfeed)r"   ÚdecoderÚeventsÚevents       r   Ú__iter__Ú_ChannelProjection.__iter__=   sW   é € ä˜dŸm™mÓ,ˆÜ”hœs‘m T§Y¡Y×%8Ñ%8¸$¿-¹-¸Ó%IÓJˆÛˆEØ—|‘| EÓ*×*Ò*ò Ù*ùs   ‚A*A8Á,A6Á-
A8c                ó"   • U R                  5       $ r   )Ú_aiter©r"   s    r   Ú	__aiter__Ú_ChannelProjection.__aiter__D   s   € Ø�{‰{‹}Ðr   c               ó  #   • [        U R                  5      n[        [        [           U R
                  R                  U R                  /5      5      nU  S h  v•N nUR                  U5       H  nU7v •  M
     M(   N#
 g 7fr   )r   r!   r	   r   r   r    r(   r)   )r"   r*   r+   r,   Úitems        r   r0   Ú_ChannelProjection._aiterG   sb   é € ä˜dŸm™mÓ,ˆÜ”m¤CÑ(¨$¯)©)×*=Ñ*=¸t¿}¹}¸oÓ*NÓOˆÙ!÷ 	�%ØŸ™ UÖ+�Ø•
ó ,ñ	™6ùs*   ‚ABÁA?ÁA=ÁA?Á BÁ=A?Á?B)r!   r    N)r#   ú$AsyncThreadStream | SyncThreadStreamr$   ÚstrÚreturnÚNone©r9   zIterator[Any]©r9   zAsyncIterator[Any])
Ú__name__Ú
__module__Ú__qualname__Ú__firstlineno__Ú__doc__r%   r-   r2   r0   Ú__static_attributes__© r   r   r   r   ,   s   † ñ
ô ô+ô÷r   r   c                  óR   • \ rS rSrSrSrSr\\-   rSS jrSS jr	SS jr
SS jrS	rg
)Ú_ProjectionRegistryéP   ui  Read-only name -> projection registry mirroring local `GraphRunStream.extensions`.

Resolution follows the langchain-protocol wire channels, and every entry
yields the same decoded item shape local does (`params.data`):

- `values` / `messages` / `tool_calls` / `subgraphs` resolve to the SDK's
  decoded typed projections. `tool_calls` is the `tools` channel â€” tool
  *execution* events, distinct from the tool-call *inputs* inside `messages`.
- `updates` / `checkpoints` / `tasks` / `custom` have no typed SDK
  projection, so they resolve to a `_ChannelProjection` that subscribes to
  the channel and yields `params.data` â€” matching the local transformer
  output for those channels.
- any other name is a specific custom-extension channel
  (`thread.extensions[name]`, i.e. `custom:<name>`).

`lifecycle` is intentionally absent: local derives a status payload from it
rather than yielding `params.data`, and the SDK consumes it as control-plane
(driving `output` / `interrupted`), so its shape can't be matched â€” it
remains reachable via the raw `events` iterator. `debug` is absent too: it
is not a v3 wire channel.
)ÚvaluesÚmessagesÚ
tool_callsÚ	subgraphs)ÚupdatesÚcheckpointsÚtasksÚcustomc                ó   • Xl         g r   ©r    )r"   r#   s     r   r%   Ú_ProjectionRegistry.__init__m   s   € Ø�	r   c                óÈ   • XR                   ;   a  [        U R                  U5      $ XR                  ;   a  [	        U R                  U5      $ U R                  R
                  U   $ r   )Ú_TYPEDÚgetattrr    Ú_DECODEDr   Ú
extensions)r"   Únames     r   Ú__getitem__Ú_ProjectionRegistry.__getitem__p   sM   € Ø—;‘;ÓÜ˜4Ÿ9™9 dÓ+Ð+Ø—=‘=Ó Ü% d§i¡i°Ó6Ð6Ø�y‰y×#Ñ# DÑ)Ð)r   c                ó,   • [        U R                  5      $ r   )ÚiterÚ_NATIVEr1   s    r   r-   Ú_ProjectionRegistry.__iter__w   s   € Ü�D—L‘LÓ!Ð!r   c                ó,   • [        U R                  5      $ r   )Úlenr\   r1   s    r   Ú__len__Ú_ProjectionRegistry.__len__z   s   € Ü�4—<‘<Ó Ð r   rP   N)r#   r7   r9   r:   )rW   r8   r9   r   )r9   zIterator[str])r9   Úint)r=   r>   r?   r@   rA   rS   rU   r\   r%   rX   r-   r`   rB   rC   r   r   rE   rE   P   s1   † ñð. ?€Fà<€HØ�xÑ€Gôô*ô"÷!r   rE   c                  ó  • \ rS rSrSr            SS jrSS jr        SS jr\SS j5       r	\SS j5       r
\SS j5       r\SS	 j5       r\SS
 j5       r\SS j5       r\SS j5       r\SS j5       rSS jrSS jrSS jrSrg)Ú_RemoteGraphRunStreamé~   z9Sync adapter: SyncThreadStream -> GraphRunStream surface.c               ón   • Xl         X l        [        U5      UUS.U l        S U l        SU l        S U l        g ©N)r   ÚconfigÚmetadataF)Ú_clientr    r   Ú_start_kwargsÚ_run_idÚ_closedÚ_events_iter)r"   Úsync_clientÚ
sdk_threadr   rh   ri   s         r   r%   Ú_RemoteGraphRunStream.__init__�   s>   € ð #ŒØŒ	ä-¨eÓ4ØØ ñ.
ˆÔð
 $(ˆŒØˆŒØ26ˆÕr   c                ó^  • U R                   (       a  [        S5      eU R                  R                  5          U R                  R                  R
                  " S0 U R                  D6nUS   U l        U $ ! [         a.    U R                  R                  " [        R                  " 5       6   e f = f)Nz$_RemoteGraphRunStream already closedÚrun_idrC   )rm   ÚRuntimeErrorr    Ú	__enter__ÚrunÚstartrk   ÚBaseExceptionÚ__exit__ÚsysÚexc_inforl   ©r"   Úresults     r   ru   Ú_RemoteGraphRunStream.__enter__•   sŠ   € Ø�<�<ÜÐEÓFÐFØ�	‰	×ÑÔð	Ø—Y‘Y—]‘]×(Ò(Ñ>¨4×+=Ñ+=Ñ>ˆFð ˜hÑ'ˆŒØˆøô	 ó 	Ø�I‰I×Ò¤§¢£Ñ/Øð	ús   ¸0A4 Á48B,c                ón   • U R                   (       a  g SU l         U R                  R                  XU5        g ©NT)rm   r    ry   ©r"   Úexc_typeÚexcÚtbs       r   ry   Ú_RemoteGraphRunStream.__exit__¡   s)   € ð �<�<ØØˆŒØ�	‰	×Ñ˜8¨"Õ-r   c                ó.   • U R                   R                  $ r   ©r    Úoutputr1   s    r   rˆ   Ú_RemoteGraphRunStream.output¬   s   € à�y‰y×ÑÐr   c                ó.   • U R                   R                  $ )aH  Whether the remote run is currently paused at an interrupt.

Reads the SDK's current value without blocking. This differs from
local `GraphRunStream.interrupted`, which drives the run to terminal
before returning the flag. Sync callers needing a wait-for-interrupt
pattern should switch to the async API and drain a projection.
©r    Úinterruptedr1   s    r   rŒ   Ú!_RemoteGraphRunStream.interrupted°   s   € ð �y‰y×$Ñ$Ð$r   c                ó@   • [        U R                  R                  5      $ )z?Current outstanding interrupt payloads (non-blocking snapshot).©Úlistr    Ú
interruptsr1   s    r   r‘   Ú _RemoteGraphRunStream.interrupts»   s   € ô �D—I‘I×(Ñ(Ó)Ð)r   c                ó.   • U R                   R                  $ ©z<Live state-snapshot projection (mirrors local `run.values`).©r    rG   r1   s    r   rG   Ú_RemoteGraphRunStream.valuesÀ   ó   € ð �y‰y×ÑÐr   c                ó.   • U R                   R                  $ ©z>Live message-stream projection (mirrors local `run.messages`).©r    rH   r1   s    r   rH   Ú_RemoteGraphRunStream.messagesÅ   ó   € ð �y‰y×!Ñ!Ð!r   c                ó.   • U R                   R                  $ ©z;Subgraph-handle projection (mirrors local `run.subgraphs`).©r    rJ   r1   s    r   rJ   Ú_RemoteGraphRunStream.subgraphsÊ   ó   € ð �y‰y×"Ñ"Ð"r   c                ó.   • U R                   R                  $ ©z³Tool-execution projection (the `tools` channel).

These are tool *execution* events (started / output / finished),
distinct from the tool-call *inputs* carried inside `messages`.
©r    rI   r1   s    r   rI   Ú _RemoteGraphRunStream.tool_callsÏ   ó   € ð �y‰y×#Ñ#Ð#r   c                ó,   • [        U R                  5      $ ©z=Name -> projection registry (mirrors local `run.extensions`).©rE   r    r1   s    r   rV   Ú _RemoteGraphRunStream.extensionsØ   ó   € ô # 4§9¡9Ó-Ð-r   c                óž  • U R                   (       a  g SU l         U R                  bD   U R                  R                  R	                  U R
                  R                  U R                  SS9   U R
                  R                  5         g ! [         a    [        R                  SSS9   N<f = f! [         a    [        R                  SSS9   g f = f©NTF)Úwaitzabort: runs.cancel failed)r{   zabort: sdk.close failed©rm   rl   rj   ÚrunsÚcancelr    Ú	thread_idÚ	ExceptionÚloggerÚdebugÚcloser1   s    r   ÚabortÚ_RemoteGraphRunStream.abortÝ   s¬   € Ø�<�<ØØˆŒØ�<‰<Ñ#ðIØ—‘×!Ñ!×(Ñ(¨¯©×)<Ñ)<¸d¿l¹lÐQVÐ(ÑWð	CØ�I‰I�O‰OÕøô ó IÜ—‘Ð8À4�ÓHðIûô ó 	CÜ�L‰LÐ2¸TˆLÓBð	Cús$   ¨AB Á,B+ ÂB(Â'B(Â+CÃCc                ó|   • U R                   c$  [        U R                  R                  5      U l         U R                   $ r   )rn   r[   r    r+   r1   s    r   r-   Ú_RemoteGraphRunStream.__iter__ë   s1   € Ø×ÑÑ$Ü $ T§Y¡Y×%5Ñ%5Ó 6ˆDÔØ× Ñ Ð r   c              '  óh   #   • U R                   R                  [        U5      5       S h  v•N   g  N7fr   )r    Úinterleave_projectionsr�   )r"   Únamess     r   Ú
interleaveÚ _RemoteGraphRunStream.interleaveð   s!   é € Ø—9‘9×3Ñ3´D¸³KÓ@×@Ó@ùs   ‚(2ª0«2)rj   rm   rn   rl   r    rk   N)ro   r   rp   r   r   r   rh   úRunnableConfig | Noneri   údict[str, Any] | Noner9   r:   )r9   rd   ©r‚   ztype[BaseException] | Nonerƒ   zBaseException | Noner„   zTracebackType | Noner9   r:   ©r9   r   ©r9   Úbool©r9   z	list[Any]©r9   zMapping[str, Any]©r9   r:   r;   )r½   r8   r9   zIterator[tuple[str, Any]])r=   r>   r?   r@   rA   r%   ru   ry   Úpropertyrˆ   rŒ   r‘   rG   rH   rJ   rI   rV   r·   r-   r¾   rB   rC   r   r   rd   rd   ~   s'  † ÙCð7ð )ð7ð %ð	7ð
 ð7ð &ð7ð (ð7ð 
ô7ô(
ð	.à,ð	.ð "ð	.ð !ð		.ð
 
ô	.ð ó ó ð ð ó%ó ð%ð ó*ó ð*ð ó ó ð ð ó"ó ð"ð ó#ó ð#ð ó$ó ð$ð ó.ó ð.ôCô!÷
Ar   rd   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S jr	SS jr
\SS	 j5       r\SS
 j5       r\SS j5       r\SS j5       r\SS j5       rSS jrSS jrSrg)Ú_AsyncRemoteGraphRunStreaméô   z@Async adapter: AsyncThreadStream -> AsyncGraphRunStream surface.c               ón   • Xl         X l        [        U5      UUS.U l        S U l        SU l        S U l        g rg   )rj   r    r   rk   rl   rm   Ú_events_aiter)r"   Úclientrp   r   rh   ri   s         r   r%   Ú#_AsyncRemoteGraphRunStream.__init__÷   s>   € ð ŒØŒ	ä-¨eÓ4ØØ ñ.
ˆÔð
 $(ˆŒØˆŒØ8<ˆÕr   c              ƒ  ó   #   • U R                   (       a  [        S5      eU R                  R                  5       I S h  v•N    U R                  R                  R
                  " S0 U R                  D6I S h  v•N nUS   U l        U $  NI N! [         a7    U R                  R                  " [        R                  " 5       6 I S h  v•N    e f = f7f)Nz)_AsyncRemoteGraphRunStream already closedrs   rC   )rm   rt   r    Ú
__aenter__rv   rw   rk   rx   Ú	__aexit__rz   r{   rl   r|   s     r   rÒ   Ú%_AsyncRemoteGraphRunStream.__aenter__  s¡   é € Ø�<�<ÜÐJÓKÐKØ�i‰i×"Ñ"Ó$×$Ð$ð	ØŸ9™9Ÿ=™=×.Ò.ÑD°×1CÑ1CÑD×DˆFð ˜hÑ'ˆŒØˆñ 	%áDøÜó 	Ø—)‘)×%Ò%¤s§|¢|£~Ð6×6Ñ6Øð	üsE   ‚:C¼B½CÁ3B
 Á5BÁ6B
 Á:CÂB
 Â
9CÃCÃCÃCc              ƒ  óŠ   #   • U R                   (       a  g SU l         U R                  R                  XU5      I S h  v•N   g  N7fr€   )rm   r    rÓ   r�   s       r   rÓ   Ú$_AsyncRemoteGraphRunStream.__aexit__  s2   é € ð �<�<ØØˆŒØ�i‰i×!Ñ! (°Ó4×4Ó4ùs   ‚9A»A¼Ac              ƒ  óJ   #   • U R                   R                  I Sh  v•N $  N7f)a  Drive the remote run to completion and return the final state.

Awaits the SDK's terminal-state awaitable, matching local
`AsyncGraphRunStream.output()` (a method, not a property, so
`run.output` without `await` fails at type-check time rather than
silently yielding a coroutine).
Nr‡   r1   s    r   rˆ   Ú!_AsyncRemoteGraphRunStream.output"  s   é € ð —Y‘Y×%Ñ%×%Ð%Ñ%ùs   ‚#œ!�#c              ƒ  ó6   #   • U R                   R                  $ 7f)aœ  Whether the remote run is currently paused at an interrupt.

Reads the SDK's current value without blocking. This differs from
local `AsyncGraphRunStream.interrupted()`, which drives the run to
terminal before returning the flag. Callers that need a
wait-for-interrupt pattern should drain a projection (e.g.,
`async for snap in stream._sdk.values`) until the SDK's paused
sentinel fires, then call this method.
r‹   r1   s    r   rŒ   Ú&_AsyncRemoteGraphRunStream.interrupted,  s   é € ð �y‰y×$Ñ$Ð$ùs   ‚c              ƒ  óH   #   • [        U R                  R                  5      $ 7f)z—Current outstanding interrupt payloads.

Non-blocking; reads the SDK's current snapshot. See `interrupted`
for the divergence from local v3 semantics.
r�   r1   s    r   r‘   Ú%_AsyncRemoteGraphRunStream.interrupts8  s   é € ô �D—I‘I×(Ñ(Ó)Ð)ùs   ‚ "c                ó.   • U R                   R                  $ r”   r•   r1   s    r   rG   Ú!_AsyncRemoteGraphRunStream.values@  r—   r   c                ó.   • U R                   R                  $ r™   rš   r1   s    r   rH   Ú#_AsyncRemoteGraphRunStream.messagesE  rœ   r   c                ó.   • U R                   R                  $ rž   rŸ   r1   s    r   rJ   Ú$_AsyncRemoteGraphRunStream.subgraphsJ  r¡   r   c                ó.   • U R                   R                  $ r£   r¤   r1   s    r   rI   Ú%_AsyncRemoteGraphRunStream.tool_callsO  r¦   r   c                ó,   • [        U R                  5      $ r¨   r©   r1   s    r   rV   Ú%_AsyncRemoteGraphRunStream.extensionsX  r«   r   c              ƒ  óÎ  #   • U R                   (       a  g SU l         U R                  bL   U R                  R                  R	                  U R
                  R                  U R                  SS9I S h  v•N    U R
                  R                  5       I S h  v•N   g  N(! [         a    [        R                  SSS9   NFf = f N+! [         a    [        R                  SSS9   g f = f7fr­   r¯   r1   s    r   r·   Ú _AsyncRemoteGraphRunStream.abort]  sË   é € Ø�<�<ØØˆŒØ�<‰<Ñ#ðIØ—l‘l×'Ñ'×.Ñ.Ø—I‘I×'Ñ'¨¯©¸Eð /ð ÷ ð ð
	CØ—)‘)—/‘/Ó#×#Ñ#ñøô ó IÜ—‘Ð8À4�ÓHðIúñ $øÜó 	CÜ�L‰LÐ2¸TˆLÓBð	Cüsk   ‚'C%ªAB Á0BÁ1B Á6C ÂB?ÂC ÂC%ÂB ÂB<Â9C%Â;B<Â<C%Â?C ÃC"ÃC%Ã!C"Ã"C%c                ó†   • U R                   c)  U R                  R                  R                  5       U l         U R                   $ r   )rÎ   r    r+   r2   r1   s    r   r2   Ú$_AsyncRemoteGraphRunStream.__aiter__m  s5   € Ø×ÑÑ%Ø!%§¡×!1Ñ!1×!;Ñ!;Ó!=ˆDÔØ×!Ñ!Ð!r   )rj   rm   rÎ   rl   r    rk   N)rÏ   r   rp   r   r   r   rh   rÀ   ri   rÁ   r9   r:   )r9   rË   rÂ   rÃ   rÄ   rÆ   rÇ   rÈ   r<   )r=   r>   r?   r@   rA   r%   rÒ   rÓ   rˆ   rŒ   r‘   rÉ   rG   rH   rJ   rI   rV   r·   r2   rB   rC   r   r   rË   rË   ô   sô   † ÙJð=ð  ð=ð &ð	=ð
 ð=ð &ð=ð (ð=ð 
ô=ô(
ð	5à,ð	5ð "ð	5ð !ð		5ð
 
ô	5ô&ô
%ô*ð ó ó ð ð ó"ó ð"ð ó#ó ð#ð ó$ó ð$ð ó.ó ð.ôC÷ "r   rË   )r   r   r9   r   )#Ú
__future__r   Úloggingrz   Úcollections.abcr   r   r   Útypesr   Útypingr   r	   Úlangchain_core.runnablesr
   Úlanggraph_sdk._async.streamr   Úlanggraph_sdk._sync.streamr   Úlanggraph_sdk.clientr   r   Úlanggraph_sdk.stream.decodersr   Úlanggraph.typesr   Ú	getLoggerr=   r´   r   r   r8   rE   rd   rË   rC   r   r   Ú<module>r÷      s}   ðå "ã Û 
ß <Ñ <Ý ß å 3Ý 9Ý 7ß EÝ 5å #à	×	Ò	˜8Ó	$€ô÷.!ñ !ôH+!˜' # s (Ñ+ô +!÷\sAñ sA÷l|"ò |"r   