diff --git a/bin/SHA256SUMS b/bin/SHA256SUMS index 8506586..8734577 100644 --- a/bin/SHA256SUMS +++ b/bin/SHA256SUMS @@ -1,4 +1,4 @@ -5703aa9367f814fd1c02ab8ee821f4ebea28934387b084e0f843efe1505530b7 control-api +118e09f98944ac59593af603c04a1d5f0e7cafb6549b5a562300bc31d9c769da control-api 48c9b99fa88be751d9badfba8b7d80f894e326a743a90e85478f1d2251ce4ab6 validator-agent -43fb660b78179b204b57388241a508345881aa16f3531d77b8fd75af60e7f2d8 prober -ef796e473d647a173954d40c646904898d16bf2c831b3f87bea3e1e5b809888b admin-dashboard +8877cf7c6d2bdd58faf70865be602b3db5f18508f4607c8942bc8899ce0fb1e6 prober +73bdf0ccda4912f25710a09bd37eaa0620f1b6e5002cb50e6b8b057cecb17f97 admin-dashboard diff --git a/bin/admin-dashboard b/bin/admin-dashboard index 73b7350..8b34732 100755 Binary files a/bin/admin-dashboard and b/bin/admin-dashboard differ diff --git a/bin/control-api b/bin/control-api index a71bd5f..eb03d87 100755 Binary files a/bin/control-api and b/bin/control-api differ diff --git a/bin/prober b/bin/prober index 5a6d5b3..c52c44c 100755 Binary files a/bin/prober and b/bin/prober differ diff --git a/cmd/control-api/main.go b/cmd/control-api/main.go index a1d586e..ac6a3f2 100644 --- a/cmd/control-api/main.go +++ b/cmd/control-api/main.go @@ -102,6 +102,9 @@ func runOrchestratorLoop(ctx context.Context, orch *orchestrator.Orchestrator, c if err := orch.SweepStaleHeartbeats(ctx); err != nil { log.Error("sweep stale heartbeats", "err", err) } + if err := orch.SweepStaleSiteHeartbeats(ctx); err != nil { + log.Error("sweep stale site heartbeats", "err", err) + } } } } diff --git a/configs/control-api.example.yaml b/configs/control-api.example.yaml index 2e3a4b3..d6c6669 100644 --- a/configs/control-api.example.yaml +++ b/configs/control-api.example.yaml @@ -67,11 +67,11 @@ validators: os_port_id: "REPLACE_WITH_NEUTRON_PORT_ID_2" # Inbound (prober) checks are OPTIONAL: list here only the external sites -# you actually run a `prober` on, up to 3 slots (index 1-3, matches -# ip_queue.siteN_complete — the schema is capped at 3). Leave this list -# empty to disable inbound checks entirely — the overall result is then -# based on egress checks alone, and aggregation doesn't wait on any -# prober. A partial list (e.g. just index 1) only waits on that one site. +# you actually run a `prober` on — index is any integer >= 1, no cap on +# how many slots you configure. Leave this list empty to disable inbound +# checks entirely — the overall result is then based on egress checks +# alone, and aggregation doesn't wait on any prober. A partial list (e.g. +# just index 1) only waits on that one site. sites: - site_id: "site-1" index: 1 diff --git a/docs/API.md b/docs/API.md index 2bcd190..95a3e0c 100644 --- a/docs/API.md +++ b/docs/API.md @@ -209,7 +209,12 @@ target)` в рамках текущей попытки — безопасна и ### `POST /api/v1/probers/register` Регистрация пробера. `site_id` должен присутствовать в конфиге control-api -(`sites[].site_id`), иначе — `400`. +(`sites[].site_id`), иначе — `400`. Помимо привычной проверки, запоминает +`hostname` и переводит площадку в состояние `idle` (если она была +`unregistered`/`unreachable`) — то же самое, что делает `validator-agent` +при своей регистрации (см. +[«Управление площадками»](USAGE.md#управление-площадками-проберами) в +USAGE.md про состояния площадки). Запрос: ```json @@ -218,6 +223,17 @@ target)` в рамках текущей попытки — безопасна и Ответ: `{"ok": true, "poll_interval_seconds": 5}`. +### `POST /api/v1/probers/{site_id}/heartbeat` + +«Я жив». Обновляет `last_heartbeat_at` площадки. Если площадка не +присылает heartbeat дольше `orchestrator.heartbeat_timeout_seconds` +(тот же параметр, что и для валидаторов), она помечается `unreachable`. +Полный аналог `POST /api/v1/agents/{id}/heartbeat` для пробера. + +Запрос: тело не требуется. + +Ответ: `{"ok": true}`. `404`, если `site_id` не сконфигурирован. + ### `GET /api/v1/probers/{site_id}/assignments` Список всех IP, которые сейчас находятся в состоянии `checking` — то есть @@ -424,15 +440,23 @@ curl -s -X POST "$BASE/api/v1/admin/ips/clear" ### Площадки: `/api/v1/admin/config/sites` -Слотов ровно три (`index` ∈ {1, 2, 3}) — это ограничение схемы БД -(`ip_queue.site{1,2,3}_complete`), а не искусственное. Пустой список слотов -— штатный сценарий, отключающий inbound-проверки целиком (см. +Число слотов не ограничено — `index` может быть любым целым `>= 1`, +столько площадок, сколько нужно оператору. Пустой список слотов — штатный +сценарий, отключающий inbound-проверки целиком (см. [USAGE.md](USAGE.md#управление-площадками-проберами)). +Помимо `index`/`site_id`, объект площадки несёт состояние подключения +пробера — `hostname`/`state`/`last_heartbeat_at`, тот же смысл, что у +аналогичных полей валидатора (`unregistered`/`idle`/`unreachable`, +обновляются через `POST /api/v1/probers/register` и `POST +/api/v1/probers/{site_id}/heartbeat`, см. выше). `PUT` на слот всегда +сбрасывает эти три поля к значениям «ещё не подключался» — новый (или +даже тот же) `site_id` трактуется как новая идентичность пробера. + | Метод | Путь | Тело | Успех | Ошибки | |---|---|---|---|---| -| GET | `/api/v1/admin/config/sites` | — | `[{"index","site_id"}]` (до 3 строк) | | -| PUT | `/api/v1/admin/config/sites/{index}` | `{"site_id"}` | `200` | `400`, если `index` не 1..3; `409`, если `site_id` уже занят другим слотом | +| GET | `/api/v1/admin/config/sites` | — | `[{"index","site_id","hostname","state","last_heartbeat_at"}]` | | +| PUT | `/api/v1/admin/config/sites/{index}` | `{"site_id"}` | `200` | `400`, если `index < 1`; `409`, если `site_id` уже занят другим слотом | | DELETE | `/api/v1/admin/config/sites/{index}` | — | `200` | `404` | ### Группы целей: `/api/v1/admin/config/targets` @@ -570,8 +594,8 @@ queued ──(control-api сам, без вызова API)──▶ assigning_fi в БД при этом уже `awaiting_self_check`. Площадки (`siteN_complete`) — опциональны: сколько их учитывается, -целиком определяется текущим списком `sites` (0–3 записи, управляется -через `/api/v1/admin/config/sites` — см. +целиком определяется текущим списком `sites` (без ограничения по числу +записей, управляется через `/api/v1/admin/config/sites` — см. [выше](#управление-очередью-и-конфигурацией)). Пустой список — агрегация ждёт только `egress_complete`, ни одна площадка не требуется. Подробнее — [USAGE.md](USAGE.md#управление-площадками-проберами). diff --git a/docs/DASHBOARD.md b/docs/DASHBOARD.md index 12bd1cc..01d499a 100644 --- a/docs/DASHBOARD.md +++ b/docs/DASHBOARD.md @@ -52,7 +52,7 @@ admin-dashboard -config /etc/cloud-ip-validator/admin-dashboard.yaml | `/ips` | Полная очередь. Форма сверху принимает список адресов (по одному на строке или через запятую) и отправляет их в `POST /api/v1/admin/ips` — **один и тот же вызов** добавляет новые адреса и принудительно перезапускает уже завершённые (см. ниже). У каждого адреса — кнопка «Перепроверить» (для `done`/`failed`) или «Отменить» (для активных состояний), и всегда — «Удалить» (безвозвратно, в отличие от «Отменить», см. ниже). Чекбоксы у строк + кнопка «Удалить выбранные» удаляют список одним вызовом; «Очистить всё» удаляет вообще всё, включая активные проверки — обе операции требуют явного подтверждения. Пока не истекла настроенная на `/settings` пауза (`fip_settle_seconds`), только что привязавший Floating IP адрес показывает отдельный бейдж «прогрев FIP» вместо обычного статуса. | | `/ips/{ip}` | Детали одного адреса: все проверки текущей попытки и вся история событий. | | `/validators` | Список валидаторов + создание/изменение `os_port_id`/удаление. | -| `/sites` | Три фиксированных слота площадок (1/2/3) — назначить/сменить/освободить `site_id`. | +| `/sites` | Площадки — число слотов не ограничено, форма сверху добавляет новый слот, назначить/сменить/освободить `site_id` в каждой строке; колонка «Статус» показывает бейдж подключения пробера (`unregistered`/`idle`/`unreachable`, по аналогии с `/validators`), см. [USAGE.md](USAGE.md#состояния-площадки). | | `/targets` | Группы целей для egress-проверок — создание/редактирование/удаление. | | `/check-types` | Типы проверок (`https`/`icmp`/`ssh`/...), включение/выключение, привязка к группам целей. | | `/settings` | Две формы: `fip_settle_seconds` — пауза (в секундах) между привязкой Floating IP и началом self-check («прогрев» дата-плейна OpenStack, см. [USAGE.md](USAGE.md#пауза-перед-self-check-fip_settle_seconds)); и типы проверок пробера — TCP-порты (через запятую) + чекбокс ICMP, общие для всех площадок (см. [USAGE.md](USAGE.md#управление-типами-проверок-пробера)). | diff --git a/docs/DIAGRAMS.md b/docs/DIAGRAMS.md index c699dca..e72bcaf 100644 --- a/docs/DIAGRAMS.md +++ b/docs/DIAGRAMS.md @@ -222,9 +222,9 @@ flowchart TB проверки — подробнее о значениях полей см. [USAGE.md](USAGE.md#значения-полей-ip). -> На диаграмме показан полный вариант с тремя площадками — это не -> обязательный минимум. Площадки опциональны: сколько их ожидать (0–3), -> определяется списком `sites` в конфиге control-api. При пустом списке -> четвёртый поток телеметрии (`inbound-site-N`) просто отсутствует, и -> агрегация ждёт только `egress_complete`. См. +> На диаграмме показан пример с тремя площадками — это лишь иллюстрация, +> не ограничение. Площадки опциональны и их число не ограничено: сколько +> их ожидать, определяется списком `sites` в конфиге control-api (пустой +> список — ни одного потока `inbound-site-N`, агрегация ждёт только +> `egress_complete`; N площадок — N параллельных потоков телеметрии). См. > [USAGE.md](USAGE.md#управление-площадками-проберами). diff --git a/docs/SETUP.md b/docs/SETUP.md index 86ec27a..898369a 100644 --- a/docs/SETUP.md +++ b/docs/SETUP.md @@ -176,9 +176,9 @@ cp configs/control-api.example.yaml /etc/cloud-ip-validator/control-api.yaml соответствующего `validator-agent`) и `os_port_id` — **ID Neutron-порта** основного сетевого интерфейса ВМ-валидатора (узнать: `openstack port list --server <имя-ВМ>` или в веб-консоли облака). -- **`sites`** — три внешние площадки, `site_id` + `index` (1, 2 или 3). - `site_id` должен совпадать с `site_id` в конфиге соответствующего - `prober`. +- **`sites`** — внешние площадки, `site_id` + `index` (число слотов не + ограничено; в примере ниже — три). `site_id` должен совпадать с + `site_id` в конфиге соответствующего `prober`. - **`ip_addresses`** — список публичных IPv4-адресов на проверку, **в порядке обработки**. Адреса должны существовать в сервисном проекте как уже выделенные (allocated) floating IP — инструмент их не создаёт. diff --git a/docs/USAGE.md b/docs/USAGE.md index 2a046c7..4ac16e8 100644 --- a/docs/USAGE.md +++ b/docs/USAGE.md @@ -123,10 +123,14 @@ curl -s http://:8080/api/v1/admin/ips \ | `AttemptNumber` | Номер попытки — растёт при каждом requeue (сбой привязки, сбой self-check, реклейм по таймауту) | | `RetryCount` | Сколько раз адрес уже переставлялся в очередь заново | | `EgressComplete` | Валидатор закончил исходящие проверки | -| `Site1Complete` / `Site2Complete` / `Site3Complete` | Соответствующая площадка закончила входящие проверки | | `OverallResult` | Итог: `pass`, `partial`, `fail`, `cancelled` (принудительно остановлена, см. [«Принудительная остановка проверки»](#принудительная-остановка-проверки)), либо пусто, пока проверка не завершена | | `AssignedAt` / `AggregatedAt` / `FIPReleasedAt` | Метки времени соответствующих этапов | +> Завершённость входящих проверок по каждой конкретной площадке в этом +> списке не отображается (площадок теперь может быть сколько угодно, а не +> фиксированные три) — детали по конкретной площадке смотрите в массиве +> `checks` ответа `GET /api/v1/admin/ips/{ip}` (`Source: "inbound-site-N"`). + ## Как читать итоговый результат (pass/partial/fail) - **`pass`** — прошли все проверки: все исходящие, плюс все площадки по @@ -237,18 +241,17 @@ curl -s -X DELETE http://:8080/api/v1/admin/config/validators/valid ## Управление площадками (проберами) -**Входящие (inbound/prober) проверки полностью опциональны.** Список -`sites` в конфиге control-api и есть переключатель: пусто — inbound- -проверки выключены целиком, итоговый результат считается только по -исходящим (egress) проверкам, и агрегация не ждёт вообще ни одного -пробера. Указан один или два слота — ждём только их, остальные не -учитываются. Указаны все три — работает как в исходной схеме процесса. +**Входящие (inbound/prober) проверки полностью опциональны, а число +площадок не ограничено.** Список `sites` в конфиге control-api и есть +переключатель: пусто — inbound-проверки выключены целиком, итоговый +результат считается только по исходящим (egress) проверкам, и агрегация +не ждёт вообще ни одного пробера. Указано N слотов — ждём именно их. Явного отдельного флага "включить/выключить" нет — самого списка `sites` достаточно. Чтобы добавить площадку (без перезапуска control-api) — назначьте -`site_id` на один из трёх слотов (`index` 1, 2 или 3 — см. ограничение -ниже) через API: +`site_id` любому свободному слоту (`index` — любое целое `>= 1`, слотов +может быть сколько угодно) через API: ```bash curl -s -X PUT http://:8080/api/v1/admin/config/sites/1 \ -d '{"site_id": "site-1"}' @@ -269,12 +272,26 @@ curl -s -X DELETE http://:8080/api/v1/admin/config/sites/1 > игнорируется при всех последующих рестартах. Для стенда, который уже > хоть раз запускался, используйте API выше. -> Важно: количество *возможных* слотов площадок жёстко зашито в схему БД -> (`Site1Complete`/`Site2Complete`/`Site3Complete`) — не более **трёх**, -> как и описано в исходной схеме процесса. `index` может быть только 1, 2 -> или 3. Использовать *меньше* трёх (в том числе ноль) — штатный, -> поддерживаемый сценарий; *больше* трёх потребует доработки схемы -> данных, одной правкой конфига не обойтись. +### Состояния площадки + +По аналогии с валидаторами (см. [«Управление +валидаторами»](#управление-валидаторами) выше), у каждой площадки есть +состояние подключения — видно и в `GET /api/v1/admin/config/sites` +(поля `hostname`/`state`/`last_heartbeat_at`), и на странице `/sites` +дашборда бейджем: +- **`unregistered`** — слот сконфигурирован, но процесс `prober` на этой + площадке ещё ни разу не подключался (не вызывал `POST + /api/v1/probers/register`). +- **`idle`** — площадка на связи: `prober` зарегистрирован и присылает + heartbeat (`POST /api/v1/probers/{site_id}/heartbeat`) на каждом опросе. +- **`unreachable`** — площадка пропустила heartbeat дольше + `orchestrator.heartbeat_timeout_seconds` (тот же параметр, что и для + валидаторов) — вероятно, процесс `prober` упал или потерял сеть до + control-api. + +В отличие от валидатора, у площадки нет состояний `assigned`/`checking` +— пробер не привязан к одному IP, а на каждом опросе обрабатывает сразу +весь активный набор. ## Управление типами проверок пробера @@ -522,13 +539,17 @@ https://api.ipify.org`) и логи `journalctl -u validator-agent` на пре сети) покажет приватный адрес валидатора независимо от того, правильно ли привязан FIP, и всегда будет давать ложный провал. -**Площадка (`site-N`) никогда не отчитывается (`SiteNComplete` всегда -`false`).** -Проверьте, что `prober` на этой площадке запущен и его `site_id` в -конфиге совпадает с `site_id` в конфиге control-api. Проверьте, что -площадка имеет сетевой доступ и до `control-api`, и до проверяемого -адреса (входящий трафик на 22/80/443/8080 + ICMP — это отдельная -связность от связи с control-api, см. +**Площадка (`site-N`) никогда не отчитывается по конкретному IP.** +Сперва проверьте статус самой площадки — `GET +/api/v1/admin/config/sites`, поле `state`. `unregistered` или +`unreachable` означает, что `prober` на этой площадке вообще не на связи +с control-api (см. [«Состояния площадки»](#состояния-площадки) выше) — +проверьте, что процесс запущен и его `site_id` в конфиге совпадает с +`site_id` в конфиге control-api, и что площадка имеет сетевой доступ до +`control-api`. Если статус `idle` (площадка на связи), а конкретный IP +всё равно не получает отметку о завершении — проверьте отдельно +связность до проверяемого адреса (входящий трафик на 22/80/443/8080 + +ICMP — это отдельная связность от связи с control-api, см. [SETUP.md](SETUP.md#сетевые-доступы)). **Много адресов зависло в `checking` дольше ожидаемого.** diff --git a/internal/config/config.go b/internal/config/config.go index 5fa4363..037866d 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -91,7 +91,7 @@ type ValidatorConfig struct { type SiteConfig struct { SiteID string `yaml:"site_id"` - Index int `yaml:"index"` // 1, 2, or 3 — maps to ip_queue.siteN_complete + Index int `yaml:"index"` // any int >= 1, no cap on the number of slots } // CheckTypeConfig maps a check type (https, icmp, ssh) to the named target diff --git a/internal/dashboard/dashboard_test.go b/internal/dashboard/dashboard_test.go index c18440b..0b1dd06 100644 --- a/internal/dashboard/dashboard_test.go +++ b/internal/dashboard/dashboard_test.go @@ -26,6 +26,8 @@ type fakeControlAPI struct { ips []ipQueueItem validators []validatorDTO sites map[int]string + siteState map[int]string + siteHostname map[int]string groups map[string][]string checkTypes map[string]checkTypeDTO fipSettleSeconds int @@ -36,9 +38,11 @@ type fakeControlAPI struct { func newFakeControlAPI(t *testing.T) (*fakeControlAPI, string) { t.Helper() f := &fakeControlAPI{ - sites: map[int]string{}, - groups: map[string][]string{}, - checkTypes: map[string]checkTypeDTO{}, + sites: map[int]string{}, + siteState: map[int]string{}, + siteHostname: map[int]string{}, + groups: map[string][]string{}, + checkTypes: map[string]checkTypeDTO{}, } ts := httptest.NewServer(f.handler()) t.Cleanup(ts.Close) @@ -285,7 +289,7 @@ func (f *fakeControlAPI) handler() http.Handler { defer f.mu.Unlock() var out []siteDTO for idx, id := range f.sites { - out = append(out, siteDTO{Index: idx, SiteID: id}) + out = append(out, siteDTO{Index: idx, SiteID: id, Hostname: f.siteHostname[idx], State: f.siteState[idx]}) } writeJSON(w, http.StatusOK, out) }) @@ -305,7 +309,11 @@ func (f *fakeControlAPI) handler() http.Handler { } } f.sites[idx] = req.SiteID - writeJSON(w, http.StatusOK, siteDTO{Index: idx, SiteID: req.SiteID}) + // UpsertSite always resets connection status on write (see its + // control-api doc comment) — mirror that here. + f.siteState[idx] = "unregistered" + f.siteHostname[idx] = "" + writeJSON(w, http.StatusOK, siteDTO{Index: idx, SiteID: req.SiteID, State: "unregistered"}) }) mux.HandleFunc("DELETE /api/v1/admin/config/sites/{index}", func(w http.ResponseWriter, r *http.Request) { f.mu.Lock() @@ -313,6 +321,8 @@ func (f *fakeControlAPI) handler() http.Handler { var idx int fmt.Sscanf(r.PathValue("index"), "%d", &idx) delete(f.sites, idx) + delete(f.siteState, idx) + delete(f.siteHostname, idx) writeJSON(w, http.StatusOK, map[string]bool{"ok": true}) }) diff --git a/internal/dashboard/dto.go b/internal/dashboard/dto.go index fe6ba1b..5a17ef2 100644 --- a/internal/dashboard/dto.go +++ b/internal/dashboard/dto.go @@ -30,9 +30,6 @@ type ipQueueItem struct { RetryCount int LeaseExpiresAt *time.Time EgressComplete bool - Site1Complete bool - Site2Complete bool - Site3Complete bool OverallResult string AssignedAt *time.Time FIPAssociatedAt *time.Time @@ -106,8 +103,11 @@ type validatorDTO struct { } type siteDTO struct { - Index int `json:"index"` - SiteID string `json:"site_id"` + Index int `json:"index"` + SiteID string `json:"site_id"` + Hostname string `json:"hostname"` + State string `json:"state"` + LastHeartbeatAt *time.Time `json:"last_heartbeat_at"` } type targetGroupDTO struct { diff --git a/internal/dashboard/handlers_sites.go b/internal/dashboard/handlers_sites.go index 72dac28..14e4974 100644 --- a/internal/dashboard/handlers_sites.go +++ b/internal/dashboard/handlers_sites.go @@ -8,26 +8,12 @@ import ( type sitesPageData struct { PageData - Items []siteDTO // always exactly 3 entries, index 1..3, SiteID=="" for an empty slot -} - -// fillSlots pads control-api's response (which only lists assigned slots) -// out to all three fixed slots, so the table always renders 3 rows. -func fillSlots(sites []siteDTO) []siteDTO { - byIndex := make(map[int]string, len(sites)) - for _, s := range sites { - byIndex[s.Index] = s.SiteID - } - out := make([]siteDTO, 3) - for i := 0; i < 3; i++ { - out[i] = siteDTO{Index: i + 1, SiteID: byIndex[i+1]} - } - return out + Items []siteDTO // dynamic length — as many slots as configured, no fixed cap } func (s *Server) handleSitesPage(w http.ResponseWriter, r *http.Request) { items, err := s.CA.ListSites(r.Context()) - data := sitesPageData{Items: fillSlots(items)} + data := sitesPageData{Items: items} data.ActiveNav = "sites" data.Banner = bannerFor(err) s.renderPage(w, "sites_page", data) @@ -38,7 +24,26 @@ func (s *Server) renderSitesTable(w http.ResponseWriter, r *http.Request, action if actionErr == nil { actionErr = listErr } - s.renderFragment(w, "sites_table", sitesPageData{Items: fillSlots(items)}, actionErr) + s.renderFragment(w, "sites_table", sitesPageData{Items: items}, actionErr) +} + +// handleSiteCreate adds a brand-new slot via the "add site" form — same +// create-alongside-update pattern as /targets' handleTargetCreate, just +// backed by the same PutSite call the per-row form already uses (there's +// no separate "create" endpoint on control-api since PUT is already an +// upsert). +func (s *Server) handleSiteCreate(w http.ResponseWriter, r *http.Request) { + if err := r.ParseForm(); err != nil { + s.renderSitesTable(w, r, fmt.Errorf("invalid form: %w", err)) + return + } + index, convErr := strconv.Atoi(r.PostFormValue("index")) + if convErr != nil { + s.renderSitesTable(w, r, &apiErr{Status: http.StatusBadRequest, Message: "index должен быть числом"}) + return + } + err := s.CA.PutSite(r.Context(), index, r.PostFormValue("site_id")) + s.renderSitesTable(w, r, err) } func (s *Server) handleSitePut(w http.ResponseWriter, r *http.Request) { diff --git a/internal/dashboard/handlers_test.go b/internal/dashboard/handlers_test.go index 15704ef..60f292b 100644 --- a/internal/dashboard/handlers_test.go +++ b/internal/dashboard/handlers_test.go @@ -316,11 +316,11 @@ func TestSitesPutAndConflict(t *testing.T) { t.Fatalf("expected site-1 assigned to slot 1, got:\n%s", body) } - // Slot page always renders exactly 3 rows, including empty ones — count - // the per-slot PUT forms rather than raw (the thead row has one too). + // The table renders exactly as many rows as configured slots — no + // fixed cap, no padding to a fixed count. page := get(t, ts, "/sites") - if strings.Count(page, `hx-put="/sites/`) != 3 { - t.Fatalf("expected exactly 3 site slot rows, got:\n%s", page) + if strings.Count(page, `hx-put="/sites/`) != 1 { + t.Fatalf("expected exactly 1 site slot row, got:\n%s", page) } body = postForm(t, ts, "PUT", "/sites/2", map[string][]string{"site_id": {"site-1"}}) @@ -329,6 +329,53 @@ func TestSitesPutAndConflict(t *testing.T) { } } +// TestSiteCreateAddsSlotBeyondThree proves there's no hardcoded cap on the +// number of prober slots — the "add site" form can create slot 4 and +// beyond. +func TestSiteCreateAddsSlotBeyondThree(t *testing.T) { + _, caURL := newFakeControlAPI(t) + ts := newTestServer(t, caURL) + + postForm(t, ts, "PUT", "/sites/1", map[string][]string{"site_id": {"site-1"}}) + postForm(t, ts, "PUT", "/sites/2", map[string][]string{"site_id": {"site-2"}}) + postForm(t, ts, "PUT", "/sites/3", map[string][]string{"site_id": {"site-3"}}) + + body := postForm(t, ts, "POST", "/sites", map[string][]string{"index": {"4"}, "site_id": {"site-4"}}) + if !strings.Contains(body, "site-4") { + t.Fatalf("expected slot 4 created, got:\n%s", body) + } + if strings.Count(body, `hx-put="/sites/`) != 4 { + t.Fatalf("expected exactly 4 site slot rows, got:\n%s", body) + } + + // A malformed index is caught by the dashboard itself. + body = postForm(t, ts, "POST", "/sites", map[string][]string{"index": {"not-a-number"}, "site_id": {"site-5"}}) + if !strings.Contains(body, "alert-warning") { + t.Fatalf("expected client error banner for non-numeric index, got:\n%s", body) + } +} + +func TestSitesPageShowsProberBadge(t *testing.T) { + fake, caURL := newFakeControlAPI(t) + fake.sites[1] = "site-1" + fake.siteState[1] = "idle" + fake.siteHostname[1] = "probe-host-1" + fake.sites[2] = "site-2" + fake.siteState[2] = "unreachable" + ts := newTestServer(t, caURL) + + page := get(t, ts, "/sites") + if !strings.Contains(page, "probe-host-1") { + t.Fatalf("expected hostname shown, got:\n%s", page) + } + if !strings.Contains(page, "pill-success") || !strings.Contains(page, ">idle<") { + t.Fatalf("expected idle badge for site-1, got:\n%s", page) + } + if !strings.Contains(page, "pill-danger") || !strings.Contains(page, ">unreachable<") { + t.Fatalf("expected unreachable badge for site-2, got:\n%s", page) + } +} + func TestTargetsAndCheckTypesRoundTrip(t *testing.T) { _, caURL := newFakeControlAPI(t) ts := newTestServer(t, caURL) diff --git a/internal/dashboard/render.go b/internal/dashboard/render.go index bd1f764..1de9d5b 100644 --- a/internal/dashboard/render.go +++ b/internal/dashboard/render.go @@ -101,10 +101,14 @@ func joinInts(ints []int, sep string) string { var funcMap = template.FuncMap{ "ipBadge": ipBadge, "validatorBadge": validatorBadge, - "fmtTime": fmtTime, - "deref": derefStr, - "join": strings.Join, - "joinInts": joinInts, + // proberBadge reuses validatorBadge as-is: prober states + // (unregistered/idle/unreachable) are a subset of validator states + // with identical labels/styles, so no separate function is needed. + "proberBadge": validatorBadge, + "fmtTime": fmtTime, + "deref": derefStr, + "join": strings.Join, + "joinInts": joinInts, } func parseTemplates() (*template.Template, error) { diff --git a/internal/dashboard/routes.go b/internal/dashboard/routes.go index a0a3850..d9f22a6 100644 --- a/internal/dashboard/routes.go +++ b/internal/dashboard/routes.go @@ -23,6 +23,7 @@ func (s *Server) routes(mux *http.ServeMux) { mux.HandleFunc("DELETE /validators/{id}", s.handleValidatorDelete) mux.HandleFunc("GET /sites", s.handleSitesPage) + mux.HandleFunc("POST /sites", s.handleSiteCreate) mux.HandleFunc("PUT /sites/{index}", s.handleSitePut) mux.HandleFunc("DELETE /sites/{index}", s.handleSiteDelete) diff --git a/internal/dashboard/templates/sites.html b/internal/dashboard/templates/sites.html index d0faf93..d4c3659 100644 --- a/internal/dashboard/templates/sites.html +++ b/internal/dashboard/templates/sites.html @@ -20,8 +20,27 @@ {{define "sites_content"}}

