diff --git a/agent/conversation_compression.py b/agent/conversation_compression.py index 047325e97b..bef36be775 100644 --- a/agent/conversation_compression.py +++ b/agent/conversation_compression.py @@ -3524,6 +3524,7 @@ def compress_context( profile_name=_profile_for_child, compression_lock_holder=_lock_holder, require_compression_lease=_lock_holder is not None, + watermark=_commit_watermark, ) agent.session_id = new_session_id try: diff --git a/hermes_state.py b/hermes_state.py index 5148a8e29e..b40537a0bd 100644 --- a/hermes_state.py +++ b/hermes_state.py @@ -5590,12 +5590,21 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin) profile_name: str = None, compression_lock_holder: str = None, require_compression_lease: bool = True, + watermark: Optional[int] = None, ) -> None: """Atomically close a parent and publish its durable compression child. The parent closure, child row, and compacted handoff become visible in one transaction. Readers can therefore observe either the live parent or a complete child, never an ended parent with a missing/empty child. + + Concurrent-append safety (#75316): when *watermark* is provided (the + parent's :meth:`get_active_message_watermark` captured at compression + start), parent rows that arrived during the slow summary call + (``id > watermark``) are cloned into the child AFTER the handoff — + same pure-SQL column clone as :meth:`archive_and_compact`, with the + session id rewritten — so a mid-compression append survives rotation + instead of stranding in the closed parent. """ def _do(conn): lock_row = conn.execute( @@ -5663,6 +5672,39 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin) total_messages, total_tool_calls = self._insert_message_rows( conn, child_session_id, messages ) + if watermark is not None: + # Clone the parent's concurrent tail (rows landed after the + # watermark) into the child, after the handoff. Column-exact + # except id/session_id; originals stay in the (closed) parent + # for lineage recovery. + tail_rows = conn.execute( + "SELECT id, tool_calls FROM messages " + "WHERE session_id = ? AND active = 1 AND id > ? ORDER BY id", + (parent_session_id, int(watermark)), + ).fetchall() + if tail_rows: + tail_ids = [int(r["id"]) for r in tail_rows] + placeholders = ",".join("?" for _ in tail_ids) + clone_cols = [ + c for c in self._message_column_names(conn) + if c not in ("id", "session_id", "active", "compacted") + ] + col_list = ", ".join(clone_cols) + conn.execute( + f"INSERT INTO messages ({col_list}, session_id, active, compacted) " + f"SELECT {col_list}, ?, 1, 0 FROM messages " + f"WHERE id IN ({placeholders}) ORDER BY id", + [child_session_id, *tail_ids], + ) + total_messages += len(tail_ids) + for r in tail_rows: + raw = r["tool_calls"] + if raw: + try: + parsed = json.loads(raw) if isinstance(raw, str) else raw + total_tool_calls += len(parsed) if isinstance(parsed, list) else 0 + except (TypeError, ValueError): + pass conn.execute( "UPDATE sessions SET message_count = ?, tool_call_count = ? WHERE id = ?", (total_messages, total_tool_calls, child_session_id), diff --git a/tests/state/test_compression_lineage_guard.py b/tests/state/test_compression_lineage_guard.py index ec8d2da182..eb64dbec43 100644 --- a/tests/state/test_compression_lineage_guard.py +++ b/tests/state/test_compression_lineage_guard.py @@ -305,16 +305,21 @@ def test_publish_compression_child_rejects_lost_or_expired_lease(db: SessionDB) def test_compression_lease_blocks_non_owner_but_allows_owner_flush( db: SessionDB, ) -> None: + """Contract flipped by the watermark commit (#75316): a live lease no + longer fences ordinary appends — both the owner's flush and a concurrent + turn land immediately, and the commit-side watermark decides what + survives compaction (see test_compression_watermark_commit.py).""" db.create_session("leased", source="webui") assert db.try_acquire_compression_lock("leased", "winner", ttl_seconds=60) - with pytest.raises(RuntimeError, match="being compressed"): - db.append_message("leased", "user", "late stale turn") - + db.append_message("leased", "user", "late concurrent turn") db.append_message( "leased", "assistant", "winner flush", compression_lock_holder="winner", ) - assert [m["content"] for m in db.get_messages("leased")] == ["winner flush"] + assert [m["content"] for m in db.get_messages("leased")] == [ + "late concurrent turn", + "winner flush", + ] diff --git a/tests/test_compression_watermark_commit.py b/tests/test_compression_watermark_commit.py index f870452269..5d91f33143 100644 --- a/tests/test_compression_watermark_commit.py +++ b/tests/test_compression_watermark_commit.py @@ -226,3 +226,54 @@ class TestConcurrentAppendDuringCompaction: 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" + + +class TestRotationPathWatermark: + """Legacy (non-in-place) compression rotates to a child session — + the concurrent tail must follow the rotation instead of stranding in + the closed parent.""" + + def test_tail_clones_into_the_child(self, db: SessionDB) -> None: + _seed(db) + watermark = db.get_active_message_watermark("sess1") + assert db.try_acquire_compression_lock("sess1", "rotator") is True + db.append_message("sess1", role="user", content="mid-rotation steer") + + db.publish_compression_child( + parent_session_id="sess1", + child_session_id="child1", + source="test", + messages=SUMMARY, + compression_lock_holder="rotator", + require_compression_lease=True, + watermark=watermark, + ) + + child = db.get_messages_as_conversation("child1") + assert [m["content"] for m in child] == [ + SUMMARY[0]["content"], + SUMMARY[1]["content"], + "mid-rotation steer", + ] + info = db.get_session("child1") + assert info["message_count"] == 3 + # Parent keeps its copy for lineage recovery; parent is closed. + parent_info = db.get_session("sess1") + assert parent_info["end_reason"] == "compression" + + def test_no_watermark_keeps_historical_rotation(self, db: SessionDB) -> None: + _seed(db) + assert db.try_acquire_compression_lock("sess1", "rotator") is True + db.append_message("sess1", role="user", content="stranded either way") + db.publish_compression_child( + parent_session_id="sess1", + child_session_id="child1", + source="test", + messages=SUMMARY, + compression_lock_holder="rotator", + require_compression_lease=True, + ) + child = db.get_messages_as_conversation("child1") + assert [m["content"] for m in child] == [ + SUMMARY[0]["content"], SUMMARY[1]["content"], + ]