diff --git a/bin/SHA256SUMS b/bin/SHA256SUMS index 5834d8e..90cd071 100644 --- a/bin/SHA256SUMS +++ b/bin/SHA256SUMS @@ -1,4 +1,4 @@ -4562d3cf3db8075b5a1213b582359991e3f58c819151cf7b25276cba0a4e5985 control-api +4302c5d936b0f6567e4b3b7af3de0123afd862b3cdda53859fad0c7eec443c33 control-api 091ad94b5706b1778b181542de7a4b251cd16c421144269d0955769603d51e06 validator-agent 3e9e14dbb361ee76aaad7c1da6864b3ea111e0ed151403f904b12485631bbf75 prober d72234688eb1954dbff420ca1c8e83b83ceab561b18a0336af0f73fe4d9ac8be admin-dashboard diff --git a/bin/control-api b/bin/control-api index 1c08640..d39fc43 100755 Binary files a/bin/control-api and b/bin/control-api differ diff --git a/internal/db/queries_ipqueue.go b/internal/db/queries_ipqueue.go index b10685e..abd0588 100644 --- a/internal/db/queries_ipqueue.go +++ b/internal/db/queries_ipqueue.go @@ -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 { diff --git a/internal/openstack/client.go b/internal/openstack/client.go index 13b6ddd..c721ad6 100644 --- a/internal/openstack/client.go +++ b/internal/openstack/client.go @@ -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) } diff --git a/internal/openstack/interface.go b/internal/openstack/interface.go index 385147e..731a359 100644 --- a/internal/openstack/interface.go +++ b/internal/openstack/interface.go @@ -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 diff --git a/internal/openstack/mock.go b/internal/openstack/mock.go index 8a9c67f..4c03952 100644 --- a/internal/openstack/mock.go +++ b/internal/openstack/mock.go @@ -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() diff --git a/internal/orchestrator/orchestrator.go b/internal/orchestrator/orchestrator.go index 11f4c30..00ac042 100644 --- a/internal/orchestrator/orchestrator.go +++ b/internal/orchestrator/orchestrator.go @@ -187,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)) @@ -328,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 @@ -359,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", "") @@ -387,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 @@ -412,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 } @@ -440,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 } diff --git a/internal/orchestrator/orchestrator_test.go b/internal/orchestrator/orchestrator_test.go index 61fe1d4..e00ecd5 100644 --- a/internal/orchestrator/orchestrator_test.go +++ b/internal/orchestrator/orchestrator_test.go @@ -996,3 +996,93 @@ func TestValidatorsStartInParallel(t *testing.T) { } } } + +// 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) + } +}