mirror of
https://github.com/langchain-ai/langgraph.git
synced 2026-09-08 10:47:52 +02:00
cr
This commit is contained in:
@@ -149,15 +149,15 @@ class InMemorySaver(
|
||||
they find their own terminator or hit the root.
|
||||
|
||||
The seed value (whether a `_DeltaSnapshot` or a plain pre-delta
|
||||
migration blob) is the value AT that ancestor, prior to its own
|
||||
pending writes that produce the child. Those on-path writes —
|
||||
including the ones stored on the terminating ancestor — are always
|
||||
collected and replayed on top of the seed, so a thread migrated from
|
||||
a pre-delta channel does not drop the writes saved under the
|
||||
migration boundary checkpoint.
|
||||
migration blob) is the value AT that ancestor. Pending writes on a
|
||||
plain seed checkpoint are replayed when they are the first deltas
|
||||
after migration (no delta-era ancestor already contributed writes).
|
||||
Stale pre-delta writes on the same checkpoint as the blob are skipped
|
||||
once a delta-era ancestor in the walk has already contributed writes.
|
||||
"""
|
||||
if not channels:
|
||||
return {}
|
||||
from langgraph.checkpoint.serde.types import _DeltaSnapshot
|
||||
|
||||
thread_id = config["configurable"]["thread_id"]
|
||||
checkpoint_ns = config["configurable"].get("checkpoint_ns", "")
|
||||
@@ -178,6 +178,8 @@ class InMemorySaver(
|
||||
collected_by_ch: dict[str, list[PendingWrite]] = {c: [] for c in channels}
|
||||
seed_by_ch: dict[str, Any] = {}
|
||||
remaining: set[str] = set(channels)
|
||||
# True once a delta-era (empty-blob) ancestor has contributed a write.
|
||||
seen_delta_write: dict[str, bool] = dict.fromkeys(channels, False)
|
||||
|
||||
for cp_id in chain:
|
||||
if not remaining:
|
||||
@@ -187,14 +189,17 @@ class InMemorySaver(
|
||||
|
||||
terminated_here: set[str] = set()
|
||||
blob_value_by_ch: dict[str, Any] = {}
|
||||
delta_era_by_ch: dict[str, bool] = {}
|
||||
if ckpt is not None:
|
||||
versions = ckpt.get("channel_versions", {})
|
||||
for ch in remaining:
|
||||
ver = versions.get(ch)
|
||||
if ver is None:
|
||||
delta_era_by_ch[ch] = True
|
||||
continue
|
||||
blob_entry = self.blobs.get((thread_id, checkpoint_ns, ch, ver))
|
||||
if blob_entry is None or blob_entry[0] == "empty":
|
||||
delta_era_by_ch[ch] = True
|
||||
continue
|
||||
blob_value_by_ch[ch] = self.serde.loads_typed(blob_entry)
|
||||
terminated_here.add(ch)
|
||||
@@ -205,15 +210,22 @@ class InMemorySaver(
|
||||
):
|
||||
if ch not in remaining:
|
||||
continue
|
||||
# Collect on-path writes regardless of seed type. A plain
|
||||
# (pre-delta migration) blob is the settled value AT that
|
||||
# ancestor; its own pending writes produce the child and must
|
||||
# still be replayed, just like a `_DeltaSnapshot` seed.
|
||||
# Skipping them would drop post-migration writes saved under
|
||||
# the migration boundary checkpoint.
|
||||
blob_value = blob_value_by_ch.get(ch)
|
||||
if (
|
||||
blob_value is not None
|
||||
and not isinstance(blob_value, _DeltaSnapshot)
|
||||
and seen_delta_write.get(ch, False)
|
||||
):
|
||||
# Pre-delta blob already settled; stale writes on the same
|
||||
# checkpoint are subsumed. Post-migration writes on the
|
||||
# migration boundary are collected when no delta-era
|
||||
# ancestor has contributed writes yet.
|
||||
continue
|
||||
collected_by_ch[ch].append(
|
||||
(tid, ch, self.serde.loads_typed(serialized))
|
||||
)
|
||||
if delta_era_by_ch.get(ch, False):
|
||||
seen_delta_write[ch] = True
|
||||
|
||||
for ch in terminated_here:
|
||||
seed_by_ch[ch] = blob_value_by_ch[ch]
|
||||
|
||||
Reference in New Issue
Block a user