fix: memory leaks and unoptimized websocket broadcasts
This commit is contained in:
1 parent
9b349e6459
commit
573c9586fd
2 files changed
+22
-17
No files matched your search
+2
-2
@@ -514,7 +514,7 @@ async fn handle_socket(socket: WebSocket, state: Arc<AppState>, client_type: Str
|
|||||||
.collect();
|
.collect();
|
||||||
|
|
||||||
for client_tx in senders {
|
for client_tx in senders {
|
||||||
let _ = client_tx.send(event.to_string()).await;
|
let _ = client_tx.try_send(event.to_string());
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
} // End if proxy
|
} // End if proxy
|
||||||
@@ -642,7 +642,7 @@ async fn nvim_telemetry_handler(
|
|||||||
let msg_str = ws_msg.to_string();
|
let msg_str = ws_msg.to_string();
|
||||||
let senders: Vec<_> = state.clients.read().unwrap().values().cloned().collect();
|
let senders: Vec<_> = state.clients.read().unwrap().values().cloned().collect();
|
||||||
for tx in senders {
|
for tx in senders {
|
||||||
let _ = tx.send(msg_str.clone()).await;
|
let _ = tx.try_send(msg_str.clone());
|
||||||
}
|
}
|
||||||
|
|
||||||
axum::Json(serde_json::json!({"status": "ok"}))
|
axum::Json(serde_json::json!({"status": "ok"}))
|
||||||
|
|||||||
+20
-15
@@ -1,8 +1,17 @@
|
|||||||
use serde_json::{Value, json};
|
use serde_json::{Value, json};
|
||||||
use std::io::{BufRead, BufReader, Write};
|
use std::io::{BufRead, BufReader, Write};
|
||||||
use std::process::{Command, Stdio};
|
use std::process::{Child, Command, Stdio};
|
||||||
use std::time::Duration;
|
use std::time::Duration;
|
||||||
|
|
||||||
|
struct ChildGuard(Child);
|
||||||
|
|
||||||
|
impl Drop for ChildGuard {
|
||||||
|
fn drop(&mut self) {
|
||||||
|
let _ = self.0.kill();
|
||||||
|
let _ = self.0.wait();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
fn send_message(stdin: &mut std::process::ChildStdin, msg: Value) {
|
fn send_message(stdin: &mut std::process::ChildStdin, msg: Value) {
|
||||||
let s = serde_json::to_string(&msg).unwrap();
|
let s = serde_json::to_string(&msg).unwrap();
|
||||||
let payload = format!("{}\n", s);
|
let payload = format!("{}\n", s);
|
||||||
@@ -52,7 +61,7 @@ async fn test_full_system_e2e_performance() {
|
|||||||
assert!(stub_exe.exists(), "Stub not found at {:?}", stub_exe);
|
assert!(stub_exe.exists(), "Stub not found at {:?}", stub_exe);
|
||||||
|
|
||||||
// 1. Start Server
|
// 1. Start Server
|
||||||
let mut server = Command::new(&server_exe)
|
let mut server = ChildGuard(Command::new(&server_exe)
|
||||||
.arg("--daemon")
|
.arg("--daemon")
|
||||||
.env("MCP_PORT", test_port)
|
.env("MCP_PORT", test_port)
|
||||||
.env("RUST_LOG", "debug")
|
.env("RUST_LOG", "debug")
|
||||||
@@ -62,7 +71,7 @@ async fn test_full_system_e2e_performance() {
|
|||||||
.stdout(Stdio::inherit())
|
.stdout(Stdio::inherit())
|
||||||
.stderr(Stdio::inherit())
|
.stderr(Stdio::inherit())
|
||||||
.spawn()
|
.spawn()
|
||||||
.expect("Failed to start server");
|
.expect("Failed to start server"));
|
||||||
|
|
||||||
// Give server time to generate TLS cert and start
|
// Give server time to generate TLS cert and start
|
||||||
let client = reqwest::Client::builder()
|
let client = reqwest::Client::builder()
|
||||||
@@ -85,7 +94,7 @@ async fn test_full_system_e2e_performance() {
|
|||||||
assert!(started, "Server failed to start in time");
|
assert!(started, "Server failed to start in time");
|
||||||
|
|
||||||
// 2. Start Stub
|
// 2. Start Stub
|
||||||
let mut stub = Command::new(&stub_exe)
|
let mut stub = ChildGuard(Command::new(&stub_exe)
|
||||||
.arg("--target")
|
.arg("--target")
|
||||||
.arg(format!("http://127.0.0.1:{}", test_port))
|
.arg(format!("http://127.0.0.1:{}", test_port))
|
||||||
.env("MCP_MEMORY_STORE_DIR", temp_dir.to_str().unwrap())
|
.env("MCP_MEMORY_STORE_DIR", temp_dir.to_str().unwrap())
|
||||||
@@ -95,21 +104,21 @@ async fn test_full_system_e2e_performance() {
|
|||||||
.stdout(Stdio::piped())
|
.stdout(Stdio::piped())
|
||||||
.stderr(Stdio::inherit())
|
.stderr(Stdio::inherit())
|
||||||
.spawn()
|
.spawn()
|
||||||
.expect("Failed to start stub");
|
.expect("Failed to start stub"));
|
||||||
|
|
||||||
let mut stub_stdin = stub.stdin.take().unwrap();
|
let mut stub_stdin = stub.0.stdin.take().unwrap();
|
||||||
let mut stub_stdout = BufReader::new(stub.stdout.take().unwrap());
|
let mut stub_stdout = BufReader::new(stub.0.stdout.take().unwrap());
|
||||||
|
|
||||||
// 3. Start Nvim Bridge
|
// 3. Start Nvim Bridge
|
||||||
let mut nvim = Command::new(&nvim_exe)
|
let mut nvim = ChildGuard(Command::new(&nvim_exe)
|
||||||
.stdin(Stdio::piped())
|
.stdin(Stdio::piped())
|
||||||
.stdout(Stdio::piped())
|
.stdout(Stdio::piped())
|
||||||
.stderr(Stdio::inherit())
|
.stderr(Stdio::inherit())
|
||||||
.spawn()
|
.spawn()
|
||||||
.expect("Failed to start nvim bridge");
|
.expect("Failed to start nvim bridge"));
|
||||||
|
|
||||||
let mut nvim_stdin = nvim.stdin.take().unwrap();
|
let mut nvim_stdin = nvim.0.stdin.take().unwrap();
|
||||||
let mut nvim_stdout = BufReader::new(nvim.stdout.take().unwrap());
|
let mut nvim_stdout = BufReader::new(nvim.0.stdout.take().unwrap());
|
||||||
|
|
||||||
println!("Server, stub, and nvim spawned successfully");
|
println!("Server, stub, and nvim spawned successfully");
|
||||||
|
|
||||||
@@ -175,9 +184,5 @@ async fn test_full_system_e2e_performance() {
|
|||||||
println!("Stub 100 requests: {:?}", stub_duration);
|
println!("Stub 100 requests: {:?}", stub_duration);
|
||||||
println!("Win-Nvim 100 requests: {:?}", nvim_duration);
|
println!("Win-Nvim 100 requests: {:?}", nvim_duration);
|
||||||
|
|
||||||
// Cleanup
|
|
||||||
let _ = server.kill();
|
|
||||||
let _ = stub.kill();
|
|
||||||
let _ = nvim.kill();
|
|
||||||
let _ = std::fs::remove_dir_all(temp_dir);
|
let _ = std::fs::remove_dir_all(temp_dir);
|
||||||
}
|
}
|
||||||
Reference in new issue
Block a user