"""Centralized structured event store for live trading monitoring. All structlog events from the PEAD engine, ORB engine, pipeline subprocesses, and AutoScheduler flow into a single SQLite table here. This gives the user a queryable log of every fallback, timeout, error, and phase lifecycle event without touching any engine emit sites. """ from __future__ import annotations import json import logging import queue import sqlite3 import threading import time from datetime import datetime, timezone from pathlib import Path from typing import Any _DB_PATH = "journal/events.db" _SCHEMA = """ PRAGMA journal_mode=WAL; CREATE TABLE IF NOT EXISTS events ( id INTEGER PRIMARY KEY AUTOINCREMENT, ts_utc TEXT NOT NULL, job_run_id TEXT, source TEXT NOT NULL, level TEXT NOT NULL, category TEXT, event_name TEXT NOT NULL, session_id TEXT, ticker TEXT, message TEXT, details TEXT ); CREATE INDEX IF NOT EXISTS idx_events_ts ON events(ts_utc DESC); CREATE INDEX IF NOT EXISTS idx_events_run ON events(job_run_id, ts_utc DESC); CREATE INDEX IF NOT EXISTS idx_events_session ON events(session_id, ts_utc DESC); CREATE INDEX IF NOT EXISTS idx_events_cat_lvl ON events(category, level, ts_utc DESC); CREATE INDEX IF NOT EXISTS idx_events_name ON events(event_name, ts_utc DESC); """ # event_name suffix → level promotion to WARN _WARN_SUFFIXES = ( "_fallback", "_skipped", "_timeout", "_stale", "_failed", "_unavailable", "_not_found", "_soft_fallback", ) # event_name prefix → category _CATEGORY_MAP: list[tuple[str, str]] = [ ("paper_engine_macro_", "macro"), ("paper_engine_clock_", "health"), ("paper_engine_kill_switch_", "kill_switch"), ("paper_engine_buy_", "order"), ("paper_engine_close_", "order"), ("paper_engine_no_bar", "order"), ("paper_engine_", "engine"), ("orb_engine_buy_", "order"), ("orb_engine_close_", "order"), ("orb_engine_kill_switch_", "kill_switch"), ("orb_engine_circuit_breaker", "kill_switch"), ("orb_engine_oracle_", "broker"), ("orb_engine_bars_", "broker"), ("orb_engine_no_", "health"), ("orb_engine_", "engine"), ("snapshot_store_", "snapshot"), ("incremental_update_canonical_", "snapshot"), ("event_detector_", "snapshot"), ("oracle_", "broker"), ("multi_daily_bars", "broker"), ("multi_intraday_bars", "broker"), ("phase_started", "lifecycle"), ("phase_completed", "lifecycle"), ("scheduler_", "lifecycle"), ("pipeline_", "pipeline"), ] def _infer_category(event_name: str, explicit: str | None) -> str | None: if explicit: return explicit for prefix, cat in _CATEGORY_MAP: if event_name.startswith(prefix): return cat return None def _promote_level(level: str, event_name: str, category: str | None) -> str: """Promote INFO → WARN for known-fallback event names.""" if level in ("warning", "warn", "WARN", "WARNING"): return "WARN" if level in ("error", "critical", "ERROR", "CRITICAL"): return "ERROR" low = event_name.lower() if any(low.endswith(s) for s in _WARN_SUFFIXES): return "WARN" # fallback-category items that are INFO → WARN if category in ("kill_switch",) and level in ("info", "INFO"): return "WARN" return "INFO" class EventsStore: """Thread-safe event store backed by SQLite. Uses a background writer thread and a bounded queue so that SQLite I/O never blocks the trading hot-path. The writer thread is resilient to transient SQLite errors (locked / disk full) — it drops the failing row and counts the drop rather than dying. """ _instance: "EventsStore | None" = None _lock = threading.Lock() @classmethod def get(cls) -> "EventsStore": """Return the process-level singleton, creating it on first call.""" if cls._instance is None: with cls._lock: if cls._instance is None: cls._instance = cls(_DB_PATH) return cls._instance def __init__(self, db_path: str = _DB_PATH) -> None: self._db_path = str(Path(db_path)) self._queue: queue.Queue[dict[str, Any] | None] = queue.Queue(maxsize=10_000) self._drop_count = 0 self._schema_done = False self._writer = threading.Thread(target=self._writer_loop, daemon=True, name="events-writer") self._writer.start() def ensure_schema(self) -> None: """Create DB tables if they don't exist. Safe to call multiple times.""" if self._schema_done: return try: Path(self._db_path).parent.mkdir(parents=True, exist_ok=True) conn = sqlite3.connect(self._db_path, timeout=10) conn.executescript(_SCHEMA) conn.commit() conn.close() self._schema_done = True except Exception as exc: logging.getLogger(__name__).error("events_store schema init failed: %s", exc) def write(self, record: dict[str, Any]) -> None: """Enqueue a record for background writing. Non-blocking; drops on full queue.""" try: self._queue.put_nowait(record) except queue.Full: self._drop_count += 1 if self._drop_count % 100 == 1: print(f"[events_store] queue full, {self._drop_count} drops", flush=True, file=__import__("sys").stderr) def _writer_loop(self) -> None: conn: sqlite3.Connection | None = None BATCH = 50 FLUSH_EVERY = 2.0 # seconds def _connect() -> sqlite3.Connection | None: try: Path(self._db_path).parent.mkdir(parents=True, exist_ok=True) c = sqlite3.connect(self._db_path, timeout=15, check_same_thread=False) c.execute("PRAGMA journal_mode=WAL") return c except Exception as exc: print(f"[events_store] connect failed: {exc}", file=__import__("sys").stderr, flush=True) return None pending: list[dict[str, Any]] = [] last_flush = time.monotonic() while True: # Drain up to BATCH items from the queue try: record = self._queue.get(timeout=FLUSH_EVERY) if record is None: # poison pill break pending.append(record) # Drain additional items without waiting while len(pending) < BATCH: try: r = self._queue.get_nowait() if r is None: break pending.append(r) except queue.Empty: break except queue.Empty: pass now = time.monotonic() if not pending and now - last_flush < FLUSH_EVERY: continue if not pending: last_flush = now continue if conn is None: conn = _connect() if conn is None: pending.clear() continue rows_to_insert = [] for rec in pending: try: rows_to_insert.append(_build_row(rec)) except Exception: pass pending.clear() if not rows_to_insert: last_flush = now continue try: conn.executemany( "INSERT INTO events (ts_utc,job_run_id,source,level,category," "event_name,session_id,ticker,message,details) " "VALUES (?,?,?,?,?,?,?,?,?,?)", rows_to_insert, ) conn.commit() except Exception as exc: self._drop_count += len(rows_to_insert) print(f"[events_store] write failed ({len(rows_to_insert)} rows dropped): {exc}", file=__import__("sys").stderr, flush=True) try: conn.close() except Exception: pass conn = None last_flush = time.monotonic() def query( self, source: str | None = None, level: str | None = None, category: str | None = None, session_id: str | None = None, job_run_id: str | None = None, event_name: str | None = None, q: str | None = None, since: str | None = None, until: str | None = None, limit: int = 200, offset: int = 0, system: str | None = None, ) -> tuple[list[dict[str, Any]], int]: """Return (rows, total_count) matching the given filters.""" clauses: list[str] = [] params: list[Any] = [] if system == "pead": srcs = ("paper_engine", "snapshot_store", "event_detector", "oracle", "pipeline") ph = ",".join("?" * len(srcs)) clauses.append(f"(source IN ({ph}) OR (source='auto_scheduler' AND details NOT LIKE ?))") params.extend(srcs) params.append('%"scheduler": "orb"%') elif system == "orb": clauses.append("(source IN (?) OR (source='auto_scheduler' AND details LIKE ?))") params.append("orb_engine") params.append('%"scheduler": "orb"%') if source: clauses.append("source = ?") params.append(source) if level: clauses.append("level = ?") params.append(level.upper()) if category: clauses.append("category = ?") params.append(category) if session_id: clauses.append("session_id = ?") params.append(session_id) if job_run_id: clauses.append("job_run_id = ?") params.append(job_run_id) if event_name: clauses.append("event_name LIKE ?") params.append(f"%{event_name}%") if q: clauses.append("(message LIKE ? OR event_name LIKE ? OR details LIKE ?)") params += [f"%{q}%", f"%{q}%", f"%{q}%"] if since: clauses.append("ts_utc >= ?") params.append(since) if until: clauses.append("ts_utc <= ?") params.append(until) where = ("WHERE " + " AND ".join(clauses)) if clauses else "" try: conn = sqlite3.connect(self._db_path, timeout=10) conn.row_factory = sqlite3.Row total = conn.execute(f"SELECT COUNT(*) FROM events {where}", params).fetchone()[0] rows = conn.execute( f"SELECT * FROM events {where} ORDER BY ts_utc DESC LIMIT ? OFFSET ?", params + [limit, offset], ).fetchall() conn.close() return [dict(r) for r in rows], total except Exception as exc: logging.getLogger(__name__).warning("events_store query failed: %s", exc) return [], 0 def recent_runs(self, limit: int = 20) -> list[dict[str, Any]]: """Return recent job_run_id summaries ordered by start time.""" sql = """ SELECT job_run_id, MIN(ts_utc) AS started_at, MAX(ts_utc) AS ended_at, COUNT(*) AS event_count, SUM(CASE WHEN level='ERROR' THEN 1 ELSE 0 END) AS error_count, SUM(CASE WHEN level='WARN' THEN 1 ELSE 0 END) AS warn_count, GROUP_CONCAT(DISTINCT source) AS sources, MAX(CASE WHEN event_name='phase_started' THEN json_extract(details,'$.phase') END) AS phase FROM events WHERE job_run_id IS NOT NULL AND job_run_id != '' GROUP BY job_run_id ORDER BY started_at DESC LIMIT ? """ try: conn = sqlite3.connect(self._db_path, timeout=10) conn.row_factory = sqlite3.Row rows = conn.execute(sql, [limit]).fetchall() conn.close() return [dict(r) for r in rows] except Exception as exc: logging.getLogger(__name__).warning("events_store recent_runs failed: %s", exc) return [] def health_summary(self) -> dict[str, Any]: """Aggregate today's events into a health snapshot for the UI.""" today_start = datetime.now(timezone.utc).strftime("%Y-%m-%dT00:00:00") phase_sql = """ SELECT json_extract(details,'$.phase') AS phase, json_extract(details,'$.status') AS status, MAX(ts_utc) AS last_ts FROM events WHERE event_name='phase_completed' AND ts_utc >= ? GROUP BY phase, status """ fallback_sql = """ SELECT category, level, COUNT(*) AS cnt FROM events WHERE ts_utc >= ? AND ( level IN ('WARN','ERROR') OR category IN ('fallback','snapshot','macro','regime','broker','kill_switch') ) AND category NOT IN ('lifecycle','pipeline') GROUP BY category, level """ recent_errors_sql = """ SELECT id, ts_utc, source, event_name, message, session_id, job_run_id FROM events WHERE level='ERROR' AND ts_utc >= ? ORDER BY ts_utc DESC LIMIT 5 """ snapshot_sql = """ SELECT session_id, json_extract(details,'$.snapshot_id') AS snapshot_id, MAX(ts_utc) AS last_update FROM events WHERE event_name LIKE 'incremental_update_canonical%done' AND ts_utc >= date('now','-7 days') GROUP BY session_id, snapshot_id """ try: conn = sqlite3.connect(self._db_path, timeout=10) conn.row_factory = sqlite3.Row phases = [dict(r) for r in conn.execute(phase_sql, [today_start]).fetchall()] fallbacks = [dict(r) for r in conn.execute(fallback_sql, [today_start]).fetchall()] recent_errors = [dict(r) for r in conn.execute(recent_errors_sql, [today_start]).fetchall()] snapshots = [dict(r) for r in conn.execute(snapshot_sql, []).fetchall()] conn.close() except Exception as exc: logging.getLogger(__name__).warning("events_store health_summary failed: %s", exc) phases, fallbacks, recent_errors, snapshots = [], [], [], [] phase_map: dict[str, dict[str, Any]] = {} for row in phases: ph = row["phase"] or "unknown" if ph not in phase_map or row["last_ts"] > phase_map[ph].get("last_ts", ""): phase_map[ph] = row fallback_map: dict[str, dict[str, Any]] = {} for row in fallbacks: cat = row["category"] or "other" if cat not in fallback_map: fallback_map[cat] = {"warn": 0, "error": 0} if row["level"] == "ERROR": fallback_map[cat]["error"] += row["cnt"] else: fallback_map[cat]["warn"] += row["cnt"] return { "phases": phase_map, "fallbacks_today": fallback_map, "recent_errors": recent_errors, "snapshot_freshness": snapshots, "drop_count": self._drop_count, } def distinct_sources(self) -> list[str]: try: conn = sqlite3.connect(self._db_path, timeout=10) rows = conn.execute( "SELECT DISTINCT source FROM events ORDER BY source" ).fetchall() conn.close() return [r[0] for r in rows] except Exception: return [] def purge_before(self, before_iso: str) -> int: try: conn = sqlite3.connect(self._db_path, timeout=10) cur = conn.execute("DELETE FROM events WHERE ts_utc < ?", [before_iso]) conn.commit() deleted = cur.rowcount conn.execute("VACUUM") conn.close() return deleted except Exception as exc: logging.getLogger(__name__).warning("events_store purge failed: %s", exc) return 0 def _build_row(rec: dict[str, Any]) -> tuple: """Convert a structlog event_dict or freeform dict into a DB row tuple.""" ts = rec.get("timestamp") or rec.get("ts_utc") or datetime.now(timezone.utc).isoformat() event_name = str(rec.get("event") or rec.get("event_name") or "") raw_level = str(rec.get("level") or "info").lower() explicit_cat = rec.get("category") category = _infer_category(event_name, explicit_cat) level = _promote_level(raw_level, event_name, category) # "fallback" refinement: any event_name containing "fallback" → category fallback if "fallback" in event_name.lower() and category not in ("lifecycle",): category = "fallback" message = rec.get("message") or rec.get("msg") or event_name details_dict = {k: v for k, v in rec.items() if k not in ("timestamp", "level", "event", "event_name", "message", "msg", "ts_utc", "source", "job_run_id", "session_id", "ticker", "category")} try: details = json.dumps(details_dict, default=str) except Exception: details = str(details_dict) return ( ts, rec.get("job_run_id") or "", str(rec.get("source") or "unknown"), level, category, event_name, rec.get("session_id"), rec.get("ticker"), message, details, ) # ── structlog processor ─────────────────────────────────────────────────────── def structlog_sink_processor(logger: Any, method: str, event_dict: dict[str, Any]) -> dict[str, Any]: """structlog processor that tees events to EventsStore without blocking.""" try: store = EventsStore.get() # Infer source from logger name if not already set source = event_dict.get("source") if not source: logger_name = str(logger.name if hasattr(logger, "name") else "") source = _infer_source(logger_name, event_dict.get("event", "")) record = dict(event_dict) record["source"] = source store.write(record) except Exception: pass return event_dict def _infer_source(logger_name: str, event_name: str) -> str: if "paper_trader" in logger_name or event_name.startswith("paper_engine_"): return "paper_engine" if "orb_trader" in logger_name or event_name.startswith("orb_"): return "orb_engine" if "snapshot_store" in logger_name or event_name.startswith("snapshot_store_"): return "snapshot_store" if "event_detector" in logger_name or event_name.startswith("event_detector_"): return "event_detector" if "oracle" in logger_name or event_name.startswith(("oracle_", "multi_")): return "oracle" if "canonical_snapshot" in logger_name or event_name.startswith("incremental_update"): return "snapshot_store" return logger_name.split(".")[-1] if logger_name else "unknown" # ── stdlib logging bridge ───────────────────────────────────────────────────── _STDLIB_LEVEL_MAP = { logging.DEBUG: "INFO", logging.INFO: "INFO", logging.WARNING: "WARN", logging.ERROR: "ERROR", logging.CRITICAL: "ERROR", } class EventsLoggingHandler(logging.Handler): """stdlib logging.Handler that writes to EventsStore. Captures ORB engine (stdlib), oracle client (stdlib), and web.main auto-restart warnings that structlog doesn't see. """ # Loggers to skip (already handled by structlog or too noisy) _SKIP_LOGGERS = { "uvicorn", "uvicorn.error", "uvicorn.access", "fastapi", "asyncio", "multiprocessing", "sqlalchemy", } def emit(self, record: logging.LogRecord) -> None: logger_name = record.name or "" # Skip loggers that are too noisy or already handled by structlog root = logger_name.split(".")[0] if root in self._SKIP_LOGGERS: return # Skip DEBUG unless it's a known important logger if record.levelno < logging.WARNING and record.levelno == logging.DEBUG: return try: store = EventsStore.get() event_name = f"stdlib_{logger_name.replace('.', '_')}" message = record.getMessage() level = _STDLIB_LEVEL_MAP.get(record.levelno, "INFO") # Infer a better event_name from known patterns in the message low_msg = message.lower() if "fallback" in low_msg: event_name = "stdlib_fallback" elif "failed" in low_msg or "error" in low_msg: event_name = f"stdlib_error_{root}" elif "auto-restart" in low_msg or "auto_restart" in low_msg: event_name = "scheduler_auto_restart" store.write({ "ts_utc": datetime.now(timezone.utc).isoformat(), "source": _infer_source(logger_name, event_name), "level": level, "event": event_name, "message": message, "logger": logger_name, }) except Exception: pass