Files
ayurishchevandClaude Sonnet 5.5 abbee9a08a Add self-check via control-api (self_check.methods)
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>
2026-10-02 03:24:20 +03:00

187 lines
6.2 KiB
Go

package httpapi
import (
"database/sql"
"errors"
"fmt"
"net"
"net/http"
"net/netip"
"time"
"cloudipvalidator/internal/db"
)
func (s *Server) handleAgentRegister(w http.ResponseWriter, r *http.Request) {
var req registerAgentRequest
if err := readJSON(r, &req); err != nil {
writeError(w, http.StatusBadRequest, "invalid body: "+err.Error())
return
}
if req.ValidatorID == "" {
writeError(w, http.StatusBadRequest, "validator_id is required")
return
}
// os_port_id is supplied via control-api's own config (config.ValidatorConfig),
// not by the agent, so registration only touches hostname/version here;
// RegisterValidator preserves any existing os_port_id row.
existing, _ := s.DB.GetValidator(r.Context(), req.ValidatorID)
osPortID := ""
if existing != nil {
osPortID = existing.OSPortID
}
if err := s.DB.RegisterValidator(r.Context(), req.ValidatorID, req.Hostname, osPortID, req.AgentVersion); err != nil {
writeError(w, http.StatusInternalServerError, err.Error())
return
}
s.Orch.RecordEvent(r.Context(), "validator-agent", req.ValidatorID, nil, "registered", "")
writeJSON(w, http.StatusOK, registerAgentResponse{OK: true, PollIntervalSeconds: s.Orch.Cfg.PollIntervalSeconds})
}
func (s *Server) handleAgentHeartbeat(w http.ResponseWriter, r *http.Request) {
id := r.PathValue("id")
var req heartbeatRequest
_ = readJSON(r, &req) // heartbeat body is informational only; tolerate empty/missing
if err := s.DB.Heartbeat(r.Context(), id); err != nil {
writeError(w, http.StatusNotFound, "unknown validator: "+id)
return
}
writeJSON(w, http.StatusOK, okResponse{OK: true})
}
func (s *Server) handleAgentAssignment(w http.ResponseWriter, r *http.Request) {
id := r.PathValue("id")
item, checks, err := s.Orch.AssignmentForValidator(r.Context(), id)
if err != nil {
writeError(w, http.StatusNotFound, "unknown validator: "+id)
return
}
if item == nil {
w.WriteHeader(http.StatusNoContent)
return
}
var cfg []checkConfigDTO
for _, c := range checks {
cfg = append(cfg, checkConfigDTO{Type: c.Type, Targets: c.Targets})
}
writeJSON(w, http.StatusOK, assignmentResponse{
IPID: item.ID, IPAddress: item.IPAddress, Phase: item.State, CheckConfig: cfg,
})
}
// handleAgentObservedIP tells a validator which source address control-api
// sees for its connection, so the agent's self-check can confirm the floating
// IP without a third-party IP-echo service. The route is open, so it is
// limited to known validators to keep it from becoming a public "what is my
// IP" service.
func (s *Server) handleAgentObservedIP(w http.ResponseWriter, r *http.Request) {
id := r.PathValue("id")
if _, err := s.DB.GetValidator(r.Context(), id); err != nil {
if errors.Is(err, sql.ErrNoRows) || errors.Is(err, db.ErrNotFound) {
writeError(w, http.StatusNotFound, "unknown validator: "+id)
return
}
writeError(w, http.StatusInternalServerError, err.Error())
return
}
ip, err := clientIP(r)
if err != nil {
writeError(w, http.StatusInternalServerError, err.Error())
return
}
writeJSON(w, http.StatusOK, observedIPResponse{IP: ip, Source: "remote_addr"})
}
// clientIP returns the peer address of the TCP connection in canonical form
// (IPv4-mapped IPv6 unmapped to plain IPv4). It deliberately ignores
// X-Forwarded-For / X-Real-IP: validators connect directly, and trusting a
// client-supplied header would let a validator forge the address and pass
// the self-check.
func clientIP(r *http.Request) (string, error) {
host, _, err := net.SplitHostPort(r.RemoteAddr)
if err != nil {
return "", fmt.Errorf("parse remote address %q: %w", r.RemoteAddr, err)
}
addr, err := netip.ParseAddr(host)
if err != nil {
return "", fmt.Errorf("parse remote address %q: %w", r.RemoteAddr, err)
}
return addr.Unmap().String(), nil
}
func (s *Server) handleAgentSelfCheck(w http.ResponseWriter, r *http.Request) {
id := r.PathValue("id")
var req selfCheckRequest
if err := readJSON(r, &req); err != nil {
writeError(w, http.StatusBadRequest, "invalid body: "+err.Error())
return
}
detail := req.Detail
if req.DetectedEgress != "" {
detail = "detected_egress_ip=" + req.DetectedEgress + " " + detail
}
if err := s.Orch.SelfCheckResult(r.Context(), id, req.IPID, req.Success, detail); err != nil {
writeError(w, http.StatusInternalServerError, err.Error())
return
}
writeJSON(w, http.StatusOK, okResponse{OK: true})
}
func (s *Server) handleAgentEvent(w http.ResponseWriter, r *http.Request) {
id := r.PathValue("id")
var req agentEventRequest
if err := readJSON(r, &req); err != nil {
writeError(w, http.StatusBadRequest, "invalid body: "+err.Error())
return
}
if req.EventType == "" {
writeError(w, http.StatusBadRequest, "event_type is required")
return
}
s.Orch.RecordEvent(r.Context(), "validator-agent", id, req.IPID, req.EventType, req.Payload)
writeJSON(w, http.StatusOK, okResponse{OK: true})
}
func (s *Server) handleAgentResults(w http.ResponseWriter, r *http.Request) {
id := r.PathValue("id")
var req agentResultsRequest
if err := readJSON(r, &req); err != nil {
writeError(w, http.StatusBadRequest, "invalid body: "+err.Error())
return
}
for _, res := range req.Results {
item, err := s.DB.GetIP(r.Context(), res.IPID)
if err != nil {
writeError(w, http.StatusNotFound, "unknown ip_id")
return
}
checkedAt, err := time.Parse(time.RFC3339Nano, res.CheckedAt)
if err != nil {
checkedAt = db.Now()
}
err = s.Orch.RecordCheck(r.Context(), db.Check{
IPID: item.ID, IPAddress: item.IPAddress, AttemptNumber: item.AttemptNumber,
ValidatorID: id, Source: db.SourceEgress, CheckType: res.CheckType, Target: res.Target,
Success: res.Success, LatencyMS: res.LatencyMS, Detail: res.Detail, CheckedAt: checkedAt,
})
if err != nil {
writeError(w, http.StatusInternalServerError, err.Error())
return
}
}
writeJSON(w, http.StatusOK, okResponse{OK: true})
}
func (s *Server) handleAgentComplete(w http.ResponseWriter, r *http.Request) {
var req agentCompleteRequest
if err := readJSON(r, &req); err != nil {
writeError(w, http.StatusBadRequest, "invalid body: "+err.Error())
return
}
if err := s.Orch.MarkEgressComplete(r.Context(), req.IPID); err != nil {
writeError(w, http.StatusInternalServerError, err.Error())
return
}
writeJSON(w, http.StatusOK, okResponse{OK: true})
}