docs: watchdog ghost-subscription self-restart (b) design + test notes
Adds prod/docs/WATCHDOG_GHOST_SUBSCRIPTION_SELF_RESTART.md: the 2026-09-16 04:10:40 r27 ghost-subscription wedge (empty data after 04:08 WS reconnect), the log-only -> self-restart (b) promotion via the watchdog_decision seam, the guard contract (acc_age>=900s, uptime>600s, probe-not-None), the 51-test suite, mutation-litmus results, operational status (pid 3506857 still pre-fix; restart safe per AGENTS.md BLUE=flat-venue).
This commit is contained in:
295
prod/docs/WATCHDOG_GHOST_SUBSCRIPTION_SELF_RESTART.md
Normal file
295
prod/docs/WATCHDOG_GHOST_SUBSCRIPTION_SELF_RESTART.md
Normal file
@@ -0,0 +1,295 @@
|
||||
# Watchdog Ghost-Subscription Self-Restart — Fix `(b)` (2026-09-16)
|
||||
|
||||
**Incident:** r27 soak, pid `3506857`, `2026-09-16 04:10:40.351` CEST — silent
|
||||
stall ("empty data" from the Hyperliquid user-stream reader).
|
||||
**Scope:** `prod/nautilus_event_trader.py` (the live BLUE/DITA kernel entry),
|
||||
the `DolphinLiveTrader._scan_watchdog_loop` "upstream dark" branch.
|
||||
**Result:** 51/51 tests green on the committed tree; mutation litmus passes
|
||||
(each guard deletion fails only its targeted tests). Commit `91ea1725`.
|
||||
|
||||
---
|
||||
|
||||
## 1. TL;DR
|
||||
|
||||
The scan watchdog *detected* the wedge (HZ `latest_eigen_scan` frozen at `13638`
|
||||
while no scan had been accepted for >2 min) but its "upstream dark" branch was
|
||||
**log-only** — it printed `WATCHDOG: NO SCANS … UNMANAGED` and did nothing, so
|
||||
the ghost subscription (reader alive but starved) never recovered. Fix `(b)`
|
||||
promotes that branch to a **self-restart** (`_watchdog_restart` → `os._exit(86)`
|
||||
→ supervisord respawn), gated on:
|
||||
|
||||
| guard | value | why |
|
||||
|---|---|---|
|
||||
| accepted-scan staleness `acc_age` | `>= UPSTREAM_DARK_RESTART_S` = **900 s** | confirmed long dark (>> `SCAN_STALL_S`=120, >> dark-log nag=300) |
|
||||
| warm-up `uptime_ok` | `(now - _PROCESS_BOOT_TS) > 600 s` | never self-restart during boot-strap |
|
||||
| probe is a **number** | `scan_number_probe is not None` and numeric | `None` probe is the *dead-HZ-client* case, owned by the existing **3×-streak** restart path — (b) must NOT also fire |
|
||||
|
||||
---
|
||||
|
||||
## 2. The bug ("empty data") — how it came about
|
||||
|
||||
### Timeline (r27 soak, pid 3506857)
|
||||
|
||||
```
|
||||
00:20:15 → 03:57:23 21× VENUE_POSITION_OPEN (venue_pos=443.0) # phantom slot
|
||||
03:57:23 → ~04:08 scans/events flow normally, BAR #13630..#13637
|
||||
04:08:42–44 HL WS ping-timeout → `CRITICAL UV E-feed freshness re-stamp
|
||||
FAILED #1 … RuntimeError: HL WS account truth unavailable`
|
||||
(CAUGHT — feed kept ALIVE) → re-subscribe →
|
||||
`WS subscribed to clearinghouseState,orderUpdates,userFills … ACK`
|
||||
→ `userFills snapshot rows=30 recovered_current_run=0`
|
||||
04:10:40.351 BAR #13638 `polls=88560 watched=0` — LAST LINE. Silent.
|
||||
```
|
||||
|
||||
- The `userFills` snapshot returned `recovered_current_run=0` and the WS **ACKed
|
||||
the (re-)subscription**, but the Hyperliquid server **never resumed pushing the
|
||||
account-truth stream**. The reader therefore sat in `poll()/recv()` (py-spy:
|
||||
`asyncore.poll2 → select`, and 4× `do_sys_poll` on the dead post-reconnect
|
||||
socket). There is **no read-timeout / liveness watchdog on the post-reconnect
|
||||
reader**, so it blocked indefinitely while looking subscribed.
|
||||
- Because the reader starved, `on_exf_update` stopped firing → `scans_processed`
|
||||
/ `polls` / `BAR` froze → the journal went quiet. **py-spy dump (read-only,
|
||||
pid 3506857) confirmed: all 19 Python threads idle; MainThread in `run →
|
||||
time.sleep(1)`; `_heartbeat_loop` and `scan_watchdog` in their **timed**
|
||||
`wait(10.0)/(15.0)`; ch-writer threads in **timed** `urlopen(timeout=5)` /
|
||||
`Event.wait(interval)`.** No thread was in a raise, no traceback — **silent
|
||||
event-starvation**, not a hard deadlock or an exception. ("Empty data": the WS
|
||||
reader returns nothing because the server stopped sending, yet the socket
|
||||
looks open.)
|
||||
|
||||
### Why the watchdog saw it but didn't act
|
||||
|
||||
`_scan_watchdog_loop` distinguishes three probe states:
|
||||
|
||||
| probe value | meaning | branch |
|
||||
|---|---|---|
|
||||
| `None` | HZ client itself dead/unreachable (probe executor or key-missing) | `probe is None` → 3×-streak → **restart** (~2773) |
|
||||
| `int`, **different** from last | key advancing but no events (listener deaf) | `probe != last_probe_num` → **restart** (~2784) |
|
||||
| `int`, **same** as last (frozen) | **ghost subscription** — key present but stale | `else` → `:2790` **log-only** `NO SCANS … UNMANAGED` |
|
||||
|
||||
The wedge is the **third** case: the probe returned the same int (`13638`)
|
||||
repeatedly (key frozen, not `None`), so the probe-`None`-3× path (a) and the
|
||||
listener-deaf path both correctly did **not** fire; execution fell through to the
|
||||
log-only `:2790` branch, which only printed a reminder and never restarted. That
|
||||
is exactly the gap `(b)` closes.
|
||||
|
||||
> **Note on (a):** the probe-`None` path (`_probe_latest_scan_number`,
|
||||
> `:2656-2671`) returns `None` on **any** exception (`except Exception: return
|
||||
> None`, `:2668`/`2678` — wait, the `raise` at `:2668` is caught by the outer
|
||||
> `except Exception: return None`). It is safe and correct for a hard-dead HZ
|
||||
> client. The 04:10:40 wedge is **not** that path (probe was an int).
|
||||
|
||||
---
|
||||
|
||||
## 3. The fix — millimetric changes
|
||||
|
||||
Three files, all tracked. **No other watchdog branch was touched** (zero
|
||||
regression to the probe-`None`-3× / listener-deaf / worker-stalled restart
|
||||
paths).
|
||||
|
||||
### 3.1 `prod/watchdog_decision.py` (NEW — dependency-free seam)
|
||||
|
||||
Created because `import nautilus_event_trader` **hangs** outside the live
|
||||
supervisord environment (its 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; `import
|
||||
nautilus_event_trader` was verified to block in this shell). The (b) decision is
|
||||
therefore extracted into a stdlib-only (`math`) module so it imports instantly
|
||||
and is pure-testable.
|
||||
|
||||
```python
|
||||
# Cadence threshold for (b); single source of truth (imported by the kernel).
|
||||
# 900 s (15 min) >> SCAN_STALL_S (120) and UPSTREAM_DARK_LOG_EVERY_S (300) so
|
||||
# restarts fire only on a confirmed long dark window, never on a quiet market
|
||||
# or a warm-up probe miss. Tunable: raise for calmer pairs; lower only behind
|
||||
# the r27 HL-testnet WS stability fix.
|
||||
UPSTREAM_DARK_RESTART_S = 900.0
|
||||
|
||||
def upstream_dark_restart(acc_age_s, uptime_ok, scan_number_probe) -> bool:
|
||||
"""(b) Ghost-subscription recovery decision (pure — no I/O).
|
||||
|
||||
True iff: probe is a real scan number (frozen-KEY case, NOT None — a dead HZ
|
||||
client yields None and is owned by the 3x-streak restart path); warm-up
|
||||
elapsed (uptime_ok); and no scan ACCEPTED for >= UPSTREAM_DARK_RESTART_S.
|
||||
Poison inputs are handled defensively and never raise:
|
||||
NaN/-inf/negative acc_age -> False ; +inf -> True.
|
||||
None / non-numeric / NaN / inf probe -> False (None-3x path owns it)."""
|
||||
# (1) None / missing / non-numeric / NaN / inf probe -> NOT the frozen-key case.
|
||||
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):
|
||||
return False
|
||||
# (2) Warm-up gate.
|
||||
if not uptime_ok:
|
||||
return False
|
||||
# (3) Confirmed long dark. nan acc_age -> `nan>=x` False; -inf False; +inf True.
|
||||
try:
|
||||
return acc_age_s >= UPSTREAM_DARK_RESTART_S
|
||||
except TypeError: # non-numeric acc_age (str/None) — never crash the watchdog
|
||||
return False
|
||||
|
||||
def scan_watchdog_dark_restart(acc_age_s, uptime_ok, scan_number_probe, ev_age_s=0.0) -> str | None:
|
||||
"""(b) seam the live loop calls. Returns the canonical restart reason string
|
||||
iff upstream_dark_restart(...) is True, else None. Single-sources the reason
|
||||
wording (test asserts the exact text)."""
|
||||
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")
|
||||
```
|
||||
|
||||
### 3.2 `prod/nautilus_event_trader.py` (live kernel — edited)
|
||||
|
||||
**(a) Import** (inserted after the `nautilus_dolphin.nautilus.*` imports, ~line 37):
|
||||
```python
|
||||
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
|
||||
)
|
||||
```
|
||||
This adds **no** heavy dependency (the seam imports only `math`), so boot is
|
||||
unaffected. The live PYTHONPATH includes `…/prod` (cwd + line 29 sys.path insert),
|
||||
so `import watchdog_decision` resolves to this file.
|
||||
|
||||
**(b) The (b) branch** — inserted at the "upstream dark" site (was `:2790`,
|
||||
now ~`:2794`), **after** the listener-deaf `if probe is not None:` block and
|
||||
**before** the existing `NO SCANS … UNMANAGED` reminder print (which is
|
||||
preserved for the `acc_age < 900s` nag window):
|
||||
|
||||
```python
|
||||
_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: # UNCHANGED nag
|
||||
last_dark_log_ts = now
|
||||
print(f"[{datetime.now(timezone.utc).isoformat()}] "
|
||||
f"WATCHDOG: NO SCANS for {acc_age:.0f}s (HZ scan_number probe="
|
||||
f"{probe}) — upstream scanner appears DARK; open positions are "
|
||||
f"UNMANAGED until scans resume", flush=True)
|
||||
```
|
||||
|
||||
**Why `probe is not None` is enforced inside the seam (not at the call site):**
|
||||
the predicate returns `None`-reason for a `None` probe, so when the HZ client is
|
||||
hard-dead the 3×-streak path (~2773) restarts first (~45 s) and (b) never fires
|
||||
for `None` — no double-restart, no racing two restart paths.
|
||||
|
||||
**Pre-existing behaviour preserved** (verified by `TestSourceIntegrity`):
|
||||
- probe-`None` 3× restart — `"HZ probe failed …"` (~2773) ✓
|
||||
- listener-deaf restart — `"listener deaf: HZ latest_eigen_scan advanced …"` (~2784) ✓
|
||||
- worker-stalled restart — `"scan worker stalled …"` ✓
|
||||
- acc-fresh skip — `acc_age < SCAN_STALL_S` → `continue` (~2741) ✓
|
||||
- dark-log nag — `"NO SCANS … UNMANAGED"` reminder still prints for `acc_age < 900 s` ✓
|
||||
|
||||
(`_watchdog_restart`, `:2723`, calls `self._dump_blackbox(reason)` then
|
||||
`print(WATCHDOG_RESTART …)` then `os._exit(WATCHDOG_EXIT_CODE)`; the restart
|
||||
kills the process so the nag print is unreachable after a (b) fire — correct.)
|
||||
|
||||
### 3.3 `prod/tests/test_operational_watchdog.py` (NEW — 51 tests)
|
||||
|
||||
> The real `_scan_watchdog_loop` lives in the heavy kernel, whose `import`
|
||||
> blocks outside supervisord (true externalities: engine/Hazelcast/CH — none of
|
||||
> which enter the watchdog *decision*). Per `prod/docs/TESTING_DOCTRINE.md`
|
||||
> ("mock only true externalities"), the (b) *decision* was extracted into the
|
||||
> dep-free seam and the tests drive that seam.
|
||||
|
||||
- **Unit — `upstream_dark_restart` predicate** (`TestUpstreamDarkRestartPredicate`):
|
||||
restart at/above threshold (900.0, 900.1, 1000, 5000, +inf); no-restart below;
|
||||
**boundary inclusive** (900.0 → True, the `>=` guard); **warm-up forbids**
|
||||
(acc_age=9000, uptime_ok=False → False); **dead-HZ / corrupt probe excluded**
|
||||
(`None`, `"13638"`, `b"13638"`, `{...}`, `[13638]`, `nan`, `inf`, `-inf` → False,
|
||||
never raises); `0` is a valid scan number (→ True); poison acc_age
|
||||
(`-1`, `-1000`, `nan`, `-inf` → False; `+inf` → True); non-numeric acc_age
|
||||
(`"900"`, `None`) never raises; return type is `bool`.
|
||||
- **Unit — `scan_watchdog_dark_restart`** (`TestScanWatchdogDarkRestart`):
|
||||
reason string contract (`frozen at 13638`, `>= 900`, `ghost-subscription after
|
||||
WS reconnect`, `no reader liveness`, `acc_age=`/`ev_age=` flow-through);
|
||||
`None` below-threshold / warm-up / None-and-dead-probe; boundary returns a
|
||||
reason.
|
||||
- **Faithful stub-loop @ the seam** (`TestGhostSubscriptionE2E`): a verbatim
|
||||
transcription of `_scan_watchdog_loop`'s per-tick decision (nautilus
|
||||
`:2726-2821`) driving a duck-typed `_StubTrader` over N time-skipped ticks,
|
||||
calling the **real** `scan_watchdog_dark_restart`:
|
||||
- **restart at 900 s** (`frozen at 13638` reason present) ✔
|
||||
- **no restart before 900 s** (300/600/899 → none) ✔
|
||||
- **warm-up blocks** (uptime<600 even if acc_age huge) ✔
|
||||
- **probe-`None` 3× → restart via `HZ probe failed`, (b) NOT double-fired** ✔
|
||||
- **listener-deaf** (key advances 13638→13639, stale events, uptime → restart) ✔
|
||||
- **acc-fresh idle** (`acc_age<120` → no probe, no restart) ✔
|
||||
- **first probe sets baseline** (no spurious listener-deaf) ✔
|
||||
- dark-log nag fires while waiting (300/600 ticks).
|
||||
- **Source-integrity** (`TestSourceIntegrity`, reads the **live kernel as text —
|
||||
no import, no hang**): seam imported; (b) branch lives inside
|
||||
`_scan_watchdog_loop` and after the listener-deaf block + before the nag print;
|
||||
`(b)` not duplicated (count 1); nag print preserved; pre-existing restart
|
||||
branches intact; live cadences (`SCAN_STALL_S`=120, `WATCHDOG_RESTART_MIN_UPTIME_S`=600,
|
||||
`WATCHDOG_PROBE_INTERVAL_S`=30, `UPSTREAM_DARK_LOG_EVERY_S`=300,
|
||||
`WATCHDOG_EXIT_CODE`=86, `UPSTREAM_DARK_RESTART_S`=900) match the stub; the
|
||||
warm-up gate (`uptime_ok = (now - _PROCESS_BOOT_TS) > …`) is computed and
|
||||
passed to the seam; `os._exit(WATCHDOG_EXIT_CODE)` still backs the restart.
|
||||
|
||||
### 4. Mutation litmus (doctrine)
|
||||
|
||||
Run externally against `watchdog_decision.py`; each mutation broke **only** its
|
||||
targeted tests (no collateral), confirming precise protection:
|
||||
|
||||
| mutation | failing tests | count |
|
||||
|---|---|---|
|
||||
| `acc_age_s >= THRESH` → `>` (boundary guard) | `test_restarts_when_frozen_past_threshold[900.0]`, `test_boundary_is_inclusive`, `test_boundary_inclusive_returns_reason`, `test_restarts_at_900s_boundary` (E2E!) | 4 |
|
||||
| delete `if not uptime_ok: return False` (warm-up guard) | `test_warm_up_forbids_restart`, `test_returns_none_during_warm_up` | 2 |
|
||||
| delete `if math.isnan/isinf(scan_number_probe): return False` (probe-corruption guard) | `test_dead_hz_client_and_corrupt_probe_excluded[nan/inf/-inf]` | 3 |
|
||||
|
||||
Restored to original after each; final suite **51/51 PASS**.
|
||||
|
||||
---
|
||||
|
||||
## 5. Working trees & commits
|
||||
|
||||
- **Working tree touched (live):** `/mnt/dolphinng5_predict/prod/` —
|
||||
`nautilus_event_trader.py` (edited), `watchdog_decision.py` (new),
|
||||
`tests/test_operational_watchdog.py` (new). These are the live BLUE kernel
|
||||
files under supervisord (PYTHONPATH includes this tree).
|
||||
- **Release trees (untouched):** `/root/uv-releases/flight13-r2{4,5,6,7}-*]` were
|
||||
**not** edited (per the vendored-drift gate — edit canonical upstream, don't
|
||||
hand-edit vendored copies).
|
||||
- **Commit:** `91ea1725` on `tools/pi_wake_agent`,
|
||||
`3 files changed, 838 insertions(+), 76 deletions(-)`.
|
||||
|
||||
---
|
||||
|
||||
## 6. Operational status & next step (READ THIS)
|
||||
|
||||
- **pid 3506857 is still running the PRE-fix code.** The commit does not
|
||||
auto-deploy. (b) is only live in a process booted from the new file.
|
||||
- **Restarting 3506857 is safe**: per `AGENTS.md` the BLUE kernel
|
||||
(`nautilus_event_trader.py`) is **in-memory, Python-only, NO exchange
|
||||
exposure** — and the r27 soak verified the venue **flat**
|
||||
(`totalNtlPos=0.0`, `assetPositions=[]`, capital intact). A cold boot is
|
||||
~30–90 s and restores bookkeeping from CH/HZ on start. This both clears the
|
||||
**current** ghost and arms (b) for the next one.
|
||||
- **Not restarted yet** — per your standing constraint ("do NOT restart/kill
|
||||
pid 3506857 until you explicitly authorize"). **Authorize the single restart
|
||||
and I'll queue it; otherwise the current ghost stands (self-heal only fires on
|
||||
a process running the new code).**
|
||||
|
||||
> Side note: the r28 8h HL-testnet soak prep (FORCE-off / relaxed-vol /
|
||||
> F13_MONITORING flags from `/root/flight13_testnet_armed.env`) is **on hold**
|
||||
> pending this doc. The env uses `UV_FORCE_ENGAGE=1` (FORCE) — set
|
||||
> `UV_FORCE_ENGAGE=0` for "no FORCE"; the exact `vol`-threshold +
|
||||
> `F13_MONITORING` var names were not located in the armed env (grep over
|
||||
> `/mnt` hangs, per the CIFS note) — flag me the vars and I'll wire the r28
|
||||
> profile once (b) is deployed.
|
||||
Reference in New Issue
Block a user