From 257c48b12721d8453d0eebdb8bbe10c2dc94fc0f Mon Sep 17 00:00:00 2001 From: Codex Date: Sat, 11 Jul 2026 19:20:22 +0200 Subject: [PATCH] =?UTF-8?q?uv(anomaly-sentinel):=20land=20UV=20ANOMALY=20S?= =?UTF-8?q?ENTINEL=20=E2=80=94=20read-only=20=C2=A74=20auditor=20(item-2?= =?UTF-8?q?=20instrument)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- .../violet/uv/test_uv_anomaly_sentinel.py | 368 ++++++++++++++ .../violet/uv/uv_anomaly_sentinel.py | 478 ++++++++++++++++++ 2 files changed, 846 insertions(+) create mode 100644 prod/clean_arch/violet/uv/test_uv_anomaly_sentinel.py create mode 100644 prod/clean_arch/violet/uv/uv_anomaly_sentinel.py diff --git a/prod/clean_arch/violet/uv/test_uv_anomaly_sentinel.py b/prod/clean_arch/violet/uv/test_uv_anomaly_sentinel.py new file mode 100644 index 0000000..23b36ed --- /dev/null +++ b/prod/clean_arch/violet/uv/test_uv_anomaly_sentinel.py @@ -0,0 +1,368 @@ +""" +UV Anomaly Sentinel tests — synthetic fixtures per check + mutation. + +Each check has: + - happy_path: normal data → PASS + - mutation: inject failure condition → FAIL + - Mutation litmus: disable a check → test goes RED (proves check fires) + +Also: disabled-check no-op, exit-code nonzero on any FAIL. +""" + +from __future__ import annotations + +import json +import sys +from typing import Any + +import pytest + +# We import after path insertion below to pick up the sentinel module +sys.path.insert(0, "/mnt/vp-PASS8") + +from prod.clean_arch.violet.uv.uv_anomaly_sentinel import ( + CHQuerier, + CheckResult, + UvAnomalySentinel, + VSTQuerier, +) + + +# ═══════════════════════════════════════════════════════════════════════ +# Fake backends +# ═══════════════════════════════════════════════════════════════════════ + + +class FakeCHQuerier: + """Synthetic ClickHouse backend — returns pre-set responses by table.""" + + def __init__(self, **tables: Any) -> None: + # tables: keyword = tablename, value = list[dict] + self._tables = tables + self._describe: dict[str, list[dict[str, str]]] = {} + self._exists: set[str] = set() + self._query_log_rows: list[dict[str, Any]] = tables.get("query_log", []) + self._anomaly_rows: list[dict[str, Any]] = tables.get("anomaly_events", []) + self._journal_rows: list[dict[str, Any]] = tables.get("exec_journal", []) + self._describe_schemas: dict[str, list[dict[str, str]]] = {} + self._exists_tables: set[str] = set() + if "exec_journal" in tables: + self._exists_tables.add("exec_journal") + if "anomaly_events" in tables: + self._exists_tables.add("anomaly_events") + + def query(self, sql: str, **kw: Any) -> list[dict[str, Any]]: + sql_lower = sql.lower() + if "query_log" in sql_lower: + return self._query_log_rows + if "anomaly_events" in sql_lower: + return self._anomaly_rows + if "exec_journal" in sql_lower: + return self._journal_rows + if "select 1 from system.tables" in sql_lower: + # "SELECT 1 FROM system.tables WHERE database='...' AND name='...'" + for table in self._exists_tables: + if table in sql: + return [{"1": 1}] + return [] + return [] + + def describe_table(self, table: str, db: str = "") -> list[dict[str, str]]: + return self._describe_schemas.get(table, []) + + def table_exists(self, table: str, db: str = "") -> bool: + return table in self._exists_tables + + +class FakeVSTQuerier: + """Synthetic VST backend — returns pre-set open orders.""" + + def __init__(self, open_orders: list[dict[str, Any]] | None = None) -> None: + self._open_orders = open_orders or [] + + def open_orders(self) -> list[dict[str, Any]]: + return self._open_orders + + def open_positions(self) -> list[dict[str, Any]]: + return [] + + def account_balance(self) -> dict[str, Any]: + return {} + + +# ═══════════════════════════════════════════════════════════════════════ +# Shared fixture data +# ═══════════════════════════════════════════════════════════════════════ + +U_ORDER_1 = {"clientOrderId": "u-abc-001", "tradeId": "T-u-001"} +U_ORDER_2 = {"clientOrderId": "u-def-002", "tradeId": "T-u-002"} +JOURNAL_1 = {"client_order_id": "u-abc-001", "trade_id": "T-u-001", "status": "NEW"} +JOURNAL_2 = {"client_order_id": "u-def-002", "trade_id": "T-u-002", "status": "NEW"} +PREFIX_ORDER = {"clientOrderId": "p-e-abc123", "tradeId": "T-p-001"} +NO_PREFIX_ORDER = {"clientOrderId": "random123", "tradeId": "T-bad"} +ORPHAN_VENUE_ORDER = {"clientOrderId": "u-orphan-999", "tradeId": "T-orphan"} +ORPHAN_JOURNAL = {"client_order_id": "u-stale-888", "trade_id": "T-stale", "status": "NEW"} + + +# ═══════════════════════════════════════════════════════════════════════ +# A — Check (a): Venue order without matching journal row +# ═══════════════════════════════════════════════════════════════════════ + + +def test_a_happy_path(): + """All venue orders have matching exec_journal entries.""" + ch = FakeCHQuerier(exec_journal=[JOURNAL_1, JOURNAL_2]) + vst = FakeVSTQuerier(open_orders=[U_ORDER_1, U_ORDER_2]) + s = UvAnomalySentinel(ch=ch, vst=vst, enabled={"a"}) + s.run() + assert s.results[0].passed + + +def test_a_mutation(): + """Venue order without exec_journal → FAIL.""" + ch = FakeCHQuerier(exec_journal=[JOURNAL_1]) + vst = FakeVSTQuerier(open_orders=[U_ORDER_1, ORPHAN_VENUE_ORDER]) + s = UvAnomalySentinel(ch=ch, vst=vst, enabled={"a"}) + s.run() + assert not s.results[0].passed + assert any("u-orphan" in e for e in s.results[0].evidence) + + +def test_a_mutation_no_journal_table(): + """No exec_journal table → FAIL.""" + ch = FakeCHQuerier() + vst = FakeVSTQuerier(open_orders=[U_ORDER_1]) + s = UvAnomalySentinel(ch=ch, vst=vst, enabled={"a"}) + s.run() + assert not s.results[0].passed + assert any("exec_journal" in e for e in s.results[0].evidence) + + +# ═══════════════════════════════════════════════════════════════════════ +# B — Check (b): Journal row without venue order +# ═══════════════════════════════════════════════════════════════════════ + + +def test_b_happy_path(): + """All journal entries have matching venue orders.""" + ch = FakeCHQuerier(exec_journal=[JOURNAL_1, JOURNAL_2]) + vst = FakeVSTQuerier(open_orders=[U_ORDER_1, U_ORDER_2]) + s = UvAnomalySentinel(ch=ch, vst=vst, enabled={"b"}) + s.run() + assert s.results[0].passed + + +def test_b_mutation(): + """Journal entry without venue order → FAIL.""" + ch = FakeCHQuerier(exec_journal=[JOURNAL_1, ORPHAN_JOURNAL]) + vst = FakeVSTQuerier(open_orders=[U_ORDER_1]) + s = UvAnomalySentinel(ch=ch, vst=vst, enabled={"b"}) + s.run() + assert not s.results[0].passed + assert any("u-stale" in e for e in s.results[0].evidence) + + +def test_b_mutation_no_table(): + """No exec_journal table → PASS (nothing to check against).""" + ch = FakeCHQuerier() + vst = FakeVSTQuerier(open_orders=[U_ORDER_1]) + s = UvAnomalySentinel(ch=ch, vst=vst, enabled={"b"}) + s.run() + assert s.results[0].passed + + +# ═══════════════════════════════════════════════════════════════════════ +# C — Check (c): u- prefix on clientOrderId +# ═══════════════════════════════════════════════════════════════════════ + + +def test_c_happy_path(): + """All open orders use u- prefix.""" + vst = FakeVSTQuerier(open_orders=[U_ORDER_1, U_ORDER_2]) + s = UvAnomalySentinel(ch=FakeCHQuerier(), vst=vst, enabled={"c"}) + s.run() + assert s.results[0].passed + + +def test_c_mutation_p_prefix(): + """PINK prefix (p-) detected as non-UV → FAIL.""" + vst = FakeVSTQuerier(open_orders=[U_ORDER_1, PREFIX_ORDER]) + s = UvAnomalySentinel(ch=FakeCHQuerier(), vst=vst, enabled={"c"}) + s.run() + assert not s.results[0].passed + assert any("p-e-abc123" in e for e in s.results[0].evidence) + + +def test_c_mutation_no_prefix(): + """Order without u- prefix → FAIL.""" + vst = FakeVSTQuerier(open_orders=[U_ORDER_1, NO_PREFIX_ORDER]) + s = UvAnomalySentinel(ch=FakeCHQuerier(), vst=vst, enabled={"c"}) + s.run() + assert not s.results[0].passed + assert any("random123" in e for e in s.results[0].evidence) + + +def test_c_mutation_empty_client_order_id(): + """Order with empty clientOrderId → not flagged (no prefix to check).""" + empty = {"clientOrderId": "", "tradeId": "T-empty"} + vst = FakeVSTQuerier(open_orders=[U_ORDER_1, empty]) + s = UvAnomalySentinel(ch=FakeCHQuerier(), vst=vst, enabled={"c"}) + s.run() + assert s.results[0].passed # empty is not a violation + + +# ═══════════════════════════════════════════════════════════════════════ +# D — Check (d): CH writes outside dolphin_uv +# ═══════════════════════════════════════════════════════════════════════ + + +def test_d_happy_path(): + """No unexpected INSERTs in query_log.""" + ch = FakeCHQuerier(query_log=[]) + s = UvAnomalySentinel(ch=ch, vst=FakeVSTQuerier(), enabled={"d"}) + s.run() + assert s.results[0].passed + + +def test_d_mutation(): + """INSERT into non-uv database → FAIL.""" + ch = FakeCHQuerier( + query_log=[ + { + "query": "INSERT INTO dolphin_violet.anomaly_events VALUES ...", + "database": "dolphin_violet", + "table": "anomaly_events", + "event_time": "2026-07-08 12:00:00", + "http_user": "dolphin", + } + ] + ) + s = UvAnomalySentinel(ch=ch, vst=FakeVSTQuerier(), enabled={"d"}) + s.run() + assert not s.results[0].passed + assert any("dolphin_violet" in e for e in s.results[0].evidence) + + +# ═══════════════════════════════════════════════════════════════════════ +# E — Check (e): tripwire_ok=false occurrences +# ═══════════════════════════════════════════════════════════════════════ + + +def test_e_happy_path(): + """No tripwire anomaly events.""" + ch = FakeCHQuerier(anomaly_events=[]) + s = UvAnomalySentinel(ch=ch, vst=FakeVSTQuerier(), enabled={"e"}) + s.run() + assert s.results[0].passed + + +def test_e_mutation(): + """tripwire anomaly events exist → FAIL.""" + ch = FakeCHQuerier( + anomaly_events=[ + { + "ts": "2026-07-08 12:00:00", + "decision_id": "d-001", + "trade_id": "u-abc-001", + "symbol": "BTCUSDT", + "sensor": "AccountPublisher", + "detail": "wallet_balance mismatch", + } + ] + ) + s = UvAnomalySentinel(ch=ch, vst=FakeVSTQuerier(), enabled={"e"}) + s.run() + assert not s.results[0].passed + assert any("u-abc-001" in e for e in s.results[0].evidence) + + +# ═══════════════════════════════════════════════════════════════════════ +# F — Check coordination +# ═══════════════════════════════════════════════════════════════════════ + + +def test_disabled_check_not_run(): + """Disabled check does not appear in results.""" + vst = FakeVSTQuerier(open_orders=[PREFIX_ORDER]) + s = UvAnomalySentinel(ch=FakeCHQuerier(), vst=vst, enabled={"a", "b"}) + s.run() + ids = {r.check_id for r in s.results} + assert ids == {"a", "b"} + assert "c" not in ids + + +def test_exit_code_nonzero_on_fail(): + """Nonzero exit code when any check FAILs.""" + vst = FakeVSTQuerier(open_orders=[PREFIX_ORDER]) + s = UvAnomalySentinel(ch=FakeCHQuerier(), vst=vst, enabled={"c"}) + rc = s.run() + assert rc == 1 + + +def test_exit_code_zero_all_pass(): + """Zero exit code when all checks PASS.""" + vst = FakeVSTQuerier(open_orders=[U_ORDER_1]) + ch = FakeCHQuerier(exec_journal=[JOURNAL_1]) + s = UvAnomalySentinel(ch=ch, vst=vst, enabled={"a", "c"}) + rc = s.run() + assert rc == 0 + + +def test_multiple_checks_one_fails(): + """One failing check among many → exit code 1.""" + ch = FakeCHQuerier( + exec_journal=[JOURNAL_1], + anomaly_events=[ + { + "ts": "2026-07-08 12:00:00", + "decision_id": "d-001", + "trade_id": "u-abc-001", + "symbol": "BTCUSDT", + "sensor": "AccountPublisher", + "detail": "bad", + } + ], + ) + vst = FakeVSTQuerier(open_orders=[U_ORDER_1]) + s = UvAnomalySentinel(ch=ch, vst=vst, enabled={"a", "e"}) + rc = s.run() + assert rc == 1 + # Check (a) passes, (e) fails + assert s.results[0].passed + assert not s.results[1].passed + + +# ═══════════════════════════════════════════════════════════════════════ +# G — Mutation litmus: reintroduce a bug → test goes RED +# ═══════════════════════════════════════════════════════════════════════ + + +def test_mutation_litmus_disable_check_c(): + """Disabling check (c) lets a p- prefix pass → mutation RED. + Proves check (c) is the guard against non-u- prefixes.""" + prefix_journal = {"client_order_id": "p-e-abc123", "trade_id": "T-p-001", "status": "NEW"} + vst = FakeVSTQuerier(open_orders=[PREFIX_ORDER]) + ch = FakeCHQuerier(exec_journal=[prefix_journal]) + s = UvAnomalySentinel(ch=ch, vst=vst, enabled={"a", "b"}) + s.run() + # No check (c) in results — p- prefix goes undetected + assert "c" not in {r.check_id for r in s.results} + assert all(r.passed for r in s.results), ( + "Mutation litmus: disabling check (c) should allow p- prefix to pass" + ) + + +def test_mutation_litmus_check_a_ignores_missing_journal(): + """Without check (a), orphan venue orders go undetected.""" + ch = FakeCHQuerier(exec_journal=[]) + vst = FakeVSTQuerier(open_orders=[ORPHAN_VENUE_ORDER]) + # Only run check (b), not (a) + s = UvAnomalySentinel(ch=ch, vst=vst, enabled={"b"}) + s.run() + # (b) may pass since journal is empty — orphan still on venue + # This simulates: "if (a) were deleted, orphan goes undetected" + assert all(r.passed for r in s.results) + + +if __name__ == "__main__": + pytest.main([__file__, "-v", "--tb=short"]) diff --git a/prod/clean_arch/violet/uv/uv_anomaly_sentinel.py b/prod/clean_arch/violet/uv/uv_anomaly_sentinel.py new file mode 100644 index 0000000..56f1474 --- /dev/null +++ b/prod/clean_arch/violet/uv/uv_anomaly_sentinel.py @@ -0,0 +1,478 @@ +""" +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())