208 lines
7.2 KiB
Go
208 lines
7.2 KiB
Go
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},
|
|
}
|
|
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: "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")
|
|
}
|
|
}
|