fix: stamp _db_persisted at row load time so resumed transcripts never re-append (#92231)
Resumed sessions loaded message dicts from state.db WITHOUT the _DB_PERSISTED_MARKER, so any flush that lost the identity boundary (compression durable-snapshot adoption, incremental tool-call persists, rotation preflight on cold resume) re-appended the ENTIRE loaded transcript as new rows. Compression cycles then doubled the copies: the incident session grew 998 -> 1995 -> 3990 -> 7981 rows across three aborted rotations (15,962 active rows, only 472 distinct). Fix at the architectural chokepoint: SessionDB._rows_to_conversation (shared by get_messages_as_conversation and get_resume_conversations) now stamps the marker at row materialization time - a dict built FROM a durable row is persisted by construction, regardless of which caller loads it or how the list is later handed to a flush. Safety: - Wire-safe: every transport strips underscore-prefixed keys before the API request (chat_completion_helpers, anthropic_adapter), same contract as the existing _row_id stamp in the same function. - Rotation handoffs still write: compression's assembly copies strip the marker (_fresh_compaction_message_copy + the terminal _strip_persistence_markers sweep), so compacted transcripts still flush to the child session (#57491 invariant preserved). - Branch/seed copies unaffected: /branch and _persist_branch_seed build fresh field-projected dicts and write via append_messages_batch directly, not through the marker-gated flush. Tests: new regression suite (marker sync, load stamping, 3-cycle amplification repro, new-tail write guard, compaction-copy handoff); updated the #68454 control test that asserted the old double-write behavior and the ACP restore shape test.
This commit is contained in:
@@ -776,6 +776,14 @@ def _strip_background_review_harness(
|
||||
# Matches a bare protocol/tool-name marker such as "[memory]" or "[skill_manage]".
|
||||
_STALE_TOOL_CALL_MARKER_RE = re.compile(r"^\[[A-Za-z_][A-Za-z0-9_.-]*\]$")
|
||||
|
||||
# Intrinsic persistence marker stamped on message dicts that are known-durable.
|
||||
# MUST stay in sync with run_agent._DB_PERSISTED_MARKER and
|
||||
# agent.context_compressor._DB_PERSISTED_MARKER (same literal; defined locally
|
||||
# because hermes_state must not import run_agent — circular import — and
|
||||
# importing the agent layer for one constant inverts the state/agent layering).
|
||||
# Drift is guarded by test_marker_constant_in_sync (#92231).
|
||||
_DB_PERSISTED_MARKER_KEY = "_db_persisted"
|
||||
|
||||
|
||||
def _is_stale_tool_call_marker_message(msg: Dict[str, Any]) -> bool:
|
||||
"""True when ``msg`` is a persisted assistant turn whose content is a bare
|
||||
@@ -11090,6 +11098,18 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin)
|
||||
if row["role"] in {"user", "assistant"} and isinstance(content, str):
|
||||
content = sanitize_context(content).strip()
|
||||
msg = {"role": row["role"], "content": content}
|
||||
# Born durable (#92231): this dict is materialized FROM a durable
|
||||
# row, so stamp the persistence marker at the source instead of
|
||||
# relying on every restore caller to thread the loaded list back
|
||||
# through a flush as ``conversation_history=`` — any
|
||||
# identity-losing handoff (compression's durable-snapshot
|
||||
# adoption, incremental persists with no history arg) would
|
||||
# otherwise re-append the ENTIRE transcript on flush.
|
||||
# Underscore-prefixed like ``_row_id``: every transport strips it
|
||||
# before the wire, and compression's assembly copies deliberately
|
||||
# strip it so rotated child handoffs still flush (see
|
||||
# _fresh_compaction_message_copy).
|
||||
msg[_DB_PERSISTED_MARKER_KEY] = True
|
||||
# Durable per-message identity for surfaces that need to address a
|
||||
# specific row later (desktop reactions). OPT-IN: only the gateway
|
||||
# asks for it — every other consumer (ACP restore, export,
|
||||
|
||||
@@ -326,6 +326,9 @@ class TestPersistence:
|
||||
assert restored is not None
|
||||
msg = restored.history[0]
|
||||
assert isinstance(msg.pop("timestamp", None), (int, float))
|
||||
# Load-time durability stamp (#92231): rows materialized from the DB
|
||||
# are marked persisted so a later flush can't re-append them.
|
||||
assert msg.pop("_db_persisted", None) is True
|
||||
assert restored.history == [{
|
||||
"role": "assistant",
|
||||
"content": "hello",
|
||||
|
||||
147
tests/agent/test_load_time_durability_stamp_92231.py
Normal file
147
tests/agent/test_load_time_durability_stamp_92231.py
Normal file
@@ -0,0 +1,147 @@
|
||||
"""Regression (#92231): resumed transcripts must not re-append to state.db.
|
||||
|
||||
Root cause: ``get_messages_as_conversation`` / ``get_resume_conversations``
|
||||
returned plain dicts WITHOUT ``_DB_PERSISTED_MARKER``. Any flush that received
|
||||
that loaded history without a matching ``conversation_history=`` identity
|
||||
boundary (compression durable-snapshot adoption, incremental tool-call
|
||||
persists, rotation preflight on cold resume) treated every loaded row as new
|
||||
and re-appended the ENTIRE transcript. Compression cycles then doubled the
|
||||
copies: the incident session grew 998 → 1995 → 3990 → 7981 rows across three
|
||||
aborted rotations (15,962 active rows / 472 distinct contents).
|
||||
|
||||
Fix: rows are stamped durable at materialization time in
|
||||
``SessionDB._rows_to_conversation`` — a dict built FROM a durable row is
|
||||
persisted by construction, no matter which caller loads it or how it is later
|
||||
handed to a flush.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from pathlib import Path
|
||||
from types import SimpleNamespace
|
||||
|
||||
from hermes_state import SessionDB, _DB_PERSISTED_MARKER_KEY
|
||||
from run_agent import AIAgent
|
||||
|
||||
|
||||
def _make_flush_agent(db: SessionDB, session_id: str):
|
||||
"""Minimal agent shell that owns the real flush implementation."""
|
||||
agent = SimpleNamespace(
|
||||
_session_db=db,
|
||||
_session_db_created=True,
|
||||
_persist_disabled=False,
|
||||
session_id=session_id,
|
||||
_session_persist_lock=None,
|
||||
_flushed_db_message_ids=set(),
|
||||
_flushed_db_message_session_id=None,
|
||||
_last_flushed_db_idx=0,
|
||||
_persist_user_message_idx=None,
|
||||
_persist_user_message_override=None,
|
||||
_persist_user_message_timestamp=None,
|
||||
_pending_cli_user_message=None,
|
||||
)
|
||||
agent._ensure_db_session = lambda: None
|
||||
agent._flush_messages_to_session_db = (
|
||||
AIAgent._flush_messages_to_session_db.__get__(agent, AIAgent)
|
||||
)
|
||||
agent._flush_messages_to_session_db_unlocked = (
|
||||
AIAgent._flush_messages_to_session_db_unlocked.__get__(agent, AIAgent)
|
||||
)
|
||||
return agent
|
||||
|
||||
|
||||
def _seed_session(db: SessionDB, sid: str, turns: int = 3) -> None:
|
||||
db.create_session(sid, source="cli")
|
||||
for i in range(turns):
|
||||
db.append_message(sid, "user", f"question {i}")
|
||||
db.append_message(sid, "assistant", f"answer {i}")
|
||||
|
||||
|
||||
def test_marker_constant_in_sync() -> None:
|
||||
"""hermes_state cannot import run_agent (circular) — literals must match."""
|
||||
import agent.context_compressor as cc
|
||||
import run_agent
|
||||
|
||||
assert _DB_PERSISTED_MARKER_KEY == run_agent._DB_PERSISTED_MARKER
|
||||
assert _DB_PERSISTED_MARKER_KEY == cc._DB_PERSISTED_MARKER
|
||||
|
||||
|
||||
def test_loaded_rows_are_stamped_durable(tmp_path: Path) -> None:
|
||||
db = SessionDB(db_path=tmp_path / "state.db")
|
||||
_seed_session(db, "S1")
|
||||
|
||||
loaded = db.get_messages_as_conversation("S1")
|
||||
assert loaded
|
||||
assert all(m.get(_DB_PERSISTED_MARKER_KEY) is True for m in loaded)
|
||||
|
||||
model_history, display_history = db.get_resume_conversations("S1")
|
||||
assert model_history and display_history
|
||||
assert all(m.get(_DB_PERSISTED_MARKER_KEY) is True for m in model_history)
|
||||
assert all(m.get(_DB_PERSISTED_MARKER_KEY) is True for m in display_history)
|
||||
|
||||
|
||||
def test_repeated_identityless_flushes_do_not_amplify(tmp_path: Path) -> None:
|
||||
"""The #92231 shape: reload + flush cycles must keep the row count flat.
|
||||
|
||||
Each cycle simulates what compression's durable-snapshot adoption (or any
|
||||
identity-losing handoff) used to do: load the transcript fresh from the DB
|
||||
and flush it with NO ``conversation_history`` boundary. Pre-fix each cycle
|
||||
doubled the row count (998 → 1995 → 3990 → 7981 in the incident session).
|
||||
"""
|
||||
db = SessionDB(db_path=tmp_path / "state.db")
|
||||
_seed_session(db, "S2", turns=4)
|
||||
baseline = len(db.get_messages("S2"))
|
||||
|
||||
for _ in range(3):
|
||||
loaded = db.get_messages_as_conversation("S2")
|
||||
agent = _make_flush_agent(db, "S2") # fresh agent: no identity state
|
||||
agent._flush_messages_to_session_db(loaded)
|
||||
|
||||
assert len(db.get_messages("S2")) == baseline
|
||||
|
||||
|
||||
def test_new_tail_after_loaded_history_still_flushes(tmp_path: Path) -> None:
|
||||
"""Guard against over-skipping: only loaded rows are exempt, new turns write."""
|
||||
db = SessionDB(db_path=tmp_path / "state.db")
|
||||
_seed_session(db, "S3", turns=1)
|
||||
|
||||
loaded = db.get_messages_as_conversation("S3")
|
||||
live = [
|
||||
*loaded,
|
||||
{"role": "user", "content": "new question"},
|
||||
{"role": "assistant", "content": "new answer"},
|
||||
]
|
||||
agent = _make_flush_agent(db, "S3")
|
||||
agent._flush_messages_to_session_db(live)
|
||||
|
||||
contents = [m.get("content") for m in db.get_messages("S3")]
|
||||
assert contents == ["question 0", "answer 0", "new question", "new answer"]
|
||||
|
||||
|
||||
def test_compaction_copy_strips_stamp_so_child_flush_writes(tmp_path: Path) -> None:
|
||||
"""Rotation handoff: compression copies must stay flushable to the child.
|
||||
|
||||
``_fresh_compaction_message_copy`` / ``_strip_persistence_markers`` remove
|
||||
the marker from assembled compaction output so the rotation flush WRITES
|
||||
the compacted transcript to the child session (#57491). Load-stamping must
|
||||
not defeat that: a loaded row that goes through the compaction copy is
|
||||
written to the child exactly once.
|
||||
"""
|
||||
from agent.context_compressor import _fresh_compaction_message_copy
|
||||
|
||||
db = SessionDB(db_path=tmp_path / "state.db")
|
||||
_seed_session(db, "PARENT", turns=2)
|
||||
db.create_session("CHILD", source="cli", parent_session_id="PARENT")
|
||||
|
||||
loaded = db.get_messages_as_conversation("PARENT")
|
||||
compacted = [_fresh_compaction_message_copy(m) for m in loaded]
|
||||
assert all(_DB_PERSISTED_MARKER_KEY not in m for m in compacted)
|
||||
|
||||
agent = _make_flush_agent(db, "CHILD")
|
||||
agent._flush_messages_to_session_db(compacted)
|
||||
child_rows = db.get_messages("CHILD")
|
||||
assert len(child_rows) == len(loaded)
|
||||
|
||||
# Idempotent thereafter: the flush stamped the copies on write.
|
||||
agent._flush_messages_to_session_db(compacted)
|
||||
assert len(db.get_messages("CHILD")) == len(loaded)
|
||||
@@ -50,8 +50,15 @@ def _contents(rows):
|
||||
return [r.get("content") for r in rows]
|
||||
|
||||
|
||||
def test_rotation_flush_without_history_boundary_duplicates(tmp_path: Path) -> None:
|
||||
"""Control: bare flush of unstamped cold-resume rows double-writes (#68454)."""
|
||||
def test_rotation_flush_without_history_boundary_is_safe(tmp_path: Path) -> None:
|
||||
"""Bare flush of cold-resumed rows must NOT double-write (#68454 → #92231).
|
||||
|
||||
Historically this was a control test asserting the double-write: loaded
|
||||
rows carried no ``_DB_PERSISTED_MARKER``, so a flush without
|
||||
``conversation_history=`` re-appended the whole transcript. Since #92231
|
||||
the loaders stamp the marker at row-materialization time, so even the
|
||||
"wrong" call shape (no history boundary) is idempotent.
|
||||
"""
|
||||
db = SessionDB(db_path=tmp_path / "state.db")
|
||||
sid = "COLD_ROTATE_DUP"
|
||||
db.create_session(sid, source="cli")
|
||||
@@ -63,7 +70,7 @@ def test_rotation_flush_without_history_boundary_duplicates(tmp_path: Path) -> N
|
||||
agent._flush_messages_to_session_db(loaded) # missing conversation_history=
|
||||
|
||||
rows = db.get_messages_as_conversation(sid, include_inactive=True)
|
||||
assert _contents(rows).count("persisted question") == 2
|
||||
assert _contents(rows) == ["persisted question", "persisted answer"]
|
||||
|
||||
|
||||
def test_rotation_flush_with_history_boundary_is_noop(tmp_path: Path) -> None:
|
||||
|
||||
@@ -164,12 +164,20 @@ def test_row_id_is_opt_in_and_never_reaches_the_provider(session, db):
|
||||
"""Only include_row_ids=True consumers see _row_id — and it's underscore-
|
||||
prefixed so transports strip it before the wire even for them. Default
|
||||
consumers (ACP restore, export) get the transcript in its historical shape.
|
||||
|
||||
``_db_persisted`` is the other sanctioned underscore key: stamped on every
|
||||
loaded row (#92231) so a flush can never re-append a resumed transcript.
|
||||
Like ``_row_id`` it is stripped before the wire by every transport.
|
||||
"""
|
||||
key, _rows = session
|
||||
|
||||
for message in db.get_messages_as_conversation(key):
|
||||
assert "_row_id" not in message
|
||||
assert message.get("_db_persisted") is True
|
||||
|
||||
for message in db.get_messages_as_conversation(key, include_row_ids=True):
|
||||
assert "_row_id" in message
|
||||
assert all(not k.startswith("_") or k == "_row_id" for k in message)
|
||||
assert all(
|
||||
not k.startswith("_") or k in {"_row_id", "_db_persisted"}
|
||||
for k in message
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user