mirror of
https://github.com/langchain-ai/langgraph.git
synced 2026-10-01 05:55:14 +02:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
3129e521ea | ||
|
|
4be610c6bc | ||
|
|
98b10ba0ff | ||
|
|
ffe4e8cd6d | ||
|
|
f75add0ee9 | ||
|
|
902be1eb2b | ||
|
|
3b969324db | ||
|
|
e43f0f5f1e | ||
|
|
8ccdedcc1b | ||
|
|
7ea31e5e45 | ||
|
|
313e0e7e0a | ||
|
|
b9919ace32 | ||
|
|
4217373227 | ||
|
|
65ab10f017 |
Generated
+3
-3
@@ -943,11 +943,11 @@ wheels = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "urllib3"
|
name = "urllib3"
|
||||||
version = "2.7.0"
|
version = "2.8.0"
|
||||||
source = { registry = "https://pypi.org/simple" }
|
source = { registry = "https://pypi.org/simple" }
|
||||||
sdist = { url = "https://files.pythonhosted.org/packages/53/0c/06f8b233b8fd13b9e5ee11424ef85419ba0d8ba0b3138bf360be2ff56953/urllib3-2.7.0.tar.gz", hash = "sha256:231e0ec3b63ceb14667c67be60f2f2c40a518cb38b03af60abc813da26505f4c", size = 433602, upload-time = "2026-05-07T16:13:18.596Z" }
|
sdist = { url = "https://files.pythonhosted.org/packages/e3/05/b17359e1cefb4f909b5e40b1b90a496d987258916dbbf88e842c729f510e/urllib3-2.8.0.tar.gz", hash = "sha256:63bf2ead4c879426ebf22ef2a781eeb4aa3b4ae798a0435506f8687fd5bb9b63", size = 458972, upload-time = "2026-09-15T19:29:36.253Z" }
|
||||||
wheels = [
|
wheels = [
|
||||||
{ url = "https://files.pythonhosted.org/packages/7f/3e/5db95bcf282c52709639744ca2a8b149baccf648e39c8cc87553df9eae0c/urllib3-2.7.0-py3-none-any.whl", hash = "sha256:9fb4c81ebbb1ce9531cce37674bbc6f1360472bc18ca9a553ede278ef7276897", size = 131087, upload-time = "2026-05-07T16:13:17.151Z" },
|
{ url = "https://files.pythonhosted.org/packages/92/9d/c4e665119135114480843e7ab388fa94d8480650450e6f8e26b70d323a4c/urllib3-2.8.0-py3-none-any.whl", hash = "sha256:0cf3cae568d36aa9576b28dfb35f11328f1cb974ca7647d9475ebb86c75ac6e3", size = 135717, upload-time = "2026-09-15T19:29:34.577Z" },
|
||||||
]
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
|
|||||||
Generated
+3
-3
@@ -1122,11 +1122,11 @@ wheels = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "urllib3"
|
name = "urllib3"
|
||||||
version = "2.7.0"
|
version = "2.8.0"
|
||||||
source = { registry = "https://pypi.org/simple" }
|
source = { registry = "https://pypi.org/simple" }
|
||||||
sdist = { url = "https://files.pythonhosted.org/packages/53/0c/06f8b233b8fd13b9e5ee11424ef85419ba0d8ba0b3138bf360be2ff56953/urllib3-2.7.0.tar.gz", hash = "sha256:231e0ec3b63ceb14667c67be60f2f2c40a518cb38b03af60abc813da26505f4c", size = 433602, upload-time = "2026-05-07T16:13:18.596Z" }
|
sdist = { url = "https://files.pythonhosted.org/packages/e3/05/b17359e1cefb4f909b5e40b1b90a496d987258916dbbf88e842c729f510e/urllib3-2.8.0.tar.gz", hash = "sha256:63bf2ead4c879426ebf22ef2a781eeb4aa3b4ae798a0435506f8687fd5bb9b63", size = 458972, upload-time = "2026-09-15T19:29:36.253Z" }
|
||||||
wheels = [
|
wheels = [
|
||||||
{ url = "https://files.pythonhosted.org/packages/7f/3e/5db95bcf282c52709639744ca2a8b149baccf648e39c8cc87553df9eae0c/urllib3-2.7.0-py3-none-any.whl", hash = "sha256:9fb4c81ebbb1ce9531cce37674bbc6f1360472bc18ca9a553ede278ef7276897", size = 131087, upload-time = "2026-05-07T16:13:17.151Z" },
|
{ url = "https://files.pythonhosted.org/packages/92/9d/c4e665119135114480843e7ab388fa94d8480650450e6f8e26b70d323a4c/urllib3-2.8.0-py3-none-any.whl", hash = "sha256:0cf3cae568d36aa9576b28dfb35f11328f1cb974ca7647d9475ebb86c75ac6e3", size = 135717, upload-time = "2026-09-15T19:29:34.577Z" },
|
||||||
]
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
|
|||||||
Generated
+3
-3
@@ -1057,11 +1057,11 @@ wheels = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "urllib3"
|
name = "urllib3"
|
||||||
version = "2.7.0"
|
version = "2.8.0"
|
||||||
source = { registry = "https://pypi.org/simple" }
|
source = { registry = "https://pypi.org/simple" }
|
||||||
sdist = { url = "https://files.pythonhosted.org/packages/53/0c/06f8b233b8fd13b9e5ee11424ef85419ba0d8ba0b3138bf360be2ff56953/urllib3-2.7.0.tar.gz", hash = "sha256:231e0ec3b63ceb14667c67be60f2f2c40a518cb38b03af60abc813da26505f4c", size = 433602, upload-time = "2026-05-07T16:13:18.596Z" }
|
sdist = { url = "https://files.pythonhosted.org/packages/e3/05/b17359e1cefb4f909b5e40b1b90a496d987258916dbbf88e842c729f510e/urllib3-2.8.0.tar.gz", hash = "sha256:63bf2ead4c879426ebf22ef2a781eeb4aa3b4ae798a0435506f8687fd5bb9b63", size = 458972, upload-time = "2026-09-15T19:29:36.253Z" }
|
||||||
wheels = [
|
wheels = [
|
||||||
{ url = "https://files.pythonhosted.org/packages/7f/3e/5db95bcf282c52709639744ca2a8b149baccf648e39c8cc87553df9eae0c/urllib3-2.7.0-py3-none-any.whl", hash = "sha256:9fb4c81ebbb1ce9531cce37674bbc6f1360472bc18ca9a553ede278ef7276897", size = 131087, upload-time = "2026-05-07T16:13:17.151Z" },
|
{ url = "https://files.pythonhosted.org/packages/92/9d/c4e665119135114480843e7ab388fa94d8480650450e6f8e26b70d323a4c/urllib3-2.8.0-py3-none-any.whl", hash = "sha256:0cf3cae568d36aa9576b28dfb35f11328f1cb974ca7647d9475ebb86c75ac6e3", size = 135717, upload-time = "2026-09-15T19:29:34.577Z" },
|
||||||
]
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
|
|||||||
Generated
+3
-3
@@ -1338,11 +1338,11 @@ wheels = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "urllib3"
|
name = "urllib3"
|
||||||
version = "2.7.0"
|
version = "2.8.0"
|
||||||
source = { registry = "https://pypi.org/simple" }
|
source = { registry = "https://pypi.org/simple" }
|
||||||
sdist = { url = "https://files.pythonhosted.org/packages/53/0c/06f8b233b8fd13b9e5ee11424ef85419ba0d8ba0b3138bf360be2ff56953/urllib3-2.7.0.tar.gz", hash = "sha256:231e0ec3b63ceb14667c67be60f2f2c40a518cb38b03af60abc813da26505f4c", size = 433602, upload-time = "2026-05-07T16:13:18.596Z" }
|
sdist = { url = "https://files.pythonhosted.org/packages/e3/05/b17359e1cefb4f909b5e40b1b90a496d987258916dbbf88e842c729f510e/urllib3-2.8.0.tar.gz", hash = "sha256:63bf2ead4c879426ebf22ef2a781eeb4aa3b4ae798a0435506f8687fd5bb9b63", size = 458972, upload-time = "2026-09-15T19:29:36.253Z" }
|
||||||
wheels = [
|
wheels = [
|
||||||
{ url = "https://files.pythonhosted.org/packages/7f/3e/5db95bcf282c52709639744ca2a8b149baccf648e39c8cc87553df9eae0c/urllib3-2.7.0-py3-none-any.whl", hash = "sha256:9fb4c81ebbb1ce9531cce37674bbc6f1360472bc18ca9a553ede278ef7276897", size = 131087, upload-time = "2026-05-07T16:13:17.151Z" },
|
{ url = "https://files.pythonhosted.org/packages/92/9d/c4e665119135114480843e7ab388fa94d8480650450e6f8e26b70d323a4c/urllib3-2.8.0-py3-none-any.whl", hash = "sha256:0cf3cae568d36aa9576b28dfb35f11328f1cb974ca7647d9475ebb86c75ac6e3", size = 135717, upload-time = "2026-09-15T19:29:34.577Z" },
|
||||||
]
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
|
|||||||
Generated
+3
-3
@@ -732,11 +732,11 @@ wheels = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "urllib3"
|
name = "urllib3"
|
||||||
version = "2.7.0"
|
version = "2.8.0"
|
||||||
source = { registry = "https://pypi.org/simple" }
|
source = { registry = "https://pypi.org/simple" }
|
||||||
sdist = { url = "https://files.pythonhosted.org/packages/53/0c/06f8b233b8fd13b9e5ee11424ef85419ba0d8ba0b3138bf360be2ff56953/urllib3-2.7.0.tar.gz", hash = "sha256:231e0ec3b63ceb14667c67be60f2f2c40a518cb38b03af60abc813da26505f4c", size = 433602, upload-time = "2026-05-07T16:13:18.596Z" }
|
sdist = { url = "https://files.pythonhosted.org/packages/e3/05/b17359e1cefb4f909b5e40b1b90a496d987258916dbbf88e842c729f510e/urllib3-2.8.0.tar.gz", hash = "sha256:63bf2ead4c879426ebf22ef2a781eeb4aa3b4ae798a0435506f8687fd5bb9b63", size = 458972, upload-time = "2026-09-15T19:29:36.253Z" }
|
||||||
wheels = [
|
wheels = [
|
||||||
{ url = "https://files.pythonhosted.org/packages/7f/3e/5db95bcf282c52709639744ca2a8b149baccf648e39c8cc87553df9eae0c/urllib3-2.7.0-py3-none-any.whl", hash = "sha256:9fb4c81ebbb1ce9531cce37674bbc6f1360472bc18ca9a553ede278ef7276897", size = 131087, upload-time = "2026-05-07T16:13:17.151Z" },
|
{ url = "https://files.pythonhosted.org/packages/92/9d/c4e665119135114480843e7ab388fa94d8480650450e6f8e26b70d323a4c/urllib3-2.8.0-py3-none-any.whl", hash = "sha256:0cf3cae568d36aa9576b28dfb35f11328f1cb974ca7647d9475ebb86c75ac6e3", size = 135717, upload-time = "2026-09-15T19:29:34.577Z" },
|
||||||
]
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
|
|||||||
Generated
+3
-3
@@ -672,11 +672,11 @@ wheels = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "urllib3"
|
name = "urllib3"
|
||||||
version = "2.7.0"
|
version = "2.8.0"
|
||||||
source = { registry = "https://pypi.org/simple" }
|
source = { registry = "https://pypi.org/simple" }
|
||||||
sdist = { url = "https://files.pythonhosted.org/packages/53/0c/06f8b233b8fd13b9e5ee11424ef85419ba0d8ba0b3138bf360be2ff56953/urllib3-2.7.0.tar.gz", hash = "sha256:231e0ec3b63ceb14667c67be60f2f2c40a518cb38b03af60abc813da26505f4c", size = 433602, upload-time = "2026-05-07T16:13:18.596Z" }
|
sdist = { url = "https://files.pythonhosted.org/packages/e3/05/b17359e1cefb4f909b5e40b1b90a496d987258916dbbf88e842c729f510e/urllib3-2.8.0.tar.gz", hash = "sha256:63bf2ead4c879426ebf22ef2a781eeb4aa3b4ae798a0435506f8687fd5bb9b63", size = 458972, upload-time = "2026-09-15T19:29:36.253Z" }
|
||||||
wheels = [
|
wheels = [
|
||||||
{ url = "https://files.pythonhosted.org/packages/7f/3e/5db95bcf282c52709639744ca2a8b149baccf648e39c8cc87553df9eae0c/urllib3-2.7.0-py3-none-any.whl", hash = "sha256:9fb4c81ebbb1ce9531cce37674bbc6f1360472bc18ca9a553ede278ef7276897", size = 131087, upload-time = "2026-05-07T16:13:17.151Z" },
|
{ url = "https://files.pythonhosted.org/packages/92/9d/c4e665119135114480843e7ab388fa94d8480650450e6f8e26b70d323a4c/urllib3-2.8.0-py3-none-any.whl", hash = "sha256:0cf3cae568d36aa9576b28dfb35f11328f1cb974ca7647d9475ebb86c75ac6e3", size = 135717, upload-time = "2026-09-15T19:29:34.577Z" },
|
||||||
]
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
|
|||||||
Generated
+6
-6
@@ -1777,11 +1777,11 @@ wheels = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "pyjwt"
|
name = "pyjwt"
|
||||||
version = "2.13.0"
|
version = "2.15.0"
|
||||||
source = { registry = "https://pypi.org/simple" }
|
source = { registry = "https://pypi.org/simple" }
|
||||||
sdist = { url = "https://files.pythonhosted.org/packages/3b/81/58d0ac84e1ef3a3843791d6954d94c0b33d526c75eeb1efbce9d0a4c4077/pyjwt-2.13.0.tar.gz", hash = "sha256:41571c89ca91598c79e8ef18a2d07367d4810fbbd6f637794879baf1b7703423", size = 107515, upload-time = "2026-05-21T19:54:36.618Z" }
|
sdist = { url = "https://files.pythonhosted.org/packages/02/a5/5197bfd06417837ac079921c66fa6393f1dea3557272a263cebfef69e432/pyjwt-2.15.0.tar.gz", hash = "sha256:b11c5f9791d7bf51c2b39a81ed669f6b2dbbd669df2942f6c60167e9e3d1abe4", size = 120513, upload-time = "2026-09-23T16:56:00.689Z" }
|
||||||
wheels = [
|
wheels = [
|
||||||
{ url = "https://files.pythonhosted.org/packages/a3/5e/ecf12fdb62546d64385c158514e9b2b671f7832108ef2ecd2020ce0af2d1/pyjwt-2.13.0-py3-none-any.whl", hash = "sha256:66adcc2aff09b3f1bbd95fc1e1577df8ac8723c978552fd43304c8a290ac5728", size = 31274, upload-time = "2026-05-21T19:54:35.362Z" },
|
{ url = "https://files.pythonhosted.org/packages/e8/55/40e45bf052ee8ee12a4dfd785519660f8effa7b065442b91646ec6828619/pyjwt-2.15.0-py3-none-any.whl", hash = "sha256:7a3742debf6b879e912dbb9819ceec1594be812452b78c5f2e2dfc56564954f8", size = 33680, upload-time = "2026-09-23T16:55:59.241Z" },
|
||||||
]
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
@@ -2243,11 +2243,11 @@ wheels = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "urllib3"
|
name = "urllib3"
|
||||||
version = "2.7.0"
|
version = "2.8.0"
|
||||||
source = { registry = "https://pypi.org/simple" }
|
source = { registry = "https://pypi.org/simple" }
|
||||||
sdist = { url = "https://files.pythonhosted.org/packages/53/0c/06f8b233b8fd13b9e5ee11424ef85419ba0d8ba0b3138bf360be2ff56953/urllib3-2.7.0.tar.gz", hash = "sha256:231e0ec3b63ceb14667c67be60f2f2c40a518cb38b03af60abc813da26505f4c", size = 433602, upload-time = "2026-05-07T16:13:18.596Z" }
|
sdist = { url = "https://files.pythonhosted.org/packages/e3/05/b17359e1cefb4f909b5e40b1b90a496d987258916dbbf88e842c729f510e/urllib3-2.8.0.tar.gz", hash = "sha256:63bf2ead4c879426ebf22ef2a781eeb4aa3b4ae798a0435506f8687fd5bb9b63", size = 458972, upload-time = "2026-09-15T19:29:36.253Z" }
|
||||||
wheels = [
|
wheels = [
|
||||||
{ url = "https://files.pythonhosted.org/packages/7f/3e/5db95bcf282c52709639744ca2a8b149baccf648e39c8cc87553df9eae0c/urllib3-2.7.0-py3-none-any.whl", hash = "sha256:9fb4c81ebbb1ce9531cce37674bbc6f1360472bc18ca9a553ede278ef7276897", size = 131087, upload-time = "2026-05-07T16:13:17.151Z" },
|
{ url = "https://files.pythonhosted.org/packages/92/9d/c4e665119135114480843e7ab388fa94d8480650450e6f8e26b70d323a4c/urllib3-2.8.0-py3-none-any.whl", hash = "sha256:0cf3cae568d36aa9576b28dfb35f11328f1cb974ca7647d9475ebb86c75ac6e3", size = 135717, upload-time = "2026-09-15T19:29:34.577Z" },
|
||||||
]
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
|
|||||||
@@ -114,6 +114,26 @@ def create_metadata_for_update_state_api(
|
|||||||
return new_counters
|
return new_counters
|
||||||
|
|
||||||
|
|
||||||
|
def advance_delta_counters(
|
||||||
|
channels: Mapping[str, BaseChannel],
|
||||||
|
updated_channels: set[str],
|
||||||
|
*,
|
||||||
|
prev_metadata: Mapping[str, Any] | None,
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
"""The `counters_since_delta_snapshot` entry for an update_state
|
||||||
|
checkpoint saved one superstep after `prev_metadata`'s, for the paths
|
||||||
|
that skip `create_checkpoint_plan_for_update_state_api`.
|
||||||
|
|
||||||
|
Without it, the next checkpoint restarts every delta channel's snapshot
|
||||||
|
cadence from zero.
|
||||||
|
"""
|
||||||
|
counters = create_metadata_for_update_state_api(
|
||||||
|
channels, updated_channels, prev_metadata=prev_metadata
|
||||||
|
)
|
||||||
|
non_zero = {k: v for k, v in counters.items() if v != (0, 0)}
|
||||||
|
return {"counters_since_delta_snapshot": non_zero} if non_zero else {}
|
||||||
|
|
||||||
|
|
||||||
def create_checkpoint_plan_for_update_state_api(
|
def create_checkpoint_plan_for_update_state_api(
|
||||||
channels: Mapping[str, BaseChannel],
|
channels: Mapping[str, BaseChannel],
|
||||||
updated_channels: set[str],
|
updated_channels: set[str],
|
||||||
|
|||||||
@@ -129,6 +129,7 @@ from langgraph.pregel._algo import (
|
|||||||
from langgraph.pregel._call import identifier
|
from langgraph.pregel._call import identifier
|
||||||
from langgraph.pregel._checkpoint import (
|
from langgraph.pregel._checkpoint import (
|
||||||
achannels_from_checkpoint,
|
achannels_from_checkpoint,
|
||||||
|
advance_delta_counters,
|
||||||
channels_from_checkpoint,
|
channels_from_checkpoint,
|
||||||
copy_checkpoint,
|
copy_checkpoint,
|
||||||
create_checkpoint,
|
create_checkpoint,
|
||||||
@@ -1681,6 +1682,7 @@ class Pregel(
|
|||||||
"Cannot apply multiple updates when clearing state"
|
"Cannot apply multiple updates when clearing state"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
updated_channels: set[str] = set()
|
||||||
if saved is not None:
|
if saved is not None:
|
||||||
# tasks for this checkpoint
|
# tasks for this checkpoint
|
||||||
next_tasks = prepare_next_tasks(
|
next_tasks = prepare_next_tasks(
|
||||||
@@ -1703,7 +1705,7 @@ class Pregel(
|
|||||||
for w in saved.pending_writes or []
|
for w in saved.pending_writes or []
|
||||||
if w[0] == NULL_TASK_ID
|
if w[0] == NULL_TASK_ID
|
||||||
]:
|
]:
|
||||||
apply_writes(
|
updated_channels |= apply_writes(
|
||||||
checkpoint,
|
checkpoint,
|
||||||
channels,
|
channels,
|
||||||
[PregelTaskWrites((), INPUT, null_writes, [])],
|
[PregelTaskWrites((), INPUT, null_writes, [])],
|
||||||
@@ -1718,7 +1720,7 @@ class Pregel(
|
|||||||
continue
|
continue
|
||||||
next_tasks[tid].writes.append((k, v))
|
next_tasks[tid].writes.append((k, v))
|
||||||
# clear all current tasks
|
# clear all current tasks
|
||||||
apply_writes(
|
updated_channels |= apply_writes(
|
||||||
checkpoint,
|
checkpoint,
|
||||||
channels,
|
channels,
|
||||||
next_tasks.values(),
|
next_tasks.values(),
|
||||||
@@ -1733,6 +1735,13 @@ class Pregel(
|
|||||||
"source": "update",
|
"source": "update",
|
||||||
"step": step + 1,
|
"step": step + 1,
|
||||||
"parents": saved.metadata.get("parents", {}) if saved else {},
|
"parents": saved.metadata.get("parents", {}) if saved else {},
|
||||||
|
**(
|
||||||
|
advance_delta_counters(
|
||||||
|
channels, updated_channels, prev_metadata=saved.metadata
|
||||||
|
)
|
||||||
|
if saved
|
||||||
|
else {}
|
||||||
|
),
|
||||||
},
|
},
|
||||||
get_new_channel_versions(
|
get_new_channel_versions(
|
||||||
checkpoint_previous_versions,
|
checkpoint_previous_versions,
|
||||||
@@ -1751,7 +1760,7 @@ class Pregel(
|
|||||||
)
|
)
|
||||||
|
|
||||||
if input_writes := deque(map_input(self.input_channels, values)):
|
if input_writes := deque(map_input(self.input_channels, values)):
|
||||||
apply_writes(
|
updated_channels = apply_writes(
|
||||||
checkpoint,
|
checkpoint,
|
||||||
channels,
|
channels,
|
||||||
[PregelTaskWrites((), INPUT, input_writes, [])],
|
[PregelTaskWrites((), INPUT, input_writes, [])],
|
||||||
@@ -1774,6 +1783,15 @@ class Pregel(
|
|||||||
"parents": saved.metadata.get("parents", {})
|
"parents": saved.metadata.get("parents", {})
|
||||||
if saved
|
if saved
|
||||||
else {},
|
else {},
|
||||||
|
**(
|
||||||
|
advance_delta_counters(
|
||||||
|
channels,
|
||||||
|
updated_channels,
|
||||||
|
prev_metadata=saved.metadata,
|
||||||
|
)
|
||||||
|
if saved
|
||||||
|
else {}
|
||||||
|
),
|
||||||
},
|
},
|
||||||
get_new_channel_versions(
|
get_new_channel_versions(
|
||||||
checkpoint_previous_versions,
|
checkpoint_previous_versions,
|
||||||
@@ -1819,6 +1837,12 @@ class Pregel(
|
|||||||
"source": "fork",
|
"source": "fork",
|
||||||
"step": step + 1,
|
"step": step + 1,
|
||||||
"parents": saved.metadata.get("parents", {}),
|
"parents": saved.metadata.get("parents", {}),
|
||||||
|
# The copy has the same values and the same parent.
|
||||||
|
**{
|
||||||
|
k: v
|
||||||
|
for k, v in saved.metadata.items()
|
||||||
|
if k == "counters_since_delta_snapshot"
|
||||||
|
},
|
||||||
},
|
},
|
||||||
{},
|
{},
|
||||||
)
|
)
|
||||||
@@ -2145,6 +2169,7 @@ class Pregel(
|
|||||||
raise InvalidUpdateError(
|
raise InvalidUpdateError(
|
||||||
"Cannot apply multiple updates when clearing state"
|
"Cannot apply multiple updates when clearing state"
|
||||||
)
|
)
|
||||||
|
updated_channels: set[str] = set()
|
||||||
if saved is not None:
|
if saved is not None:
|
||||||
# tasks for this checkpoint
|
# tasks for this checkpoint
|
||||||
next_tasks = prepare_next_tasks(
|
next_tasks = prepare_next_tasks(
|
||||||
@@ -2167,7 +2192,7 @@ class Pregel(
|
|||||||
for w in saved.pending_writes or []
|
for w in saved.pending_writes or []
|
||||||
if w[0] == NULL_TASK_ID
|
if w[0] == NULL_TASK_ID
|
||||||
]:
|
]:
|
||||||
apply_writes(
|
updated_channels |= apply_writes(
|
||||||
checkpoint,
|
checkpoint,
|
||||||
channels,
|
channels,
|
||||||
[PregelTaskWrites((), INPUT, null_writes, [])],
|
[PregelTaskWrites((), INPUT, null_writes, [])],
|
||||||
@@ -2182,7 +2207,7 @@ class Pregel(
|
|||||||
continue
|
continue
|
||||||
next_tasks[tid].writes.append((k, v))
|
next_tasks[tid].writes.append((k, v))
|
||||||
# clear all current tasks
|
# clear all current tasks
|
||||||
apply_writes(
|
updated_channels |= apply_writes(
|
||||||
checkpoint,
|
checkpoint,
|
||||||
channels,
|
channels,
|
||||||
next_tasks.values(),
|
next_tasks.values(),
|
||||||
@@ -2197,6 +2222,13 @@ class Pregel(
|
|||||||
"source": "update",
|
"source": "update",
|
||||||
"step": step + 1,
|
"step": step + 1,
|
||||||
"parents": saved.metadata.get("parents", {}) if saved else {},
|
"parents": saved.metadata.get("parents", {}) if saved else {},
|
||||||
|
**(
|
||||||
|
advance_delta_counters(
|
||||||
|
channels, updated_channels, prev_metadata=saved.metadata
|
||||||
|
)
|
||||||
|
if saved
|
||||||
|
else {}
|
||||||
|
),
|
||||||
},
|
},
|
||||||
get_new_channel_versions(
|
get_new_channel_versions(
|
||||||
checkpoint_previous_versions, checkpoint["channel_versions"]
|
checkpoint_previous_versions, checkpoint["channel_versions"]
|
||||||
@@ -2214,7 +2246,7 @@ class Pregel(
|
|||||||
)
|
)
|
||||||
|
|
||||||
if input_writes := deque(map_input(self.input_channels, values)):
|
if input_writes := deque(map_input(self.input_channels, values)):
|
||||||
apply_writes(
|
updated_channels = apply_writes(
|
||||||
checkpoint,
|
checkpoint,
|
||||||
channels,
|
channels,
|
||||||
[PregelTaskWrites((), INPUT, input_writes, [])],
|
[PregelTaskWrites((), INPUT, input_writes, [])],
|
||||||
@@ -2237,6 +2269,15 @@ class Pregel(
|
|||||||
"parents": saved.metadata.get("parents", {})
|
"parents": saved.metadata.get("parents", {})
|
||||||
if saved
|
if saved
|
||||||
else {},
|
else {},
|
||||||
|
**(
|
||||||
|
advance_delta_counters(
|
||||||
|
channels,
|
||||||
|
updated_channels,
|
||||||
|
prev_metadata=saved.metadata,
|
||||||
|
)
|
||||||
|
if saved
|
||||||
|
else {}
|
||||||
|
),
|
||||||
},
|
},
|
||||||
get_new_channel_versions(
|
get_new_channel_versions(
|
||||||
checkpoint_previous_versions,
|
checkpoint_previous_versions,
|
||||||
@@ -2282,6 +2323,12 @@ class Pregel(
|
|||||||
"source": "fork",
|
"source": "fork",
|
||||||
"step": step + 1,
|
"step": step + 1,
|
||||||
"parents": saved.metadata.get("parents", {}),
|
"parents": saved.metadata.get("parents", {}),
|
||||||
|
# The copy has the same values and the same parent.
|
||||||
|
**{
|
||||||
|
k: v
|
||||||
|
for k, v in saved.metadata.items()
|
||||||
|
if k == "counters_since_delta_snapshot"
|
||||||
|
},
|
||||||
},
|
},
|
||||||
{},
|
{},
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -147,6 +147,60 @@ async def test_predicate_fires_on_supersteps_overflow() -> None:
|
|||||||
assert "x" not in result2
|
assert "x" not in result2
|
||||||
|
|
||||||
|
|
||||||
|
def _delta_counters(saver: InMemorySaver, config: Any) -> dict[str, list[int]]:
|
||||||
|
tup = saver.get_tuple(config)
|
||||||
|
assert tup is not None
|
||||||
|
counters = tup.metadata.get("counters_since_delta_snapshot") or {}
|
||||||
|
return {ch: list(c) for ch, c in counters.items()}
|
||||||
|
|
||||||
|
|
||||||
|
_UPDATE_PATHS_WITHOUT_THE_SNAPSHOT_PLAN = pytest.mark.parametrize(
|
||||||
|
("values", "as_node", "supersteps"),
|
||||||
|
[
|
||||||
|
(None, END, 1),
|
||||||
|
(None, "__copy__", 0),
|
||||||
|
({"a": []}, "__input__", 1),
|
||||||
|
],
|
||||||
|
ids=["clear as END", "copy", "update as input"],
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
@_UPDATE_PATHS_WITHOUT_THE_SNAPSHOT_PLAN
|
||||||
|
def test_update_state_path_keeps_delta_counters(
|
||||||
|
values: Any, as_node: str, supersteps: int
|
||||||
|
) -> None:
|
||||||
|
saver = InMemorySaver()
|
||||||
|
graph = _build_two_channel_graph(saver)
|
||||||
|
config = {"configurable": {"thread_id": "counters"}}
|
||||||
|
graph.invoke({"a": ["seed-a"], "b": ["seed-b"]}, config)
|
||||||
|
before = _delta_counters(saver, config)
|
||||||
|
assert set(before) == {"a", "b"}, f"both channels need live counters: {before}"
|
||||||
|
|
||||||
|
updated = graph.update_state(config, values, as_node=as_node)
|
||||||
|
|
||||||
|
assert _delta_counters(saver, updated) == {
|
||||||
|
ch: [u, s + supersteps] for ch, (u, s) in before.items()
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
@_UPDATE_PATHS_WITHOUT_THE_SNAPSHOT_PLAN
|
||||||
|
async def test_aupdate_state_path_keeps_delta_counters(
|
||||||
|
values: Any, as_node: str, supersteps: int
|
||||||
|
) -> None:
|
||||||
|
saver = InMemorySaver()
|
||||||
|
graph = _build_two_channel_graph(saver)
|
||||||
|
config = {"configurable": {"thread_id": "counters"}}
|
||||||
|
await graph.ainvoke({"a": ["seed-a"], "b": ["seed-b"]}, config)
|
||||||
|
before = _delta_counters(saver, config)
|
||||||
|
assert set(before) == {"a", "b"}, f"both channels need live counters: {before}"
|
||||||
|
|
||||||
|
updated = await graph.aupdate_state(config, values, as_node=as_node)
|
||||||
|
|
||||||
|
assert _delta_counters(saver, updated) == {
|
||||||
|
ch: [u, s + supersteps] for ch, (u, s) in before.items()
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
async def test_counter_reset_after_supersteps_snapshot() -> None:
|
async def test_counter_reset_after_supersteps_snapshot() -> None:
|
||||||
"""After the supersteps bound triggers a snapshot, the counters for
|
"""After the supersteps bound triggers a snapshot, the counters for
|
||||||
that channel reset. Verify by using a bound higher than one run's
|
that channel reset. Verify by using a bound higher than one run's
|
||||||
|
|||||||
Generated
+6
-6
@@ -2841,11 +2841,11 @@ wheels = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "pyjwt"
|
name = "pyjwt"
|
||||||
version = "2.13.0"
|
version = "2.15.1"
|
||||||
source = { registry = "https://pypi.org/simple" }
|
source = { registry = "https://pypi.org/simple" }
|
||||||
sdist = { url = "https://files.pythonhosted.org/packages/3b/81/58d0ac84e1ef3a3843791d6954d94c0b33d526c75eeb1efbce9d0a4c4077/pyjwt-2.13.0.tar.gz", hash = "sha256:41571c89ca91598c79e8ef18a2d07367d4810fbbd6f637794879baf1b7703423", size = 107515, upload-time = "2026-05-21T19:54:36.618Z" }
|
sdist = { url = "https://files.pythonhosted.org/packages/43/ea/5194e52748b0da83d71e082d75496eaec6e58f419f5e184786ded517e6a9/pyjwt-2.15.1.tar.gz", hash = "sha256:4f259e80cdfb6b3fc18a7de51fd1ef9ec79652f25019bae68975ca2468a34df8", size = 121252, upload-time = "2026-09-28T18:40:42.598Z" }
|
||||||
wheels = [
|
wheels = [
|
||||||
{ url = "https://files.pythonhosted.org/packages/a3/5e/ecf12fdb62546d64385c158514e9b2b671f7832108ef2ecd2020ce0af2d1/pyjwt-2.13.0-py3-none-any.whl", hash = "sha256:66adcc2aff09b3f1bbd95fc1e1577df8ac8723c978552fd43304c8a290ac5728", size = 31274, upload-time = "2026-05-21T19:54:35.362Z" },
|
{ url = "https://files.pythonhosted.org/packages/50/ca/44de4e75f8aadc457f0634be3b542815078ded46dca30efb960edeecad6e/pyjwt-2.15.1-py3-none-any.whl", hash = "sha256:42d59d631f7768a1028a64c7ff581a9bf7519804daf91fc5b6c56e30eec5e193", size = 33860, upload-time = "2026-09-28T18:40:41.429Z" },
|
||||||
]
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
@@ -3687,11 +3687,11 @@ wheels = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "urllib3"
|
name = "urllib3"
|
||||||
version = "2.7.0"
|
version = "2.8.0"
|
||||||
source = { registry = "https://pypi.org/simple" }
|
source = { registry = "https://pypi.org/simple" }
|
||||||
sdist = { url = "https://files.pythonhosted.org/packages/53/0c/06f8b233b8fd13b9e5ee11424ef85419ba0d8ba0b3138bf360be2ff56953/urllib3-2.7.0.tar.gz", hash = "sha256:231e0ec3b63ceb14667c67be60f2f2c40a518cb38b03af60abc813da26505f4c", size = 433602, upload-time = "2026-05-07T16:13:18.596Z" }
|
sdist = { url = "https://files.pythonhosted.org/packages/e3/05/b17359e1cefb4f909b5e40b1b90a496d987258916dbbf88e842c729f510e/urllib3-2.8.0.tar.gz", hash = "sha256:63bf2ead4c879426ebf22ef2a781eeb4aa3b4ae798a0435506f8687fd5bb9b63", size = 458972, upload-time = "2026-09-15T19:29:36.253Z" }
|
||||||
wheels = [
|
wheels = [
|
||||||
{ url = "https://files.pythonhosted.org/packages/7f/3e/5db95bcf282c52709639744ca2a8b149baccf648e39c8cc87553df9eae0c/urllib3-2.7.0-py3-none-any.whl", hash = "sha256:9fb4c81ebbb1ce9531cce37674bbc6f1360472bc18ca9a553ede278ef7276897", size = 131087, upload-time = "2026-05-07T16:13:17.151Z" },
|
{ url = "https://files.pythonhosted.org/packages/92/9d/c4e665119135114480843e7ab388fa94d8480650450e6f8e26b70d323a4c/urllib3-2.8.0-py3-none-any.whl", hash = "sha256:0cf3cae568d36aa9576b28dfb35f11328f1cb974ca7647d9475ebb86c75ac6e3", size = 135717, upload-time = "2026-09-15T19:29:34.577Z" },
|
||||||
]
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
|
|||||||
Generated
+3
-3
@@ -1363,11 +1363,11 @@ wheels = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "urllib3"
|
name = "urllib3"
|
||||||
version = "2.7.0"
|
version = "2.8.0"
|
||||||
source = { registry = "https://pypi.org/simple" }
|
source = { registry = "https://pypi.org/simple" }
|
||||||
sdist = { url = "https://files.pythonhosted.org/packages/53/0c/06f8b233b8fd13b9e5ee11424ef85419ba0d8ba0b3138bf360be2ff56953/urllib3-2.7.0.tar.gz", hash = "sha256:231e0ec3b63ceb14667c67be60f2f2c40a518cb38b03af60abc813da26505f4c", size = 433602, upload-time = "2026-05-07T16:13:18.596Z" }
|
sdist = { url = "https://files.pythonhosted.org/packages/e3/05/b17359e1cefb4f909b5e40b1b90a496d987258916dbbf88e842c729f510e/urllib3-2.8.0.tar.gz", hash = "sha256:63bf2ead4c879426ebf22ef2a781eeb4aa3b4ae798a0435506f8687fd5bb9b63", size = 458972, upload-time = "2026-09-15T19:29:36.253Z" }
|
||||||
wheels = [
|
wheels = [
|
||||||
{ url = "https://files.pythonhosted.org/packages/7f/3e/5db95bcf282c52709639744ca2a8b149baccf648e39c8cc87553df9eae0c/urllib3-2.7.0-py3-none-any.whl", hash = "sha256:9fb4c81ebbb1ce9531cce37674bbc6f1360472bc18ca9a553ede278ef7276897", size = 131087, upload-time = "2026-05-07T16:13:17.151Z" },
|
{ url = "https://files.pythonhosted.org/packages/92/9d/c4e665119135114480843e7ab388fa94d8480650450e6f8e26b70d323a4c/urllib3-2.8.0-py3-none-any.whl", hash = "sha256:0cf3cae568d36aa9576b28dfb35f11328f1cb974ca7647d9475ebb86c75ac6e3", size = 135717, upload-time = "2026-09-15T19:29:34.577Z" },
|
||||||
]
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
|
|||||||
Generated
+3
-3
@@ -1175,11 +1175,11 @@ wheels = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "urllib3"
|
name = "urllib3"
|
||||||
version = "2.7.0"
|
version = "2.8.0"
|
||||||
source = { registry = "https://pypi.org/simple" }
|
source = { registry = "https://pypi.org/simple" }
|
||||||
sdist = { url = "https://files.pythonhosted.org/packages/53/0c/06f8b233b8fd13b9e5ee11424ef85419ba0d8ba0b3138bf360be2ff56953/urllib3-2.7.0.tar.gz", hash = "sha256:231e0ec3b63ceb14667c67be60f2f2c40a518cb38b03af60abc813da26505f4c", size = 433602, upload-time = "2026-05-07T16:13:18.596Z" }
|
sdist = { url = "https://files.pythonhosted.org/packages/e3/05/b17359e1cefb4f909b5e40b1b90a496d987258916dbbf88e842c729f510e/urllib3-2.8.0.tar.gz", hash = "sha256:63bf2ead4c879426ebf22ef2a781eeb4aa3b4ae798a0435506f8687fd5bb9b63", size = 458972, upload-time = "2026-09-15T19:29:36.253Z" }
|
||||||
wheels = [
|
wheels = [
|
||||||
{ url = "https://files.pythonhosted.org/packages/7f/3e/5db95bcf282c52709639744ca2a8b149baccf648e39c8cc87553df9eae0c/urllib3-2.7.0-py3-none-any.whl", hash = "sha256:9fb4c81ebbb1ce9531cce37674bbc6f1360472bc18ca9a553ede278ef7276897", size = 131087, upload-time = "2026-05-07T16:13:17.151Z" },
|
{ url = "https://files.pythonhosted.org/packages/92/9d/c4e665119135114480843e7ab388fa94d8480650450e6f8e26b70d323a4c/urllib3-2.8.0-py3-none-any.whl", hash = "sha256:0cf3cae568d36aa9576b28dfb35f11328f1cb974ca7647d9475ebb86c75ac6e3", size = 135717, upload-time = "2026-09-15T19:29:34.577Z" },
|
||||||
]
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
|
|||||||
Reference in New Issue
Block a user