Start the background SSE listener. Must be called before `send`/`receive`.
(&self)
| 296 | |
| 297 | /// Start the background SSE listener. Must be called before `send`/`receive`. |
| 298 | pub async fn connect(&self) -> Result<(), McpClientError> { |
| 299 | use futures::StreamExt; |
| 300 | use reqwest_eventsource::{Event, EventSource}; |
| 301 | |
| 302 | let mut builder = self.client.get(&self.url); |
| 303 | builder = builder.header("Accept", "text/event-stream"); |
| 304 | for (key, value) in &self.headers { |
| 305 | builder = builder.header(key.as_str(), value.as_str()); |
| 306 | } |
| 307 | |
| 308 | let mut es = EventSource::new(builder).map_err(|e| { |
| 309 | McpClientError::TransportError(format!("Failed to create SSE connection: {}", e)) |
| 310 | })?; |
| 311 | |
| 312 | let tx = self.response_tx.clone(); |
| 313 | |
| 314 | let handle = tokio::spawn(async move { |
| 315 | while let Some(event) = es.next().await { |
| 316 | match event { |
| 317 | Ok(Event::Message(msg)) => { |
| 318 | let data = msg.data.trim().to_string(); |
| 319 | if data.is_empty() || data == "[DONE]" { |
| 320 | continue; |
| 321 | } |
| 322 | match JsonRpcMessage::from_str(&data) { |
| 323 | Ok(msg) => { |
| 324 | if tx.send(msg).is_err() { |
| 325 | break; |
| 326 | } |
| 327 | } |
| 328 | Err(e) => { |
| 329 | tracing::warn!("SSE: failed to parse message: {}", e); |
| 330 | } |
| 331 | } |
| 332 | } |
| 333 | Ok(Event::Open) => { |
| 334 | tracing::debug!("SSE connection opened"); |
| 335 | } |
| 336 | Err(e) => { |
| 337 | tracing::error!("SSE error: {}", e); |
| 338 | break; |
| 339 | } |
| 340 | } |
| 341 | } |
| 342 | }); |
| 343 | |
| 344 | let mut task = self.sse_task.lock().await; |
| 345 | *task = Some(handle); |
| 346 | |
| 347 | Ok(()) |
| 348 | } |
| 349 | } |
| 350 | |
| 351 | #[async_trait] |