Files
ripe-cidr-collector/db.py
T

487 lines
24 KiB
Python
Raw Normal View History

2026-09-21 07:29:38 +03:00
"""Хранение собранных адресов в SQLite (ripe.db). Конфигурация и heartbeat остаются в JSON."""
import datetime
2026-09-21 08:55:42 +03:00
import glob
2026-09-21 07:29:38 +03:00
import logging
import os
2026-09-21 08:55:42 +03:00
import shutil
2026-09-21 07:29:38 +03:00
import sqlite3
import time
from contextlib import contextmanager
2026-09-21 08:55:42 +03:00
from storage import StorageError, file_lock, load_json, save_json_atomic
2026-09-21 07:29:38 +03:00
logger = logging.getLogger(__name__)
SCHEMA_VERSION = 2
2026-09-21 08:55:42 +03:00
# После отката к копии номера журнала сдвигаются на это число: курсоры, выданные до отката, становятся недействительными
RESTORE_JOURNAL_JUMP = 1_000_000
PRESENCE_BATCH = 500 # значений в одном запросе проверки наличия (лимит переменных SQLite - не менее 999)
2026-09-21 07:29:38 +03:00
_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}"
2026-09-21 08:55:42 +03:00
def _open(path, cc):
"""Открывает базу и приводит схему к актуальной. Ошибки sqlite3 пробрасываются как есть."""
2026-09-21 07:29:38 +03:00
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)
2026-09-21 08:55:42 +03:00
except BaseException:
2026-09-21 07:29:38 +03:00
conn.close()
2026-09-21 08:55:42 +03:00
raise
return conn
def connect(path=None):
"""Открывает базу, создаёт схему и однократно импортирует старые data.json / fqdn_data.json.
Повреждённую базу убирает в карантин и восстанавливает из последней исправной копии (см. _recover).
"""
import cidr_collector as cc # пути читаем при вызове (их подменяют тесты); импорт отложен из-за цикла
path = path or cc.DB_FILE
try:
return _open(path, cc)
except sqlite3.OperationalError as e: # например, database is locked: базу не трогаем
2026-09-21 07:29:38 +03:00
raise StorageError(f"Cannot open {path}: {e}") from e
except sqlite3.DatabaseError as e: # not a database / malformed
2026-09-21 08:55:42 +03:00
return _recover(path, e, cc)
def list_backups(backup_dir):
"""Копии от новых к старым (имя содержит время UTC, поэтому сортировка по имени = по времени)."""
return sorted(glob.glob(os.path.join(backup_dir, "ripe-*.db")), reverse=True)
def _quick_check_file(path):
"""True, если файл открывается как база и проходит PRAGMA quick_check."""
try:
conn = sqlite3.connect(f"file:{path}?mode=ro", uri=True, timeout=5.0)
try:
return conn.execute("PRAGMA quick_check").fetchone()[0] == "ok"
finally:
conn.close()
except sqlite3.Error:
return False
def _recover(path, error, cc):
"""Порча базы при открытии: карантин и восстановление из последней исправной копии.
API и демон могут обнаружить порчу одновременно, поэтому всё выполняется под блокировкой,
а внутри неё база проверяется повторно (другой процесс мог уже восстановить или пересоздать).
Нет исправных копий - создаётся новая пустая база и ставится метка db_recreated.json: пока данные
не собраны заново, API не отдаёт пустой список за настоящий (см. recreated_pending).
2026-09-21 08:55:42 +03:00
"""
with file_lock(path + ".restore"):
try:
return _open(path, cc)
except sqlite3.OperationalError as e:
raise StorageError(f"Cannot open {path}: {e}") from e
except sqlite3.DatabaseError:
pass
quarantine = _quarantine(path)
logger.error("Corrupted database %s (%s); moved to %s", path, error, quarantine)
for candidate in list_backups(cc.BACKUP_DIR):
if not _quick_check_file(candidate):
logger.warning("Backup %s failed the integrity check, skipping", candidate)
continue
if _backup_is_empty(candidate) and _sources_configured(cc):
# Пустая копия при настроенных источниках: метка ставится до подмены файла, как при пересоздании
logger.error("Backup %s has no data although sources are configured; the API withholds "
"addresses until the collector gathers data again", candidate)
_mark_recreated(cc, quarantine, error, "restored backup has no data for the configured sources")
2026-09-21 08:55:42 +03:00
try:
conn = _restore_backup(candidate, path, cc)
except (sqlite3.Error, OSError) as e:
logger.warning("Cannot restore from %s: %s", candidate, e)
continue
logger.error("Database restored from backup %s; data collected after it is lost "
"and will be gathered again by the collector", candidate)
save_json_atomic(cc.RESTORE_FILE, {"at": datetime.datetime.now().isoformat(timespec="seconds"),
"backup": candidate, "quarantine": quarantine, "error": str(error)})
return conn
logger.error("No valid backup for %s: a new empty database is created; the API withholds "
"addresses until the collector gathers data again (remove %s to override)", path, cc.RECREATED_FILE)
# Метка ставится до создания базы: сбой между шагами не оставит пустую базу без метки
_mark_recreated(cc, quarantine, error, "no valid backup")
conn = _open(path, cc)
try:
_reset_journal(conn)
except BaseException:
conn.close()
raise
return conn
def _mark_recreated(cc, quarantine, error, reason):
"""Метка «данные потеряны»: пока сборщик не соберёт их заново, API не отдаёт пустой список (recreated_pending)."""
save_json_atomic(cc.RECREATED_FILE, {"at": datetime.datetime.now().isoformat(timespec="seconds"),
"quarantine": quarantine, "error": str(error), "reason": reason})
def _backup_is_empty(path):
"""Копия не содержит ни одного адреса."""
try:
conn = sqlite3.connect(f"file:{path}?mode=ro", uri=True, timeout=5.0)
try:
return conn.execute("SELECT COUNT(*) FROM addresses").fetchone()[0] == 0
finally:
conn.close()
except sqlite3.Error:
return False
def _sources_configured(cc):
"""В config.json есть хотя бы один ASN или FQDN. Нечитаемый конфиг - True (блокировка по умолчанию)."""
try:
config = cc.load_full_config()
except StorageError:
return True
return bool(config.get("asns") or config.get("fqdns"))
def empty_despite_sources(conn):
"""В базе нет ни одного значения, хотя источники настроены (пустая база не должна выдаваться и копироваться)."""
import cidr_collector as cc
return count_values(conn) == 0 and _sources_configured(cc)
def _reset_journal(conn):
"""Очищает журнал и сдвигает счётчик: курсоры, выданные до потери или отката базы, станут недействительными."""
conn.execute("BEGIN IMMEDIATE")
try:
row = conn.execute("SELECT seq FROM sqlite_sequence WHERE name = 'changes'").fetchone()
horizon = int(conn.execute("SELECT value FROM meta WHERE key = 'horizon_id'").fetchone()[0])
new_seq = max(row[0] if row else 0, horizon) + RESTORE_JOURNAL_JUMP
conn.execute("DELETE FROM changes")
if row:
conn.execute("UPDATE sqlite_sequence SET seq = ? WHERE name = 'changes'", (new_seq,))
else:
conn.execute("INSERT INTO sqlite_sequence (name, seq) VALUES ('changes', ?)", (new_seq,))
conn.execute("UPDATE meta SET value = ? WHERE key = 'horizon_id'", (str(new_seq),))
conn.execute(f"UPDATE meta SET value = {_TS} WHERE key = 'horizon_ts'")
except BaseException:
conn.execute("ROLLBACK")
raise
conn.execute("COMMIT")
2026-09-21 08:55:42 +03:00
def _restore_backup(candidate, path, cc):
"""Кладёт копию на место базы и сбрасывает журнал изменений (его номера после точки копии уже выдавались)."""
tmp = path + ".restore.tmp"
shutil.copyfile(candidate, tmp)
os.replace(tmp, path)
conn = _open(path, cc)
try:
_reset_journal(conn)
2026-09-21 08:55:42 +03:00
except BaseException:
2026-09-21 07:29:38 +03:00
conn.close()
2026-09-21 08:55:42 +03:00
raise
2026-09-21 07:29:38 +03:00
return conn
def recreated_pending(conn, kinds):
"""True, если базу пересоздали после порчи без копий и по запрошенным типам данные ещё не собраны заново.
Тип считается неготовым, когда для него в config.json есть источники, а в базе нет ни одного значения;
тип без источников готов (пустой список законен). Нечитаемый конфиг - блокировка сохраняется.
"""
import cidr_collector as cc
if not os.path.exists(cc.RECREATED_FILE):
return False
try:
config = cc.load_full_config()
except StorageError:
return True
configured = {"asn": config.get("asns"), "fqdn": config.get("fqdns")}
return any(configured[kind] and count_values(conn, kind) == 0 for kind in kinds)
def settle_recreated(conn):
"""Снимает метку пересоздания, когда данные по всем типам собраны заново. Вызывается после сбора."""
import cidr_collector as cc
if os.path.exists(cc.RECREATED_FILE) and not recreated_pending(conn, ("asn", "fqdn")):
os.remove(cc.RECREATED_FILE)
logger.info("Data gathered again after the database loss; %s removed", cc.RECREATED_FILE)
2026-09-21 08:55:42 +03:00
def backup_database(conn, now, keep, backup_dir=None):
"""Онлайн-копия базы с проверкой и ротацией. Возвращает путь копии или None, если копировать нечего.
2026-09-21 08:55:42 +03:00
Живая база и копия проверяются PRAGMA quick_check; при неудаче копия не появляется, старые не трогаются
(исключение StorageError). Пустая база при настроенных источниках не копируется (None): пустая копия
вытеснила бы хорошие при ротации, а её восстановление дало бы пустую выдачу. Лишние старые копии
(сверх keep) удаляются только после успешной новой.
2026-09-21 08:55:42 +03:00
"""
import cidr_collector as cc
backup_dir = backup_dir or cc.BACKUP_DIR
if conn.execute("PRAGMA quick_check").fetchone()[0] != "ok":
raise StorageError("Live database failed the integrity check, backup skipped")
if empty_despite_sources(conn):
return None
2026-09-21 08:55:42 +03:00
os.makedirs(backup_dir, exist_ok=True)
for stale in glob.glob(os.path.join(backup_dir, "ripe-*.db.tmp")): # хвосты прерванной копии
os.unlink(stale)
stamp = now.astimezone(datetime.timezone.utc).strftime("%Y%m%dT%H%M%SZ")
path = os.path.join(backup_dir, f"ripe-{stamp}.db")
tmp = path + ".tmp"
dst = sqlite3.connect(tmp)
try:
conn.backup(dst)
dst.execute("PRAGMA journal_mode=DELETE") # копия - один самодостаточный файл без -wal
finally:
dst.close()
if not _quick_check_file(tmp):
os.unlink(tmp)
raise StorageError("Backup failed the integrity check and was discarded")
os.replace(tmp, path)
for old in list_backups(backup_dir)[keep:]:
os.unlink(old)
return path
2026-09-21 07:29:38 +03:00
@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",
)
@contextmanager
def read_transaction(conn):
"""Один снимок базы (WAL) для нескольких запросов чтения: данные внутри блока согласованы между собой."""
conn.execute("BEGIN")
try:
yield
finally:
conn.execute("COMMIT")
2026-09-21 07:29:38 +03:00
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).
"""
with read_transaction(conn): # один снимок для всех запросов
2026-09-21 07:29:38 +03:00
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)
# Наличие проверяется пакетами: один запрос на PRESENCE_BATCH значений вместо запроса на каждое
candidates = list(first)
present = set()
for i in range(0, len(candidates), PRESENCE_BATCH):
values = list({value for _, value in candidates[i:i + PRESENCE_BATCH]})
marks = ",".join("?" * len(values))
present.update(conn.execute(
f"SELECT DISTINCT kind, value FROM addresses WHERE value IN ({marks})", values))
2026-09-21 07:29:38 +03:00
added, removed = set(), set()
for (kind, value), action in first.items():
if action == "add" and (kind, value) in present:
2026-09-21 07:29:38 +03:00
added.add(value)
elif action == "del" and (kind, value) not in present:
2026-09-21 07:29:38 +03:00
removed.add(value)
return added, removed, head