throughput: the upstream direction, counted by the only party that can
server-release / image (push) Successful in 15s
server-test / test (push) Successful in 33s
server-release / release (push) Successful in 33s

The client generates the traffic and the server counts it. No grant is involved
- the client is sending its own packets, so there is nothing to amplify - but it
does need the server's tally, because only the far end knows how much arrived.
Without that number a sender measures how fast it can transmit, which is usually
just the speed of the local NIC and is not the question being asked.

A new wire type the server counts and deliberately never answers: a reply would
double the traffic and drag the return path into a measurement that is
specifically about the outbound one.

The tally is a counter, not a list, and short-circuits before the observation
log. A five-second run at 20 Mbps is around ten thousand packets; one struct
each would turn a measurement into an allocation storm on a shared server, and
nothing needs the per-packet detail since the client holds the send-side record.
The gap between the two counts is the loss.

direction=up on the throughput action sends nothing - it zeroes the counter, so
a second run in one session measures itself instead of inheriting the first.

Same honesty rule as downstream: measures_network is false when what arrived
matches what was offered, because then the path was never the constraint.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
mrambossek
2026-08-01 15:54:09 +02:00
co-authored by Claude Fable 5
parent 8646bab52d
commit 892e952a8e
6 changed files with 277 additions and 2 deletions
@@ -147,6 +147,117 @@ class ThroughputMeasurement(private val ids: IdSource) {
) to findings
}
/**
* Upstream throughput: the client sends, the server counts.
*
* The mirror image of the downstream case, and it needs no grant — the client is generating
* its own traffic, so there is no amplification to gate. What it does need is the server's
* count: only the far end knows how much arrived, and without that number a sender can
* measure how fast it can *transmit*, which is not the same question and is usually just the
* speed of the local NIC.
*/
fun runUpstream(
credential: String,
sessionId: String,
control: ControlClient,
probe: ProbeSession,
sessionRef: String,
durationS: Int = 5,
kbps: Int = 20_000,
sizeBytes: Int = 1200,
): Pair<Test, List<Finding>> {
val testId = ids.uuid()
val started = ids.monoNs()
// Zeroes the server's counter so this run measures itself rather than inheriting the
// packets of an earlier one on the same session.
val reply = runCatching {
control.action(credential, sessionId, """{"action":"throughput","direction":"up"}""")
}
if (reply.isFailure) {
return Test(
id = testId, type = TestType.PERF_THROUGHPUT_UDP, sessionRef = sessionRef, tier = Tier.APP,
startedMonoNs = started, endedMonoNs = ids.monoNs(),
status = TestStatus.UNSUPPORTED,
error = TestError("action_refused", reply.exceptionOrNull()?.message ?: "refused"),
) to emptyList()
}
val sent = probe.sendThroughput(durationS * 1000L, kbps, sizeBytes)
// A moment for the tail of the run to arrive; counting still-in-flight packets as lost
// would inflate the loss figure by whatever the path's delay happens to be.
Thread.sleep(500)
val seen = upstreamCount(control, credential, sessionId)
val lossPct = if (sent.packets > 0 && seen != null) {
round2((sent.packets - seen.packets).coerceAtLeast(0) * 100.0 / sent.packets)
} else {
null
}
// The receiver's rate is the measurement. The sender's is what we managed to emit, which
// is a property of this phone and its radio, not of the network.
val achievedKbps = seen?.kbps ?: 0
val metrics = json.encodeToJsonElement(
UpstreamThroughputMetrics(
requestedKbps = kbps,
sentPackets = sent.packets,
sentBytes = sent.bytes,
sentKbps = sent.kbps,
receivedPackets = seen?.packets,
receivedBytes = seen?.bytes,
receivedKbps = achievedKbps,
lossPct = lossPct,
// Same honesty rule as downstream: if what arrived matches what we offered, the
// path was never the constraint and this number says nothing about it.
measuresNetwork = seen != null && achievedKbps > 0 && achievedKbps < sent.kbps * 9 / 10,
),
) as JsonObject
val findings = ArrayList<Finding>()
if (seen != null && seen.packets == 0 && sent.packets > 0) {
findings.add(
finding(
FindingRegistry.THROUGHPUT_NO_DELIVERY, testId,
"No upstream traffic reached the server",
"This device sent ${sent.packets} packets and the server received none. " +
"That is a connectivity fault on the outbound path rather than a slow link.",
),
)
} else if (lossPct != null && lossPct >= 2.0) {
findings.add(
finding(
FindingRegistry.THROUGHPUT_BELOW_OFFERED, testId,
"Upstream loss of $lossPct % at ${sent.kbps / 1000} Mbit/s",
"The server received ${seen?.packets} of the ${sent.packets} packets this " +
"device sent. The outbound path could not carry what was offered.",
),
)
}
return Test(
id = testId, type = TestType.PERF_THROUGHPUT_UDP, sessionRef = sessionRef, tier = Tier.APP,
startedMonoNs = started, endedMonoNs = ids.monoNs(),
status = if (seen == null || seen.packets == 0) TestStatus.FAILED else TestStatus.OK,
metrics = metrics,
) to findings
}
private data class UpstreamCount(val packets: Int, val bytes: Long, val kbps: Int)
/** The server's tally for this session's upstream run. */
private fun upstreamCount(
control: ControlClient, credential: String, sessionId: String,
): UpstreamCount? = runCatching {
val o = Json.parseToJsonElement(control.observations(credential, sessionId))
.jsonObject["throughput_up"]?.jsonObject ?: return null
UpstreamCount(
packets = o["packets"]?.jsonPrimitive?.content?.toIntOrNull() ?: 0,
bytes = o["bytes"]?.jsonPrimitive?.content?.toLongOrNull() ?: 0,
kbps = o["kbps"]?.jsonPrimitive?.content?.toIntOrNull() ?: 0,
)
}.getOrNull()
private data class SenderReport(
val packets: Int, val bytes: Long, val kbps: Int, val limitedBy: String,
)
@@ -186,9 +297,27 @@ class ThroughputMeasurement(private val ids: IdSource) {
private fun round2(v: Double) = Math.round(v * 100.0) / 100.0
}
/** Metrics for perf.throughput_udp in the upstream direction. */
@Serializable
data class UpstreamThroughputMetrics(
val direction: String = "up",
@SerialName("requested_kbps") val requestedKbps: Int,
@SerialName("sent_packets") val sentPackets: Int,
@SerialName("sent_bytes") val sentBytes: Long,
/** What this device managed to emit — a property of the phone and its radio, not the path. */
@SerialName("sent_kbps") val sentKbps: Int,
@SerialName("received_packets") val receivedPackets: Int? = null,
@SerialName("received_bytes") val receivedBytes: Long? = null,
/** What arrived, measured by the only party that can measure it. This is the result. */
@SerialName("received_kbps") val receivedKbps: Int,
@SerialName("loss_pct") val lossPct: Double? = null,
@SerialName("measures_network") val measuresNetwork: Boolean,
)
/** Metrics for perf.throughput_udp. */
@Serializable
data class ThroughputMetrics(
val direction: String = "down",
@SerialName("requested_kbps") val requestedKbps: Int,
@SerialName("planned_duration_ms") val plannedDurationMs: Int,
@SerialName("packets_received") val packetsReceived: Int,
@@ -109,6 +109,49 @@ class ProbeSession(
return out
}
/**
* Sends paced upstream traffic for [durationMs] and reports what was put on the wire.
*
* Paced rather than flat out, for the same reason the server paces: an unpaced burst measures
* the local NIC and the first queue it meets, then collapses into loss that reads as a network
* fault. The schedule is absolute rather than sleep-per-packet, which accumulates the
* scheduler's error and drifts the achieved rate below target over a multi-second run.
*
* Nothing comes back — the server counts and stays silent — so the result here is only the
* send side. The measurement is the gap between this and the server's tally.
*/
fun sendThroughput(durationMs: Long, kbps: Int, sizeBytes: Int = 1200): Sent {
val size = sizeBytes.coerceIn(Wire.HEADER_SIZE + 16, 1472)
val payload = ByteArray(size - Wire.HEADER_SIZE)
val perPacketNs = (size.toLong() * 8 * 1_000_000 / kbps.coerceAtLeast(1)).coerceAtLeast(1_000)
val start = System.nanoTime()
val deadline = start + durationMs * 1_000_000
var next = start
var packets = 0
var bytes = 0L
while (System.nanoTime() < deadline) {
val pkt = Wire.build(Wire.TYPE_THROUGHPUT_UP, prefix, ++seq, nowNs(), key, payload)
try {
socket.send(DatagramPacket(pkt, pkt.size, server))
} catch (e: java.io.IOException) {
// A local send failure is our condition, not the path's. Stop and report what
// actually left, rather than counting the remainder as loss on the network.
break
}
packets++
bytes += pkt.size
next += perPacketNs
val sleepNs = next - System.nanoTime()
if (sleepNs > 0) Thread.sleep(sleepNs / 1_000_000, (sleepNs % 1_000_000).toInt())
}
val elapsedMs = (System.nanoTime() - start) / 1_000_000
return Sent(packets, bytes, elapsedMs, if (elapsedMs > 0) (bytes * 8 / elapsedMs).toInt() else 0)
}
/** What one upstream run put on the wire locally. */
data class Sent(val packets: Int, val bytes: Long, val durationMs: Long, val kbps: Int)
/** 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)
@@ -41,6 +41,13 @@ object Wire {
/** One packet of a sustained-rate downstream run. */
const val TYPE_THROUGHPUT_DATA: Int = 0x0E
/**
* One packet of a client-driven upstream run. The server counts it and does not answer:
* a reply would double the traffic and drag the return path into a measurement that is
* specifically about the outbound one.
*/
const val TYPE_THROUGHPUT_UP: Int = 0x0F
/** 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" }