Compare commits

...
Author SHA1 Message Date
mrambossekandClaude Fable 5 a7dccf7da2 frag_send: crafted IP fragments, so ordering can be tested and not just delivery
server-release / image (push) Successful in 15s
server-test / test (push) Successful in 32s
server-release / release (push) Successful in 32s
Letting the kernel fragment an oversized datagram answers one question — do
fragments get through. It cannot answer the more interesting one, because the
kernel always emits them in order, first one first.

The classic middlebox fault is exactly about that ordering. Only the first
fragment carries the UDP header, and therefore the ports; a stateful firewall
or NAT that has not seen it has no flow to match the rest against, and many
drop them. That is invisible to any in-order test and shows up in the field as
"large DNS answers fail on this network" or "the tunnel breaks when the MTU
drops" — it works until the network reorders, then fails intermittently, which
is the hardest kind of fault to chase.

So the server now builds the fragments itself (raw socket, IP_HDRINCL) and
controls their order: in_order as a baseline, reversed, and first-fragment-last.
The datagram is assembled and signed whole before being cut up, so what the
client reassembles is indistinguishable from an ordinary packet — otherwise it
would be measuring our sender rather than the path.

Two details that would silently produce wrong answers:
  - The UDP checksum is computed rather than left zero. A zero-checksum datagram
    is dropped by some middleboxes, and that drop would be recorded as a
    fragmentation failure, which is the wrong conclusion entirely.
  - Fragment offsets are in 8-byte units, so non-final fragments are rounded to
    a multiple of 8. A 100-byte fragment is not an error, it is a datagram no
    host will ever reassemble.

frag-send is advertised only when a raw socket can actually be opened — checked
by opening one, since a permission model has more ways to say no than a
capability bit has to say yes.

Fragment header arithmetic is unit-tested (reassembly coverage, MF flags, shared
IP ID, 8-byte offsets, checksum verification), cross-compiled and run on Linux
since the code is build-tagged.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-08-01 13:45:09 +02:00
mrambossekandClaude Fable 5 4ffa6e4ae2 engine: split packet loss by direction using the server's observations
"3 % loss" sends an engineer looking in both directions at once. The server
records every packet it received per sequence number, so the two cases are
distinguishable: sent-but-never-seen is upstream loss, seen-but-no-reply is
downstream. The findings say which, and say what is not implicated.

Downstream loss is measured against what reached the server, not against what
was sent — the other denominator counts every upstream loss twice and
overstates the return path.

Per-direction jitter comes out of the same records without needing synchronised
clocks: (server_rx - client_tx) carries a constant unknown offset, and
differencing successive samples cancels it, so RFC 3393 variation is honestly
attributable to a direction even though absolute latency is not.

Correlation is by wire sequence number, not loop index — the counter is shared
with every packet type on the session. ProbeSession exposes it even for a lost
probe, since that is precisely the packet whose direction is in question.

