diff --git a/bin/SHA256SUMS b/bin/SHA256SUMS index e6a8783..8506586 100644 --- a/bin/SHA256SUMS +++ b/bin/SHA256SUMS @@ -1,4 +1,4 @@ -b5d1ae5c625c1ab12a42b91d0e7dbfe5eb15897fbc46e1b94db21aca5f84e77f control-api +5703aa9367f814fd1c02ab8ee821f4ebea28934387b084e0f843efe1505530b7 control-api 48c9b99fa88be751d9badfba8b7d80f894e326a743a90e85478f1d2251ce4ab6 validator-agent 43fb660b78179b204b57388241a508345881aa16f3531d77b8fd75af60e7f2d8 prober -4e1e550bf026b9f9b4aa60a9bf38e1e409693f7ed8d4fd1853175f882ba349ca admin-dashboard +ef796e473d647a173954d40c646904898d16bf2c831b3f87bea3e1e5b809888b admin-dashboard diff --git a/bin/admin-dashboard b/bin/admin-dashboard index e40d372..73b7350 100755 Binary files a/bin/admin-dashboard and b/bin/admin-dashboard differ diff --git a/bin/control-api b/bin/control-api index 8fedd04..a71bd5f 100755 Binary files a/bin/control-api and b/bin/control-api differ diff --git a/configs/control-api.example.yaml b/configs/control-api.example.yaml index b979ce0..2e3a4b3 100644 --- a/configs/control-api.example.yaml +++ b/configs/control-api.example.yaml @@ -100,7 +100,11 @@ targets: - "https://packages.ubuntu.com" # Inbound checks the 3 external-site probers run directly against each -# validator's currently-assigned floating IP. +# validator's currently-assigned floating IP. Like validators/sites/targets/ +# check_types above, this is only a bootstrap seed for a fresh, empty +# database; after that it's managed at runtime via PUT +# /api/v1/admin/config/inbound-checks (or the dashboard's /settings page) +# and this field is ignored. inbound_checks: ports: [22, 80, 443, 8080] icmp: true diff --git a/docs/API.md b/docs/API.md index a2db31d..2bcd190 100644 --- a/docs/API.md +++ b/docs/API.md @@ -231,7 +231,11 @@ target)` в рамках текущей попытки — безопасна и ] ``` -Пустой список `[]`, если сейчас нечего проверять. +Пустой список `[]`, если сейчас нечего проверять. `ports`/`icmp` — текущая +конфигурация типов проверок пробера, одна и та же для каждого IP в ответе; +управляется через +[`/api/v1/admin/config/inbound-checks`](#типы-проверок-пробера-apiv1adminconfiginbound-checks) +и меняется без рестарта control-api. ### `POST /api/v1/probers/{site_id}/results` @@ -477,6 +481,33 @@ fip_settle_seconds` в `control-api.yaml` — только одноразовы для пустой БД; дальше источник истины — сама база, менять значение нужно через `PUT` выше (или страницу `/settings` в дашборде). +### Типы проверок пробера: `/api/v1/admin/config/inbound-checks` + +Единый глобальный набор TCP-портов и флага ICMP, которые `prober` +проверяет на каждой настроенной площадке для каждого адреса в состоянии +`checking` — то же самое `ports`/`icmp`, что отдаётся в ответе `GET +/api/v1/probers/{site_id}/assignments` (см. +[«Методы для prober»](#методы-для-prober) выше). Один набор общий для всех +площадок; список [«площадок»](#площадки-apiv1adminconfigsites) отдельно +решает, *сколько* точек его применяют, а не что именно они проверяют. +Пустой список портов и `icmp: false` одновременно — допустимая +конфигурация: временно отключает inbound-проверки, не трогая список +`sites`. + +| Метод | Путь | Тело | Успех | Ошибки | +|---|---|---|---|---| +| GET | `/api/v1/admin/config/inbound-checks` | — | `{"ports":[...],"icmp":bool}` | | +| PUT | `/api/v1/admin/config/inbound-checks` | `{"ports":[...],"icmp":bool}` | `200` | `400`, если какой-то порт вне диапазона `1..65535` или порты повторяются | + +Изменение вступает в силу немедленно — как для следующего ответа `GET +/api/v1/probers/{site_id}/assignments`, так и для агрегации уже идущих +проверок (см. [«Управление площадками»](USAGE.md#управление-площадками-проберами) +в USAGE.md про тот же принцип для `sites`/`targets`/`check_types`). Как и +остальные разделы этой группы, YAML-поле `orchestrator.inbound_checks` в +`control-api.yaml` — только одноразовый bootstrap для пустой БД; дальше +источник истины — сама база, менять значение нужно через `PUT` выше (или +страницу `/settings` в дашборде). + ### Пример: конфигурация целиком через API, без единой строки в YAML ```bash diff --git a/docs/DASHBOARD.md b/docs/DASHBOARD.md index a9821cc..12bd1cc 100644 --- a/docs/DASHBOARD.md +++ b/docs/DASHBOARD.md @@ -55,7 +55,7 @@ admin-dashboard -config /etc/cloud-ip-validator/admin-dashboard.yaml | `/sites` | Три фиксированных слота площадок (1/2/3) — назначить/сменить/освободить `site_id`. | | `/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)). | +| `/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/USAGE.md b/docs/USAGE.md index 6e8f1e0..2a046c7 100644 --- a/docs/USAGE.md +++ b/docs/USAGE.md @@ -17,6 +17,7 @@ - [Просмотр деталей и истории по конкретному адресу](#просмотр-деталей-и-истории-по-конкретному-адресу) - [Управление валидаторами](#управление-валидаторами) - [Управление площадками (проберами)](#управление-площадками-проберами) +- [Управление типами проверок пробера](#управление-типами-проверок-пробера) - [Управление целями проверки](#управление-целями-проверки) - [Повторная проверка адреса](#повторная-проверка-адреса) - [Принудительная остановка проверки](#принудительная-остановка-проверки) @@ -275,6 +276,44 @@ curl -s -X DELETE http://:8080/api/v1/admin/config/sites/1 > поддерживаемый сценарий; *больше* трёх потребует доработки схемы > данных, одной правкой конфига не обойтись. +## Управление типами проверок пробера + +Список TCP-портов и флаг ICMP, которые `prober` проверяет на каждой +настроенной площадке — единый глобальный набор, общий для всех площадок +сразу (не то же самое, что список `sites` выше: `sites` решает, *сколько* +точек его применяют, а этот набор — *что именно* они проверяют). + +Посмотреть текущий набор: +```bash +curl -s http://:8080/api/v1/admin/config/inbound-checks | python3 -m json.tool +``` + +Изменить набор портов и/или ICMP (без перезапуска control-api — новое +значение сразу видно и следующему опросу пробера, и уже идущей агрегации): +```bash +curl -s -X PUT http://:8080/api/v1/admin/config/inbound-checks \ + -d '{"ports": [22, 80, 443, 8080], "icmp": true}' +``` + +Порты должны быть в диапазоне `1..65535` и не повторяться — иначе `400`. +Пустой список портов вместе с `"icmp": false` — штатный способ временно +отключить inbound-проверки целиком, не трогая список площадок: +```bash +curl -s -X PUT http://:8080/api/v1/admin/config/inbound-checks \ + -d '{"ports": [], "icmp": false}' +``` + +То же самое — на странице `/settings` дашборда, второй формой рядом с +паузой перед self-check (см. ниже): текстовое поле с портами через запятую +и чекбокс ICMP. + +> Правки через `orchestrator.inbound_checks` в `control-api.yaml` тоже +> поддерживаются, но только как bootstrap пустой базы данных при самом +> первом старте — как только в БД есть эта настройка (а она появляется +> сразу же при первом старте, значение по умолчанию — из YAML), YAML для +> этой секции игнорируется при всех последующих рестартах. Для стенда, +> который уже хоть раз запускался, используйте API выше. + ## Управление целями проверки Набор egress-целей (`targets`) и типов проверок (`check_types`, diff --git a/internal/dashboard/client.go b/internal/dashboard/client.go index e8f9578..4a7c41c 100644 --- a/internal/dashboard/client.go +++ b/internal/dashboard/client.go @@ -213,3 +213,16 @@ func (c *client) PutOrchestratorSettings(ctx context.Context, fipSettleSeconds i orchestratorSettingsDTO{FIPSettleSeconds: fipSettleSeconds}, &out) return out, err } + +func (c *client) GetInboundChecks(ctx context.Context) (inboundChecksDTO, error) { + var out inboundChecksDTO + err := c.do(ctx, http.MethodGet, "/api/v1/admin/config/inbound-checks", nil, &out) + return out, err +} + +func (c *client) PutInboundChecks(ctx context.Context, ports []int, icmp bool) (inboundChecksDTO, error) { + var out inboundChecksDTO + err := c.do(ctx, http.MethodPut, "/api/v1/admin/config/inbound-checks", + inboundChecksDTO{Ports: ports, ICMP: icmp}, &out) + return out, err +} diff --git a/internal/dashboard/dashboard_test.go b/internal/dashboard/dashboard_test.go index 6e518b5..c18440b 100644 --- a/internal/dashboard/dashboard_test.go +++ b/internal/dashboard/dashboard_test.go @@ -29,6 +29,8 @@ type fakeControlAPI struct { groups map[string][]string checkTypes map[string]checkTypeDTO fipSettleSeconds int + inboundPorts []int + inboundICMP bool } func newFakeControlAPI(t *testing.T) (*fakeControlAPI, string) { @@ -206,6 +208,27 @@ func (f *fakeControlAPI) handler() http.Handler { writeJSON(w, http.StatusOK, orchestratorSettingsDTO{FIPSettleSeconds: f.fipSettleSeconds}) }) + mux.HandleFunc("GET /api/v1/admin/config/inbound-checks", func(w http.ResponseWriter, r *http.Request) { + f.mu.Lock() + defer f.mu.Unlock() + writeJSON(w, http.StatusOK, inboundChecksDTO{Ports: f.inboundPorts, ICMP: f.inboundICMP}) + }) + mux.HandleFunc("PUT /api/v1/admin/config/inbound-checks", func(w http.ResponseWriter, r *http.Request) { + f.mu.Lock() + defer f.mu.Unlock() + var req inboundChecksDTO + _ = json.NewDecoder(r.Body).Decode(&req) + for _, p := range req.Ports { + if p < 1 || p > 65535 { + writeAPIErr(w, http.StatusBadRequest, "port out of range") + return + } + } + f.inboundPorts = req.Ports + f.inboundICMP = req.ICMP + writeJSON(w, http.StatusOK, inboundChecksDTO{Ports: f.inboundPorts, ICMP: f.inboundICMP}) + }) + mux.HandleFunc("GET /api/v1/admin/config/validators", func(w http.ResponseWriter, r *http.Request) { f.mu.Lock() defer f.mu.Unlock() diff --git a/internal/dashboard/dto.go b/internal/dashboard/dto.go index ad2af4c..fe6ba1b 100644 --- a/internal/dashboard/dto.go +++ b/internal/dashboard/dto.go @@ -128,3 +128,8 @@ type errorResponse struct { type orchestratorSettingsDTO struct { FIPSettleSeconds int `json:"fip_settle_seconds"` } + +type inboundChecksDTO struct { + Ports []int `json:"ports"` + ICMP bool `json:"icmp"` +} diff --git a/internal/dashboard/handlers_settings.go b/internal/dashboard/handlers_settings.go index 7a3e883..fa57816 100644 --- a/internal/dashboard/handlers_settings.go +++ b/internal/dashboard/handlers_settings.go @@ -4,16 +4,22 @@ import ( "fmt" "net/http" "strconv" + "strings" ) type settingsPageData struct { PageData Settings orchestratorSettingsDTO + Inbound inboundChecksDTO } func (s *Server) handleSettingsPage(w http.ResponseWriter, r *http.Request) { settings, err := s.CA.GetOrchestratorSettings(r.Context()) - data := settingsPageData{Settings: settings} + inbound, inboundErr := s.CA.GetInboundChecks(r.Context()) + if err == nil { + err = inboundErr + } + data := settingsPageData{Settings: settings, Inbound: inbound} data.ActiveNav = "settings" data.Banner = bannerFor(err) s.renderPage(w, "settings_page", data) @@ -28,7 +34,11 @@ func (s *Server) renderSettingsForm(w http.ResponseWriter, r *http.Request, acti if actionErr == nil { actionErr = getErr } - s.renderFragment(w, "settings_form", settingsPageData{Settings: settings}, actionErr) + inbound, inboundErr := s.CA.GetInboundChecks(r.Context()) + if actionErr == nil { + actionErr = inboundErr + } + s.renderFragment(w, "settings_form", settingsPageData{Settings: settings, Inbound: inbound}, actionErr) } func (s *Server) handleSettingsPut(w http.ResponseWriter, r *http.Request) { @@ -44,3 +54,34 @@ func (s *Server) handleSettingsPut(w http.ResponseWriter, r *http.Request) { _, err = s.CA.PutOrchestratorSettings(r.Context(), seconds) s.renderSettingsForm(w, r, err) } + +// handleInboundChecksPut parses the comma-separated ports field of the +// second form on /settings and saves it via PUT +// /api/v1/admin/config/inbound-checks. A separate form/handler from +// fip_settle_seconds above — the two settings are unrelated and shouldn't +// share one submit. +func (s *Server) handleInboundChecksPut(w http.ResponseWriter, r *http.Request) { + if err := r.ParseForm(); err != nil { + s.renderSettingsForm(w, r, fmt.Errorf("invalid form: %w", err)) + return + } + raw := strings.TrimSpace(r.PostFormValue("ports")) + var ports []int + if raw != "" { + for _, part := range strings.Split(raw, ",") { + part = strings.TrimSpace(part) + if part == "" { + continue + } + p, err := strconv.Atoi(part) + if err != nil { + s.renderSettingsForm(w, r, &apiErr{Status: http.StatusBadRequest, Message: "порты должны быть целыми числами через запятую"}) + return + } + ports = append(ports, p) + } + } + icmp := r.PostFormValue("icmp") == "on" + _, err := s.CA.PutInboundChecks(r.Context(), ports, icmp) + s.renderSettingsForm(w, r, err) +} diff --git a/internal/dashboard/handlers_test.go b/internal/dashboard/handlers_test.go index 6b27cbb..15704ef 100644 --- a/internal/dashboard/handlers_test.go +++ b/internal/dashboard/handlers_test.go @@ -1,6 +1,7 @@ package dashboard import ( + "reflect" "strings" "testing" "time" @@ -207,6 +208,56 @@ func TestSettingsGetAndPut(t *testing.T) { } } +func TestSettingsPageShowsInboundChecks(t *testing.T) { + fake, caURL := newFakeControlAPI(t) + fake.inboundPorts = []int{22, 443} + fake.inboundICMP = true + ts := newTestServer(t, caURL) + + page := get(t, ts, "/settings") + if !strings.Contains(page, `value="22, 443"`) { + t.Fatalf("expected seeded ports in the form, got:\n%s", page) + } + if !strings.Contains(page, "checked") { + t.Fatalf("expected icmp checkbox checked, got:\n%s", page) + } +} + +func TestInboundChecksPutUpdatesForm(t *testing.T) { + fake, caURL := newFakeControlAPI(t) + ts := newTestServer(t, caURL) + + body := postForm(t, ts, "PUT", "/settings/inbound-checks", map[string][]string{"ports": {"8080, 8443"}}) + if !strings.Contains(body, `value="8080, 8443"`) { + t.Fatalf("expected updated ports in re-rendered form, got:\n%s", body) + } + if !reflect.DeepEqual(fake.inboundPorts, []int{8080, 8443}) || fake.inboundICMP { + t.Fatalf("expected fake control-api updated to {[8080 8443] false}, got ports=%v icmp=%v", fake.inboundPorts, fake.inboundICMP) + } + + // icmp checkbox checked, empty ports field. + body = postForm(t, ts, "PUT", "/settings/inbound-checks", map[string][]string{"ports": {""}, "icmp": {"on"}}) + if len(fake.inboundPorts) != 0 || !fake.inboundICMP { + t.Fatalf("expected fake control-api updated to {[] true}, got ports=%v icmp=%v", fake.inboundPorts, fake.inboundICMP) + } + if !strings.Contains(body, "checked") { + t.Fatalf("expected icmp checkbox checked in re-rendered form, got:\n%s", body) + } + + // A malformed ports field is caught by the dashboard itself. + body = postForm(t, ts, "PUT", "/settings/inbound-checks", map[string][]string{"ports": {"22, not-a-port"}}) + if !strings.Contains(body, "alert-warning") { + t.Fatalf("expected client error banner for malformed ports, got:\n%s", body) + } + + // A control-api validation error (port out of range) surfaces via the + // banner too. + body = postForm(t, ts, "PUT", "/settings/inbound-checks", map[string][]string{"ports": {"70000"}}) + if !strings.Contains(body, "alert-warning") { + t.Fatalf("expected client error banner for out-of-range port, got:\n%s", body) + } +} + func TestIPsPageShowsSettleBadge(t *testing.T) { fake, caURL := newFakeControlAPI(t) now := time.Now() diff --git a/internal/dashboard/render.go b/internal/dashboard/render.go index 79c2655..bd1f764 100644 --- a/internal/dashboard/render.go +++ b/internal/dashboard/render.go @@ -4,6 +4,7 @@ import ( "errors" "html/template" "net/http" + "strconv" "strings" "time" ) @@ -87,12 +88,23 @@ func derefStr(s *string) string { return *s } +// joinInts renders a []int as a comma-separated string, for pre-filling the +// inbound-checks ports form field on /settings. +func joinInts(ints []int, sep string) string { + parts := make([]string, len(ints)) + for i, n := range ints { + parts[i] = strconv.Itoa(n) + } + return strings.Join(parts, sep) +} + var funcMap = template.FuncMap{ "ipBadge": ipBadge, "validatorBadge": 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 f9cb3d1..a0a3850 100644 --- a/internal/dashboard/routes.go +++ b/internal/dashboard/routes.go @@ -38,6 +38,7 @@ func (s *Server) routes(mux *http.ServeMux) { mux.HandleFunc("GET /settings", s.handleSettingsPage) mux.HandleFunc("PUT /settings", s.handleSettingsPut) + mux.HandleFunc("PUT /settings/inbound-checks", s.handleInboundChecksPut) mux.Handle("GET /static/", http.StripPrefix("/static/", http.FileServerFS(staticSubFS()))) } diff --git a/internal/dashboard/templates/settings.html b/internal/dashboard/templates/settings.html index 9ffcf67..eed6583 100644 --- a/internal/dashboard/templates/settings.html +++ b/internal/dashboard/templates/settings.html @@ -45,4 +45,26 @@ + +
+
+

