|
| 1 | +/** |
| 2 | + * Graceful close — opt-in drain of in-flight requests (issue #1231). |
| 3 | + * |
| 4 | + * `close({ drainPendingRequests })` waits for in-flight requests to settle |
| 5 | + * before the transport closes, so a completed-but-still-reading HTTP response |
| 6 | + * is not torn down by the transport's abort (which OpenTelemetry's undici |
| 7 | + * instrumentation reports as UND_ERR_ABORTED on a 200 OK). Default close |
| 8 | + * behavior is unchanged. |
| 9 | + */ |
| 10 | +import type { JSONRPCMessage } from '@modelcontextprotocol/core-internal'; |
| 11 | +import { InMemoryTransport } from '@modelcontextprotocol/core-internal'; |
| 12 | +import { describe, expect, it } from 'vitest'; |
| 13 | + |
| 14 | +import { Client } from '../../src/client/client'; |
| 15 | + |
| 16 | +const flush = () => new Promise(r => setTimeout(r, 10)); |
| 17 | + |
| 18 | +type ScriptedServer = { |
| 19 | + clientTx: InMemoryTransport; |
| 20 | + serverTx: InMemoryTransport; |
| 21 | + written: JSONRPCMessage[]; |
| 22 | + /** Replies to the oldest outstanding non-initialize request. */ |
| 23 | + reply: (result: Record<string, unknown>) => void; |
| 24 | + /** Replies to the oldest outstanding non-initialize request on a delay. */ |
| 25 | + replyAfter: (ms: number, result: Record<string, unknown>) => Promise<void>; |
| 26 | +}; |
| 27 | + |
| 28 | +/** |
| 29 | + * A linked in-memory pair where the server auto-answers the legacy |
| 30 | + * `initialize` handshake (so `connect()` resolves) but holds every other |
| 31 | + * request until the test calls `reply()` / `replyAfter()`. |
| 32 | + */ |
| 33 | +async function scriptedLegacyServer(): Promise<ScriptedServer> { |
| 34 | + const [clientTx, serverTx] = InMemoryTransport.createLinkedPair(); |
| 35 | + const written: JSONRPCMessage[] = []; |
| 36 | + const pendingIds: (number | string)[] = []; |
| 37 | + serverTx.onmessage = message => { |
| 38 | + written.push(message); |
| 39 | + const req = message as { id?: number | string; method?: string; params?: { protocolVersion?: string } }; |
| 40 | + if (req.method === 'initialize' && req.id !== undefined) { |
| 41 | + void serverTx.send({ |
| 42 | + jsonrpc: '2.0', |
| 43 | + id: req.id, |
| 44 | + result: { |
| 45 | + protocolVersion: req.params?.protocolVersion ?? '2025-06-18', |
| 46 | + capabilities: {}, |
| 47 | + serverInfo: { name: 'scripted', version: '1' } |
| 48 | + } |
| 49 | + }); |
| 50 | + return; |
| 51 | + } |
| 52 | + if (req.method === 'notifications/initialized') { |
| 53 | + return; |
| 54 | + } |
| 55 | + if (req.id !== undefined) { |
| 56 | + pendingIds.push(req.id); |
| 57 | + } |
| 58 | + }; |
| 59 | + await serverTx.start(); |
| 60 | + const reply = (result: Record<string, unknown>) => { |
| 61 | + const id = pendingIds.shift(); |
| 62 | + if (id === undefined) { |
| 63 | + throw new Error('no pending request to reply to'); |
| 64 | + } |
| 65 | + void serverTx.send({ jsonrpc: '2.0', id, result }); |
| 66 | + }; |
| 67 | + return { |
| 68 | + clientTx, |
| 69 | + serverTx, |
| 70 | + written, |
| 71 | + reply, |
| 72 | + replyAfter: async (ms: number, result: Record<string, unknown>) => { |
| 73 | + await new Promise(r => setTimeout(r, ms)); |
| 74 | + reply(result); |
| 75 | + } |
| 76 | + }; |
| 77 | +} |
| 78 | + |
| 79 | +async function connectClient(options?: ConstructorParameters<typeof Client>[1]): Promise<{ client: Client; server: ScriptedServer }> { |
| 80 | + const server = await scriptedLegacyServer(); |
| 81 | + const client = new Client({ name: 'test-client', version: '1.0.0' }, options); |
| 82 | + await client.connect(server.clientTx); |
| 83 | + return { client, server }; |
| 84 | +} |
| 85 | + |
| 86 | +/** Spies on the client transport's close() without changing behavior. */ |
| 87 | +function spyTransportClose(client: Client): { closed: () => boolean } { |
| 88 | + let closed = false; |
| 89 | + const transport = client.transport!; |
| 90 | + const originalClose = transport.close.bind(transport); |
| 91 | + transport.close = async () => { |
| 92 | + closed = true; |
| 93 | + await originalClose(); |
| 94 | + }; |
| 95 | + return { closed: () => closed }; |
| 96 | +} |
| 97 | + |
| 98 | +describe('Client.close graceful drain', () => { |
| 99 | + it('default close() is unchanged: transport closes with a request in flight', async () => { |
| 100 | + const { client } = await connectClient(); |
| 101 | + const inFlight = client.request({ method: 'ping' }).catch(e => e); |
| 102 | + await flush(); |
| 103 | + await client.close(); |
| 104 | + const settled = (await inFlight) as Error; |
| 105 | + // The request is settled by the close itself, not by a response. |
| 106 | + expect(settled).toBeInstanceOf(Error); |
| 107 | + expect((settled as Error).message).toMatch(/closed/i); |
| 108 | + }); |
| 109 | + |
| 110 | + it('close({ drainPendingRequests: true }) waits for the in-flight response before closing', async () => { |
| 111 | + const { client, server } = await connectClient(); |
| 112 | + let settled: unknown; |
| 113 | + const inFlight = client |
| 114 | + .request({ method: 'ping' }) |
| 115 | + .then(r => (settled = r)) |
| 116 | + .catch(e => (settled = e)); |
| 117 | + await flush(); |
| 118 | + |
| 119 | + const spy = spyTransportClose(client); |
| 120 | + const closing = client.close({ drainPendingRequests: true }); |
| 121 | + |
| 122 | + // The transport must still be open while the request is outstanding. |
| 123 | + await flush(); |
| 124 | + expect(spy.closed()).toBe(false); |
| 125 | + |
| 126 | + // The response lands on the still-open connection; the drain then |
| 127 | + // completes and the transport closes. |
| 128 | + await server.replyAfter(20, {}); |
| 129 | + await inFlight; |
| 130 | + await closing; |
| 131 | + expect(spy.closed()).toBe(true); |
| 132 | + expect(settled).toBeDefined(); |
| 133 | + }); |
| 134 | + |
| 135 | + it('multiple in-flight requests all drain before the transport closes', async () => { |
| 136 | + const { client, server } = await connectClient(); |
| 137 | + const first = client.request({ method: 'ping' }).catch(e => e); |
| 138 | + const second = client.request({ method: 'ping' }).catch(e => e); |
| 139 | + await flush(); |
| 140 | + |
| 141 | + const spy = spyTransportClose(client); |
| 142 | + const closing = client.close({ drainPendingRequests: true }); |
| 143 | + await flush(); |
| 144 | + expect(spy.closed()).toBe(false); |
| 145 | + |
| 146 | + server.reply({}); |
| 147 | + await first; |
| 148 | + await flush(); |
| 149 | + // One of two requests still outstanding: no close yet. |
| 150 | + expect(spy.closed()).toBe(false); |
| 151 | + |
| 152 | + server.reply({}); |
| 153 | + await second; |
| 154 | + await closing; |
| 155 | + expect(spy.closed()).toBe(true); |
| 156 | + }); |
| 157 | + |
| 158 | + it('falls back to a hard close after the drain timeout and requests settle with the close', async () => { |
| 159 | + const { client } = await connectClient(); |
| 160 | + const inFlight = client.request({ method: 'ping' }).catch(e => e); |
| 161 | + await flush(); |
| 162 | + |
| 163 | + await client.close({ drainPendingRequests: { timeoutMs: 30 } }); |
| 164 | + const settled = (await inFlight) as Error; |
| 165 | + expect(settled).toBeInstanceOf(Error); |
| 166 | + expect((settled as Error).message).toMatch(/closed/i); |
| 167 | + }); |
| 168 | + |
| 169 | + it('drain resolves immediately when nothing is in flight', async () => { |
| 170 | + const { client } = await connectClient(); |
| 171 | + const start = Date.now(); |
| 172 | + await client.close({ drainPendingRequests: true }); |
| 173 | + expect(Date.now() - start).toBeLessThan(500); |
| 174 | + }); |
| 175 | + |
| 176 | + it('ClientOptions.gracefulClose applies to parameterless close()', async () => { |
| 177 | + const { client, server } = await connectClient({ gracefulClose: true }); |
| 178 | + const inFlight = client.request({ method: 'ping' }).catch(e => e); |
| 179 | + await flush(); |
| 180 | + |
| 181 | + const spy = spyTransportClose(client); |
| 182 | + const closing = client.close(); |
| 183 | + await flush(); |
| 184 | + expect(spy.closed()).toBe(false); |
| 185 | + |
| 186 | + await server.replyAfter(20, {}); |
| 187 | + await inFlight; |
| 188 | + await closing; |
| 189 | + expect(spy.closed()).toBe(true); |
| 190 | + }); |
| 191 | + |
| 192 | + it('an explicit close({ drainPendingRequests: false }) overrides the constructor default', async () => { |
| 193 | + const { client } = await connectClient({ gracefulClose: true }); |
| 194 | + const inFlight = client.request({ method: 'ping' }).catch(e => e); |
| 195 | + await flush(); |
| 196 | + |
| 197 | + const spy = spyTransportClose(client); |
| 198 | + await client.close({ drainPendingRequests: false }); |
| 199 | + expect(spy.closed()).toBe(true); |
| 200 | + const settled = (await inFlight) as Error; |
| 201 | + expect(settled).toBeInstanceOf(Error); |
| 202 | + expect((settled as Error).message).toMatch(/closed/i); |
| 203 | + }); |
| 204 | +}); |
0 commit comments