Compare commits

..
Author SHA1 Message Date
mrambossekandClaude Opus 5 d5e15816b5 server: fix egress-MTU probe — connect the socket before reading IP_MTU
server-release / image (push) Successful in 15s
server-test / test (push) Successful in 27s
server-release / release (push) Successful in 27s
IP_MTU getsockopt returns ENOTCONN on an unconnected socket; the v0.3.4
probe set IP_MTU_DISCOVER and Sendto but never Connect'd, so every probe
errored. UDP-connect (no handshake) pins the route so IP_MTU reflects the
path; switched to Write (two return values). Sysctl audit already flagged
the four real fmr issues in v0.3.4; this makes the MTU proof report.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-07-31 20:46:47 +02:00
mrambossekandClaude Opus 5 4ae744aae5 server: self-test — sysctl audit + egress-MTU self-proof ("server proven good")
server-release / image (push) Successful in 15s
server-test / test (push) Successful in 27s
server-release / release (push) Successful in 28s
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
mrambossekandClaude Opus 5 c9e0d06ea2 server: MTU probe (MTU_PROBE/MTU_ACK) — path-MTU / black-hole measurement
server-release / image (push) Successful in 14s
server-test / test (push) Successful in 27s
server-release / release (push) Successful in 27s
Server ACKs each DF-flagged probe with a tiny MTU_ACK carrying the size it
received; the client binary-searches the path MTU. Non-amplifying by
construction. Tested.

Also records: v0.3.2 (http-echo + tls-reference) verified live on fmr, and
the finding that upstream trains are already observable via the
observations API (dedicated TRAIN_REPORT deferred — needs an
anti-amplification grant + columnar encoding).

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-07-31 20:37:20 +02:00
mrambossekandClaude Opus 5 38fb73c34e server: HTTP echo + TLS reference (control-plane security measurements)
server-release / image (push) Successful in 15s
server-test / test (push) Successful in 27s
server-release / release (push) Successful in 28s
- POST /v1/echo: returns the received request head + body (base64) and the
  observed TLS parameters (version, cipher, SNI, ALPN, resumed). The client
  diffs against what it sent to detect header injection/stripping,
  transparent proxying, or TLS interception (sec.http_echo). http-echo
  added to the capability set.
- GET /v1/tls-reference: the served leaf-first DER chain + pin, so the app
  can compare an out-of-band copy against its own handshake (sec.tls_reference).
  Always available, no auth — public handshake info.
- Optional CLEARTEXT http-echo listener (ECHOLOT_HTTP_ECHO_LISTEN, default
  off) exposing only /v1/echo for the plaintext-path tampering test.

Live-smoke-tested (HTTPS echo reflected an injected header + observed
TLS1.3; cleartext variant reports tls:none); httptest unit tests added.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-07-31 20:34:03 +02:00
mrambossekandClaude Opus 5 379153219e build-status: canary DNS live on fmr — session attribution + 0x20 finding
Zone delegated + authoritative, verified via public recursion; per-session
nonce queries attributed in the observations API. First test caught
Google's 0x20 case randomization vs Cloudflare's plain case.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-07-31 20:18:50 +02:00
mrambossekandClaude Opus 5 35baf70cdb server: canary DNS — authoritative zone with frozen §6.1 reference records
server-release / image (push) Successful in 15s
server-test / test (push) Successful in 26s
server-release / release (push) Successful in 27s
Stdlib DNS responder (no external deps): parses single-question queries
with EDNS OPT (bufsize, DO, ECS), serves the spec's frozen reference
records (ttl-{5,60,3600,86400} A/AAAA/TXT, many-rr 8×A in order, big-txt
~1800B), and per-query <nonce>.<session>.<zone> answers in 192.0.2.0/24.
UDP truncation sets TC past 512 (or the EDNS bufsize); TCP never
truncates — the EDNS-bufsize / TCP-fallback test. Every query is logged
(qname, resolver, transport, EDNS, ECS, case) and surfaced per session
prefix in GET /v1/sessions/{id}/observations as dns_canary. Profile gains
canary_zone + the canary-dns capability when configured.

Wire format validated against an independent client (correct rcodes,
answer counts, TC behavior, full EDNS response); unit tests cover
references, truncation-vs-EDNS, logging, NXDOMAIN.

Versioning: patch-first convention recorded in CLAUDE.md.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-07-31 20:08:30 +02:00
mrambossekandClaude Opus 5 4f5499198b build-status: server v0.3.0 live on fmr — STUN/TCP-echo/observations/actions verified
Deployed via self-update (first real run). External checks: stun-5780
advertised, STUN binding OK v4+v6 with OTHER-ADDRESS, TCP echo mss=1440
over IPv6.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-07-31 19:56:22 +02:00
mrambossekandClaude Opus 5 7b676e666e server: STUN, TCP echo, observations API, delayed-echo + connect-back actions
server-test / test (push) Successful in 27s
server-release / image (push) Successful in 14s
server-release / release (push) Successful in 27s
- stun: RFC 5389 binding responder + RFC 5780 attributes (OTHER-ADDRESS,
  RESPONSE-ORIGIN, CHANGE-REQUEST) on a primary/alt-port socket grid per
  address; advertises stun-5780 with >=2 same-family addrs, else
  stun-basic. Unmodified framing for tooling interop. Tested.
- tcpecho: JSON greeting with observed src + TCP_INFO MSS/options
  (Linux getsockopt; zeroed elsewhere via build tags), then byte echo.
- session: per-packet UDP observations + connect-back results, ByID lookup.
- control: GET /v1/sessions/{id}/observations, POST .../actions
  (delayed_echo → DELAYED_ECHO at the observed data-plane source;
  connect_back → dial the control-plane source, record connected/refused/
  timeout+rtt). Capabilities computed from what is actually wired.
- config/main: comma-separated STUN listeners; all planes bind explicit
  addresses; graceful shutdown of the new listeners.

