From a5eb266bd5bd4b874aa79eea5bb19e14cf404ef8 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 18:27:34 -0700 Subject: [PATCH] refactor(agent/review_*,side_question): hoist review goal constant, single-lock pending stamp, compact docstrings --- agent/review_engine.py | 80 +++++++++++++++----------------- agent/review_idle_queue.py | 93 ++++++++++++++++---------------------- agent/side_question.py | 45 ++++++------------ 3 files changed, 90 insertions(+), 128 deletions(-) diff --git a/agent/review_engine.py b/agent/review_engine.py index d6b5635082..5e8430901a 100644 --- a/agent/review_engine.py +++ b/agent/review_engine.py @@ -1,18 +1,12 @@ """Shared engine for the /review command — every surface calls this. -/review spawns an independent, full-privilege background subagent (the same -async rail as ``delegate_task(background=true)``) to thoroughly review whatever -the recent conversation presented (PR, diff, code, docs). Its result re-enters -the spawning session as a normal async-delegation completion. - -Model routing: ``auxiliary.review`` (provider/model/base_url/api_key/api_mode) -when configured, else the parent agent's credentials (main-model-first). It is -passed as ``credentials_cfg`` to ``delegate_task`` so native-SDK providers, -api_mode detection and credential pools behave identically to -``delegation.provider`` pins. - -Surfaces (CLI/gateway ``/review``, TUI/Desktop) are thin adapters: snapshot the -conversation, call :func:`start_review`, print the dispatch note. +/review spawns an independent, full-privilege background subagent (the same async rail as +``delegate_task(background=true)``) to review whatever the recent conversation presented; +its result re-enters the spawning session as a normal async-delegation completion. +Model routing: ``auxiliary.review`` when configured, else the parent agent's credentials, +passed as ``credentials_cfg`` to ``delegate_task`` so native-SDK providers, api_mode +detection and credential pools behave identically to ``delegation.provider`` pins. +Surfaces (CLI/gateway ``/review``, TUI/Desktop) snapshot, call :func:`start_review`, print the note. """ from __future__ import annotations @@ -27,10 +21,23 @@ logger = logging.getLogger(__name__) # How many recent chat messages (user + assistant turns) the reviewer gets. DEFAULT_CONTEXT_MESSAGES = 10 -# Per-message excerpt cap: generous (a PR summary/diff excerpt is exactly what -# the reviewer needs) but bounded against a pathological turn. +# Per-message excerpt cap: generous (a PR summary/diff excerpt is exactly what the +# reviewer needs) but bounded against a pathological turn. _MESSAGE_CHAR_CAP = 12_000 +_REVIEW_GOAL = ( + "Act as an independent senior reviewer. Thoroughly review the work " + "presented in the conversation excerpt provided in your context: " + "investigate any code, pull request, branch, commit, documentation, " + "design, or other artifact it references (open the PR, read the " + "diff, run the code or tests where feasible) rather than judging " + "from the excerpt alone. Produce a full, structured review: what " + "the work does, whether it is correct and complete, concrete " + "defects or risks found (with file/line references where possible), " + "what was verified vs. only read, and a clear final verdict with " + "recommended next steps." +) + def _message_text(message: Dict[str, Any]) -> str: """Display text of a message; multimodal parts are joined, non-text parts noted.""" @@ -54,8 +61,7 @@ def snapshot_recent_messages( ) -> List[Dict[str, str]]: """Last ``limit`` user/assistant messages as {role, text} dicts, oldest first. - System messages, tool results and empty-text messages (pure tool-call - assistant stubs) are excluded. + System messages, tool results and empty-text messages (pure tool-call stubs) are excluded. """ out: List[Dict[str, str]] = [] for message in reversed(list(messages or [])): @@ -82,10 +88,9 @@ def collect_parent_loaded_skills( """Names of skills the parent agent was operating under. Launch-preloaded skills come from the stable marker in the parent's - ``ephemeral_system_prompt`` (``build_preloaded_skills_prompt``); mid-session - loads from ``skill_view`` tool calls in the history. Preloaded first, then - history loads, deduped, capped at ``limit`` (a reviewer told to load 30 - skills would burn its budget before working). + ``ephemeral_system_prompt``; mid-session loads from ``skill_view`` tool calls in the + history. Preloaded first, then history loads, deduped, capped at ``limit`` (a reviewer + told to load 30 skills would burn its budget before working). """ names: List[str] = [] prompt = str(getattr(parent_agent, "ephemeral_system_prompt", "") or "") @@ -101,8 +106,8 @@ def collect_parent_loaded_skills( args = json.loads(fn.get("arguments") or "{}") except Exception: continue - # Only whole-skill loads seed the reviewer; a reference-file read - # is a detail of the parent's task covered by loading the SKILL.md. + # Only whole-skill loads seed the reviewer; a reference-file read is a detail + # of the parent's task covered by loading the SKILL.md. if isinstance(args, dict) and not args.get("file_path"): candidates.append(str(args.get("name") or "")) for name in candidates: @@ -118,19 +123,6 @@ def build_review_task( loaded_skills: Optional[List[str]] = None, ) -> tuple: """Compose the reviewer subagent's (goal, context) pair.""" - goal = ( - "Act as an independent senior reviewer. Thoroughly review the work " - "presented in the conversation excerpt provided in your context: " - "investigate any code, pull request, branch, commit, documentation, " - "design, or other artifact it references (open the PR, read the " - "diff, run the code or tests where feasible) rather than judging " - "from the excerpt alone. Produce a full, structured review: what " - "the work does, whether it is correct and complete, concrete " - "defects or risks found (with file/line references where possible), " - "what was verified vs. only read, and a clear final verdict with " - "recommended next steps." - ) - lines = [ "You were spawned by the /review command. The following is an " "excerpt of the most recent conversation between the user and " @@ -161,14 +153,14 @@ def build_review_task( "the primary agent and its user. Be direct and specific; do not " "soften findings.", ] - return goal, "\n".join(lines) + return _REVIEW_GOAL, "\n".join(lines) def _load_review_credentials_cfg() -> Optional[Dict[str, Any]]: """Read ``auxiliary.review`` into a delegation-credentials-shaped dict. - None when nothing is configured (provider=auto/empty and no model/base_url): - the reviewer then inherits the parent agent's credentials. + None when nothing is configured (provider=auto/empty and no model/base_url): the + reviewer then inherits the parent agent's credentials. """ try: from hermes_cli.config import load_config_readonly @@ -194,10 +186,10 @@ def start_review( ) -> Dict[str, Any]: """Dispatch the reviewer subagent in the background. - Returns the parsed ``delegate_task`` dispatch dict (``status: "dispatched"`` - with a ``delegation_id``, or the synchronous result dict on channels that - cannot route async completions). Raises ValueError when there is nothing - to review or the dispatch is rejected/errored. + Returns the parsed ``delegate_task`` dispatch dict (``status: "dispatched"`` with a + ``delegation_id``, or the synchronous result dict on channels that cannot route async + completions). Raises ValueError when there is nothing to review or the dispatch is + rejected/errored. """ if parent_agent is None: raise ValueError("No active agent — send a message first.") @@ -222,7 +214,7 @@ def start_review( try: result = json.loads(raw) except Exception: - raise ValueError(f"Review dispatch failed: {raw!r}") + result = None if isinstance(result, dict) and result.get("error"): raise ValueError(str(result["error"])) if not isinstance(result, dict): diff --git a/agent/review_idle_queue.py b/agent/review_idle_queue.py index 033276fdd1..6f4ec3fe13 100644 --- a/agent/review_idle_queue.py +++ b/agent/review_idle_queue.py @@ -1,30 +1,20 @@ """Idle deferral for background reviews on the managed local runtime. -When the review runtime IS the managed llama-server, the post-turn review fork -monopolizes the GPU the user's next prompt needs, for minutes — and the next -live turn cancels it, so an active session pays the decode cost AND loses the -learning. This module keeps the decision to learn where it was (turn end, nudge -intervals, full model, full transcript) and moves only the execution moment: -reviews bound for the managed local endpoint are queued and dispatched when the -machine is quiet. Everything else runs immediately. +When the review runtime IS the managed llama-server, the post-turn review fork monopolizes +the GPU the user's next prompt needs, and the next live turn cancels it — an active +session pays the decode cost AND loses the learning. Only the execution moment moves: +reviews bound for the managed local endpoint are queued and dispatched when the machine +is quiet; everything else runs immediately. Policy ``auxiliary.background_review.defer``: +``auto`` (default) defers exactly when the review runtime targets the managed local +server; ``never`` is the old behavior. Explicit /refine (focus set) never defers. -Policy (auxiliary.background_review.defer): ``auto`` (default) defers exactly -when the resolved review runtime targets the managed local server; ``never`` is -the old behavior. Explicit /refine (focus set) never defers. - -Queue semantics: -- One slot per session, newest snapshot wins (a review replays the whole - conversation, so coalescing is deduplication, not loss). -- Preempted (cancelled-by-live-turn) reviews are requeued by the spawn wrapper - observing the run token's cancel flag. -- Aged-out events (defer_max_age_s, default 30 min) dispatch regardless of - idleness — deferral may delay learning, never lose it. -- In-memory, best-effort: dropped on process exit, like the immediate fork. - -Idle truth comes from the supervisor's /slots (machine-level, sees every client -incl. other profiles) and must hold for a settle window so a review is not -launched between two quick prompts. In-process turn liveness is tracked via -note_turn_started/note_turn_finished from run_conversation. +Queue semantics: one slot per session, newest snapshot wins (a review replays the whole +conversation, so coalescing is deduplication, not loss); preempted reviews are requeued by +the spawn wrapper observing the run token's cancel flag; aged-out events (defer_max_age_s, +default 30 min) dispatch regardless of idleness — deferral may delay learning, never lose +it; in-memory best-effort, dropped on process exit like the immediate fork. Idle truth +comes from the supervisor's /slots (machine-level, sees every client incl. other profiles) +and must hold for a settle window so a review is not launched between two quick prompts. """ from __future__ import annotations @@ -38,10 +28,10 @@ from typing import Any, Callable, Dict, Optional logger = logging.getLogger(__name__) -# Sustained-quiet window before dispatch: long enough that two back-to-back -# prompts do not look idle, short enough that a coffee break runs the queue. +# Sustained-quiet window before dispatch: two back-to-back prompts must not look idle, +# a coffee break must run the queue. _IDLE_SETTLE_S = 15.0 -# Poll cadence while the queue is non-empty. The thread parks when empty. +# Poll cadence while the queue is non-empty; the thread parks when empty. _POLL_INTERVAL_S = 5.0 # Age at which a queued review dispatches regardless of idleness. _MAX_AGE_DEFAULT_S = 30.0 * 60.0 @@ -66,11 +56,11 @@ def review_targets_managed_local(agent: Any, task_cfg: Optional[Dict[str, Any]]) -> bool: """Would this review fork decode on the llama-server WE manage? - Resolves the review runtime as the fork will and exact-matches its netloc - against the supervisor state file (cannot false-positive on external local - servers). Any failure reads False: immediate spawn is the safe default. - The netloc probe (one TTL-cached state-file read) runs FIRST so cloud-only - installs return False without resolving the runtime on the turn's tail. + Resolves the review runtime as the fork will and exact-matches its netloc against + the supervisor state file (cannot false-positive on external local servers). Any + failure reads False: immediate spawn is the safe default. The netloc probe (one + TTL-cached state-file read) runs FIRST so cloud-only installs return False without + resolving the runtime on the turn's tail. """ try: from agent.auxiliary_client import ( @@ -91,11 +81,11 @@ def review_targets_managed_local(agent: Any, class _PendingReview: __slots__ = ("agent", "kwargs", "enqueued_at", "session_key") - def __init__(self, agent: Any, session_key: str, kwargs: Dict[str, Any]): + def __init__(self, agent: Any, session_key: str, kwargs: Dict[str, Any], enqueued_at: float): self.agent = agent self.session_key = session_key self.kwargs = kwargs - self.enqueued_at = time.monotonic() + self.enqueued_at = enqueued_at class ReviewIdleQueue: @@ -130,16 +120,15 @@ class ReviewIdleQueue: def enqueue(self, agent: Any, session_key: str, kwargs: Dict[str, Any]) -> None: - """Add (or replace — newest snapshot wins) a session's pending review.""" + """Add (or replace — newest snapshot wins) a session's pending review. + + Keeps the ORIGINAL enqueue time on coalesce so a busy session cannot push its + review's age-out forever. + """ with self._lock: existing = self._pending.get(session_key) - item = _PendingReview(agent, session_key, kwargs) - # Stamp through the queue's clock (test seam); keep the ORIGINAL - # enqueue time on coalesce so a busy session cannot push its - # review's age-out forever. - item.enqueued_at = (existing.enqueued_at if existing is not None - else self._now()) - self._pending[session_key] = item + enqueued_at = existing.enqueued_at if existing is not None else self._now() + self._pending[session_key] = _PendingReview(agent, session_key, kwargs, enqueued_at) self._ensure_thread() self._wake.set() logger.info("Background review deferred (session=%s, queued=%d)", @@ -166,15 +155,13 @@ class ReviewIdleQueue: return self._now() - self._quiet_since def _pop_dispatchable(self) -> Optional[_PendingReview]: - """Oldest aged-out item, else any item once quiet+idle hold.""" + """Oldest aged-out item, else the oldest item once quiet+idle hold.""" with self._lock: if not self._pending: return None - items = sorted(self._pending.values(), - key=lambda p: p.enqueued_at) - aged = [p for p in items - if self._now() - p.enqueued_at - >= defer_max_age_s(p.kwargs.get("task_cfg"))] + now = self._now() + aged = [p for p in sorted(self._pending.values(), key=lambda p: p.enqueued_at) + if now - p.enqueued_at >= defer_max_age_s(p.kwargs.get("task_cfg"))] candidate = aged[0] if aged else None if candidate is None: if self._quiet_for() < _IDLE_SETTLE_S or not self._server_idle(): @@ -219,9 +206,8 @@ class ReviewIdleQueue: @staticmethod def _still_enabled(item: _PendingReview) -> bool: - """Re-check the enabled gate at DISPATCH time: minutes may pass in the - queue, and disabling reviews meanwhile must not be resurrected. Fail-open - like the gate itself.""" + """Re-check the enabled gate at DISPATCH time (minutes may pass in the queue; + disabling reviews meanwhile must not be resurrected). Fail-open like the gate.""" try: from agent.background_review import load_background_review_settings @@ -231,9 +217,8 @@ class ReviewIdleQueue: def _managed_server_idle() -> bool: - """Machine-level idle: no processing slot on any loaded model of the - managed router. Unreachable/no state file reads idle (nothing to - contend with). One /models + one /slots call per loaded model.""" + """Machine-level idle: no processing slot on any loaded model of the managed router. + Unreachable/no state file reads idle (nothing to contend with).""" try: from hermes_cli.local_runtime.supervisor import state_path from urllib.parse import quote diff --git a/agent/side_question.py b/agent/side_question.py index e5f5c88227..75bc74eba8 100644 --- a/agent/side_question.py +++ b/agent/side_question.py @@ -1,18 +1,11 @@ -"""Context-aware side questions (``/btw``). +"""Context-aware side questions (``/btw``): answer a question ABOUT the conversation without +touching it (no synthetic turns, no role-alternation risk, no prompt-cache invalidation). -Answers a quick question ABOUT the current conversation without touching it (no -synthetic turns, no role-alternation risk, no prompt-cache invalidation). Two -paths, picked automatically: - -* **Cache-parity fork (preferred).** With a live parent ``AIAgent``, a detached - fork from :func:`agent.background_review.build_cache_parity_fork` replays the - parent's snapshot verbatim against the warm prefix cache. Tool calls are denied - at dispatch, persistence is detached, usage goes to the parent. -* **One-shot digest (fallback).** With no live parent (e.g. the gateway evicted - the cached agent), a rendered transcript goes through :func:`agent.oneshot.run_oneshot`. - -``auxiliary.side_question.provider`` / ``.model`` route the fork to another model -and replay a compact digest (cold cache on a different model). +Preferred path: a detached cache-parity fork of the live parent ``AIAgent`` replays its +snapshot verbatim against the warm prefix cache (tools denied at dispatch, persistence +detached, usage attributed to the parent). Fallback (no live parent, e.g. gateway evicted +the agent): a rendered transcript through :func:`agent.oneshot.run_oneshot`. +``auxiliary.side_question.provider``/``.model`` route the fork elsewhere with a compact digest. """ import logging @@ -65,10 +58,8 @@ _ROLE_LABELS = {"user": "USER", "assistant": "ASSISTANT", "tool": "TOOL RESULT"} def trim_snapshot_for_fork(history: Optional[List[Dict[str, Any]]]) -> List[Dict[str, Any]]: """Drop trailing messages until the snapshot ends with a completed assistant text. - A mid-turn snapshot can end on unresolved ``tool_calls``, a tool result, or - the in-flight user message; appending the side question after any of those - breaks role alternation on strict providers. Trimming only the TAIL keeps - the warm prefix-cache property. + A mid-turn tail (unresolved ``tool_calls``, tool result, in-flight user message) breaks + role alternation on strict providers; trimming only the TAIL keeps the warm prefix cache. """ msgs = list(history or []) while msgs: @@ -83,12 +74,8 @@ def render_history_for_side_question( history: Optional[List[Dict[str, Any]]], char_budget: int = _TRANSCRIPT_CHAR_BUDGET, ) -> str: - """Render a snapshot as a plain-text transcript (fallback path only). - - Newest-biased fit to ``char_budget``. Tool calls are summarized by name, tool - results included truncated (so "what did that output" stays answerable), the - system prompt skipped. - """ + """Plain-text transcript for the fallback path: newest-biased fit to ``char_budget``, + tool calls summarized by name, tool results truncated, system prompt skipped.""" lines: List[str] = [] for msg in history or []: if not isinstance(msg, dict): @@ -134,9 +121,8 @@ def _side_question_task_config() -> Dict[str, Any]: def _answer_via_fork(parent_agent: Any, question: str, history: Optional[List[Dict[str, Any]]]) -> str: """Answer via a cache-parity fork of ``parent_agent`` on the calling thread. - The thread-scoped tool whitelist is emptied so any tool call is denied at - dispatch: ``tools[]`` stays byte-identical for cache parity, but the side - question can never mutate anything. + An empty thread-scoped tool whitelist denies every tool call at dispatch: ``tools[]`` + stays byte-identical for cache parity, but the side question can never mutate anything. """ from agent.background_review import ( _digest_history, @@ -208,9 +194,8 @@ def answer_side_question( temperature: Optional[float] = 0.3, timeout: float = 180.0, ) -> str: - """Answer ``question`` against a snapshot of ``history``: cache-parity fork when - ``parent_agent`` is live, else (or on empty answer / failure) the one-shot digest. - Raises on failure — callers surface the error on their own UI.""" + """Fork when ``parent_agent`` is live, else (or on empty answer / failure) the one-shot + digest. Raises on failure — callers surface the error on their own UI.""" question = (question or "").strip() if not question: raise ValueError("answer_side_question requires a non-empty question")