2026-09-18 10:43:00 +03:00
|
|
|
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)
|
|
|
|
|
}
|
|
|
|
|
}
|
2026-10-01 11:35:24 +03:00
|
|
|
|
|
|
|
|
// TestTokenGoesOnlyToControlAPI: the bearer token authenticates calls to
|
|
|
|
|
// control-api, and must never leak to the external IP-echo service.
|
|
|
|
|
func TestTokenGoesOnlyToControlAPI(t *testing.T) {
|
|
|
|
|
var apiAuth, echoAuth atomic.Value
|
|
|
|
|
api := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
|
|
|
|
apiAuth.Store(r.Header.Get("Authorization"))
|
|
|
|
|
w.Write([]byte(`{"ok":true,"poll_interval_seconds":5}`))
|
|
|
|
|
}))
|
|
|
|
|
defer api.Close()
|
|
|
|
|
echo := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
|
|
|
|
echoAuth.Store(r.Header.Get("Authorization"))
|
|
|
|
|
w.Write([]byte("203.0.113.7"))
|
|
|
|
|
}))
|
|
|
|
|
defer echo.Close()
|
|
|
|
|
|
|
|
|
|
a := New(&config.ValidatorAgent{ValidatorID: "val-1", ControlAPIURL: api.URL}, testLogger()).WithToken("agent-secret")
|
|
|
|
|
ctx := context.Background()
|
|
|
|
|
if err := a.registerWithRetry(ctx); err != nil {
|
|
|
|
|
t.Fatalf("register: %v", err)
|
|
|
|
|
}
|
|
|
|
|
if _, err := fetchIPEcho(ctx, echo.URL); err != nil {
|
|
|
|
|
t.Fatalf("fetchIPEcho: %v", err)
|
|
|
|
|
}
|
|
|
|
|
if got, _ := apiAuth.Load().(string); got != "Bearer agent-secret" {
|
|
|
|
|
t.Fatalf("control-api Authorization = %q, want bearer token", got)
|
|
|
|
|
}
|
|
|
|
|
if got, _ := echoAuth.Load().(string); got != "" {
|
|
|
|
|
t.Fatalf("IP-echo request carried Authorization %q, want none", got)
|
|
|
|
|
}
|
|
|
|
|
}
|