refactor(agent/review_*,side_question): hoist review goal constant, single-lock pending stamp, compact docstrings
This commit is contained in:
+36
-44
@@ -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):
|
||||
|
||||
+39
-54
@@ -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
|
||||
|
||||
+15
-30
@@ -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")
|
||||
|
||||
Reference in New Issue
Block a user