mirror of
https://github.com/langchain-ai/langgraph.git
synced 2026-09-30 13:35:09 +02:00
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:
@@ -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() == []
|
||||
|
||||
Reference in New Issue
Block a user