Sitelet https://github.com/category-labs/manytrace/commit/278e42bb6c82ba8515505817cbf15822cf3f84a2
Skip to content

Commit 278e42b

Browse files
committed
usertrace: implement bpftrace binary
1 parent f790fcb commit 278e42b

13 files changed

Lines changed: 626 additions & 45 deletions

File tree

‎Cargo.lock‎

Lines changed: 2 additions & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

‎agent/src/agent.rs‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -269,6 +269,7 @@ mod tests {
269269
tid: 1,
270270
pid: 2,
271271
labels: Cow::Owned(Labels::new()),
272+
unit: None,
272273
});
273274

274275
agent.submit(&event).unwrap();
@@ -294,6 +295,7 @@ mod tests {
294295
tid: 1,
295296
pid: 2,
296297
labels: Cow::Owned(Labels::new()),
298+
unit: None,
297299
});
298300

299301
agent.submit(&event).unwrap();

‎agent/src/mpsc.rs‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -108,6 +108,7 @@ mod tests {
108108
tid: 123,
109109
pid: 456,
110110
labels: Cow::Owned(Labels::new()),
111+
unit: None
111112
});
112113

113114
producer.submit(&counter).unwrap();

‎bpf/src/bpf/cpuutil.bpf.c‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -177,7 +177,7 @@ int handle_boundary_event(void *ctx)
177177
u32 tid = pid_tgid;
178178
u32 tgid = pid_tgid >> 32;
179179

180-
if (state->tid != tid || tid == 0) {
180+
if (state->tid != tid || tid == 0 || !should_track_tgid(tgid)) {
181181
return 0;
182182
}
183183

‎bpf/src/cpuutil.rs‎

Lines changed: 18 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -56,12 +56,12 @@ impl Object {
5656
where
5757
F: for<'a> FnMut(Event<'a>) + 'bd,
5858
{
59-
let mut util = CpuUtil::new(&mut self.object, callback, self.config.interval_ms)?;
60-
61-
for pid in &self.config.pid_filters {
62-
util.add_pid_filter(*pid)?;
63-
}
64-
59+
let util = CpuUtil::new(
60+
&mut self.object,
61+
self.config.clone(),
62+
callback,
63+
self.config.interval_ms,
64+
)?;
6565
Ok(util)
6666
}
6767
}
@@ -112,6 +112,7 @@ where
112112
{
113113
fn new(
114114
open_object: &'this mut MaybeUninit<OpenObject>,
115+
config: CpuUtilConfig,
115116
callback: F,
116117
interval_ms: u64,
117118
) -> Result<Self, BpfError> {
@@ -125,6 +126,15 @@ where
125126
.load()
126127
.map_err(|e| BpfError::LoadError(format!("failed to load bpf program: {}", e)))?;
127128

129+
for &pid in &config.pid_filters {
130+
let key = pid.to_ne_bytes();
131+
let value = 1u32.to_ne_bytes();
132+
skel.maps
133+
.tracked_tgids
134+
.update(&key, &value, libbpf_rs::MapFlags::ANY)
135+
.map_err(|e| BpfError::MapError(format!("failed to update filter map: {}", e)))?;
136+
}
137+
128138
skel.attach()
129139
.map_err(|e| BpfError::AttachError(format!("failed to attach bpf programs: {}", e)))?;
130140

@@ -248,6 +258,7 @@ where
248258
tid,
249259
pid,
250260
labels: Cow::Owned(Labels::new()),
261+
unit: Some("ns"),
251262
});
252263
callback(cpu_counter);
253264
}
@@ -260,6 +271,7 @@ where
260271
tid,
261272
pid,
262273
labels: Cow::Owned(Labels::new()),
274+
unit: Some("ns"),
263275
});
264276
callback(kernel_counter);
265277
}

‎bpf/src/profiler.rs‎

Lines changed: 9 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -162,17 +162,13 @@ where
162162
.load()
163163
.map_err(|e| BpfError::LoadError(format!("failed to load bpf program: {}", e)))?;
164164

