ó
    ýÞ j   ã                  ó¤   • S r SSKJr  SSKrSSKrSSKJrJr  SSKJ	r	J
r
  SSKrSSKrSSKJr  SSKJr  SSKJrJr  SS	KJrJr  \r " S
 S5      rg)aµ  HTTP/SSE transport for the v3 thread-centric protocol.

Direct port of `libs/sdk/src/client/stream/transport/http.ts`.

`ProtocolSseTransport` is bound to a single `thread_id` at construction. Commands
go to `POST /threads/{thread_id}/commands` (JSON in, JSON out). Each
`open_event_stream(params)` opens an independent filtered SSE connection at
`POST /threads/{thread_id}/stream/events` with the `SubscribeParams` in the
request body.
é    )ÚannotationsN)ÚAsyncIteratorÚMapping)ÚAnyÚcast)ÚEvent)Ú_quote_path_param)ÚBytesLineDecoderÚ
SSEDecoder)ÚEventStreamHandleÚbuild_event_stream_bodyc                  óh   • \ rS rSrSrSSSSS.             SS jjrSS jrSS jrSS	 jrS
r	g)ÚProtocolSseTransporté!   záv3 protocol transport bound to a single `thread_id`.

Commands go to `POST /threads/{thread_id}/commands` (JSON in, JSON out).
`open_event_stream` opens filtered SSE streams against
`POST /threads/{thread_id}/stream/events`.
Ni   )Úcommands_pathÚstream_pathÚheadersÚmax_queue_sizec               óú   • Xl         X l        U=(       d    S[        U5       S3U l        U=(       d    S[        U5       S3U l        [        U=(       d    0 5      U l        X`l        SU l        [        5       U l
        g )Nz	/threads/z	/commandsz/stream/eventsF)Ú_clientÚ	thread_idr	   Ú_commands_urlÚ_stream_urlÚdictÚ_default_headersÚ_max_queue_sizeÚ_closedÚsetÚ_event_streams)ÚselfÚclientr   r   r   r   r   s          Ú]/var/www/html/gaurav/venv/lib/python3.13/site-packages/langgraph_sdk/stream/transport/http.pyÚ__init__ÚProtocolSseTransport.__init__)   sz   € ð ŒØ"Œà×P˜yÔ):¸9Ó)EÐ(FÀiÐPð 	Ôð ×S˜YÔ'8¸Ó'CÐ&DÀNÐSð 	Ôô 15°W·]ÀÓ0CˆÔØ-ÔØˆŒÜ7:³uˆÕó    c              ƒ  óH  #   • U R                   (       a  [        S5      e0 U R                  ESS0EnU R                  R	                  U R
                  [        R                  " U5      US9I Sh  v•N nUR                  5         UR                  S;   a  gUR                  (       d  [        S5      e [        R                  " UR                  5      n[        U[        5      (       a  SU;  a  [        S5      eU$  NŠ! [        R                   a  n[        S5      UeSnAff = f7f)	a	  POST a command. Returns the response JSON, or `None` for 202/204.

Raises:
    httpx.HTTPStatusError: server returned >= 400.
    RuntimeError: the transport has been closed via `close()`.
    RuntimeError: server returned a response missing the protocol envelope.
úProtocol transport is closed.úcontent-typeúapplication/json©Úcontentr   N)éÊ   éÌ   z1Protocol command did not return a valid response.Úid)r   ÚRuntimeErrorr   r   Úpostr   ÚorjsonÚdumpsÚraise_for_statusÚstatus_coder+   ÚloadsÚJSONDecodeErrorÚ
isinstancer   )r    ÚcommandÚmerged_headersÚresponseÚpayloadÚerrs         r"   Úsend_commandÚ!ProtocolSseTransport.send_command@   s  é € ð �<�<ÜÐ>Ó?Ð?àV˜D×1Ñ1ÐV°>ÐCUÑVˆØŸ™×*Ñ*Ø×ÑÜ—L’L Ó)Ø"ð +ð 
÷ 
ˆð
 	×!Ñ!Ô#Ø×Ñ :Ó-ØØ××ÜÐRÓSÐSð	Ü—l’l 8×#3Ñ#3Ó4ˆGô
 ˜'¤4×(Ñ(¨D¸Ó,?ÜÐRÓSÐSØˆñ%
øô ×%Ñ%ó 	ÜØCóàðûð	üs7   ‚A+D"Á-C8Á.AD"Â0 C: Ã)D"Ã:DÄDÄDÄD"c                ó0  ^ ^^^^^	^
• T R                   (       a  [        S5      e[        R                  " 5       nUR	                  5       m	UR	                  5       m[        R
                  " T R                  S9m[        R                  " 5       mSUUUUU	U 4S jjn[        R                  " U" 5       5      m
T R                  R                  T
5        T
R                  T R                  R                  5        SUU4S jjnSUUU
4S jjn[        U" 5       T	TUS9$ )	aé  Open an independent filtered SSE event stream.

Posts `params` as a SubscribeParams body to `/threads/{thread_id}/stream/events`.
Returns an `EventStreamHandle` whose `events` async iterator yields typed
`Event` dicts as the server emits them. `handle.ready` resolves on a 2xx
response (rejects on HTTP error or transport failure before headers).

Reconnect: pass `params["since"]` to filter outbound seqs server-side. The
cursor goes in the request body, not as a `Last-Event-ID` header.
r'   )Úmaxsizec            	   “  óŒ  >#   •  0 TR                   ESSSS.En TR                  R                  STR                  [        R
                  " [        T
5      5      U S9 IS h  v•N nUR                  5         TR                  5       (       d  TR                  S 5        [        5       n[        5       nUR                  5         S h  v•N nTR                  5       (       a    O‡UR                  U5       Hp  nUR                  [        U5      5      nUc  M"  [!        UR"                  [$        5      (       d  MC  TR'                  [)        SUR"                  5      5      I S h  v•N   Mr     M§  TR                  5       (       dä  UR+                  5        Hp  nUR                  [        U5      5      nUc  M"  [!        UR"                  [$        5      (       d  MC  TR'                  [)        SUR"                  5      5      I S h  v•N   Mr     UR                  S5      nUbL  [!        UR"                  [$        5      (       a-  TR'                  [)        SUR"                  5      5      I S h  v•N   S S S 5      IS h  v•N   T	R                  5       (       d  T	R                  S 5        TR'                  S 5      I S h  v•N   g  GNO GNï GNZ
 GNU NÄ Nb NT! , IS h  v•N  (       d  f       Ni= f! [,        R.                   a,  nT	R                  5       (       d  T	R                  U5        e S nAf[0         aW  nTR                  5       (       d  TR3                  U5        T	R                  5       (       d  T	R                  U5         S nAGNS nAff = f NÓ! T	R                  5       (       d  T	R                  S 5        TR'                  S 5      I S h  v•N    f = f7f)	Nr)   ztext/event-streamzno-store)r(   Úacceptzcache-controlÚPOSTr*   r   r%   )r   r   Ústreamr   r1   r2   r   r3   ÚdoneÚ
set_resultr
   r   Úaiter_bytesÚis_setÚdecodeÚbytesr7   Údatar   Úputr   ÚflushÚasyncioÚCancelledErrorÚBaseExceptionÚset_exception)Ússe_headersr:   Úline_decoderÚsse_decoderÚchunkÚlineÚpartr<   Úcancel_eventrE   ÚparamsÚqueueÚreadyr    s           €€€€€€r"   ÚpumpÚ4ProtocolSseTransport.open_event_stream.<locals>.pumpt   sã  øé € ð1&ðØ×+Ñ+ðà$6Ø1Ø%/ò	�ð  Ÿ<™<×.Ñ.ØØ×$Ñ$Ü"ŸLšLÔ)@ÀÓ)HÓIØ'ð	 /÷ ñ ð
 Ø×-Ñ-Ô/Ø Ÿ:™:Ÿ<™<Ø×(Ñ(¨Ô.Ü#3Ó#5�LÜ",£,�KØ'/×';Ñ';Ô'=÷ J˜eØ'×.Ñ.×0Ñ0Ù!Ø$0×$7Ñ$7¸Ö$>˜DØ#.×#5Ñ#5´e¸D³kÓ#B˜DØ#™|Ù (Ü)¨$¯)©)´T×:Ó:Ø&+§i¡i´°W¸d¿i¹iÓ0HÓ&I× IÒ Ió %?ð (×.Ñ.×0Ñ0Ø$0×$6Ñ$6Ö$8˜DØ#.×#5Ñ#5´e¸D³kÓ#B˜DØ#Ó/´J¸t¿y¹yÌ$×4OÓ4OØ&+§i¡i´°W¸d¿i¹iÓ0HÓ&I× IÒ Iñ %9ð  +×1Ñ1°#Ó6˜ØÑ+´
¸4¿9¹9Äd×0KÑ0KØ"'§)¡)¬D°¸$¿)¹)Ó,DÓ"E×EÐE÷9÷ ðN —y‘y—{‘{Ø—O‘O DÔ)Ø—i‘i “o×%Ñ%òSòJò !Jò (>ñ !Jñ F÷9÷ ÷ ð ûô: ×)Ñ)ó Ø—y‘y—{‘{Ø—O‘O CÔ(ØûÜ ó )Ø—z‘z—|‘|Ø×'Ñ'¨Ô,Ø—y‘y—{‘{Ø—O‘O CÔ(ÿùð	)úñ &øð —y‘y—{‘{Ø—O‘O DÔ)Ø—i‘i “o×%Ò%üs*  ƒO…AK Á J/Á!K Á$AKÂ?J8ÃJ2ÃJ8ÃA(KÄ3(KÅJ5
ÅAKÆ-KÇ(KÇ6J;Ç7A#KÉJ=ÉKÉK É*J?É+K É/:OÊ)M=Ê*OÊ/K Ê2J8Ê5KÊ8KÊ;KÊ=KÊ?K ËKËK
ËKËK ËM? ËK ËM:Ë/'LÌM:Ì#AM5Í/M? Í5M:Í:M? Í=OÍ?;OÎ:N=Î;OÏOc                ó‚   >#   •  TR                  5       I S h  v•N n U b  TR                  5       (       a  g U 7v •  M8   N$7f©N)ÚgetrH   )ÚitemrX   rZ   s    €€r"   ÚaiterÚ5ProtocolSseTransport.open_event_stream.<locals>.aiter¬   s:   øé € ØØ"ŸY™Y›[×(�Ø‘< <×#6Ñ#6×#8Ñ#8ØØ“
ñ	 Ù(ùs   ƒ?˜=™%?c               “  ó  >#   • T R                  5         TR                  S 5        TR                  5         [        R                  " [
        R                  [        5         TI S h  v•N   S S S 5        g  N! , (       d  f       g = f7fr_   )r   Ú
put_nowaitÚcancelÚ
contextlibÚsuppressrN   rO   Ú	Exception)rX   rZ   Útasks   €€€r"   ÚcloseÚ5ProtocolSseTransport.open_event_stream.<locals>.close³   s[   øé € Ø×ÑÔà×Ñ˜TÔ"Ø�K‰KŒMÜ×$Ò$¤W×%;Ñ%;¼YÕGØ—
�
÷ HÐGÙ÷ HÕGüs0   ƒABÁA4Á$A2Á%A4Á)	BÁ2A4Á4
BÁ>B)Úeventsr[   rE   rk   ©ÚreturnÚNone)ro   zAsyncIterator[Event])r   r/   rN   Úget_running_loopÚcreate_futureÚQueuer   r   Úcreate_taskr   ÚaddÚadd_done_callbackÚdiscardr   )r    rY   Úloopr\   rb   rk   rX   rE   rZ   r[   rj   s   ``    @@@@@r"   Úopen_event_streamÚ&ProtocolSseTransport.open_event_stream`   s×   þ€ ð �<�<ÜÐ>Ó?Ð?ä×'Ò'Ó)ˆØ&*×&8Ñ&8Ó&:ˆØ59×5GÑ5GÓ5IˆÜ-4¯]ª]À4×CWÑCWÑ-XˆÜ—}’}“ˆ÷2	&ô 2	&ôh ×"Ò"¡4£6Ó*ˆØ×Ñ×Ñ Ô%Ø×Ñ˜t×2Ñ2×:Ñ:Ô;÷	ð 	÷	ñ 	ô !©«°uÀ4ÈuÑUÐUr%   c              ƒ  óp  #   • U R                   (       a  gSU l         [        U R                  5      nU H  nUR                  5         M     U(       aQ  [        R
                  " [        [        R                  5         [        R                  " USS06I Sh  v•N   SSS5        gg N! , (       d  f       g= f7f)zHCancel any open event streams and mark the transport closed. Idempotent.NTÚreturn_exceptions)
