§
    šŠtj•0  ã                  óf  — d dl mZ d dlZd dlmZmZm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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 d dlmZ d dlmZ d dlm Z m!Z! dZ"ee
dge
f         Z#dId„Z$dJd„Z%dKd„Z&dLd!„Z'dMd"„Z(dNd'„Z)dOd.„Z*ddddd/œdPd8„Z+dQd=„Z,ddd>œdRdF„Z-ddd>œdRdG„Z.dSdH„Z/dS )Té    )ÚannotationsN)ÚCallableÚIterableÚMapping)ÚdatetimeÚtimezone)ÚAnyÚcast)ÚRunnableConfig)ÚBaseCheckpointSaverÚ
Checkpoint)Úuuid6)Ú_DeltaSnapshot)Ú#DELTA_MAX_SUPERSTEPS_SINCE_SNAPSHOT)ÚPUSH)ÚMISSING)ÚBaseChannel)ÚDeltaChannel)ÚManagedValueMappingÚManagedValueSpecé   Úreturnr   c                 óÈ   — t          t          t          t          d¬¦  «        ¦  «        t	          j        t          j        ¦  «                             ¦   «         i i i ¬¦  «        S )Néþÿÿÿ©Ú	clock_seq)ÚvÚidÚtsÚchannel_valuesÚchannel_versionsÚversions_seen)	r   ÚLATEST_VERSIONÚstrr   r   Únowr   ÚutcÚ	isoformat© ó    úZ/var/www/html/CA-Chatbot/venv/lib/python3.11/site-packages/langgraph/pregel/_checkpoint.pyÚempty_checkpointr+      sV   € ÝÝ
Ý�u˜rÐ"Ñ"Ô"Ñ#Ô#ÝŒ<�œÑ%Ô%×/Ò/Ñ1Ô1ØØØðñ ô ð r)   ÚstepÚintÚtask_idr$   c           
     ó¼   — t          t          j        |¦  «        ¦  «                             d¦  «        }| d›d|d         › d|d         › d|d         › d|d         › �	S )a  Synthetic task id for exit-mode DeltaChannel writes.

    Embeds the superstep in the first UUID group so `ORDER BY task_id, idx`
    preserves chronological order while remaining a valid RFC UUID (required by
    Postgres `checkpoint_writes.task_id uuid` columns).
    ú-Ú08dé   é   é   r   )r$   ÚuuidÚUUIDÚsplit)r,   r.   Úpartss      r*   Úexit_delta_task_idr9   '   sh   € õ •”	˜'Ñ"Ô"Ñ#Ô#×)Ò)¨#Ñ.Ô.€EØÐDÐDÐD˜˜qœÐDÐD E¨!¤HÐDÐD¨u°Q¬xÐDÐD¸%À¼(ÐDÐDÐDr)   ÚchannelsúMapping[str, BaseChannel]Úcounters_since_delta_snapshotúMapping[str, tuple[int, int]]úset[str]c                ó2  — t          ¦   «         }|                      ¦   «         D ]s\  }}t          |t          ¦  «        r|                     ¦   «         sŒ/|                     |d¦  «        \  }}||j        k    s|t          k    r|                     |¦  «         Œt|S )u7  Return the set of DeltaChannel names that should snapshot now.

    A channel snapshots when EITHER its accumulated update count reaches
    `snapshot_frequency` OR the total supersteps since its last snapshot
    reaches `DELTA_MAX_SUPERSTEPS_SINCE_SNAPSHOT`. This is a pure
    predicate â€” no mutation.
    ©r   r   )	ÚsetÚitemsÚ
