ó
    ýÞ j\  ã            	      óÚ   • 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Jr  S SKJ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Jr  S SKJrJrJrJr  Sr " S S\
\   \\	\	\	4   5      r g)é    )ÚannotationsN)ÚCallableÚSequence)ÚAnyÚGeneric)ÚPendingWrite)Ú_DeltaSnapshot)ÚSelf©ÚMISSING)ÚBaseChannelÚValue)Ú_get_overwriteÚ_operators_equalÚ_strip_extras)ÚEmptyChannelErrorÚ	ErrorCodeÚInvalidUpdateErrorÚcreate_error_message)ÚDeltaChannelc                  óÔ   ^ • \ rS rSr% SrSrS\S'    SSS.       SU 4S jjjjrSS	 jr\	SS
 j5       r
\	SS j5       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rU =r$ )r   é   ag  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   iè  ©r   c               ó"  >• US::  a  [        SU 35      eUc  [        n[        TU ]  U5        Xl        X0l        [        U5      nU[        R                  R                  [        R                  R                  4;   a  [        nU[        R                  R                  [        R                  R                  4;   a  [        nU[        R                  R                  [        R                  R                  4;   a  [         nX l        [$        U l        g )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ÚdictÚtypr   r   )Úselfr   r+   r   Ú	__class__s       €ÚR/var/www/html/gaurav/venv/lib/python3.13/site-packages/langgraph/channels/delta.pyr!   ÚDeltaChannel.__init__E   sÒ   ø€ ð  Ó"ÜØAÐBTÐAUÐVóð ð ‰;ÜˆCÜ‰Ñ˜ÔØŒØ"4ÔÜ˜CÓ ˆØ”;—?‘?×+Ñ+¬[¯_©_×-LÑ-LÐMÓMÜˆCØ”;—?‘?×&Ñ&¬¯©×(BÑ(BÐCÓCÜˆCØ”;—?‘?×*Ñ*¬K¯O©O×,JÑ,JÐKÓKÜˆCØŒÜ!ˆ�
ó    c                ó¤   • [        U[        5      (       d  gU R                  UR                  :w  a  g[        U R                  UR                  5      $ )NF)Ú
isinstancer   r   r   r   )r,   Úothers     r.   Ú__eq__ÚDeltaChannel.__eq___   s>   € Ü˜%¤×.Ñ.ØØ×"Ñ" e×&>Ñ&>Ó>ØÜ §¡¨e¯m©mÓ<Ð<r0   c                ó   • U R                   $ ©N©r+   ©r,   s    r.   Ú	ValueTypeÚDeltaChannel.ValueTypef   ó   € à�x‰xˆr0   c                ó   • U R                   $ r7   r8   r9   s    r.   Ú
UpdateTypeÚDeltaChannel.UpdateTypej   r<   r0   c                ó  • U R                  U R                  U R                  U R                  S9nU R                  Ul        U R
                  [        L a  U R
                  Ul        U$ [        R                  " U R
                  5      Ul        U$ )Nr   )	r-   r   r+   r   Úkeyr   r   Ú_copyÚcopy)r,   Únews     r.   rC   ÚDeltaChannel.copyn   sn   € Ø�n‰nØ�L‰L˜$Ÿ(™(°t×7NÑ7Nð ð 
ˆð —(‘(ˆŒØ"&§*¡*´Ò"7�D—J‘JˆŒ	Øˆ
ô >C¿ZºZÈÏ
É
Ó=SˆŒ	Øˆ
r0   c                ó"  • U R                  U R                  U R                  U R                  S9nU R                  Ul        U[
        L a  U R                  5       Ul        U$ [        U[        5      (       a  UR                  Ul        U$ Xl        U$ )zú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   )	r-   r   r+   r   rA   r   r   r2   r	   )r,   Ú
checkpointrD   s      r.   Úfrom_checkpointÚDeltaChannel.from_checkpointv   s„   € ð �n‰nØ�L‰L˜$Ÿ(™(°t×7NÑ7Nð ð 
ˆð —(‘(ˆŒØœÒ ØŸ™›
ˆCŒIð
 ˆ
ô	 ˜
¤N×3Ñ3Ø"×(Ñ(ˆCŒIð ˆ
ð #ŒIØˆ
r0   c                ój  • U VVs/ sH  u    p#UPM
     nnnU(       d  gU R                   nSn[        U5       HI  u  ps[        U5      u  p‰U(       d  M  U	b  [        R                  " U	5      OU R                  5       nUS-   nMK     XFS n
U
(       a  U R                  XZ5      U l         gUU l         gs  snnf )zê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.
Nr   é   )r   Ú	enumerater   rB   rC   r+   r   )r,   ÚwritesÚ_ÚvÚvaluesÚbaseÚstartÚiÚis_owÚow_valueÚ	remainings              r.   Úreplay_writesÚDeltaChannel.replay_writes‹   sž   € ñ $*Ô*¡6™˜˜1“!¡6ˆÑ*ÞØØ�z‰zˆØˆÜ˜fÖ%‰DˆAÜ,¨QÓ/‰OˆEßˆuØ/7Ñ/C”u—z’z (Ô+ÈÏÉË�Ø˜A™’ñ	 &ð
 ˜6�Nˆ	Þ6?�T—\‘\ $Ó2ˆ�
ÀTˆ�
ùó +s   †B/c                óÜ  • U(       d  gS n[        U5       HC  u  p4[        U5      u  pVU(       d  M  Ub#  [        S[        R                  S9n[        U5      eUnME     Ub>  [        X   5      u  phUb  [        R                  " U5      OU R                  5       U l	        gU R                  [        L a  U R                  5       OU R                  n	U R                  U	[        U5      5      U l	        g)NFz4Can receive only one Overwrite value per super-step.)ÚmessageÚ
error_codeT)rL   r   r   r   ÚINVALID_CONCURRENT_GRAPH_UPDATEr   rB   rC   r+   r   r   r   r   )
r,   rP   Úoverwrite_idxrS   rO   rT   rN   ÚmsgÚoverwrite_valuerQ   s
             r.   ÚupdateÚDeltaChannel.updateŸ   sÕ   € ÞØØ$(ˆÜ˜fÖ%‰DˆAÜ% aÓ(‰HˆEßˆuØ Ñ,Ü.Ø VÜ#,×#LÑ#Lñ�Cô -¨SÓ1Ð1Ø !’ñ &ð Ñ$Ü!/°Ñ0EÓ!FÑˆAð #Ñ.ô —
’
˜?Ô+à—X‘X“Zð ŒJð
 Ø!ŸZ™Z¬7Ò2ˆt�x‰xŒz¸¿
¹
ˆØ—\‘\ $¬¨V«Ó5ˆŒ
Ør0   c                óT   • U R                   [        L a
  [        5       eU R                   $ r7   )r   r   r   r9   s    r.   ÚgetÚDeltaChannel.get¹   s!   € Ø�:‰:œÒ Ü#Ó%Ð%Ø�z‰zÐr0   c                ó&   • U R                   [        L$ r7   )r   r   r9   s    r.   Úis_availableÚDeltaChannel.is_available¾   s   € Ø�z‰z¤Ð(Ð(r0   c                ó   • [         $ )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   r9   s    r.   rG   ÚDeltaChannel.checkpointÁ   s	   € ô ˆr0   )r   r   r+   r   r7   )r   z#Callable[[Any, Sequence[Any]], Any]r+   ztype[Value] | Noner   ÚintÚreturnÚNone)r3   Úobjectrk   Úbool)rk   r   )rk   r
   )rG   r   rk   r
   )rM   zSequence[PendingWrite]rk   rl   )rP   zSequence[Any]rk   rn   )rk   rn   )Ú__name__Ú
__module__Ú__qualname__Ú__firstlineno__Ú__doc__Ú	__slots__Ú__annotations__r!   r4   Úpropertyr:   r>   rC   rH   rW   r`   rc   rf   rG   Ú__static_attributes__Ú__classcell__)r-   s   @r.   r   r      s«   ø‡ ñ&ðP ;€IØÓð
 #'ð"ð
 #'ñ"à4ð"ð  ð"ð
  ð"ð 
÷"ñ "ô4=ð óó ðð óó ðôôô*Jô(ô4ô
)÷	ò 	r0   r   )!Ú
__future__r   Úcollections.abcr"   rC   rB   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   © r0   r.   Ú<module>r…      sY   ðÝ "ã Û ß .ß å 2Ý ;Ý "å /ß 6ß TÑ T÷ó ð €ôq�7˜5‘> ;¨s°C¸¨}Ñ#=õ qr0   