Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
61 changes: 61 additions & 0 deletions gateway/lifecycle_ledger.py
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@
import asyncio
import json
import logging
import math
import os
import re
import subprocess
Expand All @@ -53,6 +54,7 @@
logger = logging.getLogger(__name__)

_LIFECYCLE_RELATIVE = ("state", "gateway.lifecycle.json")
_TEARDOWN_TIMING_RELATIVE = ("state", "gateway.teardown.json")
_EXIT_DIAG_RELATIVE = ("logs", "gateway-exit-diag.log")

# Total wall-clock budget for the kill-attribution probe. Boot must never
Expand Down Expand Up @@ -86,6 +88,65 @@ def get_lifecycle_sentinel_path(home: Optional[Path] = None) -> Path:
return base.joinpath(*_LIFECYCLE_RELATIVE)


def get_teardown_timing_path(home: Optional[Path] = None) -> Path:
"""Return the last completed post-drain teardown timing path."""
base = home if home is not None else _process_hermes_home()
return base.joinpath(*_TEARDOWN_TIMING_RELATIVE)


def read_last_teardown_seconds(home: Optional[Path] = None) -> Optional[float]:
"""Read the last completed post-drain teardown duration, if valid."""
data = _read_json(get_teardown_timing_path(home))
raw = (data or {}).get("teardown_seconds")
if raw is None:
return None
try:
value = float(raw)
except (TypeError, ValueError):
return None
# ``inf`` survives a bare ``>= 0.0`` check and is not a measurement any
# completed teardown can produce β€” it only arrives from a corrupt or
# hand-edited file. Letting it through makes the teardown reserve
# unbounded, which silently drives the next shutdown's drain budget to
# zero. ``nan`` already fails the comparison; reject both explicitly.
if not math.isfinite(value):
return None
return value if value >= 0.0 else None


def record_teardown_timing(
teardown_seconds: float,
*,
total_shutdown_seconds: float,
drain_seconds: float,
home: Optional[Path] = None,
) -> None:
"""Persist and diagnose one completed post-drain teardown measurement."""
try:
teardown = max(float(teardown_seconds), 0.0)
total = max(float(total_shutdown_seconds), 0.0)
drain = max(float(drain_seconds), 0.0)
except (TypeError, ValueError):
return
record = {
"ts": datetime.now(timezone.utc).isoformat(),
"tag": "gateway.shutdown_teardown_timing",
"pid": os.getpid(),
"teardown_seconds": teardown,
"total_shutdown_seconds": total,
"drain_seconds": drain,
}
path = get_teardown_timing_path(home)
try:
from utils import atomic_json_write

path.parent.mkdir(parents=True, exist_ok=True)
atomic_json_write(path, record, indent=None)
except Exception:
logger.debug("Failed to persist teardown timing", exc_info=True)
_append_exit_diag(record, home)


def sample_memory() -> Dict[str, Any]:
"""Cheap memory snapshot: own RSS + system availability + swap.

Expand Down
177 changes: 151 additions & 26 deletions gateway/restart.py
Original file line number Diff line number Diff line change
Expand Up @@ -70,7 +70,16 @@
# ``restart_drain_timeout: 180`` vs live 60 β†’ SIGKILL at +60s on every
# busy restart).
LAUNCHD_GUI_EXIT_TIMEOUT_CLAMP_S = 60
LAUNCHD_STOP_CLEANUP_RESERVE_S = 10.0
LAUNCHD_STOP_CLEANUP_RESERVE_S = 15.0
LAUNCHD_HARD_EXIT_RESERVE_S = 10.0

# Fallback reserve for a live ``ExitTimeOut`` at or below
# ``LAUNCHD_HARD_EXIT_RESERVE_S``. Subtracting the fixed reserve from such a
# budget yields a zero-second watchdog, i.e. hard-exit the instant shutdown
# starts β€” no drain, no persistence, strictly worse than the SIGKILL it is
# avoiding. A short budget keeps the same shape (most of it for work, a
# slice held back for launchd) at proportional scale.
LAUNCHD_SHORT_BUDGET_RESERVE_FRACTION = 0.25

_LAUNCHD_EXIT_TIMEOUT_RE = re.compile(r"^\s*exit timeout\s*=\s*(\d+)\s*$", re.MULTILINE)

Expand Down Expand Up @@ -136,22 +145,78 @@ def read_launchd_exit_timeout_s(
return parse_launchd_exit_timeout(getattr(proc, "stdout", ""))


def resolve_launchd_drain_deadline_s(
launchd_exit_timeout_s: float | None,
*,
cleanup_reserve_s: float = LAUNCHD_STOP_CLEANUP_RESERVE_S,
last_teardown_s: float | None = None,
hard_exit_reserve_s: float = LAUNCHD_HARD_EXIT_RESERVE_S,
) -> float | None:
"""Seconds from shutdown start by which EVERY drain must have returned.

The single absolute deadline the stop path is timed against under
launchd: the armed hard-exit wall (``exit_timeout -
hard_exit_reserve_s``) minus the teardown reserve this host has
actually measured (``max(cleanup_reserve_s, last_teardown_s)``).

