Stage A of review fixes: single read snapshot, redacted /health, batched diff

/addresses reads the cursor and data in one session and one snapshot,
/health exposes only time and backup file name of the last restore,
get_changes checks presence in batches instead of per value (6.8 s -> 88 ms
on a 12k-entry journal). Adds the plan for review fixes 5-10.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
This commit is contained in:
ayurishchevandClaude Sonnet 5 committed 2026-09-21 10:06:31 +03:00
1 parent 0155fd3f38
commit d6b69d0842
6 files changed
+103 -33

No files matched your search

+2 -2
View File
@@ -340,7 +340,7 @@ curl "http://localhost:8000/addresses/diff?since=42" # next sync: use the r
``` ```
- `since` is a **cursor** from the previous answer (recommended: exact, independent of clocks) or an ISO 8601 time (no time zone = UTC; write `+` in a URL as `%2B`). Time is compared with millisecond precision, borders are inclusive, so an entry may be reported twice - repeating it is harmless. A cursor is at most 30 ASCII digits; anything else that is not an ISO 8601 time (including out-of-range dates) is `400`. - `since` is a **cursor** from the previous answer (recommended: exact, independent of clocks) or an ISO 8601 time (no time zone = UTC; write `+` in a URL as `%2B`). Time is compared with millisecond precision, borders are inclusive, so an entry may be reported twice - repeating it is harmless. A cursor is at most 30 ASCII digits; anything else that is not an ISO 8601 time (including out-of-range dates) is `400`.
- The result is the net effect: an address added and removed (or the other way) within the interval is not reported; an address that another source still holds is not reported as removed. Both CIDRs and IPs of FQDNs count. Output is JSON only, no aggregation. - The result is the net effect: an address added and removed (or the other way) within the interval is not reported; an address that another source still holds is not reported as removed. Both CIDRs and IPs of FQDNs count. Output is JSON only, no aggregation.
- The first sync: take the full `/addresses` list and the cursor from its response header `X-Changes-Cursor` (read before the data), then call `/addresses/diff?since=<that cursor>` regularly. - The first sync: take the full `/addresses` list and the cursor from its response header `X-Changes-Cursor` (taken in the same database snapshot as the data, so the two match exactly), then call `/addresses/diff?since=<that cursor>` regularly.
- `400` - invalid `since`; `410 Gone` - `since` is older than the journal (see `changes_retention_days`) or the cursor does not belong to this database: fetch the full `/addresses` list and continue from the new `cursor`. - `400` - invalid `since`; `410 Gone` - `since` is older than the journal (see `changes_retention_days`) or the cursor does not belong to this database: fetch the full `/addresses` list and continue from the new `cursor`.
### Endpoint: Manage Schedule ### Endpoint: Manage Schedule
@@ -396,7 +396,7 @@ Body is optional: `type` is `asn`, `fqdn` or `all` (default). The API does not c
### Endpoint: Health ### Endpoint: Health
**GET** `/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, the number of stored addresses and `last_restore` (`null`, or the record of the last automatic restore of the database from a backup) and `db_recreated` (`null`, or `{at, pending}` after the database was recreated without a backup; `pending: true` makes `status` `degraded`), 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). 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 `{at, backup}` - time and file name of the last automatic restore from a backup; directories and the error text stay in `last_restore.json`) and `db_recreated` (`null`, or `{at, pending}` after the database was recreated without a backup; `pending: true` makes `status` `degraded`), 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).
--- ---
+13 -23
View File
@@ -6,7 +6,7 @@ import os
import re import re
import secrets import secrets
import sqlite3 import sqlite3
from typing import List, Literal, Optional from typing import Literal, Optional
from apscheduler.triggers.cron import CronTrigger from apscheduler.triggers.cron import CronTrigger
from fastapi import Depends, FastAPI, Header, HTTPException, Path, Query, Response from fastapi import Depends, FastAPI, Header, HTTPException, Path, Query, Response
@@ -121,16 +121,6 @@ def ensure_data_ready(conn, kinds):
detail="The database was recreated after corruption; data is being collected again.") detail="The database was recreated after corruption; data is being collected again.")
def get_cidrs() -> List[str]:
with db.session() as conn:
return db.get_values(conn, "asn")
def get_fqdn_ips() -> List[str]:
with db.session() as conn:
return db.get_values(conn, "fqdn")
@app.get("/addresses", response_model=None, responses={ @app.get("/addresses", response_model=None, responses={
200: {"description": "JSON list for format=json, text/plain configuration script for other formats", 200: {"description": "JSON list for format=json, text/plain configuration script for other formats",
"content": {"application/json": {"schema": {"type": "array", "items": {"type": "string"}}}, "content": {"application/json": {"schema": {"type": "array", "items": {"type": "string"}}},
@@ -144,17 +134,14 @@ def get_addresses(
description="List/set name used in generated configuration"), description="List/set name used in generated configuration"),
response: Response = None, response: Response = None,
): ):
# Курсор читаем до данных: изменения между чтением курсора и данных повторятся в diff, что безвредно # Одна сессия и один снимок: курсор в заголовке точно соответствует выданным данным
with db.session() as conn: kinds = kinds_of(type)
ensure_data_ready(conn, kinds_of(type))
headers = {"X-Changes-Cursor": str(db.journal_head(conn))}
results = set() results = set()
with db.session() as conn, db.read_transaction(conn):
if type in [AddressType.cidr, AddressType.all_types]: ensure_data_ready(conn, kinds)
results.update(get_cidrs()) headers = {"X-Changes-Cursor": str(db.journal_head(conn))}
for kind in kinds:
if type in [AddressType.fqdn, AddressType.all_types]: results.update(db.get_values(conn, kind))
results.update(get_fqdn_ips())
output = formatters.build_output(results, format.value, ip_version.value, aggregate, name) output = formatters.build_output(results, format.value, ip_version.value, aggregate, name)
if format == OutputFormat.json: if format == OutputFormat.json:
@@ -233,9 +220,12 @@ def health():
pending = db.recreated_pending(conn, ("asn", "fqdn")) pending = db.recreated_pending(conn, ("asn", "fqdn"))
try: try:
last_restore = load_json(cc.RESTORE_FILE, None) # след автовосстановления базы из копии restore = load_json(cc.RESTORE_FILE, None) # след автовосстановления базы из копии
except StorageError: except StorageError:
last_restore = None restore = None
# Наружу - только время и имя файла копии; каталоги и текст ошибки остаются в last_restore.json
last_restore = ({"at": restore.get("at"), "backup": os.path.basename(str(restore.get("backup") or ""))}
if isinstance(restore, dict) else None)
try: try:
recreated = load_json(cc.RECREATED_FILE, None) # база пересоздана после порчи без копий recreated = load_json(cc.RECREATED_FILE, None) # база пересоздана после порчи без копий
+24 -8
View File
@@ -15,6 +15,7 @@ logger = logging.getLogger(__name__)
SCHEMA_VERSION = 2 SCHEMA_VERSION = 2
# После отката к копии номера журнала сдвигаются на это число: курсоры, выданные до отката, становятся недействительными # После отката к копии номера журнала сдвигаются на это число: курсоры, выданные до отката, становятся недействительными
RESTORE_JOURNAL_JUMP = 1_000_000 RESTORE_JOURNAL_JUMP = 1_000_000
PRESENCE_BATCH = 500 # значений в одном запросе проверки наличия (лимит переменных SQLite - не менее 999)
_UPSERT = ("INSERT INTO addresses (kind, source, value, first_seen, last_seen) VALUES (?, ?, ?, ?, ?) " _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") "ON CONFLICT (kind, source, value) DO UPDATE SET last_seen = excluded.last_seen")
@@ -262,6 +263,16 @@ _JOURNAL_SCHEMA = (
) )
@contextmanager
def read_transaction(conn):
"""Один снимок базы (WAL) для нескольких запросов чтения: данные внутри блока согласованы между собой."""
conn.execute("BEGIN")
try:
yield
finally:
conn.execute("COMMIT")
def _ensure_schema(conn, cc): def _ensure_schema(conn, cc):
if conn.execute("PRAGMA user_version").fetchone()[0] >= SCHEMA_VERSION: if conn.execute("PRAGMA user_version").fetchone()[0] >= SCHEMA_VERSION:
return return
@@ -399,8 +410,7 @@ def get_changes(conn, kinds, cursor=None, since_ts=None):
добавленное и удалённое (или наоборот) в этом интервале, в результат не попадает. добавленное и удалённое (или наоборот) в этом интервале, в результат не попадает.
Возвращает (added, removed, cursor) или None, если журнал не покрывает точку отсчёта (ответ 410). Возвращает (added, removed, cursor) или None, если журнал не покрывает точку отсчёта (ответ 410).
""" """
conn.execute("BEGIN") # один снимок для всех запросов with read_transaction(conn): # один снимок для всех запросов
try:
meta = dict(conn.execute("SELECT key, value FROM meta")) meta = dict(conn.execute("SELECT key, value FROM meta"))
horizon_id = int(meta["horizon_id"]) horizon_id = int(meta["horizon_id"])
head = conn.execute("SELECT COALESCE(MAX(id), ?) FROM changes", (horizon_id,)).fetchone()[0] head = conn.execute("SELECT COALESCE(MAX(id), ?) FROM changes", (horizon_id,)).fetchone()[0]
@@ -416,14 +426,20 @@ def get_changes(conn, kinds, cursor=None, since_ts=None):
"SELECT kind, value, action FROM changes WHERE id > ? ORDER BY id", (cursor,)): "SELECT kind, value, action FROM changes WHERE id > ? ORDER BY id", (cursor,)):
if kind in kinds: if kind in kinds:
first.setdefault((kind, value), action) 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() added, removed = set(), set()
for (kind, value), action in first.items(): for (kind, value), action in first.items():
present = conn.execute("SELECT 1 FROM addresses WHERE kind = ? AND value = ? LIMIT 1", if action == "add" and (kind, value) in present:
(kind, value)).fetchone() is not None
if action == "add" and present:
added.add(value) added.add(value)
elif action == "del" and not present: elif action == "del" and (kind, value) not in present:
removed.add(value) removed.add(value)
return added, removed, head return added, removed, head
finally:
conn.execute("COMMIT")
+40
View File
@@ -0,0 +1,40 @@
# План: исправления 5-10 по результатам ревью
Источник: `docs/review-2026-09-21.md`. Находки 1-4 закрыты (`summary-input-hardening.md`, `summary-loss-guard.md`). Здесь - низкие по серьёзности находки 5-10 и связанный пробел тестов (задание `backup`, `/health`). Влияние проверено по графу и `grep`: `get_cidrs`/`get_fqdn_ips` вызываются только из `get_addresses`; `run_job` вызывают планировщик, `check_collect_requests` и один тест.
## Два этапа (каждый со своими тестами и коммитом)
### Этап А: путь чтения (находки 5, 6, 10)
| # | Исправление | Файлы |
|---|---|---|
| 5 | **Одна сессия и один снимок в `/addresses`.** Новый контекст `db.read_transaction(conn)` (`BEGIN` ... `COMMIT`, единый снимок WAL). Внутри одной сессии: проверка `ensure_data_ready`, курсор `journal_head` и данные по всем запрошенным типам; форматирование ответа - после закрытия сессии. Функции `get_cidrs` и `get_fqdn_ips` удаляются. Итог: одно соединение вместо трёх, а курсор в заголовке `X-Changes-Cursor` точно соответствует выданным данным (комментарий «курсор читаем до данных» становится не нужен). | `db.py`, `api_server.py` |
| 6 | **`/health` не отдаёт внутренние пути.** `last_restore` в ответе - только `{"at", "backup"}` (`backup` - имя файла копии без каталога); полная запись (карантин, текст ошибки) остаётся в `last_restore.json` для оператора. | `api_server.py` |
| 10 | **Проверка наличия в `get_changes` пакетами.** Вместо запроса на каждое значение - `SELECT kind, value FROM addresses WHERE value IN (...)` порциями по 500 (в пределах лимита переменных SQLite), пары `(kind, value)` сверяются в Python. Результат тот же, число запросов падает с N до N/500. | `db.py` |
Тесты этапа А (новых 1, дополняется 1 существующий): `/health` не содержит путей (запись `last_restore.json` подставляется тестом); в `test_change_journal` добавляется пакет из 1200 значений (переход через границу порции), результат сверяется с ожидаемым.
### Этап Б: демон, образ, внешний источник (находки 7, 8, 9)
| # | Исправление | Файлы |
|---|---|---|
| 7 | **Состояние заданий под одной блокировкой.** `_state_lock` становится `RLock`, все изменения `job_state` (`run_job`, `schedule_jobs`, `sync_schedule`) выполняются `with _state_lock`; `write_status` вызывается после короткого обновления, как сейчас. Поведение не меняется, гонка между записью статуса и обновлением состояния исчезает. | `collector_daemon.py` |
| 8 | **`Dockerfile` не зависит от списка модулей.** `COPY *.py ./` (в корне лежат только модули приложения; тесты и прочее - в подкаталогах или исключены `.dockerignore`) и проверка на этапе сборки `RUN python -c "import api_server, collector_daemon, db, healthcheck"`: забытый или сломанный модуль ломает сборку, а не запуск. | `Dockerfile` |
| 9 | **Повторы и `sourceapp` для RIPEstat.** Сессия `requests` с `urllib3.Retry`: до 3 повторов при сбоях соединения и кодах 429/500/502/503/504, экспоненциальная пауза (уважается `Retry-After`); в запрос добавляется `sourceapp` (по умолчанию `ripe-cidr-collector`, переопределяется ключом `ripestat_sourceapp` в `config.json`, например с контактом). Сессия создаётся на один запуск сбора. После исчерпания повторов поведение прежнее: источник пропускается, ничего не удаляется. Худший случай на один ASN: около 50 с (4 попытки по 10 с и паузы), в расписании `*/15` это допустимо. | `cidr_collector.py`, `README.md` |
Тесты этапа Б (новых 2): задание `backup` в демоне (`run_backup` создаёт копию и уважает `backup_keep`; закрывает пробел из ревью); сессия RIPEstat (в запросе есть `sourceapp`, включённые повторы и коды, сбой после повторов даёт `None`).
## Итого
Новых тестов 3, всего 27 (плюс расширение существующего). README: `ripestat_sourceapp`, поведение повторов, `/health`.
## Порядок и проверка
1. Этап А -> тесты в контейнере -> вручную: параллельные запросы `/addresses` во время записи сборщика (данные и курсор согласованы), `/health` без путей, большой `since=0` на базе с несколькими тысячами изменений (время ответа до и после) -> коммит.
2. Этап Б -> тесты в контейнере -> вручную: сборка образа (проверка импорта проходит; намеренно удалённый модуль ломает сборку), сбой RIPEstat имитируется локальным сервером, который дважды отвечает 503, затем 200 (повторы срабатывают, в лог пишется предупреждение), демон и API в отдельных процессах, Docker Compose -> коммит.
3. Итоги `docs/summary-review-fixes-5-10.md`, статусы в отчёте ревью, README.
## Не входит
- Архитектурные пункты ревью (вынос путей в `settings.py`, разделение `api_server.py` и `cidr_collector.py`): отдельная доработка после наблюдаемости, чтобы не смешивать с исправлениями поведения.
- Ограничение частоты `POST /collect` (риск 8 анализа) и повторные запросы DNS.
## Риски и откат
- Этап А: изменения только на чтении, результат проверяется тестами и сравнением ответов до и после (ответы `/addresses` и `/addresses/diff` должны совпасть побайтно на одной базе).
- Этап Б: повторы увеличивают худшее время одного запуска сбора (описано выше); блокировка `RLock` безопасна для существующих вызовов.
- Откат - возврат к предыдущему коммиту этапа; схема данных не затрагивается.
+11
View File
@@ -83,6 +83,17 @@ def test_change_journal(files):
assert db.get_changes(conn, {"asn"}, cursor=c3 + 1) is None assert db.get_changes(conn, {"asn"}, cursor=c3 + 1) is None
assert db.get_changes(conn, {"asn"}, cursor=c3) == (set(), set(), c3) assert db.get_changes(conn, {"asn"}, cursor=c3) == (set(), set(), c3)
# Пакет больше порции проверки наличия (PRESENCE_BATCH): результат полный и точный
head = db.journal_head(conn)
bulk = {f"10.{i // 250}.{i % 250}.0/24" for i in range(2 * db.PRESENCE_BATCH + 200)}
with db.transaction(conn):
db.merge_source(conn, "asn", "3", bulk, now, 90)
assert db.get_changes(conn, {"asn"}, cursor=head)[:2] == (bulk, set())
head = db.journal_head(conn)
with db.transaction(conn):
db.purge_source(conn, "asn", "3")
assert db.get_changes(conn, {"asn"}, cursor=head)[:2] == (set(), bulk)
def test_backup_rotation_and_integrity(files): def test_backup_rotation_and_integrity(files):
now = datetime.datetime.now() now = datetime.datetime.now()
+13
View File
@@ -132,3 +132,16 @@ def test_recreated_database_is_withheld(env):
db.settle_recreated(conn) db.settle_recreated(conn)
assert env.get("/addresses").json() == ["1.0.0.0/24"] assert env.get("/addresses").json() == ["1.0.0.0/24"]
assert env.get("/health").json()["db_recreated"] is None assert env.get("/health").json()["db_recreated"] is None
def test_health_hides_internal_paths(env):
save_json_atomic(cc.RESTORE_FILE, {"at": "2026-09-21T10:00:00", "backup": "/data/backups/ripe-20260921T100000Z.db",
"quarantine": "/data/ripe.db.corrupt-1", "error": "file is not a database"})
# Наружу - время и имя файла копии, без каталогов, карантина и текста ошибки
assert env.get("/health").json()["last_restore"] == {"at": "2026-09-21T10:00:00",
"backup": "ripe-20260921T100000Z.db"}
assert json_dump_has_no_paths(env.get("/health").text)
def json_dump_has_no_paths(text):
return "/data" not in text and "corrupt" not in text and "not a database" not in text