turul-http-mcp-server 0.3.36

HTTP transport layer for Model Context Protocol (MCP) servers
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
//! Simple HTTP MCP Server Tests
//!
//! Basic tests for HTTP transport and SSE functionality that use the actual available APIs.

use serde_json::json;
use std::sync::Arc;

use crate::StreamManager;
use crate::server::{HttpMcpServerBuilder, ServerConfig};
use crate::sse::{SseEvent, SseManager};
use crate::streamable_http::{McpProtocolVersion, StreamableHttpHandler};
use turul_mcp_json_rpc_server::JsonRpcDispatcher;
use turul_mcp_protocol::McpError;
use turul_mcp_session_storage::InMemorySessionStorage;

/// Test basic server configuration
#[cfg(test)]
mod basic_tests {
    use super::*;

    #[tokio::test]
    async fn test_server_config_creation() {
        let config = ServerConfig::default();

        // Config should have default values
        assert!(config.enable_cors);
        assert_eq!(config.mcp_path, "/mcp");
        assert_eq!(config.max_body_size, 1024 * 1024);

        println!("Server config created with defaults");
    }

    #[tokio::test]
    async fn test_server_config_customization() {
        let config = ServerConfig {
            enable_cors: false,
            mcp_path: "/custom".to_string(),
            max_body_size: 2 * 1024 * 1024,
            ..Default::default()
        };

        assert!(!config.enable_cors);
        assert_eq!(config.mcp_path, "/custom");
        assert_eq!(config.max_body_size, 2 * 1024 * 1024);

        println!("Server config customized successfully");
    }

    #[tokio::test]
    async fn test_server_builder_creation() {
        let _builder = HttpMcpServerBuilder::new();

        // Builder should be created successfully
        println!("HTTP MCP server builder created");
    }

    #[tokio::test]
    async fn test_streamable_handler_creation() {
        let config = ServerConfig::default();
        let dispatcher = Arc::new(JsonRpcDispatcher::<McpError>::new());
        let session_storage = Arc::new(InMemorySessionStorage::new());
        let stream_manager = Arc::new(StreamManager::new(session_storage.clone()));
        let middleware_stack = Arc::new(crate::middleware::MiddlewareStack::new());

        let _handler = StreamableHttpHandler::new(
            Arc::new(config),
            dispatcher,
            session_storage,
            stream_manager,
            turul_mcp_protocol::ServerCapabilities::default(),
            middleware_stack,
            None,
        );

        // Handler should be created successfully
        println!("Streamable HTTP handler created");
    }
}

/// Test SSE functionality
#[cfg(test)]
mod sse_tests {
    use super::*;

    #[tokio::test]
    async fn test_sse_manager_creation() {
        let manager = SseManager::new();

        assert_eq!(manager.connection_count().await, 0);
        println!("SSE manager created successfully");
    }

    #[tokio::test]
    async fn test_sse_connection_management() {
        let manager = SseManager::new();

        // Create connections
        let _conn1 = manager.create_connection("conn1".to_string()).await;
        let _conn2 = manager.create_connection("conn2".to_string()).await;

        assert_eq!(manager.connection_count().await, 2);

        // Remove a connection
        manager.remove_connection("conn1").await;
        assert_eq!(manager.connection_count().await, 1);

        println!("SSE connection management tested");
    }

    #[tokio::test]
    async fn test_sse_event_formatting() {
        // Test events with "event: message"
        let message_events = vec![
            SseEvent::Connected,
            SseEvent::Data(json!({"message": "test"})),
            SseEvent::Error("Test error".to_string()),
        ];

        for event in message_events {
            let formatted = event.format();

            // All message events should use "event: message" for MCP Inspector compatibility
            assert!(formatted.contains("event: message"));
            assert!(formatted.contains("data: "));
            assert!(formatted.ends_with("\n\n"));

            println!("Event formatted: {}", formatted.replace('\n', "\\n"));
        }

        // Test keepalive separately - it uses SSE comment syntax (no event line)
        let keepalive = SseEvent::KeepAlive;
        let formatted = keepalive.format();
        assert!(!formatted.contains("event:")); // Keepalives don't have event line
        assert!(formatted.starts_with(":")); // SSE comment syntax
        assert!(formatted.ends_with("\n\n"));
        println!("Keepalive formatted: {}", formatted.replace('\n', "\\n"));
    }

    #[tokio::test]
    async fn test_sse_broadcasting() {
        let manager = SseManager::new();

        // Create some connections
        let _conn1 = manager.create_connection("conn1".to_string()).await;
        let _conn2 = manager.create_connection("conn2".to_string()).await;

        // Test different broadcast methods
        manager.send_data(json!({"test": "data"})).await;
        manager.send_error("Test error".to_string()).await;
        manager.send_keep_alive().await;

        // Direct broadcast
        manager.broadcast(SseEvent::Connected).await;

        println!("SSE broadcasting tested successfully");
    }
}

/// Test MCP protocol version handling
#[cfg(test)]
mod protocol_tests {
    use super::*;

    #[tokio::test]
    async fn test_protocol_version_parsing() {
        let versions = vec![
            ("2024-11-05", Some(McpProtocolVersion::V2024_11_05)),
            ("2025-03-26", Some(McpProtocolVersion::V2025_03_26)),
            ("2025-06-18", Some(McpProtocolVersion::V2025_06_18)),
            ("invalid", None),
        ];

        for (input, expected) in versions {
            let parsed = McpProtocolVersion::parse_version(input);
            assert_eq!(parsed, expected);

            if let Some(version) = parsed {
                assert_eq!(version.as_str(), input);
            }
        }

        println!("Protocol version parsing tested successfully");
    }

    #[tokio::test]
    async fn test_protocol_version_comparison() {
        let v1 = McpProtocolVersion::V2024_11_05;
        let v2 = McpProtocolVersion::V2025_03_26;
        let v3 = McpProtocolVersion::V2025_06_18;

        // Test equality
        assert_eq!(v1, McpProtocolVersion::V2024_11_05);
        assert_eq!(v2, McpProtocolVersion::V2025_03_26);
        assert_eq!(v3, McpProtocolVersion::V2025_06_18);

        // Test inequality
        assert_ne!(v1, v2);
        assert_ne!(v2, v3);
        assert_ne!(v1, v3);

        println!("Protocol version comparison tested successfully");
    }
}

/// Test concurrent operations
#[cfg(test)]
mod concurrency_tests {
    use super::*;

    #[tokio::test]
    async fn test_concurrent_sse_connections() {
        let manager = Arc::new(SseManager::new());

        let num_connections = 20;
        let mut handles = Vec::new();

        // Create connections concurrently
        for i in 0..num_connections {
            let manager_clone = manager.clone();
            let handle = tokio::spawn(async move {
                let connection_id = format!("conn_{}", i);
                manager_clone.create_connection(connection_id).await
            });
            handles.push(handle);
        }

        // Wait for all connections to be created
        let _connections = futures::future::join_all(handles).await;

        assert_eq!(manager.connection_count().await, num_connections);

        // Test concurrent broadcasting
        let num_events = 10;
        let mut broadcast_handles = Vec::new();

        for i in 0..num_events {
            let manager_clone = manager.clone();
            let handle = tokio::spawn(async move {
                let data = json!({"event": i, "message": format!("Event {}", i)});
                manager_clone.send_data(data).await;
            });
            broadcast_handles.push(handle);
        }

        futures::future::join_all(broadcast_handles).await;

        println!("Concurrent SSE operations tested successfully");
    }

    #[tokio::test]
    async fn test_concurrent_handler_creation() {
        let num_handlers = 10;
        let mut handles = Vec::new();

        for i in 0..num_handlers {
            let handle = tokio::spawn(async move {
                let config = ServerConfig {
                    mcp_path: format!("/mcp_{}", i),
                    ..Default::default()
                };

                let dispatcher = Arc::new(JsonRpcDispatcher::<McpError>::new());
                let session_storage = Arc::new(InMemorySessionStorage::new());
                let stream_manager = Arc::new(StreamManager::new(session_storage.clone()));
                let middleware_stack = Arc::new(crate::middleware::MiddlewareStack::new());

                let _handler = StreamableHttpHandler::new(
                    Arc::new(config),
                    dispatcher,
                    session_storage,
                    stream_manager,
                    turul_mcp_protocol::ServerCapabilities::default(),
                    middleware_stack,
                    None,
                );
                format!("Handler {} created", i)
            });
            handles.push(handle);
        }

        let results = futures::future::join_all(handles).await;

        // All handlers should be created successfully
        assert_eq!(results.len(), num_handlers);
        for result in results {
            assert!(result.is_ok());
            println!("{}", result.unwrap());
        }
    }
}

