Files
hermes-agent/tests/test_compression_watermark_commit.py
T
Teknium 21d3e63702 fix(compression): watermark commit — appends flow freely, concurrent tail survives compaction
Redesign of the #75316 class (supersedes the approach in PR #87307).

Root cause family: the compression lock fenced ORDINARY transcript appends
for the whole slow provider-summary call. Turns died as
session_persistence_failed whenever a message overlapped a compression
(#74568, #77386, #75083), stale dead-PID locks blocked writes for the full
TTL, and the busy-wait mitigation (#75264) was an order of magnitude shorter
than real summaries. Separately, the commit archived from a pre-call
snapshot, so rows appended mid-compression were swept into the archive.

Design: the commit transaction is already exclusive — no lock phases needed.

1. Appends never check compression_locks. The lock's only job is stopping
   two compressions colliding; it keeps that job. The whole stale-lock /
   busy-wait symptom family dies as a class.
2. Watermark captured in the DB at compression start
   (get_active_message_watermark = MAX(id) of active rows) — not from
   in-memory message dicts, which carry no row ids in production.
3. archive_and_compact(watermark=, lock_holder=): one transaction verifies
   the holder still owns an unexpired lease (a reclaimed lease cannot
   publish a stale compaction), archives the snapshot, inserts the compacted
   set, and re-sequences the concurrent tail (id > watermark) via a
   pure-SQL column clone — every column except id survives byte-exact
   (api_content, platform_message_id, reasoning sidecars, token counts),
   FTS triggers index the clones naturally, originals stay archived and
   recoverable. watermark=None preserves the historical behavior.

Removed: the append-side compression fence in _check_transcript_write_guards
(with rationale note), making the _COMPRESSION_BUSY_WAIT_S retry lane
unreachable from append paths (kept for other callers).

Tests: 12 new (watermark contract, column-exact clone, commit fence incl.
lease-lost/expired/rollback failure injection, append-vs-commit race);
busy-retry suite flipped to pin the new contract; sabotage-verified (5 fail
with the watermark disabled, 12 pass restored); E2E through the real
compress_context seam with a mid-summary append landing and surviving.
2026-08-15 23:37:22 -07:00

229 lines
9.2 KiB
Python

"""Watermark commit: concurrent appends survive in-place compaction (#75316).
The provider summary call is external and slow. Messages that arrive while it
runs must (a) persist immediately — appends are not fenced by the compression
lock — and (b) survive the commit: ``archive_and_compact(watermark=...)``
re-sequences every active row above the watermark after the compacted set
instead of archiving it. The commit is holder-fenced: a compression whose
lease was reclaimed cannot publish a stale compaction.
"""
from __future__ import annotations
import json
import sqlite3
import threading
import time
from pathlib import Path
import pytest
from hermes_state import SessionCompressionInProgressError, SessionDB
@pytest.fixture
def db(tmp_path: Path) -> SessionDB:
d = SessionDB(tmp_path / "state.db")
d.create_session("sess1", source="test")
return d
def _seed(db: SessionDB, n: int = 6) -> None:
for i in range(n):
role = "user" if i % 2 == 0 else "assistant"
db.append_message("sess1", role=role, content=f"turn {i}")
SUMMARY = [
{"role": "user", "content": "[CONTEXT COMPACTION] summary of turns 0-5"},
{"role": "assistant", "content": "Continuing from the summary."},
]
class TestWatermarkCommit:
def test_concurrent_tail_survives_compaction(self, db: SessionDB) -> None:
_seed(db)
watermark = db.get_active_message_watermark("sess1")
# Simulate the slow summary window: two messages land after capture.
db.append_message("sess1", role="user", content="mid-compression steer")
db.append_message("sess1", role="assistant", content="mid-compression reply")
count = db.archive_and_compact("sess1", SUMMARY, watermark=watermark)
live = db.get_messages("sess1")
contents = [r["content"] for r in live]
assert contents == [
"[CONTEXT COMPACTION] summary of turns 0-5",
"Continuing from the summary.",
"mid-compression steer",
"mid-compression reply",
], "tail must follow the summary, in arrival order"
assert count == 4
def test_tail_clone_preserves_every_column(self, db: SessionDB) -> None:
"""The pure-SQL clone must carry sidecar fields byte-exact."""
_seed(db, 2)
watermark = db.get_active_message_watermark("sess1")
db.append_message(
"sess1",
role="assistant",
content="tool caller",
tool_calls=[{"id": "c1", "type": "function",
"function": {"name": "terminal", "arguments": "{}"}}],
)
db.append_message(
"sess1", role="tool", content="tool output",
tool_call_id="c1", tool_name="terminal",
)
db.archive_and_compact("sess1", SUMMARY, watermark=watermark)
live = db.get_messages("sess1")
by_content = {r["content"]: r for r in live}
caller = by_content["tool caller"]
result = by_content["tool output"]
parsed = caller["tool_calls"]
if isinstance(parsed, str):
parsed = json.loads(parsed)
assert parsed and parsed[0]["id"] == "c1"
assert result["tool_call_id"] == "c1"
assert result["tool_name"] == "terminal"
def test_conversation_load_is_correct_after_commit(self, db: SessionDB) -> None:
"""The live conversation projection sees summary + tail, in order."""
_seed(db)
watermark = db.get_active_message_watermark("sess1")
db.append_message("sess1", role="user", content="late arrival")
db.archive_and_compact("sess1", SUMMARY, watermark=watermark)
convo = db.get_messages_as_conversation("sess1")
assert [m["content"] for m in convo] == [
"[CONTEXT COMPACTION] summary of turns 0-5",
"Continuing from the summary.",
"late arrival",
]
def test_no_tail_behaves_identically_to_legacy(self, db: SessionDB) -> None:
_seed(db)
watermark = db.get_active_message_watermark("sess1")
count = db.archive_and_compact("sess1", SUMMARY, watermark=watermark)
assert count == 2
assert [r["content"] for r in db.get_messages("sess1")] == [
SUMMARY[0]["content"], SUMMARY[1]["content"],
]
def test_none_watermark_preserves_historical_behavior(self, db: SessionDB) -> None:
"""watermark=None archives everything — the pre-#75316 contract."""
_seed(db)
db.append_message("sess1", role="user", content="gets archived")
count = db.archive_and_compact("sess1", SUMMARY, watermark=None)
assert count == 2
contents = [r["content"] for r in db.get_messages("sess1")]
assert "gets archived" not in contents
def test_archived_rows_stay_recoverable(self, db: SessionDB) -> None:
"""Originals (snapshot AND tail source rows) survive as archived."""
_seed(db, 4)
watermark = db.get_active_message_watermark("sess1")
db.append_message("sess1", role="user", content="tail row")
db.archive_and_compact("sess1", SUMMARY, watermark=watermark)
everything = db.get_messages("sess1", include_inactive=True)
archived = [r for r in everything if not r["active"]]
assert sum(1 for r in archived if r["content"] == "turn 0") == 1
# The tail original is archived; its clone is the live copy.
tail_rows = [r for r in everything if r["content"] == "tail row"]
assert sorted(bool(r["active"]) for r in tail_rows) == [False, True]
def test_session_counters_include_tail(self, db: SessionDB) -> None:
_seed(db)
watermark = db.get_active_message_watermark("sess1")
db.append_message(
"sess1", role="assistant", content="tail with tools",
tool_calls=[{"id": "t1", "type": "function",
"function": {"name": "x", "arguments": "{}"}}],
)
db.archive_and_compact("sess1", SUMMARY, watermark=watermark)
info = db.get_session("sess1")
assert info["message_count"] == 3
assert info["tool_call_count"] == 1
class TestCommitFence:
def test_commit_refused_when_lease_lost(self, db: SessionDB) -> None:
_seed(db)
watermark = db.get_active_message_watermark("sess1")
assert db.try_acquire_compression_lock("sess1", "worker-A") is True
# Lease reclaimed by another writer while worker-A's summary ran.
db.release_compression_lock("sess1", "worker-A")
assert db.try_acquire_compression_lock("sess1", "worker-B") is True
with pytest.raises(SessionCompressionInProgressError):
db.archive_and_compact(
"sess1", SUMMARY, watermark=watermark, lock_holder="worker-A"
)
# Nothing committed: original transcript intact.
assert [r["content"] for r in db.get_messages("sess1")] == [
f"turn {i}" for i in range(6)
]
def test_commit_refused_when_lease_expired(self, db: SessionDB) -> None:
_seed(db)
assert db.try_acquire_compression_lock(
"sess1", "worker-A", ttl_seconds=0.05
) is True
time.sleep(0.1)
with pytest.raises(SessionCompressionInProgressError):
db.archive_and_compact("sess1", SUMMARY, lock_holder="worker-A")
def test_commit_allowed_for_live_holder(self, db: SessionDB) -> None:
_seed(db)
watermark = db.get_active_message_watermark("sess1")
assert db.try_acquire_compression_lock("sess1", "worker-A") is True
count = db.archive_and_compact(
"sess1", SUMMARY, watermark=watermark, lock_holder="worker-A"
)
assert count == 2
def test_refused_commit_rolls_back_atomically(self, db: SessionDB) -> None:
"""Failure injection: the fence raise must leave zero partial writes."""
_seed(db)
before = db.get_messages("sess1", include_inactive=True)
with pytest.raises(SessionCompressionInProgressError):
db.archive_and_compact("sess1", SUMMARY, lock_holder="never-held")
after = db.get_messages("sess1", include_inactive=True)
assert len(before) == len(after)
assert all(r["active"] for r in after)
class TestConcurrentAppendDuringCompaction:
def test_append_racing_the_commit_transaction(self, db: SessionDB) -> None:
"""An append serialized behind the commit lands AFTER it — never lost.
SQLite's write lock serializes the two transactions; whichever side
wins, the append must end up in the live transcript.
"""
_seed(db)
watermark = db.get_active_message_watermark("sess1")
barrier = threading.Barrier(2, timeout=10)
append_err: list = []
def _racer():
barrier.wait()
try:
db.append_message("sess1", role="user", content="racer")
except Exception as exc: # pragma: no cover
append_err.append(exc)
t = threading.Thread(target=_racer, daemon=True)
t.start()
barrier.wait()
db.archive_and_compact("sess1", SUMMARY, watermark=watermark)
t.join(timeout=10)
assert not append_err, f"append died during commit race: {append_err}"
contents = [r["content"] for r in db.get_messages("sess1")]
assert "racer" in contents, "racing append was lost"