feat(contextvm): CEP-22 bounded oversized payload transfer

Build items 8 and 10. Frame codec, bounded-reassembly receiver and the
fragmenting sender, over the notifications/progress envelope both transfer
profiles share.

Two wire details the spec only shows by example, and both are digest-critical:
chunk `data` is a raw substring of the serialized JSON-RPC message rather than
base64, and `digest` is algorithm-prefixed (`sha256:<hex>`). Fragments are
therefore concatenated as text and encoded to UTF-8 once; the payload is never
parsed and re-serialized in transit, which would reorder keys and break the
hash.

Receiver enforcement: admission control on declared totalBytes/totalChunks
before any state is committed, chunk/end position validated against start,
duplicate progress rejected, bounded out-of-order buffering, and nothing
surfaced upward until the chunk count, byte length and digest all agree.
Sender never splits a surrogate pair, since both halves would survive JSON
encoding individually but rejoin into a different codepoint.

One design bug caught by its own test: the receiver initially rejected frames
whose progress did not increase on arrival, conflating "the sender emits
monotonic progress" with "frames arrive in order". CEP-22 says the opposite --
progress is the assembly index and explicitly not a guarantee of arrival order,
and receivers may buffer out-of-order chunks. The check is now positional
(chunks sort after start, end after every chunk) rather than arrival-ordered.