165-
if !config.pid_filters.is_empty() {
166-
for &pid in &config.pid_filters {
167-
let key = pid.to_ne_bytes();
168-
let value = 1u8.to_ne_bytes();
169-
skel.maps
170-
.filter_tgid
171-
.update(&key, &value, libbpf_rs::MapFlags::ANY)
172-
.map_err(|e| {
173-
BpfError::MapError(format!("failed to update filter map: {}", e))
174-
})?;
175-
}
165+
for &pid in &config.pid_filters {
166+
let key = pid.to_ne_bytes();
167+
let value = 1u8.to_ne_bytes();
168+
skel.maps
169+
.filter_tgid
170+
.update(&key, &value, libbpf_rs::MapFlags::ANY)
171+
.map_err(|e| BpfError::MapError(format!("failed to update filter map: {}", e)))?;
176172
}
177173

178174
let perf_type = libbpf_sys::PERF_TYPE_SOFTWARE;
@@ -357,19 +353,13 @@ struct InternState<'a> {
357353

358354
impl<'a> InternState<'a> {
359355
fn new() -> Self {
360-
let mut state = Self {
361-
string_id_counter: 1,
356+
let state: InternState<'a> = Self {
357+
string_id_counter: 0,
362358
frame_id_counter: 0,
363359
function_names: Vec::new(),
364360
frames: Vec::new(),
365361
frame_ids: Vec::new(),
366362
};
367-
368-
state.function_names.push(InternedString {
369-
iid: 0,
370-
str: Cow::Borrowed("<unknown>"),
371-
});
372-
373363
state
374364
}
375365

‎config.toml‎

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,11 @@
1+
[profiler]
2+
sample_freq = 999
3+
kernel_samples = true
4+
user_samples = true
5+
pid_filters = [264373]
6+
7+
[thread_tracker]
8+
9+
[cpu_util]
10+
interval_ms = 100
11+
pid_filters = [264373]

‎perfetto-format/examples/streaming.rs‎

Lines changed: 107 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,8 @@
1-
use perfetto_format::{
2-
create_callstack_interned_data, create_debug_annotation, DebugValue, PerfettoStreamWriter,
3-
};
1+
use perfetto_format::{create_debug_annotation, DebugValue, PerfettoStreamWriter};
42
use std::fs::File;
53
use std::io::BufWriter;
6-
use std::thread;
74
use std::time::{Duration, SystemTime, UNIX_EPOCH};
5+
use std::{process, thread};
86

97
fn timestamp_nanos() -> u64 {
108
SystemTime::now()
@@ -24,21 +22,51 @@ fn main() -> anyhow::Result<()> {
2422
.write_process_descriptor(std::process::id(), Some("streaming_example".to_string()))?;
2523

2624
let main_thread_track_uuid = writer.write_thread_descriptor(
25+
process::id(),
2726
thread_id::get() as u32,
2827
Some("main_thread".to_string()),
2928
process_track_uuid,
3029
)?;
3130

32-
let callstack_data = create_callstack_interned_data(
33-
0,
34-
vec![
35-
("main", "main.rs", 0x400000),
36-
("process_batch", "processor.rs", 0x401000),
37-
("compute", "compute.rs", 0x402000),
38-
("inner_loop", "compute.rs", 0x402100),
39-
],
40-
);
41-
writer.write_interned_data(callstack_data)?;
31+
writer.write_callstack_interned_data(0, vec![("AA", "A", 0x400000), ("BB", "B", 0x401000)])?;
32+
33+
writer.write_callstack_interned_data(1, vec![("ZZ", "Z", 0x500000), ("XX", "X", 0x501000)])?;
34+
35+
let memory_counter_track = writer.write_counter_track(
36+
"Memory Usage".to_string(),
37+
Some("bytes".to_string()),
38+
main_thread_track_uuid,
39+
)?;
40+
41+
let cpu_counter_track = writer.write_counter_track(
42+
"CPU Usage".to_string(),
43+
Some("%".to_string()),
44+
main_thread_track_uuid,
45+
)?;
46+
47+
let items_processed_track =
48+
writer.write_counter_track("Items Processed".to_string(), None, main_thread_track_uuid)?;
49+
50+
let network_bytes_track = writer.write_counter_track(
51+
"Network Bytes Sent".to_string(),
52+
Some("bytes".to_string()),
53+
main_thread_track_uuid,
54+
)?;
55+
56+
let queue_size_track = writer.write_counter_track(
57+
"Queue Size".to_string(),
58+
Some("items".to_string()),
59+
main_thread_track_uuid,
60+
)?;
61+
62+
let error_count_track =
63+
writer.write_counter_track("Error Count".to_string(), None, main_thread_track_uuid)?;
64+
65+
let latency_track = writer.write_counter_track(
66+
"Processing Latency".to_string(),
67+
Some("ms".to_string()),
68+
main_thread_track_uuid,
69+
)?;
4270

4371
println!("Writing events in real-time...");
4472

@@ -53,9 +81,38 @@ fn main() -> anyhow::Result<()> {
5381
)],
5482
)?;
5583

84+
let mut base_memory = 1024 * 1024 * 100;
85+
let mut items_processed = 0i64;
86+
let mut network_bytes = 0i64;
87+
let mut error_count = 0i64;
88+
5689
for i in 0..100 {
5790
let timestamp = timestamp_nanos();
5891

92+
items_processed += 1;
93+
writer.write_counter_value(items_processed_track, items_processed, timestamp)?;
94+
95+
let memory_usage = base_memory + (i * 1024 * 512);
96+
writer.write_counter_value(memory_counter_track, memory_usage, timestamp)?;
97+
98+
let cpu_usage = 20.0 + 30.0 * ((i as f64 * 0.1).sin() + 1.0);
99+
writer.write_double_counter_value(cpu_counter_track, cpu_usage, timestamp)?;
100+
101+
let bytes_sent = 1024 + (i * 100);
102+
network_bytes += bytes_sent;
103+
writer.write_counter_value(network_bytes_track, network_bytes, timestamp)?;
104+
105+
let queue_size = ((i as f64 * 0.2).sin() * 50.0 + 50.0) as i64;
106+
writer.write_counter_value(queue_size_track, queue_size, timestamp)?;
107+
108+
if i % 17 == 0 && i > 0 {
109+
error_count += 1;
110+
writer.write_counter_value(error_count_track, error_count, timestamp)?;
111+
}
112+
113+
let latency = 10.0 + 5.0 * ((i as f64 * 0.15).cos() + 1.0);
114+
writer.write_double_counter_value(latency_track, latency, timestamp)?;
115+
59116
if i % 10 == 0 {
60117
let batch_start = timestamp;
61118
writer.write_slice_begin(
@@ -100,12 +157,24 @@ fn main() -> anyhow::Result<()> {
100157
0,
101158
perfetto_format::perfetto::profiling::CpuMode::ModeUser,
102159
)?;
160+
writer.write_perf_sample(
161+
0,
162+
std::process::id(),
163+
thread_id::get() as u32,
164+
timestamp + 1_500_000,
165+
1,
166+
perfetto_format::perfetto::profiling::CpuMode::ModeUser,
167+
)?;
103168

104169
writer.write_slice_end(main_thread_track_uuid, timestamp + 2_000_000)?;
105170
}
106171

107172
if i % 10 == 9 {
108173
writer.write_slice_end(main_thread_track_uuid, timestamp + 5_000_000)?;
174+
175+
base_memory += 1024 * 1024 * 5;
176+
let spike_timestamp = timestamp + 5_500_000;
177+
writer.write_counter_value(memory_counter_track, base_memory, spike_timestamp)?;
109178
}
110179

111180
if i % 25 == 0 && i > 0 {
@@ -125,12 +194,35 @@ fn main() -> anyhow::Result<()> {
125194

126195
writer.write_slice_end(main_thread_track_uuid, timestamp_nanos())?;
127196

197+
let final_timestamp = timestamp_nanos();
198+
199+
writer.write_counter_value(items_processed_track, items_processed, final_timestamp)?;
200+
writer.write_counter_value(memory_counter_track, base_memory, final_timestamp)?;
201+
writer.write_double_counter_value(cpu_counter_track, 0.0, final_timestamp)?;
202+
writer.write_counter_value(network_bytes_track, network_bytes, final_timestamp)?;
203+
writer.write_counter_value(queue_size_track, 0, final_timestamp)?;
204+
writer.write_counter_value(error_count_track, error_count, final_timestamp)?;
205+
writer.write_double_counter_value(latency_track, 0.0, final_timestamp)?;
206+
128207
writer.write_instant_event(
129208
main_thread_track_uuid,
130209
"trace_complete".to_string(),
131-
timestamp_nanos(),
210+
final_timestamp,
132211
vec![
133212
create_debug_annotation("total_samples".to_string(), DebugValue::Int(100)),
213+
create_debug_annotation(
214+
"items_processed".to_string(),
215+
DebugValue::Int(items_processed),
216+
),
217+
create_debug_annotation(
218+
"final_memory_mb".to_string(),
219+
DebugValue::Int(base_memory / (1024 * 1024)),
220+
),
221+
create_debug_annotation(
222+
"network_mb_sent".to_string(),
223+
DebugValue::Double(network_bytes as f64 / (1024.0 * 1024.0)),
224+
),
225+
create_debug_annotation("error_count".to_string(), DebugValue::Int(error_count)),
134226
create_debug_annotation(
135227
"status".to_string(),
136228
DebugValue::String("success".to_string()),

0 commit comments

Comments
 (0)