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 }