You cannot select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
577 lines
22 KiB
Python
577 lines
22 KiB
Python
"""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
|