fix(proxy): fully non-blocking mcp-memory-stub proxy
This makes the stub purely non-blocking via tokio::spawn and fixes cross-OS compilation boundaries in build.cmd
This commit is contained in:
1 parent
67b6a0407e
commit
1fd1d119e6
6 files changed
+617
-681
No files matched your search
+28
-25
@@ -24,48 +24,51 @@ fn main() -> Result<(), Box<dyn std::error::Error>> {
|
||||
let (msg_tx, mut msg_rx) = mpsc::channel::<String>(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();
|
||||
tokio::spawn(async move {
|
||||
let mut stdin = tokio::io::BufReader::new(tokio::io::stdin());
|
||||
let mut buffer = String::new();
|
||||
while let Ok(bytes) = std::io::BufRead::read_line(&mut handle, &mut buffer) {
|
||||
while let Ok(bytes) = stdin.read_line(&mut buffer).await {
|
||||
if bytes == 0 {
|
||||
break;
|
||||
}
|
||||
let _ = msg_tx.blocking_send(buffer.clone());
|
||||
let _ = msg_tx.send(buffer.clone()).await;
|
||||
buffer.clear();
|
||||
}
|
||||
let _ = shutdown_tx.blocking_send(());
|
||||
let _ = shutdown_tx.send(()).await;
|
||||
});
|
||||
|
||||
let target_url = cli.target;
|
||||
let post_url = Arc::new(RwLock::new(String::new()));
|
||||
let post_url_clone = Arc::clone(&post_url);
|
||||
let post_url_proxy = 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;
|
||||
let post_url_clone = Arc::clone(&post_url_proxy);
|
||||
let client = client.clone();
|
||||
tokio::spawn(async move {
|
||||
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; } }
|
||||
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...");
|
||||
}
|
||||
}
|
||||
tokio::time::sleep(tokio::time::Duration::from_millis(10)).await;
|
||||
attempts += 1;
|
||||
if attempts % 10 == 0 {
|
||||
eprintln!("[PROXY] Waiting for server to accept messages...");
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
});
|
||||
|
||||
|
||||
Reference in new issue
Block a user