package httpapi import ( "bytes" "context" "encoding/json" "io" "log/slog" "net/http" "net/http/httptest" "os" "path/filepath" "testing" "time" "cloudipvalidator/internal/config" "cloudipvalidator/internal/db" "cloudipvalidator/internal/openstack" "cloudipvalidator/internal/orchestrator" ) // fakeClient plays both a validator-agent and the three probers against a // real httptest server, driving the full protocol exactly as the real // binaries would, to prove the HTTP layer and orchestrator agree on state // transitions end-to-end. type fakeClient struct { t *testing.T base string client *http.Client } func (f *fakeClient) do(method, path string, body interface{}) (*http.Response, []byte) { f.t.Helper() var reader io.Reader if body != nil { b, err := json.Marshal(body) if err != nil { f.t.Fatalf("marshal body: %v", err) } reader = bytes.NewReader(b) } req, err := http.NewRequest(method, f.base+path, reader) if err != nil { f.t.Fatalf("new request: %v", err) } req.Header.Set("Content-Type", "application/json") resp, err := f.client.Do(req) if err != nil { f.t.Fatalf("%s %s: %v", method, path, err) } defer resp.Body.Close() respBody, _ := io.ReadAll(resp.Body) return resp, respBody } func TestEndToEndHTTPFlow(t *testing.T) { 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) } defer d.Close() mock := openstack.NewMockClient() mock.Seed("fip-1", "1.2.3.4", "svc-project") cfg := &config.ControlAPI{ Orchestrator: config.OrchestratorConfig{ PollIntervalSeconds: 1, SelfCheckTimeoutSeconds: 10, MaxSelfCheckRetries: 3, CheckingWindowSeconds: 120, MaxRetries: 3, LeaseTTLSeconds: 180, HeartbeatTimeoutSeconds: 30, }, Aggregation: config.AggregationConfig{MissingCountsAsFail: true}, Sites: []config.SiteConfig{ {SiteID: "site-1", Index: 1}, {SiteID: "site-2", Index: 2}, {SiteID: "site-3", Index: 3}, }, CheckTypes: []config.CheckTypeConfig{{Name: "https", Enabled: true, 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})) orch := orchestrator.New(d, mock, cfg, log) if err := d.SeedQueue(ctx, []string{"1.2.3.4"}); err != nil { t.Fatalf("seed queue: %v", err) } srv := New(d, orch, log) ts := httptest.NewServer(srv.Handler()) defer ts.Close() fc := &fakeClient{t: t, base: ts.URL, client: ts.Client()} // Register the validator directly via DB (os_port_id comes from // control-api config, not the agent's own registration call). if err := d.RegisterValidator(ctx, "validator-1", "host-1", "port-1", "v0.1"); err != nil { t.Fatalf("register validator: %v", err) } // Agent re-registers over HTTP (as the real binary would at startup). resp, body := fc.do(http.MethodPost, "/api/v1/agents/register", registerAgentRequest{ ValidatorID: "validator-1", Hostname: "host-1", AgentVersion: "v0.1", }) if resp.StatusCode != http.StatusOK { t.Fatalf("register: status=%d body=%s", resp.StatusCode, body) } // Orchestrator claims the IP and associates the FIP. orch.Tick(ctx) resp, body = fc.do(http.MethodGet, "/api/v1/agents/validator-1/assignment", nil) if resp.StatusCode != http.StatusOK { t.Fatalf("assignment: status=%d body=%s", resp.StatusCode, body) } var assignment assignmentResponse if err := json.Unmarshal(body, &assignment); err != nil { t.Fatalf("unmarshal assignment: %v", err) } if assignment.Phase != db.IPAwaitingSelfCheck { t.Fatalf("expected awaiting_self_check, got %s", assignment.Phase) } if assignment.IPAddress != "1.2.3.4" { t.Fatalf("expected 1.2.3.4, got %s", assignment.IPAddress) } // Self-check: in the real deployment the agent queries an external // IP-echo service (see internal/agentcore.detectPublicIP) and compares // the result to the assigned FIP; the HTTP layer here just accepts // whatever outcome the caller reports. resp, body = fc.do(http.MethodPost, "/api/v1/agents/validator-1/self-check", selfCheckRequest{ IPID: assignment.IPID, DetectedEgress: "1.2.3.4", Success: true, Detail: "matched", }) if resp.StatusCode != http.StatusOK { t.Fatalf("self-check: status=%d body=%s", resp.StatusCode, body) } // Agent runs its configured egress check and reports the result. resp, body = fc.do(http.MethodPost, "/api/v1/agents/validator-1/results", agentResultsRequest{ Results: []checkResultDTO{{ IPID: assignment.IPID, CheckType: "https", Target: "https://example.test", Success: true, LatencyMS: 12, CheckedAt: time.Now().Format(time.RFC3339Nano), }}, }) if resp.StatusCode != http.StatusOK { t.Fatalf("results: status=%d body=%s", resp.StatusCode, body) } resp, body = fc.do(http.MethodPost, "/api/v1/agents/validator-1/complete", agentCompleteRequest{IPID: assignment.IPID}) if resp.StatusCode != http.StatusOK { t.Fatalf("complete: status=%d body=%s", resp.StatusCode, body) } // Three probers register, poll, and report inbound results. for _, site := range []string{"site-1", "site-2", "site-3"} { resp, body = fc.do(http.MethodPost, "/api/v1/probers/register", registerProberRequest{SiteID: site}) if resp.StatusCode != http.StatusOK { t.Fatalf("prober register %s: status=%d body=%s", site, resp.StatusCode, body) } resp, body = fc.do(http.MethodGet, "/api/v1/probers/"+site+"/assignments", nil) if resp.StatusCode != http.StatusOK { t.Fatalf("prober assignments %s: status=%d body=%s", site, resp.StatusCode, body) } var assignments []proberAssignment if err := json.Unmarshal(body, &assignments); err != nil { t.Fatalf("unmarshal assignments: %v", err) } if len(assignments) != 1 || assignments[0].IPAddress != "1.2.3.4" { t.Fatalf("expected 1 assignment for 1.2.3.4, got %+v", assignments) } now := time.Now().Format(time.RFC3339Nano) resp, body = fc.do(http.MethodPost, "/api/v1/probers/"+site+"/results", proberResultsRequest{ Results: []proberResultDTO{ {IPID: assignments[0].IPID, IPAddress: "1.2.3.4", CheckType: "tcp-22", Success: true, CheckedAt: now}, {IPID: assignments[0].IPID, IPAddress: "1.2.3.4", CheckType: "ssh", Success: true, CheckedAt: now}, {IPID: assignments[0].IPID, IPAddress: "1.2.3.4", CheckType: "tcp-80", Success: true, CheckedAt: now}, {IPID: assignments[0].IPID, IPAddress: "1.2.3.4", CheckType: "icmp", Success: true, CheckedAt: now, Complete: true}, }, }) if resp.StatusCode != http.StatusOK { t.Fatalf("prober results %s: status=%d body=%s", site, resp.StatusCode, body) } } // Orchestrator sweep should now aggregate and release. orch.Tick(ctx) resp, body = fc.do(http.MethodGet, "/api/v1/admin/ips/1.2.3.4", nil) if resp.StatusCode != http.StatusOK { t.Fatalf("admin ip detail: status=%d body=%s", resp.StatusCode, body) } var detail struct { IP db.IPQueueItem `json:"ip"` } if err := json.Unmarshal(body, &detail); err != nil { t.Fatalf("unmarshal detail: %v", err) } if detail.IP.State != db.IPDone { t.Fatalf("expected done, got %s", detail.IP.State) } if detail.IP.OverallResult != db.ResultPass { t.Fatalf("expected pass, got %s", detail.IP.OverallResult) } if fip, _ := mock.GetFloatingIPByAddress(ctx, "1.2.3.4"); fip.PortID != "" { t.Fatalf("expected fip disassociated at end of run") } }