§
    šŠtjù1  ã                  ó    — d dl mZ d dlZd dlmZ d dlmZmZmZm	Z	 d dl
mZmZmZ erd dlmZ  ed¦  «        Z G d„ d	ee         ¦  «        ZdS )
é    )ÚannotationsN©Údeque)ÚAsyncIteratorÚ	AwaitableÚCallableÚIterator)ÚTYPE_CHECKINGÚGenericÚTypeVar)Ú	StreamMuxÚTc                  ó†   — e Zd ZdZd(ddœd)d
„Zd*d„Zd+d„Zd,d„Zd-d„Zd.d„Z	d/d„Z
d0d„Zd0d„Zd1d„Zd1d „Zd2d3d%„Zd2d4d'„ZdS )5ÚStreamChannelu.  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)ÚmaxlenÚnameú
str | Noner   ú
int | NoneÚreturnÚNonec               óæ   — |�|dk    rt          d¦  «        ‚|| _        t          ¦   «         | _        || _        d| _        d| _        d| _        d| _        d| _	        d| _
        d| _        d| _        dS )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)Ú
ValueErrorr   r   Ú_itemsÚ_maxlenÚ_closedÚ_errorÚ	_is_asyncÚ_subscribedÚ_request_moreÚ_arequest_moreÚ_wire_fnÚ_mux)Úselfr   r   s      ú]/var/www/html/CA-Chatbot/venv/lib/python3.11/site-packages/langgraph/stream/stream_channel.pyÚ__init__zStreamChannel.__init__1   sy   € ð$ Ð &¨A¢+ +ÝÐRÑSÔSÐSØˆŒ	Ý,1©G¬GˆŒØ#)ˆŒØˆŒØ,0ˆŒà&*ˆŒà ˆÔà8<ˆÔØDHˆÔà48ˆŒØ&*ˆŒ	ˆ	ˆ	ó    Úmuxr   c                ó   — || _         d S ©N)r"   )r#   r'   s     r$   Ú	_bind_muxzStreamChannel._bind_muxY   s   € ØˆŒ	ˆ	ˆ	r&   Úis_asyncÚboolc               ó@   — | j         �t          d¦  «        ‚|| _         dS )aS  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#   r+   s     r$   Ú_bindzStreamChannel._bind\   s&   € ð Œ>Ð%ÝÐ?Ñ@Ô@Ð@Ø!ˆŒˆˆr&   ÚfnúCallable[[T], None]c                ó   — || _         dS )z8Install the auto-forward callback (called by StreamMux).N)r!   )r#   r0   s     r$   Ú_wirezStreamChannel._wirep   s   € àˆŒˆˆr&   Úitemr   c                óø   — | j         rT| j        rt          d¦  «        ‚| j        �| j                             ¦   «         nd}| j                             ||f¦  «         | j        �|                      |¦  «         dS dS )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#   r4   Ústamps      r$   ÚpushzStreamChannel.pushx   sŠ   € ð Ôð 	.ØŒ|ð LÝ"Ð#JÑKÔKÐKØ26´)Ð2G�D”I×,Ò,Ñ.Ô.Ð.ÈQˆEØŒK×Ò  t˜}Ñ-Ô-Ð-ØŒ=Ð$Ø�MŠM˜$ÑÔÐÐÐð %Ð$r&   c                ó   — d| _         dS )zMark the channel as complete.TN)r   ©r#   s    r$   ÚclosezStreamChannel.closeŽ   s   € àˆŒˆˆr&   ÚerrÚBaseExceptionc                ó"   — || _         d| _        dS )zqMark the channel as errored.

        Args:
            err: The exception to surface to the subscriber.
        TN)r   r   )r#   r=   s     r$   ÚfailzStreamChannel.fail’   s   € ð ˆŒØˆŒˆˆr&   úIterator[T]c                ó¼   — | j         €t          d¦  «        ‚| j         rt          d¦  «        ‚| j        rt          d¦  «        ‚d| _        |                      ¦   «         S )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.
        Nú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__zStreamChannel.__iter__Ÿ   sƒ   € ð Œ>Ð!ÝðCñô ð ð Œ>ð 	ÝØXñô ð ð Ôð 	ÝØRñô ð ð  ˆÔØ× Ò Ñ"Ô"Ð"r&   c              #  óä   K  — 	 | j         r!| j                              ¦   «         \  }}|V — nE| j        r| j        �| j        ‚d S | j        �%|                      ¦   «         s| j         s	| j        sd S nd S Œnr)   )r   Úpopleftr   r   r   ©r#   Ú_stampr4   s      r$   rE   zStreamChannel._sync_cursor¶   s–   è è € ð	ØŒ{ð Ø#œ{×2Ò2Ñ4Ô4‘�˜Ø�
