§
    šŠtj\  ã            	      óî   — d dl mZ d dlZd dlZd dlmZmZ d dlm	Z	m
Z
 d dlmZ d dlmZ d dlmZ d dlmZ d d	lmZmZ d d
lmZmZmZ d dlmZmZmZmZ dZ G d„ de
e         ee	e	e	f         ¦  «        Z dS )é    )ÚannotationsN)ÚCallableÚSequence)ÚAnyÚGeneric)ÚPendingWrite)Ú_DeltaSnapshot)ÚSelf©ÚMISSING)ÚBaseChannelÚValue)Ú_get_overwriteÚ_operators_equalÚ_strip_extras)ÚEmptyChannelErrorÚ	ErrorCodeÚInvalidUpdateErrorÚcreate_error_message)ÚDeltaChannelc                  ó®   ‡ — e Zd ZU dZdZded<   	 d%ddœd&ˆ fd„Zd'd„Zed(d„¦   «         Z	ed(d„¦   «         Z
d)d„Zd*d„Zd+d„Zd,d!„Zd(d"„Zd-d#„Zd(d$„Zˆ xZS ).r   aß  Reducer channel that stores only a sentinel in checkpoint blobs and
    reconstructs state by replaying ancestor writes through the reducer.

    !!! warning "Beta"

        `DeltaChannel` is in beta. The API and on-disk representation may
        change in future releases. Threads written with `DeltaChannel` today
        are expected to remain readable, but the surrounding contract
        (`BaseCheckpointSaver.get_delta_channel_history`, the
        `_DeltaSnapshot` blob shape, the `counters_since_delta_snapshot`
        metadata field) is not yet stable.

    The reducer receives the current accumulated value and a batch of writes
    in one call: `reducer(state, [write1, write2, ...]) -> new_state`.

    Reducers must be deterministic and batching-invariant (associative across
    folds): applying two consecutive write batches separately must produce the
    same state as applying their concatenation once:

        reducer(reducer(state, xs), ys) == reducer(state, xs + ys)

    This lets LangGraph replay checkpointed writes in larger batches than they
    were originally produced without changing reconstructed state.

    Snapshot cadence is driven by two counters: per-channel update count and
    total supersteps since last snapshot. `create_checkpoint` writes a full
    `_DeltaSnapshot` blob when EITHER the update count reaches
    `snapshot_frequency` OR the supersteps count reaches the system-wide
    `DELTA_MAX_SUPERSTEPS_SINCE_SNAPSHOT` bound (default 5000), bounding
    replay depth even for channels that stop receiving writes.

    Parameters:
        reducer: `(state, list[writes]) -> new_state`. Must be deterministic
            and batching-invariant as described above.
        typ: The value type (e.g. `list`, `dict`). Inferred automatically
            from the outer type when used inside `Annotated[T, DeltaChannel(...)]`.
        snapshot_frequency: Every Nth update to this channel writes a snapshot
            blob (default `1000`). Must be a positive int.
    )ÚvalueÚreducerÚsnapshot_frequencyzValue | Anyr   Niè  ©r   r   ú#Callable[[Any, Sequence[Any]], Any]Útypútype[Value] | Noner   ÚintÚreturnÚNonec               óî  •— |dk    rt          d|› �¦  «        ‚|€t          }t          ¦   «                              |¦  «         || _        || _        t          |¦  «        }|t          j        j	        t          j        j
        fv rt          }|t          j        j        t          j        j        fv rt          }|t          j        j        t          j        j        fv rt           }|| _        t$          | _        d S )Nr   z/snapshot_frequency must be a positive int, got )Ú
