| 1 | // Real official SDK HTTP/SSE transports in the committed builtin bundle. |
| 2 | // FetchProxy network/ticket authority is simulated here; Rust acceptance tests |
| 3 | // separately run the real broker and guarded McpHttpClient against loopback. |
| 4 | import { test } from 'node:test' |
| 5 | import assert from 'node:assert/strict' |
| 6 | import { randomUUID } from 'node:crypto' |
| 7 | import { createMcpModule } from '../dist/builtin/mcp.mjs' |
| 8 | import { validateMessage } from '../dist/protocol.mjs' |
| 9 | |
| 10 | const owner = { plugin_id: 'host:mcp', generation: 1, owner_token: 'builtin-test-token' } |
| 11 | const initParams = { protocolVersion: '2025-06-18', capabilities: {}, clientInfo: { name: 'codewhale-tui', version: '0.10.1' } } |
| 12 | class FakeFetchProxy { |
| 13 | session = randomUUID() |
| 14 | grants = new Map() |
| 15 | bodies = new Map() |
| 16 | frames = [] |
| 17 | fetches = [] |
| 18 | closed = false |
| 19 | sequence = 0 |
| 20 | failWrite = false |
| 21 | status = undefined |
| 22 | silent = false |
| 23 | constructor(legacy = false, ssePost = false) { this.legacy = legacy; this.ssePost = ssePost } |
| 24 | grant(method, params) { |
| 25 | const grant = { ticket: randomUUID(), operation_id: randomUUID(), method, ...(method.startsWith('notifications/') ? {} : { wire_id: String(this.sequence++) }), params: structuredClone(params) } |
| 26 | this.grants.set(grant.ticket, structuredClone(grant)); return grant |
| 27 | } |
| 28 | body(bytes, open = false) { const id = randomUUID(); this.bodies.set(id, { bytes: Array.from(Buffer.from(bytes)), open, waiting: undefined }); return id } |
| 29 | append(id, bytes) { const body = this.bodies.get(id); if (!body) return; body.bytes.push(...Buffer.from(bytes)); if (body.waiting) { const take = body.waiting; body.waiting = undefined; take() } } |
| 30 | async request(method, params, signal) { |
| 31 | assert.deepEqual(params.owner, owner); assert.equal(params.session_id, this.session) |
| 32 | if (method === 'net/start') { assert.equal(params.ticket, 'launch-once'); return {} } |
| 33 | if (method === 'net/close') { this.closed = true; for (const body of this.bodies.values()) body.waiting?.(); this.bodies.clear(); return {} } |
| 34 | if (method === 'net/release') { this.bodies.delete(params.response_id); return {} } |
| 35 | if (method === 'net/read') { |
| 36 | const body = this.bodies.get(params.response_id); assert.ok(body) |
| 37 | if (!body.bytes.length && body.open) await new Promise((resolve, reject) => { |
| 38 | const abort = () => { body.waiting = undefined; reject(new Error('fake response cancelled')) } |
| 39 | body.waiting = () => { signal?.removeEventListener('abort', abort); resolve() } |
| 40 | if (signal?.aborted) abort(); else signal?.addEventListener('abort', abort, { once: true }) |
| 41 | }) |
| 42 | if (this.closed) throw new Error('fake response closed') |
| 43 | const data = body.bytes.splice(0, 32768); return { data, done: !body.open && !body.bytes.length } |
| 44 | } |
| 45 | assert.equal(method, 'net/fetch') |
| 46 | this.fetches.push(structuredClone(params)) |
| 47 | assert.equal(new URL(params.url).origin, 'https://mcp-proxy.invalid') |
| 48 | assert.ok(Object.keys(params.headers).every(key => ['accept', 'content_type', 'mcp_session_id', 'mcp_protocol_version'].includes(key))) |
| 49 | if (params.method === 'GET') { |
| 50 | assert.equal(params.url, `https://mcp-proxy.invalid/${this.session}`) |
| 51 | if (!this.legacy) return { status: 405, headers: {} } |
| 52 | this.stream = this.body(this.holdEndpoint ? '' : `event: endpoint\ndata: https://mcp-proxy.invalid/${this.session}/endpoint\n\n`, true) |
| 53 | return { status: 200, headers: { 'content-type': 'text/event-stream' }, response_id: this.stream } |
| 54 | } |
| 55 | assert.equal(params.method, 'POST') |
| 56 | assert.equal(params.url, `https://mcp-proxy.invalid/${this.session}${this.legacy ? '/endpoint' : ''}`) |
| 57 | const frame = params.frame, grant = this.grants.get(params.ticket) |
| 58 | assert.ok(grant, 'operation grant is single-use') |
| 59 | assert.equal(params.operation_id, grant.operation_id); assert.equal(frame.method, grant.method) |
| 60 | assert.deepEqual(frame.params ?? {}, grant.params) |
| 61 | if (grant.wire_id !== undefined) assert.equal(frame.id, grant.wire_id) |
| 62 | this.grants.delete(params.ticket); this.frames.push(structuredClone(frame)) |
| 63 | if (this.negotiate && !this.legacy) { |
| 64 | this.legacy = true |
| 65 | const replacement = { ...structuredClone(grant), ticket: randomUUID() } |
| 66 | if (this.mutateNegotiation) replacement.params = { forged: true } |
| 67 | this.grants.set(replacement.ticket, replacement) |
| 68 | return { status: 405, headers: {}, legacy_grant: replacement } |
| 69 | } |
| 70 | if (this.failWrite && frame.method === 'tools/call') throw new Error('fake unknown partial HTTP write') |
| 71 | if (this.status !== undefined && frame.method !== 'initialize') return { status: this.status, headers: {} } |
| 72 | if (!Object.hasOwn(frame, 'id')) return { status: 202, headers: {} } |
| 73 | assert.equal(typeof frame.id, 'string', 'real numeric SDK IDs must not reach the peer') |
| 74 | const result = this.customResult?.(frame.method) ?? (frame.method === 'initialize' ? { protocolVersion: '2025-06-18', capabilities: { tools: {}, resources: {}, prompts: {} }, serverInfo: { name: 'fetch-fixture', version: '1' } } |
| 75 | : frame.method === 'tools/list' ? { tools: [{ name: 'echo', inputSchema: { type: 'object' } }], nextCursor: 'next-page' } |
| 76 | : { content: [{ type: 'text', text: 'ok' }] }) |
| 77 | const response = JSON.stringify({ jsonrpc: '2.0', id: frame.id, result }) |
| 78 | if (this.legacy) { |
| 79 | if (!this.silent) this.append(this.stream, `event: message\ndata: ${response}\n\n`) |
| 80 | return { status: 202, headers: {} } |
| 81 | } |
| 82 | const headers = { 'content-type': this.ssePost ? 'text/event-stream' : 'application/json', ...(frame.method === 'initialize' ? { 'mcp-session-id': 'rust-observed-session' } : {}) } |
| 83 | return { status: 200, headers, response_id: this.body(this.silent ? '' : this.ssePost ? `event: message\ndata: ${response}\n\n` : response, this.silent) } |
| 84 | } |
| 85 | } |
| 86 | async function fixture(legacy = false, ssePost = false) { |
| 87 | const broker = new FakeFetchProxy(legacy, ssePost), module = createMcpModule(broker, owner) |
| 88 | const open = { owner, session_id: broker.session, transport: legacy ? 'sse' : 'http', launch_ticket: 'launch-once', initialize_grant: broker.grant('initialize', initParams), initialized_grant: broker.grant('notifications/initialized', {}), client_version: '0.10.1', deadline_ms: 1000 } |
| 89 | const signal = new AbortController().signal |
| 90 | await module.open(open, signal); return { broker, module, open, signal } |
| 91 | } |
| 92 | for (const [name, legacy, ssePost] of [['HTTP JSON', false, false], ['HTTP SSE response', false, true], ['legacy SSE', true, false]]) { |
| 93 | test(`official ${name} transport uses opaque FetchProxy and one exact per-page Rust wire ID`, async () => { |
| 94 | const f = await fixture(legacy, ssePost) |
| 95 | try { |
| 96 | const grant = f.broker.grant('tools/list', { cursor: 'one-page' }) |
| 97 | const result = await f.module.request({ owner, session_id: f.open.session_id, grant, deadline_ms: 1000 }, f.signal) |
| 98 | assert.equal(result.nextCursor, 'next-page') |
| 99 | assert.equal(f.broker.frames.filter(frame => frame.method === 'tools/list').length, 1) |
| 100 | assert.equal(f.broker.frames.find(frame => frame.method === 'tools/list').id, grant.wire_id) |
| 101 | assert.ok(f.broker.frames.filter(frame => Object.hasOwn(frame, 'id')).every(frame => typeof frame.id === 'string')) |
| 102 | if (!legacy) assert.equal(f.broker.fetches.find(params => params.frame?.method === 'tools/list').headers.mcp_session_id, 'rust-observed-session') |
| 103 | } finally { await f.module.dispose() } |
| 104 | assert.equal(f.broker.closed, true) |
| 105 | }) |
| 106 | } |
| 107 | |
| 108 | test('HTTP decoded grant mismatch fails before a FetchProxy write and closes the session', async () => { |
| 109 | const f = await fixture() |
| 110 | const grant = f.broker.grant('tools/call', { name: 'echo', arguments: {} }) |
| 111 | grant.params.arguments = { forged: true } |
| 112 | await assert.rejects(f.module.request({ owner, session_id: f.open.session_id, grant, deadline_ms: 1000 }, f.signal), /exact|Connection closed/) |
| 113 | assert.equal(f.broker.frames.some(frame => frame.method === 'tools/call'), false) |
| 114 | assert.equal(f.broker.closed, true); await f.module.dispose() |
| 115 | }) |
| 116 | test('unknown partial HTTP write retires once and never replays', async () => { |
| 117 | const f = await fixture(); f.broker.failWrite = true |
| 118 | await assert.rejects(f.module.request({ owner, session_id: f.open.session_id, grant: f.broker.grant('tools/call', { name: 'echo', arguments: {} }), deadline_ms: 1000 }, f.signal), /partial HTTP write|Connection closed/) |
| 119 | assert.equal(f.broker.frames.filter(frame => frame.method === 'tools/call').length, 1) |
| 120 | assert.equal(f.broker.closed, true); await f.module.dispose() |
| 121 | }) |
| 122 | test('401 cannot make SDK fetch replay a consumed operation', async () => { |
| 123 | const f = await fixture(); f.broker.status = 401 |
| 124 | await assert.rejects(f.module.request({ owner, session_id: f.open.session_id, grant: f.broker.grant('tools/list', {}), deadline_ms: 1000 }, f.signal)) |
| 125 | assert.equal(f.broker.frames.filter(frame => frame.method === 'tools/list').length, 1) |
| 126 | assert.equal(f.broker.closed, true); await f.module.dispose() |
| 127 | }) |
| 128 | test('plugin tier cannot redeem any FetchProxy frame', () => { |
| 129 | const rows = [['net/start', { ticket: 't' }], ['net/fetch', { url: 'https://mcp-proxy.invalid/s', method: 'GET', headers: {} }], ['net/read', { response_id: 'r' }], ['net/release', { response_id: 'r' }], ['net/close', {}]] |
| 130 | for (const [method, rest] of rows) { |
| 131 | const frame = { jsonrpc: '2.0', id: 1, method, params: { owner, session_id: 's', ...rest } } |
| 132 | assert.throws(() => validateMessage(frame, 'host_to_core', 'plugin'), /not allowed/) |
| 133 | validateMessage(frame, 'host_to_core', 'builtin') |
| 134 | } |
| 135 | }) |
| 136 | test('cancelled HTTP body releases its response and refuses queued work', async () => { |
| 137 | const f = await fixture(); f.broker.silent = true |
| 138 | const first = f.module.request({ owner, session_id: f.open.session_id, grant: f.broker.grant('tools/list', {}), deadline_ms: 20 }, f.signal) |
| 139 | const queued = f.module.request({ owner, session_id: f.open.session_id, grant: f.broker.grant('prompts/list', {}), deadline_ms: 1000 }, f.signal) |
| 140 | const firstResult = assert.rejects(first, /timed out|closed|cancelled/) |
| 141 | const queuedResult = assert.rejects(queued, /closed|cancelled|stale/) |
| 142 | await firstResult; await queuedResult |
| 143 | assert.equal(f.broker.frames.filter(frame => frame.method === 'tools/list').length, 1) |
| 144 | assert.equal(f.broker.frames.some(frame => frame.method === 'prompts/list'), false) |
| 145 | assert.equal(f.broker.closed, true); assert.equal(f.broker.bodies.size, 0) |
| 146 | await f.module.dispose() |
| 147 | }) |
| 148 | |
| 149 | // Actual official Streamable HTTP -> SSE transport change; the FetchProxy |
| 150 | // simulates one explicit Rust refusal and a fresh exact single-use grant. |
| 151 | async function negotiatedFixture(options = {}) { |
| 152 | const broker = new FakeFetchProxy(), module = createMcpModule(broker, owner) |
| 153 | Object.assign(broker, { negotiate: true }, options) |
| 154 | const open = { owner, session_id: broker.session, transport: 'http', launch_ticket: 'launch-once', initialize_grant: broker.grant('initialize', initParams), initialized_grant: broker.grant('notifications/initialized', {}), client_version: '0.10.1', deadline_ms: 100 } |
| 155 | return { broker, module, open } |
| 156 | } |
| 157 | test('official SDK negotiates legacy SSE only with a fresh exact Rust grant and original wire ID', async () => { |
| 158 | const f = await negotiatedFixture() |
| 159 | try { |
| 160 | await f.module.open(f.open, new AbortController().signal) |
| 161 | const attempts = f.broker.fetches.filter(p => p.frame?.method === 'initialize') |
| 162 | assert.equal(attempts.length, 2) |
| 163 | assert.equal(attempts[0].frame.id, f.open.initialize_grant.wire_id) |
| 164 | assert.deepEqual(attempts[0].frame, attempts[1].frame) |
| 165 | assert.notEqual(attempts[0].ticket, attempts[1].ticket) |
| 166 | assert.equal(attempts[0].operation_id, attempts[1].operation_id) |
| 167 | assert.equal(new URL(attempts[1].url).pathname.endsWith('/endpoint'), true) |
| 168 | assert.equal(f.broker.frames.filter(p => p.method === 'notifications/initialized').length, 1) |
| 169 | } finally { await f.module.dispose() } |
| 170 | }) |
| 171 | test('changed negotiation params refuse before SDK opens legacy channel or rewrites a request', async () => { |
| 172 | const f = await negotiatedFixture({ mutateNegotiation: true }) |
| 173 | await assert.rejects(f.module.open(f.open, new AbortController().signal), /mismatch|closed/) |
| 174 | assert.equal(f.broker.fetches.some(p => p.method === 'GET'), false) |
| 175 | assert.equal(f.broker.frames.filter(p => p.method === 'initialize').length, 1) |
| 176 | assert.equal(f.broker.closed, true); await f.module.dispose() |
| 177 | }) |
| 178 | test('cancellation during negotiated endpoint discovery releases body without replaying the admitted request', async () => { |
| 179 | const f = await negotiatedFixture({ holdEndpoint: true }) |
| 180 | await assert.rejects(f.module.open(f.open, new AbortController().signal), /closed|cancelled/) |
| 181 | assert.equal(f.broker.frames.filter(p => p.method === 'initialize').length, 1) |
| 182 | assert.equal(f.broker.closed, true); assert.equal(f.broker.bodies.size, 0) |
| 183 | await f.module.dispose() |
| 184 | }) |
| 185 | |
| 186 | test('correlated malformed catalog reply preserves SDK session for Rust per-item admission and the next page', async () => { |
| 187 | const f = await fixture() |
| 188 | f.broker.customResult = method => method === 'tools/list' ? { tools: [{ description: 'missing name', inputSchema: { type: 'object' } }, { name: 'valid', inputSchema: { type: 'object' } }], nextCursor: 'next' } : undefined |
| 189 | await assert.rejects(f.module.request({ owner, session_id: f.open.session_id, grant: f.broker.grant('tools/list', {}), deadline_ms: 1000 }, f.signal)) |
| 190 | assert.equal(f.broker.closed, false, 'the exact correlated reply is available to Rust; schema refusal is not an unknown write') |
| 191 | f.broker.customResult = undefined |
| 192 | const result = await f.module.request({ owner, session_id: f.open.session_id, grant: f.broker.grant('tools/list', { cursor: 'next' }), deadline_ms: 1000 }, f.signal) |
| 193 | assert.equal(result.nextCursor, 'next-page') |
| 194 | assert.equal(f.broker.frames.filter(p => p.method === 'tools/list').length, 2) |
| 195 | await f.module.dispose() |
| 196 | }) |
| 197 |