refactor(delegation): carry the dispatcher's Context to stale finalization
Store the dispatching contextvars.Context on the record instead of a home string, and run the monitor's forced _finalize under it. That keeps the whole scope (home override + secret scope) rather than only the home, and removes the set/reset override dance from _push_completion_event, which now stays scope-agnostic.
This commit is contained in:
@@ -9,6 +9,7 @@ crash-recovery wiring. Only the async lifecycle lives here; the child run is an
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import contextvars
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
@@ -555,9 +556,9 @@ def _dispatch(
|
||||
"slot_key": slot_key or delegation_id,
|
||||
# Which of the call's ``goals`` this unit runs (None = all of them).
|
||||
**({"task_indexes": list(task_indexes)} if task_indexes is not None else {}),
|
||||
# Durable finalization can run on the one unscoped stale-monitor thread;
|
||||
# retain the dispatching profile so that thread updates the same state.db.
|
||||
"_profile_home": str(get_hermes_home()),
|
||||
# The one stale-monitor thread serves every profile and starts with an empty Context;
|
||||
# a forced finalization runs under the dispatcher's so it settles the same state.db.
|
||||
"_context": contextvars.copy_context(),
|
||||
# Stale-monitor bookkeeping (see _stale_monitor_loop).
|
||||
"_progress_token": None, "_progress_ts": dispatched_at, "_interrupted_at": None}
|
||||
with _records_lock:
|
||||
@@ -727,19 +728,11 @@ def _push_completion_event(record: Dict[str, Any], result: Dict[str, Any], statu
|
||||
**({} if is_batch else {"exit_reason": result.get("exit_reason")}),
|
||||
**{k: record[k] for k in _ROUTING_KEYS if record.get(k)},
|
||||
**{k: result[k] for k in _STALL_META_KEYS if k in result}}
|
||||
token = None
|
||||
try:
|
||||
profile_home = record.get("_profile_home")
|
||||
if profile_home:
|
||||
from hermes_constants import reset_hermes_home_override, set_hermes_home_override
|
||||
token = set_hermes_home_override(profile_home)
|
||||
_persist_completion(evt, result)
|
||||
except Exception as exc: # noqa: BLE001 — a lost durable row is recoverable; a lost result + leaked slot is not
|
||||
logger.error(f"Async delegation{label} %s: durable completion write failed; delivering in-memory "
|
||||
"only (a restart may report this unit as unknown): %s", record.get("delegation_id"), exc)
|
||||
finally:
|
||||
if token is not None:
|
||||
reset_hermes_home_override(token)
|
||||
try:
|
||||
process_registry.completion_queue.put(evt)
|
||||
except Exception as exc: # pragma: no cover
|
||||
@@ -864,7 +857,9 @@ def _stale_monitor_loop() -> None:
|
||||
fn = (_records.get(delegation_id) or {}).get("interrupt_fn")
|
||||
_call_interrupt(fn, "Async delegation %s stall interrupt failed: %s", delegation_id)
|
||||
for delegation_id in expired:
|
||||
_finalize(delegation_id, lambda rec, d=delegation_id: _stalled_result(d, rec), "stalled")
|
||||
with _records_lock:
|
||||
ctx = (_records.get(delegation_id) or {}).get("_context") or contextvars.copy_context()
|
||||
ctx.run(_finalize, delegation_id, lambda rec, d=delegation_id: _stalled_result(d, rec), "stalled")
|
||||
if not any_monitorable:
|
||||
return
|
||||
|
||||
|
||||
Reference in New Issue
Block a user