Types 0x03/0x04/0x05 land with a bounded columnar train buffer (head kept, truncation declared) and grant-free multi-part reports - a report row is smaller than the packet it answers, so $3.4 holds without a grant. The read loop now collects TTL/TOS cmsgs on Linux, replacing the 0xFF stubs in the observation block with what the kernel saw; downtrain gained a dscp parameter, so DSCP survival is measurable in both directions. Rate limiting ($2.5) exists now: per-credential AND per-source buckets, 429 on the control plane, silent drop on the data plane after the HMAC gate and before the replay window. UDP ceilings default above the largest legitimate run - a limit that clips a real measurement produces a confidently wrong number. Every granted packet carries its action_id at payload[8:16]; overlapping actions were unattributable before. Canary DNS logs now honor the stated 24h privacy default. /admin/enroll-tokens answers the spec's JSON shape. protocol_version 1.0.1 (additive). Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
372 lines
13 KiB
Go
372 lines
13 KiB
Go
// SPDX-FileCopyrightText: 2026 Echolot contributors
|
|
// SPDX-License-Identifier: GPL-3.0-or-later
|
|
|
|
// Package dataplane implements the binary UDP probe protocol (spec §3):
|
|
// 32-byte header, HMAC gate, anti-replay, ECHO with observation block,
|
|
// upstream trains with columnar reports, and the granted server->client sends.
|
|
package dataplane
|
|
|
|
import (
|
|
"crypto/hmac"
|
|
"crypto/sha256"
|
|
"encoding/binary"
|
|
"encoding/hex"
|
|
"fmt"
|
|
"log/slog"
|
|
"net"
|
|
"net/netip"
|
|
"sync"
|
|
"time"
|
|
|
|
"echo-lot.app/server/internal/ratelimit"
|
|
"echo-lot.app/server/internal/session"
|
|
)
|
|
|
|
const (
|
|
Magic = "ELT1"
|
|
HeaderSize = 32
|
|
|
|
TypeEchoReq = 0x01
|
|
TypeEchoResp = 0x02
|
|
// Upstream trains (spec §3.2): DATA gets no per-packet response; REPORT_REQ fetches the
|
|
// server's received view as one or more REPORT datagrams (train.go).
|
|
TypeTrainData = 0x03
|
|
TypeTrainReportReq = 0x04
|
|
TypeTrainReport = 0x05
|
|
TypeTimesyncReq = 0x07
|
|
TypeTimesyncRsp = 0x08
|
|
TypeMtuProbe = 0x09
|
|
TypeMtuAck = 0x0A
|
|
TypeDelayedEcho = 0x0B
|
|
// Server->client under an asymmetric grant (spec §3.4/§5).
|
|
TypeDownTrainData = 0x06
|
|
TypeBigSend = 0x0C
|
|
// TypeFragData is delivered only after IP reassembly, so its arrival IS the measurement.
|
|
TypeFragData = 0x0D
|
|
// TypeThroughputData is one packet of a sustained-rate downstream run.
|
|
TypeThroughputData = 0x0E
|
|
// TypeThroughputUp is one packet of a client-driven upstream run. The server counts it and
|
|
// deliberately does not answer: a reply would double the traffic and measure the return
|
|
// path at the same time, which is the one thing this test is trying not to do.
|
|
TypeThroughputUp = 0x0F
|
|
)
|
|
|
|
type Server struct {
|
|
Sessions *session.Manager
|
|
// Spec §2.5 ceilings on verified traffic, silent-drop (nil = no ceiling). Charged after the
|
|
// HMAC gate so an unauthenticated flood cannot spend anyone's budget, keyed per source
|
|
// address AND per device credential so neither one hot address nor one hot credential can
|
|
// crowd out the rest.
|
|
PacketRate *ratelimit.Limiter // tokens are packets
|
|
ByteRate *ratelimit.Limiter // tokens are bytes
|
|
// Epoch for server-side t_rx/t_tx: process start; observation consumers
|
|
// only need differences plus the timesync exchange, not absolute time.
|
|
start time.Time
|
|
|
|
mu sync.Mutex
|
|
conns []*net.UDPConn
|
|
|
|
// dfMu serialises DF windows: the listening socket is shared by every session on that
|
|
// family, so two concurrent big_sends must not overlap their DF on/off transitions.
|
|
dfMu sync.Mutex
|
|
}
|
|
|
|
// pktMeta is what the kernel told us about one received datagram beyond its bytes (spec §3.3:
|
|
// received TTL, DSCP, ECN). 0xFF means "not observed": non-Linux hosts and datagrams whose
|
|
// cmsg never arrived keep the sentinel rather than inventing a value.
|
|
type pktMeta struct {
|
|
TTL uint8
|
|
TOS uint8 // the whole DSCP/ECN byte; DSCP = TOS>>2, ECN = TOS&3
|
|
}
|
|
|
|
const metaUnavailable = 0xFF
|
|
|
|
func (m pktMeta) dscp() uint8 {
|
|
if m.TOS == metaUnavailable {
|
|
return metaUnavailable
|
|
}
|
|
return m.TOS >> 2
|
|
}
|
|
|
|
func (m pktMeta) ecn() uint8 {
|
|
if m.TOS == metaUnavailable {
|
|
return metaUnavailable
|
|
}
|
|
return m.TOS & 0x3
|
|
}
|
|
|
|
// oobCap fits the two cmsgs (TTL + TOS, each ≤ CMSG_SPACE(4)) with headroom for whatever else
|
|
// the kernel decides to attach.
|
|
const oobCap = 64
|
|
|
|
// 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 {
|
|
s.mu.Lock()
|
|
if s.start.IsZero() {
|
|
s.start = time.Now()
|
|
}
|
|
s.conns = append(s.conns, conn)
|
|
s.mu.Unlock()
|
|
enableRecvMeta(conn)
|
|
buf := make([]byte, 65535)
|
|
oob := make([]byte, oobCap)
|
|
for {
|
|
n, oobn, _, raddr, err := conn.ReadMsgUDPAddrPort(buf, oob)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
tRx := time.Since(s.start).Nanoseconds()
|
|
s.handle(conn, raddr, buf[:n], tRx, parseMeta(oob[:oobn]))
|
|
}
|
|
}
|
|
|
|
// connFor picks a retained socket whose family matches the target.
|
|
// connFor picks the socket to send to target from.
|
|
//
|
|
// When the session recorded which local address it has been talking to (local), that socket wins
|
|
// outright. Falling back to "any socket of the right family" is only correct for a single-homed
|
|
// server: on a multi-homed one it sends from a sibling address the client's NAT has no mapping
|
|
// for, the packets are dropped in transit, and the client reports downstream loss that does not
|
|
// exist. That bug is invisible in a lab with one address, which is exactly why this is explicit.
|
|
func (s *Server) connFor(target, local netip.AddrPort) *net.UDPConn {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
if local.IsValid() {
|
|
for _, c := range s.conns {
|
|
if c.LocalAddr().(*net.UDPAddr).AddrPort() == local {
|
|
return c
|
|
}
|
|
}
|
|
}
|
|
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, sess.DataLocal())
|
|
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,
|
|
// replayed seq → silent drop, never a response.
|
|
func (s *Server) handle(conn *net.UDPConn, raddr netip.AddrPort, pkt []byte, tRxNs int64, meta pktMeta) {
|
|
if len(pkt) < HeaderSize || string(pkt[0:4]) != Magic {
|
|
return
|
|
}
|
|
typ := pkt[4]
|
|
payloadLen := binary.BigEndian.Uint16(pkt[6:8])
|
|
if int(HeaderSize+payloadLen) > len(pkt) {
|
|
return
|
|
}
|
|
var prefix [8]byte
|
|
copy(prefix[:], pkt[8:16])
|
|
seq := binary.BigEndian.Uint32(pkt[16:20])
|
|
|
|
sess := s.Sessions.ByWirePrefix(prefix)
|
|
if sess == nil {
|
|
return
|
|
}
|
|
mac := hmac.New(sha256.New, sess.Key[:])
|
|
mac.Write(pkt[0:28])
|
|
mac.Write(pkt[HeaderSize : HeaderSize+int(payloadLen)])
|
|
if !hmac.Equal(mac.Sum(nil)[:4], pkt[28:32]) {
|
|
return
|
|
}
|
|
// Spec §2.5: over-ceiling traffic is silently dropped (probes tolerate loss by design).
|
|
// After the HMAC gate so a spoofed flood cannot drain a victim's budget; before the replay
|
|
// window so a dropped packet's seq stays usable for a resend.
|
|
if !s.allowUDP(raddr, sess.Device, len(pkt)) {
|
|
return
|
|
}
|
|
if !sess.CheckSeq(seq) {
|
|
return
|
|
}
|
|
sess.NoteDataSource(raddr)
|
|
if la, ok := conn.LocalAddr().(*net.UDPAddr); ok {
|
|
sess.NoteDataLocal(la.AddrPort())
|
|
}
|
|
|
|
// Upstream throughput short-circuits before the observation log. Recording one struct per
|
|
// packet here would mean tens of thousands of allocations for a single run; the counter is
|
|
// all anyone needs, since the client holds the send-side record.
|
|
if typ == TypeThroughputUp {
|
|
sess.CountUpstream(len(pkt), tRxNs)
|
|
return
|
|
}
|
|
sess.RecordUDP(session.UDPObservation{
|
|
Seq: seq, TRxNs: tRxNs, TTxNs: time.Since(s.start).Nanoseconds(),
|
|
Src: raddr.String(), Size: len(pkt), Type: typ,
|
|
})
|
|
|
|
payload := pkt[HeaderSize : HeaderSize+int(payloadLen)]
|
|
switch typ {
|
|
case TypeEchoReq:
|
|
s.echoResp(conn, raddr, sess, pkt, seq, tRxNs, meta)
|
|
case TypeTrainData:
|
|
// No response (spec §3.2): the train is upstream-only; its received view is fetched
|
|
// afterwards via TRAIN_REPORT_REQ or the observations API.
|
|
recordTrain(sess, payload, seq, len(pkt), tRxNs, meta)
|
|
case TypeTrainReportReq:
|
|
s.trainReport(conn, raddr, sess, payload)
|
|
case TypeTimesyncReq:
|
|
s.timesyncResp(conn, raddr, sess, pkt, seq, tRxNs)
|
|
case TypeMtuProbe:
|
|
s.mtuAck(conn, raddr, sess, seq, len(pkt))
|
|
default:
|
|
slog.Debug("unhandled data-plane type", "type", typ)
|
|
}
|
|
}
|
|
|
|
// allowUDP charges the §2.5 packet and byte buckets, per source address and per credential.
|
|
func (s *Server) allowUDP(raddr netip.AddrPort, device string, size int) bool {
|
|
ipKey, credKey := "ip:"+raddr.Addr().String(), "cred:"+device
|
|
okA, _ := s.PacketRate.Allow(ipKey)
|
|
okC, _ := s.PacketRate.Allow(credKey)
|
|
okAB, _ := s.ByteRate.AllowN(ipKey, float64(size))
|
|
okCB, _ := s.ByteRate.AllowN(credKey, float64(size))
|
|
return okA && okC && okAB && okCB
|
|
}
|
|
|
|
// 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:
|
|
//
|
|
// 0 8 t_rx_ns (server clock, process epoch)
|
|
// 8 8 t_tx_ns
|
|
// 16 16 observed source IP (v4-mapped when v4)
|
|
// 32 2 observed source port
|
|
// 34 1 received TTL (0xFF = not observed; cmsgs unavailable on this host)
|
|
// 35 1 received DSCP/ECN byte (0xFF = not observed)
|
|
// 36 4 received size
|
|
func observation(tRxNs, tTxNs int64, src netip.AddrPort, rcvd int, meta pktMeta) []byte {
|
|
b := make([]byte, 40)
|
|
binary.BigEndian.PutUint64(b[0:8], uint64(tRxNs))
|
|
binary.BigEndian.PutUint64(b[8:16], uint64(tTxNs))
|
|
a16 := src.Addr().As16()
|
|
copy(b[16:32], a16[:])
|
|
binary.BigEndian.PutUint16(b[32:34], src.Port())
|
|
b[34], b[35] = meta.TTL, meta.TOS
|
|
binary.BigEndian.PutUint32(b[36:40], uint32(rcvd))
|
|
return b
|
|
}
|
|
|
|
// putActionID writes a grant's action id into payload[8:16] — the correlation the spec promises
|
|
// (§5: "an action_id echoed in resulting data-plane packets"), consumed by the client as
|
|
// test.params.action_id (§9). Bytes [0:8] stay with the packet type; [8:16] is reserved for this
|
|
// across every granted type, so the client needs one rule, not five.
|
|
func putActionID(payload []byte, actionID string) {
|
|
if len(payload) < 16 {
|
|
return
|
|
}
|
|
raw, err := hex.DecodeString(actionID)
|
|
if err != nil || len(raw) != 8 {
|
|
return // a malformed id yields zero bytes, not a crash mid-burst
|
|
}
|
|
copy(payload[8:16], raw)
|
|
}
|
|
|
|
// echoResp mirrors the request header (type flipped), appends the observation
|
|
// block, and re-HMACs with the session key. Anti-amplification: the response
|
|
// is capped at the request size (spec §3.4) — the observation block replaces
|
|
// padding rather than growing the datagram; if the request was smaller than
|
|
// header+observation, the block is truncated to fit.
|
|
func (s *Server) echoResp(conn *net.UDPConn, raddr netip.AddrPort, sess *session.Session, req []byte, seq uint32, tRxNs int64, meta pktMeta) {
|
|
obs := observation(tRxNs, time.Since(s.start).Nanoseconds(), raddr, len(req), meta)
|
|
max := len(req)
|
|
if max < HeaderSize {
|
|
return
|
|
}
|
|
payload := obs
|
|
if HeaderSize+len(payload) > max {
|
|
payload = payload[:max-HeaderSize]
|
|
}
|
|
s.send(conn, raddr, sess, TypeEchoResp, seq, payload)
|
|
}
|
|
|
|
// timesyncResp: payload = client t1 (echoed back) + t2 (rx) + t3 (tx), spec §3.2.
|
|
func (s *Server) timesyncResp(conn *net.UDPConn, raddr netip.AddrPort, sess *session.Session, req []byte, seq uint32, tRxNs int64) {
|
|
payload := make([]byte, 24)
|
|
copy(payload[0:8], req[20:28]) // client's t_ns from the request header
|
|
binary.BigEndian.PutUint64(payload[8:16], uint64(tRxNs))
|
|
binary.BigEndian.PutUint64(payload[16:24], uint64(time.Since(s.start).Nanoseconds()))
|
|
s.send(conn, raddr, sess, TypeTimesyncRsp, seq, payload)
|
|
}
|
|
|
|
func (s *Server) send(conn *net.UDPConn, raddr netip.AddrPort, sess *session.Session, typ byte, seq uint32, payload []byte) {
|
|
_ = s.sendErr(conn, raddr, sess, typ, seq, payload)
|
|
}
|
|
|
|
// sendErr is send with the write error surfaced. Only the DF-mode big_send cares: there an
|
|
// EMSGSIZE means our own egress MTU refused the datagram, which is a different fact from the
|
|
// client not receiving it.
|
|
func (s *Server) sendErr(conn *net.UDPConn, raddr netip.AddrPort, sess *session.Session, typ byte, seq uint32, payload []byte) error {
|
|
pkt := s.buildPacket(sess, typ, seq, payload)
|
|
_, err := conn.WriteToUDPAddrPort(pkt, raddr)
|
|
return err
|
|
}
|
|
|
|
// buildPacket assembles and signs an ELT1 packet without sending it.
|
|
//
|
|
// Split out for the crafted-fragment path, which needs the bytes so it can cut them up itself.
|
|
// What arrives after reassembly must be indistinguishable from an ordinary packet, or the client
|
|
// would be measuring our sender rather than the path — so it goes through exactly this function.
|
|
func (s *Server) buildPacket(sess *session.Session, typ byte, seq uint32, payload []byte) []byte {
|
|
pkt := make([]byte, HeaderSize+len(payload))
|
|
copy(pkt[0:4], Magic)
|
|
pkt[4] = typ
|
|
binary.BigEndian.PutUint16(pkt[6:8], uint16(len(payload)))
|
|
idBytes := sess.ID[:16] // hex chars of the 8-byte prefix
|
|
for i := 0; i < 8; i++ {
|
|
pkt[8+i] = hexByte(idBytes[i*2], idBytes[i*2+1])
|
|
}
|
|
binary.BigEndian.PutUint32(pkt[16:20], seq)
|
|
binary.BigEndian.PutUint64(pkt[20:28], uint64(time.Since(s.start).Nanoseconds()))
|
|
copy(pkt[HeaderSize:], payload)
|
|
mac := hmac.New(sha256.New, sess.Key[:])
|
|
mac.Write(pkt[0:28])
|
|
mac.Write(payload)
|
|
copy(pkt[28:32], mac.Sum(nil)[:4])
|
|
return pkt
|
|
}
|
|
|
|
func hexByte(hi, lo byte) byte {
|
|
h := func(c byte) byte {
|
|
switch {
|
|
case c >= '0' && c <= '9':
|
|
return c - '0'
|
|
case c >= 'a' && c <= 'f':
|
|
return c - 'a' + 10
|
|
}
|
|
return 0
|
|
}
|
|
return h(hi)<<4 | h(lo)
|
|
}
|