diff --git a/contextvm/src/commonMain/kotlin/com/vitorpamplona/contextvm/transfer/stream/OpenStreamFrame.kt b/contextvm/src/commonMain/kotlin/com/vitorpamplona/contextvm/transfer/stream/OpenStreamFrame.kt new file mode 100644 index 0000000000..8511a4a60c --- /dev/null +++ b/contextvm/src/commonMain/kotlin/com/vitorpamplona/contextvm/transfer/stream/OpenStreamFrame.kt @@ -0,0 +1,234 @@ +/* + * 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.contextvm.transfer.stream + +import com.vitorpamplona.contextvm.transfer.ProgressEnvelope +import com.vitorpamplona.contextvm.transfer.ProgressEnvelope.Companion.asLongOrNull +import com.vitorpamplona.contextvm.transfer.ProgressEnvelope.Companion.asStringOrNull +import com.vitorpamplona.contextvm.transfer.ProgressToken +import com.vitorpamplona.contextvm.transfer.TransferFrameException +import kotlinx.serialization.json.JsonObjectBuilder +import kotlinx.serialization.json.JsonPrimitive +import kotlinx.serialization.json.buildJsonObject + +/** + * A CEP-41 open-ended stream frame. + * + * The profile carries **two** independent ordering fields and conflating them is + * the classic implementation bug: + * + * - `progress` (on the envelope) orders *all* frames, control frames included, + * and is explicitly not a chunk counter. Each peer numbers its own outbound + * frames from 1, so a value may only be compared against frames from the same + * sender. + * - [Chunk.chunkIndex] starts at 0, increases contiguously, and is what + * validates payload completeness. + */ +sealed interface OpenStreamFrame { + val envelope: ProgressEnvelope + val token: ProgressToken get() = envelope.token + val progress: Double get() = envelope.progress + + data class Start( + override val envelope: ProgressEnvelope, + ) : OpenStreamFrame + + data class Accept( + override val envelope: ProgressEnvelope, + ) : OpenStreamFrame + + data class Chunk( + override val envelope: ProgressEnvelope, + val data: String, + val chunkIndex: Long, + ) : OpenStreamFrame + + data class Ping( + override val envelope: ProgressEnvelope, + val nonce: String, + ) : OpenStreamFrame + + data class Pong( + override val envelope: ProgressEnvelope, + val nonce: String, + ) : OpenStreamFrame + + /** + * Successful closure. + * + * [lastChunkIndex] is present only when the sender is declaring a + * completeness bound; a live, open-ended feed omits it, and a stream that + * carried no chunks MUST omit it. + */ + data class Close( + override val envelope: ProgressEnvelope, + val lastChunkIndex: Long? = null, + ) : OpenStreamFrame + + data class Abort( + override val envelope: ProgressEnvelope, + val reason: String? = null, + ) : OpenStreamFrame + + companion object { + const val START = "start" + const val ACCEPT = "accept" + const val CHUNK = "chunk" + const val PING = "ping" + const val PONG = "pong" + const val CLOSE = "close" + const val ABORT = "abort" + + const val DATA = "data" + const val CHUNK_INDEX = "chunkIndex" + const val NONCE = "nonce" + const val LAST_CHUNK_INDEX = "lastChunkIndex" + const val REASON = "reason" + + /** Receivers SHOULD enforce this ceiling and MAY reject oversized nonces. */ + const val MAX_NONCE_BYTES = 64 + + /** Parses a CEP-41 frame, or null when the envelope is a different profile. */ + fun parseOrNull(envelope: ProgressEnvelope): OpenStreamFrame? { + if (envelope.type != ProgressEnvelope.TYPE_OPEN_STREAM) return null + val cvm = envelope.cvm + + return when (val frameType = envelope.frameType) { + START -> Start(envelope) + ACCEPT -> Accept(envelope) + + CHUNK -> + Chunk( + envelope = envelope, + data = + cvm[DATA]?.asStringOrNull() + ?: throw TransferFrameException("chunk requires string data"), + chunkIndex = + cvm[CHUNK_INDEX]?.asLongOrNull() + ?: throw TransferFrameException("chunk requires chunkIndex"), + ) + + PING -> Ping(envelope, requireNonce(cvm[NONCE]?.asStringOrNull(), PING)) + PONG -> Pong(envelope, requireNonce(cvm[NONCE]?.asStringOrNull(), PONG)) + + CLOSE -> Close(envelope, cvm[LAST_CHUNK_INDEX]?.asLongOrNull()) + + ABORT -> Abort(envelope, cvm[REASON]?.asStringOrNull()) + + else -> throw TransferFrameException("unknown open-stream frameType: $frameType") + } + } + + private fun requireNonce( + nonce: String?, + frameType: String, + ): String { + if (nonce == null) throw TransferFrameException("$frameType requires a nonce") + if (nonce.encodeToByteArray().size > MAX_NONCE_BYTES) { + throw TransferFrameException("$frameType nonce exceeds $MAX_NONCE_BYTES bytes") + } + return nonce + } + + fun start( + token: ProgressToken, + progress: Double, + ) = Start(envelopeOf(START, token, progress)) + + fun accept( + token: ProgressToken, + progress: Double, + ) = Accept(envelopeOf(ACCEPT, token, progress)) + + fun chunk( + token: ProgressToken, + progress: Double, + chunkIndex: Long, + data: String, + ) = Chunk( + envelope = + envelopeOf(CHUNK, token, progress) { + put(CHUNK_INDEX, JsonPrimitive(chunkIndex)) + put(DATA, JsonPrimitive(data)) + }, + data = data, + chunkIndex = chunkIndex, + ) + + fun ping( + token: ProgressToken, + progress: Double, + nonce: String, + ) = Ping( + envelope = envelopeOf(PING, token, progress) { put(NONCE, JsonPrimitive(nonce)) }, + nonce = nonce, + ) + + fun pong( + token: ProgressToken, + progress: Double, + nonce: String, + ) = Pong( + envelope = envelopeOf(PONG, token, progress) { put(NONCE, JsonPrimitive(nonce)) }, + nonce = nonce, + ) + + fun close( + token: ProgressToken, + progress: Double, + lastChunkIndex: Long? = null, + ) = Close( + envelope = + envelopeOf(CLOSE, token, progress) { + lastChunkIndex?.let { put(LAST_CHUNK_INDEX, JsonPrimitive(it)) } + }, + lastChunkIndex = lastChunkIndex, + ) + + fun abort( + token: ProgressToken, + progress: Double, + reason: String? = null, + ) = Abort( + envelope = + envelopeOf(ABORT, token, progress) { + reason?.let { put(REASON, JsonPrimitive(it)) } + }, + reason = reason, + ) + + private fun envelopeOf( + frameType: String, + token: ProgressToken, + progress: Double, + body: JsonObjectBuilder.() -> Unit = {}, + ) = ProgressEnvelope( + token = token, + progress = progress, + cvm = + buildJsonObject { + put(ProgressEnvelope.TYPE, JsonPrimitive(ProgressEnvelope.TYPE_OPEN_STREAM)) + put(ProgressEnvelope.FRAME_TYPE, JsonPrimitive(frameType)) + body() + }, + ) + } +} diff --git a/contextvm/src/commonMain/kotlin/com/vitorpamplona/contextvm/transfer/stream/OpenStreamReceiver.kt b/contextvm/src/commonMain/kotlin/com/vitorpamplona/contextvm/transfer/stream/OpenStreamReceiver.kt new file mode 100644 index 0000000000..99ce01efdd --- /dev/null +++ b/contextvm/src/commonMain/kotlin/com/vitorpamplona/contextvm/transfer/stream/OpenStreamReceiver.kt @@ -0,0 +1,255 @@ +/* + * 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.contextvm.transfer.stream + +import com.vitorpamplona.contextvm.transfer.ProgressToken + +/** Local resource and keepalive policy for one stream. */ +data class OpenStreamPolicy( + /** + * Idle time before the peer must be probed with a `ping`. + * + * 30s matches what deployed clients use. A shorter window (the 20s some + * SDKs default to) turns a single lost relay round-trip into a stream abort. + */ + val idleTimeoutMs: Long = 30_000, + /** How long a `ping` may go unanswered before the stream is failed. */ + val probeTimeoutMs: Long = 30_000, + /** Bounded buffering for chunks that arrive before their predecessors. */ + val maxPendingChunks: Int = 256, +) + +/** Outcome of feeding one inbound frame to [OpenStreamReceiver]. */ +sealed interface OpenStreamEvent { + /** Frame accepted; nothing to deliver or answer. */ + data object Continue : OpenStreamEvent + + /** The receiver must send an `accept` before the peer may send chunks. */ + data object SendAccept : OpenStreamEvent + + /** Payload fragments that became contiguous, in `chunkIndex` order. */ + data class Delivered( + val fragments: List, + ) : OpenStreamEvent + + /** + * The peer probed us and MUST receive a `pong` carrying the same nonce. + * + * The response rides our own outbound progress sequence, which is why the + * caller supplies the progress value rather than this class. + */ + data class Pong( + val nonce: String, + ) : OpenStreamEvent + + /** + * The peer closed the stream. + * + * This does **not** complete the originating JSON-RPC request: CEP-41 + * requires a final response as well, and a client must never synthesize + * success from `close` alone. + */ + data class Closed( + val lastChunkIndex: Long?, + ) : OpenStreamEvent + + data class Aborted( + val reason: String?, + ) : OpenStreamEvent +} + +/** Thrown when a stream violates CEP-41. The stream is terminal once this is raised. */ +class OpenStreamException( + message: String, +) : IllegalStateException(message) + +/** + * Receives one CEP-41 open-ended stream, keyed by its `progressToken`. + * + * Rules are ยง6.5 `CVM-41-*` in `quartz/plans/2026-09-17-cordn-interop.md`. + * + * @param requireAccept true in stateless bootstrap, where the peer must wait for + * our `accept` before its first chunk. + * @param now injectable clock (epoch millis) so keepalive is testable without + * real time passing. + */ +class OpenStreamReceiver( + val token: ProgressToken, + private val policy: OpenStreamPolicy = OpenStreamPolicy(), + private val requireAccept: Boolean = true, + private val now: () -> Long = { 0L }, +) { + private var started = false + private var startProgress: Double? = null + private var accepted = false + private var terminal = false + + /** Every inbound `progress` seen, to reject replays without assuming arrival order. */ + private val seenProgress = mutableSetOf() + + private val pending = mutableMapOf() + private var nextIndex = 0L + private var highestIndex: Long? = null + + /** Nonces of pings we sent that are still awaiting a pong, with their send time. */ + private val outstandingPings = mutableMapOf() + private var lastActivityMs = 0L + + val isTerminal get() = terminal + + fun accept(frame: OpenStreamFrame): OpenStreamEvent { + if (terminal) throw OpenStreamException("frame received after the stream terminated") + if (frame.token != token) throw OpenStreamException("frame belongs to another stream") + + lastActivityMs = now() + + if (frame is OpenStreamFrame.Abort) { + terminal = true + return OpenStreamEvent.Aborted(frame.reason) + } + + // A repeated progress value from the same peer is a replay or a bug. + // Arrival *order* is deliberately not checked -- relays reorder, and the + // CEP makes chunkIndex, not progress, the completeness test. + if (!seenProgress.add(frame.progress)) { + fail("duplicate progress ${frame.progress} from the peer") + } + + return when (frame) { + is OpenStreamFrame.Start -> onStart(frame) + is OpenStreamFrame.Accept -> OpenStreamEvent.Continue + is OpenStreamFrame.Chunk -> onChunk(frame) + is OpenStreamFrame.Ping -> OpenStreamEvent.Pong(frame.nonce) + is OpenStreamFrame.Pong -> onPong(frame) + is OpenStreamFrame.Close -> onClose(frame) + is OpenStreamFrame.Abort -> error("handled above") + } + } + + private fun onStart(frame: OpenStreamFrame.Start): OpenStreamEvent { + // A second start for a live token MUST fail the stream. + if (started) fail("a second start arrived for a live stream") + started = true + startProgress = frame.progress + return if (requireAccept) OpenStreamEvent.SendAccept else OpenStreamEvent.Continue + } + + private fun onChunk(frame: OpenStreamFrame.Chunk): OpenStreamEvent { + if (!started) fail("chunk arrived before start") + if (requireAccept && !accepted) fail("chunk arrived before accept in a bootstrap stream") + + val start = startProgress + if (start != null && frame.progress <= start) { + fail("chunk progress ${frame.progress} does not follow start at $start") + } + if (frame.chunkIndex < nextIndex) { + fail("chunkIndex ${frame.chunkIndex} was already delivered") + } + if (pending.containsKey(frame.chunkIndex)) { + fail("duplicate chunkIndex ${frame.chunkIndex}") + } + if (pending.size >= policy.maxPendingChunks) { + fail("buffered chunk count exceeds the limit ${policy.maxPendingChunks}") + } + + pending[frame.chunkIndex] = frame.data + highestIndex = maxOf(highestIndex ?: frame.chunkIndex, frame.chunkIndex) + + // Release the contiguous run. A gap is not an error while the stream is + // live -- it is a provisional gap the peer may still fill. + val ready = mutableListOf() + while (true) { + val next = pending.remove(nextIndex) ?: break + ready += next + nextIndex += 1 + } + + return if (ready.isEmpty()) OpenStreamEvent.Continue else OpenStreamEvent.Delivered(ready) + } + + private fun onPong(frame: OpenStreamFrame.Pong): OpenStreamEvent { + // A pong for a nonce we never sent, or already retired, is not evidence + // of liveness. It is ignored rather than fatal: the CEP allows local + // anti-abuse policy, and failing the stream would let a third party + // disrupt it by replaying a stale pong. + outstandingPings.remove(frame.nonce) + return OpenStreamEvent.Continue + } + + private fun onClose(frame: OpenStreamFrame.Close): OpenStreamEvent { + if (!started) fail("close arrived before start") + + val bound = frame.lastChunkIndex + if (bound != null) { + if (highestIndex == null) { + // "If the stream included no chunk frames, close.lastChunkIndex + // MUST be omitted." + fail("close declared lastChunkIndex $bound but the stream carried no chunks") + } + if (bound != (nextIndex - 1)) { + fail("close declared lastChunkIndex $bound but ${nextIndex - 1} is the last contiguous index") + } + if (pending.isNotEmpty()) { + fail("close declared a completeness bound while ${pending.size} chunks remain unresolved") + } + } + + terminal = true + return OpenStreamEvent.Closed(bound) + } + + // --- keepalive --- + + /** True when the idle timeout has elapsed and the peer must be probed. */ + fun needsProbe(atMs: Long = now()) = !terminal && outstandingPings.isEmpty() && atMs - lastActivityMs >= policy.idleTimeoutMs + + /** Records that we sent a `ping` with [nonce], starting its probe window. */ + fun markProbeSent( + nonce: String, + atMs: Long = now(), + ) { + require(nonce.encodeToByteArray().size <= OpenStreamFrame.MAX_NONCE_BYTES) { + "nonce exceeds ${OpenStreamFrame.MAX_NONCE_BYTES} bytes" + } + outstandingPings[nonce] = atMs + } + + /** + * True when a probe went unanswered past the probe timeout. + * + * The caller MUST then fail the stream, and SHOULD send `abort` if it still + * can. + */ + fun probeExpired(atMs: Long = now()) = outstandingPings.values.any { atMs - it >= policy.probeTimeoutMs } + + /** Records that we sent `accept`, unblocking the peer's chunk frames. */ + fun markAccepted() { + accepted = true + } + + /** Marks the stream terminal after a local policy failure. */ + fun failLocally(reason: String): Nothing = fail(reason) + + private fun fail(reason: String): Nothing { + terminal = true + throw OpenStreamException(reason) + } +} diff --git a/contextvm/src/commonTest/kotlin/com/vitorpamplona/contextvm/transfer/stream/OpenStreamTest.kt b/contextvm/src/commonTest/kotlin/com/vitorpamplona/contextvm/transfer/stream/OpenStreamTest.kt new file mode 100644 index 0000000000..c4daf810b6 --- /dev/null +++ b/contextvm/src/commonTest/kotlin/com/vitorpamplona/contextvm/transfer/stream/OpenStreamTest.kt @@ -0,0 +1,320 @@ +/* + * 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.contextvm.transfer.stream + +import com.vitorpamplona.contextvm.jsonrpc.JsonRpcCodec +import com.vitorpamplona.contextvm.jsonrpc.JsonRpcNotification +import com.vitorpamplona.contextvm.transfer.ProgressEnvelope +import com.vitorpamplona.contextvm.transfer.ProgressToken +import com.vitorpamplona.contextvm.transfer.TransferFrameException +import kotlinx.serialization.json.JsonPrimitive +import kotlinx.serialization.json.buildJsonObject +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertFailsWith +import kotlin.test.assertFalse +import kotlin.test.assertIs +import kotlin.test.assertTrue + +/** `CVM-41-*`: open-ended streams. */ +class OpenStreamTest { + private val token = ProgressToken.Text("req-123") + + private fun receiver( + requireAccept: Boolean = false, + policy: OpenStreamPolicy = OpenStreamPolicy(), + clock: () -> Long = { 0L }, + ) = OpenStreamReceiver(token, policy, requireAccept, clock) + + // --- happy paths --- + + @Test + fun `CVM-41-01 delivers contiguous chunks incrementally`() { + val rx = receiver() + rx.accept(OpenStreamFrame.start(token, 1.0)) + + val first = rx.accept(OpenStreamFrame.chunk(token, 2.0, 0, "Hello")) + assertIs(first) + assertEquals(listOf("Hello"), first.fragments) + + val second = rx.accept(OpenStreamFrame.chunk(token, 3.0, 1, " world")) + assertIs(second) + assertEquals(listOf(" world"), second.fragments) + } + + @Test + fun `CVM-41-02 buffers a gap and releases the run once it closes`() { + val rx = receiver() + rx.accept(OpenStreamFrame.start(token, 1.0)) + + // index 1 arrives first; a gap is not an error while the stream is live. + assertIs(rx.accept(OpenStreamFrame.chunk(token, 2.0, 1, "b"))) + assertIs(rx.accept(OpenStreamFrame.chunk(token, 3.0, 2, "c"))) + + val released = rx.accept(OpenStreamFrame.chunk(token, 4.0, 0, "a")) + assertIs(released) + assertEquals(listOf("a", "b", "c"), released.fragments) + } + + @Test + fun `CVM-41-03 a zero-chunk stream is valid`() { + // close straight after start, lastChunkIndex omitted. + val rx = receiver() + rx.accept(OpenStreamFrame.start(token, 1.0)) + val closed = rx.accept(OpenStreamFrame.close(token, 2.0)) + assertIs(closed) + assertEquals(null, closed.lastChunkIndex) + } + + @Test + fun `CVM-41-04 close with a satisfied lastChunkIndex succeeds`() { + val rx = receiver() + rx.accept(OpenStreamFrame.start(token, 1.0)) + rx.accept(OpenStreamFrame.chunk(token, 2.0, 0, "a")) + rx.accept(OpenStreamFrame.chunk(token, 3.0, 1, "b")) + + val closed = rx.accept(OpenStreamFrame.close(token, 4.0, lastChunkIndex = 1)) + assertIs(closed) + assertEquals(1L, closed.lastChunkIndex) + } + + @Test + fun `CVM-41-05 close without lastChunkIndex succeeds on an open-ended feed`() { + val rx = receiver() + rx.accept(OpenStreamFrame.start(token, 1.0)) + rx.accept(OpenStreamFrame.chunk(token, 2.0, 0, "tick")) + assertIs(rx.accept(OpenStreamFrame.close(token, 3.0))) + } + + @Test + fun `CVM-41-06 a ping must be answered with a pong carrying the same nonce`() { + val rx = receiver() + rx.accept(OpenStreamFrame.start(token, 1.0)) + val event = rx.accept(OpenStreamFrame.ping(token, 2.0, "n-1")) + assertIs(event) + assertEquals("n-1", event.nonce) + } + + @Test + fun `CVM-41-07 pong progress has no ordering relationship to the ping`() { + // Pongs are matched by nonce only; each peer numbers its own sequence, + // so a pong may legitimately carry a lower progress than the ping. + val rx = receiver() + rx.accept(OpenStreamFrame.start(token, 10.0)) + rx.markProbeSent("n-1") + assertIs(rx.accept(OpenStreamFrame.pong(token, 1.0, "n-1"))) + assertFalse(rx.probeExpired(atMs = 1_000_000)) + } + + @Test + fun `CVM-41-08 keepalive fails the stream when a probe goes unanswered`() { + var clock = 0L + val rx = receiver(policy = OpenStreamPolicy(idleTimeoutMs = 30_000, probeTimeoutMs = 30_000)) { clock } + rx.accept(OpenStreamFrame.start(token, 1.0)) + + clock = 30_000 + assertTrue(rx.needsProbe(clock), "idle timeout elapsed, the peer must be probed") + + rx.markProbeSent("n-1", clock) + clock = 59_000 + assertFalse(rx.probeExpired(clock), "still inside the probe window") + + clock = 60_000 + assertTrue(rx.probeExpired(clock), "probe window elapsed with no matching pong") + } + + @Test + fun `CVM-41-09 a bootstrap stream asks the receiver to send accept`() { + val rx = receiver(requireAccept = true) + assertIs(rx.accept(OpenStreamFrame.start(token, 1.0))) + } + + @Test + fun `CVM-41-10 abort is terminal from either peer`() { + val rx = receiver() + rx.accept(OpenStreamFrame.start(token, 1.0)) + val aborted = rx.accept(OpenStreamFrame.abort(token, 2.0, "resource exhaustion")) + assertIs(aborted) + assertEquals("resource exhaustion", aborted.reason) + assertTrue(rx.isTerminal) + } + + // --- rejections --- + + @Test + fun `CVM-41-20 rejects a second start on a live stream`() { + val rx = receiver() + rx.accept(OpenStreamFrame.start(token, 1.0)) + assertFailsWith { rx.accept(OpenStreamFrame.start(token, 2.0)) } + } + + @Test + fun `CVM-41-21 rejects frames after close`() { + val rx = receiver() + rx.accept(OpenStreamFrame.start(token, 1.0)) + rx.accept(OpenStreamFrame.close(token, 2.0)) + assertFailsWith { rx.accept(OpenStreamFrame.chunk(token, 3.0, 0, "late")) } + } + + @Test + fun `CVM-41-22 rejects close with a lastChunkIndex that is not satisfied`() { + val rx = receiver() + rx.accept(OpenStreamFrame.start(token, 1.0)) + rx.accept(OpenStreamFrame.chunk(token, 2.0, 0, "a")) + // Declares index 5 as the bound while only 0 arrived. + assertFailsWith { + rx.accept(OpenStreamFrame.close(token, 3.0, lastChunkIndex = 5)) + } + } + + @Test + fun `CVM-41-23 rejects close with a bound while a gap remains unresolved`() { + val rx = receiver() + rx.accept(OpenStreamFrame.start(token, 1.0)) + rx.accept(OpenStreamFrame.chunk(token, 2.0, 0, "a")) + rx.accept(OpenStreamFrame.chunk(token, 3.0, 2, "c")) // index 1 missing + assertFailsWith { + rx.accept(OpenStreamFrame.close(token, 4.0, lastChunkIndex = 2)) + } + } + + @Test + fun `CVM-41-24 rejects lastChunkIndex on a stream that carried no chunks`() { + val rx = receiver() + rx.accept(OpenStreamFrame.start(token, 1.0)) + assertFailsWith { + rx.accept(OpenStreamFrame.close(token, 2.0, lastChunkIndex = 0)) + } + } + + @Test + fun `CVM-41-25 rejects a duplicate chunkIndex`() { + val rx = receiver() + rx.accept(OpenStreamFrame.start(token, 1.0)) + rx.accept(OpenStreamFrame.chunk(token, 2.0, 1, "b")) + assertFailsWith { rx.accept(OpenStreamFrame.chunk(token, 3.0, 1, "b again")) } + } + + @Test + fun `CVM-41-26 rejects a chunkIndex that was already delivered`() { + val rx = receiver() + rx.accept(OpenStreamFrame.start(token, 1.0)) + rx.accept(OpenStreamFrame.chunk(token, 2.0, 0, "a")) + assertFailsWith { rx.accept(OpenStreamFrame.chunk(token, 3.0, 0, "a again")) } + } + + @Test + fun `CVM-41-27 rejects a chunk before start`() { + val rx = receiver() + assertFailsWith { rx.accept(OpenStreamFrame.chunk(token, 1.0, 0, "a")) } + } + + @Test + fun `CVM-41-28 rejects a chunk before accept in a bootstrap stream`() { + val rx = receiver(requireAccept = true) + rx.accept(OpenStreamFrame.start(token, 1.0)) + assertFailsWith { rx.accept(OpenStreamFrame.chunk(token, 2.0, 0, "a")) } + } + + @Test + fun `CVM-41-29 rejects a repeated progress value from the peer`() { + val rx = receiver() + rx.accept(OpenStreamFrame.start(token, 1.0)) + rx.accept(OpenStreamFrame.chunk(token, 2.0, 0, "a")) + assertFailsWith { rx.accept(OpenStreamFrame.chunk(token, 2.0, 1, "b")) } + } + + @Test + fun `CVM-41-30 rejects a frame belonging to another stream`() { + val rx = receiver() + assertFailsWith { + rx.accept(OpenStreamFrame.start(ProgressToken.Text("other"), 1.0)) + } + } + + @Test + fun `CVM-41-31 rejects an oversized nonce at parse time`() { + val envelope = + ProgressEnvelope( + token = token, + progress = 1.0, + cvm = + buildJsonObject { + put(ProgressEnvelope.TYPE, JsonPrimitive(ProgressEnvelope.TYPE_OPEN_STREAM)) + put(ProgressEnvelope.FRAME_TYPE, JsonPrimitive(OpenStreamFrame.PING)) + put(OpenStreamFrame.NONCE, JsonPrimitive("x".repeat(65))) + }, + ) + assertFailsWith { OpenStreamFrame.parseOrNull(envelope) } + } + + @Test + fun `CVM-41-32 rejects a chunk missing its chunkIndex`() { + // progress is not a chunk counter, so a chunk without chunkIndex cannot + // be positioned at all. + val envelope = + ProgressEnvelope( + token = token, + progress = 1.0, + cvm = + buildJsonObject { + put(ProgressEnvelope.TYPE, JsonPrimitive(ProgressEnvelope.TYPE_OPEN_STREAM)) + put(ProgressEnvelope.FRAME_TYPE, JsonPrimitive(OpenStreamFrame.CHUNK)) + put(OpenStreamFrame.DATA, JsonPrimitive("a")) + }, + ) + assertFailsWith { OpenStreamFrame.parseOrNull(envelope) } + } + + @Test + fun `CVM-41-33 ignores a pong for a nonce we never sent`() { + // Not fatal: failing here would let a third party disrupt the stream by + // replaying a stale pong. It simply is not liveness evidence. + val rx = receiver() + rx.accept(OpenStreamFrame.start(token, 1.0)) + rx.markProbeSent("real") + assertIs(rx.accept(OpenStreamFrame.pong(token, 2.0, "forged"))) + assertTrue(rx.probeExpired(atMs = 1_000_000), "the real probe is still outstanding") + } + + @Test + fun `CVM-41-34 the frame round-trips through a progress notification`() { + val frame = OpenStreamFrame.chunk(token, 2.0, 7, "payload") + val decoded = JsonRpcCodec.decode(JsonRpcCodec.encode(frame.envelope.toNotification())) + val envelope = ProgressEnvelope.parseOrNull(decoded as JsonRpcNotification) + assertEquals(frame, OpenStreamFrame.parseOrNull(envelope!!)) + } + + @Test + fun `CVM-41-35 ignores an envelope belonging to another transfer profile`() { + val envelope = + ProgressEnvelope( + token = token, + progress = 1.0, + cvm = + buildJsonObject { + put(ProgressEnvelope.TYPE, JsonPrimitive(ProgressEnvelope.TYPE_OVERSIZED)) + put(ProgressEnvelope.FRAME_TYPE, JsonPrimitive("start")) + }, + ) + assertEquals(null, OpenStreamFrame.parseOrNull(envelope)) + } +}