diff --git a/agent/relay_runtime.py b/agent/relay_runtime.py index ebe1c25758..233967f1fb 100644 --- a/agent/relay_runtime.py +++ b/agent/relay_runtime.py @@ -385,27 +385,14 @@ class RelayRuntime: with self._execution_consumers_lock: return bool(self._execution_consumers) - def _push_session_scope( - 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).""" - args = (self.relay.scope.push, SESSION_SCOPE, self.relay.ScopeType.Agent) - try: - future = _scope_op_executor().submit(context.run, *args, input={}, **push_kwargs) - return future.result(timeout=_SCOPE_OP_TIMEOUT) - except RuntimeError: - if not exit_fallback: - raise - return context.run(*args, input={}, **push_kwargs) - def _open_session_scope( - self, session: RelaySession, scope_metadata: dict[str, Any], *, resolve_parent: bool, **push_kwargs: Any, + self, session: RelaySession, scope_metadata: dict[str, Any], *, resolve_parent: bool, + exit_fallback: bool = False, **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.""" + """Push a fresh SESSION_SCOPE for ``session`` (bounded by ``_SCOPE_OP_TIMEOUT``); record handle + context. + Subagents parent under their spawning turn/session handle; ``resolve_parent`` creates the parent + session when its handle is unknown. ``exit_fallback``: at interpreter shutdown the executor refuses + new futures; push synchronously instead (no agent turn waits at exit).""" parent_handle = None if session.parent_session_id: with self._sessions_lock: @@ -416,7 +403,15 @@ class RelayRuntime: parent_handle = parent.handle scope_metadata["nemo_relay_scope_role"] = "subagent" context = contextvars.Context() - session.handle = self._push_session_scope(context, handle=parent_handle, metadata=scope_metadata, **push_kwargs) + args = (self.relay.scope.push, SESSION_SCOPE, self.relay.ScopeType.Agent) + push_kwargs.update(handle=parent_handle, metadata=scope_metadata, input={}) + try: + future = _scope_op_executor().submit(context.run, *args, **push_kwargs) + session.handle = future.result(timeout=_SCOPE_OP_TIMEOUT) + except RuntimeError: + if not exit_fallback: + raise + session.handle = context.run(*args, **push_kwargs) session.context = context def ensure_session( diff --git a/agent/transports/__init__.py b/agent/transports/__init__.py index e7864ce079..83f31ce9ac 100644 --- a/agent/transports/__init__.py +++ b/agent/transports/__init__.py @@ -2,6 +2,7 @@ transport = get_transport("anthropic_messages") result = transport.normalize_response(raw_response)""" +import contextlib import importlib from agent.transports.types import ( # noqa: F401 @@ -24,13 +25,10 @@ def register_transport(api_mode: str, transport_cls: type) -> None: def get_transport(api_mode: str): """Return a transport instance for ``api_mode``, or None so callers can fall back to the legacy path.""" - if not _discovered: + # A directly-imported transport leaves the registry partial; (re)discover on first use and on misses. + if not _discovered or api_mode not in _REGISTRY: _discover_transports() cls = _REGISTRY.get(api_mode) - if cls is None: - # A directly-imported transport leaves the registry partial; rediscover on misses. - _discover_transports() - cls = _REGISTRY.get(api_mode) return None if cls is None else cls() @@ -39,7 +37,5 @@ def _discover_transports() -> None: global _discovered _discovered = True for name in _TRANSPORT_MODULES: - try: + with contextlib.suppress(ImportError): importlib.import_module(f"agent.transports.{name}") - except ImportError: - pass