app: core-engine — run engine; full server-facing vertical proven vs fmr
Composes core-protocol probes into core-measurement documents. Injected clock/UUID source keeps it pure and unit-testable. Runs a server ECHO train and derives RTT distribution, loss, and NAT-rebinding detection (from the server's observed source port) as train.udp_updown, then findings + a §7.3 summary. Verified end-to-end against fmr: 20-packet train, 0% loss, RTT 1.7/2.5/6.9ms, no rebinding → valid MeasurementDocument (2.3kB), overall GREEN. The whole server-facing stack (protocol → engine → schema → verdict) now produces the real product artifact against the live server, no device required. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5
parent
1f8860f7f8
commit
49c6197aff
@@ -0,0 +1,29 @@
|
||||
// SPDX-FileCopyrightText: 2026 Echolot contributors
|
||||
// SPDX-License-Identifier: GPL-3.0-or-later
|
||||
|
||||
plugins {
|
||||
alias(libs.plugins.kotlin.jvm)
|
||||
alias(libs.plugins.kotlin.serialization)
|
||||
}
|
||||
|
||||
// The measurement run engine: composes core-protocol probes into
|
||||
// core-measurement documents. Pure Kotlin/JVM, so it is unit-testable and can
|
||||
// run a full server-facing measurement against a live server.
|
||||
dependencies {
|
||||
implementation(project(":core-protocol"))
|
||||
implementation(project(":core-measurement"))
|
||||
implementation(libs.kotlinx.serialization.json)
|
||||
testImplementation(kotlin("test"))
|
||||
}
|
||||
|
||||
kotlin {
|
||||
jvmToolchain(21)
|
||||
compilerOptions { jvmTarget.set(org.jetbrains.kotlin.gradle.dsl.JvmTarget.JVM_17) }
|
||||
}
|
||||
java { sourceCompatibility = JavaVersion.VERSION_17; targetCompatibility = JavaVersion.VERSION_17 }
|
||||
|
||||
tasks.test {
|
||||
useJUnitPlatform()
|
||||
listOf("ECHOLOT_LIVE_URL","ECHOLOT_LIVE_PIN","ECHOLOT_LIVE_CRED","ECHOLOT_LIVE_UDP","ECHOLOT_LIVE_TARGET")
|
||||
.forEach { k -> System.getenv(k)?.let { environment(k, it) } }
|
||||
}
|
||||
@@ -0,0 +1,166 @@
|
||||
// SPDX-FileCopyrightText: 2026 Echolot contributors
|
||||
// SPDX-License-Identifier: GPL-3.0-or-later
|
||||
|
||||
// Package engine composes core-protocol probes into core-measurement documents — the run engine
|
||||
// the app drives. This file covers the server-facing vertical (control plane + UDP data plane);
|
||||
// device-tier probes (link snapshot, Shizuku, local discovery) plug in from the Android modules.
|
||||
package app.echo_lot.engine
|
||||
|
||||
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.JsonObject
|
||||
import kotlinx.serialization.json.encodeToJsonElement
|
||||
|
||||
/**
|
||||
* Runs the server-facing measurements against one target and assembles a [MeasurementDocument]:
|
||||
* an ECHO train (RTT distribution, loss, and NAT-rebinding detection from the server's observed
|
||||
* source port) as `train.udp_updown`. Everything is real evidence with recomputable metrics, and
|
||||
* findings are derived deterministically. IDs/timestamps are injected so the engine stays pure
|
||||
* (no clocks/UUIDs of its own) and unit-testable.
|
||||
*/
|
||||
class ServerMeasurement(
|
||||
private val ids: IdSource,
|
||||
private val app: AppInfo,
|
||||
private val device: DeviceInfo,
|
||||
) {
|
||||
private val json = Json { encodeDefaults = true; explicitNulls = true }
|
||||
|
||||
data class Config(
|
||||
val controlUrl: String,
|
||||
val pins: Set<String>,
|
||||
val credential: String,
|
||||
val target: String,
|
||||
val udpHost: String,
|
||||
val udpPort: Int,
|
||||
val echoCount: Int = 20,
|
||||
val echoPaddingBytes: Int = 64,
|
||||
)
|
||||
|
||||
fun run(cfg: Config): MeasurementDocument {
|
||||
val runId = ids.uuid()
|
||||
val startWall = ids.nowWall()
|
||||
val startMono = ids.monoNs()
|
||||
|
||||
val control = ControlClient(cfg.controlUrl, cfg.pins)
|
||||
val profile = control.profile(cfg.credential)
|
||||
val session = control.createSession(cfg.credential, cfg.target)
|
||||
|
||||
val serverSession = ServerSession(
|
||||
id = "sess-1",
|
||||
profileName = profile.name,
|
||||
controlUrl = cfg.controlUrl,
|
||||
serverVersion = profile.serverVersion,
|
||||
capabilities = profile.capabilities,
|
||||
sessionId = session.sessionId,
|
||||
target = SessionTarget(ip4 = cfg.udpHost, udpPort = cfg.udpPort),
|
||||
)
|
||||
|
||||
val (test, findings) = echoTrain(cfg, control, session, startMono)
|
||||
|
||||
control.deleteSession(cfg.credential, session.sessionId)
|
||||
|
||||
val summary = Verdicts.derive(listOf(test), findings)
|
||||
return MeasurementDocument(
|
||||
run = Run(
|
||||
id = runId, trigger = Trigger.MANUAL, startedAt = startWall, endedAt = ids.nowWall(),
|
||||
clock = Clock(monoOriginWall = startWall),
|
||||
app = app, device = device,
|
||||
tiers = Tiers(app = true),
|
||||
),
|
||||
serverSessions = listOf(serverSession),
|
||||
tests = listOf(test),
|
||||
findings = findings,
|
||||
summary = summary,
|
||||
)
|
||||
}
|
||||
|
||||
private fun echoTrain(
|
||||
cfg: Config, control: ControlClient, session: app.echo_lot.protocol.SessionResponse, startMono: Long,
|
||||
): Pair<Test, List<Finding>> {
|
||||
val testId = ids.uuid()
|
||||
val seqs = ArrayList<Int>()
|
||||
val tTx = ArrayList<Long?>()
|
||||
val tRx = ArrayList<Long?>()
|
||||
val sizes = ArrayList<Int>()
|
||||
val rtts = ArrayList<Double>()
|
||||
val observedPorts = LinkedHashSet<Int>()
|
||||
|
||||
ProbeSession(cfg.credential, session, cfg.udpHost, cfg.udpPort).use { ps ->
|
||||
for (i in 0 until cfg.echoCount) {
|
||||
val txMono = ids.monoNs() - startMono
|
||||
val r = ps.echo(cfg.echoPaddingBytes)
|
||||
seqs.add(i)
|
||||
tTx.add(txMono)
|
||||
sizes.add(Wire_HEADER + cfg.echoPaddingBytes)
|
||||
if (r != null) {
|
||||
tRx.add(ids.monoNs() - startMono)
|
||||
rtts.add(r.rttMs)
|
||||
r.observation?.observedPort?.let { observedPorts.add(it) }
|
||||
} else {
|
||||
tRx.add(null)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
val sent = cfg.echoCount
|
||||
val received = rtts.size
|
||||
val lossPct = if (sent == 0) 0.0 else (sent - received) * 100.0 / sent
|
||||
val natRebinding = observedPorts.size > 1
|
||||
|
||||
val evidence: JsonObject = TrainEvidence(
|
||||
epochMonoNs = startMono, seq = seqs, tTxNs = tTx, tRxNs = tRx, sizeBytes = sizes,
|
||||
).toEvidence()
|
||||
|
||||
val metrics: JsonObject = json.encodeToJsonElement(
|
||||
EchoMetrics(
|
||||
sent = sent, received = received, lossPct = round1(lossPct),
|
||||
rttMsMin = rtts.minOrNull()?.let(::round1),
|
||||
rttMsAvg = rtts.average().takeIf { received > 0 }?.let(::round1),
|
||||
rttMsMax = rtts.maxOrNull()?.let(::round1),
|
||||
observedPorts = observedPorts.toList(),
|
||||
natRebindingDetected = natRebinding,
|
||||
)
|
||||
) as JsonObject
|
||||
|
||||
val status = when {
|
||||
received == 0 -> TestStatus.FAILED
|
||||
received < sent -> TestStatus.PARTIAL
|
||||
else -> TestStatus.OK
|
||||
}
|
||||
val test = Test(
|
||||
id = testId, type = TestType.TRAIN_UDP_UPDOWN, sessionRef = "sess-1", tier = Tier.APP,
|
||||
startedMonoNs = startMono, endedMonoNs = ids.monoNs(), status = status,
|
||||
evidence = evidence, metrics = metrics,
|
||||
)
|
||||
|
||||
val findings = ArrayList<Finding>()
|
||||
if (received == 0) {
|
||||
findings.add(finding("nat.udp_unreachable", Category.CONNECTIVITY, Severity.HIGH, testId,
|
||||
"No UDP echo replies from the server",
|
||||
"Every ECHO probe to the server's UDP data plane was lost — the path blocks or drops the session's UDP traffic."))
|
||||
} else if (lossPct >= 20.0) {
|
||||
findings.add(finding("connectivity.udp_loss", Category.CONNECTIVITY, Severity.MEDIUM, testId,
|
||||
"High UDP loss to the server (${round1(lossPct)}%)",
|
||||
"A large fraction of ECHO probes were lost, indicating an unreliable UDP path."))
|
||||
}
|
||||
if (natRebinding) {
|
||||
findings.add(finding("nat.udp_rebinding", Category.NAT, Severity.MEDIUM, testId,
|
||||
"NAT remapped the UDP source port mid-flow",
|
||||
"The server observed more than one source port for this session (${observedPorts.joinToString()}), i.e. a NAT with a short UDP mapping or per-packet remapping."))
|
||||
}
|
||||
return test to findings
|
||||
}
|
||||
|
||||
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,
|
||||
title = title, description = desc, evidenceRefs = listOf(EvidenceRef(testId)),
|
||||
)
|
||||
|
||||
private companion object {
|
||||
const val Wire_HEADER = 32
|
||||
fun round1(v: Double) = Math.round(v * 10.0) / 10.0
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,38 @@
|
||||
// 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
|
||||
import java.time.Instant
|
||||
import java.util.UUID
|
||||
|
||||
/**
|
||||
* Clock/ID source, injected so the engine has no hidden nondeterminism and stays unit-testable.
|
||||
* The default uses wall + monotonic clocks and random UUIDs; tests supply deterministic ones.
|
||||
*/
|
||||
interface IdSource {
|
||||
fun uuid(): String
|
||||
fun monoNs(): Long
|
||||
fun nowWall(): String
|
||||
}
|
||||
|
||||
class SystemIdSource : IdSource {
|
||||
override fun uuid(): String = UUID.randomUUID().toString()
|
||||
override fun monoNs(): Long = System.nanoTime()
|
||||
override fun nowWall(): String = Instant.now().toString()
|
||||
}
|
||||
|
||||
/** Metrics for train.udp_updown; recomputable from the columnar evidence. */
|
||||
@Serializable
|
||||
data class EchoMetrics(
|
||||
val sent: Int,
|
||||
val received: Int,
|
||||
@SerialName("loss_pct") val lossPct: Double,
|
||||
@SerialName("rtt_ms_min") val rttMsMin: Double? = null,
|
||||
@SerialName("rtt_ms_avg") val rttMsAvg: Double? = null,
|
||||
@SerialName("rtt_ms_max") val rttMsMax: Double? = null,
|
||||
@SerialName("observed_ports") val observedPorts: List<Int> = emptyList(),
|
||||
@SerialName("nat_rebinding_detected") val natRebindingDetected: Boolean = false,
|
||||
)
|
||||
@@ -0,0 +1,63 @@
|
||||
// SPDX-FileCopyrightText: 2026 Echolot contributors
|
||||
// SPDX-License-Identifier: GPL-3.0-or-later
|
||||
|
||||
package app.echo_lot.engine
|
||||
|
||||
import app.echo_lot.measurement.*
|
||||
import kotlinx.serialization.json.Json
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertTrue
|
||||
|
||||
/**
|
||||
* Runs the full server-facing engine against a live server and validates the produced
|
||||
* MeasurementDocument. Self-skips without ECHOLOT_LIVE_* (same contract as core-protocol's live
|
||||
* test). This is the whole vertical: protocol client → engine → schema document → verdict.
|
||||
*/
|
||||
class LiveMeasurementTest {
|
||||
|
||||
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 producesValidDocumentFromLiveServer() {
|
||||
if (url == null || pin == null || cred == null || udp == null) {
|
||||
println("LiveMeasurementTest skipped (no ECHOLOT_LIVE_* env)")
|
||||
return
|
||||
}
|
||||
val (host, port) = udp.split(":").let { it[0] to it[1].toInt() }
|
||||
val engine = ServerMeasurement(
|
||||
ids = SystemIdSource(),
|
||||
app = AppInfo(version = "0.1.0", build = 1, flavor = "test"),
|
||||
device = DeviceInfo("test", "jvm", 0, "n/a"),
|
||||
)
|
||||
val doc = engine.run(
|
||||
ServerMeasurement.Config(
|
||||
controlUrl = url, pins = setOf(pin), credential = cred,
|
||||
target = target, udpHost = host, udpPort = port, echoCount = 20,
|
||||
)
|
||||
)
|
||||
|
||||
// The document must round-trip and carry the expected structure.
|
||||
val encoded = Json { encodeDefaults = true }.encodeToString(MeasurementDocument.serializer(), doc)
|
||||
println("document (${encoded.length} bytes): overall=${doc.summary?.overall}")
|
||||
|
||||
assertEquals(1, doc.serverSessions.size)
|
||||
assertTrue(doc.serverSessions[0].capabilities.contains("udp-probe"))
|
||||
val test = doc.tests.single()
|
||||
assertEquals(TestType.TRAIN_UDP_UPDOWN, test.type)
|
||||
assertTrue(test.status == TestStatus.OK || test.status == TestStatus.PARTIAL,
|
||||
"expected replies from live server, got ${test.status}")
|
||||
|
||||
val metrics = Json.parseToJsonElement(test.metrics.toString())
|
||||
println("metrics: $metrics")
|
||||
assertTrue(metrics.toString().contains("rtt_ms_avg"))
|
||||
|
||||
assertTrue(doc.summary != null)
|
||||
// A healthy local->fmr path should be green (no loss, no rebinding) or yellow.
|
||||
println("summary: ${doc.summary}")
|
||||
}
|
||||
}
|
||||
@@ -23,3 +23,4 @@ rootProject.name = "echolot-app"
|
||||
// (core-probe, core-shizuku, app) join as they land.
|
||||
include(":core-protocol")
|
||||
include(":core-measurement")
|
||||
include(":core-engine")
|
||||
|
||||
Reference in New Issue
Block a user