refactor(run_agent): re-wrap compacted docstrings, restore lost issue-ref rationale

This commit is contained in:
Teknium
2026-09-02 12:11:33 -07:00
parent e564e60a3d
commit 1e778ae7c6
9 changed files with 73 additions and 114 deletions
+4 -6
View File
@@ -40,12 +40,10 @@ class ActivityTrackingMixin:
"""Update the last-activity timestamp and description (thread-safe).
Bumps a monotonic generation under ``_liveness_activity_lock`` so the watchdog can bind a stall
observation
to the exact ``(generation, timestamp)`` it sampled. Also bridges (rate-limited, best-effort) to the
kanban
heartbeat when this is a dispatcher-spawned worker, and to the durable SessionDB activity projection.
``provenance`` names special writers (compression); ``force_persist`` bypasses the SessionDB rate
limit.
observation to the exact ``(generation, timestamp)`` it sampled. Also bridges (rate-limited,
best-effort) to the kanban heartbeat when this is a dispatcher-spawned worker, and to the durable
SessionDB activity projection. ``provenance`` names special writers (compression); ``force_persist``
bypasses the SessionDB rate limit.
"""
from agent.session_activity import (
bound_activity_description,
+3 -4
View File
@@ -22,10 +22,9 @@ class ApiErrorSummaryMixin:
"""Detect subscription/entitlement 401/403s that masquerade as auth failures.
Refreshing a token cannot fix an unsubscribed account, so callers surface the error instead of looping
the
pool. xAI returns the same permission-denied text for BOTH cases; a ``[WKE=unauthenticated:...]``
suffix (or
"access token could not be validated") means stale token → return False so the refresh path runs.
the pool. xAI returns the same permission-denied text for BOTH cases; a ``[WKE=unauthenticated:...]``
suffix (or "access token could not be validated") means stale token → return False so the refresh path
runs.
"""
if status_code not in {401, 403, None}:
return False
+18 -30
View File
@@ -89,8 +89,7 @@ class ClientLifecycleMixin:
"""Check if an OpenAI client is closed.
``is_closed`` is a bool property on httpx.Client but a method on openai.OpenAI; a bare getattr
returned
the always-truthy bound method and recreated the client on every call.
returned the always-truthy bound method and recreated the client on every call.
"""
from unittest.mock import Mock
@@ -151,10 +150,8 @@ class ClientLifecycleMixin:
``close()`` releases raw FDs from the calling thread; the shared client has no owning thread and other
threads may still hold its fd in an SSL BIO. A recycled fd then gets a TLS record written into an
unrelated
file (the SQLite-header corruption family). So: ``shutdown()`` the sockets (FD-safe from any thread)
and
let GC release the FDs once every borrower has unwound.
unrelated file (the SQLite-header corruption family: #29507 / #67142 / #70773). So: ``shutdown()``
the sockets (FD-safe from any thread) and let GC release the FDs once every borrower has unwound.
"""
if client is None:
return
@@ -179,9 +176,8 @@ class ClientLifecycleMixin:
"""FD-safe transport drain for an abandoned (timed-out) worker; returns sockets shut down.
The worker may be blocked in an OpenSSL read; hard-closing from the timeout thread releases FDs under
a live
BIO (native corruption / SIGSEGV). Only ``shutdown()`` so the read settles with EOF and the worker
closes itself.
a live BIO (native corruption / SIGSEGV, #94248). Only ``shutdown()`` so the read settles with EOF and
the worker closes itself.
"""
drained = 0
# Shared primary client (codex-direct / MoA stream on it directly).
@@ -449,9 +445,8 @@ class ClientLifecycleMixin:
"""Cross-thread abort: shut sockets down without releasing FDs.
For stranger-thread callers (interrupt loop, stale detector). ``close()`` from a non-owning thread
raced the
live SSL BIO and corrupted unrelated FDs; ``shutdown(SHUT_RDWR)`` unblocks the owner's recv/send so it
closes from its own context.
raced the live SSL BIO and corrupted unrelated FDs; ``shutdown(SHUT_RDWR)`` unblocks the owner's
recv/send so it closes from its own context.
"""
if client is None:
return
@@ -513,12 +508,10 @@ class ClientLifecycleMixin:
"""Build (or reuse) a request-local Anthropic client for one in-flight call.
The stale/interrupt watchdog must never ``close()`` the client a worker is still reading (fd recycled
under
a live SSL BIO → TLS record in a SQLite header). A per-request client lets the stranger ``shutdown()``
while
the owner closes. Single-slot cache keyed as ``_request_anthropic_client_key``; ``in_use`` gives a
second
concurrent call a fresh untracked client. Mirrors ``_rebuild_anthropic_client`` construction.
under a live SSL BIO → TLS record in a SQLite header). A per-request client lets the stranger
``shutdown()`` while the owner closes. Single-slot cache keyed as ``_request_anthropic_client_key``;
``in_use`` gives a second concurrent call a fresh untracked client. Mirrors
``_rebuild_anthropic_client`` construction.
"""
if self.api_mode == "anthropic_messages":
self._try_refresh_anthropic_client_credentials()
@@ -808,12 +801,10 @@ class ClientLifecycleMixin:
"""Adopt ~/.hermes/.env credential/base-url edits at the turn boundary.
A Settings save updates ``.env`` but a live worker keeps init-time values, so an open chat kept
calling the
old endpoint. Reacts only to env *edits* (resolved value changed since last look), never to divergence
from
the agent's current values — pool rotation/failover and a config ``model.base_url`` legitimately move
the
session and must not flap. Covers registry providers and named custom providers with ``key_env``.
calling the old endpoint. Reacts only to env *edits* (resolved value changed since last look), never
to divergence from the agent's current values — pool rotation/failover and a config ``model.base_url``
legitimately move the session and must not flap. Covers registry providers and named custom providers
with ``key_env``.
"""
if self.api_mode != "chat_completions":
return False
@@ -995,10 +986,8 @@ class ClientLifecycleMixin:
"""Refresh Copilot credentials and rebuild the shared OpenAI client.
The raw GitHub token is stable, but the short-TTL *exchanged* IDE token is what authenticates and
expires
mid-turn (``401 IDE token expired``). Re-resolving the raw token leaves the same expired JWT on the
wire, so
force a fresh exchange. Caller enforces the single-shot guard.
expires mid-turn (``401 IDE token expired``). Re-resolving the raw token leaves the same expired JWT
on the wire, so force a fresh exchange. Caller enforces the single-shot guard.
"""
if not self._is_copilot_provider():
return False
@@ -1199,8 +1188,7 @@ class ClientLifecycleMixin:
Lets custom endpoints behind a WAF that rejects the SDK's identifying headers (``User-Agent``,
``X-Stainless-*``) work. Delegates to ``agent.auxiliary_client._apply_user_default_headers`` so main
and
auxiliary clients cannot drift. No-op for Anthropic/Bedrock modes.
and auxiliary clients cannot drift. No-op for Anthropic/Bedrock modes.
"""
if self.api_mode in ("anthropic_messages", "bedrock_converse"):
return
+7 -10
View File
@@ -28,13 +28,11 @@ class InterruptControlMixin:
"""Request the agent to interrupt its current tool-calling loop (call from another thread).
``message``: new message to include in the response context. ``hard_cancel``: explicit stop;
compression
may honor it even while ordinary interrupts are masked. ``tool_reason``: trusted fixed category safe
for
tool output. ``require_generation``: activity-generation claim — the interrupt is published only if
the
turn's generation still matches at the final mutation edge (claim reserved under the activity lock,
consumed together with the first observable publication); returns False if the turn resumed meanwhile.
compression may honor it even while ordinary interrupts are masked. ``tool_reason``: trusted fixed
category safe for tool output. ``require_generation``: activity-generation claim — the interrupt is
published only if the turn's generation still matches at the final mutation edge (claim reserved under
the activity lock, consumed together with the first observable publication); returns False if the turn
resumed meanwhile.
"""
if require_generation is not None:
# RESERVE the abort's generation claim under the SAME lock `_touch_activity` stamps with. Real
@@ -320,9 +318,8 @@ class InterruptControlMixin:
During a model request this cancels only that request: completed messages/tool results are kept, the
displayed partial reasoning becomes assistant context, the correction is appended as a real user
message,
and the loop retries. During tool execution it degrades to ``steer()``; Codex app-server uses native
``turn/steer``. Returns False when there is no live turn or the text is empty.
message, and the loop retries. During tool execution it degrades to ``steer()``; Codex app-server uses
native ``turn/steer``. Returns False when there is no live turn or the text is empty.
"""
if not text or not text.strip():
return False
+3 -5
View File
@@ -51,8 +51,7 @@ class RateLimitCreditsMixin:
"""Parse x-nous-credits-* headers, cache CreditsState, fire threshold notices.
The PARSE is swallowed (miss → keep last-known); the notice EVALUATION is a separate block that WARNS
on
failure so a depletion-notice bug cannot vanish silently.
on failure so a depletion-notice bug cannot vanish silently.
"""
# Dev test fixture (HERMES_DEV_CREDITS_FIXTURE): inject a chosen notice state
# each turn for repeatable testing, bypassing real headers. Throwaway scaffolding.
@@ -132,9 +131,8 @@ class RateLimitCreditsMixin:
"""Run the threshold policy on the current credits state and emit notices.
Shared by the warm path and the cold-start seed so an already-depleted session warns immediately. Runs
only
when a notice consumer is bound. WARNS on failure. Emits clears FIRST so depleted lands last (latest-
wins slot).
only when a notice consumer is bound. WARNS on failure. Emits clears FIRST so depleted lands last
(latest- wins slot).
"""
if getattr(self, "notice_callback", None) is None and getattr(self, "notice_clear_callback", None) is None:
return
+4 -7
View File
@@ -85,8 +85,7 @@ class SessionPersistenceMixin:
"""Rewrite the current-turn user message before persistence/return.
Some paths use an API-only user-message variant that must not leak into transcripts or resumed
history;
mutate the in-memory list in place so both persistence and returned history stay clean.
history; mutate the in-memory list in place so both persistence and returned history stay clean.
"""
idx = getattr(self, "_persist_user_message_idx", None)
override = getattr(self, "_persist_user_message_override", None)
@@ -214,9 +213,8 @@ class SessionPersistenceMixin:
"""Persist any un-flushed messages to the SQLite session store.
Dedup is an intrinsic ``_DB_PERSISTED_MARKER`` on each written dict — not positional slices (drift
after
sequence repair) nor a retained ``id(msg)`` set (address reuse). ``_flushed_db_message_ids`` is only a
one-shot seed translated to markers and cleared each flush.
after sequence repair) nor a retained ``id(msg)`` set (address reuse). ``_flushed_db_message_ids`` is
only a one-shot seed translated to markers and cleared each flush.
"""
# Persistence-isolated agents (background review fork) share the parent's session_id for cache
# warmth; a write here would land the curator's harness turn in the user's real history. Hard-stop.
@@ -583,8 +581,7 @@ class SessionPersistenceMixin:
"""Optional per-session JSON snapshot writer (``sessions.write_json_snapshots``, default False).
state.db is canonical; this exists for external tooling reading ``session_{sid}.json``. Rewrites the
full
list after every persistence point, never overwriting a larger log with fewer messages.
full list after every persistence point, never overwriting a larger log with fewer messages.
"""
if not getattr(self, "_session_json_enabled", False):
return
+1 -2
View File
@@ -48,8 +48,7 @@ class StatusOutputMixin:
"""Return True when quiet-mode spinner output has a safe sink.
A raw spinner falling back to ``sys.stdout`` can corrupt protocol streams (ACP JSON-RPC); allow it
only
when output is rerouted via ``_print_fn`` or stdout is a real TTY.
only when output is rerouted via ``_print_fn`` or stdout is a real TTY.
"""
if self._print_fn is not None:
return True
+2 -4
View File
@@ -263,10 +263,8 @@ class StreamDeliveryMixin:
token.
Every attempt (each provider path, each retry) claims right before consuming. Claiming bumps the
shared
token, so an earlier attempt still alive on another thread is superseded and its late chunks fenced
out.
Stored per-thread: a thread that never claimed is never a writer and can never be fenced.
shared token, so an earlier attempt still alive on another thread is superseded and its late chunks
fenced out. Stored per-thread: a thread that never claimed is never a writer and can never be fenced.
"""
self._ensure_stream_writer_state()
with self._stream_writer_lock:
+31 -46
View File
@@ -518,8 +518,7 @@ class AIAgent(
"""Notify the active context engine about a host session transition.
The built-in compressor keeps its reset behavior; plugin engines with richer hooks (``on_session_end``
/
``on_session_reset`` / ``on_session_start`` / ``carry_over_new_session_context``) can flush, rebind
/ ``on_session_reset`` / ``on_session_start`` / ``carry_over_new_session_context``) can flush, rebind
and carry context.
"""
engine = getattr(self, "context_compressor", None)
@@ -581,8 +580,7 @@ class AIAgent(
"""Reset all session-scoped token/cost counters and compressor state for a fresh session.
When ``previous_messages`` / ``old_session_id`` / ``carry_over_context`` are given, the context engine
gets
the full transition lifecycle (``_transition_context_engine_session``) instead of a bare reset.
gets the full transition lifecycle (``_transition_context_engine_session``) instead of a bare reset.
"""
# Token usage counters
self.session_total_tokens = 0
@@ -704,9 +702,8 @@ class AIAgent(
"""Disable Responses encrypted reasoning replay and strip cached state.
Called on HTTP 400 ``invalid_encrypted_content``. Sets ``_codex_reasoning_replay_enabled=False``
(consumed by
the codex adapter/transport) and pops ``codex_reasoning_items`` from every assistant message.
Returns ``{"messages": int, "items": int}`` for diagnostic logging.
(consumed by the codex adapter/transport) and pops ``codex_reasoning_items`` from every assistant
message. Returns ``{"messages": int, "items": int}`` for diagnostic logging.
"""
stripped_messages = 0
stripped_items = 0
@@ -742,8 +739,7 @@ class AIAgent(
"""Return True for malformed provider streaming data from SDK parsers.
The Anthropic SDK surfaces a malformed event-stream frame as a plain ``ValueError``; that is wire-
format
trouble, not local validation, so it follows the truncated-JSON retry path.
format trouble, not local validation, so it follows the truncated-JSON retry path.
"""
if getattr(self, "api_mode", None) != "anthropic_messages":
return False
@@ -1104,9 +1100,9 @@ class AIAgent(
"""Detect Ollama-hosted GLM models affected by finish_reason='stop' misreports.
Matches only explicit Ollama signatures (port 11434, "ollama" in URL, provider ollama) — never
arbitrary
local proxies, which report correctly. Excludes Ollama Cloud (``ollama.com`` host, ``:cloud`` suffix):
rewriting its stop→length manufactures false truncations and burns the continuation budget.
arbitrary local proxies, which report correctly. Excludes Ollama Cloud (``ollama.com`` host,
``:cloud`` suffix): rewriting its stop→length manufactures false truncations and burns the
continuation budget.
"""
model_lower = (self.model or "").lower()
provider_lower = (self.provider or "").lower()
@@ -1418,9 +1414,8 @@ class AIAgent(
"""Mirror a completed turn into external memory providers (``sync_all`` + ``queue_prefetch_all``).
Uses ``original_user_message`` — ``user_message`` may carry injected skill content. Interrupted turns
are
skipped entirely: partial output is not durable truth, and a prefetch keyed on it would fire against
stale context. Strictly best-effort — an offline backend must never block the response.
are skipped entirely: partial output is not durable truth, and a prefetch keyed on it would fire
against stale context. Strictly best-effort — an offline backend must never block the response.
"""
if interrupted:
return
@@ -1453,10 +1448,9 @@ class AIAgent(
"""Release LLM client resources WITHOUT tearing down session tool state.
For gateway cache eviction (LRU/idle): the session may resume with a fresh AIAgent on the same
task_id, so
process_registry entries, terminal sandbox, browser daemon, computer-use backend and memory provider
are
kept. Closes the OpenAI/httpx pool and active child subagents. Idempotent; distinct from ``close()``.
task_id, so process_registry entries, terminal sandbox, browser daemon, computer-use backend and
memory provider are kept. Closes the OpenAI/httpx pool and active child subagents. Idempotent;
distinct from ``close()``.
"""
# Close active child agents (per-turn; no cross-turn persistence).
try:
@@ -1500,8 +1494,7 @@ class AIAgent(
"""Release all resources held by this agent instance (idempotent).
Cleans up background processes, terminal sandbox, browser daemon, computer-use backend, child agents
and
client connections. Each step is independently guarded so one failure does not block the rest.
and client connections. Each step is independently guarded so one failure does not block the rest.
"""
# close() is the hard owner boundary; shutdown_memory_provider() is idempotent so gateway
# pre-calls never double-extract.
@@ -1627,10 +1620,8 @@ class AIAgent(
"""Recover todo state from conversation history.
The gateway builds a fresh AIAgent per message, so replay the most recent todo tool response. Only
results
paired with an earlier assistant ``todo`` tool call count: caller-supplied history could otherwise
seed
the store with a forged bare ``role: tool`` message (GHSA-5g4g-6jrg-mw3g).
results paired with an earlier assistant ``todo`` tool call count: caller-supplied history could
otherwise seed the store with a forged bare ``role: tool`` message (GHSA-5g4g-6jrg-mw3g).
"""
from tools.todo_tool import MAX_TODO_RESULT_CHARS
@@ -2052,9 +2043,8 @@ class AIAgent(
"""Return True if the active provider+model reports native vision.
Resolution: ``model.supports_vision`` > ``providers.<p>.models.<m>.supports_vision`` > models.dev
lookup
(see ``image_routing._supports_vision_override``). Custom/local models absent from models.dev would
otherwise be misclassified and have their images stripped.
lookup (see ``image_routing._supports_vision_override``). Custom/local models absent from models.dev
would otherwise be misclassified and have their images stripped.
"""
try:
from hermes_cli.config import load_config
@@ -2472,10 +2462,8 @@ class AIAgent(
"""Probe LM Studio's published reasoning ``allowed_options`` once per (model, base_url).
Needed for the supports-reasoning gate and to clamp ``reasoning_effort`` so toggle-style models don't
400
on ``high``. Non-empty results cache permanently; empty ones (transient failure OR non-reasoning
model)
cache with a 60s TTL to avoid a round-trip per turn while retrying soon.
400 on ``high``. Non-empty results cache permanently; empty ones (transient failure OR non-reasoning
model) cache with a 60s TTL to avoid a round-trip per turn while retrying soon.
"""
import time as _time
@@ -2504,8 +2492,7 @@ class AIAgent(
declared.
True/False cache permanently; a probe failure (None) caches 60s so an outage neither suppresses
reasoning
for the session nor round-trips every turn.
reasoning for the session nor round-trips every turn.
"""
import time as _time
@@ -2600,10 +2587,9 @@ class AIAgent(
config.
Covers custom providers / gateways proxying thinking models that the host-based
``_REASONING_ECHO_RULES``
miss. Per-active-provider: primary from ``model.reasoning_echo``, fallback from the fallback entry's
field,
restored by ``restore_primary_runtime()`` — so falling back to a strict provider still strips it.
``_REASONING_ECHO_RULES`` miss. Per-active-provider: primary from ``model.reasoning_echo``, fallback
from the fallback entry's field, restored by ``restore_primary_runtime()`` — so falling back to a
strict provider still strips it.
"""
return bool(getattr(self, "_reasoning_echo_flag", False))
@@ -2622,8 +2608,8 @@ class AIAgent(
"""Return True when the current provider is Kimi / Moonshot thinking mode (requires
``reasoning_content`` echo).
Host-driven, not model-name-driven: aggregators re-exporting Kimi reject the echo. Rule table:
``message_sanitization.reasoning_echo_family``.
Host-driven, not model-name-driven: aggregators re-exporting Kimi reject the echo (#17400). Rule
table: ``message_sanitization.reasoning_echo_family``.
"""
from agent.message_sanitization import matches_reasoning_echo_family
return matches_reasoning_echo_family(
@@ -2634,7 +2620,8 @@ class AIAgent(
"""Return True when the current provider is DeepSeek thinking mode (requires ``reasoning_content``
echo).
Rule table: ``message_sanitization.reasoning_echo_family``.
Omitting the echo on replayed assistant tool-call turns is an HTTP 400 (#15250). Rule table:
``message_sanitization.reasoning_echo_family``.
"""
from agent.message_sanitization import matches_reasoning_echo_family
return matches_reasoning_echo_family(
@@ -3118,9 +3105,8 @@ class AIAgent(
"""Execute tool calls from the assistant message and append results to messages.
The segment planner splits the batch into maximal runs of parallel-safe calls (read-only, non-
overlapping
file targets, opted-in MCP) separated by sequential barriers; mixed batches run segment by segment in
emission order so safe subsets stay concurrent while side-effect ordering is preserved.
overlapping file targets, opted-in MCP) separated by sequential barriers; mixed batches run segment by
segment in emission order so safe subsets stay concurrent while side-effect ordering is preserved.
"""
tool_calls = assistant_message.tool_calls
@@ -3587,8 +3573,7 @@ class AIAgent(
"""Stop lease renewal after a committed liveness abort.
A wedge the hard interrupt cannot unwind must not keep the lease alive forever; TTL expiry
lets
stale-turn cleanup reclaim the row.
lets stale-turn cleanup reclaim the row.
"""
nonlocal durable_turn_lease_turn_active
with durable_turn_lease_activity_lock: