Skip to main content

servers/otlp/trace/
span.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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    // the following are tags
32    pub service_name: Option<String>,
33    pub trace_id: String,
34    pub span_id: String,
35    pub parent_span_id: Option<String>,
36
37    // the following are fields
38    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,  // TODO(yuanbohan): List in the future
49    pub span_links: SpanLinks,    // TODO(yuanbohan): List in the future
50    pub start_in_nanosecond: u64, // this is also the Timestamp Index
51    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, // TODO(yuanbohan): Map in the future
74}
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, // TODO(yuanbohan): Map in the future
148}
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
253/// Convert OpenTelemetry traces to SpanTraces
254///
255/// See
256/// <https://github.com/open-telemetry/opentelemetry-proto/blob/main/opentelemetry/proto/trace/v1/trace.proto>
257/// for data structure of OTLP traces.
258pub 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        // "service.name" is required; SDKs MUST provide a fallback when it is unset.
272        // This is in the specification:
273        // https://opentelemetry.io/docs/specs/semconv/resource/service/
274        // Tolerate missing names from senders without using unrelated attributes as names.
275        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        // Scopes of one resource share a single copy of the resource attributes,
443        // spans share their group's attributes, and unrelated resources and
444        // scopes stay separate.
445        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}