diff --git a/.changeset/mcp-http-session-termination.md b/.changeset/mcp-http-session-termination.md new file mode 100644 index 00000000000..46bd7c70b9a --- /dev/null +++ b/.changeset/mcp-http-session-termination.md @@ -0,0 +1,7 @@ +--- +"effect": patch +--- + +Add opt-in `allowSessionTermination` to `McpServer.layerHttp`. DELETE ends the session and interrupts its active requests; later requests with that session id return `404`. + +Fix an RPC cancellation race by registering request fibers before their handlers run. diff --git a/packages/effect/src/ai/McpServer.ts b/packages/effect/src/ai/McpServer.ts index f76029cfb18..1b79c680889 100644 --- a/packages/effect/src/ai/McpServer.ts +++ b/packages/effect/src/ai/McpServer.ts @@ -654,6 +654,7 @@ const cancelledResponses = new WeakMap() const requestKey = (requestId: string | number): string => `${typeof requestId}:${requestId}` interface ActiveRequest { + readonly requestId: RpcMessage.RequestId readonly prepared: McpRuntime.PreparedRequest readonly cancelled: boolean } @@ -805,6 +806,7 @@ const runWithRuntime = Effect.fnUntraced(function*( payload: { requestId, reason } }) }) + let writeFromClient!: (clientId: number, message: RpcMessage.FromClientEncoded) => Effect.Effect const handlers = yield* runtime.installHandlers({ core: internalState.get(server)!.core, subscribeServerNotifications: PubSub.subscribe(serverNotifications), @@ -976,8 +978,9 @@ const runWithRuntime = Effect.fnUntraced(function*( } return protocol.send(clientId, response) }, - run: (f) => - protocol.run((clientId, request_) => { + run: (f) => { + writeFromClient = f + return protocol.run((clientId, request_) => { const fiber = Fiber.getCurrent()! const request = request_ as unknown as | RpcMessage.FromServerEncoded @@ -1045,11 +1048,13 @@ const runWithRuntime = Effect.fnUntraced(function*( if (httpRequest !== undefined && session !== undefined) { appendPreResponseHandlerUnsafe(httpRequest, (_, res) => Effect.succeed( - HttpServerResponse.setHeader( - res, - MCP_PROTOCOL_VERSION_HEADER, - session.protocol.protocolVersion - ) + runtime.resolveRequest(clientId, headers) === undefined + ? HttpServerResponse.empty({ status: 404 }) + : HttpServerResponse.setHeader( + res, + MCP_PROTOCOL_VERSION_HEADER, + session.protocol.protocolVersion + ) )) } const routedRequest = runtime.routeClientRequest(selectedProtocol, request) @@ -1130,7 +1135,11 @@ const runWithRuntime = Effect.fnUntraced(function*( } if (request.isNotification !== true) { const requests = activeRequests.get(clientId) ?? new Map() - requests.set(requestKey(request.id), { prepared, cancelled: false }) + requests.set(requestKey(request.id), { + requestId: RpcMessage.RequestId(request.id), + prepared, + cancelled: false + }) activeRequests.set(clientId, requests) } const handled = f(clientId, routedRequest) @@ -1205,7 +1214,22 @@ const runWithRuntime = Effect.fnUntraced(function*( } } }) + } }) + if (isHttp) { + // Stop requests that can no longer receive client replies after termination. + yield* runtime.onSessionTerminated((binding) => { + const interrupts: Array> = [] + for (const [clientId, requests] of activeRequests) { + for (const { prepared, requestId } of requests.values()) { + if (prepared.binding === binding && cancelRequest(clientId, requestId)) { + interrupts.push(writeFromClient(clientId, { _tag: "Interrupt", requestId })) + } + } + } + return Effect.all(interrupts, { discard: true }) + }) + } const { notificationDelivery, notifications } = internalState.get(server)! yield* Effect.acquireRelease( @@ -1534,11 +1558,18 @@ const mcpStdioSerialization = ( * remain valid. The surrounding HTTP server remains responsible for binding * to an appropriate interface and installing authentication. * + * With `allowSessionTermination`, a DELETE carrying `Mcp-Session-Id` ends that + * session with `204` and interrupts its in-flight requests; later requests with + * that id get `404`. DELETE validates the session and `MCP-Protocol-Version` + * headers like POST. Without the option, or when only sessionless revisions + * such as `v2026_07_28` are configured, DELETE returns `405`. Any caller + * holding a session id can end that session, so authenticate requests in the + * surrounding router. + * * `layerHttp` always implements the single-endpoint Streamable HTTP topology. * Using `v2024_11_05` here is a custom compatibility transport for that * revision's schema. It does not implement the historical two-endpoint - * HTTP+SSE transport, GET SSE, event resumption, session expiry, or client - * session termination. + * HTTP+SSE transport, GET SSE, event resumption, or session expiry. * * @see {@link layerStdio} for exposing the server over stdio * @see {@link layer} for the base MCP server layer without a transport protocol @@ -1558,24 +1589,10 @@ export const layerHttp = (options: { readonly protocols: Arr.NonEmptyReadonlyArray readonly extensions?: ServerExtensions | undefined readonly allowedOrigins?: ReadonlyArray | undefined + readonly allowSessionTermination?: boolean | undefined }): Layer.Layer => { const runtime = McpRuntime.layer(options.protocols) - const methodNotAllowedResponse = HttpServerResponse.empty({ - status: 405, - headers: { allow: "POST" } - }) - const methodNotAllowed = (request: HttpServerRequest.HttpServerRequest) => - isAllowedMcpOrigin(request, options.allowedOrigins) - ? Effect.succeed(methodNotAllowedResponse) - : Effect.succeed(HttpServerResponse.empty({ status: 403 })) - const routes = Layer.mergeAll( - HttpRouter.add("GET", options.path, methodNotAllowed), - HttpRouter.add("PUT", options.path, methodNotAllowed), - HttpRouter.add("PATCH", options.path, methodNotAllowed), - HttpRouter.add("DELETE", options.path, methodNotAllowed), - HttpRouter.add("OPTIONS", options.path, methodNotAllowed) - ) - return Layer.merge(layerWithRuntime(options, "http"), routes).pipe( + return layerWithRuntime(options, "http").pipe( Layer.provide(layerMcpProtocolHttp(options)), Layer.provide(runtime), Layer.provide(RpcSerialization.layerJsonRpc()) @@ -1585,6 +1602,7 @@ export const layerHttp = (options: { const layerMcpProtocolHttp = (options: { readonly path: HttpRouter.PathInput readonly allowedOrigins?: ReadonlyArray | undefined + readonly allowSessionTermination?: boolean | undefined }): Layer.Layer< RpcServer.Protocol, never, @@ -1596,6 +1614,29 @@ const layerMcpProtocolHttp = (options: { Effect.provideService(RpcSerialization.RpcSerialization, mcpHttpSerialization) ) const router = yield* HttpRouter.HttpRouter + const allowSessionTermination = options.allowSessionTermination === true && + runtime.protocols.some((protocol) => protocol.runtime._tag === "Stateful") + const forbidden = Effect.succeed(HttpServerResponse.empty({ status: 403 })) + const withAllowedOrigin = ( + handler: (request: HttpServerRequest.HttpServerRequest) => Effect.Effect + ) => + (request: HttpServerRequest.HttpServerRequest) => + isAllowedMcpOrigin(request, options.allowedOrigins) ? handler(request) : forbidden + const methodNotAllowedResponse = Effect.succeed(HttpServerResponse.empty({ + status: 405, + headers: { allow: allowSessionTermination ? "POST, DELETE" : "POST" } + })) + const methodNotAllowed = withAllowedOrigin(() => methodNotAllowedResponse) + for (const method of ["GET", "PUT", "PATCH", "OPTIONS"] as const) { + yield* router.add(method, options.path, methodNotAllowed) + } + yield* router.add( + "DELETE", + options.path, + allowSessionTermination + ? withAllowedOrigin((request) => runtime.terminateHttpSession(request.headers)) + : methodNotAllowed + ) yield* router.add("POST", options.path, (request) => { if (!isAllowedMcpOrigin(request, options.allowedOrigins)) { return Effect.succeed(HttpServerResponse.empty({ status: 403 })) diff --git a/packages/effect/src/ai/internal/mcpRuntime.ts b/packages/effect/src/ai/internal/mcpRuntime.ts index b388d178053..2f520b5f66b 100644 --- a/packages/effect/src/ai/internal/mcpRuntime.ts +++ b/packages/effect/src/ai/internal/mcpRuntime.ts @@ -17,6 +17,7 @@ import * as Predicate from "../../Predicate.ts" import * as Result from "../../Result.ts" import * as RpcGroup from "../../rpc/RpcGroup.ts" import type * as RpcMessage from "../../rpc/RpcMessage.ts" +import type * as Scope from "../../Scope.ts" import type * as PublicMcpProtocol from "../McpProtocol.ts" import * as PublicMcpSchema from "../McpSchema.ts" import type * as McpCore from "./mcpCore.ts" @@ -175,6 +176,10 @@ export interface ServerRuntimeShape { fallback: LogLevel.LogLevel ) => LogLevel.LogLevel readonly disconnect: (clientId: number) => void + readonly terminateHttpSession: (headers: Headers.Headers) => Effect.Effect + readonly onSessionTerminated: ( + listener: (binding: RequestBinding) => Effect.Effect + ) => Effect.Effect readonly deliveryClientIds: () => Iterable readonly canDeliver: ( clientId: number, @@ -228,9 +233,32 @@ export const make = Effect.fnUntraced(function*( statelessProtocol = protocol } const registry = yield* McpProtocolRegistry.make(protocols) - const selectHttpProtocol = (headers: Headers.Headers, input: unknown): HttpProtocolSelection => { + const sessionTerminationListeners = new Set<(binding: RequestBinding) => Effect.Effect>() + const selectHttpSession = (headers: Headers.Headers, isInitialize: boolean): HttpProtocolSelection => { const protocolVersion = headers[MCP_PROTOCOL_VERSION_HEADER] const sessionId = headers[MCP_SESSION_ID_HEADER] + const binding = sessionId === undefined ? undefined : stateful?.resolveSessionId(sessionId) + if (sessionId !== undefined && binding === undefined) { + return { _tag: "Rejected", status: 404 } + } + if ( + !isInitialize && + protocolVersion !== undefined && + !registry.protocols.some((protocol) => protocol.protocolVersion === protocolVersion) + ) { + return { _tag: "Rejected", status: 400 } + } + if ( + !isInitialize && + binding?.protocol.runtime.transport.http.requiresVersionHeader === true && + protocolVersion !== binding.protocol.protocolVersion + ) { + return { _tag: "Rejected", status: 400 } + } + return { _tag: "Accepted", binding, protocol: binding?.protocol } + } + const selectHttpProtocol = (headers: Headers.Headers, input: unknown): HttpProtocolSelection => { + const protocolVersion = headers[MCP_PROTOCOL_VERSION_HEADER] const inputRecord = asRecord(input) const metadata = asRecord(asRecord(inputRecord?.params)?._meta) const claim = protocolVersionClaim(metadata) @@ -299,25 +327,7 @@ export const make = Effect.fnUntraced(function*( } return { _tag: "Accepted", binding: undefined, protocol: statelessProtocol } } - const binding = sessionId === undefined ? undefined : stateful?.resolveSessionId(sessionId) - if (sessionId !== undefined && binding === undefined) { - return { _tag: "Rejected", status: 404 } - } - if ( - !isInitialize && - protocolVersion !== undefined && - !registry.protocols.some((protocol) => protocol.protocolVersion === protocolVersion) - ) { - return { _tag: "Rejected", status: 400 } - } - if ( - !isInitialize && - binding?.protocol.runtime.transport.http.requiresVersionHeader === true && - protocolVersion !== binding.protocol.protocolVersion - ) { - return { _tag: "Rejected", status: 400 } - } - return { _tag: "Accepted", binding, protocol: binding?.protocol } + return selectHttpSession(headers, isInitialize) } return ServerRuntime.of({ protocols: registry.protocols, @@ -464,6 +474,33 @@ export const make = Effect.fnUntraced(function*( }, effectLogLevel: (clientId, headers, fallback) => stateful?.effectLogLevel(clientId, headers, fallback) ?? fallback, disconnect: (clientId) => stateful?.disconnect(clientId), + terminateHttpSession: (headers) => + Effect.suspend(() => { + const sessionId = headers[MCP_SESSION_ID_HEADER] + if (sessionId === undefined) { + return Effect.succeed(HttpServerResponse.empty({ status: 400 })) + } + const selection = selectHttpSession(headers, false) + if (selection._tag === "Rejected") { + return Effect.succeed(HttpServerResponse.empty({ status: selection.status })) + } + const binding = selection.binding! + stateful!.terminateSession(sessionId) + return Effect.as( + Effect.forEach(sessionTerminationListeners, (listener) => listener(binding), { discard: true }), + HttpServerResponse.empty({ status: 204 }) + ) + }), + onSessionTerminated: (listener) => + Effect.acquireRelease( + Effect.sync(() => { + sessionTerminationListeners.add(listener) + }), + () => + Effect.sync(() => { + sessionTerminationListeners.delete(listener) + }) + ), deliveryClientIds: () => stateful?.initializedClientIds() ?? [], canDeliver: (clientId, headers, notification, fallback) => stateful?.canDeliver(clientId, headers, notification, fallback) ?? true, diff --git a/packages/effect/src/ai/internal/mcpStatefulRuntime.ts b/packages/effect/src/ai/internal/mcpStatefulRuntime.ts index 19c41ab07b1..bcede31008d 100644 --- a/packages/effect/src/ai/internal/mcpStatefulRuntime.ts +++ b/packages/effect/src/ai/internal/mcpStatefulRuntime.ts @@ -47,6 +47,7 @@ export interface StatefulRuntime { readonly registerConnection: (clientId: number, registration: Registration) => Binding readonly resolve: (clientId: number, headers: Headers.Headers) => Binding | undefined readonly resolveSessionId: (sessionId: string) => Binding | undefined + readonly terminateSession: (sessionId: string) => boolean readonly setLogLevel: ( level: PublicMcpSchema.LoggingLevel, clientId: number, @@ -123,6 +124,7 @@ export const make = (): StatefulRuntime => { }, resolve: resolveSession, resolveSessionId: (sessionId) => bySessionId.get(sessionId), + terminateSession: (sessionId) => bySessionId.delete(sessionId), setLogLevel: (level, clientId, headers) => Effect.sync(() => { const session = resolveSession(clientId, headers) diff --git a/packages/effect/src/rpc/RpcServer.ts b/packages/effect/src/rpc/RpcServer.ts index a6d06f709f3..c5e260d81c3 100644 --- a/packages/effect/src/rpc/RpcServer.ts +++ b/packages/effect/src/rpc/RpcServer.ts @@ -359,11 +359,14 @@ export const makeNoSerialization: ( ) const fiber = trackFiber( runFork( - effect, + // Register before the handler runs to catch synchronous cancellation. + Effect.withFiber((fiber) => { + client.fibers.set(request.id, fiber) + return effect + }), isUninterruptible ? { uninterruptible: true } : undefined ) ) - client.fibers.set(request.id, fiber) fiber.addObserver(function onExit(exit: Exit.Exit): void { if (deferred) { const fiber = trackFiber(runFork(Effect.onExit(Deferred.await(deferred), (exit) => diff --git a/packages/effect/test/ai/McpServer/McpServer.test.ts b/packages/effect/test/ai/McpServer/McpServer.test.ts index 4f235c71143..a85313a40d1 100644 --- a/packages/effect/test/ai/McpServer/McpServer.test.ts +++ b/packages/effect/test/ai/McpServer/McpServer.test.ts @@ -2180,6 +2180,176 @@ describe("McpServer", () => { yield* client.ping({}) })) + it.effect("terminates an HTTP session on DELETE when allowSessionTermination is set", () => + Effect.gen(function*() { + const harness = yield* makeHttpHarness(makeServerLayer({ + name: "SessionDelete", + protocols: [McpProtocol.v2025_11_25], + allowSessionTermination: true + })) + const headers = yield* initializeHttpSession(harness, McpProtocol.v2025_11_25) + strictEqual((yield* harness.delete(headers)).status, 204) + strictEqual((yield* harness.post({ jsonrpc: "2.0", id: 3, method: "ping" }, headers)).status, 404) + strictEqual((yield* harness.delete(headers)).status, 404) + })) + + it.effect("rejects DELETE with a missing or unknown session id when termination is enabled", () => + Effect.gen(function*() { + const harness = yield* makeHttpHarness(makeServerLayer({ + name: "SessionDelete", + allowSessionTermination: true + })) + strictEqual((yield* harness.delete()).status, 400) + strictEqual((yield* harness.delete({ "Mcp-Session-Id": "unknown" })).status, 404) + })) + + it.effect("rejects disallowed DELETE Origins without terminating the session", () => + Effect.gen(function*() { + const harness = yield* makeHttpHarness(makeServerLayer({ + name: "SessionDelete", + protocols: [McpProtocol.v2025_11_25], + allowSessionTermination: true + })) + const headers = yield* initializeHttpSession(harness, McpProtocol.v2025_11_25) + strictEqual((yield* harness.delete({ ...headers, origin: "https://blocked.example" })).status, 403) + strictEqual((yield* harness.post({ jsonrpc: "2.0", id: 2, method: "ping" }, headers)).status, 200) + strictEqual((yield* harness.delete({ ...headers, origin: "https://allowed.example" })).status, 204) + })) + + it.effect("rejects unsupported or mismatched DELETE protocol versions without terminating the session", () => + Effect.gen(function*() { + const harness = yield* makeHttpHarness(makeServerLayer({ + name: "SessionDelete", + protocols: [McpProtocol.v2025_06_18, McpProtocol.v2025_11_25], + allowSessionTermination: true + })) + const headers = yield* initializeHttpSession(harness, McpProtocol.v2025_11_25) + for (const version of ["9999-01-01", "2025-06-18"]) { + strictEqual((yield* harness.delete({ ...headers, "Mcp-Protocol-Version": version })).status, 400) + const ping = yield* harness.post({ jsonrpc: "2.0", id: 2, method: "ping" }, headers) + strictEqual(ping.status, 200) + assert.deepStrictEqual(yield* readMcpHttpResponse(ping), { jsonrpc: "2.0", id: 2, result: {} }) + } + strictEqual((yield* harness.delete(headers)).status, 204) + })) + + it.effect("refuses DELETE for a stateless-only server even when termination is enabled", () => + Effect.gen(function*() { + const harness = yield* makeHttpHarness(makeServerLayer({ + name: "StatelessDelete", + protocols: [McpProtocol.v2026_07_28], + allowSessionTermination: true + })) + const response = yield* harness.delete() + strictEqual(response.status, 405) + strictEqual(response.headers.get("Allow"), "POST") + })) + + it.effect("advertises DELETE on unsupported methods when stateful termination is enabled", () => + Effect.gen(function*() { + const harness = yield* makeHttpHarness(makeServerLayer({ + name: "SessionDelete", + protocols: [McpProtocol.v2026_07_28, McpProtocol.v2025_11_25], + allowSessionTermination: true + })) + for (const method of ["GET", "PUT", "PATCH", "HEAD"] as const) { + const response = yield* Effect.promise(() => harness.handler(new Request("http://localhost/mcp", { method }))) + strictEqual(response.status, 405) + strictEqual(response.headers.get("Allow"), "POST, DELETE") + } + })) + + it.effect("interrupts an in-flight tool on DELETE and returns 404 for its pending POST", () => + Effect.gen(function*() { + const started = yield* Deferred.make() + const interrupted = yield* Deferred.make() + const registration = Layer.effectDiscard(Effect.gen(function*() { + const server = yield* McpServer.McpServer + yield* server.addTool({ + tool: new McpSchema.Tool({ name: "Blocked", inputSchema: { type: "object" } }), + annotations: Context.empty(), + handle: () => + Deferred.succeed(started, undefined).pipe( + Effect.andThen(Effect.never), + Effect.onInterrupt(() => Deferred.succeed(interrupted, undefined)) + ) + }) + })) + const harness = yield* makeHttpHarness(registration.pipe(Layer.provideMerge( + makeServerLayer({ + name: "SessionDelete", + protocols: [McpProtocol.v2025_11_25], + allowSessionTermination: true + }) + ))) + const headers = yield* initializeHttpSession(harness, McpProtocol.v2025_11_25) + const pending = yield* harness.post({ + jsonrpc: "2.0", + id: "blocked-tool", + method: "tools/call", + params: { name: "Blocked", arguments: {} } + }, headers).pipe(Effect.forkScoped) + yield* Deferred.await(started) + strictEqual((yield* harness.delete(headers)).status, 204) + yield* Deferred.await(interrupted) + strictEqual((yield* Fiber.join(pending)).status, 404) + })) + + it.effect("closes an already-started elicitation stream on DELETE", () => + Effect.gen(function*() { + const interrupted = yield* Deferred.make() + const registration = Layer.effectDiscard(Effect.gen(function*() { + const server = yield* McpServer.McpServer + yield* server.addTool({ + tool: new McpSchema.Tool({ name: "Authorize", inputSchema: { type: "object" } }), + annotations: Context.empty(), + handle: () => + Effect.gen(function*() { + const client = yield* Effect.serviceOption(McpSchema.McpServerClient).pipe( + Effect.flatMap(Effect.fromOption) + ) + const reverseClient = yield* client.getClient + yield* reverseClient.elicit( + Schema.decodeUnknownSync(McpSchema.Elicit.payloadSchema)({ + mode: "url", + message: "Authorize access", + url: "https://example.com/authorize", + elicitationId: "authorization-1" + }) + ) + return new McpSchema.CallToolResult({ content: [] }) + }).pipe( + Effect.scoped, + Effect.orDie, + Effect.onInterrupt(() => Deferred.succeed(interrupted, undefined)) + ) + }) + })) + const harness = yield* makeHttpHarness(registration.pipe(Layer.provideMerge( + makeServerLayer({ + name: "SessionDelete", + protocols: [McpProtocol.v2025_11_25], + allowSessionTermination: true + }) + ))) + const headers = yield* initializeHttpSession(harness, McpProtocol.v2025_11_25, { elicitation: { url: {} } }) + const response = yield* harness.post({ + jsonrpc: "2.0", + id: "authorize-tool", + method: "tools/call", + params: { name: "Authorize", arguments: {} } + }, headers) + const stream = makeMcpSseReader(response) + yield* Effect.addFinalizer(() => stream.cancel) + const elicitation = yield* stream.take() + strictEqual(elicitation.method, "elicitation/create") + assert.isDefined(elicitation.id) + strictEqual((yield* harness.delete(headers)).status, 204) + yield* Deferred.await(interrupted) + const remaining = yield* stream.drain() + assert.isFalse(remaining.some((message) => message.id === "authorize-tool")) + })) + it.effect("returns an empty 202 for notifications and responses and remains successful for request POSTs", () => Effect.gen(function*() { const { client, httpClient } = yield* makeRouterTestClient(HttpRouter.cors()) diff --git a/packages/effect/test/ai/McpServer/TestUtils/McpHttpHarness.ts b/packages/effect/test/ai/McpServer/TestUtils/McpHttpHarness.ts index e08c8e09232..2493776177f 100644 --- a/packages/effect/test/ai/McpServer/TestUtils/McpHttpHarness.ts +++ b/packages/effect/test/ai/McpServer/TestUtils/McpHttpHarness.ts @@ -52,19 +52,23 @@ export const makeHttpHarness = Effect.fnUntraced(function*( ) const post = (body: unknown, headers?: HeadersInit) => postText(JSON.stringify(body), headers) + const deleteSession = (headers?: HeadersInit) => + Effect.promise(() => handler(new Request(MCP_ENDPOINT, { method: "DELETE", headers: headers ?? {} }))) return { handler, fetch, post, postText, + delete: deleteSession, responses } as const }) export const initializeHttpSession = Effect.fnUntraced(function*( harness: Effect.Success>, - selectedProtocol: McpProtocol.ProtocolAdapter + selectedProtocol: McpProtocol.ProtocolAdapter, + capabilities: Record = {} ) { const response = yield* harness.post({ jsonrpc: "2.0", @@ -72,7 +76,7 @@ export const initializeHttpSession = Effect.fnUntraced(function*( method: "initialize", params: { protocolVersion: selectedProtocol.protocolVersion, - capabilities: {}, + capabilities, clientInfo: { name: "test", version: "1.0.0" } } }) diff --git a/packages/effect/test/ai/McpServer/TestUtils/McpServerLayer.ts b/packages/effect/test/ai/McpServer/TestUtils/McpServerLayer.ts index 8995b0a51e5..625d38b68d3 100644 --- a/packages/effect/test/ai/McpServer/TestUtils/McpServerLayer.ts +++ b/packages/effect/test/ai/McpServer/TestUtils/McpServerLayer.ts @@ -20,6 +20,7 @@ export const makeServerLayer = (options: { | undefined readonly extensions?: Readonly> | undefined readonly allowedOrigins?: ReadonlyArray | undefined + readonly allowSessionTermination?: boolean | undefined }) => McpServer.layerHttp({ name: options.name, @@ -28,7 +29,8 @@ export const makeServerLayer = (options: { path: "/mcp", protocols: options.protocols ?? [McpProtocol.v2025_06_18], allowedOrigins: ["https://allowed.example"], - extensions: options.extensions + extensions: options.extensions, + allowSessionTermination: options.allowSessionTermination }).pipe( Layer.provideMerge(Layer.succeed( References.CurrentLoggers,