diff --git a/prod/clean_arch/violet/uv/test_uv_flight_recorder.py b/prod/clean_arch/violet/uv/test_uv_flight_recorder.py new file mode 100644 index 0000000..42279fb --- /dev/null +++ b/prod/clean_arch/violet/uv/test_uv_flight_recorder.py @@ -0,0 +1,245 @@ +""" +Tests for uv_flight_recorder — real classifier/parser kernel, fakes only at +the CH / venue / log-file seams (true externalities). + +Log-line fixtures are verbatim copies from /root/dolphin_logs/uv_prime_live.log +(2026-07-11), so the parser is tested against production reality, not invented +strings. Mutation-sensitivity: flipping the kernel-verdict comparison in +classify_leg, the CET offset sign in parse_log_ts, or the slip sign in +enter_slip_bps goes RED here. +""" + +from __future__ import annotations + +from datetime import datetime, timezone + + +def utc_ms(dt: datetime) -> int: + """Naive-UTC datetime -> venue epoch ms (avoids local-tz .timestamp() skew).""" + return int(dt.replace(tzinfo=timezone.utc).timestamp() * 1000) + +import pytest + +from uv_flight_recorder import ( + LIVE, + LOG_GAP, + PHANTOM_KERNEL_REFUSED, + SUBMITTED_NO_VENUE_EVIDENCE, + SUBMITTED_VENUE_DEAD, + VENUE_SKIPPED, + FLAG_CROSS_SLOT_CLOSE, + FLAG_PHANTOM_WATCH, + FlightRecorder, + LogIndex, + classify_leg, + dash_symbol, + enter_slip_bps, + match_venue_order, + parse_log_ts, +) + +# ── verbatim production log lines (CET) ──────────────────────────────────── +L_PROMO = ("2026-07-11 15:39:06,964 INFO uv.promotion PROMOTION: bridged " + "intent=uv-promo-a7d4ddf2410e asset=BANDUSDT side=SHORT size=2919.7711 lev=1 raw_conv=1.1043") +L_WATCH = ("2026-07-11 15:39:06,965 INFO uv.tpsl_ticker TPSL WATCH: BANDUSDT 86ada259 " + "qty=2919.771058 entry=0.1718 tp=0.3500% sl=1.2000%") +L_DRAIN = ("2026-07-11 15:39:06,965 INFO uv.seam_drain SEAM DRAIN: intent=uv-promo-a7d4ddf2410e " + "-> u- clientOrderId=u-86ada259__uv_promo_a7d4ddf asset=BANDUSDT side=SHORT qty=2919.7711 lev=1.00") +L_OFFLINE = ("2026-07-11 15:39:07,769 WARNING prod.clean_arch.adapters.bingx_direct leverage POST " + "failed symbol=BAND-USDT lev=1: BAND-USDT is offline currently,all validted symbols in " + "api:/openApi/swap/v2/quote/contracts, please verify it — cache NOT updated, will retry on next submit") +L_OK = "2026-07-11 15:39:08,384 INFO uv.seam_drain KERNEL SUBMIT: intent=uv-promo-a7d4ddf2410e -> OK" +L_BUSY = "2026-07-10 14:48:06,251 INFO uv.seam_drain KERNEL SUBMIT: intent=uv-promo-3daa020eb347 -> SLOT_BUSY" +L_NOPOS = "2026-07-10 15:00:53,298 INFO uv.seam_drain KERNEL SUBMIT: intent=uv-promo-d6f57a98ebe7 -> NO_OPEN_POSITION" + + +class FakeLogIndex(LogIndex): + def __init__(self, lines: list[str]) -> None: + super().__init__(path="/nonexistent") + self._lines = list(lines) + + +class FakeCH: + def __init__(self, journal=None, decision=None, scan=None): + self._j, self._d, self._s = journal or [], decision, scan + + def query(self, sql: str, **kw): + if "exec_journal" in sql and "DISTINCT" not in sql: + return self._j + if "prime_decisions" in sql: + return [self._d] if self._d else [] + if "prime_scans" in sql: + return [self._s] if self._s else [] + if "DISTINCT trade_id" in sql: + return [{"trade_id": r["trade_id"]} for r in self._j[:1]] + return [] + + +class FakeVenue: + def __init__(self, orders=None, income=None): + self._orders, self._income = orders or [], income or [] + self.posted = False # read-only guard: nothing here can flip this + + def all_orders(self, symbol, t0, t1): + return self._orders + + def income(self, symbol, t0, t1): + return self._income + + +# ── parse_log_ts: CET -> UTC (mutation: flip the offset sign -> RED) ────── +def test_parse_log_ts_cet_to_utc(): + ts = parse_log_ts(L_OK) + assert ts == datetime(2026, 7, 11, 13, 39, 8, 384000) # 15:39 CET = 13:39 UTC + + +def test_parse_log_ts_rejects_non_ts_line(): + assert parse_log_ts(" c /= stddev[None, :]") is None + + +# ── classify_leg: the verdict kernel ─────────────────────────────────────── +@pytest.mark.parametrize( + "log_present,kernel,offline,order,checked,expect", + [ + (False, None, False, None, True, LOG_GAP), + (True, "SLOT_BUSY", False, None, True, PHANTOM_KERNEL_REFUSED), + (True, "NO_OPEN_POSITION", False, None, True, PHANTOM_KERNEL_REFUSED), + (True, "OK", True, None, True, SUBMITTED_VENUE_DEAD), + (True, "OK", False, {"orderId": 1}, True, LIVE), + (True, "OK", False, None, True, SUBMITTED_NO_VENUE_EVIDENCE), + (True, "OK", False, None, False, VENUE_SKIPPED), + # venue order trumps a stale offline warn (order proves it flew) + (True, "OK", True, {"orderId": 1}, True, LIVE), + ], +) +def test_classify_leg(log_present, kernel, offline, order, checked, expect): + assert classify_leg(log_present, kernel, offline, order, checked) == expect + + +# ── slip sign convention (mutation: swap fill/ref -> RED) ────────────────── +def test_enter_slip_negative_when_filled_below_decision(): + assert enter_slip_bps(100.0, 99.76) == pytest.approx(-24.0, abs=0.01) + assert enter_slip_bps(0.0, 99.0) == 0.0 + + +def test_dash_symbol(): + assert dash_symbol("BANDUSDT") == "BAND-USDT" + assert dash_symbol("BTCUSDT") == "BTC-USDT" + + +# ── LogIndex against verbatim production lines ───────────────────────────── +def test_log_index_kernel_verdicts(): + idx = FakeLogIndex([L_PROMO, L_DRAIN, L_OK, L_BUSY, L_NOPOS]) + assert idx.kernel_verdict("uv-promo-a7d4ddf2410e") == "OK" + assert idx.kernel_verdict("uv-promo-3daa020eb347") == "SLOT_BUSY" + assert idx.kernel_verdict("uv-promo-d6f57a98ebe7") == "NO_OPEN_POSITION" + assert idx.kernel_verdict("uv-promo-deadbeef0000") is None + + +def test_log_index_offline_warn_near(): + idx = FakeLogIndex([L_OFFLINE]) + at = datetime(2026, 7, 11, 13, 39, 7) # UTC of the warn + assert idx.offline_warn_near("BAND-USDT", at) + assert not idx.offline_warn_near("THETA-USDT", at) + assert not idx.offline_warn_near("BAND-USDT", datetime(2026, 7, 11, 13, 50, 0)) + + +# ── venue matching ───────────────────────────────────────────────────────── +def test_match_venue_order_prefers_client_id_prefix(): + orders = [ + {"clientOrderId": "u-86ada259__uv_promo_a7d4ddf", "executedQty": "2919.77", "time": 0}, + {"clientOrderId": "u-otherid__x", "executedQty": "2919.77", "time": 0}, + ] + got = match_venue_order(orders, "86ada259", 2919.77, datetime(2026, 7, 11, 13, 39, 7)) + assert got is orders[0] + + +def test_match_venue_order_time_qty_fallback_and_reject(): + ts = datetime(2026, 7, 11, 13, 39, 7) + near = {"clientOrderId": "", "executedQty": "2919.77", + "time": utc_ms(ts) + 5000} + far = {"clientOrderId": "", "executedQty": "2919.77", + "time": utc_ms(ts) + 3_600_000} + assert match_venue_order([near], "86ada259", 2919.771, ts) is near + assert match_venue_order([far], "86ada259", 2919.771, ts) is None + wrong_qty = {"clientOrderId": "", "executedQty": "100.0", + "time": utc_ms(ts)} + assert match_venue_order([wrong_qty], "86ada259", 2919.771, ts) is None + + +# ── end-to-end card on the BAND phantom (venue-offline) ──────────────────── +def _band_journal(): + return [ + {"ts": "2026-07-11 13:39:06.961", "trade_id": "86ada259", "asset": "BANDUSDT", + "side": "SHORT", "action": "ENTER", "reference_price": 0.1718, + "target_size": 2919.7710577361354, "leverage": 1.0, + "intent_id": "uv-promo-a7d4ddf2410e", "suppressed": 0}, + {"ts": "2026-07-11 14:00:53.287", "trade_id": "86ada259", "asset": "BANDUSDT", + "side": "SHORT", "action": "EXIT", "reference_price": 0.1718, + "target_size": 2919.7710577361354, "leverage": 1.0, + "intent_id": "uv-promo-61f9e4394ab3", "suppressed": 0}, + ] + + +def test_card_band_phantom_venue_offline(): + fr = FlightRecorder( + ch=FakeCH(journal=_band_journal()), + log_index=FakeLogIndex([L_PROMO, L_WATCH, L_DRAIN, L_OFFLINE, L_OK]), + venue=FakeVenue(orders=[]), + ) + card = fr.trace("86ada259") + assert card.asset == "BANDUSDT" + enter, exit_ = card.legs + assert enter.verdict == SUBMITTED_VENUE_DEAD + # EXIT intent has no log lines in this fixture -> LOG_GAP, not LIVE + assert exit_.verdict == LOG_GAP + assert FLAG_CROSS_SLOT_CLOSE not in card.flags + + +def test_card_cross_slot_close_flagged(): + """ENTER refused SLOT_BUSY + EXIT kernel-OK with a venue fill = stolen close.""" + journal = [ + {"ts": "2026-07-11 13:36:40.814", "trade_id": "30a8763b", "asset": "THETAUSDT", + "side": "SHORT", "action": "ENTER", "reference_price": 0.1515, + "target_size": 3300.330679, "leverage": 1.0, + "intent_id": "uv-promo-a5b021245578", "suppressed": 0}, + {"ts": "2026-07-11 13:37:13.290", "trade_id": "30a8763b", "asset": "THETAUSDT", + "side": "SHORT", "action": "EXIT", "reference_price": 0.1509, + "target_size": 3300.330679, "leverage": 1.0, + "intent_id": "uv-promo-f02ff2fbe073", "suppressed": 0}, + ] + lines = [ + "2026-07-11 15:36:40,818 INFO uv.tpsl_ticker TPSL WATCH: THETAUSDT 30a8763b qty=3300.330679 entry=0.1515 tp=0.3500% sl=1.2000%", + "2026-07-11 15:36:40,818 INFO uv.seam_drain SEAM DRAIN: intent=uv-promo-a5b021245578 -> u- clientOrderId=u-30a8763b__uv_promo_a5b0212 asset=THETAUSDT side=SHORT qty=3300.3307 lev=1.00", + "2026-07-11 15:36:40,820 INFO uv.seam_drain KERNEL SUBMIT: intent=uv-promo-a5b021245578 -> SLOT_BUSY", + "2026-07-11 15:37:16,984 INFO uv.seam_drain SEAM DRAIN: intent=uv-promo-f02ff2fbe073 -> u- clientOrderId=u-30a8763b__uv_promo_f02ff2f asset=THETAUSDT side=SHORT qty=3300.3307 lev=1.00", + "2026-07-11 15:37:18,330 INFO uv.seam_drain KERNEL SUBMIT: intent=uv-promo-f02ff2fbe073 -> OK", + ] + exit_ts = datetime(2026, 7, 11, 13, 37, 18) + venue = FakeVenue(orders=[{ + "clientOrderId": "", "executedQty": "3300.44", # someone else's qty, within tol + "avgPrice": "0.1509", "time": utc_ms(exit_ts), + }]) + fr = FlightRecorder(ch=FakeCH(journal=journal), log_index=FakeLogIndex(lines), venue=venue) + card = fr.trace("30a8763b") + enter, exit_ = card.legs + assert enter.verdict == PHANTOM_KERNEL_REFUSED + assert exit_.verdict == LIVE + assert FLAG_CROSS_SLOT_CLOSE in card.flags + assert FLAG_PHANTOM_WATCH in card.flags + + +def test_recorder_is_read_only_against_venue(): + venue = FakeVenue() + fr = FlightRecorder(ch=FakeCH(journal=_band_journal()), + log_index=FakeLogIndex([]), venue=venue) + fr.trace("86ada259") + assert venue.posted is False + + +def test_no_journal_rows_yields_note_not_crash(): + fr = FlightRecorder(ch=FakeCH(journal=[]), log_index=FakeLogIndex([]), + venue=FakeVenue()) + card = fr.trace("deadbeef") + assert card.legs == [] + assert any("no exec_journal rows" in n for n in card.notes) diff --git a/prod/clean_arch/violet/uv/uv_flight_recorder.py b/prod/clean_arch/violet/uv/uv_flight_recorder.py new file mode 100644 index 0000000..e7d3f57 --- /dev/null +++ b/prod/clean_arch/violet/uv/uv_flight_recorder.py @@ -0,0 +1,434 @@ +""" +UV Flight Recorder — step-debugger over recorded evidence for one order's flight. + +Correlates, for a given trade_id (or a sweep window), every stage of the +UV-PRIME pipeline, each against its authoritative evidence source: + + 1 SCAN dvol / eigenscan tape dolphin_uv.prime_scans + 2 GATE vol_ok + posture dolphin_uv.prime_decisions + 3 PICK IRP asset pick prime_decisions.entry_json (trade_id link) + 4 SIZE sizing + dual-lev decomposition exec_journal + PROMOTION log (raw_conv) + 5 BRIDGE promotion intent dolphin_uv.exec_journal + 6 DRAIN u- clientOrderId mint runner log (uv.seam_drain) + 7 KERNEL doorman verdict runner log (KERNEL SUBMIT -> ...) + 8 VENUE exchange contact / warnings runner log (bingx_direct / bingx.http) + 9 FILL venue order + fill VST allOrders (signed GET, read-only) + 10 FLIGHT TP/SL/MAX_HOLD watch + exit runner log (uv.tpsl_ticker / uv.tpsl) + 11 CLOSE exit leg (same ladder 5-9) journal EXIT intent + 12 CAPITAL qty/price integrity + pair PnL venue fills vs journal intent + 13 FRICTION fees + funding VST income endpoint + +Read-only everywhere: no CH writes, no venue POSTs, no state mutation. + +Usage: + uv_flight_recorder.py --trade-id 86ada259 + uv_flight_recorder.py --sweep --since "2026-07-10 06:43:47" --until "2026-07-11 16:00:00" + uv_flight_recorder.py --trade-id 86ada259 --no-venue # offline (journal+log only) + +Timebase law: exec_journal ts = UTC; runner log = CET (UTC+LOG_UTC_OFFSET_H, default 2). +All correlation is done in UTC. +""" + +from __future__ import annotations + +import argparse +import json +import os +import re +import sys +from dataclasses import dataclass, field +from datetime import datetime, timedelta, timezone +from typing import Any + +from uv_anomaly_sentinel import CHQuerier, VSTQuerier, CH_DB_UV + +LOG_PATH = os.environ.get("UV_RUNNER_LOG", "/root/dolphin_logs/uv_prime_live.log") +LOG_UTC_OFFSET_H = int(os.environ.get("LOG_UTC_OFFSET_H", "2")) +VENUE_MATCH_TOL_S = int(os.environ.get("FLIGHT_VENUE_TOL_S", "90")) +QTY_MISMATCH_PCT = float(os.environ.get("FLIGHT_QTY_TOL_PCT", "0.5")) + +# ─── Verdicts (leg-level) ────────────────────────────────────────────────── +LIVE = "LIVE" # kernel OK + venue order found +PHANTOM_KERNEL_REFUSED = "PHANTOM_KERNEL_REFUSED" # doorman said no; never flew +SUBMITTED_VENUE_DEAD = "SUBMITTED_VENUE_DEAD" # kernel OK, venue rejected/offline +SUBMITTED_NO_VENUE_EVIDENCE = "SUBMITTED_NO_VENUE_EVIDENCE" # kernel OK, no fill found +LOG_GAP = "LOG_GAP" # intent absent from runner log window +VENUE_SKIPPED = "VENUE_SKIPPED" # --no-venue / no keys + +# Flags (cross-cutting pathologies) +FLAG_PHANTOM_WATCH = "PHANTOM_WATCH" # TPSL watched an entry the kernel refused +FLAG_CROSS_SLOT_CLOSE = "CROSS_SLOT_CLOSE" # EXIT kernel-OK while own ENTER was refused + + +_LOG_TS_RE = re.compile(r"^(\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}),(\d{3})") +_KERNEL_RE = re.compile(r"KERNEL SUBMIT: intent=(\S+) -> (\S+)") +_OFFLINE_RE = re.compile(r"symbol=(\S+).*offline currently") + + +def parse_log_ts(line: str) -> datetime | None: + """Runner-log timestamp (CET) -> UTC datetime.""" + m = _LOG_TS_RE.match(line) + if not m: + return None + local = datetime.strptime(f"{m.group(1)}.{m.group(2)}", "%Y-%m-%d %H:%M:%S.%f") + return local - timedelta(hours=LOG_UTC_OFFSET_H) + + +def dash_symbol(asset: str) -> str: + """Journal asset (BANDUSDT) -> venue symbol (BAND-USDT).""" + return f"{asset[:-4]}-{asset[-4:]}" if asset.endswith("USDT") else asset + + +def classify_leg( + log_present: bool, + kernel_verdict: str | None, + venue_offline: bool, + venue_order: dict[str, Any] | None, + venue_checked: bool, +) -> str: + """Pure classifier: one leg's evidence -> verdict. Order of tests matters.""" + if not log_present: + return LOG_GAP + if kernel_verdict is not None and kernel_verdict != "OK": + return PHANTOM_KERNEL_REFUSED + # kernel OK (or submit line missing but drain present — treat as OK path) + if venue_order is not None: + return LIVE + if venue_offline: + return SUBMITTED_VENUE_DEAD + if not venue_checked: + return VENUE_SKIPPED + return SUBMITTED_NO_VENUE_EVIDENCE + + +def enter_slip_bps(ref_px: float, fill_px: float) -> float: + """Signed slip for a SHORT enter: negative = filled below decision = adverse.""" + if ref_px <= 0: + return 0.0 + return (fill_px - ref_px) / ref_px * 1e4 + + +@dataclass +class Stage: + name: str + status: str # PRESENT / MISSING / WARN / SKIP / MISMATCH + detail: str = "" + + +@dataclass +class Leg: + action: str # ENTER / EXIT + intent_id: str + ts_utc: str + qty: float + ref_px: float + stages: list[Stage] = field(default_factory=list) + verdict: str = "" + venue_order: dict[str, Any] | None = None + + +@dataclass +class FlightCard: + trade_id: str + asset: str + side: str + legs: list[Leg] = field(default_factory=list) + flags: list[str] = field(default_factory=list) + notes: list[str] = field(default_factory=list) + + +class LogIndex: + """One streaming pass over the runner log; lines indexed by needle.""" + + def __init__(self, path: str = LOG_PATH) -> None: + self._path = path + self._lines: list[str] | None = None + + def _load(self) -> list[str]: + if self._lines is None: + try: + with open(self._path, encoding="utf-8", errors="replace") as fh: + self._lines = [ + ln.rstrip("\n") for ln in fh + if ("uv." in ln or "bingx" in ln) and _LOG_TS_RE.match(ln) + ] + except OSError: + self._lines = [] + return self._lines + + def lines_with(self, *needles: str) -> list[str]: + return [ln for ln in self._load() if any(n in ln for n in needles)] + + def kernel_verdict(self, intent_id: str) -> str | None: + for ln in self.lines_with(intent_id): + m = _KERNEL_RE.search(ln) + if m and m.group(1) == intent_id: + return m.group(2) + return None + + def offline_warn_near(self, symbol_dash: str, ts_utc: datetime, tol_s: int = 30) -> bool: + for ln in self.lines_with("offline currently"): + m = _OFFLINE_RE.search(ln) + ts = parse_log_ts(ln) + if m and ts and m.group(1) == symbol_dash and abs((ts - ts_utc).total_seconds()) <= tol_s: + return True + return False + + +class VenueHistory(VSTQuerier): + """Adds historical order lookup to the sentinel's read-only client.""" + + def all_orders(self, symbol_dash: str, start_utc: datetime, end_utc: datetime) -> list[dict[str, Any]]: + # NOTE: server-side startTime/endTime on VST allOrders silently returns + # an empty set (observed 2026-07-11) — fetch unfiltered, filter here. + rows = self._signed_get( + "/openApi/swap/v2/trade/allOrders", + {"symbol": symbol_dash, "limit": 500}, + ) + # allOrders nests as data.orders; sentinel fallback may hand back the dict + if rows and isinstance(rows, dict): # type: ignore[unreachable] + rows = rows.get("orders", []) + out = [] + for r in rows: + if isinstance(r, dict) and isinstance(r.get("orders"), list): + out.extend(x for x in r["orders"] if isinstance(x, dict)) + elif isinstance(r, dict): + out.append(r) + elif isinstance(r, list): + out.extend(x for x in r if isinstance(x, dict)) + lo = int(start_utc.replace(tzinfo=timezone.utc).timestamp() * 1000) + hi = int(end_utc.replace(tzinfo=timezone.utc).timestamp() * 1000) + return [o for o in out if lo <= int(o.get("time") or 0) <= hi] + + def income(self, symbol_dash: str, start_utc: datetime, end_utc: datetime) -> list[dict[str, Any]]: + return self._signed_get( + "/openApi/swap/v2/user/income", + { + "symbol": symbol_dash, + "startTime": int(start_utc.timestamp() * 1000), + "endTime": int(end_utc.timestamp() * 1000), + "limit": 1000, + }, + ) + + +def match_venue_order( + orders: list[dict[str, Any]], + trade_id: str, + qty: float, + ts_utc: datetime, + tol_s: int = VENUE_MATCH_TOL_S, +) -> dict[str, Any] | None: + """Prefer clientOrderId prefix match; fall back to time+qty tolerance.""" + prefix = f"u-{trade_id}" + for o in orders: + cid = str(o.get("clientOrderId") or o.get("clientOrderID") or "") + if cid.startswith(prefix): + return o + for o in orders: + try: + ots = datetime.fromtimestamp(int(o.get("time", 0)) / 1000, timezone.utc).replace(tzinfo=None) + oqty = float(o.get("executedQty") or o.get("origQty") or 0) + except (TypeError, ValueError): + continue + if abs((ots - ts_utc).total_seconds()) <= tol_s and qty > 0 and abs(oqty - qty) / qty <= QTY_MISMATCH_PCT / 100: + return o + return None + + +class FlightRecorder: + def __init__( + self, + ch: CHQuerier | None = None, + log_index: LogIndex | None = None, + venue: VenueHistory | None = None, + use_venue: bool = True, + ) -> None: + self.ch = ch or CHQuerier() + self.log = log_index or LogIndex() + self.venue = venue if venue is not None else (VenueHistory() if use_venue else None) + self.use_venue = use_venue and self.venue is not None + + # ── evidence pulls ──────────────────────────────────────────────── + def journal_rows(self, trade_id: str) -> list[dict[str, Any]]: + return self.ch.query( + f"SELECT toString(timestamp) AS ts, trade_id, asset, side, action, reference_price," + f" target_size, leverage, intent_id, suppressed" + f" FROM {CH_DB_UV}.exec_journal WHERE trade_id='{trade_id}'" + f" AND kind='BRIDGE' ORDER BY timestamp FORMAT JSONEachRow" + ) + + def decision_row(self, trade_id: str) -> dict[str, Any] | None: + rows = self.ch.query( + f"SELECT toString(ts) AS ts, scan_number, vel_div, posture, vol_ok," + f" trade_direction, asset, side, leverage, entry_json" + f" FROM {CH_DB_UV}.prime_decisions WHERE has_entry=1" + f" AND entry_json LIKE '%{trade_id}%' LIMIT 1 FORMAT JSONEachRow" + ) + return rows[0] if rows else None + + def scan_row(self, scan_number: int) -> dict[str, Any] | None: + rows = self.ch.query( + f"SELECT toString(ts) AS ts, scan_number, vel_div, w50_velocity," + f" w750_velocity, instability_50 FROM {CH_DB_UV}.prime_scans" + f" WHERE scan_number={scan_number} LIMIT 1 FORMAT JSONEachRow" + ) + return rows[0] if rows else None + + # ── the card ────────────────────────────────────────────────────── + def trace(self, trade_id: str) -> FlightCard: + jrows = self.journal_rows(trade_id) + if not jrows: + return FlightCard(trade_id, "?", "?", notes=[f"no exec_journal rows for trade_id={trade_id}"]) + asset, side = jrows[0]["asset"], jrows[0]["side"] + card = FlightCard(trade_id, asset, side) + sym = dash_symbol(asset) + + # stages 1-3: scan / gate / pick + dec = self.decision_row(trade_id) + if dec: + card.notes.append( + f"decision scan={dec['scan_number']} vol_ok={dec['vol_ok']}" + f" posture={dec['posture']} vel_div={float(dec['vel_div']):.4f}" + ) + scan = self.scan_row(int(dec["scan_number"])) + if scan: + card.notes.append( + f"scan tape w50={float(scan['w50_velocity']):.4f}" + f" w750={float(scan['w750_velocity']):.4f}" + f" instab50={float(scan['instability_50']):.4f}" + ) + else: + card.notes.append("no prime_decisions row (pick evidence missing)") + + # venue order pool for the whole trade window + orders: list[dict[str, Any]] = [] + t0 = datetime.strptime(jrows[0]["ts"][:19], "%Y-%m-%d %H:%M:%S") - timedelta(minutes=2) + t1 = datetime.strptime(jrows[-1]["ts"][:19], "%Y-%m-%d %H:%M:%S") + timedelta(minutes=5) + if self.use_venue: + orders = self.venue.all_orders(sym, t0, t1) + + enter_refused = False + for jr in jrows: + leg = self._trace_leg(jr, sym, orders) + if leg.action == "ENTER" and leg.verdict == PHANTOM_KERNEL_REFUSED: + enter_refused = True + if self.log.lines_with(f"TPSL WATCH: {asset} {trade_id}"): + card.flags.append(FLAG_PHANTOM_WATCH) + if leg.action == "EXIT" and enter_refused and leg.verdict == LIVE: + card.flags.append(FLAG_CROSS_SLOT_CLOSE) + card.legs.append(leg) + + self._capital_and_friction(card, sym, t0, t1) + return card + + def _trace_leg(self, jr: dict[str, Any], sym: str, orders: list[dict[str, Any]]) -> Leg: + intent = jr["intent_id"] + ts_utc = datetime.strptime(jr["ts"][:23], "%Y-%m-%d %H:%M:%S.%f") + leg = Leg(jr["action"], intent, jr["ts"], float(jr["target_size"]), float(jr["reference_price"])) + leg.stages.append(Stage("BRIDGE", "PRESENT", + f"size={leg.qty:.4f} lev={float(jr['leverage']):.2f} suppressed={jr['suppressed']}")) + + ilines = self.log.lines_with(intent) + log_present = bool(ilines) + drain = next((ln for ln in ilines if "SEAM DRAIN" in ln), None) + leg.stages.append(Stage("DRAIN", "PRESENT" if drain else "MISSING", + drain.split("-> ")[-1][:60] if drain else "no seam_drain line")) + kv = self.log.kernel_verdict(intent) + leg.stages.append(Stage("KERNEL", "PRESENT" if kv else "MISSING", kv or "no KERNEL SUBMIT line")) + + offline = self.log.offline_warn_near(sym, ts_utc) + leg.stages.append(Stage("VENUE_CALL", "WARN" if offline else ("PRESENT" if kv == "OK" else "SKIP"), + f"{sym} offline at venue" if offline else "")) + + vo = match_venue_order(orders, jr["trade_id"], leg.qty, ts_utc) if orders else None + leg.venue_order = vo + if vo: + fill_px = float(vo.get("avgPrice") or 0) + slip = enter_slip_bps(leg.ref_px, fill_px) + leg.stages.append(Stage("FILL", "PRESENT", + f"qty={vo.get('executedQty')} avg={fill_px} slip={slip:+.1f}bps")) + else: + leg.stages.append(Stage("FILL", "MISSING" if self.use_venue else "SKIP", + "no venue order matched" if self.use_venue else "venue not queried")) + + leg.verdict = classify_leg(log_present, kv, offline, vo, self.use_venue) + return leg + + def _capital_and_friction(self, card: FlightCard, sym: str, t0: datetime, t1: datetime) -> None: + fills = [(lg, lg.venue_order) for lg in card.legs if lg.venue_order] + for lg, vo in fills: + vq = float(vo.get("executedQty") or 0) + if lg.qty > 0 and abs(vq - lg.qty) / lg.qty * 100 > QTY_MISMATCH_PCT: + card.notes.append( + f"CAPITAL MISMATCH {lg.action}: journal qty={lg.qty:.4f} venue qty={vq:.4f}") + if len(fills) == 2: + (e, evo), (x, xvo) = fills[0], fills[1] + ep, xp = float(evo.get("avgPrice") or 0), float(xvo.get("avgPrice") or 0) + if ep > 0 and card.side == "SHORT": + gross = (ep - xp) * float(evo.get("executedQty") or 0) + card.notes.append(f"PAIR GROSS (short, venue px) = {gross:+.2f} USDT") + if self.use_venue and fills: + inc = self.venue.income(sym, t0, t1) + fees = sum(float(r.get("income", 0)) for r in inc + if str(r.get("incomeType", "")).upper() in ("COMMISSION", "FUNDING_FEE")) + card.notes.append(f"FRICTION (commission+funding in window) = {fees:+.4f} USDT") + + def sweep(self, since: str, until: str) -> list[FlightCard]: + rows = self.ch.query( + f"SELECT DISTINCT trade_id FROM {CH_DB_UV}.exec_journal" + f" WHERE kind='BRIDGE' AND suppressed=0" + f" AND timestamp >= '{since}' AND timestamp <= '{until}'" + f" ORDER BY trade_id FORMAT JSONEachRow" + ) + return [self.trace(r["trade_id"]) for r in rows] + + +def render(card: FlightCard) -> str: + out = [f"TRADE {card.trade_id} {card.asset} {card.side}" + + (f" FLAGS: {','.join(card.flags)}" if card.flags else "")] + for n in card.notes: + out.append(f" · {n}") + for leg in card.legs: + out.append(f" [{leg.action:5}] {leg.ts_utc}Z intent={leg.intent_id} → {leg.verdict}") + for st in leg.stages: + out.append(f" {st.name:10} {st.status:8} {st.detail}") + return "\n".join(out) + + +def main(argv: list[str] | None = None) -> int: + p = argparse.ArgumentParser(description="UV flight recorder (read-only)") + p.add_argument("--trade-id") + p.add_argument("--sweep", action="store_true") + p.add_argument("--since", default="2026-07-10 06:43:47") + p.add_argument("--until", default="2100-01-01 00:00:00") + p.add_argument("--no-venue", action="store_true") + p.add_argument("--json", action="store_true") + a = p.parse_args(argv) + + fr = FlightRecorder(use_venue=not a.no_venue) + cards = [fr.trace(a.trade_id)] if a.trade_id else (fr.sweep(a.since, a.until) if a.sweep else []) + if not cards: + p.print_help() + return 2 + for c in cards: + if a.json: + print(json.dumps({ + "trade_id": c.trade_id, "asset": c.asset, "side": c.side, "flags": c.flags, + "notes": c.notes, + "legs": [{"action": l.action, "intent": l.intent_id, "ts": l.ts_utc, + "verdict": l.verdict, + "stages": [[s.name, s.status, s.detail] for s in l.stages]} + for l in c.legs]})) + else: + print(render(c)) + print() + if a.sweep: + from collections import Counter + verdicts = Counter(l.verdict for c in cards for l in c.legs) + print("SWEEP SUMMARY:", dict(verdicts)) + return 0 + + +if __name__ == "__main__": + sys.exit(main())