Files
echolot/server/internal/dataplane/udp.go
T
mrambossekandClaude Fable 5 0c5b021b63
server-release / image (push) Successful in 14s
server-test / test (push) Successful in 30s
server-release / release (push) Successful in 30s
compat: SemVer version windows between app and server
Both sides now declare what they will talk to, and enforce it. Two axes kept
deliberately separate, because conflating them is the trap:

  protocol_version  — CAN these builds talk. The correctness axis. Below 1.0.0
                      the minor is the breaking axis, per SemVer §4.
  release window    — MAY they, per policy. [min, max), advertised in the
                      profile, overridable by the operator.

The server refuses out-of-window apps with 426 and a body naming both versions
and the accepted range; the app checks the profile in both directions before a
run rather than discovering mid-measurement that it will be refused.

Three rules that shape the rest:

  - GET /v1/profile is never gated. It is where a refused client learns which
    version it needs; gating it leaves the user with a network error instead of
    an answer, which is precisely the confusion this exists to remove.
  - An unparseable or absent version is "unknown", and is allowed. Development
    builds report "dev", and a client too old to send the header cannot be
    identified anyway.
  - Bounds sit at breaking boundaries, not at releases, so shipping a patch
    never requires editing a range. The app's server minimum is 0.4.2 for a
    stated reason: earlier multi-homed servers mis-addressed granted sends and
    the client measured 100% downstream loss that never happened.

The app's versionCode is now derived from its SemVer instead of being a second
number someone has to remember to bump.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-08-01 11:36:34 +02:00

266 lines
8.7 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.
// Skeleton scope: ECHO_REQ/ECHO_RESP and TIMESYNC only; trains, MTU probes
// and delayed echo land with the corresponding client tests.
package dataplane
import (
"crypto/hmac"
"crypto/sha256"
"encoding/binary"
"fmt"
"log/slog"
"net"
"net/netip"
"sync"
"time"
"echo-lot.app/server/internal/session"
)
const (
Magic = "ELT1"
HeaderSize = 32
TypeEchoReq = 0x01
TypeEchoResp = 0x02
TypeTimesyncReq = 0x07
TypeTimesyncRsp = 0x08
TypeMtuProbe = 0x09
TypeMtuAck = 0x0A
TypeDelayedEcho = 0x0B
// Server->client under an asymmetric grant (spec §3.4/§5).
TypeDownTrainData = 0x06
TypeBigSend = 0x0C
)
type Server struct {
Sessions *session.Manager
// 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
}
// 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()
buf := make([]byte, 65535)
for {
n, raddr, err := conn.ReadFromUDPAddrPort(buf)
if err != nil {
return err
}
tRx := time.Since(s.start).Nanoseconds()
s.handle(conn, raddr, buf[:n], tRx)
}
}
// 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) {
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
}
if !sess.CheckSeq(seq) {
return
}
sess.NoteDataSource(raddr)
if la, ok := conn.LocalAddr().(*net.UDPAddr); ok {
sess.NoteDataLocal(la.AddrPort())
}
sess.RecordUDP(session.UDPObservation{
Seq: seq, TRxNs: tRxNs, TTxNs: time.Since(s.start).Nanoseconds(),
Src: raddr.String(), Size: len(pkt), Type: typ,
})
switch typ {
case TypeEchoReq:
s.echoResp(conn, raddr, sess, pkt, seq, tRxNs)
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)
}
}
// 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 yet; needs recvmsg cmsgs)
// 35 1 received DSCP/ECN byte (0xFF = not observed)
// 36 4 received size
func observation(tRxNs, tTxNs int64, src netip.AddrPort, rcvd int) []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] = 0xFF, 0xFF
binary.BigEndian.PutUint32(b[36:40], uint32(rcvd))
return b
}
// 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) {
obs := observation(tRxNs, time.Since(s.start).Nanoseconds(), raddr, len(req))
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 := 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])
_, err := conn.WriteToUDPAddrPort(pkt, raddr)
return err
}
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)
}