diff --git a/prod/nautilus_event_trader.py b/prod/nautilus_event_trader.py index 95d03d7..bb36844 100644 --- a/prod/nautilus_event_trader.py +++ b/prod/nautilus_event_trader.py @@ -34,10 +34,6 @@ 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, @@ -45,8 +41,6 @@ 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 @@ -74,10 +68,9 @@ try: except Exception: BounceAdvisor = None try: - from adaptive_exit.post_win_long_overlay import PostWinExecutionFSM, PostWinExecutionFSMConfig + from adaptive_exit.post_win_long_overlay import PostWinExecutionFSM except Exception: PostWinExecutionFSM = None - PostWinExecutionFSMConfig = None try: from nautilus_dolphin.nautilus.alpha_exit_v7_engine import AlphaExitEngineV7, TradeContextV7 except Exception: @@ -309,12 +302,6 @@ 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") @@ -413,7 +400,6 @@ 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". @@ -487,15 +473,7 @@ 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( - PostWinExecutionFSMConfig( - require_context_gate=True, - ) - ) - if PostWinExecutionFSM is not None and PostWinExecutionFSMConfig is not None - else None - ) + self._efsm = PostWinExecutionFSM() if PostWinExecutionFSM 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() @@ -536,51 +514,19 @@ class DolphinLiveTrader: raw = self.features_map.blocking().get("maras_latest") if not raw: return {} - return extract_maras_context(raw) + 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), + } 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) @@ -825,88 +771,13 @@ 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. @@ -2791,20 +2662,6 @@ 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()}] " @@ -3844,37 +3701,6 @@ 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, @@ -4086,22 +3912,12 @@ 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: - 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): + if scan_number < self.last_scan_number - SCAN_NUMBER_RESET_GAP: log(f"WARN scanner restart detected: scan_number {self.last_scan_number} → " - f"{scan_number} (gap={backwards_gap}, consec={self._consec_backward_drops}) " - f"— re-anchoring dedup ratchet") - self._consec_backward_drops = 0 + f"{scan_number} — re-anchoring dedup ratchet") 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 @@ -4240,14 +4056,10 @@ 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, - **efsm_context, - }, + metadata={"trade_id": tid}, ) overlay_flip = bool(efsm_decision and efsm_decision.action == "TAG" and efsm_decision.side == "LONG") self._pending_entries[tid] = { @@ -4673,7 +4485,6 @@ 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 ""), @@ -4683,10 +4494,7 @@ 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")), - **efsm_context, - }, + metadata={"exit_reason": str(x.get("reason", "UNKNOWN"))}, ) if _efsm_out.action in {"ARMED", "TAG", "RESET"}: log(f"EFSM { _efsm_out.action }: { _efsm_out.to_dict() }") @@ -4730,33 +4538,37 @@ class DolphinLiveTrader: ) self._persist_trade_execution_quality(execution_quality) pending.update(self._tp_curve_context(notional=float(pending.get("notional", 0) or 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_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), + }) ch_put("trade_reconstruction", { "ts": _ch_ts_us(), "trade_id": str(tid or ""), @@ -5078,33 +4890,35 @@ class DolphinLiveTrader: ) self._persist_trade_execution_quality(execution_quality) pending.update(self._tp_curve_context(notional=float(pending.get("notional", 0) or 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) + 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), + }) 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 deleted file mode 100644 index f9a051e..0000000 --- a/prod/tests/test_operational_watchdog.py +++ /dev/null @@ -1,435 +0,0 @@ -""" -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 deleted file mode 100644 index c34f8c6..0000000 --- a/prod/watchdog_decision.py +++ /dev/null @@ -1,141 +0,0 @@ -""" -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" - )