diff --git a/agent/transcript_repair.py b/agent/transcript_repair.py index f245088abc..f9b1500823 100644 --- a/agent/transcript_repair.py +++ b/agent/transcript_repair.py @@ -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: diff --git a/tests/agent/test_tool_call_incremental_persistence.py b/tests/agent/test_tool_call_incremental_persistence.py index 59898c34e5..cb3d99d6d6 100644 --- a/tests/agent/test_tool_call_incremental_persistence.py +++ b/tests/agent/test_tool_call_incremental_persistence.py @@ -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.