fix(agent): keep history replay pending after any failure following a heal

When the heal recreated the session row but its single retry write then
failed (database locked, turn lease, disk), the replay was carried only by
the per-call _replay_history argument. The next flush found a live row, hit
no FK error, ran no heal, and stamped the history prefix durable without
writing it: the durable transcript silently lost everything but the tail
(scratch probe: rows ['a2'] instead of ['one','a1','two','a2']).

Set _session_row_replay_pending in the heal branch right after the markers
are stripped, so any exit after a heal (recreate failed, retry failed)
leaves the replay pending until a write succeeds. That makes the separate
recreate-failed assignment and the _replay_history parameter redundant;
the retry is a plain _adoption_budget=0 call. The fails-closed test now
covers FK -> heal -> locked retry -> full replay on the next flush.

Co-authored-by: 赵桂雄 <daniel21436@hotmail.com>
This commit is contained in:
kshitijk4poor
2026-09-26 21:36:40 +05:30
committed by kshitij
parent 44cfc57158
commit ff05a54ecd
2 changed files with 30 additions and 8 deletions

View File

@@ -335,6 +335,9 @@ def _db_flush_failed(agent, e: Exception, batch_rows: List[Dict[str, Any]], adop
if adoption_budget <= 0 or not _db_flush_session_row_gone(agent, agent.session_id):
return None
_strip_persistence_markers(messages)
# The durable history prefix is gone with the row: keep a replay pending until a write succeeds, so
# a failed recreate or a failed retry write can't let a later flush stamp that prefix durable.
agent._session_row_replay_pending = agent.session_id
agent._flushed_db_message_ids = set()
agent._last_flushed_db_idx = 0
agent._session_db_created = False
@@ -351,8 +354,7 @@ def _db_flush_failed(agent, e: Exception, batch_rows: List[Dict[str, Any]], adop
if not agent._session_db_created:
# Row creation failed too (transient store trouble): don't append into a guaranteed
# rollback — keep the batch unmarked so the next flush retries the whole thing. That flush
# recreates the row up front (no FK error, no heal), so it must replay the history prefix itself.
agent._session_row_replay_pending = agent.session_id
# recreates the row up front (no FK error, no heal); the pending replay covers the history prefix.
logger.warning("Session DB row for %s is missing and could not be recreated; will retry next flush",
getattr(agent, "session_id", None))
return None
@@ -440,7 +442,6 @@ class SessionPersistenceMixin:
def _flush_messages_to_session_db_unlocked(
self, messages: List[Dict], conversation_history: Optional[List[Dict]] = None, _adoption_budget: int = 1,
_replay_history: bool = False,
):
"""Persist un-flushed messages to SQLite. Dedup is the intrinsic ``_DB_PERSISTED_MARKER`` on each written
dict — not positional slices (drift after sequence repair) nor an ``id(msg)`` set (address reuse). The
@@ -462,7 +463,8 @@ class SessionPersistenceMixin:
try:
if not self._session_db_created: # retry row creation if the earlier attempt failed transiently
self._ensure_db_session()
replay = _replay_history or getattr(self, "_session_row_replay_pending", None) == self.session_id
# getattr: object.__new__ test agents flush without running AIAgent init.
replay = getattr(self, "_session_row_replay_pending", None) == self.session_id
batch_rows, batch_msgs = _db_flush_collect(self, messages, conversation_history, replay)
_db_flush_write(self, batch_rows, batch_msgs, messages)
self._session_row_replay_pending = None
@@ -476,10 +478,8 @@ class SessionPersistenceMixin:
retry = _db_flush_failed(self, e, batch_rows, _adoption_budget, messages)
if retry is None:
return False
# A recreated row lost the history prefix too: replay it rather than stamping it durable, but keep
# the history set so a muted notification turn hides only its own rows.
return self._flush_messages_to_session_db_unlocked(
messages, conversation_history, _adoption_budget=0, _replay_history=retry == "healed")
# After a heal the pending replay re-sends the history prefix instead of stamping it durable.
return self._flush_messages_to_session_db_unlocked(messages, conversation_history, _adoption_budget=0)
def _get_messages_up_to_last_assistant(self, messages: List[Dict]) -> List[Dict]:
"""Messages before the last assistant turn (rollback point for a malformed final answer); all if none."""

View File

@@ -102,6 +102,28 @@ def test_flush_fails_closed_when_row_cannot_be_recreated(monkeypatch):
assert agent._flush_messages_to_session_db(turn3, turn2) is True
assert [r["content"] for r in db.get_messages("sess-gone")] == ["a", "b", "c"]
# The heal recreates the row but its single retry write fails (lock/lease/disk): the next
# flush finds a live row (no FK, no heal) and must still replay the history prefix.
retry = _make_agent(db, "sess-retry")
t1 = [{"role": "user", "content": "one"}, {"role": "assistant", "content": "a1"}]
assert retry._flush_messages_to_session_db(t1, []) is True
assert db.delete_session("sess-retry") is True
real_append, calls = db.append_messages_batch, []
def _fk_then_locked(*a, **kw):
calls.append(1)
if len(calls) == 2:
raise _sqlite3.OperationalError("database is locked")
return real_append(*a, **kw) # call 1 hits the real FK error -> heal
monkeypatch.setattr(db, "append_messages_batch", _fk_then_locked)
t2 = t1 + [{"role": "user", "content": "two"}]
assert retry._flush_messages_to_session_db(t2, t1) is False
t3 = t2 + [{"role": "assistant", "content": "a2"}]
assert retry._flush_messages_to_session_db(t3, t2) is True
assert [r["content"] for r in db.get_messages("sess-retry")] == ["one", "a1", "two", "a2"]
monkeypatch.undo()
# FK failure while the session row still exists (e.g. a sessions-table FK): no heal, and
# no replay of the history prefix onto the live transcript.
live = _make_agent(db, "sess-fk-live")