Files
ayurishchevandClaude Sonnet 5 cb935fef0d Add per-source status tracking, /health sources block and /sources
Collectors record the result of every source (schema v3, table
source_status); /health reports failing sources (failures in a row >=
source_failure_threshold) and turns degraded; GET /sources shows the full
state; runs end with a summary log line instead of "data saved".

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

545 lines
27 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Хранение собранных адресов в SQLite (ripe.db). Конфигурация и heartbeat остаются в JSON."""
import datetime
import glob
import logging
import os
import shutil
import sqlite3
import time
from contextlib import contextmanager
from storage import StorageError, file_lock, load_json, save_json_atomic
logger = logging.getLogger(__name__)
SCHEMA_VERSION = 3
# После отката к копии номера журнала сдвигаются на это число: курсоры, выданные до отката, становятся недействительными
RESTORE_JOURNAL_JUMP = 1_000_000
PRESENCE_BATCH = 500 # значений в одном запросе проверки наличия (лимит переменных SQLite - не менее 999)
_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 _open(path, cc):
"""Открывает базу и приводит схему к актуальной. Ошибки sqlite3 пробрасываются как есть."""
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 BaseException:
conn.close()
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: базу не трогаем
raise StorageError(f"Cannot open {path}: {e}") from e
except sqlite3.DatabaseError as e: # not a database / malformed
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).
"""
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")
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")
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)
except BaseException:
conn.close()
raise
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)
def backup_database(conn, now, keep, backup_dir=None):
"""Онлайн-копия базы с проверкой и ротацией. Возвращает путь копии или None, если копировать нечего.
Живая база и копия проверяются PRAGMA quick_check; при неудаче копия не появляется, старые не трогаются
(исключение StorageError). Пустая база при настроенных источниках не копируется (None): пустая копия
вытеснила бы хорошие при ротации, а её восстановление дало бы пустую выдачу. Лишние старые копии
(сверх keep) удаляются только после успешной новой.
"""
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
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
@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")
# Результат последнего опроса каждого источника (см. record_source_result)
_SOURCE_STATUS_SCHEMA = (
"CREATE TABLE source_status ("
" kind TEXT NOT NULL CHECK (kind IN ('asn', 'fqdn')), source TEXT NOT NULL,"
" last_attempt TEXT NOT NULL, last_success TEXT,"
" failures INTEGER NOT NULL DEFAULT 0, error_kind TEXT,"
" PRIMARY KEY (kind, source))",
)
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}")
if version < 3:
for statement in _SOURCE_STATUS_SCHEMA:
conn.execute(statement)
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, если были адреса."""
conn.execute("DELETE FROM source_status WHERE kind = ? AND source = ?", (kind, source))
return conn.execute("DELETE FROM addresses WHERE kind = ? AND source = ?", (kind, source)).rowcount > 0
def record_source_result(conn, kind, source, now, error_kind=None):
"""Результат опроса источника: успех (error_kind=None) сбрасывает счётчик сбоев, сбой увеличивает его.
Вызывается внутри transaction(). error_kind - короткая категория сбоя (полный текст остаётся в логе).
"""
stamp = _iso(now)
if error_kind is None:
conn.execute(
"INSERT INTO source_status (kind, source, last_attempt, last_success, failures, error_kind)"
" VALUES (?, ?, ?, ?, 0, NULL) ON CONFLICT (kind, source) DO UPDATE SET"
" last_attempt = excluded.last_attempt, last_success = excluded.last_success,"
" failures = 0, error_kind = NULL", (kind, source, stamp, stamp))
else:
conn.execute(
"INSERT INTO source_status (kind, source, last_attempt, last_success, failures, error_kind)"
" VALUES (?, ?, ?, NULL, 1, ?) ON CONFLICT (kind, source) DO UPDATE SET"
" last_attempt = excluded.last_attempt, failures = failures + 1,"
" error_kind = excluded.error_kind", (kind, source, stamp, error_kind))
def forget_unconfigured(conn, kind, configured):
"""Удаляет состояние источников, которых уже нет в конфигурации. Вызывается внутри transaction()."""
query, params = "DELETE FROM source_status WHERE kind = ?", [kind]
if configured:
query += f" AND source NOT IN ({','.join('?' * len(configured))})"
params += sorted(configured)
conn.execute(query, params)
def source_statuses(conn):
"""Состояние источников: {(kind, source): {last_attempt, last_success, failures, error_kind}}."""
return {(kind, source): {"last_attempt": attempt, "last_success": success, "failures": failures,
"error_kind": error_kind}
for kind, source, attempt, success, failures, error_kind in conn.execute(
"SELECT kind, source, last_attempt, last_success, failures, error_kind FROM source_status")}
def address_counts(conn):
"""Число значений по источникам: {(kind, source): n}."""
return {(kind, source): n for kind, source, n in conn.execute(
"SELECT kind, source, COUNT(*) FROM addresses GROUP BY kind, source")}
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): # один снимок для всех запросов
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))
added, removed = set(), set()
for (kind, value), action in first.items():
if action == "add" and (kind, value) in present:
added.add(value)
elif action == "del" and (kind, value) not in present:
removed.add(value)
return added, removed, head