Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
153 changes: 141 additions & 12 deletions agent/relay_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
import logging
import threading
import uuid
from concurrent.futures import TimeoutError as FuturesTimeoutError
from dataclasses import dataclass, field
from typing import Any, Callable

Expand All @@ -25,6 +26,73 @@
RUNTIME_INSTANCE_KEY = "hermes.relay.runtime_instance"
_PROFILE_KEY_CACHE: dict[str, str] = {}

# Bound for native scope lifecycle operations (push/pop/flush) that gate
# turn/session completion. Healthy operations complete in microseconds;
# only a wedged native pipeline breaches this, and the correct trade there
# is one lost span, never a blocked agent (2026-08-10 delegation stall).
_SCOPE_OP_TIMEOUT = 10.0

_SCOPE_OP_EXECUTOR: Any = None
_SCOPE_OP_EXECUTOR_LOCK = threading.Lock()


def _scope_op_executor():
"""Shared daemon executor for bounded native scope operations.

Daemon workers (tools.daemon_pool) so a wedged native call abandoned at
timeout cannot block interpreter exit. Sized generously: workers are
only consumed for the duration of healthy (microsecond) operations plus
any wedged calls, and ``Future.result(timeout=...)`` bounds callers even
when every worker is consumed by wedged calls — an unstarted future
still honors the result timeout, so exhaustion degrades to fast
timeouts, never a new hang.
"""
global _SCOPE_OP_EXECUTOR
if _SCOPE_OP_EXECUTOR is None:
with _SCOPE_OP_EXECUTOR_LOCK:
if _SCOPE_OP_EXECUTOR is None:
from tools.daemon_pool import DaemonThreadPoolExecutor

_SCOPE_OP_EXECUTOR = DaemonThreadPoolExecutor(
max_workers=8,
thread_name_prefix="relay-scope-op",
)
return _SCOPE_OP_EXECUTOR


def _run_bounded_on_exit_thread(fn: Callable[[], Any], timeout: float) -> Any:
"""Bounded fallback lane for interpreter shutdown.

When the shared executor refuses new futures (interpreter shutdown),
the operation still must not run unbounded on the calling thread: a
wedged native call would block process exit forever — the same defect
class this module exists to prevent, on the exit lane. Run it on a
fresh daemon thread with a bounded join; on breach the daemon worker
is abandoned exactly like the executor lane abandons its worker.
"""
result: list[Any] = []
error: list[BaseException] = []

def _target() -> None:
try:
result.append(fn())
except BaseException as exc: # noqa: BLE001 - propagated below
error.append(exc)

worker = threading.Thread(
target=_target, daemon=True, name="relay-scope-op-exit"
)
worker.start()
worker.join(timeout)
if worker.is_alive():
raise TimeoutError(
f"Relay scope operation exceeded {timeout}s during interpreter "
"shutdown; abandoning the native call so process exit can proceed"
)
if error:
raise error[0]
return result[0] if result else None


