From 569f2d2d14a8fcbc08d1b5eb0245807931dfb373 Mon Sep 17 00:00:00 2001 From: Quanzheng Long Date: Thu, 7 May 2026 10:58:13 -0700 Subject: [PATCH] done --- docs/delta-channel-checkpointer-guide.md | 402 ++++++++++++++++++ libs/checkpoint-conformance/README.md | 3 + .../checkpoint/conformance/capabilities.py | 9 + .../checkpoint/conformance/spec/__init__.py | 12 + .../conformance/spec/_delta_fixtures.py | 98 +++++ .../spec/test_delta_channel_history.py | 183 ++++++++ .../spec/test_delta_channel_keepset.py | 172 ++++++++ .../spec/test_delta_channel_reconstruction.py | 166 ++++++++ .../checkpoint/conformance/validate.py | 12 + libs/checkpoint-conformance/pyproject.toml | 5 + .../langgraph/checkpoint/base/__init__.py | 98 +++++ 11 files changed, 1160 insertions(+) create mode 100644 docs/delta-channel-checkpointer-guide.md create mode 100644 libs/checkpoint-conformance/langgraph/checkpoint/conformance/spec/_delta_fixtures.py create mode 100644 libs/checkpoint-conformance/langgraph/checkpoint/conformance/spec/test_delta_channel_history.py create mode 100644 libs/checkpoint-conformance/langgraph/checkpoint/conformance/spec/test_delta_channel_keepset.py create mode 100644 libs/checkpoint-conformance/langgraph/checkpoint/conformance/spec/test_delta_channel_reconstruction.py diff --git a/docs/delta-channel-checkpointer-guide.md b/docs/delta-channel-checkpointer-guide.md new file mode 100644 index 000000000..9c8be8466 --- /dev/null +++ b/docs/delta-channel-checkpointer-guide.md @@ -0,0 +1,402 @@ +# DeltaChannel Checkpointer Guide + +> **Beta** — `DeltaChannel` and all APIs described here (`_DeltaSnapshot`, +> `get_delta_channel_history`, `get_delta_channel_keepset`, +> `delta_updates_since_snapshot` metadata) are in beta. The on-disk +> representation and method signatures may change in future releases. + +This guide is for third-party `BaseCheckpointSaver` authors who need to +support graphs using `DeltaChannel`. It covers: + +1. How `DeltaChannel` state is stored across checkpoint tables. +2. Which saver methods interact with delta channels and what the OSS + implementations look like (working references you can copy). +3. How to correctly implement cleanup methods (`prune`, `delete_for_runs`, + `copy_thread`) without silently corrupting delta history. + +--- + +## Part I — Background + +### How DeltaChannel is stored + +Unlike regular channels, `DeltaChannel` does **not** write its accumulated +value into every checkpoint. Instead: + +- `DeltaChannel.checkpoint()` always returns `MISSING` + (`libs/langgraph/langgraph/channels/delta.py`). The channel's live value + is never placed in `channel_values` by default. + +- A full **snapshot** (`_DeltaSnapshot(value)`) is written into + `channel_values[k]` only when the channel's update count reaches + `snapshot_frequency` (default 1000). This decision is made by + `delta_channels_to_snapshot()` in `libs/langgraph/langgraph/pregel/_checkpoint.py`. + +- Between snapshots, the channel's writes are stored as ordinary + `checkpoint_writes` rows (Postgres) or `writes` rows (SQLite), keyed by + `(thread_id, checkpoint_ns, checkpoint_id, task_id, idx)`. + +- Reconstruction at read time walks the parent chain backward from the + target checkpoint, collecting writes, until it finds an ancestor whose + `channel_values[ch]` is populated (the "seed"). It then replays all + collected writes through the channel's reducer on top of that seed. + +#### Storage layout per backend + +| Backend | Snapshot blob | Writes | Key structure | +|---------|--------------|--------|---------------| +| **Postgres** | `checkpoint_blobs` table; checkpoint JSON has `True` sentinel | `checkpoint_writes` table | `(thread_id, checkpoint_ns, channel, version)` for blobs | +| **SQLite** | Inline in checkpoint blob | `writes` table | `(thread_id, checkpoint_ns, checkpoint_id, task_id, idx)` | +| **InMemory** | `self.blobs` dict | `self.writes` dict | Same logical key structure | + +#### The metadata hint + +`CheckpointMetadata.delta_updates_since_snapshot` (`dict[str, int]`) tracks +how many supersteps have written to each delta channel since the last +snapshot. It resets to 0 when a snapshot fires. This is a **hint** for the +pregel loop — the source of truth for "is there a snapshot here?" is whether +`channel_values[ch]` is populated in the checkpoint. + +### The reconstruction walk + +```mermaid +flowchart LR + subgraph parentChain [Parent Chain] + direction RL + HEAD["HEAD checkpoint
(target)"] + P1["Parent 1
writes: [w4]"] + P2["Parent 2
writes: [w3]"] + P3["Parent 3
_DeltaSnapshot(seed)
writes: [w2]"] + P4["Parent 4
writes: [w1]"] + end + HEAD --> P1 + P1 --> P2 + P2 --> P3 + P3 --> P4 + + subgraph result [History result] + SEED["seed = _DeltaSnapshot(...)"] + WRITES["writes = [w3, w4]
(oldest → newest)"] + end + P3 -.-> SEED + P2 -.-> WRITES + P1 -.-> WRITES +``` + +The walk: +1. Starts at `target.parent_config` (target's own `pending_writes` are + excluded — they belong to the *next* super-step). +2. At each ancestor, collects `pending_writes` for the requested channels. +3. Stops per-channel when `channel_values[ch]` is non-empty (= seed found). +4. Returns `{"writes": [...], "seed": ...}` per channel. If no seed is + found (walk reaches root), `"seed"` is omitted — the consumer treats + this as "start empty." + +**Default implementation:** +`BaseCheckpointSaver.aget_delta_channel_history` at +`libs/checkpoint/langgraph/checkpoint/base/__init__.py` — one `aget_tuple` +call per ancestor. Correct but O(chain_length) round-trips. + +--- + +## Part II — Read & Write Surface + +These are the saver methods that interact with `DeltaChannel`. The base +class provides correct defaults; overrides exist only as performance +optimizations. + +### `aput` — accepting `_DeltaSnapshot` in `channel_values` + +Snapshot blobs reach `aput` as `_DeltaSnapshot(value)` instances inside +`checkpoint["channel_values"]`. Round-tripping them through +`JsonPlusSerializer` is sufficient (msgpack ext code `EXT_DELTA_SNAPSHOT`). +No special-casing is required for correctness. + +**Optimization (Postgres):** hoist non-primitive values into a side blob +table (`checkpoint_blobs`) and replace the JSON value with `True`. + +**References:** +- `InMemorySaver.put` — simplest, at `libs/checkpoint/langgraph/checkpoint/memory/__init__.py` +- `PostgresSaver.put` — sentinel + side-blob, at `libs/checkpoint-postgres/langgraph/checkpoint/postgres/__init__.py` +- `SqliteSaver.put` — inline, at `libs/checkpoint-sqlite/langgraph/checkpoint/sqlite/__init__.py` + +### `aput_writes` — appending pending writes + +One row per `(channel, idx)` keyed by +`(thread_id, checkpoint_ns, checkpoint_id, task_id, idx)`. No +DeltaChannel-specific logic needed — these rows are what the parent-chain +walk reads. + +### `aget_tuple` and `alist` — joining writes + +Must return `pending_writes` populated (oldest-first per task) so the +default `aget_delta_channel_history` walk can collect them from each +ancestor. + +### `aget_delta_channel_history` — performance optimization + +The base implementation works out of the box but does one `aget_tuple` per +ancestor. For long chains (up to `snapshot_frequency` = 1000 steps), this +becomes expensive. Override for performance. + +**Working references (progressively richer):** +- **InMemory** — walks dict keys directly: + `libs/checkpoint/langgraph/checkpoint/memory/__init__.py` +- **SQLite** — paged DESC scan + per-channel UNION ALL writes fetch: + `libs/checkpoint-sqlite/langgraph/checkpoint/sqlite/_delta.py` and + `libs/checkpoint-sqlite/langgraph/checkpoint/sqlite/__init__.py` +- **Postgres** — JSONB-aware two-stage SQL (stage 1: paged `checkpoint_id` + DESC with parallel `channel_values->key IS NOT NULL` checks; stage 2: + per-channel UNION ALL writes + seed blob fetch): + `libs/checkpoint-postgres/langgraph/checkpoint/postgres/base.py` and + `libs/checkpoint-postgres/langgraph/checkpoint/postgres/__init__.py` + +### `adelete_thread` — safe cleanup + +Wholesale delete of all rows for a `thread_id`. No DeltaChannel risk +because nothing survives. Copy the pattern from any OSS saver. + +--- + +## Part III — Cleanup Surface + +These optional methods (`prune`, `delete_for_runs`, `copy_thread`) have +**no OSS implementation** yet. The docstrings on `BaseCheckpointSaver` flag +the danger; this section provides concrete recipes. + +### The pitfall + +- The "latest" checkpoint is rarely a snapshot point (`snapshot_frequency` + defaults to 1000). +- A naive `keep_latest` that drops intermediate ancestors and their + `checkpoint_writes` severs the parent chain. +- The surviving head's reconstruction walk hits `parent_config = None` + early, finds no seed, and `from_checkpoint(MISSING)` produces an empty + value. **Silent data loss — no exception raised.** +- The same trap applies to `delete_for_runs` (a run's writes may be + ancestors of a live head) and `copy_thread` (head-only copy strands the + target). + +### The `get_delta_channel_keepset` helper + +`BaseCheckpointSaver` provides a helper that returns the minimum set of +`checkpoint_id`s that must survive deletion: + +```python +keep = await saver.aget_delta_channel_keepset( + config=head_config, channels=["messages", "events"], +) +# `keep` contains the head + every ancestor back to the nearest +# _DeltaSnapshot for each listed channel. Delete everything else. +``` + +Pass `channels=[]` for graphs without `DeltaChannel` — returns +`{head_id}` only. + +### Implementing `prune` + +Three strategies, ordered by recommendation: + +#### Strategy A — Compact-and-prune (requires graph access) + +1. Resolve current value via + `saver.aget_delta_channel_history(config=head, channels=delta_channels)`. +2. Fold writes through each channel's `reducer` to get the live value. +3. Rewrite `channel_values[ch] = _DeltaSnapshot(value)` on the kept head + (update the checkpoint row + blob table). +4. Now `aget_delta_channel_keepset(channels=[])` returns `{head_id}` — + delete everything else. + +This is the most space-efficient but requires access to the graph's +reducers, so it must be done at the application layer (not purely +saver-side). + +#### Strategy B — Walk-to-boundary (purely saver-side, recommended) + +```python +async def aprune(self, thread_ids, *, strategy="keep_latest"): + for tid in thread_ids: + heads = await self._list_heads(tid) # your backend code + keep: set[str] = set() + for h in heads: + keep |= await self.aget_delta_channel_keepset( + config=h, channels=DELTA_CHANNELS, + ) + all_ids = await self._list_all_ids(tid) # your backend code + to_delete = all_ids - keep + await self._delete_checkpoint_rows(tid, to_delete) + await self._delete_writes(tid, to_delete) +``` + +No rewrite needed. Preserves more rows than Strategy A but is correct and +simple. + +**Postgres SQL sketch** (alternative to the Python helper — pure SQL +recursive CTE): + +```sql +WITH RECURSIVE ancestors AS ( + SELECT checkpoint_id, parent_checkpoint_id, + (checkpoint -> 'channel_values' -> 'messages') IS NOT NULL AS has_snap + FROM checkpoints + WHERE thread_id = $1 AND checkpoint_ns = $2 AND checkpoint_id = $3 + + UNION ALL + + SELECT c.checkpoint_id, c.parent_checkpoint_id, + (c.checkpoint -> 'channel_values' -> 'messages') IS NOT NULL + FROM checkpoints c + JOIN ancestors a ON c.checkpoint_id = a.parent_checkpoint_id + WHERE NOT a.has_snap +) +SELECT checkpoint_id FROM ancestors; +-- This returns the keep-set for channel 'messages'. +``` + +#### Strategy C — Refuse + +Detect `DeltaChannel` usage (e.g. check +`delta_updates_since_snapshot` in metadata) and skip pruning those threads. +Cheapest correct option for projects that don't need pruning yet. + +### Implementing `copy_thread` + +**Default safe approach:** copy *all* rows for +`(source_thread_id, *)` to `(target_thread_id, *)` across all three +tables. Single transaction, simple `INSERT ... SELECT`. No +DeltaChannel-specific logic needed — the full chain is preserved. + +```sql +-- Postgres example +INSERT INTO checkpoints (thread_id, checkpoint_ns, checkpoint_id, ...) +SELECT $target, checkpoint_ns, checkpoint_id, ... +FROM checkpoints WHERE thread_id = $source; + +INSERT INTO checkpoint_blobs (thread_id, checkpoint_ns, channel, version, ...) +SELECT $target, checkpoint_ns, channel, version, ... +FROM checkpoint_blobs WHERE thread_id = $source; + +INSERT INTO checkpoint_writes (thread_id, checkpoint_ns, checkpoint_id, ...) +SELECT $target, checkpoint_ns, checkpoint_id, ... +FROM checkpoint_writes WHERE thread_id = $source; +``` + +**Optimization:** if you must copy only a sub-range, apply Strategy A +(compact a snapshot onto the chosen head first) then copy that single +checkpoint. + +### Implementing `delete_for_runs` + +`run_id` lives in `metadata` JSON, not as a column. Implementations need a +JSON predicate: + +```sql +-- Postgres +SELECT checkpoint_id FROM checkpoints +WHERE thread_id = $1 AND metadata->>'run_id' = ANY($2); +``` + +**Risk:** deleting a run's checkpoints/writes can strand surviving heads in +the same thread whose reconstruction walk would have traversed those rows. + +**Recipe:** +1. Identify all `(thread_id, checkpoint_id)` belonging to the run. +2. For every still-live head in those threads, compute the keep-set: + `keep |= await saver.aget_delta_channel_keepset(config=head, channels=...)` +3. Only delete run rows that fall **outside** the union of all keep-sets. +4. If a row is both "belongs to run" and "in keep-set," it must survive. + +--- + +## Part IV — Validation + +### Conformance suite + +Run the existing conformance tests for any methods you implement: + +```python +from langgraph.checkpoint.conformance import validate + +report = await validate(my_checkpointer, capabilities={ + "delta_channel_history", + "delta_channel_keepset", + "delta_channel_reconstruction", +}) +report.print_report() +assert report.passed_all_base() +``` + +The three delta capabilities test: +- **`delta_channel_history`** — walk contract (writes oldest→newest, seed + is nearest populated `channel_values[ch]`, target's writes excluded). +- **`delta_channel_keepset`** — keep-set contract (empty channels → + `{target}`, snapshot N-back → full chain, multi-channel → union). +- **`delta_channel_reconstruction`** — end-to-end round-trip: write + `_DeltaSnapshot` + writes via `aput` / `aput_writes`, read via + `aget_delta_channel_history`, reconstruct via + `DeltaChannel.from_checkpoint(seed) + replay_writes(writes)`, assert + value equality. + +### Recommended additional test + +Build a graph with `DeltaChannel(snapshot_frequency=10)`, drive it for +~25 supersteps, run your `prune` / `delete_for_runs` / `copy_thread`, then +assert reconstructed value equals the pre-cleanup value. This is the single +test that catches the silent-corruption mode (conformance alone is +necessary but not sufficient for cleanup methods). + +See existing patterns in: +- `libs/langgraph/tests/test_delta_channel_migration.py` +- `libs/langgraph/tests/test_delta_channel_exit_mode.py` + +--- + +## Part V — Future Work + +- **`DeltaChannelAwarePruneMixin`** — wraps `aget_delta_channel_keepset` + together with abstract backend hooks into ready-to-use `aprune` / + `acopy_thread` / `adelete_for_runs` bodies. Deferred because the hook + shape is opinionated and the keepset helper already factors out the + dangerous logic. + + Sketch: + ```python + class DeltaChannelAwarePruneMixin: + DELTA_CHANNELS: ClassVar[Sequence[str]] = () + + async def _alist_thread_heads(self, tid: str) -> list[RunnableConfig]: ... + async def _alist_thread_checkpoint_ids(self, tid: str) -> set[str]: ... + async def _adelete_checkpoint_rows(self, tid: str, ids: set[str]) -> None: ... + + async def aprune(self, thread_ids, *, strategy="keep_latest"): + for tid in thread_ids: + heads = await self._alist_thread_heads(tid) + keep: set[str] = set() + for h in heads: + keep |= await self.aget_delta_channel_keepset( + config=h, channels=self.DELTA_CHANNELS, + ) + all_ids = await self._alist_thread_checkpoint_ids(tid) + await self._adelete_checkpoint_rows(tid, all_ids - keep) + ``` + +- **OSS `prune` / `delete_for_runs` / `copy_thread` implementations** on + Postgres / SQLite savers. Planned as separate PRs. + +- **Pregel-side Strategy A compaction helper** — needs graph + reducers; + lives in `langgraph` package, not the saver. + +- **API stabilization** — all methods marked Beta here will move to stable + once the design proves out in production. + +--- + +## Cross-references + +- Docstring warnings on `BaseCheckpointSaver`: `prune`, `aprune`, + `delete_for_runs`, `adelete_for_runs`, `copy_thread`, `acopy_thread` + (`libs/checkpoint/langgraph/checkpoint/base/__init__.py`). +- `DeltaChannel` Beta warning + (`libs/langgraph/langgraph/channels/delta.py`). +- Conformance suite + (`libs/checkpoint-conformance/langgraph/checkpoint/conformance/`). diff --git a/libs/checkpoint-conformance/README.md b/libs/checkpoint-conformance/README.md index a963d84fc..61807e60b 100644 --- a/libs/checkpoint-conformance/README.md +++ b/libs/checkpoint-conformance/README.md @@ -63,6 +63,9 @@ The suite tests **base** capabilities (required) and **extended** capabilities ( | `delete_for_runs` | no | `adelete_for_runs` | | `copy_thread` | no | `acopy_thread` | | `prune` | no | `aprune` | +| `delta_channel_history` | no | `aget_delta_channel_history` | +| `delta_channel_keepset` | no | `aget_delta_channel_keepset` | +| `delta_channel_reconstruction` | no | `aput` | Extended capabilities are detected by checking whether the method is overridden from `BaseCheckpointSaver`. If not overridden, those tests are skipped. diff --git a/libs/checkpoint-conformance/langgraph/checkpoint/conformance/capabilities.py b/libs/checkpoint-conformance/langgraph/checkpoint/conformance/capabilities.py index 1d5a1dd93..995b3cca9 100644 --- a/libs/checkpoint-conformance/langgraph/checkpoint/conformance/capabilities.py +++ b/libs/checkpoint-conformance/langgraph/checkpoint/conformance/capabilities.py @@ -23,6 +23,9 @@ class Capability(str, Enum): DELETE_FOR_RUNS = "delete_for_runs" COPY_THREAD = "copy_thread" PRUNE = "prune" + DELTA_CHANNEL_HISTORY = "delta_channel_history" + DELTA_CHANNEL_KEEPSET = "delta_channel_keepset" + DELTA_CHANNEL_RECONSTRUCTION = "delta_channel_reconstruction" # Capabilities that every checkpointer must support. @@ -42,6 +45,9 @@ EXTENDED_CAPABILITIES = frozenset( Capability.DELETE_FOR_RUNS, Capability.COPY_THREAD, Capability.PRUNE, + Capability.DELTA_CHANNEL_HISTORY, + Capability.DELTA_CHANNEL_KEEPSET, + Capability.DELTA_CHANNEL_RECONSTRUCTION, } ) @@ -57,6 +63,9 @@ _CAPABILITY_METHOD_MAP: dict[Capability, str] = { Capability.DELETE_FOR_RUNS: "adelete_for_runs", Capability.COPY_THREAD: "acopy_thread", Capability.PRUNE: "aprune", + Capability.DELTA_CHANNEL_HISTORY: "aget_delta_channel_history", + Capability.DELTA_CHANNEL_KEEPSET: "aget_delta_channel_keepset", + Capability.DELTA_CHANNEL_RECONSTRUCTION: "aput", } diff --git a/libs/checkpoint-conformance/langgraph/checkpoint/conformance/spec/__init__.py b/libs/checkpoint-conformance/langgraph/checkpoint/conformance/spec/__init__.py index deadb8310..915964b5d 100644 --- a/libs/checkpoint-conformance/langgraph/checkpoint/conformance/spec/__init__.py +++ b/libs/checkpoint-conformance/langgraph/checkpoint/conformance/spec/__init__.py @@ -9,6 +9,15 @@ from langgraph.checkpoint.conformance.spec.test_delete_for_runs import ( from langgraph.checkpoint.conformance.spec.test_delete_thread import ( run_delete_thread_tests, ) +from langgraph.checkpoint.conformance.spec.test_delta_channel_history import ( + run_delta_channel_history_tests, +) +from langgraph.checkpoint.conformance.spec.test_delta_channel_keepset import ( + run_delta_channel_keepset_tests, +) +from langgraph.checkpoint.conformance.spec.test_delta_channel_reconstruction import ( + run_delta_channel_reconstruction_tests, +) from langgraph.checkpoint.conformance.spec.test_get_tuple import run_get_tuple_tests from langgraph.checkpoint.conformance.spec.test_list import run_list_tests from langgraph.checkpoint.conformance.spec.test_prune import run_prune_tests @@ -24,4 +33,7 @@ __all__ = [ "run_delete_for_runs_tests", "run_copy_thread_tests", "run_prune_tests", + "run_delta_channel_history_tests", + "run_delta_channel_keepset_tests", + "run_delta_channel_reconstruction_tests", ] diff --git a/libs/checkpoint-conformance/langgraph/checkpoint/conformance/spec/_delta_fixtures.py b/libs/checkpoint-conformance/langgraph/checkpoint/conformance/spec/_delta_fixtures.py new file mode 100644 index 000000000..43aec7486 --- /dev/null +++ b/libs/checkpoint-conformance/langgraph/checkpoint/conformance/spec/_delta_fixtures.py @@ -0,0 +1,98 @@ +"""Shared fixtures for delta-channel conformance tests. + +Builds a parent chain with `_DeltaSnapshot` blobs at known positions via +direct `aput` / `aput_writes` calls. No langgraph or Pregel dependency. +""" + +from __future__ import annotations + +from collections.abc import Sequence +from typing import Any +from uuid import uuid4 + +from langchain_core.runnables import RunnableConfig +from langgraph.checkpoint.base import BaseCheckpointSaver, Checkpoint +from langgraph.checkpoint.base.id import uuid6 +from langgraph.checkpoint.serde.types import _DeltaSnapshot + +from langgraph.checkpoint.conformance.test_utils import generate_metadata + + +async def build_delta_chain( + saver: BaseCheckpointSaver, + *, + thread_id: str | None = None, + checkpoint_ns: str = "", + channel: str = "messages", + snapshots_at_steps: Sequence[int] = (0,), + total_steps: int = 6, + write_value_fn: Any | None = None, +) -> list[RunnableConfig]: + """Build a parent chain with `_DeltaSnapshot` at known positions. + + Args: + saver: Checkpointer instance. + thread_id: Defaults to a random UUID. + checkpoint_ns: Namespace (default root). + channel: Channel name used for snapshots and writes. + snapshots_at_steps: Steps at which a `_DeltaSnapshot` blob is stored + in `channel_values[channel]`. Step 0 is the oldest checkpoint. + total_steps: Number of checkpoints in the chain. + write_value_fn: Callable(step) -> write value. Defaults to step index. + + Returns: + List of stored configs (oldest first), one per step. + """ + if write_value_fn is None: + + def write_value_fn(step: int) -> Any: + return step + + thread_id = thread_id or str(uuid4()) + snapshot_set = set(snapshots_at_steps) + stored: list[RunnableConfig] = [] + parent_cfg: RunnableConfig | None = None + + for step in range(total_steps): + config: RunnableConfig = { + "configurable": { + "thread_id": thread_id, + "checkpoint_ns": checkpoint_ns, + } + } + if parent_cfg: + config["configurable"]["checkpoint_id"] = parent_cfg["configurable"][ + "checkpoint_id" + ] + + channel_values: dict[str, Any] = {} + channel_versions: dict[str, int] = {} + if step in snapshot_set: + channel_values[channel] = _DeltaSnapshot( + write_value_fn(step), + ) + channel_versions[channel] = step + 1 + + cp = Checkpoint( + v=1, + id=str(uuid6(clock_seq=-1)), + ts="", + channel_values=channel_values, + channel_versions=channel_versions, + versions_seen={}, + updated_channels=None, + ) + new_versions = dict(channel_versions) + parent_cfg = await saver.aput( + config, cp, generate_metadata(step=step), new_versions + ) + stored.append(parent_cfg) + + # Write a pending write for non-snapshot steps so the walk has + # something to collect. + if step not in snapshot_set: + await saver.aput_writes( + parent_cfg, [(channel, write_value_fn(step))], str(uuid4()) + ) + + return stored diff --git a/libs/checkpoint-conformance/langgraph/checkpoint/conformance/spec/test_delta_channel_history.py b/libs/checkpoint-conformance/langgraph/checkpoint/conformance/spec/test_delta_channel_history.py new file mode 100644 index 000000000..31f889847 --- /dev/null +++ b/libs/checkpoint-conformance/langgraph/checkpoint/conformance/spec/test_delta_channel_history.py @@ -0,0 +1,183 @@ +"""DELTA_CHANNEL_HISTORY capability tests — aget_delta_channel_history contract.""" + +from __future__ import annotations + +import traceback +from collections.abc import Callable +from uuid import uuid4 + +from langgraph.checkpoint.base import BaseCheckpointSaver + +from langgraph.checkpoint.conformance.spec._delta_fixtures import build_delta_chain + + +async def test_history_returns_writes_oldest_first( + saver: BaseCheckpointSaver, +) -> None: + """Writes are returned oldest-to-newest.""" + tid = str(uuid4()) + configs = await build_delta_chain( + saver, thread_id=tid, channel="ch", snapshots_at_steps=[0], total_steps=5 + ) + head = configs[-1] + result = await saver.aget_delta_channel_history(config=head, channels=["ch"]) + writes = result["ch"]["writes"] + values = [w[2] for w in writes] + assert values == [1, 2, 3, 4], f"Expected [1,2,3,4], got {values}" + + +async def test_history_seed_is_nearest_snapshot( + saver: BaseCheckpointSaver, +) -> None: + """Seed is the value from the nearest ancestor with channel_values populated.""" + tid = str(uuid4()) + configs = await build_delta_chain( + saver, + thread_id=tid, + channel="ch", + snapshots_at_steps=[0, 3], + total_steps=6, + ) + head = configs[-1] + result = await saver.aget_delta_channel_history(config=head, channels=["ch"]) + assert "seed" in result["ch"], "Expected seed from snapshot at step 3" + seed = result["ch"]["seed"] + from langgraph.checkpoint.serde.types import _DeltaSnapshot + + actual_value = seed.value if isinstance(seed, _DeltaSnapshot) else seed + assert actual_value == 3, f"Expected seed value 3 (step 3), got {actual_value}" + writes = result["ch"]["writes"] + values = [w[2] for w in writes] + assert values == [4, 5], f"Expected [4,5], got {values}" + + +async def test_history_excludes_target_pending_writes( + saver: BaseCheckpointSaver, +) -> None: + """Target's own pending_writes are NOT included in the history.""" + tid = str(uuid4()) + configs = await build_delta_chain( + saver, thread_id=tid, channel="ch", snapshots_at_steps=[0], total_steps=3 + ) + head = configs[-1] + # Add writes directly to the head checkpoint + await saver.aput_writes(head, [("ch", "extra")], str(uuid4())) + result = await saver.aget_delta_channel_history(config=head, channels=["ch"]) + writes = result["ch"]["writes"] + values = [w[2] for w in writes] + assert "extra" not in values, f"Target's writes should be excluded, got {values}" + + +async def test_history_multi_channel( + saver: BaseCheckpointSaver, +) -> None: + """Multiple channels have independent walk termination.""" + tid = str(uuid4()) + configs: list = [] + parent_cfg = None + from langgraph.checkpoint.base import Checkpoint + from langgraph.checkpoint.base.id import uuid6 + from langgraph.checkpoint.serde.types import _DeltaSnapshot + + from langgraph.checkpoint.conformance.test_utils import generate_metadata + + for step in range(5): + config = {"configurable": {"thread_id": tid, "checkpoint_ns": ""}} + if parent_cfg: + config["configurable"]["checkpoint_id"] = parent_cfg["configurable"][ + "checkpoint_id" + ] + cv: dict = {} + cvs: dict = {} + if step == 1: + cv["a"] = _DeltaSnapshot("snap_a") + cvs["a"] = step + 1 + if step == 3: + cv["b"] = _DeltaSnapshot("snap_b") + cvs["b"] = step + 1 + cp = Checkpoint( + v=1, + id=str(uuid6(clock_seq=-1)), + ts="", + channel_values=cv, + channel_versions=cvs, + versions_seen={}, + updated_channels=None, + ) + parent_cfg = await saver.aput(config, cp, generate_metadata(step=step), cvs) + configs.append(parent_cfg) + await saver.aput_writes(parent_cfg, [("a", step), ("b", step)], str(uuid4())) + + head = configs[-1] + result = await saver.aget_delta_channel_history(config=head, channels=["a", "b"]) + a_writes = [w[2] for w in result["a"]["writes"]] + b_writes = [w[2] for w in result["b"]["writes"]] + assert a_writes == [1, 2, 3], f"Expected a writes [1,2,3], got {a_writes}" + assert b_writes == [3], f"Expected b writes [3], got {b_writes}" + + +async def test_history_empty_channels_returns_empty( + saver: BaseCheckpointSaver, +) -> None: + """Empty channels list returns empty mapping.""" + tid = str(uuid4()) + configs = await build_delta_chain( + saver, thread_id=tid, channel="ch", snapshots_at_steps=[0], total_steps=3 + ) + result = await saver.aget_delta_channel_history(config=configs[-1], channels=[]) + assert result == {} + + +async def test_history_walk_to_root_no_seed( + saver: BaseCheckpointSaver, +) -> None: + """Walk reaches root without finding seed — no 'seed' key in result.""" + tid = str(uuid4()) + configs = await build_delta_chain( + saver, + thread_id=tid, + channel="ch", + snapshots_at_steps=[], + total_steps=4, + ) + head = configs[-1] + result = await saver.aget_delta_channel_history(config=head, channels=["ch"]) + assert "seed" not in result["ch"], f"Expected no seed, got {result['ch']}" + + +ALL_DELTA_CHANNEL_HISTORY_TESTS = [ + test_history_returns_writes_oldest_first, + test_history_seed_is_nearest_snapshot, + test_history_excludes_target_pending_writes, + test_history_multi_channel, + test_history_empty_channels_returns_empty, + test_history_walk_to_root_no_seed, +] + + +async def run_delta_channel_history_tests( + saver: BaseCheckpointSaver, + on_test_result: Callable[[str, str, bool, str | None], None] | None = None, +) -> tuple[int, int, list[str]]: + """Run all delta_channel_history tests. Returns (passed, failed, failure_names).""" + passed = 0 + failed = 0 + failures: list[str] = [] + for test_fn in ALL_DELTA_CHANNEL_HISTORY_TESTS: + try: + await test_fn(saver) + passed += 1 + if on_test_result: + on_test_result("delta_channel_history", test_fn.__name__, True, None) + except Exception: + failed += 1 + msg = f"{test_fn.__name__}: {traceback.format_exc()}" + failures.append(msg) + if on_test_result: + on_test_result( + "delta_channel_history", + test_fn.__name__, + False, + traceback.format_exc(), + ) + return passed, failed, failures diff --git a/libs/checkpoint-conformance/langgraph/checkpoint/conformance/spec/test_delta_channel_keepset.py b/libs/checkpoint-conformance/langgraph/checkpoint/conformance/spec/test_delta_channel_keepset.py new file mode 100644 index 000000000..2311f04d9 --- /dev/null +++ b/libs/checkpoint-conformance/langgraph/checkpoint/conformance/spec/test_delta_channel_keepset.py @@ -0,0 +1,172 @@ +"""DELTA_CHANNEL_KEEPSET capability tests — aget_delta_channel_keepset contract.""" + +from __future__ import annotations + +import traceback +from collections.abc import Callable +from uuid import uuid4 + +from langgraph.checkpoint.base import BaseCheckpointSaver + +from langgraph.checkpoint.conformance.spec._delta_fixtures import build_delta_chain + + +async def test_keepset_empty_channels_returns_target_only( + saver: BaseCheckpointSaver, +) -> None: + """Empty channels → keep-set is just {target_id}.""" + tid = str(uuid4()) + configs = await build_delta_chain( + saver, thread_id=tid, channel="ch", snapshots_at_steps=[0], total_steps=4 + ) + head = configs[-1] + keep = await saver.aget_delta_channel_keepset(config=head, channels=[]) + head_id = head["configurable"]["checkpoint_id"] + assert keep == {head_id}, f"Expected only target, got {keep}" + + +async def test_keepset_snapshot_at_target( + saver: BaseCheckpointSaver, +) -> None: + """When target itself has a snapshot, keep-set is just {target_id}.""" + tid = str(uuid4()) + configs = await build_delta_chain( + saver, + thread_id=tid, + channel="ch", + snapshots_at_steps=[0, 3], + total_steps=4, + ) + head = configs[3] + keep = await saver.aget_delta_channel_keepset(config=head, channels=["ch"]) + head_id = head["configurable"]["checkpoint_id"] + assert keep == {head_id}, f"Snapshot at target should yield only target, got {keep}" + + +async def test_keepset_snapshot_n_back( + saver: BaseCheckpointSaver, +) -> None: + """Snapshot N steps back → target + intermediates + snapshot ancestor.""" + tid = str(uuid4()) + configs = await build_delta_chain( + saver, thread_id=tid, channel="ch", snapshots_at_steps=[0], total_steps=5 + ) + head = configs[-1] + keep = await saver.aget_delta_channel_keepset(config=head, channels=["ch"]) + expected_ids = {c["configurable"]["checkpoint_id"] for c in configs} + assert keep == expected_ids, f"Expected all ancestors, got {keep}" + + +async def test_keepset_multi_channel_union( + saver: BaseCheckpointSaver, +) -> None: + """Multi-channel keep-set is the union (max chain per channel).""" + tid = str(uuid4()) + from langgraph.checkpoint.base import Checkpoint + from langgraph.checkpoint.base.id import uuid6 + from langgraph.checkpoint.serde.types import _DeltaSnapshot + + from langgraph.checkpoint.conformance.test_utils import generate_metadata + + configs: list = [] + parent_cfg = None + for step in range(6): + config = {"configurable": {"thread_id": tid, "checkpoint_ns": ""}} + if parent_cfg: + config["configurable"]["checkpoint_id"] = parent_cfg["configurable"][ + "checkpoint_id" + ] + cv: dict = {} + cvs: dict = {} + # Channel "a" has snapshot at step 3 (recent) + if step == 3: + cv["a"] = _DeltaSnapshot("snap_a") + cvs["a"] = step + 1 + # Channel "b" has snapshot at step 1 (further back) + if step == 1: + cv["b"] = _DeltaSnapshot("snap_b") + cvs["b"] = step + 1 + cp = Checkpoint( + v=1, + id=str(uuid6(clock_seq=-1)), + ts="", + channel_values=cv, + channel_versions=cvs, + versions_seen={}, + updated_channels=None, + ) + parent_cfg = await saver.aput(config, cp, generate_metadata(step=step), cvs) + configs.append(parent_cfg) + + head = configs[-1] + keep = await saver.aget_delta_channel_keepset(config=head, channels=["a", "b"]) + # Union: b needs back to step 1, so steps 1..5 (all except step 0) plus head + expected_ids = {c["configurable"]["checkpoint_id"] for c in configs[1:]} + expected_ids.add(head["configurable"]["checkpoint_id"]) + assert keep == expected_ids, f"Expected union, got {keep} vs {expected_ids}" + + +async def test_keepset_walk_to_root( + saver: BaseCheckpointSaver, +) -> None: + """No snapshot anywhere → entire chain to root is in keep-set.""" + tid = str(uuid4()) + configs = await build_delta_chain( + saver, thread_id=tid, channel="ch", snapshots_at_steps=[], total_steps=4 + ) + head = configs[-1] + keep = await saver.aget_delta_channel_keepset(config=head, channels=["ch"]) + all_ids = {c["configurable"]["checkpoint_id"] for c in configs} + assert keep == all_ids, f"Expected full chain, got {keep}" + + +async def test_keepset_deterministic( + saver: BaseCheckpointSaver, +) -> None: + """Same inputs return identical sets.""" + tid = str(uuid4()) + configs = await build_delta_chain( + saver, thread_id=tid, channel="ch", snapshots_at_steps=[0], total_steps=5 + ) + head = configs[-1] + keep1 = await saver.aget_delta_channel_keepset(config=head, channels=["ch"]) + keep2 = await saver.aget_delta_channel_keepset(config=head, channels=["ch"]) + assert keep1 == keep2 + + +ALL_DELTA_CHANNEL_KEEPSET_TESTS = [ + test_keepset_empty_channels_returns_target_only, + test_keepset_snapshot_at_target, + test_keepset_snapshot_n_back, + test_keepset_multi_channel_union, + test_keepset_walk_to_root, + test_keepset_deterministic, +] + + +async def run_delta_channel_keepset_tests( + saver: BaseCheckpointSaver, + on_test_result: Callable[[str, str, bool, str | None], None] | None = None, +) -> tuple[int, int, list[str]]: + """Run all delta_channel_keepset tests. Returns (passed, failed, failure_names).""" + passed = 0 + failed = 0 + failures: list[str] = [] + for test_fn in ALL_DELTA_CHANNEL_KEEPSET_TESTS: + try: + await test_fn(saver) + passed += 1 + if on_test_result: + on_test_result("delta_channel_keepset", test_fn.__name__, True, None) + except Exception: + failed += 1 + msg = f"{test_fn.__name__}: {traceback.format_exc()}" + failures.append(msg) + if on_test_result: + on_test_result( + "delta_channel_keepset", + test_fn.__name__, + False, + traceback.format_exc(), + ) + return passed, failed, failures diff --git a/libs/checkpoint-conformance/langgraph/checkpoint/conformance/spec/test_delta_channel_reconstruction.py b/libs/checkpoint-conformance/langgraph/checkpoint/conformance/spec/test_delta_channel_reconstruction.py new file mode 100644 index 000000000..f95ef60c9 --- /dev/null +++ b/libs/checkpoint-conformance/langgraph/checkpoint/conformance/spec/test_delta_channel_reconstruction.py @@ -0,0 +1,166 @@ +"""DELTA_CHANNEL_RECONSTRUCTION capability tests — end-to-end round-trip. + +Exercises: aput + aput_writes + aget_delta_channel_history + from_checkpoint + +replay_writes. This catches the most common silent-corruption mode: failing to +round-trip `_DeltaSnapshot` blobs through serialization. +""" + +from __future__ import annotations + +import traceback +from collections.abc import Callable +from uuid import uuid4 + +from langgraph.checkpoint.base import BaseCheckpointSaver + +from langgraph.checkpoint.conformance.spec._delta_fixtures import build_delta_chain + + +def _list_reducer(state: list, writes: list) -> list: + """Simple append reducer for testing.""" + return state + writes + + +async def test_reconstruction_basic( + saver: BaseCheckpointSaver, +) -> None: + """Reconstruct delta channel value from history matches expected.""" + from langgraph.channels.delta import DeltaChannel + + tid = str(uuid4()) + configs = await build_delta_chain( + saver, + thread_id=tid, + channel="msgs", + snapshots_at_steps=[0], + total_steps=5, + write_value_fn=lambda step: [step], + ) + head = configs[-1] + result = await saver.aget_delta_channel_history(config=head, channels=["msgs"]) + history = result["msgs"] + + seed = history.get("seed") + from langgraph._internal._typing import MISSING + + if seed is None: + seed = MISSING + + ch = DeltaChannel(_list_reducer, list) + replay_ch = ch.from_checkpoint(seed) + replay_ch.replay_writes(history["writes"]) + reconstructed = replay_ch.get() + + # Expected: snapshot at step 0 = [0], then writes at steps 1,2,3,4 + expected = [0] + [1] + [2] + [3] + [4] + assert reconstructed == expected, ( + f"Reconstructed {reconstructed} != expected {expected}" + ) + + +async def test_reconstruction_mid_chain_snapshot( + saver: BaseCheckpointSaver, +) -> None: + """Reconstruction works when snapshot is mid-chain.""" + from langgraph.channels.delta import DeltaChannel + + tid = str(uuid4()) + configs = await build_delta_chain( + saver, + thread_id=tid, + channel="msgs", + snapshots_at_steps=[0, 3], + total_steps=6, + write_value_fn=lambda step: [step], + ) + head = configs[-1] + result = await saver.aget_delta_channel_history(config=head, channels=["msgs"]) + history = result["msgs"] + + seed = history.get("seed") + from langgraph._internal._typing import MISSING + + if seed is None: + seed = MISSING + + ch = DeltaChannel(_list_reducer, list) + replay_ch = ch.from_checkpoint(seed) + replay_ch.replay_writes(history["writes"]) + reconstructed = replay_ch.get() + + # Snapshot at step 3 = [3], writes at steps 4,5 + expected = [3] + [4] + [5] + assert reconstructed == expected, ( + f"Reconstructed {reconstructed} != expected {expected}" + ) + + +async def test_reconstruction_no_snapshot( + saver: BaseCheckpointSaver, +) -> None: + """Reconstruction from root (no snapshot) gives all writes accumulated.""" + from langgraph.channels.delta import DeltaChannel + + tid = str(uuid4()) + configs = await build_delta_chain( + saver, + thread_id=tid, + channel="msgs", + snapshots_at_steps=[], + total_steps=4, + write_value_fn=lambda step: [step], + ) + head = configs[-1] + result = await saver.aget_delta_channel_history(config=head, channels=["msgs"]) + history = result["msgs"] + + from langgraph._internal._typing import MISSING + + seed = history.get("seed", MISSING) + + ch = DeltaChannel(_list_reducer, list) + replay_ch = ch.from_checkpoint(seed) + replay_ch.replay_writes(history["writes"]) + reconstructed = replay_ch.get() + + expected = [0] + [1] + [2] + [3] + assert reconstructed == expected, ( + f"Reconstructed {reconstructed} != expected {expected}" + ) + + +ALL_DELTA_CHANNEL_RECONSTRUCTION_TESTS = [ + test_reconstruction_basic, + test_reconstruction_mid_chain_snapshot, + test_reconstruction_no_snapshot, +] + + +async def run_delta_channel_reconstruction_tests( + saver: BaseCheckpointSaver, + on_test_result: Callable[[str, str, bool, str | None], None] | None = None, +) -> tuple[int, int, list[str]]: + """Run all reconstruction tests. Returns (passed, failed, failure_names).""" + passed = 0 + failed = 0 + failures: list[str] = [] + for test_fn in ALL_DELTA_CHANNEL_RECONSTRUCTION_TESTS: + try: + await test_fn(saver) + passed += 1 + if on_test_result: + on_test_result( + "delta_channel_reconstruction", test_fn.__name__, True, None + ) + except Exception: + failed += 1 + msg = f"{test_fn.__name__}: {traceback.format_exc()}" + failures.append(msg) + if on_test_result: + on_test_result( + "delta_channel_reconstruction", + test_fn.__name__, + False, + traceback.format_exc(), + ) + return passed, failed, failures diff --git a/libs/checkpoint-conformance/langgraph/checkpoint/conformance/validate.py b/libs/checkpoint-conformance/langgraph/checkpoint/conformance/validate.py index 34e81d89c..e1019e6a6 100644 --- a/libs/checkpoint-conformance/langgraph/checkpoint/conformance/validate.py +++ b/libs/checkpoint-conformance/langgraph/checkpoint/conformance/validate.py @@ -19,6 +19,15 @@ from langgraph.checkpoint.conformance.spec.test_delete_for_runs import ( from langgraph.checkpoint.conformance.spec.test_delete_thread import ( run_delete_thread_tests, ) +from langgraph.checkpoint.conformance.spec.test_delta_channel_history import ( + run_delta_channel_history_tests, +) +from langgraph.checkpoint.conformance.spec.test_delta_channel_keepset import ( + run_delta_channel_keepset_tests, +) +from langgraph.checkpoint.conformance.spec.test_delta_channel_reconstruction import ( + run_delta_channel_reconstruction_tests, +) from langgraph.checkpoint.conformance.spec.test_get_tuple import run_get_tuple_tests from langgraph.checkpoint.conformance.spec.test_list import run_list_tests from langgraph.checkpoint.conformance.spec.test_prune import run_prune_tests @@ -35,6 +44,9 @@ _RUNNERS = { Capability.DELETE_FOR_RUNS: run_delete_for_runs_tests, Capability.COPY_THREAD: run_copy_thread_tests, Capability.PRUNE: run_prune_tests, + Capability.DELTA_CHANNEL_HISTORY: run_delta_channel_history_tests, + Capability.DELTA_CHANNEL_KEEPSET: run_delta_channel_keepset_tests, + Capability.DELTA_CHANNEL_RECONSTRUCTION: run_delta_channel_reconstruction_tests, } diff --git a/libs/checkpoint-conformance/pyproject.toml b/libs/checkpoint-conformance/pyproject.toml index dac2f550c..9661f9bc9 100644 --- a/libs/checkpoint-conformance/pyproject.toml +++ b/libs/checkpoint-conformance/pyproject.toml @@ -43,9 +43,14 @@ asyncio_mode = "auto" # The extended methods (acopy_thread, adelete_for_runs, aprune) are checked # at runtime via capability detection and may not exist on the installed # base class. Dict literal inference is also overly strict for RunnableConfig. +# Delta-channel tests import from `langgraph` (not a declared dep of this +# package — at test time it is installed alongside); private `_DeltaSnapshot` +# imports are intentional (beta surface). unresolved-attribute = "ignore" +unresolved-import = "ignore" invalid-argument-type = "ignore" invalid-return-type = "ignore" +missing-typed-dict-key = "ignore" [tool.ruff] lint.select = [ diff --git a/libs/checkpoint/langgraph/checkpoint/base/__init__.py b/libs/checkpoint/langgraph/checkpoint/base/__init__.py index 7b727bc1e..b3bd610ba 100644 --- a/libs/checkpoint/langgraph/checkpoint/base/__init__.py +++ b/libs/checkpoint/langgraph/checkpoint/base/__init__.py @@ -680,6 +680,104 @@ class BaseCheckpointSaver(Generic[V]): result[ch] = entry return result + def get_delta_channel_keepset( + self, + *, + config: RunnableConfig, + channels: Sequence[str], + ) -> set[str]: + """Return ancestor checkpoint_ids that must survive deletion. + + !!! warning "Beta" + + This method is part of the `DeltaChannel` support surface and is + in beta. The signature may change while the delta-channel design + stabilizes. + + Walks the parent chain from `config` backward, collecting visited + checkpoint_ids (inclusive of the target), and terminates per-channel + when that channel has a populated `channel_values[ch]` (a + `_DeltaSnapshot` blob or a pre-migration plain value). The returned + set is the minimum keep-set: every checkpoint_id whose removal would + break reconstruction of the listed channels at `config`. + + Pass `channels=[]` to return just `{config.checkpoint_id}` — useful + for graphs that don't use `DeltaChannel`. + + Compose this into custom `prune` / `delete_for_runs` / `copy_thread`:: + + keep = saver.get_delta_channel_keepset( + config=head_config, channels=delta_channels, + ) + delete_rows_not_in(keep) + + Args: + config: Configuration identifying the target checkpoint. + channels: Channel names whose delta history must be preserved. + Empty sequence means only the target checkpoint_id is kept. + + Returns: + Set of checkpoint_ids that must not be deleted. + """ + target_tuple = self.get_tuple(config) + if target_tuple is None: + return set() + target_id = target_tuple.config["configurable"]["checkpoint_id"] + keep: set[str] = {target_id} + if not channels: + return keep + remaining: set[str] = set(channels) + cursor_config: RunnableConfig | None = target_tuple.parent_config + while cursor_config is not None and remaining: + tup = self.get_tuple(cursor_config) + if tup is None: + break + cid = tup.config["configurable"]["checkpoint_id"] + keep.add(cid) + for ch in list(remaining): + if ch in tup.checkpoint["channel_values"]: + remaining.discard(ch) + if not remaining: + break + cursor_config = tup.parent_config + return keep + + async def aget_delta_channel_keepset( + self, + *, + config: RunnableConfig, + channels: Sequence[str], + ) -> set[str]: + """Async version of `get_delta_channel_keepset`. + + !!! warning "Beta" + + This method is part of the `DeltaChannel` support surface and is + in beta. See `get_delta_channel_keepset` for full documentation. + """ + target_tuple = await self.aget_tuple(config) + if target_tuple is None: + return set() + target_id = target_tuple.config["configurable"]["checkpoint_id"] + keep: set[str] = {target_id} + if not channels: + return keep + remaining: set[str] = set(channels) + cursor_config: RunnableConfig | None = target_tuple.parent_config + while cursor_config is not None and remaining: + tup = await self.aget_tuple(cursor_config) + if tup is None: + break + cid = tup.config["configurable"]["checkpoint_id"] + keep.add(cid) + for ch in list(remaining): + if ch in tup.checkpoint["channel_values"]: + remaining.discard(ch) + if not remaining: + break + cursor_config = tup.parent_config + return keep + def get_next_version(self, current: V | None, channel: None) -> V: """Generate the next version ID for a channel.