diff --git a/agent/activity_tracking.py b/agent/activity_tracking.py index 1ce2cf160a..5b0e654d00 100644 --- a/agent/activity_tracking.py +++ b/agent/activity_tracking.py @@ -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, diff --git a/agent/api_error_summary.py b/agent/api_error_summary.py index 27ce350f1e..62c3bf15d7 100644 --- a/agent/api_error_summary.py +++ b/agent/api_error_summary.py @@ -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 diff --git a/agent/client_lifecycle.py b/agent/client_lifecycle.py index bc19d74936..fd1ae1f395 100644 --- a/agent/client_lifecycle.py +++ b/agent/client_lifecycle.py @@ -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 diff --git a/agent/interrupt_control.py b/agent/interrupt_control.py index f852f0a2fe..e6cf998acf 100644 --- a/agent/interrupt_control.py +++ b/agent/interrupt_control.py @@ -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 diff --git a/agent/rate_limit_credits.py b/agent/rate_limit_credits.py index 6b3643530b..8fa051d037 100644 --- a/agent/rate_limit_credits.py +++ b/agent/rate_limit_credits.py @@ -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 diff --git a/agent/session_persistence.py b/agent/session_persistence.py index 99bb85fa53..ac2808cc86 100644 --- a/agent/session_persistence.py +++ b/agent/session_persistence.py @@ -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 diff --git a/agent/status_output.py b/agent/status_output.py index 03c28257c8..72b35d7a41 100644 --- a/agent/status_output.py +++ b/agent/status_output.py @@ -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 diff --git a/agent/stream_delivery.py b/agent/stream_delivery.py index cbca070871..f1388afda9 100644 --- a/agent/stream_delivery.py +++ b/agent/stream_delivery.py @@ -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: diff --git a/run_agent.py b/run_agent.py index 6563e4d43a..2e1ad348ea 100644 --- a/run_agent.py +++ b/run_agent.py @@ -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.
.models.