mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 03:38:23 +00:00
fix(contextvm): a CEP-22 response never completed its own call
Pointing the live harness at oversized transfer found three bugs that no unit test could have found, because each one was a thing the fixture did that a real server does not. 1. We declared no CEP-35 surface at all. `initialize`'s capabilityTags defaulted to empty and nothing ever passed any, so a spec-correct ContextVM server learned nothing about us — and it only chunks a CEP-22 response, or opens a CEP-41 stream, for a client that declared it can take one. Both profiles were dead on every live session, which is also why msg_sub_many has never had a live proof. The surface now rides the session's first message whatever that message is, since a session that opens with a tool call rather than a handshake would otherwise still declare nothing. 2. A completed transfer could not end the call. `onNotification` returned Unit, so after reassembling the response the transport went back to waiting for a direct one — which a chunking server never sends, the chunked response being the response. It returns JsonRpcMessage? now; non-null completes the call. CEP-41 keeps returning null, because `close` is not an answer. The unit test missed this by having the fixture ALSO answer normally, so the call was ended by that answer and reassembly only had to win a tie-break. The fixture now sends frames and nothing else, and the test fails without the fix. 3. The deadline bounded the whole call. Both profiles deliver one response as a run of notifications, each its own signed, wrapped, published relay event, so a 100 KB transfer ran past the 20 s default and was killed mid-run with every frame arriving and parsing correctly. timeoutMs now bounds silence, which is what a caller actually wants bounded. Subscriptions keep the old meaning through an explicit TimeoutMode.TOTAL: their timeout is a budget the sync loop re-opens on, and a busy stream bounded by silence would never return. Harness: tier-b.sh gains the three tools the two-party lifecycle never reaches — join_request_store, join_request_take_many and kp_remove, via a third account who asks to join rather than being invited — plus a backlog big enough to trip the coordinator's 48000-byte threshold. `cordn fetch --json` reports oversized_transfers so the run can tell a server that chunked from one that never had to; without it a reassembled response is indistinguishable from a direct one, by design. Also corrects an assertion I had backwards: reading a join request does not consume it. Retirement rides the next call after an accept or decline, so a reader that lost its process between listing and accepting has to still find the request — two amy runs are exactly that case. Verified live against the reference coordinator: TIER B PASSED with oversized_transfers=1, and CLIENT INTEROP PASSED unchanged. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_012BfD4txdnsaPRXmNXbup9n
This commit is contained in:
@@ -616,6 +616,10 @@ internal object CordnGroupCommands {
|
||||
"epoch_changes" to epochs,
|
||||
"echoes" to echoes,
|
||||
"undecryptable" to undecryptable,
|
||||
// A fetch is the response most likely to outgrow one relay
|
||||
// event, so this is where CEP-22 shows up if it shows up at
|
||||
// all. Zero is the normal answer and not a warning.
|
||||
"oversized_transfers" to scope.session.oversizedTransfers,
|
||||
),
|
||||
)
|
||||
0
|
||||
|
||||
@@ -42,6 +42,7 @@ bad() { echo " FAIL: $*"; fail=1; }
|
||||
# the only thing carrying continuity between them.
|
||||
alice() { HOME="$WORK/alice" "$AMY" --account alice --secret-backend ncryptsec "$@" 2>/dev/null; }
|
||||
bob() { HOME="$WORK/bob" "$AMY" --account bob --secret-backend ncryptsec "$@" 2>/dev/null; }
|
||||
carol() { HOME="$WORK/carol" "$AMY" --account carol --secret-backend ncryptsec "$@" 2>/dev/null; }
|
||||
field() { python3 -c "import json,sys; d=json.load(sys.stdin); print(json.dumps(d$1) if not isinstance(d$1,str) else d$1)"; }
|
||||
|
||||
trap stack_down EXIT
|
||||
@@ -52,18 +53,24 @@ step "boot geode on $RELAY, and the reference coordinator"
|
||||
stack_up
|
||||
ok "coordinator $COORD"
|
||||
|
||||
step "two accounts"
|
||||
mkdir -p "$WORK/alice" "$WORK/bob"
|
||||
step "three accounts"
|
||||
# Carol is here for the door nobody knocks on in the two-party flow: bob is
|
||||
# invited, so he never sends a join request.
|
||||
mkdir -p "$WORK/alice" "$WORK/bob" "$WORK/carol"
|
||||
alice create --json >/dev/null
|
||||
bob create --json >/dev/null
|
||||
carol create --json >/dev/null
|
||||
ALICE_PK=$(alice whoami --json | field "['hex']")
|
||||
BOB_PK=$(bob whoami --json | field "['hex']")
|
||||
CAROL_PK=$(carol whoami --json | field "['hex']")
|
||||
ok "alice $ALICE_PK"
|
||||
ok "bob $BOB_PK"
|
||||
ok "carol $CAROL_PK"
|
||||
|
||||
step "both remember the coordinator"
|
||||
step "all three remember the coordinator"
|
||||
alice cordn coordinator add --coordinator "$COORD" --relay "$RELAY" --label tier-b --json >/dev/null
|
||||
bob cordn coordinator add --coordinator "$COORD" --relay "$RELAY" --label tier-b --json >/dev/null
|
||||
carol cordn coordinator add --coordinator "$COORD" --relay "$RELAY" --label tier-b --json >/dev/null
|
||||
|
||||
step "the MCP handshake"
|
||||
# The first thing to break if the transport regresses, and the cheapest to
|
||||
@@ -119,6 +126,75 @@ echo "$GOT2" | grep -q "$BACK_ID" && ok "alice read the reply" || bad "alice did
|
||||
# and recognising it needs bookkeeping that survives the process exit.
|
||||
[ "$(echo "$GOT2" | field "['undecryptable']")" = "[]" ] && ok "no self-inflicted gaps" || bad "own traffic came back undecryptable: $(echo "$GOT2" | field "['undecryptable']")"
|
||||
|
||||
# ── The three tools the two-party flow never reaches, and CEP-22 ────────────
|
||||
# Every other step above exercises a tool as a side effect of the lifecycle.
|
||||
# These four do not happen on that path, so they are driven directly.
|
||||
|
||||
step "carol asks to join — join_request_store"
|
||||
REQ=$(carol cordn request --gid "$GID" --json)
|
||||
echo " $REQ"
|
||||
[ "$(echo "$REQ" | field "['gid']")" = "$GID" ] && ok "request stored" || bad "join request not stored"
|
||||
# A ref is an invitation to ask, not membership (spec/01.md §5.3).
|
||||
[ "$(echo "$REQ" | field "['member']")" = "false" ] && ok "asking is not joining" || bad "request reported membership"
|
||||
|
||||
step "alice reads it — join_request_take_many"
|
||||
PEND=$(alice cordn requests list --json)
|
||||
echo " $PEND"
|
||||
[ "$(echo "$PEND" | field "['requests'][0]['pubkey']")" = "$CAROL_PK" ] && ok "sees carol" || bad "no pending request"
|
||||
[ "$(echo "$PEND" | field "['requests'][0]['gid']")" = "$GID" ] && ok "for $GID" || bad "wrong gid on the request"
|
||||
# Listing must NOT consume. A request is retired by answering it, not by
|
||||
# reading it, and the ack rides the next call — so a reader that lost the
|
||||
# process between list and accept has to still find it. Two separate amy
|
||||
# runs is exactly that case.
|
||||
AGAIN=$(alice cordn requests list --json)
|
||||
[ "$(echo "$AGAIN" | field "['requests'][0]['pubkey']")" = "$CAROL_PK" ] && ok "still there for a second reader" || bad "reading a join request consumed it"
|
||||
|
||||
step "alice accepts, carol joins"
|
||||
ACC=$(alice cordn requests accept --pubkey "$CAROL_PK" --json)
|
||||
echo " $ACC"
|
||||
[ "$(echo "$ACC" | field "['answered'][0]['pubkey']")" = "$CAROL_PK" ] && ok "accepted" || bad "accept failed"
|
||||
CJOIN=$(carol cordn join --all --json)
|
||||
[ "$(echo "$CJOIN" | field "['joined']")" = "[\"$GID\"]" ] && ok "carol joined" || bad "carol did not join"
|
||||
# Now it is retired — and the ack for it rode a later call, so this also
|
||||
# proves the retirement survived the process that issued it.
|
||||
GONE=$(alice cordn requests list --json)
|
||||
[ "$(echo "$GONE" | field "['requests']")" = "[]" ] && ok "answering retired it" || bad "an answered request is still pending: $(echo "$GONE" | field "['requests']")"
|
||||
|
||||
step "bob withdraws a KeyPackage — kp_remove"
|
||||
# Two, so there is something left to prove the removal was targeted rather
|
||||
# than a wipe.
|
||||
PUB=$(bob cordn keypackage publish --count 2 --json)
|
||||
KEEP=$(echo "$PUB" | field "['published'][0]['kp_ref']")
|
||||
DROP=$(echo "$PUB" | field "['published'][1]['kp_ref']")
|
||||
WD=$(bob cordn keypackage withdraw --kp-ref "$DROP" --json)
|
||||
echo " $WD"
|
||||
[ "$(echo "$WD" | field "['removed']")" = "[\"$DROP\"]" ] && ok "coordinator confirms the removal" || bad "kp_remove did not confirm $DROP"
|
||||
LIST=$(bob cordn keypackage list --json)
|
||||
echo "$LIST" | grep -q "$DROP" && bad "the withdrawn package is still served" || ok "gone from kp_list"
|
||||
echo "$LIST" | grep -q "$KEEP" && ok "the other one survived" || bad "kp_remove took more than it was asked for"
|
||||
|
||||
step "a response too big for one event — CEP-22"
|
||||
# The reference coordinator switches to oversized transfer at a 48000-byte
|
||||
# published envelope. One message does not reach it; a backlog of them does,
|
||||
# and a fetch is where a backlog is delivered. Nothing below reads the
|
||||
# messages differently — a reassembled response is meant to be invisible —
|
||||
# so the count is the only thing that can tell us the profile ran.
|
||||
BIG=$(python3 -c "print('x' * 6000)")
|
||||
for i in $(seq 12); do
|
||||
alice cordn send --text "chunk-$i $BIG" --json >/dev/null
|
||||
done
|
||||
FETCH=$(bob cordn fetch --json)
|
||||
COUNT=$(echo "$FETCH" | python3 -c "import json,sys; print(len(json.load(sys.stdin)['messages']))")
|
||||
CHUNKED=$(echo "$FETCH" | field "['oversized_transfers']")
|
||||
echo " $COUNT messages, oversized_transfers=$CHUNKED"
|
||||
[ "$COUNT" = "12" ] && ok "all 12 arrived" || bad "expected 12 messages, got $COUNT"
|
||||
[ "${CHUNKED:-0}" -ge 1 ] && ok "reassembled over CEP-22" || bad "the response was never chunked — CEP-22 went unexercised"
|
||||
# Carol was added at a later epoch, so she must read them too: an oversized
|
||||
# response is still one response, not a per-member special case.
|
||||
CGOT=$(carol cordn fetch --json)
|
||||
CCOUNT=$(echo "$CGOT" | python3 -c "import json,sys; print(len(json.load(sys.stdin)['messages']))")
|
||||
[ "$CCOUNT" = "12" ] && ok "carol read the same 12" || bad "carol got $CCOUNT of 12"
|
||||
|
||||
step "both sides agree"
|
||||
A=$(alice cordn group info --json)
|
||||
B=$(bob cordn group info --json)
|
||||
|
||||
+3
@@ -88,6 +88,9 @@ class CordnSession(
|
||||
/** What this coordinator says about itself. Claims, never identity (§8.5). */
|
||||
suspend fun serverInfo(): CoordinatorServerInfo? = scope.serverInfo()
|
||||
|
||||
/** How many of this session's responses arrived reassembled over CEP-22. */
|
||||
val oversizedTransfers: Int get() = scope.coordinator.oversizedTransfers
|
||||
|
||||
internal suspend fun close() = scope.close()
|
||||
}
|
||||
|
||||
|
||||
+17
@@ -85,6 +85,23 @@ class CvmGiftWrap(
|
||||
private val encryptionMode: EncryptionMode = EncryptionMode.REQUIRED,
|
||||
private val giftWrapMode: GiftWrapMode = GiftWrapMode.EPHEMERAL,
|
||||
) {
|
||||
/**
|
||||
* What this side can *receive*, as CEP-35 tag names.
|
||||
*
|
||||
* Receive, not prefer: the transport subscribes to both wrap kinds, so a
|
||||
* client that emits 1059 still accepts 21059 and says so. Declaring only
|
||||
* what we emit would tell a peer to withhold something we can read.
|
||||
*
|
||||
* Empty under [EncryptionMode.DISABLED] - there, a wrap is not something
|
||||
* this client will open at all.
|
||||
*/
|
||||
fun receivableWrapTags(): List<String> =
|
||||
if (encryptionMode == EncryptionMode.DISABLED) {
|
||||
emptyList()
|
||||
} else {
|
||||
listOf(CvmTags.SUPPORT_ENCRYPTION, CvmTags.SUPPORT_ENCRYPTION_EPHEMERAL)
|
||||
}
|
||||
|
||||
/** The wrap kind this client emits. */
|
||||
fun wrapKind(): Kind =
|
||||
when (giftWrapMode) {
|
||||
|
||||
+37
-7
@@ -38,6 +38,7 @@ import com.vitorpamplona.quartz.contextvm.transfer.ProgressEnvelope
|
||||
import com.vitorpamplona.quartz.contextvm.transfer.ProgressToken
|
||||
import com.vitorpamplona.quartz.contextvm.transport.CvmTransport
|
||||
import com.vitorpamplona.quartz.contextvm.transport.DualSigner
|
||||
import com.vitorpamplona.quartz.contextvm.transport.TimeoutMode
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Tag
|
||||
import kotlinx.serialization.json.JsonElement
|
||||
import kotlinx.serialization.json.JsonObject
|
||||
@@ -50,6 +51,15 @@ data class ToolCallResult(
|
||||
val error: com.vitorpamplona.quartz.contextvm.jsonrpc.JsonRpcError? = null,
|
||||
/** Fragments a CEP-41 stream delivered while the call was in flight. */
|
||||
val streamed: List<String> = emptyList(),
|
||||
/**
|
||||
* Whether the response arrived reassembled over CEP-22 rather than in the
|
||||
* single event that carried the rest.
|
||||
*
|
||||
* A caller does not need this to read the result - that is the point of the
|
||||
* profile. It is here so a live interop run can assert the profile actually
|
||||
* fired, instead of passing whether or not the server ever chunked.
|
||||
*/
|
||||
val viaOversizedTransfer: Boolean = false,
|
||||
) {
|
||||
val isError get() = error != null
|
||||
}
|
||||
@@ -83,7 +93,7 @@ class CvmMcpClient(
|
||||
* afford the round trip should do it.
|
||||
*/
|
||||
suspend fun initialize(
|
||||
capabilityTags: List<Tag> = emptyList(),
|
||||
capabilityTags: List<Tag> = transport.selfDiscoveryTags(),
|
||||
protocolVersion: String = PROTOCOL_VERSION,
|
||||
): JsonRpcMessage {
|
||||
val response =
|
||||
@@ -134,6 +144,7 @@ class CvmMcpClient(
|
||||
arguments: JsonObject = buildJsonObject {},
|
||||
identity: DualSigner.Identity = DualSigner.Identity.EPHEMERAL,
|
||||
timeoutMs: Long = CvmTransport.DEFAULT_TIMEOUT_MS,
|
||||
timeoutMode: TimeoutMode = TimeoutMode.IDLE,
|
||||
onStreamFragment: (String) -> Unit = {},
|
||||
): ToolCallResult {
|
||||
val id = nextId()
|
||||
@@ -164,9 +175,10 @@ class CvmMcpClient(
|
||||
),
|
||||
identity = identity,
|
||||
timeoutMs = timeoutMs,
|
||||
timeoutMode = timeoutMode,
|
||||
) { notification ->
|
||||
val envelope = ProgressEnvelope.parseOrNull(notification) ?: return@request
|
||||
if (envelope.token != token) return@request
|
||||
val envelope = ProgressEnvelope.parseOrNull(notification) ?: return@request null
|
||||
if (envelope.token != token) return@request null
|
||||
|
||||
when (envelope.type) {
|
||||
ProgressEnvelope.TYPE_OVERSIZED -> {
|
||||
@@ -175,7 +187,10 @@ class CvmMcpClient(
|
||||
.also { oversized = it }
|
||||
OversizedFrame.parseOrNull(envelope)?.let { frame ->
|
||||
val result = receiver.accept(frame)
|
||||
if (result is OversizedProgressResult.Completed) reassembled = result.message
|
||||
if (result is OversizedProgressResult.Completed) {
|
||||
reassembled = result.message
|
||||
oversizedTransfersCompleted++
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -194,20 +209,35 @@ class CvmMcpClient(
|
||||
|
||||
else -> Unit
|
||||
}
|
||||
|
||||
// A completed CEP-22 transfer IS the response, so handing it
|
||||
// back ends the call. A CEP-41 stream never does: `close` says
|
||||
// no more frames, not that the request is answered.
|
||||
reassembled
|
||||
}
|
||||
|
||||
// A CEP-22 transfer replaces the response that could not be published
|
||||
// directly; a CEP-41 stream does not, since `close` never completes the
|
||||
// JSON-RPC request.
|
||||
val chunked = reassembled != null
|
||||
return when (val effective = reassembled ?: response) {
|
||||
is JsonRpcSuccess -> ToolCallResult(effective.result, streamed = streamed)
|
||||
is JsonRpcFailure -> ToolCallResult(null, effective.error, streamed)
|
||||
else -> ToolCallResult(null, streamed = streamed)
|
||||
is JsonRpcSuccess -> ToolCallResult(effective.result, streamed = streamed, viaOversizedTransfer = chunked)
|
||||
is JsonRpcFailure -> ToolCallResult(null, effective.error, streamed, chunked)
|
||||
else -> ToolCallResult(null, streamed = streamed, viaOversizedTransfer = chunked)
|
||||
}
|
||||
}
|
||||
|
||||
private fun nextId(): JsonRpcId.Num = JsonRpcId.Num(nextId++)
|
||||
|
||||
/**
|
||||
* How many responses this client has reassembled over CEP-22.
|
||||
*
|
||||
* Diagnostics, not control flow. Nothing decides anything on it; it exists
|
||||
* so a live run can tell a server that chunked from one that never had to.
|
||||
*/
|
||||
var oversizedTransfersCompleted: Int = 0
|
||||
private set
|
||||
|
||||
companion object {
|
||||
/** The MCP revision the ContextVM spec's examples use. */
|
||||
const val PROTOCOL_VERSION = "2025-07-02"
|
||||
|
||||
+90
-9
@@ -24,6 +24,7 @@ import com.vitorpamplona.quartz.contextvm.cep04Encryption.CvmGiftWrap
|
||||
import com.vitorpamplona.quartz.contextvm.cep35Discovery.SessionDiscovery
|
||||
import com.vitorpamplona.quartz.contextvm.core.CvmKinds
|
||||
import com.vitorpamplona.quartz.contextvm.core.CvmMessageEvent
|
||||
import com.vitorpamplona.quartz.contextvm.core.CvmTags
|
||||
import com.vitorpamplona.quartz.contextvm.jsonrpc.JsonRpcFailure
|
||||
import com.vitorpamplona.quartz.contextvm.jsonrpc.JsonRpcId
|
||||
import com.vitorpamplona.quartz.contextvm.jsonrpc.JsonRpcMessage
|
||||
@@ -63,6 +64,29 @@ class DualSigner(
|
||||
}
|
||||
}
|
||||
|
||||
/** What a call's `timeoutMs` bounds. */
|
||||
enum class TimeoutMode {
|
||||
/**
|
||||
* Silence. The clock restarts on every message the peer sends, so a
|
||||
* response still being delivered is never cut off.
|
||||
*
|
||||
* The right default for a request/response call: both transfer profiles
|
||||
* deliver one logical response as a run of notifications, each its own
|
||||
* signed, wrapped, published relay event, so how long a response takes is
|
||||
* a function of its size and not something a caller can predict.
|
||||
*/
|
||||
IDLE,
|
||||
|
||||
/**
|
||||
* The whole call, from publish to return.
|
||||
*
|
||||
* For an open-ended subscription, where the timeout is not a failure but
|
||||
* the budget after which the caller re-opens - a busy stream under [IDLE]
|
||||
* would simply never come back.
|
||||
*/
|
||||
TOTAL,
|
||||
}
|
||||
|
||||
/** Thrown when a request cannot be completed at the transport layer. */
|
||||
class CvmTransportException(
|
||||
message: String,
|
||||
@@ -102,6 +126,23 @@ class CvmTransport(
|
||||
private val assumePeerSupportsEncryption: Boolean = true,
|
||||
private val assumePeerSupportsEphemeralWrap: Boolean = true,
|
||||
) {
|
||||
/**
|
||||
* This side's CEP-35 surface: what a peer may use against us.
|
||||
*
|
||||
* CEP-35 is symmetric, and the half we were not doing is the expensive one
|
||||
* to skip. A spec-correct server only chunks a CEP-22 response, or opens a
|
||||
* CEP-41 stream, for a client that declared it can take one — so declaring
|
||||
* nothing does not make us conservative, it makes those profiles dead on
|
||||
* every session and caps every response at one relay event.
|
||||
*
|
||||
* Both transfer profiles are unconditional: [CvmMcpClient] reassembles
|
||||
* CEP-22 and collects CEP-41 fragments on every call, with no flag to turn
|
||||
* either off.
|
||||
*/
|
||||
fun selfDiscoveryTags(): List<Tag> =
|
||||
(crypto.receivableWrapTags() + CvmTags.SUPPORT_OVERSIZED_TRANSFER + CvmTags.SUPPORT_OPEN_STREAM)
|
||||
.map { arrayOf(it) }
|
||||
|
||||
/** The peer's learned discovery baseline, once its first message has arrived. */
|
||||
val peer get() = discovery.peer
|
||||
|
||||
@@ -128,15 +169,25 @@ class CvmTransport(
|
||||
*
|
||||
* @param identity which key signs the request. Anything not required to be
|
||||
* attributable should stay [DualSigner.Identity.EPHEMERAL].
|
||||
* @param onNotification every notification that arrives while the call is
|
||||
* open. Returning a message **completes the call with it**, which is how
|
||||
* CEP-22 finishes: a chunked response replaces the direct one rather than
|
||||
* preceding it, so nothing else would ever arrive to end the wait.
|
||||
* Returning null keeps waiting - CEP-41 is explicit that a stream's
|
||||
* `close` does not complete the request.
|
||||
* @param discoveryTags this side's CEP-35 baseline, sent on the session's
|
||||
* first direct message only.
|
||||
* first direct message only. Defaults to [selfDiscoveryTags] so the
|
||||
* surface rides whatever that first message turns out to be - a session
|
||||
* that opens with a tool call rather than a handshake still declares.
|
||||
* Pass `emptyList()` to declare nothing.
|
||||
*/
|
||||
suspend fun request(
|
||||
message: JsonRpcRequest,
|
||||
identity: DualSigner.Identity = DualSigner.Identity.EPHEMERAL,
|
||||
timeoutMs: Long = DEFAULT_TIMEOUT_MS,
|
||||
discoveryTags: List<Tag> = emptyList(),
|
||||
onNotification: (JsonRpcNotification) -> Unit = {},
|
||||
timeoutMode: TimeoutMode = TimeoutMode.IDLE,
|
||||
discoveryTags: List<Tag> = selfDiscoveryTags(),
|
||||
onNotification: (JsonRpcNotification) -> JsonRpcMessage? = { null },
|
||||
): JsonRpcMessage {
|
||||
val signer = signers.signerFor(identity)
|
||||
val inbound = Channel<Event>(Channel.UNLIMITED)
|
||||
@@ -160,8 +211,12 @@ class CvmTransport(
|
||||
|
||||
relays.publish(outbound(request))
|
||||
|
||||
return withTimeout(timeoutMs) {
|
||||
awaitResponse(inbound, signer, request.id, message.id, onNotification)
|
||||
return when (timeoutMode) {
|
||||
TimeoutMode.IDLE -> awaitResponse(inbound, signer, request.id, message.id, timeoutMs, onNotification)
|
||||
TimeoutMode.TOTAL ->
|
||||
withTimeout(timeoutMs) {
|
||||
awaitResponse(inbound, signer, request.id, message.id, Long.MAX_VALUE, onNotification)
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
subscription.close()
|
||||
@@ -178,16 +233,43 @@ class CvmTransport(
|
||||
relays.publish(outbound(CvmMessageEvent.create(message, serverPubKey, signer)))
|
||||
}
|
||||
|
||||
/**
|
||||
* Waits for the answer, allowing [idleMs] of **silence** between the peer's
|
||||
* messages. [TimeoutMode.TOTAL] passes `Long.MAX_VALUE` and wraps the whole
|
||||
* thing instead.
|
||||
*
|
||||
* A flat deadline cannot express what a CEP-22 or CEP-41 response is. Both
|
||||
* arrive as a run of notifications, each its own signed, wrapped, published
|
||||
* relay event, so a large response takes as long as it takes: a 100 KB
|
||||
* oversized transfer from the reference coordinator ran well past the 20 s
|
||||
* default and was killed mid-run, with the frames arriving and parsing
|
||||
* correctly right up to the cancellation.
|
||||
*
|
||||
* So the clock measures a stalled peer, which is what a caller actually
|
||||
* wants bounded, and only the peer can reset it - events from anyone else
|
||||
* are dropped before this point, or a stranger could hold a call open
|
||||
* indefinitely by publishing noise addressed to us.
|
||||
*/
|
||||
private suspend fun awaitResponse(
|
||||
inbound: Channel<Event>,
|
||||
signer: NostrSigner,
|
||||
requestEventId: HexKey,
|
||||
requestId: JsonRpcId,
|
||||
onNotification: (JsonRpcNotification) -> Unit,
|
||||
idleMs: Long,
|
||||
onNotification: (JsonRpcNotification) -> JsonRpcMessage?,
|
||||
): JsonRpcMessage {
|
||||
for (event in inbound) {
|
||||
while (true) {
|
||||
// withTimeout, not a null-returning variant, so silence still
|
||||
// surfaces as the TimeoutCancellationException callers already
|
||||
// handle - only when the clock starts has changed.
|
||||
// receiveCatching tells a closed subscription from a silent one.
|
||||
val event =
|
||||
withTimeout(idleMs) { inbound.receiveCatching() }
|
||||
.getOrNull()
|
||||
?: throw CvmTransportException("subscription closed before a response arrived")
|
||||
val plain = decryptOrNull(event, signer) ?: continue
|
||||
val wrapped = CvmMessageEvent.fromOrNull(plain) ?: continue
|
||||
if (wrapped.pubKey != serverPubKey) continue
|
||||
|
||||
discovery.observe(wrapped.discoveryTags().toTypedArray())
|
||||
|
||||
@@ -201,7 +283,7 @@ class CvmTransport(
|
||||
}
|
||||
|
||||
when (decoded) {
|
||||
is JsonRpcNotification -> onNotification(decoded)
|
||||
is JsonRpcNotification -> onNotification(decoded)?.let { return it }
|
||||
|
||||
// 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
|
||||
@@ -216,7 +298,6 @@ class CvmTransport(
|
||||
else -> Unit
|
||||
}
|
||||
}
|
||||
throw CvmTransportException("subscription closed before a response arrived")
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
+6
@@ -22,6 +22,7 @@ package com.vitorpamplona.quartz.cordn.spec00Coordinator
|
||||
|
||||
import com.vitorpamplona.quartz.contextvm.mcp.CvmMcpClient
|
||||
import com.vitorpamplona.quartz.contextvm.transport.CvmTransport
|
||||
import com.vitorpamplona.quartz.contextvm.transport.TimeoutMode
|
||||
import com.vitorpamplona.quartz.nip01Core.core.Event
|
||||
import com.vitorpamplona.quartz.nip01Core.core.HexKey
|
||||
import com.vitorpamplona.quartz.nip01Core.core.OptimizedJsonMapper
|
||||
@@ -68,6 +69,8 @@ class CoordinatorClient(
|
||||
/** Performs the MCP handshake. Optional, but it is where CEP-35 tags ride. */
|
||||
suspend fun initialize() = mcp.initialize()
|
||||
|
||||
override val oversizedTransfers get() = mcp.oversizedTransfersCompleted
|
||||
|
||||
/**
|
||||
* Handshakes and reports what the coordinator says about itself.
|
||||
*
|
||||
@@ -316,6 +319,9 @@ class CoordinatorClient(
|
||||
arguments = args,
|
||||
identity = CoordinatorMethod.MSG_SUB_MANY.identity,
|
||||
timeoutMs = timeoutMs,
|
||||
// A budget, not a failure: the caller re-opens on it, and a
|
||||
// busy stream bounded by silence would never come back.
|
||||
timeoutMode = TimeoutMode.TOTAL,
|
||||
onStreamFragment = { fragment ->
|
||||
// A malformed frame is the coordinator's problem, not a
|
||||
// reason to tear down a live subscription over other groups.
|
||||
|
||||
+13
@@ -78,4 +78,17 @@ interface ICoordinator {
|
||||
timeoutMs: Long,
|
||||
onMessage: (GroupMessage) -> Unit,
|
||||
)
|
||||
|
||||
/**
|
||||
* How many responses this coordinator chunked over CEP-22 so far.
|
||||
*
|
||||
* Diagnostics, and the one thing here that is not a tool. Nothing in the
|
||||
* binding reads it; a live interop run does, to tell a server that had to
|
||||
* chunk from one that never did - which is otherwise invisible, because a
|
||||
* reassembled response is indistinguishable from a direct one by design.
|
||||
*
|
||||
* Defaulted so a substitute implementation is not forced to fake a number
|
||||
* about a transfer profile it does not have.
|
||||
*/
|
||||
val oversizedTransfers: Int get() = 0
|
||||
}
|
||||
|
||||
+8
-9
@@ -221,12 +221,13 @@ class CvmMcpClientTest {
|
||||
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)) })
|
||||
}
|
||||
// NO direct response, and that is the whole point: a chunked
|
||||
// response REPLACES the direct one. An earlier version of this test
|
||||
// had the fixture also answer normally, which meant the call was
|
||||
// ended by that answer and the reassembly only had to win a
|
||||
// tie-break. Against the reference coordinator, which sends frames
|
||||
// and nothing else, the same code hung until its deadline.
|
||||
val fixture = fixture { error("a chunked response is the only response") }
|
||||
fixture.start()
|
||||
|
||||
val result =
|
||||
@@ -234,11 +235,9 @@ class CvmMcpClientTest {
|
||||
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()
|
||||
}
|
||||
|
||||
@@ -247,7 +246,7 @@ class CvmMcpClientTest {
|
||||
result.result!!
|
||||
.jsonObject["text"]!!
|
||||
.jsonPrimitive.content,
|
||||
"the reassembled payload replaces the placeholder",
|
||||
"the reassembled payload is the response",
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
+77
@@ -244,6 +244,8 @@ class CvmTransportTest {
|
||||
async {
|
||||
transport.request(JsonRpcRequest(JsonRpcId.Num(0), "ping"), timeoutMs = 5_000) {
|
||||
seen += it.method
|
||||
// null: a notification does not answer the call.
|
||||
null
|
||||
}
|
||||
}
|
||||
yield()
|
||||
@@ -476,6 +478,81 @@ class CvmTransportTest {
|
||||
return ours.map { it.kind }.distinct()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `CVM-35-12 the session's first message declares this side's surface`() =
|
||||
runTest {
|
||||
// CEP-35 is symmetric, and a peer only offers a profile the other
|
||||
// side declared: the reference ContextVM server chunks a CEP-22
|
||||
// response, and opens a CEP-41 stream, only for a client that said
|
||||
// it can take one. A session that declared nothing capped every
|
||||
// response at a single relay event.
|
||||
//
|
||||
// Read off the inner event, because that is where the server reads
|
||||
// it: the wrap's own tags are just routing, and a surface put there
|
||||
// would be both unread and public.
|
||||
val seen = mutableListOf<List<String>>()
|
||||
|
||||
// Deliberately NOT the handshake: the surface has to ride whatever
|
||||
// the first message is, or a client that opens with a tool call
|
||||
// never declares at all.
|
||||
declaring(seen, CvmGiftWrap(encryptionMode = EncryptionMode.DISABLED), "msg_fetch_many")
|
||||
|
||||
val declared = seen.single()
|
||||
assertTrue(CvmTags.SUPPORT_OVERSIZED_TRANSFER in declared, "no CEP-22 flag in $declared")
|
||||
assertTrue(CvmTags.SUPPORT_OPEN_STREAM in declared, "no CEP-41 flag in $declared")
|
||||
// Encryption is off here, so there is no wrap we could open.
|
||||
assertFalse(CvmTags.SUPPORT_ENCRYPTION in declared, "declared a wrap it will not open")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `CVM-35-13 an encrypting client declares both wrap kinds`() =
|
||||
runTest {
|
||||
// Receive, not prefer. This client emits 1059, and still says it
|
||||
// takes 21059, because the transport subscribes to both - saying
|
||||
// otherwise would tell the peer to withhold something we can read.
|
||||
val seen = mutableListOf<List<String>>()
|
||||
|
||||
declaring(seen, CvmGiftWrap(giftWrapMode = GiftWrapMode.PERSISTENT, encryptionMode = EncryptionMode.OPTIONAL), "ping")
|
||||
|
||||
val declared = seen.single()
|
||||
assertTrue(CvmTags.SUPPORT_ENCRYPTION in declared, "no CEP-4 flag in $declared")
|
||||
assertTrue(CvmTags.SUPPORT_ENCRYPTION_EPHEMERAL in declared, "no CEP-19 flag in $declared")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `CVM-35-14 the surface is declared once, not on every message`() =
|
||||
runTest {
|
||||
val seen = mutableListOf<List<String>>()
|
||||
|
||||
declaring(seen, CvmGiftWrap(), "ping", "ping")
|
||||
|
||||
assertEquals(2, seen.size, "both requests should have reached the server")
|
||||
assertTrue(seen.first().isNotEmpty(), "the first message declared nothing")
|
||||
assertEquals(emptyList<String>(), seen.last(), "re-declared on a later message")
|
||||
}
|
||||
|
||||
/**
|
||||
* Runs [methods] in one session, collecting each request's declared
|
||||
* surface - the single-element tags on the event the server unwrapped.
|
||||
*/
|
||||
private suspend fun declaring(
|
||||
into: MutableList<List<String>>,
|
||||
crypto: CvmGiftWrap,
|
||||
vararg methods: String,
|
||||
) {
|
||||
val fixture =
|
||||
server(handler = { request ->
|
||||
into +=
|
||||
request.event.tags
|
||||
.filter { it.size == 1 }
|
||||
.map { it[0] }
|
||||
JsonRpcSuccess(JsonRpcId.Num(0), buildJsonObject { put("ok", JsonPrimitive(true)) })
|
||||
})
|
||||
fixture.start()
|
||||
val transport = transport(crypto)
|
||||
methods.forEach { exchange(fixture, transport, JsonRpcRequest(JsonRpcId.Num(0), it)) }
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `CVM-35-11 the stable identity is used only when asked for`() =
|
||||
runTest {
|
||||
|
||||
Reference in New Issue
Block a user