2026-08-21 07:34:45 +03:00
|
|
|
package httpapi
|
|
|
|
|
|
|
|
|
|
import (
|
|
|
|
|
"net/http"
|
|
|
|
|
"time"
|
|
|
|
|
|
|
|
|
|
"cloudipvalidator/internal/db"
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
func (s *Server) handleProberRegister(w http.ResponseWriter, r *http.Request) {
|
|
|
|
|
var req registerProberRequest
|
|
|
|
|
if err := readJSON(r, &req); err != nil {
|
|
|
|
|
writeError(w, http.StatusBadRequest, "invalid body: "+err.Error())
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-08-23 20:39:22 +03:00
|
|
|
idx, err := s.Orch.SiteIndexForID(r.Context(), req.SiteID)
|
|
|
|
|
if err != nil {
|
|
|
|
|
writeError(w, http.StatusInternalServerError, err.Error())
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
if idx == 0 {
|
2026-08-21 07:34:45 +03:00
|
|
|
writeError(w, http.StatusBadRequest, "unknown site_id: "+req.SiteID)
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-08-26 20:47:54 +03:00
|
|
|
if err := s.DB.RegisterSiteProber(r.Context(), req.SiteID, req.Hostname); err != nil {
|
|
|
|
|
writeError(w, http.StatusInternalServerError, err.Error())
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-08-21 07:34:45 +03:00
|
|
|
s.Orch.RecordEvent(r.Context(), "prober", req.SiteID, nil, "registered", "")
|
|
|
|
|
writeJSON(w, http.StatusOK, registerAgentResponse{OK: true, PollIntervalSeconds: s.Orch.Cfg.PollIntervalSeconds})
|
|
|
|
|
}
|
|
|
|
|
|
2026-08-26 20:47:54 +03:00
|
|
|
func (s *Server) handleProberHeartbeat(w http.ResponseWriter, r *http.Request) {
|
|
|
|
|
siteID := r.PathValue("site_id")
|
|
|
|
|
var req heartbeatRequest
|
|
|
|
|
_ = readJSON(r, &req) // heartbeat body is informational only; tolerate empty/missing
|
|
|
|
|
if err := s.DB.SiteHeartbeat(r.Context(), siteID); err != nil {
|
|
|
|
|
writeError(w, http.StatusNotFound, "unknown site_id: "+siteID)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
writeJSON(w, http.StatusOK, okResponse{OK: true})
|
|
|
|
|
}
|
|
|
|
|
|
2026-08-21 07:34:45 +03:00
|
|
|
// handleProberAssignments returns every IP currently in the checking
|
|
|
|
|
// state — probers work the whole active set each poll, not one IP at a
|
|
|
|
|
// time, since multiple validators run in parallel.
|
|
|
|
|
func (s *Server) handleProberAssignments(w http.ResponseWriter, r *http.Request) {
|
|
|
|
|
siteID := r.PathValue("site_id")
|
2026-08-23 20:39:22 +03:00
|
|
|
idx, err := s.Orch.SiteIndexForID(r.Context(), siteID)
|
|
|
|
|
if err != nil {
|
|
|
|
|
writeError(w, http.StatusInternalServerError, err.Error())
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
if idx == 0 {
|
2026-08-21 07:34:45 +03:00
|
|
|
writeError(w, http.StatusNotFound, "unknown site_id: "+siteID)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
items, err := s.DB.ListChecking(r.Context())
|
|
|
|
|
if err != nil {
|
|
|
|
|
writeError(w, http.StatusInternalServerError, err.Error())
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-08-26 19:45:48 +03:00
|
|
|
inbound, err := s.DB.GetInboundChecks(r.Context())
|
|
|
|
|
if err != nil {
|
|
|
|
|
writeError(w, http.StatusInternalServerError, err.Error())
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-08-21 07:34:45 +03:00
|
|
|
out := make([]proberAssignment, 0, len(items))
|
|
|
|
|
for _, item := range items {
|
|
|
|
|
out = append(out, proberAssignment{
|
|
|
|
|
IPID: item.ID, IPAddress: item.IPAddress,
|
2026-08-26 19:45:48 +03:00
|
|
|
Ports: inbound.Ports, ICMP: inbound.ICMP,
|
2026-08-21 07:34:45 +03:00
|
|
|
})
|
|
|
|
|
}
|
|
|
|
|
writeJSON(w, http.StatusOK, out)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (s *Server) handleProberResults(w http.ResponseWriter, r *http.Request) {
|
|
|
|
|
siteID := r.PathValue("site_id")
|
2026-08-23 20:39:22 +03:00
|
|
|
siteIndex, err := s.Orch.SiteIndexForID(r.Context(), siteID)
|
|
|
|
|
if err != nil {
|
|
|
|
|
writeError(w, http.StatusInternalServerError, err.Error())
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-08-21 07:34:45 +03:00
|
|
|
if siteIndex == 0 {
|
|
|
|
|
writeError(w, http.StatusNotFound, "unknown site_id: "+siteID)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
var req proberResultsRequest
|
|
|
|
|
if err := readJSON(r, &req); err != nil {
|
|
|
|
|
writeError(w, http.StatusBadRequest, "invalid body: "+err.Error())
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
completed := map[int64]bool{}
|
|
|
|
|
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,
|
|
|
|
|
Source: db.InboundSource(siteIndex), CheckType: res.CheckType, Target: res.IPAddress,
|
|
|
|
|
Success: res.Success, LatencyMS: res.LatencyMS, Detail: res.Detail, CheckedAt: checkedAt,
|
|
|
|
|
})
|
|
|
|
|
if err != nil {
|
|
|
|
|
writeError(w, http.StatusInternalServerError, err.Error())
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
if res.Complete {
|
|
|
|
|
completed[res.IPID] = true
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
for ipID := range completed {
|
|
|
|
|
if err := s.Orch.MarkSiteComplete(r.Context(), ipID, siteIndex); err != nil {
|
|
|
|
|
writeError(w, http.StatusInternalServerError, err.Error())
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
writeJSON(w, http.StatusOK, okResponse{OK: true})
|
|
|
|
|
}
|