mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 19:28:25 +00:00
feat(contextvm): transport, MCP client and the Tier C fixture server
Build items 2, 4, 7, 14 and 16 -- the last of Stage 2. CvmTransport exposes request() and deliberately no public publish/subscribe pair. Kind 25910 is ephemeral, so a subscription opened after publishing has missed the response permanently and the failure looks exactly like a flaky relay; making the ordering the transport's job removes the whole class of bug rather than documenting it. Correlation checks both layers, the `e` tag against our request event and the JSON-RPC id against our call, and a notification never resolves a request -- CEP-41 is explicit that close does not complete it, and CEP-8 payment notifications arrive mid-request too. DualSigner makes the stable/ephemeral split explicit in the type so a caller cannot leak the account identity by omission. Everything defaults to ephemeral. CvmMcpClient covers the lifecycle, tool listing and tool calling, and always attaches a progressToken -- without one a server MUST NOT start either transfer profile, so omitting it would silently cap every response at one relay event. Both profiles are wired in: a CEP-22 transfer replaces the response that could not be published directly, while a CEP-41 stream delivers fragments alongside it and leaves the request to its own response. The fixture server is the piece the plan called out as missing. FixtureFaults can mismatch a response id, drop or forge the correlation tag, send a malformed body, go silent, or bury the answer under noise -- none of which a real server does, yet all of which a client MUST handle. InMemoryRelayPool drops an event nobody is subscribed to, which is what makes the subscribe-before-publish property genuinely tested instead of assumed; CVM-CORE-12 asserts the drop so CVM-CORE-11 cannot pass vacuously. It also implements CEP-16 injection, since that is a server obligation with no client-side test. 20 new tests. Module suite at 172 on jvm and 143 on the Android target (the difference is the secp256k1-dependent tests, which live in src/jvmTest). Verified by mutation: ignoring the reassembled CEP-22 payload, and letting a notification resolve a request, each kill exactly their guarding tests. Also adds contextvm/README.md recording the implemented spec revision (contextvm-docs e63bce6) as §6.7 requires, and updates the plan: Stage 2 is landed, and the §4.1 gate it carried was too broad -- ContextVM is credential-agnostic, so only Stage 3 depends on that decision. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_012BfD4txdnsaPRXmNXbup9n
This commit is contained in:
@@ -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.
|
||||
+284
@@ -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<Kind>,
|
||||
val onEvent: (Event) -> Unit,
|
||||
)
|
||||
|
||||
private val listeners = mutableListOf<Listener>()
|
||||
|
||||
/** Every event published, in order. For assertions about what went on the wire. */
|
||||
val published = mutableListOf<Event>()
|
||||
|
||||
/** Events dropped because nothing was listening for them. */
|
||||
val dropped = mutableListOf<Event>()
|
||||
|
||||
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<Tag> = 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<JsonObject?>()
|
||||
|
||||
/** 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<Event>()
|
||||
|
||||
/**
|
||||
* 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<Event>(
|
||||
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"}"""
|
||||
}
|
||||
}
|
||||
@@ -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<String> = 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<Tag> = 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<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,
|
||||
) { 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"
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
+235
@@ -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<Tag> = emptyList(),
|
||||
onNotification: (JsonRpcNotification) -> Unit = {},
|
||||
): JsonRpcMessage {
|
||||
val signer = signers.signerFor(identity)
|
||||
val inbound = Channel<Event>(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<Event>,
|
||||
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
|
||||
}
|
||||
}
|
||||
@@ -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<String>()
|
||||
|
||||
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",
|
||||
)
|
||||
}
|
||||
}
|
||||
+326
@@ -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<Array<String>> = 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<JsonRpcSuccess>(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<TimeoutCancellationException> {
|
||||
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<TimeoutCancellationException> {
|
||||
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<TimeoutCancellationException> {
|
||||
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<JsonRpcSuccess>(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<TimeoutCancellationException> {
|
||||
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<String>()
|
||||
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<JsonRpcSuccess>(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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user