diff --git a/bin/SHA256SUMS b/bin/SHA256SUMS index 3c58bfb..5834d8e 100644 --- a/bin/SHA256SUMS +++ b/bin/SHA256SUMS @@ -1,4 +1,4 @@ -26ff297b66a0c974b7143179e44a58f476265482cd1929bc01a95ad29e50ca7a control-api +4562d3cf3db8075b5a1213b582359991e3f58c819151cf7b25276cba0a4e5985 control-api 091ad94b5706b1778b181542de7a4b251cd16c421144269d0955769603d51e06 validator-agent 3e9e14dbb361ee76aaad7c1da6864b3ea111e0ed151403f904b12485631bbf75 prober d72234688eb1954dbff420ca1c8e83b83ceab561b18a0336af0f73fe4d9ac8be admin-dashboard diff --git a/bin/control-api b/bin/control-api index bac8e17..1c08640 100755 Binary files a/bin/control-api and b/bin/control-api differ diff --git a/cmd/control-api/main.go b/cmd/control-api/main.go index aa3f4c3..4aff1d3 100644 --- a/cmd/control-api/main.go +++ b/cmd/control-api/main.go @@ -59,6 +59,7 @@ func run(configPath string, log *slog.Logger) error { } orch := orchestrator.New(database, osClient, cfg, log) + orch.Async = true // slow OpenStack calls run per address, Tick never waits for them // Background jobs (the floating-IP scan) live as long as the process, not // as long as the HTTP request or loop iteration that started them. orch.SetContext(ctx) diff --git a/internal/orchestrator/orchestrator.go b/internal/orchestrator/orchestrator.go index 2a1f4d3..11f4c30 100644 --- a/internal/orchestrator/orchestrator.go +++ b/internal/orchestrator/orchestrator.go @@ -49,6 +49,36 @@ type Orchestrator struct { // scan is the background floating-IP scan job (see scanjob.go); its zero // value is ready to use. scan scanJob + + // Async makes Tick return without waiting for the slow per-address work + // (OpenStack calls): each address is handled in its own goroutine, so one + // slow cloud call never delays another validator. With Async=false (the + // default, used by tests) Tick runs the same work in parallel but waits + // for all of it before returning. + Async bool + + // inflight holds the keys of work already running in a goroutine, so the + // next Tick does not start the same work twice. + inflight sync.Map + wg sync.WaitGroup +} + +// Wait blocks until all background work started by Tick has finished. +func (o *Orchestrator) Wait() { o.wg.Wait() } + +// spawn runs fn in its own goroutine unless work with the same key is +// already running. It returns immediately; callers that need the result +// wait with o.wg (see runAll). +func (o *Orchestrator) spawn(key string, fn func()) { + if _, busy := o.inflight.LoadOrStore(key, struct{}{}); busy { + return + } + o.wg.Add(1) + go func() { + defer o.wg.Done() + defer o.inflight.Delete(key) + fn() + }() } // New constructs an Orchestrator. Egress check types/targets, prober sites, @@ -79,6 +109,10 @@ func (o *Orchestrator) leaseTTL() time.Duration { // validators, sweep the checking window for ready-to-aggregate IPs, and // reclaim expired leases. Intended to be called on a fixed interval // (Cfg.PollIntervalSeconds) by the caller (cmd/control-api/main.go). +// +// The slow part of every step (calls to OpenStack) runs per address in its +// own goroutine, so validators never wait for each other. In Async mode +// Tick does not wait for those goroutines; otherwise it waits for them. func (o *Orchestrator) Tick(ctx context.Context) { if err := o.assignIdleValidators(ctx); err != nil { o.Log.Error("assign idle validators", "err", err) @@ -89,6 +123,9 @@ func (o *Orchestrator) Tick(ctx context.Context) { if err := o.sweepExpiredLeases(ctx); err != nil { o.Log.Error("sweep expired leases", "err", err) } + if !o.Async { + o.wg.Wait() + } } // assignIdleValidators claims the next queued IP for every currently idle @@ -108,9 +145,12 @@ func (o *Orchestrator) assignIdleValidators(ctx context.Context) error { continue // no work available for this validator right now } o.Log.Info("claimed ip", "validator", v.ValidatorID, "ip", item.IPAddress, "ip_id", item.ID) - if err := o.associateFIP(ctx, v.ValidatorID, v.OSPortID, item); err != nil { - o.Log.Error("associate fip", "validator", v.ValidatorID, "ip", item.IPAddress, "err", err) - } + v, item := v, item + o.spawn("assign:"+v.ValidatorID, func() { + if err := o.associateFIP(ctx, v.ValidatorID, v.OSPortID, item); err != nil { + o.Log.Error("associate fip", "validator", v.ValidatorID, "ip", item.IPAddress, "err", err) + } + }) } return nil } @@ -453,9 +493,12 @@ func (o *Orchestrator) sweepCheckingWindow(ctx context.Context) error { if !ready { continue } - if err := o.aggregateAndRelease(ctx, item); err != nil { - o.Log.Error("aggregate and release", "ip_id", item.ID, "err", err) - } + item := item + o.spawn(fmt.Sprintf("aggregate:%d", item.ID), func() { + if err := o.aggregateAndRelease(ctx, item); err != nil { + o.Log.Error("aggregate and release", "ip_id", item.ID, "err", err) + } + }) } return nil } @@ -606,14 +649,17 @@ func (o *Orchestrator) sweepExpiredLeases(ctx context.Context) error { if item.OwnerValidatorID != nil { validatorID = *item.OwnerValidatorID } - o.Log.Info("lease expired, reclaiming", "ip_id", item.ID, "ip", item.IPAddress, "validator", validatorID) - if item.FIPID != "" { - if err := o.OS.DisassociateFloatingIP(ctx, item.FIPID); err != nil { - o.Log.Error("disassociate fip on lease reclaim", "ip_id", item.ID, "err", err) + item, validatorID := item, validatorID + o.spawn(fmt.Sprintf("lease:%d", item.ID), func() { + o.Log.Info("lease expired, reclaiming", "ip_id", item.ID, "ip", item.IPAddress, "validator", validatorID) + if item.FIPID != "" { + if err := o.OS.DisassociateFloatingIP(ctx, item.FIPID); err != nil { + o.Log.Error("disassociate fip on lease reclaim", "ip_id", item.ID, "err", err) + } } - } - o.event(ctx, "control-api", "", &item.ID, "lease_expired", fmt.Sprintf(`{"validator_id":%q}`, validatorID)) - o.requeueOrFail(ctx, item.ID, validatorID, "lease expired") + o.event(ctx, "control-api", "", &item.ID, "lease_expired", fmt.Sprintf(`{"validator_id":%q}`, validatorID)) + o.requeueOrFail(ctx, item.ID, validatorID, "lease expired") + }) } return nil } diff --git a/internal/orchestrator/orchestrator_test.go b/internal/orchestrator/orchestrator_test.go index c862ff2..61fe1d4 100644 --- a/internal/orchestrator/orchestrator_test.go +++ b/internal/orchestrator/orchestrator_test.go @@ -3,6 +3,7 @@ package orchestrator import ( "context" "errors" + "fmt" "log/slog" "os" "path/filepath" @@ -942,3 +943,56 @@ func TestScanFloatingIPsNoFreeAddressesIsNotAnError(t *testing.T) { t.Fatalf("expected nothing added, got %+v", result) } } + +// slowOS adds a fixed delay to every OpenStack call, like a loaded cloud. +type slowOS struct { + *openstack.MockClient + delay time.Duration +} + +func (s slowOS) GetFloatingIPByAddress(ctx context.Context, a string) (*openstack.FloatingIP, error) { + time.Sleep(s.delay) + return s.MockClient.GetFloatingIPByAddress(ctx, a) +} + +func (s slowOS) AssociateFloatingIP(ctx context.Context, fipID, portID string) error { + time.Sleep(s.delay) + return s.MockClient.AssociateFloatingIP(ctx, fipID, portID) +} + +// All idle validators must start in the same tick: one slow cloud call may +// not delay the others (20 validators x 2 calls x 200ms = 8s if sequential). +func TestValidatorsStartInParallel(t *testing.T) { + ctx := context.Background() + o, d, mock := newTestOrchestrator(t, 180) + o.OS = slowOS{MockClient: mock, delay: 200 * time.Millisecond} + + const n = 20 + var addrs []string + for i := 1; i <= n; i++ { + addr := fmt.Sprintf("10.0.0.%d", i) + addrs = append(addrs, addr) + mock.Seed(fmt.Sprintf("fip-%d", i), addr, "svc-project") + if err := d.RegisterValidator(ctx, fmt.Sprintf("validator-%d", i), "h", fmt.Sprintf("port-%d", i), "v"); err != nil { + t.Fatalf("register validator: %v", err) + } + } + if err := d.SeedQueue(ctx, addrs); err != nil { + t.Fatalf("seed queue: %v", err) + } + + start := time.Now() + o.Tick(ctx) + if elapsed := time.Since(start); elapsed > 2*time.Second { + t.Fatalf("tick took %s: validators are served one after another", elapsed) + } + for _, a := range addrs { + ip, err := d.GetIPByAddress(ctx, a) + if err != nil { + t.Fatalf("get ip %s: %v", a, err) + } + if ip.State != db.IPAwaitingSelfCheck { + t.Fatalf("ip %s: expected awaiting_self_check, got %s", a, ip.State) + } + } +}