1use std::collections::{HashMap, HashSet};
16
17use api::v1::column_data_type_extension::TypeExt;
18use api::v1::helper::time_index_column_schema;
19use api::v1::value::ValueData;
20use api::v1::{
21 ColumnDataType, ColumnDataTypeExtension, ColumnSchema, JsonTypeExtension, RowInsertRequests,
22 SemanticType, Value,
23};
24use common_catalog::consts::{trace_operations_table_name, trace_services_table_name};
25use common_grpc::precision::Precision;
26use opentelemetry_proto::tonic::common::v1::KeyValue;
27use opentelemetry_proto::tonic::common::v1::any_value::Value as OtlpValue;
28
29use crate::error::Result;
30use crate::otlp::trace::attributes::Attributes;
31use crate::otlp::trace::span::TraceSpan;
32use crate::otlp::trace::{
33 DURATION_NANO_COLUMN, KEY_SERVICE_NAME, PARENT_SPAN_ID_COLUMN, SCOPE_NAME_COLUMN,
34 SCOPE_VERSION_COLUMN, SERVICE_NAME_COLUMN, SPAN_EVENTS_COLUMN, SPAN_ID_COLUMN,
35 SPAN_KIND_COLUMN, SPAN_NAME_COLUMN, SPAN_STATUS_CODE, SPAN_STATUS_MESSAGE_COLUMN,
36 TIMESTAMP_COLUMN, TIMESTAMP_END_COLUMN, TRACE_ID_COLUMN, TRACE_STATE_COLUMN, TraceAuxData,
37};
38use crate::otlp::utils::any_value_to_jsonb;
39use crate::row_writer::{self, MultiTableData, TableData};
40
41const APPROXIMATE_COLUMN_COUNT: usize = 30;
42
43const MAX_TIMESTAMP: i64 = 4102444800000000000;
45
46struct FixedTraceColumnIndexes {
47 timestamp: usize,
48 timestamp_end: usize,
49 duration_nano: usize,
50 parent_span_id: usize,
51 trace_id: usize,
52 span_id: usize,
53 span_kind: usize,
54 span_name: usize,
55 span_status_code: usize,
56 span_status_message: usize,
57 trace_state: usize,
58 scope_name: usize,
59 scope_version: usize,
60}
61
62impl FixedTraceColumnIndexes {
63 fn resolve(writer: &mut TableData) -> Result<Self> {
64 let timestamp = writer.ensure_column(time_index_column_schema(
65 TIMESTAMP_COLUMN,
66 ColumnDataType::TimestampNanosecond,
67 ))?;
68 let mut field = |name: &str, datatype: ColumnDataType| {
69 writer.ensure_column(ColumnSchema {
70 column_name: name.to_string(),
71 datatype: datatype as i32,
72 semantic_type: SemanticType::Field as i32,
73 ..Default::default()
74 })
75 };
76
77 Ok(Self {
78 timestamp,
79 timestamp_end: field(TIMESTAMP_END_COLUMN, ColumnDataType::TimestampNanosecond)?,
80 duration_nano: field(DURATION_NANO_COLUMN, ColumnDataType::Int64)?,
81 parent_span_id: field(PARENT_SPAN_ID_COLUMN, ColumnDataType::String)?,
82 trace_id: field(TRACE_ID_COLUMN, ColumnDataType::String)?,
83 span_id: field(SPAN_ID_COLUMN, ColumnDataType::String)?,
84 span_kind: field(SPAN_KIND_COLUMN, ColumnDataType::String)?,
85 span_name: field(SPAN_NAME_COLUMN, ColumnDataType::String)?,
86 span_status_code: field(SPAN_STATUS_CODE, ColumnDataType::String)?,
87 span_status_message: field(SPAN_STATUS_MESSAGE_COLUMN, ColumnDataType::String)?,
88 trace_state: field(TRACE_STATE_COLUMN, ColumnDataType::String)?,
89 scope_name: field(SCOPE_NAME_COLUMN, ColumnDataType::String)?,
90 scope_version: field(SCOPE_VERSION_COLUMN, ColumnDataType::String)?,
91 })
92 }
93}
94
95#[derive(Debug, Clone, Copy, PartialEq, Eq)]
97pub enum TraceBinaryType {
98 Binary,
100 Json,
102}
103
104impl TraceBinaryType {
105 pub fn apply_to_schema(self, schema: &mut ColumnSchema) {
107 schema.datatype = ColumnDataType::Binary as i32;
108 schema.datatype_extension = match self {
109 Self::Binary => None,
110 Self::Json => Some(ColumnDataTypeExtension {
111 type_ext: Some(TypeExt::JsonType(JsonTypeExtension::JsonBinary.into())),
112 }),
113 };
114 }
115}
116
117#[derive(Default)]
119struct TraceBatchColumnSchema {
120 value_types: Vec<ColumnDataType>,
121 present_rows: Vec<usize>,
122 binary_types: Vec<(usize, TraceBinaryType)>,
123}
124
125impl TraceBatchColumnSchema {
126 fn observe_row(&mut self, row_index: usize) {
128 if self.present_rows.last() != Some(&row_index) {
129 self.present_rows.push(row_index);
130 }
131 }
132
133 fn has_incompatible_logical_types(&self) -> bool {
135 let Some((_, first_type)) = self.binary_types.first() else {
136 return false;
137 };
138 self.binary_types
139 .iter()
140 .any(|(_, binary_type)| binary_type != first_type)
141 }
142}
143
144pub struct TraceRetryColumn {
146 pub present_rows: Vec<usize>,
148 pub binary_types: Vec<(usize, TraceBinaryType)>,
150}
151
152pub type TraceRetryColumns = HashMap<String, TraceRetryColumn>;
154
155#[derive(Default)]
157pub struct TraceBatchSchema {
158 columns: HashMap<String, TraceBatchColumnSchema>,
159}
160
161impl TraceBatchSchema {
162 fn observe_value_type(&mut self, name: &str, row_index: usize, value_type: ColumnDataType) {
164 let column = self.columns.entry(name.to_string()).or_default();
165 column.observe_row(row_index);
166 if !column.value_types.contains(&value_type) {
167 column.value_types.push(value_type);
168 }
169 }
170
171 fn observe_binary_type(&mut self, name: &str, row_index: usize, binary_type: TraceBinaryType) {
173 let column = self.columns.entry(name.to_string()).or_default();
174 column.observe_row(row_index);
175 if !column.value_types.contains(&ColumnDataType::Binary) {
176 column.value_types.push(ColumnDataType::Binary);
177 }
178 if let Some((last_row_index, last_type)) = column.binary_types.last_mut()
179 && *last_row_index == row_index
180 {
181 *last_type = binary_type;
182 } else {
183 column.binary_types.push((row_index, binary_type));
184 }
185 }
186
187 fn observe_present_column(&mut self, name: &str, row_index: usize) {
189 self.columns
190 .entry(name.to_string())
191 .or_default()
192 .observe_row(row_index);
193 }
194
195 pub fn value_types(&self, name: &str) -> Option<&[ColumnDataType]> {
197 self.columns
198 .get(name)
199 .map(|column| column.value_types.as_slice())
200 }
201
202 pub fn has_incompatible_logical_types(&self) -> bool {
204 self.columns
205 .values()
206 .any(TraceBatchColumnSchema::has_incompatible_logical_types)
207 }
208
209 pub fn has_incompatible_logical_types_for(&self, name: &str) -> bool {
211 self.columns
212 .get(name)
213 .is_some_and(TraceBatchColumnSchema::has_incompatible_logical_types)
214 }
215
216 pub fn into_retry_columns(self) -> TraceRetryColumns {
218 self.columns
219 .into_iter()
220 .filter(|(_, column)| !column.present_rows.is_empty())
221 .map(|(name, column)| {
222 let binary_types = if column.has_incompatible_logical_types()
223 || column
224 .value_types
225 .iter()
226 .any(|datatype| *datatype != ColumnDataType::Binary)
227 {
228 column.binary_types
229 } else {
230 Vec::new()
231 };
232 (
233 name,
234 TraceRetryColumn {
235 present_rows: column.present_rows,
236 binary_types,
237 },
238 )
239 })
240 .collect()
241 }
242}
243
244pub fn v1_to_grpc_main_insert_requests(
249 spans: &[TraceSpan],
250 table_name: &str,
251) -> Result<(RowInsertRequests, usize)> {
252 let requests = v1_to_grpc_main_insert_requests_from_iter(spans.iter().cloned(), table_name)?;
253 Ok((requests, spans.len()))
254}
255
256pub fn v1_to_main_table_data_with_schema(
258 spans: Vec<TraceSpan>,
259) -> Result<(TableData, TraceBatchSchema)> {
260 build_trace_table_data_with_schema(spans.into_iter())
261}
262
263fn v1_to_grpc_main_insert_requests_from_iter(
265 spans: impl ExactSizeIterator<Item = TraceSpan>,
266 table_name: &str,
267) -> Result<RowInsertRequests> {
268 let mut multi_table_writer = MultiTableData::default();
269 let trace_writer = build_trace_table_data_from_iter(spans, None)?;
270 multi_table_writer.add_table_data(table_name, trace_writer);
271
272 Ok(multi_table_writer.into_row_insert_requests().0)
273}
274
275pub fn build_trace_table_data(spans: &[TraceSpan]) -> Result<TableData> {
277 build_trace_table_data_from_iter(spans.iter().cloned(), None)
278}
279
280fn build_trace_table_data_with_schema(
282 spans: impl ExactSizeIterator<Item = TraceSpan>,
283) -> Result<(TableData, TraceBatchSchema)> {
284 let mut batch_schema = TraceBatchSchema::default();
285 let trace_writer = build_trace_table_data_from_iter(spans, Some(&mut batch_schema))?;
286 Ok((trace_writer, batch_schema))
287}
288
289fn build_trace_table_data_from_iter(
291 spans: impl ExactSizeIterator<Item = TraceSpan>,
292 mut batch_schema: Option<&mut TraceBatchSchema>,
293) -> Result<TableData> {
294 let mut trace_writer = TableData::new(APPROXIMATE_COLUMN_COUNT, spans.len());
295 if spans.len() == 0 {
296 return Ok(trace_writer);
297 }
298 let fixed_columns = FixedTraceColumnIndexes::resolve(&mut trace_writer)?;
299 for span in spans {
300 let row_index = trace_writer.num_rows();
301 write_span_to_row_inner(
302 &mut trace_writer,
303 span,
304 row_index,
305 &fixed_columns,
306 batch_schema.as_deref_mut(),
307 )?;
308 }
309 Ok(trace_writer)
310}
311
312pub fn build_aux_table_requests(
314 aux_data: TraceAuxData,
315 table_name: &str,
316) -> Result<(RowInsertRequests, usize)> {
317 let mut multi_table_writer = MultiTableData::default();
318 let mut trace_services_writer = TableData::new(APPROXIMATE_COLUMN_COUNT, 1);
319 let mut trace_operations_writer = TableData::new(APPROXIMATE_COLUMN_COUNT, 1);
320
321 write_trace_services_to_row(&mut trace_services_writer, aux_data.services)?;
322 write_trace_operations_to_row(&mut trace_operations_writer, aux_data.operations)?;
323
324 multi_table_writer.add_table_data(trace_services_table_name(table_name), trace_services_writer);
325 multi_table_writer.add_table_data(
326 trace_operations_table_name(table_name),
327 trace_operations_writer,
328 );
329
330 Ok(multi_table_writer.into_row_insert_requests())
331}
332
333pub fn write_span_to_row(writer: &mut TableData, span: TraceSpan) -> Result<()> {
334 let row_index = writer.num_rows();
335 let fixed_columns = FixedTraceColumnIndexes::resolve(writer)?;
336 write_span_to_row_inner(writer, span, row_index, &fixed_columns, None)
337}
338
339pub(super) fn span_duration_nano(span: &TraceSpan) -> i64 {
350 span.end_in_nanosecond
351 .saturating_sub(span.start_in_nanosecond)
352 .min(i64::MAX as u64) as i64
353}
354
355fn write_span_to_row_inner(
357 writer: &mut TableData,
358 span: TraceSpan,
359 row_index: usize,
360 fixed_columns: &FixedTraceColumnIndexes,
361 mut batch_schema: Option<&mut TraceBatchSchema>,
362) -> Result<()> {
363 let mut row = writer.alloc_one_row();
364
365 for (index, value) in [
366 (
367 fixed_columns.timestamp,
368 Some(ValueData::TimestampNanosecondValue(
369 span.start_in_nanosecond as i64,
370 )),
371 ),
372 (
373 fixed_columns.timestamp_end,
374 Some(ValueData::TimestampNanosecondValue(
375 span.end_in_nanosecond as i64,
376 )),
377 ),
378 (
379 fixed_columns.duration_nano,
380 Some(ValueData::I64Value(span_duration_nano(&span))),
381 ),
382 (
383 fixed_columns.parent_span_id,
384 span.parent_span_id.map(ValueData::StringValue),
385 ),
386 (
387 fixed_columns.trace_id,
388 Some(ValueData::StringValue(span.trace_id)),
389 ),
390 (
391 fixed_columns.span_id,
392 Some(ValueData::StringValue(span.span_id)),
393 ),
394 (
395 fixed_columns.span_kind,
396 Some(ValueData::StringValue(span.span_kind)),
397 ),
398 (
399 fixed_columns.span_name,
400 Some(ValueData::StringValue(span.span_name)),
401 ),
402 (
403 fixed_columns.span_status_code,
404 Some(ValueData::StringValue(span.span_status_code)),
405 ),
406 (
407 fixed_columns.span_status_message,
408 Some(ValueData::StringValue(span.span_status_message)),
409 ),
410 (
411 fixed_columns.trace_state,
412 Some(ValueData::StringValue(span.trace_state)),
413 ),
414 (
415 fixed_columns.scope_name,
416 Some(ValueData::StringValue(span.scope_name)),
417 ),
418 (
419 fixed_columns.scope_version,
420 Some(ValueData::StringValue(span.scope_version)),
421 ),
422 ] {
423 row[index].value_data = value;
424 }
425
426 if let Some(service_name) = span.service_name {
427 if let Some(batch_schema) = batch_schema.as_deref_mut() {
428 batch_schema.observe_present_column(SERVICE_NAME_COLUMN, row_index);
429 }
430 row_writer::write_tags(
431 writer,
432 std::iter::once((SERVICE_NAME_COLUMN.to_string(), service_name)),
433 &mut row,
434 )?;
435 }
436
437 write_attributes_with_schema(
438 writer,
439 "span_attributes",
440 span.span_attributes.take(),
441 &mut row,
442 row_index,
443 batch_schema.as_deref_mut(),
444 )?;
445 write_shared_attributes_with_schema(
446 writer,
447 "scope_attributes",
448 span.scope_attributes.as_ref(),
449 &mut row,
450 row_index,
451 batch_schema.as_deref_mut(),
452 )?;
453 write_shared_attributes_with_schema(
454 writer,
455 "resource_attributes",
456 span.resource_attributes.as_ref(),
457 &mut row,
458 row_index,
459 batch_schema,
460 )?;
461
462 row_writer::write_json(
463 writer,
464 SPAN_EVENTS_COLUMN,
465 span.span_events.into(),
466 &mut row,
467 )?;
468 row_writer::write_json(writer, "span_links", span.span_links.into(), &mut row)?;
469
470 writer.add_row(row);
471
472 Ok(())
473}
474
475fn write_trace_services_to_row(writer: &mut TableData, services: HashSet<String>) -> Result<()> {
476 for service_name in services {
477 let mut row = writer.alloc_one_row();
478 row_writer::write_ts_to_nanos(
480 writer,
481 TIMESTAMP_COLUMN,
482 Some(MAX_TIMESTAMP),
483 Precision::Nanosecond,
484 &mut row,
485 )?;
486
487 row_writer::write_tags(
489 writer,
490 std::iter::once((SERVICE_NAME_COLUMN.to_string(), service_name)),
491 &mut row,
492 )?;
493 writer.add_row(row);
494 }
495
496 Ok(())
497}
498
499fn write_trace_operations_to_row(
500 writer: &mut TableData,
501 operations: HashSet<(String, String, String)>,
502) -> Result<()> {
503 for (service_name, span_name, span_kind) in operations {
504 let mut row = writer.alloc_one_row();
505 row_writer::write_ts_to_nanos(
507 writer,
508 TIMESTAMP_COLUMN,
509 Some(MAX_TIMESTAMP),
510 Precision::Nanosecond,
511 &mut row,
512 )?;
513
514 row_writer::write_tags(
516 writer,
517 vec![
518 (SERVICE_NAME_COLUMN.to_string(), service_name),
519 (SPAN_NAME_COLUMN.to_string(), span_name),
520 (SPAN_KIND_COLUMN.to_string(), span_kind),
521 ]
522 .into_iter(),
523 &mut row,
524 )?;
525 writer.add_row(row);
526 }
527
528 Ok(())
529}
530
531#[cfg(test)]
532pub(crate) fn write_attributes(
533 writer: &mut TableData,
534 prefix: &str,
535 attributes: Attributes,
536 row: &mut Vec<Value>,
537) -> Result<()> {
538 let row_index = writer.num_rows();
539 write_attributes_with_schema(writer, prefix, attributes.take(), row, row_index, None)
540}
541
542fn skipped_attribute(prefix: &str, key: &str) -> bool {
545 prefix == "resource_attributes" && key == KEY_SERVICE_NAME
546}
547
548fn write_attributes_with_schema(
550 writer: &mut TableData,
551 prefix: &str,
552 attributes: Vec<KeyValue>,
553 row: &mut Vec<Value>,
554 row_index: usize,
555 mut batch_schema: Option<&mut TraceBatchSchema>,
556) -> Result<()> {
557 for KeyValue { key, value, .. } in attributes {
558 if skipped_attribute(prefix, &key) {
559 continue;
560 }
561 write_attribute_with_schema(
562 writer,
563 prefix,
564 &key,
565 value.and_then(|v| v.value),
566 row,
567 row_index,
568 batch_schema.as_deref_mut(),
569 );
570 }
571
572 Ok(())
573}
574
575fn write_shared_attributes_with_schema(
580 writer: &mut TableData,
581 prefix: &str,
582 attributes: &Attributes,
583 row: &mut Vec<Value>,
584 row_index: usize,
585 mut batch_schema: Option<&mut TraceBatchSchema>,
586) -> Result<()> {
587 for attr in attributes.get_ref() {
588 if skipped_attribute(prefix, &attr.key) {
589 continue;
590 }
591 write_attribute_with_schema(
592 writer,
593 prefix,
594 &attr.key,
595 attr.value.as_ref().and_then(|v| v.value.clone()),
596 row,
597 row_index,
598 batch_schema.as_deref_mut(),
599 );
600 }
601
602 Ok(())
603}
604
605fn write_attribute_with_schema(
607 writer: &mut TableData,
608 prefix: &str,
609 key_suffix: &str,
610 value: Option<OtlpValue>,
611 row: &mut Vec<Value>,
612 row_index: usize,
613 batch_schema: Option<&mut TraceBatchSchema>,
614) {
615 let key = format!("{}.{}", prefix, key_suffix);
616 match value {
617 Some(OtlpValue::StringValue(v)) => {
618 if let Some(batch_schema) = batch_schema {
619 batch_schema.observe_value_type(&key, row_index, ColumnDataType::String);
620 }
621 writer.write_field_unchecked(
624 &key,
625 ColumnDataType::String,
626 Some(ValueData::StringValue(v)),
627 row,
628 );
629 }
630 Some(OtlpValue::BoolValue(v)) => {
631 if let Some(batch_schema) = batch_schema {
632 batch_schema.observe_value_type(&key, row_index, ColumnDataType::Boolean);
633 }
634 writer.write_field_unchecked(
636 &key,
637 ColumnDataType::Boolean,
638 Some(ValueData::BoolValue(v)),
639 row,
640 );
641 }
642 Some(OtlpValue::IntValue(v)) => {
643 if let Some(batch_schema) = batch_schema {
644 batch_schema.observe_value_type(&key, row_index, ColumnDataType::Int64);
645 }
646 writer.write_field_unchecked(
648 &key,
649 ColumnDataType::Int64,
650 Some(ValueData::I64Value(v)),
651 row,
652 );
653 }
654 Some(OtlpValue::DoubleValue(v)) => {
655 if let Some(batch_schema) = batch_schema {
656 batch_schema.observe_value_type(&key, row_index, ColumnDataType::Float64);
657 }
658 writer.write_field_unchecked(
659 &key,
660 ColumnDataType::Float64,
661 Some(ValueData::F64Value(v)),
662 row,
663 );
664 }
665 Some(OtlpValue::ArrayValue(v)) => {
666 if let Some(batch_schema) = batch_schema {
667 batch_schema.observe_binary_type(&key, row_index, TraceBinaryType::Json);
668 }
669 writer.write_column_unchecked(
670 row_writer::build_json_column_schema(key),
671 Some(ValueData::BinaryValue(
672 any_value_to_jsonb(OtlpValue::ArrayValue(v)).to_vec(),
673 )),
674 row,
675 );
676 }
677 Some(OtlpValue::KvlistValue(v)) => {
678 if let Some(batch_schema) = batch_schema {
679 batch_schema.observe_binary_type(&key, row_index, TraceBinaryType::Json);
680 }
681 writer.write_column_unchecked(
682 row_writer::build_json_column_schema(key),
683 Some(ValueData::BinaryValue(
684 any_value_to_jsonb(OtlpValue::KvlistValue(v)).to_vec(),
685 )),
686 row,
687 );
688 }
689 Some(OtlpValue::BytesValue(v)) => {
690 if let Some(batch_schema) = batch_schema {
691 batch_schema.observe_binary_type(&key, row_index, TraceBinaryType::Binary);
692 }
693 writer.write_field_unchecked(
694 key,
695 ColumnDataType::Binary,
696 Some(ValueData::BinaryValue(v)),
697 row,
698 );
699 }
700 Some(OtlpValue::StringValueStrindex(_)) => {}
706 None => {}
707 }
708}
709
710#[cfg(test)]
711mod tests {
712 use api::v1::value::ValueData;
713 use opentelemetry_proto::tonic::common::v1::any_value::Value as OtlpValue;
714 use opentelemetry_proto::tonic::common::v1::{AnyValue, ArrayValue, KeyValue};
715
716 use super::*;
717 use crate::otlp::trace::TraceAuxData;
718 use crate::otlp::trace::attributes::Attributes;
719 use crate::otlp::trace::span::{SpanEvents, SpanLinks};
720 use crate::row_writer::TableData;
721
722 fn make_kv(key: &str, value: OtlpValue) -> KeyValue {
723 KeyValue {
724 key: key.to_string(),
725 value: Some(AnyValue { value: Some(value) }),
726 ..Default::default()
727 }
728 }
729
730 fn make_span(service_name: &str, trace_id: &str, span_id: &str) -> TraceSpan {
731 TraceSpan {
732 service_name: Some(service_name.to_string()),
733 trace_id: trace_id.to_string(),
734 span_id: span_id.to_string(),
735 parent_span_id: None,
736 resource_attributes: Attributes::from(vec![]).into(),
737 scope_name: "scope".to_string(),
738 scope_version: "v1".to_string(),
739 scope_attributes: Attributes::from(vec![]).into(),
740 trace_state: String::new(),
741 span_name: "op".to_string(),
742 span_kind: "SPAN_KIND_SERVER".to_string(),
743 span_status_code: "STATUS_CODE_UNSET".to_string(),
744 span_status_message: String::new(),
745 span_attributes: Attributes::from(vec![]),
746 span_events: SpanEvents::from(vec![]),
747 span_links: SpanLinks::from(vec![]),
748 start_in_nanosecond: 1,
749 end_in_nanosecond: 2,
750 }
751 }
752
753 #[test]
754 fn test_span_end_before_start_records_zero_duration() {
755 let mut span = make_span("svc", "trace", "span");
760 span.start_in_nanosecond = 200;
761 span.end_in_nanosecond = 100;
762
763 let (schema, rows) = build_trace_table_data(&[span])
764 .unwrap()
765 .into_schema_and_rows();
766
767 let idx = schema
768 .iter()
769 .position(|c| c.column_name == DURATION_NANO_COLUMN)
770 .unwrap();
771 assert_eq!(rows[0].values[idx].value_data, Some(ValueData::I64Value(0)));
772 }
773
774 #[test]
775 fn test_span_duration_above_i64_max_saturates() {
776 let mut span = make_span("svc", "trace", "span");
779 span.start_in_nanosecond = 0;
780 span.end_in_nanosecond = u64::MAX;
781
782 let (schema, rows) = build_trace_table_data(&[span])
783 .unwrap()
784 .into_schema_and_rows();
785
786 let idx = schema
787 .iter()
788 .position(|c| c.column_name == DURATION_NANO_COLUMN)
789 .unwrap();
790 assert_eq!(
791 rows[0].values[idx].value_data,
792 Some(ValueData::I64Value(i64::MAX))
793 );
794 }
795
796 #[test]
797 fn test_fixed_trace_columns_keep_schema_and_values_aligned() {
798 let table_data = build_trace_table_data(&[make_span("svc", "trace", "span")]).unwrap();
799 let (schema, rows) = table_data.into_schema_and_rows();
800 let expected = [
801 (
802 TIMESTAMP_COLUMN,
803 Some(ValueData::TimestampNanosecondValue(1)),
804 ),
805 (
806 TIMESTAMP_END_COLUMN,
807 Some(ValueData::TimestampNanosecondValue(2)),
808 ),
809 (DURATION_NANO_COLUMN, Some(ValueData::I64Value(1))),
810 (PARENT_SPAN_ID_COLUMN, None),
811 (
812 TRACE_ID_COLUMN,
813 Some(ValueData::StringValue("trace".to_string())),
814 ),
815 (
816 SPAN_ID_COLUMN,
817 Some(ValueData::StringValue("span".to_string())),
818 ),
819 (
820 SPAN_KIND_COLUMN,
821 Some(ValueData::StringValue("SPAN_KIND_SERVER".to_string())),
822 ),
823 (
824 SPAN_NAME_COLUMN,
825 Some(ValueData::StringValue("op".to_string())),
826 ),
827 (
828 SPAN_STATUS_CODE,
829 Some(ValueData::StringValue("STATUS_CODE_UNSET".to_string())),
830 ),
831 (
832 SPAN_STATUS_MESSAGE_COLUMN,
833 Some(ValueData::StringValue(String::new())),
834 ),
835 (
836 TRACE_STATE_COLUMN,
837 Some(ValueData::StringValue(String::new())),
838 ),
839 (
840 SCOPE_NAME_COLUMN,
841 Some(ValueData::StringValue("scope".to_string())),
842 ),
843 (
844 SCOPE_VERSION_COLUMN,
845 Some(ValueData::StringValue("v1".to_string())),
846 ),
847 ];
848
849 for (index, (name, value)) in expected.into_iter().enumerate() {
850 assert_eq!(schema[index].column_name, name);
851 assert_eq!(rows[0].values[index].value_data, value);
852 }
853 }
854
855 #[test]
856 fn test_batch_schema_preserves_column_order_and_attribute_types() {
857 let mut span1 = make_span("svc-a", "trace-a", "span-a");
858 span1.span_attributes = Attributes::from(vec![make_kv("val", OtlpValue::IntValue(10))]);
859 let mut span2 = make_span("svc-a", "trace-a", "span-b");
860 span2.span_attributes = Attributes::from(vec![make_kv("val", OtlpValue::DoubleValue(1.5))]);
861
862 let (table_data, batch_schema) =
863 build_trace_table_data_with_schema(vec![span1, span2].into_iter()).unwrap();
864
865 assert_eq!(
866 batch_schema.value_types("span_attributes.val").unwrap(),
867 [ColumnDataType::Int64, ColumnDataType::Float64]
868 );
869 let timestamp_index = table_data
870 .columns()
871 .iter()
872 .position(|column| column.column_name == TIMESTAMP_COLUMN)
873 .unwrap();
874 let attribute_index = table_data
875 .columns()
876 .iter()
877 .position(|column| column.column_name == "span_attributes.val")
878 .unwrap();
879 assert!(timestamp_index < attribute_index);
880 }
881
882 #[test]
883 fn test_optional_batch_schema_observation_preserves_rows() {
884 let mut span = make_span("svc-a", "trace-a", "span-a");
885 span.span_attributes = Attributes::from(vec![make_kv("val", OtlpValue::IntValue(10))]);
886
887 let without_schema = build_trace_table_data(std::slice::from_ref(&span)).unwrap();
888 let (with_schema, batch_schema) =
889 build_trace_table_data_with_schema(vec![span].into_iter()).unwrap();
890
891 assert_eq!(
892 without_schema.into_schema_and_rows(),
893 with_schema.into_schema_and_rows()
894 );
895 assert_eq!(
896 batch_schema.value_types("span_attributes.val").unwrap(),
897 [ColumnDataType::Int64]
898 );
899 let retry_columns = batch_schema.into_retry_columns();
900 let retry_column = &retry_columns["span_attributes.val"];
901 assert_eq!(retry_column.present_rows, [0]);
902 assert!(retry_column.binary_types.is_empty());
903 }
904
905 #[test]
906 fn test_batch_schema_tracks_binary_and_json_rows() {
907 let mut binary_span = make_span("svc-a", "trace-a", "span-a");
908 binary_span.span_attributes = Attributes::from(vec![make_kv(
909 "val",
910 OtlpValue::BytesValue(vec![1_u8, 2, 3]),
911 )]);
912 let mut json_span = make_span("svc-a", "trace-a", "span-b");
913 json_span.span_attributes = Attributes::from(vec![make_kv(
914 "val",
915 OtlpValue::ArrayValue(ArrayValue {
916 values: vec![AnyValue {
917 value: Some(OtlpValue::IntValue(1)),
918 }],
919 }),
920 )]);
921
922 let (_, batch_schema) =
923 build_trace_table_data_with_schema(vec![binary_span, json_span].into_iter()).unwrap();
924
925 assert!(batch_schema.has_incompatible_logical_types());
926 assert!(batch_schema.has_incompatible_logical_types_for("span_attributes.val"));
927 let retry_columns = batch_schema.into_retry_columns();
928 let retry_column = &retry_columns["span_attributes.val"];
929 assert_eq!(retry_column.present_rows, [0, 1]);
930 assert_eq!(
931 retry_column.binary_types,
932 [(0, TraceBinaryType::Binary), (1, TraceBinaryType::Json)]
933 );
934 }
935
936 #[test]
937 fn test_batch_schema_preserves_scalar_then_binary_values() {
938 let mut scalar_span = make_span("svc-a", "trace-a", "span-a");
939 scalar_span.span_attributes = Attributes::from(vec![
940 make_kv("bytes", OtlpValue::StringValue("text".to_string())),
941 make_kv("json", OtlpValue::StringValue("text".to_string())),
942 ]);
943 let mut binary_span = make_span("svc-a", "trace-a", "span-b");
944 binary_span.span_attributes = Attributes::from(vec![
945 make_kv("bytes", OtlpValue::BytesValue(vec![1_u8, 2, 3])),
946 make_kv(
947 "json",
948 OtlpValue::ArrayValue(ArrayValue {
949 values: vec![AnyValue {
950 value: Some(OtlpValue::IntValue(1)),
951 }],
952 }),
953 ),
954 ]);
955
956 let (_, batch_schema) =
957 build_trace_table_data_with_schema(vec![scalar_span, binary_span].into_iter()).unwrap();
958
959 assert_eq!(
960 batch_schema.value_types("span_attributes.bytes").unwrap(),
961 [ColumnDataType::String, ColumnDataType::Binary]
962 );
963 assert_eq!(
964 batch_schema.value_types("span_attributes.json").unwrap(),
965 [ColumnDataType::String, ColumnDataType::Binary]
966 );
967 let retry_columns = batch_schema.into_retry_columns();
968 assert_eq!(
969 retry_columns["span_attributes.bytes"].binary_types,
970 [(1, TraceBinaryType::Binary)]
971 );
972 assert_eq!(
973 retry_columns["span_attributes.json"].binary_types,
974 [(1, TraceBinaryType::Json)]
975 );
976 }
977
978 #[test]
979 fn test_keep_mixed_numeric_values_until_frontend_reconciliation() {
980 let mut writer = TableData::new(4, 2);
981
982 let attrs1 = Attributes::from(vec![make_kv("val", OtlpValue::DoubleValue(1.5))]);
983 let mut row1 = writer.alloc_one_row();
984 write_attributes(&mut writer, "attr", attrs1, &mut row1).unwrap();
985 writer.add_row(row1);
986
987 let attrs2 = Attributes::from(vec![make_kv("val", OtlpValue::IntValue(42))]);
988 let mut row2 = writer.alloc_one_row();
989 write_attributes(&mut writer, "attr", attrs2, &mut row2).unwrap();
990 writer.add_row(row2);
991
992 let (schema, rows) = writer.into_schema_and_rows();
993
994 let col_idx = schema
995 .iter()
996 .position(|c| c.column_name == "attr.val")
997 .unwrap();
998 assert_eq!(schema[col_idx].datatype, ColumnDataType::Float64 as i32);
999
1000 assert_eq!(
1001 rows[0].values[col_idx].value_data,
1002 Some(ValueData::F64Value(1.5))
1003 );
1004 assert_eq!(
1005 rows[1].values[col_idx].value_data,
1006 Some(ValueData::I64Value(42))
1007 );
1008 }
1009
1010 #[test]
1011 fn test_keep_mixed_string_and_int_values_until_frontend_reconciliation() {
1012 let mut writer = TableData::new(4, 2);
1013
1014 let attrs1 = Attributes::from(vec![make_kv("val", OtlpValue::IntValue(10))]);
1015 let mut row1 = writer.alloc_one_row();
1016 write_attributes(&mut writer, "attr", attrs1, &mut row1).unwrap();
1017 writer.add_row(row1);
1018
1019 let attrs2 = Attributes::from(vec![make_kv(
1020 "val",
1021 OtlpValue::StringValue("20".to_string()),
1022 )]);
1023 let mut row2 = writer.alloc_one_row();
1024 write_attributes(&mut writer, "attr", attrs2, &mut row2).unwrap();
1025 writer.add_row(row2);
1026
1027 let (schema, rows) = writer.into_schema_and_rows();
1028 let col_idx = schema
1029 .iter()
1030 .position(|c| c.column_name == "attr.val")
1031 .unwrap();
1032 assert_eq!(schema[col_idx].datatype, ColumnDataType::Int64 as i32);
1033 assert_eq!(
1034 rows[1].values[col_idx].value_data,
1035 Some(ValueData::StringValue("20".to_string()))
1036 );
1037 }
1038
1039 #[test]
1040 fn test_keep_first_seen_schema_until_frontend_reconciliation() {
1041 let mut writer = TableData::new(4, 2);
1042
1043 let attrs1 = Attributes::from(vec![make_kv(
1044 "val",
1045 OtlpValue::StringValue("10".to_string()),
1046 )]);
1047 let mut row1 = writer.alloc_one_row();
1048 write_attributes(&mut writer, "attr", attrs1, &mut row1).unwrap();
1049 writer.add_row(row1);
1050
1051 let attrs2 = Attributes::from(vec![make_kv("val", OtlpValue::IntValue(20))]);
1052 let mut row2 = writer.alloc_one_row();
1053 write_attributes(&mut writer, "attr", attrs2, &mut row2).unwrap();
1054 writer.add_row(row2);
1055
1056 let (schema, rows) = writer.into_schema_and_rows();
1057 let col_idx = schema
1058 .iter()
1059 .position(|c| c.column_name == "attr.val")
1060 .unwrap();
1061 assert_eq!(schema[col_idx].datatype, ColumnDataType::String as i32);
1062 assert_eq!(
1063 rows[0].values[col_idx].value_data,
1064 Some(ValueData::StringValue("10".to_string()))
1065 );
1066 assert_eq!(
1067 rows[1].values[col_idx].value_data,
1068 Some(ValueData::I64Value(20))
1069 );
1070 }
1071
1072 #[test]
1073 fn test_keep_mixed_string_and_float_values_until_frontend_reconciliation() {
1074 let mut writer = TableData::new(4, 2);
1075
1076 let attrs1 = Attributes::from(vec![make_kv("val", OtlpValue::DoubleValue(1.5))]);
1077 let mut row1 = writer.alloc_one_row();
1078 write_attributes(&mut writer, "attr", attrs1, &mut row1).unwrap();
1079 writer.add_row(row1);
1080
1081 let attrs2 = Attributes::from(vec![make_kv(
1082 "val",
1083 OtlpValue::StringValue("1.5".to_string()),
1084 )]);
1085 let mut row2 = writer.alloc_one_row();
1086 write_attributes(&mut writer, "attr", attrs2, &mut row2).unwrap();
1087 writer.add_row(row2);
1088
1089 let (schema, rows) = writer.into_schema_and_rows();
1090 let col_idx = schema
1091 .iter()
1092 .position(|c| c.column_name == "attr.val")
1093 .unwrap();
1094 assert_eq!(schema[col_idx].datatype, ColumnDataType::Float64 as i32);
1095 assert_eq!(
1096 rows[1].values[col_idx].value_data,
1097 Some(ValueData::StringValue("1.5".to_string()))
1098 );
1099 }
1100
1101 #[test]
1102 fn test_keep_mixed_string_and_bool_values_until_frontend_reconciliation() {
1103 let mut writer = TableData::new(4, 2);
1104
1105 let attrs1 = Attributes::from(vec![make_kv(
1106 "val",
1107 OtlpValue::StringValue("true".to_string()),
1108 )]);
1109 let mut row1 = writer.alloc_one_row();
1110 write_attributes(&mut writer, "attr", attrs1, &mut row1).unwrap();
1111 writer.add_row(row1);
1112
1113 let attrs2 = Attributes::from(vec![make_kv("val", OtlpValue::BoolValue(false))]);
1114 let mut row2 = writer.alloc_one_row();
1115 write_attributes(&mut writer, "attr", attrs2, &mut row2).unwrap();
1116 writer.add_row(row2);
1117
1118 let (schema, rows) = writer.into_schema_and_rows();
1119 let col_idx = schema
1120 .iter()
1121 .position(|c| c.column_name == "attr.val")
1122 .unwrap();
1123 assert_eq!(schema[col_idx].datatype, ColumnDataType::String as i32);
1124 assert_eq!(
1125 rows[0].values[col_idx].value_data,
1126 Some(ValueData::StringValue("true".to_string()))
1127 );
1128 assert_eq!(
1129 rows[1].values[col_idx].value_data,
1130 Some(ValueData::BoolValue(false))
1131 );
1132 }
1133
1134 #[test]
1135 fn test_keep_mixed_binary_and_string_values_until_frontend_reconciliation() {
1136 let mut writer = TableData::new(4, 2);
1137
1138 let attrs1 = Attributes::from(vec![make_kv(
1139 "val",
1140 OtlpValue::BytesValue(vec![1_u8, 2, 3]),
1141 )]);
1142 let mut row1 = writer.alloc_one_row();
1143 write_attributes(&mut writer, "attr", attrs1, &mut row1).unwrap();
1144 writer.add_row(row1);
1145
1146 let attrs2 = Attributes::from(vec![make_kv(
1147 "val",
1148 OtlpValue::StringValue("false".to_string()),
1149 )]);
1150 let mut row2 = writer.alloc_one_row();
1151 write_attributes(&mut writer, "attr", attrs2, &mut row2).unwrap();
1152 writer.add_row(row2);
1153
1154 let (schema, rows) = writer.into_schema_and_rows();
1155 let col_idx = schema
1156 .iter()
1157 .position(|c| c.column_name == "attr.val")
1158 .unwrap();
1159 assert_eq!(schema[col_idx].datatype, ColumnDataType::Binary as i32);
1160 assert_eq!(
1161 rows[0].values[col_idx].value_data,
1162 Some(ValueData::BinaryValue(vec![1_u8, 2, 3]))
1163 );
1164 assert_eq!(
1165 rows[1].values[col_idx].value_data,
1166 Some(ValueData::StringValue("false".to_string()))
1167 );
1168 }
1169
1170 #[test]
1171 fn test_keep_mixed_binary_and_json_values_until_frontend_reconciliation() {
1172 let mut writer = TableData::new(4, 2);
1173
1174 let attrs1 = Attributes::from(vec![make_kv(
1175 "val",
1176 OtlpValue::BytesValue(vec![1_u8, 2, 3]),
1177 )]);
1178 let mut row1 = writer.alloc_one_row();
1179 write_attributes(&mut writer, "attr", attrs1, &mut row1).unwrap();
1180 writer.add_row(row1);
1181
1182 let attrs2 = Attributes::from(vec![make_kv(
1183 "val",
1184 OtlpValue::ArrayValue(ArrayValue {
1185 values: vec![AnyValue {
1186 value: Some(OtlpValue::IntValue(1)),
1187 }],
1188 }),
1189 )]);
1190 let mut row2 = writer.alloc_one_row();
1191 write_attributes(&mut writer, "attr", attrs2, &mut row2).unwrap();
1192 writer.add_row(row2);
1193
1194 let (schema, rows) = writer.into_schema_and_rows();
1195 let col_idx = schema
1196 .iter()
1197 .position(|c| c.column_name == "attr.val")
1198 .unwrap();
1199 assert_eq!(schema[col_idx].datatype, ColumnDataType::Binary as i32);
1200 assert!(matches!(
1201 rows[1].values[col_idx].value_data.as_ref(),
1202 Some(ValueData::BinaryValue(_))
1203 ));
1204 }
1205
1206 #[test]
1207 fn test_build_aux_table_requests_deduplicates_services_and_operations() {
1208 let spans = vec![
1209 make_span("svc-a", "trace-a", "span-a"),
1210 make_span("svc-a", "trace-b", "span-b"),
1211 ];
1212 let mut aux_data = TraceAuxData::default();
1213 for span in &spans {
1214 aux_data.observe_span(span);
1215 }
1216
1217 let (requests, total_rows) =
1218 build_aux_table_requests(aux_data, "opentelemetry_traces").unwrap();
1219 assert_eq!(requests.inserts.len(), 2);
1220 assert_eq!(total_rows, 2);
1221 }
1222 }