fix(streaming): join stream worker and drain Relay scopes on interrupt

Empty-stream stalls that trip interrupt were raising InterruptedError
before the stream worker closed its physical LLM scope, corrupting the
Relay LIFO stack and cascading into a CLI EIO redraw storm (#81521).
This commit is contained in:
HexLab98
2026-08-08 15:26:46 +07:00
committed by Teknium
parent b32aa0be8f
commit 0c9d8ab0ce
3 changed files with 240 additions and 60 deletions
+148 -58
View File
@@ -543,6 +543,114 @@ class RelayRuntime:
)
return result if isinstance(result, dict) else args
def _close_scope_handle(
self,
session: RelaySession,
handle: Any,
*,
output: dict[str, Any] | None = None,
allow_closing: bool = False,
failure_label: str = "scope close failed",
drain_limit: int = 32,
) -> 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 (#81521). Drain and
close must run inside one ``run_in_session`` callback so ContextVar
stack views stay consistent across pops.
"""
if handle is None:
return None
metadata = {
RUNTIME_SCHEMA_KEY: RUNTIME_SCHEMA_VERSION,
RUNTIME_INSTANCE_KEY: self.runtime_id,
}
close_output = output or {}
session_root = session.handle
drained_holder = {"count": 0}
error_holder: dict[str, BaseException] = {}
def close_with_drain() -> None:
def current_top() -> Any:
top = self.relay.get_scope_stack()
# Some Relay builds return the live stack (list). Others
# return the top handle directly — including tuple handles
# like ("scope", name, serial) from the test fake. Only
# unwrap real list stacks; never index a handle tuple.
if isinstance(top, list):
return top[-1] if top else None
return top
try:
self.relay.scope.pop(
handle,
output=close_output,
metadata=metadata,
)
return
except Exception as first_exc:
error_holder["first"] = first_exc
for _ in range(drain_limit):
top = current_top()
if top is None or top == handle:
break
# Never pop the session root while draining for a nested handle.
if (
session_root is not None
and top == session_root
and handle is not session_root
):
break
try:
self.relay.scope.pop(
top,
output={
"outcome": "cancelled",
"hermes.orphan_drain": True,
},
metadata=metadata,
)
drained_holder["count"] += 1
except Exception as drain_exc:
error_holder["drain"] = drain_exc
logger.warning(
"Hermes Relay orphaned scope drain failed",
exc_info=True,
)
break
if drained_holder["count"]:
logger.warning(
"Hermes Relay drained %d orphaned scope(s) before closing %s",
drained_holder["count"],
handle,
)
try:
self.relay.scope.pop(
handle,
output=close_output,
metadata=metadata,
)
error_holder.pop("first", None)
error_holder.pop("drain", None)
except Exception as retry_exc:
error_holder["retry"] = retry_exc
try:
self.run_in_session(
session,
close_with_drain,
allow_closing=allow_closing,
)
except Exception as exc:
return f"{failure_label}: {exc}"
retry_exc = error_holder.get("retry") or error_holder.get("first")
if retry_exc is not None:
return f"{failure_label}: {retry_exc}"
return None
def close_session(self, event: dict[str, Any]) -> None:
"""Close one session scope and remove it from the core registry."""
session_id = _session_id(event)
@@ -559,21 +667,15 @@ class RelayRuntime:
return
session.closing = True
if session.handle is not None:
try:
self.run_in_session(
session,
self.relay.scope.pop,
session.handle,
output={},
metadata={
RUNTIME_SCHEMA_KEY: RUNTIME_SCHEMA_VERSION,
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}")
failure = self._close_scope_handle(
session,
session.handle,
output={},
allow_closing=True,
failure_label="session scope close failed",
)
if failure:
failures.append(failure)
try:
try:
_scope_op_executor().submit(
@@ -964,21 +1066,16 @@ class RelaySessionCoordinator:
if isinstance(lease.host, RelayRuntime) and lease.session is not None:
self._finish_logical_calls(turn, outcome=outcome)
if turn.handle is not None:
try:
lease.host.run_in_session(
lease.session,
lease.host.relay.scope.pop,
turn.handle,
output={"outcome": outcome},
metadata={
RUNTIME_SCHEMA_KEY: RUNTIME_SCHEMA_VERSION,
RUNTIME_INSTANCE_KEY: lease.host.runtime_id,
},
timeout=_SCOPE_OP_TIMEOUT,
)
except Exception:
failure = lease.host._close_scope_handle(
lease.session,
turn.handle,
output={"outcome": outcome},
failure_label="turn scope close failed",
)
if failure:
logger.warning(
"Hermes Relay turn finalization failed", exc_info=True
"Hermes Relay turn finalization failed: %s",
failure,
)
finally:
try:
@@ -1153,35 +1250,28 @@ class RelaySessionCoordinator:
turn.logical_llm_calls.clear()
for index in range(len(logical_calls) - 1, -1, -1):
request_id, logical_handle = logical_calls[index]
try:
lease.host.run_in_session(
lease.session,
lease.host.relay.scope.pop,
logical_handle,
output={"outcome": outcome},
metadata={
RUNTIME_SCHEMA_KEY: RUNTIME_SCHEMA_VERSION,
RUNTIME_INSTANCE_KEY: lease.host.runtime_id,
},
timeout=_SCOPE_OP_TIMEOUT,
)
except Exception:
with turn.logical_llm_lock:
# Relay scopes are stack-owned. If the newest remaining
# handle cannot close, older handles cannot close safely
# either, so retain the unclosed prefix for diagnostics.
for pending_request_id, pending_handle in logical_calls[
: index + 1
]:
turn.logical_llm_calls.setdefault(
pending_request_id,
pending_handle,
)
logger.warning(
"Hermes Relay logical LLM finalization failed",
exc_info=True,
)
break
failure = lease.host._close_scope_handle(
lease.session,
logical_handle,
output={"outcome": outcome},
failure_label="logical LLM scope close failed",
)
if failure is None:
continue
with turn.logical_llm_lock:
# Relay scopes are stack-owned. If the newest remaining
# handle cannot close even after orphan drain, older
# handles cannot close safely either — retain the
# unclosed prefix for diagnostics (#81521).
for pending_request_id, pending_handle in logical_calls[
: index + 1
]:
turn.logical_llm_calls.setdefault(
pending_request_id,
pending_handle,
)
logger.warning("Hermes Relay logical LLM finalization failed: %s", failure)
break
@staticmethod
def _reset_turn_context(turn: RelayTurnContext) -> None: