fix(compression): archive the rows the compressor held, not every row up to the newest

A gap below the newest held id, and turns appended above an unpersisted current
turn, were summarized away without being read. Name those held ids and clone
the rest.
This commit is contained in:
brooklyn!
2026-09-24 11:57:56 -05:00
parent a7c44e9ede
commit 1674499d00
9 changed files with 283 additions and 6 deletions
+18
View File
@@ -448,6 +448,21 @@ def _merge_assistant_into(prev: Dict, msg: Dict) -> None:
prev.pop(_DB_PERSISTED_MARKER, None)
def _remember_absorbed_row(survivor: Dict[str, Any], dropped: Dict[str, Any]) -> None:
"""Record durable ids a merge folded into *survivor* and then dropped from the list.
The survivor keeps one ``_row_id``. Without the absorbed ids, an archive capped at
that id clones the folded row beside content the summary already contains.
"""
absorbed = survivor.setdefault("_absorbed_row_ids", [])
row_id = dropped.get("_row_id")
if isinstance(row_id, int) and not isinstance(row_id, bool) and row_id > 0 and row_id not in absorbed:
absorbed.append(row_id)
for older in dropped.get("_absorbed_row_ids") or ():
if isinstance(older, int) and not isinstance(older, bool) and older > 0 and older not in absorbed:
absorbed.append(older)
def _merge_consecutive_assistants(messages: List[Dict]) -> Tuple[List[Dict], int]:
"""Pass 0: merge consecutive assistant turns (codex interims exempt)."""
repairs = 0
@@ -461,9 +476,11 @@ def _merge_consecutive_assistants(messages: List[Dict]) -> Tuple[List[Dict], int
):
# A provisional verification candidate is superseded, not unioned.
if prev.get("finish_reason") in {"verification_required", "verify_hook_continue"}:
_remember_absorbed_row(msg, prev)
collapsed[-1] = msg
else:
_merge_assistant_into(prev, msg)
_remember_absorbed_row(prev, msg)
repairs += 1
continue
collapsed.append(msg)
@@ -583,6 +600,7 @@ def _merge_consecutive_users(messages: List[Dict]) -> Tuple[List[Dict], int]:
# reproduces the persisted bytes (e.g. an empty incoming turn) keeps its stamp.
if merged_content != prev_content or had_api_sidecar:
prev.pop(_DB_PERSISTED_MARKER, None)
_remember_absorbed_row(prev, msg)
repairs += 1
continue
merged.append(msg)
+3
View File
@@ -3296,10 +3296,13 @@ class ContextCompressor(SummaryDispatchMixin, MicroCompactionMixin, ContextEngin
next_rearm_tokens = after + runway
if session_db and session_id:
try:
from agent.conversation_compression_archive import coverage_for_commit
covered_ids, unresolved_held = coverage_for_commit(session_db, session_id, messages)
session_db.archive_and_compact(
session_id, pruned_msgs,
model_config_patch={PROACTIVE_PRUNE_REARM_MODEL_CONFIG_KEY: next_rearm_tokens},
watermark=_archive_watermark_for(session_db, session_id, messages),
covered_ids=covered_ids, unresolved_held=unresolved_held,
)
except StaleHeldHistory:
# Another compaction already committed this session's history; a lease-less prune of the
+6
View File
@@ -3733,10 +3733,16 @@ def _commit_compaction(
from hermes_cli.partial_compress import rejoin_compressed_head_and_tail
persisted = rejoin_compressed_head_and_tail(compressed, verbatim_tail)
tail_count += len(verbatim_tail)
from agent.conversation_compression_archive import coverage_for_commit
covered_ids, unresolved_held = coverage_for_commit(
agent._session_db, agent.session_id,
messages_before_compression if messages_before_compression is not None else messages,
verbatim_tail)
agent._session_db.archive_and_compact(
agent.session_id, persisted, model_config_patch={PROACTIVE_PRUNE_REARM_MODEL_CONFIG_KEY: None},
watermark=_held_watermark(agent, lease.watermark, messages, verbatim_tail),
lock_holder=lease.holder, tail_count=tail_count, carried_messages=carried_messages,
covered_ids=covered_ids, unresolved_held=unresolved_held,
)
compressed = persisted
split_status = "in_place_committed"
+99
View File
@@ -0,0 +1,99 @@
"""Exact held-row coverage for an in-place archive.
A watermark cap archives every active row up to the newest id the caller held.
That deletes a gap below that id, and a trailing unpersisted turn makes the cap
fall back to the lease watermark. When the caller can name the rows the
compressor actually held, the commit archives those and clones the rest.
"""
from __future__ import annotations
from typing import Any, Dict, List, Optional, Sequence, Tuple
ABSORBED_ROW_IDS = "_absorbed_row_ids"
def _positive_id(value: Any) -> Optional[int]:
if isinstance(value, int) and not isinstance(value, bool) and value > 0:
return value
return None
def held_archive_coverage(
messages: Sequence[Any], verbatim_tail: Optional[Sequence[Any]] = None,
) -> Tuple[List[int], List[Dict[str, Any]]]:
"""``(covered ids, unresolved dicts)`` from the history the compressor was handed.
A positive ``_row_id`` is covered, including ids a repair merged into that dict.
A dict with no id is unresolved: the commit matches one durable row, and a
marker-less miss is an unpersisted turn rather than a reason to archive the
lease watermark.
"""
covered: List[int] = []
unresolved: List[Dict[str, Any]] = []
for batch in (messages or (), verbatim_tail or ()):
for message in batch:
if not isinstance(message, dict):
continue
row_id = _positive_id(message.get("_row_id"))
if row_id is None:
unresolved.append(message)
else:
covered.append(row_id)
for absorbed in message.get(ABSORBED_ROW_IDS) or ():
absorbed_id = _positive_id(absorbed)
if absorbed_id is not None:
covered.append(absorbed_id)
return list(dict.fromkeys(covered)), unresolved
def newest_exact_held_id(
messages: Sequence[Any], verbatim_tail: Optional[Sequence[Any]] = None,
) -> Optional[int]:
"""Newest held id that still names the row the compressor saw.
A head row counts only while it still carries the persist marker. A ``here N``
tail is marker-swept copies; its id counts when the copy kept one.
"""
from agent.context_compressor import _DB_PERSISTED_MARKER
exact: List[int] = []
for message in messages or ():
if not isinstance(message, dict) or not message.get(_DB_PERSISTED_MARKER):
continue
row_id = _positive_id(message.get("_row_id"))
if row_id is not None:
exact.append(row_id)
for message in verbatim_tail or ():
if isinstance(message, dict):
row_id = _positive_id(message.get("_row_id"))
if row_id is not None:
exact.append(row_id)
return max(exact) if exact else None
def coverage_for_commit(
session_db: Any, session_id: str, messages: Sequence[Any],
verbatim_tail: Optional[Sequence[Any]] = None,
) -> Tuple[Optional[List[int]], Optional[List[Dict[str, Any]]]]:
"""Coverage to pass into ``archive_and_compact``, or ``(None, None)`` to keep the watermark.
``None`` when the newest exact held row is already inactive (another compaction
won: archiving only the held ids would clone the winner) or when nothing held
is a durable row. A trailing unpersisted turn does not take this branch: the
rows above it stay unnamed and are cloned.
"""
from agent.context_compressor import _DB_PERSISTED_MARKER
newest = newest_exact_held_id(messages, verbatim_tail)
role_of = getattr(session_db, "get_message_role", None)
if newest is not None and callable(role_of) and role_of(session_id, newest) is None:
return None, None
covered, unresolved = held_archive_coverage(messages, verbatim_tail)
marked = [
message for message in unresolved
if isinstance(message, dict) and message.get(_DB_PERSISTED_MARKER)
]
if not covered and not marked:
return None, None
return covered, unresolved
+9 -1
View File
@@ -112,7 +112,15 @@ def compress_now(
for m in tail:
row = _fresh_compaction_message_copy(m)
if not m.get(_DB_PERSISTED_MARKER):
row.pop("_row_id", None)
# A rewritten row's id does not bound the archive, but it and any row a merge
# folded into it were still in the compressor's input. Keep them named so the
# commit archives those originals instead of cloning them beside the tail.
dropped = row.pop("_row_id", None)
if isinstance(dropped, int) and not isinstance(dropped, bool) and dropped > 0:
absorbed = [*(row.get("_absorbed_row_ids") or ())]
if dropped not in absorbed:
absorbed.append(dropped)
row["_absorbed_row_ids"] = absorbed
tail_rows.append(row)
try:
compressed, _ = agent._compress_context(
+5 -1
View File
@@ -420,10 +420,14 @@ class MicroCompactionMixin:
if isinstance(message, dict) and message.get(_cc()._DB_PERSISTED_MARKER)
]
watermark = None
covered_ids = unresolved_held = None
if held is not None and start_watermark is not None:
watermark = _cc()._archive_watermark_for(session_db, session_id, held, start_watermark)
from agent.conversation_compression_archive import coverage_for_commit
covered_ids, unresolved_held = coverage_for_commit(session_db, session_id, held)
session_db.archive_and_compact(
session_id, compacted_messages, carried_messages=carried_messages, watermark=watermark)
session_id, compacted_messages, carried_messages=carried_messages, watermark=watermark,
covered_ids=covered_ids, unresolved_held=unresolved_held)
# Shared post-commit stamp site with batch commit and proactive prune.
# See #98450.
_cc().stamp_db_persisted_markers(compacted_messages)
+93 -2
View File
@@ -732,15 +732,101 @@ class SessionMessagesMixin:
resolved.append(matches[0])
return list(dict.fromkeys(resolved))
def _matching_active_ids(self, conn, session_id: str, message: Dict[str, Any]) -> List[int]:
"""Active row ids whose stored role and content equal *message*. Empty when it was never persisted."""
content = message.get("content")
if not isinstance(content, str):
return []
stored = self._encode_content(self._loaded_view_content(message.get("role", "unknown"), content))
return [int(row["id"]) for row in conn.execute(
"SELECT id FROM messages WHERE session_id = ? AND active = 1 AND role = ? AND content = ?",
(session_id, message.get("role"), stored)).fetchall()]
def _proved_coverage(
self, conn, session_id: str, covered_ids: Optional[List[int]],
unresolved_held: Optional[List[Dict[str, Any]]],
) -> Optional[List[int]]:
"""Ids safe to archive as summarized, or None when a durable held row cannot be named.
An unresolved dict that still carries the persist marker was loaded from the DB.
Failing to name it means the watermark path, which archives the rows the compressor
saw, including ones whose ids were stripped. A marker-less miss is an unpersisted
turn: it names nothing, and it is not a reason to abandon the ids we do have.
Several active rows with the same content are ambiguous, so that also abandons.
"""
if covered_ids is None:
return None
from agent.context_compressor import _DB_PERSISTED_MARKER
proved = [int(row_id) for row_id in covered_ids if isinstance(row_id, int) and row_id > 0]
for message in unresolved_held or ():
if not isinstance(message, dict):
continue
matches = self._matching_active_ids(conn, session_id, message)
if len(matches) > 1 or (message.get(_DB_PERSISTED_MARKER) and len(matches) != 1):
return None
proved.extend(matches)
return list(dict.fromkeys(proved))
def _archive_named_rows(
self, conn, session_id: str, compacted_messages: List[Dict[str, Any]], covered: List[int], *,
tail_count: int, carried_messages: Optional[List[Dict[str, Any]]], patched_model_config: Any,
patch: bool,
) -> int:
"""Archive *covered* as summarized and clone every other active row after the new set.
A gap below the newest held id, and rows appended after an unpersisted turn, are not
in *covered*. They take the concurrent-append path: rewind the original, insert the
compacted transcript, then clone them so they stay live and searchable once each.
"""
active_ids = [int(row["id"]) for row in conn.execute(_ACTIVE_IDS_SQL, (session_id,)).fetchall()]
covered_set = set(covered)
carried_ids = self._resolve_carried_row_ids(conn, session_id, carried_messages or [])
covered_set.update(carried_ids)
unseen = [row_id for row_id in active_ids if row_id not in covered_set]
covered_active = [row_id for row_id in active_ids if row_id in covered_set]
rewind_ids = list(carried_ids)
if tail_count > 0:
rewind_ids += covered_active[-int(tail_count):]
rewind_ids += unseen
rewind_ids = list(dict.fromkeys(rewind_ids))
if rewind_ids:
placeholders = _placeholders(rewind_ids)
conn.execute(
"UPDATE messages SET active = 0, compacted = 0 "
f"WHERE session_id = ? AND id IN ({placeholders})",
[session_id, *rewind_ids])
conn.execute(_ARCHIVE_ACTIVE_SQL, (session_id,))
inserted, tool_calls_total = self._insert_message_rows(conn, session_id, compacted_messages)
if unseen:
_ids, unseen_tool_calls = self._tail_rows_after_watermark(
conn,
"SELECT id, tool_calls FROM messages WHERE id IN ({}) ORDER BY id".format(
_placeholders(unseen)),
tuple(unseen))
self._clone_message_rows(conn, unseen)
inserted += len(unseen)
tool_calls_total += unseen_tool_calls
conn.execute(
f"{_SET_COUNTERS_SQL}{', model_config = ?' if patch else ''} WHERE id = ?",
(inserted, tool_calls_total, *((patched_model_config,) if patch else ()), session_id))
return inserted
def archive_and_compact(self, session_id: str, compacted_messages: List[Dict[str, Any]],
model_config_patch: Optional[Dict[str, Any]] = None, watermark: Optional[int] = None,
lock_holder: Optional[str] = None, tail_count: int = 0,
carried_messages: Optional[List[Dict[str, Any]]] = None) -> int:
carried_messages: Optional[List[Dict[str, Any]]] = None,
covered_ids: Optional[List[int]] = None,
unresolved_held: Optional[List[Dict[str, Any]]] = None) -> int:
"""Non-destructive in-place compaction under ONE session id: soft-archive the active rows (``active=0,
compacted=1``: summarized away, still searchable) and insert *compacted_messages* as fresh active
rows, atomically; returns the new ACTIVE count (= ``message_count``). *watermark* (compression
START): rows ``id > watermark`` arrived during the slow summary and are re-sequenced after the
compacted set by a pure-SQL clone (fresh ids); ``None`` archives everything. *lock_holder*: verified
compacted set by a pure-SQL clone (fresh ids); ``None`` archives everything. *covered_ids*: the
rows the compressor actually held. When proved, only those are summarized; every other active
row is cloned after the new set, so a gap below the newest held id is not archived unseen.
``None`` keeps the watermark path. *unresolved_held*: held dicts with no row id, matched inside
the transaction. *lock_holder*: verified
in-txn so a reclaimed lease fails instead of clobbering the winner. *tail_count*: the LAST N compacted
rows are the verbatim carried tail; *carried_messages* names exact durable originals carried forward
verbatim when they are not a contiguous suffix (micro-compaction's prefix + marker + suffix shape).
@@ -768,6 +854,11 @@ class SessionMessagesMixin:
# on_missing="raise": never commit against a vanished session row (caller keeps the original).
patched_model_config = self._merge_model_config_json(
conn, session_id, model_config_patch, on_missing="raise") if patch else None
proved = self._proved_coverage(conn, session_id, covered_ids, unresolved_held)
if proved is not None:
return self._archive_named_rows(
conn, session_id, compacted_messages, proved, tail_count=tail_count,
carried_messages=carried_messages, patched_model_config=patched_model_config, patch=patch)
tail_ids, tail_tool_calls = ([], 0) if watermark is None else self._tail_rows_after_watermark(
conn, "SELECT id, tool_calls FROM messages WHERE session_id = ? AND active = 1 AND id > ? ORDER BY id",
(session_id, int(watermark)))
@@ -286,3 +286,48 @@ def test_in_place_compress_never_clones_a_row_a_merge_already_carried(session_db
contents = [m["content"] for m in session_db.get_messages_as_conversation("sid")]
assert sum("second half 9931" in c for c in contents) == 1
def test_in_place_compress_keeps_a_gap_below_the_newest_held_row(session_db):
"""A surface can hold rows 1–N, miss another surface's rows, then persist its own later rows.
The newest held id is then the lease watermark, so a cap at that id archives the unseen gap.
Those rows were never summarized; they must stay live, once each."""
from agent.context_compressor import _DB_PERSISTED_MARKER
agent, _ = _stored_agent(session_db, _exchanges(10))
held = session_db.get_resume_conversations("sid")[0]
for role, content in FOREIGN_TURN:
session_db.append_message("sid", role, content)
own_id = session_db.append_message("sid", "user", "continued on this surface")
held.append({
"role": "user", "content": "continued on this surface",
"_row_id": own_id, _DB_PERSISTED_MARKER: True,
})
assert _compress(agent, held, "").status == "compressed"
model_history, display_history = session_db.get_resume_conversations("sid")
for _role, content in FOREIGN_TURN:
assert any(content in (m.get("content") or "") for m in model_history)
assert any(content in (m.get("content") or "") for m in display_history)
assert _flags(session_db, content) == [(0, 0), (1, 0)]
assert session_db.search_messages("vault 7741")
def test_in_place_compress_keeps_foreign_rows_above_an_unpersisted_turn(session_db):
"""Turn preflight appends the current user message before compression and persists it afterward,
so the last dict has no row id. That must not fall back to the lease watermark and archive turns
another surface appended since this process loaded."""
agent, _ = _stored_agent(session_db, _exchanges(10))
held = session_db.get_resume_conversations("sid")[0]
for role, content in FOREIGN_TURN:
session_db.append_message("sid", role, content)
held.append({"role": "user", "content": "this turn is not persisted yet"})
assert _compress(agent, held, "").status == "compressed"
model_history, _display = session_db.get_resume_conversations("sid")
for _role, content in FOREIGN_TURN:
assert any(content in (m.get("content") or "") for m in model_history)
assert _flags(session_db, content) == [(0, 0), (1, 0)]
assert session_db.search_messages("vault 7741")
+5 -2
View File
@@ -1713,7 +1713,7 @@ def _append_model_switch_marker(session: dict | None, *, model: str, provider: s
"metadata when answering questions about what model/provider is active.]")
# A user message, not system: strict OpenAI-compatible providers (vLLM, Qwen) reject non-leading system messages.
# See #48338.
entry = {"role": "user", "content": marker, "display_kind": "model_switch"}
entry: dict[str, Any] = {"role": "user", "content": marker, "display_kind": "model_switch"}
with session.get("history_lock") or contextlib.nullcontext():
history = session.setdefault("history", [])
history[:] = [h for h in history if not _is_model_switch_marker(h)]
@@ -1726,7 +1726,10 @@ def _append_model_switch_marker(session: dict | None, *, model: str, provider: s
_ensure_session_db_row(session)
with (contextlib.nullcontext(db) if db is not None else _session_db(session)) as db:
if db is not None:
db.append_message(session_id=session_key, role="user", content=marker, display_kind="model_switch")
from agent.context_compressor import _DB_PERSISTED_MARKER
entry["_row_id"] = db.append_message(
session_id=session_key, role="user", content=marker, display_kind="model_switch")
entry[_DB_PERSISTED_MARKER] = True
except Exception:
logger.debug("failed to persist model switch marker", exc_info=True)