Skip to content

Commit 78e2e13

Browse files
committed
fix(memtrack): collect events in thread to avoid mutex overhead
1 parent d244c5a commit 78e2e13

1 file changed

Lines changed: 6 additions & 12 deletions

File tree

crates/memtrack/src/main.rs

Lines changed: 6 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@ use anyhow::{Context, Result, anyhow};
22
use clap::Parser;
33
use ipc_channel::ipc::{self};
44
use log::{debug, info};
5-
use memtrack::{Event, MemtrackIpcMessage, Tracker, handle_ipc_message};
5+
use memtrack::{MemtrackIpcMessage, Tracker, handle_ipc_message};
66
use runner_shared::artifacts::{ArtifactExt, MemtrackArtifact, MemtrackEvent};
77
use std::path::PathBuf;
88
use std::process::Command;
@@ -65,7 +65,6 @@ fn track_command(
6565
cmd_string: &str,
6666
ipc_server_name: Option<String>,
6767
) -> anyhow::Result<(u32, Vec<MemtrackEvent>, std::process::ExitStatus)> {
68-
let events = Arc::new(Mutex::new(Vec::new()));
6968
let tracker = Tracker::new()?;
7069

7170
let tracker_arc = Arc::new(Mutex::new(tracker));
@@ -97,10 +96,10 @@ fn track_command(
9796
info!("Spawned child with pid {root_pid}");
9897

9998
// Spawn event processing thread
100-
let events_clone = events.clone();
10199
let process_events = Arc::new(AtomicBool::new(true));
102100
let process_events_clone = process_events.clone();
103101
let processing_thread = thread::spawn(move || {
102+
let mut events = Vec::new();
104103
loop {
105104
if !process_events_clone.load(Ordering::Relaxed) {
106105
break;
@@ -110,12 +109,9 @@ fn track_command(
110109
continue;
111110
};
112111

113-
let Ok(mut e) = events_clone.lock() else {
114-
continue;
115-
};
116-
117-
e.push(event.into());
112+
events.push(event.into());
118113
}
114+
events
119115
});
120116

121117
// Wait for the command to complete
@@ -124,14 +120,12 @@ fn track_command(
124120

125121
info!("Waiting for the event processing thread to finish");
126122
process_events.store(false, Ordering::Relaxed);
127-
processing_thread
123+
let events = processing_thread
128124
.join()
129125
.map_err(|_| anyhow::anyhow!("Failed to join event thread"))?;
130126

131127
// IPC thread will exit when channel closes
132128
drop(ipc_handle);
133129

134-
info!("Unwrapping and returning events");
135-
let events = Arc::try_unwrap(events).map_err(|_| anyhow::anyhow!("Failed to unwrap events"))?;
136-
Ok((root_pid as u32, events.into_inner()?, status))
130+
Ok((root_pid as u32, events, status))
137131
}

0 commit comments

Comments
 (0)