Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
135 changes: 129 additions & 6 deletions hermes_state.py
Original file line number Diff line number Diff line change
Expand Up @@ -8431,8 +8431,9 @@ def get_messages_as_conversation(
LIVE REPLAY should pass it: a durable alternation violation (e.g. a
``user;user`` pair left by a turn that persisted no assistant row)
otherwise re-triggers the pre-request defensive repair on every
single request for the rest of the session's life — the repair
mutates only the per-request list, never the stored transcript.
single request for the rest of the session's life. A repair-enabled
load reconciles the changed survivors and archived-away rows back to
the stored transcript before returning.
Inspection/export consumers keep the default and see the transcript
verbatim.
"""
Expand Down Expand Up @@ -8464,6 +8465,7 @@ def get_messages_as_conversation(
include_ancestors=include_ancestors,
repair_alternation=repair_alternation,
include_row_ids=include_row_ids,
reconcile_repair=not include_inactive,
)

# Columns every conversation projection decodes. Shared by
Expand All @@ -8484,6 +8486,7 @@ def _rows_to_conversation(
include_ancestors: bool,
repair_alternation: bool,
include_row_ids: bool = False,
reconcile_repair: bool = True,
) -> List[Dict[str, Any]]:
"""Decode fetched message rows into the OpenAI conversation format.

Expand All @@ -8504,7 +8507,7 @@ def _rows_to_conversation(
# inspection) gets the transcript in its historical shape.
# Underscore-prefixed so every transport's convert_messages()
# strips it before the wire.
if include_row_ids and row["id"] is not None:
if (include_row_ids or repair_alternation) and row["id"] is not None:
msg["_row_id"] = row["id"]
# api_content is the byte-fidelity sidecar: the exact string sent
# to the API when it differed from the clean content. Returned
Expand Down Expand Up @@ -8593,22 +8596,142 @@ def _rows_to_conversation(
# to keep emitting the marker. No-op for unaffected sessions.
messages = _strip_stale_tool_call_markers(messages)
if repair_alternation and messages:
# Snapshot AFTER load-only sanitation so reconciliation covers
# only mutations made by repair_message_sequence. Surviving
# messages retain object/row identity; removed identities are the
# exact durable rows to archive.
pre_repair = {
msg["_row_id"]: dict(msg)
for msg in messages
if isinstance(msg, dict) and "_row_id" in msg
}
# Lazy import: hermes_state already depends on agent.* (see
# sanitize_context above), but keep this optional path from
# widening the import surface at module load.
from agent.agent_runtime_helpers import repair_message_sequence

repaired = repair_message_sequence(None, messages)
if repaired:
if reconcile_repair:
self._reconcile_repaired_conversation(pre_repair, messages)
logger.info(
"Repaired %d message-alternation violation(s) while "
"restoring session %s — durable transcript kept them, "
"see repair_message_sequence",
"Repaired %d message-alternation violation(s) while restoring "
"session %s%s",
repaired,
session_id,
(
" and reconciled the durable transcript"
if reconcile_repair
else " without changing inactive audit rows"
),
)
if not include_row_ids:
for msg in messages:
if isinstance(msg, dict):
msg.pop("_row_id", None)
return messages

def _reconcile_repaired_conversation(
self,
pre_repair: Dict[int, Dict[str, Any]],
repaired_messages: List[Dict[str, Any]],
) -> None:
"""Persist one identity-preserving alternation repair atomically.

Rows removed by the repair are soft-archived for audit. Surviving
rows keep their durable identity while repair-mutated provider fields
are updated in place. Updating content/tool_calls also drives the
existing FTS triggers, while an ``active``-only write deliberately
leaves archived text searchable only through inactive/audit paths.
"""
survivors = {
msg["_row_id"]: msg
for msg in repaired_messages
if isinstance(msg, dict) and msg.get("_row_id") in pre_repair
}
# Verification candidates are collapsed only from MODEL replay; their
# substantive answer must remain active in the persisted DISPLAY
# transcript (#65919). repair_message_sequence deliberately replaces
# the candidate with the verified assistant without mutating either
# durable row, so do not treat that model-only replacement as archival.
verification_finish_reasons = {
"verification_required",
"verify_hook_continue",
}
removed_ids = sorted(
row_id
for row_id, original in pre_repair.items()
if row_id not in survivors
and original.get("finish_reason") not in verification_finish_reasons
)
updates = []

for row_id, msg in survivors.items():
original = pre_repair[row_id]
content_changed = msg.get("content") != original.get("content")
tool_calls_changed = msg.get("tool_calls") != original.get("tool_calls")
reasoning_changed = (
msg.get("reasoning_content") != original.get("reasoning_content")
)

# api_content is the exact wire form of the old content. Any
# content merge invalidates it even on assistant rows, whose
# direct repair path historically retained the stale sidecar.
if content_changed:
msg.pop("api_content", None)
api_content_changed = (
msg.get("api_content") != original.get("api_content")
)

if not (
content_changed
or tool_calls_changed
or reasoning_changed
or api_content_changed
):
continue

tool_calls = msg.get("tool_calls")
if isinstance(tool_calls, str):
try:
tool_calls = json.loads(tool_calls)
except (json.JSONDecodeError, TypeError):
tool_calls = []
updates.append(
(
self._encode_content(msg.get("content")),
json.dumps(tool_calls) if tool_calls else None,
_scrub_surrogates(msg.get("reasoning_content")),
(
_scrub_surrogates(msg.get("api_content"))
if isinstance(msg.get("api_content"), str)
else None
),
row_id,
)
)

if not removed_ids and not updates:
return

def _do(conn):
if removed_ids:
placeholders = ",".join("?" for _ in removed_ids)
conn.execute(
f"UPDATE messages SET active = 0 "
f"WHERE active = 1 AND id IN ({placeholders})",
removed_ids,
)
if updates:
conn.executemany(
"UPDATE messages SET content = ?, tool_calls = ?, "
"reasoning_content = ?, api_content = ? "
"WHERE id = ? AND active = 1",
updates,
)

self._execute_write(_do, patience_s=self._TRANSCRIPT_WRITE_PATIENCE_S)

def get_resume_conversations(
self, session_id: str
) -> Tuple[List[Dict[str, Any]], List[Dict[str, Any]]]:
Expand Down
105 changes: 103 additions & 2 deletions tests/hermes_state/test_restore_alternation_repair.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,8 +7,9 @@
pre-request ``repair_message_sequence`` re-fires on EVERY request for the
rest of the session's life, because it mutates only the per-request list.

Default (``repair_alternation=False``) must stay verbatim: inspection and
export consumers (trace upload, context guard) read the transcript as-is.
Default (``repair_alternation=False``) never initiates repair: inspection and
export consumers read the current active transcript as stored. Rows archived
by an earlier live repair remain available through ``include_inactive=True``.
"""

import pytest
Expand Down Expand Up @@ -57,6 +58,106 @@ def test_repaired_load_is_stable_under_prerequest_repair(db):
assert repair_message_sequence(None, messages) == 0


def test_repaired_load_reconciles_durable_rows_and_fts(db):
"""A repair-enabled restore heals state.db, not only its return value."""
from agent.agent_runtime_helpers import repair_message_sequence

db.create_session("durable", "system prompt")
db.append_message("durable", "user", "first ask")
db.append_message("durable", "assistant", "first reply")
survivor_id = db.append_message(
"durable",
"user",
"unanswered turn",
api_content="WIRE-SIDECAR: unanswered turn + injected context",
)
merged_away_id = db.append_message("durable", "user", "next turn")
db.append_message(
"durable",
"assistant",
"calling tool",
tool_calls=[
{
"id": "call-1",
"type": "function",
"function": {"name": "demo", "arguments": "{}"},
}
],
)
db.append_message(
"durable", "tool", "valid result", tool_call_id="call-1"
)
orphan_id = db.append_message(
"durable", "tool", "duplicate result", tool_call_id="call-1"
)
assistant_survivor_id = db.append_message(
"durable",
"assistant",
"next reply",
reasoning_content="first reasoning",
api_content="WIRE-SIDECAR: next reply",
)
assistant_merged_away_id = db.append_message(
"durable",
"assistant",
"continued reply",
tool_calls=[
{
"id": "call-2",
"type": "function",
"function": {"name": "next_demo", "arguments": "{}"},
}
],
reasoning_content="later reasoning",
)

first = db.get_messages_as_conversation(
"durable", repair_alternation=True, include_row_ids=True
)

assert [message["_row_id"] for message in first] == [
row["id"]
for row in db.get_messages("durable")
]
assert first[2]["_row_id"] == survivor_id
assert first[2]["content"] == "unanswered turn\n\nnext turn"
assert "api_content" not in first[2]
assert first[-1]["_row_id"] == assistant_survivor_id
assert first[-1]["content"] == "next reply\ncontinued reply"
assert [call["id"] for call in first[-1]["tool_calls"]] == ["call-2"]
assert first[-1]["reasoning_content"] == "first reasoning"
assert "api_content" not in first[-1]

all_rows = {
row["id"]: row
for row in db.get_messages("durable", include_inactive=True)
}
assert all_rows[survivor_id]["content"] == "unanswered turn\n\nnext turn"
assert all_rows[survivor_id]["api_content"] is None
assert all_rows[merged_away_id]["active"] == 0
assert all_rows[orphan_id]["active"] == 0
assert all_rows[assistant_survivor_id]["content"] == (
"next reply\ncontinued reply"
)
assert [
call["id"] for call in all_rows[assistant_survivor_id]["tool_calls"]
] == ["call-2"]
assert all_rows[assistant_survivor_id]["reasoning_content"] == (
"first reasoning"
)
assert all_rows[assistant_survivor_id]["api_content"] is None
assert all_rows[assistant_merged_away_id]["active"] == 0

raw_reload = db.get_messages_as_conversation("durable")
assert repair_message_sequence(None, raw_reload) == 0
assert raw_reload == db.get_messages_as_conversation(
"durable", repair_alternation=True
)

search_hits = db.search_messages("unanswered next")
assert [hit["id"] for hit in search_hits] == [survivor_id]




# ---------------------------------------------------------------------------
Expand Down
Loading