use futures_util::StreamExt; use reqwest::Client; use std::env; use std::time::Duration; #[tokio::main] async fn main() -> Result<(), Box> { tracing_subscriber::fmt::init(); let target = env::var("MCP_TARGET").unwrap_or_else(|_| "http://127.0.0.1:3000".to_string()); let token = env::var("MCP_AUTH_TOKEN") .unwrap_or_else(|_| "jP76lUJ5DtFRZmcvXH8LKdCTIkp29eAf".to_string()); tracing::info!("Starting skeletal client to {}", target); let client = Client::new(); let sse_url = format!("{target}/sse"); tracing::info!("Connecting to SSE: {}", sse_url); let res = client.get(&sse_url).bearer_auth(&token).send().await?; if !res.status().is_success() { tracing::error!("Failed to connect to SSE: {}", res.status()); return Err("SSE connection failed".into()); } tracing::info!("SSE Connected. Reading stream..."); let mut stream = res.bytes_stream(); let mut buffer = Vec::new(); let mut post_endpoint = None; // Read the initial event containing the POST endpoint while let Some(chunk) = stream.next().await { let bytes = chunk?; buffer.extend_from_slice(&bytes); while let Some(pos) = buffer.windows(2).position(|w| w == b"\n\n" || w == b"\r\n") { let msg_bytes = buffer.drain(..pos).collect::>(); buffer.drain(..2); let text = String::from_utf8_lossy(&msg_bytes); let mut is_endpoint = false; let mut data_content = String::new(); for line in text.lines() { if line.starts_with("event: endpoint") { is_endpoint = true; } else if let Some(data) = line.strip_prefix("data: ") { data_content.push_str(data); } } if is_endpoint && !data_content.is_empty() { tracing::info!("Received POST endpoint: {}", data_content); post_endpoint = Some(data_content); break; } tracing::info!("Received early SSE data: {}", text); } if post_endpoint.is_some() { break; } } let post_endpoint = post_endpoint.ok_or("Did not receive endpoint from SSE stream")?; let post_url = format!("{target}{post_endpoint}"); let payload = r#"{"jsonrpc":"2.0","id":999,"method":"server/discover","params":{}}"#; tracing::info!("Sending test payload to {}", post_url); tracing::info!("Payload: {}", payload); let post_res = client .post(&post_url) .bearer_auth(&token) .header("Content-Type", "application/json") .body(payload.to_string()) .send() .await?; tracing::info!("POST Response Status: {}", post_res.status()); let post_body = post_res.text().await?; tracing::info!("POST Response Body: {}", post_body); // Wait for the SSE stream to deliver the response tracing::info!("Waiting 2 seconds for SSE response delivery..."); let mut timeout = tokio::time::interval(Duration::from_secs(2)); timeout.tick().await; // first tick is immediate tokio::select! { _ = timeout.tick() => { tracing::warn!("Timed out waiting for SSE response."); } () = async { while let Some(chunk) = stream.next().await { if let Ok(bytes) = chunk { tracing::info!("Received SSE Chunk: {}", String::from_utf8_lossy(&bytes)); break; } } } => { tracing::info!("Successfully read SSE response from stream."); } } tracing::info!("Skeletal client test complete."); Ok(()) }