isinstancer   Úis_availableÚgetÚsnapshot_frequencyr   Úadd)r:   r<   ÚresultÚnameÚchÚupdatesÚ
superstepss          r*   Údelta_channels_to_snapshotrM   2   sœ   € õ ‘u”u€FØ—N’NÑ$Ô$ð ð ‰ˆˆbÝ˜"�lÑ+Ô+ð 	°2·?²?Ñ3DÔ3Dð 	ØØ;×?Ò?ÀÀfÑMÔMÑˆ�à�rÔ,Ò,Ð,ØÕ@Ò@Ð@à�JŠJ�tÑÔÐøØ€Mr)   Ú	run_tasksúIterable[Any]c                ó   — d„ | D ¦   «         S )zDChannel names written by an update_state superstep (excluding PUSH).c                óB   — h | ]}|j         D ]\  }}|t          k    ¯|’ŒŒS r(   )Úwritesr   )Ú.0ÚtaskÚcÚ_s       r*   ú	<setcomp>z2get_updated_channels_from_tasks.<locals>.<setcomp>N   s/   € ÐIÐIÐI�$°´ÐIÐI©¨¨1¸qÅDºy¸yˆA¸y¸y¸y¸yr)   r(   )rN   s    r*   Úget_updated_channels_from_tasksrX   J   s   € ð JÐI˜)ÐIÑIÔIÐIr)   c                ó>   — d„ |                       ¦   «         D ¦   «         S )zFDeltaChannels to snapshot on the first update_state of a fresh thread.c                ój   — h | ]0\  }}t          |t          ¦  «        ¯|                     ¦   «         ¯.|’Œ1S r(   )rC   r   rD   )rS   ÚkrJ   s      r*   rW   z7get_delta_channels_from_all_channels.<locals>.<setcomp>U   sP   € ð ð ð áˆAˆrÝ�b�,Ñ'Ô'ðð -/¯OªOÑ,=Ô,=ðØ	ðð ð r)   )rB   )r:   s    r*   Ú$get_delta_channels_from_all_channelsr\   Q   s-   € ðð à—^’^Ñ%Ô%ðñ ô ð r)   Úupdated_channelsÚprev_metadataúMapping[str, Any] | Noneúdict[str, tuple[int, int]]c               ó  — t          |pi                      d¦  «        pi ¦  «        }i }|                      ¦   «         D ]I\  }}t          |t          ¦  «        sŒ|                     |d¦  «        \  }}|dz  }||v r|dz  }||f||<   ŒJ|S )z Advance ``counters_since_delta_snapshot`` for update_state on a non-fresh thread.

    Mirrors the per-superstep counter bump in ``_loop._put_checkpoint``.
    r<   r@   r2   )ÚdictrE   rB   rC   r   )	r:   r]   r^   Úprev_countersÚnew_countersÚch_namerJ   ÚuÚss	            r*   Ú$create_metadata_for_update_state_apirh   \   s¶   € õ Ø	Ð	˜"×!Ò!Ð"AÑBÔBÐHÀbñô €Mð 02€LØ—~’~Ñ'Ô'ð 'ð '‰ˆ�Ý˜"�lÑ+Ô+ð 	ØØ× Ò  ¨&Ñ1Ô1‰ˆˆ1Ø	ˆQ‰ˆØÐ&Ð&Ð&Ø�‰FˆAØ!" A ˆ�WÑÐØÐr)   Úparentsúdict[str, Any]Úsaved_metadataÚis_fresh_threadÚboolútuple[set[str], dict[str, Any]]c               óÞ   — d||dœ}|rt          | ¦  «        |fS t          | ||¬¦  «        }t          | |¦  «        }|D ]}	d||	<   Œd„ |                     ¦   «         D ¦   «         }
|
r|
|d<   ||fS )zEReturn ``(channels_to_snapshot, metadata)`` for an update_state head.Úupdate)Úsourcer,   ri   )r^   r@   c                ó&   — i | ]\  }}|d k    ¯||“ŒS )r@   r(   ©rS   r[   r   s      r*   ú
<dictcomp>z?create_checkpoint_plan_for_update_state_api.<locals>.<dictcomp>�   s#   € ÐEÐEÐE™˜˜A¸¸fº¸��1¸¸¸r)   r<   )r\   rh   rM   rB   )r:   r]   r,   ri   rk   rl   Úmetadatard   Úchannels_to_snapshotr[   Únon_zeros              r*   Ú+create_checkpoint_plan_for_update_state_apirx   u   s½   € ð ØØð ð  €Hð
 ð HÝ3°HÑ=Ô=¸xÐGÐGå7ØØØ$ðñ ô €Lõ
 6°hÀÑMÔMÐØ!ð !ð !ˆØ ˆ�Q‰ˆØEÐE ×!3Ò!3Ñ!5Ô!5ÐEÑEÔE€HØð =Ø4<ˆÐ0Ñ1Ø Ð)Ð)r)   )r   r]   Úget_next_versionrv   Ú
