From 1136f135dd640fff082af7859386637d347dbcd2 Mon Sep 17 00:00:00 2001 From: teknium1 <127238744+teknium1@users.noreply.github.com> Date: Wed, 23 Sep 2026 17:23:28 +0000 Subject: [PATCH] fix(gateway): crash recovery never redelivers a suppressed reply; tz-safe turn marker Review follow-ups on the crash-left reply adoption: - A persisted reply is now judged the way live delivery would have judged it. A bare silence marker ([SILENT] / SILENT / NO_REPLY ...) on an internal turn, or the reply to a diagnostic wake whose chat policy mutes diagnostics (read in the routed profile's scope, as the adapter does), is owed nothing: the marker is cleared and nothing is sent or resumed. A human turn's bare silence marker becomes the same "returned only a silence marker" notice the live path sends. Before, the raw "NO_REPLY" reached the user as a "Recovered reply". - The active-turn marker start is written as aware UTC and compared as epoch seconds (startup adoption and recover_interrupted_turns). A naive local wall clock read by a process in another zone (DST, container vs unit TZ) was hours off: a fresh in-flight turn was dropped as stale, or a previous turn's reply could be adopted as this one's. updated_at stays naive local for the older binary's recency heuristic; a pre-upgrade naive marker still reads as local time. - Reply timestamps go through coerce_epoch instead of float(), so one odd transcript row cannot abort the whole recovery pass. --- gateway/run_startup.py | 75 ++++++++++++------- gateway/session_lifecycle.py | 30 ++++---- tests/gateway/test_active_turn_recovery.py | 70 ++++++++++++++++- .../gateway-session-lifecycle.md | 11 ++- 4 files changed, 143 insertions(+), 43 deletions(-) diff --git a/gateway/run_startup.py b/gateway/run_startup.py index e7f46ac769..743335d9e6 100644 --- a/gateway/run_startup.py +++ b/gateway/run_startup.py @@ -14,7 +14,7 @@ import logging import os import signal import time -from contextlib import suppress +from contextlib import nullcontext, suppress from contextvars import copy_context from datetime import datetime from pathlib import Path @@ -723,14 +723,12 @@ class GatewayStartupMixin: return resumed, ledgered async def _ledger_crash_left_replies(self, max_age_seconds: int) -> int: - """Hand every marked turn whose final reply was persisted to the delivery ledger and clear - its marker, so the boot sweep delivers the stored reply instead of auto-resume - regenerating it. Without the ledger such turns stay marked and resume.""" + """Settle every marked turn whose final reply was persisted and clear its marker, so + auto-resume does not regenerate it: a reply live delivery would have suppressed is owed + nothing, any other goes to the delivery ledger for the boot sweep. Without the ledger a + presentable reply stays marked and resumes.""" from gateway.delivery_ledger import compute_obligation_id, ledger_enabled, record_crash_left_reply - from gateway.platforms.base import _strip_media_directives - from gateway.run import _sanitize_gateway_final_response - if not await asyncio.to_thread(ledger_enabled): - return 0 + ledger_on = await asyncio.to_thread(ledger_enabled) cutoff = time.time() - max_age_seconds # older markers are cleared, never acted on with self.session_store._lock: # noqa: SLF001 — snapshot under lock self.session_store._ensure_loaded_locked() # noqa: SLF001 @@ -742,27 +740,54 @@ class GatewayStartupMixin: ] ledgered = 0 for key, session_id, token, started_at, origin, profile in marked: - history = await self.async_session_store.load_transcript(session_id) - last = next((m for m in reversed(history) if m.get("role") not in ("session_meta", "system")), None) - if (started_at.timestamp() < cutoff or not last or last.get("role") != "assistant" - or last.get("tool_calls") - or not isinstance(last.get("content"), str) - or float(last.get("timestamp") or 0) < started_at.timestamp()): - continue # the turn never produced its final reply: it resumes - text = _strip_media_directives( - _sanitize_gateway_final_response(origin.platform, last["content"])).strip() - if not text: + started = started_at.timestamp() # aware UTC marker; a pre-upgrade naive one reads as local + if started < cutoff: continue - await asyncio.to_thread( - record_crash_left_reply, - obligation_id=compute_obligation_id(key, f"crash:{token}", text), session_key=key, - platform=str(getattr(origin.platform, "value", origin.platform)), chat_id=origin.chat_id, - thread_id=origin.thread_id, content=text, since=started_at.timestamp(), - adapter_profile=profile) - if await self.async_session_store.clear_turn_active(key, token): + text = self._crash_left_reply(await self.async_session_store.load_transcript(session_id), + started, origin) + if text is None or (text and not ledger_on): + continue # no final reply to deliver: the turn resumes + if text: + await asyncio.to_thread( + record_crash_left_reply, + obligation_id=compute_obligation_id(key, f"crash:{token}", text), session_key=key, + platform=str(getattr(origin.platform, "value", origin.platform)), chat_id=origin.chat_id, + thread_id=origin.thread_id, content=text, since=started, adapter_profile=profile) + if await self.async_session_store.clear_turn_active(key, token) and text: ledgered += 1 return ledgered + def _crash_left_reply(self, history: list, started: float, origin) -> Optional[str]: + """What a crash-left turn owes, judged as live delivery would have: ``None`` when it never + persisted a final reply after *started*; ``""`` when nothing would have been presented (a + silence marker on a machinery turn, a muted diagnostic wake); else the text to send, with a + human turn's bare silence marker replaced by the same notice the live path sends.""" + from gateway.platforms.base import _strip_media_directives + from gateway.response_filters import is_intentional_silence_response, is_machinery_display_kind + from gateway.run import _sanitize_gateway_final_response + from gateway.run_turn import _UNEXPECTED_SILENCE_REPLY + from gateway.warning_notifications import diagnostic_turn_muted + from hermes_cli.timefmt import coerce_epoch + visible = [m for m in history if m.get("role") not in ("session_meta", "system")] + last = visible[-1] if visible else {} + if (last.get("role") != "assistant" or last.get("tool_calls") or not isinstance(last.get("content"), str) + or (coerce_epoch(last.get("timestamp")) or 0) < started): + return None + prompt = next((m for m in reversed(visible) if m.get("role") == "user"), {}) + machinery = is_machinery_display_kind(prompt.get("display_kind")) + if machinery: + try: # the owning profile's display policy, as the adapter reads it at delivery + scope = self._media_delivery_scope_for_source(origin) + except Exception: + logger.debug("Crash-left reply: no routed scope for %s", origin.chat_id, exc_info=True) + scope = nullcontext() + with scope: + if diagnostic_turn_muted(prompt.get("display_metadata"), origin.platform): + return "" + if is_intentional_silence_response(last["content"]): + return "" if machinery else _UNEXPECTED_SILENCE_REPLY + return _strip_media_directives(_sanitize_gateway_final_response(origin.platform, last["content"])).strip() or None + @staticmethod def _start_hosted_room_worker_sync(): """Start the local Group Chat worker without importing the dashboard.""" diff --git a/gateway/session_lifecycle.py b/gateway/session_lifecycle.py index e1815afde7..1c70f89daf 100644 --- a/gateway/session_lifecycle.py +++ b/gateway/session_lifecycle.py @@ -4,8 +4,9 @@ from __future__ import annotations import logging import os +import time import uuid -from datetime import datetime, timedelta +from datetime import datetime, timedelta, timezone from typing import TYPE_CHECKING, Optional from hermes_state_ids import new_session_id @@ -117,15 +118,16 @@ class SessionLifecycleMixin: candidate = entry.to_dict() candidate["active_turn_token"] = token candidate["active_turn_started_at"] = _iso(started_at) - if started_at is not None: + touched = _now() if started_at is not None else None + if touched is not None: # Keeps the legacy 120s startup heuristic working for an older binary during a rolling # downgrade/upgrade window. - candidate["updated_at"] = started_at.isoformat() + candidate["updated_at"] = touched.isoformat() self._save_entry(session_key, entry_data=candidate, lock_held=True) entry.active_turn_token = token entry.active_turn_started_at = started_at - if started_at is not None: - entry.updated_at = started_at + if touched is not None: + entry.updated_at = touched def mark_turn_active(self, session_key: str) -> Optional[str]: """Persist exact ownership of the running agent turn; returns the opaque token for @@ -136,7 +138,9 @@ class SessionLifecycleMixin: entry = self._entry_locked(session_key) if entry is None: return None - self._set_turn_marker_locked(session_key, entry, token, _now()) + # Aware UTC, unlike the local wall clock elsewhere: the next process compares it with + # epoch transcript timestamps and may run in another zone (DST, container vs unit TZ). + self._set_turn_marker_locked(session_key, entry, token, datetime.now(timezone.utc)) return token def clear_turn_active(self, session_key: str, token: str) -> bool: @@ -153,8 +157,7 @@ class SessionLifecycleMixin: """Promote crash-left turn markers into ``resume_pending`` (unclean startup only). Old/invalid markers are cleared without resuming; suspended sessions are never re-armed. Returns the number of newly promoted sessions.""" - now = _now() - max_age = timedelta(seconds=max(0, max_age_seconds)) + now, epoch_now = _now(), time.time() promoted = 0 def _promote(entry: SessionEntry) -> bool: @@ -162,13 +165,10 @@ class SessionLifecycleMixin: if not entry.active_turn_token: return False started_at = entry.active_turn_started_at - try: - marker_is_stale = started_at is None or ( - max_age_seconds > 0 and now - started_at > max_age - ) - except TypeError: - # Mixed aware/naive timestamps: clear rather than risk an unsafe old resume. - marker_is_stale = True + # Epoch arithmetic: a pre-upgrade naive marker reads as local time, an aware one exactly. + marker_is_stale = started_at is None or ( + max_age_seconds > 0 and epoch_now - started_at.timestamp() > max_age_seconds + ) if not marker_is_stale and not entry.suspended: if entry.resume_pending: # A drain-timeout marker is more specific; keep it. diff --git a/tests/gateway/test_active_turn_recovery.py b/tests/gateway/test_active_turn_recovery.py index 3c2ecd1bdb..8c781ec2f9 100644 --- a/tests/gateway/test_active_turn_recovery.py +++ b/tests/gateway/test_active_turn_recovery.py @@ -6,7 +6,10 @@ marker, compare-and-swap cleanup, and promotion into the existing ``resume_pending`` recovery path after an unclean exit. """ +import os +import time from datetime import datetime, timedelta +from pathlib import Path from types import SimpleNamespace from typing import Any, cast from unittest.mock import AsyncMock, MagicMock, PropertyMock, patch @@ -464,12 +467,12 @@ def _db_runner(tmp_path) -> tuple[GatewayRunner, SessionStore]: return runner, runner.session_store -def _turn(store: SessionStore, chat_id: str, *, marked: bool, reply: str | None) -> SessionSource: +def _turn(store: SessionStore, chat_id: str, *, marked: bool, reply: str | None, **prompt: Any) -> SessionSource: source = _make_source(chat_id) entry = store.get_or_create_session(source) if marked: store.mark_turn_active(entry.session_key) - store.append_to_transcript(entry.session_id, {"role": "user", "content": f"question {chat_id}"}) + store.append_to_transcript(entry.session_id, {"role": "user", "content": f"question {chat_id}", **prompt}) if reply is not None: store.append_to_transcript(entry.session_id, {"role": "assistant", "content": reply}) return source @@ -508,3 +511,66 @@ async def test_unclean_restart_delivers_a_persisted_unledgered_reply_instead_of_ assert [(r["content"], r["needs_marker"], r["chat_id"], r["thread_id"]) for r in rows] == [ ("the stored answer", True, "replied", "thread-1")] _close_store_db(store) + + +_WAKE = {"display_kind": "internal_notification"} + + +@pytest.mark.asyncio +@pytest.mark.parametrize(("reply", "prompt", "owed"), [ + ("[SILENT]", _WAKE, []), + ("NO_REPLY", _WAKE, []), + ("disk is 91% full", {**_WAKE, "display_metadata": {"notification_category": "diagnostic"}}, []), + ("NO_REPLY", {}, ["⚠️ The model returned only a silence marker for a message that needed a reply. " + "Try again or rephrase."]), +]) +async def test_unclean_restart_never_redelivers_a_reply_live_delivery_suppressed(tmp_path, reply, prompt, owed): + """A crash-left reply is owed exactly what live delivery would have sent: nothing for a silence + marker on a machinery turn or a muted diagnostic wake (and the finished turn is not resumed), the + unexpected-silence notice for a human turn, never the raw marker.""" + from gateway.delivery_ledger import sweep_recoverable + + (Path(os.environ["HERMES_HOME"]) / "config.yaml").write_text("display: {suppress_warning_notifications: true}\n", encoding="utf-8") + runner, store = _db_runner(tmp_path) + source = _turn(store, "quiet", marked=True, reply=reply, **prompt) + + assert await runner._recover_unclean_sessions() == (0, len(owed)) + + entry = _entry_for(store, source) + assert (entry.resume_pending, entry.active_turn_token) == (False, None) + assert [r["content"] for r in sweep_recoverable(deliverable_platforms={"discord"})] == owed + _close_store_db(store) + + +@pytest.mark.asyncio +@pytest.mark.skipif(not hasattr(time, "tzset"), reason="needs a POSIX process timezone switch") +async def test_turn_marker_start_survives_a_timezone_change_across_the_crash(tmp_path): + """The dead process's local zone is not the new one's (DST, container vs unit TZ): the marked + turn still resumes, and the previous turn's answer persisted a minute before it is not re-sent + (a naive wall-clock marker read in the new zone was 7 h off and dropped the turn as stale).""" + from gateway.delivery_ledger import sweep_recoverable + + runner, store = _db_runner(tmp_path) + source = _make_source("tz") + entry = store.get_or_create_session(source) + store.append_to_transcript(entry.session_id, {"role": "user", "content": "earlier question"}) + store.append_to_transcript(entry.session_id, {"role": "assistant", "content": "earlier answer", + "timestamp": time.time() - 60}) + original_tz = os.environ.get("TZ") + try: + os.environ["TZ"] = "Etc/GMT+7" # UTC-7 when the turn starts ... + time.tzset() + store.mark_turn_active(entry.session_key) + os.environ["TZ"] = "UTC" # ... UTC when the gateway comes back + time.tzset() + assert await runner._recover_unclean_sessions() == (1, 0) + finally: + if original_tz is None: + os.environ.pop("TZ", None) + else: + os.environ["TZ"] = original_tz + time.tzset() + + assert _entry_for(store, source).resume_pending is True + assert sweep_recoverable(deliverable_platforms={"discord"}) == [] + _close_store_db(store) diff --git a/website/docs/developer-guide/gateway-session-lifecycle.md b/website/docs/developer-guide/gateway-session-lifecycle.md index d92802bc52..b9edf3f9c6 100644 --- a/website/docs/developer-guide/gateway-session-lifecycle.md +++ b/website/docs/developer-guide/gateway-session-lifecycle.md @@ -366,10 +366,19 @@ one of two things: - **The reply was persisted but never ledgered.** The stored transcript reply is recorded as an unowned ledger row and the marker is cleared; the boot sweep delivers it once with the - "Recovered reply" notice. The turn is not regenerated. + "Recovered reply" notice. The turn is not regenerated. The reply is judged the way live + delivery would have judged it: a bare silence marker (`[SILENT]`, `NO_REPLY`, ...) on an + internal turn, or the reply to a diagnostic wake the chat's policy mutes, is owed nothing + (the marker is cleared, nothing is sent or resumed). A human turn's bare silence marker + becomes the same "returned only a silence marker" notice the live path sends. - **No reply was persisted.** `recover_interrupted_turns()` sets `resume_pending=True`, `resume_reason="restart_interrupted"`, and the turn auto-resumes once. +The marker's start time is stored as aware UTC and compared as epoch seconds, so a restart +in a different local zone (DST change, container vs. unit `TZ`) neither drops a fresh marker +as stale nor adopts the previous turn's reply as this one's. A marker written by an older +build (naive local time) is read as host-local time. + A turn already in the ledger is redelivered by the ledger sweep, which also clears any `resume_pending` for that session, so it is never both delivered and re-answered.