From 2da18859755673a759486e75666d5513d8a3ff55 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 1 May 2026 02:42:41 +0000 Subject: [PATCH] Add diagnostic sweep matrix across cadence, frame count, payload, group packing, fan-out MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The previous run pinned the production failure to "82/100 frames at 20 ms cadence two-users, all 18 trailing groups lost; publisher.send returned true for every frame". To narrow the cause the next run needs to vary one dimension at a time and dump rich per-frame data, so we can attribute the loss to: - QUIC stream-credit (RTT-bound) vs absolute count cap - Per-stream-creation cost vs per-byte cost - Listener-side buffer overflow vs relay-side drop - Steady-state vs initial-burst behaviour - Shared end-of-stream FIN flush race vs prod-only loss SendTraceScenario.kt: shared per-frame instrumented pump driver. Records publisher.send Boolean per frame, send latency, endGroup exceptions, arrivals per subscriber with wall-clock timing, missing-objectId run lengths, send-latency p50/p99/max + count of >1ms sends. Verbose logging (InteropDebug.checkpoint) emits one line per tx and one per rx so the JUnit XML system-out doubles as a timing trace. Sweep methods (mirrored prod-mode and harness-mode): - baseline 100×20ms (reproduce known cliff) - cadence sweep: 5 / 10 / 40 / 80 / 200 ms (RTT-dependent?) - burst (no cadence) 100 frames (max inflight) - frame-count sweep: 50 / 200 / 400 - payload sweep: 1 KB / 4 KB / 16 KB - frames-per-group: 5 / 20 / 100 (uni-stream vs bytes) - late subscribe at frame 25, frame 50 (from-latest semantics) - mid-pump 5 s pause after frame 50 - 2-, 3-subscriber fan-out - 50 ms slow-consumer listener (per-subscription buffer overflow) - long-run 30 s and 120 s (steady-state) Each test method is a one-line Scenario literal; setup helpers withProdSpeakerAndListeners / withHarnessSpeakerAndListeners build publisher + N listeners, run the scenario, and tear down in order. Smoke-tested against the local kixelated/moq reference relay: every scenario except framesPerGroup=100 reproduces the same trailing-frame artifact (received=99/100, missing=[99]) the original test surfaced. framesPerGroup=100 receives every frame, confirming the local issue is specifically per-uni-stream FIN flushing — independent of the production 82/100 cliff. --- ...trNestsSustainedSendOutcomesInteropTest.kt | 250 ++++++++++ .../NostrnestsProdAudioTransmissionTest.kt | 287 +++++++++++ .../nestsclient/interop/SendTraceScenario.kt | 463 ++++++++++++++++++ 3 files changed, 1000 insertions(+) create mode 100644 nestsClient/src/jvmTest/kotlin/com/vitorpamplona/nestsclient/interop/SendTraceScenario.kt diff --git a/nestsClient/src/jvmTest/kotlin/com/vitorpamplona/nestsclient/interop/NostrNestsSustainedSendOutcomesInteropTest.kt b/nestsClient/src/jvmTest/kotlin/com/vitorpamplona/nestsclient/interop/NostrNestsSustainedSendOutcomesInteropTest.kt index 9c1ce9d34..22266fcf9 100644 --- a/nestsClient/src/jvmTest/kotlin/com/vitorpamplona/nestsclient/interop/NostrNestsSustainedSendOutcomesInteropTest.kt +++ b/nestsClient/src/jvmTest/kotlin/com/vitorpamplona/nestsclient/interop/NostrNestsSustainedSendOutcomesInteropTest.kt @@ -229,6 +229,256 @@ class NostrNestsSustainedSendOutcomesInteropTest { Unit } + // ==================================================================== + // Same sweep matrix as [NostrnestsProdAudioTransmissionTest], but + // wired to the local Docker harness so we can compare prod vs. + // reference-relay behaviour for each scenario in one run. The + // [SendTraceScenario] runner is shared. + // ==================================================================== + + @Test + fun harness_sweep_baseline_100x20ms() = runHarnessScenarioOrSkip("h-baseline", Scenario(frameCount = 100, cadenceMs = 20L)) + + @Test fun harness_sweep_cadence_5ms() = runHarnessScenarioOrSkip("h-cad5", Scenario(frameCount = 100, cadenceMs = 5L)) + + @Test fun harness_sweep_cadence_10ms() = runHarnessScenarioOrSkip("h-cad10", Scenario(frameCount = 100, cadenceMs = 10L)) + + @Test fun harness_sweep_cadence_40ms() = runHarnessScenarioOrSkip("h-cad40", Scenario(frameCount = 100, cadenceMs = 40L)) + + @Test fun harness_sweep_cadence_80ms() = runHarnessScenarioOrSkip("h-cad80", Scenario(frameCount = 100, cadenceMs = 80L)) + + @Test fun harness_sweep_cadence_200ms() = runHarnessScenarioOrSkip("h-cad200", Scenario(frameCount = 100, cadenceMs = 200L)) + + @Test + fun harness_sweep_burst_no_cadence() = + runHarnessScenarioOrSkip( + "h-burst", + Scenario(frameCount = 100, cadenceMs = 0L, receiveGraceMs = 30_000L), + ) + + @Test fun harness_sweep_frames_50() = runHarnessScenarioOrSkip("h-frames50", Scenario(frameCount = 50, cadenceMs = 20L)) + + @Test + fun harness_sweep_frames_200() = + runHarnessScenarioOrSkip( + "h-frames200", + Scenario(frameCount = 200, cadenceMs = 20L, receiveGraceMs = 45_000L), + ) + + @Test + fun harness_sweep_frames_400() = + runHarnessScenarioOrSkip( + "h-frames400", + Scenario(frameCount = 400, cadenceMs = 20L, receiveGraceMs = 60_000L), + ) + + @Test + fun harness_sweep_payload_1kb() = + runHarnessScenarioOrSkip( + "h-1kb", + Scenario(frameCount = 100, cadenceMs = 20L, payloadBytes = 1024), + ) + + @Test + fun harness_sweep_payload_4kb() = + runHarnessScenarioOrSkip( + "h-4kb", + Scenario(frameCount = 100, cadenceMs = 20L, payloadBytes = 4096), + ) + + @Test + fun harness_sweep_payload_16kb() = + runHarnessScenarioOrSkip( + "h-16kb", + Scenario(frameCount = 100, cadenceMs = 20L, payloadBytes = 16384), + ) + + @Test + fun harness_sweep_frames_per_group_5() = + runHarnessScenarioOrSkip( + "h-fpg5", + Scenario(frameCount = 100, cadenceMs = 20L, framesPerGroup = 5), + ) + + @Test + fun harness_sweep_frames_per_group_20() = + runHarnessScenarioOrSkip( + "h-fpg20", + Scenario(frameCount = 100, cadenceMs = 20L, framesPerGroup = 20), + ) + + @Test + fun harness_sweep_frames_per_group_all() = + runHarnessScenarioOrSkip( + "h-fpg-all", + Scenario(frameCount = 100, cadenceMs = 20L, framesPerGroup = 100), + ) + + @Test + fun harness_sweep_late_subscribe_after_25() = + runHarnessScenarioOrSkip( + "h-late25", + Scenario(frameCount = 100, cadenceMs = 20L, subscribeAtFrame = 25), + ) + + @Test + fun harness_sweep_late_subscribe_after_50() = + runHarnessScenarioOrSkip( + "h-late50", + Scenario(frameCount = 100, cadenceMs = 20L, subscribeAtFrame = 50), + ) + + @Test + fun harness_sweep_mid_pause_50_5s() = + runHarnessScenarioOrSkip( + "h-pause", + Scenario( + frameCount = 100, + cadenceMs = 20L, + pauseAfterFrame = 50, + pauseDurationMs = 5_000L, + receiveGraceMs = 45_000L, + ), + ) + + @Test + fun harness_sweep_two_subscribers() = + runHarnessScenarioOrSkip( + "h-2subs", + Scenario(frameCount = 100, cadenceMs = 20L, parallelSubscriptions = 2), + ) + + @Test + fun harness_sweep_three_subscribers() = + runHarnessScenarioOrSkip( + "h-3subs", + Scenario(frameCount = 100, cadenceMs = 20L, parallelSubscriptions = 3), + ) + + @Test + fun harness_sweep_slow_consumer_50ms() = + runHarnessScenarioOrSkip( + "h-slowconsumer", + Scenario( + frameCount = 100, + cadenceMs = 20L, + listenerSlowConsumerMs = 50L, + receiveGraceMs = 60_000L, + ), + ) + + @Test + fun harness_sweep_long_run_30s() = + runHarnessScenarioOrSkip( + "h-30s", + Scenario( + frameCount = 1500, + cadenceMs = 20L, + receiveGraceMs = 60_000L, + verbosePerFrame = false, + ), + ) + + private fun runHarnessScenarioOrSkip( + scope: String, + scenario: Scenario, + expectAllReceived: Boolean = false, + ) = runBlocking { + NostrNestsHarness.assumeNestsInterop() + val harness = harnessOrNull ?: return@runBlocking + withHarnessSpeakerAndListeners(scope, harness, scenario.parallelSubscriptions) { publisher, listeners, hostPub, pumpScope -> + val result = + SendTraceScenario.run( + scope = scope, + publisher = publisher, + listeners = listeners, + speakerPubkeyHex = hostPub, + scenario = scenario, + pumpScope = pumpScope, + ) + SendTraceScenario.reportAndAssert(scope, result, expectAllReceived) + } + Unit + } + + private suspend fun withHarnessSpeakerAndListeners( + scope: String, + harness: NostrNestsHarness, + listenerCount: Int, + block: suspend ( + publisher: com.vitorpamplona.nestsclient.moq.lite.MoqLitePublisherHandle, + listeners: List, + speakerPubkeyHex: String, + pumpScope: CoroutineScope, + ) -> Unit, + ) { + val hostSigner = NostrSignerInternal(KeyPair()) + val room = + com.vitorpamplona.nestsclient + .NestsRoomConfig( + authBaseUrl = harness.authBaseUrl, + endpoint = harness.moqEndpoint, + hostPubkey = hostSigner.pubKey, + roomId = "h-sweep-${System.currentTimeMillis()}-${(0..9999).random()}", + ) + + val httpClient = OkHttpNestsClient { OkHttpClient() } + val transport = + QuicWebTransportFactory(certificateValidator = PermissiveCertificateValidator()) + + val supervisor = SupervisorJob() + val pumpScope = CoroutineScope(supervisor + Dispatchers.IO) + + InteropDebug.checkpoint(scope, "auth=${room.authBaseUrl} endpoint=${room.endpoint} ns=${room.moqNamespace()}") + + val publishToken = + InteropDebug.stepSuspending(scope, "host: mintToken(publish=true)") { + httpClient.mintToken(room = room, publish = true, signer = hostSigner) + } + val (authority, path) = + buildRelayConnectTarget( + endpoint = room.endpoint, + namespace = room.moqNamespace(), + token = publishToken, + ) + val speakerWt = + InteropDebug.stepSuspending(scope, "host: WebTransport.connect") { + transport.connect(authority = authority, path = path, bearerToken = null) + } + val speakerSession = MoqLiteSession.client(speakerWt, pumpScope) + val publisher = + InteropDebug.stepSuspending(scope, "host: session.publish(broadcastSuffix=hostPub)") { + speakerSession.publish(broadcastSuffix = hostSigner.pubKey) + } + + val listeners = mutableListOf() + try { + for (i in 0 until listenerCount) { + val audienceSigner = NostrSignerInternal(KeyPair()) + val listener = + InteropDebug.stepSuspending(scope, "audience[$i]: connectNestsListener") { + connectNestsListener( + httpClient = httpClient, + transport = transport, + scope = pumpScope, + room = room, + signer = audienceSigner, + ) + } + InteropDebug.assertListenerReached(scope, "Connected", listener.state.value) + listeners += listener + } + block(publisher, listeners, hostSigner.pubKey, pumpScope) + } finally { + for (listener in listeners) { + runCatching { listener.close() } + } + runCatching { publisher.close() } + runCatching { speakerWt.close(0, "test done") } + supervisor.cancelAndJoin() + } + } + companion object { private const val N_FRAMES = 100 private const val FRAME_CADENCE_MS = 20L diff --git a/nestsClient/src/jvmTest/kotlin/com/vitorpamplona/nestsclient/interop/NostrnestsProdAudioTransmissionTest.kt b/nestsClient/src/jvmTest/kotlin/com/vitorpamplona/nestsclient/interop/NostrnestsProdAudioTransmissionTest.kt index e919054ec..94db1551a 100644 --- a/nestsClient/src/jvmTest/kotlin/com/vitorpamplona/nestsclient/interop/NostrnestsProdAudioTransmissionTest.kt +++ b/nestsClient/src/jvmTest/kotlin/com/vitorpamplona/nestsclient/interop/NostrnestsProdAudioTransmissionTest.kt @@ -751,6 +751,293 @@ class NostrnestsProdAudioTransmissionTest { Unit } + // ==================================================================== + // Sweep matrix — calls SendTraceScenario with one dimension varied at + // a time so a single test run on a real network produces a diagnosis + // table covering cadence, frame count, payload, group packing, + // late-join, mid-stream pause, fan-out, and slow consumer. + // + // Every scenario is wrapped by [withProdSpeakerAndListeners], which + // mints publish + listen JWTs through the production auth side, opens + // the QUIC + WebTransport + moq-lite session the same way + // [connectNestsSpeaker] / [connectNestsListener] do, and tears down + // in the right order on success / failure. + // + // None of these scenarios assert "all frames received" by default — + // we want the cliff to surface in the report instead of failing fast, + // so the diagnostic dump is visible for every scenario in one run. + // The few scenarios where we DO expect success (e.g. very slow + // cadence) pass `expectAllReceived = true` explicitly. + // ==================================================================== + + @Test + fun sweep_baseline_100x20ms() = runProdScenarioOrSkip("sweep-baseline", Scenario(frameCount = 100, cadenceMs = 20L)) + + @Test fun sweep_cadence_5ms() = runProdScenarioOrSkip("sweep-cad5", Scenario(frameCount = 100, cadenceMs = 5L)) + + @Test fun sweep_cadence_10ms() = runProdScenarioOrSkip("sweep-cad10", Scenario(frameCount = 100, cadenceMs = 10L)) + + @Test fun sweep_cadence_40ms() = runProdScenarioOrSkip("sweep-cad40", Scenario(frameCount = 100, cadenceMs = 40L)) + + @Test fun sweep_cadence_80ms() = runProdScenarioOrSkip("sweep-cad80", Scenario(frameCount = 100, cadenceMs = 80L)) + + @Test fun sweep_cadence_200ms() = runProdScenarioOrSkip("sweep-cad200", Scenario(frameCount = 100, cadenceMs = 200L)) + + @Test + fun sweep_burst_no_cadence() = + runProdScenarioOrSkip( + "sweep-burst", + Scenario(frameCount = 100, cadenceMs = 0L, receiveGraceMs = 30_000L), + ) + + @Test + fun sweep_frames_50() = + runProdScenarioOrSkip( + "sweep-frames50", + Scenario(frameCount = 50, cadenceMs = 20L), + ) + + @Test + fun sweep_frames_200() = + runProdScenarioOrSkip( + "sweep-frames200", + Scenario(frameCount = 200, cadenceMs = 20L, receiveGraceMs = 45_000L), + ) + + @Test + fun sweep_frames_400() = + runProdScenarioOrSkip( + "sweep-frames400", + Scenario(frameCount = 400, cadenceMs = 20L, receiveGraceMs = 60_000L), + ) + + @Test + fun sweep_payload_1kb() = + runProdScenarioOrSkip( + "sweep-1kb", + Scenario(frameCount = 100, cadenceMs = 20L, payloadBytes = 1024), + ) + + @Test + fun sweep_payload_4kb() = + runProdScenarioOrSkip( + "sweep-4kb", + Scenario(frameCount = 100, cadenceMs = 20L, payloadBytes = 4096), + ) + + @Test + fun sweep_payload_16kb() = + runProdScenarioOrSkip( + "sweep-16kb", + Scenario(frameCount = 100, cadenceMs = 20L, payloadBytes = 16384), + ) + + @Test + fun sweep_frames_per_group_5() = + runProdScenarioOrSkip( + "sweep-fpg5", + Scenario(frameCount = 100, cadenceMs = 20L, framesPerGroup = 5), + ) + + @Test + fun sweep_frames_per_group_20() = + runProdScenarioOrSkip( + "sweep-fpg20", + Scenario(frameCount = 100, cadenceMs = 20L, framesPerGroup = 20), + ) + + @Test + fun sweep_frames_per_group_all() = + runProdScenarioOrSkip( + "sweep-fpg-all", + Scenario(frameCount = 100, cadenceMs = 20L, framesPerGroup = 100), + ) + + @Test + fun sweep_late_subscribe_after_25() = + runProdScenarioOrSkip( + "sweep-late25", + Scenario(frameCount = 100, cadenceMs = 20L, subscribeAtFrame = 25), + ) + + @Test + fun sweep_late_subscribe_after_50() = + runProdScenarioOrSkip( + "sweep-late50", + Scenario(frameCount = 100, cadenceMs = 20L, subscribeAtFrame = 50), + ) + + @Test + fun sweep_mid_pause_50_5s() = + runProdScenarioOrSkip( + "sweep-pause", + Scenario( + frameCount = 100, + cadenceMs = 20L, + pauseAfterFrame = 50, + pauseDurationMs = 5_000L, + receiveGraceMs = 45_000L, + ), + ) + + @Test + fun sweep_two_subscribers() = + runProdScenarioOrSkip( + "sweep-2subs", + Scenario(frameCount = 100, cadenceMs = 20L, parallelSubscriptions = 2), + ) + + @Test + fun sweep_three_subscribers() = + runProdScenarioOrSkip( + "sweep-3subs", + Scenario(frameCount = 100, cadenceMs = 20L, parallelSubscriptions = 3), + ) + + @Test + fun sweep_slow_consumer_50ms() = + runProdScenarioOrSkip( + "sweep-slowconsumer", + Scenario( + frameCount = 100, + cadenceMs = 20L, + listenerSlowConsumerMs = 50L, + receiveGraceMs = 60_000L, + ), + ) + + @Test + fun sweep_long_run_30s() = + runProdScenarioOrSkip( + "sweep-30s", + Scenario( + frameCount = 1500, + cadenceMs = 20L, + receiveGraceMs = 60_000L, + verbosePerFrame = false, + ), + ) + + @Test + fun sweep_extreme_long_run_120s() = + runProdScenarioOrSkip( + "sweep-120s", + Scenario( + frameCount = 6000, + cadenceMs = 20L, + receiveGraceMs = 180_000L, + verbosePerFrame = false, + ), + ) + + /** + * Wraps one scenario in `runBlocking { assumeProd(); withProd(…); … }` + * boilerplate so each `@Test` method is just a single Scenario + * literal. Skips cleanly when `-DnestsProd=true` isn't set. + */ + private fun runProdScenarioOrSkip( + scope: String, + scenario: Scenario, + expectAllReceived: Boolean = false, + ) = runBlocking { + assumeProd() + withProdSpeakerAndListeners(scope, scenario.parallelSubscriptions) { publisher, listeners, hostPub, pumpScope -> + val result = + SendTraceScenario.run( + scope = scope, + publisher = publisher, + listeners = listeners, + speakerPubkeyHex = hostPub, + scenario = scenario, + pumpScope = pumpScope, + ) + SendTraceScenario.reportAndAssert(scope, result, expectAllReceived) + } + Unit + } + + /** + * Build a host publisher (manual MoqLiteSession + session.publish so + * we bypass [com.vitorpamplona.nestsclient.audio.NestMoqLiteBroadcaster] + * and can call [com.vitorpamplona.nestsclient.moq.lite.MoqLitePublisherHandle.send] + * directly) plus N listener-side [com.vitorpamplona.nestsclient.NestsListener]s, + * each with their own keypair so the relay sees them as distinct + * audience members. Tears everything down (in the right order) on + * success or failure. + */ + private suspend fun withProdSpeakerAndListeners( + scope: String, + listenerCount: Int, + block: suspend ( + publisher: com.vitorpamplona.nestsclient.moq.lite.MoqLitePublisherHandle, + listeners: List, + speakerPubkeyHex: String, + pumpScope: CoroutineScope, + ) -> Unit, + ) { + val hostSigner = NostrSignerInternal(KeyPair()) + val room = freshRoom(hostPubkey = hostSigner.pubKey) + val httpClient = OkHttpNestsClient { OkHttpClient() } + val transport = QuicWebTransportFactory() + + val supervisor = SupervisorJob() + val pumpScope = CoroutineScope(supervisor + Dispatchers.IO) + + InteropDebug.checkpoint(scope, "auth=${room.authBaseUrl} endpoint=${room.endpoint} ns=${room.moqNamespace()}") + + val publishToken = + InteropDebug.stepSuspending(scope, "host: mintToken(publish=true)") { + httpClient.mintToken(room = room, publish = true, signer = hostSigner) + } + val (authority, path) = + com.vitorpamplona.nestsclient + .buildRelayConnectTarget( + endpoint = room.endpoint, + namespace = room.moqNamespace(), + token = publishToken, + ) + val speakerWt = + InteropDebug.stepSuspending(scope, "host: WebTransport.connect") { + transport.connect(authority = authority, path = path, bearerToken = null) + } + val speakerSession = + com.vitorpamplona.nestsclient.moq.lite + .MoqLiteSession + .client(speakerWt, pumpScope) + val publisher = + InteropDebug.stepSuspending(scope, "host: session.publish(broadcastSuffix=hostPub)") { + speakerSession.publish(broadcastSuffix = hostSigner.pubKey) + } + + val listeners = mutableListOf() + try { + for (i in 0 until listenerCount) { + val audienceSigner = NostrSignerInternal(KeyPair()) + val listener = + InteropDebug.stepSuspending(scope, "audience[$i]: connectNestsListener") { + connectNestsListener( + httpClient = httpClient, + transport = transport, + scope = pumpScope, + room = room, + signer = audienceSigner, + ) + } + InteropDebug.assertListenerReached(scope, "Connected", listener.state.value) + listeners += listener + } + + block(publisher, listeners, hostSigner.pubKey, pumpScope) + } finally { + for (listener in listeners) { + runCatching { listener.close() } + } + runCatching { publisher.close() } + runCatching { speakerWt.close(0, "test done") } + supervisor.cancelAndJoin() + } + } + // ----- helpers ----- private fun freshRoom(hostPubkey: String): NestsRoomConfig = diff --git a/nestsClient/src/jvmTest/kotlin/com/vitorpamplona/nestsclient/interop/SendTraceScenario.kt b/nestsClient/src/jvmTest/kotlin/com/vitorpamplona/nestsclient/interop/SendTraceScenario.kt new file mode 100644 index 000000000..4d0b8e8b3 --- /dev/null +++ b/nestsClient/src/jvmTest/kotlin/com/vitorpamplona/nestsclient/interop/SendTraceScenario.kt @@ -0,0 +1,463 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.nestsclient.interop + +import com.vitorpamplona.nestsclient.NestsListener +import com.vitorpamplona.nestsclient.moq.lite.MoqLitePublisherHandle +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Job +import kotlinx.coroutines.cancelAndJoin +import kotlinx.coroutines.delay +import kotlinx.coroutines.flow.onEach +import kotlinx.coroutines.flow.take +import kotlinx.coroutines.flow.toList +import kotlinx.coroutines.launch +import kotlinx.coroutines.withTimeoutOrNull +import kotlin.test.fail + +/** + * Knobs for one diagnostic pump scenario. Holds the parameters that the + * sweep matrix in [NostrnestsProdAudioTransmissionTest] / + * [NostrNestsSustainedSendOutcomesInteropTest] varies. + * + * Defaults match the production-realistic baseline (same shape that + * [NostrnestsProdAudioTransmissionTest.sustained_real_time_cadence_two_users] + * already uses): 100 Opus-sized frames, 20 ms cadence, one moq-lite + * group per frame. + */ +data class Scenario( + /** Number of frames the speaker pushes. */ + val frameCount: Int = 100, + /** Wall-clock delay between consecutive [send] calls. */ + val cadenceMs: Long = 20L, + /** Bytes per frame. Production Opus @ 32 kbps / 20 ms ≈ 80. */ + val payloadBytes: Int = 80, + /** + * How many frames per moq-lite group. The production broadcaster + * uses `1` (one Opus packet → one uni stream). Higher values pack + * multiple `varint(size)+payload` frames into the same uni stream + * to test whether the loss is per-stream-creation or per-byte. + */ + val framesPerGroup: Int = 1, + /** + * Wait this long between issuing the SUBSCRIBE and starting to + * pump frames so SUBSCRIBE_OK lands at the speaker before the + * first send (otherwise [MoqLitePublisherHandle.send] returns false + * for early frames because `inboundSubs.isEmpty()`). + */ + val subscribeSettleMs: Long = 500L, + /** + * Total wall-clock budget for the listener flow to drain + * [Scenario.frameCount] frames. Includes the pump time itself, so + * give plenty of grace beyond `frameCount * cadenceMs`. + */ + val receiveGraceMs: Long = 30_000L, + /** + * Pause this long after pushing [pauseAfterFrame] frames before + * resuming the pump. -1 disables the pause. Lets us see whether + * the connection wedges during idle time (mid-stream uni-stream + * cancellation, RTO PTO, etc). + */ + val pauseAfterFrame: Int = -1, + /** Pause duration when [pauseAfterFrame] >= 0. */ + val pauseDurationMs: Long = 0L, + /** + * Subscribe AFTER the speaker has already pushed this many frames + * (or run for this many milliseconds — see [subscribeAtFrame]). + * Tests "from-latest" join semantics — the listener should pick up + * from the next group boundary, not from the very first. + */ + val subscribeAtFrame: Int = -1, + /** Listener consumer adds this delay between each frame consumed. */ + val listenerSlowConsumerMs: Long = 0L, + /** + * Number of independent listener subscriptions to attach. > 1 + * exercises fan-out; the relay forwards each group to every + * subscriber over an independent uni stream. + */ + val parallelSubscriptions: Int = 1, + /** Larger payloads (e.g. 4 KB) test stream-write vs stream-creation cost. */ + val verbosePerFrame: Boolean = true, +) + +/** + * One-frame arrival record. Captures wall-clock timing relative to the + * start of the collector so we can graph cadence drift / stalls + * post-hoc. + */ +data class FrameArrival( + val subscriberIndex: Int, + val wallMs: Long, + val groupId: Long, + val objectId: Long, + val firstByte: Int, +) + +/** + * Outcome of one [SendTraceScenario.run] call. The fields are + * deliberately denormalised so a single object dumps cleanly through + * [InteropDebug] and is easy to read in the failure message. + */ +data class ScenarioResult( + val scenario: Scenario, + val sendOutcomes: BooleanArray, + val sendDurationsMicros: LongArray, + val endGroupErrors: Array, + val arrivalsPerSubscriber: List>, + val pumpStartedAtMs: Long, + val pumpDurationMs: Long, + val collectStartedAtMs: Long, +) { + val sendTrueCount: Int get() = sendOutcomes.count { it } + val firstFalseSendIndex: Int get() = sendOutcomes.indexOfFirst { !it } + + /** + * Per subscriber: `{objectId -> arrival}`. We index by objectId + * (the per-session monotonic counter from + * [com.vitorpamplona.nestsclient.MoqLiteNestsListener]) rather than + * groupId so the sweep stays correct when multiple frames pack + * into the same moq-lite group (`framesPerGroup > 1`). Each frame + * the test pushes corresponds to exactly one object, and objectIds + * arrive 0..N-1 in send order. + */ + val arrivalIndexPerSubscriber: List> by lazy { + arrivalsPerSubscriber.map { list -> list.associateBy { it.objectId } } + } + + fun missingObjectIdsFor(subscriberIndex: Int): List { + val seen = arrivalIndexPerSubscriber[subscriberIndex].keys + return (0L until scenario.frameCount.toLong()).filter { it !in seen } + } + + fun firstSentButLost(subscriberIndex: Int): Int { + val seen = arrivalIndexPerSubscriber[subscriberIndex].keys + return (0 until scenario.frameCount).firstOrNull { sendOutcomes[it] && it.toLong() !in seen } ?: -1 + } +} + +/** + * Diagnostic per-frame instrumented pump shared by the production and + * local-harness end-to-end tests. Bypasses + * [com.vitorpamplona.nestsclient.audio.NestMoqLiteBroadcaster] so we + * can capture the [MoqLitePublisherHandle.send] boolean per frame + * (the production broadcaster swallows it via `runCatching {…}`). + * + * Logs everything via [InteropDebug.checkpoint]: + * - Scenario parameters at start + * - Per-frame: send result, send duration, endGroup exception (if any) + * - Per-arrival: subscriber index, wall-clock ms, groupId, objectId, first payload byte + * - Final summary: counts, missing-group run-length, first false-send, + * first sent-but-lost + * + * Caller responsibilities: + * - Provide a [MoqLitePublisherHandle] that has already + * `MoqLiteSession.publish()`'d a broadcast suffix. + * - Provide one or more [NestsListener]s, NOT yet subscribed (the + * scenario calls [NestsListener.subscribeSpeaker] internally so it + * can time the subscription relative to the pump for late-join + * scenarios). + * - Tear down all of those after [run] returns. + */ +object SendTraceScenario { + suspend fun run( + scope: String, + publisher: MoqLitePublisherHandle, + listeners: List, + speakerPubkeyHex: String, + scenario: Scenario, + pumpScope: CoroutineScope, + ): ScenarioResult { + require(listeners.size == scenario.parallelSubscriptions) { + "expected ${scenario.parallelSubscriptions} listener(s), got ${listeners.size}" + } + InteropDebug.checkpoint(scope, "scenario=$scenario speaker=${speakerPubkeyHex.take(8)}…") + + val sendOutcomes = BooleanArray(scenario.frameCount) + val sendDurationsMicros = LongArray(scenario.frameCount) + val endGroupErrors = arrayOfNulls(scenario.frameCount) + val arrivalsPerSubscriber = + List(scenario.parallelSubscriptions) { + java.util.concurrent.CopyOnWriteArrayList() + } + val collectStart = System.currentTimeMillis() + + // Subscribe before the pump starts UNLESS subscribeAtFrame says + // we should wait. We launch one collector per listener; each + // takes [scenario.frameCount] frames or runs to grace timeout. + val collectorJobs = mutableListOf() + if (scenario.subscribeAtFrame < 0) { + for ((idx, listener) in listeners.withIndex()) { + collectorJobs += spawnCollector(scope, idx, listener, speakerPubkeyHex, scenario, arrivalsPerSubscriber[idx], collectStart, pumpScope) + } + delay(scenario.subscribeSettleMs) + InteropDebug.checkpoint(scope, "settled (${scenario.subscribeSettleMs}ms); starting pump") + } else { + InteropDebug.checkpoint( + scope, + "deferred subscribe — listeners will subscribe AFTER frame ${scenario.subscribeAtFrame}", + ) + } + + val pumpStart = System.currentTimeMillis() + val payloadPrefix = ByteArray(scenario.payloadBytes - 1) { 0x4F.toByte() } + for (i in 0 until scenario.frameCount) { + // Late-join hook: drop the subscriptions in *now* if we + // crossed the configured frame index. + if (scenario.subscribeAtFrame == i) { + for ((idx, listener) in listeners.withIndex()) { + collectorJobs += + spawnCollector( + scope, + idx, + listener, + speakerPubkeyHex, + scenario, + arrivalsPerSubscriber[idx], + collectStart, + pumpScope, + ) + } + InteropDebug.checkpoint( + scope, + "subscribed at frame $i (t=${System.currentTimeMillis() - collectStart}ms)", + ) + } + if (scenario.pauseAfterFrame == i && scenario.pauseDurationMs > 0) { + InteropDebug.checkpoint( + scope, + "pause: idling ${scenario.pauseDurationMs}ms after frame $i (t=${System.currentTimeMillis() - collectStart}ms)", + ) + delay(scenario.pauseDurationMs) + InteropDebug.checkpoint(scope, "pause ended (t=${System.currentTimeMillis() - collectStart}ms)") + } + + val payload = payloadPrefix + byteArrayOf(i.toByte()) + val sendStartedNs = System.nanoTime() + sendOutcomes[i] = publisher.send(payload) + val sendElapsedNs = System.nanoTime() - sendStartedNs + sendDurationsMicros[i] = sendElapsedNs / 1_000 + + // endGroup boundary: only end at the configured cadence + // (default = every frame, which matches production). + val isGroupBoundary = ((i + 1) % scenario.framesPerGroup) == 0 || i == scenario.frameCount - 1 + if (isGroupBoundary) { + runCatching { publisher.endGroup() } + .onFailure { endGroupErrors[i] = it::class.simpleName + ": " + it.message } + } + if (scenario.verbosePerFrame) { + val tag = if (sendOutcomes[i]) "ok" else "send=false" + val err = endGroupErrors[i]?.let { ", endGroup err=$it" } ?: "" + val groupBoundary = if (isGroupBoundary) " ⏎" else "" + InteropDebug.checkpoint( + scope, + "tx i=$i $tag dt=${sendDurationsMicros[i]}us$groupBoundary$err " + + "(t=${System.currentTimeMillis() - collectStart}ms)", + ) + } + if (i < scenario.frameCount - 1 && scenario.cadenceMs > 0) { + delay(scenario.cadenceMs) + } + } + val pumpDuration = System.currentTimeMillis() - pumpStart + InteropDebug.checkpoint( + scope, + "pump done: ${scenario.frameCount} frames in ${pumpDuration}ms " + + "(target=${scenario.frameCount * scenario.cadenceMs}ms) " + + "sendTrue=${sendOutcomes.count { it }}/${scenario.frameCount}", + ) + + // Wait for collectors. If they hit `take(N)` they exit naturally; + // otherwise the per-collector withTimeoutOrNull cancels them. + for (job in collectorJobs) { + withTimeoutOrNull(scenario.receiveGraceMs - (System.currentTimeMillis() - collectStart)) { + job.join() + } + if (job.isActive) job.cancelAndJoin() + } + + return ScenarioResult( + scenario = scenario, + sendOutcomes = sendOutcomes, + sendDurationsMicros = sendDurationsMicros, + endGroupErrors = endGroupErrors, + arrivalsPerSubscriber = arrivalsPerSubscriber.map { it.toList() }, + pumpStartedAtMs = pumpStart - collectStart, + pumpDurationMs = pumpDuration, + collectStartedAtMs = collectStart, + ) + } + + private fun spawnCollector( + scope: String, + subscriberIndex: Int, + listener: NestsListener, + speakerPubkeyHex: String, + scenario: Scenario, + sink: java.util.concurrent.CopyOnWriteArrayList, + collectStart: Long, + pumpScope: CoroutineScope, + ): Job = + pumpScope.launch { + val sub = listener.subscribeSpeaker(speakerPubkeyHex) + try { + withTimeoutOrNull(scenario.receiveGraceMs) { + sub.objects + .onEach { obj -> + val now = System.currentTimeMillis() + sink += + FrameArrival( + subscriberIndex = subscriberIndex, + wallMs = now - collectStart, + groupId = obj.groupId, + objectId = obj.objectId, + firstByte = + obj.payload + .firstOrNull() + ?.toInt() + ?.and(0xFF) ?: -1, + ) + if (scenario.verbosePerFrame) { + InteropDebug.checkpoint( + scope, + "rx[$subscriberIndex] gid=${obj.groupId} oid=${obj.objectId} " + + "firstByte=0x${( + obj.payload + .firstOrNull() + ?.toInt() + ?.and(0xFF) ?: -1 + ).toString(16)} " + + "(t=${now - collectStart}ms)", + ) + } + if (scenario.listenerSlowConsumerMs > 0) { + delay(scenario.listenerSlowConsumerMs) + } + }.take(scenario.frameCount) + .toList() + } + } finally { + runCatching { sub.unsubscribe() } + } + } + + /** + * Pretty-print the result + assert the expected outcome. + * + * @param expectAllReceived true → every frame must have arrived at + * every subscriber. false → just dump the diagnostic summary + * (lets a "reproduce the cliff" scenario succeed even when the + * relay drops trailing frames). + */ + fun reportAndAssert( + scope: String, + result: ScenarioResult, + expectAllReceived: Boolean, + ) { + val s = result.scenario + for (subscriberIndex in 0 until s.parallelSubscriptions) { + val arrivals = result.arrivalsPerSubscriber[subscriberIndex] + val missing = result.missingObjectIdsFor(subscriberIndex) + val firstSentButLost = result.firstSentButLost(subscriberIndex) + val firstArrival = arrivals.firstOrNull()?.wallMs ?: -1 + val lastArrival = arrivals.lastOrNull()?.wallMs ?: -1 + InteropDebug.checkpoint( + scope, + "sub[$subscriberIndex] received=${arrivals.size}/${s.frameCount} " + + "firstArrival=${firstArrival}ms lastArrival=${lastArrival}ms " + + "firstFalseSend=${result.firstFalseSendIndex} " + + "firstSentButLost=$firstSentButLost " + + "missing=${runs(missing)}", + ) + } + + // Average / min / max / p99 send latency, plus count of sends + // that took > 1 ms (a hint that publisher.send() blocked on a + // mutex / IO). + val nonZero = result.sendDurationsMicros.filter { it > 0 } + if (nonZero.isNotEmpty()) { + val sorted = nonZero.sorted() + val p50 = sorted[sorted.size / 2] + val p99 = sorted[(sorted.size * 99 / 100).coerceAtMost(sorted.size - 1)] + val max = sorted.last() + val sloweries = sorted.count { it > 1_000 } + InteropDebug.checkpoint( + scope, + "send latency: p50=${p50}us p99=${p99}us max=${max}us sends>1ms=$sloweries", + ) + } + if (result.endGroupErrors.any { it != null }) { + val errs = result.endGroupErrors.withIndex().filter { it.value != null } + InteropDebug.checkpoint( + scope, + "endGroup errors: ${errs.joinToString { "i=${it.index} ${it.value}" }}", + ) + } + + if (expectAllReceived) { + for (subIdx in 0 until s.parallelSubscriptions) { + val arrivals = result.arrivalsPerSubscriber[subIdx] + if (arrivals.size != s.frameCount) { + fail( + "[$scope] subscriber $subIdx received ${arrivals.size}/${s.frameCount} — " + + "missing=${runs(result.missingObjectIdsFor(subIdx))} " + + "firstSentButLost=${result.firstSentButLost(subIdx)} " + + "firstFalseSend=${result.firstFalseSendIndex}", + ) + } + } + } + } + + /** Compress a sorted list of integers into run-length ranges, e.g. `[3-7,42,90-99]`. */ + private fun runs(values: List): String { + if (values.isEmpty()) return "[]" + return buildString { + append('[') + var run: Pair? = null + + fun flush() { + run?.let { + if (length > 1) append(',') + if (it.first == it.second) append(it.first) else append("${it.first}-${it.second}") + } + } + for (v in values) { + run = + when { + run == null -> { + v to v + } + + v == run.second + 1 -> { + run.first to v + } + + else -> { + flush() + v to v + } + } + } + flush() + append(']') + } + } +}