ssePollingExample.js 4.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107
  1. "use strict";
  2. var __importDefault = (this && this.__importDefault) || function (mod) {
  3. return (mod && mod.__esModule) ? mod : { "default": mod };
  4. };
  5. Object.defineProperty(exports, "__esModule", { value: true });
  6. const node_crypto_1 = require("node:crypto");
  7. const mcp_js_1 = require("../../server/mcp.js");
  8. const express_js_1 = require("../../server/express.js");
  9. const streamableHttp_js_1 = require("../../server/streamableHttp.js");
  10. const inMemoryEventStore_js_1 = require("../shared/inMemoryEventStore.js");
  11. const cors_1 = __importDefault(require("cors"));
  12. // Factory to create a new MCP server per session.
  13. // Each session needs its own server+transport pair to avoid cross-session contamination.
  14. const getServer = () => {
  15. const server = new mcp_js_1.McpServer({
  16. name: 'sse-polling-example',
  17. version: '1.0.0'
  18. }, {
  19. capabilities: { logging: {} }
  20. });
  21. // Register a long-running tool that demonstrates server-initiated disconnect
  22. server.tool('long-task', 'A long-running task that sends progress updates. Server will disconnect mid-task to demonstrate polling.', {}, async (_args, extra) => {
  23. const sleep = (ms) => new Promise(resolve => setTimeout(resolve, ms));
  24. console.log(`[${extra.sessionId}] Starting long-task...`);
  25. // Send first progress notification
  26. await server.sendLoggingMessage({
  27. level: 'info',
  28. data: 'Progress: 25% - Starting work...'
  29. }, extra.sessionId);
  30. await sleep(1000);
  31. // Send second progress notification
  32. await server.sendLoggingMessage({
  33. level: 'info',
  34. data: 'Progress: 50% - Halfway there...'
  35. }, extra.sessionId);
  36. await sleep(1000);
  37. // Server decides to disconnect the client to free resources
  38. // Client will reconnect via GET with Last-Event-ID after the transport's retryInterval
  39. // Use extra.closeSSEStream callback - available when eventStore is configured
  40. if (extra.closeSSEStream) {
  41. console.log(`[${extra.sessionId}] Closing SSE stream to trigger client polling...`);
  42. extra.closeSSEStream();
  43. }
  44. // Continue processing while client is disconnected
  45. // Events are stored in eventStore and will be replayed on reconnect
  46. await sleep(500);
  47. await server.sendLoggingMessage({
  48. level: 'info',
  49. data: 'Progress: 75% - Almost done (sent while client disconnected)...'
  50. }, extra.sessionId);
  51. await sleep(500);
  52. await server.sendLoggingMessage({
  53. level: 'info',
  54. data: 'Progress: 100% - Complete!'
  55. }, extra.sessionId);
  56. console.log(`[${extra.sessionId}] Task complete`);
  57. return {
  58. content: [
  59. {
  60. type: 'text',
  61. text: 'Long task completed successfully!'
  62. }
  63. ]
  64. };
  65. });
  66. return server;
  67. };
  68. // Set up Express app
  69. const app = (0, express_js_1.createMcpExpressApp)();
  70. app.use((0, cors_1.default)());
  71. // Create event store for resumability
  72. const eventStore = new inMemoryEventStore_js_1.InMemoryEventStore();
  73. // Track transports by session ID for session reuse
  74. const transports = new Map();
  75. // Handle all MCP requests
  76. app.all('/mcp', async (req, res) => {
  77. const sessionId = req.headers['mcp-session-id'];
  78. // Reuse existing transport or create new one
  79. let transport = sessionId ? transports.get(sessionId) : undefined;
  80. if (!transport) {
  81. transport = new streamableHttp_js_1.StreamableHTTPServerTransport({
  82. sessionIdGenerator: () => (0, node_crypto_1.randomUUID)(),
  83. eventStore,
  84. retryInterval: 2000, // Default retry interval for priming events
  85. onsessioninitialized: id => {
  86. console.log(`[${id}] Session initialized`);
  87. transports.set(id, transport);
  88. }
  89. });
  90. // Create a new server per session and connect it to the transport
  91. const server = getServer();
  92. await server.connect(transport);
  93. }
  94. await transport.handleRequest(req, res, req.body);
  95. });
  96. // Start the server
  97. const PORT = 3001;
  98. app.listen(PORT, () => {
  99. console.log(`SSE Polling Example Server running on http://localhost:${PORT}/mcp`);
  100. console.log('');
  101. console.log('This server demonstrates SEP-1699 SSE polling:');
  102. console.log('- retryInterval: 2000ms (client waits 2s before reconnecting)');
  103. console.log('- eventStore: InMemoryEventStore (events are persisted for replay)');
  104. console.log('');
  105. console.log('Try calling the "long-task" tool to see server-initiated disconnect in action.');
  106. });
  107. //# sourceMappingURL=ssePollingExample.js.map