diff --git a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/CordnGroupCommands.kt b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/CordnGroupCommands.kt index 4513db1de8..53e3868b3f 100644 --- a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/CordnGroupCommands.kt +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/CordnGroupCommands.kt @@ -616,6 +616,10 @@ internal object CordnGroupCommands { "epoch_changes" to epochs, "echoes" to echoes, "undecryptable" to undecryptable, + // A fetch is the response most likely to outgrow one relay + // event, so this is where CEP-22 shows up if it shows up at + // all. Zero is the normal answer and not a warning. + "oversized_transfers" to scope.session.oversizedTransfers, ), ) 0 diff --git a/cli/tests/cordn/tier-b.sh b/cli/tests/cordn/tier-b.sh index 62c8432083..6c02532b7d 100755 --- a/cli/tests/cordn/tier-b.sh +++ b/cli/tests/cordn/tier-b.sh @@ -42,6 +42,7 @@ bad() { echo " FAIL: $*"; fail=1; } # the only thing carrying continuity between them. alice() { HOME="$WORK/alice" "$AMY" --account alice --secret-backend ncryptsec "$@" 2>/dev/null; } bob() { HOME="$WORK/bob" "$AMY" --account bob --secret-backend ncryptsec "$@" 2>/dev/null; } +carol() { HOME="$WORK/carol" "$AMY" --account carol --secret-backend ncryptsec "$@" 2>/dev/null; } field() { python3 -c "import json,sys; d=json.load(sys.stdin); print(json.dumps(d$1) if not isinstance(d$1,str) else d$1)"; } trap stack_down EXIT @@ -52,18 +53,24 @@ step "boot geode on $RELAY, and the reference coordinator" stack_up ok "coordinator $COORD" -step "two accounts" -mkdir -p "$WORK/alice" "$WORK/bob" +step "three accounts" +# Carol is here for the door nobody knocks on in the two-party flow: bob is +# invited, so he never sends a join request. +mkdir -p "$WORK/alice" "$WORK/bob" "$WORK/carol" alice create --json >/dev/null bob create --json >/dev/null +carol create --json >/dev/null ALICE_PK=$(alice whoami --json | field "['hex']") BOB_PK=$(bob whoami --json | field "['hex']") +CAROL_PK=$(carol whoami --json | field "['hex']") ok "alice $ALICE_PK" ok "bob $BOB_PK" +ok "carol $CAROL_PK" -step "both remember the coordinator" +step "all three remember the coordinator" alice cordn coordinator add --coordinator "$COORD" --relay "$RELAY" --label tier-b --json >/dev/null bob cordn coordinator add --coordinator "$COORD" --relay "$RELAY" --label tier-b --json >/dev/null +carol cordn coordinator add --coordinator "$COORD" --relay "$RELAY" --label tier-b --json >/dev/null step "the MCP handshake" # The first thing to break if the transport regresses, and the cheapest to @@ -119,6 +126,75 @@ echo "$GOT2" | grep -q "$BACK_ID" && ok "alice read the reply" || bad "alice did # and recognising it needs bookkeeping that survives the process exit. [ "$(echo "$GOT2" | field "['undecryptable']")" = "[]" ] && ok "no self-inflicted gaps" || bad "own traffic came back undecryptable: $(echo "$GOT2" | field "['undecryptable']")" +# ── The three tools the two-party flow never reaches, and CEP-22 ──────────── +# Every other step above exercises a tool as a side effect of the lifecycle. +# These four do not happen on that path, so they are driven directly. + +step "carol asks to join — join_request_store" +REQ=$(carol cordn request --gid "$GID" --json) +echo " $REQ" +[ "$(echo "$REQ" | field "['gid']")" = "$GID" ] && ok "request stored" || bad "join request not stored" +# A ref is an invitation to ask, not membership (spec/01.md §5.3). +[ "$(echo "$REQ" | field "['member']")" = "false" ] && ok "asking is not joining" || bad "request reported membership" + +step "alice reads it — join_request_take_many" +PEND=$(alice cordn requests list --json) +echo " $PEND" +[ "$(echo "$PEND" | field "['requests'][0]['pubkey']")" = "$CAROL_PK" ] && ok "sees carol" || bad "no pending request" +[ "$(echo "$PEND" | field "['requests'][0]['gid']")" = "$GID" ] && ok "for $GID" || bad "wrong gid on the request" +# Listing must NOT consume. A request is retired by answering it, not by +# reading it, and the ack rides the next call — so a reader that lost the +# process between list and accept has to still find it. Two separate amy +# runs is exactly that case. +AGAIN=$(alice cordn requests list --json) +[ "$(echo "$AGAIN" | field "['requests'][0]['pubkey']")" = "$CAROL_PK" ] && ok "still there for a second reader" || bad "reading a join request consumed it" + +step "alice accepts, carol joins" +ACC=$(alice cordn requests accept --pubkey "$CAROL_PK" --json) +echo " $ACC" +[ "$(echo "$ACC" | field "['answered'][0]['pubkey']")" = "$CAROL_PK" ] && ok "accepted" || bad "accept failed" +CJOIN=$(carol cordn join --all --json) +[ "$(echo "$CJOIN" | field "['joined']")" = "[\"$GID\"]" ] && ok "carol joined" || bad "carol did not join" +# Now it is retired — and the ack for it rode a later call, so this also +# proves the retirement survived the process that issued it. +GONE=$(alice cordn requests list --json) +[ "$(echo "$GONE" | field "['requests']")" = "[]" ] && ok "answering retired it" || bad "an answered request is still pending: $(echo "$GONE" | field "['requests']")" + +step "bob withdraws a KeyPackage — kp_remove" +# Two, so there is something left to prove the removal was targeted rather +# than a wipe. +PUB=$(bob cordn keypackage publish --count 2 --json) +KEEP=$(echo "$PUB" | field "['published'][0]['kp_ref']") +DROP=$(echo "$PUB" | field "['published'][1]['kp_ref']") +WD=$(bob cordn keypackage withdraw --kp-ref "$DROP" --json) +echo " $WD" +[ "$(echo "$WD" | field "['removed']")" = "[\"$DROP\"]" ] && ok "coordinator confirms the removal" || bad "kp_remove did not confirm $DROP" +LIST=$(bob cordn keypackage list --json) +echo "$LIST" | grep -q "$DROP" && bad "the withdrawn package is still served" || ok "gone from kp_list" +echo "$LIST" | grep -q "$KEEP" && ok "the other one survived" || bad "kp_remove took more than it was asked for" + +step "a response too big for one event — CEP-22" +# The reference coordinator switches to oversized transfer at a 48000-byte +# published envelope. One message does not reach it; a backlog of them does, +# and a fetch is where a backlog is delivered. Nothing below reads the +# messages differently — a reassembled response is meant to be invisible — +# so the count is the only thing that can tell us the profile ran. +BIG=$(python3 -c "print('x' * 6000)") +for i in $(seq 12); do + alice cordn send --text "chunk-$i $BIG" --json >/dev/null +done +FETCH=$(bob cordn fetch --json) +COUNT=$(echo "$FETCH" | python3 -c "import json,sys; print(len(json.load(sys.stdin)['messages']))") +CHUNKED=$(echo "$FETCH" | field "['oversized_transfers']") +echo " $COUNT messages, oversized_transfers=$CHUNKED" +[ "$COUNT" = "12" ] && ok "all 12 arrived" || bad "expected 12 messages, got $COUNT" +[ "${CHUNKED:-0}" -ge 1 ] && ok "reassembled over CEP-22" || bad "the response was never chunked — CEP-22 went unexercised" +# Carol was added at a later epoch, so she must read them too: an oversized +# response is still one response, not a per-member special case. +CGOT=$(carol cordn fetch --json) +CCOUNT=$(echo "$CGOT" | python3 -c "import json,sys; print(len(json.load(sys.stdin)['messages']))") +[ "$CCOUNT" = "12" ] && ok "carol read the same 12" || bad "carol got $CCOUNT of 12" + step "both sides agree" A=$(alice cordn group info --json) B=$(bob cordn group info --json) diff --git a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnCoordinatorRegistry.kt b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnCoordinatorRegistry.kt index 6235214f35..29cf8c16e9 100644 --- a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnCoordinatorRegistry.kt +++ b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnCoordinatorRegistry.kt @@ -88,6 +88,9 @@ class CordnSession( /** What this coordinator says about itself. Claims, never identity (§8.5). */ suspend fun serverInfo(): CoordinatorServerInfo? = scope.serverInfo() + /** How many of this session's responses arrived reassembled over CEP-22. */ + val oversizedTransfers: Int get() = scope.coordinator.oversizedTransfers + internal suspend fun close() = scope.close() } diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/contextvm/cep04Encryption/CvmGiftWrap.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/contextvm/cep04Encryption/CvmGiftWrap.kt index d086d45a73..a758ad2a8c 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/contextvm/cep04Encryption/CvmGiftWrap.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/contextvm/cep04Encryption/CvmGiftWrap.kt @@ -85,6 +85,23 @@ class CvmGiftWrap( private val encryptionMode: EncryptionMode = EncryptionMode.REQUIRED, private val giftWrapMode: GiftWrapMode = GiftWrapMode.EPHEMERAL, ) { + /** + * What this side can *receive*, as CEP-35 tag names. + * + * Receive, not prefer: the transport subscribes to both wrap kinds, so a + * client that emits 1059 still accepts 21059 and says so. Declaring only + * what we emit would tell a peer to withhold something we can read. + * + * Empty under [EncryptionMode.DISABLED] - there, a wrap is not something + * this client will open at all. + */ + fun receivableWrapTags(): List = + if (encryptionMode == EncryptionMode.DISABLED) { + emptyList() + } else { + listOf(CvmTags.SUPPORT_ENCRYPTION, CvmTags.SUPPORT_ENCRYPTION_EPHEMERAL) + } + /** The wrap kind this client emits. */ fun wrapKind(): Kind = when (giftWrapMode) { diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/contextvm/mcp/CvmMcpClient.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/contextvm/mcp/CvmMcpClient.kt index afc5fdc803..9e323f3171 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/contextvm/mcp/CvmMcpClient.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/contextvm/mcp/CvmMcpClient.kt @@ -38,6 +38,7 @@ import com.vitorpamplona.quartz.contextvm.transfer.ProgressEnvelope import com.vitorpamplona.quartz.contextvm.transfer.ProgressToken import com.vitorpamplona.quartz.contextvm.transport.CvmTransport import com.vitorpamplona.quartz.contextvm.transport.DualSigner +import com.vitorpamplona.quartz.contextvm.transport.TimeoutMode import com.vitorpamplona.quartz.nip01Core.core.Tag import kotlinx.serialization.json.JsonElement import kotlinx.serialization.json.JsonObject @@ -50,6 +51,15 @@ data class ToolCallResult( val error: com.vitorpamplona.quartz.contextvm.jsonrpc.JsonRpcError? = null, /** Fragments a CEP-41 stream delivered while the call was in flight. */ val streamed: List = emptyList(), + /** + * Whether the response arrived reassembled over CEP-22 rather than in the + * single event that carried the rest. + * + * A caller does not need this to read the result - that is the point of the + * profile. It is here so a live interop run can assert the profile actually + * fired, instead of passing whether or not the server ever chunked. + */ + val viaOversizedTransfer: Boolean = false, ) { val isError get() = error != null } @@ -83,7 +93,7 @@ class CvmMcpClient( * afford the round trip should do it. */ suspend fun initialize( - capabilityTags: List = emptyList(), + capabilityTags: List = transport.selfDiscoveryTags(), protocolVersion: String = PROTOCOL_VERSION, ): JsonRpcMessage { val response = @@ -134,6 +144,7 @@ class CvmMcpClient( arguments: JsonObject = buildJsonObject {}, identity: DualSigner.Identity = DualSigner.Identity.EPHEMERAL, timeoutMs: Long = CvmTransport.DEFAULT_TIMEOUT_MS, + timeoutMode: TimeoutMode = TimeoutMode.IDLE, onStreamFragment: (String) -> Unit = {}, ): ToolCallResult { val id = nextId() @@ -164,9 +175,10 @@ class CvmMcpClient( ), identity = identity, timeoutMs = timeoutMs, + timeoutMode = timeoutMode, ) { notification -> - val envelope = ProgressEnvelope.parseOrNull(notification) ?: return@request - if (envelope.token != token) return@request + val envelope = ProgressEnvelope.parseOrNull(notification) ?: return@request null + if (envelope.token != token) return@request null when (envelope.type) { ProgressEnvelope.TYPE_OVERSIZED -> { @@ -175,7 +187,10 @@ class CvmMcpClient( .also { oversized = it } OversizedFrame.parseOrNull(envelope)?.let { frame -> val result = receiver.accept(frame) - if (result is OversizedProgressResult.Completed) reassembled = result.message + if (result is OversizedProgressResult.Completed) { + reassembled = result.message + oversizedTransfersCompleted++ + } } } @@ -194,20 +209,35 @@ class CvmMcpClient( else -> Unit } + + // A completed CEP-22 transfer IS the response, so handing it + // back ends the call. A CEP-41 stream never does: `close` says + // no more frames, not that the request is answered. + reassembled } // A CEP-22 transfer replaces the response that could not be published // directly; a CEP-41 stream does not, since `close` never completes the // JSON-RPC request. + val chunked = reassembled != null return when (val effective = reassembled ?: response) { - is JsonRpcSuccess -> ToolCallResult(effective.result, streamed = streamed) - is JsonRpcFailure -> ToolCallResult(null, effective.error, streamed) - else -> ToolCallResult(null, streamed = streamed) + is JsonRpcSuccess -> ToolCallResult(effective.result, streamed = streamed, viaOversizedTransfer = chunked) + is JsonRpcFailure -> ToolCallResult(null, effective.error, streamed, chunked) + else -> ToolCallResult(null, streamed = streamed, viaOversizedTransfer = chunked) } } private fun nextId(): JsonRpcId.Num = JsonRpcId.Num(nextId++) + /** + * How many responses this client has reassembled over CEP-22. + * + * Diagnostics, not control flow. Nothing decides anything on it; it exists + * so a live run can tell a server that chunked from one that never had to. + */ + var oversizedTransfersCompleted: Int = 0 + private set + companion object { /** The MCP revision the ContextVM spec's examples use. */ const val PROTOCOL_VERSION = "2025-07-02" diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/contextvm/transport/CvmTransport.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/contextvm/transport/CvmTransport.kt index a344e4fd26..cfda426607 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/contextvm/transport/CvmTransport.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/contextvm/transport/CvmTransport.kt @@ -24,6 +24,7 @@ import com.vitorpamplona.quartz.contextvm.cep04Encryption.CvmGiftWrap import com.vitorpamplona.quartz.contextvm.cep35Discovery.SessionDiscovery import com.vitorpamplona.quartz.contextvm.core.CvmKinds import com.vitorpamplona.quartz.contextvm.core.CvmMessageEvent +import com.vitorpamplona.quartz.contextvm.core.CvmTags import com.vitorpamplona.quartz.contextvm.jsonrpc.JsonRpcFailure import com.vitorpamplona.quartz.contextvm.jsonrpc.JsonRpcId import com.vitorpamplona.quartz.contextvm.jsonrpc.JsonRpcMessage @@ -63,6 +64,29 @@ class DualSigner( } } +/** What a call's `timeoutMs` bounds. */ +enum class TimeoutMode { + /** + * Silence. The clock restarts on every message the peer sends, so a + * response still being delivered is never cut off. + * + * The right default for a request/response call: both transfer profiles + * deliver one logical response as a run of notifications, each its own + * signed, wrapped, published relay event, so how long a response takes is + * a function of its size and not something a caller can predict. + */ + IDLE, + + /** + * The whole call, from publish to return. + * + * For an open-ended subscription, where the timeout is not a failure but + * the budget after which the caller re-opens - a busy stream under [IDLE] + * would simply never come back. + */ + TOTAL, +} + /** Thrown when a request cannot be completed at the transport layer. */ class CvmTransportException( message: String, @@ -102,6 +126,23 @@ class CvmTransport( private val assumePeerSupportsEncryption: Boolean = true, private val assumePeerSupportsEphemeralWrap: Boolean = true, ) { + /** + * This side's CEP-35 surface: what a peer may use against us. + * + * CEP-35 is symmetric, and the half we were not doing is the expensive one + * to skip. A spec-correct server only chunks a CEP-22 response, or opens a + * CEP-41 stream, for a client that declared it can take one — so declaring + * nothing does not make us conservative, it makes those profiles dead on + * every session and caps every response at one relay event. + * + * Both transfer profiles are unconditional: [CvmMcpClient] reassembles + * CEP-22 and collects CEP-41 fragments on every call, with no flag to turn + * either off. + */ + fun selfDiscoveryTags(): List = + (crypto.receivableWrapTags() + CvmTags.SUPPORT_OVERSIZED_TRANSFER + CvmTags.SUPPORT_OPEN_STREAM) + .map { arrayOf(it) } + /** The peer's learned discovery baseline, once its first message has arrived. */ val peer get() = discovery.peer @@ -128,15 +169,25 @@ class CvmTransport( * * @param identity which key signs the request. Anything not required to be * attributable should stay [DualSigner.Identity.EPHEMERAL]. + * @param onNotification every notification that arrives while the call is + * open. Returning a message **completes the call with it**, which is how + * CEP-22 finishes: a chunked response replaces the direct one rather than + * preceding it, so nothing else would ever arrive to end the wait. + * Returning null keeps waiting - CEP-41 is explicit that a stream's + * `close` does not complete the request. * @param discoveryTags this side's CEP-35 baseline, sent on the session's - * first direct message only. + * first direct message only. Defaults to [selfDiscoveryTags] so the + * surface rides whatever that first message turns out to be - a session + * that opens with a tool call rather than a handshake still declares. + * Pass `emptyList()` to declare nothing. */ suspend fun request( message: JsonRpcRequest, identity: DualSigner.Identity = DualSigner.Identity.EPHEMERAL, timeoutMs: Long = DEFAULT_TIMEOUT_MS, - discoveryTags: List = emptyList(), - onNotification: (JsonRpcNotification) -> Unit = {}, + timeoutMode: TimeoutMode = TimeoutMode.IDLE, + discoveryTags: List = selfDiscoveryTags(), + onNotification: (JsonRpcNotification) -> JsonRpcMessage? = { null }, ): JsonRpcMessage { val signer = signers.signerFor(identity) val inbound = Channel(Channel.UNLIMITED) @@ -160,8 +211,12 @@ class CvmTransport( relays.publish(outbound(request)) - return withTimeout(timeoutMs) { - awaitResponse(inbound, signer, request.id, message.id, onNotification) + return when (timeoutMode) { + TimeoutMode.IDLE -> awaitResponse(inbound, signer, request.id, message.id, timeoutMs, onNotification) + TimeoutMode.TOTAL -> + withTimeout(timeoutMs) { + awaitResponse(inbound, signer, request.id, message.id, Long.MAX_VALUE, onNotification) + } } } finally { subscription.close() @@ -178,16 +233,43 @@ class CvmTransport( relays.publish(outbound(CvmMessageEvent.create(message, serverPubKey, signer))) } + /** + * Waits for the answer, allowing [idleMs] of **silence** between the peer's + * messages. [TimeoutMode.TOTAL] passes `Long.MAX_VALUE` and wraps the whole + * thing instead. + * + * A flat deadline cannot express what a CEP-22 or CEP-41 response is. Both + * arrive as a run of notifications, each its own signed, wrapped, published + * relay event, so a large response takes as long as it takes: a 100 KB + * oversized transfer from the reference coordinator ran well past the 20 s + * default and was killed mid-run, with the frames arriving and parsing + * correctly right up to the cancellation. + * + * So the clock measures a stalled peer, which is what a caller actually + * wants bounded, and only the peer can reset it - events from anyone else + * are dropped before this point, or a stranger could hold a call open + * indefinitely by publishing noise addressed to us. + */ private suspend fun awaitResponse( inbound: Channel, signer: NostrSigner, requestEventId: HexKey, requestId: JsonRpcId, - onNotification: (JsonRpcNotification) -> Unit, + idleMs: Long, + onNotification: (JsonRpcNotification) -> JsonRpcMessage?, ): JsonRpcMessage { - for (event in inbound) { + while (true) { + // withTimeout, not a null-returning variant, so silence still + // surfaces as the TimeoutCancellationException callers already + // handle - only when the clock starts has changed. + // receiveCatching tells a closed subscription from a silent one. + val event = + withTimeout(idleMs) { inbound.receiveCatching() } + .getOrNull() + ?: throw CvmTransportException("subscription closed before a response arrived") val plain = decryptOrNull(event, signer) ?: continue val wrapped = CvmMessageEvent.fromOrNull(plain) ?: continue + if (wrapped.pubKey != serverPubKey) continue discovery.observe(wrapped.discoveryTags().toTypedArray()) @@ -201,7 +283,7 @@ class CvmTransport( } when (decoded) { - is JsonRpcNotification -> onNotification(decoded) + is JsonRpcNotification -> onNotification(decoded)?.let { return it } // Correlate on both layers: the `e` tag ties the response to our // request event, and the JSON-RPC id ties it to our call. Either @@ -216,7 +298,6 @@ class CvmTransport( else -> Unit } } - throw CvmTransportException("subscription closed before a response arrived") } /** diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/cordn/spec00Coordinator/CoordinatorClient.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/cordn/spec00Coordinator/CoordinatorClient.kt index 7606add362..8e2044f6de 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/cordn/spec00Coordinator/CoordinatorClient.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/cordn/spec00Coordinator/CoordinatorClient.kt @@ -22,6 +22,7 @@ package com.vitorpamplona.quartz.cordn.spec00Coordinator import com.vitorpamplona.quartz.contextvm.mcp.CvmMcpClient import com.vitorpamplona.quartz.contextvm.transport.CvmTransport +import com.vitorpamplona.quartz.contextvm.transport.TimeoutMode import com.vitorpamplona.quartz.nip01Core.core.Event import com.vitorpamplona.quartz.nip01Core.core.HexKey import com.vitorpamplona.quartz.nip01Core.core.OptimizedJsonMapper @@ -68,6 +69,8 @@ class CoordinatorClient( /** Performs the MCP handshake. Optional, but it is where CEP-35 tags ride. */ suspend fun initialize() = mcp.initialize() + override val oversizedTransfers get() = mcp.oversizedTransfersCompleted + /** * Handshakes and reports what the coordinator says about itself. * @@ -316,6 +319,9 @@ class CoordinatorClient( arguments = args, identity = CoordinatorMethod.MSG_SUB_MANY.identity, timeoutMs = timeoutMs, + // A budget, not a failure: the caller re-opens on it, and a + // busy stream bounded by silence would never come back. + timeoutMode = TimeoutMode.TOTAL, onStreamFragment = { fragment -> // A malformed frame is the coordinator's problem, not a // reason to tear down a live subscription over other groups. diff --git a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/cordn/spec00Coordinator/ICoordinator.kt b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/cordn/spec00Coordinator/ICoordinator.kt index d9558aefa3..36746a4d04 100644 --- a/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/cordn/spec00Coordinator/ICoordinator.kt +++ b/quartz/src/commonMain/kotlin/com/vitorpamplona/quartz/cordn/spec00Coordinator/ICoordinator.kt @@ -78,4 +78,17 @@ interface ICoordinator { timeoutMs: Long, onMessage: (GroupMessage) -> Unit, ) + + /** + * How many responses this coordinator chunked over CEP-22 so far. + * + * Diagnostics, and the one thing here that is not a tool. Nothing in the + * binding reads it; a live interop run does, to tell a server that had to + * chunk from one that never did - which is otherwise invisible, because a + * reassembled response is indistinguishable from a direct one by design. + * + * Defaulted so a substitute implementation is not forced to fake a number + * about a transfer profile it does not have. + */ + val oversizedTransfers: Int get() = 0 } diff --git a/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/contextvm/mcp/CvmMcpClientTest.kt b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/contextvm/mcp/CvmMcpClientTest.kt index a72fadb05d..6b693a8699 100644 --- a/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/contextvm/mcp/CvmMcpClientTest.kt +++ b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/contextvm/mcp/CvmMcpClientTest.kt @@ -221,12 +221,13 @@ class CvmMcpClientTest { val big = buildJsonObject { put("text", JsonPrimitive("x".repeat(400))) } val serialized = JsonRpcCodec.encode(JsonRpcSuccess(JsonRpcId.Num(0), big)) - val fixture = - fixture { - // A placeholder direct response: the real payload arrives - // through the frames, which is the point of the profile. - JsonRpcSuccess(JsonRpcId.Num(0), buildJsonObject { put("placeholder", JsonPrimitive(true)) }) - } + // NO direct response, and that is the whole point: a chunked + // response REPLACES the direct one. An earlier version of this test + // had the fixture also answer normally, which meant the call was + // ended by that answer and the reassembly only had to win a + // tie-break. Against the reference coordinator, which sends frames + // and nothing else, the same code hung until its deadline. + val fixture = fixture { error("a chunked response is the only response") } fixture.start() val result = @@ -234,11 +235,9 @@ class CvmMcpClientTest { val pending = async { client().callTool("big", timeoutMs = 5_000) } yield() - // Frames first, then the direct response. OversizedTransferSender(chunkChars = 64).frame(firstCallToken, serialized).forEach { frame -> fixture.reply(frame.envelope.toNotification(), clientSigner.pubKey, "0".repeat(64)) } - fixture.pump() pending.await() } @@ -247,7 +246,7 @@ class CvmMcpClientTest { result.result!! .jsonObject["text"]!! .jsonPrimitive.content, - "the reassembled payload replaces the placeholder", + "the reassembled payload is the response", ) } diff --git a/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/contextvm/transport/CvmTransportTest.kt b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/contextvm/transport/CvmTransportTest.kt index bc0238587e..60bf2981a8 100644 --- a/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/contextvm/transport/CvmTransportTest.kt +++ b/quartz/src/jvmAndroidTest/kotlin/com/vitorpamplona/quartz/contextvm/transport/CvmTransportTest.kt @@ -244,6 +244,8 @@ class CvmTransportTest { async { transport.request(JsonRpcRequest(JsonRpcId.Num(0), "ping"), timeoutMs = 5_000) { seen += it.method + // null: a notification does not answer the call. + null } } yield() @@ -476,6 +478,81 @@ class CvmTransportTest { return ours.map { it.kind }.distinct() } + @Test + fun `CVM-35-12 the session's first message declares this side's surface`() = + runTest { + // CEP-35 is symmetric, and a peer only offers a profile the other + // side declared: the reference ContextVM server chunks a CEP-22 + // response, and opens a CEP-41 stream, only for a client that said + // it can take one. A session that declared nothing capped every + // response at a single relay event. + // + // Read off the inner event, because that is where the server reads + // it: the wrap's own tags are just routing, and a surface put there + // would be both unread and public. + val seen = mutableListOf>() + + // Deliberately NOT the handshake: the surface has to ride whatever + // the first message is, or a client that opens with a tool call + // never declares at all. + declaring(seen, CvmGiftWrap(encryptionMode = EncryptionMode.DISABLED), "msg_fetch_many") + + val declared = seen.single() + assertTrue(CvmTags.SUPPORT_OVERSIZED_TRANSFER in declared, "no CEP-22 flag in $declared") + assertTrue(CvmTags.SUPPORT_OPEN_STREAM in declared, "no CEP-41 flag in $declared") + // Encryption is off here, so there is no wrap we could open. + assertFalse(CvmTags.SUPPORT_ENCRYPTION in declared, "declared a wrap it will not open") + } + + @Test + fun `CVM-35-13 an encrypting client declares both wrap kinds`() = + runTest { + // Receive, not prefer. This client emits 1059, and still says it + // takes 21059, because the transport subscribes to both - saying + // otherwise would tell the peer to withhold something we can read. + val seen = mutableListOf>() + + declaring(seen, CvmGiftWrap(giftWrapMode = GiftWrapMode.PERSISTENT, encryptionMode = EncryptionMode.OPTIONAL), "ping") + + val declared = seen.single() + assertTrue(CvmTags.SUPPORT_ENCRYPTION in declared, "no CEP-4 flag in $declared") + assertTrue(CvmTags.SUPPORT_ENCRYPTION_EPHEMERAL in declared, "no CEP-19 flag in $declared") + } + + @Test + fun `CVM-35-14 the surface is declared once, not on every message`() = + runTest { + val seen = mutableListOf>() + + declaring(seen, CvmGiftWrap(), "ping", "ping") + + assertEquals(2, seen.size, "both requests should have reached the server") + assertTrue(seen.first().isNotEmpty(), "the first message declared nothing") + assertEquals(emptyList(), seen.last(), "re-declared on a later message") + } + + /** + * Runs [methods] in one session, collecting each request's declared + * surface - the single-element tags on the event the server unwrapped. + */ + private suspend fun declaring( + into: MutableList>, + crypto: CvmGiftWrap, + vararg methods: String, + ) { + val fixture = + server(handler = { request -> + into += + request.event.tags + .filter { it.size == 1 } + .map { it[0] } + JsonRpcSuccess(JsonRpcId.Num(0), buildJsonObject { put("ok", JsonPrimitive(true)) }) + }) + fixture.start() + val transport = transport(crypto) + methods.forEach { exchange(fixture, transport, JsonRpcRequest(JsonRpcId.Num(0), it)) } + } + @Test fun `CVM-35-11 the stable identity is used only when asked for`() = runTest {