Skip to content

Commit 0a53be9

Browse files
[SEP-2567] core: serverStatelessRouter + stdio/InMemory per-message routers
NEW core/shared/serverStatelessRouter.ts: routeServerStateless() — per-message branch on isStatelessRequest to StatelessHandlers {dispatch,listen}; notifications/cancelled aborts matching _inflight controller. stdio + InMemory server-side: receive paths route via routeServerStateless; close() aborts all in-flight. stdio + InMemory client-side: sendAndReceive threads opts?.signal to StreamDriver. Satisfies: 2567-R1 (pipe), 2567-R2
1 parent 502065d commit 0a53be9

5 files changed

Lines changed: 147 additions & 7 deletions

File tree

packages/client/src/client/stdio.ts

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -107,8 +107,8 @@ export class StdioClientTransport implements Transport {
107107
* Sends one request and returns the server's messages for it. Backed by
108108
* {@linkcode StreamDriver}; bypasses `Protocol.request()`.
109109
*/
110-
sendAndReceive(request: Omit<JSONRPCRequest, 'jsonrpc' | 'id'>): AsyncIterable<JSONRPCMessage> {
111-
return this._driver.sendAndReceive(request);
110+
sendAndReceive(request: Omit<JSONRPCRequest, 'jsonrpc' | 'id'>, opts?: { signal?: AbortSignal }): AsyncIterable<JSONRPCMessage> {
111+
return this._driver.sendAndReceive(request, opts);
112112
}
113113

114114
setProtocolVersion(version: string): void {

packages/core/src/index.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ export * from './shared/authUtils.js';
55
export * from './shared/dispatcher.js';
66
export * from './shared/metadataUtils.js';
77
export * from './shared/protocol.js';
8+
export * from './shared/serverStatelessRouter.js';
89
export * from './shared/stateless.js';
910
export * from './shared/stdio.js';
1011
export * from './shared/streamDriver.js';
Lines changed: 74 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,74 @@
1+
import type { JSONRPCMessage, JSONRPCRequest, RequestId } from '../types/index.js';
2+
import { isJSONRPCNotification, ProtocolErrorCode } from '../types/index.js';
3+
import { errorResponse } from './dispatcher.js';
4+
import type { ListenContext, StatelessHandlers } from './stateless.js';
5+
import { isStatelessRequest } from './stateless.js';
6+
7+
/**
8+
* Per-message router for pipe-shaped server transports (stdio, in-memory).
9+
* Call once per inbound message. Returns `true` if the message was claimed by
10+
* the stateless path (so the caller should NOT pass it to legacy `onmessage`).
11+
*
12+
* Stateless requests are dispatched via {@linkcode StatelessHandlers};
13+
* `notifications/cancelled` aborts a tracked in-flight request.
14+
*/
15+
export function routeServerStateless(
16+
message: JSONRPCMessage,
17+
handlers: StatelessHandlers,
18+
inflight: Map<RequestId, AbortController>,
19+
write: (m: JSONRPCMessage) => void,
20+
ctx: ListenContext,
21+
onerror?: (e: Error) => void
22+
): boolean {
23+
if (isStatelessRequest(message)) {
24+
const ac = new AbortController();
25+
inflight.set(message.id, ac);
26+
void handleOne(handlers, message, ac, write, ctx)
27+
.catch(error => onerror?.(error instanceof Error ? error : new Error(String(error))))
28+
.finally(() => inflight.delete(message.id));
29+
return true;
30+
}
31+
if (isJSONRPCNotification(message) && message.method === 'notifications/cancelled') {
32+
const requestId = (message.params as { requestId?: RequestId } | undefined)?.requestId;
33+
const ac = requestId === undefined ? undefined : inflight.get(requestId);
34+
if (ac) {
35+
ac.abort();
36+
return true;
37+
}
38+
}
39+
return false;
40+
}
41+
42+
async function handleOne(
43+
handlers: StatelessHandlers,
44+
req: JSONRPCRequest,
45+
ac: AbortController,
46+
write: (m: JSONRPCMessage) => void,
47+
ctx: ListenContext
48+
): Promise<void> {
49+
if (req.method === 'subscriptions/listen') {
50+
let listenStream;
51+
try {
52+
listenStream = handlers.listen(req, ctx);
53+
} catch (error) {
54+
write(
55+
errorResponse(req.id, ProtocolErrorCode.InvalidParams, error instanceof Error ? error.message : 'Invalid listen request')
56+
);
57+
return;
58+
}
59+
const { stream, close } = listenStream;
60+
ac.signal.addEventListener('abort', close, { once: true });
61+
try {
62+
for await (const m of stream) {
63+
if (ac.signal.aborted) break;
64+
write(m);
65+
}
66+
} finally {
67+
ac.signal.removeEventListener('abort', close);
68+
close();
69+
}
70+
} else {
71+
const response = await handlers.dispatch(req, { signal: ac.signal, authInfo: ctx.authInfo, notify: write });
72+
write(response);
73+
}
74+
}

packages/core/src/util/inMemory.ts

Lines changed: 35 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,6 @@
11
import { SdkError, SdkErrorCode } from '../errors/sdkErrors.js';
2+
import { routeServerStateless } from '../shared/serverStatelessRouter.js';
3+
import type { StatelessHandlers } from '../shared/stateless.js';
24
import { StreamDriver } from '../shared/streamDriver.js';
35
import type { Transport } from '../shared/transport.js';
46
import type { AuthInfo, JSONRPCMessage, JSONRPCRequest, RequestId } from '../types/index.js';
@@ -19,6 +21,9 @@ export class InMemoryTransport implements Transport {
1921
private _messageQueue: QueuedMessage[] = [];
2022
private _closed = false;
2123

24+
private _statelessHandlers?: StatelessHandlers;
25+
private readonly _inflight = new Map<RequestId, AbortController>();
26+
2227
/* eslint-disable-next-line unicorn/consistent-function-scoping */
2328
private readonly _driver = new StreamDriver(m => this.send(m));
2429

@@ -28,8 +33,13 @@ export class InMemoryTransport implements Transport {
2833
sessionId?: string;
2934

3035
/** Client-side: backed by {@linkcode StreamDriver}. */
31-
sendAndReceive(request: Omit<JSONRPCRequest, 'jsonrpc' | 'id'>): AsyncIterable<JSONRPCMessage> {
32-
return this._driver.sendAndReceive(request);
36+
sendAndReceive(request: Omit<JSONRPCRequest, 'jsonrpc' | 'id'>, opts?: { signal?: AbortSignal }): AsyncIterable<JSONRPCMessage> {
37+
return this._driver.sendAndReceive(request, opts);
38+
}
39+
40+
/** Server-side: installed by `Server.connect()`. */
41+
setStatelessHandlers(h: StatelessHandlers): void {
42+
this._statelessHandlers = h;
3343
}
3444

3545
/**
@@ -51,15 +61,37 @@ export class InMemoryTransport implements Transport {
5161
}
5262
}
5363

54-
/** Receive path: route to the StreamDriver first; fall through to `onmessage` for unclaimed. */
64+
/**
65+
* Receive path. Per-message router (mirrors stdio): server-side stateless
66+
* requests go to {@linkcode StatelessHandlers}; client-side
67+
* {@linkcode StreamDriver} claims responses for pending iterators; unclaimed
68+
* messages fall through to legacy `onmessage`.
69+
*/
5570
private _receive(message: JSONRPCMessage, extra?: { authInfo?: AuthInfo }): void {
71+
if (
72+
this._statelessHandlers &&
73+
routeServerStateless(
74+
message,
75+
this._statelessHandlers,
76+
this._inflight,
77+
m => {
78+
this.send(m).catch(error => this.onerror?.(error as Error));
79+
},
80+
{ authInfo: extra?.authInfo },
81+
e => this.onerror?.(e)
82+
)
83+
) {
84+
return;
85+
}
5686
if (this._driver.onMessage(message)) return;
5787
this.onmessage?.(message, extra);
5888
}
5989

6090
async close(): Promise<void> {
6191
if (this._closed) return;
6292
this._closed = true;
93+
for (const ac of this._inflight.values()) ac.abort();
94+
this._inflight.clear();
6395
this._driver.close();
6496

6597
const other = this._otherTransport;

packages/server/src/server/stdio.ts

Lines changed: 35 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
import type { Readable, Writable } from 'node:stream';
22

3-
import type { JSONRPCMessage, Transport } from '@modelcontextprotocol/core';
4-
import { ReadBuffer, serializeMessage } from '@modelcontextprotocol/core';
3+
import type { JSONRPCMessage, RequestId, StatelessHandlers, Transport } from '@modelcontextprotocol/core';
4+
import { ReadBuffer, routeServerStateless, serializeMessage } from '@modelcontextprotocol/core';
55
import { process } from '@modelcontextprotocol/server/_shims';
66

77
/**
@@ -21,6 +21,9 @@ export class StdioServerTransport implements Transport {
2121
private _started = false;
2222
private _closed = false;
2323

24+
private _statelessHandlers?: StatelessHandlers;
25+
private readonly _inflight = new Map<RequestId, AbortController>();
26+
2427
constructor(
2528
private _stdin: Readable = process.stdin,
2629
private _stdout: Writable = process.stdout
@@ -30,6 +33,15 @@ export class StdioServerTransport implements Transport {
3033
onerror?: (error: Error) => void;
3134
onmessage?: (message: JSONRPCMessage) => void;
3235

36+
/**
37+
* Installed by `Server.connect()`. When present, {@linkcode processReadBuffer}
38+
* routes 2026-06 requests via {@linkcode StatelessHandlers}; everything else
39+
* falls through to legacy `onmessage`.
40+
*/
41+
setStatelessHandlers(h: StatelessHandlers): void {
42+
this._statelessHandlers = h;
43+
}
44+
3345
// Arrow functions to bind `this` properly, while maintaining function identity.
3446
_ondata = (chunk: Buffer) => {
3547
this._readBuffer.append(chunk);
@@ -69,6 +81,25 @@ export class StdioServerTransport implements Transport {
6981
break;
7082
}
7183

84+
// Per-message router. Stateless requests (carrying _meta.protocolVersion
85+
// for a 2026-06 version) go to the dispatch/listen handlers; everything
86+
// else flows to the legacy onmessage path (Protocol._onrequest).
87+
if (
88+
this._statelessHandlers &&
89+
routeServerStateless(
90+
message,
91+
this._statelessHandlers,
92+
this._inflight,
93+
m => {
94+
this.send(m).catch(error => this.onerror?.(error as Error));
95+
},
96+
{},
97+
e => this.onerror?.(e)
98+
)
99+
) {
100+
continue;
101+
}
102+
72103
this.onmessage?.(message);
73104
} catch (error) {
74105
this.onerror?.(error as Error);
@@ -81,6 +112,8 @@ export class StdioServerTransport implements Transport {
81112
return;
82113
}
83114
this._closed = true;
115+
for (const ac of this._inflight.values()) ac.abort();
116+
this._inflight.clear();
84117

85118
// Remove our event listeners first
86119
this._stdin.off('data', this._ondata);

0 commit comments

Comments
 (0)