2026-10-01 19:31:11 +03:00
|
|
|
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))]
|
2026-10-03 18:36:03 +03:00
|
|
|
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)
|
2026-10-01 19:31:11 +03:00
|
|
|
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
|
|
|
|
|
}
|