diff --git a/echolot-app/core-engine/bin/test/app/echo_lot/engine/LiveGrantedTest.kt b/echolot-app/core-engine/bin/test/app/echo_lot/engine/LiveGrantedTest.kt new file mode 100644 index 0000000..f8327be --- /dev/null +++ b/echolot-app/core-engine/bin/test/app/echo_lot/engine/LiveGrantedTest.kt @@ -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) + } +} diff --git a/echolot-app/core-engine/src/test/kotlin/app/echo_lot/engine/LiveGrantedTest.kt b/echolot-app/core-engine/src/test/kotlin/app/echo_lot/engine/LiveGrantedTest.kt new file mode 100644 index 0000000..f8327be --- /dev/null +++ b/echolot-app/core-engine/src/test/kotlin/app/echo_lot/engine/LiveGrantedTest.kt @@ -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) + } +} diff --git a/echolot-app/core-protocol/bin/main/app/echo_lot/protocol/ControlClient.kt b/echolot-app/core-protocol/bin/main/app/echo_lot/protocol/ControlClient.kt index 8275afc..2ff261d 100644 --- a/echolot-app/core-protocol/bin/main/app/echo_lot/protocol/ControlClient.kt +++ b/echolot-app/core-protocol/bin/main/app/echo_lot/protocol/ControlClient.kt @@ -80,6 +80,19 @@ class ControlClient(private val controlUrl: String, pins: Set) { 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}" } diff --git a/echolot-app/core-protocol/bin/main/app/echo_lot/protocol/ProbeSession.kt b/echolot-app/core-protocol/bin/main/app/echo_lot/protocol/ProbeSession.kt index 9961cb4..6266726 100644 --- a/echolot-app/core-protocol/bin/main/app/echo_lot/protocol/ProbeSession.kt +++ b/echolot-app/core-protocol/bin/main/app/echo_lot/protocol/ProbeSession.kt @@ -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 { + val out = ArrayList() + 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 { diff --git a/echolot-app/core-protocol/bin/main/app/echo_lot/protocol/Wire.kt b/echolot-app/core-protocol/bin/main/app/echo_lot/protocol/Wire.kt index 4a9a79d..2462f33 100644 --- a/echolot-app/core-protocol/bin/main/app/echo_lot/protocol/Wire.kt +++ b/echolot-app/core-protocol/bin/main/app/echo_lot/protocol/Wire.kt @@ -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 { diff --git a/echolot-app/core-protocol/src/main/kotlin/app/echo_lot/protocol/ControlClient.kt b/echolot-app/core-protocol/src/main/kotlin/app/echo_lot/protocol/ControlClient.kt index 8275afc..2ff261d 100644 --- a/echolot-app/core-protocol/src/main/kotlin/app/echo_lot/protocol/ControlClient.kt +++ b/echolot-app/core-protocol/src/main/kotlin/app/echo_lot/protocol/ControlClient.kt @@ -80,6 +80,19 @@ class ControlClient(private val controlUrl: String, pins: Set) { 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}" } diff --git a/echolot-app/core-protocol/src/main/kotlin/app/echo_lot/protocol/ProbeSession.kt b/echolot-app/core-protocol/src/main/kotlin/app/echo_lot/protocol/ProbeSession.kt index 9961cb4..6266726 100644 --- a/echolot-app/core-protocol/src/main/kotlin/app/echo_lot/protocol/ProbeSession.kt +++ b/echolot-app/core-protocol/src/main/kotlin/app/echo_lot/protocol/ProbeSession.kt @@ -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 { + val out = ArrayList() + 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 { diff --git a/echolot-app/core-protocol/src/main/kotlin/app/echo_lot/protocol/Wire.kt b/echolot-app/core-protocol/src/main/kotlin/app/echo_lot/protocol/Wire.kt index 4a9a79d..2462f33 100644 --- a/echolot-app/core-protocol/src/main/kotlin/app/echo_lot/protocol/Wire.kt +++ b/echolot-app/core-protocol/src/main/kotlin/app/echo_lot/protocol/Wire.kt @@ -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 { diff --git a/echolot-app/scripts/test-fmr.sh b/echolot-app/scripts/test-fmr.sh index 052fb8d..9740317 100644 --- a/echolot-app/scripts/test-fmr.sh +++ b/echolot-app/scripts/test-fmr.sh @@ -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 diff --git a/server/cmd/echolot-server/main.go b/server/cmd/echolot-server/main.go index 1fe0866..e5e0b33 100644 --- a/server/cmd/echolot-server/main.go +++ b/server/cmd/echolot-server/main.go @@ -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() diff --git a/server/internal/config/config.go b/server/internal/config/config.go index 3b0a3ca..9052a89 100644 --- a/server/internal/config/config.go +++ b/server/internal/config/config.go @@ -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_ 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_ 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") diff --git a/server/internal/control/bigsend_df_test.go b/server/internal/control/bigsend_df_test.go new file mode 100644 index 0000000..1f29a5b --- /dev/null +++ b/server/internal/control/bigsend_df_test.go @@ -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) + } + }) + } +} diff --git a/server/internal/control/control.go b/server/internal/control/control.go index ee7cd72..e74540d 100644 --- a/server/internal/control/control.go +++ b/server/internal/control/control.go @@ -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) +} diff --git a/server/internal/dataplane/df_linux.go b/server/internal/dataplane/df_linux.go new file mode 100644 index 0000000..d88b976 --- /dev/null +++ b/server/internal/dataplane/df_linux.go @@ -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 diff --git a/server/internal/dataplane/df_other.go b/server/internal/dataplane/df_other.go new file mode 100644 index 0000000..4ad98b8 --- /dev/null +++ b/server/internal/dataplane/df_other.go @@ -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 diff --git a/server/internal/dataplane/granted.go b/server/internal/dataplane/granted.go index 1e155c2..c3567b5 100644 --- a/server/internal/dataplane/granted.go +++ b/server/internal/dataplane/granted.go @@ -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,24 +79,46 @@ 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)) - for i, size := range sizes { - if size < HeaderSize+8 { - size = HeaderSize + 8 + + results := make([]BigSendResult, 0, len(sizes)) + burst := func() error { + for i, size := range sizes { + if size < HeaderSize+8 { + size = HeaderSize + 8 + } + if size > 9000 { // jumbo ceiling; beyond this the kernel will refuse anyway + size = 9000 + } + if !g.Allow(size) { + break + } + payload := make([]byte, size-HeaderSize) + // 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)) + 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 } - if size > 9000 { // jumbo ceiling; beyond this the kernel will refuse anyway - size = 9000 - } - if !g.Allow(size) { - break - } - payload := make([]byte, size-HeaderSize) - // 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) - time.Sleep(20 * time.Millisecond) // keep bursts from being read as congestion loss + return nil } - return attempted, 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() } diff --git a/server/internal/dataplane/udp.go b/server/internal/dataplane/udp.go index f48782f..d0c03be 100644 --- a/server/internal/dataplane/udp.go +++ b/server/internal/dataplane/udp.go @@ -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 { diff --git a/server/internal/runs/runs.go b/server/internal/runs/runs.go new file mode 100644 index 0000000..f6541ad --- /dev/null +++ b/server/internal/runs/runs.go @@ -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 +} diff --git a/server/internal/runs/runs_test.go b/server/internal/runs/runs_test.go new file mode 100644 index 0000000..c0db420 --- /dev/null +++ b/server/internal/runs/runs_test.go @@ -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) + } + } +}