21 CVM-22-* tests, 14 of them negative. Module suite at 50.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_012BfD4txdnsaPRXmNXbup9n
This commit is contained in:
Claude
2026-09-18 01:31:23 +00:00
parent a9a883cb66
commit 34d82ad17a
6 changed files with 1073 additions and 0 deletions
@@ -0,0 +1,70 @@
/*
* 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.mcp
/**
* The MCP methods ContextVM carries.
*
* ContextVM is a transport: these names are MCP's, not ContextVM's, and travel
* unmodified inside the event `content`.
*/
object McpMethods {
const val INITIALIZE = "initialize"
const val INITIALIZED = "notifications/initialized"
const val PING = "ping"
const val TOOLS_LIST = "tools/list"
const val TOOLS_CALL = "tools/call"
const val RESOURCES_LIST = "resources/list"
const val RESOURCE_TEMPLATES_LIST = "resources/templates/list"
const val PROMPTS_LIST = "prompts/list"
/**
* The envelope CEP-22 and CEP-41 both ride.
*
* Their frames are ordinary MCP progress notifications with an extra `cvm`
* object in `params`; a peer that does not understand ContextVM transfer
* profiles still sees valid MCP.
*/
const val PROGRESS = "notifications/progress"
/** CEP-8 transparent lifecycle notifications. */
const val PAYMENT_REQUIRED = "notifications/payment_required"
const val PAYMENT_ACCEPTED = "notifications/payment_accepted"
const val PAYMENT_REJECTED = "notifications/payment_rejected"
}
/**
* Well-known member names inside MCP `params`.
*
* `_meta` matters beyond convenience: CEP-8 excludes it from the canonical
* invocation identity (because MCP regenerates `progressToken` on every call)
* while still requiring it to reach the handler at execution time.
*/
object McpParams {
const val META = "_meta"
const val PROGRESS_TOKEN = "progressToken"
const val NAME = "name"
const val ARGUMENTS = "arguments"
/** CEP-16: the caller identity a server transport injects into `_meta`. */
const val CLIENT_PUBKEY = "clientPubkey"
}
@@ -0,0 +1,142 @@
/*
* 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
import com.vitorpamplona.contextvm.jsonrpc.JsonRpcNotification
import com.vitorpamplona.contextvm.mcp.McpMethods
import com.vitorpamplona.contextvm.mcp.McpParams
import kotlinx.serialization.json.JsonObject
import kotlinx.serialization.json.JsonPrimitive
import kotlinx.serialization.json.buildJsonObject
import kotlinx.serialization.json.doubleOrNull
import kotlinx.serialization.json.longOrNull
/**
* An MCP progress token: the transfer/stream identifier for CEP-22 and CEP-41.
*
* MCP allows a string or a number. The distinction is kept rather than
* normalised for the same reason as `JsonRpcId`: a faithful round trip.
*/
sealed interface ProgressToken {
data class Num(
val value: Long,
) : ProgressToken
data class Text(
val value: String,
) : ProgressToken
}
/** Thrown when a `notifications/progress` payload is not a usable transfer frame. */
class TransferFrameException(
message: String,
) : IllegalArgumentException(message)
/**
* The `notifications/progress` envelope shared by CEP-22 and CEP-41.
*
* MCP owns `progressToken`, `progress`, `total` and `message`; ContextVM adds
* the `cvm` object carrying the frame. `total` and `message` are UX hints and
* explicitly do not define transfer correctness, so nothing below reads them
* for control flow.
*/
data class ProgressEnvelope(
val token: ProgressToken,
val progress: Double,
val total: Double? = null,
val message: String? = null,
val cvm: JsonObject,
) {
val type: String? get() = cvm[TYPE]?.asStringOrNull()
val frameType: String? get() = cvm[FRAME_TYPE]?.asStringOrNull()
fun toNotification(): JsonRpcNotification =
JsonRpcNotification(
McpMethods.PROGRESS,
buildJsonObject {
put(McpParams.PROGRESS_TOKEN, token.toPrimitive())
put(PROGRESS, JsonPrimitive(progress))
total?.let { put(TOTAL, JsonPrimitive(it)) }
message?.let { put(MESSAGE, JsonPrimitive(it)) }
put(CVM, cvm)
},
)
companion object {
const val PROGRESS = "progress"
const val TOTAL = "total"
const val MESSAGE = "message"
const val CVM = "cvm"
const val TYPE = "type"
const val FRAME_TYPE = "frameType"
/** CEP-22's `cvm.type`. */
const val TYPE_OVERSIZED = "oversized-transfer"
/** CEP-41's `cvm.type`. */
const val TYPE_OPEN_STREAM = "open-stream"
/**
* Reads the envelope out of a notification, or returns null when this is
* an ordinary MCP progress notification rather than a ContextVM frame.
*
* Returning null rather than throwing is deliberate: a peer may send
* plain MCP progress for the same request, and that is not an error.
*/
fun parseOrNull(notification: JsonRpcNotification): ProgressEnvelope? {
if (notification.method != McpMethods.PROGRESS) return null
val params = notification.params ?: return null
val cvm = params[CVM] as? JsonObject ?: return null
val token = params[McpParams.PROGRESS_TOKEN]?.let { parseToken(it) } ?: return null
val progress =
(params[PROGRESS] as? JsonPrimitive)?.doubleOrNull
?: throw TransferFrameException("progress must be a number")
return ProgressEnvelope(
token = token,
progress = progress,
total = (params[TOTAL] as? JsonPrimitive)?.doubleOrNull,
message = params[MESSAGE]?.asStringOrNull(),
cvm = cvm,
)
}
private fun parseToken(element: kotlinx.serialization.json.JsonElement): ProgressToken? {
val primitive = element as? JsonPrimitive ?: return null
if (primitive.isString) return ProgressToken.Text(primitive.content)
return primitive.longOrNull?.let { ProgressToken.Num(it) }
}
internal fun ProgressToken.toPrimitive() =
when (this) {
is ProgressToken.Num -> JsonPrimitive(value)
is ProgressToken.Text -> JsonPrimitive(value)
}
internal fun kotlinx.serialization.json.JsonElement.asStringOrNull(): String? {
val primitive = this as? JsonPrimitive ?: return null
return if (primitive.isString) primitive.content else null
}
internal fun kotlinx.serialization.json.JsonElement.asLongOrNull(): Long? = (this as? JsonPrimitive)?.takeIf { !it.isString }?.longOrNull
}
}
@@ -0,0 +1,215 @@
/*
* 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.oversized
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-22 bounded oversized-transfer frame.
*
* The payload is reassembled from the `data` fragments of [Chunk] frames, which
* are **raw substrings of the serialized JSON-RPC message**, not base64. That is
* what makes the digest rule work: concatenate the fragments, encode the result
* as UTF-8, and hash those bytes.
*/
sealed interface OversizedFrame {
val envelope: ProgressEnvelope
val token: ProgressToken get() = envelope.token
val progress: Double get() = envelope.progress
data class Start(
override val envelope: ProgressEnvelope,
val completionMode: String,
val digest: String,
val totalBytes: Long,
val totalChunks: Long,
) : OversizedFrame
data class Accept(
override val envelope: ProgressEnvelope,
) : OversizedFrame
data class Chunk(
override val envelope: ProgressEnvelope,
val data: String,
) : OversizedFrame
data class End(
override val envelope: ProgressEnvelope,
) : OversizedFrame
data class Abort(
override val envelope: ProgressEnvelope,
val reason: String? = null,
) : OversizedFrame
companion object {
const val START = "start"
const val ACCEPT = "accept"
const val CHUNK = "chunk"
const val END = "end"
const val ABORT = "abort"
/** The only completion mode this CEP version defines. */
const val COMPLETION_MODE_RENDER = "render"
/** Digest values are algorithm-prefixed on the wire, e.g. `sha256:ab12…`. */
const val DIGEST_PREFIX_SHA256 = "sha256:"
const val COMPLETION_MODE = "completionMode"
const val DIGEST = "digest"
const val TOTAL_BYTES = "totalBytes"
const val TOTAL_CHUNKS = "totalChunks"
const val DATA = "data"
const val REASON = "reason"
/**
* Parses a CEP-22 frame, or returns null when the envelope belongs to a
* different transfer profile (CEP-41 shares this envelope).
*/
fun parseOrNull(envelope: ProgressEnvelope): OversizedFrame? {
if (envelope.type != ProgressEnvelope.TYPE_OVERSIZED) return null
val cvm = envelope.cvm
return when (val frameType = envelope.frameType) {
START -> {
val completionMode =
cvm[COMPLETION_MODE]?.asStringOrNull()
?: throw TransferFrameException("start requires completionMode")
// Receivers MUST reject unknown or unsupported completion
// modes; the field exists as an extension point for future
// CEPs, so silently treating anything else as render would
// be the wrong kind of tolerance.
if (completionMode != COMPLETION_MODE_RENDER) {
throw TransferFrameException("unsupported completionMode: $completionMode")
}
OversizedFrame.Start(
envelope = envelope,
completionMode = completionMode,
digest =
cvm[DIGEST]?.asStringOrNull()
?: throw TransferFrameException("start requires digest"),
totalBytes =
cvm[TOTAL_BYTES]?.asLongOrNull()
?: throw TransferFrameException("start requires totalBytes"),
totalChunks =
cvm[TOTAL_CHUNKS]?.asLongOrNull()
?: throw TransferFrameException("start requires totalChunks"),
)
}
ACCEPT -> OversizedFrame.Accept(envelope)
CHUNK ->
OversizedFrame.Chunk(
envelope = envelope,
data =
cvm[DATA]?.asStringOrNull()
?: throw TransferFrameException("chunk requires string data"),
)
END -> OversizedFrame.End(envelope)
ABORT -> OversizedFrame.Abort(envelope, cvm[REASON]?.asStringOrNull())
else -> throw TransferFrameException("unknown oversized frameType: $frameType")
}
}
fun start(
token: ProgressToken,
progress: Double,
digest: String,
totalBytes: Long,
totalChunks: Long,
message: String? = null,
) = OversizedFrame.Start(
envelope =
envelopeOf(START, token, progress, message) {
put(COMPLETION_MODE, JsonPrimitive(COMPLETION_MODE_RENDER))
put(DIGEST, JsonPrimitive(digest))
put(TOTAL_BYTES, JsonPrimitive(totalBytes))
put(TOTAL_CHUNKS, JsonPrimitive(totalChunks))
},
completionMode = COMPLETION_MODE_RENDER,
digest = digest,
totalBytes = totalBytes,
totalChunks = totalChunks,
)
fun accept(
token: ProgressToken,
progress: Double,
) = OversizedFrame.Accept(envelopeOf(ACCEPT, token, progress))
fun chunk(
token: ProgressToken,
progress: Double,
data: String,
) = OversizedFrame.Chunk(
envelope = envelopeOf(CHUNK, token, progress) { put(DATA, JsonPrimitive(data)) },
data = data,
)
fun end(
token: ProgressToken,
progress: Double,
message: String? = null,
) = OversizedFrame.End(envelopeOf(END, token, progress, message))
fun abort(
token: ProgressToken,
progress: Double,
reason: String? = null,
) = OversizedFrame.Abort(
envelope =
envelopeOf(ABORT, token, progress) {
reason?.let { put(REASON, JsonPrimitive(it)) }
},
reason = reason,
)
private fun envelopeOf(
frameType: String,
token: ProgressToken,
progress: Double,
message: String? = null,
body: JsonObjectBuilder.() -> Unit = {},
) = ProgressEnvelope(
token = token,
progress = progress,
message = message,
cvm =
buildJsonObject {
put(ProgressEnvelope.TYPE, JsonPrimitive(ProgressEnvelope.TYPE_OVERSIZED))
put(ProgressEnvelope.FRAME_TYPE, JsonPrimitive(frameType))
body()
},
)
}
}
@@ -0,0 +1,214 @@
/*
* 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.oversized
import com.vitorpamplona.contextvm.jsonrpc.JsonRpcCodec
import com.vitorpamplona.contextvm.jsonrpc.JsonRpcMessage
import com.vitorpamplona.contextvm.transfer.ProgressToken
import com.vitorpamplona.quartz.nip01Core.core.toHexKey
import com.vitorpamplona.quartz.utils.sha256.sha256
/**
* Admission-control limits for one receiver.
*
* CEP-22 requires evaluating the declared `totalBytes` and `totalChunks` against
* local policy **before** committing reassembly state, so a peer cannot make us
* allocate by announcing a huge transfer.
*/
data class OversizedLimits(
val maxTotalBytes: Long = 8L * 1024 * 1024,
val maxTotalChunks: Long = 4_096,
/**
* How many out-of-order chunks may be buffered while waiting for earlier
* ones. Relay delivery can reorder, so some buffering is required, but it
* has to be bounded or a peer can park unbounded memory by withholding one
* low-`progress` chunk.
*/
val maxPendingChunks: Int = 256,
)
/** Outcome of feeding one frame to [OversizedTransferReceiver]. */
sealed interface OversizedProgressResult {
/** Frame accepted; the transfer continues. Nothing is surfaced upward yet. */
data object Continue : OversizedProgressResult
/** The receiver should send an `accept` frame before chunks may flow. */
data object SendAccept : OversizedProgressResult
/** The transfer completed and validated. [message] is the reassembled payload. */
data class Completed(
val message: JsonRpcMessage,
val raw: String,
) : OversizedProgressResult
/** The transfer ended without producing a payload. */
data class Aborted(
val reason: String?,
) : OversizedProgressResult
}
/** Thrown when a transfer violates CEP-22. The transfer is terminal once this is raised. */
class OversizedTransferException(
message: String,
) : IllegalStateException(message)
/**
* Receives one CEP-22 bounded transfer, keyed by its `progressToken`.
*
* The rules this enforces are §6.5 `CVM-22-*` in
* `quartz/plans/2026-09-17-cordn-interop.md`. Two are worth stating here because
* they shape the design:
*
* - **Nothing is surfaced upward before validation succeeds.** The reassembled
* string only becomes a message after the chunk count, byte length and digest
* all agree, so a partially delivered payload can never be mistaken for a
* complete one.
* - **The digest covers the exact serialized string.** Fragments are
* concatenated and encoded to UTF-8 once; the payload is never parsed and
* re-serialized on the way, because that would reorder keys or change
* whitespace and break the hash.
*
* @param requireAccept true in stateless bootstrap, where the sender must wait
* for `accept` before the first chunk. When peer support is already known the
* sender may go straight from `start` to `chunk`, and this is false.
*/
class OversizedTransferReceiver(
val token: ProgressToken,
private val limits: OversizedLimits = OversizedLimits(),
private val requireAccept: Boolean = true,
) {
private var started: OversizedFrame.Start? = null
private var accepted = false
private var terminal = false
/** Fragments keyed by `progress`, so out-of-order arrivals assemble correctly. */
private val chunks = mutableMapOf<Double, String>()
val isTerminal get() = terminal
fun accept(frame: OversizedFrame): OversizedProgressResult {
if (terminal) throw OversizedTransferException("frame received after the transfer ended")
if (frame.token != token) throw OversizedTransferException("frame belongs to another transfer")
if (frame is OversizedFrame.Abort) {
terminal = true
return OversizedProgressResult.Aborted(frame.reason)
}
return when (frame) {
is OversizedFrame.Start -> onStart(frame)
is OversizedFrame.Accept -> OversizedProgressResult.Continue
is OversizedFrame.Chunk -> onChunk(frame)
is OversizedFrame.End -> onEnd(frame)
is OversizedFrame.Abort -> error("handled above")
}
}
private fun onStart(frame: OversizedFrame.Start): OversizedProgressResult {
if (started != null) fail("a second start frame arrived for this transfer")
// Admission control runs before any state is committed.
if (frame.totalBytes < 0) fail("totalBytes must not be negative")
if (frame.totalChunks < 0) fail("totalChunks must not be negative")
if (frame.totalBytes > limits.maxTotalBytes) {
fail("declared totalBytes ${frame.totalBytes} exceeds the limit ${limits.maxTotalBytes}")
}
if (frame.totalChunks > limits.maxTotalChunks) {
fail("declared totalChunks ${frame.totalChunks} exceeds the limit ${limits.maxTotalChunks}")
}
if (!frame.digest.startsWith(OversizedFrame.DIGEST_PREFIX_SHA256)) {
fail("unsupported digest algorithm: ${frame.digest.substringBefore(':')}")
}
started = frame
return if (requireAccept) OversizedProgressResult.SendAccept else OversizedProgressResult.Continue
}
private fun onChunk(frame: OversizedFrame.Chunk): OversizedProgressResult {
val start = started ?: fail("chunk arrived before start")
if (requireAccept && !accepted) {
// In stateless bootstrap the sender must not transmit before the
// receiver confirms. Seeing a chunk first means the peer skipped a
// step the profile requires, so the transfer cannot be trusted.
fail("chunk arrived before accept in a bootstrap transfer")
}
if (chunks.size >= limits.maxPendingChunks) {
fail("buffered chunk count exceeds the limit ${limits.maxPendingChunks}")
}
if (chunks.size.toLong() >= start.totalChunks) {
fail("more chunks arrived than the declared totalChunks ${start.totalChunks}")
}
// `progress` is the assembly index, so a chunk must sort after `start`.
// Arrival order is NOT checked: relays reorder, and the CEP explicitly
// says progress is not a guarantee of arrival order. Out-of-order
// frames are buffered and assembled by progress at `end`.
if (frame.progress <= start.progress) {
fail("chunk progress ${frame.progress} does not follow start at ${start.progress}")
}
if (chunks.put(frame.progress, frame.data) != null) {
fail("duplicate chunk at progress ${frame.progress}")
}
return OversizedProgressResult.Continue
}
private fun onEnd(frame: OversizedFrame.End): OversizedProgressResult {
val start = started ?: fail("end arrived before start")
val highestChunk = chunks.keys.maxOrNull()
if (highestChunk != null && frame.progress <= highestChunk) {
fail("end progress ${frame.progress} does not follow the last chunk at $highestChunk")
}
if (frame.progress <= start.progress) {
fail("end progress ${frame.progress} does not follow start at ${start.progress}")
}
if (chunks.size.toLong() != start.totalChunks) {
fail("expected ${start.totalChunks} chunks but assembled ${chunks.size}")
}
// `progress` is the canonical assembly index, not arrival order.
val raw = chunks.entries.sortedBy { it.key }.joinToString("") { it.value }
val bytes = raw.encodeToByteArray()
if (bytes.size.toLong() != start.totalBytes) {
fail("reassembled ${bytes.size} bytes but start declared ${start.totalBytes}")
}
val actual = OversizedFrame.DIGEST_PREFIX_SHA256 + sha256(bytes).toHexKey()
if (!actual.equals(start.digest, ignoreCase = true)) {
fail("digest mismatch: expected ${start.digest} but reassembled $actual")
}
terminal = true
return OversizedProgressResult.Completed(JsonRpcCodec.decode(raw), raw)
}
/** Records that we sent `accept`, unblocking chunk frames. */
fun markAccepted() {
accepted = true
}
private fun fail(reason: String): Nothing {
terminal = true
throw OversizedTransferException(reason)
}
}
@@ -0,0 +1,112 @@
/*
* 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.oversized
import com.vitorpamplona.contextvm.transfer.ProgressToken
import com.vitorpamplona.quartz.nip01Core.core.toHexKey
import com.vitorpamplona.quartz.utils.sha256.sha256
/**
* Splits an oversized serialized JSON-RPC message into CEP-22 frames.
*
* Two subtleties the chunking has to respect:
*
* - **Relay limits apply to the whole serialized event**, not just `content`,
* and roughly 64 KiB is the practical ceiling. [DEFAULT_CHUNK_CHARS] leaves
* generous room for the JSON-RPC envelope, the event's tags and signature,
* and any gift wrap around it.
* - **A chunk boundary must not split a surrogate pair.** Fragments are
* concatenated as text and only then encoded to UTF-8 for the digest, so
* cutting a pair in half would corrupt the reassembled bytes even though each
* fragment still looks like a valid JSON string.
*/
class OversizedTransferSender(
private val chunkChars: Int = DEFAULT_CHUNK_CHARS,
) {
init {
require(chunkChars > 1) { "chunkChars must leave room for a surrogate pair" }
}
/**
* Frames [serialized] for [token], starting at `progress` 1.
*
* The returned list is `start`, then the chunks, then `end`. When the peer's
* support is not yet known the caller must withhold everything after `start`
* until the peer's `accept` arrives; the frames themselves are unchanged.
*/
fun frame(
token: ProgressToken,
serialized: String,
): List<OversizedFrame> {
val pieces = split(serialized)
val bytes = serialized.encodeToByteArray()
val digest = OversizedFrame.DIGEST_PREFIX_SHA256 + sha256(bytes).toHexKey()
var progress = 1.0
val frames = mutableListOf<OversizedFrame>()
frames +=
OversizedFrame.start(
token = token,
progress = progress,
digest = digest,
totalBytes = bytes.size.toLong(),
totalChunks = pieces.size.toLong(),
)
pieces.forEach { piece ->
progress += 1.0
frames += OversizedFrame.chunk(token, progress, piece)
}
frames += OversizedFrame.end(token, progress + 1.0)
return frames
}
/** True when [serialized] is large enough to be worth fragmenting. */
fun shouldFragment(serialized: String) = serialized.length > chunkChars
private fun split(text: String): List<String> {
if (text.isEmpty()) return emptyList()
val pieces = mutableListOf<String>()
var index = 0
while (index < text.length) {
var end = minOf(index + chunkChars, text.length)
// Never cut between a high and low surrogate: the two halves would
// each survive JSON encoding but rejoin into a different codepoint.
if (end < text.length && text[end - 1].isHighSurrogate()) end -= 1
pieces += text.substring(index, end)
index = end
}
return pieces
}
companion object {
/**
* Conservative fragment size in UTF-16 characters.
*
* Well below the ~64 KiB practical relay ceiling because the fragment is
* only part of what ships: it is wrapped in a progress notification, then
* an event with tags and a signature, then possibly a gift wrap.
*/
const val DEFAULT_CHUNK_CHARS = 16 * 1024
}
}
@@ -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.oversized
import com.vitorpamplona.contextvm.jsonrpc.JsonRpcCodec
import com.vitorpamplona.contextvm.jsonrpc.JsonRpcId
import com.vitorpamplona.contextvm.jsonrpc.JsonRpcSuccess
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.assertIs
import kotlin.test.assertTrue
/**
* `CVM-22-*`: bounded oversized payload transfer.
*
* The CEP is written almost entirely as failure conditions, so most of this is
* negative. A suite of happy paths would prove nothing.
*/
class OversizedTransferTest {
private val token = ProgressToken.Text("req-123")
private val payload =
JsonRpcCodec.encode(
JsonRpcSuccess(
JsonRpcId.Num(1),
buildJsonObject { put("text", JsonPrimitive("a".repeat(200))) },
),
)
private fun sender(chunkChars: Int = 64) = OversizedTransferSender(chunkChars)
private fun receiver(
limits: OversizedLimits = OversizedLimits(),
requireAccept: Boolean = false,
) = OversizedTransferReceiver(token, limits, requireAccept)
// --- happy paths ---
@Test
fun `CVM-22-01 round-trips a fragmented message through the frames`() {
val frames = sender().frame(token, payload)
val rx = receiver()
var completed: OversizedProgressResult.Completed? = null
frames.forEach { frame ->
val result = rx.accept(frame)
if (result is OversizedProgressResult.Completed) completed = result
}
assertEquals(payload, completed?.raw)
assertEquals(JsonRpcCodec.decode(payload), completed?.message)
}
@Test
fun `CVM-22-02 assembles out-of-order chunks by progress, not arrival order`() {
// Relays may reorder. `progress` is the canonical assembly index.
val frames = sender().frame(token, payload)
val start = frames.first()
val end = frames.last()
val chunks = frames.drop(1).dropLast(1).reversed()
val rx = receiver()
rx.accept(start)
chunks.forEach { rx.accept(it) }
val result = rx.accept(end)
assertIs<OversizedProgressResult.Completed>(result)
assertEquals(payload, result.raw)
}
@Test
fun `CVM-22-03 surfaces nothing before validation succeeds`() {
val frames = sender().frame(token, payload)
val rx = receiver()
frames.dropLast(1).forEach {
assertTrue(
rx.accept(it) is OversizedProgressResult.Continue,
"no payload may be surfaced before end validates",
)
}
}
@Test
fun `CVM-22-04 a bootstrap transfer asks the receiver to send accept`() {
val frames = sender().frame(token, payload)
val rx = receiver(requireAccept = true)
assertIs<OversizedProgressResult.SendAccept>(rx.accept(frames.first()))
}
@Test
fun `CVM-22-05 abort is terminal`() {
val rx = receiver()
rx.accept(sender().frame(token, payload).first())
val aborted = rx.accept(OversizedFrame.abort(token, 99.0, "peer gave up"))
assertIs<OversizedProgressResult.Aborted>(aborted)
assertEquals("peer gave up", aborted.reason)
assertTrue(rx.isTerminal)
assertFailsWith<OversizedTransferException> {
rx.accept(OversizedFrame.chunk(token, 100.0, "x"))
}
}
@Test
fun `CVM-22-06 the sender never splits a surrogate pair`() {
// Each emoji is a surrogate pair. Splitting one would still produce two
// valid JSON strings, so only the reassembled bytes catch it.
val emoji = "🚀".repeat(40)
val text = JsonRpcCodec.encode(JsonRpcSuccess(JsonRpcId.Num(1), JsonPrimitive(emoji)))
val frames = sender(chunkChars = 9).frame(token, text)
frames.filterIsInstance<OversizedFrame.Chunk>().forEach { chunk ->
assertTrue(
chunk.data.isEmpty() || !chunk.data.last().isHighSurrogate(),
"a fragment must not end on an unpaired high surrogate",
)
}
val rx = receiver()
val result = frames.map { rx.accept(it) }.last()
assertIs<OversizedProgressResult.Completed>(result)
assertEquals(text, result.raw)
}
// --- rejections ---
@Test
fun `CVM-22-10 rejects a digest mismatch`() {
val frames = sender().frame(token, payload).toMutableList()
val start = frames.first() as OversizedFrame.Start
frames[0] =
OversizedFrame.start(
token,
start.progress,
OversizedFrame.DIGEST_PREFIX_SHA256 + "00".repeat(32),
start.totalBytes,
start.totalChunks,
)
val rx = receiver()
assertFailsWith<OversizedTransferException> { frames.forEach { rx.accept(it) } }
}
@Test
fun `CVM-22-11 rejects a totalBytes mismatch`() {
val frames = sender().frame(token, payload).toMutableList()
val start = frames.first() as OversizedFrame.Start
frames[0] =
OversizedFrame.start(token, start.progress, start.digest, start.totalBytes + 1, start.totalChunks)
val rx = receiver()
assertFailsWith<OversizedTransferException> { frames.forEach { rx.accept(it) } }
}
@Test
fun `CVM-22-12 rejects a totalChunks mismatch`() {
val frames = sender().frame(token, payload).toMutableList()
val start = frames.first() as OversizedFrame.Start
frames[0] =
OversizedFrame.start(token, start.progress, start.digest, start.totalBytes, start.totalChunks + 1)
val rx = receiver()
assertFailsWith<OversizedTransferException> { frames.forEach { rx.accept(it) } }
}
@Test
fun `CVM-22-13 rejects a chunk before accept in a bootstrap transfer`() {
val frames = sender().frame(token, payload)
val rx = receiver(requireAccept = true)
rx.accept(frames.first())
assertFailsWith<OversizedTransferException> { rx.accept(frames[1]) }
}
@Test
fun `CVM-22-14 rejects non-monotonic progress`() {
val rx = receiver()
rx.accept(OversizedFrame.start(token, 5.0, digestOf(""), 0, 0))
assertFailsWith<OversizedTransferException> {
rx.accept(OversizedFrame.chunk(token, 4.0, "x"))
}
}
@Test
fun `CVM-22-15 rejects end with unresolved gaps`() {
val frames = sender().frame(token, payload)
val rx = receiver()
rx.accept(frames.first())
// Deliver every chunk but the last, then end.
frames.drop(1).dropLast(2).forEach { rx.accept(it) }
assertFailsWith<OversizedTransferException> { rx.accept(frames.last()) }
}
@Test
fun `CVM-22-16 rejects an unknown completionMode at parse time`() {
val envelope =
ProgressEnvelope(
token = token,
progress = 1.0,
cvm =
buildJsonObject {
put(ProgressEnvelope.TYPE, JsonPrimitive(ProgressEnvelope.TYPE_OVERSIZED))
put(ProgressEnvelope.FRAME_TYPE, JsonPrimitive(OversizedFrame.START))
put(OversizedFrame.COMPLETION_MODE, JsonPrimitive("stream-someday"))
put(OversizedFrame.DIGEST, JsonPrimitive(digestOf("")))
put(OversizedFrame.TOTAL_BYTES, JsonPrimitive(0))
put(OversizedFrame.TOTAL_CHUNKS, JsonPrimitive(0))
},
)
assertFailsWith<TransferFrameException> { OversizedFrame.parseOrNull(envelope) }
}
@Test
fun `CVM-22-17 rejects declared totals over local policy at start`() {
val rx = receiver(limits = OversizedLimits(maxTotalBytes = 10, maxTotalChunks = 2))
assertFailsWith<OversizedTransferException> {
rx.accept(OversizedFrame.start(token, 1.0, digestOf(""), 1_000_000, 1))
}
}
@Test
fun `CVM-22-18 rejects a chunk arriving before start`() {
val rx = receiver()
assertFailsWith<OversizedTransferException> {
rx.accept(OversizedFrame.chunk(token, 1.0, "x"))
}
}
@Test
fun `CVM-22-19 rejects a second start on the same transfer`() {
val rx = receiver()
rx.accept(OversizedFrame.start(token, 1.0, digestOf(""), 0, 0))
assertFailsWith<OversizedTransferException> {
rx.accept(OversizedFrame.start(token, 2.0, digestOf(""), 0, 0))
}
}
@Test
fun `CVM-22-20 rejects a frame belonging to another transfer`() {
val rx = receiver()
assertFailsWith<OversizedTransferException> {
rx.accept(OversizedFrame.start(ProgressToken.Text("other"), 1.0, digestOf(""), 0, 0))
}
}
@Test
fun `CVM-22-21 rejects an unsupported digest algorithm`() {
val rx = receiver()
assertFailsWith<OversizedTransferException> {
rx.accept(OversizedFrame.start(token, 1.0, "md5:" + "00".repeat(16), 0, 0))
}
}
@Test
fun `CVM-22-22 rejects a duplicate chunk at the same progress`() {
val rx = receiver()
rx.accept(OversizedFrame.start(token, 1.0, digestOf("ab"), 2, 2))
rx.accept(OversizedFrame.chunk(token, 2.0, "a"))
assertFailsWith<OversizedTransferException> {
// Same progress value, so this cannot be a distinct fragment.
rx.accept(OversizedFrame.chunk(token, 2.0, "b"))
}
}
@Test
fun `CVM-22-23 the envelope round-trips through a progress notification`() {
val frame = sender().frame(token, payload).first()
val notification = frame.envelope.toNotification()
val decoded = JsonRpcCodec.decode(JsonRpcCodec.encode(notification))
val envelope = ProgressEnvelope.parseOrNull(decoded as com.vitorpamplona.contextvm.jsonrpc.JsonRpcNotification)
assertEquals(frame, OversizedFrame.parseOrNull(envelope!!))
}
@Test
fun `CVM-22-24 ignores an envelope belonging to another transfer profile`() {
val envelope =
ProgressEnvelope(
token = token,
progress = 1.0,
cvm =
buildJsonObject {
put(ProgressEnvelope.TYPE, JsonPrimitive(ProgressEnvelope.TYPE_OPEN_STREAM))
put(ProgressEnvelope.FRAME_TYPE, JsonPrimitive("start"))
},
)
assertEquals(null, OversizedFrame.parseOrNull(envelope))
}
private fun digestOf(text: String) =
OversizedFrame.DIGEST_PREFIX_SHA256 +
com.vitorpamplona.quartz.utils.sha256
.sha256(text.encodeToByteArray())
.let { bytes -> bytes.joinToString("") { b -> (b.toInt() and 0xFF).toString(16).padStart(2, '0') } }
}