Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/fix-double-onerror-startOrAuthSse.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
'@modelcontextprotocol/client': patch
---

Fix double `onerror` invocation when `_startOrAuthSse` fails. The internal catch block fired `onerror` then threw, and all callers already `.catch(onerror)`, causing every failure to fire twice. Removed the redundant internal call.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 The changeset description inaccurately states "all callers already .catch(onerror)" -- resumeStream did NOT have .catch(onerror) before this PR and was silently dropping errors. The PR makes two distinct changes (removing double-fire for three fire-and-forget callers AND adding new error handling to resumeStream), but the changeset only documents the first, leaving the behavioral change to resumeStream undocumented.

Extended reasoning...

What the bug is

The changeset at .changeset/fix-double-onerror-startOrAuthSse.md line 5 reads:

Fix double onerror invocation when _startOrAuthSse fails. The internal catch block fired onerror then threw, and all callers already .catch(onerror), causing every failure to fire twice. Removed the redundant internal call.

The claim "all callers already .catch(onerror)" is factually incorrect. There are four callers of _startOrAuthSse, not three: lines 344, 531/532, 642, and resumeStream (~line 750). The first three are fire-and-forget callers that each chain .catch(error => this.onerror?.(error)). But resumeStream was different -- it simply did await this._startOrAuthSse({...}) with no try-catch and no .catch(onerror) at all.

The specific code path

Before this PR, when resumeStream failed:

  1. _startOrAuthSse threw (after its internal catch called onerror and rethrew)
  2. resumeStream had no handler, so the exception propagated to the caller
  3. onerror was called exactly once -- by the internal catch in _startOrAuthSse

After this PR removes the internal catch from _startOrAuthSse, the PR correctly compensates by adding a try-catch to resumeStream itself (lines 747-758 in the diff). So the code is correct. But the changeset description omits this second change entirely.

Why existing documentation does not cover it

The changeset is the primary artifact consumers read to understand what changed between versions. By describing the fix purely as "removing the redundant internal call" and saying "all callers already .catch(onerror)", a reader would not know that resumeStream now has onerror-reporting behavior it lacked before. Prior to this PR, a resumeStream failure would call onerror once (via the internal catch). After this PR, it still calls onerror once -- but via a different mechanism, and the behavior is preserved only because the PR added a new try-catch. The changeset implies no behavioral change occurred for any caller.

Step-by-step proof

  1. Before this PR: resumeStream calls await this._startOrAuthSse({resumptionToken: lastEventId})
  2. Server returns HTTP 500; _startOrAuthSse throws SdkError(ClientHttpFailedToOpenStream, ...)
  3. _startOrAuthSse internal catch fires: calls this.onerror?.(error), then rethrows
  4. resumeStream has no handler -- exception propagates to caller
  5. onerror fires once (from step 3)
  6. After this PR: step 3 no longer calls onerror internally
  7. The new try-catch in resumeStream catches the error, calls this.onerror?.(error), rethrows
  8. onerror fires once (now from step 7 instead of step 3)
  9. The changeset says only "removed the redundant internal call" -- it does not mention step 7 was added

Impact and fix

The impact is documentation-only -- the code behavior is correct. A changelog reader consulting this entry would not know resumeStream's error-reporting pathway changed. The fix is to update the changeset to accurately describe both changes: (1) removing the redundant double-fire for the three fire-and-forget callers, and (2) adding explicit onerror handling to resumeStream which previously relied on the internal catch that was removed.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The trace in this comment shows resumeStream fired onerror exactly once before and exactly once after, with the error still propagating. Observable behavior is identical, so there's nothing consumer-facing to document. The 'all callers' phrasing describes a private method consumers can't call.

117 changes: 58 additions & 59 deletions packages/client/src/client/streamableHttp.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -233,72 +233,66 @@ export class StreamableHTTPClientTransport implements Transport {
private async _startOrAuthSse(options: StartSSEOptions, isAuthRetry = false): Promise<void> {
const { resumptionToken } = options;

try {
// Try to open an initial SSE stream with GET to listen for server messages
// This is optional according to the spec - server may not support it
const headers = await this._commonHeaders();
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'text/event-stream'];
headers.set('accept', [...new Set(types)].join(', '));

// Include Last-Event-ID header for resumable streams if provided
if (resumptionToken) {
headers.set('last-event-id', resumptionToken);
}
// Try to open an initial SSE stream with GET to listen for server messages
// This is optional according to the spec - server may not support it
const headers = await this._commonHeaders();
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'text/event-stream'];
headers.set('accept', [...new Set(types)].join(', '));

// Include Last-Event-ID header for resumable streams if provided
if (resumptionToken) {
headers.set('last-event-id', resumptionToken);
}

const response = await (this._fetch ?? fetch)(this._url, {
...this._requestInit,
method: 'GET',
headers,
signal: this._abortController?.signal
});
const response = await (this._fetch ?? fetch)(this._url, {
...this._requestInit,
method: 'GET',
headers,
signal: this._abortController?.signal
});

if (!response.ok) {
if (response.status === 401 && this._authProvider) {
if (response.headers.has('www-authenticate')) {
const { resourceMetadataUrl, scope } = extractWWWAuthenticateParams(response);
this._resourceMetadataUrl = resourceMetadataUrl;
this._scope = scope;
}
if (!response.ok) {
if (response.status === 401 && this._authProvider) {
if (response.headers.has('www-authenticate')) {
const { resourceMetadataUrl, scope } = extractWWWAuthenticateParams(response);
this._resourceMetadataUrl = resourceMetadataUrl;
this._scope = scope;
}

if (this._authProvider.onUnauthorized && !isAuthRetry) {
await this._authProvider.onUnauthorized({
response,
serverUrl: this._url,
fetchFn: this._fetchWithInit
});
await response.text?.().catch(() => {});
// Purposely _not_ awaited, so we don't call onerror twice
return this._startOrAuthSse(options, true);
}
if (this._authProvider.onUnauthorized && !isAuthRetry) {
await this._authProvider.onUnauthorized({
response,
serverUrl: this._url,
fetchFn: this._fetchWithInit
});
await response.text?.().catch(() => {});
if (isAuthRetry) {
throw new SdkError(SdkErrorCode.ClientHttpAuthentication, 'Server returned 401 after re-authentication', {
status: 401
});
}
throw new UnauthorizedError();
return this._startOrAuthSse(options, true);
}

await response.text?.().catch(() => {});

// 405 indicates that the server does not offer an SSE stream at GET endpoint
// This is an expected case that should not trigger an error
if (response.status === 405) {
return;
if (isAuthRetry) {
throw new SdkError(SdkErrorCode.ClientHttpAuthentication, 'Server returned 401 after re-authentication', {
status: 401
});
}
throw new UnauthorizedError();
}

throw new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
status: response.status,
statusText: response.statusText
});
await response.text?.().catch(() => {});

// 405 indicates that the server does not offer an SSE stream at GET endpoint
// This is an expected case that should not trigger an error
if (response.status === 405) {
return;
}

this._handleSseStream(response.body, options, true);
} catch (error) {
this.onerror?.(error as Error);
throw error;
throw new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
status: response.status,
statusText: response.statusText
});
}
Comment thread
claude[bot] marked this conversation as resolved.

this._handleSseStream(response.body, options, true);
}

/**
Expand DownExpand Up@@ -753,9 +747,14 @@ export class StreamableHTTPClientTransport implements Transport {
* @param options Optional callback to receive new resumption tokens
*/
async resumeStream(lastEventId: string, options?: { onresumptiontoken?: (token: string) => void }): Promise<void> {
await this._startOrAuthSse({
resumptionToken: lastEventId,
onresumptiontoken: options?.onresumptiontoken
});
try {
await this._startOrAuthSse({
resumptionToken: lastEventId,
onresumptiontoken: options?.onresumptiontoken
});
} catch (error) {
this.onerror?.(error as Error);
throw error;
}
}
}
59 changes: 59 additions & 0 deletions packages/client/test/client/streamableHttp.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -1224,6 +1224,65 @@ describe('StreamableHTTPClientTransport', () => {
expect(fetchMock.mock.calls[0]![1]?.method).toBe('POST');
});

it('should fire onerror exactly once when _startOrAuthSse fails (not double-fire from catch + caller)', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'));

const errorSpy = vi.fn();
transport.onerror = errorSpy;

const fetchMock = globalThis.fetch as Mock;

// POST returns 202, which triggers _startOrAuthSse when the outbound
// message is an initialized notification (streamableHttp.ts:642)
fetchMock.mockResolvedValueOnce({
ok: true,
status: 202,
headers: new Headers(),
text: async () => ''
});

// The subsequent GET (_startOrAuthSse) fails with a non-ok status
fetchMock.mockResolvedValueOnce({
ok: false,
status: 500,
statusText: 'Internal Server Error',
headers: new Headers(),
text: async () => 'server error'
});

await transport.start();
// Sending an initialized notification triggers the _startOrAuthSse path
await transport.send({ jsonrpc: '2.0', method: 'notifications/initialized' });

// Let the fire-and-forget _startOrAuthSse().catch() settle
await vi.runAllTimersAsync();

expect(errorSpy).toHaveBeenCalledTimes(1);
expect(errorSpy.mock.calls[0]![0].message).toContain('Failed to open SSE stream');
});

it('should fire onerror and reject when resumeStream fails', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'));

const errorSpy = vi.fn();
transport.onerror = errorSpy;

const fetchMock = globalThis.fetch as Mock;
fetchMock.mockResolvedValueOnce({
ok: false,
status: 500,
statusText: 'Internal Server Error',
headers: new Headers(),
text: async () => 'server error'
});

await transport.start();
await expect(transport.resumeStream('event-123')).rejects.toThrow('Failed to open SSE stream');

expect(errorSpy).toHaveBeenCalledTimes(1);
expect(errorSpy.mock.calls[0]![0].message).toContain('Failed to open SSE stream');
});

