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.
This commit is contained in:
@@ -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."""
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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.
|
||||
|
||||
|
||||
Reference in New Issue
Block a user