/// Test error handling
#[cfg(test)]
mod error_tests {
    use super::*;

    #[tokio::test]
    async fn test_sse_error_event_handling() {
        let manager = SseManager::new();
        let _conn = manager.create_connection("test".to_string()).await;

        // Test various error scenarios
        manager.send_error("".to_string()).await; // Empty error
        manager
            .send_error("Test error with special chars: 🚀".to_string())
            .await;
        manager
            .send_error("Error with \"quotes\" and \n newlines".to_string())
            .await;

        println!("SSE error handling tested successfully");
    }

    #[tokio::test]
    async fn test_large_data_handling() {
        let manager = SseManager::new();
        let _conn = manager
            .create_connection("large_data_test".to_string())
            .await;

        // Test with large JSON data
        let large_data = json!({
            "large_string": "x".repeat(10_000),
            "large_array": (0..1000).collect::<Vec<i32>>(),
            "nested": {
                "deep": {
                    "structure": "test"
                }
            }
        });

        manager.send_data(large_data).await;

        println!("Large data handling tested successfully");
    }

    #[tokio::test]
    async fn test_invalid_protocol_versions() {
        let invalid_versions = vec!["", "2024", "invalid-version", "2025-99-99", "1.0.0"];

        for version in invalid_versions {
            let parsed = McpProtocolVersion::parse_version(version);
            assert!(parsed.is_none(), "Version '{}' should not parse", version);
        }

        println!("Invalid protocol version handling tested successfully");
    }
}

/// Test performance characteristics
#[cfg(test)]
mod performance_tests {
    use super::*;
    use std::time::Instant;

    #[tokio::test]
    async fn test_sse_performance() {
        let manager = SseManager::new();

        // Create connections
        let num_connections = 50;
        for i in 0..num_connections {
            let _conn = manager.create_connection(format!("perf_conn_{}", i)).await;
        }

        assert_eq!(manager.connection_count().await, num_connections);

        // Test broadcast performance
        let num_events = 100;
        let start = Instant::now();

        for i in 0..num_events {
            let data = json!({"event_id": i, "timestamp": start.elapsed().as_millis()});
            manager.send_data(data).await;
        }

        let duration = start.elapsed();
        println!(
            "Sent {} events to {} connections in {:?}",
            num_events, num_connections, duration
        );

        // Should complete reasonably quickly
        assert!(duration.as_millis() < 1000);
    }

    #[tokio::test]
    async fn test_handler_creation_performance() {
        let num_handlers = 100;
        let start = Instant::now();

        let mut handlers = Vec::new();
        for i in 0..num_handlers {
            let config = ServerConfig {
                mcp_path: format!("/perf_{}", i),
                ..Default::default()
            };

            let dispatcher = Arc::new(JsonRpcDispatcher::<McpError>::new());
            let session_storage = Arc::new(InMemorySessionStorage::new());
            let stream_manager = Arc::new(StreamManager::new(session_storage.clone()));
            let middleware_stack = Arc::new(crate::middleware::MiddlewareStack::new());

            let handler = StreamableHttpHandler::new(
                Arc::new(config),
                dispatcher,
                session_storage,
                stream_manager,
                turul_mcp_protocol::ServerCapabilities::default(),
                middleware_stack,
                None,
            );
            handlers.push(handler);
        }

        let duration = start.elapsed();
        println!("Created {} handlers in {:?}", num_handlers, duration);

        assert_eq!(handlers.len(), num_handlers);
        assert!(duration.as_millis() < 100);
    }
}