it('should not throw JSON parse error on priming events with empty data', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'));

Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Add copy buttons to all
 blocks
(function() {
function addCopyButtons() {
document.querySelectorAll('pre code').forEach(function(codeBlock) {
if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;
codeBlock.parentElement.setAttribute('data-copy-added', 'true');
var btn = document.createElement('button');
btn.textContent = 'Copy';
btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';
btn.onmouseover = function() { this.style.opacity = '1'; };
btn.onmouseout = function() { this.style.opacity = '0.7'; };
btn.onclick = function() {
navigator.clipboard.writeText(codeBlock.textContent).then(function() {
btn.textContent = 'Copied!';
setTimeout(function() { btn.textContent = 'Copy'; }, 1500);
});
};
codeBlock.parentElement.style.position = 'relative';
codeBlock.parentElement.appendChild(btn);
});
}
addCopyButtons();
// Re-run on dynamic content
var observer = new MutationObserver(addCopyButtons);
observer.observe(document.body, { childList: true, subtree: true });
})();
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
fix(client): remove redundant onerror in _startOrAuthSse catch by felixweinberger · Pull Request #1826 · modelcontextprotocol/typescript-sdk · GitHub
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/fix-double-onerror-startOrAuthSse.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
'@modelcontextprotocol/client': patch
---

Fix double `onerror` invocation when `_startOrAuthSse` fails. The internal catch block fired `onerror` then threw, and all callers already `.catch(onerror)`, causing every failure to fire twice. Removed the redundant internal call.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 The changeset description inaccurately states "all callers already .catch(onerror)" -- resumeStream did NOT have .catch(onerror) before this PR and was silently dropping errors. The PR makes two distinct changes (removing double-fire for three fire-and-forget callers AND adding new error handling to resumeStream), but the changeset only documents the first, leaving the behavioral change to resumeStream undocumented.

Extended reasoning...

What the bug is

The changeset at .changeset/fix-double-onerror-startOrAuthSse.md line 5 reads:

Fix double onerror invocation when _startOrAuthSse fails. The internal catch block fired onerror then threw, and all callers already .catch(onerror), causing every failure to fire twice. Removed the redundant internal call.

The claim "all callers already .catch(onerror)" is factually incorrect. There are four callers of _startOrAuthSse, not three: lines 344, 531/532, 642, and resumeStream (~line 750). The first three are fire-and-forget callers that each chain .catch(error => this.onerror?.(error)). But resumeStream was different -- it simply did await this._startOrAuthSse({...}) with no try-catch and no .catch(onerror) at all.

The specific code path

Before this PR, when resumeStream failed:

  1. _startOrAuthSse threw (after its internal catch called onerror and rethrew)
  2. resumeStream had no handler, so the exception propagated to the caller
  3. onerror was called exactly once -- by the internal catch in _startOrAuthSse

After this PR removes the internal catch from _startOrAuthSse, the PR correctly compensates by adding a try-catch to resumeStream itself (lines 747-758 in the diff). So the code is correct. But the changeset description omits this second change entirely.

Why existing documentation does not cover it

The changeset is the primary artifact consumers read to understand what changed between versions. By describing the fix purely as "removing the redundant internal call" and saying "all callers already .catch(onerror)", a reader would not know that resumeStream now has onerror-reporting behavior it lacked before. Prior to this PR, a resumeStream failure would call onerror once (via the internal catch). After this PR, it still calls onerror once -- but via a different mechanism, and the behavior is preserved only because the PR added a new try-catch. The changeset implies no behavioral change occurred for any caller.

Step-by-step proof

  1. Before this PR: resumeStream calls await this._startOrAuthSse({resumptionToken: lastEventId})
  2. Server returns HTTP 500; _startOrAuthSse throws SdkError(ClientHttpFailedToOpenStream, ...)
  3. _startOrAuthSse internal catch fires: calls this.onerror?.(error), then rethrows
  4. resumeStream has no handler -- exception propagates to caller
  5. onerror fires once (from step 3)
  6. After this PR: step 3 no longer calls onerror internally
  7. The new try-catch in resumeStream catches the error, calls this.onerror?.(error), rethrows
  8. onerror fires once (now from step 7 instead of step 3)
  9. The changeset says only "removed the redundant internal call" -- it does not mention step 7 was added

Impact and fix

The impact is documentation-only -- the code behavior is correct. A changelog reader consulting this entry would not know resumeStream's error-reporting pathway changed. The fix is to update the changeset to accurately describe both changes: (1) removing the redundant double-fire for the three fire-and-forget callers, and (2) adding explicit onerror handling to resumeStream which previously relied on the internal catch that was removed.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The trace in this comment shows resumeStream fired onerror exactly once before and exactly once after, with the error still propagating. Observable behavior is identical, so there's nothing consumer-facing to document. The 'all callers' phrasing describes a private method consumers can't call.

117 changes: 58 additions & 59 deletions packages/client/src/client/streamableHttp.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -233,72 +233,66 @@ export class StreamableHTTPClientTransport implements Transport {
private async _startOrAuthSse(options: StartSSEOptions, isAuthRetry = false): Promise<void> {
const { resumptionToken } = options;

try {
// Try to open an initial SSE stream with GET to listen for server messages
// This is optional according to the spec - server may not support it
const headers = await this._commonHeaders();
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'text/event-stream'];
headers.set('accept', [...new Set(types)].join(', '));

// Include Last-Event-ID header for resumable streams if provided
if (resumptionToken) {
headers.set('last-event-id', resumptionToken);
}
// Try to open an initial SSE stream with GET to listen for server messages
// This is optional according to the spec - server may not support it
const headers = await this._commonHeaders();
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'text/event-stream'];
headers.set('accept', [...new Set(types)].join(', '));

// Include Last-Event-ID header for resumable streams if provided
if (resumptionToken) {
headers.set('last-event-id', resumptionToken);
}

const response = await (this._fetch ?? fetch)(this._url, {
...this._requestInit,
method: 'GET',
headers,
signal: this._abortController?.signal
});
const response = await (this._fetch ?? fetch)(this._url, {
...this._requestInit,
method: 'GET',
headers,
signal: this._abortController?.signal
});

if (!response.ok) {
if (response.status === 401 && this._authProvider) {
if (response.headers.has('www-authenticate')) {
const { resourceMetadataUrl, scope } = extractWWWAuthenticateParams(response);
this._resourceMetadataUrl = resourceMetadataUrl;
this._scope = scope;
}
if (!response.ok) {
if (response.status === 401 && this._authProvider) {
if (response.headers.has('www-authenticate')) {
const { resourceMetadataUrl, scope } = extractWWWAuthenticateParams(response);
this._resourceMetadataUrl = resourceMetadataUrl;
this._scope = scope;
}

if (this._authProvider.onUnauthorized && !isAuthRetry) {
await this._authProvider.onUnauthorized({
response,
serverUrl: this._url,
fetchFn: this._fetchWithInit
});
await response.text?.().catch(() => {});
// Purposely _not_ awaited, so we don't call onerror twice
return this._startOrAuthSse(options, true);
}
if (this._authProvider.onUnauthorized && !isAuthRetry) {
await this._authProvider.onUnauthorized({
response,
serverUrl: this._url,
fetchFn: this._fetchWithInit
});
await response.text?.().catch(() => {});
if (isAuthRetry) {
throw new SdkError(SdkErrorCode.ClientHttpAuthentication, 'Server returned 401 after re-authentication', {
status: 401
});
}
throw new UnauthorizedError();
return this._startOrAuthSse(options, true);
}

await response.text?.().catch(() => {});

// 405 indicates that the server does not offer an SSE stream at GET endpoint
// This is an expected case that should not trigger an error
if (response.status === 405) {
return;
if (isAuthRetry) {
throw new SdkError(SdkErrorCode.ClientHttpAuthentication, 'Server returned 401 after re-authentication', {
status: 401
});
}
throw new UnauthorizedError();
}

throw new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
status: response.status,
statusText: response.statusText
});
await response.text?.().catch(() => {});

// 405 indicates that the server does not offer an SSE stream at GET endpoint
// This is an expected case that should not trigger an error
if (response.status === 405) {
return;
}

this._handleSseStream(response.body, options, true);
} catch (error) {
this.onerror?.(error as Error);
throw error;
throw new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
status: response.status,
statusText: response.statusText
});
}
Comment thread
claude[bot] marked this conversation as resolved.

this._handleSseStream(response.body, options, true);
}

/**
Expand DownExpand Up@@ -753,9 +747,14 @@ export class StreamableHTTPClientTransport implements Transport {
* @param options Optional callback to receive new resumption tokens
*/
async resumeStream(lastEventId: string, options?: { onresumptiontoken?: (token: string) => void }): Promise<void> {
await this._startOrAuthSse({
resumptionToken: lastEventId,
onresumptiontoken: options?.onresumptiontoken
});
try {
await this._startOrAuthSse({
resumptionToken: lastEventId,
onresumptiontoken: options?.onresumptiontoken
});
} catch (error) {
this.onerror?.(error as Error);
throw error;
}
}
}
59 changes: 59 additions & 0 deletions packages/client/test/client/streamableHttp.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -1224,6 +1224,65 @@ describe('StreamableHTTPClientTransport', () => {
expect(fetchMock.mock.calls[0]![1]?.method).toBe('POST');
});

it('should fire onerror exactly once when _startOrAuthSse fails (not double-fire from catch + caller)', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'));

const errorSpy = vi.fn();
transport.onerror = errorSpy;

const fetchMock = globalThis.fetch as Mock;

// POST returns 202, which triggers _startOrAuthSse when the outbound
// message is an initialized notification (streamableHttp.ts:642)
fetchMock.mockResolvedValueOnce({
ok: true,
status: 202,
headers: new Headers(),
text: async () => ''
});

// The subsequent GET (_startOrAuthSse) fails with a non-ok status
fetchMock.mockResolvedValueOnce({
ok: false,
status: 500,
statusText: 'Internal Server Error',
headers: new Headers(),
text: async () => 'server error'
});

await transport.start();
// Sending an initialized notification triggers the _startOrAuthSse path
await transport.send({ jsonrpc: '2.0', method: 'notifications/initialized' });

// Let the fire-and-forget _startOrAuthSse().catch() settle
await vi.runAllTimersAsync();

expect(errorSpy).toHaveBeenCalledTimes(1);
expect(errorSpy.mock.calls[0]![0].message).toContain('Failed to open SSE stream');
});

it('should fire onerror and reject when resumeStream fails', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'));

const errorSpy = vi.fn();
transport.onerror = errorSpy;

const fetchMock = globalThis.fetch as Mock;
fetchMock.mockResolvedValueOnce({
ok: false,
status: 500,
statusText: 'Internal Server Error',
headers: new Headers(),
text: async () => 'server error'
});

await transport.start();
await expect(transport.resumeStream('event-123')).rejects.toThrow('Failed to open SSE stream');

expect(errorSpy).toHaveBeenCalledTimes(1);
expect(errorSpy.mock.calls[0]![0].message).toContain('Failed to open SSE stream');
});

