Files
sentiment-engine/prod/clean_arch/violet/uv/uv_anomaly_sentinel.py
Codex b8cfd5df67 uv(sentinel): live-schema fix rounds 1+2 (PASS8 fb08a9f+361e6de, QA-verified 5/5 live)
Real colnames (timestamp/u_prefix_client_id, anomaly_events.ts), query/VST
errors FAIL never PASS (no green-by-error), openOrders wrapper unwrap,
check(d) CH_ALLOWLIST_DBS attribution (BLUE OBF writer no longer flagged).
Maiden live runs by Fable caught all 6; PASS8 fixed same-night.
2026-07-12 03:28:18 +02:00

546 lines
19 KiB
Python

"""
UV Anomaly Sentinel — read-only auditor for the VIOLET UV namespace.
Implements 5 checks from arming-checklist §4:
(a) Venue order without matching BRIDGE exec_journal row
(b) BRIDGE exec_journal row without venue order
(c) Venue order without 'u-' prefix in clientOrderId
(d) CH write outside dolphin_uv.* from runner window
(e) tripwire_ok=false occurrences
Usage:
uv_anomaly_sentinel.py # run all checks
uv_anomaly_sentinel.py --check a c # run subset
uv_anomaly_sentinel.py --verbose # show evidence rows
Exit code 0 = all PASS; 1 = any FAIL.
"""
from __future__ import annotations
import hashlib
import hmac
import json
import logging
import os
import sys
import time
import urllib.error
import urllib.parse
import urllib.request
from collections.abc import Mapping
from dataclasses import dataclass, field
from typing import Any
log = logging.getLogger(__name__)
# ─── Environment ──────────────────────────────────────────────────────────
CH_URL = os.environ.get("CH_URL", "http://localhost:8123")
CH_USER = os.environ.get("CH_USER", "dolphin")
CH_PASS = os.environ.get("CH_PASS", "dolphin_ch_2026")
CH_DB_UV = os.environ.get("CH_DB_UV", "dolphin_uv")
CH_DB_VIOLET = os.environ.get("CH_DB_VIOLET", "dolphin_violet")
VST_BASE_URL = os.environ.get(
"VST_BASE_URL", "https://open-api-vst.bingx.com"
)
VST_API_KEY = os.environ.get("BINGX_API_KEY", "")
VST_SECRET_KEY = os.environ.get("BINGX_SECRET_KEY", "")
LOOKBACK_HOURS = int(os.environ.get("SENTINEL_LOOKBACK_HOURS", "24"))
EVIDENCE_LIMIT = int(os.environ.get("SENTINEL_EVIDENCE_LIMIT", "10"))
QUERY_TIMEOUT_S = int(os.environ.get("SENTINEL_QUERY_TIMEOUT_S", "15"))
# Known non-UV databases that the sentinel must NOT flag (BLUE OBF, etc.)
CH_ALLOWLIST_DBS = os.environ.get(
"CH_ALLOWLIST_DBS", "dolphin,dolphin_malkhut,default"
).split(",")
# ─── Data types ────────────────────────────────────────────────────────────
@dataclass
class CheckResult:
check_id: str
description: str
passed: bool = True
evidence: list[str] = field(default_factory=list)
error: str = "" # IF non-empty, the check failed with an error (not a detection)
def add(self, line: str) -> None:
if len(self.evidence) < EVIDENCE_LIMIT:
self.evidence.append(line)
def summary(self) -> list[str]:
lines = [f" {'PASS' if self.passed else 'FAIL'} [{self.check_id}] {self.description}"]
if self.error:
lines.append(f" Error: {self.error}")
if not self.passed and self.evidence:
lines.append(" Evidence:")
for e in self.evidence:
lines.append(f" - {e}")
return lines
# ─── ClickHouse querier (read-only HTTP, error-tracked) ────────────────────
@dataclass
class CHQueryResult:
"""Result of a CH query — rows + error status."""
rows: list[dict[str, Any]]
error: str # empty = success; non-empty = error message
class CHQuerier:
"""Read-only ClickHouse HTTP client (JSONEachRow).
Every query method returns (rows, error): error is '' on success.
A caller that gets error != '' MUST report FAIL, never silent PASS.
"""
def __init__(
self,
url: str = CH_URL,
user: str = CH_USER,
password: str = CH_PASS,
) -> None:
self._url = url.rstrip("/")
self._user = user
self._password = password
def query(self, sql: str, **kw: Any) -> CHQueryResult:
timeout = kw.pop("timeout_s", QUERY_TIMEOUT_S)
endpoint = f"{self._url}/?default_format=JSONEachRow"
data = sql.encode("utf-8")
req = urllib.request.Request(endpoint, data=data, method="POST")
b64 = __import__("base64")
auth_raw = f"{self._user}:{self._password}"
encoded = b64.b64encode(auth_raw.encode()).decode()
req.add_header("Authorization", f"Basic {encoded}")
req.add_header("Content-Type", "text/plain; charset=utf-8")
try:
resp = urllib.request.urlopen(req, timeout=timeout)
body = resp.read().decode("utf-8")
except urllib.error.HTTPError as exc:
detail = exc.read().decode()[:300]
return CHQueryResult([], f"CH HTTP {exc.code}: {detail}")
except OSError as exc:
return CHQueryResult([], f"CH connection error: {exc}")
if not body.strip():
return CHQueryResult([], "")
rows: list[dict[str, Any]] = []
for line in body.strip().split("\n"):
line = line.strip()
if not line:
continue
try:
rows.append(json.loads(line))
except json.JSONDecodeError:
return CHQueryResult([], f"CH non-JSON response line: {line[:100]}")
return CHQueryResult(rows, "")
def describe_table(self, table: str, db: str = CH_DB_UV) -> CHQueryResult:
return self.query(f"DESCRIBE TABLE {db}.{table}")
def table_exists(self, table: str, db: str = CH_DB_UV) -> CHQueryResult:
r = self.query(
f"SELECT 1 FROM system.tables "
f"WHERE database = '{db}' AND name = '{table}'"
)
if r.error:
return r
return CHQueryResult(r.rows, "")
# ─── VST querier (read-only: open orders, positions) ───────────────────────
def _canonical_query(params: Mapping[str, object]) -> str:
filtered = {k: v for k, v in params.items() if v is not None and v != ""}
ordered = sorted(filtered.items(), key=lambda x: x[0])
return urllib.parse.urlencode(ordered, doseq=True)
def _sign_query(secret: str, query: str) -> str:
return hmac.new(
secret.encode("utf-8"), query.encode("utf-8"), hashlib.sha256
).hexdigest()
def _build_signed_params(
params: dict[str, object], secret: str
) -> dict[str, object]:
signed = dict(params)
signed["timestamp"] = int(time.time() * 1000)
signed["recvWindow"] = 5000
q = _canonical_query(signed)
signed["signature"] = _sign_query(secret, q)
return signed
@dataclass
class VSTResult:
"""Result of a VST API call — items + error status."""
items: list[dict[str, Any]]
error: str # empty = success; non-empty = error message
class VSTQuerier:
"""Read-only VST venue client — open orders, positions, balance."""
def __init__(
self,
base_url: str = VST_BASE_URL,
api_key: str = VST_API_KEY,
secret_key: str = VST_SECRET_KEY,
) -> None:
self._base = base_url.rstrip("/")
self._api_key = api_key
self._secret_key = secret_key
def _signed_get(self, path: str, params: dict[str, object] | None = None) -> VSTResult:
params = _build_signed_params(params or {}, self._secret_key)
qs = _canonical_query(params)
url = f"{self._base}{path}?{qs}"
req = urllib.request.Request(url, method="GET")
req.add_header("X-BX-APIKEY", self._api_key)
try:
resp = urllib.request.urlopen(req, timeout=QUERY_TIMEOUT_S)
body = resp.read().decode("utf-8")
except urllib.error.HTTPError as exc:
detail = exc.read().decode()[:300]
return VSTResult([], f"VST HTTP {exc.code}: {detail}")
except OSError as exc:
return VSTResult([], f"VST connection error: {exc}")
try:
parsed = json.loads(body)
except json.JSONDecodeError:
return VSTResult([], f"VST non-JSON response: {body[:200]}")
# BingX shape: {"code":0,"msg":"ok","data":{"orders":[...]}}
# openOrders → data.orders, positions → data.positions
code = parsed.get("code", -1) if isinstance(parsed, dict) else -1
if code != 0 and code != -1:
msg = parsed.get("msg", "unknown") if isinstance(parsed, dict) else "?"
return VSTResult([], f"VST API error code={code} msg={msg}")
data = parsed.get("data") if isinstance(parsed, dict) else None
if isinstance(data, dict):
items = data.get("orders") or data.get("positions") or []
elif isinstance(data, list):
items = data
else:
items = []
# Filter out non-dict items (empty wrappers, etc.)
clean = [r for r in items if isinstance(r, dict)]
return VSTResult(clean, "")
def open_orders(self) -> VSTResult:
return self._signed_get("/openApi/swap/v2/trade/openOrders")
def open_positions(self) -> VSTResult:
return self._signed_get("/openApi/swap/v2/user/positions")
def account_balance(self) -> VSTResult:
r = self._signed_get("/openApi/swap/v2/user/balance")
return r
# ─── Sentinel ──────────────────────────────────────────────────────────────
class UvAnomalySentinel:
"""Runs the 5 UV anomaly checks against CH + VST."""
def __init__(
self,
ch: CHQuerier | None = None,
vst: VSTQuerier | None = None,
enabled: set[str] | None = None,
lookback_hours: int = LOOKBACK_HOURS,
verbose: bool = False,
allowlist_dbs: list[str] | None = None,
) -> None:
self.ch = ch or CHQuerier()
self.vst = vst or VSTQuerier()
self.enabled = enabled or {"a", "b", "c", "d", "e"}
self.lookback_hours = lookback_hours
self.verbose = verbose
self.allowlist_dbs = allowlist_dbs or CH_ALLOWLIST_DBS
self.results: list[CheckResult] = []
def run(self) -> int:
checks: list[tuple[str, Any]] = [
("a", self._check_venue_without_journal),
("b", self._check_journal_without_venue),
("c", self._check_u_prefix),
("d", self._check_ch_write_outside_uv),
("e", self._check_tripwire_ok),
]
for cid, method in checks:
if cid not in self.enabled:
continue
self.results.append(method())
return self._report()
# ── Check implementations ───────────────────────────────────────────
def _check_venue_without_journal(self) -> CheckResult:
"""(a) Venue order without matching BRIDGE exec_journal row."""
r = CheckResult("a", "Venue order ↔ BRIDGE exec_journal match")
# 1. Fetch venue open orders
vst_r = self.vst.open_orders()
if vst_r.error:
r.passed = False
r.error = vst_r.error
return r
orders = vst_r.items
if not orders:
return r
# 2. Check exec_journal exists
exists_r = self.ch.table_exists("exec_journal", CH_DB_UV)
if exists_r.error:
r.passed = False
r.error = exists_r.error
return r
if not exists_r.rows:
r.passed = False
r.add("exec_journal table not found in dolphin_uv — no journal to match against")
return r
# 3. Collect venue clientOrderIds
venue_ids: set[str] = set()
for o in orders:
cid = str(
o.get("clientOrderId")
or o.get("clientOrderID")
or o.get("c")
or ""
)
if cid:
venue_ids.add(cid)
if not venue_ids:
return r
# 4. Query exec_journal for matching u_prefix_client_id
ids_escaped = "','".join(venue_ids)
sql = (
f"SELECT u_prefix_client_id, trade_id "
f"FROM {CH_DB_UV}.exec_journal "
f"WHERE u_prefix_client_id IN ('{ids_escaped}')"
)
jr = self.ch.query(sql)
if jr.error:
r.passed = False
r.error = jr.error
return r
journal_ids: set[str] = set()
for row in jr.rows:
cid = str(row.get("u_prefix_client_id") or "")
if cid:
journal_ids.add(cid)
orphan = venue_ids - journal_ids
if orphan:
r.passed = False
for cid in sorted(orphan):
r.add(f"clientOrderId={cid} has no exec_journal row")
return r
def _check_journal_without_venue(self) -> CheckResult:
"""(b) BRIDGE exec_journal row without venue order."""
r = CheckResult("b", "BRIDGE exec_journal row ↔ venue order")
# 1. Check exec_journal exists
exists_r = self.ch.table_exists("exec_journal", CH_DB_UV)
if exists_r.error:
r.passed = False
r.error = exists_r.error
return r
if not exists_r.rows:
return r
# 2. Get recent journal rows (not suppressed = still live)
sql = (
f"SELECT u_prefix_client_id, trade_id, suppressed "
f"FROM {CH_DB_UV}.exec_journal "
f"WHERE timestamp >= now() - INTERVAL {self.lookback_hours} HOUR "
f"AND suppressed = 0 "
f"ORDER BY timestamp DESC"
)
jr = self.ch.query(sql)
if jr.error:
r.passed = False
r.error = jr.error
return r
if not jr.rows:
return r
journal_ids: set[str] = set()
for row in jr.rows:
cid = str(row.get("u_prefix_client_id") or "")
if cid:
journal_ids.add(cid)
# 3. Fetch venue open orders
vst_r = self.vst.open_orders()
if vst_r.error:
r.passed = False
r.error = vst_r.error
return r
venue_ids: set[str] = set()
for o in vst_r.items:
cid = str(
o.get("clientOrderId")
or o.get("clientOrderID")
or o.get("c")
or ""
)
if cid:
venue_ids.add(cid)
missing = journal_ids - venue_ids
if missing:
r.passed = False
for cid in sorted(missing):
r.add(f"exec_journal u_prefix_client_id={cid} not on venue")
return r
def _check_u_prefix(self) -> CheckResult:
"""(c) Venue order without 'u-' prefix in clientOrderId."""
r = CheckResult("c", "Venue order u- prefix on clientOrderId")
vst_r = self.vst.open_orders()
if vst_r.error:
r.passed = False
r.error = vst_r.error
return r
bad: list[str] = []
for o in vst_r.items:
cid = str(
o.get("clientOrderId")
or o.get("clientOrderID")
or o.get("c")
or ""
)
if cid and not cid.startswith("u-"):
tid = o.get("tradeId") or o.get("orderId") or "(no tradeId)"
bad.append(f"clientOrderId={cid} tradeId/orderId={tid}")
if bad:
r.passed = False
for b in bad:
r.add(b)
return r
def _check_ch_write_outside_uv(self) -> CheckResult:
"""(d) CH write outside dolphin_uv.* from runner window."""
r = CheckResult("d", "CH writes only within dolphin_uv.*")
col = "current_database"
# Exclude known non-UV databases (BLUE OBF writer, system, etc.)
allowed = [CH_DB_UV] + self.allowlist_dbs
db_filter = " AND ".join(f"{col} != '{db}'" for db in allowed)
sql = (
f"SELECT query, {col} AS db, event_time "
f"FROM system.query_log "
f"WHERE type = 'QueryFinish' "
f"AND query LIKE 'INSERT%' "
f"AND {db_filter} "
f"AND event_time >= now() - INTERVAL {self.lookback_hours} HOUR "
f"ORDER BY event_time DESC "
f"LIMIT {EVIDENCE_LIMIT}"
)
qr = self.ch.query(sql)
if qr.error:
r.passed = False
r.error = qr.error
return r
if qr.rows:
r.passed = False
for row in qr.rows:
db = row.get("db") or row.get("current_database") or "?"
q = (row.get("query") or "")[:120]
r.add(f"INSERT into {db}: {q}")
return r
def _check_tripwire_ok(self) -> CheckResult:
"""(e) tripwire_ok=false occurrences in anomaly_events.
Real schema: dolphin_violet.anomaly_events uses 'ts' (DateTime64(6)),
not 'timestamp'. Verified by DESCRIBE TABLE.
"""
r = CheckResult("e", "tripwire_ok=false occurrences")
sql = (
f"SELECT ts, decision_id, trade_id, symbol, sensor, detail "
f"FROM {CH_DB_VIOLET}.anomaly_events "
f"WHERE anomaly = 'tripwire' "
f"AND ts >= now() - INTERVAL {self.lookback_hours} HOUR "
f"ORDER BY ts DESC "
f"LIMIT {EVIDENCE_LIMIT}"
)
qr = self.ch.query(sql)
if qr.error:
r.passed = False
r.error = qr.error
return r
if qr.rows:
r.passed = False
for row in qr.rows:
ts = row.get("ts", "?")
tid = row.get("trade_id", "?")
det = (row.get("detail") or "")[:80]
r.add(f"ts={ts} trade_id={tid} detail={det}")
return r
# ── Reporting ───────────────────────────────────────────────────────
def _report(self) -> int:
passed = 0
total = len(self.results)
for r in self.results:
if r.passed:
passed += 1
for line in r.summary():
print(line)
ok = passed == total
print(f"\nResult: {'PASS' if ok else 'FAIL'} ({passed}/{total} checks passed)")
return 0 if ok else 1
# ─── CLI entrypoint ─────────────────────────────────────────────────────────
def main(argv: list[str] | None = None) -> int:
import argparse
parser = argparse.ArgumentParser(
description="UV Anomaly Sentinel — read-only auditor (§4 checks)"
)
parser.add_argument(
"--check",
nargs="+",
choices=["a", "b", "c", "d", "e"],
default=["a", "b", "c", "d", "e"],
help="Which checks to run",
)
parser.add_argument(
"--verbose", "-v", action="store_true", help="Show evidence rows"
)
parser.add_argument(
"--lookback",
type=int,
default=LOOKBACK_HOURS,
help=f"Lookback hours (default {LOOKBACK_HOURS})",
)
parsed = parser.parse_args(argv)
sentinel = UvAnomalySentinel(
enabled=set(parsed.check),
lookback_hours=parsed.lookback,
verbose=parsed.verbose,
)
return sentinel.run()
if __name__ == "__main__":
sys.exit(main())