Files
cloud-ip-validator/internal/orchestrator/scanjob.go
T
ayurishchevandClaude Sonnet 5.5 b7669c9e41 Add the Analytics section: check runs, analytics API and page
Runs (migration 0011): a run groups the cycles of one launch. It opens when an
address enters an idle queue, takes everything submitted or re-checked while it
is open and is finalized when all its addresses are done; a re-check after that
opens a new run, so results of different runs never mix. check_runs,
run_results (one result per address and run, with the verdict and the expected
and stored check counts), subnets, run_id on ip_queue and checks. Existing data
is split into runs at pauses of more than an hour; ingress checks get the
validator that held the address (also at write time from now on).

Analytics (internal/analytics): figures computed from the stored checks of the
latest cycle of each address in the run, as facts next to the verdict: summary,
reasons of partial, data quality, subnets, targets and the subnet x target
matrix by check type, ingress by site, error classes, validators, and the
address lists behind the indicators and error classes. API: analytics runs,
report, lists (JSON or CSV), subnet list; run and subnet filters for the
registry.

Dashboard: /analytics matching the approved mockup (run selector, indicators
with address lists and CSV, error-class dialogs, drill-down to the registry),
subnet list on /settings. Sidebar: the control-api link state, theme toggle and
logout moved to the top, the three dots next to the logo removed, sections
grouped.

Rebuilt bin/control-api and bin/admin-dashboard to match. Plan, summary and the
updated README, API, USAGE, DASHBOARD and ADMIN_CLEANUP docs are in docs/.

Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
2026-10-03 18:36:03 +03:00

375 lines
11 KiB
Go

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))]
kind := db.RunManual
if opts.ClearFirst {
kind = db.RunAuto // the auto-cycle's scan; a manual scan never clears first
}
r, err := o.DB.SubmitIPsAs(ctx, chunk, kind)
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
}