460 lines
16 KiB
Go
460 lines
16 KiB
Go
package orchestrator
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"log/slog"
|
|
"os"
|
|
"path/filepath"
|
|
"testing"
|
|
"time"
|
|
|
|
"cloudipvalidator/internal/config"
|
|
"cloudipvalidator/internal/db"
|
|
"cloudipvalidator/internal/openstack"
|
|
)
|
|
|
|
var threeSites = []config.SiteConfig{
|
|
{SiteID: "site-1", Index: 1},
|
|
{SiteID: "site-2", Index: 2},
|
|
{SiteID: "site-3", Index: 3},
|
|
}
|
|
|
|
func newTestOrchestrator(t *testing.T, leaseTTLSeconds int) (*Orchestrator, *db.DB, *openstack.MockClient) {
|
|
t.Helper()
|
|
return newTestOrchestratorWithSites(t, leaseTTLSeconds, threeSites)
|
|
}
|
|
|
|
func newTestOrchestratorWithSites(t *testing.T, leaseTTLSeconds int, sites []config.SiteConfig) (*Orchestrator, *db.DB, *openstack.MockClient) {
|
|
t.Helper()
|
|
ctx := context.Background()
|
|
dbPath := filepath.Join(t.TempDir(), "test.db")
|
|
d, err := db.Open(ctx, dbPath)
|
|
if err != nil {
|
|
t.Fatalf("open db: %v", err)
|
|
}
|
|
t.Cleanup(func() { d.Close() })
|
|
|
|
mock := openstack.NewMockClient()
|
|
|
|
cfg := &config.ControlAPI{
|
|
Orchestrator: config.OrchestratorConfig{
|
|
PollIntervalSeconds: 1,
|
|
SelfCheckTimeoutSeconds: 10,
|
|
MaxSelfCheckRetries: 3,
|
|
CheckingWindowSeconds: 120,
|
|
MaxRetries: 3,
|
|
LeaseTTLSeconds: leaseTTLSeconds,
|
|
HeartbeatTimeoutSeconds: 30,
|
|
},
|
|
Aggregation: config.AggregationConfig{MissingCountsAsFail: true},
|
|
Sites: sites,
|
|
CheckTypes: []config.CheckTypeConfig{
|
|
{Name: "https", Enabled: true, Targets: []string{"web"}},
|
|
{Name: "ssh", Enabled: false, Targets: []string{"web"}},
|
|
},
|
|
Targets: map[string][]string{
|
|
"web": {"https://example.test"},
|
|
},
|
|
Inbound: config.InboundConfig{Ports: []int{22, 80}, ICMP: true},
|
|
}
|
|
|
|
if err := d.BootstrapFromConfig(ctx, cfg); err != nil {
|
|
t.Fatalf("bootstrap from config: %v", err)
|
|
}
|
|
|
|
log := slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{Level: slog.LevelError}))
|
|
return New(d, mock, cfg, log), d, mock
|
|
}
|
|
|
|
func TestHappyPath(t *testing.T) {
|
|
ctx := context.Background()
|
|
o, d, mock := newTestOrchestrator(t, 180)
|
|
|
|
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.Fatalf("register validator: %v", err)
|
|
}
|
|
if err := d.SeedQueue(ctx, []string{"1.2.3.4"}); err != nil {
|
|
t.Fatalf("seed queue: %v", err)
|
|
}
|
|
|
|
// 1. claim + associate
|
|
o.Tick(ctx)
|
|
|
|
ip, err := d.GetIPByAddress(ctx, "1.2.3.4")
|
|
if err != nil {
|
|
t.Fatalf("get ip: %v", err)
|
|
}
|
|
if ip.State != db.IPAwaitingSelfCheck {
|
|
t.Fatalf("expected awaiting_self_check, got %s", ip.State)
|
|
}
|
|
if ip.FIPID != "fip-1" {
|
|
t.Fatalf("expected fip-1 associated, got %q", ip.FIPID)
|
|
}
|
|
if fip, _ := mock.GetFloatingIPByAddress(ctx, "1.2.3.4"); fip.PortID != "port-1" {
|
|
t.Fatalf("expected fip associated to port-1, got %q", fip.PortID)
|
|
}
|
|
|
|
v, err := d.GetValidator(ctx, "validator-1")
|
|
if err != nil {
|
|
t.Fatalf("get validator: %v", err)
|
|
}
|
|
if v.State != db.ValidatorAssigned || v.CurrentIPID == nil || *v.CurrentIPID != ip.ID {
|
|
t.Fatalf("expected validator assigned to ip %d, got state=%s current_ip=%v", ip.ID, v.State, v.CurrentIPID)
|
|
}
|
|
|
|
// 2. self-check success
|
|
if err := o.SelfCheckResult(ctx, "validator-1", ip.ID, true, "egress matched"); err != nil {
|
|
t.Fatalf("self check result: %v", err)
|
|
}
|
|
ip, _ = d.GetIP(ctx, ip.ID)
|
|
if ip.State != db.IPChecking {
|
|
t.Fatalf("expected checking, got %s", ip.State)
|
|
}
|
|
|
|
// 3. egress result + completion
|
|
if err := o.RecordCheck(ctx, db.Check{
|
|
IPID: ip.ID, IPAddress: ip.IPAddress, AttemptNumber: ip.AttemptNumber,
|
|
ValidatorID: "validator-1", Source: db.SourceEgress, CheckType: "https",
|
|
Target: "https://example.test", Success: true, CheckedAt: db.Now(),
|
|
}); err != nil {
|
|
t.Fatalf("record egress check: %v", err)
|
|
}
|
|
if err := o.MarkEgressComplete(ctx, ip.ID); err != nil {
|
|
t.Fatalf("mark egress complete: %v", err)
|
|
}
|
|
|
|
// 4. inbound results from all 3 sites
|
|
for site := 1; site <= 3; site++ {
|
|
for _, ct := range []string{"tcp-22", "tcp-80", "icmp"} {
|
|
if err := o.RecordCheck(ctx, db.Check{
|
|
IPID: ip.ID, IPAddress: ip.IPAddress, AttemptNumber: ip.AttemptNumber,
|
|
Source: db.InboundSource(site), CheckType: ct, Target: ip.IPAddress,
|
|
Success: true, CheckedAt: db.Now(),
|
|
}); err != nil {
|
|
t.Fatalf("record inbound check site %d: %v", site, err)
|
|
}
|
|
}
|
|
if err := o.MarkSiteComplete(ctx, ip.ID, site); err != nil {
|
|
t.Fatalf("mark site %d complete: %v", site, err)
|
|
}
|
|
}
|
|
|
|
// 5. sweep should now aggregate + release
|
|
o.Tick(ctx)
|
|
|
|
ip, _ = d.GetIP(ctx, ip.ID)
|
|
if ip.State != db.IPDone {
|
|
t.Fatalf("expected done, got %s", ip.State)
|
|
}
|
|
if ip.OverallResult != db.ResultPass {
|
|
t.Fatalf("expected pass, got %s", ip.OverallResult)
|
|
}
|
|
if ip.FIPReleasedAt == nil {
|
|
t.Fatalf("expected fip_released_at to be set")
|
|
}
|
|
|
|
v, _ = d.GetValidator(ctx, "validator-1")
|
|
if v.State != db.ValidatorIdle || v.CurrentIPID != nil {
|
|
t.Fatalf("expected validator idle with no current ip, got state=%s current_ip=%v", v.State, v.CurrentIPID)
|
|
}
|
|
|
|
if fip, _ := mock.GetFloatingIPByAddress(ctx, "1.2.3.4"); fip.PortID != "" {
|
|
t.Fatalf("expected fip disassociated, still on port %q", fip.PortID)
|
|
}
|
|
}
|
|
|
|
func TestPartialResult(t *testing.T) {
|
|
ctx := context.Background()
|
|
o, d, mock := newTestOrchestrator(t, 180)
|
|
mock.Seed("fip-1", "1.2.3.4", "svc-project")
|
|
_ = d.RegisterValidator(ctx, "validator-1", "host-1", "port-1", "v0.1")
|
|
_ = d.SeedQueue(ctx, []string{"1.2.3.4"})
|
|
|
|
o.Tick(ctx)
|
|
ip, _ := d.GetIPByAddress(ctx, "1.2.3.4")
|
|
_ = o.SelfCheckResult(ctx, "validator-1", ip.ID, true, "ok")
|
|
ip, _ = d.GetIP(ctx, ip.ID)
|
|
|
|
// Egress passes...
|
|
_ = o.RecordCheck(ctx, db.Check{
|
|
IPID: ip.ID, IPAddress: ip.IPAddress, AttemptNumber: ip.AttemptNumber,
|
|
ValidatorID: "validator-1", Source: db.SourceEgress, CheckType: "https",
|
|
Target: "https://example.test", Success: true, CheckedAt: db.Now(),
|
|
})
|
|
_ = o.MarkEgressComplete(ctx, ip.ID)
|
|
// ...but only site-1 reports, and one of its checks fails.
|
|
_ = o.RecordCheck(ctx, db.Check{
|
|
IPID: ip.ID, IPAddress: ip.IPAddress, AttemptNumber: ip.AttemptNumber,
|
|
Source: db.InboundSource(1), CheckType: "tcp-22", Target: ip.IPAddress, Success: false, CheckedAt: db.Now(),
|
|
})
|
|
_ = o.RecordCheck(ctx, db.Check{
|
|
IPID: ip.ID, IPAddress: ip.IPAddress, AttemptNumber: ip.AttemptNumber,
|
|
Source: db.InboundSource(1), CheckType: "tcp-80", Target: ip.IPAddress, Success: true, CheckedAt: db.Now(),
|
|
})
|
|
_ = o.RecordCheck(ctx, db.Check{
|
|
IPID: ip.ID, IPAddress: ip.IPAddress, AttemptNumber: ip.AttemptNumber,
|
|
Source: db.InboundSource(1), CheckType: "icmp", Target: ip.IPAddress, Success: true, CheckedAt: db.Now(),
|
|
})
|
|
_ = o.MarkSiteComplete(ctx, ip.ID, 1)
|
|
|
|
// Force the checking window to have elapsed so aggregation proceeds
|
|
// even though site-2/site-3 never reported.
|
|
o.Cfg.CheckingWindowSeconds = 0
|
|
time.Sleep(5 * time.Millisecond)
|
|
o.Tick(ctx)
|
|
|
|
ip, _ = d.GetIP(ctx, ip.ID)
|
|
if ip.State != db.IPDone {
|
|
t.Fatalf("expected done, got %s", ip.State)
|
|
}
|
|
if ip.OverallResult != db.ResultPartial {
|
|
t.Fatalf("expected partial, got %s", ip.OverallResult)
|
|
}
|
|
if fip, _ := mock.GetFloatingIPByAddress(ctx, "1.2.3.4"); fip.PortID != "" {
|
|
t.Fatalf("expected fip disassociated even on partial result")
|
|
}
|
|
}
|
|
|
|
func TestLeaseReclaim(t *testing.T) {
|
|
ctx := context.Background()
|
|
// A 1s lease (rather than 0) avoids a race within the very first Tick:
|
|
// with a 0s TTL the item's lease can already look expired by the time
|
|
// the same Tick's lease-sweep phase runs, depending on how much
|
|
// wall-clock time the claim+associate phase happened to take.
|
|
o, d, mock := newTestOrchestrator(t, 1)
|
|
mock.Seed("fip-1", "1.2.3.4", "svc-project")
|
|
_ = d.RegisterValidator(ctx, "validator-1", "host-1", "port-1", "v0.1")
|
|
_ = d.SeedQueue(ctx, []string{"1.2.3.4"})
|
|
|
|
o.Tick(ctx) // claims + associates; validator never self-checks
|
|
ip, _ := d.GetIPByAddress(ctx, "1.2.3.4")
|
|
if ip.State != db.IPAwaitingSelfCheck {
|
|
t.Fatalf("expected awaiting_self_check, got %s", ip.State)
|
|
}
|
|
|
|
time.Sleep(1100 * time.Millisecond) // let the 1s lease expire
|
|
o.Tick(ctx) // should reclaim via lease sweep
|
|
|
|
ip, _ = d.GetIP(ctx, ip.ID)
|
|
if ip.State != db.IPQueued {
|
|
t.Fatalf("expected requeued after lease reclaim, got %s (retry_count=%d)", ip.State, ip.RetryCount)
|
|
}
|
|
if ip.RetryCount != 1 {
|
|
t.Fatalf("expected retry_count=1, got %d", ip.RetryCount)
|
|
}
|
|
|
|
v, _ := d.GetValidator(ctx, "validator-1")
|
|
if v.State != db.ValidatorIdle || v.CurrentIPID != nil {
|
|
t.Fatalf("expected validator freed, got state=%s current_ip=%v", v.State, v.CurrentIPID)
|
|
}
|
|
|
|
if fip, _ := mock.GetFloatingIPByAddress(ctx, "1.2.3.4"); fip.PortID != "" {
|
|
t.Fatalf("expected fip disassociated on reclaim, still on port %q", fip.PortID)
|
|
}
|
|
|
|
// A subsequent tick should re-claim and re-associate the same IP for
|
|
// the now-idle validator, proving the queue keeps making progress.
|
|
// Give this attempt a real lease so it isn't immediately re-expired by
|
|
// the same tick's lease sweep (a 0s TTL, as above, expires instantly).
|
|
o.Cfg.LeaseTTLSeconds = 180
|
|
o.Tick(ctx)
|
|
ip, _ = d.GetIP(ctx, ip.ID)
|
|
if ip.State != db.IPAwaitingSelfCheck {
|
|
t.Fatalf("expected re-claimed ip to be awaiting_self_check again, got %s", ip.State)
|
|
}
|
|
if ip.AttemptNumber != 2 {
|
|
t.Fatalf("expected attempt_number=2 after reclaim+reassign, got %d", ip.AttemptNumber)
|
|
}
|
|
}
|
|
|
|
func TestMaxRetriesExhausted(t *testing.T) {
|
|
ctx := context.Background()
|
|
// A 1s lease (rather than 0) avoids the same intra-tick race noted in
|
|
// TestLeaseReclaim: with a 0s TTL, whether a freshly claimed item is
|
|
// reclaimed within the very same Tick (making each iteration's timing
|
|
// unpredictable) depends on how much wall-clock time claim+associate
|
|
// happened to take, which made this test flaky under load.
|
|
o, d, mock := newTestOrchestrator(t, 1)
|
|
o.Cfg.MaxRetries = 1
|
|
mock.Seed("fip-1", "1.2.3.4", "svc-project")
|
|
_ = d.RegisterValidator(ctx, "validator-1", "host-1", "port-1", "v0.1")
|
|
_ = d.SeedQueue(ctx, []string{"1.2.3.4"})
|
|
|
|
// Each (claim, sleep past 1s lease) pair reclaims once, bumping
|
|
// retry_count by 1; with MaxRetries=1 that needs two full reclaim
|
|
// cycles (retry_count 0->1 requeues, 1->2 fails) — four ticks total:
|
|
// claim, reclaim+requeue, re-claim, reclaim+fail.
|
|
for i := 0; i < 4; i++ {
|
|
o.Tick(ctx)
|
|
time.Sleep(1100 * time.Millisecond)
|
|
}
|
|
|
|
ip, _ := d.GetIPByAddress(ctx, "1.2.3.4")
|
|
if ip.State != db.IPFailed {
|
|
t.Fatalf("expected failed after exhausting retries, got %s (retry_count=%d)", ip.State, ip.RetryCount)
|
|
}
|
|
}
|
|
|
|
// TestInboundChecksDisabled confirms inbound (prober) checks are genuinely
|
|
// optional: with no sites configured, an IP must aggregate as soon as
|
|
// egress completes, without ever waiting on siteN_complete flags that
|
|
// nothing will ever set — and, critically, without waiting out the full
|
|
// checking_window_seconds timeout to get there (newTestOrchestrator uses
|
|
// 120s; this test never sleeps, so a pass here proves the "all required
|
|
// sources complete" path fired, not the timeout fallback).
|
|
func TestInboundChecksDisabled(t *testing.T) {
|
|
ctx := context.Background()
|
|
o, d, mock := newTestOrchestratorWithSites(t, 180, nil)
|
|
mock.Seed("fip-1", "1.2.3.4", "svc-project")
|
|
_ = d.RegisterValidator(ctx, "validator-1", "host-1", "port-1", "v0.1")
|
|
_ = d.SeedQueue(ctx, []string{"1.2.3.4"})
|
|
|
|
o.Tick(ctx)
|
|
ip, _ := d.GetIPByAddress(ctx, "1.2.3.4")
|
|
_ = o.SelfCheckResult(ctx, "validator-1", ip.ID, true, "ok")
|
|
ip, _ = d.GetIP(ctx, ip.ID)
|
|
|
|
_ = o.RecordCheck(ctx, db.Check{
|
|
IPID: ip.ID, IPAddress: ip.IPAddress, AttemptNumber: ip.AttemptNumber,
|
|
ValidatorID: "validator-1", Source: db.SourceEgress, CheckType: "https",
|
|
Target: "https://example.test", Success: true, CheckedAt: db.Now(),
|
|
})
|
|
_ = o.MarkEgressComplete(ctx, ip.ID)
|
|
|
|
o.Tick(ctx)
|
|
|
|
ip, _ = d.GetIP(ctx, ip.ID)
|
|
if ip.State != db.IPDone {
|
|
t.Fatalf("expected done immediately after egress completed with no sites configured, got %s", ip.State)
|
|
}
|
|
if ip.OverallResult != db.ResultPass {
|
|
t.Fatalf("expected pass, got %s", ip.OverallResult)
|
|
}
|
|
}
|
|
|
|
// TestDeleteIPDisassociatesFIP proves DeleteIP disassociates a currently
|
|
// attached floating IP (via the mock) before permanently removing the
|
|
// address — the same resource-freeing ForceCancel does, but going straight
|
|
// to physical deletion instead of a `cancelled` record.
|
|
func TestDeleteIPDisassociatesFIP(t *testing.T) {
|
|
ctx := context.Background()
|
|
o, d, mock := newTestOrchestrator(t, 180)
|
|
mock.Seed("fip-1", "1.2.3.4", "svc-project")
|
|
_ = d.RegisterValidator(ctx, "validator-1", "host-1", "port-1", "v0.1")
|
|
_ = d.SeedQueue(ctx, []string{"1.2.3.4"})
|
|
|
|
o.Tick(ctx) // claim + associate -> awaiting_self_check, fip attached
|
|
|
|
ip, err := d.GetIPByAddress(ctx, "1.2.3.4")
|
|
if err != nil {
|
|
t.Fatalf("get ip: %v", err)
|
|
}
|
|
if ip.FIPID == "" {
|
|
t.Fatalf("expected fip associated before delete")
|
|
}
|
|
|
|
if err := o.DeleteIP(ctx, "1.2.3.4"); err != nil {
|
|
t.Fatalf("delete ip: %v", err)
|
|
}
|
|
|
|
if fip, _ := mock.GetFloatingIPByAddress(ctx, "1.2.3.4"); fip.PortID != "" {
|
|
t.Fatalf("expected fip disassociated on delete, still on port %q", fip.PortID)
|
|
}
|
|
if _, err := d.GetIPByAddress(ctx, "1.2.3.4"); err == nil {
|
|
t.Fatalf("expected ip row gone after delete")
|
|
}
|
|
v, err := d.GetValidator(ctx, "validator-1")
|
|
if err != nil {
|
|
t.Fatalf("get validator: %v", err)
|
|
}
|
|
if v.State != db.ValidatorIdle || v.CurrentIPID != nil {
|
|
t.Fatalf("expected validator freed, got state=%s current_ip=%v", v.State, v.CurrentIPID)
|
|
}
|
|
|
|
if err := o.DeleteIP(ctx, "1.2.3.4"); !errors.Is(err, db.ErrNotFound) {
|
|
t.Fatalf("expected ErrNotFound deleting already-gone ip, got %v", err)
|
|
}
|
|
}
|
|
|
|
// TestClearQueueDisassociatesAllFIPs proves ClearQueue deletes every
|
|
// address regardless of state and disassociates any attached floating IPs
|
|
// along the way.
|
|
func TestClearQueueDisassociatesAllFIPs(t *testing.T) {
|
|
ctx := context.Background()
|
|
o, d, mock := newTestOrchestrator(t, 180)
|
|
mock.Seed("fip-1", "1.2.3.4", "svc-project")
|
|
_ = d.RegisterValidator(ctx, "validator-1", "host-1", "port-1", "v0.1")
|
|
_ = d.SeedQueue(ctx, []string{"1.2.3.4", "5.6.7.8"})
|
|
|
|
o.Tick(ctx) // claims + associates 1.2.3.4; 5.6.7.8 stays queued
|
|
|
|
result, err := o.ClearQueue(ctx)
|
|
if err != nil {
|
|
t.Fatalf("clear queue: %v", err)
|
|
}
|
|
if len(result.Deleted) != 2 {
|
|
t.Fatalf("expected both addresses deleted, got %+v", result)
|
|
}
|
|
if fip, _ := mock.GetFloatingIPByAddress(ctx, "1.2.3.4"); fip.PortID != "" {
|
|
t.Fatalf("expected fip disassociated on clear, still on port %q", fip.PortID)
|
|
}
|
|
ips, err := d.ListIPs(ctx)
|
|
if err != nil {
|
|
t.Fatalf("list ips: %v", err)
|
|
}
|
|
if len(ips) != 0 {
|
|
t.Fatalf("expected empty queue after clear, got %+v", ips)
|
|
}
|
|
}
|
|
|
|
// TestInboundChecksPartialSites confirms a partially-configured sites list
|
|
// (fewer than 3 slots assigned) only waits on the sites actually
|
|
// configured — the two unassigned slots are never expected.
|
|
func TestInboundChecksPartialSites(t *testing.T) {
|
|
ctx := context.Background()
|
|
sites := []config.SiteConfig{{SiteID: "site-1", Index: 1}}
|
|
o, d, mock := newTestOrchestratorWithSites(t, 180, sites)
|
|
mock.Seed("fip-1", "1.2.3.4", "svc-project")
|
|
_ = d.RegisterValidator(ctx, "validator-1", "host-1", "port-1", "v0.1")
|
|
_ = d.SeedQueue(ctx, []string{"1.2.3.4"})
|
|
|
|
o.Tick(ctx)
|
|
ip, _ := d.GetIPByAddress(ctx, "1.2.3.4")
|
|
_ = o.SelfCheckResult(ctx, "validator-1", ip.ID, true, "ok")
|
|
ip, _ = d.GetIP(ctx, ip.ID)
|
|
|
|
_ = o.RecordCheck(ctx, db.Check{
|
|
IPID: ip.ID, IPAddress: ip.IPAddress, AttemptNumber: ip.AttemptNumber,
|
|
ValidatorID: "validator-1", Source: db.SourceEgress, CheckType: "https",
|
|
Target: "https://example.test", Success: true, CheckedAt: db.Now(),
|
|
})
|
|
_ = o.MarkEgressComplete(ctx, ip.ID)
|
|
|
|
// Egress is done but the one configured site (site-1) hasn't reported
|
|
// yet — must not aggregate.
|
|
o.Tick(ctx)
|
|
ip, _ = d.GetIP(ctx, ip.ID)
|
|
if ip.State != db.IPChecking {
|
|
t.Fatalf("expected still checking (site-1 pending), got %s", ip.State)
|
|
}
|
|
|
|
for _, ct := range []string{"tcp-22", "tcp-80", "icmp"} {
|
|
_ = o.RecordCheck(ctx, db.Check{
|
|
IPID: ip.ID, IPAddress: ip.IPAddress, AttemptNumber: ip.AttemptNumber,
|
|
Source: db.InboundSource(1), CheckType: ct, Target: ip.IPAddress, Success: true, CheckedAt: db.Now(),
|
|
})
|
|
}
|
|
_ = o.MarkSiteComplete(ctx, ip.ID, 1)
|
|
|
|
o.Tick(ctx)
|
|
ip, _ = d.GetIP(ctx, ip.ID)
|
|
if ip.State != db.IPDone {
|
|
t.Fatalf("expected done once the single configured site reported, got %s", ip.State)
|
|
}
|
|
if ip.OverallResult != db.ResultPass {
|
|
t.Fatalf("expected pass, got %s", ip.OverallResult)
|
|
}
|
|
}
|