Add system alerts log with nav badge and Admin Alerts tab.
Persist job and ingestion warnings/errors for 7 days, surface a dismissible top-nav badge, treat stale OHLCV as a warning (e.g. ticker renames), and show market bar age on the ticker freshness chip.
This commit is contained in:
@@ -13,6 +13,7 @@ from app.models.paper_trade import PaperTrade
|
||||
from app.models.regime_snapshot import RegimeSnapshot
|
||||
from app.models.benchmark_price import BenchmarkPrice
|
||||
from app.models.signal_context_snapshot import SignalContextSnapshot
|
||||
from app.models.system_event import SystemEvent
|
||||
|
||||
__all__ = [
|
||||
"Ticker",
|
||||
@@ -32,4 +33,5 @@ __all__ = [
|
||||
"RegimeSnapshot",
|
||||
"BenchmarkPrice",
|
||||
"SignalContextSnapshot",
|
||||
"SystemEvent",
|
||||
]
|
||||
|
||||
@@ -0,0 +1,37 @@
|
||||
"""Operational system events (warnings/errors) for the admin UI and top-nav badge."""
|
||||
|
||||
from datetime import datetime
|
||||
|
||||
from sqlalchemy import DateTime, Index, String, Text
|
||||
from sqlalchemy.orm import Mapped, mapped_column
|
||||
|
||||
from app.database import Base
|
||||
|
||||
|
||||
class SystemEvent(Base):
|
||||
"""Durable warning/error record (jobs, ingestion, pipelines).
|
||||
|
||||
``acknowledged_at`` is set when a user dismisses the nav badge / clears
|
||||
events — history still shows in Admin → Jobs for the retention window.
|
||||
"""
|
||||
|
||||
__tablename__ = "system_events"
|
||||
__table_args__ = (
|
||||
Index("ix_system_events_created_at", "created_at"),
|
||||
Index("ix_system_events_ack_created", "acknowledged_at", "created_at"),
|
||||
Index("ix_system_events_dedup_created", "dedup_key", "created_at"),
|
||||
)
|
||||
|
||||
id: Mapped[int] = mapped_column(primary_key=True)
|
||||
severity: Mapped[str] = mapped_column(String(16), nullable=False) # warning | error
|
||||
source: Mapped[str] = mapped_column(String(64), nullable=False)
|
||||
code: Mapped[str] = mapped_column(String(64), nullable=False)
|
||||
message: Mapped[str] = mapped_column(Text, nullable=False)
|
||||
symbol: Mapped[str | None] = mapped_column(String(20), nullable=True)
|
||||
dedup_key: Mapped[str | None] = mapped_column(String(200), nullable=True)
|
||||
created_at: Mapped[datetime] = mapped_column(
|
||||
DateTime(timezone=True), default=datetime.utcnow, nullable=False
|
||||
)
|
||||
acknowledged_at: Mapped[datetime | None] = mapped_column(
|
||||
DateTime(timezone=True), nullable=True
|
||||
)
|
||||
+49
-1
@@ -6,7 +6,7 @@ All endpoints require admin role.
|
||||
from fastapi import APIRouter, Depends, Query
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from app.dependencies import get_db, require_admin
|
||||
from app.dependencies import get_db, require_access, require_admin
|
||||
from app.models.user import User
|
||||
from app.schemas.admin import (
|
||||
ActivationConfigUpdate,
|
||||
@@ -29,6 +29,7 @@ from app.schemas.common import APIEnvelope
|
||||
from app.services import admin_service
|
||||
from app.services import alert_service
|
||||
from app.services import sentiment_provider_service
|
||||
from app.services import system_event_service
|
||||
from app.services import ticker_universe_service
|
||||
|
||||
router = APIRouter(tags=["admin"])
|
||||
@@ -403,3 +404,50 @@ async def toggle_job(
|
||||
status="success",
|
||||
data={"key": setting.key, "value": setting.value},
|
||||
)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# System events (operational warnings / errors)
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
@router.get("/admin/system-events", response_model=APIEnvelope)
|
||||
async def list_system_events(
|
||||
days: int = Query(7, ge=1, le=30),
|
||||
severity: str | None = Query(None, description="warning | error"),
|
||||
unacknowledged_only: bool = Query(False),
|
||||
_user: User = Depends(require_access),
|
||||
db: AsyncSession = Depends(get_db),
|
||||
):
|
||||
"""List recent system events (default last 7 days)."""
|
||||
rows = await system_event_service.list_events(
|
||||
db,
|
||||
days=days,
|
||||
severity=severity,
|
||||
unacknowledged_only=unacknowledged_only,
|
||||
)
|
||||
return APIEnvelope(
|
||||
status="success",
|
||||
data=[system_event_service.event_to_dict(r) for r in rows],
|
||||
)
|
||||
|
||||
|
||||
@router.get("/admin/system-events/summary", response_model=APIEnvelope)
|
||||
async def system_events_summary(
|
||||
days: int = Query(7, ge=1, le=30),
|
||||
_user: User = Depends(require_access),
|
||||
db: AsyncSession = Depends(get_db),
|
||||
):
|
||||
"""Badge counts for the top nav."""
|
||||
data = await system_event_service.summary(db, days=days)
|
||||
return APIEnvelope(status="success", data=data)
|
||||
|
||||
|
||||
@router.post("/admin/system-events/acknowledge", response_model=APIEnvelope)
|
||||
async def acknowledge_system_events(
|
||||
days: int = Query(7, ge=1, le=30),
|
||||
_user: User = Depends(require_access),
|
||||
db: AsyncSession = Depends(get_db),
|
||||
):
|
||||
"""Dismiss unacknowledged events in the lookback window (clears the badge)."""
|
||||
count = await system_event_service.acknowledge_all(db, days=days)
|
||||
return APIEnvelope(status="success", data={"acknowledged": count})
|
||||
|
||||
@@ -117,12 +117,32 @@ async def fetch_symbol(
|
||||
result = await ingestion_service.fetch_and_ingest(
|
||||
db, provider, symbol_upper, start_date, end_date
|
||||
)
|
||||
status_map = {"complete": "ok", "partial": "ok", "no_data": "warning"}
|
||||
# "stale" = provider returned nothing but our last bar is old
|
||||
# (rename/delist/halt) — must not look like a successful refresh.
|
||||
status_map = {
|
||||
"complete": "ok",
|
||||
"partial": "ok",
|
||||
"no_data": "warning",
|
||||
"stale": "warning",
|
||||
}
|
||||
sources_out["ohlcv"] = {
|
||||
"status": status_map.get(result.status, "error"),
|
||||
"records": result.records_ingested,
|
||||
"message": result.message,
|
||||
"last_date": result.last_date.isoformat() if result.last_date else None,
|
||||
}
|
||||
if result.status in ("stale", "no_data", "error"):
|
||||
from app.services.system_event_service import log_event
|
||||
|
||||
await log_event(
|
||||
db,
|
||||
severity="warning" if result.status != "error" else "error",
|
||||
source="ingestion",
|
||||
code=f"ohlcv_{result.status}",
|
||||
message=result.message or f"OHLCV fetch {result.status} for {symbol_upper}",
|
||||
symbol=symbol_upper,
|
||||
dedup_key=f"ohlcv_{result.status}:{symbol_upper}",
|
||||
)
|
||||
except Exception as exc:
|
||||
logger.error("OHLCV fetch failed for %s: %s", symbol_upper, exc)
|
||||
sources_out["ohlcv"] = {"status": "error", "records": 0, "message": str(exc)}
|
||||
|
||||
@@ -152,6 +152,28 @@ def _log_job_error(job_name: str, ticker: str, error: Exception) -> None:
|
||||
)
|
||||
|
||||
|
||||
async def _record_system_event(
|
||||
*,
|
||||
severity: str,
|
||||
source: str,
|
||||
code: str,
|
||||
message: str,
|
||||
symbol: str | None = None,
|
||||
dedup_key: str | None = None,
|
||||
) -> None:
|
||||
"""Best-effort durable event for Admin → Jobs and the top-nav badge."""
|
||||
from app.services.system_event_service import log_event_standalone
|
||||
|
||||
await log_event_standalone(
|
||||
severity=severity,
|
||||
source=source,
|
||||
code=code,
|
||||
message=message,
|
||||
symbol=symbol,
|
||||
dedup_key=dedup_key,
|
||||
)
|
||||
|
||||
|
||||
def _runtime_start(job_name: str, total: int | None = None, message: str | None = None) -> None:
|
||||
_job_runtime[job_name] = {
|
||||
**_idle_runtime(),
|
||||
@@ -206,6 +228,22 @@ def _runtime_finish(
|
||||
"message": message,
|
||||
})
|
||||
_job_runtime[job_name] = runtime
|
||||
# Durable event for error / rate-limit finishes (badge + Admin → Jobs panel).
|
||||
if status in ("error", "rate_limited"):
|
||||
severity = "error" if status == "error" else "warning"
|
||||
try:
|
||||
loop = asyncio.get_running_loop()
|
||||
loop.create_task(
|
||||
_record_system_event(
|
||||
severity=severity,
|
||||
source=job_name,
|
||||
code=f"job_{status}",
|
||||
message=message or f"Job {job_name} finished with status {status}",
|
||||
dedup_key=f"job:{job_name}:{status}:{(message or '')[:80]}"[:200],
|
||||
)
|
||||
)
|
||||
except RuntimeError:
|
||||
pass
|
||||
|
||||
|
||||
def get_job_runtime_snapshot(job_name: str | None = None) -> dict[str, dict[str, object]] | dict[str, object]:
|
||||
@@ -489,6 +527,15 @@ async def collect_ohlcv(full_backfill: bool = False, job_name: str = "data_colle
|
||||
processed += 1
|
||||
_runtime_progress(job_name, processed=processed, total=total, current_ticker=symbol)
|
||||
_log_event(logging.INFO, "ticker_collected", job=job_name, ticker=symbol, status=result.status, records=result.records_ingested)
|
||||
if result.status == "stale":
|
||||
await _record_system_event(
|
||||
severity="warning",
|
||||
source=job_name,
|
||||
code="ohlcv_stale",
|
||||
message=result.message or f"No new OHLCV bars for {symbol}",
|
||||
symbol=symbol,
|
||||
dedup_key=f"ohlcv_stale:{symbol}",
|
||||
)
|
||||
if result.status == "partial":
|
||||
# Rate limited — stop and resume next run
|
||||
_log_event(logging.WARNING, "rate_limited", job=job_name, ticker=symbol, processed=processed)
|
||||
@@ -496,6 +543,14 @@ async def collect_ohlcv(full_backfill: bool = False, job_name: str = "data_colle
|
||||
return
|
||||
except Exception as exc:
|
||||
_log_job_error(job_name, symbol, exc)
|
||||
await _record_system_event(
|
||||
severity="error",
|
||||
source=job_name,
|
||||
code="job_ticker_error",
|
||||
message=f"{type(exc).__name__}: {exc}",
|
||||
symbol=symbol,
|
||||
dedup_key=f"job_ticker_error:{job_name}:{symbol}:{type(exc).__name__}",
|
||||
)
|
||||
|
||||
# Reset resume pointer on full completion
|
||||
_last_successful[job_name] = None
|
||||
|
||||
@@ -59,6 +59,19 @@ async def _get_ohlcv_bar_count(db: AsyncSession, ticker_id: int) -> int:
|
||||
return int(result.scalar() or 0)
|
||||
|
||||
|
||||
async def _get_latest_ohlcv_date(db: AsyncSession, ticker_id: int) -> date | None:
|
||||
result = await db.execute(
|
||||
select(func.max(OHLCVRecord.date)).where(OHLCVRecord.ticker_id == ticker_id)
|
||||
)
|
||||
return result.scalar_one_or_none()
|
||||
|
||||
|
||||
# If the provider returns no bars but our last stored session is older than this,
|
||||
# treat the run as stale (not "success / up to date"). Common causes: ticker
|
||||
# rename, delisting, or multi-day halt — SATS→ECHO is the canonical example.
|
||||
_STALE_OHLCV_GAP_DAYS = 5
|
||||
|
||||
|
||||
async def _update_progress(
|
||||
db: AsyncSession, ticker_id: int, last_date: date
|
||||
) -> None:
|
||||
@@ -145,8 +158,9 @@ async def fetch_and_ingest(
|
||||
|
||||
# Provider returned nothing. With no history at all this almost always means
|
||||
# the provider doesn't cover this symbol (Alpaca = US listings only) — surface
|
||||
# that instead of a misleading "success". With existing bars it just means
|
||||
# there were no new bars in the requested window.
|
||||
# that instead of a misleading "success". With recent history, an empty window
|
||||
# usually means weekends/holidays. With a multi-day gap, the symbol is likely
|
||||
# halted, delisted, or *renamed* (e.g. SATS → ECHO) and we must not claim success.
|
||||
if not records:
|
||||
existing = await _get_ohlcv_bar_count(db, ticker.id)
|
||||
if existing == 0:
|
||||
@@ -160,10 +174,24 @@ async def fetch_and_ingest(
|
||||
"(Alpaca serves US-listed securities only)."
|
||||
),
|
||||
)
|
||||
latest = await _get_latest_ohlcv_date(db, ticker.id)
|
||||
gap_days = (end_date - latest).days if latest is not None else None
|
||||
if gap_days is not None and gap_days > _STALE_OHLCV_GAP_DAYS:
|
||||
return IngestionResult(
|
||||
symbol=ticker.symbol,
|
||||
records_ingested=0,
|
||||
last_date=latest,
|
||||
status="stale",
|
||||
message=(
|
||||
f"No new bars since {latest.isoformat()} ({gap_days}d gap). "
|
||||
"The symbol may be halted, delisted, or renamed under a new ticker — "
|
||||
"check the listing and add/fetch the current symbol if it changed."
|
||||
),
|
||||
)
|
||||
return IngestionResult(
|
||||
symbol=ticker.symbol,
|
||||
records_ingested=0,
|
||||
last_date=None,
|
||||
last_date=latest,
|
||||
status="complete",
|
||||
message="Already up to date — no new bars.",
|
||||
)
|
||||
|
||||
@@ -0,0 +1,191 @@
|
||||
"""Persist and query operational system events (warnings / errors)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from datetime import datetime, timedelta, timezone
|
||||
|
||||
from sqlalchemy import select, update
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from app.models.system_event import SystemEvent
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
DEFAULT_LOOKBACK_DAYS = 7
|
||||
DEFAULT_DEDUP_HOURS = 24
|
||||
SEVERITIES = frozenset({"warning", "error"})
|
||||
|
||||
|
||||
async def log_event(
|
||||
db: AsyncSession,
|
||||
*,
|
||||
severity: str,
|
||||
source: str,
|
||||
code: str,
|
||||
message: str,
|
||||
symbol: str | None = None,
|
||||
dedup_key: str | None = None,
|
||||
dedup_hours: int = DEFAULT_DEDUP_HOURS,
|
||||
) -> SystemEvent | None:
|
||||
"""Insert a system event, optionally de-duplicating recent identical keys.
|
||||
|
||||
Returns the new row, or None when a recent dedup_key already exists.
|
||||
Commits the session.
|
||||
"""
|
||||
severity = (severity or "").strip().lower()
|
||||
if severity not in SEVERITIES:
|
||||
severity = "warning"
|
||||
source = (source or "system")[:64]
|
||||
code = (code or "unknown")[:64]
|
||||
message = (message or "").strip() or code
|
||||
symbol = symbol.strip().upper()[:20] if symbol else None
|
||||
dedup_key = dedup_key[:200] if dedup_key else None
|
||||
|
||||
if dedup_key:
|
||||
cutoff = datetime.now(timezone.utc) - timedelta(hours=max(1, dedup_hours))
|
||||
existing = await db.execute(
|
||||
select(SystemEvent.id)
|
||||
.where(
|
||||
SystemEvent.dedup_key == dedup_key,
|
||||
SystemEvent.created_at >= cutoff,
|
||||
)
|
||||
.limit(1)
|
||||
)
|
||||
if existing.scalar_one_or_none() is not None:
|
||||
return None
|
||||
|
||||
row = SystemEvent(
|
||||
severity=severity,
|
||||
source=source,
|
||||
code=code,
|
||||
message=message[:4000],
|
||||
symbol=symbol,
|
||||
dedup_key=dedup_key,
|
||||
created_at=datetime.now(timezone.utc),
|
||||
)
|
||||
db.add(row)
|
||||
await db.commit()
|
||||
await db.refresh(row)
|
||||
return row
|
||||
|
||||
|
||||
async def log_event_standalone(
|
||||
*,
|
||||
severity: str,
|
||||
source: str,
|
||||
code: str,
|
||||
message: str,
|
||||
symbol: str | None = None,
|
||||
dedup_key: str | None = None,
|
||||
) -> None:
|
||||
"""Open a short-lived session and log an event (for scheduler / fire-and-forget)."""
|
||||
try:
|
||||
from app.database import async_session_factory
|
||||
|
||||
async with async_session_factory() as db:
|
||||
await log_event(
|
||||
db,
|
||||
severity=severity,
|
||||
source=source,
|
||||
code=code,
|
||||
message=message,
|
||||
symbol=symbol,
|
||||
dedup_key=dedup_key,
|
||||
)
|
||||
except Exception:
|
||||
logger.exception("Failed to persist system event %s/%s", source, code)
|
||||
|
||||
|
||||
async def list_events(
|
||||
db: AsyncSession,
|
||||
*,
|
||||
days: int = DEFAULT_LOOKBACK_DAYS,
|
||||
severity: str | None = None,
|
||||
unacknowledged_only: bool = False,
|
||||
limit: int = 200,
|
||||
) -> list[SystemEvent]:
|
||||
days = max(1, min(int(days), 30))
|
||||
cutoff = datetime.now(timezone.utc) - timedelta(days=days)
|
||||
stmt = (
|
||||
select(SystemEvent)
|
||||
.where(SystemEvent.created_at >= cutoff)
|
||||
.order_by(SystemEvent.created_at.desc())
|
||||
.limit(max(1, min(limit, 500)))
|
||||
)
|
||||
if severity in SEVERITIES:
|
||||
stmt = stmt.where(SystemEvent.severity == severity)
|
||||
if unacknowledged_only:
|
||||
stmt = stmt.where(SystemEvent.acknowledged_at.is_(None))
|
||||
result = await db.execute(stmt)
|
||||
return list(result.scalars().all())
|
||||
|
||||
|
||||
async def summary(
|
||||
db: AsyncSession,
|
||||
*,
|
||||
days: int = DEFAULT_LOOKBACK_DAYS,
|
||||
) -> dict:
|
||||
"""Counts for badge + admin header."""
|
||||
days = max(1, min(int(days), 30))
|
||||
cutoff = datetime.now(timezone.utc) - timedelta(days=days)
|
||||
rows = (
|
||||
await db.execute(
|
||||
select(SystemEvent.severity, SystemEvent.acknowledged_at).where(
|
||||
SystemEvent.created_at >= cutoff
|
||||
)
|
||||
)
|
||||
).all()
|
||||
total = len(rows)
|
||||
unacked = 0
|
||||
errors = 0
|
||||
warnings = 0
|
||||
for severity, acknowledged_at in rows:
|
||||
if acknowledged_at is not None:
|
||||
continue
|
||||
unacked += 1
|
||||
if severity == "error":
|
||||
errors += 1
|
||||
elif severity == "warning":
|
||||
warnings += 1
|
||||
return {
|
||||
"days": days,
|
||||
"total": total,
|
||||
"unacknowledged": unacked,
|
||||
"unacknowledged_errors": errors,
|
||||
"unacknowledged_warnings": warnings,
|
||||
}
|
||||
|
||||
|
||||
async def acknowledge_all(
|
||||
db: AsyncSession,
|
||||
*,
|
||||
days: int = DEFAULT_LOOKBACK_DAYS,
|
||||
) -> int:
|
||||
"""Mark unacknowledged events in the lookback window as dismissed. Returns count."""
|
||||
days = max(1, min(int(days), 30))
|
||||
cutoff = datetime.now(timezone.utc) - timedelta(days=days)
|
||||
now = datetime.now(timezone.utc)
|
||||
result = await db.execute(
|
||||
update(SystemEvent)
|
||||
.where(
|
||||
SystemEvent.created_at >= cutoff,
|
||||
SystemEvent.acknowledged_at.is_(None),
|
||||
)
|
||||
.values(acknowledged_at=now)
|
||||
)
|
||||
await db.commit()
|
||||
return int(result.rowcount or 0)
|
||||
|
||||
|
||||
def event_to_dict(row: SystemEvent) -> dict:
|
||||
return {
|
||||
"id": row.id,
|
||||
"severity": row.severity,
|
||||
"source": row.source,
|
||||
"code": row.code,
|
||||
"message": row.message,
|
||||
"symbol": row.symbol,
|
||||
"created_at": row.created_at.isoformat() if row.created_at else None,
|
||||
"acknowledged_at": row.acknowledged_at.isoformat() if row.acknowledged_at else None,
|
||||
}
|
||||
Reference in New Issue
Block a user