Files

361 lines
12 KiB
Go
Raw Permalink Normal View History

2026-08-21 07:34:45 +03:00
package httpapi
import (
"context"
"errors"
"fmt"
2026-08-21 07:34:45 +03:00
"net/http"
"net/url"
"strconv"
"strings"
"time"
2026-08-21 07:34:45 +03:00
"cloudipvalidator/internal/db"
"cloudipvalidator/internal/orchestrator"
2026-08-21 07:34:45 +03:00
)
// 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)
}
2026-08-21 07:34:45 +03:00
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())
2026-08-21 07:34:45 +03:00
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
2026-08-21 07:34:45 +03:00
}
writeJSON(w, http.StatusOK, map[string]interface{}{
"total_ips": total,
"ips_by_state": byState,
"total_validators": len(validators),
"results_by_overall": overall,
2026-08-21 07:34:45 +03:00
})
}
// 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).
2026-08-21 07:34:45 +03:00
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)
2026-08-21 07:34:45 +03:00
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)
2026-08-21 07:34:45 +03:00
}
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)
}
2026-08-23 20:39:22 +03:00
// 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,
}
}
2026-08-23 20:39:22 +03:00
// 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 {
2026-08-23 20:39:22 +03:00
writeDBError(w, err)
return
}
writeJSON(w, http.StatusOK, okResponse{OK: true})
}
2026-08-23 22:24:55 +03:00
// 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 {
2026-08-23 22:24:55 +03:00
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)
2026-08-23 22:24:55 +03:00
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)
2026-08-23 22:24:55 +03:00
if err != nil {
writeDBError(w, err)
return
}
writeJSON(w, http.StatusOK, clearQueueResponse{Deleted: emptyIfNil(result.Deleted), Count: len(result.Deleted)})
2026-08-23 22:24:55 +03:00
}
2026-08-23 20:39:22 +03:00
// 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
}