Skip to main content

zelos_trace/
segment.rs

1use std::collections::HashMap;
2
3use chrono::{DateTime, Utc};
4use rpds::HashTrieMapSync;
5use uuid::Uuid;
6use zelos_trace_types::{
7    PathSegment, Signal, SignalKey, Value,
8    ipc::{self, TraceEventFieldMetadata},
9};
10
11#[derive(Debug, Clone)]
12pub struct TraceEventSchemaRef<'a> {
13    pub segment: &'a TraceSegment,
14    pub event_schema: &'a TraceEventSchema,
15}
16
17#[derive(Debug, Clone)]
18pub struct TraceEventFieldRef<'a> {
19    pub segment: &'a TraceSegment,
20    pub event_schema: &'a TraceEventSchema,
21    pub field: &'a TraceEventField,
22}
23
24impl TraceEventFieldRef<'_> {
25    pub fn as_signal(&self) -> Signal {
26        let &TraceEventFieldRef {
27            segment,
28            event_schema,
29            field,
30        } = self;
31
32        Signal {
33            data_segment_id: segment.id,
34            source: segment.source.clone(),
35            message: event_schema.name.clone(),
36            signal: field.metadata.name.clone(),
37            data_type: field.metadata.data_type.clone(),
38            unit: field.metadata.unit.clone(),
39            value_table: if field.values.is_empty() {
40                None
41            } else {
42                Some(
43                    field
44                        .values
45                        .iter()
46                        .filter_map(|(k, v)| k.as_number().map(|n| (n, v.clone())))
47                        .collect(),
48                )
49            },
50        }
51    }
52}
53
54impl TraceEventFieldRef<'_> {
55    pub fn table_key(&self) -> String {
56        format!(
57            "{}/{}/{}",
58            self.segment.id, self.segment.source, self.event_schema.name
59        )
60    }
61}
62
63#[derive(Debug, Clone)]
64pub struct TraceEventField {
65    pub metadata: TraceEventFieldMetadata,
66    pub values: HashMap<Value, String>,
67}
68
69impl TraceEventField {
70    pub fn from_ipc(msg: ipc::TraceEventFieldMetadata) -> Self {
71        Self {
72            metadata: msg,
73            values: HashMap::new(),
74        }
75    }
76}
77
78#[derive(Clone, Debug)]
79pub struct TraceEventSchema {
80    pub name: String,
81    pub fields: Vec<TraceEventField>,
82}
83
84impl TraceEventSchema {
85    pub fn from_ipc(msg: ipc::TraceEventSchema) -> Self {
86        Self {
87            name: msg.name,
88            fields: msg
89                .fields
90                .into_iter()
91                .map(TraceEventField::from_ipc)
92                .collect(),
93        }
94    }
95
96    pub fn get_field(&self, field_name: &str) -> Option<&TraceEventField> {
97        self.fields
98            .iter()
99            .find(|field| field.metadata.name == field_name)
100    }
101
102    pub fn get_field_mut(&mut self, field_name: &str) -> Option<&mut TraceEventField> {
103        self.fields
104            .iter_mut()
105            .find(|field| field.metadata.name == field_name)
106    }
107
108    pub fn metadata(&self) -> impl Iterator<Item = &TraceEventFieldMetadata> {
109        self.fields.iter().map(|field| &field.metadata)
110    }
111}
112
113#[derive(Clone, Debug)]
114pub struct TraceSegment {
115    pub id: Uuid,
116    pub source: String,
117    pub start_time: Option<DateTime<Utc>>,
118    pub end_time: Option<DateTime<Utc>>,
119    pub schemas: HashTrieMapSync<String, TraceEventSchema>,
120}
121
122impl TraceSegment {
123    /// Create an empty trace segment when we have not received a start message
124    pub fn empty(id: Uuid, source_name: String) -> Self {
125        Self {
126            id,
127            source: source_name,
128            start_time: None,
129            end_time: None,
130            schemas: HashTrieMapSync::new_sync(),
131        }
132    }
133
134    /// Create a trace segment from a start message.
135    pub fn from_ipc(id: Uuid, start: &ipc::TraceSegmentStart) -> Self {
136        Self {
137            id,
138            source: start.source_name.clone(),
139            start_time: Some(DateTime::from_timestamp_nanos(start.time_ns)),
140            end_time: None,
141            schemas: HashTrieMapSync::new_sync(),
142        }
143    }
144
145    pub fn update_mut(&mut self, msg: &ipc::IpcMessage) {
146        match msg {
147            ipc::IpcMessage::TraceSegmentStart(m) => {
148                self.source = m.source_name.clone();
149
150                // Update the start time if it's earlier than the existing one
151                let start_time = DateTime::from_timestamp_nanos(m.time_ns);
152                if let Some(existing_start_time) = self.start_time {
153                    if start_time < existing_start_time {
154                        self.start_time = Some(start_time);
155                    }
156                } else {
157                    self.start_time = Some(start_time);
158                }
159            }
160            ipc::IpcMessage::TraceSegmentEnd(m) => {
161                self.end_time = Some(DateTime::from_timestamp_nanos(m.time_ns));
162            }
163            ipc::IpcMessage::TraceEventSchema(m) => {
164                if !self.schemas.contains_key(&m.name) {
165                    self.schemas
166                        .insert_mut(m.name.clone(), TraceEventSchema::from_ipc(m.clone()));
167                }
168            }
169            ipc::IpcMessage::TraceEventFieldNamedValues(m) => {
170                // Update our event schema in place
171                if let Some(mut event_schema) = self.schemas.get(&m.event_name).cloned() {
172                    if let Some(field) = event_schema.get_field_mut(&m.field_name) {
173                        field.values.extend(m.values.clone());
174                    }
175                    self.schemas.insert_mut(m.event_name.clone(), event_schema);
176                }
177            }
178            ipc::IpcMessage::TraceEvent(_m) => {
179                // Do nothing
180            }
181        }
182    }
183
184    pub fn update(&self, msg: &ipc::IpcMessage) -> Self {
185        let mut new = self.clone();
186        new.update_mut(msg);
187        new
188    }
189
190    pub fn maybe_event<'a>(&'a self, event_name: &str) -> Option<TraceEventSchemaRef<'a>> {
191        self.schemas
192            .get(event_name)
193            .map(|event_schema| TraceEventSchemaRef {
194                segment: self,
195                event_schema,
196            })
197    }
198
199    pub fn field_refs(&self) -> impl Iterator<Item = TraceEventFieldRef<'_>> {
200        self.schemas.values().flat_map(|event_schema| {
201            event_schema.fields.iter().map(|field| TraceEventFieldRef {
202                segment: self,
203                event_schema,
204                field,
205            })
206        })
207    }
208
209    pub fn field_refs_matching<'a>(
210        &'a self,
211        signal_keys: &[SignalKey],
212    ) -> impl Iterator<Item = TraceEventFieldRef<'a>> {
213        signal_keys
214            .iter()
215            .flat_map(|k| self.maybe_field_ref_matching(k))
216    }
217
218    pub fn maybe_field_ref_matching<'a>(
219        &'a self,
220        key: &SignalKey,
221    ) -> Option<TraceEventFieldRef<'a>> {
222        // Return early if this signal key is for a specific uuid and we don't match it
223        if let PathSegment::Uuid { uuid } = key.data_segment_id {
224            if uuid != self.id {
225                return None;
226            }
227        }
228
229        // Return early if this signal key is for a specific source and we don't match it
230        if key.source != self.source {
231            return None;
232        }
233
234        // Attempt to get the message, and map it to the column
235        self.schemas.get(&key.message).and_then(|event_schema| {
236            event_schema
237                .fields
238                .iter()
239                .find(|k| k.metadata.name == key.signal)
240                .map(|field| TraceEventFieldRef {
241                    segment: self,
242                    event_schema,
243                    field,
244                })
245        })
246    }
247
248    pub fn signals(&self) -> impl Iterator<Item = Signal> {
249        self.field_refs().map(|r| r.as_signal())
250    }
251
252    pub fn signals_matching(&self, signal_keys: &[SignalKey]) -> impl Iterator<Item = Signal> {
253        self.field_refs_matching(signal_keys).map(|r| r.as_signal())
254    }
255
256    /// Represent this trace segment as ipc messages
257    pub fn as_ipc(&self) -> Vec<ipc::IpcMessage> {
258        let mut msgs = Vec::new();
259
260        // Send start, if we have a start timestamp
261        // NOTE(jbott): this is somewhat weird, should we store trace segments if we don't have a timestamp? should we fixup timestamps from the data contained within?
262        if let Some(start_time_ns) = self.start_time.and_then(|t| t.timestamp_nanos_opt()) {
263            let start = ipc::TraceSegmentStart {
264                time_ns: start_time_ns,
265                source_name: self.source.clone(),
266            };
267            msgs.push(start.into());
268        }
269
270        // Iterate over all schemas and send messages as required
271        for (event_name, schema) in &self.schemas {
272            // Send the schema
273            let event_schema = ipc::TraceEventSchema {
274                name: schema.name.clone(),
275                fields: schema.fields.iter().map(|f| f.metadata.clone()).collect(),
276            };
277            msgs.push(event_schema.into());
278
279            // For each field with values, send the hashmap
280            for (field_name, values) in schema
281                .fields
282                .iter()
283                .filter(|f| !f.values.is_empty())
284                .map(|f| (f.metadata.name.clone(), f.values.clone()))
285            {
286                let event_field_named_values = ipc::TraceEventFieldNamedValues {
287                    event_name: event_name.clone(),
288                    field_name: field_name.clone(),
289                    values,
290                };
291                msgs.push(event_field_named_values.into());
292            }
293        }
294
295        // Send end, if we have an end timestamp
296        if let Some(end_time_ns) = self.end_time.and_then(|t| t.timestamp_nanos_opt()) {
297            let end = ipc::TraceSegmentEnd {
298                time_ns: end_time_ns,
299            };
300            msgs.push(end.into());
301        }
302
303        msgs
304    }
305}