"""Хранение собранных адресов в 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