fix: Resolve SSE reading deadlocks in Stdio proxies
This commit is contained in:
1 parent
0ecd862e89
commit
0de78e8fae
2 files changed
+62
-38
No files matched your search
+31
-19
@@ -69,27 +69,39 @@ pub fn run_proxy(target_url: &str) -> Result<bool, Box<dyn std::error::Error + S
|
|||||||
let mut is_message = false;
|
let mut is_message = false;
|
||||||
let mut is_endpoint = false;
|
let mut is_endpoint = false;
|
||||||
|
|
||||||
while let Ok(bytes) = reader.read_line(&mut line).await {
|
loop {
|
||||||
if bytes == 0 { break; }
|
tokio::select! {
|
||||||
let trimmed = line.trim();
|
_ = shutdown_rx.recv() => {
|
||||||
if trimmed.starts_with("event: message") {
|
return Ok(false);
|
||||||
is_message = true;
|
}
|
||||||
is_endpoint = false;
|
res = reader.read_line(&mut line) => {
|
||||||
} else if trimmed.starts_with("event: endpoint") {
|
match res {
|
||||||
is_endpoint = true;
|
Ok(bytes) => {
|
||||||
is_message = false;
|
if bytes == 0 { break; }
|
||||||
} else if trimmed.starts_with("data: ") {
|
let trimmed = line.trim();
|
||||||
if is_message {
|
if trimmed.starts_with("event: message") {
|
||||||
println!("{}", &trimmed[6..]);
|
is_message = true;
|
||||||
is_message = false;
|
is_endpoint = false;
|
||||||
} else if is_endpoint {
|
} else if trimmed.starts_with("event: endpoint") {
|
||||||
let ep = &trimmed[6..];
|
is_endpoint = true;
|
||||||
let mut p = post_url.write().await;
|
is_message = false;
|
||||||
*p = format!("{}{}", target_url, ep);
|
} else if trimmed.starts_with("data: ") {
|
||||||
is_endpoint = false;
|
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();
|
*post_url.write().await = String::new();
|
||||||
return Ok(true);
|
return Ok(true);
|
||||||
|
|||||||
+31
-19
@@ -87,27 +87,39 @@ fn main() -> Result<(), Box<dyn std::error::Error>> {
|
|||||||
let mut is_message = false;
|
let mut is_message = false;
|
||||||
let mut is_endpoint = false;
|
let mut is_endpoint = false;
|
||||||
|
|
||||||
while let Ok(bytes) = reader.read_line(&mut line).await {
|
loop {
|
||||||
if bytes == 0 { break; }
|
tokio::select! {
|
||||||
let trimmed = line.trim();
|
_ = shutdown_rx.recv() => {
|
||||||
if trimmed.starts_with("event: message") {
|
return Ok(()); // Stdin closed, exit entirely
|
||||||
is_message = true;
|
}
|
||||||
is_endpoint = false;
|
res = reader.read_line(&mut line) => {
|
||||||
} else if trimmed.starts_with("event: endpoint") {
|
match res {
|
||||||
is_endpoint = true;
|
Ok(bytes) => {
|
||||||
is_message = false;
|
if bytes == 0 { break; }
|
||||||
} else if trimmed.starts_with("data: ") {
|
let trimmed = line.trim();
|
||||||
if is_message {
|
if trimmed.starts_with("event: message") {
|
||||||
println!("{}", &trimmed[6..]);
|
is_message = true;
|
||||||
is_message = false;
|
is_endpoint = false;
|
||||||
} else if is_endpoint {
|
} else if trimmed.starts_with("event: endpoint") {
|
||||||
let ep = &trimmed[6..];
|
is_endpoint = true;
|
||||||
let mut p = post_url.write().await;
|
is_message = false;
|
||||||
*p = format!("{}{}", target_url, ep);
|
} else if trimmed.starts_with("data: ") {
|
||||||
is_endpoint = false;
|
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();
|
*post_url.write().await = String::new();
|
||||||
tokio::time::sleep(tokio::time::Duration::from_millis(500)).await;
|
tokio::time::sleep(tokio::time::Duration::from_millis(500)).await;
|
||||||
|
|||||||
Reference in new issue
Block a user