Files
sentiment-engine/prod/clean_arch/violet/uv/uv_anomaly_sentinel.py
Codex 257c48b127 uv(anomaly-sentinel): land UV ANOMALY SENTINEL — read-only §4 auditor (item-2 instrument)
Read-only auditor for arming-checklist §4 anomalies, the instrument for the
72h-clean finish-line clock. 5 checks against CH + VST (all read-only):
  (a) venue order without matching BRIDGE exec_journal row
  (b) exec_journal row without venue order
  (c) venue order without 'u-' clientOrderId prefix
  (d) CH write outside dolphin_uv.* in runner window
  (e) tripwire_ok=false occurrences

Authored by cmd-PASS1.2 (PASS8), provenance /mnt/vp-PASS8 bf55777
(agent/oa-violetPASS8). Landed by content (cross-clone fetch hung); files
SHA256-identical to source. Fable verification: 20/20 pass; mutation litmus
RED (inverting the check-(c) prefix guard fails 6 tests incl. PASS8's own
test_c_mutation_* guards — the doctrine's litmus baked into the suite);
read-only confirmed (no writes on any path).

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-11 19:20:22 +02:00

479 lines
17 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"))
# ─── Data types ────────────────────────────────────────────────────────────
@dataclass
class CheckResult:
check_id: str
description: str
passed: bool = True
evidence: list[str] = field(default_factory=list)
@property
def status(self) -> str:
return "PASS" if self.passed else "FAIL"
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 not self.passed and self.evidence:
lines.append(" Evidence:")
for e in self.evidence:
lines.append(f" - {e}")
hidden = max(0, 0) # we cap at EVIDENCE_LIMIT already
# Actually count surplus that would have been added
return lines
# ─── ClickHouse querier (read-only HTTP) ───────────────────────────────────
class CHQuerier:
"""Read-only ClickHouse HTTP client (JSONEachRow)."""
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) -> list[dict[str, Any]]:
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:
log.error("CH HTTP %s: %s", exc.code, exc.read().decode()[:200])
return []
except OSError as exc:
log.error("CH connection error: %s", exc)
return []
if not body.strip():
return []
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:
log.warning("CH non-JSON line: %.100s", line)
return rows
def describe_table(self, table: str, db: str = CH_DB_UV) -> list[dict[str, str]]:
return self.query(f"DESCRIBE TABLE {db}.{table}")
def table_exists(self, table: str, db: str = CH_DB_UV) -> bool:
rows = self.query(
f"SELECT 1 FROM system.tables "
f"WHERE database = '{db}' AND name = '{table}'"
)
return len(rows) > 0
# ─── 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
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) -> list[dict[str, Any]]:
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:
log.error("VST HTTP %s: %s", exc.code, exc.read().decode()[:200])
return []
except OSError as exc:
log.error("VST connection error: %s", exc)
return []
try:
parsed = json.loads(body)
except json.JSONDecodeError:
log.error("VST non-JSON response: %.200s", body)
return []
# BingX wraps data. openOrders uses "orders", positions uses "positions".
# Fallback chain: orders > data > positions > bare list.
if isinstance(parsed, dict):
data = (
parsed.get("orders")
or parsed.get("data")
or parsed.get("positions")
or []
)
else:
data = parsed
if isinstance(data, list):
return data
return [data] if data else []
def open_orders(self) -> list[dict[str, Any]]:
"""Return all open orders (across all symbols)."""
# VST uses same endpoint as LIVE for swap
rows = self._signed_get("/openApi/swap/v2/trade/openOrders")
orders: list[dict[str, Any]] = []
for r in rows:
if not isinstance(r, dict):
continue
orders.append(r)
return orders
def open_positions(self) -> list[dict[str, Any]]:
return self._signed_get("/openApi/swap/v2/user/positions")
def account_balance(self) -> dict[str, Any]:
rows = self._signed_get("/openApi/swap/v2/user/balance")
if rows:
return rows[0]
return {}
# ─── 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,
) -> 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.results: list[CheckResult] = []
# ── Public API ──────────────────────────────────────────────────────
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")
orders = self.vst.open_orders()
if not orders:
r.passed = True
return r
if not self.ch.table_exists("exec_journal", CH_DB_UV):
r.passed = False
r.add("exec_journal table does not exist in dolphin_uv")
return r
# 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)
else:
tid = o.get("tradeId") or o.get("orderId") or ""
r.add(f"venue order without clientOrderId: orderId={tid}")
if not venue_ids:
r.passed = len(r.evidence) == 0
return r
# Query exec_journal for any matching client_order_id
ids_escaped = "','".join(venue_ids)
sql = (
f"SELECT client_order_id, trade_id FROM {CH_DB_UV}.exec_journal "
f"WHERE client_order_id IN ('{ids_escaped}')"
)
journal_rows = self.ch.query(sql)
journal_ids: set[str] = set()
for jr in journal_rows:
cid = str(jr.get("client_order_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")
if not self.ch.table_exists("exec_journal", CH_DB_UV):
r.passed = True
return r
# Get recent journal rows with open state
sql = (
f"SELECT client_order_id, trade_id, status "
f"FROM {CH_DB_UV}.exec_journal "
f"WHERE ts >= now() - INTERVAL {self.lookback_hours} HOUR "
f"ORDER BY ts DESC"
)
journal_rows = self.ch.query(sql)
if not journal_rows:
return r
journal_ids: set[str] = set()
for jr in journal_rows:
cid = str(jr.get("client_order_id") or "")
if cid:
journal_ids.add(cid)
orders = self.vst.open_orders()
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)
missing = journal_ids - venue_ids
if missing:
r.passed = False
for cid in sorted(missing):
r.add(f"exec_journal client_order_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")
orders = self.vst.open_orders()
bad: list[str] = []
for o in orders:
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.*")
# Query system.query_log for INSERTs not targeting dolphin_uv
sql = (
f"SELECT query, database, table, event_time, http_user "
f"FROM system.query_log "
f"WHERE type = 'QueryFinish' "
f"AND query LIKE 'INSERT%' "
f"AND database != '{CH_DB_UV}' "
f"AND event_time >= now() - INTERVAL {self.lookback_hours} HOUR "
f"ORDER BY event_time DESC "
f"LIMIT {EVIDENCE_LIMIT}"
)
rows = self.ch.query(sql)
if rows:
r.passed = False
for row in rows:
db = row.get("database", "?")
tbl = row.get("table", "?")
q = (row.get("query") or "")[:120]
r.add(f"INSERT into {db}.{tbl}: {q}")
return r
def _check_tripwire_ok(self) -> CheckResult:
"""(e) tripwire_ok=false occurrences in anomaly_events."""
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}"
)
rows = self.ch.query(sql)
if rows:
r.passed = False
for row in 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())