feat: implement dual-transport WAL architecture with permanent background leader and lightweight stubs
This commit is contained in:
1 parent
a2febd1b00
commit
e4ff476b6d
16 files changed
+1910
-52
No files matched your search
@@ -0,0 +1,144 @@
|
||||
use clap::Parser;
|
||||
use futures_util::StreamExt;
|
||||
use std::sync::Arc;
|
||||
use tokio::io::AsyncBufReadExt;
|
||||
use tokio::sync::RwLock;
|
||||
use tokio_util::io::StreamReader;
|
||||
|
||||
#[derive(Parser)]
|
||||
#[command(name = "mcp-memory-stub")]
|
||||
struct Cli {
|
||||
#[arg(long, default_value = "http://localhost:3000")]
|
||||
target: String,
|
||||
#[arg(long)]
|
||||
wake_cmd: Option<String>,
|
||||
}
|
||||
|
||||
fn run_proxy(target_url: &str) -> Result<bool, Box<dyn std::error::Error>> {
|
||||
let rt = tokio::runtime::Runtime::new()?;
|
||||
rt.block_on(async {
|
||||
let client = reqwest::Client::builder().build()?;
|
||||
let sse_url = format!("{}/sse", target_url);
|
||||
|
||||
eprintln!("[PROXY] Connecting to SSE: {}", sse_url);
|
||||
let resp = match client.get(&sse_url).send().await {
|
||||
Ok(r) => r,
|
||||
Err(e) => {
|
||||
eprintln!("[PROXY] SSE connection failed: {}", e);
|
||||
return Ok(true);
|
||||
}
|
||||
};
|
||||
|
||||
let post_url = Arc::new(RwLock::new(format!("{}/messages", target_url)));
|
||||
let post_url_clone = Arc::clone(&post_url);
|
||||
|
||||
let (tx, mut rx) = tokio::sync::mpsc::channel(1);
|
||||
|
||||
let tx_clone = tx.clone();
|
||||
let client_clone = client.clone();
|
||||
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 body = buffer.clone();
|
||||
buffer.clear();
|
||||
eprintln!("[PROXY] Stdin read: {}", body.trim());
|
||||
let client = client_clone.clone();
|
||||
let url_arc = Arc::clone(&post_url_clone);
|
||||
|
||||
tokio::spawn(async move {
|
||||
let mut url = "".to_string();
|
||||
for _ in 0..50 {
|
||||
let u = url_arc.read().await.clone();
|
||||
if u.contains("sessionId") {
|
||||
url = u;
|
||||
break;
|
||||
}
|
||||
tokio::time::sleep(tokio::time::Duration::from_millis(100)).await;
|
||||
}
|
||||
if url.is_empty() {
|
||||
url = url_arc.read().await.clone();
|
||||
}
|
||||
eprintln!("[PROXY] POSTing to {}", url);
|
||||
let res = client.post(&url).header("Content-Type", "application/json").body(body).send().await;
|
||||
eprintln!("[PROXY] POST response: {:?}", res.map(|r| r.status()));
|
||||
});
|
||||
}
|
||||
eprintln!("[PROXY] Stdin closed");
|
||||
let _ = tx_clone.blocking_send(false);
|
||||
});
|
||||
|
||||
let target_url = target_url.to_string();
|
||||
let tx_clone2 = tx.clone();
|
||||
tokio::spawn(async move {
|
||||
let stream = resp.bytes_stream().map(|res| res.map_err(|e| std::io::Error::new(std::io::ErrorKind::Other, e)));
|
||||
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;
|
||||
while let Ok(bytes) = reader.read_line(&mut line).await {
|
||||
if bytes == 0 { break; }
|
||||
let trimmed = line.trim();
|
||||
eprintln!("[PROXY] SSE read: {}", trimmed);
|
||||
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);
|
||||
eprintln!("[PROXY] Endpoint updated: {}", *p);
|
||||
is_endpoint = false;
|
||||
}
|
||||
}
|
||||
line.clear();
|
||||
}
|
||||
eprintln!("[PROXY] SSE stream closed");
|
||||
let _ = tx_clone2.send(true).await;
|
||||
});
|
||||
|
||||
let dropped = rx.recv().await.unwrap_or(true);
|
||||
Ok(dropped)
|
||||
})
|
||||
}
|
||||
|
||||
fn main() {
|
||||
let cli = Cli::parse();
|
||||
loop {
|
||||
if let Err(_) = std::net::TcpStream::connect(cli.target.replace("http://", "").replace("https://", "")) {
|
||||
eprintln!("[PROXY] Target {} offline, sleeping", cli.target);
|
||||
|
||||
if let Some(ref cmd) = cli.wake_cmd {
|
||||
eprintln!("[PROXY] Executing wake command...");
|
||||
let parts: Vec<&str> = cmd.split_whitespace().collect();
|
||||
if !parts.is_empty() {
|
||||
let _ = std::process::Command::new(parts[0])
|
||||
.args(&parts[1..])
|
||||
.spawn();
|
||||
}
|
||||
}
|
||||
|
||||
std::thread::sleep(std::time::Duration::from_secs(1));
|
||||
}
|
||||
match run_proxy(&cli.target) {
|
||||
Ok(true) => {
|
||||
std::thread::sleep(std::time::Duration::from_millis(50));
|
||||
}
|
||||
Ok(false) => {
|
||||
break;
|
||||
}
|
||||
Err(_) => {
|
||||
std::thread::sleep(std::time::Duration::from_millis(1000));
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user