control-api is hosted outside the cloud and validators reach it directly,
so it sees the floating IP as the connection's source address. New open
route GET /api/v1/agents/{id}/observed-ip returns that address (taken only
from the TCP peer; forwarding headers are ignored so a validator cannot
forge it).
The agent gets self_check.methods, a priority-ordered list of ip_echo
(unchanged) and control_api; the default stays [ip_echo]. The self-check
passes when any method confirms the address; the next method is tried on
no answer and on a mismatch. Each method has its own timeout so a hung
first method cannot starve the fallback, and control_api uses a new TCP
connection per call (a connection opened before the floating IP was
attached would keep reporting the old address).
Also: docker agent template/env, example config, docs, plan in
docs/changes, e2e script switch E2E_SELF_CHECK_METHODS, rebuilt
bin/control-api and bin/validator-agent.
Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
233 lines
8.2 KiB
Go
233 lines
8.2 KiB
Go
package agentcore
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"net"
|
|
"net/http"
|
|
"net/http/httptest"
|
|
"strings"
|
|
"sync/atomic"
|
|
"testing"
|
|
"time"
|
|
|
|
"cloudipvalidator/internal/config"
|
|
)
|
|
|
|
// echoServer answers every request with body (an IP-echo stand-in).
|
|
func echoServer(t *testing.T, body string) *httptest.Server {
|
|
t.Helper()
|
|
ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
|
fmt.Fprint(w, body)
|
|
}))
|
|
t.Cleanup(ts.Close)
|
|
return ts
|
|
}
|
|
|
|
// controlAPIServer plays control-api's observed-ip route: it answers with ip
|
|
// (status 200) or with the given error status when ip is empty.
|
|
func controlAPIServer(t *testing.T, ip string, status int) *httptest.Server {
|
|
t.Helper()
|
|
ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
|
if r.URL.Path != "/api/v1/agents/val-1/observed-ip" {
|
|
http.NotFound(w, r)
|
|
return
|
|
}
|
|
if ip == "" {
|
|
w.WriteHeader(status)
|
|
return
|
|
}
|
|
fmt.Fprintf(w, `{"ip":%q,"source":"remote_addr"}`, ip)
|
|
}))
|
|
t.Cleanup(ts.Close)
|
|
return ts
|
|
}
|
|
|
|
func selfCheckAgent(controlAPIURL string, methods []string, echoURLs ...string) *Agent {
|
|
return &Agent{
|
|
cfg: &config.ValidatorAgent{
|
|
ValidatorID: "val-1",
|
|
ControlAPIURL: controlAPIURL,
|
|
SelfCheck: config.SelfCheckCfg{TimeoutSeconds: 2, Methods: methods, IPEchoURLs: echoURLs},
|
|
},
|
|
log: testLogger(),
|
|
}
|
|
}
|
|
|
|
func TestSelfCheckControlAPIFirstWins(t *testing.T) {
|
|
capi := controlAPIServer(t, "1.2.3.4", 0)
|
|
var echoCalls int32
|
|
echo := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
|
atomic.AddInt32(&echoCalls, 1)
|
|
fmt.Fprint(w, "1.2.3.4")
|
|
}))
|
|
defer echo.Close()
|
|
|
|
a := selfCheckAgent(capi.URL, []string{config.SelfCheckControlAPI, config.SelfCheckIPEcho}, echo.URL)
|
|
ip, method, detail, ok := a.runSelfCheckMethods(context.Background(), "1.2.3.4")
|
|
if !ok || ip != "1.2.3.4" || method != config.SelfCheckControlAPI || detail != "matched (control_api)" {
|
|
t.Fatalf("got ip=%q method=%q detail=%q ok=%v", ip, method, detail, ok)
|
|
}
|
|
if n := atomic.LoadInt32(&echoCalls); n != 0 {
|
|
t.Fatalf("ip_echo was called %d times although control_api already confirmed", n)
|
|
}
|
|
}
|
|
|
|
// Any one confirming method is enough: control-api answers with a different
|
|
// address, ip_echo confirms.
|
|
func TestSelfCheckFallsThroughOnMismatch(t *testing.T) {
|
|
capi := controlAPIServer(t, "5.5.5.5", 0)
|
|
echo := echoServer(t, "1.2.3.4")
|
|
|
|
a := selfCheckAgent(capi.URL, []string{config.SelfCheckControlAPI, config.SelfCheckIPEcho}, echo.URL)
|
|
ip, method, _, ok := a.runSelfCheckMethods(context.Background(), "1.2.3.4")
|
|
if !ok || ip != "1.2.3.4" || method != config.SelfCheckIPEcho {
|
|
t.Fatalf("got ip=%q method=%q ok=%v, want a pass via ip_echo", ip, method, ok)
|
|
}
|
|
}
|
|
|
|
// An old control-api without the route (404) or a failing one (5xx) must not
|
|
// stop the self-check: the next method decides.
|
|
func TestSelfCheckFallsThroughOnControlAPIError(t *testing.T) {
|
|
for _, status := range []int{http.StatusNotFound, http.StatusInternalServerError} {
|
|
capi := controlAPIServer(t, "", status)
|
|
echo := echoServer(t, "1.2.3.4")
|
|
a := selfCheckAgent(capi.URL, []string{config.SelfCheckControlAPI, config.SelfCheckIPEcho}, echo.URL)
|
|
_, method, _, ok := a.runSelfCheckMethods(context.Background(), "1.2.3.4")
|
|
if !ok || method != config.SelfCheckIPEcho {
|
|
t.Fatalf("status %d: method=%q ok=%v, want a pass via ip_echo", status, method, ok)
|
|
}
|
|
}
|
|
}
|
|
|
|
func TestSelfCheckFailsWhenNoMethodConfirms(t *testing.T) {
|
|
capi := controlAPIServer(t, "10.0.0.5", 0) // private address: internal-network case
|
|
echo := echoServer(t, "6.6.6.6")
|
|
|
|
a := selfCheckAgent(capi.URL, []string{config.SelfCheckControlAPI, config.SelfCheckIPEcho}, echo.URL)
|
|
ip, method, detail, ok := a.runSelfCheckMethods(context.Background(), "1.2.3.4")
|
|
if ok || method != "" {
|
|
t.Fatalf("expected a failure, got ok=%v method=%q", ok, method)
|
|
}
|
|
if ip != "6.6.6.6" {
|
|
t.Fatalf("detected ip = %q, want the last reported address 6.6.6.6", ip)
|
|
}
|
|
for _, want := range []string{
|
|
`control_api: egress ip "10.0.0.5" does not match assigned fip "1.2.3.4"`,
|
|
"private address", // the hint for the internal-network case
|
|
`ip_echo: egress ip "6.6.6.6" does not match`,
|
|
} {
|
|
if !strings.Contains(detail, want) {
|
|
t.Fatalf("detail %q does not contain %q", detail, want)
|
|
}
|
|
}
|
|
}
|
|
|
|
func TestSelfCheckPrivateHintOnlyForControlAPI(t *testing.T) {
|
|
echo := echoServer(t, "10.1.1.1")
|
|
a := selfCheckAgent("http://unused", []string{config.SelfCheckIPEcho}, echo.URL)
|
|
_, _, detail, ok := a.runSelfCheckMethods(context.Background(), "1.2.3.4")
|
|
if ok || strings.Contains(detail, "private address") {
|
|
t.Fatalf("ok=%v detail=%q: the control_api hint must not appear for ip_echo", ok, detail)
|
|
}
|
|
}
|
|
|
|
// With nothing configured (a config built without the loader) the agent
|
|
// behaves as before: ip_echo only, control-api is never asked.
|
|
func TestSelfCheckDefaultsToIPEcho(t *testing.T) {
|
|
var capiCalls int32
|
|
capi := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
|
atomic.AddInt32(&capiCalls, 1)
|
|
}))
|
|
defer capi.Close()
|
|
echo := echoServer(t, "1.2.3.4")
|
|
|
|
a := selfCheckAgent(capi.URL, nil, echo.URL)
|
|
_, method, _, ok := a.runSelfCheckMethods(context.Background(), "1.2.3.4")
|
|
if !ok || method != config.SelfCheckIPEcho || atomic.LoadInt32(&capiCalls) != 0 {
|
|
t.Fatalf("method=%q ok=%v control-api calls=%d", method, ok, atomic.LoadInt32(&capiCalls))
|
|
}
|
|
}
|
|
|
|
// A hung control-api must not use up the time of the fallback: each method
|
|
// has its own timeout.
|
|
func TestSelfCheckHungControlAPIDoesNotStarveFallback(t *testing.T) {
|
|
release := make(chan struct{})
|
|
capi := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
|
<-release
|
|
}))
|
|
defer capi.Close()
|
|
defer close(release)
|
|
echo := echoServer(t, "1.2.3.4")
|
|
|
|
a := selfCheckAgent(capi.URL, []string{config.SelfCheckControlAPI, config.SelfCheckIPEcho}, echo.URL)
|
|
a.cfg.SelfCheck.TimeoutSeconds = 1
|
|
|
|
start := time.Now()
|
|
// The outer context mirrors handleSelfCheckAndRun: timeout x methods.
|
|
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
|
|
defer cancel()
|
|
_, method, detail, ok := a.runSelfCheckMethods(ctx, "1.2.3.4")
|
|
if !ok || method != config.SelfCheckIPEcho {
|
|
t.Fatalf("method=%q ok=%v detail=%q, want a pass via ip_echo after control_api timed out", method, ok, detail)
|
|
}
|
|
if elapsed := time.Since(start); elapsed > 1900*time.Millisecond {
|
|
t.Fatalf("took %s: the hung method consumed the fallback's time", elapsed)
|
|
}
|
|
}
|
|
|
|
// Every control-api request must use a new TCP connection: a connection
|
|
// opened before the floating IP was attached would report the old address.
|
|
func TestDetectViaControlAPIDialsNewConnectionEachTime(t *testing.T) {
|
|
var conns int32
|
|
ts := httptest.NewUnstartedServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
|
fmt.Fprint(w, `{"ip":"1.2.3.4","source":"remote_addr"}`)
|
|
}))
|
|
ts.Config.ConnState = func(_ net.Conn, s http.ConnState) {
|
|
if s == http.StateNew {
|
|
atomic.AddInt32(&conns, 1)
|
|
}
|
|
}
|
|
ts.Start()
|
|
defer ts.Close()
|
|
|
|
a := selfCheckAgent(ts.URL, nil)
|
|
for i := 0; i < 3; i++ {
|
|
if _, err := a.detectViaControlAPI(context.Background()); err != nil {
|
|
t.Fatalf("call %d: %v", i, err)
|
|
}
|
|
}
|
|
if n := atomic.LoadInt32(&conns); n != 3 {
|
|
t.Fatalf("3 calls opened %d connections, want 3 (no keep-alive reuse)", n)
|
|
}
|
|
}
|
|
|
|
func TestDetectViaControlAPIRejectsBadAnswers(t *testing.T) {
|
|
for name, body := range map[string]string{"not json": "oops", "not an ip": `{"ip":"abc"}`, "empty": `{}`} {
|
|
ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { fmt.Fprint(w, body) }))
|
|
a := selfCheckAgent(ts.URL, nil)
|
|
if ip, err := a.detectViaControlAPI(context.Background()); err == nil {
|
|
t.Fatalf("%s: expected an error, got %q", name, ip)
|
|
}
|
|
ts.Close()
|
|
}
|
|
}
|
|
|
|
// The agent token must never be sent on this request (the route is open and
|
|
// the token is meant for control-api writes only).
|
|
func TestDetectViaControlAPISendsNoToken(t *testing.T) {
|
|
var auth string
|
|
ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
|
auth = r.Header.Get("Authorization")
|
|
fmt.Fprint(w, `{"ip":"1.2.3.4"}`)
|
|
}))
|
|
defer ts.Close()
|
|
a := selfCheckAgent(ts.URL, nil)
|
|
if _, err := a.detectViaControlAPI(context.Background()); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if auth != "" {
|
|
t.Fatalf("Authorization header sent: %q", auth)
|
|
}
|
|
}
|