it('should not throw JSON parse error on priming events with empty data', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'));

Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Force GitHub README to respect dark mode (function() { var style = document.createElement('style'); style.textContent = ' .markdown-body { color-scheme: dark light; } .markdown-body pre { background: #161b22 !important; } .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; } .markdown-body table th, .markdown-body table td { border-color: #30363d !important; } .markdown-body img { background: #0d1117; } .markdown-body blockquote { border-left-color: #8b949e; } .markdown-body hr { border-color: #30363d; } '; document.head.appendChild(style); })(); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' fix(client): remove redundant onerror in _startOrAuthSse catch by felixweinberger · Pull Request #1826 · modelcontextprotocol/typescript-sdk · GitHub
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/fix-double-onerror-startOrAuthSse.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
'@modelcontextprotocol/client': patch
---

Fix double `onerror` invocation when `_startOrAuthSse` fails. The internal catch block fired `onerror` then threw, and all callers already `.catch(onerror)`, causing every failure to fire twice. Removed the redundant internal call.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 The changeset description inaccurately states "all callers already .catch(onerror)" -- resumeStream did NOT have .catch(onerror) before this PR and was silently dropping errors. The PR makes two distinct changes (removing double-fire for three fire-and-forget callers AND adding new error handling to resumeStream), but the changeset only documents the first, leaving the behavioral change to resumeStream undocumented.

Extended reasoning...

What the bug is

The changeset at .changeset/fix-double-onerror-startOrAuthSse.md line 5 reads:

Fix double onerror invocation when _startOrAuthSse fails. The internal catch block fired onerror then threw, and all callers already .catch(onerror), causing every failure to fire twice. Removed the redundant internal call.

The claim "all callers already .catch(onerror)" is factually incorrect. There are four callers of _startOrAuthSse, not three: lines 344, 531/532, 642, and resumeStream (~line 750). The first three are fire-and-forget callers that each chain .catch(error => this.onerror?.(error)). But resumeStream was different -- it simply did await this._startOrAuthSse({...}) with no try-catch and no .catch(onerror) at all.

The specific code path

Before this PR, when resumeStream failed:

  1. _startOrAuthSse threw (after its internal catch called onerror and rethrew)
  2. resumeStream had no handler, so the exception propagated to the caller
  3. onerror was called exactly once -- by the internal catch in _startOrAuthSse

After this PR removes the internal catch from _startOrAuthSse, the PR correctly compensates by adding a try-catch to resumeStream itself (lines 747-758 in the diff). So the code is correct. But the changeset description omits this second change entirely.

Why existing documentation does not cover it

The changeset is the primary artifact consumers read to understand what changed between versions. By describing the fix purely as "removing the redundant internal call" and saying "all callers already .catch(onerror)", a reader would not know that resumeStream now has onerror-reporting behavior it lacked before. Prior to this PR, a resumeStream failure would call onerror once (via the internal catch). After this PR, it still calls onerror once -- but via a different mechanism, and the behavior is preserved only because the PR added a new try-catch. The changeset implies no behavioral change occurred for any caller.

Step-by-step proof

  1. Before this PR: resumeStream calls await this._startOrAuthSse({resumptionToken: lastEventId})
  2. Server returns HTTP 500; _startOrAuthSse throws SdkError(ClientHttpFailedToOpenStream, ...)
  3. _startOrAuthSse internal catch fires: calls this.onerror?.(error), then rethrows
  4. resumeStream has no handler -- exception propagates to caller
  5. onerror fires once (from step 3)
  6. After this PR: step 3 no longer calls onerror internally
  7. The new try-catch in resumeStream catches the error, calls this.onerror?.(error), rethrows
  8. onerror fires once (now from step 7 instead of step 3)
  9. The changeset says only "removed the redundant internal call" -- it does not mention step 7 was added

Impact and fix

The impact is documentation-only -- the code behavior is correct. A changelog reader consulting this entry would not know resumeStream's error-reporting pathway changed. The fix is to update the changeset to accurately describe both changes: (1) removing the redundant double-fire for the three fire-and-forget callers, and (2) adding explicit onerror handling to resumeStream which previously relied on the internal catch that was removed.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The trace in this comment shows resumeStream fired onerror exactly once before and exactly once after, with the error still propagating. Observable behavior is identical, so there's nothing consumer-facing to document. The 'all callers' phrasing describes a private method consumers can't call.

117 changes: 58 additions & 59 deletions packages/client/src/client/streamableHttp.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -233,72 +233,66 @@ export class StreamableHTTPClientTransport implements Transport {
private async _startOrAuthSse(options: StartSSEOptions, isAuthRetry = false): Promise<void> {
const { resumptionToken } = options;

try {
// Try to open an initial SSE stream with GET to listen for server messages
// This is optional according to the spec - server may not support it
const headers = await this._commonHeaders();
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'text/event-stream'];
headers.set('accept', [...new Set(types)].join(', '));

// Include Last-Event-ID header for resumable streams if provided
if (resumptionToken) {
headers.set('last-event-id', resumptionToken);
}
// Try to open an initial SSE stream with GET to listen for server messages
// This is optional according to the spec - server may not support it
const headers = await this._commonHeaders();
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'text/event-stream'];
headers.set('accept', [...new Set(types)].join(', '));

// Include Last-Event-ID header for resumable streams if provided
if (resumptionToken) {
headers.set('last-event-id', resumptionToken);
}

const response = await (this._fetch ?? fetch)(this._url, {
...this._requestInit,
method: 'GET',
headers,
signal: this._abortController?.signal
});
const response = await (this._fetch ?? fetch)(this._url, {
...this._requestInit,
method: 'GET',
headers,
signal: this._abortController?.signal
});

if (!response.ok) {
if (response.status === 401 && this._authProvider) {
if (response.headers.has('www-authenticate')) {
const { resourceMetadataUrl, scope } = extractWWWAuthenticateParams(response);
this._resourceMetadataUrl = resourceMetadataUrl;
this._scope = scope;
}
if (!response.ok) {
if (response.status === 401 && this._authProvider) {
if (response.headers.has('www-authenticate')) {
const { resourceMetadataUrl, scope } = extractWWWAuthenticateParams(response);
this._resourceMetadataUrl = resourceMetadataUrl;
this._scope = scope;
}

if (this._authProvider.onUnauthorized && !isAuthRetry) {
await this._authProvider.onUnauthorized({
response,
serverUrl: this._url,
fetchFn: this._fetchWithInit
});
await response.text?.().catch(() => {});
// Purposely _not_ awaited, so we don't call onerror twice
return this._startOrAuthSse(options, true);
}
if (this._authProvider.onUnauthorized && !isAuthRetry) {
await this._authProvider.onUnauthorized({
response,
serverUrl: this._url,
fetchFn: this._fetchWithInit
});
await response.text?.().catch(() => {});
if (isAuthRetry) {
throw new SdkError(SdkErrorCode.ClientHttpAuthentication, 'Server returned 401 after re-authentication', {
status: 401
});
}
throw new UnauthorizedError();
return this._startOrAuthSse(options, true);
}

await response.text?.().catch(() => {});

// 405 indicates that the server does not offer an SSE stream at GET endpoint
// This is an expected case that should not trigger an error
if (response.status === 405) {
return;
if (isAuthRetry) {
throw new SdkError(SdkErrorCode.ClientHttpAuthentication, 'Server returned 401 after re-authentication', {
status: 401
});
}
throw new UnauthorizedError();
}

throw new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
status: response.status,
statusText: response.statusText
});
await response.text?.().catch(() => {});

// 405 indicates that the server does not offer an SSE stream at GET endpoint
// This is an expected case that should not trigger an error
if (response.status === 405) {
return;
}

this._handleSseStream(response.body, options, true);
} catch (error) {
this.onerror?.(error as Error);
throw error;
throw new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
status: response.status,
statusText: response.statusText
});
}
Comment thread
claude[bot] marked this conversation as resolved.

this._handleSseStream(response.body, options, true);
}

/**
Expand DownExpand Up@@ -753,9 +747,14 @@ export class StreamableHTTPClientTransport implements Transport {
* @param options Optional callback to receive new resumption tokens
*/
async resumeStream(lastEventId: string, options?: { onresumptiontoken?: (token: string) => void }): Promise<void> {
await this._startOrAuthSse({
resumptionToken: lastEventId,
onresumptiontoken: options?.onresumptiontoken
});
try {
await this._startOrAuthSse({
resumptionToken: lastEventId,
onresumptiontoken: options?.onresumptiontoken
});
} catch (error) {
this.onerror?.(error as Error);
throw error;
}
}
}
59 changes: 59 additions & 0 deletions packages/client/test/client/streamableHttp.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -1224,6 +1224,65 @@ describe('StreamableHTTPClientTransport', () => {
expect(fetchMock.mock.calls[0]![1]?.method).toBe('POST');
});

it('should fire onerror exactly once when _startOrAuthSse fails (not double-fire from catch + caller)', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'));

const errorSpy = vi.fn();
transport.onerror = errorSpy;

const fetchMock = globalThis.fetch as Mock;

// POST returns 202, which triggers _startOrAuthSse when the outbound
// message is an initialized notification (streamableHttp.ts:642)
fetchMock.mockResolvedValueOnce({
ok: true,
status: 202,
headers: new Headers(),
text: async () => ''
});

// The subsequent GET (_startOrAuthSse) fails with a non-ok status
fetchMock.mockResolvedValueOnce({
ok: false,
status: 500,
statusText: 'Internal Server Error',
headers: new Headers(),
text: async () => 'server error'
});

await transport.start();
// Sending an initialized notification triggers the _startOrAuthSse path
await transport.send({ jsonrpc: '2.0', method: 'notifications/initialized' });

// Let the fire-and-forget _startOrAuthSse().catch() settle
await vi.runAllTimersAsync();

expect(errorSpy).toHaveBeenCalledTimes(1);
expect(errorSpy.mock.calls[0]![0].message).toContain('Failed to open SSE stream');
});

it('should fire onerror and reject when resumeStream fails', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'));

const errorSpy = vi.fn();
transport.onerror = errorSpy;

const fetchMock = globalThis.fetch as Mock;
fetchMock.mockResolvedValueOnce({
ok: false,
status: 500,
statusText: 'Internal Server Error',
headers: new Headers(),
text: async () => 'server error'
});

await transport.start();
await expect(transport.resumeStream('event-123')).rejects.toThrow('Failed to open SSE stream');

expect(errorSpy).toHaveBeenCalledTimes(1);
expect(errorSpy.mock.calls[0]![0].message).toContain('Failed to open SSE stream');
});

it('should not throw JSON parse error on priming events with empty data', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'));

Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Highlight search terms from Google/DuckDuckGo/Bing referrer (function() { var ref = document.referrer; var terms = []; if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) { var url = new URL(ref); var q = url.searchParams.get('q') || url.searchParams.get('p'); if (q) { terms = q.split(/\s+/).filter(function(t) { return t.length > 2; }); } } if (terms.length === 0) return; var style = document.createElement('style'); style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }'; document.head.appendChild(style); function highlight(node) { if (node.nodeType === 3) { // text node var text = node.textContent; var found = false; terms.forEach(function(term) { var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\]\\]/g, '\\') + ')', 'gi'); if (regex.test(text)) { found = true; var frag = document.createDocumentFragment(); var parts = text.split(regex); parts.forEach(function(part, i) { if (i % 2 === 0) { frag.appendChild(document.createTextNode(part)); } else { var span = document.createElement('span'); span.className = 'userscript-highlight'; span.textContent = part; frag.appendChild(span); } }); node.parentNode.replaceChild(frag, node); } }); } else if (node.nodeType === 1 && node.childNodes) { // element var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT']; if (!skipTags.includes(node.tagName)) { Array.from(node.childNodes).forEach(highlight); } } } highlight(document.body); // Re-highlight on dynamic content var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1 || node.nodeType === 3) highlight(node); }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' fix(client): remove redundant onerror in _startOrAuthSse catch by felixweinberger · Pull Request #1826 · modelcontextprotocol/typescript-sdk · GitHub
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/fix-double-onerror-startOrAuthSse.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
'@modelcontextprotocol/client': patch
---

