From 71cf99dc228684bc4573248f89a5c2aedd471e3f Mon Sep 17 00:00:00 2001 From: Claude Date: Sun, 26 Apr 2026 18:01:00 +0000 Subject: [PATCH] feat(audio-rooms): moq-lite speaker side end-to-end (phase 5c-speaker) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Production speaker path now runs on moq-lite, so connectNestsSpeaker exchanges real moq-lite framing with the nostrnests reference relay. Transport layer: - WebTransportSession.incomingBidiStreams() — peer-initiated bidi flow. moq-lite publishers receive Announce + Subscribe bidis from the relay (rs/moq-lite/src/lite/publisher.rs:40 uses Stream::accept(session)), so the abstraction grew the accept-bidi-from-peer surface. - WebTransportSession.openUniStream() — locally-opened uni stream for group push (rs/moq-lite/src/lite/publisher.rs:338 uses session.open_uni()). - :quic WtPeerStreamDemux StrippedWtStream now carries optional send/finish closures. The demux takes the QuicConnectionDriver so wakeups fire after each app-level write on a peer-initiated bidi. - FakeWebTransport now exposes incomingBidiStreams + openUniStream directly; the openPeerUniStream test helper went away (production flow covers it). Session layer: - MoqLiteSession.publish(suffix) — claims a broadcast suffix and lazily launches a relay→us bidi pump. ControlType=Announce reads AnnouncePlease, replies Active(suffix). ControlType=Subscribe reads body, replies SubscribeOk, registers the inbound subscription. - MoqLitePublisherHandle — startGroup / send / endGroup / close semantics. send opens a uni stream per group with DataType=0 + GroupHeader and pushes varint(size)+payload frames. close emits Announce(Ended) on every active announce bidi, FINs the uni. Application layer: - AudioRoomMoqLiteBroadcaster — sibling of AudioRoomBroadcaster but drives MoqLitePublisherHandle (keeps IETF broadcaster intact for its unit tests). - MoqLiteNestsSpeaker — NestsSpeaker adapter, mirror of MoqLiteNestsListener on the publish side. - connectNestsSpeaker now opens a MoqLiteSession (no SETUP) and returns MoqLiteNestsSpeaker. Tests: - 4 new MoqLiteSessionTest cases: publisher_replies_to_announcePlease_with_active_announce, publisher_acks_subscribe_and_pushes_group_data_on_uni_stream, publisher_send_returns_false_when_no_inbound_subscriber, publisher_close_emits_ended_announce. Verified :commons:compileKotlinJvm + :amethyst:compilePlayDebugKotlin both still compile against the swap. Docs (plans + CLAUDE.md) refreshed to reflect speaker-side landing. --- .claude/CLAUDE.md | 7 +- .../2026-04-26-audio-rooms-completion.md | 20 +- nestsClient/plans/2026-04-26-moq-lite-gap.md | 42 +- .../nestsclient/MoqLiteNestsSpeaker.kt | 169 ++++++++ .../vitorpamplona/nestsclient/NestsConnect.kt | 41 +- .../audio/AudioRoomMoqLiteBroadcaster.kt | 140 ++++++ .../nestsclient/moq/lite/MoqLiteSession.kt | 409 +++++++++++++++++- .../nestsclient/transport/FakeWebTransport.kt | 34 +- .../transport/WebTransportSession.kt | 21 + .../moq/lite/MoqLiteSessionTest.kt | 140 +++++- .../transport/QuicWebTransportFactory.kt | 70 +++ .../QuicWebTransportSessionState.kt | 2 +- .../quic/webtransport/WtPeerStreamDemux.kt | 46 ++ 13 files changed, 1052 insertions(+), 89 deletions(-) create mode 100644 nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/MoqLiteNestsSpeaker.kt create mode 100644 nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/audio/AudioRoomMoqLiteBroadcaster.kt diff --git a/.claude/CLAUDE.md b/.claude/CLAUDE.md index 6e1924027..aedfbdf9d 100644 --- a/.claude/CLAUDE.md +++ b/.claude/CLAUDE.md @@ -14,9 +14,10 @@ humans, agents, and interop tests. `quic` is a from-scratch pure-Kotlin QUIC v1 WebTransport client (no JNI, no BouncyCastle), built because no Android-compatible Java QUIC library exists. `nestsClient` runs the audio-room protocol on top of `:quic` for the NIP-53 audio-rooms feature. It implements both IETF `draft-ietf-moq-transport-17` (under -`moq/`) and **moq-lite Lite-03** (kixelated's variant, under `moq/lite/`) — the -production listener path uses moq-lite to interop with the nostrnests reference -relay. Speaker-side moq-lite is still pending; see +`moq/`) and **moq-lite Lite-03** (kixelated's variant, under `moq/lite/`); the +production listener AND speaker paths both run on moq-lite to interop with the +nostrnests reference relay. The IETF code is kept as a reference + unit-test +implementation for any future IETF target; see `nestsClient/plans/2026-04-26-moq-lite-gap.md`. ## Architecture diff --git a/nestsClient/plans/2026-04-26-audio-rooms-completion.md b/nestsClient/plans/2026-04-26-audio-rooms-completion.md index 99de1e17d..41df2580c 100644 --- a/nestsClient/plans/2026-04-26-audio-rooms-completion.md +++ b/nestsClient/plans/2026-04-26-audio-rooms-completion.md @@ -2,17 +2,15 @@ What's left between today's code and shippable audio rooms in Amethyst. -> **STATUS UPDATE (2026-04-26 PM):** the interop test suite (phases 1–5d -> of the nostrnests work, commits `3283d30` → `41f4dcd`) uncovered that -> nostrnests runs on **moq-lite** (kixelated's variant), not IETF -> `draft-ietf-moq-transport-17`. **Listener-side moq-lite is now landed** -> — `connectNestsListener` opens a `MoqLiteSession` and adapts frames -> back to `MoqObject` so `AudioRoomPlayer` / `AudioRoomViewModel` keep -> working unchanged. The Phase M1 manual validation against -> `nostrnests.com` should now actually work end-to-end (modulo Docker -> harness running). Speaker-side moq-lite is still pending — needs a -> small `WebTransportSession.acceptBidiStream` extension; tracked in -> [`2026-04-26-moq-lite-gap.md`](2026-04-26-moq-lite-gap.md). +> **STATUS UPDATE (2026-04-26 PM):** the interop test suite (phases 1–5 +> of the nostrnests work) uncovered that nostrnests runs on **moq-lite** +> (kixelated's variant), not IETF `draft-ietf-moq-transport-17`. Both +> the listener and speaker sides are now wired through moq-lite +> (`MoqLiteNestsListener` / `MoqLiteNestsSpeaker`), the WebTransport +> abstraction grew `incomingBidiStreams` + `openUniStream`, and Phase M1 +> manual validation against `nostrnests.com` should work end-to-end. +> See [`2026-04-26-moq-lite-gap.md`](2026-04-26-moq-lite-gap.md) for +> the exact wire spec + landing summary. ## Where we are diff --git a/nestsClient/plans/2026-04-26-moq-lite-gap.md b/nestsClient/plans/2026-04-26-moq-lite-gap.md index 587cbd114..bf4924df8 100644 --- a/nestsClient/plans/2026-04-26-moq-lite-gap.md +++ b/nestsClient/plans/2026-04-26-moq-lite-gap.md @@ -1,8 +1,9 @@ # Plan: bridge the moq-lite protocol gap -**Status:** **listener side done** (phase 5a → 5d, commits `fb47a4c` → -`41f4dcd`); speaker side blocked on a `WebTransportSession` API -extension. Default `:nestsClient:jvmTest` suite (124 tests) passes. +**Status:** **listener AND speaker sides done** (phases 5a → 5d, commit +range `fb47a4c` → `5914e9e` + the speaker phase 5c-speaker landing in +this PR). Default `:nestsClient:jvmTest` suite passes; downstream +`:commons` and `:amethyst` compile clean. **Origin:** discovered while writing the nostrnests interop test suite (phases 1–4). @@ -206,30 +207,21 @@ range. `priority` is a plain byte. Strings = `varint length + UTF-8`. | 5c | `4e136ca` | `MoqLiteSession.client(...)` (no SETUP), `announce`, `subscribe`, group uni-stream demux, framing helpers, `MoqLiteSessionTest` | | 5d | `41f4dcd` | `connectNestsListener` swap — `MoqLiteNestsListener` adapts `MoqLiteFrame` → `MoqObject` for downstream `AudioRoomPlayer` / `AudioRoomViewModel`. WT URL path = `/?jwt=`. | -**Pending (speaker path — phase 5c-speaker):** +**Speaker path landed (phase 5c-speaker):** -The agent's clarifying lookup confirmed (publisher.rs:40 / connection.js:130) -that moq-lite *publishers* run via `Stream::accept(session)` — the **relay** -opens both Announce and Subscribe bidi streams *to* the publisher. The -publisher only initiates uni streams (one per group of audio data). +| Layer | Surface | +| --------------------------- | ---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | +| `WebTransportSession` | `incomingBidiStreams()` (relay → us) + `openUniStream()` (group push) added. | +| `:quic` `WtPeerStreamDemux` | `StrippedWtStream` carries optional `send` + `finish` closures (driver-aware) so peer-initiated bidis can write back; demux now takes the `QuicConnectionDriver` so wakeups fire after each app-level write. | +| `FakeWebTransport` | now exposes `incomingBidiStreams` + `openUniStream` directly, so production callers and tests use the same API. | +| `MoqLiteSession` | New `publish(suffix)` returns a `MoqLitePublisherHandle`. The session lazily launches a relay→us bidi pump that dispatches by `ControlType`: `Announce` → reply `Active(suffix=ours)`; `Subscribe` → reply `SubscribeOk` and register an inbound subscription. `MoqLitePublisherHandle.send` opens a uni stream per group with `DataType=0` + group header + `varint(size)+payload` frames; `endGroup` FINs the current uni; `close` emits `Announce(Ended)` on every active announce bidi. | +| `MoqLiteNestsSpeaker` | NestsSpeaker adapter that wraps a `MoqLiteSession` + `AudioRoomMoqLiteBroadcaster` (Opus encode + push loop, mirror of the IETF `AudioRoomBroadcaster` but driving the moq-lite publisher). | +| `connectNestsSpeaker` | now opens a `MoqLiteSession` (no SETUP) and returns `MoqLiteNestsSpeaker`. | -That requires `WebTransportSession.acceptBidiStream(): Flow` -which is **not** currently exposed by `:nestsClient`'s WT abstraction -(it has `incomingUniStreams` and `openBidiStream` but no -`incomingBidiStreams`). The underlying `:quic` stack already has -`QuicConnection.awaitIncomingPeerStream` (commonMain:397), so wiring -this through is mechanical — but it's a real API addition and worth a -separate phase. - -Once that lands, the speaker side adds: -- `MoqLiteSession.runPublisher(suffix, onAnnouncePlease, onSubscribe)` - that loops on `acceptBidi` and dispatches by ControlType -- A new `MoqLiteNestsSpeaker` that wraps the session and feeds - Opus frames to one uni-stream-per-group writer - -Then `connectNestsSpeaker` switches the same way `connectNestsListener` -just did, and the existing integration tests (round-trip, multi-peer) -should pass against the real Docker'd nostrnests stack. +The existing integration tests (round-trip, multi-peer, fan-out, +subscribe-before-announce) drive the end-to-end public API — they +should now pass against the real Docker'd nostrnests stack when the +host runs them with `-DnestsInterop=true`. ## Implementation plan (original spec — kept for reference) diff --git a/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/MoqLiteNestsSpeaker.kt b/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/MoqLiteNestsSpeaker.kt new file mode 100644 index 000000000..4e37e980c --- /dev/null +++ b/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/MoqLiteNestsSpeaker.kt @@ -0,0 +1,169 @@ +/* + * 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 + +import com.vitorpamplona.nestsclient.audio.AudioCapture +import com.vitorpamplona.nestsclient.audio.AudioRoomMoqLiteBroadcaster +import com.vitorpamplona.nestsclient.audio.OpusEncoder +import com.vitorpamplona.nestsclient.moq.lite.MoqLitePublisherHandle +import com.vitorpamplona.nestsclient.moq.lite.MoqLiteSession +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.flow.MutableStateFlow +import kotlinx.coroutines.flow.StateFlow +import kotlinx.coroutines.flow.asStateFlow +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock + +/** + * Moq-lite-backed [NestsSpeaker]. Mirrors [MoqLiteNestsListener] on the + * publish side: takes a connected [MoqLiteSession] and exposes the + * existing [NestsSpeaker] API so [connectNestsSpeaker] can swap the + * framing layer without changing any downstream consumers. + * + * Wire-flow per [MoqLiteSession.publish]: + * - the session opens a publisher state when [startBroadcasting] is + * called, then services every relay-opened Announce / Subscribe + * bidi automatically. + * - frames pushed via [MoqLitePublisherHandle.send] go on a fresh + * uni stream per group, framed as `varint(size) + payload`. + */ +class MoqLiteNestsSpeaker internal constructor( + private val session: MoqLiteSession, + private val speakerPubkeyHex: String, + private val captureFactory: () -> AudioCapture, + private val encoderFactory: () -> OpusEncoder, + private val scope: CoroutineScope, + private val mutableState: MutableStateFlow, +) : NestsSpeaker { + override val state: StateFlow = mutableState.asStateFlow() + + private val gate = Mutex() + private var activeHandle: MoqLiteBroadcastHandle? = null + + override suspend fun startBroadcasting(): BroadcastHandle { + gate.withLock { + val current = state.value + check(current is NestsSpeakerState.Connected) { + "startBroadcasting requires Connected state, was $current" + } + check(activeHandle == null) { "speaker is already broadcasting" } + + // Per the audio-rooms NIP draft + JS reference + // (`@moq/publish/screen-B680RFft.js:5641`), publishers + // claim a broadcast suffix equal to their pubkey hex. + val publisher = + try { + session.publish(broadcastSuffix = speakerPubkeyHex) + } catch (t: Throwable) { + throw t + } + val broadcaster = + AudioRoomMoqLiteBroadcaster( + capture = captureFactory(), + encoder = encoderFactory(), + publisher = publisher, + scope = scope, + ) + broadcaster.start() + mutableState.value = + NestsSpeakerState.Broadcasting( + room = current.room, + negotiatedMoqVersion = current.negotiatedMoqVersion, + isMuted = false, + ) + val handle = + MoqLiteBroadcastHandle( + broadcaster = broadcaster, + publisher = publisher, + parent = this, + ) + activeHandle = handle + return handle + } + } + + /** + * Compare-and-clear that runs from inside [close] (already holds + * [gate]) and from [MoqLiteBroadcastHandle.close] (doesn't). + * Mirrors [DefaultNestsSpeaker.broadcastClosed]. + */ + internal fun broadcastClosed(handle: MoqLiteBroadcastHandle) { + if (activeHandle !== handle) return + activeHandle = null + val current = mutableState.value + if (current is NestsSpeakerState.Broadcasting) { + mutableState.value = + NestsSpeakerState.Connected(current.room, current.negotiatedMoqVersion) + } + } + + internal fun reportMuteState(muted: Boolean) { + val current = mutableState.value + if (current is NestsSpeakerState.Broadcasting) { + mutableState.value = current.copy(isMuted = muted) + } + } + + override suspend fun close() { + // Take + clear under [gate] so a concurrent `startBroadcasting` + // can't observe a half-closed state, then run the long-running + // suspends (handle.close + session.close) outside the lock. + val handle: MoqLiteBroadcastHandle? + gate.withLock { + if (state.value is NestsSpeakerState.Closed) return + handle = activeHandle + activeHandle = null + mutableState.value = NestsSpeakerState.Closed + } + handle?.runCatching { close() } + runCatching { session.close() } + } +} + +internal class MoqLiteBroadcastHandle( + private val broadcaster: AudioRoomMoqLiteBroadcaster, + private val publisher: MoqLitePublisherHandle, + private val parent: MoqLiteNestsSpeaker, +) : BroadcastHandle { + @Volatile private var muted: Boolean = false + + @Volatile private var closed: Boolean = false + + override val isMuted: Boolean get() = muted + + override suspend fun setMuted(muted: Boolean) { + if (closed) return + this.muted = muted + broadcaster.setMuted(muted) + parent.reportMuteState(muted) + } + + override suspend fun close() { + if (closed) return + closed = true + runCatching { broadcaster.stop() } + // broadcaster.stop() already calls publisher.close(); call again + // defensively to make this method idempotent against partial + // failures on the broadcaster.stop path. + runCatching { publisher.close() } + parent.broadcastClosed(this) + } +} diff --git a/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/NestsConnect.kt b/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/NestsConnect.kt index 773656e8d..34a71ddb5 100644 --- a/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/NestsConnect.kt +++ b/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/NestsConnect.kt @@ -22,10 +22,8 @@ package com.vitorpamplona.nestsclient import com.vitorpamplona.nestsclient.audio.AudioCapture import com.vitorpamplona.nestsclient.audio.OpusEncoder -import com.vitorpamplona.nestsclient.moq.MoqSession import com.vitorpamplona.nestsclient.moq.MoqVersion import com.vitorpamplona.nestsclient.moq.SubscribeHandle -import com.vitorpamplona.nestsclient.moq.TrackNamespace import com.vitorpamplona.nestsclient.transport.WebTransportException import com.vitorpamplona.nestsclient.transport.WebTransportFactory import com.vitorpamplona.quartz.nip01Core.signers.NostrSigner @@ -163,12 +161,16 @@ private fun failedListener(state: MutableStateFlow): NestsLi /** * Speaker / host counterpart of [connectNestsListener]. Walks the same - * HTTP → WebTransport → MoQ handshake; the difference is the post-setup - * step is `announce(...)` (driven by [NestsSpeaker.startBroadcasting]) - * instead of `subscribe(...)`. + * three-step HTTP → WebTransport → moq-lite session sequence the + * listener does; the difference is downstream: [NestsSpeaker.startBroadcasting] + * claims a broadcast suffix on the moq-lite session and starts pumping + * Opus frames out as one uni stream per group. * - * @param speakerPubkeyHex this user's pubkey hex, used as the MoQ track - * name when we ANNOUNCE — listeners look us up by exactly that name. + * @param speakerPubkeyHex this user's pubkey hex. Used as the moq-lite + * broadcast suffix the relay routes to subscribers + * (`MoqLiteSubscribe.broadcast == speakerPubkeyHex`); the JS reference + * mints the same value via `Path.from(identity)` + * (`@moq/publish/screen-B680RFft.js`). * @param captureFactory builds an [AudioCapture] (one per broadcast). * Android passes `{ AudioRecordCapture() }`. * @param encoderFactory builds an [OpusEncoder] (one per broadcast). @@ -228,28 +230,23 @@ suspend fun connectNestsSpeaker( state.value = NestsSpeakerState.Connecting(NestsSpeakerState.Connecting.ConnectStep.MoqHandshake) + // moq-lite Lite-03 has NO setup message. Same logic as the listener + // path — `supportedMoqVersions` retained for backward compat but + // currently a no-op because version is selected by ALPN. val moq = try { - MoqSession.client(webTransport, scope).also { it.setup(supportedMoqVersions) } + com.vitorpamplona.nestsclient.moq.lite.MoqLiteSession + .client(webTransport, scope) } catch (t: Throwable) { - runCatching { webTransport.close(0, "moq setup failed") } - state.value = NestsSpeakerState.Failed("MoQ handshake failed: ${t.message}", t) + runCatching { webTransport.close(0, "moq-lite session init failed") } + state.value = NestsSpeakerState.Failed("moq-lite session init failed: ${t.message}", t) return failedSpeaker(state) } - val negotiatedVersion = - moq.selectedVersion ?: run { - runCatching { moq.close() } - state.value = NestsSpeakerState.Failed("MoQ session reported no negotiated version") - return failedSpeaker(state) - } - - state.value = NestsSpeakerState.Connected(room, negotiatedVersion) - return DefaultNestsSpeaker( + state.value = NestsSpeakerState.Connected(room, MOQ_LITE_03_VERSION) + return MoqLiteNestsSpeaker( session = moq, - // Same single-segment shape as the listener path; see comment there. - roomNamespace = TrackNamespace.of(room.moqNamespace()), - speakerTrackName = speakerPubkeyHex.encodeToByteArray(), + speakerPubkeyHex = speakerPubkeyHex, captureFactory = captureFactory, encoderFactory = encoderFactory, scope = scope, diff --git a/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/audio/AudioRoomMoqLiteBroadcaster.kt b/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/audio/AudioRoomMoqLiteBroadcaster.kt new file mode 100644 index 000000000..e08e51c2e --- /dev/null +++ b/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/audio/AudioRoomMoqLiteBroadcaster.kt @@ -0,0 +1,140 @@ +/* + * 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.audio + +import com.vitorpamplona.nestsclient.moq.lite.MoqLitePublisherHandle +import kotlinx.coroutines.CancellationException +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Job +import kotlinx.coroutines.cancelAndJoin +import kotlinx.coroutines.launch + +/** + * Mirror of [AudioRoomBroadcaster] but driving a moq-lite + * [MoqLitePublisherHandle] instead of an IETF `MoqSession.TrackPublisher`. + * Keeps the IETF broadcaster intact for the IETF unit-test suite while + * letting the production speaker path use moq-lite. + * + * Lifecycle and audit comments mirror the IETF version 1:1 — only the + * sink type changes. + */ +class AudioRoomMoqLiteBroadcaster( + private val capture: AudioCapture, + private val encoder: OpusEncoder, + private val publisher: MoqLitePublisherHandle, + private val scope: CoroutineScope, +) { + private var job: Job? = null + + @Volatile private var stopped: Boolean = false + + @Volatile private var muted: Boolean = false + + /** + * Start capturing + encoding + publishing in the background. + * Returns immediately. Calling twice is an error. If + * [AudioCapture.start] throws, the broadcaster is left stopped and + * the exception propagates. + */ + fun start(onError: (AudioException) -> Unit = { /* swallow */ }) { + check(!stopped) { "AudioRoomMoqLiteBroadcaster already stopped" } + check(job == null) { "AudioRoomMoqLiteBroadcaster.start already called" } + + try { + capture.start() + } catch (t: Throwable) { + stopped = true + runCatching { capture.stop() } + throw t + } + job = + scope.launch { + try { + while (true) { + val pcm = capture.readFrame() ?: break + val opus = + try { + encoder.encode(pcm) + } catch (ce: CancellationException) { + throw ce + } catch (t: Throwable) { + onError( + AudioException( + AudioException.Kind.DecoderError, + "Opus encode failed for a frame", + t, + ), + ) + continue + } + if (opus.isEmpty()) continue + if (muted) continue + runCatching { publisher.send(opus) } + .onFailure { t -> + if (t is CancellationException) throw t + onError( + AudioException( + AudioException.Kind.PlaybackFailed, + "publisher.send failed", + t, + ), + ) + } + } + } catch (ce: CancellationException) { + throw ce + } catch (t: Throwable) { + onError( + AudioException( + AudioException.Kind.PlaybackFailed, + "audio capture pipeline failed", + t, + ), + ) + } + } + } + + /** + * Toggle whether captured frames reach the wire. Mic stays open and + * encoder keeps state consistent so unmute is sample-accurate. + */ + fun setMuted(muted: Boolean) { + if (stopped) return + this.muted = muted + } + + /** + * Stop the loop, release the mic, release the encoder, close the + * moq-lite publisher (which sends `Announce(Ended)` on every active + * announce bidi). Idempotent. We `cancelAndJoin` the loop before + * releasing the encoder + publisher so the loop's last + * `encoder.encode` / `publisher.send` can't race the close path. + */ + suspend fun stop() { + if (stopped) return + stopped = true + job?.cancelAndJoin() + runCatching { capture.stop() } + runCatching { encoder.release() } + runCatching { publisher.close() } + } +} diff --git a/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/moq/lite/MoqLiteSession.kt b/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/moq/lite/MoqLiteSession.kt index 00b0fdffc..de801f335 100644 --- a/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/moq/lite/MoqLiteSession.kt +++ b/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/moq/lite/MoqLiteSession.kt @@ -38,9 +38,10 @@ import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.withLock /** - * Listener-side moq-lite (Lite-03) session. Wraps a connected - * [WebTransportSession] and exposes: + * moq-lite (Lite-03) session for both listener and publisher roles. + * Wraps a connected [WebTransportSession] and exposes: * + * Listener side: * - [announce] — open an Announce bidi with a prefix, observe live * `Active` / `Ended` updates from the relay. * - [subscribe] — open a Subscribe bidi for a `(broadcast, track)` @@ -48,14 +49,20 @@ import kotlinx.coroutines.sync.withLock * [MoqLiteSubscribeHandle.frames] yields each frame the publisher * pushes. * - * Speaker-side is **not yet implemented** — the publisher direction - * needs server-initiated bidi acceptance (relay → us) which the - * current `WebTransportSession` interface does not surface. Tracked in - * `nestsClient/plans/2026-04-26-moq-lite-gap.md` phase-5c-speaker. + * Publisher side: + * - [publish] — claim a broadcast suffix; the session then services + * every relay-opened Announce / Subscribe bidi automatically and + * returns a [MoqLitePublisherHandle] the application can push + * Opus frames into. Group rollover is the application's call — + * [MoqLitePublisherHandle.send] auto-starts a group on first call, + * [MoqLitePublisherHandle.endGroup] FINs and starts a new one on + * the next send. * * Wire-protocol scope: Lite-03 — no SETUP, no datagrams, one fresh - * client-initiated bidi per request, group data on uni streams. See - * the gap doc for the byte-level layout. + * bidi per request (subscriber → publisher OR publisher accepts from + * relay), group data on uni streams. See + * `nestsClient/plans/2026-04-26-moq-lite-gap.md` for the byte-level + * layout. */ class MoqLiteSession internal constructor( private val transport: WebTransportSession, @@ -66,6 +73,12 @@ class MoqLiteSession internal constructor( private var nextSubscribeId: Long = 0L private var groupPump: Job? = null + /** Lazily-launched relay→us inbound bidi pump; only runs while a publisher is active. */ + private var bidiPump: Job? = null + + /** Single active publisher per session (moq-lite doesn't model multi-broadcast publishers). */ + private var activePublisher: PublisherStateImpl? = null + @Volatile private var closed: Boolean = false val isClosed: Boolean get() = closed @@ -277,19 +290,222 @@ class MoqLiteSession internal constructor( sub.frames.close() } + // ==================================================================== + // Publisher side + // ==================================================================== + + /** + * Begin publishing under [broadcastSuffix]. Returns a + * [MoqLitePublisherHandle] that: + * - announces `Active` to every relay-opened Announce bidi whose + * `AnnouncePlease.prefix` matches our suffix + * - opens a fresh group uni stream when a relay-opened Subscribe + * bidi for our broadcast arrives, then forwards [MoqLitePublisherHandle.send] + * bytes as `varint(size) + payload` frames until [MoqLitePublisherHandle.endGroup] + * or [MoqLitePublisherHandle.close] + * + * Wire-flow per `kixelated/moq-rs/rs/moq-lite/src/lite/publisher.rs:40+`: + * - the relay opens Announce + Subscribe bidi streams *to* us + * (`Stream::accept(session)`), so the session's + * [WebTransportSession.incomingBidiStreams] pump is the entry + * point + * - we open uni streams ourselves to push group data + * (`session.open_uni()`) + * + * Only one [publish] is supported per session for now. Calling + * [publish] twice on the same session is rejected with [IllegalStateException]. + */ + suspend fun publish(broadcastSuffix: String): MoqLitePublisherHandle { + ensureOpen() + val normalised = MoqLitePath.normalize(broadcastSuffix) + val publisher: PublisherStateImpl + state.withLock { + check(!closed) { "session is closed" } + check(activePublisher == null) { + "MoqLiteSession.publish called twice — only one broadcast per session is supported" + } + publisher = PublisherStateImpl(suffix = normalised) + activePublisher = publisher + // Lazy launch — the inbound-bidi pump needs to keep running + // for the lifetime of any active publisher. + if (bidiPump == null) bidiPump = scope.launch { pumpInboundBidis() } + } + return publisher + } + + /** + * Drain inbound bidi streams (relay → us) and dispatch each by + * its leading [MoqLiteControlType] varint. The relay opens + * Announce / Subscribe bidis on its own initiative; we read the + * control type, the request body, and reply on the same bidi. + * + * One pump per session — started lazily on the first [publish]. + */ + private suspend fun pumpInboundBidis() { + try { + transport.incomingBidiStreams().collect { bidi -> + scope.launch { handleInboundBidi(bidi) } + } + } catch (ce: CancellationException) { + throw ce + } catch (_: Throwable) { + // Transport closed. + } + } + + private suspend fun handleInboundBidi(bidi: com.vitorpamplona.nestsclient.transport.WebTransportBidiStream) { + val buffer = MoqLiteFrameBuffer() + val publisher = state.withLock { activePublisher } ?: return + try { + // Read the leading ControlType varint from the first chunk. + val first = + bidi.incoming().firstOrNull() ?: return + buffer.push(first) + val controlCode = buffer.readVarint() ?: return + val controlType = MoqLiteControlType.fromCode(controlCode) ?: return + when (controlType) { + MoqLiteControlType.Announce -> { + handleAnnounceRequest(bidi, buffer, publisher) + } + + MoqLiteControlType.Subscribe -> { + handleSubscribeRequest(bidi, buffer, publisher) + } + + else -> { + // Lite-03 treats Session/Fetch/Probe as separate flows; + // we don't implement them here. Drop the bidi. + runCatching { bidi.finish() } + } + } + } catch (ce: CancellationException) { + throw ce + } catch (_: Throwable) { + runCatching { bidi.finish() } + } + } + + private suspend fun handleAnnounceRequest( + bidi: com.vitorpamplona.nestsclient.transport.WebTransportBidiStream, + seedBuffer: MoqLiteFrameBuffer, + publisher: PublisherStateImpl, + ) { + val pleasePayload = readSizePrefixedFromBidiInto(bidi.incoming(), seedBuffer) + val please = MoqLiteCodec.decodeAnnouncePlease(pleasePayload) + // The relay sets the prefix to the namespace it expects us to + // publish under (typically `claims.root`). Our broadcast path + // (after stripping the prefix) is `publisher.suffix`. moq-lite + // requires the suffix on the wire to be the *remaining* part + // after `please.prefix` — so strip it. + val emittedSuffix = MoqLitePath.stripPrefix(please.prefix, publisher.suffix) ?: publisher.suffix + bidi.write( + MoqLiteCodec.encodeAnnounce( + MoqLiteAnnounce( + status = MoqLiteAnnounceStatus.Active, + suffix = emittedSuffix, + hops = 0L, + ), + ), + ) + // Hold the bidi open until the publisher closes; if/when the + // application stops broadcasting, send `Ended`. + publisher.registerAnnounceBidi(bidi, emittedSuffix) + } + + private suspend fun handleSubscribeRequest( + bidi: com.vitorpamplona.nestsclient.transport.WebTransportBidiStream, + seedBuffer: MoqLiteFrameBuffer, + publisher: PublisherStateImpl, + ) { + val subPayload = readSizePrefixedFromBidiInto(bidi.incoming(), seedBuffer) + val sub = MoqLiteCodec.decodeSubscribe(subPayload) + // Reply Ok right away — moq-lite is permissive on the publisher + // side; the relay decides whether the subscriber is allowed to + // see this broadcast. + bidi.write( + MoqLiteCodec.encodeSubscribeOk( + MoqLiteSubscribeOk( + priority = sub.priority, + ordered = sub.ordered, + maxLatencyMillis = sub.maxLatencyMillis, + startGroup = null, + endGroup = null, + ), + ), + ) + publisher.registerInboundSubscription(sub) + } + + /** + * Open a new uni stream for one group's frames and push the + * `DataType=Group` byte + size-prefixed [MoqLiteGroupHeader]. + * Returns a write stream the caller frames each Opus packet onto. + */ + internal suspend fun openGroupStream( + subscribeId: Long, + sequence: Long, + ): com.vitorpamplona.nestsclient.transport.WebTransportWriteStream { + val uni = transport.openUniStream() + uni.write(Varint.encode(MoqLiteDataType.Group.code)) + uni.write(MoqLiteCodec.encodeGroupHeader(MoqLiteGroupHeader(subscribeId, sequence))) + return uni + } + + /** + * Read a size-prefixed payload from a bidi, seeded with whatever's + * already in [buffer] (the ControlType byte may have arrived with + * extra bytes). Used internally by the publisher inbound dispatch. + */ + private suspend fun readSizePrefixedFromBidiInto( + incoming: kotlinx.coroutines.flow.Flow, + buffer: MoqLiteFrameBuffer, + ): ByteArray { + // Try the buffer first — first chunk often contains the whole + // body since moq-lite messages are small and arrive as single + // QUIC sends. + buffer.readSizePrefixed()?.let { return it } + var done: ByteArray? = null + try { + incoming.collect { chunk -> + buffer.push(chunk) + buffer.readSizePrefixed()?.let { + done = it + throw EarlyExit + } + } + } catch (_: EarlyExit) { + // expected + } catch (ce: CancellationException) { + throw ce + } + return done + ?: throw MoqCodecException("incoming bidi closed before a complete size-prefixed body arrived") + } + + private object EarlyExit : RuntimeException() { + private fun readResolve(): Any = EarlyExit + + override fun fillInStackTrace(): Throwable = this + } + suspend fun close() { if (closed) return closed = true val toClose: List + val publisherToClose: PublisherStateImpl? state.withLock { toClose = subscriptionsBySubscribeId.values.toList() subscriptionsBySubscribeId.clear() + publisherToClose = activePublisher + activePublisher = null } for (sub in toClose) { runCatching { sub.bidi.finish() } sub.frames.close() } + runCatching { publisherToClose?.close() } groupPump?.cancelAndJoin() + bidiPump?.cancelAndJoin() runCatching { transport.close() } } @@ -331,6 +547,132 @@ class MoqLiteSession internal constructor( val frames: Channel, ) + /** + * Publisher state. Tracks the announce bidis the relay opened to us + * + the inbound subscriptions a relay (or peer) opened against our + * broadcast, and owns the current group's uni stream. + * + * `gate` serialises access to per-group state so concurrent + * `send` / `startGroup` / `endGroup` / `close` can't race. + */ + private inner class PublisherStateImpl( + override val suffix: String, + ) : MoqLitePublisherHandle { + private val gate = Mutex() + private val announceBidis = mutableListOf() + private val inboundSubs = mutableListOf() + private var currentGroup: GroupOutbound? = null + private var nextSequence: Long = 0L + + @Volatile private var publisherClosed = false + + suspend fun registerAnnounceBidi( + bidi: com.vitorpamplona.nestsclient.transport.WebTransportBidiStream, + emittedSuffix: String, + ) { + gate.withLock { + if (publisherClosed) { + runCatching { bidi.finish() } + return + } + announceBidis += AnnounceBidiEntry(bidi, emittedSuffix) + } + } + + suspend fun registerInboundSubscription(sub: MoqLiteSubscribe) { + gate.withLock { + if (publisherClosed) return + inboundSubs += sub + } + } + + override suspend fun startGroup() { + gate.withLock { + if (publisherClosed) return + runCatching { currentGroup?.uni?.finish() } + currentGroup = openNextGroupLocked() + } + } + + override suspend fun send(payload: ByteArray): Boolean { + gate.withLock { + if (publisherClosed) return false + if (inboundSubs.isEmpty()) return false + val group = currentGroup ?: openNextGroupLocked().also { currentGroup = it } + val framed = Varint.encode(payload.size.toLong()) + payload + runCatching { group.uni.write(framed) } + } + return true + } + + override suspend fun endGroup() { + gate.withLock { + if (publisherClosed) return + val group = currentGroup ?: return + currentGroup = null + runCatching { group.uni.finish() } + } + } + + override suspend fun close() { + val toFinalise: List + val groupToFinish: GroupOutbound? + gate.withLock { + if (publisherClosed) return + publisherClosed = true + toFinalise = announceBidis.toList() + announceBidis.clear() + inboundSubs.clear() + groupToFinish = currentGroup + currentGroup = null + } + for (entry in toFinalise) { + runCatching { + entry.bidi.write( + MoqLiteCodec.encodeAnnounce( + MoqLiteAnnounce( + status = MoqLiteAnnounceStatus.Ended, + suffix = entry.emittedSuffix, + hops = 0L, + ), + ), + ) + } + runCatching { entry.bidi.finish() } + } + runCatching { groupToFinish?.uni?.finish() } + // Detach from the session so a subsequent `publish` can run. + state.withLock { + if (activePublisher === this) activePublisher = null + } + } + + /** Caller holds [gate]. */ + private suspend fun openNextGroupLocked(): GroupOutbound { + // moq-lite groups are addressed by `subscribeId` on the wire — + // each inbound subscription gets its own group stream. For + // simplicity we open one stream per group keyed off the + // *first* subscription's id; relay-side multi-subscriber + // fan-out happens above us. Inbound subscription set is + // expected to be small (1 in nests's listener-per-room + // model), so this is fine. + val sub = inboundSubs.first() + val sequence = nextSequence++ + val uni = openGroupStream(subscribeId = sub.id, sequence = sequence) + return GroupOutbound(sequence = sequence, uni = uni) + } + } + + private data class AnnounceBidiEntry( + val bidi: com.vitorpamplona.nestsclient.transport.WebTransportBidiStream, + val emittedSuffix: String, + ) + + private data class GroupOutbound( + val sequence: Long, + val uni: com.vitorpamplona.nestsclient.transport.WebTransportWriteStream, + ) + companion object { /** moq-lite priority byte midpoint — neutral default. */ const val DEFAULT_PRIORITY: Int = 0x80 @@ -408,3 +750,54 @@ class MoqLiteSubscribeException( message: String, cause: Throwable? = null, ) : RuntimeException(message, cause) + +/** + * Active publisher handle returned by [MoqLiteSession.publish]. + * + * Lifecycle: + * 1. Call [startGroup] (or [send] which auto-starts a fresh group on + * first call) to begin pushing frames for one Opus group. + * 2. Call [send] for each frame (one Opus packet = one frame). + * 3. Call [endGroup] to FIN the current group's uni stream and start + * a fresh group on the next [send]. Group rollover is the + * publisher's call — typically every N seconds or every keyframe. + * 4. Call [close] when the broadcast ends — sends `Announce(Ended)` + * on every active announce bidi and FINs every group stream. + */ +interface MoqLitePublisherHandle { + /** + * The broadcast suffix this publisher claimed at [MoqLiteSession.publish]. + * Always normalised per [MoqLitePath]. + */ + val suffix: String + + /** + * Start a new group. Allocates a fresh sequence id and opens a new + * uni stream pre-loaded with `DataType=Group + GroupHeader`. Idempotent + * — calling [startGroup] when the previous group hasn't been ended + * is treated as an implicit [endGroup] then a new start. + */ + suspend fun startGroup() + + /** + * Push one [payload] (one Opus packet) as a `varint(size) + payload` + * frame on the current group's uni stream. Auto-starts a group if + * none is active. + * + * Returns false if no inbound subscriber is currently attached. + * Subscriber-less sends silently drop on the wire — the relay keeps + * the publisher's announce active either way, so unmute is + * sample-accurate. + */ + suspend fun send(payload: ByteArray): Boolean + + /** FIN the current group's uni stream. The next [send] starts a fresh group. */ + suspend fun endGroup() + + /** + * Stop publishing. Sends `Announce(Ended)` on every active announce + * bidi, FINs the current group, and releases all per-publisher + * resources. Idempotent. + */ + suspend fun close() +} diff --git a/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/transport/FakeWebTransport.kt b/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/transport/FakeWebTransport.kt index 6cc32c967..3e33c3aa5 100644 --- a/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/transport/FakeWebTransport.kt +++ b/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/transport/FakeWebTransport.kt @@ -66,6 +66,20 @@ class FakeWebTransport private constructor( override fun incomingUniStreams(): Flow = inboundUniStreams.receiveAsFlow() + override fun incomingBidiStreams(): Flow = inboundBidiStreams.receiveAsFlow() + + /** + * Open a uni stream toward the paired peer — production flow used by + * the moq-lite publisher path to push group data. The peer side + * receives the new stream via [incomingUniStreams]. + */ + override suspend fun openUniStream(): WebTransportWriteStream { + stateLock.withLock { check(open) { "session closed" } } + val pipe = Channel(Channel.BUFFERED) + outboundUniStreams.send(FakeReadStream(pipe)) + return ChannelWriteStream(pipe) + } + override suspend fun sendDatagram(payload: ByteArray): Boolean { if (!open) return false outboundDatagrams.send(payload) @@ -96,19 +110,6 @@ class FakeWebTransport private constructor( */ fun peerOpenedBidiStreams(): Flow = inboundBidiStreams.receiveAsFlow() - /** - * Open a uni stream from this side toward the paired peer. Tests use this - * to simulate a publisher pushing a moq-lite group uni stream — write the - * full sequence of chunks via [WebTransportWriteStream.write] then call - * [WebTransportWriteStream.finish] to FIN. - */ - suspend fun openPeerUniStream(): WebTransportWriteStream { - stateLock.withLock { check(open) { "session closed" } } - val pipe = Channel(Channel.BUFFERED) - outboundUniStreams.send(FakeReadStream(pipe)) - return ChannelWriteStream(pipe) - } - companion object { /** * Create two linked fakes that act as "client" and "server" endpoints of @@ -168,9 +169,10 @@ class FakeReadStream internal constructor( } /** - * Write-only adapter over a [Channel]. Used by [FakeWebTransport.openPeerUniStream] - * so a test can drive a peer-initiated uni stream by writing chunks then - * FIN'ing via [finish]. + * Write-only adapter over a [Channel]. Backs both + * [FakeWebTransport.openUniStream] (production: locally-opened uni + * stream that the paired peer reads via incomingUniStreams) and any + * test that wants to drive uni-stream bytes through a known channel. */ private class ChannelWriteStream( private val channel: Channel, diff --git a/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/transport/WebTransportSession.kt b/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/transport/WebTransportSession.kt index 3a42e17a0..32ae06d25 100644 --- a/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/transport/WebTransportSession.kt +++ b/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/transport/WebTransportSession.kt @@ -48,6 +48,16 @@ interface WebTransportSession { */ suspend fun openBidiStream(): WebTransportBidiStream + /** + * Open a new client-initiated unidirectional WebTransport stream. + * + * Required by the moq-lite (Lite-03) publisher path: each group of + * audio frames is pushed on a fresh uni stream that the publisher + * opens — see `rs/moq-lite/src/lite/publisher.rs:338` + * (`session.open_uni()`). + */ + suspend fun openUniStream(): WebTransportWriteStream + /** * Flow of inbound unidirectional streams initiated by the peer. * @@ -55,6 +65,17 @@ interface WebTransportSession { */ fun incomingUniStreams(): Flow + /** + * Flow of inbound bidirectional streams initiated by the peer. + * + * Required by the moq-lite (Lite-03) publisher path: kixelated/moq-rs + * relays open Announce and Subscribe bidi streams *to* the publisher + * — see `rs/moq-lite/src/lite/publisher.rs:40` (`Stream::accept(session)`). + * The flow completes when [close] is called or the peer tears down + * the session. + */ + fun incomingBidiStreams(): Flow + /** Send a QUIC datagram; returns false if the datagram was dropped by congestion control. */ suspend fun sendDatagram(payload: ByteArray): Boolean diff --git a/nestsClient/src/commonTest/kotlin/com/vitorpamplona/nestsclient/moq/lite/MoqLiteSessionTest.kt b/nestsClient/src/commonTest/kotlin/com/vitorpamplona/nestsclient/moq/lite/MoqLiteSessionTest.kt index 3542aee09..09ff40cd0 100644 --- a/nestsClient/src/commonTest/kotlin/com/vitorpamplona/nestsclient/moq/lite/MoqLiteSessionTest.kt +++ b/nestsClient/src/commonTest/kotlin/com/vitorpamplona/nestsclient/moq/lite/MoqLiteSessionTest.kt @@ -37,6 +37,7 @@ import kotlin.test.Test import kotlin.test.assertContentEquals import kotlin.test.assertEquals import kotlin.test.assertFailsWith +import kotlin.test.assertIs /** * Drives a [MoqLiteSession] from the listener side against a fake peer @@ -99,7 +100,7 @@ class MoqLiteSessionTest { assertEquals(0L, handle.id, "first subscribe id is 0") // Now push one group with two frames from the server side. - val uni = serverSide.openPeerUniStream() + val uni = serverSide.openUniStream() uni.write(Varint.encode(MoqLiteDataType.Group.code)) uni.write(MoqLiteCodec.encodeGroupHeader(MoqLiteGroupHeader(subscribeId = handle.id, sequence = 7L))) uni.write(framePayload(byteArrayOf(0x10, 0x11))) @@ -226,13 +227,13 @@ class MoqLiteSessionTest { assertEquals(1L, handleB.id) // Push one group for A with payload "a", one for B with payload "b". - val uniB = serverSide.openPeerUniStream() + val uniB = serverSide.openUniStream() uniB.write(Varint.encode(MoqLiteDataType.Group.code)) uniB.write(MoqLiteCodec.encodeGroupHeader(MoqLiteGroupHeader(subscribeId = handleB.id, sequence = 0L))) uniB.write(framePayload("b".encodeToByteArray())) uniB.finish() - val uniA = serverSide.openPeerUniStream() + val uniA = serverSide.openUniStream() uniA.write(Varint.encode(MoqLiteDataType.Group.code)) uniA.write(MoqLiteCodec.encodeGroupHeader(MoqLiteGroupHeader(subscribeId = handleA.id, sequence = 0L))) uniA.write(framePayload("a".encodeToByteArray())) @@ -248,6 +249,139 @@ class MoqLiteSessionTest { session.close() } + @Test + fun publisher_replies_to_announcePlease_with_active_announce() = + runBlocking { + val (clientSide, serverSide) = FakeWebTransport.pair() + val session = MoqLiteSession.client(clientSide, pumpScope) + + val publisher = session.publish(broadcastSuffix = "speakerPubkey") + + // Relay (serverSide) opens an Announce bidi to us with + // AnnouncePlease(prefix=""). + val relayBidi = serverSide.openBidiStream() + relayBidi.write(Varint.encode(MoqLiteControlType.Announce.code)) + relayBidi.write(MoqLiteCodec.encodeAnnouncePlease(MoqLiteAnnouncePlease(prefix = ""))) + + // We reply on the same bidi with Announce(active=true, + // suffix="speakerPubkey"). + val resp = withTimeout(2_000) { relayBidi.incoming().first() } + val announce = MoqLiteCodec.decodeAnnounce(MoqLiteFrameBuffer().apply { push(resp) }.readSizePrefixed()!!) + assertEquals(MoqLiteAnnounceStatus.Active, announce.status) + assertEquals("speakerPubkey", announce.suffix) + + publisher.close() + session.close() + } + + @Test + fun publisher_acks_subscribe_and_pushes_group_data_on_uni_stream() = + runBlocking { + val (clientSide, serverSide) = FakeWebTransport.pair() + val session = MoqLiteSession.client(clientSide, pumpScope) + + val publisher = session.publish(broadcastSuffix = "speakerPubkey") + + // Step 1: relay opens Subscribe bidi. + val subBidi = serverSide.openBidiStream() + subBidi.write(Varint.encode(MoqLiteControlType.Subscribe.code)) + subBidi.write( + MoqLiteCodec.encodeSubscribe( + MoqLiteSubscribe( + id = 7L, + broadcast = "speakerPubkey", + track = "audio/data", + priority = 0x80, + ordered = true, + maxLatencyMillis = 0L, + startGroup = null, + endGroup = null, + ), + ), + ) + + // Step 2: we reply SubscribeOk. + val ackChunk = withTimeout(2_000) { subBidi.incoming().first() } + val resp = + MoqLiteCodec.decodeSubscribeResponse( + MoqLiteFrameBuffer().apply { push(ackChunk) }.readSizePrefixed()!!, + ) + assertIs(resp) + + // Step 3: publisher pushes one frame, which opens a uni + // stream with DataType=Group + group header + frame. + assertEquals(true, publisher.send("opus-1".encodeToByteArray())) + publisher.endGroup() + + // The uni stream surfaces on the relay side via incomingUniStreams. + val relayUni = withTimeout(2_000) { serverSide.incomingUniStreams().first() } + val uniChunks = relayUni.incoming().toList() + // Concatenate all chunks then parse: type + group header + + // first frame. The buffer reader handles arbitrary chunk + // boundaries. + val buf = MoqLiteFrameBuffer() + uniChunks.forEach { buf.push(it) } + assertEquals(MoqLiteDataType.Group.code, buf.readVarint(), "uni stream starts with Group type byte") + val header = + MoqLiteCodec.decodeGroupHeader( + buf.readSizePrefixed() ?: error("group header missing"), + ) + assertEquals(7L, header.subscribeId) + assertEquals(0L, header.sequence, "first group is sequence 0") + val firstFrame = + buf.readSizePrefixed() + ?: error("first frame missing") + assertContentEquals("opus-1".encodeToByteArray(), firstFrame) + + publisher.close() + session.close() + } + + @Test + fun publisher_send_returns_false_when_no_inbound_subscriber() = + runBlocking { + val (clientSide, _) = FakeWebTransport.pair() + val session = MoqLiteSession.client(clientSide, pumpScope) + + val publisher = session.publish(broadcastSuffix = "speakerPubkey") + // No relay-opened Subscribe bidi → no subscribers → send is + // a silent no-op (returns false), matching the listener + // semantics where the speaker keeps capturing audio even + // when nobody is listening. + assertEquals(false, publisher.send("ignored".encodeToByteArray())) + + publisher.close() + session.close() + } + + @Test + fun publisher_close_emits_ended_announce() = + runBlocking { + val (clientSide, serverSide) = FakeWebTransport.pair() + val session = MoqLiteSession.client(clientSide, pumpScope) + + val publisher = session.publish(broadcastSuffix = "speakerPubkey") + + // Relay opens an announce bidi. + val relayBidi = serverSide.openBidiStream() + relayBidi.write(Varint.encode(MoqLiteControlType.Announce.code)) + relayBidi.write(MoqLiteCodec.encodeAnnouncePlease(MoqLiteAnnouncePlease(prefix = ""))) + // Drain the Active announce so the next .first() picks up Ended. + withTimeout(2_000) { relayBidi.incoming().first() } + + publisher.close() + + val endedChunk = withTimeout(2_000) { relayBidi.incoming().first() } + val ended = + MoqLiteCodec.decodeAnnounce( + MoqLiteFrameBuffer().apply { push(endedChunk) }.readSizePrefixed()!!, + ) + assertEquals(MoqLiteAnnounceStatus.Ended, ended.status) + assertEquals("speakerPubkey", ended.suffix) + + session.close() + } + @Test fun unsubscribe_FINs_the_subscribe_bidi() = runBlocking { diff --git a/nestsClient/src/jvmAndroid/kotlin/com/vitorpamplona/nestsclient/transport/QuicWebTransportFactory.kt b/nestsClient/src/jvmAndroid/kotlin/com/vitorpamplona/nestsclient/transport/QuicWebTransportFactory.kt index c40910a57..41a5e017f 100644 --- a/nestsClient/src/jvmAndroid/kotlin/com/vitorpamplona/nestsclient/transport/QuicWebTransportFactory.kt +++ b/nestsClient/src/jvmAndroid/kotlin/com/vitorpamplona/nestsclient/transport/QuicWebTransportFactory.kt @@ -234,6 +234,11 @@ class QuicWebTransportSession( return QuicBidiStreamAdapter(s, state.driver) } + override suspend fun openUniStream(): WebTransportWriteStream { + val s = state.openUniStream() + return QuicUniWriteStreamAdapter(s, state.driver) + } + override fun incomingUniStreams(): Flow = flow { // Surface only unidirectional WT streams whose prefix bytes @@ -247,6 +252,18 @@ class QuicWebTransportSession( } } + override fun incomingBidiStreams(): Flow = + flow { + // Surface only peer-initiated bidi streams whose WT_BIDI_STREAM + // prefix has been stripped by the demux. send/finish are + // wired through the driver wakeup. + state.incomingStrippedStreams.collect { stripped -> + if (!stripped.isUnidirectional) { + emit(StrippedWtBidiStreamAdapter(stripped)) + } + } + } + override suspend fun sendDatagram(payload: ByteArray): Boolean { state.sendDatagram(payload) return true @@ -292,9 +309,62 @@ private class QuicReadStreamAdapter( override fun incoming(): Flow = stream.incoming } +/** + * Write-only adapter over a locally-opened uni QUIC stream. The + * underlying [QuicWebTransportSessionState.openUniStream] has already + * pushed the WT framing prefix (0x54 + quarter session id), so the + * caller's [write] payload goes straight onto the wire. + */ +private class QuicUniWriteStreamAdapter( + private val stream: QuicStream, + private val driver: com.vitorpamplona.quic.connection.QuicConnectionDriver, +) : WebTransportWriteStream { + override suspend fun write(chunk: ByteArray) { + stream.send.enqueue(chunk) + driver.wakeup() + } + + override suspend fun finish() { + stream.send.finish() + driver.wakeup() + } +} + /** Adapter for a WT peer-initiated uni stream whose prefix has been stripped. */ private class StrippedWtReadStreamAdapter( private val stripped: com.vitorpamplona.quic.webtransport.StrippedWtStream, ) : WebTransportReadStream { override fun incoming(): Flow = stripped.data } + +/** + * Adapter for a peer-initiated bidi WT stream whose WT_BIDI_STREAM prefix + * has been stripped. Routes [write] / [finish] through the demux's + * driver-aware closures so application bytes actually leave the + * connection. + */ +private class StrippedWtBidiStreamAdapter( + private val stripped: com.vitorpamplona.quic.webtransport.StrippedWtStream, +) : WebTransportBidiStream { + init { + check(!stripped.isUnidirectional) { + "StrippedWtBidiStreamAdapter requires a bidi stream, got uni" + } + } + + override fun incoming(): Flow = stripped.data + + override suspend fun write(chunk: ByteArray) { + val send = + stripped.send + ?: error("peer-initiated bidi stream has no send half — demux didn't wire one") + send(chunk) + } + + override suspend fun finish() { + val finish = + stripped.finish + ?: error("peer-initiated bidi stream has no finish — demux didn't wire one") + finish() + } +} diff --git a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/webtransport/QuicWebTransportSessionState.kt b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/webtransport/QuicWebTransportSessionState.kt index 81358bd7b..7ae02f24d 100644 --- a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/webtransport/QuicWebTransportSessionState.kt +++ b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/webtransport/QuicWebTransportSessionState.kt @@ -54,7 +54,7 @@ class QuicWebTransportSessionState( val isOpen: Boolean get() = connection.status == QuicConnection.Status.CONNECTED - private val demux: WtPeerStreamDemux = WtPeerStreamDemux(connectStreamId, scope) + private val demux: WtPeerStreamDemux = WtPeerStreamDemux(connectStreamId, scope, driver) /** * Completes once a WT_CLOSE_SESSION capsule arrives on the CONNECT bidi. diff --git a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/webtransport/WtPeerStreamDemux.kt b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/webtransport/WtPeerStreamDemux.kt index 14f76b186..a2f48a19e 100644 --- a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/webtransport/WtPeerStreamDemux.kt +++ b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/webtransport/WtPeerStreamDemux.kt @@ -37,11 +37,25 @@ import kotlinx.coroutines.launch /** * A peer-initiated WebTransport stream whose framing prefix has already been * stripped. The [data] flow yields only application-level bytes. + * + * For peer-initiated **bidirectional** streams (where [isUnidirectional] + * is false), [send] and [finish] are wired to the stream's outbound + * half so the application can write its response. They are null on + * unidirectional streams. */ class StrippedWtStream( val streamId: Long, val isUnidirectional: Boolean, val data: Flow, + /** + * Suspends until [chunk] is queued on the stream's send half. Null + * for unidirectional streams (write side is the peer's, not ours). + */ + val send: (suspend (chunk: ByteArray) -> Unit)? = null, + /** + * Half-close the send side (FIN). Null for unidirectional streams. + */ + val finish: (suspend () -> Unit)? = null, ) /** @@ -60,6 +74,14 @@ class StrippedWtStream( class WtPeerStreamDemux( private val expectedConnectStreamId: Long, private val scope: CoroutineScope, + /** + * Optional driver hook. When non-null, every write on a surfaced + * [StrippedWtStream]'s send half wakes the QUIC pump so frames + * actually leave the connection. Tests that drive the demux + * without a real driver (the stream's own send queue is enough) + * can leave this null. + */ + private val driver: com.vitorpamplona.quic.connection.QuicConnectionDriver? = null, ) { private val readyStreams = Channel(Channel.UNLIMITED) @@ -280,11 +302,35 @@ class WtPeerStreamDemux( for (b in prebuffered) if (b.isNotEmpty()) emit(b) for (chunk in chunkChannel) if (chunk.isNotEmpty()) emit(chunk) } + // For peer-initiated bidi streams the application also needs to + // write back. We expose stream.send.enqueue + driver.wakeup so + // the surfaced stream looks symmetric to a locally-opened bidi. + // For uni streams (peer is the sender) we leave both null. + val send: (suspend (ByteArray) -> Unit)? = + if (isUni) { + null + } else { + { chunk -> + stream.send.enqueue(chunk) + driver?.wakeup() + } + } + val finish: (suspend () -> Unit)? = + if (isUni) { + null + } else { + { + stream.send.finish() + driver?.wakeup() + } + } readyStreams.trySend( StrippedWtStream( streamId = stream.streamId, isUnidirectional = isUni, data = data, + send = send, + finish = finish, ), ) }