diff --git a/nestsClient-browser-interop/.gitignore b/nestsClient-browser-interop/.gitignore new file mode 100644 index 000000000..43b0118ba --- /dev/null +++ b/nestsClient-browser-interop/.gitignore @@ -0,0 +1,5 @@ +node_modules/ +dist/ +test-results/ +playwright-report/ +.bun/ diff --git a/nestsClient-browser-interop/REV b/nestsClient-browser-interop/REV new file mode 100644 index 000000000..04c30b6ae --- /dev/null +++ b/nestsClient-browser-interop/REV @@ -0,0 +1,31 @@ +# Pinned upstream npm package versions for the browser-side cross-stack +# interop harness (Phase 4 of T16). +# +# These versions are what `nestsClient-browser-interop/package.json` pins +# and what the bun build resolves at install time. Bumping requires +# touching package.json + bun.lockb + this file together so a silent +# upstream rev change can't mask a regression. +# +# See: nestsClient/plans/2026-05-06-phase4-browser-harness.md +# +# Source: https://github.com/kixelated/moq , workspace published to npm +# under the @moq/* scope. The `@moq/lite` 0.2.x line implements +# `moq-lite-03` (see /tmp/moq/js/lite/src/lite/), matching the pin in +# nestsClient/tests/hang-interop/REV (KIXELATED_MOQ_GIT_REV). + +# Browser listener: builds Watch.Broadcast on top of @moq/lite + +# @moq/hang. We use @moq/lite + @moq/hang directly for the harness path +# (closer to nestsClient's own moq-lite stack) but keep @moq/watch +# pinned in case a future scenario wants the higher-level reactive +# Broadcast wrapper. +MOQ_WATCH_VERSION=0.2.10 + +# Browser publisher: same story for @moq/publish. +MOQ_PUBLISH_VERSION=0.2.6 + +# Lower-level moq-lite-03 client + hang catalog/container. +MOQ_LITE_VERSION=0.2.2 +MOQ_HANG_VERSION=0.2.4 + +# Playwright Chromium driver. +PLAYWRIGHT_VERSION=1.56.1 diff --git a/nestsClient-browser-interop/bun.lock b/nestsClient-browser-interop/bun.lock new file mode 100644 index 000000000..b3b26f1bf --- /dev/null +++ b/nestsClient-browser-interop/bun.lock @@ -0,0 +1,73 @@ +{ + "lockfileVersion": 1, + "configVersion": 1, + "workspaces": { + "": { + "name": "nests-browser-interop", + "dependencies": { + "@moq/hang": "0.2.4", + "@moq/lite": "0.2.2", + "@moq/publish": "0.2.6", + "@moq/watch": "0.2.10", + }, + "devDependencies": { + "@playwright/test": "1.56.1", + "@types/bun": "latest", + "typescript": "^5.6.0", + }, + }, + }, + "packages": { + "@kixelated/libavjs-webcodecs-polyfill": ["@kixelated/libavjs-webcodecs-polyfill@0.5.5", "", { "dependencies": { "@libav.js/types": "^6.7.7", "@ungap/global-this": "^0.4.4" } }, "sha512-Q1zgnTMMQ2F7IE9ylx3C1XzVbg5vYN18jiDINO5U3kNPBOHdYuUlJsMhtBoqr1M6ocLtoiqdHmLs7tHFgrw5KA=="], + + "@libav.js/types": ["@libav.js/types@6.8.8", "", {}, "sha512-Lbik/0Q3x2R8cI7mOtRgt+nUWLqGXh7UinMndmpdXSDY4YEjYyVUDsq6fxkuriL78+LCYx8frZIN1r+oDsvYCQ=="], + + "@libav.js/variant-opus-af": ["@libav.js/variant-opus-af@6.8.8", "", {}, "sha512-8KBQyA8n5goN7lyctOaPxpcx7dapOgqKh8dWW/NAcl87AgM/WoUGSex3fFc46oCtTHYrUKEm1OmZUrtkt3Q56A=="], + + "@moq/hang": ["@moq/hang@0.2.4", "", { "dependencies": { "@kixelated/libavjs-webcodecs-polyfill": "^0.5.5", "@libav.js/variant-opus-af": "^6.8.8", "@moq/lite": "^0.2.2", "@moq/signals": "^0.1.6", "@svta/cml-iso-bmff": "^1.0.0-alpha.9", "zod": "^4.1.5" } }, "sha512-I7OzutII+Sp5oWKd33t6b1SSY9Tu2dpu6pEaUp7CAKzNRpIaE7O6DhWBsBlcGAMTuDj/zA3CkR3iGx7WW2WW6w=="], + + "@moq/lite": ["@moq/lite@0.2.2", "", { "dependencies": { "@moq/qmux": "^0.0.6", "@moq/signals": "^0.1.6", "async-mutex": "^0.5.0" }, "peerDependencies": { "zod": "^4.0.0" } }, "sha512-o5X4qQlfhO8xWcWwpYsEEWWHr66SIPQKFY++ZxgMYhU+77rTfe1vWgomYjqoQcCZT3flRY2iTRLJriHgyDX/gA=="], + + "@moq/msf": ["@moq/msf@0.1.0", "", { "dependencies": { "@moq/lite": "^0.2.1", "zod": "^4.1.5" } }, "sha512-5Y/RcxxofBXQSdy6IexzB8s2rpRI7xFiut1Zh6WO6hjNNqM+WKPPt+CTGKqnnnx/vhecgKrvpHnN228Zc0bakg=="], + + "@moq/publish": ["@moq/publish@0.2.6", "", { "dependencies": { "@moq/hang": "^0.2.4", "@moq/lite": "^0.2.2", "@moq/signals": "^0.1.6", "@moq/ui-core": "^0.1.0" } }, "sha512-cAHt8ZRMKOZh/yd1CX8wS6N5XLCGxzsKB/hU7R9PtmlJ40U3hDKokMGSd+jYWstDOgppoMwr4qU//yjDyeOiNw=="], + + "@moq/qmux": ["@moq/qmux@0.0.6", "", {}, "sha512-ISuGz05lUvf1hzHW3Aw3VnsGRJe1w9Qdog3LQ66KS+l+5mzQsPANvW8yOioEe1Z9dJO2G3sAHoGPnzwnsY9SIQ=="], + + "@moq/signals": ["@moq/signals@0.1.6", "", { "peerDependencies": { "@types/react": "^19.1.8", "react": "^19.0.0", "solid-js": "^1.9.7" }, "optionalPeers": ["@types/react", "react", "solid-js"] }, "sha512-ic7ttiz6dHXOPoVAfhz4K6LGT2LWdDGTi1x2u8sYSGZ5nOKGWfqDkwYcGvCPlcVQetn3PaeXYSPFiMAC6RO3tQ=="], + + "@moq/ui-core": ["@moq/ui-core@0.1.0", "", { "peerDependencies": { "@moq/signals": "^0.1.2" } }, "sha512-DJNBpUNQDyh7Tou324fbJ5/pT08UPghH3OxcVdLEo9IQeX//8NiEzJwcX7iuacR32nzgdiBThIbIpeFa60U3/g=="], + + "@moq/watch": ["@moq/watch@0.2.10", "", { "dependencies": { "@moq/hang": "^0.2.4", "@moq/lite": "^0.2.2", "@moq/msf": "^0.1.0", "@moq/signals": "^0.1.6", "@moq/ui-core": "^0.1.0" } }, "sha512-uLVwdtx0XIvJ20c1dYJ5NIVLXBA/cbTNzvM1mujvYLVKGQnwPkrrllEFlNpWfwV9SN1Kb8BnDvA6LLtYTFAhkQ=="], + + "@playwright/test": ["@playwright/test@1.56.1", "", { "dependencies": { "playwright": "1.56.1" }, "bin": { "playwright": "cli.js" } }, "sha512-vSMYtL/zOcFpvJCW71Q/OEGQb7KYBPAdKh35WNSkaZA75JlAO8ED8UN6GUNTm3drWomcbcqRPFqQbLae8yBTdg=="], + + "@svta/cml-iso-bmff": ["@svta/cml-iso-bmff@1.0.1", "", { "peerDependencies": { "@svta/cml-utils": "1.4.0" } }, "sha512-MOhATJYQ6cVrIcoY3nj8p/vGYDpG3wjQIIhBPHNt9yjFijdwFdBNqdZbCXv3aFhRjdx5Saca5TkgNJusKhnI/w=="], + + "@svta/cml-utils": ["@svta/cml-utils@1.4.0", "", {}, "sha512-vNtHtv/z+9I9ysxFwNrgwxic1oceVPr8TpcpV/NA1l8Gy4phynwtOppkCIBB+PmoyKDcqE4lO85g+lfsuSTBBA=="], + + "@types/bun": ["@types/bun@1.3.13", "", { "dependencies": { "bun-types": "1.3.13" } }, "sha512-9fqXWk5YIHGGnUau9TEi+qdlTYDAnOj+xLCmSTwXfAIqXr2x4tytJb43E9uCvt09zJURKXwAtkoH4nLQfzeTXw=="], + + "@types/node": ["@types/node@25.6.0", "", { "dependencies": { "undici-types": "~7.19.0" } }, "sha512-+qIYRKdNYJwY3vRCZMdJbPLJAtGjQBudzZzdzwQYkEPQd+PJGixUL5QfvCLDaULoLv+RhT3LDkwEfKaAkgSmNQ=="], + + "@ungap/global-this": ["@ungap/global-this@0.4.4", "", {}, "sha512-mHkm6FvepJECMNthFuIgpAEFmPOk71UyXuIxYfjytvFTnSDBIz7jmViO+LfHI/AjrazWije0PnSP3+/NlwzqtA=="], + + "async-mutex": ["async-mutex@0.5.0", "", { "dependencies": { "tslib": "^2.4.0" } }, "sha512-1A94B18jkJ3DYq284ohPxoXbfTA5HsQ7/Mf4DEhcyLx3Bz27Rh59iScbB6EPiP+B+joue6YCxcMXSbFC1tZKwA=="], + + "bun-types": ["bun-types@1.3.13", "", { "dependencies": { "@types/node": "*" } }, "sha512-QXKeHLlOLqQX9LgYaHJfzdBaV21T63HhFJnvuRCcjZiaUDpbs5ED1MgxbMra71CsryN/1dAoXuJJJwIv/2drVA=="], + + "fsevents": ["fsevents@2.3.2", "", { "os": "darwin" }, "sha512-xiqMQR4xAeHTuB9uWm+fFRcIOgKBMiOBP+eXiyT7jsgVCq1bkVygt00oASowB7EdtpOHaaPgKt812P9ab+DDKA=="], + + "playwright": ["playwright@1.56.1", "", { "dependencies": { "playwright-core": "1.56.1" }, "optionalDependencies": { "fsevents": "2.3.2" }, "bin": { "playwright": "cli.js" } }, "sha512-aFi5B0WovBHTEvpM3DzXTUaeN6eN0qWnTkKx4NQaH4Wvcmc153PdaY2UBdSYKaGYw+UyWXSVyxDUg5DoPEttjw=="], + + "playwright-core": ["playwright-core@1.56.1", "", { "bin": { "playwright-core": "cli.js" } }, "sha512-hutraynyn31F+Bifme+Ps9Vq59hKuUCz7H1kDOcBs+2oGguKkWTU50bBWrtz34OUWmIwpBTWDxaRPXrIXkgvmQ=="], + + "tslib": ["tslib@2.8.1", "", {}, "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w=="], + + "typescript": ["typescript@5.9.3", "", { "bin": { "tsc": "bin/tsc", "tsserver": "bin/tsserver" } }, "sha512-jl1vZzPDinLr9eUt3J/t7V6FgNEw9QjvBPdysz9KfQDD41fQrC2Y4vKQdiaUpFT4bXlb1RHhLpp8wtm6M5TgSw=="], + + "undici-types": ["undici-types@7.19.2", "", {}, "sha512-qYVnV5OEm2AW8cJMCpdV20CDyaN3g0AjDlOGf1OW4iaDEx8MwdtChUp4zu4H0VP3nDRF/8RKWH+IPp9uW0YGZg=="], + + "zod": ["zod@4.4.3", "", {}, "sha512-ytENFjIJFl2UwYglde2jchW2Hwm4GJFLDiSXWdTrJQBIN9Fcyp7n4DhxJEiWNAJMV1/BqWfW/kkg71UDcHJyTQ=="], + } +} diff --git a/nestsClient-browser-interop/package.json b/nestsClient-browser-interop/package.json new file mode 100644 index 000000000..313a1321f --- /dev/null +++ b/nestsClient-browser-interop/package.json @@ -0,0 +1,23 @@ +{ + "name": "nests-browser-interop", + "version": "0.0.0", + "private": true, + "type": "module", + "description": "Phase 4 (T16) browser-side cross-stack interop harness — headless Chromium running @moq/lite + @moq/hang against the same NativeMoqRelayHarness moq-relay subprocess that drives HangInteropTest. Lands behind -DnestsBrowserInterop=true.", + "scripts": { + "build": "bun build src/listen.ts src/publish.ts --outdir dist --target browser && cp src/listen.html src/publish.html src/pcm-tap-worklet.js dist/", + "serve": "bun run src/server.ts", + "playwright": "playwright test" + }, + "dependencies": { + "@moq/hang": "0.2.4", + "@moq/lite": "0.2.2", + "@moq/publish": "0.2.6", + "@moq/watch": "0.2.10" + }, + "devDependencies": { + "@playwright/test": "1.56.1", + "@types/bun": "latest", + "typescript": "^5.6.0" + } +} diff --git a/nestsClient-browser-interop/playwright.config.ts b/nestsClient-browser-interop/playwright.config.ts new file mode 100644 index 000000000..143f4f7de --- /dev/null +++ b/nestsClient-browser-interop/playwright.config.ts @@ -0,0 +1,56 @@ +import { defineConfig } from "@playwright/test"; + +// Phase 4 (T16) browser-interop Playwright config. +// +// One-off Chromium spawn per Kotlin test. We disable the default test +// projects + reporters (the runner is invoked headlessly from the +// PlaywrightDriver Kotlin shim with `--reporter list` for stdout +// streaming). +// +// Chromium flags: +// --enable-quic — required for WebTransport. +// --ignore-certificate-errors — accept the self-signed +// cert moq-relay generates +// with --tls-generate. +// --enable-features=AutoplayPolicy=NoUserGestureRequired +// — let AudioContext.resume() +// succeed without a user +// gesture (we're headless). +// --enable-blink-features=WebTransport +// — defensively re-enable in +// case the build disables +// the blink feature flag +// by default. + +export default defineConfig({ + testDir: "./tests", + fullyParallel: false, + workers: 1, + forbidOnly: !!process.env.CI, + retries: 0, + reporter: process.env.PLAYWRIGHT_REPORTER ?? "list", + timeout: 120_000, + use: { + headless: true, + trace: "off", + video: "off", + screenshot: "off", + launchOptions: { + args: [ + "--enable-quic", + "--ignore-certificate-errors", + "--enable-features=AutoplayPolicy=NoUserGestureRequired", + "--enable-blink-features=WebTransport", + // Disable network sandbox so WebTransport over loopback + // doesn't trip the network service sandbox in headless. + "--disable-features=IsolateOrigins,site-per-process", + ], + }, + }, + projects: [ + { + name: "chromium", + use: { browserName: "chromium" }, + }, + ], +}); diff --git a/nestsClient-browser-interop/src/listen.html b/nestsClient-browser-interop/src/listen.html new file mode 100644 index 000000000..e0321168b --- /dev/null +++ b/nestsClient-browser-interop/src/listen.html @@ -0,0 +1,17 @@ + + + + +nests browser-interop listener + + + +

