Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
delta.cpython-313.pyc115 linesDownload Raw Back to __pycache__
1�

2��j��	���SSKJr SSKrSSKrSSKJrJr SSKJ	r	J3r4 SSKJr SSK
Jr SSKJr SSKJr SS	KJrJr SS5KJrJrJr SSKJrJrJrJr Sr"S
S\6\\\	\	\	45r 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��^�\rSrSr%SrSrS\S'SSS.SU4SjjjjrSS	jr\	SS7j5r8\	SSj5rSSjrSS
jr
SSjrSSjrSSjrSSjrSSjrSrU=r$)r�agReducer channel that stores only a sentinel in checkpoint blobs and9reconstructs state by replaying ancestor writes through the reducer.10 11!!! warning "Beta"12 13    `DeltaChannel` is in beta. The API and on-disk representation may14    change in future releases. Threads written with `DeltaChannel` today15    are expected to remain readable, but the surrounding contract16    (`BaseCheckpointSaver.get_delta_channel_history`, the17    `_DeltaSnapshot` blob shape, the `counters_since_delta_snapshot`18    metadata field) is not yet stable.19 20The reducer receives the current accumulated value and a batch of writes21in one call: `reducer(state, [write1, write2, ...]) -> new_state`.22 23Reducers must be deterministic and batching-invariant (associative across24folds): applying two consecutive write batches separately must produce the25same state as applying their concatenation once:26 27    reducer(reducer(state, xs), ys) == reducer(state, xs + ys)28 29This lets LangGraph replay checkpointed writes in larger batches than they30were originally produced without changing reconstructed state.31 32Snapshot cadence is driven by two counters: per-channel update count and33total supersteps since last snapshot. `create_checkpoint` writes a full34`_DeltaSnapshot` blob when EITHER the update count reaches35`snapshot_frequency` OR the supersteps count reaches the system-wide36`DELTA_MAX_SUPERSTEPS_SINCE_SNAPSHOT` bound (default 5000), bounding37replay depth even for channels that stop receiving writes.38 39Parameters:40    reducer: `(state, list[writes]) -> new_state`. Must be deterministic41        and batching-invariant as described above.42    typ: The value type (e.g. `list`, `dict`). Inferred automatically43        from the outer type when used inside `Annotated[T, DeltaChannel(...)]`.44    snapshot_frequency: Every Nth update to this channel writes a snapshot45        blob (default `1000`). Must be a positive int.46)�value�reducer�snapshot_frequencyzValue | Anyri��rc�">�US::a[SU35eUc[n[TU]
U5 XlX0l[
U5nU[RR[RR4;a[nU[RR[RR4;a[nU[RR[RR4;a[ nX l[$Ulg)Nrz/snapshot_frequency must be a positive int, got )�47ValueError�list�super�__init__rrr�collections�abcr�MutableSequence�Set�48MutableSet�set�Mapping�MutableMapping�dict�typrr)�selfrr+r�	__class__s    ��^D:\code\apps\devtools\python\user_packages\Python313\site-packages\langgraph/channels/delta.pyr!�DeltaChannel.__init__Es������"��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���!��49�c��[U[5(dgURUR:wag[URUR5$)NF)�50isinstancerrrr)r,�others  r.�__eq__�DeltaChannel.__eq___s>���%��.�.���"�"�e�&>�&>�>������e�m�m�<�<r0c��UR$�N�r+�r,s r.�	ValueType�DeltaChannel.ValueTypef����x�x�r0c��UR$r7r8r9s r.�51UpdateType�DeltaChannel.UpdateTypejr<r0c��URURURURS9nURUlUR52[LaUR53UlU$[R"UR545UlU$)Nr)	r-rr+r�keyrr�_copy�copy)r,�news  r.rC�DeltaChannel.copynsn���n�n��L�L�$�(�(�t�7N�7N��55���(�(���"&�*�*��"7�D�J�J��	��56�>C�Z�Z��57�58�=S��	��59r0c�"�URURURURS9nURUlU[60LaUR5UlU$[U[5(aURUlU$XlU$)z�Initialize from a stored blob.61 62Blob types:63  * `MISSING`: start empty; caller replays writes.64  * `_DeltaSnapshot(value)`: restore value directly from snapshot.65  * plain value (migration from old `BinaryOperatorAggregate` blobs):66    use directly.67r)	r-rr+rrArrr2r	)r,�68checkpointrDs   r.�from_checkpoint�DeltaChannel.from_checkpointvs����n�n��L�L�$�(�(�t�7N�7N��69���(�(����� ����70�C�I�71�72�	�73�N�
3�
3�"�(�(�C�I��74�#�I��75r0c�l�UVVs/sHu p#UPM76 nnnU(dgURnSn[U5HIups[U5up�U(dMU	b[R"U	5OUR5nUS-nMK XFSn77U78(aUR
XZ5UlgUUlgs snnf)z�Apply ancestor writes oldest-to-newest via a single reducer call.79 80If any write is an Overwrite, the last one in the sequence acts as81the reset point: its value becomes the new base and only writes82after it are passed to the reducer.83Nr�)r�	enumeraterrBrCr+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���	&�84�6�N�	�6?�T�\�\�$�2��85�T��86��+s�B0c�h�U(dgSn[U5HCup4[U5upVU(dMUb#[S[RS9n[U5eUnME Ub~[X5uphUb[R"U5OUR5n	[U5VVs/sHup4X2:wdMUPM n87nnU88(aURX�5OU	Ul89gUR[LaUR5OURn	URU	[U55Ul90gs snnf)NFz4Can receive only one Overwrite value per super-step.)�message�91error_codeT)
rLrrr�INVALID_CONCURRENT_GRAPH_UPDATErrBrCr+rrrr)r,rP�
overwrite_idxrSrOrTrN�msg�overwrite_valuerQrVs           r.�update�DeltaChannel.update�s����$(�
��f�%�D�A�%�a�(�H�E��u� �,�.� V�#,�#L�#L��C�-�S�1�1� !�
�&��$�!/��0E�!F��A�#�.��92�93�?�+��X�X�Z�
�94(1��'8�O�'8�t�q�A�<N��'8�I�O�:C����d�6��D�J��!�Z�Z�7�2�t�x�x�z��95�96���\�\�$��V��5��97���Ps�&D.�5D.c�T�UR[La98[5eUR$r7)rrrr9s r.�get�DeltaChannel.get�s!���:�:�� �#�%�%��z�z�r0c�&�UR[L$r7)rrr9s r.�is_available�DeltaChannel.is_available�s���z�z��(�(r0c��[$)a_Return stored representation: always `MISSING`.99 100Snapshot decisions live in `create_checkpoint` (which has the channel101version) and write `_DeltaSnapshot(ch.get())` directly into102`channel_values`. For non-snapshot steps the channel does not appear103in `channel_values`; reconstruction walks ancestor writes via the104saver's `get_delta_channel_history`.105rr9s r.rG�DeltaChannel.checkpoint�s	���r0)rrr+rr7)rz#Callable[[Any, Sequence[Any]], Any]r+ztype[Value] | Noner�int�return�None)r3�objectrk�bool)rkr)rkr106)rGrrkr107)rMzSequence[PendingWrite]rkrl)rPz
Sequence[Any]rkrn)rkrn)�__name__�108__module__�__qualname__�__firstlineno__�__doc__�	__slots__�__annotations__r!r4�propertyr:r>rCrHrWr`rcrfrG�__static_attributes__�
__classcell__)r-s@r.rrs����&�P;�I���109#'�"�110#'�"�4�"� �"�111 �"�112�
"�"�4=�����������*J�(�8�113)�	�	r0r)!�114__future__r�collections.abcr"rCrBrr�typingrr�langgraph.checkpoint.baser� langgraph.checkpoint.serde.typesr	�typing_extensionsr115�langgraph._internal._typingr�langgraph.channels.baser
r�langgraph.channels.binoprrr�langgraph.errorsrrrr�__all__r�r0r.�<module>r�sY��"���.��2�;�"�/�6�T�T�����s�7�5�>�;�s�C��}�#=�sr0
codekingpro/portable-devtools · Team Ai