@@ -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 } ;
1010use serde:: { Deserialize , Serialize } ;
1111use std:: marker:: PhantomData ;
1212use 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
116123impl < ' obj , F > Profiler < ' obj , F >
117124where
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
260274impl < ' obj , F > Filterable for Profiler < ' obj , F >
261275where
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