ó
    ýÞ jù1  ã                  ó–   • S SK J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Jr  \(       a  S SKJr  \" S5      r " S S	\\   5      rg)
é    )ÚannotationsN)Údeque)ÚAsyncIteratorÚ	AwaitableÚCallableÚIterator)ÚTYPE_CHECKINGÚGenericÚTypeVar)Ú	StreamMuxÚTc                  ó®   • \ rS rSrSrSSS.SS j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 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S jjrSrg) ÚStreamChannelé   uÊ  Single-consumer drainable queue for streaming events, with optional
protocol auto-forwarding.

When constructed with a `name`, the StreamMux auto-wires every
`push()` to also inject a `ProtocolEvent` into the main event stream
using the channel's name as the method. When constructed without a
name, the channel is local-only â€” items are only visible to
in-process consumers that iterate the channel directly.

Items are popped off the front as the consumer advances â€” there is
no retention beyond what's currently queued. A channel accepts
exactly one subscriber; a second `__iter__` / `__aiter__` call
raises. Use `tee(n)` / `atee(n)` for fan-out.

Starts unbound â€” neither `__iter__` nor `__aiter__` is available
until the StreamMux calls `_bind(is_async)`. After binding, only
the matching iteration protocol works; the other raises `TypeError`.

Pump wiring (set by the run stream, not by `_bind`):
    - `_request_more`: sync pump callable, returns True if a new
      event was produced.
    - `_arequest_more`: async pump coroutine factory, same contract.

Memory is bounded by caller pace: both sync and async use caller-
driven pumps, so each cursor advance produces at most one event.

Lazy-subscribe: `push` appends to the local buffer only when a
subscriber has registered. Auto-forward via `_wire_fn` always fires
regardless of subscription state.

Lifecycle (`close` / `fail`) is managed by the mux â€” transformers
don't need to close their channels manually.
N)Úmaxlenc               óÒ   • Ub  US::  a  [        S5      eXl        [        5       U l        X 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        g)a„  Initialize the channel.

Args:
    name: Optional protocol channel name. When set, the
        StreamMux wires every `push()` to also inject a
        `ProtocolEvent` into the main event stream. Surfaced
        on the wire as `custom:<name>` for user-defined
        transformers, or as `<name>` for channels owned by a
        native transformer (`_native = True`). When `None`,
        the channel is local-only.
    maxlen: Accepted for forward compatibility; currently
        unused. The caller-driven pump bounds memory naturally
        for single-consumer use.

Raises:
    ValueError: If `maxlen` is not a positive integer or `None`.
Nr   z3StreamChannel maxlen must be a positive int or NoneF)Ú
ValueErrorÚnamer   Ú_itemsÚ_maxlenÚ_closedÚ_errorÚ	_is_asyncÚ_subscribedÚ_request_moreÚ_arequest_moreÚ_wire_fnÚ_mux)Úselfr   r   s      ÚY/var/www/html/gaurav/venv/lib/python3.13/site-packages/langgraph/stream/stream_channel.pyÚ__init__ÚStreamChannel.__init__1   sj   € ð$ Ñ &¨A£+ÜÐRÓSÐSØŒ	Ü,1«GˆŒØ#)ŒØˆŒØ,0ˆŒà&*ˆŒà ˆÔà8<ˆÔØDHˆÔà48ˆŒØ&*ˆ�	ó    c                ó   • Xl         g ©N)r   )r   Úmuxs     r    Ú	_bind_muxÚStreamChannel._bind_muxY   s   € Ø�	r#   c               ó@   • U R                   b  [        S5      eXl         g)a  Bind this channel to sync or async mode.

Called by the StreamMux after transformer registration. Must be
called exactly once before any iteration.

Args:
    is_async: True to enable async iteration, False for sync.

Raises:
    RuntimeError: If the channel has already been bound.
