From b8cfd5df675f2e1726185ed0d045d0e331c751b1 Mon Sep 17 00:00:00 2001 From: Codex Date: Sun, 12 Jul 2026 03:28:18 +0200 Subject: [PATCH] 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. --- .../violet/uv/test_uv_anomaly_sentinel.py | 437 ++++++++++++++---- .../violet/uv/uv_anomaly_sentinel.py | 273 ++++++----- 2 files changed, 510 insertions(+), 200 deletions(-) diff --git a/prod/clean_arch/violet/uv/test_uv_anomaly_sentinel.py b/prod/clean_arch/violet/uv/test_uv_anomaly_sentinel.py index 23b36ed..e2cdbed 100644 --- a/prod/clean_arch/violet/uv/test_uv_anomaly_sentinel.py +++ b/prod/clean_arch/violet/uv/test_uv_anomaly_sentinel.py @@ -1,107 +1,236 @@ """ -UV Anomaly Sentinel tests — synthetic fixtures per check + mutation. +UV Anomaly Sentinel tests — synthetic fixtures per check + mutation + error. Each check has: - happy_path: normal data → PASS - mutation: inject failure condition → FAIL + - error: query error → FAIL (never silent PASS) - Mutation litmus: disable a check → test goes RED (proves check fires) -Also: disabled-check no-op, exit-code nonzero on any FAIL. +Also: schema-pinned fixture mirroring DESCRIBE dolphin_uv.exec_journal, +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, + CHQueryResult, CheckResult, UvAnomalySentinel, VSTQuerier, + VSTResult, ) # ═══════════════════════════════════════════════════════════════════════ -# Fake backends +# Fake backends — real schema: exec_journal(timestamp, u_prefix_client_id, +# suppressed, suppress_reason, ...), VST: {"data": {"orders": [...]}} # ═══════════════════════════════════════════════════════════════════════ class FakeCHQuerier: - """Synthetic ClickHouse backend — returns pre-set responses by table.""" + """Synthetic CH backend — returns pre-set data 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() + self._anomaly_rows: list[dict[str, Any]] = tables.get("anomaly_events", []) + self._query_log_rows: list[dict[str, Any]] = tables.get("query_log", []) + self._schema_rows: dict[str, list[dict[str, str]]] = {} + self._describe_schema: dict[str, list[dict[str, str]]] = {} if "exec_journal" in tables: - self._exists_tables.add("exec_journal") + self._exists.add("exec_journal") if "anomaly_events" in tables: - self._exists_tables.add("anomaly_events") + self._exists.add("anomaly_events") + # error_mode: if True, all queries return error + self._error_mode: bool = tables.get("_error_mode", False) + self._error_msg: str = tables.get("_error_msg", "simulated query error") - def query(self, sql: str, **kw: Any) -> list[dict[str, Any]]: + def query(self, sql: str, **kw: Any) -> CHQueryResult: + if self._error_mode: + return CHQueryResult([], self._error_msg) sql_lower = sql.lower() - if "query_log" in sql_lower: - return self._query_log_rows + if "exec_journal" in sql_lower and "describe" not in sql_lower: + # Simulate SQL "suppressed = 0" filter if present + if "suppressed = 0" in sql or "suppressed=0" in sql: + return CHQueryResult( + [r for r in self._journal_rows if r.get("suppressed") == 0], + "", + ) + return CHQueryResult(self._journal_rows, "") if "anomaly_events" in sql_lower: - return self._anomaly_rows - if "exec_journal" in sql_lower: - return self._journal_rows + return CHQueryResult(self._anomaly_rows, "") + if "query_log" in sql_lower: + # Simulate current_database != exclusion filters + # Extract db names to exclude from SQL + excluded: set[str] = set() + for part in sql.split("current_database != "): + if "'" not in part: + continue + db = part.split("'")[1] if "'" in part else "" + if db: + excluded.add(db) + if excluded: + filtered = [ + r for r in self._query_log_rows + if r.get("db", "") not in excluded + ] + return CHQueryResult(filtered, "") + return CHQueryResult(self._query_log_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: + for table in self._exists: if table in sql: - return [{"1": 1}] - return [] - return [] + return CHQueryResult([{"1": 1}], "") + return CHQueryResult([], "") + if "describe table" in sql_lower: + table_name = "" + for t in self._describe_schema: + if t in sql_lower: + return CHQueryResult(self._describe_schema[t], "") + return CHQueryResult([], "") + # default: match by looking up in _tables + for table_name, rows in self._tables.items(): + if table_name.lower() in sql_lower: + return CHQueryResult(rows, "") + return CHQueryResult([], "") - def describe_table(self, table: str, db: str = "") -> list[dict[str, str]]: - return self._describe_schemas.get(table, []) + def describe_table(self, table: str, db: str = "") -> CHQueryResult: + if table in self._describe_schema: + return CHQueryResult(self._describe_schema[table], "") + # Use journal_rows first row keys if available + if self._journal_rows: + keys = list(self._journal_rows[0].keys()) + schema = [{"name": k, "type": "String"} for k in keys] + self._describe_schema[table] = schema + return CHQueryResult(schema, "") + return CHQueryResult([], "") - def table_exists(self, table: str, db: str = "") -> bool: - return table in self._exists_tables + def table_exists(self, table: str, db: str = "") -> CHQueryResult: + if self._error_mode: + return CHQueryResult([], self._error_msg) + if table in self._exists: + return CHQueryResult([{"1": 1}], "") + return CHQueryResult([], "") class FakeVSTQuerier: - """Synthetic VST backend — returns pre-set open orders.""" + """Synthetic VST backend — returns pre-set open orders or error.""" - def __init__(self, open_orders: list[dict[str, Any]] | None = None) -> None: + def __init__( + self, + open_orders: list[dict[str, Any]] | None = None, + error: str = "", + ) -> None: self._open_orders = open_orders or [] + self._error = error - def open_orders(self) -> list[dict[str, Any]]: - return self._open_orders + def open_orders(self) -> VSTResult: + if self._error: + return VSTResult([], self._error) + return VSTResult(self._open_orders, "") - def open_positions(self) -> list[dict[str, Any]]: - return [] + def open_positions(self) -> VSTResult: + return VSTResult([], "") - def account_balance(self) -> dict[str, Any]: - return {} + def account_balance(self) -> VSTResult: + return VSTResult([], "") # ═══════════════════════════════════════════════════════════════════════ -# Shared fixture data +# Schema-pinned fixtures — mirror real DESCRIBE output +# ═══════════════════════════════════════════════════════════════════════ + +EXEC_JOURNAL_SCHEMA = [ + {"name": "kind", "type": "String"}, + {"name": "suppressed", "type": "UInt8"}, + {"name": "timestamp", "type": "DateTime64(3)"}, + {"name": "scan_number", "type": "UInt64"}, + {"name": "intent_id", "type": "String"}, + {"name": "trade_id", "type": "String"}, + {"name": "slot_id", "type": "UInt32"}, + {"name": "asset", "type": "String"}, + {"name": "side", "type": "String"}, + {"name": "action", "type": "String"}, + {"name": "reference_price", "type": "Float64"}, + {"name": "target_size", "type": "Float64"}, + {"name": "leverage", "type": "Float64"}, + {"name": "has_uv_prefix", "type": "UInt8"}, + {"name": "u_prefix_client_id", "type": "String"}, + {"name": "suppress_reason", "type": "String"}, + {"name": "promo_metadata", "type": "String"}, +] + + +ANOMALY_EVENTS_SCHEMA = [ + {"name": "ts", "type": "DateTime64(6, 'UTC')"}, + {"name": "ts_day", "type": "Date"}, + {"name": "decision_id", "type": "String"}, + {"name": "trade_id", "type": "String"}, + {"name": "symbol", "type": "LowCardinality(String)"}, + {"name": "anomaly", "type": "LowCardinality(String)"}, + {"name": "origin", "type": "LowCardinality(String)"}, + {"name": "sensor", "type": "LowCardinality(String)"}, + {"name": "detail", "type": "String"}, + {"name": "rm_meta", "type": "Float32"}, +] + + +def test_schema_pinned_exec_journal(): + """Fixture mirrors live DESCRIBE dolphin_uv.exec_journal — 17 columns.""" + assert len(EXEC_JOURNAL_SCHEMA) == 17 + names = {c["name"] for c in EXEC_JOURNAL_SCHEMA} + assert "timestamp" in names, "REAL schema uses timestamp not ts" + assert "u_prefix_client_id" in names, "REAL schema uses u_prefix_client_id" + assert "suppressed" in names + assert "suppress_reason" in names + + +def test_schema_pinned_anomaly_events(): + """Fixture mirrors live DESCRIBE dolphin_violet.anomaly_events — 10 columns.""" + assert len(ANOMALY_EVENTS_SCHEMA) == 10 + names = {c["name"] for c in ANOMALY_EVENTS_SCHEMA} + assert "ts" in names, "REAL schema uses ts not timestamp" + assert "anomaly" in names + assert "trade_id" in names + assert "detail" in names + + +# ═══════════════════════════════════════════════════════════════════════ +# Shared fixture data (real column names) # ═══════════════════════════════════════════════════════════════════════ 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"} +JOURNAL_1 = {"u_prefix_client_id": "u-abc-001", "trade_id": "T-u-001", + "suppressed": 0, "timestamp": "2026-07-08 12:00:00"} +JOURNAL_2 = {"u_prefix_client_id": "u-def-002", "trade_id": "T-u-002", + "suppressed": 0, "timestamp": "2026-07-08 12:01:00"} 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"} +ORPHAN_JOURNAL = {"u_prefix_client_id": "u-stale-888", "trade_id": "T-stale", + "suppressed": 0, "timestamp": "2026-07-08 10:00:00"} +SUPPRESSED_JOURNAL = {"u_prefix_client_id": "u-sup-777", "trade_id": "T-sup", + "suppressed": 1, "timestamp": "2026-07-08 11:00:00"} +# anomaly_events rows — real schema uses 'ts' not 'timestamp' +TRIPWIRE_EVENT = { + "ts": "2026-07-08 12:00:00", + "decision_id": "d-001", + "trade_id": "u-abc-001", + "symbol": "BTCUSDT", + "sensor": "AccountPublisher", + "detail": "wallet_balance mismatch", +} +EMPTY_WRAPPER = {} # VST may return empty dict wrappers that must be filtered # ═══════════════════════════════════════════════════════════════════════ @@ -124,12 +253,13 @@ def test_a_mutation(): 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) + r = s.results[0] + assert not r.passed + assert any("u-orphan" in e for e in r.evidence) def test_a_mutation_no_journal_table(): - """No exec_journal table → FAIL.""" + """No exec_journal table → FAIL with evidence.""" ch = FakeCHQuerier() vst = FakeVSTQuerier(open_orders=[U_ORDER_1]) s = UvAnomalySentinel(ch=ch, vst=vst, enabled={"a"}) @@ -138,13 +268,43 @@ def test_a_mutation_no_journal_table(): assert any("exec_journal" in e for e in s.results[0].evidence) +def test_a_error_ch(): + """CH query error → FAIL, never silent PASS.""" + ch = FakeCHQuerier(_error_mode=True, _error_msg="CH unreachable") + vst = FakeVSTQuerier(open_orders=[U_ORDER_1]) + s = UvAnomalySentinel(ch=ch, vst=vst, enabled={"a"}) + s.run() + assert not s.results[0].passed + assert "error" in s.results[0].error.lower() or "unreachable" in s.results[0].error + + +def test_a_error_vst(): + """VST error → FAIL, never silent PASS.""" + ch = FakeCHQuerier() + vst = FakeVSTQuerier(open_orders=[U_ORDER_1], error="VST connection refused") + s = UvAnomalySentinel(ch=ch, vst=vst, enabled={"a"}) + s.run() + assert not s.results[0].passed + assert "error" in s.results[0].error.lower() or "connection" in s.results[0].error + + +def test_a_empty_wrapper_filtered(): + """Empty dict wrapper in VST response → filtered, not flagged as order.""" + ch = FakeCHQuerier(exec_journal=[JOURNAL_1]) + vst = FakeVSTQuerier(open_orders=[U_ORDER_1, EMPTY_WRAPPER]) + s = UvAnomalySentinel(ch=ch, vst=vst, enabled={"a"}) + s.run() + # Empty wrapper must not trigger false-positive orphan detection + assert s.results[0].passed + + # ═══════════════════════════════════════════════════════════════════════ # B — Check (b): Journal row without venue order # ═══════════════════════════════════════════════════════════════════════ def test_b_happy_path(): - """All journal entries have matching venue orders.""" + """All non-suppressed 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"}) @@ -158,8 +318,18 @@ def test_b_mutation(): 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) + r = s.results[0] + assert not r.passed + assert any("u-stale" in e for e in r.evidence) + + +def test_b_suppressed_not_flagged(): + """Suppressed journal entry (suppressed=1) → not flagged.""" + ch = FakeCHQuerier(exec_journal=[JOURNAL_1, SUPPRESSED_JOURNAL]) + vst = FakeVSTQuerier(open_orders=[U_ORDER_1]) + s = UvAnomalySentinel(ch=ch, vst=vst, enabled={"b"}) + s.run() + assert s.results[0].passed def test_b_mutation_no_table(): @@ -171,6 +341,16 @@ def test_b_mutation_no_table(): assert s.results[0].passed +def test_b_error_ch(): + """CH query error → FAIL, never silent PASS.""" + ch = FakeCHQuerier(_error_mode=True) + vst = FakeVSTQuerier(open_orders=[U_ORDER_1]) + s = UvAnomalySentinel(ch=ch, vst=vst, enabled={"b"}) + s.run() + assert not s.results[0].passed + assert s.results[0].error + + # ═══════════════════════════════════════════════════════════════════════ # C — Check (c): u- prefix on clientOrderId # ═══════════════════════════════════════════════════════════════════════ @@ -203,12 +383,21 @@ def test_c_mutation_no_prefix(): def test_c_mutation_empty_client_order_id(): - """Order with empty clientOrderId → not flagged (no prefix to check).""" + """Order with empty clientOrderId → not flagged.""" 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 + assert s.results[0].passed + + +def test_c_error_vst(): + """VST error → FAIL, never silent PASS.""" + vst = FakeVSTQuerier(open_orders=[], error="VST API error code=-1000") + s = UvAnomalySentinel(ch=FakeCHQuerier(), vst=vst, enabled={"c"}) + s.run() + assert not s.results[0].passed + assert s.results[0].error # ═══════════════════════════════════════════════════════════════════════ @@ -217,7 +406,23 @@ def test_c_mutation_empty_client_order_id(): def test_d_happy_path(): - """No unexpected INSERTs in query_log.""" + """INSERTs only into dolphin_uv or allowlisted DBs → PASS.""" + ch = FakeCHQuerier( + query_log=[ + { + "query": "INSERT INTO dolphin.obf_universe VALUES ...", + "event_time": "2026-07-08 12:00:00", + "db": "dolphin", + } + ] + ) + s = UvAnomalySentinel(ch=ch, vst=FakeVSTQuerier(), enabled={"d"}) + s.run() + assert s.results[0].passed + + +def test_d_happy_path_no_inserts_at_all(): + """No INSERTs (or only dolphin_uv INSERTs) → PASS.""" ch = FakeCHQuerier(query_log=[]) s = UvAnomalySentinel(ch=ch, vst=FakeVSTQuerier(), enabled={"d"}) s.run() @@ -225,22 +430,50 @@ def test_d_happy_path(): def test_d_mutation(): - """INSERT into non-uv database → FAIL.""" + """INSERT into non-allowlisted database → FAIL.""" ch = FakeCHQuerier( query_log=[ { - "query": "INSERT INTO dolphin_violet.anomaly_events VALUES ...", - "database": "dolphin_violet", - "table": "anomaly_events", + "query": "INSERT INTO dolphin_violet.anomaly_events ...", "event_time": "2026-07-08 12:00:00", - "http_user": "dolphin", + "db": "dolphin_violet", } ] ) 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) + assert any("INSERT" in e for e in s.results[0].evidence) + + +def test_d_allowlist_not_flagged(): + """INSERT into allowlisted DB (dolphin) → not flagged.""" + ch = FakeCHQuerier( + query_log=[ + { + "query": "INSERT INTO dolphin.obf_universe VALUES ...", + "event_time": "2026-07-08 12:00:00", + "db": "dolphin", + }, + { + "query": "INSERT INTO dolphin_malkhut.something VALUES ...", + "event_time": "2026-07-08 12:00:00", + "db": "dolphin_malkhut", + }, + ] + ) + s = UvAnomalySentinel(ch=ch, vst=FakeVSTQuerier(), enabled={"d"}) + s.run() + assert s.results[0].passed + + +def test_d_error_ch(): + """CH query_log error → FAIL, never silent PASS.""" + ch = FakeCHQuerier(_error_mode=True) + s = UvAnomalySentinel(ch=ch, vst=FakeVSTQuerier(), enabled={"d"}) + s.run() + assert not s.results[0].passed + assert s.results[0].error # ═══════════════════════════════════════════════════════════════════════ @@ -258,24 +491,22 @@ def test_e_happy_path(): 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", - } - ] - ) + ch = FakeCHQuerier(anomaly_events=[TRIPWIRE_EVENT]) 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) +def test_e_error_ch(): + """CH anomaly_events error → FAIL, never silent PASS.""" + ch = FakeCHQuerier(_error_mode=True) + s = UvAnomalySentinel(ch=ch, vst=FakeVSTQuerier(), enabled={"e"}) + s.run() + assert not s.results[0].passed + assert s.results[0].error + + # ═══════════════════════════════════════════════════════════════════════ # F — Check coordination # ═══════════════════════════════════════════════════════════════════════ @@ -284,7 +515,8 @@ def test_e_mutation(): 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"}) + ch = FakeCHQuerier(exec_journal=[JOURNAL_1]) + s = UvAnomalySentinel(ch=ch, vst=vst, enabled={"a", "b"}) s.run() ids = {r.check_id for r in s.results} assert ids == {"a", "b"} @@ -312,56 +544,67 @@ 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", - } - ], + anomaly_events=[TRIPWIRE_EVENT], ) 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 + assert s.results[0].passed # check (a) passes + assert not s.results[1].passed # check (e) fails # ═══════════════════════════════════════════════════════════════════════ -# G — Mutation litmus: reintroduce a bug → test goes RED +# G — Error handling: query errors never silent PASS +# ═══════════════════════════════════════════════════════════════════════ + + +def test_error_makes_check_fail(): + """Every check: if CH query errors → FAIL, never PASS.""" + for check_id in ("a", "b", "c", "d", "e"): + ch = FakeCHQuerier(_error_mode=True) + vst = FakeVSTQuerier(open_orders=[U_ORDER_1], error="") + # For check (c), VST error is the primary path + if check_id == "c": + vst = FakeVSTQuerier(open_orders=[], error="VST error") + s = UvAnomalySentinel( + ch=ch, vst=vst, enabled={check_id} + ) + s.run() + r = s.results[0] + assert not r.passed, f"Check {check_id} must FAIL on error, got PASS" + assert r.error, f"Check {check_id} must have error message" + + +# ═══════════════════════════════════════════════════════════════════════ +# H — 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"} + """Disabling check (c) lets a p- prefix pass → mutation RED.""" + prefix_journal = {"u_prefix_client_id": "p-e-abc123", "trade_id": "T-p-001", + "suppressed": 0, "timestamp": "2026-07-08 12:00:00"} 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" + "Mutation litmus: disabling (c) should let p- prefix 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"}) +def test_mutation_litmus_remove_error_check(): + """If error check is removed → query error becomes silent PASS.""" + # This litmus proves the error guard works: with it, error → FAIL. + ch = FakeCHQuerier(_error_mode=True) + vst = FakeVSTQuerier(open_orders=[U_ORDER_1]) + s = UvAnomalySentinel(ch=ch, vst=vst, enabled={"a"}) 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) + # Error guard is active → FAIL + assert not s.results[0].passed + assert s.results[0].error if __name__ == "__main__": diff --git a/prod/clean_arch/violet/uv/uv_anomaly_sentinel.py b/prod/clean_arch/violet/uv/uv_anomaly_sentinel.py index 56f1474..38c7b97 100644 --- a/prod/clean_arch/violet/uv/uv_anomaly_sentinel.py +++ b/prod/clean_arch/violet/uv/uv_anomaly_sentinel.py @@ -48,6 +48,10 @@ 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 ──────────────────────────────────────────────────────────── @@ -59,10 +63,7 @@ class CheckResult: description: str passed: bool = True evidence: list[str] = field(default_factory=list) - - @property - def status(self) -> str: - return "PASS" if self.passed else "FAIL" + 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: @@ -70,20 +71,31 @@ class CheckResult: 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}") - 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) ─────────────────────────────────── +# ─── 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).""" + """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, @@ -95,7 +107,7 @@ class CHQuerier: self._user = user self._password = password - def query(self, sql: str, **kw: Any) -> list[dict[str, Any]]: + 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") @@ -109,13 +121,12 @@ class CHQuerier: 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 [] + detail = exc.read().decode()[:300] + return CHQueryResult([], f"CH HTTP {exc.code}: {detail}") except OSError as exc: - log.error("CH connection error: %s", exc) - return [] + return CHQueryResult([], f"CH connection error: {exc}") if not body.strip(): - return [] + return CHQueryResult([], "") rows: list[dict[str, Any]] = [] for line in body.strip().split("\n"): line = line.strip() @@ -124,18 +135,20 @@ class CHQuerier: try: rows.append(json.loads(line)) except json.JSONDecodeError: - log.warning("CH non-JSON line: %.100s", line) - return rows + return CHQueryResult([], f"CH non-JSON response line: {line[:100]}") + return CHQueryResult(rows, "") - def describe_table(self, table: str, db: str = CH_DB_UV) -> list[dict[str, str]]: + 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) -> bool: - rows = self.query( + 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}'" ) - return len(rows) > 0 + if r.error: + return r + return CHQueryResult(r.rows, "") # ─── VST querier (read-only: open orders, positions) ─────────────────────── @@ -164,6 +177,13 @@ def _build_signed_params( 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.""" @@ -177,7 +197,7 @@ class VSTQuerier: 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]]: + 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}" @@ -187,50 +207,40 @@ class VSTQuerier: 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 [] + detail = exc.read().decode()[:300] + return VSTResult([], f"VST HTTP {exc.code}: {detail}") except OSError as exc: - log.error("VST connection error: %s", exc) - return [] + return VSTResult([], f"VST connection error: {exc}") 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 [] - ) + 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: - data = parsed - if isinstance(data, list): - return data - return [data] if 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) -> 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_orders(self) -> VSTResult: + return self._signed_get("/openApi/swap/v2/trade/openOrders") - def open_positions(self) -> list[dict[str, Any]]: + def open_positions(self) -> VSTResult: 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 {} + def account_balance(self) -> VSTResult: + r = self._signed_get("/openApi/swap/v2/user/balance") + return r # ─── Sentinel ────────────────────────────────────────────────────────────── @@ -246,16 +256,16 @@ class UvAnomalySentinel: 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] = [] - # ── Public API ────────────────────────────────────────────────────── - def run(self) -> int: checks: list[tuple[str, Any]] = [ ("a", self._check_venue_without_journal), @@ -275,15 +285,29 @@ class UvAnomalySentinel: 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): + + # 1. Fetch venue open orders + vst_r = self.vst.open_orders() + if vst_r.error: r.passed = False - r.add("exec_journal table does not exist in dolphin_uv") + r.error = vst_r.error return r - # Collect venue clientOrderIds + 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( @@ -294,24 +318,28 @@ class UvAnomalySentinel: ) 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 + + # 4. Query exec_journal for matching u_prefix_client_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}')" + f"SELECT u_prefix_client_id, trade_id " + f"FROM {CH_DB_UV}.exec_journal " + f"WHERE u_prefix_client_id IN ('{ids_escaped}')" ) - journal_rows = self.ch.query(sql) + jr = self.ch.query(sql) + if jr.error: + r.passed = False + r.error = jr.error + return r + journal_ids: set[str] = set() - for jr in journal_rows: - cid = str(jr.get("client_order_id") or "") + 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 @@ -322,27 +350,47 @@ class UvAnomalySentinel: 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 + + # 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 - # Get recent journal rows with open state + if not exists_r.rows: + return r + + # 2. Get recent journal rows (not suppressed = still live) sql = ( - f"SELECT client_order_id, trade_id, status " + f"SELECT u_prefix_client_id, trade_id, suppressed " f"FROM {CH_DB_UV}.exec_journal " - f"WHERE ts >= now() - INTERVAL {self.lookback_hours} HOUR " - f"ORDER BY ts DESC" + f"WHERE timestamp >= now() - INTERVAL {self.lookback_hours} HOUR " + f"AND suppressed = 0 " + f"ORDER BY timestamp DESC" ) - journal_rows = self.ch.query(sql) - if not journal_rows: + 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 jr in journal_rows: - cid = str(jr.get("client_order_id") or "") + for row in jr.rows: + cid = str(row.get("u_prefix_client_id") or "") if cid: journal_ids.add(cid) - orders = self.vst.open_orders() + + # 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 orders: + for o in vst_r.items: cid = str( o.get("clientOrderId") or o.get("clientOrderID") @@ -351,19 +399,24 @@ class UvAnomalySentinel: ) 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") + 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") - orders = self.vst.open_orders() + 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 orders: + for o in vst_r.items: cid = str( o.get("clientOrderId") or o.get("clientOrderID") @@ -382,29 +435,39 @@ class UvAnomalySentinel: 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 + 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, database, table, event_time, http_user " + f"SELECT query, {col} AS db, event_time " f"FROM system.query_log " f"WHERE type = 'QueryFinish' " f"AND query LIKE 'INSERT%' " - f"AND database != '{CH_DB_UV}' " + f"AND {db_filter} " 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: + qr = self.ch.query(sql) + if qr.error: r.passed = False - for row in rows: - db = row.get("database", "?") - tbl = row.get("table", "?") + 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}.{tbl}: {q}") + r.add(f"INSERT into {db}: {q}") return r def _check_tripwire_ok(self) -> CheckResult: - """(e) tripwire_ok=false occurrences in anomaly_events.""" + """(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 " @@ -414,10 +477,14 @@ class UvAnomalySentinel: f"ORDER BY ts DESC " f"LIMIT {EVIDENCE_LIMIT}" ) - rows = self.ch.query(sql) - if rows: + qr = self.ch.query(sql) + if qr.error: r.passed = False - for row in rows: + 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]