(
principal: Principal,
resource: McpResource,
request: Request,
)
| 447 | |
| 448 | /** Open a new session: build the server, connect a transport, drive the request. */ |
| 449 | const openSession = ( |
| 450 | principal: Principal, |
| 451 | resource: McpResource, |
| 452 | request: Request, |
| 453 | ): Effect.Effect<McpDispatchResult> => { |
| 454 | let createdSessionId: string | null = null; |
| 455 | return buildServer(principal, { |
| 456 | ...buildOptionsFor(request, () => createdSessionId), |
| 457 | resource, |
| 458 | }).pipe( |
| 459 | Effect.flatMap(({ mcpServer, engine, executor, close }) => |
| 460 | Effect.gen(function* () { |
| 461 | const transport = new WebStandardStreamableHTTPServerTransport({ |
| 462 | sessionIdGenerator: () => crypto.randomUUID(), |
| 463 | enableJsonResponse: true, |
| 464 | onsessioninitialized: (sid) => { |
| 465 | createdSessionId = sid; |
| 466 | transports.set(sid, transport); |
| 467 | servers.set(sid, mcpServer); |
| 468 | owners.set(sid, { principal, resource }); |
| 469 | engines.set(sid, engine); |
| 470 | if (executor) executors.set(sid, executor); |
| 471 | if (close) closers.set(sid, close); |
| 472 | lastSeen.set(sid, Date.now()); |
| 473 | }, |
| 474 | onsessionclosed: (sid) => void dispose(sid, { server: true }), |
| 475 | }); |
| 476 | transport.onclose = () => { |
| 477 | const sid = transport.sessionId; |
| 478 | if (sid) void dispose(sid, { server: true }); |
| 479 | }; |
| 480 | yield* Effect.promise(() => mcpServer.connect(transport)); |
| 481 | // The session id is minted on the first (initialize) request, so we |
| 482 | // drive `handleRequest` here; if no id results we close eagerly. |
| 483 | return yield* runHandleRequest( |
| 484 | transport, |
| 485 | request, |
| 486 | orgWriteAccessForPrincipal(principal), |
| 487 | () => { |
| 488 | // Nothing was ever registered under a session id, so `dispose` has |
| 489 | // no entry to work from — release the handles by hand, engine |
| 490 | // and executor included. |
| 491 | void ignoreClose(null, "transport", () => transport.close()); |
| 492 | void ignoreClose(null, "server", () => mcpServer.close()); |
| 493 | void shutdownEngine(null, engine); |
| 494 | void shutdownExecutor(null, executor); |
| 495 | if (close) void ignoreClose(null, "session", close); |
| 496 | }, |
| 497 | ); |
| 498 | }), |
| 499 | ), |
| 500 | // A build failure has nowhere typed to go in the envelope; render a 500. |
| 501 | Effect.catchTags({ |
| 502 | McpEngineBuildError: () => |
| 503 | Effect.succeed(jsonRpcError(500, -32603, "Internal server error")), |
| 504 | McpPassthroughUnavailableError: () => |
| 505 | Effect.succeed(jsonRpcError(500, -32603, "Internal server error")), |
| 506 | }), |
no test coverage detected