diff --git a/agent/relay_llm.py b/agent/relay_llm.py index 199e8b73c4..bbe138adf5 100644 --- a/agent/relay_llm.py +++ b/agent/relay_llm.py @@ -95,10 +95,8 @@ class _ManagedAttempt: def run_callback(self, callback: Callable[..., Any], *args: Any) -> Any: """Run a Hermes callback in a fresh copy of the captured context. - Relay can invoke callbacks while another still owns the captured Context (hence the - copy); nested relay calls run unmanaged — see relay_runtime.managed_callback_guard. - """ + copy); nested relay calls run unmanaged — see relay_runtime.managed_callback_guard.""" def guarded() -> Any: with relay_runtime.managed_callback_guard(): return callback(*args) @@ -142,9 +140,7 @@ class _ManagedAttempt: def resolve_failure(self, exc: BaseException, defer_logical_completion: bool) -> Any: """Re-raise the provider's own error, or recover a completed provider result. - - Must be called from the ``except`` handling ``exc`` (bare ``raise``). - """ + Must be called from the ``except`` handling ``exc`` (bare ``raise``).""" callback_error = self.raw_response.get("error") if ( callback_error is not None @@ -246,13 +242,11 @@ def stream_current( defer_logical_completion: bool = False, completed_response_predicate: Callable[[Any], bool] | None = None, ) -> Any: """Run a provider stream under the inherited Hermes turn when present. - With ``completed_response_predicate`` set, a factory that ignores ``stream=True`` and returns a complete response is unwrapped and returned directly (pre-Relay behavior) instead of staying trapped as ``final_response``. Detecting that primes the lazy pipeline: a genuine first chunk is buffered, but provider latency and pre-first-yield - errors may surface before this returns. - """ + errors may surface before this returns.""" session_id = _current_session_id() if session_id is None: return stream_factory(request) diff --git a/agent/relay_runtime.py b/agent/relay_runtime.py index 301ddc8f9b..4593cc8843 100644 --- a/agent/relay_runtime.py +++ b/agent/relay_runtime.py @@ -48,10 +48,8 @@ def runtime_metadata(runtime_id: str, **extra: Any) -> dict[str, Any]: def _scope_op_executor(): """Shared daemon executor for bounded native scope ops. - Daemon workers so a wedged call abandoned at timeout cannot block interpreter exit; - ``Future.result(timeout=...)`` still bounds callers when every worker is wedged. - """ + ``Future.result(timeout=...)`` still bounds callers when every worker is wedged.""" global _SCOPE_OP_EXECUTOR if _SCOPE_OP_EXECUTOR is None: with _SCOPE_OP_EXECUTOR_LOCK: @@ -67,10 +65,8 @@ def _run_on_daemon_thread( fn: Callable[[], Any], *, name: str, timeout: float | None = None, timeout_message: str = "" ) -> Any: """Run ``fn`` on a fresh daemon thread; re-raise its error or return its result. - With ``timeout`` a still-running worker is abandoned with ``TimeoutError`` (daemon: - cannot block interpreter exit). - """ + cannot block interpreter exit).""" outcome: dict[str, Any] = {} def _target() -> None: @@ -93,9 +89,7 @@ def pop_relay_scope( relay: Any, handle: Any, *, output: Any = None, metadata: Any = None, timestamp: Any = None ) -> Any: """Pop a Relay scope, forwarding only the kwargs the live binding accepts. - - ``scope.pop`` gained ``metadata`` in nemo-relay 0.4+; older wheels raise TypeError. - """ + ``scope.pop`` gained ``metadata`` in nemo-relay 0.4+; older wheels raise TypeError.""" pop = relay.scope.pop candidates = (("output", output), ("metadata", metadata), ("timestamp", timestamp)) kwargs = {key: value for key, value in candidates if value is not None} @@ -424,10 +418,8 @@ class RelayRuntime: self, context: contextvars.Context, *, exit_fallback: bool = False, **push_kwargs: Any ) -> Any: """Push a SESSION_SCOPE Agent scope inside ``context``, bounded by ``_SCOPE_OP_TIMEOUT``. - ``exit_fallback``: at interpreter shutdown the executor refuses new futures; push - synchronously instead (no agent turn waits at exit). - """ + synchronously instead (no agent turn waits at exit).""" args = (self.relay.scope.push, SESSION_SCOPE, self.relay.ScopeType.Agent) try: future = _scope_op_executor().submit(context.run, *args, input={}, **push_kwargs) @@ -442,10 +434,8 @@ class RelayRuntime: **push_kwargs: Any, ) -> None: """Push a fresh session scope for ``session`` and record its handle + context. - Subagents parent under their spawning turn/session handle; ``resolve_parent`` - creates the parent session when its handle is unknown. - """ + creates the parent session when its handle is unknown.""" parent_handle = None if session.parent_session_id: with self._sessions_lock: @@ -492,11 +482,9 @@ class RelayRuntime: def rotate_session_scope(self, session: RelaySession, *, reason: str) -> None: """Close the current session scope and open the next segment. - Called ONLY at a turn boundary: the stack is LIFO and rotating under a live child would close a parent out of order. Bookkeeping advances even when a native call - fails so a degraded rotation cannot retry on every turn. - """ + fails so a degraded rotation cannot retry on every turn.""" with session.lock: if session.closing or session.handle is None: return @@ -596,11 +584,9 @@ class RelayRuntime: allow_closing: bool = False, timeout: float | None = None, **kwargs: Any, ) -> Any: """Run a Relay operation against a session's isolated scope stack. - ``timeout`` bounds the native call on the daemon executor (``TimeoutError`` on breach); ``None`` runs synchronously. Lifecycle ops gating turn/session completion - pass ``_SCOPE_OP_TIMEOUT``: a wedged pipeline must cost one span, never the agent. - """ + pass ``_SCOPE_OP_TIMEOUT``: a wedged pipeline must cost one span, never the agent.""" self._begin_operation() try: return self._run_in_session_untracked( @@ -701,10 +687,8 @@ class RelayRuntime: drain_limit: int, ) -> BaseException | None: """Pop ``handle``; if that fails, drain orphans above it and retry once. - Returns the retry's error (None on success). Must run inside ONE ``run_in_session`` - callback so ContextVar stack views stay consistent. - """ + callback so ContextVar stack views stay consistent.""" with contextlib.suppress(Exception): pop_relay_scope(self.relay, handle, output=output, metadata=metadata) return None @@ -737,11 +721,9 @@ class RelayRuntime: drain_limit: int = 32, operation_already_held: bool = False, ) -> str | None: """Pop ``handle``, draining orphaned children in the same session context. - Relay scopes are strict LIFO; empty-stream retries + interrupt can abandon a physical LLM scope above TURN/SESSION. Drain+close is bounded so a wedged pipeline - never blocks turn/session completion. Returns a failure string or None. - """ + never blocks turn/session completion. Returns a failure string or None.""" if handle is None: return None run_in_session = (self._run_in_session_untracked if operation_already_held else self.run_in_session) @@ -953,10 +935,8 @@ _MANAGED_CALLBACK_DEPTH: contextvars.ContextVar[int] = contextvars.ContextVar( class managed_callback_guard: """Mark the current context as inside a managed Relay callback. - Wrap the ``invoke()`` callbacks handed to the native pipeline; everything they - transitively call (incl. work forwarded via copy_context()) runs unmanaged. - """ + transitively call (incl. work forwarded via copy_context()) runs unmanaged.""" def __enter__(self) -> "managed_callback_guard": self._token = _MANAGED_CALLBACK_DEPTH.set(_MANAGED_CALLBACK_DEPTH.get() + 1) @@ -1128,11 +1108,9 @@ class RelaySessionCoordinator: def _consume_deferred_close(self, lease: Any) -> None: """Close a session whose rotating-compaction close was deferred. - ``notify_session_compacted`` sets ``close_pending`` when the old session had a live turn (closing then breaks LIFO). The last live turn consumes it here after its own - scope popped and it left the active-turn table. - """ + scope popped and it left the active-turn table.""" # Telemetry must never block end_turn. _warn_on_error("deferred session close", self._consume_deferred_close_unguarded, lease) @@ -1149,13 +1127,11 @@ class RelaySessionCoordinator: self, *, profile_key: str, session_id: str, old_session_id: str = "" ) -> None: """React to a completed compaction, per compaction mode. - In-place (``old_session_id`` empty/equal): flag rotation for the next turn boundary — never rotate immediately, a turn may be live and rotating under it breaks LIFO. Rotating (ids differ): the next turn gets a fresh session under the new id, so close the OLD session now or its scope stays an unexported orphan. Unknown sessions and - disabled config are silent no-ops. - """ + disabled config are silent no-ops.""" # Telemetry must never block compaction. _warn_on_error( "compaction notification", self._notify_session_compacted_unguarded, diff --git a/agent/transports/__init__.py b/agent/transports/__init__.py index 8164e2b319..e7864ce079 100644 --- a/agent/transports/__init__.py +++ b/agent/transports/__init__.py @@ -1,8 +1,6 @@ """Transport registry for provider response normalization. - transport = get_transport("anthropic_messages") - result = transport.normalize_response(raw_response) -""" + result = transport.normalize_response(raw_response)""" import importlib diff --git a/agent/transports/anthropic.py b/agent/transports/anthropic.py index 2d2f1a9c10..49360a91da 100644 --- a/agent/transports/anthropic.py +++ b/agent/transports/anthropic.py @@ -11,11 +11,9 @@ _THINKING_TYPES = ("thinking", "redacted_thinking") def _unprefix_oauth_tool_name(name: str) -> str: """Reverse the OAuth-wire ``mcp__`` prefix back to the registered tool name. - Two originals map onto one wire name (``read_file`` / ``mcp_linear_get_issue``), so resolve by registry lookup, never rewriting a name that already resolves natively. - OAuth wire aliases are checked LAST so a real tool under the wire name still wins. - """ + OAuth wire aliases are checked LAST so a real tool under the wire name still wins.""" from agent.anthropic_adapter import _OAUTH_TOOL_NAME_REVERSE_ALIASES from tools.registry import registry as _tool_registry bare = name[len(_MCP_PREFIX):] diff --git a/agent/transports/base.py b/agent/transports/base.py index 07b8a88ac3..e53a7265ce 100644 --- a/agent/transports/base.py +++ b/agent/transports/base.py @@ -1,9 +1,7 @@ """Abstract base for provider transports. - A transport owns one api_mode's data path (convert_messages -> convert_tools -> build_kwargs -> normalize_response), NOT client construction, streaming, credentials, caching, interrupts -or retries — those stay on AIAgent. -""" +or retries — those stay on AIAgent.""" from abc import ABC, abstractmethod from typing import Any, Dict, List, Optional