fix(persistence): repair inactive transcript rows in place
(cherry picked from commit 1d33010a27fbd7f053e1037269b6059edaef77ef)
This commit is contained in:
@@ -29,8 +29,9 @@ def resolve_and_repair_transcript_batch(
|
||||
decode_content_fn: Callable[[Any], Any],
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""Partition a message batch within an active write transaction. An assistant message carrying an
|
||||
existing integer ``_row_id`` targets its active SQLite row (or the active clone a watermark compaction
|
||||
made of it): a blank row is updated in place; a non-blank one (concurrent winner) has its canonical
|
||||
existing integer ``_row_id`` targets that SQLite row, or the active clone a watermark compaction made
|
||||
of it. An inactive row without a clone is repaired in place without changing its archive/rewind state;
|
||||
an active blank row is filled, while an active non-blank row (concurrent winner) has its canonical
|
||||
content adopted without overwrite. Returns the messages that must be inserted as fresh rows."""
|
||||
inserted_rows: List[Dict[str, Any]] = []
|
||||
for msg in messages:
|
||||
@@ -44,7 +45,15 @@ def resolve_and_repair_transcript_batch(
|
||||
target_id = int(target_row["id"])
|
||||
decoded = decode_content_fn(target_row["content"])
|
||||
msg["_row_id"] = target_id
|
||||
if is_content_blank(decoded):
|
||||
if int(target_row["active"] or 0) == 0:
|
||||
# The row identity still belongs to this session, but compaction/rewind removed it from the
|
||||
# model projection. Persist an in-place sanitizer rewrite without resurrecting the row or
|
||||
# appending a second display identity.
|
||||
conn.execute(
|
||||
"UPDATE messages SET content = ? WHERE id = ? AND session_id = ? AND active = 0",
|
||||
(encode_content_fn(msg.get("content")), target_id, session_id),
|
||||
)
|
||||
elif is_content_blank(decoded):
|
||||
conn.execute(
|
||||
"UPDATE messages SET content = ? "
|
||||
"WHERE id = ? AND session_id = ? AND active = 1",
|
||||
@@ -56,7 +65,7 @@ def resolve_and_repair_transcript_batch(
|
||||
|
||||
|
||||
def _active_assistant_row(conn: sqlite3.Connection, session_id: str, row_id: int):
|
||||
"""The active assistant row for ``row_id``, or the active clone a watermark compaction made of it."""
|
||||
"""The active clone for ``row_id``, or the addressed inactive assistant when no clone exists."""
|
||||
row = conn.execute(
|
||||
"SELECT id, role, active, timestamp, content FROM messages "
|
||||
"WHERE id = ? AND session_id = ?",
|
||||
@@ -66,14 +75,16 @@ def _active_assistant_row(conn: sqlite3.Connection, session_id: str, row_id: int
|
||||
return None
|
||||
if int(row["active"] or 0) == 1:
|
||||
return row
|
||||
# Watermark compaction soft-archived the concurrent tail and cloned it.
|
||||
return conn.execute(
|
||||
# Watermark compaction soft-archived the concurrent tail and cloned it. Prefer that live identity;
|
||||
# otherwise keep the addressed inactive row so a later sanitizer pass cannot append it as new.
|
||||
clone = conn.execute(
|
||||
"SELECT id, role, active, timestamp, content FROM messages "
|
||||
"WHERE session_id = ? AND active = 1 AND role = 'assistant' "
|
||||
"AND timestamp IS ? AND id != ? "
|
||||
"ORDER BY id DESC LIMIT 1",
|
||||
(session_id, row["timestamp"], row["id"]),
|
||||
).fetchone()
|
||||
return clone if clone is not None else row
|
||||
|
||||
|
||||
def sync_flushed_message_markers(batch_msgs: List[Dict[str, Any]], batch_rows: List[Dict[str, Any]]) -> None:
|
||||
|
||||
@@ -748,6 +748,82 @@ def test_flush_stale_row_id_from_other_session_does_not_fill_child_blank(tmp_pat
|
||||
]
|
||||
|
||||
|
||||
def test_flush_sanitized_archived_row_does_not_append_duplicate(tmp_path):
|
||||
"""A sanitizer rewrite of an archived row must preserve its durable identity.
|
||||
|
||||
SessionDB already scrubs lone surrogates on the initial write. The outbound sanitizer later makes the
|
||||
same repair on the live dict and pops ``_db_persisted``. Re-appending that inactive row creates a second
|
||||
assistant with the same content and microsecond timestamp, and both rows enter the display projection.
|
||||
"""
|
||||
from agent.message_sanitization import _sanitize_messages_surrogates
|
||||
|
||||
agent = _make_agent()
|
||||
db_path = tmp_path / "state.db"
|
||||
session_id = "sess-sanitized-archived-row"
|
||||
db = _attach_real_session_db(agent, db_path, session_id)
|
||||
messages = [
|
||||
{"role": "user", "content": "summarize"},
|
||||
{"role": "assistant", "content": "answer \ud800 tail"},
|
||||
]
|
||||
agent._flush_messages_to_session_db(messages)
|
||||
assistant_id = messages[-1]["_row_id"]
|
||||
assistant_timestamp = messages[-1]["timestamp"]
|
||||
|
||||
db.archive_and_compact(
|
||||
session_id,
|
||||
compacted_messages=[{"role": "user", "content": "prior turns summarized"}],
|
||||
)
|
||||
assert _sanitize_messages_surrogates(messages) is True
|
||||
agent._db_flush_scan_prefix = None
|
||||
|
||||
assert agent._flush_messages_to_session_db(messages) is True
|
||||
|
||||
all_rows = db.get_messages(session_id, include_inactive=True)
|
||||
assistants = [row for row in all_rows if row.get("role") == "assistant"]
|
||||
assert len(assistants) == 1
|
||||
assert assistants[0]["id"] == assistant_id
|
||||
assert assistants[0]["timestamp"] == assistant_timestamp
|
||||
assert assistants[0]["active"] in (0, False)
|
||||
assert assistants[0]["compacted"] in (1, True)
|
||||
assert assistants[0]["content"] == "answer \ufffd tail"
|
||||
display = db.get_messages(session_id, include_compacted=True)
|
||||
assert [row["content"] for row in display].count("answer \ufffd tail") == 1
|
||||
assert messages[-1]["_row_id"] == assistant_id
|
||||
assert messages[-1]["_db_persisted"] is True
|
||||
|
||||
|
||||
def test_flush_ascii_repair_updates_archived_row_without_resurrecting_it(tmp_path):
|
||||
"""A changed archived payload is updated in place while its active/compacted state is preserved."""
|
||||
from agent.message_sanitization import _sanitize_messages_non_ascii
|
||||
|
||||
agent = _make_agent()
|
||||
db_path = tmp_path / "state.db"
|
||||
session_id = "sess-ascii-repaired-archived-row"
|
||||
db = _attach_real_session_db(agent, db_path, session_id)
|
||||
messages = [
|
||||
{"role": "user", "content": "summarize"},
|
||||
{"role": "assistant", "content": "caf\u00e9 answer"},
|
||||
]
|
||||
agent._flush_messages_to_session_db(messages)
|
||||
assistant_id = messages[-1]["_row_id"]
|
||||
|
||||
db.archive_and_compact(
|
||||
session_id,
|
||||
compacted_messages=[{"role": "user", "content": "prior turns summarized"}],
|
||||
)
|
||||
assert _sanitize_messages_non_ascii(messages) is True
|
||||
agent._db_flush_scan_prefix = None
|
||||
|
||||
assert agent._flush_messages_to_session_db(messages) is True
|
||||
|
||||
all_rows = db.get_messages(session_id, include_inactive=True)
|
||||
assistant = next(row for row in all_rows if row.get("id") == assistant_id)
|
||||
assert assistant["content"] == "caf answer"
|
||||
assert assistant["active"] in (0, False)
|
||||
assert assistant["compacted"] in (1, True)
|
||||
assert len([row for row in all_rows if row.get("role") == "assistant"]) == 1
|
||||
assert not any(row.get("role") == "assistant" for row in db.get_messages(session_id))
|
||||
|
||||
def test_flush_archived_same_session_row_id_fills_active_clone(tmp_path):
|
||||
"""Watermark compaction clones the tail; stale `_row_id` must not win.
|
||||
|
||||
|
||||
Reference in New Issue
Block a user