@dataclass
class RelaySession:
Expand Down Expand Up @@ -113,15 +181,29 @@ def ensure_session(
scope_metadata["nemo_relay_scope_role"] = "subagent"
context = contextvars.Context()
try:
session.handle = context.run(
self.relay.scope.push,
SESSION_SCOPE,
self.relay.ScopeType.Agent,
handle=parent_handle,
data=data,
input={},
metadata=scope_metadata,
)
try:
session.handle = _scope_op_executor().submit(
context.run,
self.relay.scope.push,
SESSION_SCOPE,
self.relay.ScopeType.Agent,
handle=parent_handle,
data=data,
input={},
metadata=scope_metadata,
).result(timeout=_SCOPE_OP_TIMEOUT)
except RuntimeError:
# Interpreter shutdown: executor refuses new futures;
# push synchronously (no agent turn waits at exit).
session.handle = context.run(
self.relay.scope.push,
SESSION_SCOPE,
self.relay.ScopeType.Agent,
handle=parent_handle,
data=data,
input={},
metadata=scope_metadata,
)
except Exception:
session.context = None
raise
Expand Down Expand Up @@ -194,9 +276,22 @@ def run_in_session(
callback: Callable[..., Any],
*args: Any,
allow_closing: bool = False,
timeout: float | None = None,
**kwargs: Any,
) -> Any:
"""Run a Relay operation against a session's isolated scope stack."""
"""Run a Relay operation against a session's isolated scope stack.

``timeout`` (seconds) bounds the native call by running it on a
shared daemon executor; ``TimeoutError`` propagates to the caller's
existing exception handling on breach. ``None`` (default) preserves
the historical synchronous behavior. Scope lifecycle operations
that gate turn/session completion pass ``_SCOPE_OP_TIMEOUT``: the
native binding's ``scope.pop`` "returns after the scope is closed
successfully" — unbounded — and a wedged native pipeline (proven
live 2026-08-10 in the delegation topology) must cost at most one
span, never the agent. The abandoned daemon worker cannot block
process exit (tools.daemon_pool contract).
"""
with session.lock:
if session.closing and not allow_closing:
raise RuntimeError("Hermes Relay session is closing")
Expand All @@ -214,7 +309,27 @@ def invoke() -> Any:

# A copy permits a helper called by an existing Relay callback to
# re-enter the same logical session without re-entering Context.
return context.run(invoke)
if timeout is None:
return context.run(invoke)
try:
future = _scope_op_executor().submit(context.run, invoke)
except RuntimeError:
# Interpreter shutdown: the executor refuses new futures, but
# the atexit close path must still flush cleanly. Still
# bounded — a wedged native call must not block process exit
# (the CI runner hang, 2026-08-12: 6 tests passed in 4s, then
# the file-timeout SIGKILL'd a process stuck in this lane).
return _run_bounded_on_exit_thread(
lambda: context.run(invoke), timeout
)
try:
return future.result(timeout=timeout)
except FuturesTimeoutError as exc:
raise TimeoutError(
f"Relay scope operation exceeded {timeout}s "
f"(session={session.session_id}); abandoning the native call "
"so the agent can continue — the span for this scope is lost"
) from exc

async def run_in_session_async(
self,
Expand Down Expand Up @@ -323,11 +438,22 @@ def close_session(self, event: dict[str, Any]) -> None:
RUNTIME_INSTANCE_KEY: self.runtime_id,
},
allow_closing=True,
timeout=_SCOPE_OP_TIMEOUT,
)
except Exception as exc:
failures.append(f"session scope close failed: {exc}")
try:
self.relay.subscribers.flush()
try:
_scope_op_executor().submit(
self.relay.subscribers.flush
).result(timeout=_SCOPE_OP_TIMEOUT)
except RuntimeError:
# Interpreter shutdown: executor refuses new futures; flush
# on a bounded exit thread so a wedged pipeline cannot
# block process exit.
_run_bounded_on_exit_thread(
self.relay.subscribers.flush, _SCOPE_OP_TIMEOUT
)
except Exception as exc:
failures.append(f"subscriber flush failed: {exc}")
with self._sessions_lock:
Expand Down Expand Up @@ -661,6 +787,7 @@ def begin_turn(
RUNTIME_INSTANCE_KEY: lease.host.runtime_id,
"hermes.execution_surface": lease.platform or "unknown",
},
timeout=_SCOPE_OP_TIMEOUT,
)
except Exception:
logger.warning("Hermes Relay turn initialization failed", exc_info=True)
Expand Down Expand Up @@ -694,6 +821,7 @@ def end_turn(
RUNTIME_SCHEMA_KEY: RUNTIME_SCHEMA_VERSION,
RUNTIME_INSTANCE_KEY: lease.host.runtime_id,
},
timeout=_SCOPE_OP_TIMEOUT,
)
except Exception:
logger.warning(
Expand Down Expand Up @@ -778,6 +906,7 @@ def _finish_logical_calls(
RUNTIME_SCHEMA_KEY: RUNTIME_SCHEMA_VERSION,
RUNTIME_INSTANCE_KEY: lease.host.runtime_id,
},
timeout=_SCOPE_OP_TIMEOUT,
)
except Exception:
with turn.logical_llm_lock:
Expand Down
Loading
Loading