ó
    ýÞ jH=  ã                  ó  • S r SSKJr  SSKrSSKrSSKrSSKrSSKJr  SSK	J
r
JrJr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Jr   " S
 S5      r\R4                  " \5      r\ " S S5      5       rSS.SS jjr " S S5      rg)aÉ  Stream controller: subscription registry and fan-out for AsyncThreadStream.

`StreamController` manages the set of active subscriptions against one shared
SSE connection, routing events from the shared stream to per-subscription
queues.  It is the centralised place for:

- subscription registration / teardown
- shared-stream lifecycle (open, rotate, close)
- dedup of replayed events across rotations
- fan-out from the shared stream to subscriber queues
é    )ÚannotationsN)ÚOrderedDict)ÚAsyncGeneratorÚAsyncIteratorÚ	AwaitableÚCallable)Ú	dataclassÚfield)ÚAny)ÚEventÚSubscribeParams)ÚAsyncProtocolTransportÚEventStreamHandlec                  óD   • \ rS rSrSrSrS	S
S jjrSS jrSS jrS r	Sr
g)Ú_SeenEventIdsé!   z)LRU set of event ids with bounded memory.)Ú_dataÚ_maxsizec                ó.   • [        5       U l        Xl        g ©N)r   r   r   )ÚselfÚmaxsizes     ÚY/var/www/html/gaurav/venv/lib/python3.13/site-packages/langgraph_sdk/stream/controller.pyÚ__init__Ú_SeenEventIds.__init__&   s   € Ü-8«]ˆŒ
Ø�ó    c                óò   • XR                   ;   a  U R                   R                  U5        g S U R                   U'   [        U R                   5      U R                  :”  a  U R                   R	                  SS9  g g )NF)Úlast)r   Úmove_to_endÚlenr   Úpopitem©r   Úevent_ids     r   ÚaddÚ_SeenEventIds.add*   s]   € Ø—z‘zÓ!Ø�J‰J×"Ñ" 8Ô,ØØ#ˆ�
‰
�8ÑÜˆt�z‰z‹?˜TŸ]™]Ó*Ø�J‰J×Ñ EÐÒ*ð +r   c                ó   • XR                   ;   $ r   )r   r"   s     r   Ú__contains__Ú_SeenEventIds.__contains__2   s   € ØŸ:™:Ñ%Ð%r   c                ó,   • [        U R                  5      $ r   )Úiterr   ©r   s    r   Ú__iter__Ú_SeenEventIds.__iter__5   s   € Ü�D—J‘JÓÐr   N)é'  )r   ÚintÚreturnÚNone)r#   Ústrr0   r1   )r#   Úobjectr0   Úbool)Ú__name__Ú
__module__Ú__qualname__Ú__firstlineno__Ú__doc__Ú	__slots__r   r$   r'   r,   Ú__static_attributes__© r   r   r   r   !   s   † Ù3à%€Iö ô+ô&õ r   r   c                  óX   • \ rS rSr% SrS\S'   S\S'   \" \R                  S9r	S\S	'   S
r
g)Ú_Subscriptioné@   zDInternal record for one active subscription on a `StreamController`.r/   Úidr   Úparams)Údefault_factoryzasyncio.QueueÚqueuer<   N)r5   r6   r7   r8   r9   Ú__annotations__r
   ÚasyncioÚQueuerC   r;   r<   r   r   r>   r>   @   s#   ‡ áNàƒGØÓÙ °·±Ñ?€Eˆ=Ö?r   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)rE   ÚsleepÚclose)ÚhandlerG   s     r   Ú_close_afterrL   P   s3   é € ö Ü�mŠm˜EÓ"×"Ð"Ø
�,‰,‹.×Ññ 	#Ùùs   ‚!A£A¤A»A¼AÁAc                  ó4  • \ rS rSrSrSSSSSSS	.               SS
 jjrSSS.       SS jjrS S jrS!S jrS"S jr	\r
\	r    S#S jrS S jr\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)S jrS*S jrS)S jrS+S jrS,S jrSrg)-ÚStreamControlleré`   aÓ  Manages subscriptions and fan-out against one shared SSE connection.

