import os import sqlite3 import time from contextlib import closing SCHEMA = """ CREATE TABLE IF NOT EXISTS samples ( ts REAL NOT NULL, metric TEXT NOT NULL, value REAL ); CREATE INDEX IF NOT EXISTS samples_metric_ts ON samples(metric, ts); CREATE TABLE IF NOT EXISTS events ( id INTEGER PRIMARY KEY AUTOINCREMENT, ts REAL NOT NULL, kind TEXT NOT NULL, message TEXT NOT NULL ); CREATE INDEX IF NOT EXISTS events_ts ON events(ts); CREATE TABLE IF NOT EXISTS kv ( key TEXT PRIMARY KEY, value TEXT ); """ def _like_prefix(prefix): escaped = prefix.replace("\\", "\\\\").replace("%", "\\%").replace("_", "\\_") return escaped + "%" class Storage: def __init__(self, path): directory = os.path.dirname(path) if directory: os.makedirs(directory, exist_ok=True) self.path = path with closing(self._connect()) as conn: conn.executescript(SCHEMA) def _connect(self): conn = sqlite3.connect(self.path, timeout=10) conn.execute("PRAGMA journal_mode=WAL") return conn def add_samples(self, pairs, ts=None): ts = ts or time.time() with closing(self._connect()) as conn, conn: conn.executemany( "INSERT INTO samples (ts, metric, value) VALUES (?, ?, ?)", [(ts, metric, value) for metric, value in pairs], ) def add_sample(self, metric, value, ts=None): self.add_samples([(metric, value)], ts) def add_event(self, kind, message, ts=None): with closing(self._connect()) as conn, conn: conn.execute( "INSERT INTO events (ts, kind, message) VALUES (?, ?, ?)", (ts or time.time(), kind, message), ) def events(self, limit=200): with closing(self._connect()) as conn: rows = conn.execute( "SELECT ts, kind, message FROM events ORDER BY ts DESC, id DESC LIMIT ?", (limit,), ).fetchall() return [{"ts": ts, "kind": kind, "message": message} for ts, kind, message in rows] def samples_by_prefix(self, prefix, since, bucket): bucket = max(1, int(bucket)) with closing(self._connect()) as conn: rows = conn.execute( """ SELECT metric, CAST(ts / ? AS INTEGER) * ? AS bucket, AVG(value) FROM samples WHERE metric LIKE ? ESCAPE '\\' AND ts >= ? GROUP BY metric, bucket ORDER BY metric, bucket """, (bucket, bucket, _like_prefix(prefix), since), ).fetchall() series = {} for metric, ts, value in rows: series.setdefault(metric, []).append([ts, value]) return series def get(self, key): with closing(self._connect()) as conn: row = conn.execute("SELECT value FROM kv WHERE key = ?", (key,)).fetchone() return row[0] if row else None def set(self, key, value): with closing(self._connect()) as conn, conn: conn.execute( "INSERT INTO kv (key, value) VALUES (?, ?) " "ON CONFLICT(key) DO UPDATE SET value = excluded.value", (key, value), ) def prune_samples(self, days): cutoff = time.time() - days * 86400 with closing(self._connect()) as conn, conn: conn.execute("DELETE FROM samples WHERE ts < ?", (cutoff,))