ValueErrorÚlistÚsuperÚ__init__r   r   r   ÚcollectionsÚabcr   ÚMutableSequenceÚSetÚ
MutableSetÚsetÚMappingÚMutableMappingÚdictr   r   r   )Úselfr   r   r   Ú	__class__s       €úV/var/www/html/CA-Chatbot/venv/lib/python3.11/site-packages/langgraph/channels/delta.pyr&   zDeltaChannel.__init__E   sÝ   ø€ ð  Ò"Ð"ÝØVÐBTÐVÐVñô ð ð ˆ;ÝˆCÝ‰Œ×Ò˜ÑÔÐØˆŒØ"4ˆÔÝ˜CÑ Ô ˆØ•;”?Ô+­[¬_Ô-LÐMÐMÐMÝˆCØ•;”?Ô&­¬Ô(BÐCÐCÐCÝˆCØ•;”?Ô*­K¬OÔ,JÐKÐKÐKÝˆCØˆŒÝ!ˆŒ
ˆ
ˆ
ó    ÚotherÚobjectÚboolc                óˆ   — t          |t          ¦  «        sdS | j        |j        k    rdS t          | j        |j        ¦  «        S )NF)Ú
isinstancer   r   r   r   )r0   r4   s     r2   Ú__eq__zDeltaChannel.__eq___   sC   € Ý˜%¥Ñ.Ô.ð 	Ø�5ØÔ" eÔ&>Ò>Ð>Ø�5Ý ¤¨e¬mÑ<Ô<Ð<r3   r   c                ó   — | j         S ©N©r   ©r0   s    r2   Ú	ValueTypezDeltaChannel.ValueTypef   ó	   € àŒxˆr3   c                ó   — | j         S r;   r<   r=   s    r2   Ú
UpdateTypezDeltaChannel.UpdateTypej   r?   r3   r
   c                óÒ   — |                       | j        | j        | j        ¬¦  «        }| j        |_        | j        t          u r| j        nt          j        | j        ¦  «        |_        |S )Nr   )	r1   r   r   r   Úkeyr   r   Ú_copyÚcopy)r0   Únews     r2   rE   zDeltaChannel.copyn   s]   € Ø�nŠnØŒL˜$œ(°tÔ7Nð ñ 
ô 
ˆð ”(ˆŒØ"&¤*µÐ"7Ð"7�D”J�J½U¼ZÈÌ
Ñ=SÔ=SˆŒ	Øˆ
r3   Ú
checkpointc                ó  — |                       | j        | j        | j        ¬¦  «        }| j        |_        |t
          u r|                      ¦   «         |_        n)t          |t          ¦  «        r|j        |_        n||_        |S )a*  Initialize from a stored blob.

        Blob types:
          * `MISSING`: start empty; caller replays writes.
          * `_DeltaSnapshot(value)`: restore value directly from snapshot.
          * plain value (migration from old `BinaryOperatorAggregate` blobs):
            use directly.
        r   )	r1   r   r   r   rC   r   r   r8   r	   )r0   rG   rF   s      r2   Úfrom_checkpointzDeltaChannel.from_checkpointv   s{   € ð �nŠnØŒL˜$œ(°tÔ7Nð ñ 
ô 
ˆð ”(ˆŒØ�Ð Ð ØŸš™
œ
ˆCŒIˆIÝ˜
¥NÑ3Ô3ð 	#Ø"Ô(ˆCŒIˆIà"ˆCŒIØˆ
r3   ÚwritesúSequence[PendingWrite]c                ó:  — d„ |D ¦   «         }|sdS | j         }d}t          |¦  «        D ]H\  }}t          |¦  «        \  }}|r/|�t          j        |¦  «        n|                      ¦   «         }|dz   }ŒI||d…         }	|	r|                      ||	¦  «        n|| _         dS )a
  Apply ancestor writes oldest-to-newest via a single reducer call.

        If any write is an Overwrite, the last one in the sequence acts as
        the reset point: its value becomes the new base and only writes
        after it are passed to the reducer.
        c                ó   — g | ]\  }}}|‘Œ	S © rN   )Ú.0Ú_Úvs      r2   ú