NzStreamChannel is already bound)r   ÚRuntimeError)r   Úis_asyncs     r    Ú_bindÚStreamChannel._bind\   s   € ð �>‰>Ñ%ÜÐ?Ó@Ð@Ø!�r#   c                ó   • Xl         g)z8Install the auto-forward callback (called by StreamMux).N)r   )r   Úfns     r    Ú_wireÚStreamChannel._wirep   s   € à�r#   c                ó&  • U R                   (       aa  U R                  (       a  [        S5      eU R                  b  U R                  R	                  5       OSnU R
                  R                  X!45        U R                  b  U R                  U5        gg)aÜ  Append an item. Auto-forwards if wired.

The local buffer append is a no-op when no subscriber is
registered, but auto-forwarding always fires so wired events
reach the main event log regardless of subscription state.

Items are stored as `(stamp, item)` tuples where stamp is a
monotonic counter from the owning mux. Stamps are stripped by
the default cursors; raw stamped tuples are visible on `_items`.

Raises:
    RuntimeError: If the channel is closed (and subscribed).
z%Cannot push to a closed StreamChannelNr   )r   r   r*   r   Ú_next_push_seqr   Úappendr   )r   ÚitemÚstamps      r    ÚpushÚStreamChannel.pushx   sl   € ð ××Ø�|�|Ü"Ð#JÓKÐKØ26·)±)Ñ2G�D—I‘I×,Ñ,Ô.ÈQˆEØ�K‰K×Ñ ˜}Ô-Ø�=‰=Ñ$Ø�M‰M˜$Õð %r#   c                ó   • SU l         g)zMark the channel as complete.TN)r   ©r   s    r    ÚcloseÚStreamChannel.closeŽ   s	   € àˆ�r#   c                ó   • Xl         SU l        g)zYMark the channel as errored.

Args:
    err: The exception to surface to the subscriber.
TN)r   r   )r   Úerrs     r    ÚfailÚStreamChannel.fail’   s   € ð ŒØˆ�r#   c                óÐ   • U R                   c  [        S5      eU R                   (       a  [        S5      eU R                  (       a  [        S5      eSU l        U R	                  5       $ )zÂSubscribe and return a sync cursor. Can be called only once.

Raises:
    TypeError: If the channel is unbound or bound to async mode.
    RuntimeError: If the channel already has a subscriber.
úVStreamChannel has not been bound yet. Register the transformer with a StreamMux first.uF   This StreamChannel is bound to async mode â€” use 'async for' instead.z@StreamChannel already has a subscriber; use .tee(n) for fan-out.T)r   Ú	TypeErrorr   r*   Ú_sync_cursorr:   s    r    Ú__iter__ÚStreamChannel.__iter__Ÿ   sn   € ð �>‰>Ñ!ÜðCóð ð �>�>ÜØXóð ð ××ÜØRóð ð  ˆÔØ× Ñ Ó"Ð"r#   c              #  óX  #   •  U R                   (       a!  U R                   R                  5       u  pUv •  OrU R                  (       a  U R                  b  U R                  eg U R                  b9  U R	                  5       (       d#  U R                   (       d  U R                  (       d  g Og M¦  7fr%   )r   Úpopleftr   r   r   ©r   Ú_stampr5   s      r    rD   ÚStreamChannel._sync_cursor¶   sz   é € ØØ�{�{Ø#Ÿ{™{×2Ñ2Ó4‘�Ø“
Ø——Ø—;‘;Ñ*ØŸ+™+Ð%ØØ×#Ñ#Ñ/Ø×)Ñ)×+Ñ+ØŸ;Ÿ;¨t¯|¯|Øøàñ ùs   ‚B(B*c                óÐ   • U R                   c  [        S5      eU R                   (       d  [        S5      eU R                  (       a  [        S5      eSU l        U R	                  5       $ )zÃSubscribe and return an async cursor. Can be called only once.

Raises:
    TypeError: If the channel is unbound or bound to sync mode.
    RuntimeError: If the channel already has a subscriber.
