watchdog: ghost-subscription self-restart (b) + seam + tests

Promote the log-only 'upstream dark' branch in _scan_watchdog_loop (nautilus_event_trader.py) to _watchdog_restart when the HZ latest_eigen_scan key is frozen past UPSTREAM_DARK_RESTART_S=900s with uptime elapsed.

- prod/watchdog_decision.py (NEW, dep-free seam): UPSTREAM_DARK_RESTART_S + upstream_dark_restart(predicate) + scan_watchdog_dark_restart((b) branch seam). Testable without importing the heavy kernel (module-level engine/HZ import blocks outside supervisord).
- prod/nautilus_event_trader.py: import the seam; dark-log branch calls scan_watchdog_dark_restart(acc_age, uptime_ok, probe, ev_age) -> _watchdog_restart. Guarded acc_age>=900s, uptime>600s, probe NOT None (None owned by existing 3x-streak path). Pre-existing branches (probe-None-3x, listener-deaf, worker-stuck) and dark-log reminder print UNCHANGED.
- prod/tests/test_operational_watchdog.py (NEW, 51 tests): predicate unit (all branches/edges/poison/NaN/inf/warm-up), wrapper seam, faithful stub-loop E2E (frozen key + time-skipped ticks -> restart at 900s; not before; warm-up blocks; probe-None-3x not double-fired; listener-deaf; acc-fresh idle), source-integrity pin on live file. Mutation litmus: each guard deletion fails only its targeted tests (>= -> >: 4; uptime: 2; nan/inf probe: 3).
This commit is contained in:
codex
2026-09-16 19:16:26 +02:00
parent b76fbe5042
commit 91ea1725a8
3 changed files with 838 additions and 76 deletions

View File