Типы проверок, которые прогоняет prober на каждой +настроенной площадке для каждого проверяемого адреса — глобальный набор, общий для всех площадок. Отдельно от +списка площадок (/sites), который решает, сколько площадок опрашивать, а не +что именно они проверяют. Пустой список портов и выключенный ICMP — штатный способ временно свести +inbound-проверки к нулю, не трогая список площадок.

+
+
+
+ + +
+
+ +
+ +
+
+
+
{{end}} diff --git a/internal/db/bootstrap.go b/internal/db/bootstrap.go index 4e3dc02..a8be8d7 100644 --- a/internal/db/bootstrap.go +++ b/internal/db/bootstrap.go @@ -2,6 +2,7 @@ package db import ( "context" + "encoding/json" "fmt" "cloudipvalidator/internal/config" @@ -36,6 +37,9 @@ func (d *DB) BootstrapFromConfig(ctx context.Context, cfg *config.ControlAPI) er if err := d.bootstrapSettings(ctx, cfg.Orchestrator.FIPSettleSeconds); err != nil { return fmt.Errorf("bootstrap settings: %w", err) } + if err := d.bootstrapInboundChecks(ctx, cfg.Inbound.Ports, cfg.Inbound.ICMP); err != nil { + return fmt.Errorf("bootstrap inbound checks: %w", err) + } if err := d.SeedQueue(ctx, cfg.IPAddresses); err != nil { return fmt.Errorf("seed ip queue: %w", err) } @@ -120,3 +124,25 @@ func (d *DB) bootstrapSettings(ctx context.Context, fipSettleSeconds int) error `, fipSettleSeconds, now, now) return err } + +func (d *DB) bootstrapInboundChecks(ctx context.Context, ports []int, icmp bool) error { + var count int + if err := d.QueryRowContext(ctx, `SELECT COUNT(*) FROM inbound_checks_settings`).Scan(&count); err != nil { + return err + } + if count > 0 { + return nil + } + if ports == nil { + ports = []int{} + } + payload, err := json.Marshal(ports) + if err != nil { + return err + } + now := timeToDB(Now()) + _, err = d.ExecContext(ctx, ` + INSERT INTO inbound_checks_settings (id, ports, icmp, created_at, updated_at) VALUES (1, ?, ?, ?, ?) + `, string(payload), icmp, now, now) + return err +} diff --git a/internal/db/db.go b/internal/db/db.go index afdc453..da38304 100644 --- a/internal/db/db.go +++ b/internal/db/db.go @@ -22,6 +22,9 @@ var dynamicConfigSchema string //go:embed migrations/0003_fip_settle_delay.sql var fipSettleDelaySchema string +//go:embed migrations/0004_inbound_checks_admin.sql +var inboundChecksAdminSchema 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 @@ -33,6 +36,7 @@ var migrations = []struct { {1, initSchema}, {2, dynamicConfigSchema}, {3, fipSettleDelaySchema}, + {4, inboundChecksAdminSchema}, } type DB struct { diff --git a/internal/db/migrations/0004_inbound_checks_admin.sql b/internal/db/migrations/0004_inbound_checks_admin.sql new file mode 100644 index 0000000..ad9112c --- /dev/null +++ b/internal/db/migrations/0004_inbound_checks_admin.sql @@ -0,0 +1,16 @@ +-- Support for admin-API-configurable prober check types (see docs/USAGE.md, +-- "Управление типами проверок пробера"). Previously orchestrator.inbound_checks +-- was a YAML-only value baked into the process at startup. + +-- Singleton row, deliberately separate from `settings`: unlike +-- fip_settle_seconds, inbound checks need no cross-field validation against +-- other orchestrator config, and ports is a list rather than a scalar, so it +-- doesn't fit the "add a column to settings" guidance left in +-- 0003_fip_settle_delay.sql. +CREATE TABLE inbound_checks_settings ( + id INTEGER PRIMARY KEY CHECK (id = 1), + ports TEXT NOT NULL, -- JSON array of int, e.g. "[22,80,443,8080]" + icmp BOOLEAN NOT NULL DEFAULT 0, + created_at TIMESTAMP NOT NULL, + updated_at TIMESTAMP NOT NULL +); diff --git a/internal/db/models.go b/internal/db/models.go index ad60176..58e8263 100644 --- a/internal/db/models.go +++ b/internal/db/models.go @@ -171,3 +171,13 @@ type Settings struct { CreatedAt time.Time UpdatedAt time.Time } + +// InboundChecksSettings is the singleton row describing what the prober +// checks on every site for every in-flight IP (TCP ports + optional ICMP). +// Admin-configurable at runtime (see queries_inbound.go). +type InboundChecksSettings struct { + Ports []int + ICMP bool + CreatedAt time.Time + UpdatedAt time.Time +} diff --git a/internal/db/queries_inbound.go b/internal/db/queries_inbound.go new file mode 100644 index 0000000..9cfaf30 --- /dev/null +++ b/internal/db/queries_inbound.go @@ -0,0 +1,65 @@ +package db + +import ( + "context" + "encoding/json" + "fmt" +) + +// GetInboundChecks returns the singleton inbound-checks settings row. +// BootstrapFromConfig guarantees it exists before any other code path can +// observe it, so sql.ErrNoRows here would indicate a bootstrap bug, not a +// normal condition. +func (d *DB) GetInboundChecks(ctx context.Context) (InboundChecksSettings, error) { + var s InboundChecksSettings + var portsJSON, createdAt, updatedAt string + err := d.QueryRowContext(ctx, ` + SELECT ports, icmp, created_at, updated_at FROM inbound_checks_settings WHERE id=1 + `).Scan(&portsJSON, &s.ICMP, &createdAt, &updatedAt) + if err != nil { + return InboundChecksSettings{}, err + } + if err := json.Unmarshal([]byte(portsJSON), &s.Ports); err != nil { + return InboundChecksSettings{}, fmt.Errorf("decode inbound check ports: %w", err) + } + if s.Ports == nil { + s.Ports = []int{} + } + if s.CreatedAt, err = dbToTime(createdAt); err != nil { + return InboundChecksSettings{}, err + } + if s.UpdatedAt, err = dbToTime(updatedAt); err != nil { + return InboundChecksSettings{}, err + } + return s, nil +} + +// SetInboundChecks persists a new prober check configuration. Validation is +// purely local (port range, no duplicates) — unlike fip_settle_seconds, +// there's no cross-field dependency on other orchestrator config, so this +// is called directly by the HTTP handler without going through +// orchestrator.Orchestrator. +func (d *DB) SetInboundChecks(ctx context.Context, ports []int, icmp bool) error { + seen := make(map[int]bool, len(ports)) + for _, p := range ports { + if p < 1 || p > 65535 { + return fmt.Errorf("port %d out of range 1..65535: %w", p, ErrValidation) + } + if seen[p] { + return fmt.Errorf("duplicate port %d: %w", p, ErrValidation) + } + seen[p] = true + } + if ports == nil { + ports = []int{} + } + payload, err := json.Marshal(ports) + if err != nil { + return err + } + now := timeToDB(Now()) + _, err = d.ExecContext(ctx, ` + UPDATE inbound_checks_settings SET ports=?, icmp=?, updated_at=? WHERE id=1 + `, string(payload), icmp, now) + return err +} diff --git a/internal/db/queries_inbound_test.go b/internal/db/queries_inbound_test.go new file mode 100644 index 0000000..c7ca7da --- /dev/null +++ b/internal/db/queries_inbound_test.go @@ -0,0 +1,103 @@ +package db + +import ( + "errors" + "reflect" + "testing" + + "cloudipvalidator/internal/config" +) + +func TestBootstrapInboundChecksSeedsOnceFromConfig(t *testing.T) { + d, ctx := newTestDB(t) + + cfg := &config.ControlAPI{Inbound: config.InboundConfig{Ports: []int{22, 80}, ICMP: true}} + if err := d.BootstrapFromConfig(ctx, cfg); err != nil { + t.Fatalf("first bootstrap: %v", err) + } + settings, err := d.GetInboundChecks(ctx) + if err != nil { + t.Fatalf("get inbound checks: %v", err) + } + if !reflect.DeepEqual(settings.Ports, []int{22, 80}) || !settings.ICMP { + t.Fatalf("expected seeded {[22 80] true}, got %+v", settings) + } + + // A second bootstrap with a different YAML value must not overwrite the + // now-non-empty table — same "DB is source of truth once seeded" + // semantics as validators/sites/targets/check_types/settings. + cfg2 := &config.ControlAPI{Inbound: config.InboundConfig{Ports: []int{443}, ICMP: false}} + if err := d.BootstrapFromConfig(ctx, cfg2); err != nil { + t.Fatalf("second bootstrap: %v", err) + } + settings, err = d.GetInboundChecks(ctx) + if err != nil { + t.Fatalf("get inbound checks after second bootstrap: %v", err) + } + if !reflect.DeepEqual(settings.Ports, []int{22, 80}) || !settings.ICMP { + t.Fatalf("expected YAML to be ignored on non-empty table, got %+v", settings) + } +} + +func TestSetInboundChecksRejectsPortOutOfRange(t *testing.T) { + d, ctx := newTestDB(t) + if err := d.BootstrapFromConfig(ctx, &config.ControlAPI{}); err != nil { + t.Fatalf("bootstrap: %v", err) + } + for _, p := range []int{0, -1, 65536, 70000} { + if err := d.SetInboundChecks(ctx, []int{p}, false); !errors.Is(err, ErrValidation) { + t.Fatalf("port %d: expected ErrValidation, got %v", p, err) + } + } +} + +func TestSetInboundChecksRejectsDuplicatePorts(t *testing.T) { + d, ctx := newTestDB(t) + if err := d.BootstrapFromConfig(ctx, &config.ControlAPI{}); err != nil { + t.Fatalf("bootstrap: %v", err) + } + if err := d.SetInboundChecks(ctx, []int{22, 80, 22}, false); !errors.Is(err, ErrValidation) { + t.Fatalf("expected ErrValidation for duplicate port, got %v", err) + } +} + +func TestGetSetInboundChecksRoundTrip(t *testing.T) { + d, ctx := newTestDB(t) + if err := d.BootstrapFromConfig(ctx, &config.ControlAPI{}); err != nil { + t.Fatalf("bootstrap: %v", err) + } + before, err := d.GetInboundChecks(ctx) + if err != nil { + t.Fatalf("get inbound checks: %v", err) + } + if len(before.Ports) != 0 || before.ICMP { + t.Fatalf("expected default {[] false}, got %+v", before) + } + + if err := d.SetInboundChecks(ctx, []int{22, 443, 8080}, true); err != nil { + t.Fatalf("set inbound checks: %v", err) + } + after, err := d.GetInboundChecks(ctx) + if err != nil { + t.Fatalf("get inbound checks after set: %v", err) + } + if !reflect.DeepEqual(after.Ports, []int{22, 443, 8080}) || !after.ICMP { + t.Fatalf("expected {[22 443 8080] true}, got %+v", after) + } + if after.UpdatedAt.Before(before.UpdatedAt) { + t.Fatalf("expected updated_at not to go backwards, before=%v after=%v", before.UpdatedAt, after.UpdatedAt) + } + + // Setting an empty port list + icmp:false is legal — it's the "disable + // all inbound checks without touching sites" configuration. + if err := d.SetInboundChecks(ctx, nil, false); err != nil { + t.Fatalf("set empty inbound checks: %v", err) + } + cleared, err := d.GetInboundChecks(ctx) + if err != nil { + t.Fatalf("get inbound checks after clear: %v", err) + } + if len(cleared.Ports) != 0 || cleared.ICMP { + t.Fatalf("expected {[] false} after clearing, got %+v", cleared) + } +} diff --git a/internal/httpapi/dto_admin.go b/internal/httpapi/dto_admin.go index 3f8209b..afd1fe0 100644 --- a/internal/httpapi/dto_admin.go +++ b/internal/httpapi/dto_admin.go @@ -80,3 +80,11 @@ type putCheckTypeRequest struct { type orchestratorSettingsDTO struct { FIPSettleSeconds int `json:"fip_settle_seconds"` } + +// inboundChecksDTO doubles as both the GET response and the PUT request +// body for /api/v1/admin/config/inbound-checks. Same field shape as +// proberAssignment.Ports/ICMP (dto.go). +type inboundChecksDTO struct { + Ports []int `json:"ports"` + ICMP bool `json:"icmp"` +} diff --git a/internal/httpapi/handlers_config.go b/internal/httpapi/handlers_config.go index df513c1..a445ed4 100644 --- a/internal/httpapi/handlers_config.go +++ b/internal/httpapi/handlers_config.go @@ -209,3 +209,30 @@ func (s *Server) handleConfigPutOrchestratorSettings(w http.ResponseWriter, r *h } writeJSON(w, http.StatusOK, orchestratorSettingsDTO{FIPSettleSeconds: req.FIPSettleSeconds}) } + +// --- prober inbound checks --- + +func (s *Server) handleConfigGetInboundChecks(w http.ResponseWriter, r *http.Request) { + settings, err := s.DB.GetInboundChecks(r.Context()) + if err != nil { + writeDBError(w, err) + return + } + writeJSON(w, http.StatusOK, inboundChecksDTO{Ports: settings.Ports, ICMP: settings.ICMP}) +} + +// handleConfigPutInboundChecks writes directly to the DB (unlike the +// orchestrator-settings PUT above) because port-range/duplicate validation +// is purely local — no cross-field dependency on other orchestrator config. +func (s *Server) handleConfigPutInboundChecks(w http.ResponseWriter, r *http.Request) { + var req inboundChecksDTO + if err := readJSON(r, &req); err != nil { + writeError(w, http.StatusBadRequest, "invalid body: "+err.Error()) + return + } + if err := s.DB.SetInboundChecks(r.Context(), req.Ports, req.ICMP); err != nil { + writeDBError(w, err) + return + } + writeJSON(w, http.StatusOK, inboundChecksDTO{Ports: req.Ports, ICMP: req.ICMP}) +} diff --git a/internal/httpapi/handlers_config_test.go b/internal/httpapi/handlers_config_test.go index 13ac30c..d38c387 100644 --- a/internal/httpapi/handlers_config_test.go +++ b/internal/httpapi/handlers_config_test.go @@ -500,3 +500,126 @@ func TestFIPSettleDelayGatesAssignmentEndpoint(t *testing.T) { t.Fatalf("expected 200 after settle window elapsed, status=%d body=%s", resp.StatusCode, body) } } + +// TestInboundChecksGetPut proves the prober check config round-trips +// through GET/PUT /api/v1/admin/config/inbound-checks. +func TestInboundChecksGetPut(t *testing.T) { + fc, _, _, _ := newConfigTestHarness(t) + + resp, body := fc.do(http.MethodGet, "/api/v1/admin/config/inbound-checks", nil) + if resp.StatusCode != http.StatusOK { + t.Fatalf("get inbound checks: status=%d body=%s", resp.StatusCode, body) + } + var got inboundChecksDTO + if err := json.Unmarshal(body, &got); err != nil { + t.Fatalf("unmarshal get response: %v", err) + } + // newConfigTestHarness seeds Ports:[22,80], ICMP:true. + if len(got.Ports) != 2 || !got.ICMP { + t.Fatalf("expected seeded {[22 80] true}, got %+v", got) + } + + resp, body = fc.do(http.MethodPut, "/api/v1/admin/config/inbound-checks", inboundChecksDTO{Ports: []int{443, 8080}, ICMP: false}) + if resp.StatusCode != http.StatusOK { + t.Fatalf("put inbound checks: status=%d body=%s", resp.StatusCode, body) + } + if err := json.Unmarshal(body, &got); err != nil { + t.Fatalf("unmarshal put response: %v", err) + } + if len(got.Ports) != 2 || got.ICMP { + t.Fatalf("expected {[443 8080] false}, got %+v", got) + } + + resp, body = fc.do(http.MethodGet, "/api/v1/admin/config/inbound-checks", nil) + if resp.StatusCode != http.StatusOK { + t.Fatalf("get inbound checks after put: status=%d body=%s", resp.StatusCode, body) + } + if err := json.Unmarshal(body, &got); err != nil { + t.Fatalf("unmarshal get-after-put response: %v", err) + } + if got.Ports[0] != 443 || got.Ports[1] != 8080 || got.ICMP { + t.Fatalf("expected {[443 8080] false} to persist, got %+v", got) + } +} + +// TestInboundChecksPutValidation proves out-of-range and duplicate ports +// are rejected with 400. +func TestInboundChecksPutValidation(t *testing.T) { + fc, _, _, _ := newConfigTestHarness(t) + + resp, body := fc.do(http.MethodPut, "/api/v1/admin/config/inbound-checks", inboundChecksDTO{Ports: []int{0}, ICMP: false}) + if resp.StatusCode != http.StatusBadRequest { + t.Fatalf("expected 400 for port 0, status=%d body=%s", resp.StatusCode, body) + } + + resp, body = fc.do(http.MethodPut, "/api/v1/admin/config/inbound-checks", inboundChecksDTO{Ports: []int{70000}, ICMP: false}) + if resp.StatusCode != http.StatusBadRequest { + t.Fatalf("expected 400 for port 70000, status=%d body=%s", resp.StatusCode, body) + } + + resp, body = fc.do(http.MethodPut, "/api/v1/admin/config/inbound-checks", inboundChecksDTO{Ports: []int{22, 22}, ICMP: false}) + if resp.StatusCode != http.StatusBadRequest { + t.Fatalf("expected 400 for duplicate port, status=%d body=%s", resp.StatusCode, body) + } +} + +// TestInboundChecksReflectedInProberAssignmentsWithoutRestart proves the +// fix to the formerly-static handlers_prober.go read: a PUT to +// /api/v1/admin/config/inbound-checks changes what GET +// /api/v1/probers/{site_id}/assignments hands back to an already-registered +// prober, for an IP already in `checking`, with no control-api restart. +func TestInboundChecksReflectedInProberAssignmentsWithoutRestart(t *testing.T) { + fc, _, orch, mock := newConfigTestHarness(t) + ctx := context.Background() + mock.Seed("fip-1", "9.9.9.9", "svc-project") + + fc.do(http.MethodPut, "/api/v1/admin/config/sites/1", putSiteRequest{SiteID: "site-1"}) + fc.do(http.MethodPost, "/api/v1/admin/config/validators", createValidatorRequest{ValidatorID: "validator-1", OSPortID: "port-1"}) + fc.do(http.MethodPost, "/api/v1/agents/register", registerAgentRequest{ValidatorID: "validator-1"}) + fc.do(http.MethodPost, "/api/v1/admin/ips", submitIPsRequest{Addresses: []string{"9.9.9.9"}}) + + orch.Tick(ctx) // claim + associate -> awaiting_self_check + + resp, body := fc.do(http.MethodGet, "/api/v1/agents/validator-1/assignment", nil) + if resp.StatusCode != http.StatusOK { + t.Fatalf("assignment: status=%d body=%s", resp.StatusCode, body) + } + var assignment assignmentResponse + if err := json.Unmarshal(body, &assignment); err != nil { + t.Fatalf("unmarshal assignment: %v", err) + } + fc.do(http.MethodPost, "/api/v1/agents/validator-1/self-check", selfCheckRequest{ + IPID: assignment.IPID, DetectedEgress: "9.9.9.9", Success: true, Detail: "matched", + }) + // Now item.State == "checking" — a prober assignment target. + + fc.do(http.MethodPost, "/api/v1/probers/register", registerProberRequest{SiteID: "site-1"}) + + resp, body = fc.do(http.MethodGet, "/api/v1/probers/site-1/assignments", nil) + if resp.StatusCode != http.StatusOK { + t.Fatalf("prober assignments: status=%d body=%s", resp.StatusCode, body) + } + var assignments []proberAssignment + if err := json.Unmarshal(body, &assignments); err != nil { + t.Fatalf("unmarshal assignments: %v", err) + } + if len(assignments) != 1 || len(assignments[0].Ports) != 2 || !assignments[0].ICMP { + t.Fatalf("expected seeded {[22 80] true} before PUT, got %+v", assignments) + } + + resp, body = fc.do(http.MethodPut, "/api/v1/admin/config/inbound-checks", inboundChecksDTO{Ports: []int{8080}, ICMP: false}) + if resp.StatusCode != http.StatusOK { + t.Fatalf("put inbound checks: status=%d body=%s", resp.StatusCode, body) + } + + resp, body = fc.do(http.MethodGet, "/api/v1/probers/site-1/assignments", nil) + if resp.StatusCode != http.StatusOK { + t.Fatalf("prober assignments after put: status=%d body=%s", resp.StatusCode, body) + } + if err := json.Unmarshal(body, &assignments); err != nil { + t.Fatalf("unmarshal assignments after put: %v", err) + } + if len(assignments) != 1 || len(assignments[0].Ports) != 1 || assignments[0].Ports[0] != 8080 || assignments[0].ICMP { + t.Fatalf("expected updated {[8080] false} without restart, got %+v", assignments) + } +} diff --git a/internal/httpapi/handlers_prober.go b/internal/httpapi/handlers_prober.go index b0febcd..b7de846 100644 --- a/internal/httpapi/handlers_prober.go +++ b/internal/httpapi/handlers_prober.go @@ -45,11 +45,16 @@ func (s *Server) handleProberAssignments(w http.ResponseWriter, r *http.Request) writeError(w, http.StatusInternalServerError, err.Error()) return } + inbound, err := s.DB.GetInboundChecks(r.Context()) + if err != nil { + writeError(w, http.StatusInternalServerError, err.Error()) + return + } out := make([]proberAssignment, 0, len(items)) for _, item := range items { out = append(out, proberAssignment{ IPID: item.ID, IPAddress: item.IPAddress, - Ports: s.Orch.Inbound.Ports, ICMP: s.Orch.Inbound.ICMP, + Ports: inbound.Ports, ICMP: inbound.ICMP, }) } writeJSON(w, http.StatusOK, out) diff --git a/internal/httpapi/routes.go b/internal/httpapi/routes.go index 9e02761..6a13e80 100644 --- a/internal/httpapi/routes.go +++ b/internal/httpapi/routes.go @@ -46,4 +46,7 @@ func (s *Server) routes(mux *http.ServeMux) { mux.HandleFunc("GET /api/v1/admin/config/orchestrator", s.handleConfigGetOrchestratorSettings) mux.HandleFunc("PUT /api/v1/admin/config/orchestrator", s.handleConfigPutOrchestratorSettings) + + mux.HandleFunc("GET /api/v1/admin/config/inbound-checks", s.handleConfigGetInboundChecks) + mux.HandleFunc("PUT /api/v1/admin/config/inbound-checks", s.handleConfigPutInboundChecks) } diff --git a/internal/orchestrator/orchestrator.go b/internal/orchestrator/orchestrator.go index 415016f..29d64b8 100644 --- a/internal/orchestrator/orchestrator.go +++ b/internal/orchestrator/orchestrator.go @@ -31,28 +31,28 @@ type CheckConfig struct { } type Orchestrator struct { - DB *db.DB - OS openstack.FloatingIPClient - Cfg config.OrchestratorConfig - Agg config.AggregationConfig - Inbound config.InboundConfig - Log *slog.Logger + DB *db.DB + OS openstack.FloatingIPClient + Cfg config.OrchestratorConfig + Agg config.AggregationConfig + Log *slog.Logger } -// New constructs an Orchestrator. Egress check types/targets and prober -// sites are no longer taken from cfg — they're read from the database on -// every use (see AssignmentForValidator, expectedCheckCount, -// isReadyToAggregate) so admin API changes to them take effect without a -// restart. cfg.Validators/.Sites/.CheckTypes/.Targets/.IPAddresses are only -// consulted once, at process startup, by db.BootstrapFromConfig. +// New constructs an Orchestrator. Egress check types/targets, prober sites, +// and prober inbound check config (ports/icmp) are no longer taken from +// cfg — they're read from the database on every use (see +// AssignmentForValidator, expectedCheckCount, isReadyToAggregate, +// httpapi.handleProberAssignments) so admin API changes to them take effect +// without a restart. cfg.Validators/.Sites/.CheckTypes/.Targets/.Inbound/ +// .IPAddresses are only consulted once, at process startup, by +// db.BootstrapFromConfig. func New(d *db.DB, osClient openstack.FloatingIPClient, cfg *config.ControlAPI, log *slog.Logger) *Orchestrator { return &Orchestrator{ - DB: d, - OS: osClient, - Cfg: cfg.Orchestrator, - Agg: cfg.Aggregation, - Inbound: cfg.Inbound, - Log: log, + DB: d, + OS: osClient, + Cfg: cfg.Orchestrator, + Agg: cfg.Aggregation, + Log: log, } } @@ -509,10 +509,10 @@ func (o *Orchestrator) aggregateAndRelease(ctx context.Context, item db.IPQueueI // expectedCheckCount is the number of check rows a fully-reported IP should // have: one per (egress check-type x target) plus one per (site x inbound -// port/icmp probe). Reads the current check_types/targets/sites from the -// database, so a config change between assignment and aggregation is -// reflected in this specific aggregation (see the "accepted tradeoff" note -// in docs/PLAN_API_CONFIG_MANAGEMENT.md). +// port/icmp probe). Reads the current check_types/targets/sites/inbound +// checks from the database, so a config change between assignment and +// aggregation is reflected in this specific aggregation (see the "accepted +// tradeoff" note in docs/PLAN_API_CONFIG_MANAGEMENT.md). func (o *Orchestrator) expectedCheckCount(ctx context.Context) (int, error) { resolved, err := o.DB.ListResolvedCheckTypes(ctx) if err != nil { @@ -526,8 +526,12 @@ func (o *Orchestrator) expectedCheckCount(ctx context.Context) (int, error) { if err != nil { return 0, err } - inboundPerSite := len(o.Inbound.Ports) - if o.Inbound.ICMP { + inbound, err := o.DB.GetInboundChecks(ctx) + if err != nil { + return 0, err + } + inboundPerSite := len(inbound.Ports) + if inbound.ICMP { inboundPerSite++ } return egress + inboundPerSite*len(sites), nil diff --git a/internal/orchestrator/orchestrator_test.go b/internal/orchestrator/orchestrator_test.go index 70543db..8c34662 100644 --- a/internal/orchestrator/orchestrator_test.go +++ b/internal/orchestrator/orchestrator_test.go @@ -575,3 +575,52 @@ func TestInboundChecksPartialSites(t *testing.T) { t.Fatalf("expected pass, got %s", ip.OverallResult) } } + +// TestExpectedCheckCountReflectsInboundChecksConfigChange confirms +// expectedCheckCount reads inbound_checks (ports/icmp) from the database on +// every use, not a cached copy — same "accepted tradeoff" already true of +// check_types/targets/sites (see expectedCheckCount's doc comment). +func TestExpectedCheckCountReflectsInboundChecksConfigChange(t *testing.T) { + ctx := context.Background() + sites := []config.SiteConfig{{SiteID: "site-1", Index: 1}} + o, d, mock := newTestOrchestratorWithSites(t, 180, sites) + // newTestOrchestratorWithSites seeds Inbound: {Ports: [22, 80], ICMP: true} + // (3 inbound checks expected per site). + 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) + + // Shrink the inbound check config to a single TCP port, no ICMP — an + // admin API change made between assignment and aggregation. + if err := d.SetInboundChecks(ctx, []int{22}, false); err != nil { + t.Fatalf("set inbound checks: %v", err) + } + + // Report only the one now-expected inbound check. + _ = o.RecordCheck(ctx, db.Check{ + IPID: ip.ID, IPAddress: ip.IPAddress, AttemptNumber: ip.AttemptNumber, + Source: db.InboundSource(1), CheckType: "tcp-22", Target: ip.IPAddress, Success: true, CheckedAt: db.Now(), + }) + _ = o.MarkSiteComplete(ctx, ip.ID, 1) + + o.Tick(ctx) + ip, _ = d.GetIP(ctx, ip.ID) + if ip.State != db.IPDone { + t.Fatalf("expected done once the (shrunk) inbound config was fully reported, got %s", ip.State) + } + if ip.OverallResult != db.ResultPass { + t.Fatalf("expected pass, got %s", ip.OverallResult) + } +}