�
�
�
Ø”ð 	Ø”;Ð*Øœ+Ð%Ø�ØÔ#Ð/Ø×)Ò)Ñ+Ô+ð Øœ;ð ¨t¬|ð Ø˜øà�ð	r&   úAsyncIterator[T]c                ó¼   — | j         €t          d¦  «        ‚| j         st          d¦  «        ‚| j        rt          d¦  «        ‚d| _        |                      ¦   «         S )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.
        NrC   u?   This StreamChannel is bound to sync mode â€” use 'for' instead.zAStreamChannel already has a subscriber; use .atee(n) for fan-out.T)r   rD   r   r.   Ú_async_cursorr;   s    r$   Ú	__aiter__zStreamChannel.__aiter__Ê   sƒ   € ð Œ>Ð!ÝðCñô ð ð Œ~ð 	ÝØQñô ð ð Ôð 	ÝØSñô ð ð  ˆÔØ×!Ò!Ñ#Ô#Ð#r&   c               óò   K  — 	 | j         r"| j                              ¦   «         \  }}|W V — nK| j        r| j        �| j        ‚d S | j        �+|                      ¦   «         ƒ d {V —†s| j         s	| j        sd S nd S Œur)   )r   rH   r   r   r    rI   s      r$   rM   zStreamChannel._async_cursorá   sª   è è € ð	ØŒ{ð Ø#œ{×2Ò2Ñ4Ô4‘�˜Ø�
