diff --git a/bin/SHA256SUMS b/bin/SHA256SUMS index 48219e8..2dd96bc 100644 --- a/bin/SHA256SUMS +++ b/bin/SHA256SUMS @@ -1,3 +1,3 @@ -a7c481f3c8421dc8b3985d7c3eda13a0ad1e5d37c7a58a4aaf37c1b18254ff54 control-api -2e34dce04889a25c564bbd0dd9dfc26499c9e63486f8419ff622ca34f7a9bed8 validator-agent -e6ec2444854f624b11f917b7d3cf9fa32f9ee49aff0da637fe2e59140e02133c prober +ea275ac1b75e2b1aa16f8c2424be04983fda12b23bc99866e2b56e2dc25dda31 control-api +aa82eca68cf5abb6e91091418c728fa4407d306126d0915799cacf7c6c2a1dac validator-agent +645c69e72d5fcd2ff385eaa558c0e5849f5ee28ffbbc46090dd3be6244b7e36c prober diff --git a/bin/control-api b/bin/control-api index db73f1d..d37fcf3 100755 Binary files a/bin/control-api and b/bin/control-api differ diff --git a/bin/prober b/bin/prober index cf4725e..f906b4a 100755 Binary files a/bin/prober and b/bin/prober differ diff --git a/bin/validator-agent b/bin/validator-agent index 5d1a746..864d8cb 100755 Binary files a/bin/validator-agent and b/bin/validator-agent differ diff --git a/configs/control-api.example.yaml b/configs/control-api.example.yaml index f4317ac..c6e44ad 100644 --- a/configs/control-api.example.yaml +++ b/configs/control-api.example.yaml @@ -57,7 +57,12 @@ validators: - validator_id: "validator_02" os_port_id: "REPLACE_WITH_NEUTRON_PORT_ID_2" -# The three external prober sites, indexed 1-3 (matches ip_queue.siteN_complete). +# 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. sites: - site_id: "site-1" index: 1 diff --git a/docs/API.md b/docs/API.md index 666676b..d291369 100644 --- a/docs/API.md +++ b/docs/API.md @@ -309,11 +309,11 @@ queued ──(control-api сам, без вызова API)──▶ assigning_fi POST .../self-check {success:true} ▼ checking - │ POST .../results (agent) │ POST .../results (prober, x3 площадки) + │ POST .../results (agent) │ POST .../results (prober, по числу настроенных площадок) ▼ ▼ - egress_complete=true siteN_complete=true (N=1,2,3) + egress_complete=true siteN_complete=true (только для N, перечисленных в sites конфига) │ - все complete=true ИЛИ истекло checking_window_seconds + egress + все НАСТРОЕННЫЕ площадки complete=true ИЛИ истекло checking_window_seconds ▼ aggregating │ @@ -326,6 +326,11 @@ queued ──(control-api сам, без вызова API)──▶ assigning_fi `orchestrator.poll_interval_seconds`), явного HTTP-метода для их запуска нет — это фоновый цикл (`Tick`), а не запрос/ответ. +Площадки (`siteN_complete`) — опциональны: сколько их учитывается, +целиком определяется списком `sites` в конфиге control-api (0–3 записи). +Пустой список — агрегация ждёт только `egress_complete`, ни одна площадка +не требуется. Подробнее — [USAGE.md](USAGE.md#управление-площадками-проберами). + ## Сквозной пример работы (curl) Ниже — минимальный ручной прогон одного IP через API, как если бы вы diff --git a/docs/DIAGRAMS.md b/docs/DIAGRAMS.md index c76353b..157b5c3 100644 --- a/docs/DIAGRAMS.md +++ b/docs/DIAGRAMS.md @@ -207,10 +207,17 @@ flowchart TB источника (`source`: `egress` для валидатора, `inbound-site-1/2/3` для площадок). Отдельно, по мере поступления данных, выставляются флаги завершения (`egress_complete`, `site1/2/3_complete`) в `ip_queue`. Как -только все четыре флага выставлены — либо истекло время ожидания +только все ожидаемые флаги выставлены — либо истекло время ожидания (`checking_window_seconds`) — фоновая агрегация суммирует все строки `checks` по текущей попытке и записывает итог (`pass`/`partial`/`fail`) обратно в `ip_queue`. Оператор в любой момент читает уже накопленные данные через административные `GET`-методы, не дожидаясь завершения проверки — подробнее о значениях полей см. [USAGE.md](USAGE.md#значения-полей-ip). + +> На диаграмме показан полный вариант с тремя площадками — это не +> обязательный минимум. Площадки опциональны: сколько их ожидать (0–3), +> определяется списком `sites` в конфиге control-api. При пустом списке +> четвёртый поток телеметрии (`inbound-site-N`) просто отсутствует, и +> агрегация ждёт только `egress_complete`. См. +> [USAGE.md](USAGE.md#управление-площадками-проберами). diff --git a/docs/USAGE.md b/docs/USAGE.md index 4caa4e1..ff94cad 100644 --- a/docs/USAGE.md +++ b/docs/USAGE.md @@ -121,8 +121,12 @@ curl -s http://:8080/api/v1/admin/ips \ ## Как читать итоговый результат (pass/partial/fail) -- **`pass`** — прошли все проверки (все исходящие + все три площадки по - всем портам и ICMP). Адрес можно считать пригодным к повторной выдаче. +- **`pass`** — прошли все проверки: все исходящие, плюс все площадки по + всем портам и ICMP — если площадки вообще настроены (`sites` в конфиге + control-api может быть пустым, см. + [«Управление площадками»](#управление-площадками-проберами) — тогда + учитываются только исходящие). Адрес можно считать пригодным к + повторной выдаче. - **`partial`** — часть проверок прошла, часть — нет (например, площадка site-2 не смогла достучаться по 8080/tcp, но остальное в порядке). Означает частичную деградацию — например, адрес заблокирован в @@ -194,16 +198,28 @@ curl -s http://:8080/api/v1/admin/validators | python3 -m json.tool ## Управление площадками (проберами) -Аналогично валидаторам: чтобы добавить площадку, добавьте `site_id` + -`index` (свободный от 1 до 3, или больше — но текущая схема БД -рассчитана ровно на 3 площадки, см. ниже) в `sites` конфига control-api, -разверните на площадке `prober` с тем же `site_id`. +**Входящие (inbound/prober) проверки полностью опциональны.** Список +`sites` в конфиге control-api и есть переключатель: пусто — inbound- +проверки выключены целиком, итоговый результат считается только по +исходящим (egress) проверкам, и агрегация не ждёт вообще ни одного +пробера. Указан один или два слота — ждём только их, остальные не +учитываются. Указаны все три — работает как в исходной схеме процесса. +Явного отдельного флага "включить/выключить" нет — самого списка `sites` +достаточно. -> Важно: количество площадок в текущей версии жёстко зашито в схему БД -> (`Site1Complete`/`Site2Complete`/`Site3Complete`) — система рассчитана -> ровно на **три** внешние площадки, как и описано в исходной схеме -> процесса. Изменение их числа потребует доработки схемы данных, это не -> делается только правкой конфига. +Чтобы добавить площадку: добавьте `site_id` + `index` (1, 2 или 3 — см. +ограничение ниже) в `sites` конфига control-api и разверните на площадке +`prober` с тем же `site_id`. Чтобы отключить конкретную площадку — +уберите соответствующую запись из `sites` и перезапустите control-api; +процесс `prober` на ней можно не останавливать (он просто перестанет +получать назначения). + +> Важно: количество *возможных* слотов площадок жёстко зашито в схему БД +> (`Site1Complete`/`Site2Complete`/`Site3Complete`) — не более **трёх**, +> как и описано в исходной схеме процесса. `index` может быть только 1, 2 +> или 3. Использовать *меньше* трёх (в том числе ноль) — штатный, +> поддерживаемый сценарий; *больше* трёх потребует доработки схемы +> данных, одной правкой конфига не обойтись. ## Повторная проверка адреса diff --git a/internal/db/queries_ipqueue.go b/internal/db/queries_ipqueue.go index 75268c9..9567722 100644 --- a/internal/db/queries_ipqueue.go +++ b/internal/db/queries_ipqueue.go @@ -267,21 +267,6 @@ func (d *DB) ListChecking(ctx context.Context) ([]IPQueueItem, error) { return scanIPQueueItems(rows) } -// ListReadyToAggregate returns checking-state IPs where every source has -// reported completion, or whose checking window has expired. -func (d *DB) ListReadyToAggregate(ctx context.Context, windowDeadline time.Time) ([]IPQueueItem, error) { - rows, err := d.QueryContext(ctx, ipQueueSelect+` - WHERE state=? AND ( - (egress_complete=1 AND site1_complete=1 AND site2_complete=1 AND site3_complete=1) - OR assigned_at < ? - )`, IPChecking, timeToDB(windowDeadline)) - if err != nil { - return nil, err - } - defer rows.Close() - return scanIPQueueItems(rows) -} - // ListExpiredLeases returns non-terminal IPs whose lease has expired — // candidates for the lease sweep (crash recovery + stuck-validator reclaim). func (d *DB) ListExpiredLeases(ctx context.Context, now time.Time) ([]IPQueueItem, error) { diff --git a/internal/orchestrator/orchestrator.go b/internal/orchestrator/orchestrator.go index 8cf4087..0b48ca3 100644 --- a/internal/orchestrator/orchestrator.go +++ b/internal/orchestrator/orchestrator.go @@ -207,11 +207,14 @@ func (o *Orchestrator) MarkSiteComplete(ctx context.Context, ipID int64, siteInd // every source, or hit the checking-window deadline, into aggregation. func (o *Orchestrator) sweepCheckingWindow(ctx context.Context) error { deadline := db.Now().Add(-time.Duration(o.Cfg.CheckingWindowSeconds) * time.Second) - ready, err := o.DB.ListReadyToAggregate(ctx, deadline) + checking, err := o.DB.ListChecking(ctx) if err != nil { - return fmt.Errorf("list ready to aggregate: %w", err) + return fmt.Errorf("list checking: %w", err) } - for _, item := range ready { + for _, item := range checking { + if !o.isReadyToAggregate(item, deadline) { + continue + } if err := o.aggregateAndRelease(ctx, item); err != nil { o.Log.Error("aggregate and release", "ip_id", item.ID, "err", err) } @@ -219,6 +222,41 @@ func (o *Orchestrator) sweepCheckingWindow(ctx context.Context) error { return nil } +// isReadyToAggregate reports whether an in-progress IP has either finished +// reporting from every source it's actually expecting, or hit the +// checking-window deadline. Which inbound sources it's expecting is driven +// entirely by o.Sites — inbound checks are optional: an empty (or +// partial) `sites` config in control-api.yaml 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) bool { + if item.AssignedAt != nil && item.AssignedAt.Before(deadline) { + return true + } + if !item.EgressComplete { + return false + } + for _, s := range o.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 + } + } + } + return true +} + func (o *Orchestrator) aggregateAndRelease(ctx context.Context, item db.IPQueueItem) error { if err := o.DB.SetAggregating(ctx, item.ID); err != nil { return err diff --git a/internal/orchestrator/orchestrator_test.go b/internal/orchestrator/orchestrator_test.go index fa50f84..2dd6e5e 100644 --- a/internal/orchestrator/orchestrator_test.go +++ b/internal/orchestrator/orchestrator_test.go @@ -13,7 +13,18 @@ import ( "cloudipvalidator/internal/openstack" ) +var threeSites = []config.SiteConfig{ + {SiteID: "site-1", Index: 1}, + {SiteID: "site-2", Index: 2}, + {SiteID: "site-3", Index: 3}, +} + func newTestOrchestrator(t *testing.T, leaseTTLSeconds int) (*Orchestrator, *db.DB, *openstack.MockClient) { + t.Helper() + return newTestOrchestratorWithSites(t, leaseTTLSeconds, threeSites) +} + +func newTestOrchestratorWithSites(t *testing.T, leaseTTLSeconds int, sites []config.SiteConfig) (*Orchestrator, *db.DB, *openstack.MockClient) { t.Helper() ctx := context.Background() dbPath := filepath.Join(t.TempDir(), "test.db") @@ -36,11 +47,7 @@ func newTestOrchestrator(t *testing.T, leaseTTLSeconds int) (*Orchestrator, *db. HeartbeatTimeoutSeconds: 30, }, Aggregation: config.AggregationConfig{MissingCountsAsFail: true}, - Sites: []config.SiteConfig{ - {SiteID: "site-1", Index: 1}, - {SiteID: "site-2", Index: 2}, - {SiteID: "site-3", Index: 3}, - }, + Sites: sites, CheckTypes: []config.CheckTypeConfig{ {Name: "https", Enabled: true, Targets: []string{"web"}}, {Name: "ssh", Enabled: false, Targets: []string{"web"}}, @@ -259,15 +266,24 @@ func TestLeaseReclaim(t *testing.T) { func TestMaxRetriesExhausted(t *testing.T) { ctx := context.Background() - o, d, mock := newTestOrchestrator(t, 0) + // A 1s lease (rather than 0) avoids the same intra-tick race noted in + // TestLeaseReclaim: with a 0s TTL, whether a freshly claimed item is + // reclaimed within the very same Tick (making each iteration's timing + // unpredictable) depends on how much wall-clock time claim+associate + // happened to take, which made this test flaky under load. + o, d, mock := newTestOrchestrator(t, 1) o.Cfg.MaxRetries = 1 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"}) - for i := 0; i < 3; i++ { + // Each (claim, sleep past 1s lease) pair reclaims once, bumping + // retry_count by 1; with MaxRetries=1 that needs two full reclaim + // cycles (retry_count 0->1 requeues, 1->2 fails) — four ticks total: + // claim, reclaim+requeue, re-claim, reclaim+fail. + for i := 0; i < 4; i++ { o.Tick(ctx) - time.Sleep(5 * time.Millisecond) + time.Sleep(1100 * time.Millisecond) } ip, _ := d.GetIPByAddress(ctx, "1.2.3.4") @@ -275,3 +291,89 @@ func TestMaxRetriesExhausted(t *testing.T) { t.Fatalf("expected failed after exhausting retries, got %s (retry_count=%d)", ip.State, ip.RetryCount) } } + +// TestInboundChecksDisabled confirms inbound (prober) checks are genuinely +// optional: with no sites configured, an IP must aggregate as soon as +// egress completes, without ever waiting on siteN_complete flags that +// nothing will ever set — and, critically, without waiting out the full +// checking_window_seconds timeout to get there (newTestOrchestrator uses +// 120s; this test never sleeps, so a pass here proves the "all required +// sources complete" path fired, not the timeout fallback). +func TestInboundChecksDisabled(t *testing.T) { + ctx := context.Background() + o, d, mock := newTestOrchestratorWithSites(t, 180, nil) + 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) + + o.Tick(ctx) + + ip, _ = d.GetIP(ctx, ip.ID) + if ip.State != db.IPDone { + t.Fatalf("expected done immediately after egress completed with no sites configured, got %s", ip.State) + } + if ip.OverallResult != db.ResultPass { + t.Fatalf("expected pass, got %s", ip.OverallResult) + } +} + +// TestInboundChecksPartialSites confirms a partially-configured sites list +// (fewer than 3 slots assigned) only waits on the sites actually +// configured — the two unassigned slots are never expected. +func TestInboundChecksPartialSites(t *testing.T) { + ctx := context.Background() + sites := []config.SiteConfig{{SiteID: "site-1", Index: 1}} + 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) + + // Egress is done but the one configured site (site-1) hasn't reported + // yet — must not aggregate. + o.Tick(ctx) + ip, _ = d.GetIP(ctx, ip.ID) + if ip.State != db.IPChecking { + t.Fatalf("expected still checking (site-1 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(1), CheckType: ct, 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 single configured site reported, got %s", ip.State) + } + if ip.OverallResult != db.ResultPass { + t.Fatalf("expected pass, got %s", ip.OverallResult) + } +}