ssePollingExample.js 4.3 KB

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