Fix double `onerror` invocation when `_startOrAuthSse` fails. The internal catch block fired `onerror` then threw, and all callers already `.catch(onerror)`, causing every failure to fire twice. Removed the redundant internal call.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 The changeset description inaccurately states "all callers already .catch(onerror)" -- resumeStream did NOT have .catch(onerror) before this PR and was silently dropping errors. The PR makes two distinct changes (removing double-fire for three fire-and-forget callers AND adding new error handling to resumeStream), but the changeset only documents the first, leaving the behavioral change to resumeStream undocumented.

Extended reasoning...

What the bug is

The changeset at .changeset/fix-double-onerror-startOrAuthSse.md line 5 reads:

Fix double onerror invocation when _startOrAuthSse fails. The internal catch block fired onerror then threw, and all callers already .catch(onerror), causing every failure to fire twice. Removed the redundant internal call.

The claim "all callers already .catch(onerror)" is factually incorrect. There are four callers of _startOrAuthSse, not three: lines 344, 531/532, 642, and resumeStream (~line 750). The first three are fire-and-forget callers that each chain .catch(error => this.onerror?.(error)). But resumeStream was different -- it simply did await this._startOrAuthSse({...}) with no try-catch and no .catch(onerror) at all.

The specific code path

Before this PR, when resumeStream failed:

  1. _startOrAuthSse threw (after its internal catch called onerror and rethrew)
  2. resumeStream had no handler, so the exception propagated to the caller
  3. onerror was called exactly once -- by the internal catch in _startOrAuthSse

After this PR removes the internal catch from _startOrAuthSse, the PR correctly compensates by adding a try-catch to resumeStream itself (lines 747-758 in the diff). So the code is correct. But the changeset description omits this second change entirely.

Why existing documentation does not cover it

The changeset is the primary artifact consumers read to understand what changed between versions. By describing the fix purely as "removing the redundant internal call" and saying "all callers already .catch(onerror)", a reader would not know that resumeStream now has onerror-reporting behavior it lacked before. Prior to this PR, a resumeStream failure would call onerror once (via the internal catch). After this PR, it still calls onerror once -- but via a different mechanism, and the behavior is preserved only because the PR added a new try-catch. The changeset implies no behavioral change occurred for any caller.

Step-by-step proof

  1. Before this PR: resumeStream calls await this._startOrAuthSse({resumptionToken: lastEventId})
  2. Server returns HTTP 500; _startOrAuthSse throws SdkError(ClientHttpFailedToOpenStream, ...)
  3. _startOrAuthSse internal catch fires: calls this.onerror?.(error), then rethrows
  4. resumeStream has no handler -- exception propagates to caller
  5. onerror fires once (from step 3)
  6. After this PR: step 3 no longer calls onerror internally
  7. The new try-catch in resumeStream catches the error, calls this.onerror?.(error), rethrows
  8. onerror fires once (now from step 7 instead of step 3)
  9. The changeset says only "removed the redundant internal call" -- it does not mention step 7 was added

Impact and fix

The impact is documentation-only -- the code behavior is correct. A changelog reader consulting this entry would not know resumeStream's error-reporting pathway changed. The fix is to update the changeset to accurately describe both changes: (1) removing the redundant double-fire for the three fire-and-forget callers, and (2) adding explicit onerror handling to resumeStream which previously relied on the internal catch that was removed.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The trace in this comment shows resumeStream fired onerror exactly once before and exactly once after, with the error still propagating. Observable behavior is identical, so there's nothing consumer-facing to document. The 'all callers' phrasing describes a private method consumers can't call.

117 changes: 58 additions & 59 deletions packages/client/src/client/streamableHttp.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -233,72 +233,66 @@ export class StreamableHTTPClientTransport implements Transport {
private async _startOrAuthSse(options: StartSSEOptions, isAuthRetry = false): Promise<void> {
const { resumptionToken } = options;

try {
// Try to open an initial SSE stream with GET to listen for server messages
// This is optional according to the spec - server may not support it
const headers = await this._commonHeaders();
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'text/event-stream'];
headers.set('accept', [...new Set(types)].join(', '));

// Include Last-Event-ID header for resumable streams if provided
if (resumptionToken) {
headers.set('last-event-id', resumptionToken);
}
// Try to open an initial SSE stream with GET to listen for server messages
// This is optional according to the spec - server may not support it
const headers = await this._commonHeaders();
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'text/event-stream'];
headers.set('accept', [...new Set(types)].join(', '));

// Include Last-Event-ID header for resumable streams if provided
if (resumptionToken) {
headers.set('last-event-id', resumptionToken);
}

const response = await (this._fetch ?? fetch)(this._url, {
...this._requestInit,
method: 'GET',
headers,
signal: this._abortController?.signal
});
const response = await (this._fetch ?? fetch)(this._url, {
...this._requestInit,
method: 'GET',
headers,
signal: this._abortController?.signal
});

if (!response.ok) {
if (response.status === 401 && this._authProvider) {
if (response.headers.has('www-authenticate')) {
const { resourceMetadataUrl, scope } = extractWWWAuthenticateParams(response);
this._resourceMetadataUrl = resourceMetadataUrl;
this._scope = scope;
}
if (!response.ok) {
if (response.status === 401 && this._authProvider) {
if (response.headers.has('www-authenticate')) {
const { resourceMetadataUrl, scope } = extractWWWAuthenticateParams(response);
this._resourceMetadataUrl = resourceMetadataUrl;
this._scope = scope;
}

if (this._authProvider.onUnauthorized && !isAuthRetry) {
await this._authProvider.onUnauthorized({
response,
serverUrl: this._url,
fetchFn: this._fetchWithInit
});
await response.text?.().catch(() => {});
// Purposely _not_ awaited, so we don't call onerror twice
return this._startOrAuthSse(options, true);
}
if (this._authProvider.onUnauthorized && !isAuthRetry) {
await this._authProvider.onUnauthorized({
response,
serverUrl: this._url,
fetchFn: this._fetchWithInit
});
await response.text?.().catch(() => {});
if (isAuthRetry) {
throw new SdkError(SdkErrorCode.ClientHttpAuthentication, 'Server returned 401 after re-authentication', {
status: 401
});
}
throw new UnauthorizedError();
return this._startOrAuthSse(options, true);
}

await response.text?.().catch(() => {});

// 405 indicates that the server does not offer an SSE stream at GET endpoint
// This is an expected case that should not trigger an error
if (response.status === 405) {
return;
if (isAuthRetry) {
throw new SdkError(SdkErrorCode.ClientHttpAuthentication, 'Server returned 401 after re-authentication', {
status: 401
});
}
throw new UnauthorizedError();
}

throw new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
status: response.status,
statusText: response.statusText
});
await response.text?.().catch(() => {});

// 405 indicates that the server does not offer an SSE stream at GET endpoint
// This is an expected case that should not trigger an error
if (response.status === 405) {
return;
}

this._handleSseStream(response.body, options, true);
} catch (error) {
this.onerror?.(error as Error);
throw error;
throw new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
status: response.status,
statusText: response.statusText
});
}
Comment thread
claude[bot] marked this conversation as resolved.

this._handleSseStream(response.body, options, true);
}

/**
Expand DownExpand Up@@ -753,9 +747,14 @@ export class StreamableHTTPClientTransport implements Transport {
* @param options Optional callback to receive new resumption tokens
*/
async resumeStream(lastEventId: string, options?: { onresumptiontoken?: (token: string) => void }): Promise<void> {
await this._startOrAuthSse({
resumptionToken: lastEventId,
onresumptiontoken: options?.onresumptiontoken
});
try {
await this._startOrAuthSse({
resumptionToken: lastEventId,
onresumptiontoken: options?.onresumptiontoken
});
} catch (error) {
this.onerror?.(error as Error);
throw error;
}
}
}
59 changes: 59 additions & 0 deletions packages/client/test/client/streamableHttp.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -1224,6 +1224,65 @@ describe('StreamableHTTPClientTransport', () => {
expect(fetchMock.mock.calls[0]![1]?.method).toBe('POST');
});

it('should fire onerror exactly once when _startOrAuthSse fails (not double-fire from catch + caller)', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'));

const errorSpy = vi.fn();
transport.onerror = errorSpy;

const fetchMock = globalThis.fetch as Mock;

// POST returns 202, which triggers _startOrAuthSse when the outbound
// message is an initialized notification (streamableHttp.ts:642)
fetchMock.mockResolvedValueOnce({
ok: true,
status: 202,
headers: new Headers(),
text: async () => ''
});

// The subsequent GET (_startOrAuthSse) fails with a non-ok status
fetchMock.mockResolvedValueOnce({
ok: false,
status: 500,
statusText: 'Internal Server Error',
headers: new Headers(),
text: async () => 'server error'
});

await transport.start();
// Sending an initialized notification triggers the _startOrAuthSse path
await transport.send({ jsonrpc: '2.0', method: 'notifications/initialized' });

// Let the fire-and-forget _startOrAuthSse().catch() settle
await vi.runAllTimersAsync();

expect(errorSpy).toHaveBeenCalledTimes(1);
expect(errorSpy.mock.calls[0]![0].message).toContain('Failed to open SSE stream');
});

it('should fire onerror and reject when resumeStream fails', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'));

const errorSpy = vi.fn();
transport.onerror = errorSpy;

const fetchMock = globalThis.fetch as Mock;
fetchMock.mockResolvedValueOnce({
ok: false,
status: 500,
statusText: 'Internal Server Error',
headers: new Headers(),
text: async () => 'server error'
});

await transport.start();
await expect(transport.resumeStream('event-123')).rejects.toThrow('Failed to open SSE stream');

expect(errorSpy).toHaveBeenCalledTimes(1);
expect(errorSpy.mock.calls[0]![0].message).toContain('Failed to open SSE stream');
});

