diff --git a/README.md b/README.md index bd4ccc1..4f50f93 100644 --- a/README.md +++ b/README.md @@ -200,6 +200,7 @@ scripts/run-local-e2e.sh # сквозной прог | Дата | Веха | Документ | |---|---|---| +| 2026-10-02 | Валидатор не получает второй адрес при потере heartbeat (иначе адреса уходили в `fail` без проверок); heartbeat агента в отдельном потоке; быстрая очистка очереди, не зависящая от соединения клиента | [план](docs/changes/2026-10-02_09-02_orchestrator-validator-state-and-clear-plan.md) · [анализ инцидента](analysis/2026-10-02_08-56_1026-addresses_mass-check-analysis.md) · [USAGE](docs/USAGE.md#управление-валидаторами) | | 2026-10-02 | Самопроверка через ручку control-api (`self_check.methods`, способ `control_api` рядом с IP-echo); отвязка Floating IP при провале self-check | [план](docs/changes/2026-10-02_03-06_self-check-control-api-plan.md) · [API](docs/API.md#get-apiv1agentsidobserved-ip) | | 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#аутентификация) | diff --git a/analysis/2026-10-02_08-56_1026-addresses_mass-check-analysis.md b/analysis/2026-10-02_08-56_1026-addresses_mass-check-analysis.md new file mode 100644 index 0000000..9c8bda9 --- /dev/null +++ b/analysis/2026-10-02_08-56_1026-addresses_mass-check-analysis.md @@ -0,0 +1,167 @@ +# Аналитический разбор массовой проверки: 1026 адресов + +> Время отчёта: 2026-10-02 08:56 UTC · Данные: снимок БД control-api на 08:46:53 UTC и лог control-api за 4 часа до остановки +> Проверено адресов: **1026** (984 завершены `done` + 42 завершены `failed`) из 6440 в очереди +> Окно прогона: 07:14:58 – 08:46:53 UTC (1 ч 32 мин). Остановлен вручную в 08:53:46 UTC. + +## 1. Остановка проверок + +- Проверки остановлены в 08:53:46 UTC операцией «Очистить всё» (`POST /api/v1/admin/ips/clear`). +- После остановки: очередь пуста, все 20 валидаторов `idle`, последняя выдача адреса в 08:53:43, новых выдач нет. + Реестр с историей адресов сохранён (6445 записей). +- Отвязка Floating IP с портов валидаторов при очистке: 13 зависших привязок снято прямым опросом портов + (`detached floating ip from validator port`), ещё 8 привязок, завершавшихся во время очистки, система сняла сама + (`address removed during association`). Состояние портов в самом OpenStack на момент отчёта не проверялось. +- **Первая попытка очистки не сработала.** Очистка шла дольше таймаута клиента (120 с) и оборвалась: 600 отвязок завершились + с ошибкой `context canceled`, ничего не удалилось, проверки продолжались. Повторная очистка без таймаута выполнялась 256 с. +- **Причина долгой очистки (дефект кода, не исправлен):** у завершённых адресов (`done`) в БД остаётся `fip_id`, и очистка + последовательно отвязывает все такие FIP, хотя они уже свободны (более 1000 вызовов OpenStack по ~0,2 с). + +## 2. Итоги прогона + +| Показатель | Значение | +|---|---| +| Адресов в очереди | 6440 | +| Завершено | 1026 (984 `done` + 42 `failed`) | +| `pass` | 227 (23% от завершённых) | +| `partial` | 757 (77%) | +| `fail` | 42 (все `failed`, вердикта по существу нет, см. раздел 4) | +| Остались в очереди / в работе на момент снимка | 5392 в очереди, 22 в работе | + +Пропускная способность по 10-минутным окнам (завершено адресов за минуту): 07:10 — 5,8; 07:20 — 13,4; 07:30 — 11,6; +07:40 — 10,1; 07:50 — 10,6; 08:00 — 9,4; 08:10 — 9,7; 08:20 — 10,5; 08:30 — 10,6; 08:40 — 6,7 (окно неполное). + +Цикл одного адреса (от выдачи валидатору до итога): минимум 65 с, медиана 85 с, p90 100 с, p99 115 с, максимум 125 с, +среднее 84 с (по 984 адресам `done`). Для 20 валидаторов это теоретически ~14 адресов в минуту; фактически ~10,4, +потеря около 27% из-за залипших валидаторов (раздел 4). + +## 3. Фактура по накопившимся ошибкам + +### 3.1. Классы ошибок (с 07:10) + +В логе control-api за 4 часа на уровне `ERROR` только один вид сообщений: 41 ошибка привязки Floating IP. + +| Ошибка / событие | Число | Что это | +|---|---|---| +| `lease expired` (возврат адреса по истечении лизинга) | 198 | Валидатор не подхватил задание за время лизинга. Затронуто минимум 77 адресов (часть событий без привязки к адресу) | +| Привязка FIP: `409 Cannot associate floating IP … fixed IP already has a floating IP` | 41 | На порту валидатора уже висит другой FIP. Порты: v1 — 13, v7 — 6, v13 — 6, v3 — 5, v12 — 5, v16 — 4, v17 — 2 | +| `validator_unreachable` | 7 | По одному разу: v1 и v12 (07:18:03), v13 (07:33:53), v3 (07:34:33), v7 (07:40:03), v16 (07:42:33), v17 (08:17:53) | +| `site_unreachable` | 5 | rxmsk (08:09, 08:45) и misha-v (08:09, 08:20, 08:45) | +| «Очистить всё»: `context canceled` | 600 | Последствие обрыва первой очистки (08:49), см. раздел 1 | + +Чего не было: **провалов self-check — 0 из 1015** результатов; все 1015 прошли способом `control_api` (запасной `ip_echo` +не понадобился). Событий `fip_occupied` — 0, адресов `occupied` — 0. + +Журнал событий за прогон: `self_check_result` 2030 (по две записи на адрес), `fip_associated` 1120, `config_received` 1015, +`aggregated` 1004, `retry_or_fail` 239 (198 лизинг + 41 привязка), `lease_expired` 198, `validator_unreachable` 7, +`site_unreachable` 5. + +### 3.2. По валидаторам + +| Валидатор | Завершено (`done`) | `failed` | Сбросов лизинга | `unreachable` | +|---|---|---|---|---| +| vkiplab-v1 | 2 | 17 | 49 | 1 | +| vkiplab-v12 | 2 | 13 | 51 | 1 | +| vkiplab-v13 | 13 | 10 | 41 | 1 | +| vkiplab-v16 | 38 | 2 | 19 | 1 | +| vkiplab-v7 | 36 | 0 | 21 | 1 | +| vkiplab-v17 | 50 | 0 | 10 | 1 | +| vkiplab-v3 | 50 | 0 | 7 | 1 | +| остальные 13 (v2, v4–v6, v8–v11, v14, v15, v18–v20) | 60–63 | 0 | 0 | 0 | + +(Сбросы лизинга и `unreachable` в таблице — только за прогон, с 07:10. У v14, v18 и v2 в истории БД есть сбросы лизинга +за 1 октября, к этому прогону они не относятся.) + +- **v1 и v12** не подхватили ни одного задания после 07:17 (последний подхват 07:17:29 и 07:17:26). +- **v13** — после 07:33:16. +- **v7, v16, v17, v3** залипали временно и затем восстановились; механизм восстановления не выяснен. + +### 3.3. 42 адреса `fail` + +- По валидаторам: v1 — 17, v12 — 13, v13 — 10, v16 — 2. Все 42 исчерпали повторы: `retry_count` = 4 и `attempt_number` = 4 у каждого. +- Причины неудачных попыток по этим адресам: 143 сброса лизинга и 25 ошибок привязки `409`. +- **Ни одной проверки по ним не выполнено** (в таблице проверок у этих адресов 0 записей): адреса не получили вердикта, + и `fail` здесь не характеризует сами адреса. +- Появлялись равномерно с 07:31 до 08:46 (4–10 за 10 минут). +- Подсети: 37.139.x, 79.137.x, 83.166.x и другие, без концентрации. + +## 4. Причины + +Подтверждены по БД и логам control-api. Логи агентов на самих валидаторах не изучались. + +**Причина 1. Валидатор получает два адреса сразу и застревает.** Два дефекта вместе: +- Heartbeat (`queries_validators.go`) возвращает валидатор из `unreachable` в `idle`, не проверяя, что за ним числится адрес. + Адрес с долгими внешними проверками (3–4 таймаута по 10 с) блокирует агента больше 30 с (порог + `heartbeat_timeout_seconds`), control-api помечает валидатор недоступным, затем возвращает в `idle` занятым. +- При завершении старого адреса `ReleaseFIP` (`queries_ipqueue.go`) освобождает валидатор по его имени, а не по адресу. +- Подтверждение: **35 двойных выдач** (два адреса одному валидатору с интервалом ~5 с) в логе: v1 — 10, v12 — 9, v7 — 6, + v13 — 4, v17 — 3, v16 — 2, v3 — 1. Из 1265 выдач в логе. + +**Причина 2. Регрессия моей правки с параллельной привязкой.** Защита от дублей в `orchestrator.go:149` ключуется по +валидатору (`assign:`). Вторая выдача того же валидатора пропускает привязку, и адрес стоит в `assigning_fip` +до истечения лизинга. В БД у валидатора одно поле `current_ip_id`, оно указывает на последний выданный адрес, агент по нему +получает пустое задание, лизинг истекает, валидатор берёт новый адрес, и круг повторяется. Ошибки `409` — следствие: на +порту остаётся FIP первого адреса, привязать второй нельзя. + +**Причина 3. Агент молчит во время долгих проверок.** Heartbeat отправляется только между заданиями; адреса с несколькими +таймаутами ведут к `unreachable` (пусковой механизм причины 1). Пример: v1 проверял `37.139.32.1`, v12 — `37.139.32.4`; +у обоих проваливались все четыре внешних HTTPS-цели (~40 с таймаутов), в 07:18:03 оба помечены `unreachable`. + +**Дефект очистки.** См. раздел 1: отвязка всех `done`-адресов последовательно. + +## 5. Результаты проверок самих адресов + +**Почему 77% `partial`:** 735 адресов проваливают только исходящие HTTPS, ещё 22 — исходящие и входящие. + +| Цель (egress HTTPS) | Адресов с провалом (из 984) | +|---|---| +| `packages.ubuntu.com` | 698 (71%) | +| `repo.almalinux.org/almalinux/` | 304 (31%) | +| `github.com` | 251 (26%) | +| `hub.docker.com` | 17 (2%) | + +Число провалов на один `partial`-адрес: 1 — 392 адреса, 2 — 202, 3 — 136, 4 и больше — 27. + +По подсетям (доля адресов с провалом цели): + +| Подсеть | Адресов | `packages.ubuntu.com` | `repo.almalinux.org` | `github.com` | `hub.docker.com` | Доля `partial` | +|---|---|---|---|---|---|---| +| 83.166.x | 435 | 80% | 53% | 42% | 0% | 92% | +| 37.139.x | 418 | 62% | 11% | 10% | 4% | 63% | +| 79.137.x | 109 | 66% | 19% | 18% | 0% | 68% | +| 5.188.x | 22 | 59% | 0% | 0% | 0% | 59% | + +- `packages.ubuntu.com` проваливается у всех подсетей и во все окна. Доля проваленных строк egress для этой цели по 10-минутным + окнам выросла с 26% до 40% (строки включают HTTPS и ICMP, поэтому реальная доля HTTPS вдвое выше). +- Доля `partial` почти одинакова у всех валидаторов (72–84%; v1 — 100% по 2 адресам, v12 — 50% по 2): причина в самом адресе + или во внешнем ресурсе, а не в валидаторе. +- Провалы `repo.almalinux.org` и `github.com` сильно зависят от подсети (83.166.x — 53% и 42%, 37.139.x — 11% и 10%). +- Успешные HTTPS-проверки: медиана 200 мс, p90 6707 мс (много ответов близко к таймауту 10 с). + +**Входящие проверки.** Провалы только у `inbound-site-1` (33 пробы из 2934) и `inbound-site-3` (50 из 2931); `inbound-site-2` — 0 из 2952, +`inbound-site-4` — 1 из 2952. Провалы всплесками в 08:00, 08:20 и 08:40 по ICMP, SSH и TCP 22; у 20 адресов упали все пробы +одной площадки. Пробер rxmsk и misha-v в эти же минуты отмечены `unreachable`. Причина недоступности проберов не выяснена. + +## 6. Рекомендации + +1. **Исправить оркестратор** (до повторного запуска): защита от дублей по адресу, `unreachable` → `assigned` вместо `idle`, + освобождение валидатора только по текущему адресу, heartbeat агента в отдельном потоке. +2. **Исправить очистку:** отвязывать только действительно привязанные FIP (без `fip_released_at`) и параллельно. +3. **Перепроверить 42 адреса** после исправления: вердикта по существу они не получили. Все `fail` с причиной «lease expired» считать недействительными. +4. **Решить по `packages.ubuntu.com`:** 71% `partial` и 10 с таймаута на каждый такой адрес; если ресурс нестабилен, убрать или заменить. +5. **Проверить причины `site_unreachable`** у проберов (возможна перегрузка при большой очереди). +6. Проверить в OpenStack, что порты валидаторов свободны от Floating IP. + +## 7. Что не проверено + +- Состояние портов и Floating IP в самом OpenStack. +- Логи агентов `validator-agent` на валидаторах и логи проберов (причины блокировок heartbeat подтверждены по времени и + событиям control-api, не по логам агентов). +- Почему залипшие v7, v16, v17, v3 восстановились. +- Причины провалов внешних HTTPS-целей (ресурс, сеть облака или фильтрация) и недоступности проберов. + +## 8. Источники данных + +- Снимок БД control-api: `.backup` на 08:46:53 UTC (таблицы `ip_queue`, `checks`, `events`, `validators`, `sites`). +- Лог контейнера control-api за 4 часа до остановки (1504 строки) и за период очистки. +- Статус и результат очистки: HTTP 200 за 256,7 с. diff --git a/analysis/2026-10-02_08-56_failed-addresses.txt b/analysis/2026-10-02_08-56_failed-addresses.txt new file mode 100644 index 0000000..285320f --- /dev/null +++ b/analysis/2026-10-02_08-56_failed-addresses.txt @@ -0,0 +1,42 @@ +37.139.32.70 +37.139.32.71 +37.139.32.72 +37.139.32.151 +37.139.33.219 +37.139.33.227 +37.139.33.230 +37.139.34.75 +37.139.34.94 +37.139.34.128 +37.139.40.122 +37.139.41.103 +37.139.41.129 +37.139.41.157 +37.139.41.202 +37.139.42.126 +37.139.43.9 +79.137.174.150 +79.137.174.172 +79.137.174.205 +79.137.174.210 +79.137.175.14 +79.137.175.51 +79.137.175.71 +79.137.175.162 +83.166.232.116 +83.166.232.117 +83.166.232.159 +83.166.233.37 +83.166.233.75 +83.166.233.78 +83.166.234.136 +83.166.235.75 +83.166.235.132 +83.166.235.247 +83.166.237.59 +83.166.237.228 +83.166.238.19 +83.166.248.57 +83.166.248.60 +83.166.248.92 +83.166.248.188 diff --git a/bin/SHA256SUMS b/bin/SHA256SUMS index c598d1e..6e142ba 100644 --- a/bin/SHA256SUMS +++ b/bin/SHA256SUMS @@ -1,4 +1,4 @@ -6fc8aef1671a6e17c60461431b90e86ababc3f154042dff205a3ad9dd11d83c3 control-api -b0e75dce3a840024af40d9544ab17524fb85431e7944582893fe3f9aae1b3401 validator-agent +9447699d9eed4b9fb5d417769fc868986a10357ec633ec3742883b3002aaafc3 control-api +9fb6608b84143f7c4f318f3cc92dcd9f95c7831d627b23d67cce5a5908ced704 validator-agent 3e9e14dbb361ee76aaad7c1da6864b3ea111e0ed151403f904b12485631bbf75 prober d72234688eb1954dbff420ca1c8e83b83ceab561b18a0336af0f73fe4d9ac8be admin-dashboard diff --git a/bin/control-api b/bin/control-api index ea1b7e6..4425ef8 100755 Binary files a/bin/control-api and b/bin/control-api differ diff --git a/bin/validator-agent b/bin/validator-agent index b02bdfe..5bed3d9 100755 Binary files a/bin/validator-agent and b/bin/validator-agent differ diff --git a/docs/USAGE.md b/docs/USAGE.md index 8b61623..1c70f4b 100644 --- a/docs/USAGE.md +++ b/docs/USAGE.md @@ -408,6 +408,17 @@ curl -s http://:8080/api/v1/admin/validators | python3 -m json.tool (занят), `unreachable` (пропустил heartbeat дольше `orchestrator.heartbeat_timeout_seconds`). +Правила, которые держат состояние валидатора согласованным: +- валидатор держит **не более одного адреса**; адрес освобождает валидатор + только пока он остаётся его текущим (запоздалое завершение старого адреса + чужого валидатора не освобождает); +- после пропущенного heartbeat валидатор, у которого есть адрес, возвращается + в `assigned`, а не в `idle`, и не получает второй адрес; без адреса — в `idle`; +- `unreachable`-валидатор не получает адресов, пока не пришлёт heartbeat + (даже если лизинг его адреса истёк и адрес вернулся в очередь); +- на каждом такте оркестратор сверяет валидаторы с очередью и исправляет + расхождения (в логе `repaired validators that disagreed with the queue`). + **Добавление нового валидатора (без перезапуска control-api):** 1. Поднимите новую ВМ в сервисном проекте облака, узнайте её Neutron `port_id`. @@ -690,6 +701,12 @@ curl -s -X POST http://:8080/api/v1/admin/ips/delete \ curl -s -X POST http://:8080/api/v1/admin/ips/clear ``` +Очистка отвязывает Floating IP только у адресов, которые ещё в работе +(завершённые уже свободны), затем сама опрашивает порты валидаторов в облаке и +снимает оставшиеся привязки адресов из реестра. Занимает секунды. Она не +прерывается разрывом соединения (таймаутом клиента или дашборда): операция +доводится до конца на стороне control-api, предел — 10 минут. + В `admin-dashboard` то же самое доступно на странице `/ips`: чекбоксы у каждой строки + кнопка «Удалить выбранные» для точечного/массового удаления, кнопка «Удалить» в каждой строке, и отдельная кнопка «Очистить diff --git a/docs/changes/2026-10-02_09-02_orchestrator-validator-state-and-clear-plan.md b/docs/changes/2026-10-02_09-02_orchestrator-validator-state-and-clear-plan.md new file mode 100644 index 0000000..c1f8e21 --- /dev/null +++ b/docs/changes/2026-10-02_09-02_orchestrator-validator-state-and-clear-plan.md @@ -0,0 +1,146 @@ +# План: исправление оркестратора (двойная выдача валидатору) и очистки очереди + +> Дата: 2026-10-02 09:02 UTC · Статус: **реализовано и проверено** (юнит-тесты с `-race`, локальный e2e; выкладка и приёмка на стенде — ниже). Решения пользователя: предел очистки 10 минут и 8 параллельных отвязок приняты как базовые +> Основание: [analysis/2026-10-02_08-56_1026-addresses_mass-check-analysis.md](../../analysis/2026-10-02_08-56_1026-addresses_mass-check-analysis.md) + +## Context + +Массовая проверка 2 октября остановилась на 1026 из 6440 адресов. 7 из 20 валидаторов «залипли»: v1, v12, v13 не взяли +ни одного задания после 07:17 и 07:33; v7, v16, v17, v3 залипали временно. Итог: 198 сбросов лизинга, 41 ошибка привязки +Floating IP (`409 fixed IP already has a floating IP`), **42 адреса в `fail` без единой выполненной проверки**, потеря +~27% пропускной способности. Отдельно «Очистить всё» не уложилась в таймаут клиента (256 с вместо секунд) и сначала +оборвалась с 600 ошибками `context canceled`. + +Цель: валидатор в любой момент держит не больше одного адреса; потеря heartbeat не приводит к двойной выдаче; +«Очистить всё» выполняется за секунды и не прерывается разрывом соединения. + +## Причины (по коду, подтверждены логами и БД) + +| № | Причина | Где | +|---|---|---| +| 1 | Heartbeat возвращает `unreachable` → `idle`, не глядя на `current_ip_id`: занятый валидатор снова считается свободным | `internal/db/queries_validators.go:41` (`Heartbeat`), `:31` (`RegisterValidator`, возврат из `unreachable`) | +| 2 | Освобождение валидатора идёт **по имени**, а не по адресу, который он держит: завершение старого адреса освобождает валидатор, уже взявший новый. Так же `RequeueOrFail` (в т. ч. при сбросе лизинга) и `MarkFIPOccupied`, `FreeValidator`. Освобождённый ставится в `idle` даже если он `unreachable` — мёртвый валидатор получает новые адреса каждые 3 минуты | `internal/db/queries_ipqueue.go` (`ReleaseFIP`, `RequeueOrFail`, `MarkFIPOccupied`), `queries_validators.go:143` (`FreeValidator`), `internal/orchestrator/orchestrator.go:469` | +| 3 | Моя регрессия: защита от дублей привязки ключуется по валидатору (`assign:`). Вторая выдача того же валидатора не запускает привязку и стоит в `assigning_fip` до конца лизинга | `internal/orchestrator/orchestrator.go:149` | +| 4 | Агент шлёт heartbeat только между заданиями. Адрес с 3–4 таймаутами внешних проверок (~40 с) блокирует его дольше порога 30 с | `internal/agentcore/agentcore.go` (`Run`, `pollOnce`) | +| 5 | «Очистить всё» отвязывает FIP у **всех** строк с непустым `fip_id`, а он остаётся у `done`/`failed`. Больше 1000 последовательных вызовов OpenStack | `internal/db/queries_ipqueue.go` (`ListFIPRefs`, `ListFIPRefsByAddresses`), `internal/orchestrator/orchestrator.go` (`ClearQueue`, `DeleteIPs`) | +| 6 | Очистка работает на контексте HTTP-запроса: разрыв соединения клиентом обрывает её посреди дела (отвязано часть, БД не очищена) | `internal/httpapi/handlers_admin.go:322` | + +## Инварианты, которые вводим + +- **I1.** Валидатор держит не более одного адреса: `validators.current_ip_id = X` тогда и только тогда, когда у строки `X` + `owner_validator_id` равен этому валидатору и состояние не терминальное (`done`, `failed`, `occupied`). +- **I2.** `idle` означает `current_ip_id IS NULL`. Состояния `unreachable` и `unregistered` не затираются освобождением. +- **I3.** Валидатор освобождает только тот адрес, который он сейчас держит. Освобождение и возврат по лизингу чужого или + устаревшего адреса состояние валидатора не меняют. + +## Изменения + +### 1. База данных (без миграций, только запросы) + +`internal/db/queries_validators.go`, `queries_ipqueue.go`: + +- **`Heartbeat`:** `unreachable` → `assigned`, если `current_ip_id IS NOT NULL`, иначе `idle`. + То же в `RegisterValidator` (возврат из `unreachable`/`unregistered`). +- **Общая функция освобождения** `freeValidatorTx(tx, validatorID, ipID)`: + `UPDATE validators SET current_ip_id=NULL, state = CASE WHEN state='unreachable' THEN state ELSE 'idle' END WHERE validator_id=? AND current_ip_id=?`. + Используют `ReleaseFIP`, `RequeueOrFail`, `MarkFIPOccupied`, `FreeValidator` (получает второй аргумент — id адреса). + `deleteIPTx` (уже по `current_ip_id`) и `ClearAllIPs` переводятся на тот же `CASE`, чтобы не затирать `unreachable`. +- **`ClaimNextQueued`:** условие обновления валидатора дополняется `AND current_ip_id IS NULL`. +- **`ReconcileValidators`** (новая, вызывается из `Orchestrator.Tick`): лечит нарушение инвариантов, если они всё же возникли + (падение процесса, старые строки): валидатор с `current_ip_id`, чья строка не существует, терминальна или принадлежит другому + валидатору, освобождается; строка в `assigning_fip`/`awaiting_self_check`/`checking`, чей владелец не ссылается на неё, + не трогается (её вернёт сброс лизинга). Один короткий запрос на такт. +- **`ListFIPRefs` и `ListFIPRefsByAddresses`:** только строки в нетерминальных состояниях (`state NOT IN done, failed, occupied`) + с непустым `fip_id`: у терминальных FIP уже отвязан (агрегация, возврат по лизингу, отмена и `occupied` делают это до записи + состояния). Значения `fip_id` в строках не меняются (дашборд их показывает). + +### 2. Оркестратор (`internal/orchestrator/orchestrator.go`) + +- Ключ защиты от дублей привязки: `assign:`, а не `assign:` (строка 149). Дублирующий запуск привязки + одного и того же адреса по-прежнему исключён. +- `ForceCancel`, `sweepExpiredLeases`, `aggregateAndRelease`, `SelfCheckResult`: передают id адреса в освобождение (I3). +- **`ClearQueue`/`DeleteIPs`:** отвязка FIP из `ListFIPRefs` (теперь ≤ числа валидаторов) выполняется параллельно, не более + 8 одновременных вызовов; затем прямой опрос портов (`releaseValidatorPorts`) как сейчас. +- `Tick`: вызывает `ReconcileValidators` первым шагом. + +### 3. HTTP (`internal/httpapi/handlers_admin.go`) + +- «Очистить всё», удаление списка и отмена: контекст отвязан от отмены запроса (`context.WithoutCancel`) с собственным пределом + времени (10 минут). Разрыв соединения клиентом (в т. ч. таймаут дашборда) больше не обрывает операцию на середине. + +### 4. Агент (`internal/agentcore/agentcore.go`) + +- Heartbeat уходит из `pollOnce` в **отдельную горутину** со своим тикером (период `poll_interval_seconds`); останавливается по + отмене контекста. Долгие внешние проверки больше не блокируют heartbeat. Ошибки heartbeat пишутся в лог (предупреждение). +- Протокол и конфигурация агента не меняются (старый агент с новым сервером и наоборот работают). + +### 5. Документация + +`docs/USAGE.md`/`docs/DIAGRAMS.md` (состояния валидатора и правила освобождения), запись в истории изменений `README.md`. + +## Тесты + +Пишу сам (агенты тесты не делают); каждый тест должен падать без исправления. + +- **db:** + - `Heartbeat` из `unreachable` с адресом даёт `assigned`, без адреса `idle`; `RegisterValidator` аналогично; + - `ReleaseFIP`/`RequeueOrFail`/`MarkFIPOccupied` старого адреса не освобождают валидатор, который уже держит другой адрес; + - освобождение `unreachable`-валидатора оставляет `unreachable`; + - `ClaimNextQueued` не выдаёт адрес валидатору с `current_ip_id`; + - `ReconcileValidators` лечит битые строки и не трогает корректные; + - `ListFIPRefs`/`ListFIPRefsByAddresses` не содержат `done`/`failed`/`occupied`. +- **orchestrator:** + - сценарий инцидента: валидатор помечен `unreachable`, держа адрес A с идущими проверками → heartbeat → такт оркестратора + **не выдаёт** ему B; после завершения A валидатор получает B; + - мёртвый валидатор (`unreachable`, лизинг истёк) не получает новых адресов; + - `ClearQueue` при 1000 строк `done` с `fip_id` и 5 активных не вызывает `Disassociate` для `done` (счётчик вызовов в обёртке OpenStack); + - защита привязки по адресу (два адреса одного валидатора, искусственно, оба привязываются); + - **случайный сценарий под нагрузкой** (20 валидаторов, 300 адресов, случайная потеря heartbeat, долгие проверки, часть адресов + с отказом привязки): после каждого такта проверяются I1–I3; в конце все адреса завершены, сбросов лизинга нет. +- **httpapi:** «Очистить всё» доживает до конца при отмене контекста запроса посередине (БД очищена, FIP отвязаны). +- **agentcore:** во время долгой проверки (цель отвечает 3 с) heartbeat уходит по расписанию; остановка по отмене контекста. +- **Общий прогон:** `go vet`, `go test -race ./...`, `scripts/run-local-e2e.sh` (с `ip_echo` и `control_api`). + +## Выкладка + +1. Правки control-api: сборка `bin/control-api`, образ `civ-capi`, перезапуск (БД не затрагивается; миграций нет). + Одного этого достаточно, чтобы остановить двойные выдачи и бесконечное «залипание». +2. Агент: сборка `bin/validator-agent`, коммит, раскатка Ansible-сценарием `deploy/ansible` (запускает пользователь): heartbeat в + отдельном потоке убирает ложные `unreachable`. +3. Перепроверка 42 адресов: список выгружается из снимка анализа в файл `analysis/2026-10-02_08-56_failed-addresses.txt` (уже выгружен, 42 адреса); ставятся в очередь + через `POST /api/v1/admin/ips` (или «Перепроверка» в дашборде) после выкладки. + +## Приёмка на стенде + +Контрольная группа из 20 адресов, затем 400 адресов (при тех же внешних целях с долгими таймаутами) с наблюдением 30 минут: + +| Показатель | Критерий | +|---|---| +| `lease_expired` | 0 | +| Ошибки привязки `409` | 0 | +| Двойные выдачи (два `claimed ip` одному валидатору за <10 с) | 0 | +| Завершено каждым валидатором | отклонение от среднего не больше 15% | +| `unreachable` при долгих проверках | нет (после раскатки нового агента) | +| «Очистить всё» при >1000 строк `done` | ответ не дольше 10 с | +| Адреса `fail` | только по существу (не из-за лизинга) | + +## Риски и откат + +- Изменения только в запросах и логике, без миграций; схема БД и протокол агента не меняются. Откат — предыдущий образ `civ-capi` + (`docker tag`/предыдущий коммит) и прежний бинарник агента. +- Риск: условные `UPDATE` по `current_ip_id` могут оставить валидатор занятым, если строка адреса пропала. Страхует `ReconcileValidators`. +- Валидатор `unreachable` теперь не получает адресов до первого heartbeat — это намеренно; если агент жив, но heartbeat по сети + не проходит, он простаивает (раньше брал адреса и терял их). + +## Не входит в эту правку + +- Параллельное выполнение внешних проверок в агенте (3–4 таймаута сейчас идут подряд): сократило бы цикл с ~84 с и снизило бы + нагрузку на heartbeat; отдельное решение. +- Судьба цели `packages.ubuntu.com` (71% `partial`) и причины недоступности проберов — отдельные вопросы из анализа. +- Сокращение `fip_settle_seconds` (30 с) ради темпа. + +## Вопросы к согласованию + +1. Предел времени «Очистить всё» (10 минут) и число параллельных отвязок (8) — подходят? +2. Выкладывать control-api сразу после тестов (до раскатки агента) — да, как в разделе «Выкладка»? +3. 42 адреса `fail` перепроверять сразу после выкладки или вместе со следующим большим прогоном? diff --git a/internal/agentcore/agentcore.go b/internal/agentcore/agentcore.go index 35618e8..8ed9f27 100644 --- a/internal/agentcore/agentcore.go +++ b/internal/agentcore/agentcore.go @@ -19,6 +19,7 @@ import ( neturl "net/url" "os" "strings" + "sync/atomic" "time" "cloudipvalidator/internal/apiclient" @@ -33,6 +34,10 @@ type Agent struct { lastHandledIPID int64 + // busy is true while an assignment is being worked on; it is reported in + // the heartbeat body (informational on the control-api side). + busy atomic.Bool + // registerRetryInitial/Max govern the backoff used while waiting for a // successful registration (see registerWithRetry): control-api may not // be up yet at agent boot, or may come and go across a redeploy, and the @@ -70,6 +75,15 @@ func (a *Agent) Run(ctx context.Context) error { } interval := time.Duration(a.cfg.PollIntervalSeconds) * time.Second + + // Heartbeats run on their own schedule. Sent from the poll loop they + // stopped for as long as a slow assignment took (an address whose + // outbound targets all time out keeps the loop busy for ~40 s), which + // control-api reads as a lost validator after heartbeat_timeout_seconds. + hbCtx, stopHeartbeat := context.WithCancel(ctx) + defer stopHeartbeat() + go a.heartbeatLoop(hbCtx, interval) + ticker := time.NewTicker(interval) defer ticker.Stop() @@ -148,12 +162,32 @@ type checkConfigDTO struct { Targets []string `json:"targets"` } -func (a *Agent) pollOnce(ctx context.Context) { - if _, err := a.client.Do(ctx, "POST", "/api/v1/agents/"+a.cfg.ValidatorID+"/heartbeat", heartbeatReq{LocalState: "idle"}, nil); err != nil { - a.log.Error("heartbeat", "err", err) - return +// heartbeatLoop sends a heartbeat now and then every interval until ctx is +// cancelled. +func (a *Agent) heartbeatLoop(ctx context.Context, interval time.Duration) { + ticker := time.NewTicker(interval) + defer ticker.Stop() + for { + a.sendHeartbeat(ctx) + select { + case <-ctx.Done(): + return + case <-ticker.C: + } } +} +func (a *Agent) sendHeartbeat(ctx context.Context) { + state := "idle" + if a.busy.Load() { + state = "checking" + } + if _, err := a.client.Do(ctx, "POST", "/api/v1/agents/"+a.cfg.ValidatorID+"/heartbeat", heartbeatReq{LocalState: state}, nil); err != nil && ctx.Err() == nil { + a.log.Error("heartbeat", "err", err) + } +} + +func (a *Agent) pollOnce(ctx context.Context) { var assignment assignmentResp ok, err := a.client.Do(ctx, "GET", "/api/v1/agents/"+a.cfg.ValidatorID+"/assignment", nil, &assignment) if err != nil { @@ -169,6 +203,8 @@ func (a *Agent) pollOnce(ctx context.Context) { return // already handled this IP's work this attempt } + a.busy.Store(true) + defer a.busy.Store(false) switch assignment.Phase { case "awaiting_self_check": a.handleSelfCheckAndRun(ctx, assignment) diff --git a/internal/agentcore/heartbeat_test.go b/internal/agentcore/heartbeat_test.go new file mode 100644 index 0000000..ae2ac80 --- /dev/null +++ b/internal/agentcore/heartbeat_test.go @@ -0,0 +1,72 @@ +package agentcore + +import ( + "context" + "fmt" + "net/http" + "net/http/httptest" + "sync/atomic" + "testing" + "time" + + "cloudipvalidator/internal/config" +) + +// While the agent is busy with a slow assignment (an address whose outbound +// targets time out keeps it occupied for tens of seconds) it must keep sending +// heartbeats; control-api marks a validator that stays silent for +// heartbeat_timeout_seconds as unreachable. +func TestHeartbeatContinuesDuringSlowChecks(t *testing.T) { + slowTarget := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + time.Sleep(3500 * time.Millisecond) + })) + defer slowTarget.Close() + + var heartbeats, assignments int32 + mux := http.NewServeMux() + mux.HandleFunc("POST /api/v1/agents/register", func(w http.ResponseWriter, r *http.Request) { + fmt.Fprint(w, `{"ok":true}`) + }) + mux.HandleFunc("POST /api/v1/agents/val-1/heartbeat", func(w http.ResponseWriter, r *http.Request) { + atomic.AddInt32(&heartbeats, 1) + fmt.Fprint(w, `{"ok":true}`) + }) + mux.HandleFunc("GET /api/v1/agents/val-1/assignment", func(w http.ResponseWriter, r *http.Request) { + if atomic.AddInt32(&assignments, 1) > 1 { + w.WriteHeader(http.StatusNoContent) + return + } + fmt.Fprintf(w, `{"ip_id":1,"ip_address":"1.1.1.1","phase":"checking","check_config":[{"type":"https","targets":[%q]}]}`, slowTarget.URL) + }) + mux.HandleFunc("POST /api/v1/agents/val-1/results", func(w http.ResponseWriter, r *http.Request) { fmt.Fprint(w, `{"ok":true}`) }) + mux.HandleFunc("POST /api/v1/agents/val-1/complete", func(w http.ResponseWriter, r *http.Request) { fmt.Fprint(w, `{"ok":true}`) }) + capi := httptest.NewServer(mux) + defer capi.Close() + + a := New(&config.ValidatorAgent{ + ValidatorID: "val-1", ControlAPIURL: capi.URL, PollIntervalSeconds: 1, + Checks: config.AgentChecks{HTTPSTimeoutSeconds: 10, ICMPTimeoutSeconds: 1, ICMPCount: 1}, + }, testLogger()) + + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan struct{}) + go func() { _ = a.Run(ctx); close(done) }() + + time.Sleep(3 * time.Second) // the slow check (3.5 s) is still running + during := atomic.LoadInt32(&heartbeats) + cancel() + select { + case <-done: + case <-time.After(10 * time.Second): + t.Fatal("Run did not stop after the context was cancelled") + } + + // One per second plus the first: 3-4 in 3 s. With heartbeats in the poll + // loop there is exactly one, sent before the slow assignment started. + if during < 3 { + t.Fatalf("%d heartbeats in 3 s while a check was running, want at least 3", during) + } + if atomic.LoadInt32(&assignments) < 1 { + t.Fatal("the assignment was never fetched, the test did not exercise a busy agent") + } +} diff --git a/internal/db/queries_ipqueue.go b/internal/db/queries_ipqueue.go index abd0588..9d179c9 100644 --- a/internal/db/queries_ipqueue.go +++ b/internal/db/queries_ipqueue.go @@ -101,7 +101,7 @@ func (d *DB) ClaimNextQueued(ctx context.Context, validatorID string, leaseTTL t res, err = tx.ExecContext(ctx, ` UPDATE validators SET state=?, current_ip_id=?, updated_at=? - WHERE validator_id=? AND state=? + WHERE validator_id=? AND state=? AND current_ip_id IS NULL `, ValidatorAssigned, item.ID, timeToDB(now), validatorID, ValidatorIdle) if err != nil { return nil, err @@ -207,7 +207,8 @@ func (d *DB) FinishIP(ctx context.Context, ipID int64, result string) error { } // ReleaseFIP records that the floating IP has been disassociated and frees -// the owning validator back to idle, in one transaction. +// the owning validator (if this address is still its current one, see +// freeValidatorSQL), in one transaction. func (d *DB) ReleaseFIP(ctx context.Context, ipID int64, validatorID string) error { tx, err := d.BeginTx(ctx, nil) if err != nil { @@ -219,10 +220,7 @@ func (d *DB) ReleaseFIP(ctx context.Context, ipID int64, validatorID string) err if _, err := tx.ExecContext(ctx, `UPDATE ip_queue SET fip_released_at=?, updated_at=? WHERE id=?`, now, now, ipID); err != nil { return err } - if _, err := tx.ExecContext(ctx, ` - UPDATE validators SET state=?, current_ip_id=NULL, updated_at=? - WHERE validator_id=? - `, ValidatorIdle, now, validatorID); err != nil { + if _, err := tx.ExecContext(ctx, freeValidatorSQL, now, validatorID, ipID); err != nil { return err } return tx.Commit() @@ -252,10 +250,7 @@ func (d *DB) MarkFIPOccupied(ctx context.Context, ipID int64, validatorID string return err } if validatorID != "" { - if _, err := tx.ExecContext(ctx, ` - UPDATE validators SET state=?, current_ip_id=NULL, updated_at=? - WHERE validator_id=? - `, ValidatorIdle, now, validatorID); err != nil { + if _, err := tx.ExecContext(ctx, freeValidatorSQL, now, validatorID, ipID); err != nil { return err } } @@ -315,10 +310,7 @@ func (d *DB) RequeueOrFail(ctx context.Context, ipID int64, validatorID string, } if validatorID != "" { - if _, err := tx.ExecContext(ctx, ` - UPDATE validators SET state=?, current_ip_id=NULL, updated_at=? - WHERE validator_id=? - `, ValidatorIdle, now, validatorID); err != nil { + if _, err := tx.ExecContext(ctx, freeValidatorSQL, now, validatorID, ipID); err != nil { return err } } @@ -520,9 +512,11 @@ func (d *DB) DeleteIPs(ctx context.Context, addresses []string) (DeleteIPsResult func deleteIPTx(ctx context.Context, tx *sql.Tx, ipID int64) error { now := timeToDB(Now()) if _, err := tx.ExecContext(ctx, ` - UPDATE validators SET state=?, current_ip_id=NULL, updated_at=? + UPDATE validators SET current_ip_id=NULL, + state = CASE WHEN state=? THEN state ELSE ? END, + updated_at=? WHERE current_ip_id=? - `, ValidatorIdle, now, ipID); err != nil { + `, ValidatorUnreachable, ValidatorIdle, now, ipID); err != nil { return fmt.Errorf("free owning validator: %w", err) } if _, err := tx.ExecContext(ctx, `UPDATE checks SET ip_id=NULL WHERE ip_id=?`, ipID); err != nil { @@ -579,9 +573,11 @@ func (d *DB) ClearAllIPs(ctx context.Context) ([]string, error) { now := timeToDB(Now()) if _, err := tx.ExecContext(ctx, ` - UPDATE validators SET state=?, current_ip_id=NULL, updated_at=? + UPDATE validators SET current_ip_id=NULL, + state = CASE WHEN state=? THEN state ELSE ? END, + updated_at=? WHERE current_ip_id IS NOT NULL - `, ValidatorIdle, now); err != nil { + `, ValidatorUnreachable, ValidatorIdle, now); err != nil { return nil, fmt.Errorf("free owning validators: %w", err) } if _, err := tx.ExecContext(ctx, `UPDATE checks SET ip_id=NULL WHERE ip_id IS NOT NULL`); err != nil { @@ -610,11 +606,15 @@ type FIPRef struct { FIPID string } -// ListFIPRefs returns every queue row with an attached floating IP (fip_id -// set) — typically at most one per validator — so a bulk clear can -// disassociate them without loading the whole queue. +// ListFIPRefs returns the queue rows that may still hold a floating IP: an +// fip_id is set and the row is not finished. A finished row (done, failed, +// occupied) keeps its fip_id for display, but its floating IP was already +// disassociated before the final state was written, so listing it would only +// make a bulk clear issue thousands of pointless cloud calls. At most one row +// per validator qualifies, so a clear does not need to load the whole queue. func (d *DB) ListFIPRefs(ctx context.Context) ([]FIPRef, error) { - rows, err := d.QueryContext(ctx, `SELECT id, ip_address, fip_id FROM ip_queue WHERE fip_id<>'' ORDER BY id`) + rows, err := d.QueryContext(ctx, `SELECT id, ip_address, fip_id FROM ip_queue WHERE fip_id<>'' AND state NOT IN (?, ?, ?) ORDER BY id`, + IPDone, IPFailed, IPOccupied) if err != nil { return nil, err } @@ -631,8 +631,7 @@ func (d *DB) ListFIPRefs(ctx context.Context) ([]FIPRef, error) { } // ListFIPRefsByAddresses is ListFIPRefs restricted to the given addresses -// (unknown addresses and rows without an attached floating IP are simply -// absent), using a handful of IN (...) queries instead of one lookup per +// (unknown addresses and rows that hold no floating IP are simply absent), using a handful of IN (...) queries instead of one lookup per // address. func (d *DB) ListFIPRefsByAddresses(ctx context.Context, addresses []string) ([]FIPRef, error) { const chunk = 500 @@ -643,12 +642,12 @@ func (d *DB) ListFIPRefsByAddresses(ctx context.Context, addresses []string) ([] end = len(addresses) } part := addresses[start:end] - args := make([]any, len(part)) - for i, a := range part { - args[i] = a + args := []any{IPDone, IPFailed, IPOccupied} + for _, a := range part { + args = append(args, a) } rows, err := d.QueryContext(ctx, - `SELECT id, ip_address, fip_id FROM ip_queue WHERE fip_id<>'' AND ip_address IN (`+placeholders(len(part))+`)`, args...) + `SELECT id, ip_address, fip_id FROM ip_queue WHERE fip_id<>'' AND state NOT IN (?, ?, ?) AND ip_address IN (`+placeholders(len(part))+`)`, args...) if err != nil { return nil, err } diff --git a/internal/db/queries_validator_state_test.go b/internal/db/queries_validator_state_test.go new file mode 100644 index 0000000..f8bebdb --- /dev/null +++ b/internal/db/queries_validator_state_test.go @@ -0,0 +1,186 @@ +package db + +import ( + "context" + "fmt" + "testing" + "time" +) + +func validatorState(t *testing.T, d *DB, id string) *Validator { + t.Helper() + v, err := d.GetValidator(testCtx(t), id) + if err != nil { + t.Fatalf("get validator %s: %v", id, err) + } + return v +} + +func testCtx(t *testing.T) context.Context { + t.Helper() + return context.Background() +} + +// claimFor2 seeds an address and returns its id without claiming it. +func claimFor2(t *testing.T, d *DB, addr string) int64 { + t.Helper() + ctx := testCtx(t) + if err := d.SeedQueue(ctx, []string{addr}); err != nil { + t.Fatalf("seed %s: %v", addr, err) + } + ip, err := d.GetIPByAddress(ctx, addr) + if err != nil { + t.Fatalf("get %s: %v", addr, err) + } + return ip.ID +} + +// claimFor seeds an address and claims it for the validator. +func claimFor(t *testing.T, d *DB, addr, validatorID string) *IPQueueItem { + t.Helper() + ctx := testCtx(t) + if err := d.SeedQueue(ctx, []string{addr}); err != nil { + t.Fatalf("seed %s: %v", addr, err) + } + item, err := d.ClaimNextQueued(ctx, validatorID, time.Minute) + if err != nil || item == nil { + t.Fatalf("claim %s for %s: item=%v err=%v", addr, validatorID, item, err) + } + return item +} + +func TestHeartbeatKeepsAssignedWhenValidatorStillHoldsAnAddress(t *testing.T) { + d, ctx := newTestDB(t) + _ = d.AdminCreateValidator(ctx, "v1", "p1") + _ = d.AdminCreateValidator(ctx, "v2", "p2") + item := claimFor(t, d, "1.1.1.1", "v1") + _ = d.MarkValidatorUnreachable(ctx, "v1") + _ = d.MarkValidatorUnreachable(ctx, "v2") + + if err := d.Heartbeat(ctx, "v1"); err != nil { + t.Fatal(err) + } + if v := validatorState(t, d, "v1"); v.State != ValidatorAssigned || v.CurrentIPID == nil || *v.CurrentIPID != item.ID { + t.Fatalf("v1 after heartbeat: state=%s current_ip=%v, want assigned to %d", v.State, v.CurrentIPID, item.ID) + } + if err := d.Heartbeat(ctx, "v2"); err != nil { + t.Fatal(err) + } + if v := validatorState(t, d, "v2"); v.State != ValidatorIdle { + t.Fatalf("v2 after heartbeat: state=%s, want idle", v.State) + } +} + +func TestRegisterValidatorReactivationKeepsAssignedAddress(t *testing.T) { + d, ctx := newTestDB(t) + _ = d.AdminCreateValidator(ctx, "v1", "p1") + item := claimFor(t, d, "1.1.1.1", "v1") + _ = d.MarkValidatorUnreachable(ctx, "v1") + + if err := d.RegisterValidator(ctx, "v1", "host", "p1", "v"); err != nil { + t.Fatal(err) + } + if v := validatorState(t, d, "v1"); v.State != ValidatorAssigned || v.CurrentIPID == nil || *v.CurrentIPID != item.ID { + t.Fatalf("after re-register: state=%s current_ip=%v, want assigned to %d", v.State, v.CurrentIPID, item.ID) + } +} + +func TestStaleReleasesLeaveTheCurrentAddressAlone(t *testing.T) { + cases := map[string]func(d *DB, ipID int64) error{ + "ReleaseFIP": func(d *DB, id int64) error { return d.ReleaseFIP(context.Background(), id, "v1") }, + "RequeueOrFail": func(d *DB, id int64) error { return d.RequeueOrFail(context.Background(), id, "v1", 3) }, + "MarkFIPOccupied": func(d *DB, id int64) error { return d.MarkFIPOccupied(context.Background(), id, "v1") }, + "FreeValidator": func(d *DB, id int64) error { return d.FreeValidator(context.Background(), "v1", id) }, + } + for name, release := range cases { + t.Run(name, func(t *testing.T) { + d, ctx := newTestDB(t) + _ = d.AdminCreateValidator(ctx, "v1", "p1") + stale := claimFor(t, d, "1.1.1.1", "v1") + // v1 has moved on to another address (set directly: the claim path + // itself refuses a validator that is still busy). + cur := claimFor2(t, d, "2.2.2.2") + if _, err := d.ExecContext(ctx, `UPDATE ip_queue SET state='checking', owner_validator_id='v1' WHERE id=?`, cur); err != nil { + t.Fatal(err) + } + if _, err := d.ExecContext(ctx, `UPDATE validators SET state='assigned', current_ip_id=? WHERE validator_id='v1'`, cur); err != nil { + t.Fatal(err) + } + + if err := release(d, stale.ID); err != nil { + t.Fatalf("%s: %v", name, err) + } + v := validatorState(t, d, "v1") + if v.State != ValidatorAssigned || v.CurrentIPID == nil || *v.CurrentIPID != cur { + t.Fatalf("%s freed a validator that holds another address: state=%s current_ip=%v", name, v.State, v.CurrentIPID) + } + }) + } +} + +// The right address frees the validator, but an unreachable one stays so. +func TestReleaseKeepsUnreachableValidatorUnreachable(t *testing.T) { + d, ctx := newTestDB(t) + _ = d.AdminCreateValidator(ctx, "v1", "p1") + _ = d.AdminCreateValidator(ctx, "v2", "p2") + a := claimFor(t, d, "1.1.1.1", "v1") + b := claimFor(t, d, "2.2.2.2", "v2") + _ = d.MarkValidatorUnreachable(ctx, "v1") + + if err := d.ReleaseFIP(ctx, a.ID, "v1"); err != nil { + t.Fatal(err) + } + if v := validatorState(t, d, "v1"); v.State != ValidatorUnreachable || v.CurrentIPID != nil { + t.Fatalf("v1: state=%s current_ip=%v, want unreachable and empty", v.State, v.CurrentIPID) + } + if err := d.RequeueOrFail(ctx, b.ID, "v2", 3); err != nil { + t.Fatal(err) + } + if v := validatorState(t, d, "v2"); v.State != ValidatorIdle || v.CurrentIPID != nil { + t.Fatalf("v2: state=%s current_ip=%v, want idle and empty", v.State, v.CurrentIPID) + } +} + +func TestClaimRefusesValidatorThatStillHoldsAnAddress(t *testing.T) { + d, ctx := newTestDB(t) + _ = d.AdminCreateValidator(ctx, "v1", "p1") + first := claimFor(t, d, "1.1.1.1", "v1") + // An old version could leave a validator idle while it still pointed at an address. + if _, err := d.ExecContext(ctx, `UPDATE validators SET state='idle' WHERE validator_id='v1'`); err != nil { + t.Fatal(err) + } + _ = d.SeedQueue(ctx, []string{"2.2.2.2"}) + got, err := d.ClaimNextQueued(ctx, "v1", time.Minute) + if err != nil || got != nil { + t.Fatalf("claim for a validator that holds %d: item=%v err=%v, want nothing", first.ID, got, err) + } + if ip, _ := d.GetIPByAddress(ctx, "2.2.2.2"); ip.State != IPQueued { + t.Fatalf("2.2.2.2 is %s, want queued", ip.State) + } +} + +func TestListFIPRefsSkipsFinishedAddresses(t *testing.T) { + d, ctx := newTestDB(t) + var addrs []string + for i := 0; i < 6; i++ { + addrs = append(addrs, fmt.Sprintf("10.0.0.%d", i+1)) + } + _ = d.SeedQueue(ctx, addrs) + states := []string{"done", "failed", "occupied", "awaiting_self_check", "checking", "aggregating"} + for i, a := range addrs { + if _, err := d.ExecContext(ctx, `UPDATE ip_queue SET state=?, fip_id=? WHERE ip_address=?`, states[i], fmt.Sprintf("fip-%d", i), a); err != nil { + t.Fatal(err) + } + } + refs, err := d.ListFIPRefs(ctx) + if err != nil { + t.Fatal(err) + } + if len(refs) != 3 { + t.Fatalf("ListFIPRefs returned %d rows, want the 3 unfinished ones: %+v", len(refs), refs) + } + byAddr, err := d.ListFIPRefsByAddresses(ctx, addrs) + if err != nil || len(byAddr) != 3 { + t.Fatalf("ListFIPRefsByAddresses returned %d rows (err %v), want 3", len(byAddr), err) + } +} diff --git a/internal/db/queries_validators.go b/internal/db/queries_validators.go index 8b94a29..3ff810e 100644 --- a/internal/db/queries_validators.go +++ b/internal/db/queries_validators.go @@ -27,24 +27,32 @@ func (d *DB) RegisterValidator(ctx context.Context, validatorID, hostname, osPor } // A brand-new row already lands in ValidatorIdle via the INSERT branch; // a re-registering validator that was 'unregistered' or 'unreachable' - // (but not mid-assignment) should also come back to idle. + // comes back: to idle when it holds no address, to assigned when it still + // does (its current_ip_id is kept, so it must not be handed another). _, err = d.ExecContext(ctx, ` - UPDATE validators SET state=?, updated_at=? + UPDATE validators SET + state = CASE WHEN current_ip_id IS NULL THEN ? ELSE ? END, + updated_at=? WHERE validator_id=? AND state IN (?, ?) - `, ValidatorIdle, now, validatorID, ValidatorUnregistered, ValidatorUnreachable) + `, ValidatorIdle, ValidatorAssigned, now, validatorID, ValidatorUnregistered, ValidatorUnreachable) if err != nil { return fmt.Errorf("register validator (reactivate): %w", err) } return nil } +// Heartbeat records a sign of life. A validator that was marked unreachable +// returns to idle if it holds no address, but to assigned if it still does: +// it was only silent (for example busy with slow checks), and handing it a +// second address while it works on the first would leave the second one +// without an owner that can ever pick it up. func (d *DB) Heartbeat(ctx context.Context, validatorID string) error { now := timeToDB(Now()) res, err := d.ExecContext(ctx, ` UPDATE validators SET last_heartbeat_at=?, updated_at=?, - state = CASE WHEN state=? THEN ? ELSE state END + state = CASE WHEN state=? THEN (CASE WHEN current_ip_id IS NULL THEN ? ELSE ? END) ELSE state END WHERE validator_id=? - `, now, now, ValidatorUnreachable, ValidatorIdle, validatorID) + `, now, now, ValidatorUnreachable, ValidatorIdle, ValidatorAssigned, validatorID) if err != nil { return fmt.Errorf("heartbeat: %w", err) } @@ -138,16 +146,59 @@ func (d *DB) MarkValidatorUnreachable(ctx context.Context, validatorID string) e return err } -// FreeValidator returns a validator to idle with no assigned IP. Used after -// an IP finishes (success or failure) or is reclaimed by the lease sweep. -func (d *DB) FreeValidator(ctx context.Context, validatorID string) error { - _, err := d.ExecContext(ctx, ` - UPDATE validators SET state=?, current_ip_id=NULL, updated_at=? - WHERE validator_id=? - `, ValidatorIdle, timeToDB(Now()), validatorID) +// freeValidatorSQL releases a validator from the address it holds. It only +// applies when the validator's current address is the one being released +// (args: now, validator id, ip id): a late release of an old address must not +// free a validator that has already moved on to another one. An unreachable +// validator stays unreachable until its next heartbeat, so a dead validator +// is not handed new addresses just because its lease was reclaimed. +var freeValidatorSQL = fmt.Sprintf(` + UPDATE validators SET current_ip_id=NULL, + state = CASE WHEN state='%s' THEN state ELSE '%s' END, + updated_at=? + WHERE validator_id=? AND current_ip_id=?`, ValidatorUnreachable, ValidatorIdle) + +// FreeValidator releases a validator from the given address (see +// freeValidatorSQL). Used after an IP finishes (success or failure) or is +// reclaimed by the lease sweep. +func (d *DB) FreeValidator(ctx context.Context, validatorID string, ipID int64) error { + _, err := d.ExecContext(ctx, freeValidatorSQL, timeToDB(Now()), validatorID, ipID) return err } +// ReconcileValidators repairs validators whose state disagrees with the +// queue: a validator pointing at an address that no longer exists, is +// finished, or belongs to another validator is released; an "assigned" +// validator that holds nothing goes back to idle. The invariants normally +// hold by construction (every change is one transaction); this heals what a +// crash or an older version left behind. It returns the number of validators +// repaired. +func (d *DB) ReconcileValidators(ctx context.Context) (int64, error) { + now := timeToDB(Now()) + res, err := d.ExecContext(ctx, fmt.Sprintf(` + UPDATE validators SET current_ip_id=NULL, + state = CASE WHEN state='%s' THEN state ELSE '%s' END, + updated_at=? + WHERE current_ip_id IS NOT NULL AND NOT EXISTS ( + SELECT 1 FROM ip_queue q + WHERE q.id = validators.current_ip_id + AND q.owner_validator_id = validators.validator_id + AND q.state NOT IN ('%s','%s','%s'))`, + ValidatorUnreachable, ValidatorIdle, IPDone, IPFailed, IPOccupied), now) + if err != nil { + return 0, fmt.Errorf("reconcile validators: %w", err) + } + n, _ := res.RowsAffected() + res, err = d.ExecContext(ctx, ` + UPDATE validators SET state=?, updated_at=? + WHERE state=? AND current_ip_id IS NULL`, ValidatorIdle, now, ValidatorAssigned) + if err != nil { + return n, fmt.Errorf("reconcile validators: %w", err) + } + m, _ := res.RowsAffected() + return n + m, nil +} + // AdminCreateValidator registers a brand-new validator via the admin API. // Unlike RegisterValidator (used by the agent's self-registration call), // this refuses to upsert over an existing row. diff --git a/internal/httpapi/handlers_admin.go b/internal/httpapi/handlers_admin.go index a234089..76b014e 100644 --- a/internal/httpapi/handlers_admin.go +++ b/internal/httpapi/handlers_admin.go @@ -1,17 +1,32 @@ package httpapi import ( + "context" "errors" "fmt" "net/http" "net/url" "strconv" "strings" + "time" "cloudipvalidator/internal/db" "cloudipvalidator/internal/orchestrator" ) +// destructiveOpTimeout bounds a cancel, delete or clear. Such an operation +// talks to the cloud as well as the database, and abandoning it halfway +// leaves floating IPs detached from rows that still exist (or the reverse), +// so it must not die with the client connection: a client that gives up +// (a dashboard or curl timeout) only stops waiting for the answer. +const destructiveOpTimeout = 10 * time.Minute + +// detachedContext returns a context that ignores cancellation of the request +// but keeps its values, with destructiveOpTimeout as the upper bound. +func detachedContext(r *http.Request) (context.Context, context.CancelFunc) { + return context.WithTimeout(context.WithoutCancel(r.Context()), destructiveOpTimeout) +} + func (s *Server) handleHealthz(w http.ResponseWriter, r *http.Request) { writeJSON(w, http.StatusOK, okResponse{OK: true}) } @@ -273,7 +288,9 @@ func toScanStatusDTO(st orchestrator.ScanStatus) scanStatusDTO { // need to be disassociated in OpenStack. func (s *Server) handleAdminCancelIP(w http.ResponseWriter, r *http.Request) { address := r.PathValue("ip") - if err := s.Orch.ForceCancel(r.Context(), address); err != nil { + ctx, cancel := detachedContext(r) + defer cancel() + if err := s.Orch.ForceCancel(ctx, address); err != nil { writeDBError(w, err) return } @@ -286,7 +303,9 @@ func (s *Server) handleAdminCancelIP(w http.ResponseWriter, r *http.Request) { // associated floating IP needs disassociating first. func (s *Server) handleAdminDeleteIP(w http.ResponseWriter, r *http.Request) { address := r.PathValue("ip") - if err := s.Orch.DeleteIP(r.Context(), address); err != nil { + ctx, cancel := detachedContext(r) + defer cancel() + if err := s.Orch.DeleteIP(ctx, address); err != nil { writeDBError(w, err) return } @@ -306,7 +325,9 @@ func (s *Server) handleAdminDeleteIPs(w http.ResponseWriter, r *http.Request) { writeError(w, http.StatusBadRequest, "addresses must not be empty") return } - result, err := s.Orch.DeleteIPs(r.Context(), req.Addresses) + ctx, cancel := detachedContext(r) + defer cancel() + result, err := s.Orch.DeleteIPs(ctx, req.Addresses) if err != nil { writeDBError(w, err) return @@ -320,7 +341,9 @@ func (s *Server) handleAdminDeleteIPs(w http.ResponseWriter, r *http.Request) { // handleAdminClearQueue permanently removes every address currently in the // queue, including those actively being checked. func (s *Server) handleAdminClearQueue(w http.ResponseWriter, r *http.Request) { - result, err := s.Orch.ClearQueue(r.Context()) + ctx, cancel := detachedContext(r) + defer cancel() + result, err := s.Orch.ClearQueue(ctx) if err != nil { writeDBError(w, err) return diff --git a/internal/httpapi/handlers_detached_test.go b/internal/httpapi/handlers_detached_test.go new file mode 100644 index 0000000..e99dfc8 --- /dev/null +++ b/internal/httpapi/handlers_detached_test.go @@ -0,0 +1,65 @@ +package httpapi + +import ( + "context" + "net/http" + "testing" + "time" + + "cloudipvalidator/internal/openstack" +) + +// slowDetachOS delays every Disassociate so the client can give up first. +type slowDetachOS struct { + *openstack.MockClient + delay time.Duration +} + +func (s slowDetachOS) DisassociateFloatingIP(ctx context.Context, fipID string) error { + time.Sleep(s.delay) + return s.MockClient.DisassociateFloatingIP(ctx, fipID) +} + +// A client that gives up (a dashboard or curl timeout) must not abort +// "clear queue" halfway: on 2026-10-02 the request context was cancelled after +// part of the floating IPs were detached, the database was left untouched and +// the checks kept running. +func TestClearQueueSurvivesClientDisconnect(t *testing.T) { + fc, d, orch, mock := newConfigTestHarness(t) + ctx := context.Background() + mock.Seed("fip-1", "1.2.3.4", "svc") + if err := d.RegisterValidator(ctx, "validator-1", "host", "port-1", "v"); err != nil { + t.Fatal(err) + } + if err := d.SeedQueue(ctx, []string{"1.2.3.4", "5.6.7.8"}); err != nil { + t.Fatal(err) + } + orch.Tick(ctx) // validator-1 takes 1.2.3.4 and attaches its floating IP + if fip, _ := mock.GetFloatingIPByAddress(ctx, "1.2.3.4"); fip.PortID != "port-1" { + t.Fatalf("setup: floating ip not attached (port %q)", fip.PortID) + } + orch.OS = slowDetachOS{MockClient: mock, delay: 400 * time.Millisecond} + + reqCtx, cancel := context.WithTimeout(ctx, 100*time.Millisecond) + defer cancel() + // No body: the server only notices that the client has left (and cancels + // the request context) when it is not waiting for a request body. + req, _ := http.NewRequestWithContext(reqCtx, http.MethodPost, fc.base+"/api/v1/admin/ips/clear", nil) + if resp, err := fc.client.Do(req); err == nil { + resp.Body.Close() + t.Fatalf("the client was expected to time out, got status %d", resp.StatusCode) + } + + deadline := time.Now().Add(5 * time.Second) + for { + ips, _ := d.ListIPs(ctx) + fip, _ := mock.GetFloatingIPByAddress(ctx, "1.2.3.4") + if len(ips) == 0 && fip.PortID == "" { + return + } + if time.Now().After(deadline) { + t.Fatalf("clear did not finish after the client left: %d rows, floating ip on port %q", len(ips), fip.PortID) + } + time.Sleep(50 * time.Millisecond) + } +} diff --git a/internal/orchestrator/orchestrator.go b/internal/orchestrator/orchestrator.go index 7aa9664..c09744a 100644 --- a/internal/orchestrator/orchestrator.go +++ b/internal/orchestrator/orchestrator.go @@ -114,6 +114,11 @@ func (o *Orchestrator) leaseTTL() time.Duration { // own goroutine, so validators never wait for each other. In Async mode // Tick does not wait for those goroutines; otherwise it waits for them. func (o *Orchestrator) Tick(ctx context.Context) { + if n, err := o.DB.ReconcileValidators(ctx); err != nil { + o.Log.Error("reconcile validators", "err", err) + } else if n > 0 { + o.Log.Warn("repaired validators that disagreed with the queue", "count", n) + } if err := o.assignIdleValidators(ctx); err != nil { o.Log.Error("assign idle validators", "err", err) } @@ -146,7 +151,11 @@ func (o *Orchestrator) assignIdleValidators(ctx context.Context) error { } o.Log.Info("claimed ip", "validator", v.ValidatorID, "ip", item.IPAddress, "ip_id", item.ID) v, item := v, item - o.spawn("assign:"+v.ValidatorID, func() { + // Keyed by address, not validator: the key only guards against + // starting the same association twice. A validator-wide key made a + // second address claimed while the first was still associating skip + // its association and wait for the lease to expire. + o.spawn(fmt.Sprintf("assign:%d", item.ID), func() { if err := o.associateFIP(ctx, v.ValidatorID, v.OSPortID, item); err != nil { o.Log.Error("associate fip", "validator", v.ValidatorID, "ip", item.IPAddress, "err", err) } @@ -466,7 +475,7 @@ func (o *Orchestrator) ForceCancel(ctx context.Context, ipAddress string) error } if item.OwnerValidatorID != nil { - if err := o.DB.FreeValidator(ctx, *item.OwnerValidatorID); err != nil { + if err := o.DB.FreeValidator(ctx, *item.OwnerValidatorID, item.ID); err != nil { return fmt.Errorf("free validator: %w", err) } o.releaseValidatorPorts(ctx, o.validatorsByID(ctx, *item.OwnerValidatorID)) @@ -522,11 +531,7 @@ func (o *Orchestrator) DeleteIPs(ctx context.Context, addresses []string) (db.De if err != nil { o.Log.Error("list attached fips before delete", "err", err) } - for _, ref := range refs { - if err := o.OS.DisassociateFloatingIP(ctx, ref.FIPID); err != nil { - o.Log.Error("disassociate fip on delete", "ip_id", ref.IPID, "fip_id", ref.FIPID, "err", err) - } - } + o.disassociateAll(ctx, refs, "delete") owners := o.busyValidatorsFor(ctx, addresses) result, err := o.DB.DeleteIPs(ctx, addresses) @@ -538,6 +543,30 @@ func (o *Orchestrator) DeleteIPs(ctx context.Context, addresses []string) (db.De return result, nil } +// maxParallelDetach bounds how many floating IPs a bulk delete or clear +// detaches at the same time (each is one Neutron call). +const maxParallelDetach = 8 + +// disassociateAll detaches the given floating IPs best-effort, at most +// maxParallelDetach at a time; a failure is logged and does not stop the rest. +func (o *Orchestrator) disassociateAll(ctx context.Context, refs []db.FIPRef, what string) { + var wg sync.WaitGroup + sem := make(chan struct{}, maxParallelDetach) + for _, ref := range refs { + ref := ref + sem <- struct{}{} + wg.Add(1) + go func() { + defer wg.Done() + defer func() { <-sem }() + if err := o.OS.DisassociateFloatingIP(ctx, ref.FIPID); err != nil { + o.Log.Error("disassociate fip on "+what, "ip_id", ref.IPID, "fip_id", ref.FIPID, "err", err) + } + }() + } + wg.Wait() +} + // ClearQueue deletes every address currently in the queue, regardless of // state — the "delete everything" operation. It is set-based (see // db.ClearAllIPs): O(1) statements however many rows there are. Floating IPs @@ -548,11 +577,7 @@ func (o *Orchestrator) ClearQueue(ctx context.Context) (db.DeleteIPsResult, erro if err != nil { return db.DeleteIPsResult{}, fmt.Errorf("list attached fips: %w", err) } - for _, ref := range refs { - if err := o.OS.DisassociateFloatingIP(ctx, ref.FIPID); err != nil { - o.Log.Error("disassociate fip on clear queue", "ip_id", ref.IPID, "fip_id", ref.FIPID, "err", err) - } - } + o.disassociateAll(ctx, refs, "clear queue") deleted, err := o.DB.ClearAllIPs(ctx) if err != nil { diff --git a/internal/orchestrator/validator_state_test.go b/internal/orchestrator/validator_state_test.go new file mode 100644 index 0000000..b62d011 --- /dev/null +++ b/internal/orchestrator/validator_state_test.go @@ -0,0 +1,401 @@ +package orchestrator + +import ( + "context" + "fmt" + "math/rand" + "sync/atomic" + "testing" + "time" + + "cloudipvalidator/internal/db" + "cloudipvalidator/internal/openstack" +) + +// finishChecks reports a successful egress run and all three inbound sites +// for the address, so the next Tick aggregates it and releases the validator. +func finishChecks(t *testing.T, o *Orchestrator, ip *db.IPQueueItem, validatorID string) { + t.Helper() + ctx := context.Background() + if err := o.RecordCheck(ctx, db.Check{ + IPID: ip.ID, IPAddress: ip.IPAddress, AttemptNumber: ip.AttemptNumber, + ValidatorID: validatorID, Source: db.SourceEgress, CheckType: "https", + Target: "https://example.test", Success: true, CheckedAt: db.Now(), + }); err != nil { + t.Fatalf("record egress check: %v", err) + } + if err := o.MarkEgressComplete(ctx, ip.ID); err != nil { + t.Fatalf("mark egress complete: %v", err) + } + for site := 1; site <= 3; site++ { + for _, ct := range []string{"tcp-22", "ssh", "tcp-80", "icmp"} { + if err := o.RecordCheck(ctx, db.Check{ + IPID: ip.ID, IPAddress: ip.IPAddress, AttemptNumber: ip.AttemptNumber, + Source: db.InboundSource(site), CheckType: ct, Target: ip.IPAddress, + Success: true, CheckedAt: db.Now(), + }); err != nil { + t.Fatalf("record inbound check: %v", err) + } + } + if err := o.MarkSiteComplete(ctx, ip.ID, site); err != nil { + t.Fatalf("mark site complete: %v", err) + } + } +} + +// The incident of 2026-10-02: a validator busy with slow checks goes silent, +// is marked unreachable, and its next heartbeat used to return it to idle +// although it still held the address. It was then handed a second address, +// whose association never ran, and both stalled until their leases expired. +func TestSilentValidatorKeepsItsAddressAfterHeartbeat(t *testing.T) { + ctx := context.Background() + o, d, mock := newTestOrchestrator(t, 180) + mock.Seed("fip-a", "1.1.1.1", "svc") + mock.Seed("fip-b", "2.2.2.2", "svc") + if err := d.RegisterValidator(ctx, "validator-1", "host", "port-1", "v"); err != nil { + t.Fatal(err) + } + if err := d.SeedQueue(ctx, []string{"1.1.1.1", "2.2.2.2"}); err != nil { + t.Fatal(err) + } + + o.Tick(ctx) // validator-1 claims 1.1.1.1 + a, _ := d.GetIPByAddress(ctx, "1.1.1.1") + if err := o.SelfCheckResult(ctx, "validator-1", a.ID, true, "ok"); err != nil { + t.Fatal(err) + } + + // The agent goes silent for longer than heartbeat_timeout_seconds ... + if err := d.MarkValidatorUnreachable(ctx, "validator-1"); err != nil { + t.Fatal(err) + } + // ... and then speaks again while the checks are still running. + if err := d.Heartbeat(ctx, "validator-1"); err != nil { + t.Fatal(err) + } + v, _ := d.GetValidator(ctx, "validator-1") + if v.State != db.ValidatorAssigned || v.CurrentIPID == nil || *v.CurrentIPID != a.ID { + t.Fatalf("after heartbeat: state=%s current_ip=%v, want assigned to %d", v.State, v.CurrentIPID, a.ID) + } + + o.Tick(ctx) + if b, _ := d.GetIPByAddress(ctx, "2.2.2.2"); b.State != db.IPQueued { + t.Fatalf("the busy validator was given a second address: 2.2.2.2 is %s", b.State) + } + + // The first address finishes: only now the validator may take the next one. + a, _ = d.GetIP(ctx, a.ID) + finishChecks(t, o, a, "validator-1") + o.Tick(ctx) // aggregates and releases 1.1.1.1 + o.Tick(ctx) // validator-1 is idle again and claims 2.2.2.2 + if b, _ := d.GetIPByAddress(ctx, "2.2.2.2"); b.State != db.IPAwaitingSelfCheck { + t.Fatalf("2.2.2.2 is %s, want awaiting_self_check once the validator is free", b.State) + } +} + +// A dead validator whose lease was reclaimed must stay out of rotation until +// it speaks again; before, reclaiming set it idle and it was handed a fresh +// address every lease period, burning each address's retries. +func TestUnreachableValidatorGetsNoAddressesAfterLeaseReclaim(t *testing.T) { + ctx := context.Background() + o, d, mock := newTestOrchestrator(t, 1) + mock.Seed("fip-a", "1.1.1.1", "svc") + mock.Seed("fip-b", "2.2.2.2", "svc") + _ = d.RegisterValidator(ctx, "validator-1", "host", "port-1", "v") + _ = d.SeedQueue(ctx, []string{"1.1.1.1", "2.2.2.2"}) + + o.Tick(ctx) // claims 1.1.1.1, the agent never answers + if err := d.MarkValidatorUnreachable(ctx, "validator-1"); err != nil { + t.Fatal(err) + } + time.Sleep(1100 * time.Millisecond) + o.Cfg.LeaseTTLSeconds = 180 + o.Tick(ctx) // lease sweep reclaims 1.1.1.1 + + a, _ := d.GetIPByAddress(ctx, "1.1.1.1") + if a.RetryCount != 1 { + t.Fatalf("retry_count = %d, want the lease to be reclaimed once", a.RetryCount) + } + v, _ := d.GetValidator(ctx, "validator-1") + if v.State != db.ValidatorUnreachable || v.CurrentIPID != nil { + t.Fatalf("validator state=%s current_ip=%v, want unreachable and empty", v.State, v.CurrentIPID) + } + o.Tick(ctx) + if a, _ := d.GetIPByAddress(ctx, "1.1.1.1"); a.State != db.IPQueued { + t.Fatalf("an unreachable validator was handed an address: 1.1.1.1 is %s", a.State) + } + + if err := d.Heartbeat(ctx, "validator-1"); err != nil { // it is back + t.Fatal(err) + } + if v, _ := d.GetValidator(ctx, "validator-1"); v.State != db.ValidatorIdle { + t.Fatalf("after heartbeat state=%s, want idle (it holds nothing)", v.State) + } + o.Tick(ctx) + if a, _ := d.GetIPByAddress(ctx, "1.1.1.1"); a.State != db.IPAwaitingSelfCheck { + t.Fatalf("1.1.1.1 is %s, want it picked up again", a.State) + } +} + +// countingOS counts Disassociate calls. +type countingOS struct { + *openstack.MockClient + disassociations int32 +} + +func (c *countingOS) DisassociateFloatingIP(ctx context.Context, fipID string) error { + atomic.AddInt32(&c.disassociations, 1) + return c.MockClient.DisassociateFloatingIP(ctx, fipID) +} + +// Finished addresses keep their fip_id for display but their floating IP is +// already free; clearing the queue must not make a cloud call for each of +// them (on 2026-10-02 that was >1000 sequential calls: 256 s). +func TestClearQueueDoesNotTouchFinishedAddresses(t *testing.T) { + ctx := context.Background() + o, d, mock := newTestOrchestrator(t, 180) + cos := &countingOS{MockClient: mock} + o.OS = cos + + const finished = 300 + var addrs []string + for i := 0; i < finished; i++ { + addrs = append(addrs, fmt.Sprintf("10.9.%d.%d", i/200, i%200+1)) + } + if err := d.SeedQueue(ctx, addrs); err != nil { + t.Fatal(err) + } + for i, a := range addrs { + state := []string{db.IPDone, db.IPFailed, db.IPOccupied}[i%3] + if _, err := d.ExecContext(ctx, `UPDATE ip_queue SET state=?, fip_id=? WHERE ip_address=?`, state, fmt.Sprintf("fip-old-%d", i), a); err != nil { + t.Fatal(err) + } + } + // One address really holds a floating IP. + mock.Seed("fip-live", "1.2.3.4", "svc") + _ = d.RegisterValidator(ctx, "validator-1", "host", "port-1", "v") + _ = d.SeedQueue(ctx, []string{"1.2.3.4"}) + live, _ := d.GetIPByAddress(ctx, "1.2.3.4") + if _, err := d.ExecContext(ctx, `UPDATE ip_queue SET sequence=-1 WHERE id=?`, live.ID); err != nil { + t.Fatal(err) + } + o.Tick(ctx) // validator-1 takes the live address (lowest sequence) and attaches it + if got := portFIPs(t, mock, "port-1"); len(got) != 1 { + t.Fatalf("setup: the live address is not attached (%d floating ips on port-1)", len(got)) + } + atomic.StoreInt32(&cos.disassociations, 0) + + start := time.Now() + if _, err := o.ClearQueue(ctx); err != nil { + t.Fatalf("clear queue: %v", err) + } + // One call for the live address; the port sweep finds nothing left. + if n := atomic.LoadInt32(&cos.disassociations); n != 1 { + t.Fatalf("clear queue made %d disassociate calls, want 1 (finished addresses must be skipped)", n) + } + if got := portFIPs(t, mock, "port-1"); len(got) != 0 { + t.Fatalf("the live floating ip is still attached: %+v", got) + } + if left, _ := d.ListIPs(ctx); len(left) != 0 { + t.Fatalf("%d rows left after clear", len(left)) + } + t.Logf("clear of %d rows took %s", finished+1, time.Since(start)) +} + +// A bulk detach runs in parallel but never above maxParallelDetach at once. +func TestDisassociateAllIsBoundedAndParallel(t *testing.T) { + o, _, mock := newTestOrchestrator(t, 180) + g := &gaugeOS{MockClient: mock, delay: 30 * time.Millisecond} + o.OS = g + var refs []db.FIPRef + for i := 0; i < 40; i++ { + id := fmt.Sprintf("fip-%d", i) + mock.Seed(id, fmt.Sprintf("10.8.0.%d", i+1), "svc") + refs = append(refs, db.FIPRef{IPID: int64(i), IPAddress: id, FIPID: id}) + } + start := time.Now() + o.disassociateAll(context.Background(), refs, "test") + elapsed := time.Since(start) + if peak := atomic.LoadInt32(&g.peak); peak < 2 || peak > maxParallelDetach { + t.Fatalf("peak concurrency %d, want between 2 and %d", peak, maxParallelDetach) + } + if elapsed > 600*time.Millisecond { // 40 x 30 ms sequentially is 1.2 s + t.Fatalf("took %s: the detach is not parallel", elapsed) + } +} + +type gaugeOS struct { + *openstack.MockClient + delay time.Duration + cur, peak int32 +} + +func (g *gaugeOS) DisassociateFloatingIP(ctx context.Context, fipID string) error { + n := atomic.AddInt32(&g.cur, 1) + for { + p := atomic.LoadInt32(&g.peak) + if n <= p || atomic.CompareAndSwapInt32(&g.peak, p, n) { + break + } + } + time.Sleep(g.delay) + atomic.AddInt32(&g.cur, -1) + return g.MockClient.DisassociateFloatingIP(ctx, fipID) +} + +// Rows that disagree with the queue are repaired by ReconcileValidators (run on every Tick). +func TestReconcileRepairsInconsistentValidators(t *testing.T) { + ctx := context.Background() + _, d, _ := newTestOrchestrator(t, 180) + _ = d.RegisterValidator(ctx, "validator-1", "host", "port-1", "v") + _ = d.RegisterValidator(ctx, "validator-2", "host", "port-2", "v") + _ = d.SeedQueue(ctx, []string{"1.1.1.1", "2.2.2.2"}) + + finished, _ := d.GetIPByAddress(ctx, "1.1.1.1") + foreign, _ := d.GetIPByAddress(ctx, "2.2.2.2") + // validator-1 still points at a finished address; validator-2 at an + // address owned by somebody else (what an older version could leave behind). + for _, q := range []string{ + fmt.Sprintf(`UPDATE ip_queue SET state='done' WHERE id=%d`, finished.ID), + fmt.Sprintf(`UPDATE ip_queue SET state='checking', owner_validator_id='validator-1' WHERE id=%d`, foreign.ID), + fmt.Sprintf(`UPDATE validators SET state='assigned', current_ip_id=%d WHERE validator_id='validator-1'`, finished.ID), + fmt.Sprintf(`UPDATE validators SET state='assigned', current_ip_id=%d WHERE validator_id='validator-2'`, foreign.ID), + } { + if _, err := d.ExecContext(ctx, q); err != nil { + t.Fatal(err) + } + } + n, err := d.ReconcileValidators(ctx) + if err != nil || n != 2 { + t.Fatalf("reconcile repaired %d validators (err %v), want 2", n, err) + } + for _, id := range []string{"validator-1", "validator-2"} { + v, _ := d.GetValidator(ctx, id) + if v.State != db.ValidatorIdle || v.CurrentIPID != nil { + t.Fatalf("%s: state=%s current_ip=%v, want idle", id, v.State, v.CurrentIPID) + } + } + // A consistent validator is left alone. + if n, _ := d.ReconcileValidators(ctx); n != 0 { + t.Fatalf("second reconcile repaired %d, want 0", n) + } +} + +// Randomised run: many validators, random silences and association failures. +// After every tick each validator holds at most one address, the validator +// and the address agree on who holds what, and no lease is ever reclaimed. +func TestRandomFlowKeepsOneAddressPerValidator(t *testing.T) { + ctx := context.Background() + rng := rand.New(rand.NewSource(42)) + o, d, mock := newTestOrchestrator(t, 180) + + const validators, addresses = 12, 60 + var addrs []string + for i := 0; i < validators; i++ { + _ = d.RegisterValidator(ctx, fmt.Sprintf("validator-%d", i+1), "host", fmt.Sprintf("port-%d", i+1), "v") + } + for i := 0; i < addresses; i++ { + a := fmt.Sprintf("10.7.%d.%d", i/200, i%200+1) + addrs = append(addrs, a) + mock.Seed(fmt.Sprintf("fip-%d", i), a, "svc") + } + if err := d.SeedQueue(ctx, addrs); err != nil { + t.Fatal(err) + } + + checkInvariants := func(tick int) { + t.Helper() + vs, _ := d.ListValidators(ctx) + holders := map[int64]string{} + for _, v := range vs { + if v.CurrentIPID == nil { + if v.State == db.ValidatorAssigned { + t.Fatalf("tick %d: %s is assigned but holds nothing", tick, v.ValidatorID) + } + continue + } + ip, err := d.GetIP(ctx, *v.CurrentIPID) + if err != nil { + t.Fatalf("tick %d: %s points at a missing address %d", tick, v.ValidatorID, *v.CurrentIPID) + } + if ip.OwnerValidatorID == nil || *ip.OwnerValidatorID != v.ValidatorID { + t.Fatalf("tick %d: %s holds %s but its owner is %v", tick, v.ValidatorID, ip.IPAddress, ip.OwnerValidatorID) + } + if ip.State == db.IPDone || ip.State == db.IPFailed || ip.State == db.IPOccupied { + t.Fatalf("tick %d: %s still holds the finished address %s", tick, v.ValidatorID, ip.IPAddress) + } + if prev, ok := holders[ip.ID]; ok { + t.Fatalf("tick %d: %s is held by both %s and %s", tick, ip.IPAddress, prev, v.ValidatorID) + } + holders[ip.ID] = v.ValidatorID + } + // Every owned, unfinished address is the current one of its owner. + ips, _ := d.ListIPs(ctx) + perOwner := map[string]int{} + for _, ip := range ips { + if ip.OwnerValidatorID != nil && ip.State != db.IPDone && ip.State != db.IPFailed && ip.State != db.IPOccupied { + perOwner[*ip.OwnerValidatorID]++ + } + } + for owner, n := range perOwner { + if n > 1 { + t.Fatalf("tick %d: %s owns %d unfinished addresses", tick, owner, n) + } + } + } + + checkTicks := map[int64]int{} // address id -> ticks spent in checking + done := false + for tick := 1; tick <= 600 && !done; tick++ { + // Occasionally an association fails (409 on a port). + if rng.Intn(15) == 0 { + mock.AssociateFailures = map[string]error{fmt.Sprintf("fip-%d", rng.Intn(addresses)): fmt.Errorf("409 conflict")} + } + o.Tick(ctx) + + vs, _ := d.ListValidators(ctx) + for _, v := range vs { + // A busy agent sometimes goes silent, then speaks again. + if v.CurrentIPID != nil && v.State == db.ValidatorAssigned && rng.Intn(10) == 0 { + _ = d.MarkValidatorUnreachable(ctx, v.ValidatorID) + } + if v.State == db.ValidatorUnreachable && rng.Intn(3) == 0 { + _ = d.Heartbeat(ctx, v.ValidatorID) + } + item, _, err := o.AssignmentForValidator(ctx, v.ValidatorID) + if err != nil || item == nil { + continue + } + if item.State == db.IPAwaitingSelfCheck { + _ = o.SelfCheckResult(ctx, v.ValidatorID, item.ID, true, "ok") + continue + } + checkTicks[item.ID]++ + if checkTicks[item.ID] >= 1+rng.Intn(4) { + finishChecks(t, o, item, v.ValidatorID) + delete(checkTicks, item.ID) + } + } + checkInvariants(tick) + + counts, _, _ := d.CountIPsByState(ctx) + unfinished := 0 + for st, n := range counts { + if st != db.IPDone && st != db.IPFailed && st != db.IPOccupied { + unfinished += n + } + } + done = unfinished == 0 + } + if !done { + counts, _, _ := d.CountIPsByState(ctx) + t.Fatalf("not all addresses finished: %v", counts) + } + var expired int + if err := d.QueryRowContext(ctx, `SELECT COUNT(*) FROM events WHERE event_type='lease_expired'`).Scan(&expired); err != nil { + t.Fatal(err) + } + if expired != 0 { + t.Fatalf("%d leases were reclaimed, want 0", expired) + } +}