diag(nests-interop): per-test moq-relay trace capture for routing race
`NativeMoqRelayHarness` gains a `testTag` parameter on `shared` /
`resetShared` and writes the relay subprocess's combined stdout/stderr
to `nestsClient/build/relay-logs/<tag>-<seq>-<ts>.log` whenever
`-DnestsHangInteropTraceRelay=true` is set. The subprocess runs with
`RUST_LOG=info,moq_relay=trace,moq_lite=trace,moq_native=debug` so the
per-broadcast subscribe-routing path investigated in
`nestsClient/plans/2026-05-07-moq-relay-routing-investigation.md`
(Step 1) becomes visible.
`HangInteropTest` and `BrowserInteropTest` now expose a JUnit 4
`TestName` rule and forward the method name into `resetShared`, so
each scenario's per-method log is easy to locate by name when
cross-referencing with the speaker-side `Log.d("NestTx")` lines
already captured in JUnit XML's `<system-err>`.
This commit is contained in:
@@ -227,6 +227,18 @@ tasks.withType<Test>().configureEach {
|
|||||||
val cargoBin = hangInteropCacheDir.dir("bin").asFile
|
val cargoBin = hangInteropCacheDir.dir("bin").asFile
|
||||||
systemProperty("nestsHangInteropSidecarsDir", sidecarRelease.absolutePath)
|
systemProperty("nestsHangInteropSidecarsDir", sidecarRelease.absolutePath)
|
||||||
systemProperty("nestsHangInteropCargoBinDir", cargoBin.absolutePath)
|
systemProperty("nestsHangInteropCargoBinDir", cargoBin.absolutePath)
|
||||||
|
// Per-method moq-relay trace log dir for the routing-race
|
||||||
|
// investigation (plan 2026-05-07-moq-relay-routing-investigation.md).
|
||||||
|
// Off by default; opt in via -DnestsHangInteropTraceRelay=true so a
|
||||||
|
// routine sweep doesn't generate ~MBs of trace per run.
|
||||||
|
if (System.getProperty("nestsHangInteropTraceRelay") == "true") {
|
||||||
|
val relayLogDir =
|
||||||
|
layout.buildDirectory
|
||||||
|
.dir("relay-logs")
|
||||||
|
.get()
|
||||||
|
.asFile
|
||||||
|
systemProperty("nestsHangInteropRelayLogDir", relayLogDir.absolutePath)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// ---- Cross-stack interop: BROWSER (Phase 4 of T16) --------------------------
|
// ---- Cross-stack interop: BROWSER (Phase 4 of T16) --------------------------
|
||||||
|
|||||||
@@ -1,6 +1,29 @@
|
|||||||
# Plan: investigate moq-relay 0.10.x per-broadcast subscribe-routing race
|
# Plan: investigate moq-relay 0.10.x per-broadcast subscribe-routing race
|
||||||
|
|
||||||
**Status:** specced — pickup ready.
|
**Status:** Step 1 instrumentation landed; sweep + analysis in progress.
|
||||||
|
|
||||||
|
## Progress log (2026-05-07)
|
||||||
|
|
||||||
|
- Step 1 instrumentation landed. `NativeMoqRelayHarness` now accepts an
|
||||||
|
optional `testTag` and writes the relay subprocess's combined
|
||||||
|
stdout/stderr to
|
||||||
|
`nestsClient/build/relay-logs/<methodName>-<seq>-<ts>.log` whenever
|
||||||
|
`-DnestsHangInteropTraceRelay=true` is set. The subprocess runs with
|
||||||
|
`RUST_LOG=info,moq_relay=trace,moq_lite=trace,moq_native=debug` so
|
||||||
|
the per-broadcast subscribe-routing path is observable; quinn /
|
||||||
|
rustls / h3 stay at info to keep the file < ~10 MB per scenario.
|
||||||
|
`HangInteropTest` and `BrowserInteropTest` now expose a JUnit 4
|
||||||
|
`TestName` rule and pass the method name into `resetShared` so each
|
||||||
|
scenario's per-method log is easy to locate by name.
|
||||||
|
- Step 4 ruled out — `cargo info moq-relay` confirms `0.10.25` is the
|
||||||
|
current `crates.io` release; no newer minor exists. The plan's
|
||||||
|
Step 4 path (next minor on crates.io) does not apply. The fallback
|
||||||
|
is a `cargo install --git https://github.com/moq-dev/moq.git --rev
|
||||||
|
<main-head>` against `bdda6bd19a37ccdf7f7b66f3d760d8892ea8db59`
|
||||||
|
(main HEAD as of investigation) — moq-relay-v0.10.25 tag matches
|
||||||
|
the published crate, so post-0.10.25 work lives only on `main`.
|
||||||
|
Also note the upstream moved from `kixelated/moq` to `moq-dev/moq`;
|
||||||
|
REV's `KIXELATED_MOQ_GIT_REV` predates the move.
|
||||||
|
|
||||||
**Owns:** the residual flake that affects four T16 scenarios:
|
**Owns:** the residual flake that affects four T16 scenarios:
|
||||||
`late_join_listener_still_decodes_tail`,
|
`late_join_listener_still_decodes_tail`,
|
||||||
|
|||||||
+10
-1
@@ -52,6 +52,8 @@ import java.util.UUID
|
|||||||
import kotlin.test.BeforeTest
|
import kotlin.test.BeforeTest
|
||||||
import kotlin.test.Test
|
import kotlin.test.Test
|
||||||
import kotlin.test.assertTrue
|
import kotlin.test.assertTrue
|
||||||
|
import org.junit.Rule
|
||||||
|
import org.junit.rules.TestName
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Phase 4 (T16) — browser-side cross-stack interop scenarios.
|
* Phase 4 (T16) — browser-side cross-stack interop scenarios.
|
||||||
@@ -78,6 +80,13 @@ import kotlin.test.assertTrue
|
|||||||
* `NativeMoqRelayHarness`).
|
* `NativeMoqRelayHarness`).
|
||||||
*/
|
*/
|
||||||
class BrowserInteropTest {
|
class BrowserInteropTest {
|
||||||
|
/**
|
||||||
|
* Tags the per-method moq-relay log file when trace capture is
|
||||||
|
* enabled. See `HangInteropTest.testName`.
|
||||||
|
*/
|
||||||
|
@Rule @JvmField
|
||||||
|
val testName: TestName = TestName()
|
||||||
|
|
||||||
@BeforeTest
|
@BeforeTest
|
||||||
fun gate() {
|
fun gate() {
|
||||||
PlaywrightDriver.assumeBrowserInterop()
|
PlaywrightDriver.assumeBrowserInterop()
|
||||||
@@ -97,7 +106,7 @@ class BrowserInteropTest {
|
|||||||
// browser-publisher tests run alongside browser-listener
|
// browser-publisher tests run alongside browser-listener
|
||||||
// tests). Per-method reboot costs ~500 ms (cargo binaries are
|
// tests). Per-method reboot costs ~500 ms (cargo binaries are
|
||||||
// cached); acceptable for the stability gain.
|
// cached); acceptable for the stability gain.
|
||||||
NativeMoqRelayHarness.resetShared()
|
NativeMoqRelayHarness.resetShared(testTag = testName.methodName)
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
+12
-1
@@ -58,6 +58,8 @@ import kotlin.test.Test
|
|||||||
import kotlin.test.assertEquals
|
import kotlin.test.assertEquals
|
||||||
import kotlin.test.assertFalse
|
import kotlin.test.assertFalse
|
||||||
import kotlin.test.assertTrue
|
import kotlin.test.assertTrue
|
||||||
|
import org.junit.Rule
|
||||||
|
import org.junit.rules.TestName
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Cross-stack interop scenarios driving the reference `kixelated/moq`
|
* Cross-stack interop scenarios driving the reference `kixelated/moq`
|
||||||
@@ -86,6 +88,15 @@ import kotlin.test.assertTrue
|
|||||||
* Gated by `-DnestsHangInterop=true`.
|
* Gated by `-DnestsHangInterop=true`.
|
||||||
*/
|
*/
|
||||||
class HangInteropTest {
|
class HangInteropTest {
|
||||||
|
/**
|
||||||
|
* JUnit 4 rule that exposes the running test method's name. Used
|
||||||
|
* to tag the per-method moq-relay log file when trace capture is
|
||||||
|
* enabled (`-DnestsHangInteropTraceRelay=true`) — see
|
||||||
|
* `NativeMoqRelayHarness.RELAY_LOG_DIR_PROPERTY`.
|
||||||
|
*/
|
||||||
|
@Rule @JvmField
|
||||||
|
val testName: TestName = TestName()
|
||||||
|
|
||||||
@BeforeTest
|
@BeforeTest
|
||||||
fun gate() {
|
fun gate() {
|
||||||
NativeMoqRelayHarness.assumeHangInterop()
|
NativeMoqRelayHarness.assumeHangInterop()
|
||||||
@@ -97,7 +108,7 @@ class HangInteropTest {
|
|||||||
// that don't reproduce in isolation. Per-method reboot
|
// that don't reproduce in isolation. Per-method reboot
|
||||||
// costs ~500 ms (cargo binaries are cached) — acceptable
|
// costs ~500 ms (cargo binaries are cached) — acceptable
|
||||||
// for the stability gain.
|
// for the stability gain.
|
||||||
NativeMoqRelayHarness.resetShared()
|
NativeMoqRelayHarness.resetShared(testTag = testName.methodName)
|
||||||
}
|
}
|
||||||
|
|
||||||
/** I1: 5 s 440 Hz mono sine, asserted via FFT peak + ZCR. */
|
/** I1: 5 s 440 Hz mono sine, asserted via FFT peak + ZCR. */
|
||||||
|
|||||||
+119
-10
@@ -26,8 +26,11 @@ import java.net.InetSocketAddress
|
|||||||
import java.net.ServerSocket
|
import java.net.ServerSocket
|
||||||
import java.nio.file.Files
|
import java.nio.file.Files
|
||||||
import java.nio.file.Path
|
import java.nio.file.Path
|
||||||
|
import java.time.LocalDateTime
|
||||||
|
import java.time.format.DateTimeFormatter
|
||||||
import java.util.concurrent.ConcurrentLinkedQueue
|
import java.util.concurrent.ConcurrentLinkedQueue
|
||||||
import java.util.concurrent.TimeUnit
|
import java.util.concurrent.TimeUnit
|
||||||
|
import java.util.concurrent.atomic.AtomicInteger
|
||||||
import java.util.concurrent.locks.ReentrantLock
|
import java.util.concurrent.locks.ReentrantLock
|
||||||
import kotlin.concurrent.withLock
|
import kotlin.concurrent.withLock
|
||||||
|
|
||||||
@@ -59,6 +62,13 @@ class NativeMoqRelayHarness private constructor(
|
|||||||
private val relayPort: Int,
|
private val relayPort: Int,
|
||||||
private val sidecarsDir: Path,
|
private val sidecarsDir: Path,
|
||||||
private val cargoBinDir: Path,
|
private val cargoBinDir: Path,
|
||||||
|
/**
|
||||||
|
* File the relay's combined stdout/stderr was tee'd to for this
|
||||||
|
* boot, when trace-log capture was enabled. Useful for tests that
|
||||||
|
* want to attach the per-method relay log to a failure assertion.
|
||||||
|
* `null` when the per-method log dir wasn't configured.
|
||||||
|
*/
|
||||||
|
val relayLogFile: Path?,
|
||||||
) : AutoCloseable {
|
) : AutoCloseable {
|
||||||
private var stopped = false
|
private var stopped = false
|
||||||
|
|
||||||
@@ -107,9 +117,34 @@ class NativeMoqRelayHarness private constructor(
|
|||||||
*/
|
*/
|
||||||
const val CARGO_BIN_DIR_PROPERTY = "nestsHangInteropCargoBinDir"
|
const val CARGO_BIN_DIR_PROPERTY = "nestsHangInteropCargoBinDir"
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Optional dir where each relay subprocess boot writes its
|
||||||
|
* combined stdout/stderr to a file. Forwarded by
|
||||||
|
* `:nestsClient`'s test task to
|
||||||
|
* `nestsClient/build/relay-logs/`. When set, the relay also
|
||||||
|
* runs with `RUST_LOG=moq_relay=trace,moq_lite=trace` so the
|
||||||
|
* captured file contains the per-broadcast subscribe-routing
|
||||||
|
* trace investigated in
|
||||||
|
* `nestsClient/plans/2026-05-07-moq-relay-routing-investigation.md`.
|
||||||
|
*
|
||||||
|
* One file per relay boot; `resetShared(testTag = "<name>")`
|
||||||
|
* tags filenames so a sweep produces
|
||||||
|
* `<name>-<seq>-<timestamp>.log` and a failed run is easy to
|
||||||
|
* locate by test method name. Without the property, the relay
|
||||||
|
* runs with `--log-level info` and no per-boot file is
|
||||||
|
* produced — keeps the harness's existing behaviour for
|
||||||
|
* non-investigatory runs.
|
||||||
|
*/
|
||||||
|
const val RELAY_LOG_DIR_PROPERTY = "nestsHangInteropRelayLogDir"
|
||||||
|
|
||||||
private const val PORT_READY_TIMEOUT_MS = 30_000L
|
private const val PORT_READY_TIMEOUT_MS = 30_000L
|
||||||
private const val PORT_PROBE_INTERVAL_MS = 200L
|
private const val PORT_PROBE_INTERVAL_MS = 200L
|
||||||
|
|
||||||
|
/** Monotonic counter used to disambiguate same-tag boots in one JVM. */
|
||||||
|
private val bootSequence = AtomicInteger(0)
|
||||||
|
private val LOG_TIMESTAMP_FMT =
|
||||||
|
DateTimeFormatter.ofPattern("yyyyMMdd-HHmmss-SSS")
|
||||||
|
|
||||||
fun isEnabled(): Boolean = System.getProperty(ENABLE_PROPERTY) == "true"
|
fun isEnabled(): Boolean = System.getProperty(ENABLE_PROPERTY) == "true"
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -139,12 +174,18 @@ class NativeMoqRelayHarness private constructor(
|
|||||||
* Bring the relay up if not already running; reuses the same
|
* Bring the relay up if not already running; reuses the same
|
||||||
* subprocess across test classes within one JVM run. Mirrors
|
* subprocess across test classes within one JVM run. Mirrors
|
||||||
* the singleton pattern in [com.vitorpamplona.nestsclient.interop.NostrNestsHarness].
|
* the singleton pattern in [com.vitorpamplona.nestsclient.interop.NostrNestsHarness].
|
||||||
|
*
|
||||||
|
* If a relay log dir is configured (see [RELAY_LOG_DIR_PROPERTY])
|
||||||
|
* and no shared relay exists yet, the boot is tagged with
|
||||||
|
* [testTag]; otherwise the existing relay is reused regardless
|
||||||
|
* of tag. Use [resetShared] to force a fresh boot with a
|
||||||
|
* specific tag.
|
||||||
*/
|
*/
|
||||||
fun shared(): NativeMoqRelayHarness {
|
fun shared(testTag: String? = null): NativeMoqRelayHarness {
|
||||||
shared?.let { return it }
|
shared?.let { return it }
|
||||||
synchronized(sharedLock) {
|
synchronized(sharedLock) {
|
||||||
shared?.let { return it }
|
shared?.let { return it }
|
||||||
val instance = doStart()
|
val instance = doStart(testTag)
|
||||||
Runtime.getRuntime().addShutdownHook(
|
Runtime.getRuntime().addShutdownHook(
|
||||||
Thread({ runCatching { instance.close() } }, "NativeMoqRelayHarness-shutdown"),
|
Thread({ runCatching { instance.close() } }, "NativeMoqRelayHarness-shutdown"),
|
||||||
)
|
)
|
||||||
@@ -169,16 +210,27 @@ class NativeMoqRelayHarness private constructor(
|
|||||||
* are paid). At 11 scenarios × 500 ms that's ~5.5 s added
|
* are paid). At 11 scenarios × 500 ms that's ~5.5 s added
|
||||||
* to the suite wallclock — acceptable trade for stability.
|
* to the suite wallclock — acceptable trade for stability.
|
||||||
*/
|
*/
|
||||||
fun resetShared() {
|
fun resetShared(testTag: String? = null) {
|
||||||
synchronized(sharedLock) {
|
synchronized(sharedLock) {
|
||||||
shared?.let {
|
shared?.let {
|
||||||
runCatching { it.close() }
|
runCatching { it.close() }
|
||||||
}
|
}
|
||||||
shared = null
|
shared = null
|
||||||
}
|
}
|
||||||
|
// Pre-warm so the next caller observes the relay already up
|
||||||
|
// tagged with this test method's name. Without this, the
|
||||||
|
// first call to shared() after resetShared() picks up the
|
||||||
|
// tag of whoever wins the race — usually the test body
|
||||||
|
// calling `shared()`, which is fine, but a concurrent
|
||||||
|
// listener-side helper may race in first under suite
|
||||||
|
// mode. Pre-warming is cheap (~500 ms cargo cache hit)
|
||||||
|
// and keeps the per-method log filename stable.
|
||||||
|
if (testTag != null && System.getProperty(RELAY_LOG_DIR_PROPERTY) != null) {
|
||||||
|
shared(testTag)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private fun doStart(): NativeMoqRelayHarness {
|
private fun doStart(testTag: String?): NativeMoqRelayHarness {
|
||||||
check(isEnabled()) {
|
check(isEnabled()) {
|
||||||
"NativeMoqRelayHarness.shared() called without -D$ENABLE_PROPERTY=true."
|
"NativeMoqRelayHarness.shared() called without -D$ENABLE_PROPERTY=true."
|
||||||
}
|
}
|
||||||
@@ -196,6 +248,25 @@ class NativeMoqRelayHarness private constructor(
|
|||||||
}
|
}
|
||||||
|
|
||||||
val port = reservePort()
|
val port = reservePort()
|
||||||
|
val relayLogDir = System.getProperty(RELAY_LOG_DIR_PROPERTY)?.let { File(it) }
|
||||||
|
val relayLogFile: File? =
|
||||||
|
if (relayLogDir != null) {
|
||||||
|
relayLogDir.mkdirs()
|
||||||
|
val seq = bootSequence.incrementAndGet().toString().padStart(3, '0')
|
||||||
|
val ts = LocalDateTime.now().format(LOG_TIMESTAMP_FMT)
|
||||||
|
val tag = sanitiseTag(testTag ?: "boot")
|
||||||
|
File(relayLogDir, "$tag-$seq-$ts.log")
|
||||||
|
} else {
|
||||||
|
null
|
||||||
|
}
|
||||||
|
// Keep `--log-level info` as the baseline; the relay's
|
||||||
|
// tracing_subscriber EnvFilter honours `RUST_LOG`, which
|
||||||
|
// we set to trace on `moq_relay` + `moq_lite` only when
|
||||||
|
// capture is enabled. That gives us the per-broadcast
|
||||||
|
// subscribe-routing trace investigated in plan
|
||||||
|
// `2026-05-07-moq-relay-routing-investigation.md` while
|
||||||
|
// keeping quinn/h3/tls noise at info — full-tree trace
|
||||||
|
// is ~100s of MB per test, way more than needed.
|
||||||
val pb =
|
val pb =
|
||||||
ProcessBuilder(
|
ProcessBuilder(
|
||||||
moqRelay.toString(),
|
moqRelay.toString(),
|
||||||
@@ -219,9 +290,14 @@ class NativeMoqRelayHarness private constructor(
|
|||||||
"--log-level",
|
"--log-level",
|
||||||
"info",
|
"info",
|
||||||
).redirectErrorStream(true)
|
).redirectErrorStream(true)
|
||||||
|
if (relayLogFile != null) {
|
||||||
|
pb.environment()["RUST_LOG"] =
|
||||||
|
"info,moq_relay=trace,moq_lite=trace,moq_native=debug"
|
||||||
|
}
|
||||||
|
|
||||||
val process = pb.start()
|
val process = pb.start()
|
||||||
val drainer = ProcessOutputDrainer(process, "moq-relay").also { it.start() }
|
val drainer =
|
||||||
|
ProcessOutputDrainer(process, "moq-relay", relayLogFile).also { it.start() }
|
||||||
|
|
||||||
try {
|
try {
|
||||||
// moq-relay logs `addr=… listening` on bind. Wait for
|
// moq-relay logs `addr=… listening` on bind. Wait for
|
||||||
@@ -251,9 +327,20 @@ class NativeMoqRelayHarness private constructor(
|
|||||||
relayPort = port,
|
relayPort = port,
|
||||||
sidecarsDir = sidecarsDir,
|
sidecarsDir = sidecarsDir,
|
||||||
cargoBinDir = cargoBinDir,
|
cargoBinDir = cargoBinDir,
|
||||||
|
relayLogFile = relayLogFile?.toPath(),
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Strip filesystem-unfriendly characters from a JUnit test
|
||||||
|
* method name so it can be used directly in a log filename.
|
||||||
|
*/
|
||||||
|
private fun sanitiseTag(raw: String): String =
|
||||||
|
raw
|
||||||
|
.replace(Regex("[^A-Za-z0-9._-]"), "_")
|
||||||
|
.take(80)
|
||||||
|
.ifBlank { "boot" }
|
||||||
|
|
||||||
private fun requireDirProperty(name: String): Path {
|
private fun requireDirProperty(name: String): Path {
|
||||||
val raw = System.getProperty(name)
|
val raw = System.getProperty(name)
|
||||||
check(!raw.isNullOrBlank()) {
|
check(!raw.isNullOrBlank()) {
|
||||||
@@ -350,6 +437,16 @@ class NativeMoqRelayHarness private constructor(
|
|||||||
private class ProcessOutputDrainer(
|
private class ProcessOutputDrainer(
|
||||||
private val process: Process,
|
private val process: Process,
|
||||||
private val name: String,
|
private val name: String,
|
||||||
|
/**
|
||||||
|
* Optional sink for the full subprocess output. When non-null,
|
||||||
|
* every line is also written here verbatim — used by the
|
||||||
|
* routing-race investigation (see plan
|
||||||
|
* `2026-05-07-moq-relay-routing-investigation.md`) to keep the
|
||||||
|
* trace-level log around for post-hoc analysis. The in-memory
|
||||||
|
* ring is kept regardless so `tail()` still feeds failure
|
||||||
|
* messages.
|
||||||
|
*/
|
||||||
|
private val sinkFile: File? = null,
|
||||||
) {
|
) {
|
||||||
private val ring = ConcurrentLinkedQueue<String>()
|
private val ring = ConcurrentLinkedQueue<String>()
|
||||||
private val maxLines = 64
|
private val maxLines = 64
|
||||||
@@ -360,12 +457,24 @@ private class ProcessOutputDrainer(
|
|||||||
fun start() {
|
fun start() {
|
||||||
thread =
|
thread =
|
||||||
Thread({
|
Thread({
|
||||||
process.inputStream.bufferedReader().useLines { lines ->
|
val writer = sinkFile?.bufferedWriter()
|
||||||
for (line in lines) {
|
try {
|
||||||
ring.add(line)
|
process.inputStream.bufferedReader().useLines { lines ->
|
||||||
while (ring.size > maxLines) ring.poll()
|
for (line in lines) {
|
||||||
lock.withLock { newLineCond.signalAll() }
|
ring.add(line)
|
||||||
|
while (ring.size > maxLines) ring.poll()
|
||||||
|
if (writer != null) {
|
||||||
|
runCatching {
|
||||||
|
writer.write(line)
|
||||||
|
writer.newLine()
|
||||||
|
writer.flush()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
lock.withLock { newLineCond.signalAll() }
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
} finally {
|
||||||
|
runCatching { writer?.close() }
|
||||||
}
|
}
|
||||||
}, "NativeMoqRelayHarness-$name").apply {
|
}, "NativeMoqRelayHarness-$name").apply {
|
||||||
isDaemon = true
|
isDaemon = true
|
||||||
|
|||||||
Reference in New Issue
Block a user