fix(compression): keep the /compress here N tail in state.db under in-place compaction

compress_now() hands only the head to _compress_context() and rejoins the
kept exchanges in memory afterwards. With compression.in_place (the
default), the commit's archive_and_compact() archives every active row at
or below the lease watermark, which includes the kept tail's rows, and
inserts only the compacted head. Its tail_count rewind then lands on the
newest rows under the watermark, which are the kept tail itself, so those
rows end up active=0, compacted=0 with no live copy. No surface writes them
back: the CLI re-flushes only after a rotation, the gateway skips in-place
on purpose, and the TUI only swaps its in-memory history. The exchanges the
user asked to keep verbatim are gone on resume, and on the gateway from the
very next message.

compress_now() now passes copies of the tail as verbatim_tail. The in-place
commit stores head + tail in the same archive_and_compact() transaction,
joined with the same seam rejoin_compressed_head_and_tail() builds in
memory, and adds the tail to tail_count. The rewind flags then land on the
kept tail's originals and on compress()'s own carried rows, which were
left compacted=1 and shown twice in the resumed display history. The
copies are stamped as persisted and compress_now() returns the stored list
instead of rejoining the tail a second time. Rotation, no-op and
rolled-back commits are unchanged: the copies stay unstamped and the
caller's tail is rejoined as before.
This commit is contained in:
John Paul Soliva
2026-09-23 11:03:00 +09:00
committed by Teknium
parent 81a095cb1d
commit 526d135a96
4 changed files with 132 additions and 12 deletions

View File

