1pub mod attributes;
16pub mod span;
17pub mod v0;
18pub mod v1;
19mod v2;
20
21use std::collections::HashSet;
22
23use api::v1::RowInsertRequests;
24pub use common_catalog::consts::{
27 DURATION_NANO_COLUMN, SERVICE_NAME_COLUMN, SPAN_KIND_COLUMN,
28 SPAN_STATUS_CODE_COLUMN as SPAN_STATUS_CODE, TRACE_TIMESTAMP_COLUMN as TIMESTAMP_COLUMN,
29};
30pub use common_catalog::consts::{
31 PARENT_SPAN_ID_COLUMN, RESOURCE_ATTRIBUTES_COLUMN, SCOPE_ATTRIBUTES_COLUMN,
32 SPAN_ATTRIBUTES_COLUMN, SPAN_ID_COLUMN, SPAN_NAME_COLUMN, TRACE_ID_COLUMN,
33};
34use pipeline::PipelineWay;
35
36use crate::error::{NotSupportedSnafu, Result};
37use crate::otlp::trace::span::TraceSpan;
38
39pub const TIMESTAMP_END_COLUMN: &str = "timestamp_end";
40pub const SPAN_STATUS_MESSAGE_COLUMN: &str = "span_status_message";
41pub const SPAN_EVENTS_COLUMN: &str = "span_events";
42pub const SPAN_LINKS_COLUMN: &str = "span_links";
43pub const SCOPE_NAME_COLUMN: &str = "scope_name";
44pub const SCOPE_VERSION_COLUMN: &str = "scope_version";
45pub const TRACE_STATE_COLUMN: &str = "trace_state";
46
47pub const KEY_SERVICE_NAME: &str = "service.name";
49pub const KEY_SERVICE_NAMESPACE: &str = "service.namespace";
50pub const KEY_SERVICE_INSTANCE_ID: &str = "service.instance.id";
51pub const KEY_HOST_ID: &str = "host.id";
52pub const KEY_HOST_NAME: &str = "host.name";
53pub const KEY_CONTAINER_ID: &str = "container.id";
54pub const KEY_CONTAINER_NAME: &str = "container.name";
55pub const KEY_K8S_POD_UID: &str = "k8s.pod.uid";
56pub const KEY_K8S_POD_NAME: &str = "k8s.pod.name";
57pub const KEY_K8S_CONTAINER_NAME: &str = "k8s.container.name";
58pub const KEY_K8S_NAMESPACE_NAME: &str = "k8s.namespace.name";
59pub const KEY_K8S_NODE_NAME: &str = "k8s.node.name";
60pub const KEY_SPAN_KIND: &str = "span.kind";
61
62pub const KEY_OTEL_SCOPE_NAME: &str = "otel.scope.name";
64pub const KEY_OTEL_SCOPE_VERSION: &str = "otel.scope.version";
65pub const KEY_OTEL_STATUS_CODE: &str = "otel.status_code";
66pub const KEY_OTEL_STATUS_MESSAGE: &str = "otel.status_description";
67pub const KEY_OTEL_STATUS_ERROR_KEY: &str = "error";
68pub const KEY_OTEL_TRACE_STATE: &str = "w3c.tracestate";
69
70pub const SPAN_KIND_PREFIX: &str = "SPAN_KIND_";
73
74pub const SPAN_STATUS_PREFIX: &str = "STATUS_CODE_";
76pub const SPAN_STATUS_UNSET: &str = "STATUS_CODE_UNSET";
77pub use common_catalog::consts::SPAN_STATUS_ERROR;
78
79#[derive(Debug, Default)]
86pub struct TraceAuxData {
87 pub services: HashSet<String>,
88 pub operations: HashSet<(String, String, String)>,
89}
90
91impl TraceAuxData {
92 pub fn observe_span(&mut self, span: &TraceSpan) {
95 if let Some(service_name) = &span.service_name {
96 self.services.insert(service_name.clone());
97 self.operations.insert((
98 service_name.clone(),
99 span.span_name.clone(),
100 span.span_kind.clone(),
101 ));
102 }
103 }
104
105 pub fn is_empty(&self) -> bool {
107 self.services.is_empty() && self.operations.is_empty()
108 }
109}
110
111pub fn to_grpc_insert_requests_from_spans(
113 spans: &[TraceSpan],
114 pipeline: &PipelineWay,
115 table_name: &str,
116) -> Result<(RowInsertRequests, usize)> {
117 match pipeline {
118 PipelineWay::OtlpTraceDirectV0 => v0::v0_to_grpc_main_insert_requests(spans, table_name),
119 PipelineWay::OtlpTraceDirectV1 => v1::v1_to_grpc_main_insert_requests(spans, table_name),
120 PipelineWay::OtlpTraceDirectV2 => v2::v2_to_grpc_main_insert_requests(spans, table_name),
121 _ => NotSupportedSnafu {
122 feat: "Unsupported pipeline for trace",
123 }
124 .fail(),
125 }
126}
127
128pub fn to_grpc_insert_requests_for_aux_tables(
134 aux_data: TraceAuxData,
135 pipeline: &PipelineWay,
136 table_name: &str,
137) -> Result<(RowInsertRequests, usize)> {
138 match pipeline {
139 PipelineWay::OtlpTraceDirectV0 => v0::build_aux_table_requests(aux_data, table_name),
140 PipelineWay::OtlpTraceDirectV1 | PipelineWay::OtlpTraceDirectV2 => {
141 v1::build_aux_table_requests(aux_data, table_name)
142 }
143 _ => NotSupportedSnafu {
144 feat: "Unsupported pipeline for trace",
145 }
146 .fail(),
147 }
148}