From 02f7b49e126a75523417424582ca2a57595c01ed Mon Sep 17 00:00:00 2001 From: ayurishchev Date: Mon, 21 Sep 2026 08:55:42 +0300 Subject: [PATCH] Add database backups and automatic restore Daemon job "backup" (online copy, quick_check, rotation), restore of the newest valid copy when the database cannot be opened, change journal reset after restore, last_restore in /health, docs and tests. Co-Authored-By: Claude Sonnet 5 --- .dockerignore | 2 + .gitignore | 2 + README.md | 22 ++++-- api_server.py | 8 ++- cidr_collector.py | 4 ++ collector_daemon.py | 21 ++++-- db.py | 144 +++++++++++++++++++++++++++++++++++--- docs/plan-db-backup.md | 37 ++++++++++ docs/summary-db-backup.md | 24 +++++++ tests/test_db.py | 51 ++++++++++++++ 10 files changed, 294 insertions(+), 21 deletions(-) create mode 100644 docs/plan-db-backup.md create mode 100644 docs/summary-db-backup.md diff --git a/.dockerignore b/.dockerignore index f46e663..c6a1143 100644 --- a/.dockerignore +++ b/.dockerignore @@ -15,3 +15,5 @@ status.json *.migrated-* collect_request.json graphify-out/ +backups/ +last_restore.json diff --git a/.gitignore b/.gitignore index 9e750be..56a627d 100644 --- a/.gitignore +++ b/.gitignore @@ -14,3 +14,5 @@ collect_request.json graphify-out/ data.json fqdn_data.json +backups/ +last_restore.json diff --git a/README.md b/README.md index 77db7c5..7c25486 100644 --- a/README.md +++ b/README.md @@ -45,6 +45,7 @@ The system consists of **two independent processes**: the collector daemon (`col } ``` `ttl_days` - how long an address is kept after it was last seen (default `90`, `0` = keep forever). + Optional `backup_keep` - how many database backups to keep (default `7`); the backup cron is `schedule.backup` (default `30 4 * * *`), see section 9. Optional `changes_retention_days` - how long the change journal for `/addresses/diff` is kept (default `30`, `0` = forever). 5. **Tests** (optional) run in a container, not in the local `venv`: @@ -358,7 +359,7 @@ Body: "cron": "*/15 * * * *" } ``` -*Note: `type` can be `asn` or `fqdn`. The change is written to `config.json`; the collector daemon applies it within 30 seconds (`applied_within_seconds` in the response). An invalid cron expression is rejected by the API (400) and ignored by the daemon.* +*Note: `type` can be `asn`, `fqdn` or `backup` (database backup, see section 9). The change is written to `config.json`; the collector daemon applies it within 30 seconds (`applied_within_seconds` in the response). An invalid cron expression is rejected by the API (400) and ignored by the daemon.* ### Endpoints: Manage ASNs and FQDNs The lists of monitored sources are stored in `config.json` and can be managed through the API. Reading is open; changing requires the `X-API-Key` header (same token as `POST /schedule`). @@ -394,7 +395,7 @@ Body is optional: `type` is `asn`, `fqdn` or `all` (default). The API does not c ### Endpoint: Health **GET** `/health` -Reports the state of the collector daemon (read from `status.json`): `collector_alive`, the cron / `running` / last run / `last_finished` / last error / next run of each job, and the number of stored addresses. The daemon writes a heartbeat every 30 seconds; `collector_alive` is `false` if it is older than 120 seconds or the daemon never ran. `status` is `ok` only if the daemon is alive and no job failed in its last run; otherwise `degraded` (HTTP code is still 200). +Reports the state of the collector daemon (read from `status.json`): `collector_alive`, the cron / `running` / last run / `last_finished` / last error / next run of each job, the number of stored addresses and `last_restore` (`null`, or the record of the last automatic restore of the database from a backup, see section 9). The daemon writes a heartbeat every 30 seconds; `collector_alive` is `false` if it is older than 120 seconds or the daemon never ran. `status` is `ok` only if the daemon is alive and no job failed in its last run; otherwise `degraded` (HTTP code is still 200). --- @@ -433,7 +434,7 @@ When the collector runs (whether manually or via schedule): ### Scheduler Logic `collector_daemon.py` uses `APScheduler` (`BlockingScheduler`) in its own process; the API has no scheduler. -1. **Startup**: the daemon takes an exclusive lock (`collector.daemon.lock`, single instance), loads the `schedule` block from `config.json` and creates two independent jobs (`asn_job`, `fqdn_job`; defaults `0 2 * * *` and `0 3 * * *`). Overlapping runs of the same job are not allowed. +1. **Startup**: the daemon takes an exclusive lock (`collector.daemon.lock`, single instance), loads the `schedule` block from `config.json` and creates three independent jobs (`asn_job`, `fqdn_job`, `backup_job`; defaults `0 2 * * *`, `0 3 * * *` and `30 4 * * *`). Overlapping runs of the same job are not allowed. 2. **Runtime updates (POST /schedule)**: the API validates the cron expression and writes `config.json` atomically. Every 30 seconds the daemon compares the `schedule` block with the active one and reschedules the changed job (`reschedule_job`). An invalid cron expression is logged and ignored, the previous schedule stays. 3. **Status**: at the start and end of every run and every 30 seconds the daemon atomically writes `status.json` (jobs state + `updated_at` heartbeat); `GET /health` reads it. **Manual runs**: every 5 seconds the daemon checks `collect_request.json` (written by `POST /collect`) and starts the requested collections as one-off jobs; one run per type at a time. @@ -443,7 +444,7 @@ When the collector runs (whether manually or via schedule): ## 8. Application Setup: Docker Compose -One image (`Dockerfile`), two services: `api` (uvicorn) and `collector` (`collector_daemon.py`). They share the named volume `ripe_data` mounted at `/data` (`RIPE_DATA_DIR`), which holds `ripe.db` (with `-wal`/`-shm`), `config.json`, `status.json` and lock files. Containers run as non-root (uid 10001) with a read-only root filesystem, dropped capabilities and rotated logs. Both have healthchecks (`healthcheck.py`). +One image (`Dockerfile`), two services: `api` (uvicorn) and `collector` (`collector_daemon.py`). They share the named volume `ripe_data` mounted at `/data` (`RIPE_DATA_DIR`), which holds `ripe.db` (with `-wal`/`-shm`), `backups/` (database copies, see section 9), `config.json`, `status.json` and lock files. Containers run as non-root (uid 10001) with a read-only root filesystem, dropped capabilities and rotated logs. Both have healthchecks (`healthcheck.py`). ### Start ```bash @@ -517,5 +518,14 @@ Automatic: on the first start of any process (API, collector or CLI) `data.json` ### Rollback to the JSON version Stop both services, rename the `*.migrated-*` files back to `data.json` / `fqdn_data.json`, remove `ripe.db*`, start the previous version of the code. Addresses collected after the migration are lost in that case. -### Backup -Copy the database consistently with `sqlite3 ripe.db ".backup ripe.db.bak"` (do not copy `ripe.db` alone while the services are running - part of the data may still be in the `-wal` file). +### Backup and automatic restore +**Backup job.** The collector daemon runs the job `backup` (default `30 4 * * *`, change it with `POST /schedule` and `"type": "backup"` or in `schedule.backup`). It makes an online copy of `ripe.db` (safe while the services run), checks the live database and the copy with `PRAGMA quick_check`, and keeps the last `backup_keep` copies (default 7; older ones are deleted only after a new copy succeeded). If the live database fails the check or the copy is bad, no copy is written, the older copies stay, and `GET /health` shows `degraded` (`jobs.backup.last_error`). + +Copies are named `ripe-.db` and stored in `RIPE_BACKUP_DIR` (default `/backups`, i.e. `/data/backups` in Docker). **By default they are on the same volume as the database**: this protects against a corrupted file, not against losing the volume. For that, mount a separate volume/host directory and set `RIPE_BACKUP_DIR` to it, or copy the directory elsewhere regularly (e.g. `docker compose cp collector:/data/backups ./backups`). + +**Automatic restore.** If `ripe.db` cannot be opened as a database (any process: API, daemon, CLI), the file is moved to `ripe.db.corrupt-` and the newest copy that passes the integrity check is put in its place; concurrent processes are serialized with a lock. The event is logged (ERROR), written to `last_restore.json` and shown in `GET /health` as `last_restore`. Data collected after that copy is lost; the collector gathers it again on the next runs. If there is no valid copy, the behaviour is as before: the API answers `503`. + +- The change journal goes back with the copy, so after a restore all cursors and times issued earlier answer `410` on `/addresses/diff` (the client makes a full download). The journal counter is shifted by 1,000,000 for that (a heuristic: it assumes fewer changes than that between two copies). +- Only corruption detected **when the database is opened** is restored automatically. Damage inside the file shows up as `503` on reads and is caught by the `backup` job (`quick_check`); restore it by hand: stop both services, keep the damaged `ripe.db*`, copy the chosen `backups/ripe-*.db` to `ripe.db`, start the services. + +Manual consistent copy: `sqlite3 ripe.db ".backup ripe.db.bak"` (do not copy `ripe.db` alone while the services are running - part of the data may still be in the `-wal` file). diff --git a/api_server.py b/api_server.py index cda8115..fc7e11c 100644 --- a/api_server.py +++ b/api_server.py @@ -52,7 +52,7 @@ class CollectRequest(BaseModel): class ScheduleUpdate(BaseModel): - type: Literal["asn", "fqdn"] + type: Literal["asn", "fqdn", "backup"] cron: str model_config = {"json_schema_extra": {"example": {"type": "asn", "cron": "*/10 * * * *"}}} @@ -211,6 +211,11 @@ def health(): with db.session() as conn: counts = {"cidrs": db.count_values(conn, "asn"), "fqdn_ips": db.count_values(conn, "fqdn")} + try: + last_restore = load_json(cc.RESTORE_FILE, None) # след автовосстановления базы из копии + except StorageError: + last_restore = None + healthy = alive and not any(j.get("last_error") for j in jobs.values()) return { "status": "ok" if healthy else "degraded", @@ -218,6 +223,7 @@ def health(): "collector_updated_at": updated_at, "jobs": jobs, "counts": counts, + "last_restore": last_restore, } diff --git a/cidr_collector.py b/cidr_collector.py index 478b46d..753bdd7 100644 --- a/cidr_collector.py +++ b/cidr_collector.py @@ -19,9 +19,13 @@ DATA_FILE = os.path.join(DATA_DIR, "data.json") FQDN_DATA_FILE = os.path.join(DATA_DIR, "fqdn_data.json") STATUS_FILE = os.path.join(DATA_DIR, "status.json") # пишет collector_daemon, читает API (/health) COLLECT_REQUEST_FILE = os.path.join(DATA_DIR, "collect_request.json") # API -> демон: POST /collect +# Резервные копии базы (задание backup в демоне); по умолчанию на том же томе, каталог можно вынести +BACKUP_DIR = os.environ.get("RIPE_BACKUP_DIR", os.path.join(DATA_DIR, "backups")) +RESTORE_FILE = os.path.join(DATA_DIR, "last_restore.json") # след автовосстановления базы (читает /health) BASE_URL = "https://stat.ripe.net/data/announced-prefixes/data.json" DEFAULT_TTL_DAYS = 90 +DEFAULT_BACKUP_KEEP = 7 # сколько копий базы хранить DEFAULT_CHANGES_RETENTION_DAYS = 30 # срок хранения журнала изменений для /addresses/diff # Демон сверяет расписание с config.json и пишет heartbeat в status.json с этим периодом; # API считает демон живым, пока heartbeat не старше STATUS_STALE_AFTER секунд diff --git a/collector_daemon.py b/collector_daemon.py index 0506d24..5e8be8a 100644 --- a/collector_daemon.py +++ b/collector_daemon.py @@ -14,13 +14,27 @@ from apscheduler.schedulers.blocking import BlockingScheduler from apscheduler.triggers.cron import CronTrigger import cidr_collector as cc +import db from cidr_collector import CIDRCollector, FQDNCollector, load_full_config from storage import StorageError, save_json_atomic, try_lock logger = logging.getLogger("collector_daemon") -DEFAULT_CRONS = {"asn": "0 2 * * *", "fqdn": "0 3 * * *"} -COLLECTORS = {"asn": CIDRCollector, "fqdn": FQDNCollector} +DEFAULT_CRONS = {"asn": "0 2 * * *", "fqdn": "0 3 * * *", "backup": "30 4 * * *"} + + +def run_backup(): + """Задание backup: онлайн-копия базы, проверка, ротация (число копий - backup_keep в config.json).""" + keep = max(1, int(load_full_config().get("backup_keep", cc.DEFAULT_BACKUP_KEEP))) + with db.session() as conn: + path = db.backup_database(conn, datetime.datetime.now(), keep) + logger.info("Database backup written to %s (keeping the last %d)", path, keep) + + +# Задания: имя -> вызываемый объект (новый экземпляр на каждый запуск - свежий config.json) +COLLECTORS = {"asn": lambda: CIDRCollector().run_collection(), + "fqdn": lambda: FQDNCollector().run_collection(), + "backup": run_backup} LOCK_FILE = os.path.join(cc.DATA_DIR, "collector.daemon.lock") # Состояние заданий: действующий cron, результат последнего запуска, отклонённый cron (чтобы не спамить логом) @@ -58,8 +72,7 @@ def run_job(name, scheduler): state["running"] = True write_status(scheduler) try: - # Новый экземпляр на каждый запуск - свежий config.json - COLLECTORS[name]().run_collection() + COLLECTORS[name]() state["last_error"] = None except Exception as e: logger.exception("%s collection failed", name) diff --git a/db.py b/db.py index a12f6bc..2683e9b 100644 --- a/db.py +++ b/db.py @@ -1,16 +1,20 @@ """Хранение собранных адресов в 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, load_json +from storage import StorageError, file_lock, load_json, save_json_atomic logger = logging.getLogger(__name__) SCHEMA_VERSION = 2 +# После отката к копии номера журнала сдвигаются на это число: курсоры, выданные до отката, становятся недействительными +RESTORE_JOURNAL_JUMP = 1_000_000 _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") @@ -29,28 +33,148 @@ def _quarantine(path): 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 +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 sqlite3.OperationalError as e: # например, database is locked: базу не трогаем + 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 и демон могут обнаружить порчу одновременно, поэтому всё выполняется под блокировкой, + а внутри неё база проверяется повторно (другой процесс мог уже восстановить). Нет копий - StorageError. + """ + 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 + 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 + raise StorageError(f"{path} is corrupted and no valid backup was found") from error + + +def _restore_backup(candidate, path, cc): + """Кладёт копию на место базы и сбрасывает журнал изменений (его номера после точки копии уже выдавались).""" + tmp = path + ".restore.tmp" + shutil.copyfile(candidate, tmp) + os.replace(tmp, path) + conn = _open(path, cc) + try: + 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") + except BaseException: 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 + raise return conn +def backup_database(conn, now, keep, backup_dir=None): + """Онлайн-копия базы с проверкой и ротацией. Возвращает путь копии. + + Живая база и копия проверяются PRAGMA quick_check; при неудаче копия не появляется, старые не трогаются + (исключение StorageError). Лишние старые копии (сверх 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") + 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(): """Соединение на время блока (закрывается по выходу).""" diff --git a/docs/plan-db-backup.md b/docs/plan-db-backup.md new file mode 100644 index 0000000..2b85ddf --- /dev/null +++ b/docs/plan-db-backup.md @@ -0,0 +1,37 @@ +# План: резервные копии базы и автовосстановление (доработка Б, п. 1.3-1.4) + +## Цель +База `ripe.db` копируется автоматически по расписанию, а при порче файла сервис сам поднимает последнюю исправную копию вместо пустой базы (риск 3 анализа). Откат к копии заметен в логе и в `/health`, а не происходит молча. + +## Дизайн + +### Резервное копирование +- **Задание `backup` в демоне** (третье после `asn` и `fqdn`, те же расписание, состояние в `status.json`, защита от наложения). Расписание: `schedule.backup` в `config.json`, по умолчанию `30 4 * * *`; менять можно через `POST /schedule` (тип `backup`). +- **Как копируем:** онлайн-копия средствами SQLite (`Connection.backup`), безопасная при работающих API и сборщике (WAL, чтение снимка). Файл сначала пишется как `*.tmp`, затем переименовывается: неполная копия не появляется. +- **Проверка:** перед копированием `PRAGMA quick_check` живой базы, после копирования она же на копии. Если проверка не пройдена, копия не создаётся или удаляется, старые хорошие копии не трогаются, задание получает `last_error`, `/health` показывает `degraded`. +- **Хранение:** `RIPE_BACKUP_DIR` (по умолчанию `/backups`), файлы `ripe-.db`. Ключ `backup_keep` в `config.json` (по умолчанию 7): лишние старые копии удаляются только после успешной новой. +- **Ограничение:** по умолчанию копии лежат на том же томе `/data`, что и база. Это защищает от порчи файла, но не от потери тома. Вынос на отдельный том или хост описывается в README (переменная `RIPE_BACKUP_DIR` и копирование каталога), не автоматизируется. + +### Автовосстановление +- **Где:** в `db.connect()`, там, где сейчас порченая база убирается в `*.corrupt-*`. После этого берётся самая свежая копия, проходящая `quick_check`; копия кладётся на место базы (через временный файл и `os.replace`), соединение открывается заново. Если исправных копий нет, поведение прежнее: `StorageError`, API отвечает 503. +- **Гонка API и демона:** оба процесса могут обнаружить порчу одновременно. Восстановление идёт под файловой блокировкой (`file_lock`), внутри блокировки база проверяется повторно: если другой процесс уже восстановил, повторных действий нет. +- **Курсоры diff после отката:** журнал изменений откатывается вместе с базой, и номера записей после точки копии могли уже быть выданы клиентам. Чтобы старые курсоры не совпали с новой историей, после восстановления счётчик журнала сдвигается на 1 000 000 и «горизонт» ставится на новое значение (время горизонта = момент восстановления). Все прежние курсоры и времена дают `410`, клиенты делают полную выгрузку. Это эвристика: она надёжна, пока между копиями накапливается меньше миллиона изменений (при суточной копии и редких изменениях RIPE запас большой). +- **Заметность:** запись в лог уровня ERROR и файл `last_restore.json` (время, использованная копия, путь карантина); `/health` показывает поле `last_restore`. Статус `degraded` восстановление не вызывает (сервис работает), но событие не теряется. +- **Ограничение:** автоматически ловится порча, обнаруженная при открытии базы. Порча страниц внутри файла проявится при чтении (ответ 503) и будет замечена заданием `backup` (`quick_check` не пройдёт, `/health` = `degraded`); восстановление в этом случае ручное (порядок в README). + +### Что не входит +Копии вне сервера, шифрование, восстановление на момент времени, команда CLI `restore`, ручной запуск копии через `POST /collect` (можно добавить отдельно). + +## Изменения +1. `cidr_collector.py`: `BACKUP_DIR`, `RESTORE_FILE`, `DEFAULT_BACKUP_KEEP`. +2. `db.py`: `backup_database()` (копия, проверка, ротация), `_restore_latest_backup()` под блокировкой, вызов из `connect()`, сдвиг журнала. +3. `collector_daemon.py`: задание `backup` (реестр `COLLECTORS` становится реестром заданий с вызываемыми объектами, `DEFAULT_CRONS["backup"]`). +4. `api_server.py`: `ScheduleUpdate.type` допускает `backup`; `/health` показывает `last_restore`. +5. `README.md`: расписание и ключи (`schedule.backup`, `backup_keep`), `RIPE_BACKUP_DIR`, ручное восстановление, вынос копий с тома. +6. Тесты (2): копия проходит проверку, ротация оставляет `backup_keep` штук, повреждённая живая база копию не создаёт; порченая база при открытии восстанавливается из копии (карантин сохранён, старые курсоры дают `410`), без копий остаётся `StorageError`. + +## Проверка +Тесты в контейнере. Вручную на копии данных: задание `backup` создаёт файл, порча `ripe.db` (запись мусора) при работающих API и демоне приводит к восстановлению, `/health` показывает `last_restore`, запрос `since` со старым курсором даёт `410`. Проверка в Docker Compose (общий том). + +## Откат +Удалить задание `backup` и вызов восстановления в `connect()`: схема базы не меняется (журнал и `meta` остаются), файлы копий можно удалить вручную. diff --git a/docs/summary-db-backup.md b/docs/summary-db-backup.md new file mode 100644 index 0000000..99e09b2 --- /dev/null +++ b/docs/summary-db-backup.md @@ -0,0 +1,24 @@ +# Итоги: резервные копии базы и автовосстановление (доработка Б) + +План: `docs/plan-db-backup.md`. + +## Сделано +- **Задание `backup` в демоне** (`collector_daemon.py`): третье задание рядом с `asn` и `fqdn` (то же расписание, защита от наложения, состояние в `status.json`); по умолчанию `30 4 * * *`, меняется в `schedule.backup` или через `POST /schedule` (`type: backup`). Реестр `COLLECTORS` теперь хранит вызываемые объекты, существующие тесты не менялись. +- **Копия** (`db.backup_database`): онлайн-копия средствами SQLite, `quick_check` живой базы и копии, запись через `*.tmp` и переименование, ротация (`backup_keep`, по умолчанию 7) только после успешной новой. Если проверка не пройдена, копия не создаётся, старые остаются, `/health` показывает `degraded`. +- **Автовосстановление** (`db.connect` -> `_recover`): порченая база уходит в `*.corrupt-*`, на её место встаёт самая свежая копия, прошедшая проверку (битые копии пропускаются); под файловой блокировкой, с повторной проверкой внутри (гонка API и демона). Без исправных копий прежнее поведение (`StorageError`, 503). +- **Журнал diff после отката:** журнал очищается, счётчик сдвигается на 1 000 000, горизонт на момент восстановления: все выданные ранее курсоры и времена дают `410`. +- **Заметность:** лог ERROR, файл `last_restore.json`, поле `last_restore` в `/health`. +- **Прочее:** `RIPE_BACKUP_DIR` (по умолчанию `/backups`), `backups/` и `last_restore.json` в `.gitignore`/`.dockerignore`, `POST /schedule` принимает `backup`. +- **README:** расписание и `backup_keep`, раздел «Backup and automatic restore» (вынос копий с тома, ручное восстановление), `last_restore` в `/health`. +- **Тесты:** 2 новых (копия, ротация и проверка; восстановление с пропуском битой копии, сброс курсоров, работа после отката, случай без копий). Всего 20, в контейнере 20 passed. + +## Проверка +- Отдельные процессы (демон с `backup` раз в минуту и `backup_keep=2`, API): скопировано и ротировано до двух файлов; после записи мусора в `ripe.db` API поднял базу из копии, `/addresses` вернул данные, `/health` показал `last_restore`, карантинный файл сохранён; новый курсор принимается (200). +- Docker Compose (отдельный проект, порт 18000, стенд убран): расписание `backup` выставлено через `POST /schedule`, копия появилась в `/data/backups`, порча базы в томе привела к восстановлению, оба сервиса `healthy`. + +## Замечания +- Копии по умолчанию лежат на том же томе, что и база: это защита от порчи файла, но не от потери тома. +- Автоматически восстанавливается только порча, видимая при открытии; порча внутри файла даёт 503, её замечает задание `backup`, восстановление ручное (порядок в README). +- Сдвиг журнала на 1 000 000 - эвристика: старый курсор мог бы совпасть с новой историей, только если между копиями накопится больше миллиона изменений. +- Отказ `410` для старого курсора после восстановления подтверждён тестом; вручную в отдельных процессах эту проверку не повторял (команда не выполнилась из-за экранирования в оболочке). +- Не сделано: копии вне сервера, шифрование, команда CLI `restore`, ручной запуск копии через `POST /collect`. diff --git a/tests/test_db.py b/tests/test_db.py index daf8cf6..3aaaef3 100644 --- a/tests/test_db.py +++ b/tests/test_db.py @@ -1,6 +1,8 @@ import datetime +import glob import json import os +import sqlite3 import sys import pytest @@ -15,6 +17,8 @@ import db def files(tmp_path, monkeypatch): for attr, name in (("DB_FILE", "ripe.db"), ("DATA_FILE", "data.json"), ("FQDN_DATA_FILE", "fqdn_data.json")): monkeypatch.setattr(cc, attr, str(tmp_path / name)) + monkeypatch.setattr(cc, "BACKUP_DIR", str(tmp_path / "backups")) + monkeypatch.setattr(cc, "RESTORE_FILE", str(tmp_path / "last_restore.json")) return tmp_path @@ -76,3 +80,50 @@ def test_change_journal(files): assert db.get_changes(conn, {"asn"}, cursor=0) is None assert db.get_changes(conn, {"asn"}, cursor=c3 + 1) is None assert db.get_changes(conn, {"asn"}, cursor=c3) == (set(), set(), c3) + + +def test_backup_rotation_and_integrity(files): + now = datetime.datetime.now() + with db.session() as conn: + with db.transaction(conn): + db.merge_source(conn, "asn", "1", {"a"}, now, 90) + for minutes in range(3): + db.backup_database(conn, now + datetime.timedelta(minutes=minutes), keep=2) + + # Ротация оставляет две самые свежие копии; копия самодостаточна и содержит данные + backups = db.list_backups(cc.BACKUP_DIR) + assert len(backups) == 2 and not glob.glob(os.path.join(cc.BACKUP_DIR, "*.tmp")) + assert backups[0].endswith(f"{(now + datetime.timedelta(minutes=2)).astimezone(datetime.timezone.utc):%Y%m%dT%H%M%SZ}.db") + assert sqlite3.connect(backups[0]).execute("SELECT COUNT(*) FROM addresses").fetchone()[0] == 1 + + +def test_restore_from_backup(files): + now = datetime.datetime.now() + with db.session() as conn: + with db.transaction(conn): + db.merge_source(conn, "asn", "1", {"a", "b"}, now, 90) + old_cursor = db.get_changes(conn, {"asn"}, cursor=0)[2] + db.backup_database(conn, now, keep=3) + + # «Свежая» копия с мусором пропускается; живая база портится + (files / "backups" / "ripe-99999999T999999Z.db").write_bytes(b"garbage" * 100) + (files / "ripe.db").write_bytes(b"garbage" * 1000) + + with db.session() as conn: + assert sorted(db.get_values(conn, "asn")) == ["a", "b"] + # Курсор, выданный до отката, недействителен (клиент сделает полную выгрузку); новый работает + assert db.get_changes(conn, {"asn"}, cursor=old_cursor) is None + head = db.journal_head(conn) + assert db.get_changes(conn, {"asn"}, cursor=head) == (set(), set(), head) + with db.transaction(conn): + db.merge_source(conn, "asn", "1", {"a", "b", "c"}, now, 90) + assert db.get_changes(conn, {"asn"}, cursor=head)[:2] == ({"c"}, set()) + assert len(glob.glob(str(files / "ripe.db.corrupt-*"))) == 1 + assert json.loads((files / "last_restore.json").read_text())["backup"].endswith(".db") + + # Без исправных копий - прежнее поведение: ошибка хранилища, база в карантине + for backup in db.list_backups(cc.BACKUP_DIR): + os.unlink(backup) + (files / "ripe.db").write_bytes(b"garbage" * 1000) + with pytest.raises(db.StorageError): + db.connect()