diff --git a/README.md b/README.md index d8fb438..d49fa6e 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 `source_failure_threshold` (default `3`) - how many runs in a row a source may fail before `GET /health` reports `degraded` (see the sources status below). 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. @@ -399,6 +400,11 @@ Body is optional: `type` is `asn`, `fqdn` or `all` (default). The API does not c **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 `{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). +`sources` shows the state of the data sources: for `asn` and `fqdn` the number of configured sources (`total`) and the list of **failing** ones (`failing`: `source`, `failures`, `last_success`, `error_kind`). A source is failing when its last `source_failure_threshold` runs in a row failed (default 3, so a single RIPEstat hiccup does not turn the status yellow); any failing source makes `status` `degraded`. `error_kind` is a short category without details (`timeout`, `connection`, `http_429`, `http_4xx`, `http_5xx`, `invalid_response`, `dns`, `no_global_addresses`, `unknown`); the full error text stays in the collector log. + +### Endpoint: Sources status +**GET** `/sources` (open) - the state of every configured source: `kind`, `source`, `addresses` (how many values it gave), `last_attempt`, `last_success`, `failures` (runs failed in a row), `error_kind`, plus `failure_threshold`. A source that was never polled has empty state; a source that "succeeds" but gives `addresses: 0` is easy to spot here. The state is kept in the database (table `source_status`), survives restarts and is included in backups; it is removed with the source (`purge` or the next collection after removal). + --- ## 6. Advanced CLI Usage @@ -424,7 +430,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`) 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. + * **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; the failure is recorded for `GET /health` and `GET /sources`. Every run ends with one summary line in the log (`ASN collection finished: 6 sources, 5 ok, 1 failed, +12/-3 prefixes`). * **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`. @@ -508,7 +514,7 @@ CREATE TABLE addresses ( The schema version is stored in `PRAGMA user_version`. The database runs in WAL mode (files `ripe.db-wal`, `ripe.db-shm` next to it); API and collector (also two containers sharing one local volume) can work with it concurrently. Do not place it on a network file system. ### Change journal -For `/addresses/diff` the database keeps the tables `changes(id, ts, kind, value, action add|del)` and `meta` (the journal "horizon"). SQL triggers on `addresses` write to `changes`, so every path (collection, TTL, `purge`, removal of a source) is covered: `add` - the value appeared and no other source had it, `del` - the last entry of the value is gone. `ts` is UTC. The import from JSON is not written to the journal. Old records are deleted by the collector (`changes_retention_days`); the schema version is 2 (a version 1 database is upgraded automatically on the first start, data is kept). +For `/addresses/diff` the database keeps the tables `changes(id, ts, kind, value, action add|del)` and `meta` (the journal "horizon"). SQL triggers on `addresses` write to `changes`, so every path (collection, TTL, `purge`, removal of a source) is covered: `add` - the value appeared and no other source had it, `del` - the last entry of the value is gone. `ts` is UTC. The import from JSON is not written to the journal. Old records are deleted by the collector (`changes_retention_days`); the schema version is 3 (an older database is upgraded automatically on the first start, data is kept). Version 3 adds the table `source_status(kind, source, last_attempt, last_success, failures, error_kind)`. Inspect the data: ```bash diff --git a/api_server.py b/api_server.py index f545cb6..761b4eb 100644 --- a/api_server.py +++ b/api_server.py @@ -210,6 +210,28 @@ def collector_state(): return alive, updated_at, (status.get("jobs", {}) if status else {}) +def source_health(): + """Сводка по источникам для /health: всего и только неисправные (сбоев подряд не меньше порога).""" + try: + config = load_full_config() + except StorageError: + return None, False + threshold = cc.failure_threshold(config) + with db.session() as conn: + statuses = db.source_statuses(conn) + summary, any_failing = {}, False + for kind, key in (("asn", "asns"), ("fqdn", "fqdns")): + configured = [str(source) for source in config.get(key, [])] + failing = [{"source": source, "failures": statuses[kind, source]["failures"], + "last_success": statuses[kind, source]["last_success"], + "error_kind": statuses[kind, source]["error_kind"]} + for source in configured + if (kind, source) in statuses and statuses[kind, source]["failures"] >= threshold] + summary[kind] = {"total": len(configured), "failing": failing} + any_failing = any_failing or bool(failing) + return summary, any_failing + + @app.get("/health") def health(): """Состояние сборщика: демон пишет status.json (задания + heartbeat), API только читает.""" @@ -232,13 +254,15 @@ def health(): except StorageError: recreated = None - healthy = alive and not pending and not any(j.get("last_error") for j in jobs.values()) + sources, any_failing = source_health() + healthy = alive and not pending and not any_failing and not any(j.get("last_error") for j in jobs.values()) return { "status": "ok" if healthy else "degraded", "collector_alive": alive, "collector_updated_at": updated_at, "jobs": jobs, "counts": counts, + "sources": sources, "last_restore": last_restore, # Только время и признак ожидания (пути карантина наружу не отдаём) "db_recreated": {"at": recreated.get("at"), "pending": pending} if recreated else None, @@ -275,6 +299,19 @@ def _purge(kind, source): return db.purge_source(conn, kind, source) +@app.get("/sources") +def list_sources(): + """Состояние всех настроенных источников: сколько адресов дал, когда опрашивался, сколько сбоев подряд.""" + config = load_full_config() + with db.session() as conn: + statuses, counts = db.source_statuses(conn), db.address_counts(conn) + empty = {"last_attempt": None, "last_success": None, "failures": 0, "error_kind": None} + return {"failure_threshold": cc.failure_threshold(config), "sources": [ + {"kind": kind, "source": str(source), "addresses": counts.get((kind, str(source)), 0), + **statuses.get((kind, str(source)), empty)} + for kind, key in (("asn", "asns"), ("fqdn", "fqdns")) for source in config.get(key, [])]} + + @app.get("/asns") def list_asns(): return {"asns": load_full_config().get("asns", [])} diff --git a/cidr_collector.py b/cidr_collector.py index 3b96ddb..dd421e4 100644 --- a/cidr_collector.py +++ b/cidr_collector.py @@ -2,6 +2,7 @@ import requests import datetime import ipaddress import os +import re import argparse import logging import socket @@ -35,6 +36,7 @@ MAX_RETRY_AFTER = 30 # потолок паузы по заголовку Retry DEFAULT_TTL_DAYS = 90 DEFAULT_BACKUP_KEEP = 7 # сколько копий базы хранить +DEFAULT_SOURCE_FAILURE_THRESHOLD = 3 # сколько запусков подряд источник может сбоить, прежде чем /health станет degraded DEFAULT_CHANGES_RETENTION_DAYS = 30 # срок хранения журнала изменений для /addresses/diff # Демон сверяет расписание с config.json и пишет heartbeat в status.json с этим периодом; # API считает демон живым, пока heartbeat не старше STATUS_STALE_AFTER секунд @@ -113,6 +115,35 @@ class _BoundedRetry(Retry): return None if value is None else min(value, MAX_RETRY_AFTER) +def classify_error(exc): + """Категория сбоя источника без деталей: полный текст ошибки остаётся в логе и наружу не отдаётся.""" + text = str(exc) + code = None + if isinstance(exc, requests.exceptions.HTTPError) and getattr(exc.response, "status_code", None): + code = exc.response.status_code + elif isinstance(exc, requests.exceptions.RetryError): + match = re.search(r"too many (\d{3}) error responses", text) + code = int(match.group(1)) if match else None + if code: + return "http_429" if code == 429 else "http_5xx" if code >= 500 else "http_4xx" + if isinstance(exc, requests.exceptions.Timeout) or "timed out" in text.lower(): + return "timeout" + if isinstance(exc, requests.exceptions.ConnectionError): + return "connection" + if isinstance(exc, ValueError): # включая нечитаемый JSON + return "invalid_response" + return "unknown" + + +def failure_threshold(config): + """Порог сбоев подряд для /health (ключ source_failure_threshold; неверное значение - по умолчанию).""" + try: + value = int(config.get("source_failure_threshold", DEFAULT_SOURCE_FAILURE_THRESHOLD)) + except (TypeError, ValueError): + return DEFAULT_SOURCE_FAILURE_THRESHOLD + return value if value >= 1 else DEFAULT_SOURCE_FAILURE_THRESHOLD + + def make_session(): """Сессия для RIPEstat: повторы с нарастающей паузой при сбоях соединения и кодах 429/500/502/503/504.""" retry = _BoundedRetry(total=RIPESTAT_RETRIES, backoff_factor=1, status_forcelist=(429, 500, 502, 503, 504), @@ -138,6 +169,7 @@ class CIDRCollector: self.ttl_days = self.config.get("ttl_days", DEFAULT_TTL_DAYS) self.sourceapp = self.config.get("ripestat_sourceapp", DEFAULT_SOURCEAPP) self.session = None # создаётся при первом запросе, закрывается в конце сбора + self.errors = {} # источник -> категория сбоя последнего опроса (для source_status) def add_asn(self, asn): added, self.asns = add_to_config_list("asns", asn) @@ -167,6 +199,7 @@ class CIDRCollector: return prefixes except Exception as e: logger.error("Error fetching data for AS%s: %s", asn, e) + self.errors[str(asn)] = classify_error(e) return None def run_collection(self): @@ -185,19 +218,28 @@ class CIDRCollector: self.session = None now = datetime.datetime.now() + ok = failed = total_added = total_removed = 0 with db.session() as conn, db.transaction(conn): # Конфиг перечитываем внутри транзакции: источник, удалённый во время сбора, не воскресает configured = {str(a) for a in load_full_config().get("asns", [])} - for str_asn, prefixes in fetched.items(): + for str_asn in (str(a) for a in self.asns): if str_asn not in configured: continue - added, removed = db.merge_source(conn, "asn", str_asn, prefixes, now, self.ttl_days) - logger.info("AS%s: +%d / -%d prefixes", str_asn, len(added), len(removed)) + if str_asn in fetched: + added, removed = db.merge_source(conn, "asn", str_asn, fetched[str_asn], now, self.ttl_days) + db.record_source_result(conn, "asn", str_asn, now) + logger.info("AS%s: +%d / -%d prefixes", str_asn, len(added), len(removed)) + ok, total_added, total_removed = ok + 1, total_added + len(added), total_removed + len(removed) + else: + db.record_source_result(conn, "asn", str_asn, now, self.errors.get(str_asn, "unknown")) + failed += 1 swept = db.sweep_unconfigured(conn, "asn", configured, now, self.ttl_days) if swept: logger.info("Expired %d prefixes of unconfigured ASNs", swept) + db.forget_unconfigured(conn, "asn", configured) db.prune_changes(conn, now, self.config.get("changes_retention_days", DEFAULT_CHANGES_RETENTION_DAYS)) - logger.info("CIDR data saved to %s", DB_FILE) + logger.info("ASN collection finished: %d sources, %d ok, %d failed, +%d/-%d prefixes", + ok + failed, ok, failed, total_added, total_removed) with db.session() as conn: db.settle_recreated(conn) @@ -209,6 +251,7 @@ class FQDNCollector: self.ttl_days = self.config.get("ttl_days", DEFAULT_TTL_DAYS) # По умолчанию в список попадают только глобальные адреса; true - для внутренних имён self.allow_non_global_ips = bool(self.config.get("allow_non_global_ips", False)) + self.errors = {} # имя -> категория сбоя последнего опроса (для source_status) def add_fqdn(self, fqdn): added, self.fqdns = add_to_config_list("fqdns", fqdn) @@ -229,12 +272,15 @@ class FQDNCollector: addresses = {result[4][0].split("%")[0] for result in results} except socket.gaierror as e: logger.error("Error resolving %s: %s", fqdn, e) + self.errors[fqdn] = "dns" return [] if self.allow_non_global_ips: return list(addresses) kept = {a for a in addresses if _is_global(a)} if kept != addresses: logger.warning("%s: ignoring non-global addresses: %s", fqdn, ", ".join(sorted(addresses - kept))) + if not kept: + self.errors[fqdn] = "no_global_addresses" return list(kept) def run_collection(self): @@ -250,18 +296,27 @@ class FQDNCollector: fetched[fqdn] = set(resolved_ips) now = datetime.datetime.now() + ok = failed = total_added = total_removed = 0 with db.session() as conn, db.transaction(conn): configured = set(load_full_config().get("fqdns", [])) - for fqdn, ips in fetched.items(): + for fqdn in self.fqdns: if fqdn not in configured: continue - added, removed = db.merge_source(conn, "fqdn", fqdn, ips, now, self.ttl_days) - logger.info("%s: +%d / -%d IPs", fqdn, len(added), len(removed)) + if fqdn in fetched: + added, removed = db.merge_source(conn, "fqdn", fqdn, fetched[fqdn], now, self.ttl_days) + db.record_source_result(conn, "fqdn", fqdn, now) + logger.info("%s: +%d / -%d IPs", fqdn, len(added), len(removed)) + ok, total_added, total_removed = ok + 1, total_added + len(added), total_removed + len(removed) + else: + db.record_source_result(conn, "fqdn", fqdn, now, self.errors.get(fqdn, "dns")) + failed += 1 swept = db.sweep_unconfigured(conn, "fqdn", configured, now, self.ttl_days) if swept: logger.info("Expired %d IPs of unconfigured FQDNs", swept) + db.forget_unconfigured(conn, "fqdn", configured) db.prune_changes(conn, now, self.config.get("changes_retention_days", DEFAULT_CHANGES_RETENTION_DAYS)) - logger.info("FQDN data saved to %s", DB_FILE) + logger.info("FQDN collection finished: %d sources, %d ok, %d failed, +%d/-%d IPs", + ok + failed, ok, failed, total_added, total_removed) with db.session() as conn: db.settle_recreated(conn) diff --git a/db.py b/db.py index 5e4541a..f3b2799 100644 --- a/db.py +++ b/db.py @@ -12,7 +12,7 @@ from storage import StorageError, file_lock, load_json, save_json_atomic logger = logging.getLogger(__name__) -SCHEMA_VERSION = 2 +SCHEMA_VERSION = 3 # После отката к копии номера журнала сдвигаются на это число: курсоры, выданные до отката, становятся недействительными RESTORE_JOURNAL_JUMP = 1_000_000 PRESENCE_BATCH = 500 # значений в одном запросе проверки наличия (лимит переменных SQLite - не менее 999) @@ -315,6 +315,16 @@ def read_transaction(conn): conn.execute("COMMIT") +# Результат последнего опроса каждого источника (см. record_source_result) +_SOURCE_STATUS_SCHEMA = ( + "CREATE TABLE source_status (" + " kind TEXT NOT NULL CHECK (kind IN ('asn', 'fqdn')), source TEXT NOT NULL," + " last_attempt TEXT NOT NULL, last_success TEXT," + " failures INTEGER NOT NULL DEFAULT 0, error_kind TEXT," + " PRIMARY KEY (kind, source))", +) + + def _ensure_schema(conn, cc): if conn.execute("PRAGMA user_version").fetchone()[0] >= SCHEMA_VERSION: return @@ -336,6 +346,9 @@ def _ensure_schema(conn, cc): conn.execute(statement) conn.execute("INSERT INTO meta VALUES ('horizon_id', '0')") conn.execute(f"INSERT INTO meta SELECT 'horizon_ts', {_TS}") + if version < 3: + for statement in _SOURCE_STATUS_SCHEMA: + conn.execute(statement) conn.execute(f"PRAGMA user_version = {SCHEMA_VERSION}") conn.execute("COMMIT") except BaseException: @@ -403,10 +416,54 @@ def sweep_unconfigured(conn, kind, configured, now, ttl_days): def purge_source(conn, kind, source): - """Удаляет все адреса источника. True, если что-то было.""" + """Удаляет все адреса источника и его состояние. True, если были адреса.""" + conn.execute("DELETE FROM source_status WHERE kind = ? AND source = ?", (kind, source)) return conn.execute("DELETE FROM addresses WHERE kind = ? AND source = ?", (kind, source)).rowcount > 0 +def record_source_result(conn, kind, source, now, error_kind=None): + """Результат опроса источника: успех (error_kind=None) сбрасывает счётчик сбоев, сбой увеличивает его. + + Вызывается внутри transaction(). error_kind - короткая категория сбоя (полный текст остаётся в логе). + """ + stamp = _iso(now) + if error_kind is None: + conn.execute( + "INSERT INTO source_status (kind, source, last_attempt, last_success, failures, error_kind)" + " VALUES (?, ?, ?, ?, 0, NULL) ON CONFLICT (kind, source) DO UPDATE SET" + " last_attempt = excluded.last_attempt, last_success = excluded.last_success," + " failures = 0, error_kind = NULL", (kind, source, stamp, stamp)) + else: + conn.execute( + "INSERT INTO source_status (kind, source, last_attempt, last_success, failures, error_kind)" + " VALUES (?, ?, ?, NULL, 1, ?) ON CONFLICT (kind, source) DO UPDATE SET" + " last_attempt = excluded.last_attempt, failures = failures + 1," + " error_kind = excluded.error_kind", (kind, source, stamp, error_kind)) + + +def forget_unconfigured(conn, kind, configured): + """Удаляет состояние источников, которых уже нет в конфигурации. Вызывается внутри transaction().""" + query, params = "DELETE FROM source_status WHERE kind = ?", [kind] + if configured: + query += f" AND source NOT IN ({','.join('?' * len(configured))})" + params += sorted(configured) + conn.execute(query, params) + + +def source_statuses(conn): + """Состояние источников: {(kind, source): {last_attempt, last_success, failures, error_kind}}.""" + return {(kind, source): {"last_attempt": attempt, "last_success": success, "failures": failures, + "error_kind": error_kind} + for kind, source, attempt, success, failures, error_kind in conn.execute( + "SELECT kind, source, last_attempt, last_success, failures, error_kind FROM source_status")} + + +def address_counts(conn): + """Число значений по источникам: {(kind, source): n}.""" + return {(kind, source): n for kind, source, n in conn.execute( + "SELECT kind, source, COUNT(*) FROM addresses GROUP BY kind, source")} + + def get_values(conn, kind=None): query, params = "SELECT DISTINCT value FROM addresses", () if kind: diff --git a/docs/plan-observability.md b/docs/plan-observability.md new file mode 100644 index 0000000..fa6a3d8 --- /dev/null +++ b/docs/plan-observability.md @@ -0,0 +1,47 @@ +# План: наблюдаемость (сбои источников, здоровье, логи) + +Источник: `docs/analysis-2026-09-21.md` (наблюдения 2 и 5, риск «сбои RIPE/DNS не видны»), пункт 2 рекомендуемого порядка. **Метрики Prometheus (`/metrics`) в этот план не входят** (решение пользователя); всё, что описано ниже, доступно через `/health`, `/sources`, лог и базу. + +## Проблема +Недоступный RIPEstat или DNS пишется только в лог: задание считается успешным, `/health` показывает `ok`, а данные источника не обновляются (до истечения TTL - 90 дней). Сбой одного из шести ASN невидим, пока адреса не начнут пропадать. В логе после каждого запуска стоит «data saved» даже без изменений, а итог запуска по источникам не сводится. + +## Дизайн + +### 1. Состояние источников (SQLite, схема версии 3) +- Таблица `source_status(kind, source, last_attempt, last_success, failures, error_kind)`, ключ `(kind, source)`. Хранится в базе, а не в `status.json`: переживает перезапуск демона, пишется в той же транзакции, что и данные, попадает в резервные копии. Миграция 2 -> 3 автоматическая (как 1 -> 2), данные не затрагиваются. +- Сборщик записывает результат по каждому источнику при каждом запуске: успех - `last_success = last_attempt = сейчас`, `failures = 0`, `error_kind = NULL`; сбой - `last_attempt = сейчас`, `failures + 1`, `error_kind`; `last_success` не меняется. +- `error_kind` - короткая категория без деталей: `timeout`, `connection`, `http_429`, `http_4xx`, `http_5xx`, `invalid_response` (нечитаемый ответ), `dns`, `no_global_addresses` (имя разрешилось, но все адреса отфильтрованы), `unknown`. Полный текст ошибки остаётся в логе; наружу он не отдаётся (как и пути в находке 6 ревью). +- Пустой, но корректный ответ RIPEstat считается успехом (число адресов видно в `/sources`). +- Строки состояния источников, которых уже нет в конфигурации, удаляются в очередном запуске сбора и при `purge`. + +### 2. `/health` +- В ответ добавляется блок `sources`: по типам (`asn`, `fqdn`) число источников (`total`) и список **только неисправных** (`failing`): `source`, `failures`, `last_success`, `error_kind`. Компактно и без лишних деталей. +- Источник считается неисправным при `failures >= source_failure_threshold` (ключ в `config.json`, по умолчанию 3): один сбой RIPEstat не «краснит» систему. При наличии неисправных источников `status` = `degraded`. +- Существующие поля и поведение (`collector_alive`, `jobs`, `db_recreated`, `last_restore`) не меняются. + +### 3. `GET /sources` (чтение открыто, как `/asns`) +Полная картина по всем источникам: `kind`, `source`, `addresses` (сколько значений в базе), `last_attempt`, `last_success`, `failures`, `error_kind`. Позволяет увидеть источник, который «успешно» вернул ноль адресов или давно не обновлялся. Существующие `/asns` и `/fqdns` не меняются. + +### 4. Логи +- Итог запуска одной строкой: «ASN collection finished: 6 sources, 5 ok, 1 failed, +12/-3 prefixes» (для FQDN так же). +- Сообщение «... data saved to ...» после каждой транзакции убирается (наблюдение 5 анализа); подробности по источникам остаются на уровне INFO. + +## Изменения +1. `db.py`: схема v3 и миграция, `record_source_result`, `source_statuses`, очистка состояния (`sweep`, `purge`). +2. `cidr_collector.py`: классификация ошибок (`error_kind`), запись результата по каждому источнику в обоих `run_collection`, итоговая строка лога, ключ `source_failure_threshold`. +3. `api_server.py`: блок `sources` и статус в `/health`, `GET /sources`. +4. `README.md`: `/health`, `/sources`, ключ конфигурации, лог, схема базы. +5. Тесты (2 новых, всего 30): сборщик (успех и сбой источников записываются, категории ошибок, сброс `failures` при успехе, очистка удалённого источника); API (`/health` `ok` до порога и `degraded` после, ответ без текста ошибок, `/sources` с числом адресов). + +## Не входит +- Метрики Prometheus и `/metrics` (исключены по решению пользователя). +- Уведомления (webhook, Telegram) и алерты: `/health` и `/sources` пригодны для внешних проверок, сами уведомления - отдельная доработка. +- Свежесть резервных копий как отдельный признак здоровья, команда CLI `status`, ограничение частоты `POST /collect`. +- Проверка `docker healthcheck` не меняется: `degraded` по-прежнему HTTP 200, контейнер из-за внешнего сбоя перезапускать нечего. + +## Проверка +Тесты в контейнере. Вручную: миграция базы версии 2 (созданной кодом до изменения) с сохранением данных; локальный «RIPEstat» (как при проверке повторов) отвечает 503: после трёх запусков `/health` показывает `degraded` и источник в `failing`, после восстановления сервера `ok`; сравнение ответов существующих эндпоинтов до и после; Docker Compose (демон и API, ручной сбор). + +## Риски и откат +- Схема базы v3: таблица добавляется, существующие таблицы не меняются; откат кода оставляет лишнюю таблицу без последствий (старый код её не читает, `user_version` 3 у старого кода вызовет только пропуск миграции). +- Порог 3 на редких расписаниях (сбор FQDN раз в сутки) означает три дня до `degraded`; порог настраивается. diff --git a/docs/summary-observability.md b/docs/summary-observability.md new file mode 100644 index 0000000..47a6baa --- /dev/null +++ b/docs/summary-observability.md @@ -0,0 +1,24 @@ +# Итоги: наблюдаемость (сбои источников, здоровье, логи) + +План: `docs/plan-observability.md`. Закрывает наблюдения 2 и 5 анализа (`analysis-2026-09-21.md`). **Метрики Prometheus не входили** (решение пользователя). + +## Сделано +- **Состояние источников (`db.py`):** схема версии 3, таблица `source_status(kind, source, last_attempt, last_success, failures, error_kind)`; автоматическая миграция 2 -> 3. Функции `record_source_result`, `forget_unconfigured`, `source_statuses`, `address_counts`; `purge_source` удаляет и состояние. Запись идёт в той же транзакции, что и данные. +- **Сборщики (`cidr_collector.py`):** результат по каждому источнику при каждом запуске (успех обнуляет `failures`, сбой увеличивает, `last_success` при сбое не меняется); `classify_error` даёт категорию без деталей (`timeout`, `connection`, `http_429/4xx/5xx`, `invalid_response`, `dns`, `no_global_addresses`, `unknown`); состояние удалённых из конфигурации источников очищается при очередном сборе; порог `failure_threshold` (`source_failure_threshold`, по умолчанию 3). +- **`/health` (`api_server.py`):** блок `sources` (по типам: `total` и только неисправные), `degraded` при неисправных источниках; остальные поля не менялись. **`GET /sources`:** все настроенные источники с числом адресов, временем попытки и успеха, сбоями и категорией. +- **Логи:** итоговая строка запуска («ASN collection finished: 6 sources, 5 ok, 1 failed, +12/-3 prefixes», для FQDN так же); сообщения «data saved» убраны (наблюдение 5). +- **README:** `/health`, `/sources`, ключ `source_failure_threshold`, схема базы, лог. +- **Тесты:** 2 новых (сборщики: запись успеха и сбоев, рост и сброс счётчика, очистка удалённого источника, категории ошибок и миграция с версии 2; API: `ok` до порога и `degraded` после, ответ без текста ошибки, `/sources`, настройка порога). Всего 30, в контейнере 30 passed. + +## Проверка +- **Миграция:** база версии 2, созданная кодом до изменения (`SCHEMA_VERSION = 2` в предыдущем коммите), открыта новым кодом: `user_version` 3, все адреса и журнал на месте, состояние источников пусто. +- **Отдельные процессы:** демон с подменённым адресом RIPEstat (локальный сервер отвечает 503) и API: после первых двух запусков `/health` `ok`, после третьего `degraded` с `{"source": "62041", "failures": 3, "error_kind": "http_5xx"}`; `/sources` показывает 0 адресов и 3 сбоя; после «починки» сервера `ok`, 1 адрес, `failures: 0`. В логе итоговые строки запусков, сообщений «data saved» нет. +- **Docker Compose** (отдельный проект, порт 18000, стенд убран): FQDN `example.com` и `nonexistent.invalid`; после трёх ручных сборов `degraded`, неисправный `nonexistent.invalid` с `error_kind: dns`, `example.com` с 4 адресами; оба сервиса `healthy`. +- В первом запуске проверки миграции сценарий упал на моей опечатке (обращение к закрытой сессии после записи), сама миграция прошла; версию «до» подтвердил по коду предыдущего коммита. + +## Замечания +- Порог 3 при суточном расписании FQDN означает три дня до `degraded`; порог настраивается (`source_failure_threshold`). +- Пустой корректный ответ RIPEstat считается успехом: такой источник виден в `/sources` по `addresses: 0`, но `degraded` не вызывает. +- Категория `no_global_addresses` считается сбоем: имя, разрешившееся только в частные адреса при выключенном `allow_non_global_ips`, через три запуска даст `degraded`. +- `/health` по-прежнему отвечает HTTP 200 при `degraded`; Docker healthcheck из-за внешних сбоев контейнеры не перезапускает. +- Не вошло: Prometheus, уведомления и алерты, свежесть резервных копий как признак здоровья, команда CLI `status`. diff --git a/tests/test_core.py b/tests/test_core.py index c545bc3..69eb9b7 100644 --- a/tests/test_core.py +++ b/tests/test_core.py @@ -140,3 +140,71 @@ def test_ripestat_session_retries_and_sourceapp(files, monkeypatch): class Response: # сервер просит подождать час: пауза ограничена потолком headers = {"Retry-After": "3600"} assert retry.get_retry_after(Response()) == cc.MAX_RETRY_AFTER + + +def test_source_status_is_recorded(files, monkeypatch): + (files / "config.json").write_text('{"asns": [1, 2], "fqdns": ["ok.example", "dead.example", "private.example"]}') + + def statuses(): + with db.session() as conn: + return db.source_statuses(conn) + + # ASN 1 отвечает, ASN 2 сбоит трижды подряд: счётчик растёт, last_success не появляется + def fetch_all(collector, broken): + def fetch(asn): + if asn in broken: + collector.errors[str(asn)] = "http_5xx" + return None + return ["1.0.0.0/24"] + return fetch + + for _ in range(3): + collector = cc.CIDRCollector() + monkeypatch.setattr(collector, "fetch_prefixes", fetch_all(collector, {2})) + collector.run_collection() + state = statuses() + assert state["asn", "1"]["failures"] == 0 and state["asn", "1"]["last_success"] + assert state["asn", "2"]["failures"] == 3 and state["asn", "2"]["error_kind"] == "http_5xx" + assert state["asn", "2"]["last_success"] is None + + # Успех сбрасывает счётчик; источник, удалённый из конфигурации, теряет состояние + collector = cc.CIDRCollector() + monkeypatch.setattr(collector, "fetch_prefixes", fetch_all(collector, set())) + collector.run_collection() + assert statuses()["asn", "2"]["failures"] == 0 and statuses()["asn", "2"]["error_kind"] is None + (files / "config.json").write_text('{"asns": [1], "fqdns": ["ok.example", "dead.example", "private.example"]}') + collector = cc.CIDRCollector() + monkeypatch.setattr(collector, "fetch_prefixes", lambda asn: ["1.0.0.0/24"]) + collector.run_collection() + assert ("asn", "2") not in statuses() + + # FQDN: DNS-сбой и «только неглобальные адреса» - разные категории + def getaddrinfo(name, *args, **kwargs): + if name == "dead.example": + raise cc.socket.gaierror("no such host") + address = "93.184.216.34" if name == "ok.example" else "10.0.0.5" + return [(2, 1, 6, "", (address, 0))] + + monkeypatch.setattr(cc.socket, "getaddrinfo", getaddrinfo) + cc.FQDNCollector().run_collection() + state = statuses() + assert state["fqdn", "ok.example"]["failures"] == 0 + assert state["fqdn", "dead.example"]["error_kind"] == "dns" + assert state["fqdn", "private.example"]["error_kind"] == "no_global_addresses" + + # Категории ошибок (текст наружу не отдаётся) + requests = cc.requests.exceptions + assert cc.classify_error(requests.ConnectTimeout("x")) == "timeout" + assert cc.classify_error(requests.ConnectionError("Read timed out")) == "timeout" + assert cc.classify_error(requests.ConnectionError("refused")) == "connection" + assert cc.classify_error(requests.RetryError("too many 429 error responses")) == "http_429" + assert cc.classify_error(requests.HTTPError(response=type("R", (), {"status_code": 503})())) == "http_5xx" + assert cc.classify_error(requests.HTTPError(response=type("R", (), {"status_code": 404})())) == "http_4xx" + assert cc.classify_error(ValueError("bad json")) == "invalid_response" and cc.classify_error(KeyError()) == "unknown" + + # База версии 2 (без таблицы состояния) обновляется при открытии, данные остаются + with db.session() as conn: + conn.execute("DROP TABLE source_status") + conn.execute("PRAGMA user_version = 2") + with db.session() as conn: + assert db.source_statuses(conn) == {} and db.get_values(conn, "asn") == ["1.0.0.0/24"] diff --git a/tests/test_sources_api.py b/tests/test_sources_api.py index d15c4f6..e9e6b43 100644 --- a/tests/test_sources_api.py +++ b/tests/test_sources_api.py @@ -145,3 +145,47 @@ def test_health_hides_internal_paths(env): def json_dump_has_no_paths(text): return "/data" not in text and "corrupt" not in text and "not a database" not in text + + +def test_health_reports_failing_sources(env): + cc.add_to_config_list("asns", 62041) + cc.add_to_config_list("fqdns", "example.com") + save_json_atomic(cc.STATUS_FILE, {"updated_at": datetime.datetime.now().isoformat(), "jobs": {}}) # демон "жив" + now = datetime.datetime.now() + + def fail(times): + for _ in range(times): + with db.session() as conn, db.transaction(conn): + db.record_source_result(conn, "asn", "62041", now, "http_5xx") + + # До порога (3 сбоя подряд) источник не считается неисправным + fail(2) + health = env.get("/health").json() + assert health["status"] == "ok" and health["sources"] == {"asn": {"total": 1, "failing": []}, + "fqdn": {"total": 1, "failing": []}} + + # На пороге: degraded, источник в списке, наружу - только категория ошибки + fail(1) + health = env.get("/health").json() + assert health["status"] == "degraded" + assert health["sources"]["asn"]["failing"] == [{"source": "62041", "failures": 3, "last_success": None, + "error_kind": "http_5xx"}] + + # Полная картина: адреса источника и состояние; неопрошенный источник - пустое состояние + with db.session() as conn, db.transaction(conn): + db.merge_source(conn, "asn", "62041", {"1.0.0.0/24", "2.0.0.0/24"}, now, 90) + report = env.get("/sources").json() + assert report["failure_threshold"] == 3 + asn, fqdn = report["sources"] + assert (asn["kind"], asn["source"], asn["addresses"], asn["failures"]) == ("asn", "62041", 2, 3) + assert (fqdn["source"], fqdn["addresses"], fqdn["last_attempt"], fqdn["failures"]) == ("example.com", 0, None, 0) + + # Успех сбрасывает счётчик, /health снова ok + with db.session() as conn, db.transaction(conn): + db.record_source_result(conn, "asn", "62041", now) + assert env.get("/health").json()["status"] == "ok" + + # Порог настраивается ключом source_failure_threshold + cc.update_config(lambda config: config.update(source_failure_threshold=1)) + fail(1) + assert env.get("/health").json()["status"] == "degraded"