Run per-validator OpenStack work in parallel so all validators start at once

The orchestrator claimed idle validators one by one and associated each
floating IP synchronously (~30 s per address), so 20 validators started
about 30 s apart. Aggregation/disassociation and lease reclaim were
sequential in the same way.

The slow OpenStack calls now run in one goroutine per address, guarded by
an in-flight set against duplicates. control-api runs Tick in Async mode
(Tick does not wait); tests keep the waiting mode.

Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
This commit is contained in:
ayurishchevandClaude Sonnet 5.5 committed 2026-10-01 20:37:32 +03:00
1 parent aff8fe38b5
commit c0e300f71c
5 files changed
+115 -14

No files matched your search

+1 -1
View File
@@ -1,4 +1,4 @@
26ff297b66a0c974b7143179e44a58f476265482cd1929bc01a95ad29e50ca7a control-api 4562d3cf3db8075b5a1213b582359991e3f58c819151cf7b25276cba0a4e5985 control-api
091ad94b5706b1778b181542de7a4b251cd16c421144269d0955769603d51e06 validator-agent 091ad94b5706b1778b181542de7a4b251cd16c421144269d0955769603d51e06 validator-agent
3e9e14dbb361ee76aaad7c1da6864b3ea111e0ed151403f904b12485631bbf75 prober 3e9e14dbb361ee76aaad7c1da6864b3ea111e0ed151403f904b12485631bbf75 prober
d72234688eb1954dbff420ca1c8e83b83ceab561b18a0336af0f73fe4d9ac8be admin-dashboard d72234688eb1954dbff420ca1c8e83b83ceab561b18a0336af0f73fe4d9ac8be admin-dashboard
BIN
View File
Binary file not shown.
+1
View File
@@ -59,6 +59,7 @@ func run(configPath string, log *slog.Logger) error {
} }
orch := orchestrator.New(database, osClient, cfg, log) 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 // 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. // as long as the HTTP request or loop iteration that started them.
orch.SetContext(ctx) orch.SetContext(ctx)
+59 -13
View File
@@ -49,6 +49,36 @@ type Orchestrator struct {
// scan is the background floating-IP scan job (see scanjob.go); its zero // scan is the background floating-IP scan job (see scanjob.go); its zero
// value is ready to use. // value is ready to use.
scan scanJob 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, // 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 // validators, sweep the checking window for ready-to-aggregate IPs, and
// reclaim expired leases. Intended to be called on a fixed interval // reclaim expired leases. Intended to be called on a fixed interval
// (Cfg.PollIntervalSeconds) by the caller (cmd/control-api/main.go). // (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) { func (o *Orchestrator) Tick(ctx context.Context) {
if err := o.assignIdleValidators(ctx); err != nil { if err := o.assignIdleValidators(ctx); err != nil {
o.Log.Error("assign idle validators", "err", err) 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 { if err := o.sweepExpiredLeases(ctx); err != nil {
o.Log.Error("sweep expired leases", "err", err) o.Log.Error("sweep expired leases", "err", err)
} }
if !o.Async {
o.wg.Wait()
}
} }
// assignIdleValidators claims the next queued IP for every currently idle // 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 continue // no work available for this validator right now
} }
o.Log.Info("claimed ip", "validator", v.ValidatorID, "ip", item.IPAddress, "ip_id", item.ID) 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 { v, item := v, item
o.Log.Error("associate fip", "validator", v.ValidatorID, "ip", item.IPAddress, "err", err) 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 return nil
} }
@@ -453,9 +493,12 @@ func (o *Orchestrator) sweepCheckingWindow(ctx context.Context) error {
if !ready { if !ready {
continue continue
} }
if err := o.aggregateAndRelease(ctx, item); err != nil { item := item
o.Log.Error("aggregate and release", "ip_id", item.ID, "err", err) 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 return nil
} }
@@ -606,14 +649,17 @@ func (o *Orchestrator) sweepExpiredLeases(ctx context.Context) error {
if item.OwnerValidatorID != nil { if item.OwnerValidatorID != nil {
validatorID = *item.OwnerValidatorID validatorID = *item.OwnerValidatorID
} }
o.Log.Info("lease expired, reclaiming", "ip_id", item.ID, "ip", item.IPAddress, "validator", validatorID) item, validatorID := item, validatorID
if item.FIPID != "" { o.spawn(fmt.Sprintf("lease:%d", item.ID), func() {
if err := o.OS.DisassociateFloatingIP(ctx, item.FIPID); err != nil { o.Log.Info("lease expired, reclaiming", "ip_id", item.ID, "ip", item.IPAddress, "validator", validatorID)
o.Log.Error("disassociate fip on lease reclaim", "ip_id", item.ID, "err", err) 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.event(ctx, "control-api", "", &item.ID, "lease_expired", fmt.Sprintf(`{"validator_id":%q}`, validatorID)) o.requeueOrFail(ctx, item.ID, validatorID, "lease expired")
o.requeueOrFail(ctx, item.ID, validatorID, "lease expired") })
} }
return nil return nil
} }
@@ -3,6 +3,7 @@ package orchestrator
import ( import (
"context" "context"
"errors" "errors"
"fmt"
"log/slog" "log/slog"
"os" "os"
"path/filepath" "path/filepath"
@@ -942,3 +943,56 @@ func TestScanFloatingIPsNoFreeAddressesIsNotAnError(t *testing.T) {
t.Fatalf("expected nothing added, got %+v", result) 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)
}
}
}