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/gentle-maps-smile.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
'@modelcontextprotocol/client': patch
---

Clear stale Streamable HTTP client sessions when a session-bound request receives HTTP 404 by clearing the stored session ID, so the next initialize flow can proceed without an MCP session header.
33 changes: 31 additions & 2 deletions packages/client/src/client/streamableHttp.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,6 +25,19 @@ const DEFAULT_STREAMABLE_HTTP_RECONNECTION_OPTIONS: StreamableHTTPReconnectionOp
maxRetries: 2
};

const SESSION_BOUND_404_ERROR = Symbol('sessionBound404Error');

type SessionBound404Error = Error & { [SESSION_BOUND_404_ERROR]?: true };

function markSessionBound404Error(error: Error): Error {
(error as SessionBound404Error)[SESSION_BOUND_404_ERROR] = true;
return error;
}

function isSessionBound404Error(error: unknown): boolean {
return Boolean(error && typeof error === 'object' && (error as SessionBound404Error)[SESSION_BOUND_404_ERROR] === true);
}

/**
* Options for starting or authenticating an SSE connection
*/
Expand DownExpand Up@@ -237,6 +250,7 @@ export class StreamableHTTPClientTransport implements Transport {
// 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 sentSessionId = headers.get('mcp-session-id');
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'text/event-stream'];
headers.set('accept', [...new Set(types)].join(', '));
Expand All@@ -254,6 +268,11 @@ export class StreamableHTTPClientTransport implements Transport {
});

if (!response.ok) {
const shouldClearSessionFor404 = response.status === 404 && sentSessionId !== null && this._sessionId === sentSessionId;
if (shouldClearSessionFor404) {
this._sessionId = undefined;
}

if (response.status === 401 && this._authProvider) {
if (response.headers.has('www-authenticate')) {
const { resourceMetadataUrl, scope } = extractWWWAuthenticateParams(response);
Comment thread
Maverick-666 marked this conversation as resolved.
Expand DownExpand Up@@ -288,10 +307,11 @@ export class StreamableHTTPClientTransport implements Transport {
return;
}

throw new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
const error = new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
status: response.status,
statusText: response.statusText
});
throw shouldClearSessionFor404 ? markSessionBound404Error(error) : error;
}

this._handleSseStream(response.body, options, true);
Expand DownExpand Up@@ -345,7 +365,11 @@ export class StreamableHTTPClientTransport implements Transport {
this._cancelReconnection = undefined;
if (this._abortController?.signal.aborted) return;
this._startOrAuthSse(options).catch(error => {
this.onerror?.(new Error(`Failed to reconnect SSE stream: ${error instanceof Error ? error.message : String(error)}`));
const reconnectError = error instanceof Error ? error : new Error(String(error));
this.onerror?.(new Error(`Failed to reconnect SSE stream: ${reconnectError.message}`));
if (isSessionBound404Error(reconnectError)) {
return;
}
try {
this._scheduleReconnection(options, attemptCount + 1);
} catch (scheduleError) {
Expand DownExpand Up@@ -539,6 +563,7 @@ export class StreamableHTTPClientTransport implements Transport {
}

const headers = await this._commonHeaders();
const sentSessionId = headers.get('mcp-session-id');
headers.set('content-type', 'application/json');
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'application/json', 'text/event-stream'];
Expand All@@ -561,6 +586,10 @@ export class StreamableHTTPClientTransport implements Transport {
}

if (!response.ok) {
if (response.status === 404 && sentSessionId !== null && this._sessionId === sentSessionId) {
this._sessionId = undefined;
}
Comment thread
Maverick-666 marked this conversation as resolved.

if (response.status === 401 && this._authProvider) {
// Store WWW-Authenticate params for interactive finishAuth() path
if (response.headers.has('www-authenticate')) {
Comment thread
Maverick-666 marked this conversation as resolved.
Expand Down
203 changes: 202 additions & 1 deletion packages/client/test/client/streamableHttp.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -220,7 +220,7 @@ describe('StreamableHTTPClientTransport', () => {
await expect(transport.terminateSession()).resolves.not.toThrow();
});

it('should handle 404 response when session expires', async () => {
it('should preserve existing 404 behavior when request is not session-bound', async () => {
const message: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'test',
Expand DownExpand Up@@ -248,6 +248,104 @@ describe('StreamableHTTPClientTransport', () => {
expect(errorSpy).toHaveBeenCalled();
});

it('should clear session ID on 404 for session-bound POST requests', async () => {
const initializeMessage: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'initialize',
params: {
clientInfo: { name: 'test-client', version: '1.0' },
protocolVersion: '2025-03-26'
},
id: 'init-id'
};
const message: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'tools/list',
params: {},
id: 'test-id'
};

(globalThis.fetch as Mock)
.mockResolvedValueOnce({
ok: true,
status: 202,
headers: new Headers({ 'mcp-session-id': 'stale-session-id' }),
text: () => Promise.resolve('')
})
.mockResolvedValueOnce({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
})
.mockResolvedValueOnce({
ok: true,
status: 202,
headers: new Headers(),
text: () => Promise.resolve('')
});

await transport.send(initializeMessage);
expect(transport.sessionId).toBe('stale-session-id');

await expect(transport.send(message)).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpNotImplemented,
data: expect.objectContaining({
status: 404,
text: 'Session not found'
})
});
expect(transport.sessionId).toBeUndefined();

await transport.send({ jsonrpc: '2.0', method: 'notifications/ping' } as JSONRPCMessage);
const lastCall = (globalThis.fetch as Mock).mock.calls.at(-1)!;
expect(lastCall[1].headers.get('mcp-session-id')).toBeNull();
});

it('should not clear a newer session ID when a stale session-bound POST request returns 404', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-A'
});

const message: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'tools/list',
params: {},
id: 'test-id'
};

let resolveFetch!: (value: unknown) => void;
const deferredFetch = new Promise(resolve => {
resolveFetch = resolve;
});

(globalThis.fetch as Mock).mockImplementationOnce(() => {
// Simulate another in-flight request establishing a fresh session while this request is pending.
(transport as unknown as { _sessionId?: string })._sessionId = 'fresh-session-B';
return deferredFetch;
});

const sendPromise = transport.send(message);

resolveFetch({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
});

await expect(sendPromise).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpNotImplemented,
data: expect.objectContaining({
status: 404
})
});

expect(transport.sessionId).toBe('fresh-session-B');
});

it('should handle non-streaming JSON response', async () => {
const message: JSONRPCMessage = {
jsonrpc: '2.0',
Expand DownExpand Up@@ -309,6 +407,75 @@ describe('StreamableHTTPClientTransport', () => {
expect(globalThis.fetch).toHaveBeenCalledTimes(2);
});

it('should clear session ID when GET SSE stream returns 404 for a session-bound request', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-id'
});
await transport.start();

(globalThis.fetch as Mock).mockResolvedValueOnce({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
});

await expect(
(transport as unknown as { _startOrAuthSse: (opts: StartSSEOptions) => Promise<void> })._startOrAuthSse({})
).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpFailedToOpenStream,
data: expect.objectContaining({
status: 404,
statusText: 'Not Found'
})
});

expect(transport.sessionId).toBeUndefined();

const getCall = (globalThis.fetch as Mock).mock.calls[0]!;
expect(getCall[1].method).toBe('GET');
expect(getCall[1].headers.get('mcp-session-id')).toBe('stale-session-id');
});

it('should not clear a newer session ID when a stale session-bound GET request returns 404', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-A'
});
await transport.start();

let resolveFetch!: (value: unknown) => void;
const deferredFetch = new Promise(resolve => {
resolveFetch = resolve;
});

(globalThis.fetch as Mock).mockImplementationOnce(() => {
// Simulate another in-flight request establishing a fresh session while this request is pending.
(transport as unknown as { _sessionId?: string })._sessionId = 'fresh-session-B';
return deferredFetch;
});

const startPromise = (transport as unknown as { _startOrAuthSse: (opts: StartSSEOptions) => Promise<void> })._startOrAuthSse({});

resolveFetch({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
});

await expect(startPromise).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpFailedToOpenStream,
data: expect.objectContaining({
status: 404,
statusText: 'Not Found'
})
});

expect(transport.sessionId).toBe('fresh-session-B');
});

it('should handle successful initial GET connection for SSE', async () => {
// Set up readable stream for SSE events
const encoder = new TextEncoder();
Expand DownExpand Up@@ -936,6 +1103,40 @@ describe('StreamableHTTPClientTransport', () => {
expect(fetchMock.mock.calls[1]![1]?.method).toBe('GET');
});

it('should stop retrying GET reconnection after a session-bound 404 clears the stale session', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-id',
reconnectionOptions: {
initialReconnectionDelay: 10,
maxRetries: 3,
maxReconnectionDelay: 1000,
reconnectionDelayGrowFactor: 1
}
});
await transport.start();

const fetchMock = globalThis.fetch as Mock;
fetchMock.mockResolvedValue({
ok: false,
status: 404,
statusText: 'Not Found',
headers: new Headers(),
text: () => Promise.resolve('Session not found')
});

(
transport as unknown as {
_scheduleReconnection: (opts: StartSSEOptions, attemptCount?: number) => void;
}
)._scheduleReconnection({}, 0);

await vi.advanceTimersByTimeAsync(20);
await vi.advanceTimersByTimeAsync(100);

expect(fetchMock).toHaveBeenCalledTimes(1);
expect(transport.sessionId).toBeUndefined();
});

