Scan floating IPs in the background, page by page, so thousands of addresses work

The "Scan Floating IP" button failed with a client timeout: the project now
holds ~6.4k floating IPs and the scan listed them all in one unpaginated,
timeout-less Neutron request on the HTTP request context.

openstack: ListFreeFloatingIPs reads marker-based pages (fields= keeps them
small) with per-page retry/backoff on transport errors, 5xx and 429, and every
request now has a timeout (also ends hangs inside the orchestrator tick).

orchestrator: the scan is a single-flight background job on the process
context with progress (clearing/listing/enqueuing/done/error), dry_run, full
discovery before anything is enqueued, then SubmitIPs in chunks of 500 in
ascending IP order; a failed read leaves the queue untouched. The auto-cycle
gets a "scanning" phase that polls the job, so the control loop and
autoCycleMu are never held across OpenStack/DB work; it recovers after a
restart and waits for (instead of adopting) a scan started by someone else.

db: migration 0009 (indexes), paged ListIPsPage/ListRegistryPage, GROUP BY
counters, EXISTS completion check, set-based ClearAllIPs.

API: POST /admin/ips/scan -> 202 (dry_run, wait), GET /admin/ips/scan, paging
and filters on /admin/ips and /admin/registry (bare arrays without limit),
results_by_overall in /admin/status.

dashboard: scan progress panel and dry-run button, paginated /ips and
/registry with server-side filters, Overview on counters and capped lists
with progress/ETA, "select all N by filter", hx-params fix for per-row
buttons, real counts in confirmations.

Also: docs (API, USAGE, DASHBOARD, README), plan and review under
docs/changes/, bin/ rebuilt with new SHA256SUMS.

Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
This commit is contained in:
ayurishchevandClaude Sonnet 5.5 committed 2026-10-01 19:31:11 +03:00
1 parent debf2afed2
commit aff8fe38b5
61 files changed
+5833 -536

No files matched your search

