diff --git a/README.md b/README.md index 4f14a6e..4ac375b 100644 --- a/README.md +++ b/README.md @@ -38,6 +38,8 @@ docker compose up -d --build # весь стенд на одной | `orchestrator.heartbeat_timeout_seconds` | После скольких секунд тишины валидатор или площадка считаются потерянными (30) | | `orchestrator.fip_settle_seconds` | Только начальное значение: пауза между привязкой FIP и self-check; далее управляется на лету. Должно выполняться `fip_settle_seconds + self_check_timeout_seconds < lease_ttl_seconds` | | `orchestrator.fip_scan_interval_seconds` | Периодический скан Floating IP (0 — выключен). Не выполняется, пока включён автоматический цикл | +| `orchestrator.fip_scan_timeout_seconds` | Предел всего сканирования Floating IP (1800) | +| `openstack.list_page_size`, `list_page_retries`, `request_timeout_seconds` | Скан читает Floating IP страницами (200), повторяя страницу при обрыве/5xx/429 (5 раз); таймаут одного запроса к OpenStack (60 с) | | `auth.admin_token_env`, `auth.agent_token_env` | Имена переменных окружения с токеном администратора (`CONTROL_API_ADMIN_TOKEN`) и токеном агентов (`CONTROL_API_AGENT_TOKEN`). Значения в YAML не хранятся; пустой токен — соответствующий уровень API открыт (с предупреждением в логе) | | `aggregation.missing_counts_as_fail` | Отсутствие ответа источника засчитывается как провал (`true`) | | `validators`, `sites`, `check_types`, `targets`, `inbound_checks` | Начальная загрузка пустой БД: валидаторы (`validator_id` + `os_port_id`), внешние площадки, типы и цели egress-проверок, порты inbound-проверок. Дальше источник истины — БД, правки через API/UI | @@ -103,22 +105,23 @@ docs/ документация и планы доработок - История проверок хранится по циклам; `history_retention_cycles` ограничивает глубину (0 — без ограничения), сама запись реестра остаётся. **Автоматический цикл** (опционально, по умолчанию выключен) -- Один цикл: очистить очередь → просканировать Floating IP → дождаться, пока все адреса станут терминальными → пауза `interval_seconds` → заново. Пауза считается от завершения цикла. +- Один цикл: очистить очередь → просканировать Floating IP (фаза `scanning`, фоновое задание) → дождаться, пока все адреса станут терминальными → пауза `interval_seconds` → заново. Пауза считается от завершения цикла. Сканирование не блокирует оркестратор; чужое идущее сканирование цикл дожидается, а не присоединяется к нему. - `interval_seconds` — по умолчанию 3600, минимум 60; `max_run_seconds` — максимальное ожидание проверок (0 — без лимита), по истечении исход `timeout`. - Исходы цикла: `completed`, `no_free_ips`, `timeout`, `error`, `stopped`. Опустевшая очередь посреди цикла считается завершением. - Состояние хранится в БД и переживает перезапуск; выключение не прерывает идущие проверки. Подробности — [docs/USAGE.md](docs/USAGE.md#автоматический-цикл-проверок). **Удаление и сканирование** - Удаление (точечное, списком, «очистить всё») убирает строку очереди, но не историю в реестре. -- Скан Floating IP ставит в очередь только свободные адреса (не привязанные ни к одному порту); уже идущие проверки не трогаются. +- Скан Floating IP ставит в очередь только свободные адреса (не привязанные ни к одному порту); уже идущие проверки не трогаются. Он идёт **в фоне и читает облако страницами**: подходит и для тысяч адресов (на стенде 6441 Floating IP читаются ≈ 1,5–2 мин). Адреса ставятся в очередь только после полного обнаружения (кусками по 500, по возрастанию IP); при сбое чтения очередь не меняется. `dry_run=true` / «Пробное сканирование» считает адреса, не меняя очередь. +- Пропускная способность: ≈ 50 с на адрес на валидатор — очередь из 6440 адресов это ≈ 18 ч на 5 валидаторах, ≈ 9 ч на 10 (рычаги: число валидаторов и `fip_settle_seconds`). ## API (`/api/v1`) | Область | Эндпоинты | |---|---| | Валидатор | `POST /agents/register`, `POST /agents/{id}/heartbeat`, `GET /agents/{id}/assignment`, `POST /agents/{id}/self-check\|events\|results\|complete` | | Пробер | `POST /probers/register`, `POST /probers/{site_id}/heartbeat`, `GET /probers/{site_id}/assignments`, `POST /probers/{site_id}/results` | -| Очередь | `GET /admin/status`, `GET\|POST /admin/ips`, `GET /admin/ips/{ip}`, `POST /admin/ips/{ip}/cancel`, `DELETE /admin/ips/{ip}`, `POST /admin/ips/delete\|clear\|scan` | -| Реестр | `GET /admin/registry`, `GET /admin/registry/{ip}` | +| Очередь | `GET /admin/status`, `GET\|POST /admin/ips` (`limit/offset/state/q/result/order` — постранично), `GET /admin/ips/{ip}`, `POST /admin/ips/{ip}/cancel`, `DELETE /admin/ips/{ip}`, `POST /admin/ips/delete\|clear`, `POST\|GET /admin/ips/scan` (фоновый скан: `202`, `dry_run`, `wait`; статус и прогресс) | +| Реестр | `GET /admin/registry` (`limit/offset/q/last_result` — постранично), `GET /admin/registry/{ip}` | | Автоцикл | `GET\|PUT /admin/auto-cycle`, `POST /admin/auto-cycle/start\|stop` | | Конфигурация | `/admin/config/validators`, `/sites`, `/targets`, `/check-types`, `GET\|PUT /admin/config/orchestrator`, `GET\|PUT /admin/config/inbound-checks` | | Служебное | `GET /admin/validators`, `GET /healthz` | @@ -155,9 +158,9 @@ docs/ документация и планы доработок Страницы `admin-dashboard` (подробно — [docs/DASHBOARD.md](docs/DASHBOARD.md)): | Страница | Назначение | |---|---| -| `/overview` | Счётчики по состояниям, «текущая» и «последние завершённые» проверки, поиск по IP и фильтр по статусу, индикатор автоцикла; обновляется без перезагрузки | -| `/ips`, `/ips/{ip}` | Очередь: добавление адресов, «Сканировать Floating IP», перепроверка, отмена, удаление (в том числе списком и «Очистить всё»); детали и события адреса | -| `/registry`, `/registry/{ip}` | Реестр всех адресов и полная история проверок адреса; поиск и фильтр сохраняются в адресной строке | +| `/overview` | Счётчики и прогресс («Готово D из T», оценка времени), «в работе», «в очереди: Q», «последние завершённые», поиск по IP и фильтр по статусу, индикатор скана и автоцикла; работает на счётчиках и ограниченных списках, поэтому быстрый и при тысячах адресов | +| `/ips`, `/ips/{ip}` | Очередь **постранично** с поиском и фильтром на сервере: добавление адресов, «Сканировать Floating IP» (панель прогресса) и «Пробное сканирование», перепроверка, отмена, удаление (страница или «все N по фильтру», «Очистить всё»); детали и события адреса | +| `/registry`, `/registry/{ip}` | Реестр всех адресов (постранично) и полная история проверок адреса; поиск, фильтр и страница сохраняются в адресной строке | | `/validators`, `/sites`, `/targets`, `/check-types` | Управление валидаторами, внешними площадками, группами целей и типами проверок | | `/settings` | Панель «Автоматический цикл», пауза перед self-check, глубина истории, TCP-порты и ICMP для inbound-проверок | - Порядок блоков на `/overview` фиксирован: статистика → фильтр → таблицы; поллится только блок таблиц, поэтому набранный в фильтре текст не сбрасывается. @@ -181,7 +184,7 @@ scripts/run-local-e2e.sh # сквозной прог | [docs/DASHBOARD.md](docs/DASHBOARD.md) | Устройство `admin-dashboard`: страницы, поиск и фильтр, обработка ошибок | | [docs/DIAGRAMS.md](docs/DIAGRAMS.md) | Диаграммы потоков данных: control plane, egress-проверка, телеметрия | | [docs/LOCAL_E2E.md](docs/LOCAL_E2E.md) | Полностью офлайн-прогон всей системы одним скриптом | -| [docs/changes/](docs/changes/) | Планы доработок и отчёты ревью с отметкой времени в имени файла (последняя: [аутентификация](docs/changes/2026-10-01_11-31_authentication-review.md)) | +| [docs/changes/](docs/changes/) | Планы доработок и отчёты ревью с отметкой времени в имени файла (последняя: [скан при тысячах адресов](docs/changes/2026-10-01_18-59_fip-scan-at-scale-review.md)) | | [docs/CONTROL_DATA_PLANE.html](docs/CONTROL_DATA_PLANE.html) | Презентационные схемы control/data plane для docker-compose-деплоя — открыть в браузере | ## История изменений @@ -197,6 +200,7 @@ scripts/run-local-e2e.sh # сквозной прог | Дата | Веха | Документ | |---|---|---| +| 2026-10-01 | Скан Floating IP при тысячах адресов: фоновый постраничный скан с прогрессом, фаза `scanning` в автоцикле, постраничные `/ips` и `/registry`, «Обзор» на счётчиках | [план](docs/changes/2026-10-01_18-19_fip-scan-at-scale-plan.md) · [ревью и тесты](docs/changes/2026-10-01_18-59_fip-scan-at-scale-review.md) · [USAGE](docs/USAGE.md#сканирование-floating-ip-из-openstack) · [API](docs/API.md#post-apiv1adminipsscan) | | 2026-10-01 | Аутентификация: токены администратора и агентов для API, логин и пароль для дашборда | [план](docs/changes/2026-10-01_11-12_authentication-plan.md) · [ревью и тесты](docs/changes/2026-10-01_11-31_authentication-review.md) · [API](docs/API.md#аутентификация) | | 2026-10-01 | Автоматический цикл проверок по сценарию: очистка → скан FIP → проверка → пауза | [USAGE](docs/USAGE.md#автоматический-цикл-проверок) · [API](docs/API.md#автоматический-цикл-проверок) | | 2026-09-23 | Сканирование Floating IP и устойчивый реестр адресов с настраиваемой глубиной истории | [USAGE](docs/USAGE.md#реестр-адресов-и-глубина-истории) | diff --git a/bin/SHA256SUMS b/bin/SHA256SUMS index 1035e4b..3c58bfb 100644 --- a/bin/SHA256SUMS +++ b/bin/SHA256SUMS @@ -1,4 +1,4 @@ -a0eac5409a2929d641ef2a217c31f1b6a974a8679da866ebbb50ff1d5b57df8d control-api -5d7fb2f476871843ae4179dddb61e77d3f71e8eb73a00e331f0119d8893ba813 validator-agent -23c687a988350c61f486d400af654bd563865fff59421c437eb679e0e4e10500 prober -90921ad25f16a922124368770e439d1228c5b0bf7126ba6c726f3fb98bb22f7a admin-dashboard +26ff297b66a0c974b7143179e44a58f476265482cd1929bc01a95ad29e50ca7a control-api +091ad94b5706b1778b181542de7a4b251cd16c421144269d0955769603d51e06 validator-agent +3e9e14dbb361ee76aaad7c1da6864b3ea111e0ed151403f904b12485631bbf75 prober +d72234688eb1954dbff420ca1c8e83b83ceab561b18a0336af0f73fe4d9ac8be admin-dashboard diff --git a/bin/admin-dashboard b/bin/admin-dashboard index 924bf0d..bcb5ee5 100755 Binary files a/bin/admin-dashboard and b/bin/admin-dashboard differ diff --git a/bin/control-api b/bin/control-api index 1cfb7e7..bac8e17 100755 Binary files a/bin/control-api and b/bin/control-api differ diff --git a/bin/prober b/bin/prober index 3387f97..9d03692 100755 Binary files a/bin/prober and b/bin/prober differ diff --git a/bin/validator-agent b/bin/validator-agent index b8da4ad..0aec448 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 eb6ef11..aa3f4c3 100644 --- a/cmd/control-api/main.go +++ b/cmd/control-api/main.go @@ -59,6 +59,9 @@ func run(configPath string, log *slog.Logger) error { } orch := orchestrator.New(database, osClient, cfg, log) + // Background jobs (the floating-IP scan) live as long as the process, not + // as long as the HTTP request or loop iteration that started them. + orch.SetContext(ctx) adminToken := os.Getenv(cfg.Auth.AdminTokenEnv) agentToken := os.Getenv(cfg.Auth.AgentTokenEnv) @@ -135,8 +138,10 @@ func runOrchestratorLoop(ctx context.Context, orch *orchestrator.Orchestrator, c } else if ac.Enabled { continue } - if _, _, err := orch.ScanFloatingIPs(ctx); err != nil { - log.Error("scan floating ips", "err", err) + // Non-blocking: the scan runs in the background (single-flight, so + // a still-running scan is simply joined) and must not stall Tick. + if st, started := orch.StartScan(orchestrator.ScanOptions{}); !started { + log.Info("periodic floating ip scan skipped: a scan is already running", "state", st.State) } } } @@ -158,11 +163,18 @@ func newOpenStackClient(ctx context.Context, cfg *config.ControlAPI) (openstack. } func newRealOpenStackClient(ctx context.Context, cfg *config.ControlAPI) (openstack.FloatingIPClient, error) { + retries := cfg.OpenStack.ListPageRetries + if retries < 0 { + retries = 0 // negative in the config disables retries + } clientCfg := openstack.ClientConfig{ AuthURL: os.Getenv(cfg.OpenStack.AuthURLEnv), ProjectID: os.Getenv(cfg.OpenStack.ProjectIDEnv), Region: os.Getenv(cfg.OpenStack.RegionEnv), Interface: os.Getenv(cfg.OpenStack.InterfaceEnv), + + RequestTimeout: time.Duration(cfg.OpenStack.RequestTimeoutSeconds) * time.Second, + ListPageRetries: retries, } switch cfg.OpenStack.AuthMethod { diff --git a/configs/control-api.example.yaml b/configs/control-api.example.yaml index 7edafb7..caefba4 100644 --- a/configs/control-api.example.yaml +++ b/configs/control-api.example.yaml @@ -37,6 +37,20 @@ openstack: user_domain_name_env: "OS_USER_DOMAIN_NAME" password_env: "OS_PASSWORD" + # Постраничное чтение Floating IP (скан при тысячах адресов): сколько + # адресов запрашивать у Neutron за один запрос. Default 200. + # Page size of the paged floating-IP listing. Default 200. + list_page_size: 200 + # Таймаут каждого HTTP-запроса к Keystone/Neutron, секунд. Default 60. + # Per-request HTTP timeout (also protects the orchestrator tick from a + # hung Neutron call). Default 60. + request_timeout_seconds: 60 + # Сколько раз повторять неудавшуюся страницу (сетевая ошибка, EOF/ + # RemoteDisconnected, 5xx, 429) с паузами 1,2,4,8,16 с. Default 5; + # отрицательное значение отключает повторы. + # Retries per failed listing page. Default 5; negative disables retries. + list_page_retries: 5 + # Аутентификация API: здесь только ИМЕНА переменных окружения, значения # (статические bearer-токены) задаются окружением процесса — см. # deploy/systemd/control-api.service (EnvironmentFile=). Генерация: @@ -76,6 +90,10 @@ orchestrator: # scan on demand via POST /api/v1/admin/ips/scan or the dashboard's # "Scan Floating IPs" button. fip_scan_interval_seconds: 0 + # Общий таймаут одного фонового скана Floating IP (очистка + чтение всех + # страниц + постановка в очередь), секунд. Default 1800. + # Overall deadline of one background floating-IP scan. Default 1800. + fip_scan_timeout_seconds: 1800 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 3e3d916..5e5034e 100644 --- a/deploy/docker/control-api/control-api.docker.example.yaml +++ b/deploy/docker/control-api/control-api.docker.example.yaml @@ -16,6 +16,8 @@ database: openstack: mode: "mock" + # Постраничное чтение Floating IP / page size of the paged listing (default 200) + list_page_size: 200 # Статические bearer-токены читаются из переменных окружения контейнера с # этими именами (задаются в .env, см. deploy/docker/.env.example). Пусто = @@ -34,6 +36,7 @@ orchestrator: heartbeat_timeout_seconds: 30 fip_settle_seconds: 0 fip_scan_interval_seconds: 0 + fip_scan_timeout_seconds: 1800 # общий таймаут фонового скана / scan deadline aggregation: missing_counts_as_fail: true diff --git a/docs/API.md b/docs/API.md index a49ed79..96e2210 100644 --- a/docs/API.md +++ b/docs/API.md @@ -330,15 +330,37 @@ IP на данном проходе". До этого момента control-api { "total_ips": 25, "ips_by_state": {"queued": 10, "checking": 3, "done": 11, "failed": 1}, + "results_by_overall": {"pass": 8, "partial": 3, "fail": 0, "cancelled": 0}, "total_validators": 4 } ``` +`results_by_overall` — сколько адресов с каким итогом (всегда все четыре ключа). Счётчики считаются +запросами `GROUP BY` на стороне БД, а не загрузкой всей очереди, поэтому метод быстрый и при тысячах адресов. + ### `GET /api/v1/admin/ips` -Полный список всех IP из очереди со всеми полями (см. +Список IP из очереди со всеми полями (см. [USAGE.md](USAGE.md#значения-полей-ip) — расшифровка полей и статусов). +**Без параметров** — как раньше: весь список одним массивом (при тысячах адресов это мегабайты — для больших очередей +используйте постраничный режим). **С `limit`** — постраничный режим: ответ — конверт + +```json +{"items": [ ... ], "total": 6440, "limit": 50, "offset": 0} +``` + +| Параметр | Значение | +|---|---| +| `limit` | размер страницы, `1`…`1000` (иначе `400`); включает постраничный режим | +| `offset` | смещение, `>= 0` (без `limit` — `400`) | +| `state` | одно или несколько состояний через запятую (`queued`, `assigning_fip`, `awaiting_self_check`, `checking`, `aggregating`, `done`, `failed`, `occupied`) | +| `q` | подстрока адреса | +| `result` | итог: `pass`, `partial`, `fail`, `cancelled` | +| `order` | `sequence` (по умолчанию, порядок очереди) или `aggregated_at_desc` (последние завершённые) | + +`total` — число записей после фильтров. Параметры фильтров без `limit` возвращают отфильтрованный массив. + ### `GET /api/v1/admin/ips/{ip}` Детали по одному адресу: сам объект IP, все проверки текущей попытки и @@ -470,33 +492,52 @@ YAML для этой секции больше не перечитывается ### `POST /api/v1/admin/ips/scan` -Сканирует текущий проект OpenStack на предмет свободных (не привязанных ни -к одному порту) Floating IP и сразу передаёт найденный список в `POST -/api/v1/admin/ips` — тот же add/requeue/reorder-вызов, как если бы -оператор ввёл эти адреса вручную. Не принимает тело запроса. +Запускает **фоновое** сканирование проекта OpenStack: находит все свободные (не привязанные ни к одному порту) Floating IP +и ставит их в очередь — тот же add/requeue/reorder, что и `POST /api/v1/admin/ips`. Не принимает тело запроса и **сразу отвечает** +`202` со статусом задания; ход сканирования смотрите через `GET /api/v1/admin/ips/scan`. + +Почему в фоне: в проекте может быть тысячи Floating IP (на стенде — около 6,4 тыс.), Neutron отдаёт такой список минуты. Control-api читает +его **страницами** (по `openstack.list_page_size`, по умолчанию 200, по `marker`), повторяет страницу при обрыве соединения/5xx/429, +сначала обнаруживает **все** адреса и только потом ставит их в очередь кусками по 500 в порядке возрастания IP. Если чтение не удалось +(после повторов), в очередь не попадает ничего — очередь остаётся как была, а статус задания — `error`. + +| Параметр | Значение | +|---|---| +| `dry_run=true` | только найти и посчитать свободные адреса; очередь не меняется (безопасная проверка, итог — в статусе) | +| `wait=true` | дождаться окончания и ответить `200` прежним телом `{scanned_free, added[], requeued[], reordered[], skipped_in_progress[]}` (для curl и скриптов; при ошибке `502`) | + +Одновременно идёт одно сканирование: повторный запрос во время работы **присоединяется** к текущему и тоже отвечает `202` с его статусом. + +### `GET /api/v1/admin/ips/scan` + +Статус и прогресс сканирования (admin-токен). -Ответ (`200`): ```json { - "scanned_free": 3, - "added": ["203.0.113.20"], - "requeued": [], - "reordered": ["203.0.113.10", "203.0.113.11"], - "skipped_in_progress": [] + "state": "listing", + "running": true, + "dry_run": false, + "pages": 12, + "discovered": 2400, + "free": 2399, + "added": 0, + "requeued": 0, + "reordered": 0, + "skipped_in_progress": 0, + "started_at": "2026-10-01T15:47:40.759Z", + "finished_at": null, + "error": "" } ``` -`scanned_free` — сколько свободных Floating IP нашлось в проекте всего -(включая уже стоящие в очереди — они попадут в `reordered`, а не -`added`). Если свободных адресов нет вообще, это не ошибка: ответ будет -`{"scanned_free": 0, "added": [], ...}`. +`state`: `idle` (в этом процессе сканирования ещё не было), `clearing` (очистка очереди — только в автоцикле), `listing` (чтение страниц), +`enqueuing` (постановка в очередь), `done`, `error` (причина в `error`), `cancelled`. `discovered` — сколько Floating IP прочитано +(свободных и занятых), `free` — из них свободных, `added`/`requeued`/`reordered`/`skipped_in_progress` — итог постановки в очередь +(как в `POST /admin/ips`). Статус хранится в памяти процесса: после перезапуска control-api он снова `idle`. -Помимо ручного вызова, сканирование можно включить по расписанию — -`orchestrator.fip_scan_interval_seconds` в `control-api.yaml` (0, по -умолчанию, — только по запросу через эту ручку или кнопку «Сканировать -Floating IP» в дашборде). Пока включён -[автоматический цикл](#автоматический-цикл-проверок), периодический скан -не выполняется. +Помимо ручного вызова, сканирование можно включить по расписанию — `orchestrator.fip_scan_interval_seconds` в `control-api.yaml` +(0, по умолчанию, — только по запросу через эту ручку или кнопку «Сканировать Floating IP» в дашборде). Пока включён +[автоматический цикл](#автоматический-цикл-проверок), периодический скан не выполняется. ## Автоматический цикл проверок @@ -575,7 +616,10 @@ curl -s -X POST http://:8080/api/v1/admin/auto-cycle/stop ### `GET /api/v1/admin/registry` -Список всех адресов реестра с краткой сводкой по каждому. +Список адресов реестра с краткой сводкой по каждому. Без параметров — все адреса одним массивом; **с `limit`** (`1`…`1000`) — +постраничный конверт `{"items": [...], "total": N, "limit": L, "offset": O}`, параметры `offset`, `q` (подстрока адреса) и +`last_result` (`pass`/`partial`/`fail`/`cancelled`). Страница и фильтры применяются в SQL до расчёта сводки, поэтому +реестр из тысяч адресов отдаётся за доли секунды. ```json [ diff --git a/docs/DASHBOARD.md b/docs/DASHBOARD.md index 0c4cbd2..6e2f119 100644 --- a/docs/DASHBOARD.md +++ b/docs/DASHBOARD.md @@ -49,7 +49,7 @@ 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` секунд без перезагрузки страницы. Поиск по IP и фильтр по статусу (`pass`/`partial`/`fail`/`cancelled`) над обеими таблицами — набранное/выбранное не сбрасывается очередным обновлением. Пока включён [автоматический цикл](USAGE.md#автоматический-цикл-проверок), под счётчиками показывается индикатор «Автоцикл активен» с текущей фазой и временем следующего запуска; управляется цикл на `/settings`. | -| `/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` | Очередь **постранично** (по 50 адресов; 25/50/100/200) с поиском по IP и фильтром по состоянию/итогу на сервере; кнопка «Сканировать Floating IP» запускает фоновое сканирование с панелью прогресса, «Пробное сканирование» ничего не ставит в очередь (подробности — «Очередь из тысяч адресов» ниже). Форма сверху принимает список адресов (по одному на строке или через запятую) и отправляет их в `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` и повторное добавление того же адреса позже (см. «Реестр адресов» ниже). Поиск по IP и фильтр по статусу — то же самое, что на `/overview`, плюс отражается в адресной строке (`?q=&status=`), так что отфильтрованную ссылку можно сохранить/переслать. | | `/registry/{ip}` | Полная сохранённая история проверок одного адреса по всем циклам (не только текущему) — в отличие от `/ips/{ip}`, которая показывает только текущую попытку. | @@ -59,41 +59,37 @@ admin-dashboard -config /etc/cloud-ip-validator/admin-dashboard.yaml | `/check-types` | Типы проверок (`https`/`icmp`/`ssh`/...), включение/выключение, привязка к группам целей. | | `/settings` | Четыре блока. Первый — панель **«Автоматический цикл»**: статус и фаза, время последнего/следующего запуска, результат последнего цикла, поля «Интервал между циклами (мин)» и «Максимальная длительность проверки (мин, 0 = без лимита)» с кнопкой «Сохранить» и кнопка «Включить»/«Выключить» (показывается та, что сейчас применима). Значения вводятся в минутах (допустимы дробные), в control-api уходят секундами; минимум интервала — 1 минута (`60` с), нарушение приходит предупреждением в баннере. Подробности — [USAGE.md](USAGE.md#автоматический-цикл-проверок), API — [API.md](API.md#автоматический-цикл-проверок). Далее три формы: `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#управление-типами-проверок-пробера)). | -### «Текущая» и «последняя завершённая» проверка +### «В работе», «в очереди» и «последняя завершённая» проверка В `control-api` нет понятия «запуска»/«цикла проверки» как отдельной сущности — есть только общая очередь IP-адресов (`docs/PLAN_ADMIN_DASHBOARD.md`). Дашборд ничего не меняет в этом -устройстве и не заводит своего состояния: +устройстве и не заводит своего состояния. Очередь может содержать тысячи +адресов, поэтому `/overview` **никогда не загружает её целиком** — на каждое +обновление запрашиваются счётчики и несколько ограниченных списков: -- **Текущая проверка** — все адреса, которые прямо сейчас не в - состоянии `done`/`failed` (`queued`, `assigning_fip`, - `awaiting_self_check`, `checking`, `aggregating`), вычисляется заново на - каждый запрос из `GET /api/v1/admin/status` + `GET /api/v1/admin/ips`. -- **Последняя завершённая проверка** — последние N адресов, перешедших в - `done`/`failed`, отсортированные по `AggregatedAt` по убыванию (не - «последний запуск», а именно скользящее окно последних по времени - завершений). +- **Счётчики** — `GET /api/v1/admin/status` (по состояниям и `results_by_overall`). +- **В работе** — адреса в `assigning_fip`, `awaiting_self_check`, `checking`, `aggregating` + (не более 100; естественный предел — число валидаторов). +- **В очереди: Q** — счётчик `queued` со ссылкой на `/ips?state=queued` и несколько ближайших адресов. +- **Последние N завершённых** — последние N адресов в `done`/`failed` по `AggregatedAt` по убыванию + (не «последний запуск», а скользящее окно последних по времени завершений). +- **Прогресс** в блоке статистики: «Готово D из T (P%) · в работе A · в очереди Q» с полосой и оценкой + оставшегося времени (по скорости последних завершений, когда их не меньше пяти). `occupied` считается + завершённым состоянием. ### Поиск по IP и фильтр по статусу -На `/overview` и `/registry` есть форма из двух полей — поиск по IP -(подстрока, без учёта регистра) и выпадающий список статуса -(`pass`/`partial`/`fail`/`cancelled`). Оба поля работают вместе (И, а не -ИЛИ) и применяются целиком на стороне дашборда — `client.ListIPs`/ -`client.ListRegistry` всегда получают от `control-api` полный список, -`internal/httpapi`/`internal/db` про фильтр вообще не знают. +На `/overview`, `/ips` и `/registry` есть поиск по IP (подстрока) и фильтр по статусу/итогу (`pass`/`partial`/`fail`/`cancelled`; +на `/ips` — ещё по состоянию очереди). Поля работают вместе (И, а не ИЛИ). Фильтрация выполняется **на стороне `control-api`** +(параметры `q`, `state`, `result`/`last_result` у `GET /admin/ips` и `GET /admin/registry`), а дашборд получает только нужную страницу, поэтому +фильтр работает быстро при любом размере очереди. -- **`/overview`** — фильтр действует на обе таблицы сразу («Текущая - проверка» и «Последние N завершённых»). Статус — это фильтр по - итоговому результату (`OverallResult`), поэтому выбор конкретного - статуса скрывает «Текущую проверку» целиком: у ещё идущих проверок - результата попросту нет. Панель статистики (счётчики сверху) фильтру не - подчиняется — это агрегаты по всей очереди, а не по видимым строкам. -- **`/registry`** — тот же принцип, но по одной таблице (`LastResult`), и - значения полей отражаются в адресной строке (`?q=&status=`) через - `hx-replace-url` — отфильтрованную ссылку можно сохранить или переслать, - а обновление страницы (F5) сохраняет применённый фильтр. +- **`/overview`** — `q` и статус передаются в списки «В работе» и «Последние N завершённых». Статус — это фильтр по + итоговому результату (`OverallResult`), поэтому выбор конкретного статуса скрывает «В работе» и «В очереди»: у ещё идущих проверок + результата попросту нет. Панель статистики (счётчики сверху) фильтру не подчиняется — это агрегаты по всей очереди, а не по видимым строкам. +- **`/registry` и `/ips`** — таблица постраничная; значения полей и страница отражаются в адресной строке (`?q=&status=&page=`) через + `hx-replace-url` — отфильтрованную ссылку можно сохранить или переслать, а обновление страницы (F5) сохраняет применённый фильтр. **Раскладка `/overview` сверху вниз**: панель статистики → форма фильтра → таблицы. Панель статистики и форма фильтра физически лежат @@ -157,6 +153,23 @@ auto-refresh на `/ips`, см. git-историю). Опрашивается т циклов) остаётся всегда. Подробнее — [API.md](API.md#реестр-адресов-и-история-проверок). +## Очередь из тысяч адресов + +После сканирования проекта в очереди может оказаться несколько тысяч адресов, поэтому тяжёлые страницы работают постранично: + +- **`/ips` и `/registry`** — параметры `page` и `per_page` (по умолчанию 50; допустимо 25/50/100/200), «Показано a–b из N» и кнопки ‹ ›. + Поиск (`q`), состояние/итог и размер страницы применяются **на стороне control-api** (`GET /admin/ips?limit=…`, `GET /admin/registry?limit=…`), + поэтому страница весит десятки килобайт независимо от длины очереди. Фильтры и страница отражены в адресной строке. +- **Массовые операции.** Чекбоксы выбирают строки текущей страницы (счётчик «Выбрано на странице: k из P»). Если отмечен заголовок таблицы и записей + больше страницы, появляется ссылка «Выбрать все N по фильтру»: тогда «Перепроверить»/«Удалить» применяются ко **всем** адресам по текущему фильтру + (адреса разрешаются на сервере и отправляются кусками по 500). Подтверждения показывают реальное число: «Удалить ВСЕ 6440 адресов…». + «Очистить всё» очищает очередь целиком одной быстрой операцией. +- **Сканирование.** Кнопка «Сканировать Floating IP» мгновенно возвращает панель прогресса под кнопкой; пока задание идёт, панель сама + обновляется каждые 2 секунды, по окончании опрос прекращается и таблица перезагружается. Во время сканирования кнопки заблокированы; повторное + нажатие присоединяется к идущему заданию. Ошибка (например, OpenStack недоступен) показывается в панели с причиной. + Саму таблицу `/ips` по таймеру по-прежнему не обновляем — она не сбрасывает ввод оператора. +- Долгие операции (`Очистить всё`, массовое удаление/перепроверка) выполняются с увеличенным таймаутом (120 с), остальные запросы к control-api — с `control_api.timeout_seconds`. + ## Вход и сессия Если заданы `ADMIN_DASHBOARD_USERNAME` и `ADMIN_DASHBOARD_PASSWORD`, все страницы, кроме `/login` и `/static/*`, требуют входа. diff --git a/docs/LOCAL_E2E.md b/docs/LOCAL_E2E.md index db7c42c..ad33870 100644 --- a/docs/LOCAL_E2E.md +++ b/docs/LOCAL_E2E.md @@ -66,7 +66,7 @@ It will: outcome is `completed`, the phase is `waiting` with `runs_total=1`, and that the registry's `total_cycles` for `127.0.0.1` grew (the cycle cleared the queue, re-scanned the mock floating IP and re-checked it). - Finally it `stop`s the cycle and asserts it is `idle`. The second cycle + The cycle now passes through the background scan (`scanning` phase) before the checks run. Finally it `stop`s the cycle and asserts it is `idle`. The second cycle (the interval wait) is covered by unit tests, so the script does not sit through the 60s pause. The script exits non-zero if any assertion fails. diff --git a/docs/USAGE.md b/docs/USAGE.md index 4b49a5c..45772d0 100644 --- a/docs/USAGE.md +++ b/docs/USAGE.md @@ -90,7 +90,8 @@ curl -s -X POST http://:8080/api/v1/admin/ips \ самому найти их в облаке: ```bash -curl -s -X POST http://:8080/api/v1/admin/ips/scan +curl -s -X POST http://:8080/api/v1/admin/ips/scan # 202: сканирование запущено в фоне +curl -s http://:8080/api/v1/admin/ips/scan # ход и результат ``` Сканируются все Floating IP текущего проекта OpenStack, но в очередь @@ -101,8 +102,23 @@ curl -s -X POST http://:8080/api/v1/admin/ips/scan встают в очередь, уже завершённые перезапускаются, активно проверяемые не трогаются (см. [выше](#добавление-новых-ip-в-очередь)). -В `admin-dashboard` то же самое — кнопка «Сканировать Floating IP» на -странице `/ips`. +**Сколько адресов — не важно.** Сканирование работает в фоне и читает список из OpenStack +**страницами** (по 200 адресов, с повторами при обрывах), поэтому подходит и для проекта с +тысячами Floating IP: на стенде с 6441 адресом чтение занимает около 1,5–2 минут. Сначала +обнаруживаются **все** адреса, и только потом они ставятся в очередь (кусками по 500, в порядке +возрастания IP); сразу после этого начинаются проверки. Если чтение сорвалось даже после повторов, +в очередь не попадает ничего — очередь остаётся как была, а на панели виден `error` и причина. +Одновременно идёт одно сканирование: повторное нажатие присоединяется к текущему. + +В `admin-dashboard` — кнопка «Сканировать Floating IP» на странице `/ips`: она сразу отвечает, а под +кнопкой появляется панель прогресса (читаются страницы → ставятся в очередь → готово: прочитано +страниц, найдено, свободных, добавлено, время), по окончании таблица обновляется сама. Рядом — +«Пробное сканирование»: оно проходит все страницы и показывает, сколько свободных адресов нашлось, +**не меняя очередь** (удобно проверить, что облако отвечает и сколько адресов будет поставлено). + +Параметры чтения (`control-api.yaml`): `openstack.list_page_size` (200), `openstack.request_timeout_seconds` +(60 — таймаут одного запроса к OpenStack), `openstack.list_page_retries` (5), `orchestrator.fip_scan_timeout_seconds` +(1800 — предел всего сканирования). Если хочется, чтобы сканирование происходило само по расписанию, а не только по запросу — задайте `orchestrator.fip_scan_interval_seconds` @@ -111,6 +127,11 @@ curl -s -X POST http://:8080/api/v1/admin/ips/scan [автоматический цикл](#автоматический-цикл-проверок), это периодическое сканирование не выполняется — цикл сам управляет очередью. +> **Сколько займут проверки.** Один адрес занимает около 50 секунд на валидаторе (из них 30 с — пауза +> `fip_settle_seconds`). Поэтому очередь из 6440 адресов — примерно 18 часов на 5 валидаторах, +> 9 часов на 10, 4,5 часа на 20. Ускорить можно числом валидаторов и (осторожно) `fip_settle_seconds`; +> на странице «Обзор» виден прогресс «Готово D из T» и оценка оставшегося времени. + ## Автоматический цикл проверок Опциональный режим, который сам повторяет то, что оператор делает руками: @@ -122,7 +143,9 @@ curl -s -X POST http://:8080/api/v1/admin/ips/scan в [реестре](#реестр-адресов-и-глубина-истории) при этом сохраняется. 2. Control-api находит все свободные Floating IP и ставит их в очередь — то же, что «Сканировать Floating IP» (см. - [выше](#сканирование-floating-ip-из-openstack)). + [выше](#сканирование-floating-ip-из-openstack)). Шаги 1–2 выполняются одним + фоновым заданием, поэтому долгое чтение тысяч адресов не блокирует работу + оркестратора (назначение валидаторов, лизинги, heartbeat). 3. Проверки запускаются сами — как для любого адреса в очереди. 4. Цикл ждёт, пока **все** адреса очереди дойдут до конечного состояния (`done`, `failed` или `occupied`). К этому моменту результат каждого @@ -136,7 +159,7 @@ curl -s -X POST http://:8080/api/v1/admin/ips/scan | Параметр | По умолчанию | Смысл | |---|---|---| | `interval_seconds` | `3600` (1 час) | Пауза между циклами. Не меньше `60`: слишком частые сканы нагружают API OpenStack. | -| `max_run_seconds` | `0` (без лимита) | Сколько максимум ждать на шаге 4. По истечении цикл фиксирует `timeout` и переходит к паузе — защита от зависания (нет свободных валидаторов, недоступна площадка). Очередь при этом не трогается: следующий цикл её очистит, а до тех пор видно, что именно не дошло до конца. | +| `max_run_seconds` | `0` (без лимита) | Сколько максимум ждать на шаге 4 (отсчёт — **от конца сканирования**). По истечении цикл фиксирует `timeout` и переходит к паузе — защита от зависания (нет свободных валидаторов, недоступна площадка). Очередь при этом не трогается: следующий цикл её очистит, а до тех пор видно, что именно не дошло до конца. **При тысячах адресов оставьте `0`** (или задайте больше расчётного времени: 6440 адресов — часы). | Параметры хранятся в базе и меняются на лету, без перезапуска; в `control-api.yaml` ничего задавать не нужно. Новый `interval_seconds` применяется к паузе @@ -170,7 +193,8 @@ curl -s http://:8080/api/v1/admin/auto-cycle ### Фазы и результат последнего цикла `phase` показывает, что происходит сейчас: `idle` (автоцикл выключен или ещё не -стартовал), `running` (идут проверки — шаги 3–4) и `waiting` (пауза между циклами, +стартовал), `scanning` (шаги 1–2: очистка очереди и чтение/постановка Floating IP), +`running` (идут проверки — шаги 3–4) и `waiting` (пауза между циклами, шаг 5; время следующего запуска — `next_run_at`). Результат последнего цикла (`last_outcome`): @@ -185,7 +209,8 @@ curl -s http://:8080/api/v1/admin/auto-cycle ### Что важно знать - **Выключение не прерывает проверки**, которые уже идут: они закончатся и попадут - в реестр, остановится только повторение. + в реестр, остановится только повторение. Если выключить цикл во время фазы `scanning`, + сканирование отменяется (уже поставленные в очередь куски остаются). - Автоцикл **владеет очередью**: каждый цикл начинается с её полной очистки, поэтому адреса, добавленные вручную, будут удалены (их история в реестре остаётся). Ручные «Очистить всё» и «Сканировать Floating IP» во время цикла @@ -195,6 +220,12 @@ curl -s http://:8080/api/v1/admin/auto-cycle - События цикла (`auto_cycle_started`, `auto_cycle_completed`, `auto_cycle_timeout`, `auto_cycle_error`, `auto_cycle_stopped`) пишутся в журнал событий вместе с `queue_cleared` и `fip_scan`. +- Если в момент старта цикла уже идёт чужое сканирование (ручное, пробное или по расписанию), + цикл **дожидается** его окончания и запускает собственное (с очисткой очереди) — присоединяться + к чужому нельзя: оно могло ничего не поставить в очередь. +- Если сканирование завершилось ошибкой, очередь уже очищена (шаг 1) и остаётся пустой до следующего + цикла (`interval_seconds`); исход цикла — `error` с причиной в `last_error`. +- Цикл на тысячах адресов длится часы; пауза `interval_seconds` отсчитывается после его завершения. - В реальном OpenStack отвязка Floating IP после очистки очереди может отразиться с задержкой; если скан сразу после неё не увидел свободных адресов, цикл завершится с `no_free_ips` и повторится через `interval_seconds`. diff --git a/docs/changes/2026-10-01_18-19_fip-scan-at-scale-plan.md b/docs/changes/2026-10-01_18-19_fip-scan-at-scale-plan.md new file mode 100644 index 0000000..0c92223 --- /dev/null +++ b/docs/changes/2026-10-01_18-19_fip-scan-at-scale-plan.md @@ -0,0 +1,166 @@ +# План: сканирование Floating IP и автоцикл при тысячах адресов + +> Дата: 2026-10-01 18:19 MSK · Статус: **реализовано** — результаты ревью и тестов: [2026-10-01_18-59_fip-scan-at-scale-review.md](2026-10-01_18-59_fip-scan-at-scale-review.md) + +## Context + +Нажатие «Сканировать Floating IP» на живом стенде падает: `control-api недоступен: … context deadline exceeded`. Расследование +(2026-10-01) показало: в проекте OpenStack **6441 Floating IP, 6440 свободны** (раньше было 5). `ListFloatingIPs` запрашивает весь +список одним запросом без `limit` и без таймаута, Neutron отвечает >60 с (постранично: 200 адресов ≈ 2,4 с, весь список ≈ 70–80 с), +а дашборд ждёт 10 с. Запрос к тому же привязан к `r.Context()`: при обрыве соединения скан отменяется и не может завершиться. + +Цель пользователя прежняя: **одной кнопкой подключить к проверке все доступные в проекте адреса, даже если их тысячи**, после чего +проверка запускается автоматически; автоматический цикл (очистка → скан → проверка → пауза) должен работать в этих условиях. +Решение пользователя по объёму: **полная адаптация** — фоновый постраничный скан + автоцикл + постраничные страницы дашборда. + +Ожидание по времени (оценка по реальным данным стенда: слот на адрес ≈ 50 с, из них 30 с — `fip_settle_seconds`): +6440 адресов ≈ 18 ч на 5 валидаторах, ≈ 9 ч на 10, ≈ 4,5 ч на 20. Это ограничение пропускной способности, а не кода; рычаги — +число валидаторов и `fip_settle_seconds`. Для автоцикла это значит: `max_run_seconds` должен быть `0` (без лимита) или > 20 ч. + +## Карта кода (по графу и разведке) + +- `openstack.FloatingIPClient` (`internal/openstack/interface.go`): `GetFloatingIPByAddress`, `ListFloatingIPs`, `Associate…`, `Disassociate…`; + реализации — реальный `Client` (`client.go`, **нет таймаутов и ретраев**) и `MockClient` (`mock.go`, есть `ListFailure`, нет пагинации). +- `Orchestrator.ScanFloatingIPs` (`internal/orchestrator/orchestrator.go:412`) — синхронный: `OS.ListFloatingIPs` → фильтр `PortID==""` → + один `DB.SubmitIPs`. Вызывается из `handleAdminScanFloatingIPs` (`internal/httpapi/handlers_admin.go:107`, на `r.Context()`), + из `autoCycleStartRun` (`autocycle.go`, **в горутине цикла оркестратора под `autoCycleMu`** — минутный скан остановит `Tick` + и sweeps) и из периодического `scanTickerC` (`cmd/control-api/main.go`). +- БД: SQLite, `SetMaxOpenConns(1)`; `SubmitIPs` — одна транзакция, ~6 запросов на новый адрес (6440 ≈ 38 тыс. запросов); + `DeleteIPs`/`ClearQueue` — одна транзакция, ~6 запросов на адрес; `ListRegistry` — N+1 и **O(n²)**: нет индекса `ip_queue(registry_id)`. +- Full-table загрузчики: `GET /admin/ips`, `GET /admin/status` (грузит все строки ради счёта), `GET /admin/registry`, `autoCycleCheckRun` + (`ListIPs` каждый тик), дашборд `/overview` (опрос каждые 5 с: ~3 МБ JSON и таблица «текущая проверка» из ~6000 `queued`-строк), + `/ips` (~8 МБ HTML), `/registry`. Пагинации и фильтров на сервере нет. +- Дашборд: `client.do` с таймаутом 10 с; фрагменты и опрос htmx (`overview.html`, `overview_fragment.html`), GET-фильтр + `registry.html` (`hx-select` + `hx-replace-url`) — идиома для переиспользования. Per-row кнопки `hx-delete` при «выбрать все» + кладут все отмеченные адреса в URL (уже сейчас дефект, при тысячах — фатальный). +- Тик оркестратора не читает всю очередь (`ClaimNextQueued`, `ListChecking` — по индексу), пропускная способность не зависит от размера очереди. + +## Дизайн + +### 1. OpenStack: постраничное чтение, таймауты, ретраи (`internal/openstack`, `internal/config`) +- Новый метод интерфейса `ListFreeFloatingIPs(ctx, pageSize int, onPage func(page []FloatingIP) error) (pages int, err error)`: + цикл «страница → `onPage`»; для каждой страницы один запрос `floatingips.List(ListOpts{Limit, Marker=lastID})` с `EachPage` + (возврат `false` после первой страницы) — собственная пагинация по `marker`, а не `next`-ссылка (за прокси она может указывать на + внутренний хост). Свой `ListOptsBuilder`, добавляющий `fields=id&fields=floating_ip_address&fields=port_id&fields=project_id` + (в gophercloud `ListOpts.Fields` нет; на стенде проверено: `fields` + `marker` работают). Фильтр свободных — на клиенте + (`PortID==""`); серверный `status=DOWN` не используем (надмножество, возможны гонки статуса). +- **Ретраи страницы** с backoff (по умолчанию 5 попыток, 1→2→4→8→16 с) на сетевые ошибки, `EOF/RemoteDisconnected`, 5xx и 429 + (на стенде уже наблюдался `RemoteDisconnected` на второй странице); 4xx (кроме 429) — без ретрая. Контекст отменяет ретраи. +- **Таймаут на запрос**: `provider.HTTPClient.Timeout` (`openstack.request_timeout_seconds`, 60) — закрывает и вечные зависания в `Tick` + (`GetFloatingIPByAddress`/`Associate`/`Disassociate`), ключевой побочный эффект. +- Конфиг: `openstack.list_page_size` (200), `openstack.request_timeout_seconds` (60), `orchestrator.fip_scan_timeout_seconds` (1800), + `openstack.list_page_retries` (5); дефолты в `LoadControlAPI`, примеры в `configs/*.example.yaml` и `rxprod-compose/sources/`. +- `ListFloatingIPs` (полный список) остаётся для совместимости и тестов (реализован поверх нового метода). +- `MockClient`: пагинация (`PageSize`), счётчик вызовов/страниц, очередь ошибок `ListFailures []error` (по одной на запрос), + опциональная задержка страницы, `SeedMany(n)` для тестов на тысячи. + +### 2. Фоновое задание скана (`internal/orchestrator/scanjob.go`) +- `ScanJob` в `Orchestrator`, **нулевое значение пригодно** (тесты строят `&Orchestrator{…}` литералом): `sync.Mutex`, текущий прогресс, + `cancel`. Метод `StartScan(opts) (ScanStatus, started bool)` — single-flight: если скан уже идёт, возвращает его статус + (`started=false`). Горутина работает на контексте жизни процесса (хранится в `Orchestrator`, задаётся из `main`, по умолчанию + `context.Background()`), **не** на `r.Context()`; общий дедлайн `fip_scan_timeout_seconds`. +- Опции: `ClearFirst bool` (для автоцикла), `DryRun bool` (только обнаружить и посчитать, очередь не трогать — безопасная проверка + на живом стенде и полезная функция для оператора). +- Фазы и прогресс: `idle → clearing → listing → enqueuing → done|error|cancelled`; поля `pages`, `discovered`, `free`, `added`, + `requeued`, `reordered`, `skipped_in_progress`, `started_at`, `finished_at`, `error`. +- **Алгоритм:** (1) `ClearFirst` → `ClearQueue`; (2) чтение всех страниц в память (6440 строк — килобайты), `free = PortID==""`; + (3) сортировка по IPv4 по возрастанию (детерминированный порядок очереди); (4) **только после полного обнаружения** — `SubmitIPs` + кусками по 500 в этом порядке (`base=MAX+1` пересчитывается на вызов ⇒ порядок сохраняется; транзакции короткие, единственное + соединение освобождается между кусками); проверки стартуют, как только появляются первые `queued`; (5) одно событие `fip_scan` + с итоговыми счётчиками. Ошибка чтения после ретраев ⇒ **ничего не ставится в очередь** (для ручного скана очередь не меняется), + статус `error` с причиной; повтор — кнопкой или следующим циклом. Ошибка БД посередине ⇒ уже поставленные куски остаются + (повтор идемпотентен: `SubmitIPs` переупорядочивает/пропускает). +- `ScanFloatingIPs(ctx)` остаётся тонкой синхронной обёрткой («запустить и дождаться») для существующих тестов/скриптов. +- Периодический `scanTickerC` вызывает неблокирующий `StartScan`. + +### 3. Масштабирование БД и запросов (`internal/db`, миграция `0009`) +- Миграция `0009_scale_indexes.sql`: `idx_ip_queue_registry ON ip_queue(registry_id)` (убирает O(n²) в реестре), + `idx_ip_queue_state_aggregated ON ip_queue(state, aggregated_at)` (список «последние завершённые»). +- Новые запросы: `CountIPsByState`, `CountIPsByResult` (GROUP BY — вместо загрузки всех строк в `/admin/status`), + `AnyNonTerminalIP` (`SELECT EXISTS … state NOT IN (done,failed,occupied)`), `ListIPsPage(filter{states[], q, result, order}, + limit, offset) → (items, total)`, `ListRegistryPage(filter{q, lastResult}, limit, offset) → (items, total)` — **LIMIT/OFFSET до** + `fillRegistrySummary`, поэтому 3–4 запроса на строку платят только строки страницы. Фильтр `lastResult` реализуется одним SQL: + `ip_registry r LEFT JOIN ip_queue q ON q.registry_id=r.id`, условие `(q.id IS NOT NULL AND q.overall_result=?) OR (q.id IS NULL AND + <подзапрос по checks последнего цикла: pass/fail/partial>=?)` — та же семантика, что `fillRegistrySummary`/`lastCycleResultFromChecks`, + без денормализации и миграции данных. +- `ClearQueue`: set-based очистка без цикла по адресам — `UPDATE validators SET current_ip_id=NULL…`, `UPDATE checks SET ip_id=NULL`, + `UPDATE events SET ip_id=NULL`, `DELETE ip_site_checks`, `DELETE ip_queue` (5 запросов, O(n)); disassociate FIP только для строк с `FIPID`; + событие `queue_cleared` — счётчик и усечённый список (не 6440 адресов). `Orchestrator.DeleteIPs` — выбор строк без N `GetIPByAddress`. + +### 4. HTTP API (`internal/httpapi`) — обратная совместимость сохраняется +| Метод | Путь | Изменение | +|---|---|---| +| POST | `/admin/ips/scan` | `202 {state, started_at, …}` (запуск или уже идущий скан — `202` с текущим статусом); `?dry_run=true`; `?wait=true` — старая синхронная семантика (`200` + счётчики) для curl/скриптов | +| GET | `/admin/ips/scan` | **новый**: статус и прогресс скана (admin-токен) | +| GET | `/admin/ips` | без параметров — как раньше (массив); с `limit` — конверт `{items,total,limit,offset}`; фильтры `state` (csv), `q`, `result`, `order` | +| GET | `/admin/registry` | то же: `limit/offset/q/last_result` → конверт с `total` | +| GET | `/admin/status` | + `results_by_overall`, счёт через `GROUP BY` | +| GET | `/admin/overview` | **опционально** одним запросом: счётчики, активные (≤100), последние завершённые (N), ближайшие в очереди (≤10), статус скана и автоцикла | + +Новые admin-маршруты попадают в таблицу `routes.go` с `accessAdmin`; `TestRouteTableClassification` (`auth_test.go`) обновить (+1–2 admin). + +### 5. Автоцикл (`internal/orchestrator/autocycle.go`, `queries_autocycle.go`, `dashboard/dto.go`) +- Новая фаза **`scanning`** (миграция не нужна — валидатор фаз в `UpdateAutoCycleState` расширить; подписи `PhaseLabel` — «сканирование Floating IP»). +- `autoCycleStartRun` перестаёт блокировать цикл: запускает `StartScan{ClearFirst:true}` (очистка + скан целиком в фоне, **литерал + сценария пользователя сохранён: очистка → скан**) и сразу переводит фазу в `scanning`; `autoCycleMu` держится только на время + чтения/записи состояния, а не на всё время скана ⇒ `Tick` и `Start/Stop` не блокируются. +- Шаг `scanning`: опрос статуса задания. `running` → выход; `error` → существующий путь `fail()` (исход `error`, повтор через + `interval_seconds`); `done` и `free==0` → `no_free_ips`; `done` → фаза `running`, `last_scanned_free`, **`run_started_at` = конец скана** + (лимит `max_run_seconds` считается от конца скана). Таймаут самого скана — `fip_scan_timeout_seconds`. +- **Восстановление после рестарта:** фаза `scanning` без живого задания ⇒ заново `StartScan{ClearFirst:true}` (идемпотентно). +- `Stop` отменяет задание скана (исход `stopped`, если шёл скан или проверка). +- `autoCycleCheckRun`: проверка завершения — `AnyNonTerminalIP` вместо `ListIPs` каждый тик; `COUNT` только при завершении. +- Документировать: при тысячах адресов `max_run_seconds=0`; цикл длится часы; интервал отсчитывается от завершения. + +### 6. Дашборд (`internal/dashboard`, шаблоны, CSS) +- **Скан-кнопка и прогресс:** `hx-post="/ips/scan"` возвращает панель `scan_progress` (вне `#ips-form`): стадия, ``, + прочитано/свободных/добавлено, время, ошибка; пока `running` панель сама опрашивает `GET /ips/scan/status` (`hx-trigger="every 2s"`), + по завершении — без триггера и с `HX-Trigger: scan-finished`, по которому таблица перезагружается (`hx-get` + `hx-select`); + кнопка блокируется на время скана; `409`/ошибки — штатным баннером (`bannerFor`). Скан-старт возвращается мгновенно, таймаут + 10 с больше не проблема. +- **Пагинация `/ips` и `/registry`:** `page`, `per_page` (50 по умолчанию; 25/50/100/200), partial `pager` («Показано a–b из N», ‹ ›, + `url.Values` для экранирования); фильтры на сервере: `/registry` — `q`, `status`; `/ips` — новая форма `q` + состояние + (все / в очереди / в работе / done / failed / occupied / результат), идиома GET + `hx-select` + `hx-replace-url`. + Скрытые `page/q/state` внутри `#ips-form`, чтобы мутации возвращали ту же страницу. +- **Массовые операции:** чекбоксы — только строки страницы + счётчик «Выбрано на странице k из 50»; при отмеченном заголовке и + `total > per_page` — ссылка «Выбрать все N по фильтру» (`scope=all`): адреса разрешаются на сервере постранично и уходят в + `DeleteIPs`/`SubmitIPs` кусками по ~500. Per-row и scan/clear-кнопки получают `hx-params="page,q,state"` (чинит URL из тысяч адресов). + `hx-confirm` с реальным числом (`Удалить ВСЕ {{.Total}} адресов…`). +- **«Обзор»:** без `ListIPs`; `Status` (+`results_by_overall`) и ограниченные списки: «В работе» (активные состояния), «В очереди: Q» + (счётчик + ссылка на `/ips?state=queued`, ≤10 ближайших), «Последние N завершённых»; `occupied` — терминальное состояние. + Индикатор в блоке статистики (OOB-обновление, как сейчас): «Готово D из T (P%) · в работе A · в очереди Q» + ``, + оценка времени по скорости последних завершённых; статус скана («Сканирование: прочитано X») и автоцикла. +- `client.go`: `ListIPsPage`, `ListRegistryPage`, `ScanStatus`; длинный таймаут/отдельный клиент для clear и массовых операций. + +### 7. Прочее +- `routes`, `dto_admin.go`, `docs/API.md` (202/статус/пагинация/`dry_run`/`wait`), `docs/USAGE.md` (скан, автоцикл: фаза `scanning`, + ожидание ≈ N×50 с/валидаторов, рычаги), `docs/DASHBOARD.md`, `docs/LOCAL_E2E.md`, `README.md`, `configs/*.example.yaml` + (+ копии в `rxprod-compose/sources/` и docker-примеры). +- `rxprod-compose/control-api.yaml` (живой конфиг) — при необходимости задать `max_run_seconds` автоцикла = 0 (уже 0) и + `openstack.list_page_size`; по умолчанию достаточно дефолтов. +- Опционально (не входит): переупорядочить `Tick` (сначала sweeps, потом `assignIdleValidators`) — экономит до 5 с на адрес (~10 %). + +## Тесты (минимальные, в стиле существующих) +- `internal/openstack`: пагинация по marker на моке (N=2500, `PageSize=200`), ретрай страницы при `ListFailures`, отмена контекста; + классификация ретраемых ошибок. +- `internal/orchestrator`: задание скана — single-flight, прогресс, `DryRun`, ошибка чтения ⇒ очередь не изменена, куски по 500, + порядок по возрастанию IP, 6440 адресов за разумное время; автоцикл: фаза `scanning`, шаги с явным `now`, ошибка скана, + `no_free_ips`, рестарт в фазе `scanning`, `Stop` отменяет скан, `Tick` не блокируется во время долгого скана (мок с задержкой). +- `internal/db`: `ListIPsPage`/`ListRegistryPage` (фильтры, total, LIMIT до summary), `CountIPsBy*`, `AnyNonTerminalIP`, + set-based `ClearQueue`, тест на 6440 строк (время и отсутствие N+1). +- `internal/httpapi`: 202/409-семантика скана, `?wait=true`, `?dry_run=true`, `GET /ips/scan`, конверт пагинации, обратная + совместимость без `limit`; обновить `TestScanFloatingIPsEndpoint` и счётчики в `auth_test.go`. +- `internal/dashboard`: панель прогресса и остановка опроса, пагинация/pager, фильтры `/ips`, `scope=all`, Overview на ограниченных + списках при тысячах `queued`, `hx-params`; обновить `fakeControlAPI` и затронутые тесты (`TestIPsScan` и др.). +- `scripts/run-local-e2e.sh`: автоцикл проходит через фазу `scanning` (мок-пагинация), полный сценарий остаётся зелёным за 90 с. + +## Верификация +1. `go build ./... && go vet ./... && go test ./...` и `go test -race` для `orchestrator`, `openstack`, `httpapi`, `dashboard`, `db`. +2. `scripts/run-local-e2e.sh` (с токенами) — зелёный, автоцикл проходит `scanning → running → waiting`. +3. **Живой стенд, безопасно (чтение):** `POST /admin/ips/scan?dry_run=true` — скан проходит все страницы реального Neutron, прогресс + идёт, итог ≈ 6440 свободных, очередь не меняется, дашборд не получает таймаута. Замер времени скана и нагрузки. +4. **Живой стенд, по согласованию:** кнопка «Сканировать Floating IP» — адреса поставлены в очередь, проверки идут, `/ips`, + `/registry`, «Обзор» остаются быстрыми (проверка размера ответов и времени), «Очистить всё» отрабатывает за секунды. + Решение о реальной постановке 6440 адресов (≈18 ч проверок на 5 валидаторах) принимает пользователь; перед этим — бэкап БД. +5. Автоцикл на живом стенде — только после п. 4 и с согласия пользователя; наблюдать фазу `scanning`, отсутствие блокировки `Tick`. +6. Реальный браузер (Playwright из venv): прогресс скана, пагинация, фильтры, выбор «все N по фильтру», отсутствие JS-ошибок. diff --git a/docs/changes/2026-10-01_18-59_fip-scan-at-scale-review.md b/docs/changes/2026-10-01_18-59_fip-scan-at-scale-review.md new file mode 100644 index 0000000..df5fe1f --- /dev/null +++ b/docs/changes/2026-10-01_18-59_fip-scan-at-scale-review.md @@ -0,0 +1,81 @@ +# Ревью и тестирование: скан Floating IP и автоцикл при тысячах адресов + +> Дата: 2026-10-01 18:59 MSK · План: [2026-10-01_18-19_fip-scan-at-scale-plan.md](2026-10-01_18-19_fip-scan-at-scale-plan.md) +> Статус: **реализовано и проверено на живом окружении**; реальная постановка 6440 адресов в очередь и включение автоцикла на живом стенде **не выполнялись** — ждут решения пользователя. + +## Итог + +Исходная ошибка («control-api недоступен: context deadline exceeded» при нажатии «Сканировать Floating IP») устранена: на живом стенде кнопка +возвращает панель прогресса за 0,7 с, а сканирование реального Neutron (6441 Floating IP, 6440 свободных) проходит за ≈ 1,5–2 минуты в фоне +с видимым прогрессом. Код написан двумя агентами параллельно (Sonnet 5.5): серверная часть и дашборд, ревью и все проверки — независимо (Sonnet 5.5, high). +Найден и исправлен один дефект автоцикла и одна мелочь в клиенте OpenStack; остальные замечания — ограничения дизайна, перечислены ниже. + +## Причина исходной ошибки + +`ListFloatingIPs` запрашивал весь список одним запросом без `limit` и без таймаута: при ≈ 6,4 тыс. адресов Neutron отвечал дольше минуты, дашборд ждал 10 с, +а запрос был привязан к `r.Context()` — при обрыве соединения скан отменялся и не мог завершиться. Ошибка нигде не логировалась. +Побочные открытия: у клиента OpenStack вообще не было таймаутов и ретраев; `ListRegistry` был O(n²) (нет индекса `ip_queue(registry_id)`); +`/ips` отдавал ≈ 8 МБ HTML, а «Обзор» каждые 5 с тянул ≈ 3 МБ JSON. + +## Что реализовано + +| Область | Изменения | +|---|---| +| OpenStack | `ListFreeFloatingIPs` — постраничное чтение по `marker` (200 на страницу, `fields=` сокращает ответ), повтор страницы с backoff на обрывы/`RemoteDisconnected`/5xx/429, таймаут запроса (`openstack.request_timeout_seconds`, 60 с — закрывает и зависания в тике оркестратора); `MockClient` с пагинацией, `ListFailures`, `PageDelay`, `SeedMany` | +| Скан-задание | `orchestrator/scanjob.go`: single-flight фоновое задание на контексте процесса (не `r.Context()`), фазы `clearing → listing → enqueuing → done/error/cancelled`, прогресс; сначала полное обнаружение, затем `SubmitIPs` кусками по 500 по возрастанию IP; сбой чтения ⇒ очередь не меняется; `dry_run`; синхронная обёртка `ScanFloatingIPs` сохранена | +| Автоцикл | фаза `scanning`: очистка и скан — одно фоновое задание, цикл оркестратора и `autoCycleMu` не блокируются; восстановление после рестарта; `Stop` отменяет скан; завершение определяется `EXISTS`, а не чтением всей очереди каждый тик | +| БД | миграция `0009` (индексы `ip_queue(registry_id)`, `ip_queue(state, aggregated_at)`); `ListIPsPage`, `ListRegistryPage` (LIMIT/OFFSET до расчёта сводки, фильтр итога одним SQL), `CountIPsByState/Result`, `AnyNonTerminalIP`; `ClearAllIPs` — 5 запросов вместо цикла по адресам | +| API | `POST /admin/ips/scan` → `202` (`dry_run`, `wait`), новый `GET /admin/ips/scan`, пагинация и фильтры у `GET /admin/ips` и `/admin/registry` (без `limit` — прежний массив), `results_by_overall` в `/admin/status`, `count` у `clear` | +| Дашборд | панель прогресса скана и «Пробное сканирование»; постраничные `/ips` и `/registry` с серверными фильтрами; «Обзор» на счётчиках и ограниченных списках (прогресс «Готово D из T», оценка времени); «Выбрать все N по фильтру»; `hx-params` на кнопках (исправлен дефект — отмеченные адреса попадали в URL `hx-delete`); подтверждения с реальным числом; длинный таймаут для массовых операций | +| Конфиг и документы | `openstack.list_page_size/request_timeout_seconds/list_page_retries`, `orchestrator.fip_scan_timeout_seconds`; `API.md`, `USAGE.md`, `DASHBOARD.md`, `README.md`, примеры конфигов | + +## Результаты проверок + +| Проверка | Результат | +|---|---| +| `gofmt`, `go build ./...`, `go vet ./...` | чисто | +| `go test ./...` (включая тесты на 6440 адресов) | все пакеты зелёные | +| `go test -race -short` (openstack, orchestrator, db, httpapi, dashboard, config) | зелёные (тесты на 6440 адресов под `-race` слишком долгие, пропускаются по `-short`; агент прогонял их полностью) | +| `scripts/run-local-e2e.sh` (с токенами) | exit 0: автоцикл прошёл через `scanning`, проверки реальными агентом и пробером — `pass` | +| **Живой стенд — пробный скан реального Neutron** (`POST …/scan?dry_run=true`) | 33 страницы, **6441 найдено / 6440 свободных за 93 с**, `202` за 2 мс, повторный `POST` присоединился к идущему заданию, **API отвечал за 2–4 мс всё время скана**, очередь осталась пустой | +| **Живой дашборд (настоящий Chromium)**, кнопка «Пробное сканирование» | панель за 0,7 с без баннера ошибки, прогресс по страницам, итог «готово: 33 страницы, 6441, 6440, время 1 мин 48 с» | +| Изолированный mock-стенд на **6440 адресах**, настоящий Chromium (17 из 19 автопроверок, 2 — ложные, см. ниже) | `/ips` **68 КБ за 0,18 с** (было ≈ 8 МБ), фрагмент «Обзора» **3 КБ за 0,05 с**, `/registry` 34 КБ за 0,12 с; пагинация и серверный поиск; прогресс скана (читаются страницы → ставятся в очередь → готово); «Очистить всё» с реальным числом в подтверждении — **0,5 с**; «Выбрать все 6440 по фильтру» + массовое удаление — **8 с**; JS-ошибок нет; ошибок в логе control-api нет | +| Миграция `0009` на живой БД | `user_version = 9`, оба индекса созданы, данные не тронуты | + +Примечание: два «FAIL» в браузерном скрипте mock-стенда — ошибка самого скрипта: панель показывает состояние заглавными («ГОТОВО», CSS), а скрипт искал строчные. +Выведенный текст панели подтверждает успех (`добавлено6440`, `найдено адресов6440`). Реальных провалов нет. + +## Замечания ревью + +| № | Серьёзность | Замечание | Статус | +|---|---|---|---| +| 1 | средняя | **Автоцикл «усыновлял» чужое сканирование.** Если в момент старта цикла уже шло ручное/периодическое/**пробное** сканирование, `StartScan` возвращал `started=false`, а цикл переходил в `scanning` и ждал чужое задание. Пробное ничего не ставит в очередь, ручное не очищает очередь ⇒ цикл переходил в `running` над пустой/нетронутой очередью и сразу отчитывался `completed` (`runs_total+1`) без единой проверки | **исправлено**: если собственный скан не стартовал, цикл ничего не меняет и пробует снова на следующем такте (после окончания чужого); добавлен тест `TestAutoCycleWaitsForForeignScanInsteadOfFollowingIt` (падал до правки) | +| 2 | низкая | Клиент OpenStack заполнял `ProjectID` только из `tenant_id`; при `fields=` Neutron может вернуть лишь `project_id` | **исправлено** (запасной вариант `project_id`); поле нигде не влияет на логику | +| 3 | низкая | `aggregated_at_desc` сортирует по `strftime(...)` — временные метки хранятся как RFC3339Nano, и сырая сортировка текстом неверна (поймал тест агента); индекс `(state, aggregated_at)` помогает фильтру по состоянию, но не сортировке | принято; на 6440 строк незаметно | +| 4 | низкая | `StopAutoCycle` в фазе `scanning` вызывает `CancelScan`, который ждёт до 5 с под `autoCycleMu`: «Выключить» может занять до 5 с, а следующий шаг цикла — подождать | принято | +| 5 | низкая | Если скан завершился ошибкой после «Очистить» (шаг 1 цикла), очередь остаётся пустой до следующего цикла (`interval_seconds`); исход — `error` с причиной | принято, описано в `USAGE.md`; при желании — отдельная доработка (повтор скана сразу) | +| 6 | низкая | Статус скана хранится в памяти: после рестарта control-api он `idle`; автоцикл в фазе `scanning` при этом корректно перезапускает скан | принято | +| 7 | инфо | Отмена сканирования (`CancelScan`) вызывается только из `Stop` автоцикла; ручной кнопки/эндпоинта отмены нет | не входило в план | +| 8 | инфо | Дашборд: мутации заменяют `#ips-table-wrap` целиком (`outerHTML`), чтобы `hx-get` обёртки всегда указывал на текущую страницу/фильтр; убраны функции `filterQueueItems/filterRegistryItems/currentlyChecking/lastCompleted` вместе с тестами (фильтрация перенесена на сервер); сводка «последние N» считается по показанному (возможно, отфильтрованному) окну, а общие итоги — отдельной строкой | принято | +| 9 | инфо | Тесты на 6440 адресов под `-race` занимают 40–90 с на пакет — пропускаются по `-short` | принято | + +## Пропускная способность (важно для автоцикла) + +Проверка не стала быстрее — стало возможным её запустить. По фактическим данным стенда слот на адрес ≈ 50 с на валидатор (из них 30 с — `fip_settle_seconds`): +**6440 адресов ≈ 18 ч на 5 валидаторах, ≈ 9 ч на 10, ≈ 4,5 ч на 20.** Для автоцикла `max_run_seconds` должен оставаться `0`. Рычаги — число валидаторов и (осторожно) +`fip_settle_seconds`. На странице «Обзор» виден прогресс и оценка времени. + +## Состояние живого стенда (`rxprod-compose`) + +- Развёрнуты новые образы `civ-capi`, `civ-adash`, `civ-prober` (и пересобран `civ-agent`); миграция `0009` применена. Предыдущие образы сохранены под тегом `:pre-scale` + (откат: `docker tag civ-capi:pre-scale civ-capi:latest` и `docker compose up -d`; на `0009` откат БД не нужен — это только индексы). Бэкап БД перед обновлением — в каталоге scratchpad сессии (`/tmp`, временный). +- Очередь пуста (5 адресов прежней работы остались в реестре). **Реальная постановка 6440 адресов и автоцикл на живом стенде не запускались**: это ≈ 18 ч реальных проверок, + решение за пользователем. Перед запуском рекомендую «Пробное сканирование» (уже отработало штатно) и бэкап БД. +- Токен агентов по-прежнему не включён (внешние валидаторы и пробер `rxyc` со старыми бинарниками) — см. [ревью аутентификации](2026-10-01_11-31_authentication-review.md). + Новые бинарники для внешних валидаторов — в `bin/` (после раскатки токена агентов их можно обновить одновременно). +- `bin/` пересобран (`CGO_ENABLED=0`, `-trimpath -ldflags="-s -w"`), `SHA256SUMS` обновлён. Временные контейнеры `civ-scale-*` удалены. + +## Что осталось + +- Решение пользователя: поставить 6440 адресов в очередь («Сканировать Floating IP») и/или включить автоцикл на живом стенде. +- По желанию: кнопка/эндпоинт отмены скана (замечание 7), повтор скана внутри цикла при ошибке (замечание 5), `Cache-Control: no-store` (из ревью аутентификации). diff --git a/internal/config/config.go b/internal/config/config.go index 30010cc..99b1e7c 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -69,6 +69,17 @@ type OpenStackConfig struct { UsernameEnv string `yaml:"username_env"` // default OS_USERNAME — used when auth_method: password UserDomainNameEnv string `yaml:"user_domain_name_env"` // default OS_USER_DOMAIN_NAME PasswordEnv string `yaml:"password_env"` // default OS_PASSWORD + + // ListPageSize is how many floating IPs one Neutron list request asks + // for (the scan reads the project page by page). Default 200. + ListPageSize int `yaml:"list_page_size"` + // RequestTimeoutSeconds bounds every single HTTP request to Keystone and + // Neutron. Default 60. + RequestTimeoutSeconds int `yaml:"request_timeout_seconds"` + // ListPageRetries is how many times one failed page of the listing is + // retried (backoff 1s,2s,4s,...) on network errors, 5xx and 429. Default + // 5; a negative value disables retries. + ListPageRetries int `yaml:"list_page_retries"` } type OrchestratorConfig struct { @@ -93,6 +104,9 @@ type OrchestratorConfig struct { // POST /api/v1/admin/ips/scan or the dashboard's "Scan Floating IPs" // button either way. FIPScanIntervalSeconds int `yaml:"fip_scan_interval_seconds"` + // FIPScanTimeoutSeconds is the overall deadline of one background + // floating-IP scan (clear + paged read + enqueue). Default 1800. + FIPScanTimeoutSeconds int `yaml:"fip_scan_timeout_seconds"` } type AggregationConfig struct { @@ -164,6 +178,15 @@ func LoadControlAPI(path string) (*ControlAPI, error) { if c.OpenStack.PasswordEnv == "" { c.OpenStack.PasswordEnv = "OS_PASSWORD" } + if c.OpenStack.ListPageSize == 0 { + c.OpenStack.ListPageSize = 200 + } + if c.OpenStack.RequestTimeoutSeconds == 0 { + c.OpenStack.RequestTimeoutSeconds = 60 + } + if c.OpenStack.ListPageRetries == 0 { + c.OpenStack.ListPageRetries = 5 + } if c.Auth.AdminTokenEnv == "" { c.Auth.AdminTokenEnv = "CONTROL_API_ADMIN_TOKEN" } @@ -191,6 +214,9 @@ func LoadControlAPI(path string) (*ControlAPI, error) { if c.Orchestrator.HeartbeatTimeoutSeconds == 0 { c.Orchestrator.HeartbeatTimeoutSeconds = 30 } + if c.Orchestrator.FIPScanTimeoutSeconds == 0 { + c.Orchestrator.FIPScanTimeoutSeconds = 1800 + } return &c, nil } diff --git a/internal/config/config_test.go b/internal/config/config_test.go new file mode 100644 index 0000000..6f033e1 --- /dev/null +++ b/internal/config/config_test.go @@ -0,0 +1,66 @@ +package config + +import ( + "bytes" + "os" + "path/filepath" + "testing" +) + +func TestLoadControlAPIScanDefaults(t *testing.T) { + path := filepath.Join(t.TempDir(), "c.yaml") + if err := os.WriteFile(path, []byte("server:\n listen_addr: \":8080\"\n"), 0o600); err != nil { + t.Fatal(err) + } + c, err := LoadControlAPI(path) + if err != nil { + t.Fatalf("load: %v", err) + } + if c.OpenStack.ListPageSize != 200 || c.OpenStack.RequestTimeoutSeconds != 60 || + c.OpenStack.ListPageRetries != 5 || c.Orchestrator.FIPScanTimeoutSeconds != 1800 { + t.Fatalf("unexpected defaults: openstack=%+v orchestrator=%+v", c.OpenStack, c.Orchestrator) + } +} + +func TestLoadControlAPIScanOverrides(t *testing.T) { + path := filepath.Join(t.TempDir(), "c.yaml") + yaml := "openstack:\n list_page_size: 50\n request_timeout_seconds: 10\n list_page_retries: -1\n" + + "orchestrator:\n fip_scan_timeout_seconds: 99\n" + if err := os.WriteFile(path, []byte(yaml), 0o600); err != nil { + t.Fatal(err) + } + c, err := LoadControlAPI(path) + if err != nil { + t.Fatalf("load: %v", err) + } + if c.OpenStack.ListPageSize != 50 || c.OpenStack.RequestTimeoutSeconds != 10 || + c.OpenStack.ListPageRetries != -1 || c.Orchestrator.FIPScanTimeoutSeconds != 99 { + t.Fatalf("overrides lost: openstack=%+v orchestrator=%+v", c.OpenStack, c.Orchestrator) + } +} + +// The shipped example must load and carry the scan settings, and the rxprod +// copy must stay byte-identical to it. +func TestControlAPIExampleConfigs(t *testing.T) { + c, err := LoadControlAPI("../../configs/control-api.example.yaml") + if err != nil { + t.Fatalf("load example: %v", err) + } + if c.OpenStack.ListPageSize != 200 || c.Orchestrator.FIPScanTimeoutSeconds != 1800 { + t.Fatalf("example scan settings: %+v %+v", c.OpenStack, c.Orchestrator) + } + a, err := os.ReadFile("../../configs/control-api.example.yaml") + if err != nil { + t.Fatal(err) + } + b, err := os.ReadFile("../../rxprod-compose/sources/control-api.example.yaml") + if err != nil { + t.Fatal(err) + } + if !bytes.Equal(a, b) { + t.Fatalf("rxprod-compose/sources/control-api.example.yaml differs from configs/control-api.example.yaml") + } + if _, err := LoadControlAPI("../../deploy/docker/control-api/control-api.docker.example.yaml"); err != nil { + t.Fatalf("load docker example: %v", err) + } +} diff --git a/internal/dashboard/client.go b/internal/dashboard/client.go index 45684e6..ebe152d 100644 --- a/internal/dashboard/client.go +++ b/internal/dashboard/client.go @@ -8,6 +8,8 @@ import ( "io" "net/http" "net/url" + "strconv" + "strings" "time" ) @@ -39,15 +41,44 @@ func (e *apiErr) Error() string { type client struct { baseURL string http *http.Client + // long is used for clear / bulk operations, which legitimately take far + // longer than a plain read (thousands of rows in one transaction): same + // transport, but a longer whole-call timeout. + long *http.Client // token, when non-empty, is sent to control-api as a Bearer credential. token string } +// longCallTimeout is the minimum whole-call timeout for clear and bulk +// operations (ClearQueue, DeleteIPs, SubmitIPs). +const longCallTimeout = 120 * time.Second + func newClient(baseURL string, timeout time.Duration) *client { - return &client{baseURL: baseURL, http: &http.Client{Timeout: timeout}} + longT := longCallTimeout + if timeout > longT { + longT = timeout + } + return &client{ + baseURL: baseURL, + http: &http.Client{Timeout: timeout}, + long: &http.Client{Timeout: longT}, + } } func (c *client) do(ctx context.Context, method, path string, body, out interface{}) error { + return c.doWith(ctx, c.http, method, path, body, out) +} + +// doLong is do with the long (clear/bulk) timeout. +func (c *client) doLong(ctx context.Context, method, path string, body, out interface{}) error { + hc := c.long + if hc == nil { + hc = c.http + } + return c.doWith(ctx, hc, method, path, body, out) +} + +func (c *client) doWith(ctx context.Context, hc *http.Client, method, path string, body, out interface{}) error { var reader io.Reader if body != nil { b, err := json.Marshal(body) @@ -67,7 +98,7 @@ func (c *client) do(ctx context.Context, method, path string, body, out interfac req.Header.Set("Authorization", "Bearer "+c.token) } - resp, err := c.http.Do(req) + resp, err := hc.Do(req) if err != nil { return &apiErr{Status: 0, Message: err.Error()} } @@ -96,9 +127,58 @@ func (c *client) Status(ctx context.Context) (statusResponse, error) { return out, err } -func (c *client) ListIPs(ctx context.Context) ([]ipQueueItem, error) { - var out []ipQueueItem - err := c.do(ctx, http.MethodGet, "/api/v1/admin/ips", nil, &out) +// maxPageLimit is control-api's cap on `limit`. +const maxPageLimit = 1000 + +// clampLimit keeps limit within 1..maxPageLimit: a request without `limit` +// would make control-api answer with the legacy unbounded bare array. +func clampLimit(limit int) int { + if limit < 1 { + return 1 + } + if limit > maxPageLimit { + return maxPageLimit + } + return limit +} + +// ipsQuery selects one page of GET /admin/ips: server-side filters plus +// limit/offset. Order is "sequence" (default) or "aggregated_at_desc". +type ipsQuery struct { + States []string + Q string + Result string + Order string + Limit int + Offset int +} + +func (q ipsQuery) values() url.Values { + v := url.Values{} + v.Set("limit", strconv.Itoa(clampLimit(q.Limit))) + if q.Offset > 0 { + v.Set("offset", strconv.Itoa(q.Offset)) + } + if len(q.States) > 0 { + v.Set("state", strings.Join(q.States, ",")) + } + if q.Q != "" { + v.Set("q", q.Q) + } + if q.Result != "" { + v.Set("result", q.Result) + } + if q.Order != "" { + v.Set("order", q.Order) + } + return v +} + +// ListIPsPage returns one page of the check queue plus the total number of +// rows matching the filter. Never loads the whole queue. +func (c *client) ListIPsPage(ctx context.Context, q ipsQuery) (ipsPage, error) { + var out ipsPage + err := c.do(ctx, http.MethodGet, "/api/v1/admin/ips?"+q.values().Encode(), nil, &out) return out, err } @@ -112,7 +192,7 @@ func (c *client) GetIP(ctx context.Context, ip string) (ipDetailResponse, error) // forcing a recheck of already-finished ones — see docs/API.md. func (c *client) SubmitIPs(ctx context.Context, addresses []string) (submitIPsResponse, error) { var out submitIPsResponse - err := c.do(ctx, http.MethodPost, "/api/v1/admin/ips", map[string][]string{"addresses": addresses}, &out) + err := c.doLong(ctx, http.MethodPost, "/api/v1/admin/ips", map[string][]string{"addresses": addresses}, &out) return out, err } @@ -129,7 +209,7 @@ func (c *client) DeleteIP(ctx context.Context, ip string) error { // DeleteIPs permanently removes a specific list of addresses in one call. func (c *client) DeleteIPs(ctx context.Context, addresses []string) (deleteIPsResponse, error) { var out deleteIPsResponse - err := c.do(ctx, http.MethodPost, "/api/v1/admin/ips/delete", map[string][]string{"addresses": addresses}, &out) + err := c.doLong(ctx, http.MethodPost, "/api/v1/admin/ips/delete", map[string][]string{"addresses": addresses}, &out) return out, err } @@ -137,25 +217,55 @@ func (c *client) DeleteIPs(ctx context.Context, addresses []string) (deleteIPsRe // including those actively being checked. func (c *client) ClearQueue(ctx context.Context) (clearQueueResponse, error) { var out clearQueueResponse - err := c.do(ctx, http.MethodPost, "/api/v1/admin/ips/clear", nil, &out) + err := c.doLong(ctx, http.MethodPost, "/api/v1/admin/ips/clear", nil, &out) 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) +// StartScan starts control-api's background floating-IP scan (or joins the +// one already running) and returns immediately with its current status. +func (c *client) StartScan(ctx context.Context, dryRun bool) (scanStatusDTO, error) { + var out scanStatusDTO + path := "/api/v1/admin/ips/scan" + if dryRun { + path += "?dry_run=true" + } + err := c.do(ctx, http.MethodPost, path, 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) +// ScanStatus returns the progress of the background scan job. +func (c *client) ScanStatus(ctx context.Context) (scanStatusDTO, error) { + var out scanStatusDTO + err := c.do(ctx, http.MethodGet, "/api/v1/admin/ips/scan", nil, &out) + return out, err +} + +// registryQuery selects one page of GET /admin/registry. +type registryQuery struct { + Q string + LastResult string + Limit int + Offset int +} + +// ListRegistryPage returns one page of the registry (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) plus the total number of rows matching the filter. +func (c *client) ListRegistryPage(ctx context.Context, q registryQuery) (registryPage, error) { + v := url.Values{} + v.Set("limit", strconv.Itoa(clampLimit(q.Limit))) + if q.Offset > 0 { + v.Set("offset", strconv.Itoa(q.Offset)) + } + if q.Q != "" { + v.Set("q", q.Q) + } + if q.LastResult != "" { + v.Set("last_result", q.LastResult) + } + var out registryPage + err := c.do(ctx, http.MethodGet, "/api/v1/admin/registry?"+v.Encode(), nil, &out) return out, err } diff --git a/internal/dashboard/dashboard_test.go b/internal/dashboard/dashboard_test.go index ed6c5fb..3488ffc 100644 --- a/internal/dashboard/dashboard_test.go +++ b/internal/dashboard/dashboard_test.go @@ -9,6 +9,8 @@ import ( "net/http/httptest" "net/url" "os" + "sort" + "strconv" "strings" "sync" "testing" @@ -42,6 +44,27 @@ type fakeControlAPI struct { registry map[string]registryItem registryChecks map[string][]check + // Scan job state machine (see the scan handlers): POST starts a job that + // stays "running" for scanRunPolls GET polls (0 = finishes at once), then + // ends as done — or as error when scanFinalError is set. scan is the status + // served by GET; tests may also set it directly (with scanPollsLeft == 0 it + // stays as is). scanStartStatus != 0 makes POST fail with that HTTP status. + scan scanStatusDTO + scanRunPolls int + scanPollsLeft int + scanFinalError string + scanStartStatus int + + // Requests seen on the list endpoints, for "never loads everything" checks: + // the raw query of every GET /ips (ipsQueries) and the number of GET + // /registry calls without `limit` (bare), plus the sizes of the DeleteIPs / + // SubmitIPs bulk calls. + ipsQueries []string + registryQueries []string + bareIPsCalls int + deleteChunks []int + submitChunks []int + // autoCycle is the state served by /api/v1/admin/auto-cycle*; // autoCycleDown makes all four endpoints answer 500 (unavailable API). autoCycle autoCycleDTO @@ -86,16 +109,72 @@ func (f *fakeControlAPI) handler() http.Handler { f.mu.Lock() defer f.mu.Unlock() byState := map[string]int{} + results := map[string]int{"pass": 0, "partial": 0, "fail": 0, "cancelled": 0} for _, ip := range f.ips { byState[ip.State]++ + if ip.OverallResult != "" { + results[ip.OverallResult]++ + } } - writeJSON(w, http.StatusOK, statusResponse{TotalIPs: len(f.ips), IPsByState: byState, TotalValidators: len(f.validators)}) + writeJSON(w, http.StatusOK, statusResponse{TotalIPs: len(f.ips), IPsByState: byState, TotalValidators: len(f.validators), ResultsByOverall: results}) }) mux.HandleFunc("GET /api/v1/admin/ips", func(w http.ResponseWriter, r *http.Request) { f.mu.Lock() defer f.mu.Unlock() - writeJSON(w, http.StatusOK, f.ips) + f.ipsQueries = append(f.ipsQueries, r.URL.RawQuery) + qv := r.URL.Query() + if qv.Get("limit") == "" { + // Legacy shape: the whole queue as a bare array. + f.bareIPsCalls++ + writeJSON(w, http.StatusOK, f.ips) + return + } + limit, err := strconv.Atoi(qv.Get("limit")) + if err != nil || limit < 1 || limit > 1000 { + writeAPIErr(w, http.StatusBadRequest, "limit must be 1..1000") + return + } + offset, _ := strconv.Atoi(qv.Get("offset")) + var states map[string]bool + if st := qv.Get("state"); st != "" { + states = map[string]bool{} + for _, x := range strings.Split(st, ",") { + states[x] = true + } + } + var matched []ipQueueItem + for _, ip := range f.ips { + if states != nil && !states[ip.State] { + continue + } + if q := qv.Get("q"); q != "" && !strings.Contains(strings.ToLower(ip.IPAddress), strings.ToLower(q)) { + continue + } + if res := qv.Get("result"); res != "" && ip.OverallResult != res { + continue + } + matched = append(matched, ip) + } + if qv.Get("order") == "aggregated_at_desc" { + sort.SliceStable(matched, func(i, j int) bool { + a, b := matched[i].AggregatedAt, matched[j].AggregatedAt + if a == nil || b == nil { + return a != nil && b == nil + } + return a.After(*b) + }) + } else { + sort.SliceStable(matched, func(i, j int) bool { return matched[i].Sequence < matched[j].Sequence }) + } + page := []ipQueueItem{} + if offset < len(matched) { + page = matched[offset:] + if len(page) > limit { + page = page[:limit] + } + } + writeJSON(w, http.StatusOK, ipsPage{Items: page, Total: len(matched), Limit: limit, Offset: offset}) }) mux.HandleFunc("GET /api/v1/admin/ips/{ip}", func(w http.ResponseWriter, r *http.Request) { @@ -122,6 +201,7 @@ func (f *fakeControlAPI) handler() http.Handler { writeAPIErr(w, http.StatusBadRequest, "addresses must not be empty") return } + f.submitChunks = append(f.submitChunks, len(req.Addresses)) resp := submitIPsResponse{} for _, addr := range req.Addresses { idx := f.findIP(addr) @@ -190,6 +270,7 @@ func (f *fakeControlAPI) handler() http.Handler { writeAPIErr(w, http.StatusBadRequest, "addresses must not be empty") return } + f.deleteChunks = append(f.deleteChunks, len(req.Addresses)) resp := deleteIPsResponse{} for _, addr := range req.Addresses { idx := f.findIP(addr) @@ -210,6 +291,7 @@ func (f *fakeControlAPI) handler() http.Handler { for _, ip := range f.ips { resp.Deleted = append(resp.Deleted, ip.IPAddress) } + resp.Count = len(resp.Deleted) f.ips = nil writeJSON(w, http.StatusOK, resp) }) @@ -246,18 +328,32 @@ func (f *fakeControlAPI) handler() http.Handler { 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) + if f.scanStartStatus != 0 { + writeAPIErr(w, f.scanStartStatus, "scan refused") + return } - writeJSON(w, http.StatusOK, resp) + if !f.scan.Running { + now := time.Now() + f.scan = scanStatusDTO{State: "listing", Running: true, DryRun: r.URL.Query().Get("dry_run") == "true", StartedAt: &now} + f.scanPollsLeft = f.scanRunPolls + if f.scanPollsLeft == 0 { + f.finishScan() + } + } + writeJSON(w, http.StatusAccepted, f.scan) + }) + mux.HandleFunc("GET /api/v1/admin/ips/scan", func(w http.ResponseWriter, r *http.Request) { + f.mu.Lock() + defer f.mu.Unlock() + if f.scan.Running && f.scanPollsLeft > 0 { + f.scanPollsLeft-- + if f.scanPollsLeft == 0 { + f.finishScan() + } else { + f.scan.Pages++ + } + } + writeJSON(w, http.StatusOK, f.scan) }) mux.HandleFunc("GET /api/v1/admin/auto-cycle", func(w http.ResponseWriter, r *http.Request) { @@ -329,11 +425,37 @@ func (f *fakeControlAPI) handler() http.Handler { 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)) + qv := r.URL.Query() + f.registryQueries = append(f.registryQueries, r.URL.RawQuery) + matched := make([]registryItem, 0, len(f.registry)) for _, item := range f.registry { - out = append(out, item) + if q := qv.Get("q"); q != "" && !strings.Contains(strings.ToLower(item.IPAddress), strings.ToLower(q)) { + continue + } + if lr := qv.Get("last_result"); lr != "" && item.LastResult != lr { + continue + } + matched = append(matched, item) } - writeJSON(w, http.StatusOK, out) + sort.Slice(matched, func(i, j int) bool { return matched[i].IPAddress < matched[j].IPAddress }) + if qv.Get("limit") == "" { + writeJSON(w, http.StatusOK, matched) + return + } + limit, err := strconv.Atoi(qv.Get("limit")) + if err != nil || limit < 1 || limit > 1000 { + writeAPIErr(w, http.StatusBadRequest, "limit must be 1..1000") + return + } + offset, _ := strconv.Atoi(qv.Get("offset")) + page := []registryItem{} + if offset < len(matched) { + page = matched[offset:] + if len(page) > limit { + page = page[:limit] + } + } + writeJSON(w, http.StatusOK, registryPage{Items: page, Total: len(matched), Limit: limit, Offset: offset}) }) mux.HandleFunc("GET /api/v1/admin/registry/{ip}", func(w http.ResponseWriter, r *http.Request) { f.mu.Lock() @@ -544,6 +666,51 @@ func (f *fakeControlAPI) handler() http.Handler { }) } +// finishScan ends the running scan job (caller holds f.mu): the discovered +// free addresses are queued unless it was a dry run, or the job fails with +// scanFinalError. +func (f *fakeControlAPI) finishScan() { + now := time.Now() + f.scan.Running = false + f.scan.FinishedAt = &now + f.scan.Discovered = len(f.scanFreeAddresses) + f.scan.Free = len(f.scanFreeAddresses) + if f.scanFinalError != "" { + f.scan.State = "error" + f.scan.Error = f.scanFinalError + return + } + f.scan.State = "done" + if f.scan.DryRun { + return + } + for _, addr := range f.scanFreeAddresses { + if f.findIP(addr) < 0 { + f.ips = append(f.ips, ipQueueItem{IPAddress: addr, State: "queued", Sequence: len(f.ips) + 1, CreatedAt: now, UpdatedAt: now}) + f.scan.Added++ + } else { + f.scan.Reordered++ + } + } +} + +// seedIPs appends n queued addresses 10.x.y.z (distinct, in sequence order) +// and returns them. Use it for tests that need pages' worth of rows. +func (f *fakeControlAPI) seedIPs(n int) []string { + f.mu.Lock() + defer f.mu.Unlock() + now := time.Now() + out := make([]string, 0, n) + base := len(f.ips) + for i := 0; i < n; i++ { + k := base + i + 1 + addr := fmt.Sprintf("10.%d.%d.%d", k/65536, (k/256)%256, k%256) + f.ips = append(f.ips, ipQueueItem{IPAddress: addr, State: "queued", Sequence: k, CreatedAt: now, UpdatedAt: now}) + out = append(out, addr) + } + return out +} + func (f *fakeControlAPI) findIP(addr string) int { for i, ip := range f.ips { if ip.IPAddress == addr { @@ -596,3 +763,14 @@ func postForm(t *testing.T, ts *httptest.Server, method, path string, form url.V body, _ := io.ReadAll(resp.Body) return string(body) } + +// newSlowAPI serves an empty JSON object for every request after delay. +func newSlowAPI(t *testing.T, delay time.Duration) string { + t.Helper() + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + time.Sleep(delay) + writeJSON(w, http.StatusOK, map[string]interface{}{}) + })) + t.Cleanup(ts.Close) + return ts.URL +} diff --git a/internal/dashboard/dto.go b/internal/dashboard/dto.go index 84416c0..81ab8d1 100644 --- a/internal/dashboard/dto.go +++ b/internal/dashboard/dto.go @@ -1,6 +1,7 @@ package dashboard import ( + "fmt" "strconv" "time" ) @@ -20,6 +21,27 @@ type statusResponse struct { TotalIPs int `json:"total_ips"` IPsByState map[string]int `json:"ips_by_state"` TotalValidators int `json:"total_validators"` + // ResultsByOverall counts finished addresses by overall result + // (pass/partial/fail/cancelled). + ResultsByOverall map[string]int `json:"results_by_overall"` +} + +// Terminal queue states: the address needs no further processing. Shared by +// the overview progress indicator; "occupied" counts as terminal too (the +// check cycle never ran because the floating IP was already bound). +var terminalStates = []string{"done", "failed", "occupied"} + +// activeStates are the states of an address that is being worked on right +// now (everything between "queued" and a terminal state). +var activeStates = []string{"assigning_fip", "awaiting_self_check", "checking", "aggregating"} + +// sumStates adds up the counts of the given states in a status breakdown. +func sumStates(byState map[string]int, states []string) int { + n := 0 + for _, s := range states { + n += byState[s] + } + return n } type ipQueueItem struct { @@ -99,14 +121,135 @@ type deleteIPsResponse struct { type clearQueueResponse struct { Deleted []string `json:"deleted"` + Count int `json:"count"` } -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"` +// ipsPage is the paginated envelope of GET /admin/ips (sent when the request +// carries `limit`). +type ipsPage struct { + Items []ipQueueItem `json:"items"` + Total int `json:"total"` + Limit int `json:"limit"` + Offset int `json:"offset"` +} + +// registryPage is the paginated envelope of GET /admin/registry. +type registryPage struct { + Items []registryItem `json:"items"` + Total int `json:"total"` + Limit int `json:"limit"` + Offset int `json:"offset"` +} + +// scanStatusDTO is the state of control-api's background floating-IP scan job +// (POST/GET /api/v1/admin/ips/scan). +type scanStatusDTO struct { + State string `json:"state"` + Running bool `json:"running"` + DryRun bool `json:"dry_run"` + Pages int `json:"pages"` + Discovered int `json:"discovered"` + Free int `json:"free"` + Added int `json:"added"` + Requeued int `json:"requeued"` + Reordered int `json:"reordered"` + SkippedInProgress int `json:"skipped_in_progress"` + StartedAt *time.Time `json:"started_at"` + FinishedAt *time.Time `json:"finished_at"` + Error string `json:"error"` +} + +// Finished reports a job that has run to a terminal state (as opposed to +// "idle" = never started, or still running). +func (s scanStatusDTO) Finished() bool { + if s.Running { + return false + } + switch s.State { + case "done", "error", "cancelled": + return true + } + return false +} + +// StateLabel is the Russian description of the job's state. +func (s scanStatusDTO) StateLabel() string { + switch s.State { + case "clearing": + return "очистка" + case "listing": + return "читаются страницы" + case "enqueuing": + return "ставятся в очередь" + case "done": + return "готово" + case "error": + return "ошибка" + case "cancelled": + return "отменено" + case "idle", "": + return "нет активного сканирования" + default: + return s.State + } +} + +// PillClass picks the pill style for the state. +func (s scanStatusDTO) PillClass() string { + switch s.State { + case "done": + return "pill-success" + case "error": + return "pill-danger" + case "cancelled": + return "pill-cancel" + case "clearing", "listing", "enqueuing": + return "pill-info" + default: + return "pill-neutral" + } +} + +// Handled is how many of the free addresses the enqueuing phase has already +// processed. +func (s scanStatusDTO) Handled() int { + return s.Added + s.Requeued + s.Reordered + s.SkippedInProgress +} + +// Indeterminate is true while the amount of work is not known yet. +func (s scanStatusDTO) Indeterminate() bool { + return s.State == "clearing" || s.State == "listing" || (s.State == "enqueuing" && s.Free <= 0) +} + +// Elapsed is the human-readable run time: until now while running, until +// finished_at afterwards; empty when the job never started. +func (s scanStatusDTO) Elapsed() string { + if s.StartedAt == nil { + return "" + } + end := time.Now() + if !s.Running && s.FinishedAt != nil { + end = *s.FinishedAt + } + return fmtDuration(end.Sub(*s.StartedAt)) +} + +// fmtDuration renders a duration in Russian as "2 ч 05 мин", "3 мин 07 с" or +// "42 с" — coarse on purpose (progress/ETA display). +func fmtDuration(d time.Duration) string { + if d < 0 { + d = 0 + } + sec := int(d.Round(time.Second) / time.Second) + h, m, sc := sec/3600, (sec%3600)/60, sec%60 + switch { + case h > 0: + return fmt.Sprintf("%d ч %02d мин", h, m) + case m > 0: + return fmt.Sprintf("%d мин %02d с", m, sc) + default: + return fmt.Sprintf("%d с", sc) + } } // registryItem is one row of the durable per-address registry — see @@ -200,6 +343,8 @@ func (a autoCycleDTO) MaxRunMinutes() string { return secondsToMinutes(a.MaxRu // PhaseLabel is the Russian description of the current phase. func (a autoCycleDTO) PhaseLabel() string { switch a.Phase { + case "scanning": + return "сканирование Floating IP" case "running": return "идёт проверка" case "waiting": diff --git a/internal/dashboard/handlers_ips.go b/internal/dashboard/handlers_ips.go index 9644689..4f95281 100644 --- a/internal/dashboard/handlers_ips.go +++ b/internal/dashboard/handlers_ips.go @@ -1,14 +1,122 @@ package dashboard import ( + "context" + "errors" "fmt" "net/http" + "net/url" + "strings" ) +// bulkChunk is how many addresses go into one DeleteIPs/SubmitIPs call when +// an operation spans a whole filter ("scope=all"). +const bulkChunk = 500 + +// ipStates are the queue states control-api accepts in the `state` filter. +var ipStates = []string{"queued", "assigning_fip", "awaiting_self_check", "checking", "aggregating", "done", "failed", "occupied"} + +var ipResults = []string{"pass", "partial", "fail", "cancelled"} + +func containsStr(list []string, s string) bool { + for _, v := range list { + if v == s { + return true + } + } + return false +} + +// ipsFilter is the server-side filter of the /ips list. State is a filter +// token: one of "", "queued", "active" (every in-progress state), "done", +// "failed", "occupied", or a comma-separated list of raw queue states; the +// result filter is separate (Result) and may be combined with it. +type ipsFilter struct { + Q string + State string + Result string +} + +// parseIPsFilter reads q/state/result from the request (query or form body). +// The state option matching the filter. +func (f ipsFilter) Token() string { + if f.State == "" && f.Result != "" { + return "result:" + f.Result + } + return f.State +} + +func (f ipsFilter) values() url.Values { + v := url.Values{} + if f.Q != "" { + v.Set("q", f.Q) + } + if f.State != "" { + v.Set("state", f.State) + } + if f.Result != "" { + v.Set("result", f.Result) + } + return v +} + type ipsPageData struct { PageData Items []ipQueueItem FIPSettleSeconds int + Filter ipsFilter + Page, PerPage int + // Total is the number of rows matching the filter; QueueTotal is the whole + // queue (what "Очистить всё" would remove). + Total, QueueTotal int + Pager pagerData + // SelfURL is this page's own URL (filter + page), re-requested to reload + // the table when a scan finishes. + SelfURL string + PerPageOptions []int + Scan scanProgressData } type ipDetailData struct { @@ -16,13 +124,67 @@ type ipDetailData struct { Detail ipDetailResponse } -func (s *Server) handleIPsPage(w http.ResponseWriter, r *http.Request) { - items, err := s.CA.ListIPs(r.Context()) - settings, settingsErr := s.CA.GetOrchestratorSettings(r.Context()) +// loadIPsData fetches the page of the queue selected by the request's +// page/per_page/q/state/result params. A page that became empty (e.g. after +// deleting its last rows) is clamped to the last non-empty page. +func (s *Server) loadIPsData(r *http.Request) (ipsPageData, error) { + ctx := r.Context() + f := parseIPsFilter(r) + perPage := parsePerPage(r.FormValue("per_page")) + page := parsePage(r.FormValue("page")) + q := ipsQuery{States: f.States(), Q: f.Q, Result: f.Result, Order: "sequence", Limit: perPage, Offset: (page - 1) * perPage} + + res, err := s.CA.ListIPsPage(ctx, q) if err == nil { - err = settingsErr + if clamped := clampPage(page, res.Total, perPage); clamped != page { + page = clamped + q.Offset = (page - 1) * perPage + res, err = s.CA.ListIPsPage(ctx, q) + } + } + data := ipsPageData{ + Items: res.Items, + Filter: f, + Page: page, + PerPage: perPage, + Total: res.Total, + QueueTotal: res.Total, + PerPageOptions: perPageOptions, + } + data.Pager = newPager("/ips", "ips-table-wrap", f.values(), page, perPage, res.Total) + data.SelfURL = pageURL("/ips", f.values(), page, perPage) + + if settings, settingsErr := s.CA.GetOrchestratorSettings(ctx); settingsErr != nil { + if err == nil { + err = settingsErr + } + } else { + data.FIPSettleSeconds = settings.FIPSettleSeconds + } + if err == nil && f.Active() { + // "Очистить всё" ignores the filter: show the real queue size in its + // confirmation. Non-fatal — the label just falls back to the filtered total. + if st, stErr := s.CA.Status(ctx); stErr == nil { + data.QueueTotal = st.TotalIPs + } + } + return data, err +} + +func (s *Server) handleIPsPage(w http.ResponseWriter, r *http.Request) { + data, err := s.loadIPsData(r) + // A filter/pager request from htmx swaps only #ips-table-wrap (hx-select), + // so there is no need to re-render the whole page (and re-query the scan + // status). A history-restore fetch needs the full page. + if r.Header.Get("HX-Request") == "true" && r.Header.Get("HX-History-Restore-Request") != "true" { + s.renderFragment(w, "ips_table_wrap", data, err) + return + } + if st, scanErr := s.CA.ScanStatus(r.Context()); scanErr != nil { + s.Log.Warn("ips: scan status unavailable", "err", scanErr) + } else { + data.Scan = newScanProgress(st) } - data := ipsPageData{Items: items, FIPSettleSeconds: settings.FIPSettleSeconds} data.ActiveNav = "ips" data.Banner = bannerFor(err) s.renderPage(w, r, "ips_page", data) @@ -37,20 +199,18 @@ func (s *Server) handleIPDetail(w http.ResponseWriter, r *http.Request) { s.renderPage(w, r, "ip_detail_page", data) } -// renderIPsTable re-fetches the current queue and renders the ips_table -// fragment, tagging actionErr (if any) on the shared error banner. Called -// after every mutating /ips/* request so the table always reflects true -// current state regardless of whether the mutation itself succeeded. +// renderIPsTable re-fetches the current page of the queue (same page/filter as +// the request, carried in hidden #ips-form inputs) and renders the +// ips_table_wrap fragment, tagging actionErr (if any) on the shared error +// banner. Called after every mutating /ips/* request so the table always +// reflects true current state regardless of whether the mutation itself +// succeeded. func (s *Server) renderIPsTable(w http.ResponseWriter, r *http.Request, actionErr error) { - items, listErr := s.CA.ListIPs(r.Context()) + data, err := s.loadIPsData(r) if actionErr == nil { - actionErr = listErr + actionErr = err } - settings, settingsErr := s.CA.GetOrchestratorSettings(r.Context()) - if actionErr == nil { - actionErr = settingsErr - } - s.renderFragment(w, "ips_table", ipsPageData{Items: items, FIPSettleSeconds: settings.FIPSettleSeconds}, actionErr) + s.renderFragment(w, "ips_table_wrap", data, actionErr) } func (s *Server) handleIPsSubmit(w http.ResponseWriter, r *http.Request) { @@ -63,7 +223,7 @@ func (s *Server) handleIPsSubmit(w http.ResponseWriter, r *http.Request) { s.renderIPsTable(w, r, &apiErr{Status: http.StatusBadRequest, Message: "укажите хотя бы один адрес"}) return } - _, err := s.CA.SubmitIPs(r.Context(), addresses) + err := s.submitChunked(r.Context(), addresses) s.renderIPsTable(w, r, err) } @@ -85,31 +245,76 @@ func (s *Server) handleIPDelete(w http.ResponseWriter, r *http.Request) { s.renderIPsTable(w, r, err) } -func (s *Server) handleIPsDeleteSelected(w http.ResponseWriter, r *http.Request) { +// submitChunked feeds addresses to SubmitIPs in chunks of bulkChunk so one +// huge list never becomes one huge request/transaction. +func (s *Server) submitChunked(ctx context.Context, addresses []string) error { + for _, part := range chunk(addresses, bulkChunk) { + if _, err := s.CA.SubmitIPs(ctx, part); err != nil { + return err + } + } + return nil +} + +// resolveFilterAddresses lists every address matching filter by paging +// ListIPsPage (≤ maxPageLimit rows per call), for "select all N by filter". +func (s *Server) resolveFilterAddresses(ctx context.Context, f ipsFilter) ([]string, error) { + var out []string + for offset := 0; ; { + page, err := s.CA.ListIPsPage(ctx, ipsQuery{States: f.States(), Q: f.Q, Result: f.Result, Order: "sequence", Limit: maxPageLimit, Offset: offset}) + if err != nil { + return nil, err + } + for _, it := range page.Items { + out = append(out, it.IPAddress) + } + offset += len(page.Items) + if len(page.Items) == 0 || offset >= page.Total { + return out, nil + } + } +} + +// bulkAddresses returns the addresses a bulk delete/recheck acts on: the +// checked rows of the current page, or — with scope=all — everything that +// matches the current filter, resolved server-side. +func (s *Server) bulkAddresses(r *http.Request) ([]string, error) { if err := r.ParseForm(); err != nil { - s.renderIPsTable(w, r, fmt.Errorf("invalid form: %w", err)) - return + return nil, fmt.Errorf("invalid form: %w", err) + } + var addresses []string + if r.FormValue("scope") == "all" { + var err error + addresses, err = s.resolveFilterAddresses(r.Context(), parseIPsFilter(r)) + if err != nil { + return nil, err + } + } else { + addresses = r.Form["addresses"] } - addresses := r.Form["addresses"] if len(addresses) == 0 { - s.renderIPsTable(w, r, &apiErr{Status: http.StatusBadRequest, Message: "ничего не выбрано"}) - return + return nil, &apiErr{Status: http.StatusBadRequest, Message: "ничего не выбрано"} + } + return addresses, nil +} + +func (s *Server) handleIPsDeleteSelected(w http.ResponseWriter, r *http.Request) { + addresses, err := s.bulkAddresses(r) + if err == nil { + for _, part := range chunk(addresses, bulkChunk) { + if _, err = s.CA.DeleteIPs(r.Context(), part); err != nil { + break + } + } } - _, err := s.CA.DeleteIPs(r.Context(), addresses) s.renderIPsTable(w, r, err) } func (s *Server) handleIPsRecheckSelected(w http.ResponseWriter, r *http.Request) { - if err := r.ParseForm(); err != nil { - s.renderIPsTable(w, r, fmt.Errorf("invalid form: %w", err)) - return + addresses, err := s.bulkAddresses(r) + if err == nil { + err = s.submitChunked(r.Context(), addresses) } - addresses := r.Form["addresses"] - if len(addresses) == 0 { - s.renderIPsTable(w, r, &apiErr{Status: http.StatusBadRequest, Message: "ничего не выбрано"}) - return - } - _, err := s.CA.SubmitIPs(r.Context(), addresses) s.renderIPsTable(w, r, err) } @@ -118,10 +323,46 @@ func (s *Server) handleIPsClear(w http.ResponseWriter, r *http.Request) { 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. +// scanProgressData drives the scan_progress partial. Poll keeps the partial's +// own hx-trigger="every 2s" alive: while the job runs, or while control-api is +// transiently unreachable (so one failed poll doesn't freeze the panel). +type scanProgressData struct { + Status scanStatusDTO + Poll bool +} + +func newScanProgress(st scanStatusDTO) scanProgressData { + return scanProgressData{Status: st, Poll: st.Running} +} + +// renderScanProgress renders the scan panel fragment. A finished job also +// sends `HX-Trigger: scan-finished`, which makes #ips-table-wrap reload. +// Errors go to the shared banner; only transient ones (transport/5xx) keep +// polling when keepPolling is set. +func (s *Server) renderScanProgress(w http.ResponseWriter, st scanStatusDTO, err error, keepPolling bool) { + data := newScanProgress(st) + if err != nil { + data = scanProgressData{} + var ae *apiErr + if keepPolling && errors.As(err, &ae) && (ae.Status == 0 || ae.Status >= 500) { + data.Poll = true + } + } else if st.Finished() { + w.Header().Set("HX-Trigger", "scan-finished") + } + s.renderFragment(w, "scan_progress", data, err) +} + +// handleIPsScan starts the background floating-IP scan (or joins the running +// one) and returns the progress panel at once — the job itself runs in +// control-api, so this never waits for OpenStack. ?dry_run=true only counts. func (s *Server) handleIPsScan(w http.ResponseWriter, r *http.Request) { - _, err := s.CA.ScanFloatingIPs(r.Context()) - s.renderIPsTable(w, r, err) + st, err := s.CA.StartScan(r.Context(), r.URL.Query().Get("dry_run") == "true") + s.renderScanProgress(w, st, err, false) +} + +// handleIPsScanStatus is the progress panel's poll target. +func (s *Server) handleIPsScanStatus(w http.ResponseWriter, r *http.Request) { + st, err := s.CA.ScanStatus(r.Context()) + s.renderScanProgress(w, st, err, true) } diff --git a/internal/dashboard/handlers_overview.go b/internal/dashboard/handlers_overview.go index 6eec7af..5062ec5 100644 --- a/internal/dashboard/handlers_overview.go +++ b/internal/dashboard/handlers_overview.go @@ -2,23 +2,51 @@ package dashboard import ( "net/http" - "sort" "strings" + "time" ) +const ( + // overviewActiveLimit caps the "В работе" table; overviewQueueLimit caps + // the compact list of the next queued addresses. The queue itself can hold + // thousands of rows, so the overview never lists more than these. + overviewActiveLimit = 100 + overviewQueueLimit = 10 + // etaMinSamples is how many completed rows (with AggregatedAt) the ETA + // needs to derive a rate from. + etaMinSamples = 5 +) + +// overviewProgress is the "Готово D из T" indicator of the stats block. +type overviewProgress struct { + Total, Done, Active, Queued, Percent int + // ETA is a rough time-to-finish estimate ("" when not enough data). + ETA string +} + type overviewData struct { PageData - Status statusResponse - CurrentItems []ipQueueItem + Status statusResponse + // ActiveItems are the addresses being checked right now (≤ overviewActiveLimit + // of ActiveTotal); QueuedItems are the next few waiting ones (QueuedTotal + // in total). Both stay empty under a result filter: an address without a + // verdict can't match one. + ActiveItems []ipQueueItem + ActiveTotal int + QueuedItems []ipQueueItem + QueuedTotal int LastCompleted []ipQueueItem Breakdown map[string]int LastN int PollSeconds int Query string StatusFilter string + Progress overviewProgress // AutoCycle is nil when the auto-cycle status could not be fetched; the - // indicator is then simply hidden (a non-fatal failure). + // indicator is then simply hidden (a non-fatal failure). Scan is likewise + // nil when the scan status is unavailable. AutoCycle *autoCycleDTO + Scan *scanStatusDTO } func (s *Server) loadOverview(r *http.Request) (overviewData, error) { @@ -27,30 +55,58 @@ func (s *Server) loadOverview(r *http.Request) (overviewData, error) { if err != nil { return overviewData{}, err } - ips, err := s.CA.ListIPs(ctx) + q := strings.TrimSpace(r.URL.Query().Get("q")) + resultFilter := r.URL.Query().Get("status") + if !containsStr(ipResults, resultFilter) { + resultFilter = "" + } + lastN := s.Cfg.LastCompletedCount + if lastN < 1 { + lastN = 20 + } + + data := overviewData{ + Status: status, + LastN: lastN, + PollSeconds: s.Cfg.OverviewPollIntervalS, + Query: q, + StatusFilter: resultFilter, + } + + if resultFilter == "" { + active, err := s.CA.ListIPsPage(ctx, ipsQuery{States: activeStates, Q: q, Order: "sequence", Limit: overviewActiveLimit}) + if err != nil { + return overviewData{}, err + } + queued, err := s.CA.ListIPsPage(ctx, ipsQuery{States: []string{"queued"}, Q: q, Order: "sequence", Limit: overviewQueueLimit}) + if err != nil { + return overviewData{}, err + } + data.ActiveItems, data.ActiveTotal = active.Items, active.Total + data.QueuedItems, data.QueuedTotal = queued.Items, queued.Total + } + completed, err := s.CA.ListIPsPage(ctx, ipsQuery{States: []string{"done", "failed"}, Q: q, Result: resultFilter, Order: "aggregated_at_desc", Limit: lastN}) if err != nil { return overviewData{}, err } - var autoCycle *autoCycleDTO + data.LastCompleted = completed.Items + data.Breakdown = resultBreakdown(completed.Items) + + // ETA from the unfiltered window only: a filtered list is not a sample of + // the checks' real throughput. + data.Progress = computeProgress(status, completed.Items, q == "" && resultFilter == "") + if ac, acErr := s.CA.GetAutoCycle(ctx); acErr != nil { s.Log.Warn("overview: auto-cycle status unavailable", "err", acErr) } else { - autoCycle = &ac + data.AutoCycle = &ac } - q := strings.TrimSpace(r.URL.Query().Get("q")) - resultFilter := r.URL.Query().Get("status") - last := lastCompleted(ips, s.Cfg.LastCompletedCount) - return overviewData{ - Status: status, - CurrentItems: filterQueueItems(currentlyChecking(ips), q, resultFilter), - LastCompleted: filterQueueItems(last, q, resultFilter), - Breakdown: resultBreakdown(last), - LastN: s.Cfg.LastCompletedCount, - PollSeconds: s.Cfg.OverviewPollIntervalS, - Query: q, - StatusFilter: resultFilter, - AutoCycle: autoCycle, - }, nil + if sc, scErr := s.CA.ScanStatus(ctx); scErr != nil { + s.Log.Warn("overview: scan status unavailable", "err", scErr) + } else if sc.Running { + data.Scan = &sc + } + return data, nil } func (s *Server) handleOverview(w http.ResponseWriter, r *http.Request) { @@ -75,40 +131,45 @@ func (s *Server) handleOverviewFragment(w http.ResponseWriter, r *http.Request) } } -// currentlyChecking is every IP not yet in a terminal state, ordered by -// queue position — the "текущая проверка" live snapshot. No backend -// concept of a "run" exists; this is computed fresh on every request. -func currentlyChecking(ips []ipQueueItem) []ipQueueItem { - var out []ipQueueItem - for _, ip := range ips { - if ip.State != "done" && ip.State != "failed" { - out = append(out, ip) +// computeProgress derives the done/active/queued counters from the status +// breakdown (done = every terminal state, "occupied" included) and, when +// useETA is set and the newest-first list of completed rows has at least +// etaMinSamples with AggregatedAt, a rough ETA from their completion rate. +func computeProgress(st statusResponse, completed []ipQueueItem, useETA bool) overviewProgress { + p := overviewProgress{ + Total: st.TotalIPs, + Done: sumStates(st.IPsByState, terminalStates), + Active: sumStates(st.IPsByState, activeStates), + Queued: st.IPsByState["queued"], + } + if p.Total > 0 { + p.Percent = p.Done * 100 / p.Total + } + remaining := p.Active + p.Queued + if !useETA || remaining == 0 { + return p + } + var stamps []time.Time + for _, ip := range completed { + if ip.AggregatedAt != nil { + stamps = append(stamps, *ip.AggregatedAt) } } - sort.Slice(out, func(i, j int) bool { return out[i].Sequence < out[j].Sequence }) - return out -} - -// lastCompleted returns the n most recently completed (done/failed) IPs by -// AggregatedAt descending — the "последняя завершённая проверка" summary -// window. This is an operational definition, not a real "batch": resubmit -// n if the window size needs tuning (overview.last_completed_count). -func lastCompleted(ips []ipQueueItem, n int) []ipQueueItem { - var done []ipQueueItem - for _, ip := range ips { - if (ip.State == "done" || ip.State == "failed") && ip.AggregatedAt != nil { - done = append(done, ip) - } + if len(stamps) < etaMinSamples { + return p } - sort.Slice(done, func(i, j int) bool { return done[i].AggregatedAt.After(*done[j].AggregatedAt) }) - if len(done) > n { - done = done[:n] + // completed is newest-first: stamps[0] is the newest, the last the oldest. + span := stamps[0].Sub(stamps[len(stamps)-1]) + if span <= 0 { + return p } - return done + perItem := span / time.Duration(len(stamps)-1) + p.ETA = fmtDuration(perItem * time.Duration(remaining)) + return p } // resultBreakdown counts OverallResult values across exactly the given -// items (normally the output of lastCompleted) — pass/partial/fail/cancelled. +// items (the "последние N завершённых" window) — pass/partial/fail/cancelled. func resultBreakdown(items []ipQueueItem) map[string]int { out := map[string]int{"pass": 0, "partial": 0, "fail": 0, "cancelled": 0} for _, ip := range items { @@ -116,26 +177,3 @@ func resultBreakdown(items []ipQueueItem) map[string]int { } return out } - -// filterQueueItems narrows items to those whose address contains q -// (case-insensitive substring) and, if status is set, whose OverallResult -// matches it exactly. A still-in-progress item always has an empty -// OverallResult, so picking any specific status hides it — the intended -// behavior for "Текущая проверка", which has no verdict yet. -func filterQueueItems(items []ipQueueItem, q, status string) []ipQueueItem { - if q == "" && status == "" { - return items - } - q = strings.ToLower(q) - out := make([]ipQueueItem, 0, len(items)) - for _, ip := range items { - if q != "" && !strings.Contains(strings.ToLower(ip.IPAddress), q) { - continue - } - if status != "" && ip.OverallResult != status { - continue - } - out = append(out, ip) - } - return out -} diff --git a/internal/dashboard/handlers_registry.go b/internal/dashboard/handlers_registry.go index 45d1811..56d9f89 100644 --- a/internal/dashboard/handlers_registry.go +++ b/internal/dashboard/handlers_registry.go @@ -2,14 +2,19 @@ package dashboard import ( "net/http" + "net/url" "strings" ) type registryPageData struct { PageData - Items []registryItem - Query string - StatusFilter string + Items []registryItem + Query string + StatusFilter string + Page, PerPage int + Total int + Pager pagerData + PerPageOptions []int } type registryDetailData struct { @@ -22,43 +27,48 @@ type registryDetailData struct { // record that survives an address being deleted from /ips and later // re-added. See internal/db/migrations/0007_ip_registry.sql. Optional // ?q=&status= query params narrow the list by address substring and by -// LastResult — see filterRegistryItems. +// LastResult, and ?page=&per_page= select a page — all applied server-side +// (control-api's ListRegistryPage), so only the visible rows are transferred. func (s *Server) handleRegistryPage(w http.ResponseWriter, r *http.Request) { - items, err := s.CA.ListRegistry(r.Context()) q := strings.TrimSpace(r.URL.Query().Get("q")) status := r.URL.Query().Get("status") + if !containsStr(ipResults, status) { + status = "" + } + perPage := parsePerPage(r.URL.Query().Get("per_page")) + page := parsePage(r.URL.Query().Get("page")) + query := registryQuery{Q: q, LastResult: status, Limit: perPage, Offset: (page - 1) * perPage} + + res, err := s.CA.ListRegistryPage(r.Context(), query) + if err == nil { + if clamped := clampPage(page, res.Total, perPage); clamped != page { + page = clamped + query.Offset = (page - 1) * perPage + res, err = s.CA.ListRegistryPage(r.Context(), query) + } + } + params := url.Values{} + if q != "" { + params.Set("q", q) + } + if status != "" { + params.Set("status", status) + } data := registryPageData{ - Items: filterRegistryItems(items, q, status), - Query: q, - StatusFilter: status, + Items: res.Items, + Query: q, + StatusFilter: status, + Page: page, + PerPage: perPage, + Total: res.Total, + Pager: newPager("/registry", "registry-table-wrap", params, page, perPage, res.Total), + PerPageOptions: perPageOptions, } data.ActiveNav = "registry" data.Banner = bannerFor(err) s.renderPage(w, r, "registry_page", data) } -// filterRegistryItems narrows items to those whose address contains q -// (case-insensitive substring) and, if status is set, whose LastResult -// matches it exactly — the registry list's search-by-IP and -// filter-by-status, mirroring filterQueueItems in handlers_overview.go. -func filterRegistryItems(items []registryItem, q, status string) []registryItem { - if q == "" && status == "" { - return items - } - q = strings.ToLower(q) - out := make([]registryItem, 0, len(items)) - for _, it := range items { - if q != "" && !strings.Contains(strings.ToLower(it.IPAddress), q) { - continue - } - if status != "" && it.LastResult != status { - continue - } - out = append(out, it) - } - return out -} - // 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. diff --git a/internal/dashboard/handlers_test.go b/internal/dashboard/handlers_test.go index 92508d9..e7f1057 100644 --- a/internal/dashboard/handlers_test.go +++ b/internal/dashboard/handlers_test.go @@ -558,20 +558,6 @@ 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() @@ -667,64 +653,6 @@ func TestControlAPIUnreachable(t *testing.T) { } } -func TestFilterQueueItems(t *testing.T) { - items := []ipQueueItem{ - {IPAddress: "1.1.1.1", OverallResult: "pass"}, - {IPAddress: "1.1.1.2", OverallResult: "fail"}, - {IPAddress: "2.2.2.2", OverallResult: "pass"}, - } - - if got := filterQueueItems(items, "", ""); len(got) != 3 { - t.Fatalf("expected no-op with empty q/status, got %+v", got) - } - if got := filterQueueItems(items, "1.1.1", ""); len(got) != 2 { - t.Fatalf("expected 2 matches for q=1.1.1, got %+v", got) - } - if got := filterQueueItems(items, "1.1.1.1", ""); len(got) != 1 || got[0].IPAddress != "1.1.1.1" { - t.Fatalf("expected exact-substring match, got %+v", got) - } - if got := filterQueueItems(items, "1.1.1.1", ""); len(got) != 1 { - t.Fatalf("expected search to be case/substring based, got %+v", got) - } - if got := filterQueueItems(items, "", "pass"); len(got) != 2 { - t.Fatalf("expected 2 matches for status=pass, got %+v", got) - } - if got := filterQueueItems(items, "1.1.1", "pass"); len(got) != 1 || got[0].IPAddress != "1.1.1.1" { - t.Fatalf("expected q+status combined with AND, got %+v", got) - } - if got := filterQueueItems(items, "9.9.9.9", ""); len(got) != 0 { - t.Fatalf("expected no matches, got %+v", got) - } - - // Case-insensitivity, via a query with mixed-case letters (IP octets - // are numeric, so exercise it through IPv6-shaped input instead). - mixed := []ipQueueItem{{IPAddress: "fe80::AbCd"}} - if got := filterQueueItems(mixed, "abcd", ""); len(got) != 1 { - t.Fatalf("expected case-insensitive search to match, got %+v", got) - } -} - -func TestFilterRegistryItems(t *testing.T) { - items := []registryItem{ - {IPAddress: "1.1.1.1", LastResult: "pass"}, - {IPAddress: "1.1.1.2", LastResult: "partial"}, - {IPAddress: "2.2.2.2", LastResult: ""}, - } - - if got := filterRegistryItems(items, "", ""); len(got) != 3 { - t.Fatalf("expected no-op with empty q/status, got %+v", got) - } - if got := filterRegistryItems(items, "1.1.1", ""); len(got) != 2 { - t.Fatalf("expected 2 matches for q=1.1.1, got %+v", got) - } - if got := filterRegistryItems(items, "", "partial"); len(got) != 1 || got[0].IPAddress != "1.1.1.2" { - t.Fatalf("expected exactly the partial-result address, got %+v", got) - } - if got := filterRegistryItems(items, "2.2.2", "partial"); len(got) != 0 { - t.Fatalf("expected q+status combined with AND to exclude non-matching, got %+v", got) - } -} - func TestAutoCyclePanelRenders(t *testing.T) { fake, caURL := newFakeControlAPI(t) finished := time.Now().Add(-time.Hour) diff --git a/internal/dashboard/paging.go b/internal/dashboard/paging.go new file mode 100644 index 0000000..866f8ef --- /dev/null +++ b/internal/dashboard/paging.go @@ -0,0 +1,118 @@ +package dashboard + +import ( + "net/url" + "strconv" + "strings" +) + +// Server-side pagination shared by /ips and /registry: `page` (1-based, +// clamped to the last page) and `per_page` (one of perPageOptions, else +// defaultPerPage). + +const defaultPerPage = 50 + +var perPageOptions = []int{25, 50, 100, 200} + +// parsePerPage accepts only the whitelisted page sizes. +func parsePerPage(raw string) int { + n, err := strconv.Atoi(strings.TrimSpace(raw)) + if err != nil { + return defaultPerPage + } + for _, o := range perPageOptions { + if n == o { + return n + } + } + return defaultPerPage +} + +// parsePage returns the 1-based page number; anything invalid is page 1. +func parsePage(raw string) int { + n, err := strconv.Atoi(strings.TrimSpace(raw)) + if err != nil || n < 1 { + return 1 + } + return n +} + +func pageCount(total, perPage int) int { + if total <= 0 || perPage <= 0 { + return 1 + } + return (total + perPage - 1) / perPage +} + +// clampPage limits page to the last non-empty page for the given total. +func clampPage(page, total, perPage int) int { + if last := pageCount(total, perPage); page > last { + page = last + } + if page < 1 { + page = 1 + } + return page +} + +// pagerData drives the shared "pager" template partial. +type pagerData struct { + // Wrap is the id of the swap-target element holding the table (the + // ‹ › links replace it via hx-select + hx-target). + Wrap string + Page, PerPage, Total, Pages, From, To int + PrevURL, NextURL string +} + +// newPager builds the pager for page (already clamped) of total rows. base is +// the page path; params are the filters to preserve in the links (without +// page/per_page, which the pager sets itself). +func newPager(base, wrap string, params url.Values, page, perPage, total int) pagerData { + p := pagerData{Wrap: wrap, Page: page, PerPage: perPage, Total: total, Pages: pageCount(total, perPage)} + if total > 0 { + p.From = (page-1)*perPage + 1 + p.To = page * perPage + if p.To > total { + p.To = total + } + } + if page > 1 { + p.PrevURL = pageURL(base, params, page-1, perPage) + } + if page < p.Pages { + p.NextURL = pageURL(base, params, page+1, perPage) + } + return p +} + +// pageURL builds base?...&page=N&per_page=M with every value escaped by +// url.Values (so a search string like "a&b" can't smuggle in a parameter). +// Page 1 omits `page`. +func pageURL(base string, params url.Values, page, perPage int) string { + v := url.Values{} + for k, vals := range params { + for _, val := range vals { + if val != "" { + v.Add(k, val) + } + } + } + v.Set("per_page", strconv.Itoa(perPage)) + if page > 1 { + v.Set("page", strconv.Itoa(page)) + } + return base + "?" + v.Encode() +} + +// chunk splits list into slices of at most size elements. +func chunk(list []string, size int) [][]string { + var out [][]string + for len(list) > size { + out = append(out, list[:size]) + list = list[size:] + } + if len(list) > 0 { + out = append(out, list) + } + return out +} diff --git a/internal/dashboard/routes.go b/internal/dashboard/routes.go index d5cff68..c162e15 100644 --- a/internal/dashboard/routes.go +++ b/internal/dashboard/routes.go @@ -18,6 +18,7 @@ func (s *Server) routes(mux *http.ServeMux) { mux.HandleFunc("GET /ips/{ip}", s.handleIPDetail) mux.HandleFunc("POST /ips", s.handleIPsSubmit) mux.HandleFunc("POST /ips/scan", s.handleIPsScan) + mux.HandleFunc("GET /ips/scan/status", s.handleIPsScanStatus) mux.HandleFunc("POST /ips/{ip}/recheck", s.handleIPRecheck) mux.HandleFunc("POST /ips/{ip}/cancel", s.handleIPCancel) mux.HandleFunc("DELETE /ips/{ip}", s.handleIPDelete) diff --git a/internal/dashboard/scale_test.go b/internal/dashboard/scale_test.go new file mode 100644 index 0000000..ec01be3 --- /dev/null +++ b/internal/dashboard/scale_test.go @@ -0,0 +1,685 @@ +package dashboard + +import ( + "fmt" + "html" + "net/http" + "net/url" + "regexp" + "strconv" + "strings" + "testing" + "time" +) + +// Tests for the scale features: background scan panel, server-side +// pagination/filters on /ips and /registry, bulk selection by filter, and the +// bounded overview. + +func ipRowLink(ip string) string { return `href="/ips/` + ip + `"` } + +func countRows(body string) int { return strings.Count(body, `name="addresses" value="`) } + +func hasHXTriggerEvery(body string) bool { return strings.Contains(body, `hx-trigger="every 2s"`) } + +var hrefRe = regexp.MustCompile(`href="([^"]*)"[^>]*rel="(prev|next)"`) + +// pagerLink extracts the (unescaped, parsed) prev/next link of the pager. +func pagerLink(t *testing.T, body, rel string) *url.URL { + t.Helper() + for _, m := range hrefRe.FindAllStringSubmatch(body, -1) { + if m[2] == rel { + u, err := url.Parse(html.UnescapeString(m[1])) + if err != nil { + t.Fatalf("parse pager link %q: %v", m[1], err) + } + return u + } + } + return nil +} + +func TestIPsScan(t *testing.T) { + fake, caURL := newFakeControlAPI(t) + fake.scanFreeAddresses = []string{"5.5.5.5", "5.5.5.6"} + fake.scanRunPolls = 2 + ts := newTestServer(t, caURL) + + // Start: answers at once with the panel; it polls itself, the scan + // buttons are disabled and nothing is queued yet. + resp, body := doReq(t, ts, reqOpts{method: http.MethodPost, path: "/ips/scan"}) + if resp.StatusCode != http.StatusOK { + t.Fatalf("status = %d, want 200", resp.StatusCode) + } + for _, want := range []string{ + `id="scan-progress"`, `hx-get="/ips/scan/status"`, `hx-trigger="every 2s"`, "читаются страницы", + ``, ` disabled`, + } { + if !strings.Contains(body, want) { + t.Fatalf("expected %q in the running panel, got:\n%s", want, body) + } + } + if resp.Header.Get("HX-Trigger") != "" { + t.Fatalf("a running scan must not fire scan-finished, got %q", resp.Header.Get("HX-Trigger")) + } + if len(fake.ips) != 0 { + t.Fatalf("nothing should be queued before the job finishes, got %+v", fake.ips) + } + + // First poll: still running, keeps polling. + resp, body = doReq(t, ts, reqOpts{path: "/ips/scan/status"}) + if !hasHXTriggerEvery(body) || resp.Header.Get("HX-Trigger") != "" { + t.Fatalf("expected a still-polling panel without HX-Trigger, got %q:\n%s", resp.Header.Get("HX-Trigger"), body) + } + + // Second poll: finished — no polling trigger, HX-Trigger tells the table to reload. + resp, body = doReq(t, ts, reqOpts{path: "/ips/scan/status"}) + if hasHXTriggerEvery(body) || strings.Contains(body, `hx-get="/ips/scan/status"`) { + t.Fatalf("a finished panel must stop polling, got:\n%s", body) + } + if got := resp.Header.Get("HX-Trigger"); got != "scan-finished" { + t.Fatalf("HX-Trigger = %q, want scan-finished", got) + } + for _, want := range []string{"готово", "
добавлено
2
", `прочитано страниц
33
"} { + if !strings.Contains(page, want) { + t.Fatalf("expected %q in the page, got:\n%s", want, page) + } + } + if strings.Index(page, `id="scan-progress"`) > strings.Index(page, `