it('should NOT reconnect a POST-initiated stream that fails', async () => {
// ARRANGE
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" + '
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/gentle-maps-smile.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
'@modelcontextprotocol/client': patch
---

Clear stale Streamable HTTP client sessions when a session-bound request receives HTTP 404 by clearing the stored session ID, so the next initialize flow can proceed without an MCP session header.
33 changes: 31 additions & 2 deletions packages/client/src/client/streamableHttp.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,6 +25,19 @@ const DEFAULT_STREAMABLE_HTTP_RECONNECTION_OPTIONS: StreamableHTTPReconnectionOp
maxRetries: 2
};

const SESSION_BOUND_404_ERROR = Symbol('sessionBound404Error');

type SessionBound404Error = Error & { [SESSION_BOUND_404_ERROR]?: true };

function markSessionBound404Error(error: Error): Error {
(error as SessionBound404Error)[SESSION_BOUND_404_ERROR] = true;
return error;
}

function isSessionBound404Error(error: unknown): boolean {
return Boolean(error && typeof error === 'object' && (error as SessionBound404Error)[SESSION_BOUND_404_ERROR] === true);
}

/**
* Options for starting or authenticating an SSE connection
*/
Expand DownExpand Up@@ -237,6 +250,7 @@ export class StreamableHTTPClientTransport implements Transport {
// 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 sentSessionId = headers.get('mcp-session-id');
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'text/event-stream'];
headers.set('accept', [...new Set(types)].join(', '));
Expand All@@ -254,6 +268,11 @@ export class StreamableHTTPClientTransport implements Transport {
});

if (!response.ok) {
const shouldClearSessionFor404 = response.status === 404 && sentSessionId !== null && this._sessionId === sentSessionId;
if (shouldClearSessionFor404) {
this._sessionId = undefined;
}

if (response.status === 401 && this._authProvider) {
if (response.headers.has('www-authenticate')) {
const { resourceMetadataUrl, scope } = extractWWWAuthenticateParams(response);
Comment thread
Maverick-666 marked this conversation as resolved.
Expand DownExpand Up@@ -288,10 +307,11 @@ export class StreamableHTTPClientTransport implements Transport {
return;
}

throw new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
const error = new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
status: response.status,
statusText: response.statusText
});
throw shouldClearSessionFor404 ? markSessionBound404Error(error) : error;
}

this._handleSseStream(response.body, options, true);
Expand DownExpand Up@@ -345,7 +365,11 @@ export class StreamableHTTPClientTransport implements Transport {
this._cancelReconnection = undefined;
if (this._abortController?.signal.aborted) return;
this._startOrAuthSse(options).catch(error => {
this.onerror?.(new Error(`Failed to reconnect SSE stream: ${error instanceof Error ? error.message : String(error)}`));
const reconnectError = error instanceof Error ? error : new Error(String(error));
this.onerror?.(new Error(`Failed to reconnect SSE stream: ${reconnectError.message}`));
if (isSessionBound404Error(reconnectError)) {
return;
}
try {
this._scheduleReconnection(options, attemptCount + 1);
} catch (scheduleError) {
Expand DownExpand Up@@ -539,6 +563,7 @@ export class StreamableHTTPClientTransport implements Transport {
}

const headers = await this._commonHeaders();
const sentSessionId = headers.get('mcp-session-id');
headers.set('content-type', 'application/json');
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'application/json', 'text/event-stream'];
Expand All@@ -561,6 +586,10 @@ export class StreamableHTTPClientTransport implements Transport {
}

if (!response.ok) {
if (response.status === 404 && sentSessionId !== null && this._sessionId === sentSessionId) {
this._sessionId = undefined;
}
Comment thread
Maverick-666 marked this conversation as resolved.

if (response.status === 401 && this._authProvider) {
// Store WWW-Authenticate params for interactive finishAuth() path
if (response.headers.has('www-authenticate')) {
Comment thread
Maverick-666 marked this conversation as resolved.
Expand Down
203 changes: 202 additions & 1 deletion packages/client/test/client/streamableHttp.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -220,7 +220,7 @@ describe('StreamableHTTPClientTransport', () => {
await expect(transport.terminateSession()).resolves.not.toThrow();
});

it('should handle 404 response when session expires', async () => {
it('should preserve existing 404 behavior when request is not session-bound', async () => {
const message: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'test',
Expand DownExpand Up@@ -248,6 +248,104 @@ describe('StreamableHTTPClientTransport', () => {
expect(errorSpy).toHaveBeenCalled();
});

it('should clear session ID on 404 for session-bound POST requests', async () => {
const initializeMessage: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'initialize',
params: {
clientInfo: { name: 'test-client', version: '1.0' },
protocolVersion: '2025-03-26'
},
id: 'init-id'
};
const message: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'tools/list',
params: {},
id: 'test-id'
};

(globalThis.fetch as Mock)
.mockResolvedValueOnce({
ok: true,
status: 202,
headers: new Headers({ 'mcp-session-id': 'stale-session-id' }),
text: () => Promise.resolve('')
})
.mockResolvedValueOnce({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
})
.mockResolvedValueOnce({
ok: true,
status: 202,
headers: new Headers(),
text: () => Promise.resolve('')
});

await transport.send(initializeMessage);
expect(transport.sessionId).toBe('stale-session-id');

await expect(transport.send(message)).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpNotImplemented,
data: expect.objectContaining({
status: 404,
text: 'Session not found'
})
});
expect(transport.sessionId).toBeUndefined();

await transport.send({ jsonrpc: '2.0', method: 'notifications/ping' } as JSONRPCMessage);
const lastCall = (globalThis.fetch as Mock).mock.calls.at(-1)!;
expect(lastCall[1].headers.get('mcp-session-id')).toBeNull();
});

it('should not clear a newer session ID when a stale session-bound POST request returns 404', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-A'
});

const message: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'tools/list',
params: {},
id: 'test-id'
};

let resolveFetch!: (value: unknown) => void;
const deferredFetch = new Promise(resolve => {
resolveFetch = resolve;
});

(globalThis.fetch as Mock).mockImplementationOnce(() => {
// Simulate another in-flight request establishing a fresh session while this request is pending.
(transport as unknown as { _sessionId?: string })._sessionId = 'fresh-session-B';
return deferredFetch;
});

const sendPromise = transport.send(message);

resolveFetch({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
});

await expect(sendPromise).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpNotImplemented,
data: expect.objectContaining({
status: 404
})
});

expect(transport.sessionId).toBe('fresh-session-B');
});

it('should handle non-streaming JSON response', async () => {
const message: JSONRPCMessage = {
jsonrpc: '2.0',
Expand DownExpand Up@@ -309,6 +407,75 @@ describe('StreamableHTTPClientTransport', () => {
expect(globalThis.fetch).toHaveBeenCalledTimes(2);
});

it('should clear session ID when GET SSE stream returns 404 for a session-bound request', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-id'
});
await transport.start();

(globalThis.fetch as Mock).mockResolvedValueOnce({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
});

await expect(
(transport as unknown as { _startOrAuthSse: (opts: StartSSEOptions) => Promise<void> })._startOrAuthSse({})
).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpFailedToOpenStream,
data: expect.objectContaining({
status: 404,
statusText: 'Not Found'
})
});

expect(transport.sessionId).toBeUndefined();

const getCall = (globalThis.fetch as Mock).mock.calls[0]!;
expect(getCall[1].method).toBe('GET');
expect(getCall[1].headers.get('mcp-session-id')).toBe('stale-session-id');
});

it('should not clear a newer session ID when a stale session-bound GET request returns 404', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-A'
});
await transport.start();

let resolveFetch!: (value: unknown) => void;
const deferredFetch = new Promise(resolve => {
resolveFetch = resolve;
});

(globalThis.fetch as Mock).mockImplementationOnce(() => {
// Simulate another in-flight request establishing a fresh session while this request is pending.
(transport as unknown as { _sessionId?: string })._sessionId = 'fresh-session-B';
return deferredFetch;
});

const startPromise = (transport as unknown as { _startOrAuthSse: (opts: StartSSEOptions) => Promise<void> })._startOrAuthSse({});

resolveFetch({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
});

await expect(startPromise).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpFailedToOpenStream,
data: expect.objectContaining({
status: 404,
statusText: 'Not Found'
})
});

expect(transport.sessionId).toBe('fresh-session-B');
});

it('should handle successful initial GET connection for SSE', async () => {
// Set up readable stream for SSE events
const encoder = new TextEncoder();
Expand DownExpand Up@@ -936,6 +1103,40 @@ describe('StreamableHTTPClientTransport', () => {
expect(fetchMock.mock.calls[1]![1]?.method).toBe('GET');
});

it('should stop retrying GET reconnection after a session-bound 404 clears the stale session', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-id',
reconnectionOptions: {
initialReconnectionDelay: 10,
maxRetries: 3,
maxReconnectionDelay: 1000,
reconnectionDelayGrowFactor: 1
}
});
await transport.start();

const fetchMock = globalThis.fetch as Mock;
fetchMock.mockResolvedValue({
ok: false,
status: 404,
statusText: 'Not Found',
headers: new Headers(),
text: () => Promise.resolve('Session not found')
});

(
transport as unknown as {
_scheduleReconnection: (opts: StartSSEOptions, attemptCount?: number) => void;
}
)._scheduleReconnection({}, 0);

await vi.advanceTimersByTimeAsync(20);
await vi.advanceTimersByTimeAsync(100);

expect(fetchMock).toHaveBeenCalledTimes(1);
expect(transport.sessionId).toBeUndefined();
});

