server: DF-mode big_send + uploaded-run storage with an operator policy

big_send now forces the Don't-Fragment bit for the whole burst by default, so
the largest size that arrives IS the downstream path MTU rather than "fragments
got through" — two different measurements the schema already separates. Sizes
above our own egress MTU (from the startup self-test) are refused up front and
reported as max_df_bytes, because absence caused by our kernel must not be read
as a limit of the client's path.

Uploads: one JSON file per run under the state dir, with the policy the operator
actually cares about — who may upload (off / anonymous / account), how large,
how long to keep, and the least anonymization accepted. The profile advertises
all of it so the app can present the switch honestly instead of discovering the
rules by failing. `account` refuses today rather than falling back to anonymous:
picking the strict setting before OIDC lands must not silently mean the loose one.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
mrambossek
2026-08-01 10:26:19 +02:00
co-authored by Claude Fable 5
parent 7e1015c211
commit 2521d39989
19 changed files with 1172 additions and 31 deletions
@@ -0,0 +1,73 @@
// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
package app.echo_lot.engine
import app.echo_lot.protocol.ControlClient
import app.echo_lot.protocol.ProbeSession
import app.echo_lot.protocol.Wire
import kotlin.test.Test
import kotlin.test.assertTrue
/**
* Exercises the server's §5 granted sends against a LIVE server: downtrain (downstream loss /
* ordering) and big_send (downstream MTU). Self-skips without ECHOLOT_LIVE_*.
*
* This is the direction a client cannot measure alone — only the far end can push large or
* numerous packets toward it — so it is also the direction that needs the anti-amplification
* grant, and this test is the proof that the grant path works end to end.
*/
class LiveGrantedTest {
private val url = System.getenv("ECHOLOT_LIVE_URL")
private val pin = System.getenv("ECHOLOT_LIVE_PIN")
private val cred = System.getenv("ECHOLOT_LIVE_CRED")
private val udp = System.getenv("ECHOLOT_LIVE_UDP")
private val target = System.getenv("ECHOLOT_LIVE_TARGET") ?: "fmr"
@Test
fun downstreamTrainAndBigSend() {
if (url == null || pin == null || cred == null || udp == null) {
println("LiveGrantedTest skipped (no ECHOLOT_LIVE_* env)"); return
}
val control = ControlClient(url, setOf(pin))
val profile = control.profile(cred)
println("capabilities: ${profile.capabilities}")
val session = control.createSession(cred, target)
val (host, port) = udp.split(":").let { it[0] to it[1].toInt() }
ProbeSession(cred, session, host, port).use { ps ->
// The grant is bound to the OBSERVED data-plane source, so we must be seen first.
val echo = ps.echo()
println("primed with echo rtt=${echo?.rttMs}")
// --- downtrain: 50 packets of 300 bytes, 5ms apart ---
val dtResp = control.action(
cred, session.sessionId,
"""{"action":"downtrain","count":50,"size_bytes":300,"interval_us":5000}""",
)
println("downtrain accepted: ${dtResp.take(160)}")
val down = ps.collectGranted(windowMs = 4000)
.filter { it.type == Wire.TYPE_DOWNTRAIN_DATA }
val seqs = down.map { it.seq }.toSet()
println("downtrain received ${down.size}/50 packets, distinct seqs=${seqs.size}, " +
"sizes=${down.map { it.sizeBytes }.distinct()}")
assertTrue(down.isNotEmpty(), "no DOWNTRAIN_DATA arrived — granted send path is broken")
// --- big_send: which downstream sizes survive? ---
val sizes = listOf(600, 1200, 1400, 1472, 1500, 2000, 4000)
val bsResp = control.action(
cred, session.sessionId,
"""{"action":"big_send","sizes_bytes":${sizes}}""",
)
println("big_send accepted: ${bsResp.take(160)}")
val big = ps.collectGranted(windowMs = 4000)
.filter { it.type == Wire.TYPE_BIG_SEND }
val arrived = big.map { it.sizeBytes }.sorted()
println("big_send arrived sizes: $arrived (requested $sizes)")
assertTrue(big.isNotEmpty(), "no BIG_SEND packets arrived")
println("largest downstream datagram delivered: ${arrived.maxOrNull()}")
}
control.deleteSession(cred, session.sessionId)
}
}
@@ -0,0 +1,73 @@
// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
package app.echo_lot.engine
import app.echo_lot.protocol.ControlClient
import app.echo_lot.protocol.ProbeSession
import app.echo_lot.protocol.Wire
import kotlin.test.Test
import kotlin.test.assertTrue
/**
* Exercises the server's §5 granted sends against a LIVE server: downtrain (downstream loss /
* ordering) and big_send (downstream MTU). Self-skips without ECHOLOT_LIVE_*.
*
* This is the direction a client cannot measure alone — only the far end can push large or
* numerous packets toward it — so it is also the direction that needs the anti-amplification
* grant, and this test is the proof that the grant path works end to end.
*/
class LiveGrantedTest {
private val url = System.getenv("ECHOLOT_LIVE_URL")
private val pin = System.getenv("ECHOLOT_LIVE_PIN")
private val cred = System.getenv("ECHOLOT_LIVE_CRED")
private val udp = System.getenv("ECHOLOT_LIVE_UDP")
private val target = System.getenv("ECHOLOT_LIVE_TARGET") ?: "fmr"
@Test
fun downstreamTrainAndBigSend() {
if (url == null || pin == null || cred == null || udp == null) {
println("LiveGrantedTest skipped (no ECHOLOT_LIVE_* env)"); return
}
val control = ControlClient(url, setOf(pin))
val profile = control.profile(cred)
println("capabilities: ${profile.capabilities}")
val session = control.createSession(cred, target)
val (host, port) = udp.split(":").let { it[0] to it[1].toInt() }
ProbeSession(cred, session, host, port).use { ps ->
// The grant is bound to the OBSERVED data-plane source, so we must be seen first.
val echo = ps.echo()
println("primed with echo rtt=${echo?.rttMs}")
// --- downtrain: 50 packets of 300 bytes, 5ms apart ---
val dtResp = control.action(
cred, session.sessionId,
"""{"action":"downtrain","count":50,"size_bytes":300,"interval_us":5000}""",
)
println("downtrain accepted: ${dtResp.take(160)}")
val down = ps.collectGranted(windowMs = 4000)
.filter { it.type == Wire.TYPE_DOWNTRAIN_DATA }
val seqs = down.map { it.seq }.toSet()
println("downtrain received ${down.size}/50 packets, distinct seqs=${seqs.size}, " +
"sizes=${down.map { it.sizeBytes }.distinct()}")
assertTrue(down.isNotEmpty(), "no DOWNTRAIN_DATA arrived — granted send path is broken")
// --- big_send: which downstream sizes survive? ---
val sizes = listOf(600, 1200, 1400, 1472, 1500, 2000, 4000)
val bsResp = control.action(
cred, session.sessionId,
"""{"action":"big_send","sizes_bytes":${sizes}}""",
)
println("big_send accepted: ${bsResp.take(160)}")
val big = ps.collectGranted(windowMs = 4000)
.filter { it.type == Wire.TYPE_BIG_SEND }
val arrived = big.map { it.sizeBytes }.sorted()
println("big_send arrived sizes: $arrived (requested $sizes)")
assertTrue(big.isNotEmpty(), "no BIG_SEND packets arrived")
println("largest downstream datagram delivered: ${arrived.maxOrNull()}")
}
control.deleteSession(cred, session.sessionId)
}
}
@@ -80,6 +80,19 @@ class ControlClient(private val controlUrl: String, pins: Set<String>) {
return json.decodeFromString(SessionResponse.serializer(), body(conn))
}
/**
* Requests a §5 action. The server creates an asymmetric grant for the granted ones
* (downtrain / big_send) and starts sending toward the session's observed data-plane source,
* so the caller must already have sent at least one ECHO. Returns the raw JSON reply.
*/
fun action(credential: String, sessionId: String, bodyJson: String): String {
val conn = open("/v1/sessions/$sessionId/actions", "POST", credential)
writeJson(conn, bodyJson)
val body = body(conn)
check(conn.responseCode in 200..299) { "action failed: ${conn.responseCode} $body" }
return body
}
fun observations(credential: String, sessionId: String): String {
val conn = open("/v1/sessions/$sessionId/observations", "GET", credential)
check(conn.responseCode == 200) { "observations failed: ${conn.responseCode}" }
@@ -60,6 +60,40 @@ class ProbeSession(
(resp.payload[3].toInt() and 0xFF)
}
/**
* Collects packets the SERVER sends under a grant (downtrain / big_send) for [windowMs].
* These arrive unsolicited after a control-plane action, so this just drains the socket and
* keeps every HMAC-verified packet — anything that fails verification is not ours and is
* silently ignored (an injected packet must not be able to fake a measurement).
*/
fun collectGranted(windowMs: Long): List<Received> {
val out = ArrayList<Received>()
val deadline = System.nanoTime() + windowMs * 1_000_000
val buf = ByteArray(9200)
val prevTimeout = socket.soTimeout
try {
while (System.nanoTime() < deadline) {
val remainMs = ((deadline - System.nanoTime()) / 1_000_000).toInt()
if (remainMs <= 0) break
socket.soTimeout = remainMs.coerceAtMost(2000)
val dp = DatagramPacket(buf, buf.size)
try {
socket.receive(dp)
} catch (e: java.net.SocketTimeoutException) {
continue
}
val pkt = Wire.parseVerified(buf, dp.length, key) ?: continue
out.add(Received(pkt.type, pkt.seq, dp.length, (System.nanoTime() - epochNanos)))
}
} finally {
socket.soTimeout = prevTimeout
}
return out
}
/** One packet received from the server, with the wire size actually delivered. */
data class Received(val type: Int, val seq: Int, val sizeBytes: Int, val tRxNs: Long)
private fun receive(wantType: Int): Wire.Packet? {
val buf = ByteArray(2048)
return try {
@@ -28,6 +28,9 @@ object Wire {
const val TYPE_MTU_PROBE: Int = 0x09
const val TYPE_MTU_ACK: Int = 0x0A
const val TYPE_DELAYED_ECHO: Int = 0x0B
/** Server->client under an asymmetric grant (spec §3.4/§5). */
const val TYPE_DOWNTRAIN_DATA: Int = 0x06
const val TYPE_BIG_SEND: Int = 0x0C
/** The 8-byte on-the-wire prefix = first 16 hex chars of the session id, decoded. */
fun wirePrefix(sessionId: String): ByteArray {
@@ -80,6 +80,19 @@ class ControlClient(private val controlUrl: String, pins: Set<String>) {
return json.decodeFromString(SessionResponse.serializer(), body(conn))
}
/**
* Requests a §5 action. The server creates an asymmetric grant for the granted ones
* (downtrain / big_send) and starts sending toward the session's observed data-plane source,
* so the caller must already have sent at least one ECHO. Returns the raw JSON reply.
*/
fun action(credential: String, sessionId: String, bodyJson: String): String {
val conn = open("/v1/sessions/$sessionId/actions", "POST", credential)
writeJson(conn, bodyJson)
val body = body(conn)
check(conn.responseCode in 200..299) { "action failed: ${conn.responseCode} $body" }
return body
}
fun observations(credential: String, sessionId: String): String {
val conn = open("/v1/sessions/$sessionId/observations", "GET", credential)
check(conn.responseCode == 200) { "observations failed: ${conn.responseCode}" }
@@ -60,6 +60,40 @@ class ProbeSession(
(resp.payload[3].toInt() and 0xFF)
}
/**
* Collects packets the SERVER sends under a grant (downtrain / big_send) for [windowMs].
* These arrive unsolicited after a control-plane action, so this just drains the socket and
* keeps every HMAC-verified packet — anything that fails verification is not ours and is
* silently ignored (an injected packet must not be able to fake a measurement).
*/
fun collectGranted(windowMs: Long): List<Received> {
val out = ArrayList<Received>()
val deadline = System.nanoTime() + windowMs * 1_000_000
val buf = ByteArray(9200)
val prevTimeout = socket.soTimeout
try {
while (System.nanoTime() < deadline) {
val remainMs = ((deadline - System.nanoTime()) / 1_000_000).toInt()
if (remainMs <= 0) break
socket.soTimeout = remainMs.coerceAtMost(2000)
val dp = DatagramPacket(buf, buf.size)
try {
socket.receive(dp)
} catch (e: java.net.SocketTimeoutException) {
continue
}
val pkt = Wire.parseVerified(buf, dp.length, key) ?: continue
out.add(Received(pkt.type, pkt.seq, dp.length, (System.nanoTime() - epochNanos)))
}
} finally {
socket.soTimeout = prevTimeout
}
return out
}
/** One packet received from the server, with the wire size actually delivered. */
data class Received(val type: Int, val seq: Int, val sizeBytes: Int, val tRxNs: Long)
private fun receive(wantType: Int): Wire.Packet? {
val buf = ByteArray(2048)
return try {
@@ -28,6 +28,9 @@ object Wire {
const val TYPE_MTU_PROBE: Int = 0x09
const val TYPE_MTU_ACK: Int = 0x0A
const val TYPE_DELAYED_ECHO: Int = 0x0B
/** Server->client under an asymmetric grant (spec §3.4/§5). */
const val TYPE_DOWNTRAIN_DATA: Int = 0x06
const val TYPE_BIG_SEND: Int = 0x0C
/** The 8-byte on-the-wire prefix = first 16 hex chars of the session id, decoded. */
fun wirePrefix(sessionId: String): ByteArray {
+7 -4
View File
@@ -8,7 +8,8 @@
# whole lot to the Gradle test. Proves the Kotlin client talks to the real
# server over the wire.
#
# Usage: JAVA_HOME=... echolot-app/scripts/test-fmr.sh
# Usage: JAVA_HOME=... echolot-app/scripts/test-fmr.sh [gradle-task] [test-filter]
# e.g. ... test-fmr.sh :core-engine:test '*LiveGrantedTest*'
set -euo pipefail
SSH_HOST="${ECHOLOT_SSH:-claude-echolot}"
@@ -33,12 +34,14 @@ PIN=$(echo | openssl s_client -connect "${CTL_HOST}:${CTL_PORT}" 2>/dev/null \
| openssl dgst -sha256 -binary | openssl base64)
echo "· pin=${PIN}"
echo "· running LiveServerTest ..."
TASK="${1:-:core-protocol:test}"
FILTER="${2:-*LiveServerTest*}"
echo "· running ${TASK} ${FILTER} ..."
cd "$(dirname "$0")/.."
ECHOLOT_LIVE_URL="$CTL_URL" \
ECHOLOT_LIVE_PIN="$PIN" \
ECHOLOT_LIVE_CRED="$CRED" \
ECHOLOT_LIVE_UDP="${CTL_HOST}:${UDP_PORT}" \
ECHOLOT_LIVE_TARGET="${ECHOLOT_LIVE_TARGET:-fmr}" \
./gradlew :core-protocol:test --tests '*LiveServerTest*' --info --rerun-tasks --console=plain \
2>&1 | grep -E "profile:|session:|echo |mtu probe|observations bytes|LiveServerTest|BUILD|FAIL|PASS" || true
./gradlew "$TASK" --tests "$FILTER" --info --rerun-tasks --console=plain \
2>&1 | grep -E "profile:|capabilities:|session:|echo |primed|mtu probe|downtrain|big_send|largest|observations bytes|Live[A-Za-z]*Test|BUILD|FAIL|PASS|^e:" || true
+30
View File
@@ -39,6 +39,7 @@ import (
"echo-lot.app/server/internal/config"
"echo-lot.app/server/internal/control"
"echo-lot.app/server/internal/dataplane"
"echo-lot.app/server/internal/runs"
"echo-lot.app/server/internal/selftest"
"echo-lot.app/server/internal/selfupdate"
"echo-lot.app/server/internal/session"
@@ -113,6 +114,23 @@ func serve(cfg *config.Config) error {
caps = append(caps, "tcp-echo", "tls-echo")
}
// Uploaded-run storage. A failure here is not fatal: measurement still works, uploads just
// stay unavailable and say so in the profile.
runStore, err := runs.Open(cfg.StateDir, runs.Policy{
Mode: runs.Mode(cfg.UploadsMode),
MaxBytes: cfg.UploadMaxBytes,
RetentionDays: cfg.UploadRetentionDays,
MaxRunsPerDevice: cfg.UploadMaxRuns,
MinAnonymization: cfg.UploadMinAnon,
})
if err != nil {
slog.Warn("uploaded-run storage unavailable — uploads disabled", "err", err)
runStore = nil
} else {
slog.Info("uploads", "mode", cfg.UploadsMode, "min_anonymization", cfg.UploadMinAnon,
"retention_days", cfg.UploadRetentionDays, "max_runs_per_device", cfg.UploadMaxRuns)
}
ctl := &control.Server{
Store: st, Sessions: sessions, Name: cfg.Name,
UDPPort: mustPort(firstAddr(cfg.UDPListen)), TCPPort: mustPort(firstAddr(cfg.TCPListen)),
@@ -121,6 +139,7 @@ func serve(cfg *config.Config) error {
DownTrain: dp.DownTrain,
BigSend: dp.BigSend,
TCPRecent: func(ip string) any { return tcpSrv.RecentFor(ip) },
Runs: runStore,
}
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
@@ -172,6 +191,17 @@ func serve(cfg *config.Config) error {
r := selftestPtr.Load()
return r.MTUOK, r.SysctlOK
}
// The smallest egress MTU we measured is the ceiling for DF-mode big_send: above it our own
// kernel refuses the datagram, which would otherwise look like a downstream path limit.
ctl.EgressMTU = func() int {
best := 0
for _, m := range selftestPtr.Load().EgressMTU {
if m.DiscoveredMTU > 0 && (best == 0 || m.DiscoveredMTU < best) {
best = m.DiscoveredMTU
}
}
return best
}
// Admin/health (plain HTTP, localhost by default; spec §7)
admin := http.NewServeMux()
+24
View File
@@ -11,6 +11,7 @@ import (
"flag"
"fmt"
"os"
"strconv"
"strings"
)
@@ -49,10 +50,28 @@ type Config struct {
// e.g. https://git.example.net/api/v1/repos/owner/repo
SelfUpdateAPI string // ECHOLOT_SELF_UPDATE_API / --self-update-api
// Uploaded-run storage. The default is "anonymous": any enrolled device may upload,
// which is what a self-hosted server wants. Operators of shared servers turn it down.
UploadsMode string // ECHOLOT_UPLOADS / --uploads (off|anonymous|account)
UploadMaxBytes int64 // ECHOLOT_UPLOAD_MAX_BYTES / --upload-max-bytes
UploadRetentionDays int // ECHOLOT_UPLOAD_RETENTION_DAYS / --upload-retention-days
UploadMaxRuns int // ECHOLOT_UPLOAD_MAX_RUNS / --upload-max-runs (per device)
UploadMinAnon string // ECHOLOT_UPLOAD_MIN_ANONYMIZATION / --upload-min-anonymization
// Mode
Docker bool // --docker (or autodetected; env ECHOLOT_DOCKER=1 forces)
}
// envInt reads ECHOLOT_<key> as an integer with a fallback.
func envInt(key string, def int) int {
if v := envOr(key, ""); v != "" {
if n, err := strconv.Atoi(v); err == nil {
return n
}
}
return def
}
// envOr reads ECHOLOT_<key> with a fallback.
func envOr(key, def string) string {
if v, ok := os.LookupEnv("ECHOLOT_" + key); ok {
@@ -82,6 +101,11 @@ func Load(args []string) (*Config, *Actions, error) {
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.SelfUpdateAPI, "self-update-api", envOr("SELF_UPDATE_API", ""), "Gitea repo API base for self-update; empty disables")
fs.StringVar(&c.UploadsMode, "uploads", envOr("UPLOADS", "anonymous"), "who may upload measurement runs: off|anonymous|account")
fs.Int64Var(&c.UploadMaxBytes, "upload-max-bytes", int64(envInt("UPLOAD_MAX_BYTES", 4<<20)), "largest accepted uploaded run, bytes")
fs.IntVar(&c.UploadRetentionDays, "upload-retention-days", envInt("UPLOAD_RETENTION_DAYS", 90), "delete uploaded runs older than this; 0 disables")
fs.IntVar(&c.UploadMaxRuns, "upload-max-runs", envInt("UPLOAD_MAX_RUNS", 200), "keep at most this many runs per device; 0 disables")
fs.StringVar(&c.UploadMinAnon, "upload-min-anonymization", envOr("UPLOAD_MIN_ANONYMIZATION", "full"), "least anonymization accepted: full|balanced|strict")
fs.BoolVar(&c.Docker, "docker", envOr("DOCKER", "") == "1", "force container mode (config from env, no systemd/self-update)")
fs.BoolVar(&a.InstallSystemd, "install-systemd", false, "install a systemd unit for this binary and exit")
@@ -0,0 +1,52 @@
// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
package control
import (
"net/netip"
"testing"
"time"
"echo-lot.app/server/internal/session"
)
// The DF ceiling is the difference between "the client's path cannot carry this" and "we could
// never have sent it in the first place". Getting the header arithmetic wrong would silently
// attribute a server limit to the client's network, so it is pinned here.
func TestMaxDFPayload(t *testing.T) {
mgr := session.NewManager(time.Minute)
newSess := func(src string) *session.Session {
s, _, err := mgr.New("dev", "cred", netip.MustParseAddr("203.0.113.9"))
if err != nil {
t.Fatalf("new session: %v", err)
}
if src != "" {
s.NoteDataSource(netip.MustParseAddrPort(src))
}
return s
}
cases := []struct {
name string
mtu func() int
src string
want int
}{
{"no egress mtu hook means no clamp", nil, "198.51.100.4:5000", 0},
{"unknown egress mtu means no clamp", func() int { return 0 }, "198.51.100.4:5000", 0},
{"ipv4 subtracts ip+udp", func() int { return 1500 }, "198.51.100.4:5000", 1472},
{"ipv6 subtracts the larger header", func() int { return 1500 }, "[2001:db8::4]:5000", 1452},
{"pppoe-style 1492 egress", func() int { return 1492 }, "198.51.100.4:5000", 1464},
{"no data source yet falls back to ipv4 overhead", func() int { return 1500 }, "", 1472},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
s := &Server{EgressMTU: tc.mtu}
if got := s.maxDFPayload(newSess(tc.src)); got != tc.want {
t.Fatalf("maxDFPayload = %d, want %d", got, tc.want)
}
})
}
}
+167 -4
View File
@@ -15,6 +15,7 @@ import (
"encoding/hex"
"encoding/json"
"errors"
"io"
"log/slog"
"net"
"net/http"
@@ -23,6 +24,8 @@ import (
"strings"
"time"
"echo-lot.app/server/internal/dataplane"
"echo-lot.app/server/internal/runs"
"echo-lot.app/server/internal/session"
"echo-lot.app/server/internal/store"
)
@@ -51,7 +54,13 @@ type Server struct {
DelayedEcho func(sess *session.Session, actionID string) error
// Granted server->client sends (spec §5). Both consume an asymmetric grant.
DownTrain func(sess *session.Session, g *session.Grant, count, sizeBytes, intervalUs int) (int, error)
BigSend func(sess *session.Session, g *session.Grant, sizes []int) ([]int, error)
BigSend func(sess *session.Session, g *session.Grant, sizes []int, df bool) ([]dataplane.BigSendResult, error)
// Runs stores uploaded measurement documents (may be nil: uploads unsupported).
Runs *runs.Store
// EgressMTU reports the server's own measured egress path MTU (0 = unknown). With DF set
// we cannot emit a datagram larger than this, so requested sizes above it are refused up
// front and reported as such — the client must not read that as a downstream path limit.
EgressMTU func() int
// 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.
@@ -73,6 +82,10 @@ func (s *Server) Handler() http.Handler {
mux.HandleFunc("POST /v1/sessions/{id}/actions", s.actions)
mux.HandleFunc("POST /v1/echo", s.httpEcho)
mux.HandleFunc("GET /v1/tls-reference", s.tlsReference)
mux.HandleFunc("POST /v1/runs", s.uploadRun)
mux.HandleFunc("GET /v1/runs", s.listRuns)
mux.HandleFunc("GET /v1/runs/{id}", s.getRun)
mux.HandleFunc("DELETE /v1/runs/{id}", s.deleteRun)
// TODO(spec §5): frag_send, throughput (both build on the same grant machinery)
return mux
}
@@ -156,6 +169,7 @@ func (s *Server) actions(w http.ResponseWriter, r *http.Request) {
SizeBytes int `json:"size_bytes"`
IntervalUs int `json:"interval_us"`
SizesBytes []int `json:"sizes_bytes"`
DF *bool `json:"df"`
}
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "bad body"})
@@ -241,6 +255,35 @@ func (s *Server) actions(w http.ResponseWriter, r *http.Request) {
if len(sizes) > 32 {
sizes = sizes[:32]
}
// DF on by default: an unfragmented burst is what makes the result a path-MTU
// measurement rather than a fragment-delivery one. Callers opt out explicitly.
df := true
if req.DF != nil {
df = *req.DF
}
// With DF we can only emit up to our own egress MTU minus IP+UDP headers. Drop the
// rest here and say so, rather than sending nothing and letting the client blame
// the path.
maxDF := 0
if df {
maxDF = s.maxDFPayload(sess)
if maxDF > 0 {
kept := sizes[:0]
for _, x := range sizes {
if x <= maxDF {
kept = append(kept, x)
}
}
sizes = kept
}
}
if len(sizes) == 0 {
writeJSON(w, http.StatusBadRequest, map[string]any{
"error": "every requested size exceeds the server's own egress MTU with DF set",
"max_df_bytes": maxDF,
})
return
}
total := 0
for _, x := range sizes {
total += clamp(x, dataMinPacket, 9000)
@@ -251,11 +294,11 @@ func (s *Server) actions(w http.ResponseWriter, r *http.Request) {
return
}
go func() {
attempted, err := s.BigSend(sess, g, sizes)
slog.Info("big_send finished", "action", actionID, "attempted", attempted, "err", err)
results, err := s.BigSend(sess, g, sizes, df)
slog.Info("big_send finished", "action", actionID, "results", results, "df", df, "err", err)
}()
writeJSON(w, http.StatusAccepted, map[string]any{
"action_id": actionID, "sizes_bytes": sizes,
"action_id": actionID, "sizes_bytes": sizes, "df": df, "max_df_bytes": maxDF,
"grant": map[string]any{"max_bytes": g.MaxBytes, "max_kbps": g.MaxKbps},
})
@@ -264,6 +307,24 @@ func (s *Server) actions(w http.ResponseWriter, r *http.Request) {
}
}
// maxDFPayload is the largest UDP payload the server can emit toward this session without
// fragmenting: its own egress MTU less the IP and UDP headers of the session's address family.
// Returns 0 when the egress MTU is unknown, meaning "do not clamp".
func (s *Server) maxDFPayload(sess *session.Session) int {
if s.EgressMTU == nil {
return 0
}
mtu := s.EgressMTU()
if mtu <= 0 {
return 0
}
overhead := 28 // IPv4 (20) + UDP (8)
if src := sess.DataSource(); src.IsValid() && !src.Addr().Unmap().Is4() {
overhead = 48 // IPv6 (40) + UDP (8)
}
return mtu - overhead
}
// dataMinPacket is the smallest datagram that still carries a header + a little payload.
const dataMinPacket = 40
@@ -360,6 +421,9 @@ func (s *Server) profile(w http.ResponseWriter, r *http.Request) {
"canary_zone": s.CanaryZone,
"server_selftest": selftestSignal(s.ProvenGood),
"limits": map[string]any{"max_kbps": 50000, "max_session_s": 900},
// The app needs the upload rules before it offers the switch: whether uploads are
// accepted at all, and how much identifying detail it must strip first.
"uploads": s.uploadPolicy(),
})
}
@@ -411,3 +475,102 @@ func SpkiPinB64(cert tls.Certificate) (string, error) {
sum := sha256.Sum256(leaf.RawSubjectPublicKeyInfo)
return base64.StdEncoding.EncodeToString(sum[:]), nil
}
// uploadPolicy is the profile's advertisement of the operator's upload rules.
func (s *Server) uploadPolicy() map[string]any {
if s.Runs == nil {
return map[string]any{"mode": string(runs.ModeOff), "reason": "not configured"}
}
p := s.Runs.Policy()
return map[string]any{
"mode": string(p.Mode),
"max_bytes": p.MaxBytes,
"retention_days": p.RetentionDays,
"max_runs_per_device": p.MaxRunsPerDevice,
"min_anonymization": p.MinAnonymization,
}
}
// uploadRun stores one measurement document for the calling device.
func (s *Server) uploadRun(w http.ResponseWriter, r *http.Request) {
dev := s.Store.DeviceByCredential(bearer(r))
if dev == nil {
writeJSON(w, http.StatusUnauthorized, map[string]string{"error": "unknown credential"})
return
}
if s.Runs == nil {
writeJSON(w, http.StatusForbidden, map[string]string{"error": runs.ErrDisabled.Error()})
return
}
limit := s.Runs.Policy().MaxBytes
if limit <= 0 {
limit = 4 << 20
}
// +1 so a body exactly at the limit is distinguishable from one over it.
body, err := io.ReadAll(io.LimitReader(r.Body, limit+1))
if err != nil {
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "read failed"})
return
}
meta, err := s.Runs.Put(dev.ID, body)
switch {
case err == nil:
slog.Info("run uploaded", "device", dev.ID, "run", meta.ID,
"bytes", meta.SizeBytes, "anon", meta.Anonymization, "findings", meta.FindingCount)
writeJSON(w, http.StatusCreated, meta)
case errors.Is(err, runs.ErrDisabled), errors.Is(err, runs.ErrNeedAccount),
errors.Is(err, runs.ErrNotAnonEnough):
writeJSON(w, http.StatusForbidden, map[string]any{
"error": err.Error(), "uploads": s.uploadPolicy(),
})
case errors.Is(err, runs.ErrTooLarge):
writeJSON(w, http.StatusRequestEntityTooLarge, map[string]any{
"error": err.Error(), "max_bytes": s.Runs.Policy().MaxBytes,
})
default:
writeJSON(w, http.StatusBadRequest, map[string]string{"error": err.Error()})
}
}
func (s *Server) listRuns(w http.ResponseWriter, r *http.Request) {
dev := s.Store.DeviceByCredential(bearer(r))
if dev == nil || s.Runs == nil {
writeJSON(w, http.StatusUnauthorized, map[string]string{"error": "unknown credential"})
return
}
list := s.Runs.List(dev.ID)
if list == nil {
list = []runs.Meta{}
}
writeJSON(w, http.StatusOK, map[string]any{"runs": list})
}
func (s *Server) getRun(w http.ResponseWriter, r *http.Request) {
dev := s.Store.DeviceByCredential(bearer(r))
if dev == nil || s.Runs == nil {
writeJSON(w, http.StatusUnauthorized, map[string]string{"error": "unknown credential"})
return
}
// Scoped to the calling device's own directory: one device cannot read another's runs by
// guessing a run id.
b, err := s.Runs.Get(dev.ID, r.PathValue("id"))
if err != nil {
writeJSON(w, http.StatusNotFound, map[string]string{"error": "no such run"})
return
}
w.Header().Set("Content-Type", "application/json")
_, _ = w.Write(b)
}
func (s *Server) deleteRun(w http.ResponseWriter, r *http.Request) {
dev := s.Store.DeviceByCredential(bearer(r))
if dev == nil || s.Runs == nil {
writeJSON(w, http.StatusUnauthorized, map[string]string{"error": "unknown credential"})
return
}
if err := s.Runs.Delete(dev.ID, r.PathValue("id")); err != nil {
writeJSON(w, http.StatusInternalServerError, map[string]string{"error": err.Error()})
return
}
w.WriteHeader(http.StatusNoContent)
}
+61
View File
@@ -0,0 +1,61 @@
// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
//go:build linux
package dataplane
import (
"net"
"syscall"
)
// PMTUD socket-option values. Go's syscall package exports IP_MTU_DISCOVER and
// IPV6_MTU_DISCOVER but not the IP_PMTUDISC_* values, so they are spelled out
// here (include/linux/in.h, in6.h — stable ABI, same reasoning as the client's
// OsAbi.kt).
const (
pmtudiscWant = 0 // per-route default: fragment locally when needed
pmtudiscDo = 2 // always set DF: oversized sends fail with EMSGSIZE, never fragment
)
// withDF runs fn with the Don't-Fragment bit forced on for conn, then puts the
// socket back the way it was found.
//
// The socket is shared by every session on that address family, so the caller
// must hold Server.dfMu: a concurrent big_send must not silently ride along
// with someone else's DF window (or, worse, clear it mid-flight).
func withDF(conn *net.UDPConn, fn func() error) error {
raw, err := conn.SyscallConn()
if err != nil {
return err
}
v4 := conn.LocalAddr().(*net.UDPAddr).IP.To4() != nil
level, opt := syscall.IPPROTO_IPV6, syscall.IPV6_MTU_DISCOVER
if v4 {
level, opt = syscall.IPPROTO_IP, syscall.IP_MTU_DISCOVER
}
var setErr error
prev := pmtudiscWant
if err := raw.Control(func(fd uintptr) {
if p, e := syscall.GetsockoptInt(int(fd), level, opt); e == nil {
prev = p
}
setErr = syscall.SetsockoptInt(int(fd), level, opt, pmtudiscDo)
}); err != nil {
return err
}
if setErr != nil {
return setErr
}
defer func() {
_ = raw.Control(func(fd uintptr) {
_ = syscall.SetsockoptInt(int(fd), level, opt, prev)
})
}()
return fn()
}
// dfSupported reports whether withDF can actually set the DF bit here.
const dfSupported = true
+16
View File
@@ -0,0 +1,16 @@
// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
//go:build !linux
package dataplane
import "net"
// Forcing DF per-socket is Linux-specific (IP_MTU_DISCOVER). Off Linux the
// send still happens — just without the guarantee that nothing fragmented it,
// so the caller must report the result as fragment-delivery evidence rather
// than a path-MTU measurement. See dfSupported.
func withDF(conn *net.UDPConn, fn func() error) error { return fn() }
const dfSupported = false
+45 -8
View File
@@ -52,10 +52,25 @@ func (s *Server) DownTrain(sess *session.Session, g *session.Grant, count, sizeB
return sent, nil
}
// BigSend transmits one datagram per requested size, largest-first metadata intact, so the client
// can see which sizes survive the *downstream* path — the mtu.pmtud_down / mtu.blackhole evidence.
// The client cannot produce this itself: only the far end can emit a large packet toward it.
func (s *Server) BigSend(sess *session.Session, g *session.Grant, sizes []int) ([]int, error) {
// BigSendResult records what happened to one requested size. `Sent` false with an EMSGSIZE-ish
// Err means *we* could not put it on the wire (the datagram exceeds our own egress MTU with DF
// set) — the client must not read its absence as a path limit, so this is reported, not hidden.
type BigSendResult struct {
SizeBytes int `json:"size_bytes"`
Seq int `json:"seq"`
Sent bool `json:"sent"`
Err string `json:"err,omitempty"`
}
// BigSend transmits one datagram per requested size so the client can see which sizes survive the
// *downstream* path — the mtu.pmtud_down / mtu.frag_delivery evidence. The client cannot produce
// this itself: only the far end can emit a large packet toward it.
//
// With df set, the DF bit is forced for the whole burst, so nothing fragments and the largest
// size that arrives IS the downstream path MTU. Without it, the kernel fragments freely and the
// result only says whether fragments get through — a different (also useful) measurement, and
// the reason the two are separate test types.
func (s *Server) BigSend(sess *session.Session, g *session.Grant, sizes []int, df bool) ([]BigSendResult, error) {
target := sess.DataSource()
if !target.IsValid() {
return nil, fmt.Errorf("no observed data-plane source")
@@ -64,7 +79,9 @@ func (s *Server) BigSend(sess *session.Session, g *session.Grant, sizes []int) (
if conn == nil {
return nil, fmt.Errorf("no data-plane socket matches target family")
}
attempted := make([]int, 0, len(sizes))
results := make([]BigSendResult, 0, len(sizes))
burst := func() error {
for i, size := range sizes {
if size < HeaderSize+8 {
size = HeaderSize + 8
@@ -79,9 +96,29 @@ func (s *Server) BigSend(sess *session.Session, g *session.Grant, sizes []int) (
// Echo the intended size into the payload so a truncated/fragmented arrival is
// still attributable to the size we meant to send.
binary.BigEndian.PutUint32(payload[0:4], uint32(size))
s.send(conn, target, sess, TypeBigSend, uint32(i), payload)
attempted = append(attempted, size)
err := s.sendErr(conn, target, sess, TypeBigSend, uint32(i), payload)
results = append(results, BigSendResult{
SizeBytes: size, Seq: i, Sent: err == nil, Err: errString(err),
})
time.Sleep(20 * time.Millisecond) // keep bursts from being read as congestion loss
}
return attempted, nil
return nil
}
if df && dfSupported {
s.dfMu.Lock()
defer s.dfMu.Unlock()
if err := withDF(conn, burst); err != nil {
return results, err
}
return results, nil
}
return results, burst()
}
func errString(err error) string {
if err == nil {
return ""
}
return err.Error()
}
+13 -1
View File
@@ -45,6 +45,10 @@ type Server struct {
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.
@@ -203,6 +207,13 @@ func (s *Server) timesyncResp(conn *net.UDPConn, raddr netip.AddrPort, sess *ses
}
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
@@ -218,7 +229,8 @@ func (s *Server) send(conn *net.UDPConn, raddr netip.AddrPort, sess *session.Ses
mac.Write(pkt[0:28])
mac.Write(payload)
copy(pkt[28:32], mac.Sum(nil)[:4])
_, _ = conn.WriteToUDPAddrPort(pkt, raddr)
_, err := conn.WriteToUDPAddrPort(pkt, raddr)
return err
}
func hexByte(hi, lo byte) byte {
+300
View File
@@ -0,0 +1,300 @@
// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
// Package runs stores uploaded measurement documents.
//
// The premise (and the reason uploads exist at all) is that an engineer runs their own server:
// uploading a run there is how history, diffing and "it was fine last Tuesday" work. That makes
// the storage deliberately dumb — one JSON file per run, on disk, greppable, deletable with rm —
// and puts the interesting policy in two places instead:
//
// - who may upload (Policy.Mode), because a public server is a different proposition from a
// private one; and
// - how much identifying detail the client must strip first (Policy.MinAnonymization), because
// someone measuring against a stranger's server should not be shipping their SSIDs there.
//
// Retention is enforced on every upload, not by a sweeper, so a server left alone does not grow.
package runs
import (
"encoding/json"
"errors"
"fmt"
"os"
"path/filepath"
"sort"
"strings"
"sync"
"time"
)
// Mode says who may upload.
type Mode string
const (
// ModeOff refuses every upload. The endpoint still answers, with 403 and a reason, so the
// app can say "this server does not accept uploads" instead of showing a network error.
ModeOff Mode = "off"
// ModeAnonymous accepts uploads from any enrolled device. The default: enrollment already
// required an admin-minted token, so "anyone enrolled" is not "anyone".
ModeAnonymous Mode = "anonymous"
// ModeAccount accepts uploads only from a device tied to a signed-in account. The account
// system (OIDC) is not built yet, so today this refuses everything with a distinct reason —
// it exists so operators can pick the strict setting now and have it mean the right thing
// when accounts land, rather than silently loosening on upgrade.
ModeAccount Mode = "account"
)
// Anonymization levels, mirroring the client's redaction levels (measurement-schema.md §8).
// Ordered: full < balanced < strict.
const (
AnonFull = "full" // nothing removed — for your own server
AnonBalanced = "balanced" // network names and device identity pseudonymized, neighbours dropped
AnonStrict = "strict" // metrics and findings only
)
func anonRank(level string) int {
switch level {
case AnonStrict:
return 2
case AnonBalanced:
return 1
case AnonFull:
return 0
}
return -1 // unknown
}
// Policy is the operator's upload configuration.
type Policy struct {
Mode Mode `json:"mode"`
MaxBytes int64 `json:"max_bytes"`
RetentionDays int `json:"retention_days"`
MaxRunsPerDevice int `json:"max_runs_per_device"`
MinAnonymization string `json:"min_anonymization"`
}
func DefaultPolicy() Policy {
return Policy{
Mode: ModeAnonymous,
MaxBytes: 4 << 20,
RetentionDays: 90,
MaxRunsPerDevice: 200,
MinAnonymization: AnonFull,
}
}
var (
ErrDisabled = errors.New("uploads are disabled on this server")
ErrNeedAccount = errors.New("this server only accepts uploads from signed-in accounts")
ErrTooLarge = errors.New("run exceeds the server's upload size limit")
ErrNotAnonEnough = errors.New("run is less anonymized than this server requires")
ErrMalformed = errors.New("run is not a measurement document")
)
// Meta is the index entry for one stored run — enough to list history without opening the files.
type Meta struct {
ID string `json:"id"`
DeviceID string `json:"device_id"`
UploadedAt time.Time `json:"uploaded_at"`
StartedAt string `json:"started_at,omitempty"`
Anonymization string `json:"anonymization"`
SizeBytes int64 `json:"size_bytes"`
Verdict string `json:"verdict,omitempty"`
FindingCount int `json:"finding_count"`
}
type Store struct {
mu sync.Mutex
dir string
policy Policy
}
func Open(stateDir string, p Policy) (*Store, error) {
dir := filepath.Join(stateDir, "runs")
if err := os.MkdirAll(dir, 0o700); err != nil {
return nil, err
}
return &Store{dir: dir, policy: p}, nil
}
func (s *Store) Policy() Policy { return s.policy }
// Accepts reports whether an upload would be allowed at all, so callers can answer the
// capability question without a body.
func (s *Store) Accepts() error {
switch s.policy.Mode {
case ModeOff:
return ErrDisabled
case ModeAccount:
return ErrNeedAccount
}
return nil
}
// Put validates and stores one uploaded document. body is the raw JSON as received: it is stored
// byte-for-byte so what the device signed off on is what sits on disk.
func (s *Store) Put(deviceID string, body []byte) (Meta, error) {
if err := s.Accepts(); err != nil {
return Meta{}, err
}
if s.policy.MaxBytes > 0 && int64(len(body)) > s.policy.MaxBytes {
return Meta{}, ErrTooLarge
}
// Peek at the parts we index on. Unknown fields are ignored: the server must not become a
// second schema authority that rejects documents a newer client legitimately produces.
var doc struct {
Run struct {
ID string `json:"id"`
StartedAt string `json:"started_at"`
Privacy struct {
Anonymization string `json:"anonymization"`
} `json:"privacy"`
} `json:"run"`
Findings []json.RawMessage `json:"findings"`
Summary struct {
Verdict string `json:"verdict"`
} `json:"summary"`
}
if err := json.Unmarshal(body, &doc); err != nil || doc.Run.ID == "" {
return Meta{}, ErrMalformed
}
level := doc.Run.Privacy.Anonymization
if level == "" {
level = AnonFull // no declaration means nothing was stripped
}
if anonRank(level) < anonRank(s.policy.MinAnonymization) {
return Meta{}, fmt.Errorf("%w: got %q, need at least %q",
ErrNotAnonEnough, level, s.policy.MinAnonymization)
}
id := sanitizeID(doc.Run.ID)
if id == "" {
return Meta{}, ErrMalformed
}
s.mu.Lock()
defer s.mu.Unlock()
devDir := filepath.Join(s.dir, sanitizeID(deviceID))
if err := os.MkdirAll(devDir, 0o700); err != nil {
return Meta{}, err
}
if err := os.WriteFile(filepath.Join(devDir, id+".json"), body, 0o600); err != nil {
return Meta{}, err
}
meta := Meta{
ID: id, DeviceID: deviceID, UploadedAt: time.Now().UTC(),
StartedAt: doc.Run.StartedAt, Anonymization: level,
SizeBytes: int64(len(body)), Verdict: doc.Summary.Verdict,
FindingCount: len(doc.Findings),
}
if err := os.WriteFile(filepath.Join(devDir, id+".meta.json"), mustJSON(meta), 0o600); err != nil {
return Meta{}, err
}
s.enforceRetentionLocked(devDir)
return meta, nil
}
// List returns one device's runs, newest first.
func (s *Store) List(deviceID string) []Meta {
s.mu.Lock()
defer s.mu.Unlock()
return s.listLocked(filepath.Join(s.dir, sanitizeID(deviceID)))
}
// Get returns the stored document bytes for one run.
func (s *Store) Get(deviceID, runID string) ([]byte, error) {
s.mu.Lock()
defer s.mu.Unlock()
return os.ReadFile(filepath.Join(s.dir, sanitizeID(deviceID), sanitizeID(runID)+".json"))
}
// Delete removes one run. Missing is not an error: delete is idempotent so a client retrying
// after a dropped response does not see a spurious failure.
func (s *Store) Delete(deviceID, runID string) error {
s.mu.Lock()
defer s.mu.Unlock()
dev, run := sanitizeID(deviceID), sanitizeID(runID)
for _, suffix := range []string{".json", ".meta.json"} {
if err := os.Remove(filepath.Join(s.dir, dev, run+suffix)); err != nil && !errors.Is(err, os.ErrNotExist) {
return err
}
}
return nil
}
func (s *Store) listLocked(devDir string) []Meta {
entries, err := os.ReadDir(devDir)
if err != nil {
return nil
}
out := make([]Meta, 0, len(entries))
for _, e := range entries {
if !strings.HasSuffix(e.Name(), ".meta.json") {
continue
}
b, err := os.ReadFile(filepath.Join(devDir, e.Name()))
if err != nil {
continue
}
var m Meta
if json.Unmarshal(b, &m) == nil {
out = append(out, m)
}
}
sort.Slice(out, func(i, j int) bool { return out[i].UploadedAt.After(out[j].UploadedAt) })
return out
}
// enforceRetentionLocked drops runs past the age limit, then past the count limit. Age first, so
// a burst of uploads cannot push out runs that are still inside the retention window.
func (s *Store) enforceRetentionLocked(devDir string) {
metas := s.listLocked(devDir)
drop := func(m Meta) {
_ = os.Remove(filepath.Join(devDir, m.ID+".json"))
_ = os.Remove(filepath.Join(devDir, m.ID+".meta.json"))
}
kept := metas[:0]
if s.policy.RetentionDays > 0 {
cutoff := time.Now().Add(-time.Duration(s.policy.RetentionDays) * 24 * time.Hour)
for _, m := range metas {
if m.UploadedAt.Before(cutoff) {
drop(m)
continue
}
kept = append(kept, m)
}
} else {
kept = metas
}
if s.policy.MaxRunsPerDevice > 0 && len(kept) > s.policy.MaxRunsPerDevice {
for _, m := range kept[s.policy.MaxRunsPerDevice:] { // listLocked is newest-first
drop(m)
}
}
}
// sanitizeID keeps ids to characters that cannot escape the directory or collide with the
// .meta.json suffix convention. Ids are uuids and device ids in practice; anything else is
// truncated to nothing and rejected upstream.
func sanitizeID(s string) string {
var b strings.Builder
for _, r := range s {
switch {
case r >= 'a' && r <= 'z', r >= 'A' && r <= 'Z', r >= '0' && r <= '9', r == '-', r == '_':
b.WriteRune(r)
}
if b.Len() >= 64 {
break
}
}
return b.String()
}
func mustJSON(v any) []byte {
b, _ := json.Marshal(v)
return b
}
+197
View File
@@ -0,0 +1,197 @@
// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
package runs
import (
"encoding/json"
"errors"
"fmt"
"os"
"path/filepath"
"strings"
"testing"
"time"
)
func doc(id, anon string) []byte {
return []byte(fmt.Sprintf(
`{"run":{"id":%q,"started_at":"2026-08-01T10:00:00Z","privacy":{"anonymization":%q}},`+
`"findings":[{"id":"f1"},{"id":"f2"}],"summary":{"verdict":"warn"}}`, id, anon))
}
func open(t *testing.T, p Policy) (*Store, string) {
t.Helper()
dir := t.TempDir()
s, err := Open(dir, p)
if err != nil {
t.Fatalf("open: %v", err)
}
return s, dir
}
func TestModeOffRefusesEverything(t *testing.T) {
p := DefaultPolicy()
p.Mode = ModeOff
s, _ := open(t, p)
if _, err := s.Put("dev1", doc("run-1", AnonFull)); !errors.Is(err, ErrDisabled) {
t.Fatalf("want ErrDisabled, got %v", err)
}
}
// ModeAccount must refuse today rather than fall back to anonymous: an operator who selects the
// strict setting before accounts exist must not be silently running the permissive one.
func TestModeAccountRefusesUntilAccountsExist(t *testing.T) {
p := DefaultPolicy()
p.Mode = ModeAccount
s, _ := open(t, p)
if _, err := s.Put("dev1", doc("run-1", AnonFull)); !errors.Is(err, ErrNeedAccount) {
t.Fatalf("want ErrNeedAccount, got %v", err)
}
}
func TestMinAnonymizationEnforced(t *testing.T) {
p := DefaultPolicy()
p.MinAnonymization = AnonBalanced
s, _ := open(t, p)
if _, err := s.Put("dev1", doc("run-full", AnonFull)); !errors.Is(err, ErrNotAnonEnough) {
t.Fatalf("full should be refused when balanced is required, got %v", err)
}
// An undeclared level means nothing was stripped, so it must be treated as "full".
if _, err := s.Put("dev1", []byte(`{"run":{"id":"run-bare"},"summary":{}}`)); !errors.Is(err, ErrNotAnonEnough) {
t.Fatalf("undeclared level should be treated as full, got %v", err)
}
for _, lvl := range []string{AnonBalanced, AnonStrict} {
if _, err := s.Put("dev1", doc("run-"+lvl, lvl)); err != nil {
t.Fatalf("%s should be accepted: %v", lvl, err)
}
}
}
func TestSizeLimit(t *testing.T) {
p := DefaultPolicy()
p.MaxBytes = 200
s, _ := open(t, p)
big := append(doc("run-1", AnonFull), make([]byte, 400)...)
if _, err := s.Put("dev1", big); !errors.Is(err, ErrTooLarge) {
t.Fatalf("want ErrTooLarge, got %v", err)
}
}
func TestRetentionByCountKeepsNewest(t *testing.T) {
p := DefaultPolicy()
p.MaxRunsPerDevice = 3
s, _ := open(t, p)
for i := 0; i < 6; i++ {
if _, err := s.Put("dev1", doc(fmt.Sprintf("run-%d", i), AnonFull)); err != nil {
t.Fatalf("put %d: %v", i, err)
}
time.Sleep(2 * time.Millisecond) // distinct UploadedAt so "newest" is well defined
}
got := s.List("dev1")
if len(got) != 3 {
t.Fatalf("kept %d runs, want 3", len(got))
}
for i, want := range []string{"run-5", "run-4", "run-3"} {
if got[i].ID != want {
t.Fatalf("kept[%d] = %s, want %s (newest first)", i, got[i].ID, want)
}
}
// The documents themselves must be gone too, not just their index entries.
if _, err := s.Get("dev1", "run-0"); err == nil {
t.Fatal("purged run is still readable")
}
}
func TestRetentionByAge(t *testing.T) {
p := DefaultPolicy()
p.RetentionDays = 7
p.MaxRunsPerDevice = 0
s, dir := open(t, p)
if _, err := s.Put("dev1", doc("run-old", AnonFull)); err != nil {
t.Fatal(err)
}
// Backdate the index entry past the retention window.
metaPath := filepath.Join(dir, "runs", "dev1", "run-old.meta.json")
b, _ := os.ReadFile(metaPath)
var m Meta
_ = json.Unmarshal(b, &m)
m.UploadedAt = time.Now().Add(-30 * 24 * time.Hour)
nb, _ := json.Marshal(m)
if err := os.WriteFile(metaPath, nb, 0o600); err != nil {
t.Fatal(err)
}
if _, err := s.Put("dev1", doc("run-new", AnonFull)); err != nil {
t.Fatal(err)
}
got := s.List("dev1")
if len(got) != 1 || got[0].ID != "run-new" {
t.Fatalf("age retention did not drop the old run: %+v", got)
}
}
// Run and device ids reach the filesystem, so a hostile one must not be able to climb out of the
// store directory or overwrite another device's data.
func TestIDsCannotEscapeTheStoreDirectory(t *testing.T) {
s, dir := open(t, DefaultPolicy())
if _, err := s.Put("../../etc", doc("../../../passwd", AnonFull)); err != nil {
t.Fatalf("put: %v", err)
}
var found []string
_ = filepath.Walk(dir, func(p string, info os.FileInfo, err error) error {
if err == nil && !info.IsDir() {
rel, _ := filepath.Rel(dir, p)
found = append(found, filepath.ToSlash(rel))
}
return nil
})
for _, f := range found {
if strings.Contains(f, "..") {
t.Fatalf("path escaped the store: %s", f)
}
}
if len(found) == 0 {
t.Fatal("nothing written at all")
}
}
func TestListIsPerDevice(t *testing.T) {
s, _ := open(t, DefaultPolicy())
if _, err := s.Put("devA", doc("run-a", AnonFull)); err != nil {
t.Fatal(err)
}
if _, err := s.Put("devB", doc("run-b", AnonFull)); err != nil {
t.Fatal(err)
}
if got := s.List("devA"); len(got) != 1 || got[0].ID != "run-a" {
t.Fatalf("devA sees %+v", got)
}
if _, err := s.Get("devA", "run-b"); err == nil {
t.Fatal("devA could read devB's run")
}
}
func TestMetaSummarisesTheDocument(t *testing.T) {
s, _ := open(t, DefaultPolicy())
m, err := s.Put("dev1", doc("run-1", AnonBalanced))
if err != nil {
t.Fatal(err)
}
if m.FindingCount != 2 || m.Verdict != "warn" || m.Anonymization != AnonBalanced {
t.Fatalf("meta not extracted: %+v", m)
}
if m.StartedAt != "2026-08-01T10:00:00Z" {
t.Fatalf("started_at = %q", m.StartedAt)
}
}
func TestMalformedRejected(t *testing.T) {
s, _ := open(t, DefaultPolicy())
for _, body := range [][]byte{[]byte("not json"), []byte(`{"run":{}}`), []byte(`{}`)} {
if _, err := s.Put("dev1", body); !errors.Is(err, ErrMalformed) {
t.Fatalf("body %q: want ErrMalformed, got %v", body, err)
}
}
}