1use std::collections::HashSet;
16
17use api::v1::value::ValueData;
18use api::v1::{ColumnDataType, RowInsertRequests};
19use common_catalog::consts::{trace_operations_table_name, trace_services_table_name};
20use common_grpc::precision::Precision;
21use pipeline::{GreptimePipelineParams, PipelineWay};
22use session::context::QueryContextRef;
23
24use crate::error::Result;
25use crate::otlp::trace::span::TraceSpan;
26use crate::otlp::trace::{
27 DURATION_NANO_COLUMN, PARENT_SPAN_ID_COLUMN, SERVICE_NAME_COLUMN, SPAN_ATTRIBUTES_COLUMN,
28 SPAN_EVENTS_COLUMN, SPAN_ID_COLUMN, SPAN_KIND_COLUMN, SPAN_NAME_COLUMN, SPAN_STATUS_CODE,
29 SPAN_STATUS_MESSAGE_COLUMN, TIMESTAMP_COLUMN, TRACE_ID_COLUMN, TRACE_STATE_COLUMN,
30 TraceAuxData,
31};
32use crate::otlp::utils::{make_column_data, make_string_column_data};
33use crate::query_handler::PipelineHandlerRef;
34use crate::row_writer::{self, MultiTableData, TableData};
35
36const APPROXIMATE_COLUMN_COUNT: usize = 24;
37
38const MAX_TIMESTAMP: i64 = 4102444800000000000;
40
41pub fn v0_to_grpc_main_insert_requests(
46 spans: &[TraceSpan],
47 _pipeline: &PipelineWay,
48 _pipeline_params: &GreptimePipelineParams,
49 table_name: &str,
50 _query_ctx: &QueryContextRef,
51 _pipeline_handler: PipelineHandlerRef,
52) -> Result<(RowInsertRequests, usize)> {
53 let mut multi_table_writer = MultiTableData::default();
54 let trace_writer = build_trace_table_data(spans)?;
55 multi_table_writer.add_table_data(table_name, trace_writer);
56
57 Ok(multi_table_writer.into_row_insert_requests())
58}
59
60pub fn build_trace_table_data(spans: &[TraceSpan]) -> Result<TableData> {
62 let mut trace_writer = TableData::new(APPROXIMATE_COLUMN_COUNT, spans.len());
63 for span in spans.iter().cloned() {
64 write_span_to_row(&mut trace_writer, span)?;
65 }
66
67 Ok(trace_writer)
68}
69
70pub fn build_aux_table_requests(
72 aux_data: TraceAuxData,
73 table_name: &str,
74) -> Result<(RowInsertRequests, usize)> {
75 let mut multi_table_writer = MultiTableData::default();
76 let mut trace_services_writer = TableData::new(APPROXIMATE_COLUMN_COUNT, 1);
77 let mut trace_operations_writer = TableData::new(APPROXIMATE_COLUMN_COUNT, 1);
78
79 write_trace_services_to_row(&mut trace_services_writer, aux_data.services)?;
80 write_trace_operations_to_row(&mut trace_operations_writer, aux_data.operations)?;
81
82 multi_table_writer.add_table_data(trace_services_table_name(table_name), trace_services_writer);
83 multi_table_writer.add_table_data(
84 trace_operations_table_name(table_name),
85 trace_operations_writer,
86 );
87 Ok(multi_table_writer.into_row_insert_requests())
88}
89
90pub fn write_span_to_row(writer: &mut TableData, span: TraceSpan) -> Result<()> {
91 let mut row = writer.alloc_one_row();
92
93 row_writer::write_ts_to_nanos(
95 writer,
96 TIMESTAMP_COLUMN,
97 Some(span.start_in_nanosecond as i64),
98 Precision::Nanosecond,
99 &mut row,
100 )?;
101
102 let fields = vec![
104 make_column_data(
105 "timestamp_end",
106 ColumnDataType::TimestampNanosecond,
107 Some(ValueData::TimestampNanosecondValue(
108 span.end_in_nanosecond as i64,
109 )),
110 ),
111 make_column_data(
117 DURATION_NANO_COLUMN,
118 ColumnDataType::Uint64,
119 Some(ValueData::U64Value(
120 span.end_in_nanosecond - span.start_in_nanosecond,
121 )),
122 ),
123 make_string_column_data(TRACE_ID_COLUMN, Some(span.trace_id)),
124 make_string_column_data(SPAN_ID_COLUMN, Some(span.span_id)),
125 make_string_column_data(PARENT_SPAN_ID_COLUMN, span.parent_span_id),
126 make_string_column_data(SPAN_KIND_COLUMN, Some(span.span_kind)),
127 make_string_column_data(SPAN_NAME_COLUMN, Some(span.span_name)),
128 make_string_column_data(SPAN_STATUS_CODE, Some(span.span_status_code)),
129 make_string_column_data(SPAN_STATUS_MESSAGE_COLUMN, Some(span.span_status_message)),
130 make_string_column_data(TRACE_STATE_COLUMN, Some(span.trace_state)),
131 ];
132 row_writer::write_fields(writer, fields.into_iter(), &mut row)?;
133
134 if let Some(service_name) = span.service_name {
135 row_writer::write_tag(writer, SERVICE_NAME_COLUMN, service_name, &mut row)?;
136 }
137
138 row_writer::write_json(
139 writer,
140 SPAN_ATTRIBUTES_COLUMN,
141 span.span_attributes.into(),
142 &mut row,
143 )?;
144 row_writer::write_json(
145 writer,
146 SPAN_EVENTS_COLUMN,
147 span.span_events.into(),
148 &mut row,
149 )?;
150 row_writer::write_json(writer, "span_links", span.span_links.into(), &mut row)?;
151
152 let fields = vec![
154 make_string_column_data("scope_name", Some(span.scope_name)),
155 make_string_column_data("scope_version", Some(span.scope_version)),
156 ];
157 row_writer::write_fields(writer, fields.into_iter(), &mut row)?;
158
159 row_writer::write_json(
160 writer,
161 "scope_attributes",
162 span.scope_attributes.into(),
163 &mut row,
164 )?;
165
166 row_writer::write_json(
167 writer,
168 "resource_attributes",
169 span.resource_attributes.into(),
170 &mut row,
171 )?;
172
173 writer.add_row(row);
174
175 Ok(())
176}
177
178fn write_trace_services_to_row(writer: &mut TableData, services: HashSet<String>) -> Result<()> {
179 for service_name in services {
180 let mut row = writer.alloc_one_row();
181 row_writer::write_ts_to_nanos(
183 writer,
184 TIMESTAMP_COLUMN,
185 Some(MAX_TIMESTAMP),
186 Precision::Nanosecond,
187 &mut row,
188 )?;
189
190 row_writer::write_fields(
192 writer,
193 std::iter::once(make_string_column_data(
194 SERVICE_NAME_COLUMN,
195 Some(service_name),
196 )),
197 &mut row,
198 )?;
199 writer.add_row(row);
200 }
201
202 Ok(())
203}
204
205fn write_trace_operations_to_row(
206 writer: &mut TableData,
207 operations: HashSet<(String, String, String)>,
208) -> Result<()> {
209 for (service_name, span_name, span_kind) in operations {
210 let mut row = writer.alloc_one_row();
211 row_writer::write_ts_to_nanos(
213 writer,
214 TIMESTAMP_COLUMN,
215 Some(MAX_TIMESTAMP),
216 Precision::Nanosecond,
217 &mut row,
218 )?;
219
220 row_writer::write_fields(
222 writer,
223 vec![
224 make_string_column_data(SERVICE_NAME_COLUMN, Some(service_name)),
225 make_string_column_data(SPAN_NAME_COLUMN, Some(span_name)),
226 make_string_column_data(SPAN_KIND_COLUMN, Some(span_kind)),
227 ]
228 .into_iter(),
229 &mut row,
230 )?;
231 writer.add_row(row);
232 }
233
234 Ok(())
235}
236
237#[cfg(test)]
238mod tests {
239 use api::v1::ColumnDataType;
240 use api::v1::value::ValueData;
241
242 use super::{build_aux_table_requests, build_trace_table_data};
243 use crate::otlp::trace::attributes::Attributes;
244 use crate::otlp::trace::span::{SpanEvents, SpanLinks, TraceSpan};
245 use crate::otlp::trace::{DURATION_NANO_COLUMN, TraceAuxData};
246
247 fn make_span(service_name: &str, trace_id: &str, span_id: &str) -> TraceSpan {
248 TraceSpan {
249 service_name: Some(service_name.to_string()),
250 trace_id: trace_id.to_string(),
251 span_id: span_id.to_string(),
252 parent_span_id: None,
253 resource_attributes: Attributes::from(vec![]),
254 scope_name: "scope".to_string(),
255 scope_version: "v1".to_string(),
256 scope_attributes: Attributes::from(vec![]),
257 trace_state: String::new(),
258 span_name: "op".to_string(),
259 span_kind: "SPAN_KIND_SERVER".to_string(),
260 span_status_code: "STATUS_CODE_UNSET".to_string(),
261 span_status_message: String::new(),
262 span_attributes: Attributes::from(vec![]),
263 span_events: SpanEvents::from(vec![]),
264 span_links: SpanLinks::from(vec![]),
265 start_in_nanosecond: 1,
266 end_in_nanosecond: 2,
267 }
268 }
269
270 #[test]
271 fn test_build_trace_table_data_from_span_subset() {
272 let spans = [
273 make_span("svc-a", "trace-a", "span-a"),
274 make_span("svc-b", "trace-b", "span-b"),
275 ];
276
277 let writer = build_trace_table_data(&spans[..1]).unwrap();
278 let (_, rows) = writer.into_schema_and_rows();
279 assert_eq!(rows.len(), 1);
280 }
281
282 #[test]
283 fn test_v0_duration_nano_stays_uint64() {
284 let writer = build_trace_table_data(&[make_span("svc-a", "trace-a", "span-a")]).unwrap();
289 let (schema, rows) = writer.into_schema_and_rows();
290
291 let idx = schema
292 .iter()
293 .position(|c| c.column_name == DURATION_NANO_COLUMN)
294 .unwrap();
295 assert_eq!(schema[idx].datatype, ColumnDataType::Uint64 as i32);
296 assert_eq!(rows[0].values[idx].value_data, Some(ValueData::U64Value(1)));
297 }
298
299 #[test]
300 fn test_build_aux_table_requests_deduplicates_services_and_operations() {
301 let spans = vec![
302 make_span("svc-a", "trace-a", "span-a"),
303 make_span("svc-a", "trace-b", "span-b"),
304 ];
305 let mut aux_data = TraceAuxData::default();
306 for span in &spans {
307 aux_data.observe_span(span);
308 }
309
310 let (requests, total_rows) =
311 build_aux_table_requests(aux_data, "opentelemetry_traces").unwrap();
312 assert_eq!(requests.inserts.len(), 2);
313 assert_eq!(total_rows, 2);
314 }
315}