""" prod/tests/test_watchdog_decision.py ===================================== Comprehensive tests for fix (b): the 2026-09-16 04:10:40 *ghost-subscription* watchdog self-restart in nautilus_event_trader.py::_scan_watchdog_loop. What (b) does ------------- The scan watchdog already *detected* the r27 wedge (HZ latest_eigen_scan key frozen at 13638 while no scan was accepted), but its "upstream dark" branch only printed a reminder ("NO SCANS ... UNMANAGED"). (b) promotes that branch to self-restart via the existing _watchdog_restart -> os._exit(86) -> supervisord respawn path, gated on: * acc_age >= UPSTREAM_DARK_RESTART_S (900s) -- confirmed long dark * uptime_ok (past WATCHDOG_RESTART_MIN_UPTIME_S=600s warm-up) * probe is NOT None -- frozen-KEY case, NOT the dead-HZ-client case (which the probe-None-3x path at ~2773 owns) Why these tests don't `import nautilus_event_trader` --------------------------------------------------- nautilus_event_trader.py does a module-level `from nautilus_dolphin.nautilus.proxy_boost_engine import create_d_liq_engine` and companion imports that BLOCK at import time outside the live supervisord env (verified: `import nautilus_event_trader` hangs in this shell). Importing the heavy kernel would therefore poison the test session. Per TESTING_DOCTRINE ("prefer the real kernel; mock only true externalities") the true externalities here ARE the engine/Hazelcast/CH stacks -- and they do not participate in the watchdog *decision*. So (b)'s decision was extracted into the dependency-free module `prod/watchdog_decision.py` (only stdlib `math`), which the live loop now calls. These tests: * unit-test the pure predicate `upstream_dark_restart` (all branches + edges + poison + warm-up) -- the core of the fix -- * drive a *faithful* stub of `_scan_watchdog_loop`'s per-tick decision (verbatim transcription of nautilus 2726-2821) that calls the REAL `scan_watchdog_dark_restart` seam, over N time-skipped ticks (unit-to-unit @ the seam + E2E shape) -- * source-integrity-pin the LIVE file as text (no import) so the stub can never drift from the loop that actually runs. Mutation litmus (doctrine) is run externally (see README in this file's dir): flipping `>=`->`>` / removing the warm-up gate / removing the probe-None guard each fails a test below (boundary / warmup / dead-probe tests respectively). """ from __future__ import annotations import re import sys from dataclasses import dataclass, field from pathlib import Path from typing import Optional import pytest PROD_DIR = "/mnt/dolphinng5_predict/prod" if PROD_DIR not in sys.path: sys.path.insert(0, PROD_DIR) NAUTILUS_PATH = str(Path(PROD_DIR) / "nautilus_event_trader.py") from watchdog_decision import ( # noqa: E402 (sys.path set above) UPSTREAM_DARK_RESTART_S, upstream_dark_restart, scan_watchdog_dark_restart, ) # Cadences mirrored from nautilus_event_trader.py (296-304). Pinned by # TestSourceIntegrity::test_live_cadences_match so a divergence is caught. THRESH = UPSTREAM_DARK_RESTART_S # 900.0 SCAN_STALL_S = 120.0 WARMUP_S = 600.0 PROBE_INTERVAL_S = 30.0 LOG_EVERY_S = 300.0 WATCHDOG_EXIT_CODE = 86 # --------------------------------------------------------------------------- # UNIT: the (b) predicate -- all branches + edges + poison + warm-up # --------------------------------------------------------------------------- class TestUpstreamDarkRestartPredicate: """upstream_dark_restart(acc_age_s, uptime_ok, scan_number_probe) -> bool.""" @pytest.mark.parametrize("acc_age", [900.0, 900.1, 1000.0, 5000.0, float("inf")]) def test_restarts_when_frozen_past_threshold(self, acc_age): # real scan-number probe (frozen key) + uptime elapsed + acc>=900s assert upstream_dark_restart(acc_age, True, 13638) is True @pytest.mark.parametrize("acc_age", [0.0, 1.0, 120.0, 300.0, 899.0, 899.9, 899.999]) def test_no_restart_below_threshold(self, acc_age): assert upstream_dark_restart(acc_age, True, 13638) is False def test_boundary_is_inclusive(self): # `>=` (not `>`): exactly at the threshold MUST restart. # MUTATION GUARD: flipping `>=` -> `>` makes this False -> test fails. assert upstream_dark_restart(THRESH, True, 13638) is True def test_warm_up_forbids_restart(self): # MUTATION GUARD: removing the `if not uptime_ok: return False` gate # makes this True -> test fails. assert upstream_dark_restart(THRESH * 10, False, 13638) is False @pytest.mark.parametrize( "probe", [ None, # dead HZ client (probe-None-3x path owns it) "13638", # corrupt: JSON parse produced a str, not int b"13638", # corrupt: bytes {"scan_number": 13638}, # corrupt: dict [13638], # corrupt: list float("nan"), # corrupt probe float("inf"), # corrupt probe float("-inf"), # corrupt probe ], ) def test_dead_hz_client_and_corrupt_probe_excluded(self, probe): # (b) must NOT fire for a None/corrupt probe -- those are owned by the # probe-None-3x path or are defensive no-ops. Must return False, NEVER # raise. assert upstream_dark_restart(THRESH * 10, True, probe) is False def test_probe_zero_is_a_real_scan_number(self): # 0 is a valid int scan_number; the loop's `if not raw: return None` # maps 0->None upstream, but IF the predicate ever sees 0 it must treat # it as a real (frozen) key, not a miss. assert upstream_dark_restart(THRESH * 10, True, 0) is True @pytest.mark.parametrize("acc_age", [-1.0, -1000.0, float("nan"), float("-inf")]) def test_corrupt_or_backward_clock_no_restart(self, acc_age): assert upstream_dark_restart(acc_age, True, 13638) is False def test_infinite_age_is_really_dead(self): # +inf acc_age -> nan>=x is False but inf>=x is True -> restart. assert upstream_dark_restart(float("inf"), True, 13638) is True def test_non_numeric_acc_age_never_raises(self): # Defensive: a corrupt acc_age must not crash the watchdog thread. assert upstream_dark_restart("900", True, 13638) is False assert upstream_dark_restart(None, True, 13638) is False assert upstream_dark_restart(900.0, True, {"x": 1}) is False def test_return_type_is_bool(self): # Guard against `return acc_age_s >= THRESH` leaking non-bool for weird # inputs (e.g. a numpy type); downstream `if reason:` must be clean. assert isinstance(upstream_dark_restart(1000.0, True, 13638), bool) assert isinstance(upstream_dark_restart(100.0, True, 13638), bool) # --------------------------------------------------------------------------- # UNIT: the (b) branch seam (reason-string contract the live loop fires) # --------------------------------------------------------------------------- class TestScanWatchdogDarkRestart: def test_returns_reason_when_predicate_true(self): r = scan_watchdog_dark_restart(1000.0, True, 13638, ev_age_s=600.0) assert r is not None assert "frozen at 13638" in r assert ">= 900" in r assert "ghost-subscription after WS reconnect" in r assert "no reader liveness" in r assert "acc_age=1000" in r assert "ev_age=600" in r def test_returns_none_below_threshold(self): assert scan_watchdog_dark_restart(100.0, True, 13638, 0.0) is None def test_returns_none_during_warm_up(self): assert scan_watchdog_dark_restart(THRESH * 10, False, 13638, 0.0) is None def test_returns_none_for_dead_hz_probe(self): # probe None -> 3x-streak path owns it; (b) seam must yield None. assert scan_watchdog_dark_restart(THRESH * 10, True, None, 0.0) is None assert scan_watchdog_dark_restart(THRESH * 10, True, "bad", 0.0) is None def test_ev_age_flows_through(self): r = scan_watchdog_dark_restart(1000.0, True, 13638, ev_age_s=1234.0) assert "ev_age=1234" in r def test_boundary_inclusive_returns_reason(self): # `>=` boundary: acc == THRESH -> reason (not None). MUTATION GUARD. assert scan_watchdog_dark_restart(THRESH, True, 13638, 0.0) is not None # --------------------------------------------------------------------------- # Faithful stub of _scan_watchdog_loop's per-tick decision (nautilus 2726-2821) # --------------------------------------------------------------------------- @dataclass class _ProbeState: """Loop-local mutable state for the stub watchdog tick.""" last_probe_num: Optional[int] = None last_probe_ts: float = 0.0 last_dark_log_ts: float = 0.0 dupes_at_stall: Optional[int] = None probe_fail_streak: int = 0 @dataclass class _StubTrader: """Duck-typed stand-in for DolphinLiveTrader exposing ONLY what the watchdog decision touches (engine/hazelcast/CH are the true externalities and do not enter the decision). _watchdog_restart is captured (NOT os._exit-ed).""" boot_ts: float accept_ts: float event_ts: float scan_number: int probe_seq: list = field(default_factory=list) dupe_drops_total: int = 0 logs: list = field(default_factory=list) # captured reminder/log prints restarts: list = field(default_factory=list) # captured _watchdog_restart reasons _probe_calls: int = 0 def _probe_latest_scan_number(self): if self._probe_calls < len(self.probe_seq): v = self.probe_seq[self._probe_calls] self._probe_calls += 1 return v return self.scan_number # frozen key once the explicit seq is exhausted def _watchdog_restart(self, reason): # Real loop calls os._exit here; stub captures only. self.restarts.append(reason) def _watchdog_tick(state: _ProbeState, trader: _StubTrader, now: float) -> Optional[str]: """Verbatim transcription of nautilus_event_trader._scan_watchdog_loop body (2726-2821). Only the (b) dark-restart decision delegates to the REAL ``scan_watchdog_dark_restart`` seam; all other branches are transcribed from the live file so the stub pins the documented behaviour.""" acc_age = now - trader.accept_ts ev_age = now - trader.event_ts uptime_ok = (now - trader.boot_ts) > WARMUP_S if acc_age < SCAN_STALL_S: state.last_probe_num = None state.dupes_at_stall = None state.probe_fail_streak = 0 return None if ev_age < SCAN_STALL_S: if state.dupes_at_stall is None: state.dupes_at_stall = trader.dupe_drops_total return None if trader.dupe_drops_total > state.dupes_at_stall: if now - state.last_dark_log_ts > LOG_EVERY_S: state.last_dark_log_ts = now trader.logs.append("dupe") elif uptime_ok: trader.restarts.append( f"scan worker stalled {acc_age:.0f}s with events still arriving") return None probe = trader._probe_latest_scan_number() if probe is None: state.probe_fail_streak += 1 if state.probe_fail_streak >= 3 and uptime_ok: trader.restarts.append(f"HZ probe failed {state.probe_fail_streak}x") else: state.probe_fail_streak = 0 if probe is not None: if state.last_probe_num is None: state.last_probe_num = probe state.last_probe_ts = now elif (now - state.last_probe_ts) >= PROBE_INTERVAL_S: if probe != state.last_probe_num and uptime_ok: trader.restarts.append( f"listener deaf: HZ latest_eigen_scan advanced " f"{state.last_probe_num} -> {probe}") state.last_probe_num = probe state.last_probe_ts = now # ----- (b) seam: ghost-subscription frozen-key restart ----- reason = scan_watchdog_dark_restart(acc_age, uptime_ok, probe, ev_age) if reason: trader.restarts.append(reason) # ----- dark-log reminder (only when (b) did NOT fire) ----- if reason is None and now - state.last_dark_log_ts > LOG_EVERY_S: state.last_dark_log_ts = now trader.logs.append("dark-log") return reason def _drive(ticks): """Drive N ticks: ticks = [(now, )...] for a frozen-key trader. Returns state.""" t0 = ticks[0] - 300.0 # accept/event happen 300s before first tick trader = _StubTrader(boot_ts=t0 - 1000.0, accept_ts=t0, event_ts=t0, scan_number=13638, probe_seq=[13638] * len(ticks)) state = _ProbeState() for now in ticks: _watchdog_tick(state, trader, now) return trader, state # --------------------------------------------------------------------------- # E2E-shape: frozen features_map key + time-skipped loop -> _watchdog_restart # --------------------------------------------------------------------------- class TestGhostSubscriptionE2E: def test_restarts_at_900s_boundary(self): trader, _ = _drive([300, 600, 899, 900]) # acc_age crosses 900 at tick 4 ghosts = [r for r in trader.restarts if "ghost-subscription" in r] assert len(ghosts) == 1 assert "frozen at 13638" in ghosts[0] def test_no_restart_before_900s(self): trader, _ = _drive([300, 600, 899]) assert trader.restarts == [] def test_dark_log_reminder_while_waiting(self): # acc_age in [300,600): (b) off, dark-log reminder fires every 300s. trader, _ = _drive([300, 600]) assert trader.restarts == [] assert "dark-log" in trader.logs def test_probe_none_3x_self_restart_not_double_fired(self): # Dead HZ client: probe None 3x with stale events + uptime -> the # probe-None-3x path restarts; (b) must NOT also fire for None. t0 = 1000.0 trader = _StubTrader(boot_ts=t0 - 1000, accept_ts=t0 - 1000, event_ts=t0 - 1000, scan_number=13638, probe_seq=[None, None, None]) state = _ProbeState() for now in [t0 + 300, t0 + 330, t0 + 360]: # 3 ticks, 30s apart, past stall _watchdog_tick(state, trader, now) probes = [r for r in trader.restarts if r.startswith("HZ probe failed")] ghosts = [r for r in trader.restarts if "ghost-subscription" in r] assert len(probes) == 1 assert ghosts == [] # (b) correctly suppressed for None probe def test_listener_deaf_restarts_on_key_advance(self): t0 = 1000.0 trader = _StubTrader(boot_ts=t0 - 1000, accept_ts=t0, event_ts=t0 - 1000, scan_number=13639, probe_seq=[13638, 13639]) # key advances once @ +30s state = _ProbeState() ticks = [t0 + 300, t0 + 330] # 1st: baseline; 2nd: advance -> listener deaf for now in ticks: _watchdog_tick(state, trader, now) deafs = [r for r in trader.restarts if r.startswith("listener deaf")] assert len(deafs) == 1 assert "13638 -> 13639" in deafs[0] def test_accept_fresh_is_idle(self): # acc_age < SCAN_STALL_S: loop skips probing entirely -> no restart. t0 = 1000.0 trader = _StubTrader(boot_ts=t0 - 1000, accept_ts=t0, event_ts=t0, scan_number=13638, probe_seq=[13638]) state = _ProbeState() _watchdog_tick(state, trader, t0 + 10) # acc_age=10s assert trader.restarts == [] def test_first_probe_only_sets_baseline(self): # First non-None probe sets last_probe_num without a listener-deaf restart. t0 = 1000.0 trader = _StubTrader(boot_ts=t0 - 1000, accept_ts=t0, event_ts=t0 - 1000, scan_number=13638, probe_seq=[13638, 13638, 13638]) state = _ProbeState() for now in [t0 + 300, t0 + 330, t0 + 360]: # frozen key, 30s apart _watchdog_tick(state, trader, now) assert not any(r.startswith("listener deaf") for r in trader.restarts) assert not any("ghost-subscription" in r for r in trader.restarts) # acc<900 # --------------------------------------------------------------------------- # SOURCE-INTEGRITY: the LIVE file must wire the seam (text-only, no import) # --------------------------------------------------------------------------- def _nautilus_src() -> str: return Path(NAUTILUS_PATH).read_text(encoding="utf-8") class TestSourceIntegrity: def test_seam_module_imported_by_live_kernel(self): src = _nautilus_src() assert "from watchdog_decision import (" in src assert "UPSTREAM_DARK_RESTART_S" in src assert "scan_watchdog_dark_restart" in src def test_b_branch_lives_inside_scan_watchdog_loop(self): src = _nautilus_src() i = src.index(" def _scan_watchdog_loop(self):") j = src.index("\n def ", i + 5) # next method def loop = src[i:j] assert "scan_watchdog_dark_restart(" in loop assert "_dark_restart_reason" in loop assert "self._watchdog_restart(_dark_restart_reason)" in loop # (b) must sit AFTER the listener-deaf block and BEFORE the dark-log # reminder print. assert loop.index("listener deaf: HZ latest_eigen_scan advanced") \ < loop.index("scan_watchdog_dark_restart(") assert loop.index("scan_watchdog_dark_restart(") \ < loop.index("UNMANAGED until scans resume") def test_reminder_print_preserved(self): # acc_age < 900s window still reminder-logs (not restart-only). src = _nautilus_src() assert "WATCHDOG: NO SCANS for" in src assert "UNMANAGED until scans resume" in src def test_preexisting_restart_branches_untouched(self): # (b) adds ONE branch; the old restart paths must remain intact. src = _nautilus_src() assert "HZ probe failed" in src # probe-None 3x -> restart (~2773) assert "listener deaf: HZ latest_eigen_scan advanced" in src # ~2784 assert "scan worker stalled" in src # worker-stuck -> restart def test_b_not_duplicated(self): src = _nautilus_src() assert src.count("scan_watchdog_dark_restart(") == 1 def test_live_cadences_match_stub(self): # Pin the stub's hardcoded cadences to the LIVE constant values so the # faithful stub-tick can't silently drift from the running kernel. src = _nautilus_src() def _val(name): m = re.search(rf"\b{name}\s*=\s*([0-9.]+)", src) assert m, f"{name} not found in live kernel" return float(m.group(1)) assert _val("SCAN_STALL_S") == SCAN_STALL_S assert _val("WATCHDOG_RESTART_MIN_UPTIME_S") == WARMUP_S assert _val("WATCHDOG_PROBE_INTERVAL_S") == PROBE_INTERVAL_S assert _val("UPSTREAM_DARK_LOG_EVERY_S") == LOG_EVERY_S assert _val("WATCHDOG_EXIT_CODE") == WATCHDOG_EXIT_CODE # nautilus imports UPSTREAM_DARK_RESTART_S from the light module; # the live value is pinned by test_seam_module_imported_by_live_kernel # + the light-module constant itself (see below). assert THRESH == 900.0 def test_live_uptime_guard_exists_before_b(self): # The warm-up gate (uptime_ok from _PROCESS_BOOT_TS) must gate the (b) # restart; confirm the live loop computes uptime_ok and passes it to the # seam (not a bare call). src = _nautilus_src() i = src.index(" def _scan_watchdog_loop(self):") j = src.index("\n def ", i + 5) loop = src[i:j] assert "uptime_ok = (now - _PROCESS_BOOT_TS)" in loop assert "scan_watchdog_dark_restart(" in loop def test_live_restart_calls_exit(self): # Sanity: _watchdog_restart actually exits with the watchdog code so a # (b) decision kills the process for supervisord to respawn. assert "os._exit(WATCHDOG_EXIT_CODE)" in _nautilus_src()