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 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 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 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 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 }
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 if let PathSegment::Uuid { uuid } = key.data_segment_id {
224 if uuid != self.id {
225 return None;
226 }
227 }
228
229 if key.source != self.source {
231 return None;
232 }
233
234 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 pub fn as_ipc(&self) -> Vec<ipc::IpcMessage> {
258 let mut msgs = Vec::new();
259
260 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 for (event_name, schema) in &self.schemas {
272 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 (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 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}