Sitelet https://github.com/category-labs/manytrace/commit/2a7f8d5bbb13326aa84736ae7198806db6f64527
Skip to content

Commit 2a7f8d5

Browse files
committed
bpf, manytrace, perfetto-format, protocol: complete sequence_id removal
remove self.sequence_id from perfetto-format writer add stream_id parameter to methods that write packets update all callers to pass stream_id when available
1 parent 3fe01da commit 2a7f8d5

10 files changed

Lines changed: 218 additions & 166 deletions

File tree

‎bpf/examples/cpuutil.rs‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
use bpf::{cpuutil, CpuUtilConfig};
2-
use protocol::Event;
2+
use protocol::{Event, Message};
33
use std::sync::atomic::{AtomicBool, Ordering};
44
use std::sync::Arc;
55
use std::thread::sleep;
@@ -25,8 +25,8 @@ fn main() -> Result<(), Box<dyn std::error::Error>> {
2525
pid_filters: vec![std::process::id()],
2626
filter_process: vec![],
2727
});
28-
let mut tracker = builder.build(|event: Event| {
29-
if let Event::Counter(counter) = event {
28+
let mut tracker = builder.build(|message: Message| {
29+
if let Message::Event(Event::Counter(counter)) = message {
3030
match counter.name {
3131
"cpu_time_ns" => {
3232
println!(

‎bpf/examples/threadtrack.rs‎

Lines changed: 20 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
use bpf::threadtrack;
2-
use protocol::Event;
2+
use protocol::{Event, Message};
33
use std::sync::atomic::{AtomicBool, Ordering};
44
use std::sync::Arc;
55
use std::time::Duration;
@@ -17,20 +17,26 @@ fn main() -> Result<(), Box<dyn std::error::Error>> {
1717
println!("{:-<80}", "");
1818

1919
let mut builder = threadtrack::Object::new();
20-
let mut tracker = builder.build(|event: Event| match event {
21-
Event::ThreadName(thread) => {
22-
println!(
23-
"{:<8} {:<8} {:<64} {:<40}",
24-
thread.pid, thread.tid, thread.name, "Thread"
25-
);
20+
let mut tracker = builder.build(|message: Message| {
21+
let event = match message {
22+
Message::Event(e) => e,
23+
_ => return,
24+
};
25+
match event {
26+
Event::ThreadName(thread) => {
27+
println!(
28+
"{:<8} {:<8} {:<64} {:<40}",
29+
thread.pid, thread.tid, thread.name, "Thread"
30+
);
31+
}
32+
Event::ProcessName(process) => {
33+
println!(
34+
"{:<8} {:<8} {:<64} {:<40}",
35+
process.pid, "-", process.name, "Process"
36+
);
37+
}
38+
_ => {}
2639
}
27-
Event::ProcessName(process) => {
28-
println!(
29-
"{:<8} {:<8} {:<64} {:<40}",
30-
process.pid, "-", process.name, "Process"
31-
);
32-
}
33-
_ => {}
3440
})?;
3541

3642
while running.load(Ordering::SeqCst) {

‎bpf/src/cpuutil.rs‎

Lines changed: 12 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,7 @@ use crate::{perf_event, BpfError, Filterable};
88
use libbpf_rs::skel::{OpenSkel, Skel, SkelBuilder};
99
use libbpf_rs::{MapCore, MapFlags, OpenObject, RingBufferBuilder};
1010
use libbpf_sys::{PERF_COUNT_SW_CPU_CLOCK, PERF_TYPE_SOFTWARE};
11-
use protocol::{Counter, Event, Labels};
11+
use protocol::{Counter, Event, Labels, Message};
1212
use serde::{Deserialize, Serialize};
1313
use std::borrow::Cow;
1414
use std::collections::HashMap;
@@ -47,7 +47,7 @@ impl Object {
4747

4848
pub fn build<'bd, F>(&'bd mut self, callback: F) -> Result<CpuUtil<'bd, F>, BpfError>
4949
where
50-
F: for<'a> FnMut(Event<'a>) + 'bd,
50+
F: for<'a> FnMut(Message<'a>) + 'bd,
5151
{
5252
let util = CpuUtil::new(
5353
&mut self.object,
@@ -107,7 +107,7 @@ pub struct CpuUtil<'this, F> {
107107

108108
impl<'this, F> CpuUtil<'this, F>
109109
where
110-
F: for<'a> FnMut(Event<'a>) + 'this,
110+
F: for<'a> FnMut(Message<'a>) + 'this,
111111
{
112112
fn new(
113113
open_object: &'this mut MaybeUninit<OpenObject>,
@@ -180,15 +180,15 @@ where
180180
elapsed_ns = elapsed_ns as u64,
181181
"emitting cpu_time counter"
182182
);
183-
let cpu_counter = Event::Counter(Counter {
183+
let cpu_counter = Message::Event(Event::Counter(Counter {
184184
name: "cpu_time",
185185
value: cpu_percent,
186186
timestamp: thread_stats.min_timestamp,
187187
tid: *tid,
188188
pid: *pid,
189189
labels: Cow::Owned(Labels::new()),
190190
unit: Some("%"),
191-
});
191+
}));
192192
callback(cpu_counter);
193193

194194
let kernel_percent =
@@ -201,15 +201,15 @@ where
201201
elapsed_ns = elapsed_ns as u64,
202202
"emitting kernel_time counter"
203203
);
204-
let kernel_counter = Event::Counter(Counter {
204+
let kernel_counter = Message::Event(Event::Counter(Counter {
205205
name: "kernel_time",
206206
value: kernel_percent,
207207
timestamp: thread_stats.min_timestamp,
208208
tid: *tid,
209209
pid: *pid,
210210
labels: Cow::Owned(Labels::new()),
211211
unit: Some("%"),
212-
});
212+
}));
213213
callback(kernel_counter);
214214
thread_stats.cpu_time_ns = 0;
215215
thread_stats.kernel_time_ns = 0;
@@ -276,7 +276,7 @@ where
276276

277277
impl<'this, F> Filterable for CpuUtil<'this, F>
278278
where
279-
F: for<'a> FnMut(Event<'a>) + 'this,
279+
F: for<'a> FnMut(Message<'a>) + 'this,
280280
{
281281
fn filter(&mut self, pid: i32) -> Result<(), BpfError> {
282282
self.add_pid_filter(pid as u32)
@@ -336,8 +336,8 @@ mod root_tests {
336336

337337
let mut object = Object::new(config);
338338
let mut cpuutil = object
339-
.build(move |event| {
340-
if let Event::Counter(c) = event {
339+
.build(move |message| {
340+
if let Message::Event(Event::Counter(c)) = message {
341341
let test_counter = TestCounter {
342342
name: c.name.to_string(),
343343
value: c.value,
@@ -405,8 +405,8 @@ mod root_tests {
405405

406406
let mut object = Object::new(config);
407407
let mut cpuutil = object
408-
.build(move |event| {
409-
if let Event::Counter(c) = event {
408+
.build(move |message| {
409+
if let Message::Event(Event::Counter(c)) = message {
410410
let test_counter = TestCounter {
411411
name: c.name.to_string(),
412412
value: c.value,

‎bpf/src/lib.rs‎

Lines changed: 19 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,12 @@
11
use blazesym::symbolize::Symbolizer;
2-
use protocol::Event;
2+
use protocol::{Message, StreamId, StreamIdAllocator};
33
use serde::{Deserialize, Serialize};
44
use std::{cell::RefCell, collections::HashSet, path::Path, rc::Rc};
55
use thiserror::Error;
66
use tracing::debug;
77

8+
pub use protocol::Message as BpfMessage;
9+
810
pub mod cpuutil;
911
mod perf_event;
1012
pub mod profiler;
@@ -154,19 +156,24 @@ pub struct BpfObject {
154156
}
155157

156158
impl BpfObject {
157-
pub fn consumer<'this, F>(&'this mut self, callback: F) -> Result<BpfConsumer<'this>, BpfError>
159+
pub fn consumer<'this, F>(
160+
&'this mut self,
161+
callback: F,
162+
stream_allocator: &mut StreamIdAllocator,
163+
) -> Result<BpfConsumer<'this>, BpfError>
158164
where
159-
F: for<'a> FnMut(Event<'a>) + Clone + 'this,
165+
F: for<'a> FnMut(Message<'a>) + Clone + 'this,
160166
{
161167
let cpuutil = if let Some(ref mut obj) = self.cpuutils {
162-
Some(obj.build(Box::new(callback.clone()) as Box<dyn for<'a> FnMut(Event<'a>)>)?)
168+
Some(obj.build(Box::new(callback.clone()) as Box<dyn for<'a> FnMut(Message<'a>)>)?)
163169
} else {
164170
None
165171
};
166172

167173
let profiler = if let Some(ref mut obj) = self.profiler {
168-
let callback = Box::new(callback.clone()) as Box<dyn for<'a> FnMut(Event<'a>)>;
169-
Some(obj.build(callback, &self.symbolizer)?)
174+
let callback = Box::new(callback.clone()) as Box<dyn for<'a> FnMut(Message<'a>)>;
175+
let stream_id = stream_allocator.allocate();
176+
Some(obj.build(callback, &self.symbolizer, stream_id)?)
170177
} else {
171178
None
172179
};
@@ -185,8 +192,8 @@ impl BpfObject {
185192
let cpuutil_ref = cpuutil_rc.clone();
186193
let profiler_ref = profiler_rc.clone();
187194

188-
let wrapper_callback = move |event: Event<'_>| {
189-
if let Event::ProcessName(ref pn) = event {
195+
let wrapper_callback = move |message: Message<'_>| {
196+
if let Message::Event(protocol::Event::ProcessName(ref pn)) = message {
190197
let pid = pn.pid;
191198
let name = pn.name;
192199

@@ -216,13 +223,13 @@ impl BpfObject {
216223
}
217224
}
218225

219-
user_callback(event);
226+
user_callback(message);
220227
};
221228

222-
let callback = Box::new(wrapper_callback) as Box<dyn for<'a> FnMut(Event<'a>)>;
229+
let callback = Box::new(wrapper_callback) as Box<dyn for<'a> FnMut(Message<'a>)>;
223230
Some(obj.build(callback)?)
224231
} else {
225-
let callback = Box::new(callback) as Box<dyn for<'a> FnMut(Event<'a>)>;
232+
let callback = Box::new(callback) as Box<dyn for<'a> FnMut(Message<'a>)>;
226233
Some(obj.build(callback)?)
227234
}
228235
} else {
@@ -237,7 +244,7 @@ impl BpfObject {
237244
}
238245
}
239246

240-
type Callback<'cb> = Box<dyn for<'a> FnMut(Event<'a>) + 'cb>;
247+
type Callback<'cb> = Box<dyn for<'a> FnMut(Message<'a>) + 'cb>;
241248

242249
pub struct BpfConsumer<'this> {
243250
threadtrack: Option<threadtrack::ThreadTracker<'this, Callback<'this>>>,

‎bpf/src/profiler.rs‎

Lines changed: 39 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ use libbpf_rs::{
66
skel::{OpenSkel, SkelBuilder},
77
MapCore, RingBuffer, RingBufferBuilder,
88
};
9-
use protocol::{CpuMode, Event, Sample};
9+
use protocol::{CpuMode, Event, Message, Sample};
1010
use serde::{Deserialize, Serialize};
1111
use std::marker::PhantomData;
1212
use std::mem::MaybeUninit;
@@ -98,11 +98,18 @@ impl Object {
9898
&'obj mut self,
9999
callback: F,
100100
symbolizer: &'obj Symbolizer,
101+
stream_id: crate::StreamId,
101102
) -> Result<Profiler<'obj, F>, BpfError>
102103
where
103-
F: for<'a> FnMut(Event<'a>) + 'obj,
104+
F: for<'a> FnMut(Message<'a>) + 'obj,
104105
{
105-
Profiler::new(&mut self.object, callback, self.config.clone(), symbolizer)
106+
Profiler::new(
107+
&mut self.object,
108+
callback,
109+
self.config.clone(),
110+
symbolizer,
111+
stream_id,
112+
)
106113
}
107114
}
108115

@@ -115,13 +122,14 @@ pub struct Profiler<'obj, F> {
115122

116123
impl<'obj, F> Profiler<'obj, F>
117124
where
118-
F: for<'a> FnMut(Event<'a>) + 'obj,
125+
F: for<'a> FnMut(Message<'a>) + 'obj,
119126
{
120127
fn new(
121128
open_object: &'obj mut MaybeUninit<libbpf_rs::OpenObject>,
122129
mut callback: F,
123130
config: ProfilerConfig,
124131
symbolizer: &'obj Symbolizer,
132+
stream_id: crate::StreamId,
125133
) -> Result<Self, BpfError> {
126134
let skel_builder = ProfilerSkelBuilder::default();
127135
let mut open_skel = skel_builder
@@ -213,11 +221,17 @@ where
213221
let (callstack_iid, interned_data_opt) = interned.data();
214222

215223
if let Some(interned_data) = interned_data_opt {
216-
callback(Event::InternedData(interned_data));
224+
callback(Message::Stream {
225+
stream_id,
226+
event: Event::InternedData(interned_data),
227+
});
217228
}
218229

219230
let sample = create_sample(event, callstack_iid);
220-
callback(Event::Sample(sample));
231+
callback(Message::Stream {
232+
stream_id,
233+
event: Event::Sample(sample),
234+
});
221235
}
222236
Err(err) => {
223237
warn!(err = %err, pid = %pid, "failed to symbolize stack");
@@ -259,7 +273,7 @@ where
259273

260274
impl<'obj, F> Filterable for Profiler<'obj, F>
261275
where
262-
F: for<'a> FnMut(Event<'a>) + 'obj,
276+
F: for<'a> FnMut(Message<'a>) + 'obj,
263277
{
264278
fn filter(&mut self, pid: i32) -> Result<(), BpfError> {
265279
let key = pid.to_ne_bytes();
@@ -350,21 +364,28 @@ mod root_tests {
350364
let mut object = Object::new(config);
351365
let mut profiler = object
352366
.build(
353-
move |event| match &event {
354-
Event::Sample(sample) => {
355-
*sample_count_ref += 1;
356-
thread_ids_ref.insert(sample.tid);
357-
}
358-
Event::InternedData(data) => {
359-
*interned_data_count_ref += 1;
360-
for func in &data.function_names {
361-
let name = func.str.to_string();
362-
function_names_ref.push(name.clone());
367+
move |message| {
368+
let event = match &message {
369+
Message::Stream { event, .. } => event,
370+
Message::Event(e) => e,
371+
};
372+
match event {
373+
Event::Sample(sample) => {
374+
*sample_count_ref += 1;
375+
thread_ids_ref.insert(sample.tid);
376+
}
377+
Event::InternedData(data) => {
378+
*interned_data_count_ref += 1;
379+
for func in &data.function_names {
380+
let name = func.str.to_string();
381+
function_names_ref.push(name.clone());
382+
}
363383
}
384+
_ => {}
364385
}
365-
_ => {}
366386
},
367387
&symbolizer,
388+
0,
368389
)
369390
.expect("failed to create profiler");
370391

0 commit comments

Comments
 (0)