diff --git a/Dockerfile b/Dockerfile index d2ec47c..7f91115 100644 --- a/Dockerfile +++ b/Dockerfile @@ -9,7 +9,10 @@ WORKDIR /app COPY requirements.txt constraints.txt ./ RUN pip install --no-cache-dir -c constraints.txt -r requirements.txt -COPY api_server.py cidr_collector.py collector_daemon.py db.py formatters.py storage.py healthcheck.py ./ +# Все модули приложения из корня (тесты и прочее лежат в подкаталогах или исключены .dockerignore) +COPY *.py ./ +# Забытый или сломанный модуль ломает сборку, а не запуск +RUN python -c "import api_server, collector_daemon, db, healthcheck" # Без root; том /data создаётся с владельцем ripe (именованный том наследует его при первом создании) RUN useradd --system --uid 10001 --no-create-home ripe \ diff --git a/README.md b/README.md index a7bdaf0..9bd4b45 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 `ripestat_sourceapp` (default `ripe-cidr-collector`) - the application name sent to RIPEstat as `sourceapp`; RIPEstat asks clients to identify themselves, you can add a contact (`my-collector admin@example.org`). Optional `allow_non_global_ips` (default `false`) - by default only public addresses resolved from FQDNs are stored; loopback, private, link-local and unspecified addresses (`127.0.0.1`, `10.x`, `0.0.0.0`, ...) are ignored and logged. Set `true` for internal names. 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). @@ -423,7 +424,7 @@ python3 cidr_collector.py run --mode fqdn When the collector runs (whether manually or via schedule): 1. **Instantiation**: Creates a new instance of `CIDRCollector` or `FQDNCollector`. This forces a fresh read of `config.json`, ensuring any added ASNs/FQDNs are immediately processed. 2. **Fetching**: - * **ASN**: Queries RIPE NCC API (`stat.ripe.net`). + * **ASN**: Queries RIPE NCC API (`stat.ripe.net`) with the `sourceapp` parameter. Connection failures and codes 429/500/502/503/504 are retried up to 3 times with growing pauses (0, 2, 4 s; a `Retry-After` pause is capped at 30 s), 10 s per attempt - at worst about 50 s per ASN. If all attempts fail the ASN is skipped for this run and nothing is deleted. * **FQDN**: Uses Python's `socket.getaddrinfo` to resolve A and AAAA records. Non-public addresses are dropped (see `allow_non_global_ips`); if nothing is left the run is treated like a DNS failure (nothing is deleted). 3. **Merge in one transaction**: the fetched addresses are merged into the SQLite table `addresses` (`db.py`) in a single write transaction, so readers (the API) never see a half-updated state. 4. **Accumulation with TTL**: Each address has `first_seen`/`last_seen`. @@ -483,6 +484,7 @@ docker compose down # stop; data stays in the volume ### Notes and risks - **Port 8000 is published without TLS**, so `X-API-Key` travels in clear text. Restrict access with a firewall or put a TLS reverse proxy in front (bind the port to `127.0.0.1` by changing `ports` in `docker-compose.yml`). +- The image copies all `*.py` modules from the project root and imports the main ones during the build, so a missing or broken module fails the build instead of the start. - Dependencies are pinned (see "Dependency versions" below), so rebuilds give the same libraries. The base image `python:3.11-slim` is not pinned by digest: Python patch releases arrive on rebuild. - Files written by the services in the volume (`config.json`, `status.json`) have mode `600`; both services run as the same user. - `docker compose` uses the image name `ripe-cidr-collector`; `Dockerfile.test` is used only for running the tests (section 1, step 5). diff --git a/cidr_collector.py b/cidr_collector.py index e7f853f..3b96ddb 100644 --- a/cidr_collector.py +++ b/cidr_collector.py @@ -7,6 +7,9 @@ import logging import socket import sys +from requests.adapters import HTTPAdapter +from urllib3.util.retry import Retry + import db from storage import StorageError, file_lock, load_json, save_json_atomic @@ -25,6 +28,10 @@ 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) RECREATED_FILE = os.path.join(DATA_DIR, "db_recreated.json") # база пересоздана после порчи без копий (см. db.recreated_pending) BASE_URL = "https://stat.ripe.net/data/announced-prefixes/data.json" +DEFAULT_SOURCEAPP = "ripe-cidr-collector" # RIPEstat просит указывать приложение (ключ ripestat_sourceapp в config.json) +RIPESTAT_TIMEOUT = 10 # секунд на одну попытку +RIPESTAT_RETRIES = 3 # повторов при сбоях соединения и кодах 429/5xx +MAX_RETRY_AFTER = 30 # потолок паузы по заголовку Retry-After, секунд: сборщик не должен зависать на часы DEFAULT_TTL_DAYS = 90 DEFAULT_BACKUP_KEEP = 7 # сколько копий базы хранить @@ -98,6 +105,24 @@ def pop_collection_requests(): return [t for t in types if t in COLLECT_TYPES] +class _BoundedRetry(Retry): + """Retry, у которого пауза по Retry-After ограничена MAX_RETRY_AFTER.""" + + def get_retry_after(self, response): + value = super().get_retry_after(response) + return None if value is None else min(value, MAX_RETRY_AFTER) + + +def make_session(): + """Сессия для RIPEstat: повторы с нарастающей паузой при сбоях соединения и кодах 429/500/502/503/504.""" + retry = _BoundedRetry(total=RIPESTAT_RETRIES, backoff_factor=1, status_forcelist=(429, 500, 502, 503, 504), + allowed_methods=("GET",)) + session = requests.Session() + session.mount("https://", HTTPAdapter(max_retries=retry)) + session.mount("http://", HTTPAdapter(max_retries=retry)) + return session + + def _is_global(address): """Публичный адрес (не loopback, не частный, не 0.0.0.0, не link-local и т. п.).""" try: @@ -111,6 +136,8 @@ class CIDRCollector: self.config = load_full_config() self.asns = self.config.get("asns", []) self.ttl_days = self.config.get("ttl_days", DEFAULT_TTL_DAYS) + self.sourceapp = self.config.get("ripestat_sourceapp", DEFAULT_SOURCEAPP) + self.session = None # создаётся при первом запросе, закрывается в конце сбора def add_asn(self, asn): added, self.asns = add_to_config_list("asns", asn) @@ -124,9 +151,11 @@ class CIDRCollector: print("Current ASNs:", self.asns) def fetch_prefixes(self, asn): - params = {'resource': f'AS{asn}'} + params = {'resource': f'AS{asn}', 'sourceapp': self.sourceapp} try: - response = requests.get(BASE_URL, params=params, timeout=10) + if self.session is None: + self.session = make_session() + response = self.session.get(BASE_URL, params=params, timeout=RIPESTAT_TIMEOUT) response.raise_for_status() data = response.json() @@ -144,11 +173,16 @@ class CIDRCollector: logger.info("Starting ASN CIDR collection...") # Сетевые запросы выполняем до транзакции, чтобы держать блокировку записи минимально fetched = {} - for asn in self.asns: - logger.info("Processing AS%s...", asn) - prefixes = self.fetch_prefixes(asn) - if prefixes is not None: - fetched[str(asn)] = set(prefixes) + try: + for asn in self.asns: + logger.info("Processing AS%s...", asn) + prefixes = self.fetch_prefixes(asn) + if prefixes is not None: + fetched[str(asn)] = set(prefixes) + finally: + if self.session is not None: + self.session.close() + self.session = None now = datetime.datetime.now() with db.session() as conn, db.transaction(conn): diff --git a/collector_daemon.py b/collector_daemon.py index 5e8be8a..9d69415 100644 --- a/collector_daemon.py +++ b/collector_daemon.py @@ -41,7 +41,8 @@ LOCK_FILE = os.path.join(cc.DATA_DIR, "collector.daemon.lock") job_state = {name: {"cron": None, "last_run": None, "last_finished": None, "running": False, "last_error": None, "rejected_cron": None} for name in COLLECTORS} -_state_lock = threading.Lock() +# RLock: изменения состояния и write_status (читает под той же блокировкой) не должны мешать друг другу +_state_lock = threading.RLock() # Один запуск за раз на тип: наложение планового и ручного сбора пропускается _run_locks = {name: threading.Lock() for name in COLLECTORS} @@ -68,18 +69,22 @@ def run_job(name, scheduler): try: logger.info("Running %s collection...", name) state = job_state[name] - state["last_run"] = datetime.datetime.now().isoformat() - state["running"] = True + with _state_lock: + error = state["last_error"] # при прерывании (SystemExit) прежняя ошибка остаётся + state["last_run"] = datetime.datetime.now().isoformat() + state["running"] = True write_status(scheduler) try: COLLECTORS[name]() - state["last_error"] = None + error = None except Exception as e: logger.exception("%s collection failed", name) - state["last_error"] = str(e) + error = str(e) finally: - state["running"] = False - state["last_finished"] = datetime.datetime.now().isoformat() + with _state_lock: + state["last_error"] = error + state["running"] = False + state["last_finished"] = datetime.datetime.now().isoformat() write_status(scheduler) finally: lock.release() @@ -104,7 +109,8 @@ def schedule_jobs(scheduler): cron = schedule.get(name, DEFAULT_CRONS[name]) scheduler.add_job(run_job, CronTrigger.from_crontab(cron), args=[name, scheduler], id=f"{name}_job", replace_existing=True, max_instances=1, coalesce=True) - job_state[name]["cron"] = cron + with _state_lock: + job_state[name]["cron"] = cron logger.info("Job %s scheduled: %s", name, cron) @@ -125,12 +131,14 @@ def sync_schedule(scheduler): trigger = CronTrigger.from_crontab(cron) except ValueError as e: logger.error("Invalid cron for %s (%r), keeping %r: %s", name, cron, state["cron"], e) - state["rejected_cron"] = cron + with _state_lock: + state["rejected_cron"] = cron continue scheduler.reschedule_job(f"{name}_job", trigger=trigger) logger.info("Job %s rescheduled: %s -> %s", name, state["cron"], cron) - state["cron"] = cron - state["rejected_cron"] = None + with _state_lock: + state["cron"] = cron + state["rejected_cron"] = None write_status(scheduler) diff --git a/docs/review-2026-09-21.md b/docs/review-2026-09-21.md index af5e6e0..32ec80f 100644 --- a/docs/review-2026-09-21.md +++ b/docs/review-2026-09-21.md @@ -16,12 +16,12 @@ | 2 | Средняя | **500 вместо 401 при нелатинском `X-API-Key` (подтверждено).** `secrets.compare_digest` для `str` работает только с ASCII. Обхода авторизации нет. | `api_server.py:105` | Исправлено (`summary-input-hardening.md`) | | 3 | Средняя | **500 при некорректном `since` (подтверждено):** `since=²` (`str.isdigit()` истинно для символов юникода, а `int()` их не принимает), число длиннее 4300 цифр (лимит `int`), крайние даты с поясом (`0001-01-01T00:00:00+05:00`, `9999-12-31T23:59:59-05:00`: `OverflowError`). | `api_server.py:152` | Исправлено (`summary-input-hardening.md`) | | 4 | Средняя | **DNS-адреса без фильтрации (подтверждено).** В списки попадают любые ответы: `127.0.0.1`, `10.x`, `0.0.0.0`. Оставлять нужно только глобальные адреса. | `cidr_collector.py:176-179` | Исправлено (`summary-input-hardening.md`) | -| 5 | Низкая | `/addresses` открывает три соединения и три разных снимка (курсор, ASN, FQDN): части ответа могут относиться к разным моментам, соединения тратят время на PRAGMA и проверку схемы. Достаточно одной сессии с одной транзакцией. | `api_server.py:133-141` | Не начато | -| 6 | Низкая | `/health` без авторизации отдаёт внутренние пути и текст ошибки (`last_restore`). | `api_server.py:214-226` | Не начато | -| 7 | Низкая | `run_job` меняет `job_state` вне `_state_lock`, хотя `write_status` читает под блокировкой (гонка безвредна благодаря GIL). | `collector_daemon.py:72` | Не начато | -| 8 | Низкая | `Dockerfile` копирует модули явным списком: новый модуль не попадёт в образ, ошибка проявится при запуске (healthcheck поймает, но поздно). | `Dockerfile:12` | Не начато | -| 9 | Низкая | Запросы к RIPEstat без повторов и без параметра `sourceapp`. | `cidr_collector.py:119` | Не начато | -| 10 | Низкая | `get_changes`: по запросу на каждое значение при проверке наличия; на десятках тысяч изменений станет заметно. | `db.py` | Не начато | +| 5 | Низкая | `/addresses` открывает три соединения и три разных снимка (курсор, ASN, FQDN): части ответа могут относиться к разным моментам, соединения тратят время на PRAGMA и проверку схемы. Достаточно одной сессии с одной транзакцией. | `api_server.py:133-141` | Исправлено (`summary-review-fixes-5-10.md`) | +| 6 | Низкая | `/health` без авторизации отдаёт внутренние пути и текст ошибки (`last_restore`). | `api_server.py:214-226` | Исправлено (`summary-review-fixes-5-10.md`) | +| 7 | Низкая | `run_job` меняет `job_state` вне `_state_lock`, хотя `write_status` читает под блокировкой (гонка безвредна благодаря GIL). | `collector_daemon.py:72` | Исправлено (`summary-review-fixes-5-10.md`) | +| 8 | Низкая | `Dockerfile` копирует модули явным списком: новый модуль не попадёт в образ, ошибка проявится при запуске (healthcheck поймает, но поздно). | `Dockerfile:12` | Исправлено (`summary-review-fixes-5-10.md`) | +| 9 | Низкая | Запросы к RIPEstat без повторов и без параметра `sourceapp`. | `cidr_collector.py:119` | Исправлено (`summary-review-fixes-5-10.md`) | +| 10 | Низкая | `get_changes`: по запросу на каждое значение при проверке наличия; на десятках тысяч изменений станет заметно. | `db.py` | Исправлено (`summary-review-fixes-5-10.md`) | ## Архитектура и сопровождение - Цикл `db` ↔ `cidr_collector` обходится отложенным импортом; чище вынести пути и значения по умолчанию в `settings.py`. @@ -36,7 +36,7 @@ - Порча данных не приводит к потере: карантин и копии. ## Пробелы тестов -Нет проверок задания `backup` в демоне, поля `last_restore` в `/health` и граничных значений `X-API-Key` и `since`. +Не было проверок задания `backup` в демоне, поля `last_restore` в `/health` и граничных значений `X-API-Key` и `since`; закрыты доработками `input-hardening` и `review-fixes-5-10`. ## Ограничения ревью - Граф не заменяет чтение кода: 14 связей типа INFERRED между модулями по отдельности не проверялись. @@ -45,4 +45,4 @@ ## Рекомендуемый порядок 1. Доработка «Безопасность ввода и данных»: п. 2, 3, 4 и тесты (`plan-input-hardening.md`). 2. Находка 1: после порчи без копий отвечать 503 до первого успешного сбора (`plan-loss-guard.md`, выполнено). -3. Пункты 5-10 объединить с доработкой наблюдаемости (п. 2 плана из анализа). +3. Пункты 5-10 исправлены отдельной доработкой (`plan-review-fixes-5-10.md`, выполнено); из ревью остаются только архитектурные замечания (вынос путей в `settings.py`, разделение `api_server.py` и `cidr_collector.py`). diff --git a/docs/summary-review-fixes-5-10.md b/docs/summary-review-fixes-5-10.md new file mode 100644 index 0000000..a9ade9b --- /dev/null +++ b/docs/summary-review-fixes-5-10.md @@ -0,0 +1,28 @@ +# Итоги: исправления 5-10 по результатам ревью + +План: `docs/plan-review-fixes-5-10.md`. Источник: `docs/review-2026-09-21.md`. Два этапа, два коммита; тестов стало 27 (было 24). + +## Этап А: путь чтения (находки 5, 6, 10), коммит `d6b69d0` +- **5. Один снимок в `/addresses`:** `db.read_transaction`, одна сессия вместо трёх; курсор в `X-Changes-Cursor` берётся в том же снимке, что и данные. Функции `get_cidrs`/`get_fqdn_ips` удалены. +- **6. `/health` без путей:** в `last_restore` только `{at, backup}` (имя файла копии); каталоги, карантин и текст ошибки остаются в `last_restore.json`. +- **10. `get_changes` пакетами:** проверка наличия порциями по `PRESENCE_BATCH = 500` вместо запроса на каждое значение; `get_changes` использует `read_transaction`. +- **Проверка:** ответы шести запросов (`/addresses` с разными параметрами, `/addresses/diff`) до и после совпали побайтно; `get_changes(cursor=0)` на журнале из 12 тыс. записей: 6832 мс -> 88 мс; при параллельной записи пар ASN/FQDN старый код дал 303 несогласованных ответа из 411, новый - 0 из 601. + +## Этап Б: демон, образ, внешний источник (находки 7, 8, 9) +- **7. Состояние заданий:** `_state_lock` стал `RLock`, все изменения `job_state` (`run_job`, `schedule_jobs`, `sync_schedule`) выполняются под ним; при прерывании (`SystemExit`) прежняя `last_error` сохраняется. +- **8. `Dockerfile`:** `COPY *.py ./` вместо списка модулей и проверка `import api_server, collector_daemon, db, healthcheck` на этапе сборки. +- **9. RIPEstat:** сессия `requests` с повторами (до 3, при сбоях соединения и кодах 429/500/502/503/504, паузы 0, 2, 4 с), параметр `sourceapp` (по умолчанию `ripe-cidr-collector`, ключ `ripestat_sourceapp` в `config.json`), сессия создаётся на запуск сбора и закрывается в конце. После исчерпания повторов источник пропускается, ничего не удаляется. +- **Тесты:** 2 новых (задание `backup` в демоне: регистрация, ротация по `backup_keep`; сессия RIPEstat: `sourceapp`, ключ из конфига, `None` после сбоя, настройка повторов и потолка `Retry-After`). +- **README:** `ripestat_sourceapp`, логика повторов, проверка образа. + +## Проверка этапа Б +- Локальный «RIPEstat»: два ответа 503, затем 200 -> результат за 2,0 с, три запроса, в каждом `sourceapp` (значение из конфига); постоянный 429 с `Retry-After: 3600` (потолок на время проверки 1 с) -> 4 попытки за 3,0 с, источник пропущен с записью в лог. +- Сборка образа: текущий код собирается; копия проекта без `formatters.py` не собирается (падает импорт); новый модуль в корне попадает в образ без правки `Dockerfile`. +- Docker Compose (отдельный проект, порт 18000, стенд убран): три задания в `/health`, ручной сбор FQDN отработал, расписание `backup` изменено через API и применено демоном (копия появилась), `status: ok`, сервисы `healthy`, ошибок в логе нет. +- Первая попытка негативной проверки образа не выполнилась (в среде нет `rsync`): повторена через `tar`. + +## Отличия от плана и замечания +- **Паузы 0, 2, 4 с, а не 1, 2, 4:** первая повторная попытка в `urllib3` идёт без паузы. Худший случай на один ASN около 46 с (4 попытки по 10 с и паузы). +- **`Retry-After` не «уважается» безусловно, а ограничен потолком `MAX_RETRY_AFTER = 30` с:** иначе ответ с `Retry-After: 3600` блокировал бы сборщик на час. В худшем случае с длинным `Retry-After` один ASN займёт до 130 с. +- `sourceapp` виден RIPE: не помещайте в него секреты; контакт администратора допустим. +- Оставлено: архитектурные замечания ревью (вынос путей в `settings.py`, разделение `api_server.py` и `cidr_collector.py`), ограничение частоты `POST /collect`, повторные запросы DNS. diff --git a/tests/test_core.py b/tests/test_core.py index 4f48dc2..c545bc3 100644 --- a/tests/test_core.py +++ b/tests/test_core.py @@ -97,3 +97,46 @@ def test_fqdn_keeps_only_global_addresses(files, monkeypatch): monkeypatch.setattr(cc.socket, "getaddrinfo", lambda *args, **kwargs: [answers[1]]) (files / "config.json").write_text('{"fqdns": ["x.example"]}') assert cc.FQDNCollector().resolve_fqdn("x.example") == [] + + +def test_ripestat_session_retries_and_sourceapp(files, monkeypatch): + calls = [] + real_make_session = cc.make_session + + class Response: + def raise_for_status(self): + pass + + def json(self): + return {"data": {"prefixes": [{"prefix": "1.0.0.0/24"}]}} + + class Session: + def get(self, url, params, timeout): + calls.append((params, timeout)) + return Response() + + def close(self): + pass + + monkeypatch.setattr(cc, "make_session", lambda: Session()) + assert cc.CIDRCollector().fetch_prefixes(1) == ["1.0.0.0/24"] + assert calls[0][0] == {"resource": "AS1", "sourceapp": "ripe-cidr-collector"} and calls[0][1] == cc.RIPESTAT_TIMEOUT + + # Своё имя приложения из config.json (например, с контактом) + (files / "config.json").write_text('{"asns": [], "ripestat_sourceapp": "my-app admin@example.org"}') + assert cc.CIDRCollector().fetch_prefixes(2) == ["1.0.0.0/24"] + assert calls[1][0]["sourceapp"] == "my-app admin@example.org" + + # После исчерпания повторов источник пропускается (None), сбор ничего не удаляет + collector = cc.CIDRCollector() + monkeypatch.setattr(collector, "sourceapp", "x") + Session.get = lambda self, url, params, timeout: (_ for _ in ()).throw(cc.requests.exceptions.RetryError("x")) + assert collector.fetch_prefixes(3) is None + + # Настройка сессии: повторы на сбои и коды 429/5xx, пауза Retry-After ограничена потолком + retry = real_make_session().get_adapter("https://stat.ripe.net").max_retries + assert retry.total == cc.RIPESTAT_RETRIES and 503 in retry.status_forcelist and 429 in retry.status_forcelist + + class Response: # сервер просит подождать час: пауза ограничена потолком + headers = {"Retry-After": "3600"} + assert retry.get_retry_after(Response()) == cc.MAX_RETRY_AFTER diff --git a/tests/test_daemon.py b/tests/test_daemon.py index ba6b82b..f9eda2b 100644 --- a/tests/test_daemon.py +++ b/tests/test_daemon.py @@ -10,6 +10,7 @@ sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) import api_server import cidr_collector as cc import collector_daemon as daemon +import db from storage import load_json, save_json_atomic @@ -95,3 +96,20 @@ def test_daemon_runs_requested_collection_once(env, monkeypatch): daemon.run_job("asn", scheduler) finally: daemon._run_locks["asn"].release() + + +def test_backup_job_rotates_and_is_registered(): + # Задание backup зарегистрировано рядом с asn и fqdn и использует расписание по умолчанию + assert daemon.COLLECTORS["backup"] is daemon.run_backup and "backup" in daemon.DEFAULT_CRONS + assert "backup" in daemon.job_state + + with open(cc.CONFIG_FILE, "w") as f: + f.write('{"asns": [], "fqdns": [], "backup_keep": 2}') + now = datetime.datetime.now() + with db.session() as conn: + for days in (1, 2, 3): + db.backup_database(conn, now - datetime.timedelta(days=days), keep=10) + daemon.run_backup() # keep=2 из config.json: остаются свежая копия и самая новая из старых + backups = db.list_backups(cc.BACKUP_DIR) + yesterday = (now - datetime.timedelta(days=1)).astimezone(datetime.timezone.utc).strftime("%Y%m%dT%H%M%SZ") + assert len(backups) == 2 and backups[1].endswith(f"ripe-{yesterday}.db") and backups[0] > backups[1]