164 lines
4.5 KiB
Go
164 lines
4.5 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"
|
|
"fmt"
|
|
"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
|
|
}
|
|
|
|
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,
|
|
}
|
|
}
|
|
|
|
func (p *Prober) Run(ctx context.Context) error {
|
|
if err := p.register(ctx); err != nil {
|
|
return fmt.Errorf("register: %w", 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
|
|
}
|
|
|
|
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)
|
|
}
|
|
}
|