it('should NOT reconnect a POST-initiated stream that fails', async () => {
// ARRANGE
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('^' + ".*" + '
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/gentle-maps-smile.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
'@modelcontextprotocol/client': patch
---

Clear stale Streamable HTTP client sessions when a session-bound request receives HTTP 404 by clearing the stored session ID, so the next initialize flow can proceed without an MCP session header.
33 changes: 31 additions & 2 deletions packages/client/src/client/streamableHttp.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,6 +25,19 @@ const DEFAULT_STREAMABLE_HTTP_RECONNECTION_OPTIONS: StreamableHTTPReconnectionOp
maxRetries: 2
};

const SESSION_BOUND_404_ERROR = Symbol('sessionBound404Error');

type SessionBound404Error = Error & { [SESSION_BOUND_404_ERROR]?: true };

function markSessionBound404Error(error: Error): Error {
(error as SessionBound404Error)[SESSION_BOUND_404_ERROR] = true;
return error;
}

function isSessionBound404Error(error: unknown): boolean {
return Boolean(error && typeof error === 'object' && (error as SessionBound404Error)[SESSION_BOUND_404_ERROR] === true);
}

/**
* Options for starting or authenticating an SSE connection
*/
Expand DownExpand Up@@ -237,6 +250,7 @@ export class StreamableHTTPClientTransport implements Transport {
// 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 sentSessionId = headers.get('mcp-session-id');
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'text/event-stream'];
headers.set('accept', [...new Set(types)].join(', '));
Expand All@@ -254,6 +268,11 @@ export class StreamableHTTPClientTransport implements Transport {
});

if (!response.ok) {
const shouldClearSessionFor404 = response.status === 404 && sentSessionId !== null && this._sessionId === sentSessionId;
if (shouldClearSessionFor404) {
this._sessionId = undefined;
}

if (response.status === 401 && this._authProvider) {
if (response.headers.has('www-authenticate')) {
const { resourceMetadataUrl, scope } = extractWWWAuthenticateParams(response);
Comment thread
Maverick-666 marked this conversation as resolved.
Expand DownExpand Up@@ -288,10 +307,11 @@ export class StreamableHTTPClientTransport implements Transport {
return;
}

throw new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
const error = new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
status: response.status,
statusText: response.statusText
});
throw shouldClearSessionFor404 ? markSessionBound404Error(error) : error;
}

this._handleSseStream(response.body, options, true);
Expand DownExpand Up@@ -345,7 +365,11 @@ export class StreamableHTTPClientTransport implements Transport {
this._cancelReconnection = undefined;
if (this._abortController?.signal.aborted) return;
this._startOrAuthSse(options).catch(error => {
this.onerror?.(new Error(`Failed to reconnect SSE stream: ${error instanceof Error ? error.message : String(error)}`));
const reconnectError = error instanceof Error ? error : new Error(String(error));
this.onerror?.(new Error(`Failed to reconnect SSE stream: ${reconnectError.message}`));
if (isSessionBound404Error(reconnectError)) {
return;
}
try {
this._scheduleReconnection(options, attemptCount + 1);
} catch (scheduleError) {
Expand DownExpand Up@@ -539,6 +563,7 @@ export class StreamableHTTPClientTransport implements Transport {
}

const headers = await this._commonHeaders();
const sentSessionId = headers.get('mcp-session-id');
headers.set('content-type', 'application/json');
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'application/json', 'text/event-stream'];
Expand All@@ -561,6 +586,10 @@ export class StreamableHTTPClientTransport implements Transport {
}

if (!response.ok) {
if (response.status === 404 && sentSessionId !== null && this._sessionId === sentSessionId) {
this._sessionId = undefined;
}
Comment thread
Maverick-666 marked this conversation as resolved.

if (response.status === 401 && this._authProvider) {
// Store WWW-Authenticate params for interactive finishAuth() path
if (response.headers.has('www-authenticate')) {
Comment thread
Maverick-666 marked this conversation as resolved.
Expand Down
203 changes: 202 additions & 1 deletion packages/client/test/client/streamableHttp.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -220,7 +220,7 @@ describe('StreamableHTTPClientTransport', () => {
await expect(transport.terminateSession()).resolves.not.toThrow();
});

it('should handle 404 response when session expires', async () => {
it('should preserve existing 404 behavior when request is not session-bound', async () => {
const message: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'test',
Expand DownExpand Up@@ -248,6 +248,104 @@ describe('StreamableHTTPClientTransport', () => {
expect(errorSpy).toHaveBeenCalled();
});

it('should clear session ID on 404 for session-bound POST requests', async () => {
const initializeMessage: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'initialize',
params: {
clientInfo: { name: 'test-client', version: '1.0' },
protocolVersion: '2025-03-26'
},
id: 'init-id'
};
const message: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'tools/list',
params: {},
id: 'test-id'
};

(globalThis.fetch as Mock)
.mockResolvedValueOnce({
ok: true,
status: 202,
headers: new Headers({ 'mcp-session-id': 'stale-session-id' }),
text: () => Promise.resolve('')
})
.mockResolvedValueOnce({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
})
.mockResolvedValueOnce({
ok: true,
status: 202,
headers: new Headers(),
text: () => Promise.resolve('')
});

await transport.send(initializeMessage);
expect(transport.sessionId).toBe('stale-session-id');

await expect(transport.send(message)).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpNotImplemented,
data: expect.objectContaining({
status: 404,
text: 'Session not found'
})
});
expect(transport.sessionId).toBeUndefined();

await transport.send({ jsonrpc: '2.0', method: 'notifications/ping' } as JSONRPCMessage);
const lastCall = (globalThis.fetch as Mock).mock.calls.at(-1)!;
expect(lastCall[1].headers.get('mcp-session-id')).toBeNull();
});

it('should not clear a newer session ID when a stale session-bound POST request returns 404', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-A'
});

const message: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'tools/list',
params: {},
id: 'test-id'
};

let resolveFetch!: (value: unknown) => void;
const deferredFetch = new Promise(resolve => {
resolveFetch = resolve;
});

(globalThis.fetch as Mock).mockImplementationOnce(() => {
// Simulate another in-flight request establishing a fresh session while this request is pending.
(transport as unknown as { _sessionId?: string })._sessionId = 'fresh-session-B';
return deferredFetch;
});

const sendPromise = transport.send(message);

resolveFetch({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
});

await expect(sendPromise).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpNotImplemented,
data: expect.objectContaining({
status: 404
})
});

expect(transport.sessionId).toBe('fresh-session-B');
});

it('should handle non-streaming JSON response', async () => {
const message: JSONRPCMessage = {
jsonrpc: '2.0',
Expand DownExpand Up@@ -309,6 +407,75 @@ describe('StreamableHTTPClientTransport', () => {
expect(globalThis.fetch).toHaveBeenCalledTimes(2);
});

it('should clear session ID when GET SSE stream returns 404 for a session-bound request', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-id'
});
await transport.start();

(globalThis.fetch as Mock).mockResolvedValueOnce({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
});

await expect(
(transport as unknown as { _startOrAuthSse: (opts: StartSSEOptions) => Promise<void> })._startOrAuthSse({})
).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpFailedToOpenStream,
data: expect.objectContaining({
status: 404,
statusText: 'Not Found'
})
});

expect(transport.sessionId).toBeUndefined();

const getCall = (globalThis.fetch as Mock).mock.calls[0]!;
expect(getCall[1].method).toBe('GET');
expect(getCall[1].headers.get('mcp-session-id')).toBe('stale-session-id');
});

it('should not clear a newer session ID when a stale session-bound GET request returns 404', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-A'
});
await transport.start();

let resolveFetch!: (value: unknown) => void;
const deferredFetch = new Promise(resolve => {
resolveFetch = resolve;
});

(globalThis.fetch as Mock).mockImplementationOnce(() => {
// Simulate another in-flight request establishing a fresh session while this request is pending.
(transport as unknown as { _sessionId?: string })._sessionId = 'fresh-session-B';
return deferredFetch;
});

const startPromise = (transport as unknown as { _startOrAuthSse: (opts: StartSSEOptions) => Promise<void> })._startOrAuthSse({});

resolveFetch({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
});

await expect(startPromise).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpFailedToOpenStream,
data: expect.objectContaining({
status: 404,
statusText: 'Not Found'
})
});

expect(transport.sessionId).toBe('fresh-session-B');
});

it('should handle successful initial GET connection for SSE', async () => {
// Set up readable stream for SSE events
const encoder = new TextEncoder();
Expand DownExpand Up@@ -936,6 +1103,40 @@ describe('StreamableHTTPClientTransport', () => {
expect(fetchMock.mock.calls[1]![1]?.method).toBe('GET');
});

it('should stop retrying GET reconnection after a session-bound 404 clears the stale session', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-id',
reconnectionOptions: {
initialReconnectionDelay: 10,
maxRetries: 3,
maxReconnectionDelay: 1000,
reconnectionDelayGrowFactor: 1
}
});
await transport.start();

const fetchMock = globalThis.fetch as Mock;
fetchMock.mockResolvedValue({
ok: false,
status: 404,
statusText: 'Not Found',
headers: new Headers(),
text: () => Promise.resolve('Session not found')
});

(
transport as unknown as {
_scheduleReconnection: (opts: StartSSEOptions, attemptCount?: number) => void;
}
)._scheduleReconnection({}, 0);

await vi.advanceTimersByTimeAsync(20);
await vi.advanceTimersByTimeAsync(100);

expect(fetchMock).toHaveBeenCalledTimes(1);
expect(transport.sessionId).toBeUndefined();
});