Площадки

-

Ровно три слота (1, 2, 3) — жёсткое ограничение схемы БД control-api. -Пустой слот — площадка не используется, входящие проверки для неё не ожидаются.

+

Внешние площадки (проберы), проверяющие адреса «снаружи». Число площадок +не ограничено — добавляйте столько слотов, сколько нужно. «Статус» показывает, подключён ли прямо сейчас +процесс prober на этой площадке (по heartbeat, аналогично валидаторам).

+ +
+
+
+
+
+ + +
+
+ + +
+ +
+
+
+
{{template "sites_table" .}} @@ -32,9 +51,10 @@
- + {{range .Items}} +{{$b := proberBadge .State}} + + + {{end}}
Слотsite_id
Слотsite_idХостСтатусHeartbeat
{{.Index}} @@ -43,17 +63,19 @@ {{if .Hostname}}{{.Hostname}}{{else}}—{{end}}{{$b.Label}}{{fmtTime .LastHeartbeatAt}}
-{{if .SiteID}} - -{{end}} +
+{{if not .Items}}

Площадок пока нет — добавьте первую формой выше.

{{end}}
{{end}} diff --git a/internal/db/db.go b/internal/db/db.go index da38304..d1d2e08 100644 --- a/internal/db/db.go +++ b/internal/db/db.go @@ -25,6 +25,12 @@ var fipSettleDelaySchema string //go:embed migrations/0004_inbound_checks_admin.sql var inboundChecksAdminSchema string +//go:embed migrations/0005_unbounded_sites.sql +var unboundedSitesSchema string + +//go:embed migrations/0006_prober_heartbeat.sql +var proberHeartbeatSchema string + // migrations is the ordered list of schema versions. Each entry's SQL is // applied, in order, for any version greater than the database's current // PRAGMA user_version — so a fresh database walks the whole list and an @@ -37,6 +43,8 @@ var migrations = []struct { {2, dynamicConfigSchema}, {3, fipSettleDelaySchema}, {4, inboundChecksAdminSchema}, + {5, unboundedSitesSchema}, + {6, proberHeartbeatSchema}, } type DB struct { diff --git a/internal/db/migrations/0005_unbounded_sites.sql b/internal/db/migrations/0005_unbounded_sites.sql new file mode 100644 index 0000000..27bb45c --- /dev/null +++ b/internal/db/migrations/0005_unbounded_sites.sql @@ -0,0 +1,19 @@ +-- Remove the hardcoded 3-site limit: ip_queue.site{1,2,3}_complete are +-- replaced by a child table with no cap on site_idx, keyed by +-- attempt_number the same way `checks` already is (see docs/USAGE.md). + +ALTER TABLE ip_queue DROP COLUMN site1_complete; +ALTER TABLE ip_queue DROP COLUMN site2_complete; +ALTER TABLE ip_queue DROP COLUMN site3_complete; + +-- Per-IP-per-site inbound-check completion. No fixed cap on site_idx, and +-- no reset needed on retry — a new attempt_number simply has no rows yet, +-- unlike the old fixed columns which had to be explicitly zeroed. +CREATE TABLE ip_site_checks ( + ip_id INTEGER NOT NULL REFERENCES ip_queue(id), + attempt_number INTEGER NOT NULL, + site_idx INTEGER NOT NULL, + complete BOOLEAN NOT NULL DEFAULT 0, + completed_at TIMESTAMP, + PRIMARY KEY (ip_id, attempt_number, site_idx) +); diff --git a/internal/db/migrations/0006_prober_heartbeat.sql b/internal/db/migrations/0006_prober_heartbeat.sql new file mode 100644 index 0000000..775a76b --- /dev/null +++ b/internal/db/migrations/0006_prober_heartbeat.sql @@ -0,0 +1,8 @@ +-- Prober availability indication (see docs/USAGE.md, "Управление +-- площадками"): sites gain the same hostname/state/heartbeat tracking +-- validators already have, so the dashboard can show whether a prober is +-- actually connected and polling. + +ALTER TABLE sites ADD COLUMN hostname TEXT NOT NULL DEFAULT ''; +ALTER TABLE sites ADD COLUMN state TEXT NOT NULL DEFAULT 'unregistered'; +ALTER TABLE sites ADD COLUMN last_heartbeat_at TIMESTAMP; diff --git a/internal/db/models.go b/internal/db/models.go index 58e8263..dd32acd 100644 --- a/internal/db/models.go +++ b/internal/db/models.go @@ -12,6 +12,12 @@ const ( ValidatorChecking = "checking" ValidatorUnreachable = "unreachable" + // Site (prober) availability states — no assigned/checking equivalent, + // since a prober isn't bound to one IP at a time. + SiteUnregistered = "unregistered" + SiteIdle = "idle" + SiteUnreachable = "unreachable" + IPQueued = "queued" IPAssigningFIP = "assigning_fip" IPAwaitingSelfCheck = "awaiting_self_check" @@ -79,9 +85,6 @@ type IPQueueItem struct { RetryCount int LeaseExpiresAt *time.Time EgressComplete bool - Site1Complete bool - Site2Complete bool - Site3Complete bool OverallResult string AssignedAt *time.Time FIPAssociatedAt *time.Time @@ -119,10 +122,13 @@ type Event struct { } type Site struct { - Index int - SiteID string - CreatedAt time.Time - UpdatedAt time.Time + Index int + SiteID string + Hostname string + State string + LastHeartbeatAt *time.Time + CreatedAt time.Time + UpdatedAt time.Time } type TargetGroup struct { diff --git a/internal/db/queries_dynconfig_test.go b/internal/db/queries_dynconfig_test.go index dd6ab0e..c8101ce 100644 --- a/internal/db/queries_dynconfig_test.go +++ b/internal/db/queries_dynconfig_test.go @@ -28,13 +28,18 @@ func TestUpsertSiteValidation(t *testing.T) { if err := d.UpsertSite(ctx, 0, "site-1"); !errors.Is(err, ErrValidation) { t.Fatalf("expected ErrValidation for idx=0, got %v", err) } - if err := d.UpsertSite(ctx, 4, "site-1"); !errors.Is(err, ErrValidation) { - t.Fatalf("expected ErrValidation for idx=4, got %v", err) - } if err := d.UpsertSite(ctx, 1, ""); !errors.Is(err, ErrValidation) { t.Fatalf("expected ErrValidation for empty site_id, got %v", err) } + // No cap on the number of slots — idx=4 and idx=100 must both succeed. + if err := d.UpsertSite(ctx, 4, "site-4"); err != nil { + t.Fatalf("upsert site 4: %v", err) + } + if err := d.UpsertSite(ctx, 100, "site-100"); err != nil { + t.Fatalf("upsert site 100: %v", err) + } + if err := d.UpsertSite(ctx, 1, "site-1"); err != nil { t.Fatalf("upsert site 1: %v", err) } @@ -352,6 +357,9 @@ func TestDeleteIP(t *testing.T) { if err := d.InsertEvent(ctx, Event{SourceType: "control-api", IPID: &ip.ID, EventType: "test_event", OccurredAt: Now()}); err != nil { t.Fatalf("insert event: %v", err) } + if err := d.SetSiteComplete(ctx, ip.ID, 1); err != nil { + t.Fatalf("set site complete: %v", err) + } if err := d.DeleteIP(ctx, ip.ID); err != nil { t.Fatalf("delete ip: %v", err) @@ -366,6 +374,13 @@ func TestDeleteIP(t *testing.T) { if len(checks) != 0 { t.Fatalf("expected checks gone, got %+v", checks) } + completed, err := d.ListCompletedSiteIndices(ctx, ip.ID, ip.AttemptNumber) + if err != nil { + t.Fatalf("list completed site indices: %v", err) + } + if len(completed) != 0 { + t.Fatalf("expected ip_site_checks rows gone, got %+v", completed) + } events, err := d.ListEventsForIP(ctx, ip.ID) if err != nil { t.Fatalf("list events: %v", err) diff --git a/internal/db/queries_ipqueue.go b/internal/db/queries_ipqueue.go index 8c74304..09638a6 100644 --- a/internal/db/queries_ipqueue.go +++ b/internal/db/queries_ipqueue.go @@ -196,7 +196,7 @@ func (d *DB) RequeueOrFail(ctx context.Context, ipID int64, validatorID string, _, err = tx.ExecContext(ctx, ` UPDATE ip_queue SET state=?, owner_validator_id=NULL, fip_id='', retry_count=?, attempt_number=attempt_number+1, - lease_expires_at=NULL, egress_complete=0, site1_complete=0, site2_complete=0, site3_complete=0, + lease_expires_at=NULL, egress_complete=0, overall_result='', assigned_at=NULL, fip_associated_at=NULL, updated_at=? WHERE id=? `, nextState, retryCount, now, ipID) @@ -282,7 +282,7 @@ func (d *DB) SubmitIPs(ctx context.Context, addresses []string) (SubmitIPsResult UPDATE ip_queue SET state=?, sequence=?, owner_validator_id=NULL, fip_id='', retry_count=0, attempt_number=attempt_number+1, lease_expires_at=NULL, egress_complete=0, - site1_complete=0, site2_complete=0, site3_complete=0, overall_result='', + overall_result='', assigned_at=NULL, fip_associated_at=NULL, aggregated_at=NULL, fip_released_at=NULL, updated_at=? WHERE ip_address=? `, IPQueued, seq, now, addr); err != nil { @@ -409,6 +409,9 @@ func deleteIPTx(ctx context.Context, tx *sql.Tx, ipID int64) error { if _, err := tx.ExecContext(ctx, `DELETE FROM events WHERE ip_id=?`, ipID); err != nil { return fmt.Errorf("delete events: %w", err) } + if _, err := tx.ExecContext(ctx, `DELETE FROM ip_site_checks WHERE ip_id=?`, ipID); err != nil { + return fmt.Errorf("delete ip_site_checks: %w", err) + } res, err := tx.ExecContext(ctx, `DELETE FROM ip_queue WHERE id=?`, ipID) if err != nil { return fmt.Errorf("delete ip_queue row: %w", err) @@ -424,16 +427,6 @@ func (d *DB) SetEgressComplete(ctx context.Context, ipID int64) error { return err } -// SetSiteComplete marks completion for prober site 1, 2, or 3. -func (d *DB) SetSiteComplete(ctx context.Context, ipID int64, siteIndex int) error { - col := map[int]string{1: "site1_complete", 2: "site2_complete", 3: "site3_complete"}[siteIndex] - if col == "" { - return fmt.Errorf("invalid site index %d", siteIndex) - } - _, err := d.ExecContext(ctx, fmt.Sprintf(`UPDATE ip_queue SET %s=1, updated_at=? WHERE id=?`, col), timeToDB(Now()), ipID) - return err -} - func (d *DB) GetIP(ctx context.Context, ipID int64) (*IPQueueItem, error) { row := d.QueryRowContext(ctx, ipQueueSelect+`WHERE id=?`, ipID) return scanIPQueueItem(row) @@ -479,7 +472,7 @@ func (d *DB) ListExpiredLeases(ctx context.Context, now time.Time) ([]IPQueueIte const ipQueueSelect = ` SELECT id, ip_address, sequence, state, owner_validator_id, fip_id, attempt_number, retry_count, - lease_expires_at, egress_complete, site1_complete, site2_complete, site3_complete, overall_result, + lease_expires_at, egress_complete, overall_result, assigned_at, fip_associated_at, aggregated_at, fip_released_at, created_at, updated_at FROM ip_queue ` @@ -504,7 +497,7 @@ func scanIPQueueItem(row rowScanner) (*IPQueueItem, error) { if err := row.Scan( &item.ID, &item.IPAddress, &item.Sequence, &item.State, &owner, &item.FIPID, &item.AttemptNumber, &item.RetryCount, &leaseExpires, - &item.EgressComplete, &item.Site1Complete, &item.Site2Complete, &item.Site3Complete, + &item.EgressComplete, &item.OverallResult, &assignedAt, &fipAssociatedAt, &aggregatedAt, &fipReleasedAt, &createdAt, &updatedAt, ); err != nil { return nil, err diff --git a/internal/db/queries_ipsitechecks.go b/internal/db/queries_ipsitechecks.go new file mode 100644 index 0000000..f5e8550 --- /dev/null +++ b/internal/db/queries_ipsitechecks.go @@ -0,0 +1,39 @@ +package db + +import "context" + +// SetSiteComplete marks a given site's inbound checks complete for an IP's +// *current* attempt. attempt_number is resolved via subquery from +// ip_queue rather than taken as a parameter, so the call signature stays +// the same as before the fixed-3-columns design was replaced by this +// table. No cap on siteIndex. +func (d *DB) SetSiteComplete(ctx context.Context, ipID int64, siteIndex int) error { + _, err := d.ExecContext(ctx, ` + INSERT INTO ip_site_checks (ip_id, attempt_number, site_idx, complete, completed_at) + SELECT ?, attempt_number, ?, 1, ? FROM ip_queue WHERE id=? + ON CONFLICT(ip_id, attempt_number, site_idx) DO UPDATE SET complete=1, completed_at=excluded.completed_at + `, ipID, siteIndex, timeToDB(Now()), ipID) + return err +} + +// ListCompletedSiteIndices returns the set of site indices that have +// reported completion for the given IP's attempt. +func (d *DB) ListCompletedSiteIndices(ctx context.Context, ipID int64, attemptNumber int) (map[int]bool, error) { + rows, err := d.QueryContext(ctx, ` + SELECT site_idx FROM ip_site_checks WHERE ip_id=? AND attempt_number=? AND complete=1 + `, ipID, attemptNumber) + if err != nil { + return nil, err + } + defer rows.Close() + + out := map[int]bool{} + for rows.Next() { + var idx int + if err := rows.Scan(&idx); err != nil { + return nil, err + } + out[idx] = true + } + return out, rows.Err() +} diff --git a/internal/db/queries_ipsitechecks_test.go b/internal/db/queries_ipsitechecks_test.go new file mode 100644 index 0000000..86557e4 --- /dev/null +++ b/internal/db/queries_ipsitechecks_test.go @@ -0,0 +1,92 @@ +package db + +import "testing" + +func TestSetSiteCompleteAndListCompletedSiteIndices(t *testing.T) { + d, ctx := newTestDB(t) + + if err := d.SeedQueue(ctx, []string{"1.2.3.4"}); err != nil { + t.Fatalf("seed queue: %v", err) + } + ip, err := d.GetIPByAddress(ctx, "1.2.3.4") + if err != nil { + t.Fatalf("get ip: %v", err) + } + + for _, idx := range []int{1, 2, 4} { + if err := d.SetSiteComplete(ctx, ip.ID, idx); err != nil { + t.Fatalf("set site %d complete: %v", idx, err) + } + } + + completed, err := d.ListCompletedSiteIndices(ctx, ip.ID, ip.AttemptNumber) + if err != nil { + t.Fatalf("list completed site indices: %v", err) + } + for _, idx := range []int{1, 2, 4} { + if !completed[idx] { + t.Fatalf("expected site %d complete, got %+v", idx, completed) + } + } + if completed[3] { + t.Fatalf("expected site 3 not complete, got %+v", completed) + } + + // Idempotent re-mark doesn't error or duplicate. + if err := d.SetSiteComplete(ctx, ip.ID, 1); err != nil { + t.Fatalf("re-mark site 1 complete: %v", err) + } +} + +// TestNewAttemptStartsWithNoCompletedSites confirms a bumped +// attempt_number (via SubmitIPs requeue) starts with an empty completed +// set even though the previous attempt had rows — proving no explicit +// reset is needed on retry, unlike the old fixed-3-column design. +func TestNewAttemptStartsWithNoCompletedSites(t *testing.T) { + d, ctx := newTestDB(t) + + if err := d.SeedQueue(ctx, []string{"1.2.3.4"}); err != nil { + t.Fatalf("seed queue: %v", err) + } + ip, err := d.GetIPByAddress(ctx, "1.2.3.4") + if err != nil { + t.Fatalf("get ip: %v", err) + } + if err := d.SetSiteComplete(ctx, ip.ID, 1); err != nil { + t.Fatalf("set site complete: %v", err) + } + + // Force the IP to a terminal state, then requeue via SubmitIPs — the + // same path an admin-triggered recheck takes — to bump attempt_number. + if _, err := d.ExecContext(ctx, `UPDATE ip_queue SET state=? WHERE id=?`, IPFailed, ip.ID); err != nil { + t.Fatalf("force failed: %v", err) + } + if _, err := d.SubmitIPs(ctx, []string{"1.2.3.4"}); err != nil { + t.Fatalf("resubmit: %v", err) + } + + requeued, err := d.GetIP(ctx, ip.ID) + if err != nil { + t.Fatalf("get ip after requeue: %v", err) + } + if requeued.AttemptNumber != ip.AttemptNumber+1 { + t.Fatalf("expected attempt_number bumped, got %d", requeued.AttemptNumber) + } + + completed, err := d.ListCompletedSiteIndices(ctx, ip.ID, requeued.AttemptNumber) + if err != nil { + t.Fatalf("list completed site indices for new attempt: %v", err) + } + if len(completed) != 0 { + t.Fatalf("expected no completed sites for the new attempt, got %+v", completed) + } + + // The old attempt's row is still there as history, not wiped. + old, err := d.ListCompletedSiteIndices(ctx, ip.ID, ip.AttemptNumber) + if err != nil { + t.Fatalf("list completed site indices for old attempt: %v", err) + } + if !old[1] { + t.Fatalf("expected old attempt's completion to remain as history, got %+v", old) + } +} diff --git a/internal/db/queries_sites.go b/internal/db/queries_sites.go index 0d9c74f..9ee2dde 100644 --- a/internal/db/queries_sites.go +++ b/internal/db/queries_sites.go @@ -4,12 +4,13 @@ import ( "context" "database/sql" "fmt" + "time" ) // ListSites returns all configured prober sites, ordered by slot index. func (d *DB) ListSites(ctx context.Context) ([]Site, error) { rows, err := d.QueryContext(ctx, ` - SELECT idx, site_id, created_at, updated_at FROM sites ORDER BY idx + SELECT idx, site_id, hostname, state, last_heartbeat_at, created_at, updated_at FROM sites ORDER BY idx `) if err != nil { return nil, err @@ -18,24 +19,37 @@ func (d *DB) ListSites(ctx context.Context) ([]Site, error) { var out []Site for rows.Next() { - var s Site - var createdAt, updatedAt string - if err := rows.Scan(&s.Index, &s.SiteID, &createdAt, &updatedAt); err != nil { + s, err := scanSite(rows) + if err != nil { return nil, err } - var err error - if s.CreatedAt, err = dbToTime(createdAt); err != nil { - return nil, err - } - if s.UpdatedAt, err = dbToTime(updatedAt); err != nil { - return nil, err - } - out = append(out, s) + out = append(out, *s) } return out, rows.Err() } -// GetSiteIndex resolves a configured site_id to its 1/2/3 slot index. It +func scanSite(row rowScanner) (*Site, error) { + var s Site + var lastHeartbeat sql.NullString + var createdAt, updatedAt string + if err := row.Scan(&s.Index, &s.SiteID, &s.Hostname, &s.State, &lastHeartbeat, &createdAt, &updatedAt); err != nil { + return nil, err + } + hb, err := nullStringToTimePtr(lastHeartbeat) + if err != nil { + return nil, err + } + s.LastHeartbeatAt = hb + if s.CreatedAt, err = dbToTime(createdAt); err != nil { + return nil, err + } + if s.UpdatedAt, err = dbToTime(updatedAt); err != nil { + return nil, err + } + return &s, nil +} + +// GetSiteIndex resolves a configured site_id to its slot index. It // returns (0, nil) — not an error — if no site with that ID is configured, // matching the "0 means unconfigured" convention used throughout the // prober-facing handlers. @@ -52,12 +66,21 @@ func (d *DB) GetSiteIndex(ctx context.Context, siteID string) (int, error) { } // UpsertSite assigns (or renames) the site occupying the given slot. idx -// must be 1, 2, or 3 — the schema hard-caps the number of prober slots at -// three (see ip_queue.site{1,2,3}_complete). site_id must be unique across +// has no upper bound — an admin can configure as many prober sites as they +// need — only idx >= 1 is enforced (slot numbering starts at 1, matching +// the existing admin-facing convention). site_id must be unique across // slots. +// +// Re-pointing an existing slot to a (possibly different) site_id always +// resets hostname/state/last_heartbeat_at to their "never connected" +// defaults: a new site_id is a new prober identity, and even a same-value +// PUT is treated as idempotent-intent rather than a heartbeat, so the old +// prober's connection status must not silently linger. The prober +// re-establishes state on its next register/heartbeat within +// poll_interval_seconds. func (d *DB) UpsertSite(ctx context.Context, idx int, siteID string) error { - if idx < 1 || idx > 3 { - return fmt.Errorf("site index must be 1, 2, or 3, got %d: %w", idx, ErrValidation) + if idx < 1 { + return fmt.Errorf("site index must be >= 1, got %d: %w", idx, ErrValidation) } if siteID == "" { return fmt.Errorf("site_id must not be empty: %w", ErrValidation) @@ -76,8 +99,10 @@ func (d *DB) UpsertSite(ctx context.Context, idx int, siteID string) error { _, err = d.ExecContext(ctx, ` INSERT INTO sites (idx, site_id, created_at, updated_at) VALUES (?, ?, ?, ?) - ON CONFLICT(idx) DO UPDATE SET site_id=excluded.site_id, updated_at=excluded.updated_at - `, idx, siteID, now, now) + ON CONFLICT(idx) DO UPDATE SET + site_id=excluded.site_id, updated_at=excluded.updated_at, + hostname='', state=?, last_heartbeat_at=NULL + `, idx, siteID, now, now, SiteUnregistered) if err != nil { return fmt.Errorf("upsert site: %w", err) } @@ -95,3 +120,75 @@ func (d *DB) DeleteSite(ctx context.Context, idx int) error { } return nil } + +// RegisterSiteProber records that a prober has (re)connected for this +// site: sets hostname, stamps last_heartbeat_at, and reactivates state to +// idle if it was unregistered or unreachable — mirrors RegisterValidator's +// reactivation UPDATE. Called once at prober startup. Returns +// sql.ErrNoRows if site_id is unconfigured. +func (d *DB) RegisterSiteProber(ctx context.Context, siteID, hostname string) error { + now := timeToDB(Now()) + res, err := d.ExecContext(ctx, ` + UPDATE sites SET hostname=?, last_heartbeat_at=?, updated_at=?, + state = CASE WHEN state IN (?, ?) THEN ? ELSE state END + WHERE site_id=? + `, hostname, now, now, SiteUnregistered, SiteUnreachable, SiteIdle, siteID) + if err != nil { + return fmt.Errorf("register site prober: %w", err) + } + n, _ := res.RowsAffected() + if n == 0 { + return sql.ErrNoRows + } + return nil +} + +// SiteHeartbeat is the periodic "I'm alive" touch — mirrors Heartbeat for +// validators exactly: stamps last_heartbeat_at and flips state from +// unreachable back to idle (does not touch hostname). +func (d *DB) SiteHeartbeat(ctx context.Context, siteID string) error { + now := timeToDB(Now()) + res, err := d.ExecContext(ctx, ` + UPDATE sites SET last_heartbeat_at=?, updated_at=?, + state = CASE WHEN state=? THEN ? ELSE state END + WHERE site_id=? + `, now, now, SiteUnreachable, SiteIdle, siteID) + if err != nil { + return fmt.Errorf("site heartbeat: %w", err) + } + n, _ := res.RowsAffected() + if n == 0 { + return sql.ErrNoRows + } + return nil +} + +// ListStaleSiteHeartbeats returns sites whose last heartbeat predates the +// given cutoff and that aren't already marked unreachable. +func (d *DB) ListStaleSiteHeartbeats(ctx context.Context, cutoff time.Time) ([]Site, error) { + rows, err := d.QueryContext(ctx, ` + SELECT idx, site_id, hostname, state, last_heartbeat_at, created_at, updated_at + FROM sites + WHERE state != ? AND last_heartbeat_at IS NOT NULL AND last_heartbeat_at < ? + `, SiteUnreachable, timeToDB(cutoff)) + if err != nil { + return nil, err + } + defer rows.Close() + + var out []Site + for rows.Next() { + s, err := scanSite(rows) + if err != nil { + return nil, err + } + out = append(out, *s) + } + return out, rows.Err() +} + +func (d *DB) MarkSiteUnreachable(ctx context.Context, siteID string) error { + _, err := d.ExecContext(ctx, `UPDATE sites SET state=?, updated_at=? WHERE site_id=?`, + SiteUnreachable, timeToDB(Now()), siteID) + return err +} diff --git a/internal/db/queries_sites_heartbeat_test.go b/internal/db/queries_sites_heartbeat_test.go new file mode 100644 index 0000000..680158f --- /dev/null +++ b/internal/db/queries_sites_heartbeat_test.go @@ -0,0 +1,172 @@ +package db + +import ( + "database/sql" + "errors" + "testing" + "time" +) + +func TestRegisterSiteProberUnknownSiteReturnsErrNoRows(t *testing.T) { + d, ctx := newTestDB(t) + if err := d.RegisterSiteProber(ctx, "unknown-site", "host-1"); !errors.Is(err, sql.ErrNoRows) { + t.Fatalf("expected sql.ErrNoRows for unknown site, got %v", err) + } +} + +func TestRegisterSiteProberSetsHostnameAndIdleState(t *testing.T) { + d, ctx := newTestDB(t) + if err := d.UpsertSite(ctx, 1, "site-1"); err != nil { + t.Fatalf("upsert site: %v", err) + } + + sites, err := d.ListSites(ctx) + if err != nil { + t.Fatalf("list sites: %v", err) + } + if sites[0].State != SiteUnregistered { + t.Fatalf("expected freshly-created slot unregistered, got %s", sites[0].State) + } + + if err := d.RegisterSiteProber(ctx, "site-1", "probe-host-1"); err != nil { + t.Fatalf("register site prober: %v", err) + } + sites, err = d.ListSites(ctx) + if err != nil { + t.Fatalf("list sites after register: %v", err) + } + if sites[0].State != SiteIdle { + t.Fatalf("expected idle after register, got %s", sites[0].State) + } + if sites[0].Hostname != "probe-host-1" { + t.Fatalf("expected hostname set, got %q", sites[0].Hostname) + } + if sites[0].LastHeartbeatAt == nil { + t.Fatalf("expected last_heartbeat_at stamped") + } +} + +func TestRegisterSiteProberReactivatesFromUnreachable(t *testing.T) { + d, ctx := newTestDB(t) + if err := d.UpsertSite(ctx, 1, "site-1"); err != nil { + t.Fatalf("upsert site: %v", err) + } + if err := d.MarkSiteUnreachable(ctx, "site-1"); err != nil { + t.Fatalf("mark unreachable: %v", err) + } + if err := d.RegisterSiteProber(ctx, "site-1", "probe-host-1"); err != nil { + t.Fatalf("register site prober: %v", err) + } + sites, err := d.ListSites(ctx) + if err != nil { + t.Fatalf("list sites: %v", err) + } + if sites[0].State != SiteIdle { + t.Fatalf("expected reactivated to idle, got %s", sites[0].State) + } +} + +func TestSiteHeartbeatRoundTrip(t *testing.T) { + d, ctx := newTestDB(t) + if err := d.UpsertSite(ctx, 1, "site-1"); err != nil { + t.Fatalf("upsert site: %v", err) + } + if err := d.SiteHeartbeat(ctx, "unknown-site"); !errors.Is(err, sql.ErrNoRows) { + t.Fatalf("expected sql.ErrNoRows for unknown site, got %v", err) + } + + if err := d.SiteHeartbeat(ctx, "site-1"); err != nil { + t.Fatalf("heartbeat: %v", err) + } + sites, err := d.ListSites(ctx) + if err != nil { + t.Fatalf("list sites: %v", err) + } + // A plain heartbeat (no prior register) doesn't flip an unregistered + // site to idle — only reactivates from unreachable, mirroring + // validators' Heartbeat. + if sites[0].State != SiteUnregistered { + t.Fatalf("expected still unregistered after bare heartbeat, got %s", sites[0].State) + } + if sites[0].LastHeartbeatAt == nil { + t.Fatalf("expected last_heartbeat_at stamped") + } + if sites[0].Hostname != "" { + t.Fatalf("expected heartbeat not to set hostname, got %q", sites[0].Hostname) + } + + if err := d.MarkSiteUnreachable(ctx, "site-1"); err != nil { + t.Fatalf("mark unreachable: %v", err) + } + if err := d.SiteHeartbeat(ctx, "site-1"); err != nil { + t.Fatalf("heartbeat after unreachable: %v", err) + } + sites, err = d.ListSites(ctx) + if err != nil { + t.Fatalf("list sites: %v", err) + } + if sites[0].State != SiteIdle { + t.Fatalf("expected heartbeat to reactivate from unreachable to idle, got %s", sites[0].State) + } +} + +func TestListStaleSiteHeartbeatsAndMarkUnreachable(t *testing.T) { + d, ctx := newTestDB(t) + if err := d.UpsertSite(ctx, 1, "site-1"); err != nil { + t.Fatalf("upsert site: %v", err) + } + if err := d.UpsertSite(ctx, 2, "site-2"); err != nil { + t.Fatalf("upsert site: %v", err) + } + // site-3 never heartbeats at all — must not show up as "stale" (it's + // simply unregistered, which is a distinct, expected state). + if err := d.UpsertSite(ctx, 3, "site-3"); err != nil { + t.Fatalf("upsert site: %v", err) + } + + if err := d.SiteHeartbeat(ctx, "site-1"); err != nil { + t.Fatalf("heartbeat site-1: %v", err) + } + if err := d.SiteHeartbeat(ctx, "site-2"); err != nil { + t.Fatalf("heartbeat site-2: %v", err) + } + + cutoff := Now().Add(1 * time.Hour) // everything heartbeated is "before" this + stale, err := d.ListStaleSiteHeartbeats(ctx, cutoff) + if err != nil { + t.Fatalf("list stale site heartbeats: %v", err) + } + if len(stale) != 2 { + t.Fatalf("expected 2 stale sites (site-1, site-2), got %+v", stale) + } + + for _, s := range stale { + if err := d.MarkSiteUnreachable(ctx, s.SiteID); err != nil { + t.Fatalf("mark unreachable %s: %v", s.SiteID, err) + } + } + + // Now already-unreachable sites are excluded from a subsequent sweep. + stale, err = d.ListStaleSiteHeartbeats(ctx, cutoff) + if err != nil { + t.Fatalf("list stale site heartbeats after mark: %v", err) + } + if len(stale) != 0 { + t.Fatalf("expected no stale sites left, got %+v", stale) + } + + sites, err := d.ListSites(ctx) + if err != nil { + t.Fatalf("list sites: %v", err) + } + byID := map[string]Site{} + for _, s := range sites { + byID[s.SiteID] = s + } + if byID["site-1"].State != SiteUnreachable || byID["site-2"].State != SiteUnreachable { + t.Fatalf("expected site-1/site-2 unreachable, got %+v", sites) + } + if byID["site-3"].State != SiteUnregistered { + t.Fatalf("expected site-3 still unregistered (never heartbeated), got %s", byID["site-3"].State) + } +} diff --git a/internal/httpapi/dto_admin.go b/internal/httpapi/dto_admin.go index afd1fe0..90e14b7 100644 --- a/internal/httpapi/dto_admin.go +++ b/internal/httpapi/dto_admin.go @@ -1,5 +1,7 @@ package httpapi +import "time" + // DTOs for the admin queue-management and dynamic-config endpoints // (/api/v1/admin/ips, /api/v1/admin/config/*). Unlike the older read-only // admin endpoints (which marshal internal/db model structs directly, in @@ -46,8 +48,11 @@ type updateValidatorRequest struct { } type siteDTO struct { - Index int `json:"index"` - SiteID string `json:"site_id"` + Index int `json:"index"` + SiteID string `json:"site_id"` + Hostname string `json:"hostname"` + State string `json:"state"` + LastHeartbeatAt *time.Time `json:"last_heartbeat_at"` } type putSiteRequest struct { diff --git a/internal/httpapi/handlers_config.go b/internal/httpapi/handlers_config.go index a445ed4..7caa451 100644 --- a/internal/httpapi/handlers_config.go +++ b/internal/httpapi/handlers_config.go @@ -3,6 +3,8 @@ package httpapi import ( "net/http" "strconv" + + "cloudipvalidator/internal/db" ) // --- validators --- @@ -70,7 +72,10 @@ func (s *Server) handleConfigListSites(w http.ResponseWriter, r *http.Request) { } out := make([]siteDTO, len(sites)) for i, site := range sites { - out[i] = siteDTO{Index: site.Index, SiteID: site.SiteID} + out[i] = siteDTO{ + Index: site.Index, SiteID: site.SiteID, Hostname: site.Hostname, + State: site.State, LastHeartbeatAt: site.LastHeartbeatAt, + } } writeJSON(w, http.StatusOK, out) } @@ -90,7 +95,9 @@ func (s *Server) handleConfigPutSite(w http.ResponseWriter, r *http.Request) { writeDBError(w, err) return } - writeJSON(w, http.StatusOK, siteDTO{Index: idx, SiteID: req.SiteID}) + // UpsertSite always resets hostname/state/last_heartbeat_at to their + // "never connected" defaults on write (see its doc comment). + writeJSON(w, http.StatusOK, siteDTO{Index: idx, SiteID: req.SiteID, State: db.SiteUnregistered}) } func (s *Server) handleConfigDeleteSite(w http.ResponseWriter, r *http.Request) { diff --git a/internal/httpapi/handlers_config_test.go b/internal/httpapi/handlers_config_test.go index d38c387..5fd9f57 100644 --- a/internal/httpapi/handlers_config_test.go +++ b/internal/httpapi/handlers_config_test.go @@ -623,3 +623,110 @@ func TestInboundChecksReflectedInProberAssignmentsWithoutRestart(t *testing.T) { t.Fatalf("expected updated {[8080] false} without restart, got %+v", assignments) } } + +// TestProberRegisterSetsHostnameAndIdleState proves POST +// /api/v1/probers/register persists the calling prober's hostname and +// flips the site's state to idle, visible via GET +// /api/v1/admin/config/sites — the prober-availability analog of +// validator registration. +func TestProberRegisterSetsHostnameAndIdleState(t *testing.T) { + fc, _, _, _ := newConfigTestHarness(t) + + fc.do(http.MethodPut, "/api/v1/admin/config/sites/1", putSiteRequest{SiteID: "site-1"}) + + resp, body := fc.do(http.MethodGet, "/api/v1/admin/config/sites", nil) + if resp.StatusCode != http.StatusOK { + t.Fatalf("get sites: status=%d body=%s", resp.StatusCode, body) + } + var sites []siteDTO + if err := json.Unmarshal(body, &sites); err != nil { + t.Fatalf("unmarshal sites: %v", err) + } + if len(sites) != 1 || sites[0].State != "unregistered" { + t.Fatalf("expected freshly-created slot unregistered, got %+v", sites) + } + + resp, body = fc.do(http.MethodPost, "/api/v1/probers/register", registerProberRequest{SiteID: "site-1", Hostname: "probe-host-1"}) + if resp.StatusCode != http.StatusOK { + t.Fatalf("register prober: status=%d body=%s", resp.StatusCode, body) + } + + resp, body = fc.do(http.MethodGet, "/api/v1/admin/config/sites", nil) + if resp.StatusCode != http.StatusOK { + t.Fatalf("get sites after register: status=%d body=%s", resp.StatusCode, body) + } + if err := json.Unmarshal(body, &sites); err != nil { + t.Fatalf("unmarshal sites: %v", err) + } + if sites[0].State != "idle" || sites[0].Hostname != "probe-host-1" { + t.Fatalf("expected {state:idle, hostname:probe-host-1}, got %+v", sites[0]) + } + if sites[0].LastHeartbeatAt == nil { + t.Fatalf("expected last_heartbeat_at stamped") + } +} + +// TestProberHeartbeat404UnknownSite proves the heartbeat endpoint rejects +// an unconfigured site_id, mirroring the validator heartbeat's 404. +func TestProberHeartbeat404UnknownSite(t *testing.T) { + fc, _, _, _ := newConfigTestHarness(t) + resp, body := fc.do(http.MethodPost, "/api/v1/probers/unknown-site/heartbeat", nil) + if resp.StatusCode != http.StatusNotFound { + t.Fatalf("expected 404 for unknown site_id, status=%d body=%s", resp.StatusCode, body) + } +} + +// TestSweepStaleSiteHeartbeatsMarksUnreachable is the end-to-end proof that +// a prober that stops heartbeating gets marked unreachable by the sweep, +// visible via the admin API — the prober-side analog of +// TestFIPSettleDelayGatesAssignmentEndpoint's wire-level style. +func TestSweepStaleSiteHeartbeatsMarksUnreachable(t *testing.T) { + fc, d, orch, _ := newConfigTestHarness(t) + ctx := context.Background() + + fc.do(http.MethodPut, "/api/v1/admin/config/sites/1", putSiteRequest{SiteID: "site-1"}) + resp, body := fc.do(http.MethodPost, "/api/v1/probers/register", registerProberRequest{SiteID: "site-1", Hostname: "probe-host-1"}) + if resp.StatusCode != http.StatusOK { + t.Fatalf("register prober: status=%d body=%s", resp.StatusCode, body) + } + + // Backdate last_heartbeat_at past HeartbeatTimeoutSeconds (30s, per + // newConfigTestHarness) directly in the DB, rather than sleeping 30+ + // real seconds in the test. + old := db.Now().Add(-time.Hour).UTC().Format(time.RFC3339Nano) + if _, err := d.ExecContext(ctx, `UPDATE sites SET last_heartbeat_at=? WHERE site_id=?`, old, "site-1"); err != nil { + t.Fatalf("backdate heartbeat: %v", err) + } + + if err := orch.SweepStaleSiteHeartbeats(ctx); err != nil { + t.Fatalf("sweep stale site heartbeats: %v", err) + } + + resp, body = fc.do(http.MethodGet, "/api/v1/admin/config/sites", nil) + if resp.StatusCode != http.StatusOK { + t.Fatalf("get sites: status=%d body=%s", resp.StatusCode, body) + } + var sites []siteDTO + if err := json.Unmarshal(body, &sites); err != nil { + t.Fatalf("unmarshal sites: %v", err) + } + if len(sites) != 1 || sites[0].State != "unreachable" { + t.Fatalf("expected site-1 marked unreachable after sweep, got %+v", sites) + } + + // A fresh heartbeat brings it back to idle. + resp, body = fc.do(http.MethodPost, "/api/v1/probers/site-1/heartbeat", nil) + if resp.StatusCode != http.StatusOK { + t.Fatalf("heartbeat: status=%d body=%s", resp.StatusCode, body) + } + resp, body = fc.do(http.MethodGet, "/api/v1/admin/config/sites", nil) + if resp.StatusCode != http.StatusOK { + t.Fatalf("get sites after heartbeat: status=%d body=%s", resp.StatusCode, body) + } + if err := json.Unmarshal(body, &sites); err != nil { + t.Fatalf("unmarshal sites: %v", err) + } + if sites[0].State != "idle" { + t.Fatalf("expected site-1 back to idle after heartbeat, got %+v", sites) + } +} diff --git a/internal/httpapi/handlers_prober.go b/internal/httpapi/handlers_prober.go index b7de846..138bccd 100644 --- a/internal/httpapi/handlers_prober.go +++ b/internal/httpapi/handlers_prober.go @@ -22,10 +22,25 @@ func (s *Server) handleProberRegister(w http.ResponseWriter, r *http.Request) { writeError(w, http.StatusBadRequest, "unknown site_id: "+req.SiteID) return } + if err := s.DB.RegisterSiteProber(r.Context(), req.SiteID, req.Hostname); err != nil { + writeError(w, http.StatusInternalServerError, err.Error()) + return + } s.Orch.RecordEvent(r.Context(), "prober", req.SiteID, nil, "registered", "") writeJSON(w, http.StatusOK, registerAgentResponse{OK: true, PollIntervalSeconds: s.Orch.Cfg.PollIntervalSeconds}) } +func (s *Server) handleProberHeartbeat(w http.ResponseWriter, r *http.Request) { + siteID := r.PathValue("site_id") + var req heartbeatRequest + _ = readJSON(r, &req) // heartbeat body is informational only; tolerate empty/missing + if err := s.DB.SiteHeartbeat(r.Context(), siteID); err != nil { + writeError(w, http.StatusNotFound, "unknown site_id: "+siteID) + return + } + writeJSON(w, http.StatusOK, okResponse{OK: true}) +} + // handleProberAssignments returns every IP currently in the checking // state — probers work the whole active set each poll, not one IP at a // time, since multiple validators run in parallel. diff --git a/internal/httpapi/routes.go b/internal/httpapi/routes.go index 6a13e80..49a7972 100644 --- a/internal/httpapi/routes.go +++ b/internal/httpapi/routes.go @@ -14,6 +14,7 @@ func (s *Server) routes(mux *http.ServeMux) { mux.HandleFunc("POST /api/v1/agents/{id}/complete", s.handleAgentComplete) mux.HandleFunc("POST /api/v1/probers/register", s.handleProberRegister) + mux.HandleFunc("POST /api/v1/probers/{site_id}/heartbeat", s.handleProberHeartbeat) mux.HandleFunc("GET /api/v1/probers/{site_id}/assignments", s.handleProberAssignments) mux.HandleFunc("POST /api/v1/probers/{site_id}/results", s.handleProberResults) diff --git a/internal/orchestrator/orchestrator.go b/internal/orchestrator/orchestrator.go index 29d64b8..806d0e2 100644 --- a/internal/orchestrator/orchestrator.go +++ b/internal/orchestrator/orchestrator.go @@ -401,7 +401,12 @@ func (o *Orchestrator) sweepCheckingWindow(ctx context.Context) error { return fmt.Errorf("list sites: %w", err) } for _, item := range checking { - if !o.isReadyToAggregate(item, deadline, sites) { + ready, err := o.isReadyToAggregate(ctx, item, deadline, sites) + if err != nil { + o.Log.Error("check ready to aggregate", "ip_id", item.ID, "err", err) + continue + } + if !ready { continue } if err := o.aggregateAndRelease(ctx, item); err != nil { @@ -416,34 +421,26 @@ func (o *Orchestrator) sweepCheckingWindow(ctx context.Context) error { // checking-window deadline. Which inbound sources it's expecting is driven // entirely by the currently configured sites — inbound checks are optional: // an empty (or partial) sites configuration means this IP is ready as soon -// as egress completes (or after the corresponding subset of siteN_complete -// flags), with no need to wait on a prober that will never exist. This is -// what makes inbound checks genuinely opt-in rather than a hardcoded -// expectation of exactly three sites. -func (o *Orchestrator) isReadyToAggregate(item db.IPQueueItem, deadline time.Time, sites []db.Site) bool { +// as egress completes (or after the corresponding subset of sites has +// reported in ip_site_checks), with no need to wait on a prober that will +// never exist, and no cap on how many sites can be configured. +func (o *Orchestrator) isReadyToAggregate(ctx context.Context, item db.IPQueueItem, deadline time.Time, sites []db.Site) (bool, error) { if item.AssignedAt != nil && item.AssignedAt.Before(deadline) { - return true + return true, nil } if !item.EgressComplete { - return false + return false, nil + } + completed, err := o.DB.ListCompletedSiteIndices(ctx, item.ID, item.AttemptNumber) + if err != nil { + return false, err } for _, s := range sites { - switch s.Index { - case 1: - if !item.Site1Complete { - return false - } - case 2: - if !item.Site2Complete { - return false - } - case 3: - if !item.Site3Complete { - return false - } + if !completed[s.Index] { + return false, nil } } - return true + return true, nil } func (o *Orchestrator) aggregateAndRelease(ctx context.Context, item db.IPQueueItem) error { @@ -585,6 +582,25 @@ func (o *Orchestrator) SweepStaleHeartbeats(ctx context.Context) error { return nil } +// SweepStaleSiteHeartbeats marks prober sites unreachable if they haven't +// heartbeated within HeartbeatTimeoutSeconds — mirrors SweepStaleHeartbeats +// exactly, for sites instead of validators. +func (o *Orchestrator) SweepStaleSiteHeartbeats(ctx context.Context) error { + cutoff := db.Now().Add(-time.Duration(o.Cfg.HeartbeatTimeoutSeconds) * time.Second) + stale, err := o.DB.ListStaleSiteHeartbeats(ctx, cutoff) + if err != nil { + return err + } + for _, s := range stale { + if err := o.DB.MarkSiteUnreachable(ctx, s.SiteID); err != nil { + o.Log.Error("mark site unreachable", "site_id", s.SiteID, "err", err) + continue + } + o.event(ctx, "control-api", "", nil, "site_unreachable", fmt.Sprintf(`{"site_id":%q}`, s.SiteID)) + } + return nil +} + // RecordEvent is the exported entry point httpapi uses to log // agent/prober-reported audit events (config_received, fip_changed, // error, etc.) through the same path as internally generated events. diff --git a/internal/orchestrator/orchestrator_test.go b/internal/orchestrator/orchestrator_test.go index 8c34662..62751ee 100644 --- a/internal/orchestrator/orchestrator_test.go +++ b/internal/orchestrator/orchestrator_test.go @@ -624,3 +624,65 @@ func TestExpectedCheckCountReflectsInboundChecksConfigChange(t *testing.T) { t.Fatalf("expected pass, got %s", ip.OverallResult) } } + +// TestFourSitesAllMustReportBeforeAggregation confirms there's no hardcoded +// cap of three sites: with 4 sites configured, aggregation must wait on +// all 4, not silently treat the 4th as always-complete. +func TestFourSitesAllMustReportBeforeAggregation(t *testing.T) { + ctx := context.Background() + sites := []config.SiteConfig{ + {SiteID: "site-1", Index: 1}, {SiteID: "site-2", Index: 2}, + {SiteID: "site-3", Index: 3}, {SiteID: "site-4", Index: 4}, + } + o, d, mock := newTestOrchestratorWithSites(t, 180, sites) + mock.Seed("fip-1", "1.2.3.4", "svc-project") + _ = d.RegisterValidator(ctx, "validator-1", "host-1", "port-1", "v0.1") + _ = d.SeedQueue(ctx, []string{"1.2.3.4"}) + + o.Tick(ctx) + ip, _ := d.GetIPByAddress(ctx, "1.2.3.4") + _ = o.SelfCheckResult(ctx, "validator-1", ip.ID, true, "ok") + ip, _ = d.GetIP(ctx, ip.ID) + + _ = o.RecordCheck(ctx, db.Check{ + IPID: ip.ID, IPAddress: ip.IPAddress, AttemptNumber: ip.AttemptNumber, + ValidatorID: "validator-1", Source: db.SourceEgress, CheckType: "https", + Target: "https://example.test", Success: true, CheckedAt: db.Now(), + }) + _ = o.MarkEgressComplete(ctx, ip.ID) + + for _, site := range []int{1, 2, 3} { + for _, ct := range []string{"tcp-22", "tcp-80", "icmp"} { + _ = 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(), + }) + } + _ = o.MarkSiteComplete(ctx, ip.ID, site) + } + + // Sites 1-3 reported, site 4 (the case beyond the old fixed-3 cap) + // hasn't — must not aggregate yet. + o.Tick(ctx) + ip, _ = d.GetIP(ctx, ip.ID) + if ip.State != db.IPChecking { + t.Fatalf("expected still checking (site-4 pending), got %s", ip.State) + } + + for _, ct := range []string{"tcp-22", "tcp-80", "icmp"} { + _ = o.RecordCheck(ctx, db.Check{ + IPID: ip.ID, IPAddress: ip.IPAddress, AttemptNumber: ip.AttemptNumber, + Source: db.InboundSource(4), CheckType: ct, Target: ip.IPAddress, Success: true, CheckedAt: db.Now(), + }) + } + _ = o.MarkSiteComplete(ctx, ip.ID, 4) + + o.Tick(ctx) + ip, _ = d.GetIP(ctx, ip.ID) + if ip.State != db.IPDone { + t.Fatalf("expected done once all 4 sites reported, got %s", ip.State) + } + if ip.OverallResult != db.ResultPass { + t.Fatalf("expected pass, got %s", ip.OverallResult) + } +} diff --git a/internal/probercore/probercore.go b/internal/probercore/probercore.go index dbbf159..342e02e 100644 --- a/internal/probercore/probercore.go +++ b/internal/probercore/probercore.go @@ -89,6 +89,11 @@ type resultDTO struct { } func (p *Prober) pollOnce(ctx context.Context) { + if _, err := p.client.Do(ctx, "POST", "/api/v1/probers/"+p.cfg.SiteID+"/heartbeat", nil, nil); err != nil { + p.log.Error("heartbeat", "err", err) + return + } + var assignments []assignment ok, err := p.client.Do(ctx, "GET", "/api/v1/probers/"+p.cfg.SiteID+"/assignments", nil, &assignments) if err != nil { diff --git a/internal/probercore/probercore_test.go b/internal/probercore/probercore_test.go new file mode 100644 index 0000000..7e6d668 --- /dev/null +++ b/internal/probercore/probercore_test.go @@ -0,0 +1,96 @@ +package probercore + +import ( + "context" + "log/slog" + "net/http" + "net/http/httptest" + "os" + "sync" + "testing" + "time" + + "cloudipvalidator/internal/apiclient" + "cloudipvalidator/internal/config" +) + +func testLogger() *slog.Logger { + return slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{Level: slog.LevelError})) +} + +// TestPollOnceCallsHeartbeatBeforeAssignments mirrors agentcore's +// heartbeat-then-assignment ordering: the prober must announce it's alive +// on every poll cycle before fetching work, not just once at startup. +func TestPollOnceCallsHeartbeatBeforeAssignments(t *testing.T) { + var mu sync.Mutex + var calls []string + + mux := http.NewServeMux() + mux.HandleFunc("POST /api/v1/probers/site-1/heartbeat", func(w http.ResponseWriter, r *http.Request) { + mu.Lock() + calls = append(calls, "heartbeat") + mu.Unlock() + w.WriteHeader(http.StatusOK) + w.Write([]byte(`{"ok":true}`)) + }) + mux.HandleFunc("GET /api/v1/probers/site-1/assignments", func(w http.ResponseWriter, r *http.Request) { + mu.Lock() + calls = append(calls, "assignments") + mu.Unlock() + w.WriteHeader(http.StatusOK) + w.Write([]byte(`[]`)) + }) + ts := httptest.NewServer(mux) + defer ts.Close() + + p := &Prober{ + cfg: &config.Prober{SiteID: "site-1"}, + client: apiclient.New(ts.URL, 5*time.Second), + log: testLogger(), + } + p.pollOnce(context.Background()) + + mu.Lock() + defer mu.Unlock() + if len(calls) != 2 || calls[0] != "heartbeat" || calls[1] != "assignments" { + t.Fatalf("expected [heartbeat assignments], got %v", calls) + } +} + +// TestPollOnceSkipsAssignmentsWhenHeartbeatFails confirms a failed +// heartbeat aborts the poll cycle early — mirrors agentcore.pollOnce's +// same early-return behavior. +func TestPollOnceSkipsAssignmentsWhenHeartbeatFails(t *testing.T) { + var mu sync.Mutex + var calls []string + + mux := http.NewServeMux() + mux.HandleFunc("POST /api/v1/probers/site-1/heartbeat", func(w http.ResponseWriter, r *http.Request) { + mu.Lock() + calls = append(calls, "heartbeat") + mu.Unlock() + w.WriteHeader(http.StatusNotFound) + }) + mux.HandleFunc("GET /api/v1/probers/site-1/assignments", func(w http.ResponseWriter, r *http.Request) { + mu.Lock() + calls = append(calls, "assignments") + mu.Unlock() + w.WriteHeader(http.StatusOK) + w.Write([]byte(`[]`)) + }) + ts := httptest.NewServer(mux) + defer ts.Close() + + p := &Prober{ + cfg: &config.Prober{SiteID: "site-1"}, + client: apiclient.New(ts.URL, 5*time.Second), + log: testLogger(), + } + p.pollOnce(context.Background()) + + mu.Lock() + defer mu.Unlock() + if len(calls) != 1 || calls[0] != "heartbeat" { + t.Fatalf("expected only [heartbeat] (assignments skipped on heartbeat failure), got %v", calls) + } +}