Both drain allowances β€” the chat drain and the cron floor β€” are
derived from this one value, so no drain can consume time the
teardown after it has already been promised. ``None`` means no
launchd budget applies and no deadline is imposed; callers must treat
that as fail-open. The value may be ``0.0`` when the measured
teardown already saturates the budget: that is a real zero-drain
answer, not a floor to be raised.
"""

def _seconds(value: object) -> float:
try:
return max(float(value), 0.0) # type: ignore[arg-type]
except (TypeError, ValueError):
return 0.0

if launchd_exit_timeout_s is None:
return None
try:
budget = float(launchd_exit_timeout_s)
except (TypeError, ValueError):
return None
if budget <= 0.0:
return None
hard_exit = resolve_launchd_shutdown_watchdog_delay(
budget,
budget,
signal_driven=True,
hard_exit_reserve_s=hard_exit_reserve_s,
)
reserve = max(_seconds(cleanup_reserve_s), _seconds(last_teardown_s))
return max(hard_exit - reserve, 0.0)


def resolve_launchd_capped_drain(
drain_timeout: float,
launchd_exit_timeout_s: float | None,
*,
cleanup_reserve_s: float = LAUNCHD_STOP_CLEANUP_RESERVE_S,
last_teardown_s: float | None = None,
hard_exit_reserve_s: float = LAUNCHD_HARD_EXIT_RESERVE_S,
) -> float:
"""Clamp a SIGTERM-driven stop drain to what launchd will actually allow.

``launchd_exit_timeout_s`` is the live ``exit timeout`` for this job (see
:func:`read_launchd_exit_timeout_s`); ``None`` means no launchd budget
applies and the configured drain is returned untouched. Otherwise the
drain may use at most ``exit_timeout - cleanup_reserve_s`` so the
post-drain teardown (interrupt agents, disconnect adapters, checkpoint
and close SQLite) still completes before launchd escalates to SIGKILL.
Never *extends* the drain β€” an operator who configured a short one
keeps it.
applies and the configured drain is returned untouched. Never *extends*
the drain β€” an operator who configured a short one keeps it.

The two reserves are ADDITIVE, not overlapping. The process does not run
until launchd's SIGKILL at ``exit_timeout``; the shutdown watchdog
hard-exits ``hard_exit_reserve_s`` earlier (see
:func:`resolve_launchd_shutdown_watchdog_delay`). Sizing the teardown
window against the SIGKILL wall therefore over-promises by exactly that
reserve β€” at clamp 60 a 45s drain left only 5s before os._exit for a
teardown allocated 15s, and a recorded 22s teardown got 12s. The drain
may use at most ``exit_timeout - hard_exit_reserve_s -
max(cleanup_reserve_s, last_teardown_s)`` so the teardown completes
before the hard exit that actually fires.
"""

def _seconds(value: object) -> float:
Expand All @@ -161,16 +226,54 @@ def _seconds(value: object) -> float:
return 0.0

drain = _seconds(drain_timeout)
if launchd_exit_timeout_s is None:
deadline = resolve_launchd_drain_deadline_s(
launchd_exit_timeout_s,
cleanup_reserve_s=cleanup_reserve_s,
last_teardown_s=last_teardown_s,
hard_exit_reserve_s=hard_exit_reserve_s,
)
if deadline is None:
return drain
return min(drain, deadline)


def resolve_launchd_shutdown_watchdog_delay(
watchdog_delay_s: float,
launchd_exit_timeout_s: float | None,
*,
signal_driven: bool,
hard_exit_reserve_s: float = LAUNCHD_HARD_EXIT_RESERVE_S,
) -> float:
"""Return the hard-exit deadline for a launchd-timed shutdown.

Signal-driven shutdown must stop itself before launchd's uncatchable
SIGKILL. The final reserve leaves launchd headroom even when persistence
or adapter teardown wedges. Other shutdown paths keep their normal
watchdog leash.