r   Úlistr   rf   rg   rh   ri   rN   rO   Úgather)r    Útasksrj   s      r"   rk   ÚProtocolSseTransport.close½   sƒ   é € à�<�<ØØˆŒÜ�T×(Ñ(Ó)ˆÛˆDØ�K‰KŽMñ æÜ×$Ò$¤Y´×0FÑ0FÕGÜ—n’n eÐD¸tÑD×DÐD÷ HÐGð áD÷ HÕGüs0   ‚A8B6Á:B%ÂB#ÂB%Â
B6Â#B%Â%
B3Â/B6)r   r   r   r   r   r   r   r   )r!   zhttpx.AsyncClientr   Ústrr   ú
str | Noner   r‚   r   zMapping[str, str] | Noner   Úintro   rp   )r8   údict[str, Any]ro   zdict[str, Any] | None)rY   r„   ro   r   rn   )
Ú__name__Ú
__module__Ú__qualname__Ú__firstlineno__Ú__doc__r#   r=   ry   rk   Ú__static_attributes__© r%   r"   r   r   !   st   † ñð %)Ø"&Ø,0Ø"ñ=ð "ð=ð ð	=ð
 "ð=ð  ð=ð *ð=ð ð=ð 
õ=ô.ô@[V÷z
Er%   r   )r‰   Ú
__future__r   rN   rg   Úcollections.abcr   r   Útypingr   r   Úhttpxr1   Úlangchain_protocolr   Úlanggraph_sdk._shared.utilitiesr	   Úlanggraph_sdk.sser
   r   Ú#langgraph_sdk.stream.transport.baser   r   Ú_build_event_stream_bodyr   r‹   r%   r"   Ú<module>r•      sE   ðñ	õ #ã Û ß 2ß ã Û Ý $å =ß :÷ð
 3Ð ÷fEò fEr%   