Live against fmr: 0.08 ms upstream jitter vs 0.85 ms downstream, an asymmetry a
round-trip test cannot see.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-08-01 13:09:59 +02:00
mrambossekandClaude Fable 5 3e7e3b8d33 scripts: one command to mint an enrollment link, QR included
Scanning beats pasting a 200-character string onto a phone, and with a device
attached the deep link can be delivered by adb with no typing at all.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-08-01 12:12:28 +02:00
mrambossekandClaude Fable 5 199807a8c9 docs: enrollment link encoding rules in the spec, session log in build-status
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-08-01 12:11:37 +02:00
mrambossekandClaude Fable 5 5291bdd045 enrollment: actually emit enroll_uri from the admin endpoint
server-release / image (push) Successful in 15s
server-test / test (push) Successful in 30s
server-release / release (push) Successful in 31s
The previous commit's edit to the admin handler silently did not apply, so the
endpoint still returned just the token. Caught by deploying and looking at the
response rather than by trusting the build to have picked it up.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-08-01 12:08:11 +02:00
mrambossekandClaude Fable 5 fe3658e009 chore: ignore the Kotlin compiler's .kotlin scratch directory
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-08-01 12:06:47 +02:00
20 changed files with 1217 additions and 86 deletions
+3
View File
@@ -43,3 +43,6 @@ web/.wrangler/
# Eclipse/JDT output from the VSCodium Java extension — not a build artifact we own
echolot-app/*/bin/
# Kotlin compiler scratch/error logs
echolot-app/.kotlin/
+8
View File
@@ -162,3 +162,11 @@ First build downloads AGP/Compose/Shizuku from Google Maven + Maven Central.
are SUPPORTED on both known devices via `Os.recvmsg` + `StructMsghdr` reflection.
3. Fold the confirmed capabilities + Shizuku dump-format samples back into the production
`core-probe` / `core-shizuku` modules.
## Enrolling a device with a server
`echolot-app/scripts/enroll-link.sh [note]` mints a §2.1 bootstrap link on fmr over SSH and prints
it (plus a QR if `qrencode` is installed, plus the `adb shell am start -a …VIEW -d '<uri>'` command
when a device is attached). The link carries a single-use token — treat it as a secret until spent.
Never hand-assemble one: the base64 pin needs percent-encoding, and a pin wrong by one character
fails as an inscrutable TLS error rather than as a bad pin.
+68
View File
@@ -660,3 +660,71 @@ One user-visible bug caught in the process: Go's JSON encoder HTML-escapes `<`,
default, so the refusal reached the client as `needs \u003e= 0.2.0`. Disabled at the encoder (this
is an API, not a page), and the client now *parses* the error field instead of pattern-matching it,
so it survives whatever a future encoder decides to escape.
### Enrollment: the server mints the bootstrap link (server-v0.5.3 … v0.5.4, 2026-08-01)
Until now a device was configured by hand-typing a control URL, a base64 SPKI pin and a
credential. That is the step that goes wrong, and it goes wrong quietly: a pin off by one
character does not fail loudly, it just never matches, and surfaces days later as an inscrutable
TLS error.
`POST /admin/enroll-tokens` now returns the whole §2.1 bootstrap link alongside the token, because
the server is the only party holding all three parts at once. The app takes it from a paste or an
`echolot://enroll` deep link (so a QR scan configures a server in one action) and writes URL, pin
and credential **together or not at all** — a half-applied server fails later, somewhere else,
with an error pointing at the wrong thing.
The control URL comes from `ECHOLOT_PUBLIC_URL` (set on fmr to `https://fmr-1.echo-lot.app:8443`),
falling back to the first control listen address; a wildcard bind warns rather than emitting a
link to `0.0.0.0`.
**The encoding trap, which is the whole reason this is tested across both languages.** The pin is
base64, so it contains `+`, `/` and `=` — each of which means something else in a query string. An
unencoded `+` decodes to a space, leaving the pin wrong by exactly one character. Base64 has no
spaces, so the parser restores them; that cannot damage a correctly-encoded pin and it rescues
every hand-assembled link. `LiveEnrollmentTest` redeems a link the *server* produced, which is the
only way to catch a disagreement between the Go assembler and the Kotlin parser — a unit test on
either side alone cannot see it. It also asserts the token is refused the second time.
Also fixed a spec divergence found while reading §2.1: the spec names the field
`device_credential`, the first implementation shipped `credential`. The server now sends both and
the client prefers the spec's; the alias goes once nothing reads it.
Two process notes from this round:
- An edit to the admin handler silently failed to apply and the endpoint kept returning just the
token. Caught by deploying and *looking at the response*, not by trusting a green build.
- The live suite is now six tests (`LiveServerTest`, `LiveMeasurement`, `LiveGranted`,
`LiveUpload`, `LiveCompat`, `LiveEnrollment`), all green against fmr from the PC with no device.
### Directional loss: which way is the packet loss? (2026-08-01)
A round trip can only report that *something* was lost somewhere, which is the least useful form
of the answer — "3 % loss" sends an engineer looking in both directions at once. The server
already records every packet it received per sequence number (§6), so the two cases are actually
distinguishable, and `train.udp_updown` now reports them separately:
- sent, never seen by the server → **upstream** loss
- seen by the server, reply never arrived → **downstream** loss
Findings name the direction and say what is *not* implicated, which is half the value:
`connectivity.loss_upstream` ("the return path is not implicated: replies came back for everything
that arrived"), `connectivity.loss_downstream`, `nat.udp_unreachable_upstream`.
Two things the implementation gets deliberately right:
- **Downstream loss is measured against what reached the server**, not against what was sent.
Using "sent" as the denominator counts every upstream loss a second time and overstates the
return path. Pinned by a test with loss in both directions at once.
- **Per-direction jitter without synchronised clocks.** Absolute one-way delay would need clock
sync and we deliberately have none (the two-clock rule). But `server_rx client_tx` carries a
constant unknown offset, and differencing successive samples cancels it — so RFC 3393 one-way
delay variation *is* honestly attributable to a direction even though latency is not. A test
pins that a 10-second clock offset changes nothing.
Correlation is by **wire sequence number**, which is not the loop index: the counter is shared
with every other packet type on the session, so "the nth echo" is not "sequence n". `ProbeSession`
now exposes `lastSeq`, including for a probe that was lost — a lost packet still has a sequence
number, and that number is exactly what tells you which way it was lost.
Live against fmr: 20/20 both ways, and jitter of **0.08 ms upstream vs 0.85 ms downstream** — a
tenfold asymmetry that a round-trip measurement cannot see at all.
10 unit tests on the arithmetic (a wrong denominator here does not crash, it produces a plausible
number pointing at the wrong half of the network) plus the live correlation check.
+26 -3
View File
@@ -24,11 +24,34 @@ echolot://enroll?v=1&u=<control-URL, urlencoded>&p=pin-sha256:<b64 SPKI hash>&t=
```
POST /v1/enroll Authorization: Bearer <enrollment-token>
→ 200 { "device_credential": "<random 256-bit, b64url>",
"device_id": "uuid",
"profile": { ... §2.2 ... } }
→ 201 { "device_credential": "<random 256-bit, b64url>",
"device_id": "uuid" }
```
The **server assembles the bootstrap link**, because it is the only party holding all three parts
at once, and the part an operator gets wrong by hand is the base64 pin — which does not fail
loudly, it just never matches, and surfaces later as an inscrutable TLS error:
```
POST /admin/enroll-tokens
→ { "token": "…", "expires_in_s": 86400,
"enroll_uri": "echolot://enroll?v=1&u=…&p=…&t=…" }
```
The control URL in the link comes from `ECHOLOT_PUBLIC_URL`, falling back to the first control
listen address. A wildcard bind has no single right answer, so it warns rather than guessing.
Encoding notes that matter in practice:
- `u`, `p` and `t` are **percent-encoded**. The pin is base64, so it contains `+`, `/` and `=`,
every one of which means something else in a query string.
- A `+` that was *not* encoded decodes to a space. Base64 contains no spaces, so a parser SHOULD
restore them — the alternative is a pin wrong by one character and a failure that points nowhere
near the cause.
- The control URL MUST be `https://`. The pin only protects a TLS connection; a cleartext URL
would hand the token to anyone on the path.
- **The link is a secret** while it is live: it carries a bearer token, so anyone who sees it
before the device does can enroll instead.
Enrollment tokens are single-use with expiry, created in the admin UI, scoped `enroll`. The device credential is a long-lived bearer secret, scoped `run-tests`; it is also the HKDF input for session keys. Revocation = deleting the device in the admin UI.
### 2.2 Profile
@@ -1,75 +0,0 @@
kotlin version: 2.2.10
error message: Daemon compilation failed: null
java.lang.Exception
at org.jetbrains.kotlin.daemon.common.CompileService$CallResult$Error.get(CompileService.kt:69)
at org.jetbrains.kotlin.daemon.common.CompileService$CallResult$Error.get(CompileService.kt:65)
at org.jetbrains.kotlin.compilerRunner.GradleKotlinCompilerWork.compileWithDaemon(GradleKotlinCompilerWork.kt:240)
at org.jetbrains.kotlin.compilerRunner.GradleKotlinCompilerWork.compileWithDaemonOrFallbackImpl(GradleKotlinCompilerWork.kt:159)
at org.jetbrains.kotlin.compilerRunner.GradleKotlinCompilerWork.run(GradleKotlinCompilerWork.kt:111)
at org.jetbrains.kotlin.compilerRunner.GradleCompilerRunnerWithWorkers$GradleKotlinCompilerWorkAction.execute(GradleCompilerRunnerWithWorkers.kt:74)
at org.gradle.workers.internal.DefaultWorkerServer.execute(DefaultWorkerServer.java:68)
at org.gradle.workers.internal.NoIsolationWorkerFactory$1$1.create(NoIsolationWorkerFactory.java:64)
at org.gradle.workers.internal.NoIsolationWorkerFactory$1$1.create(NoIsolationWorkerFactory.java:61)
at org.gradle.internal.classloader.ClassLoaderUtils.executeInClassloader(ClassLoaderUtils.java:102)
at org.gradle.workers.internal.NoIsolationWorkerFactory$1.lambda$execute$0(NoIsolationWorkerFactory.java:61)
at org.gradle.workers.internal.AbstractWorker$1.call(AbstractWorker.java:44)
at org.gradle.workers.internal.AbstractWorker$1.call(AbstractWorker.java:41)
at org.gradle.internal.operations.DefaultBuildOperationRunner$CallableBuildOperationWorker.execute(DefaultBuildOperationRunner.java:210)
at org.gradle.internal.operations.DefaultBuildOperationRunner$CallableBuildOperationWorker.execute(DefaultBuildOperationRunner.java:205)
at org.gradle.internal.operations.DefaultBuildOperationRunner$2.execute(DefaultBuildOperationRunner.java:67)
at org.gradle.internal.operations.DefaultBuildOperationRunner$2.execute(DefaultBuildOperationRunner.java:60)
at org.gradle.internal.operations.DefaultBuildOperationRunner.execute(DefaultBuildOperationRunner.java:167)
at org.gradle.internal.operations.DefaultBuildOperationRunner.execute(DefaultBuildOperationRunner.java:60)
at org.gradle.internal.operations.DefaultBuildOperationRunner.call(DefaultBuildOperationRunner.java:54)
at org.gradle.workers.internal.AbstractWorker.executeWrappedInBuildOperation(AbstractWorker.java:41)
at org.gradle.workers.internal.NoIsolationWorkerFactory$1.execute(NoIsolationWorkerFactory.java:58)
at org.gradle.workers.internal.DefaultWorkerExecutor.lambda$submitWork$0(DefaultWorkerExecutor.java:174)
at java.base/java.util.concurrent.FutureTask.run(Unknown Source)
at org.gradle.internal.work.DefaultConditionalExecutionQueue$ExecutionRunner.runExecution(DefaultConditionalExecutionQueue.java:191)
at org.gradle.internal.work.DefaultConditionalExecutionQueue$ExecutionRunner.access$500(DefaultConditionalExecutionQueue.java:112)
at org.gradle.internal.work.DefaultConditionalExecutionQueue$ExecutionRunner$1.run(DefaultConditionalExecutionQueue.java:168)
at org.gradle.internal.Factories$1.create(Factories.java:30)
at org.gradle.internal.work.DefaultWorkerLeaseService.lambda$runAndReleaseLocks$0(DefaultWorkerLeaseService.java:300)
at org.gradle.internal.work.ResourceLockStatistics$1.measure(ResourceLockStatistics.java:43)
at org.gradle.internal.work.DefaultWorkerLeaseService.runAndReleaseLocks(DefaultWorkerLeaseService.java:298)
at org.gradle.internal.work.DefaultWorkerLeaseService.withLocksAcquired(DefaultWorkerLeaseService.java:294)
at org.gradle.internal.work.DefaultWorkerLeaseService.withLocks(DefaultWorkerLeaseService.java:286)
at org.gradle.internal.work.DefaultWorkerLeaseService.runAsWorkerThread(DefaultWorkerLeaseService.java:130)
at org.gradle.internal.work.DefaultWorkerLeaseService.runAsWorkerThread(DefaultWorkerLeaseService.java:135)
at org.gradle.internal.work.DefaultConditionalExecutionQueue$ExecutionRunner.runBatch(DefaultConditionalExecutionQueue.java:163)
at org.gradle.internal.work.DefaultConditionalExecutionQueue$ExecutionRunner.run(DefaultConditionalExecutionQueue.java:125)
at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Unknown Source)
at java.base/java.util.concurrent.FutureTask.run(Unknown Source)
at org.gradle.internal.concurrent.ExecutorPolicy$CatchAndRecordFailures.onExecute(ExecutorPolicy.java:64)
at org.gradle.internal.concurrent.AbstractManagedExecutor$1.run(AbstractManagedExecutor.java:47)
at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(Unknown Source)
at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(Unknown Source)
at java.base/java.lang.Thread.run(Unknown Source)
Caused by: java.nio.file.NoSuchFileException: C:\Users\mram.AD\AppData\Local\Temp\kotlin-backups9735132086774701577\24.backup -> C:\Users\mram.AD\dev\echolot\echolot-app\core-protocol\build\classes\kotlin\main\META-INF\core-protocol.kotlin_module
at java.base/sun.nio.fs.WindowsException.translateToIOException(Unknown Source)
at java.base/sun.nio.fs.WindowsException.rethrowAsIOException(Unknown Source)
at java.base/sun.nio.fs.WindowsFileCopy.move(Unknown Source)
at java.base/sun.nio.fs.WindowsFileSystemProvider.move(Unknown Source)
at java.base/java.nio.file.Files.move(Unknown Source)
at org.jetbrains.kotlin.incremental.RecoverableCompilationTransaction.revertChanges(CompilationTransaction.kt:231)
at org.jetbrains.kotlin.incremental.RecoverableCompilationTransaction.close(CompilationTransaction.kt:256)
at org.jetbrains.kotlin.incremental.IncrementalCompilerRunner.tryCompileIncrementally(IncrementalCompilerRunner.kt:740)
at org.jetbrains.kotlin.incremental.IncrementalCompilerRunner.compile(IncrementalCompilerRunner.kt:124)
at org.jetbrains.kotlin.daemon.CompileServiceImplBase.execIncrementalCompiler(CompileServiceImpl.kt:679)
at org.jetbrains.kotlin.daemon.CompileServiceImplBase.access$execIncrementalCompiler(CompileServiceImpl.kt:93)
at org.jetbrains.kotlin.daemon.CompileServiceImpl.compile(CompileServiceImpl.kt:1806)
at java.base/jdk.internal.reflect.DirectMethodHandleAccessor.invoke(Unknown Source)
at java.base/java.lang.reflect.Method.invoke(Unknown Source)
at java.rmi/sun.rmi.server.UnicastServerRef.dispatch(Unknown Source)
at java.rmi/sun.rmi.transport.Transport$1.run(Unknown Source)
at java.rmi/sun.rmi.transport.Transport$1.run(Unknown Source)
at java.base/java.security.AccessController.doPrivileged(Unknown Source)
at java.rmi/sun.rmi.transport.Transport.serviceCall(Unknown Source)
at java.rmi/sun.rmi.transport.tcp.TCPTransport.handleMessages(Unknown Source)
at java.rmi/sun.rmi.transport.tcp.TCPTransport$ConnectionHandler.run0(Unknown Source)
at java.rmi/sun.rmi.transport.tcp.TCPTransport$ConnectionHandler.lambda$run$0(Unknown Source)
at java.base/java.security.AccessController.doPrivileged(Unknown Source)
at java.rmi/sun.rmi.transport.tcp.TCPTransport$ConnectionHandler.run(Unknown Source)
... 3 more
@@ -0,0 +1,107 @@
// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
package app.echo_lot.engine
import kotlinx.serialization.SerialName
import kotlinx.serialization.Serializable
/**
* Splits a round-trip train into its two directions using what the server witnessed.
*
* A round trip can only report that *something* was lost somewhere. That is the least useful form
* of the answer: "3 % loss" sends an engineer looking in both directions at once. The server
* records every packet it received, per sequence number (probe-protocol.md §6), so the two cases
* are actually distinguishable:
*
* - sent, never seen by the server → **upstream** loss
* - seen by the server, reply never arrived → **downstream** loss
*
* The same records give one-way delay *variation* per direction. Absolute one-way delay would
* need synchronised clocks and we deliberately have none (measurement-schema.md's two-clock rule),
* but the variation does not: (server_rx client_tx) contains an unknown constant clock offset,
* and differencing successive samples cancels it. So jitter is honestly attributable to a
* direction even though latency is not.
*/
object Directional {
/** One probe as the client saw it. [tRxNs] null means no reply came back. */
data class Sample(val seq: Int, val tTxNs: Long, val tRxNs: Long?)
/** One probe as the server saw it: its own receive and transmit stamps, on its own clock. */
data class ServerSighting(val seq: Int, val tRxNs: Long, val tTxNs: Long)
fun analyse(sent: List<Sample>, seen: List<ServerSighting>): DirectionalMetrics {
val byServerSeq = seen.associateBy { it.seq }
// Only sequences we actually sent count. A server record for a sequence we have no note
// of is not evidence about this train — it is a bug or a stray, and silently folding it
// in would produce loss percentages above 100 or below zero.
val relevant = sent.filter { byServerSeq.containsKey(it.seq) }
val nSent = sent.size
val nSeen = relevant.size
val nReplied = sent.count { it.tRxNs != null }
// A reply can only exist if the request arrived, so downstream loss is measured against
// what the server saw, not against what we sent — otherwise upstream loss is counted twice.
val lostUp = nSent - nSeen
val lostDown = (nSeen - nReplied).coerceAtLeast(0)
val upDeltas = relevant.sortedBy { it.seq }
.map { byServerSeq.getValue(it.seq).tRxNs - it.tTxNs }
val downDeltas = sent.filter { it.tRxNs != null && byServerSeq.containsKey(it.seq) }
.sortedBy { it.seq }
.map { it.tRxNs!! - byServerSeq.getValue(it.seq).tTxNs }
return DirectionalMetrics(
sent = nSent,
seenByServer = nSeen,
repliesReceived = nReplied,
lostUpstream = lostUp,
lostDownstream = lostDown,
lossUpstreamPct = pct(lostUp, nSent),
// Denominator is what reached the server: of the packets that got there, how many
// replies came back.
lossDownstreamPct = pct(lostDown, nSeen),
jitterUpstreamMs = jitterMs(upDeltas),
jitterDownstreamMs = jitterMs(downDeltas),
/** True when the server saw nothing at all, which is a different fault from loss. */
noneReachedServer = nSent > 0 && nSeen == 0,
)
}
/**
* Mean absolute difference between consecutive one-way samples (RFC 3393 IPDV, averaged).
*
* Differencing is what makes this legitimate without synchronised clocks: each sample carries
* the same unknown offset between the two clocks, and the difference cancels it. Fewer than
* two samples yields null rather than zero — "no jitter" and "not enough data to say" are
* different claims and only one of them is true here.
*/
private fun jitterMs(oneWayNs: List<Long>): Double? {
if (oneWayNs.size < 2) return null
val deltas = oneWayNs.zipWithNext { a, b -> kotlin.math.abs(b - a) }
return round2(deltas.average() / 1_000_000.0)
}
private fun pct(part: Int, whole: Int): Double =
if (whole <= 0) 0.0 else round2(part * 100.0 / whole)
private fun round2(v: Double) = Math.round(v * 100.0) / 100.0
}
/** Directional metrics for train.udp_updown; recomputable from the columnar evidence. */
@Serializable
data class DirectionalMetrics(
val sent: Int,
@SerialName("seen_by_server") val seenByServer: Int,
@SerialName("replies_received") val repliesReceived: Int,
@SerialName("lost_upstream") val lostUpstream: Int,
@SerialName("lost_downstream") val lostDownstream: Int,
@SerialName("loss_upstream_pct") val lossUpstreamPct: Double,
@SerialName("loss_downstream_pct") val lossDownstreamPct: Double,
/** One-way delay variation (RFC 3393), per direction. Null when there were too few samples. */
@SerialName("jitter_upstream_ms") val jitterUpstreamMs: Double? = null,
@SerialName("jitter_downstream_ms") val jitterDownstreamMs: Double? = null,
@SerialName("none_reached_server") val noneReachedServer: Boolean = false,
)
@@ -37,6 +37,120 @@ class DownstreamMeasurement(private val ids: IdSource) {
/** How long to wait for a granted burst after the server accepts the action. */
private val collectWindowMs = 4_000L
/**
* Shorter, but long enough to cover the first_last mode's deliberate 250 ms hold plus a
* reassembly. A fragment burst is one datagram: it is here quickly or not at all.
*/
private val fragWindowMs = 1_500L
/**
* Asks the server to send one deliberately-fragmented datagram per ordering, and reports
* which orderings survive the path.
*
* Kernel fragmentation always emits fragments in order, first one first, so an oversized
* datagram can only answer "do fragments get through at all". The interesting fault is about
* ordering: only the *first* fragment carries the UDP ports, so a stateful firewall that has
* not seen it has nothing to match the rest against, and many drop them. That failure is
* invisible to every in-order test and shows up in the field as "large DNS answers fail here"
* or "the tunnel breaks when the MTU drops".
*/
fun fragmentOrdering(
credential: String,
sessionId: String,
control: ControlClient,
probe: ProbeSession,
sessionRef: String,
sizeBytes: Int = 2000,
fragBytes: Int = 576,
): Pair<Test, List<Finding>> {
val testId = ids.uuid()
val started = ids.monoNs()
val delivered = LinkedHashMap<String, Boolean>()
val fragmentCounts = LinkedHashMap<String, Int>()
var unsupported = false
for (mode in FRAG_MODES) {
val reply = runCatching {
control.action(
credential, sessionId,
"""{"action":"frag_send","size_bytes":$sizeBytes,"mode":"$mode","frag_bytes":$fragBytes}""",
)
}
if (reply.isFailure) {
// A server without a raw socket says so; that is a missing capability, not a
// property of the network, and must not be recorded as a failed delivery.
unsupported = true
break
}
parseInt(reply.getOrNull(), "fragments")?.let { fragmentCounts[mode] = it }
// The burst is already on the wire when the action returns (it is sent
// synchronously), so anything that survived is either here or lost.
val got = probe.collectGranted(fragWindowMs).any { it.type == Wire.TYPE_FRAG_DATA }
delivered[mode] = got
}
if (unsupported) {
return Test(
id = testId, type = TestType.MTU_FRAG_ORDERING, sessionRef = sessionRef, tier = Tier.APP,
startedMonoNs = started, endedMonoNs = ids.monoNs(),
status = TestStatus.UNSUPPORTED,
error = TestError("no_raw_socket", "this server cannot craft fragments"),
) to emptyList()
}
val metrics = json.encodeToJsonElement(
FragOrderingMetrics(
sizeBytes = sizeBytes,
fragBytes = fragBytes,
fragmentsPerBurst = fragmentCounts,
deliveredByMode = delivered,
inOrderDelivered = delivered[FRAG_IN_ORDER] == true,
reorderedDelivered = delivered[FRAG_REVERSED] == true,
delayedFirstDelivered = delivered[FRAG_FIRST_LAST] == true,
),
) as JsonObject
val findings = ArrayList<Finding>()
val inOrder = delivered[FRAG_IN_ORDER] == true
val reversed = delivered[FRAG_REVERSED] == true
val firstLast = delivered[FRAG_FIRST_LAST] == true
if (!inOrder) {
findings.add(
finding(
"mtu.fragments_blocked", Category.MTU, Severity.MEDIUM, testId,
"IP fragments do not reach this device",
"A fragmented datagram sent in the normal order never arrived. Anything that " +
"relies on fragmentation — large DNS answers over UDP, some VPN traffic — " +
"will fail here rather than slow down.",
),
)
} else if (!reversed || !firstLast) {
// The precise and useful finding: fragments work, but only if they arrive tidily.
val which = buildList {
if (!reversed) add("out of order")
if (!firstLast) add("with the first fragment delayed")
}.joinToString(" or ")
findings.add(
finding(
"mtu.fragment_reorder_sensitive", Category.MTU, Severity.LOW, testId,
"Fragments are dropped when they arrive $which",
"In-order fragments are delivered, but the same datagram sent $which is not. " +
"Something on the path only reassembles when the first fragment (the one " +
"carrying the UDP ports) arrives first — typical of a stateful firewall " +
"or NAT. It works until the network reorders, then fails intermittently, " +
"which is the hardest kind of fault to chase.",
),
)
}
return Test(
id = testId, type = TestType.MTU_FRAG_ORDERING, sessionRef = sessionRef, tier = Tier.APP,
startedMonoNs = started, endedMonoNs = ids.monoNs(),
status = if (inOrder) TestStatus.OK else TestStatus.PARTIAL,
metrics = metrics,
) to findings
}
/**
* Runs all three against an already-primed session.
*
@@ -66,6 +180,16 @@ class DownstreamMeasurement(private val ids: IdSource) {
tests.add(df.test); tests.add(frag.test); tests.add(train.test)
// Fragment ordering only makes sense once we know fragments arrive at all; when they do
// not, the ordering variants would all report "not delivered" and read as three faults
// instead of one.
if (frag.largestDelivered != null) {
val (fragTest, fragFindings) =
fragmentOrdering(credential, sessionId, control, probe, sessionRef)
tests.add(fragTest)
findings.addAll(fragFindings)
}
// A downstream MTU below the classic 1500-byte Ethernet payload is worth saying out loud:
// it is the usual cause of "small requests work, large responses hang".
val pathMtu = df.largestDelivered
@@ -310,6 +434,11 @@ class DownstreamMeasurement(private val ids: IdSource) {
/** IPv4 (20) + UDP (8). The v6 case is 48; reported per-family once v6 sessions land. */
const val IP_UDP_OVERHEAD4 = 28
const val FRAG_IN_ORDER = "in_order"
const val FRAG_REVERSED = "reversed"
const val FRAG_FIRST_LAST = "first_last"
val FRAG_MODES = listOf(FRAG_IN_ORDER, FRAG_REVERSED, FRAG_FIRST_LAST)
/** Straddles the usual suspects: 1500 Ethernet, 1492 PPPoE, 1400-ish tunnels. */
val DEFAULT_SIZES = listOf(600, 1200, 1372, 1400, 1450, 1472, 1500, 2000, 4000)
@@ -330,6 +459,18 @@ data class BigSendMetrics(
@SerialName("path_mtu_bytes") val pathMtuBytes: Int? = null,
)
/** Metrics for mtu.frag_ordering. */
@Serializable
data class FragOrderingMetrics(
@SerialName("size_bytes") val sizeBytes: Int,
@SerialName("frag_bytes") val fragBytes: Int,
@SerialName("fragments_per_burst") val fragmentsPerBurst: Map<String, Int>,
@SerialName("delivered_by_mode") val deliveredByMode: Map<String, Boolean>,
@SerialName("in_order_delivered") val inOrderDelivered: Boolean,
@SerialName("reordered_delivered") val reorderedDelivered: Boolean,
@SerialName("delayed_first_delivered") val delayedFirstDelivered: Boolean,
)
/** Metrics for train.udp_downstream. */
@Serializable
data class DownTrainMetrics(
@@ -10,8 +10,13 @@ import app.echo_lot.measurement.*
import app.echo_lot.protocol.ControlClient
import app.echo_lot.protocol.ProbeSession
import kotlinx.serialization.json.Json
import kotlinx.serialization.json.JsonArray
import kotlinx.serialization.json.JsonObject
import kotlinx.serialization.json.encodeToJsonElement
import kotlinx.serialization.json.intOrNull
import kotlinx.serialization.json.jsonObject
import kotlinx.serialization.json.jsonPrimitive
import kotlinx.serialization.json.longOrNull
/**
* Runs the server-facing measurements against one target and assembles a [MeasurementDocument]:
@@ -71,7 +76,7 @@ class ServerMeasurement(
// re-primed source is never recorded and every granted send goes to the old, closed port.
// Session identity lives on the server; the socket must live as long as it does.
ProbeSession(cfg.credential, session, cfg.udpHost, cfg.udpPort).use { ps ->
val (test, findings) = echoTrain(cfg, ps, startMono)
val (test, findings) = echoTrain(cfg, ps, startMono, control, session.sessionId)
tests.add(test)
allFindings.addAll(findings)
@@ -105,6 +110,7 @@ class ServerMeasurement(
private fun echoTrain(
cfg: Config, ps: ProbeSession, startMono: Long,
control: ControlClient? = null, sessionId: String? = null,
): Pair<Test, List<Finding>> {
val testId = ids.uuid()
val seqs = ArrayList<Int>()
@@ -114,9 +120,15 @@ class ServerMeasurement(
val rtts = ArrayList<Double>()
val observedPorts = LinkedHashSet<Int>()
// Wire sequence numbers, kept so the server's observations can be correlated packet by
// packet. They are not 0..n-1: the counter is shared with every other packet type on the
// session, so "the nth echo" is not "sequence n".
val wireSeqs = ArrayList<Int>()
for (i in 0 until cfg.echoCount) {
val txMono = ids.monoNs() - startMono
val r = ps.echo(cfg.echoPaddingBytes)
wireSeqs.add(ps.lastSeq)
seqs.add(i)
tTx.add(txMono)
sizes.add(Wire_HEADER + cfg.echoPaddingBytes)
@@ -129,6 +141,20 @@ class ServerMeasurement(
}
}
// Ask the server what it actually received. This is what turns "3 % loss somewhere" into
// "3 % loss upstream" - the least useful form of the answer into a usable one.
val directional: DirectionalMetrics? =
if (control != null && sessionId != null) {
runCatching {
val samples = wireSeqs.indices.map {
Directional.Sample(wireSeqs[it], tTx[it] ?: 0L, tRx[it])
}
Directional.analyse(samples, serverSightings(control, cfg, sessionId))
}.getOrNull() // an older server without the endpoint simply yields no split
} else {
null
}
val sent = cfg.echoCount
val received = rtts.size
val lossPct = if (sent == 0) 0.0 else (sent - received) * 100.0 / sent
@@ -138,6 +164,9 @@ class ServerMeasurement(
epochMonoNs = startMono, seq = seqs, tTxNs = tTx, tRxNs = tRx, sizeBytes = sizes,
).toEvidence()
val directionalJson = directional?.let {
json.encodeToJsonElement(DirectionalMetrics.serializer(), it) as JsonObject
}
val metrics: JsonObject = json.encodeToJsonElement(
EchoMetrics(
sent = sent, received = received, lossPct = round1(lossPct),
@@ -147,7 +176,7 @@ class ServerMeasurement(
observedPorts = observedPorts.toList(),
natRebindingDetected = natRebinding,
)
) as JsonObject
).let { base -> JsonObject((base as JsonObject) + (directionalJson ?: JsonObject(emptyMap()))) }
val status = when {
received == 0 -> TestStatus.FAILED
@@ -170,6 +199,35 @@ class ServerMeasurement(
"High UDP loss to the server (${round1(lossPct)}%)",
"A large fraction of ECHO probes were lost, indicating an unreliable UDP path."))
}
// Naming the direction is the entire value of the split, so the findings do.
directional?.let { d ->
when {
d.noneReachedServer && received == 0 -> findings.add(
finding("nat.udp_unreachable_upstream", Category.CONNECTIVITY, Severity.HIGH, testId,
"Nothing reached the server",
"The server received none of the ${d.sent} probes, so the traffic is being " +
"dropped on the way out, not on the way back. A firewall or NAT on " +
"this side of the path is the place to look."),
)
d.lossUpstreamPct >= 2.0 -> findings.add(
finding("connectivity.loss_upstream", Category.CONNECTIVITY, Severity.MEDIUM, testId,
"${d.lossUpstreamPct} % of probes were lost on the way to the server",
"${d.lostUpstream} of ${d.sent} probes never reached the server. The " +
"return path is not implicated: replies came back for everything that " +
"arrived."),
)
}
if (d.lossDownstreamPct >= 2.0) {
findings.add(
finding("connectivity.loss_downstream", Category.CONNECTIVITY, Severity.MEDIUM, testId,
"${d.lossDownstreamPct} % of replies were lost on the way back",
"The server received ${d.seenByServer} probes and answered them, but " +
"${d.lostDownstream} of those replies never arrived. The outbound path " +
"is fine; the fault is on the return leg."),
)
}
}
if (natRebinding) {
findings.add(finding("nat.udp_rebinding", Category.NAT, Severity.MEDIUM, testId,
"NAT remapped the UDP source port mid-flow",
@@ -178,6 +236,29 @@ class ServerMeasurement(
return test to findings
}
/**
* The server's per-packet record of this session's echoes (spec section 6). Filtered to
* ECHO_REQ, because the observation list also holds MTU probes and anything else we sent -
* counting those as train packets would invent loss that is not there.
*/
private fun serverSightings(
control: ControlClient, cfg: Config, sessionId: String,
): List<Directional.ServerSighting> {
val body = control.observations(cfg.credential, sessionId)
val packets = Json.parseToJsonElement(body).jsonObject["udp"]
?.jsonObject?.get("packets") as? JsonArray ?: return emptyList()
return packets.mapNotNull { el ->
val o = el as? JsonObject ?: return@mapNotNull null
val type = o["type"]?.jsonPrimitive?.intOrNull ?: return@mapNotNull null
if (type != ECHO_REQ_TYPE) return@mapNotNull null
Directional.ServerSighting(
seq = o["seq"]?.jsonPrimitive?.intOrNull ?: return@mapNotNull null,
tRxNs = o["t_rx_ns"]?.jsonPrimitive?.longOrNull ?: return@mapNotNull null,
tTxNs = o["t_tx_ns"]?.jsonPrimitive?.longOrNull ?: return@mapNotNull null,
)
}
}
private fun finding(code: String, cat: Category, sev: Severity, testId: String, title: String, desc: String) =
Finding(
id = ids.uuid(), code = code, category = cat, severity = sev, confidence = Confidence.HIGH,
@@ -186,6 +267,7 @@ class ServerMeasurement(
private companion object {
const val Wire_HEADER = 32
const val ECHO_REQ_TYPE = 0x01
fun round1(v: Double) = Math.round(v * 10.0) / 10.0
}
}
@@ -0,0 +1,155 @@
// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
package app.echo_lot.engine
import app.echo_lot.engine.Directional.Sample
import app.echo_lot.engine.Directional.ServerSighting
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFalse
import kotlin.test.assertNotNull
import kotlin.test.assertNull
import kotlin.test.assertTrue
/**
* The arithmetic that turns "3 % loss somewhere" into "3 % loss upstream". Getting a denominator
* wrong here does not crash anything — it produces a plausible number pointing at the wrong half
* of the network, which is worse than no number at all. Hence a test per claim.
*/
class DirectionalTest {
/** A clean train: every packet sent, seen and answered. Server clock offset by a constant. */
private fun clean(n: Int, offsetNs: Long = 5_000_000_000L): Pair<List<Sample>, List<ServerSighting>> {
val sent = (1..n).map { Sample(it, tTxNs = it * 10_000_000L, tRxNs = it * 10_000_000L + 4_000_000L) }
val seen = (1..n).map {
ServerSighting(it, tRxNs = offsetNs + it * 10_000_000L + 2_000_000L,
tTxNs = offsetNs + it * 10_000_000L + 2_100_000L)
}
return sent to seen
}
@Test
fun aCleanTrainReportsNoLossInEitherDirection() {
val (sent, seen) = clean(10)
val m = Directional.analyse(sent, seen)
assertEquals(10, m.sent)
assertEquals(10, m.seenByServer)
assertEquals(10, m.repliesReceived)
assertEquals(0.0, m.lossUpstreamPct)
assertEquals(0.0, m.lossDownstreamPct)
assertFalse(m.noneReachedServer)
}
// The whole point: a packet the server never saw was lost on the way there.
@Test
fun packetsTheServerNeverSawAreUpstreamLoss() {
val (sent, seen) = clean(10)
val m = Directional.analyse(sent, seen.filter { it.seq !in setOf(3, 7) })
assertEquals(2, m.lostUpstream)
assertEquals(0, m.lostDownstream)
assertEquals(20.0, m.lossUpstreamPct)
assertEquals(0.0, m.lossDownstreamPct, "a packet that never arrived cannot be lost coming back")
}
@Test
fun repliesThatNeverArrivedAreDownstreamLoss() {
val (sent, seen) = clean(10)
val withHoles = sent.map { if (it.seq in setOf(2, 5)) it.copy(tRxNs = null) else it }
val m = Directional.analyse(withHoles, seen)
assertEquals(0, m.lostUpstream)
assertEquals(2, m.lostDownstream)
assertEquals(20.0, m.lossDownstreamPct)
}
// Downstream loss is measured against what actually reached the server. Using "sent" as the
// denominator would count every upstream loss a second time and overstate the return path.
@Test
fun downstreamLossIsRelativeToWhatReachedTheServer() {
val (sent, seen) = clean(10)
// 5 lost on the way there; of the 5 that arrived, 1 reply is lost coming back.
val seenPartial = seen.filter { it.seq > 5 }
val withHole = sent.map {
when {
it.seq <= 5 -> it.copy(tRxNs = null) // never got there, so never came back
it.seq == 6 -> it.copy(tRxNs = null) // arrived, reply lost
else -> it
}
}
val m = Directional.analyse(withHole, seenPartial)
assertEquals(5, m.lostUpstream)
assertEquals(50.0, m.lossUpstreamPct)
assertEquals(1, m.lostDownstream)
assertEquals(20.0, m.lossDownstreamPct, "1 of the 5 that arrived, not 1 of 10")
}
@Test
fun aServerThatSawNothingIsCalledOutSeparately() {
val (sent, _) = clean(6)
val m = Directional.analyse(sent.map { it.copy(tRxNs = null) }, emptyList())
assertTrue(m.noneReachedServer)
assertEquals(100.0, m.lossUpstreamPct)
assertEquals(0.0, m.lossDownstreamPct, "with nothing arriving there is no return path to blame")
}
// Jitter is legitimate without synchronised clocks because the offset cancels when successive
// one-way samples are differenced. This pins that: a huge constant offset must not show up.
@Test
fun jitterIsUnaffectedByTheClockOffsetBetweenTheTwoMachines() {
val (sent, near) = clean(10, offsetNs = 0)
val (_, far) = clean(10, offsetNs = 9_999_999_999L)
val a = Directional.analyse(sent, near)
val b = Directional.analyse(sent, far)
assertEquals(a.jitterUpstreamMs, b.jitterUpstreamMs,
"a constant clock offset must cancel when consecutive samples are differenced")
assertEquals(0.0, assertNotNull(a.jitterUpstreamMs), "an evenly spaced train has no jitter")
}
@Test
fun jitterReflectsUnevenArrival() {
val sent = listOf(
Sample(1, 0, 10_000_000),
Sample(2, 10_000_000, 20_000_000),
Sample(3, 20_000_000, 30_000_000),
)
// Server receive times drift: +2ms, +7ms, +3ms relative to send.
val seen = listOf(
ServerSighting(1, 2_000_000, 2_100_000),
ServerSighting(2, 17_000_000, 17_100_000),
ServerSighting(3, 23_000_000, 23_100_000),
)
val m = Directional.analyse(sent, seen)
// one-way samples: 2ms, 7ms, 3ms → |7-2| and |3-7| → mean 4.5ms
assertEquals(4.5, assertNotNull(m.jitterUpstreamMs))
}
// "No jitter" and "not enough data to say" are different claims, and only one is true here.
@Test
fun tooFewSamplesReportsNoJitterRatherThanZero() {
val m = Directional.analyse(
listOf(Sample(1, 0, 10_000_000)),
listOf(ServerSighting(1, 2_000_000, 2_100_000)),
)
assertNull(m.jitterUpstreamMs)
assertNull(m.jitterDownstreamMs)
}
// A server record for a sequence we never sent is not evidence about this train; folding it
// in would yield loss percentages outside 0100.
@Test
fun strayServerRecordsAreIgnored() {
val (sent, seen) = clean(5)
val m = Directional.analyse(sent, seen + ServerSighting(99, 1, 2) + ServerSighting(100, 3, 4))
assertEquals(5, m.seenByServer)
assertEquals(0.0, m.lossUpstreamPct)
assertTrue(m.lossDownstreamPct in 0.0..100.0)
}
@Test
fun anEmptyTrainDoesNotDivideByZero() {
val m = Directional.analyse(emptyList(), emptyList())
assertEquals(0.0, m.lossUpstreamPct)
assertEquals(0.0, m.lossDownstreamPct)
assertFalse(m.noneReachedServer, "nothing sent is not the same as nothing arriving")
}
}
@@ -7,6 +7,7 @@ import app.echo_lot.measurement.*
import kotlinx.serialization.json.Json
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertNotNull
import kotlin.test.assertTrue
/**
@@ -59,6 +60,18 @@ class LiveMeasurementTest {
println("metrics: $metrics")
assertTrue(metrics.toString().contains("rtt_ms_avg"))
// The directional split is the point of asking the server what it saw: without it a
// lossy path is reported as "loss" with no direction, which sends an engineer looking
// in both at once. Correlation is by wire sequence number, so a mismatch here means the
// two sides disagree about which packet is which.
val m = metrics.toString()
assertTrue(m.contains("seen_by_server"), "no directional split in the metrics: $m")
val seen = Regex(""""seen_by_server":(\d+)""").find(m)?.groupValues?.get(1)?.toInt()
assertNotNull(seen, "seen_by_server missing")
assertEquals(20, seen, "the server should have seen every probe on a healthy path")
assertTrue(m.contains("jitter_upstream_ms"), "no per-direction jitter: $m")
println("directional: $m")
assertTrue(doc.summary != null)
// A healthy local->fmr path should be green (no loss, no rebinding) or yellow.
println("summary: ${doc.summary}")
@@ -77,6 +77,8 @@ object TestType {
const val MTU_BLACKHOLE = "mtu.blackhole"
const val MTU_MSS_OBSERVED = "mtu.mss_observed"
const val MTU_FRAG_DELIVERY = "mtu.frag_delivery"
/** Whether fragments survive arriving out of order, not merely whether they survive. */
const val MTU_FRAG_ORDERING = "mtu.frag_ordering"
// nat
const val NAT_STUN_5780 = "nat.stun_5780"
const val NAT_MAPPING_LIFETIME_UDP = "nat.mapping_lifetime_udp"
@@ -44,13 +44,26 @@ class ProbeSession(
*/
fun echo(paddingBytes: Int = 40): EchoResult? {
val t0 = System.nanoTime()
val pkt = Wire.build(Wire.TYPE_ECHO_REQ, prefix, ++seq, nowNs(), key, ByteArray(paddingBytes))
val wireSeq = ++seq
val pkt = Wire.build(Wire.TYPE_ECHO_REQ, prefix, wireSeq, nowNs(), key, ByteArray(paddingBytes))
socket.send(DatagramPacket(pkt, pkt.size, server))
// A lost probe still has a sequence number, and that number is what lets the server's
// observations say whether it was lost going out or coming back — so report it either way.
lastSeq = wireSeq
val resp = receive(Wire.TYPE_ECHO_RESP) ?: return null
val rttMs = (System.nanoTime() - t0) / 1_000_000.0
return EchoResult(rttMs, Observation.parse(resp.payload))
return EchoResult(rttMs, Observation.parse(resp.payload), wireSeq)
}
/**
* The wire sequence number of the most recent [echo], including one that was lost.
*
* Exposed because the caller cannot derive it: the counter is shared with every other packet
* type on this session, so "the nth echo" is not "sequence n".
*/
var lastSeq: Int = 0
private set
/** One MTU probe of [totalSize] bytes (DF is set by the OS on the socket where supported).
* Returns the size the server acknowledged receiving, or null if the probe was lost. */
fun mtuProbe(totalSize: Int): Int? {
@@ -112,5 +125,5 @@ class ProbeSession(
override fun close() = socket.close()
data class EchoResult(val rttMs: Double, val observation: Observation?)
data class EchoResult(val rttMs: Double, val observation: Observation?, val seq: Int = 0)
}
@@ -32,6 +32,12 @@ object Wire {
const val TYPE_DOWNTRAIN_DATA: Int = 0x06
const val TYPE_BIG_SEND: Int = 0x0C
/**
* A datagram the server deliberately fragmented. Its arrival IS the measurement: it can only
* be delivered if every fragment survived the path and the local stack reassembled them.
*/
const val TYPE_FRAG_DATA: Int = 0x0D
/** The 8-byte on-the-wire prefix = first 16 hex chars of the session id, decoded. */
fun wirePrefix(sessionId: String): ByteArray {
require(sessionId.length >= 16) { "session id too short" }
+42
View File
@@ -0,0 +1,42 @@
#!/usr/bin/env bash
# SPDX-FileCopyrightText: 2026 Echolot contributors
# SPDX-License-Identifier: GPL-3.0-or-later
#
# Mints an enrollment link on the probe server and prints it — as text, as a QR code if
# `qrencode` is around, and as an adb command if a device is attached.
#
# The admin listener is localhost-only by design, so this goes over SSH. The link carries a
# single-use bearer token: treat it like a password until it is redeemed.
#
# Usage: echolot-app/scripts/enroll-link.sh [note]
set -euo pipefail
SSH_HOST="${ECHOLOT_SSH:-claude-echolot}"
NOTE="${1:-manual}"
MINTED=$(ssh -o BatchMode=yes "$SSH_HOST" \
"curl -s -X POST 'http://127.0.0.1:8444/admin/enroll-tokens?note=$NOTE'")
URI=$(printf '%s' "$MINTED" | python -c 'import json,sys;print(json.load(sys.stdin).get("enroll_uri",""))')
if [ -z "$URI" ]; then
echo "server returned no enroll_uri (needs server-v0.5.4+):" >&2
echo "$MINTED" >&2
exit 1
fi
echo "$URI"
echo
# A QR is the point of the format: scanning beats pasting a 200-character string onto a phone.
if command -v qrencode >/dev/null 2>&1; then
qrencode -t ANSIUTF8 "$URI"
else
echo "(install qrencode to get a scannable QR here)"
fi
# With a device attached, the deep link can be delivered straight to the app — no typing at all.
if command -v adb >/dev/null 2>&1 && [ -n "$(adb devices | sed -n '2p')" ]; then
echo
echo "attached device — deliver it directly with:"
echo " adb shell am start -a android.intent.action.VIEW -d '$URI'"
fi
+24 -1
View File
@@ -112,6 +112,14 @@ func serve(cfg *config.Config) error {
}
caps := []string{"udp-probe", "delayed-echo", "connect-back", "http-echo", "downtrain", "big-send"}
// Crafted fragments need a raw socket. Advertised only when one can actually be opened —
// a capability we cannot deliver turns a missing feature into a failed measurement.
rawFrag := dataplane.RawFragSupported()
if rawFrag {
caps = append(caps, "frag-send")
} else {
slog.Info("frag-send unavailable: no raw socket (needs CAP_NET_RAW)")
}
if len(config.Addrs(cfg.TCPListen)) > 0 {
caps = append(caps, "tcp-echo", "tls-echo")
}
@@ -154,6 +162,11 @@ func serve(cfg *config.Config) error {
AppRange: appRange,
PublicControlURL: publicControlURL(cfg),
}
// Left nil when there is no raw socket, so the handler answers "not implemented" with a
// reason rather than failing somewhere deeper.
if rawFrag {
ctl.FragSend = dp.FragSend
}
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer stop()
@@ -234,7 +247,17 @@ func serve(cfg *config.Config) error {
http.Error(w, err.Error(), 500)
return
}
fmt.Fprintf(w, `{"token":%q,"expires_in_s":86400}`+"\n", tok)
// The whole bootstrap, not just the token: this is what gets pasted or turned into a
// QR code, and assembling it here is what keeps an operator from transcribing a pin by
// hand — a pin wrong by one character fails as an inscrutable TLS error days later.
w.Header().Set("Content-Type", "application/json")
enc := json.NewEncoder(w)
enc.SetEscapeHTML(false) // the link is full of / and =; escaping them helps nobody
_ = enc.Encode(map[string]any{
"token": tok,
"expires_in_s": 86400,
"enroll_uri": ctl.EnrollmentLink(tok),
})
})
adminSrv := &http.Server{Addr: cfg.AdminListen, Handler: admin, ReadHeaderTimeout: 10 * time.Second}
go func() { errCh <- fmt.Errorf("admin: %w", adminSrv.ListenAndServe()) }()
+44
View File
@@ -59,6 +59,9 @@ type Server struct {
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
// FragSend emits one datagram as hand-built IP fragments in a chosen order (may be nil:
// needs a raw socket, so it is unavailable to an unprivileged server).
FragSend func(sess *session.Session, g *session.Grant, sizeBytes int, mode dataplane.FragMode, fragSize int) (dataplane.FragResult, error)
// 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.
@@ -241,6 +244,8 @@ func (s *Server) actions(w http.ResponseWriter, r *http.Request) {
IntervalUs int `json:"interval_us"`
SizesBytes []int `json:"sizes_bytes"`
DF *bool `json:"df"`
Mode string `json:"mode"`
FragBytes int `json:"frag_bytes"`
}
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "bad body"})
@@ -373,6 +378,45 @@ func (s *Server) actions(w http.ResponseWriter, r *http.Request) {
"grant": map[string]any{"max_bytes": g.MaxBytes, "max_kbps": g.MaxKbps},
})
case "frag_send":
if s.FragSend == nil {
writeJSON(w, http.StatusNotImplemented, map[string]string{
"error": "frag_send needs a raw socket, which this server does not have",
})
return
}
size := clamp(req.SizeBytes, 1600, 8000) // must exceed the path MTU or nothing fragments
mode := dataplane.FragMode(req.Mode)
switch mode {
case dataplane.FragInOrder, dataplane.FragReversed, dataplane.FragFirstLast:
default:
mode = dataplane.FragInOrder
}
fragBytes := clamp(req.FragBytes, 8, 1400)
g := sess.NewGrant(actionID, int64(size), 0, session.DefaultGrantLimits)
if g == nil {
writeJSON(w, http.StatusConflict, noDataPlaneYet)
return
}
// Synchronous: the whole burst is a few kB and at most a few hundred milliseconds, and
// the caller wants to know it was actually emitted before it starts listening. An
// asynchronous send would make "nothing arrived" ambiguous between a path drop and a
// send that never happened — the one distinction this test exists to make.
result, err := s.FragSend(sess, g, size, mode, fragBytes)
slog.Info("frag_send finished", "action", actionID, "mode", mode,
"size", size, "fragments", result.Fragments, "err", err)
if err != nil {
writeJSON(w, http.StatusConflict, map[string]any{
"error": err.Error(), "action_id": actionID, "result": result,
})
return
}
writeJSON(w, http.StatusAccepted, map[string]any{
"action_id": actionID, "mode": string(mode), "size_bytes": size,
"frag_bytes": fragBytes, "fragments": result.Fragments,
"grant": map[string]any{"max_bytes": g.MaxBytes, "max_kbps": g.MaxKbps},
})
default:
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "unknown or unimplemented action"})
}
+263
View File
@@ -0,0 +1,263 @@
// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
//go:build linux
package dataplane
import (
"encoding/binary"
"fmt"
"net/netip"
"sync/atomic"
"syscall"
"time"
"echo-lot.app/server/internal/session"
)
// Crafted IPv4 fragmentation (spec §5 frag_send).
//
// Letting the kernel fragment an oversized datagram — which is what big_send with df=false does —
// answers one question: do fragments get through at all. It cannot answer the more interesting
// one, because the kernel always emits fragments in order, first one first.
//
// The classic middlebox fault is precisely about that ordering. Only the *first* fragment carries
// the UDP header, and therefore the ports; a stateful firewall or NAT that has not seen it has no
// flow to match later fragments against. Plenty of implementations drop them. Others hold them
// briefly and reassemble; others leak. The difference is invisible to any test that sends
// fragments in order, and it shows up in the real world as "large DNS answers fail on this
// network" or "the VPN works until the MTU drops".
//
// So this builds the fragments by hand and controls their order and timing. That needs a raw
// socket (CAP_NET_RAW); when we do not have one the capability is not advertised, rather than
// advertised and failing later.
// FragMode is how a fragmented datagram is put on the wire.
type FragMode string
const (
// FragInOrder is the baseline: first fragment first, as the kernel would. A path that fails
// this fails everything, and it tells the others apart from a path that drops all fragments.
FragInOrder FragMode = "in_order"
// FragReversed sends the last fragment first. This is the one that finds stateful devices
// which need the first fragment to build state.
FragReversed FragMode = "reversed"
// FragFirstLast holds the first fragment back until the others have arrived, which tests
// whether the path buffers non-first fragments at all and for how long.
FragFirstLast FragMode = "first_last"
)
var fragIPID atomic.Uint32
// RawFragSupported reports whether crafted fragments can actually be sent here.
//
// Checked by opening the socket rather than by inspecting capabilities: the question is "will
// this work", and a permission model has more ways to say no than a capability bit has to say yes
// (user namespaces, seccomp, LSM). Advertising a capability we cannot deliver would turn a
// missing feature into a failed measurement.
func RawFragSupported() bool {
fd, err := syscall.Socket(syscall.AF_INET, syscall.SOCK_RAW, syscall.IPPROTO_RAW)
if err != nil {
return false
}
_ = syscall.Close(fd)
return true
}
// FragResult is what happened to one crafted fragment burst.
type FragResult struct {
Mode FragMode `json:"mode"`
SizeBytes int `json:"size_bytes"`
Fragments int `json:"fragments"`
Sent bool `json:"sent"`
Err string `json:"err,omitempty"`
}
// FragSend emits one ELT1 packet of sizeBytes as hand-built IPv4 fragments, in the given order.
//
// The datagram is assembled whole and then cut up, so what the client reassembles — if it
// reassembles — is a normal, HMAC-valid packet indistinguishable from any other. That matters:
// the client must not be able to tell a crafted fragment burst from a kernel one, or it would be
// measuring our sender rather than the path.
func (s *Server) FragSend(
sess *session.Session, g *session.Grant, sizeBytes int, mode FragMode, fragSize int,
) (FragResult, error) {
res := FragResult{Mode: mode, SizeBytes: sizeBytes}
target := sess.DataSource()
if !target.IsValid() {
return res, fmt.Errorf("no observed data-plane source")
}
if !target.Addr().Unmap().Is4() {
// IPv6 has no in-network fragmentation: only the source may fragment, via an extension
// header. Worth building, but it is a different mechanism and belongs in its own code
// path rather than pretending this one covers it.
return res, fmt.Errorf("crafted fragmentation is IPv4-only for now")
}
conn := s.connFor(target, sess.DataLocal())
if conn == nil {
return res, fmt.Errorf("no data-plane socket matches target family")
}
local := sess.DataLocal()
if !local.IsValid() {
return res, fmt.Errorf("session has no recorded local address")
}
if sizeBytes < HeaderSize+8 {
sizeBytes = HeaderSize + 8
}
if sizeBytes > 8000 {
sizeBytes = 8000
}
if !g.Allow(sizeBytes) {
return res, fmt.Errorf("grant exhausted")
}
// The ELT1 packet, signed exactly as any other, then wrapped in UDP.
payload := make([]byte, sizeBytes-HeaderSize)
binary.BigEndian.PutUint32(payload[0:4], uint32(sizeBytes))
copy(payload[4:], mode)
elt := s.buildPacket(sess, TypeFragData, 0, payload)
udp := buildUDP(local, target, elt)
// Fragment offsets are in 8-byte units, so every fragment except the last must be a multiple
// of 8. A payload that is not is not an error — it is a fragment that no host will reassemble.
if fragSize <= 0 {
fragSize = 576
}
fragSize = (fragSize / 8) * 8
if fragSize < 8 {
fragSize = 8
}
fragments := splitIPv4(local.Addr(), target.Addr(), udp, fragSize, uint16(fragIPID.Add(1)))
res.Fragments = len(fragments)
fd, err := syscall.Socket(syscall.AF_INET, syscall.SOCK_RAW, syscall.IPPROTO_RAW)
if err != nil {
res.Err = err.Error()
return res, err
}
defer syscall.Close(fd)
if err := syscall.SetsockoptInt(fd, syscall.IPPROTO_IP, syscall.IP_HDRINCL, 1); err != nil {
res.Err = err.Error()
return res, err
}
dst := syscall.SockaddrInet4{}
copy(dst.Addr[:], target.Addr().Unmap().AsSlice())
send := func(pkt []byte) error { return syscall.Sendto(fd, pkt, 0, &dst) }
switch mode {
case FragReversed:
for i := len(fragments) - 1; i >= 0; i-- {
if err := send(fragments[i]); err != nil {
res.Err = err.Error()
return res, err
}
time.Sleep(time.Millisecond)
}
case FragFirstLast:
for i := 1; i < len(fragments); i++ {
if err := send(fragments[i]); err != nil {
res.Err = err.Error()
return res, err
}
time.Sleep(time.Millisecond)
}
// Long enough to be a real test of whether anything holds fragments, short enough to stay
// inside the usual 30-second reassembly timeout by a wide margin.
time.Sleep(250 * time.Millisecond)
if err := send(fragments[0]); err != nil {
res.Err = err.Error()
return res, err
}
default:
for _, f := range fragments {
if err := send(f); err != nil {
res.Err = err.Error()
return res, err
}
time.Sleep(time.Millisecond)
}
}
res.Sent = true
return res, nil
}
// buildUDP wraps a payload in a UDP header with a computed checksum.
//
// The checksum is optional in IPv4 and it would be less code to send zero, but a zero-checksum
// datagram is dropped by some middleboxes — and that drop would be recorded as a fragmentation
// failure, which is exactly the wrong conclusion.
func buildUDP(src, dst netip.AddrPort, payload []byte) []byte {
out := make([]byte, 8+len(payload))
binary.BigEndian.PutUint16(out[0:2], src.Port())
binary.BigEndian.PutUint16(out[2:4], dst.Port())
binary.BigEndian.PutUint16(out[4:6], uint16(8+len(payload)))
copy(out[8:], payload)
// Pseudo-header + UDP header + data, per RFC 768.
var sum uint32
s4, d4 := src.Addr().Unmap().As4(), dst.Addr().Unmap().As4()
for _, b := range [][]byte{s4[:], d4[:]} {
sum += uint32(binary.BigEndian.Uint16(b[0:2]))
sum += uint32(binary.BigEndian.Uint16(b[2:4]))
}
sum += uint32(syscall.IPPROTO_UDP)
sum += uint32(len(out))
for i := 0; i+1 < len(out); i += 2 {
sum += uint32(binary.BigEndian.Uint16(out[i : i+2]))
}
if len(out)%2 == 1 {
sum += uint32(out[len(out)-1]) << 8
}
for sum>>16 != 0 {
sum = (sum & 0xFFFF) + (sum >> 16)
}
ck := ^uint16(sum)
if ck == 0 {
ck = 0xFFFF // 0 means "no checksum" in IPv4; the all-ones form is the same value
}
binary.BigEndian.PutUint16(out[6:8], ck)
return out
}
// splitIPv4 cuts a UDP datagram into IPv4 fragments of at most fragSize payload bytes each.
//
// Every fragment carries the same IP ID — that is what marks them as one datagram — and every one
// but the last sets MF. The kernel fills in the header checksum and total length for us under
// IP_HDRINCL (raw(7)); the ID it only fills when zero, which is why it is set explicitly here.
func splitIPv4(src, dst netip.Addr, udp []byte, fragSize int, id uint16) [][]byte {
s4, d4 := src.Unmap().As4(), dst.Unmap().As4()
var out [][]byte
for off := 0; off < len(udp); off += fragSize {
end := off + fragSize
if end > len(udp) {
end = len(udp)
}
chunk := udp[off:end]
more := end < len(udp)
hdr := make([]byte, 20, 20+len(chunk))
hdr[0] = 0x45 // IPv4, 5 words of header
hdr[1] = 0 // DSCP/ECN
binary.BigEndian.PutUint16(hdr[2:4], uint16(20+len(chunk)))
binary.BigEndian.PutUint16(hdr[4:6], id)
flagsOff := uint16(off / 8)
if more {
flagsOff |= 0x2000 // MF
}
binary.BigEndian.PutUint16(hdr[6:8], flagsOff)
hdr[8] = 64 // TTL
hdr[9] = syscall.IPPROTO_UDP
// hdr[10:12] checksum left zero: the kernel computes it under IP_HDRINCL.
copy(hdr[12:16], s4[:])
copy(hdr[16:20], d4[:])
out = append(out, append(hdr, chunk...))
}
return out
}
@@ -0,0 +1,160 @@
// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
//go:build linux
package dataplane
import (
"encoding/binary"
"net/netip"
"testing"
)
// Fragment headers are the kind of thing that is either exactly right or silently useless: a
// wrong offset unit, a missing MF bit or a bad checksum produces packets that leave the machine
// and are dropped by the receiver's IP stack without a word. Nothing downstream would notice —
// the client would simply record "fragments do not get through", which is a wrong answer rather
// than a missing one. Hence these check the bytes.
func testAddrs() (netip.AddrPort, netip.AddrPort) {
return netip.MustParseAddrPort("192.0.2.1:8442"), netip.MustParseAddrPort("198.51.100.9:41000")
}
func TestSplitCoversThePayloadExactlyOnce(t *testing.T) {
src, dst := testAddrs()
udp := buildUDP(src, dst, make([]byte, 2000))
frags := splitIPv4(src.Addr(), dst.Addr(), udp, 576, 0x1234)
if len(frags) < 3 {
t.Fatalf("expected several fragments for %d bytes, got %d", len(udp), len(frags))
}
// Reassemble the way a receiver would: place each fragment's payload at its offset.
rebuilt := make([]byte, len(udp))
covered := make([]bool, len(udp))
for _, f := range frags {
flagsOff := binary.BigEndian.Uint16(f[6:8])
off := int(flagsOff&0x1FFF) * 8
body := f[20:]
if off+len(body) > len(udp) {
t.Fatalf("fragment at offset %d overruns the datagram", off)
}
for i, b := range body {
if covered[off+i] {
t.Fatalf("byte %d delivered twice", off+i)
}
covered[off+i] = true
rebuilt[off+i] = b
}
}
for i, c := range covered {
if !c {
t.Fatalf("byte %d was never sent", i)
}
}
for i := range udp {
if rebuilt[i] != udp[i] {
t.Fatalf("reassembled byte %d differs", i)
}
}
}
func TestFragmentHeadersAreWellFormed(t *testing.T) {
src, dst := testAddrs()
udp := buildUDP(src, dst, make([]byte, 3000))
frags := splitIPv4(src.Addr(), dst.Addr(), udp, 800, 0xBEEF)
for i, f := range frags {
if got := f[0]; got != 0x45 {
t.Errorf("fragment %d: version/IHL = %#x, want 0x45", i, got)
}
if got := f[9]; got != 17 {
t.Errorf("fragment %d: protocol = %d, want 17 (UDP)", i, got)
}
if got := binary.BigEndian.Uint16(f[4:6]); got != 0xBEEF {
t.Errorf("fragment %d: IP ID = %#x — all fragments of one datagram must share it", i, got)
}
if got := binary.BigEndian.Uint16(f[2:4]); int(got) != len(f) {
t.Errorf("fragment %d: total length = %d, actual %d", i, got, len(f))
}
flagsOff := binary.BigEndian.Uint16(f[6:8])
mf := flagsOff&0x2000 != 0
wantMF := i < len(frags)-1
if mf != wantMF {
t.Errorf("fragment %d: MF = %v, want %v", i, mf, wantMF)
}
}
}
// Offsets are counted in 8-byte units, so every fragment but the last must be a multiple of 8.
// A 100-byte "fragment size" that silently becomes 100 bytes on the wire produces a datagram no
// host will ever reassemble.
func TestNonFinalFragmentsAreEightByteMultiples(t *testing.T) {
src, dst := testAddrs()
udp := buildUDP(src, dst, make([]byte, 2500))
for _, size := range []int{8, 100, 576, 999, 1400} {
frags := splitIPv4(src.Addr(), dst.Addr(), udp, (size/8)*8, 1)
for i, f := range frags[:len(frags)-1] {
if body := len(f) - 20; body%8 != 0 {
t.Errorf("size %d: non-final fragment %d carries %d bytes, not a multiple of 8",
size, i, body)
}
}
}
}
// The UDP checksum is optional in IPv4, and sending zero would be less code — but a
// zero-checksum datagram is dropped by some middleboxes, and that drop would be recorded as a
// fragmentation failure. So it must be present and correct.
func TestUDPChecksumVerifies(t *testing.T) {
src, dst := testAddrs()
for _, n := range []int{0, 1, 7, 8, 100, 1001} { // odd lengths exercise the tail-byte path
udp := buildUDP(src, dst, make([]byte, n))
if got := binary.BigEndian.Uint16(udp[6:8]); got == 0 {
t.Fatalf("payload %d: checksum is zero, which means 'not computed'", n)
}
if sum := verifyUDPChecksum(src.Addr(), dst.Addr(), udp); sum != 0xFFFF {
t.Errorf("payload %d: checksum does not verify (one's complement sum %#x)", n, sum)
}
if got := binary.BigEndian.Uint16(udp[4:6]); int(got) != len(udp) {
t.Errorf("payload %d: UDP length field %d, actual %d", n, got, len(udp))
}
}
}
func TestUDPPortsComeFromTheSessionAddresses(t *testing.T) {
src, dst := testAddrs()
udp := buildUDP(src, dst, []byte("x"))
if got := binary.BigEndian.Uint16(udp[0:2]); got != src.Port() {
t.Errorf("source port = %d, want %d", got, src.Port())
}
// The destination port must be the client's observed source port, or the datagram arrives
// at the machine and is discarded before any socket sees it.
if got := binary.BigEndian.Uint16(udp[2:4]); got != dst.Port() {
t.Errorf("destination port = %d, want %d", got, dst.Port())
}
}
// Recomputes the one's complement sum over the pseudo-header and datagram; a correct checksum
// makes the total 0xFFFF.
func verifyUDPChecksum(src, dst netip.Addr, udp []byte) uint16 {
var sum uint32
s4, d4 := src.Unmap().As4(), dst.Unmap().As4()
for _, b := range [][]byte{s4[:], d4[:]} {
sum += uint32(binary.BigEndian.Uint16(b[0:2]))
sum += uint32(binary.BigEndian.Uint16(b[2:4]))
}
sum += 17
sum += uint32(len(udp))
for i := 0; i+1 < len(udp); i += 2 {
sum += uint32(binary.BigEndian.Uint16(udp[i : i+2]))
}
if len(udp)%2 == 1 {
sum += uint32(udp[len(udp)-1]) << 8
}
for sum>>16 != 0 {
sum = (sum & 0xFFFF) + (sum >> 16)
}
return uint16(sum)
}
+41
View File
@@ -0,0 +1,41 @@
// SPDX-FileCopyrightText: 2026 Echolot contributors
// SPDX-License-Identifier: GPL-3.0-or-later
//go:build !linux
package dataplane
import (
"fmt"
"echo-lot.app/server/internal/session"
)
// Crafting IP fragments needs a raw socket and Linux's IP_HDRINCL semantics. Off Linux the
// capability is simply not advertised, so a client never asks for it — better than answering
// with a measurement we cannot actually make.
type FragMode string
const (
FragInOrder FragMode = "in_order"
FragReversed FragMode = "reversed"
FragFirstLast FragMode = "first_last"
)
type FragResult struct {
Mode FragMode `json:"mode"`
SizeBytes int `json:"size_bytes"`
Fragments int `json:"fragments"`
Sent bool `json:"sent"`
Err string `json:"err,omitempty"`
}
func RawFragSupported() bool { return false }
func (s *Server) FragSend(
sess *session.Session, g *session.Grant, sizeBytes int, mode FragMode, fragSize int,
) (FragResult, error) {
return FragResult{Mode: mode, SizeBytes: sizeBytes},
fmt.Errorf("crafted fragmentation is only implemented on Linux")
}
+14 -2
View File
@@ -35,6 +35,8 @@ const (
// Server->client under an asymmetric grant (spec §3.4/§5).
TypeDownTrainData = 0x06
TypeBigSend = 0x0C
// TypeFragData is delivered only after IP reassembly, so its arrival IS the measurement.
TypeFragData = 0x0D
)
type Server struct {
@@ -232,6 +234,17 @@ func (s *Server) send(conn *net.UDPConn, raddr netip.AddrPort, sess *session.Ses
// EMSGSIZE means our own egress MTU refused the datagram, which is a different fact from the
// client not receiving it.
func (s *Server) sendErr(conn *net.UDPConn, raddr netip.AddrPort, sess *session.Session, typ byte, seq uint32, payload []byte) error {
pkt := s.buildPacket(sess, typ, seq, payload)
_, err := conn.WriteToUDPAddrPort(pkt, raddr)
return err
}
// buildPacket assembles and signs an ELT1 packet without sending it.
//
// Split out for the crafted-fragment path, which needs the bytes so it can cut them up itself.
// What arrives after reassembly must be indistinguishable from an ordinary packet, or the client
// would be measuring our sender rather than the path — so it goes through exactly this function.
func (s *Server) buildPacket(sess *session.Session, typ byte, seq uint32, payload []byte) []byte {
pkt := make([]byte, HeaderSize+len(payload))
copy(pkt[0:4], Magic)
pkt[4] = typ
@@ -247,8 +260,7 @@ func (s *Server) sendErr(conn *net.UDPConn, raddr netip.AddrPort, sess *session.
mac.Write(pkt[0:28])
mac.Write(payload)
copy(pkt[28:32], mac.Sum(nil)[:4])
_, err := conn.WriteToUDPAddrPort(pkt, raddr)
return err
return pkt
}
func hexByte(hi, lo byte) byte {