perf(server): Resolve tokio executor starvation and optimize broadcast loops
- Offloaded the blocking idx.commit() inside index_committer_worker to okio::task::spawn_blocking to prevent the tokio async thread from being starved by blocking I/O. - Removed expensive full HashMap .clone() operations during UI client websocket broadcasts in both handle_socket and vim_telemetry_handler, filtering and collecting only the Sender handles. - Migrated blocking std::fs::write calls across the WSL 9P boundary in vim_telemetry_handler over to asynchronous okio::fs::write().await to prevent severe async thread blockages on native file I/O delays.
This commit is contained in:
1 parent
b9e44e1fa3
commit
b98f95d9d8
2 files changed
+218
-18
No files matched your search
+27
-13
@@ -88,11 +88,14 @@ enum GateCommands {
|
||||
|
||||
async fn index_committer_worker(state: Arc<MemoryState>) {
|
||||
loop {
|
||||
sleep(Duration::from_secs(5)).await;
|
||||
tokio::time::sleep(Duration::from_secs(5)).await;
|
||||
// Periodically commit the search index to persist inline indexing operations
|
||||
if let Ok(idx) = state.search_index.read() {
|
||||
let _ = idx.commit();
|
||||
}
|
||||
let state_clone = Arc::clone(&state);
|
||||
let _ = tokio::task::spawn_blocking(move || {
|
||||
if let Ok(idx) = state_clone.search_index.read() {
|
||||
let _ = idx.commit();
|
||||
}
|
||||
}).await;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -500,11 +503,22 @@ async fn handle_socket(socket: WebSocket, state: Arc<AppState>, client_type: Str
|
||||
"data": activity_msg
|
||||
});
|
||||
|
||||
let clients_map = state_clone.clients.read().unwrap().clone();
|
||||
for (id, client_tx) in clients_map.iter() {
|
||||
if id != &session_id_clone {
|
||||
let _ = client_tx.send(event.to_string()).await;
|
||||
}
|
||||
let senders: Vec<_> = state_clone
|
||||
.clients
|
||||
.read()
|
||||
.unwrap()
|
||||
.iter()
|
||||
.filter_map(|(id, tx)| {
|
||||
if id != &session_id_clone {
|
||||
Some(tx.clone())
|
||||
} else {
|
||||
None
|
||||
}
|
||||
})
|
||||
.collect();
|
||||
|
||||
for client_tx in senders {
|
||||
let _ = client_tx.send(event.to_string()).await;
|
||||
}
|
||||
}
|
||||
} // End if proxy
|
||||
@@ -617,10 +631,10 @@ async fn nvim_telemetry_handler(
|
||||
let profile =
|
||||
std::env::var("USERPROFILE").unwrap_or_else(|_| "C:\\Users\\reazul.ashraf".into());
|
||||
let win_path = format!("{}\\.gemini\\active_nvim.txt", profile);
|
||||
let _ = std::fs::write(&win_path, &payload.session_id);
|
||||
let _ = tokio::fs::write(&win_path, &payload.session_id).await;
|
||||
|
||||
let wsl_path = "\\\\wsl.localhost\\Ubuntu\\home\\riz\\.gemini\\active_nvim.txt";
|
||||
let _ = std::fs::write(wsl_path, &payload.session_id);
|
||||
let _ = tokio::fs::write(wsl_path, &payload.session_id).await;
|
||||
}
|
||||
|
||||
// 2. Broadcast to UI WebSockets
|
||||
@@ -630,8 +644,8 @@ async fn nvim_telemetry_handler(
|
||||
});
|
||||
|
||||
let msg_str = ws_msg.to_string();
|
||||
let clients = state.clients.read().unwrap().clone();
|
||||
for tx in clients.values() {
|
||||
let senders: Vec<_> = state.clients.read().unwrap().values().cloned().collect();
|
||||
for tx in senders {
|
||||
let _ = tx.send(msg_str.clone()).await;
|
||||
}
|
||||
|
||||
|
||||
Reference in new issue
Block a user