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>
362 lines
12 KiB
Go
362 lines
12 KiB
Go
package httpapi
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"net/http"
|
|
"net/url"
|
|
"strconv"
|
|
"strings"
|
|
"time"
|
|
|
|
"cloudipvalidator/internal/db"
|
|
"cloudipvalidator/internal/orchestrator"
|
|
)
|
|
|
|
// destructiveOpTimeout bounds a cancel, delete or clear. Such an operation
|
|
// talks to the cloud as well as the database, and abandoning it halfway
|
|
// leaves floating IPs detached from rows that still exist (or the reverse),
|
|
// so it must not die with the client connection: a client that gives up
|
|
// (a dashboard or curl timeout) only stops waiting for the answer.
|
|
const destructiveOpTimeout = 10 * time.Minute
|
|
|
|
// detachedContext returns a context that ignores cancellation of the request
|
|
// but keeps its values, with destructiveOpTimeout as the upper bound.
|
|
func detachedContext(r *http.Request) (context.Context, context.CancelFunc) {
|
|
return context.WithTimeout(context.WithoutCancel(r.Context()), destructiveOpTimeout)
|
|
}
|
|
|
|
func (s *Server) handleHealthz(w http.ResponseWriter, r *http.Request) {
|
|
writeJSON(w, http.StatusOK, okResponse{OK: true})
|
|
}
|
|
|
|
func (s *Server) handleAdminStatus(w http.ResponseWriter, r *http.Request) {
|
|
// GROUP BY counts instead of loading every row: the dashboard polls this.
|
|
byState, total, err := s.DB.CountIPsByState(r.Context())
|
|
if err != nil {
|
|
writeError(w, http.StatusInternalServerError, err.Error())
|
|
return
|
|
}
|
|
byResult, err := s.DB.CountIPsByResult(r.Context())
|
|
if err != nil {
|
|
writeError(w, http.StatusInternalServerError, err.Error())
|
|
return
|
|
}
|
|
validators, err := s.DB.ListValidators(r.Context())
|
|
if err != nil {
|
|
writeError(w, http.StatusInternalServerError, err.Error())
|
|
return
|
|
}
|
|
overall := map[string]int{
|
|
db.ResultPass: 0, db.ResultPartial: 0, db.ResultFail: 0, db.ResultCancelled: 0,
|
|
}
|
|
for res, n := range byResult {
|
|
overall[res] = n
|
|
}
|
|
writeJSON(w, http.StatusOK, map[string]interface{}{
|
|
"total_ips": total,
|
|
"ips_by_state": byState,
|
|
"total_validators": len(validators),
|
|
"results_by_overall": overall,
|
|
})
|
|
}
|
|
|
|
// handleAdminIPs lists the queue. Without `limit` it returns the bare array of
|
|
// every row (the original contract); with `limit` (1..1000) it returns the
|
|
// envelope {items,total,limit,offset}. Filters: state (csv of valid ip
|
|
// states), q (substring of the address), result (pass|partial|fail|
|
|
// cancelled), order (sequence|aggregated_at_desc); offset >= 0 (needs limit).
|
|
func (s *Server) handleAdminIPs(w http.ResponseWriter, r *http.Request) {
|
|
q := r.URL.Query()
|
|
limit, offset, paged, err := parsePaging(q)
|
|
if err != nil {
|
|
writeError(w, http.StatusBadRequest, err.Error())
|
|
return
|
|
}
|
|
filter := db.IPFilter{Query: strings.TrimSpace(q.Get("q"))}
|
|
for _, st := range strings.Split(q.Get("state"), ",") {
|
|
st = strings.TrimSpace(st)
|
|
if st == "" {
|
|
continue
|
|
}
|
|
if !db.IsValidIPState(st) {
|
|
writeError(w, http.StatusBadRequest, "invalid state "+strconv.Quote(st)+" (valid: "+strings.Join(db.IPStates, ", ")+")")
|
|
return
|
|
}
|
|
filter.States = append(filter.States, st)
|
|
}
|
|
if res := q.Get("result"); res != "" {
|
|
if !db.IsValidResult(res) {
|
|
writeError(w, http.StatusBadRequest, "invalid result "+strconv.Quote(res)+" (valid: pass, partial, fail, cancelled)")
|
|
return
|
|
}
|
|
filter.Result = res
|
|
}
|
|
switch order := q.Get("order"); order {
|
|
case "", db.IPOrderSequence:
|
|
filter.Order = db.IPOrderSequence
|
|
case db.IPOrderAggregatedAtDesc:
|
|
filter.Order = order
|
|
default:
|
|
writeError(w, http.StatusBadRequest, "invalid order "+strconv.Quote(order)+" (valid: sequence, aggregated_at_desc)")
|
|
return
|
|
}
|
|
|
|
if len(q) == 0 {
|
|
ips, err := s.DB.ListIPs(r.Context())
|
|
if err != nil {
|
|
writeError(w, http.StatusInternalServerError, err.Error())
|
|
return
|
|
}
|
|
writeJSON(w, http.StatusOK, ips)
|
|
return
|
|
}
|
|
items, total, err := s.DB.ListIPsPage(r.Context(), filter, limit, offset)
|
|
if err != nil {
|
|
writeError(w, http.StatusInternalServerError, err.Error())
|
|
return
|
|
}
|
|
if !paged {
|
|
writeJSON(w, http.StatusOK, items)
|
|
return
|
|
}
|
|
writeJSON(w, http.StatusOK, ipsPageResponse{Items: items, Total: total, Limit: limit, Offset: offset})
|
|
}
|
|
|
|
// maxPageLimit is the largest allowed `limit` of the paginated list endpoints.
|
|
const maxPageLimit = 1000
|
|
|
|
// parsePaging reads limit/offset. paged is true when `limit` was given (the
|
|
// envelope form); `offset` without `limit` is rejected.
|
|
func parsePaging(q url.Values) (limit, offset int, paged bool, err error) {
|
|
if v := q.Get("limit"); v != "" {
|
|
limit, err = strconv.Atoi(v)
|
|
if err != nil || limit < 1 || limit > maxPageLimit {
|
|
return 0, 0, false, fmt.Errorf("limit must be an integer in 1..%d", maxPageLimit)
|
|
}
|
|
paged = true
|
|
}
|
|
if v := q.Get("offset"); v != "" {
|
|
if !paged {
|
|
return 0, 0, false, errors.New("offset requires limit")
|
|
}
|
|
offset, err = strconv.Atoi(v)
|
|
if err != nil || offset < 0 {
|
|
return 0, 0, false, errors.New("offset must be an integer >= 0")
|
|
}
|
|
}
|
|
return limit, offset, paged, nil
|
|
}
|
|
|
|
// parseBoolParam reads an optional boolean query parameter.
|
|
func parseBoolParam(q url.Values, name string) (bool, error) {
|
|
v := strings.ToLower(q.Get(name))
|
|
switch v {
|
|
case "", "0", "false":
|
|
return false, nil
|
|
case "1", "true":
|
|
return true, nil
|
|
}
|
|
return false, fmt.Errorf("%s must be true or false", name)
|
|
}
|
|
|
|
func (s *Server) handleAdminIPDetail(w http.ResponseWriter, r *http.Request) {
|
|
address := r.PathValue("ip")
|
|
item, err := s.DB.GetIPByAddress(r.Context(), address)
|
|
if err != nil {
|
|
writeError(w, http.StatusNotFound, "unknown ip: "+address)
|
|
return
|
|
}
|
|
checks, err := s.DB.ListChecksForAttempt(r.Context(), item.ID, item.AttemptNumber)
|
|
if err != nil {
|
|
writeError(w, http.StatusInternalServerError, err.Error())
|
|
return
|
|
}
|
|
events, err := s.DB.ListEventsForIP(r.Context(), item.ID)
|
|
if err != nil {
|
|
writeError(w, http.StatusInternalServerError, err.Error())
|
|
return
|
|
}
|
|
writeJSON(w, http.StatusOK, struct {
|
|
IP *db.IPQueueItem `json:"ip"`
|
|
Checks []db.Check `json:"checks"`
|
|
Events []db.Event `json:"events"`
|
|
}{item, checks, events})
|
|
}
|
|
|
|
func (s *Server) handleAdminValidators(w http.ResponseWriter, r *http.Request) {
|
|
validators, err := s.DB.ListValidators(r.Context())
|
|
if err != nil {
|
|
writeError(w, http.StatusInternalServerError, err.Error())
|
|
return
|
|
}
|
|
writeJSON(w, http.StatusOK, validators)
|
|
}
|
|
|
|
// handleAdminSubmitIPs is the single entry point for both adding new
|
|
// addresses to the queue and forcing a re-check of already-finished ones —
|
|
// see db.SubmitIPs for the exact per-address rules. It's a direct DB call
|
|
// (no OpenStack interaction is needed to merely queue work), matching the
|
|
// existing admin handlers above which also bypass the orchestrator for
|
|
// reads.
|
|
func (s *Server) handleAdminSubmitIPs(w http.ResponseWriter, r *http.Request) {
|
|
var req submitIPsRequest
|
|
if err := readJSON(r, &req); err != nil {
|
|
writeError(w, http.StatusBadRequest, "invalid body: "+err.Error())
|
|
return
|
|
}
|
|
result, err := s.DB.SubmitIPs(r.Context(), req.Addresses)
|
|
if err != nil {
|
|
writeDBError(w, err)
|
|
return
|
|
}
|
|
writeJSON(w, http.StatusOK, submitIPsResponse{
|
|
Added: emptyIfNil(result.Added),
|
|
Requeued: emptyIfNil(result.Requeued),
|
|
Reordered: emptyIfNil(result.Reordered),
|
|
SkippedInProgress: emptyIfNil(result.SkippedInProgress),
|
|
})
|
|
}
|
|
|
|
// handleAdminScanFloatingIPs starts the background floating-IP scan (see
|
|
// orchestrator.StartScan) and answers 202 with its status at once; a scan that
|
|
// is already running is joined (202 with the running job's status, no second
|
|
// job). Query: dry_run=true only discovers and counts, leaving the queue
|
|
// untouched; wait=true blocks until the job finishes and answers 200 with the
|
|
// classic synchronous body {scanned_free, added, requeued, reordered,
|
|
// skipped_in_progress} (502 if the scan failed). Takes no body; POST is used
|
|
// (rather than GET) because it mutates the queue, matching handleAdminSubmitIPs.
|
|
func (s *Server) handleAdminScanFloatingIPs(w http.ResponseWriter, r *http.Request) {
|
|
q := r.URL.Query()
|
|
dryRun, err := parseBoolParam(q, "dry_run")
|
|
if err != nil {
|
|
writeError(w, http.StatusBadRequest, err.Error())
|
|
return
|
|
}
|
|
wait, err := parseBoolParam(q, "wait")
|
|
if err != nil {
|
|
writeError(w, http.StatusBadRequest, err.Error())
|
|
return
|
|
}
|
|
opts := orchestrator.ScanOptions{DryRun: dryRun}
|
|
|
|
if !wait {
|
|
st, _ := s.Orch.StartScan(opts)
|
|
writeJSON(w, http.StatusAccepted, toScanStatusDTO(st))
|
|
return
|
|
}
|
|
result, scannedFree, err := s.Orch.ScanAndWait(r.Context(), opts)
|
|
if err != nil {
|
|
writeError(w, http.StatusBadGateway, err.Error())
|
|
return
|
|
}
|
|
writeJSON(w, http.StatusOK, scanIPsResponse{
|
|
ScannedFree: scannedFree,
|
|
Added: emptyIfNil(result.Added),
|
|
Requeued: emptyIfNil(result.Requeued),
|
|
Reordered: emptyIfNil(result.Reordered),
|
|
SkippedInProgress: emptyIfNil(result.SkippedInProgress),
|
|
})
|
|
}
|
|
|
|
// handleAdminScanStatus returns the current status/progress of the scan job.
|
|
func (s *Server) handleAdminScanStatus(w http.ResponseWriter, r *http.Request) {
|
|
writeJSON(w, http.StatusOK, toScanStatusDTO(s.Orch.ScanStatus()))
|
|
}
|
|
|
|
func toScanStatusDTO(st orchestrator.ScanStatus) scanStatusDTO {
|
|
return scanStatusDTO{
|
|
State: string(st.State),
|
|
Running: st.Running,
|
|
DryRun: st.DryRun,
|
|
Pages: st.Pages,
|
|
Discovered: st.Discovered,
|
|
Free: st.Free,
|
|
Added: st.Added,
|
|
Requeued: st.Requeued,
|
|
Reordered: st.Reordered,
|
|
SkippedInProgress: st.SkippedInProgress,
|
|
StartedAt: st.StartedAt,
|
|
FinishedAt: st.FinishedAt,
|
|
Error: st.Error,
|
|
}
|
|
}
|
|
|
|
// handleAdminCancelIP force-stops a check in progress (or still-queued) for
|
|
// the given address. Requires the orchestrator, since a floating IP may
|
|
// need to be disassociated in OpenStack.
|
|
func (s *Server) handleAdminCancelIP(w http.ResponseWriter, r *http.Request) {
|
|
address := r.PathValue("ip")
|
|
ctx, cancel := detachedContext(r)
|
|
defer cancel()
|
|
if err := s.Orch.ForceCancel(ctx, address); err != nil {
|
|
writeDBError(w, err)
|
|
return
|
|
}
|
|
writeJSON(w, http.StatusOK, okResponse{OK: true})
|
|
}
|
|
|
|
// handleAdminDeleteIP permanently removes one address and its full history
|
|
// — unlike handleAdminCancelIP, there's nothing left to look up afterward.
|
|
// Goes through the orchestrator (not a direct DB call) since a currently
|
|
// associated floating IP needs disassociating first.
|
|
func (s *Server) handleAdminDeleteIP(w http.ResponseWriter, r *http.Request) {
|
|
address := r.PathValue("ip")
|
|
ctx, cancel := detachedContext(r)
|
|
defer cancel()
|
|
if err := s.Orch.DeleteIP(ctx, address); err != nil {
|
|
writeDBError(w, err)
|
|
return
|
|
}
|
|
writeJSON(w, http.StatusOK, okResponse{OK: true})
|
|
}
|
|
|
|
// handleAdminDeleteIPs permanently removes a specific list of addresses in
|
|
// one call — unknown addresses are reported in not_found rather than
|
|
// failing the whole request.
|
|
func (s *Server) handleAdminDeleteIPs(w http.ResponseWriter, r *http.Request) {
|
|
var req deleteIPsRequest
|
|
if err := readJSON(r, &req); err != nil {
|
|
writeError(w, http.StatusBadRequest, "invalid body: "+err.Error())
|
|
return
|
|
}
|
|
if len(req.Addresses) == 0 {
|
|
writeError(w, http.StatusBadRequest, "addresses must not be empty")
|
|
return
|
|
}
|
|
ctx, cancel := detachedContext(r)
|
|
defer cancel()
|
|
result, err := s.Orch.DeleteIPs(ctx, req.Addresses)
|
|
if err != nil {
|
|
writeDBError(w, err)
|
|
return
|
|
}
|
|
writeJSON(w, http.StatusOK, deleteIPsResponse{
|
|
Deleted: emptyIfNil(result.Deleted),
|
|
NotFound: emptyIfNil(result.NotFound),
|
|
})
|
|
}
|
|
|
|
// handleAdminClearQueue permanently removes every address currently in the
|
|
// queue, including those actively being checked.
|
|
func (s *Server) handleAdminClearQueue(w http.ResponseWriter, r *http.Request) {
|
|
ctx, cancel := detachedContext(r)
|
|
defer cancel()
|
|
result, err := s.Orch.ClearQueue(ctx)
|
|
if err != nil {
|
|
writeDBError(w, err)
|
|
return
|
|
}
|
|
writeJSON(w, http.StatusOK, clearQueueResponse{Deleted: emptyIfNil(result.Deleted), Count: len(result.Deleted)})
|
|
}
|
|
|
|
// emptyIfNil turns a nil slice into an empty one so these fields always
|
|
// marshal as `[]` rather than `null`.
|
|
func emptyIfNil(s []string) []string {
|
|
if s == nil {
|
|
return []string{}
|
|
}
|
|
return s
|
|
}
|