diff --git a/agent/background_review.py b/agent/background_review.py index 70c83dea14..85373db03a 100644 --- a/agent/background_review.py +++ b/agent/background_review.py @@ -1,15 +1,11 @@ """Background memory/skill review — fork the agent to evaluate the turn. -After every turn ``AIAgent.run_conversation`` may spawn a daemon thread that -replays the conversation snapshot in a forked :class:`AIAgent` and asks -"should any skill/memory be saved or updated?". Writes go straight to the -memory + skill stores; the main conversation and prompt cache are never touched. -The fork inherits the parent's live runtime (provider, model, credentials, -cached system prompt) so it hits the same prefix cache, and runs under a +After every turn ``AIAgent.run_conversation`` may spawn a daemon thread that replays the +conversation snapshot in a forked :class:`AIAgent` and asks "should any skill/memory be saved +or updated?". Writes go straight to the memory + skill stores; the main conversation and prompt +cache are never touched. The fork inherits the parent's live runtime (provider, model, +credentials, cached system prompt) so it hits the same prefix cache, and runs under a dispatch-side tool whitelist limited to memory/skill tools. - -See the ``hermes-agent-dev`` skill (``references/self-improvement-loop.md``) -for invariants and PR review criteria. """ from __future__ import annotations @@ -20,6 +16,7 @@ import logging import os import threading from contextlib import contextmanager +from dataclasses import dataclass, field from typing import Any, Dict, Iterator, List, Optional, Tuple from agent.thread_scoped_output import thread_scoped_silence @@ -96,10 +93,7 @@ def prepare_background_review_run(agent: Any) -> Optional[_BackgroundReviewRun]: return run -def finish_background_review_run( - agent: Any, - run: Optional[_BackgroundReviewRun], -) -> None: +def finish_background_review_run(agent: Any, run: Optional[_BackgroundReviewRun]) -> None: """Publish one run's request exit without clearing a successor (ABA-safe).""" if run is None or not run.mark_request_finished(): return @@ -110,38 +104,31 @@ def finish_background_review_run( def _interrupt_background_review(review_agent: Any) -> None: - """Request abort off-thread so a wedged abort hook cannot stall the live turn - (the bounded ``request_done`` wait in the canceller relies on this returning fast).""" + """Request abort off-thread so a wedged abort hook cannot stall the live turn (the bounded + ``request_done`` wait in the canceller relies on this returning fast).""" def _interrupt() -> None: try: from agent.interrupt_compat import request_hard_interrupt request_hard_interrupt( - review_agent, - "superseded by a new live turn", + review_agent, "superseded by a new live turn", tool_reason="background review superseded", ) except Exception: - logger.debug( - "Failed to cancel in-flight background review for a new turn", - exc_info=True, - ) + logger.debug("Failed to cancel in-flight background review for a new turn", exc_info=True) try: threading.Thread(target=_interrupt, daemon=True, name="bg-review-cancel").start() except Exception: - logger.debug( - "Failed to start background-review cancellation thread", - exc_info=True, - ) + logger.debug("Failed to start background-review cancellation thread", exc_info=True) def cancel_background_review_for_live_turn(agent: Any) -> None: """Cancel the current review and await its request-phase acknowledgement. - Foreground priority: past the bounded deadline, warn and let the live turn - proceed — self-improvement work must never block a user-facing turn. + Foreground priority: past the bounded deadline, warn and let the live turn proceed — + self-improvement work must never block a user-facing turn. """ with _optional_lock(agent, "_background_review_lock"): run = getattr(agent, "_background_review_run", None) @@ -164,19 +151,17 @@ def cancel_background_review_for_live_turn(agent: Any) -> None: ) -# Aux-model routing: by default ("auto") the fork runs on the MAIN model and -# replays the full conversation as warm cache reads. When -# auxiliary.background_review.{provider,model} routes it to a DIFFERENT model -# the cache is cold anyway, so the fork replays a compact digest instead. +# Aux-model routing: by default ("auto") the fork runs on the MAIN model and replays the full +# conversation as warm cache reads. When auxiliary.background_review.{provider,model} routes it +# to a DIFFERENT model the cache is cold anyway, so the fork replays a compact digest instead. _REVIEW_MAX_ITERATIONS = 16 # Aggregate INPUT-token budget for one review fork (checked in conversation_loop's -# ``_review_input_budget_exhausted``). Request #1 replays the full snapshot as a -# warm cache read (both compression gates deferred until the first response); -# compaction then bounds each request, but nothing else caps the SUM across the -# tool loop. 2x the historical 300k foreground trigger. Override via -# ``auxiliary.background_review.max_input_tokens``; <= 0 disables. +# ``_review_input_budget_exhausted``). Request #1 replays the full snapshot as a warm cache read +# (both compression gates deferred until the first response); compaction then bounds each +# request, but nothing else caps the SUM across the tool loop. 2x the historical 300k foreground +# trigger. Override via ``auxiliary.background_review.max_input_tokens``; <= 0 disables. _REVIEW_MAX_INPUT_TOKENS_DEFAULT = 600_000 @@ -187,14 +172,9 @@ def _task_block(cfg: Any) -> Dict[str, Any]: return task if isinstance(task, dict) else {} -def _background_review_task_config( - task_cfg: Optional[Dict[str, Any]] = None, -) -> Dict[str, Any]: - """Return ``auxiliary.background_review`` (or ``{}`` on any failure). - - Pass ``task_cfg`` when the caller already loaded the block so the spawn / - resolve / prompt paths do not re-read config on every turn. - """ +def _background_review_task_config(task_cfg: Optional[Dict[str, Any]] = None) -> Dict[str, Any]: + """``auxiliary.background_review`` (or ``{}`` on any failure); pass a pre-loaded ``task_cfg`` + so the spawn / resolve / prompt paths do not re-read config on every turn.""" if task_cfg is not None: return task_cfg if isinstance(task_cfg, dict) else {} try: @@ -205,13 +185,9 @@ def _background_review_task_config( return {} -def _review_input_token_budget( - task_cfg: Optional[Dict[str, Any]] = None, -) -> Optional[int]: +def _review_input_token_budget(task_cfg: Optional[Dict[str, Any]] = None) -> Optional[int]: """Aggregate input-token budget for one review fork (None = unlimited; <= 0 disables).""" - raw = _background_review_task_config(task_cfg).get( - "max_input_tokens", _REVIEW_MAX_INPUT_TOKENS_DEFAULT - ) + raw = _background_review_task_config(task_cfg).get("max_input_tokens", _REVIEW_MAX_INPUT_TOKENS_DEFAULT) try: budget = int(raw) except (TypeError, ValueError): @@ -220,11 +196,8 @@ def _review_input_token_budget( def load_background_review_settings() -> tuple[bool, Dict[str, Any]]: - """Single config read for the automatic-review gate + task block. - - Returns ``(enabled, task_cfg)``. Fail-open (``enabled=True``) so a broken - config never silently disables reviews — but WARN so the cost is visible. - """ + """Single config read -> ``(enabled, task_cfg)``. Fail-open (``enabled=True``) so a broken + config never silently disables reviews — but WARN so the cost is visible.""" try: from hermes_cli.config import load_config_readonly from utils import is_truthy_value @@ -240,15 +213,10 @@ def load_background_review_settings() -> tuple[bool, Dict[str, Any]]: return True, {} -def is_background_review_enabled( - task_cfg: Optional[Dict[str, Any]] = None, -) -> bool: - """Whether automatic post-turn review may spawn (``enabled``, default true). - - Explicit ``/refine`` (``focus`` set) bypasses this gate. Prefer - :func:`load_background_review_settings` at the spawn site so the block is - not re-read on the same turn. - """ +def is_background_review_enabled(task_cfg: Optional[Dict[str, Any]] = None) -> bool: + """Whether automatic post-turn review may spawn (``enabled``, default true). Explicit + ``/refine`` (``focus`` set) bypasses this gate; prefer :func:`load_background_review_settings` + at the spawn site so the block is not re-read on the same turn.""" if task_cfg is None: return load_background_review_settings()[0] try: @@ -264,16 +232,13 @@ def is_background_review_enabled( return True -def _resolve_review_runtime( - agent: Any, - task_cfg: Optional[Dict[str, Any]] = None, -) -> Dict[str, Any]: +def _resolve_review_runtime(agent: Any, task_cfg: Optional[Dict[str, Any]] = None) -> Dict[str, Any]: """Resolve provider/model/credentials for the review fork. - Default (auto / unset / same as parent): the parent's live runtime with - ``routed=False`` (codex_app_server -> codex_responses downgrade applied). - When ``auxiliary.background_review.{provider,model}`` names a different - concrete model, resolve that runtime and set ``routed=True``. + Default (auto / unset / same as parent): the parent's live runtime with ``routed=False`` + (codex_app_server -> codex_responses downgrade applied). When + ``auxiliary.background_review.{provider,model}`` names a different concrete model, resolve + that runtime and set ``routed=True``. """ parent_runtime = agent._current_main_runtime() parent_api_mode = parent_runtime.get("api_mode") or None @@ -294,8 +259,7 @@ def _resolve_review_runtime( } task = _background_review_task_config(task_cfg) task_provider, task_model, task_base_url, task_api_key = ( - str(task.get(key, "")).strip() or None - for key in ("provider", "model", "base_url", "api_key") + str(task.get(key, "")).strip() or None for key in ("provider", "model", "base_url", "api_key") ) if not (task_provider and task_provider != "auto" and task_model): return parent @@ -304,10 +268,8 @@ def _resolve_review_runtime( try: from hermes_cli.runtime_provider import resolve_runtime_provider rp = resolve_runtime_provider( - requested=task_provider, - target_model=task_model, - explicit_api_key=task_api_key, - explicit_base_url=task_base_url, + requested=task_provider, target_model=task_model, + explicit_api_key=task_api_key, explicit_base_url=task_base_url, ) return { "provider": rp.get("provider") or task_provider, @@ -328,13 +290,9 @@ def _resolve_review_runtime( def _parent_can_emit_tool_calls(agent: Any) -> bool: - """Whether a fork inheriting ``agent``'s runtime could act at all. - - An agent-as-provider client shim that cannot carry Hermes tool calls back - declares ``SUPPORTS_HERMES_TOOL_CALLS = False`` (instance or class) and is - skipped — the fork would be a guaranteed no-op that still pays a full - spawn. Anything that doesn't say otherwise is assumed capable. - """ + """Whether a fork inheriting ``agent``'s runtime could act at all: an agent-as-provider client + shim declaring ``SUPPORTS_HERMES_TOOL_CALLS = False`` (instance or class) is skipped — the fork + would be a guaranteed no-op that still pays a full spawn. Silence means capable.""" client = getattr(agent, "client", None) for candidate in (client, type(client) if client is not None else None): supported = getattr(candidate, "SUPPORTS_HERMES_TOOL_CALLS", None) @@ -353,13 +311,9 @@ def _msg_text(m: Dict) -> str: def _digest_history(messages_snapshot: List[Dict], tail: int = 24) -> List[Dict]: - """Compact replay for the routed (different-model) path only. - - Keeps the recent ``tail`` messages verbatim (extended so the kept run never - starts on a tool result) and collapses older turns into one synthetic - user-role digest, preserving role alternation. Never used on the - main-model path, where the full replay stays warm. - """ + """Compact replay for the routed (different-model) path only: keep the recent ``tail`` + messages verbatim (extended so the kept run never starts on a tool result) and collapse older + turns into one synthetic user-role digest, preserving role alternation.""" msgs = list(messages_snapshot or []) if len(msgs) <= tail: return msgs @@ -380,9 +334,7 @@ def _digest_history(messages_snapshot: List[Dict], tail: int = 24) -> List[Dict] elif role == "assistant": if m.get("tool_calls"): names = [ - (tc.get("function") or {}).get("name", "?") - for tc in m["tool_calls"] - if isinstance(tc, dict) + (tc.get("function") or {}).get("name", "?") for tc in m["tool_calls"] if isinstance(tc, dict) ] lines.append(f"ASSISTANT[tools: {', '.join(names)}]") if text: @@ -663,16 +615,13 @@ _COMBINED_REVIEW_PROMPT = ( ) - def _preview(text: str, limit: int) -> str: return text[:limit] + ("…" if len(text) > limit else "") # Memory op -> (glyph, which field carries the preview, preview length). _MEMORY_OP_FORMATS: Dict[str, Tuple[str, str, int]] = { - "add": ("➕", "content", 120), - "replace": ("✏️", "content", 120), - "remove": ("➖", "old_text", 60), + "add": ("➕", "content", 120), "replace": ("✏️", "content", 120), "remove": ("➖", "old_text", 60) } @@ -681,8 +630,8 @@ def _memory_op_line(label: str, action: str, fields: Dict[str, str]) -> Optional fmt = _MEMORY_OP_FORMATS.get(action) if fmt is None: return None - glyph, field, limit = fmt - text = fields.get(field) or "" + glyph, field_name, limit = fmt + text = fields.get(field_name) or "" return f"{label} {glyph} {_preview(text, limit)}" if text else None @@ -711,11 +660,7 @@ def _verbose_memory_lines(label: str, detail: Dict) -> List[str]: ops_raw = detail.get("operations") operations: list = ops_raw if isinstance(ops_raw, list) else [] if operations: - lines = [ - _memory_op_line(label, op.get("action", ""), op) - for op in operations - if isinstance(op, dict) - ] + lines = [_memory_op_line(label, op.get("action", ""), op) for op in operations if isinstance(op, dict)] return [line for line in lines if line] line = _memory_op_line(label, detail.get("action", ""), detail) return [line or f"{label} updated"] @@ -724,9 +669,9 @@ def _verbose_memory_lines(label: str, detail: Dict) -> List[str]: def _collect_review_call_details(review_messages: List[Dict]) -> Tuple[set, dict]: """Map review-agent tool_call ids -> parsed call arguments for notify tools. - Result JSON only says "Entry added"; the call arguments carry action, - target and content previews. Restricting to notify tools keeps helper - tools from surfacing as memory work just because they succeeded. + Result JSON only says "Entry added"; the call arguments carry action, target and content + previews. Restricting to notify tools keeps helper tools from surfacing as memory work just + because they succeeded. """ notify_tools = {"memory", "skill_manage"} all_tool_call_ids: set = set() @@ -763,37 +708,61 @@ def _collect_review_call_details(review_messages: List[Dict]) -> Tuple[set, dict return all_tool_call_ids, call_details -def summarize_background_review_actions( - review_messages: List[Dict], - prior_snapshot: List[Dict], - notification_mode: str = "on", -) -> List[str]: - """Build the human-facing action summary for a background review pass. - - Collects successful memory / skill-management tool results from the review - agent's messages, skipping tool messages already present in - ``prior_snapshot`` so inherited results are not re-surfaced as fresh work. - - ``notification_mode``: ``off`` -> no actions; ``on`` -> generic - "Memory updated"/tool messages; ``verbose`` -> content previews from the - tool-call arguments. - """ - mode = str(notification_mode or "on").lower() - if mode == "off": - return [] - verbose = mode == "verbose" - - existing_tool_call_ids = set() - existing_tool_contents = set() +def _prior_tool_keys(prior_snapshot: List[Dict]) -> Tuple[set, set]: + """``(tool_call_ids, contents)`` of tool messages already in the parent snapshot.""" + ids: set = set() + contents: set = set() for prior in prior_snapshot or []: if not isinstance(prior, dict) or prior.get("role") != "tool": continue tcid = prior.get("tool_call_id") if tcid: - existing_tool_call_ids.add(tcid) + ids.add(tcid) elif isinstance(prior.get("content"), str): - existing_tool_contents.add(prior["content"]) + contents.add(prior["content"]) + return ids, contents + +def _action_lines(data: Dict, detail: Dict, verbose: bool) -> List[str]: + """Summary line(s) for one successful notify-tool result (``[]`` when nothing to report).""" + message = data.get("message", "") + target = data.get("target", "") or detail.get("target", "") + is_skill = detail.get("tool") == "skill_manage" + message_lower = message.lower() + if not verbose and ( + "created" in message_lower or "updated" in message_lower or (is_skill and "patched" in message_lower) + ): + return [message] + if is_skill: + label = "Skill" + elif target: + label = "Memory" if target == "memory" else "User profile" if target == "user" else target + else: + return [] + if verbose: + return [_verbose_skill_line(data, detail, message)] if is_skill else _verbose_memory_lines(label, detail) + if any(k in message_lower for k in ("added", "replaced", "removed", "applied")) or ( + target and "add" in message_lower + ): + return [f"{label} updated"] + return [] + + +def summarize_background_review_actions( + review_messages: List[Dict], prior_snapshot: List[Dict], notification_mode: str = "on" +) -> List[str]: + """Build the human-facing action summary for a background review pass. + + Collects successful memory / skill-management tool results from the review agent's messages, + skipping tool messages already present in ``prior_snapshot`` so inherited results are not + re-surfaced as fresh work. ``notification_mode``: ``off`` -> no actions; ``on`` -> generic + "Memory updated"/tool messages; ``verbose`` -> content previews from the tool-call arguments. + """ + mode = str(notification_mode or "on").lower() + if mode == "off": + return [] + verbose = mode == "verbose" + existing_tool_call_ids, existing_tool_contents = _prior_tool_keys(prior_snapshot) all_tool_call_ids, call_details = _collect_review_call_details(review_messages) actions: List[str] = [] @@ -813,65 +782,22 @@ def summarize_background_review_actions( data = json.loads(msg.get("content", "{}")) except (json.JSONDecodeError, TypeError): continue - # Wrapper MCP servers may return a top-level list/scalar; only dict - # payloads carry ``success``/``_change``. + # Wrapper MCP servers may return a top-level list/scalar; only dict payloads carry + # ``success``/``_change``. if not isinstance(data, dict) or not data.get("success"): continue - message = data.get("message", "") - detail = call_details.get(tcid) or {} - if not isinstance(detail, dict): - detail = {} - target = data.get("target", "") or detail.get("target", "") - is_skill = detail.get("tool") == "skill_manage" - - message_lower = message.lower() - if not verbose and ( - "created" in message_lower - or "updated" in message_lower - or (is_skill and "patched" in message_lower) - ): - actions.append(message) - continue - - if is_skill: - label = "Skill" - elif target: - label = "Memory" if target == "memory" else "User profile" if target == "user" else target - else: - continue - - if verbose: - if is_skill: - actions.append(_verbose_skill_line(data, detail, message)) - else: - actions.extend(_verbose_memory_lines(label, detail)) - elif ( - "added" in message_lower - or "replaced" in message_lower - or "removed" in message_lower - or "applied" in message_lower - or (target and "add" in message.lower()) - or "Entry added" in message - ): - actions.append(f"{label} updated") + actions.extend(_action_lines(data, call_details.get(tcid) or {}, verbose)) return actions def build_memory_write_metadata( - agent: Any, - *, - write_origin: Optional[str] = None, - execution_context: Optional[str] = None, - task_id: Optional[str] = None, - tool_call_id: Optional[str] = None, + agent: Any, *, write_origin: Optional[str] = None, execution_context: Optional[str] = None, + task_id: Optional[str] = None, tool_call_id: Optional[str] = None, ) -> Dict[str, Any]: """Build provenance metadata for external memory-provider mirrors.""" metadata: Dict[str, Any] = { "write_origin": write_origin or getattr(agent, "_memory_write_origin", "assistant_tool"), - "execution_context": ( - execution_context - or getattr(agent, "_memory_write_context", "foreground") - ), + "execution_context": execution_context or getattr(agent, "_memory_write_context", "foreground"), "session_id": agent.session_id or "", "parent_session_id": agent._parent_session_id or "", "platform": agent.platform or os.environ.get("HERMES_SESSION_SOURCE", "cli"), @@ -885,36 +811,25 @@ def build_memory_write_metadata( _USAGE_COUNTERS = ( - "input_tokens", - "output_tokens", - "cache_read_tokens", - "cache_write_tokens", - "reasoning_tokens", - "api_calls", + "input_tokens", "output_tokens", "cache_read_tokens", "cache_write_tokens", "reasoning_tokens", "api_calls", ) def _snapshot_review_usage(review_agent: Any) -> Dict[str, Any]: """Snapshot in-memory usage counters from a review fork (pre-close).""" - usage: Dict[str, Any] = { - key: getattr(review_agent, key, None) - for key in ("model", "provider", "base_url") - } + usage: Dict[str, Any] = {key: getattr(review_agent, key, None) for key in ("model", "provider", "base_url")} for key in _USAGE_COUNTERS: usage[key] = int(getattr(review_agent, f"session_{key}", 0) or 0) usage["estimated_cost_usd"] = getattr(review_agent, "session_estimated_cost_usd", None) return usage -def _record_review_usage_to_parent( - parent_agent: Any, - usage: Dict[str, Any], -) -> None: +def _record_review_usage_to_parent(parent_agent: Any, usage: Dict[str, Any]) -> None: """Record a fork's usage against the parent session (best-effort, never raises). - The fork has ``_session_db = None`` so conversation_loop's DB-gated accounting - never sees its calls; route them through the aux-accounting chokepoint, which - writes only ``session_model_usage`` — never the transcript or ``sessions`` row. + The fork has ``_session_db = None`` so conversation_loop's DB-gated accounting never sees its + calls; route them through the aux-accounting chokepoint, which writes only + ``session_model_usage`` — never the transcript or ``sessions`` row. """ try: session_db = getattr(parent_agent, "_session_db", None) @@ -925,27 +840,21 @@ def _record_review_usage_to_parent( if not any(counts.values()): return # fork made no successful API calls (e.g. failed at spawn) session_db.record_auxiliary_usage( - session_id, - task="background_review", - model=usage.get("model"), - billing_provider=usage.get("provider"), - billing_base_url=usage.get("base_url"), + session_id, task="background_review", model=usage.get("model"), + billing_provider=usage.get("provider"), billing_base_url=usage.get("base_url"), estimated_cost_usd=usage.get("estimated_cost_usd"), - api_call_count=counts.pop("api_calls"), - **counts, + api_call_count=counts.pop("api_calls"), **counts, ) except Exception as e: - logger.debug( - "Background review usage recording failed (non-fatal): %s", e - ) + logger.debug("Background review usage recording failed (non-fatal): %s", e) def _classify_review_result(actions: List[str]) -> str: """Map a review action summary to ``none`` / ``skill`` / ``memory`` / ``skill+memory``. - Prefix-based on the formats :func:`summarize_background_review_actions` - emits (``Skill …``, ``📝 Skill …``, ``Memory …``, ``User profile …``), so a - free-text line like ``Skipped: no skill worth saving`` stays ``none``. + Prefix-based on the formats :func:`summarize_background_review_actions` emits (``Skill …``, + ``📝 Skill …``, ``Memory …``, ``User profile …``), so a free-text line like ``Skipped: no + skill worth saving`` stays ``none``. """ has_skill = has_memory = False for action in actions or []: @@ -957,9 +866,7 @@ def _classify_review_result(actions: List[str]) -> str: has_skill = True elif lower.startswith("memory") or lower.startswith("user profile"): has_memory = True - return "+".join( - kind for kind, hit in (("skill", has_skill), ("memory", has_memory)) if hit - ) or "none" + return "+".join(kind for kind, hit in (("skill", has_skill), ("memory", has_memory)) if hit) or "none" def _log_review_completion(usage: Dict[str, Any], result: str) -> None: @@ -975,37 +882,32 @@ def _log_review_completion(usage: Dict[str, Any], result: str) -> None: ) -# OpenRouter provider-routing pins: prompt caches live per UPSTREAM provider, -# so a fork without the parent's pins can land on a different upstream and -# miss the warm cache even with byte-identical prompt/tools bytes. +# OpenRouter provider-routing pins: prompt caches live per UPSTREAM provider, so a fork without +# the parent's pins can land on a different upstream and miss the warm cache even with +# byte-identical prompt/tools bytes. _PROVIDER_PIN_ATTRS = ( - "providers_allowed", - "providers_ignored", - "providers_order", - "provider_sort", - "provider_require_parameters", - "provider_data_collection", + "providers_allowed", "providers_ignored", "providers_order", "provider_sort", + "provider_require_parameters", "provider_data_collection", ) def _same_model_parity_kwargs(agent: Any) -> Dict[str, Any]: """AIAgent kwargs that keep a SAME-model fork's request bytes identical to the parent's. - Only for the un-routed path: on a different model the cache is cold - anyway, and the parent's reasoning-effort vocabulary may be invalid for - the routed provider (OpenRouter forwards ``reasoning.effort`` unclamped; - codex_responses passes ``max``/``ultra`` through unmapped). + Only for the un-routed path: on a different model the cache is cold anyway, and the parent's + reasoning-effort vocabulary may be invalid for the routed provider (OpenRouter forwards + ``reasoning.effort`` unclamped; codex_responses passes ``max``/``ultra`` through unmapped). """ kwargs: Dict[str, Any] = { # Anthropic's cache key is namespaced by ``thinking`` presence. "reasoning_config": getattr(agent, "reasoning_config", None), - # Gateway session context appended to the cached system prompt at - # API-call time; without it the effective system prompt diverges. + # Gateway session context appended to the cached system prompt at API-call time; + # without it the effective system prompt diverges. "ephemeral_system_prompt": getattr(agent, "ephemeral_system_prompt", None), } - # Prefill sits right after the system message, so a parent with prefill - # would diverge at index 1. Deep copy: unicode-error recovery sanitizes - # prefill entries IN PLACE and must not rewrite the parent's bytes. + # Prefill sits right after the system message, so a parent with prefill would diverge at + # index 1. Deep copy: unicode-error recovery sanitizes prefill entries IN PLACE and must not + # rewrite the parent's bytes. parent_prefill = copy.deepcopy(getattr(agent, "prefill_messages", None) or []) if parent_prefill: kwargs["prefill_messages"] = parent_prefill @@ -1019,25 +921,24 @@ def _same_model_parity_kwargs(agent: Any) -> Dict[str, Any]: def _detach_fork_compression(review_agent: Any) -> None: """Detached in-memory compaction for a fork sharing the parent's session_id. - Disabling compression (the old guard against compacting the parent's live - session) removed the only bound on the review's snapshot. Persistence is - already off, so compaction can only rewrite the fork's transcript — but the - compressor's own SessionDB/session_id binding must be severed too, or - cooldown/streak counters land on the parent's row. Force in-place mode and - re-enable compression ONLY after the rebind succeeded (fail-closed); gates - stay deferred until the first response so request #1 is a warm cache read. + Disabling compression (the old guard against compacting the parent's live session) removed + the only bound on the review's snapshot. Persistence is already off, so compaction can only + rewrite the fork's transcript — but the compressor's own SessionDB/session_id binding must be + severed too, or cooldown/streak counters land on the parent's row. Force in-place mode and + re-enable compression ONLY after the rebind succeeded (fail-closed); gates stay deferred + until the first response so request #1 is a warm cache read. """ bind = getattr(getattr(review_agent, "context_compressor", None), "bind_session_state", None) detached = False if callable(bind): try: - # Plugin/third-party context engines may reject these kwargs; they - # own their persistence policy, so a failed rebind never aborts the review. + # Plugin/third-party context engines may reject these kwargs; they own their + # persistence policy, so a failed rebind never aborts the review. bind(session_db=None, session_id="") detached = True except Exception: - # FAIL-CLOSED: the compressor may still point at the parent's - # SessionDB; enabling compression would re-open the sibling race. + # FAIL-CLOSED: the compressor may still point at the parent's SessionDB; enabling + # compression would re-open the sibling race. logger.warning( "background-review compressor detachment failed; " "keeping compression DISABLED on this review fork " @@ -1050,102 +951,98 @@ def _detach_fork_compression(review_agent: Any) -> None: review_agent._review_defer_compaction_before_first_response = True +def _fork_init_kwargs(agent: Any, rt: Dict[str, Any], routed: bool, max_iterations: int) -> Dict[str, Any]: + """AIAgent constructor kwargs for the review fork. + + skip_memory=True: an external memory plugin scoped to the parent's session_id would leak the + harness prompt into the user's real memory namespace; built-in MEMORY.md/USER.md state is + re-bound by the caller. Toolsets match the parent so ``tools[]`` is byte-identical + (Anthropic's cache key includes it); the runtime whitelist restricts dispatch. + """ + kwargs: Dict[str, Any] = { + "model": rt.get("model") or agent.model, + "max_iterations": max_iterations, + "quiet_mode": True, + "platform": agent.platform, + "provider": rt.get("provider") or agent.provider, + "api_mode": rt.get("api_mode"), + "base_url": rt.get("base_url") or None, + "api_key": rt.get("api_key") or None, + "credential_pool": rt.get("credential_pool"), + "request_overrides": rt.get("request_overrides") or {}, + "parent_session_id": agent.session_id, + "enabled_toolsets": getattr(agent, "enabled_toolsets", None), + "disabled_toolsets": getattr(agent, "disabled_toolsets", None), + "skip_memory": True, + } + if isinstance(rt.get("max_tokens"), int): + kwargs["max_tokens"] = rt["max_tokens"] + if isinstance(rt.get("command"), str) and rt["command"]: + kwargs["acp_command"] = rt["command"] + kwargs["acp_args"] = rt.get("args") or [] + if not routed: + kwargs.update(_same_model_parity_kwargs(agent)) + return kwargs + + def build_cache_parity_fork( - agent: Any, - task_cfg: Optional[Dict[str, Any]] = None, - *, - max_iterations: int, + agent: Any, task_cfg: Optional[Dict[str, Any]] = None, *, max_iterations: int, write_origin: str = "background_review", ) -> Tuple[Any, Dict[str, Any], bool]: """Construct a detached AIAgent fork with warm prompt-cache parity (shared with ``/btw``). - Same runtime/credentials as the parent, byte-identical system prompt / - tools[] / reasoning config on the same-model path, shared session_id for - prefix warmth, full persistence detachment (no state.db writes, rotation, - or external memory providers; in-place-only compaction). - Returns ``(fork_agent, runtime_dict, routed)``; ``routed`` means a different - model (cache cold — replay a digest). The caller owns registration, - whitelisting, running, usage attribution and teardown. + Same runtime/credentials as the parent, byte-identical system prompt / tools[] / reasoning + config on the same-model path, shared session_id for prefix warmth, full persistence + detachment (no state.db writes, rotation, or external memory providers; in-place-only + compaction). Returns ``(fork_agent, runtime_dict, routed)``; ``routed`` means a different + model (cache cold — replay a digest). The caller owns registration, whitelisting, running, + usage attribution and teardown. """ from run_agent import AIAgent # local: avoids a circular import at load - # Inherit the parent's live runtime: AIAgent.__init__'s env auto-resolution - # fails for OAuth-only providers, session-scoped creds and credential pools. + # Inherit the parent's live runtime: AIAgent.__init__'s env auto-resolution fails for + # OAuth-only providers, session-scoped creds and credential pools. _rt = _resolve_review_runtime(agent, task_cfg) _routed = bool(_rt.get("routed")) - _fork_kwargs: Dict[str, Any] = {} - if isinstance(_rt.get("max_tokens"), int): - _fork_kwargs["max_tokens"] = _rt["max_tokens"] - if isinstance(_rt.get("command"), str) and _rt["command"]: - _fork_kwargs["acp_command"] = _rt["command"] - _fork_kwargs["acp_args"] = _rt.get("args") or [] - if not _routed: - _fork_kwargs.update(_same_model_parity_kwargs(agent)) - # skip_memory=True: an external memory plugin scoped to the parent's - # session_id would leak the harness prompt into the user's real memory - # namespace; built-in MEMORY.md/USER.md state is re-bound below. Toolsets - # match the parent so ``tools[]`` is byte-identical (Anthropic's cache key - # includes it); the runtime whitelist restricts dispatch. - review_agent = AIAgent( - model=_rt.get("model") or agent.model, - max_iterations=max_iterations, - quiet_mode=True, - platform=agent.platform, - provider=_rt.get("provider") or agent.provider, - api_mode=_rt.get("api_mode"), - base_url=_rt.get("base_url") or None, - api_key=_rt.get("api_key") or None, - credential_pool=_rt.get("credential_pool"), - request_overrides=_rt.get("request_overrides") or {}, - parent_session_id=agent.session_id, - enabled_toolsets=getattr(agent, "enabled_toolsets", None), - disabled_toolsets=getattr(agent, "disabled_toolsets", None), - skip_memory=True, - **_fork_kwargs, - ) + review_agent = AIAgent(**_fork_init_kwargs(agent, _rt, _routed, max_iterations)) review_agent._memory_write_origin = write_origin review_agent._memory_write_context = write_origin - # The between-turns MCP refresh would add late-connecting MCP tools and - # break tools[] parity, so opt out. + # The between-turns MCP refresh would add late-connecting MCP tools and break tools[] parity. review_agent._skip_mcp_refresh = True review_agent._memory_store = agent._memory_store review_agent._memory_enabled = agent._memory_enabled review_agent._user_profile_enabled = agent._user_profile_enabled review_agent._memory_nudge_interval = 0 review_agent._skill_nudge_interval = 0 - # PERSISTENCE ISOLATION (curator-takeover root cause): sharing the parent's - # session_id, the fork would otherwise write its harness turn into the REAL - # session, which the next live turn re-reads as a standing instruction. + # PERSISTENCE ISOLATION (curator-takeover root cause): sharing the parent's session_id, the + # fork would otherwise write its harness turn into the REAL session, which the next live turn + # re-reads as a standing instruction. review_agent._persist_disabled = True review_agent._session_db = None review_agent._session_json_enabled = False - # Fork status/warning emits go via _print_fn/status_callback, which bypass - # the stdout redirect — suppress them. + # Fork status/warning emits go via _print_fn/status_callback, which bypass the stdout redirect. review_agent.suppress_status_output = True - # Same model only: share the warm cached system prompt (~26% cost cut; a - # rebuilt prompt misses the byte-exact prefix key) and pin session_start so - # any re-render (compression, plugin hooks) stays byte-identical. + # Same model only: share the warm cached system prompt (~26% cost cut; a rebuilt prompt misses + # the byte-exact prefix key) and pin session_start so any re-render (compression, plugin + # hooks) stays byte-identical. if not _routed: review_agent._cached_system_prompt = agent._cached_system_prompt review_agent.session_start = agent.session_start review_agent.session_id = agent.session_id - # Single-lifecycle fork sharing the live session_id: close() must not - # finalize the parent's still-active session row. + # Single-lifecycle fork sharing the live session_id: close() must not finalize the parent's + # still-active session row. review_agent._end_session_on_close = False _detach_fork_compression(review_agent) - # Compaction bounds a single request; this bounds the WHOLE review - # (checked in conversation_loop via _review_input_budget_exhausted). + # Compaction bounds a single request; this bounds the WHOLE review (checked in + # conversation_loop via _review_input_budget_exhausted). review_agent._review_input_token_budget = _review_input_token_budget(task_cfg) return review_agent, _rt, _routed def _bg_review_auto_deny(command, description, **kwargs): - """Non-interactive approval: dangerous-command guards resolve to "deny" - instead of input(), which would deadlock against the parent's TUI.""" - logger.warning( - "Background review auto-denied dangerous command: %s (%s)", - command, description, - ) + """Non-interactive approval: dangerous-command guards resolve to "deny" instead of input(), + which would deadlock against the parent's TUI.""" + logger.warning("Background review auto-denied dangerous command: %s (%s)", command, description) return "deny" @@ -1160,10 +1057,10 @@ def _set_thread_approval_callback(callback: Any) -> None: def _track_review_fork(agent: Any, review_agent: Any, *, register: bool) -> None: """Add (``register=True``) or remove the fork on the PARENT's tracking slots: - ``_background_review_agent`` (direct pointer the next live turn interrupts) - and ``_active_children`` (interrupt() fan-out). Removal is identity-scoped - and idempotent; both are best-effort for direct test stubs — the prepared - run token is the live-turn cancellation authority.""" + ``_background_review_agent`` (direct pointer the next live turn interrupts) and + ``_active_children`` (interrupt() fan-out). Removal is identity-scoped and idempotent; both + are best-effort for direct test stubs — the prepared run token is the live-turn cancellation + authority.""" if review_agent is None: return if hasattr(agent, "_background_review_agent"): @@ -1184,72 +1081,159 @@ def _track_review_fork(agent: Any, review_agent: Any, *, register: bool) -> None raise -def _review_tool_whitelist( - review_agent: Any, task_cfg: Optional[Dict[str, Any]] -) -> Tuple[set, set]: - """Return ``(whitelist, configured_extra_tools)`` for the review fork. - - DISPATCH-side only: the advertised ``tools[]`` stays byte-identical to the - parent's, so prompt-cache parity is untouched. - """ +def _review_tool_whitelist(review_agent: Any, task_cfg: Optional[Dict[str, Any]]) -> Tuple[set, set]: + """``(whitelist, configured_extra_tools)`` for the review fork — DISPATCH-side only, so the + advertised ``tools[]`` stays byte-identical to the parent's (prompt-cache parity).""" from model_tools import get_tool_definitions - # Gate the built-in memory tool on the profile's memory flags so a - # memory-disabled profile is never contaminated by the review LLM. + # Gate the built-in memory tool on the profile's memory flags so a memory-disabled profile + # is never contaminated by the review LLM. review_toolsets = ["skills"] if review_agent._memory_enabled or review_agent._user_profile_enabled: review_toolsets.insert(0, "memory") - whitelist = { - t["function"]["name"] - for t in get_tool_definitions(enabled_toolsets=review_toolsets, quiet_mode=True) - } - # Read-only file tools: denying read_file/search_files caused a per-review - # denial storm that starved the loop (read_file also registers the read - # with the read-before-write guard). Write tools stay denied — autonomous - # maintenance goes through skill_manage's validation. + whitelist = {t["function"]["name"] for t in get_tool_definitions(enabled_toolsets=review_toolsets, quiet_mode=True)} + # Read-only file tools: denying read_file/search_files caused a per-review denial storm that + # starved the loop (read_file also registers the read with the read-before-write guard). + # Write tools stay denied — autonomous maintenance goes through skill_manage's validation. whitelist |= {"read_file", "search_files"} - # ``extra_tools`` admits named parent tools (e.g. a human-gated proposal - # tool). The whitelist can only admit, never advertise: a listed tool must - # already exist in the inherited schema. + # ``extra_tools`` admits named parent tools (e.g. a human-gated proposal tool). The whitelist + # can only admit, never advertise: a listed tool must already exist in the inherited schema. configured_extra_tools: set = set() try: _extra_raw = _background_review_task_config(task_cfg).get("extra_tools", []) if isinstance(_extra_raw, list): - configured_extra_tools = { - name.strip() - for name in _extra_raw - if isinstance(name, str) and name.strip() - } + configured_extra_tools = {name.strip() for name in _extra_raw if isinstance(name, str) and name.strip()} whitelist |= configured_extra_tools except Exception: logger.debug("background_review extra_tools parse failed", exc_info=True) return whitelist, configured_extra_tools -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, +@dataclass +class _ReviewForkState: + """Mutable hand-off between the fork phase and the outer worker's error/cleanup paths.""" + + review_agent: Any = None + review_messages: List[Dict] = field(default_factory=list) + review_usage: Dict[str, Any] = field(default_factory=dict) + + +def _release_fork_clients(review_agent: Any) -> None: + """The fork shares the foreground session ID: close() / shutdown_memory_provider() are + session-bound (close() kills that session's terminal processes), so release only clients.""" + try: + review_agent.release_clients() + except Exception: + pass + + +def _finish_request_phase(agent: Any, review_agent: Any, review_run: Optional[_BackgroundReviewRun]) -> None: + """Unregister the fork and publish request completion (identity-scoped, idempotent).""" + _track_review_fork(agent, review_agent, register=False) + finish_background_review_run(agent, review_run) + + +def _run_review_fork( + agent: Any, messages_snapshot: List[Dict], prompt: str, task_cfg: Optional[Dict[str, Any]], + review_run: Optional[_BackgroundReviewRun], st: _ReviewForkState, ) -> None: - """Daemon-thread worker: build the fork, run the prompt, surface the action - summary via ``agent._safe_print`` / ``background_review_callback``. - ``review_run`` (from :func:`prepare_background_review_run`) cancelled before - the first provider call aborts without entering ``run_conversation()``.""" + """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) + _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) + extra_list = ", ".join(sorted(configured_extra_tools)) + set_thread_tool_whitelist( + review_whitelist, + deny_msg_fmt=( + "Background review denied non-whitelisted tool: " + "{tool_name}. Allowed here: skill_view/skills_list/" + "read_file/search_files to read, " + "skill_manage(action='patch'|...) to change skills, and " + "memory for notes." + + ( + " Configured extra tools also allowed: " + extra_list + "." + if configured_extra_tools + else "" + ) + + " Do not retry {tool_name}." + ), + ) + try: + from tools.skill_manager_tool import _reset_background_review_read_marks + + _reset_background_review_read_marks() + except Exception: + pass + + try: + if review_run is None or review_run.begin_request(st.review_agent): + # Routed -> digest (cache cold anyway); same model -> full snapshot (warm cache reads). + st.review_agent.run_conversation( + user_message=( + prompt + + "\n\nYou can only call memory and skill " + "management tools. Other tools will be denied " + "at runtime — do not attempt them." + + ( + " Exception — these configured tools are " + "also allowed: " + extra_list + "." + if configured_extra_tools + else "" + ) + ), + conversation_history=_digest_history(messages_snapshot) if _routed else messages_snapshot, + ) + finally: + clear_thread_tool_whitelist() + # Attribute usage to the PARENT session. Snapshot BEFORE unregister/close so counters + # survive teardown, and in this finally so a fork that consumed tokens then raised is + # still attributed. The recorder never raises. + if st.review_agent is not None: + st.review_usage.update(_snapshot_review_usage(st.review_agent)) + _record_review_usage_to_parent(agent, st.review_usage) + # Publish completion as soon as the provider-capable phase has returned or startup + # cancellation has fenced it out. + _finish_request_phase(agent, st.review_agent, review_run) + + st.review_messages = list(getattr(st.review_agent, "_session_messages", [])) + _release_fork_clients(st.review_agent) + st.review_agent = None + + +def _publish_review_summary(agent: Any, actions: List[str]) -> None: + summary = " · ".join(dict.fromkeys(actions)) + agent._safe_print(f" 💾 Self-improvement review: {summary}") + _bg_cb = agent.background_review_callback + if _bg_cb: + try: + _bg_cb(f"💾 Self-improvement review: {summary}") + except Exception: + pass + + +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, +) -> None: + """Daemon-thread worker: build the fork, run the prompt, surface the action summary via + ``agent._safe_print`` / ``background_review_callback``. ``review_run`` (from + :func:`prepare_background_review_run`) cancelled before the first provider call aborts + without entering ``run_conversation()``.""" if review_run is not None and review_run.cancel_requested.is_set(): finish_background_review_run(agent, review_run) return _set_thread_approval_callback(_bg_review_auto_deny) - # A client that can't carry Hermes tool calls back would spawn a fork that - # cannot write anything. Checked BEFORE the thread-scoped silence so the - # warning is not swallowed; cheap check first so the normal path never - # resolves the runtime twice. - if not _parent_can_emit_tool_calls(agent) and not bool( - _resolve_review_runtime(agent, task_cfg).get("routed") - ): + # A client that can't carry Hermes tool calls back would spawn a fork that cannot write + # anything. Checked BEFORE the thread-scoped silence so the warning is not swallowed; cheap + # check first so the normal path never resolves the runtime twice. + if not _parent_can_emit_tool_calls(agent) and not bool(_resolve_review_runtime(agent, task_cfg).get("routed")): logger.warning( "Background review skipped: provider %r cannot emit Hermes tool calls, " "so the review fork could not write memories or skills. Set " @@ -1260,112 +1244,18 @@ def _run_review_in_thread( _set_thread_approval_callback(None) return - review_agent = None - review_messages: List[Dict] = [] - review_usage: Dict[str, Any] = {} - - def _finish_request_phase(agent_ref) -> None: - _track_review_fork(agent, agent_ref, register=False) - finish_background_review_run(agent, review_run) - - def _release(agent_ref) -> None: - try: - agent_ref.release_clients() - except Exception: - pass - + st = _ReviewForkState() try: - # Silence stdout/stderr for THIS thread only: a process-global redirect - # would blank every other thread's console for the whole review. + # Silence stdout/stderr for THIS thread only: a process-global redirect would blank every + # other thread's console for the whole review. with thread_scoped_silence(): - review_agent, _rt, _routed = build_cache_parity_fork( - agent, task_cfg, max_iterations=_REVIEW_MAX_ITERATIONS - ) - _track_review_fork(agent, review_agent, register=True) + _run_review_fork(agent, messages_snapshot, prompt, task_cfg, review_run, st) - from hermes_cli.plugins import ( - set_thread_tool_whitelist, - clear_thread_tool_whitelist, - ) - - review_whitelist, configured_extra_tools = _review_tool_whitelist( - review_agent, task_cfg - ) - extra_list = ", ".join(sorted(configured_extra_tools)) - set_thread_tool_whitelist( - review_whitelist, - deny_msg_fmt=( - "Background review denied non-whitelisted tool: " - "{tool_name}. Allowed here: skill_view/skills_list/" - "read_file/search_files to read, " - "skill_manage(action='patch'|...) to change skills, and " - "memory for notes." - + ( - " Configured extra tools also allowed: " + extra_list + "." - if configured_extra_tools - else "" - ) - + " Do not retry {tool_name}." - ), - ) - try: - from tools.skill_manager_tool import _reset_background_review_read_marks - - _reset_background_review_read_marks() - except Exception: - pass - - try: - if review_run is None or review_run.begin_request(review_agent): - # Routed -> digest (cache cold anyway); same model -> full - # snapshot (warm cache reads). - review_agent.run_conversation( - user_message=( - prompt - + "\n\nYou can only call memory and skill " - "management tools. Other tools will be denied " - "at runtime — do not attempt them." - + ( - " Exception — these configured tools are " - "also allowed: " + extra_list + "." - if configured_extra_tools - else "" - ) - ), - conversation_history=( - _digest_history(messages_snapshot) if _routed - else messages_snapshot - ), - ) - finally: - clear_thread_tool_whitelist() - # Attribute usage to the PARENT session. Snapshot BEFORE - # unregister/close so counters survive teardown, and in this - # finally so a fork that consumed tokens then raised is still - # attributed. The recorder never raises. - if review_agent is not None: - review_usage.update(_snapshot_review_usage(review_agent)) - _record_review_usage_to_parent(agent, review_usage) - # Publish completion as soon as the provider-capable phase has - # returned or startup cancellation has fenced it out. - _finish_request_phase(review_agent) - - # Snapshot review actions before teardown. - review_messages = list(getattr(review_agent, "_session_messages", [])) - - # The fork shares the foreground session ID: close() / - # shutdown_memory_provider() are session-bound (close() kills that - # session's terminal processes), so release only this fork's clients. - _release(review_agent) - review_agent = None - - # 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. + # 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: actions = summarize_background_review_actions( - review_messages, - messages_snapshot, + st.review_messages, messages_snapshot, notification_mode=getattr(agent, "memory_notifications", "on"), ) except Exception as e: @@ -1377,33 +1267,24 @@ def _run_review_in_thread( ) actions = [] - _log_review_completion(review_usage, _classify_review_result(actions)) - + _log_review_completion(st.review_usage, _classify_review_result(actions)) if actions: - summary = " · ".join(dict.fromkeys(actions)) - agent._safe_print(f" 💾 Self-improvement review: {summary}") - _bg_cb = agent.background_review_callback - if _bg_cb: - try: - _bg_cb(f"💾 Self-improvement review: {summary}") - except Exception: - pass + _publish_review_summary(agent, actions) except Exception as e: logger.warning("Background memory/skill review failed: %s", e) - if review_usage: - _log_review_completion(review_usage, "error") + if st.review_usage: + _log_review_completion(st.review_usage, "error") agent._emit_auxiliary_failure("background review", e) finally: - # Safety net for the exception path (setup failures before the - # request-phase finally). Both cleanups are identity-scoped and - # idempotent; re-enter thread-scoped silence so cleanup output stays - # quiet without blanking other threads. - _finish_request_phase(review_agent) - if review_agent is not None: + # Safety net for the exception path (setup failures before the request-phase finally). + # Both cleanups are identity-scoped and idempotent; re-enter thread-scoped silence so + # cleanup output stays quiet without blanking other threads. + _finish_request_phase(agent, st.review_agent, review_run) + if st.review_agent is not None: try: with thread_scoped_silence(): - _release(review_agent) + _release_fork_clients(st.review_agent) except Exception: pass # Clear the approval callback so a recycled thread-id doesn't inherit it. @@ -1411,20 +1292,16 @@ def _run_review_in_thread( 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, + 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, ): - """Return ``(target, prompt)``; the caller builds the ``threading.Thread`` so - test patches of ``run_agent.threading.Thread`` keep working. + """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. + ``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. """ if task_cfg is None: task_cfg = _background_review_task_config() @@ -1445,14 +1322,8 @@ def spawn_background_review_thread( f"{focus}" ) - def _target() -> None: - _run_review_in_thread( - agent, - messages_snapshot, - prompt, - task_cfg=task_cfg, - review_run=review_run, - ) + 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) return _target, prompt diff --git a/tests/test_background_review_list_shapes.py b/tests/test_background_review_list_shapes.py index 79b003a2bb..ada4c151d9 100644 --- a/tests/test_background_review_list_shapes.py +++ b/tests/test_background_review_list_shapes.py @@ -286,8 +286,7 @@ def test_e_call_does_not_unwind_module_callables(): "summarize_background_review_actions returned partial results" in src ), "expected partial-results guard message present" - # And the prior-tonon-dict guard for the call_details lookup. - assert "if not isinstance(detail, dict):" in src + # And the non-dict guards on free-form tool payload fields. assert "if isinstance(ops_raw, list)" in src assert "if isinstance(change_raw, dict)" in src