coinpush / app /core /database.py
zt p
Persist monitor state across restarts
3be8894
Raw
History Blame Contribute Delete
12.2 kB
import json
import os
import sqlite3
import threading
from datetime import datetime, timezone
def utc_now_iso():
return datetime.now(timezone.utc).isoformat(timespec="seconds")
class MonitorStore:
"""Small SQLite state store for health, alerts and event history."""
def __init__(self, db_path):
self.db_path = db_path
os.makedirs(os.path.dirname(db_path), exist_ok=True)
self._lock = threading.RLock()
self._conn = sqlite3.connect(db_path, check_same_thread=False)
self._conn.row_factory = sqlite3.Row
self._conn.execute("PRAGMA journal_mode=WAL")
self._conn.execute("PRAGMA foreign_keys=ON")
self.init_schema()
def init_schema(self):
with self._lock, self._conn:
self._conn.executescript(
"""
CREATE TABLE IF NOT EXISTS monitors (
id TEXT PRIMARY KEY,
type TEXT NOT NULL,
name TEXT NOT NULL,
enabled INTEGER NOT NULL DEFAULT 1,
paused INTEGER NOT NULL DEFAULT 0,
schedule_seconds INTEGER,
status TEXT NOT NULL DEFAULT 'unknown',
last_started_at TEXT,
last_success_at TEXT,
last_error_at TEXT,
last_error TEXT,
consecutive_failures INTEGER NOT NULL DEFAULT 0,
next_run_at TEXT,
last_value_json TEXT,
metadata_json TEXT,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS alerts (
id INTEGER PRIMARY KEY AUTOINCREMENT,
monitor_id TEXT,
category TEXT,
title TEXT,
direction TEXT,
text TEXT NOT NULL,
delivered INTEGER NOT NULL DEFAULT 0,
error TEXT,
created_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS events (
id INTEGER PRIMARY KEY AUTOINCREMENT,
monitor_id TEXT,
level TEXT NOT NULL,
message TEXT NOT NULL,
created_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS kv (
key TEXT PRIMARY KEY,
value TEXT NOT NULL,
updated_at TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_alerts_created_at ON alerts(created_at DESC);
CREATE INDEX IF NOT EXISTS idx_events_created_at ON events(created_at DESC);
"""
)
self._ensure_column("monitors", "paused", "INTEGER NOT NULL DEFAULT 0")
def _ensure_column(self, table, column, definition):
columns = {
row["name"]
for row in self._conn.execute(f"PRAGMA table_info({table})").fetchall()
}
if column not in columns:
self._conn.execute(f"ALTER TABLE {table} ADD COLUMN {column} {definition}")
def upsert_monitor(
self,
monitor_id,
monitor_type,
name,
enabled=True,
schedule_seconds=None,
metadata=None,
):
now = utc_now_iso()
metadata_json = json.dumps(metadata or {}, ensure_ascii=False)
with self._lock, self._conn:
self._conn.execute(
"""
INSERT INTO monitors (
id, type, name, enabled, schedule_seconds, metadata_json,
created_at, updated_at
)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(id) DO UPDATE SET
type=excluded.type,
name=excluded.name,
enabled=excluded.enabled,
schedule_seconds=excluded.schedule_seconds,
metadata_json=excluded.metadata_json,
updated_at=excluded.updated_at
""",
(
monitor_id,
monitor_type,
name,
1 if enabled else 0,
schedule_seconds,
metadata_json,
now,
now,
),
)
def apply_paused_ids(self, monitor_ids):
ids = sorted(set(monitor_ids or []))
if not ids:
return
now = utc_now_iso()
with self._lock, self._conn:
self._conn.executemany(
"""
UPDATE monitors
SET paused=1, status='paused', next_run_at=NULL, updated_at=?
WHERE id=?
""",
[(now, monitor_id) for monitor_id in ids],
)
def get_monitor(self, monitor_id):
with self._lock:
row = self._conn.execute(
"SELECT * FROM monitors WHERE id=?",
(monitor_id,),
).fetchone()
return dict(row) if row else None
def set_monitor_paused(self, monitor_id, paused):
row = self.get_monitor(monitor_id)
if row is None:
raise KeyError(f"未知监控项: {monitor_id}")
now = utc_now_iso()
status = "paused" if paused else row.get("status") or "unknown"
if not paused and status == "paused":
if row.get("last_error_at") and row.get("last_error_at") > (row.get("last_success_at") or ""):
status = "error"
elif row.get("last_success_at"):
status = "ok"
else:
status = "unknown"
with self._lock, self._conn:
self._conn.execute(
"""
UPDATE monitors
SET paused=?, status=?, next_run_at=CASE WHEN ? THEN NULL ELSE next_run_at END,
updated_at=?
WHERE id=?
""",
(1 if paused else 0, status, 1 if paused else 0, now, monitor_id),
)
self.add_event(monitor_id, "info", "监控项已暂停" if paused else "监控项已恢复")
return self.get_monitor(monitor_id)
def is_monitor_paused(self, monitor_id):
row = self.get_monitor(monitor_id)
if row is None:
return False
return bool(row.get("paused"))
def mark_paused(self, monitor_id):
now = utc_now_iso()
with self._lock, self._conn:
self._conn.execute(
"""
UPDATE monitors
SET status='paused', next_run_at=NULL, updated_at=?
WHERE id=?
""",
(now, monitor_id),
)
def mark_started(self, monitor_id):
now = utc_now_iso()
with self._lock, self._conn:
self._conn.execute(
"""
UPDATE monitors
SET status='running', last_started_at=?, updated_at=?
WHERE id=?
""",
(now, now, monitor_id),
)
def mark_success(self, monitor_id, value=None, next_run_at=None, metadata=None):
now = utc_now_iso()
value_json = json.dumps(value, ensure_ascii=False) if value is not None else None
metadata_json = json.dumps(metadata, ensure_ascii=False) if metadata is not None else None
with self._lock, self._conn:
if metadata_json is None:
self._conn.execute(
"""
UPDATE monitors
SET status='ok', last_success_at=?, last_error=NULL,
consecutive_failures=0, next_run_at=?,
last_value_json=COALESCE(?, last_value_json), updated_at=?
WHERE id=?
""",
(now, next_run_at, value_json, now, monitor_id),
)
else:
self._conn.execute(
"""
UPDATE monitors
SET status='ok', last_success_at=?, last_error=NULL,
consecutive_failures=0, next_run_at=?,
last_value_json=COALESCE(?, last_value_json),
metadata_json=?, updated_at=?
WHERE id=?
""",
(now, next_run_at, value_json, metadata_json, now, monitor_id),
)
def mark_failure(self, monitor_id, error):
now = utc_now_iso()
with self._lock, self._conn:
row = self._conn.execute(
"SELECT consecutive_failures FROM monitors WHERE id=?",
(monitor_id,),
).fetchone()
failures = int(row["consecutive_failures"] if row else 0) + 1
self._conn.execute(
"""
UPDATE monitors
SET status='error', last_error_at=?, last_error=?,
consecutive_failures=?, updated_at=?
WHERE id=?
""",
(now, str(error), failures, now, monitor_id),
)
def add_alert(
self,
monitor_id,
text,
category=None,
title=None,
direction=None,
delivered=False,
error=None,
):
with self._lock, self._conn:
self._conn.execute(
"""
INSERT INTO alerts (
monitor_id, category, title, direction, text,
delivered, error, created_at
)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
""",
(
monitor_id,
category,
title,
direction,
text,
1 if delivered else 0,
error,
utc_now_iso(),
),
)
def add_event(self, monitor_id, level, message):
with self._lock, self._conn:
self._conn.execute(
"""
INSERT INTO events (monitor_id, level, message, created_at)
VALUES (?, ?, ?, ?)
""",
(monitor_id, level, message, utc_now_iso()),
)
def set_kv(self, key, value):
now = utc_now_iso()
payload = json.dumps(value, ensure_ascii=False)
with self._lock, self._conn:
self._conn.execute(
"""
INSERT INTO kv (key, value, updated_at)
VALUES (?, ?, ?)
ON CONFLICT(key) DO UPDATE SET
value=excluded.value,
updated_at=excluded.updated_at
""",
(key, payload, now),
)
def get_kv(self, key, default=None):
with self._lock:
row = self._conn.execute("SELECT value FROM kv WHERE key=?", (key,)).fetchone()
if not row:
return default
try:
return json.loads(row["value"])
except json.JSONDecodeError:
return default
def list_monitors(self):
return self._fetch_all(
"""
SELECT * FROM monitors
ORDER BY
CASE status WHEN 'error' THEN 0 WHEN 'running' THEN 1 WHEN 'unknown' THEN 2 WHEN 'paused' THEN 3 ELSE 4 END,
id
"""
)
def recent_alerts(self, limit=50):
return self._fetch_all(
"SELECT * FROM alerts ORDER BY id DESC LIMIT ?",
(int(limit),),
)
def recent_events(self, limit=80):
return self._fetch_all(
"SELECT * FROM events ORDER BY id DESC LIMIT ?",
(int(limit),),
)
def _fetch_all(self, sql, params=()):
with self._lock:
rows = self._conn.execute(sql, params).fetchall()
return [dict(row) for row in rows]
def close(self):
with self._lock:
self._conn.close()