it('should not throw JSON parse error on priming events with empty data', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'));

Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Strip utm_, fbclid, gclid, etc. from all links on page (function() { var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content', 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid', 'ref', 'ref_src', 'source', 'medium', 'campaign']; function cleanUrl(url) { try { var u = new URL(url, window.location.origin); var changed = false; trackingParams.forEach(function(p) { if (u.searchParams.has(p)) { u.searchParams.delete(p); changed = true; } }); return changed ? u.toString() : url; } catch (e) { return url; } } function cleanLinks() { document.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } cleanLinks(); var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1) { if (node.tagName === 'A') cleanLinks(); node.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + ' fix(client): remove redundant onerror in _startOrAuthSse catch by felixweinberger · Pull Request #1826 · modelcontextprotocol/typescript-sdk · GitHub
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/fix-double-onerror-startOrAuthSse.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
'@modelcontextprotocol/client': patch
---

Fix double `onerror` invocation when `_startOrAuthSse` fails. The internal catch block fired `onerror` then threw, and all callers already `.catch(onerror)`, causing every failure to fire twice. Removed the redundant internal call.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 The changeset description inaccurately states "all callers already .catch(onerror)" -- resumeStream did NOT have .catch(onerror) before this PR and was silently dropping errors. The PR makes two distinct changes (removing double-fire for three fire-and-forget callers AND adding new error handling to resumeStream), but the changeset only documents the first, leaving the behavioral change to resumeStream undocumented.

Extended reasoning...

What the bug is

The changeset at .changeset/fix-double-onerror-startOrAuthSse.md line 5 reads:

Fix double onerror invocation when _startOrAuthSse fails. The internal catch block fired onerror then threw, and all callers already .catch(onerror), causing every failure to fire twice. Removed the redundant internal call.

The claim "all callers already .catch(onerror)" is factually incorrect. There are four callers of _startOrAuthSse, not three: lines 344, 531/532, 642, and resumeStream (~line 750). The first three are fire-and-forget callers that each chain .catch(error => this.onerror?.(error)). But resumeStream was different -- it simply did await this._startOrAuthSse({...}) with no try-catch and no .catch(onerror) at all.

The specific code path

Before this PR, when resumeStream failed:

  1. _startOrAuthSse threw (after its internal catch called onerror and rethrew)
  2. resumeStream had no handler, so the exception propagated to the caller
  3. onerror was called exactly once -- by the internal catch in _startOrAuthSse

After this PR removes the internal catch from _startOrAuthSse, the PR correctly compensates by adding a try-catch to resumeStream itself (lines 747-758 in the diff). So the code is correct. But the changeset description omits this second change entirely.

Why existing documentation does not cover it

The changeset is the primary artifact consumers read to understand what changed between versions. By describing the fix purely as "removing the redundant internal call" and saying "all callers already .catch(onerror)", a reader would not know that resumeStream now has onerror-reporting behavior it lacked before. Prior to this PR, a resumeStream failure would call onerror once (via the internal catch). After this PR, it still calls onerror once -- but via a different mechanism, and the behavior is preserved only because the PR added a new try-catch. The changeset implies no behavioral change occurred for any caller.

Step-by-step proof

  1. Before this PR: resumeStream calls await this._startOrAuthSse({resumptionToken: lastEventId})
  2. Server returns HTTP 500; _startOrAuthSse throws SdkError(ClientHttpFailedToOpenStream, ...)
  3. _startOrAuthSse internal catch fires: calls this.onerror?.(error), then rethrows
  4. resumeStream has no handler -- exception propagates to caller
  5. onerror fires once (from step 3)
  6. After this PR: step 3 no longer calls onerror internally
  7. The new try-catch in resumeStream catches the error, calls this.onerror?.(error), rethrows
  8. onerror fires once (now from step 7 instead of step 3)
  9. The changeset says only "removed the redundant internal call" -- it does not mention step 7 was added

Impact and fix

The impact is documentation-only -- the code behavior is correct. A changelog reader consulting this entry would not know resumeStream's error-reporting pathway changed. The fix is to update the changeset to accurately describe both changes: (1) removing the redundant double-fire for the three fire-and-forget callers, and (2) adding explicit onerror handling to resumeStream which previously relied on the internal catch that was removed.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The trace in this comment shows resumeStream fired onerror exactly once before and exactly once after, with the error still propagating. Observable behavior is identical, so there's nothing consumer-facing to document. The 'all callers' phrasing describes a private method consumers can't call.

117 changes: 58 additions & 59 deletions packages/client/src/client/streamableHttp.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -233,72 +233,66 @@ export class StreamableHTTPClientTransport implements Transport {
private async _startOrAuthSse(options: StartSSEOptions, isAuthRetry = false): Promise<void> {
const { resumptionToken } = options;

try {
// Try to open an initial SSE stream with GET to listen for server messages
// This is optional according to the spec - server may not support it
const headers = await this._commonHeaders();
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'text/event-stream'];
headers.set('accept', [...new Set(types)].join(', '));

// Include Last-Event-ID header for resumable streams if provided
if (resumptionToken) {
headers.set('last-event-id', resumptionToken);
}
// Try to open an initial SSE stream with GET to listen for server messages
// This is optional according to the spec - server may not support it
const headers = await this._commonHeaders();
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'text/event-stream'];
headers.set('accept', [...new Set(types)].join(', '));

// Include Last-Event-ID header for resumable streams if provided
if (resumptionToken) {
headers.set('last-event-id', resumptionToken);
}

const response = await (this._fetch ?? fetch)(this._url, {
...this._requestInit,
method: 'GET',
headers,
signal: this._abortController?.signal
});
const response = await (this._fetch ?? fetch)(this._url, {
...this._requestInit,
method: 'GET',
headers,
signal: this._abortController?.signal
});

if (!response.ok) {
if (response.status === 401 && this._authProvider) {
if (response.headers.has('www-authenticate')) {
const { resourceMetadataUrl, scope } = extractWWWAuthenticateParams(response);
this._resourceMetadataUrl = resourceMetadataUrl;
this._scope = scope;
}
if (!response.ok) {
if (response.status === 401 && this._authProvider) {
if (response.headers.has('www-authenticate')) {
const { resourceMetadataUrl, scope } = extractWWWAuthenticateParams(response);
this._resourceMetadataUrl = resourceMetadataUrl;
this._scope = scope;
}

if (this._authProvider.onUnauthorized && !isAuthRetry) {
await this._authProvider.onUnauthorized({
response,
serverUrl: this._url,
fetchFn: this._fetchWithInit
});
await response.text?.().catch(() => {});
// Purposely _not_ awaited, so we don't call onerror twice
return this._startOrAuthSse(options, true);
}
if (this._authProvider.onUnauthorized && !isAuthRetry) {
await this._authProvider.onUnauthorized({
response,
serverUrl: this._url,
fetchFn: this._fetchWithInit
});
await response.text?.().catch(() => {});
if (isAuthRetry) {
throw new SdkError(SdkErrorCode.ClientHttpAuthentication, 'Server returned 401 after re-authentication', {
status: 401
});
}
throw new UnauthorizedError();
return this._startOrAuthSse(options, true);
}

await response.text?.().catch(() => {});

// 405 indicates that the server does not offer an SSE stream at GET endpoint
// This is an expected case that should not trigger an error
if (response.status === 405) {
return;
if (isAuthRetry) {
throw new SdkError(SdkErrorCode.ClientHttpAuthentication, 'Server returned 401 after re-authentication', {
status: 401
});
}
throw new UnauthorizedError();
}

throw new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
status: response.status,
statusText: response.statusText
});
await response.text?.().catch(() => {});

// 405 indicates that the server does not offer an SSE stream at GET endpoint
// This is an expected case that should not trigger an error
if (response.status === 405) {
return;
}

this._handleSseStream(response.body, options, true);
} catch (error) {
this.onerror?.(error as Error);
throw error;
throw new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
status: response.status,
statusText: response.statusText
});
}
Comment thread
claude[bot] marked this conversation as resolved.

this._handleSseStream(response.body, options, true);
}

/**
Expand DownExpand Up@@ -753,9 +747,14 @@ export class StreamableHTTPClientTransport implements Transport {
* @param options Optional callback to receive new resumption tokens
*/
async resumeStream(lastEventId: string, options?: { onresumptiontoken?: (token: string) => void }): Promise<void> {
await this._startOrAuthSse({
resumptionToken: lastEventId,
onresumptiontoken: options?.onresumptiontoken
});
try {
await this._startOrAuthSse({
resumptionToken: lastEventId,
onresumptiontoken: options?.onresumptiontoken
});
} catch (error) {
this.onerror?.(error as Error);
throw error;
}
}
}
59 changes: 59 additions & 0 deletions packages/client/test/client/streamableHttp.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -1224,6 +1224,65 @@ describe('StreamableHTTPClientTransport', () => {
expect(fetchMock.mock.calls[0]![1]?.method).toBe('POST');
});

it('should fire onerror exactly once when _startOrAuthSse fails (not double-fire from catch + caller)', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'));

const errorSpy = vi.fn();
transport.onerror = errorSpy;

const fetchMock = globalThis.fetch as Mock;

// POST returns 202, which triggers _startOrAuthSse when the outbound
// message is an initialized notification (streamableHttp.ts:642)
fetchMock.mockResolvedValueOnce({
ok: true,
status: 202,
headers: new Headers(),
text: async () => ''
});

// The subsequent GET (_startOrAuthSse) fails with a non-ok status
fetchMock.mockResolvedValueOnce({
ok: false,
status: 500,
statusText: 'Internal Server Error',
headers: new Headers(),
text: async () => 'server error'
});

await transport.start();
// Sending an initialized notification triggers the _startOrAuthSse path
await transport.send({ jsonrpc: '2.0', method: 'notifications/initialized' });

// Let the fire-and-forget _startOrAuthSse().catch() settle
await vi.runAllTimersAsync();

expect(errorSpy).toHaveBeenCalledTimes(1);
expect(errorSpy.mock.calls[0]![0].message).toContain('Failed to open SSE stream');
});

it('should fire onerror and reject when resumeStream fails', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'));

const errorSpy = vi.fn();
transport.onerror = errorSpy;

const fetchMock = globalThis.fetch as Mock;
fetchMock.mockResolvedValueOnce({
ok: false,
status: 500,
statusText: 'Internal Server Error',
headers: new Headers(),
text: async () => 'server error'
});

await transport.start();
await expect(transport.resumeStream('event-123')).rejects.toThrow('Failed to open SSE stream');

expect(errorSpy).toHaveBeenCalledTimes(1);
expect(errorSpy.mock.calls[0]![0].message).toContain('Failed to open SSE stream');
});

