|
@@ -17,8 +17,10 @@ export interface HttpServerOptions {
|
|
|
export class HttpMcpServer {
|
|
export class HttpMcpServer {
|
|
|
private app: express.Application;
|
|
private app: express.Application;
|
|
|
private server: Server;
|
|
private server: Server;
|
|
|
|
|
+ private httpServer: any; // HTTP server instance from Express
|
|
|
private port: number;
|
|
private port: number;
|
|
|
private host: string;
|
|
private host: string;
|
|
|
|
|
+ private transports: Map<string, SSEServerTransport> = new Map();
|
|
|
|
|
|
|
|
constructor(options: HttpServerOptions) {
|
|
constructor(options: HttpServerOptions) {
|
|
|
this.server = options.server;
|
|
this.server = options.server;
|
|
@@ -48,37 +50,111 @@ export class HttpMcpServer {
|
|
|
this.app.get('/sse', async (req, res) => {
|
|
this.app.get('/sse', async (req, res) => {
|
|
|
console.error('New SSE connection established');
|
|
console.error('New SSE connection established');
|
|
|
|
|
|
|
|
- const transport = new SSEServerTransport('/message', res);
|
|
|
|
|
- await this.server.connect(transport);
|
|
|
|
|
|
|
+ const transport = new SSEServerTransport('/messages', res);
|
|
|
|
|
+
|
|
|
|
|
+ // Store transport by session ID
|
|
|
|
|
+ this.transports.set(transport.sessionId, transport);
|
|
|
|
|
+ console.error(`Transport stored with session ID: ${transport.sessionId}`);
|
|
|
|
|
|
|
|
- // Keep connection alive
|
|
|
|
|
- req.on('close', () => {
|
|
|
|
|
- console.error('SSE connection closed');
|
|
|
|
|
|
|
+ // Clean up when connection closes
|
|
|
|
|
+ res.on('close', () => {
|
|
|
|
|
+ console.error(`SSE connection closed for session: ${transport.sessionId}`);
|
|
|
|
|
+ this.transports.delete(transport.sessionId);
|
|
|
});
|
|
});
|
|
|
|
|
+
|
|
|
|
|
+ await this.server.connect(transport);
|
|
|
});
|
|
});
|
|
|
|
|
|
|
|
- // Message endpoint for client requests
|
|
|
|
|
- this.app.post('/message', async (req, res) => {
|
|
|
|
|
|
|
+ // Message endpoint for client requests (must match the path given to SSEServerTransport)
|
|
|
|
|
+ this.app.post('/messages', async (req, res) => {
|
|
|
|
|
+ const sessionId = req.query.sessionId as string;
|
|
|
|
|
+
|
|
|
|
|
+ if (!sessionId) {
|
|
|
|
|
+ res.status(400).json({
|
|
|
|
|
+ jsonrpc: '2.0',
|
|
|
|
|
+ error: {
|
|
|
|
|
+ code: -32000,
|
|
|
|
|
+ message: 'Bad Request: No sessionId provided'
|
|
|
|
|
+ },
|
|
|
|
|
+ id: null
|
|
|
|
|
+ });
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ const transport = this.transports.get(sessionId);
|
|
|
|
|
+
|
|
|
|
|
+ if (!transport) {
|
|
|
|
|
+ res.status(400).json({
|
|
|
|
|
+ jsonrpc: '2.0',
|
|
|
|
|
+ error: {
|
|
|
|
|
+ code: -32000,
|
|
|
|
|
+ message: 'Bad Request: No transport found for sessionId'
|
|
|
|
|
+ },
|
|
|
|
|
+ id: null
|
|
|
|
|
+ });
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
try {
|
|
try {
|
|
|
- // The SSE transport will handle the message
|
|
|
|
|
- res.status(200).json({ received: true });
|
|
|
|
|
|
|
+ await transport.handlePostMessage(req, res, req.body);
|
|
|
} catch (error) {
|
|
} catch (error) {
|
|
|
console.error('Error handling message:', error);
|
|
console.error('Error handling message:', error);
|
|
|
- res.status(500).json({
|
|
|
|
|
- error: error instanceof Error ? error.message : 'Unknown error'
|
|
|
|
|
- });
|
|
|
|
|
|
|
+ if (!res.headersSent) {
|
|
|
|
|
+ res.status(500).json({
|
|
|
|
|
+ jsonrpc: '2.0',
|
|
|
|
|
+ error: {
|
|
|
|
|
+ code: -32603,
|
|
|
|
|
+ message: error instanceof Error ? error.message : 'Internal server error'
|
|
|
|
|
+ },
|
|
|
|
|
+ id: null
|
|
|
|
|
+ });
|
|
|
|
|
+ }
|
|
|
}
|
|
}
|
|
|
});
|
|
});
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
public start(): Promise<void> {
|
|
public start(): Promise<void> {
|
|
|
- return new Promise((resolve) => {
|
|
|
|
|
- this.app.listen(this.port, this.host, () => {
|
|
|
|
|
- console.error(`Gogs MCP HTTP Server running on http://${this.host}:${this.port}`);
|
|
|
|
|
- console.error(`SSE endpoint: http://${this.host}:${this.port}/sse`);
|
|
|
|
|
- console.error(`Health check: http://${this.host}:${this.port}/health`);
|
|
|
|
|
- resolve();
|
|
|
|
|
- });
|
|
|
|
|
|
|
+ return new Promise((resolve, reject) => {
|
|
|
|
|
+ try {
|
|
|
|
|
+ this.httpServer = this.app.listen(this.port, this.host, () => {
|
|
|
|
|
+ console.error(`Gogs MCP HTTP Server running on http://${this.host}:${this.port}`);
|
|
|
|
|
+ console.error(`SSE endpoint: http://${this.host}:${this.port}/sse`);
|
|
|
|
|
+ console.error(`Messages endpoint: http://${this.host}:${this.port}/messages`);
|
|
|
|
|
+ console.error(`Health check: http://${this.host}:${this.port}/health`);
|
|
|
|
|
+ resolve();
|
|
|
|
|
+ });
|
|
|
|
|
+
|
|
|
|
|
+ this.httpServer.on('error', (error: Error) => {
|
|
|
|
|
+ console.error('HTTP server error:', error);
|
|
|
|
|
+ reject(error);
|
|
|
|
|
+ });
|
|
|
|
|
+ } catch (error) {
|
|
|
|
|
+ console.error('Failed to start HTTP server:', error);
|
|
|
|
|
+ reject(error);
|
|
|
|
|
+ }
|
|
|
});
|
|
});
|
|
|
}
|
|
}
|
|
|
|
|
+
|
|
|
|
|
+ public async cleanup(): Promise<void> {
|
|
|
|
|
+ console.error('Cleaning up HTTP server transports...');
|
|
|
|
|
+ for (const [sessionId, transport] of this.transports) {
|
|
|
|
|
+ try {
|
|
|
|
|
+ console.error(`Closing transport for session ${sessionId}`);
|
|
|
|
|
+ await transport.close();
|
|
|
|
|
+ } catch (error) {
|
|
|
|
|
+ console.error(`Error closing transport for session ${sessionId}:`, error);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ this.transports.clear();
|
|
|
|
|
+
|
|
|
|
|
+ // Close HTTP server
|
|
|
|
|
+ if (this.httpServer) {
|
|
|
|
|
+ return new Promise((resolve) => {
|
|
|
|
|
+ this.httpServer.close(() => {
|
|
|
|
|
+ console.error('HTTP server closed');
|
|
|
|
|
+ resolve();
|
|
|
|
|
+ });
|
|
|
|
|
+ });
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
}
|
|
}
|