@@ -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.adaptive_circuit_breaker import AdaptiveCircuitBreaker
from nautilus_dolphin.nautilus.ob_features import OBFeatureEngine from nautilus_dolphin.nautilus.ob_features import OBFeatureEngine
from nautilus_dolphin.nautilus.ob_provider import MockOBProvider 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 ( from nautilus_dolphin.nautilus.esof_size_gate import (
parse_esof_payload, esof_gate_from_payload, esof_score_from_payload, parse_esof_payload, esof_gate_from_payload, esof_score_from_payload,
esof_size_mult_from_score, ESOF_STALE_FALLBACK_MULT, ESOF_FRESHNESS_S, 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.adapters.eigen_scan_normalizer import normalize_ng7_scan
from prod.clean_arch.obf_tp_observation import inject_obf_midprice 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 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: try:
sys.path.insert(0, '/mnt/dolphinng5_predict/Observability') sys.path.insert(0, '/mnt/dolphinng5_predict/Observability')
from esof_advisor import compute_esof as _compute_esof_inline from esof_advisor import compute_esof as _compute_esof_inline
@@ -68,9 +74,10 @@ try:
except Exception: except Exception:
BounceAdvisor = None BounceAdvisor = None
try: try:
from adaptive_exit.post_win_long_overlay import PostWinExecutionFSM from adaptive_exit.post_win_long_overlay import PostWinExecutionFSM, PostWinExecutionFSMConfig
except Exception: except Exception:
PostWinExecutionFSM = None PostWinExecutionFSM = None
PostWinExecutionFSMConfig = None
try: try:
from nautilus_dolphin.nautilus.alpha_exit_v7_engine import AlphaExitEngineV7, TradeContextV7 from nautilus_dolphin.nautilus.alpha_exit_v7_engine import AlphaExitEngineV7, TradeContextV7
except Exception: except Exception:
@@ -302,6 +309,12 @@ WATCHDOG_EXIT_CODE = 86
# Scanner restarts reset scan_number to 0. A backwards jump larger than this # 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). # is a restart (accept + re-anchor ratchet), not a stale duplicate (drop).
SCAN_NUMBER_RESET_GAP = 1000 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]: def _trade_log_paths(ts_dt: datetime) -> tuple[str, str]:
log_date = ts_dt.strftime("%Y%m%d") 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._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._scan_executor = ThreadPoolExecutor(max_workers=1, thread_name_prefix="scan")
self.last_scan_number = -1 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; # Scan-flow watchdog state. Event ts proves the HZ listener is alive;
# accept ts proves the worker thread is draining; the dupe counter # accept ts proves the worker thread is draining; the dupe counter
# separates "worker stuck" from "upstream flooding duplicates". # separates "worker stuck" from "upstream flooding duplicates".
@@ -473,7 +487,15 @@ class DolphinLiveTrader:
self.trade_direction: int = _direction_from_env() self.trade_direction: int = _direction_from_env()
self.vol_p60_threshold: float = _vol_p60_threshold_from_env() self.vol_p60_threshold: float = _vol_p60_threshold_from_env()
self._runtime_direction: int = self.trade_direction 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._trade_announcement_center = None
self._processed_retract_commands: deque = deque(maxlen=5000) self._processed_retract_commands: deque = deque(maxlen=5000)
self._processed_retract_set: set[str] = set() self._processed_retract_set: set[str] = set()
@@ -514,19 +536,51 @@ class DolphinLiveTrader:
raw = self.features_map.blocking().get("maras_latest") raw = self.features_map.blocking().get("maras_latest")
if not raw: if not raw:
return {} return {}
payload = json.loads(raw) if isinstance(raw, str) else raw return extract_maras_context(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: except Exception:
return {} 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: def _resolve_runtime_direction(self) -> int:
"""Resolve active trade direction for the next eligible entry.""" """Resolve active trade direction for the next eligible entry."""
base = int(self.trade_direction) base = int(self.trade_direction)
@@ -771,13 +825,88 @@ class DolphinLiveTrader:
bundle = {} bundle = {}
if self._market_state_runtime is not None and getattr(self._market_state_runtime, "latest_bundle_dict", None): 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) bundle = dict(self._market_state_runtime.latest_bundle_dict)
maras_ctx = self._latest_maras_context()
bundle = merge_bundle_with_maras(bundle, maras_ctx)
return { return {
"tp_base_pct": float(self._tp_base_pct), "tp_base_pct": float(self._tp_base_pct),
"tp_effective_pct": float(tp_effective_pct), "tp_effective_pct": float(tp_effective_pct),
"our_leverage": float(our_leverage), "our_leverage": float(our_leverage),
"market_state_bundle_json": json.dumps(bundle, default=str, sort_keys=True) if bundle else "{}", "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: def _sync_tp_threshold(self) -> None:
"""Read live TP threshold from HZ control plane and propagate to engine. """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") f"{last_probe_num} → {probe} but no events for {ev_age:.0f}s")
last_probe_num = probe last_probe_num = probe
last_probe_ts = now 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: if now - last_dark_log_ts > UPSTREAM_DARK_LOG_EVERY_S:
last_dark_log_ts = now last_dark_log_ts = now
print(f"[{datetime.now(timezone.utc).isoformat()}] " 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), "pnl_realized_total": float(pending.get("realized_pnl_legs_total", 0.0) or 0.0),
"bars_held": bars_held, "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", { ch_put("trade_reconstruction", {
"ts": _ch_ts_us(), "ts": _ch_ts_us(),
"trade_id": tid, "trade_id": tid,
@@ -3912,12 +4086,22 @@ class DolphinLiveTrader:
# scan until manually restarted (near-miss on 2026-06-09/10). # scan until manually restarted (near-miss on 2026-06-09/10).
with self._dedup_lock: with self._dedup_lock:
if scan_number > 0 and scan_number <= self.last_scan_number: 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} → " 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: else:
self._dupe_drops_total += 1 self._dupe_drops_total += 1
return return
else:
self._consec_backward_drops = 0
self.last_scan_number = scan_number self.last_scan_number = scan_number
self._last_scan_accept_ts = time.time() self._last_scan_accept_ts = time.time()
self.scans_processed += 1 self.scans_processed += 1
@@ -4056,10 +4240,14 @@ class DolphinLiveTrader:
efsm_decision = None efsm_decision = None
overlay_flip = False overlay_flip = False
if self._efsm is not None and int(e.get('direction', -1)) == 1 and int(self.trade_direction) == -1: 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( efsm_decision = self._efsm.tag_next_entry(
asset=str(e.get('asset', '') or ''), asset=str(e.get('asset', '') or ''),
entry_ts=datetime.now(timezone.utc), 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") overlay_flip = bool(efsm_decision and efsm_decision.action == "TAG" and efsm_decision.side == "LONG")
self._pending_entries[tid] = { self._pending_entries[tid] = {
@@ -4485,6 +4673,7 @@ class DolphinLiveTrader:
log(f" MarketStateRuntime outcome update failed for {tid}: {e}") log(f" MarketStateRuntime outcome update failed for {tid}: {e}")
if self._efsm is not None: if self._efsm is not None:
try: try:
efsm_context = self._efsm_long_overlay_context(vel_div=vel_div, prices_dict=prices_dict)
_efsm_out = self._efsm.observe_closed_trade( _efsm_out = self._efsm.observe_closed_trade(
trade_id=str(tid or ""), trade_id=str(tid or ""),
asset=str(pending.get("asset", "") or ""), asset=str(pending.get("asset", "") or ""),
@@ -4494,7 +4683,10 @@ class DolphinLiveTrader:
leverage=float(pending.get("leverage", 0) or 0), leverage=float(pending.get("leverage", 0) or 0),
closed_ts=datetime.now(timezone.utc), closed_ts=datetime.now(timezone.utc),
was_overlay_flip=bool(pending.get("overlay_flip", False)), 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"}: if _efsm_out.action in {"ARMED", "TAG", "RESET"}:
log(f"EFSM { _efsm_out.action }: { _efsm_out.to_dict() }") log(f"EFSM { _efsm_out.action }: { _efsm_out.to_dict() }")
@@ -4538,37 +4730,33 @@ class DolphinLiveTrader:
) )
self._persist_trade_execution_quality(execution_quality) self._persist_trade_execution_quality(execution_quality)
pending.update(self._tp_curve_context(notional=float(pending.get("notional", 0) or 0))) pending.update(self._tp_curve_context(notional=float(pending.get("notional", 0) or 0)))
ch_put("trade_events", { te = self._trade_event_payload(
"ts": _ch_ts_us(), pending={
"date": pending['entry_date'], **pending,
"strategy": "blue",
"trade_id": tid, "trade_id": tid,
"asset": pending['asset'], "chain_root_trade_id": str(pending.get("chain_root_trade_id", tid) or tid),
"side": pending['side'], "chain_head_leg_id": str(pending.get("chain_head_leg_id", f"{tid}:open") or f"{tid}:open"),
"entry_price": pending['entry_price'], "chain_prev_leg_id": str(pending.get("chain_prev_leg_id", "") or ""),
"exit_price": exit_price, "chain_seq": int(pending.get("chain_seq", pending.get("retraction_legs", 0)) or 0),
"quantity": pending['quantity'], "chain_token": str(pending.get("chain_token", "") or ""),
"capital_before": capital_before, "chain_mode": str(pending.get("chain_mode", "LIVE") or "LIVE"),
"capital_after": capital_after, },
"pnl": realized_pnl, event_type="CLOSE",
"pnl_pct": float(x.get('pnl_pct', 0) or 0), event_id=f"{tid}:close",
"exit_reason": str(x.get('reason', 'UNKNOWN')), exit_reason=str(x.get('reason', 'UNKNOWN')),
"vel_div_entry": pending['vel_div_entry'], exit_price=exit_price,
"boost_at_entry": pending['boost_at_entry'], quantity=pending['quantity'],
"beta_at_entry": pending['beta_at_entry'], pnl=realized_pnl,
"posture": pending['posture'], pnl_pct=float(x.get('pnl_pct', 0) or 0),
"leverage": pending['leverage'], capital_before=capital_before,
# CH column is UInt16 — a negative value poisons the spool capital_after=capital_after,
# (head-of-line jam, incident 2026-06-12: bars_held=-106) bars_held=max(0, int(x.get('bars_held', 0) or 0)),
"bars_held": max(0, int(x.get('bars_held', 0) or 0)), scan_uuid=str(pending.get("scan_uuid", "") or ""),
"regime_signal": 0, pnl_leg=realized_pnl,
"tp_threshold": float(self.eng.exit_manager.fixed_tp_pct), pnl_realized_total=float(pending.get("realized_pnl_legs_total", realized_pnl) or realized_pnl),
"execution_quality_json": json.dumps(execution_quality, default=str), )
"market_state_bundle_json": str(pending.get("market_state_bundle_json", "") or ""), te["execution_quality_json"] = json.dumps(execution_quality, default=str)
"tp_base_pct": float(pending.get("tp_base_pct", 0.0) or 0.0), ch_put("trade_events", te)
"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", { ch_put("trade_reconstruction", {
"ts": _ch_ts_us(), "ts": _ch_ts_us(),
"trade_id": str(tid or ""), "trade_id": str(tid or ""),
@@ -4890,35 +5078,33 @@ class DolphinLiveTrader:
) )
self._persist_trade_execution_quality(execution_quality) self._persist_trade_execution_quality(execution_quality)
pending.update(self._tp_curve_context(notional=float(pending.get("notional", 0) or 0))) pending.update(self._tp_curve_context(notional=float(pending.get("notional", 0) or 0)))
ch_put("trade_events", { te = self._trade_event_payload(
"ts": _ch_ts_us(), pending={
"date": self.current_day or '', **pending,
"strategy": "blue",
"trade_id": tid, "trade_id": tid,
"asset": pending.get('asset', subday_exit.get('asset', '')), "chain_root_trade_id": str(pending.get("chain_root_trade_id", tid) or tid),
"side": pending.get('side', 'SHORT'), "chain_head_leg_id": str(pending.get("chain_head_leg_id", f"{tid}:open") or f"{tid}:open"),
"entry_price": pending.get('entry_price', 0), "chain_prev_leg_id": str(pending.get("chain_prev_leg_id", "") or ""),
"exit_price": float(subday_exit.get('exit_price', 0) or 0), "chain_seq": int(pending.get("chain_seq", pending.get("retraction_legs", 0)) or 0),
"quantity": round(float(pending.get('notional', 0) or 0) / max(float(pending.get('entry_price', 1) or 1), 1e-12), 6), "chain_token": str(pending.get("chain_token", "") or ""),
"capital_before": capital_before, "chain_mode": str(pending.get("chain_mode", "LIVE") or "LIVE"),
"capital_after": capital_after, },
"pnl": realized_pnl, event_type="SUBDAY_EXIT",
"pnl_pct": float(subday_exit.get('pnl_pct', 0) or 0), event_id=f"{tid}:{str(subday_exit.get('reason', 'SUBDAY_ACB_NORMALIZATION')).lower()}",
"exit_reason": str(subday_exit.get('reason', 'SUBDAY_ACB_NORMALIZATION')), exit_reason=str(subday_exit.get('reason', 'SUBDAY_ACB_NORMALIZATION')),
"vel_div_entry": float(pending.get('vel_div_entry', 0) or 0), exit_price=float(subday_exit.get('exit_price', 0) or 0),
"boost_at_entry": float(pending.get('boost_at_entry', 0) or 0), quantity=round(float(pending.get('notional', 0) or 0) / max(float(pending.get('entry_price', 1) or 1), 1e-12), 6),
"beta_at_entry": float(pending.get('beta_at_entry', 0) or 0), pnl=realized_pnl,
"posture": pending.get('posture', ''), pnl_pct=float(subday_exit.get('pnl_pct', 0) or 0),
"leverage": float(pending.get('leverage', 0) or 0), capital_before=capital_before,
# CH column is UInt16 — negative poisons the spool capital_after=capital_after,
"bars_held": max(0, int(subday_exit.get('bars_held', 0) or 0)), bars_held=max(0, int(subday_exit.get('bars_held', 0) or 0)),
"regime_signal": 0, scan_uuid=str(pending.get("scan_uuid", "") or ""),
"execution_quality_json": json.dumps(execution_quality, default=str), pnl_leg=realized_pnl,
"market_state_bundle_json": str(pending.get("market_state_bundle_json", "") or ""), pnl_realized_total=float(pending.get("realized_pnl_legs_total", realized_pnl) or realized_pnl),
"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), te["execution_quality_json"] = json.dumps(execution_quality, default=str)
"our_leverage": float(pending.get("our_leverage", 0.0) or 0.0), ch_put("trade_events", te)
})
self._announce_position_event( self._announce_position_event(
kind="trade_exit", kind="trade_exit",
severity="info" if float(subday_exit.get("pnl_pct", 0) or 0) >= 0 else "warning", severity="info" if float(subday_exit.get("pnl_pct", 0) or 0) >= 0 else "warning",

View File

@@ -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()

141
prod/watchdog_decision.py Normal file
View File

@@ -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"
)