Full flow smoke-tested; go test green (stun binding/change-port,
dataplane wire format).

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-07-31 19:53:36 +02:00
mrambossekandClaude Opus 5 507a8bfc1f build-status: production server v0.2.0 live on the fmr VM
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-07-31 19:42:56 +02:00
24 changed files with 2436 additions and 24 deletions
+7
View File
@@ -30,6 +30,13 @@ is the mutable one — update it as work lands.
Keep prober result IDs aligned with the measurement-schema test-type registry. Keep prober result IDs aligned with the measurement-schema test-type registry.
## Versioning
Prefer **patch** bumps (`server-v0.3.1`) for additive/incremental work; reserve **minor** bumps
for real milestones. Don't burn through minor versions. Tags are namespaced: `server-v*` for the
Go server, `v*` for the app. Pushing a `server-v*` tag runs CI → binaries + Gitea release +
registry image; the server on fmr can `--self-update` from those releases.
## Layout ## Layout
Monorepo. The prober is one deliverable; the Go server and the production app land as siblings. Monorepo. The prober is one deliverable; the Go server and the production app land as siblings.
+57
View File
@@ -214,3 +214,60 @@ Two collection-loop gotchas found while driving the phone over USB:
3. Start the Go server skeleton (enrollment + profile + sessions + UDP echo with observation 3. Start the Go server skeleton (enrollment + profile + sessions + UDP echo with observation
blocks + canary-DNS reference records) per probe-protocol.md. blocks + canary-DNS reference records) per probe-protocol.md.
4. Fold confirmed capabilities into the production `core-probe` / `core-shizuku` modules. 4. Fold confirmed capabilities into the production `core-probe` / `core-shizuku` modules.
## Production probe server — LIVE on dedicated VM "fmr" (2026-07-31)
`echolot-server v0.2.0` runs natively (systemd, no docker) on a dedicated VM: 2×IPv4 + 2×IPv6
service addresses (fmr-1/fmr-2.echo-lot.app, dual-stack DNS), a third IPv6 (`::2`) reserved for
SSH only — verified untouched by the daemon (explicit multi-address binds, no wildcard).
Control: fmr-1:8443 (SPKI pin `zRV9qkiLnRexAeh4RrSfJzbPWO+U/2Oj2/NVM/KfXlg=`, verified
externally over v4+v6). UDP data plane on all four service addresses :8442 — the second IP is
the stun-5780 substrate. Daily randomized self-update timer installed (checksum-verified
against SHA256SUMS; signature verification still TODO before treating the source as untrusted).
Host config in `/etc/echolot-server.env`. SSH access for sessions: `ssh claude-echolot`.
## Server v0.3.0 — STUN + TCP echo + observations + actions (2026-07-31)
Shipped and deployed to fmr via the server's own `--self-update` (first real exercise:
checksum-verified download v0.2.0→v0.3.0, atomic replace, restart — worked). Added over v0.2.0:
- **STUN** (RFC 5389 + 5780): 4 service addrs × primary/alt-port grid. Externally verified on
v4 AND v6 — binding success with XOR-MAPPED, RESPONSE-ORIGIN, OTHER-ADDRESS present, so the
profile now advertises **`stun-5780`** (the second IP earns its keep).
- **TCP echo** (:8441): JSON greeting with observed src + real Linux TCP_INFO — verified
externally `mss:1440` (v6, 150060), options `[sack,wscale]`, then byte-echo.
- **Observations API** `GET /v1/sessions/{id}/observations` (per-packet UDP view, connect-back
results, TCP records correlated by source IP).
- **Actions** `POST /v1/sessions/{id}/actions`: `delayed_echo` (DELAYED_ECHO at the observed
data-plane source — NAT-lifetime primitive) and `connect_back` (dials the control-plane
source, records connected/refused/timeout+rtt).
- Capabilities computed from what's actually wired: `udp-probe, delayed-echo, connect-back,
tcp-echo, stun-5780`.
Still not implemented: TLS-echo/JA4, HTTP echo, tls-reference, canary DNS (§6.1 reference
records), and the train/big-send/frag/throughput actions. Admin UI still token-mint + health only.
## Canary DNS live — server v0.3.1 on fmr (2026-07-31)
Zone `c.echo-lot.app` delegated (NS → fmr-1/fmr-2) and authoritative on all 4 service IPs
udp+tcp/53. Verified through full public recursion: `ttl-5` A→192.0.2.5 (Cloudflare), `ttl-3600`
AAAA→2001:db8::3600 (Google), `big-txt` TXT returned (TCP fallback, truncated over UDP as
designed). End-to-end session attribution works: a `<nonce>.<session-prefix>.c.echo-lot.app`
query resolved via a public resolver shows up in `GET /v1/sessions/{id}/observations` →
`dns_canary` with the resolver's real egress IP, transport, and EDNS. First real test already
caught a finding: **Google applies 0x20 case randomization** (mixed-case qname), Cloudflare does
not — captured via `case_preserved`. Capabilities now: udp-probe, delayed-echo, connect-back,
tcp-echo, stun-5780, canary-dns. Kept the hand-rolled stdlib DNS (no miekg/dns) — validated
against independent clients. Deployed via `--self-update` (v0.3.0→v0.3.1, checksum-verified).
## Server v0.3.2 + v0.3.3 (2026-07-31)
- **v0.3.2 — control-plane security (live on fmr, externally verified):** `POST /v1/echo`
reflects the received request head+body (b64) and observed TLS (version/cipher/SNI/ALPN) —
captured real SNI `fmr-1.echo-lot.app` and an injected header over public TLS1.3; `GET
/v1/tls-reference` returns the served DER chain + pin (cross-checked against the openssl-derived
pin). Optional cleartext echo listener (default off). Capability `http-echo`.
- **v0.3.3 — MTU probe (data plane):** MTU_PROBE (0x09) → small MTU_ACK (0x0A) carrying the
received datagram size; client DF-probes increasing sizes to find path MTU / black holes. ACK
is tiny → never amplifies. Tested.
- **Note on trains:** upstream trains (TRAIN_DATA 0x03) are already observable — every HMAC-valid
packet is recorded (seq/t_rx/size/type) with no per-packet response, so loss/reordering/inter-
arrival are visible via GET observations. The dedicated data-plane TRAIN_REPORT (0x05) is
deferred: §3.4 anti-amplification means it needs an asymmetric grant + columnar multi-datagram
encoding — a focused batch, not a corner to rush.
Remaining spec: tls-echo (ClientHello+JA4), TRAIN_REPORT, big/frag-send, throughput, downtrain;
real admin UI.
+31 -3
View File
@@ -36,15 +36,43 @@ All configurable via `ECHOLOT_*_LISTEN`. Plus:
2. **Second IP (optional):** full RFC 5780 NAT-behavior discovery (`stun-5780`) needs an 2. **Second IP (optional):** full RFC 5780 NAT-behavior discovery (`stun-5780`) needs an
alternate reply address; without it the profile advertises `stun-basic` and clients degrade alternate reply address; without it the profile advertises `stun-basic` and clients degrade
gracefully. gracefully.
3. **Delegated DNS subzone (optional, later):** the `canary-dns` capability needs port 53 on 3. **Delegated DNS subzone (for `canary-dns`):** set `ECHOLOT_DNS_LISTEN` (udp+tcp/53 on the
some IP + an NS delegation (mind systemd-resolved on 127.0.0.53). Absent → capability simply service IPs) and `ECHOLOT_CANARY_ZONE` (e.g. `c.echo-lot.app`), then delegate the zone to
not advertised. this host in your DNS provider:
```
c.echo-lot.app. NS fmr-1.echo-lot.app.
c.echo-lot.app. NS fmr-2.echo-lot.app.
```
The server is authoritative for that zone only, serving the spec §6.1 reference records
(frozen in `internal/canarydns/dns_reference.go`) plus per-query `<nonce>.<session>.<zone>`
lookups it logs. Binding :53 on the public IPs is fine even with systemd-resolved (it only
claims 127.0.0.53). Absent config → capability simply not advertised.
4. **Outbound freedom** for connect-back / delayed-echo actions — no extra inbound ports; 4. **Outbound freedom** for connect-back / delayed-echo actions — no extra inbound ports;
generated traffic goes only to the session's observed source. generated traffic goes only to the session's observed source.
Deliberately out of scope here: an echo listener on 443 (to detect port-based egress filtering) Deliberately out of scope here: an echo listener on 443 (to detect port-based egress filtering)
— that genuinely needs 443 and belongs on a dedicated IP, not on a host running a reverse proxy. — that genuinely needs 443 and belongs on a dedicated IP, not on a host running a reverse proxy.
## Host tuning (measurement fidelity)
A measurement server must not let the kernel distort what clients observe. Apply the
recommended sysctls and the daemon will confirm the host is clean:
```sh
sudo cp deploy/99-echolot-sysctl.conf /etc/sysctl.d/ && sudo sysctl --system
```
The daemon **self-tests at startup and via `GET /admin/selftest`** (localhost):
- **sysctl audit** — flags settings that would distort results (RA acceptance on a static host,
ICMP redirects, ICMP rate-limiting of the server's own errors, disabled TCP options).
- **egress-MTU self-proof** — DF-probes external anchors (`ECHOLOT_MTU_PROBE_TARGETS`,
default 1.1.1.1 + a v6 anchor) and reads the discovered path MTU. If the server's *own* uplink
can't carry 1500, client MTU results would measure this server, not the client — so the profile
exposes `server_selftest.mtu_ok` and the log warns loudly.
Both signals ride in `GET /v1/profile` as `server_selftest` so a client can trust — or skip —
MTU testing accordingly.
## Run in Docker (config via env) ## Run in Docker (config via env)
```sh ```sh
+152 -3
View File
@@ -18,6 +18,7 @@ import (
"crypto/tls" "crypto/tls"
"crypto/x509" "crypto/x509"
"crypto/x509/pkix" "crypto/x509/pkix"
"encoding/json"
"encoding/pem" "encoding/pem"
"errors" "errors"
"fmt" "fmt"
@@ -25,20 +26,26 @@ import (
"math/big" "math/big"
"net" "net"
"net/http" "net/http"
"net/netip"
"os" "os"
"os/signal" "os/signal"
"path/filepath" "path/filepath"
"strconv" "strconv"
"sync/atomic"
"syscall" "syscall"
"time" "time"
"echo-lot.app/server/internal/canarydns"
"echo-lot.app/server/internal/config" "echo-lot.app/server/internal/config"
"echo-lot.app/server/internal/control" "echo-lot.app/server/internal/control"
"echo-lot.app/server/internal/dataplane" "echo-lot.app/server/internal/dataplane"
"echo-lot.app/server/internal/selftest"
"echo-lot.app/server/internal/selfupdate" "echo-lot.app/server/internal/selfupdate"
"echo-lot.app/server/internal/session" "echo-lot.app/server/internal/session"
"echo-lot.app/server/internal/store" "echo-lot.app/server/internal/store"
"echo-lot.app/server/internal/stun"
"echo-lot.app/server/internal/system" "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. // Version is stamped via -ldflags "-X main.Version=v1.2.3" in CI.
@@ -95,9 +102,20 @@ func serve(cfg *config.Config) error {
slog.Info("control-plane certificate", "pin-sha256", pin) slog.Info("control-plane certificate", "pin-sha256", pin)
sessions := session.NewManager(15 * time.Minute) 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{ ctl := &control.Server{
Store: st, Sessions: sessions, Name: cfg.Name, Store: st, Sessions: sessions, Name: cfg.Name,
UDPPort: mustPort(firstAddr(cfg.UDPListen)), TCPPort: mustPort(firstAddr(cfg.TCPListen)), PinB64: pin, 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) ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
@@ -123,11 +141,42 @@ func serve(cfg *config.Config) error {
}(addr, ln) }(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/health (plain HTTP, localhost by default; spec §7)
admin := http.NewServeMux() admin := http.NewServeMux()
admin.HandleFunc("GET /healthz", func(w http.ResponseWriter, _ *http.Request) { admin.HandleFunc("GET /healthz", func(w http.ResponseWriter, _ *http.Request) {
fmt.Fprintf(w, `{"ok":true,"version":%q}`, Version) 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 // TODO(spec §7): enrollment token management + device list. Until the
// admin UI exists, mint tokens with: echolot-admin (or curl on this // admin UI exists, mint tokens with: echolot-admin (or curl on this
// listener once the endpoint lands). // listener once the endpoint lands).
@@ -156,14 +205,84 @@ func serve(cfg *config.Config) error {
return fmt.Errorf("udp listen %s: %w", addr, err) return fmt.Errorf("udp listen %s: %w", addr, err)
} }
udpConns = append(udpConns, conn) udpConns = append(udpConns, conn)
dp := &dataplane.Server{Sessions: sessions}
go func(a string, c *net.UDPConn) { go func(a string, c *net.UDPConn) {
errCh <- fmt.Errorf("udp %s: %w", a, dp.Serve(c)) errCh <- fmt.Errorf("udp %s: %w", a, dp.Serve(c))
}(addr, conn) }(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", slog.Info("listening",
"control", ctlAddrs, "admin", cfg.AdminListen, "udp", udpAddrs) "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 { select {
case <-ctx.Done(): case <-ctx.Done():
@@ -175,12 +294,42 @@ func serve(cfg *config.Config) error {
for _, c := range udpConns { for _, c := range udpConns {
_ = c.Close() _ = 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 return nil
case err := <-errCh: case err := <-errCh:
return err 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 { func firstAddr(spec string) string {
if a := config.Addrs(spec); len(a) > 0 { if a := config.Addrs(spec); len(a) > 0 {
return a[0] return a[0]
+40
View File
@@ -0,0 +1,40 @@
# SPDX-FileCopyrightText: 2026 Echolot contributors
# SPDX-License-Identifier: GPL-3.0-or-later
#
# Recommended sysctls for an Echolot probe-server host: keep the kernel from
# silently altering what clients measure. Install with:
# sudo cp 99-echolot-sysctl.conf /etc/sysctl.d/
# sudo sysctl --system
# The daemon audits these at startup and via GET /admin/selftest; anything not
# set here shows up as a "not measurement-clean" warning.
# Static-addressed host: never let a Router Advertisement mutate our routing.
# (Echolot's whole job is detecting broken RAs — the server must be immune.)
net.ipv6.conf.all.accept_ra = 0
net.ipv6.conf.default.accept_ra = 0
# Don't let ICMP redirects rewrite our routing mid-measurement, and don't
# emit redirects (we're an endpoint, not a router).
net.ipv4.conf.all.accept_redirects = 0
net.ipv4.conf.default.accept_redirects = 0
net.ipv6.conf.all.accept_redirects = 0
net.ipv4.conf.all.send_redirects = 0
net.ipv4.conf.default.send_redirects = 0
# Don't throttle the server's own ICMP errors (dest-unreachable/frag-needed/
# time-exceeded) — throttling produces false loss/black-hole readings when
# clients probe toward this server.
net.ipv4.icmp_ratelimit = 0
# These are usually already correct; pinned so the server can honestly
# negotiate/reflect them (a missing option in a client's evidence is then the
# path's fault, not ours).
net.ipv4.tcp_sack = 1
net.ipv4.tcp_timestamps = 1
net.ipv4.tcp_window_scaling = 1
net.ipv4.ip_no_pmtu_disc = 0
net.ipv4.icmp_echo_ignore_all = 0
# Loose reverse-path filtering suits a multi-IP measurement host (strict mode
# can drop alt-address / asymmetric replies used by STUN 5780).
net.ipv4.conf.all.rp_filter = 2
+287
View File
@@ -0,0 +1,287 @@
// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
package canarydns
import (
"encoding/binary"
"hash/fnv"
"net"
"net/netip"
"strings"
"sync"
"time"
)
// DNS constants (RFC 1035 + RFC 6891 EDNS).
const (
typeA = 1
typeNS = 2
typeTXT = 16
typeAAAA = 28
typeOPT = 41
classIN = 1
rcodeNoError = 0
rcodeNXDomain = 3
flagQR = 0x8000
flagAA = 0x0400
flagTC = 0x0200
flagRD = 0x0100
flagRA = 0x0080
udpMaxNoEDNS = 512
ednsDO = 0x8000 // DO bit lives in the OPT TTL field's high half
optECS = 8 // EDNS Client Subnet option code
)
// Query is one logged canary lookup (spec §6 dns_canary shape).
type Query struct {
QName string `json:"qname"`
At time.Time `json:"at"`
ResolverIP string `json:"resolver_ip"`
Transport string `json:"transport"` // "udp" | "tcp"
EDNS edns `json:"edns"`
ECS string `json:"ecs,omitempty"`
CasePreserved bool `json:"case_preserved"`
// qname_minimized is not reliably detectable authoritative-side without
// cross-query correlation; left false (TODO) rather than guessed.
QNameMinimized bool `json:"qname_minimized"`
}
type edns struct {
Present bool `json:"present"`
Bufsize int `json:"bufsize"`
Flags []string `json:"flags"`
}
// Server is the authoritative responder for one canary zone.
type Server struct {
zone string // fully-qualified, lowercase, trailing dot, e.g. "c.echo-lot.app."
nsName string // this server's own name for NS/authority answers
primaryV4 netip.Addr
primaryV6 netip.Addr
mu sync.Mutex
log []Query // ring, newest last
retainTo time.Time
}
const logCap = 8192
// New creates a server for zone (with or without trailing dot). nsName is the
// server's own hostname (for the zone's NS record); primary v4/v6 are this
// host's addresses used to answer the zone apex / NS glue.
func New(zone, nsName string, v4, v6 netip.Addr) *Server {
z := strings.ToLower(strings.TrimSuffix(zone, ".")) + "."
return &Server{zone: z, nsName: strings.TrimSuffix(nsName, ".") + ".", primaryV4: v4, primaryV6: v6}
}
// RecentForPrefix returns logged queries whose qname contains ".<prefix>."
// (the session prefix the app embeds: <nonce>.<session-prefix>.<zone>).
func (s *Server) RecentForPrefix(prefix string) []Query {
s.mu.Lock()
defer s.mu.Unlock()
needle := "." + strings.ToLower(prefix) + "."
var out []Query
for _, q := range s.log {
if strings.Contains(strings.ToLower(q.QName), needle) {
out = append(out, q)
}
}
return out
}
func (s *Server) record(q Query) {
s.mu.Lock()
defer s.mu.Unlock()
if len(s.log) >= logCap {
s.log = s.log[1:]
}
s.log = append(s.log, q)
}
// ServeUDP / ServeTCP run read loops; call one per bound address.
func (s *Server) ServeUDP(conn *net.UDPConn) error {
buf := make([]byte, 1500)
for {
n, raddr, err := conn.ReadFromUDPAddrPort(buf)
if err != nil {
return err
}
resp := s.handle(buf[:n], raddr.Addr(), "udp")
if resp != nil {
_, _ = conn.WriteToUDPAddrPort(resp, raddr)
}
}
}
func (s *Server) ServeTCP(ln net.Listener) error {
for {
c, err := ln.Accept()
if err != nil {
return err
}
go s.handleTCP(c)
}
}
func (s *Server) handleTCP(c net.Conn) {
defer c.Close()
_ = c.SetDeadline(time.Now().Add(10 * time.Second))
var lenBuf [2]byte
if _, err := readFull(c, lenBuf[:]); err != nil {
return
}
msg := make([]byte, binary.BigEndian.Uint16(lenBuf[:]))
if _, err := readFull(c, msg); err != nil {
return
}
ra, _ := netip.ParseAddrPort(c.RemoteAddr().String())
resp := s.handle(msg, ra.Addr(), "tcp")
if resp == nil {
return
}
// TCP has no 512 limit; never truncate.
out := make([]byte, 2+len(resp))
binary.BigEndian.PutUint16(out[0:2], uint16(len(resp)))
copy(out[2:], resp)
_, _ = c.Write(out)
}
func readFull(c net.Conn, b []byte) (int, error) {
got := 0
for got < len(b) {
n, err := c.Read(b[got:])
got += n
if err != nil {
return got, err
}
}
return got, nil
}
// handle parses one query, logs it, and returns the wire response (nil to drop).
func (s *Server) handle(pkt []byte, resolver netip.Addr, transport string) []byte {
if len(pkt) < 12 {
return nil
}
id := binary.BigEndian.Uint16(pkt[0:2])
qdcount := binary.BigEndian.Uint16(pkt[4:6])
arcount := binary.BigEndian.Uint16(pkt[10:12])
if qdcount != 1 {
return s.errorResponse(id, rcodeNoError, nil) // we only answer single-question queries
}
qnameRaw, qtype, _, qEnd, ok := parseQuestion(pkt, 12)
if !ok {
return nil
}
// EDNS OPT is an additional-section RR; scan for it after the question.
opt := parseOPT(pkt, qEnd, arcount)
// Log every query — this is the whole point of the canary zone.
q := Query{
QName: strings.TrimSuffix(qnameRaw, "."), At: time.Now().UTC(),
ResolverIP: resolver.Unmap().String(), Transport: transport,
EDNS: opt.edns,
ECS: opt.ecs,
CasePreserved: qnameRaw == strings.ToLower(qnameRaw), // mixed case ⇒ 0x20 randomization
}
s.record(q)
name := strings.ToLower(qnameRaw)
if !strings.HasSuffix(name, s.zone) {
return s.errorResponse(id, rcodeNXDomain, &opt)
}
sub := strings.TrimSuffix(name, s.zone) // e.g. "ttl-5." or "" for apex
return s.answer(id, pkt, qEnd, sub, qtype, &opt, transport)
}
// answer builds the response for a name known to be in-zone.
func (s *Server) answer(id uint16, pkt []byte, qEnd int, sub string, qtype uint16, opt *optInfo, transport string) []byte {
labels := splitLabels(sub) // e.g. ["ttl-5"], [], ["<nonce>","miss"], ["<nonce>","<sessprefix>"]
var rrs []rr
switch {
case len(labels) == 0: // zone apex
if qtype == typeNS {
rrs = append(rrs, rr{ttl: 3600, typ: typeNS, ns: s.nsName})
} else if qtype == typeA && s.primaryV4.IsValid() {
rrs = append(rrs, rr{ttl: 3600, typ: typeA, addr: s.primaryV4})
} else if qtype == typeAAAA && s.primaryV6.IsValid() {
rrs = append(rrs, rr{ttl: 3600, typ: typeAAAA, addr: s.primaryV6})
}
case len(labels) == 1:
if ref := findReference(labels[0]); ref != nil {
rrs = referenceAnswers(ref, qtype)
}
default:
// Per-query names: <nonce>.miss.<zone> and <nonce>.<session-prefix>.<zone>.
// Deterministic A derived from the leftmost label (the nonce), TTL 3600,
// documentation range — ground truth that can never be pre-cached.
if qtype == typeA {
rrs = append(rrs, rr{ttl: 3600, typ: typeA, addr: nonceAddr(labels[0])})
}
}
if len(rrs) == 0 {
// In-zone but no such record/type → NOERROR/NODATA (or NXDOMAIN at apex miss).
return s.buildResponse(id, pkt, qEnd, nil, opt, transport, rcodeNoError)
}
return s.buildResponse(id, pkt, qEnd, rrs, opt, transport, rcodeNoError)
}
// nonceAddr maps a nonce label into 192.0.2.0/24 deterministically.
func nonceAddr(nonce string) netip.Addr {
h := fnv.New32a()
_, _ = h.Write([]byte(nonce))
return netip.AddrFrom4([4]byte{192, 0, 2, byte(h.Sum32()%254 + 1)})
}
func referenceAnswers(ref *refRecord, qtype uint16) []rr {
var rrs []rr
switch qtype {
case typeA:
for _, a := range ref.a {
if a.Is4() {
rrs = append(rrs, rr{ttl: ref.ttl, typ: typeA, addr: a})
}
}
case typeAAAA:
for _, a := range ref.a {
if a.Is6() && !a.Is4In6() {
rrs = append(rrs, rr{ttl: ref.ttl, typ: typeAAAA, addr: a})
}
}
case typeTXT:
if len(ref.txt) > 0 {
rrs = append(rrs, rr{ttl: ref.ttl, typ: typeTXT, txt: ref.txt})
}
}
return rrs
}
func splitLabels(sub string) []string {
sub = strings.TrimSuffix(sub, ".")
if sub == "" {
return nil
}
return strings.Split(sub, ".")
}
func (s *Server) errorResponse(id uint16, rcode int, opt *optInfo) []byte {
hdr := make([]byte, 12)
binary.BigEndian.PutUint16(hdr[0:2], id)
binary.BigEndian.PutUint16(hdr[2:4], uint16(flagQR|flagAA|rcode))
if opt != nil && opt.edns.Present {
binary.BigEndian.PutUint16(hdr[10:12], 1)
return append(hdr, buildOPT(opt)...)
}
return hdr
}
+161
View File
@@ -0,0 +1,161 @@
// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
package canarydns
import (
"encoding/binary"
"net"
"net/netip"
"testing"
)
// buildQuery makes a single-question DNS query, optionally with an EDNS OPT.
func buildQuery(name string, qtype uint16, ednsBufsize int) []byte {
msg := make([]byte, 12)
binary.BigEndian.PutUint16(msg[0:2], 0x1234)
binary.BigEndian.PutUint16(msg[2:4], flagRD)
binary.BigEndian.PutUint16(msg[4:6], 1) // QDCOUNT
msg = append(msg, encodeName(name)...)
msg = binary.BigEndian.AppendUint16(msg, qtype)
msg = binary.BigEndian.AppendUint16(msg, classIN)
if ednsBufsize > 0 {
binary.BigEndian.PutUint16(msg[10:12], 1) // ARCOUNT
msg = append(msg, 0) // root name
msg = binary.BigEndian.AppendUint16(msg, typeOPT)
msg = binary.BigEndian.AppendUint16(msg, uint16(ednsBufsize))
msg = binary.BigEndian.AppendUint32(msg, 0)
msg = binary.BigEndian.AppendUint16(msg, 0)
}
return msg
}
// parseAnswers pulls (type, ttl, rdata) tuples from a response.
type ans struct {
typ uint16
ttl uint32
data []byte
}
func parseResponse(t *testing.T, resp []byte) (flags uint16, answers []ans) {
t.Helper()
flags = binary.BigEndian.Uint16(resp[2:4])
qd := binary.BigEndian.Uint16(resp[4:6])
an := binary.BigEndian.Uint16(resp[6:8])
off := 12
for i := uint16(0); i < qd; i++ {
_, next, ok := readName(resp, off)
if !ok {
t.Fatal("bad question name")
}
off = next + 4
}
for i := uint16(0); i < an; i++ {
_, next, ok := readName(resp, off)
if !ok {
t.Fatal("bad answer name")
}
typ := binary.BigEndian.Uint16(resp[next : next+2])
ttl := binary.BigEndian.Uint32(resp[next+4 : next+8])
rdlen := int(binary.BigEndian.Uint16(resp[next+8 : next+10]))
answers = append(answers, ans{typ, ttl, resp[next+10 : next+10+rdlen]})
off = next + 10 + rdlen
}
return
}
func newTestServer() *Server {
return New("c.echo-lot.app", "fmr", netip.MustParseAddr("192.0.2.1"), netip.MustParseAddr("2001:db8::1"))
}
func TestReferenceRecords(t *testing.T) {
s := newTestServer()
resolver := netip.MustParseAddr("198.51.100.7")
// ttl-5 A → 192.0.2.5, TTL 5
resp := s.handle(buildQuery("ttl-5.c.echo-lot.app", typeA, 0), resolver, "udp")
_, answers := parseResponse(t, resp)
if len(answers) != 1 || answers[0].ttl != 5 || !netip.AddrFrom4([4]byte(answers[0].data)).IsValid() {
t.Fatalf("ttl-5 A: %+v", answers)
}
if got := net.IP(answers[0].data).String(); got != "192.0.2.5" {
t.Fatalf("ttl-5 A = %s, want 192.0.2.5", got)
}
// many-rr → exactly 8 A records, in order .101..108
resp = s.handle(buildQuery("many-rr.c.echo-lot.app", typeA, 0), resolver, "udp")
_, answers = parseResponse(t, resp)
if len(answers) != 8 {
t.Fatalf("many-rr: got %d A records, want 8", len(answers))
}
for i, a := range answers {
if a.data[3] != byte(101+i) {
t.Fatalf("many-rr order: record %d = .%d, want .%d", i, a.data[3], 101+i)
}
}
}
func TestBigTxtTruncationVsEDNS(t *testing.T) {
s := newTestServer()
resolver := netip.MustParseAddr("198.51.100.7")
// No EDNS → 512 cap → TC set, answers dropped.
resp := s.handle(buildQuery("big-txt.c.echo-lot.app", typeTXT, 0), resolver, "udp")
flags, answers := parseResponse(t, resp)
if flags&flagTC == 0 {
t.Fatal("big-txt over plain UDP should set TC")
}
if len(answers) != 0 {
t.Fatalf("truncated response should carry no answers, got %d", len(answers))
}
// EDNS bufsize 4096 → full answer, no TC.
resp = s.handle(buildQuery("big-txt.c.echo-lot.app", typeTXT, 4096), resolver, "udp")
flags, answers = parseResponse(t, resp)
if flags&flagTC != 0 {
t.Fatal("big-txt with EDNS 4096 should not truncate")
}
if len(answers) != 1 {
t.Fatalf("want 1 TXT answer, got %d", len(answers))
}
// TCP → never truncates.
resp = s.handle(buildQuery("big-txt.c.echo-lot.app", typeTXT, 0), resolver, "tcp")
flags, _ = parseResponse(t, resp)
if flags&flagTC != 0 {
t.Fatal("TCP must not truncate")
}
}
func TestQueryLogAndPerPrefix(t *testing.T) {
s := newTestServer()
s.handle(buildQuery("abc123.SESSPREFIX1.c.echo-lot.app", typeA, 1232), netip.MustParseAddr("198.51.100.7"), "udp")
s.handle(buildQuery("def456.other.c.echo-lot.app", typeA, 0), netip.MustParseAddr("203.0.113.9"), "udp")
all := s.RecentForPrefix("sessprefix1")
if len(all) != 1 {
t.Fatalf("per-prefix filter: got %d, want 1", len(all))
}
q := all[0]
if q.ResolverIP != "198.51.100.7" || q.Transport != "udp" {
t.Fatalf("logged resolver/transport wrong: %+v", q)
}
if !q.EDNS.Present || q.EDNS.Bufsize != 1232 {
t.Fatalf("EDNS not captured: %+v", q.EDNS)
}
// nonce answer is deterministic + in doc range
resp := s.handle(buildQuery("abc123.SESSPREFIX1.c.echo-lot.app", typeA, 0), netip.MustParseAddr("198.51.100.7"), "udp")
_, answers := parseResponse(t, resp)
if len(answers) != 1 || answers[0].data[0] != 192 || answers[0].data[1] != 0 || answers[0].data[2] != 2 {
t.Fatalf("nonce answer not in 192.0.2.0/24: %+v", answers)
}
}
func TestOutOfZoneNXDomain(t *testing.T) {
s := newTestServer()
resp := s.handle(buildQuery("example.com", typeA, 0), netip.MustParseAddr("198.51.100.7"), "udp")
flags, _ := parseResponse(t, resp)
if flags&0x000F != rcodeNXDomain {
t.Fatalf("out-of-zone should be NXDOMAIN, flags=%#x", flags)
}
}
@@ -0,0 +1,78 @@
// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
// Package canarydns serves the authoritative canary zone (spec §6.1). The
// reference records below are FROZEN by the protocol spec — names, TTLs, and
// RDATA are ground truth the client compares against, so they must never
// change without a spec revision and a matching update in the app. All
// addresses are documentation-range (RFC 5737 192.0.2.0/24, RFC 3849
// 2001:db8::/32).
package canarydns
import "net/netip"
// refRecord is one frozen reference name (relative to the zone) with its
// per-type answers. A zero value in a field means "no record of that type".
type refRecord struct {
label string
ttl uint32
a []netip.Addr // A / AAAA answers (order preserved)
txt []string // one string per TXT record
}
// referenceRecords are the spec §6.1 fixed records. The RDATA constants are
// frozen HERE (the spec calls this file the source of truth) and mirrored in
// the app. Order within many-rr is part of the test (order/stripping check).
var referenceRecords = []refRecord{
{label: "ttl-5", ttl: 5,
a: []netip.Addr{netip.MustParseAddr("192.0.2.5"), netip.MustParseAddr("2001:db8::5")},
txt: []string{"echolot-ref ttl=5"}},
{label: "ttl-60", ttl: 60,
a: []netip.Addr{netip.MustParseAddr("192.0.2.60"), netip.MustParseAddr("2001:db8::60")},
txt: []string{"echolot-ref ttl=60"}},
{label: "ttl-3600", ttl: 3600,
a: []netip.Addr{netip.MustParseAddr("192.0.2.36"), netip.MustParseAddr("2001:db8::3600")},
txt: []string{"echolot-ref ttl=3600"}},
{label: "ttl-86400", ttl: 86400,
a: []netip.Addr{netip.MustParseAddr("192.0.2.86"), netip.MustParseAddr("2001:db8::8640")},
txt: []string{"echolot-ref ttl=86400"}},
// many-rr: exactly 8 A records in defined order.
{label: "many-rr", ttl: 300, a: []netip.Addr{
netip.MustParseAddr("192.0.2.101"), netip.MustParseAddr("192.0.2.102"),
netip.MustParseAddr("192.0.2.103"), netip.MustParseAddr("192.0.2.104"),
netip.MustParseAddr("192.0.2.105"), netip.MustParseAddr("192.0.2.106"),
netip.MustParseAddr("192.0.2.107"), netip.MustParseAddr("192.0.2.108"),
}},
// big-txt: ~1800 bytes, exercises EDNS bufsize / TCP fallback.
{label: "big-txt", ttl: 300, txt: bigTxt()},
}
// bigTxt builds a deterministic ~1800-byte TXT payload as a sequence of
// 255-byte character-strings (the DNS TXT chunk limit). The content is fixed
// so the client can verify integrity, not just length.
func bigTxt() []string {
const total = 1800
const pattern = "echolot-big-txt-reference-0123456789abcdef-"
buf := make([]byte, 0, total)
for len(buf) < total {
buf = append(buf, pattern...)
}
buf = buf[:total]
var out []string
for len(buf) > 0 {
n := min(255, len(buf))
out = append(out, string(buf[:n]))
buf = buf[n:]
}
return out
}
// findReference returns the reference record for a label, or nil.
func findReference(label string) *refRecord {
for i := range referenceRecords {
if referenceRecords[i].label == label {
return &referenceRecords[i]
}
}
return nil
}
+266
View File
@@ -0,0 +1,266 @@
// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
package canarydns
import (
"encoding/binary"
"net/netip"
"strings"
)
// rr is a resource record to encode into the answer section.
type rr struct {
ttl uint32
typ uint16
addr netip.Addr // for A/AAAA
txt []string // for TXT
ns string // for NS
}
// optInfo is the parsed EDNS OPT plus the derived observation fields.
type optInfo struct {
edns edns
ecs string
}
// parseQuestion reads a single question starting at off. Returns the raw
// (case-preserved) qname with trailing dot, qtype, qclass, and the offset
// just past the question.
func parseQuestion(pkt []byte, off int) (qname string, qtype, qclass uint16, end int, ok bool) {
name, next, ok := readName(pkt, off)
if !ok || next+4 > len(pkt) {
return "", 0, 0, 0, false
}
qtype = binary.BigEndian.Uint16(pkt[next : next+2])
qclass = binary.BigEndian.Uint16(pkt[next+2 : next+4])
return name, qtype, qclass, next + 4, true
}
// readName decodes a DNS name (with compression pointers) into a
// dot-terminated string, preserving label case.
func readName(pkt []byte, off int) (string, int, bool) {
var sb strings.Builder
end := -1
jumps := 0
for {
if off >= len(pkt) {
return "", 0, false
}
l := int(pkt[off])
switch {
case l == 0:
off++
if end < 0 {
end = off
}
if sb.Len() == 0 {
return ".", end, true
}
return sb.String(), end, true
case l&0xC0 == 0xC0: // compression pointer
if off+1 >= len(pkt) {
return "", 0, false
}
if end < 0 {
end = off + 2
}
off = int(binary.BigEndian.Uint16(pkt[off:off+2]) & 0x3FFF)
jumps++
if jumps > 16 {
return "", 0, false
}
default:
if off+1+l > len(pkt) {
return "", 0, false
}
sb.Write(pkt[off+1 : off+1+l])
sb.WriteByte('.')
off += 1 + l
}
}
}
// parseOPT scans the additional section for an EDNS OPT RR and extracts
// bufsize, the DO flag, and any ECS option.
func parseOPT(pkt []byte, off int, arcount uint16) optInfo {
var info optInfo
for i := uint16(0); i < arcount && off < len(pkt); i++ {
_, next, ok := readName(pkt, off)
if !ok || next+10 > len(pkt) {
return info
}
typ := binary.BigEndian.Uint16(pkt[next : next+2])
class := binary.BigEndian.Uint16(pkt[next+2 : next+4]) // OPT: requester bufsize
ttl := binary.BigEndian.Uint32(pkt[next+4 : next+8]) // OPT: extended-rcode/version/flags
rdlen := int(binary.BigEndian.Uint16(pkt[next+8 : next+10]))
rdata := next + 10
if rdata+rdlen > len(pkt) {
return info
}
if typ == typeOPT {
info.edns.Present = true
info.edns.Bufsize = int(class)
if ttl&ednsDO != 0 {
info.edns.Flags = append(info.edns.Flags, "do")
}
info.ecs = parseECS(pkt[rdata : rdata+rdlen])
return info
}
off = rdata + rdlen
}
return info
}
// parseECS extracts an EDNS Client Subnet option (RFC 7871) as "ip/scope".
func parseECS(rdata []byte) string {
for len(rdata) >= 4 {
code := binary.BigEndian.Uint16(rdata[0:2])
olen := int(binary.BigEndian.Uint16(rdata[2:4]))
if 4+olen > len(rdata) {
return ""
}
if code == optECS && olen >= 4 {
fam := binary.BigEndian.Uint16(rdata[4:6])
srcPrefix := rdata[6]
addrBytes := rdata[8 : 4+olen]
var ip netip.Addr
if fam == 1 {
var b [4]byte
copy(b[:], addrBytes)
ip = netip.AddrFrom4(b)
} else if fam == 2 {
var b [16]byte
copy(b[:], addrBytes)
ip = netip.AddrFrom16(b)
}
if ip.IsValid() {
return ip.String() + "/" + itoa(int(srcPrefix))
}
}
rdata = rdata[4+olen:]
}
return ""
}
func itoa(n int) string {
if n == 0 {
return "0"
}
var b [4]byte
i := len(b)
for n > 0 {
i--
b[i] = byte('0' + n%10)
n /= 10
}
return string(b[i:])
}
// buildResponse assembles the answer, sets TC when a UDP response exceeds the
// negotiated buffer, and appends the OPT RR when the query used EDNS.
func (s *Server) buildResponse(id uint16, pkt []byte, qEnd int, answers []rr, opt *optInfo, transport string, rcode int) []byte {
msg := make([]byte, 12)
binary.BigEndian.PutUint16(msg[0:2], id)
// question is copied verbatim (case preserved) from the query
msg = append(msg, pkt[12:qEnd]...)
body := make([]byte, 0, 512)
for _, a := range answers {
body = append(body, encodeRR(a)...)
}
extra := 0
if opt != nil && opt.edns.Present {
extra = 1
}
flags := uint16(flagQR|flagAA) | (binary.BigEndian.Uint16(pkt[2:4]) & flagRD) | uint16(rcode)
if opt != nil && opt.edns.Present {
body = append(body, buildOPT(opt)...)
}
// UDP truncation: without EDNS the limit is 512; with EDNS it's the
// requester's bufsize (floored at 512). Drop the answer section and set TC.
if transport == "udp" {
limit := udpMaxNoEDNS
if opt != nil && opt.edns.Present && opt.edns.Bufsize > udpMaxNoEDNS {
limit = opt.edns.Bufsize
}
if 12+(qEnd-12)+len(body) > limit {
flags |= flagTC
// Keep only the OPT RR (if any); drop answers.
body = body[:0]
if opt != nil && opt.edns.Present {
body = append(body, buildOPT(opt)...)
answers = nil
} else {
answers = nil
}
}
}
binary.BigEndian.PutUint16(msg[2:4], uint16(flags))
binary.BigEndian.PutUint16(msg[4:6], 1) // QDCOUNT
binary.BigEndian.PutUint16(msg[6:8], uint16(len(answers)))
binary.BigEndian.PutUint16(msg[10:12], uint16(extra))
return append(msg, body...)
}
// encodeRR encodes one answer RR, using a compression pointer (0xC00C) to the
// question name at offset 12.
func encodeRR(a rr) []byte {
var rdata []byte
switch a.typ {
case typeA:
b := a.addr.As4()
rdata = b[:]
case typeAAAA:
b := a.addr.As16()
rdata = b[:]
case typeTXT:
for _, s := range a.txt {
for len(s) > 0 {
n := len(s)
if n > 255 {
n = 255
}
rdata = append(rdata, byte(n))
rdata = append(rdata, s[:n]...)
s = s[n:]
}
}
case typeNS:
rdata = encodeName(a.ns)
}
out := make([]byte, 0, 12+len(rdata))
out = append(out, 0xC0, 0x0C) // name → pointer to question
out = binary.BigEndian.AppendUint16(out, a.typ)
out = binary.BigEndian.AppendUint16(out, classIN)
out = binary.BigEndian.AppendUint32(out, a.ttl)
out = binary.BigEndian.AppendUint16(out, uint16(len(rdata)))
return append(out, rdata...)
}
func encodeName(name string) []byte {
var out []byte
for _, label := range strings.Split(strings.TrimSuffix(name, "."), ".") {
if label == "" {
continue
}
out = append(out, byte(len(label)))
out = append(out, label...)
}
return append(out, 0)
}
// buildOPT emits a minimal EDNS OPT RR echoing our own bufsize (advertise a
// generous 4096) with DO cleared — we serve no DNSSEC.
func buildOPT(*optInfo) []byte {
out := []byte{0} // root name
out = binary.BigEndian.AppendUint16(out, typeOPT)
out = binary.BigEndian.AppendUint16(out, 4096) // our bufsize
out = binary.BigEndian.AppendUint32(out, 0) // ext-rcode/version/flags
out = binary.BigEndian.AppendUint16(out, 0) // rdlen
return out
}
+13
View File
@@ -26,6 +26,14 @@ type Config struct {
// Data plane // Data plane
UDPListen string // ECHOLOT_UDP_LISTEN / --udp-listen (spec default port 8442) UDPListen string // ECHOLOT_UDP_LISTEN / --udp-listen (spec default port 8442)
TCPListen string // ECHOLOT_TCP_LISTEN / --tcp-listen (spec default port 8441) TCPListen string // ECHOLOT_TCP_LISTEN / --tcp-listen (spec default port 8441)
StunListen string // ECHOLOT_STUN_LISTEN / --stun-listen (spec default 3478; empty disables)
DNSListen string // ECHOLOT_DNS_LISTEN / --dns-listen (canary zone; empty disables)
CanaryZone string // ECHOLOT_CANARY_ZONE / --canary-zone (e.g. c.echo-lot.app)
// Optional cleartext HTTP-echo listener (spec §4 plaintext-path test).
// Default empty = off; it exposes only POST /v1/echo, no auth, no secrets.
HTTPEchoListen string // ECHOLOT_HTTP_ECHO_LISTEN / --http-echo-listen
// Comma-separated anchors for the egress-MTU self-proof (host or ip).
MTUProbeTargets string // ECHOLOT_MTU_PROBE_TARGETS / --mtu-probe-targets
// Admin UI / health listener (spec §7: localhost-only by default) // Admin UI / health listener (spec §7: localhost-only by default)
AdminListen string // ECHOLOT_ADMIN_LISTEN / --admin-listen AdminListen string // ECHOLOT_ADMIN_LISTEN / --admin-listen
@@ -65,6 +73,11 @@ func Load(args []string) (*Config, *Actions, error) {
fs.StringVar(&c.TLSKey, "tls-key", envOr("TLS_KEY", ""), "TLS key path (empty: self-signed in state dir)") fs.StringVar(&c.TLSKey, "tls-key", envOr("TLS_KEY", ""), "TLS key path (empty: self-signed in state dir)")
fs.StringVar(&c.UDPListen, "udp-listen", envOr("UDP_LISTEN", ":8442"), "UDP data-plane listen address(es), comma-separated") fs.StringVar(&c.UDPListen, "udp-listen", envOr("UDP_LISTEN", ":8442"), "UDP data-plane listen address(es), comma-separated")
fs.StringVar(&c.TCPListen, "tcp-listen", envOr("TCP_LISTEN", ":8441"), "TCP echo listen address(es), comma-separated") fs.StringVar(&c.TCPListen, "tcp-listen", envOr("TCP_LISTEN", ":8441"), "TCP echo listen address(es), comma-separated")
fs.StringVar(&c.StunListen, "stun-listen", envOr("STUN_LISTEN", ":3478"), "STUN listen address(es), comma-separated; empty disables (spec §4)")
fs.StringVar(&c.DNSListen, "dns-listen", envOr("DNS_LISTEN", ""), "canary-DNS listen address(es) udp+tcp/53, comma-separated; empty disables (spec §6.1)")
fs.StringVar(&c.CanaryZone, "canary-zone", envOr("CANARY_ZONE", ""), "authoritative canary zone, e.g. c.echo-lot.app")
fs.StringVar(&c.HTTPEchoListen, "http-echo-listen", envOr("HTTP_ECHO_LISTEN", ""), "optional CLEARTEXT http-echo listen address(es); empty disables (spec §4)")
fs.StringVar(&c.MTUProbeTargets, "mtu-probe-targets", envOr("MTU_PROBE_TARGETS", "1.1.1.1,2606:4700:4700::1111"), "egress-MTU self-proof anchors, comma-separated")
fs.StringVar(&c.AdminListen, "admin-listen", envOr("ADMIN_LISTEN", "127.0.0.1:8444"), "admin/health listen address (keep localhost)") fs.StringVar(&c.AdminListen, "admin-listen", envOr("ADMIN_LISTEN", "127.0.0.1:8444"), "admin/health listen address (keep localhost)")
fs.StringVar(&c.StateDir, "state-dir", envOr("STATE_DIR", defaultStateDir()), "state directory (device store, generated TLS)") fs.StringVar(&c.StateDir, "state-dir", envOr("STATE_DIR", defaultStateDir()), "state directory (device store, generated TLS)")
fs.StringVar(&c.Name, "name", envOr("NAME", "echolot"), "server profile name") fs.StringVar(&c.Name, "name", envOr("NAME", "echolot"), "server profile name")
+175 -4
View File
@@ -7,14 +7,19 @@
package control package control
import ( import (
"crypto/rand"
"crypto/sha256" "crypto/sha256"
"crypto/tls" "crypto/tls"
"crypto/x509" "crypto/x509"
"encoding/base64" "encoding/base64"
"encoding/hex"
"encoding/json" "encoding/json"
"errors"
"log/slog" "log/slog"
"net"
"net/http" "net/http"
"net/netip" "net/netip"
"strconv"
"strings" "strings"
"time" "time"
@@ -29,12 +34,30 @@ type Server struct {
Store *store.Store Store *store.Store
Sessions *session.Manager Sessions *session.Manager
Name string Name string
// Targets/capabilities for the profile response. The skeleton offers only // Targets/capabilities for the profile response. The server offers only
// what it actually implements; the registry grows with the code. // what it actually implements; the registry grows with the code.
UDPPort int UDPPort int
TCPPort int TCPPort int
StunPort int
// SPKI pin of the serving cert, for the profile's pins[] field. // SPKI pin of the serving cert, for the profile's pins[] field.
PinB64 string PinB64 string
// CertChain is the served leaf-first DER chain, for GET /v1/tls-reference.
CertChain [][]byte
// Capabilities as computed at startup from what is actually wired up.
Capabilities []string
// TCPRecent returns recent TCP-echo connections for a source IP (may be nil).
TCPRecent func(ip string) any
// DelayedEcho schedules/sends a DELAYED_ECHO for a session (may be nil).
DelayedEcho func(sess *session.Session, actionID string) error
// CanaryQueries returns logged canary lookups for a session prefix (may be nil).
CanaryQueries func(sessionPrefix string) any
// CanaryZone is surfaced in the profile so the app knows what to query.
CanaryZone string
// ProvenGood reports the server's self-test signal (may be nil). Surfaced
// in the profile so a client can trust — or skip — MTU tests: if the
// server's own egress isn't full-MTU, client MTU results measure the
// server, not the client.
ProvenGood func() (mtuOK, sysctlOK bool)
} }
func (s *Server) Handler() http.Handler { func (s *Server) Handler() http.Handler {
@@ -43,11 +66,156 @@ func (s *Server) Handler() http.Handler {
mux.HandleFunc("GET /v1/profile", s.profile) mux.HandleFunc("GET /v1/profile", s.profile)
mux.HandleFunc("POST /v1/sessions", s.newSession) mux.HandleFunc("POST /v1/sessions", s.newSession)
mux.HandleFunc("DELETE /v1/sessions/{id}", s.deleteSession) mux.HandleFunc("DELETE /v1/sessions/{id}", s.deleteSession)
// TODO(spec §5, §6): /v1/sessions/{id}/actions, /v1/sessions/{id}/observations mux.HandleFunc("GET /v1/sessions/{id}/observations", s.observations)
// TODO(spec §4): POST /v1/echo, GET /v1/tls-reference mux.HandleFunc("POST /v1/sessions/{id}/actions", s.actions)
mux.HandleFunc("POST /v1/echo", s.httpEcho)
mux.HandleFunc("GET /v1/tls-reference", s.tlsReference)
// TODO(spec §4): TLS-echo/JA4 (tls-echo capability, needs ClientHello capture)
// TODO(spec §5): downtrain, big_send, frag_send, throughput
return mux return mux
} }
// EchoHandler exposes just the HTTP-echo endpoint for the optional cleartext
// listener (spec §4: plaintext-path tampering test).
func (s *Server) EchoHandler() http.Handler {
mux := http.NewServeMux()
mux.HandleFunc("POST /v1/echo", s.httpEcho)
return mux
}
// selftestSignal is the compact "server proven good" object for the profile.
// mtu_ok=false tells a client its MTU results would measure this server.
func selftestSignal(f func() (bool, bool)) map[string]any {
if f == nil {
return map[string]any{"mtu_ok": nil, "sysctl_ok": nil}
}
mtuOK, sysctlOK := f()
return map[string]any{"mtu_ok": mtuOK, "sysctl_ok": sysctlOK}
}
// sessionAuth resolves {id} and requires the bearer to be the owning device.
func (s *Server) sessionAuth(w http.ResponseWriter, r *http.Request) *session.Session {
dev := s.Store.DeviceByCredential(bearer(r))
if dev == nil {
writeJSON(w, http.StatusUnauthorized, map[string]string{"error": "unknown credential"})
return nil
}
sess := s.Sessions.ByID(r.PathValue("id"))
if sess == nil || sess.Device != dev.ID {
writeJSON(w, http.StatusNotFound, map[string]string{"error": "no such session"})
return nil
}
return sess
}
// observations is spec §6 — everything the server witnessed for a session.
func (s *Server) observations(w http.ResponseWriter, r *http.Request) {
sess := s.sessionAuth(w, r)
if sess == nil {
return
}
packetsSeen, udp, cb := sess.Observations()
var tcp any
if s.TCPRecent != nil {
// Correlate by source IP: TCP echo carries no session id on the wire.
if ds := sess.DataSource(); ds.IsValid() {
tcp = s.TCPRecent(ds.Addr().Unmap().String())
} else if sess.ControlSource.IsValid() {
tcp = s.TCPRecent(sess.ControlSource.Unmap().String())
}
}
var dnsCanary any
if s.CanaryQueries != nil {
dnsCanary = s.CanaryQueries(sess.ID[:16]) // the session's wire prefix
}
writeJSON(w, http.StatusOK, map[string]any{
"udp": map[string]any{"packets_seen": packetsSeen, "packets": udp},
"tcp": tcp,
"connect_back": cb,
"dns_canary": dnsCanary,
// TODO(spec §6): http echo records
})
}
// actions is spec §5 — authenticated asymmetric operations. Implemented:
// delayed_echo, connect_back. Destination is ALWAYS the session's observed
// source (data-plane source; connect_back uses the control-plane source).
func (s *Server) actions(w http.ResponseWriter, r *http.Request) {
sess := s.sessionAuth(w, r)
if sess == nil {
return
}
var req struct {
Action string `json:"action"`
DelayS int `json:"delay_s"`
Protocol string `json:"protocol"`
Port int `json:"port"`
}
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "bad body"})
return
}
actionID := randomID()
switch req.Action {
case "delayed_echo":
if s.DelayedEcho == nil {
writeJSON(w, http.StatusNotImplemented, map[string]string{"error": "delayed_echo not wired"})
return
}
delay := min(max(req.DelayS, 1), 600)
if !sess.DataSource().IsValid() {
writeJSON(w, http.StatusConflict, map[string]string{"error": "no data-plane traffic seen yet — send an ECHO first"})
return
}
time.AfterFunc(time.Duration(delay)*time.Second, func() {
if err := s.DelayedEcho(sess, actionID); err != nil {
slog.Debug("delayed echo failed", "err", err)
}
})
writeJSON(w, http.StatusAccepted, map[string]any{"action_id": actionID, "delay_s": delay})
case "connect_back":
if req.Port < 1 || req.Port > 65535 || (req.Protocol != "tcp" && req.Protocol != "udp") {
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "connect_back needs protocol tcp|udp and a port"})
return
}
target := net.JoinHostPort(sess.ControlSource.Unmap().String(), strconv.Itoa(req.Port))
go func() {
start := time.Now()
conn, err := net.DialTimeout(req.Protocol, target, 5*time.Second)
res := session.ConnectBackResult{ActionID: actionID, RttMs: float64(time.Since(start).Microseconds()) / 1000}
switch {
case err == nil:
res.Result = "connected"
if req.Protocol == "udp" {
// UDP "dial" always succeeds locally; send one datagram
// so the client actually observes something.
_, _ = conn.Write([]byte("echolot-connect-back " + actionID))
}
conn.Close()
case isTimeout(err):
res.Result = "timeout"
default:
res.Result = "refused"
}
sess.RecordConnectBack(res)
}()
writeJSON(w, http.StatusAccepted, map[string]any{"action_id": actionID, "target": target})
default:
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "unknown or unimplemented action"})
}
}
func isTimeout(err error) bool {
var ne net.Error
return errors.As(err, &ne) && ne.Timeout()
}
func randomID() string {
var b [8]byte
_, _ = rand.Read(b[:])
return hex.EncodeToString(b[:])
}
func bearer(r *http.Request) string { func bearer(r *http.Request) string {
h := r.Header.Get("Authorization") h := r.Header.Get("Authorization")
if v, ok := strings.CutPrefix(h, "Bearer "); ok { if v, ok := strings.CutPrefix(h, "Bearer "); ok {
@@ -103,15 +271,18 @@ func (s *Server) profile(w http.ResponseWriter, r *http.Request) {
// source_url makes GPL §6 compliance mechanical for operators of // source_url makes GPL §6 compliance mechanical for operators of
// modified builds and gives clients provenance for the measurement. // modified builds and gives clients provenance for the measurement.
"source_url": "", // TODO: stamp from build metadata "source_url": "", // TODO: stamp from build metadata
"capabilities": []string{"udp-probe"}, "capabilities": s.Capabilities,
"targets": []map[string]any{{ "targets": []map[string]any{{
"id": s.Name, "id": s.Name,
"ip4": host, // TODO: explicit configured addresses, v6, second STUN addr "ip4": host, // TODO: explicit configured addresses, v6, second STUN addr
"udp_port": s.UDPPort, "udp_port": s.UDPPort,
"tcp_port": s.TCPPort, "tcp_port": s.TCPPort,
"stun_port": s.StunPort,
}}, }},
"pins": []string{"pin-sha256:" + s.PinB64}, "pins": []string{"pin-sha256:" + s.PinB64},
"next_pins": []string{}, "next_pins": []string{},
"canary_zone": s.CanaryZone,
"server_selftest": selftestSignal(s.ProvenGood),
"limits": map[string]any{"max_kbps": 50000, "max_session_s": 900}, "limits": map[string]any{"max_kbps": 50000, "max_session_s": 900},
}) })
} }
+94
View File
@@ -0,0 +1,94 @@
// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
package control
import (
"crypto/tls"
"encoding/base64"
"io"
"net/http"
"strings"
)
// httpEcho implements spec §4 HTTP echo: return the exact received request
// (request line + headers + body, base64) plus the TLS parameters the server
// observed. The client diffs this against what it sent to detect header
// injection/stripping, transparent proxying, or TLS interception
// (sec.http_echo). Served on the control HTTPS listener and, optionally, on a
// cleartext listener to test plaintext-path tampering.
func (s *Server) httpEcho(w http.ResponseWriter, r *http.Request) {
body, _ := io.ReadAll(io.LimitReader(r.Body, 1<<20))
// Reconstruct the received request head verbatim (as close as net/http
// exposes it — header order is lost, but names/values and the request
// line survive, which is what tampering changes).
var head strings.Builder
head.WriteString(r.Method + " " + r.RequestURI + " " + r.Proto + "\r\n")
head.WriteString("Host: " + r.Host + "\r\n")
for name, vals := range r.Header {
for _, v := range vals {
head.WriteString(name + ": " + v + "\r\n")
}
}
head.WriteString("\r\n")
resp := map[string]any{
"observed_src": r.RemoteAddr,
"request_head_b64": base64.StdEncoding.EncodeToString([]byte(head.String())),
"body_b64": base64.StdEncoding.EncodeToString(body),
"body_len": len(body),
"scheme": schemeOf(r),
}
if r.TLS != nil {
resp["tls"] = tlsParams(r.TLS)
}
writeJSON(w, http.StatusOK, resp)
}
func schemeOf(r *http.Request) string {
if r.TLS != nil {
return "https"
}
return "http"
}
func tlsParams(cs *tls.ConnectionState) map[string]any {
return map[string]any{
"version": tlsVersionName(cs.Version),
"cipher": tls.CipherSuiteName(cs.CipherSuite),
"sni": cs.ServerName,
"alpn": cs.NegotiatedProtocol,
"resumed": cs.DidResume,
}
}
func tlsVersionName(v uint16) string {
switch v {
case tls.VersionTLS13:
return "TLS1.3"
case tls.VersionTLS12:
return "TLS1.2"
case tls.VersionTLS11:
return "TLS1.1"
case tls.VersionTLS10:
return "TLS1.0"
}
return "unknown"
}
// tlsReference implements spec §4: return the exact certificate chain this
// server serves (DER, base64), so the app can compare it against a copy it
// obtained out-of-band and against what its own direct handshake yielded
// (sec.tls_reference). Not a capability — always available on the control
// plane. No auth: the chain is public information a handshake already reveals.
func (s *Server) tlsReference(w http.ResponseWriter, r *http.Request) {
chain := make([]string, 0, len(s.CertChain))
for _, der := range s.CertChain {
chain = append(chain, base64.StdEncoding.EncodeToString(der))
}
writeJSON(w, http.StatusOK, map[string]any{
"pin_sha256": s.PinB64,
"chain_der": chain, // leaf first, as served
})
}
+63
View File
@@ -0,0 +1,63 @@
// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
package control
import (
"encoding/base64"
"encoding/json"
"net/http/httptest"
"strings"
"testing"
)
func TestHTTPEchoReflectsRequest(t *testing.T) {
s := &Server{}
req := httptest.NewRequest("POST", "/v1/echo", strings.NewReader("payload-bytes"))
req.Header.Set("X-Injected", "canary")
rr := httptest.NewRecorder()
s.httpEcho(rr, req)
var resp struct {
RequestHeadB64 string `json:"request_head_b64"`
BodyB64 string `json:"body_b64"`
BodyLen int `json:"body_len"`
Scheme string `json:"scheme"`
}
if err := json.Unmarshal(rr.Body.Bytes(), &resp); err != nil {
t.Fatal(err)
}
head, _ := base64.StdEncoding.DecodeString(resp.RequestHeadB64)
if !strings.Contains(string(head), "X-Injected: canary") {
t.Fatalf("echo did not reflect the injected header:\n%s", head)
}
body, _ := base64.StdEncoding.DecodeString(resp.BodyB64)
if string(body) != "payload-bytes" || resp.BodyLen != 13 {
t.Fatalf("body mismatch: %q len=%d", body, resp.BodyLen)
}
if resp.Scheme != "http" { // httptest requests carry no TLS
t.Fatalf("scheme = %s, want http", resp.Scheme)
}
}
func TestTLSReferenceReturnsChain(t *testing.T) {
s := &Server{PinB64: "TESTPIN", CertChain: [][]byte{{0x30, 0x82, 0x01}, {0xAA, 0xBB}}}
rr := httptest.NewRecorder()
s.tlsReference(rr, httptest.NewRequest("GET", "/v1/tls-reference", nil))
var resp struct {
PinSHA256 string `json:"pin_sha256"`
ChainDER []string `json:"chain_der"`
}
if err := json.Unmarshal(rr.Body.Bytes(), &resp); err != nil {
t.Fatal(err)
}
if resp.PinSHA256 != "TESTPIN" || len(resp.ChainDER) != 2 {
t.Fatalf("bad tls-reference: %+v", resp)
}
first, _ := base64.StdEncoding.DecodeString(resp.ChainDER[0])
if len(first) != 3 || first[0] != 0x30 {
t.Fatalf("leaf DER not round-tripped: %x", first)
}
}
+63 -2
View File
@@ -11,9 +11,11 @@ import (
"crypto/hmac" "crypto/hmac"
"crypto/sha256" "crypto/sha256"
"encoding/binary" "encoding/binary"
"fmt"
"log/slog" "log/slog"
"net" "net"
"net/netip" "net/netip"
"sync"
"time" "time"
"echo-lot.app/server/internal/session" "echo-lot.app/server/internal/session"
@@ -27,6 +29,9 @@ const (
TypeEchoResp = 0x02 TypeEchoResp = 0x02
TypeTimesyncReq = 0x07 TypeTimesyncReq = 0x07
TypeTimesyncRsp = 0x08 TypeTimesyncRsp = 0x08
TypeMtuProbe = 0x09
TypeMtuAck = 0x0A
TypeDelayedEcho = 0x0B
) )
type Server struct { type Server struct {
@@ -34,13 +39,22 @@ type Server struct {
// Epoch for server-side t_rx/t_tx: process start; observation consumers // Epoch for server-side t_rx/t_tx: process start; observation consumers
// only need differences plus the timesync exchange, not absolute time. // only need differences plus the timesync exchange, not absolute time.
start time.Time start time.Time
mu sync.Mutex
conns []*net.UDPConn
} }
// Serve runs the read loop for one socket; call once per bound address.
// The socket is retained so actions (delayed echo) can pick a family-matching
// sender later.
func (s *Server) Serve(conn *net.UDPConn) error { func (s *Server) Serve(conn *net.UDPConn) error {
s.mu.Lock()
if s.start.IsZero() {
s.start = time.Now() s.start = time.Now()
}
s.conns = append(s.conns, conn)
s.mu.Unlock()
buf := make([]byte, 65535) buf := make([]byte, 65535)
oob := make([]byte, 0)
_ = oob // TODO: recvmsg w/ IP_RECVTOS+IP_RECVTTL via golang.org/x/net for TTL/DSCP/ECN observation
for { for {
n, raddr, err := conn.ReadFromUDPAddrPort(buf) n, raddr, err := conn.ReadFromUDPAddrPort(buf)
if err != nil { if err != nil {
@@ -51,6 +65,36 @@ func (s *Server) Serve(conn *net.UDPConn) error {
} }
} }
// connFor picks a retained socket whose family matches the target.
func (s *Server) connFor(target netip.AddrPort) *net.UDPConn {
s.mu.Lock()
defer s.mu.Unlock()
want4 := target.Addr().Unmap().Is4()
for _, c := range s.conns {
la := c.LocalAddr().(*net.UDPAddr).AddrPort()
if la.Addr().Unmap().Is4() == want4 {
return c
}
}
return nil
}
// SendDelayedEcho fires one DELAYED_ECHO packet at the session's observed
// data-plane source (spec §5: the NAT-mapping-lifetime primitive). The
// payload carries the action id for correlation.
func (s *Server) SendDelayedEcho(sess *session.Session, actionID string) error {
target := sess.DataSource()
if !target.IsValid() {
return fmt.Errorf("session has no observed data-plane source yet")
}
conn := s.connFor(target)
if conn == nil {
return fmt.Errorf("no data-plane socket matches target family")
}
s.send(conn, target, sess, TypeDelayedEcho, 0, []byte(actionID))
return nil
}
// handle enforces spec §3.1/§3.4: unknown prefix, bad HMAC, expired session, // handle enforces spec §3.1/§3.4: unknown prefix, bad HMAC, expired session,
// replayed seq → silent drop, never a response. // replayed seq → silent drop, never a response.
func (s *Server) handle(conn *net.UDPConn, raddr netip.AddrPort, pkt []byte, tRxNs int64) { func (s *Server) handle(conn *net.UDPConn, raddr netip.AddrPort, pkt []byte, tRxNs int64) {
@@ -80,17 +124,34 @@ func (s *Server) handle(conn *net.UDPConn, raddr netip.AddrPort, pkt []byte, tRx
return return
} }
sess.NoteDataSource(raddr) sess.NoteDataSource(raddr)
sess.RecordUDP(session.UDPObservation{
Seq: seq, TRxNs: tRxNs, TTxNs: time.Since(s.start).Nanoseconds(),
Src: raddr.String(), Size: len(pkt), Type: typ,
})
switch typ { switch typ {
case TypeEchoReq: case TypeEchoReq:
s.echoResp(conn, raddr, sess, pkt, seq, tRxNs) s.echoResp(conn, raddr, sess, pkt, seq, tRxNs)
case TypeTimesyncReq: case TypeTimesyncReq:
s.timesyncResp(conn, raddr, sess, pkt, seq, tRxNs) s.timesyncResp(conn, raddr, sess, pkt, seq, tRxNs)
case TypeMtuProbe:
s.mtuAck(conn, raddr, sess, seq, len(pkt))
default: default:
slog.Debug("unhandled data-plane type", "type", typ) slog.Debug("unhandled data-plane type", "type", typ)
} }
} }
// mtuAck replies to an MTU_PROBE with a small MTU_ACK carrying the total
// datagram size the server actually received (spec §3.2). The client sends
// DF-flagged probes of increasing size and binary-searches the path MTU / a
// black hole from which sizes stop being acknowledged. The ACK is tiny, so it
// can never amplify regardless of probe size.
func (s *Server) mtuAck(conn *net.UDPConn, raddr netip.AddrPort, sess *session.Session, seq uint32, received int) {
var payload [4]byte
binary.BigEndian.PutUint32(payload[:], uint32(received))
s.send(conn, raddr, sess, TypeMtuAck, seq, payload[:])
}
// Observation block (spec §3.3), fixed 40 bytes appended to the RESP header: // Observation block (spec §3.3), fixed 40 bytes appended to the RESP header:
// 0 8 t_rx_ns (server clock, process epoch) // 0 8 t_rx_ns (server clock, process epoch)
// 8 8 t_tx_ns // 8 8 t_tx_ns
+34
View File
@@ -102,6 +102,40 @@ func TestEchoRoundtripObservationAndAntiAmplification(t *testing.T) {
} }
} }
func TestMtuProbeAckReportsReceivedSizeAndDoesNotAmplify(t *testing.T) {
mgr, addr := startServer(t)
sess, _, err := mgr.New("dev1", "credential-ikm", netip.MustParseAddr("127.0.0.1"))
if err != nil {
t.Fatal(err)
}
client, err := net.DialUDP("udp", nil, net.UDPAddrFromAddrPort(addr))
if err != nil {
t.Fatal(err)
}
defer client.Close()
client.SetDeadline(time.Now().Add(2 * time.Second))
// A large probe: 32 header + 1400 payload.
probe := craft(t, sess, TypeMtuProbe, 1, make([]byte, 1400))
if _, err := client.Write(probe); err != nil {
t.Fatal(err)
}
buf := make([]byte, 2000)
n, err := client.Read(buf)
if err != nil {
t.Fatalf("no MTU_ACK: %v", err)
}
if buf[4] != TypeMtuAck {
t.Fatalf("type = %#x, want MTU_ACK", buf[4])
}
if n >= len(probe) {
t.Fatalf("MTU_ACK (%d) must be far smaller than the probe (%d)", n, len(probe))
}
if got := binary.BigEndian.Uint32(buf[HeaderSize:n]); int(got) != len(probe) {
t.Fatalf("acked size %d, want %d", got, len(probe))
}
}
func TestDropsReplayBadHmacAndUnknownPrefix(t *testing.T) { func TestDropsReplayBadHmacAndUnknownPrefix(t *testing.T) {
mgr, addr := startServer(t) mgr, addr := startServer(t)
sess, _, err := mgr.New("dev1", "credential-ikm", netip.MustParseAddr("127.0.0.1")) sess, _, err := mgr.New("dev1", "credential-ikm", netip.MustParseAddr("127.0.0.1"))
+112
View File
@@ -0,0 +1,112 @@
// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
//go:build linux
package selftest
import (
"net"
"net/netip"
"syscall"
"time"
)
// Linux IP-level constants for PMTU discovery. Not all are exported by the
// stdlib syscall package across versions, so they are pinned here (stable
// kernel ABI) — same rationale as the prober's OsAbi.
const (
ipMTUDiscover = 10 // IP_MTU_DISCOVER
ipMTU = 14 // IP_MTU
ipPMTUDiscDo = 2 // IP_PMTUDISC_DO (set DF, honor PMTU)
ipv6MTUDiscover = 23 // IPV6_MTU_DISCOVER
ipv6MTU = 24 // IPV6_MTU
ipv6PMTUDiscDo = 2 // IPV6_PMTUDISC_DO
)
// probeEgressMTU sends a DF-flagged full-size UDP datagram toward target and
// reads back the kernel's discovered path MTU. A reduction below 1500 means
// the SERVER's own uplink can't carry full-size packets — so client MTU
// results would measure the server, not the client. No root, no raw socket:
// IP_MTU_DISCOVER + a getsockopt on IP_MTU, mirroring the prober's approach.
func probeEgressMTU(target string) MTUResult {
res := MTUResult{Target: target}
addr, err := netip.ParseAddr(target)
if err != nil {
// allow "host" that resolves
ips, e := net.LookupIP(target)
if e != nil || len(ips) == 0 {
res.Err = "resolve: " + errStr(err)
return res
}
addr, _ = netip.AddrFromSlice(ips[0])
}
addr = addr.Unmap()
is4 := addr.Is4()
fam := syscall.AF_INET6
if is4 {
fam = syscall.AF_INET
}
fd, err := syscall.Socket(fam, syscall.SOCK_DGRAM, 0)
if err != nil {
res.Err = "socket: " + errStr(err)
return res
}
defer syscall.Close(fd)
if is4 {
_ = syscall.SetsockoptInt(fd, syscall.IPPROTO_IP, ipMTUDiscover, ipPMTUDiscDo)
} else {
_ = syscall.SetsockoptInt(fd, syscall.IPPROTO_IPV6, ipv6MTUDiscover, ipv6PMTUDiscDo)
}
// IP_MTU reflects the CONNECTED path's MTU, so the socket must be connected
// (an unconnected socket returns ENOTCONN). No handshake — UDP connect just
// pins the destination and resolves the route.
sa := sockaddr(addr, 33434)
if err := syscall.Connect(fd, sa); err != nil {
res.Err = "connect: " + errStr(err)
return res
}
// Full-size probe: 1500 total IP/UDP headers (28 v4, 48 v6). A DF send
// larger than the local MTU fails immediately with EMSGSIZE; a path
// reduction updates IP_MTU after the ICMP frag-needed returns, so we send,
// briefly wait, and read the discovered MTU.
payload := 1472
if !is4 {
payload = 1452
}
probe := make([]byte, payload)
_, _ = syscall.Write(fd, probe)
time.Sleep(700 * time.Millisecond)
_, _ = syscall.Write(fd, probe) // second send observes any reduction
level, opt := syscall.IPPROTO_IP, ipMTU
if !is4 {
level, opt = syscall.IPPROTO_IPV6, ipv6MTU
}
mtu, err := syscall.GetsockoptInt(fd, level, opt)
if err != nil || mtu <= 0 {
res.Err = "getsockopt IP_MTU: " + errStr(err)
return res
}
res.DiscoveredMTU = mtu
res.FullMTU = mtu >= 1500
return res
}
func sockaddr(a netip.Addr, port int) syscall.Sockaddr {
if a.Is4() {
return &syscall.SockaddrInet4{Port: port, Addr: a.As4()}
}
return &syscall.SockaddrInet6{Port: port, Addr: a.As16()}
}
func errStr(err error) string {
if err == nil {
return "nil"
}
return err.Error()
}
+13
View File
@@ -0,0 +1,13 @@
// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
//go:build !linux
package selftest
// probeEgressMTU: PMTUD via IP_MTU_DISCOVER is Linux-specific. Off-Linux the
// self-test reports MTU as unproven rather than guessing (the daemon runs on
// Linux in production; this keeps dev builds compiling).
func probeEgressMTU(target string) MTUResult {
return MTUResult{Target: target, Err: "egress MTU probe is Linux-only"}
}
+126
View File
@@ -0,0 +1,126 @@
// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
// Package selftest lets the daemon prove its own host is a clean measurement
// target: the kernel isn't silently altering what clients measure, and the
// server's own egress reaches full MTU. If the server side is already broken,
// client-side results (especially MTU/PMTUD) measure the server, not the
// client — so the daemon says so.
package selftest
import (
"os"
"strconv"
"strings"
)
// Severity of a check result.
type Severity string
const (
OK Severity = "ok"
Warn Severity = "warn"
)
// Check is one sysctl (or derived) assertion.
type Check struct {
Name string `json:"name"`
Got string `json:"got"`
Want string `json:"want"`
Severity Severity `json:"severity"`
Why string `json:"why"`
}
// MTUResult is one egress path-MTU probe outcome.
type MTUResult struct {
Target string `json:"target"`
DiscoveredMTU int `json:"discovered_mtu"`
FullMTU bool `json:"full_mtu"` // >= 1500
Err string `json:"err,omitempty"`
}
// Report is the whole self-test.
type Report struct {
Sysctls []Check `json:"sysctls"`
EgressMTU []MTUResult `json:"egress_mtu"`
// SysctlOK / MTUOK are the compact "server proven good" signals; the
// profile surfaces these so a client can skip MTU tests the server can't
// support honestly.
SysctlOK bool `json:"sysctl_ok"`
MTUOK bool `json:"mtu_ok"`
}
// readSysctl reads /proc/sys/<dotted.name>. Empty string if unavailable.
func readSysctl(name string) string {
p := "/proc/sys/" + strings.ReplaceAll(name, ".", "/")
b, err := os.ReadFile(p)
if err != nil {
return ""
}
return strings.TrimSpace(string(b))
}
// sysctlChecks are the measurement-fidelity assertions. Each closure returns
// OK/Warn given the read value; a missing value (non-Linux / restricted) is
// reported as Warn "unreadable" but never fatal.
var sysctlChecks = []struct {
name string
want string
why string
ok func(v string) bool
}{
{"net.ipv6.conf.all.accept_ra", "0", "static v6 host must not let RAs mutate routing (the very thing Echolot detects)", eq("0")},
{"net.ipv4.conf.all.accept_redirects", "0", "ICMP redirects could alter routing mid-measurement", eq("0")},
{"net.ipv4.conf.all.send_redirects", "0", "an endpoint should not emit ICMP redirects", eq("0")},
{"net.ipv4.icmp_echo_ignore_all", "0", "server must answer ping so clients can measure to it", eq("0")},
{"net.ipv4.ip_no_pmtu_disc", "0", "server must honor path MTU on its own sends", eq("0")},
{"net.ipv4.tcp_sack", "1", "so a missing SACK in mss_observed is the path's fault, not the server's", eq("1")},
{"net.ipv4.tcp_timestamps", "1", "so TCP-timestamp absence reflects the path, not the server", eq("1")},
{"net.ipv4.tcp_window_scaling", "1", "so wscale absence reflects the path, not the server", eq("1")},
{"net.ipv4.icmp_ratelimit", "0", "nonzero throttles the server's ICMP errors → false loss/black-hole readings", eq("0")},
}
func eq(want string) func(string) bool { return func(v string) bool { return v == want } }
// Sysctls runs the sysctl audit.
func Sysctls() []Check {
out := make([]Check, 0, len(sysctlChecks))
for _, c := range sysctlChecks {
got := readSysctl(c.name)
sev := Warn
switch {
case got == "":
got = "(unreadable)"
case c.ok(got):
sev = OK
}
out = append(out, Check{Name: c.name, Got: got, Want: c.want, Severity: sev, Why: c.why})
}
return out
}
// Run performs the full self-test: sysctl audit + egress MTU probes to the
// given targets (each "host" — port is irrelevant for PMTUD).
func Run(mtuTargets []string) Report {
r := Report{Sysctls: Sysctls()}
r.SysctlOK = true
for _, c := range r.Sysctls {
if c.Severity == Warn {
r.SysctlOK = false
}
}
r.MTUOK = true
for _, t := range mtuTargets {
res := probeEgressMTU(t)
r.EgressMTU = append(r.EgressMTU, res)
if !res.FullMTU {
r.MTUOK = false
}
}
if len(r.EgressMTU) == 0 {
r.MTUOK = false // couldn't prove it
}
return r
}
var _ = strconv.Atoi
+72
View File
@@ -31,6 +31,64 @@ type Session struct {
// a bitmask of the 1024 preceding. // a bitmask of the 1024 preceding.
maxSeq uint32 maxSeq uint32
window [16]uint64 window [16]uint64
// Observations (spec §6): per-packet UDP view + connect-back results.
packetsSeen uint64
udpObs []UDPObservation // ring, newest last, cap obsCap
connectBack []ConnectBackResult
}
const obsCap = 4096
// UDPObservation is the server's witnessed view of one data-plane packet.
type UDPObservation struct {
Seq uint32 `json:"seq"`
TRxNs int64 `json:"t_rx_ns"`
TTxNs int64 `json:"t_tx_ns"`
Src string `json:"src"`
Size int `json:"size"`
Type uint8 `json:"type"`
}
// ConnectBackResult records one connect-back action outcome.
type ConnectBackResult struct {
ActionID string `json:"action_id"`
Result string `json:"result"` // connected | refused | timeout
RttMs float64 `json:"rtt_ms"`
}
// RecordUDP appends a packet observation (ring-capped).
func (s *Session) RecordUDP(o UDPObservation) {
s.mu.Lock()
defer s.mu.Unlock()
s.packetsSeen++
if len(s.udpObs) >= obsCap {
s.udpObs = s.udpObs[1:]
}
s.udpObs = append(s.udpObs, o)
}
// RecordConnectBack appends a connect-back outcome.
func (s *Session) RecordConnectBack(r ConnectBackResult) {
s.mu.Lock()
defer s.mu.Unlock()
s.connectBack = append(s.connectBack, r)
}
// Observations returns a copy of everything witnessed so far.
func (s *Session) Observations() (packetsSeen uint64, udp []UDPObservation, cb []ConnectBackResult) {
s.mu.Lock()
defer s.mu.Unlock()
return s.packetsSeen, append([]UDPObservation(nil), s.udpObs...),
append([]ConnectBackResult(nil), s.connectBack...)
}
// DataSource returns the last verified data-plane source (invalid when the
// session has not sent data-plane traffic yet).
func (s *Session) DataSource() netip.AddrPort {
s.mu.Lock()
defer s.mu.Unlock()
return s.dataSource
} }
// KeySalt returns nothing — the salt is not retained after derivation; it is // KeySalt returns nothing — the salt is not retained after derivation; it is
@@ -87,6 +145,20 @@ func (m *Manager) ByWirePrefix(prefix [8]byte) *Session {
return s return s
} }
// ByID resolves a full session id (sessions are keyed by their wire prefix).
func (m *Manager) ByID(id string) *Session {
if len(id) < 16 {
return nil
}
m.mu.Lock()
defer m.mu.Unlock()
s := m.byPrefix[id[:16]]
if s == nil || s.ID != id || time.Now().After(s.Expires) {
return nil
}
return s
}
func (m *Manager) Delete(id string) { func (m *Manager) Delete(id string) {
m.mu.Lock() m.mu.Lock()
defer m.mu.Unlock() defer m.mu.Unlock()
+266
View File
@@ -0,0 +1,266 @@
// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
// Package stun implements an unmodified RFC 5389 STUN binding responder with
// the RFC 5780 NAT-behavior-discovery attributes (OTHER-ADDRESS,
// RESPONSE-ORIGIN, CHANGE-REQUEST) when alternate addresses are available.
// No custom framing — interop with existing STUN tooling is a feature
// (spec §4). Each configured primary address gets two sockets: the given
// port and port+1 (the RFC 5780 alternate-port convention).
package stun
import (
"crypto/rand"
"encoding/binary"
"log/slog"
"net"
"net/netip"
)
const (
magicCookie = 0x2112A442
typeBindingRequest = 0x0001
typeBindingSuccess = 0x0101
attrChangeRequest = 0x0003
attrXorMapped = 0x0020
attrSoftware = 0x8022
attrResponseOrigin = 0x802B
attrOtherAddress = 0x802C
changeIP = 0x04
changePort = 0x02
)
// sock is one bound socket, addressable by (address index, port index).
type sock struct {
conn *net.UDPConn
addr netip.AddrPort
}
// Server holds the socket grid: addrs × {primary, alternate} ports.
type Server struct {
// socks[i][0] = primary port, socks[i][1] = alt port for address i.
socks [][2]*sock
}
// Listen binds primary+alternate sockets for every address. Addresses are
// "ip:port" specs; the alternate port is port+1.
func Listen(addrs []string) (*Server, error) {
s := &Server{}
for _, spec := range addrs {
ap, err := netip.ParseAddrPort(spec)
if err != nil {
return nil, err
}
var pair [2]*sock
for i, port := range []uint16{ap.Port(), ap.Port() + 1} {
bind := netip.AddrPortFrom(ap.Addr(), port)
conn, err := net.ListenUDP("udp", net.UDPAddrFromAddrPort(bind))
if err != nil {
s.Close()
return nil, err
}
pair[i] = &sock{conn: conn, addr: bind}
}
s.socks = append(s.socks, pair)
}
return s, nil
}
func (s *Server) Close() {
for _, pair := range s.socks {
for _, sk := range pair {
if sk != nil {
sk.conn.Close()
}
}
}
}
// Has5780 reports whether any address family has ≥2 addresses — the
// prerequisite for full NAT behavior discovery.
func (s *Server) Has5780() bool {
var v4, v6 int
for _, pair := range s.socks {
if pair[0].addr.Addr().Is4() || pair[0].addr.Addr().Is4In6() {
v4++
} else {
v6++
}
}
return v4 >= 2 || v6 >= 2
}
// Serve starts one read loop per socket and blocks until the first error.
func (s *Server) Serve() error {
errCh := make(chan error, len(s.socks)*2)
for ai := range s.socks {
for pi := range s.socks[ai] {
go func(ai, pi int) { errCh <- s.loop(ai, pi) }(ai, pi)
}
}
return <-errCh
}
func (s *Server) loop(ai, pi int) error {
sk := s.socks[ai][pi]
buf := make([]byte, 1500)
for {
n, raddr, err := sk.conn.ReadFromUDPAddrPort(buf)
if err != nil {
return err
}
s.handle(ai, pi, buf[:n], raddr)
}
}
// otherAddr finds the "diagonal" alternate for RFC 5780: different address
// (same family), different port. Returns nil when there is none.
func (s *Server) other(ai int, sameFamily bool, fam4 bool) int {
for i, pair := range s.socks {
if i == ai {
continue
}
is4 := pair[0].addr.Addr().Is4() || pair[0].addr.Addr().Is4In6()
if !sameFamily || is4 == fam4 {
return i
}
}
return -1
}
func (s *Server) handle(ai, pi int, pkt []byte, raddr netip.AddrPort) {
if len(pkt) < 20 || binary.BigEndian.Uint16(pkt[0:2]) != typeBindingRequest {
return
}
if binary.BigEndian.Uint32(pkt[4:8]) != magicCookie {
return
}
msgLen := int(binary.BigEndian.Uint16(pkt[2:4]))
if 20+msgLen > len(pkt) {
return
}
var txid [12]byte
copy(txid[:], pkt[8:20])
// Parse CHANGE-REQUEST if present (RFC 5780 §7.2).
var change byte
for off := 20; off+4 <= 20+msgLen; {
at := binary.BigEndian.Uint16(pkt[off : off+2])
al := int(binary.BigEndian.Uint16(pkt[off+2 : off+4]))
if off+4+al > len(pkt) {
break
}
if at == attrChangeRequest && al >= 4 {
change = pkt[off+7]
}
off += 4 + al + (4-al%4)%4 // attributes are 32-bit aligned
}
// Pick the responding socket per CHANGE-REQUEST.
fam4 := raddr.Addr().Is4() || raddr.Addr().Is4In6()
rai, rpi := ai, pi
if change&changeIP != 0 {
if o := s.other(ai, true, fam4); o >= 0 {
rai = o
} else {
return // cannot honor — RFC says error response; silence is safer for a probe target
}
}
if change&changePort != 0 {
rpi = 1 - pi
}
responder := s.socks[rai][rpi]
resp := buildResponse(txid, raddr, responder.addr, s.otherAddress(ai, fam4))
if _, err := responder.conn.WriteToUDPAddrPort(resp, raddr); err != nil {
slog.Debug("stun write failed", "to", raddr, "err", err)
}
}
// otherAddress computes the OTHER-ADDRESS attribute value (alt IP, alt port)
// for the client's family, or an invalid AddrPort when unavailable.
func (s *Server) otherAddress(ai int, fam4 bool) netip.AddrPort {
if o := s.other(ai, true, fam4); o >= 0 {
return s.socks[o][1].addr
}
return netip.AddrPort{}
}
func buildResponse(txid [12]byte, mapped, origin, other netip.AddrPort) []byte {
attrs := xorMappedAttr(attrXorMapped, mapped, txid)
attrs = append(attrs, addrAttr(attrResponseOrigin, origin)...)
if other.IsValid() {
attrs = append(attrs, addrAttr(attrOtherAddress, other)...)
}
sw := []byte("echolot")
attrs = append(attrs, attrHeader(attrSoftware, len(sw))...)
attrs = append(attrs, pad4(sw)...)
msg := make([]byte, 20, 20+len(attrs))
binary.BigEndian.PutUint16(msg[0:2], typeBindingSuccess)
binary.BigEndian.PutUint16(msg[2:4], uint16(len(attrs)))
binary.BigEndian.PutUint32(msg[4:8], magicCookie)
copy(msg[8:20], txid[:])
return append(msg, attrs...)
}
func attrHeader(typ uint16, valLen int) []byte {
h := make([]byte, 4)
binary.BigEndian.PutUint16(h[0:2], typ)
binary.BigEndian.PutUint16(h[2:4], uint16(valLen))
return h
}
func pad4(b []byte) []byte {
for len(b)%4 != 0 {
b = append(b, 0)
}
return b
}
// addrValue encodes the RFC 5389 address structure (family, port, address).
func addrValue(ap netip.AddrPort) []byte {
addr := ap.Addr().Unmap()
if addr.Is4() {
v := make([]byte, 8)
v[1] = 0x01
binary.BigEndian.PutUint16(v[2:4], ap.Port())
a4 := addr.As4()
copy(v[4:], a4[:])
return v
}
v := make([]byte, 20)
v[1] = 0x02
binary.BigEndian.PutUint16(v[2:4], ap.Port())
a16 := addr.As16()
copy(v[4:], a16[:])
return v
}
func addrAttr(typ uint16, ap netip.AddrPort) []byte {
v := addrValue(ap)
return append(attrHeader(typ, len(v)), v...)
}
// xorMappedAttr encodes XOR-MAPPED-ADDRESS (RFC 5389 §15.2).
func xorMappedAttr(typ uint16, ap netip.AddrPort, txid [12]byte) []byte {
v := addrValue(ap)
binary.BigEndian.PutUint16(v[2:4], ap.Port()^uint16(magicCookie>>16))
var key [16]byte
binary.BigEndian.PutUint32(key[0:4], magicCookie)
copy(key[4:], txid[:])
for i := 4; i < len(v); i++ {
v[i] ^= key[i-4]
}
return append(attrHeader(typ, len(v)), v...)
}
// NewTxID is exported for tests and client code.
func NewTxID() [12]byte {
var t [12]byte
_, _ = rand.Read(t[:])
return t
}
+137
View File
@@ -0,0 +1,137 @@
// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
package stun
import (
"encoding/binary"
"net"
"net/netip"
"testing"
"time"
)
// bindingRequest builds a minimal RFC 5389 binding request.
func bindingRequest(txid [12]byte, change byte) []byte {
var attrs []byte
if change != 0 {
attrs = append(attrs, attrHeader(attrChangeRequest, 4)...)
attrs = append(attrs, 0, 0, 0, change)
}
msg := make([]byte, 20, 20+len(attrs))
binary.BigEndian.PutUint16(msg[0:2], typeBindingRequest)
binary.BigEndian.PutUint16(msg[2:4], uint16(len(attrs)))
binary.BigEndian.PutUint32(msg[4:8], magicCookie)
copy(msg[8:20], txid[:])
return append(msg, attrs...)
}
// parseXorMapped extracts XOR-MAPPED-ADDRESS from a binding success.
func parseXorMapped(t *testing.T, resp []byte, txid [12]byte) netip.AddrPort {
t.Helper()
if binary.BigEndian.Uint16(resp[0:2]) != typeBindingSuccess {
t.Fatalf("type = %#x, want binding success", resp[0:2])
}
msgLen := int(binary.BigEndian.Uint16(resp[2:4]))
for off := 20; off+4 <= 20+msgLen; {
at := binary.BigEndian.Uint16(resp[off : off+2])
al := int(binary.BigEndian.Uint16(resp[off+2 : off+4]))
if at == attrXorMapped {
v := append([]byte(nil), resp[off+4:off+4+al]...)
port := binary.BigEndian.Uint16(v[2:4]) ^ uint16(magicCookie>>16)
var key [16]byte
binary.BigEndian.PutUint32(key[0:4], magicCookie)
copy(key[4:], txid[:])
for i := 4; i < len(v); i++ {
v[i] ^= key[i-4]
}
if v[1] == 0x01 {
return netip.AddrPortFrom(netip.AddrFrom4([4]byte(v[4:8])), port)
}
return netip.AddrPortFrom(netip.AddrFrom16([16]byte(v[4:20])), port)
}
off += 4 + al + (4-al%4)%4
}
t.Fatal("no XOR-MAPPED-ADDRESS in response")
return netip.AddrPort{}
}
func TestBindingAndChangePort(t *testing.T) {
// Two loopback "addresses" is not possible portably, so exercise one
// address (basic binding + change-port); the change-IP path needs the
// two-address grid of a real deployment.
srv, err := Listen([]string{"127.0.0.1:0"})
if err != nil {
t.Fatal(err)
}
// port 0 twice would collide at 1; rebind explicitly on free ports
srv.Close()
base := freePort(t)
srv, err = Listen([]string{netip.AddrPortFrom(netip.MustParseAddr("127.0.0.1"), base).String()})
if err != nil {
t.Skipf("cannot bind %d/%d: %v", base, base+1, err)
}
defer srv.Close()
go srv.Serve()
client, err := net.DialUDP("udp", nil, srv.socks[0][0].conn.LocalAddr().(*net.UDPAddr))
if err != nil {
t.Fatal(err)
}
defer client.Close()
client.SetDeadline(time.Now().Add(2 * time.Second))
txid := NewTxID()
client.Write(bindingRequest(txid, 0))
buf := make([]byte, 1500)
n, err := client.Read(buf)
if err != nil {
t.Fatal(err)
}
mapped := parseXorMapped(t, buf[:n], txid)
want := client.LocalAddr().(*net.UDPAddr).AddrPort()
if mapped.Port() != want.Port() {
t.Fatalf("mapped port %d, want %d", mapped.Port(), want.Port())
}
// CHANGE-REQUEST(port): response must come from the alternate port.
// Dial-connected sockets drop packets from other sources, so use an
// unconnected socket and inspect the reply's source.
uc, err := net.ListenUDP("udp", &net.UDPAddr{IP: net.IPv4(127, 0, 0, 1)})
if err != nil {
t.Fatal(err)
}
defer uc.Close()
uc.SetDeadline(time.Now().Add(2 * time.Second))
txid2 := NewTxID()
uc.WriteToUDPAddrPort(bindingRequest(txid2, changePort), srv.socks[0][0].addr)
n, from, err := uc.ReadFromUDPAddrPort(buf)
if err != nil {
t.Fatal(err)
}
if from.Port() != srv.socks[0][1].addr.Port() {
t.Fatalf("change-port reply came from %v, want alt port %d", from, srv.socks[0][1].addr.Port())
}
parseXorMapped(t, buf[:n], txid2)
}
func freePort(t *testing.T) uint16 {
t.Helper()
// Find two adjacent free ports for the primary/alternate pair.
for tries := 0; tries < 20; tries++ {
l, err := net.ListenUDP("udp", &net.UDPAddr{IP: net.IPv4(127, 0, 0, 1)})
if err != nil {
continue
}
p := l.LocalAddr().(*net.UDPAddr).Port
l.Close()
l2, err := net.ListenUDP("udp", &net.UDPAddr{IP: net.IPv4(127, 0, 0, 1), Port: p + 1})
if err != nil {
continue
}
l2.Close()
return uint16(p)
}
t.Skip("no adjacent free UDP ports found")
return 0
}
+91
View File
@@ -0,0 +1,91 @@
// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
// Package tcpecho implements the spec §4 TCP echo: after connect the server
// sends one JSON line with what it observed (source address/port, negotiated
// MSS and TCP options from TCP_INFO), then byte-echoes until FIN. This is the
// evidence source for mtu.mss_observed.
//
// The TLS/ALPN "elt-echo" variant (ClientHello capture + JA4) is not
// implemented yet.
package tcpecho
import (
"encoding/json"
"io"
"net"
"sync"
"time"
)
// ConnRecord is what the observations API reports per connection (spec §6).
type ConnRecord struct {
ConnectedAt time.Time `json:"connected_at"`
Src string `json:"src"`
MSS int `json:"mss"`
Options []string `json:"options"`
}
type Server struct {
mu sync.Mutex
recent []ConnRecord // ring, newest last
}
const recentCap = 1024
func (s *Server) record(r ConnRecord) {
s.mu.Lock()
defer s.mu.Unlock()
if len(s.recent) >= recentCap {
s.recent = s.recent[1:]
}
s.recent = append(s.recent, r)
}
// RecentFor returns records whose source IP matches ip.
func (s *Server) RecentFor(ip string) []ConnRecord {
s.mu.Lock()
defer s.mu.Unlock()
var out []ConnRecord
for _, r := range s.recent {
if h, _, err := net.SplitHostPort(r.Src); err == nil && h == ip {
out = append(out, r)
}
}
return out
}
func (s *Server) Serve(ln net.Listener) error {
for {
conn, err := ln.Accept()
if err != nil {
return err
}
go s.handle(conn)
}
}
func (s *Server) handle(conn net.Conn) {
defer conn.Close()
_ = conn.SetDeadline(time.Now().Add(5 * time.Minute))
info := tcpInfo(conn) // platform-specific; zero values off-Linux
rec := ConnRecord{
ConnectedAt: time.Now().UTC(),
Src: conn.RemoteAddr().String(),
MSS: info.MSS,
Options: info.Options,
}
s.record(rec)
greeting, _ := json.Marshal(map[string]any{
"observed_src": rec.Src,
"mss": rec.MSS,
"options": rec.Options,
})
if _, err := conn.Write(append(greeting, '\n')); err != nil {
return
}
// Byte-echo until FIN; the client's data is its own to interpret.
_, _ = io.Copy(conn, conn)
}
+69
View File
@@ -0,0 +1,69 @@
// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
//go:build linux
package tcpecho
import (
"encoding/binary"
"net"
"syscall"
"unsafe"
)
type connInfo struct {
MSS int
Options []string
}
// tcpInfo reads TCP_INFO via getsockopt. Only the head of struct tcp_info is
// needed: 8 header bytes (state..wscale flags) then u32 rto, ato, snd_mss,
// rcv_mss — layout is part of the kernel ABI and stable.
func tcpInfo(conn net.Conn) connInfo {
tc, ok := conn.(*net.TCPConn)
if !ok {
return connInfo{}
}
raw, err := tc.SyscallConn()
if err != nil {
return connInfo{}
}
var buf [104]byte
var got bool
_ = raw.Control(func(fd uintptr) {
l := uint32(len(buf))
_, _, errno := syscall.Syscall6(syscall.SYS_GETSOCKOPT, fd,
uintptr(syscall.SOL_TCP), uintptr(syscall.TCP_INFO),
uintptr(unsafe.Pointer(&buf[0])), uintptr(unsafe.Pointer(&l)), 0)
got = errno == 0 && l >= 24
})
if !got {
return connInfo{}
}
// tcpi_options bit flags (include/uapi/linux/tcp.h)
const (
optTimestamps = 1
optSACK = 2
optWscale = 4
optECN = 8
)
var opts []string
ob := buf[5]
if ob&optTimestamps != 0 {
opts = append(opts, "timestamps")
}
if ob&optSACK != 0 {
opts = append(opts, "sack")
}
if ob&optWscale != 0 {
opts = append(opts, "wscale")
}
if ob&optECN != 0 {
opts = append(opts, "ecn")
}
return connInfo{
MSS: int(binary.LittleEndian.Uint32(buf[16:20])), // tcpi_snd_mss
Options: opts,
}
}
+17
View File
@@ -0,0 +1,17 @@
// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
//go:build !linux
package tcpecho
import "net"
type connInfo struct {
MSS int
Options []string
}
// tcpInfo: TCP_INFO is Linux-only; other platforms report zero values and
// the greeting says mss:0 — honest absence rather than a guess.
func tcpInfo(net.Conn) connInfo { return connInfo{} }