1pub mod attributes;
16pub mod span;
17pub mod v0;
18pub mod v1;
19
20use std::collections::HashSet;
21
22use api::v1::RowInsertRequests;
23pub use common_catalog::consts::{
26 DURATION_NANO_COLUMN, SERVICE_NAME_COLUMN, SPAN_KIND_COLUMN,
27 SPAN_STATUS_CODE_COLUMN as SPAN_STATUS_CODE, TRACE_TIMESTAMP_COLUMN as TIMESTAMP_COLUMN,
28};
29pub use common_catalog::consts::{
30 PARENT_SPAN_ID_COLUMN, SPAN_ID_COLUMN, SPAN_NAME_COLUMN, TRACE_ID_COLUMN,
31};
32use pipeline::{GreptimePipelineParams, PipelineWay};
33use session::context::QueryContextRef;
34
35use crate::error::{NotSupportedSnafu, Result};
36use crate::otlp::trace::span::TraceSpan;
37use crate::query_handler::PipelineHandlerRef;
38
39pub const SPAN_STATUS_MESSAGE_COLUMN: &str = "span_status_message";
40pub const SPAN_ATTRIBUTES_COLUMN: &str = "span_attributes";
41pub const SPAN_EVENTS_COLUMN: &str = "span_events";
42pub const SCOPE_NAME_COLUMN: &str = "scope_name";
43pub const SCOPE_VERSION_COLUMN: &str = "scope_version";
44pub const RESOURCE_ATTRIBUTES_COLUMN: &str = "resource_attributes";
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 pipeline_params: &GreptimePipelineParams,
116 table_name: &str,
117 query_ctx: &QueryContextRef,
118 pipeline_handler: PipelineHandlerRef,
119) -> Result<(RowInsertRequests, usize)> {
120 match pipeline {
121 PipelineWay::OtlpTraceDirectV0 => v0::v0_to_grpc_main_insert_requests(
122 spans,
123 pipeline,
124 pipeline_params,
125 table_name,
126 query_ctx,
127 pipeline_handler,
128 ),
129 PipelineWay::OtlpTraceDirectV1 => v1::v1_to_grpc_main_insert_requests(
130 spans,
131 pipeline,
132 pipeline_params,
133 table_name,
134 query_ctx,
135 pipeline_handler,
136 ),
137 _ => NotSupportedSnafu {
138 feat: "Unsupported pipeline for trace",
139 }
140 .fail(),
141 }
142}
143
144pub fn to_grpc_insert_requests_for_aux_tables(
150 aux_data: TraceAuxData,
151 pipeline: &PipelineWay,
152 table_name: &str,
153) -> Result<(RowInsertRequests, usize)> {
154 match pipeline {
155 PipelineWay::OtlpTraceDirectV0 => v0::build_aux_table_requests(aux_data, table_name),
156 PipelineWay::OtlpTraceDirectV1 => v1::build_aux_table_requests(aux_data, table_name),
157 _ => NotSupportedSnafu {
158 feat: "Unsupported pipeline for trace",
159 }
160 .fail(),
161 }
162}