-
Notifications
You must be signed in to change notification settings - Fork 2.2k
fix(core): await the notification send so a failed send is never briefly unhandled #2885
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Merged
Merged
Changes from all commits
Commits
Show all changes
3 commits
Select commit
Hold shift + click to select a range
File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,6 @@ | ||
| --- | ||
| '@modelcontextprotocol/client': patch | ||
| '@modelcontextprotocol/server': patch | ||
| --- | ||
|
|
||
| Sending a notification on a closed connection no longer produces a briefly unhandled promise rejection (seen as `unhandledrejection` on Cloudflare Workers) in addition to the returned rejection. |
82 changes: 82 additions & 0 deletions
82
packages/client/test/client/legacyHandshakeCloseAfterInitialize.test.ts
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,82 @@ | ||
| /** | ||
| * The trigger reported in #2864: the transport closes during the legacy | ||
| * `initialize` handshake, right after the server's result is delivered, so the | ||
| * client's `notifications/initialized` is sent on a connection that is already | ||
| * gone. `connect()` must reject through its returned promise with the | ||
| * not-connected send failure, and nothing else may escape. | ||
| * | ||
| * Node's rejection tracker does not report the one-microtask handler gap that | ||
| * workerd does (see core-internal's notificationSendRejection test for the | ||
| * timing observation itself); this end-to-end test pins the trigger path and | ||
| * the error `connect()` rejects with. | ||
| */ | ||
| import type { JSONRPCMessage, Transport } from '@modelcontextprotocol/core-internal'; | ||
| import { isJSONRPCRequest, SdkError, SdkErrorCode } from '@modelcontextprotocol/core-internal'; | ||
| import { afterEach, beforeEach, describe, expect, test } from 'vitest'; | ||
|
|
||
| import { Client } from '../../src/client/client'; | ||
|
|
||
| /** Answers `initialize`, then closes in the same tick — before the client can send `notifications/initialized`. */ | ||
| class ReplyThenCloseTransport implements Transport { | ||
| onclose?: () => void; | ||
| onerror?: (error: Error) => void; | ||
| onmessage?: (message: JSONRPCMessage) => void; | ||
|
|
||
| sent: JSONRPCMessage[] = []; | ||
| closeCalls = 0; | ||
|
|
||
| async start(): Promise<void> {} | ||
|
|
||
| async send(message: JSONRPCMessage): Promise<void> { | ||
| this.sent.push(message); | ||
| if (!isJSONRPCRequest(message) || message.method !== 'initialize') return; | ||
| queueMicrotask(() => { | ||
| this.onmessage?.({ | ||
| jsonrpc: '2.0', | ||
| id: message.id, | ||
| result: { protocolVersion: '2025-03-26', capabilities: {}, serverInfo: { name: 'flaky', version: '0' } } | ||
| }); | ||
| this.onclose?.(); | ||
| }); | ||
| } | ||
|
|
||
| async close(): Promise<void> { | ||
| this.closeCalls++; | ||
| } | ||
| } | ||
|
|
||
| describe('legacy handshake: transport closes between the initialize result and notifications/initialized', () => { | ||
| const unhandled: unknown[] = []; | ||
| const onUnhandled = (reason: unknown): void => { | ||
| unhandled.push(reason); | ||
| }; | ||
|
|
||
| beforeEach(() => { | ||
| unhandled.length = 0; | ||
| process.on('unhandledRejection', onUnhandled); | ||
| }); | ||
|
|
||
| afterEach(() => { | ||
| process.off('unhandledRejection', onUnhandled); | ||
| }); | ||
|
|
||
| test('connect() rejects with SdkError NotConnected from the initialized-notification send and no rejection escapes', async () => { | ||
| const transport = new ReplyThenCloseTransport(); | ||
| const client = new Client({ name: 'c', version: '0' }, { versionNegotiation: { mode: 'legacy' } }); | ||
|
|
||
| const rejection = await client.connect(transport).then( | ||
| () => undefined, | ||
| (error: unknown) => error | ||
| ); | ||
|
|
||
| expect(rejection).toBeInstanceOf(SdkError); | ||
| expect((rejection as SdkError).code).toBe(SdkErrorCode.NotConnected); | ||
| // The handshake got as far as the send that failed: initialize went out, | ||
| // the initialized notification never did. | ||
| expect(transport.sent.map(m => ('method' in m ? m.method : 'response'))).toEqual(['initialize']); | ||
|
|
||
| // Let any stray rejection surface before asserting none did. | ||
| await new Promise<void>(resolve => setTimeout(resolve, 0)); | ||
| expect(unhandled).toEqual([]); | ||
| }); | ||
| }); |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
82 changes: 82 additions & 0 deletions
82
packages/core-internal/test/shared/notificationSendRejection.test.ts
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,82 @@ | ||
| /** | ||
| * `Protocol.notification()` must hand its caller the ONLY rejection a failed | ||
| * send produces. The send funnel (`_notificationViaCodec`) is an async method | ||
| * that throws synchronously when there is no transport, so the promise it | ||
| * returns is already rejected. If `notification()` returns that promise | ||
| * instead of awaiting it, the async function resolves with a thenable: its | ||
| * `then` is read synchronously at return, but the call is deferred to the | ||
| * thenable job, so the inner rejection sits with no handler for one | ||
| * microtask. Node's tracker forgives that; | ||
| * workerd (Cloudflare Workers) reports it as `unhandledrejection` followed by | ||
| * `rejectionhandled`, which surfaces as noise in Vitest runs on that platform | ||
| * (#2864). | ||
| * | ||
| * The observation below is the handler-attachment timing itself: with | ||
| * `return await`, `await` attaches its reaction synchronously through the | ||
| * internal promise path and the inner promise's own `then` property is never | ||
| * read; with a bare `return`, `then` is read at return and called one | ||
| * microtask later. | ||
| */ | ||
| import { describe, expect, test } from 'vitest'; | ||
|
|
||
| import { SdkError, SdkErrorCode } from '../../src/errors/sdkErrors'; | ||
| import type { BaseContext } from '../../src/shared/protocol'; | ||
| import { Protocol } from '../../src/shared/protocol'; | ||
|
|
||
| class TestProtocolImpl extends Protocol<BaseContext> { | ||
| protected assertCapabilityForMethod(): void {} | ||
| protected assertNotificationCapability(): void {} | ||
| protected assertRequestHandlerCapability(): void {} | ||
| protected buildContext(ctx: BaseContext): BaseContext { | ||
| return ctx; | ||
| } | ||
| } | ||
|
|
||
| /** | ||
| * A natively rejected promise whose `then` property records every read. | ||
| * Resolving an async function with a thenable reaches `then` via property | ||
| * lookup; `await` on a native promise does not. | ||
| */ | ||
| function instrumentedRejection(error: Error): { promise: Promise<void>; thenReads: () => number } { | ||
| const promise = Promise.reject<void>(error); | ||
| const originalThen = promise.then.bind(promise); | ||
| let reads = 0; | ||
| Object.defineProperty(promise, 'then', { | ||
| get() { | ||
| reads++; | ||
| return originalThen; | ||
| } | ||
| }); | ||
| return { promise, thenReads: () => reads }; | ||
| } | ||
|
|
||
| describe('Protocol.notification(): a failed send rejects only through the returned promise', () => { | ||
| test('when not connected, the inner send rejection is handled in the same microtask (no thenable-job hop)', async () => { | ||
| const protocol = new TestProtocolImpl(); | ||
| const notConnected = new SdkError(SdkErrorCode.NotConnected, 'Not connected'); | ||
| const inner = instrumentedRejection(notConnected); | ||
|
|
||
| // Stand in for the real funnel with a rejection we can observe. The real | ||
| // one throws synchronously on `!this._transport`, i.e. it also returns an | ||
| // already-rejected promise — the shape that matters here. | ||
| // Plain instance override rather than `vi.spyOn`: the spy wrapper itself | ||
| // reads `then` on any promise a spied call returns, which would mask the | ||
| // observation below. | ||
| (protocol as unknown as { _notificationViaCodec: () => Promise<void> })._notificationViaCodec = () => inner.promise; | ||
|
|
||
| await expect(protocol.notification({ method: 'notifications/initialized' })).rejects.toBe(notConnected); | ||
|
|
||
| // A bare `return innerPromise` from the async method reads `then` | ||
| // synchronously at return and calls it one microtask later — the window | ||
| // in which workerd reports the inner rejection unhandled. | ||
| expect(inner.thenReads()).toBe(0); | ||
| }); | ||
|
|
||
| test('when not connected, the returned promise still rejects with SdkError NotConnected (unstubbed path)', async () => { | ||
| const protocol = new TestProtocolImpl(); | ||
|
|
||
| await expect(protocol.notification({ method: 'notifications/initialized' })).rejects.toSatisfy( | ||
| (error: unknown) => error instanceof SdkError && error.code === SdkErrorCode.NotConnected | ||
| ); | ||
| }); | ||
| }); | ||
52 changes: 52 additions & 0 deletions
52
packages/server/test/server/notificationSendRejection.test.ts
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,52 @@ | ||
| /** | ||
| * Server-side twin of core-internal's notificationSendRejection test: a | ||
| * `sendLoggingMessage()` on a server with no transport must reject ONLY | ||
| * through the returned promise. The send funnel returns an already-rejected | ||
| * promise; if `Protocol.notification()` returned it instead of awaiting it, | ||
| * the inner rejection would sit unhandled for one microtask (reported by | ||
| * workerd as `unhandledrejection` + `rejectionhandled`, #2864). | ||
| */ | ||
| import { SdkError, SdkErrorCode } from '@modelcontextprotocol/core-internal'; | ||
| import { describe, expect, test } from 'vitest'; | ||
|
|
||
| import { Server } from '../../src/server/server'; | ||
|
|
||
| function instrumentedRejection(error: Error): { promise: Promise<void>; thenReads: () => number } { | ||
| const promise = Promise.reject<void>(error); | ||
| const originalThen = promise.then.bind(promise); | ||
| let reads = 0; | ||
| Object.defineProperty(promise, 'then', { | ||
| get() { | ||
| reads++; | ||
| return originalThen; | ||
| } | ||
| }); | ||
| return { promise, thenReads: () => reads }; | ||
| } | ||
|
|
||
| describe('Server notification sends on a closed connection', () => { | ||
| test('sendLoggingMessage() when not connected: the inner send rejection is handled in the same microtask', async () => { | ||
| const server = new Server({ name: 'test', version: '1.0.0' }, { capabilities: { logging: {} } }); | ||
| const notConnected = new SdkError(SdkErrorCode.NotConnected, 'Not connected'); | ||
| const inner = instrumentedRejection(notConnected); | ||
|
|
||
| // Plain instance override rather than `vi.spyOn`: the spy wrapper itself | ||
| // reads `then` on any promise a spied call returns, which would mask the | ||
| // observation below. | ||
| (server as unknown as { _notificationViaCodec: () => Promise<void> })._notificationViaCodec = () => inner.promise; | ||
|
|
||
| await expect(server.sendLoggingMessage({ level: 'info', data: 'hello' })).rejects.toBe(notConnected); | ||
|
|
||
|
claude[bot] marked this conversation as resolved.
|
||
| // `then` is read (synchronously at return, called one microtask later) | ||
| // only when `notification()` returns the inner promise without awaiting it. | ||
| expect(inner.thenReads()).toBe(0); | ||
| }); | ||
|
|
||
| test('sendLoggingMessage() when not connected still rejects with SdkError NotConnected (unstubbed path)', async () => { | ||
| const server = new Server({ name: 'test', version: '1.0.0' }, { capabilities: { logging: {} } }); | ||
|
|
||
| await expect(server.sendLoggingMessage({ level: 'info', data: 'hello' })).rejects.toSatisfy( | ||
| (error: unknown) => error instanceof SdkError && error.code === SdkErrorCode.NotConnected | ||
| ); | ||
| }); | ||
| }); | ||
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.