it('should NOT reconnect a POST-initiated stream that fails', async () => {
// ARRANGE
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('^' + ".*" + '
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/gentle-maps-smile.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
'@modelcontextprotocol/client': patch
---

Clear stale Streamable HTTP client sessions when a session-bound request receives HTTP 404 by clearing the stored session ID, so the next initialize flow can proceed without an MCP session header.
33 changes: 31 additions & 2 deletions packages/client/src/client/streamableHttp.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,6 +25,19 @@ const DEFAULT_STREAMABLE_HTTP_RECONNECTION_OPTIONS: StreamableHTTPReconnectionOp
maxRetries: 2
};

const SESSION_BOUND_404_ERROR = Symbol('sessionBound404Error');

type SessionBound404Error = Error & { [SESSION_BOUND_404_ERROR]?: true };

function markSessionBound404Error(error: Error): Error {
(error as SessionBound404Error)[SESSION_BOUND_404_ERROR] = true;
return error;
}

function isSessionBound404Error(error: unknown): boolean {
return Boolean(error && typeof error === 'object' && (error as SessionBound404Error)[SESSION_BOUND_404_ERROR] === true);
}

/**
* Options for starting or authenticating an SSE connection
*/
Expand DownExpand Up@@ -237,6 +250,7 @@ export class StreamableHTTPClientTransport implements Transport {
// 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 sentSessionId = headers.get('mcp-session-id');
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'text/event-stream'];
headers.set('accept', [...new Set(types)].join(', '));
Expand All@@ -254,6 +268,11 @@ export class StreamableHTTPClientTransport implements Transport {
});

if (!response.ok) {
const shouldClearSessionFor404 = response.status === 404 && sentSessionId !== null && this._sessionId === sentSessionId;
if (shouldClearSessionFor404) {
this._sessionId = undefined;
}

if (response.status === 401 && this._authProvider) {
if (response.headers.has('www-authenticate')) {
const { resourceMetadataUrl, scope } = extractWWWAuthenticateParams(response);
Comment thread
Maverick-666 marked this conversation as resolved.
Expand DownExpand Up@@ -288,10 +307,11 @@ export class StreamableHTTPClientTransport implements Transport {
return;
}

throw new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
const error = new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
status: response.status,
statusText: response.statusText
});
throw shouldClearSessionFor404 ? markSessionBound404Error(error) : error;
}

this._handleSseStream(response.body, options, true);
Expand DownExpand Up@@ -345,7 +365,11 @@ export class StreamableHTTPClientTransport implements Transport {
this._cancelReconnection = undefined;
if (this._abortController?.signal.aborted) return;
this._startOrAuthSse(options).catch(error => {
this.onerror?.(new Error(`Failed to reconnect SSE stream: ${error instanceof Error ? error.message : String(error)}`));
const reconnectError = error instanceof Error ? error : new Error(String(error));
this.onerror?.(new Error(`Failed to reconnect SSE stream: ${reconnectError.message}`));
if (isSessionBound404Error(reconnectError)) {
return;
}
try {
this._scheduleReconnection(options, attemptCount + 1);
} catch (scheduleError) {
Expand DownExpand Up@@ -539,6 +563,7 @@ export class StreamableHTTPClientTransport implements Transport {
}

const headers = await this._commonHeaders();
const sentSessionId = headers.get('mcp-session-id');
headers.set('content-type', 'application/json');
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'application/json', 'text/event-stream'];
Expand All@@ -561,6 +586,10 @@ export class StreamableHTTPClientTransport implements Transport {
}

if (!response.ok) {
if (response.status === 404 && sentSessionId !== null && this._sessionId === sentSessionId) {
this._sessionId = undefined;
}
Comment thread
Maverick-666 marked this conversation as resolved.

if (response.status === 401 && this._authProvider) {
// Store WWW-Authenticate params for interactive finishAuth() path
if (response.headers.has('www-authenticate')) {
Comment thread
Maverick-666 marked this conversation as resolved.
Expand Down
203 changes: 202 additions & 1 deletion packages/client/test/client/streamableHttp.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -220,7 +220,7 @@ describe('StreamableHTTPClientTransport', () => {
await expect(transport.terminateSession()).resolves.not.toThrow();
});

it('should handle 404 response when session expires', async () => {
it('should preserve existing 404 behavior when request is not session-bound', async () => {
const message: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'test',
Expand DownExpand Up@@ -248,6 +248,104 @@ describe('StreamableHTTPClientTransport', () => {
expect(errorSpy).toHaveBeenCalled();
});

it('should clear session ID on 404 for session-bound POST requests', async () => {
const initializeMessage: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'initialize',
params: {
clientInfo: { name: 'test-client', version: '1.0' },
protocolVersion: '2025-03-26'
},
id: 'init-id'
};
const message: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'tools/list',
params: {},
id: 'test-id'
};

(globalThis.fetch as Mock)
.mockResolvedValueOnce({
ok: true,
status: 202,
headers: new Headers({ 'mcp-session-id': 'stale-session-id' }),
text: () => Promise.resolve('')
})
.mockResolvedValueOnce({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
})
.mockResolvedValueOnce({
ok: true,
status: 202,
headers: new Headers(),
text: () => Promise.resolve('')
});

await transport.send(initializeMessage);
expect(transport.sessionId).toBe('stale-session-id');

await expect(transport.send(message)).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpNotImplemented,
data: expect.objectContaining({
status: 404,
text: 'Session not found'
})
});
expect(transport.sessionId).toBeUndefined();

await transport.send({ jsonrpc: '2.0', method: 'notifications/ping' } as JSONRPCMessage);
const lastCall = (globalThis.fetch as Mock).mock.calls.at(-1)!;
expect(lastCall[1].headers.get('mcp-session-id')).toBeNull();
});

it('should not clear a newer session ID when a stale session-bound POST request returns 404', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-A'
});

const message: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'tools/list',
params: {},
id: 'test-id'
};

let resolveFetch!: (value: unknown) => void;
const deferredFetch = new Promise(resolve => {
resolveFetch = resolve;
});

(globalThis.fetch as Mock).mockImplementationOnce(() => {
// Simulate another in-flight request establishing a fresh session while this request is pending.
(transport as unknown as { _sessionId?: string })._sessionId = 'fresh-session-B';
return deferredFetch;
});

const sendPromise = transport.send(message);

resolveFetch({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
});

await expect(sendPromise).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpNotImplemented,
data: expect.objectContaining({
status: 404
})
});

expect(transport.sessionId).toBe('fresh-session-B');
});

it('should handle non-streaming JSON response', async () => {
const message: JSONRPCMessage = {
jsonrpc: '2.0',
Expand DownExpand Up@@ -309,6 +407,75 @@ describe('StreamableHTTPClientTransport', () => {
expect(globalThis.fetch).toHaveBeenCalledTimes(2);
});

it('should clear session ID when GET SSE stream returns 404 for a session-bound request', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-id'
});
await transport.start();

(globalThis.fetch as Mock).mockResolvedValueOnce({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
});

await expect(
(transport as unknown as { _startOrAuthSse: (opts: StartSSEOptions) => Promise<void> })._startOrAuthSse({})
).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpFailedToOpenStream,
data: expect.objectContaining({
status: 404,
statusText: 'Not Found'
})
});

expect(transport.sessionId).toBeUndefined();

const getCall = (globalThis.fetch as Mock).mock.calls[0]!;
expect(getCall[1].method).toBe('GET');
expect(getCall[1].headers.get('mcp-session-id')).toBe('stale-session-id');
});

it('should not clear a newer session ID when a stale session-bound GET request returns 404', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-A'
});
await transport.start();

let resolveFetch!: (value: unknown) => void;
const deferredFetch = new Promise(resolve => {
resolveFetch = resolve;
});

(globalThis.fetch as Mock).mockImplementationOnce(() => {
// Simulate another in-flight request establishing a fresh session while this request is pending.
(transport as unknown as { _sessionId?: string })._sessionId = 'fresh-session-B';
return deferredFetch;
});

const startPromise = (transport as unknown as { _startOrAuthSse: (opts: StartSSEOptions) => Promise<void> })._startOrAuthSse({});

resolveFetch({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
});

await expect(startPromise).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpFailedToOpenStream,
data: expect.objectContaining({
status: 404,
statusText: 'Not Found'
})
});

expect(transport.sessionId).toBe('fresh-session-B');
});

it('should handle successful initial GET connection for SSE', async () => {
// Set up readable stream for SSE events
const encoder = new TextEncoder();
Expand DownExpand Up@@ -936,6 +1103,40 @@ describe('StreamableHTTPClientTransport', () => {
expect(fetchMock.mock.calls[1]![1]?.method).toBe('GET');
});

it('should stop retrying GET reconnection after a session-bound 404 clears the stale session', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-id',
reconnectionOptions: {
initialReconnectionDelay: 10,
maxRetries: 3,
maxReconnectionDelay: 1000,
reconnectionDelayGrowFactor: 1
}
});
await transport.start();

const fetchMock = globalThis.fetch as Mock;
fetchMock.mockResolvedValue({
ok: false,
status: 404,
statusText: 'Not Found',
headers: new Headers(),
text: () => Promise.resolve('Session not found')
});

(
transport as unknown as {
_scheduleReconnection: (opts: StartSSEOptions, attemptCount?: number) => void;
}
)._scheduleReconnection({}, 0);

await vi.advanceTimersByTimeAsync(20);
await vi.advanceTimersByTimeAsync(100);

expect(fetchMock).toHaveBeenCalledTimes(1);
expect(transport.sessionId).toBeUndefined();
});

