Files
echolot/server/cmd/echolot-server/main.go
T
mrambossekandClaude Opus 5 4ae744aae5
server-release / image (push) Successful in 15s
server-test / test (push) Successful in 27s
server-release / release (push) Successful in 28s
server: self-test — sysctl audit + egress-MTU self-proof ("server proven good")
A measurement server must prove its own host isn't distorting results:
- sysctl audit (/proc/sys): flags accept_ra on a static host, ICMP
  redirects, ICMP rate-limiting of the server's own errors, and disabled
  TCP options — each a measurement-fidelity hazard, with the "why".
- egress-MTU self-proof: DF PMTUD probe (IP_MTU_DISCOVER + getsockopt
  IP_MTU, no root — Linux-only, stub elsewhere) to external anchors. If the
  server's own uplink is below 1500, client MTU tests measure THIS server,
  so we say so.
Exposed at GET /admin/selftest (full report) and as server_selftest
{mtu_ok, sysctl_ok} in the profile so clients can trust or skip MTU tests.
Recommended deploy/99-echolot-sysctl.conf + README section.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-07-31 20:44:38 +02:00

403 lines
13 KiB
Go

// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
// echolot-server: the probe server (spec: docs/probe-protocol.md).
//
// Modes:
// - container ("docker mode"): autodetected (or --docker / ECHOLOT_DOCKER=1);
// config comes from ECHOLOT_* env vars; systemd/self-update are refused.
// - native: same flags/env, plus --install-systemd / --uninstall-systemd
// and opt-in --self-update against a Gitea releases API.
package main
import (
"context"
"crypto/ecdsa"
"crypto/elliptic"
"crypto/rand"
"crypto/tls"
"crypto/x509"
"crypto/x509/pkix"
"encoding/json"
"encoding/pem"
"errors"
"fmt"
"log/slog"
"math/big"
"net"
"net/http"
"net/netip"
"os"
"os/signal"
"path/filepath"
"strconv"
"sync/atomic"
"syscall"
"time"
"echo-lot.app/server/internal/canarydns"
"echo-lot.app/server/internal/config"
"echo-lot.app/server/internal/control"
"echo-lot.app/server/internal/dataplane"
"echo-lot.app/server/internal/selftest"
"echo-lot.app/server/internal/selfupdate"
"echo-lot.app/server/internal/session"
"echo-lot.app/server/internal/store"
"echo-lot.app/server/internal/stun"
"echo-lot.app/server/internal/system"
"echo-lot.app/server/internal/tcpecho"
)
// Version is stamped via -ldflags "-X main.Version=v1.2.3" in CI.
var Version = "dev"
func main() {
if err := run(); err != nil {
slog.Error("fatal", "err", err)
os.Exit(1)
}
}
func run() error {
cfg, actions, err := config.Load(os.Args[1:])
if err != nil {
return err
}
control.Version = Version
switch {
case actions.Version:
fmt.Println(Version)
return nil
case actions.InstallSystemd:
// The unit runs this same binary in serve mode; host config comes
// from /etc/echolot-server.env. A self-update timer is installed
// only when an update API is configured.
return system.InstallSystemd(cfg.SelfUpdateAPI)
case actions.UninstallSystemd:
return system.UninstallSystemd()
case actions.SelfUpdate:
return selfupdate.Run(cfg.SelfUpdateAPI, Version)
}
return serve(cfg)
}
func serve(cfg *config.Config) error {
slog.Info("echolot-server starting", "version", Version, "mode",
map[bool]string{true: "container", false: "native"}[cfg.Docker],
"state_dir", cfg.StateDir)
st, err := store.Open(cfg.StateDir)
if err != nil {
return fmt.Errorf("state store: %w", err)
}
cert, err := loadOrCreateCert(cfg)
if err != nil {
return fmt.Errorf("tls: %w", err)
}
pin, err := control.SpkiPinB64(cert)
if err != nil {
return err
}
slog.Info("control-plane certificate", "pin-sha256", pin)
sessions := session.NewManager(15 * time.Minute)
dp := &dataplane.Server{Sessions: sessions}
tcpSrv := &tcpecho.Server{}
caps := []string{"udp-probe", "delayed-echo", "connect-back", "http-echo"}
if len(config.Addrs(cfg.TCPListen)) > 0 {
caps = append(caps, "tcp-echo")
}
ctl := &control.Server{
Store: st, Sessions: sessions, Name: cfg.Name,
UDPPort: mustPort(firstAddr(cfg.UDPListen)), TCPPort: mustPort(firstAddr(cfg.TCPListen)),
StunPort: mustPort(firstAddr(cfg.StunListen)), PinB64: pin, CertChain: cert.Certificate,
DelayedEcho: dp.SendDelayedEcho,
TCPRecent: func(ip string) any { return tcpSrv.RecentFor(ip) },
}
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer stop()
ctlAddrs := config.Addrs(cfg.ControlListen)
udpAddrs := config.Addrs(cfg.UDPListen)
errCh := make(chan error, len(ctlAddrs)+len(udpAddrs)+2)
// Control plane (HTTPS, pin-based trust) — one shared server, one
// listener per configured address; Shutdown closes them all.
ctlSrv := &http.Server{
Handler: ctl.Handler(),
TLSConfig: &tls.Config{Certificates: []tls.Certificate{cert}, MinVersion: tls.VersionTLS12},
ReadHeaderTimeout: 10 * time.Second,
}
for _, addr := range ctlAddrs {
ln, err := net.Listen("tcp", addr)
if err != nil {
return fmt.Errorf("control listen %s: %w", addr, err)
}
go func(a string, l net.Listener) {
errCh <- fmt.Errorf("control %s: %w", a, ctlSrv.ServeTLS(l, "", ""))
}(addr, ln)
}
// Self-test: prove the host is a clean measurement target. Sysctl audit is
// instant; the egress-MTU proof does network round trips, so publish the
// sysctl-only report immediately and swap in the full one when it lands.
var selftestPtr atomic.Pointer[selftest.Report]
initial := selftest.Report{Sysctls: selftest.Sysctls()}
selftestPtr.Store(&initial)
for _, c := range initial.Sysctls {
if c.Severity == selftest.Warn {
slog.Warn("sysctl not measurement-clean", "sysctl", c.Name, "got", c.Got, "want", c.Want, "why", c.Why)
}
}
go func() {
r := selftest.Run(config.Addrs(cfg.MTUProbeTargets))
selftestPtr.Store(&r)
for _, m := range r.EgressMTU {
if !m.FullMTU {
slog.Warn("egress MTU below 1500 — client MTU results measure THIS server, not the client",
"target", m.Target, "discovered_mtu", m.DiscoveredMTU, "err", m.Err)
}
}
slog.Info("self-test complete", "sysctl_ok", r.SysctlOK, "mtu_ok", r.MTUOK)
}()
ctl.ProvenGood = func() (mtuOK, sysctlOK bool) {
r := selftestPtr.Load()
return r.MTUOK, r.SysctlOK
}
// Admin/health (plain HTTP, localhost by default; spec §7)
admin := http.NewServeMux()
admin.HandleFunc("GET /healthz", func(w http.ResponseWriter, _ *http.Request) {
fmt.Fprintf(w, `{"ok":true,"version":%q}`, Version)
})
admin.HandleFunc("GET /admin/selftest", func(w http.ResponseWriter, _ *http.Request) {
w.Header().Set("Content-Type", "application/json")
_ = json.NewEncoder(w).Encode(selftestPtr.Load())
})
// TODO(spec §7): enrollment token management + device list. Until the
// admin UI exists, mint tokens with: echolot-admin (or curl on this
// listener once the endpoint lands).
admin.HandleFunc("POST /admin/enroll-tokens", func(w http.ResponseWriter, r *http.Request) {
tok, err := st.NewEnrollToken(24*time.Hour, r.URL.Query().Get("note"))
if err != nil {
http.Error(w, err.Error(), 500)
return
}
fmt.Fprintf(w, `{"token":%q,"expires_in_s":86400}`+"\n", tok)
})
adminSrv := &http.Server{Addr: cfg.AdminListen, Handler: admin, ReadHeaderTimeout: 10 * time.Second}
go func() { errCh <- fmt.Errorf("admin: %w", adminSrv.ListenAndServe()) }()
// UDP data plane — one socket per configured address. Distinct sockets
// (not wildcard) also guarantee responses leave from the address the
// request arrived on, which stun-5780 will rely on.
var udpConns []*net.UDPConn
for _, addr := range udpAddrs {
udpAddr, err := net.ResolveUDPAddr("udp", addr)
if err != nil {
return fmt.Errorf("udp addr %s: %w", addr, err)
}
conn, err := net.ListenUDP("udp", udpAddr)
if err != nil {
return fmt.Errorf("udp listen %s: %w", addr, err)
}
udpConns = append(udpConns, conn)
go func(a string, c *net.UDPConn) {
errCh <- fmt.Errorf("udp %s: %w", a, dp.Serve(c))
}(addr, conn)
}
// TCP echo (spec §4)
var tcpLns []net.Listener
for _, addr := range config.Addrs(cfg.TCPListen) {
ln, err := net.Listen("tcp", addr)
if err != nil {
return fmt.Errorf("tcp listen %s: %w", addr, err)
}
tcpLns = append(tcpLns, ln)
go func(a string, l net.Listener) {
errCh <- fmt.Errorf("tcp %s: %w", a, tcpSrv.Serve(l))
}(addr, ln)
}
// Optional cleartext HTTP-echo (spec §4 plaintext-path test) — only
// POST /v1/echo, no auth, no secrets. Off unless configured.
var httpEchoSrvs []*http.Server
for _, addr := range config.Addrs(cfg.HTTPEchoListen) {
hs := &http.Server{Addr: addr, Handler: ctl.EchoHandler(), ReadHeaderTimeout: 10 * time.Second}
httpEchoSrvs = append(httpEchoSrvs, hs)
go func(a string, srv *http.Server) { errCh <- fmt.Errorf("http-echo %s: %w", a, srv.ListenAndServe()) }(addr, hs)
}
// STUN (spec §4) — advertises stun-5780 only with ≥2 same-family addrs.
var stunSrv *stun.Server
if stunAddrs := config.Addrs(cfg.StunListen); len(stunAddrs) > 0 {
stunSrv, err = stun.Listen(stunAddrs)
if err != nil {
return fmt.Errorf("stun listen: %w", err)
}
go func() { errCh <- fmt.Errorf("stun: %w", stunSrv.Serve()) }()
if stunSrv.Has5780() {
ctl.Capabilities = append(caps, "stun-5780")
} else {
ctl.Capabilities = append(caps, "stun-basic")
}
} else {
ctl.Capabilities = caps
}
// Canary DNS (spec §6.1) — authoritative for CanaryZone, udp+tcp per addr.
var dnsUDP []*net.UDPConn
var dnsTCP []net.Listener
if dnsAddrs := config.Addrs(cfg.DNSListen); len(dnsAddrs) > 0 && cfg.CanaryZone != "" {
v4, v6 := firstByFamily(dnsAddrs)
cd := canarydns.New(cfg.CanaryZone, cfg.Name, v4, v6)
for _, addr := range dnsAddrs {
ua, err := net.ResolveUDPAddr("udp", addr)
if err != nil {
return fmt.Errorf("dns udp addr %s: %w", addr, err)
}
uc, err := net.ListenUDP("udp", ua)
if err != nil {
return fmt.Errorf("dns udp listen %s: %w", addr, err)
}
dnsUDP = append(dnsUDP, uc)
go func(a string, c *net.UDPConn) { errCh <- fmt.Errorf("dns-udp %s: %w", a, cd.ServeUDP(c)) }(addr, uc)
tl, err := net.Listen("tcp", addr)
if err != nil {
return fmt.Errorf("dns tcp listen %s: %w", addr, err)
}
dnsTCP = append(dnsTCP, tl)
go func(a string, l net.Listener) { errCh <- fmt.Errorf("dns-tcp %s: %w", a, cd.ServeTCP(l)) }(addr, tl)
}
ctl.CanaryZone = cfg.CanaryZone
ctl.CanaryQueries = func(prefix string) any { return cd.RecentForPrefix(prefix) }
ctl.Capabilities = append(ctl.Capabilities, "canary-dns")
}
slog.Info("listening",
"control", ctlAddrs, "admin", cfg.AdminListen, "udp", udpAddrs,
"tcp", config.Addrs(cfg.TCPListen), "stun", config.Addrs(cfg.StunListen),
"dns", config.Addrs(cfg.DNSListen), "capabilities", ctl.Capabilities)
select {
case <-ctx.Done():
slog.Info("shutting down")
shutCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
_ = ctlSrv.Shutdown(shutCtx)
_ = adminSrv.Shutdown(shutCtx)
for _, c := range udpConns {
_ = c.Close()
}
for _, l := range tcpLns {
_ = l.Close()
}
if stunSrv != nil {
stunSrv.Close()
}
for _, c := range dnsUDP {
_ = c.Close()
}
for _, l := range dnsTCP {
_ = l.Close()
}
for _, hs := range httpEchoSrvs {
_ = hs.Shutdown(shutCtx)
}
return nil
case err := <-errCh:
return err
}
}
// firstByFamily returns the first v4 and first v6 address from a list of
// "ip:port" specs — used for the canary zone's apex/NS answers.
func firstByFamily(addrs []string) (v4, v6 netip.Addr) {
for _, a := range addrs {
if ap, err := netip.ParseAddrPort(a); err == nil {
if ap.Addr().Unmap().Is4() && !v4.IsValid() {
v4 = ap.Addr().Unmap()
} else if ap.Addr().Is6() && !ap.Addr().Is4In6() && !v6.IsValid() {
v6 = ap.Addr()
}
}
}
return
}
func firstAddr(spec string) string {
if a := config.Addrs(spec); len(a) > 0 {
return a[0]
}
return ""
}
func mustPort(listen string) int {
_, p, err := net.SplitHostPort(listen)
if err != nil {
return 0
}
n, _ := strconv.Atoi(p)
return n
}
// loadOrCreateCert loads the configured cert/key, or generates a self-signed
// ECDSA P-256 cert into the state dir on first start. Self-signed is
// first-class per spec §1 — clients pin the SPKI, they don't chase CAs.
func loadOrCreateCert(cfg *config.Config) (tls.Certificate, error) {
certPath, keyPath := cfg.TLSCert, cfg.TLSKey
if certPath == "" {
certPath = filepath.Join(cfg.StateDir, "tls-cert.pem")
keyPath = filepath.Join(cfg.StateDir, "tls-key.pem")
}
if c, err := tls.LoadX509KeyPair(certPath, keyPath); err == nil {
return c, nil
} else if cfg.TLSCert != "" {
// Explicitly configured paths must exist — do not silently overwrite.
return tls.Certificate{}, fmt.Errorf("loading configured cert %s: %w", certPath, err)
}
key, err := ecdsa.GenerateKey(elliptic.P256(), rand.Reader)
if err != nil {
return tls.Certificate{}, err
}
serial, _ := rand.Int(rand.Reader, big.NewInt(1).Lsh(big.NewInt(1), 62))
tmpl := x509.Certificate{
SerialNumber: serial,
Subject: pkix.Name{CommonName: "echolot-server"},
NotBefore: time.Now().Add(-time.Hour),
// Long-lived on purpose: trust is the SPKI pin from enrollment, not
// cert validity; rotation happens via next_pins in the profile.
NotAfter: time.Now().AddDate(10, 0, 0),
KeyUsage: x509.KeyUsageDigitalSignature,
ExtKeyUsage: []x509.ExtKeyUsage{x509.ExtKeyUsageServerAuth},
BasicConstraintsValid: true,
}
der, err := x509.CreateCertificate(rand.Reader, &tmpl, &tmpl, &key.PublicKey, key)
if err != nil {
return tls.Certificate{}, err
}
keyDer, err := x509.MarshalECPrivateKey(key)
if err != nil {
return tls.Certificate{}, err
}
if err := os.MkdirAll(cfg.StateDir, 0o700); err != nil && !errors.Is(err, os.ErrExist) {
return tls.Certificate{}, err
}
certPem := pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: der})
keyPem := pem.EncodeToMemory(&pem.Block{Type: "EC PRIVATE KEY", Bytes: keyDer})
if err := os.WriteFile(certPath, certPem, 0o644); err != nil {
return tls.Certificate{}, err
}
if err := os.WriteFile(keyPath, keyPem, 0o600); err != nil {
return tls.Certificate{}, err
}
slog.Info("generated self-signed certificate", "cert", certPath)
return tls.X509KeyPair(certPem, keyPem)
}