1use crate::metadata::StreamMetadata;
2use crate::net_block_processing::{NetBlockProcessor, parse_net_block};
3use crate::net_spans_table::{NetSpanRecord, NetSpanRecordBuilder};
4use crate::time::ConvertTicks;
5use anyhow::Result;
6use micromegas_telemetry::blob_storage::BlobStorage;
7use micromegas_telemetry::types::block::BlockMetadata;
8use micromegas_tracing::prelude::*;
9use std::sync::Arc;
10
11lazy_static::lazy_static! {
12 static ref KIND_CONNECTION: Arc<String> = Arc::new(String::from("connection"));
13 static ref KIND_OBJECT: Arc<String> = Arc::new(String::from("object"));
14 static ref KIND_PROPERTY: Arc<String> = Arc::new(String::from("property"));
15 static ref KIND_RPC: Arc<String> = Arc::new(String::from("rpc"));
16 static ref EMPTY_NAME: Arc<String> = Arc::new(String::new());
17}
18
19pub const ROOT_PARENT_SPAN_ID: i64 = -1;
23
24#[derive(Debug)]
25struct OpenSpan {
26 span_id: i64,
27 parent_span_id: i64,
28 depth: u32,
29 kind: Arc<String>,
30 name: Arc<String>,
31 connection_name: Arc<String>,
32 is_outgoing: bool,
33 begin_time_ns: i64,
34 child_bits_consumed: i64,
36}
37
38pub struct NetSpanTreeBuilder<'a> {
44 record_builder: &'a mut NetSpanRecordBuilder,
45 stack: Vec<OpenSpan>,
46 process_id: Arc<String>,
47 stream_id: Arc<String>,
48 convert_ticks: ConvertTicks,
49}
50
51impl<'a> NetSpanTreeBuilder<'a> {
52 pub fn new(
53 record_builder: &'a mut NetSpanRecordBuilder,
54 process_id: Arc<String>,
55 stream_id: Arc<String>,
56 convert_ticks: ConvertTicks,
57 ) -> Self {
58 Self {
59 record_builder,
60 stack: Vec::new(),
61 process_id,
62 stream_id,
63 convert_ticks,
64 }
65 }
66
67 fn connection_context(&self) -> (Arc<String>, bool) {
70 if let Some(root) = self.stack.first() {
71 (root.connection_name.clone(), root.is_outgoing)
72 } else {
73 (EMPTY_NAME.clone(), false)
74 }
75 }
76
77 fn parent_of_new_child(&self) -> (i64, u32, i64) {
78 if let Some(top) = self.stack.last() {
79 (top.span_id, top.depth + 1, top.child_bits_consumed)
80 } else {
81 (ROOT_PARENT_SPAN_ID, 0, 0)
83 }
84 }
85
86 fn close_span(
91 &mut self,
92 expected_kind: &Arc<String>,
93 event_time_ns: i64,
94 bit_size: i64,
95 ) -> Result<bool> {
96 match self.stack.last() {
97 None => {
98 debug!(
100 "net span end event with no matching begin (expected kind={})",
101 expected_kind
102 );
103 return Ok(true);
104 }
105 Some(top) if !Arc::ptr_eq(&top.kind, expected_kind) && *top.kind != **expected_kind => {
106 debug!(
110 "net span stack mismatch: expected {}, got {}; skipping end event",
111 expected_kind, top.kind
112 );
113 return Ok(true);
114 }
115 Some(_) => {}
116 }
117 let open = self.stack.pop().expect("peeked above");
118 let begin_bits = if let Some(parent) = self.stack.last() {
119 parent.child_bits_consumed
120 } else {
121 0
122 };
123 let end_bits = begin_bits + bit_size;
124 let connection_name = if self.stack.is_empty() {
125 open.connection_name.clone()
126 } else {
127 self.stack[0].connection_name.clone()
128 };
129 let is_outgoing = if self.stack.is_empty() {
130 open.is_outgoing
131 } else {
132 self.stack[0].is_outgoing
133 };
134 let record = NetSpanRecord {
135 process_id: self.process_id.clone(),
136 stream_id: self.stream_id.clone(),
137 span_id: open.span_id,
138 parent_span_id: open.parent_span_id,
139 depth: open.depth,
140 kind: open.kind.clone(),
141 name: open.name.clone(),
142 connection_name,
143 is_outgoing,
144 begin_bits,
145 end_bits,
146 bit_size,
147 begin_time: open.begin_time_ns,
148 end_time: event_time_ns,
149 };
150 self.record_builder.append(&record)?;
151 if let Some(parent) = self.stack.last_mut() {
152 parent.child_bits_consumed += bit_size;
153 }
154 Ok(true)
155 }
156
157 pub fn finish(self) {
160 if !self.stack.is_empty() {
161 debug!(
162 "net span tree finishing with {} unclosed span(s); dropping",
163 self.stack.len()
164 );
165 }
166 }
167}
168
169impl<'a> NetBlockProcessor for NetSpanTreeBuilder<'a> {
170 fn on_connection_begin(
171 &mut self,
172 event_id: i64,
173 time: i64,
174 connection_name: &str,
175 is_outgoing: bool,
176 ) -> Result<bool> {
177 let begin_time_ns = self.convert_ticks.ticks_to_nanoseconds(time);
178 let connection_name = Arc::new(connection_name.to_owned());
180 self.stack.push(OpenSpan {
181 span_id: event_id,
182 parent_span_id: ROOT_PARENT_SPAN_ID,
183 depth: 0,
184 kind: KIND_CONNECTION.clone(),
185 name: connection_name.clone(),
186 connection_name,
187 is_outgoing,
188 begin_time_ns,
189 child_bits_consumed: 0,
190 });
191 Ok(true)
192 }
193
194 fn on_connection_end(&mut self, _event_id: i64, time: i64, bit_size: i64) -> Result<bool> {
195 let end_time_ns = self.convert_ticks.ticks_to_nanoseconds(time);
196 self.close_span(&KIND_CONNECTION, end_time_ns, bit_size)
197 }
198
199 fn on_object_begin(&mut self, event_id: i64, time: i64, object_name: &str) -> Result<bool> {
200 let begin_time_ns = self.convert_ticks.ticks_to_nanoseconds(time);
201 let (parent_span_id, depth, _) = self.parent_of_new_child();
202 let (connection_name, is_outgoing) = self.connection_context();
203 self.stack.push(OpenSpan {
204 span_id: event_id,
205 parent_span_id,
206 depth,
207 kind: KIND_OBJECT.clone(),
208 name: Arc::new(object_name.to_owned()),
209 connection_name,
210 is_outgoing,
211 begin_time_ns,
212 child_bits_consumed: 0,
213 });
214 Ok(true)
215 }
216
217 fn on_object_end(&mut self, _event_id: i64, time: i64, bit_size: i64) -> Result<bool> {
218 let end_time_ns = self.convert_ticks.ticks_to_nanoseconds(time);
219 self.close_span(&KIND_OBJECT, end_time_ns, bit_size)
220 }
221
222 fn on_property(
223 &mut self,
224 event_id: i64,
225 time: i64,
226 property_name: &str,
227 bit_size: i64,
228 ) -> Result<bool> {
229 let event_time_ns = self.convert_ticks.ticks_to_nanoseconds(time);
230 let (parent_span_id, depth, begin_bits) = self.parent_of_new_child();
231 let end_bits = begin_bits + bit_size;
232 let (connection_name, is_outgoing) = self.connection_context();
233 let record = NetSpanRecord {
234 process_id: self.process_id.clone(),
235 stream_id: self.stream_id.clone(),
236 span_id: event_id,
237 parent_span_id,
238 depth,
239 kind: KIND_PROPERTY.clone(),
240 name: Arc::new(property_name.to_owned()),
241 connection_name,
242 is_outgoing,
243 begin_bits,
244 end_bits,
245 bit_size,
246 begin_time: event_time_ns,
247 end_time: event_time_ns,
248 };
249 self.record_builder.append(&record)?;
250 if let Some(parent) = self.stack.last_mut() {
251 parent.child_bits_consumed += bit_size;
252 }
253 Ok(true)
254 }
255
256 fn on_rpc_begin(&mut self, event_id: i64, time: i64, function_name: &str) -> Result<bool> {
257 let begin_time_ns = self.convert_ticks.ticks_to_nanoseconds(time);
258 let (parent_span_id, depth, _) = self.parent_of_new_child();
259 let (connection_name, is_outgoing) = self.connection_context();
260 self.stack.push(OpenSpan {
261 span_id: event_id,
262 parent_span_id,
263 depth,
264 kind: KIND_RPC.clone(),
265 name: Arc::new(function_name.to_owned()),
266 connection_name,
267 is_outgoing,
268 begin_time_ns,
269 child_bits_consumed: 0,
270 });
271 Ok(true)
272 }
273
274 fn on_rpc_end(&mut self, _event_id: i64, time: i64, bit_size: i64) -> Result<bool> {
275 let end_time_ns = self.convert_ticks.ticks_to_nanoseconds(time);
276 self.close_span(&KIND_RPC, end_time_ns, bit_size)
277 }
278}
279
280#[span_fn]
283pub async fn make_net_span_tree(
284 blocks: &[BlockMetadata],
285 record_builder: &mut NetSpanRecordBuilder,
286 blob_storage: Arc<BlobStorage>,
287 stream: &StreamMetadata,
288 process_id: Arc<String>,
289 convert_ticks: ConvertTicks,
290) -> Result<()> {
291 let stream_id = Arc::new(stream.stream_id.to_string());
292 let mut builder = NetSpanTreeBuilder::new(record_builder, process_id, stream_id, convert_ticks);
293 for block in blocks {
294 parse_net_block(
295 blob_storage.clone(),
296 stream,
297 block.block_id,
298 block.object_offset,
299 &mut builder,
300 )
301 .await?;
302 }
303 builder.finish();
304 Ok(())
305}