Skip to content

Commit c3ec4a8

Browse files
fix(client): settle() drops _listenState by the captured id, not via parked
When a server-cancel (or any termination) is delivered synchronously inside _parkRequest's send — in-process transports — settle() runs before `parked` is assigned, so the previous `_listenState.delete(parked.messageId)` was skipped and the entry registered by onBeforeSend leaked. The catch-up after _parkRequest only called parked.unpark(), not the delete. Capture the messageId in the onBeforeSend closure as `listenMessageId`; settle() and wireTeardown() key off that, so the entry is dropped on every exit path regardless of whether `parked` has been assigned yet.
1 parent 9a8bf5b commit c3ec4a8

2 files changed

Lines changed: 54 additions & 6 deletions

File tree

packages/client/src/client/client.ts

Lines changed: 15 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1274,6 +1274,12 @@ export class Client extends Protocol<ClientContext> {
12741274
// before-ack / close-before-ack hangs are impossible by construction.
12751275
let state: 'opening' | 'open' | 'closed' = 'opening';
12761276
let parked: ReturnType<typeof this._parkRequest> | undefined;
1277+
// The listen request id, captured by `onBeforeSend` BEFORE the request
1278+
// goes out. `settle()` deletes `_listenState` by this id (not via
1279+
// `parked.messageId`), so a synchronously-delivered termination during
1280+
// `_parkRequest`'s send — when `parked` is still unassigned — does not
1281+
// leak the entry.
1282+
let listenMessageId: number | undefined;
12771283
let ackTimer: ReturnType<typeof setTimeout> | undefined;
12781284
let resolveOpening!: (honored: SubscriptionFilter) => void;
12791285
let rejectOpening!: (error: Error) => void;
@@ -1297,10 +1303,10 @@ export class Client extends Protocol<ClientContext> {
12971303
return;
12981304
}
12991305
state = 'closed';
1300-
if (parked !== undefined) {
1301-
this._listenState.delete(parked.messageId);
1302-
parked.unpark();
1306+
if (listenMessageId !== undefined) {
1307+
this._listenState.delete(listenMessageId);
13031308
}
1309+
parked?.unpark();
13041310
if (wasOpening) {
13051311
rejectOpening(
13061312
outcome === 'closed'
@@ -1321,7 +1327,7 @@ export class Client extends Protocol<ClientContext> {
13211327
// cancelled).
13221328
const wireTeardown = async (): Promise<void> => {
13231329
requestAbort.abort();
1324-
const id = parked?.messageId;
1330+
const id = listenMessageId;
13251331
if (id !== undefined) {
13261332
await this.transport
13271333
?.send({ jsonrpc: '2.0', method: 'notifications/cancelled', params: { requestId: id } })
@@ -1359,6 +1365,7 @@ export class Client extends Protocol<ClientContext> {
13591365
},
13601366
{ requestSignal: requestAbort.signal },
13611367
messageId => {
1368+
listenMessageId = messageId;
13621369
this._listenState.set(messageId, {
13631370
onAck: honored => settle({ ack: honored }),
13641371
onServerCancel: () => {
@@ -1369,8 +1376,10 @@ export class Client extends Protocol<ClientContext> {
13691376
);
13701377
// A synchronously-delivered termination during `send()` (an
13711378
// in-process transport) ran `settle()` before `parked` was
1372-
// assigned — unpark now so the handler does not leak. (Cast: TS
1373-
// control-flow narrowing does not track closure mutation.)
1379+
// assigned; `settle()` already cleared `_listenState` via
1380+
// `listenMessageId` — unpark now so the response handler does not
1381+
// leak either. (Cast: TS control-flow narrowing does not track
1382+
// closure mutation.)
13741383
if ((state as 'opening' | 'open' | 'closed') === 'closed') parked.unpark();
13751384
// Pre-ack capacity / params rejection arrives as a JSON-RPC error
13761385
// for the listen id; transport close is delivered the same way.

packages/client/test/client/listen.test.ts

Lines changed: 39 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -228,6 +228,42 @@ describe('Client.listen()', () => {
228228
await client.close();
229229
});
230230

231+
it('a synchronously-delivered server-cancel during send does not leak a _listenState entry', async () => {
232+
// In-process delivery: the server's notifications/cancelled arrives
233+
// inside `_parkRequest`'s send (before `parked` is assigned). settle()
234+
// must still drop the `_listenState` entry it registered via
235+
// onBeforeSend.
236+
const [clientTx, serverTx] = InMemoryTransport.createLinkedPair();
237+
serverTx.onmessage = m => {
238+
const req = m as { id?: number | string; method?: string };
239+
if (req.method === 'server/discover' && req.id !== undefined) {
240+
void serverTx.send({
241+
jsonrpc: '2.0',
242+
id: req.id,
243+
result: {
244+
resultType: 'complete',
245+
supportedVersions: [MODERN],
246+
capabilities: {},
247+
serverInfo: { name: 's', version: '1' }
248+
}
249+
});
250+
}
251+
if (req.method === 'subscriptions/listen' && req.id !== undefined) {
252+
void serverTx.send({ jsonrpc: '2.0', method: 'notifications/cancelled', params: { requestId: req.id } });
253+
}
254+
};
255+
await serverTx.start();
256+
const client = new Client({ name: 'c', version: '1' }, { versionNegotiation: { mode: 'auto' } });
257+
await client.connect(clientTx);
258+
const listenState = (client as unknown as { _listenState: Map<unknown, unknown> })._listenState;
259+
const before = listenState.size;
260+
const error = await client.listen({ toolsListChanged: true }).catch(e => e as Error);
261+
expect((error as Error).message).toContain('server cancelled before acknowledging');
262+
// No leaked _listenState entry for the listen id.
263+
expect(listenState.size).toBe(before);
264+
await client.close();
265+
});
266+
231267
it('a synchronous transport.send throw does not leak a _responseHandlers entry', async () => {
232268
const { clientTx } = await scriptedModern();
233269
const client = new Client({ name: 'c', version: '1' }, { versionNegotiation: { mode: 'auto' } });
@@ -242,6 +278,9 @@ describe('Client.listen()', () => {
242278
expect((error as Error).message).toContain('send blew up');
243279
// The park primitive unregistered before rethrowing — no leak.
244280
expect(handlers.size).toBe(before);
281+
// settle() in the catch path also dropped the _listenState entry that
282+
// onBeforeSend registered before send threw.
283+
expect((client as unknown as { _listenState: Map<unknown, unknown> })._listenState.size).toBe(0);
245284
clientTx.send = realSend;
246285
await client.close();
247286
});

0 commit comments

Comments
 (0)