diff --git a/bin/SHA256SUMS b/bin/SHA256SUMS index 7582eb7..ecffbe7 100644 --- a/bin/SHA256SUMS +++ b/bin/SHA256SUMS @@ -1,4 +1,4 @@ -8f2ad37b7819131a9b25eee9d35a57880c7aa233bbb741f4fd8dc0af33f6241c control-api -3288930ce4e09e6793048103e25991f025f71dbedced39dee9bccf821eb72a3d validator-agent -586cb533631dfcef2246cdb69e20777f882506755d15afe475399356c1ed2d02 prober -4729c33a2bce584cba29a8833b0701c0a1cd98bb17bf331aabffc27c4cde940c admin-dashboard +dcb137ae6a3f6a59465f2772215d48f98fd303dd1c266645b20ec70194ab7915 control-api +b7a6068db1d095ae7b73cbd9b7273d6a1629ae017e9c3c9601e103432894a247 validator-agent +be8e9576c6fbe5e5798b5d4b758a5617b63bea6abe80e326998d10583edbbf5c prober +579253ac22d6b9537870cf45aaa85437a4a9dba0c87eab9ba8c49c28e50ee438 admin-dashboard diff --git a/bin/admin-dashboard b/bin/admin-dashboard index 4510b8d..f8add6f 100755 Binary files a/bin/admin-dashboard and b/bin/admin-dashboard differ diff --git a/bin/control-api b/bin/control-api index 81270de..ed06e42 100755 Binary files a/bin/control-api and b/bin/control-api differ diff --git a/bin/prober b/bin/prober index 03b4eba..5293591 100755 Binary files a/bin/prober and b/bin/prober differ diff --git a/bin/validator-agent b/bin/validator-agent index ce183ed..e73751d 100755 Binary files a/bin/validator-agent and b/bin/validator-agent differ diff --git a/internal/agentcore/agentcore.go b/internal/agentcore/agentcore.go index 0c937c6..a90e208 100644 --- a/internal/agentcore/agentcore.go +++ b/internal/agentcore/agentcore.go @@ -30,6 +30,13 @@ type Agent struct { log *slog.Logger lastHandledIPID int64 + + // registerRetryInitial/Max govern the backoff used while waiting for a + // successful registration (see registerWithRetry): control-api may not + // be up yet at agent boot, or may come and go across a redeploy, and the + // agent should keep waiting rather than exit. + registerRetryInitial time.Duration + registerRetryMax time.Duration } func New(cfg *config.ValidatorAgent, log *slog.Logger) *Agent { @@ -38,17 +45,19 @@ func New(cfg *config.ValidatorAgent, log *slog.Logger) *Agent { timeout = 10 * time.Second } return &Agent{ - cfg: cfg, - client: apiclient.New(cfg.ControlAPIURL, timeout+5*time.Second), - log: log, + cfg: cfg, + client: apiclient.New(cfg.ControlAPIURL, timeout+5*time.Second), + log: log, + registerRetryInitial: 3 * time.Second, + registerRetryMax: 30 * time.Second, } } // Run registers with the Control API and polls forever until ctx is // cancelled. func (a *Agent) Run(ctx context.Context) error { - if err := a.register(ctx); err != nil { - return fmt.Errorf("register: %w", err) + if err := a.registerWithRetry(ctx); err != nil { + return err } interval := time.Duration(a.cfg.PollIntervalSeconds) * time.Second @@ -88,6 +97,32 @@ func (a *Agent) register(ctx context.Context) error { return nil } +// registerWithRetry retries register with capped exponential backoff until +// it succeeds or ctx is cancelled. Control-api may not be reachable yet at +// agent boot (started before control-api, or a network blip), or may reject +// the request until an admin adds this validator_id to its config — either +// way the agent should keep waiting rather than exit, since both conditions +// can resolve on their own after the agent has already started. +func (a *Agent) registerWithRetry(ctx context.Context) error { + delay := a.registerRetryInitial + for { + err := a.register(ctx) + if err == nil { + return nil + } + a.log.Warn("registration failed, will retry", "err", err, "retry_in", delay) + select { + case <-ctx.Done(): + return ctx.Err() + case <-time.After(delay): + } + delay *= 2 + if delay > a.registerRetryMax { + delay = a.registerRetryMax + } + } +} + type heartbeatReq struct { LocalState string `json:"local_state"` } diff --git a/internal/agentcore/agentcore_test.go b/internal/agentcore/agentcore_test.go new file mode 100644 index 0000000..fcecb80 --- /dev/null +++ b/internal/agentcore/agentcore_test.go @@ -0,0 +1,100 @@ +package agentcore + +import ( + "context" + "log/slog" + "net/http" + "net/http/httptest" + "os" + "sync/atomic" + "testing" + "time" + + "cloudipvalidator/internal/apiclient" + "cloudipvalidator/internal/config" +) + +func testLogger() *slog.Logger { + return slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{Level: slog.LevelError})) +} + +// TestRegisterWithRetrySucceedsAfterControlAPIComesUp mirrors starting the +// validator-agent before control-api is up (or before an admin has +// configured this validator_id): the first attempts fail, and the agent +// must keep retrying — not give up — until registration succeeds. +func TestRegisterWithRetrySucceedsAfterControlAPIComesUp(t *testing.T) { + var attempts int32 + + mux := http.NewServeMux() + mux.HandleFunc("POST /api/v1/agents/register", func(w http.ResponseWriter, r *http.Request) { + n := atomic.AddInt32(&attempts, 1) + if n < 3 { + w.WriteHeader(http.StatusServiceUnavailable) + return + } + w.WriteHeader(http.StatusOK) + w.Write([]byte(`{"ok":true,"poll_interval_seconds":5}`)) + }) + ts := httptest.NewServer(mux) + defer ts.Close() + + a := &Agent{ + cfg: &config.ValidatorAgent{ValidatorID: "val-1", ControlAPIURL: ts.URL}, + client: apiclient.New(ts.URL, 5*time.Second), + log: testLogger(), + registerRetryInitial: time.Millisecond, + registerRetryMax: 5 * time.Millisecond, + } + + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + + if err := a.registerWithRetry(ctx); err != nil { + t.Fatalf("expected eventual success, got err: %v", err) + } + if got := atomic.LoadInt32(&attempts); got != 3 { + t.Fatalf("expected exactly 3 attempts, got %d", got) + } +} + +// TestRegisterWithRetryStopsOnCancel confirms a validator-agent waiting on +// an unreachable/unconfigured control-api can still be shut down promptly +// (e.g. via SIGTERM) instead of retrying forever with no way out. +func TestRegisterWithRetryStopsOnCancel(t *testing.T) { + mux := http.NewServeMux() + mux.HandleFunc("POST /api/v1/agents/register", func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusBadRequest) + }) + ts := httptest.NewServer(mux) + defer ts.Close() + + a := &Agent{ + cfg: &config.ValidatorAgent{ValidatorID: "val-1", ControlAPIURL: ts.URL}, + client: apiclient.New(ts.URL, 5*time.Second), + log: testLogger(), + registerRetryInitial: 10 * time.Millisecond, + registerRetryMax: 10 * time.Millisecond, + } + + ctx, cancel := context.WithCancel(context.Background()) + go func() { + time.Sleep(30 * time.Millisecond) + cancel() + }() + + var err error + done := make(chan struct{}) + go func() { + defer close(done) + err = a.registerWithRetry(ctx) + }() + + select { + case <-done: + case <-time.After(2 * time.Second): + t.Fatal("registerWithRetry did not return after ctx cancellation") + } + if err != context.Canceled { + t.Fatalf("expected context.Canceled, got %v", err) + } +} diff --git a/internal/dashboard/dto.go b/internal/dashboard/dto.go index 5a17ef2..bc6941d 100644 --- a/internal/dashboard/dto.go +++ b/internal/dashboard/dto.go @@ -97,9 +97,11 @@ type clearQueueResponse struct { } type validatorDTO struct { - ValidatorID string `json:"validator_id"` - OSPortID string `json:"os_port_id"` - State string `json:"state"` + ValidatorID string `json:"validator_id"` + Hostname string `json:"hostname"` + OSPortID string `json:"os_port_id"` + State string `json:"state"` + LastHeartbeatAt *time.Time `json:"last_heartbeat_at"` } type siteDTO struct { diff --git a/internal/dashboard/handlers_test.go b/internal/dashboard/handlers_test.go index 60f292b..50dab2a 100644 --- a/internal/dashboard/handlers_test.go +++ b/internal/dashboard/handlers_test.go @@ -376,6 +376,17 @@ func TestSitesPageShowsProberBadge(t *testing.T) { } } +func TestValidatorsPageShowsHostname(t *testing.T) { + fake, caURL := newFakeControlAPI(t) + fake.validators = append(fake.validators, validatorDTO{ValidatorID: "val-1", Hostname: "validator-host-1", OSPortID: "port-1", State: "idle"}) + ts := newTestServer(t, caURL) + + page := get(t, ts, "/validators") + if !strings.Contains(page, "validator-host-1") { + t.Fatalf("expected hostname shown, got:\n%s", page) + } +} + func TestTargetsAndCheckTypesRoundTrip(t *testing.T) { _, caURL := newFakeControlAPI(t) ts := newTestServer(t, caURL) diff --git a/internal/dashboard/templates/validators.html b/internal/dashboard/templates/validators.html index 8d9c849..e5b43b7 100644 --- a/internal/dashboard/templates/validators.html +++ b/internal/dashboard/templates/validators.html @@ -50,12 +50,13 @@
| ID | OS Port ID | Состояние | |||
|---|---|---|---|---|---|
| ID | Хост | OS Port ID | Состояние | Heartbeat | |
| {{.ValidatorID}} | +{{if .Hostname}}{{.Hostname}}{{else}}—{{end}} | {{$b.Label}} | +{{fmtTime .LastHeartbeatAt}} |
diff --git a/internal/httpapi/dto_admin.go b/internal/httpapi/dto_admin.go
index 90e14b7..60f9ccb 100644
--- a/internal/httpapi/dto_admin.go
+++ b/internal/httpapi/dto_admin.go
@@ -33,9 +33,11 @@ type clearQueueResponse struct {
}
type validatorDTO struct {
- ValidatorID string `json:"validator_id"`
- OSPortID string `json:"os_port_id"`
- State string `json:"state"`
+ ValidatorID string `json:"validator_id"`
+ Hostname string `json:"hostname"`
+ OSPortID string `json:"os_port_id"`
+ State string `json:"state"`
+ LastHeartbeatAt *time.Time `json:"last_heartbeat_at"`
}
type createValidatorRequest struct {
diff --git a/internal/httpapi/handlers_config.go b/internal/httpapi/handlers_config.go
index 7caa451..680e266 100644
--- a/internal/httpapi/handlers_config.go
+++ b/internal/httpapi/handlers_config.go
@@ -17,7 +17,10 @@ func (s *Server) handleConfigListValidators(w http.ResponseWriter, r *http.Reque
}
out := make([]validatorDTO, len(validators))
for i, v := range validators {
- out[i] = validatorDTO{ValidatorID: v.ValidatorID, OSPortID: v.OSPortID, State: v.State}
+ out[i] = validatorDTO{
+ ValidatorID: v.ValidatorID, Hostname: v.Hostname, OSPortID: v.OSPortID,
+ State: v.State, LastHeartbeatAt: v.LastHeartbeatAt,
+ }
}
writeJSON(w, http.StatusOK, out)
}
diff --git a/internal/httpapi/handlers_config_test.go b/internal/httpapi/handlers_config_test.go
index 5fd9f57..3ce0598 100644
--- a/internal/httpapi/handlers_config_test.go
+++ b/internal/httpapi/handlers_config_test.go
@@ -666,6 +666,35 @@ func TestProberRegisterSetsHostnameAndIdleState(t *testing.T) {
}
}
+// TestValidatorRegisterSetsHostname proves the admin config listing surfaces
+// a validator's hostname (and heartbeat time) once its agent registers,
+// mirroring TestProberRegisterSetsHostnameAndIdleState for sites.
+func TestValidatorRegisterSetsHostname(t *testing.T) {
+ fc, _, _, _ := newConfigTestHarness(t)
+
+ fc.do(http.MethodPost, "/api/v1/admin/config/validators", createValidatorRequest{ValidatorID: "val-1", OSPortID: "port-1"})
+
+ resp, body := fc.do(http.MethodPost, "/api/v1/agents/register", registerAgentRequest{ValidatorID: "val-1", Hostname: "validator-host-1", AgentVersion: "test"})
+ if resp.StatusCode != http.StatusOK {
+ t.Fatalf("register agent: status=%d body=%s", resp.StatusCode, body)
+ }
+
+ resp, body = fc.do(http.MethodGet, "/api/v1/admin/config/validators", nil)
+ if resp.StatusCode != http.StatusOK {
+ t.Fatalf("get validators: status=%d body=%s", resp.StatusCode, body)
+ }
+ var validators []validatorDTO
+ if err := json.Unmarshal(body, &validators); err != nil {
+ t.Fatalf("unmarshal validators: %v", err)
+ }
+ if len(validators) != 1 || validators[0].Hostname != "validator-host-1" {
+ t.Fatalf("expected hostname validator-host-1, got %+v", validators)
+ }
+ if validators[0].LastHeartbeatAt != nil {
+ t.Fatalf("expected no heartbeat yet (only register was called), got %+v", validators[0].LastHeartbeatAt)
+ }
+}
+
// TestProberHeartbeat404UnknownSite proves the heartbeat endpoint rejects
// an unconfigured site_id, mirroring the validator heartbeat's 404.
func TestProberHeartbeat404UnknownSite(t *testing.T) {
diff --git a/internal/probercore/probercore.go b/internal/probercore/probercore.go
index cabc614..1cff71b 100644
--- a/internal/probercore/probercore.go
+++ b/internal/probercore/probercore.go
@@ -8,7 +8,6 @@ package probercore
import (
"context"
- "fmt"
"log/slog"
"os"
"time"
@@ -22,6 +21,13 @@ type Prober struct {
cfg *config.Prober
client *apiclient.Client
log *slog.Logger
+
+ // registerRetryInitial/Max govern the backoff used while waiting for a
+ // successful registration (see registerWithRetry): control-api may not
+ // be up yet at prober boot, or may reject an unconfigured site_id until
+ // an admin adds it — the prober should keep waiting rather than exit.
+ registerRetryInitial time.Duration
+ registerRetryMax time.Duration
}
func New(cfg *config.Prober, log *slog.Logger) *Prober {
@@ -30,15 +36,17 @@ func New(cfg *config.Prober, log *slog.Logger) *Prober {
timeout = 10 * time.Second
}
return &Prober{
- cfg: cfg,
- client: apiclient.New(cfg.ControlAPIURL, timeout+5*time.Second),
- log: log,
+ cfg: cfg,
+ client: apiclient.New(cfg.ControlAPIURL, timeout+5*time.Second),
+ log: log,
+ registerRetryInitial: 3 * time.Second,
+ registerRetryMax: 30 * time.Second,
}
}
func (p *Prober) Run(ctx context.Context) error {
- if err := p.register(ctx); err != nil {
- return fmt.Errorf("register: %w", err)
+ if err := p.registerWithRetry(ctx); err != nil {
+ return err
}
interval := time.Duration(p.cfg.PollIntervalSeconds) * time.Second
@@ -70,6 +78,30 @@ func (p *Prober) register(ctx context.Context) error {
return nil
}
+// registerWithRetry retries register with capped exponential backoff until
+// it succeeds or ctx is cancelled — see agentcore.Agent.registerWithRetry
+// for the identical rationale (control-api may start later, or an admin may
+// add this site_id to its config later).
+func (p *Prober) registerWithRetry(ctx context.Context) error {
+ delay := p.registerRetryInitial
+ for {
+ err := p.register(ctx)
+ if err == nil {
+ return nil
+ }
+ p.log.Warn("registration failed, will retry", "err", err, "retry_in", delay)
+ select {
+ case <-ctx.Done():
+ return ctx.Err()
+ case <-time.After(delay):
+ }
+ delay *= 2
+ if delay > p.registerRetryMax {
+ delay = p.registerRetryMax
+ }
+ }
+}
+
type assignment struct {
IPID int64 `json:"ip_id"`
IPAddress string `json:"ip_address"`
diff --git a/internal/probercore/probercore_test.go b/internal/probercore/probercore_test.go
index 1ac6eb8..eb07076 100644
--- a/internal/probercore/probercore_test.go
+++ b/internal/probercore/probercore_test.go
@@ -7,6 +7,7 @@ import (
"net/http/httptest"
"os"
"sync"
+ "sync/atomic"
"testing"
"time"
@@ -95,6 +96,86 @@ func TestPollOnceSkipsAssignmentsWhenHeartbeatFails(t *testing.T) {
}
}
+// TestRegisterWithRetrySucceedsAfterSiteConfigured mirrors an admin adding
+// this prober's site_id to control-api's config after the prober has
+// already started: registration is rejected with 400 until then, and the
+// prober must keep retrying rather than give up.
+func TestRegisterWithRetrySucceedsAfterSiteConfigured(t *testing.T) {
+ var attempts int32
+
+ mux := http.NewServeMux()
+ mux.HandleFunc("POST /api/v1/probers/register", func(w http.ResponseWriter, r *http.Request) {
+ n := atomic.AddInt32(&attempts, 1)
+ if n < 3 {
+ w.WriteHeader(http.StatusBadRequest)
+ w.Write([]byte(`{"error":"unknown site_id: site-1"}`))
+ return
+ }
+ w.WriteHeader(http.StatusOK)
+ })
+ ts := httptest.NewServer(mux)
+ defer ts.Close()
+
+ p := &Prober{
+ cfg: &config.Prober{SiteID: "site-1", ControlAPIURL: ts.URL},
+ client: apiclient.New(ts.URL, 5*time.Second),
+ log: testLogger(),
+ registerRetryInitial: time.Millisecond,
+ registerRetryMax: 5 * time.Millisecond,
+ }
+
+ ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
+ defer cancel()
+
+ if err := p.registerWithRetry(ctx); err != nil {
+ t.Fatalf("expected eventual success, got err: %v", err)
+ }
+ if got := atomic.LoadInt32(&attempts); got != 3 {
+ t.Fatalf("expected exactly 3 attempts, got %d", got)
+ }
+}
+
+// TestRegisterWithRetryStopsOnCancel confirms a prober waiting on an
+// unreachable/unconfigured control-api can still be shut down promptly.
+func TestRegisterWithRetryStopsOnCancel(t *testing.T) {
+ mux := http.NewServeMux()
+ mux.HandleFunc("POST /api/v1/probers/register", func(w http.ResponseWriter, r *http.Request) {
+ w.WriteHeader(http.StatusBadRequest)
+ })
+ ts := httptest.NewServer(mux)
+ defer ts.Close()
+
+ p := &Prober{
+ cfg: &config.Prober{SiteID: "site-1", ControlAPIURL: ts.URL},
+ client: apiclient.New(ts.URL, 5*time.Second),
+ log: testLogger(),
+ registerRetryInitial: 10 * time.Millisecond,
+ registerRetryMax: 10 * time.Millisecond,
+ }
+
+ ctx, cancel := context.WithCancel(context.Background())
+ go func() {
+ time.Sleep(30 * time.Millisecond)
+ cancel()
+ }()
+
+ var err error
+ done := make(chan struct{})
+ go func() {
+ defer close(done)
+ err = p.registerWithRetry(ctx)
+ }()
+
+ select {
+ case <-done:
+ case <-time.After(2 * time.Second):
+ t.Fatal("registerWithRetry did not return after ctx cancellation")
+ }
+ if err != context.Canceled {
+ t.Fatalf("expected context.Canceled, got %v", err)
+ }
+}
+
// TestExtraChecksForPort confirms only ports 443 (TLS) and 22 (SSH) get an
// extra protocol check beyond the baseline TCPConnect every port gets.
func TestExtraChecksForPort(t *testing.T) {
|