1use std::fmt::Display;
16
17use common_time::timestamp::Timestamp;
18use itertools::Itertools;
19use opentelemetry_proto::tonic::collector::trace::v1::ExportTraceServiceRequest;
20use opentelemetry_proto::tonic::common::v1::{InstrumentationScope, any_value};
21use opentelemetry_proto::tonic::trace::v1::span::{Event, Link};
22use opentelemetry_proto::tonic::trace::v1::{Span, Status};
23use serde::Serialize;
24
25use crate::otlp::trace::KEY_SERVICE_NAME;
26use crate::otlp::trace::attributes::{Attributes, SharedAttributes};
27use crate::otlp::utils::bytes_to_hex_string;
28
29#[derive(Debug, Clone)]
30pub struct TraceSpan {
31 pub service_name: Option<String>,
33 pub trace_id: String,
34 pub span_id: String,
35 pub parent_span_id: Option<String>,
36
37 pub resource_attributes: SharedAttributes,
39 pub scope_name: String,
40 pub scope_version: String,
41 pub scope_attributes: SharedAttributes,
42 pub trace_state: String,
43 pub span_name: String,
44 pub span_kind: String,
45 pub span_status_code: String,
46 pub span_status_message: String,
47 pub span_attributes: Attributes,
48 pub span_events: SpanEvents, pub span_links: SpanLinks, pub start_in_nanosecond: u64, pub end_in_nanosecond: u64,
52}
53
54pub type TraceSpans = Vec<TraceSpan>;
55
56#[derive(Debug, Clone)]
57pub struct TraceSpanGroup {
58 pub service_name: Option<String>,
59 pub resource_attributes: SharedAttributes,
60 pub scope_name: String,
61 pub scope_version: String,
62 pub scope_attributes: SharedAttributes,
63 pub spans: TraceSpans,
64}
65
66pub type TraceSpanGroups = Vec<TraceSpanGroup>;
67
68#[derive(Debug, Clone, Serialize)]
69pub struct SpanLink {
70 pub trace_id: String,
71 pub span_id: String,
72 pub trace_state: String,
73 pub attributes: Attributes, }
75
76impl From<Link> for SpanLink {
77 fn from(link: Link) -> Self {
78 Self {
79 trace_id: bytes_to_hex_string(&link.trace_id),
80 span_id: bytes_to_hex_string(&link.span_id),
81 trace_state: link.trace_state,
82 attributes: Attributes::from(link.attributes),
83 }
84 }
85}
86
87impl From<SpanLink> for jsonb::Value<'static> {
88 fn from(value: SpanLink) -> jsonb::Value<'static> {
89 jsonb::Value::Object(
90 vec![
91 (
92 "trace_id".to_string(),
93 jsonb::Value::String(value.trace_id.into()),
94 ),
95 (
96 "span_id".to_string(),
97 jsonb::Value::String(value.span_id.into()),
98 ),
99 (
100 "trace_state".to_string(),
101 jsonb::Value::String(value.trace_state.into()),
102 ),
103 ("attributes".to_string(), value.attributes.into()),
104 ]
105 .into_iter()
106 .collect(),
107 )
108 }
109}
110
111#[derive(Debug, Clone, Serialize)]
112pub struct SpanLinks(Vec<SpanLink>);
113
114impl From<Vec<Link>> for SpanLinks {
115 fn from(value: Vec<Link>) -> Self {
116 let links = value.into_iter().map(SpanLink::from).collect_vec();
117 Self(links)
118 }
119}
120
121impl From<SpanLinks> for jsonb::Value<'static> {
122 fn from(value: SpanLinks) -> jsonb::Value<'static> {
123 jsonb::Value::Array(value.0.into_iter().map(Into::into).collect())
124 }
125}
126
127impl Display for SpanLinks {
128 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
129 write!(f, "{}", serde_json::to_string(self).unwrap_or_default())
130 }
131}
132
133impl SpanLinks {
134 pub fn get_ref(&self) -> &Vec<SpanLink> {
135 &self.0
136 }
137
138 pub fn get_mut(&mut self) -> &mut Vec<SpanLink> {
139 &mut self.0
140 }
141}
142
143#[derive(Debug, Clone, Serialize)]
144pub struct SpanEvent {
145 pub name: String,
146 pub time: String,
147 pub attributes: Attributes, }
149
150impl From<Event> for SpanEvent {
151 fn from(event: Event) -> Self {
152 Self {
153 name: event.name,
154 time: Timestamp::new_nanosecond(event.time_unix_nano as i64).to_iso8601_string(),
155 attributes: Attributes::from(event.attributes),
156 }
157 }
158}
159
160impl From<SpanEvent> for jsonb::Value<'static> {
161 fn from(value: SpanEvent) -> jsonb::Value<'static> {
162 jsonb::Value::Object(
163 vec![
164 ("name".to_string(), jsonb::Value::String(value.name.into())),
165 ("time".to_string(), jsonb::Value::String(value.time.into())),
166 ("attributes".to_string(), value.attributes.into()),
167 ]
168 .into_iter()
169 .collect(),
170 )
171 }
172}
173
174#[derive(Debug, Clone, Serialize)]
175pub struct SpanEvents(Vec<SpanEvent>);
176
177impl From<Vec<Event>> for SpanEvents {
178 fn from(value: Vec<Event>) -> Self {
179 let events = value.into_iter().map(SpanEvent::from).collect_vec();
180 Self(events)
181 }
182}
183
184impl From<SpanEvents> for jsonb::Value<'static> {
185 fn from(value: SpanEvents) -> jsonb::Value<'static> {
186 jsonb::Value::Array(value.0.into_iter().map(Into::into).collect())
187 }
188}
189
190impl Display for SpanEvents {
191 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
192 write!(f, "{}", serde_json::to_string(self).unwrap_or_default())
193 }
194}
195
196impl SpanEvents {
197 pub fn get_ref(&self) -> &Vec<SpanEvent> {
198 &self.0
199 }
200
201 pub fn get_mut(&mut self) -> &mut Vec<SpanEvent> {
202 &mut self.0
203 }
204}
205
206pub fn parse_span(
207 service_name: Option<String>,
208 resource_attrs: &SharedAttributes,
209 scope_name: &str,
210 scope_version: &str,
211 scope_attrs: &SharedAttributes,
212 span: Span,
213) -> TraceSpan {
214 let (span_status_code, span_status_message) = status_to_string(&span.status);
215 let span_kind = span.kind().as_str_name().into();
216 TraceSpan {
217 service_name,
218 trace_id: bytes_to_hex_string(&span.trace_id),
219 span_id: bytes_to_hex_string(&span.span_id),
220 parent_span_id: if span.parent_span_id.is_empty() {
221 None
222 } else {
223 Some(bytes_to_hex_string(&span.parent_span_id))
224 },
225
226 resource_attributes: resource_attrs.clone(),
227 trace_state: span.trace_state,
228
229 scope_name: scope_name.to_string(),
230 scope_version: scope_version.to_string(),
231 scope_attributes: scope_attrs.clone(),
232
233 span_name: span.name,
234 span_kind,
235 span_status_code,
236 span_status_message,
237 span_attributes: Attributes::from(span.attributes),
238 span_events: SpanEvents::from(span.events),
239 span_links: SpanLinks::from(span.links),
240
241 start_in_nanosecond: span.start_time_unix_nano,
242 end_in_nanosecond: span.end_time_unix_nano,
243 }
244}
245
246pub fn status_to_string(status: &Option<Status>) -> (String, String) {
247 match status {
248 Some(status) => (status.code().as_str_name().into(), status.message.clone()),
249 None => ("".into(), "".into()),
250 }
251}
252
253pub fn parse(request: ExportTraceServiceRequest) -> TraceSpanGroups {
259 let group_size = request
260 .resource_spans
261 .iter()
262 .flat_map(|res| res.scope_spans.iter())
263 .count();
264 let mut groups = Vec::with_capacity(group_size);
265 for resource_spans in request.resource_spans {
266 let resource_attrs = resource_spans
267 .resource
268 .map(|r| r.attributes)
269 .unwrap_or_default();
270
271 let service_name = resource_attrs
276 .iter()
277 .find(|kv| kv.key == KEY_SERVICE_NAME)
278 .and_then(|kv| kv.value.clone())
279 .and_then(|v| match v.value {
280 Some(any_value::Value::StringValue(s)) => Some(s),
281 Some(any_value::Value::BytesValue(b)) => {
282 Some(String::from_utf8_lossy(&b).to_string())
283 }
284 _ => None,
285 });
286
287 let resource_attrs = SharedAttributes::from(Attributes::from(resource_attrs));
288 for scope_spans in resource_spans.scope_spans {
289 let InstrumentationScope {
290 name: scope_name,
291 version: scope_version,
292 attributes: scope_attrs,
293 ..
294 } = scope_spans.scope.unwrap_or_default();
295 let scope_attrs = SharedAttributes::from(Attributes::from(scope_attrs));
296 let mut spans = Vec::with_capacity(scope_spans.spans.len());
297 for span in scope_spans.spans {
298 spans.push(parse_span(
299 service_name.clone(),
300 &resource_attrs,
301 &scope_name,
302 &scope_version,
303 &scope_attrs,
304 span,
305 ));
306 }
307 groups.push(TraceSpanGroup {
308 service_name: service_name.clone(),
309 resource_attributes: resource_attrs.clone(),
310 scope_name,
311 scope_version,
312 scope_attributes: scope_attrs,
313 spans,
314 });
315 }
316 }
317 groups
318}
319
320#[cfg(test)]
321mod tests {
322 use opentelemetry_proto::tonic::collector::trace::v1::ExportTraceServiceRequest;
323 use opentelemetry_proto::tonic::common::v1::{
324 AnyValue, InstrumentationScope, KeyValue, any_value,
325 };
326 use opentelemetry_proto::tonic::resource::v1::Resource;
327 use opentelemetry_proto::tonic::trace::v1::{ResourceSpans, ScopeSpans, Span, Status};
328
329 use crate::otlp::trace::KEY_SERVICE_NAME;
330 use crate::otlp::trace::span::{bytes_to_hex_string, parse, status_to_string};
331
332 fn make_kv(key: &str, value: &str) -> KeyValue {
333 KeyValue {
334 key: key.to_string(),
335 value: Some(AnyValue {
336 value: Some(any_value::Value::StringValue(value.to_string())),
337 }),
338 ..Default::default()
339 }
340 }
341
342 fn make_span(trace_id: u8, span_id: u8) -> Span {
343 Span {
344 trace_id: vec![trace_id; 16],
345 span_id: vec![span_id; 8],
346 ..Default::default()
347 }
348 }
349
350 #[test]
351 fn test_bytes_to_hex_string() {
352 assert_eq!(
353 "24fe79948641b110a29bc27859307e8d",
354 bytes_to_hex_string(&[
355 36, 254, 121, 148, 134, 65, 177, 16, 162, 155, 194, 120, 89, 48, 126, 141,
356 ])
357 );
358
359 assert_eq!(
360 "baffeedd7b8debc0",
361 bytes_to_hex_string(&[186, 255, 238, 221, 123, 141, 235, 192,])
362 );
363 }
364
365 #[test]
366 fn test_status_to_string() {
367 let message = String::from("status message");
368 let status = Status {
369 code: 1,
370 message: message.clone(),
371 };
372
373 assert_eq!(
374 ("STATUS_CODE_OK".into(), message),
375 status_to_string(&Some(status)),
376 );
377 }
378
379 #[test]
380 fn test_parse_preserves_resource_scope_groups() {
381 let request = ExportTraceServiceRequest {
382 resource_spans: vec![
383 ResourceSpans {
384 resource: Some(Resource {
385 attributes: vec![make_kv(KEY_SERVICE_NAME, "svc-a")],
386 ..Default::default()
387 }),
388 scope_spans: vec![
389 ScopeSpans {
390 scope: Some(InstrumentationScope {
391 name: "scope-1".to_string(),
392 attributes: vec![make_kv("scope.key", "scope-1-value")],
393 ..Default::default()
394 }),
395 spans: vec![make_span(0x11, 0x21), make_span(0x12, 0x22)],
396 ..Default::default()
397 },
398 ScopeSpans {
399 scope: Some(InstrumentationScope {
400 name: "scope-2".to_string(),
401 attributes: vec![make_kv("scope.key", "scope-2-value")],
402 ..Default::default()
403 }),
404 spans: vec![make_span(0x13, 0x23)],
405 ..Default::default()
406 },
407 ],
408 ..Default::default()
409 },
410 ResourceSpans {
411 resource: Some(Resource {
412 attributes: vec![make_kv(KEY_SERVICE_NAME, "svc-b")],
413 ..Default::default()
414 }),
415 scope_spans: vec![ScopeSpans {
416 scope: Some(InstrumentationScope {
417 name: "scope-3".to_string(),
418 ..Default::default()
419 }),
420 spans: vec![make_span(0x14, 0x24)],
421 ..Default::default()
422 }],
423 ..Default::default()
424 },
425 ],
426 };
427
428 let groups = parse(request);
429 assert_eq!(groups.len(), 3);
430 assert_eq!(groups[0].service_name.as_deref(), Some("svc-a"));
431 assert_eq!(groups[0].scope_name, "scope-1");
432 assert_eq!(groups[0].spans.len(), 2);
433 assert_eq!(groups[1].scope_name, "scope-2");
434 assert_eq!(groups[1].spans.len(), 1);
435 assert_eq!(groups[2].service_name.as_deref(), Some("svc-b"));
436 assert_eq!(groups[2].scope_name, "scope-3");
437 assert_eq!(
438 groups[0].scope_attributes.as_ref().get_ref(),
439 &[make_kv("scope.key", "scope-1-value")]
440 );
441
442 assert!(std::ptr::eq(
446 groups[0].resource_attributes.as_ref(),
447 groups[1].resource_attributes.as_ref(),
448 ));
449 assert!(!std::ptr::eq(
450 groups[0].resource_attributes.as_ref(),
451 groups[2].resource_attributes.as_ref(),
452 ));
453 assert!(!std::ptr::eq(
454 groups[0].scope_attributes.as_ref(),
455 groups[1].scope_attributes.as_ref(),
456 ));
457 for group in &groups {
458 for span in &group.spans {
459 assert!(std::ptr::eq(
460 group.resource_attributes.as_ref(),
461 span.resource_attributes.as_ref()
462 ));
463 assert!(std::ptr::eq(
464 group.scope_attributes.as_ref(),
465 span.scope_attributes.as_ref()
466 ));
467 }
468 }
469 }
470}