it('should NOT reconnect a POST-initiated stream that fails', async () => {
// ARRANGE
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" + '
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/gentle-maps-smile.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
'@modelcontextprotocol/client': patch
---

Clear stale Streamable HTTP client sessions when a session-bound request receives HTTP 404 by clearing the stored session ID, so the next initialize flow can proceed without an MCP session header.
33 changes: 31 additions & 2 deletions packages/client/src/client/streamableHttp.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,6 +25,19 @@ const DEFAULT_STREAMABLE_HTTP_RECONNECTION_OPTIONS: StreamableHTTPReconnectionOp
maxRetries: 2
};

const SESSION_BOUND_404_ERROR = Symbol('sessionBound404Error');

type SessionBound404Error = Error & { [SESSION_BOUND_404_ERROR]?: true };

function markSessionBound404Error(error: Error): Error {
(error as SessionBound404Error)[SESSION_BOUND_404_ERROR] = true;
return error;
}

function isSessionBound404Error(error: unknown): boolean {
return Boolean(error && typeof error === 'object' && (error as SessionBound404Error)[SESSION_BOUND_404_ERROR] === true);
}

/**
* Options for starting or authenticating an SSE connection
*/
Expand DownExpand Up@@ -237,6 +250,7 @@ export class StreamableHTTPClientTransport implements Transport {
// 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 sentSessionId = headers.get('mcp-session-id');
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'text/event-stream'];
headers.set('accept', [...new Set(types)].join(', '));
Expand All@@ -254,6 +268,11 @@ export class StreamableHTTPClientTransport implements Transport {
});

if (!response.ok) {
const shouldClearSessionFor404 = response.status === 404 && sentSessionId !== null && this._sessionId === sentSessionId;
if (shouldClearSessionFor404) {
this._sessionId = undefined;
}

if (response.status === 401 && this._authProvider) {
if (response.headers.has('www-authenticate')) {
const { resourceMetadataUrl, scope } = extractWWWAuthenticateParams(response);
Comment thread
Maverick-666 marked this conversation as resolved.
Expand DownExpand Up@@ -288,10 +307,11 @@ export class StreamableHTTPClientTransport implements Transport {
return;
}

throw new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
const error = new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
status: response.status,
statusText: response.statusText
});
throw shouldClearSessionFor404 ? markSessionBound404Error(error) : error;
}

this._handleSseStream(response.body, options, true);
Expand DownExpand Up@@ -345,7 +365,11 @@ export class StreamableHTTPClientTransport implements Transport {
this._cancelReconnection = undefined;
if (this._abortController?.signal.aborted) return;
this._startOrAuthSse(options).catch(error => {
this.onerror?.(new Error(`Failed to reconnect SSE stream: ${error instanceof Error ? error.message : String(error)}`));
const reconnectError = error instanceof Error ? error : new Error(String(error));
this.onerror?.(new Error(`Failed to reconnect SSE stream: ${reconnectError.message}`));
if (isSessionBound404Error(reconnectError)) {
return;
}
try {
this._scheduleReconnection(options, attemptCount + 1);
} catch (scheduleError) {
Expand DownExpand Up@@ -539,6 +563,7 @@ export class StreamableHTTPClientTransport implements Transport {
}

const headers = await this._commonHeaders();
const sentSessionId = headers.get('mcp-session-id');
headers.set('content-type', 'application/json');
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'application/json', 'text/event-stream'];
Expand All@@ -561,6 +586,10 @@ export class StreamableHTTPClientTransport implements Transport {
}

if (!response.ok) {
if (response.status === 404 && sentSessionId !== null && this._sessionId === sentSessionId) {
this._sessionId = undefined;
}
Comment thread
Maverick-666 marked this conversation as resolved.

if (response.status === 401 && this._authProvider) {
// Store WWW-Authenticate params for interactive finishAuth() path
if (response.headers.has('www-authenticate')) {
Comment thread
Maverick-666 marked this conversation as resolved.
Expand Down
203 changes: 202 additions & 1 deletion packages/client/test/client/streamableHttp.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -220,7 +220,7 @@ describe('StreamableHTTPClientTransport', () => {
await expect(transport.terminateSession()).resolves.not.toThrow();
});

it('should handle 404 response when session expires', async () => {
it('should preserve existing 404 behavior when request is not session-bound', async () => {
const message: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'test',
Expand DownExpand Up@@ -248,6 +248,104 @@ describe('StreamableHTTPClientTransport', () => {
expect(errorSpy).toHaveBeenCalled();
});

it('should clear session ID on 404 for session-bound POST requests', async () => {
const initializeMessage: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'initialize',
params: {
clientInfo: { name: 'test-client', version: '1.0' },
protocolVersion: '2025-03-26'
},
id: 'init-id'
};
const message: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'tools/list',
params: {},
id: 'test-id'
};

(globalThis.fetch as Mock)
.mockResolvedValueOnce({
ok: true,
status: 202,
headers: new Headers({ 'mcp-session-id': 'stale-session-id' }),
text: () => Promise.resolve('')
})
.mockResolvedValueOnce({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
})
.mockResolvedValueOnce({
ok: true,
status: 202,
headers: new Headers(),
text: () => Promise.resolve('')
});

await transport.send(initializeMessage);
expect(transport.sessionId).toBe('stale-session-id');

await expect(transport.send(message)).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpNotImplemented,
data: expect.objectContaining({
status: 404,
text: 'Session not found'
})
});
expect(transport.sessionId).toBeUndefined();

await transport.send({ jsonrpc: '2.0', method: 'notifications/ping' } as JSONRPCMessage);
const lastCall = (globalThis.fetch as Mock).mock.calls.at(-1)!;
expect(lastCall[1].headers.get('mcp-session-id')).toBeNull();
});

it('should not clear a newer session ID when a stale session-bound POST request returns 404', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-A'
});

const message: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'tools/list',
params: {},
id: 'test-id'
};

let resolveFetch!: (value: unknown) => void;
const deferredFetch = new Promise(resolve => {
resolveFetch = resolve;
});

(globalThis.fetch as Mock).mockImplementationOnce(() => {
// Simulate another in-flight request establishing a fresh session while this request is pending.
(transport as unknown as { _sessionId?: string })._sessionId = 'fresh-session-B';
return deferredFetch;
});

const sendPromise = transport.send(message);

resolveFetch({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
});

await expect(sendPromise).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpNotImplemented,
data: expect.objectContaining({
status: 404
})
});

expect(transport.sessionId).toBe('fresh-session-B');
});

it('should handle non-streaming JSON response', async () => {
const message: JSONRPCMessage = {
jsonrpc: '2.0',
Expand DownExpand Up@@ -309,6 +407,75 @@ describe('StreamableHTTPClientTransport', () => {
expect(globalThis.fetch).toHaveBeenCalledTimes(2);
});

it('should clear session ID when GET SSE stream returns 404 for a session-bound request', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-id'
});
await transport.start();

(globalThis.fetch as Mock).mockResolvedValueOnce({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
});

await expect(
(transport as unknown as { _startOrAuthSse: (opts: StartSSEOptions) => Promise<void> })._startOrAuthSse({})
).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpFailedToOpenStream,
data: expect.objectContaining({
status: 404,
statusText: 'Not Found'
})
});

expect(transport.sessionId).toBeUndefined();

const getCall = (globalThis.fetch as Mock).mock.calls[0]!;
expect(getCall[1].method).toBe('GET');
expect(getCall[1].headers.get('mcp-session-id')).toBe('stale-session-id');
});

it('should not clear a newer session ID when a stale session-bound GET request returns 404', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-A'
});
await transport.start();

let resolveFetch!: (value: unknown) => void;
const deferredFetch = new Promise(resolve => {
resolveFetch = resolve;
});

(globalThis.fetch as Mock).mockImplementationOnce(() => {
// Simulate another in-flight request establishing a fresh session while this request is pending.
(transport as unknown as { _sessionId?: string })._sessionId = 'fresh-session-B';
return deferredFetch;
});

const startPromise = (transport as unknown as { _startOrAuthSse: (opts: StartSSEOptions) => Promise<void> })._startOrAuthSse({});

resolveFetch({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
});

await expect(startPromise).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpFailedToOpenStream,
data: expect.objectContaining({
status: 404,
statusText: 'Not Found'
})
});

expect(transport.sessionId).toBe('fresh-session-B');
});

it('should handle successful initial GET connection for SSE', async () => {
// Set up readable stream for SSE events
const encoder = new TextEncoder();
Expand DownExpand Up@@ -936,6 +1103,40 @@ describe('StreamableHTTPClientTransport', () => {
expect(fetchMock.mock.calls[1]![1]?.method).toBe('GET');
});

it('should stop retrying GET reconnection after a session-bound 404 clears the stale session', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-id',
reconnectionOptions: {
initialReconnectionDelay: 10,
maxRetries: 3,
maxReconnectionDelay: 1000,
reconnectionDelayGrowFactor: 1
}
});
await transport.start();

const fetchMock = globalThis.fetch as Mock;
fetchMock.mockResolvedValue({
ok: false,
status: 404,
statusText: 'Not Found',
headers: new Headers(),
text: () => Promise.resolve('Session not found')
});

(
transport as unknown as {
_scheduleReconnection: (opts: StartSSEOptions, attemptCount?: number) => void;
}
)._scheduleReconnection({}, 0);

await vi.advanceTimersByTimeAsync(20);
await vi.advanceTimersByTimeAsync(100);

expect(fetchMock).toHaveBeenCalledTimes(1);
expect(transport.sessionId).toBeUndefined();
});

