Compare commits

...
2 Commits
Author SHA1 Message Date
ayurishchevandClaude Sonnet 5.5 1ee5757004 Free validator ports in the cloud on clear/cancel/delete
An address still being associated (assigning_fip) has no fip_id in the
database, so "clear queue" did not detach it, and the association then
finished after the row was gone, leaving the floating IP on the validator
port for good.

- After clear/cancel/delete, ask the cloud which floating IPs sit on the
  affected validator ports (new ListFloatingIPsByPort) and detach those
  that this system queued (known in ip_registry); foreign ones are left.
- SetFIPAssociated applies only to a row still in assigning_fip; if the
  address was removed meanwhile, associateFIP detaches the floating IP.

Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
2026-10-01 21:00:31 +03:00
ayurishchevandClaude Sonnet 5.5 c0e300f71c 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>
2026-10-01 20:37:32 +03:00
9 changed files with 401 additions and 18 deletions

No files matched your search

+1 -1
View File
@@ -1,4 +1,4 @@
26ff297b66a0c974b7143179e44a58f476265482cd1929bc01a95ad29e50ca7a control-api
4302c5d936b0f6567e4b3b7af3de0123afd862b3cdda53859fad0c7eec443c33 control-api
091ad94b5706b1778b181542de7a4b251cd16c421144269d0955769603d51e06 validator-agent
3e9e14dbb361ee76aaad7c1da6864b3ea111e0ed151403f904b12485631bbf75 prober
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.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)
+50 -4
View File
@@ -122,13 +122,59 @@ func (d *DB) ClaimNextQueued(ctx context.Context, validatorID string, leaseTTL t
return &item, nil
}
// SetFIPAssociated records that the floating IP is now attached. It only
// applies to a row still in assigning_fip: if the address was deleted or
// cancelled while the cloud call was in flight, it returns ErrInvalidState
// and the caller must detach the floating IP again.
func (d *DB) SetFIPAssociated(ctx context.Context, ipID int64, fipID string, leaseTTL time.Duration) error {
now := Now()
_, err := d.ExecContext(ctx, `
res, err := d.ExecContext(ctx, `
UPDATE ip_queue SET state=?, fip_id=?, fip_associated_at=?, lease_expires_at=?, updated_at=?
WHERE id=?
`, IPAwaitingSelfCheck, fipID, timeToDB(now), timeToDB(now.Add(leaseTTL)), timeToDB(now), ipID)
return err
WHERE id=? AND state=?
`, IPAwaitingSelfCheck, fipID, timeToDB(now), timeToDB(now.Add(leaseTTL)), timeToDB(now), ipID, IPAssigningFIP)
if err != nil {
return err
}
if n, _ := res.RowsAffected(); n == 0 {
return fmt.Errorf("ip_id %d is no longer assigning_fip: %w", ipID, ErrInvalidState)
}
return nil
}
// KnownAddresses returns the subset of addresses present in ip_registry,
// i.e. addresses this system has ever queued.
func (d *DB) KnownAddresses(ctx context.Context, addresses []string) (map[string]bool, error) {
const chunk = 500
known := make(map[string]bool, len(addresses))
for start := 0; start < len(addresses); start += chunk {
end := start + chunk
if end > len(addresses) {
end = len(addresses)
}
part := addresses[start:end]
args := make([]any, len(part))
for i, a := range part {
args[i] = a
}
rows, err := d.QueryContext(ctx, `SELECT ip_address FROM ip_registry WHERE ip_address IN (`+placeholders(len(part))+`)`, args...)
if err != nil {
return nil, err
}
for rows.Next() {
var a string
if err := rows.Scan(&a); err != nil {
rows.Close()
return nil, err
}
known[a] = true
}
if err := rows.Err(); err != nil {
rows.Close()
return nil, err
}
rows.Close()
}
return known, nil
}
func (d *DB) SetChecking(ctx context.Context, ipID int64, leaseTTL time.Duration) error {
+21
View File
@@ -239,6 +239,27 @@ func (c *Client) ListFreeFloatingIPs(ctx context.Context, pageSize int, onPage f
return paginate(ctx, pageSize, c.retry, c.fetchPage, onPage)
}
// ListFloatingIPsByPort asks Neutron for the floating IPs attached to portID.
func (c *Client) ListFloatingIPsByPort(ctx context.Context, portID string) ([]FloatingIP, error) {
pages, err := floatingips.List(c.networking, floatingips.ListOpts{PortID: portID}).AllPages(ctx)
if err != nil {
return nil, fmt.Errorf("openstack: list floating ips of port %s: %w", portID, err)
}
list, err := floatingips.ExtractFloatingIPs(pages)
if err != nil {
return nil, fmt.Errorf("openstack: extract floating ips: %w", err)
}
out := make([]FloatingIP, 0, len(list))
for _, f := range list {
proj := f.TenantID
if proj == "" {
proj = f.ProjectID
}
out = append(out, FloatingIP{ID: f.ID, Address: f.FloatingIP, PortID: f.PortID, ProjectID: proj})
}
return out, nil
}
func (c *Client) ListFloatingIPs(ctx context.Context) ([]FloatingIP, error) {
return listAll(ctx, c)
}
+5
View File
@@ -39,6 +39,11 @@ type FloatingIPClient interface {
// onPage stops the walk and is returned as-is.
ListFreeFloatingIPs(ctx context.Context, pageSize int, onPage func(page []FloatingIP) error) (pages int, err error)
// ListFloatingIPsByPort returns the floating IPs currently attached to
// the given port — read from the cloud, not from our database, so it
// also finds attachments we lost track of.
ListFloatingIPsByPort(ctx context.Context, portID string) ([]FloatingIP, error)
// AssociateFloatingIP attaches the floating IP to the given Neutron
// port (the validator's primary NIC port).
AssociateFloatingIP(ctx context.Context, fipID, portID string) error
+12
View File
@@ -159,6 +159,18 @@ func (m *MockClient) ListFloatingIPs(ctx context.Context) ([]FloatingIP, error)
return listAll(ctx, m)
}
func (m *MockClient) ListFloatingIPsByPort(ctx context.Context, portID string) ([]FloatingIP, error) {
m.mu.Lock()
defer m.mu.Unlock()
var out []FloatingIP
for _, f := range m.fips {
if f.PortID == portID {
out = append(out, *f)
}
}
return out, nil
}
func (m *MockClient) AssociateFloatingIP(ctx context.Context, fipID, portID string) error {
m.mu.Lock()
defer m.mu.Unlock()
+167 -13
View File
@@ -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
}
@@ -147,6 +187,16 @@ func (o *Orchestrator) associateFIP(ctx context.Context, validatorID, osPortID s
return err
}
if err := o.DB.SetFIPAssociated(ctx, item.ID, fip.ID, o.leaseTTL()); err != nil {
if errors.Is(err, db.ErrInvalidState) {
// The address was deleted or cancelled while the cloud call was
// running: nobody owns this attachment any more, so detach it,
// otherwise it stays on the validator's port forever.
o.Log.Info("address removed during association, detaching floating ip", "ip", item.IPAddress, "fip_id", fip.ID)
if derr := o.OS.DisassociateFloatingIP(ctx, fip.ID); derr != nil {
o.Log.Error("disassociate fip of removed address", "ip", item.IPAddress, "fip_id", fip.ID, "err", derr)
}
return nil
}
return fmt.Errorf("set fip associated: %w", err)
}
o.event(ctx, "control-api", "", &item.ID, "fip_associated", fmt.Sprintf(`{"fip_id":%q,"validator_id":%q}`, fip.ID, validatorID))
@@ -288,6 +338,87 @@ func (o *Orchestrator) MarkSiteComplete(ctx context.Context, ipID int64, siteInd
return o.DB.SetSiteComplete(ctx, ipID, siteIndex)
}
// releaseValidatorPorts detaches, in the cloud, every floating IP that sits on
// the given validators' ports and belongs to an address this system has
// queued. The cloud is asked directly because our database does not know
// about an attachment that is still being created (state assigning_fip has no
// fip_id yet) or that was lost earlier. Floating IPs this system never
// queued are left alone. Ports are processed in parallel.
func (o *Orchestrator) releaseValidatorPorts(ctx context.Context, validators []db.Validator) {
var wg sync.WaitGroup
for _, v := range validators {
if v.OSPortID == "" {
continue
}
wg.Add(1)
go func(v db.Validator) {
defer wg.Done()
fips, err := o.OS.ListFloatingIPsByPort(ctx, v.OSPortID)
if err != nil {
o.Log.Error("list floating ips of validator port", "validator", v.ValidatorID, "port", v.OSPortID, "err", err)
return
}
addrs := make([]string, len(fips))
for i, f := range fips {
addrs[i] = f.Address
}
known, err := o.DB.KnownAddresses(ctx, addrs)
if err != nil {
o.Log.Error("check known addresses", "validator", v.ValidatorID, "err", err)
return
}
for _, f := range fips {
if !known[f.Address] {
o.Log.Warn("floating ip on validator port was not queued by this system, left attached",
"validator", v.ValidatorID, "ip", f.Address)
continue
}
if err := o.OS.DisassociateFloatingIP(ctx, f.ID); err != nil {
o.Log.Error("detach floating ip from validator port", "validator", v.ValidatorID, "ip", f.Address, "err", err)
continue
}
o.Log.Info("detached floating ip from validator port", "validator", v.ValidatorID, "ip", f.Address)
}
}(v)
}
wg.Wait()
}
// validatorsByID returns the validators with the given ids (unknown ids are skipped).
func (o *Orchestrator) validatorsByID(ctx context.Context, ids ...string) []db.Validator {
var out []db.Validator
for _, id := range ids {
if v, err := o.DB.GetValidator(ctx, id); err == nil {
out = append(out, *v)
}
}
return out
}
// busyValidatorsFor returns the validators currently working on one of the
// given addresses.
func (o *Orchestrator) busyValidatorsFor(ctx context.Context, addresses []string) []db.Validator {
want := make(map[string]bool, len(addresses))
for _, a := range addresses {
want[a] = true
}
all, err := o.DB.ListValidators(ctx)
if err != nil {
o.Log.Error("list validators", "err", err)
return nil
}
var out []db.Validator
for _, v := range all {
if v.CurrentIPID == nil {
continue
}
if ip, err := o.DB.GetIP(ctx, *v.CurrentIPID); err == nil && want[ip.IPAddress] {
out = append(out, v)
}
}
return out
}
// ForceCancel stops an in-progress (or still-queued) check for the given
// address on admin request, even though it was never going to finish on
// its own within the checking window. If a floating IP is currently
@@ -319,6 +450,7 @@ func (o *Orchestrator) ForceCancel(ctx context.Context, ipAddress string) error
if err := o.DB.FreeValidator(ctx, *item.OwnerValidatorID); err != nil {
return fmt.Errorf("free validator: %w", err)
}
o.releaseValidatorPorts(ctx, o.validatorsByID(ctx, *item.OwnerValidatorID))
}
o.event(ctx, "control-api", "", &item.ID, "force_cancel", "")
@@ -347,9 +479,14 @@ func (o *Orchestrator) DeleteIP(ctx context.Context, ipAddress string) error {
}
}
var owners []db.Validator
if item.OwnerValidatorID != nil {
owners = o.validatorsByID(ctx, *item.OwnerValidatorID)
}
if err := o.DB.DeleteIP(ctx, item.ID); err != nil {
return err
}
o.releaseValidatorPorts(ctx, owners)
o.event(ctx, "control-api", "", nil, "ip_deleted", fmt.Sprintf(`{"ip_address":%q}`, ipAddress))
return nil
@@ -372,10 +509,12 @@ func (o *Orchestrator) DeleteIPs(ctx context.Context, addresses []string) (db.De
}
}
owners := o.busyValidatorsFor(ctx, addresses)
result, err := o.DB.DeleteIPs(ctx, addresses)
if err != nil {
return result, err
}
o.releaseValidatorPorts(ctx, owners)
o.event(ctx, "control-api", "", nil, "ips_deleted", deletedAddressesPayload(result.Deleted))
return result, nil
}
@@ -400,6 +539,15 @@ func (o *Orchestrator) ClearQueue(ctx context.Context) (db.DeleteIPsResult, erro
if err != nil {
return db.DeleteIPsResult{}, err
}
// After the rows are gone, free every validator port in the cloud: this
// catches attachments that were still being created (no fip_id yet) and
// ones left over from earlier runs. An association that finishes after
// this point detaches itself (see associateFIP).
if validators, err := o.DB.ListValidators(ctx); err != nil {
o.Log.Error("list validators to release ports", "err", err)
} else {
o.releaseValidatorPorts(ctx, validators)
}
o.event(ctx, "control-api", "", nil, "queue_cleared", deletedAddressesPayload(deleted))
return db.DeleteIPsResult{Deleted: deleted}, nil
}
@@ -453,9 +601,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 +757,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
}
+144
View File
@@ -3,6 +3,7 @@ package orchestrator
import (
"context"
"errors"
"fmt"
"log/slog"
"os"
"path/filepath"
@@ -942,3 +943,146 @@ 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)
}
}
}
// gatedOS blocks AssociateFloatingIP until release is closed, to hold an
// association "in flight" while the test clears the queue.
type gatedOS struct {
*openstack.MockClient
started chan struct{}
release chan struct{}
}
func (g gatedOS) AssociateFloatingIP(ctx context.Context, fipID, portID string) error {
g.started <- struct{}{}
<-g.release
return g.MockClient.AssociateFloatingIP(ctx, fipID, portID)
}
func portFIPs(t *testing.T, mock *openstack.MockClient, port string) []openstack.FloatingIP {
t.Helper()
got, err := mock.ListFloatingIPsByPort(context.Background(), port)
if err != nil {
t.Fatalf("list by port: %v", err)
}
return got
}
// Clear queue while an association is still running in the cloud: the
// validator port must end up free once the association finishes.
func TestClearQueueFreesPortOfInFlightAssociation(t *testing.T) {
ctx := context.Background()
o, d, mock := newTestOrchestrator(t, 180)
g := gatedOS{MockClient: mock, started: make(chan struct{}, 1), release: make(chan struct{})}
o.OS = g
o.Async = true
mock.Seed("fip-1", "1.2.3.4", "svc-project")
if err := d.RegisterValidator(ctx, "validator-1", "host-1", "port-1", "v0.1"); err != nil {
t.Fatal(err)
}
if err := d.SeedQueue(ctx, []string{"1.2.3.4"}); err != nil {
t.Fatal(err)
}
o.Tick(ctx) // claims the address and starts the association
<-g.started
if _, err := o.ClearQueue(ctx); err != nil {
t.Fatalf("clear queue: %v", err)
}
close(g.release) // the association completes after the clear
o.Wait()
if got := portFIPs(t, mock, "port-1"); len(got) != 0 {
t.Fatalf("port-1 still holds %d floating ip(s): %+v", len(got), got)
}
}
// Clear queue frees ports that already hold a queued address, including an
// attachment left over from earlier, but never touches a floating IP the
// system did not queue.
func TestClearQueueFreesAllValidatorPorts(t *testing.T) {
ctx := context.Background()
o, d, mock := newTestOrchestrator(t, 180)
mock.Seed("fip-1", "1.2.3.4", "svc-project")
mock.SeedWithPort("fip-leak", "1.2.3.9", "svc-project", "port-2") // leaked earlier
mock.SeedWithPort("fip-foreign", "9.9.9.9", "svc-project", "port-2")
if err := d.RegisterValidator(ctx, "validator-1", "h", "port-1", "v"); err != nil {
t.Fatal(err)
}
if err := d.SeedQueue(ctx, []string{"1.2.3.4", "1.2.3.9"}); err != nil {
t.Fatal(err)
}
o.Tick(ctx) // validator-1 gets 1.2.3.4 attached; 1.2.3.9 stays queued, its stale attachment on port-2 is unknown to the queue
if err := d.RegisterValidator(ctx, "validator-2", "h", "port-2", "v"); err != nil {
t.Fatal(err)
}
if got := portFIPs(t, mock, "port-1"); len(got) != 1 {
t.Fatalf("setup: expected 1 floating ip on port-1, got %d", len(got))
}
if _, err := o.ClearQueue(ctx); err != nil {
t.Fatalf("clear queue: %v", err)
}
if got := portFIPs(t, mock, "port-1"); len(got) != 0 {
t.Fatalf("port-1 not free: %+v", got)
}
got := portFIPs(t, mock, "port-2")
if len(got) != 1 || got[0].Address != "9.9.9.9" {
t.Fatalf("port-2: expected only the foreign 9.9.9.9 to remain, got %+v", got)
}
}