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