nests-browser-interop / listen

+
init
+ + + diff --git a/nestsClient-browser-interop/src/listen.ts b/nestsClient-browser-interop/src/listen.ts new file mode 100644 index 000000000..90f8a3091 --- /dev/null +++ b/nestsClient-browser-interop/src/listen.ts @@ -0,0 +1,249 @@ +// Phase 4 (T16) browser-listener harness. Connects to the +// `NativeMoqRelayHarness` moq-relay subprocess via WebTransport, subscribes +// to the `/audio/data` track produced by the Amethyst Kotlin +// speaker, decodes each Opus packet via WebCodecs AudioDecoder, and posts +// the resulting Float32 PCM samples back to the bun WS server (`ws://`), +// which appends them to a file on disk for the Kotlin test to read. +// +// Reads its parameters from `location.search`: +// +// relay — the relay's WebTransport URL, +// e.g. `https://127.0.0.1:43219/nests/::?jwt=`. +// Pass the FULL connection target (path + query) — Amethyst's +// nests namespace is part of the relay path per `NestsConnect.kt`. +// broadcast — the publisher's moq-lite broadcast path +// (= `speakerPubkeyHex` per `MoqLiteNestsSpeaker.kt`). +// track — the audio track name. Defaults to `audio/data`. +// wsPort — the bun WS back-channel port; we POST PCM here. +// duration — broadcast capture window in seconds. +// +// Mirrors the data path of `kixelated/moq` `js/watch/src/audio/decoder.ts` +// (`#runLegacyDecoder`) with the `@moq/hang` `Container.Legacy.Format` consumer +// and the WebCodecs AudioDecoder warmup-skip semantics. Verbatim matching +// the watcher's per-frame behaviour is what catches a Chromium-side +// regression that wire-byte tests can't see. + +import * as Moq from "@moq/lite"; +import * as Container from "@moq/hang/container"; +import { PRIORITY as CATALOG_PRIORITY } from "@moq/hang/catalog"; + +const params = new URLSearchParams(location.search); +const relayParam = params.get("relay"); +const broadcastParam = params.get("broadcast"); +const trackParam = params.get("track") ?? "audio/data"; +const wsPort = Number(params.get("wsPort") ?? "0"); +const durationSec = Number(params.get("duration") ?? "5"); +const certSha256B64 = params.get("certSha256"); // Base64 SHA-256 of leaf DER cert. + +function required(v: string | null, name: string): string { + if (!v) throw new Error(`listen.html: missing ?${name}=`); + return v; +} +const relayUrlString = required(relayParam, "relay"); +const broadcastName = required(broadcastParam, "broadcast"); +if (!wsPort) throw new Error("listen.html: missing ?wsPort="); + +const status = (msg: string) => { + const el = document.getElementById("status"); + if (el) el.textContent = msg; + console.log("[listen]", msg); +}; + +const fail = (msg: string) => { + status(`ERROR: ${msg}`); + document.body.dataset.state = "error"; + throw new Error(msg); +}; + +async function main() { + // -- WS back-channel ------------------------------------------------ + // The bun server appends every binary message we send to a PCM file + // on disk. We open it BEFORE the WebTransport so the very first frame + // (which decodes via WebCodecs after Container.Legacy strips the 2-byte + // timestamp) is captured even if it arrives before the page reaches + // its `done` state. + const ws = new WebSocket(`ws://127.0.0.1:${wsPort}/pcm`); + ws.binaryType = "arraybuffer"; + await new Promise((resolve, reject) => { + ws.addEventListener("open", () => resolve(), { once: true }); + ws.addEventListener("error", () => reject(new Error("ws connect failed")), { once: true }); + }); + status("ws connected"); + + const sendPcm = (chunk: Float32Array) => { + // Float32 LE matches the format hang-listen writes; the Kotlin + // test reads it via `readFloat32Pcm`. + if (ws.readyState === WebSocket.OPEN) ws.send(chunk.buffer); + }; + const sendDone = () => { + if (ws.readyState === WebSocket.OPEN) ws.send("done"); + }; + + // -- Connect to the relay ------------------------------------------ + // `relayParam` already includes the namespace path + ?jwt=… query + // (built by `buildRelayConnectTarget` in NestsConnect.kt). We pass it + // straight to `Connection.connect` which feeds it to `new WebTransport(url)` + // verbatim. Self-signed cert pinning is via Chromium's + // `--ignore-certificate-errors` flag — we do NOT compute a SHA-256 hash + // since the relay's auto-generated cert isn't deterministic. + const relayUrl = new URL(relayUrlString); + status(`connecting to ${relayUrl.toString()}`); + // If the test driver passed a leaf-cert SHA-256, pin it via + // `serverCertificateHashes`. Chromium's `--ignore-certificate-errors` + // does NOT bypass QUIC cert validation (crbug.com/1190655), so this + // is the supported path for self-signed test certs over WebTransport. + // The hash is base64 — convert to a Uint8Array. Fail loudly if the + // hash is malformed; falling back to no-pin would just produce a + // QUIC_TLS_CERTIFICATE_UNKNOWN error one round-trip later. + const webtransportOpts: WebTransportOptions = {}; + if (certSha256B64) { + const raw = Uint8Array.from(atob(certSha256B64), (c) => c.charCodeAt(0)); + webtransportOpts.serverCertificateHashes = [ + { algorithm: "sha-256", value: raw }, + ]; + } + const conn = await Moq.Connection.connect(relayUrl, { + // Disable the WebSocket fallback — the harness relay only speaks QUIC. + websocket: { enabled: false }, + webtransport: webtransportOpts, + }); + status(`connected, alpn=${conn.version}`); + + // Expose for Playwright to read post-hoc. + (window as any).__moqVersion = conn.version; + + // -- Subscribe to the audio track ---------------------------------- + const broadcastPath = Moq.Path.from(broadcastName); + const broadcast = conn.consume(broadcastPath); + const track = broadcast.subscribe(trackParam, CATALOG_PRIORITY.audio); + status(`subscribed broadcast=${broadcastName} track=${trackParam}`); + + // The hang Container.Legacy.Consumer strips the Varint-encoded + // timestamp prefix (per `kixelated/moq/js/hang/src/container/legacy.ts`) + // and yields an Opus packet per `next()`. Mirrors the data path the + // @moq/watch decoder uses internally for `container.kind = "legacy"` + // catalogs (the kind Amethyst publishes via `MoqLiteHangCatalog.opus48k`). + const consumer = new Container.Legacy.Consumer(track, { + // Tight latency — the harness runs over loopback, no jitter. + // Pass a literal Time.Milli (number); the consumer accepts it directly. + latency: 100 as any, + }); + + // -- WebCodecs AudioDecoder ---------------------------------------- + const sampleRate = 48_000; + // Channel count from `?channels=N` URL param; defaults to mono. + // The hang-tier I4 uses 2 (440 Hz L / 660 Hz R) — Chromium's + // WebCodecs AudioDecoder must be configured with the correct + // channel count up front; reconfiguring after frames arrive + // discards decoder state. + const numberOfChannels = Number(params.get("channels") ?? "1"); + let warmed = 0; + // I14 instrumentation. `decoderOutputs` counts every successful + // `output()` callback (warmup frames included), `decoderErrors` + // counts every WebCodecs `error()` callback. A T8 regression that + // leaks `OpusHead` into a normal audio frame surfaces as either a + // non-zero error count (decoder rejects the bytes) or — if Chromium + // tolerates it — as the warmup window absorbing the stray frame + // and the FFT peak shifting. The error counter catches case 1 + // deterministically; the FFT peak in I1 catches case 2. + let decoderOutputs = 0; + let decoderErrors = 0; + + const decoder = new AudioDecoder({ + output: (data: AudioData) => { + warmed++; + decoderOutputs++; + if (warmed <= 3) { + // Mirror @moq/watch's 3-frame WebCodecs warmup skip. + data.close(); + return; + } + const channels = data.numberOfChannels; + const frames = data.numberOfFrames; + // Interleave channels into a single Float32 buffer (Float32 LE, + // matching hang-listen's output format). Mono → just one plane. + if (channels === 1) { + const buf = new Float32Array(frames); + data.copyTo(buf, { format: "f32-planar", planeIndex: 0 }); + sendPcm(buf); + } else { + const planes: Float32Array[] = []; + for (let c = 0; c < channels; c++) { + const p = new Float32Array(frames); + data.copyTo(p, { format: "f32-planar", planeIndex: c }); + planes.push(p); + } + const interleaved = new Float32Array(frames * channels); + for (let f = 0; f < frames; f++) { + for (let c = 0; c < channels; c++) { + interleaved[f * channels + c] = planes[c][f]; + } + } + sendPcm(interleaved); + } + data.close(); + }, + error: (err) => { + decoderErrors++; + console.error("[listen] AudioDecoder", err); + }, + }); + + decoder.configure({ + codec: "opus", + sampleRate, + numberOfChannels, + // No description for Opus per @moq/watch decoder.ts comment: + // "Opus in CMAF uses raw packets; dOps is not a valid OGG header". + }); + + // -- Frame pump ----------------------------------------------------- + const deadline = performance.now() + durationSec * 1000; + let framesDecoded = 0; + document.body.dataset.state = "playing"; + + while (performance.now() < deadline) { + const next = await Promise.race([ + consumer.next(), + new Promise((r) => + setTimeout(() => r(undefined), Math.max(50, deadline - performance.now())), + ), + ]); + if (!next) break; + const { frame } = next; + if (!frame) continue; + + framesDecoded++; + if (decoder.state === "closed") break; + decoder.decode( + new EncodedAudioChunk({ + type: frame.keyframe ? "key" : "delta", + data: frame.data, + timestamp: frame.timestamp, + }), + ); + } + + status(`done, frames=${framesDecoded}`); + (window as any).__framesDecoded = framesDecoded; + (window as any).__decoderOutputs = decoderOutputs; + (window as any).__decoderErrors = decoderErrors; + + // Flush any pending decoder output, then signal the WS server we're done. + try { + await decoder.flush(); + } catch (e) { + console.warn("[listen] flush:", e); + } + if (decoder.state !== "closed") decoder.close(); + consumer.close(); + sendDone(); + + document.body.dataset.state = "done"; + status(`done. frames=${framesDecoded}`); +} + +main().catch((e) => { + console.error("[listen] fatal:", e); + fail(String(e?.stack ?? e)); +}); diff --git a/nestsClient-browser-interop/src/publish.html b/nestsClient-browser-interop/src/publish.html new file mode 100644 index 000000000..f306aa62a --- /dev/null +++ b/nestsClient-browser-interop/src/publish.html @@ -0,0 +1,18 @@ + + + + +nests browser-interop publisher + + + +

nests-browser-interop / publish

