diff --git a/prod/nautilus_event_trader.py b/prod/nautilus_event_trader.py index bb36844..95d03d7 100644 --- a/prod/nautilus_event_trader.py +++ b/prod/nautilus_event_trader.py @@ -34,6 +34,10 @@ from nautilus_dolphin.nautilus.esf_alpha_orchestrator import NDPosition from nautilus_dolphin.nautilus.adaptive_circuit_breaker import AdaptiveCircuitBreaker from nautilus_dolphin.nautilus.ob_features import OBFeatureEngine from nautilus_dolphin.nautilus.ob_provider import MockOBProvider +from watchdog_decision import ( + UPSTREAM_DARK_RESTART_S, # (b) 900s: frozen HZ scan_number key -> self-restart + scan_watchdog_dark_restart, # pure (b) seam: acc/uptime/probe/ev -> reason|None +) from nautilus_dolphin.nautilus.esof_size_gate import ( parse_esof_payload, esof_gate_from_payload, esof_score_from_payload, esof_size_mult_from_score, ESOF_STALE_FALLBACK_MULT, ESOF_FRESHNESS_S, @@ -41,6 +45,8 @@ from nautilus_dolphin.nautilus.esof_size_gate import ( from prod.clean_arch.adapters.eigen_scan_normalizer import normalize_ng7_scan from prod.clean_arch.obf_tp_observation import inject_obf_midprice from prod.clean_arch.tp_curve import compute_our_leverage, compute_soft_tp_pct +from adaptive_exit.maras_persistence import extract_maras_context, merge_bundle_with_maras +from adaptive_exit.trade_event_schema import normalize_trade_event_type try: sys.path.insert(0, '/mnt/dolphinng5_predict/Observability') from esof_advisor import compute_esof as _compute_esof_inline @@ -68,9 +74,10 @@ try: except Exception: BounceAdvisor = None try: - from adaptive_exit.post_win_long_overlay import PostWinExecutionFSM + from adaptive_exit.post_win_long_overlay import PostWinExecutionFSM, PostWinExecutionFSMConfig except Exception: PostWinExecutionFSM = None + PostWinExecutionFSMConfig = None try: from nautilus_dolphin.nautilus.alpha_exit_v7_engine import AlphaExitEngineV7, TradeContextV7 except Exception: @@ -302,6 +309,12 @@ WATCHDOG_EXIT_CODE = 86 # Scanner restarts reset scan_number to 0. A backwards jump larger than this # is a restart (accept + re-anchor ratchet), not a stale duplicate (drop). SCAN_NUMBER_RESET_GAP = 1000 +# A scanner counter reset SMALLER than the gap above (e.g. NG7 reset ~263 on +# 2026-06-22) is otherwise read as stale duplicates -> BLUE drops every scan and +# dark-freezes until manually restarted. Re-anchor after this many CONSECUTIVE +# sub-high-water scans: a genuine stale duplicate is a one-off; a scanner reset +# is a sustained run. +SCAN_RESET_CONSEC = 3 def _trade_log_paths(ts_dt: datetime) -> tuple[str, str]: log_date = ts_dt.strftime("%Y%m%d") @@ -400,6 +413,7 @@ class DolphinLiveTrader: self._dedup_lock = threading.Lock() # guards atomic check-and-set on last_scan_number self._scan_executor = ThreadPoolExecutor(max_workers=1, thread_name_prefix="scan") self.last_scan_number = -1 + self._consec_backward_drops = 0 # consecutive sub-high-water scans (small scanner-reset detector) # Scan-flow watchdog state. Event ts proves the HZ listener is alive; # accept ts proves the worker thread is draining; the dupe counter # separates "worker stuck" from "upstream flooding duplicates". @@ -473,7 +487,15 @@ class DolphinLiveTrader: self.trade_direction: int = _direction_from_env() self.vol_p60_threshold: float = _vol_p60_threshold_from_env() self._runtime_direction: int = self.trade_direction - self._efsm = PostWinExecutionFSM() if PostWinExecutionFSM is not None else None + self._efsm = ( + PostWinExecutionFSM( + PostWinExecutionFSMConfig( + require_context_gate=True, + ) + ) + if PostWinExecutionFSM is not None and PostWinExecutionFSMConfig is not None + else None + ) self._trade_announcement_center = None self._processed_retract_commands: deque = deque(maxlen=5000) self._processed_retract_set: set[str] = set() @@ -514,19 +536,51 @@ class DolphinLiveTrader: raw = self.features_map.blocking().get("maras_latest") if not raw: return {} - payload = json.loads(raw) if isinstance(raw, str) else raw - if not isinstance(payload, dict): - return {} - return { - "composite_hash": payload.get("composite_hash", payload.get("hash", 0)), - "scalar_hash": payload.get("scalar_hash", 0), - "regime": payload.get("regime", ""), - "final_score": payload.get("final_score", 0.0), - "confidence": payload.get("confidence", 0.0), - } + return extract_maras_context(raw) except Exception: return {} + def _efsm_long_overlay_context(self, *, vel_div: float, prices_dict: dict[str, float]) -> dict[str, object]: + """Build the live context required before arming a post-win LONG overlay.""" + btc_price = float(prices_dict.get("BTCUSDT") or 0.0) + btc_window = [float(x) for x in self.btc_prices if math.isfinite(float(x))] + btc_not_in_freefall = False + btc_guard_reason = "btc_context_unavailable" + btc_drawdown_frac = None + btc_window_return_frac = None + if btc_window: + first = btc_window[0] + last = btc_window[-1] + peak = max(btc_window) + if first > 0 and peak > 0: + btc_window_return_frac = (last - first) / first + btc_drawdown_frac = (peak - last) / peak + btc_not_in_freefall = ( + btc_drawdown_frac <= 0.015 + and btc_window_return_frac >= -0.01 + and last > min(btc_window) * 1.001 + ) + btc_guard_reason = "btc_ok" if btc_not_in_freefall else "btc_freefall" + + maras_ctx = self._latest_maras_context() + return { + "vel_div_now": float(vel_div or 0.0), + "btc_price": btc_price, + "btc_not_in_freefall": btc_not_in_freefall, + "btc_guard_reason": btc_guard_reason, + "btc_drawdown_frac": btc_drawdown_frac, + "btc_window_return_frac": btc_window_return_frac, + "maras_composite_hash": maras_ctx.get("maras_composite_hash", 0), + "maras_scalar_hash": maras_ctx.get("maras_scalar_hash", 0), + "maras_confidence": maras_ctx.get("maras_confidence", 0.0), + "maras_conflict_level": maras_ctx.get("maras_conflict_level", 0.0), + "maras_final_score": maras_ctx.get("maras_final_score", 0.0), + "maras_tier_eigen": maras_ctx.get("maras_tier_eigen", 0.0), + "maras_tier_btc": maras_ctx.get("maras_tier_btc", 0.0), + "maras_tier_esof": maras_ctx.get("maras_tier_esof", 0.0), + "maras_tier_micro": maras_ctx.get("maras_tier_micro", 0.0), + } + def _resolve_runtime_direction(self) -> int: """Resolve active trade direction for the next eligible entry.""" base = int(self.trade_direction) @@ -771,13 +825,88 @@ class DolphinLiveTrader: bundle = {} if self._market_state_runtime is not None and getattr(self._market_state_runtime, "latest_bundle_dict", None): bundle = dict(self._market_state_runtime.latest_bundle_dict) + maras_ctx = self._latest_maras_context() + bundle = merge_bundle_with_maras(bundle, maras_ctx) return { "tp_base_pct": float(self._tp_base_pct), "tp_effective_pct": float(tp_effective_pct), "our_leverage": float(our_leverage), "market_state_bundle_json": json.dumps(bundle, default=str, sort_keys=True) if bundle else "{}", + **maras_ctx, } + def _trade_event_payload( + self, + *, + pending: Mapping[str, Any], + event_type: str, + event_id: str, + exit_reason: str, + exit_price: float, + quantity: float, + pnl: float, + pnl_pct: float, + capital_before: float, + capital_after: float, + bars_held: int, + scan_uuid: str = "", + exit_leg_id: str = "", + exit_seq: int = 0, + exit_notional: float = 0.0, + remaining_notional: float = 0.0, + remaining_qty: float = 0.0, + pnl_leg: float = 0.0, + pnl_realized_total: float = 0.0, + ) -> dict[str, Any]: + event = { + "ts": _ch_ts_us(), + "date": pending.get("entry_date", self.current_day or ""), + "strategy": "blue", + "trade_id": str(pending.get("trade_id", "") or ""), + "asset": str(pending.get("asset", "") or ""), + "side": str(pending.get("side", "") or ""), + "entry_price": float(pending.get("entry_price", 0.0) or 0.0), + "exit_price": float(exit_price), + "quantity": float(quantity), + "capital_before": float(capital_before), + "capital_after": float(capital_after), + "pnl": float(pnl), + "pnl_pct": float(pnl_pct), + "exit_reason": str(exit_reason or ""), + "vel_div_entry": float(pending.get("vel_div_entry", 0.0) or 0.0), + "boost_at_entry": float(pending.get("boost_at_entry", 0.0) or 0.0), + "beta_at_entry": float(pending.get("beta_at_entry", 0.0) or 0.0), + "posture": str(pending.get("posture", "") or ""), + "leverage": float(pending.get("leverage", 0.0) or 0.0), + "bars_held": max(0, int(bars_held or 0)), + "regime_signal": 0, + "tp_threshold": float(self.eng.exit_manager.fixed_tp_pct), + "execution_quality_json": "", + "market_state_bundle_json": str(pending.get("market_state_bundle_json", "") or ""), + "tp_base_pct": float(pending.get("tp_base_pct", 0.0) or 0.0), + "tp_effective_pct": float(pending.get("tp_effective_pct", 0.0) or 0.0), + "our_leverage": float(pending.get("our_leverage", 0.0) or 0.0), + "event_type": normalize_trade_event_type(event_type), + "event_id": str(event_id or ""), + "chain_root_trade_id": str(pending.get("chain_root_trade_id", pending.get("trade_id", "")) or ""), + "chain_head_leg_id": str(pending.get("chain_head_leg_id", f"{pending.get('trade_id', '')}:open") or ""), + "chain_prev_leg_id": str(pending.get("chain_prev_leg_id", "") or ""), + "chain_seq": int(pending.get("chain_seq", pending.get("retraction_legs", 0)) or 0), + "chain_token": str(pending.get("chain_token", "") or ""), + "chain_mode": str(pending.get("chain_mode", "LIVE") or "LIVE"), + "exit_leg_id": str(exit_leg_id or ""), + "exit_seq": int(exit_seq or 0), + "retraction_legs": int(pending.get("retraction_legs", 0) or 0), + "retraction_realized_total": float(pending.get("realized_pnl_legs_total", 0.0) or 0.0), + "pnl_leg": float(pnl_leg), + "pnl_realized_total": float(pnl_realized_total), + "exit_notional": float(exit_notional), + "remaining_notional": float(remaining_notional), + "remaining_qty": float(remaining_qty), + "scan_uuid": str(scan_uuid or pending.get("scan_uuid", "") or ""), + } + return event + def _sync_tp_threshold(self) -> None: """Read live TP threshold from HZ control plane and propagate to engine. @@ -2662,6 +2791,20 @@ class DolphinLiveTrader: f"{last_probe_num} → {probe} but no events for {ev_age:.0f}s") last_probe_num = probe last_probe_ts = now + _dark_restart_reason = scan_watchdog_dark_restart( + acc_age, uptime_ok, probe, ev_age) + if _dark_restart_reason: + # (b) 2026-09-16 04:10:40 ghost-subscription wedge recovery: + # HZ latest_eigen_scan frozen (probe == last_probe_num, NOT None) + # for >= UPSTREAM_DARK_RESTART_S with uptime elapsed. The 0408 WS + # reconnect re-subscribed+ACKed but the server never resumed the + # stream -> reader blocks in poll()/recv() w/ no liveness watchdog + # -> key never advances -> silent starvation (r27 py-spy: all 19 + # threads idle, only timed waits ticking). Self-restart -> + # supervisord respawn w/ fresh WS session. Safe: r27 venue flat + # (zero fills, capital intact). probe None owned by 3x-streak + # path (~2773) -> scan_watchdog_dark_restart returns None for None. + self._watchdog_restart(_dark_restart_reason) if now - last_dark_log_ts > UPSTREAM_DARK_LOG_EVERY_S: last_dark_log_ts = now print(f"[{datetime.now(timezone.utc).isoformat()}] " @@ -3701,6 +3844,37 @@ class DolphinLiveTrader: "pnl_realized_total": float(pending.get("realized_pnl_legs_total", 0.0) or 0.0), "bars_held": bars_held, }) + ch_put("trade_events", self._trade_event_payload( + pending={ + **pending, + "trade_id": tid, + "chain_root_trade_id": str(chain_state.get("chain_root_trade_id", tid) or tid), + "chain_head_leg_id": str(chain_state.get("chain_head_leg_id", leg_id) or leg_id), + "chain_prev_leg_id": str(chain_state.get("chain_prev_leg_id", "") or ""), + "chain_seq": int(chain_state.get("chain_seq", leg_seq) or leg_seq), + "chain_token": str(chain_state.get("chain_token", "") or ""), + "chain_mode": str(chain_state.get("chain_mode", "LIVE") or "LIVE"), + "retraction_legs": int(pending.get("retraction_legs", 0) or 0), + "realized_pnl_legs_total": float(pending.get("realized_pnl_legs_total", 0.0) or 0.0), + }, + event_type="PARTIAL_EXIT", + event_id=leg_id, + exit_reason=str(cmd.get("reason", "RETRACT")), + exit_price=current_price, + quantity=remaining_qty, + pnl=net_pnl_leg, + pnl_pct=pnl_pct_now, + capital_before=capital_before, + capital_after=capital_after, + bars_held=bars_held, + exit_leg_id=leg_id, + exit_seq=leg_seq, + exit_notional=reduce_notional, + remaining_notional=remaining_notional, + remaining_qty=remaining_qty, + pnl_leg=net_pnl_leg, + pnl_realized_total=float(pending.get("realized_pnl_legs_total", 0.0) or 0.0), + )) ch_put("trade_reconstruction", { "ts": _ch_ts_us(), "trade_id": tid, @@ -3912,12 +4086,22 @@ class DolphinLiveTrader: # scan until manually restarted (near-miss on 2026-06-09/10). with self._dedup_lock: if scan_number > 0 and scan_number <= self.last_scan_number: - if scan_number < self.last_scan_number - SCAN_NUMBER_RESET_GAP: + backwards_gap = self.last_scan_number - scan_number + self._consec_backward_drops += 1 + # Re-anchor on a LARGE jump (full reset to ~0) OR a SUSTAINED run + # of smaller sub-high-water scans (scanner reset by < gap — the + # 2026-06-22 dark-freeze: NG7 reset ~263 < 1000, every scan dropped). + if (scan_number < self.last_scan_number - SCAN_NUMBER_RESET_GAP + or self._consec_backward_drops >= SCAN_RESET_CONSEC): log(f"WARN scanner restart detected: scan_number {self.last_scan_number} → " - f"{scan_number} — re-anchoring dedup ratchet") + f"{scan_number} (gap={backwards_gap}, consec={self._consec_backward_drops}) " + f"— re-anchoring dedup ratchet") + self._consec_backward_drops = 0 else: self._dupe_drops_total += 1 return + else: + self._consec_backward_drops = 0 self.last_scan_number = scan_number self._last_scan_accept_ts = time.time() self.scans_processed += 1 @@ -4056,10 +4240,14 @@ class DolphinLiveTrader: efsm_decision = None overlay_flip = False if self._efsm is not None and int(e.get('direction', -1)) == 1 and int(self.trade_direction) == -1: + efsm_context = self._efsm_long_overlay_context(vel_div=vel_div, prices_dict=prices_dict) efsm_decision = self._efsm.tag_next_entry( asset=str(e.get('asset', '') or ''), entry_ts=datetime.now(timezone.utc), - metadata={"trade_id": tid}, + metadata={ + "trade_id": tid, + **efsm_context, + }, ) overlay_flip = bool(efsm_decision and efsm_decision.action == "TAG" and efsm_decision.side == "LONG") self._pending_entries[tid] = { @@ -4485,6 +4673,7 @@ class DolphinLiveTrader: log(f" MarketStateRuntime outcome update failed for {tid}: {e}") if self._efsm is not None: try: + efsm_context = self._efsm_long_overlay_context(vel_div=vel_div, prices_dict=prices_dict) _efsm_out = self._efsm.observe_closed_trade( trade_id=str(tid or ""), asset=str(pending.get("asset", "") or ""), @@ -4494,7 +4683,10 @@ class DolphinLiveTrader: leverage=float(pending.get("leverage", 0) or 0), closed_ts=datetime.now(timezone.utc), was_overlay_flip=bool(pending.get("overlay_flip", False)), - metadata={"exit_reason": str(x.get("reason", "UNKNOWN"))}, + metadata={ + "exit_reason": str(x.get("reason", "UNKNOWN")), + **efsm_context, + }, ) if _efsm_out.action in {"ARMED", "TAG", "RESET"}: log(f"EFSM { _efsm_out.action }: { _efsm_out.to_dict() }") @@ -4538,37 +4730,33 @@ class DolphinLiveTrader: ) self._persist_trade_execution_quality(execution_quality) pending.update(self._tp_curve_context(notional=float(pending.get("notional", 0) or 0))) - ch_put("trade_events", { - "ts": _ch_ts_us(), - "date": pending['entry_date'], - "strategy": "blue", - "trade_id": tid, - "asset": pending['asset'], - "side": pending['side'], - "entry_price": pending['entry_price'], - "exit_price": exit_price, - "quantity": pending['quantity'], - "capital_before": capital_before, - "capital_after": capital_after, - "pnl": realized_pnl, - "pnl_pct": float(x.get('pnl_pct', 0) or 0), - "exit_reason": str(x.get('reason', 'UNKNOWN')), - "vel_div_entry": pending['vel_div_entry'], - "boost_at_entry": pending['boost_at_entry'], - "beta_at_entry": pending['beta_at_entry'], - "posture": pending['posture'], - "leverage": pending['leverage'], - # CH column is UInt16 — a negative value poisons the spool - # (head-of-line jam, incident 2026-06-12: bars_held=-106) - "bars_held": max(0, int(x.get('bars_held', 0) or 0)), - "regime_signal": 0, - "tp_threshold": float(self.eng.exit_manager.fixed_tp_pct), - "execution_quality_json": json.dumps(execution_quality, default=str), - "market_state_bundle_json": str(pending.get("market_state_bundle_json", "") or ""), - "tp_base_pct": float(pending.get("tp_base_pct", 0.0) or 0.0), - "tp_effective_pct": float(pending.get("tp_effective_pct", 0.0) or 0.0), - "our_leverage": float(pending.get("our_leverage", 0.0) or 0.0), - }) + te = self._trade_event_payload( + pending={ + **pending, + "trade_id": tid, + "chain_root_trade_id": str(pending.get("chain_root_trade_id", tid) or tid), + "chain_head_leg_id": str(pending.get("chain_head_leg_id", f"{tid}:open") or f"{tid}:open"), + "chain_prev_leg_id": str(pending.get("chain_prev_leg_id", "") or ""), + "chain_seq": int(pending.get("chain_seq", pending.get("retraction_legs", 0)) or 0), + "chain_token": str(pending.get("chain_token", "") or ""), + "chain_mode": str(pending.get("chain_mode", "LIVE") or "LIVE"), + }, + event_type="CLOSE", + event_id=f"{tid}:close", + exit_reason=str(x.get('reason', 'UNKNOWN')), + exit_price=exit_price, + quantity=pending['quantity'], + pnl=realized_pnl, + pnl_pct=float(x.get('pnl_pct', 0) or 0), + capital_before=capital_before, + capital_after=capital_after, + bars_held=max(0, int(x.get('bars_held', 0) or 0)), + scan_uuid=str(pending.get("scan_uuid", "") or ""), + pnl_leg=realized_pnl, + pnl_realized_total=float(pending.get("realized_pnl_legs_total", realized_pnl) or realized_pnl), + ) + te["execution_quality_json"] = json.dumps(execution_quality, default=str) + ch_put("trade_events", te) ch_put("trade_reconstruction", { "ts": _ch_ts_us(), "trade_id": str(tid or ""), @@ -4890,35 +5078,33 @@ class DolphinLiveTrader: ) self._persist_trade_execution_quality(execution_quality) pending.update(self._tp_curve_context(notional=float(pending.get("notional", 0) or 0))) - ch_put("trade_events", { - "ts": _ch_ts_us(), - "date": self.current_day or '', - "strategy": "blue", - "trade_id": tid, - "asset": pending.get('asset', subday_exit.get('asset', '')), - "side": pending.get('side', 'SHORT'), - "entry_price": pending.get('entry_price', 0), - "exit_price": float(subday_exit.get('exit_price', 0) or 0), - "quantity": round(float(pending.get('notional', 0) or 0) / max(float(pending.get('entry_price', 1) or 1), 1e-12), 6), - "capital_before": capital_before, - "capital_after": capital_after, - "pnl": realized_pnl, - "pnl_pct": float(subday_exit.get('pnl_pct', 0) or 0), - "exit_reason": str(subday_exit.get('reason', 'SUBDAY_ACB_NORMALIZATION')), - "vel_div_entry": float(pending.get('vel_div_entry', 0) or 0), - "boost_at_entry": float(pending.get('boost_at_entry', 0) or 0), - "beta_at_entry": float(pending.get('beta_at_entry', 0) or 0), - "posture": pending.get('posture', ''), - "leverage": float(pending.get('leverage', 0) or 0), - # CH column is UInt16 — negative poisons the spool - "bars_held": max(0, int(subday_exit.get('bars_held', 0) or 0)), - "regime_signal": 0, - "execution_quality_json": json.dumps(execution_quality, default=str), - "market_state_bundle_json": str(pending.get("market_state_bundle_json", "") or ""), - "tp_base_pct": float(pending.get("tp_base_pct", 0.0) or 0.0), - "tp_effective_pct": float(pending.get("tp_effective_pct", 0.0) or 0.0), - "our_leverage": float(pending.get("our_leverage", 0.0) or 0.0), - }) + te = self._trade_event_payload( + pending={ + **pending, + "trade_id": tid, + "chain_root_trade_id": str(pending.get("chain_root_trade_id", tid) or tid), + "chain_head_leg_id": str(pending.get("chain_head_leg_id", f"{tid}:open") or f"{tid}:open"), + "chain_prev_leg_id": str(pending.get("chain_prev_leg_id", "") or ""), + "chain_seq": int(pending.get("chain_seq", pending.get("retraction_legs", 0)) or 0), + "chain_token": str(pending.get("chain_token", "") or ""), + "chain_mode": str(pending.get("chain_mode", "LIVE") or "LIVE"), + }, + event_type="SUBDAY_EXIT", + event_id=f"{tid}:{str(subday_exit.get('reason', 'SUBDAY_ACB_NORMALIZATION')).lower()}", + exit_reason=str(subday_exit.get('reason', 'SUBDAY_ACB_NORMALIZATION')), + exit_price=float(subday_exit.get('exit_price', 0) or 0), + quantity=round(float(pending.get('notional', 0) or 0) / max(float(pending.get('entry_price', 1) or 1), 1e-12), 6), + pnl=realized_pnl, + pnl_pct=float(subday_exit.get('pnl_pct', 0) or 0), + capital_before=capital_before, + capital_after=capital_after, + bars_held=max(0, int(subday_exit.get('bars_held', 0) or 0)), + scan_uuid=str(pending.get("scan_uuid", "") or ""), + pnl_leg=realized_pnl, + pnl_realized_total=float(pending.get("realized_pnl_legs_total", realized_pnl) or realized_pnl), + ) + te["execution_quality_json"] = json.dumps(execution_quality, default=str) + ch_put("trade_events", te) self._announce_position_event( kind="trade_exit", severity="info" if float(subday_exit.get("pnl_pct", 0) or 0) >= 0 else "warning", diff --git a/prod/tests/test_operational_watchdog.py b/prod/tests/test_operational_watchdog.py new file mode 100644 index 0000000..f9a051e --- /dev/null +++ b/prod/tests/test_operational_watchdog.py @@ -0,0 +1,435 @@ +""" +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() diff --git a/prod/watchdog_decision.py b/prod/watchdog_decision.py new file mode 100644 index 0000000..c34f8c6 --- /dev/null +++ b/prod/watchdog_decision.py @@ -0,0 +1,141 @@ +""" +Lightweight, dependency-free decision seam for the scan-flow watchdog. + +Why a separate, dependency-free module: + ``nautilus_event_trader.py`` drags in the engine / Hazelcast / CH-writer + stack at *import* time (the module-level + ``from nautilus_dolphin.nautilus.proxy_boost_engine import create_d_liq_engine`` + and companion imports connect to infra that only exists under supervisord; + outside that environment the import blocks). The watchdog's restart + decision therefore can't be unit-tested in isolation there. This module + holds the **only** new logic for the 2026-09-16 04:10:40 ghost-subscription + wedge recovery — so it has zero dependencies and imports instantly. + +``nautilus_event_trader.py`` imports ``UPSTREAM_DARK_RESTART_S`` and +``upstream_dark_restart`` from here and consults the predicate in the previously +log-only ``"NO SCANS ... UNMANAGED"`` branch, promoting a long-frozen HZ +``latest_eigen_scan`` key (ghost-subscription after a WS reconnect) from a +reminder print to a self-restart via the existing ``_watchdog_restart`` -> +``os._exit(WATCHDOG_EXIT_CODE=86)`` -> supervisord respawn path. + +Cadence invariants (enforced by ``prod/tests/test_watchdog_decision.py``): + + SCAN_STALL_S (120s) < UPSTREAM_DARK_LOG_EVERY_S (300s) + < UPSTREAM_DARK_RESTART_S (900s) [this module] + <= warm-up-gated (uptime_ok checked before the + predicate in _scan_watchdog_loop) + +See prod/docs/SYSTEM_BIBLE_v7.md §38.10 and the r27 py-spy report (pid 3506857, +2026-09-16 04:10:40). +""" +from __future__ import annotations + +import math + +# --------------------------------------------------------------------------- +# (b) 2026-09-16 ghost-subscription self-heal threshold. +# +# The scan watchdog only *logs* "upstream dark / UNMANAGED" while the HZ +# latest_eigen_scan key is frozen (probe == last_probe_num, i.e. NOT None — a +# None probe is the "HZ client dead" case handled by the 3x-failure restart +# path in _scan_watchdog_loop). A frozen-but-not-None key is the hallmark of a +# *ghost subscription*: the reconnect re-subscribed + ACKed but the server never +# resumed the event stream, so the WS reader blocks in poll()/recv() (no +# liveness watchdog) and the key never advances. After this much +# accepted-scan staleness we stop nagging and self-restart, because the engine +# is starved and -- per the r27 soak -- the venue is FLAT (zero fills, capital +# intact), so a restart is free of position side-effects. +# +# Tunable: 900s (15 min) >> SCAN_STALL_S (120) and UPSTREAM_DARK_LOG_EVERY_S +# (300) so restarts only fire on a *confirmed* long dark window, never on a +# quiet market or a warm-up probe miss. Raise for calmer pairs/markets; lower +# only behind the r27 HL-testnet WS stability fix. +UPSTREAM_DARK_RESTART_S = 900.0 + +# Sentinel "no probe read" returned by _probe_latest_scan_number when the HZ +# key is missing/empty/corrupt. Kept here (vs the loop) so the seam is the +# single authority for the (b) restart condition. +_PROBE_MISSING = object() + + +def upstream_dark_restart( + acc_age_s: float, + uptime_ok: bool, + scan_number_probe: object, +) -> bool: + """(b) Ghost-subscription recovery decision (pure — no I/O, no self state). + + Returns True iff the watchdog should self-restart for the + "upstream dark (HZ key frozen)" case: + + * ``scan_number_probe`` is a real number (the HZ ``latest_eigen_scan`` + probe SUCCEEDED and returned an int). This is the *frozen-key* + ghost-subscription case (probe == last_probe_num). A *None / falsy* + probe means the HZ client itself is dead/unreachable; that is owned by + the separate 3x-failure restart path in ``_scan_watchdog_loop``, so + this predicate MUST return False for it (avoids a double-restart / + racing two restart paths). + * ``uptime_ok`` -- warm-up window elapsed; never self-restart during the + first ``WATCHDOG_RESTART_MIN_UPTIME_S`` to dodge boot-strap flakes. + * ``acc_age_s >= UPSTREAM_DARK_RESTART_S`` -- no scan ACCEPTED for at + least the dark-restart threshold (the engine has been starved long + enough that "it will come back" is an assumption, not evidence). + + Poison inputs are handled defensively (never raises): + * NaN acc_age -> False (corrupt clock; `nan >= x` is False) + * negative acc_age -> False (clock skew backward) + * +inf acc_age -> True (definitely dead) + * non-numeric probe (str/dict/etc.) -> treated as None (safe: the loop + still owns a None probe; (b) stays off) + """ + # (1) Probe must be a real scan number (frozen-key case). None / falsy / + # non-numeric probe is owned by the probe-None-3x path -> do NOT fire. + if scan_number_probe is None or scan_number_probe is _PROBE_MISSING: + return False + if not isinstance(scan_number_probe, (int, float)): + return False + if math.isnan(scan_number_probe) or math.isinf(scan_number_probe): + # A scan number that is NaN/inf is a corrupt probe, not a frozen key; + # let the probe-None-3x path handle it. (b) stays off. + return False + + # (2) Warm-up: never self-restart during boot. + if not uptime_ok: + return False + + # (3) Accepted-scan staleness past the ghost-subscription threshold. + # NaN acc_age -> `nan >= x` is False -> no spurious restart on a + # corrupt acc_age clock. -inf -> False. +inf -> True (== definitely + # dead). negative -> False (clock skew). + try: + return acc_age_s >= UPSTREAM_DARK_RESTART_S + except TypeError: + # Non-numeric acc_age (str/dict) -> don't crash the watchdog; treat as + # "not stale enough" and keep logging dark instead. + return False + +def scan_watchdog_dark_restart( + acc_age_s: float, + uptime_ok: bool, + scan_number_probe: object, + ev_age_s: float = 0.0, +) -> str | None: + """(b) Seam the live ``_scan_watchdog_loop`` calls for the ghost-subscription + restart decision. + + Pure (no I/O / no kernel state). Returns the canonical restart-reason + string iff :func:`upstream_dark_restart` says restart, else ``None``. + Centralising the reason text here (instead of building it inline in the + heavy ``nautilus_event_trader`` module) keeps the entire (b) branch + contract -- predicate + reason wording -- unit-testable without importing + ``nautilus_event_trader`` (whose module-level engine/HZ import blocks + outside the live supervisord environment). + """ + if not upstream_dark_restart(acc_age_s, uptime_ok, scan_number_probe): + return None + return ( + f"upstream dark: HZ latest_eigen_scan frozen at {scan_number_probe} " + f"for {acc_age_s:.0f}s (>= {UPSTREAM_DARK_RESTART_S}s) -- " + "ghost-subscription after WS reconnect (no reader liveness " + f"watchdog); acc_age={acc_age_s:.0f}s ev_age={float(ev_age_s):.0f}s" + )