it('should NOT reconnect a POST-initiated stream that fails', async () => {
// ARRANGE
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('^' + ".*" + '
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/gentle-maps-smile.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
'@modelcontextprotocol/client': patch
---

Clear stale Streamable HTTP client sessions when a session-bound request receives HTTP 404 by clearing the stored session ID, so the next initialize flow can proceed without an MCP session header.
33 changes: 31 additions & 2 deletions packages/client/src/client/streamableHttp.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,6 +25,19 @@ const DEFAULT_STREAMABLE_HTTP_RECONNECTION_OPTIONS: StreamableHTTPReconnectionOp
maxRetries: 2
};

const SESSION_BOUND_404_ERROR = Symbol('sessionBound404Error');

type SessionBound404Error = Error & { [SESSION_BOUND_404_ERROR]?: true };

function markSessionBound404Error(error: Error): Error {
(error as SessionBound404Error)[SESSION_BOUND_404_ERROR] = true;
return error;
}

function isSessionBound404Error(error: unknown): boolean {
return Boolean(error && typeof error === 'object' && (error as SessionBound404Error)[SESSION_BOUND_404_ERROR] === true);
}

/**
* Options for starting or authenticating an SSE connection
*/
Expand DownExpand Up@@ -237,6 +250,7 @@ export class StreamableHTTPClientTransport implements Transport {
// 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 sentSessionId = headers.get('mcp-session-id');
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'text/event-stream'];
headers.set('accept', [...new Set(types)].join(', '));
Expand All@@ -254,6 +268,11 @@ export class StreamableHTTPClientTransport implements Transport {
});

if (!response.ok) {
const shouldClearSessionFor404 = response.status === 404 && sentSessionId !== null && this._sessionId === sentSessionId;
if (shouldClearSessionFor404) {
this._sessionId = undefined;
}

if (response.status === 401 && this._authProvider) {
if (response.headers.has('www-authenticate')) {
const { resourceMetadataUrl, scope } = extractWWWAuthenticateParams(response);
Comment thread
Maverick-666 marked this conversation as resolved.
Expand DownExpand Up@@ -288,10 +307,11 @@ export class StreamableHTTPClientTransport implements Transport {
return;
}

throw new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
const error = new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
status: response.status,
statusText: response.statusText
});
throw shouldClearSessionFor404 ? markSessionBound404Error(error) : error;
}

this._handleSseStream(response.body, options, true);
Expand DownExpand Up@@ -345,7 +365,11 @@ export class StreamableHTTPClientTransport implements Transport {
this._cancelReconnection = undefined;
if (this._abortController?.signal.aborted) return;
this._startOrAuthSse(options).catch(error => {
this.onerror?.(new Error(`Failed to reconnect SSE stream: ${error instanceof Error ? error.message : String(error)}`));
const reconnectError = error instanceof Error ? error : new Error(String(error));
this.onerror?.(new Error(`Failed to reconnect SSE stream: ${reconnectError.message}`));
if (isSessionBound404Error(reconnectError)) {
return;
}
try {
this._scheduleReconnection(options, attemptCount + 1);
} catch (scheduleError) {
Expand DownExpand Up@@ -539,6 +563,7 @@ export class StreamableHTTPClientTransport implements Transport {
}

const headers = await this._commonHeaders();
const sentSessionId = headers.get('mcp-session-id');
headers.set('content-type', 'application/json');
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'application/json', 'text/event-stream'];
Expand All@@ -561,6 +586,10 @@ export class StreamableHTTPClientTransport implements Transport {
}

if (!response.ok) {
if (response.status === 404 && sentSessionId !== null && this._sessionId === sentSessionId) {
this._sessionId = undefined;
}
Comment thread
Maverick-666 marked this conversation as resolved.

if (response.status === 401 && this._authProvider) {
// Store WWW-Authenticate params for interactive finishAuth() path
if (response.headers.has('www-authenticate')) {
Comment thread
Maverick-666 marked this conversation as resolved.
Expand Down
203 changes: 202 additions & 1 deletion packages/client/test/client/streamableHttp.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -220,7 +220,7 @@ describe('StreamableHTTPClientTransport', () => {
await expect(transport.terminateSession()).resolves.not.toThrow();
});

it('should handle 404 response when session expires', async () => {
it('should preserve existing 404 behavior when request is not session-bound', async () => {
const message: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'test',
Expand DownExpand Up@@ -248,6 +248,104 @@ describe('StreamableHTTPClientTransport', () => {
expect(errorSpy).toHaveBeenCalled();
});

it('should clear session ID on 404 for session-bound POST requests', async () => {
const initializeMessage: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'initialize',
params: {
clientInfo: { name: 'test-client', version: '1.0' },
protocolVersion: '2025-03-26'
},
id: 'init-id'
};
const message: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'tools/list',
params: {},
id: 'test-id'
};

(globalThis.fetch as Mock)
.mockResolvedValueOnce({
ok: true,
status: 202,
headers: new Headers({ 'mcp-session-id': 'stale-session-id' }),
text: () => Promise.resolve('')
})
.mockResolvedValueOnce({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
})
.mockResolvedValueOnce({
ok: true,
status: 202,
headers: new Headers(),
text: () => Promise.resolve('')
});

await transport.send(initializeMessage);
expect(transport.sessionId).toBe('stale-session-id');

await expect(transport.send(message)).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpNotImplemented,
data: expect.objectContaining({
status: 404,
text: 'Session not found'
})
});
expect(transport.sessionId).toBeUndefined();

await transport.send({ jsonrpc: '2.0', method: 'notifications/ping' } as JSONRPCMessage);
const lastCall = (globalThis.fetch as Mock).mock.calls.at(-1)!;
expect(lastCall[1].headers.get('mcp-session-id')).toBeNull();
});

it('should not clear a newer session ID when a stale session-bound POST request returns 404', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-A'
});

const message: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'tools/list',
params: {},
id: 'test-id'
};

let resolveFetch!: (value: unknown) => void;
const deferredFetch = new Promise(resolve => {
resolveFetch = resolve;
});

(globalThis.fetch as Mock).mockImplementationOnce(() => {
// Simulate another in-flight request establishing a fresh session while this request is pending.
(transport as unknown as { _sessionId?: string })._sessionId = 'fresh-session-B';
return deferredFetch;
});

const sendPromise = transport.send(message);

resolveFetch({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
});

await expect(sendPromise).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpNotImplemented,
data: expect.objectContaining({
status: 404
})
});

expect(transport.sessionId).toBe('fresh-session-B');
});

it('should handle non-streaming JSON response', async () => {
const message: JSONRPCMessage = {
jsonrpc: '2.0',
Expand DownExpand Up@@ -309,6 +407,75 @@ describe('StreamableHTTPClientTransport', () => {
expect(globalThis.fetch).toHaveBeenCalledTimes(2);
});

it('should clear session ID when GET SSE stream returns 404 for a session-bound request', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-id'
});
await transport.start();

(globalThis.fetch as Mock).mockResolvedValueOnce({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
});

await expect(
(transport as unknown as { _startOrAuthSse: (opts: StartSSEOptions) => Promise<void> })._startOrAuthSse({})
).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpFailedToOpenStream,
data: expect.objectContaining({
status: 404,
statusText: 'Not Found'
})
});

expect(transport.sessionId).toBeUndefined();

const getCall = (globalThis.fetch as Mock).mock.calls[0]!;
expect(getCall[1].method).toBe('GET');
expect(getCall[1].headers.get('mcp-session-id')).toBe('stale-session-id');
});

it('should not clear a newer session ID when a stale session-bound GET request returns 404', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-A'
});
await transport.start();

let resolveFetch!: (value: unknown) => void;
const deferredFetch = new Promise(resolve => {
resolveFetch = resolve;
});

(globalThis.fetch as Mock).mockImplementationOnce(() => {
// Simulate another in-flight request establishing a fresh session while this request is pending.
(transport as unknown as { _sessionId?: string })._sessionId = 'fresh-session-B';
return deferredFetch;
});

const startPromise = (transport as unknown as { _startOrAuthSse: (opts: StartSSEOptions) => Promise<void> })._startOrAuthSse({});

resolveFetch({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
});

await expect(startPromise).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpFailedToOpenStream,
data: expect.objectContaining({
status: 404,
statusText: 'Not Found'
})
});

expect(transport.sessionId).toBe('fresh-session-B');
});

it('should handle successful initial GET connection for SSE', async () => {
// Set up readable stream for SSE events
const encoder = new TextEncoder();
Expand DownExpand Up@@ -936,6 +1103,40 @@ describe('StreamableHTTPClientTransport', () => {
expect(fetchMock.mock.calls[1]![1]?.method).toBe('GET');
});

it('should stop retrying GET reconnection after a session-bound 404 clears the stale session', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-id',
reconnectionOptions: {
initialReconnectionDelay: 10,
maxRetries: 3,
maxReconnectionDelay: 1000,
reconnectionDelayGrowFactor: 1
}
});
await transport.start();

const fetchMock = globalThis.fetch as Mock;
fetchMock.mockResolvedValue({
ok: false,
status: 404,
statusText: 'Not Found',
headers: new Headers(),
text: () => Promise.resolve('Session not found')
});

(
transport as unknown as {
_scheduleReconnection: (opts: StartSSEOptions, attemptCount?: number) => void;
}
)._scheduleReconnection({}, 0);

await vi.advanceTimersByTimeAsync(20);
await vi.advanceTimersByTimeAsync(100);

expect(fetchMock).toHaveBeenCalledTimes(1);
expect(transport.sessionId).toBeUndefined();
});

