Manage external site-prober actions va API and Dashboard
This commit is contained in:
1 parent
f22ad569b6
commit
42f584dd3c
28 files changed
+717
-32
No files matched your search
+2
-2
@@ -1,4 +1,4 @@
|
||||
b5d1ae5c625c1ab12a42b91d0e7dbfe5eb15897fbc46e1b94db21aca5f84e77f control-api
|
||||
5703aa9367f814fd1c02ab8ee821f4ebea28934387b084e0f843efe1505530b7 control-api
|
||||
48c9b99fa88be751d9badfba8b7d80f894e326a743a90e85478f1d2251ce4ab6 validator-agent
|
||||
43fb660b78179b204b57388241a508345881aa16f3531d77b8fd75af60e7f2d8 prober
|
||||
4e1e550bf026b9f9b4aa60a9bf38e1e409693f7ed8d4fd1853175f882ba349ca admin-dashboard
|
||||
ef796e473d647a173954d40c646904898d16bf2c831b3f87bea3e1e5b809888b admin-dashboard
|
||||
Binary file not shown.
Binary file not shown.
@@ -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
|
||||
|
||||
+32
-1
@@ -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
|
||||
|
||||
+1
-1
@@ -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#управление-типами-проверок-пробера)). |
|
||||
|
||||
### «Текущая» и «последняя завершённая» проверка
|
||||
|
||||
|
||||
@@ -17,6 +17,7 @@
|
||||
- [Просмотр деталей и истории по конкретному адресу](#просмотр-деталей-и-истории-по-конкретному-адресу)
|
||||
- [Управление валидаторами](#управление-валидаторами)
|
||||
- [Управление площадками (проберами)](#управление-площадками-проберами)
|
||||
- [Управление типами проверок пробера](#управление-типами-проверок-пробера)
|
||||
- [Управление целями проверки](#управление-целями-проверки)
|
||||
- [Повторная проверка адреса](#повторная-проверка-адреса)
|
||||
- [Принудительная остановка проверки](#принудительная-остановка-проверки)
|
||||
@@ -275,6 +276,44 @@ curl -s -X DELETE http://<control-api>:8080/api/v1/admin/config/sites/1
|
||||
> поддерживаемый сценарий; *больше* трёх потребует доработки схемы
|
||||
> данных, одной правкой конфига не обойтись.
|
||||
|
||||
## Управление типами проверок пробера
|
||||
|
||||
Список TCP-портов и флаг ICMP, которые `prober` проверяет на каждой
|
||||
настроенной площадке — единый глобальный набор, общий для всех площадок
|
||||
сразу (не то же самое, что список `sites` выше: `sites` решает, *сколько*
|
||||
точек его применяют, а этот набор — *что именно* они проверяют).
|
||||
|
||||
Посмотреть текущий набор:
|
||||
```bash
|
||||
curl -s http://<control-api>:8080/api/v1/admin/config/inbound-checks | python3 -m json.tool
|
||||
```
|
||||
|
||||
Изменить набор портов и/или ICMP (без перезапуска control-api — новое
|
||||
значение сразу видно и следующему опросу пробера, и уже идущей агрегации):
|
||||
```bash
|
||||
curl -s -X PUT http://<control-api>: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://<control-api>: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`,
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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()
|
||||
|
||||
@@ -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"`
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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()
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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())))
|
||||
}
|
||||
|
||||
@@ -45,4 +45,26 @@
|
||||
</form>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
<div class="panel" style="margin-top:16px">
|
||||
<div class="panel-body">
|
||||
<p class="muted" style="margin-bottom:16px">Типы проверок, которые прогоняет <code>prober</code> на каждой
|
||||
настроенной площадке для каждого проверяемого адреса — глобальный набор, общий для всех площадок. Отдельно от
|
||||
списка площадок (<a href="/sites">/sites</a>), который решает, <em>сколько</em> площадок опрашивать, а не
|
||||
<em>что</em> именно они проверяют. Пустой список портов и выключенный ICMP — штатный способ временно свести
|
||||
inbound-проверки к нулю, не трогая список площадок.</p>
|
||||
<form hx-put="/settings/inbound-checks" hx-target="#settings-form-wrap" hx-swap="innerHTML">
|
||||
<div class="field-row">
|
||||
<div class="field">
|
||||
<label for="inbound_ports">TCP-порты (через запятую)</label>
|
||||
<input type="text" id="inbound_ports" name="ports" placeholder="22, 80, 443, 8080" value="{{joinInts .Inbound.Ports ", "}}">
|
||||
</div>
|
||||
<div class="field">
|
||||
<label for="inbound_icmp"><input type="checkbox" id="inbound_icmp" name="icmp" {{if .Inbound.ICMP}}checked{{end}}> ICMP</label>
|
||||
</div>
|
||||
<button type="submit" class="btn btn-primary">Сохранить</button>
|
||||
</div>
|
||||
</form>
|
||||
</div>
|
||||
</div>
|
||||
{{end}}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
);
|
||||
@@ -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
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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"`
|
||||
}
|
||||
@@ -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})
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user