rB   u?   This StreamChannel is bound to sync mode â€” use 'for' instead.zAStreamChannel already has a subscriber; use .atee(n) for fan-out.T)r   rC   r   r*   Ú_async_cursorr:   s    r    Ú	__aiter__ÚStreamChannel.__aiter__Ê   sn   € ð �>‰>Ñ!ÜðCóð ð �~�~ÜØQóð ð ××ÜØSóð ð  ˆÔØ×!Ñ!Ó#Ð#r#   c               ón  #   •  U R                   (       a"  U R                   R                  5       u  pU7v •  OzU R                  (       a  U R                  b  U R                  eg U R                  bA  U R	                  5       I S h  v•N (       d#  U R                   (       d  U R                  (       d  g Og M¯   N07fr%   )r   rH   r   r   r   rI   s      r    rM   ÚStreamChannel._async_cursorá   s‚   é € ØØ�{�{Ø#Ÿ{™{×2Ñ2Ó4‘�Ø”
Ø——Ø—;‘;Ñ*ØŸ+™+Ð%ØØ×$Ñ$Ñ0Ø!×0Ñ0Ó2×2Õ2ØŸ;Ÿ;¨t¯|¯|Øøàñ ñ 3ùs   ‚B B5ÂB3Â1B5c                óò   ^^^^• US:  a  [        S5      eU R                  5       m[        U5       Vs/ sH  n[        5       PM     snmS/mSUUU4S jjm[	        U4S j[        U5       5       5      $ s  snf )a  Subscribe and return `n` independent sync iterators.

Each branch has its own buffer; items pulled from the
underlying cursor are copied into every branch. Branches are
naturally bounded by caller pace since the sync pump is
caller-driven.

Args:
    n: Number of branches to create. Must be >= 1.

Returns:
    A tuple of `n` iterators over the same underlying stream.

Raises:
    TypeError: If the channel is unbound or bound to async mode.
    RuntimeError: If the channel already has a subscriber.
    ValueError: If `n` < 1.
é   ztee() requires n >= 1Fc              3  óÜ   >#   • TU    n U(       a  UR                  5       v •  O1TS   (       a  g  [        T5      nT H  nUR                  U5        M     MM  ! [         a    STS'    g f = f7f©NTr   )rH   ÚnextÚStopIterationr4   )ÚiÚbufr5   ÚbÚbuffersÚ	exhaustedÚsources       €€€r    ÚbranchÚ!StreamChannel.tee.<locals>.branch  sq   øé € Ø˜!‘*ˆCØÞØŸ+™+›-Ó'Ø˜q—\ØðÜ# F›|˜ó %˜ØŸ™ žñ %ñ øô )ó Ø'+˜	 !™Ùðüs'   ƒ,A,°A »A,ÁA)Á&A,Á(A)Á)A,c              3  ó2   >#   • U H  nT" U5      v •  M     g 7fr%   © ©Ú.0rX   r^   s     €r    Ú	<genexpr>Ú$StreamChannel.tee.<locals>.<genexpr>  ó   øé € Ð1© 1‘V˜A—Y�Yªùó   ƒ)rX   ÚintÚreturnúIterator[T])r   rE   Úranger   Útuple)r   ÚnÚ_r^   r[   r\   r]   s      @@@@r    ÚteeÚStreamChannel.teeõ   sj   û€ ð& ˆq‹5ÜÐ4Ó5Ð5Ø—‘“ˆÜ49¸!´HÓ"=±H¨q¤5¦7±HÑ"=ˆØ�Gˆ	÷	'ñ 	'ô  Ô1¬¨a¬Ó1Ó1Ð1ùò' #>s   ³A4c                ó*  ^^^^^^• US:  a  [        S5      eU R                  5       m[        U5       Vs/ sH  n[        5       PM     snmS/mS/m[        R
                  " 5       mSUUUUU4S jjm[        U4S j[        U5       5       5      $ s  snf )a.  Subscribe and return `n` independent async iterators.

Caller-driven fan-out: each branch's `__anext__` either pops
from its own buffer or, under a shared `asyncio.Lock`, pulls
one item from the underlying cursor and distributes it to
every branch's buffer.

Args:
    n: Number of branches to create. Must be >= 1.

Returns:
    A tuple of `n` async iterators over the same underlying
    stream.

Raises:
    TypeError: If the channel is unbound or bound to sync mode.
    RuntimeError: If the channel already has a subscriber.
    ValueError: If `n` < 1.
rS   zatee() requires n >= 1FNc               óR  >#   • TU    n U(       a  UR                  5       7v •  M  TS   (       a  TS   b  TS   eg T IS h  v•N   U(       d
  TS   (       a   S S S 5      IS h  v•N   Mb   T	R                  5       I S h  v•N nT H  nUR	                  U5        M     S S S 5      IS h  v•N   M¦   Nm NM N4! [         a    STS'    S S S 5      IS h  v•N    MÐ  [         a&  nUTS'   STS'    S nAS S S 5      IS h  v•N    Mú  S nAff = f Na! , IS h  v•N  (       d  f       Nv= f7frU   )rH   Ú	__anext__ÚStopAsyncIterationÚ	Exceptionr4   )
rX   rY   r5   ÚerZ   r[   Úerrorr\   Úlockr]   s
        €€€€€r    r^   Ú"StreamChannel.atee.<locals>.branch<  sö   øé € Ø˜!‘*ˆCØÞØŸ+™+›-Ó'ÙØ˜Q—<Ø˜Q‘xÑ+Ø# A™h˜Øßš4Þ˜i¨ŸlØ ÷  Ÿ4™4ð!Ø%+×%5Ñ%5Ó%7×7˜ó %˜ØŸ™ žñ %÷  Ÿ4ñ ó  ñ  8øÜ-ó !Ø'+˜	 !™Ø ÷  Ÿ4š4ô %ó !Ø#$˜˜a™Ø'+˜	 !™Û ÷  Ÿ4š4ûð!ú÷  Ÿ4Ÿ4˜4üsË   ƒ?D'ÁB0ÁD'ÁDÁD'Á$B2Á%D'Á,B6Á?B4Â B6ÂDÂD'Â)DÂ*D'Â2D'Â4B6Â6DÃDÃD'ÃCÃD'Ã	DÃ"
DÃ,DÃ0D'Ã;C>Ã<D'ÄDÄDÄD'ÄD$ÄDÄD$Ä D'c              3  ó2   >#   • U H  nT" U5      v •  M     g 7fr%   ra   rb   s     €r    rd   Ú%StreamChannel.atee.<locals>.<genexpr>U  rf   rg   )rX   rh   ri   úAsyncIterator[T])r   rN   rk   r   ÚasyncioÚLockrl   )	r   rm   rn   r^   r[   rw   r\   rx   r]   s	      @@@@@@r    ÚateeÚStreamChannel.atee   s}   ý€ ð( ˆq‹5ÜÐ5Ó6Ð6Ø—‘Ó!ˆÜ49¸!´HÓ"=±H¨q¤5¦7±HÑ"=ˆØ�Gˆ	Ø-1¨FˆÜ�|Š|‹~ˆ÷	'ó 	'ô2 Ô1¬¨a¬Ó1Ó1Ð1ùò= #>s   µB)r   r   r   r   r   r   r   r   r   r   r   r%   )r   z
str | Noner   z
int | Noneri   ÚNone)r&   r   ri   r�   )r+   Úboolri   r�   )r/   zCallable[[T], None]ri   r�   )r5   r   ri   r�   )ri   r�   )r>   ÚBaseExceptionri   r�   )ri   rj   )ri   r|   )é   )rm   rh   ri   ztuple[Iterator[T], ...])rm   rh   ri   ztuple[AsyncIterator[T], ...])Ú__name__Ú
__module__Ú__qualname__Ú__firstlineno__Ú__doc__r!   r'   r,   r0   r7   r;   r?   rE   rD   rN   rM   ro   r   Ú__static_attributes__ra   r#   r    r   r      s\   † ñ ðD"+È÷ "+ð "+ôPô"ô(ô ô,ôô#ô.ô($ô.ö()2÷V52ñ 52r#   r   )Ú
__future__r   r}   Úcollectionsr   Úcollections.abcr   r   r   r   Útypingr	   r
   r   Úlanggraph.stream._muxr   r   r   ra   r#   r    Ú<module>r�      s;   ðÝ "ã Ý ß HÓ Hß 2Ñ 2æÝ/áˆCƒL€ôG2�G˜A‘Jõ G2r#   