From d5d86270b965f524f2a5fa7a40c6aec291c16e83 Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 23 Sep 2026 21:53:26 +0000 Subject: [PATCH] fix(contextvm): a subscription ended in an exception, always MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit msg_sub_many was the last of the eleven coordinator tools with no live proof, and pointing a harness at it showed why: it could not succeed. CEP-41 says a stream's `close` does not complete the JSON-RPC request, so an open-ended subscription has no response to wait for — its budget running out is the only way it can ever end. callTool let that timeout propagate, so every subscription, however healthy, finished by throwing. CordnSyncLoop read that as a failure and went into backoff after each one. Under TimeoutMode.TOTAL the budget is now a normal return with whatever the stream delivered; IDLE still throws, because there silence really is a stalled peer. The fixture missed it by answering the subscribe call like any other request, which no coordinator does. CordnGroupManager.subscribe also persists on the way out now, like catchUp and for the same reason: delivering a message advances the ratchet and the cursor, so a subscription that ended without writing left that progress only in memory. A long-lived host survives it; a process-per-command client does not, and neither does a crash. New verb: `amy cordn watch [--timeout MS]`. It calls msg_sub_many and nothing else, so what it prints arrived over an open stream rather than a poll — the only way to exercise the one tool a request/response client never reaches. `fetch` and `watch` now share one delivery collector instead of two copies of the same when-branch. tier-b.sh gains the step, ordered so it proves something: the watcher goes up FIRST and only then does alice send, because a message sent beforehand would be backlog the subscription replays. A second step checks the cursor survived the watching process — removing the persist makes that one fail, so it is not decoration. Verified live: TIER B PASSED, with the message arriving under "via": "msg_sub_many" and not coming back on the next fetch. All eleven tools of spec/00.md are now exercised against the reference coordinator. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_012BfD4txdnsaPRXmNXbup9n --- .../amethyst/cli/commands/CordnCommands.kt | 4 +- .../cli/commands/CordnGroupCommands.kt | 117 +++++++++++++----- cli/tests/cordn/tier-b.sh | 41 +++++- .../commons/cordn/CordnGroupManager.kt | 14 ++- .../quartz/contextvm/mcp/CvmMcpClient.kt | 115 +++++++++-------- .../quartz/contextvm/mcp/CvmMcpClientTest.kt | 40 ++++++ 6 files changed, 246 insertions(+), 85 deletions(-) diff --git a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/CordnCommands.kt b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/CordnCommands.kt index cf4f817e5f..b3492b9156 100644 --- a/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/CordnCommands.kt +++ b/cli/src/main/kotlin/com/vitorpamplona/amethyst/cli/commands/CordnCommands.kt @@ -91,6 +91,7 @@ object CordnCommands { | cordn send --text "…" [--gid GID] a kind-9 chat message | [--reply-to ID] [--react-to ID] | cordn fetch drain the stream and print it + | cordn watch [--timeout MS] hold a live subscription | |A group ref is a locator, not an invitation: holding one lets you ASK to |join, it does not make you a member. Relays say where to reach the @@ -104,7 +105,7 @@ object CordnCommands { route( "cordn", tail, - "cordn ", + "cordn ", help = USAGE, routes = mapOf( @@ -122,6 +123,7 @@ object CordnCommands { "decline" to { rest -> CordnGroupCommands.decline(dataDir, rest) }, "send" to { rest -> CordnGroupCommands.send(dataDir, rest) }, "fetch" to { rest -> CordnGroupCommands.fetch(dataDir, rest) }, + "watch" to { rest -> CordnGroupCommands.watch(dataDir, rest) }, ), ) 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 53e3868b3f..a0dc614122 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 @@ -576,37 +576,7 @@ internal object CordnGroupCommands { val echoes = mutableListOf>() val undecryptable = mutableListOf>() - val drained = - scope.manager.catchUp { delivery -> - when (delivery) { - is CordnGroupManager.Delivery.Message -> - messages += - mapOf( - "gid" to delivery.gid, - "cursor" to delivery.cursor, - "id" to delivery.received.envelope.id, - // What MLS authenticated, not what the - // envelope claims: the envelope is unsigned - // (spec/02.md), so its pubKey field is a - // claim and this is the fact. - "sender" to delivery.received.sender, - "epoch" to delivery.received.epoch, - "kind" to delivery.received.envelope.kind, - "created_at" to delivery.received.envelope.createdAt, - "content" to delivery.received.envelope.content, - ) - - is CordnGroupManager.Delivery.EpochAdvanced -> - epochs += mapOf("gid" to delivery.gid, "cursor" to delivery.cursor, "epoch" to delivery.epoch) - - is CordnGroupManager.Delivery.Echo -> - echoes += mapOf("gid" to delivery.gid, "cursor" to delivery.cursor) - - is CordnGroupManager.Delivery.Undecryptable -> - undecryptable += - mapOf("gid" to delivery.gid, "cursor" to delivery.cursor, "reason" to delivery.reason) - } - } + val drained = scope.manager.catchUp { it.collectInto(messages, epochs, echoes, undecryptable) } Output.emit( mapOf( @@ -626,6 +596,91 @@ internal object CordnGroupCommands { } } + /** + * `cordn watch` - holds a live subscription and prints what arrives. + * + * The sibling of `fetch`, and deliberately not a superset of it: this + * calls `msg_sub_many` and nothing else, so what it prints arrived over + * an open CEP-41 stream rather than a poll. That distinction is the + * point - it is the only way to exercise the one coordinator tool a + * request/response client never reaches. + * + * `--timeout` is a budget, not a failure: the call returns when the + * coordinator closes the stream or the budget runs out, whichever comes + * first, and either way what was delivered has been ingested and saved. + */ + suspend fun watch( + dataDir: DataDir, + tail: Array, + ): Int { + val args = Args(tail) + val timeoutMs = args.longFlag("timeout", DEFAULT_WATCH_MS) + args.rejectUnknown("coordinator", "relay", "timeout") + + if (timeoutMs < 1) return Output.error("bad_args", "--timeout must be positive") + + return CordnRun.withSession(dataDir, args) { _, scope -> + val messages = mutableListOf>() + val epochs = mutableListOf>() + val echoes = mutableListOf>() + val undecryptable = mutableListOf>() + + scope.manager.subscribe(timeoutMs) { delivery -> delivery.collectInto(messages, epochs, echoes, undecryptable) } + + Output.emit( + mapOf( + "coordinator" to scope.config.pubKey, + "watched_ms" to timeoutMs, + "messages" to messages, + "epoch_changes" to epochs, + "echoes" to echoes, + "undecryptable" to undecryptable, + // Everything above came over the subscription, so this is + // a fact about the run rather than a label. + "via" to "msg_sub_many", + ), + ) + 0 + } + } + + /** Files one delivery under the list its kind belongs in. */ + private fun CordnGroupManager.Delivery.collectInto( + messages: MutableList>, + epochs: MutableList>, + echoes: MutableList>, + undecryptable: MutableList>, + ) { + when (this) { + is CordnGroupManager.Delivery.Message -> + messages += + mapOf( + "gid" to gid, + "cursor" to cursor, + "id" to received.envelope.id, + // What MLS authenticated, not what the envelope + // claims: the envelope is unsigned (spec/02.md), so + // its pubKey field is a claim and this is the fact. + "sender" to received.sender, + "epoch" to received.epoch, + "kind" to received.envelope.kind, + "created_at" to received.envelope.createdAt, + "content" to received.envelope.content, + ) + + is CordnGroupManager.Delivery.EpochAdvanced -> + epochs += mapOf("gid" to gid, "cursor" to cursor, "epoch" to epoch) + + is CordnGroupManager.Delivery.Echo -> echoes += mapOf("gid" to gid, "cursor" to cursor) + + is CordnGroupManager.Delivery.Undecryptable -> + undecryptable += mapOf("gid" to gid, "cursor" to cursor, "reason" to reason) + } + } + + /** Long enough for a round trip, short enough to be a foreground command. */ + private const val DEFAULT_WATCH_MS = 15_000L + private fun JoinRequest.toMap() = mapOf( "gid" to gid, diff --git a/cli/tests/cordn/tier-b.sh b/cli/tests/cordn/tier-b.sh index 6c02532b7d..ce85b43bd0 100755 --- a/cli/tests/cordn/tier-b.sh +++ b/cli/tests/cordn/tier-b.sh @@ -3,11 +3,18 @@ # tier-b.sh — the cordn binding, end to end against the REFERENCE coordinator. # # Tier B of quartz/plans/2026-09-17-cordn-interop.md §6.4: a live -# counterparty, not a fixture. Two amy accounts, one local relay (geode), one -# reference coordinator, and the whole lifecycle — publish a KeyPackage, +# counterparty, not a fixture. Three amy accounts, one local relay (geode), +# one reference coordinator, and the whole lifecycle — publish a KeyPackage, # create a group, invite, open the Welcome without joining, join, talk in both # directions, and check both sides agree on epoch and membership. # +# Then the parts that lifecycle never reaches, each driven directly, so that +# all eleven coordinator tools of `spec/00.md` are exercised against a real +# coordinator: a third account who ASKS to join rather than being invited +# (join_request_store, join_request_take_many), a withdrawn KeyPackage +# (kp_remove), a backlog big enough to be chunked (CEP-22), and a message +# pushed down an open stream (msg_sub_many). +# # The reference coordinator it runs against is UNLICENSED — read the header of # stack.sh, which boots it, before running this. Nothing here is wired into a # build, and it must not become so. @@ -195,6 +202,36 @@ 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 "a message delivered over an open stream — msg_sub_many" +# The eleventh tool, and the only one a request/response client never +# reaches: `cordn watch` calls msg_sub_many and nothing else, so anything +# it prints arrived over an open CEP-41 stream rather than a poll. +# +# Ordering is the whole test. The watcher goes up FIRST and we wait for it +# to be listening; only then does alice send. A message sent beforehand +# would be backlog the subscription replays, which proves a stream opened +# but not that anything was pushed down it. +WATCH_OUT="$WORK/watch.json" +bob cordn watch --timeout 25000 --json >"$WATCH_OUT" 2>/dev/null & +WATCH_PID=$! +sleep 8 +LIVE="live-$$-$(date +%s)" +alice cordn send --text "$LIVE" --json >/dev/null +wait "$WATCH_PID" +WATCHED=$(cat "$WATCH_OUT") +echo " $(echo "$WATCHED" | head -c 400)" +[ "$(echo "$WATCHED" | field "['via']")" = "msg_sub_many" ] && ok "the stream was the source" || bad "cordn watch did not report a subscription" +echo "$WATCHED" | grep -q "$LIVE" && ok "pushed live, not polled for" || bad "the subscription never delivered $LIVE" +[ "$(echo "$WATCHED" | field "['messages'][0]['sender']")" = "$ALICE_PK" ] && ok "MLS authenticated the sender over the stream too" || bad "wrong sender on the streamed message" + +step "the watcher saved what it ingested" +# Decrypting a streamed message advances the ratchet and the cursor. A +# subscription that ended without writing would leave that only in memory, +# and the next process would re-read its own progress as a gap. +AFTER=$(bob cordn fetch --json) +echo "$AFTER" | grep -q "$LIVE" && bad "the streamed message came back on the next fetch" || ok "the cursor survived the watching process" +[ "$(echo "$AFTER" | field "['undecryptable']")" = "[]" ] && ok "no gaps after the stream closed" || bad "gaps after the stream: $(echo "$AFTER" | field "['undecryptable']")" + 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/CordnGroupManager.kt b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnGroupManager.kt index 32933a2e97..6b4b90cfde 100644 --- a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnGroupManager.kt +++ b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/cordn/CordnGroupManager.kt @@ -574,13 +574,25 @@ class CordnGroupManager( /** * Subscribes to live delivery for every group. Suspends until the * coordinator closes the stream, so give it its own coroutine. + * + * Persists on the way out, like [catchUp], and for the same reason: + * delivering a message advances the ratchet and the cursor, so a + * subscription that ended without writing would leave that progress only + * in memory. A long-lived host survives that; a process-per-command client + * does not, and neither does a crash. */ override suspend fun subscribe( timeoutMs: Long, onDelivery: (Delivery) -> Unit, ) { if (groups.isEmpty()) return - call { sync.subscribe(groups.keys.toList(), timeoutMs) { gid, ingestion -> onDelivery(ingest(gid, ingestion)) } } + try { + call { sync.subscribe(groups.keys.toList(), timeoutMs) { gid, ingestion -> onDelivery(ingest(gid, ingestion)) } } + } finally { + // finally, not after: a subscription normally ends by timing out or + // being cancelled, and those are the cases with progress to keep. + persistAll() + } } private fun ingest( 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 9e323f3171..7f21cd7f32 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 @@ -40,6 +40,7 @@ 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.coroutines.TimeoutCancellationException import kotlinx.serialization.json.JsonElement import kotlinx.serialization.json.JsonObject import kotlinx.serialization.json.JsonPrimitive @@ -156,71 +157,85 @@ class CvmMcpClient( val streamed = mutableListOf() val response = - transport.request( - message = - JsonRpcRequest( - id = id, - method = McpMethods.TOOLS_CALL, - params = - buildJsonObject { - put(McpParams.NAME, JsonPrimitive(name)) - put(McpParams.ARGUMENTS, arguments) - put( - McpParams.META, - buildJsonObject { - put(McpParams.PROGRESS_TOKEN, JsonPrimitive("call-${id.value}")) - }, - ) - }, - ), - identity = identity, - timeoutMs = timeoutMs, - timeoutMode = timeoutMode, - ) { notification -> - val envelope = ProgressEnvelope.parseOrNull(notification) ?: return@request null - if (envelope.token != token) return@request null + try { + transport.request( + message = + JsonRpcRequest( + id = id, + method = McpMethods.TOOLS_CALL, + params = + buildJsonObject { + put(McpParams.NAME, JsonPrimitive(name)) + put(McpParams.ARGUMENTS, arguments) + put( + McpParams.META, + buildJsonObject { + put(McpParams.PROGRESS_TOKEN, JsonPrimitive("call-${id.value}")) + }, + ) + }, + ), + identity = identity, + timeoutMs = timeoutMs, + timeoutMode = timeoutMode, + ) { notification -> + val envelope = ProgressEnvelope.parseOrNull(notification) ?: return@request null + if (envelope.token != token) return@request null - when (envelope.type) { - ProgressEnvelope.TYPE_OVERSIZED -> { - val receiver = - oversized ?: OversizedTransferReceiver(token, oversizedLimits, requireAccept = false) - .also { oversized = it } - OversizedFrame.parseOrNull(envelope)?.let { frame -> - val result = receiver.accept(frame) - if (result is OversizedProgressResult.Completed) { - reassembled = result.message - oversizedTransfersCompleted++ + when (envelope.type) { + ProgressEnvelope.TYPE_OVERSIZED -> { + val receiver = + oversized ?: OversizedTransferReceiver(token, oversizedLimits, requireAccept = false) + .also { oversized = it } + OversizedFrame.parseOrNull(envelope)?.let { frame -> + val result = receiver.accept(frame) + if (result is OversizedProgressResult.Completed) { + reassembled = result.message + oversizedTransfersCompleted++ + } } } - } - ProgressEnvelope.TYPE_OPEN_STREAM -> { - val receiver = - stream ?: OpenStreamReceiver(token, streamPolicy, requireAccept = false) - .also { stream = it } - OpenStreamFrame.parseOrNull(envelope)?.let { frame -> - val event = receiver.accept(frame) - if (event is OpenStreamEvent.Delivered) { - streamed += event.fragments - event.fragments.forEach(onStreamFragment) + ProgressEnvelope.TYPE_OPEN_STREAM -> { + val receiver = + stream ?: OpenStreamReceiver(token, streamPolicy, requireAccept = false) + .also { stream = it } + OpenStreamFrame.parseOrNull(envelope)?.let { frame -> + val event = receiver.accept(frame) + if (event is OpenStreamEvent.Delivered) { + streamed += event.fragments + event.fragments.forEach(onStreamFragment) + } } } + + else -> Unit } - 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 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 + } catch (e: TimeoutCancellationException) { + // Under TOTAL the timeout IS how the call ends. CEP-41 says a + // stream's `close` does not complete the JSON-RPC request, so + // an open-ended subscription has no other exit - treating the + // budget running out as a failure meant every subscription, + // however healthy, ended in an exception. + // + // Fragments were handed to onStreamFragment as they arrived, + // so nothing is lost by returning here. + if (timeoutMode == TimeoutMode.TOTAL) null else throw e } // 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) { + // No response and no transfer: a subscription that ran its budget. + val effective = reassembled ?: response ?: return ToolCallResult(null, streamed = streamed) + return when (effective) { is JsonRpcSuccess -> ToolCallResult(effective.result, streamed = streamed, viaOversizedTransfer = chunked) is JsonRpcFailure -> ToolCallResult(null, effective.error, streamed, chunked) else -> ToolCallResult(null, streamed = streamed, viaOversizedTransfer = chunked) 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 6b693a8699..1baf6c2801 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 @@ -36,6 +36,7 @@ import com.vitorpamplona.quartz.contextvm.jsonrpc.JsonRpcSuccess 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.crypto.KeyPair import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerInternal import kotlinx.coroutines.async @@ -250,6 +251,45 @@ class CvmMcpClientTest { ) } + @Test + fun `a CEP-41 subscription ends on its budget instead of throwing`() = + runTest { + // The other half of the rule above, and the one that made every + // live subscription an exception: if `close` does not complete the + // request, an open-ended subscription has NO response to wait for, + // so its budget running out is the only way it can end. Under + // TimeoutMode.TOTAL that is a normal return, not a failure. + val fixture = fixture { error("a subscription has no response to give") } + fixture.start() + + val live = mutableListOf() + + val result = + coroutineScope { + val pending = + async { + client().callTool("sub", timeoutMs = 5_000, timeoutMode = TimeoutMode.TOTAL) { live += it } + } + yield() + + listOf( + OpenStreamFrame.start(firstCallToken, 1.0), + OpenStreamFrame.chunk(firstCallToken, 2.0, 0, "pushed"), + OpenStreamFrame.close(firstCallToken, 3.0, lastChunkIndex = 0), + ).forEach { frame -> + fixture.reply(frame.envelope.toNotification(), clientSigner.pubKey, "0".repeat(64)) + } + + // No pump: nothing answers, and the budget is the exit. + pending.await() + } + + assertEquals(listOf("pushed"), live, "fragments arrive as they are pushed") + assertEquals(listOf("pushed"), result.streamed) + assertNull(result.result, "there was no response, and that is not an error") + assertFalse(result.isError) + } + @Test fun `a CEP-41 stream delivers fragments but close does not complete the call`() = runTest {