deer-flow/backend/tests/test_context_compaction.py
Vanzeren 42baed8c8c
feat(checkpoint): dual-mode checkpoint storage with LangGraph DeltaChannel (#4292)
* feat(checkpoint): dual-mode checkpoint storage with LangGraph DeltaChannel

Add a restart-required database.checkpoint_channel_mode ("full" default,
"delta") that stores the messages channel via LangGraph 1.2 DeltaChannel,
cutting checkpoint storage from O(n^2) to O(n) for append-only history.
Existing full checkpoints seed delta state transparently; no data migration.

- config: mode schema + freeze-on-first-use with
  CheckpointModeReconfigurationError; mode marker persisted in checkpoint
  metadata; unsafe delta->full downgrade rejected fail-closed with
  CheckpointModeMismatchError (run-level error, failed state read)
- state: delta message state schema; CheckpointStateAccessor centralizes
  materialized reads for all consumers (threads API, branches,
  regeneration, compaction, state updates, memory, goal workers)
- runtime: raw writers (run durations, interrupted title, thread goal)
  parent their checkpoints to the checkpoint they derive from, preserving
  delta ancestry; rollback forks the pre-run lineage through a state
  mutation graph with Overwrite restores; InMemorySaver delta-history
  override delegates to the base walk (fixes dropped first write after
  migration, also present upstream)
- tests: conformance suite over {memory, sqlite, postgres} covering
  migration replay, stable message IDs, storage shape and writer
  preservation; conftest fixture isolates the frozen mode between tests;
  stale config fakes refreshed
- ci: backend unit tests gain a postgres service

* fix(checkpoint): close materialization gaps in goal flow, guard public factory

- Route goal-continuation message reads through CheckpointStateAccessor:
  raw channel_values reads see the delta sentinel in delta mode, which
  disabled goal continuation (stand_down=no_durable_end_of_turn) after
  durable assistant turns. Raw tuples remain for tuple-only metadata
  (checkpoint id, pending_writes).
- Reject checkpoint_channel_mode='delta' + checkpointer in
  create_deerflow_agent at construction: factory-built persisted graphs
  bypass mode-marker injection and the fail-closed gate, reproducing
  silent mixed-mode state loss. Delta without persistence stays allowed.
- Import the postgres saver lazily (pytest.importorskip in the fixture)
  so the documented default install collects the suite; add a CI job
  running pytest --collect-only on uv sync --group dev without extras.
- Fix test_checkpointer fallback test to patch get_app_config at its
  use site (provider module), making it deterministic when a local
  config.yaml selects a persistent backend.

* fix(gateway): preserve extension-owned channels in state mutations, bump config version

- build_state_mutation_graph / build_checkpoint_state_mutation_accessor
  accept an explicit state_schema; branch and POST /state now compile the
  mutation graph from the thread's effective schema
  (graph_state_schema on the assistant graph). The base-ThreadState
  fallback silently discarded channels contributed by custom
  AgentMiddleware.state_schema on branch (data loss) and returned a
  false-success 200 on POST /state.
- POST /state validates values keys against the mutation graph's
  channels and rejects unknown fields with 422 instead of ignoring
  them; reducer detection covers extension channels
  (BinaryOperatorAggregate or DeltaChannel) so Overwrite replace
  semantics work for middleware reducers in both modes.
- Endpoint regression: custom AgentMiddleware.state_schema value
  survives branch, updates through POST /state, and an unknown field
  receives 422.
- config_version 26 -> 27 for the new database.checkpoint_channel_mode
  (example, Helm chart values + README, support-bundle fixture), so
  existing installs get the outdated-config warning and
  make config-upgrade merges the field; covered by a test driving the
  real example file and the real config-upgrade script.

* fix(gateway): resolve assistant schema via one boundary, copy branch reducer values with Overwrite

GET /threads/{id}/state now resolves the thread's assistant_id through a
single reusable boundary (thread metadata -> assistant_id -> effective
graph), so channels contributed by a custom AgentMiddleware.state_schema
are materialized instead of dropped by the default lead schema. POST
/state uses the same boundary instead of resolving the schema ad hoc.

Branch writes wrap every copied reducer channel in Overwrite (derived
from the effective mutation graph: BinaryOperatorAggregate + DeltaChannel),
not just messages, so already-aggregated values are never re-merged.

Regression tests use a real AgentMiddleware.state_schema with a
non-identity reducer in both full and delta modes: GET /state returns the
extension value, POST /state replaces it, branch preserves it
byte-for-byte; the unknown-field 422 is a separate assertion.

* refactor(checkpoint): collapse read-path round-trips and ship dual-mode parity tests

Address review round 4 on PR #4292:

- Push ahistory/history limit through Pregel into checkpointer.alist
  (SQL LIMIT) instead of materializing all rows and breaking in Python
- Fold the read-side mode-compat gate onto the returned snapshot's
  metadata; only writes keep the pre-write tuple fetch (fail-closed)
- Cache factory-built accessor graphs per (assistant_id, mode) with
  factory-identity revalidation; state reads no longer build a lead
  agent per request
- get_thread: one snapshot fetch + one raw pending_writes fetch on the
  resolved checkpoint (post-checkpoint __error__ writes never surface
  in snapshot.tasks; verified empirically)
- DeerFlowClient.get_thread: single checkpointer.list walk collects
  pending_writes per checkpoint instead of N get_tuple calls
- InMemorySaver delta-history patch: stand-down when the upstream
  override disappears, try/except guard, validated-version warning,
  guard tests
- make_lead_agent mode precedence: first freeze is owned by app_config
  (client-supplied configurable key ignored); once frozen, injected
  key/app_config must match or fail closed
- Rollback: lock in non-message channel restoration via fork
  inheritance with a dedicated reducer-channel test
- Add tests/test_threads_checkpoint_mode.py and
  tests/test_gateway_checkpoint_mode.py referenced by AGENTS.md and
  the PR validation section: lifecycle parity (memory + sqlite),
  per-step blob-count storage guard, gateway endpoint parity

Counted-saver tests pin checkpoint round-trips for aget/ahistory so
these regressions cannot silently return.

* fix(checkpoint): precise mode-mismatch HTTP mapping, gate E2E, and accessor resilience

- threads router: map CheckpointModeMismatchError to 409 (with cause and
  thread id) and CheckpointModeReconfigurationError to 503 across all state
  endpoints instead of swallowing both into a generic 500
- gate coverage: seed a real delta checkpoint into AsyncSqliteSaver and
  assert aget/aupdate/ahistory fail closed; assert 409 at the HTTP boundary
  through the real route stack
- rollback: compile the restore mutation graph with the thread's effective
  state schema per the build_state_mutation_graph contract
- inheritance contract locks: rollback and manual compaction preserve
  middleware-contributed channels via checkpoint fork cloning
- services: revalidate the accessor-graph cache against app_config identity
  so config.yaml hot-reloads never serve a stale compiled graph
- services: degrade full-mode state reads to raw checkpointer reads when the
  agent factory is unavailable (delta gate still applies; delta mode has no
  fallback)
- deps: override websockets==16.0 (langgraph-sdk 0.4.2's <16 pin silently
  downgraded 16.0 -> 15.0.1; pin is not grounded in any API incompatibility)
  and bump the langchain lower bound to what the lockfile actually resolves

* fix(checkpoint): include anchor checkpoint in degraded history walk + cover get_thread

- _RawCheckpointReadAccessor.ahistory: alist(before=...) is exclusive while
  pregel's get_state_history treats config.checkpoint_id as the inclusive
  start; fetch the anchor explicitly so both read paths paginate identically
- extend the degraded-path gateway test: GET /thread returns raw values, and
  POST /history with before starts at the anchor checkpoint

* fix(gateway): preserve degraded checkpoint timestamps

* fix(gateway): harden degraded checkpoint access

* fix(gateway): resolve assistants for checkpoint reads

---------

Co-authored-by: Willem Jiang <willem.jiang@gmail.com>
2026-07-22 08:33:29 +08:00

258 lines
8.7 KiB
Python

from __future__ import annotations
from types import SimpleNamespace
from typing import Annotated, NotRequired, TypedDict
import pytest
from langchain_core.messages import AIMessage, HumanMessage
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph.message import add_messages
from langgraph.types import Overwrite
from app.gateway import services as gateway_services
from deerflow.runtime import context_compaction
from deerflow.runtime.checkpoint_state import CheckpointStateAccessor
from deerflow.runtime.context_compaction import compact_thread_context
class _FakeAccessor:
def __init__(self, values: dict) -> None:
self.snapshot = SimpleNamespace(
values=values,
config={
"configurable": {
"thread_id": "thread-1",
"checkpoint_id": "ckpt-old",
"checkpoint_ns": "",
}
},
metadata={"step": 4, "created_at": "2026-07-06T00:00:00+00:00"},
)
self.update_args = None
async def aget(self, _config):
return self.snapshot
async def aupdate(self, config, values, *, as_node=None):
self.update_args = (config, values, as_node)
return {
"configurable": {
"thread_id": config["configurable"]["thread_id"],
"checkpoint_ns": "",
"checkpoint_id": "ckpt-compacted",
}
}
class _FakeCompactionMiddleware:
def __init__(self, *, should_compact: bool = True) -> None:
self.should_compact = should_compact
self.prepare_calls = 0
self.runtime_contexts: list[dict] = []
def _prepare_compaction(self, state, *, force=False):
self.prepare_calls += 1
if not self.should_compact:
return None
return (state["messages"][:-1], state["messages"][-1:], state.get("summary_text"), 123)
async def acompact_state(self, state, runtime, *, force=False):
self.runtime_contexts.append(dict(runtime.context))
prepared = self._prepare_compaction(state, force=force)
if prepared is None:
return None
messages_to_summarize, preserved_messages, _previous_summary, total_tokens = prepared
return SimpleNamespace(
summary_text="COMPRESSED SUMMARY",
messages_to_summarize=tuple(messages_to_summarize),
preserved_messages=tuple(preserved_messages),
total_tokens=total_tokens,
)
@pytest.mark.asyncio
async def test_compact_thread_context_reads_materialized_state_and_overwrites_messages(monkeypatch):
messages = [
HumanMessage(content="old question"),
AIMessage(content="old answer"),
HumanMessage(content="latest question"),
]
accessor = _FakeAccessor(
{
"messages": messages,
"summary_text": "OLD SUMMARY",
"sandbox": object(),
}
)
middleware = _FakeCompactionMiddleware()
monkeypatch.setattr(
context_compaction,
"_create_compaction_middleware",
lambda **_kwargs: middleware,
)
result = await compact_thread_context(
accessor,
"thread-1",
app_config=SimpleNamespace(),
user_id="user-1",
agent_name="research-agent",
)
assert result.compacted is True
assert result.removed_message_count == 2
assert result.preserved_message_count == 1
assert result.summary_updated is True
assert result.checkpoint_id == "ckpt-compacted"
assert result.total_tokens == 123
assert accessor.update_args is not None
update_config, written_values, as_node = accessor.update_args
assert update_config == accessor.snapshot.config
assert isinstance(written_values["messages"], Overwrite)
assert written_values["messages"].value == [messages[-1]]
assert written_values["summary_text"] == "COMPRESSED SUMMARY"
assert as_node == "manual_compaction"
assert middleware.prepare_calls == 1
assert middleware.runtime_contexts == [
{"thread_id": "thread-1", "user_id": "user-1", "agent_name": "research-agent"},
]
@pytest.mark.asyncio
async def test_compact_thread_context_real_mutation_graph_finishes_without_scheduling(monkeypatch):
messages = [
HumanMessage(id="h1", content="old question"),
AIMessage(id="a1", content="old answer"),
HumanMessage(id="h2", content="latest question"),
]
request = SimpleNamespace(
app=SimpleNamespace(
state=SimpleNamespace(
checkpointer=InMemorySaver(),
checkpoint_channel_mode="delta",
store=None,
)
)
)
seed_accessor, seed_config = gateway_services.build_checkpoint_state_mutation_accessor(
request,
thread_id="thread-real-compaction",
as_node="seed",
)
await seed_accessor.aupdate(
seed_config,
{
"messages": Overwrite(messages),
"summary_text": "OLD SUMMARY",
},
as_node="seed",
)
accessor, config = gateway_services.build_checkpoint_state_mutation_accessor(
request,
thread_id="thread-real-compaction",
as_node="manual_compaction",
)
monkeypatch.setattr(
context_compaction,
"_create_compaction_middleware",
lambda **_kwargs: _FakeCompactionMiddleware(),
)
result = await compact_thread_context(
accessor,
"thread-real-compaction",
app_config=SimpleNamespace(),
)
snapshot = await accessor.aget(config)
assert result.compacted is True
assert [message.id for message in snapshot.values["messages"]] == ["h2"]
assert snapshot.values["summary_text"] == "COMPRESSED SUMMARY"
assert snapshot.next == ()
@pytest.mark.asyncio
async def test_compact_thread_context_preserves_middleware_contributed_channels(monkeypatch):
"""Compaction must not drop channels the base ThreadState does not know.
Contract lock for fork inheritance: the compaction write carries only
messages + summary_text (base channels), so a middleware-contributed
channel survives even though the mutation graph compiles with the base
schema. If LangGraph ever stops cloning unknown channels into forked
checkpoints, this test fails before production state is lost.
"""
from deerflow.runtime.checkpoint_state import build_state_mutation_graph
class ExtensionState(TypedDict):
messages: Annotated[list, add_messages]
memory_notes: NotRequired[str]
messages = [
HumanMessage(id="h1", content="old question"),
AIMessage(id="a1", content="old answer"),
HumanMessage(id="h2", content="latest question"),
]
request = SimpleNamespace(
app=SimpleNamespace(
state=SimpleNamespace(
checkpointer=InMemorySaver(),
checkpoint_channel_mode="full",
store=None,
)
)
)
seed_graph = build_state_mutation_graph("seed", "full", ExtensionState)
seed_accessor = CheckpointStateAccessor.bind(seed_graph, request.app.state.checkpointer, mode="full")
seed_config = {"configurable": {"thread_id": "thread-ext-compaction", "checkpoint_ns": ""}}
await seed_accessor.aupdate(
seed_config,
{"messages": messages, "memory_notes": "extension-value"},
as_node="seed",
)
# Production path: compact uses the base-schema mutation accessor.
accessor, config = gateway_services.build_checkpoint_state_mutation_accessor(
request,
thread_id="thread-ext-compaction",
as_node="manual_compaction",
)
monkeypatch.setattr(
context_compaction,
"_create_compaction_middleware",
lambda **_kwargs: _FakeCompactionMiddleware(),
)
result = await compact_thread_context(
accessor,
"thread-ext-compaction",
app_config=SimpleNamespace(),
)
snapshot = await seed_accessor.aget(config)
assert result.compacted is True
assert [message.id for message in snapshot.values["messages"]] == ["h2"]
assert snapshot.values["memory_notes"] == "extension-value"
@pytest.mark.asyncio
async def test_compact_thread_context_returns_noop_without_writing(monkeypatch):
accessor = _FakeAccessor(
{
"messages": [HumanMessage(content="latest only")],
}
)
middleware = _FakeCompactionMiddleware(should_compact=False)
monkeypatch.setattr(
context_compaction,
"_create_compaction_middleware",
lambda **_kwargs: middleware,
)
result = await compact_thread_context(accessor, "thread-1", app_config=SimpleNamespace())
assert result.compacted is False
assert result.reason == "not_enough_messages"
assert accessor.update_args is None
assert middleware.prepare_calls == 1