Responsibilities:
  - subscription registry (register / unregister)
  - shared-stream lifecycle (open on first subscribe, rotate on filter widen)
  - dedup of replayed events via a bounded LRU `_SeenEventIds`
  - fan-out from the shared stream to per-subscription queues

Args:
    transport: the `AsyncProtocolTransport` bound to this thread session.
    run_start_gate: zero-argument async callable that resolves once the
        current `run.start` has committed server-side (no-op when no
        run is in flight).
    max_queue_size: per-subscription queue bound (default 1024).
    seen_event_ids_max: LRU cap for the dedup set (default 10_000).
Ni   r.   é   gš™™™™™¹?g       @)Úrun_start_gateÚmax_queue_sizeÚseen_event_ids_maxÚmax_reconnect_attemptsÚreconnect_backoff_baseÚreconnect_backoff_capc               óÜ   • Xl         X0l        [        US9U l        SU l        0 U l        S U l        S U l        S U l        [        5       U l
        SU l        S U l        XPl        X`l        Xpl        g )N©r   é   F)Ú
_transportÚ_max_queue_sizer   Ú_seen_event_idsÚ_next_subscription_idÚ_subscriptionsÚ_shared_streamÚ_shared_stream_filterÚ_fanout_taskÚsetÚ_rotation_close_tasksÚ_closedÚ_cursorÚ_max_reconnect_attemptsÚ_reconnect_backoff_baseÚ_reconnect_backoff_cap)r   Ú	transportrQ   rR   rS   rT   rU   rV   s           r   r   ÚStreamController.__init__r   so   € ð $ŒØ-ÔÜ,Ð5GÑHˆÔØ%&ˆÔ"Ø8:ˆÔØ8<ˆÔØ<@ˆÔ"Ø7;ˆÔÜ>A»eˆÔ"ØˆŒØ#'ˆŒØ'=Ô$Ø'=Ô$Ø&;Õ#r   )Ú
namespacesÚdepthc               óZ   • 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.
Úchannelsrk   rl   )ÚlistÚ_subscription_iter)r   rn   rk   rl   rA   s        r   Ú	subscribeÚStreamController.subscribe�   s>   € ð $.¬t°H«~Ð">ˆØÑ!Ø#-�<Ñ ØÑØ#�7‰OØ×&Ñ& vÓ.Ð.r   c              ƒ  ó  #   • U R                   (       a  gSU l         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                  (       a)  [        R                  " U R                  SS06I Sh  v•N   gg Nv! , (       d  f       Nz= f NR N7f)z?Tear down the controller, awaiting any pending rotation closes.NTÚreturn_exceptions)rd   ra   ÚcancelÚ
contextlibÚsuppressÚ	ExceptionrE   ÚCancelledErrorr_   rJ   rc   Úgatherr+   s    r   rJ   ÚStreamController.close¤   sÆ   é € à�<�<ØØˆŒØ×ÑÑ(Ø×Ñ×$Ñ$Ô&Ü×$Ò$¤Y´×0FÑ0FÕGØ×'Ñ'×'Ð'÷ Hà×ÑÑ*Ø×%Ñ%×+Ñ+Ó-×-Ð-Ø×%×%Ü—.’. $×"<Ñ"<ÐUÐPTÑU×UÑUð &ñ (÷ HÕGúñ .áUùsN   ‚A*D
Á,C5Á<C3Á=C5Â2D
Â3DÂ48D
Ã,DÃ-D
Ã3C5Ã5
DÃ?D
ÄD
c                óÂ   • [        U R                  U[        R                  " U R                  S9S9nU =R                  S-  sl        X R
                  UR                  '   U$ )zDAllocate a subscription id, create a bounded queue, add to registry.rX   )r@   rA   rC   rY   )r>   r]   rE   rF   r[   r^   r@   )r   rA   Úsubs      r   Ú_register_subscriptionÚ'StreamController._register_subscription¶   sT   € äØ×)Ñ)ØÜ—-’-¨×(<Ñ(<Ñ=ñ
ˆð
 	×"Ò" aÑ'Õ"Ø&)×Ñ˜CŸF™FÑ#Øˆ
r   c                ó<   • U R                   R                  US5        g)zARemove a subscription from the registry. No-op if already absent.N)r^   Úpop)r   Úsubscription_ids     r   Ú_unregister_subscriptionÚ)StreamController._unregister_subscriptionÁ   s   € à×Ñ×Ñ °Õ6r   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r   )r~   rd   rƒ   r@   Ú_reconcile_streamÚ_ensure_fanout_runningrC   Úget)r   rA   r}   Úitems       r   rp   Ú#StreamController._subscription_iterÉ   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 r   )ra   ÚdonerE   Úcreate_taskÚ_fanoutr+   s    r   r‡   Ú'StreamController._ensure_fanout_runningÞ   sA   € Ø×ÑÑ$¨×(9Ñ(9×(>Ñ(>×(@Ñ(@Ü '× 3Ò 3°D·L±L³NÓ CˆDÕð )Ar   c              ƒ  ó@  #   • SSK Jn  U R                  (       Gdz  U R                  nUc  g U R	                  UR
                  5        Sh  v•N nU R                  (       a    O`[        U R                  R                  5       5       H7  nU" X4R                  5      (       d  M  UR                  R                  U5        M9     M|  U R                  UL a¯  UR                  I Sh  v•N nUb—  [!        U["        R$                  5      (       dx  U R                  (       dg  [&        R(                  " [        5         U R                  R+                  5       I Sh  v•N   SSS5        U R-                  5       I Sh  v•N nU(       a  GMw  OU R                  (       d  GMz  U R                  R                  5        H  nUR                  R                  S5        M      g GN‡
 GN! [         a!  n[        R                  SU5         SnAGN;SnAff = f G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 `_max_reconnect_attempts` times before
giving up and closing subscriber queues.
r   )Úmatches_subscriptionNztransport drop in fanout: %r)Ú!langgraph_sdk.stream.subscriptionr‘   rd   r_   Ú_dedup_iterÚeventsro   r^   ÚvaluesrA   rC   Ú
put_nowaitrx   Ú_loggerÚdebugrŒ   Ú
isinstancerE   ry   rv   rw   rJ   Ú_reconnect_shared_stream)r   r‘   ÚsharedÚeventr}   Údrop_errÚerrÚreconnecteds           r   rŽ   ÚStreamController._fanoutå   sŒ  é € õ 	Kà—,—,�,Ø×(Ñ(ˆFØ‰~ØðHØ#'×#3Ñ#3°F·M±MÔ#B÷ 8˜%Ø—|—|ÙÜ# D×$7Ñ$7×$>Ñ$>Ó$@ÖA˜Ù/°·z±z×BÓBØŸI™I×0Ñ0°Ö7ó  Bð ×"Ñ" fÒ,Ø"ŸK™K×'�à‘OÜ& s¬G×,BÑ,B×CÑCØ ŸLŸLä#×,Ò,¬YÕ7Ø"×1Ñ1×7Ñ7Ó9×9Ð9÷ 8à(,×(EÑ(EÓ(G×"G�KÞ"Ú Øð5 —,—,’,ð< ×&Ñ&×-Ñ-Ö/ˆCØ�I‰I× Ñ  Ö&ò 0ò38Ò#Bøô ó HÜ—‘Ð<¸h×GÒGûðHúò (ñ :÷ 8Õ7úá"Gùs³   ‚)H¬G ÁGÁGÁGÁG Á"HÁ#;G Â"!G ÃHÃ!HÃ"AHÄ3HÅH	ÅHÅHÅ1HÅ2 HÆ=HÇGÇG ÇHÇ
HÇ"G>Ç8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)rg   rh   ÚminÚrandomÚuniformrE   rI   )r   ÚattemptÚbaseÚcaprG   Újitters         r   Ú_reconnect_sleepÚ!StreamController._reconnect_sleep  sV   é € à×+Ñ+ˆØ×)Ñ)ˆÜ�C  G¡Ñ,Ó-ˆÜ—’  5¨4¡<Ó0ˆÜ�mŠm˜E™NÓ+×+Ó+ùs   ‚AA)Á!A'Á"A)c              ƒ  ó–  #   • U R                   nUc  g[        U R                  5       H[  nU R                  (       a    g U R                  R                  U R                  U5      5      nUR                  I Sh  v•N   X0l          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 transport drop.

Returns True if a new stream was successfully opened, False if all
reconnect attempts were exhausted or the controller was closed.
NFT)r`   Úrangerf   rd   rZ   Úopen_event_streamÚ_filter_with_sinceÚreadyrE   ry   rx   rª   r_   )r   Úbase_filterr¦   Ú
new_streams       r   rš   Ú)StreamController._reconnect_shared_stream  s½   é € ð ×0Ñ0ˆØÑØÜ˜T×9Ñ9Ö:ˆGØ�|�|Ùð	Ø!Ÿ_™_×>Ñ>Ø×+Ñ+¨KÓ8ó�
ð !×&Ñ&×&Ð&ð #-ÔÙñ ;ð ñ 'øÜ×)Ñ)ó ØÜó Ø×+Ñ+¨GÓ4×4Ñ4ÚðüsF   ‚<C	¿9B
Á8BÁ9B
Á=C	ÂB
Â
2CÂ<B?Â=CÃC	ÃCÃC	c              ƒ  ó$  #   • SSK Jn  U R                  b/  U R                  b"  U" U R                  [	        U5      5      (       a  gU R                  US9nU R                  R                  U R                  U5      5      nU R                  nX@l        X0l        UR                  I Sh  v•N   Ub`  [        R                  " [        U5      5      nU R                  R                  U5        UR                  U R                  R                   5        gg Nh7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 `_seen_event_ids`
dedupes the overlap. Awaits `new_stream.ready` so the HTTP connection
is established before returning.
r   )Úfilter_coversN)Úextra)r’   rµ   r_   r`   ÚdictÚ_compute_current_unionrZ   r®   r¯   r°   rE   r�   rL   rc   r$   Úadd_done_callbackÚdiscard)r   Úcandidate_filterrµ   Ú
new_filterr²   Ú
old_streamÚtasks          r   r†   Ú"StreamController._reconcile_stream?  sò   é € õ 	Dð ×ÑÑ+Ø×*Ñ*Ñ6Ù˜d×8Ñ8¼$Ð?OÓ:P×QÑQàà×0Ñ0Ð7GÐ0ÐHˆ
Ø—_‘_×6Ñ6Ø×#Ñ# JÓ/ó
ˆ
ð ×(Ñ(ˆ
Ø(ÔØ%/Ô"Ø×Ñ×ÐØÑ!Ü×&Ò&¤|°JÓ'?Ó@ˆDØ×&Ñ&×*Ñ*¨4Ô0Ø×"Ñ" 4×#=Ñ#=×#EÑ#EÕFð "ñ 	ùs   ‚B#DÂ%DÂ&A)Dc              ƒ  ó@   #   • U R                  U5      I Sh  v•N $  N7f)z%Public alias for `_reconcile_stream`.N)r†   )r   r»   s     r   Úreconcile_streamÚ!StreamController.reconcile_stream]  s   é € à×+Ñ+Ð,<Ó=×=Ð=Ñ=ùs   ‚—˜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" U5      $ s  snf )Nr   )Úcompute_union_filter)r’   rÄ   r^   r•   r·   rA   Úappend)r   r¶   rÄ   r}   Úfilterss        r   r¸   Ú'StreamController._compute_current_uniona  sg   € õ 	Kð )-×(;Ñ(;×(BÑ(BÔ(Dó)
Ù(D ŒD�—‘ÖÑ(Dð 	ð )
ð ÑØ�N‰Nœ4 ›;Ô'Ù# GÓ,Ð,ùò)
s   £A(c                ó&   • U R                  U5        g)zCAdvance the reconnect cursor from a command response meta sequence.N)Ú_observe_seq©r   Úseqs     r   Úobserve_applied_through_seqÚ,StreamController.observe_applied_through_seqq  s   € à×Ñ˜#Õr   c                óD   • U R                  UR                  S5      5        g )NrË   )rÉ   rˆ   )r   rœ   s     r   Ú_observe_eventÚStreamController._observe_eventu  s   € Ø×Ñ˜%Ÿ)™) EÓ*Õ+r   c                óv   • [        U[        5      (       a$  U R                  b  XR                  :”  a  Xl        g g g r   )r™   r/   re   rÊ   s     r   rÉ   ÚStreamController._observe_seqx  s/   € Ü�cœ3×Ñ T§\¡\Ñ%9¸SÇ<Á<Ó=OØ�Lð >PÐr   c                óT   • [        U5      nU R                  b  U R                  US'   U$ )NÚsince)r·   re   )r   rA   Úouts      r   r¯   Ú#StreamController._filter_with_since|  s'   € Ü�6‹lˆØ�<‰<Ñ#ØŸ<™<ˆC�‰LØˆ
r   c               óØ   #   • U  S h  v•N nUR                  S5      nUb,  X0R                  ;   a  M.  U R                  R                  U5        U R                  U5        U7v •  Ma   N\
 g 7f)Nr#   )rˆ   r\   r$   rÏ   )r   Úsourcerœ   r#   s       r   r“   ÚStreamController._dedup_iter†  s`   é € Ù!÷ 	�%Ø—y‘y Ó,ˆHØÑ#Ø×3Ñ3Ó3ÙØ×$Ñ$×(Ñ(¨Ô2Ø×Ñ Ô&Ø�Kñ	™6ùs&   ‚A*…A(‰A&ŠA(�AA*Á&A(Á(A*)rd   re   ra   r[   rf   r]   rg   rh   rc   r\   r_   r`   r^   rZ   )ri   r   rQ   z$Callable[[], Awaitable[None]] | NonerR   r/   rS   r/   rT   r/   rU   ÚfloatrV   rÚ   r0   r1   )rn   z	list[str]rk   zlist[list[str]] | Nonerl   z
int | Noner0   úAsyncIterator[Event])r0   r1   )rA   r   r0   r>   )r‚   r/   r0   r1   )rA   r   r0   zAsyncGenerator[Event, None])r¦   r/   r0   r1   )r0   r4   )r»   r   r0   r1   r   )r¶   zSubscribeParams | Noner0   údict[str, Any])rË   r   r0   r1   )rœ   r   r0   r1   )rA   rÜ   r0   rÜ   )rØ   rÛ   r0   rÛ   )r5   r6   r7   r8   r9   r   rq   rJ   r~   rƒ   Úregister_subscriptionÚunregister_subscriptionrp   r‡   Úensure_fanout_runningrŽ   rª   rš   r†   rÁ   r¸   rÌ   rÏ   rÉ   r¯   r“   r;   r<   r   r   rN   rN   `   s6  † ñð* @DØ"Ø"(Ø&'Ø(+Ø'*ñ<ð *ð<ð =ð	<ð
 ð<ð  ð<ð !$ð<ð !&ð<ð  %ð<ð 
õ<ðD .2Ø ñ/àð/ð +ð	/ð
 ð/ð 
õ/ô(Vô$	ô7ð
 3ÐØ6Ðð2Ø%ð2à	$ô2ô*Dð
 3Ðô-'ô^,ôôFGô<>ð
 /3ð
-Ø+ð
-à	õ
-ô ô,ôô÷r   rN   )rK   r   rG   rÚ   r0   r1   ) r9   Ú
__future__r   rE   rv   Úloggingr¤   Úcollectionsr   Úcollections.abcr   r   r   r   Údataclassesr	   r
   Útypingr   Úlangchain_protocolr   r   Úlanggraph_sdk.stream.transportr   r   r   Ú	getLoggerr5   r—   r>   rL   rN   r<   r   r   Ú<module>ré      s€   ðñ
õ #ã Û Û Û Ý #ß NÓ Nß (Ý ç 5ç T÷ ñ  ð8 ×
Ò
˜HÓ
%€ð ÷@ð @ó ð@ð EH÷ ÷ nò nr   