checkpointú Mapping[str, BaseChannel] | Noner   ú
str | Noneúset[str] | Nonery   úGetNextVersion | Nonerv   c               óh  — t          j        t          j        ¦  «                             ¦   «         }|pt          ¦   «         }|€| d         }| d         }	n‘i }t          | d         ¦  «        }	|D ]w}
|
|	vrŒ||
         }|
|v rB|�|�|
|vr ||	|
         d¦  «        |	|
<   t          |                     ¦   «         ¦  «        ||
<   ŒU| 	                    ¦   «         }|t          ur|||
<   Œxt          t          ||pt          t          |¬¦  «        ¦  «        ||	| d         |€dnt          |¦  «        ¬¦  «        S )uÑ  Build a new Checkpoint from the previous one and live channel state.

    For each name in `channels_to_snapshot`, a `_DeltaSnapshot(value)` blob
    is written into `channel_values[k]`. Other delta channels are omitted
    from `channel_values` â€” the ancestor walk reconstructs their state
    from `checkpoint_writes`. Callers compute the set via
    `delta_channels_to_snapshot(channels, counters)`; defaults to empty
    (no snapshots) when not provided.
    Nr    r!   r   r"   ©r   r   r   r    r!   r"   r]   )r   r%   r   r&   r'   rA   rb   r   rE   rz   r   r   r#   r$   r   Úsorted)rz   r:   r,   r   r]   ry   rv   r   Úvaluesr!   r[   rJ   r   s                r*   Úcreate_checkpointrƒ   •   sl  € õ& 
Œ•h”lÑ	#Ô	#×	-Ò	-Ñ	/Ô	/€BØ/Ð8µ3±5´5ÐØÐØÐ,Ô-ˆØ%Ð&8Ô9ÐÐàˆÝ 
Ð+=Ô >Ñ?Ô?ÐØð 	"ð 	"ˆAØÐ(Ð(Ð(ØØ˜!”ˆBØÐ(Ð(Ð(ð" $Ð/Ø$Ð,°Ð9IÐ0IÐ0Ià*:Ð*:Ð;KÈAÔ;NÐPTÑ*UÔ*UÐ$ QÑ'Ý*¨2¯6ª6©8¬8Ñ4Ô4��q‘	�	à—M’M‘O”O�Ø�GÐ#Ð#Ø !�F˜1‘IøÝÝ
ØØÐ+••U TÐ*Ñ*Ô*Ñ+Ô+ØØ)Ø  Ô1Ø!1Ð!9˜˜½vÐFVÑ?WÔ?Wðñ ô ð r)   Úspecr   ÚstoredÚobjectc                óB   — t          | t          ¦  «        sdS |t          u S )u  True if `spec` is a `DeltaChannel` and no value is stored at this
    checkpoint, requiring an ancestor walk to reconstruct.

    `_DeltaSnapshot` blobs and plain values (migration) resolve directly via
    `from_checkpoint` â€” only absence (`MISSING`) triggers replay.
    F)rC   r   r   )r„   r…   s     r*   Ú_needs_replayrˆ   Ù   s&   € õ �d�LÑ)Ô)ð ØˆuØ•WÐÐr)   )ÚsaverÚconfigÚspecsú,Mapping[str, BaseChannel | ManagedValueSpec]r‰   úBaseCheckpointSaver | NonerŠ   úRunnableConfig | Noneú5tuple[Mapping[str, BaseChannel], ManagedValueMapping]c               óŠ  ‡— i }i }|                       ¦   «         D ]%\  }}t          |t          ¦  «        r|||<   Œ |||<   Œ&ˆfd„|                      ¦   «         D ¦   «         }i }	|r|�|�|                     ||¬¦  «        }	i }
|                      ¦   «         D ]«\  }}||	v rit	          t
          |¦  «        }|	|         }|                     |                     dt          ¦  «        ¦  «        }| 	                    |d         ¦  «         |}n4|                     ‰d                              |t          ¦  «        ¦  «        }||
