Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,9 @@ jobs:
- name: Apply patches
run: npx patch-package

- name: Test
run: bun run test

- name: Build
run: bun run build

Expand Down
1 change: 1 addition & 0 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
"main": "index.js",
"scripts": {
"build": "tsc && mcp-build",
"test": "node --test tests/*.test.mjs",
"watch": "tsc --watch",
"start": "node dist/index.js",
"mcp-inspect": "mcp-inspector --transport http --server-url http://localhost:1234",
Expand Down
107 changes: 97 additions & 10 deletions patches/mcp-framework+0.2.18.patch
Original file line number Diff line number Diff line change
@@ -1,8 +1,7 @@
diff --git a/node_modules/mcp-framework/dist/transports/http/server.js b/node_modules/mcp-framework/dist/transports/http/server.js
index b535515..ad7f146 100644
--- a/node_modules/mcp-framework/dist/transports/http/server.js
+++ b/node_modules/mcp-framework/dist/transports/http/server.js
@@ -17,6 +17,8 @@ export class HttpStreamTransport extends AbstractTransport {
@@ -17,10 +17,12 @@ export class HttpStreamTransport extends AbstractTransport {
_config;
_oauthMetadata;
_transports = {};
Expand All @@ -11,9 +10,17 @@ index b535515..ad7f146 100644
constructor(config = {}) {
super();
this._config = config;
@@ -90,6 +92,23 @@ export class HttpStreamTransport extends AbstractTransport {
- this._port = config.port || 8080;
+ this._port = config.port ?? 8080;
this._endpoint = config.endpoint || '/mcp';
this._enableJsonResponse = config.responseMode === 'batch';
// Initialize OAuth metadata if OAuth provider is configured
@@ -88,12 +90,33 @@ export class HttpStreamTransport extends AbstractTransport {
this._onclose?.();
});
this._server.listen(this._port, () => {
logger.info(`HTTP server listening on port ${this._port}, endpoint ${this._endpoint}`);
- logger.info(`HTTP server listening on port ${this._port}, endpoint ${this._endpoint}`);
+ logger.info(`HTTP server listening on port ${this.port}, endpoint ${this._endpoint}`);
this._isRunning = true;
+ // Start periodic session cleanup
+ const timeoutMs = this._config.session?.sessionTimeout || 300000;
Expand All @@ -35,32 +42,101 @@ index b535515..ad7f146 100644
resolve();
});
});
@@ -117,6 +136,7 @@ export class HttpStreamTransport extends AbstractTransport {
}
+ get port() {
+ const address = this._server?.address();
+ return address && typeof address === 'object' ? address.port : this._port;
+ }
async handleMcpRequest(req, res) {
const sessionId = req.headers['mcp-session-id'];
let transport;
@@ -109,6 +132,10 @@ export class HttpStreamTransport extends AbstractTransport {
return;
authData = authResult.data || {};
}
+ if (this._config.session?.enabled === false) {
+ await this.handleStatelessRequest(req, res, body, authData);
+ return;
+ }
// Allow re-initialization even when a stale session ID is provided.
// Clients like Cline may keep sending the old session ID header after
// a session is lost (server restart, transport error, etc.).
@@ -117,6 +144,7 @@ export class HttpStreamTransport extends AbstractTransport {
if (sessionId && this._transports[sessionId]) {
// Existing session
transport = this._transports[sessionId];
+ this._sessionLastActivity[sessionId] = Date.now();
logger.debug(`Reusing existing session: ${sessionId}`);
}
else if (isInitialize || isReInitialize) {
@@ -131,6 +151,7 @@ export class HttpStreamTransport extends AbstractTransport {
@@ -131,6 +159,7 @@ export class HttpStreamTransport extends AbstractTransport {
onsessioninitialized: (sessionId) => {
logger.info(`Session initialized: ${sessionId}`);
this._transports[sessionId] = transport;
+ this._sessionLastActivity[sessionId] = Date.now();
},
enableJsonResponse: this._enableJsonResponse,
});
@@ -138,6 +159,7 @@ export class HttpStreamTransport extends AbstractTransport {
@@ -138,6 +167,7 @@ export class HttpStreamTransport extends AbstractTransport {
if (transport.sessionId) {
logger.info(`Transport closed for session: ${transport.sessionId}`);
delete this._transports[transport.sessionId];
+ delete this._sessionLastActivity[transport.sessionId];
}
};
transport.onerror = (error) => {
@@ -231,16 +253,16 @@ export class HttpStreamTransport extends AbstractTransport {
await transport.send(message);
@@ -172,6 +202,23 @@ export class HttpStreamTransport extends AbstractTransport {
await transport.handleRequest(req, res, body);
});
}
+ async handleStatelessRequest(req, res, body, authData) {
+ const transport = new StreamableHTTPServerTransport({
+ sessionIdGenerator: undefined,
+ enableJsonResponse: this._enableJsonResponse,
+ });
+ transport.onerror = (error) => {
+ logger.error(`Stateless transport error: ${error}`);
+ };
+ transport.onmessage = async (message) => {
+ if (this._onmessage) {
+ await this._onmessage(message);
+ }
+ };
+ await requestContext.run({ ...authData, httpTransport: transport }, async () => {
+ await transport.handleRequest(req, res, body);
+ });
+ }
async readRequestBody(req) {
return new Promise((resolve, reject) => {
let body = '';
@@ -214,11 +261,20 @@ export class HttpStreamTransport extends AbstractTransport {
id: null,
}));
}
- async send(message) {
+ async send(message, options) {
if (!this._isRunning) {
logger.warn('Attempted to send message, but HTTP transport is not running');
return;
}
+ if (this._config.session?.enabled === false) {
+ const transport = requestContext.getStore()?.httpTransport;
+ if (!transport) {
+ logger.warn('No active stateless request to send message to');
+ return;
+ }
+ await transport.send(message, options);
+ return;
+ }
const activeSessions = Object.entries(this._transports);
if (activeSessions.length === 0) {
logger.warn('No active sessions to send message to');
@@ -228,19 +284,19 @@ export class HttpStreamTransport extends AbstractTransport {
const failedSessions = [];
for (const [sessionId, transport] of activeSessions) {
try {
- await transport.send(message);
+ await transport.send(message, options);
}
catch (error) {
- logger.error(`Error sending message to session ${sessionId}: ${error}`);
Expand All @@ -82,7 +158,7 @@ index b535515..ad7f146 100644
}
}
async close() {
@@ -256,6 +284,11 @@ export class HttpStreamTransport extends AbstractTransport {
@@ -256,6 +312,11 @@ export class HttpStreamTransport extends AbstractTransport {
}
}
this._transports = {};
Expand All @@ -94,3 +170,14 @@ index b535515..ad7f146 100644
if (this._server) {
this._server.close();
this._server = undefined;
diff --git a/node_modules/mcp-framework/dist/transports/http/server.d.ts b/node_modules/mcp-framework/dist/transports/http/server.d.ts
--- a/node_modules/mcp-framework/dist/transports/http/server.d.ts
+++ b/node_modules/mcp-framework/dist/transports/http/server.d.ts
@@ -13,6 +13,7 @@ export declare class HttpStreamTransport extends AbstractTransport {
private _transports;
constructor(config?: HttpStreamTransportConfig);
start(): Promise<void>;
+ get port(): number;
private handleMcpRequest;
private readRequestBody;
private setCorsHeaders;
10 changes: 4 additions & 6 deletions scripts/download-content.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,13 +22,11 @@ export async function downloadDocs() {
}

export async function downloadExamples() {
console.log(`Downloading examples from appwrite/appwrite (version: ${appwriteExamplesBranch})`);
console.log(`Downloading examples from appwrite/specs (version: ${appwriteExamplesBranch})`);

const owner = "appwrite";
const repo = "appwrite";
const docsSubdirPath = `docs/examples/${appwriteExamplesBranch}`;
// The version-pinned example folders (docs/examples/<version>) only live on the
// `main` branch; they were removed from the per-version branches like 1.8.x.
const repo = "specs";
const docsSubdirPath = `examples/${appwriteExamplesBranch}`;
const ref = "main";

console.log(`Downloading examples from ${owner}/${repo}/${docsSubdirPath} to ${examplesTargetDir}`);
Expand All @@ -54,4 +52,4 @@ async function main() {
await createTableOfContents();
}

await main();
await main();
10 changes: 5 additions & 5 deletions src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,18 +22,18 @@ const server = new MCPServer({
responseMode: "stream", // Response mode: "batch" or "stream" (default: "batch")
batchTimeout: 30000, // Timeout for batch responses in ms (default: 30000)
session: {
enabled: true,
headerName: "Mcp-Session-Id",
allowClientTermination: true,
sessionTimeout: 300000, // 5 minutes
// The service runs with multiple replicas. Keeping MCP sessions in process
// memory makes follow-up requests fail whenever the load balancer sends them
// to another replica, so every request must be independently routable.
enabled: false,
},
cors: {
// CORS configuration
allowOrigin: "*",
allowMethods: "GET, POST, DELETE, OPTIONS",
allowHeaders:
"Content-Type, Accept, Authorization, x-api-key, Mcp-Session-Id, Last-Event-ID",
exposeHeaders: "Content-Type, Authorization, x-api-key, Mcp-Session-Id",
exposeHeaders: "Content-Type, Authorization, x-api-key",
maxAge: "86400",
},
},
Expand Down
11 changes: 11 additions & 0 deletions tests/fixtures/tools/ping.tool.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
import { MCPTool } from "mcp-framework";

export default class PingTool extends MCPTool {
name = "ping";
description = "Return a pong response";
schema = {};

async execute() {
return "pong";
}
}
Loading