fix(compression): retire completed todo snapshots safely
This commit is contained in:
@@ -2450,6 +2450,29 @@ def _strip_stale_todo_snapshot(content: Any) -> Any:
|
||||
return content
|
||||
|
||||
|
||||
def _todo_snapshot_is_only_content(content: Any, stripped: Any) -> bool:
|
||||
"""Return whether stripping the snapshot leaves no structured content.
|
||||
|
||||
Text snapshots are appended at the end of a string. Structured snapshots
|
||||
occupy their own text part, so only an empty remainder proves that the row
|
||||
was synthetic scaffolding alone. Text extraction is deliberately not used:
|
||||
image, audio, and future non-text parts are content that must survive.
|
||||
"""
|
||||
if isinstance(content, str) and isinstance(stripped, str):
|
||||
return not stripped.strip()
|
||||
if isinstance(content, list) and isinstance(stripped, list):
|
||||
return not stripped
|
||||
return False
|
||||
|
||||
|
||||
def _replace_message_content(message: dict, content: Any) -> None:
|
||||
"""Rewrite message content without allowing an old API sidecar to replay."""
|
||||
from agent.turn_context import drop_stale_api_content
|
||||
|
||||
message["content"] = content
|
||||
drop_stale_api_content(message)
|
||||
|
||||
|
||||
# Retention-parity notice (#84718): compaction re-injects the todo list
|
||||
# verbatim while skill instructions are pruned to [SKILL_PRUNED: ...] markers,
|
||||
# so the imperative crosses the boundary without the policy that governed it.
|
||||
@@ -2521,10 +2544,10 @@ def _merge_anchor_into_user_message(target: dict, anchor: dict) -> None:
|
||||
if isinstance(target_content, list)
|
||||
else [{"type": "text", "text": str(target_content or "")}]
|
||||
)
|
||||
target["content"] = anchor_parts + target_parts
|
||||
_replace_message_content(target, anchor_parts + target_parts)
|
||||
else:
|
||||
merged = f"{anchor_content or ''}\n\n{target_content or ''}".strip()
|
||||
target["content"] = merged
|
||||
_replace_message_content(target, merged)
|
||||
for flag in _SYNTHETIC_USER_FLAGS:
|
||||
target.pop(flag, None)
|
||||
|
||||
@@ -3824,6 +3847,53 @@ def compress_context(
|
||||
)
|
||||
|
||||
todo_snapshot = agent._todo_store.format_for_injection()
|
||||
# A non-empty store is authoritative even when every item is already
|
||||
# completed/cancelled and format_for_injection() therefore returns an
|
||||
# empty string. In that case remove the previous snapshot so completed
|
||||
# work is not resurrected. A truly empty store is different: fresh
|
||||
# gateway agents may be unable to rehydrate todo tool results after a
|
||||
# prior compaction, so the retained snapshot is the only surviving
|
||||
# record of pending work and must stay in place.
|
||||
_todo_has_items = getattr(agent._todo_store, "has_items", None)
|
||||
try:
|
||||
_todo_store_is_authoritative = bool(
|
||||
_todo_has_items()
|
||||
) if callable(_todo_has_items) else False
|
||||
except Exception:
|
||||
# A plugin/test double may implement only format_for_injection().
|
||||
# Unknown authority must preserve pending snapshot state rather than
|
||||
# risk deleting it during compression.
|
||||
_todo_store_is_authoritative = False
|
||||
if _todo_store_is_authoritative:
|
||||
for _todo_idx in range(len(compressed) - 1, -1, -1):
|
||||
_todo_message = compressed[_todo_idx]
|
||||
if not isinstance(_todo_message, dict) or _todo_message.get("role") != "user":
|
||||
continue
|
||||
_todo_content = _todo_message.get("content")
|
||||
_todo_stripped = _strip_stale_todo_snapshot(_todo_content)
|
||||
if _todo_stripped == _todo_content:
|
||||
continue
|
||||
if (
|
||||
_todo_message.get("_todo_snapshot_synthetic")
|
||||
and _todo_snapshot_is_only_content(
|
||||
_todo_content, _todo_stripped
|
||||
)
|
||||
):
|
||||
compressed.pop(_todo_idx)
|
||||
if _todo_idx < len(compressed):
|
||||
# A standalone snapshot can move away from the tail
|
||||
# after later turns arrive. Deleting it may expose two
|
||||
# assistant rows; use the normal replay repair so their
|
||||
# content/tool-call metadata is preserved consistently.
|
||||
agent._repair_message_sequence(compressed)
|
||||
else:
|
||||
_replace_message_content(_todo_message, _todo_stripped)
|
||||
# The row is no longer todo-only scaffolding. Other
|
||||
# synthetic flags, if any, remain authoritative and
|
||||
# _is_real_user_message() recomputes provenance from the
|
||||
# surviving content plus those flags.
|
||||
_todo_message.pop("_todo_snapshot_synthetic", None)
|
||||
break
|
||||
if todo_snapshot:
|
||||
# Retention parity (#84718): the snapshot below re-injects the
|
||||
# imperative verbatim. If this same boundary pruned skill bodies
|
||||
@@ -3865,8 +3935,9 @@ def compress_context(
|
||||
if isinstance(_stripped, str) and _stripped
|
||||
else todo_snapshot
|
||||
)
|
||||
_tail["content"] = _append_text_to_content(
|
||||
_stripped, _snapshot_text
|
||||
_replace_message_content(
|
||||
_tail,
|
||||
_append_text_to_content(_stripped, _snapshot_text),
|
||||
)
|
||||
merged = True
|
||||
elif _stripped != _tail.get("content") and not _message_text(
|
||||
@@ -3874,7 +3945,7 @@ def compress_context(
|
||||
).strip():
|
||||
# The tail was nothing but an earlier snapshot row —
|
||||
# refresh it in place instead of stacking a duplicate.
|
||||
_tail["content"] = todo_snapshot
|
||||
_replace_message_content(_tail, todo_snapshot)
|
||||
_tail["_todo_snapshot_synthetic"] = True
|
||||
merged = True
|
||||
if not merged:
|
||||
|
||||
@@ -1566,7 +1566,11 @@ class TestTodoSnapshotScaffoldingTails:
|
||||
agent = self._agent_with_todo(
|
||||
db,
|
||||
"PARENT_TODO_RESTRIP",
|
||||
{"role": "user", "content": previously_merged},
|
||||
{
|
||||
"role": "user",
|
||||
"content": previously_merged,
|
||||
"api_content": "stale wire copy containing the old task",
|
||||
},
|
||||
)
|
||||
|
||||
compressed, _ = agent._compress_context(
|
||||
@@ -1579,6 +1583,7 @@ class TestTodoSnapshotScaffoldingTails:
|
||||
assert "task A" in tail["content"]
|
||||
assert "old finished task" not in tail["content"]
|
||||
assert tail["content"].count(TODO_INJECTION_HEADER) == 1
|
||||
assert "api_content" not in tail
|
||||
assert not any(
|
||||
previous.get("role") == current.get("role") == "user"
|
||||
for previous, current in zip(compressed, compressed[1:])
|
||||
@@ -1620,6 +1625,288 @@ class TestTodoSnapshotScaffoldingTails:
|
||||
for message in compressed
|
||||
)
|
||||
|
||||
def test_empty_todo_store_removes_previous_compaction_snapshot(
|
||||
self, tmp_path: Path
|
||||
):
|
||||
"""Completed items make the store authoritative and clear old work."""
|
||||
from tools.todo_tool import TODO_INJECTION_HEADER
|
||||
|
||||
db = SessionDB(db_path=tmp_path / "state.db")
|
||||
session_id = "PARENT_TODO_CLEARED"
|
||||
db.create_session(session_id, source="cli")
|
||||
agent = _build_agent_with_db(db, session_id, platform="cli")
|
||||
stale_task = "- [ ] stale-task. This task was already cleared"
|
||||
agent.context_compressor.compress.return_value = [
|
||||
{"role": "user", "content": "[CONTEXT COMPACTION] summary"},
|
||||
{"role": "assistant", "content": "acknowledged"},
|
||||
{
|
||||
"role": "user",
|
||||
"content": (
|
||||
"Keep this user context.\n\n"
|
||||
f"{TODO_INJECTION_HEADER}\n{stale_task}"
|
||||
),
|
||||
"api_content": "stale wire copy containing the old snapshot",
|
||||
},
|
||||
]
|
||||
agent._todo_store.write(
|
||||
[{"id": "done", "content": "done thing", "status": "completed"}]
|
||||
)
|
||||
|
||||
compressed, _ = agent._compress_context(
|
||||
_msgs(), "sys", approx_tokens=120_000
|
||||
)
|
||||
|
||||
combined = "\n".join(str(m.get("content") or "") for m in compressed)
|
||||
assert "Keep this user context." in combined
|
||||
assert TODO_INJECTION_HEADER not in combined
|
||||
assert stale_task not in combined
|
||||
retained = next(
|
||||
message
|
||||
for message in compressed
|
||||
if "Keep this user context." in str(message.get("content") or "")
|
||||
)
|
||||
assert "api_content" not in retained
|
||||
|
||||
def test_cancelled_only_store_is_authoritative(self, tmp_path: Path):
|
||||
"""Cancelled work retires the old snapshot just like completed work."""
|
||||
from tools.todo_tool import TODO_INJECTION_HEADER
|
||||
|
||||
db = SessionDB(db_path=tmp_path / "state.db")
|
||||
session_id = "PARENT_TODO_CANCELLED"
|
||||
db.create_session(session_id, source="cli")
|
||||
agent = _build_agent_with_db(db, session_id, platform="cli")
|
||||
agent.context_compressor.compress.return_value = [
|
||||
{"role": "user", "content": "Keep the request"},
|
||||
{"role": "assistant", "content": "acknowledged"},
|
||||
{
|
||||
"role": "user",
|
||||
"content": f"{TODO_INJECTION_HEADER}\n- [ ] obsolete task",
|
||||
"_todo_snapshot_synthetic": True,
|
||||
},
|
||||
]
|
||||
agent._todo_store.write(
|
||||
[{"id": "nope", "content": "obsolete task", "status": "cancelled"}]
|
||||
)
|
||||
|
||||
compressed, _ = agent._compress_context(
|
||||
_msgs(), "sys", approx_tokens=120_000
|
||||
)
|
||||
|
||||
assert [message.get("role") for message in compressed] == [
|
||||
"user",
|
||||
"assistant",
|
||||
]
|
||||
assert all(
|
||||
TODO_INJECTION_HEADER not in str(message.get("content") or "")
|
||||
for message in compressed
|
||||
)
|
||||
|
||||
def test_structured_snapshot_only_tail_is_removed(self, tmp_path: Path):
|
||||
"""A flagged list row is removable only when no other part remains."""
|
||||
from tools.todo_tool import TODO_INJECTION_HEADER
|
||||
|
||||
db = SessionDB(db_path=tmp_path / "state.db")
|
||||
session_id = "PARENT_TODO_STRUCTURED_ONLY"
|
||||
db.create_session(session_id, source="cli")
|
||||
agent = _build_agent_with_db(db, session_id, platform="cli")
|
||||
agent.context_compressor.compress.return_value = [
|
||||
{"role": "user", "content": "Keep the request"},
|
||||
{"role": "assistant", "content": "acknowledged"},
|
||||
{
|
||||
"role": "user",
|
||||
"content": [
|
||||
{
|
||||
"type": "text",
|
||||
"text": f"{TODO_INJECTION_HEADER}\n- [ ] stale task",
|
||||
}
|
||||
],
|
||||
"_todo_snapshot_synthetic": True,
|
||||
},
|
||||
]
|
||||
agent._todo_store.write(
|
||||
[{"id": "done", "content": "stale task", "status": "completed"}]
|
||||
)
|
||||
|
||||
compressed, _ = agent._compress_context(
|
||||
_msgs(), "sys", approx_tokens=120_000
|
||||
)
|
||||
|
||||
assert [message.get("role") for message in compressed] == [
|
||||
"user",
|
||||
"assistant",
|
||||
]
|
||||
|
||||
def test_non_tail_snapshot_deletion_repairs_assistant_alternation(
|
||||
self, tmp_path: Path
|
||||
):
|
||||
from tools.todo_tool import TODO_INJECTION_HEADER
|
||||
|
||||
db = SessionDB(db_path=tmp_path / "state.db")
|
||||
session_id = "PARENT_TODO_MIDDLE"
|
||||
db.create_session(session_id, source="cli")
|
||||
agent = _build_agent_with_db(db, session_id, platform="cli")
|
||||
agent.context_compressor.compress.return_value = [
|
||||
{"role": "user", "content": "Keep the request"},
|
||||
{
|
||||
"role": "assistant",
|
||||
"content": "before snapshot",
|
||||
"api_content": "stale assistant wire copy",
|
||||
},
|
||||
{
|
||||
"role": "user",
|
||||
"content": f"{TODO_INJECTION_HEADER}\n- [ ] stale task",
|
||||
"_todo_snapshot_synthetic": True,
|
||||
},
|
||||
{"role": "assistant", "content": "after snapshot"},
|
||||
{"role": "user", "content": "new request"},
|
||||
]
|
||||
agent._todo_store.write(
|
||||
[{"id": "done", "content": "stale task", "status": "completed"}]
|
||||
)
|
||||
|
||||
compressed, _ = agent._compress_context(
|
||||
_msgs(), "sys", approx_tokens=120_000
|
||||
)
|
||||
|
||||
assert [message.get("role") for message in compressed] == [
|
||||
"user",
|
||||
"assistant",
|
||||
"user",
|
||||
]
|
||||
repaired = compressed[1]
|
||||
assert "before snapshot" in repaired["content"]
|
||||
assert "after snapshot" in repaired["content"]
|
||||
assert "api_content" not in repaired
|
||||
assert not any(
|
||||
previous.get("role") == current.get("role")
|
||||
for previous, current in zip(compressed, compressed[1:])
|
||||
)
|
||||
|
||||
def test_multimodal_content_survives_and_synthetic_provenance_clears(
|
||||
self, tmp_path: Path
|
||||
):
|
||||
from tools.todo_tool import TODO_INJECTION_HEADER
|
||||
|
||||
db = SessionDB(db_path=tmp_path / "state.db")
|
||||
session_id = "PARENT_TODO_MULTIMODAL_RETIRE"
|
||||
db.create_session(session_id, source="cli")
|
||||
agent = _build_agent_with_db(db, session_id, platform="cli")
|
||||
surviving_parts = [
|
||||
{"type": "text", "text": "Keep this caption"},
|
||||
{
|
||||
"type": "image_url",
|
||||
"image_url": {"url": "https://example.com/context.png"},
|
||||
},
|
||||
{"type": "input_audio", "input_audio": {"data": "audio-data"}},
|
||||
]
|
||||
agent.context_compressor.compress.return_value = [
|
||||
{"role": "user", "content": "summary"},
|
||||
{"role": "assistant", "content": "acknowledged"},
|
||||
{
|
||||
"role": "user",
|
||||
"content": surviving_parts
|
||||
+ [
|
||||
{
|
||||
"type": "text",
|
||||
"text": f"{TODO_INJECTION_HEADER}\n- [ ] stale task",
|
||||
}
|
||||
],
|
||||
"api_content": "stale wire copy containing the snapshot",
|
||||
"_todo_snapshot_synthetic": True,
|
||||
"unrelated_metadata": "keep-me",
|
||||
},
|
||||
]
|
||||
agent._todo_store.write(
|
||||
[{"id": "done", "content": "stale task", "status": "completed"}]
|
||||
)
|
||||
|
||||
input_msgs = [
|
||||
{
|
||||
"role": "user" if index % 2 == 0 else "assistant",
|
||||
"content": f"message {index} " + "x" * 1000,
|
||||
}
|
||||
for index in range(40)
|
||||
]
|
||||
compressed, _ = agent._compress_context(
|
||||
input_msgs, "sys", approx_tokens=120_000
|
||||
)
|
||||
|
||||
retained = next(
|
||||
message for message in compressed if isinstance(message.get("content"), list)
|
||||
)
|
||||
assert retained["content"] == surviving_parts
|
||||
assert retained["unrelated_metadata"] == "keep-me"
|
||||
assert "api_content" not in retained
|
||||
assert "_todo_snapshot_synthetic" not in retained
|
||||
|
||||
@pytest.mark.parametrize("authority_mode", ["missing", "raises"])
|
||||
def test_unknown_store_authority_preserves_snapshot(
|
||||
self, tmp_path: Path, authority_mode: str
|
||||
):
|
||||
"""Legacy or broken stores fail conservative, never erasing pending work."""
|
||||
from tools.todo_tool import TODO_INJECTION_HEADER
|
||||
|
||||
db = SessionDB(db_path=tmp_path / "state.db")
|
||||
session_id = f"PARENT_TODO_COMPAT_{authority_mode.upper()}"
|
||||
db.create_session(session_id, source="telegram")
|
||||
agent = _build_agent_with_db(db, session_id, platform="telegram")
|
||||
pending_task = "- [ ] pending-task. Preserve conservatively"
|
||||
agent.context_compressor.compress.return_value = [
|
||||
{"role": "user", "content": "Keep the request"},
|
||||
{"role": "assistant", "content": "acknowledged"},
|
||||
{
|
||||
"role": "user",
|
||||
"content": f"{TODO_INJECTION_HEADER}\n{pending_task}",
|
||||
"_todo_snapshot_synthetic": True,
|
||||
},
|
||||
]
|
||||
|
||||
class LegacyStore:
|
||||
def format_for_injection(self):
|
||||
return None
|
||||
|
||||
store = LegacyStore()
|
||||
if authority_mode == "raises":
|
||||
store.has_items = MagicMock(side_effect=RuntimeError("store unavailable"))
|
||||
agent._todo_store = store
|
||||
|
||||
compressed, _ = agent._compress_context(
|
||||
_msgs(), "sys", approx_tokens=120_000
|
||||
)
|
||||
|
||||
combined = "\n".join(str(message.get("content") or "") for message in compressed)
|
||||
assert TODO_INJECTION_HEADER in combined
|
||||
assert pending_task in combined
|
||||
|
||||
def test_unhydrated_empty_todo_store_preserves_pending_snapshot(
|
||||
self, tmp_path: Path
|
||||
):
|
||||
"""A fresh empty store must not erase pending work retained by compression."""
|
||||
from tools.todo_tool import TODO_INJECTION_HEADER
|
||||
|
||||
db = SessionDB(db_path=tmp_path / "state.db")
|
||||
session_id = "PARENT_TODO_UNHYDRATED"
|
||||
db.create_session(session_id, source="telegram")
|
||||
agent = _build_agent_with_db(db, session_id, platform="telegram")
|
||||
pending_task = "- [ ] pending-task. Continue after the next compaction"
|
||||
getattr(agent, "context_compressor").compress.return_value = [
|
||||
{"role": "user", "content": "[CONTEXT COMPACTION] summary"},
|
||||
{"role": "assistant", "content": "acknowledged"},
|
||||
{
|
||||
"role": "user",
|
||||
"content": f"{TODO_INJECTION_HEADER}\n{pending_task}",
|
||||
"_todo_snapshot_synthetic": True,
|
||||
},
|
||||
]
|
||||
|
||||
compressed, _ = agent._compress_context(
|
||||
_msgs(), "sys", approx_tokens=120_000
|
||||
)
|
||||
|
||||
combined = "\n".join(str(m.get("content") or "") for m in compressed)
|
||||
assert TODO_INJECTION_HEADER in combined
|
||||
assert pending_task in combined
|
||||
|
||||
|
||||
class TestArchivedParentActivityLabelsCleared:
|
||||
def test_parent_labels_cleared_after_rotation_child_lineage_intact(
|
||||
|
||||
Reference in New Issue
Block a user