fix(contextvm): a subscription ended in an exception, always

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 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_012BfD4txdnsaPRXmNXbup9n
This commit is contained in:
Claude
2026-09-23 21:53:26 +00:00
parent 57ee66c51f
commit d5d86270b9
6 changed files with 246 additions and 85 deletions
@@ -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 <coordinator|keypackage|migrate|group|invite|request|requests|welcomes|join|decline|send|fetch|ref|exposure>",
"cordn <coordinator|keypackage|migrate|group|invite|request|requests|welcomes|join|decline|send|fetch|watch|ref|exposure>",
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) },
),
)
@@ -576,37 +576,7 @@ internal object CordnGroupCommands {
val echoes = mutableListOf<Map<String, Any?>>()
val undecryptable = mutableListOf<Map<String, Any?>>()
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<String>,
): 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<Map<String, Any?>>()
val epochs = mutableListOf<Map<String, Any?>>()
val echoes = mutableListOf<Map<String, Any?>>()
val undecryptable = mutableListOf<Map<String, Any?>>()
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<Map<String, Any?>>,
epochs: MutableList<Map<String, Any?>>,
echoes: MutableList<Map<String, Any?>>,
undecryptable: MutableList<Map<String, Any?>>,
) {
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,
+39 -2
View File
@@ -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)
@@ -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(
@@ -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<String>()
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)
@@ -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<String>()
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 {