it('should NOT reconnect a POST-initiated stream that fails', async () => {
// ARRANGE
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('^' + ".*" + '
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/gentle-maps-smile.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
'@modelcontextprotocol/client': patch
---

Clear stale Streamable HTTP client sessions when a session-bound request receives HTTP 404 by clearing the stored session ID, so the next initialize flow can proceed without an MCP session header.
33 changes: 31 additions & 2 deletions packages/client/src/client/streamableHttp.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,6 +25,19 @@ const DEFAULT_STREAMABLE_HTTP_RECONNECTION_OPTIONS: StreamableHTTPReconnectionOp
maxRetries: 2
};

const SESSION_BOUND_404_ERROR = Symbol('sessionBound404Error');

type SessionBound404Error = Error & { [SESSION_BOUND_404_ERROR]?: true };

function markSessionBound404Error(error: Error): Error {
(error as SessionBound404Error)[SESSION_BOUND_404_ERROR] = true;
return error;
}

function isSessionBound404Error(error: unknown): boolean {
return Boolean(error && typeof error === 'object' && (error as SessionBound404Error)[SESSION_BOUND_404_ERROR] === true);
}

/**
* Options for starting or authenticating an SSE connection
*/
Expand DownExpand Up@@ -237,6 +250,7 @@ export class StreamableHTTPClientTransport implements Transport {
// 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 sentSessionId = headers.get('mcp-session-id');
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'text/event-stream'];
headers.set('accept', [...new Set(types)].join(', '));
Expand All@@ -254,6 +268,11 @@ export class StreamableHTTPClientTransport implements Transport {
});

if (!response.ok) {
const shouldClearSessionFor404 = response.status === 404 && sentSessionId !== null && this._sessionId === sentSessionId;
if (shouldClearSessionFor404) {
this._sessionId = undefined;
}

if (response.status === 401 && this._authProvider) {
if (response.headers.has('www-authenticate')) {
const { resourceMetadataUrl, scope } = extractWWWAuthenticateParams(response);
Comment thread
Maverick-666 marked this conversation as resolved.
Expand DownExpand Up@@ -288,10 +307,11 @@ export class StreamableHTTPClientTransport implements Transport {
return;
}

throw new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
const error = new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
status: response.status,
statusText: response.statusText
});
throw shouldClearSessionFor404 ? markSessionBound404Error(error) : error;
}

this._handleSseStream(response.body, options, true);
Expand DownExpand Up@@ -345,7 +365,11 @@ export class StreamableHTTPClientTransport implements Transport {
this._cancelReconnection = undefined;
if (this._abortController?.signal.aborted) return;
this._startOrAuthSse(options).catch(error => {
this.onerror?.(new Error(`Failed to reconnect SSE stream: ${error instanceof Error ? error.message : String(error)}`));
const reconnectError = error instanceof Error ? error : new Error(String(error));
this.onerror?.(new Error(`Failed to reconnect SSE stream: ${reconnectError.message}`));
if (isSessionBound404Error(reconnectError)) {
return;
}
try {
this._scheduleReconnection(options, attemptCount + 1);
} catch (scheduleError) {
Expand DownExpand Up@@ -539,6 +563,7 @@ export class StreamableHTTPClientTransport implements Transport {
}

const headers = await this._commonHeaders();
const sentSessionId = headers.get('mcp-session-id');
headers.set('content-type', 'application/json');
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'application/json', 'text/event-stream'];
Expand All@@ -561,6 +586,10 @@ export class StreamableHTTPClientTransport implements Transport {
}

if (!response.ok) {
if (response.status === 404 && sentSessionId !== null && this._sessionId === sentSessionId) {
this._sessionId = undefined;
}
Comment thread
Maverick-666 marked this conversation as resolved.

if (response.status === 401 && this._authProvider) {
// Store WWW-Authenticate params for interactive finishAuth() path
if (response.headers.has('www-authenticate')) {
Comment thread
Maverick-666 marked this conversation as resolved.
Expand Down
203 changes: 202 additions & 1 deletion packages/client/test/client/streamableHttp.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -220,7 +220,7 @@ describe('StreamableHTTPClientTransport', () => {
await expect(transport.terminateSession()).resolves.not.toThrow();
});

it('should handle 404 response when session expires', async () => {
it('should preserve existing 404 behavior when request is not session-bound', async () => {
const message: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'test',
Expand DownExpand Up@@ -248,6 +248,104 @@ describe('StreamableHTTPClientTransport', () => {
expect(errorSpy).toHaveBeenCalled();
});

it('should clear session ID on 404 for session-bound POST requests', async () => {
const initializeMessage: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'initialize',
params: {
clientInfo: { name: 'test-client', version: '1.0' },
protocolVersion: '2025-03-26'
},
id: 'init-id'
};
const message: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'tools/list',
params: {},
id: 'test-id'
};

(globalThis.fetch as Mock)
.mockResolvedValueOnce({
ok: true,
status: 202,
headers: new Headers({ 'mcp-session-id': 'stale-session-id' }),
text: () => Promise.resolve('')
})
.mockResolvedValueOnce({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
})
.mockResolvedValueOnce({
ok: true,
status: 202,
headers: new Headers(),
text: () => Promise.resolve('')
});

await transport.send(initializeMessage);
expect(transport.sessionId).toBe('stale-session-id');

await expect(transport.send(message)).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpNotImplemented,
data: expect.objectContaining({
status: 404,
text: 'Session not found'
})
});
expect(transport.sessionId).toBeUndefined();

await transport.send({ jsonrpc: '2.0', method: 'notifications/ping' } as JSONRPCMessage);
const lastCall = (globalThis.fetch as Mock).mock.calls.at(-1)!;
expect(lastCall[1].headers.get('mcp-session-id')).toBeNull();
});

it('should not clear a newer session ID when a stale session-bound POST request returns 404', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-A'
});

const message: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'tools/list',
params: {},
id: 'test-id'
};

let resolveFetch!: (value: unknown) => void;
const deferredFetch = new Promise(resolve => {
resolveFetch = resolve;
});

(globalThis.fetch as Mock).mockImplementationOnce(() => {
// Simulate another in-flight request establishing a fresh session while this request is pending.
(transport as unknown as { _sessionId?: string })._sessionId = 'fresh-session-B';
return deferredFetch;
});

const sendPromise = transport.send(message);

resolveFetch({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
});

await expect(sendPromise).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpNotImplemented,
data: expect.objectContaining({
status: 404
})
});

expect(transport.sessionId).toBe('fresh-session-B');
});

it('should handle non-streaming JSON response', async () => {
const message: JSONRPCMessage = {
jsonrpc: '2.0',
Expand DownExpand Up@@ -309,6 +407,75 @@ describe('StreamableHTTPClientTransport', () => {
expect(globalThis.fetch).toHaveBeenCalledTimes(2);
});

it('should clear session ID when GET SSE stream returns 404 for a session-bound request', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-id'
});
await transport.start();

(globalThis.fetch as Mock).mockResolvedValueOnce({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
});

await expect(
(transport as unknown as { _startOrAuthSse: (opts: StartSSEOptions) => Promise<void> })._startOrAuthSse({})
).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpFailedToOpenStream,
data: expect.objectContaining({
status: 404,
statusText: 'Not Found'
})
});

expect(transport.sessionId).toBeUndefined();

const getCall = (globalThis.fetch as Mock).mock.calls[0]!;
expect(getCall[1].method).toBe('GET');
expect(getCall[1].headers.get('mcp-session-id')).toBe('stale-session-id');
});

it('should not clear a newer session ID when a stale session-bound GET request returns 404', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-A'
});
await transport.start();

let resolveFetch!: (value: unknown) => void;
const deferredFetch = new Promise(resolve => {
resolveFetch = resolve;
});

(globalThis.fetch as Mock).mockImplementationOnce(() => {
// Simulate another in-flight request establishing a fresh session while this request is pending.
(transport as unknown as { _sessionId?: string })._sessionId = 'fresh-session-B';
return deferredFetch;
});

const startPromise = (transport as unknown as { _startOrAuthSse: (opts: StartSSEOptions) => Promise<void> })._startOrAuthSse({});

resolveFetch({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
});

await expect(startPromise).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpFailedToOpenStream,
data: expect.objectContaining({
status: 404,
statusText: 'Not Found'
})
});

expect(transport.sessionId).toBe('fresh-session-B');
});

it('should handle successful initial GET connection for SSE', async () => {
// Set up readable stream for SSE events
const encoder = new TextEncoder();
Expand DownExpand Up@@ -936,6 +1103,40 @@ describe('StreamableHTTPClientTransport', () => {
expect(fetchMock.mock.calls[1]![1]?.method).toBe('GET');
});

it('should stop retrying GET reconnection after a session-bound 404 clears the stale session', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-id',
reconnectionOptions: {
initialReconnectionDelay: 10,
maxRetries: 3,
maxReconnectionDelay: 1000,
reconnectionDelayGrowFactor: 1
}
});
await transport.start();

const fetchMock = globalThis.fetch as Mock;
fetchMock.mockResolvedValue({
ok: false,
status: 404,
statusText: 'Not Found',
headers: new Headers(),
text: () => Promise.resolve('Session not found')
});

(
transport as unknown as {
_scheduleReconnection: (opts: StartSSEOptions, attemptCount?: number) => void;
}
)._scheduleReconnection({}, 0);

await vi.advanceTimersByTimeAsync(20);
await vi.advanceTimersByTimeAsync(100);

expect(fetchMock).toHaveBeenCalledTimes(1);
expect(transport.sessionId).toBeUndefined();
});

it('should NOT reconnect a POST-initiated stream that fails', async () => {
// ARRANGE
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); } })(); })();
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/gentle-maps-smile.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
'@modelcontextprotocol/client': patch
---

Clear stale Streamable HTTP client sessions when a session-bound request receives HTTP 404 by clearing the stored session ID, so the next initialize flow can proceed without an MCP session header.
33 changes: 31 additions & 2 deletions packages/client/src/client/streamableHttp.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,6 +25,19 @@ const DEFAULT_STREAMABLE_HTTP_RECONNECTION_OPTIONS: StreamableHTTPReconnectionOp
maxRetries: 2
};

const SESSION_BOUND_404_ERROR = Symbol('sessionBound404Error');

