Merge simp/r3-37 (late tail) into hermes/simplify-codebase
This commit is contained in:
+667
-1114
File diff suppressed because it is too large
Load Diff
@@ -8,10 +8,10 @@ import contextlib
|
||||
|
||||
from .method_ctx import bind_module
|
||||
|
||||
# A concluded turn (success, handled error, interrupt) clears its durable marker (turn_marker.py) in
|
||||
# _run_prompt_submit's finally; only a process death leaves it behind, so a marker at session.resume
|
||||
# proves the turn never finished AND the client never saw a terminal frame. Fresh: re-submit
|
||||
# automatically (as the messaging gateway does). Stale: clear it and let the partial transcript speak.
|
||||
# A concluded turn (success, handled error, interrupt) clears its durable marker (turn_marker.py) in _run_prompt_submit's
|
||||
# finally; only a process death leaves it behind, so a marker at session.resume proves the turn never finished AND the
|
||||
# client never saw a terminal frame. Fresh: re-submit automatically (as the messaging gateway does). Stale: clear it
|
||||
# and let the partial transcript speak.
|
||||
_AUTO_CONTINUE_FRESHNESS_MINUTES_DEFAULT = 15
|
||||
|
||||
|
||||
@@ -24,11 +24,8 @@ def _auto_continue_config() -> tuple[bool, float, int]:
|
||||
minutes = float(cfg.get("freshness_minutes", _AUTO_CONTINUE_FRESHNESS_MINUTES_DEFAULT))
|
||||
except (TypeError, ValueError):
|
||||
minutes = float(_AUTO_CONTINUE_FRESHNESS_MINUTES_DEFAULT)
|
||||
return (
|
||||
is_truthy_value(cfg.get("enabled"), default=True),
|
||||
max(0.0, minutes) * 60.0,
|
||||
_coerce_int_config_value(cfg.get("max_attempts"), 2, min_value=0),
|
||||
)
|
||||
return (is_truthy_value(cfg.get("enabled"), default=True), max(0.0, minutes) * 60.0,
|
||||
_coerce_int_config_value(cfg.get("max_attempts"), 2, min_value=0))
|
||||
|
||||
|
||||
def _session_home(session: dict) -> Path:
|
||||
@@ -37,9 +34,9 @@ def _session_home(session: dict) -> Path:
|
||||
|
||||
|
||||
def _retire_turn_marker(session: dict, *keys: str) -> None:
|
||||
"""Drop the crash marker right before the terminal frame (not at turn-thread end: post-turn work
|
||||
outlives the client's answer, and quitting in that window would leave a marker that re-runs a
|
||||
finished turn). Extra ``keys`` cover a session_key that compression rotated mid-turn."""
|
||||
"""Drop the crash marker right before the terminal frame (not at turn-thread end: post-turn work outlives the
|
||||
client's answer, and quitting in that window would leave a marker that re-runs a finished turn). Extra ``keys``
|
||||
cover a session_key that compression rotated mid-turn."""
|
||||
home = _session_home(session)
|
||||
for key in dict.fromkeys((*keys, str(session.get("session_key") or ""))):
|
||||
if key:
|
||||
@@ -47,35 +44,23 @@ def _retire_turn_marker(session: dict, *keys: str) -> None:
|
||||
|
||||
|
||||
def _auto_continue_note(prompt: str) -> str:
|
||||
# Same opening as the gateway's recovery notes (transcript tooling recognizes both). The prompt is
|
||||
# embedded: a hard crash persists nothing else of the turn.
|
||||
return (
|
||||
f"{_AUTO_CONTINUE_NOTE_PREFIX} — the app or its backend process "
|
||||
"stopped before the turn could finish. Some of the work may already "
|
||||
"be complete; check the current state before redoing anything, then "
|
||||
"finish the task. The interrupted request was:]\n\n"
|
||||
f"{prompt}"
|
||||
)
|
||||
|
||||
|
||||
def _ac_release_turn(session: dict, *, unschedule: bool = False) -> None:
|
||||
with session["history_lock"]:
|
||||
session["running"] = False
|
||||
if unschedule:
|
||||
session["_auto_continue_scheduled"] = False
|
||||
# Same opening as the gateway's recovery notes (transcript tooling recognizes both). The prompt is embedded: a hard
|
||||
# crash persists nothing else of the turn.
|
||||
return (f"{_AUTO_CONTINUE_NOTE_PREFIX} — the app or its backend process stopped before the turn could finish. "
|
||||
"Some of the work may already be complete; check the current state before redoing anything, then "
|
||||
f"finish the task. The interrupted request was:]\n\n{prompt}")
|
||||
|
||||
|
||||
def _maybe_schedule_auto_continue(sid: str, session: dict, session_key: str) -> dict | None:
|
||||
"""Kick off a continuation turn for a crash-interrupted session (session.resume cold paths). Returns a
|
||||
descriptor for the resume payload when scheduled, else None. The turn runs on a background thread
|
||||
after the deferred agent build via _run_prompt_submit, so the client that just resumed streams it."""
|
||||
# Hosted room turns are recovered by their durable task/lease state machine; generic auto-continue
|
||||
# would bypass its execution generation and duplicate work.
|
||||
"""Kick off a continuation turn for a crash-interrupted session (session.resume cold paths). Returns a descriptor
|
||||
for the resume payload when scheduled, else None. The turn runs on a background thread after the deferred agent
|
||||
build via _run_prompt_submit, so the client that just resumed streams it."""
|
||||
# Hosted room turns are recovered by their durable task/lease state machine; generic auto-continue would bypass
|
||||
# its execution generation and duplicate work.
|
||||
if session.get("source") == "bot_room":
|
||||
return None
|
||||
home = _session_home(session)
|
||||
marker = read_turn_marker(home, session_key)
|
||||
if marker is None:
|
||||
if (marker := read_turn_marker(home, session_key)) is None:
|
||||
return None
|
||||
enabled, freshness_secs, max_attempts = _auto_continue_config()
|
||||
age = time.time() - marker["started_at"]
|
||||
@@ -104,16 +89,17 @@ def _maybe_schedule_auto_continue(sid: str, session: dict, session_key: str) ->
|
||||
return
|
||||
session["running"] = True
|
||||
session["last_active"] = time.time()
|
||||
# Ownership admission BEFORE message.start: a sibling backend sharing this HERMES_HOME may have
|
||||
# written the marker and still be mid-turn. Leave the marker so a later resume retries.
|
||||
# Ownership admission BEFORE message.start: a sibling backend sharing this HERMES_HOME may have written the
|
||||
# marker and still be mid-turn. Leave the marker so a later resume retries.
|
||||
if _ensure_active_session_slot(sid, session) is not None:
|
||||
logger.info("auto-continue for %s refused: session has another live owner", session_key)
|
||||
_ac_release_turn(session, unschedule=True)
|
||||
with session["history_lock"]:
|
||||
session["running"] = False
|
||||
session["_auto_continue_scheduled"] = False
|
||||
return
|
||||
with session["history_lock"]:
|
||||
# Marker inputs read back by _run_prompt_submit: attempt count (crash breaker) and the ORIGINAL
|
||||
# prompt (no nested notes). Set here, not at schedule time, so a bail above leaves nothing
|
||||
# for a racing user turn.
|
||||
# Marker inputs read back by _run_prompt_submit: attempt count (crash breaker) and the ORIGINAL prompt (no
|
||||
# nested notes). Set here, not at schedule time, so a bail above leaves nothing for a racing user turn.
|
||||
session["_auto_continue_attempt"], session["_auto_continue_prompt"] = attempt, marker["prompt"]
|
||||
try:
|
||||
_emit("status.update", sid, {"kind": "process", "text": "Resuming interrupted turn…"})
|
||||
@@ -121,7 +107,7 @@ def _maybe_schedule_auto_continue(sid: str, session: dict, session_key: str) ->
|
||||
_run_prompt_submit(rid, sid, session, text, display_kind="auto_continue")
|
||||
except Exception as exc:
|
||||
_notif_log_failure("auto-continue dispatch failed", exc)
|
||||
_ac_release_turn(session)
|
||||
_notif_release_turn(session) # rebound from session_notifications
|
||||
threading.Thread(target=kickoff, daemon=True).start()
|
||||
logger.info("auto-continue scheduled for session %s (attempt %d, interrupted %.0fs ago)", session_key, attempt, age)
|
||||
return {"attempt": attempt, "interrupted_at": marker["started_at"]}
|
||||
@@ -134,11 +120,11 @@ def _ac_inflight_original(session: dict) -> str:
|
||||
|
||||
def _enqueue_prompt(session: dict, text: Any, transport: Any, image_paths: list[str] | None = None) -> None:
|
||||
"""Queue a message for the next turn. Text-only arrivals share a slot and merge losslessly (like the
|
||||
consecutive-user merge in ``repair_message_sequence``); image-bearing ones stay separate envelopes so
|
||||
attachment chronology survives. ``transport`` is pinned so the drained turn streams to its sender."""
|
||||
consecutive-user merge in ``repair_message_sequence``); image-bearing ones stay separate envelopes so attachment
|
||||
chronology survives. ``transport`` is pinned so the drained turn streams to its sender."""
|
||||
image_paths = list(image_paths or [])
|
||||
# Scrub live-turn self-duplicates first so the text merge below can't glue "{original}\n\n{later}"
|
||||
# and re-fire the original after a correction settles.
|
||||
# Scrub live-turn self-duplicates first so the text merge below can't glue "{original}\n\n{later}" and re-fire the
|
||||
# original after a correction settles.
|
||||
_drop_queued_duplicates_of_inflight_user(session)
|
||||
text_only = not image_paths and isinstance(text, str)
|
||||
# Never queue a text-only self-copy of the live prompt: draining it would restart it.
|
||||
@@ -146,10 +132,8 @@ def _enqueue_prompt(session: dict, text: Any, transport: Any, image_paths: list[
|
||||
return
|
||||
queued = {"text": text, "transport": transport, **({"image_paths": image_paths} if image_paths else {})}
|
||||
existing = session.get("queued_prompt")
|
||||
if (
|
||||
existing and text_only and isinstance(existing.get("text"), str)
|
||||
and not existing.get("image_paths") and not session.get("queued_prompts")
|
||||
):
|
||||
if (existing and text_only and isinstance(existing.get("text"), str)
|
||||
and not existing.get("image_paths") and not session.get("queued_prompts")):
|
||||
prev = existing["text"]
|
||||
existing["text"] = f"{prev}\n\n{text}" if prev and text else (prev or text)
|
||||
elif existing:
|
||||
@@ -160,8 +144,8 @@ def _enqueue_prompt(session: dict, text: Any, transport: Any, image_paths: list[
|
||||
|
||||
def _sanitize_queued_entry_vs_inflight_user(entry: Any, original: str) -> dict | None:
|
||||
"""Drop (``None``) a text-only self-duplicate of the live user text, or rewrite a merged slot
|
||||
``"{original}\\n\\n{later}"`` to ``later`` so the correction survives without re-firing the original.
|
||||
Image-bearing envelopes are left alone (chronology is load-bearing)."""
|
||||
``"{original}\\n\\n{later}"`` to ``later`` so the correction survives without re-firing the original. Image-bearing
|
||||
envelopes are left alone (chronology is load-bearing)."""
|
||||
if not isinstance(entry, dict):
|
||||
return None
|
||||
text = entry.get("text")
|
||||
@@ -173,16 +157,13 @@ def _sanitize_queued_entry_vs_inflight_user(entry: Any, original: str) -> dict |
|
||||
|
||||
|
||||
def _drop_queued_duplicates_of_inflight_user(session: dict) -> None:
|
||||
"""Remove server-queue copies of the live turn's original user text: a mid-turn ``prompt.submit`` of
|
||||
the same text queued while redirect was unavailable must not drain and restart the original."""
|
||||
original = _ac_inflight_original(session)
|
||||
if not original:
|
||||
"""Remove server-queue copies of the live turn's original user text: a mid-turn ``prompt.submit`` of the same text
|
||||
queued while redirect was unavailable must not drain and restart the original."""
|
||||
if not (original := _ac_inflight_original(session)):
|
||||
return
|
||||
head = session.get("queued_prompt")
|
||||
cleaned = (
|
||||
_sanitize_queued_entry_vs_inflight_user(e, original)
|
||||
for e in ([head] if head else []) + list(session.get("queued_prompts") or [])
|
||||
)
|
||||
cleaned = (_sanitize_queued_entry_vs_inflight_user(e, original)
|
||||
for e in ([head] if head else []) + list(session.get("queued_prompts") or []))
|
||||
_ac_set_queue(session, [c for c in cleaned if c is not None])
|
||||
|
||||
|
||||
@@ -196,9 +177,9 @@ def _ac_set_queue(session: dict, entries: list) -> None:
|
||||
|
||||
|
||||
def _interrupt_busy_session(sid: str, session: dict, agent: Any) -> None:
|
||||
"""Interrupt a busy turn on a worker thread, never under ``history_lock`` (some providers can't apply
|
||||
``interrupt()`` until a blocking call returns; inline it stalled ``session.resume``). At most one
|
||||
interrupt worker per session so repeated steering can't leak threads."""
|
||||
"""Interrupt a busy turn on a worker thread, never under ``history_lock`` (some providers can't apply ``interrupt()``
|
||||
until a blocking call returns; inline it stalled ``session.resume``). At most one interrupt worker per session so
|
||||
repeated steering can't leak threads."""
|
||||
use_agent = agent is not None and hasattr(agent, "interrupt")
|
||||
if not use_agent and not _session_uses_compute_host(session):
|
||||
return
|
||||
@@ -209,12 +190,8 @@ def _interrupt_busy_session(sid: str, session: dict, agent: Any) -> None:
|
||||
|
||||
def interrupt() -> None:
|
||||
try:
|
||||
if use_agent:
|
||||
agent.interrupt()
|
||||
else:
|
||||
_get_compute_host_supervisor().interrupt(sid)
|
||||
except Exception:
|
||||
pass
|
||||
with contextlib.suppress(Exception):
|
||||
agent.interrupt() if use_agent else _get_compute_host_supervisor().interrupt(sid)
|
||||
finally:
|
||||
with session["history_lock"]:
|
||||
session["_busy_interrupt_pending"] = False
|
||||
@@ -222,9 +199,9 @@ def _interrupt_busy_session(sid: str, session: dict, agent: Any) -> None:
|
||||
|
||||
|
||||
def _ac_try_correction(rid, session: dict, agent: Any, method: str, plain_text: str, status: str) -> dict | None:
|
||||
"""Apply ``agent.<method>(plain_text)`` (steer/redirect); on acceptance record the correction, scrub
|
||||
stale self-duplicates so the live turn's original text is not re-fired after settle, and return the
|
||||
``status`` reply. None → caller falls through to the queue path."""
|
||||
"""Apply ``agent.<method>(plain_text)`` (steer/redirect); on acceptance record the correction, scrub stale
|
||||
self-duplicates so the live turn's original text is not re-fired after settle, and return the ``status`` reply.
|
||||
None → caller falls through to the queue path."""
|
||||
try:
|
||||
if not getattr(agent, method)(plain_text):
|
||||
return None
|
||||
@@ -238,10 +215,10 @@ def _ac_try_correction(rid, session: dict, agent: Any, method: str, plain_text:
|
||||
|
||||
|
||||
def _handle_busy_submit(rid, sid: str, session: dict, text: Any, transport: Any, queued: bool = False) -> dict | None:
|
||||
"""Apply ``display.busy_input_mode`` to a mid-turn prompt instead of rejecting it (rejection made
|
||||
clients busy-retry and drop sends): ``interrupt`` (default) → redirect, falling back to hard interrupt
|
||||
+ queue; ``queue`` → queue only; ``steer`` → inject after the current atomic action. ``queued=True``
|
||||
(client queue drain) forces queue mode: a "run after" message must NEVER become a live correction."""
|
||||
"""Apply ``display.busy_input_mode`` to a mid-turn prompt instead of rejecting it (rejection made clients busy-retry
|
||||
and drop sends): ``interrupt`` (default) → redirect, falling back to hard interrupt + queue; ``queue`` → queue only;
|
||||
``steer`` → inject after the current atomic action. ``queued=True`` (client queue drain) forces queue mode: a "run
|
||||
after" message must NEVER become a live correction."""
|
||||
mode = "queue" if queued else _load_busy_input_mode()
|
||||
agent = session.get("agent")
|
||||
with session["history_lock"]:
|
||||
@@ -250,22 +227,19 @@ def _handle_busy_submit(rid, sid: str, session: dict, text: Any, transport: Any,
|
||||
image_paths = list(session.get("attached_images", []))
|
||||
if image_paths:
|
||||
session["attached_images"] = [] # claim now so a later paste isn't consumed when the turn yields
|
||||
text_only = not image_paths and _is_text_only_busy_payload(text)
|
||||
plain_text = _coerce_message_text(text).strip() if text_only else ""
|
||||
# Text-only corrections steer/redirect in place when supported; media payloads and older agents fall
|
||||
# through to the proven interrupt + queue path.
|
||||
plain_text = _coerce_message_text(text).strip() if not image_paths and _is_text_only_busy_payload(text) else ""
|
||||
# Text-only corrections steer/redirect in place when supported; media payloads and older agents fall through to
|
||||
# the proven interrupt + queue path.
|
||||
if plain_text and agent is not None:
|
||||
supported = {
|
||||
"steer": hasattr(agent, "steer"),
|
||||
"interrupt": getattr(agent, "_supports_active_turn_redirect", False) is True and hasattr(agent, "redirect"),
|
||||
}
|
||||
"interrupt": getattr(agent, "_supports_active_turn_redirect", False) is True and hasattr(agent, "redirect")}
|
||||
method, status = {"steer": ("steer", "steered"), "interrupt": ("redirect", "redirected")}.get(mode, (None, None))
|
||||
if method and supported[mode]:
|
||||
resp = _ac_try_correction(rid, session, agent, method, plain_text, status)
|
||||
if resp is not None:
|
||||
return resp
|
||||
# Queue before asking the live turn to stop. Never call a provider/compute-host method under
|
||||
# history_lock: an interrupt can wait behind the op it cancels.
|
||||
if (method and supported[mode]
|
||||
and (resp := _ac_try_correction(rid, session, agent, method, plain_text, status)) is not None):
|
||||
return resp
|
||||
# Queue before asking the live turn to stop. Never call a provider/compute-host method under history_lock: an
|
||||
# interrupt can wait behind the op it cancels.
|
||||
with session["history_lock"]:
|
||||
if not session.get("running"):
|
||||
if image_paths:
|
||||
@@ -273,20 +247,19 @@ def _handle_busy_submit(rid, sid: str, session: dict, text: Any, transport: Any,
|
||||
return None
|
||||
_enqueue_prompt(session, text, transport, image_paths=image_paths)
|
||||
session["last_active"] = time.time()
|
||||
# Attachments need their own model invocation: queue without cancelling so the user gets both results
|
||||
# in order. ``steer`` must NEVER escalate to a hard interrupt: it would kill the live turn AND drop
|
||||
# ``AIAgent._pending_steer``, destroying earlier accepted steers; steer fall-throughs stay FIFO-queued.
|
||||
# Attachments need their own model invocation: queue without cancelling so the user gets both results in order.
|
||||
# ``steer`` must NEVER escalate to a hard interrupt: it would kill the live turn AND drop ``AIAgent._pending_steer``
|
||||
# (earlier accepted steers); steer fall-throughs stay FIFO-queued.
|
||||
if mode == "interrupt" and not image_paths:
|
||||
_interrupt_busy_session(sid, session, agent)
|
||||
return _ok(rid, {"status": "queued"})
|
||||
|
||||
|
||||
def _drain_queued_prompt(rid, sid: str, session: dict) -> bool:
|
||||
"""Fire a queued next-turn prompt if one is waiting and the session is idle. True when dispatched: the
|
||||
caller skips lower-priority follow-ups this cycle (the user's message wins)."""
|
||||
"""Fire a queued next-turn prompt if one is waiting and the session is idle. True when dispatched: the caller
|
||||
skips lower-priority follow-ups this cycle (the user's message wins)."""
|
||||
with session["history_lock"]:
|
||||
queued = session.get("queued_prompt")
|
||||
if session.get("_closing") or not queued or session.get("running"):
|
||||
if session.get("_closing") or not (queued := session.get("queued_prompt")) or session.get("running"):
|
||||
return False
|
||||
queue_generation = int(session.get("_queued_prompt_generation", 0))
|
||||
_ac_set_queue(session, session.get("queued_prompts") or [])
|
||||
@@ -296,9 +269,8 @@ def _drain_queued_prompt(rid, sid: str, session: dict) -> bool:
|
||||
use_compute_host = _session_uses_compute_host(session)
|
||||
with session["history_lock"]:
|
||||
if int(session.get("_queued_prompt_generation", 0)) != queue_generation:
|
||||
# Generation bump cancelled the claim (Stop, compress re-anchor, …): don't dispatch, but
|
||||
# restore the envelope (claimed head first, then whatever advanced into the slot) so a
|
||||
# legitimate follow-up isn't dropped.
|
||||
# Generation bump cancelled the claim (Stop, compress re-anchor, …): don't dispatch, but restore the
|
||||
# envelope (claimed head first, then whatever advanced into the slot) so a legitimate follow-up isn't dropped.
|
||||
advanced = session.get("queued_prompt")
|
||||
_ac_set_queue(session, [queued, *([advanced] if advanced else []), *(session.get("queued_prompts") or [])])
|
||||
session["running"] = False
|
||||
@@ -318,7 +290,7 @@ def _drain_queued_prompt(rid, sid: str, session: dict) -> bool:
|
||||
dispatch_failed = True
|
||||
except Exception as exc:
|
||||
_notif_log_failure("queued prompt dispatch failed", exc)
|
||||
_ac_release_turn(session)
|
||||
_notif_release_turn(session)
|
||||
dispatch_failed = True
|
||||
if dispatch_failed:
|
||||
with session["history_lock"]:
|
||||
@@ -334,67 +306,58 @@ def _inflight_snapshot(session: dict) -> dict | None:
|
||||
return None
|
||||
user, assistant = str(turn.get("user") or "").strip(), str(turn.get("assistant") or "")
|
||||
streaming, error = bool(turn.get("streaming")), str(turn.get("error") or "").strip()
|
||||
if not user and not assistant and not streaming and not error:
|
||||
if not (user or assistant or streaming or error):
|
||||
return None
|
||||
snapshot = {"assistant": assistant, "streaming": streaming, "user": user}
|
||||
raw_offsets = turn.get("correction_offsets") or []
|
||||
correction_pairs = [
|
||||
(str(c), raw_offsets[i] if i < len(raw_offsets) else None)
|
||||
for i, c in enumerate(turn.get("corrections") or []) if str(c).strip()
|
||||
]
|
||||
correction_pairs = [(str(c), raw_offsets[i] if i < len(raw_offsets) else None)
|
||||
for i, c in enumerate(turn.get("corrections") or []) if str(c).strip()]
|
||||
if correction_pairs:
|
||||
# Mid-turn redirects alongside (not over) the original prompt so resume can rebuild every user
|
||||
# bubble; offsets only when every correction has one so clients can trust the pairing.
|
||||
# Mid-turn redirects alongside (not over) the original prompt so resume can rebuild every user bubble; offsets
|
||||
# only when every correction has one so clients can trust the pairing.
|
||||
snapshot["corrections"] = [c for c, _ in correction_pairs]
|
||||
if all(isinstance(offset, int) and offset >= 0 for _, offset in correction_pairs):
|
||||
snapshot["correction_offsets"] = [int(offset) for _, offset in correction_pairs] # type: ignore[arg-type]
|
||||
if error:
|
||||
# Retained failed turn (_fail_inflight_turn): a resuming client must rebuild the failed bubble,
|
||||
# not render the partial text as a healthy reply.
|
||||
# Retained failed turn (_fail_inflight_turn): a resuming client must rebuild the failed bubble, not render the
|
||||
# partial text as a healthy reply.
|
||||
snapshot.update(error=error, status=str(turn.get("status") or "error"), recoverable=bool(turn.get("recoverable")))
|
||||
surface = turn.get("error_surface")
|
||||
if isinstance(surface, dict) and surface:
|
||||
if isinstance(surface := turn.get("error_surface"), dict) and surface:
|
||||
snapshot["error_surface"] = surface
|
||||
return snapshot
|
||||
|
||||
|
||||
def _emit_terminal_turn_error(
|
||||
sid: str, session: dict, error: Any, error_surface: Optional[dict] = None, *, retire_marker: bool = True
|
||||
) -> None:
|
||||
"""Close a failed turn with the same ``status: "error"`` ``message.complete`` frame as the returned-error
|
||||
path, retaining the turn so a client that missed the frame recovers it from ``session.resume``'s
|
||||
``inflight``. ``error_surface`` ({layer, code, retryable}) is classified from an exception if absent."""
|
||||
sid: str, session: dict, error: Any, error_surface: Optional[dict] = None, *, retire_marker: bool = True) -> None:
|
||||
"""Close a failed turn with the same ``status: "error"`` ``message.complete`` frame as the returned-error path,
|
||||
retaining the turn so a client that missed the frame recovers it from ``session.resume``'s ``inflight``.
|
||||
``error_surface`` ({layer, code, retryable}) is classified from an exception if absent."""
|
||||
agent = session.get("agent")
|
||||
if error_surface is None and isinstance(error, BaseException):
|
||||
with contextlib.suppress(Exception):
|
||||
from agent.error_surface import build_error_surface_from_exception
|
||||
error_surface = build_error_surface_from_exception(
|
||||
error, provider=str(getattr(agent, "provider", "") or ""), model=str(getattr(agent, "model", "") or "")
|
||||
)
|
||||
error, provider=str(getattr(agent, "provider", "") or ""), model=str(getattr(agent, "model", "") or ""))
|
||||
with session["history_lock"]:
|
||||
_fail_inflight_turn(session, error, error_surface=error_surface)
|
||||
turn = session.get("inflight_turn") or {}
|
||||
message, partial = str(turn.get("error") or "turn failed"), str(turn.get("assistant") or "")
|
||||
cols = int(session.get("cols", 80))
|
||||
text = partial or f"Error: {message}"
|
||||
try:
|
||||
rendered = ""
|
||||
with contextlib.suppress(Exception):
|
||||
rendered = render_message(text, cols)
|
||||
except Exception:
|
||||
rendered = ""
|
||||
payload = {
|
||||
"text": text, "usage": _get_usage(agent) if agent is not None else {}, "status": "error",
|
||||
"error": message, "recoverable": True,
|
||||
**({"error_surface": error_surface} if error_surface else {}),
|
||||
**({"partial": True} if partial else {}), **({"rendered": rendered} if rendered else {}),
|
||||
}
|
||||
payload = {"text": text, "usage": _get_usage(agent) if agent is not None else {}, "status": "error",
|
||||
"error": message, "recoverable": True, **({"error_surface": error_surface} if error_surface else {}),
|
||||
**({"partial": True} if partial else {}), **({"rendered": rendered} if rendered else {})}
|
||||
if retire_marker:
|
||||
_retire_turn_marker(session)
|
||||
_emit("message.complete", sid, payload)
|
||||
|
||||
|
||||
def _restore_agent_history_after_turn_error(session: dict, agent) -> bool:
|
||||
"""Keep a failed turn's working transcript: ``AIAgent`` persists its messages independently, so after
|
||||
a raise the next prompt must see them, not the pre-turn snapshot."""
|
||||
"""Keep a failed turn's working transcript: ``AIAgent`` persists its messages independently, so after a raise the
|
||||
next prompt must see them, not the pre-turn snapshot."""
|
||||
agent_messages = getattr(agent, "_session_messages", None)
|
||||
if not isinstance(agent_messages, list):
|
||||
return False
|
||||
@@ -405,8 +368,8 @@ def _restore_agent_history_after_turn_error(session: dict, agent) -> bool:
|
||||
|
||||
|
||||
def _queued_prompt_snapshot(session: dict) -> dict | None:
|
||||
"""The accepted next-turn prompt without its transport handle, for the live-session projection
|
||||
(Desktop may reconnect while it is still queued)."""
|
||||
"""The accepted next-turn prompt without its transport handle, for the live-session projection (Desktop may
|
||||
reconnect while it is still queued)."""
|
||||
queued = session.get("queued_prompt")
|
||||
user = _inflight_text(queued.get("text")) if isinstance(queued, dict) else ""
|
||||
return {"user": user} if user else None
|
||||
|
||||
@@ -162,11 +162,9 @@ def _apply_pending_model_switch(sid: str, session: dict) -> None:
|
||||
if not pending or session.get("agent") is None:
|
||||
return
|
||||
try:
|
||||
result = _apply_model_switch(
|
||||
sid, session, pending["raw"], confirm_expensive_model=bool(pending.get("confirm_expensive_model"))
|
||||
)
|
||||
# Honour the expensive-model confirm: surface the warning and drop the switch rather than
|
||||
# spend on a model the user never confirmed.
|
||||
result = _apply_model_switch(sid, session, pending["raw"], confirm_expensive_model=bool(pending.get("confirm_expensive_model")))
|
||||
# Honour the expensive-model confirm: surface the warning and drop the switch rather than spend
|
||||
# on a model the user never confirmed.
|
||||
if result.get("confirm_required"):
|
||||
_emit("error", sid, {"message": result.get("confirm_message") or result.get("warning") or ""})
|
||||
except Exception as e:
|
||||
|
||||
@@ -35,8 +35,7 @@ def _claim_active_session_slot(
|
||||
track_liveness=str(surface or "").strip().lower() == "desktop")
|
||||
except Exception as exc:
|
||||
logger.warning("Failed to claim active session slot: %s", exc)
|
||||
# Fail CLOSED regardless of surface: an errored claim has NOT proven the session unowned, and
|
||||
# proceeding lease-less is a silent double-writer hole.
|
||||
# Fail CLOSED: an errored claim has NOT proven the session unowned; lease-less = silent double-writer hole.
|
||||
return (None, _SESSION_OWNERSHIP_UNAVAILABLE)
|
||||
|
||||
|
||||
@@ -55,8 +54,8 @@ def _ensure_active_session_slot(sid: str, session: dict) -> str | None:
|
||||
|
||||
|
||||
def _lease_retry(attempts: int, fn) -> Exception | None:
|
||||
"""Call ``fn`` up to ``attempts`` times (50ms*n backoff, registry writes contend across processes); the
|
||||
last exception when every try failed, else None."""
|
||||
"""Call ``fn`` up to ``attempts`` times (50ms*n backoff: registry writes contend across processes); last
|
||||
exception when every try failed, else None."""
|
||||
for attempt in range(attempts):
|
||||
try:
|
||||
fn()
|
||||
@@ -92,9 +91,8 @@ def _own_live_lease_ids(*, exclude=None) -> set[str]:
|
||||
|
||||
@contextlib.contextmanager
|
||||
def _other_runtime_lease_guard(session_id: str, session: dict):
|
||||
"""Release this runtime and lock sibling ownership through the DB write. Yields True (another runtime
|
||||
owns the lifecycle -> preserve) when the guard cannot be loaded or entered after 3 tries: an unknown
|
||||
ownership state must never end a row."""
|
||||
"""Release this runtime and lock sibling ownership through the DB write. Yields True (another runtime owns
|
||||
the lifecycle -> preserve) when the guard can't be loaded/entered in 3 tries: unknown ownership never ends a row."""
|
||||
lease = session.get("active_session_lease")
|
||||
try:
|
||||
from hermes_cli.active_sessions import active_session_liveness_guard, release_active_session_liveness_guard
|
||||
@@ -142,17 +140,14 @@ def _transfer_active_session_slot(sid: str, session: dict, *, new_session_id: st
|
||||
logger.debug("Failed to transfer active session slot", exc_info=True)
|
||||
if getattr(lease, "track_liveness", False):
|
||||
return False
|
||||
# Fallback (entry pruned / pid-check transiently failed): reserve the new slot BEFORE releasing the old one
|
||||
# so a gateway at the cap can't grab the freed slot and leave this session lease-less. On reserve failure
|
||||
# KEEP the old lease.
|
||||
# Fallback (entry pruned / pid-check transiently failed): reserve the new slot BEFORE releasing the old one so
|
||||
# a gateway at the cap can't grab the freed slot and leave this session lease-less; on failure KEEP the old lease.
|
||||
new_lease, limit_message = _claim_active_session_slot(
|
||||
new_session_id, live_session_id=sid, surface=_session_source(session), profile_home=session.get("profile_home"))
|
||||
if new_lease is None:
|
||||
if limit_message:
|
||||
logger.warning(
|
||||
"Compression session lease re-anchor failed (kept old lease): "
|
||||
"sid=%s new_session_id=%s reason=%s",
|
||||
sid, new_session_id, limit_message)
|
||||
logger.warning("Compression session lease re-anchor failed (kept old lease): sid=%s new_session_id=%s reason=%s",
|
||||
sid, new_session_id, limit_message)
|
||||
return False
|
||||
if (old := session.pop("active_session_lease", None)) is not None and (err := _lease_retry(1, old.release)):
|
||||
logger.debug("Failed to release stale active session slot", exc_info=err)
|
||||
@@ -160,18 +155,16 @@ def _transfer_active_session_slot(sid: str, session: dict, *, new_session_id: st
|
||||
return True
|
||||
|
||||
|
||||
# Sources this backend must never end in state.db: the messaging gateway owns those sessions and the TUI is
|
||||
# only a viewer (ending one causes the Groundhog Day routing loop, see _finalize_session). Self-created and
|
||||
# CLI sources are NOT gateway-owned.
|
||||
# Sources this backend must never end in state.db: the messaging gateway owns those sessions and the TUI is only
|
||||
# a viewer (ending one causes the Groundhog Day loop, see _finalize_session). Self-created/CLI sources are NOT gateway-owned.
|
||||
_NON_GATEWAY_SOURCES = frozenset({
|
||||
"", "tui", "cli", "webui", "desktop", "cron", "kanban", "subagent", "test",
|
||||
"local", "acp", "webhook", "api_server", "msgraph_webhook"})
|
||||
|
||||
|
||||
def _is_gateway_owned_source(source: str) -> bool:
|
||||
"""True when ``source`` resolves to a gateway ``Platform`` (enum member or plugin via ``Platform._missing_``),
|
||||
so new platforms are covered automatically; self-owned sources that are Platform members
|
||||
(``local``/``webhook``/``api_server``) are excluded explicitly."""
|
||||
"""True when ``source`` resolves to a gateway ``Platform`` (enum member or plugin via ``Platform._missing_``, so
|
||||
new platforms are covered automatically); self-owned Platform members (local/webhook/api_server) are excluded."""
|
||||
src = (source or "").strip().lower()
|
||||
if src in _NON_GATEWAY_SOURCES:
|
||||
return False
|
||||
@@ -198,26 +191,22 @@ def _finalize_session(session: dict | None, end_reason: str = "tui_close") -> No
|
||||
if not session or session.get("_finalized"):
|
||||
return
|
||||
session["_finalized"] = True
|
||||
history_ready = session.get("resume_history_ready")
|
||||
if history_ready is not None and not history_ready.is_set():
|
||||
if (history_ready := session.get("resume_history_ready")) is not None and not history_ready.is_set():
|
||||
session["resume_history_error"] = "session resume cancelled"
|
||||
history_ready.set()
|
||||
_desktop_automatic_cleanup = (
|
||||
end_reason in _AUTOMATIC_SESSION_END_REASONS and _session_source(session).strip().lower() == "desktop")
|
||||
# Automatic Desktop cleanup releases its lease inside the lock-held lifecycle guard below; explicit close
|
||||
# and non-Desktop paths keep force/end semantics.
|
||||
# Automatic Desktop cleanup releases its lease inside the lifecycle guard below; other paths keep force/end semantics.
|
||||
if not _desktop_automatic_cleanup:
|
||||
_release_active_session_slot(session)
|
||||
if (stop_event := session.get("_notif_stop")) is not None:
|
||||
stop_event.set()
|
||||
agent = session.get("agent")
|
||||
lock = session.get("history_lock")
|
||||
with (lock if lock is not None else contextlib.nullcontext()):
|
||||
with (session.get("history_lock") or contextlib.nullcontext()):
|
||||
history = list(session.get("history", []))
|
||||
# Persist via ``_persist_session``'s marker-based dedup (same contract as the gateway-shutdown flush). Do
|
||||
# NOT pass ``conversation_history``: ``session["history"]`` and ``_session_messages`` alias the SAME list
|
||||
# after a turn, so the flush would treat every message as durable and skip it — data loss when finalize
|
||||
# is the sole persist path after a WS disconnect/restart.
|
||||
# Persist via ``_persist_session``'s marker-based dedup (gateway-shutdown flush contract). Do NOT pass
|
||||
# ``conversation_history``: ``session["history"]`` and ``_session_messages`` alias the SAME list after a turn, so
|
||||
# the flush would treat every message as durable and skip it — data loss when finalize is the sole persist path.
|
||||
if hasattr(agent, "_persist_session") and (snapshot := getattr(agent, "_session_messages", None)):
|
||||
with contextlib.suppress(Exception):
|
||||
agent._persist_session(snapshot)
|
||||
@@ -237,53 +226,46 @@ def _finalize_session(session: dict | None, end_reason: str = "tui_close") -> No
|
||||
session_id = getattr(agent, "session_id", None) or session_key
|
||||
_notify_session_boundary("on_session_finalize", session_id, _session_source(session))
|
||||
# End the state.db row so it doesn't linger as a ghost in /resume. Use session_id (agent.session_id), not
|
||||
# session_key: after compression the key may be the stale ended parent while session_id is the live
|
||||
# continuation.
|
||||
# session_key: after compression the key may be the stale ended parent while session_id is the live continuation.
|
||||
if _desktop_automatic_cleanup and not session_id:
|
||||
_release_active_session_slot(session)
|
||||
_lifecycle_guard = (
|
||||
_other_runtime_lease_guard(session_id, session)
|
||||
if _desktop_automatic_cleanup and session_id else contextlib.nullcontext(False))
|
||||
_lifecycle_guard = (_other_runtime_lease_guard(session_id, session)
|
||||
if _desktop_automatic_cleanup and session_id else contextlib.nullcontext(False))
|
||||
with _lifecycle_guard as _other_runtime_owns_lifecycle:
|
||||
_tui_owns_lifecycle = not _other_runtime_owns_lifecycle
|
||||
if _other_runtime_owns_lifecycle:
|
||||
logger.info(
|
||||
"Preserving session %s during %s: another backend owns an active lease", session_id, end_reason)
|
||||
logger.info("Preserving session %s during %s: another backend owns an active lease", session_id, end_reason)
|
||||
if session_id:
|
||||
# The *session's* profile state.db (app-global remote mode), not the launch profile's.
|
||||
with contextlib.suppress(Exception), _session_db(session) as db:
|
||||
if db is not None:
|
||||
# Never end gateway-originated sessions (Groundhog Day loop: the gateway's self-heal
|
||||
# recovers to the parent, compression splits back to the reaped child, repeat forever).
|
||||
# Never end gateway-originated sessions: Groundhog Day loop (gateway self-heals to the parent,
|
||||
# compression splits back to the reaped child, forever).
|
||||
if _is_gateway_owned_source((db.get_session(session_id) or {}).get("source", "")):
|
||||
_tui_owns_lifecycle = False
|
||||
elif _tui_owns_lifecycle:
|
||||
db.end_session(session_id, end_reason)
|
||||
# In-flight async delegations end WITH the session (no return address left). Always interrupt by THIS
|
||||
# live UI sid; by durable session_key only when the TUI owns the lifecycle — closing a viewer tab must
|
||||
# not kill the gateway's own background work.
|
||||
# In-flight async delegations end WITH the session (no return address left). Always interrupt by THIS live UI
|
||||
# sid; by durable session_key only when the TUI owns the lifecycle — a viewer tab must not kill gateway work.
|
||||
with contextlib.suppress(Exception):
|
||||
from tools.async_delegation import interrupt_for_session
|
||||
interrupt_for_session(
|
||||
session_key=str(session_key or "") if _tui_owns_lifecycle else "",
|
||||
origin_ui_session_id=_lifecycle_own_sid(session), reason=end_reason)
|
||||
# Close the slash-worker in this single ``_finalized``-guarded chokepoint so a direct _finalize_session
|
||||
# caller can't leak it. Idempotent: close() is poll()-guarded.
|
||||
# Close the slash-worker in this single ``_finalized``-guarded chokepoint (a direct caller can't leak it); idempotent.
|
||||
with contextlib.suppress(Exception):
|
||||
if worker := session.get("slash_worker"):
|
||||
worker.close()
|
||||
|
||||
|
||||
# End reasons where the BACKEND reclaimed a session the client never asked to close; without a signal the
|
||||
# client's next prompt fails against a forgotten id. Client-initiated reasons (``tui_close`` etc.) are
|
||||
# deliberately absent.
|
||||
# End reasons where the BACKEND reclaimed a session the client never asked to close (else its next prompt fails
|
||||
# against a forgotten id). Client-initiated reasons (``tui_close`` etc.) are deliberately absent.
|
||||
_RECLAIM_END_REASONS = frozenset({"idle_timeout", "lru_evict", "ws_orphan_reap"})
|
||||
|
||||
|
||||
def _announce_session_reclaimed(session: dict, end_reason: str) -> None:
|
||||
"""Tell connected clients a session was reclaimed out from under them. Broadcast, not session-targeted:
|
||||
reap paths run on timer threads with no contextvar binding and the WS-orphan case has lost its
|
||||
transport, so ``_emit`` would bottom out on stdio. Best-effort; never breaks teardown."""
|
||||
"""Tell connected clients a session was reclaimed out from under them. Broadcast, not session-targeted: reap
|
||||
paths run on timer threads with no contextvar binding and no live transport, so ``_emit`` would hit stdio."""
|
||||
if end_reason not in _RECLAIM_END_REASONS:
|
||||
return
|
||||
try:
|
||||
@@ -296,9 +278,8 @@ def _announce_session_reclaimed(session: dict, end_reason: str) -> None:
|
||||
|
||||
|
||||
def _teardown_session(session: dict | None, *, end_reason: str = "tui_close") -> None:
|
||||
"""Fully tear down a session: finalize, unregister notifier, close agent. Shared by ``session.close`` and
|
||||
the orphaned-WS reaper. The slash-worker is closed in ``_finalize_session`` (the single chokepoint), NOT
|
||||
here. Idempotent via ``_finalized``."""
|
||||
"""Fully tear down a session: finalize, unregister notifier, close agent (``session.close`` + WS reaper). The
|
||||
slash-worker is closed in ``_finalize_session`` (the single chokepoint), NOT here. Idempotent via ``_finalized``."""
|
||||
if not session:
|
||||
return
|
||||
_finalize_session(session, end_reason=end_reason)
|
||||
@@ -313,8 +294,7 @@ def _teardown_session(session: dict | None, *, end_reason: str = "tui_close") ->
|
||||
|
||||
|
||||
def _attach_worker(sid: str, session: dict, worker) -> None:
|
||||
"""Store worker on session iff sid still maps to it, else close it — a concurrent teardown already popped
|
||||
the session and would orphan the worker."""
|
||||
"""Store worker on session iff sid still maps to it, else close it (a concurrent teardown popped the session)."""
|
||||
with _sessions_lock:
|
||||
if _sessions.get(sid) is session:
|
||||
session["slash_worker"] = worker
|
||||
@@ -324,14 +304,12 @@ def _attach_worker(sid: str, session: dict, worker) -> None:
|
||||
|
||||
def _pop_session_by_id(sid: str) -> dict | None:
|
||||
"""Atomically detach one live session from the registry — the ownership claim for teardown (a concurrent
|
||||
close/reaper then no-ops). Separate from ``_teardown_session`` because finalization does slow external
|
||||
work that must not run under ``_session_resume_lock``."""
|
||||
close/reaper no-ops). Separate from ``_teardown_session``: slow finalization must not run under the resume lock."""
|
||||
with _sessions_lock:
|
||||
session = _sessions.pop(sid, None)
|
||||
if session is not None:
|
||||
session["_closing"] = True
|
||||
if session is not None:
|
||||
session["_sid"] = sid # out of _sessions now, so teardown can't recover the live id by scanning
|
||||
session["_sid"] = sid # out of _sessions now, so teardown can't recover the live id by scanning
|
||||
return session
|
||||
|
||||
|
||||
@@ -354,8 +332,7 @@ def _teardown_popped_session(session: dict | None, *, end_reason: str = "tui_clo
|
||||
|
||||
|
||||
def _close_session_by_id(
|
||||
sid: str, *, end_reason: str = "tui_close", predicate: Callable[[dict], bool] | None = None
|
||||
) -> bool:
|
||||
sid: str, *, end_reason: str = "tui_close", predicate: Callable[[dict], bool] | None = None) -> bool:
|
||||
"""Idempotent teardown funnel for callers with no resume race (resume-sensitive callers pop under
|
||||
``_session_resume_lock`` and call ``_teardown_popped_session`` after releasing it). Automatic reapers pass
|
||||
``predicate`` to revalidate under ``_sessions_lock`` right before the claim, so a stale scan can't close a
|
||||
@@ -374,22 +351,19 @@ def _ws_session_is_detached(session: dict | None) -> bool:
|
||||
|
||||
|
||||
def _ws_session_is_orphaned(session: dict | None) -> bool:
|
||||
"""True if a WS session sits on ``_detached_ws_transport`` (where ``handle_ws`` parks disconnected clients)
|
||||
with no in-flight turn."""
|
||||
"""True if a WS session sits on ``_detached_ws_transport`` (where ``handle_ws`` parks disconnected clients), idle."""
|
||||
return bool(_ws_session_is_detached(session) and not session.get("running"))
|
||||
|
||||
|
||||
def _interrupt_session_turn(sid: str, session: dict, *, request_id: str | None = None) -> bool:
|
||||
"""Apply the shared ``session.interrupt`` contract to one claimed session; returns whether the compute-host
|
||||
control channel was used. The WS orphan reaper reuses this so a dead client gets the same partial-history
|
||||
and queued-prompt semantics."""
|
||||
"""Apply the shared ``session.interrupt`` contract to one claimed session; returns whether the compute-host control
|
||||
channel was used. The WS orphan reaper reuses this so a dead client gets the same partial-history/queue semantics."""
|
||||
use_compute_host = _session_uses_compute_host(session)
|
||||
should_interrupt = bool(session.get("running"))
|
||||
run_thread_alive = False
|
||||
if use_compute_host:
|
||||
# The host owns the live turn (parent `running` can lag a blocked tool), so let it decide. Gate on
|
||||
# `_compute_host_active`: HostSupervisor.interrupt() calls start(), so forwarding blindly for a lazy
|
||||
# session would spawn a child just to interrupt.
|
||||
# `_compute_host_active`: HostSupervisor.interrupt() calls start(), so a lazy session would spawn a child to interrupt.
|
||||
if should_interrupt or session.get("_compute_host_active"):
|
||||
_get_compute_host_supervisor().interrupt(sid, request_id=request_id)
|
||||
else:
|
||||
@@ -416,9 +390,8 @@ def _interrupt_session_turn(sid: str, session: dict, *, request_id: str | None =
|
||||
|
||||
|
||||
def _session_has_active_delegations(sid: str, session: dict | None = None) -> bool:
|
||||
"""True when UI session ``sid`` still owns live background work — by live UI sid AND, when the TUI owns the
|
||||
durable lifecycle (never for gateway-viewer tabs), by durable session_key so a delegation from an earlier
|
||||
tab of the same session keeps it alive."""
|
||||
"""True when UI session ``sid`` still owns live background work — by live UI sid AND, when the TUI owns the durable
|
||||
lifecycle (never for gateway-viewer tabs), by session_key so a delegation from an earlier tab keeps it alive."""
|
||||
if session is None:
|
||||
with _sessions_lock:
|
||||
session = _sessions.get(sid)
|
||||
@@ -428,8 +401,8 @@ def _session_has_active_delegations(sid: str, session: dict | None = None) -> bo
|
||||
owned_session_key = session_key = str(session.get("session_key") or "")
|
||||
session_id = getattr(session.get("agent"), "session_id", None) or session_key
|
||||
if session_id:
|
||||
# Only when this TUI/desktop session may end its durable row by key — never for gateway-originated
|
||||
# sessions (the TUI is only a viewer there). Unknown DB state -> assume ownership.
|
||||
# Only when this session may end its durable row by key — never for gateway-originated sessions (TUI is a
|
||||
# viewer there). Unknown DB state -> assume ownership.
|
||||
with contextlib.suppress(Exception):
|
||||
db = _get_db()
|
||||
if db is not None and _is_gateway_owned_source((db.get_session(session_id) or {}).get("source", "")):
|
||||
@@ -444,16 +417,14 @@ def _session_has_active_delegations(sid: str, session: dict | None = None) -> bo
|
||||
return True # a transient registry/import failure must not become destructive cleanup
|
||||
|
||||
|
||||
# One pending WS-orphan reap Timer per live sid; guarded by _sessions_lock. Cancelled by _cancel_ws_orphan_reap
|
||||
# from every resume/reuse/transport-rebind path — otherwise a reap could fire on a reattached session and
|
||||
# trigger a reap->broadcast->resume storm.
|
||||
# One pending WS-orphan reap Timer per live sid; guarded by _sessions_lock. Cancelled by _cancel_ws_orphan_reap from
|
||||
# every resume/reuse/transport-rebind path — else a reap on a reattached session triggers a reap->broadcast->resume storm.
|
||||
_pending_ws_reaps: dict[str, threading.Timer] = {}
|
||||
|
||||
|
||||
def _cancel_ws_orphan_reap(sid: str) -> None:
|
||||
"""Cancel a pending WS-orphan reap for ``sid`` (client came back). Called from every path that re-binds a
|
||||
live transport; closes the fired-but-not-run Timer race and stops dead Timers accumulating on flappy
|
||||
clients."""
|
||||
"""Cancel a pending WS-orphan reap for ``sid`` (client came back). Called from every path that re-binds a live
|
||||
transport; closes the fired-but-not-run Timer race and stops dead Timers accumulating on flappy clients."""
|
||||
with _sessions_lock:
|
||||
timer = _pending_ws_reaps.pop(sid, None)
|
||||
if timer is not None:
|
||||
@@ -462,14 +433,12 @@ def _cancel_ws_orphan_reap(sid: str) -> None:
|
||||
|
||||
|
||||
def _ws_orphan_turn_activity_is_fresh(session: dict) -> bool:
|
||||
"""Whether a detached RUNNING turn's activity clock (``_touch_activity``, the watchdog's clock) is still
|
||||
fresh — the reaper must NOT interrupt healthy detached work (closed laptop, backgrounded app).
|
||||
Conservative fallbacks keep the wedged-turn safety net: disabled threshold, missing/opaque agent,
|
||||
unreadable summary or never-stamped clock all report NOT fresh, i.e. eligible for interrupt-at-grace."""
|
||||
"""Whether a detached RUNNING turn's activity clock (``_touch_activity``) is still fresh — the reaper must NOT
|
||||
interrupt healthy detached work (closed laptop). Conservative: disabled threshold, missing/opaque agent, unreadable
|
||||
summary or never-stamped clock all report NOT fresh (eligible for interrupt-at-grace) to keep the wedged-turn net."""
|
||||
if _WS_ORPHAN_ACTIVITY_STALE_S <= 0:
|
||||
return False
|
||||
summary_fn = getattr(session.get("agent"), "get_activity_summary", None)
|
||||
if not callable(summary_fn):
|
||||
if not callable(summary_fn := getattr(session.get("agent"), "get_activity_summary", None)):
|
||||
return False
|
||||
try:
|
||||
elapsed = summary_fn().get("seconds_since_activity")
|
||||
@@ -485,13 +454,12 @@ def _schedule_ws_orphan_reap(sid: str, *, delay_s: float | None = None) -> None:
|
||||
return
|
||||
|
||||
def _reap() -> None:
|
||||
# Serialize the re-check against session.resume (rebinds under _session_resume_lock). Claim teardown
|
||||
# by popping under both locks, then release the resume lock before slow finalization. Ordering is
|
||||
# always resume_lock -> sessions_lock (RLock).
|
||||
# Serialize the re-check against session.resume (rebinds under _session_resume_lock). Claim teardown by popping
|
||||
# under both locks, then release the resume lock before slow finalization. Order: resume_lock -> sessions_lock.
|
||||
reschedule_delay = interrupt_session = session = None
|
||||
with _session_resume_lock:
|
||||
# Drop this Timer's registration so a concurrent _cancel_ws_orphan_reap can't cancel a dead Timer
|
||||
# while a rescheduled one (registered below) is the owner.
|
||||
# Drop this Timer's registration so a concurrent _cancel_ws_orphan_reap can't cancel a dead Timer while a
|
||||
# rescheduled one (registered below) is the owner.
|
||||
with _sessions_lock:
|
||||
_pending_ws_reaps.pop(sid, None)
|
||||
current = _sessions.get(sid)
|
||||
@@ -502,22 +470,18 @@ def _schedule_ws_orphan_reap(sid: str, *, delay_s: float | None = None) -> None:
|
||||
elif not current.get("running"):
|
||||
session = _pop_session_by_id(sid)
|
||||
elif not current.get("_client_gone_interrupt_requested") and _ws_orphan_turn_activity_is_fresh(current):
|
||||
# Client-absent but producing: keep running detached (the sentinel buffers emits) and re-check
|
||||
# each grace interval.
|
||||
logger.debug(
|
||||
"client_gone sid=%s action=defer (turn activity fresh; stale threshold %.0fs)",
|
||||
sid, _WS_ORPHAN_ACTIVITY_STALE_S)
|
||||
# Client-absent but producing: keep running detached (the sentinel buffers emits), re-check each grace.
|
||||
logger.debug("client_gone sid=%s action=defer (turn activity fresh; stale threshold %.0fs)",
|
||||
sid, _WS_ORPHAN_ACTIVITY_STALE_S)
|
||||
reschedule_delay = _WS_ORPHAN_REAP_GRACE_S
|
||||
else:
|
||||
# Mid-turn detached sessions must never drop the single Timer: interrupt once after grace, then
|
||||
# poll until turn-finalization settles.
|
||||
polls = int(current.get("_client_gone_interrupt_polls") or 0) + 1
|
||||
current["_client_gone_interrupt_polls"] = polls
|
||||
# Mid-turn detached sessions must never drop the single Timer: interrupt once after grace, then poll
|
||||
# until turn-finalization settles.
|
||||
polls = current["_client_gone_interrupt_polls"] = int(current.get("_client_gone_interrupt_polls") or 0) + 1
|
||||
if polls > _WS_ORPHAN_INTERRUPT_REAP_MAX_POLLS:
|
||||
# Never settled inside the budget — force-reap rather than park forever.
|
||||
logger.error(
|
||||
"client_gone sid=%s: turn did not settle after %d "
|
||||
"interrupt polls (%.0fs) — force-reaping detached session",
|
||||
"client_gone sid=%s: turn did not settle after %d interrupt polls (%.0fs) — force-reaping detached session",
|
||||
sid, polls - 1, (polls - 1) * _WS_ORPHAN_INTERRUPT_REAP_POLL_S)
|
||||
session = _pop_session_by_id(sid)
|
||||
else:
|
||||
@@ -553,17 +517,17 @@ def _schedule_ws_orphan_reap(sid: str, *, delay_s: float | None = None) -> None:
|
||||
|
||||
|
||||
def _close_sessions_for_transport(transport, *, end_reason: str = "ws_disconnect") -> tuple[int, int]:
|
||||
"""Single WS-disconnect teardown entry point: reap close_on_disconnect sessions (sidecar/dashboard)
|
||||
immediately and re-point the rest at the detached transport so later emits don't hit a dead socket;
|
||||
those go to the grace-windowed WS-orphan reaper. Returns ``(reaped, detached)`` counts."""
|
||||
"""Single WS-disconnect teardown entry point: reap close_on_disconnect sessions (sidecar/dashboard) immediately;
|
||||
re-point the rest at the detached transport (later emits miss the dead socket) for the grace-windowed WS-orphan
|
||||
reaper. Returns ``(reaped, detached)`` counts."""
|
||||
with _sessions_lock:
|
||||
owned = [(sid, s) for sid, s in _sessions.items() if s.get("transport") is transport]
|
||||
reaped = detached = 0
|
||||
for sid, session in owned:
|
||||
claimed_for_teardown = None
|
||||
should_schedule_reap = False
|
||||
# session.resume fast-path rebinds under _session_resume_lock: take it before re-checking so a
|
||||
# reconnect can't move the transport between check and claim.
|
||||
# session.resume fast-path rebinds under _session_resume_lock: take it so a reconnect can't move the transport
|
||||
# between check and claim.
|
||||
with _session_resume_lock, _sessions_lock:
|
||||
current = _sessions.get(sid)
|
||||
if current is not session:
|
||||
@@ -575,13 +539,12 @@ def _close_sessions_for_transport(transport, *, end_reason: str = "ws_disconnect
|
||||
if current.get("close_on_disconnect"):
|
||||
claimed_for_teardown = _pop_session_by_id(sid)
|
||||
else:
|
||||
# Point at the drop sentinel (NOT real stdio) so _ws_session_is_orphaned recognizes it;
|
||||
# standalone `hermes --tui` keeps real _stdio. UNLESS another window (multi-window pop-out
|
||||
# viewer) still shows the session: re-bind to the most recent surviving viewer instead.
|
||||
# Point at the drop sentinel (NOT real stdio) so _ws_session_is_orphaned recognizes it; standalone
|
||||
# `hermes --tui` keeps real _stdio. UNLESS another window (pop-out viewer) still shows the session:
|
||||
# re-bind to the most recent surviving viewer instead.
|
||||
viewers = current.get("viewers") or {}
|
||||
viewers.pop(transport, None)
|
||||
live = [vt for vt, ts in sorted(viewers.items(), key=lambda kv: kv[1])
|
||||
if vt is not transport and not _transport_is_dead(vt)]
|
||||
live = [vt for vt, ts in sorted(viewers.items(), key=lambda kv: kv[1]) if not _transport_is_dead(vt)]
|
||||
if live:
|
||||
current["transport"] = live[-1]
|
||||
else:
|
||||
|
||||
@@ -27,14 +27,12 @@ def _notif_session_matches(s: dict, keys) -> bool:
|
||||
|
||||
|
||||
def _notif_live_session_matches(keys, exclude: dict | None = None) -> bool:
|
||||
"""Any non-finalized live session (other than ``exclude``) matches ``keys``; False if the registry
|
||||
can't be read (fail open rather than drop the event)."""
|
||||
"""Any non-finalized live session (other than ``exclude``) matches ``keys``; False if the registry can't be read
|
||||
(fail open rather than drop the event)."""
|
||||
return _notif_locked_sessions(
|
||||
lambda ss: any(
|
||||
s is not exclude and not s.get("_finalized") and _notif_session_matches(s, keys) for s in ss.values()
|
||||
),
|
||||
False,
|
||||
)
|
||||
lambda ss: any(s is not exclude and not s.get("_finalized") and _notif_session_matches(s, keys)
|
||||
for s in ss.values()),
|
||||
False)
|
||||
|
||||
|
||||
def _notif_resolve_event_key(evt_key: str) -> str:
|
||||
@@ -47,25 +45,23 @@ def _notif_resolve_event_key(evt_key: str) -> str:
|
||||
|
||||
|
||||
def _notification_event_belongs_elsewhere(sid: str, session: dict, evt: dict) -> bool:
|
||||
"""True if ``evt`` is owned by a *different* live session. Background completions carry the
|
||||
``session_key`` of the session that started the work; async delegation completions also carry
|
||||
``origin_ui_session_id`` (the live TUI tab that commissioned them)."""
|
||||
"""True if ``evt`` is owned by a *different* live session. Background completions carry the ``session_key`` of the
|
||||
session that started the work; async delegation completions also carry ``origin_ui_session_id`` (the live TUI tab)."""
|
||||
evt_ui_sid = str(evt.get("origin_ui_session_id") or "")
|
||||
if evt_ui_sid:
|
||||
if evt_ui_sid == str(sid or "") and not session.get("_finalized"):
|
||||
return False
|
||||
if _notif_locked_sessions(lambda ss: evt_ui_sid in ss and not ss[evt_ui_sid].get("_finalized"), False):
|
||||
return True
|
||||
# Exact UI tab gone: fall through to durable session_key routing so a resumed
|
||||
# continuation with the same key/lineage can still claim it.
|
||||
# Exact UI tab gone: fall through to durable session_key routing so a resumed continuation with the same
|
||||
# key/lineage can still claim it.
|
||||
evt_key = str(evt.get("session_key") or "")
|
||||
if not evt_key:
|
||||
return False
|
||||
current_keys = _notif_current_keys(sid, session)
|
||||
# Compression can rotate AIAgent.session_id while the detached child is still running: map the
|
||||
# event's original key to its continuation tip so it reaches the live session instead of becoming
|
||||
# an orphan any poller may consume. A live continuation wins over the compressed parent, else a
|
||||
# stale parent tab could consume the event before the current conversation.
|
||||
# Compression can rotate AIAgent.session_id while the detached child is still running: map the event's original
|
||||
# key to its continuation tip so it reaches the live session instead of becoming an orphan any poller may consume.
|
||||
# A live continuation wins over the compressed parent, else a stale parent tab could consume the event first.
|
||||
resolved_key = _notif_resolve_event_key(evt_key)
|
||||
if resolved_key != evt_key:
|
||||
if resolved_key in current_keys:
|
||||
@@ -78,9 +74,8 @@ def _notification_event_belongs_elsewhere(sid: str, session: dict, evt: dict) ->
|
||||
|
||||
|
||||
def _session_owns_notification_event(sid: str, session: dict, evt: dict) -> bool:
|
||||
"""True iff *this* session PROVABLY owns ``evt`` (UI origin is this live session, or ``session_key``
|
||||
raw/compression-resolved matches) — the fail-closed gate for every addressed notification, without
|
||||
``_notification_event_belongs_elsewhere``'s orphan-adoption fallback."""
|
||||
"""True iff *this* session PROVABLY owns ``evt`` (UI origin is this live session, or ``session_key`` raw/compression-
|
||||
resolved matches) — the fail-closed gate for addressed notifications, without the orphan-adoption fallback."""
|
||||
if session.get("_finalized"):
|
||||
return False
|
||||
if str(evt.get("origin_ui_session_id") or "") == str(sid or ""):
|
||||
@@ -92,13 +87,11 @@ def _session_owns_notification_event(sid: str, session: dict, evt: dict) -> bool
|
||||
|
||||
def _notification_event_requires_owner(evt: dict) -> bool:
|
||||
"""Whether ``evt`` must be positively claimed before TUI delivery."""
|
||||
return evt.get("type") == "async_delegation" or bool(
|
||||
str(evt.get("origin_ui_session_id") or "") or str(evt.get("session_key") or "")
|
||||
)
|
||||
return evt.get("type") == "async_delegation" or bool(evt.get("origin_ui_session_id") or evt.get("session_key"))
|
||||
|
||||
|
||||
# Extra dedup fields per event type. Completions are terminal (one-shot per process session); watch
|
||||
# events are not — one process can match patterns many times, so their content is part of the key.
|
||||
# Extra dedup fields per event type. Completions are terminal (one-shot per process session); watch events are not —
|
||||
# one process can match patterns many times, so their content is part of the key.
|
||||
_DEDUP_EXTRA_FIELDS = {
|
||||
"watch_match": ("command", "pattern", "output", "suppressed", "message_id"),
|
||||
"watch_disabled": ("command", "message", "suppressed"),
|
||||
@@ -110,18 +103,16 @@ def _notification_event_dedup_key(evt: dict) -> tuple:
|
||||
"""UI-emission identity for a process notification event."""
|
||||
evt_type = evt.get("type", "completion")
|
||||
if evt_type == "async_delegation":
|
||||
# No process session_id: without this every completion keys as ("", "async_delegation")
|
||||
# and the second one is suppressed forever.
|
||||
# No process session_id: else every completion keys as ("", "async_delegation") and the second is suppressed forever.
|
||||
return (evt.get("delegation_id", ""), evt_type)
|
||||
extra = _DEDUP_EXTRA_FIELDS.get("watch_overflow_" if evt_type.startswith("watch_overflow_") else evt_type, ())
|
||||
return (evt.get("session_id", ""), evt_type, *(evt.get(f, 0 if f == "suppressed" else "") for f in extra))
|
||||
|
||||
|
||||
# Mirror gateway/kanban_watchers.py TERMINAL_KINDS: claim silent kinds (archived/unblocked) too so the
|
||||
# cursor advances past them and they can't wedge a later completed/blocked event behind an unclaimed row.
|
||||
# Mirror gateway/kanban_watchers.py TERMINAL_KINDS: claim silent kinds (archived/unblocked) too so the cursor advances
|
||||
# past them and they can't wedge a later completed/blocked event behind an unclaimed row.
|
||||
_KANBAN_NOTIFY_KINDS = ("completed", "blocked", "gave_up", "crashed", "timed_out", "status", "archived", "unblocked")
|
||||
_KANBAN_POLL_SECONDS = 5.0
|
||||
_LOOP_POLL_SECONDS = 5.0
|
||||
_KANBAN_POLL_SECONDS = _LOOP_POLL_SECONDS = 5.0
|
||||
|
||||
|
||||
def _notif_release_turn(session: dict) -> None:
|
||||
@@ -132,10 +123,9 @@ def _notif_release_turn(session: dict) -> None:
|
||||
def _notif_claim_turn(session: dict) -> bool:
|
||||
"""Claim the idle session (running=True) under history_lock; False if a turn is live."""
|
||||
with session["history_lock"]:
|
||||
if session.get("running"):
|
||||
return False
|
||||
claimed = not session.get("running")
|
||||
session["running"] = True
|
||||
return True
|
||||
return claimed
|
||||
|
||||
|
||||
def _notif_log_failure(what: str, exc: BaseException) -> None:
|
||||
@@ -158,18 +148,16 @@ def _notif_loop_status(sid: str, text: str) -> None:
|
||||
|
||||
|
||||
def _notif_slash_loop_tick(rid: str, sid: str, session: dict, mgr, wakeup: str) -> None:
|
||||
"""Slash-command /loop wakeup: route through the slash pipeline, not the model. No model reply to
|
||||
evaluate, so the tick completes immediately — unless the command resolves to a prompt (skill command
|
||||
etc.), which runs as a normal turn whose post-turn hook completes the tick."""
|
||||
"""Slash-command /loop wakeup: route through the slash pipeline, not the model. No model reply to evaluate, so the
|
||||
tick completes immediately — unless the command resolves to a prompt (skill command etc.), which runs as a normal
|
||||
turn whose post-turn hook completes the tick."""
|
||||
_notif_release_turn(session)
|
||||
try:
|
||||
parts = wakeup.lstrip()[1:].split(None, 1)
|
||||
resp = _methods["command.dispatch"](
|
||||
rid, {"name": parts[0] if parts else "", "arg": parts[1] if len(parts) > 1 else "", "session_id": sid}
|
||||
)
|
||||
rid, {"name": parts[0] if parts else "", "arg": parts[1] if len(parts) > 1 else "", "session_id": sid})
|
||||
payload = (resp or {}).get("result") or {}
|
||||
out = str(payload.get("output") or "").strip()
|
||||
if out:
|
||||
if out := str(payload.get("output") or "").strip():
|
||||
_notif_loop_status(sid, out)
|
||||
if payload.get("type") == "send" and payload.get("message"):
|
||||
if not _notif_claim_turn(session):
|
||||
@@ -186,21 +174,18 @@ def _notif_slash_loop_tick(rid: str, sid: str, session: dict, mgr, wakeup: str)
|
||||
|
||||
|
||||
def _maybe_fire_tui_loop_tick(sid: str, session: dict) -> None:
|
||||
"""Fire a due /loop wakeup for an idle TUI/Desktop/dashboard session (per-session poller, coarse
|
||||
cadence). Claims the session (running=True) before dispatching so a racing user prompt wins
|
||||
cleanly; the turn dispatcher's post-turn hook completes the tick."""
|
||||
"""Fire a due /loop wakeup for an idle TUI/Desktop/dashboard session (per-session poller, coarse cadence). Claims
|
||||
the session (running=True) before dispatching so a racing user prompt wins; the post-turn hook completes the tick."""
|
||||
try:
|
||||
from hermes_cli.loops import LoopManager, goal_blocks_loop_tick
|
||||
except Exception:
|
||||
return
|
||||
sid_key = session.get("session_key") or ""
|
||||
if not sid_key:
|
||||
if not (sid_key := session.get("session_key") or ""):
|
||||
return
|
||||
mgr = LoopManager(session_id=sid_key)
|
||||
if not mgr.is_due() or goal_blocks_loop_tick(sid_key) or not _notif_claim_turn(session):
|
||||
return # busy — stays due, next poll retries
|
||||
wakeup = mgr.fire_tick()
|
||||
if not wakeup:
|
||||
if not (wakeup := mgr.fire_tick()):
|
||||
_notif_release_turn(session)
|
||||
return
|
||||
rid = f"__loop__{int(time.time() * 1000)}"
|
||||
@@ -224,10 +209,8 @@ def _kb_first_line(value: Any, limit: int) -> str:
|
||||
|
||||
|
||||
def _kb_completed(task, payload: dict, title: str) -> str:
|
||||
if payload.get("summary"):
|
||||
handoff = _kb_first_line(payload["summary"], 200)
|
||||
else:
|
||||
handoff = _kb_first_line(task.result, 160) if getattr(task, "result", None) else ""
|
||||
handoff = (_kb_first_line(payload["summary"], 200) if payload.get("summary")
|
||||
else _kb_first_line(task.result, 160) if getattr(task, "result", None) else "")
|
||||
return f" done — {title}{handoff}"
|
||||
|
||||
|
||||
@@ -250,10 +233,9 @@ _KANBAN_EVENT_FORMATTERS = {
|
||||
|
||||
|
||||
def _format_kanban_event_text(sub: dict, task, ev, board_slug: str) -> Optional[str]:
|
||||
"""Single-line notification text for one kanban event; wording mirrors gateway/kanban_watchers.py
|
||||
so it reads the same as on Telegram. None for silent kinds."""
|
||||
entry = _KANBAN_EVENT_FORMATTERS.get(getattr(ev, "kind", ""))
|
||||
if entry is None:
|
||||
"""Single-line notification text for one kanban event; wording mirrors gateway/kanban_watchers.py (reads the same
|
||||
as on Telegram). None for silent kinds."""
|
||||
if (entry := _KANBAN_EVENT_FORMATTERS.get(getattr(ev, "kind", ""))) is None:
|
||||
return None
|
||||
glyph, fmt = entry
|
||||
task_id = sub.get("task_id", "")
|
||||
@@ -274,20 +256,18 @@ def _kb_board_key(_kb, board_meta) -> tuple[str, str]:
|
||||
|
||||
|
||||
def _kb_poll_board(_kb, slug: str, session_key: str) -> list:
|
||||
"""Claim + format this session's unseen events on one board. One poller per live session: the board
|
||||
is not opened writable unless it has a subscription owned by this exact session (if the read-only
|
||||
probe fails — locked/corrupt DB — delivery is preserved by falling through)."""
|
||||
try:
|
||||
"""Claim + format this session's unseen events on one board. One poller per live session: the board is not opened
|
||||
writable unless it has a subscription owned by this exact session (a failed read-only probe — locked/corrupt DB —
|
||||
falls through so delivery is preserved)."""
|
||||
with contextlib.suppress(Exception):
|
||||
if _kb.count_notify_subs(board=slug, platform="tui", chat_id=session_key) == 0:
|
||||
return []
|
||||
except Exception:
|
||||
pass
|
||||
try:
|
||||
conn = _kb.connect(board=slug)
|
||||
except Exception:
|
||||
return []
|
||||
texts: list = []
|
||||
try:
|
||||
with contextlib.closing(conn):
|
||||
try:
|
||||
subs = _kb.list_notify_subs(conn)
|
||||
except Exception:
|
||||
@@ -295,29 +275,25 @@ def _kb_poll_board(_kb, slug: str, session_key: str) -> list:
|
||||
for sub in subs:
|
||||
if (sub.get("platform") or "").lower() != "tui" or sub.get("chat_id") != session_key:
|
||||
continue
|
||||
sub_ident = dict(
|
||||
task_id=sub["task_id"], platform=sub["platform"], chat_id=sub["chat_id"], thread_id=sub.get("thread_id") or ""
|
||||
)
|
||||
sub_ident = dict(task_id=sub["task_id"], platform=sub["platform"], chat_id=sub["chat_id"],
|
||||
thread_id=sub.get("thread_id") or "")
|
||||
_old, _new, events = _kb.claim_unseen_events_for_sub(conn, kinds=_KANBAN_NOTIFY_KINDS, **sub_ident)
|
||||
if not events:
|
||||
continue
|
||||
task = _kb.get_task(conn, sub["task_id"])
|
||||
texts.extend(t for t in (_format_kanban_event_text(sub, task, ev, slug) for ev in events) if t)
|
||||
# Unsubscribe only on archive: ``done`` is reversible in review/controller flows, so keeping the
|
||||
# sub lets a later reopen notify the same session. The claimed cursor prevents replay.
|
||||
# Unsubscribe only on archive: ``done`` is reversible in review/controller flows, so keeping the sub lets a
|
||||
# later reopen notify the same session. The claimed cursor prevents replay.
|
||||
if task and getattr(task, "status", "") == "archived":
|
||||
with contextlib.suppress(Exception):
|
||||
_kb.remove_notify_sub(conn, **sub_ident)
|
||||
finally:
|
||||
conn.close()
|
||||
return texts
|
||||
|
||||
|
||||
def _collect_kanban_notifications(session: dict) -> list:
|
||||
"""Claim unseen terminal kanban events for this session's ``platform="tui"`` subscriptions
|
||||
(``kanban_create`` auto-subscribes with ``chat_id=HERMES_SESSION_KEY``; no "tui" messaging adapter
|
||||
exists, so this poller is the delivery path). Same atomic cursor-claim as the gateway notifier, so
|
||||
a sub is delivered exactly once even if a gateway and a TUI poll the same board DB."""
|
||||
"""Claim unseen terminal kanban events for this session's ``platform="tui"`` subscriptions (``kanban_create``
|
||||
auto-subscribes with ``chat_id=HERMES_SESSION_KEY``; no "tui" messaging adapter exists, so this poller is the
|
||||
delivery path). Same atomic cursor-claim as the gateway notifier: exactly-once even if a gateway polls the same DB."""
|
||||
session_key = str(session.get("session_key") or "")
|
||||
if not session_key or session.get("_finalized"):
|
||||
return []
|
||||
@@ -340,35 +316,32 @@ def _collect_kanban_notifications(session: dict) -> list:
|
||||
|
||||
|
||||
def _notif_poll_kanban(sid: str, session: dict) -> None:
|
||||
"""One kanban poll: emit new texts, buffer them, and run the buffered batch as a turn if idle. Events
|
||||
are cursor-claimed (never re-queued), so they wait in the buffer instead of dropping the agent turn."""
|
||||
"""One kanban poll: emit new texts, buffer them, and run the buffered batch as a turn if idle. Events are
|
||||
cursor-claimed (never re-queued), so they wait in the buffer instead of dropping the agent turn."""
|
||||
try:
|
||||
_kanban_texts = _collect_kanban_notifications(session)
|
||||
except Exception as _kb_exc:
|
||||
_notif_log_failure("kanban notification poll failed", _kb_exc)
|
||||
_kanban_texts = []
|
||||
for _kb_text in _kanban_texts:
|
||||
_emit("status.update", sid, {"kind": "process", "text": _kb_text})
|
||||
if _kanban_texts:
|
||||
session.setdefault("_kanban_pending", []).extend(_kanban_texts)
|
||||
texts = _collect_kanban_notifications(session)
|
||||
except Exception as exc:
|
||||
_notif_log_failure("kanban notification poll failed", exc)
|
||||
texts = []
|
||||
for text in texts:
|
||||
_emit("status.update", sid, {"kind": "process", "text": text})
|
||||
if texts:
|
||||
session.setdefault("_kanban_pending", []).extend(texts)
|
||||
if not session.get("_kanban_pending") or not _notif_claim_turn(session):
|
||||
return
|
||||
with session["history_lock"]:
|
||||
_batch, session["_kanban_pending"] = list(session.get("_kanban_pending") or []), []
|
||||
batch, session["_kanban_pending"] = list(session.get("_kanban_pending") or []), []
|
||||
with contextlib.suppress(Exception):
|
||||
_notif_submit(f"__notif__{int(time.time() * 1000)}", sid, session, "\n".join(_batch), "kanban notification dispatch failed")
|
||||
_notif_submit(f"__notif__{int(time.time() * 1000)}", sid, session, "\n".join(batch), "kanban notification dispatch failed")
|
||||
|
||||
|
||||
def _notif_dispatch_event(sid: str, session: dict, evt: dict, text: str) -> None:
|
||||
"""Run the claimed (running=True) agent turn for one notification event."""
|
||||
from tools.async_delegation import claim_event_delivery, complete_event_delivery, release_event_delivery
|
||||
claim = claim_event_delivery(evt, "tui-poller")
|
||||
if claim is None:
|
||||
if (claim := claim_event_delivery(evt, "tui-poller")) is None:
|
||||
return
|
||||
kwargs = (
|
||||
{"display_kind": "async_delegation_complete", "display_metadata": _async_delegation_display_metadata(evt)}
|
||||
if evt.get("type") == "async_delegation" else {}
|
||||
)
|
||||
kwargs = ({"display_kind": "async_delegation_complete", "display_metadata": _async_delegation_display_metadata(evt)}
|
||||
if evt.get("type") == "async_delegation" else {})
|
||||
try:
|
||||
_notif_submit(f"__notif__{int(time.time() * 1000)}", sid, session, text, "notification poller dispatch failed", **kwargs)
|
||||
except Exception:
|
||||
@@ -378,10 +351,10 @@ def _notif_dispatch_event(sid: str, session: dict, evt: dict, text: str) -> None
|
||||
|
||||
|
||||
def _notif_handle_event(sid, session, evt, emitted, registry, fmt, deferred) -> bool:
|
||||
"""Route one dequeued event: foreign (another live session owns it) → requeued, or onto ``deferred``
|
||||
during the shutdown drain; unowned (addressed but unprovable — never adopt an orphan) → dropped, except
|
||||
delegation payloads deferred for a resume; ours (or ownerless legacy, kept process-global) →
|
||||
status.update once, then an agent turn if idle. False = the drain must stop (session busy)."""
|
||||
"""Route one dequeued event: foreign (another live session owns it) → requeued, or onto ``deferred`` during the
|
||||
shutdown drain; unowned (addressed but unprovable — never adopt an orphan) → dropped, except delegation payloads
|
||||
deferred for a resume; ours (or ownerless legacy, kept process-global) → status.update once, then an agent turn if
|
||||
idle. False = the drain must stop (session busy)."""
|
||||
queue = registry.completion_queue
|
||||
evt_type, is_delegation = evt.get("type", "completion"), evt.get("type") == "async_delegation"
|
||||
if _notification_event_belongs_elsewhere(sid, session, evt):
|
||||
@@ -396,8 +369,7 @@ def _notif_handle_event(sid, session, evt, emitted, registry, fmt, deferred) ->
|
||||
if deferred is None:
|
||||
(logger.warning if is_delegation else logger.debug)(
|
||||
"Dropping unowned %s notification (origin=%r key=%r) instead of delivering to session %s",
|
||||
evt_type, origin, key, sid,
|
||||
)
|
||||
evt_type, origin, key, sid)
|
||||
elif is_delegation:
|
||||
deferred.append(evt)
|
||||
else:
|
||||
@@ -408,8 +380,8 @@ def _notif_handle_event(sid, session, evt, emitted, registry, fmt, deferred) ->
|
||||
text = fmt(evt)
|
||||
if not text:
|
||||
return True
|
||||
# Emit once per dedup key: a re-queued completion would otherwise re-emit every 0.5s while the
|
||||
# session is busy, while distinct watch_match events from one process must stay visible.
|
||||
# Emit once per dedup key: a re-queued completion would otherwise re-emit every 0.5s while the session is busy,
|
||||
# while distinct watch_match events from one process must stay visible.
|
||||
dedup_key = _notification_event_dedup_key(evt)
|
||||
if dedup_key not in emitted:
|
||||
_emit("status.update", sid, {"kind": "process", "text": text})
|
||||
@@ -418,29 +390,26 @@ def _notif_handle_event(sid, session, evt, emitted, registry, fmt, deferred) ->
|
||||
queue.put(evt)
|
||||
if deferred is not None:
|
||||
return False
|
||||
# Back off: the re-queued event keeps the queue non-empty, so without a sleep this loop spins
|
||||
# at 100% CPU while busy.
|
||||
time.sleep(0.25)
|
||||
time.sleep(0.25) # back off: the re-queued event keeps the queue non-empty, else this loop spins at 100% CPU
|
||||
return True
|
||||
_notif_dispatch_event(sid, session, evt, text)
|
||||
return True
|
||||
|
||||
|
||||
def _notification_poller_loop(stop_event: threading.Event, sid: str, session: dict) -> None:
|
||||
"""Daemon thread (started by _init_session()) that drains the process-global completion_queue for
|
||||
this session (see _notif_handle_event for ownership routing) and polls ``kanban_notify_subs`` every
|
||||
``_KANBAN_POLL_SECONDS`` — the delivery path for platform="tui" rows."""
|
||||
"""Daemon thread (started by _init_session()) that drains the process-global completion_queue for this session
|
||||
(ownership routing: _notif_handle_event) and polls ``kanban_notify_subs`` every ``_KANBAN_POLL_SECONDS`` — the
|
||||
delivery path for platform="tui" rows."""
|
||||
from tools.process_registry import process_registry, format_process_notification
|
||||
queue = process_registry.completion_queue
|
||||
emitted: set = set() # dedup re-queued events so one completion isn't emitted 50 times while busy
|
||||
handle = lambda evt, deferred: _notif_handle_event( # noqa: E731
|
||||
sid, session, evt, emitted, process_registry, format_process_notification, deferred
|
||||
)
|
||||
sid, session, evt, emitted, process_registry, format_process_notification, deferred)
|
||||
last_kanban_poll = last_loop_poll = 0.0
|
||||
while not stop_event.is_set() and not session.get("_finalized"):
|
||||
now = time.monotonic()
|
||||
# /loop wakeup driver: fire a due tick for THIS session while idle (same claim-under-lock as
|
||||
# kanban dispatch). An active non-parked /goal owns the idle boundary and defers it.
|
||||
# /loop wakeup driver: fire a due tick for THIS session while idle (same claim-under-lock as kanban dispatch).
|
||||
# An active non-parked /goal owns the idle boundary and defers it.
|
||||
if now - last_loop_poll >= _LOOP_POLL_SECONDS:
|
||||
last_loop_poll = now
|
||||
try:
|
||||
@@ -455,8 +424,8 @@ def _notification_poller_loop(stop_event: threading.Event, sid: str, session: di
|
||||
except Exception:
|
||||
continue
|
||||
handle(evt, None)
|
||||
# Drain remaining events after the stop signal so nothing is lost on shutdown; foreign and
|
||||
# orphaned-delegation events are handed back to the shared queue afterwards.
|
||||
# Drain remaining events after the stop signal so nothing is lost on shutdown; foreign and orphaned-delegation
|
||||
# events are handed back to the shared queue afterwards.
|
||||
deferred: list = []
|
||||
while not queue.empty():
|
||||
try:
|
||||
@@ -477,50 +446,42 @@ def _async_delegation_display_metadata(evt: dict) -> dict:
|
||||
completed_count = sum(1 for r in results if r.get("status") in {"completed", "success"})
|
||||
failed_count = sum(1 for r in results if r.get("status") in {"failed", "error"})
|
||||
duration = evt.get("total_duration_seconds") or evt.get("duration_seconds")
|
||||
return {
|
||||
"delegation_id": str(evt.get("delegation_id") or ""), "task_count": task_count,
|
||||
"completed_count": completed_count or task_count - failed_count, "failed_count": failed_count,
|
||||
**({"duration_seconds": duration} if isinstance(duration, (int, float)) else {}),
|
||||
}
|
||||
return {"delegation_id": str(evt.get("delegation_id") or ""), "task_count": task_count,
|
||||
"completed_count": completed_count or task_count - failed_count, "failed_count": failed_count,
|
||||
**({"duration_seconds": duration} if isinstance(duration, (int, float)) else {})}
|
||||
|
||||
|
||||
_desktop_ui_wired = False
|
||||
|
||||
|
||||
def _wire_desktop_sinks() -> None:
|
||||
"""Idempotently wire process-registry and desktop-tool sinks to renderer events: `agent.terminal.output`
|
||||
and `terminal.close` (drops a tab without killing the process) route to the window owning the process;
|
||||
desktop-only tools pass the turn's ``HERMES_UI_SESSION_ID`` as ``sid``. `_emit` is thread-safe."""
|
||||
"""Idempotently wire process-registry and desktop-tool sinks to renderer events: `agent.terminal.output` and
|
||||
`terminal.close` (drops a tab without killing the process) route to the window owning the process; desktop-only
|
||||
tools pass the turn's ``HERMES_UI_SESSION_ID`` as ``sid``. `_emit` is thread-safe."""
|
||||
global _desktop_ui_wired
|
||||
from tools.process_registry import process_registry
|
||||
|
||||
def _owner_sid_for_process(session) -> str:
|
||||
session_key = str(getattr(session, "session_key", "") or "")
|
||||
def _owner_sid(session) -> str:
|
||||
# session may be None (process already finished/pruned) — the tab can still linger and be closed.
|
||||
session_key = str(getattr(session, "session_key", "") or "") if session is not None else ""
|
||||
if not session_key:
|
||||
return ""
|
||||
with _sessions_lock:
|
||||
return next((sid for sid, s in _sessions.items() if str(s.get("session_key") or "") == session_key), "")
|
||||
if getattr(process_registry, "on_output", None) is None:
|
||||
process_registry.on_output = lambda session, chunk: _emit(
|
||||
"agent.terminal.output", _owner_sid_for_process(session), {"process_id": session.id, "chunk": chunk}
|
||||
)
|
||||
"agent.terminal.output", _owner_sid(session), {"process_id": session.id, "chunk": chunk})
|
||||
if getattr(process_registry, "on_close", None) is None:
|
||||
# session may be None (process already finished/pruned) — the tab can still linger and be closed.
|
||||
process_registry.on_close = lambda session, process_id: _emit(
|
||||
"terminal.close", _owner_sid_for_process(session) if session is not None else "", {"process_id": process_id}
|
||||
)
|
||||
process_registry.on_close = lambda session, pid: _emit("terminal.close", _owner_sid(session), {"process_id": pid})
|
||||
if not _desktop_ui_wired:
|
||||
try:
|
||||
with contextlib.suppress(Exception):
|
||||
from tools import desktop_ui
|
||||
except Exception:
|
||||
return
|
||||
desktop_ui.set_emitter(lambda sid, event, payload: _emit(event, sid, payload))
|
||||
_desktop_ui_wired = True
|
||||
desktop_ui.set_emitter(lambda sid, event, payload: _emit(event, sid, payload))
|
||||
_desktop_ui_wired = True
|
||||
|
||||
|
||||
# (stop_event, thread) for every poller started in this process, pruned of dead threads on each spawn.
|
||||
# Test teardowns reap leaked pollers through it: an unjoined poller steals events off the process-global
|
||||
# completion_queue mid-assertion in a LATER test.
|
||||
# (stop_event, thread) for every poller started in this process, pruned of dead threads on each spawn. Test teardowns
|
||||
# reap leaked pollers through it: an unjoined poller steals events off the process-global queue in a LATER test.
|
||||
_notification_pollers: list = []
|
||||
|
||||
|
||||
@@ -528,11 +489,8 @@ def _start_notification_poller(sid: str, session: dict) -> threading.Event:
|
||||
"""Start the background notification poller for a TUI session (thread name is greppable)."""
|
||||
_wire_desktop_sinks()
|
||||
stop = threading.Event()
|
||||
t = threading.Thread(
|
||||
target=_notification_poller_loop, args=(stop, sid, session), daemon=True, name=f"tui-notif-poller-{sid}"
|
||||
)
|
||||
_notification_pollers[:] = [(s, th) for (s, th) in _notification_pollers if th.is_alive()]
|
||||
_notification_pollers.append((stop, t))
|
||||
t = threading.Thread(target=_notification_poller_loop, args=(stop, sid, session), daemon=True, name=f"tui-notif-poller-{sid}")
|
||||
_notification_pollers[:] = [(s, th) for (s, th) in _notification_pollers if th.is_alive()] + [(stop, t)]
|
||||
t.start()
|
||||
return stop
|
||||
|
||||
@@ -546,14 +504,12 @@ def _hud_surface_note(session: dict) -> str:
|
||||
|
||||
|
||||
def _prepend_note(run_message: Any, note: str) -> Any:
|
||||
"""Prefix a per-turn note onto the MODEL INPUT, leaving the prompt alone: everything the model must
|
||||
know that the user did not type (interrupted reply, reactions, surface) arrives this way, so no
|
||||
scaffolding reaches the transcript and no sent message is rewritten — the cached prefix survives."""
|
||||
if not note:
|
||||
return run_message
|
||||
if isinstance(run_message, str):
|
||||
"""Prefix a per-turn note onto the MODEL INPUT, leaving the prompt alone: everything the model must know that the
|
||||
user did not type (interrupted reply, reactions, surface) arrives this way, so no scaffolding reaches the
|
||||
transcript and no sent message is rewritten — the cached prefix survives."""
|
||||
if note and isinstance(run_message, str):
|
||||
return f"{note}\n\n{run_message}"
|
||||
if isinstance(run_message, list):
|
||||
if note and isinstance(run_message, list):
|
||||
return [{"type": "text", "text": note}, *run_message]
|
||||
return run_message
|
||||
|
||||
|
||||
+91
-125
@@ -22,16 +22,12 @@ def _normalize_completion_path(path_part: str) -> str:
|
||||
|
||||
def _completion_cwd(params: dict | None = None) -> str:
|
||||
params = params or {}
|
||||
raw = (
|
||||
params.get("cwd")
|
||||
or _sessions.get(params.get("session_id") or "", {}).get("cwd")
|
||||
# A session bound to another profile resolves its workspace from THAT profile's config before the
|
||||
# launch profile's env var; the dashboard's in-memory gateway does NOT inherit the PTY child's bridged
|
||||
# TERMINAL_CWD, so a configured terminal.cwd is read directly.
|
||||
or _profile_configured_cwd(_profile_home(params.get("profile")))
|
||||
or _launch_configured_cwd()
|
||||
or os.environ.get("TERMINAL_CWD")
|
||||
or os.getcwd())
|
||||
# A session bound to another profile resolves its workspace from THAT profile's config before the launch profile's
|
||||
# env var; the dashboard's in-memory gateway does NOT inherit the PTY child's bridged TERMINAL_CWD, so a configured
|
||||
# terminal.cwd is read directly.
|
||||
raw = (params.get("cwd") or _sessions.get(params.get("session_id") or "", {}).get("cwd")
|
||||
or _profile_configured_cwd(_profile_home(params.get("profile"))) or _launch_configured_cwd()
|
||||
or os.environ.get("TERMINAL_CWD") or os.getcwd())
|
||||
with contextlib.suppress(Exception):
|
||||
resolved = os.path.abspath(os.path.expanduser(str(raw)))
|
||||
if os.path.isdir(resolved):
|
||||
@@ -49,15 +45,15 @@ def _workdir_terminal_cfg(key: str) -> str:
|
||||
|
||||
|
||||
def _terminal_task_cwd(session: dict | None) -> str:
|
||||
"""The cwd terminal_tool should use for this TUI session (NOT host-validated, unlike ``_completion_cwd``:
|
||||
a non-local backend's cwd lives inside the target environment)."""
|
||||
"""The cwd terminal_tool should use for this TUI session (NOT host-validated: a non-local backend's cwd lives
|
||||
inside the target environment)."""
|
||||
return _terminal_task_cwd_with_source(session)[0]
|
||||
|
||||
|
||||
def _terminal_task_cwd_with_source(session: dict | None) -> tuple[str, str]:
|
||||
"""``(cwd, source)``: ``"session"`` for THIS session's workspace (``explicit_cwd``/tracked dir), ``"process"``
|
||||
for the process-global ``TERMINAL_CWD`` / ``terminal.cwd`` fallback — under per-session docker isolation
|
||||
that is a launch artifact of a PREVIOUS session, so terminal_tool refuses it as a bind-mount source."""
|
||||
"""``(cwd, source)``: ``"session"`` for THIS session's workspace (``explicit_cwd``/tracked dir), ``"process"`` for
|
||||
the global ``TERMINAL_CWD``/``terminal.cwd`` fallback — under per-session docker isolation that is a PREVIOUS
|
||||
session's launch artifact, so terminal_tool refuses it as a bind-mount source."""
|
||||
backend = _effective_terminal_backend()
|
||||
if backend != "local":
|
||||
# THIS session's explicit workspace beats the LAST session's env var.
|
||||
@@ -84,8 +80,7 @@ def _session_cwd(session: dict | None) -> str:
|
||||
return str(session["cwd"]) if session and session.get("cwd") else _completion_cwd()
|
||||
|
||||
|
||||
# Sources whose launch directory is an artifact of how the app was started, not a workspace the user picked
|
||||
# (everything else is a directory the user cd'd into).
|
||||
# Sources whose launch directory is an artifact of how the app was started, not a workspace the user picked.
|
||||
_LAUNCH_CWD_NOT_A_WORKSPACE = {"desktop"}
|
||||
|
||||
|
||||
@@ -95,8 +90,7 @@ def _context_cwd_is_launch_artifact(session: dict | None) -> bool:
|
||||
|
||||
|
||||
def _persisted_session_cwd(session: dict) -> str | None:
|
||||
"""The cwd to stamp on the session's DB row, or None to leave it unset (see :func:`_ensure_session_db_row`
|
||||
for the desktop vs terminal launch-dir rule)."""
|
||||
"""The cwd to stamp on the session's DB row, or None to leave it unset (launch-dir rule: ``_ensure_session_db_row``)."""
|
||||
if session.get("explicit_cwd"):
|
||||
return _session_cwd(session)
|
||||
if _session_source(session) in _LAUNCH_CWD_NOT_A_WORKSPACE:
|
||||
@@ -105,10 +99,9 @@ def _persisted_session_cwd(session: dict) -> str | None:
|
||||
|
||||
|
||||
def _heal_dead_cwd(cwd: str) -> str:
|
||||
"""Resolve a session cwd inside a now-deleted directory (e.g. a removed linked worktree, which probes to no
|
||||
branch while the sidebar folds it to the main lane): walk up to the first existing ancestor and take its
|
||||
common git root. Local backends only — a remote/SSH cwd may legitimately not exist on the host, so callers
|
||||
skip healing there."""
|
||||
"""Resolve a session cwd inside a now-deleted directory (e.g. a removed linked worktree, which probes to no branch
|
||||
while the sidebar folds it to the main lane): walk up to the first existing ancestor and take its common git root.
|
||||
Local backends only — a remote/SSH cwd may legitimately not exist on the host, so callers skip healing there."""
|
||||
raw = (cwd or "").strip()
|
||||
if not raw or os.path.isdir(raw):
|
||||
return raw
|
||||
@@ -133,20 +126,16 @@ def _is_local_terminal_backend() -> bool:
|
||||
|
||||
|
||||
def _effective_terminal_backend() -> str:
|
||||
"""Active terminal backend name (``local``, ``docker``, ``ssh``, ...): ``TERMINAL_ENV`` when set (launchers
|
||||
bridge ``terminal.backend`` into env), else the ``terminal.backend`` config key (desktop/TUI in-process
|
||||
gateways skip that bridge)."""
|
||||
"""Active terminal backend name (``local``, ``docker``, ``ssh``, ...): ``TERMINAL_ENV`` when set (launchers bridge
|
||||
``terminal.backend`` into env), else the ``terminal.backend`` config key (in-process gateways skip that bridge)."""
|
||||
backend = (os.environ.get("TERMINAL_ENV") or "").strip().lower()
|
||||
if not backend or backend == "local":
|
||||
cfg_backend = _workdir_terminal_cfg("backend").lower()
|
||||
if cfg_backend and cfg_backend != "local":
|
||||
backend = cfg_backend
|
||||
backend = _workdir_terminal_cfg("backend").lower()
|
||||
return backend or "local"
|
||||
|
||||
|
||||
def _display_session_cwd(session: dict | None) -> str:
|
||||
"""Session cwd for display/probe surfaces, healed past deleted worktrees; the healed value is persisted
|
||||
back (best-effort, local only)."""
|
||||
"""Session cwd for display/probe surfaces, healed past deleted worktrees (healed value persisted back; local only)."""
|
||||
cwd = _session_cwd(session)
|
||||
if not _is_local_terminal_backend():
|
||||
return cwd
|
||||
@@ -158,40 +147,37 @@ def _display_session_cwd(session: dict | None) -> str:
|
||||
|
||||
|
||||
def _reconcile_session_cwd_from_terminal(session: dict | None) -> bool:
|
||||
"""Re-anchor a session that SETTLED in another worktree of the SAME repo. Returns moved. An agent told to
|
||||
work in a fresh worktree `git worktree add`s and `cd`s in while the session stays pinned (labelled with the
|
||||
primary checkout's branch). A plain `cd` is deliberately NOT a workspace move (see
|
||||
``_apply_project_workspace``): a non-git workspace stepping into a repo or a visit to an unrelated repo is
|
||||
browsing, and an explicitly chosen workspace is never overridden. Local backends only (a remote cwd cannot
|
||||
be stat'ed or git-probed here)."""
|
||||
"""Re-anchor a session that SETTLED in another worktree of the SAME repo. Returns moved. An agent told to work in
|
||||
a fresh worktree `git worktree add`s and `cd`s in while the session stays pinned (labelled with the primary
|
||||
checkout's branch). A plain `cd` is deliberately NOT a workspace move (see ``_apply_project_workspace``): a non-git
|
||||
workspace stepping into a repo or a visit to an unrelated repo is browsing, and an explicitly chosen workspace is
|
||||
never overridden. Local backends only (a remote cwd cannot be stat'ed or git-probed here)."""
|
||||
# An explicit choice only moves by another explicit action; a cwd adopted HERE is marked `cwd_from_settle` so
|
||||
# successive settles keep following.
|
||||
if not session or not _is_local_terminal_backend():
|
||||
return False
|
||||
# An explicit choice only moves by another explicit action; a cwd adopted HERE is marked `cwd_from_settle`
|
||||
# so successive settles keep following.
|
||||
if session.get("explicit_cwd") and not session.get("cwd_from_settle"):
|
||||
return False
|
||||
try:
|
||||
from tools.terminal_tool import get_session_cwd
|
||||
recorded = get_session_cwd(session.get("session_key") or "")
|
||||
if not (recorded := get_session_cwd(session.get("session_key") or "")):
|
||||
return False
|
||||
except Exception:
|
||||
return False
|
||||
if not recorded:
|
||||
return False
|
||||
resolved = os.path.abspath(os.path.expanduser(str(recorded)))
|
||||
current = os.path.abspath(os.path.expanduser(_session_cwd(session)))
|
||||
if resolved == current or not os.path.isdir(resolved):
|
||||
return False
|
||||
# Worktree ROOTS (folding to the common root would hide the move), both in a git tree, different from each
|
||||
# other, sharing the SAME common .git dir.
|
||||
landed = _git_repo_root_for_cwd(resolved)
|
||||
current_root = _git_repo_root_for_cwd(current)
|
||||
# Worktree ROOTS (folding to the common root would hide the move), both in a git tree, different from each other,
|
||||
# sharing the SAME common .git dir.
|
||||
landed, current_root = _git_repo_root_for_cwd(resolved), _git_repo_root_for_cwd(current)
|
||||
if not landed or not current_root or landed == current_root:
|
||||
return False
|
||||
landed_common = _git_common_repo_root_for_cwd(resolved)
|
||||
if not landed_common or landed_common != _git_common_repo_root_for_cwd(current):
|
||||
return False
|
||||
# This is the session's workspace now (a desktop launch-artifact cwd earns a real row); the settle marker
|
||||
# keeps it overridable by the NEXT settle.
|
||||
# This is the session's workspace now (a desktop launch-artifact cwd earns a real row); the settle marker keeps it
|
||||
# overridable by the NEXT settle.
|
||||
session.update(cwd=resolved, explicit_cwd=True, cwd_from_settle=True)
|
||||
_register_session_cwd(session)
|
||||
_persist_session_cwd_and_schedule_git_meta(session, resolved)
|
||||
@@ -199,8 +185,8 @@ def _reconcile_session_cwd_from_terminal(session: dict | None) -> bool:
|
||||
|
||||
|
||||
def _emit_settled_session_info(sid: str, session: dict, agent) -> None:
|
||||
"""Emit end-of-turn ``session.info``, reconciling a settled cwd first: the agent has stopped moving, and
|
||||
riding the reconcile on the turn-end event needs no new event type/round trip."""
|
||||
"""Emit end-of-turn ``session.info``, reconciling a settled cwd first (the agent has stopped moving; riding the
|
||||
turn-end event needs no new event type/round trip)."""
|
||||
try:
|
||||
_reconcile_session_cwd_from_terminal(session)
|
||||
except Exception:
|
||||
@@ -223,18 +209,16 @@ def _register_session_cwd(session: dict | None) -> None:
|
||||
|
||||
|
||||
def _workdir_row_model_config(session: dict) -> tuple[str, dict]:
|
||||
"""``(model, model_config)`` for a fresh session row. The session's own model/effort/fast pick (composer
|
||||
override or restored /model switch) must own the row: the agent isn't built yet at first prompt.submit,
|
||||
and writing the global default here wins the INSERT-OR-IGNORE race, so a reconnect silently reverts to
|
||||
the profile default. model_config carries provider/reasoning/service_tier so resume restores effort +
|
||||
fast too."""
|
||||
override = session.get("model_override")
|
||||
override = override if isinstance(override, dict) else {}
|
||||
"""``(model, model_config)`` for a fresh session row. The session's own model/effort/fast pick (composer override
|
||||
or restored /model switch) must own the row: the agent isn't built yet at first prompt.submit, and writing the
|
||||
global default here wins the INSERT-OR-IGNORE race (a reconnect silently reverts to the profile default).
|
||||
model_config carries provider/reasoning/service_tier so resume restores effort + fast too."""
|
||||
override = raw if isinstance(raw := session.get("model_override"), dict) else {}
|
||||
row_model = str(override.get("model") or "").strip() or _resolve_model()
|
||||
model_config: dict = {k: str(v) for k in ("model", "provider", "base_url", "api_mode") if (v := override.get(k))}
|
||||
# A RESOLVED provider "custom" (named ``providers:``/``custom_providers:`` entry) persisted bare here is the
|
||||
# origin of "No LLM provider configured" rows (resume routes to OpenRouter with no key). Recover the
|
||||
# durable ``custom:<name>`` identity (matches _runtime_model_config).
|
||||
# A RESOLVED provider "custom" (named ``providers:``/``custom_providers:`` entry) persisted bare here is the origin
|
||||
# of "No LLM provider configured" rows (resume routes to OpenRouter with no key). Recover the durable
|
||||
# ``custom:<name>`` identity (matches _runtime_model_config).
|
||||
if str(model_config.get("provider") or "").strip().lower() == "custom":
|
||||
try:
|
||||
from hermes_cli.runtime_provider import canonical_custom_identity
|
||||
@@ -247,15 +231,14 @@ def _workdir_row_model_config(session: dict) -> tuple[str, dict]:
|
||||
if (reasoning := session.get("create_reasoning_override")) is not None:
|
||||
model_config["reasoning_config"] = reasoning
|
||||
if (service_tier := session.get("create_service_tier_override")) is not None:
|
||||
# "" is the in-memory sentinel for an explicit normal tier (bypasses _make_agent's profile fallback);
|
||||
# persist a durable marker so resume can tell it from an inherited tier.
|
||||
# "" is the in-memory sentinel for an explicit normal tier (bypasses _make_agent's profile fallback); persist a
|
||||
# durable marker so resume can tell it from an inherited tier.
|
||||
model_config["service_tier"] = service_tier or "normal"
|
||||
# Same ``_branched_from`` marker the TUI /branch uses (list_sessions_rich + sidebar nesting).
|
||||
if parent_session_id := session.get("parent_session_id"):
|
||||
model_config["_branched_from"] = parent_session_id
|
||||
# Bot-Mode canonical chats / room plumbing are plugin-owned scratch conversations whose runtime must ALWAYS
|
||||
# follow the member profile's CURRENT config, never the provider pinned at first write; persist that
|
||||
# contract for resume (see _stored_session_runtime_overrides).
|
||||
# Bot-Mode canonical chats / room plumbing are plugin-owned scratch conversations whose runtime must ALWAYS follow
|
||||
# the member profile's CURRENT config, never the provider pinned at first write (see _stored_session_runtime_overrides).
|
||||
for flag in ("room_plumbing", "follow_profile_config"):
|
||||
if session.get(flag):
|
||||
model_config[flag] = True
|
||||
@@ -263,46 +246,43 @@ def _workdir_row_model_config(session: dict) -> tuple[str, dict]:
|
||||
|
||||
|
||||
def _ensure_session_db_row(session: dict) -> bool:
|
||||
"""Idempotently persist the session's DB row on first real activity (prompt.submit), so abandoned drafts
|
||||
never leave an empty "Untitled" session. INSERT OR IGNORE: re-calls and the AIAgent's lazy create are
|
||||
no-ops. Returns False only when the store is unavailable (no openable state.db) — prompt.submit fails the
|
||||
send loudly instead of streaming into a store that will never save it; no key / best-effort / success are
|
||||
all True.
|
||||
"""Idempotently persist the session's DB row on first real activity (prompt.submit), so abandoned drafts never
|
||||
leave an empty "Untitled" session. INSERT OR IGNORE: re-calls and the AIAgent's lazy create are no-ops. Returns
|
||||
False only when the store is unavailable (no openable state.db) — prompt.submit fails the send loudly instead of
|
||||
streaming into a store that will never save it; no key / best-effort / success are all True.
|
||||
|
||||
A cwd the user *chose* is always persisted. Otherwise the launch directory stands in only for terminal
|
||||
sessions (the user deliberately ``cd``'d there; dropping it left the sidebar with no cwd AND no
|
||||
git_repo_root); desktop launch dirs (``/``, home) stay null -> "No workspace"."""
|
||||
key = session.get("session_key")
|
||||
if not key:
|
||||
A cwd the user *chose* is always persisted. Otherwise the launch directory stands in only for terminal sessions
|
||||
(the user deliberately ``cd``'d there; dropping it left the sidebar with no cwd AND no git_repo_root); desktop
|
||||
launch dirs (``/``, home) stay null -> "No workspace"."""
|
||||
if not (key := session.get("session_key")):
|
||||
return
|
||||
# Persist into the session's own profile db (global remote mode), not the launch profile's — otherwise the
|
||||
# unified list mis-tags the row and resume 404s ("session not found").
|
||||
# Persist into the session's own profile db (global remote mode), not the launch profile's — otherwise the unified
|
||||
# list mis-tags the row and resume 404s ("session not found").
|
||||
profile_home = session.get("profile_home")
|
||||
with _workdir_owner_db(session, "failed to open profile db for session row") as db:
|
||||
if db is _WORKDIR_DB_OPEN_FAILED:
|
||||
return False
|
||||
if db is None:
|
||||
# Fail loud ONLY when the store failed to open (_db_error records the SessionDB open exception);
|
||||
# None with no recorded error means "no store in this context" -> True.
|
||||
# Fail loud ONLY when the store failed to open (_db_error records the SessionDB open exception); None with
|
||||
# no recorded error means "no store in this context" -> True.
|
||||
return _db_error is None
|
||||
row_model, model_config = _workdir_row_model_config(session)
|
||||
try:
|
||||
db.create_session(
|
||||
key, source=_session_source(session), model=row_model, model_config=model_config or None,
|
||||
parent_session_id=session.get("parent_session_id") or None, cwd=_persisted_session_cwd(session),
|
||||
# Self-describing rows: aggregators merging several profile DBs can't rely on which file a row
|
||||
# came from; a NULL is only repaired by the one-shot backfill.
|
||||
# Self-describing rows: aggregators merging several profile DBs can't rely on which file a row came
|
||||
# from; a NULL is only repaired by the one-shot backfill.
|
||||
profile_name=Path(profile_home).name if profile_home else _current_profile_name())
|
||||
# Born hidden (session.create hidden=true, or set_hidden before the row existed): apply the
|
||||
# deferred intent now, like pending_title.
|
||||
# Born hidden (session.create hidden=true, or set_hidden before the row existed): apply the deferred intent.
|
||||
if session.get("pending_hidden"):
|
||||
try:
|
||||
db.set_session_hidden(key, True)
|
||||
except Exception:
|
||||
logger.debug("failed to apply pending hidden flag", exc_info=True)
|
||||
except Exception as exc:
|
||||
# Disk-full is not a soft failure: swallowed here, prompt.submit returns {"status":"streaming"} and
|
||||
# the message vanishes silently.
|
||||
# Disk-full is not a soft failure: swallowed here, prompt.submit returns {"status":"streaming"} and the
|
||||
# message vanishes silently.
|
||||
_workdir_reraise_disk_full(exc, "failed to persist desktop session row")
|
||||
return True
|
||||
|
||||
@@ -315,21 +295,19 @@ def _workdir_reraise_disk_full(exc: BaseException, log_msg: str) -> None:
|
||||
logger.debug(log_msg, exc_info=True)
|
||||
|
||||
|
||||
# Seed row fields copied from the parent transcript. display_kind/metadata: timeline markers ride as role=user,
|
||||
# dropping the tag re-plants them as bare user turns after a restart and corrupts the truncate ordinal address
|
||||
# space. timestamp: parent's original, not "now".
|
||||
# Seed row fields copied from the parent transcript. display_kind/metadata: timeline markers ride as role=user;
|
||||
# dropping the tag re-plants them as bare user turns after a restart and corrupts the truncate ordinal address space.
|
||||
_WORKDIR_SEED_FIELDS = (
|
||||
"content", "reasoning", "reasoning_content", "reasoning_details", "codex_reasoning_items",
|
||||
"codex_message_items", "display_kind", "display_metadata", "timestamp")
|
||||
|
||||
|
||||
def _persist_branch_seed(session: dict) -> None:
|
||||
"""First-turn persist of a branch's copied transcript. A branch is a draft until its first submit: the
|
||||
parent's messages live only in ``session["history"]`` (ridden into the agent as ``conversation_history``,
|
||||
which ``_flush_messages_to_session_db`` skips by identity), so the row would otherwise resume missing its
|
||||
pre-branch context. Runs once, after ``_ensure_session_db_row`` wrote the row + parent link."""
|
||||
key = session.get("session_key")
|
||||
if not key or not session.get("parent_session_id") or session.get("_branch_seed_persisted"):
|
||||
"""First-turn persist of a branch's copied transcript. A branch is a draft until its first submit: the parent's
|
||||
messages live only in ``session["history"]`` (ridden into the agent as ``conversation_history``, which
|
||||
``_flush_messages_to_session_db`` skips by identity), so the row would otherwise resume missing its pre-branch
|
||||
context. Runs once, after ``_ensure_session_db_row`` wrote the row + parent link."""
|
||||
if not (key := session.get("session_key")) or not session.get("parent_session_id") or session.get("_branch_seed_persisted"):
|
||||
return
|
||||
with session["history_lock"]:
|
||||
seed = [dict(msg) for msg in (session.get("history") or [])]
|
||||
@@ -339,29 +317,25 @@ def _persist_branch_seed(session: dict) -> None:
|
||||
if db is None:
|
||||
return
|
||||
try:
|
||||
# Chunked so each BEGIN IMMEDIATE stays short (a seed can be hundreds of rows); a mid-copy failure
|
||||
# leaves a partial seed with _branch_seed_persisted unset.
|
||||
# Chunked so each BEGIN IMMEDIATE stays short (a seed can be hundreds of rows); a mid-copy failure leaves a
|
||||
# partial seed with _branch_seed_persisted unset.
|
||||
db.append_messages_batch(
|
||||
key,
|
||||
[{"role": msg.get("role", "user"), **{f: msg.get(f) for f in _WORKDIR_SEED_FIELDS}} for msg in seed],
|
||||
key, [{"role": msg.get("role", "user"), **{f: msg.get(f) for f in _WORKDIR_SEED_FIELDS}} for msg in seed],
|
||||
chunk_rows=500)
|
||||
session["_branch_seed_persisted"] = True
|
||||
except Exception as exc:
|
||||
_workdir_reraise_disk_full(exc, "branch seed persist failed")
|
||||
|
||||
|
||||
# Yielded by _workdir_owner_db when the profile db failed to OPEN (as opposed to "no store in this context");
|
||||
# _ensure_session_db_row fails loud on it.
|
||||
# Yielded by _workdir_owner_db when the profile db failed to OPEN (vs "no store in this context"); row creation fails loud.
|
||||
_WORKDIR_DB_OPEN_FAILED = object()
|
||||
|
||||
|
||||
@contextlib.contextmanager
|
||||
def _workdir_owner_db(session: dict, fail_log: str):
|
||||
"""Body of :func:`_session_db`; also used directly by ``_ensure_session_db_row`` so a test-patched
|
||||
``_session_db`` does not change row creation."""
|
||||
"""Body of :func:`_session_db`; ``_ensure_session_db_row`` uses it directly so a patched ``_session_db`` can't alter rows."""
|
||||
db, close_db = None, False
|
||||
profile_home = session.get("profile_home")
|
||||
if profile_home:
|
||||
if profile_home := session.get("profile_home"):
|
||||
try:
|
||||
from hermes_state import get_shared_session_db
|
||||
db, close_db = get_shared_session_db(Path(profile_home) / "state.db"), True
|
||||
@@ -381,20 +355,17 @@ def _workdir_owner_db(session: dict, fail_log: str):
|
||||
|
||||
@contextlib.contextmanager
|
||||
def _session_db(session: dict):
|
||||
"""Yield the SessionDB that owns this session's row (profile-aware): a remote/profile session persists into
|
||||
its own profile's ``state.db`` (fresh handle, closed on exit); everything else borrows the shared
|
||||
``_get_db()`` handle (left open). Yields None when unavailable."""
|
||||
"""Yield the SessionDB that owns this session's row (profile-aware): a remote/profile session persists into its own
|
||||
profile's ``state.db`` (fresh handle, closed on exit); else the shared ``_get_db()`` handle (left open). None if unavailable."""
|
||||
with _workdir_owner_db(session, "failed to open profile db for session") as db:
|
||||
yield None if db is _WORKDIR_DB_OPEN_FAILED else db
|
||||
|
||||
|
||||
def _rewind_active_session_history(
|
||||
session: dict, user_ordinal: int, *, require_retryable: bool = False
|
||||
) -> tuple[list[dict], dict, int]:
|
||||
"""Rewind one canonical user turn while retaining carrier scaffolding. Caller holds ``history_lock``.
|
||||
Persistent sessions archive the target and tail, inserting a composite carrier's hidden handoff in the
|
||||
same transaction; memory is installed only after the durable commit, from the already-validated prefix
|
||||
plus the returned scaffold row id (no fallible post-commit reload)."""
|
||||
session: dict, user_ordinal: int, *, require_retryable: bool = False) -> tuple[list[dict], dict, int]:
|
||||
"""Rewind one canonical user turn while retaining carrier scaffolding. Caller holds ``history_lock``. Persistent
|
||||
sessions archive the target and tail, inserting a composite carrier's hidden handoff in the same transaction; memory
|
||||
is installed only after the durable commit, from the validated prefix + returned scaffold row id (no reload)."""
|
||||
from agent.context_compressor import (
|
||||
history_before_user_originated_turn, retryable_user_text, split_user_originated_turn,
|
||||
user_originated_turn_view)
|
||||
@@ -442,8 +413,7 @@ def _rewind_active_session_history(
|
||||
durable_prefix, durable_live_view = history_before_user_originated_turn(durable, durable_target_index)
|
||||
if _comparison_content(durable_live_view) != _comparison_content(live_view):
|
||||
raise RuntimeError("session history changed before the rewind could be persisted")
|
||||
target_row_id = durable_target.get("_row_id")
|
||||
if not isinstance(target_row_id, int):
|
||||
if not isinstance(target_row_id := durable_target.get("_row_id"), int):
|
||||
raise RuntimeError("rewind target has no durable row identity")
|
||||
if require_retryable:
|
||||
retryable_user_text(durable_live_view.get("content"))
|
||||
@@ -452,13 +422,12 @@ def _rewind_active_session_history(
|
||||
session_key, target_row_id, preserve_compaction_handoff=scaffold is not None,
|
||||
expected_active_ids=expected_active_ids, expected_target_content=durable_live_view.get("content"))
|
||||
if scaffold is not None:
|
||||
replacement_id = result.get("replacement_message_id")
|
||||
if not isinstance(replacement_id, int):
|
||||
if not isinstance(replacement_id := result.get("replacement_message_id"), int):
|
||||
raise RuntimeError("rewind commit did not return the replacement scaffold id")
|
||||
durable_prefix[-1].update(_row_id=replacement_id, _db_persisted=True)
|
||||
installed[-1] = durable_prefix[-1]
|
||||
# Clients address follow-ups by durable row id: keep the richer warm content but copy row
|
||||
# identities when the shapes align.
|
||||
# Clients address follow-ups by durable row id: keep the richer warm content but copy row identities when
|
||||
# the shapes align.
|
||||
if len(installed) == len(durable_prefix) and all(
|
||||
warm.get("role") == durable_message.get("role")
|
||||
and bool(warm.get("display_kind")) == bool(durable_message.get("display_kind"))
|
||||
@@ -501,8 +470,7 @@ def _workdir_valid_generation(generation) -> bool:
|
||||
def _persist_session_git_meta(session: dict, cwd: str, generation: int) -> None:
|
||||
"""Resolve + persist a session's git branch / repo root on a daemon thread: inline ``git`` probes on the
|
||||
session-init / cwd-set path would stall startup on a slow or unreachable ``cwd``. Persists via the same
|
||||
profile-aware db the caller wrote ``cwd`` to. Best-effort: a probe failure leaves the enrichment columns
|
||||
unset (project tree uses its live resolver)."""
|
||||
profile-aware db the caller wrote ``cwd`` to. Best-effort: a probe failure leaves the enrichment columns unset."""
|
||||
session_key = session.get("session_key", "")
|
||||
if not session_key or not cwd or not _workdir_valid_generation(generation):
|
||||
return
|
||||
@@ -511,8 +479,7 @@ def _persist_session_git_meta(session: dict, cwd: str, generation: int) -> None:
|
||||
|
||||
def _run() -> None:
|
||||
try:
|
||||
branch = _git_branch_for_cwd(cwd)
|
||||
root = _git_common_repo_root_for_cwd(cwd)
|
||||
branch, root = _git_branch_for_cwd(cwd), _git_common_repo_root_for_cwd(cwd)
|
||||
if not (branch or root):
|
||||
return
|
||||
with _session_db(db_session) as db:
|
||||
@@ -546,8 +513,7 @@ def _set_session_cwd(session: dict, cwd: str) -> str:
|
||||
resolved = os.path.abspath(os.path.expanduser(cwd))
|
||||
if not os.path.isdir(resolved):
|
||||
raise ValueError(f"working directory does not exist: {cwd}")
|
||||
# An explicit user choice: persisted as the workspace (not the launch-dir fallback) and superseding any
|
||||
# settle-adopted cwd.
|
||||
# An explicit user choice: persisted as the workspace (not the launch-dir fallback), superseding a settle-adopted cwd.
|
||||
session.update(cwd=resolved, explicit_cwd=True, cwd_from_settle=False)
|
||||
_register_session_cwd(session)
|
||||
# The synchronous DB write claims ordering authority; git probes may publish only for that exact generation.
|
||||
|
||||
@@ -96,10 +96,7 @@ def record_turn_start(home: Path | str, session_key: str, prompt: str, *, attemp
|
||||
def clear_turn_marker(home: Path | str, session_key: str) -> None:
|
||||
"""Remove the marker once its turn concluded (any outcome the client saw)."""
|
||||
if session_key:
|
||||
_update(
|
||||
home, session_key,
|
||||
lambda e: {k: v for k, v in e.items() if k != session_key} if session_key in e else None, "clear",
|
||||
)
|
||||
_update(home, session_key, lambda e: {k: v for k, v in e.items() if k != session_key} if session_key in e else None, "clear")
|
||||
|
||||
|
||||
def read_turn_marker(home: Path | str, session_key: str) -> dict[str, Any] | None:
|
||||
|
||||
Reference in New Issue
Block a user