diff --git a/bin/SHA256SUMS b/bin/SHA256SUMS index ce77107..2330100 100644 --- a/bin/SHA256SUMS +++ b/bin/SHA256SUMS @@ -1,4 +1,4 @@ -dcb137ae6a3f6a59465f2772215d48f98fd303dd1c266645b20ec70194ab7915 control-api -b7a6068db1d095ae7b73cbd9b7273d6a1629ae017e9c3c9601e103432894a247 validator-agent -be8e9576c6fbe5e5798b5d4b758a5617b63bea6abe80e326998d10583edbbf5c prober -61d1df7aa84528e8fd9b5e700ac6a9b77e45763b5c5d938913dbce25d0ad9940 admin-dashboard +62fb7e4e6a05f60e0e2f1d113c11b4ef46840eadb87202c3539a508ea0e996ae control-api +7c11ffe00f09cbfc33bab57d8a889cdfedacc69f0049a5881705dfabc154fa76 validator-agent +543f7df07939d576b79dc8fdc655ee6ff5577c834d5606a1e99687b99d2a8927 prober +2e9798e86ebaf92e70622cbd76ca0d3de8f90958eadbe7c71fb2118940b4a275 admin-dashboard diff --git a/bin/admin-dashboard b/bin/admin-dashboard index 84f43f2..e125397 100755 Binary files a/bin/admin-dashboard and b/bin/admin-dashboard differ diff --git a/bin/control-api b/bin/control-api index ed06e42..ca32640 100755 Binary files a/bin/control-api and b/bin/control-api differ diff --git a/bin/prober b/bin/prober index 5293591..8548125 100755 Binary files a/bin/prober and b/bin/prober differ diff --git a/bin/validator-agent b/bin/validator-agent index e73751d..8067bff 100755 Binary files a/bin/validator-agent and b/bin/validator-agent differ diff --git a/cmd/control-api/main.go b/cmd/control-api/main.go index ac6a3f2..283dab7 100644 --- a/cmd/control-api/main.go +++ b/cmd/control-api/main.go @@ -92,6 +92,17 @@ func runOrchestratorLoop(ctx context.Context, orch *orchestrator.Orchestrator, c heartbeatTicker := time.NewTicker(interval * 2) defer heartbeatTicker.Stop() + // Periodic floating-IP scanning is optional: a zero interval leaves + // scanTickerC nil, and a select on a nil channel simply never fires, so + // the loop falls back to manual-only scanning (POST + // /api/v1/admin/ips/scan) without a special-cased branch below. + var scanTickerC <-chan time.Time + if cfg.Orchestrator.FIPScanIntervalSeconds > 0 { + scanTicker := time.NewTicker(time.Duration(cfg.Orchestrator.FIPScanIntervalSeconds) * time.Second) + defer scanTicker.Stop() + scanTickerC = scanTicker.C + } + for { select { case <-ctx.Done(): @@ -105,6 +116,10 @@ func runOrchestratorLoop(ctx context.Context, orch *orchestrator.Orchestrator, c if err := orch.SweepStaleSiteHeartbeats(ctx); err != nil { log.Error("sweep stale site heartbeats", "err", err) } + case <-scanTickerC: + if _, _, err := orch.ScanFloatingIPs(ctx); err != nil { + log.Error("scan floating ips", "err", err) + } } } } diff --git a/configs/control-api.example.yaml b/configs/control-api.example.yaml index d6c6669..f9f2a1f 100644 --- a/configs/control-api.example.yaml +++ b/configs/control-api.example.yaml @@ -54,6 +54,12 @@ orchestrator: # /settings page) and this field is ignored. Must satisfy # fip_settle_seconds + self_check_timeout_seconds < lease_ttl_seconds. fip_settle_seconds: 0 + # How often (seconds) to automatically scan the OpenStack project for free + # (unassociated) floating IPs and submit them to the check queue. 0 (the + # default) disables periodic scanning — an operator can still trigger a + # scan on demand via POST /api/v1/admin/ips/scan or the dashboard's + # "Scan Floating IPs" button. + fip_scan_interval_seconds: 0 aggregation: missing_counts_as_fail: true diff --git a/deploy/docker/control-api/control-api.docker.example.yaml b/deploy/docker/control-api/control-api.docker.example.yaml index 9f7bf4f..663ee81 100644 --- a/deploy/docker/control-api/control-api.docker.example.yaml +++ b/deploy/docker/control-api/control-api.docker.example.yaml @@ -26,6 +26,7 @@ orchestrator: lease_ttl_seconds: 180 heartbeat_timeout_seconds: 30 fip_settle_seconds: 0 + fip_scan_interval_seconds: 0 aggregation: missing_counts_as_fail: true diff --git a/docs/API.md b/docs/API.md index 9adc8f3..84fc6f8 100644 --- a/docs/API.md +++ b/docs/API.md @@ -28,6 +28,7 @@ JSON, базовый префикс прикладных методов — `/ap - [Методы для prober](#методы-для-prober) - [Служебные и административные методы](#служебные-и-административные-методы) - [Управление очередью и конфигурацией](#управление-очередью-и-конфигурацией) +- [Реестр адресов и история проверок](#реестр-адресов-и-история-проверок) - [Модель состояний и связь методов с ней](#модель-состояний-и-связь-методов-с-ней) - [Сквозной пример работы (curl)](#сквозной-пример-работы-curl) @@ -416,12 +417,18 @@ YAML для этой секции больше не перечитывается ### `DELETE /api/v1/admin/ips/{ip}`, `POST /api/v1/admin/ips/delete`, `POST /api/v1/admin/ips/clear` -**Безвозвратное удаление**, в отличие от `cancel` выше: строка `ip_queue` -и вся её история (`checks`, `events`) стираются физически, без возможности -восстановления. Работает из любого состояния, включая активно -проверяемое — если Floating IP привязан, он отвязывается тем же -best-effort способом, что и при `cancel`/обычном завершении, владеющий -валидатор освобождается. +Удаляет строку `ip_queue` — адрес пропадает из очереди/`GET +/api/v1/admin/ips*` — безвозвратно, без возможности восстановить именно +эту строку. Работает из любого состояния, включая активно проверяемое — +если Floating IP привязан, он отвязывается тем же best-effort способом, +что и при `cancel`/обычном завершении, владеющий валидатор освобождается. + +**Накопленная история адреса при этом не теряется**: `checks`/`events` +остаются в реестре (`ip_registry`, см. раздел [«Реестр +адресов»](#реестр-адресов-и-история-проверок) ниже) и доступны через `GET +/api/v1/admin/registry/{ip}` даже после удаления строки из очереди — в +отличие от `ip_queue`, реестровая запись никогда не удаляется этими +методами. | Метод | Путь | Тело | Успех | Ошибки | |---|---|---|---|---| @@ -433,7 +440,95 @@ best-effort способом, что и при `cancel`/обычном заве идут в `not_found`, не ошибка — тот же терпимый стиль, что у `POST /api/v1/admin/ips`). `POST .../clear` удаляет **вообще всё**, что сейчас в очереди, включая адреса в процессе проверки — самая опасная операция -этого API, используйте с осторожностью. +этого API, используйте с осторожностью (и, опять же, ничья история при +этом физически не стирается — см. выше). + +### `POST /api/v1/admin/ips/scan` + +Сканирует текущий проект OpenStack на предмет свободных (не привязанных ни +к одному порту) Floating IP и сразу передаёт найденный список в `POST +/api/v1/admin/ips` — тот же add/requeue/reorder-вызов, как если бы +оператор ввёл эти адреса вручную. Не принимает тело запроса. + +Ответ (`200`): +```json +{ + "scanned_free": 3, + "added": ["203.0.113.20"], + "requeued": [], + "reordered": ["203.0.113.10", "203.0.113.11"], + "skipped_in_progress": [] +} +``` + +`scanned_free` — сколько свободных Floating IP нашлось в проекте всего +(включая уже стоящие в очереди — они попадут в `reordered`, а не +`added`). Если свободных адресов нет вообще, это не ошибка: ответ будет +`{"scanned_free": 0, "added": [], ...}`. + +Помимо ручного вызова, сканирование можно включить по расписанию — +`orchestrator.fip_scan_interval_seconds` в `control-api.yaml` (0, по +умолчанию, — только по запросу через эту ручку или кнопку «Сканировать +Floating IP» в дашборде). + +## Реестр адресов и история проверок + +В отличие от `ip_queue` (текущая рабочая очередь, см. выше), реестр — +`ip_registry` — это накопительная запись **обо всех адресах, когда-либо +поставленных на проверку**, вне зависимости от того, стоят ли они сейчас в +очереди. Запись в реестре переживает удаление адреса из `ip_queue` (`DELETE +/api/v1/admin/ips/{ip}` и т.п.) и повторное добавление того же адреса +позже — обе истории (до и после) остаются доступны и не перекрывают друг +друга (каждой постановке на проверку соответствует свой `cycle_id`, +уникальный в пределах адреса на всё время). + +### `GET /api/v1/admin/registry` + +Список всех адресов реестра с краткой сводкой по каждому. + +```json +[ + { + "ip_address": "203.0.113.10", + "first_seen_at": "2026-01-10T12:00:00Z", + "last_seen_at": "2026-02-01T09:00:00Z", + "total_cycles": 4, + "last_result": "pass", + "last_checked_at": "2026-02-01T09:05:00Z", + "in_queue": true, + "current_state": "done" + } +] +``` + +`in_queue`/`current_state` отражают, есть ли у адреса сейчас живая строка в +`ip_queue`, а не только в реестре. + +### `GET /api/v1/admin/registry/{ip}` + +Реестровая запись по одному адресу плюс вся сохранённая история проверок +по нему, по всем циклам (не только текущему — в отличие от `GET +/api/v1/admin/ips/{ip}`, который отдаёт проверки только текущей попытки). +Порядок — от новых циклов к старым. + +```json +{ + "registry": { "ip_address": "203.0.113.10", "total_cycles": 4, "...": "..." }, + "checks": [ + {"CycleID": 4, "Source": "egress", "CheckType": "https", "Success": true, "...": "..."}, + {"CycleID": 3, "Source": "egress", "CheckType": "https", "Success": false, "...": "..."} + ] +} +``` + +`404`, если адрес никогда не ставился на проверку. + +**Глубина хранения.** Сколько последних циклов на адрес хранится в +`checks` (и синхронно — в `events`), управляется полем +`history_retention_cycles` в `GET`/`PUT /api/v1/admin/config/orchestrator` +(0, по умолчанию, — без ограничения). Сама реестровая запись (`ip_address`, +`first_seen_at`, счётчик циклов) не удаляется никогда, независимо от этой +настройки — она лишь ограничивает глубину детальной истории проверок. ```bash curl -s -X DELETE "$BASE/api/v1/admin/ips/203.0.113.10" @@ -497,25 +592,33 @@ curl -s -X POST "$BASE/api/v1/admin/ips/clear" ### Настройки оркестратора: `/api/v1/admin/config/orchestrator` -Единственный на сегодня параметр — `fip_settle_seconds`: пауза между -привязкой Floating IP к валидатору и моментом, когда self-check по этому -адресу становится доступен агенту (`GET -/api/v1/agents/{id}/assignment` до истечения паузы отдаёт `204`, как если -бы валидатору просто нечего было делать — никаких изменений в протоколе -агента). Нужна, чтобы дать data plane OpenStack время реально начать -пропускать трафик через только что привязанный адрес, прежде чем -запускать по нему проверки. `0` — без паузы (поведение по умолчанию, как -до появления этого параметра). +Два параметра: + +- `fip_settle_seconds` — пауза между привязкой Floating IP к валидатору и + моментом, когда self-check по этому адресу становится доступен агенту + (`GET /api/v1/agents/{id}/assignment` до истечения паузы отдаёт `204`, + как если бы валидатору просто нечего было делать — никаких изменений в + протоколе агента). Нужна, чтобы дать data plane OpenStack время реально + начать пропускать трафик через только что привязанный адрес, прежде чем + запускать по нему проверки. `0` — без паузы (поведение по умолчанию, как + до появления этого параметра). +- `history_retention_cycles` — сколько последних циклов проверки хранить + на адрес в реестре (`GET /api/v1/admin/registry/{ip}`, см. + [«Реестр адресов»](#реестр-адресов-и-история-проверок)). `0` — без + ограничения (поведение по умолчанию). | Метод | Путь | Тело | Успех | Ошибки | |---|---|---|---|---| -| GET | `/api/v1/admin/config/orchestrator` | — | `{"fip_settle_seconds":N}` | | -| PUT | `/api/v1/admin/config/orchestrator` | `{"fip_settle_seconds":N}` | `200` | `400`, если `N < 0`, или если `fip_settle_seconds + self_check_timeout_seconds >= lease_ttl_seconds` (пауза не должна съедать весь лизинг адреса — иначе self-check не успеет пройти до истечения `lease_ttl_seconds`, и адрес будет вечно возвращаться в очередь) | +| GET | `/api/v1/admin/config/orchestrator` | — | `{"fip_settle_seconds":N,"history_retention_cycles":M}` | | +| PUT | `/api/v1/admin/config/orchestrator` | `{"fip_settle_seconds":N,"history_retention_cycles":M}` | `200` | `400`, если `N < 0` или `M < 0`, или если `fip_settle_seconds + self_check_timeout_seconds >= lease_ttl_seconds` (пауза не должна съедать весь лизинг адреса — иначе self-check не успеет пройти до истечения `lease_ttl_seconds`, и адрес будет вечно возвращаться в очередь) | Как и остальные разделы этой группы, YAML-поле `orchestrator. fip_settle_seconds` в `control-api.yaml` — только одноразовый bootstrap для пустой БД; дальше источник истины — сама база, менять значение нужно через `PUT` выше (или страницу `/settings` в дашборде). +`history_retention_cycles` не имеет YAML-эквивалента вообще — управляется +только через `PUT` выше/дашборд, значение по умолчанию `0` всегда +применяется на пустой БД. ### Типы проверок пробера: `/api/v1/admin/config/inbound-checks` @@ -633,12 +736,15 @@ queued ──(control-api сам, без вызова API)──▶ assigning_fi "cancelled"`)** — `POST /api/v1/admin/ips/{ip}/cancel`; - **`done`/`failed`/`occupied` → `queued` (новая попытка)** — `POST /api/v1/admin/ips` с уже завершённым (или занятым) адресом в списке; -- **любое состояние → адрес физически исчезает из очереди**, вместе со - всей историей — `DELETE /api/v1/admin/ips/{ip}`, `POST - /api/v1/admin/ips/delete`, `POST /api/v1/admin/ips/clear` (см. +- **любое состояние → адрес физически исчезает из очереди** — `DELETE + /api/v1/admin/ips/{ip}`, `POST /api/v1/admin/ips/delete`, `POST + /api/v1/admin/ips/clear` (см. [выше](#delete-apiv1adminipsip-post-apiv1adminipsdelete-post-apiv1adminipsclear)). Не путать с cancel — cancel сохраняет запись как историю (`failed`/ - `cancelled`), delete стирает её целиком без возможности восстановления. + `cancelled`) прямо в `ip_queue`; delete убирает саму строку `ip_queue` + безвозвратно, но накопленная история проверок остаётся в реестре (`GET + /api/v1/admin/registry/{ip}`) — см. + [«Реестр адресов»](#реестр-адресов-и-история-проверок). ## Сквозной пример работы (curl) diff --git a/docs/DASHBOARD.md b/docs/DASHBOARD.md index 8dc5e71..f41f203 100644 --- a/docs/DASHBOARD.md +++ b/docs/DASHBOARD.md @@ -49,13 +49,15 @@ admin-dashboard -config /etc/cloud-ip-validator/admin-dashboard.yaml | Страница | Назначение | |---|---| | `/overview` | Сводная статистика: счётчики по состояниям, «текущая проверка» (live-снимок всех IP не в терминальном состоянии) и «последние N завершённых» (по умолчанию 20, `overview.last_completed_count`) с разбивкой pass/partial/fail/cancelled. Обновляется каждые `overview.poll_interval_seconds` секунд без перезагрузки страницы. | -| `/ips` | Полная очередь. Форма сверху принимает список адресов (по одному на строке или через запятую) и отправляет их в `POST /api/v1/admin/ips` — **один и тот же вызов** добавляет новые адреса и принудительно перезапускает уже завершённые (см. ниже). У каждого адреса — кнопка «Перепроверить» (для `done`/`failed`) или «Отменить» (для активных состояний), и всегда — «Удалить» (безвозвратно, в отличие от «Отменить», см. ниже). Чекбоксы у строк + кнопка «Удалить выбранные» удаляют список одним вызовом; «Очистить всё» удаляет вообще всё, включая активные проверки — обе операции требуют явного подтверждения. Пока не истекла настроенная на `/settings` пауза (`fip_settle_seconds`), только что привязавший Floating IP адрес показывает отдельный бейдж «прогрев FIP» вместо обычного статуса. Если на момент попытки привязки Floating IP оказался уже занят другим портом (дрейф состояния облака или ошибочно переданный адрес), цикл проверки для него не запускается — адрес показывает отдельный бейдж «занят» (отличный от «fail») и строку `fip_occupied` в списке событий на его странице; кнопка «Перепроверить» ставит его в очередь заново. | -| `/ips/{ip}` | Детали одного адреса: все проверки текущей попытки и вся история событий. | +| `/ips` | Полная очередь. Форма сверху принимает список адресов (по одному на строке или через запятую) и отправляет их в `POST /api/v1/admin/ips` — **один и тот же вызов** добавляет новые адреса и принудительно перезапускает уже завершённые (см. ниже). Кнопка «Сканировать Floating IP» делает то же самое автоматически: находит в проекте OpenStack все свободные (не привязанные к порту) Floating IP и сразу ставит их в очередь (`POST /api/v1/admin/ips/scan`, см. [API.md](API.md#post-apiv1adminipsscan)) — то же сканирование можно включить по расписанию через `orchestrator.fip_scan_interval_seconds`. У каждого адреса — кнопка «Перепроверить» (для `done`/`failed`) или «Отменить» (для активных состояний), и всегда — «Удалить» (безвозвратно убирает адрес из очереди, но не из реестра — см. ниже). Чекбоксы у строк + кнопка «Удалить выбранные» удаляют список одним вызовом; «Очистить всё» удаляет вообще всё, включая активные проверки — обе операции требуют явного подтверждения. Пока не истекла настроенная на `/settings` пауза (`fip_settle_seconds`), только что привязавший Floating IP адрес показывает отдельный бейдж «прогрев FIP» вместо обычного статуса. Если на момент попытки привязки Floating IP оказался уже занят другим портом (дрейф состояния облака или ошибочно переданный адрес), цикл проверки для него не запускается — адрес показывает отдельный бейдж «занят» (отличный от «fail») и строку `fip_occupied` в списке событий на его странице; кнопка «Перепроверить» ставит его в очередь заново. | +| `/ips/{ip}` | Детали одного адреса, пока он в очереди: все проверки текущей попытки и вся история событий, плюс ссылка на полную историю в реестре (см. ниже). | +| `/registry` | **Реестр** — все адреса, когда-либо поставленные на проверку, независимо от того, стоят ли они сейчас в очереди. Переживает удаление адреса из `/ips` и повторное добавление того же адреса позже (см. «Реестр адресов» ниже). | +| `/registry/{ip}` | Полная сохранённая история проверок одного адреса по всем циклам (не только текущему) — в отличие от `/ips/{ip}`, которая показывает только текущую попытку. | | `/validators` | Список валидаторов + создание/изменение `os_port_id`/удаление. | | `/sites` | Площадки — число слотов не ограничено, форма сверху добавляет новый слот, назначить/сменить/освободить `site_id` в каждой строке; колонка «Статус» показывает бейдж подключения пробера (`unregistered`/`idle`/`unreachable`, по аналогии с `/validators`), см. [USAGE.md](USAGE.md#состояния-площадки). | | `/targets` | Группы целей для egress-проверок — создание/редактирование/удаление. | | `/check-types` | Типы проверок (`https`/`icmp`/`ssh`/...), включение/выключение, привязка к группам целей. | -| `/settings` | Две формы: `fip_settle_seconds` — пауза (в секундах) между привязкой Floating IP и началом self-check («прогрев» дата-плейна OpenStack, см. [USAGE.md](USAGE.md#пауза-перед-self-check-fip_settle_seconds)); и типы проверок пробера — TCP-порты (через запятую) + чекбокс ICMP, общие для всех площадок (см. [USAGE.md](USAGE.md#управление-типами-проверок-пробера)). | +| `/settings` | Три формы: `fip_settle_seconds` — пауза (в секундах) между привязкой Floating IP и началом self-check («прогрев» дата-плейна OpenStack, см. [USAGE.md](USAGE.md#пауза-перед-self-check-fip_settle_seconds)); `history_retention_cycles` — сколько последних циклов проверки хранить на адрес в реестре (0 — без ограничения); и типы проверок пробера — TCP-порты (через запятую) + чекбокс ICMP, общие для всех площадок (см. [USAGE.md](USAGE.md#управление-типами-проверок-пробера)). | ### «Текущая» и «последняя завершённая» проверка @@ -83,20 +85,38 @@ admin-dashboard -config /etc/cloud-ip-validator/admin-dashboard.yaml (дашборд честно показывает это в таблице, а не делает вид, что запрос ничего не значил). -### Удаление адресов — безвозвратно, в отличие от «Отменить» +### Удаление адресов — безвозвратно из очереди, но не из реестра «Отменить» (`POST .../cancel`) останавливает проверку, но сохраняет адрес и его историю как `failed`/`cancelled` — он остаётся виден в очереди. «Удалить» (кнопка в строке, «Удалить выбранные» по чекбоксам, -«Очистить всё») стирает строку и всю её историю проверок/событий -физически, без возможности восстановления — работает из любого -состояния, включая активно проверяемое (Floating IP отвязывается, +«Очистить всё») убирает строку из `/ips` безвозвратно — работает из +любого состояния, включая активно проверяемое (Floating IP отвязывается, валидатор освобождается). Все три операции удаления в UI защищены `hx-confirm` с формулировкой, отражающей необратимость — «Очистить всё» -предупреждает отдельно, так как затрагивает и активные проверки. Подробнее -— [API.md](API.md#delete-apiv1adminipsip-post-apiv1adminipsdelete-post-apiv1adminipsclear) +предупреждает отдельно, так как затрагивает и активные проверки. + +**Накопленная история при этом не теряется** — она остаётся в +[реестре](#реестр-адресов) (`/registry/{ip}`) даже после того, как адрес +пропал из `/ips`, и продолжает пополняться, если адрес позже добавят +заново. Подробнее — +[API.md](API.md#delete-apiv1adminipsip-post-apiv1adminipsdelete-post-apiv1adminipsclear) и [USAGE.md](USAGE.md#удаление-адресов-из-очереди). +### Реестр адресов + +`/registry` решает задачу, которую `/ips` принципиально не может: история +проверок конкретного адреса не должна теряться только из-за того, что его +временно вывели из очереди (например, адрес переиспользуют для другой +цели) и позже добавили обратно — возможно, под другим циклом проверки. +Каждая запись реестра живёт всё время, что адрес когда-либо существовал в +системе, и не удаляется вместе со строкой `ip_queue`. Единственное, что +можно ограничить — глубину детальной истории проверок на один адрес +(`history_retention_cycles` на `/settings`, по циклам, а не по времени); +сама запись в реестре (когда адрес впервые встречен, сколько всего было +циклов) остаётся всегда. Подробнее — +[API.md](API.md#реестр-адресов-и-история-проверок). + ## Конфигурация См. `configs/admin-dashboard.example.yaml`. Ключевые поля: diff --git a/docs/USAGE.md b/docs/USAGE.md index 82e0ebc..65865d7 100644 --- a/docs/USAGE.md +++ b/docs/USAGE.md @@ -11,10 +11,12 @@ - [Как устроена работа с системой](#как-устроена-работа-с-системой) - [Добавление новых IP в очередь](#добавление-новых-ip-в-очередь) +- [Сканирование Floating IP из OpenStack](#сканирование-floating-ip-из-openstack) - [Наблюдение за очередью](#наблюдение-за-очередью) - [Значения полей IP](#значения-полей-ip) - [Как читать итоговый результат (pass/partial/fail)](#как-читать-итоговый-результат-passpartialfail) - [Просмотр деталей и истории по конкретному адресу](#просмотр-деталей-и-истории-по-конкретному-адресу) +- [Реестр адресов и глубина истории](#реестр-адресов-и-глубина-истории) - [Управление валидаторами](#управление-валидаторами) - [Управление площадками (проберами)](#управление-площадками-проберами) - [Управление типами проверок пробера](#управление-типами-проверок-пробера) @@ -75,6 +77,31 @@ curl -s -X POST http://:8080/api/v1/admin/ips \ при разворачивании стенда — но для повседневного добавления адресов проще и быстрее пользоваться API выше. +## Сканирование Floating IP из OpenStack + +Вместо того чтобы перечислять адреса вручную, можно попросить control-api +самому найти их в облаке: + +```bash +curl -s -X POST http://:8080/api/v1/admin/ips/scan +``` + +Сканируются все Floating IP текущего проекта OpenStack, но в очередь +ставятся только **свободные** — те, что не привязаны сейчас ни к одному +порту (это и есть пул адресов, ожидающих проверки перед повторной +выдачей). Уже привязанные к чему-то Floating IP игнорируются. Внутри +вызов делает то же самое, что и обычное добавление — новые адреса +встают в очередь, уже завершённые перезапускаются, активно проверяемые не +трогаются (см. [выше](#добавление-новых-ip-в-очередь)). + +В `admin-dashboard` то же самое — кнопка «Сканировать Floating IP» на +странице `/ips`. + +Если хочется, чтобы сканирование происходило само по расписанию, а не +только по запросу — задайте `orchestrator.fip_scan_interval_seconds` +(в секундах) в `control-api.yaml`; `0` (по умолчанию) оставляет только +ручной запуск через ручку/кнопку выше. + ## Наблюдение за очередью Общая сводка: @@ -198,6 +225,41 @@ curl -s http://:8080/api/v1/admin/ips/203.0.113.10 \ | jq '.checks[] | select(.Success==false)' ``` +Это — только текущая попытка. Полная история адреса за всё время, включая +предыдущие попытки и даже периоды, когда адрес не стоял в очереди вовсе, +смотрится через реестр — см. следующий раздел. + +## Реестр адресов и глубина истории + +`GET /api/v1/admin/ips/{ip}` (и `/ips/{ip}` в дашборде) показывает только +*текущую* попытку проверки. Но если адрес удалили из очереди и позже +добавили заново, это уже новая попытка — а куда девается история старой? +Она не пропадает: control-api ведёт отдельный **реестр** (`ip_registry`) — +запись обо всех адресах, когда-либо поставленных на проверку, вместе с +полной накопленной историей проверок по каждому, независимо от того, +удалялся ли адрес из очереди и добавлялся ли повторно. + +```bash +# все адреса, когда-либо ставившиеся на проверку, с краткой сводкой +curl -s http://:8080/api/v1/admin/registry | python3 -m json.tool + +# полная история проверок одного адреса, по всем циклам, не только текущему +curl -s http://:8080/api/v1/admin/registry/203.0.113.10 | python3 -m json.tool +``` + +В `admin-dashboard` — страницы `/registry` (список) и `/registry/{ip}` +(история конкретного адреса), со ссылкой туда со страницы `/ips/{ip}`. + +**Глубина хранения.** Чтобы история не росла бесконечно на адресах, +которые перепроверяют очень часто, можно ограничить, сколько последних +циклов проверки хранить на каждый адрес — `history_retention_cycles` на +странице `/settings` (или `PUT /api/v1/admin/config/orchestrator`, см. +[API.md](API.md#настройки-оркестратора-apiv1adminconfigorchestrator)). +`0` (по умолчанию) — хранить без ограничения. Ограничение действует только +на глубину детальной истории проверок; сама запись в реестре (что адрес +существует, когда впервые встречен, сколько всего было циклов) не +удаляется никогда. + ## Управление валидаторами Список валидаторов и их текущее состояние: @@ -457,12 +519,16 @@ curl -s http://:8080/api/v1/admin/ips/203.0.113.10 | python3 -m jso ## Удаление адресов из очереди -**Отличие от отмены (`cancel`) выше: удаление безвозвратно.** Cancel -переводит адрес в `failed`/`cancelled` и сохраняет запись как историю — -её видно в очереди и в деталях адреса. Delete физически стирает строку -`ip_queue` и всю её историю проверок и событий: адрес полностью исчезает, -восстановить его нельзя. Если нужно просто остановить зависшую проверку, -но сохранить её результат в истории — используйте +**Отличие от отмены (`cancel`) выше: удаление безвозвратно убирает адрес +из очереди.** Cancel переводит адрес в `failed`/`cancelled` и сохраняет +запись как историю — её видно в очереди и в деталях адреса. Delete +физически стирает строку `ip_queue`: адрес полностью исчезает из +`/ips`/`/ips/{ip}`, восстановить именно эту строку нельзя. Накопленная +история проверок при этом **не теряется** — она остаётся в +[реестре](#реестр-адресов-и-глубина-истории) (`/registry/{ip}`) и видна +там даже после удаления адреса из очереди. Если нужно просто остановить +зависшую проверку, но сохранить её результат прямо в очереди — +используйте [«Принудительную остановку проверки»](#принудительная-остановка-проверки) выше, а не удаление. diff --git a/internal/config/config.go b/internal/config/config.go index 037866d..511a33c 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -78,6 +78,12 @@ type OrchestratorConfig struct { // field in this struct, it deliberately gets no nonzero default in // LoadControlAPI below. FIPSettleSeconds int `yaml:"fip_settle_seconds"` + // FIPScanIntervalSeconds enables periodic Orchestrator.ScanFloatingIPs + // runs on this interval when > 0. 0 (the default) disables periodic + // scanning — an operator can still trigger a scan manually via + // POST /api/v1/admin/ips/scan or the dashboard's "Scan Floating IPs" + // button either way. + FIPScanIntervalSeconds int `yaml:"fip_scan_interval_seconds"` } type AggregationConfig struct { diff --git a/internal/dashboard/client.go b/internal/dashboard/client.go index 4a7c41c..29dfd87 100644 --- a/internal/dashboard/client.go +++ b/internal/dashboard/client.go @@ -136,6 +136,32 @@ func (c *client) ClearQueue(ctx context.Context) (clearQueueResponse, error) { return out, err } +// ScanFloatingIPs lists the OpenStack project's free (unassociated) +// floating IPs and submits them to the check queue — see +// orchestrator.ScanFloatingIPs. +func (c *client) ScanFloatingIPs(ctx context.Context) (scanIPsResponse, error) { + var out scanIPsResponse + err := c.do(ctx, http.MethodPost, "/api/v1/admin/ips/scan", nil, &out) + return out, err +} + +// ListRegistry returns every address ever submitted to the check queue, +// each with a summary of its accumulated check history — survives an +// address being deleted from the queue and later re-added. +func (c *client) ListRegistry(ctx context.Context) ([]registryItem, error) { + var out []registryItem + err := c.do(ctx, http.MethodGet, "/api/v1/admin/registry", nil, &out) + return out, err +} + +// GetRegistryHistory returns one address's registry record plus its full +// retained check history across every cycle still kept. +func (c *client) GetRegistryHistory(ctx context.Context, ip string) (registryHistoryResponse, error) { + var out registryHistoryResponse + err := c.do(ctx, http.MethodGet, "/api/v1/admin/registry/"+url.PathEscape(ip), nil, &out) + return out, err +} + func (c *client) ListValidators(ctx context.Context) ([]validatorDTO, error) { var out []validatorDTO err := c.do(ctx, http.MethodGet, "/api/v1/admin/config/validators", nil, &out) @@ -207,10 +233,10 @@ func (c *client) GetOrchestratorSettings(ctx context.Context) (orchestratorSetti return out, err } -func (c *client) PutOrchestratorSettings(ctx context.Context, fipSettleSeconds int) (orchestratorSettingsDTO, error) { +func (c *client) PutOrchestratorSettings(ctx context.Context, fipSettleSeconds, historyRetentionCycles int) (orchestratorSettingsDTO, error) { var out orchestratorSettingsDTO err := c.do(ctx, http.MethodPut, "/api/v1/admin/config/orchestrator", - orchestratorSettingsDTO{FIPSettleSeconds: fipSettleSeconds}, &out) + orchestratorSettingsDTO{FIPSettleSeconds: fipSettleSeconds, HistoryRetentionCycles: historyRetentionCycles}, &out) return out, err } diff --git a/internal/dashboard/dashboard_test.go b/internal/dashboard/dashboard_test.go index 0b1dd06..29554f4 100644 --- a/internal/dashboard/dashboard_test.go +++ b/internal/dashboard/dashboard_test.go @@ -33,16 +33,26 @@ type fakeControlAPI struct { fipSettleSeconds int inboundPorts []int inboundICMP bool + + historyRetentionCycles int + // scanFreeAddresses is what POST /ips/scan "discovers" — tests set it + // directly rather than this fake reimplementing OpenStack floating-IP + // filtering (already covered by internal/orchestrator's own tests). + scanFreeAddresses []string + registry map[string]registryItem + registryChecks map[string][]check } func newFakeControlAPI(t *testing.T) (*fakeControlAPI, string) { t.Helper() f := &fakeControlAPI{ - sites: map[int]string{}, - siteState: map[int]string{}, - siteHostname: map[int]string{}, - groups: map[string][]string{}, - checkTypes: map[string]checkTypeDTO{}, + sites: map[int]string{}, + siteState: map[int]string{}, + siteHostname: map[int]string{}, + groups: map[string][]string{}, + checkTypes: map[string]checkTypeDTO{}, + registry: map[string]registryItem{}, + registryChecks: map[string][]check{}, } ts := httptest.NewServer(f.handler()) t.Cleanup(ts.Close) @@ -197,7 +207,10 @@ func (f *fakeControlAPI) handler() http.Handler { mux.HandleFunc("GET /api/v1/admin/config/orchestrator", func(w http.ResponseWriter, r *http.Request) { f.mu.Lock() defer f.mu.Unlock() - writeJSON(w, http.StatusOK, orchestratorSettingsDTO{FIPSettleSeconds: f.fipSettleSeconds}) + writeJSON(w, http.StatusOK, orchestratorSettingsDTO{ + FIPSettleSeconds: f.fipSettleSeconds, + HistoryRetentionCycles: f.historyRetentionCycles, + }) }) mux.HandleFunc("PUT /api/v1/admin/config/orchestrator", func(w http.ResponseWriter, r *http.Request) { f.mu.Lock() @@ -208,8 +221,54 @@ func (f *fakeControlAPI) handler() http.Handler { writeAPIErr(w, http.StatusBadRequest, "fip_settle_seconds must be >= 0") return } + if req.HistoryRetentionCycles < 0 { + writeAPIErr(w, http.StatusBadRequest, "history_retention_cycles must be >= 0") + return + } f.fipSettleSeconds = req.FIPSettleSeconds - writeJSON(w, http.StatusOK, orchestratorSettingsDTO{FIPSettleSeconds: f.fipSettleSeconds}) + f.historyRetentionCycles = req.HistoryRetentionCycles + writeJSON(w, http.StatusOK, orchestratorSettingsDTO{ + FIPSettleSeconds: f.fipSettleSeconds, + HistoryRetentionCycles: f.historyRetentionCycles, + }) + }) + + mux.HandleFunc("POST /api/v1/admin/ips/scan", func(w http.ResponseWriter, r *http.Request) { + f.mu.Lock() + defer f.mu.Unlock() + resp := scanIPsResponse{ScannedFree: len(f.scanFreeAddresses)} + for _, addr := range f.scanFreeAddresses { + idx := f.findIP(addr) + if idx < 0 { + now := time.Now() + f.ips = append(f.ips, ipQueueItem{IPAddress: addr, State: "queued", CreatedAt: now, UpdatedAt: now}) + resp.Added = append(resp.Added, addr) + continue + } + resp.Reordered = append(resp.Reordered, addr) + } + writeJSON(w, http.StatusOK, resp) + }) + + mux.HandleFunc("GET /api/v1/admin/registry", func(w http.ResponseWriter, r *http.Request) { + f.mu.Lock() + defer f.mu.Unlock() + out := make([]registryItem, 0, len(f.registry)) + for _, item := range f.registry { + out = append(out, item) + } + writeJSON(w, http.StatusOK, out) + }) + mux.HandleFunc("GET /api/v1/admin/registry/{ip}", func(w http.ResponseWriter, r *http.Request) { + f.mu.Lock() + defer f.mu.Unlock() + addr := r.PathValue("ip") + item, ok := f.registry[addr] + if !ok { + writeAPIErr(w, http.StatusNotFound, "unknown ip: "+addr) + return + } + writeJSON(w, http.StatusOK, registryHistoryResponse{Registry: item, Checks: f.registryChecks[addr]}) }) mux.HandleFunc("GET /api/v1/admin/config/inbound-checks", func(w http.ResponseWriter, r *http.Request) { diff --git a/internal/dashboard/dto.go b/internal/dashboard/dto.go index bc6941d..4163ecc 100644 --- a/internal/dashboard/dto.go +++ b/internal/dashboard/dto.go @@ -41,6 +41,8 @@ type ipQueueItem struct { type check struct { ID int64 + RegistryID int64 + CycleID int IPID int64 IPAddress string AttemptNumber int @@ -96,6 +98,33 @@ type clearQueueResponse struct { Deleted []string `json:"deleted"` } +type scanIPsResponse struct { + ScannedFree int `json:"scanned_free"` + Added []string `json:"added"` + Requeued []string `json:"requeued"` + Reordered []string `json:"reordered"` + SkippedInProgress []string `json:"skipped_in_progress"` +} + +// registryItem is one row of the durable per-address registry — see +// httpapi's registryDTO. Survives an address being deleted from the check +// queue and later re-added, unlike ipQueueItem above. +type registryItem struct { + IPAddress string `json:"ip_address"` + FirstSeenAt time.Time `json:"first_seen_at"` + LastSeenAt time.Time `json:"last_seen_at"` + TotalCycles int `json:"total_cycles"` + LastResult string `json:"last_result"` + LastCheckedAt *time.Time `json:"last_checked_at"` + InQueue bool `json:"in_queue"` + CurrentState string `json:"current_state"` +} + +type registryHistoryResponse struct { + Registry registryItem `json:"registry"` + Checks []check `json:"checks"` +} + type validatorDTO struct { ValidatorID string `json:"validator_id"` Hostname string `json:"hostname"` @@ -128,7 +157,8 @@ type errorResponse struct { } type orchestratorSettingsDTO struct { - FIPSettleSeconds int `json:"fip_settle_seconds"` + FIPSettleSeconds int `json:"fip_settle_seconds"` + HistoryRetentionCycles int `json:"history_retention_cycles"` } type inboundChecksDTO struct { diff --git a/internal/dashboard/handlers_ips.go b/internal/dashboard/handlers_ips.go index 53ff12c..ed91443 100644 --- a/internal/dashboard/handlers_ips.go +++ b/internal/dashboard/handlers_ips.go @@ -117,3 +117,11 @@ func (s *Server) handleIPsClear(w http.ResponseWriter, r *http.Request) { _, err := s.CA.ClearQueue(r.Context()) s.renderIPsTable(w, r, err) } + +// handleIPsScan lists the OpenStack project's free (unassociated) floating +// IPs and submits them to the check queue in one step — see +// client.ScanFloatingIPs. +func (s *Server) handleIPsScan(w http.ResponseWriter, r *http.Request) { + _, err := s.CA.ScanFloatingIPs(r.Context()) + s.renderIPsTable(w, r, err) +} diff --git a/internal/dashboard/handlers_registry.go b/internal/dashboard/handlers_registry.go new file mode 100644 index 0000000..7cf847d --- /dev/null +++ b/internal/dashboard/handlers_registry.go @@ -0,0 +1,37 @@ +package dashboard + +import "net/http" + +type registryPageData struct { + PageData + Items []registryItem +} + +type registryDetailData struct { + PageData + History registryHistoryResponse +} + +// handleRegistryPage lists every address ever submitted to the check +// queue, with a summary of its accumulated check history — the durable +// record that survives an address being deleted from /ips and later +// re-added. See internal/db/migrations/0007_ip_registry.sql. +func (s *Server) handleRegistryPage(w http.ResponseWriter, r *http.Request) { + items, err := s.CA.ListRegistry(r.Context()) + data := registryPageData{Items: items} + data.ActiveNav = "registry" + data.Banner = bannerFor(err) + s.renderPage(w, "registry_page", data) +} + +// handleRegistryDetail shows one address's full retained check history +// across every cycle it has ever run, not just the current attempt — see +// ip_detail_content in ip_detail.html for the attempt-scoped equivalent. +func (s *Server) handleRegistryDetail(w http.ResponseWriter, r *http.Request) { + ip := r.PathValue("ip") + history, err := s.CA.GetRegistryHistory(r.Context(), ip) + data := registryDetailData{History: history} + data.ActiveNav = "registry" + data.Banner = bannerFor(err) + s.renderPage(w, "registry_detail_page", data) +} diff --git a/internal/dashboard/handlers_settings.go b/internal/dashboard/handlers_settings.go index fa57816..04efdae 100644 --- a/internal/dashboard/handlers_settings.go +++ b/internal/dashboard/handlers_settings.go @@ -51,7 +51,12 @@ func (s *Server) handleSettingsPut(w http.ResponseWriter, r *http.Request) { s.renderSettingsForm(w, r, &apiErr{Status: http.StatusBadRequest, Message: "пауза должна быть целым числом секунд"}) return } - _, err = s.CA.PutOrchestratorSettings(r.Context(), seconds) + retentionCycles, err := strconv.Atoi(r.PostFormValue("history_retention_cycles")) + if err != nil { + s.renderSettingsForm(w, r, &apiErr{Status: http.StatusBadRequest, Message: "глубина истории должна быть целым числом циклов"}) + return + } + _, err = s.CA.PutOrchestratorSettings(r.Context(), seconds, retentionCycles) s.renderSettingsForm(w, r, err) } diff --git a/internal/dashboard/handlers_test.go b/internal/dashboard/handlers_test.go index 63d1205..3eb7df5 100644 --- a/internal/dashboard/handlers_test.go +++ b/internal/dashboard/handlers_test.go @@ -230,27 +230,44 @@ func TestSettingsGetAndPut(t *testing.T) { t.Fatalf("expected default 0 in the form, got:\n%s", page) } - body := postForm(t, ts, "PUT", "/settings", map[string][]string{"fip_settle_seconds": {"15"}}) - if !strings.Contains(body, `value="15"`) { - t.Fatalf("expected updated value 15 in re-rendered form, got:\n%s", body) + body := postForm(t, ts, "PUT", "/settings", map[string][]string{ + "fip_settle_seconds": {"15"}, "history_retention_cycles": {"10"}, + }) + if !strings.Contains(body, `value="15"`) || !strings.Contains(body, `value="10"`) { + t.Fatalf("expected updated values 15/10 in re-rendered form, got:\n%s", body) } if fake.fipSettleSeconds != 15 { t.Fatalf("expected fake control-api settings updated, got %d", fake.fipSettleSeconds) } + if fake.historyRetentionCycles != 10 { + t.Fatalf("expected fake control-api history_retention_cycles updated, got %d", fake.historyRetentionCycles) + } // A control-api validation error (negative value here) surfaces via // the banner, not a crash. - body = postForm(t, ts, "PUT", "/settings", map[string][]string{"fip_settle_seconds": {"-1"}}) + body = postForm(t, ts, "PUT", "/settings", map[string][]string{ + "fip_settle_seconds": {"-1"}, "history_retention_cycles": {"10"}, + }) if !strings.Contains(body, "alert-warning") { t.Fatalf("expected client error banner for invalid value, got:\n%s", body) } // A non-numeric value is caught by the dashboard itself before it ever // reaches control-api. - body = postForm(t, ts, "PUT", "/settings", map[string][]string{"fip_settle_seconds": {"not-a-number"}}) + body = postForm(t, ts, "PUT", "/settings", map[string][]string{ + "fip_settle_seconds": {"not-a-number"}, "history_retention_cycles": {"10"}, + }) if !strings.Contains(body, "alert-warning") { t.Fatalf("expected client error banner for non-numeric value, got:\n%s", body) } + + // Same for a non-numeric retention value. + body = postForm(t, ts, "PUT", "/settings", map[string][]string{ + "fip_settle_seconds": {"15"}, "history_retention_cycles": {"not-a-number"}, + }) + if !strings.Contains(body, "alert-warning") { + t.Fatalf("expected client error banner for non-numeric retention value, got:\n%s", body) + } } func TestSettingsPageShowsInboundChecks(t *testing.T) { @@ -463,6 +480,52 @@ func TestTargetsAndCheckTypesRoundTrip(t *testing.T) { } } +func TestIPsScan(t *testing.T) { + fake, caURL := newFakeControlAPI(t) + fake.scanFreeAddresses = []string{"5.5.5.5"} + ts := newTestServer(t, caURL) + + body := postForm(t, ts, "POST", "/ips/scan", nil) + if !strings.Contains(body, "5.5.5.5") { + t.Fatalf("expected scanned address in re-rendered table, got:\n%s", body) + } + if len(fake.ips) != 1 || fake.ips[0].IPAddress != "5.5.5.5" { + t.Fatalf("expected fake control-api queue to contain the scanned address, got %+v", fake.ips) + } +} + +func TestRegistryPageAndDetail(t *testing.T) { + fake, caURL := newFakeControlAPI(t) + now := time.Now() + fake.registry["9.9.9.9"] = registryItem{ + IPAddress: "9.9.9.9", FirstSeenAt: now, LastSeenAt: now, + TotalCycles: 2, LastResult: "pass", InQueue: false, + } + fake.registryChecks["9.9.9.9"] = []check{ + {CycleID: 2, Source: "egress", CheckType: "https", Target: "https://example.test", Success: true, CheckedAt: now}, + {CycleID: 1, Source: "egress", CheckType: "https", Target: "https://example.test", Success: false, CheckedAt: now}, + } + ts := newTestServer(t, caURL) + + page := get(t, ts, "/registry") + if !strings.Contains(page, "9.9.9.9") || !strings.Contains(page, "/registry/9.9.9.9") { + t.Fatalf("expected registry list to show the address with a detail link, got:\n%s", page) + } + + detail := get(t, ts, "/registry/9.9.9.9") + if !strings.Contains(detail, "9.9.9.9") { + t.Fatalf("expected detail page to show the address, got:\n%s", detail) + } + if !strings.Contains(detail, "https://example.test") { + t.Fatalf("expected detail page to list retained checks, got:\n%s", detail) + } + + notFound := get(t, ts, "/registry/1.2.3.4") + if !strings.Contains(notFound, "alert-warning") { + t.Fatalf("expected client error banner for unknown registry address, got:\n%s", notFound) + } +} + func TestControlAPIUnreachable(t *testing.T) { // Point the dashboard at an address nothing listens on, rather than a // closed httptest.Server, to get a deterministic connection-refused diff --git a/internal/dashboard/routes.go b/internal/dashboard/routes.go index 8c9f05d..1f13206 100644 --- a/internal/dashboard/routes.go +++ b/internal/dashboard/routes.go @@ -11,6 +11,7 @@ func (s *Server) routes(mux *http.ServeMux) { mux.HandleFunc("GET /ips", s.handleIPsPage) mux.HandleFunc("GET /ips/{ip}", s.handleIPDetail) mux.HandleFunc("POST /ips", s.handleIPsSubmit) + mux.HandleFunc("POST /ips/scan", s.handleIPsScan) mux.HandleFunc("POST /ips/{ip}/recheck", s.handleIPRecheck) mux.HandleFunc("POST /ips/{ip}/cancel", s.handleIPCancel) mux.HandleFunc("DELETE /ips/{ip}", s.handleIPDelete) @@ -18,6 +19,9 @@ func (s *Server) routes(mux *http.ServeMux) { mux.HandleFunc("POST /ips/recheck", s.handleIPsRecheckSelected) mux.HandleFunc("POST /ips/clear", s.handleIPsClear) + mux.HandleFunc("GET /registry", s.handleRegistryPage) + mux.HandleFunc("GET /registry/{ip}", s.handleRegistryDetail) + mux.HandleFunc("GET /validators", s.handleValidatorsPage) mux.HandleFunc("POST /validators", s.handleValidatorCreate) mux.HandleFunc("PUT /validators/{id}", s.handleValidatorUpdate) diff --git a/internal/dashboard/templates/ip_detail.html b/internal/dashboard/templates/ip_detail.html index 7ac5b7b..38289d0 100644 --- a/internal/dashboard/templates/ip_detail.html +++ b/internal/dashboard/templates/ip_detail.html @@ -25,6 +25,7 @@

{{.Detail.IP.IPAddress}} {{$b.Label}}

+

Полная история проверок этого адреса за всё время →

попытка {{.Detail.IP.AttemptNumber}} (повторов: {{.Detail.IP.RetryCount}}) · валидатор {{deref .Detail.IP.OwnerValidatorID}} diff --git a/internal/dashboard/templates/ips.html b/internal/dashboard/templates/ips.html index 2a20843..63affa4 100644 --- a/internal/dashboard/templates/ips.html +++ b/internal/dashboard/templates/ips.html @@ -34,7 +34,10 @@

Новый адрес встаёт в очередь; уже завершённый (done/failed) запускается заново -(тот же вызов); тот, что сейчас проверяется, не трогается.

+(тот же вызов); тот, что сейчас проверяется, не трогается. Кнопка «Сканировать Floating IP» ниже делает то же самое +автоматически: находит в проекте OpenStack все свободные (не привязанные к порту) Floating IP и сразу передаёт их +в очередь на проверку. Полная история проверок по каждому адресу — включая уже удалённые из очереди — доступна в +реестре.

@@ -47,6 +50,7 @@
+ diff --git a/internal/dashboard/templates/layout.html b/internal/dashboard/templates/layout.html index fd08793..3c455e8 100644 --- a/internal/dashboard/templates/layout.html +++ b/internal/dashboard/templates/layout.html @@ -52,6 +52,10 @@ Очередь IP + + +Реестр + Валидаторы diff --git a/internal/dashboard/templates/registry.html b/internal/dashboard/templates/registry.html new file mode 100644 index 0000000..661a735 --- /dev/null +++ b/internal/dashboard/templates/registry.html @@ -0,0 +1,55 @@ +{{define "registry_page"}} + + +{{template "html_head" .}} + +
+ +
+{{template "sidebar_nav" .}} +
+{{template "topbar_mobile" .}} +
+
{{template "banner_inner" .Banner}}
+{{template "registry_content" .}} +
+
+
+ + +{{end}} + +{{define "registry_content"}} +

Реестр IP

+

Все адреса, когда-либо поставленные на проверку — накопленная статистика +сохраняется здесь даже после удаления адреса из очереди и не теряется при повторном добавлении. +Глубина хранимой истории на адрес настраивается на странице настроек.

+ +{{if .Items}} +
+
+ + + +{{range .Items}} + + + + + + + + +{{end}} + +
АдресВпервые замеченПоследний раз замеченЦикловПоследний результатСейчас в очереди
{{.IPAddress}}{{fmtTime .FirstSeenAt}}{{fmtTime .LastSeenAt}}{{.TotalCycles}} +{{if eq .LastResult "pass"}}pass +{{else if eq .LastResult "fail"}}fail +{{else}}—{{end}} + +{{if .InQueue}}{{.CurrentState}}{{else}}нет{{end}} +
+
+
+{{else}}

Реестр пуст — ни один адрес ещё не ставился на проверку.

{{end}} +{{end}} diff --git a/internal/dashboard/templates/registry_detail.html b/internal/dashboard/templates/registry_detail.html new file mode 100644 index 0000000..55f9f44 --- /dev/null +++ b/internal/dashboard/templates/registry_detail.html @@ -0,0 +1,56 @@ +{{define "registry_detail_page"}} + + +{{template "html_head" .}} + +
+ +
+{{template "sidebar_nav" .}} +
+{{template "topbar_mobile" .}} +
+
{{template "banner_inner" .Banner}}
+{{template "registry_detail_content" .}} +
+
+
+ + +{{end}} + +{{define "registry_detail_content"}} +

← К реестру

+
+

{{.History.Registry.IPAddress}} +{{if .History.Registry.InQueue}}{{.History.Registry.CurrentState}} · в очереди{{else}}не в очереди{{end}} +

+
+

+впервые замечен: {{fmtTime .History.Registry.FirstSeenAt}} · последний раз замечен: {{fmtTime .History.Registry.LastSeenAt}} +· всего циклов проверки: {{.History.Registry.TotalCycles}} +

+ +

История проверок (все сохранённые циклы)

+{{if .History.Checks}} +
+
+ + + +{{range .History.Checks}} + + + + + + +{{end}} + +
ЦиклИсточникТипЦельУспехЗадержка, мсВремяДетали
{{.CycleID}}{{.Source}}{{.CheckType}}{{.Target}}{{if .Success}}ok{{else}}fail{{end}}{{.LatencyMS}}{{fmtTime .CheckedAt}}{{.Detail}}
+
+
+

Если глубина истории (см. настройки) ограничена, здесь +показаны только сохранённые циклы — более старые уже вытеснены.

+{{else}}

Проверок пока нет.

{{end}} +{{end}} diff --git a/internal/dashboard/templates/settings.html b/internal/dashboard/templates/settings.html index 213bc8c..74e3fb2 100644 --- a/internal/dashboard/templates/settings.html +++ b/internal/dashboard/templates/settings.html @@ -35,12 +35,20 @@ проверять его. 0 — без паузы (поведение по умолчанию). Значение должно оставлять запас внутри лизинга адреса: fip_settle_seconds + self_check_timeout_seconds < lease_ttl_seconds — иначе адрес не успеет пройти self-check до истечения лизинга и будет возвращён в очередь.

+

Сколько последних циклов проверки хранить в +реестре для каждого адреса. Адрес и его накопленная статистика остаются в реестре даже +после удаления из очереди — эта настройка ограничивает только глубину истории конкретных проверок, не сам реестр. +0 — хранить без ограничения.

+
+ + +
diff --git a/internal/db/db.go b/internal/db/db.go index d1d2e08..b9bc2f4 100644 --- a/internal/db/db.go +++ b/internal/db/db.go @@ -31,6 +31,9 @@ var unboundedSitesSchema string //go:embed migrations/0006_prober_heartbeat.sql var proberHeartbeatSchema string +//go:embed migrations/0007_ip_registry.sql +var ipRegistrySchema string + // migrations is the ordered list of schema versions. Each entry's SQL is // applied, in order, for any version greater than the database's current // PRAGMA user_version — so a fresh database walks the whole list and an @@ -45,6 +48,7 @@ var migrations = []struct { {4, inboundChecksAdminSchema}, {5, unboundedSitesSchema}, {6, proberHeartbeatSchema}, + {7, ipRegistrySchema}, } type DB struct { diff --git a/internal/db/migrations/0007_ip_registry.sql b/internal/db/migrations/0007_ip_registry.sql new file mode 100644 index 0000000..73bd87f --- /dev/null +++ b/internal/db/migrations/0007_ip_registry.sql @@ -0,0 +1,80 @@ +-- Durable per-address registry, decoupled from ip_queue's lifecycle: today +-- deleting an address from ip_queue (DeleteIP/DeleteIPs/ClearQueue) cascades +-- to a hard DELETE of its checks/events, so history is lost forever if an +-- address is removed and later re-added. ip_registry gives every address +-- ever submitted a durable identity that check/event history attaches to +-- instead, surviving ip_queue row deletion and recreation. + +CREATE TABLE ip_registry ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + ip_address TEXT NOT NULL UNIQUE, + first_seen_at TIMESTAMP NOT NULL, + last_seen_at TIMESTAMP NOT NULL, + next_cycle INTEGER NOT NULL DEFAULT 1, + created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +-- Backfill: one registry row per address currently (or ever) in ip_queue. +-- ip_queue.ip_address is already UNIQUE, so this is a straight 1:1 copy. +-- next_cycle starts past the current attempt_number so the first future +-- resubmission of an address gets a cycle_id that has never been used +-- before, even though it's seeded from attempt_number here. +INSERT INTO ip_registry (ip_address, first_seen_at, last_seen_at, next_cycle, created_at, updated_at) +SELECT ip_address, created_at, updated_at, attempt_number + 1, created_at, updated_at FROM ip_queue; + +ALTER TABLE ip_queue ADD COLUMN registry_id INTEGER REFERENCES ip_registry(id); +ALTER TABLE ip_queue ADD COLUMN cycle_id INTEGER NOT NULL DEFAULT 1; +UPDATE ip_queue SET + registry_id = (SELECT id FROM ip_registry r WHERE r.ip_address = ip_queue.ip_address), + cycle_id = attempt_number; + +-- checks.ip_id is today NOT NULL + REFERENCES ip_queue(id), which is exactly +-- what forces the cascading DELETE on ip_queue row removal (foreign_keys=ON +-- would otherwise block the delete once ip_queue's row disappears out from +-- under a referencing row). To let history outlive its ip_queue row, ip_id +-- must become nullable and checks must carry their own durable registry_id. +-- SQLite has no ALTER to relax a column's NOT NULL/REFERENCES, so the table +-- is rebuilt. +CREATE TABLE checks_new ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + registry_id INTEGER NOT NULL REFERENCES ip_registry(id), + cycle_id INTEGER NOT NULL, + ip_id INTEGER REFERENCES ip_queue(id), + ip_address TEXT NOT NULL, + attempt_number INTEGER NOT NULL, + validator_id TEXT NOT NULL DEFAULT '', + source TEXT NOT NULL, + check_type TEXT NOT NULL, + target TEXT NOT NULL DEFAULT '', + success BOOLEAN NOT NULL, + latency_ms INTEGER NOT NULL DEFAULT 0, + detail TEXT NOT NULL DEFAULT '', + checked_at TIMESTAMP NOT NULL, + created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + UNIQUE(registry_id, cycle_id, source, check_type, target) +); +INSERT INTO checks_new (id, registry_id, cycle_id, ip_id, ip_address, attempt_number, validator_id, + source, check_type, target, success, latency_ms, detail, checked_at, created_at) +SELECT c.id, iq.registry_id, iq.cycle_id, c.ip_id, c.ip_address, c.attempt_number, c.validator_id, + c.source, c.check_type, c.target, c.success, c.latency_ms, c.detail, c.checked_at, c.created_at +FROM checks c JOIN ip_queue iq ON iq.id = c.ip_id; +DROP TABLE checks; +ALTER TABLE checks_new RENAME TO checks; +CREATE INDEX idx_checks_registry_cycle ON checks(registry_id, cycle_id); +CREATE INDEX idx_checks_ip_attempt ON checks(ip_id, attempt_number); + +-- events.ip_id is already nullable, so no rebuild is needed there — just +-- add the durable registry_id (plus cycle_id, so retention pruning can cut +-- events at the same cycle boundary as checks) alongside it. +ALTER TABLE events ADD COLUMN registry_id INTEGER REFERENCES ip_registry(id); +ALTER TABLE events ADD COLUMN cycle_id INTEGER NOT NULL DEFAULT 0; +UPDATE events SET + registry_id = (SELECT registry_id FROM ip_queue WHERE ip_queue.id = events.ip_id), + cycle_id = (SELECT cycle_id FROM ip_queue WHERE ip_queue.id = events.ip_id) +WHERE ip_id IS NOT NULL; +CREATE INDEX idx_events_registry ON events(registry_id, cycle_id); + +-- Configurable history retention depth, in check cycles per address. 0 (the +-- default, matching today's unbounded behavior) means keep everything. +ALTER TABLE settings ADD COLUMN history_retention_cycles INTEGER NOT NULL DEFAULT 0; diff --git a/internal/db/models.go b/internal/db/models.go index da1d922..c7897c8 100644 --- a/internal/db/models.go +++ b/internal/db/models.go @@ -95,12 +95,21 @@ type IPQueueItem struct { FIPAssociatedAt *time.Time AggregatedAt *time.Time FIPReleasedAt *time.Time + RegistryID int64 + CycleID int CreatedAt time.Time UpdatedAt time.Time } type Check struct { - ID int64 + ID int64 + RegistryID int64 + CycleID int + // IPID is the ip_queue row this check was originally recorded against. + // It's cleared to 0 (SQL NULL) if that row was later deleted — history + // stays reachable via RegistryID/CycleID regardless (see + // migrations/0007_ip_registry.sql). 0 is never a valid ip_queue id + // (AUTOINCREMENT starts at 1), so it unambiguously means "orphaned." IPID int64 IPAddress string AttemptNumber int @@ -115,11 +124,33 @@ type Check struct { CreatedAt time.Time } +// RegistryItem is a durable per-address record that survives an address +// being removed from ip_queue and later re-added — see +// migrations/0007_ip_registry.sql. NextCycle is the cycle_id that will be +// assigned the next time this address is (re)submitted; it only ever +// increases, so cycle_id stays unique for this address even across +// ip_queue row deletion/recreation. +type RegistryItem struct { + ID int64 + IPAddress string + FirstSeenAt time.Time + LastSeenAt time.Time + NextCycle int + CreatedAt time.Time + UpdatedAt time.Time +} + type Event struct { ID int64 SourceType string SourceID string IPID *int64 + // RegistryID/CycleID are resolved from IPID at insert time (see + // InsertEvent) and stay set even after the ip_queue row IPID pointed to + // is later deleted, so retention pruning can cut events at the same + // cycle boundary as checks — see migrations/0007_ip_registry.sql. + RegistryID int64 + CycleID int EventType string Payload string OccurredAt time.Time @@ -179,8 +210,11 @@ type DeleteIPsResult struct { // admin-configurable at runtime (see queries_settings.go). type Settings struct { FIPSettleSeconds int - CreatedAt time.Time - UpdatedAt time.Time + // HistoryRetentionCycles caps how many recent check cycles are kept per + // registry address (see PruneRegistryHistory); 0 means unlimited. + HistoryRetentionCycles int + CreatedAt time.Time + UpdatedAt time.Time } // InboundChecksSettings is the singleton row describing what the prober diff --git a/internal/db/queries_checks.go b/internal/db/queries_checks.go index 3d5038d..6bf4525 100644 --- a/internal/db/queries_checks.go +++ b/internal/db/queries_checks.go @@ -2,50 +2,88 @@ package db import ( "context" + "database/sql" ) -// UpsertCheck records (or, on retry, overwrites) a single check result. The -// UNIQUE(ip_id, attempt_number, source, check_type, target) constraint plus +// UpsertCheck records (or, on retry, overwrites) a single check result. +// registry_id/cycle_id are resolved from c.IPID's current ip_queue row at +// write time, so callers (agentcore/probercore) never need to know about +// the registry — see migrations/0007_ip_registry.sql. The +// UNIQUE(registry_id, cycle_id, source, check_type, target) constraint plus // this upsert is what makes agent/prober result submission safely -// retryable without producing duplicate rows. +// retryable without producing duplicate rows, now scoped to the durable +// per-address cycle rather than the ip_queue row's attempt_number, so it +// survives that row being deleted and the address later resubmitted. func (d *DB) UpsertCheck(ctx context.Context, c Check) error { _, err := d.ExecContext(ctx, ` - INSERT INTO checks (ip_id, ip_address, attempt_number, validator_id, source, check_type, target, - success, latency_ms, detail, checked_at, created_at) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) - ON CONFLICT(ip_id, attempt_number, source, check_type, target) DO UPDATE SET + INSERT INTO checks (registry_id, cycle_id, ip_id, ip_address, attempt_number, validator_id, + source, check_type, target, success, latency_ms, detail, checked_at, created_at) + VALUES ((SELECT registry_id FROM ip_queue WHERE id=?), (SELECT cycle_id FROM ip_queue WHERE id=?), + ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT(registry_id, cycle_id, source, check_type, target) DO UPDATE SET validator_id=excluded.validator_id, success=excluded.success, latency_ms=excluded.latency_ms, detail=excluded.detail, checked_at=excluded.checked_at - `, c.IPID, c.IPAddress, c.AttemptNumber, c.ValidatorID, c.Source, c.CheckType, c.Target, + `, c.IPID, c.IPID, c.IPID, c.IPAddress, c.AttemptNumber, c.ValidatorID, c.Source, c.CheckType, c.Target, c.Success, c.LatencyMS, c.Detail, timeToDB(c.CheckedAt), timeToDB(Now())) return err } +const checksSelect = ` + SELECT id, registry_id, cycle_id, ip_id, ip_address, attempt_number, validator_id, source, check_type, target, + success, latency_ms, detail, checked_at, created_at + FROM checks +` + // ListChecksForAttempt returns every check recorded for an IP's current // attempt — the input to overall-result aggregation. func (d *DB) ListChecksForAttempt(ctx context.Context, ipID int64, attemptNumber int) ([]Check, error) { - rows, err := d.QueryContext(ctx, ` - SELECT id, ip_id, ip_address, attempt_number, validator_id, source, check_type, target, - success, latency_ms, detail, checked_at, created_at - FROM checks WHERE ip_id=? AND attempt_number=? + rows, err := d.QueryContext(ctx, checksSelect+` + WHERE ip_id=? AND attempt_number=? ORDER BY source, check_type, target `, ipID, attemptNumber) if err != nil { return nil, err } defer rows.Close() + return scanChecks(rows) +} +// ListChecksForRegistry returns an address's check history across every +// cycle still retained (see PruneRegistryHistory), newest cycle first. A +// nil limit returns everything currently retained. +func (d *DB) ListChecksForRegistry(ctx context.Context, registryID int64, limit *int) ([]Check, error) { + query := checksSelect + `WHERE registry_id=? ORDER BY cycle_id DESC, source, check_type, target` + args := []interface{}{registryID} + if limit != nil { + query += ` LIMIT ?` + args = append(args, *limit) + } + rows, err := d.QueryContext(ctx, query, args...) + if err != nil { + return nil, err + } + defer rows.Close() + return scanChecks(rows) +} + +func scanChecks(rows *sql.Rows) ([]Check, error) { var out []Check for rows.Next() { var c Check + var ipID sql.NullInt64 var checkedAt, createdAt string - if err := rows.Scan(&c.ID, &c.IPID, &c.IPAddress, &c.AttemptNumber, &c.ValidatorID, &c.Source, - &c.CheckType, &c.Target, &c.Success, &c.LatencyMS, &c.Detail, &checkedAt, &createdAt); err != nil { + if err := rows.Scan(&c.ID, &c.RegistryID, &c.CycleID, &ipID, &c.IPAddress, &c.AttemptNumber, + &c.ValidatorID, &c.Source, &c.CheckType, &c.Target, &c.Success, &c.LatencyMS, &c.Detail, + &checkedAt, &createdAt); err != nil { return nil, err } + if ipID.Valid { + c.IPID = ipID.Int64 + } + var err error if c.CheckedAt, err = dbToTime(checkedAt); err != nil { return nil, err } diff --git a/internal/db/queries_events.go b/internal/db/queries_events.go index 73b3d2e..3f6a70c 100644 --- a/internal/db/queries_events.go +++ b/internal/db/queries_events.go @@ -1,12 +1,19 @@ package db -import "context" +import ( + "context" + "database/sql" +) +// InsertEvent records an audit-trail row. When e.IPID is set, registry_id +// and cycle_id are resolved from that ip_queue row's current values, so the +// event stays traceable to its address (via registry_id) and prunable at +// the right cycle boundary even after the ip_queue row itself is deleted. func (d *DB) InsertEvent(ctx context.Context, e Event) error { _, err := d.ExecContext(ctx, ` - INSERT INTO events (source_type, source_id, ip_id, event_type, payload, occurred_at, created_at) - VALUES (?, ?, ?, ?, ?, ?, ?) - `, e.SourceType, e.SourceID, e.IPID, e.EventType, e.Payload, timeToDB(e.OccurredAt), timeToDB(Now())) + INSERT INTO events (source_type, source_id, ip_id, registry_id, cycle_id, event_type, payload, occurred_at, created_at) + VALUES (?, ?, ?, (SELECT registry_id FROM ip_queue WHERE id=?), COALESCE((SELECT cycle_id FROM ip_queue WHERE id=?), 0), ?, ?, ?, ?) + `, e.SourceType, e.SourceID, e.IPID, e.IPID, e.IPID, e.EventType, e.Payload, timeToDB(e.OccurredAt), timeToDB(Now())) return err } @@ -14,7 +21,7 @@ func (d *DB) InsertEvent(ctx context.Context, e Event) error { // first — used by the admin detail endpoint. func (d *DB) ListEventsForIP(ctx context.Context, ipID int64) ([]Event, error) { rows, err := d.QueryContext(ctx, ` - SELECT id, source_type, source_id, ip_id, event_type, payload, occurred_at, created_at + SELECT id, source_type, source_id, ip_id, registry_id, cycle_id, event_type, payload, occurred_at, created_at FROM events WHERE ip_id=? ORDER BY occurred_at DESC `, ipID) if err != nil { @@ -26,7 +33,7 @@ func (d *DB) ListEventsForIP(ctx context.Context, ipID int64) ([]Event, error) { func (d *DB) ListRecentEvents(ctx context.Context, limit int) ([]Event, error) { rows, err := d.QueryContext(ctx, ` - SELECT id, source_type, source_id, ip_id, event_type, payload, occurred_at, created_at + SELECT id, source_type, source_id, ip_id, registry_id, cycle_id, event_type, payload, occurred_at, created_at FROM events ORDER BY id DESC LIMIT ? `, limit) if err != nil { @@ -45,11 +52,16 @@ func scanEvents(rows interface { for rows.Next() { var e Event var ipID *int64 + var registryID sql.NullInt64 var occurredAt, createdAt string - if err := rows.Scan(&e.ID, &e.SourceType, &e.SourceID, &ipID, &e.EventType, &e.Payload, &occurredAt, &createdAt); err != nil { + if err := rows.Scan(&e.ID, &e.SourceType, &e.SourceID, &ipID, ®istryID, &e.CycleID, + &e.EventType, &e.Payload, &occurredAt, &createdAt); err != nil { return nil, err } e.IPID = ipID + if registryID.Valid { + e.RegistryID = registryID.Int64 + } var err error if e.OccurredAt, err = dbToTime(occurredAt); err != nil { return nil, err diff --git a/internal/db/queries_ipqueue.go b/internal/db/queries_ipqueue.go index 97eb2c6..f3f4d00 100644 --- a/internal/db/queries_ipqueue.go +++ b/internal/db/queries_ipqueue.go @@ -21,12 +21,26 @@ func (d *DB) SeedQueue(ctx context.Context, addresses []string) error { now := timeToDB(Now()) for i, addr := range addresses { - _, err := tx.ExecContext(ctx, ` - INSERT INTO ip_queue (ip_address, sequence, state, created_at, updated_at) - VALUES (?, ?, ?, ?, ?) - ON CONFLICT(ip_address) DO NOTHING - `, addr, i, IPQueued, now, now) + var exists bool + if err := tx.QueryRowContext(ctx, `SELECT EXISTS(SELECT 1 FROM ip_queue WHERE ip_address=?)`, addr).Scan(&exists); err != nil { + return fmt.Errorf("check existing %s: %w", addr, err) + } + if exists { + continue + } + registryID, err := findOrCreateRegistryTx(ctx, tx, addr, now) if err != nil { + return err + } + cycle, err := nextRegistryCycleTx(ctx, tx, registryID, now) + if err != nil { + return err + } + if _, err := tx.ExecContext(ctx, ` + INSERT INTO ip_queue (ip_address, sequence, state, registry_id, cycle_id, created_at, updated_at) + VALUES (?, ?, ?, ?, ?, ?, ?) + ON CONFLICT(ip_address) DO NOTHING + `, addr, i, IPQueued, registryID, cycle, now, now); err != nil { return fmt.Errorf("seed %s: %w", addr, err) } } @@ -227,13 +241,21 @@ func (d *DB) RequeueOrFail(ctx context.Context, ipID int64, validatorID string, } if nextState == IPQueued { + var registryID int64 + if err := tx.QueryRowContext(ctx, `SELECT registry_id FROM ip_queue WHERE id=?`, ipID).Scan(®istryID); err != nil { + return err + } + cycle, cErr := nextRegistryCycleTx(ctx, tx, registryID, now) + if cErr != nil { + return cErr + } _, err = tx.ExecContext(ctx, ` UPDATE ip_queue SET state=?, owner_validator_id=NULL, fip_id='', retry_count=?, attempt_number=attempt_number+1, - lease_expires_at=NULL, egress_complete=0, + cycle_id=?, lease_expires_at=NULL, egress_complete=0, overall_result='', assigned_at=NULL, fip_associated_at=NULL, updated_at=? WHERE id=? - `, nextState, retryCount, now, ipID) + `, nextState, retryCount, cycle, now, ipID) } else { _, err = tx.ExecContext(ctx, ` UPDATE ip_queue SET @@ -300,10 +322,18 @@ func (d *DB) SubmitIPs(ctx context.Context, addresses []string) (SubmitIPsResult err := tx.QueryRowContext(ctx, `SELECT state FROM ip_queue WHERE ip_address=?`, addr).Scan(&state) switch { case err == sql.ErrNoRows: + registryID, rErr := findOrCreateRegistryTx(ctx, tx, addr, now) + if rErr != nil { + return result, rErr + } + cycle, cErr := nextRegistryCycleTx(ctx, tx, registryID, now) + if cErr != nil { + return result, cErr + } if _, err := tx.ExecContext(ctx, ` - INSERT INTO ip_queue (ip_address, sequence, state, created_at, updated_at) - VALUES (?, ?, ?, ?, ?) - `, addr, seq, IPQueued, now, now); err != nil { + INSERT INTO ip_queue (ip_address, sequence, state, registry_id, cycle_id, created_at, updated_at) + VALUES (?, ?, ?, ?, ?, ?, ?) + `, addr, seq, IPQueued, registryID, cycle, now, now); err != nil { return result, fmt.Errorf("insert %s: %w", addr, err) } result.Added = append(result.Added, addr) @@ -312,14 +342,22 @@ func (d *DB) SubmitIPs(ctx context.Context, addresses []string) (SubmitIPsResult return result, err case state == IPDone || state == IPFailed || state == IPOccupied: + var registryID int64 + if err := tx.QueryRowContext(ctx, `SELECT registry_id FROM ip_queue WHERE ip_address=?`, addr).Scan(®istryID); err != nil { + return result, err + } + cycle, cErr := nextRegistryCycleTx(ctx, tx, registryID, now) + if cErr != nil { + return result, cErr + } if _, err := tx.ExecContext(ctx, ` UPDATE ip_queue SET state=?, sequence=?, owner_validator_id=NULL, fip_id='', retry_count=0, - attempt_number=attempt_number+1, lease_expires_at=NULL, egress_complete=0, + attempt_number=attempt_number+1, cycle_id=?, lease_expires_at=NULL, egress_complete=0, overall_result='', assigned_at=NULL, fip_associated_at=NULL, aggregated_at=NULL, fip_released_at=NULL, updated_at=? WHERE ip_address=? - `, IPQueued, seq, now, addr); err != nil { + `, IPQueued, seq, cycle, now, addr); err != nil { return result, fmt.Errorf("requeue %s: %w", addr, err) } result.Requeued = append(result.Requeued, addr) @@ -427,8 +465,11 @@ func (d *DB) DeleteIPs(ctx context.Context, addresses []string) (DeleteIPsResult } // deleteIPTx is the shared body of DeleteIP/DeleteIPs: free the owning -// validator, delete dependent checks/events, then the ip_queue row itself -// — the same FK-clearing order DeleteValidator uses for owner_validator_id. +// validator, detach dependent checks/events from the doomed ip_queue row +// (their registry_id already anchors them permanently — see +// migrations/0007_ip_registry.sql — so this is a hand-off, not a loss), +// delete the ephemeral per-attempt ip_site_checks progress flags, then the +// ip_queue row itself. func deleteIPTx(ctx context.Context, tx *sql.Tx, ipID int64) error { now := timeToDB(Now()) if _, err := tx.ExecContext(ctx, ` @@ -437,11 +478,11 @@ func deleteIPTx(ctx context.Context, tx *sql.Tx, ipID int64) error { `, ValidatorIdle, now, ipID); err != nil { return fmt.Errorf("free owning validator: %w", err) } - if _, err := tx.ExecContext(ctx, `DELETE FROM checks WHERE ip_id=?`, ipID); err != nil { - return fmt.Errorf("delete checks: %w", err) + if _, err := tx.ExecContext(ctx, `UPDATE checks SET ip_id=NULL WHERE ip_id=?`, ipID); err != nil { + return fmt.Errorf("detach checks: %w", err) } - if _, err := tx.ExecContext(ctx, `DELETE FROM events WHERE ip_id=?`, ipID); err != nil { - return fmt.Errorf("delete events: %w", err) + if _, err := tx.ExecContext(ctx, `UPDATE events SET ip_id=NULL WHERE ip_id=?`, ipID); err != nil { + return fmt.Errorf("detach events: %w", err) } if _, err := tx.ExecContext(ctx, `DELETE FROM ip_site_checks WHERE ip_id=?`, ipID); err != nil { return fmt.Errorf("delete ip_site_checks: %w", err) @@ -507,7 +548,8 @@ func (d *DB) ListExpiredLeases(ctx context.Context, now time.Time) ([]IPQueueIte const ipQueueSelect = ` SELECT id, ip_address, sequence, state, owner_validator_id, fip_id, attempt_number, retry_count, lease_expires_at, egress_complete, overall_result, - assigned_at, fip_associated_at, aggregated_at, fip_released_at, created_at, updated_at + assigned_at, fip_associated_at, aggregated_at, fip_released_at, + registry_id, cycle_id, created_at, updated_at FROM ip_queue ` @@ -527,18 +569,23 @@ func scanIPQueueItem(row rowScanner) (*IPQueueItem, error) { var item IPQueueItem var owner sql.NullString var leaseExpires, assignedAt, fipAssociatedAt, aggregatedAt, fipReleasedAt sql.NullString + var registryID sql.NullInt64 var createdAt, updatedAt string if err := row.Scan( &item.ID, &item.IPAddress, &item.Sequence, &item.State, &owner, &item.FIPID, &item.AttemptNumber, &item.RetryCount, &leaseExpires, &item.EgressComplete, - &item.OverallResult, &assignedAt, &fipAssociatedAt, &aggregatedAt, &fipReleasedAt, &createdAt, &updatedAt, + &item.OverallResult, &assignedAt, &fipAssociatedAt, &aggregatedAt, &fipReleasedAt, + ®istryID, &item.CycleID, &createdAt, &updatedAt, ); err != nil { return nil, err } if owner.Valid { item.OwnerValidatorID = &owner.String } + if registryID.Valid { + item.RegistryID = registryID.Int64 + } var err error if item.LeaseExpiresAt, err = nullStringToTimePtr(leaseExpires); err != nil { return nil, err diff --git a/internal/db/queries_registry.go b/internal/db/queries_registry.go new file mode 100644 index 0000000..795b4d0 --- /dev/null +++ b/internal/db/queries_registry.go @@ -0,0 +1,216 @@ +package db + +import ( + "context" + "database/sql" + "fmt" + "time" +) + +// findOrCreateRegistryTx returns the ip_registry row id for addr, creating +// it (next_cycle starting at 1) the first time this address is ever seen. +// Safe to call on every (re)submission of an address — ip_address is +// UNIQUE, so a second call for the same address is a no-op lookup. +func findOrCreateRegistryTx(ctx context.Context, tx *sql.Tx, addr string, now string) (int64, error) { + if _, err := tx.ExecContext(ctx, ` + INSERT INTO ip_registry (ip_address, first_seen_at, last_seen_at, next_cycle, created_at, updated_at) + VALUES (?, ?, ?, 1, ?, ?) + ON CONFLICT(ip_address) DO NOTHING + `, addr, now, now, now, now); err != nil { + return 0, fmt.Errorf("find or create registry for %s: %w", addr, err) + } + var id int64 + if err := tx.QueryRowContext(ctx, `SELECT id FROM ip_registry WHERE ip_address=?`, addr).Scan(&id); err != nil { + return 0, fmt.Errorf("lookup registry id for %s: %w", addr, err) + } + return id, nil +} + +// nextRegistryCycleTx allocates the next cycle_id for registryID and +// advances last_seen_at. Call once per fresh generation of check history for +// an address — a brand new ip_queue row, an admin-triggered requeue, or a +// system retry that bumps attempt_number — so every generation's checks get +// a cycle_id no other generation, past or future (even across ip_queue row +// deletion and recreation), will ever reuse. +func nextRegistryCycleTx(ctx context.Context, tx *sql.Tx, registryID int64, now string) (int, error) { + var cycle int + if err := tx.QueryRowContext(ctx, `SELECT next_cycle FROM ip_registry WHERE id=?`, registryID).Scan(&cycle); err != nil { + return 0, fmt.Errorf("read next_cycle for registry %d: %w", registryID, err) + } + if _, err := tx.ExecContext(ctx, ` + UPDATE ip_registry SET next_cycle=next_cycle+1, last_seen_at=?, updated_at=? WHERE id=? + `, now, now, registryID); err != nil { + return 0, fmt.Errorf("advance next_cycle for registry %d: %w", registryID, err) + } + return cycle, nil +} + +// RegistrySummary is one row of the full IP registry, combining the durable +// ip_registry record with a rollup of its check history and, if the address +// currently has a live ip_queue row, that row's state. +type RegistrySummary struct { + RegistryItem + TotalCycles int + LastResult string + LastCheckedAt *time.Time + InQueue bool + CurrentState string +} + +// ListRegistry returns every address ever submitted, newest first-seen +// last, each with a summary of its accumulated check history. Reads two +// simple queries plus one aggregate rather than a single large join, since +// the registry is expected to stay small enough (one row per distinct +// address ever seen) that this is simpler to reason about than a +// multi-way correlated subquery. +func (d *DB) ListRegistry(ctx context.Context) ([]RegistrySummary, error) { + rows, err := d.QueryContext(ctx, ` + SELECT id, ip_address, first_seen_at, last_seen_at, next_cycle, created_at, updated_at + FROM ip_registry ORDER BY first_seen_at + `) + if err != nil { + return nil, err + } + defer rows.Close() + + items, err := scanRegistryItems(rows) + if err != nil { + return nil, err + } + + out := make([]RegistrySummary, len(items)) + for i, item := range items { + s := RegistrySummary{RegistryItem: item} + if err := d.fillRegistrySummary(ctx, &s); err != nil { + return nil, err + } + out[i] = s + } + return out, nil +} + +// GetRegistryByAddress returns the registry row (with summary) for a single +// address, or ErrNotFound if it has never been submitted. +func (d *DB) GetRegistryByAddress(ctx context.Context, address string) (*RegistrySummary, error) { + row := d.QueryRowContext(ctx, ` + SELECT id, ip_address, first_seen_at, last_seen_at, next_cycle, created_at, updated_at + FROM ip_registry WHERE ip_address=? + `, address) + item, err := scanRegistryItem(row) + if err != nil { + if err == sql.ErrNoRows { + return nil, fmt.Errorf("ip %q: %w", address, ErrNotFound) + } + return nil, err + } + s := &RegistrySummary{RegistryItem: *item} + if err := d.fillRegistrySummary(ctx, s); err != nil { + return nil, err + } + return s, nil +} + +func (d *DB) fillRegistrySummary(ctx context.Context, s *RegistrySummary) error { + if err := d.QueryRowContext(ctx, ` + SELECT COUNT(DISTINCT cycle_id) FROM checks WHERE registry_id=? + `, s.ID).Scan(&s.TotalCycles); err != nil { + return err + } + + var lastResult sql.NullBool + var lastCheckedAt sql.NullString + if err := d.QueryRowContext(ctx, ` + SELECT success, checked_at FROM checks WHERE registry_id=? ORDER BY cycle_id DESC, checked_at DESC LIMIT 1 + `, s.ID).Scan(&lastResult, &lastCheckedAt); err != nil && err != sql.ErrNoRows { + return err + } + if lastResult.Valid { + if lastResult.Bool { + s.LastResult = "pass" + } else { + s.LastResult = "fail" + } + } + t, err := nullStringToTimePtr(lastCheckedAt) + if err != nil { + return err + } + s.LastCheckedAt = t + + var state sql.NullString + if err := d.QueryRowContext(ctx, `SELECT state FROM ip_queue WHERE registry_id=?`, s.ID).Scan(&state); err != nil && err != sql.ErrNoRows { + return err + } + if state.Valid { + s.InQueue = true + s.CurrentState = state.String + } + return nil +} + +func scanRegistryItems(rows *sql.Rows) ([]RegistryItem, error) { + var out []RegistryItem + for rows.Next() { + item, err := scanRegistryItem(rows) + if err != nil { + return nil, err + } + out = append(out, *item) + } + return out, rows.Err() +} + +func scanRegistryItem(row rowScanner) (*RegistryItem, error) { + var item RegistryItem + var firstSeenAt, lastSeenAt, createdAt, updatedAt string + if err := row.Scan(&item.ID, &item.IPAddress, &firstSeenAt, &lastSeenAt, &item.NextCycle, &createdAt, &updatedAt); err != nil { + return nil, err + } + var err error + if item.FirstSeenAt, err = dbToTime(firstSeenAt); err != nil { + return nil, err + } + if item.LastSeenAt, err = dbToTime(lastSeenAt); err != nil { + return nil, err + } + if item.CreatedAt, err = dbToTime(createdAt); err != nil { + return nil, err + } + if item.UpdatedAt, err = dbToTime(updatedAt); err != nil { + return nil, err + } + return &item, nil +} + +// PruneRegistryHistory deletes all but the newest keepCycles cycles' worth +// of check history for registryID — a no-op if keepCycles <= 0 (unlimited +// retention) or the address has that many cycles or fewer. events rows tied +// to this registry are pruned to the same cutoff; ip_registry itself and its +// next_cycle counter are never touched, so pruning never risks a future +// cycle_id collision. +func (d *DB) PruneRegistryHistory(ctx context.Context, registryID int64, keepCycles int) error { + if keepCycles <= 0 { + return nil + } + var cutoff sql.NullInt64 + err := d.QueryRowContext(ctx, ` + SELECT MIN(cycle_id) FROM ( + SELECT DISTINCT cycle_id FROM checks WHERE registry_id=? ORDER BY cycle_id DESC LIMIT ? + ) + `, registryID, keepCycles).Scan(&cutoff) + if err != nil { + return fmt.Errorf("find prune cutoff for registry %d: %w", registryID, err) + } + if !cutoff.Valid { + return nil // fewer than keepCycles cycles recorded — nothing to prune + } + if _, err := d.ExecContext(ctx, `DELETE FROM checks WHERE registry_id=? AND cycle_id= 0: %w", ErrValidation) + } + now := timeToDB(Now()) + _, err := d.ExecContext(ctx, ` + UPDATE settings SET history_retention_cycles=?, updated_at=? WHERE id=1 + `, cycles, now) + return err +} diff --git a/internal/httpapi/dto_admin.go b/internal/httpapi/dto_admin.go index 60f9ccb..eeddf1e 100644 --- a/internal/httpapi/dto_admin.go +++ b/internal/httpapi/dto_admin.go @@ -32,6 +32,27 @@ type clearQueueResponse struct { Deleted []string `json:"deleted"` } +type scanIPsResponse struct { + ScannedFree int `json:"scanned_free"` + Added []string `json:"added"` + Requeued []string `json:"requeued"` + Reordered []string `json:"reordered"` + SkippedInProgress []string `json:"skipped_in_progress"` +} + +// registryDTO is one row of the durable per-address registry — see +// db.RegistrySummary. +type registryDTO struct { + IPAddress string `json:"ip_address"` + FirstSeenAt time.Time `json:"first_seen_at"` + LastSeenAt time.Time `json:"last_seen_at"` + TotalCycles int `json:"total_cycles"` + LastResult string `json:"last_result"` + LastCheckedAt *time.Time `json:"last_checked_at"` + InQueue bool `json:"in_queue"` + CurrentState string `json:"current_state"` +} + type validatorDTO struct { ValidatorID string `json:"validator_id"` Hostname string `json:"hostname"` @@ -85,7 +106,8 @@ type putCheckTypeRequest struct { // request body for /api/v1/admin/config/orchestrator — a single-field DTO, // same shape both ways, like putSiteRequest/siteDTO. type orchestratorSettingsDTO struct { - FIPSettleSeconds int `json:"fip_settle_seconds"` + FIPSettleSeconds int `json:"fip_settle_seconds"` + HistoryRetentionCycles int `json:"history_retention_cycles"` } // inboundChecksDTO doubles as both the GET response and the PUT request diff --git a/internal/httpapi/handlers_admin.go b/internal/httpapi/handlers_admin.go index 7fdc952..9dfdd2f 100644 --- a/internal/httpapi/handlers_admin.go +++ b/internal/httpapi/handlers_admin.go @@ -99,6 +99,26 @@ func (s *Server) handleAdminSubmitIPs(w http.ResponseWriter, r *http.Request) { }) } +// handleAdminScanFloatingIPs lists every floating IP in the configured +// OpenStack project, filters to the free (unassociated) pool, and submits +// that address list to the check queue — see orchestrator.ScanFloatingIPs. +// Takes no body; POST is used (rather than GET) because it mutates the +// queue, matching handleAdminSubmitIPs. +func (s *Server) handleAdminScanFloatingIPs(w http.ResponseWriter, r *http.Request) { + result, scannedFree, err := s.Orch.ScanFloatingIPs(r.Context()) + if err != nil { + writeError(w, http.StatusBadGateway, err.Error()) + return + } + writeJSON(w, http.StatusOK, scanIPsResponse{ + ScannedFree: scannedFree, + Added: emptyIfNil(result.Added), + Requeued: emptyIfNil(result.Requeued), + Reordered: emptyIfNil(result.Reordered), + SkippedInProgress: emptyIfNil(result.SkippedInProgress), + }) +} + // handleAdminCancelIP force-stops a check in progress (or still-queued) for // the given address. Requires the orchestrator, since a floating IP may // need to be disassociated in OpenStack. diff --git a/internal/httpapi/handlers_config.go b/internal/httpapi/handlers_config.go index 680e266..1ce0a76 100644 --- a/internal/httpapi/handlers_config.go +++ b/internal/httpapi/handlers_config.go @@ -200,7 +200,10 @@ func (s *Server) handleConfigGetOrchestratorSettings(w http.ResponseWriter, r *h writeDBError(w, err) return } - writeJSON(w, http.StatusOK, orchestratorSettingsDTO{FIPSettleSeconds: settings.FIPSettleSeconds}) + writeJSON(w, http.StatusOK, orchestratorSettingsDTO{ + FIPSettleSeconds: settings.FIPSettleSeconds, + HistoryRetentionCycles: settings.HistoryRetentionCycles, + }) } // handleConfigPutOrchestratorSettings goes through the orchestrator (not a @@ -217,7 +220,14 @@ func (s *Server) handleConfigPutOrchestratorSettings(w http.ResponseWriter, r *h writeDBError(w, err) return } - writeJSON(w, http.StatusOK, orchestratorSettingsDTO{FIPSettleSeconds: req.FIPSettleSeconds}) + if err := s.DB.SetHistoryRetentionCycles(r.Context(), req.HistoryRetentionCycles); err != nil { + writeDBError(w, err) + return + } + writeJSON(w, http.StatusOK, orchestratorSettingsDTO{ + FIPSettleSeconds: req.FIPSettleSeconds, + HistoryRetentionCycles: req.HistoryRetentionCycles, + }) } // --- prober inbound checks --- diff --git a/internal/httpapi/handlers_registry.go b/internal/httpapi/handlers_registry.go new file mode 100644 index 0000000..d68308e --- /dev/null +++ b/internal/httpapi/handlers_registry.go @@ -0,0 +1,59 @@ +package httpapi + +import ( + "net/http" + + "cloudipvalidator/internal/db" +) + +// handleAdminRegistry lists every address ever submitted to the check +// queue, each with a summary of its accumulated check history — the +// durable record that survives an address being deleted from ip_queue and +// later re-added. See migrations/0007_ip_registry.sql. +func (s *Server) handleAdminRegistry(w http.ResponseWriter, r *http.Request) { + items, err := s.DB.ListRegistry(r.Context()) + if err != nil { + writeError(w, http.StatusInternalServerError, err.Error()) + return + } + out := make([]registryDTO, len(items)) + for i, it := range items { + out[i] = registrySummaryToDTO(it) + } + writeJSON(w, http.StatusOK, out) +} + +// handleAdminRegistryHistory returns one address's registry record plus its +// full retained check history (across every cycle still kept — see +// db.PruneRegistryHistory / settings.history_retention_cycles), newest +// cycle first. +func (s *Server) handleAdminRegistryHistory(w http.ResponseWriter, r *http.Request) { + address := r.PathValue("ip") + summary, err := s.DB.GetRegistryByAddress(r.Context(), address) + if err != nil { + writeDBError(w, err) + return + } + checks, err := s.DB.ListChecksForRegistry(r.Context(), summary.ID, nil) + if err != nil { + writeError(w, http.StatusInternalServerError, err.Error()) + return + } + writeJSON(w, http.StatusOK, struct { + Registry registryDTO `json:"registry"` + Checks []db.Check `json:"checks"` + }{registrySummaryToDTO(*summary), checks}) +} + +func registrySummaryToDTO(s db.RegistrySummary) registryDTO { + return registryDTO{ + IPAddress: s.IPAddress, + FirstSeenAt: s.FirstSeenAt, + LastSeenAt: s.LastSeenAt, + TotalCycles: s.TotalCycles, + LastResult: s.LastResult, + LastCheckedAt: s.LastCheckedAt, + InQueue: s.InQueue, + CurrentState: s.CurrentState, + } +} diff --git a/internal/httpapi/handlers_registry_test.go b/internal/httpapi/handlers_registry_test.go new file mode 100644 index 0000000..842a29d --- /dev/null +++ b/internal/httpapi/handlers_registry_test.go @@ -0,0 +1,150 @@ +package httpapi + +import ( + "context" + "encoding/json" + "net/http" + "testing" + "time" + + "cloudipvalidator/internal/db" +) + +// TestScanFloatingIPsEndpoint proves POST /api/v1/admin/ips/scan only +// queues floating IPs that are currently unassociated in OpenStack. +func TestScanFloatingIPsEndpoint(t *testing.T) { + fc, _, _, mock := newConfigTestHarness(t) + mock.Seed("fip-free", "5.5.5.5", "svc-project") + mock.SeedWithPort("fip-occupied", "6.6.6.6", "svc-project", "some-port") + + resp, body := fc.do(http.MethodPost, "/api/v1/admin/ips/scan", nil) + if resp.StatusCode != http.StatusOK { + t.Fatalf("scan: status=%d body=%s", resp.StatusCode, body) + } + var scanResp scanIPsResponse + if err := json.Unmarshal(body, &scanResp); err != nil { + t.Fatalf("unmarshal scan response: %v", err) + } + if scanResp.ScannedFree != 1 { + t.Fatalf("expected 1 free fip scanned, got %+v", scanResp) + } + if len(scanResp.Added) != 1 || scanResp.Added[0] != "5.5.5.5" { + t.Fatalf("expected only 5.5.5.5 added, got %+v", scanResp) + } + + _, body = fc.do(http.MethodGet, "/api/v1/admin/ips", nil) + var ips []db.IPQueueItem + if err := json.Unmarshal(body, &ips); err != nil { + t.Fatalf("unmarshal ips: %v", err) + } + if len(ips) != 1 || ips[0].IPAddress != "5.5.5.5" { + t.Fatalf("expected only the free address queued, got %+v", ips) + } +} + +// TestRegistryHistoryOutlivesIPDeletion proves that after an address +// completes a full check cycle and is then deleted from the queue, its +// history is still reachable via the registry endpoints (though the plain +// /ips/{ip} endpoint now 404s), and that resubmitting the same address adds +// a second, distinct cycle to the same registry entry. +func TestRegistryHistoryOutlivesIPDeletion(t *testing.T) { + fc, d, orch, mock := newConfigTestHarness(t) + ctx := context.Background() + mock.Seed("fip-1", "9.9.9.9", "svc-project") + + fc.do(http.MethodPost, "/api/v1/admin/config/validators", createValidatorRequest{ValidatorID: "validator-1", OSPortID: "port-1"}) + fc.do(http.MethodPut, "/api/v1/admin/config/targets/web", putTargetGroupRequest{Targets: []string{"https://example.test"}}) + fc.do(http.MethodPut, "/api/v1/admin/config/check-types/https", putCheckTypeRequest{Enabled: true, Targets: []string{"web"}}) + fc.do(http.MethodPost, "/api/v1/agents/register", registerAgentRequest{ValidatorID: "validator-1"}) + + runOneCycle := func() { + orch.Tick(ctx) + _, body := fc.do(http.MethodGet, "/api/v1/agents/validator-1/assignment", nil) + var assignment assignmentResponse + if err := json.Unmarshal(body, &assignment); err != nil { + t.Fatalf("unmarshal assignment: %v", err) + } + fc.do(http.MethodPost, "/api/v1/agents/validator-1/self-check", selfCheckRequest{ + IPID: assignment.IPID, DetectedEgress: "9.9.9.9", Success: true, + }) + fc.do(http.MethodPost, "/api/v1/agents/validator-1/results", agentResultsRequest{ + Results: []checkResultDTO{{ + IPID: assignment.IPID, CheckType: "https", Target: "https://example.test", + Success: true, CheckedAt: time.Now().Format(time.RFC3339Nano), + }}, + }) + fc.do(http.MethodPost, "/api/v1/agents/validator-1/complete", agentCompleteRequest{IPID: assignment.IPID}) + orch.Tick(ctx) + } + + fc.do(http.MethodPost, "/api/v1/admin/ips", submitIPsRequest{Addresses: []string{"9.9.9.9"}}) + runOneCycle() + + item, err := d.GetIPByAddress(ctx, "9.9.9.9") + if err != nil { + t.Fatalf("get ip: %v", err) + } + if item.State != db.IPDone { + t.Fatalf("expected done after first cycle, got %s", item.State) + } + + resp, _ := fc.do(http.MethodDelete, "/api/v1/admin/ips/9.9.9.9", nil) + if resp.StatusCode != http.StatusOK { + t.Fatalf("delete ip: status=%d", resp.StatusCode) + } + + resp, _ = fc.do(http.MethodGet, "/api/v1/admin/ips/9.9.9.9", nil) + if resp.StatusCode != http.StatusNotFound { + t.Fatalf("expected 404 for deleted address on /ips, got %d", resp.StatusCode) + } + + resp, body := fc.do(http.MethodGet, "/api/v1/admin/registry/9.9.9.9", nil) + if resp.StatusCode != http.StatusOK { + t.Fatalf("registry history: status=%d body=%s", resp.StatusCode, body) + } + var hist struct { + Registry registryDTO `json:"registry"` + Checks []db.Check `json:"checks"` + } + if err := json.Unmarshal(body, &hist); err != nil { + t.Fatalf("unmarshal registry history: %v", err) + } + if hist.Registry.TotalCycles != 1 { + t.Fatalf("expected 1 retained cycle, got %+v", hist.Registry) + } + if len(hist.Checks) == 0 { + t.Fatalf("expected retained check history, got none") + } + if hist.Registry.InQueue { + t.Fatalf("expected registry entry to report not-in-queue after delete, got %+v", hist.Registry) + } + + // Resubmitting starts a second, distinct cycle on the same registry + // entry rather than colliding with the first. + mock.Seed("fip-1", "9.9.9.9", "svc-project") + fc.do(http.MethodPost, "/api/v1/admin/ips", submitIPsRequest{Addresses: []string{"9.9.9.9"}}) + runOneCycle() + + resp, body = fc.do(http.MethodGet, "/api/v1/admin/registry/9.9.9.9", nil) + if resp.StatusCode != http.StatusOK { + t.Fatalf("registry history (2nd): status=%d body=%s", resp.StatusCode, body) + } + if err := json.Unmarshal(body, &hist); err != nil { + t.Fatalf("unmarshal registry history (2nd): %v", err) + } + if hist.Registry.TotalCycles != 2 { + t.Fatalf("expected 2 retained cycles after resubmission, got %+v", hist.Registry) + } + if !hist.Registry.InQueue { + t.Fatalf("expected registry entry to report in-queue again, got %+v", hist.Registry) + } + + _, body = fc.do(http.MethodGet, "/api/v1/admin/registry", nil) + var all []registryDTO + if err := json.Unmarshal(body, &all); err != nil { + t.Fatalf("unmarshal registry list: %v", err) + } + if len(all) != 1 || all[0].IPAddress != "9.9.9.9" { + t.Fatalf("expected a single registry entry for the address, got %+v", all) + } +} diff --git a/internal/httpapi/routes.go b/internal/httpapi/routes.go index 49a7972..f59b00b 100644 --- a/internal/httpapi/routes.go +++ b/internal/httpapi/routes.go @@ -21,6 +21,7 @@ func (s *Server) routes(mux *http.ServeMux) { mux.HandleFunc("GET /api/v1/admin/status", s.handleAdminStatus) mux.HandleFunc("GET /api/v1/admin/ips", s.handleAdminIPs) mux.HandleFunc("POST /api/v1/admin/ips", s.handleAdminSubmitIPs) + mux.HandleFunc("POST /api/v1/admin/ips/scan", s.handleAdminScanFloatingIPs) mux.HandleFunc("GET /api/v1/admin/ips/{ip}", s.handleAdminIPDetail) mux.HandleFunc("POST /api/v1/admin/ips/{ip}/cancel", s.handleAdminCancelIP) mux.HandleFunc("DELETE /api/v1/admin/ips/{ip}", s.handleAdminDeleteIP) @@ -28,6 +29,9 @@ func (s *Server) routes(mux *http.ServeMux) { mux.HandleFunc("POST /api/v1/admin/ips/clear", s.handleAdminClearQueue) mux.HandleFunc("GET /api/v1/admin/validators", s.handleAdminValidators) + mux.HandleFunc("GET /api/v1/admin/registry", s.handleAdminRegistry) + mux.HandleFunc("GET /api/v1/admin/registry/{ip}", s.handleAdminRegistryHistory) + mux.HandleFunc("GET /api/v1/admin/config/validators", s.handleConfigListValidators) mux.HandleFunc("POST /api/v1/admin/config/validators", s.handleConfigCreateValidator) mux.HandleFunc("PUT /api/v1/admin/config/validators/{id}", s.handleConfigUpdateValidator) diff --git a/internal/openstack/client.go b/internal/openstack/client.go index ddf9aec..490b4e0 100644 --- a/internal/openstack/client.go +++ b/internal/openstack/client.go @@ -150,6 +150,22 @@ func (c *Client) GetFloatingIPByAddress(ctx context.Context, address string) (*F return &FloatingIP{ID: f.ID, Address: f.FloatingIP, PortID: f.PortID, ProjectID: f.TenantID}, nil } +func (c *Client) ListFloatingIPs(ctx context.Context) ([]FloatingIP, error) { + pages, err := floatingips.List(c.networking, floatingips.ListOpts{}).AllPages(ctx) + if err != nil { + return nil, fmt.Errorf("openstack: list floating ips: %w", err) + } + list, err := floatingips.ExtractFloatingIPs(pages) + if err != nil { + return nil, fmt.Errorf("openstack: extract floating ips: %w", err) + } + out := make([]FloatingIP, 0, len(list)) + for _, f := range list { + out = append(out, FloatingIP{ID: f.ID, Address: f.FloatingIP, PortID: f.PortID, ProjectID: f.TenantID}) + } + return out, nil +} + func (c *Client) AssociateFloatingIP(ctx context.Context, fipID, portID string) error { _, err := floatingips.Update(ctx, c.networking, fipID, floatingips.UpdateOpts{ PortID: &portID, diff --git a/internal/openstack/interface.go b/internal/openstack/interface.go index 4df55c3..f7ce031 100644 --- a/internal/openstack/interface.go +++ b/internal/openstack/interface.go @@ -24,6 +24,12 @@ type FloatingIPClient interface { // is registered in the service project. GetFloatingIPByAddress(ctx context.Context, address string) (*FloatingIP, error) + // ListFloatingIPs returns every floating IP registered in the service + // project, both associated and free — filtering to the free pool (empty + // PortID) is the caller's job. Used to discover addresses to feed into + // the check queue without an admin having to enumerate them by hand. + ListFloatingIPs(ctx context.Context) ([]FloatingIP, error) + // AssociateFloatingIP attaches the floating IP to the given Neutron // port (the validator's primary NIC port). AssociateFloatingIP(ctx context.Context, fipID, portID string) error diff --git a/internal/openstack/mock.go b/internal/openstack/mock.go index 6502514..d83e491 100644 --- a/internal/openstack/mock.go +++ b/internal/openstack/mock.go @@ -56,6 +56,16 @@ func (m *MockClient) GetFloatingIPByAddress(ctx context.Context, address string) return &f, nil } +func (m *MockClient) ListFloatingIPs(ctx context.Context) ([]FloatingIP, error) { + m.mu.Lock() + defer m.mu.Unlock() + out := make([]FloatingIP, 0, len(m.fips)) + for _, f := range m.fips { + out = append(out, *f) + } + return out, nil +} + func (m *MockClient) AssociateFloatingIP(ctx context.Context, fipID, portID string) error { m.mu.Lock() defer m.mu.Unlock() diff --git a/internal/orchestrator/orchestrator.go b/internal/orchestrator/orchestrator.go index 836e2bd..1759cb8 100644 --- a/internal/orchestrator/orchestrator.go +++ b/internal/orchestrator/orchestrator.go @@ -395,6 +395,43 @@ func (o *Orchestrator) ClearQueue(ctx context.Context) (db.DeleteIPsResult, erro return result, nil } +// ScanFloatingIPs lists every floating IP in the OpenStack project, filters +// to the ones not currently associated to any port (the free pool awaiting +// validation before reissue), and submits that address list to the check +// queue via db.SubmitIPs — the same entry point the admin API's "add +// addresses" call uses, so add/requeue/reorder semantics are identical +// whether the address list came from an operator or from this scan. Returns +// the SubmitIPs outcome plus how many free floating IPs were found in total +// (which can be larger than the sum of the SubmitIPsResult slices, since +// addresses already mid-check are silently skipped — see db.SubmitIPs). +func (o *Orchestrator) ScanFloatingIPs(ctx context.Context) (db.SubmitIPsResult, int, error) { + fips, err := o.OS.ListFloatingIPs(ctx) + if err != nil { + return db.SubmitIPsResult{}, 0, fmt.Errorf("list floating ips: %w", err) + } + + var free []string + for _, f := range fips { + if f.PortID == "" { + free = append(free, f.Address) + } + } + + if len(free) == 0 { + o.event(ctx, "control-api", "", nil, "fip_scan", `{"scanned_free":0}`) + return db.SubmitIPsResult{}, 0, nil + } + + result, err := o.DB.SubmitIPs(ctx, free) + if err != nil { + return result, len(free), fmt.Errorf("submit scanned ips: %w", err) + } + o.event(ctx, "control-api", "", nil, "fip_scan", fmt.Sprintf( + `{"scanned_free":%d,"added":%d,"requeued":%d,"reordered":%d,"skipped_in_progress":%d}`, + len(free), len(result.Added), len(result.Requeued), len(result.Reordered), len(result.SkippedInProgress))) + return result, len(free), nil +} + // deletedAddressesPayload builds the event payload for the batch delete // operations — a proper JSON array via encoding/json rather than fmt's %q // slice formatting (which produces space-separated quoted strings, not @@ -508,6 +545,14 @@ func (o *Orchestrator) aggregateAndRelease(ctx context.Context, item db.IPQueueI o.event(ctx, "control-api", "", &item.ID, "aggregated", fmt.Sprintf(`{"result":%q,"checks":%d,"passed":%d,"missing":%d}`, result, len(checks), passCount, missing)) + if settings, err := o.DB.GetSettings(ctx); err != nil { + o.Log.Error("get settings for history retention", "ip_id", item.ID, "err", err) + } else if settings.HistoryRetentionCycles > 0 { + if err := o.DB.PruneRegistryHistory(ctx, item.RegistryID, settings.HistoryRetentionCycles); err != nil { + o.Log.Error("prune registry history", "ip_id", item.ID, "registry_id", item.RegistryID, "err", err) + } + } + if item.FIPID != "" { if err := o.OS.DisassociateFloatingIP(ctx, item.FIPID); err != nil { o.Log.Error("disassociate fip", "ip_id", item.ID, "fip_id", item.FIPID, "err", err) diff --git a/internal/orchestrator/orchestrator_test.go b/internal/orchestrator/orchestrator_test.go index 585a037..0591b2b 100644 --- a/internal/orchestrator/orchestrator_test.go +++ b/internal/orchestrator/orchestrator_test.go @@ -882,3 +882,60 @@ func TestAggregationWaitsForTLSAndSSHResults(t *testing.T) { t.Fatalf("expected partial (missing ssh/tls-443 counted against it), got %s", ip.OverallResult) } } + +// TestScanFloatingIPsSubmitsOnlyFreeAddresses proves ScanFloatingIPs filters +// out floating IPs already associated to a port and only submits the free +// pool to the check queue. +func TestScanFloatingIPsSubmitsOnlyFreeAddresses(t *testing.T) { + ctx := context.Background() + o, d, mock := newTestOrchestrator(t, 180) + mock.Seed("fip-free-1", "1.1.1.1", "svc-project") + mock.Seed("fip-free-2", "2.2.2.2", "svc-project") + mock.SeedWithPort("fip-occupied", "3.3.3.3", "svc-project", "some-other-port") + + result, scanned, err := o.ScanFloatingIPs(ctx) + if err != nil { + t.Fatalf("scan floating ips: %v", err) + } + if scanned != 2 { + t.Fatalf("expected 2 free fips scanned, got %d", scanned) + } + if len(result.Added) != 2 { + t.Fatalf("expected 2 addresses added, got %+v", result) + } + + ips, err := d.ListIPs(ctx) + if err != nil { + t.Fatalf("list ips: %v", err) + } + seen := map[string]bool{} + for _, ip := range ips { + seen[ip.IPAddress] = true + } + if !seen["1.1.1.1"] || !seen["2.2.2.2"] { + t.Fatalf("expected free addresses queued, got %+v", ips) + } + if seen["3.3.3.3"] { + t.Fatalf("expected occupied address not queued, got %+v", ips) + } +} + +// TestScanFloatingIPsNoFreeAddressesIsNotAnError proves scanning a project +// with no free floating IPs (or none at all) succeeds with an empty result +// rather than hitting SubmitIPs' "addresses must not be empty" validation. +func TestScanFloatingIPsNoFreeAddressesIsNotAnError(t *testing.T) { + ctx := context.Background() + o, _, mock := newTestOrchestrator(t, 180) + mock.SeedWithPort("fip-occupied", "3.3.3.3", "svc-project", "some-other-port") + + result, scanned, err := o.ScanFloatingIPs(ctx) + if err != nil { + t.Fatalf("scan floating ips: %v", err) + } + if scanned != 0 { + t.Fatalf("expected 0 free fips scanned, got %d", scanned) + } + if len(result.Added) != 0 { + t.Fatalf("expected nothing added, got %+v", result) + } +} diff --git a/rxprod-compose/sources/control-api.example.yaml b/rxprod-compose/sources/control-api.example.yaml index d6c6669..f9f2a1f 100644 --- a/rxprod-compose/sources/control-api.example.yaml +++ b/rxprod-compose/sources/control-api.example.yaml @@ -54,6 +54,12 @@ orchestrator: # /settings page) and this field is ignored. Must satisfy # fip_settle_seconds + self_check_timeout_seconds < lease_ttl_seconds. fip_settle_seconds: 0 + # How often (seconds) to automatically scan the OpenStack project for free + # (unassociated) floating IPs and submit them to the check queue. 0 (the + # default) disables periodic scanning — an operator can still trigger a + # scan on demand via POST /api/v1/admin/ips/scan or the dashboard's + # "Scan Floating IPs" button. + fip_scan_interval_seconds: 0 aggregation: missing_counts_as_fail: true