diff --git a/gateway/lifecycle_ledger.py b/gateway/lifecycle_ledger.py index e0defad934408..2bf0d4d82258f 100644 --- a/gateway/lifecycle_ledger.py +++ b/gateway/lifecycle_ledger.py @@ -41,6 +41,7 @@ import asyncio import json import logging +import math import os import re import subprocess @@ -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 @@ -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. diff --git a/gateway/restart.py b/gateway/restart.py index dcd8ff82783ea..80856256d5434 100644 --- a/gateway/restart.py +++ b/gateway/restart.py @@ -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) @@ -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: @@ -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: @@ -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( @@ -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. @@ -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: @@ -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( diff --git a/gateway/restart_loop_guard.py b/gateway/restart_loop_guard.py index cb3d23c891b7f..ecca2de0d3559 100644 --- a/gateway/restart_loop_guard.py +++ b/gateway/restart_loop_guard.py @@ -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, diff --git a/gateway/run.py b/gateway/run.py index 66e95c2207455..0297daf13001c 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -3309,6 +3309,8 @@ def _platform_has_bot_credential(platform: "Platform", platform_config: "Platfor read_launchd_exit_timeout_s, resolve_cron_drain_budget, resolve_launchd_capped_drain, + resolve_launchd_drain_deadline_s, + resolve_launchd_shutdown_watchdog_delay, resolve_replace_takeover_grace_s, ) @@ -7979,6 +7981,12 @@ def __init__(self, config: Optional[GatewayConfig] = None): self._launchd_exit_timeout_s = self._load_launchd_exit_timeout( self._restart_drain_timeout ) + try: + from gateway.lifecycle_ledger import read_last_teardown_seconds + + self._last_shutdown_teardown_s = read_last_teardown_seconds() + except Exception: + self._last_shutdown_teardown_s = None self._provider_routing = self._load_provider_routing() self._fallback_model = self._load_fallback_model() @@ -16394,9 +16402,15 @@ def _schedule_resume_pending_sessions(self, platform=None) -> int: from gateway import restart_loop_guard as _rlg _max_restarts, _window, _max_gap = self._restart_loop_guard_config() - _tripped = _rlg.check_and_record( - _max_restarts, _window, max_gap_seconds=_max_gap - ) + if not getattr(self, "_restart_loop_guard_recorded_this_boot", False): + _tripped = _rlg.check_and_record( + _max_restarts, _window, max_gap_seconds=_max_gap + ) + self._restart_loop_guard_recorded_this_boot = True + else: + _tripped = _rlg.is_restart_loop_tripped( + _max_restarts, _window, max_gap_seconds=_max_gap + ) if _tripped: # F2 is armed whenever the per-session breaker is enabled, # which (given the max(1, ...) clamp) it always is. The @@ -19642,17 +19656,33 @@ def _shutdown_watchdog_snapshot() -> dict: "restart_drain_timeout": self._restart_drain_timeout, "effective_drain_timeout": effective_stop_drain_timeout(self), "launchd_exit_timeout_s": getattr(self, "_launchd_exit_timeout_s", None), - "watchdog_delay_s": resolve_shutdown_watchdog_delay( - effective_stop_drain_timeout(self) + "watchdog_delay_s": resolve_launchd_shutdown_watchdog_delay( + resolve_shutdown_watchdog_delay( + effective_stop_drain_timeout(self) + ), + getattr(self, "_launchd_exit_timeout_s", None), + signal_driven=getattr( + self, "_stop_requested_by_signal", False + ), ), + "persistence_complete": False, "phase_elapsed_s": ( time.monotonic() - started if started is not None else None ), } if not os.environ.get("PYTEST_CURRENT_TEST"): + _watchdog_delay = resolve_shutdown_watchdog_delay( + effective_stop_drain_timeout(self) + ) + _launchd_budget = getattr(self, "_launchd_exit_timeout_s", None) + _watchdog_delay = resolve_launchd_shutdown_watchdog_delay( + _watchdog_delay, + _launchd_budget, + signal_driven=getattr(self, "_stop_requested_by_signal", False), + ) arm_shutdown_watchdog( - resolve_shutdown_watchdog_delay(effective_stop_drain_timeout(self)), + _watchdog_delay, done_event=_watchdog_done, snapshot_fn=_shutdown_watchdog_snapshot, exit_code=1, @@ -19759,15 +19789,29 @@ def _phase_elapsed() -> float: # Under launchd the real leash is launchd's own exit timeout, not # our watchdog (drain + grace): a signal-driven stop that lets # cron work push past it is SIGKILLed before cleanup runs. + # One absolute deadline governs BOTH allowances. ``timeout`` + # above is already ``exit_timeout - hard_exit_reserve - measured + # teardown``; deriving cron's ceiling from the same value is what + # stops in-flight cron work from holding the drain open past the + # teardown reserve this host has measured. A raw-clamp min() here + # would be the un-reserved wall (60s, not 60-10-teardown) and is + # never the binding term once the deadline is threaded through. _cron_leash = resolve_shutdown_watchdog_delay(timeout) _launchd_budget = getattr(self, "_launchd_exit_timeout_s", None) - if getattr(self, "_stop_requested_by_signal", False) and _launchd_budget: - _cron_leash = min(_cron_leash, float(_launchd_budget)) + _drain_deadline = ( + resolve_launchd_drain_deadline_s( + _launchd_budget, + last_teardown_s=getattr(self, "_last_shutdown_teardown_s", None), + ) + if getattr(self, "_stop_requested_by_signal", False) + else None + ) _cron_timeout = resolve_cron_drain_budget( timeout, _cron_drain_cfg, watchdog_delay=_cron_leash, elapsed=_phase_elapsed(), + deadline_s=_drain_deadline, ) if _cron_at_start and _cron_timeout > timeout: logger.info( @@ -19792,6 +19836,7 @@ def _phase_elapsed() -> float: timeout, _cron_timeout ) _drain_elapsed = time.monotonic() - _drain_started_at + _post_drain_started_at = time.monotonic() logger.info( "Shutdown phase: drain done at +%.2fs (drain took %.2fs, " "timed_out=%s, active_at_start=%d, active_now=%d, " @@ -20314,6 +20359,23 @@ def _phase_elapsed() -> float: else: self._update_runtime_status("stopped", self._exit_reason) _shutdown_gateway_health_export(self) + _teardown_elapsed = time.monotonic() - _post_drain_started_at + try: + from gateway.lifecycle_ledger import record_teardown_timing + + record_teardown_timing( + _teardown_elapsed, + total_shutdown_seconds=_phase_elapsed(), + drain_seconds=_drain_elapsed, + ) + except Exception as _e: + logger.debug("Failed to record shutdown teardown timing: %s", _e) + logger.info( + "Shutdown phase: post-drain teardown completed in %.2fs " + "(persistence complete; total %.2fs)", + _teardown_elapsed, + _phase_elapsed(), + ) logger.info("Gateway stopped (total teardown %.2fs)", _phase_elapsed()) self._stop_task = asyncio.create_task(_stop_impl()) diff --git a/tests/gateway/test_auto_continue_interrupted_turns.py b/tests/gateway/test_auto_continue_interrupted_turns.py index 4380f32d9c2a4..6467325566957 100644 --- a/tests/gateway/test_auto_continue_interrupted_turns.py +++ b/tests/gateway/test_auto_continue_interrupted_turns.py @@ -136,6 +136,30 @@ def _rowid_credits(tmp_path: Path) -> list: return json.loads(path.read_text())["attempts"] +@pytest.mark.asyncio +async def test_restart_loop_guard_records_at_most_once_per_gateway_boot( + tmp_path, monkeypatch +): + from gateway import restart_loop_guard + + runner, _adapter, db = _runner(tmp_path, monkeypatch) + entry = _entry(runner) + _mark_pending(runner, entry) + recorded = MagicMock(return_value=False) + inspected = MagicMock(return_value=False) + monkeypatch.setattr(restart_loop_guard, "check_and_record", recorded) + monkeypatch.setattr(restart_loop_guard, "is_restart_loop_tripped", inspected) + + runner._schedule_resume_pending_sessions() + await asyncio.gather(*runner._background_tasks) + runner._schedule_resume_pending_sessions() + await asyncio.gather(*runner._background_tasks) + + recorded.assert_called_once() + inspected.assert_called_once() + db.close() + + @pytest.mark.asyncio async def test_t1_prompt_default_keeps_note_bytes_and_adds_taxonomy_log( tmp_path, monkeypatch, caplog diff --git a/tests/gateway/test_cron_drain_teardown_reserve.py b/tests/gateway/test_cron_drain_teardown_reserve.py new file mode 100644 index 0000000000000..c1ffb9a5d7560 --- /dev/null +++ b/tests/gateway/test_cron_drain_teardown_reserve.py @@ -0,0 +1,328 @@ +"""Cron's drain allowance must come off the SAME deadline as the chat drain. + +The adaptive teardown reserve shrinks the chat drain by +``max(LAUNCHD_STOP_CLEANUP_RESERVE_S, last measured teardown)``. Cron's +allowance was never given ``last_teardown_s``: ``resolve_cron_drain_budget`` +subtracted only the fixed ``CRON_DRAIN_CLEANUP_RESERVE_S`` (10s), so on a host +that had MEASURED a long teardown the chat drain shrank while cron's deadline +stayed pinned at ``clamp - 10``: + + last_teardown | chat_drain | cron_due | left before SIGKILL | needed + None | 45.0 | 50.0 | 10.0 | 15.0 + 22.0 | 38.0 | 50.0 | 10.0 | 22.0 + 58.0 | 2.0 | 50.0 | 10.0 | 58.0 + +and ``_drain_active_agents._still_draining()`` returns True on the cron branch +independently of the chat ``deadline``, so in-flight cron work genuinely held +the drain open to ``cron_deadline`` after the shortened chat drain expired — +leaving less time before launchd's SIGKILL than the teardown this machine had +demonstrated it needs. + +Both allowances are now derived from one absolute deadline +(``resolve_launchd_drain_deadline_s``). These tests assert the COMPOSED budget +that ``gateway/run.py``'s stop path actually builds, because green isolated +helpers do not prove the composition: the pre-existing leash test passes an +unarmed ``watchdog_delay`` and never sees the reserve at all. + +No minimum-drain floor is asserted anywhere here. A deadline fully consumed by +a measured teardown is a real zero-second answer, not a value to be raised. +""" + +from __future__ import annotations + +import asyncio +from unittest.mock import AsyncMock, MagicMock, patch + +import pytest + +from gateway.restart import ( + CRON_DRAIN_CLEANUP_RESERVE_S, + LAUNCHD_STOP_CLEANUP_RESERVE_S, + effective_stop_drain_timeout, + resolve_cron_drain_budget, + resolve_launchd_drain_deadline_s, + resolve_launchd_shutdown_watchdog_delay, +) +from gateway.shutdown_watchdog import resolve_shutdown_watchdog_delay +from tests.gateway.restart_test_helpers import make_restart_runner + +CLAMP = 60.0 +CRON_CFG = 600.0 +TEARDOWNS = (None, 22.0, 35.0, 45.0, 58.0) + + +@pytest.fixture(autouse=True) +def _reset_cron_running_set(): + import cron.scheduler as sched + + sched._running_job_ids.clear() + sched._interrupted_job_ids.clear() + yield + sched._running_job_ids.clear() + sched._interrupted_job_ids.clear() + + +def _armed_watchdog(chat_drain: float, clamp: float) -> float: + """The hard-exit wall the stop path actually arms, per gateway/run.py.""" + return resolve_launchd_shutdown_watchdog_delay( + resolve_shutdown_watchdog_delay(chat_drain), + clamp, + signal_driven=True, + ) + + +# --------------------------------------------------------------------------- +# The shared deadline +# --------------------------------------------------------------------------- + + +class TestResolveLaunchdDrainDeadline: + def test_deadline_shrinks_with_the_measured_teardown(self): + # Cold host: the 15s floor reserve applies. + assert resolve_launchd_drain_deadline_s(CLAMP) == 35.0 + # A measured teardown longer than the floor widens the reserve. + assert resolve_launchd_drain_deadline_s(CLAMP, last_teardown_s=22.0) == 28.0 + assert resolve_launchd_drain_deadline_s(CLAMP, last_teardown_s=45.0) == 5.0 + + def test_a_shorter_measurement_never_shrinks_the_floor_reserve(self): + assert ( + resolve_launchd_drain_deadline_s(CLAMP, last_teardown_s=3.0) + == resolve_launchd_drain_deadline_s(CLAMP) + ) + + def test_saturating_teardown_yields_a_real_zero_not_a_floor(self): + """Legitimate zero-drain saturation: no minimum is imposed.""" + assert resolve_launchd_drain_deadline_s(CLAMP, last_teardown_s=58.0) == 0.0 + assert resolve_launchd_drain_deadline_s(CLAMP, last_teardown_s=900.0) == 0.0 + + def test_no_launchd_budget_imposes_no_deadline(self): + assert resolve_launchd_drain_deadline_s(None) is None + assert resolve_launchd_drain_deadline_s(0.0) is None + assert resolve_launchd_drain_deadline_s(-5.0) is None + assert resolve_launchd_drain_deadline_s("nope") is None # type: ignore[arg-type] + + def test_garbage_teardown_degrades_to_the_floor_reserve(self): + assert ( + resolve_launchd_drain_deadline_s(CLAMP, last_teardown_s="nope") # type: ignore[arg-type] + == resolve_launchd_drain_deadline_s(CLAMP) + ) + + +# --------------------------------------------------------------------------- +# Cron's allowance, bounded by that deadline +# --------------------------------------------------------------------------- + + +class TestCronBudgetHonoursTheDeadline: + @pytest.mark.parametrize("last_teardown", TEARDOWNS) + def test_cron_never_outlives_the_measured_teardown_reserve(self, last_teardown): + """The card's table, recomposed with the deadline threaded through.""" + elapsed = 1.0 + chat = effective_stop_drain_timeout( + _runner_like(drain=50.0, clamp=CLAMP, teardown=last_teardown) + ) + deadline = resolve_launchd_drain_deadline_s(CLAMP, last_teardown_s=last_teardown) + cron = resolve_cron_drain_budget( + chat, + CRON_CFG, + watchdog_delay=resolve_shutdown_watchdog_delay(chat), + elapsed=elapsed, + deadline_s=deadline, + ) + armed = _armed_watchdog(chat, CLAMP) + reserve = max(LAUNCHD_STOP_CLEANUP_RESERVE_S, last_teardown or 0.0) + assert deadline is not None, "a live launchd clamp must impose a deadline" + + # 1. Cron lands inside the one shared deadline. ``elapsed`` is time + # already spent before the drain, so what is left of the deadline is + # ``deadline - elapsed`` — clamped at 0 when the reserve already + # consumed the whole budget (a real zero, not a floor). + assert cron <= max(deadline - elapsed, 0.0) + 1e-9, ( + f"cron ({cron}) overran the remaining deadline " + f"({max(deadline - elapsed, 0.0)}) for teardown={last_teardown}" + ) + # 2. ...and the drain cannot outlive the chat drain's own budget. + assert cron <= max(chat, 0.0) + 1e-9 + + # 3. The consequence the incident is about: the window left between the + # end of cron work and the armed hard exit is the reserve — or all + # the time that exists, when the measurement saturates the budget. + window = armed - (cron + elapsed) + assert window >= min(reserve, armed - elapsed) - 1e-9, ( + f"teardown={last_teardown}: only {window}s before the hard exit, " + f"reserve is {reserve}s" + ) + + def test_the_old_fixed_reserve_formula_fails_this_invariant(self): + """Non-vacuity: without ``deadline_s`` the same row is short. + + Pins the regression itself — if the deadline stops being threaded + through, this is what the budget goes back to. + """ + elapsed = 1.0 + chat = effective_stop_drain_timeout( + _runner_like(drain=50.0, clamp=CLAMP, teardown=22.0) + ) + old = resolve_cron_drain_budget( + chat, + CRON_CFG, + watchdog_delay=min(resolve_shutdown_watchdog_delay(chat), CLAMP), + elapsed=elapsed, + # deadline_s omitted == the pre-fix call site + ) + assert old == CLAMP - CRON_DRAIN_CLEANUP_RESERVE_S - elapsed # 49.0, pinned + armed = _armed_watchdog(chat, CLAMP) + assert armed - (old + elapsed) < 22.0, "old formula must be SHORT of the reserve" + + def test_cron_floor_still_extends_a_short_drain_within_the_deadline(self): + """The #82161 intent is preserved: the floor only ever EXTENDS.""" + deadline = resolve_launchd_drain_deadline_s(CLAMP) # 35.0 + budget = resolve_cron_drain_budget( + 0.0, CRON_CFG, watchdog_delay=95.0, elapsed=0.0, deadline_s=deadline + ) + assert budget == 35.0 > 0.0 + + def test_configured_long_drain_is_never_shortened_below_the_deadline(self): + deadline = resolve_launchd_drain_deadline_s(CLAMP) # 35.0 + assert ( + resolve_cron_drain_budget( + 35.0, 0.0, watchdog_delay=95.0, elapsed=0.0, deadline_s=deadline + ) + == 35.0 + ) + + def test_no_deadline_keeps_the_pre_existing_behaviour(self): + with_none = resolve_cron_drain_budget( + 0.0, CRON_CFG, watchdog_delay=95.0, elapsed=0.0, deadline_s=None + ) + assert with_none == 95.0 - CRON_DRAIN_CLEANUP_RESERVE_S + + def test_saturated_deadline_gives_cron_zero_without_a_floor(self): + deadline = resolve_launchd_drain_deadline_s(CLAMP, last_teardown_s=58.0) + assert deadline == 0.0 + assert ( + resolve_cron_drain_budget( + 0.0, CRON_CFG, watchdog_delay=95.0, elapsed=1.0, deadline_s=deadline + ) + == 0.0 + ) + + +def _runner_like(*, drain: float, clamp: float | None, teardown: float | None): + from types import SimpleNamespace + + return SimpleNamespace( + _restart_drain_timeout=drain, + _launchd_exit_timeout_s=clamp, + _stop_requested_by_signal=True, + _last_shutdown_teardown_s=teardown, + ) + + +# --------------------------------------------------------------------------- +# The composed stop path (this is what the isolated helpers do not prove) +# --------------------------------------------------------------------------- + + +def _stop_runner(*, clamp: float | None, teardown: float | None, drain: float = 50.0): + runner, adapter = make_restart_runner() + adapter.disconnect = AsyncMock() + runner.session_store = MagicMock() + runner.session_store._entries = {} + runner.session_store.mark_resume_pending = MagicMock(return_value=True) + runner.session_store.clear_resume_pending = MagicMock(return_value=True) + runner._restart_drain_timeout = drain + runner._cron_drain_timeout = CRON_CFG + runner._launchd_exit_timeout_s = clamp + runner._stop_requested_by_signal = True + runner._last_shutdown_teardown_s = teardown + return runner + + +async def _composed_budgets(runner) -> tuple[float, float]: + """Run the real stop() and capture the (chat, cron) pair it composed.""" + captured: dict[str, float] = {} + real_drain = runner._drain_active_agents + + async def _spy(timeout, cron_timeout=None): + captured["chat"] = timeout + captured["cron"] = timeout if cron_timeout is None else cron_timeout + # Do not actually wait out the budget — the composition is the subject. + import cron.scheduler as sched + + sched._running_job_ids.clear() + return await real_drain(0.0, 0.0) + + runner._drain_active_agents = _spy + with patch("gateway.status.remove_pid_file"), patch("gateway.status.write_runtime_status"): + await runner.stop() + assert "cron" in captured, "stop() never reached the drain" + return captured["chat"], captured["cron"] + + +class TestComposedStopPath: + @pytest.mark.asyncio + @pytest.mark.parametrize("teardown", TEARDOWNS) + async def test_stop_path_leaves_the_measured_reserve_before_the_clamp(self, teardown): + """Active cron work + a live launchd clamp: the budget the stop path + actually hands ``_drain_active_agents`` must fit inside the deadline.""" + import cron.scheduler as sched + + sched._running_job_ids.add("in-flight-cron") + runner = _stop_runner(clamp=CLAMP, teardown=teardown) + + chat, cron = await _composed_budgets(runner) + + deadline = resolve_launchd_drain_deadline_s(CLAMP, last_teardown_s=teardown) + armed = _armed_watchdog(chat, CLAMP) + reserve = max(LAUNCHD_STOP_CLEANUP_RESERVE_S, teardown or 0.0) + assert deadline is not None, "a live launchd clamp must impose a deadline" + + assert chat <= deadline + 1e-9 + assert cron <= deadline + 1e-9, ( + f"teardown={teardown}: cron budget {cron}s overran the shared " + f"deadline {deadline}s — in-flight cron work can hold the drain " + f"past the reserve this host measured" + ) + assert armed - cron >= min(reserve, armed) - 1e-9 + + @pytest.mark.asyncio + async def test_cron_budget_is_pinned_to_the_deadline_not_the_raw_clamp(self): + """The binding term must be ``clamp - hard_exit_reserve - teardown``, + never the raw clamp: at teardown=22 the old call site produced 50.0.""" + import cron.scheduler as sched + + sched._running_job_ids.add("in-flight-cron") + runner = _stop_runner(clamp=CLAMP, teardown=22.0) + + _chat, cron = await _composed_budgets(runner) + + assert cron == pytest.approx(28.0, abs=0.5) + assert cron < CLAMP - CRON_DRAIN_CLEANUP_RESERVE_S + + @pytest.mark.asyncio + async def test_non_launchd_stop_keeps_the_configured_cron_floor(self): + """No clamp, no deadline: the #82161 floor behaviour is untouched.""" + import cron.scheduler as sched + + sched._running_job_ids.add("in-flight-cron") + runner = _stop_runner(clamp=None, teardown=22.0, drain=0.0) + + chat, cron = await _composed_budgets(runner) + + assert chat == 0.0 + assert cron > 0.0, "cron floor must still extend a zero chat drain" + + @pytest.mark.asyncio + async def test_in_band_restart_is_not_launchd_timed(self): + """A non-signal stop keeps the operator's configured drain.""" + import cron.scheduler as sched + + sched._running_job_ids.add("in-flight-cron") + runner = _stop_runner(clamp=CLAMP, teardown=45.0) + runner._stop_requested_by_signal = False + + chat, cron = await _composed_budgets(runner) + + assert chat == 50.0 + assert cron >= 50.0 diff --git a/tests/gateway/test_launchd_exit_timeout_drain_cap.py b/tests/gateway/test_launchd_exit_timeout_drain_cap.py index c2c1bbfe00e72..57bf1d191694a 100644 --- a/tests/gateway/test_launchd_exit_timeout_drain_cap.py +++ b/tests/gateway/test_launchd_exit_timeout_drain_cap.py @@ -12,6 +12,7 @@ from __future__ import annotations import subprocess +import threading import pytest @@ -21,9 +22,11 @@ parse_launchd_exit_timeout, read_launchd_exit_timeout_s, resolve_launchd_capped_drain, + resolve_launchd_shutdown_watchdog_delay, ) from gateway.shutdown_watchdog import ( DEFAULT_SHUTDOWN_WATCHDOG_GRACE_S, + arm_shutdown_watchdog, resolve_shutdown_watchdog_delay, ) @@ -159,9 +162,25 @@ def fake_run(argv, **_k): # --------------------------------------------------------------------------- -def test_capped_drain_fits_inside_launchd_budget_minus_reserve(): - # The incident shape: configured 180s, launchd clamps to 60s. - assert resolve_launchd_capped_drain(180.0, 60.0) == 60.0 - LAUNCHD_STOP_CLEANUP_RESERVE_S +@pytest.mark.parametrize( + ("configured", "last_teardown", "expected"), + [ + # clamp 60 - hard-exit reserve 10 - teardown reserve 15 = 35 + (50.0, None, 35.0), + (30.0, None, 30.0), + # clamp 60 - hard-exit reserve 10 - measured teardown 22 = 28 + (50.0, 22.0, 28.0), + ], +) +def test_capped_drain_preserves_teardown_headroom(configured, last_teardown, expected): + assert ( + resolve_launchd_capped_drain( + configured, + 60.0, + last_teardown_s=last_teardown, + ) + == expected + ) def test_capped_drain_never_extends_a_short_drain(): @@ -169,6 +188,84 @@ def test_capped_drain_never_extends_a_short_drain(): assert resolve_launchd_capped_drain(0.0, 60.0) == 0.0 +@pytest.mark.parametrize( + ("clamp", "configured", "last_teardown"), + [ + (60.0, 50.0, None), + (60.0, 50.0, 22.0), + (60.0, 180.0, None), + (50.0, 50.0, None), + (30.0, 50.0, None), + (60.0, 50.0, 35.0), + ], +) +def test_teardown_reserve_fits_before_the_watchdog_that_actually_fires( + clamp, configured, last_teardown +): + """The two reserves are additive, not overlapping. + + The drain cap protects a teardown window ending at launchd's SIGKILL, + but the process hard-exits ``LAUNCHD_HARD_EXIT_RESERVE_S`` earlier. The + reserve must therefore fit between the end of the drain and the *armed + watchdog*, otherwise os._exit lands mid-persistence — the state.db + corruption class this card exists to close. + """ + from gateway.restart import ( + LAUNCHD_STOP_CLEANUP_RESERVE_S, + resolve_launchd_shutdown_watchdog_delay, + ) + + drain = resolve_launchd_capped_drain( + configured, clamp, last_teardown_s=last_teardown + ) + armed = resolve_launchd_shutdown_watchdog_delay( + resolve_shutdown_watchdog_delay(drain), clamp, signal_driven=True + ) + promised_reserve = max(LAUNCHD_STOP_CLEANUP_RESERVE_S, last_teardown or 0.0) + assert armed - drain >= promised_reserve, ( + f"clamp={clamp} drain={drain} watchdog={armed}: only " + f"{armed - drain}s before hard exit for a {promised_reserve}s teardown" + ) + # The reserve is carved out of the drain, never borrowed past the wall. + assert drain + promised_reserve <= clamp + + +def test_budget_too_small_for_the_reserve_saturates_the_drain_to_zero(): + """When the clamp cannot hold the teardown, persistence wins, not the drain. + + This is honest degradation rather than a satisfiable window: an 8s + ``ExitTimeOut`` has no 15s teardown to protect, so the drain goes to + zero and the whole (short) budget is left for persistence. The property + that must survive is that the hard exit still lands before SIGKILL. + """ + from gateway.restart import resolve_launchd_shutdown_watchdog_delay + + for clamp in (8.0, 10.0, 4.0): + drain = resolve_launchd_capped_drain(50.0, clamp) + armed = resolve_launchd_shutdown_watchdog_delay( + resolve_shutdown_watchdog_delay(drain), clamp, signal_driven=True + ) + assert drain == 0.0 + assert 0.0 < armed < clamp + + +def test_short_launchd_budget_still_drains_and_persists(): + """A clamp at/below the hard-exit reserve must not mean 'exit immediately'. + + Returning a zero-second watchdog skips all drain and persistence work + instead of using the little time that genuinely exists. + """ + from gateway.restart import resolve_launchd_shutdown_watchdog_delay + + for clamp in (8.0, 10.0, 4.0): + armed = resolve_launchd_shutdown_watchdog_delay( + 95.0, clamp, signal_driven=True + ) + assert 0.0 < armed < clamp, ( + f"clamp={clamp}: watchdog {armed} leaves no time to persist" + ) + + def test_capped_drain_no_launchd_budget_returns_configured(): assert resolve_launchd_capped_drain(180.0, None) == 180.0 assert resolve_launchd_capped_drain(180.0, 0.0) == 180.0 @@ -200,7 +297,7 @@ def _runner(*, drain: float, launchd: float | None, by_signal: bool): def test_effective_drain_capped_only_for_signal_stops_under_launchd(): - assert _runner(drain=180.0, launchd=60.0, by_signal=True)._effective_stop_drain_timeout() == 50.0 + assert _runner(drain=180.0, launchd=60.0, by_signal=True)._effective_stop_drain_timeout() == 35.0 # In-band restart (SIGUSR1 → after-turn → stop()) is not launchd-timed. assert _runner(drain=180.0, launchd=60.0, by_signal=False)._effective_stop_drain_timeout() == 180.0 # Not launchd-owned (systemd, s6, foreground): configured drain stands. @@ -227,7 +324,7 @@ def test_effective_drain_getattr_guarded_for_bare_doubles(): _launchd_exit_timeout_s=60.0, ) ) - == 50.0 + == 35.0 ) @@ -268,11 +365,22 @@ def test_cron_leash_under_launchd_cannot_exceed_exit_timeout(): """The cron drain floor (#82161) is clamped to the launchd budget too. Without the clamp the cron ceiling is watchdog(drain+grace) - reserve, - which for a capped 50s drain is 100s — past launchd's 60s SIGKILL. + which for a capped drain is 95s — past launchd's 60s SIGKILL. + + SCOPE: this covers the LEASH term only, and ``<= 60.0`` is not the safety + property. The leash composition leaves exactly ``CRON_DRAIN_CLEANUP_RESERVE_S`` + (10s) before the clamp regardless of the teardown this host has measured, + so it is green while in-flight cron work can still be SIGKILLed mid-teardown. + The binding term on the real stop path is the shared absolute deadline — + see tests/gateway/test_cron_drain_teardown_reserve.py. """ - from gateway.restart import CRON_DRAIN_CLEANUP_RESERVE_S, resolve_cron_drain_budget + from gateway.restart import ( + CRON_DRAIN_CLEANUP_RESERVE_S, + resolve_cron_drain_budget, + resolve_launchd_drain_deadline_s, + ) - drain = resolve_launchd_capped_drain(180.0, 60.0) # 50 + drain = resolve_launchd_capped_drain(180.0, 60.0) leash = min(resolve_shutdown_watchdog_delay(drain), 60.0) budget = resolve_cron_drain_budget(drain, 600.0, watchdog_delay=leash, elapsed=0.0) assert budget == max(drain, 60.0 - CRON_DRAIN_CLEANUP_RESERVE_S) @@ -283,6 +391,72 @@ def test_cron_leash_under_launchd_cannot_exceed_exit_timeout(): ) assert unclamped == drain + DEFAULT_SHUTDOWN_WATCHDOG_GRACE_S - CRON_DRAIN_CLEANUP_RESERVE_S assert unclamped > 60.0 + # Reconciliation: the leash-only budget above is NOT the stop path's answer. + # Threading the shared deadline in is strictly tighter, and that is the + # term the teardown reserve rides on. + deadline = resolve_launchd_drain_deadline_s(60.0, last_teardown_s=22.0) + assert deadline is not None + composed = resolve_cron_drain_budget( + resolve_launchd_capped_drain(180.0, 60.0, last_teardown_s=22.0), + 600.0, + watchdog_delay=leash, + elapsed=0.0, + deadline_s=deadline, + ) + assert composed < budget + assert composed <= deadline + + +def test_launchd_shutdown_watchdog_hard_exits_before_supervisor_sigkill(): + assert ( + resolve_launchd_shutdown_watchdog_delay( + 95.0, + 60.0, + signal_driven=True, + ) + == 50.0 + ) + assert ( + resolve_launchd_shutdown_watchdog_delay( + 95.0, + 60.0, + signal_driven=False, + ) + == 95.0 + ) + + +def test_hard_exit_backstop_ignores_blocked_daemon_threads(monkeypatch, tmp_path): + from gateway import shutdown_watchdog + + release = threading.Event() + blockers = [ + threading.Thread(target=release.wait, daemon=True, name=f"blocked-{i}") + for i in range(25) + ] + for thread in blockers: + thread.start() + + exited = threading.Event() + codes: list[int] = [] + + def fake_exit(code): + codes.append(code) + exited.set() + + monkeypatch.setattr(shutdown_watchdog.os, "_exit", fake_exit) + arm_shutdown_watchdog( + 0.05, + exit_code=1, + dump_path=tmp_path / "watchdog.log", + ) + try: + assert exited.wait(1.0), "hard exit was blocked by unrelated daemon threads" + assert codes == [1] + finally: + release.set() + for thread in blockers: + thread.join(timeout=1.0) def test_signal_handler_marks_stop_as_signal_driven(): diff --git a/tests/gateway/test_lifecycle_ledger.py b/tests/gateway/test_lifecycle_ledger.py index 016a8b2bcd147..7b1be082a9fb8 100644 --- a/tests/gateway/test_lifecycle_ledger.py +++ b/tests/gateway/test_lifecycle_ledger.py @@ -20,8 +20,10 @@ detect_unclean_exit, get_lifecycle_sentinel_path, mark_exited, + read_last_teardown_seconds, read_prior_exit_label, record_startup, + record_teardown_timing, sample_memory, ) @@ -75,6 +77,56 @@ def test_sample_memory_has_expected_keys_on_linux() -> None: assert "mem_available_kib" in sample +# --------------------------------------------------------------------------- +# Teardown timing +# --------------------------------------------------------------------------- + + +def test_teardown_timing_round_trips_and_reaches_exit_diag(tmp_path: Path) -> None: + assert read_last_teardown_seconds(tmp_path) is None + + record_teardown_timing( + 18.25, + total_shutdown_seconds=48.5, + drain_seconds=30.0, + home=tmp_path, + ) + + assert read_last_teardown_seconds(tmp_path) == 18.25 + records = _exit_diag_records(tmp_path) + assert records[-1]["tag"] == "gateway.shutdown_teardown_timing" + assert records[-1]["teardown_seconds"] == 18.25 + assert records[-1]["total_shutdown_seconds"] == 48.5 + + +@pytest.mark.parametrize("bad", ["Infinity", "-Infinity", "NaN"]) +def test_non_finite_persisted_teardown_is_rejected(tmp_path: Path, bad: str) -> None: + """A corrupt file must not become an unbounded teardown reserve. + + ``inf`` passes a bare ``>= 0.0`` check, and the reserve flows straight + into the next shutdown's drain budget — an unbounded reserve silently + drives that budget to zero. Exercised through the real file path, not + the parameter, because the file is what feeds the live value. + """ + from gateway.lifecycle_ledger import get_teardown_timing_path + + path = get_teardown_timing_path(tmp_path) + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text(json.dumps({"teardown_seconds": float(bad)})) + + assert read_last_teardown_seconds(tmp_path) is None + + +def test_finite_persisted_teardown_still_survives_the_guard(tmp_path: Path) -> None: + from gateway.lifecycle_ledger import get_teardown_timing_path + + path = get_teardown_timing_path(tmp_path) + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text(json.dumps({"teardown_seconds": 22.0})) + + assert read_last_teardown_seconds(tmp_path) == 22.0 + + # --------------------------------------------------------------------------- # First boot / clean lifecycle # ---------------------------------------------------------------------------