it('should not throw JSON parse error on priming events with empty data', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'));

Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Auto-enable theater mode on YouTube (function() { function tryTheater() { var btn = document.querySelector('button[aria-label="Theater mode"], ytd-player #player button[title="Theater mode"]'); if (btn && !btn.classList.contains('activated')) { btn.click(); } } // Try immediately tryTheater(); // Try after navigation (SPA) var lastUrl = location.href; setInterval(function() { if (location.href !== lastUrl) { lastUrl = location.href; setTimeout(tryTheater, 500); } }, 1000); // Also try on player load var observer = new MutationObserver(tryTheater); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' fix(client): remove redundant onerror in _startOrAuthSse catch by felixweinberger · Pull Request #1826 · modelcontextprotocol/typescript-sdk · GitHub
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/fix-double-onerror-startOrAuthSse.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
'@modelcontextprotocol/client': patch
---

Fix double `onerror` invocation when `_startOrAuthSse` fails. The internal catch block fired `onerror` then threw, and all callers already `.catch(onerror)`, causing every failure to fire twice. Removed the redundant internal call.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 The changeset description inaccurately states "all callers already .catch(onerror)" -- resumeStream did NOT have .catch(onerror) before this PR and was silently dropping errors. The PR makes two distinct changes (removing double-fire for three fire-and-forget callers AND adding new error handling to resumeStream), but the changeset only documents the first, leaving the behavioral change to resumeStream undocumented.

Extended reasoning...

What the bug is

The changeset at .changeset/fix-double-onerror-startOrAuthSse.md line 5 reads:

Fix double onerror invocation when _startOrAuthSse fails. The internal catch block fired onerror then threw, and all callers already .catch(onerror), causing every failure to fire twice. Removed the redundant internal call.

The claim "all callers already .catch(onerror)" is factually incorrect. There are four callers of _startOrAuthSse, not three: lines 344, 531/532, 642, and resumeStream (~line 750). The first three are fire-and-forget callers that each chain .catch(error => this.onerror?.(error)). But resumeStream was different -- it simply did await this._startOrAuthSse({...}) with no try-catch and no .catch(onerror) at all.

The specific code path

Before this PR, when resumeStream failed:

  1. _startOrAuthSse threw (after its internal catch called onerror and rethrew)
  2. resumeStream had no handler, so the exception propagated to the caller
  3. onerror was called exactly once -- by the internal catch in _startOrAuthSse

After this PR removes the internal catch from _startOrAuthSse, the PR correctly compensates by adding a try-catch to resumeStream itself (lines 747-758 in the diff). So the code is correct. But the changeset description omits this second change entirely.

Why existing documentation does not cover it

The changeset is the primary artifact consumers read to understand what changed between versions. By describing the fix purely as "removing the redundant internal call" and saying "all callers already .catch(onerror)", a reader would not know that resumeStream now has onerror-reporting behavior it lacked before. Prior to this PR, a resumeStream failure would call onerror once (via the internal catch). After this PR, it still calls onerror once -- but via a different mechanism, and the behavior is preserved only because the PR added a new try-catch. The changeset implies no behavioral change occurred for any caller.

Step-by-step proof

  1. Before this PR: resumeStream calls await this._startOrAuthSse({resumptionToken: lastEventId})
  2. Server returns HTTP 500; _startOrAuthSse throws SdkError(ClientHttpFailedToOpenStream, ...)
  3. _startOrAuthSse internal catch fires: calls this.onerror?.(error), then rethrows
  4. resumeStream has no handler -- exception propagates to caller
  5. onerror fires once (from step 3)
  6. After this PR: step 3 no longer calls onerror internally
  7. The new try-catch in resumeStream catches the error, calls this.onerror?.(error), rethrows
  8. onerror fires once (now from step 7 instead of step 3)
  9. The changeset says only "removed the redundant internal call" -- it does not mention step 7 was added

Impact and fix

The impact is documentation-only -- the code behavior is correct. A changelog reader consulting this entry would not know resumeStream's error-reporting pathway changed. The fix is to update the changeset to accurately describe both changes: (1) removing the redundant double-fire for the three fire-and-forget callers, and (2) adding explicit onerror handling to resumeStream which previously relied on the internal catch that was removed.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The trace in this comment shows resumeStream fired onerror exactly once before and exactly once after, with the error still propagating. Observable behavior is identical, so there's nothing consumer-facing to document. The 'all callers' phrasing describes a private method consumers can't call.

117 changes: 58 additions & 59 deletions packages/client/src/client/streamableHttp.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -233,72 +233,66 @@ export class StreamableHTTPClientTransport implements Transport {
private async _startOrAuthSse(options: StartSSEOptions, isAuthRetry = false): Promise<void> {
const { resumptionToken } = options;

try {
// Try to open an initial SSE stream with GET to listen for server messages
// This is optional according to the spec - server may not support it
const headers = await this._commonHeaders();
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'text/event-stream'];
headers.set('accept', [...new Set(types)].join(', '));

// Include Last-Event-ID header for resumable streams if provided
if (resumptionToken) {
headers.set('last-event-id', resumptionToken);
}
// Try to open an initial SSE stream with GET to listen for server messages
// This is optional according to the spec - server may not support it
const headers = await this._commonHeaders();
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'text/event-stream'];
headers.set('accept', [...new Set(types)].join(', '));

// Include Last-Event-ID header for resumable streams if provided
if (resumptionToken) {
headers.set('last-event-id', resumptionToken);
}

const response = await (this._fetch ?? fetch)(this._url, {
...this._requestInit,
method: 'GET',
headers,
signal: this._abortController?.signal
});
const response = await (this._fetch ?? fetch)(this._url, {
...this._requestInit,
method: 'GET',
headers,
signal: this._abortController?.signal
});

if (!response.ok) {
if (response.status === 401 && this._authProvider) {
if (response.headers.has('www-authenticate')) {
const { resourceMetadataUrl, scope } = extractWWWAuthenticateParams(response);
this._resourceMetadataUrl = resourceMetadataUrl;
this._scope = scope;
}
if (!response.ok) {
if (response.status === 401 && this._authProvider) {
if (response.headers.has('www-authenticate')) {
const { resourceMetadataUrl, scope } = extractWWWAuthenticateParams(response);
this._resourceMetadataUrl = resourceMetadataUrl;
this._scope = scope;
}

if (this._authProvider.onUnauthorized && !isAuthRetry) {
await this._authProvider.onUnauthorized({
response,
serverUrl: this._url,
fetchFn: this._fetchWithInit
});
await response.text?.().catch(() => {});
// Purposely _not_ awaited, so we don't call onerror twice
return this._startOrAuthSse(options, true);
}
if (this._authProvider.onUnauthorized && !isAuthRetry) {
await this._authProvider.onUnauthorized({
response,
serverUrl: this._url,
fetchFn: this._fetchWithInit
});
await response.text?.().catch(() => {});
if (isAuthRetry) {
throw new SdkError(SdkErrorCode.ClientHttpAuthentication, 'Server returned 401 after re-authentication', {
status: 401
});
}
throw new UnauthorizedError();
return this._startOrAuthSse(options, true);
}

await response.text?.().catch(() => {});

// 405 indicates that the server does not offer an SSE stream at GET endpoint
// This is an expected case that should not trigger an error
if (response.status === 405) {
return;
if (isAuthRetry) {
throw new SdkError(SdkErrorCode.ClientHttpAuthentication, 'Server returned 401 after re-authentication', {
status: 401
});
}
throw new UnauthorizedError();
}

throw new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
status: response.status,
statusText: response.statusText
});
await response.text?.().catch(() => {});

// 405 indicates that the server does not offer an SSE stream at GET endpoint
// This is an expected case that should not trigger an error
if (response.status === 405) {
return;
}

this._handleSseStream(response.body, options, true);
} catch (error) {
this.onerror?.(error as Error);
throw error;
throw new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
status: response.status,
statusText: response.statusText
});
}
Comment thread
claude[bot] marked this conversation as resolved.

this._handleSseStream(response.body, options, true);
}

/**
Expand DownExpand Up@@ -753,9 +747,14 @@ export class StreamableHTTPClientTransport implements Transport {
* @param options Optional callback to receive new resumption tokens
*/
async resumeStream(lastEventId: string, options?: { onresumptiontoken?: (token: string) => void }): Promise<void> {
await this._startOrAuthSse({
resumptionToken: lastEventId,
onresumptiontoken: options?.onresumptiontoken
});
try {
await this._startOrAuthSse({
resumptionToken: lastEventId,
onresumptiontoken: options?.onresumptiontoken
});
} catch (error) {
this.onerror?.(error as Error);
throw error;
}
}
}
59 changes: 59 additions & 0 deletions packages/client/test/client/streamableHttp.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -1224,6 +1224,65 @@ describe('StreamableHTTPClientTransport', () => {
expect(fetchMock.mock.calls[0]![1]?.method).toBe('POST');
});

it('should fire onerror exactly once when _startOrAuthSse fails (not double-fire from catch + caller)', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'));

const errorSpy = vi.fn();
transport.onerror = errorSpy;

const fetchMock = globalThis.fetch as Mock;

// POST returns 202, which triggers _startOrAuthSse when the outbound
// message is an initialized notification (streamableHttp.ts:642)
fetchMock.mockResolvedValueOnce({
ok: true,
status: 202,
headers: new Headers(),
text: async () => ''
});

// The subsequent GET (_startOrAuthSse) fails with a non-ok status
fetchMock.mockResolvedValueOnce({
ok: false,
status: 500,
statusText: 'Internal Server Error',
headers: new Headers(),
text: async () => 'server error'
});

await transport.start();
// Sending an initialized notification triggers the _startOrAuthSse path
await transport.send({ jsonrpc: '2.0', method: 'notifications/initialized' });

// Let the fire-and-forget _startOrAuthSse().catch() settle
await vi.runAllTimersAsync();

expect(errorSpy).toHaveBeenCalledTimes(1);
expect(errorSpy.mock.calls[0]![0].message).toContain('Failed to open SSE stream');
});

it('should fire onerror and reject when resumeStream fails', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'));

const errorSpy = vi.fn();
transport.onerror = errorSpy;

const fetchMock = globalThis.fetch as Mock;
fetchMock.mockResolvedValueOnce({
ok: false,
status: 500,
statusText: 'Internal Server Error',
headers: new Headers(),
text: async () => 'server error'
});

await transport.start();
await expect(transport.resumeStream('event-123')).rejects.toThrow('Failed to open SSE stream');

expect(errorSpy).toHaveBeenCalledTimes(1);
expect(errorSpy.mock.calls[0]![0].message).toContain('Failed to open SSE stream');
});

