1- use perfetto_format:: {
2- create_callstack_interned_data, create_debug_annotation, DebugValue , PerfettoStreamWriter ,
3- } ;
1+ use perfetto_format:: { create_debug_annotation, DebugValue , PerfettoStreamWriter } ;
42use std:: fs:: File ;
53use std:: io:: BufWriter ;
6- use std:: thread;
74use std:: time:: { Duration , SystemTime , UNIX_EPOCH } ;
5+ use std:: { process, thread} ;
86
97fn 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