+131 -63
View File
@@ -11,11 +11,14 @@ import (
// This file implements the automatic check cycle: an optional, repeating
// "clear queue -> scan floating IPs -> wait until every queued address has
// reached a terminal state -> wait interval" scenario. Steps 1-2 reuse
// ClearQueue/ScanFloatingIPs verbatim; step 3 needs no code at all because
// Tick already picks up `queued` addresses. All state lives in the database
// (db.AutoCycle), so the cycle survives a control-api restart and the
// interval/limits can be changed at runtime.
// reached a terminal state -> wait interval" scenario. Steps 1-2 run as ONE
// background scan job (StartScan{ClearFirst:true}) so the control loop and
// autoCycleMu are never held across OpenStack/DB-heavy work: the cycle sits in
// phase `scanning` while the job runs and each step merely polls it. Step 3
// needs no code at all because Tick already picks up `queued` addresses. All
// state lives in the database (db.AutoCycle), so the cycle survives a
// control-api restart (a restart in phase `scanning` simply starts the scan
// again) and the interval/limits can be changed at runtime.
// GetAutoCycle returns the current auto-cycle configuration and state.
func (o *Orchestrator) GetAutoCycle(ctx context.Context) (db.AutoCycle, error) {
@@ -65,12 +68,18 @@ func (o *Orchestrator) StopAutoCycle(ctx context.Context) error {
// "stopped" describes an interrupted run. Stopping during the pause
// between cycles must not overwrite the result of the last finished one.
outcome := ""
if ac.Phase == db.AutoCyclePhaseRunning {
if ac.Phase == db.AutoCyclePhaseRunning || ac.Phase == db.AutoCyclePhaseScanning {
outcome = db.AutoCycleOutcomeStopped
}
if err := o.DB.SetAutoCycleEnabled(ctx, false, nil, outcome); err != nil {
return fmt.Errorf("disable auto cycle: %w", err)
}
if ac.Phase == db.AutoCyclePhaseScanning {
// Persisted first, so a step racing in right after cannot restart the
// scan; then abort the background job (queue chunks already enqueued
// stay, like in-flight checks).
o.CancelScan()
}
o.event(ctx, "control-api", "", nil, "auto_cycle_stopped", autoCyclePayload(map[string]any{
"phase": ac.Phase,
}))
@@ -106,7 +115,11 @@ func (o *Orchestrator) autoCycleStep(ctx context.Context, now time.Time) {
return
}
if ac.Phase == db.AutoCyclePhaseRunning {
switch ac.Phase {
case db.AutoCyclePhaseScanning:
o.autoCycleCheckScan(ctx, ac, now)
return
case db.AutoCyclePhaseRunning:
o.autoCycleCheckRun(ctx, ac, now)
return
}
@@ -118,95 +131,145 @@ func (o *Orchestrator) autoCycleStep(ctx context.Context, now time.Time) {
o.autoCycleStartRun(ctx, ac, now)
}
// autoCycleStartRun performs steps 1-2 of the scenario (clear the queue,
// scan floating IPs) and moves the state to running, or straight to waiting
// if there is nothing to wait for.
// autoCycleStartRun begins a cycle: it starts the background scan job (clear
// the queue, then discover and enqueue the free floating IPs) and moves to
// phase `scanning` right away. Nothing slow happens here, so autoCycleMu is
// released immediately and the control loop keeps ticking.
func (o *Orchestrator) autoCycleStartRun(ctx context.Context, ac db.AutoCycle, now time.Time) {
interval := time.Duration(ac.IntervalSeconds) * time.Second
next := now.Add(interval)
if _, started := o.StartScan(ScanOptions{ClearFirst: true}); !started {
// A scan started by someone else (an operator's manual or dry-run scan,
// the periodic scan) is in flight. It is not this cycle's scan: it may
// not clear the queue, or may not enqueue anything at all (dry run), so
// adopting it would end in a "completed" cycle over an untouched or
// empty queue. Change nothing and try again on the next step, once it
// has finished (it is bounded by fip_scan_timeout_seconds).
o.Log.Info("auto cycle: another floating ip scan is running, waiting for it to finish")
return
}
st := autoCycleStateOf(ac)
st.LastRunStartedAt = &now
st.RunStartedAt = nil
st.NextRunAt = nil
st.Phase = db.AutoCyclePhaseScanning
st.LastError = ""
if err := o.DB.UpdateAutoCycleState(ctx, st); err != nil {
o.Log.Error("auto cycle: save state", "err", err)
}
}
fail := func(step string, err error) {
o.Log.Error("auto cycle: step failed", "step", step, "err", err)
// autoCycleCheckScan handles the scanning phase by polling the scan job.
func (o *Orchestrator) autoCycleCheckScan(ctx context.Context, ac db.AutoCycle, now time.Time) {
scan := o.ScanStatus()
switch {
case scan.Running:
return
case scan.State == ScanIdle:
// Phase `scanning` but no job in this process: control-api restarted
// mid-scan. Starting again is idempotent (clear + scan).
o.Log.Info("auto cycle: restarting floating ip scan after restart")
o.StartScan(ScanOptions{ClearFirst: true})
return
}
interval := time.Duration(ac.IntervalSeconds) * time.Second
next := now.Add(interval)
st := autoCycleStateOf(ac)
switch scan.State {
case ScanError:
o.Log.Error("auto cycle: step failed", "step", "scan floating ips", "err", scan.Error)
st.Phase = db.AutoCyclePhaseWaiting
st.RunStartedAt = nil
st.NextRunAt = &next
st.LastRunFinishedAt = &now
st.LastOutcome = db.AutoCycleOutcomeError
st.LastError = fmt.Sprintf("%s: %v", step, err)
if uerr := o.DB.UpdateAutoCycleState(ctx, st); uerr != nil {
o.Log.Error("auto cycle: save state", "err", uerr)
st.LastError = "scan floating ips: " + scan.Error
if err := o.DB.UpdateAutoCycleState(ctx, st); err != nil {
o.Log.Error("auto cycle: save state", "err", err)
return
}
o.event(ctx, "control-api", "", nil, "auto_cycle_error", autoCyclePayload(map[string]any{
"step": step,
"error": err.Error(),
"step": "scan floating ips",
"error": scan.Error,
}))
}
if _, err := o.ClearQueue(ctx); err != nil {
fail("clear queue", err)
return
}
_, scanned, err := o.ScanFloatingIPs(ctx)
if err != nil {
fail("scan floating ips", err)
return
}
st.LastScannedFree = scanned
st.LastError = ""
if scanned == 0 {
// Nothing was queued; waiting for completion would never end.
o.Log.Info("auto cycle: no free floating ips, waiting for next interval")
case ScanCancelled:
if life := o.lifetimeErr(); life != nil {
// The process is shutting down: leave the phase as is so the next
// start resumes the cycle (scanning without a job restarts it).
return
}
// Cancelled by an operator (CancelScan): treat as an interrupted run.
st.Phase = db.AutoCyclePhaseWaiting
st.RunStartedAt = nil
st.NextRunAt = &next
st.LastRunFinishedAt = &now
st.LastOutcome = db.AutoCycleOutcomeNoFreeIPs
st.LastOutcome = db.AutoCycleOutcomeStopped
if err := o.DB.UpdateAutoCycleState(ctx, st); err != nil {
o.Log.Error("auto cycle: save state", "err", err)
}
return
}
st.Phase = db.AutoCyclePhaseRunning
st.RunStartedAt = &now
st.NextRunAt = nil
if err := o.DB.UpdateAutoCycleState(ctx, st); err != nil {
o.Log.Error("auto cycle: save state", "err", err)
return
default: // ScanDone
st.LastScannedFree = scan.Free
st.LastError = ""
if scan.Free == 0 {
// Nothing was queued; waiting for completion would never end.
o.Log.Info("auto cycle: no free floating ips, waiting for next interval")
st.Phase = db.AutoCyclePhaseWaiting
st.RunStartedAt = nil
st.NextRunAt = &next
st.LastRunFinishedAt = &now
st.LastOutcome = db.AutoCycleOutcomeNoFreeIPs
if err := o.DB.UpdateAutoCycleState(ctx, st); err != nil {
o.Log.Error("auto cycle: save state", "err", err)
}
return
}
// max_run_seconds counts from the end of the scan.
st.Phase = db.AutoCyclePhaseRunning
st.RunStartedAt = &now
st.NextRunAt = nil
if err := o.DB.UpdateAutoCycleState(ctx, st); err != nil {
o.Log.Error("auto cycle: save state", "err", err)
return
}
o.event(ctx, "control-api", "", nil, "auto_cycle_started", autoCyclePayload(map[string]any{
"reason": "cycle",
"scanned_free": scan.Free,
}))
}
o.event(ctx, "control-api", "", nil, "auto_cycle_started", autoCyclePayload(map[string]any{
"reason": "cycle",
"scanned_free": scanned,
}))
}
func (o *Orchestrator) lifetimeErr() error {
o.scan.mu.Lock()
defer o.scan.mu.Unlock()
return o.scan.lifetime().Err()
}
// autoCycleCheckRun handles the running phase: finish the cycle once every
// queued address is terminal, or give up after max_run_seconds.
func (o *Orchestrator) autoCycleCheckRun(ctx context.Context, ac db.AutoCycle, now time.Time) {
items, err := o.DB.ListIPs(ctx)
// An empty queue counts as finished: right after the scan it cannot be
// empty (scanned > 0), so it only happens when an operator cleared or
// deleted every address mid-cycle — and then there is nothing to wait for
// (with max_run_seconds=0 the cycle would otherwise hang forever).
// EXISTS is cheap enough to run on every tick, however long the queue.
pending, err := o.DB.AnyNonTerminalIP(ctx)
if err != nil {
o.Log.Error("auto cycle: list ips", "err", err)
o.Log.Error("auto cycle: check pending ips", "err", err)
return
}
interval := time.Duration(ac.IntervalSeconds) * time.Second
next := now.Add(interval)
// An empty queue counts as finished: right after the scan it cannot be
// empty (scanned > 0), so it only happens when an operator cleared or
// deleted every address mid-cycle — and then there is nothing to wait for
// (with max_run_seconds=0 the cycle would otherwise hang forever).
allTerminal := true
for _, it := range items {
if !isTerminalIPState(it.State) {
allTerminal = false
break
if !pending {
_, addresses, err := o.DB.CountIPsByState(ctx)
if err != nil {
o.Log.Error("auto cycle: count ips", "err", err)
return
}
}
if allTerminal {
st := autoCycleStateOf(ac)
st.Phase = db.AutoCyclePhaseWaiting
st.RunStartedAt = nil
@@ -220,7 +283,7 @@ func (o *Orchestrator) autoCycleCheckRun(ctx context.Context, ac db.AutoCycle, n
return
}
o.event(ctx, "control-api", "", nil, "auto_cycle_completed", autoCyclePayload(map[string]any{
"addresses": len(items),
"addresses": addresses,
"runs": st.RunsTotal,
}))
return
@@ -228,6 +291,11 @@ func (o *Orchestrator) autoCycleCheckRun(ctx context.Context, ac db.AutoCycle, n
if ac.MaxRunSeconds > 0 && ac.RunStartedAt != nil &&
now.Sub(*ac.RunStartedAt) > time.Duration(ac.MaxRunSeconds)*time.Second {
_, addresses, err := o.DB.CountIPsByState(ctx)
if err != nil {
o.Log.Error("auto cycle: count ips", "err", err)
return
}
// The queue is left untouched: the next cycle clears it anyway, and
// an operator can still inspect what got stuck.
st := autoCycleStateOf(ac)
@@ -243,7 +311,7 @@ func (o *Orchestrator) autoCycleCheckRun(ctx context.Context, ac db.AutoCycle, n
}
o.event(ctx, "control-api", "", nil, "auto_cycle_timeout", autoCyclePayload(map[string]any{
"max_run_seconds": ac.MaxRunSeconds,
"addresses": len(items),
"addresses": addresses,
}))
}
}
+294 -15
View File
@@ -58,6 +58,37 @@ func countEvents(t *testing.T, d *db.DB, eventType string) int {
return n
}
// waitScan blocks until the background scan job is no longer running.
func waitScan(t *testing.T, o *Orchestrator) ScanStatus {
t.Helper()
return waitScanFor(t, o, 30*time.Second)
}
func waitScanFor(t *testing.T, o *Orchestrator, limit time.Duration) ScanStatus {
t.Helper()
deadline := time.Now().Add(limit)
for {
if st := o.ScanStatus(); !st.Running {
return st
}
if time.Now().After(deadline) {
t.Fatalf("scan job did not finish in time: %+v", o.ScanStatus())
}
time.Sleep(2 * time.Millisecond)
}
}
// stepStartCycle drives one cycle start: the first step launches the
// background scan (phase scanning), then we wait for the job and the second
// step, at the same virtual time, consumes its result.
func stepStartCycle(t *testing.T, o *Orchestrator, now time.Time) {
t.Helper()
ctx := context.Background()
o.autoCycleStep(ctx, now)
waitScan(t, o)
o.autoCycleStep(ctx, now)
}
func TestAutoCycleDisabledIsNoOp(t *testing.T) {
ctx := context.Background()
o, d, mock := newTestOrchestrator(t, 180)
@@ -111,7 +142,7 @@ func TestAutoCycleStartClearsQueueAndScans(t *testing.T) {
}
now := db.Now()
o.autoCycleStep(ctx, now)
stepStartCycle(t, o, now)
got := queuedAddresses(t, d)
if _, stale := got["9.9.9.9"]; stale {
@@ -154,7 +185,7 @@ func TestAutoCycleWaitsWhileChecksInProgressAndTickPicksUpQueue(t *testing.T) {
}
t0 := db.Now()
o.autoCycleStep(ctx, t0)
stepStartCycle(t, o, t0)
// Existing Tick logic starts the checks on its own.
o.Tick(ctx)
@@ -198,7 +229,7 @@ func TestAutoCycleCompletesAndRepeatsAfterInterval(t *testing.T) {
}
t0 := db.Now()
o.autoCycleStep(ctx, t0)
stepStartCycle(t, o, t0)
// done, failed and occupied are all terminal.
if _, err := d.ExecContext(ctx, `UPDATE ip_queue SET state=? WHERE ip_address=?`, db.IPDone, "1.1.1.1"); err != nil {
@@ -241,7 +272,7 @@ func TestAutoCycleCompletesAndRepeatsAfterInterval(t *testing.T) {
}
// At next_run_at: new cycle starts, queue is rebuilt from scratch.
o.autoCycleStep(ctx, wantNext)
stepStartCycle(t, o, wantNext)
ac = getAutoCycle(t, d)
if ac.Phase != db.AutoCyclePhaseRunning {
t.Fatalf("expected running again, got %+v", ac)
@@ -267,7 +298,7 @@ func TestAutoCycleTimeout(t *testing.T) {
}
t0 := db.Now()
o.autoCycleStep(ctx, t0)
stepStartCycle(t, o, t0)
o.autoCycleStep(ctx, t0.Add(300*time.Second))
if ac := getAutoCycle(t, d); ac.Phase != db.AutoCyclePhaseRunning {
@@ -302,7 +333,7 @@ func TestAutoCycleNoLimitNeverTimesOut(t *testing.T) {
t.Fatalf("start: %v", err)
}
t0 := db.Now()
o.autoCycleStep(ctx, t0)
stepStartCycle(t, o, t0)
o.autoCycleStep(ctx, t0.Add(1000*time.Hour))
if ac := getAutoCycle(t, d); ac.Phase != db.AutoCyclePhaseRunning {
t.Fatalf("max_run_seconds=0 means no limit, got %+v", ac)
@@ -320,7 +351,7 @@ func TestAutoCycleManualClearMidCycleCompletes(t *testing.T) {
t.Fatalf("start: %v", err)
}
t0 := db.Now()
o.autoCycleStep(ctx, t0)
stepStartCycle(t, o, t0)
if ac := getAutoCycle(t, d); ac.Phase != db.AutoCyclePhaseRunning {
t.Fatalf("expected running after start, got %+v", ac)
}
@@ -350,7 +381,7 @@ func TestAutoCycleNoFreeIPs(t *testing.T) {
}
now := db.Now()
o.autoCycleStep(ctx, now)
stepStartCycle(t, o, now)
ac := getAutoCycle(t, d)
if ac.Phase != db.AutoCyclePhaseWaiting || ac.LastOutcome != db.AutoCycleOutcomeNoFreeIPs {
@@ -378,7 +409,7 @@ func TestAutoCycleOpenStackErrorRetriesNextInterval(t *testing.T) {
}
t0 := db.Now()
o.autoCycleStep(ctx, t0)
stepStartCycle(t, o, t0)
ac := getAutoCycle(t, d)
if ac.Phase != db.AutoCyclePhaseWaiting || ac.LastOutcome != db.AutoCycleOutcomeError {
@@ -400,7 +431,7 @@ func TestAutoCycleOpenStackErrorRetriesNextInterval(t *testing.T) {
// OpenStack recovers: the next interval starts a normal cycle and the
// stale error is cleared.
mock.ListFailure = nil
o.autoCycleStep(ctx, t0.Add(60*time.Second))
stepStartCycle(t, o, t0.Add(60*time.Second))
ac = getAutoCycle(t, d)
if ac.Phase != db.AutoCyclePhaseRunning || ac.LastError != "" {
t.Fatalf("expected running with cleared error, got %+v", ac)
@@ -415,7 +446,7 @@ func TestAutoCycleStartIsIdempotent(t *testing.T) {
t.Fatalf("start: %v", err)
}
t0 := db.Now()
o.autoCycleStep(ctx, t0)
stepStartCycle(t, o, t0)
if err := o.StartAutoCycle(ctx); err != nil {
t.Fatalf("second start: %v", err)
@@ -440,7 +471,7 @@ func TestAutoCycleStopMidCycle(t *testing.T) {
t.Fatalf("start: %v", err)
}
t0 := db.Now()
o.autoCycleStep(ctx, t0)
stepStartCycle(t, o, t0)
o.Tick(ctx) // check in flight
if err := o.StopAutoCycle(ctx); err != nil {
@@ -484,7 +515,7 @@ func TestAutoCycleStopWhileWaitingKeepsLastOutcome(t *testing.T) {
t.Fatalf("start: %v", err)
}
t0 := db.Now()
o.autoCycleStep(ctx, t0)
stepStartCycle(t, o, t0)
finishAllIPs(t, d, db.IPDone)
o.autoCycleStep(ctx, t0.Add(5*time.Second))
if ac := getAutoCycle(t, d); ac.Phase != db.AutoCyclePhaseWaiting || ac.LastOutcome != db.AutoCycleOutcomeCompleted {
@@ -512,7 +543,7 @@ func TestAutoCycleSurvivesRestart(t *testing.T) {
t.Fatalf("start: %v", err)
}
t0 := db.Now()
o.autoCycleStep(ctx, t0)
stepStartCycle(t, o, t0)
// A fresh Orchestrator on the same database (a control-api restart)
// continues the running phase instead of starting over.
@@ -537,8 +568,256 @@ func TestAutoCycleSurvivesRestart(t *testing.T) {
if ac := getAutoCycle(t, d); ac.Phase != db.AutoCyclePhaseWaiting {
t.Fatalf("expected waiting until next_run_at, got %+v", ac)
}
o3.autoCycleStep(ctx, t1.Add(300*time.Second))
stepStartCycle(t, o3, t1.Add(300*time.Second))
if ac := getAutoCycle(t, d); ac.Phase != db.AutoCyclePhaseRunning {
t.Fatalf("expected a new cycle at next_run_at, got %+v", ac)
}
}
func TestAutoCycleScanningPhaseThenRunning(t *testing.T) {
ctx := context.Background()
o, d, mock := newTestOrchestrator(t, 180)
mock.SeedMany("fip", 30)
mock.PageSize = 10
mock.PageDelay = 100 * time.Millisecond
setAutoCycleParams(t, d, 60, 300)
if err := d.SeedQueue(ctx, []string{"9.9.9.9"}); err != nil {
t.Fatal(err)
}
if err := o.StartAutoCycle(ctx); err != nil {
t.Fatal(err)
}
t0 := db.Now()
o.autoCycleStep(ctx, t0)
ac := getAutoCycle(t, d)
if ac.Phase != db.AutoCyclePhaseScanning || ac.RunStartedAt != nil || ac.NextRunAt != nil {
t.Fatalf("expected scanning without run_started_at, got %+v", ac)
}
if ac.LastRunStartedAt == nil || !ac.LastRunStartedAt.Equal(t0) {
t.Fatalf("expected last_run_started_at=%v, got %v", t0, ac.LastRunStartedAt)
}
if !o.ScanStatus().Running {
t.Fatalf("expected the scan job to be running")
}
// While the job runs, further steps leave the phase alone, even far past
// max_run_seconds (it only counts from the end of the scan).
o.autoCycleStep(ctx, t0.Add(time.Hour))
if ac := getAutoCycle(t, d); ac.Phase != db.AutoCyclePhaseScanning || ac.LastOutcome != "" {
t.Fatalf("expected still scanning, got %+v", ac)
}
waitScan(t, o)
t1 := t0.Add(2 * time.Hour)
o.autoCycleStep(ctx, t1)
ac = getAutoCycle(t, d)
if ac.Phase != db.AutoCyclePhaseRunning || ac.LastScannedFree != 30 {
t.Fatalf("expected running with last_scanned_free=30, got %+v", ac)
}
if ac.RunStartedAt == nil || !ac.RunStartedAt.Equal(t1) {
t.Fatalf("run_started_at must be the scan end %v, got %v", t1, ac.RunStartedAt)
}
if got := queuedAddresses(t, d); len(got) != 30 {
t.Fatalf("expected the 30 scanned addresses only, got %d", len(got))
}
// max_run_seconds is measured from t1.
o.autoCycleStep(ctx, t1.Add(300*time.Second))
if ac := getAutoCycle(t, d); ac.Phase != db.AutoCyclePhaseRunning {
t.Fatalf("expected still running at the limit, got %+v", ac)
}
o.autoCycleStep(ctx, t1.Add(301*time.Second))
if ac := getAutoCycle(t, d); ac.LastOutcome != db.AutoCycleOutcomeTimeout {
t.Fatalf("expected timeout, got %+v", ac)
}
}
func TestAutoCycleScanErrorSurfacesAndRetries(t *testing.T) {
ctx := context.Background()
o, d, mock := newTestOrchestrator(t, 180)
mock.SeedMany("fip", 5)
mock.ListFailure = errors.New("neutron is down")
setAutoCycleParams(t, d, 60, 0)
if err := o.StartAutoCycle(ctx); err != nil {
t.Fatal(err)
}
t0 := db.Now()
o.autoCycleStep(ctx, t0)
if ac := getAutoCycle(t, d); ac.Phase != db.AutoCyclePhaseScanning {
t.Fatalf("expected scanning first, got %+v", ac)
}
waitScan(t, o)
t1 := t0.Add(time.Second)
o.autoCycleStep(ctx, t1)
ac := getAutoCycle(t, d)
if ac.Phase != db.AutoCyclePhaseWaiting || ac.LastOutcome != db.AutoCycleOutcomeError ||
!strings.Contains(ac.LastError, "neutron is down") || ac.NextRunAt == nil || !ac.NextRunAt.Equal(t1.Add(time.Minute)) {
t.Fatalf("expected waiting/error with next_run_at=t1+interval, got %+v", ac)
}
if countEvents(t, d, "auto_cycle_error") != 1 {
t.Fatalf("expected one auto_cycle_error event")
}
}
func TestAutoCycleRestartInScanningPhaseStartsScanAgain(t *testing.T) {
ctx := context.Background()
o, d, mock := newTestOrchestrator(t, 180)
mock.SeedMany("fip", 12)
if err := d.SeedQueue(ctx, []string{"9.9.9.9"}); err != nil {
t.Fatal(err)
}
setAutoCycleParams(t, d, 60, 0)
if err := o.StartAutoCycle(ctx); err != nil {
t.Fatal(err)
}
t0 := db.Now()
// Persisted state of a process that died mid-scan.
if err := d.UpdateAutoCycleState(ctx, db.AutoCycleState{Phase: db.AutoCyclePhaseScanning, LastRunStartedAt: &t0}); err != nil {
t.Fatal(err)
}
o2 := &Orchestrator{DB: d, OS: mock, Cfg: o.Cfg, Agg: o.Agg, Log: o.Log}
t.Cleanup(func() { o2.CancelScan() })
if st := o2.ScanStatus(); st.State != ScanIdle {
t.Fatalf("fresh process must have no scan job, got %+v", st)
}
o2.autoCycleStep(ctx, t0.Add(time.Second))
if ac := getAutoCycle(t, d); ac.Phase != db.AutoCyclePhaseScanning {
t.Fatalf("expected to stay in scanning while the new job runs, got %+v", ac)
}
if st := o2.ScanStatus(); st.State == ScanIdle {
t.Fatalf("expected a new scan job to be started")
}
waitScan(t, o2)
t1 := t0.Add(2 * time.Second)
o2.autoCycleStep(ctx, t1)
ac := getAutoCycle(t, d)
if ac.Phase != db.AutoCyclePhaseRunning || ac.LastScannedFree != 12 || ac.RunStartedAt == nil || !ac.RunStartedAt.Equal(t1) {
t.Fatalf("expected running after the restarted scan, got %+v", ac)
}
got := queuedAddresses(t, d)
if _, stale := got["9.9.9.9"]; stale || len(got) != 12 {
t.Fatalf("restarted scan must clear and re-scan, got %d rows (stale=%v)", len(got), stale)
}
}
func TestAutoCycleStopCancelsScan(t *testing.T) {
ctx := context.Background()
o, d, mock := newTestOrchestrator(t, 180)
mock.SeedMany("fip", 30)
mock.PageSize = 10
mock.PageDelay = 5 * time.Second
if err := o.StartAutoCycle(ctx); err != nil {
t.Fatal(err)
}
t0 := db.Now()
o.autoCycleStep(ctx, t0)
if ac := getAutoCycle(t, d); ac.Phase != db.AutoCyclePhaseScanning || !o.ScanStatus().Running {
t.Fatalf("precondition: scanning with a live job, got %+v", ac)
}
if err := o.StopAutoCycle(ctx); err != nil {
t.Fatalf("stop: %v", err)
}
ac := getAutoCycle(t, d)
if ac.Enabled || ac.Phase != db.AutoCyclePhaseIdle || ac.LastOutcome != db.AutoCycleOutcomeStopped {
t.Fatalf("expected disabled/idle/stopped, got %+v", ac)
}
if st := o.ScanStatus(); st.State != ScanCancelled || st.Running {
t.Fatalf("expected the scan to be cancelled, got %+v", st)
}
// Later steps do nothing.
o.autoCycleStep(ctx, t0.Add(time.Hour))
if ac := getAutoCycle(t, d); ac.Enabled || ac.Phase != db.AutoCyclePhaseIdle {
t.Fatalf("expected no activity after stop, got %+v", ac)
}
}
// A slow scan must neither block the control loop (Tick, auto-cycle steps)
// nor Start/Stop/Get on the auto-cycle: autoCycleMu is never held across it.
func TestSlowScanDoesNotBlockTickOrAutoCycleCalls(t *testing.T) {
ctx := context.Background()
o, d, mock := newTestOrchestrator(t, 180)
mock.SeedMany("fip", 20)
mock.PageSize = 10
mock.PageDelay = 1500 * time.Millisecond
if err := d.RegisterValidator(ctx, "validator-1", "host-1", "port-1", "v0.1"); err != nil {
t.Fatal(err)
}
within := func(what string, limit time.Duration, f func()) {
t.Helper()
start := time.Now()
f()
if el := time.Since(start); el > limit {
t.Fatalf("%s took %v while a scan was running (limit %v)", what, el, limit)
}
}
within("StartAutoCycle", 500*time.Millisecond, func() {
if err := o.StartAutoCycle(ctx); err != nil {
t.Fatal(err)
}
})
t0 := db.Now()
within("first auto-cycle step", 500*time.Millisecond, func() { o.autoCycleStep(ctx, t0) })
if !o.ScanStatus().Running {
t.Fatalf("scan should be in flight")
}
within("Tick", 500*time.Millisecond, func() { o.Tick(ctx) })
within("polling auto-cycle step", 500*time.Millisecond, func() { o.autoCycleStep(ctx, t0.Add(time.Second)) })
within("GetAutoCycle", 500*time.Millisecond, func() {
if _, err := o.GetAutoCycle(ctx); err != nil {
t.Fatal(err)
}
})
within("StopAutoCycle", 3*time.Second, func() {
if err := o.StopAutoCycle(ctx); err != nil {
t.Fatal(err)
}
})
}
// A scan that is already running (an operator's dry run, a manual or periodic
// scan) is NOT a cycle's scan: it may not clear the queue or enqueue anything.
// The cycle must wait for it and then run its own clear+scan; following the
// foreign job would end in "completed" over an untouched/empty queue.
func TestAutoCycleWaitsForForeignScanInsteadOfFollowingIt(t *testing.T) {
ctx := context.Background()
o, d, mock := newTestOrchestrator(t, 180)
mock.Seed("fip-1", "1.1.1.1", "svc")
mock.PageSize = 1
mock.PageDelay = 80 * time.Millisecond // keeps the dry run in flight for a while
if err := d.SeedQueue(ctx, []string{"9.9.9.9"}); err != nil {
t.Fatalf("seed queue: %v", err)
}
setAutoCycleParams(t, d, 60, 0)
if err := o.StartAutoCycle(ctx); err != nil {
t.Fatalf("start: %v", err)
}
if _, started := o.StartScan(ScanOptions{DryRun: true}); !started {
t.Fatalf("precondition: the dry run must start")
}
t0 := db.Now()
o.autoCycleStep(ctx, t0)
ac := getAutoCycle(t, d)
if ac.Phase == db.AutoCyclePhaseScanning || ac.Phase == db.AutoCyclePhaseRunning {
t.Fatalf("cycle must not adopt the foreign scan, got phase %q", ac.Phase)
}
if got := queuedAddresses(t, d); got["9.9.9.9"] != db.IPQueued {
t.Fatalf("queue must be untouched while the foreign scan runs, got %v", got)
}
waitScan(t, o) // the dry run ends
mock.PageDelay = 0
stepStartCycle(t, o, t0.Add(time.Second))
ac = getAutoCycle(t, d)
if ac.Phase != db.AutoCyclePhaseRunning || ac.LastScannedFree != 1 {
t.Fatalf("expected the cycle's own scan to complete (running, 1 free), got %+v", ac)
}
got := queuedAddresses(t, d)
if _, stale := got["9.9.9.9"]; stale || got["1.1.1.1"] != db.IPQueued {
t.Fatalf("the cycle's own scan must clear the old queue and enqueue 1.1.1.1, got %v", got)
}
}
+50 -69
View File
@@ -41,6 +41,14 @@ type Orchestrator struct {
// autoCycleMu serializes AutoCycleStep with StartAutoCycle/StopAutoCycle
// so an API call can never interleave with a half-finished step.
autoCycleMu sync.Mutex
// ScanPageSize is the Neutron page size used by the floating-IP scan
// (config openstack.list_page_size); 0 means openstack.DefaultListPageSize.
ScanPageSize int
// scan is the background floating-IP scan job (see scanjob.go); its zero
// value is ready to use.
scan scanJob
}
// New constructs an Orchestrator. Egress check types/targets, prober sites,
@@ -58,6 +66,8 @@ func New(d *db.DB, osClient openstack.FloatingIPClient, cfg *config.ControlAPI,
Cfg: cfg.Orchestrator,
Agg: cfg.Aggregation,
Log: log,
ScanPageSize: cfg.OpenStack.ListPageSize,
}
}
@@ -349,17 +359,16 @@ func (o *Orchestrator) DeleteIP(ctx context.Context, ipAddress string) error {
// in the list that has one attached, then deletes the whole list in one
// DB.DeleteIPs call. Addresses not currently in the queue are simply
// omitted from the disassociation pass and reported back in NotFound by
// DB.DeleteIPs — not an error.
// DB.DeleteIPs — not an error. The rows holding a floating IP are found with
// a few IN (...) queries rather than one lookup per address.
func (o *Orchestrator) DeleteIPs(ctx context.Context, addresses []string) (db.DeleteIPsResult, error) {
for _, addr := range addresses {
item, err := o.DB.GetIPByAddress(ctx, addr)
if err != nil {
continue // unknown address — DB.DeleteIPs will report it in NotFound
}
if item.FIPID != "" {
if err := o.OS.DisassociateFloatingIP(ctx, item.FIPID); err != nil {
o.Log.Error("disassociate fip on delete", "ip_id", item.ID, "fip_id", item.FIPID, "err", err)
}
refs, err := o.DB.ListFIPRefsByAddresses(ctx, addresses)
if err != nil {
o.Log.Error("list attached fips before delete", "err", err)
}
for _, ref := range refs {
if err := o.OS.DisassociateFloatingIP(ctx, ref.FIPID); err != nil {
o.Log.Error("disassociate fip on delete", "ip_id", ref.IPID, "fip_id", ref.FIPID, "err", err)
}
}
@@ -372,79 +381,51 @@ func (o *Orchestrator) DeleteIPs(ctx context.Context, addresses []string) (db.De
}
// ClearQueue deletes every address currently in the queue, regardless of
// state — the "delete everything" operation, implemented as DeleteIPs over
// the full current address list rather than a separate DB code path.
// state — the "delete everything" operation. It is set-based (see
// db.ClearAllIPs): O(1) statements however many rows there are. Floating IPs
// attached to rows are disassociated first (best-effort), and only for rows
// that actually hold one.
func (o *Orchestrator) ClearQueue(ctx context.Context) (db.DeleteIPsResult, error) {
items, err := o.DB.ListIPs(ctx)
refs, err := o.DB.ListFIPRefs(ctx)
if err != nil {
return db.DeleteIPsResult{}, fmt.Errorf("list ips: %w", err)
return db.DeleteIPsResult{}, fmt.Errorf("list attached fips: %w", err)
}
addresses := make([]string, len(items))
for i, item := range items {
addresses[i] = item.IPAddress
}
for _, item := range items {
if item.FIPID != "" {
if err := o.OS.DisassociateFloatingIP(ctx, item.FIPID); err != nil {
o.Log.Error("disassociate fip on clear queue", "ip_id", item.ID, "fip_id", item.FIPID, "err", err)
}
for _, ref := range refs {
if err := o.OS.DisassociateFloatingIP(ctx, ref.FIPID); err != nil {
o.Log.Error("disassociate fip on clear queue", "ip_id", ref.IPID, "fip_id", ref.FIPID, "err", err)
}
}
result, err := o.DB.DeleteIPs(ctx, addresses)
deleted, err := o.DB.ClearAllIPs(ctx)
if err != nil {
return result, err
return db.DeleteIPsResult{}, err
}
o.event(ctx, "control-api", "", nil, "queue_cleared", deletedAddressesPayload(result.Deleted))
return result, nil
o.event(ctx, "control-api", "", nil, "queue_cleared", deletedAddressesPayload(deleted))
return db.DeleteIPsResult{Deleted: deleted}, nil
}
// ScanFloatingIPs lists every floating IP in the OpenStack project, filters
// to the ones not currently associated to any port (the free pool awaiting
// validation before reissue), and submits that address list to the check
// queue via db.SubmitIPs — the same entry point the admin API's "add
// addresses" call uses, so add/requeue/reorder semantics are identical
// whether the address list came from an operator or from this scan. Returns
// the SubmitIPs outcome plus how many free floating IPs were found in total
// (which can be larger than the sum of the SubmitIPsResult slices, since
// addresses already mid-check are silently skipped — see db.SubmitIPs).
func (o *Orchestrator) ScanFloatingIPs(ctx context.Context) (db.SubmitIPsResult, int, error) {
fips, err := o.OS.ListFloatingIPs(ctx)
if err != nil {
return db.SubmitIPsResult{}, 0, fmt.Errorf("list floating ips: %w", err)
}
var free []string
for _, f := range fips {
if f.PortID == "" {
free = append(free, f.Address)
}
}
if len(free) == 0 {
o.event(ctx, "control-api", "", nil, "fip_scan", `{"scanned_free":0}`)
return db.SubmitIPsResult{}, 0, nil
}
result, err := o.DB.SubmitIPs(ctx, free)
if err != nil {
return result, len(free), fmt.Errorf("submit scanned ips: %w", err)
}
o.event(ctx, "control-api", "", nil, "fip_scan", fmt.Sprintf(
`{"scanned_free":%d,"added":%d,"requeued":%d,"reordered":%d,"skipped_in_progress":%d}`,
len(free), len(result.Added), len(result.Requeued), len(result.Reordered), len(result.SkippedInProgress)))
return result, len(free), nil
}
// maxEventAddresses caps how many addresses a batch-delete event payload
// lists: clearing 6440 addresses must not write a 100 KB event row.
const maxEventAddresses = 50
// deletedAddressesPayload builds the event payload for the batch delete
// operations — a proper JSON array via encoding/json rather than fmt's %q
// slice formatting (which produces space-separated quoted strings, not
// valid JSON).
// operations — a proper JSON object via encoding/json: the total count plus
// at most the first maxEventAddresses addresses (and "truncated":true when
// the list was cut).
func deletedAddressesPayload(addresses []string) string {
b, err := json.Marshal(struct {
p := struct {
Count int `json:"count"`
Addresses []string `json:"addresses"`
}{addresses})
Truncated bool `json:"truncated,omitempty"`
}{Count: len(addresses), Addresses: addresses}
if p.Addresses == nil {
p.Addresses = []string{}
}
if len(p.Addresses) > maxEventAddresses {
p.Addresses = p.Addresses[:maxEventAddresses]
p.Truncated = true
}
b, err := json.Marshal(p)
if err != nil {
return "{}"
}
+4 -1
View File
@@ -64,7 +64,10 @@ func newTestOrchestratorWithSites(t *testing.T, leaseTTLSeconds int, sites []con
}
log := slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{Level: slog.LevelError}))
return New(d, mock, cfg, log), d, mock
o := New(d, mock, cfg, log)
// A background scan must never outlive the database it writes to.
t.Cleanup(func() { o.CancelScan() })
return o, d, mock
}
func TestHappyPath(t *testing.T) {
+370
View File
@@ -0,0 +1,370 @@
package orchestrator
import (
"context"
"errors"
"fmt"
"net/netip"
"sort"
"sync"
"time"
"cloudipvalidator/internal/db"
"cloudipvalidator/internal/openstack"
)
// This file implements the background floating-IP scan job. With thousands of
// floating IPs a scan takes tens of seconds (Neutron is read page by page)
// and then enqueues thousands of rows, so it can neither run inside an HTTP
// request nor under the auto-cycle mutex. StartScan launches it on the
// process-lifetime context and returns immediately; the job publishes its
// progress as a ScanStatus that anyone can poll.
// ScanState is the phase of the scan job.
type ScanState string
const (
ScanIdle ScanState = "idle" // no scan has run in this process yet
ScanClearing ScanState = "clearing"
ScanListing ScanState = "listing"
ScanEnqueuing ScanState = "enqueuing"
ScanDone ScanState = "done"
ScanError ScanState = "error"
ScanCancelled ScanState = "cancelled"
)
const (
// scanChunkSize is how many addresses go into one db.SubmitIPs call, i.e.
// one short transaction; the single DB connection is released in between
// so the orchestrator tick and the API stay responsive.
scanChunkSize = 500
// defaultScanTimeout is used when Cfg.FIPScanTimeoutSeconds is zero.
defaultScanTimeout = 1800 * time.Second
// cancelWait bounds how long CancelScan waits for the job to wind down.
cancelWait = 5 * time.Second
)
// ScanOptions selects the variant of a scan.
type ScanOptions struct {
// ClearFirst clears the whole queue before scanning (auto-cycle step 1).
// Ignored together with DryRun: a dry run never touches the queue.
ClearFirst bool
// DryRun only discovers and counts the free floating IPs; the queue is
// left untouched.
DryRun bool
}
// ScanStatus is a snapshot of the scan job's progress.
type ScanStatus struct {
State ScanState
Running bool
DryRun bool
Pages int // Neutron pages read so far
Discovered int // floating IPs seen (free and associated)
Free int // of those, free (no port) — the addresses to enqueue
Added int
Requeued int
Reordered int
SkippedInProgress int
StartedAt *time.Time
FinishedAt *time.Time
Error string
}
// scanResult is what a finished job leaves for the synchronous wrapper.
type scanResult struct {
submit db.SubmitIPsResult
free int
err error
}
// scanRun is one execution of the job.
type scanRun struct {
done chan struct{} // closed when the job has finished (result is set)
cancel context.CancelFunc
result scanResult
}
// scanJob is the zero-value-usable scan state embedded in Orchestrator
// (tests build Orchestrator as a literal, so there is no constructor-only
// initialization).
type scanJob struct {
mu sync.Mutex
lifeCtx context.Context // process lifetime; nil = context.Background()
status ScanStatus // zero value reads as idle (see snapshot)
run *scanRun // current/last run
cancelRequested bool
}
// SetContext sets the lifetime context background jobs (the scan) run on.
// Call once at startup; until then context.Background() is used. Cancelling
// it cancels a running scan.
func (o *Orchestrator) SetContext(ctx context.Context) {
o.scan.mu.Lock()
defer o.scan.mu.Unlock()
o.scan.lifeCtx = ctx
}
func (j *scanJob) lifetime() context.Context {
if j.lifeCtx != nil {
return j.lifeCtx
}
return context.Background()
}
// snapshot returns a copy of the status; caller holds j.mu.
func (j *scanJob) snapshot() ScanStatus {
st := j.status
if st.State == "" {
st.State = ScanIdle
}
return st
}
// ScanStatus returns the current scan status (state "idle" if no scan has run
// in this process yet).
func (o *Orchestrator) ScanStatus() ScanStatus {
o.scan.mu.Lock()
defer o.scan.mu.Unlock()
return o.scan.snapshot()
}
// StartScan starts the background scan job, or — single-flight — joins the one
// already running: then started is false and the returned status is the
// running job's. It never blocks on OpenStack or the database.
func (o *Orchestrator) StartScan(opts ScanOptions) (ScanStatus, bool) {
st, _, started := o.startScan(opts)
return st, started
}
func (o *Orchestrator) startScan(opts ScanOptions) (ScanStatus, *scanRun, bool) {
j := &o.scan
j.mu.Lock()
defer j.mu.Unlock()
if j.status.Running && j.run != nil {
return j.snapshot(), j.run, false
}
timeout := time.Duration(o.Cfg.FIPScanTimeoutSeconds) * time.Second
if timeout <= 0 {
timeout = defaultScanTimeout
}
ctx, cancel := context.WithTimeout(j.lifetime(), timeout)
run := &scanRun{done: make(chan struct{}), cancel: cancel}
if opts.DryRun {
opts.ClearFirst = false
}
now := db.Now()
state := ScanListing
if opts.ClearFirst {
state = ScanClearing
}
j.run = run
j.cancelRequested = false
j.status = ScanStatus{State: state, Running: true, DryRun: opts.DryRun, StartedAt: &now}
go o.runScan(ctx, run, opts)
return j.snapshot(), run, true
}
// CancelScan cancels a running scan and waits (briefly) for it to wind down.
// It reports whether a running scan was cancelled. Already-enqueued chunks
// stay in the queue.
func (o *Orchestrator) CancelScan() bool {
j := &o.scan
j.mu.Lock()
run := j.run
if !j.status.Running || run == nil {
j.mu.Unlock()
return false
}
j.cancelRequested = true
j.mu.Unlock()
run.cancel()
select {
case <-run.done:
case <-time.After(cancelWait):
}
return true
}
// update mutates the status under the lock.
func (j *scanJob) update(f func(*ScanStatus)) {
j.mu.Lock()
defer j.mu.Unlock()
f(&j.status)
}
func (o *Orchestrator) runScan(ctx context.Context, run *scanRun, opts ScanOptions) {
j := &o.scan
defer run.cancel()
res, err := o.doScan(ctx, opts)
run.result = res
j.mu.Lock()
finished := db.Now()
switch {
case err == nil:
j.status.State = ScanDone
case j.cancelRequested || errors.Is(j.lifetime().Err(), context.Canceled):
j.status.State = ScanCancelled
err = fmt.Errorf("scan cancelled: %w", context.Canceled)
j.status.Error = "cancelled"
case errors.Is(err, context.DeadlineExceeded):
j.status.State = ScanError
err = fmt.Errorf("scan timed out: %w", err)
j.status.Error = err.Error()
default:
j.status.State = ScanError
j.status.Error = err.Error()
}
j.status.Running = false
j.status.FinishedAt = &finished
state := j.status.State
j.mu.Unlock()
run.result.err = err
if state == ScanError || state == ScanCancelled {
o.Log.Warn("floating ip scan did not complete", "state", state, "err", err)
}
close(run.done)
}
// doScan is the scan algorithm: optional clear -> read every page into memory
// -> sort -> (unless dry run) enqueue in chunks -> one fip_scan event.
func (o *Orchestrator) doScan(ctx context.Context, opts ScanOptions) (scanResult, error) {
j := &o.scan
var res scanResult
if opts.ClearFirst {
if _, err := o.ClearQueue(ctx); err != nil {
return res, fmt.Errorf("clear queue: %w", err)
}
j.update(func(s *ScanStatus) { s.State = ScanListing })
}
// Read everything first: a read error after retries must leave the queue
// untouched, so nothing is enqueued until discovery is complete.
var free []string
seen := map[string]struct{}{}
pages, err := o.OS.ListFreeFloatingIPs(ctx, o.ScanPageSize, func(page []openstack.FloatingIP) error {
nFree := 0
for _, f := range page {
if f.PortID != "" || f.Address == "" {
continue
}
if _, dup := seen[f.Address]; dup {
continue
}
seen[f.Address] = struct{}{}
free = append(free, f.Address)
nFree++
}
j.update(func(s *ScanStatus) {
s.Pages++
s.Discovered += len(page)
s.Free += nFree
})
return ctx.Err()
})
if err != nil {
return res, fmt.Errorf("list floating ips: %w", err)
}
j.update(func(s *ScanStatus) { s.Pages = pages })
res.free = len(free)
sortAddressesAscending(free)
if !opts.DryRun && len(free) > 0 {
j.update(func(s *ScanStatus) { s.State = ScanEnqueuing })
for off := 0; off < len(free); off += scanChunkSize {
if err := ctx.Err(); err != nil {
return res, err
}
chunk := free[off:min(off+scanChunkSize, len(free))]
r, err := o.DB.SubmitIPs(ctx, chunk)
res.submit.Added = append(res.submit.Added, r.Added...)
res.submit.Requeued = append(res.submit.Requeued, r.Requeued...)
res.submit.Reordered = append(res.submit.Reordered, r.Reordered...)
res.submit.SkippedInProgress = append(res.submit.SkippedInProgress, r.SkippedInProgress...)
j.update(func(s *ScanStatus) {
s.Added += len(r.Added)
s.Requeued += len(r.Requeued)
s.Reordered += len(r.Reordered)
s.SkippedInProgress += len(r.SkippedInProgress)
})
if err != nil {
return res, fmt.Errorf("submit scanned ips: %w", err)
}
}
}
if !opts.DryRun {
o.event(ctx, "control-api", "", nil, "fip_scan", fmt.Sprintf(
`{"scanned_free":%d,"pages":%d,"added":%d,"requeued":%d,"reordered":%d,"skipped_in_progress":%d}`,
len(free), pages, len(res.submit.Added), len(res.submit.Requeued),
len(res.submit.Reordered), len(res.submit.SkippedInProgress)))
}
return res, nil
}
// sortAddressesAscending orders addresses numerically (10.0.0.2 before
// 10.0.0.10) so the queue order is deterministic; anything that does not parse
// as an IP goes last, in string order.
func sortAddressesAscending(addrs []string) {
type keyed struct {
s string
ip netip.Addr
ok bool
}
ks := make([]keyed, len(addrs))
for i, a := range addrs {
ip, err := netip.ParseAddr(a)
ks[i] = keyed{s: a, ip: ip, ok: err == nil}
}
sort.SliceStable(ks, func(i, j int) bool {
a, b := ks[i], ks[j]
switch {
case a.ok && b.ok:
return a.ip.Compare(b.ip) < 0
case a.ok != b.ok:
return a.ok
default:
return a.s < b.s
}
})
for i := range ks {
addrs[i] = ks[i].s
}
}
// ScanFloatingIPs lists every floating IP in the OpenStack project, filters
// to the ones not currently associated to any port (the free pool awaiting
// validation before reissue), and submits that address list to the check
// queue via db.SubmitIPs — the same entry point the admin API's "add
// addresses" call uses, so add/requeue/reorder semantics are identical
// whether the address list came from an operator or from this scan. It is the
// synchronous "start (or join) the background scan and wait for it" wrapper:
// returns the aggregated SubmitIPs outcome plus how many free floating IPs
// were found in total (which can be larger than the sum of the
// SubmitIPsResult slices, since addresses already mid-check are silently
// skipped — see db.SubmitIPs). If ctx is cancelled first it returns
// ctx.Err(); the background job keeps running.
func (o *Orchestrator) ScanFloatingIPs(ctx context.Context) (db.SubmitIPsResult, int, error) {
return o.ScanAndWait(ctx, ScanOptions{})
}
// ScanAndWait starts (or joins) the background scan with opts and waits for it
// to finish, returning the aggregated SubmitIPs outcome and the number of free
// floating IPs found. Joining a scan that is already running returns that
// scan's result regardless of opts. If ctx ends first it returns ctx.Err()
// and the job keeps running.
func (o *Orchestrator) ScanAndWait(ctx context.Context, opts ScanOptions) (db.SubmitIPsResult, int, error) {
_, run, _ := o.startScan(opts)
select {
case <-run.done:
case <-ctx.Done():
return db.SubmitIPsResult{}, 0, ctx.Err()
}
return run.result.submit, run.result.free, run.result.err
}
+398
View File
@@ -0,0 +1,398 @@
package orchestrator
import (
"context"
"encoding/json"
"errors"
"net/netip"
"strconv"
"strings"
"testing"
"time"
"cloudipvalidator/internal/db"
)
func eventPayload(t *testing.T, d *db.DB, eventType string) string {
t.Helper()
var p string
if err := d.QueryRowContext(context.Background(),
`SELECT payload FROM events WHERE event_type=? ORDER BY id DESC LIMIT 1`, eventType).Scan(&p); err != nil {
t.Fatalf("read %s event: %v", eventType, err)
}
return p
}
func TestScanStatusIdleBeforeAnyScan(t *testing.T) {
// Zero-value Orchestrator literal (as several tests build it) must work.
o := &Orchestrator{}
st := o.ScanStatus()
if st.State != ScanIdle || st.Running || st.StartedAt != nil {
t.Fatalf("expected idle, got %+v", st)
}
if o.CancelScan() {
t.Fatalf("nothing to cancel")
}
}
func TestScanJobSingleFlightAndProgress(t *testing.T) {
o, d, mock := newTestOrchestrator(t, 180)
mock.SeedMany("fip", 600)
mock.SeedWithPort("busy", "203.0.113.9", "svc", "port-x")
mock.PageSize = 100
mock.PageDelay = 40 * time.Millisecond
st, started := o.StartScan(ScanOptions{})
if !started || !st.Running || st.State != ScanListing || st.StartedAt == nil {
t.Fatalf("expected a started listing job, got started=%v %+v", started, st)
}
st2, started2 := o.StartScan(ScanOptions{DryRun: true})
if started2 || !st2.Running || st2.DryRun {
t.Fatalf("second start must join the running job (not dry-run), got started=%v %+v", started2, st2)
}
// The synchronous wrapper joins the same job instead of scanning twice.
type out struct {
res db.SubmitIPsResult
scanned int
err error
}
ch := make(chan out, 1)
go func() {
r, n, err := o.ScanFloatingIPs(context.Background())
ch <- out{r, n, err}
}()
// Progress is visible while the job runs.
sawProgress := false
for i := 0; i < 200 && o.ScanStatus().Running; i++ {
if s := o.ScanStatus(); s.Pages > 0 && s.Pages < 7 {
sawProgress = true
}
time.Sleep(5 * time.Millisecond)
}
got := <-ch
if got.err != nil || got.scanned != 600 || len(got.res.Added) != 600 {
t.Fatalf("wrapper result: %+v", got)
}
if !sawProgress {
t.Fatalf("expected to observe intermediate progress")
}
fin := waitScan(t, o)
if fin.State != ScanDone || fin.Running || fin.FinishedAt == nil || fin.Error != "" {
t.Fatalf("expected done, got %+v", fin)
}
if fin.Pages != 7 || fin.Discovered != 601 || fin.Free != 600 || fin.Added != 600 {
t.Fatalf("unexpected counters: %+v", fin)
}
if countEvents(t, d, "fip_scan") != 1 {
t.Fatalf("expected exactly one fip_scan event, got %d", countEvents(t, d, "fip_scan"))
}
var p struct {
ScannedFree int `json:"scanned_free"`
Pages int `json:"pages"`
Added int `json:"added"`
}
if err := json.Unmarshal([]byte(eventPayload(t, d, "fip_scan")), &p); err != nil || p.ScannedFree != 600 || p.Added != 600 || p.Pages != 7 {
t.Fatalf("fip_scan payload: %+v err=%v", p, err)
}
// A new scan after completion starts a fresh job; everything is requeued
// or reordered (idempotent), nothing added.
mock.PageDelay = 0
st3, started3 := o.StartScan(ScanOptions{})
if !started3 {
t.Fatalf("a finished job must not block a new one: %+v", st3)
}
fin = waitScan(t, o)
if fin.State != ScanDone || fin.Added != 0 || fin.Reordered != 600 {
t.Fatalf("rescan: %+v", fin)
}
}
func TestScanJobDryRunLeavesQueueUntouched(t *testing.T) {
ctx := context.Background()
o, d, mock := newTestOrchestrator(t, 180)
mock.SeedMany("fip", 450)
mock.PageSize = 200
if err := d.SeedQueue(ctx, []string{"9.9.9.9"}); err != nil {
t.Fatal(err)
}
o.StartScan(ScanOptions{DryRun: true, ClearFirst: true}) // ClearFirst is ignored for dry runs
st := waitScan(t, o)
if st.State != ScanDone || !st.DryRun || st.Free != 450 || st.Pages != 3 || st.Added != 0 {
t.Fatalf("dry run status: %+v", st)
}
if got := queuedAddresses(t, d); len(got) != 1 || got["9.9.9.9"] != db.IPQueued {
t.Fatalf("dry run must not touch the queue, got %d rows", len(got))
}
if countEvents(t, d, "fip_scan") != 0 || countEvents(t, d, "queue_cleared") != 0 {
t.Fatalf("dry run must not emit scan/clear events")
}
}
func TestScanJobReadErrorLeavesQueueUntouched(t *testing.T) {
ctx := context.Background()
o, d, mock := newTestOrchestrator(t, 180)
mock.SeedMany("fip", 450)
mock.PageSize = 200
if err := d.SeedQueue(ctx, []string{"9.9.9.9"}); err != nil {
t.Fatal(err)
}
// Page 1 succeeds, page 2 fails: the 200 already read must NOT be queued.
mock.ListFailures = []error{nil, errors.New("neutron exploded")}
o.StartScan(ScanOptions{})
st := waitScan(t, o)
if st.State != ScanError || !strings.Contains(st.Error, "neutron exploded") || st.FinishedAt == nil {
t.Fatalf("expected error state, got %+v", st)
}
if st.Pages != 1 || st.Added != 0 {
t.Fatalf("expected 1 page read and nothing added, got %+v", st)
}
if got := queuedAddresses(t, d); len(got) != 1 || got["9.9.9.9"] != db.IPQueued {
t.Fatalf("queue must be untouched after a read error, got %d rows", len(got))
}
if countEvents(t, d, "fip_scan") != 0 {
t.Fatalf("no fip_scan event on failure")
}
// The wrapper reports the same error.
mock.ListFailure = errors.New("still down")
if _, _, err := o.ScanFloatingIPs(ctx); err == nil || !strings.Contains(err.Error(), "still down") {
t.Fatalf("expected wrapped list error, got %v", err)
}
}
func TestScanJobClearFirst(t *testing.T) {
ctx := context.Background()
o, d, mock := newTestOrchestrator(t, 180)
mock.Seed("fip-1", "1.1.1.1", "svc")
if err := d.SeedQueue(ctx, []string{"9.9.9.9", "8.8.8.8"}); err != nil {
t.Fatal(err)
}
st, _ := o.StartScan(ScanOptions{ClearFirst: true})
if st.State != ScanClearing {
t.Fatalf("expected clearing as first state, got %s", st.State)
}
fin := waitScan(t, o)
if fin.State != ScanDone || fin.Added != 1 {
t.Fatalf("clear-first scan: %+v", fin)
}
if got := queuedAddresses(t, d); len(got) != 1 || got["1.1.1.1"] != db.IPQueued {
t.Fatalf("expected only the scanned address, got %v", got)
}
var payload struct {
Count int `json:"count"`
Addresses []string `json:"addresses"`
}
if err := json.Unmarshal([]byte(eventPayload(t, d, "queue_cleared")), &payload); err != nil || payload.Count != 2 || len(payload.Addresses) != 2 {
t.Fatalf("queue_cleared payload: %+v err=%v", payload, err)
}
}
// 6440 free addresses (the real stand's size) plus a few out-of-order ones:
// everything is queued in ascending numeric IP order within a sane time.
func TestScanJobEnqueuesAllInAscendingOrderAtScale(t *testing.T) {
if testing.Short() {
t.Skip("scale test skipped in -short mode")
}
ctx := context.Background()
o, d, mock := newTestOrchestrator(t, 180)
o.ScanPageSize = 200
mock.SeedMany("fip", 6440)
// IDs sort before "fip-*", addresses sort numerically before 198.18.*;
// 10.0.0.9 < 10.0.0.200 numerically but not as strings.
mock.Seed("aaa-2", "10.0.0.200", "svc")
mock.Seed("aaa-1", "10.0.0.9", "svc")
mock.SeedWithPort("occupied", "203.0.113.1", "svc", "port-x")
start := time.Now()
o.StartScan(ScanOptions{})
st := waitScanFor(t, o, 5*time.Minute) // -race is an order of magnitude slower
elapsed := time.Since(start)
if st.State != ScanDone || st.Free != 6442 || st.Added != 6442 || st.Discovered != 6443 || st.Pages != 33 {
t.Fatalf("scan status: %+v", st)
}
if elapsed > 3*time.Minute {
t.Fatalf("scanning 6442 addresses took %v", elapsed)
}
t.Logf("scan+enqueue of %d addresses took %v", st.Added, elapsed)
items, err := d.ListIPs(ctx)
if err != nil || len(items) != 6442 {
t.Fatalf("queue: n=%d err=%v", len(items), err)
}
var prev netip.Addr
for i, it := range items {
ip := netip.MustParseAddr(it.IPAddress)
if i > 0 && ip.Compare(prev) <= 0 {
t.Fatalf("queue not ascending at %d: %s after %s", i, it.IPAddress, prev)
}
if it.Sequence != i {
t.Fatalf("expected contiguous sequences, row %d has %d", i, it.Sequence)
}
prev = ip
}
if items[0].IPAddress != "10.0.0.9" || items[1].IPAddress != "10.0.0.200" {
t.Fatalf("numeric ordering broken: %s, %s", items[0].IPAddress, items[1].IPAddress)
}
if _, stale := queuedAddresses(t, d)["203.0.113.1"]; stale {
t.Fatalf("occupied address must not be queued")
}
}
func TestScanJobCancelAndLifetimeContext(t *testing.T) {
ctx := context.Background()
o, d, mock := newTestOrchestrator(t, 180)
mock.SeedMany("fip", 100)
mock.PageSize = 10
mock.PageDelay = 5 * time.Second
if err := d.SeedQueue(ctx, []string{"9.9.9.9"}); err != nil {
t.Fatal(err)
}
o.StartScan(ScanOptions{})
time.Sleep(20 * time.Millisecond)
start := time.Now()
if !o.CancelScan() {
t.Fatalf("expected a running scan to be cancelled")
}
if time.Since(start) > 3*time.Second {
t.Fatalf("cancel took too long")
}
st := o.ScanStatus()
if st.State != ScanCancelled || st.Running || st.FinishedAt == nil {
t.Fatalf("expected cancelled, got %+v", st)
}
if got := queuedAddresses(t, d); len(got) != 1 {
t.Fatalf("cancelled scan must not enqueue, got %d rows", len(got))
}
// Cancelling the process-lifetime context also stops a running scan.
life, cancelLife := context.WithCancel(context.Background())
o.SetContext(life)
o.StartScan(ScanOptions{})
time.Sleep(20 * time.Millisecond)
cancelLife()
if st := waitScan(t, o); st.State != ScanCancelled {
t.Fatalf("expected cancelled by lifetime ctx, got %+v", st)
}
}
func TestScanJobDeadline(t *testing.T) {
o, _, mock := newTestOrchestrator(t, 180)
o.Cfg.FIPScanTimeoutSeconds = 1
mock.SeedMany("fip", 20)
mock.PageSize = 10
mock.PageDelay = 10 * time.Second
o.StartScan(ScanOptions{})
st := waitScan(t, o)
if st.State != ScanError || !strings.Contains(st.Error, "timed out") {
t.Fatalf("expected timeout error, got %+v", st)
}
}
func TestScanFloatingIPsWrapperHonoursCallerContext(t *testing.T) {
o, _, mock := newTestOrchestrator(t, 180)
mock.SeedMany("fip", 20)
mock.PageSize = 10
mock.PageDelay = 5 * time.Second
ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond)
defer cancel()
if _, _, err := o.ScanFloatingIPs(ctx); !errors.Is(err, context.DeadlineExceeded) {
t.Fatalf("expected the caller's ctx error, got %v", err)
}
// The job itself keeps going until cancelled (cleanup cancels it).
if !o.ScanStatus().Running {
t.Fatalf("background job should still be running")
}
}
func TestSortAddressesAscending(t *testing.T) {
in := []string{"10.0.0.10", "zzz", "10.0.0.9", "2001:db8::1", "9.255.255.255", "aaa", "10.0.0.2"}
sortAddressesAscending(in)
want := []string{"9.255.255.255", "10.0.0.2", "10.0.0.9", "10.0.0.10", "2001:db8::1", "aaa", "zzz"}
for i := range want {
if in[i] != want[i] {
t.Fatalf("got %v want %v", in, want)
}
}
}
func TestClearQueueEventPayloadIsTruncated(t *testing.T) {
ctx := context.Background()
o, d, _ := newTestOrchestrator(t, 180)
var addrs []string
for i := 0; i < 120; i++ {
addrs = append(addrs, "10.1.0."+strconv.Itoa(i))
}
if err := d.SeedQueue(ctx, addrs); err != nil {
t.Fatal(err)
}
res, err := o.ClearQueue(ctx)
if err != nil || len(res.Deleted) != 120 {
t.Fatalf("clear: %+v err=%v", res, err)
}
var p struct {
Count int `json:"count"`
Addresses []string `json:"addresses"`
Truncated bool `json:"truncated"`
}
if err := json.Unmarshal([]byte(eventPayload(t, d, "queue_cleared")), &p); err != nil {
t.Fatal(err)
}
if p.Count != 120 || len(p.Addresses) != 50 || !p.Truncated || p.Addresses[0] != "10.1.0.0" {
t.Fatalf("payload: count=%d addrs=%d truncated=%v", p.Count, len(p.Addresses), p.Truncated)
}
}
func TestDeleteIPsDisassociatesOnlyAttachedFIPs(t *testing.T) {
ctx := context.Background()
o, d, mock := newTestOrchestrator(t, 180)
mock.Seed("fip-a", "1.1.1.1", "svc")
mock.Seed("fip-b", "2.2.2.2", "svc")
if _, err := d.SubmitIPs(ctx, []string{"1.1.1.1", "2.2.2.2", "3.3.3.3"}); err != nil {
t.Fatal(err)
}
if err := d.RegisterValidator(ctx, "validator-1", "host-1", "port-1", "v0.1"); err != nil {
t.Fatal(err)
}
claimed, err := d.ClaimNextQueued(ctx, "validator-1", time.Minute)
if err != nil || claimed == nil || claimed.IPAddress != "1.1.1.1" {
t.Fatalf("claim: %+v err=%v", claimed, err)
}
if err := mock.AssociateFloatingIP(ctx, "fip-a", "port-1"); err != nil {
t.Fatal(err)
}
if err := d.SetFIPAssociated(ctx, claimed.ID, "fip-a", time.Minute); err != nil {
t.Fatal(err)
}
// An unrelated, externally associated FIP must stay as it is.
if err := mock.AssociateFloatingIP(ctx, "fip-b", "other-port"); err != nil {
t.Fatal(err)
}
res, err := o.DeleteIPs(ctx, []string{"1.1.1.1", "2.2.2.2", "nope"})
if err != nil || len(res.Deleted) != 2 || len(res.NotFound) != 1 {
t.Fatalf("delete: %+v err=%v", res, err)
}
a, _ := mock.GetFloatingIPByAddress(ctx, "1.1.1.1")
b, _ := mock.GetFloatingIPByAddress(ctx, "2.2.2.2")
if a.PortID != "" {
t.Fatalf("attached fip must be disassociated, still on %q", a.PortID)
}
if b.PortID != "other-port" {
t.Fatalf("row without recorded fip_id must not trigger a disassociate, got %q", b.PortID)
}
v, _ := d.GetValidator(ctx, "validator-1")
if v.State != db.ValidatorIdle {
t.Fatalf("validator must be freed, got %s", v.State)
}
}