use clap::Parser; use futures_util::StreamExt; use std::sync::Arc; use tokio::io::AsyncBufReadExt; use tokio::sync::RwLock; use tokio::sync::mpsc; use tokio_util::io::StreamReader; #[derive(Parser)] #[command(name = "mcp-memory-stub", author, version, about = "Antigravity MCP Memory Stub / Proxy", long_about = None)] struct Cli { /// Target URL for the stub to proxy messages to #[arg(long, default_value = "http://localhost:3000")] target: String, /// Optional command to execute if the target server is unreachable #[arg(long)] wake_cmd: Option, } fn main() -> Result<(), Box> { let cli = Cli::parse(); let rt = tokio::runtime::Runtime::new()?; rt.block_on(async { let (msg_tx, mut msg_rx) = mpsc::channel::(100); let (shutdown_tx, mut shutdown_rx) = mpsc::channel::<()>(1); tokio::task::spawn_blocking(move || { let stdin = std::io::stdin(); let mut handle = stdin.lock(); let mut buffer = String::new(); while let Ok(bytes) = std::io::BufRead::read_line(&mut handle, &mut buffer) { if bytes == 0 { break; } let _ = msg_tx.blocking_send(buffer.clone()); buffer.clear(); } let _ = shutdown_tx.blocking_send(()); }); let target_url = cli.target; let post_url = Arc::new(RwLock::new(String::new())); let post_url_clone = Arc::clone(&post_url); let client = reqwest::Client::builder().build()?; tokio::spawn(async move { while let Some(msg) = msg_rx.recv().await { let mut attempts = 0; loop { let url = post_url_clone.read().await.clone(); if !url.is_empty() { let res = client .post(&url) .header("Accept", "application/json, text/event-stream") .header("Content-Type", "application/json") .body(msg.clone()) .send() .await; if let Ok(resp) = res { if resp.status().is_success() { break; } } } tokio::time::sleep(tokio::time::Duration::from_millis(10)).await; attempts += 1; if attempts % 10 == 0 { eprintln!("[PROXY] Waiting for server to accept messages..."); } } } }); let wake_cmd = cli.wake_cmd; loop { if shutdown_rx.try_recv().is_ok() { break; } let sse_url = format!("{}/sse", target_url); let client = reqwest::Client::builder().build()?; match client .get(&sse_url) .header("Accept", "text/event-stream") .send() .await { Ok(resp) => { if resp.status() == reqwest::StatusCode::GONE { eprintln!("[PROXY] Target gone, exiting."); break; } let stream = resp.bytes_stream().map(|res| { res.map_err(std::io::Error::other) }); let mut reader = tokio::io::BufReader::new(StreamReader::new(stream)); let mut line = String::new(); let mut is_message = false; let mut is_endpoint = false; loop { tokio::select! { _ = shutdown_rx.recv() => { return Ok(()); // Stdin closed, exit entirely } res = reader.read_line(&mut line) => { match res { Ok(bytes) => { if bytes == 0 { break; } let trimmed = line.trim(); if trimmed.starts_with("event: message") { is_message = true; is_endpoint = false; } else if trimmed.starts_with("event: endpoint") { is_endpoint = true; is_message = false; } else if let Some(stripped) = trimmed.strip_prefix("data: ") { if is_message { println!("{}", stripped); is_message = false; } else if is_endpoint { let mut p = post_url.write().await; *p = format!("{}{}", target_url, stripped); is_endpoint = false; } } line.clear(); } Err(_) => break, } } } } *post_url.write().await = String::new(); tokio::time::sleep(tokio::time::Duration::from_millis(10)).await; } Err(_) => { if let Some(ref cmd) = wake_cmd { let parts: Vec<&str> = cmd.split_whitespace().collect(); if !parts.is_empty() { let _ = std::process::Command::new(parts[0]) .args(&parts[1..]) .spawn(); } } tokio::time::sleep(tokio::time::Duration::from_secs(1)).await; } } } Ok(()) }) }