Skip to main content

micromegas_analytics/
async_block_processing.rs

1use crate::{metadata::StreamMetadata, payload::parse_block, scope::BorrowedScopeDesc};
2use anyhow::{Context, Result};
3use micromegas_telemetry::block_wire_format::BlockPayload;
4use micromegas_tracing::prelude::*;
5use micromegas_transit::value::{Object, Value};
6
7/// Helper function to extract async event fields
8fn on_async_event<'a, F>(obj: &Object<'a>, mut fun: F) -> Result<bool>
9where
10    F: FnMut(&'a Object<'a>, u64, u64, u32, i64) -> Result<bool>,
11{
12    let span_id = obj.get::<u64>("span_id")?;
13    let parent_span_id = obj.get::<u64>("parent_span_id")?;
14    let depth = obj.get::<u32>("depth")?;
15    let time = obj.get::<i64>("time")?;
16    let span_desc = obj.get::<&Object>("span_desc")?;
17    fun(span_desc, span_id, parent_span_id, depth, time)
18}
19
20/// Helper function to extract async named event fields
21fn on_async_named_event<'a, F>(obj: &Object<'a>, mut fun: F) -> Result<bool>
22where
23    F: FnMut(&'a Object<'a>, &'a str, u64, u64, u32, i64) -> Result<bool>,
24{
25    let span_id = obj.get::<u64>("span_id")?;
26    let parent_span_id = obj.get::<u64>("parent_span_id")?;
27    let depth = obj.get::<u32>("depth")?;
28    let time = obj.get::<i64>("time")?;
29    let span_location = obj.get::<&Object>("span_location")?;
30    let name = obj.get::<&str>("name")?;
31    fun(span_location, name, span_id, parent_span_id, depth, time)
32}
33
34/// Trait for processing async event blocks.
35pub trait AsyncBlockProcessor {
36    fn on_begin_async_scope(
37        &mut self,
38        block_id: &str,
39        scope: BorrowedScopeDesc<'_>,
40        ts: i64,
41        span_id: i64,
42        parent_span_id: i64,
43        depth: u32,
44    ) -> Result<bool>;
45    fn on_end_async_scope(
46        &mut self,
47        block_id: &str,
48        scope: BorrowedScopeDesc<'_>,
49        ts: i64,
50        span_id: i64,
51        parent_span_id: i64,
52        depth: u32,
53    ) -> Result<bool>;
54}
55
56/// Parses async span events from a thread event block payload.
57#[span_fn]
58pub fn parse_async_block_payload<Proc: AsyncBlockProcessor>(
59    block_id: &str,
60    _object_offset: i64,
61    payload: &BlockPayload,
62    stream: &StreamMetadata,
63    processor: &mut Proc,
64) -> Result<bool> {
65    parse_block(stream, payload, |val| {
66        if let Value::Object(obj) = val {
67            match obj.type_name {
68                "BeginAsyncSpanEvent" => {
69                    on_async_event(obj, |span_desc, span_id, parent_span_id, depth, ts| {
70                        let name = span_desc.get::<&str>("name")?;
71                        let filename = span_desc.get::<&str>("file")?;
72                        let target = span_desc.get::<&str>("target")?;
73                        let line = span_desc.get::<u32>("line")?;
74                        let scope_desc = BorrowedScopeDesc::new(name, filename, target, line);
75                        processor.on_begin_async_scope(
76                            block_id,
77                            scope_desc,
78                            ts,
79                            span_id as i64,
80                            parent_span_id as i64,
81                            depth,
82                        )
83                    })
84                    .with_context(|| "reading BeginAsyncSpanEvent")
85                }
86                "EndAsyncSpanEvent" => {
87                    on_async_event(obj, |span_desc, span_id, parent_span_id, depth, ts| {
88                        let name = span_desc.get::<&str>("name")?;
89                        let filename = span_desc.get::<&str>("file")?;
90                        let target = span_desc.get::<&str>("target")?;
91                        let line = span_desc.get::<u32>("line")?;
92                        let scope_desc = BorrowedScopeDesc::new(name, filename, target, line);
93                        processor.on_end_async_scope(
94                            block_id,
95                            scope_desc,
96                            ts,
97                            span_id as i64,
98                            parent_span_id as i64,
99                            depth,
100                        )
101                    })
102                    .with_context(|| "reading EndAsyncSpanEvent")
103                }
104                "BeginAsyncNamedSpanEvent" => on_async_named_event(
105                    obj,
106                    |span_location, name, span_id, parent_span_id, depth, ts| {
107                        let filename = span_location.get::<&str>("file")?;
108                        let target = span_location.get::<&str>("target")?;
109                        let line = span_location.get::<u32>("line")?;
110                        let scope_desc = BorrowedScopeDesc::new(name, filename, target, line);
111                        processor.on_begin_async_scope(
112                            block_id,
113                            scope_desc,
114                            ts,
115                            span_id as i64,
116                            parent_span_id as i64,
117                            depth,
118                        )
119                    },
120                )
121                .with_context(|| "reading BeginAsyncNamedSpanEvent"),
122                "EndAsyncNamedSpanEvent" => on_async_named_event(
123                    obj,
124                    |span_location, name, span_id, parent_span_id, depth, ts| {
125                        let filename = span_location.get::<&str>("file")?;
126                        let target = span_location.get::<&str>("target")?;
127                        let line = span_location.get::<u32>("line")?;
128                        let scope_desc = BorrowedScopeDesc::new(name, filename, target, line);
129                        processor.on_end_async_scope(
130                            block_id,
131                            scope_desc,
132                            ts,
133                            span_id as i64,
134                            parent_span_id as i64,
135                            depth,
136                        )
137                    },
138                )
139                .with_context(|| "reading EndAsyncNamedSpanEvent"),
140                _ => Ok(true),
141            }
142        } else {
143            Ok(true)
144        }
145    })
146}