返回 CodeWhale
mcp-http.test.mjs
根目录 / crates / tui / extension-host / test / mcp-http.test.mjs
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
197 lines Plain Text