diff --git a/.github/workflows/package-validation.yml b/.github/workflows/package-validation.yml index f09aa05d7..71fb73175 100644 --- a/.github/workflows/package-validation.yml +++ b/.github/workflows/package-validation.yml @@ -351,6 +351,9 @@ jobs: scripts/verify-macos-binary.sh "release-binaries/${OUTPUT}" "$EXPECTED_ARCH" -- --version echo "STANDALONE_CLI=$PWD/release-binaries/${OUTPUT}" >> "$GITHUB_ENV" + - name: Verify standalone MCP dispatches each request once + run: node scripts/verify-standalone-mcp-single-dispatch.mjs "$STANDALONE_CLI" + - name: Smoke standalone lifecycle env: AGENT_RELAY_STARTUP_DEBUG: 1 diff --git a/.github/workflows/publish.yml b/.github/workflows/publish.yml index f7dc51c42..4fa5768d6 100644 --- a/.github/workflows/publish.yml +++ b/.github/workflows/publish.yml @@ -287,6 +287,10 @@ jobs: fi echo "OK: ssh2 is bundled (ssh-userauth symbol present)" + - name: Verify standalone MCP dispatches each request once + if: matrix.target == 'bun-linux-x64' + run: node scripts/verify-standalone-mcp-single-dispatch.mjs "release-binaries/${{ matrix.binary_name }}" + - name: Sign macOS binary if: startsWith(matrix.target, 'bun-darwin-') run: scripts/sign-macos-binary.sh release-binaries/${{ matrix.binary_name }} diff --git a/CHANGELOG.md b/CHANGELOG.md index bf946055d..175398f1a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,7 +5,12 @@ All notable changes to Agent Relay will be documented in this file. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). -## [Unreleased] +## [Unreleased - Patch] + +### Fixed + +- `agent-relay mcp` standalone binaries no longer dispatch each tool call twice. +- MCP `send_dm` now forwards idempotency keys so keyed retries do not store duplicate messages. ## [13.0.1] - 2026-10-02 diff --git a/packages/cli/src/cli/agent-relay-mcp.startup.test.ts b/packages/cli/src/cli/agent-relay-mcp.startup.test.ts index f64593618..96d615bc5 100644 --- a/packages/cli/src/cli/agent-relay-mcp.startup.test.ts +++ b/packages/cli/src/cli/agent-relay-mcp.startup.test.ts @@ -545,6 +545,9 @@ describe('agent-relay-mcp startup helpers', () => { expect(retry).toEqual(first); expect(mocks.agentClients.reduce((count, client) => count + client.dm.mock.calls.length, 0)).toBe(1); + expect(mocks.agentClients.flatMap((client) => client.dm.mock.calls)).toEqual([ + [input.to, input.text, expect.objectContaining({ idempotencyKey: input.idempotency_key })], + ]); }); it('parses startup options and helper flags from the environment', async () => { diff --git a/packages/cli/src/cli/agent-relay-mcp.ts b/packages/cli/src/cli/agent-relay-mcp.ts index 764bcbb24..ecd4ac37b 100644 --- a/packages/cli/src/cli/agent-relay-mcp.ts +++ b/packages/cli/src/cli/agent-relay-mcp.ts @@ -24,6 +24,7 @@ import { } from '@agent-relay/sdk'; import { z } from 'zod'; import { declaredWorkforceMetadata } from './lib/registration-metadata.js'; +import { isBundledBunEntrypointPath } from './lib/agent-relay-mcp-command.js'; import { DEFAULT_AGENT_REGISTRATION_TIMEOUT_MS, withAgentRegistrationDeadline, @@ -496,6 +497,11 @@ export function normalizeBaseUrl(baseUrl?: string): string | undefined { function isEntrypoint(): boolean { const invocationPath = process.argv[1]; if (!invocationPath) return false; + // `bun build --compile` folds this importable module into the CLI entrypoint, + // and both modules observe the same virtual import.meta.url. Treating this + // module as an entrypoint in that environment starts a second stdio server; + // the CLI's `mcp` command owns startup for the compiled executable. + if (isBundledBunEntrypointPath(invocationPath)) return false; try { return fs.realpathSync(invocationPath) === fs.realpathSync(fileURLToPath(import.meta.url)); } catch { diff --git a/packages/cli/src/cli/mcp/messaging-tools.protocol.test.ts b/packages/cli/src/cli/mcp/messaging-tools.protocol.test.ts index 5e48b45ff..56c7b075d 100644 --- a/packages/cli/src/cli/mcp/messaging-tools.protocol.test.ts +++ b/packages/cli/src/cli/mcp/messaging-tools.protocol.test.ts @@ -3,6 +3,7 @@ import { InMemoryTransport } from '@modelcontextprotocol/sdk/inMemory.js'; import { McpServer } from '@modelcontextprotocol/sdk/server/mcp.js'; import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import { McpRequestReplay } from './request-replay.js'; import { registerMessagingTools } from './messaging-tools.js'; describe('messaging delivery receipts over MCP', () => { @@ -14,6 +15,94 @@ describe('messaging delivery receipts over MCP', () => { vi.unstubAllEnvs(); }); + it('reuses keyed receipts across two replay boundaries while unkeyed identical sends stay distinct', async () => { + const rows: { id: string; conversationId: string }[] = []; + const keyed = new Map(); + const dm = vi.fn(async (_to: string, _text: string, options: { idempotencyKey?: string }) => { + const key = options.idempotencyKey; + if (key && keyed.has(key)) return keyed.get(key)!; + const row = { id: `msg_${rows.length + 1}`, conversationId: 'dm_chief' }; + rows.push(row); + if (key) keyed.set(key, row); + return row; + }); + async function send(idempotencyKey?: string) { + const server = new McpServer({ name: 'boundary-test', version: '1.0.0' }); + registerMessagingTools( + server, + () => ({ dm }) as never, + async () => [{ name: 'chief' }], + new McpRequestReplay() + ); + const client = new Client({ name: 'boundary-client', version: '1.0.0' }); + const [clientTransport, serverTransport] = InMemoryTransport.createLinkedPair(); + try { + await server.connect(serverTransport); + await client.connect(clientTransport); + const result = await client.callTool({ + name: 'send_dm', + arguments: { + to: 'chief', + text: 'same text', + ...(idempotencyKey ? { idempotency_key: idempotencyKey } : {}), + }, + }); + expect(result.isError).not.toBe(true); + return result.structuredContent; + } finally { + await client.close(); + await server.close(); + } + } + const first = await send('logical-send'); + const retry = await send('logical-send'); + expect(rows).toHaveLength(1); + expect(first).toMatchObject({ id: 'msg_1', conversationId: 'dm_chief' }); + expect(retry).toEqual(first); + expect(dm).toHaveBeenCalledTimes(2); + const unkeyedFirst = await send(); + const unkeyedSecond = await send(); + expect(rows).toHaveLength(3); + expect(unkeyedFirst?.id).not.toBe(unkeyedSecond?.id); + }); + + it('rejects a whitespace-only idempotency key before sending and forwards it trimmed', async () => { + const dm = vi.fn(async () => ({ id: 'msg_1', conversationId: 'dm_chief' })); + const server = new McpServer({ name: 'idempotency-key-test', version: '1.0.0' }); + registerMessagingTools( + server, + () => ({ dm }) as never, + async () => [{ name: 'chief' }], + new McpRequestReplay() + ); + const client = new Client({ name: 'idempotency-key-client', version: '1.0.0' }); + const [clientTransport, serverTransport] = InMemoryTransport.createLinkedPair(); + try { + await server.connect(serverTransport); + await client.connect(clientTransport); + const rejected = await client.callTool({ + name: 'send_dm', + arguments: { to: 'chief', text: 'same text', idempotency_key: ' ' }, + }); + expect(rejected.isError).toBe(true); + expect(JSON.stringify(rejected.content)).toContain('idempotency_key'); + expect(dm).not.toHaveBeenCalled(); + const result = await client.callTool({ + name: 'send_dm', + arguments: { to: 'chief', text: 'same text', idempotency_key: ' logical-send ' }, + }); + expect(result.isError).not.toBe(true); + expect(dm).toHaveBeenCalledWith( + 'chief', + 'same text', + expect.objectContaining({ idempotencyKey: 'logical-send' }) + ); + } finally { + await client.close(); + await server.close(); + } + }); + it('exposes enqueue state on send and an explicit signal for an empty reader list', async () => { const dm = vi.fn(async () => ({ id: 'msg_1', text: 'hello' })); const readers = vi.fn(async () => []); diff --git a/packages/cli/src/cli/mcp/messaging-tools.ts b/packages/cli/src/cli/mcp/messaging-tools.ts index 46c798113..22d323524 100644 --- a/packages/cli/src/cli/mcp/messaging-tools.ts +++ b/packages/cli/src/cli/mcp/messaging-tools.ts @@ -350,11 +350,15 @@ export function registerMessagingTools( attachments: z.array(z.string()).optional().describe('File attachment IDs'), idempotency_key: z .string() + // Relaycast trims this key upstream, so trim before both the local + // replay cache and the forwarded call to keep one logical send on one + // key; a whitespace-only key would otherwise become an unkeyed send. + .trim() .min(1) .max(255) .optional() .describe( - 'Stable key for retrying this same message after a lost response; use a new key for a new message.' + 'Stable key for retrying this same message after a lost response; use a new key for a new message. Surrounding whitespace is trimmed and a whitespace-only key is rejected.' ), ...identityOverrideInputShape, }, @@ -371,6 +375,7 @@ export function registerMessagingTools( const agents = await listAgentsForRecipientResolution?.(); const resolvedRecipient = agents ? resolveExactAgentName(agents, to) : undefined; const message = await getAgentClient(as).dm(to, text, { + idempotencyKey: idempotency_key, mode, attachments, data: replayMessageMetadata(), diff --git a/packages/sdk/src/__tests__/thin-client.test.ts b/packages/sdk/src/__tests__/thin-client.test.ts index 62f7fe246..f7851158e 100644 --- a/packages/sdk/src/__tests__/thin-client.test.ts +++ b/packages/sdk/src/__tests__/thin-client.test.ts @@ -199,10 +199,15 @@ describe('createAgentClient', () => { }); expect(instance.as).toHaveBeenCalledWith('at_live_test', { autoHeartbeatMs: false }); - const agent = relaycastMocks.agentClients[0] as { send: ReturnType }; + const agent = relaycastMocks.agentClients[0] as { + send: ReturnType; + dm: ReturnType; + }; const sent = await client.send('general', 'hello', { mode: 'wait' }); expect(agent.send).toHaveBeenCalledWith('general', 'hello', { mode: 'wait' }); expect(sent).toEqual({ raw: true, args: ['general', 'hello', { mode: 'wait' }] }); + await client.dm('chief', 'dm', { idempotencyKey: 'dm-logical-1' }); + expect(agent.dm).toHaveBeenCalledWith('chief', 'dm', { idempotencyKey: 'dm-logical-1' }); }); it('honors an explicit heartbeat interval', () => { @@ -239,13 +244,16 @@ describe('createAgentClient', () => { await client.send('general', 'hello'); await client.reply('msg_parent', 'reply'); - await client.dm('chief', 'dm'); + await client.dm('chief', 'dm', { idempotencyKey: 'dm-logical-1' }); await client.dms.sendMessage('conv_group', 'group'); const replayData = { session_ref: '11111111-1111-4111-8111-111111111111' }; expect(agent.send).toHaveBeenCalledWith('general', 'hello', { data: replayData }); expect(agent.reply).toHaveBeenCalledWith('msg_parent', 'reply', { data: replayData }); - expect(agent.dm).toHaveBeenCalledWith('chief', 'dm', { data: replayData }); + expect(agent.dm).toHaveBeenCalledWith('chief', 'dm', { + data: replayData, + idempotencyKey: 'dm-logical-1', + }); expect(agent.dms.sendMessage).toHaveBeenCalledWith('conv_group', 'group', { data: replayData }); }); diff --git a/packages/sdk/src/messaging/thin-client.ts b/packages/sdk/src/messaging/thin-client.ts index 087b45726..97e225e13 100644 --- a/packages/sdk/src/messaging/thin-client.ts +++ b/packages/sdk/src/messaging/thin-client.ts @@ -119,6 +119,8 @@ export interface RelayAgentThinClient { to: string, text: string, options?: { + /** Stable key for retries of one logical send, including across processes. */ + idempotencyKey?: string; mode?: RelayMessageMode; attachments?: string[]; data?: Record | null; diff --git a/scripts/verify-standalone-mcp-single-dispatch.mjs b/scripts/verify-standalone-mcp-single-dispatch.mjs new file mode 100644 index 000000000..59c74a141 --- /dev/null +++ b/scripts/verify-standalone-mcp-single-dispatch.mjs @@ -0,0 +1,245 @@ +#!/usr/bin/env node + +import { spawn } from 'node:child_process'; +import fs from 'node:fs'; +import os from 'node:os'; +import path from 'node:path'; + +const binary = process.argv[2]; +if (!binary) { + process.stderr.write('Usage: verify-standalone-mcp-single-dispatch.mjs \n'); + process.exit(2); +} + +const STARTUP_TIMEOUT_MS = 60_000; +const RESPONSE_TIMEOUT_MS = 15_000; +const DUPLICATE_GRACE_MS = 1_000; +const SHUTDOWN_GRACE_MS = 2_000; +const DIAGNOSTIC_BYTES = 8_000; +const DIAGNOSTIC_OUTPUT_BYTES = 2_000; +const initializeId = 'single-dispatch-initialize'; +const probeId = 'single-dispatch-probe'; + +const isolatedHome = fs.mkdtempSync(path.join(os.tmpdir(), 'agent-relay-mcp-dispatch-')); +const env = { ...process.env }; +for (const key of [ + 'RELAY_WORKSPACES_JSON', + 'RELAY_WORKSPACE_KEY', + 'AGENT_RELAY_WORKSPACE_KEY', + 'RELAY_API_KEY', + 'RELAY_AGENT_TOKEN', + 'RELAY_AGENT_NAME', +]) { + delete env[key]; +} +Object.assign(env, { + HOME: isolatedHome, + AGENT_RELAY_SKIP_UPDATE_CHECK: '1', + AGENT_RELAY_TELEMETRY_DISABLED: '1', +}); + +const child = spawn(binary, ['mcp'], { env, stdio: ['pipe', 'pipe', 'pipe'] }); +const responseCounts = new Map(); +const responses = new Map(); +const waiters = new Set(); +let stdoutBuffer = ''; +let stdoutDiagnostic = ''; +let stderrDiagnostic = ''; +let startupError; +let stdinError; +let exitResult; + +const delay = (milliseconds) => new Promise((resolve) => setTimeout(resolve, milliseconds)); +const boundedAppend = (current, chunk) => `${current}${chunk}`.slice(-DIAGNOSTIC_BYTES); +const notifyWaiters = () => { + for (const waiter of [...waiters]) waiter(); +}; + +const spawned = new Promise((resolve, reject) => { + child.once('spawn', resolve); + child.once('error', reject); +}); +const exited = new Promise((resolve) => { + child.once('exit', (code, signal) => { + exitResult = { code, signal }; + notifyWaiters(); + resolve(exitResult); + }); +}); + +child.on('error', (error) => { + startupError = error; + notifyWaiters(); +}); +child.stdin.on('error', (error) => { + stdinError = error; + notifyWaiters(); +}); +child.stdout.setEncoding('utf8'); +child.stdout.on('data', (chunk) => { + stdoutDiagnostic = boundedAppend(stdoutDiagnostic, chunk); + stdoutBuffer += chunk; + let newline; + while ((newline = stdoutBuffer.indexOf('\n')) >= 0) { + const line = stdoutBuffer.slice(0, newline).trim(); + stdoutBuffer = stdoutBuffer.slice(newline + 1); + if (!line) continue; + try { + const message = JSON.parse(line); + if (message.id === initializeId || message.id === probeId) { + responseCounts.set(message.id, (responseCounts.get(message.id) ?? 0) + 1); + responses.set(message.id, message); + notifyWaiters(); + } + } catch { + // Retain bounded raw output for the redacted failure diagnostic. + } + } +}); +child.stderr.setEncoding('utf8'); +child.stderr.on('data', (chunk) => { + stderrDiagnostic = boundedAppend(stderrDiagnostic, chunk); +}); + +async function withTimeout(promise, milliseconds, message) { + let timeout; + try { + return await Promise.race([ + promise, + new Promise((_, reject) => { + timeout = setTimeout(() => reject(new Error(message)), milliseconds); + }), + ]); + } finally { + clearTimeout(timeout); + } +} + +function waitForResponse(id, milliseconds) { + return new Promise((resolve, reject) => { + const timeout = setTimeout(() => { + cleanup(); + reject(new Error(`Timed out after ${milliseconds} ms waiting for MCP response ${id}.`)); + }, milliseconds); + const check = () => { + if ((responseCounts.get(id) ?? 0) >= 1) { + cleanup(); + resolve(); + } else if (startupError) { + cleanup(); + reject(startupError); + } else if (stdinError) { + cleanup(); + reject(stdinError); + } else if (exitResult) { + cleanup(); + reject( + new Error( + `MCP child exited before response ${id} (code ${exitResult.code}, signal ${exitResult.signal}).` + ) + ); + } + }; + const cleanup = () => { + clearTimeout(timeout); + waiters.delete(check); + }; + waiters.add(check); + check(); + }); +} + +function send(message) { + return new Promise((resolve, reject) => { + child.stdin.write(`${JSON.stringify(message)}\n`, (error) => { + if (error) reject(error); + else resolve(); + }); + }); +} + +async function stopChild() { + if (exitResult || startupError) return true; + child.kill('SIGTERM'); + if (await Promise.race([exited.then(() => true), delay(SHUTDOWN_GRACE_MS).then(() => false)])) { + return true; + } + child.kill('SIGKILL'); + if (await Promise.race([exited.then(() => true), delay(SHUTDOWN_GRACE_MS).then(() => false)])) { + return true; + } + // Do not let inherited stdio keep the verifier alive if the OS never reports + // the forced exit. The failure below still makes the publish gate red. + child.stdin.destroy(); + child.stdout.destroy(); + child.stderr.destroy(); + child.unref(); + return false; +} + +function safeDiagnostic(value) { + return value + .replace(/\b(?:at|rk)_[A-Za-z0-9_-]{8,}\b/g, '[REDACTED]') + .replace(/(Bearer\s+)[^\s"']+/gi, '$1[REDACTED]') + .replace( + /("(?:token|workspace_?key|api_?key|authorization|credential|secret)"\s*:\s*")[^"]+/gi, + '$1[REDACTED]' + ) + .slice(-DIAGNOSTIC_OUTPUT_BYTES); +} + +let failure; +try { + await withTimeout( + spawned, + STARTUP_TIMEOUT_MS, + `Timed out after ${STARTUP_TIMEOUT_MS} ms starting MCP child.` + ); + await send({ + jsonrpc: '2.0', + id: initializeId, + method: 'initialize', + params: { + protocolVersion: '2025-06-18', + capabilities: {}, + clientInfo: { name: 'standalone-single-dispatch-check', version: '1' }, + }, + }); + // The initialize response is the readiness signal. Stdin buffers the request + // while optional Cloud discovery runs, so cold startup has no timing race. + await waitForResponse(initializeId, STARTUP_TIMEOUT_MS); + await send({ jsonrpc: '2.0', method: 'notifications/initialized' }); + await send({ jsonrpc: '2.0', id: probeId, method: 'tools/list', params: {} }); + await waitForResponse(probeId, RESPONSE_TIMEOUT_MS); + // Two folded stdio servers answer the same request nearly together. Keep a + // bounded grace window after the protocol-driven response to observe both. + await delay(DUPLICATE_GRACE_MS); + + for (const id of [initializeId, probeId]) { + const count = responseCounts.get(id) ?? 0; + if (count !== 1) throw new Error(`Expected one MCP response for ${id}, received ${count}.`); + if (responses.get(id)?.error) throw new Error(`MCP response ${id} returned an error.`); + } +} catch (error) { + failure = error; +} finally { + const stopped = await stopChild(); + fs.rmSync(isolatedHome, { recursive: true, force: true }); + if (!stopped && !failure) + failure = new Error('MCP child did not exit after SIGTERM and SIGKILL deadlines.'); +} + +if (failure) { + const message = safeDiagnostic(failure instanceof Error ? failure.message : String(failure)); + const stdout = safeDiagnostic(stdoutDiagnostic); + const stderr = safeDiagnostic(stderrDiagnostic); + process.stderr.write( + `Standalone MCP single-dispatch verification failed: ${message}` + + (stdout ? `\nMCP stdout (redacted, bounded):\n${stdout}` : '') + + (stderr ? `\nMCP stderr (redacted, bounded):\n${stderr}` : '') + + '\n' + ); + process.exit(1); +} + +process.stdout.write('Standalone MCP single-dispatch verification passed.\n');