diff --git a/contextvm/README.md b/contextvm/README.md new file mode 100644 index 0000000000..cbc7c42a5c --- /dev/null +++ b/contextvm/README.md @@ -0,0 +1,105 @@ +# :contextvm + +A Kotlin Multiplatform client for **ContextVM** — the Model Context Protocol +(MCP) carried over Nostr. + +This is a general MCP-over-Nostr implementation, not a client for any one +server. It exists because Amethyst needs to talk to +[cordn](https://cordn.net) coordinators (see +`quartz/plans/2026-09-17-cordn-interop.md`), but nothing in it is +cordn-specific. + +## Implemented specification revision + +Built from the specification documents in +[`ContextVM/contextvm-docs`](https://github.com/ContextVM/contextvm-docs) at +commit **`e63bce6`**, read on 2026-09-17. + +Most of these are **Draft**, including both transfer profiles, so the text can +still move. When bumping the revision, diff against the `CVM-*` rule ids in the +plan's §6.5 catalog — the test names carry them, so a spec change shows up as +named failures rather than as silent divergence. + +| Spec | Status | Implemented in | +| ---- | ------ | -------------- | +| Core draft spec | Draft | `core/`, `jsonrpc/`, `transport/` | +| CEP-4 Encryption | Final | `crypto/CvmGiftWrap` | +| CEP-6 Public Announcements | Final | `discovery/ServerAnnouncement` | +| CEP-8 Pricing and Payment | Draft | `payment/` | +| CEP-15 Common Tool Schemas | Draft | `schema/CommonToolSchema` | +| CEP-16 Client Pubkey Injection | Final | `fixture/` (server role) | +| CEP-17 Relay List Metadata | Draft | `discovery/ServerRelay` | +| CEP-19 Ephemeral Gift Wraps | Draft | `crypto/CvmGiftWrap` | +| CEP-21 PMI Recommendations | Draft | `payment/Pmi` | +| CEP-22 Oversized Transfer | Draft | `transfer/oversized/` | +| CEP-23 Server Profile Metadata | Draft | `discovery/` (kind 0 via quartz) | +| CEP-24 Server Reviews | Draft | `discovery/ServerReview` | +| CEP-35 Stateless Discovery | Draft | `discovery/SessionDiscovery` | +| CEP-41 Open Streams | Draft | `transfer/stream/` | + +RFC 8785 (JCS), required by CEP-8 and CEP-15, lives in +`quartz/…/utils/jcs/JsonCanonicalization.kt` — it is a generic primitive, not a +ContextVM concern. + +## Licensing + +Implemented **clean-room from the specification**. The reference +[`ContextVM/sdk`](https://github.com/ContextVM/sdk) is **LGPL-3.0**; Amethyst +ships under MIT, so translating that source would carry copyleft terms into +Quartz. Do not read it while working here — the plan's §6 was written so you do +not have to. + +The spec repository carries no LICENSE file. Implementing a published protocol +is fine; do not paste spec prose into this repo. + +## Three things that are easy to get wrong + +1. **Kind 25910 is ephemeral, so relays do not store it.** A subscription + opened after the peer published has missed the response permanently. This is + why `CvmTransport` exposes `request()` and no public publish/subscribe pair: + the ordering is the transport's job, not the caller's. `InMemoryRelayPool` + drops an event nobody is listening for, so the property is tested rather + than assumed. + +2. **CEP-41 has two independent ordering fields.** `progress` orders every + frame, control frames included, and is explicitly *not* a chunk counter; + `chunkIndex` (contiguous from 0) is what validates completeness. Progress + sequences are also per-sender, so a `pong`'s progress bears no relation to + the `ping` it answers — they match by nonce alone. + +3. **`close` does not complete the request.** After a CEP-41 stream closes, the + originating JSON-RPC request still needs its own response, and a client must + never synthesize success from `close`. `ToolCallResult` carries `streamed` + fragments and the `result` separately for exactly this reason. + +## Testing + +```bash +./gradlew :contextvm:jvmTest # 172 tests +./gradlew :contextvm:testAndroidHostTest # 143 tests +``` + +The Android run is smaller because the tests needing real secp256k1 and NIP-44 +— the gift wrap round trip, the transport and the MCP client — live in +`src/jvmTest` where the JVM JNI artifact is on the classpath. Everything +protocol-level is in `commonTest` and runs on both. + +The suite is **Tier A and Tier C** from the plan's §6.4: rule-derived unit +tests, plus adversarial tests against `fixture/CvmFixtureServer`, which plays +the server role and can be told to violate any rule on demand (`FixtureFaults`). +Most tests are negative, because the CEPs are written as failure conditions. + +Still open: + +- **Tier B** — live integration against `ghcr.io/cordn-msg/cordn:latest`. +- **Tier D** — cross-implementation vector exchange for CEP-4 wraps and + CEP-22/41 framing. Worth offering upstream; no official vectors exist. +- **Tier E** — a real Lightning wallet behind CEP-8 (NIP-47 NWC is already in + Quartz). + +## Not in scope + +- A production server. `fixture/` plays the server role for tests only and is + deliberately not hardened; for a real coordinator, `cordn-rs` exists. +- Full MCP. Lifecycle, tool listing and tool calling are implemented because + that is what the CEPs define; sampling, roots and elicitation are not. diff --git a/contextvm/src/commonMain/kotlin/com/vitorpamplona/contextvm/fixture/CvmFixtureServer.kt b/contextvm/src/commonMain/kotlin/com/vitorpamplona/contextvm/fixture/CvmFixtureServer.kt new file mode 100644 index 0000000000..c1a212e293 --- /dev/null +++ b/contextvm/src/commonMain/kotlin/com/vitorpamplona/contextvm/fixture/CvmFixtureServer.kt @@ -0,0 +1,284 @@ +/* + * 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.fixture + +import com.vitorpamplona.contextvm.core.CvmKinds +import com.vitorpamplona.contextvm.core.CvmMessageEvent +import com.vitorpamplona.contextvm.crypto.CvmGiftWrap +import com.vitorpamplona.contextvm.jsonrpc.JsonRpcMessage +import com.vitorpamplona.contextvm.jsonrpc.JsonRpcNotification +import com.vitorpamplona.contextvm.jsonrpc.JsonRpcRequest +import com.vitorpamplona.contextvm.mcp.McpParams +import com.vitorpamplona.contextvm.transport.CvmRelayPool +import com.vitorpamplona.contextvm.transport.CvmSubscription +import com.vitorpamplona.quartz.nip01Core.core.Event +import com.vitorpamplona.quartz.nip01Core.core.HexKey +import com.vitorpamplona.quartz.nip01Core.core.Kind +import com.vitorpamplona.quartz.nip01Core.core.Tag +import com.vitorpamplona.quartz.nip01Core.signers.NostrSigner +import kotlinx.serialization.json.JsonObject +import kotlinx.serialization.json.JsonPrimitive +import kotlinx.serialization.json.buildJsonObject + +/** + * An in-memory relay that delivers to whoever is subscribed at publish time. + * + * Its most important property is a faithful one: an event published while + * nobody is subscribed is **dropped**, exactly as a real relay treats an + * ephemeral kind. A test double that queued it would hide the subscribe-before- + * publish bug this whole layer is designed to prevent. + */ +class InMemoryRelayPool : CvmRelayPool { + private data class Listener( + val pubKey: HexKey, + val kinds: Set, + val onEvent: (Event) -> Unit, + ) + + private val listeners = mutableListOf() + + /** Every event published, in order. For assertions about what went on the wire. */ + val published = mutableListOf() + + /** Events dropped because nothing was listening for them. */ + val dropped = mutableListOf() + + override fun subscribe( + pubKey: HexKey, + kinds: IntArray, + onEvent: (Event) -> Unit, + ): CvmSubscription { + val listener = Listener(pubKey, kinds.toSet(), onEvent) + listeners += listener + return object : CvmSubscription { + override fun close() { + listeners -= listener + } + } + } + + override suspend fun publish(event: Event) { + published += event + + val recipients = event.tags.filter { it.size >= 2 && it[0] == "p" }.map { it[1] } + val matched = + listeners.filter { listener -> + listener.kinds.contains(event.kind) && recipients.contains(listener.pubKey) + } + + if (matched.isEmpty()) dropped += event + matched.forEach { it.onEvent(event) } + } +} + +/** + * How the fixture should misbehave. + * + * This is the reason the fixture exists. No real server sends a duplicate + * `progress`, a stale `pong` nonce or a digest that does not match, yet a client + * MUST handle all of them correctly — several are outright MUST-fail rules. The + * only way to test that is a counterparty that can be told to break them. + */ +data class FixtureFaults( + /** Answer with a JSON-RPC id that does not match the request. */ + val mismatchedResponseId: Boolean = false, + /** Omit the `e` tag that correlates the response to the request event. */ + val omitCorrelationTag: Boolean = false, + /** Answer a different request event id entirely. */ + val wrongCorrelationTag: Boolean = false, + /** Send a response body that is not valid JSON-RPC. */ + val malformedResponse: Boolean = false, + /** Never answer at all. */ + val silent: Boolean = false, + /** Send this many junk notifications before the real response. */ + val noisePrefix: Int = 0, +) + +/** + * A ContextVM server that plays the peer role in tests (Tier C). + * + * Not hardened for deployment and deliberately so: for a real coordinator, + * `cordn-rs` already exists. This exists to be wrong on demand. + */ +class CvmFixtureServer( + private val relays: InMemoryRelayPool, + private val signer: NostrSigner, + private val crypto: CvmGiftWrap = CvmGiftWrap(), + private val faults: FixtureFaults = FixtureFaults(), + /** CEP-16: inject the caller's pubkey into `_meta` before handling. */ + private val injectClientPubkey: Boolean = false, + /** Discovery tags sent on the first direct message back, per CEP-35. */ + private val discoveryTags: List = emptyList(), + /** Answers a request, given its params (with `_meta.clientPubkey` if injected). */ + private val handler: suspend (method: String, params: JsonObject?) -> JsonRpcMessage, +) { + private var subscription: CvmSubscription? = null + private var sentFirstMessage = false + + /** Params as the handler saw them, for asserting CEP-16 injection. */ + val handledParams = mutableListOf() + + /** Starts listening. Call before the client publishes anything. */ + fun start(): CvmSubscription { + val sub = + relays.subscribe( + pubKey = signer.pubKey, + kinds = intArrayOf(CvmKinds.MESSAGE) + CvmKinds.GIFT_WRAPS, + onEvent = { event -> pending += event }, + ) + subscription = sub + return sub + } + + /** Events received but not yet answered. Drained by [pump]. */ + private val pending = mutableListOf() + + /** + * Handles everything received so far. + * + * Explicit rather than automatic because the relay callback cannot suspend, + * and because a test usually wants to control when the answer appears. + */ + suspend fun pump() { + val batch = pending.toList() + pending.clear() + batch.forEach { handle(it) } + } + + private suspend fun handle(event: Event) { + val plain = + if (CvmKinds.isGiftWrap(event.kind)) { + try { + crypto.unwrap(event, signer) + } catch (e: IllegalStateException) { + return + } + } else { + event + } + + val message = CvmMessageEvent.fromOrNull(plain) ?: return + val request = message.message() as? JsonRpcRequest ?: return + + if (faults.silent) return + + val clientPubKey = plain.pubKey + val params = if (injectClientPubkey) inject(request.params, clientPubKey) else request.params + handledParams += params + + repeat(faults.noisePrefix) { index -> + reply( + JsonRpcNotification("notifications/message", buildJsonObject { put("seq", JsonPrimitive(index)) }), + clientPubKey, + plain.id, + ) + } + + val response = handler(request.method, params) + reply(response, clientPubKey, plain.id) + } + + /** Sends [message] back to [clientPubKey], answering [requestEventId]. */ + suspend fun reply( + message: JsonRpcMessage, + clientPubKey: HexKey, + requestEventId: HexKey, + ) { + val correlation = + when { + faults.omitCorrelationTag -> null + faults.wrongCorrelationTag -> "f".repeat(64) + else -> requestEventId + } + + val tags = if (sentFirstMessage) emptyList() else discoveryTags + sentFirstMessage = true + + val content = + if (faults.malformedResponse) { + MALFORMED + } else { + com.vitorpamplona.contextvm.jsonrpc.JsonRpcCodec + .encode(faultInjected(message)) + } + + val inner = + signer.sign( + createdAt = + com.vitorpamplona.quartz.utils.TimeUtils + .now(), + kind = CvmKinds.MESSAGE, + tags = + buildList { + add(arrayOf("p", clientPubKey)) + correlation?.let { add(arrayOf("e", it)) } + addAll(tags) + }.toTypedArray(), + content = content, + ) + + relays.publish( + if (crypto.shouldEncrypt(peerSupportsEncryption = true)) { + crypto.wrap(inner, clientPubKey) + } else { + inner + }, + ) + } + + private fun faultInjected(message: JsonRpcMessage): JsonRpcMessage = + if (faults.mismatchedResponseId && message is com.vitorpamplona.contextvm.jsonrpc.JsonRpcSuccess) { + message.copy( + id = + com.vitorpamplona.contextvm.jsonrpc.JsonRpcId + .Num(999_999), + ) + } else { + message + } + + private fun inject( + params: JsonObject?, + clientPubKey: HexKey, + ): JsonObject = + buildJsonObject { + params?.forEach { (key, value) -> if (key != McpParams.META) put(key, value) } + put( + McpParams.META, + buildJsonObject { + (params?.get(McpParams.META) as? JsonObject)?.forEach { (key, value) -> put(key, value) } + // The client never supplies this: a client-supplied value + // would be a spoof, which is exactly why CEP-16 has the + // server derive it from the event signature. + put(McpParams.CLIENT_PUBKEY, JsonPrimitive(clientPubKey)) + }, + ) + } + + fun stop() { + subscription?.close() + subscription = null + } + + private companion object { + const val MALFORMED = """{"not":"json-rpc"}""" + } +} diff --git a/contextvm/src/commonMain/kotlin/com/vitorpamplona/contextvm/mcp/CvmMcpClient.kt b/contextvm/src/commonMain/kotlin/com/vitorpamplona/contextvm/mcp/CvmMcpClient.kt new file mode 100644 index 0000000000..0b5cfe8dc9 --- /dev/null +++ b/contextvm/src/commonMain/kotlin/com/vitorpamplona/contextvm/mcp/CvmMcpClient.kt @@ -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.mcp + +import com.vitorpamplona.contextvm.jsonrpc.JsonRpcFailure +import com.vitorpamplona.contextvm.jsonrpc.JsonRpcId +import com.vitorpamplona.contextvm.jsonrpc.JsonRpcMessage +import com.vitorpamplona.contextvm.jsonrpc.JsonRpcNotification +import com.vitorpamplona.contextvm.jsonrpc.JsonRpcRequest +import com.vitorpamplona.contextvm.jsonrpc.JsonRpcSuccess +import com.vitorpamplona.contextvm.transfer.ProgressEnvelope +import com.vitorpamplona.contextvm.transfer.ProgressToken +import com.vitorpamplona.contextvm.transfer.oversized.OversizedFrame +import com.vitorpamplona.contextvm.transfer.oversized.OversizedLimits +import com.vitorpamplona.contextvm.transfer.oversized.OversizedProgressResult +import com.vitorpamplona.contextvm.transfer.oversized.OversizedTransferReceiver +import com.vitorpamplona.contextvm.transfer.stream.OpenStreamEvent +import com.vitorpamplona.contextvm.transfer.stream.OpenStreamFrame +import com.vitorpamplona.contextvm.transfer.stream.OpenStreamPolicy +import com.vitorpamplona.contextvm.transfer.stream.OpenStreamReceiver +import com.vitorpamplona.contextvm.transport.CvmTransport +import com.vitorpamplona.contextvm.transport.DualSigner +import com.vitorpamplona.quartz.nip01Core.core.Tag +import kotlinx.serialization.json.JsonElement +import kotlinx.serialization.json.JsonObject +import kotlinx.serialization.json.JsonPrimitive +import kotlinx.serialization.json.buildJsonObject + +/** A tool call's outcome, with anything a transfer profile delivered alongside it. */ +data class ToolCallResult( + val result: JsonElement?, + val error: com.vitorpamplona.contextvm.jsonrpc.JsonRpcError? = null, + /** Fragments a CEP-41 stream delivered while the call was in flight. */ + val streamed: List = emptyList(), +) { + val isError get() = error != null +} + +/** + * A minimal MCP client over ContextVM. + * + * Implements the surface the CEPs define — lifecycle, tool listing and calling — + * rather than all of MCP. Sampling, roots and elicitation are deliberately out + * of scope; nothing in ContextVM or its CEPs needs them, and a partial version + * would be worse than none. + * + * Transfer profiles are handled here because they are request-scoped: a call + * carries a `progressToken`, and CEP-22 or CEP-41 frames for that token arrive + * as notifications while the call is open. + */ +class CvmMcpClient( + private val transport: CvmTransport, + private val clientName: String = "amethyst-contextvm", + private val clientVersion: String = "0.1.0", + private val oversizedLimits: OversizedLimits = OversizedLimits(), + private val streamPolicy: OpenStreamPolicy = OpenStreamPolicy(), +) { + private var nextId = 0L + + /** + * Performs the MCP handshake. + * + * Optional per the spec — servers may operate statelessly — but it is the + * natural place to exchange CEP-35 discovery tags, so a client that can + * afford the round trip should do it. + */ + suspend fun initialize( + capabilityTags: List = emptyList(), + protocolVersion: String = PROTOCOL_VERSION, + ): JsonRpcMessage { + val response = + transport.request( + message = + JsonRpcRequest( + id = nextId(), + method = McpMethods.INITIALIZE, + params = + buildJsonObject { + put("protocolVersion", JsonPrimitive(protocolVersion)) + put("capabilities", buildJsonObject {}) + put( + "clientInfo", + buildJsonObject { + put("name", JsonPrimitive(clientName)) + put("version", JsonPrimitive(clientVersion)) + }, + ) + }, + ), + discoveryTags = capabilityTags, + ) + + // The server is only allowed to assume readiness after this. + transport.notify(JsonRpcNotification(McpMethods.INITIALIZED)) + return response + } + + suspend fun listTools(cursor: String? = null): JsonRpcMessage = + transport.request( + JsonRpcRequest( + id = nextId(), + method = McpMethods.TOOLS_LIST, + params = cursor?.let { buildJsonObject { put("cursor", JsonPrimitive(it)) } }, + ), + ) + + /** + * Calls a tool. + * + * A `progressToken` is always attached: without it a server MUST NOT start + * either transfer profile, so omitting it would silently cap every response + * at one relay event. + */ + suspend fun callTool( + name: String, + arguments: JsonObject = buildJsonObject {}, + identity: DualSigner.Identity = DualSigner.Identity.EPHEMERAL, + timeoutMs: Long = CvmTransport.DEFAULT_TIMEOUT_MS, + onStreamFragment: (String) -> Unit = {}, + ): ToolCallResult { + val id = nextId() + val token = ProgressToken.Text("call-${(id as JsonRpcId.Num).value}") + + var oversized: OversizedTransferReceiver? = null + var stream: OpenStreamReceiver? = null + var reassembled: JsonRpcMessage? = null + 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, + ) { notification -> + val envelope = ProgressEnvelope.parseOrNull(notification) ?: return@request + if (envelope.token != token) return@request + + 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 + } + } + + 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 + } + } + + // 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. + 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) + } + } + + private fun nextId(): JsonRpcId.Num = JsonRpcId.Num(nextId++) + + companion object { + /** The MCP revision the ContextVM spec's examples use. */ + const val PROTOCOL_VERSION = "2025-07-02" + } +} diff --git a/contextvm/src/commonMain/kotlin/com/vitorpamplona/contextvm/transport/CvmRelayPool.kt b/contextvm/src/commonMain/kotlin/com/vitorpamplona/contextvm/transport/CvmRelayPool.kt new file mode 100644 index 0000000000..4e672d3789 --- /dev/null +++ b/contextvm/src/commonMain/kotlin/com/vitorpamplona/contextvm/transport/CvmRelayPool.kt @@ -0,0 +1,55 @@ +/* + * 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.transport + +import com.vitorpamplona.quartz.nip01Core.core.Event +import com.vitorpamplona.quartz.nip01Core.core.HexKey + +/** A live subscription. Closing it stops delivery. */ +interface CvmSubscription { + fun close() +} + +/** + * The relay surface ContextVM needs. + * + * Deliberately tiny and transport-agnostic: everything above it is pure + * protocol, and the Tier C fixture implements this same interface to play a + * misbehaving peer without any network. The production binding wraps quartz's + * relay client. + */ +interface CvmRelayPool { + /** + * Subscribes to events addressed to [pubKey] (`#p`) of the given [kinds]. + * + * Delivery starts when this returns. Because kind 25910 is ephemeral, a + * subscription opened after a peer published has missed the event + * permanently — [CvmTransport] is built so callers cannot make that mistake. + */ + fun subscribe( + pubKey: HexKey, + kinds: IntArray, + onEvent: (Event) -> Unit, + ): CvmSubscription + + /** Publishes [event] to the configured relays. */ + suspend fun publish(event: Event) +} diff --git a/contextvm/src/commonMain/kotlin/com/vitorpamplona/contextvm/transport/CvmTransport.kt b/contextvm/src/commonMain/kotlin/com/vitorpamplona/contextvm/transport/CvmTransport.kt new file mode 100644 index 0000000000..cad22355d7 --- /dev/null +++ b/contextvm/src/commonMain/kotlin/com/vitorpamplona/contextvm/transport/CvmTransport.kt @@ -0,0 +1,235 @@ +/* + * 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.transport + +import com.vitorpamplona.contextvm.core.CvmKinds +import com.vitorpamplona.contextvm.core.CvmMessageEvent +import com.vitorpamplona.contextvm.crypto.CvmGiftWrap +import com.vitorpamplona.contextvm.discovery.SessionDiscovery +import com.vitorpamplona.contextvm.jsonrpc.JsonRpcFailure +import com.vitorpamplona.contextvm.jsonrpc.JsonRpcId +import com.vitorpamplona.contextvm.jsonrpc.JsonRpcMessage +import com.vitorpamplona.contextvm.jsonrpc.JsonRpcNotification +import com.vitorpamplona.contextvm.jsonrpc.JsonRpcRequest +import com.vitorpamplona.contextvm.jsonrpc.JsonRpcSuccess +import com.vitorpamplona.quartz.nip01Core.core.Event +import com.vitorpamplona.quartz.nip01Core.core.HexKey +import com.vitorpamplona.quartz.nip01Core.core.Tag +import com.vitorpamplona.quartz.nip01Core.signers.NostrSigner +import kotlinx.coroutines.channels.Channel +import kotlinx.coroutines.withTimeout + +/** + * The two identities a ContextVM client uses. + * + * Splitting them is a privacy measure, not plumbing: the stable identity signs + * only what must be attributable, while everything else rides a throwaway key so + * a server cannot link a session's activity to an account. Making the choice an + * explicit parameter means a caller cannot leak the stable one by omission. + */ +class DualSigner( + /** The account identity. Used only where a call must be attributable. */ + val stable: NostrSigner, + /** A per-session throwaway. Used for everything else. */ + val ephemeral: NostrSigner, +) { + enum class Identity { + STABLE, + EPHEMERAL, + } + + fun signerFor(identity: Identity) = + when (identity) { + Identity.STABLE -> stable + Identity.EPHEMERAL -> ephemeral + } +} + +/** Thrown when a request cannot be completed at the transport layer. */ +class CvmTransportException( + message: String, +) : IllegalStateException(message) + +/** + * Correlates ContextVM requests with their responses over a [CvmRelayPool]. + * + * [request] subscribes before it publishes, always. Kind 25910 is ephemeral, so + * a subscription opened afterwards has missed the response permanently and the + * failure looks exactly like a flaky relay. Making the ordering the transport's + * job rather than the caller's removes the whole class of bug, which is why + * there is no public publish/subscribe pair to get wrong. + * + * Inbound notifications that are not responses — CEP-22/41 frames, CEP-8 payment + * notifications — are handed to `onNotification` while the request is still in + * flight. A notification never resolves a request: CEP-41 is explicit that a + * stream's `close` does not complete it. + */ +class CvmTransport( + private val relays: CvmRelayPool, + private val signers: DualSigner, + private val serverPubKey: HexKey, + private val crypto: CvmGiftWrap = CvmGiftWrap(), + private val discovery: SessionDiscovery = SessionDiscovery(), + /** Whether the peer is known to accept encrypted messages. */ + private val peerSupportsEncryption: Boolean = true, + private val peerSupportsEphemeralWrap: Boolean = true, +) { + /** The peer's learned discovery baseline, once its first message has arrived. */ + val peer get() = discovery.peer + + private var sentFirstMessage = false + + /** + * Sends [message] and waits for the correlated response. + * + * @param identity which key signs the request. Anything not required to be + * attributable should stay [DualSigner.Identity.EPHEMERAL]. + * @param discoveryTags this side's CEP-35 baseline, sent on the session's + * first direct message only. + */ + suspend fun request( + message: JsonRpcRequest, + identity: DualSigner.Identity = DualSigner.Identity.EPHEMERAL, + timeoutMs: Long = DEFAULT_TIMEOUT_MS, + discoveryTags: List = emptyList(), + onNotification: (JsonRpcNotification) -> Unit = {}, + ): JsonRpcMessage { + val signer = signers.signerFor(identity) + val inbound = Channel(Channel.UNLIMITED) + + // Subscribe first, before anything is published. + val subscription = + relays.subscribe( + pubKey = signer.pubKey, + // Both wrap kinds plus the bare message: CEP-19's fallback means + // either wrap may arrive, and an unencrypted peer sends 25910. + kinds = intArrayOf(CvmKinds.MESSAGE) + CvmKinds.GIFT_WRAPS, + onEvent = { inbound.trySend(it) }, + ) + + try { + // CEP-35: the baseline rides the first direct message only; after + // that both sides omit repeated common discovery tags. + val tags = if (sentFirstMessage) emptyList() else discoveryTags + val request = CvmMessageEvent.create(message, serverPubKey, signer, extraTags = tags) + sentFirstMessage = true + + relays.publish(outbound(request)) + + return withTimeout(timeoutMs) { + awaitResponse(inbound, signer, request.id, message.id, onNotification) + } + } finally { + subscription.close() + inbound.close() + } + } + + /** Sends a notification. Nothing is awaited, so no subscription is opened. */ + suspend fun notify( + message: JsonRpcNotification, + identity: DualSigner.Identity = DualSigner.Identity.EPHEMERAL, + ) { + val signer = signers.signerFor(identity) + relays.publish(outbound(CvmMessageEvent.create(message, serverPubKey, signer))) + } + + private suspend fun awaitResponse( + inbound: Channel, + signer: NostrSigner, + requestEventId: HexKey, + requestId: JsonRpcId, + onNotification: (JsonRpcNotification) -> Unit, + ): JsonRpcMessage { + for (event in inbound) { + val plain = decryptOrNull(event, signer) ?: continue + val wrapped = CvmMessageEvent.fromOrNull(plain) ?: continue + + discovery.observe(wrapped.discoveryTags().toTypedArray()) + + val decoded = + try { + wrapped.message() + } catch (e: IllegalArgumentException) { + // A malformed payload from the peer is not our request's + // answer; keep waiting rather than failing the call on it. + continue + } + + when (decoded) { + is JsonRpcNotification -> onNotification(decoded) + + // 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 + // alone is weaker -- a peer may omit the tag on a wrap, and ids + // are only unique within a session. + is JsonRpcSuccess -> + if (matches(wrapped, requestEventId, decoded.id, requestId)) return decoded + + is JsonRpcFailure -> + if (decoded.id == null || matches(wrapped, requestEventId, decoded.id, requestId)) return decoded + + else -> Unit + } + } + throw CvmTransportException("subscription closed before a response arrived") + } + + private fun matches( + message: CvmMessageEvent, + requestEventId: HexKey, + responseId: JsonRpcId, + requestId: JsonRpcId, + ): Boolean { + val inReplyTo = message.inReplyTo() + if (inReplyTo != null && inReplyTo != requestEventId) return false + return responseId == requestId + } + + private suspend fun decryptOrNull( + event: Event, + signer: NostrSigner, + ): Event? = + if (CvmKinds.isGiftWrap(event.kind)) { + try { + crypto.unwrap(event, signer) + } catch (e: IllegalStateException) { + // Not addressed to us, or forged. Ignoring is correct: a relay + // may deliver wraps we cannot open, and failing the request on + // one would let anyone disrupt a call. + null + } + } else { + event + } + + private suspend fun outbound(inner: Event): Event = + if (crypto.shouldEncrypt(peerSupportsEncryption)) { + crypto.wrap(inner, serverPubKey, crypto.negotiatedWrapKind(peerSupportsEphemeralWrap)) + } else { + inner + } + + companion object { + /** Covers signing, the relay round trip and the peer's own work. */ + const val DEFAULT_TIMEOUT_MS = 20_000L + } +} diff --git a/contextvm/src/jvmTest/kotlin/com/vitorpamplona/contextvm/mcp/CvmMcpClientTest.kt b/contextvm/src/jvmTest/kotlin/com/vitorpamplona/contextvm/mcp/CvmMcpClientTest.kt new file mode 100644 index 0000000000..981e93d9d0 --- /dev/null +++ b/contextvm/src/jvmTest/kotlin/com/vitorpamplona/contextvm/mcp/CvmMcpClientTest.kt @@ -0,0 +1,266 @@ +/* + * 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 + +import com.vitorpamplona.contextvm.crypto.CvmGiftWrap +import com.vitorpamplona.contextvm.crypto.EncryptionMode +import com.vitorpamplona.contextvm.fixture.CvmFixtureServer +import com.vitorpamplona.contextvm.fixture.InMemoryRelayPool +import com.vitorpamplona.contextvm.jsonrpc.JsonRpcCodec +import com.vitorpamplona.contextvm.jsonrpc.JsonRpcError +import com.vitorpamplona.contextvm.jsonrpc.JsonRpcFailure +import com.vitorpamplona.contextvm.jsonrpc.JsonRpcId +import com.vitorpamplona.contextvm.jsonrpc.JsonRpcSuccess +import com.vitorpamplona.contextvm.transfer.ProgressToken +import com.vitorpamplona.contextvm.transfer.oversized.OversizedTransferSender +import com.vitorpamplona.contextvm.transfer.stream.OpenStreamFrame +import com.vitorpamplona.contextvm.transport.CvmTransport +import com.vitorpamplona.contextvm.transport.DualSigner +import com.vitorpamplona.quartz.nip01Core.crypto.KeyPair +import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerInternal +import kotlinx.coroutines.async +import kotlinx.coroutines.coroutineScope +import kotlinx.coroutines.test.runTest +import kotlinx.coroutines.yield +import kotlinx.serialization.json.JsonObject +import kotlinx.serialization.json.JsonPrimitive +import kotlinx.serialization.json.buildJsonObject +import kotlinx.serialization.json.jsonObject +import kotlinx.serialization.json.jsonPrimitive +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertFalse +import kotlin.test.assertNull +import kotlin.test.assertTrue + +/** The MCP client end to end, including both transfer profiles. */ +class CvmMcpClientTest { + private val relays = InMemoryRelayPool() + private val serverSigner = NostrSignerInternal(KeyPair()) + private val clientSigner = NostrSignerInternal(KeyPair()) + private val plaintext = CvmGiftWrap(encryptionMode = EncryptionMode.DISABLED) + + private fun client() = + CvmMcpClient( + CvmTransport( + relays = relays, + signers = DualSigner(clientSigner, clientSigner), + serverPubKey = serverSigner.pubKey, + crypto = plaintext, + ), + ) + + private fun fixture(handler: suspend (String, JsonObject?) -> com.vitorpamplona.contextvm.jsonrpc.JsonRpcMessage) = + CvmFixtureServer( + relays = relays, + signer = serverSigner, + crypto = plaintext, + handler = handler, + ) + + /** The progressToken the client derives for the first call it makes. */ + private val firstCallToken = ProgressToken.Text("call-0") + + @Test + fun `a tool call returns the server result`() = + runTest { + val fixture = + fixture { _, _ -> + JsonRpcSuccess(JsonRpcId.Num(0), buildJsonObject { put("cursor", JsonPrimitive(7)) }) + } + fixture.start() + + val result = + coroutineScope { + val pending = async { client().callTool("msg_post", timeoutMs = 5_000) } + yield() + fixture.pump() + pending.await() + } + + assertFalse(result.isError) + assertEquals( + 7, + result.result!! + .jsonObject["cursor"]!! + .jsonPrimitive.content + .toInt(), + ) + } + + @Test + fun `a tool call always carries a progressToken`() = + runTest { + // Without one a server MUST NOT start either transfer profile, so + // omitting it would silently cap every response at one relay event. + val fixture = fixture { _, _ -> JsonRpcSuccess(JsonRpcId.Num(0), buildJsonObject {}) } + fixture.start() + + coroutineScope { + val pending = async { client().callTool("x", timeoutMs = 5_000) } + yield() + fixture.pump() + pending.await() + } + + val meta = fixture.handledParams.first()!![McpParams.META]!!.jsonObject + assertEquals("call-0", meta[McpParams.PROGRESS_TOKEN]!!.jsonPrimitive.content) + } + + @Test + fun `an error response surfaces as an error rather than a throw`() = + runTest { + val fixture = + fixture { _, _ -> + JsonRpcFailure( + JsonRpcId.Num(0), + JsonRpcError(JsonRpcError.PAYMENT_REQUIRED, "Payment Required"), + ) + } + fixture.start() + + val result = + coroutineScope { + val pending = async { client().callTool("priced", timeoutMs = 5_000) } + yield() + fixture.pump() + pending.await() + } + + assertTrue(result.isError) + assertEquals(JsonRpcError.PAYMENT_REQUIRED, result.error!!.code) + assertNull(result.result) + } + + @Test + fun `a CEP-22 transfer reassembles into the effective response`() = + runTest { + 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)) }) + } + fixture.start() + + val result = + coroutineScope { + 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() + } + + assertEquals( + "x".repeat(400), + result.result!! + .jsonObject["text"]!! + .jsonPrimitive.content, + "the reassembled payload replaces the placeholder", + ) + } + + @Test + fun `a CEP-41 stream delivers fragments but close does not complete the call`() = + runTest { + // The rule worth pinning end to end: close says no more frames, and + // the request is still only finished by its own JSON-RPC response. + val fixture = + fixture { _, _ -> + JsonRpcSuccess(JsonRpcId.Num(0), buildJsonObject { put("done", JsonPrimitive(true)) }) + } + fixture.start() + + val live = mutableListOf() + + val result = + coroutineScope { + val pending = + async { + client().callTool("stream", timeoutMs = 5_000) { live += it } + } + yield() + + listOf( + OpenStreamFrame.start(firstCallToken, 1.0), + OpenStreamFrame.chunk(firstCallToken, 2.0, 0, "Hello"), + OpenStreamFrame.chunk(firstCallToken, 3.0, 1, " world"), + OpenStreamFrame.close(firstCallToken, 4.0, lastChunkIndex = 1), + ).forEach { frame -> + fixture.reply(frame.envelope.toNotification(), clientSigner.pubKey, "0".repeat(64)) + } + + // Only now does the request's own response arrive. + fixture.pump() + pending.await() + } + + assertEquals(listOf("Hello", " world"), live) + assertEquals(listOf("Hello", " world"), result.streamed) + assertEquals( + true, + result.result!! + .jsonObject["done"]!! + .jsonPrimitive.content + .toBoolean(), + "the call is completed by its JSON-RPC response, not by close", + ) + } + + @Test + fun `initialize completes the handshake and sends initialized`() = + runTest { + val fixture = + fixture { _, _ -> + JsonRpcSuccess( + JsonRpcId.Num(0), + buildJsonObject { put("protocolVersion", JsonPrimitive(CvmMcpClient.PROTOCOL_VERSION)) }, + ) + } + fixture.start() + + coroutineScope { + val pending = async { client().initialize() } + yield() + fixture.pump() + pending.await() + } + + // The notification rides after the response, unsubscribed, so it + // shows up on the wire rather than in a correlation slot. + val methods = + relays.published.mapNotNull { event -> + runCatching { JsonRpcCodec.decode(event.content) }.getOrNull() + } + assertTrue( + methods.any { it is com.vitorpamplona.contextvm.jsonrpc.JsonRpcNotification && it.method == McpMethods.INITIALIZED }, + "the client must tell the server it is ready", + ) + } +} diff --git a/contextvm/src/jvmTest/kotlin/com/vitorpamplona/contextvm/transport/CvmTransportTest.kt b/contextvm/src/jvmTest/kotlin/com/vitorpamplona/contextvm/transport/CvmTransportTest.kt new file mode 100644 index 0000000000..628d8a968c --- /dev/null +++ b/contextvm/src/jvmTest/kotlin/com/vitorpamplona/contextvm/transport/CvmTransportTest.kt @@ -0,0 +1,326 @@ +/* + * 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.transport + +import com.vitorpamplona.contextvm.core.CvmKinds +import com.vitorpamplona.contextvm.core.CvmTags +import com.vitorpamplona.contextvm.crypto.CvmGiftWrap +import com.vitorpamplona.contextvm.crypto.EncryptionMode +import com.vitorpamplona.contextvm.fixture.CvmFixtureServer +import com.vitorpamplona.contextvm.fixture.FixtureFaults +import com.vitorpamplona.contextvm.fixture.InMemoryRelayPool +import com.vitorpamplona.contextvm.jsonrpc.JsonRpcId +import com.vitorpamplona.contextvm.jsonrpc.JsonRpcRequest +import com.vitorpamplona.contextvm.jsonrpc.JsonRpcSuccess +import com.vitorpamplona.contextvm.mcp.McpParams +import com.vitorpamplona.quartz.nip01Core.crypto.KeyPair +import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerInternal +import kotlinx.coroutines.TimeoutCancellationException +import kotlinx.coroutines.async +import kotlinx.coroutines.coroutineScope +import kotlinx.coroutines.test.runTest +import kotlinx.coroutines.yield +import kotlinx.serialization.json.JsonObject +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.assertNotEquals +import kotlin.test.assertTrue + +/** + * `CVM-CORE-*` and `CVM-16-*` driven end to end against the Tier C fixture. + * + * The fixture is what makes the negative half testable: no real server produces + * a mismatched correlation tag or a malformed body on request. + */ +class CvmTransportTest { + private val relays = InMemoryRelayPool() + private val serverSigner = NostrSignerInternal(KeyPair()) + private val stableSigner = NostrSignerInternal(KeyPair()) + private val ephemeralSigner = NostrSignerInternal(KeyPair()) + private val signers = DualSigner(stableSigner, ephemeralSigner) + + private fun server( + faults: FixtureFaults = FixtureFaults(), + injectClientPubkey: Boolean = false, + discoveryTags: List> = emptyList(), + handler: suspend (String, JsonObject?) -> com.vitorpamplona.contextvm.jsonrpc.JsonRpcMessage = { _, _ -> + JsonRpcSuccess(JsonRpcId.Num(0), buildJsonObject { put("ok", JsonPrimitive(true)) }) + }, + ) = CvmFixtureServer( + relays = relays, + signer = serverSigner, + faults = faults, + injectClientPubkey = injectClientPubkey, + discoveryTags = discoveryTags, + handler = handler, + ) + + private fun transport(crypto: CvmGiftWrap = CvmGiftWrap()) = + CvmTransport( + relays = relays, + signers = signers, + serverPubKey = serverSigner.pubKey, + crypto = crypto, + ) + + /** Runs a request while pumping the fixture, so the answer arrives in-flight. */ + private suspend fun exchange( + fixture: CvmFixtureServer, + transport: CvmTransport, + request: JsonRpcRequest, + identity: DualSigner.Identity = DualSigner.Identity.EPHEMERAL, + timeoutMs: Long = 5_000, + ) = coroutineScope { + val pending = async { transport.request(request, identity, timeoutMs) } + yield() + fixture.pump() + pending.await() + } + + @Test + fun `CVM-CORE-10 completes a request against the fixture`() = + runTest { + val fixture = server() + fixture.start() + + val response = + exchange(fixture, transport(), JsonRpcRequest(JsonRpcId.Num(0), "tools/list")) + + assertIs(response) + } + + @Test + fun `CVM-CORE-11 subscribes before publishing so nothing is dropped`() = + runTest { + // The property that matters: an ephemeral kind published with nobody + // listening is gone. If the transport ever published first, the + // request event would land in `dropped`. + val fixture = server() + fixture.start() + + exchange(fixture, transport(), JsonRpcRequest(JsonRpcId.Num(0), "ping")) + + assertTrue(relays.dropped.isEmpty(), "dropped: ${relays.dropped.map { it.kind }}") + } + + @Test + fun `CVM-CORE-12 the relay drops an event nobody is subscribed to`() = + runTest { + // Guards the guard: confirms the fixture relay really does model + // ephemeral delivery, so the previous test is meaningful. + val fixture = server() + // deliberately not started + val transport = transport() + + assertFailsWith { + transport.request(JsonRpcRequest(JsonRpcId.Num(0), "ping"), timeoutMs = 50) + } + assertTrue(relays.dropped.isNotEmpty()) + } + + @Test + fun `CVM-CORE-13 ignores a response whose JSON-RPC id does not match`() = + runTest { + val fixture = server(faults = FixtureFaults(mismatchedResponseId = true)) + fixture.start() + + assertFailsWith { + exchange(fixture, transport(), JsonRpcRequest(JsonRpcId.Num(7), "ping"), timeoutMs = 200) + } + } + + @Test + fun `CVM-CORE-14 ignores a response correlated to a different request event`() = + runTest { + val fixture = server(faults = FixtureFaults(wrongCorrelationTag = true)) + fixture.start() + + assertFailsWith { + exchange(fixture, transport(), JsonRpcRequest(JsonRpcId.Num(0), "ping"), timeoutMs = 200) + } + } + + @Test + fun `CVM-CORE-15 accepts a response that omits the e tag but matches by id`() = + runTest { + // The `e` tag is the stronger signal but a peer may omit it; the + // JSON-RPC id still correlates. + val fixture = server(faults = FixtureFaults(omitCorrelationTag = true)) + fixture.start() + + val response = exchange(fixture, transport(), JsonRpcRequest(JsonRpcId.Num(0), "ping")) + assertIs(response) + } + + @Test + fun `CVM-CORE-16 keeps waiting through a malformed payload`() = + runTest { + // A peer sending garbage must not fail an unrelated in-flight call. + val fixture = server(faults = FixtureFaults(malformedResponse = true)) + fixture.start() + + assertFailsWith { + exchange(fixture, transport(), JsonRpcRequest(JsonRpcId.Num(0), "ping"), timeoutMs = 200) + } + } + + @Test + fun `CVM-CORE-17 routes notifications without resolving the request`() = + runTest { + val fixture = server(faults = FixtureFaults(noisePrefix = 3)) + fixture.start() + + val seen = mutableListOf() + val transport = transport() + + val response = + coroutineScope { + val pending = + async { + transport.request(JsonRpcRequest(JsonRpcId.Num(0), "ping"), timeoutMs = 5_000) { + seen += it.method + } + } + yield() + fixture.pump() + pending.await() + } + + assertEquals(3, seen.size, "every notification is delivered") + assertIs(response, "and none of them completed the request") + } + + @Test + fun `CVM-4-20 the request goes out encrypted by default`() = + runTest { + val fixture = server() + fixture.start() + + exchange(fixture, transport(), JsonRpcRequest(JsonRpcId.Num(0), "ping")) + + val first = relays.published.first() + assertTrue(CvmKinds.isGiftWrap(first.kind), "kind ${first.kind} is not a wrap") + assertNotEquals(ephemeralSigner.pubKey, first.pubKey, "the wrap hides the sender") + } + + @Test + fun `CVM-4-21 a disabled-encryption client publishes the bare message kind`() = + runTest { + val fixture = + CvmFixtureServer( + relays = relays, + signer = serverSigner, + crypto = CvmGiftWrap(encryptionMode = EncryptionMode.DISABLED), + handler = { _, _ -> JsonRpcSuccess(JsonRpcId.Num(0), buildJsonObject {}) }, + ) + fixture.start() + + exchange( + fixture, + transport(CvmGiftWrap(encryptionMode = EncryptionMode.DISABLED)), + JsonRpcRequest(JsonRpcId.Num(0), "ping"), + ) + + assertEquals(CvmKinds.MESSAGE, relays.published.first().kind) + } + + @Test + fun `CVM-16-01 the server derives clientPubkey rather than trusting the client`() = + runTest { + val fixture = server(injectClientPubkey = true) + fixture.start() + + exchange( + fixture, + transport(CvmGiftWrap(encryptionMode = EncryptionMode.DISABLED)), + JsonRpcRequest(JsonRpcId.Num(0), "tools/call", buildJsonObject { put("name", JsonPrimitive("x")) }), + ) + + val meta = fixture.handledParams.first()?.get(McpParams.META) as JsonObject + assertEquals( + ephemeralSigner.pubKey, + (meta[McpParams.CLIENT_PUBKEY] as JsonPrimitive).content, + "the injected identity is the event signer, not anything the client claimed", + ) + } + + @Test + fun `CVM-16-02 injection is off unless the server opts in`() = + runTest { + val fixture = server(injectClientPubkey = false) + fixture.start() + + exchange( + fixture, + transport(CvmGiftWrap(encryptionMode = EncryptionMode.DISABLED)), + JsonRpcRequest(JsonRpcId.Num(0), "ping"), + ) + + assertTrue(fixture.handledParams.first()?.containsKey(McpParams.META) != true) + } + + @Test + fun `CVM-35-10 learns the server baseline from its first direct message`() = + runTest { + val fixture = + server( + discoveryTags = + listOf( + arrayOf("name", "Fixture"), + CvmTags.flag(CvmTags.SUPPORT_OPEN_STREAM), + arrayOf("unknown_future", "keep-me"), + ), + ) + fixture.start() + + val transport = transport() + exchange(fixture, transport, JsonRpcRequest(JsonRpcId.Num(0), "ping")) + + val peer = transport.peer!! + assertEquals("Fixture", peer.name) + assertTrue(peer.supportsOpenStream) + assertTrue( + peer.unknownTags.any { it[0] == "unknown_future" }, + "CEP-35 requires unknown discovery tags to survive", + ) + } + + @Test + fun `CVM-35-11 the stable identity is used only when asked for`() = + runTest { + val fixture = server(injectClientPubkey = true) + fixture.start() + + exchange( + fixture, + transport(CvmGiftWrap(encryptionMode = EncryptionMode.DISABLED)), + JsonRpcRequest(JsonRpcId.Num(0), "kp_publish"), + identity = DualSigner.Identity.STABLE, + ) + + val meta = fixture.handledParams.first()?.get(McpParams.META) as JsonObject + assertEquals(stableSigner.pubKey, (meta[McpParams.CLIENT_PUBKEY] as JsonPrimitive).content) + } +} diff --git a/quartz/plans/2026-09-17-cordn-interop.md b/quartz/plans/2026-09-17-cordn-interop.md index aecf4c268d..0d04427c60 100644 --- a/quartz/plans/2026-09-17-cordn-interop.md +++ b/quartz/plans/2026-09-17-cordn-interop.md @@ -1,7 +1,13 @@ # Cordn interop: extract the MLS core, then add a second binding -Status: Queued. Research complete; no code written. Blocked on one upstream protocol decision -(§4.1) before any of Stage 2+ is worth starting. +Status: Stage 2 landed. `:contextvm` implements the core spec plus all 12 CEPs on the client +side, with the Tier C fixture server, at 172 tests. Stages 0 (cordn-side vectors), 1 (extract the +MLS engine) and 3-4 (the cordn binding and app integration) are open. Stage 3 is the one gated on +the §4.1 upstream decision. + +Correction to an earlier gate in this plan: §4.1 does **not** block Stage 2. ContextVM is +credential-agnostic and has no MLS dependency at all, so the transport was safe to build first; +only Stage 3's KeyPackage work depends on that decision. Sources checked on 2026-09-17: @@ -16,9 +22,10 @@ Sources checked on 2026-09-17: - `ContextVM/sdk` @ `b5d1e4e` (2026-09-17), version `0.13.17` — **LGPL-3.0**, see §7. Read only to confirm deployed defaults, never as an implementation source -Not verified by execution: this container could not run `:quartz:jvmTest` (no Gradle dependency -cache and Maven Central returns HTTP 429 through the agent proxy), so every claim below about our -own code comes from reading it, not from a green test run. Stage 0 exists to fix that. +Verification status: `:quartz:jvmTest` now passes in a container (5054 tests), so the MLS +interop claims in §3 are execution-verified rather than read-verified. `:contextvm:jvmTest` +passes at 172. `testAndroidHostTest` remains unrun — Maven Central rate-limits the Android +secp256k1 artifact through the agent proxy. ## 1. Executive summary @@ -47,16 +54,17 @@ we already have. Only **CEP-4, CEP-6 and CEP-16** are Final; the core spec and t are Draft, including CEP-22 and CEP-41 (§6.7). Because the CEPs are symmetric, compliance is not demonstrable against cordn alone. §6.4 defines -five test tiers, and the one that does not exist yet is **Tier C: a Kotlin fixture server that -misbehaves on demand.** No real server sends a non-monotonic `progress`, a stale `pong` nonce or a -mismatched digest, yet those are MUST-fail requirements — so the fixture is a first-class Stage 2 -deliverable, not scaffolding. Two CEPs (8 and 15) also need **RFC 8785 JCS**, which Quartz does not -have; it lands in `quartz/…/utils/` as a shared primitive. +five test tiers. **Tier C — a Kotlin fixture server that misbehaves on demand — is built** +(`contextvm/…/fixture/`) and is what makes the negative half testable: no real server sends a +non-monotonic `progress`, a stale `pong` nonce or a mismatched digest, yet those are MUST-fail +requirements. **RFC 8785 JCS** is built too, in `quartz/…/utils/jcs/`, shared by CEP-8 and CEP-15. +Tiers B (live coordinator), D (cross-implementation vectors) and E (a real wallet) remain open. -Recommended sequencing: **Stage 0 (vectors) → Stage 1 (extract engine) → decide → Stage 2+.** -Do not start Stage 2 before the §4.1 decision, because if it goes the wrong way every KeyPackage -is permanently ecosystem-bound and "interop" degrades to Amethyst speaking two unrelated -protocols. +Sequencing, as revised in practice: **Stage 2 (the transport) was built first**, because +ContextVM has no MLS dependency and so no dependency on the §4.1 decision. What that decision +gates is **Stage 3**, the cordn binding: if it goes the wrong way every KeyPackage is permanently +ecosystem-bound and "interop" degrades to Amethyst speaking two unrelated protocols. Remaining +order: Stage 0 (cordn-side vectors) → Stage 1 (extract the MLS engine) → decide §4.1 → Stage 3-4. ## 2. The coordinator protocol surface @@ -653,7 +661,55 @@ Risk: `MlsGroup.kt` is 4,505 lines and carries the convergence/lifecycle logic. strictly mechanical — extraction and parameterization, no logic edits — so the diff stays reviewable and the test suite is a real check. -### Stage 2 — `:contextvm` module (clean-room) +### Stage 2 — `:contextvm` module (clean-room) — LANDED + +Shipped as `:contextvm`, a KMP module (jvm + android host tests) with `:quartz` as an `api` +dependency, implemented from the specification documents rather than the LGPL SDK. 172 tests, +green on jvm. What is in: + +| Build item | Where | +| ---------- | ----- | +| 1 constants, tags, JSON-RPC codec | `core/CvmKinds`, `core/CvmTags`, `jsonrpc/` | +| 2 minimal MCP client | `mcp/CvmMcpClient`, `mcp/McpMethods` | +| 3 CEP-4/19 gift wrap | `crypto/CvmGiftWrap` (pins `REQUIRED`) | +| 4 correlation + subscription lifecycle | `transport/CvmTransport` | +| 5 CEP-35 discovery learning | `discovery/SessionDiscovery` | +| 6 CEP-6/17/23 discovery | `discovery/ServerDiscovery` | +| 7 **fixture server (Tier C)** | `fixture/CvmFixtureServer`, `fixture/InMemoryRelayPool` | +| 8 CEP-22 receiver | `transfer/oversized/OversizedTransferReceiver` | +| 9 CEP-41 receiver | `transfer/stream/OpenStreamReceiver` | +| 10 CEP-22 sender | `transfer/oversized/OversizedTransferSender` | +| 11 RFC 8785 JCS | `quartz/…/utils/jcs/JsonCanonicalization` | +| 12 CEP-15 schemas | `schema/CommonToolSchema` | +| 13 CEP-8 + CEP-21 | `payment/` | +| 14 CEP-16 injection | in the fixture's server role | +| 15 CEP-24 reviews | `discovery/ServerDiscovery.ServerReview` | +| 16 dual-signer | `transport/DualSigner` | + +Four findings worth carrying forward, all caught by tests rather than review: + +1. **CEP-22/41 ordering.** The first receiver rejected frames whose `progress` did not increase on + arrival, conflating "the sender emits monotonic progress" with "frames arrive in order". Both + CEPs say the opposite: `progress` is the assembly index and explicitly not an arrival-order + guarantee, and receivers may buffer out-of-order chunks. Validation is positional now. +2. **JCS and `Double.MIN_VALUE`.** JVM prints `4.9E-324` where ECMAScript requires `5e-324`, so + "trust the platform to already be shortest" would have hashed differently from every other + implementation. Digits are shortened explicitly until the shortest round-tripping form is + found, which removes the platform assumption entirely. +3. **Event subclassing.** `CvmMessageEvent` began as an `Event` subclass whose `create()` claimed + to return that subclass, but quartz mints subclasses through its own kind-to-class factory, + which knows nothing about 25910. It is a wrapper over `Event` now. +4. **Subscribe-before-publish is an API-shape problem, not a discipline problem.** `CvmTransport` + exposes `request()` with no public publish/subscribe pair, and `InMemoryRelayPool` drops an + event nobody is subscribed to so the property is actually tested rather than assumed. + +Remaining gaps in this stage: Tier B (live integration against +`ghcr.io/cordn-msg/cordn:latest`), Tier D (cross-implementation vectors) and Tier E (a real +wallet for CEP-8) are all unstarted — see §6.4. `testAndroidHostTest` has not been run in a +container yet; Maven Central rate-limits the Android secp256k1 artifact. + +Original scope notes follow. + New Gradle module, peer of `:quic`. Depends on `:quartz` only; no Android framework deps. **Scope, CEP inventory, the fourteen subtle rules and the ordered build list are §6** — this stage is that