type SessionBound404Error = Error & { [SESSION_BOUND_404_ERROR]?: true };

function markSessionBound404Error(error: Error): Error {
(error as SessionBound404Error)[SESSION_BOUND_404_ERROR] = true;
return error;
}

function isSessionBound404Error(error: unknown): boolean {
return Boolean(error && typeof error === 'object' && (error as SessionBound404Error)[SESSION_BOUND_404_ERROR] === true);
}

/**
* Options for starting or authenticating an SSE connection
*/
Expand DownExpand Up@@ -237,6 +250,7 @@ export class StreamableHTTPClientTransport implements Transport {
// 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 sentSessionId = headers.get('mcp-session-id');
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'text/event-stream'];
headers.set('accept', [...new Set(types)].join(', '));
Expand All@@ -254,6 +268,11 @@ export class StreamableHTTPClientTransport implements Transport {
});

if (!response.ok) {
const shouldClearSessionFor404 = response.status === 404 && sentSessionId !== null && this._sessionId === sentSessionId;
if (shouldClearSessionFor404) {
this._sessionId = undefined;
}

if (response.status === 401 && this._authProvider) {
if (response.headers.has('www-authenticate')) {
const { resourceMetadataUrl, scope } = extractWWWAuthenticateParams(response);
Comment thread
Maverick-666 marked this conversation as resolved.
Expand DownExpand Up@@ -288,10 +307,11 @@ export class StreamableHTTPClientTransport implements Transport {
return;
}

throw new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
const error = new SdkError(SdkErrorCode.ClientHttpFailedToOpenStream, `Failed to open SSE stream: ${response.statusText}`, {
status: response.status,
statusText: response.statusText
});
throw shouldClearSessionFor404 ? markSessionBound404Error(error) : error;
}

this._handleSseStream(response.body, options, true);
Expand DownExpand Up@@ -345,7 +365,11 @@ export class StreamableHTTPClientTransport implements Transport {
this._cancelReconnection = undefined;
if (this._abortController?.signal.aborted) return;
this._startOrAuthSse(options).catch(error => {
this.onerror?.(new Error(`Failed to reconnect SSE stream: ${error instanceof Error ? error.message : String(error)}`));
const reconnectError = error instanceof Error ? error : new Error(String(error));
this.onerror?.(new Error(`Failed to reconnect SSE stream: ${reconnectError.message}`));
if (isSessionBound404Error(reconnectError)) {
return;
}
try {
this._scheduleReconnection(options, attemptCount + 1);
} catch (scheduleError) {
Expand DownExpand Up@@ -539,6 +563,7 @@ export class StreamableHTTPClientTransport implements Transport {
}

const headers = await this._commonHeaders();
const sentSessionId = headers.get('mcp-session-id');
headers.set('content-type', 'application/json');
const userAccept = headers.get('accept');
const types = [...(userAccept?.split(',').map(s => s.trim().toLowerCase()) ?? []), 'application/json', 'text/event-stream'];
Expand All@@ -561,6 +586,10 @@ export class StreamableHTTPClientTransport implements Transport {
}

if (!response.ok) {
if (response.status === 404 && sentSessionId !== null && this._sessionId === sentSessionId) {
this._sessionId = undefined;
}
Comment thread
Maverick-666 marked this conversation as resolved.

if (response.status === 401 && this._authProvider) {
// Store WWW-Authenticate params for interactive finishAuth() path
if (response.headers.has('www-authenticate')) {
Comment thread
Maverick-666 marked this conversation as resolved.
Expand Down
203 changes: 202 additions & 1 deletion packages/client/test/client/streamableHttp.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -220,7 +220,7 @@ describe('StreamableHTTPClientTransport', () => {
await expect(transport.terminateSession()).resolves.not.toThrow();
});

it('should handle 404 response when session expires', async () => {
it('should preserve existing 404 behavior when request is not session-bound', async () => {
const message: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'test',
Expand DownExpand Up@@ -248,6 +248,104 @@ describe('StreamableHTTPClientTransport', () => {
expect(errorSpy).toHaveBeenCalled();
});

it('should clear session ID on 404 for session-bound POST requests', async () => {
const initializeMessage: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'initialize',
params: {
clientInfo: { name: 'test-client', version: '1.0' },
protocolVersion: '2025-03-26'
},
id: 'init-id'
};
const message: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'tools/list',
params: {},
id: 'test-id'
};

(globalThis.fetch as Mock)
.mockResolvedValueOnce({
ok: true,
status: 202,
headers: new Headers({ 'mcp-session-id': 'stale-session-id' }),
text: () => Promise.resolve('')
})
.mockResolvedValueOnce({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
})
.mockResolvedValueOnce({
ok: true,
status: 202,
headers: new Headers(),
text: () => Promise.resolve('')
});

await transport.send(initializeMessage);
expect(transport.sessionId).toBe('stale-session-id');

await expect(transport.send(message)).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpNotImplemented,
data: expect.objectContaining({
status: 404,
text: 'Session not found'
})
});
expect(transport.sessionId).toBeUndefined();

await transport.send({ jsonrpc: '2.0', method: 'notifications/ping' } as JSONRPCMessage);
const lastCall = (globalThis.fetch as Mock).mock.calls.at(-1)!;
expect(lastCall[1].headers.get('mcp-session-id')).toBeNull();
});

it('should not clear a newer session ID when a stale session-bound POST request returns 404', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-A'
});

const message: JSONRPCMessage = {
jsonrpc: '2.0',
method: 'tools/list',
params: {},
id: 'test-id'
};

let resolveFetch!: (value: unknown) => void;
const deferredFetch = new Promise(resolve => {
resolveFetch = resolve;
});

(globalThis.fetch as Mock).mockImplementationOnce(() => {
// Simulate another in-flight request establishing a fresh session while this request is pending.
(transport as unknown as { _sessionId?: string })._sessionId = 'fresh-session-B';
return deferredFetch;
});

const sendPromise = transport.send(message);

resolveFetch({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
});

await expect(sendPromise).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpNotImplemented,
data: expect.objectContaining({
status: 404
})
});

expect(transport.sessionId).toBe('fresh-session-B');
});

it('should handle non-streaming JSON response', async () => {
const message: JSONRPCMessage = {
jsonrpc: '2.0',
Expand DownExpand Up@@ -309,6 +407,75 @@ describe('StreamableHTTPClientTransport', () => {
expect(globalThis.fetch).toHaveBeenCalledTimes(2);
});

it('should clear session ID when GET SSE stream returns 404 for a session-bound request', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-id'
});
await transport.start();

(globalThis.fetch as Mock).mockResolvedValueOnce({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
});

await expect(
(transport as unknown as { _startOrAuthSse: (opts: StartSSEOptions) => Promise<void> })._startOrAuthSse({})
).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpFailedToOpenStream,
data: expect.objectContaining({
status: 404,
statusText: 'Not Found'
})
});

expect(transport.sessionId).toBeUndefined();

const getCall = (globalThis.fetch as Mock).mock.calls[0]!;
expect(getCall[1].method).toBe('GET');
expect(getCall[1].headers.get('mcp-session-id')).toBe('stale-session-id');
});

it('should not clear a newer session ID when a stale session-bound GET request returns 404', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-A'
});
await transport.start();

let resolveFetch!: (value: unknown) => void;
const deferredFetch = new Promise(resolve => {
resolveFetch = resolve;
});

(globalThis.fetch as Mock).mockImplementationOnce(() => {
// Simulate another in-flight request establishing a fresh session while this request is pending.
(transport as unknown as { _sessionId?: string })._sessionId = 'fresh-session-B';
return deferredFetch;
});

const startPromise = (transport as unknown as { _startOrAuthSse: (opts: StartSSEOptions) => Promise<void> })._startOrAuthSse({});

resolveFetch({
ok: false,
status: 404,
statusText: 'Not Found',
text: () => Promise.resolve('Session not found'),
headers: new Headers()
});

await expect(startPromise).rejects.toMatchObject({
code: SdkErrorCode.ClientHttpFailedToOpenStream,
data: expect.objectContaining({
status: 404,
statusText: 'Not Found'
})
});

expect(transport.sessionId).toBe('fresh-session-B');
});

it('should handle successful initial GET connection for SSE', async () => {
// Set up readable stream for SSE events
const encoder = new TextEncoder();
Expand DownExpand Up@@ -936,6 +1103,40 @@ describe('StreamableHTTPClientTransport', () => {
expect(fetchMock.mock.calls[1]![1]?.method).toBe('GET');
});

it('should stop retrying GET reconnection after a session-bound 404 clears the stale session', async () => {
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
sessionId: 'stale-session-id',
reconnectionOptions: {
initialReconnectionDelay: 10,
maxRetries: 3,
maxReconnectionDelay: 1000,
reconnectionDelayGrowFactor: 1
}
});
await transport.start();

const fetchMock = globalThis.fetch as Mock;
fetchMock.mockResolvedValue({
ok: false,
status: 404,
statusText: 'Not Found',
headers: new Headers(),
text: () => Promise.resolve('Session not found')
});

(
transport as unknown as {
_scheduleReconnection: (opts: StartSSEOptions, attemptCount?: number) => void;
}
)._scheduleReconnection({}, 0);

await vi.advanceTimersByTimeAsync(20);
await vi.advanceTimersByTimeAsync(100);

expect(fetchMock).toHaveBeenCalledTimes(1);
expect(transport.sessionId).toBeUndefined();
});

it('should NOT reconnect a POST-initiated stream that fails', async () => {
// ARRANGE
transport = new StreamableHTTPClientTransport(new URL('http://localhost:1234/mcp'), {
Expand Down
Loading