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

This reverts commit 91ea1725a8.
This commit is contained in:
Codex
2026-09-16 20:49:52 +02:00
parent 7fc006329a
commit f144c3c324
3 changed files with 76 additions and 838 deletions

View File

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