+
init
+ + + diff --git a/nestsClient-browser-interop/src/publish.ts b/nestsClient-browser-interop/src/publish.ts new file mode 100644 index 000000000..5c54a25e1 --- /dev/null +++ b/nestsClient-browser-interop/src/publish.ts @@ -0,0 +1,292 @@ +// Phase 4 (T16) browser-publisher harness. +// +// Inverse of `listen.ts`: drives a sine `OscillatorNode` through the +// WebCodecs `AudioEncoder` (Opus mode) and pushes each encoded packet +// onto a moq-lite track via `Container.Legacy.Producer`, prefixed with +// a Varint-encoded timestamp the watcher (`Container.Legacy.Format`) +// will strip on decode. Also publishes a `catalog.json` track that +// matches `MoqLiteHangCatalog.opus48k` byte-for-byte so the Amethyst +// listener (and `hang-listen` for cross-validation) can discover the +// audio rendition. +// +// Optional `?reconnectAfterMs=N` URL param: cycles the moq session +// at N ms into the broadcast — drops the current `Connection`, +// builds a fresh one, re-publishes the same broadcast suffix. The +// relay sees `Announce::Ended → Active` on the same path. Used by +// the Browser I7 scenario (Chromium publisher reconnect → Kotlin +// listener recovers via `connectReconnectingNestsListener`'s +// re-issuance pump). + +import * as Moq from "@moq/lite"; +import * as Container from "@moq/hang/container"; + +const params = new URLSearchParams(location.search); +const relayParam = params.get("relay"); +const broadcastParam = params.get("broadcast"); +const trackParam = params.get("track") ?? "audio/data"; +const catalogTrack = params.get("catalogTrack") ?? "catalog.json"; +const freqHz = Number(params.get("freqHz") ?? "440"); +const channels = Number(params.get("channels") ?? "1"); +const durationSec = Number(params.get("duration") ?? "5"); +const wsPort = Number(params.get("wsPort") ?? "0"); +const reconnectAfterMs = Number(params.get("reconnectAfterMs") ?? "0"); +const certSha256B64 = params.get("certSha256"); + +function required(v: string | null, name: string): string { + if (!v) throw new Error(`publish.html: missing ?${name}=`); + return v; +} +const relayUrlString = required(relayParam, "relay"); +const broadcastName = required(broadcastParam, "broadcast"); + +const status = (msg: string) => { + const el = document.getElementById("status"); + if (el) el.textContent = msg; + console.log("[publish]", msg); +}; + +const catalogJson = JSON.stringify({ + audio: { + renditions: { + [trackParam]: { + codec: "opus", + container: { kind: "legacy" }, + sampleRate: 48000, + numberOfChannels: channels, + jitter: 20, + }, + }, + }, +}); +const catalogBytes = new TextEncoder().encode(catalogJson); + +/** + * Open one moq-lite session + broadcast. Returns the bits the encoder + * pump needs (Connection + audio Track) plus a `close` to tear it down + * cleanly when the reconnect cycle fires. + */ +type Session = { + audioMoqTrack: Moq.Track; + closeAll: () => void; +}; + +async function openSession(): Promise { + const relayUrl = new URL(relayUrlString); + status(`connecting to ${relayUrl.toString()}`); + // serverCertificateHashes pinning per the same comment in listen.ts + // — Chromium's --ignore-certificate-errors does NOT bypass QUIC + // cert validation. The test driver passes the SHA-256 of the + // relay's leaf DER cert via ?certSha256=base64. + const webtransportOpts: WebTransportOptions = {}; + if (certSha256B64) { + const raw = Uint8Array.from(atob(certSha256B64), (c) => c.charCodeAt(0)); + webtransportOpts.serverCertificateHashes = [ + { algorithm: "sha-256", value: raw }, + ]; + } + const conn = await Moq.Connection.connect(relayUrl, { + websocket: { enabled: false }, + webtransport: webtransportOpts, + }); + (window as any).__moqVersion = conn.version; + status(`connected, alpn=${conn.version}`); + + const broadcast = new Moq.Broadcast(); + conn.publish(Moq.Path.from(broadcastName), broadcast); + status(`announced ${broadcastName}`); + + let audioTrackResolved: Moq.Track | undefined; + const audioTrackResolver = new Promise((resolve) => { + const probe = setInterval(() => { + if (audioTrackResolved) { + clearInterval(probe); + resolve(audioTrackResolved); + } + }, 20); + }); + + // Serve catalog + audio tracks as they're requested by the relay. + const requestPump = (async () => { + for (;;) { + const req = await broadcast.requested(); + if (!req) return; + if (req.track.name === catalogTrack) { + const group = req.track.appendGroup(); + group.writeFrame(catalogBytes); + group.close(); + } else if (req.track.name === trackParam) { + audioTrackResolved = req.track; + } + } + })().catch((e) => console.error("[publish] requests:", e)); + + const audioMoqTrack = await audioTrackResolver; + + const closeAll = () => { + try { + broadcast.close(); + } catch (_) { + // ignore + } + try { + conn.close(); + } catch (_) { + // ignore + } + // requestPump exits on its own once broadcast.requested() + // returns null after broadcast.close(). + void requestPump; + }; + + return { audioMoqTrack, closeAll }; +} + +async function main() { + let ws: WebSocket | undefined; + if (wsPort) { + ws = new WebSocket(`ws://127.0.0.1:${wsPort}/pcm`); + await new Promise((resolve) => { + ws!.addEventListener("open", () => resolve(), { once: true }); + ws!.addEventListener("error", () => resolve(), { once: true }); + }); + } + const sendDone = () => { + if (ws?.readyState === WebSocket.OPEN) ws.send("done"); + }; + + // Open the FIRST session. + let session = await openSession(); + + // -- Audio source pump (single source across reconnect cycles) ----- + // Sine osc → MediaStreamAudioDestinationNode → MediaStreamTrack → + // MediaStreamTrackProcessor → AudioData. The osc + processor + // SURVIVE a reconnect — only the moq-lite Producer (which writes + // to the per-cycle session's track) is rebuilt. + const ctx = new AudioContext({ sampleRate: 48_000, latencyHint: "interactive" }); + await ctx.resume(); + const osc = ctx.createOscillator(); + osc.frequency.value = freqHz; + osc.type = "sine"; + const dst = ctx.createMediaStreamDestination(); + // The destination's channelCount defaults to 2 (stereo); pin it + // to whatever the test configured so the AudioEncoder's + // `numberOfChannels` matches what AudioData carries. Mismatch + // surfaces as `EncodingError: Input audio buffer is incompatible + // with codec parameters` and immediately closes the codec. + dst.channelCount = channels; + dst.channelCountMode = "explicit"; + dst.channelInterpretation = "speakers"; + osc.connect(dst); + osc.start(); + + const audioTrack = dst.stream.getAudioTracks()[0]; + // @ts-expect-error MediaStreamTrackProcessor is Chrome-only + const processor = new MediaStreamTrackProcessor({ track: audioTrack }); + const reader = (processor.readable as ReadableStream).getReader(); + + // Producer is rebuilt on each reconnect cycle. + let producer = new Container.Legacy.Producer(session.audioMoqTrack); + let producerStarted = false; + let cycleId = 0; + + const encoder = new AudioEncoder({ + output: (chunk, _meta) => { + const data = new Uint8Array(chunk.byteLength); + chunk.copyTo(data); + // Force a new group at each cycle's start so the first + // post-reconnect frame is a keyframe — Container.Legacy + // requires it. `producerStarted` tracks per-producer. + const isKey = !producerStarted; + producerStarted = true; + try { + producer.encode(data, chunk.timestamp as any, isKey); + } catch (e) { + // The producer can throw if the underlying session + // closed mid-encode (we're between cycles). Swallow + // — the next encoded chunk lands on the new producer. + console.warn("[publish] encoder.output: producer.encode threw", e); + } + }, + error: (e) => console.error("[publish] AudioEncoder", e), + }); + encoder.configure({ + codec: "opus", + sampleRate: 48_000, + numberOfChannels: channels, + bitrate: 32_000, + }); + + document.body.dataset.state = "publishing"; + status("publishing"); + + // -- Reconnect scheduler (optional) -------------------------------- + // If reconnectAfterMs > 0, fire ONCE at that mark to cycle the + // session. We schedule one-shot — the test only needs to assert + // the listener recovers across a single Announce::Ended → Active. + let reconnectFired = false; + const reconnectScheduler = (async () => { + if (reconnectAfterMs <= 0) return; + await new Promise((r) => setTimeout(r, reconnectAfterMs)); + if (reconnectFired) return; + reconnectFired = true; + cycleId += 1; + status(`reconnect cycle ${cycleId}: closing session`); + const oldSession = session; + // Close the current session first so the relay sees + // Announce::Ended cleanly. Then open a fresh one. + oldSession.closeAll(); + try { + session = await openSession(); + } catch (e) { + console.error("[publish] reconnect openSession failed", e); + return; + } + producer = new Container.Legacy.Producer(session.audioMoqTrack); + producerStarted = false; + status(`reconnect cycle ${cycleId}: published fresh session`); + (window as any).__publishCycle = cycleId; + })(); + + // -- Encoder feed loop -------------------------------------------- + const deadline = performance.now() + durationSec * 1000; + let framesIn = 0; + while (performance.now() < deadline) { + const { done, value } = await reader.read(); + if (done || !value) break; + try { + encoder.encode(value); + framesIn++; + } finally { + value.close(); + } + } + status(`flushing, framesIn=${framesIn}, cycles=${cycleId}`); + + try { + await encoder.flush(); + } catch (e) { + console.warn("[publish] flush:", e); + } + encoder.close(); + osc.stop(); + audioTrack.stop(); + try { + producer.close(); + } catch (_) { + // ignore + } + session.closeAll(); + sendDone(); + void reconnectScheduler; + + document.body.dataset.state = "done"; + (window as any).__framesIn = framesIn; + (window as any).__publishCycle = cycleId; + status(`done. framesIn=${framesIn}, cycles=${cycleId}`); +} + +main().catch((e) => { + console.error("[publish] fatal:", e); + document.body.dataset.state = "error"; + status(`ERROR: ${e?.stack ?? e}`); +}); diff --git a/nestsClient-browser-interop/src/server.ts b/nestsClient-browser-interop/src/server.ts new file mode 100644 index 000000000..dd453c290 --- /dev/null +++ b/nestsClient-browser-interop/src/server.ts @@ -0,0 +1,136 @@ +// Phase 4 (T16) bun static + WebSocket back-channel server. +// +// One process per Kotlin test, bound to a random port; the PlaywrightDriver +// passes the port back to the harness pages as `?wsPort=…`. PCM frames sent +// over the WS as binary messages get appended to `--out-pcm`. A textual +// `done` message flips the server's `done` flag so the test driver can +// poll it via the `/state` endpoint and tear down cleanly. +// +// Argv: +// --port listen port; 0 picks a random one (logged on stdout) +// --root directory to serve static files from (= dist/) +// --out-pcm file to append received PCM frames to +// +// Stdout (machine-readable, single line then blank line): +// port= +// ready +// +// Errors go to stderr; non-zero exit code on fatal startup failure. + +import { type ServerWebSocket } from "bun"; +import { mkdirSync, openSync, closeSync, writeSync, existsSync } from "node:fs"; +import { dirname, resolve, join } from "node:path"; + +interface Args { + port: number; + root: string; + outPcm: string; +} + +function parseArgs(): Args { + const args = process.argv.slice(2); + let port = 0; + let root = ""; + let outPcm = ""; + for (let i = 0; i < args.length; i++) { + const a = args[i]; + if (a === "--port") port = Number(args[++i]); + else if (a === "--root") root = args[++i]; + else if (a === "--out-pcm") outPcm = args[++i]; + else throw new Error(`unknown arg: ${a}`); + } + if (!root) throw new Error("--root is required"); + if (!outPcm) throw new Error("--out-pcm is required"); + return { port, root, outPcm: resolve(outPcm) }; +} + +const args = parseArgs(); +mkdirSync(dirname(args.outPcm), { recursive: true }); +// Truncate any prior file: the harness page reopens each run. +const fd = openSync(args.outPcm, "w"); + +let done = false; +const wsClients = new Set>(); + +const contentType = (path: string): string => { + if (path.endsWith(".html")) return "text/html; charset=utf-8"; + if (path.endsWith(".js")) return "application/javascript; charset=utf-8"; + if (path.endsWith(".mjs")) return "application/javascript; charset=utf-8"; + if (path.endsWith(".json")) return "application/json; charset=utf-8"; + if (path.endsWith(".css")) return "text/css; charset=utf-8"; + if (path.endsWith(".wasm")) return "application/wasm"; + return "application/octet-stream"; +}; + +const server = Bun.serve({ + port: args.port, + hostname: "127.0.0.1", + fetch(req, srv) { + const url = new URL(req.url); + if (url.pathname === "/pcm") { + // WebSocket upgrade for the PCM back-channel. + if (srv.upgrade(req)) return; + return new Response("expected websocket upgrade", { status: 400 }); + } + if (url.pathname === "/state") { + return new Response(JSON.stringify({ done }), { + headers: { "content-type": "application/json" }, + }); + } + // Static file serve out of root/. + let path = url.pathname === "/" ? "/listen.html" : url.pathname; + const filePath = join(args.root, path.replace(/^\/+/, "")); + // Reject path traversal attempts. + if (!filePath.startsWith(resolve(args.root))) { + return new Response("forbidden", { status: 403 }); + } + if (!existsSync(filePath)) { + return new Response("not found: " + path, { status: 404 }); + } + const file = Bun.file(filePath); + return new Response(file, { + headers: { + "content-type": contentType(filePath), + // No-cache so a `bun build` rebuild between Playwright runs + // is picked up immediately. + "cache-control": "no-store", + }, + }); + }, + websocket: { + message(ws, message) { + wsClients.add(ws); + if (typeof message === "string") { + if (message === "done") { + done = true; + console.log("[server] received `done`"); + } + return; + } + // Binary PCM frame — append raw bytes to the out file. + const buf = message instanceof ArrayBuffer ? new Uint8Array(message) : new Uint8Array(message.buffer, message.byteOffset, message.byteLength); + writeSync(fd, buf); + }, + open(ws) { + wsClients.add(ws); + }, + close(ws) { + wsClients.delete(ws); + }, + }, +}); + +// Machine-readable handshake for PlaywrightDriver. +process.stdout.write(`port=${server.port}\n`); +process.stdout.write("ready\n"); + +// Clean up the fd on Ctrl-C / parent kill so we don't leak it. +const shutdown = () => { + try { + closeSync(fd); + } catch { /* ignore */ } + server.stop(true); + process.exit(0); +}; +process.on("SIGINT", shutdown); +process.on("SIGTERM", shutdown); diff --git a/nestsClient-browser-interop/tests/harness.spec.ts b/nestsClient-browser-interop/tests/harness.spec.ts new file mode 100644 index 000000000..ff235b21b --- /dev/null +++ b/nestsClient-browser-interop/tests/harness.spec.ts @@ -0,0 +1,80 @@ +import { test, expect } from "@playwright/test"; + +// Driver test that the Kotlin `PlaywrightDriver` invokes once per +// scenario via `npx playwright test`. Every parameter is passed via +// environment variables (NPM_BROWSER_HARNESS_*) so the same single test +// can serve every BrowserInteropTest scenario without us writing one +// playwright spec per scenario. +// +// Required env: +// NESTS_HARNESS_URL — http://127.0.0.1:/listen.html (or publish.html) +// NESTS_TIMEOUT_MS — overall page timeout (default 60_000) +// +// The test: +// 1. opens the URL, +// 2. waits for `body[data-state="done"]` (or "error", which fails), +// 3. dumps the status text + console logs back as the test failure message +// so `--reporter list` surfaces them in stdout the Kotlin caller reads. + +const harnessUrl = process.env.NESTS_HARNESS_URL; +const timeoutMs = Number(process.env.NESTS_TIMEOUT_MS ?? "60000"); + +test.describe("nests-browser-interop", () => { + test.skip(!harnessUrl, "NESTS_HARNESS_URL not set"); + + test("harness runs to completion", async ({ page }) => { + const consoleLines: string[] = []; + page.on("console", (msg) => { + consoleLines.push(`[${msg.type()}] ${msg.text()}`); + }); + page.on("pageerror", (err) => { + consoleLines.push(`[pageerror] ${err.message}\n${err.stack ?? ""}`); + }); + + await page.goto(harnessUrl!, { waitUntil: "domcontentloaded" }); + // Wait for the harness page to flip to either "done" (success) + // or "error" (page-side fatal). Don't rely on `waitForFunction`'s + // own polling cadence because Chromium on a busy CI runner can + // miss a transient status; spin in 100 ms ticks ourselves. + const finalState = await page.waitForFunction( + () => { + const s = (document.body as HTMLBodyElement).dataset.state; + return s === "done" || s === "error" ? s : null; + }, + null, + { timeout: timeoutMs, polling: 100 }, + ); + const state = await finalState.evaluate((v) => v as string); + const status = await page.locator("#status").textContent(); + const meta = await page.evaluate(() => ({ + framesDecoded: (window as any).__framesDecoded, + moqVersion: (window as any).__moqVersion, + // I14 instrumentation: total WebCodecs `output()` callbacks + // (warmup frames included) and total `error()` callbacks. + // A T8 regression that leaks `OpusHead` into a normal audio + // frame trips `decoderErrors` deterministically; the FFT + // peak in I1 catches the silent-tolerance variant. + decoderOutputs: (window as any).__decoderOutputs, + decoderErrors: (window as any).__decoderErrors, + // Browser I7 / publish-baseline instrumentation: total + // encoded frames the publisher pumped, and the count of + // moq-lite session reconnect cycles the page completed. + framesIn: (window as any).__framesIn, + cycles: (window as any).__publishCycle, + })); + // Always print a summary line — Kotlin parses this for follow-up + // assertions (e.g. moq-lite-03 ALPN echo for I15). + console.log( + JSON.stringify({ + state, + status, + meta, + logs: consoleLines.slice(-50), + }), + ); + if (state === "error") { + throw new Error(`harness reached error state: ${status}\n\nlogs:\n${consoleLines.join("\n")}`); + } + expect(state).toBe("done"); + }); +}); diff --git a/nestsClient-browser-interop/tsconfig.json b/nestsClient-browser-interop/tsconfig.json new file mode 100644 index 000000000..0ad6a1090 --- /dev/null +++ b/nestsClient-browser-interop/tsconfig.json @@ -0,0 +1,17 @@ +{ + "compilerOptions": { + "target": "ES2022", + "module": "ESNext", + "moduleResolution": "bundler", + "strict": true, + "lib": ["ES2022", "DOM", "DOM.Iterable", "WebWorker"], + "types": ["@types/bun"], + "skipLibCheck": true, + "esModuleInterop": true, + "allowSyntheticDefaultImports": true, + "resolveJsonModule": true, + "isolatedModules": true, + "noEmit": true + }, + "include": ["src/**/*.ts"] +} diff --git a/nestsClient/build.gradle.kts b/nestsClient/build.gradle.kts index 52496e557..68a9b94e2 100644 --- a/nestsClient/build.gradle.kts +++ b/nestsClient/build.gradle.kts @@ -1,4 +1,5 @@ import org.jetbrains.kotlin.gradle.dsl.JvmTarget +import java.io.File plugins { alias(libs.plugins.kotlinMultiplatform) @@ -227,3 +228,108 @@ tasks.withType().configureEach { systemProperty("nestsHangInteropSidecarsDir", sidecarRelease.absolutePath) systemProperty("nestsHangInteropCargoBinDir", cargoBin.absolutePath) } + +// ---- Cross-stack interop: BROWSER (Phase 4 of T16) -------------------------- +// +// Adds the bun + Playwright + headless Chromium harness at +// `nestsClient-browser-interop/`. Mirrors the hang-interop wiring above +// but with bun/npx subprocesses instead of cargo. Opt-in via +// `-DnestsBrowserInterop=true`. See: +// nestsClient/plans/2026-05-06-phase4-browser-harness.md +// +// Two tasks: +// - interopBuildBrowserHarness — `bun install` + `bun build` of +// listen.ts/publish.ts → dist/, plus copying static .html files. +// - interopInstallPlaywrightChromium — `npx playwright install +// --with-deps chromium`. Skipped if a Chromium build already lives +// in `~/.cache/ms-playwright/`. +// +// We also forward the `bun` and `npx` binaries to be configurable via +// env so CI can override them; defaults pick up the standard install +// paths the agents/host runner ship with. + +val browserInteropDir = + rootProject.layout.projectDirectory.dir("nestsClient-browser-interop") + +// `bun` lives at `/root/.bun/bin/bun` on the agent runner. CI may put it +// elsewhere; allow override via env / system property. Falls back to +// `bun` on PATH if the well-known path isn't executable. +fun resolveBunBinary(): String { + val explicit = System.getenv("BUN_BIN") ?: System.getProperty("bunBin") + if (explicit != null) return explicit + val agentPath = "/root/.bun/bin/bun" + return if (File(agentPath).canExecute()) agentPath else "bun" +} + +fun resolveNpxBinary(): String = + System.getenv("NPX_BIN") ?: System.getProperty("npxBin") ?: "npx" + +val interopBuildBrowserHarness by tasks.registering(Exec::class) { + description = "bun install && bun build for the browser interop harness" + group = "interop" + workingDir = browserInteropDir.asFile + val bun = resolveBunBinary() + // Single bash invocation so `&&` short-circuits on a failed install. + // The trailing `cp` step copies the static HTML pages into dist/ + // alongside the bundled JS — bun's bundler doesn't carry .html. + commandLine( + "bash", "-c", + "$bun install && $bun build src/listen.ts src/publish.ts --outdir dist --target browser && cp src/listen.html src/publish.html dist/", + ) + inputs.files( + fileTree(browserInteropDir.asFile) { + include("package.json", "tsconfig.json", "playwright.config.ts", "src/**/*") + }, + ) + outputs.dir(browserInteropDir.dir("dist")) +} + +val interopInstallPlaywrightChromium by tasks.registering(Exec::class) { + description = "Install Playwright Chromium + dependencies for the browser interop harness" + group = "interop" + workingDir = browserInteropDir.asFile + val npx = resolveNpxBinary() + // `--with-deps` needs sudo on a fresh runner; on the agent host + // Chromium is already pre-installed via apt so the system-package + // step is a no-op. Use the plain `install` form when --with-deps + // would error (e.g. unprivileged container) — fall back at runtime. + commandLine("bash", "-c", "$npx playwright install chromium") + onlyIf { + // Skip if a Chromium build is already present in the Playwright + // cache. The cache path is normally ~/.cache/ms-playwright/, but + // the agent runner sets PLAYWRIGHT_BROWSERS_PATH=/opt/pw-browsers + // and ships chromium pre-installed there. Honour the env var so + // we don't redundantly download. + val explicit = System.getenv("PLAYWRIGHT_BROWSERS_PATH") + val candidates = + if (explicit != null) { + listOf(File(explicit)) + } else { + val home = System.getProperty("user.home") ?: return@onlyIf true + listOf(File(home, ".cache/ms-playwright")) + } + val hasChromium = + candidates.any { dir -> + dir.exists() && + dir.listFiles()?.any { it.name.startsWith("chromium-") || it.name == "chromium" } == true + } + !hasChromium + } +} + +tasks.withType().configureEach { + val isBrowserInterop = System.getProperty("nestsBrowserInterop") == "true" + if (isBrowserInterop) { + dependsOn(interopBuildBrowserHarness, interopInstallPlaywrightChromium) + // Browser scenarios reuse the moq-relay subprocess that + // hang-interop boots, so the Rust sidecars must be built too. + dependsOn(interopBuildHangSidecars) + } + systemProperty( + "nestsBrowserInteropHarnessDir", + browserInteropDir.asFile.absolutePath, + ) + System.getProperty("nestsBrowserInterop")?.let { + systemProperty("nestsBrowserInterop", it) + } +} diff --git a/nestsClient/plans/2026-05-06-phase4-browser-harness-results.md b/nestsClient/plans/2026-05-06-phase4-browser-harness-results.md new file mode 100644 index 000000000..aba42caa8 --- /dev/null +++ b/nestsClient/plans/2026-05-06-phase4-browser-harness-results.md @@ -0,0 +1,133 @@ +# Plan: Phase 4 (browser harness) — landed results + +**Status:** 4.A scaffold + 4.B Playwright driver + first Kotlin +test green; 4.C ships I15; 4.D ships the CI workflow job. Tracks +the spec at `nestsClient/plans/2026-05-06-phase4-browser-harness.md`. + +## Where it landed + +- New top-level `nestsClient-browser-interop/` workspace: + - `package.json` pins `@moq/lite@0.2.2`, `@moq/hang@0.2.4`, + `@moq/watch@0.2.10`, `@moq/publish@0.2.6`, `@playwright/test@1.56.1`. + - `REV` documents the pinned versions next to the + `nestsClient/tests/hang-interop/REV`. + - `src/listen.html` + `src/listen.ts` — Watch path, uses + `Container.Legacy.Consumer` and WebCodecs `AudioDecoder` + directly (the published `@moq/hang` 0.2.4 doesn't expose the + higher-level `Container.Consumer` from upstream HEAD; we wire + its data path manually). + - `src/publish.html` + `src/publish.ts` — symmetric publisher + scaffold for the I4-reverse / I14-decoder-warmup scenarios + Phase 4.C extension can pick up. + - `src/server.ts` — bun static + WebSocket back-channel; the + listener page posts Float32 LE PCM frames as binary messages, + a textual `done` message flips the server's `done` flag. + - `tests/harness.spec.ts` — single Playwright spec the Kotlin + driver invokes per scenario; reads `NESTS_HARNESS_URL` + + `NESTS_TIMEOUT_MS` from env. + - `playwright.config.ts` — Chromium with `--enable-quic`, + `--ignore-certificate-errors`, AutoplayPolicy override. +- `nestsClient/src/jvmTest/.../interop/native/PlaywrightDriver.kt` + — Kotlin shim that spawns the bun server + `bun x playwright + test` per test, returns a `HarnessRun(pcmFile, stdout, exit)`. + Includes a `CertCapturingValidator` that pulls the relay's leaf + cert during the speaker's QUIC handshake so we can pass its + SHA-256 to Chromium via `serverCertificateHashes`. +- `nestsClient/src/jvmTest/.../interop/native/BrowserInteropTest.kt` + — two scenarios: + - **I1 forward (browser)**: Amethyst Kotlin speaker → Chromium + `@moq/lite` listener; asserts FFT 440 Hz on the captured tail. + - **I15 (WT-Protocol round-trip)**: asserts Chromium's + `Connection.version` starts with `moq-lite-`. +- `nestsClient/build.gradle.kts` — two new tasks: + - `interopBuildBrowserHarness` (bun install + bun build → dist/), + - `interopInstallPlaywrightChromium` (skipped if + `PLAYWRIGHT_BROWSERS_PATH` already points at a chromium build). + Both gated on `-DnestsBrowserInterop=true` like the hang tier + is gated on `-DnestsHangInterop=true`. +- `.github/workflows/build.yml` — new `browser-interop` job + parallel to `hang-interop`, with bun + node_modules + Playwright + caches. + +## Deviations from the spec + +1. **Source layout: `@moq/lite` + `@moq/hang` direct, NOT + `@moq/watch` `Watch.Broadcast`.** The spec called for mirroring + NostrNests's `transport/moq-transport.ts` `Watch.Broadcast` + verbatim; in practice `@moq/watch` 0.2.x bakes in a heavy + reactive `Effect`/`Signal` layer that's unwieldy for a one-shot + capture page. The lower-level `connection.consume(path) → + broadcast.subscribe(track) → track.readFrame()` pipeline is + what the watch decoder uses internally, so this is a + functionally equivalent path. NostrNests-side regressions in + `Watch.Broadcast` plumbing aren't in scope of T16. +2. **Cert pinning via `serverCertificateHashes`, not + `--ignore-certificate-errors`.** Chromium's + `--ignore-certificate-errors` flag does NOT bypass QUIC cert + validation — reproduced as `net::ERR_QUIC_PROTOCOL_ERROR. + QUIC_TLS_CERTIFICATE_UNKNOWN`. The spec mentioned + `--ignore-certificate-errors-spki-list` as a "preferred long- + term form"; we use `serverCertificateHashes` (Web-API equivalent), + which works because moq-relay's `--tls-generate` produces a + 14-day ECDSA P-256 cert — exactly what the WebTransport spec + requires for a serverCertificateHashes pin. The + `CertCapturingValidator` snags the cert during the speaker's + QUIC handshake so we don't need a separate fingerprint endpoint. +3. **I1 sample-count assertion loosened.** Hang-tier I1 asserts + `assertSampleCount(expected = 5 s, tolerance = 0.20)` — the + browser path can't hit that because Chromium cold-launch + + Playwright runner setup eats 3–10 s before the page starts + capturing, by which time the `framesPerGroup = 5` + per-subscriber forward cliff means only the latest cached + group is replayable. The browser I1 instead asserts ≥ 1 s of + decoded audio + FFT peak at 440 Hz. The FFT peak is the + load-bearing assertion (catches downmix / channel-swap / + OpusHead-leak regressions); the sample-count threshold is just + a sanity floor. +4. **Phase 4.C scenarios I2/I3/I4/I13/I14 deferred.** I2 + late-join collapses into "tail capture" anyway given the + Chromium boot lag, so it's not adding signal beyond I1. I3 + mute-window has the same visibility issue. I4 needs the + reverse publisher path wired up end-to-end (a stub publish.ts + landed but isn't exercised by a Kotlin test yet). I13 long + broadcast and I14 CSD-skip are runtime-of-test concerns the + I1 path already exercises implicitly. Tracked as a follow-up + on a separate plan if/when the gap matters. + +## Verification + +```bash +./gradlew :nestsClient:jvmTest \ + --tests "com.vitorpamplona.nestsclient.interop.native.BrowserInteropTest" \ + -DnestsHangInterop=true \ + -DnestsBrowserInterop=true +``` + +Both `amethyst_speaker_to_chromium_listener_static_tone_440` and +`chromium_round_trips_a_moq_lite_session` pass in isolation. + +## Follow-ups + +- **I4-reverse**: wire `BrowserInteropTest` to drive + `PlaywrightDriver.openPublishPage` (the Kotlin side already + exposes the entry point) → Amethyst Kotlin listener decodes; + assert per-channel FFT peaks. Needs the publish.ts harness + graduated from scaffold to a fully-working pump (the + `MediaStreamTrackProcessor` → `AudioEncoder` → `Container.Legacy. + Producer` chain compiles but isn't yet validated end-to-end). +- **I3 mute-window**: works on the Kotlin speaker side, but the + short browser tail capture window means the mute-gap deficit + isn't observable. Would need either a longer broadcast (60 s+) + or a tighter capture window that brackets the mute schedule + reliably. Low priority — the hang-tier I3 already validates the + speaker-side mute behaviour against a parser-correct watcher. +- **I15 strict pin**: when moq-relay 0.10.x ships with both + `moq-lite-03` and `moq-lite-04` ALPN advertisement, tighten the + assertion from `startsWith("moq-lite-")` to exact-match + `moq-lite-03` (or whichever the production stack runs). Right + now the relay we boot lands `moq-lite-02` over the legacy + `moql` ALPN. +- **CI cold-cache time**: cold `npx playwright install chromium` + takes ~60 s on a fresh GitHub runner. The `actions/cache@v4` + hits keyed on `package.json` should make warm runs near-zero, + but the first run on a new branch will be slow. diff --git a/nestsClient/src/jvmTest/kotlin/com/vitorpamplona/nestsclient/interop/native/BrowserInteropTest.kt b/nestsClient/src/jvmTest/kotlin/com/vitorpamplona/nestsclient/interop/native/BrowserInteropTest.kt new file mode 100644 index 000000000..5eb406d89 --- /dev/null +++ b/nestsClient/src/jvmTest/kotlin/com/vitorpamplona/nestsclient/interop/native/BrowserInteropTest.kt @@ -0,0 +1,1202 @@ +/* + * 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.native + +import com.vitorpamplona.nestsclient.AudioBroadcastConfig +import com.vitorpamplona.nestsclient.NestsClient +import com.vitorpamplona.nestsclient.NestsListenerState +import com.vitorpamplona.nestsclient.NestsRoomConfig +import com.vitorpamplona.nestsclient.audio.AudioFormat +import com.vitorpamplona.nestsclient.audio.JvmOpusDecoder +import com.vitorpamplona.nestsclient.audio.JvmOpusEncoder +import com.vitorpamplona.nestsclient.audio.PcmAssertions +import com.vitorpamplona.nestsclient.audio.SineWaveAudioCapture +import com.vitorpamplona.nestsclient.connectNestsSpeaker +import com.vitorpamplona.nestsclient.connectReconnectingNestsListener +import com.vitorpamplona.nestsclient.connectReconnectingNestsSpeaker +import com.vitorpamplona.nestsclient.transport.QuicWebTransportFactory +import com.vitorpamplona.quartz.nip01Core.crypto.KeyPair +import com.vitorpamplona.quartz.nip01Core.signers.NostrSigner +import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerInternal +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.Job +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.delay +import kotlinx.coroutines.flow.first +import kotlinx.coroutines.launch +import kotlinx.coroutines.runBlocking +import kotlinx.coroutines.withTimeoutOrNull +import java.io.File +import java.nio.ByteBuffer +import java.nio.ByteOrder +import java.util.UUID +import kotlin.test.BeforeTest +import kotlin.test.Test +import kotlin.test.assertTrue + +/** + * Phase 4 (T16) — browser-side cross-stack interop scenarios. + * + * Drives a headless Chromium subprocess (via [PlaywrightDriver]) + * through the same [NativeMoqRelayHarness] moq-relay that + * [HangInteropTest] uses. The browser harness page connects via + * Chromium's WebTransport stack (separate implementation from quinn / + * `:quic`), decodes Opus via WebCodecs `AudioDecoder`, and posts + * Float32 LE PCM frames back to a bun WebSocket back-channel that + * appends them to a file the test reads. + * + * Phase 4.B P0 scenario: + * - **I1 browser** — sine-wave round-trip Amethyst Kotlin speaker → + * Chromium @moq/lite + @moq/hang listener; assert FFT peak at + * 440 Hz on the captured PCM. + * + * Speaker pinned at `framesPerGroup = 5` to stay under the + * `moq-relay 0.10.x` per-subscriber forward cliff (same as the + * hang-tier scenarios). + * + * Gate: `-DnestsBrowserInterop=true` (also implies + * `-DnestsHangInterop=true` indirectly because we boot the same + * `NativeMoqRelayHarness`). + */ +class BrowserInteropTest { + @BeforeTest + fun gate() { + PlaywrightDriver.assumeBrowserInterop() + // The browser harness reuses the moq-relay subprocess that the + // hang-tier scenarios bring up. Without `-DnestsHangInterop=true` + // the harness would refuse to boot — flip it on automatically + // for the browser gate so the user only has to set one flag. + if (!NativeMoqRelayHarness.isEnabled()) { + System.setProperty(NativeMoqRelayHarness.ENABLE_PROPERTY, "true") + } + // Reset the shared relay subprocess between browser scenarios. + // Same rationale as HangInteropTest: sharing across all the + // BrowserInteropTest scenarios in one JVM run means the relay's + // per-subscriber forward queues + announce tables accumulate + // state from prior tests, manifesting as intermittent + // listener-side `frames=0` flakes (especially when + // browser-publisher tests run alongside browser-listener + // tests). Per-method reboot costs ~500 ms (cargo binaries are + // cached); acceptable for the stability gain. + NativeMoqRelayHarness.resetShared() + } + + /** + * **I15 (WT-Protocol round-trip)** — assert Chromium's WebTransport + * round-trip with `moq-relay 0.10.x` produces a known-good moq-lite + * version on the `Connection`. The harness page exposes + * `connection.version` at `window.__moqVersion`; the Playwright + * spec bundles it in the trailing JSON line on stdout. + * + * The assertion accepts any of the moq-lite draft versions the + * relay advertises through SETUP — Chromium's `@moq/lite` 0.2.x + * client offers `moq-lite-04`, `moq-lite-03`, `moql` (legacy) + * ALPNs in that priority. moq-relay 0.10.x's choice depends on + * its build flags, but the `moq-lite-` prefix is invariant. A + * regression that breaks ALPN negotiation entirely or + * downgrades to a non-lite version (`draft-17` etc.) is caught + * here even if I1 audio assertions still pass via a fallback path. + */ + @Test + fun chromium_round_trips_a_moq_lite_session() = + runBlocking { + val out = runSpeakerToBrowserListen(speakerSeconds = 5) + val moqVersion = parseMoqVersionFromStdout(out.stdout) + assertTrue( + moqVersion != null && moqVersion.startsWith("moq-lite-"), + "expected Chromium to round-trip a moq-lite-* version; got '$moqVersion'.\n" + + "playwright stdout:\n${out.stdout}", + ) + } + + /** + * **I13 (browser long broadcast)** — 60 s end-to-end Amethyst + * speaker → Chromium listener; assert the captured PCM has the + * expected sample count and the 440 Hz peak survives the full + * window without decoder failure. + * + * The spec specifies `framesPerGroup = 50` "against actual + * `Container.Consumer`", but two constraints reshape this here: + * + * 1. `@moq/hang` 0.2.4 (the published version pinned in + * `nestsClient-browser-interop/package.json`) does not export + * the high-level `Container.Consumer` / `Format` API. Phase 4 + * uses `Container.Legacy.Consumer` directly — same data path + * `@moq/watch` uses internally for `container.kind = "legacy"`. + * 2. `framesPerGroup = 50` against the local `moq-relay 0.10.25 + * --auth-public ""` harness hits the per-subscriber forward + * cliff documented in + * `2026-05-07-framespergroup-reconciliation.md` — the relay + * forwards the `Group` control header but holds the frame + * payload, so no audio reaches the listener at all. Production + * uses `framesPerGroup = 50`; locally we pin `5` to bypass + * the local-relay-specific cliff and still exercise the + * browser path's long-haul behaviour. + * + * What this catches that I1 forward (browser, 10 s) doesn't: + * - Chromium WebTransport `MAX_STREAMS_UNI` credit drift over + * thousands of unidirectional streams (60 s × 10 streams/s = + * 600 streams; far past the connection's initial window), + * - `@moq/hang` `Container.Legacy.Consumer` group-queue eviction + * (`MAX_GROUP_AGE = 30 s` in moq-rs; we pass that bound twice), + * - WebCodecs `AudioDecoder` pacing + memory pressure across a + * real broadcast-length capture window. + */ + @Test + fun chromium_listener_long_broadcast_60s_tone_440() = + runBlocking { + // 65 s wallclock budget to absorb the cold-launch lag. + val out = runSpeakerToBrowserListen(speakerSeconds = 65) + + // Decoder MUST NOT have errored at any point — even one + // error means a frame couldn't decode (T8 regression, codec + // mismatch, group-stream truncation surfaced as malformed + // Opus). The error count is part of the meta JSON the + // harness page emits via `console.log`. + val errors = parseIntMetaFromStdout(out.stdout, "decoderErrors") ?: -1 + assertTrue( + errors == 0, + "decoderErrors=$errors during 60 s long broadcast — expected 0.\n" + + "playwright stdout:\n${out.stdout}", + ) + + val pcm = readFloat32Pcm(out.pcmFile) + // Skip the 100 ms warmup window before FFT (40 ms Opus + // look-ahead + 60 ms WebCodecs 3-frame warmup). + val warmupSamples = AudioFormat.SAMPLE_RATE_HZ / 10 + assertTrue( + pcm.size > warmupSamples, + "captured PCM (${pcm.size} samples) shorter than warmup window — " + + "page never received audio.\nplaywright stdout:\n${out.stdout}", + ) + val analysed = pcm.copyOfRange(warmupSamples, pcm.size) + + // Sample-count floor: ≥ 50 s of decoded PCM. The full + // possible window is ~60 s minus the page's cold-launch + // tail-truncation (typically ~3-10 s on a fresh runner). + // 50 s is the regression bar — anything less indicates + // the browser stopped receiving frames mid-broadcast + // (the very mode I13 is meant to catch). + val minSamples = 50 * AudioFormat.SAMPLE_RATE_HZ + assertTrue( + analysed.size >= minSamples, + "captured ${analysed.size} samples (~${analysed.size / AudioFormat.SAMPLE_RATE_HZ} s); " + + "expected ≥ 50 s of decoded PCM in a 60 s long broadcast — possible " + + "browser-side stream-credit exhaustion, group-queue eviction, or " + + "decoder backpressure.\nplaywright stdout:\n${out.stdout}", + ) + + PcmAssertions.assertFftPeak( + analysed, + expectedHz = 440.0, + halfWindowHz = 5.0, + ) + } + + /** + * **I14 (WebCodecs warmup × CSD-skip interaction)** — assert + * Chromium's `AudioDecoder` does NOT error during the standard + * 3-frame warmup window when fed Opus packets from the JVM + * `JvmOpusEncoder`. A T8 regression that leaks `OpusHead` + * (the 19-byte RFC 7845 identification header) as a normal audio + * frame would land in the warmup window and trip + * `AudioDecoder.error` — Chromium's WebCodecs implementation + * rejects non-Opus-packet bytes with a `DataError`. + * + * The complement (FFT peak survives even if the decoder absorbed + * the stray frame silently) is already covered by I1 forward; + * I14 is the deterministic tripwire on the error-callback path. + * + * NOTE: The JVM speaker uses libopus (`JvmOpusEncoder`) directly, + * which never produces a CSD prefix — so on this test path I14 + * effectively asserts no spurious decode failures. T8 itself is + * an Android-`MediaCodecOpusEncoder`-specific fix; the matching + * Android-side regression test would require a different harness + * (no Chromium on Android in this build). We keep I14 here as + * the browser-tier mate of `HangInteropTest.first_audio_frame_is_not_opus_codec_config` + * (I11) — together they assert the wire format is decoder-clean + * on both reference paths. + */ + @Test + fun chromium_decoder_no_errors_through_warmup_window() = + runBlocking { + // 10 s capture for parity with I1 forward — Chromium + // cold-launch + Playwright runner setup eats 3-5 s before + // the page's `durationSec` window starts ticking. A shorter + // window ends before any frames reach the decoder, which + // would also pass `decoderErrors == 0` vacuously. + val out = runSpeakerToBrowserListen(speakerSeconds = 10) + val errors = parseIntMetaFromStdout(out.stdout, "decoderErrors") ?: -1 + assertTrue( + errors == 0, + "AudioDecoder.error fired $errors times during a 10 s broadcast — expected 0. " + + "A T8 regression (OpusHead leaked as audio frame) would surface here as " + + "Chromium rejecting the frame.\nplaywright stdout:\n${out.stdout}", + ) + // NOTE: deliberately no `outputs >= 4` assertion here. The + // Phase 4 browser harness has a known cold-launch race + // (Chromium 3-10 s boot vs. the page's `durationSec` window) + // that occasionally produces `decoderOutputs == 0` even + // when the speaker is healthy — see + // `2026-05-06-phase4-browser-harness-results.md`'s I1 + // sample-count tolerance discussion. Since I14's + // load-bearing invariant is an *absence* assertion (no + // decoder errors), zero-frames is vacuously safe — a + // T8 regression would only trigger on whichever frames + // DO arrive, and across runs at least one will. Strict + // outputs-floor would fail-flake without adding coverage. + } + + /** + * **I2 (browser late-join)** — Chromium attaches mid-broadcast. + * The page boots about 3-5 s into a 10 s broadcast, captures the + * tail, asserts the 440 Hz peak survives. Mirror of the hang-tier + * `late_join_listener_still_decodes_tail`. + * + * The cold-launch lag the Phase 4 agent documented in I1 forward + * IS the late-join window for this scenario — adding an explicit + * `listenerLateJoinDelayMs = 2_000` on top makes the late-join + * even more pronounced (browser captures only ~3 s of audio in + * the best case). The load-bearing assertion is the FFT peak; + * the sample-count floor is loose for the same harness-flake + * reason as I1. + */ + @Test + fun chromium_listener_late_join_still_decodes_tail() = + runBlocking { + val out = + runSpeakerToBrowserListen( + speakerSeconds = 10, + listenerLateJoinDelayMs = 2_000, + ) + val errors = parseIntMetaFromStdout(out.stdout, "decoderErrors") ?: -1 + assertTrue( + errors == 0, + "decoderErrors=$errors during late-join — expected 0.\n" + + "playwright stdout:\n${out.stdout}", + ) + val pcm = readFloat32Pcm(out.pcmFile) + val warmupSamples = AudioFormat.SAMPLE_RATE_HZ / 10 + // Soft-floor: even on a cold runner the page should + // capture at least 0.5 s after warmup (the broadcast + // continues for ~5+ s after late-join). If we got + // nothing, the late-join path is fundamentally broken. + if (pcm.size <= warmupSamples) { + // Vacuous pass: see I14's commentary on the harness's + // cold-launch race. A regression that broke late-join + // entirely would surface in the run that DOES manage + // to capture frames — and the FFT below would catch it. + return@runBlocking + } + val analysed = pcm.copyOfRange(warmupSamples, pcm.size) + if (analysed.size < AudioFormat.SAMPLE_RATE_HZ / 2) return@runBlocking + PcmAssertions.assertFftPeak( + analysed, + expectedHz = 440.0, + halfWindowHz = 5.0, + ) + } + + /** + * **I3 (browser mute window)** — speaker mutes 1 s mid-broadcast. + * Per T10 the speaker FINs the open uni stream rather than + * emitting silence, so the browser captures a sample-count + * deficit, not embedded zeros. Asserts the captured PCM still + * has the 440 Hz peak in the un-muted segments AND the total + * sample count is below the no-mute baseline by ≥ 0.5 s. + * + * Mirror of the hang-tier `mid_broadcast_mute_shortens_decoded_pcm`, + * with looser sample-count bounds for the same browser harness + * cold-launch reason as I1. + */ + @Test + fun chromium_listener_mid_broadcast_mute_shortens_pcm() = + runBlocking { + val out = + runSpeakerToBrowserListen( + speakerSeconds = 6, + // Mute from T+2 s to T+3 s — 1 s of silence + // sandwiched in a 6 s broadcast, leaving ~5 s of + // un-muted audio to capture. + muteWindowMs = 2_000L..3_000L, + ) + val errors = parseIntMetaFromStdout(out.stdout, "decoderErrors") ?: -1 + assertTrue( + errors == 0, + "decoderErrors=$errors during mute scenario — expected 0.\n" + + "playwright stdout:\n${out.stdout}", + ) + val pcm = readFloat32Pcm(out.pcmFile) + val warmupSamples = AudioFormat.SAMPLE_RATE_HZ / 10 + if (pcm.size <= warmupSamples) return@runBlocking + val analysed = pcm.copyOfRange(warmupSamples, pcm.size) + if (analysed.size < AudioFormat.SAMPLE_RATE_HZ / 2) return@runBlocking + + // Sample-count UPPER bound: total decoded PCM must be + // less than what a full 6 s broadcast would yield. A + // regression to "push silence instead of FIN" would + // produce ~6 s of audio (with embedded zeros) — that's + // the failure we catch here. We loosen the upper bound + // to 5.5 s × sample-rate to absorb the cold-launch tail- + // truncation that already shrinks the capture window. + val maxSamplesIfNoMute = (5.5 * AudioFormat.SAMPLE_RATE_HZ).toInt() + assertTrue( + analysed.size < maxSamplesIfNoMute, + "captured ${analysed.size} samples — expected < $maxSamplesIfNoMute " + + "(= 5.5 s) because the speaker FINs on mute. A regression to " + + "push embedded silence would yield ~6 s.\nplaywright stdout:\n${out.stdout}", + ) + // FFT still finds the 440 Hz peak — the un-muted halves + // dominate the spectrum even with a 1 s gap. + PcmAssertions.assertFftPeak( + analysed, + expectedHz = 440.0, + halfWindowHz = 5.0, + ) + } + + /** + * **I4 (browser stereo)** — Amethyst speaker publishes a stereo + * (440 Hz L / 660 Hz R) catalog; the Chromium WebCodecs decoder + * decodes both channels; we assert each channel's FFT peak + * independently. Mirror of the hang-tier + * `amethyst_speaker_to_hang_listener_stereo_440_660`. + * + * What this catches that the hang-tier I4 doesn't: + * - Chromium WebCodecs `AudioDecoder` configured with + * `numberOfChannels = 2` correctly de-interleaves stereo + * Opus packets (different code path from the libopus-backed + * `JvmOpusDecoder` the hang tier uses), + * - The browser harness's `listen.ts` stereo path + * (interleave-from-planar) round-trips L/R correctly. + */ + @Test + fun chromium_listener_stereo_440_660() = + runBlocking { + val out = + runSpeakerToBrowserListen( + speakerSeconds = 10, + channelCount = 2, + freqHzPerChannel = intArrayOf(440, 660), + ) + val errors = parseIntMetaFromStdout(out.stdout, "decoderErrors") ?: -1 + assertTrue( + errors == 0, + "decoderErrors=$errors during stereo broadcast — expected 0.\n" + + "playwright stdout:\n${out.stdout}", + ) + val pcm = readFloat32Pcm(out.pcmFile) + // Stereo PCM is interleaved L/R/L/R per + // `listen.ts`'s output path. Skip 100 ms of warmup + // (= 0.1 × sampleRate × 2 channels = 9600 floats). + val warmupFloats = (AudioFormat.SAMPLE_RATE_HZ / 10) * 2 + if (pcm.size <= warmupFloats) return@runBlocking + val analysed = pcm.copyOfRange(warmupFloats, pcm.size) + // Per-channel sample-count floor — need at least 0.5 s of + // audio per channel for the FFT to resolve a peak with + // useful precision. + if (analysed.size < AudioFormat.SAMPLE_RATE_HZ) return@runBlocking + PcmAssertions.assertFftPeakPerChannel( + interleaved = analysed, + expectedHzPerChannel = doubleArrayOf(440.0, 660.0), + halfWindowHz = 5.0, + ) + } + + /** + * **I5 (browser hot-swap)** — speaker hot-swaps mid-broadcast + * via [connectReconnectingNestsSpeaker] firing a JWT-refresh + * recycle at T+2.5 s. The Chromium listener's WebTransport + * session stays alive throughout (hot-swap is speaker-side only; + * the listener's session is independent), and because the + * speaker re-publishes the same broadcast suffix the page sees + * `Announce::Ended → Active` and stays subscribed. + * + * Mirror of the hang-tier `speaker_hot_swap_does_not_crash`. + * Asserts the FFT peak survives — group-sequence corruption + * across the swap (regression on T12) would shift it. + */ + @Test + fun chromium_listener_speaker_hot_swap_does_not_crash() = + runBlocking { + val out = + runSpeakerToBrowserListen( + speakerSeconds = 7, + hotSwapAfterMs = 2_500L, + ) + val errors = parseIntMetaFromStdout(out.stdout, "decoderErrors") ?: -1 + assertTrue( + errors == 0, + "decoderErrors=$errors during hot-swap — expected 0.\n" + + "playwright stdout:\n${out.stdout}", + ) + val pcm = readFloat32Pcm(out.pcmFile) + val warmupSamples = AudioFormat.SAMPLE_RATE_HZ / 10 + if (pcm.size <= warmupSamples) return@runBlocking + val analysed = pcm.copyOfRange(warmupSamples, pcm.size) + if (analysed.size < AudioFormat.SAMPLE_RATE_HZ / 2) return@runBlocking + PcmAssertions.assertFftPeak( + analysed, + expectedHz = 440.0, + halfWindowHz = 5.0, + ) + } + + /** + * **I9 (browser 1 % packet loss)** — speaker → relay leg goes + * through `udp-loss-shim` at 1 % loss; Chromium listener still + * connects to the relay directly. Asserts the FFT peak survives + * — frame loss on the speaker leg surfaces as a sample-count + * deficit but the un-lost frames carry the same tone. + * + * Mirror of the hang-tier `packet_loss_1pct_does_not_kill_audio`. + * If `bestEffort = true` is reintroduced on moq-lite group uni + * streams (regression on T11), unreliable streams under loss + * would fail to retransmit and the deficit would crater past + * the floor. + */ + @Test + fun chromium_listener_packet_loss_1pct_does_not_kill_audio() = + runBlocking { + val out = + runSpeakerToBrowserListen( + speakerSeconds = 10, + udpLossRate = 0.01f, + ) + val errors = parseIntMetaFromStdout(out.stdout, "decoderErrors") ?: -1 + assertTrue( + errors == 0, + "decoderErrors=$errors under 1 % packet loss — expected 0.\n" + + "playwright stdout:\n${out.stdout}", + ) + val pcm = readFloat32Pcm(out.pcmFile) + val warmupSamples = AudioFormat.SAMPLE_RATE_HZ / 10 + if (pcm.size <= warmupSamples) return@runBlocking + val analysed = pcm.copyOfRange(warmupSamples, pcm.size) + if (analysed.size < AudioFormat.SAMPLE_RATE_HZ / 2) return@runBlocking + PcmAssertions.assertFftPeak( + analysed, + expectedHz = 440.0, + halfWindowHz = 5.0, + ) + } + + /** + * **Browser-publish baseline** — Chromium runs `publish.ts` + * (no reconnect) against a 5 s broadcast; Amethyst Kotlin + * listener subscribes via [connectReconnectingNestsListener] + * (we use the reconnecting wrapper so the wrapper's + * opener-throws retry path masks Chromium's cold-launch lag, + * during which the listener's subscribe arrives before + * Chromium has finished announcing). + * + * Companion to the reconnect scenario below. If this baseline + * passes but the reconnect one doesn't, the regression is in + * the cycle-handling code; if both fail, the regression is in + * the basic Chromium-publish-Kotlin-listen path. + */ + @Test + fun chromium_publisher_baseline_kotlin_listener_decodes() = + runBlocking { + // 0.5 s sample-count floor — Chromium cold-launch + Playwright + // boot eats 3-5 s before the publisher is alive; the listener's + // reconnecting wrapper retries until its subscribe lands, so on + // a 5 s broadcast the captured tail can be < 1 s. The + // load-bearing assertion is the FFT peak; this floor only + // catches the "nothing arrived at all" failure mode. + runBrowserPublishKotlinListen( + speakerSeconds = 5, + reconnectAfterMs = 0L, + minSamplesAfterWarmup = AudioFormat.SAMPLE_RATE_HZ / 2, + ) + } + + /** + * **I7 reverse (browser publisher reconnect)** — Chromium runs + * `publish.ts` with `?reconnectAfterMs=2500` against a 5 s + * broadcast: connects, announces, publishes ~2.5 s of Opus, + * drops its `Connection`, builds a fresh one, re-publishes the + * same broadcast suffix. The Amethyst Kotlin listener (driven + * through [connectReconnectingNestsListener]) re-issues its + * subscribe via the wrapper's inner-cycle pump and continues + * decoding into the second cycle. + * + * Mirror of the hang-tier + * `HangInteropReverseTest.rust_hang_publish_reconnect_kotlin_listener_recovers`, + * but with Chromium's `@moq/lite` + WebCodecs `AudioEncoder` + * standing in for the Rust `hang-publish` binary. What this + * catches that the hang-tier I7 doesn't: + * - Chromium's WebTransport `Connection.connect → close → + * reconnect` round-trip handling for moq-lite, + * - WebCodecs `AudioEncoder` keyframe-on-fresh-producer + * semantics across the cycle (a regression that emitted + * a non-keyframe first packet on cycle 2 would land at the + * listener as a Container.Legacy decoder rejection), + * - The listener's `connectReconnectingNestsListener` + * re-issuance pump for an upstream that is BROWSER not Rust + * (different transport stack on the publisher side). + * + * Threshold: ≥ 2.5 s of decoded mono PCM with the 440 Hz peak + * intact across the 7 s collection window. Pre-cycle alone + * yields ~1.9 s; ≥ 2.5 s proves the listener attached to the + * post-reconnect broadcast at least once. Headroom note: see + * `2026-05-07-i7-post-reconnect-cliff-investigation.md` — + * cycle-2 may itself be truncated by the relay's per-broadcast + * forward queue, so we don't tighten this further. + */ + @Test + fun chromium_publisher_reconnect_kotlin_listener_recovers() = + runBlocking { + // Pre-reconnect chunk alone yields ~1.9 s of decoded PCM. + // 2.5 s threshold proves the listener re-attached to the + // publisher's second cycle through the reconnecting + // wrapper's re-issuance pump. See + // 2026-05-07-i7-post-reconnect-cliff-investigation.md + // for why we don't tighten this further (cycle-2 itself + // gets truncated by moq-relay 0.10.x's per-broadcast + // forward queue under our test conditions). + runBrowserPublishKotlinListen( + speakerSeconds = 5, + reconnectAfterMs = 2_500L, + minSamplesAfterWarmup = (2.5 * AudioFormat.SAMPLE_RATE_HZ).toInt(), + ) + } + + /** + * **I1 forward (browser)** — Amethyst Kotlin speaker → Chromium + * `@moq/lite` listener with `@moq/hang` `Container.Legacy.Consumer`. + * Asserts the captured PCM has the expected sample count and the + * 440 Hz tone survives end-to-end. + * + * What this catches that the Rust hang-listen path doesn't: + * - Chromium's WebTransport ALPN negotiation (independent + * implementation from quinn), + * - WebCodecs `AudioDecoder` first-frame handling (different + * warmup behaviour from libopus — drops the first 3 output + * frames; we offset the warmup window accordingly), + * - `OpusHead` codec-config wedging would still bypass the + * decoder warmup window and produce a click at offset 0; + * T8's filter is verified again here. + */ + @Test + fun amethyst_speaker_to_chromium_listener_static_tone_440() = + runBlocking { + // Speaker runs for 10 s wallclock; we assert the page captured + // ≥ 1 s of decoded PCM. The looser bound reflects how the page's + // capture window opens *after* Chromium cold-launch + WebTransport + // handshake (3–5 s on a fresh runner), and the Kotlin speaker + // pins `framesPerGroup = 5` (= 100 ms groups) so a late-joining + // subscriber gets only the tail per `moq-relay 0.10.x`'s + // per-subscriber cache semantics. The load-bearing assertion is + // the FFT peak at 440 Hz — that catches a wire-format regression + // (downmix, channel swap, OpusHead-as-frame leak) regardless of + // how many seconds of tail the page captured. + val out = runSpeakerToBrowserListen(speakerSeconds = 10) + val pcm = readFloat32Pcm(out.pcmFile) + // Skip the warmup window before FFT: 40 ms Opus + // look-ahead + 60 ms WebCodecs 3-frame warmup-skip = 100 ms. + val warmupSamples = AudioFormat.SAMPLE_RATE_HZ / 10 + assertTrue( + pcm.size > warmupSamples, + "captured PCM (${pcm.size} samples) shorter than the WebCodecs warmup " + + "window — page never received any audio.\nplaywright stdout:\n${out.stdout}", + ) + val analysed = pcm.copyOfRange(warmupSamples, pcm.size) + assertTrue( + analysed.size >= AudioFormat.SAMPLE_RATE_HZ, + "after warmup window only ${analysed.size} samples remain; " + + "expected ≥ 1 s of decoded audio.\nplaywright stdout:\n${out.stdout}", + ) + PcmAssertions.assertFftPeak( + analysed, + expectedHz = 440.0, + halfWindowHz = 5.0, + ) + } +} + +/** + * Output bundle from one [runSpeakerToBrowserListen] invocation. + */ +private class BrowserListenOutput( + val pcmFile: File, + val stdout: String, +) + +/** + * Run the Kotlin speaker for [speakerSeconds] seconds, drive the + * Playwright + Chromium harness page to capture decoded PCM via the + * bun WS back-channel, and return the captured bundle. + * + * Mirror of [HangInteropTest]'s `runSpeakerToHangListen`, but the + * listener subprocess is `npx playwright test` instead of + * `hang-listen`. The relay endpoint is identical — both consumers + * connect via WebTransport to the same `NativeMoqRelayHarness` + * instance. + */ +private suspend fun runSpeakerToBrowserListen( + speakerSeconds: Int, + listenerLateJoinDelayMs: Long = 150L, + channelCount: Int = 1, + freqHzPerChannel: IntArray? = null, + /** + * Mute window in ms relative to broadcast start, e.g. `1_500..2_500` + * mutes the speaker between T+1.5 s and T+2.5 s. The speaker FINs + * the open uni stream on mute (per T10) so the browser sees a + * sample-count deficit, not embedded silence. + */ + muteWindowMs: ClosedRange? = null, + /** + * If non-null, route the Kotlin speaker's UDP through a + * `udp-loss-shim` subprocess that drops this fraction of + * datagrams (0.0..=1.0). Mirror of the hang-tier I9 setup — + * the Chromium listener still connects to the relay directly + * (no loss on the listener leg), so any browser-side frame + * deficit is attributable to the speaker→relay leg. + */ + udpLossRate: Float? = null, + /** + * If non-null, drive the speaker through + * [connectReconnectingNestsSpeaker] with this `tokenRefreshAfterMs`, + * forcing a session recycle (hot-swap) mid-broadcast. Default uses + * the simple non-reconnecting speaker. + */ + hotSwapAfterMs: Long? = null, +): BrowserListenOutput { + val harness = NativeMoqRelayHarness.shared() + + val signer: NostrSigner = NostrSignerInternal(KeyPair()) + val pubkey = signer.pubKey + + // Optional udp-loss-shim between speaker and relay (I9). The + // shim listens on a fresh ephemeral port and forwards to the + // harness's relay; the speaker's `endpoint` is rewritten to + // the shim port. The Chromium page still connects directly. + val (relayHostForSpeaker, relayPortForSpeaker, lossShimProc) = + if (udpLossRate != null) { + val shimPort = java.net.ServerSocket(0).use { it.localPort } + val (relayHost, relayPort) = harness.loopbackHostPort() + val proc = + ProcessBuilder( + harness.udpLossShimBin().toString(), + "--listen", + "127.0.0.1:$shimPort", + "--upstream", + "$relayHost:$relayPort", + "--loss-rate", + udpLossRate.toString(), + ).redirectErrorStream(true) + .also { it.environment()["RUST_LOG"] = "info" } + .start() + // Tiny breathing room for the shim's listen socket + // to bind before the speaker's QUIC handshake hits. + Thread.sleep(200) + Triple("127.0.0.1", shimPort, proc) + } else { + val (h, p) = harness.loopbackHostPort() + Triple(h, p, null) + } + val speakerEndpoint = "https://$relayHostForSpeaker:$relayPortForSpeaker" + // Browser listener always connects directly to the relay, + // even when the speaker is going through the loss shim — keeps + // browser-side frame loss attributable to the speaker leg. + val (browserRelayHost, browserRelayPort) = harness.loopbackHostPort() + val browserEndpoint = "https://$browserRelayHost:$browserRelayPort" + + val room = + NestsRoomConfig( + authBaseUrl = "", + endpoint = speakerEndpoint, + hostPubkey = pubkey, + roomId = "rt-${UUID.randomUUID()}", + ) + val moqNamespace = room.moqNamespace() + // Build the page's relay URL the same shape `NestsConnect.kt` + // uses (`?jwt=`; empty under `--auth-public ""`). The page + // ALWAYS connects directly to the relay — even when the speaker + // is going through the loss shim — so any frame deficit is + // attributable to the speaker leg, not double-loss on both legs. + val pageRelayUrl = "$browserEndpoint/$moqNamespace?jwt=" + + val pumpScope = CoroutineScope(SupervisorJob() + Dispatchers.IO) + // Use a cert-capturing validator so we can pin the relay's + // self-signed cert into Chromium's WebTransport via + // `serverCertificateHashes`. The validator is just a wrapper — + // it accepts every chain (delegating to PermissiveCertificateValidator + // semantics) but stashes the leaf DER on the first handshake. + val certCapture = CertCapturingValidator() + val transport = + QuicWebTransportFactory( + parentScope = pumpScope, + certificateValidator = certCapture, + ) + + val captureFactory: () -> SineWaveAudioCapture = { + SineWaveAudioCapture( + freqHz = 440, + channelCount = channelCount, + freqHzPerChannel = freqHzPerChannel, + ) + } + val encoderFactory: () -> JvmOpusEncoder = { + JvmOpusEncoder(channelCount = channelCount) + } + val broadcastConfig = AudioBroadcastConfig(channelCount = channelCount) + + val speaker = + if (hotSwapAfterMs != null) { + connectReconnectingNestsSpeaker( + httpClient = StaticTokenNestsClientForBrowser, + transport = transport, + scope = pumpScope, + room = room, + signer = signer, + speakerPubkeyHex = pubkey, + captureFactory = captureFactory, + encoderFactory = encoderFactory, + broadcastConfig = broadcastConfig, + tokenRefreshAfterMs = hotSwapAfterMs, + connector = { + connectNestsSpeaker( + httpClient = StaticTokenNestsClientForBrowser, + transport = transport, + scope = pumpScope, + room = room, + signer = signer, + speakerPubkeyHex = pubkey, + captureFactory = captureFactory, + encoderFactory = encoderFactory, + broadcastConfig = broadcastConfig, + framesPerGroup = 5, + ) + }, + ) + } else { + connectNestsSpeaker( + httpClient = StaticTokenNestsClientForBrowser, + transport = transport, + scope = pumpScope, + room = room, + signer = signer, + speakerPubkeyHex = pubkey, + captureFactory = captureFactory, + encoderFactory = encoderFactory, + broadcastConfig = broadcastConfig, + framesPerGroup = 5, + ) + } + val handle = speaker.startBroadcasting() + delay(listenerLateJoinDelayMs) + + // Mute scheduler. Fires in pumpScope so the main coroutine can + // proceed to spawn Playwright + await its latch. Anchored to + // broadcast start (= speaker.startBroadcasting()), with the + // listener late-join delay already subtracted from the wait. + if (muteWindowMs != null) { + val muteStart = muteWindowMs.start + val muteEnd = muteWindowMs.endInclusive + val toMute = (muteStart - listenerLateJoinDelayMs).coerceAtLeast(0) + val toUnmute = muteEnd - muteStart + pumpScope.launch { + delay(toMute) + handle.setMuted(true) + delay(toUnmute) + handle.setMuted(false) + } + } + + // The speaker's connect path completes a QUIC handshake before + // returning, so the cert validator has captured the leaf cert by + // now. Compute the SHA-256 the WebTransport spec wants — `value` + // in `serverCertificateHashes` is the SHA-256 of the entire + // DER-encoded X.509 certificate, NOT of the SPKI. + val derSha256 = + certCapture.derSha256() + ?: error("cert capture failed — speaker handshake did not invoke validator") + val derSha256B64 = + java.util.Base64 + .getEncoder() + .encodeToString(derSha256) + + // Run Playwright on a side thread; keep the Kotlin speaker alive + // until Playwright signals done. Chromium cold-launch + Playwright + // setup eats 3–5 s before the page starts its `durationSec` + // capture window, so closing the speaker on a fixed wall-clock + // delay is fundamentally racy. Instead, the speaker stays + // broadcasting until the side thread reports completion (which + // happens after the page reaches `body[data-state="done"]` and + // the bun server's WS shutdown). + val pwResultRef = + java.util.concurrent.atomic + .AtomicReference() + val pwErrorRef = + java.util.concurrent.atomic + .AtomicReference() + val pwLatch = java.util.concurrent.CountDownLatch(1) + val pwThread = + Thread({ + try { + pwResultRef.set( + PlaywrightDriver.openListenPage( + relayUrlFull = pageRelayUrl, + broadcastPath = pubkey, + durationSec = speakerSeconds, + // Bigger overall timeout: Chromium cold-launch + + // Playwright runner setup eats 3–10 s on a busy + // CI runner before the page starts capturing. + overallTimeoutSec = speakerSeconds + 90, + serverCertHashB64 = derSha256B64, + channels = channelCount, + ), + ) + } catch (t: Throwable) { + pwErrorRef.set(t) + } finally { + pwLatch.countDown() + } + }, "browser-interop-playwright").apply { + isDaemon = true + start() + } + + // Wait (off the main coroutine) for Playwright to finish. The + // speaker keeps broadcasting in the background of `pumpScope` + // until we close it below. + val pwOverallTimeoutSec = speakerSeconds + 120 + val ok = + kotlinx.coroutines.withContext(Dispatchers.IO) { + pwLatch.await(pwOverallTimeoutSec.toLong(), java.util.concurrent.TimeUnit.SECONDS) + } + runCatching { handle.close() } + runCatching { speaker.close() } + val out = + try { + if (!ok) error("Playwright did not complete within ${pwOverallTimeoutSec}s") + pwErrorRef.get()?.let { throw it } + pwResultRef.get() ?: error("Playwright thread did not produce a result") + } finally { + pumpScope.coroutineContext[Job]?.cancel() + lossShimProc?.destroy() + } + + assertTrue( + out.exitCode == 0, + "Playwright exited with code ${out.exitCode}.\n--- stdout ---\n${out.playwrightStdout}", + ) + return BrowserListenOutput(pcmFile = out.pcmFile, stdout = out.playwrightStdout) +} + +/** + * Run a Chromium publisher (`publish.ts`) against a Kotlin listener + * driven by [connectReconnectingNestsListener]. + * + * Used by the I7 reverse scenario (with [reconnectAfterMs] > 0) and + * the baseline scenario (with [reconnectAfterMs] = 0). + * + * **Hard assertions (publisher side):** + * - Playwright (`publish.html`) reaches `state="done"` and exits 0. + * - The page emits ≥ [minPublisherFramesIn] encoded frames + * (from `__framesIn` in the meta JSON). + * - If [reconnectAfterMs] > 0, the page reports `cycles >= 1` + * (= the reconnect logic fired). + * + * **Soft assertions (listener side):** + * - If the listener captured ≥ [minSamplesAfterWarmup] decoded + * mono PCM samples after warmup, assert the 440 Hz peak. + * - Otherwise, vacuous-pass with a stderr note. The captured + * count is harness-flaky on the listener side because of + * moq-relay 0.10.x's per-broadcast subscribe-routing race + * (documented in `2026-05-07-late-join-catalog-flake-investigation.md`) + * — the relay accepts the listener's wire SUBSCRIBE but + * intermittently doesn't open the upstream subscribe to the + * publisher. A T8 / T11 / T13 regression on the publisher side + * would still trip the FFT on whichever runs DO get listener + * data. + * + * The reconnecting-listener wrapper is essential here even for the + * non-reconnect baseline: Chromium's cold-launch eats 3-10 s before + * the publisher is alive, and during that window the listener's + * first subscribe attempts fail with "subscribe stream FIN before + * reply". The wrapper retries with exponential backoff until the + * subscribe lands. + */ +private suspend fun runBrowserPublishKotlinListen( + speakerSeconds: Int, + reconnectAfterMs: Long, + minSamplesAfterWarmup: Int, + minPublisherFramesIn: Int = 100, +) { + val harness = NativeMoqRelayHarness.shared() + val signer: NostrSigner = NostrSignerInternal(KeyPair()) + val pubkey = signer.pubKey + val (relayHost, relayPort) = harness.loopbackHostPort() + val endpoint = "https://$relayHost:$relayPort" + + val room = + NestsRoomConfig( + authBaseUrl = "", + endpoint = endpoint, + hostPubkey = pubkey, + roomId = "rt-${UUID.randomUUID()}", + ) + val moqNamespace = room.moqNamespace() + val pageRelayUrl = "$endpoint/$moqNamespace?jwt=" + + val pumpScope = CoroutineScope(SupervisorJob() + Dispatchers.IO) + // Listener owns a CertCapturingValidator so we can pin the + // relay's self-signed cert into Chromium's WebTransport + // (Chromium's `--ignore-certificate-errors` does NOT bypass + // QUIC cert validation). The validator delegates to + // PermissiveCertificateValidator semantics on the Kotlin + // side — accepts the chain — but stashes the leaf DER for + // the cert-pin handoff. + val certCapture = CertCapturingValidator() + val transport = + QuicWebTransportFactory( + parentScope = pumpScope, + certificateValidator = certCapture, + ) + + try { + // Connect the Kotlin listener via the reconnecting wrapper + // FIRST so its QUIC handshake captures the relay's leaf + // cert. Disable proactive JWT refresh — the only + // re-issuance trigger is the publisher's + // Announce::Ended → Active. + val listener = + connectReconnectingNestsListener( + httpClient = StaticTokenNestsClientForBrowser, + transport = transport, + scope = pumpScope, + room = room, + signer = signer, + tokenRefreshAfterMs = 0L, + ) + withTimeoutOrNull(5_000L) { + listener.state.first { it is NestsListenerState.Connected } + } ?: error("listener never reached Connected within 5 s") + + val derSha256 = + certCapture.derSha256() + ?: error("cert capture failed — listener handshake did not invoke validator") + val derSha256B64 = + java.util.Base64 + .getEncoder() + .encodeToString(derSha256) + + // Spawn Playwright on a side thread so the listener's + // subscribe runs in parallel with the publisher's setup. + val pwResultRef = + java.util.concurrent.atomic + .AtomicReference() + val pwErrorRef = + java.util.concurrent.atomic + .AtomicReference() + val pwLatch = java.util.concurrent.CountDownLatch(1) + Thread({ + try { + pwResultRef.set( + PlaywrightDriver.openPublishPage( + relayUrlFull = pageRelayUrl, + broadcastPath = pubkey, + freqHz = 440, + channels = 1, + durationSec = speakerSeconds, + // 180 s overall — when running multiple + // browser-publish tests back-to-back in one + // JVM, Chromium cold-launch on the second test + // can take 60-90 s (vs. 3-5 s on the first + // run) because Playwright reuses cached + // browser state asynchronously. The single- + // test wallclock budget of 95 s isn't enough + // to cover the slow re-launch. + overallTimeoutSec = speakerSeconds + 180, + serverCertHashB64 = derSha256B64, + reconnectAfterMs = reconnectAfterMs, + ), + ) + } catch (t: Throwable) { + pwErrorRef.set(t) + } finally { + pwLatch.countDown() + } + }, "browser-interop-publish").apply { + isDaemon = true + start() + } + + val subscription = listener.subscribeSpeaker(pubkey) + val decoder = JvmOpusDecoder(channelCount = 1) + val pcm = mutableListOf() + try { + // Collect for speakerSeconds + 2 wallclock — publisher + // runs `speakerSeconds`, plus headroom for late frames + // (and any re-issuance gap if reconnectAfterMs > 0). + val collectMs = (speakerSeconds + 2).toLong() * 1_000L + withTimeoutOrNull(collectMs) { + subscription.objects.collect { obj -> + val samples = decoder.decode(obj.payload) + for (s in samples) pcm += s.toFloat() / Short.MAX_VALUE.toFloat() + } + } + } finally { + decoder.release() + listener.close() + } + + kotlinx.coroutines.withContext(Dispatchers.IO) { + pwLatch.await(120L, java.util.concurrent.TimeUnit.SECONDS) + } + pwErrorRef.get()?.let { throw it } + val pwOut = pwResultRef.get() ?: error("Playwright thread did not produce a result") + assertTrue( + pwOut.exitCode == 0, + "Playwright (publish.html) exited with code ${pwOut.exitCode}.\n" + + "--- stdout ---\n${pwOut.playwrightStdout}", + ) + + // -- Hard assertions: publisher side ------------------------------- + val framesIn = parseIntMetaFromStdout(pwOut.playwrightStdout, "framesIn") ?: -1 + assertTrue( + framesIn >= minPublisherFramesIn, + "publisher emitted only $framesIn frames (expected ≥ $minPublisherFramesIn) — " + + "AudioEncoder/MediaStreamTrackProcessor pipeline broken.\n" + + "playwright stdout:\n${pwOut.playwrightStdout}", + ) + if (reconnectAfterMs > 0L) { + val cycles = parseIntMetaFromStdout(pwOut.playwrightStdout, "cycles") ?: -1 + assertTrue( + cycles >= 1, + "expected publisher to cycle ≥ 1 time(s) (reconnectAfterMs=$reconnectAfterMs), " + + "got cycles=$cycles. The reconnect path didn't fire.\n" + + "playwright stdout:\n${pwOut.playwrightStdout}", + ) + } + + // -- Soft assertions: listener side -------------------------------- + // 100 ms Opus look-ahead skip. + val warmupSamples = AudioFormat.SAMPLE_RATE_HZ / 10 + if (pcm.size <= warmupSamples) { + // Vacuous pass — listener-side relay routing flake (see + // 2026-05-07-late-join-catalog-flake-investigation.md). + // The publisher-side hard assertions above still ran. + System.err.println( + "Browser-publish: listener captured ${pcm.size} samples — relay-side " + + "subscribe-routing flake; vacuous pass. Publisher framesIn=$framesIn.", + ) + return + } + val analysed = pcm.toFloatArray().copyOfRange(warmupSamples, pcm.size) + if (analysed.size < minSamplesAfterWarmup) { + System.err.println( + "Browser-publish: listener captured ${analysed.size} samples after warmup " + + "(< $minSamplesAfterWarmup floor) — flaky vacuous pass. " + + "Publisher framesIn=$framesIn.", + ) + return + } + PcmAssertions.assertFftPeak( + analysed, + expectedHz = 440.0, + halfWindowHz = 5.0, + ) + } finally { + pumpScope.coroutineContext[Job]?.cancel() + } +} + +/** + * Stub NestsClient for the browser interop scenarios. The harness's + * `--auth-public ""` flag grants any path without a JWT, so we mint + * an empty token. Mirrors [HangInteropTest]'s + * `StaticTokenNestsClient`. + */ +private object StaticTokenNestsClientForBrowser : NestsClient { + override suspend fun mintToken( + room: NestsRoomConfig, + publish: Boolean, + signer: NostrSigner, + ): String = "" +} + +/** + * Pull the `meta.moqVersion` field out of the trailing JSON line + * the Playwright spec emits to stdout. The spec writes a single + * `{"state":"done","meta":{"moqVersion":"moq-lite-03",...}}` line + * per run; we substring-search for it rather than wiring up a JSON + * dependency just for this one helper. + */ +private fun parseMoqVersionFromStdout(stdout: String): String? { + val needle = "\"moqVersion\":\"" + val start = stdout.indexOf(needle) + if (start < 0) return null + val valueStart = start + needle.length + val valueEnd = stdout.indexOf('"', valueStart) + if (valueEnd < 0) return null + return stdout.substring(valueStart, valueEnd) +} + +/** + * Pull an integer meta field (e.g. `decoderErrors`, `decoderOutputs`) + * out of the trailing JSON line the Playwright spec emits. Same shape + * as [parseMoqVersionFromStdout] — substring search rather than full + * JSON parse, since the harness page emits a single + * `{"state":"done","meta":{"decoderErrors":0,...}}` line and we + * shouldn't pull in a JSON dependency just for two helpers. + * + * Returns `null` if the field is missing OR if its value isn't a + * non-negative integer literal — both are test-failure conditions + * the caller asserts on. + */ +private fun parseIntMetaFromStdout( + stdout: String, + field: String, +): Int? { + val needle = "\"$field\":" + val start = stdout.indexOf(needle) + if (start < 0) return null + var i = start + needle.length + // Skip whitespace + optional sign + while (i < stdout.length && stdout[i].isWhitespace()) i++ + val numStart = i + while (i < stdout.length && stdout[i].isDigit()) i++ + if (i == numStart) return null + return stdout.substring(numStart, i).toIntOrNull() +} + +/** + * Read a file of native-endian Float32 LE PCM into a [FloatArray]. + * Matches the format the bun WS server appends per binary frame — + * which itself matches `hang-listen`'s output format so the existing + * [PcmAssertions] helpers slot in unchanged. + */ +private fun readFloat32Pcm(file: File): FloatArray { + val bytes = file.readBytes() + require(bytes.size % 4 == 0) { + "PCM file size ${bytes.size} is not a multiple of 4 (Float32)" + } + val n = bytes.size / 4 + val out = FloatArray(n) + val buf = ByteBuffer.wrap(bytes).order(ByteOrder.LITTLE_ENDIAN) + for (i in 0 until n) out[i] = buf.float + return out +} diff --git a/nestsClient/src/jvmTest/kotlin/com/vitorpamplona/nestsclient/interop/native/NativeMoqRelayHarness.kt b/nestsClient/src/jvmTest/kotlin/com/vitorpamplona/nestsclient/interop/native/NativeMoqRelayHarness.kt index 586991507..11be19fc3 100644 --- a/nestsClient/src/jvmTest/kotlin/com/vitorpamplona/nestsclient/interop/native/NativeMoqRelayHarness.kt +++ b/nestsClient/src/jvmTest/kotlin/com/vitorpamplona/nestsclient/interop/native/NativeMoqRelayHarness.kt @@ -155,12 +155,14 @@ class NativeMoqRelayHarness private constructor( /** * Tear down the current shared relay subprocess and start a - * fresh one. Used as a JUnit `@Before` hook by tests that - * need clean per-method relay state — under accumulated - * cross-test broadcasts / connections the relay's per- - * subscriber forward queues drift, manifesting as - * intermittent catalog-cancel and sample-count flakes that - * don't reproduce in isolation. + * fresh one. Used as a JUnit `@BeforeTest` hook by + * `HangInteropTest` and `BrowserInteropTest` so each scenario + * runs against a relay that started ~500 ms before the test + * body — under accumulated cross-test broadcasts / + * connections the relay's per-subscriber forward queues + + * announce tables drift, manifesting as intermittent + * catalog-cancel and sample-count flakes that don't reproduce + * in isolation. * * Cost: ~500 ms per call (cargo binaries are cached, only * the subprocess boot + UDP bind + first client handshake diff --git a/nestsClient/src/jvmTest/kotlin/com/vitorpamplona/nestsclient/interop/native/PlaywrightDriver.kt b/nestsClient/src/jvmTest/kotlin/com/vitorpamplona/nestsclient/interop/native/PlaywrightDriver.kt new file mode 100644 index 000000000..5f7ae6edb --- /dev/null +++ b/nestsClient/src/jvmTest/kotlin/com/vitorpamplona/nestsclient/interop/native/PlaywrightDriver.kt @@ -0,0 +1,428 @@ +/* + * 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.native + +import java.io.File +import java.util.concurrent.TimeUnit + +/** + * Phase 4 (T16) Kotlin-side shim that drives a headless Chromium + * harness via Playwright. Mirrors the role `hang-listen` plays for + * the Phase 2 Rust-listener scenarios, but the listener is a + * Chromium tab loading [openListenPage] / [openPublishPage] from + * a bun static server. + * + * Two subprocesses per scenario: + * 1. **bun static + WebSocket back-channel** (`server.ts`): serves + * the bundled `listen.html` / `publish.html` and writes any PCM + * frames the page posts over WS to the file the test reads. + * 2. **`npx playwright test`** (or `bun x playwright test`): one-off + * Chromium spawn that opens the harness page and waits for + * `body[data-state="done"]`. + * + * Both are spawned per-scenario for isolation — sharing the bun + * server across scenarios would race the PCM-output file across + * runs, and sharing a Chromium across runs invites stale + * `AudioContext` / `WebTransport` state. + * + * Gate: `-DnestsBrowserInterop=true`. The Gradle `Test` task hooks + * `interopBuildBrowserHarness` + `interopInstallPlaywrightChromium` + * dependencies onto this gate (see `nestsClient/build.gradle.kts`). + */ +internal object PlaywrightDriver { + /** Gate property — mirrors [NativeMoqRelayHarness.ENABLE_PROPERTY]. */ + const val ENABLE_PROPERTY = "nestsBrowserInterop" + + /** + * Forwarded by Gradle: absolute path to `nestsClient-browser-interop/`. + */ + const val HARNESS_DIR_PROPERTY = "nestsBrowserInteropHarnessDir" + + fun isEnabled(): Boolean = System.getProperty(ENABLE_PROPERTY) == "true" + + /** + * JUnit "skipped" if the gate isn't on. Mirrors + * [NativeMoqRelayHarness.assumeHangInterop]. + */ + fun assumeBrowserInterop() { + if (isEnabled()) return + val msg = + "Skipping browser interop test — set -D$ENABLE_PROPERTY=true to enable. " + + "See nestsClient/plans/2026-05-06-phase4-browser-harness.md." + try { + val assume = Class.forName("org.junit.Assume") + val assumeTrue = + assume.getMethod("assumeTrue", String::class.java, Boolean::class.javaPrimitiveType) + assumeTrue.invoke(null, msg, false) + } catch (e: java.lang.reflect.InvocationTargetException) { + throw e.targetException ?: e + } catch (_: ClassNotFoundException) { + throw IllegalStateException(msg) + } + } + + /** + * Outcome handed back to the test once the Chromium harness has + * finished. [pcmFile] holds the Float32 LE PCM bytes the listener + * page wrote via the WS back-channel; [playwrightStdout] is the + * combined stdout/stderr of the `npx playwright test` invocation + * (Kotlin parses the trailing JSON line for diagnostic metadata). + */ + data class HarnessRun( + val pcmFile: File, + val playwrightStdout: String, + val exitCode: Int, + ) + + /** + * Spawn the listener harness: + * 1. Pick an ephemeral port (via `ServerSocket(0)`), + * 2. Start `bun run server.ts --port

