From 0207ec7cffe533e1dd8d17de61eb75cb1815a5d7 Mon Sep 17 00:00:00 2001 From: Elior Nataf Lackritz Date: Mon, 28 Sep 2026 16:35:56 -0400 Subject: [PATCH] fix(langgraph): raise when a written DeltaChannel is hydrated without a saver A missing saver used to hydrate the channel as empty, silently. Raise when a DeltaChannel has a version but no stored value and no saver was passed. Keying on the version keeps a never-written channel, which has no value either, reading as empty without a saver. --- .../langgraph/langgraph/pregel/_checkpoint.py | 15 +++++++++ .../tests/test_delta_channel_subgraph.py | 32 +++++++++++++++++++ 2 files changed, 47 insertions(+) diff --git a/libs/langgraph/langgraph/pregel/_checkpoint.py b/libs/langgraph/langgraph/pregel/_checkpoint.py index c336f75a6..534c2dcb7 100644 --- a/libs/langgraph/langgraph/pregel/_checkpoint.py +++ b/libs/langgraph/langgraph/pregel/_checkpoint.py @@ -226,6 +226,19 @@ def _needs_replay(spec: BaseChannel, stored: object) -> bool: return stored is MISSING +def _require_saver_for_history( + checkpoint: Checkpoint, + delta_channels: list[str], + saver: BaseCheckpointSaver | None, +) -> None: + written = [k for k in delta_channels if k in checkpoint["channel_versions"]] + if written and saver is None: + raise ValueError( + f"DeltaChannel {written} has history to replay but no checkpointer " + "was passed to read it" + ) + + def channels_from_checkpoint( specs: Mapping[str, BaseChannel | ManagedValueSpec], checkpoint: Checkpoint, @@ -256,6 +269,7 @@ def channels_from_checkpoint( for k, spec in channel_specs.items() if _needs_replay(spec, checkpoint["channel_values"].get(k, MISSING)) ] + _require_saver_for_history(checkpoint, delta_channels, saver) histories: Mapping[str, Any] = {} if delta_channels and saver is not None and config is not None: histories = saver.get_delta_channel_history( @@ -298,6 +312,7 @@ async def achannels_from_checkpoint( for k, spec in channel_specs.items() if _needs_replay(spec, checkpoint["channel_values"].get(k, MISSING)) ] + _require_saver_for_history(checkpoint, delta_channels, saver) histories: Mapping[str, Any] = {} if delta_channels and saver is not None and config is not None: histories = await saver.aget_delta_channel_history( diff --git a/libs/langgraph/tests/test_delta_channel_subgraph.py b/libs/langgraph/tests/test_delta_channel_subgraph.py index a795ff6c3..c873ae759 100644 --- a/libs/langgraph/tests/test_delta_channel_subgraph.py +++ b/libs/langgraph/tests/test_delta_channel_subgraph.py @@ -7,6 +7,11 @@ from typing_extensions import TypedDict from langgraph.channels.delta import DeltaChannel from langgraph.graph import END, START, StateGraph +from langgraph.pregel._checkpoint import ( + achannels_from_checkpoint, + channels_from_checkpoint, + empty_checkpoint, +) pytestmark = pytest.mark.anyio @@ -251,3 +256,30 @@ def test_completed_subgraph_exposes_no_task_state( app.invoke({}, config) assert app.get_state(config, subgraphs=True).tasks == () + + +def _written_delta_checkpoint() -> Any: + checkpoint = empty_checkpoint() + checkpoint["channel_versions"]["delta"] = 1 + return checkpoint + + +def test_hydrating_written_delta_channel_without_saver_raises() -> None: + with pytest.raises(ValueError, match="no checkpointer"): + channels_from_checkpoint( + {"delta": DeltaChannel(_extend)}, _written_delta_checkpoint() + ) + + +async def test_ahydrating_written_delta_channel_without_saver_raises() -> None: + with pytest.raises(ValueError, match="no checkpointer"): + await achannels_from_checkpoint( + {"delta": DeltaChannel(_extend)}, _written_delta_checkpoint() + ) + + +def test_hydrating_unwritten_delta_channel_without_saver_is_empty() -> None: + channels, _ = channels_from_checkpoint( + {"delta": DeltaChannel(_extend)}, empty_checkpoint() + ) + assert channels["delta"].get() == []