fix(review): distinguish explicit /refine from unattended reviews and surface staged consolidations

Review follow-up on #105944 (#105921):

- explicit /refine forks now run under the refine_review write origin
  (explicit flows from the CLI/gateway handlers through
  _spawn_background_review_now and spawn_background_review_thread down
  to build_cache_parity_fork), so a user-requested review keeps the
  full memory operation set; only automatic reviews stay behind the
  unattended delete gate.
- the unattended delete gate now stages the denied replace/remove (or
  whole batch) into the pending store instead of dropping it: the
  fork's own review summary is never published, so a plain denial lost
  the consolidation request with no surfacing path. The staged proposal
  carries a proposal_staged marker that summarize surfaces as an action
  line, and a staging failure still fails closed to a plain denial.
- regression tests: explicit-path origin pass-through, refine_review
  keeping replace working, near-limit denial end to end (add rejected
  by budget -> replace staged -> proposal surfaces, store unchanged).
This commit is contained in:
liuhao1024
2026-09-09 05:48:28 +08:00
committed by kshitij
parent 1571f502a9
commit c0714575c3
6 changed files with 259 additions and 33 deletions

View File

@@ -613,7 +613,9 @@ def _prior_tool_keys(prior_snapshot: List[Dict]) -> Tuple[set, set]:
def _action_lines(data: Dict, detail: Dict, verbose: bool) -> List[str]:
"""Summary line(s) for one successful notify-tool result (``[]`` when nothing to report)."""
if data.get("staged"):
return []
# The fork's own review summary is never published back, so an unattended-review
# consolidation proposal must surface here or it is silently lost (#105921).
return [data["message"]] if data.get("proposal_staged") and data.get("message") else []
message = data.get("message", "")
target = data.get("target", "") or detail.get("target", "")
is_skill = detail.get("tool") == "skill_manage"
@@ -1014,11 +1016,16 @@ def _release_fork_clients(review_agent: Any) -> None:
def _run_review_fork(
agent: Any, messages_snapshot: List[Dict], prompt: str, task_cfg: Optional[Dict[str, Any]],
review_run: Optional[_BackgroundReviewRun], st: _ReviewForkState, review_memory: bool = False,
explicit: bool = False,
) -> None:
"""Fork phase (inside thread-scoped silence): build the fork, run the prompt under the tool
whitelist, snapshot its messages/usage, release its clients. Partial progress lands on ``st``
so the caller's error path still sees usage and the fork to clean up."""
st.review_agent, _rt, _routed = build_cache_parity_fork(agent, task_cfg, max_iterations=_REVIEW_MAX_ITERATIONS)
so the caller's error path still sees usage and the fork to clean up. ``explicit`` (/refine)
forks under the ``refine_review`` origin: a user-requested review keeps the full memory
operation set — the unattended-only delete gate must not treat it as automatic."""
st.review_agent, _rt, _routed = build_cache_parity_fork(
agent, task_cfg, max_iterations=_REVIEW_MAX_ITERATIONS,
write_origin="refine_review" if explicit else "background_review")
_track_review_fork(agent, st.review_agent, register=True)
from hermes_cli.plugins import set_thread_tool_whitelist, clear_thread_tool_whitelist
review_whitelist, configured_extra_tools = _review_tool_whitelist(st.review_agent, task_cfg, review_memory)
@@ -1081,7 +1088,7 @@ def _publish_review_summary(agent: Any, actions: List[str]) -> None:
def _run_review_in_thread(
agent: Any, messages_snapshot: List[Dict], prompt: str,
task_cfg: Optional[Dict[str, Any]] = None, review_run: Optional[_BackgroundReviewRun] = None,
review_memory: bool = False,
review_memory: bool = False, explicit: bool = False,
) -> None:
"""Daemon-thread worker: build the fork, run the prompt, surface the action summary via
``agent._safe_print`` / ``background_review_callback``. ``review_run`` (from
@@ -1116,7 +1123,7 @@ def _run_review_in_thread(
# their console output (#55769 / #55925). ``thread_scoped_silence`` routes only this thread's writes
# to devnull and leaves all other threads on the real streams.
with thread_scoped_silence():
_run_review_fork(agent, messages_snapshot, prompt, task_cfg, review_run, st, review_memory)
_run_review_fork(agent, messages_snapshot, prompt, task_cfg, review_run, st, review_memory, explicit)
# A buggy/legacy tool response shape must NOT take down the whole review (the outer
# except would discard every action the fork DID complete), so coerce to an empty list.
try:
@@ -1173,11 +1180,14 @@ def spawn_background_review_thread(
agent: Any, messages_snapshot: List[Dict], review_memory: bool = False,
review_skills: bool = False, focus: Optional[str] = None,
task_cfg: Optional[Dict[str, Any]] = None, review_run: Optional[_BackgroundReviewRun] = None,
explicit: bool = False,
):
"""Return ``(target, prompt)``; the caller builds the ``threading.Thread`` so test patches of
``run_agent.threading.Thread`` keep working. ``focus`` (``/refine [instructions]``) is appended
to the chosen prompt; automatic reviews pass ``None``. ``task_cfg`` is the pre-loaded
``auxiliary.background_review`` block; when omitted it is read once here."""
``auxiliary.background_review`` block; when omitted it is read once here. ``explicit``
(/refine) propagates to the fork's write origin so user-requested reviews keep the full
memory operation set."""
if task_cfg is None:
task_cfg = _background_review_task_config()
# Per-agent overrides (agent._MEMORY_REVIEW_PROMPT etc.) keep working.
@@ -1192,7 +1202,7 @@ def spawn_background_review_thread(
def _target() -> None: # resolves _run_review_in_thread at call time (tests patch it)
_run_review_in_thread(
agent, messages_snapshot, prompt, task_cfg=task_cfg, review_run=review_run,
review_memory=review_memory)
review_memory=review_memory, explicit=explicit)
return _target, prompt

View File

@@ -758,7 +758,8 @@ class AIAgent(
# deferral) goes through. See #100795.
from agent.turn_finalizer import _clone_background_review_messages
kwargs = dict(messages_snapshot=_clone_background_review_messages(messages_snapshot),
review_memory=review_memory, review_skills=review_skills, focus=focus, task_cfg=task_cfg)
review_memory=review_memory, review_skills=review_skills, focus=focus, task_cfg=task_cfg,
explicit=explicit)
if focus is None and not explicit and _review_should_defer(self, task_cfg):
from agent.review_idle_queue import QUEUE
QUEUE.enqueue(self, _review_queue_key(self), kwargs)
@@ -767,12 +768,15 @@ class AIAgent(
def _spawn_background_review_now(self, messages_snapshot: List[Dict], review_memory: bool = False,
review_skills: bool = False, focus: Optional[str] = None,
task_cfg: Optional[Dict[str, Any]] = None, _requeue_attempts: int = 0) -> None:
task_cfg: Optional[Dict[str, Any]] = None, _requeue_attempts: int = 0,
explicit: bool = False) -> None:
"""Spawn the background memory/skill review thread.
``threading.Thread`` is constructed here so tests patching ``run_agent.threading.Thread`` keep working.
``focus`` is /refine steering text; ``task_cfg`` is the pre-loaded config block (None on direct calls).
A deferred review preempted by a live turn is requeued (bounded) rather than lost.
``explicit`` (/refine) forks under the ``refine_review`` write origin, keeping the full
memory operation set. A deferred review preempted by a live turn is requeued (bounded)
rather than lost.
"""
from agent.background_review import (
finish_background_review_run, prepare_background_review_run, spawn_background_review_thread,
@@ -785,14 +789,15 @@ class AIAgent(
try:
target, _prompt = spawn_background_review_thread(
self, messages_snapshot, review_memory=review_memory, review_skills=review_skills,
focus=focus, task_cfg=task_cfg, review_run=review_run,
focus=focus, task_cfg=task_cfg, review_run=review_run, explicit=explicit,
)
def _target_with_requeue() -> None:
target()
self._maybe_requeue_preempted_review(review_run, dict(
messages_snapshot=messages_snapshot, review_memory=review_memory, review_skills=review_skills,
focus=focus, task_cfg=task_cfg, _requeue_attempts=_requeue_attempts + 1))
focus=focus, task_cfg=task_cfg, _requeue_attempts=_requeue_attempts + 1,
explicit=explicit))
# Carry the active profile into the review thread so MEMORY.md / skill review writes land in the
# right profile.

View File

@@ -53,7 +53,8 @@ class TestSpawnForwardsScope:
def test_target_passes_review_memory_to_worker(self):
captured = {}
def fake_worker(agent, messages_snapshot, prompt, task_cfg=None, review_run=None, review_memory=False):
def fake_worker(agent, messages_snapshot, prompt, task_cfg=None, review_run=None,
review_memory=False, explicit=False):
captured["review_memory"] = review_memory
agent = SimpleNamespace()
@@ -67,3 +68,146 @@ class TestSpawnForwardsScope:
agent, [], review_memory=True, review_skills=False)
target()
assert captured["review_memory"] is True
class TestExplicitRefineOrigin:
"""``/refine`` (explicit) must not inherit the unattended-review origin: the user asked
for that review, so its fork keeps the full memory operation set and the delete gate
does not apply (#105921 review follow-up)."""
def test_target_passes_explicit_to_worker(self):
captured = {}
def fake_worker(agent, messages_snapshot, prompt, task_cfg=None, review_run=None,
review_memory=False, explicit=False):
captured["explicit"] = explicit
with patch.object(bg, "_run_review_in_thread", fake_worker):
target, _prompt = bg.spawn_background_review_thread(
SimpleNamespace(), [], review_memory=True, explicit=True)
target()
assert captured["explicit"] is True
def test_explicit_fork_uses_refine_review_origin(self):
captured = {}
def fake_build(agent, task_cfg=None, *, max_iterations, write_origin="background_review"):
captured["write_origin"] = write_origin
fork = SimpleNamespace(
_memory_enabled=True, _user_profile_enabled=False,
run_conversation=lambda **kw: None, _session_messages=[])
return fork, {}, False
noop = lambda *a, **k: None
with patch.object(bg, "build_cache_parity_fork", fake_build), \
patch.object(bg, "_track_review_fork", noop), \
patch.object(bg, "_snapshot_review_usage", lambda a: {}), \
patch.object(bg, "_record_review_usage_to_parent", noop), \
patch.object(bg, "finish_background_review_run", noop), \
patch.object(bg, "_release_fork_clients", noop):
bg._run_review_fork(SimpleNamespace(), [], "p", None, None, bg._ReviewForkState(), True, True)
assert captured["write_origin"] == "refine_review"
def test_automatic_fork_keeps_background_review_origin(self):
captured = {}
def fake_build(agent, task_cfg=None, *, max_iterations, write_origin="background_review"):
captured["write_origin"] = write_origin
fork = SimpleNamespace(
_memory_enabled=True, _user_profile_enabled=False,
run_conversation=lambda **kw: None, _session_messages=[])
return fork, {}, False
noop = lambda *a, **k: None
with patch.object(bg, "build_cache_parity_fork", fake_build), \
patch.object(bg, "_track_review_fork", noop), \
patch.object(bg, "_snapshot_review_usage", lambda a: {}), \
patch.object(bg, "_record_review_usage_to_parent", noop), \
patch.object(bg, "finish_background_review_run", noop), \
patch.object(bg, "_release_fork_clients", noop):
bg._run_review_fork(SimpleNamespace(), [], "p", None, None, bg._ReviewForkState(), True, False)
assert captured["write_origin"] == "background_review"
class TestConsolidationProposalSurfaces:
"""The fork's own review summary is never published back, so a consolidation the delete
gate staged must surface through ``summarize_background_review_actions`` — otherwise the
near-limit denial path drops both the requested update and the proposal, silently (#105921)."""
def _store(self, tmp_path, monkeypatch):
import json as _json
from tools.memory_tool_store import MemoryStore
monkeypatch.setenv("HERMES_HOME", str(tmp_path / "home"))
monkeypatch.setattr("tools.memory_tool.get_memory_dir", lambda: tmp_path)
store = MemoryStore(memory_char_limit=500, user_char_limit=300)
store.load_from_disk()
return store
def test_staged_proposal_surfaces_in_summary(self, tmp_path, monkeypatch):
import json
from tools.memory_tool import memory_tool
from tools.skill_provenance import set_current_write_origin, reset_current_write_origin
store = self._store(tmp_path, monkeypatch)
assert store.add("memory", "standing rule entry")["success"] is True
token = set_current_write_origin("background_review")
try:
raw = memory_tool(
action="replace", old_text="standing rule", content="consolidated entry", store=store)
finally:
reset_current_write_origin(token)
result = json.loads(raw)
assert result["staged"] is True and result["proposal_staged"] is True
review_messages = [
{"role": "assistant", "tool_calls": [{"id": "c1", "function": {"name": "memory", "arguments": json.dumps(
{"action": "replace", "old_text": "standing rule", "content": "consolidated entry"})}}]},
{"role": "tool", "tool_call_id": "c1", "content": raw},
]
actions = bg.summarize_background_review_actions(review_messages, [])
assert any("staged for your approval" in a for a in actions)
def test_near_limit_denial_end_to_end(self, tmp_path, monkeypatch):
"""add rejected by the budget -> fork follows the 'consolidate now' hint with a
replace -> the delete gate stages it -> the proposal surfaces; the store never
changed and nothing was silently lost."""
import json
from tools.memory_tool import memory_tool
from tools.skill_provenance import set_current_write_origin, reset_current_write_origin
store = self._store(tmp_path, monkeypatch)
assert store.add("memory", "seed entry one")["success"] is True
# Near-limit: a further add is rejected and the store's hint says to consolidate.
assert store.add("memory", "x" * 600)["success"] is False
token = set_current_write_origin("background_review")
try:
add_raw = memory_tool(action="add", content="y" * 600, store=store)
replace_raw = memory_tool(
action="replace", old_text="seed entry one", content="merged entry", store=store)
finally:
reset_current_write_origin(token)
assert json.loads(add_raw)["success"] is False # the budget still rejects the add
replace_result = json.loads(replace_raw)
assert replace_result["staged"] is True and replace_result["proposal_staged"] is True
# Fail-closed: nothing was applied or dropped.
assert "seed entry one" in store._entries_for("memory")
assert "merged entry" not in store._entries_for("memory")
review_messages = [
{"role": "assistant", "tool_calls": [
{"id": "c1", "function": {"name": "memory", "arguments": json.dumps(
{"action": "add", "content": "y" * 600})}},
{"id": "c2", "function": {"name": "memory", "arguments": json.dumps(
{"action": "replace", "old_text": "seed entry one", "content": "merged entry"})}},
]},
{"role": "tool", "tool_call_id": "c1", "content": add_raw},
{"role": "tool", "tool_call_id": "c2", "content": replace_raw},
]
actions = bg.summarize_background_review_actions(review_messages, [])
assert any("staged for your approval" in a for a in actions)

View File

@@ -75,6 +75,23 @@ def test_cli_refine_snapshot_does_not_alias_live_history(monkeypatch):
_assert_isolated(cli.conversation_history, snapshot)
def test_cli_refine_reaches_spawn_as_explicit(monkeypatch):
"""The explicit /refine path must reach the spawn with ``explicit=True`` so the fork runs
under the refine_review origin and keeps the full memory operation set (#105921 review)."""
from hermes_cli.cli_commands_mixin import CLICommandsMixin
monkeypatch.setattr("cli._cprint", lambda *a, **k: None, raising=False)
agent = _agent_with_real_chokepoint()
cli = object.__new__(CLICommandsMixin)
cli.agent = agent
cli.conversation_history = _nested_history()
cli._handle_refine_command("/refine")
agent._spawn_background_review_now.assert_called_once()
assert agent._spawn_background_review_now.call_args.kwargs["explicit"] is True
@pytest.mark.asyncio
async def test_gateway_refine_snapshot_does_not_alias_live_history():
from gateway.run import GatewayRunner

View File

@@ -765,21 +765,33 @@ class TestBatchRefusesToEmptyNonEmptyStore:
class TestBackgroundReviewDeleteGate:
"""An unattended background-review fork may append, never delete: the near-limit
'consolidate now' hint is otherwise an instruction to decide what to forget,
executed with no human in the loop."""
executed with no human in the loop. Denied ops are staged as pending proposals
(surfaced via /memory pending) instead of silently dropped — the fork's own review
summary is never published back."""
def test_remove_denied_and_store_untouched(self, store):
def test_remove_staged_not_applied(self, store, tmp_path, monkeypatch):
monkeypatch.setenv("HERMES_HOME", str(tmp_path))
store.add("memory", "never create records without permission")
token = set_current_write_origin("background_review")
try:
result = json.loads(memory_tool(action="remove", old_text="without permission", store=store))
finally:
reset_current_write_origin(token)
assert result["success"] is False
assert "Background review may not delete" in result["error"]
assert result["success"] is True
assert result["staged"] is True
assert result["proposal_staged"] is True
assert result["pending_id"]
assert "staged for your approval" in result["message"]
# Fail-closed: the standing rule is still on disk.
assert "never create records without permission" in store._entries_for("memory")
# The proposal itself landed in the pending store for the user to approve or discard.
from tools.write_approval import MEMORY, get_pending
record = get_pending(MEMORY, result["pending_id"])
assert record["payload"]["action"] == "remove"
assert record["origin"] == "background_review"
def test_replace_denied_in_background_review(self, store):
def test_replace_staged_in_background_review(self, store, tmp_path, monkeypatch):
monkeypatch.setenv("HERMES_HOME", str(tmp_path))
store.add("memory", "entry the fork must not rewrite")
token = set_current_write_origin("background_review")
try:
@@ -787,10 +799,13 @@ class TestBackgroundReviewDeleteGate:
action="replace", old_text="entry the fork", content="rewritten by fork", store=store))
finally:
reset_current_write_origin(token)
assert result["success"] is False
assert "Background review may not delete" in result["error"]
assert result["staged"] is True
assert result["proposal_staged"] is True
# Fail-closed: the original entry is untouched.
assert "entry the fork must not rewrite" in store._entries_for("memory")
def test_batch_containing_remove_denied(self, store):
def test_batch_containing_remove_staged_whole_batch(self, store, tmp_path, monkeypatch):
monkeypatch.setenv("HERMES_HOME", str(tmp_path))
store.add("memory", "rule one")
token = set_current_write_origin("background_review")
try:
@@ -800,8 +815,8 @@ class TestBackgroundReviewDeleteGate:
], store=store))
finally:
reset_current_write_origin(token)
assert result["success"] is False
# Atomic denial: the batch's add must not land either.
assert result["staged"] is True
# Atomic: the batch is only a proposal — its add must not land either.
assert "fork consolidation" not in store._entries_for("memory")
def test_add_still_allowed_in_background_review(self, store):
@@ -818,3 +833,17 @@ class TestBackgroundReviewDeleteGate:
result = json.loads(memory_tool(action="remove", old_text="supervised turn", store=store))
assert result["success"] is True
assert "entry a supervised turn may remove" not in store._entries_for("memory")
def test_refine_review_origin_keeps_full_operation_set(self, store):
# A user-requested /refine fork runs under the refine_review origin: it is not an
# unattended review, so replace/remove keep working on that supervised surface.
store.add("memory", "entry an explicit refine may rewrite")
token = set_current_write_origin("refine_review")
try:
result = json.loads(memory_tool(
action="replace", old_text="entry an explicit", content="rewritten by refine", store=store))
finally:
reset_current_write_origin(token)
assert result["success"] is True
assert "rewritten by refine" in store._entries_for("memory")
assert "entry an explicit refine may rewrite" not in store._entries_for("memory")

View File

@@ -129,23 +129,44 @@ def _validate_single_op(store, action, target, content, old_text) -> Optional[st
_BG_DELETE_ACTIONS = ("replace", "remove")
def _background_delete_gate(action, operations) -> Optional[str]:
def _background_delete_gate(action, operations, target="memory", content=None, old_text=None) -> Optional[str]:
"""Fail-closed operation gate for unattended background-review forks (#105921): ``add``
stays available (it is all any review prompt asks for), while ``replace``/``remove`` —
single or inside a batch — are denied. The near-limit "consolidate now" hint is otherwise
an instruction to decide what to forget, executed with no human in the loop."""
single or inside a batch — are never applied unattended. The op is staged in the pending
store instead of merely denied: the fork's own review summary is never published back, so
a plain denial would drop the consolidation request with no surfacing path at all. A
staging failure fails closed to a plain denial."""
from tools.skill_provenance import is_background_review
if not is_background_review():
return None
if action in _BG_DELETE_ACTIONS or any(
isinstance(op, dict) and op.get("action") in _BG_DELETE_ACTIONS for op in (operations or [])
):
hit = action in _BG_DELETE_ACTIONS or any(
isinstance(op, dict) and op.get("action") in _BG_DELETE_ACTIONS for op in (operations or []))
if not hit:
return None
payload = ({"action": "batch", "target": target, "operations": operations}
if operations is not None else
{"action": action, "target": target, "content": content, "old_text": old_text})
detail = ("; ".join(_batch_op_line(op) for op in operations) if operations is not None
else _batch_op_line({"action": action, "content": content, "old_text": old_text}))
try:
from tools import write_approval as wa
record = wa.stage_write(
wa.MEMORY, payload,
summary=(f"background review consolidation ({'batch' if operations is not None else action} "
f"on {target}): {detail}")[:200],
origin=wa.current_origin())
return json.dumps({
"success": True, "staged": True, "proposal_staged": True, "pending_id": record["id"],
"message": ("Background review may not delete memory entries unattended. The proposed "
f"{'batch' if operations is not None else action} was staged for your approval — "
"review it with /memory pending (approve to apply, discard to drop)."),
}, ensure_ascii=False)
except Exception:
logger.warning("Failed to stage background-review consolidation; denying", exc_info=True)
return tool_error(
"Background review may not delete memory entries ('replace'/'remove', including in a "
"batch); 'add' is still available. Propose a consolidation in your review summary "
"instead.", success=False)
return None
"batch); 'add' is still available.", success=False)
def memory_tool(action: str = None, target: str = "memory", content: str = None, old_text: str = None,
@@ -163,7 +184,7 @@ def memory_tool(action: str = None, target: str = "memory", content: str = None,
target_error = _memory_target_error(store, target)
if target_error is not None:
return json.dumps(target_error)
denied = _background_delete_gate(action, operations)
denied = _background_delete_gate(action, operations, target, content, old_text)
if denied is not None:
return denied
if operations: