2026-08-21 07:34:45 +03:00
|
|
|
// Package agentcore implements the validator-agent's poll loop: register,
|
|
|
|
|
// heartbeat, wait for an assignment, self-check that egress actually flows
|
|
|
|
|
// through the newly attached FIP, run the configured outbound/egress
|
|
|
|
|
// checks, and report results — all driven entirely by the Control API, so
|
|
|
|
|
// the process itself holds no durable state (constraint: the agent must be
|
|
|
|
|
// safely restartable at any point without losing correctness, only
|
|
|
|
|
// possibly re-doing in-flight work, which the Control API's idempotent
|
|
|
|
|
// upserts tolerate).
|
|
|
|
|
package agentcore
|
|
|
|
|
|
|
|
|
|
import (
|
|
|
|
|
"context"
|
2026-10-02 03:24:20 +03:00
|
|
|
"encoding/json"
|
2026-08-21 07:34:45 +03:00
|
|
|
"fmt"
|
2026-08-21 11:04:49 +03:00
|
|
|
"io"
|
2026-08-21 07:34:45 +03:00
|
|
|
"log/slog"
|
2026-08-21 11:04:49 +03:00
|
|
|
"net"
|
|
|
|
|
"net/http"
|
2026-10-02 03:24:20 +03:00
|
|
|
neturl "net/url"
|
2026-08-21 07:34:45 +03:00
|
|
|
"os"
|
2026-08-21 11:04:49 +03:00
|
|
|
"strings"
|
2026-10-02 14:42:12 +03:00
|
|
|
"sync/atomic"
|
2026-08-21 07:34:45 +03:00
|
|
|
"time"
|
|
|
|
|
|
|
|
|
|
"cloudipvalidator/internal/apiclient"
|
|
|
|
|
"cloudipvalidator/internal/checkrunner"
|
|
|
|
|
"cloudipvalidator/internal/config"
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
type Agent struct {
|
|
|
|
|
cfg *config.ValidatorAgent
|
|
|
|
|
client *apiclient.Client
|
|
|
|
|
log *slog.Logger
|
|
|
|
|
|
|
|
|
|
lastHandledIPID int64
|
2026-09-18 10:43:00 +03:00
|
|
|
|
2026-10-02 14:42:12 +03:00
|
|
|
// busy is true while an assignment is being worked on; it is reported in
|
|
|
|
|
// the heartbeat body (informational on the control-api side).
|
|
|
|
|
busy atomic.Bool
|
|
|
|
|
|
2026-09-18 10:43:00 +03:00
|
|
|
// registerRetryInitial/Max govern the backoff used while waiting for a
|
|
|
|
|
// successful registration (see registerWithRetry): control-api may not
|
|
|
|
|
// be up yet at agent boot, or may come and go across a redeploy, and the
|
|
|
|
|
// agent should keep waiting rather than exit.
|
|
|
|
|
registerRetryInitial time.Duration
|
|
|
|
|
registerRetryMax time.Duration
|
2026-08-21 07:34:45 +03:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func New(cfg *config.ValidatorAgent, log *slog.Logger) *Agent {
|
|
|
|
|
timeout := time.Duration(cfg.Checks.HTTPSTimeoutSeconds) * time.Second
|
|
|
|
|
if timeout <= 0 {
|
|
|
|
|
timeout = 10 * time.Second
|
|
|
|
|
}
|
|
|
|
|
return &Agent{
|
2026-09-18 10:43:00 +03:00
|
|
|
cfg: cfg,
|
|
|
|
|
client: apiclient.New(cfg.ControlAPIURL, timeout+5*time.Second),
|
|
|
|
|
log: log,
|
|
|
|
|
registerRetryInitial: 3 * time.Second,
|
|
|
|
|
registerRetryMax: 30 * time.Second,
|
2026-08-21 07:34:45 +03:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-10-01 11:35:24 +03:00
|
|
|
// WithToken sets the bearer token sent to the Control API (and only to it:
|
|
|
|
|
// the IP-echo lookup and all check traffic use separate clients).
|
|
|
|
|
func (a *Agent) WithToken(token string) *Agent {
|
|
|
|
|
a.client.Token = token
|
|
|
|
|
return a
|
|
|
|
|
}
|
|
|
|
|
|
2026-08-21 07:34:45 +03:00
|
|
|
// Run registers with the Control API and polls forever until ctx is
|
|
|
|
|
// cancelled.
|
|
|
|
|
func (a *Agent) Run(ctx context.Context) error {
|
2026-09-18 10:43:00 +03:00
|
|
|
if err := a.registerWithRetry(ctx); err != nil {
|
|
|
|
|
return err
|
2026-08-21 07:34:45 +03:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
interval := time.Duration(a.cfg.PollIntervalSeconds) * time.Second
|
2026-10-02 14:42:12 +03:00
|
|
|
|
|
|
|
|
// Heartbeats run on their own schedule. Sent from the poll loop they
|
|
|
|
|
// stopped for as long as a slow assignment took (an address whose
|
|
|
|
|
// outbound targets all time out keeps the loop busy for ~40 s), which
|
|
|
|
|
// control-api reads as a lost validator after heartbeat_timeout_seconds.
|
|
|
|
|
hbCtx, stopHeartbeat := context.WithCancel(ctx)
|
|
|
|
|
defer stopHeartbeat()
|
|
|
|
|
go a.heartbeatLoop(hbCtx, interval)
|
|
|
|
|
|
2026-08-21 07:34:45 +03:00
|
|
|
ticker := time.NewTicker(interval)
|
|
|
|
|
defer ticker.Stop()
|
|
|
|
|
|
|
|
|
|
for {
|
|
|
|
|
a.pollOnce(ctx)
|
|
|
|
|
select {
|
|
|
|
|
case <-ctx.Done():
|
|
|
|
|
return ctx.Err()
|
|
|
|
|
case <-ticker.C:
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
type registerReq struct {
|
|
|
|
|
ValidatorID string `json:"validator_id"`
|
|
|
|
|
Hostname string `json:"hostname"`
|
|
|
|
|
AgentVersion string `json:"agent_version"`
|
|
|
|
|
}
|
|
|
|
|
type registerResp struct {
|
|
|
|
|
OK bool `json:"ok"`
|
|
|
|
|
PollIntervalSeconds int `json:"poll_interval_seconds"`
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (a *Agent) register(ctx context.Context) error {
|
|
|
|
|
hostname, _ := hostnameOrDefault()
|
|
|
|
|
var resp registerResp
|
|
|
|
|
_, err := a.client.Do(ctx, "POST", "/api/v1/agents/register", registerReq{
|
|
|
|
|
ValidatorID: a.cfg.ValidatorID, Hostname: hostname, AgentVersion: "dev",
|
|
|
|
|
}, &resp)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return err
|
|
|
|
|
}
|
|
|
|
|
a.log.Info("registered", "validator_id", a.cfg.ValidatorID)
|
|
|
|
|
return nil
|
|
|
|
|
}
|
|
|
|
|
|
2026-09-18 10:43:00 +03:00
|
|
|
// registerWithRetry retries register with capped exponential backoff until
|
|
|
|
|
// it succeeds or ctx is cancelled. Control-api may not be reachable yet at
|
|
|
|
|
// agent boot (started before control-api, or a network blip), or may reject
|
|
|
|
|
// the request until an admin adds this validator_id to its config — either
|
|
|
|
|
// way the agent should keep waiting rather than exit, since both conditions
|
|
|
|
|
// can resolve on their own after the agent has already started.
|
|
|
|
|
func (a *Agent) registerWithRetry(ctx context.Context) error {
|
|
|
|
|
delay := a.registerRetryInitial
|
|
|
|
|
for {
|
|
|
|
|
err := a.register(ctx)
|
|
|
|
|
if err == nil {
|
|
|
|
|
return nil
|
|
|
|
|
}
|
|
|
|
|
a.log.Warn("registration failed, will retry", "err", err, "retry_in", delay)
|
|
|
|
|
select {
|
|
|
|
|
case <-ctx.Done():
|
|
|
|
|
return ctx.Err()
|
|
|
|
|
case <-time.After(delay):
|
|
|
|
|
}
|
|
|
|
|
delay *= 2
|
|
|
|
|
if delay > a.registerRetryMax {
|
|
|
|
|
delay = a.registerRetryMax
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-08-21 07:34:45 +03:00
|
|
|
type heartbeatReq struct {
|
|
|
|
|
LocalState string `json:"local_state"`
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
type assignmentResp struct {
|
|
|
|
|
IPID int64 `json:"ip_id"`
|
|
|
|
|
IPAddress string `json:"ip_address"`
|
|
|
|
|
Phase string `json:"phase"`
|
|
|
|
|
CheckConfig []checkConfigDTO `json:"check_config"`
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
type checkConfigDTO struct {
|
|
|
|
|
Type string `json:"type"`
|
|
|
|
|
Targets []string `json:"targets"`
|
|
|
|
|
}
|
|
|
|
|
|
2026-10-02 14:42:12 +03:00
|
|
|
// heartbeatLoop sends a heartbeat now and then every interval until ctx is
|
|
|
|
|
// cancelled.
|
|
|
|
|
func (a *Agent) heartbeatLoop(ctx context.Context, interval time.Duration) {
|
|
|
|
|
ticker := time.NewTicker(interval)
|
|
|
|
|
defer ticker.Stop()
|
|
|
|
|
for {
|
|
|
|
|
a.sendHeartbeat(ctx)
|
|
|
|
|
select {
|
|
|
|
|
case <-ctx.Done():
|
|
|
|
|
return
|
|
|
|
|
case <-ticker.C:
|
|
|
|
|
}
|
2026-08-21 07:34:45 +03:00
|
|
|
}
|
2026-10-02 14:42:12 +03:00
|
|
|
}
|
2026-08-21 07:34:45 +03:00
|
|
|
|
2026-10-02 14:42:12 +03:00
|
|
|
func (a *Agent) sendHeartbeat(ctx context.Context) {
|
|
|
|
|
state := "idle"
|
|
|
|
|
if a.busy.Load() {
|
|
|
|
|
state = "checking"
|
|
|
|
|
}
|
|
|
|
|
if _, err := a.client.Do(ctx, "POST", "/api/v1/agents/"+a.cfg.ValidatorID+"/heartbeat", heartbeatReq{LocalState: state}, nil); err != nil && ctx.Err() == nil {
|
|
|
|
|
a.log.Error("heartbeat", "err", err)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (a *Agent) pollOnce(ctx context.Context) {
|
2026-08-21 07:34:45 +03:00
|
|
|
var assignment assignmentResp
|
|
|
|
|
ok, err := a.client.Do(ctx, "GET", "/api/v1/agents/"+a.cfg.ValidatorID+"/assignment", nil, &assignment)
|
|
|
|
|
if err != nil {
|
|
|
|
|
a.log.Error("get assignment", "err", err)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
if !ok {
|
|
|
|
|
a.lastHandledIPID = 0
|
|
|
|
|
return // nothing assigned right now
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
if assignment.IPID == a.lastHandledIPID {
|
|
|
|
|
return // already handled this IP's work this attempt
|
|
|
|
|
}
|
|
|
|
|
|
2026-10-02 14:42:12 +03:00
|
|
|
a.busy.Store(true)
|
|
|
|
|
defer a.busy.Store(false)
|
2026-08-21 07:34:45 +03:00
|
|
|
switch assignment.Phase {
|
|
|
|
|
case "awaiting_self_check":
|
|
|
|
|
a.handleSelfCheckAndRun(ctx, assignment)
|
|
|
|
|
case "checking":
|
|
|
|
|
// Agent restarted (or a prior response was lost) after self-check
|
|
|
|
|
// already succeeded server-side: just (re-)run checks, which is
|
|
|
|
|
// safe since results are upserted idempotently.
|
|
|
|
|
a.runChecks(ctx, assignment)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (a *Agent) handleSelfCheckAndRun(ctx context.Context, assignment assignmentResp) {
|
|
|
|
|
a.postEvent(ctx, assignment.IPID, "config_received", "")
|
|
|
|
|
|
2026-10-02 03:24:20 +03:00
|
|
|
// Each method gets the full timeout (see runSelfCheckMethods), so the
|
|
|
|
|
// overall budget scales with the number of methods.
|
|
|
|
|
timeout := time.Duration(a.cfg.SelfCheck.TimeoutSeconds) * time.Second * time.Duration(len(a.selfCheckMethods()))
|
2026-08-21 07:34:45 +03:00
|
|
|
selfCtx, cancel := context.WithTimeout(ctx, timeout)
|
|
|
|
|
defer cancel()
|
|
|
|
|
|
2026-10-02 03:24:20 +03:00
|
|
|
detectedIP, _, detail, success := a.runSelfCheckMethods(selfCtx, assignment.IPAddress)
|
2026-08-21 07:34:45 +03:00
|
|
|
|
2026-08-21 11:04:49 +03:00
|
|
|
a.postSelfCheck(ctx, assignment.IPID, detectedIP, success, detail)
|
2026-08-21 07:34:45 +03:00
|
|
|
a.postEvent(ctx, assignment.IPID, "self_check_result", fmt.Sprintf(`{"success":%t}`, success))
|
|
|
|
|
|
|
|
|
|
if !success {
|
|
|
|
|
a.log.Warn("self-check failed", "ip", assignment.IPAddress, "detail", detail)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
a.runChecks(ctx, assignment)
|
|
|
|
|
}
|
|
|
|
|
|
2026-10-02 03:24:20 +03:00
|
|
|
// selfCheckMethods returns the configured methods in priority order. Configs
|
|
|
|
|
// built without the loader (tests) may leave the list empty; that means the
|
|
|
|
|
// historical behaviour, ip_echo only.
|
|
|
|
|
func (a *Agent) selfCheckMethods() []string {
|
|
|
|
|
if len(a.cfg.SelfCheck.Methods) == 0 {
|
|
|
|
|
return []string{config.SelfCheckIPEcho}
|
|
|
|
|
}
|
|
|
|
|
return a.cfg.SelfCheck.Methods
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// runSelfCheckMethods tries the configured methods in priority order and
|
|
|
|
|
// stops at the first one that confirms assignedIP. A method that gives no
|
|
|
|
|
// answer and one that reports a different address are treated alike: the
|
|
|
|
|
// next method is tried, since either may be a limitation of that method
|
|
|
|
|
// (e.g. control-api reached over the internal network sees a private
|
|
|
|
|
// address) rather than proof the floating IP is not attached. The check
|
|
|
|
|
// fails only when no method confirms, and detail then carries the reason
|
|
|
|
|
// from every method. detectedIP is the matching address on success, else the
|
|
|
|
|
// last address any method reported (may be empty).
|
|
|
|
|
//
|
|
|
|
|
// Every method runs under its own self_check.timeout_seconds: with one shared
|
|
|
|
|
// deadline a hung first method (priority control_api) would use it all up and
|
|
|
|
|
// the fallback would never get a chance to answer.
|
|
|
|
|
func (a *Agent) runSelfCheckMethods(ctx context.Context, assignedIP string) (detectedIP, method, detail string, ok bool) {
|
|
|
|
|
var reasons []string
|
|
|
|
|
perMethod := time.Duration(a.cfg.SelfCheck.TimeoutSeconds) * time.Second
|
|
|
|
|
for _, m := range a.selfCheckMethods() {
|
|
|
|
|
mctx, cancel := ctx, context.CancelFunc(func() {})
|
|
|
|
|
if perMethod > 0 {
|
|
|
|
|
mctx, cancel = context.WithTimeout(ctx, perMethod)
|
|
|
|
|
}
|
|
|
|
|
var ip string
|
|
|
|
|
var err error
|
|
|
|
|
switch m {
|
|
|
|
|
case config.SelfCheckControlAPI:
|
|
|
|
|
ip, err = a.detectViaControlAPI(mctx)
|
|
|
|
|
case config.SelfCheckIPEcho:
|
|
|
|
|
ip, err = a.detectViaIPEcho(mctx)
|
|
|
|
|
if err != nil {
|
|
|
|
|
err = fmt.Errorf("ip echo request failed: %w", err)
|
|
|
|
|
}
|
|
|
|
|
default:
|
|
|
|
|
err = fmt.Errorf("unknown self-check method")
|
|
|
|
|
}
|
|
|
|
|
cancel()
|
|
|
|
|
if err != nil {
|
|
|
|
|
reasons = append(reasons, m+": "+err.Error())
|
|
|
|
|
continue
|
|
|
|
|
}
|
|
|
|
|
if ip == assignedIP {
|
|
|
|
|
return ip, m, fmt.Sprintf("matched (%s)", m), true
|
|
|
|
|
}
|
|
|
|
|
detectedIP = ip
|
|
|
|
|
reason := fmt.Sprintf("%s: egress ip %q does not match assigned fip %q", m, ip, assignedIP)
|
|
|
|
|
if m == config.SelfCheckControlAPI && isLocalAddr(ip) {
|
|
|
|
|
reason += " (control-api sees a private address; it is reachable over the internal network, self-check via control_api is not possible, use ip_echo)"
|
|
|
|
|
}
|
|
|
|
|
reasons = append(reasons, reason)
|
|
|
|
|
}
|
|
|
|
|
if len(reasons) == 0 {
|
|
|
|
|
reasons = append(reasons, "no self-check methods configured")
|
|
|
|
|
}
|
|
|
|
|
return detectedIP, "", strings.Join(reasons, "; "), false
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// isLocalAddr reports whether ip is a private, loopback or link-local
|
|
|
|
|
// address, i.e. one that can never be a floating IP.
|
|
|
|
|
func isLocalAddr(ip string) bool {
|
|
|
|
|
parsed := net.ParseIP(ip)
|
|
|
|
|
return parsed != nil && (parsed.IsPrivate() || parsed.IsLoopback() || parsed.IsLinkLocalUnicast())
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// detectViaControlAPI asks control-api which source address it sees for this
|
|
|
|
|
// validator. It is only meaningful when control-api is reached over the
|
|
|
|
|
// external network, where the floating IP is the visible source (see
|
|
|
|
|
// config.SelfCheckCfg).
|
|
|
|
|
//
|
|
|
|
|
// Every call dials a brand-new TCP connection through a dedicated transport
|
|
|
|
|
// with keep-alives off: a connection opened before the floating IP was
|
|
|
|
|
// attached (heartbeat, assignment polling) keeps its old NAT state and would
|
|
|
|
|
// keep reporting the previous address, so the shared apiclient must not be
|
|
|
|
|
// used. The endpoint is open, so no agent token is sent.
|
|
|
|
|
func (a *Agent) detectViaControlAPI(ctx context.Context) (string, error) {
|
|
|
|
|
url := strings.TrimRight(a.cfg.ControlAPIURL, "/") + "/api/v1/agents/" + neturl.PathEscape(a.cfg.ValidatorID) + "/observed-ip"
|
|
|
|
|
req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return "", fmt.Errorf("build request: %w", err)
|
|
|
|
|
}
|
|
|
|
|
transport := &http.Transport{DisableKeepAlives: true}
|
|
|
|
|
defer transport.CloseIdleConnections()
|
|
|
|
|
resp, err := (&http.Client{Transport: transport}).Do(req)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return "", err
|
|
|
|
|
}
|
|
|
|
|
defer resp.Body.Close()
|
|
|
|
|
if resp.StatusCode < 200 || resp.StatusCode > 299 {
|
|
|
|
|
return "", fmt.Errorf("unexpected status %d", resp.StatusCode)
|
|
|
|
|
}
|
|
|
|
|
var body struct {
|
|
|
|
|
IP string `json:"ip"`
|
|
|
|
|
}
|
|
|
|
|
if err := json.NewDecoder(io.LimitReader(resp.Body, 4096)).Decode(&body); err != nil {
|
|
|
|
|
return "", fmt.Errorf("decode response: %w", err)
|
|
|
|
|
}
|
|
|
|
|
if net.ParseIP(body.IP) == nil {
|
|
|
|
|
return "", fmt.Errorf("response is not a valid IP: %q", body.IP)
|
|
|
|
|
}
|
|
|
|
|
return body.IP, nil
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// detectViaIPEcho asks each configured IP-echo URL, in order, for the
|
2026-08-21 11:04:49 +03:00
|
|
|
// address this validator is currently seen egressing from, returning the
|
|
|
|
|
// first one that answers with a parseable IP. These must be resources
|
|
|
|
|
// genuinely outside the cloud project (see config.SelfCheckCfg) — OpenStack
|
|
|
|
|
// only applies floating-IP SNAT to traffic leaving via the external
|
|
|
|
|
// network, so anything reachable over the project's internal network would
|
|
|
|
|
// report the validator's private address instead, regardless of whether
|
|
|
|
|
// the floating IP is correctly attached.
|
2026-10-02 03:24:20 +03:00
|
|
|
func (a *Agent) detectViaIPEcho(ctx context.Context) (string, error) {
|
2026-08-21 11:04:49 +03:00
|
|
|
var lastErr error
|
|
|
|
|
for _, url := range a.cfg.SelfCheck.IPEchoURLs {
|
|
|
|
|
ip, err := fetchIPEcho(ctx, url)
|
|
|
|
|
if err != nil {
|
|
|
|
|
lastErr = fmt.Errorf("%s: %w", url, err)
|
|
|
|
|
continue
|
|
|
|
|
}
|
|
|
|
|
return ip, nil
|
|
|
|
|
}
|
|
|
|
|
if lastErr == nil {
|
|
|
|
|
lastErr = fmt.Errorf("no self_check.ip_echo_urls configured")
|
|
|
|
|
}
|
|
|
|
|
return "", lastErr
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// fetchIPEcho performs a single GET against an IP-echo endpoint that
|
|
|
|
|
// returns the caller's address as a bare string in the response body
|
|
|
|
|
// (the common contract shared by services like api.ipify.org,
|
|
|
|
|
// ifconfig.me/ip, icanhazip.com — and by the local stub used in
|
|
|
|
|
// scripts/run-local-e2e.sh).
|
|
|
|
|
func fetchIPEcho(ctx context.Context, url string) (string, error) {
|
|
|
|
|
req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return "", fmt.Errorf("build request: %w", err)
|
|
|
|
|
}
|
|
|
|
|
resp, err := http.DefaultClient.Do(req)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return "", err
|
|
|
|
|
}
|
|
|
|
|
defer resp.Body.Close()
|
|
|
|
|
if resp.StatusCode >= 400 {
|
|
|
|
|
return "", fmt.Errorf("unexpected status %d", resp.StatusCode)
|
|
|
|
|
}
|
|
|
|
|
body, err := io.ReadAll(io.LimitReader(resp.Body, 256))
|
|
|
|
|
if err != nil {
|
|
|
|
|
return "", fmt.Errorf("read response: %w", err)
|
|
|
|
|
}
|
|
|
|
|
ip := strings.TrimSpace(string(body))
|
|
|
|
|
if net.ParseIP(ip) == nil {
|
|
|
|
|
return "", fmt.Errorf("response is not a valid IP: %q", ip)
|
|
|
|
|
}
|
|
|
|
|
return ip, nil
|
|
|
|
|
}
|
|
|
|
|
|
2026-08-21 07:34:45 +03:00
|
|
|
func (a *Agent) runChecks(ctx context.Context, assignment assignmentResp) {
|
|
|
|
|
var results []checkResultDTO
|
|
|
|
|
for _, ct := range assignment.CheckConfig {
|
|
|
|
|
for _, target := range ct.Targets {
|
|
|
|
|
var fn func(context.Context) checkrunner.Result
|
|
|
|
|
switch ct.Type {
|
|
|
|
|
case "https":
|
|
|
|
|
fn = checkrunner.HTTPS(target, time.Duration(a.cfg.Checks.HTTPSTimeoutSeconds)*time.Second)
|
|
|
|
|
case "icmp":
|
|
|
|
|
fn = checkrunner.ICMPEcho(hostOnly(target), a.cfg.Checks.ICMPCount, time.Duration(a.cfg.Checks.ICMPTimeoutSeconds)*time.Second)
|
|
|
|
|
case "ssh":
|
|
|
|
|
if !a.cfg.Checks.SSH.Enabled {
|
|
|
|
|
continue
|
|
|
|
|
}
|
|
|
|
|
fn = checkrunner.SSHBanner(hostOnly(target), time.Duration(a.cfg.Checks.SSH.TimeoutSeconds)*time.Second)
|
|
|
|
|
default:
|
|
|
|
|
a.log.Warn("unknown check type", "type", ct.Type)
|
|
|
|
|
continue
|
|
|
|
|
}
|
|
|
|
|
res := fn(ctx)
|
|
|
|
|
// Record the originally configured target (a full URL for
|
|
|
|
|
// https, e.g.), not checkrunner's internal host-only value
|
|
|
|
|
// used for icmp/ssh — otherwise two configured targets that
|
|
|
|
|
// happen to share a bare host (as can occur, e.g., in the
|
|
|
|
|
// loopback-only local e2e harness) would collide on the
|
|
|
|
|
// checks table's UNIQUE(ip_id, attempt, source, type,
|
|
|
|
|
// target) key and silently overwrite each other.
|
|
|
|
|
results = append(results, checkResultDTO{
|
|
|
|
|
IPID: assignment.IPID, CheckType: res.CheckType, Target: target,
|
|
|
|
|
Success: res.Success, LatencyMS: res.LatencyMS, Detail: res.Detail,
|
|
|
|
|
CheckedAt: res.CheckedAt.Format(time.RFC3339Nano),
|
|
|
|
|
})
|
|
|
|
|
// Report progressively rather than batching until the end, so
|
|
|
|
|
// a crash mid-run doesn't lose already-completed check results.
|
|
|
|
|
a.postResults(ctx, []checkResultDTO{results[len(results)-1]})
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
a.postComplete(ctx, assignment.IPID)
|
|
|
|
|
a.lastHandledIPID = assignment.IPID
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
type checkResultDTO struct {
|
|
|
|
|
IPID int64 `json:"ip_id"`
|
|
|
|
|
CheckType string `json:"check_type"`
|
|
|
|
|
Target string `json:"target,omitempty"`
|
|
|
|
|
Success bool `json:"success"`
|
|
|
|
|
LatencyMS int64 `json:"latency_ms"`
|
|
|
|
|
Detail string `json:"detail,omitempty"`
|
|
|
|
|
CheckedAt string `json:"checked_at"`
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (a *Agent) postResults(ctx context.Context, results []checkResultDTO) {
|
|
|
|
|
body := struct {
|
|
|
|
|
Results []checkResultDTO `json:"results"`
|
|
|
|
|
}{results}
|
|
|
|
|
if _, err := a.client.Do(ctx, "POST", "/api/v1/agents/"+a.cfg.ValidatorID+"/results", body, nil); err != nil {
|
|
|
|
|
a.log.Error("post results", "err", err)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (a *Agent) postComplete(ctx context.Context, ipID int64) {
|
|
|
|
|
body := struct {
|
|
|
|
|
IPID int64 `json:"ip_id"`
|
|
|
|
|
}{ipID}
|
|
|
|
|
if _, err := a.client.Do(ctx, "POST", "/api/v1/agents/"+a.cfg.ValidatorID+"/complete", body, nil); err != nil {
|
|
|
|
|
a.log.Error("post complete", "err", err)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (a *Agent) postSelfCheck(ctx context.Context, ipID int64, detected string, success bool, detail string) {
|
|
|
|
|
body := struct {
|
|
|
|
|
IPID int64 `json:"ip_id"`
|
|
|
|
|
DetectedEgress string `json:"detected_egress_ip"`
|
|
|
|
|
Success bool `json:"success"`
|
|
|
|
|
Detail string `json:"detail"`
|
|
|
|
|
}{ipID, detected, success, detail}
|
|
|
|
|
if _, err := a.client.Do(ctx, "POST", "/api/v1/agents/"+a.cfg.ValidatorID+"/self-check", body, nil); err != nil {
|
|
|
|
|
a.log.Error("post self-check", "err", err)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (a *Agent) postEvent(ctx context.Context, ipID int64, eventType, payload string) {
|
|
|
|
|
body := struct {
|
|
|
|
|
EventType string `json:"event_type"`
|
|
|
|
|
IPID int64 `json:"ip_id"`
|
|
|
|
|
Payload string `json:"payload"`
|
|
|
|
|
}{eventType, ipID, payload}
|
|
|
|
|
if _, err := a.client.Do(ctx, "POST", "/api/v1/agents/"+a.cfg.ValidatorID+"/events", body, nil); err != nil {
|
|
|
|
|
a.log.Error("post event", "type", eventType, "err", err)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func hostnameOrDefault() (string, error) {
|
|
|
|
|
h, err := os.Hostname()
|
|
|
|
|
if err != nil || h == "" {
|
|
|
|
|
return "unknown", err
|
|
|
|
|
}
|
|
|
|
|
return h, nil
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// hostOnly strips a URL scheme (https://) from a target, since ICMP/SSH
|
|
|
|
|
// checks operate on bare hostnames while HTTPS checks take a full URL.
|
|
|
|
|
func hostOnly(target string) string {
|
|
|
|
|
for _, prefix := range []string{"https://", "http://"} {
|
|
|
|
|
if len(target) > len(prefix) && target[:len(prefix)] == prefix {
|
|
|
|
|
target = target[len(prefix):]
|
|
|
|
|
break
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
for i := 0; i < len(target); i++ {
|
|
|
|
|
if target[i] == '/' || target[i] == ':' {
|
|
|
|
|
return target[:i]
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
return target
|
|
|
|
|
}
|