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:
committed by
Teknium
parent
81a095cb1d
commit
526d135a96
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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:])
|
||||
|
||||
Reference in New Issue
Block a user