1use crate::{
2 metadata::{ProcessMetadata, StreamMetadata},
3 payload::{fetch_block_payload, parse_block},
4 properties::property_set::PropertySet,
5 time::ConvertTicks,
6};
7use anyhow::{Context, Result};
8use micromegas_telemetry::{blob_storage::BlobStorage, types::block::BlockMetadata};
9use micromegas_tracing::prelude::*;
10use micromegas_transit::value::{Object, Value};
11use std::sync::Arc;
12
13#[derive(Debug)]
19pub struct LogEntry<'a> {
20 pub process: Arc<ProcessMetadata>,
21 pub stream_id: Arc<String>,
22 pub block_id: Arc<String>,
23 pub insert_time: i64,
24 pub time: i64,
25 pub level: i32,
26 pub target: &'a str,
27 pub msg: &'a str,
28 pub properties: PropertySet<'a>,
29}
30
31#[span_fn]
33pub fn log_entry_from_value<'a>(
34 convert_ticks: &ConvertTicks,
35 process: Arc<ProcessMetadata>,
36 stream_id: Arc<String>,
37 block_id: Arc<String>,
38 block_insert_time_ns: i64,
39 val: Value<'a>,
40) -> Result<Option<LogEntry<'a>>> {
41 if let Value::Object(obj) = val {
42 match obj.type_name {
43 "LogStaticStrEvent" => {
44 let ticks = obj
45 .get::<i64>("time")
46 .with_context(|| "reading time from LogStaticStrEvent")?;
47 let desc = obj
48 .get::<&Object>("desc")
49 .with_context(|| "reading desc from LogStaticStrEvent")?;
50 let level = desc
51 .get::<u32>("level")
52 .with_context(|| "reading level from LogStaticStrEvent")?;
53 let target = desc
54 .get::<&str>("target")
55 .with_context(|| "reading target from LogStaticStrEvent")?;
56 let msg = desc
57 .get::<&str>("fmt_str")
58 .with_context(|| "reading fmt_str from LogStaticStrEvent")?;
59 Ok(Some(LogEntry {
60 process,
61 stream_id,
62 block_id,
63 insert_time: block_insert_time_ns,
64 time: convert_ticks.ticks_to_nanoseconds(ticks),
65 level: level as i32,
66 target,
67 msg,
68 properties: PropertySet::empty(),
69 }))
70 }
71 "LogStringEvent" | "LogStringEventV2" => {
72 let ticks = obj
73 .get::<i64>("time")
74 .with_context(|| "reading time from LogStringEvent")?;
75 let desc = obj
76 .get::<&Object>("desc")
77 .with_context(|| "reading desc from LogStringEvent")?;
78 let level = desc
79 .get::<u32>("level")
80 .with_context(|| "reading level from LogStringEvent")?;
81 let target = desc
82 .get::<&str>("target")
83 .with_context(|| "reading target from LogStringEvent")?;
84 let msg = obj
85 .get::<&str>("msg")
86 .with_context(|| "reading msg from LogStringEvent")?;
87 Ok(Some(LogEntry {
88 process,
89 stream_id,
90 block_id,
91 insert_time: block_insert_time_ns,
92 time: convert_ticks.ticks_to_nanoseconds(ticks),
93 level: level as i32,
94 target,
95 msg,
96 properties: PropertySet::empty(),
97 }))
98 }
99 "LogStaticStrInteropEvent" | "LogStringInteropEventV2" | "LogStringInteropEventV3" => {
100 let ticks = obj
101 .get::<i64>("time")
102 .with_context(|| format!("reading time from {}", obj.type_name))?;
103 let level = obj
104 .get::<u32>("level")
105 .with_context(|| format!("reading level from {}", obj.type_name))?;
106 let target = obj
107 .get::<&str>("target")
108 .with_context(|| format!("reading target from {}", obj.type_name))?;
109 let msg = obj
110 .get::<&str>("msg")
111 .with_context(|| format!("reading msg from {}", obj.type_name))?;
112 Ok(Some(LogEntry {
113 process,
114 stream_id,
115 block_id,
116 insert_time: block_insert_time_ns,
117 time: convert_ticks.ticks_to_nanoseconds(ticks),
118 level: level as i32,
119 target,
120 msg,
121 properties: PropertySet::empty(),
122 }))
123 }
124 "TaggedLogInteropEvent" => {
125 let ticks = obj
126 .get::<i64>("time")
127 .with_context(|| format!("reading time from {}", obj.type_name))?;
128 let level = obj
129 .get::<u32>("level")
130 .with_context(|| format!("reading level from {}", obj.type_name))?;
131 let target = obj
132 .get::<&str>("target")
133 .with_context(|| format!("reading target from {}", obj.type_name))?;
134 let msg = obj
135 .get::<&str>("msg")
136 .with_context(|| format!("reading msg from {}", obj.type_name))?;
137 let properties = obj
138 .get::<&Object>("properties")
139 .with_context(|| format!("reading properties from {}", obj.type_name))?;
140 let time = convert_ticks.ticks_to_nanoseconds(ticks);
141 Ok(Some(LogEntry {
142 process,
143 stream_id,
144 block_id,
145 insert_time: block_insert_time_ns,
146 time,
147 level: level as i32,
148 target,
149 msg,
150 properties: properties.into(),
151 }))
152 }
153 "TaggedLogString" => {
154 let ticks = obj
155 .get::<i64>("time")
156 .with_context(|| format!("reading time from {}", obj.type_name))?;
157 let msg = obj
158 .get::<&str>("msg")
159 .with_context(|| format!("reading msg from {}", obj.type_name))?;
160 let desc = obj
161 .get::<&Object>("desc")
162 .with_context(|| format!("reading desc from {}", obj.type_name))?;
163 let mut level = desc
164 .get::<u32>("level")
165 .with_context(|| format!("reading level from {}", obj.type_name))?;
166 let mut target = desc
167 .get::<&str>("target")
168 .with_context(|| format!("reading target from {}", obj.type_name))?;
169 let properties = obj
170 .get::<&Object>("properties")
171 .with_context(|| format!("reading properties from {}", obj.type_name))?;
172 for &(prop_name, prop_value) in properties.members {
173 match (prop_name, prop_value) {
174 ("target", Value::String(value_str)) => {
175 target = value_str;
176 }
177 ("level", Value::String(level_str)) => {
178 level = Level::parse(level_str).with_context(|| "parsing log level")?
179 as u32;
180 }
181 (_, _) => {}
182 }
183 }
184 Ok(Some(LogEntry {
185 process,
186 stream_id,
187 block_id,
188 insert_time: block_insert_time_ns,
189 time: convert_ticks.ticks_to_nanoseconds(ticks),
190 level: level as i32,
191 target,
192 msg,
193 properties: properties.into(),
194 }))
195 }
196
197 _ => {
198 warn!("unknown log event {:?}", obj);
199 Ok(None)
200 }
201 }
202 } else {
203 Ok(None)
204 }
205}
206
207#[span_fn]
209pub async fn for_each_log_entry_in_block<Predicate>(
210 blob_storage: Arc<BlobStorage>,
211 convert_ticks: &ConvertTicks,
212 process: Arc<ProcessMetadata>,
213 stream: &StreamMetadata,
214 block: &BlockMetadata,
215 mut fun: Predicate,
216) -> Result<bool>
217where
218 Predicate: for<'a> FnMut(LogEntry<'a>) -> Result<bool>,
219{
220 let payload = fetch_block_payload(
221 blob_storage,
222 stream.process_id,
223 stream.stream_id,
224 block.block_id,
225 )
226 .await?;
227 let stream_id = Arc::new(stream.stream_id.to_string());
228 let block_id = Arc::new(block.block_id.to_string());
229 let block_insert_time_ns = block.insert_time.timestamp_nanos_opt().unwrap_or_default();
230 let continue_iterating = parse_block(stream, &payload, |val| {
231 if let Some(log_entry) = log_entry_from_value(
232 convert_ticks,
233 process.clone(),
234 stream_id.clone(),
235 block_id.clone(),
236 block_insert_time_ns,
237 val,
238 )
239 .with_context(|| "log_entry_from_value")?
240 && !fun(log_entry)?
241 {
242 return Ok(false); }
244 Ok(true) })
246 .with_context(|| format!("parse_block {}", block.block_id))?;
247 Ok(continue_iterating)
248}