|<   Œ¬|
|fS )a   Hydrate channels from a checkpoint.

    For most channels, `spec.from_checkpoint(checkpoint["channel_values"][k])`
    is sufficient. `DeltaChannel` is the exception: when the channel is
    absent from `channel_values`, an ancestor walk via
    `saver.get_delta_channel_history` is required to find the nearest seed
    (`_DeltaSnapshot` blob or pre-migration plain value) and accumulate
    the writes between it and the target. All delta channels needing
    replay are batched into a single saver call.
    c           	     óx   •— g | ]6\  }}t          |‰d                               |t          ¦  «        ¦  «        ¯4|‘Œ7S ©r    ©rˆ   rE   r   ©rS   r[   r„   rz   s      €r*   ú
<listcomp>z,channels_from_checkpoint.<locals>.<listcomp>þ   óS   ø€ ð !ð !ð !áˆAˆtÝ˜˜zÐ*:Ô;×?Ò?ÀÅ7ÑKÔKÑLÔLð!Ø	ð!ð !ð !r)   N©rŠ   r:   ÚseedrR   r    )
rB   rC   r   Úget_delta_channel_historyr
   r   Úfrom_checkpointrE   r   Úreplay_writes©r‹   rz   r‰   rŠ   Úchannel_specsÚmanaged_specsr[   r   Údelta_channelsÚ	historiesr:   r„   Ú
delta_specÚhistoryÚ	replay_chrJ   s    `              r*   Úchannels_from_checkpointr¤   å   s‰  ø€ ð" -/€MØ13€MØ—’‘”ð !ð !‰ˆˆ1Ý�a�Ñ%Ô%ð 	!Ø ˆM˜!ÑÐà ˆM˜!ÑÐð!ð !ð !ð !à$×*Ò*Ñ,Ô,ð!ñ !ô !€Nð
 $&€IØð 
˜%Ð+°Ð0BØ×3Ò3Ø Nð 4ñ 
ô 
ˆ	ð (*€HØ ×&Ò&Ñ(Ô(ð 
ð 
‰ˆˆ4à�	ˆ>ˆ>Ý�l¨DÑ1Ô1ˆJØ ”lˆGØ"×2Ò2°7·;²;¸vÅwÑ3OÔ3OÑPÔPˆIØ×#Ò# G¨HÔ$5Ñ6Ô6Ð6ØˆBˆBà×%Ò% jÐ1AÔ&B×&FÒ&FÀqÍ'Ñ&RÔ&RÑSÔSˆBØˆ�‰ˆØ�]Ð"Ð"r)   c             ƒ  óš  ‡K  — i }i }|                       ¦   «         D ]%\  }}t          |t          ¦  «        r|||<   Œ |||<   Œ&ˆfd„|                      ¦   «         D ¦   «         }i }	|r!|�|�|                     ||¬¦  «        ƒ d{V —†}	i }
|                      ¦   «         D ]«\  }}||	v rit	          t
          |¦  «        }|	|         }|                     |                     dt          ¦  «        ¦  «        }| 	                    |d         ¦  «         |}n4|                     ‰d                              |t          ¦  «        ¦  «        }||
|<   Œ¬|
|fS )zAAsync version of `channels_from_checkpoint`. See docstring there.c           	     óx   •— g | ]6\  }}t          |‰d                               |t          ¦  «        ¦  «        ¯4|‘Œ7S r’   r“   r”   s      €r*   r•   z-achannels_from_checkpoint.<locals>.<listcomp>(  r–   r)   Nr—   r˜   rR   r    )
rB   rC   r   Úaget_delta_channel_historyr
   r   rš   rE   r   r›   rœ   s    `              r*   Úachannels_from_checkpointr¨     s«  øè è € ð -/€MØ13€MØ—’‘”ð !ð !‰ˆˆ1Ý�a�Ñ%Ô%ð 	!Ø ˆM˜!ÑÐà ˆM˜!ÑÐð!ð !ð !ð !à$×*Ò*Ñ,Ô,ð!ñ !ô !€Nð
 $&€IØð 
˜%Ð+°Ð0BØ×:Ò:Ø Nð ;ñ 
ô 
ð 
ð 
ð 
ð 
ð 
ð 
ˆ	ð (*€HØ ×&Ò&Ñ(Ô(ð 
ð 
‰ˆˆ4à�	ˆ>ˆ>Ý�l¨DÑ1Ô1ˆJØ ”lˆGØ"×2Ò2°7·;²;¸vÅwÑ3OÔ3OÑPÔPˆIØ×#Ò# G¨HÔ$5Ñ6Ô6Ð6ØˆBˆBà×%Ò% jÐ1AÔ&B×&FÒ&FÀqÍ'Ñ&RÔ&RÑSÔSˆBØˆ�‰ˆØ�]Ð"Ð"r)   c                ó  — t          | d         | d         | d         | d                              ¦   «         | d                              ¦   «         d„ | d                              ¦   «         D ¦   «         |                      dd ¦  «        ¬	¦  «        S )
Nr   r   r   r    r!   c                ó>   — i | ]\  }}||                      ¦   «         “ŒS r(   )Úcopyrs   s      r*   rt   z#copy_checkpoint.<locals>.<dictcomp>I  s&   € ÐSÐSÐS¡t q¨!�q˜!Ÿ&š&™(œ(ÐSÐSÐSr)   r"   r]   r€   )r   r«   rB   rE   )rz   s    r*   Úcopy_checkpointr¬   B  sŽ   € ÝØ
�SŒ/Ø�dÔØ�dÔØ!Ð"2Ô3×8Ò8Ñ:Ô:Ø#Ð$6Ô7×<Ò<Ñ>Ô>ØSÐS¨z¸/Ô/J×/PÒ/PÑ/RÔ/RÐSÑSÔSØ#ŸšÐ(:¸DÑAÔAðñ ô ð r)   )r   r   )r,   r-   r.   r$   r   r$   )r:   r;   r<   r=   r   r>   )rN   rO   r   r>   )r:   r;   r   r>   )r:   r;   r]   r>   r^   r_   r   r`   )r:   r;   r]   r>   r,   r-   ri   rj   rk   r_   rl   rm   r   rn   )rz   r   r:   r{   r,   r-   r   r|   r]   r}   ry   r~   rv   r}   r   r   )r„   r   r…   r†   r   rm   )
r‹   rŒ   rz   r   r‰   r�   rŠ   rŽ   r   r�   )rz   r   r   r   )0Ú
__future__r   r5   Úcollections.abcr   r   r   r   r   Útypingr	   r
   Úlangchain_core.runnablesr   Úlanggraph.checkpoint.baser   r   Úlanggraph.checkpoint.base.idr   Ú langgraph.checkpoint.serde.typesr   Úlanggraph._internal._configr   Úlanggraph._internal._constantsr   Úlanggraph._internal._typingr   Úlanggraph.channels.baser   Úlanggraph.channels.deltar   Úlanggraph.managed.baser   r   r#   ÚGetNextVersionr+   r9   rM   rX   r\   rh   rx   rƒ   rˆ   r¤   r¨   r¬   r(   r)   r*   ú<module>r»      s   ðØ "Ð "Ð "Ð "Ð "Ð "à €€€Ø 7Ð 7Ð 7Ð 7Ð 7Ð 7Ð 7Ð 7Ð 7Ð 7Ø 'Ð 'Ð 'Ð 'Ð 'Ð 'Ð 'Ð 'Ø Ð Ð Ð Ð Ð Ð Ð à 3Ð 3Ð 3Ð 3Ð 3Ð 3ðð ð ð ð ð ð ð ð /Ð .Ð .Ð .Ð .Ð .Ø ;Ð ;Ð ;Ð ;Ð ;Ð ;à KÐ KÐ KÐ KÐ KÐ KØ /Ð /Ð /Ð /Ð /Ð /Ø /Ð /Ð /Ð /Ð /Ð /Ø /Ð /Ð /Ð /Ð /Ð /Ø 1Ð 1Ð 1Ð 1Ð 1Ð 1Ø HÐ HÐ HÐ HÐ HÐ HÐ HÐ Hà€à˜3 ˜+ sÐ*Ô+€ðð ð ð ðEð Eð Eð Eðð ð ð ð0Jð Jð Jð Jðð ð ð ðð ð ð ð2*ð *ð *ð *ðJ Ø(,Ø.2Ø,0ðAð Að Að Að Að AðH	ð 	ð 	ð 	ð  )-Ø$(ð0#ð 0#ð 0#ð 0#ð 0#ð 0#ðn )-Ø$(ð'#ð '#ð '#ð '#ð '#ð '#ðT	ð 	ð 	ð 	ð 	ð 	r)   