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__)
|
|
|
|
|
|
|
2026-09-21 10:46:38 +03:00
|
|
|
|
SCHEMA_VERSION = 3
|
2026-09-21 08:55:42 +03:00
|
|
|
|
# После отката к копии номера журнала сдвигаются на это число: курсоры, выданные до отката, становятся недействительными
|
|
|
|
|
|
RESTORE_JOURNAL_JUMP = 1_000_000
|
2026-09-21 10:06:31 +03:00
|
|
|
|
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 и демон могут обнаружить порчу одновременно, поэтому всё выполняется под блокировкой,
|
2026-09-21 09:56:43 +03:00
|
|
|
|
а внутри неё база проверяется повторно (другой процесс мог уже восстановить или пересоздать).
|
|
|
|
|
|
Нет исправных копий - создаётся новая пустая база и ставится метка 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
|
2026-09-21 10:30:03 +03:00
|
|
|
|
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
|
2026-09-21 09:56:43 +03:00
|
|
|
|
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)
|
|
|
|
|
|
# Метка ставится до создания базы: сбой между шагами не оставит пустую базу без метки
|
2026-09-21 10:30:03 +03:00
|
|
|
|
_mark_recreated(cc, quarantine, error, "no valid backup")
|
2026-09-21 09:56:43 +03:00
|
|
|
|
conn = _open(path, cc)
|
|
|
|
|
|
try:
|
|
|
|
|
|
_reset_journal(conn)
|
|
|
|
|
|
except BaseException:
|
|
|
|
|
|
conn.close()
|
|
|
|
|
|
raise
|
|
|
|
|
|
return conn
|
|
|
|
|
|
|
|
|
|
|
|
|
2026-09-21 10:30:03 +03:00
|
|
|
|
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)
|
|
|
|
|
|
|
|
|
|
|
|
|
2026-09-21 09:56:43 +03:00
|
|
|
|
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:
|
2026-09-21 09:56:43 +03:00
|
|
|
|
_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
|
|
|
|
|
|
|
|
|
|
|
|
|
2026-09-21 09:56:43 +03:00
|
|
|
|
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):
|
2026-09-21 10:30:03 +03:00
|
|
|
|
"""Онлайн-копия базы с проверкой и ротацией. Возвращает путь копии или None, если копировать нечего.
|
2026-09-21 08:55:42 +03:00
|
|
|
|
|
|
|
|
|
|
Живая база и копия проверяются PRAGMA quick_check; при неудаче копия не появляется, старые не трогаются
|
2026-09-21 10:30:03 +03:00
|
|
|
|
(исключение 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")
|
2026-09-21 10:30:03 +03:00
|
|
|
|
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",
|
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
|
|
|
2026-09-21 10:06:31 +03:00
|
|
|
|
@contextmanager
|
|
|
|
|
|
def read_transaction(conn):
|
|
|
|
|
|
"""Один снимок базы (WAL) для нескольких запросов чтения: данные внутри блока согласованы между собой."""
|
|
|
|
|
|
conn.execute("BEGIN")
|
|
|
|
|
|
try:
|
|
|
|
|
|
yield
|
|
|
|
|
|
finally:
|
|
|
|
|
|
conn.execute("COMMIT")
|
|
|
|
|
|
|
|
|
|
|
|
|
2026-09-21 10:46:38 +03:00
|
|
|
|
# Результат последнего опроса каждого источника (см. 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))",
|
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
|
|
|
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}")
|
2026-09-21 10:46:38 +03:00
|
|
|
|
if version < 3:
|
|
|
|
|
|
for statement in _SOURCE_STATUS_SCHEMA:
|
|
|
|
|
|
conn.execute(statement)
|
2026-09-21 07:29:38 +03:00
|
|
|
|
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):
|
2026-09-21 10:46:38 +03:00
|
|
|
|
"""Удаляет все адреса источника и его состояние. True, если были адреса."""
|
|
|
|
|
|
conn.execute("DELETE FROM source_status WHERE kind = ? AND source = ?", (kind, source))
|
2026-09-21 07:29:38 +03:00
|
|
|
|
return conn.execute("DELETE FROM addresses WHERE kind = ? AND source = ?", (kind, source)).rowcount > 0
|
|
|
|
|
|
|
|
|
|
|
|
|
2026-09-21 10:46:38 +03:00
|
|
|
|
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")}
|
|
|
|
|
|
|
|
|
|
|
|
|
2026-09-21 07:29:38 +03:00
|
|
|
|
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).
|
|
|
|
|
|
"""
|
2026-09-21 10:06:31 +03:00
|
|
|
|
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)
|
2026-09-21 10:06:31 +03:00
|
|
|
|
|
|
|
|
|
|
# Наличие проверяется пакетами: один запрос на 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():
|
2026-09-21 10:06:31 +03:00
|
|
|
|
if action == "add" and (kind, value) in present:
|
2026-09-21 07:29:38 +03:00
|
|
|
|
added.add(value)
|
2026-09-21 10:06:31 +03:00
|
|
|
|
elif action == "del" and (kind, value) not in present:
|
2026-09-21 07:29:38 +03:00
|
|
|
|
removed.add(value)
|
|
|
|
|
|
return added, removed, head
|