Files
cloud-ip-validator/internal/probercore/probercore.go
T
ayurishchevandClaude Sonnet 5 95f8066eed Retry validator-agent/prober registration; show validator hostname
Both binaries registered with control-api exactly once at startup and
exited (os.Exit(1)) on any failure — including control-api simply not
being up yet (no ordering guarantee between the two at boot/redeploy) or
the admin not having added this validator_id/site_id to the config yet.
Run() now retries registration with capped exponential backoff (3s->30s)
until it succeeds or the process is asked to shut down, instead of
crashing; registerWithRetry is identical in agentcore and probercore
since their Run/register shape already was.

Separately, the admin dashboard's Validators page had no hostname column
even though the agent already reports one on register (mirroring the
prober) and control-api already persists it — only the admin-config read
DTO (validatorDTO in httpapi and dashboard) dropped it before it reached
the template. Added hostname + last_heartbeat_at to that DTO end-to-end
and a Хост/Heartbeat column to validators.html, matching sites.html.

Rebuilt bin/{control-api,admin-dashboard,prober,validator-agent} and
bin/SHA256SUMS per docs/SETUP.md's documented build recipe.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-09-18 10:43:00 +03:00

196 lines
5.6 KiB
Go

// Package probercore implements a prober's poll loop: register once with
// its site identity, then each cycle fetch the set of IPs currently under
// test and run inbound reachability checks (TCP connect on each configured
// port, ICMP echo) directly against each one — this is the real,
// unmediated network test; only the control channel goes through the
// Control API.
package probercore
import (
"context"
"log/slog"
"os"
"time"
"cloudipvalidator/internal/apiclient"
"cloudipvalidator/internal/checkrunner"
"cloudipvalidator/internal/config"
)
type Prober struct {
cfg *config.Prober
client *apiclient.Client
log *slog.Logger
// registerRetryInitial/Max govern the backoff used while waiting for a
// successful registration (see registerWithRetry): control-api may not
// be up yet at prober boot, or may reject an unconfigured site_id until
// an admin adds it — the prober should keep waiting rather than exit.
registerRetryInitial time.Duration
registerRetryMax time.Duration
}
func New(cfg *config.Prober, log *slog.Logger) *Prober {
timeout := time.Duration(cfg.Checks.TCPTimeoutSeconds) * time.Second
if timeout <= 0 {
timeout = 10 * time.Second
}
return &Prober{
cfg: cfg,
client: apiclient.New(cfg.ControlAPIURL, timeout+5*time.Second),
log: log,
registerRetryInitial: 3 * time.Second,
registerRetryMax: 30 * time.Second,
}
}
func (p *Prober) Run(ctx context.Context) error {
if err := p.registerWithRetry(ctx); err != nil {
return err
}
interval := time.Duration(p.cfg.PollIntervalSeconds) * time.Second
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
p.pollOnce(ctx)
select {
case <-ctx.Done():
return ctx.Err()
case <-ticker.C:
}
}
}
type registerReq struct {
SiteID string `json:"site_id"`
Hostname string `json:"hostname"`
}
func (p *Prober) register(ctx context.Context) error {
hostname, _ := os.Hostname()
_, err := p.client.Do(ctx, "POST", "/api/v1/probers/register", registerReq{SiteID: p.cfg.SiteID, Hostname: hostname}, nil)
if err != nil {
return err
}
p.log.Info("registered", "site_id", p.cfg.SiteID)
return nil
}
// registerWithRetry retries register with capped exponential backoff until
// it succeeds or ctx is cancelled — see agentcore.Agent.registerWithRetry
// for the identical rationale (control-api may start later, or an admin may
// add this site_id to its config later).
func (p *Prober) registerWithRetry(ctx context.Context) error {
delay := p.registerRetryInitial
for {
err := p.register(ctx)
if err == nil {
return nil
}
p.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 > p.registerRetryMax {
delay = p.registerRetryMax
}
}
}
type assignment struct {
IPID int64 `json:"ip_id"`
IPAddress string `json:"ip_address"`
Ports []int `json:"ports"`
ICMP bool `json:"icmp"`
}
type resultDTO struct {
IPID int64 `json:"ip_id"`
IPAddress string `json:"ip_address"`
CheckType string `json:"check_type"`
Success bool `json:"success"`
LatencyMS int64 `json:"latency_ms"`
Detail string `json:"detail,omitempty"`
CheckedAt string `json:"checked_at"`
Complete bool `json:"complete"`
}
func (p *Prober) pollOnce(ctx context.Context) {
if _, err := p.client.Do(ctx, "POST", "/api/v1/probers/"+p.cfg.SiteID+"/heartbeat", nil, nil); err != nil {
p.log.Error("heartbeat", "err", err)
return
}
var assignments []assignment
ok, err := p.client.Do(ctx, "GET", "/api/v1/probers/"+p.cfg.SiteID+"/assignments", nil, &assignments)
if err != nil {
p.log.Error("get assignments", "err", err)
return
}
if !ok || len(assignments) == 0 {
return
}
for _, a := range assignments {
p.probeOne(ctx, a)
}
}
// extraChecksForPort returns checks to run in addition to the base
// TCPConnect for ports where a bare TCP handshake doesn't actually prove
// the expected service is behind it: a real TLS handshake on 443, a real
// SSH banner exchange on 22. nil for every other port.
func extraChecksForPort(host string, port int, timeout time.Duration) []func(ctx context.Context) checkrunner.Result {
switch port {
case 443:
return []func(ctx context.Context) checkrunner.Result{
checkrunner.TLSHandshake(host, port, timeout),
}
case 22:
return []func(ctx context.Context) checkrunner.Result{
checkrunner.SSHBanner(host, timeout),
}
default:
return nil
}
}
func (p *Prober) probeOne(ctx context.Context, a assignment) {
tcpTimeout := time.Duration(p.cfg.Checks.TCPTimeoutSeconds) * time.Second
icmpTimeout := time.Duration(p.cfg.Checks.ICMPTimeoutSeconds) * time.Second
var results []resultDTO
record := func(res checkrunner.Result) {
results = append(results, resultDTO{
IPID: a.IPID, IPAddress: a.IPAddress, CheckType: res.CheckType,
Success: res.Success, LatencyMS: res.LatencyMS, Detail: res.Detail,
CheckedAt: res.CheckedAt.Format(time.RFC3339Nano),
})
}
for _, port := range a.Ports {
record(checkrunner.TCPConnect(a.IPAddress, port, tcpTimeout)(ctx))
for _, extra := range extraChecksForPort(a.IPAddress, port, tcpTimeout) {
record(extra(ctx))
}
}
if a.ICMP {
record(checkrunner.ICMPEcho(a.IPAddress, p.cfg.Checks.ICMPCount, icmpTimeout)(ctx))
}
if len(results) > 0 {
results[len(results)-1].Complete = true
}
body := struct {
Results []resultDTO `json:"results"`
}{results}
if _, err := p.client.Do(ctx, "POST", "/api/v1/probers/"+p.cfg.SiteID+"/results", body, nil); err != nil {
p.log.Error("post results", "ip", a.IPAddress, "err", err)
}
}