Compare commits
2
Commits
aff8fe38b5
...
1ee5757004
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1ee5757004 | ||
|
|
c0e300f71c |
No files matched your search
+1
-1
@@ -1,4 +1,4 @@
|
||||
26ff297b66a0c974b7143179e44a58f476265482cd1929bc01a95ad29e50ca7a control-api
|
||||
4302c5d936b0f6567e4b3b7af3de0123afd862b3cdda53859fad0c7eec443c33 control-api
|
||||
091ad94b5706b1778b181542de7a4b251cd16c421144269d0955769603d51e06 validator-agent
|
||||
3e9e14dbb361ee76aaad7c1da6864b3ea111e0ed151403f904b12485631bbf75 prober
|
||||
d72234688eb1954dbff420ca1c8e83b83ceab561b18a0336af0f73fe4d9ac8be admin-dashboard
|
||||
Binary file not shown.
@@ -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)
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user