@@ -221,7 +221,7 @@ class CompressionFacadeMixin:
def _compress_context(
self, messages: list, system_message: str, *, approx_tokens: int = None, task_id: str = "default",
focus_topic: str = None, force: bool = False, bypass_cooldown: bool = False,
defer_context_engine_notification: bool = False, commit_fence=None,
defer_context_engine_notification: bool = False, commit_fence=None, verbatim_tail: list = None,
) -> tuple:
"""Forwarder — see ``agent.conversation_compression.compress_context``.
``force=True`` (manual /compress) bypasses the summary-failure cooldown; ``bypass_cooldown=True``
@@ -281,6 +281,7 @@ class CompressionFacadeMixin:
approx_tokens=approx_tokens, task_id=task_id, focus_topic=focus_topic, force=force,
bypass_cooldown=bypass_cooldown or same_turn_fallback_recovery,
defer_context_engine_notification=(defer_context_engine_notification), commit_fence=fence,
verbatim_tail=verbatim_tail,
)
# Callers that already own a progress-aware wait (gateway session

View File

@@ -3583,12 +3583,14 @@ def _commit_compaction(
agent: Any, messages: list, compressed: list, *, in_place: bool, lease: _CompressionLease,
new_system_prompt: str, system_message: str, compressed_user_turn_outcome: str,
messages_before_compression: Optional[list], made_progress: bool, attempt: _Attempt,
verbatim_tail: Optional[list] = None,
) -> _CommitOutcome:
"""Persist the compacted transcript: memory extraction, anti-growth guard, then the
in-place archive or the parent->child rotation.
Failures roll the live list back and arm the split-failure cooldown; a refused (would-grow) candidate returns
``refused_prompt`` so the caller hands back the input unchanged.
``refused_prompt`` so the caller hands back the input unchanged. ``verbatim_tail`` (``/compress here N``) is
re-inserted after the compacted head by the in-place commit and stamped once durable; rotation ignores it.
"""
session_commit_succeeded = False
compacted_in_place = False
@@ -3619,16 +3621,27 @@ def _commit_compaction(
from agent.context_compressor import PROACTIVE_PRUNE_REARM_MODEL_CONFIG_KEY, stamp_db_persisted_markers
# Tail rows tagged by compress() are archived as superseded duplicates, not
# compacted=1. Count against the FINAL list — salvage may have dropped rows.
tail_count = sum(1 for m in compressed if id(m) in _tail_tagged_ids)
persisted = compressed
if verbatim_tail:
# The kept exchanges are durable rows under the watermark, so the archive below covers
# them too. Store them after the head in the same transaction, with the seam the caller
# would build, and count their originals as carried duplicates like compress()'s tail.
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)
agent._session_db.archive_and_compact(
agent.session_id, compressed, model_config_patch={PROACTIVE_PRUNE_REARM_MODEL_CONFIG_KEY: None},
watermark=lease.watermark, lock_holder=lease.holder,
tail_count=sum(1 for m in compressed if id(m) in _tail_tagged_ids),
agent.session_id, persisted, model_config_patch={PROACTIVE_PRUNE_REARM_MODEL_CONFIG_KEY: None},
watermark=lease.watermark, lock_holder=lease.holder, tail_count=tail_count,
)
compressed = persisted
split_status = "in_place_committed"
# compress() returned marker-swept copies; stamp them as persisted or the next
# flush re-INSERTs the whole compacted transcript, doubling the live set. Reset
# the flush identity set so next turn diffs against the COMPACTED transcript.
stamp_db_persisted_markers(compressed)
# The verbatim tail is stamped as well: a seam fold drops its first row from
# `compressed`, and the stamps tell the caller the tail is already in the list.
stamp_db_persisted_markers([*compressed, *(verbatim_tail or ())])
agent._flushed_db_message_ids = set()
# Rotation-independent signal; the gateway reads this (not an id diff) to
# re-baseline transcript handling.
@@ -3903,7 +3916,7 @@ def compress_context(
agent: Any, messages: list, system_message: str, *, approx_tokens: Optional[int] = None,
task_id: str = "default", focus_topic: Optional[str] = None, force: bool = False,
bypass_cooldown: bool = False, defer_context_engine_notification: bool = False,
commit_fence: Optional[CompressionCommitFence] = None,
commit_fence: Optional[CompressionCommitFence] = None, verbatim_tail: Optional[list] = None,
) -> Tuple[list, str]:
"""Compress conversation context and split the session in SQLite.
``force`` (manual /compress) clears the summary-failure cooldown; ``bypass_cooldown`` (provider-proven
@@ -3924,7 +3937,8 @@ def compress_context(
failed attempt records its cooldown normally. defer_context_engine_notification: Delay the existing
context-engine hook until a manual host commits its outer history transaction. commit_fence: Optional
cooperative fence for executor callers that may time out. It prevents a late worker from mutating
session state after its caller has moved on.
session state after its caller has moved on. verbatim_tail: The exchanges ``/compress here N`` keeps
after ``messages``; an in-place commit stores them after the compacted head and returns head + tail.
"""
attempt = _begin_compression_attempt(agent, force=force, defer_notification=defer_context_engine_notification)
@@ -4057,7 +4071,7 @@ def compress_context(
agent, messages, compressed, in_place=in_place, lease=lease, new_system_prompt=new_system_prompt,
system_message=system_message, compressed_user_turn_outcome=compressed_user_turn_outcome,
messages_before_compression=messages_before_compression, made_progress=_compression_made_progress,
attempt=attempt,
attempt=attempt, verbatim_tail=verbatim_tail,
)
if commit.refused_prompt is not None:
return messages, commit.refused_prompt

View File

@@ -83,6 +83,7 @@ def compress_now(
``_compress_context`` still does useful work there — codex_app_server native compaction, and the
phase-1 tool-result prune / blank-echo drop that ``ContextCompressor.compress`` commits even when no
summary window exists."""
from agent.context_compressor import _DB_PERSISTED_MARKER, _fresh_compaction_message_copy
from agent.conversation_compression import finalize_context_engine_compression_notification
from agent.manual_compression_feedback import summarize_manual_compression
from hermes_cli.partial_compress import (
@@ -103,10 +104,14 @@ def compress_now(
has_content = getattr(compressor, "has_content_to_compress", None)
if skip_without_window and callable(has_content) and has_content(head) is False:
return CompressResult("nothing_to_do", before, before, before_tokens, before_tokens, request)
# An in-place commit archives every durable row under the lease watermark, the kept tail's included, so
# it must store the tail again itself. It gets copies because the insert writes row ids onto them.
tail_rows = [_fresh_compaction_message_copy(m) for m in tail]
try:
compressed, _ = agent._compress_context(
head, system_message, approx_tokens=before_tokens, focus_topic=request.focus_topic, force=True,
defer_context_engine_notification=True, **({"task_id": task_id} if task_id != "default" else {}))
defer_context_engine_notification=True, **({"task_id": task_id} if task_id != "default" else {}),
**({"verbatim_tail": tail_rows} if tail_rows else {}))
except Exception:
finalize_context_engine_compression_notification(agent, committed=False)
raise
@@ -117,7 +122,9 @@ def compress_now(
finalize_context_engine_compression_notification(agent, committed=False)
return CompressResult("lock_skipped", before, before, before_tokens, before_tokens, request,
lock_holder=lock_signal if isinstance(lock_signal, str) else None)
if tail:
# Stamped copies mean the in-place commit stored the tail and already returned head + tail. Rotation, a
# no-op or a rolled-back commit leave them unstamped, and the tail is then only in the caller's dicts.
if tail and not all(row.get(_DB_PERSISTED_MARKER) is True for row in tail_rows):
compressed = rejoin_compressed_head_and_tail(compressed, tail)
after_tokens = estimate_request_tokens(agent, compressed)
summary = summarize_manual_compression(before, compressed, before_tokens, after_tokens, compression_state=compressor)

View File

@@ -3,8 +3,9 @@
from __future__ import annotations
import copy
import os
import threading
from unittest.mock import MagicMock
from unittest.mock import MagicMock, patch
import pytest
@@ -110,3 +111,100 @@ def _coro(value):
async def _inner(*_a, **_k):
return value
return _inner
@pytest.fixture
def session_db(tmp_path):
from hermes_state import SessionDB
db = SessionDB(db_path=tmp_path / "state.db")
yield db
db.close()
def _exchanges(n, *, unanswered=()):
history = []
for i in range(n):
history.append({"role": "user", "content": f"question {i} about fruit{i} " + " ".join(["filler"] * 40)})
if i not in unanswered:
history.append({"role": "assistant", "content": f"answer {i} " + " ".join(["lorem"] * 400)})
return history
def _stored_agent(db, history):
"""A real AIAgent (default in-place mode) over ``history`` as stored rows, loaded back like the gateway does."""
db.create_session("sid", "telegram", model="test/model")
for message in history:
db.append_message("sid", message["role"], message["content"])
with patch.dict(os.environ, {"OPENROUTER_API_KEY": "test-key"}):
from run_agent import AIAgent
agent = AIAgent(api_key="test-key", base_url="https://openrouter.ai/api/v1", model="test/model",
quiet_mode=True, session_db=db, session_id="sid", skip_context_files=True, skip_memory=True)
agent._compression_feasibility_checked = True
return agent, db.get_messages_as_conversation("sid")
def _compress_here(agent, history, keep):
response = MagicMock()
response.choices = [MagicMock()]
response.choices[0].message.content = "## Goal\nNumbered fruit questions.\n## Progress\nEarly ones answered."
with patch("agent.context_compressor.call_llm", lambda **_kw: response):
return compress_now(agent, history, parse_compress_args(f"here {keep}"), system_message="",
skip_without_window=True)
def _flags(db, content):
"""``(active, compacted)`` of every stored row with this exact content."""
rows = db._conn.execute("SELECT active, compacted FROM messages WHERE session_id = 'sid' AND content = ?",
(content,)).fetchall()
return sorted(tuple(row) for row in rows)
def _live(messages):
return [(m.get("role"), m.get("content")) for m in messages]
def test_in_place_here_n_stores_the_kept_exchanges(session_db):
"""The in-place commit archives every row under the lease watermark, the kept tail's included, so it must
store the tail again after the head: otherwise a resume (and the gateway's next turn) loses exactly the
exchanges the user asked to keep."""
history = _exchanges(10)
agent, loaded = _stored_agent(session_db, history)
frozen = copy.deepcopy(loaded)
assert agent.compression_in_place is True
result = _compress_here(agent, loaded, 2)
assert result.status == "compressed" and agent.session_id == "sid"
assert loaded == frozen
durable = session_db.get_messages_as_conversation("sid")
assert _live(durable) == _live(result.after_messages)
assert _live(durable[-4:]) == _live(history[-4:])
model_history, display_history = session_db.get_resume_conversations("sid")
for kept in history[-4:]:
assert [m["content"] for m in model_history].count(kept["content"]) == 1
# compress() carries its own recent rows after the summary, then comes the kept tail: each has one live
# row, and its original is a superseded duplicate, not a turn summarized away that resume shows again.
summary_at = next(i for i, m in enumerate(durable) if "Numbered fruit questions" in m["content"])
carried = [m["content"] for m in durable[summary_at + 1:]]
assert len(carried) > 4
for content in carried:
assert [m["content"] for m in display_history].count(content) == 1
assert _flags(session_db, content) == [(0, 0), (1, 0)]
# The kept rows are stamped as stored, so the next persist appends only the new turn.
next_turn = [*result.after_messages, {"role": "user", "content": "question 10 about fruit10"}]
agent._flush_messages_to_session_db(next_turn, None)
assert _live(session_db.get_messages_as_conversation("sid")) == _live(next_turn)
def test_in_place_here_n_folds_the_seam_once(session_db):
"""A head ending on a user turn folds the tail's first message into it; the stored transcript carries the
same fold, and the tail is not rejoined a second time."""
history = _exchanges(10, unanswered={7})
agent, loaded = _stored_agent(session_db, history)
result = _compress_here(agent, loaded, 2)
assert result.status == "compressed"
durable = session_db.get_messages_as_conversation("sid")
assert _live(durable) == _live(result.after_messages)
assert f"{history[-5]['content']}\n\n{history[-4]['content']}" in [m["content"] for m in durable]
assert _live(durable[-3:]) == _live(history[-3:])