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