mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 11:48:24 +00:00
feat(contextvm): CEP-41 open-ended streams
Build item 9. All seven frame types, the receiver state machine, and keepalive with an injectable clock. The profile carries two independent ordering fields and the whole design turns on keeping them apart: `progress` orders every frame including control frames and is explicitly not a chunk counter, while `chunkIndex` -- contiguous from 0 -- is what validates payload completeness. Chunks are buffered by chunkIndex and released as the contiguous run extends, so a gap is a provisional state while the stream is live rather than an error. Progress sequences are per-sender, so a pong's progress has no required relationship to the ping it answers; pongs match by nonce alone. A pong for a nonce we never sent is ignored rather than fatal -- failing would let a third party disrupt a stream by replaying a stale pong -- but it is not counted as liveness either, so the real probe still expires. Other enforced rules: a second start on a live token fails, close is required for success, lastChunkIndex is a completeness bound that must be satisfied and must be omitted when no chunks were sent, a zero-chunk stream is valid, frames after close or abort are refused, and nonces over 64 bytes are rejected at parse time. Closed carries the bound but deliberately does not complete the JSON-RPC request, since CEP-41 still requires a final response. 26 CVM-41-* tests, 16 negative. Verified by mutation: indexing chunks by progress instead of chunkIndex, and treating any pong as liveness, each turn the suite red (7 failures). Module suite at 76. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_012BfD4txdnsaPRXmNXbup9n
This commit is contained in:
+234
@@ -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()
|
||||
},
|
||||
)
|
||||
}
|
||||
}
|
||||
+255
@@ -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<String>,
|
||||
) : 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<Double>()
|
||||
|
||||
private val pending = mutableMapOf<Long, String>()
|
||||
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<String, Long>()
|
||||
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<String>()
|
||||
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)
|
||||
}
|
||||
}
|
||||
+320
@@ -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<OpenStreamEvent.Delivered>(first)
|
||||
assertEquals(listOf("Hello"), first.fragments)
|
||||
|
||||
val second = rx.accept(OpenStreamFrame.chunk(token, 3.0, 1, " world"))
|
||||
assertIs<OpenStreamEvent.Delivered>(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<OpenStreamEvent.Continue>(rx.accept(OpenStreamFrame.chunk(token, 2.0, 1, "b")))
|
||||
assertIs<OpenStreamEvent.Continue>(rx.accept(OpenStreamFrame.chunk(token, 3.0, 2, "c")))
|
||||
|
||||
val released = rx.accept(OpenStreamFrame.chunk(token, 4.0, 0, "a"))
|
||||
assertIs<OpenStreamEvent.Delivered>(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<OpenStreamEvent.Closed>(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<OpenStreamEvent.Closed>(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<OpenStreamEvent.Closed>(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<OpenStreamEvent.Pong>(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<OpenStreamEvent.Continue>(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<OpenStreamEvent.SendAccept>(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<OpenStreamEvent.Aborted>(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<OpenStreamException> { 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<OpenStreamException> { 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<OpenStreamException> {
|
||||
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<OpenStreamException> {
|
||||
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<OpenStreamException> {
|
||||
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<OpenStreamException> { 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<OpenStreamException> { rx.accept(OpenStreamFrame.chunk(token, 3.0, 0, "a again")) }
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `CVM-41-27 rejects a chunk before start`() {
|
||||
val rx = receiver()
|
||||
assertFailsWith<OpenStreamException> { 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<OpenStreamException> { 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<OpenStreamException> { 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<OpenStreamException> {
|
||||
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<TransferFrameException> { 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<TransferFrameException> { 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<OpenStreamEvent.Continue>(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))
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user