From 0de78e8faecd65ad2d03ca837f526bd0fbc94ab1 Mon Sep 17 00:00:00 2001 From: Riz Ashraf Date: Wed, 9 Sep 2026 08:27:56 +0100 Subject: [PATCH] fix: Resolve SSE reading deadlocks in Stdio proxies --- server/src/proxy.rs | 50 ++++++++++++++++++++++++++++----------------- stub/src/main.rs | 50 ++++++++++++++++++++++++++++----------------- 2 files changed, 62 insertions(+), 38 deletions(-) diff --git a/server/src/proxy.rs b/server/src/proxy.rs index 393df47..237bafa 100644 --- a/server/src/proxy.rs +++ b/server/src/proxy.rs @@ -69,27 +69,39 @@ pub fn run_proxy(target_url: &str) -> Result { + return Ok(false); + } + 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 trimmed.starts_with("data: ") { + if is_message { + println!("{}", &trimmed[6..]); + is_message = false; + } else if is_endpoint { + let ep = &trimmed[6..]; + let mut p = post_url.write().await; + *p = format!("{}{}", target_url, ep); + is_endpoint = false; + } + } + line.clear(); + } + Err(_) => break, + } } } - line.clear(); } *post_url.write().await = String::new(); return Ok(true); diff --git a/stub/src/main.rs b/stub/src/main.rs index acd6f34..ff7a451 100644 --- a/stub/src/main.rs +++ b/stub/src/main.rs @@ -87,27 +87,39 @@ fn main() -> Result<(), Box> { let mut is_message = false; let mut is_endpoint = false; - while let Ok(bytes) = reader.read_line(&mut line).await { - 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 trimmed.starts_with("data: ") { - if is_message { - println!("{}", &trimmed[6..]); - is_message = false; - } else if is_endpoint { - let ep = &trimmed[6..]; - let mut p = post_url.write().await; - *p = format!("{}{}", target_url, ep); - 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 trimmed.starts_with("data: ") { + if is_message { + println!("{}", &trimmed[6..]); + is_message = false; + } else if is_endpoint { + let ep = &trimmed[6..]; + let mut p = post_url.write().await; + *p = format!("{}{}", target_url, ep); + is_endpoint = false; + } + } + line.clear(); + } + Err(_) => break, + } } } - line.clear(); } *post_url.write().await = String::new(); tokio::time::sleep(tokio::time::Duration::from_millis(500)).await;