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 {