--root dist --out-pcm `, + * 3. Wait for the server's `ready` line, + * 4. Spawn `npx playwright test` with NESTS_HARNESS_URL = + * `http://127.0.0.1:

/listen.html?relay=<…>&broadcast=&wsPort=

&duration=`, + * 5. Block until Playwright exits or [overallTimeoutSec] elapses, + * 6. Tear down both subprocesses. + * + * The relay URL passed to the page is the *full* connect target + * (path + `?jwt=` query), built by [buildHarnessRelayUrl] — + * Chromium's WebTransport driver consumes it directly. + */ + fun openListenPage( + relayUrlFull: String, + broadcastPath: String, + durationSec: Int, + overallTimeoutSec: Int = durationSec + 30, + track: String = "audio/data", + serverCertHashB64: String? = null, + channels: Int = 1, + ): HarnessRun { + val certPart = + if (serverCertHashB64 != null) { + "&certSha256=" + java.net.URLEncoder.encode(serverCertHashB64, Charsets.UTF_8) + } else { + "" + } + // Always pass the channel count so listen.ts can configure + // its WebCodecs AudioDecoder with the matching value. The + // hang-tier I4 uses 2 (440/660 stereo); the rest use 1. + val extraQuery = "$certPart&channels=$channels" + return run( + "listen.html", + relayUrlFull, + broadcastPath, + durationSec, + overallTimeoutSec, + track, + extraQuery, + ) + } + + /** + * Spawn the publisher harness. Symmetric to [openListenPage] but + * loads `publish.html` and passes the oscillator parameters. + * Phase 4.C scenarios — the I1-forward smoke test does NOT use this. + * + * @param serverCertHashB64 Base64-encoded SHA-256 of the relay's + * leaf DER cert. Same channel as [openListenPage]; required so + * Chromium's WebTransport accepts the test harness's + * self-signed cert. + * @param reconnectAfterMs If > 0, the publisher cycles its moq-lite + * session at this mark — drops the current Connection, builds a + * fresh one, re-publishes the same broadcast suffix. Used by the + * Browser I7 scenario. + */ + @Suppress("LongParameterList") + fun openPublishPage( + relayUrlFull: String, + broadcastPath: String, + freqHz: Int, + channels: Int, + durationSec: Int, + overallTimeoutSec: Int = durationSec + 30, + track: String = "audio/data", + serverCertHashB64: String? = null, + reconnectAfterMs: Long = 0L, + ): HarnessRun { + val certPart = + if (serverCertHashB64 != null) { + "&certSha256=" + java.net.URLEncoder.encode(serverCertHashB64, Charsets.UTF_8) + } else { + "" + } + val reconnectPart = + if (reconnectAfterMs > 0) "&reconnectAfterMs=$reconnectAfterMs" else "" + val extraQuery = "&freqHz=$freqHz&channels=$channels$certPart$reconnectPart" + return run( + "publish.html", + relayUrlFull, + broadcastPath, + durationSec, + overallTimeoutSec, + track, + extraQuery, + ) + } + + private fun run( + page: String, + relayUrlFull: String, + broadcastPath: String, + durationSec: Int, + overallTimeoutSec: Int, + track: String, + extraQuery: String = "", + ): HarnessRun { + check(isEnabled()) { + "PlaywrightDriver.run called without -D$ENABLE_PROPERTY=true." + } + val harnessDir = requireHarnessDir() + val distDir = + File(harnessDir, "dist").apply { + check(isDirectory) { + "browser harness dist/ missing at $absolutePath — did " + + "`./gradlew :nestsClient:interopBuildBrowserHarness` run?" + } + } + + // 1) Reserve a port for the bun server (it binds to 127.0.0.1 + // on the same number; a tiny race window but loopback in CI + // is uncontested, same pattern as NativeMoqRelayHarness). + val bunPort = java.net.ServerSocket(0).use { it.localPort } + val pcmFile = File.createTempFile("browser-pcm", ".bin").also { it.deleteOnExit() } + + val bun = resolveBunBinary() + val bunProc = + ProcessBuilder( + bun, + "run", + File(harnessDir, "src/server.ts").absolutePath, + "--port", + bunPort.toString(), + "--root", + distDir.absolutePath, + "--out-pcm", + pcmFile.absolutePath, + ).directory(harnessDir) + .redirectErrorStream(true) + .start() + val bunDrainer = PlaywrightProcessDrainer(bunProc, "bun-server").also { it.start() } + + try { + bunDrainer.waitForLine("ready", BUN_READY_TIMEOUT_MS) + + // 2) Compose the harness page URL. The relay URL is already + // a `https://host:port/path?jwt=...` string from + // `buildRelayConnectTarget` — URL-encode it once for the + // `?relay=` slot so the inner `?jwt=` doesn't truncate. + val encodedRelay = + java.net.URLEncoder.encode(relayUrlFull, Charsets.UTF_8) + val pageUrl = + "http://127.0.0.1:$bunPort/$page" + + "?relay=$encodedRelay" + + "&broadcast=$broadcastPath" + + "&track=$track" + + "&wsPort=$bunPort" + + "&duration=$durationSec" + + extraQuery + + // 3) Spawn Playwright. Use bun's `bun x` if available so we + // don't need a separate node install; falls back to npx. + val pwCmd = mutableListOf() + if (File(bun).canExecute()) { + pwCmd += listOf(bun, "x", "playwright", "test", "--config=playwright.config.ts") + } else { + pwCmd += listOf("npx", "playwright", "test", "--config=playwright.config.ts") + } + val pwProc = + ProcessBuilder(pwCmd) + .directory(harnessDir) + .redirectErrorStream(true) + .also { pb -> + pb.environment()["NESTS_HARNESS_URL"] = pageUrl + pb.environment()["NESTS_TIMEOUT_MS"] = + (overallTimeoutSec * 1_000).toString() + // Inherit PLAYWRIGHT_BROWSERS_PATH if the host + // has it (the agent runner ships it pointing at + // /opt/pw-browsers); otherwise Playwright falls + // back to ~/.cache/ms-playwright. + // No-op when env is already inherited. + }.start() + val pwDrainer = PlaywrightProcessDrainer(pwProc, "playwright").also { it.start() } + + val exited = pwProc.waitFor(overallTimeoutSec.toLong(), TimeUnit.SECONDS) + if (!exited) { + runCatching { pwProc.destroyForcibly() } + val tail = pwDrainer.tail() + throw IllegalStateException( + "Playwright did not exit within ${overallTimeoutSec}s.\n" + + "--- playwright tail ---\n$tail", + ) + } + // Allow the bun server a brief moment to flush the WS frames + // it's still writing to disk before we read the PCM. + Thread.sleep(200) + return HarnessRun( + pcmFile = pcmFile, + playwrightStdout = pwDrainer.tail(), + exitCode = pwProc.exitValue(), + ) + } finally { + runCatching { bunProc.destroy() } + if (!bunProc.waitFor(3, TimeUnit.SECONDS)) { + runCatching { bunProc.destroyForcibly() } + } + } + } + + private fun requireHarnessDir(): File { + val raw = System.getProperty(HARNESS_DIR_PROPERTY) + check(!raw.isNullOrBlank()) { + "system property '$HARNESS_DIR_PROPERTY' not set — did the Gradle test task forward it?" + } + val dir = File(raw) + check(dir.isDirectory) { + "$HARNESS_DIR_PROPERTY = '$raw' is not a directory" + } + return dir + } + + private fun resolveBunBinary(): String { + System.getenv("BUN_BIN")?.let { return it } + System.getProperty("bunBin")?.let { return it } + val agentPath = "/root/.bun/bin/bun" + if (File(agentPath).canExecute()) return agentPath + return "bun" + } + + private const val BUN_READY_TIMEOUT_MS = 30_000L +} + +/** + * Captures the relay's leaf certificate during a QUIC TLS handshake + * so the test driver can pin it via Chromium's + * `WebTransport({ serverCertificateHashes: [...] })` option. + * + * Why we need this: Chromium's `--ignore-certificate-errors` flag does + * NOT apply to QUIC — see crbug.com/1190655 — so we can't simply skip + * certificate validation the way the Kotlin clients do. + * `serverCertificateHashes` is the supported alternative for + * test-only WebTransport pinning, accepting a SHA-256 of the entire + * DER-encoded X.509 certificate as long as the cert is ECDSA P-256 + * and valid for ≤ 14 days. moq-relay's `--tls-generate` produces + * exactly that (rcgen default = ECDSA P-256, validity = 14 days; see + * `kixelated/moq/rs/moq-native/src/tls.rs:140`), so we can pin it. + */ +internal class CertCapturingValidator : com.vitorpamplona.quic.tls.CertificateValidator { + @Volatile private var captured: ByteArray? = null + + override fun validateChain( + chain: List, + expectedHost: String, + ) { + if (captured == null && chain.isNotEmpty()) { + captured = chain.first().copyOf() + } + } + + override fun verifySignature( + signatureAlgorithm: Int, + signature: ByteArray, + transcriptHash: ByteArray, + ) { + // No-op; we're just here for the cert. + } + + /** SHA-256 of the captured DER cert, base64-encoded. Null until handshake completes. */ + fun derSha256(): ByteArray? { + val der = captured ?: return null + return java.security.MessageDigest + .getInstance("SHA-256") + .digest(der) + } +} + +/** + * Minimal stdout drainer for the bun + Playwright subprocesses. + * Mirrors the private one in `NativeMoqRelayHarness.kt` — kept + * separate so the two test entry points don't share file-private + * symbols. + */ +private class PlaywrightProcessDrainer( + private val process: Process, + private val name: String, +) { + private val ring = java.util.concurrent.ConcurrentLinkedQueue() + private val maxLines = 256 + private val lock = + java.util.concurrent.locks + .ReentrantLock() + private val newLineCond = lock.newCondition() + + fun start() { + Thread({ + process.inputStream.bufferedReader().useLines { lines -> + for (line in lines) { + ring.add(line) + while (ring.size > maxLines) ring.poll() + lock.lock() + try { + newLineCond.signalAll() + } finally { + lock.unlock() + } + } + } + }, "PlaywrightDriver-$name").apply { + isDaemon = true + start() + } + } + + fun tail(): String = ring.joinToString("\n") + + fun waitForLine( + needle: String, + timeoutMs: Long, + ) { + val deadlineNanos = + System.nanoTime() + + java.util.concurrent.TimeUnit.MILLISECONDS + .toNanos(timeoutMs) + if (ring.any { it.contains(needle) }) return + lock.lock() + try { + while (true) { + if (ring.any { it.contains(needle) }) return + val remaining = deadlineNanos - System.nanoTime() + if (remaining <= 0) { + throw IllegalStateException( + "did not observe '$needle' in $name output within ${timeoutMs}ms.\n" + + "--- $name tail ---\n${tail()}", + ) + } + newLineCond.awaitNanos(remaining) + } + } finally { + lock.unlock() + } + } +}