test(compaction): cover the reply anchor with invariant tests
tests/agent/test_compression_last_assistant_anchor.py holds 4 tests (twin no-op, reinsertion that keeps the reply row active with and without a todo store, older-assistant non-recovery, and the tool-chain skip with unique and reused tool-call ids). Rotation, persistence, abort-reset, concurrent-fork and todo-parity fixtures now model a conforming engine (a real user tail) instead of relying on the old stub that invented one; their assertions are unchanged.
This commit is contained in:
@@ -107,14 +107,6 @@ class TestAbortPathsResetPerAttemptState:
|
||||
{"role": "user", "content": "old question"},
|
||||
{"role": "assistant", "content": "old answer"},
|
||||
]
|
||||
# Bulk middle so the stub fold genuinely shrinks: the commit
|
||||
# guard keeps a dropped last reply live (#118900), and a 2-row
|
||||
# transcript cannot shrink once "old answer" is restored —
|
||||
# refusal (keep everything live) would be correct there, but
|
||||
# this test needs one successful compaction first.
|
||||
for i in range(10):
|
||||
original.append({"role": "user", "content": f"old bulk question {i} " + "x" * 60})
|
||||
original.append({"role": "assistant", "content": f"old bulk answer {i} " + "y" * 60})
|
||||
agent._flush_messages_to_session_db(original, [])
|
||||
compacted, history = self._in_place_success(agent, original)
|
||||
|
||||
|
||||
@@ -722,14 +722,6 @@ def test_durable_message_committed_before_lease_is_adopted(
|
||||
parent_sid = "PRE_LEASE_DURABLE_RACE"
|
||||
db.create_session(parent_sid, source="webui")
|
||||
db.append_message(parent_sid, "user", "old durable")
|
||||
# A genuinely compressible middle: the commit-site guards now keep a
|
||||
# dropped last assistant reply live (#118900), so a 2-row transcript can
|
||||
# no longer shrink once the late row is restored — refusal (keep
|
||||
# everything live) would be the correct outcome there. Bulk old rows give
|
||||
# the stub fold real middle to reclaim, as in production.
|
||||
for i in range(10):
|
||||
db.append_message(parent_sid, "user", f"old durable question {i} " + "x" * 60)
|
||||
db.append_message(parent_sid, "assistant", f"old durable answer {i} " + "y" * 60)
|
||||
|
||||
# Frontend takes its snapshot, then another producer commits before this
|
||||
# compressor acquires the lease.
|
||||
@@ -743,13 +735,10 @@ def test_durable_message_committed_before_lease_is_adopted(
|
||||
|
||||
agent.context_compressor.compress.assert_called_once()
|
||||
compressed_arg = agent.context_compressor.compress.call_args.args[0]
|
||||
assert [m["content"] for m in compressed_arg][0] == "old durable"
|
||||
assert [m["content"] for m in compressed_arg][-1] == "late committed before lease"
|
||||
assert len(compressed_arg) == 22
|
||||
# The late assistant reply the stub fold drops must stay live in the
|
||||
# committed child (#118900), not archived away with the folded middle.
|
||||
live_contents = [m["content"] for m in db.get_messages_as_conversation(agent.session_id)]
|
||||
assert "late committed before lease" in live_contents
|
||||
assert [m["content"] for m in compressed_arg] == [
|
||||
"old durable",
|
||||
"late committed before lease",
|
||||
]
|
||||
# Must not echo the stale snapshot — compression proceeded on the
|
||||
# adopted durable transcript (rotation publishes a child session).
|
||||
assert returned is not stale_snapshot
|
||||
|
||||
@@ -9,130 +9,150 @@ content stayed on disk.
|
||||
The built-in compressor keeps the latest visible reply in the tail
|
||||
(``_ensure_last_assistant_message_in_tail``), but plugin engines (e.g. LCM)
|
||||
implement their own ``compress()`` without that guard — and nothing enforces
|
||||
the invariant at the durable commit layer. ``_commit_compaction`` is the one
|
||||
the invariant at the durable commit layer. ``compress_context`` is the one
|
||||
choke point every engine and every path (threshold preflight, engine
|
||||
maintenance, manual /compress; in-place and rotation) routes through, so the
|
||||
guard lives here, mirroring ``_ensure_compressed_has_user_turn``.
|
||||
guard lives there, next to ``_ensure_compressed_has_user_turn``.
|
||||
"""
|
||||
|
||||
from agent.conversation_compression import (
|
||||
import pytest
|
||||
|
||||
from agent.context_compressor import SUMMARY_PREFIX
|
||||
from agent.conversation_compression_reply_anchor import (
|
||||
_ensure_compressed_keeps_last_assistant_reply,
|
||||
)
|
||||
|
||||
REPLY = "x" * 200 + " just-delivered verdict report the user is still reading."
|
||||
# Trailing whitespace on purpose: a conforming engine may hand the row back stripped.
|
||||
REPLY = "x" * 200 + " just-delivered verdict report the user is still reading. "
|
||||
|
||||
SUMMARY = "[Recent Summary (d0, node 490)]\n\nEarlier turns summarized."
|
||||
FOLLOWER = "Thanks, one more question."
|
||||
|
||||
|
||||
def _transcript(*, with_follower=True):
|
||||
def _assert_no_same_role_adjacency(messages):
|
||||
"""Strict alternation, tool rows included: a tool-call assistant is an assistant row
|
||||
too, so never assistant;assistant (or user;user / tool;tool) anywhere in the view."""
|
||||
roles = [m.get("role") for m in messages]
|
||||
assert not any(a == b for a, b in zip(roles, roles[1:])), f"same-role adjacency in {roles}"
|
||||
|
||||
|
||||
def _count_reply(messages):
|
||||
return sum(
|
||||
1 for m in messages
|
||||
if m.get("role") == "assistant" and (m.get("content") or "").strip() == REPLY.strip()
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"kept_shape",
|
||||
[
|
||||
pytest.param(lambda: REPLY.strip(), id="stripped"),
|
||||
pytest.param(lambda: [{"type": "text", "text": REPLY}], id="parts_list"),
|
||||
],
|
||||
)
|
||||
def test_normalized_twin_of_reply_counts_as_present(kept_shape):
|
||||
"""A conforming engine that keeps the reply whitespace-stripped or with parts-list
|
||||
content is NOT a dropped reply: exactly one copy stays, and the reinsertion must never
|
||||
manufacture assistant;assistant adjacency (strict-alternation invariant)."""
|
||||
original = [
|
||||
{"role": "user", "content": "Audit the auth module."},
|
||||
{"role": "assistant", "content": "Short ack."},
|
||||
{"role": "user", "content": "Go deeper, full report."},
|
||||
{"role": "assistant", "content": REPLY, "finish_reason": "stop"},
|
||||
{"role": "user", "content": FOLLOWER},
|
||||
]
|
||||
if with_follower:
|
||||
original.append({"role": "user", "content": "Thanks, one more question."})
|
||||
return original
|
||||
|
||||
|
||||
def test_s1_dropped_reply_reinserted_before_surviving_follower():
|
||||
"""S1 (the measured bug): engine folds the just-delivered long reply into
|
||||
the summary; the next-turn user message survives. The reply must come back
|
||||
ahead of that follower, not appended after it."""
|
||||
original = _transcript(with_follower=True)
|
||||
compressed = [
|
||||
{"role": "user", "content": SUMMARY},
|
||||
{"role": "user", "content": "Thanks, one more question."},
|
||||
{"role": "assistant", "content": kept_shape()},
|
||||
{"role": "user", "content": FOLLOWER},
|
||||
]
|
||||
|
||||
outcome = _ensure_compressed_keeps_last_assistant_reply(original, compressed)
|
||||
inserted = _ensure_compressed_keeps_last_assistant_reply(original, compressed)
|
||||
|
||||
assert outcome == "reinserted"
|
||||
kept = [m for m in compressed if m.get("role") == "assistant" and m.get("content") == REPLY]
|
||||
assert len(kept) == 1
|
||||
follower = next(m for m in compressed if m.get("content") == "Thanks, one more question.")
|
||||
assert compressed.index(kept[0]) < compressed.index(follower)
|
||||
roles = [m.get("role") for m in compressed]
|
||||
assert not any(a == b == "assistant" for a, b in zip(roles, roles[1:]))
|
||||
assert inserted is None
|
||||
assert len(compressed) == 3
|
||||
assistant_rows = [m for m in compressed if m.get("role") == "assistant"]
|
||||
assert len(assistant_rows) == 1
|
||||
_assert_no_same_role_adjacency(compressed)
|
||||
|
||||
|
||||
def test_s2_reply_already_kept_is_noop():
|
||||
"""S2: engine (built-in path) kept the reply in the tail — transcript must
|
||||
not be touched."""
|
||||
original = _transcript(with_follower=True)
|
||||
compressed = [
|
||||
{"role": "user", "content": SUMMARY},
|
||||
{"role": "assistant", "content": REPLY, "finish_reason": "stop"},
|
||||
{"role": "user", "content": "Thanks, one more question."},
|
||||
|
||||
def _tool_chain(*call_ids):
|
||||
"""Content-less tool-call assistants + their ``tool`` results, one round per id."""
|
||||
rows = []
|
||||
for call_id in call_ids:
|
||||
calls = [{"id": call_id, "type": "function", "function": {"name": "f", "arguments": "{}"}}]
|
||||
rows += [
|
||||
{"role": "assistant", "content": None, "tool_calls": calls},
|
||||
{"role": "tool", "content": "ok", "tool_call_id": call_id},
|
||||
]
|
||||
return rows
|
||||
|
||||
|
||||
def _seed(tmp_path, tail_rows, fold, follower):
|
||||
"""Real SessionDB with the persisted history, a stub engine returning ``fold``, and the
|
||||
live list ending on the just-typed ``follower`` turn (unpersisted, like production)."""
|
||||
import os
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
from hermes_state import SessionDB
|
||||
|
||||
db = SessionDB(db_path=tmp_path / "state.db")
|
||||
session_id = "E2E_118900_REPLY"
|
||||
db.create_session(session_id, source="cli")
|
||||
rows = [
|
||||
("user", "Audit the auth module."),
|
||||
("assistant", "Short ack."),
|
||||
("user", "Go deeper, full report."),
|
||||
]
|
||||
before = [dict(m) for m in compressed]
|
||||
|
||||
outcome = _ensure_compressed_keeps_last_assistant_reply(original, compressed)
|
||||
|
||||
assert outcome == "already_present"
|
||||
assert compressed == before
|
||||
|
||||
|
||||
def test_s3_trailing_reply_without_follower_appends_at_end():
|
||||
"""S3: compression ran before the next user turn; the reply is the last
|
||||
row. It must be appended back at the end."""
|
||||
original = _transcript(with_follower=False)
|
||||
compressed = [{"role": "user", "content": SUMMARY}]
|
||||
|
||||
outcome = _ensure_compressed_keeps_last_assistant_reply(original, compressed)
|
||||
|
||||
assert outcome == "reinserted"
|
||||
assert compressed[-1].get("role") == "assistant"
|
||||
assert compressed[-1].get("content") == REPLY
|
||||
|
||||
|
||||
def test_s4_empty_reasoning_only_reply_not_reinserted():
|
||||
"""S4: the #118755/#118738 family (text stranded in reasoning, empty
|
||||
content) is a different bug with its own fix — this guard stays out."""
|
||||
original = [
|
||||
{"role": "user", "content": "Think hard."},
|
||||
{"role": "assistant", "content": "", "reasoning": "long chain of thought"},
|
||||
for i in range(6):
|
||||
rows += [("user", f"bulk q {i} " + "q" * 80), ("assistant", f"bulk a {i} " + "a" * 80)]
|
||||
for role, content in [*rows, *tail_rows, ("assistant", REPLY)]:
|
||||
db.append_message(session_id, role, content)
|
||||
# Production live dicts carry _row_id + timestamp (session flush stamps both).
|
||||
messages = [
|
||||
*db.get_messages_as_conversation(session_id, include_row_ids=True),
|
||||
{"role": "user", "content": follower},
|
||||
]
|
||||
compressed = [{"role": "user", "content": SUMMARY}]
|
||||
before = [dict(m) for m in compressed]
|
||||
assert isinstance(messages[-2].get("_row_id"), int) and messages[-2].get("timestamp") is not None
|
||||
|
||||
outcome = _ensure_compressed_keeps_last_assistant_reply(original, compressed)
|
||||
with patch.dict(os.environ, {"OPENROUTER_API_KEY": "test-key"}):
|
||||
from run_agent import AIAgent
|
||||
|
||||
assert outcome == "no_reply"
|
||||
assert compressed == before
|
||||
agent = AIAgent(
|
||||
api_key="test-key",
|
||||
base_url="https://openrouter.ai/api/v1",
|
||||
model="test/model",
|
||||
platform="telegram",
|
||||
quiet_mode=True,
|
||||
session_db=db,
|
||||
session_id=session_id,
|
||||
skip_context_files=True,
|
||||
skip_memory=True,
|
||||
)
|
||||
engine = MagicMock()
|
||||
engine.compress.return_value = fold
|
||||
engine.compression_count = 1
|
||||
engine.last_prompt_tokens = 0
|
||||
engine.last_completion_tokens = 0
|
||||
engine._last_summary_error = None
|
||||
engine._last_compress_aborted = False
|
||||
engine._last_summary_auth_failure = False
|
||||
engine._last_aux_model_failure_model = None
|
||||
engine._last_aux_model_failure_error = None
|
||||
agent.context_compressor = engine
|
||||
return db, session_id, messages, agent
|
||||
|
||||
|
||||
def test_s5_tool_call_assistant_not_reinserted():
|
||||
"""S5: an assistant row carrying tool_calls must keep its atomic tool
|
||||
group — reinserting it alone would corrupt pairing."""
|
||||
original = [
|
||||
{"role": "user", "content": "Run it."},
|
||||
{
|
||||
"role": "assistant",
|
||||
"content": "Working.",
|
||||
"tool_calls": [{"id": "call-1", "function": {"name": "terminal", "arguments": "{}"}}],
|
||||
},
|
||||
]
|
||||
compressed = [{"role": "user", "content": SUMMARY}]
|
||||
before = [dict(m) for m in compressed]
|
||||
def _compress(agent, messages):
|
||||
from agent.conversation_compression import CompressionCommitFence
|
||||
|
||||
outcome = _ensure_compressed_keeps_last_assistant_reply(original, compressed)
|
||||
|
||||
assert outcome == "no_reply"
|
||||
assert compressed == before
|
||||
|
||||
|
||||
def test_s6_summary_rows_never_count_as_the_reply():
|
||||
"""S6: a compaction handoff must not satisfy the presence check — the live
|
||||
reply text must survive verbatim, not merely inside a summary."""
|
||||
original = _transcript(with_follower=False)
|
||||
compressed = [{"role": "user", "content": SUMMARY + "\n\n" + REPLY}]
|
||||
|
||||
outcome = _ensure_compressed_keeps_last_assistant_reply(original, compressed)
|
||||
|
||||
assert outcome == "reinserted"
|
||||
assert sum(1 for m in compressed if m.get("role") == "assistant" and m.get("content") == REPLY) == 1
|
||||
agent._persist_user_message_idx = len(messages) - 1
|
||||
out_messages, _ = agent._compress_context(
|
||||
messages, "sys", approx_tokens=120_000,
|
||||
commit_fence=CompressionCommitFence(),
|
||||
)
|
||||
return out_messages
|
||||
|
||||
|
||||
class TestEngineDropsReplyEndToEnd:
|
||||
@@ -141,70 +161,97 @@ class TestEngineDropsReplyEndToEnd:
|
||||
kept): the reply row must stay ACTIVE in state.db after commit, while
|
||||
genuinely old rows are still archived (compression itself keeps working)."""
|
||||
|
||||
def _agent(self, db, session_id):
|
||||
import os
|
||||
from unittest.mock import MagicMock, patch
|
||||
def _compress_and_read(self, db, session_id, agent, messages, tail_role):
|
||||
out_messages = _compress(agent, messages)
|
||||
live = db.get_messages_as_conversation(session_id)
|
||||
live_contents = [m.get("content") for m in live]
|
||||
# Compression still folds: the superseded early rows are archived, not live.
|
||||
assert "Short ack." not in live_contents
|
||||
assert "Go deeper, full report." not in live_contents
|
||||
# The model-facing list and the durable active set both alternate strictly and
|
||||
# neither ends on an assistant row: the trailing turn is the one the loop is about
|
||||
# to answer, so the old reply must never read as its answer.
|
||||
for view in (out_messages, live):
|
||||
_assert_no_same_role_adjacency(view)
|
||||
assert view[-1].get("role") == tail_role
|
||||
return live, live_contents
|
||||
|
||||
from agent.conversation_compression import CompressionCommitFence # noqa: F401 (re-export check)
|
||||
with patch.dict(os.environ, {"OPENROUTER_API_KEY": "test-key"}):
|
||||
from run_agent import AIAgent
|
||||
@pytest.mark.parametrize("todo_store", [False, True], ids=["dropped", "dropped_with_todo_store"])
|
||||
def test_reply_row_stays_active_after_commit(self, tmp_path, todo_store):
|
||||
"""LCM-shaped fold: summary + next-turn tail only; the long reply gone. With an
|
||||
active todo store the snapshot fold rewrites the trailing user row, so the reply
|
||||
must be back in place BEFORE it runs."""
|
||||
fold = [{"role": "user", "content": SUMMARY}, {"role": "user", "content": FOLLOWER}]
|
||||
db, session_id, messages, agent = _seed(tmp_path, [], fold, FOLLOWER)
|
||||
if todo_store:
|
||||
agent._todo_store._items = [{"id": "t1", "content": "inspect image", "status": "pending"}]
|
||||
|
||||
agent = AIAgent(
|
||||
api_key="test-key",
|
||||
base_url="https://openrouter.ai/api/v1",
|
||||
model="test/model",
|
||||
platform="telegram",
|
||||
quiet_mode=True,
|
||||
session_db=db,
|
||||
session_id=session_id,
|
||||
skip_context_files=True,
|
||||
skip_memory=True,
|
||||
)
|
||||
engine = MagicMock()
|
||||
# LCM-shaped fold: summary + next-turn tail only; the long reply gone.
|
||||
engine.compress.return_value = [
|
||||
live, live_contents = self._compress_and_read(db, session_id, agent, messages, tail_role="user")
|
||||
|
||||
assert _count_reply(live) == 1, "just-delivered reply left the active set after compaction"
|
||||
assert any(str(c).startswith(FOLLOWER) for c in live_contents)
|
||||
# Display projection (include_compacted): the reply renders once, not as archived
|
||||
# original + fresh twin.
|
||||
display = db.get_messages_as_conversation(
|
||||
session_id, include_ancestors=True, include_row_ids=True, include_compacted=True,
|
||||
)
|
||||
assert _count_reply(display) == 1
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"call_ids",
|
||||
[
|
||||
# Unique ids: the only slot is assistant-adjacent to the first call row.
|
||||
pytest.param(("call_1", "call_2"), id="unique_ids"),
|
||||
# A provider that reuses one id for every call (llama.cpp, see
|
||||
# ``_dedupe_tool_call_ids``): the ids cannot locate the slot, so the reply must
|
||||
# not be bound to a later round and land mid-chain.
|
||||
pytest.param(("call_0", "call_0", "call_0"), id="constant_ids"),
|
||||
],
|
||||
)
|
||||
def test_tool_chain_slot_is_skipped(self, tmp_path, call_ids, caplog):
|
||||
"""Mid-turn compaction: the engine folded the new user turn into the summary but kept
|
||||
the current tool chain (content-less tool-call assistants + tool results). The
|
||||
recovery is SKIPPED (logged) rather than committed as assistant;assistant or placed
|
||||
inside the chain: strict alternation beats recovery."""
|
||||
# Canonical summary prefix so the user-turn anchor can tell the summary from a
|
||||
# real turn and restore the dropped user turn behind the reply.
|
||||
fold = [{"role": "user", "content": SUMMARY_PREFIX + " earlier turns."}, *_tool_chain(*call_ids)]
|
||||
db, session_id, messages, agent = _seed(tmp_path, [], fold, FOLLOWER)
|
||||
# The tool rounds already live after the new user turn; the engine keeps them.
|
||||
messages += fold[1:]
|
||||
|
||||
with caplog.at_level("WARNING", logger="agent.conversation_compression_reply_anchor"):
|
||||
live, _ = self._compress_and_read(db, session_id, agent, messages, tail_role="tool")
|
||||
|
||||
# No legal slot: the reply is not recovered, the chain is untouched, and the skip is
|
||||
# logged so the operator can find the reply on disk. (The user-turn anchor's own
|
||||
# placement inside the chain is pre-existing and not asserted.)
|
||||
assert _count_reply(live) == 0
|
||||
assert all(m.get("tool_calls") for m in live if m.get("role") == "assistant")
|
||||
assert len([m for m in live if m.get("tool_calls")]) == len(call_ids)
|
||||
assert any("not reinserting" in r.getMessage() for r in caplog.records)
|
||||
|
||||
def test_older_assistant_at_slot_is_not_displaced(self, tmp_path):
|
||||
"""The engine kept an older assistant row ("noted") right where the reply would go,
|
||||
with "ok" twins around it: reinsertion would either create assistant;assistant or
|
||||
put the newest reply before an older one. Invariant beats recovery."""
|
||||
fold = [
|
||||
{"role": "user", "content": SUMMARY},
|
||||
{"role": "user", "content": "Thanks, one more question."},
|
||||
{"role": "assistant", "content": "noted"},
|
||||
{"role": "user", "content": "ok"},
|
||||
]
|
||||
engine.compression_count = 1
|
||||
engine.last_prompt_tokens = 0
|
||||
engine.last_completion_tokens = 0
|
||||
engine._last_summary_error = None
|
||||
engine._last_compress_aborted = False
|
||||
engine._last_summary_auth_failure = False
|
||||
engine._last_aux_model_failure_model = None
|
||||
engine._last_aux_model_failure_error = None
|
||||
agent.context_compressor = engine
|
||||
return agent
|
||||
|
||||
def test_reply_row_stays_active_after_commit(self, tmp_path):
|
||||
from agent.conversation_compression import CompressionCommitFence
|
||||
from hermes_state import SessionDB
|
||||
|
||||
db = SessionDB(db_path=tmp_path / "state.db")
|
||||
session_id = "E2E_118900_REPLY"
|
||||
db.create_session(session_id, source="cli")
|
||||
for role, content in [
|
||||
("user", "Audit the auth module."),
|
||||
("assistant", "Short ack."),
|
||||
("user", "Go deeper, full report."),
|
||||
("assistant", REPLY),
|
||||
]:
|
||||
db.append_message(session_id, role, content)
|
||||
|
||||
messages = [*db.get_messages_as_conversation(session_id),
|
||||
{"role": "user", "content": "Thanks, one more question."}]
|
||||
agent = self._agent(db, session_id)
|
||||
agent._persist_user_message_idx = len(messages) - 1
|
||||
|
||||
agent._compress_context(
|
||||
messages, "sys", approx_tokens=120_000,
|
||||
commit_fence=CompressionCommitFence(),
|
||||
db, session_id, messages, agent = _seed(
|
||||
tmp_path, [("user", "ok"), ("assistant", "noted"), ("user", "ok")], fold, "ok",
|
||||
)
|
||||
|
||||
live = [m.get("content") for m in db.get_messages_as_conversation(session_id)]
|
||||
assert REPLY in live, "just-delivered reply left the active set after compaction"
|
||||
assert "Thanks, one more question." in live
|
||||
# Compression still folds: the superseded early rows are archived, not live.
|
||||
assert "Short ack." not in live
|
||||
assert "Go deeper, full report." not in live
|
||||
out_messages = _compress(agent, messages)
|
||||
|
||||
live = db.get_messages_as_conversation(session_id)
|
||||
live_contents = [m.get("content") for m in live]
|
||||
for view in (out_messages, live):
|
||||
_assert_no_same_role_adjacency(view)
|
||||
assert view[-1].get("role") == "user"
|
||||
# The reply is not forced in beside "noted", and chronology holds ("noted" still
|
||||
# precedes the last "ok" that followed it).
|
||||
assert _count_reply(live) == 0
|
||||
assert live_contents.index("noted") < len(live_contents) - 1 == live_contents.index("ok")
|
||||
|
||||
@@ -211,14 +211,9 @@ class TestFlushAfterCompression:
|
||||
awaiting_real_usage_after_compression = False
|
||||
|
||||
def compress(self, _messages, **_kwargs):
|
||||
# Conforming-engine shape (#118900): the retained tail row is
|
||||
# the transcript's own last reply, kept verbatim — the commit
|
||||
# guard reinserts a dropped last reply, so inventing a new
|
||||
# tail row here would (correctly) come back with the original
|
||||
# alongside it.
|
||||
return [
|
||||
{"role": "user", "content": "[summary] earlier state"},
|
||||
{"role": "assistant", "content": "old answer"},
|
||||
{"role": "assistant", "content": "retained tail"},
|
||||
]
|
||||
|
||||
class AbortCompressor:
|
||||
@@ -274,7 +269,7 @@ class TestFlushAfterCompression:
|
||||
agent.session_id
|
||||
)] == [
|
||||
"[summary] earlier state",
|
||||
"old answer",
|
||||
"retained tail",
|
||||
"new request",
|
||||
"new answer",
|
||||
]
|
||||
|
||||
@@ -26,7 +26,7 @@ from unittest.mock import MagicMock, patch
|
||||
|
||||
import pytest
|
||||
|
||||
from agent.context_compressor import ContextCompressor, _DB_PERSISTED_MARKER
|
||||
from agent.context_compressor import SUMMARY_PREFIX, ContextCompressor, _DB_PERSISTED_MARKER
|
||||
from agent.conversation_compression import (
|
||||
CompressionCommitFence,
|
||||
_is_real_user_message,
|
||||
@@ -70,6 +70,32 @@ def _build_agent_with_db(db: SessionDB, session_id: str, platform: str = "telegr
|
||||
return agent
|
||||
|
||||
|
||||
def _seed_bulk_head(db: SessionDB, parent: str) -> None:
|
||||
"""Durable parent transcript: 10 bulk turns + the persisted question/answer pair.
|
||||
|
||||
Bulk head so the stub fold genuinely shrinks: the no-growth commit guard
|
||||
refuses a 3-row candidate that is not smaller than a 3-row original; the
|
||||
transcript's last reply stays "persisted answer" (#118900 guard shape).
|
||||
"""
|
||||
for i in range(10):
|
||||
db.append_message(parent, "user", f"bulk question {i} " + "x" * 60)
|
||||
db.append_message(parent, "assistant", f"bulk answer {i} " + "y" * 60)
|
||||
db.append_message(parent, "user", "persisted question")
|
||||
db.append_message(parent, "assistant", "persisted answer")
|
||||
|
||||
|
||||
def _conforming_fold(*tail: dict) -> list:
|
||||
"""Conforming-engine shape (#118900): a real handoff summary (recognized as
|
||||
scaffolding, not human intent) plus the transcript's own last reply kept
|
||||
verbatim ahead of the scaffolding tail, so the commit guard sees the reply
|
||||
present and the user-turn anchor still takes the MERGED branch these tests exercise."""
|
||||
return [
|
||||
{"role": "user", "content": SUMMARY_PREFIX + " earlier turns"},
|
||||
{"role": "assistant", "content": "persisted answer"},
|
||||
*tail,
|
||||
]
|
||||
|
||||
|
||||
def _msgs(n=20):
|
||||
return [{"role": "user", "content": f"m{i}"} for i in range(n)]
|
||||
|
||||
@@ -221,12 +247,6 @@ class TestRotationChildFlushDedup:
|
||||
db.create_session(parent, source="cli")
|
||||
db.append_message(parent, "user", "persisted question")
|
||||
db.append_message(parent, "assistant", "persisted answer")
|
||||
# Bulk middle so the stub fold genuinely shrinks: the commit guard
|
||||
# keeps a dropped last reply live (#118900), and without real middle
|
||||
# to reclaim the no-growth guard correctly refuses the attempt.
|
||||
for i in range(10):
|
||||
db.append_message(parent, "user", f"bulk question {i} " + "x" * 60)
|
||||
db.append_message(parent, "assistant", f"bulk answer {i} " + "y" * 60)
|
||||
|
||||
loaded = db.get_messages_as_conversation(parent)
|
||||
messages = [*loaded, {"role": "user", "content": "live question"}]
|
||||
@@ -515,8 +535,7 @@ class TestRotationChildFlushDedup:
|
||||
db = SessionDB(db_path=tmp_path / "state.db")
|
||||
parent = "PARENT_ROT_DRIFTED_MERGED"
|
||||
db.create_session(parent, source="cli")
|
||||
db.append_message(parent, "user", "persisted question")
|
||||
db.append_message(parent, "assistant", "persisted answer")
|
||||
_seed_bulk_head(db, parent)
|
||||
|
||||
loaded = db.get_messages_as_conversation(parent)
|
||||
messages = [
|
||||
@@ -531,13 +550,13 @@ class TestRotationChildFlushDedup:
|
||||
|
||||
agent = _build_agent_with_db(db, parent)
|
||||
agent._persist_user_message_idx = len(messages) - 1
|
||||
agent.context_compressor.compress.return_value = [
|
||||
agent.context_compressor.compress.return_value = _conforming_fold(
|
||||
{
|
||||
"role": "user",
|
||||
"content": "handoff scaffolding",
|
||||
"_todo_snapshot_synthetic": True,
|
||||
},
|
||||
]
|
||||
)
|
||||
|
||||
real_flush = agent._flush_messages_to_session_db
|
||||
with patch.object(
|
||||
@@ -572,21 +591,20 @@ class TestRotationChildFlushDedup:
|
||||
db = SessionDB(db_path=tmp_path / "state.db")
|
||||
parent = "PARENT_ROT_MERGED_LIVE"
|
||||
db.create_session(parent, source="cli")
|
||||
db.append_message(parent, "user", "persisted question")
|
||||
db.append_message(parent, "assistant", "persisted answer")
|
||||
_seed_bulk_head(db, parent)
|
||||
|
||||
loaded = db.get_messages_as_conversation(parent)
|
||||
messages = [*loaded, {"role": "user", "content": "live question"}]
|
||||
|
||||
agent = _build_agent_with_db(db, parent)
|
||||
agent._persist_user_message_idx = len(messages) - 1
|
||||
agent.context_compressor.compress.return_value = [
|
||||
agent.context_compressor.compress.return_value = _conforming_fold(
|
||||
{
|
||||
"role": "user",
|
||||
"content": "scaffolding",
|
||||
"_todo_snapshot_synthetic": True,
|
||||
},
|
||||
]
|
||||
)
|
||||
|
||||
real_flush = agent._flush_messages_to_session_db
|
||||
with patch.object(
|
||||
@@ -627,8 +645,7 @@ class TestRotationChildFlushDedup:
|
||||
db = SessionDB(db_path=tmp_path / "state.db")
|
||||
parent = "PARENT_ROT_ADOPT_DIVERGE"
|
||||
db.create_session(parent, source="cli")
|
||||
db.append_message(parent, "user", "persisted question")
|
||||
db.append_message(parent, "assistant", "persisted answer")
|
||||
_seed_bulk_head(db, parent)
|
||||
|
||||
# (b) Old live list object kept alive; divergence set so the twin is
|
||||
# the SAME object the guard scans (not a fresh copy).
|
||||
@@ -656,18 +673,18 @@ class TestRotationChildFlushDedup:
|
||||
# condition (durable parent longer than the caller snapshot) fires.
|
||||
db.append_message(parent, "user", "live question")
|
||||
durable_check = db.get_messages_as_conversation(parent)
|
||||
assert len(durable_check) == 3 > len(stale_snapshot) == 2
|
||||
assert len(durable_check) == len(loaded) + 1 > len(stale_snapshot) == 2
|
||||
# Sync the twin's timestamp to the committed row so the guard's
|
||||
# exact-timestamp twin scan matches the adopted anchor.
|
||||
old_live_list[-1]["timestamp"] = durable_check[-1]["timestamp"]
|
||||
|
||||
agent.context_compressor.compress.return_value = [
|
||||
agent.context_compressor.compress.return_value = _conforming_fold(
|
||||
{
|
||||
"role": "user",
|
||||
"content": "handoff scaffolding",
|
||||
"_todo_snapshot_synthetic": True,
|
||||
},
|
||||
]
|
||||
)
|
||||
|
||||
# The pre-publish flush fails, so the post-rotation flush below is the
|
||||
# only writer of the live view.
|
||||
@@ -829,8 +846,7 @@ class TestRotationChildFlushDedup:
|
||||
db = SessionDB(db_path=tmp_path / "state.db")
|
||||
parent = "PARENT_ROT_LIST_MERGED"
|
||||
db.create_session(parent, source="cli")
|
||||
db.append_message(parent, "user", "persisted question")
|
||||
db.append_message(parent, "assistant", "persisted answer")
|
||||
_seed_bulk_head(db, parent)
|
||||
|
||||
loaded = db.get_messages_as_conversation(parent)
|
||||
messages = [
|
||||
@@ -843,13 +859,13 @@ class TestRotationChildFlushDedup:
|
||||
|
||||
agent = _build_agent_with_db(db, parent)
|
||||
agent._persist_user_message_idx = len(messages) - 1
|
||||
agent.context_compressor.compress.return_value = [
|
||||
agent.context_compressor.compress.return_value = _conforming_fold(
|
||||
{
|
||||
"role": "user",
|
||||
"content": [{"type": "text", "text": "scaffolding"}],
|
||||
"_todo_snapshot_synthetic": True,
|
||||
},
|
||||
]
|
||||
)
|
||||
|
||||
real_flush = agent._flush_messages_to_session_db
|
||||
with patch.object(
|
||||
@@ -902,8 +918,7 @@ class TestRotationChildFlushDedup:
|
||||
db = SessionDB(db_path=tmp_path / "state.db")
|
||||
parent = "PARENT_ROT_STAMPED_TWIN"
|
||||
db.create_session(parent, source="cli")
|
||||
db.append_message(parent, "user", "persisted question")
|
||||
db.append_message(parent, "assistant", "persisted answer")
|
||||
_seed_bulk_head(db, parent)
|
||||
|
||||
loaded = db.get_messages_as_conversation(parent)
|
||||
messages = [
|
||||
@@ -927,13 +942,13 @@ class TestRotationChildFlushDedup:
|
||||
},
|
||||
{"role": "user", "content": "live question"},
|
||||
]
|
||||
agent.context_compressor.compress.return_value = [
|
||||
agent.context_compressor.compress.return_value = _conforming_fold(
|
||||
{
|
||||
"role": "user",
|
||||
"content": "handoff scaffolding",
|
||||
"_todo_snapshot_synthetic": True,
|
||||
},
|
||||
]
|
||||
)
|
||||
|
||||
with patch.object(
|
||||
agent,
|
||||
|
||||
@@ -292,16 +292,20 @@ class TestNoticeStripLifecycle:
|
||||
parent = "PARENT_SKILL_TODO_RESTRIP"
|
||||
db.create_session(parent, source="cli")
|
||||
agent = _build_agent_with_db(db, parent)
|
||||
original = _msgs()
|
||||
agent.context_compressor.compress.return_value = [
|
||||
{"role": "user", "content": summary},
|
||||
{"role": "assistant", "content": "acknowledged"},
|
||||
# Conforming-engine shape (#118900): the kept assistant row is the
|
||||
# transcript's own last reply, so the commit guard sees it present
|
||||
# and the snapshot still merges into the trailing human row.
|
||||
{"role": "assistant", "content": original[-1]["content"]},
|
||||
{"role": "user", "content": stale_tail},
|
||||
]
|
||||
agent._todo_store._items = [
|
||||
{"id": "t1", "content": "fresh task", "status": "pending"}
|
||||
]
|
||||
compressed, _ = agent._compress_context(
|
||||
_msgs(), "sys", approx_tokens=120_000
|
||||
original, "sys", approx_tokens=120_000
|
||||
)
|
||||
db.close()
|
||||
tail_text = str(compressed[-1]["content"])
|
||||
|
||||
Reference in New Issue
Block a user