�
�
�
�
Ø”ð 	Ø”;Ð*Øœ+Ð%Ø�ØÔ$Ð0Ø!×0Ò0Ñ2Ô2Ð2Ð2Ð2Ð2Ð2Ð2ð Øœ;ð ¨t¬|ð Ø˜øà�ð	r&   é   ÚnÚintútuple[Iterator[T], ...]c                óô   ‡‡‡‡— |dk     rt          d¦  «        ‚|                      ¦   «         Šd„ t          |¦  «        D ¦   «         ŠdgŠdˆˆˆfd	„Št          ˆfd
„t          |¦  «        D ¦   «         ¦  «        S )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 >= 1c                ó*   — g | ]}t          ¦   «         ‘ŒS © r   ©Ú.0Ú_s     r$   ú
<listcomp>z%StreamChannel.tee.<locals>.<listcomp>  ó   € Ð"=Ð"=Ð"=¨q¥5¡7¤7Ð"=Ð"=Ð"=r&   FÚirR   r   rA   c              3  óä   •K  — ‰|          }	 |r|                      ¦   «         V — nK‰d         rd S 	 t          ‰¦  «        }n# t          $ r	 d‰d<   Y d S w xY w‰D ]}|                     |¦  «         ŒŒe©NTr   )rH   ÚnextÚStopIterationr7   )r]   Úbufr4   ÚbÚbuffersÚ	exhaustedÚsources       €€€r$   Úbranchz!StreamChannel.tee.<locals>.branch  s©   øè è € Ø˜!”*ˆCð'Øð 'ØŸ+š+™-œ-Ð'Ð'Ð'Ð'Ø˜q”\ð 	'Ø�FðÝ# F™|œ|˜˜øÝ(ð ð ð Ø'+˜	 !™Ø˜˜ðøøøð %ð 'ð '˜ØŸš ™œ˜˜ð's   ±A ÁAÁAc              3  ó.   •K  — | ]} ‰|¦  «        V — Œd S r)   rW   ©rY   r]   rg   s     €r$   ú	<genexpr>z$StreamChannel.tee.<locals>.<genexpr>  ó+   øè è € Ð1Ð1 1�V�V˜A‘Y”YÐ1Ð1Ð1Ð1Ð1Ð1r&   )r]   rR   r   rA   )r   rF   ÚrangeÚtuple)r#   rQ   rg   rd   re   rf   s     @@@@r$   ÚteezStreamChannel.teeõ   s    øøøø€ ð& ˆqŠ5ˆ5ÝÐ4Ñ5Ô5Ð5Ø—’‘”ˆØ"=Ð"=µE¸!±H´HÐ"=Ñ"=Ô"=ˆØ�Gˆ	ð	'ð 	'ð 	'ð 	'ð 	'ð 	'ð 	'ð 	'õ  Ð1Ð1Ð1Ð1­¨a©¬Ð1Ñ1Ô1Ñ1Ô1Ð1r&   útuple[AsyncIterator[T], ...]c                ó(  ‡‡‡‡‡‡— |dk     rt          d¦  «        ‚|                      ¦   «         Šd„ t          |¦  «        D ¦   «         ŠdgŠdgŠt          j        ¦   «         Šdˆˆˆˆˆfd
„Št          ˆfd„t          |¦  «        D ¦   «         ¦  «        S )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.
        rU   zatee() requires n >= 1c                ó*   — g | ]}t          ¦   «         ‘ŒS rW   r   rX   s     r$   r[   z&StreamChannel.atee.<locals>.<listcomp>7  r\   r&   FNr]   rR   r   rK   c               ó,  •K  — ‰|          }	 |r|                      ¦   «         W V — Œ‰d         r‰d         �‰d         ‚d S ‰4 ƒd {V —† |s‰d         r	 d d d ¦  «        ƒd {V —† Œ[	 ‰	                     ¦   «         ƒ d {V —†}nS# t          $ r d‰d<   Y d d d ¦  «        ƒd {V —† Œ™t          $ r%}|‰d<   d‰d<   Y d }~d d d ¦  «        ƒd {V —† ŒÂd }~ww xY w‰D ]}|                     |¦  «         Œ	 d d d ¦  «        ƒd {V —† n# 1 ƒd {V —†swxY w Y   �Œ	r_   )rH   Ú	__anext__ÚStopAsyncIterationÚ	Exceptionr7   )
r]   rb   r4   Úerc   rd   Úerrorre   Úlockrf   s
        €€€€€r$   rg   z"StreamChannel.atee.<locals>.branch<  s™  øè è € Ø˜!”*ˆCð'Øð ØŸ+š+™-œ-Ð'Ð'Ð'Ð'ØØ˜Q”<ð Ø˜Q”xÐ+Ø# Aœh˜Ø�FØð 'ð 'ð 'ð 'ð 'ð 'ð 'ð 'Øð !˜i¨œlð !Ø ð'ð 'ð 'ñ 'ô 'ð 'ð 'ð 'ð 'ð 'ð 'ð 'ð 'ð!Ø%+×%5Ò%5Ñ%7Ô%7Ð7Ð7Ð7Ð7Ð7Ð7˜˜øÝ-ð !ð !ð !Ø'+˜	 !™Ø ð'ð 'ð 'ñ 'ô 'ð 'ð 'ð 'ð 'ð 'ð 'ð 'ð 'õ %ð !ð !ð !Ø#$˜˜a™Ø'+˜	 !™Ø ˜˜˜ð'ð 'ð 'ñ 'ô 'ð 'ð 'ð 'ð 'ð 'ð 'ð 'ð 'øøøøð!øøøð %ð 'ð '˜ØŸš ™œ˜˜ð'ð'ð 'ð 'ñ 'ô 'ð 'ð 'ð 'ð 'ð 'ð 'ð 'ð 'ð 'ð 'ð 'ð 'ð 'ð 'ð 'ð 'ð 'ð 'øøøð 'ð 'ð 'ð 'ñ'sH   Á	DÁ(BÂDÂCÂDÂ%	CÂ.
CÂ8DÃCÃDÄ
DÄDc              3  ó.   •K  — | ]} ‰|¦  «        V — Œd S r)   rW   ri   s     €r$   rj   z%StreamChannel.atee.<locals>.<genexpr>U  rk   r&   )r]   rR   r   rK   )r   rN   rl   ÚasyncioÚLockrm   )r#   rQ   rg   rd   rw   re   rx   rf   s     @@@@@@r$   ÚateezStreamChannel.atee   sÀ   øøøøøø€ ð( ˆqŠ5ˆ5ÝÐ5Ñ6Ô6Ð6Ø—’Ñ!Ô!ˆØ"=Ð"=µE¸!±H´HÐ"=Ñ"=Ô"=ˆØ�Gˆ	Ø-1¨FˆÝŒ|‰~Œ~ˆð	'ð 	'ð 	'ð 	'ð 	'ð 	'ð 	'ð 	'ð 	'ð 	'õ2 Ð1Ð1Ð1Ð1­¨a©¬Ð1Ñ1Ô1Ñ1Ô1Ð1r&   r)   )r   r   r   r   r   r   )r'   r   r   r   )r+   r,   r   r   )r0   r1   r   r   )r4   r   r   r   )r   r   )r=   r>   r   r   )r   rA   )r   rK   )rP   )rQ   rR   r   rS   )rQ   rR   r   ro   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__r%   r*   r/   r3   r9   r<   r@   rF   rE   rN   rM   rn   r|   rW   r&   r$   r   r      s@  € € € € € ð ð  ðD"+Èð "+ð "+ð "+ð "+ð "+ð "+ðPð ð ð ð"ð "ð "ð "ð(ð ð ð ð ð  ð  ð  ð,ð ð ð ðð ð ð ð#ð #ð #ð #ð.ð ð ð ð($ð $ð $ð $ð.ð ð ð ð()2ð )2ð )2ð )2ð )2ðV52ð 52ð 52ð 52ð 52ð 52ð 52r&   r   )Ú
__future__r   rz   Úcollectionsr   Úcollections.abcr   r   r   r	   Útypingr
   r   r   Úlanggraph.stream._muxr   r   r   rW   r&   r$   ú<module>r†      sÝ   ðØ "Ð "Ð "Ð "Ð "Ð "à €€€Ø Ð Ð Ð Ð Ð Ø HÐ HÐ HÐ HÐ HÐ HÐ HÐ HÐ HÐ HÐ HÐ HØ 2Ð 2Ð 2Ð 2Ð 2Ð 2Ð 2Ð 2Ð 2Ð 2àð 0Ø/Ð/Ð/Ð/Ð/Ð/à€GˆC�L„L€ðG2ð G2ð G2ð G2ð G2�G˜A”Jñ G2ô G2ð G2ð G2ð G2r&   