<listcomp>z.DeltaChannel.replay_writes.<locals>.<listcomp>’   s   € Ð*Ð*Ð*™˜˜1˜a�!Ð*Ð*Ð*r3   Nr   é   )r   Ú	enumerater   rD   rE   r   r   )
r0   rJ   ÚvaluesÚbaseÚstartÚirQ   Úis_owÚow_valueÚ	remainings
             r2   Úreplay_writeszDeltaChannel.replay_writes‹   s¿   € ð +Ð* 6Ð*Ñ*Ô*ˆØð 	ØˆFØŒzˆØˆÝ˜fÑ%Ô%ð 	ð 	‰DˆAˆqÝ,¨QÑ/Ô/‰OˆE�8Øð Ø/7Ð/C•u”z (Ñ+Ô+Ð+ÈÏÊÉÌ�Ø˜A™�øØ˜5˜6˜6”Nˆ	Ø6?ÐI�T—\’\ $¨	Ñ2Ô2Ð2ÀTˆŒ
ˆ
ˆ
r3   rU   úSequence[Any]c                óø  — |sdS d }t          |¦  «        D ]G\  }}t          |¦  «        \  }}|r.|�*t          dt          j        ¬¦  «        }t          |¦  «        ‚|}ŒH|�It          ||         ¦  «        \  }}|�t          j        |¦  «        n|                      ¦   «         | _	        dS | j	        t          u r|                      ¦   «         n| j	        }	|                      |	t          |¦  «        ¦  «        | _	        dS )NFz4Can receive only one Overwrite value per super-step.)ÚmessageÚ
error_codeT)rT   r   r   r   ÚINVALID_CONCURRENT_GRAPH_UPDATEr   rD   rE   r   r   r   r   r$   )
r0   rU   Úoverwrite_idxrX   rQ   rY   rP   ÚmsgÚoverwrite_valuerV   s
             r2   ÚupdatezDeltaChannel.updateŸ   s  € Øð 	Ø�5Ø$(ˆÝ˜fÑ%Ô%ð 		"ð 		"‰DˆAˆqÝ% aÑ(Ô(‰HˆE�1Øð "Ø Ð,Ý.Ø VÝ#,Ô#Lðñ ô �Cõ -¨SÑ1Ô1Ð1Ø !�øØÐ$Ý!/°°}Ô0EÑ!FÔ!FÑˆAˆð #Ð.õ ”
˜?Ñ+Ô+Ð+à—X’X‘Z”Zð ŒJð
 �4Ø!œZ­7Ð2Ð2ˆt�xŠx‰zŒzˆz¸¼
ˆØ—\’\ $­¨V©¬Ñ5Ô5ˆŒ
Øˆtr3   c                óH   — | j         t          u rt          ¦   «         ‚| j         S r;   )r   r   r   r=   s    r2   ÚgetzDeltaChannel.get¹   s#   € ØŒ:�Ð Ð Ý#Ñ%Ô%Ð%ØŒzÐr3   c                ó   — | j         t          uS r;   )r   r   r=   s    r2   Úis_availablezDeltaChannel.is_available¾   s   € ØŒz¥Ð(Ð(r3   c                ó   — t           S )a�  Return stored representation: always `MISSING`.

        Snapshot decisions live in `create_checkpoint` (which has the channel
        version) and write `_DeltaSnapshot(ch.get())` directly into
        `channel_values`. For non-snapshot steps the channel does not appear
        in `channel_values`; reconstruction walks ancestor writes via the
        saver's `get_delta_channel_history`.
        r   r=   s    r2   rG   zDeltaChannel.checkpointÁ   s	   € õ ˆr3   r;   )r   r   r   r   r   r   r    r!   )r4   r5   r    r6   )r    r   )r    r
   )rG   r   r    r
   )rJ   rK   r    r!   )rU   r]   r    r6   )r    r6   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__Ú	__slots__Ú__annotations__r&   r9   Úpropertyr>   rA   rE   rI   r\   re   rg   ri   rG   Ú__classcell__)r1   s   @r2   r   r      sZ  ø€ € € € € € ð&ð &ðP ;€IØÐÐÑð
 #'ð"ð
 #'ð"ð "ð "ð "ð "ð "ð "ð "ð4=ð =ð =ð =ð ðð ð ñ „Xðð ðð ð ñ „Xððð ð ð ðð ð ð ð*Jð Jð Jð Jð(ð ð ð ð4ð ð ð ð
)ð )ð )ð )ð	ð 	ð 	ð 	ð 	ð 	ð 	ð 	r3   r   )!Ú
__future__r   Úcollections.abcr'   rE   rD   r   r   Útypingr   r   Úlanggraph.checkpoint.baser   Ú langgraph.checkpoint.serde.typesr	   Útyping_extensionsr
   Úlanggraph._internal._typingr   Úlanggraph.channels.baser   r   Úlanggraph.channels.binopr   r   r   Úlanggraph.errorsr   r   r   r   Ú__all__r   rN   r3   r2   ú<module>r~      so  ðØ "Ð "Ð "Ð "Ð "Ð "à Ð Ð Ð Ø Ð Ð Ð Ø .Ð .Ð .Ð .Ð .Ð .Ð .Ð .Ø Ð Ð Ð Ð Ð Ð Ð à 2Ð 2Ð 2Ð 2Ð 2Ð 2Ø ;Ð ;Ð ;Ð ;Ð ;Ð ;Ø "Ð "Ð "Ð "Ð "Ð "à /Ð /Ð /Ð /Ð /Ð /Ø 6Ð 6Ð 6Ð 6Ð 6Ð 6Ð 6Ð 6Ø TÐ TÐ TÐ TÐ TÐ TÐ TÐ TÐ TÐ Tðð ð ð ð ð ð ð ð ð ð ð ð €ðqð qð qð qð q�7˜5”> ;¨s°C¸¨}Ô#=ñ qô qð qð qð qr3   