返回 CodeWhale
mcp.ts
1 /** Tier-0 SDK protocol owner. Rust owns launch, every operation grant,
2 * pipes, catalog admission, keys, credentials, cancellation and teardown. */
3 import { Client, SSEClientTransport, StreamableHTTPClientTransport, parseJSONRPCMessage, type JSONRPCMessage, type RequestMethod, type Transport } from '@modelcontextprotocol/client'
4 import type { OwnerRef } from '../protocol.generated.ts'
5
6 export interface BrokerRpc {
7 request<T>(method: string, params: unknown, signal?: AbortSignal): Promise<T>
8 }
9 export interface OperationGrant {
10 ticket: string
11 operation_id: string
12 method: string
13 wire_id?: string
14 params: Record<string, unknown>
15 }
16 export interface OpenParams {
17 owner: OwnerRef
18 session_id: string
19 launch_ticket: string
20 transport?: 'stdio' | 'http' | 'sse'
21 initialize_grant: OperationGrant
22 initialized_grant: OperationGrant
23 client_version: string
24 deadline_ms: number
25 }
26 export interface RequestParams {
27 owner: OwnerRef
28 session_id: string
29 grant: OperationGrant
30 deadline_ms: number
31 }
32 const MAX_FRAME_BYTES = 16 * 1024 * 1024
33 const MAX_SESSIONS = 64
34 const MAX_PENDING = 256
35 const MAX_DEPTH = 128
36 function requestMethod(value: string): RequestMethod {
37 switch (value) {
38 case 'tools/list': case 'resources/list': case 'resources/templates/list': case 'prompts/list':
39 case 'tools/call': case 'resources/read': case 'prompts/get': return value
40 default: throw new Error('MCP request method is unavailable')
41 }
42 }
43 const SUPPORTED = ['2025-06-18', '2025-11-25', '2025-03-26', '2024-11-05']
44 function identity(owner: OwnerRef): string {
45 return JSON.stringify([owner.plugin_id, owner.generation, owner.owner_token])
46 }
47 function jsonEqual(a: unknown, b: unknown, depth = 0): boolean {
48 if (depth > MAX_DEPTH) throw new Error('MCP operation nesting exceeds the bound')
49 if (a === b) return true
50 if (typeof a !== typeof b || a === null || b === null) return false
51 if (Array.isArray(a) || Array.isArray(b)) {
52 return Array.isArray(a) && Array.isArray(b) && a.length === b.length && a.every((value, i) => jsonEqual(value, b[i], depth + 1))
53 }
54 if (typeof a !== 'object') return false
55 const first = Object.keys(a as object), second = Object.keys(b as object)
56 if (first.length !== second.length) return false
57 return first.every(key => Object.hasOwn(b as object, key) && jsonEqual((a as Record<string, unknown>)[key], (b as Record<string, unknown>)[key], depth + 1))
58 }
59 function frameSize(value: unknown): number {
60 return Buffer.byteLength(JSON.stringify(value), 'utf8')
61 }
62 function deadline(value: number): number {
63 if (!Number.isSafeInteger(value) || value < 1 || value > 24 * 60 * 60 * 1000) throw new Error('invalid MCP request deadline')
64 return value
65 }
66
67 interface AdmittedFrame { frame: JSONRPCMessage; ticket?: string; operation_id?: string }
68 /** Real SDK IDs stay internal. Rust-issued string IDs and exact grants reach peers. */
69 class OperationFrames {
70 private readonly grants = new Map<string, OperationGrant>()
71 private readonly pending = new Map<string, { operation: string; wire: string; sdk: string | number }>()
72 private readonly wireToSdk = new Map<string, string>()
73 private readonly replied = new Set<string>()
74 grant(value: OperationGrant): void {
75 if (this.grants.size >= MAX_PENDING || this.grants.has(value.method)) throw new Error('MCP grant slot unavailable')
76 if (value.wire_id !== undefined && (!value.wire_id || value.wire_id.length > 256)) throw new Error('invalid Rust wire ID')
77 this.grants.set(value.method, structuredClone(value))
78 }
79 incoming(value: unknown): JSONRPCMessage {
80 const frame = parseJSONRPCMessage(value)
81 if ('id' in frame && !('method' in frame)) {
82 const sdkKey = this.wireToSdk.get(JSON.stringify(frame.id))
83 if (sdkKey === undefined) throw new Error('MCP reply has no admitted wire ID')
84 const binding = this.pending.get(sdkKey)!
85 this.pending.delete(sdkKey); this.wireToSdk.delete(JSON.stringify(frame.id))
86 if (this.replied.size >= MAX_PENDING) throw new Error('MCP completed exchange bound exceeded')
87 this.replied.add(binding.operation)
88 frame.id = binding.sdk
89 }
90 return frame
91 }
92 outgoing(value: JSONRPCMessage): AdmittedFrame {
93 const frame = structuredClone(value)
94 if (!('method' in frame)) {
95 const refusal = 'error' in frame && frame.error.code === -32601
96 const ping = 'result' in frame && Object.keys(frame.result).length === 0
97 if (!refusal && !ping) throw new Error('MCP server request requires Rust authority')
98 return { frame }
99 }
100 if (frame.method === 'notifications/cancelled') {
101 const params = frame.params as Record<string, unknown> | undefined
102 const binding = this.pending.get(JSON.stringify(params?.requestId))
103 if (!binding || !params) throw new Error('MCP cancellation has no admitted request')
104 params.requestId = binding.wire
105 return { frame, operation_id: binding.operation }
106 }
107 const grant = this.grants.get(frame.method)
108 if (!grant || !jsonEqual(frame.params ?? {}, grant.params)) throw new Error('MCP frame has no exact operation grant')
109 if ('id' in frame) {
110 if (grant.wire_id === undefined || this.pending.size >= MAX_PENDING || this.wireToSdk.has(JSON.stringify(grant.wire_id))) throw new Error('MCP wire ID is not admitted')
111 const sdk = frame.id, key = JSON.stringify(sdk)
112 if (this.pending.has(key)) throw new Error('duplicate SDK request ID')
113 this.pending.set(key, { operation: grant.operation_id, wire: grant.wire_id, sdk })
114 this.wireToSdk.set(JSON.stringify(grant.wire_id), key)
115 frame.id = grant.wire_id
116 } else if (grant.wire_id !== undefined) throw new Error('notification cannot redeem a request grant')
117 this.grants.delete(frame.method)
118 return { frame, ticket: grant.ticket, operation_id: grant.operation_id }
119 }
120 takeReply(operation: string): boolean { const found = this.replied.has(operation); this.replied.delete(operation); return found }
121 clear(): void { this.grants.clear(); this.pending.clear(); this.wireToSdk.clear(); this.replied.clear() }
122 }
123 interface GrantedTransport extends Transport { grant(value: OperationGrant): void; takeReply(operation: string): boolean }
124 class BrokerTransport implements GrantedTransport {
125 onclose?: () => void
126 onerror?: (error: Error) => void
127 onmessage?: Transport['onmessage']
128 private readonly stop = new AbortController()
129 private readonly frames = new OperationFrames()
130 private started = false
131 private closed = false
132 constructor(private readonly rpc: BrokerRpc, private readonly owner: OwnerRef, private readonly session: string, private readonly launch: string) {}
133 takeReply(operation: string): boolean { return this.frames.takeReply(operation) }
134 grant(value: OperationGrant): void { if (this.closed) throw new Error('MCP transport closed'); this.frames.grant(value) }
135 async start(): Promise<void> {
136 if (this.started || this.closed) throw new Error('MCP transport already started or closed')
137 this.started = true
138 await this.rpc.request('proc/launch', { owner: this.owner, session_id: this.session, ticket: this.launch }, this.stop.signal)
139 void this.readLoop().catch(async error => { if (!this.closed) this.onerror?.(error); await this.close().catch(() => undefined) })
140 }
141 private async readLoop(): Promise<void> {
142 while (!this.closed) {
143 const result = await this.rpc.request<{ frame?: unknown; closed?: boolean }>('proc/read', { owner: this.owner, session_id: this.session }, this.stop.signal)
144 if (result.closed) { this.finish(); return }
145 if (result.frame === undefined || frameSize(result.frame) > MAX_FRAME_BYTES) throw new Error('MCP broker frame exceeds the bound')
146 this.onmessage?.(this.frames.incoming(result.frame))
147 }
148 }
149 async send(message: JSONRPCMessage): Promise<void> {
150 if (this.closed || !this.started || frameSize(message) > MAX_FRAME_BYTES) throw new Error('MCP broker write unavailable')
151 const { frame, ticket, operation_id } = this.frames.outgoing(message)
152 try {
153 await this.rpc.request('proc/write', { owner: this.owner, session_id: this.session, frame, ...(ticket === undefined ? {} : { ticket }), ...(operation_id === undefined ? {} : { operation_id }) }, this.stop.signal)
154 } catch (error) { await this.close().catch(() => undefined); throw error }
155 }
156 private finish(): void { if (this.closed) return; this.closed = true; this.stop.abort(); this.frames.clear(); this.onclose?.() }
157 async close(): Promise<void> { if (this.closed) return; this.finish(); await this.rpc.request('proc/close', { owner: this.owner, session_id: this.session }) }
158 }
159 /** Official HTTP/SSE framing with an opaque Rust-authorized FetchProxy. */
160 class HttpBrokerTransport implements GrantedTransport {
161 onclose?: () => void
162 onerror?: (error: Error) => void
163 onmessage?: Transport['onmessage']
164 private readonly stop = new AbortController()
165 private readonly frames = new OperationFrames()
166 private readonly writes = new Map<string, AdmittedFrame>()
167 private inner: SSEClientTransport | StreamableHTTPClientTransport
168 private negotiation?: OperationGrant
169 private legacy: boolean
170 private started = false
171 private closed = false
172 constructor(private readonly rpc: BrokerRpc, private readonly owner: OwnerRef, private readonly session: string, private readonly launch: string, legacy: boolean) {
173 this.legacy = legacy
174 this.inner = this.createInner(legacy)
175 }
176 private createInner(legacy: boolean): SSEClientTransport | StreamableHTTPClientTransport {
177 const url = new URL(`https://mcp-proxy.invalid/${this.session}`)
178 const fetcher = this.fetch.bind(this)
179 const inner = legacy ? new SSEClientTransport(url, { fetch: fetcher }) : new StreamableHTTPClientTransport(url, { fetch: fetcher, onInsufficientScope: 'throw', reconnectionOptions: { maxRetries: 0, initialReconnectionDelay: 0, maxReconnectionDelay: 0, reconnectionDelayGrowFactor: 1 } })
180 inner.onmessage = value => { try { this.onmessage?.(this.frames.incoming(value)) } catch (error) { this.fail(error) } }
181 inner.onclose = () => { if (!this.closed) this.fail(new Error('MCP HTTP channel closed')) }
182 inner.onerror = error => { if (!this.negotiation) this.fail(error) }
183 return inner
184 }
185 private fail(error: unknown): void { if (!this.closed) this.onerror?.(error instanceof Error ? error : new Error('MCP HTTP failed')); void this.close().catch(() => undefined) }
186 takeReply(operation: string): boolean { return this.frames.takeReply(operation) }
187 grant(value: OperationGrant): void { if (this.closed) throw new Error('MCP transport closed'); this.frames.grant(value) }
188 async start(): Promise<void> {
189 if (this.started || this.closed) throw new Error('MCP HTTP already started or closed')
190 this.started = true
191 await this.rpc.request('net/start', { owner: this.owner, session_id: this.session, ticket: this.launch }, this.stop.signal)
192 await this.inner.start()
193 }
194 setProtocolVersion(version: string): void { this.inner.setProtocolVersion(version) }
195 async send(message: JSONRPCMessage, options?: Parameters<StreamableHTTPClientTransport['send']>[1]): Promise<void> {
196 if (!this.started || this.closed || frameSize(message) > MAX_FRAME_BYTES) throw new Error('MCP HTTP write unavailable')
197 const admitted = this.frames.outgoing(message), key = JSON.stringify(admitted.frame)
198 if (this.writes.size >= MAX_PENDING || this.writes.has(key)) throw new Error('MCP HTTP write slot unavailable')
199 this.writes.set(key, admitted)
200 try {
201 try { if (this.inner instanceof SSEClientTransport) await this.inner.send(admitted.frame); else await this.inner.send(admitted.frame, options) }
202 catch (error) {
203 const replacement = this.negotiation
204 if (!replacement || this.legacy || this.closed || this.stop.signal.aborted) throw error
205 this.negotiation = undefined
206 const original = admitted.frame
207 if (!('method' in original) || replacement.method !== original.method || replacement.operation_id !== admitted.operation_id
208 || replacement.wire_id !== ('id' in original ? original.id : undefined) || !jsonEqual(replacement.params, original.params ?? {})) throw new Error('MCP negotiation grant mismatch')
209 // Only Rust's explicit refusal issued this fresh, single-use grant.
210 // The original numeric SDK request remains pending; peer IDs stay exact.
211 this.inner.onclose = undefined; this.inner.onerror = undefined; this.inner.onmessage = undefined
212 await this.inner.close()
213 this.legacy = true
214 this.inner = this.createInner(true)
215 await this.inner.start()
216 this.writes.set(key, { frame: original, ticket: replacement.ticket, operation_id: replacement.operation_id })
217 await this.inner.send(original)
218 }
219 } catch (error) { await this.close().catch(() => undefined); throw error }
220 finally { this.writes.delete(key) }
221 }
222 private async fetch(input: string | URL | Request, init?: RequestInit): Promise<Response> {
223 if (this.closed || this.stop.signal.aborted) throw new Error('MCP HTTP session closed')
224 const url = input instanceof Request ? input.url : String(input)
225 const method = init?.method ?? 'GET'
226 const signal = init?.signal ? AbortSignal.any([this.stop.signal, init.signal]) : this.stop.signal
227 let admitted: AdmittedFrame | undefined
228 if (method === 'POST') {
229 if (typeof init?.body !== 'string' || Buffer.byteLength(init.body) > MAX_FRAME_BYTES) throw new Error('MCP HTTP body unavailable')
230 admitted = this.writes.get(JSON.stringify(JSON.parse(init.body)))
231 if (!admitted) throw new Error('MCP HTTP body has no exact admitted frame')
232 this.writes.delete(JSON.stringify(admitted.frame)) // one fetch, no SDK replay
233 }
234 const headers: Record<string, string> = {}
235 for (const [name, value] of new Headers(init?.headers)) {
236 if (!['accept', 'content-type', 'mcp-session-id', 'mcp-protocol-version'].includes(name) || value.length > 8192) throw new Error('MCP framing header unavailable')
237 headers[name.replaceAll('-', '_')] = value
238 }
239 const reply = await this.rpc.request<{ response_id?: string; status: number; headers: Record<string, string>; legacy_grant?: OperationGrant }>('net/fetch', { owner: this.owner, session_id: this.session, url, method, headers, ...(admitted ?? {}) }, signal)
240 if (reply.legacy_grant !== undefined) {
241 if (!admitted || this.legacy || this.negotiation || reply.response_id !== undefined
242 || ![404, 405, 406, 415, 501].includes(reply.status)
243 || typeof reply.legacy_grant.ticket !== 'string' || !reply.legacy_grant.ticket) throw new Error('MCP negotiation unavailable')
244 this.negotiation = reply.legacy_grant
245 throw new Error('Rust admitted legacy MCP negotiation')
246 }
247 if (reply.response_id === undefined) return new Response(null, { status: reply.status, headers: reply.headers })
248 const response_id = reply.response_id
249 let done = false
250 const cancel = async () => { if (done) return; done = true; await this.rpc.request('net/release', { owner: this.owner, session_id: this.session, response_id }).catch(() => undefined) }
251 const body = new ReadableStream<Uint8Array>({
252 pull: async controller => { try {
253 if (signal.aborted) throw new Error('MCP HTTP body cancelled')
254 const next = await this.rpc.request<{ data: number[]; done: boolean }>('net/read', { owner: this.owner, session_id: this.session, response_id }, signal)
255 if (!Array.isArray(next.data) || next.data.length > 32768 || next.data.some(b => !Number.isInteger(b) || b < 0 || b > 255)) throw new Error('MCP HTTP chunk exceeds bound')
256 if (next.data.length) controller.enqueue(Uint8Array.from(next.data))
257 if (next.done) { await cancel(); controller.close() }
258 } catch (error) { await cancel(); controller.error(error); this.fail(error) } },
259 cancel,
260 }, { highWaterMark: 0 })
261 signal.addEventListener('abort', () => { void cancel() }, { once: true })
262 return new Response(body, { status: reply.status, headers: reply.headers })
263 }
264 async close(): Promise<void> {
265 if (this.closed) return
266 this.closed = true; this.stop.abort(); this.frames.clear(); this.writes.clear()
267 await this.inner.close().catch(() => undefined)
268 await this.rpc.request('net/close', { owner: this.owner, session_id: this.session }).catch(() => undefined)
269 this.onclose?.()
270 }
271 }
272 interface Session {
273 client: Client
274 transport: GrantedTransport
275 tail: Promise<unknown>
276 }
277 export class McpModule {
278 private readonly sessions = new Map<string, Session>()
279 private disposed = false
280 constructor(private readonly rpc: BrokerRpc, private readonly owner: OwnerRef) {}
281 private check(owner: OwnerRef): void {
282 if (this.disposed || identity(owner) !== identity(this.owner)) throw new Error('MCP builtin owner is no longer live')
283 }
284 async open(params: OpenParams, signal: AbortSignal): Promise<unknown> {
285 this.check(params.owner)
286 const timeout = deadline(params.deadline_ms)
287 if (params.transport !== undefined && !['stdio', 'http', 'sse'].includes(params.transport)) throw new Error('MCP transport selector unavailable')
288 if (this.sessions.size >= MAX_SESSIONS || this.sessions.has(params.session_id)) throw new Error('MCP session slot unavailable')
289 const transport: GrantedTransport = params.transport === undefined || params.transport === 'stdio'
290 ? new BrokerTransport(this.rpc, this.owner, params.session_id, params.launch_ticket)
291 : new HttpBrokerTransport(this.rpc, this.owner, params.session_id, params.launch_ticket, params.transport === 'sse')
292 const client = new Client({ name: 'codewhale-tui', version: params.client_version }, { capabilities: {}, supportedProtocolVersions: SUPPORTED, versionNegotiation: { mode: 'legacy' }, listMaxPages: 64 })
293 const session: Session = { client, transport, tail: Promise.resolve() }
294 this.sessions.set(params.session_id, session)
295 transport.grant(params.initialize_grant)
296 transport.grant(params.initialized_grant)
297 const abort = () => { void transport.close().catch(() => undefined) }
298 signal.addEventListener('abort', abort, { once: true })
299 const timer = setTimeout(abort, timeout)
300 timer.unref()
301 try {
302 if (signal.aborted) throw new Error('MCP open cancelled')
303 await client.connect(transport)
304 transport.takeReply(params.initialize_grant.operation_id)
305 this.check(params.owner)
306 return { protocolVersion: client.getNegotiatedProtocolVersion(), serverInfo: client.getServerVersion(), capabilities: client.getServerCapabilities(), instructions: client.getInstructions() }
307 } catch (error) {
308 this.sessions.delete(params.session_id)
309 await transport.close().catch(() => undefined)
310 throw error
311 } finally { clearTimeout(timer); signal.removeEventListener('abort', abort) }
312 }
313 async request(params: RequestParams, signal: AbortSignal): Promise<unknown> {
314 this.check(params.owner)
315 const timeout = deadline(params.deadline_ms)
316 const method = requestMethod(params.grant.method)
317 const session = this.sessions.get(params.session_id)
318 if (!session) throw new Error('MCP session is no longer live')
319 const run = session.tail.then(async () => {
320 this.check(params.owner)
321 if (this.sessions.get(params.session_id) !== session || signal.aborted) throw new Error('MCP operation cancelled or stale')
322 session.transport.grant(params.grant)
323 // Generic request preserves per-page cursors. listTools/listResources /
324 // listPrompts convenience methods would auto-aggregate into an SDK cache.
325 try {
326 return await session.client.request({ method, params: params.grant.params }, { signal, timeout, resetTimeoutOnProgress: false })
327 } catch (error) {
328 // A correlated reply completed the exchange, even if the SDK's
329 // stricter result schema rejected a legacy malformed catalog item.
330 // Rust receives the original bytes and applies its shared admission.
331 // Every failure without that receipt retires before queued work runs.
332 if (!session.transport.takeReply(params.grant.operation_id)) {
333 this.sessions.delete(params.session_id)
334 await session.transport.close().catch(() => undefined)
335 }
336 throw error
337 } finally { session.transport.takeReply(params.grant.operation_id) }
338 })
339 session.tail = run.catch(() => undefined)
340 return run
341 }
342 async close(owner: OwnerRef, sessionId: string): Promise<void> {
343 this.check(owner)
344 const session = this.sessions.get(sessionId)
345 this.sessions.delete(sessionId)
346 await session?.client.close()
347 }
348 async dispose(): Promise<void> {
349 if (this.disposed) return
350 this.disposed = true
351 const sessions = [...this.sessions.values()]
352 this.sessions.clear()
353 await Promise.allSettled(sessions.map(session => session.client.close()))
354 }
355 }
356 export function createMcpModule(rpc: BrokerRpc, owner: OwnerRef): McpModule {
357 if (owner.plugin_id !== 'host:mcp') throw new Error('MCP protocol owner must be the pinned builtin')
358 return new McpModule(rpc, Object.freeze(structuredClone(owner)))
359 }
360
360 lines TYPESCRIPT