When the live budget is at or below the fixed reserve (a short but valid
``ExitTimeOut`` such as 8s), subtracting it outright would return zero β€”
hard-exiting the instant shutdown starts, skipping all drain *and*
persistence rather than using the time that genuinely exists. Short
budgets therefore fall back to a proportional reserve, keeping a real
persistence window on the near side of SIGKILL.
"""
try:
budget = float(launchd_exit_timeout_s)
watchdog = max(float(watchdog_delay_s), 0.0)
except (TypeError, ValueError):
return drain
if budget <= 0.0:
return drain
cap = max(budget - _seconds(cleanup_reserve_s), 0.0)
return min(drain, cap)
watchdog = 0.0
if not signal_driven or launchd_exit_timeout_s is None:
return watchdog
try:
launchd_budget = float(launchd_exit_timeout_s)
reserve = max(float(hard_exit_reserve_s), 0.0)
except (TypeError, ValueError):
return watchdog
if launchd_budget <= 0.0:
return watchdog
if reserve >= launchd_budget:
reserve = launchd_budget * LAUNCHD_SHORT_BUDGET_RESERVE_FRACTION
return min(watchdog, max(launchd_budget - reserve, 0.0))


def effective_stop_drain_timeout(runner: object) -> float:
Expand All @@ -186,7 +289,11 @@ def effective_stop_drain_timeout(runner: object) -> float:
drain = getattr(runner, "_restart_drain_timeout", DEFAULT_GATEWAY_RESTART_DRAIN_TIMEOUT)
if not getattr(runner, "_stop_requested_by_signal", False):
return drain
return resolve_launchd_capped_drain(drain, getattr(runner, "_launchd_exit_timeout_s", None))
return resolve_launchd_capped_drain(
drain,
getattr(runner, "_launchd_exit_timeout_s", None),
last_teardown_s=getattr(runner, "_last_shutdown_teardown_s", None),
)


def is_gateway_supervisor_process(
Expand Down Expand Up @@ -272,6 +379,7 @@ def resolve_cron_drain_budget(
watchdog_delay: float,
elapsed: float = 0.0,
cleanup_reserve_s: float = CRON_DRAIN_CLEANUP_RESERVE_S,
deadline_s: float | None = None,
) -> float:
"""Seconds the shutdown drain may spend waiting on in-flight cron work.

Expand All @@ -284,9 +392,23 @@ def resolve_cron_drain_budget(
would swap a cleanly-interrupted job for a SIGKILL that leaves it
wedged mid-run β€” strictly worse than the bug being fixed.

Never returns less than ``drain_timeout``: the cron floor only ever
extends the wait, so an operator who deliberately configured a long
``restart_drain_timeout`` keeps it.
``deadline_s`` is the absolute stop-path deadline from
:func:`resolve_launchd_drain_deadline_s` β€” the hard-exit wall less the
teardown reserve this host has MEASURED. It is the same value the chat
drain is derived from, so both allowances come off one budget.
``cleanup_reserve_s`` is a fixed 10s guess; a host whose recorded
teardown is longer needs the longer reserve, and without this input
in-flight cron work held the drain open to the fixed-reserve ceiling
regardless (chat 45β†’2s as the measurement grew, cron pinned at 50s).
``None`` means no launchd budget applies and only the leash ceiling
governs.

Never returns less than ``drain_timeout`` *within the deadline*: the
cron floor only ever extends the wait, so an operator who deliberately
configured a long ``restart_drain_timeout`` keeps it β€” but it cannot
extend past time the teardown has already been promised. No minimum
floor is imposed; a deadline already consumed by the reserve yields a
real zero-second answer.
"""

def _seconds(value: object, fallback: float = 0.0) -> float:
Expand All @@ -297,14 +419,17 @@ def _seconds(value: object, fallback: float = 0.0) -> float:

drain = _seconds(drain_timeout)
floor = _seconds(cron_drain_timeout)
if floor <= 0.0:
return drain
ceiling = (
_seconds(watchdog_delay)
- _seconds(elapsed)
- _seconds(cleanup_reserve_s, CRON_DRAIN_CLEANUP_RESERVE_S)
)
return max(drain, min(floor, ceiling))
budget = drain
if floor > 0.0:
ceiling = (
_seconds(watchdog_delay)
- _seconds(elapsed)
- _seconds(cleanup_reserve_s, CRON_DRAIN_CLEANUP_RESERVE_S)
)
budget = max(drain, min(floor, ceiling))
if deadline_s is None:
return budget
return min(budget, max(_seconds(deadline_s) - _seconds(elapsed), 0.0))


def resolve_systemd_timeout_stop_sec(
Expand Down
16 changes: 8 additions & 8 deletions gateway/restart_loop_guard.py
Original file line number Diff line number Diff line change
Expand Up @@ -200,14 +200,14 @@ def check_and_record(
tripped = len(boots) >= max_restarts if max_restarts > 0 else False
if tripped:
logger.warning(
"Restart-loop breaker TRIPPED: %d chained restart-interrupted "
"gateway boots (no gap wider than %ds; threshold %d). The CALLER "
"decides what to skip β€” in the gateway the per-session replay "
"breaker and the per-session auto-resume cap own the break, so "
"healthy sessions keep resuming; grep the adjacent "
"'Restart-loop guard tripped at boot: deferred_to=' line for which "
"mechanism acted (#30719, #81642). If this is a false positive, "
"delete %s.",
"Restart-loop breaker TRIPPED: reason_class=restart_interrupted_boot_chain "
"%d chained restart-interrupted gateway boots (no gap wider than "
"%ds; threshold %d). The CALLER decides what to skip β€” in the "
"gateway the per-session replay breaker and the per-session "
"auto-resume cap own the break, so healthy sessions keep resuming; "
"grep the adjacent 'Restart-loop guard tripped at boot: "
"deferred_to=' line for which mechanism acted (#30719, #81642). "
"If this is a false positive, delete %s.",
len(boots),
int(_chain_gap(window_seconds, max_gap_seconds)),
max_restarts,
Expand Down
Loading
Loading