it('should not throw JSON parse error on priming events with empty data', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'));

Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Remove or un-stick sticky/fixed headers that block content (function() { function unstick() { document.querySelectorAll('header, nav, [role="banner"], .header, .navbar, .sticky, .fixed-top, [style*="position: fixed"], [style*="position:sticky"]').forEach(function(el) { if (el.style.position === 'fixed' || el.style.position === 'sticky' || getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') { el.style.position = 'static'; el.style.top = 'auto'; el.style.zIndex = 'auto'; } }); } unstick(); var observer = new MutationObserver(unstick); observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] }); })(); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' fix(client): remove redundant onerror in _startOrAuthSse catch by felixweinberger · Pull Request #1826 · modelcontextprotocol/typescript-sdk · GitHub
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/fix-double-onerror-startOrAuthSse.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
'@modelcontextprotocol/client': patch
---

Fix double `onerror` invocation when `_startOrAuthSse` fails. The internal catch block fired `onerror` then threw, and all callers already `.catch(onerror)`, causing every failure to fire twice. Removed the redundant internal call.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 The changeset description inaccurately states "all callers already .catch(onerror)" -- resumeStream did NOT have .catch(onerror) before this PR and was silently dropping errors. The PR makes two distinct changes (removing double-fire for three fire-and-forget callers AND adding new error handling to resumeStream), but the changeset only documents the first, leaving the behavioral change to resumeStream undocumented.

Extended reasoning...

What the bug is

The changeset at .changeset/fix-double-onerror-startOrAuthSse.md line 5 reads:

Fix double onerror invocation when _startOrAuthSse fails. The internal catch block fired onerror then threw, and all callers already .catch(onerror), causing every failure to fire twice. Removed the redundant internal call.

The claim "all callers already .catch(onerror)" is factually incorrect. There are four callers of _startOrAuthSse, not three: lines 344, 531/532, 642, and resumeStream (~line 750). The first three are fire-and-forget callers that each chain .catch(error => this.onerror?.(error)). But resumeStream was different -- it simply did await this._startOrAuthSse({...}) with no try-catch and no .catch(onerror) at all.

The specific code path

Before this PR, when resumeStream failed:

  1. _startOrAuthSse threw (after its internal catch called onerror and rethrew)
  2. resumeStream had no handler, so the exception propagated to the caller
  3. onerror was called exactly once -- by the internal catch in _startOrAuthSse

After this PR removes the internal catch from _startOrAuthSse, the PR correctly compensates by adding a try-catch to resumeStream itself (lines 747-758 in the diff). So the code is correct. But the changeset description omits this second change entirely.

Why existing documentation does not cover it

The changeset is the primary artifact consumers read to understand what changed between versions. By describing the fix purely as "removing the redundant internal call" and saying "all callers already .catch(onerror)", a reader would not know that resumeStream now has onerror-reporting behavior it lacked before. Prior to this PR, a resumeStream failure would call onerror once (via the internal catch). After this PR, it still calls onerror once -- but via a different mechanism, and the behavior is preserved only because the PR added a new try-catch. The changeset implies no behavioral change occurred for any caller.

Step-by-step proof

  1. Before this PR: resumeStream calls await this._startOrAuthSse({resumptionToken: lastEventId})
  2. Server returns HTTP 500; _startOrAuthSse throws SdkError(ClientHttpFailedToOpenStream, ...)
  3. _startOrAuthSse internal catch fires: calls this.onerror?.(error), then rethrows
  4. resumeStream has no handler -- exception propagates to caller
  5. onerror fires once (from step 3)
  6. After this PR: step 3 no longer calls onerror internally
  7. The new try-catch in resumeStream catches the error, calls this.onerror?.(error), rethrows
  8. onerror fires once (now from step 7 instead of step 3)
  9. The changeset says only "removed the redundant internal call" -- it does not mention step 7 was added

Impact and fix

The impact is documentation-only -- the code behavior is correct. A changelog reader consulting this entry would not know resumeStream's error-reporting pathway changed. The fix is to update the changeset to accurately describe both changes: (1) removing the redundant double-fire for the three fire-and-forget callers, and (2) adding explicit onerror handling to resumeStream which previously relied on the internal catch that was removed.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The trace in this comment shows resumeStream fired onerror exactly once before and exactly once after, with the error still propagating. Observable behavior is identical, so there's nothing consumer-facing to document. The 'all callers' phrasing describes a private method consumers can't call.

117 changes: 58 additions & 59 deletions packages/client/src/client/streamableHttp.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -233,72 +233,66 @@ export class StreamableHTTPClientTransport implements Transport {
private async _startOrAuthSse(options: StartSSEOptions, isAuthRetry = false): Promise<void> {
const { resumptionToken } = options;

try {
// Try to open an initial SSE stream with GET to listen for server messages
// This is optional according to the spec - server may not support it
const headers = await this._commonHeaders();
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'text/event-stream'];
headers.set('accept', [...new Set(types)].join(', '));

// Include Last-Event-ID header for resumable streams if provided
if (resumptionToken) {
headers.set('last-event-id', resumptionToken);
}
// Try to open an initial SSE stream with GET to listen for server messages
// This is optional according to the spec - server may not support it
const headers = await this._commonHeaders();
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'text/event-stream'];
headers.set('accept', [...new Set(types)].join(', '));

// Include Last-Event-ID header for resumable streams if provided
if (resumptionToken) {
headers.set('last-event-id', resumptionToken);
}

const response = await (this._fetch ?? fetch)(this._url, {
...this._requestInit,
method: 'GET',
headers,
signal: this._abortController?.signal
});
const response = await (this._fetch ?? fetch)(this._url, {
...this._requestInit,
method: 'GET',
headers,
signal: this._abortController?.signal
});

if (!response.ok) {
if (response.status === 401 && this._authProvider) {
if (response.headers.has('www-authenticate')) {
const { resourceMetadataUrl, scope } = extractWWWAuthenticateParams(response);
this._resourceMetadataUrl = resourceMetadataUrl;
this._scope = scope;
}
if (!response.ok) {
if (response.status === 401 && this._authProvider) {
if (response.headers.has('www-authenticate')) {
const { resourceMetadataUrl, scope } = extractWWWAuthenticateParams(response);
this._resourceMetadataUrl = resourceMetadataUrl;
this._scope = scope;
}

if (this._authProvider.onUnauthorized && !isAuthRetry) {
await this._authProvider.onUnauthorized({
response,
serverUrl: this._url,
fetchFn: this._fetchWithInit
});
await response.text?.().catch(() => {});
// Purposely _not_ awaited, so we don't call onerror twice
return this._startOrAuthSse(options, true);
}
if (this._authProvider.onUnauthorized && !isAuthRetry) {
await this._authProvider.onUnauthorized({
response,
serverUrl: this._url,
fetchFn: this._fetchWithInit
});
await response.text?.().catch(() => {});
if (isAuthRetry) {
throw new SdkError(SdkErrorCode.ClientHttpAuthentication, 'Server returned 401 after re-authentication', {
status: 401
});
}
throw new UnauthorizedError();
return this._startOrAuthSse(options, true);
}

await response.text?.().catch(() => {});

// 405 indicates that the server does not offer an SSE stream at GET endpoint
// This is an expected case that should not trigger an error
if (response.status === 405) {
return;
if (isAuthRetry) {
throw new SdkError(SdkErrorCode.ClientHttpAuthentication, 'Server returned 401 after re-authentication', {
status: 401
});
}
throw new UnauthorizedError();
}

throw new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
status: response.status,
statusText: response.statusText
});
await response.text?.().catch(() => {});

// 405 indicates that the server does not offer an SSE stream at GET endpoint
// This is an expected case that should not trigger an error
if (response.status === 405) {
return;
}

this._handleSseStream(response.body, options, true);
} catch (error) {
this.onerror?.(error as Error);
throw error;
throw new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
status: response.status,
statusText: response.statusText
});
}
Comment thread
claude[bot] marked this conversation as resolved.

this._handleSseStream(response.body, options, true);
}

/**
Expand DownExpand Up@@ -753,9 +747,14 @@ export class StreamableHTTPClientTransport implements Transport {
* @param options Optional callback to receive new resumption tokens
*/
async resumeStream(lastEventId: string, options?: { onresumptiontoken?: (token: string) => void }): Promise<void> {
await this._startOrAuthSse({
resumptionToken: lastEventId,
onresumptiontoken: options?.onresumptiontoken
});
try {
await this._startOrAuthSse({
resumptionToken: lastEventId,
onresumptiontoken: options?.onresumptiontoken
});
} catch (error) {
this.onerror?.(error as Error);
throw error;
}
}
}
59 changes: 59 additions & 0 deletions packages/client/test/client/streamableHttp.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -1224,6 +1224,65 @@ describe('StreamableHTTPClientTransport', () => {
expect(fetchMock.mock.calls[0]![1]?.method).toBe('POST');
});

it('should fire onerror exactly once when _startOrAuthSse fails (not double-fire from catch + caller)', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'));

const errorSpy = vi.fn();
transport.onerror = errorSpy;

const fetchMock = globalThis.fetch as Mock;

// POST returns 202, which triggers _startOrAuthSse when the outbound
// message is an initialized notification (streamableHttp.ts:642)
fetchMock.mockResolvedValueOnce({
ok: true,
status: 202,
headers: new Headers(),
text: async () => ''
});

// The subsequent GET (_startOrAuthSse) fails with a non-ok status
fetchMock.mockResolvedValueOnce({
ok: false,
status: 500,
statusText: 'Internal Server Error',
headers: new Headers(),
text: async () => 'server error'
});

await transport.start();
// Sending an initialized notification triggers the _startOrAuthSse path
await transport.send({ jsonrpc: '2.0', method: 'notifications/initialized' });

// Let the fire-and-forget _startOrAuthSse().catch() settle
await vi.runAllTimersAsync();

expect(errorSpy).toHaveBeenCalledTimes(1);
expect(errorSpy.mock.calls[0]![0].message).toContain('Failed to open SSE stream');
});

it('should fire onerror and reject when resumeStream fails', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'));

const errorSpy = vi.fn();
transport.onerror = errorSpy;

const fetchMock = globalThis.fetch as Mock;
fetchMock.mockResolvedValueOnce({
ok: false,
status: 500,
statusText: 'Internal Server Error',
headers: new Headers(),
text: async () => 'server error'
});

await transport.start();
await expect(transport.resumeStream('event-123')).rejects.toThrow('Failed to open SSE stream');

expect(errorSpy).toHaveBeenCalledTimes(1);
expect(errorSpy.mock.calls[0]![0].message).toContain('Failed to open SSE stream');
});

it('should not throw JSON parse error on priming events with empty data', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'));

Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Universal Dark Mode - works on any site (function() { var enabled = true; function applyDarkMode() { if (!enabled) return; // Create style element if it doesn't exist var style = document.getElementById('universal-dark-mode-style'); if (!style) { style = document.createElement('style'); style.id = 'universal-dark-mode-style'; document.head.appendChild(style); } // Dark mode CSS - inverts colors but preserves images/video style.textContent = ' /* Invert everything except media */ html { filter: invert(1) hue-rotate(180deg) !important; background: #1a1a2e !important; } /* Restore images, videos, iframes, canvas */ img, video, iframe, canvas, svg, picture, [style*="background-image"] { filter: invert(1) hue-rotate(180deg) !important; } /* Preserve specific elements that should not be inverted */ .no-dark-mode, .no-dark-mode *, [data-theme="light"], [data-theme="light"], .ace_editor, .ace_editor *, .CodeMirror, .CodeMirror *, .monaco-editor, .monaco-editor *, .markdown-body pre, .markdown-body pre *, .highlight, .highlight *, pre code, pre code * { filter: none !important; } /* Fix common UI elements */ .modal, .popup, .dropdown-menu, .tooltip, .popover { filter: invert(1) hue-rotate(180deg) !important; background: #2d2d44 !important; border-color: #444 !important; } /* Scrollbars */ ::-webkit-scrollbar { background: #1a1a2e !important; } ::-webkit-scrollbar-thumb { background: #444 !important; } ::-webkit-scrollbar-thumb:hover { background: #555 !important; } /* Selection */ ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; } ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; } '; } function removeDarkMode() { var style = document.getElementById('universal-dark-mode-style'); if (style) style.remove(); } // Toggle with Alt+Shift+D document.addEventListener('keydown', function(e) { if (e.altKey && e.shiftKey && e.key === 'D') { e.preventDefault(); enabled = !enabled; if (enabled) { applyDarkMode(); console.log('[Universal Dark Mode] Enabled'); } else { removeDarkMode(); console.log('[Universal Dark Mode] Disabled'); } } }); // Apply on load applyDarkMode(); // Re-apply on dynamic content var observer = new MutationObserver(function(mutations) { if (enabled && !document.getElementById('universal-dark-mode-style')) { applyDarkMode(); } }); observer.observe(document.head, { childList: true }); console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle'); })(); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })(); fix(client): remove redundant onerror in _startOrAuthSse catch by felixweinberger · Pull Request #1826 · modelcontextprotocol/typescript-sdk · GitHub
Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/fix-double-onerror-startOrAuthSse.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
'@modelcontextprotocol/client': patch
---

Fix double `onerror` invocation when `_startOrAuthSse` fails. The internal catch block fired `onerror` then threw, and all callers already `.catch(onerror)`, causing every failure to fire twice. Removed the redundant internal call.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 The changeset description inaccurately states "all callers already .catch(onerror)" -- resumeStream did NOT have .catch(onerror) before this PR and was silently dropping errors. The PR makes two distinct changes (removing double-fire for three fire-and-forget callers AND adding new error handling to resumeStream), but the changeset only documents the first, leaving the behavioral change to resumeStream undocumented.

Extended reasoning...

What the bug is

The changeset at .changeset/fix-double-onerror-startOrAuthSse.md line 5 reads:

Fix double onerror invocation when _startOrAuthSse fails. The internal catch block fired onerror then threw, and all callers already .catch(onerror), causing every failure to fire twice. Removed the redundant internal call.

The claim "all callers already .catch(onerror)" is factually incorrect. There are four callers of _startOrAuthSse, not three: lines 344, 531/532, 642, and resumeStream (~line 750). The first three are fire-and-forget callers that each chain .catch(error => this.onerror?.(error)). But resumeStream was different -- it simply did await this._startOrAuthSse({...}) with no try-catch and no .catch(onerror) at all.

The specific code path

Before this PR, when resumeStream failed:

  1. _startOrAuthSse threw (after its internal catch called onerror and rethrew)
  2. resumeStream had no handler, so the exception propagated to the caller
  3. onerror was called exactly once -- by the internal catch in _startOrAuthSse

After this PR removes the internal catch from _startOrAuthSse, the PR correctly compensates by adding a try-catch to resumeStream itself (lines 747-758 in the diff). So the code is correct. But the changeset description omits this second change entirely.

Why existing documentation does not cover it

The changeset is the primary artifact consumers read to understand what changed between versions. By describing the fix purely as "removing the redundant internal call" and saying "all callers already .catch(onerror)", a reader would not know that resumeStream now has onerror-reporting behavior it lacked before. Prior to this PR, a resumeStream failure would call onerror once (via the internal catch). After this PR, it still calls onerror once -- but via a different mechanism, and the behavior is preserved only because the PR added a new try-catch. The changeset implies no behavioral change occurred for any caller.

Step-by-step proof

  1. Before this PR: resumeStream calls await this._startOrAuthSse({resumptionToken: lastEventId})
  2. Server returns HTTP 500; _startOrAuthSse throws SdkError(ClientHttpFailedToOpenStream, ...)
  3. _startOrAuthSse internal catch fires: calls this.onerror?.(error), then rethrows
  4. resumeStream has no handler -- exception propagates to caller
  5. onerror fires once (from step 3)
  6. After this PR: step 3 no longer calls onerror internally
  7. The new try-catch in resumeStream catches the error, calls this.onerror?.(error), rethrows
  8. onerror fires once (now from step 7 instead of step 3)
  9. The changeset says only "removed the redundant internal call" -- it does not mention step 7 was added

Impact and fix

The impact is documentation-only -- the code behavior is correct. A changelog reader consulting this entry would not know resumeStream's error-reporting pathway changed. The fix is to update the changeset to accurately describe both changes: (1) removing the redundant double-fire for the three fire-and-forget callers, and (2) adding explicit onerror handling to resumeStream which previously relied on the internal catch that was removed.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The trace in this comment shows resumeStream fired onerror exactly once before and exactly once after, with the error still propagating. Observable behavior is identical, so there's nothing consumer-facing to document. The 'all callers' phrasing describes a private method consumers can't call.

117 changes: 58 additions & 59 deletions packages/client/src/client/streamableHttp.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -233,72 +233,66 @@ export class StreamableHTTPClientTransport implements Transport {
private async _startOrAuthSse(options: StartSSEOptions, isAuthRetry = false): Promise<void> {
const { resumptionToken } = options;

try {
// Try to open an initial SSE stream with GET to listen for server messages
// This is optional according to the spec - server may not support it
const headers = await this._commonHeaders();
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'text/event-stream'];
headers.set('accept', [...new Set(types)].join(', '));

// Include Last-Event-ID header for resumable streams if provided
if (resumptionToken) {
headers.set('last-event-id', resumptionToken);
}
// Try to open an initial SSE stream with GET to listen for server messages
// This is optional according to the spec - server may not support it
const headers = await this._commonHeaders();
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'text/event-stream'];
headers.set('accept', [...new Set(types)].join(', '));

// Include Last-Event-ID header for resumable streams if provided
if (resumptionToken) {
headers.set('last-event-id', resumptionToken);
}

const response = await (this._fetch ?? fetch)(this._url, {
...this._requestInit,
method: 'GET',
headers,
signal: this._abortController?.signal
});
const response = await (this._fetch ?? fetch)(this._url, {
...this._requestInit,
method: 'GET',
headers,
signal: this._abortController?.signal
});

if (!response.ok) {
if (response.status === 401 && this._authProvider) {
if (response.headers.has('www-authenticate')) {
const { resourceMetadataUrl, scope } = extractWWWAuthenticateParams(response);
this._resourceMetadataUrl = resourceMetadataUrl;
this._scope = scope;
}
if (!response.ok) {
if (response.status === 401 && this._authProvider) {
if (response.headers.has('www-authenticate')) {
const { resourceMetadataUrl, scope } = extractWWWAuthenticateParams(response);
this._resourceMetadataUrl = resourceMetadataUrl;
this._scope = scope;
}

if (this._authProvider.onUnauthorized && !isAuthRetry) {
await this._authProvider.onUnauthorized({
response,
serverUrl: this._url,
fetchFn: this._fetchWithInit
});
await response.text?.().catch(() => {});
// Purposely _not_ awaited, so we don't call onerror twice
return this._startOrAuthSse(options, true);
}
if (this._authProvider.onUnauthorized && !isAuthRetry) {
await this._authProvider.onUnauthorized({
response,
serverUrl: this._url,
fetchFn: this._fetchWithInit
});
await response.text?.().catch(() => {});
if (isAuthRetry) {
throw new SdkError(SdkErrorCode.ClientHttpAuthentication, 'Server returned 401 after re-authentication', {
status: 401
});
}
throw new UnauthorizedError();
return this._startOrAuthSse(options, true);
}

await response.text?.().catch(() => {});

// 405 indicates that the server does not offer an SSE stream at GET endpoint
// This is an expected case that should not trigger an error
if (response.status === 405) {
return;
if (isAuthRetry) {
throw new SdkError(SdkErrorCode.ClientHttpAuthentication, 'Server returned 401 after re-authentication', {
status: 401
});
}
throw new UnauthorizedError();
}

throw new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
status: response.status,
statusText: response.statusText
});
await response.text?.().catch(() => {});

// 405 indicates that the server does not offer an SSE stream at GET endpoint
// This is an expected case that should not trigger an error
if (response.status === 405) {
return;
}

this._handleSseStream(response.body, options, true);
} catch (error) {
this.onerror?.(error as Error);
throw error;
throw new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
status: response.status,
statusText: response.statusText
});
}
Comment thread
claude[bot] marked this conversation as resolved.

this._handleSseStream(response.body, options, true);
}

/**
Expand DownExpand Up@@ -753,9 +747,14 @@ export class StreamableHTTPClientTransport implements Transport {
* @param options Optional callback to receive new resumption tokens
*/
async resumeStream(lastEventId: string, options?: { onresumptiontoken?: (token: string) => void }): Promise<void> {
await this._startOrAuthSse({
resumptionToken: lastEventId,
onresumptiontoken: options?.onresumptiontoken
});
try {
await this._startOrAuthSse({
resumptionToken: lastEventId,
onresumptiontoken: options?.onresumptiontoken
});
} catch (error) {
this.onerror?.(error as Error);
throw error;
}
}
}
59 changes: 59 additions & 0 deletions packages/client/test/client/streamableHttp.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -1224,6 +1224,65 @@ describe('StreamableHTTPClientTransport', () => {
expect(fetchMock.mock.calls[0]![1]?.method).toBe('POST');
});

it('should fire onerror exactly once when _startOrAuthSse fails (not double-fire from catch + caller)', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'));

const errorSpy = vi.fn();
transport.onerror = errorSpy;

const fetchMock = globalThis.fetch as Mock;

// POST returns 202, which triggers _startOrAuthSse when the outbound
// message is an initialized notification (streamableHttp.ts:642)
fetchMock.mockResolvedValueOnce({
ok: true,
status: 202,
headers: new Headers(),
text: async () => ''
});

// The subsequent GET (_startOrAuthSse) fails with a non-ok status
fetchMock.mockResolvedValueOnce({
ok: false,
status: 500,
statusText: 'Internal Server Error',
headers: new Headers(),
text: async () => 'server error'
});

await transport.start();
// Sending an initialized notification triggers the _startOrAuthSse path
await transport.send({ jsonrpc: '2.0', method: 'notifications/initialized' });

// Let the fire-and-forget _startOrAuthSse().catch() settle
await vi.runAllTimersAsync();

expect(errorSpy).toHaveBeenCalledTimes(1);
expect(errorSpy.mock.calls[0]![0].message).toContain('Failed to open SSE stream');
});

it('should fire onerror and reject when resumeStream fails', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'));

const errorSpy = vi.fn();
transport.onerror = errorSpy;

const fetchMock = globalThis.fetch as Mock;
fetchMock.mockResolvedValueOnce({
ok: false,
status: 500,
statusText: 'Internal Server Error',
headers: new Headers(),
text: async () => 'server error'
});

await transport.start();
await expect(transport.resumeStream('event-123')).rejects.toThrow('Failed to open SSE stream');

expect(errorSpy).toHaveBeenCalledTimes(1);
expect(errorSpy.mock.calls[0]![0].message).toContain('Failed to open SSE stream');
});

it('should not throw JSON parse error on priming events with empty data', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'));

Expand Down
Loading