fix(memory): review follow-ups for the pre-compress checkpoint contract
Addresses the review on #93996: - gateway: hygiene and manual /compress load the memory provider only when compression.checkpoint_required is enabled (skip_memory=not required). The historical fast path — no provider init, no best-effort hook — is back for everyone who did not opt in, so default behavior is truly unchanged. - conversation_compression: assistant messages carrying both prose and tool_calls keep their prose in the checkpoint evidence (the tool_calls payload is stripped, the original message is not mutated); pure tool-call wrappers without prose are still dropped. - tests: legacy-database regression proving the _compressed_summary column is added by the declarative _reconcile_columns() path on a plain reopen (no version-gated migration needed — append_message works right after), plus coverage for the prose-preserving filter. - docs: providers must implement idempotent, content-keyed checkpoint writes — a fail-closed block means the next attempt re-runs on_pre_compress over largely the same transcript. Refs #93986
This commit is contained in:
committed by
Teknium
parent
3c31c44880
commit
70d0b1fffb
@@ -1554,10 +1554,12 @@ class _CompressionActivityHeartbeat:
|
||||
def _direct_messages_for_pre_compress_memory(messages: Any) -> list[dict[str, Any]]:
|
||||
"""Return direct user/assistant evidence safe for memory checkpointing.
|
||||
|
||||
Compression summaries are derivative context, not new source evidence. Tool
|
||||
rows, system messages, and assistant tool-call wrappers are likewise omitted
|
||||
so memory providers receive one normalized host contract instead of having
|
||||
to infer Hermes transcript internals independently.
|
||||
Compression summaries are derivative context, not new source evidence.
|
||||
Tool rows and system messages are likewise omitted so memory providers
|
||||
receive one normalized host contract instead of having to infer Hermes
|
||||
transcript internals independently. Assistant messages that carry both
|
||||
prose and ``tool_calls`` keep their prose (the ``tool_calls`` payload is
|
||||
stripped); pure tool-call wrappers without prose are dropped.
|
||||
"""
|
||||
# Deferred import: context_compressor imports turn_context, which imports
|
||||
# this module — a module-level import here would close that cycle
|
||||
@@ -1574,7 +1576,13 @@ def _direct_messages_for_pre_compress_memory(messages: Any) -> list[dict[str, An
|
||||
if message.get(COMPRESSED_SUMMARY_METADATA_KEY):
|
||||
continue
|
||||
if role == "assistant" and message.get("tool_calls"):
|
||||
continue
|
||||
content = message.get("content")
|
||||
has_prose = bool(
|
||||
content.strip() if isinstance(content, str) else content
|
||||
)
|
||||
if not has_prose:
|
||||
continue
|
||||
message = {k: v for k, v in message.items() if k != "tool_calls"}
|
||||
direct_messages.append(message)
|
||||
return direct_messages
|
||||
|
||||
|
||||
+22
-13
@@ -558,9 +558,10 @@ def _seed_hygiene_system_prompt(
|
||||
) -> bool:
|
||||
"""Keep gateway hygiene from rebuilding a live session's system prompt.
|
||||
|
||||
The hygiene helper loads the memory provider for the pre-compress hook,
|
||||
but it still runs outside the live session's fully initialized prompt
|
||||
environment (hygiene-only platform marker, no platform context files).
|
||||
The hygiene helper runs outside the live session's fully initialized
|
||||
prompt environment (hygiene-only platform marker, no platform context
|
||||
files; the memory provider is loaded only when
|
||||
``compression.checkpoint_required`` demands it).
|
||||
Compression is allowed to persist a system prompt, so letting that helper
|
||||
rebuild one would strip external provider blocks from the live session.
|
||||
Seed the exact persisted prompt instead. When no usable prompt can be
|
||||
@@ -19720,21 +19721,29 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
|
||||
exc_info=True,
|
||||
)
|
||||
_hyg_session_db = getattr(self._session_db, "_db", self._session_db)
|
||||
# Hygiene performs the same lossy rewrite as
|
||||
# normal compression. When the operator enabled
|
||||
# compression.checkpoint_required, the memory
|
||||
# provider must be loaded so the required
|
||||
# checkpoint is created before any transcript
|
||||
# mutation; otherwise keep the historical fast
|
||||
# path (no provider init, no best-effort hook)
|
||||
# for hygiene.
|
||||
from hermes_cli.config import load_config as _load_cfg
|
||||
from utils import is_truthy_value as _is_truthy
|
||||
|
||||
_hyg_checkpoint_required = _is_truthy(
|
||||
((_load_cfg() or {}).get("compression") or {}).get(
|
||||
"checkpoint_required"
|
||||
),
|
||||
default=False,
|
||||
)
|
||||
_hyg_agent = AIAgent(
|
||||
**_hyg_runtime,
|
||||
model=_hyg_model,
|
||||
max_iterations=4,
|
||||
quiet_mode=True,
|
||||
# Hygiene performs the same lossy rewrite
|
||||
# as normal compression, so it loads the
|
||||
# memory provider unconditionally (normal
|
||||
# compression never skips it either):
|
||||
# best-effort on_pre_compress runs for
|
||||
# hygiene too, and when
|
||||
# compression.checkpoint_required is
|
||||
# enabled the required checkpoint is
|
||||
# created before any transcript mutation.
|
||||
skip_memory=False,
|
||||
skip_memory=not _hyg_checkpoint_required,
|
||||
enabled_toolsets=["memory"],
|
||||
session_id=session_entry.session_id,
|
||||
session_db=_hyg_session_db,
|
||||
|
||||
+21
-12
@@ -4392,11 +4392,12 @@ class GatewaySlashCommandsMixin:
|
||||
runtime_kwargs["platform"] = platform_key
|
||||
runtime_kwargs["gateway_session_key"] = session_key
|
||||
|
||||
# The manual compression helper loads the memory provider (see
|
||||
# skip_memory=False below) but still runs outside the live
|
||||
# session's fully initialized prompt environment, and
|
||||
# _compress_context may persist its cached system prompt. Restore
|
||||
# the exact live-session prompt so provider blocks are retained.
|
||||
# The manual compression helper runs outside the live session's
|
||||
# fully initialized prompt environment (it loads the memory
|
||||
# provider only when compression.checkpoint_required demands it),
|
||||
# and _compress_context may persist its cached system prompt.
|
||||
# Restore the exact live-session prompt so provider blocks are
|
||||
# retained.
|
||||
session_row = None
|
||||
get_session = getattr(self._session_db, "get_session", None)
|
||||
if callable(get_session):
|
||||
@@ -4412,18 +4413,26 @@ class GatewaySlashCommandsMixin:
|
||||
exc_info=True,
|
||||
)
|
||||
|
||||
# This agent performs a lossy rewrite. When the operator enabled
|
||||
# compression.checkpoint_required, the memory provider must be
|
||||
# loaded so _compress_context() can create the required
|
||||
# pre-compression checkpoint; otherwise keep the historical fast
|
||||
# path (no provider init, no best-effort hook) for this helper.
|
||||
from hermes_cli.config import load_config as _load_cfg
|
||||
from utils import is_truthy_value as _is_truthy
|
||||
|
||||
_checkpoint_required = _is_truthy(
|
||||
((_load_cfg() or {}).get("compression") or {}).get(
|
||||
"checkpoint_required"
|
||||
),
|
||||
default=False,
|
||||
)
|
||||
tmp_agent = AIAgent(
|
||||
**runtime_kwargs,
|
||||
model=model,
|
||||
max_iterations=4,
|
||||
quiet_mode=True,
|
||||
# This agent performs the same lossy rewrite as normal
|
||||
# compression, so it loads the memory provider unconditionally
|
||||
# (normal compression never skips it either): best-effort
|
||||
# on_pre_compress runs for manual compression too, and when
|
||||
# compression.checkpoint_required is enabled the required
|
||||
# pre-compression checkpoint can be created.
|
||||
skip_memory=False,
|
||||
skip_memory=not _checkpoint_required,
|
||||
enabled_toolsets=["memory"],
|
||||
session_id=session_entry.session_id,
|
||||
session_db=getattr(self._session_db, "_db", self._session_db),
|
||||
|
||||
@@ -69,7 +69,7 @@ def test_direct_messages_filter_keeps_only_direct_source_evidence():
|
||||
{"role": "system", "content": "system prompt"},
|
||||
{"role": "user", "content": "durable user decision"},
|
||||
{"role": "assistant", "content": "direct assistant answer"},
|
||||
{"role": "assistant", "content": "call", "tool_calls": [{"id": "t1"}]},
|
||||
{"role": "assistant", "content": "", "tool_calls": [{"id": "t1"}]},
|
||||
{"role": "tool", "content": "tool output", "tool_call_id": "t1"},
|
||||
{
|
||||
"role": "assistant",
|
||||
@@ -87,6 +87,29 @@ def test_direct_messages_filter_keeps_only_direct_source_evidence():
|
||||
]
|
||||
|
||||
|
||||
def test_direct_messages_filter_keeps_prose_of_tool_call_messages():
|
||||
"""Assistant prose next to tool_calls is evidence; the payload is not."""
|
||||
messages = [
|
||||
{"role": "user", "content": "please scan the network"},
|
||||
{
|
||||
"role": "assistant",
|
||||
"content": "Scanning now — the last sweep found 26 hosts.",
|
||||
"tool_calls": [{"id": "t1", "function": {"name": "terminal"}}],
|
||||
},
|
||||
{"role": "assistant", "content": " ", "tool_calls": [{"id": "t2"}]},
|
||||
]
|
||||
|
||||
direct = _direct_messages_for_pre_compress_memory(messages)
|
||||
|
||||
assert [m["content"] for m in direct] == [
|
||||
"please scan the network",
|
||||
"Scanning now — the last sweep found 26 hosts.",
|
||||
]
|
||||
assert all("tool_calls" not in m for m in direct)
|
||||
# The original message list is not mutated.
|
||||
assert messages[1]["tool_calls"]
|
||||
|
||||
|
||||
def test_manager_advertises_checkpoint_capability_only_with_capable_provider():
|
||||
# The host allows one external provider per manager, so capability is
|
||||
# probed on two separate managers.
|
||||
@@ -182,3 +205,37 @@ def test_compressed_summary_marker_survives_restart_via_resume_history(tmp_path)
|
||||
|
||||
plain = reopened.get_messages_as_conversation("s1")
|
||||
assert all("_compressed_summary" not in m for m in plain)
|
||||
|
||||
|
||||
def test_compressed_summary_column_is_added_to_legacy_databases(tmp_path):
|
||||
"""Pre-upgrade databases gain the marker column via declarative reconcile.
|
||||
|
||||
``_init_schema()`` diffs live columns against SCHEMA_SQL on every
|
||||
writable open and ADDs whatever is missing, so a database created
|
||||
before this feature must accept marker writes after a plain reopen.
|
||||
"""
|
||||
import sqlite3
|
||||
|
||||
from hermes_state import SessionDB
|
||||
|
||||
db_path = tmp_path / "state.db"
|
||||
SessionDB(db_path)
|
||||
|
||||
# Simulate a pre-upgrade database: the marker column does not exist.
|
||||
conn = sqlite3.connect(db_path)
|
||||
conn.execute("ALTER TABLE messages DROP COLUMN _compressed_summary")
|
||||
conn.commit()
|
||||
legacy_cols = {
|
||||
row[1] for row in conn.execute("PRAGMA table_info(messages)")
|
||||
}
|
||||
conn.close()
|
||||
assert "_compressed_summary" not in legacy_cols
|
||||
|
||||
upgraded = SessionDB(db_path)
|
||||
upgraded.create_session("legacy", source="cli")
|
||||
upgraded.append_message(
|
||||
"legacy", "assistant", "derivative summary", _compressed_summary=True
|
||||
)
|
||||
|
||||
model_history, _display = upgraded.get_resume_conversations("legacy")
|
||||
assert model_history[-1].get("_compressed_summary") is True
|
||||
|
||||
@@ -166,11 +166,19 @@ uncompressed transcript is preserved, the compaction attempt errors with
|
||||
recovers. With the gate off (default), nothing changes for existing providers.
|
||||
|
||||
What your provider receives in both modes is normalized direct evidence:
|
||||
user/assistant text rows only — tool results, system messages, assistant
|
||||
tool-call wrappers, and prior compaction summaries are filtered host-side.
|
||||
Prior summaries are recognized via a persistent `_compressed_summary` message
|
||||
marker that survives process restarts, so a resumed session never feeds
|
||||
derivative summaries back into your archive.
|
||||
user/assistant text rows only — tool results, system messages, the
|
||||
`tool_calls` payload of assistant messages (their prose is kept), and prior
|
||||
compaction summaries are filtered host-side. Prior summaries are recognized
|
||||
via a persistent `_compressed_summary` message marker that survives process
|
||||
restarts, so a resumed session never feeds derivative summaries back into
|
||||
your archive.
|
||||
|
||||
**Checkpoints must be idempotent.** After a fail-closed block, the next
|
||||
compaction attempt calls `on_pre_compress()` again with the same transcript —
|
||||
and a transcript that grew only slightly produces largely overlapping
|
||||
evidence. Key your archive writes by content (for example a transcript
|
||||
digest) and upsert, so retries and overlaps deduplicate instead of
|
||||
accumulating duplicate archives.
|
||||
|
||||
Contract tests: `tests/agent/test_pre_compress_checkpoint_contract.py`.
|
||||
|
||||
|
||||
Reference in New Issue
Block a user