// 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" "encoding/json" "fmt" "io" "log/slog" "net" "net/http" neturl "net/url" "os" "strings" "sync/atomic" "time" "cloudipvalidator/internal/apiclient" "cloudipvalidator/internal/checkrunner" "cloudipvalidator/internal/config" ) type Agent struct { cfg *config.ValidatorAgent client *apiclient.Client log *slog.Logger lastHandledIPID int64 // 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 // 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 } 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{ cfg: cfg, client: apiclient.New(cfg.ControlAPIURL, timeout+5*time.Second), log: log, registerRetryInitial: 3 * time.Second, registerRetryMax: 30 * time.Second, } } // 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 } // Run registers with the Control API and polls forever until ctx is // cancelled. func (a *Agent) Run(ctx context.Context) error { if err := a.registerWithRetry(ctx); err != nil { return err } interval := time.Duration(a.cfg.PollIntervalSeconds) * time.Second // 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) 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 } // 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 } } } 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"` } // 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: } } } 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) { 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 } a.busy.Store(true) defer a.busy.Store(false) 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", "") // 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())) selfCtx, cancel := context.WithTimeout(ctx, timeout) defer cancel() detectedIP, _, detail, success := a.runSelfCheckMethods(selfCtx, assignment.IPAddress) a.postSelfCheck(ctx, assignment.IPID, detectedIP, success, detail) 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) } // 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 // 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. func (a *Agent) detectViaIPEcho(ctx context.Context) (string, error) { 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 } 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 }