Files
cloud-ip-validator/internal/agentcore/agentcore.go
T
ayurishchevandClaude Sonnet 5.5 debf2afed2 Add authentication: admin/agent bearer tokens for the API, login for the dashboard
control-api: every route now carries a mandatory access level (admin / agent /
open) in a route table. All /api/v1/admin/* require the admin token; the
write calls of validator-agent and prober (self-check, events, results,
complete) require a separate static agent token; register, heartbeat and
fetching the assignment stay open. Tokens come from env vars, are compared in
constant time and never logged. An empty token leaves that level open with a
startup warning (backward compatible).

validator-agent / prober: apiclient sends the agent token only to control-api.

admin-dashboard: login/password (from env) with a stateless HMAC session
cookie, Origin-based CSRF check, per-IP brute-force throttle, HX-Redirect for
htmx polls, logout in the sidebar; the dashboard calls control-api with the
admin token. Login page layout fixed after review.

Also: env plumbing in docker-compose/rxprod-compose/systemd/config examples,
e2e script with token assertions, tests, docs (API, SETUP, USAGE, DASHBOARD,
README), plan and review under docs/changes/, bin/ rebuilt with new
SHA256SUMS.

Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
2026-10-01 11:35:24 +03:00

377 lines
12 KiB
Go

// 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"
"fmt"
"io"
"log/slog"
"net"
"net/http"
"os"
"strings"
"time"
"cloudipvalidator/internal/apiclient"
"cloudipvalidator/internal/checkrunner"
"cloudipvalidator/internal/config"
)
type Agent struct {
cfg *config.ValidatorAgent
client *apiclient.Client
log *slog.Logger
lastHandledIPID int64
// 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
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"`
}
func (a *Agent) pollOnce(ctx context.Context) {
if _, err := a.client.Do(ctx, "POST", "/api/v1/agents/"+a.cfg.ValidatorID+"/heartbeat", heartbeatReq{LocalState: "idle"}, nil); err != nil {
a.log.Error("heartbeat", "err", err)
return
}
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
}
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", "")
timeout := time.Duration(a.cfg.SelfCheck.TimeoutSeconds) * time.Second
selfCtx, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
detectedIP, err := a.detectPublicIP(selfCtx)
success := err == nil && detectedIP == assignment.IPAddress
detail := "matched"
if err != nil {
detail = "ip echo request failed: " + err.Error()
} else if !success {
detail = fmt.Sprintf("egress ip %q does not match assigned fip %q", detectedIP, 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)
}
// detectPublicIP 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) detectPublicIP(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
}