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 cb1ef868d..3073faccf 100644 --- a/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/transport/FakeWebTransport.kt +++ b/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/transport/FakeWebTransport.kt @@ -57,8 +57,40 @@ class FakeWebTransport private constructor( stateLock.withLock { check(open) { "session closed" } } val localToPeer = Channel(Channel.BUFFERED) val peerToLocal = Channel(Channel.BUFFERED) - val local = FakeBidiStream(write = localToPeer, read = peerToLocal) - val peer = FakeBidiStream(write = peerToLocal, read = localToPeer) + // Shared error-code cells: each side records its own reset() + // and stopSending() codes here so the peer can introspect. + // `localResetCell` records resets on `local.reset(code)` (= + // peer's `lastPeerResetCode`), and vice-versa. + val localResetCell = + java.util.concurrent.atomic + .AtomicLong(NO_CODE) + val peerResetCell = + java.util.concurrent.atomic + .AtomicLong(NO_CODE) + val localStopSendingCell = + java.util.concurrent.atomic + .AtomicLong(NO_CODE) + val peerStopSendingCell = + java.util.concurrent.atomic + .AtomicLong(NO_CODE) + val local = + FakeBidiStream( + write = localToPeer, + read = peerToLocal, + myResetCell = localResetCell, + myStopSendingCell = localStopSendingCell, + peerResetCell = peerResetCell, + peerStopSendingCell = peerStopSendingCell, + ) + val peer = + FakeBidiStream( + write = peerToLocal, + read = localToPeer, + myResetCell = peerResetCell, + myStopSendingCell = peerStopSendingCell, + peerResetCell = localResetCell, + peerStopSendingCell = localStopSendingCell, + ) // Send *peer* half to the other side's inbound queue. outboundBidiStreams.send(peer) return local @@ -87,8 +119,20 @@ class FakeWebTransport private constructor( val priorityCell = java.util.concurrent.atomic .AtomicInteger(0) - outboundUniStreams.send(FakeReadStream(pipe, priorityCell)) - return ChannelWriteStream(pipe, priorityCell) + // Shared reset / stopSending cells. The writer (us) records its + // `reset(code)` call here; the peer-side reader exposes it via + // `FakeReadStream.lastResetCode`. Symmetric: the reader's + // `stopSending(code)` call lands in `stopSendingCell` so the + // writer (= moq-lite publisher tests) can assert the listener + // canceled with the expected code. + val resetCell = + java.util.concurrent.atomic + .AtomicLong(NO_CODE) + val stopSendingCell = + java.util.concurrent.atomic + .AtomicLong(NO_CODE) + outboundUniStreams.send(FakeReadStream(pipe, priorityCell, resetCell, stopSendingCell)) + return ChannelWriteStream(pipe, priorityCell, resetCell, stopSendingCell) } override suspend fun sendDatagram(payload: ByteArray): Boolean { @@ -161,6 +205,20 @@ class FakeWebTransport private constructor( class FakeBidiStream internal constructor( private val write: Channel, private val read: Channel, + /** + * Cell that records this side's `reset(code)` calls. Read by + * [lastResetCode] from the same side (rare — usually it's the + * peer's [peerResetCell] you want). + */ + private val myResetCell: java.util.concurrent.atomic.AtomicLong? = null, + private val myStopSendingCell: java.util.concurrent.atomic.AtomicLong? = null, + /** + * Cell that records the *peer's* `reset(code)` calls (i.e. + * mirror of the peer's [myResetCell]). [lastPeerResetCode] reads + * this so tests can assert "the peer reset the bidi with code X." + */ + private val peerResetCell: java.util.concurrent.atomic.AtomicLong? = null, + private val peerStopSendingCell: java.util.concurrent.atomic.AtomicLong? = null, ) : WebTransportBidiStream { override fun incoming(): Flow = read.receiveAsFlow() @@ -172,6 +230,40 @@ class FakeBidiStream internal constructor( write.close() } + override suspend fun reset(errorCode: Long) { + // First call wins, mirroring `:quic`'s `QuicStream.resetStream` + // contract. Subsequent reset/finish/write are silently + // dropped to keep the wire frame stable in the eyes of any + // test that reads `lastResetCode`. + myResetCell?.compareAndSet(NO_CODE, errorCode) + write.close() + } + + override suspend fun stopSending(errorCode: Long) { + myStopSendingCell?.compareAndSet(NO_CODE, errorCode) + // The peer's send half is our receive half; close it so a + // collector sees end-of-flow. + read.close() + } + + /** + * Code our side passed to [reset], or `null` if reset wasn't + * called. Mostly useful when a test holds the peer's reference; + * [lastPeerResetCode] is the more common assertion. + */ + val lastResetCode: Long? + get() = myResetCell?.get()?.takeIf { it != NO_CODE } + + /** Code the peer side passed to [reset], or `null` if not called. */ + val lastPeerResetCode: Long? + get() = peerResetCell?.get()?.takeIf { it != NO_CODE } + + val lastStopSendingCode: Long? + get() = myStopSendingCell?.get()?.takeIf { it != NO_CODE } + + val lastPeerStopSendingCode: Long? + get() = peerStopSendingCell?.get()?.takeIf { it != NO_CODE } + override fun setPriority(priority: Int) = Unit } @@ -184,9 +276,33 @@ class FakeReadStream internal constructor( * uni-stream writer (e.g. tests that hand-roll a [Channel]). */ private val priorityCell: java.util.concurrent.atomic.AtomicInteger? = null, + /** + * Optional shared cell that records the writer's `reset(code)` + * call. Exposed via [lastResetCode] so tests can assert the + * publisher reset its uni stream with the expected typed error. + */ + private val resetCell: java.util.concurrent.atomic.AtomicLong? = null, + /** + * Optional shared cell that records this side's own + * `stopSending(code)` call. Exposed via [lastStopSendingCode] + * so the test can assert the listener actually fired the cancel. + */ + private val stopSendingCell: java.util.concurrent.atomic.AtomicLong? = null, ) : WebTransportReadStream { override fun incoming(): Flow = read.receiveAsFlow() + override suspend fun stopSending(errorCode: Long) { + // First-call-wins like the production path. Closing `read` + // here would race the production path's "publisher closes + // its write end" semantic — moq-lite calls `stopSending` + // when it's done caring about the bytes, but the publisher + // is still expected to finish what it was sending. We + // record the code without closing so a test that asserts + // `incoming().toList()` post-stopSending still sees any + // bytes the publisher had already enqueued. + stopSendingCell?.compareAndSet(NO_CODE, errorCode) + } + /** * Last priority value the paired write side applied via * [WebTransportWriteStream.setPriority], or `0` if no priority was @@ -195,6 +311,17 @@ class FakeReadStream internal constructor( */ val lastSetPriority: Int? get() = priorityCell?.get() + + /** + * Code the peer-side writer passed to [WebTransportWriteStream.reset], + * or `null` if reset wasn't called. + */ + val lastResetCode: Long? + get() = resetCell?.get()?.takeIf { it != NO_CODE } + + /** Code we passed to [stopSending], or `null` if not called. */ + val lastStopSendingCode: Long? + get() = stopSendingCell?.get()?.takeIf { it != NO_CODE } } /** @@ -210,6 +337,8 @@ private class ChannelWriteStream( * stores each [setPriority] call so the peer can introspect. */ private val priorityCell: java.util.concurrent.atomic.AtomicInteger? = null, + private val resetCell: java.util.concurrent.atomic.AtomicLong? = null, + private val stopSendingCell: java.util.concurrent.atomic.AtomicLong? = null, ) : WebTransportWriteStream { override suspend fun write(chunk: ByteArray) { channel.send(chunk) @@ -219,7 +348,28 @@ private class ChannelWriteStream( channel.close() } + override suspend fun reset(errorCode: Long) { + // First-call-wins. We close the channel so the peer's + // receive flow ends — same shape as `finish` but with a + // typed error code recorded for the peer to introspect. + resetCell?.compareAndSet(NO_CODE, errorCode) + channel.close() + } + override fun setPriority(priority: Int) { priorityCell?.set(priority) } + + @Suppress("unused") + private val unusedStopSendingCellTether: java.util.concurrent.atomic.AtomicLong? = stopSendingCell } + +/** + * Sentinel for "no reset / stopSending code recorded yet" in the + * [java.util.concurrent.atomic.AtomicLong] cells used by [FakeBidiStream] + * and [FakeReadStream] / [ChannelWriteStream]. Real moq-lite codes are + * non-negative varints (RFC 9000 application error codes are u32 + * unsigned-fitting-in-Long); `Long.MIN_VALUE` is safely outside that + * range. + */ +internal const val NO_CODE: Long = Long.MIN_VALUE 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 f581ea986..467578427 100644 --- a/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/transport/WebTransportSession.kt +++ b/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/transport/WebTransportSession.kt @@ -104,6 +104,30 @@ interface WebTransportBidiStream : interface WebTransportReadStream { /** Flow of chunks as they arrive. Completes when the peer closes its write side. */ fun incoming(): Flow + + /** + * RFC 9000 §3.5: ask the peer to stop sending on this stream. + * Causes a `STOP_SENDING(applicationErrorCode)` frame to land at + * the peer, which typically responds with `RESET_STREAM` carrying + * the same code. Call this when the application no longer needs + * the stream's bytes — it lets the peer abandon any pending + * retransmits instead of wasting bandwidth on data we'd discard. + * + * First call wins per `:quic`'s `QuicStream.stopSending` + * lock-free first-call-wins gate; subsequent calls (including + * with a different code) are silently ignored to keep the wire + * frame stable. + * + * Used by moq-lite Lite-03's group-cancel path: a listener that + * decides a specific group is too stale can `stopSending` its + * uni stream so the publisher abandons any in-flight retransmits + * instead of wasting bandwidth on bytes the listener will discard. + * + * Implementations that don't model receive-side cancellation + * (e.g. the in-memory fake when no test asserts on it) MAY + * treat this as a no-op. + */ + suspend fun stopSending(errorCode: Long) } /** Write-only WebTransport stream. */ @@ -113,6 +137,31 @@ interface WebTransportWriteStream { /** Half-close the write side (FIN). No further writes after this call. */ suspend fun finish() + /** + * RFC 9000 §3.5: send `RESET_STREAM(applicationErrorCode)` on + * this stream's send half, abandoning any pending bytes. Distinct + * from [finish], which is a graceful FIN — `reset` carries a + * typed error code the peer can act on. + * + * First call wins per `:quic`'s `QuicStream.resetStream` + * lock-free first-call-wins gate; subsequent calls (and any + * subsequent [write] / [finish]) are silently ignored to keep + * the wire frame stable. + * + * Used by moq-lite Lite-03's typed cancel paths: a publisher + * that rejects a Subscribe (e.g. broadcast / track does not + * exist) follows the SubscribeDrop body with a + * `RESET_STREAM(errorCode)` so the subscriber sees a typed + * application reason rather than the ambiguous "publisher FINed + * the bidi" signal that overlaps with a graceful publisher + * shutdown. + * + * Implementations that don't model sender-driven reset (e.g. + * the in-memory fake when no test asserts on it) MAY treat this + * as a [finish]. + */ + suspend fun reset(errorCode: Long) + /** * Hint to the transport about this stream's drain priority relative * to other streams on the same session. Higher value drains first 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 019ef6e21..f93fa5fe0 100644 --- a/nestsClient/src/jvmAndroid/kotlin/com/vitorpamplona/nestsclient/transport/QuicWebTransportFactory.kt +++ b/nestsClient/src/jvmAndroid/kotlin/com/vitorpamplona/nestsclient/transport/QuicWebTransportFactory.kt @@ -366,6 +366,16 @@ private class QuicBidiStreamAdapter( driver.wakeup() } + override suspend fun reset(errorCode: Long) { + stream.resetStream(errorCode) + driver.wakeup() + } + + override suspend fun stopSending(errorCode: Long) { + stream.stopSending(errorCode) + driver.wakeup() + } + override fun setPriority(priority: Int) { stream.priority = priority } @@ -373,8 +383,14 @@ private class QuicBidiStreamAdapter( private class QuicReadStreamAdapter( private val stream: QuicStream, + private val driver: com.vitorpamplona.quic.connection.QuicConnectionDriver, ) : WebTransportReadStream { override fun incoming(): Flow = stream.incoming + + override suspend fun stopSending(errorCode: Long) { + stream.stopSending(errorCode) + driver.wakeup() + } } /** @@ -397,6 +413,11 @@ private class QuicUniWriteStreamAdapter( driver.wakeup() } + override suspend fun reset(errorCode: Long) { + stream.resetStream(errorCode) + driver.wakeup() + } + override fun setPriority(priority: Int) { stream.priority = priority } @@ -407,13 +428,24 @@ private class StrippedWtReadStreamAdapter( private val stripped: com.vitorpamplona.quic.webtransport.StrippedWtStream, ) : WebTransportReadStream { override fun incoming(): Flow = stripped.data + + override suspend fun stopSending(errorCode: Long) { + // Demux unconditionally wires `stopSending` for every surfaced + // stream — the contract in `StrippedWtStream` is "non-null on + // streams with a read side", which is true here. The `?:` + // fallback exists only to defend against a future demux bug + // that forgets to wire it; Unit is a safe no-op for tests + // that mock a stripped stream without a real receiver. + val ss = stripped.stopSending ?: return + ss(errorCode) + } } /** * 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. + * has been stripped. Routes [write] / [finish] / [reset] / [stopSending] + * 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, @@ -440,13 +472,28 @@ private class StrippedWtBidiStreamAdapter( finish() } + override suspend fun reset(errorCode: Long) { + // Bidi streams have a send side we own, so `reset` is wired by + // the demux. Same defensive `?:` shape as in `write` / `finish`. + val r = + stripped.reset + ?: error("peer-initiated bidi stream has no reset — demux didn't wire one") + r(errorCode) + } + + override suspend fun stopSending(errorCode: Long) { + val ss = stripped.stopSending ?: return + ss(errorCode) + } + /** * No-op: peer-initiated bidi streams arrive through the demux as a * [com.vitorpamplona.quic.webtransport.StrippedWtStream] which exposes - * only `send`/`finish` closures, not the underlying [QuicStream]. The - * moq-lite priority use case targets locally-opened uni group streams - * only, so this path doesn't need to model priority — see the - * [WebTransportWriteStream.setPriority] contract. + * only `send`/`finish`/`reset`/`stopSending` closures, not the + * underlying [QuicStream]. The moq-lite priority use case targets + * locally-opened uni group streams only, so this path doesn't need + * to model priority — see the [WebTransportWriteStream.setPriority] + * contract. */ override fun setPriority(priority: Int) = Unit } 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 0cc040795..d7412b893 100644 --- a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/webtransport/WtPeerStreamDemux.kt +++ b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/webtransport/WtPeerStreamDemux.kt @@ -44,6 +44,14 @@ import kotlinx.coroutines.launch * 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. + * + * For both directions, [reset] (RFC 9000 §3.5 RESET_STREAM) and + * [stopSending] (RFC 9000 §3.5 STOP_SENDING) are exposed so application + * code (e.g. moq-lite Lite-03's typed cancel paths) can convey error + * codes instead of relying on graceful FIN. [reset] is null on uni + * streams whose write side belongs to the peer; [stopSending] is null + * when there's no read side to ask the peer to stop on (i.e. never + * — every stream we surface here has a read side). */ class StrippedWtStream( val streamId: Long, @@ -58,6 +66,22 @@ class StrippedWtStream( * Half-close the send side (FIN). Null for unidirectional streams. */ val finish: (suspend () -> Unit)? = null, + /** + * Send `RESET_STREAM(applicationErrorCode)` on the send half (RFC + * 9000 §3.5). Null when the application doesn't own the send side + * (peer-initiated uni stream). The first call wins per + * `QuicStream.resetStream`'s lock-free first-call-wins gate; + * duplicate calls with different codes are silently ignored to + * keep the wire frame stable. + */ + val reset: (suspend (errorCode: Long) -> Unit)? = null, + /** + * Send `STOP_SENDING(applicationErrorCode)` on the receive half + * (RFC 9000 §3.5) — asks the peer to RESET its corresponding + * send side. Always non-null on streams surfaced here: every + * stripped stream has a read side. First call wins. + */ + val stopSending: (suspend (errorCode: Long) -> Unit)? = null, ) /** @@ -515,6 +539,26 @@ class WtPeerStreamDemux( driver?.wakeup() } } + // RESET_STREAM is meaningful only on streams whose write side + // belongs to us. Same null pattern as [send] / [finish] for + // peer-initiated uni streams. + val reset: (suspend (Long) -> Unit)? = + if (isUni) { + null + } else { + { errorCode -> + stream.resetStream(errorCode) + driver?.wakeup() + } + } + // STOP_SENDING is meaningful on the receive half — every + // stripped stream we surface here has one (uni: peer is + // sender; bidi: peer's send side is our receive side), so + // wire unconditionally. + val stopSending: (suspend (Long) -> Unit) = { errorCode -> + stream.stopSending(errorCode) + driver?.wakeup() + } // Suspending send instead of `trySend` so a slow application // back-pressures the demux instead of silently dropping the // stream. With a bounded buffer ([readyStreamsBuffer]), @@ -529,6 +573,8 @@ class WtPeerStreamDemux( data = data, send = send, finish = finish, + reset = reset, + stopSending = stopSending, ), ) }