A mass check on 2026-10-02 stalled 7 of 20 validators and sent 42 addresses to fail without a single check. A validator busy with slow checks went silent, was marked unreachable, and its next heartbeat put it back to idle while it still held the address; it was handed a second one, whose association never ran (the in-flight guard was keyed by validator), and both waited for their leases to expire. - Heartbeat/re-register return an unreachable validator to assigned when it still holds an address, else idle. - A validator is released only from the address it currently holds (ReleaseFIP, RequeueOrFail, MarkFIPOccupied, FreeValidator); an unreachable validator stays unreachable until its next heartbeat, so a dead validator is no longer handed a new address every lease period. - ClaimNextQueued refuses a validator that still has an address; a ReconcileValidators pass on every tick repairs rows that disagree with the queue. - Association guard is keyed by address, not validator. - The agent sends heartbeats from their own goroutine. - Clear queue / delete: detach only floating IPs of unfinished rows (done, failed and occupied rows kept their fip_id and made a clear issue >1000 sequential cloud calls: 256 s), at most 8 in parallel; the operation no longer dies with the client connection (10 minute limit). Includes the incident analysis and the plan under analysis/ and docs/changes/, and rebuilt bin/control-api and bin/validator-agent. Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
73 lines
2.6 KiB
Go
73 lines
2.6 KiB
Go
package agentcore
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"net/http"
|
|
"net/http/httptest"
|
|
"sync/atomic"
|
|
"testing"
|
|
"time"
|
|
|
|
"cloudipvalidator/internal/config"
|
|
)
|
|
|
|
// While the agent is busy with a slow assignment (an address whose outbound
|
|
// targets time out keeps it occupied for tens of seconds) it must keep sending
|
|
// heartbeats; control-api marks a validator that stays silent for
|
|
// heartbeat_timeout_seconds as unreachable.
|
|
func TestHeartbeatContinuesDuringSlowChecks(t *testing.T) {
|
|
slowTarget := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
|
time.Sleep(3500 * time.Millisecond)
|
|
}))
|
|
defer slowTarget.Close()
|
|
|
|
var heartbeats, assignments int32
|
|
mux := http.NewServeMux()
|
|
mux.HandleFunc("POST /api/v1/agents/register", func(w http.ResponseWriter, r *http.Request) {
|
|
fmt.Fprint(w, `{"ok":true}`)
|
|
})
|
|
mux.HandleFunc("POST /api/v1/agents/val-1/heartbeat", func(w http.ResponseWriter, r *http.Request) {
|
|
atomic.AddInt32(&heartbeats, 1)
|
|
fmt.Fprint(w, `{"ok":true}`)
|
|
})
|
|
mux.HandleFunc("GET /api/v1/agents/val-1/assignment", func(w http.ResponseWriter, r *http.Request) {
|
|
if atomic.AddInt32(&assignments, 1) > 1 {
|
|
w.WriteHeader(http.StatusNoContent)
|
|
return
|
|
}
|
|
fmt.Fprintf(w, `{"ip_id":1,"ip_address":"1.1.1.1","phase":"checking","check_config":[{"type":"https","targets":[%q]}]}`, slowTarget.URL)
|
|
})
|
|
mux.HandleFunc("POST /api/v1/agents/val-1/results", func(w http.ResponseWriter, r *http.Request) { fmt.Fprint(w, `{"ok":true}`) })
|
|
mux.HandleFunc("POST /api/v1/agents/val-1/complete", func(w http.ResponseWriter, r *http.Request) { fmt.Fprint(w, `{"ok":true}`) })
|
|
capi := httptest.NewServer(mux)
|
|
defer capi.Close()
|
|
|
|
a := New(&config.ValidatorAgent{
|
|
ValidatorID: "val-1", ControlAPIURL: capi.URL, PollIntervalSeconds: 1,
|
|
Checks: config.AgentChecks{HTTPSTimeoutSeconds: 10, ICMPTimeoutSeconds: 1, ICMPCount: 1},
|
|
}, testLogger())
|
|
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
done := make(chan struct{})
|
|
go func() { _ = a.Run(ctx); close(done) }()
|
|
|
|
time.Sleep(3 * time.Second) // the slow check (3.5 s) is still running
|
|
during := atomic.LoadInt32(&heartbeats)
|
|
cancel()
|
|
select {
|
|
case <-done:
|
|
case <-time.After(10 * time.Second):
|
|
t.Fatal("Run did not stop after the context was cancelled")
|
|
}
|
|
|
|
// One per second plus the first: 3-4 in 3 s. With heartbeats in the poll
|
|
// loop there is exactly one, sent before the slow assignment started.
|
|
if during < 3 {
|
|
t.Fatalf("%d heartbeats in 3 s while a check was running, want at least 3", during)
|
|
}
|
|
if atomic.LoadInt32(&assignments) < 1 {
|
|
t.Fatal("the assignment was never fetched, the test did not exercise a busy agent")
|
|
}
|
|
}
|