uv(flight-recorder): stage-by-stage order tracer + dirty-25 root-cause maiden run
Read-only harness correlating scan->pick->size->bridge->drain->kernel->venue-> fill->flight->closure->capital->friction per trade_id. 20 tests, mutation-RED verified (kernel-verdict flip). Maiden sweep classifies all 25 bridged-no-venue rows: 15 PHANTOM_KERNEL_REFUSED (6 SLOT_BUSY enters, 9 NO_OPEN_POSITION exits), 5 SUBMITTED_VENUE_DEAD (BAND/CELR offline on VST), 4 LOG_GAP (pre-restart), 1 SUBMITTED_NO_VENUE_EVIDENCE (THETA phantom exit). 3 refused-exit orphans confirmed open at venue (THETA 3298.1, LTC 11.09, BNB 0.861).
This commit is contained in:
245
prod/clean_arch/violet/uv/test_uv_flight_recorder.py
Normal file
245
prod/clean_arch/violet/uv/test_uv_flight_recorder.py
Normal file
@@ -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)
|
||||
434
prod/clean_arch/violet/uv/uv_flight_recorder.py
Normal file
434
prod/clean_arch/violet/uv/uv_flight_recorder.py
Normal file
@@ -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())
|
||||
Reference in New Issue
Block a user