Files
ripe-cidr-collector/db.py
T
ayurishchevandClaude Sonnet 5 bcf8156085 Initial commit: RIPE CIDR/FQDN collector
Collector daemon, FastAPI server (addresses, diff, collect, sources),
SQLite storage with change journal, Docker Compose deployment,
tests, documentation and project rules.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-09-21 07:29:38 +03:00

261 lines
13 KiB
Python

"""Хранение собранных адресов в SQLite (ripe.db). Конфигурация и heartbeat остаются в JSON."""
import datetime
import logging
import os
import sqlite3
import time
from contextlib import contextmanager
from storage import StorageError, load_json
logger = logging.getLogger(__name__)
SCHEMA_VERSION = 2
_UPSERT = ("INSERT INTO addresses (kind, source, value, first_seen, last_seen) VALUES (?, ?, ?, ?, ?) "
"ON CONFLICT (kind, source, value) DO UPDATE SET last_seen = excluded.last_seen")
def _iso(moment):
return moment.isoformat(timespec="seconds")
def _quarantine(path):
"""Повреждённую базу (вместе с -wal/-shm) убираем в сторону, чтобы её не затёрли."""
stamp = int(time.time())
for suffix in ("", "-wal", "-shm"):
if os.path.exists(path + suffix):
os.replace(path + suffix, f"{path}.corrupt-{stamp}{suffix}")
return f"{path}.corrupt-{stamp}"
def connect(path=None):
"""Открывает базу, создаёт схему и однократно импортирует старые data.json / fqdn_data.json."""
import cidr_collector as cc # пути читаем при вызове (их подменяют тесты); импорт отложен из-за цикла
path = path or cc.DB_FILE
conn = sqlite3.connect(path, timeout=5.0, isolation_level=None) # транзакциями управляем явно
try:
conn.execute("PRAGMA journal_mode=WAL")
conn.execute("PRAGMA synchronous=NORMAL")
conn.execute("PRAGMA temp_store=MEMORY")
_ensure_schema(conn, cc)
except sqlite3.OperationalError as e: # например, database is locked: базу не трогаем
conn.close()
raise StorageError(f"Cannot open {path}: {e}") from e
except sqlite3.DatabaseError as e: # not a database / malformed
conn.close()
backup = _quarantine(path)
logger.error("Corrupted database %s (%s); moved to %s", path, e, backup)
raise StorageError(f"{path} is corrupted") from e
return conn
@contextmanager
def session():
"""Соединение на время блока (закрывается по выходу)."""
conn = connect()
try:
yield conn
finally:
conn.close()
@contextmanager
def transaction(conn):
"""Одна транзакция на запись (BEGIN IMMEDIATE: писатель один, читатели не блокируются)."""
conn.execute("BEGIN IMMEDIATE")
try:
yield
except BaseException:
conn.execute("ROLLBACK")
raise
conn.execute("COMMIT")
_TS = "strftime('%Y-%m-%dT%H:%M:%fZ', 'now')" # UTC с миллисекундами: сравнивается как строка
# Журнал изменений итогового набора значений (см. get_changes). Триггеры покрывают все пути записи и удаления.
_JOURNAL_SCHEMA = (
"CREATE TABLE changes ("
" id INTEGER PRIMARY KEY AUTOINCREMENT, ts TEXT NOT NULL, kind TEXT NOT NULL, value TEXT NOT NULL,"
" action TEXT NOT NULL CHECK (action IN ('add', 'del')))",
"CREATE INDEX changes_ts ON changes (ts)",
"CREATE TABLE meta (key TEXT PRIMARY KEY, value TEXT NOT NULL)",
# Значение появилось, если у других источников его не было
"CREATE TRIGGER addresses_added AFTER INSERT ON addresses"
" WHEN NOT EXISTS (SELECT 1 FROM addresses WHERE kind = NEW.kind AND value = NEW.value AND source <> NEW.source)"
f" BEGIN INSERT INTO changes (ts, kind, value, action) VALUES ({_TS}, NEW.kind, NEW.value, 'add'); END",
# Значение исчезло, если удалена его последняя запись
"CREATE TRIGGER addresses_removed AFTER DELETE ON addresses"
" WHEN NOT EXISTS (SELECT 1 FROM addresses WHERE kind = OLD.kind AND value = OLD.value)"
f" BEGIN INSERT INTO changes (ts, kind, value, action) VALUES ({_TS}, OLD.kind, OLD.value, 'del'); END",
)
def _ensure_schema(conn, cc):
if conn.execute("PRAGMA user_version").fetchone()[0] >= SCHEMA_VERSION:
return
imported = []
conn.execute("BEGIN IMMEDIATE") # API и демон могут стартовать одновременно: второй дождётся первого
try:
version = conn.execute("PRAGMA user_version").fetchone()[0]
if version == 0:
conn.execute(
"CREATE TABLE IF NOT EXISTS addresses ("
" kind TEXT NOT NULL CHECK (kind IN ('asn', 'fqdn')),"
" source TEXT NOT NULL, value TEXT NOT NULL,"
" first_seen TEXT NOT NULL, last_seen TEXT NOT NULL,"
" PRIMARY KEY (kind, source, value))")
conn.execute("CREATE INDEX IF NOT EXISTS addresses_value ON addresses (value)")
imported = _import_legacy_json(conn, cc) # журнал ещё не создан: импорт в него не попадает
if version < 2:
for statement in _JOURNAL_SCHEMA:
conn.execute(statement)
conn.execute("INSERT INTO meta VALUES ('horizon_id', '0')")
conn.execute(f"INSERT INTO meta SELECT 'horizon_ts', {_TS}")
conn.execute(f"PRAGMA user_version = {SCHEMA_VERSION}")
conn.execute("COMMIT")
except BaseException:
conn.execute("ROLLBACK")
raise
# Старые JSON оставляем как резервную копию (и путь отката)
stamp = int(time.time())
for legacy in imported:
os.replace(legacy, f"{legacy}.migrated-{stamp}")
logger.info("Imported %s into the database; original kept as %s.migrated-%d", legacy, legacy, stamp)
def _import_legacy_json(conn, cc):
"""Переносит data.json (asn) и fqdn_data.json (fqdn). Возвращает пути реально импортированных файлов."""
now = _iso(datetime.datetime.now())
imported = []
for kind, path, items_key in (("asn", cc.DATA_FILE, "prefixes"), ("fqdn", cc.FQDN_DATA_FILE, "ips")):
if not os.path.exists(path):
continue
rows = []
for source, entry in load_json(path, {}).items():
seen = entry.get("seen")
if seen:
rows += [(kind, source, value, info["first_seen"], info["last_seen"]) for value, info in seen.items()]
else:
# Формат до TTL: last_seen = сейчас, чтобы после перехода ничего не истекло сразу
first = entry.get("last_updated", now)
rows += [(kind, source, value, first, now) for value in entry.get(items_key, [])]
conn.executemany("INSERT OR IGNORE INTO addresses VALUES (?, ?, ?, ?, ?)", rows)
imported.append(path)
return imported
def merge_source(conn, kind, source, values, now, ttl_days):
"""Объединяет найденные значения источника с сохранёнными и применяет TTL.
Вызывается внутри transaction(). Возвращает (добавленные, удалённые по TTL).
"""
before = {v for (v,) in conn.execute("SELECT value FROM addresses WHERE kind = ? AND source = ?", (kind, source))}
stamp = _iso(now)
conn.executemany(_UPSERT, [(kind, source, value, stamp, stamp) for value in values])
removed = set()
if ttl_days > 0:
cutoff = _iso(now - datetime.timedelta(days=ttl_days))
where = "kind = ? AND source = ? AND last_seen < ?"
removed = {v for (v,) in conn.execute(f"SELECT value FROM addresses WHERE {where}", (kind, source, cutoff))}
conn.execute(f"DELETE FROM addresses WHERE {where}", (kind, source, cutoff))
return set(values) - before, removed
def sweep_unconfigured(conn, kind, configured, now, ttl_days):
"""Источники, которых уже нет в конфигурации, не опрашиваются: их адреса истекают по TTL.
Вызывается внутри transaction(). Возвращает число удалённых адресов.
"""
if ttl_days <= 0:
return 0
cutoff = _iso(now - datetime.timedelta(days=ttl_days))
where, params = "kind = ? AND last_seen < ?", [kind, cutoff]
if configured:
where += f" AND source NOT IN ({','.join('?' * len(configured))})"
params += sorted(configured)
return conn.execute(f"DELETE FROM addresses WHERE {where}", params).rowcount
def purge_source(conn, kind, source):
"""Удаляет все адреса источника. True, если что-то было."""
return conn.execute("DELETE FROM addresses WHERE kind = ? AND source = ?", (kind, source)).rowcount > 0
def get_values(conn, kind=None):
query, params = "SELECT DISTINCT value FROM addresses", ()
if kind:
query, params = query + " WHERE kind = ?", (kind,)
return [v for (v,) in conn.execute(query, params)]
def count_values(conn, kind=None):
query, params = "SELECT COUNT(DISTINCT value) FROM addresses", ()
if kind:
query, params = query + " WHERE kind = ?", (kind,)
return conn.execute(query, params).fetchone()[0]
def prune_changes(conn, now, retention_days):
"""Удаляет записи журнала старше срока хранения и сдвигает «горизонт» (граница, глубже которой diff недоступен).
Вызывается внутри transaction(). retention_days <= 0 - без очистки. Возвращает число удалённых записей.
"""
if retention_days <= 0:
return 0
cutoff = (now.astimezone(datetime.timezone.utc) - datetime.timedelta(days=retention_days)).strftime(
"%Y-%m-%dT%H:%M:%S.000Z")
last_id = conn.execute("SELECT MAX(id) FROM changes WHERE ts < ?", (cutoff,)).fetchone()[0]
if last_id is None:
return 0
removed = conn.execute("DELETE FROM changes WHERE id <= ?", (last_id,)).rowcount
conn.execute("UPDATE meta SET value = ? WHERE key = 'horizon_id'", (str(last_id),))
conn.execute("UPDATE meta SET value = ? WHERE key = 'horizon_ts'", (cutoff,))
return removed
def journal_head(conn):
"""Текущий курсор журнала (номер последней записи)."""
return conn.execute("SELECT COALESCE(MAX(id), (SELECT CAST(value AS INTEGER) FROM meta WHERE key = 'horizon_id'))"
" FROM changes").fetchone()[0]
def get_changes(conn, kinds, cursor=None, since_ts=None):
"""Изменения итогового набора значений после курсора (номер записи журнала) или времени (UTC, строка).
Берётся первое действие по значению после точки отсчёта и сверяется с текущим состоянием: значение,
добавленное и удалённое (или наоборот) в этом интервале, в результат не попадает.
Возвращает (added, removed, cursor) или None, если журнал не покрывает точку отсчёта (ответ 410).
"""
conn.execute("BEGIN") # один снимок для всех запросов
try:
meta = dict(conn.execute("SELECT key, value FROM meta"))
horizon_id = int(meta["horizon_id"])
head = conn.execute("SELECT COALESCE(MAX(id), ?) FROM changes", (horizon_id,)).fetchone()[0]
if since_ts is not None:
if since_ts < meta["horizon_ts"]:
return None
cursor = conn.execute("SELECT COALESCE(MAX(id), ?) FROM changes WHERE ts < ?",
(horizon_id, since_ts)).fetchone()[0]
if not horizon_id <= cursor <= head:
return None
first = {}
for kind, value, action in conn.execute(
"SELECT kind, value, action FROM changes WHERE id > ? ORDER BY id", (cursor,)):
if kind in kinds:
first.setdefault((kind, value), action)
added, removed = set(), set()
for (kind, value), action in first.items():
present = conn.execute("SELECT 1 FROM addresses WHERE kind = ? AND value = ? LIMIT 1",
(kind, value)).fetchone() is not None
if action == "add" and present:
added.add(value)
elif action == "del" and not present:
removed.add(value)
return added, removed, head
finally:
conn.execute("COMMIT")