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.
This commit is contained in:
Elior Nataf Lackritz
2026-09-28 16:35:56 -04:00
parent c9d5d5b317
commit f0d0813b6c
2 changed files with 47 additions and 0 deletions
@@ -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(
@@ -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() == []