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::Uint64)?,
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 write_span_to_row_inner(
347 writer: &mut TableData,
348 span: TraceSpan,
349 row_index: usize,
350 fixed_columns: &FixedTraceColumnIndexes,
351 mut batch_schema: Option<&mut TraceBatchSchema>,
352) -> Result<()> {
353 let mut row = writer.alloc_one_row();
354
355 for (index, value) in [
356 (
357 fixed_columns.timestamp,
358 Some(ValueData::TimestampNanosecondValue(
359 span.start_in_nanosecond as i64,
360 )),
361 ),
362 (
363 fixed_columns.timestamp_end,
364 Some(ValueData::TimestampNanosecondValue(
365 span.end_in_nanosecond as i64,
366 )),
367 ),
368 (
369 fixed_columns.duration_nano,
370 Some(ValueData::U64Value(
371 span.end_in_nanosecond - span.start_in_nanosecond,
372 )),
373 ),
374 (
375 fixed_columns.parent_span_id,
376 span.parent_span_id.map(ValueData::StringValue),
377 ),
378 (
379 fixed_columns.trace_id,
380 Some(ValueData::StringValue(span.trace_id)),
381 ),
382 (
383 fixed_columns.span_id,
384 Some(ValueData::StringValue(span.span_id)),
385 ),
386 (
387 fixed_columns.span_kind,
388 Some(ValueData::StringValue(span.span_kind)),
389 ),
390 (
391 fixed_columns.span_name,
392 Some(ValueData::StringValue(span.span_name)),
393 ),
394 (
395 fixed_columns.span_status_code,
396 Some(ValueData::StringValue(span.span_status_code)),
397 ),
398 (
399 fixed_columns.span_status_message,
400 Some(ValueData::StringValue(span.span_status_message)),
401 ),
402 (
403 fixed_columns.trace_state,
404 Some(ValueData::StringValue(span.trace_state)),
405 ),
406 (
407 fixed_columns.scope_name,
408 Some(ValueData::StringValue(span.scope_name)),
409 ),
410 (
411 fixed_columns.scope_version,
412 Some(ValueData::StringValue(span.scope_version)),
413 ),
414 ] {
415 row[index].value_data = value;
416 }
417
418 if let Some(service_name) = span.service_name {
419 if let Some(batch_schema) = batch_schema.as_deref_mut() {
420 batch_schema.observe_present_column(SERVICE_NAME_COLUMN, row_index);
421 }
422 row_writer::write_tags(
423 writer,
424 std::iter::once((SERVICE_NAME_COLUMN.to_string(), service_name)),
425 &mut row,
426 )?;
427 }
428
429 write_attributes_with_schema(
430 writer,
431 "span_attributes",
432 span.span_attributes,
433 &mut row,
434 row_index,
435 batch_schema.as_deref_mut(),
436 )?;
437 write_attributes_with_schema(
438 writer,
439 "scope_attributes",
440 span.scope_attributes,
441 &mut row,
442 row_index,
443 batch_schema.as_deref_mut(),
444 )?;
445 write_attributes_with_schema(
446 writer,
447 "resource_attributes",
448 span.resource_attributes,
449 &mut row,
450 row_index,
451 batch_schema,
452 )?;
453
454 row_writer::write_json(
455 writer,
456 SPAN_EVENTS_COLUMN,
457 span.span_events.into(),
458 &mut row,
459 )?;
460 row_writer::write_json(writer, "span_links", span.span_links.into(), &mut row)?;
461
462 writer.add_row(row);
463
464 Ok(())
465}
466
467fn write_trace_services_to_row(writer: &mut TableData, services: HashSet<String>) -> Result<()> {
468 for service_name in services {
469 let mut row = writer.alloc_one_row();
470 row_writer::write_ts_to_nanos(
472 writer,
473 TIMESTAMP_COLUMN,
474 Some(MAX_TIMESTAMP),
475 Precision::Nanosecond,
476 &mut row,
477 )?;
478
479 row_writer::write_tags(
481 writer,
482 std::iter::once((SERVICE_NAME_COLUMN.to_string(), service_name)),
483 &mut row,
484 )?;
485 writer.add_row(row);
486 }
487
488 Ok(())
489}
490
491fn write_trace_operations_to_row(
492 writer: &mut TableData,
493 operations: HashSet<(String, String, String)>,
494) -> Result<()> {
495 for (service_name, span_name, span_kind) in operations {
496 let mut row = writer.alloc_one_row();
497 row_writer::write_ts_to_nanos(
499 writer,
500 TIMESTAMP_COLUMN,
501 Some(MAX_TIMESTAMP),
502 Precision::Nanosecond,
503 &mut row,
504 )?;
505
506 row_writer::write_tags(
508 writer,
509 vec![
510 (SERVICE_NAME_COLUMN.to_string(), service_name),
511 (SPAN_NAME_COLUMN.to_string(), span_name),
512 (SPAN_KIND_COLUMN.to_string(), span_kind),
513 ]
514 .into_iter(),
515 &mut row,
516 )?;
517 writer.add_row(row);
518 }
519
520 Ok(())
521}
522
523#[cfg(test)]
524pub(crate) fn write_attributes(
525 writer: &mut TableData,
526 prefix: &str,
527 attributes: Attributes,
528 row: &mut Vec<Value>,
529) -> Result<()> {
530 let row_index = writer.num_rows();
531 write_attributes_with_schema(writer, prefix, attributes, row, row_index, None)
532}
533
534fn write_attributes_with_schema(
536 writer: &mut TableData,
537 prefix: &str,
538 attributes: Attributes,
539 row: &mut Vec<Value>,
540 row_index: usize,
541 mut batch_schema: Option<&mut TraceBatchSchema>,
542) -> Result<()> {
543 for attr in attributes.take().into_iter() {
544 let key_suffix = attr.key;
545 if prefix == "resource_attributes" && key_suffix == KEY_SERVICE_NAME {
548 continue;
549 }
550
551 let key = format!("{}.{}", prefix, key_suffix);
552 match attr.value.and_then(|v| v.value) {
553 Some(OtlpValue::StringValue(v)) => {
554 if let Some(batch_schema) = batch_schema.as_deref_mut() {
555 batch_schema.observe_value_type(&key, row_index, ColumnDataType::String);
556 }
557 writer.write_field_unchecked(
560 &key,
561 ColumnDataType::String,
562 Some(ValueData::StringValue(v)),
563 row,
564 );
565 }
566 Some(OtlpValue::BoolValue(v)) => {
567 if let Some(batch_schema) = batch_schema.as_deref_mut() {
568 batch_schema.observe_value_type(&key, row_index, ColumnDataType::Boolean);
569 }
570 writer.write_field_unchecked(
572 &key,
573 ColumnDataType::Boolean,
574 Some(ValueData::BoolValue(v)),
575 row,
576 );
577 }
578 Some(OtlpValue::IntValue(v)) => {
579 if let Some(batch_schema) = batch_schema.as_deref_mut() {
580 batch_schema.observe_value_type(&key, row_index, ColumnDataType::Int64);
581 }
582 writer.write_field_unchecked(
584 &key,
585 ColumnDataType::Int64,
586 Some(ValueData::I64Value(v)),
587 row,
588 );
589 }
590 Some(OtlpValue::DoubleValue(v)) => {
591 if let Some(batch_schema) = batch_schema.as_deref_mut() {
592 batch_schema.observe_value_type(&key, row_index, ColumnDataType::Float64);
593 }
594 writer.write_field_unchecked(
595 &key,
596 ColumnDataType::Float64,
597 Some(ValueData::F64Value(v)),
598 row,
599 );
600 }
601 Some(OtlpValue::ArrayValue(v)) => {
602 if let Some(batch_schema) = batch_schema.as_deref_mut() {
603 batch_schema.observe_binary_type(&key, row_index, TraceBinaryType::Json);
604 }
605 writer.write_column_unchecked(
606 row_writer::build_json_column_schema(key),
607 Some(ValueData::BinaryValue(
608 any_value_to_jsonb(OtlpValue::ArrayValue(v)).to_vec(),
609 )),
610 row,
611 );
612 }
613 Some(OtlpValue::KvlistValue(v)) => {
614 if let Some(batch_schema) = batch_schema.as_deref_mut() {
615 batch_schema.observe_binary_type(&key, row_index, TraceBinaryType::Json);
616 }
617 writer.write_column_unchecked(
618 row_writer::build_json_column_schema(key),
619 Some(ValueData::BinaryValue(
620 any_value_to_jsonb(OtlpValue::KvlistValue(v)).to_vec(),
621 )),
622 row,
623 );
624 }
625 Some(OtlpValue::BytesValue(v)) => {
626 if let Some(batch_schema) = batch_schema.as_deref_mut() {
627 batch_schema.observe_binary_type(&key, row_index, TraceBinaryType::Binary);
628 }
629 writer.write_field_unchecked(
630 key,
631 ColumnDataType::Binary,
632 Some(ValueData::BinaryValue(v)),
633 row,
634 );
635 }
636 None => {}
637 }
638 }
639
640 Ok(())
641}
642
643#[cfg(test)]
644mod tests {
645 use api::v1::value::ValueData;
646 use opentelemetry_proto::tonic::common::v1::any_value::Value as OtlpValue;
647 use opentelemetry_proto::tonic::common::v1::{AnyValue, ArrayValue, KeyValue};
648
649 use super::*;
650 use crate::otlp::trace::TraceAuxData;
651 use crate::otlp::trace::attributes::Attributes;
652 use crate::otlp::trace::span::{SpanEvents, SpanLinks};
653 use crate::row_writer::TableData;
654
655 fn make_kv(key: &str, value: OtlpValue) -> KeyValue {
656 KeyValue {
657 key: key.to_string(),
658 value: Some(AnyValue { value: Some(value) }),
659 }
660 }
661
662 fn make_span(service_name: &str, trace_id: &str, span_id: &str) -> TraceSpan {
663 TraceSpan {
664 service_name: Some(service_name.to_string()),
665 trace_id: trace_id.to_string(),
666 span_id: span_id.to_string(),
667 parent_span_id: None,
668 resource_attributes: Attributes::from(vec![]),
669 scope_name: "scope".to_string(),
670 scope_version: "v1".to_string(),
671 scope_attributes: Attributes::from(vec![]),
672 trace_state: String::new(),
673 span_name: "op".to_string(),
674 span_kind: "SPAN_KIND_SERVER".to_string(),
675 span_status_code: "STATUS_CODE_UNSET".to_string(),
676 span_status_message: String::new(),
677 span_attributes: Attributes::from(vec![]),
678 span_events: SpanEvents::from(vec![]),
679 span_links: SpanLinks::from(vec![]),
680 start_in_nanosecond: 1,
681 end_in_nanosecond: 2,
682 }
683 }
684
685 #[test]
686 fn test_fixed_trace_columns_keep_schema_and_values_aligned() {
687 let table_data = build_trace_table_data(&[make_span("svc", "trace", "span")]).unwrap();
688 let (schema, rows) = table_data.into_schema_and_rows();
689 let expected = [
690 (
691 TIMESTAMP_COLUMN,
692 Some(ValueData::TimestampNanosecondValue(1)),
693 ),
694 (
695 "timestamp_end",
696 Some(ValueData::TimestampNanosecondValue(2)),
697 ),
698 (DURATION_NANO_COLUMN, Some(ValueData::U64Value(1))),
699 (PARENT_SPAN_ID_COLUMN, None),
700 (
701 TRACE_ID_COLUMN,
702 Some(ValueData::StringValue("trace".to_string())),
703 ),
704 (
705 SPAN_ID_COLUMN,
706 Some(ValueData::StringValue("span".to_string())),
707 ),
708 (
709 SPAN_KIND_COLUMN,
710 Some(ValueData::StringValue("SPAN_KIND_SERVER".to_string())),
711 ),
712 (
713 SPAN_NAME_COLUMN,
714 Some(ValueData::StringValue("op".to_string())),
715 ),
716 (
717 SPAN_STATUS_CODE,
718 Some(ValueData::StringValue("STATUS_CODE_UNSET".to_string())),
719 ),
720 (
721 SPAN_STATUS_MESSAGE_COLUMN,
722 Some(ValueData::StringValue(String::new())),
723 ),
724 (
725 TRACE_STATE_COLUMN,
726 Some(ValueData::StringValue(String::new())),
727 ),
728 (
729 SCOPE_NAME_COLUMN,
730 Some(ValueData::StringValue("scope".to_string())),
731 ),
732 (
733 SCOPE_VERSION_COLUMN,
734 Some(ValueData::StringValue("v1".to_string())),
735 ),
736 ];
737
738 for (index, (name, value)) in expected.into_iter().enumerate() {
739 assert_eq!(schema[index].column_name, name);
740 assert_eq!(rows[0].values[index].value_data, value);
741 }
742 }
743
744 #[test]
745 fn test_batch_schema_preserves_column_order_and_attribute_types() {
746 let mut span1 = make_span("svc-a", "trace-a", "span-a");
747 span1.span_attributes = Attributes::from(vec![make_kv("val", OtlpValue::IntValue(10))]);
748 let mut span2 = make_span("svc-a", "trace-a", "span-b");
749 span2.span_attributes = Attributes::from(vec![make_kv("val", OtlpValue::DoubleValue(1.5))]);
750
751 let (table_data, batch_schema) =
752 build_trace_table_data_with_schema(vec![span1, span2].into_iter()).unwrap();
753
754 assert_eq!(
755 batch_schema.value_types("span_attributes.val").unwrap(),
756 [ColumnDataType::Int64, ColumnDataType::Float64]
757 );
758 let timestamp_index = table_data
759 .columns()
760 .iter()
761 .position(|column| column.column_name == TIMESTAMP_COLUMN)
762 .unwrap();
763 let attribute_index = table_data
764 .columns()
765 .iter()
766 .position(|column| column.column_name == "span_attributes.val")
767 .unwrap();
768 assert!(timestamp_index < attribute_index);
769 }
770
771 #[test]
772 fn test_optional_batch_schema_observation_preserves_rows() {
773 let mut span = make_span("svc-a", "trace-a", "span-a");
774 span.span_attributes = Attributes::from(vec![make_kv("val", OtlpValue::IntValue(10))]);
775
776 let without_schema = build_trace_table_data(std::slice::from_ref(&span)).unwrap();
777 let (with_schema, batch_schema) =
778 build_trace_table_data_with_schema(vec![span].into_iter()).unwrap();
779
780 assert_eq!(
781 without_schema.into_schema_and_rows(),
782 with_schema.into_schema_and_rows()
783 );
784 assert_eq!(
785 batch_schema.value_types("span_attributes.val").unwrap(),
786 [ColumnDataType::Int64]
787 );
788 let retry_columns = batch_schema.into_retry_columns();
789 let retry_column = &retry_columns["span_attributes.val"];
790 assert_eq!(retry_column.present_rows, [0]);
791 assert!(retry_column.binary_types.is_empty());
792 }
793
794 #[test]
795 fn test_batch_schema_tracks_binary_and_json_rows() {
796 let mut binary_span = make_span("svc-a", "trace-a", "span-a");
797 binary_span.span_attributes = Attributes::from(vec![make_kv(
798 "val",
799 OtlpValue::BytesValue(vec![1_u8, 2, 3]),
800 )]);
801 let mut json_span = make_span("svc-a", "trace-a", "span-b");
802 json_span.span_attributes = Attributes::from(vec![make_kv(
803 "val",
804 OtlpValue::ArrayValue(ArrayValue {
805 values: vec![AnyValue {
806 value: Some(OtlpValue::IntValue(1)),
807 }],
808 }),
809 )]);
810
811 let (_, batch_schema) =
812 build_trace_table_data_with_schema(vec![binary_span, json_span].into_iter()).unwrap();
813
814 assert!(batch_schema.has_incompatible_logical_types());
815 assert!(batch_schema.has_incompatible_logical_types_for("span_attributes.val"));
816 let retry_columns = batch_schema.into_retry_columns();
817 let retry_column = &retry_columns["span_attributes.val"];
818 assert_eq!(retry_column.present_rows, [0, 1]);
819 assert_eq!(
820 retry_column.binary_types,
821 [(0, TraceBinaryType::Binary), (1, TraceBinaryType::Json)]
822 );
823 }
824
825 #[test]
826 fn test_batch_schema_preserves_scalar_then_binary_values() {
827 let mut scalar_span = make_span("svc-a", "trace-a", "span-a");
828 scalar_span.span_attributes = Attributes::from(vec![
829 make_kv("bytes", OtlpValue::StringValue("text".to_string())),
830 make_kv("json", OtlpValue::StringValue("text".to_string())),
831 ]);
832 let mut binary_span = make_span("svc-a", "trace-a", "span-b");
833 binary_span.span_attributes = Attributes::from(vec![
834 make_kv("bytes", OtlpValue::BytesValue(vec![1_u8, 2, 3])),
835 make_kv(
836 "json",
837 OtlpValue::ArrayValue(ArrayValue {
838 values: vec![AnyValue {
839 value: Some(OtlpValue::IntValue(1)),
840 }],
841 }),
842 ),
843 ]);
844
845 let (_, batch_schema) =
846 build_trace_table_data_with_schema(vec![scalar_span, binary_span].into_iter()).unwrap();
847
848 assert_eq!(
849 batch_schema.value_types("span_attributes.bytes").unwrap(),
850 [ColumnDataType::String, ColumnDataType::Binary]
851 );
852 assert_eq!(
853 batch_schema.value_types("span_attributes.json").unwrap(),
854 [ColumnDataType::String, ColumnDataType::Binary]
855 );
856 let retry_columns = batch_schema.into_retry_columns();
857 assert_eq!(
858 retry_columns["span_attributes.bytes"].binary_types,
859 [(1, TraceBinaryType::Binary)]
860 );
861 assert_eq!(
862 retry_columns["span_attributes.json"].binary_types,
863 [(1, TraceBinaryType::Json)]
864 );
865 }
866
867 #[test]
868 fn test_keep_mixed_numeric_values_until_frontend_reconciliation() {
869 let mut writer = TableData::new(4, 2);
870
871 let attrs1 = Attributes::from(vec![make_kv("val", OtlpValue::DoubleValue(1.5))]);
872 let mut row1 = writer.alloc_one_row();
873 write_attributes(&mut writer, "attr", attrs1, &mut row1).unwrap();
874 writer.add_row(row1);
875
876 let attrs2 = Attributes::from(vec![make_kv("val", OtlpValue::IntValue(42))]);
877 let mut row2 = writer.alloc_one_row();
878 write_attributes(&mut writer, "attr", attrs2, &mut row2).unwrap();
879 writer.add_row(row2);
880
881 let (schema, rows) = writer.into_schema_and_rows();
882
883 let col_idx = schema
884 .iter()
885 .position(|c| c.column_name == "attr.val")
886 .unwrap();
887 assert_eq!(schema[col_idx].datatype, ColumnDataType::Float64 as i32);
888
889 assert_eq!(
890 rows[0].values[col_idx].value_data,
891 Some(ValueData::F64Value(1.5))
892 );
893 assert_eq!(
894 rows[1].values[col_idx].value_data,
895 Some(ValueData::I64Value(42))
896 );
897 }
898
899 #[test]
900 fn test_keep_mixed_string_and_int_values_until_frontend_reconciliation() {
901 let mut writer = TableData::new(4, 2);
902
903 let attrs1 = Attributes::from(vec![make_kv("val", OtlpValue::IntValue(10))]);
904 let mut row1 = writer.alloc_one_row();
905 write_attributes(&mut writer, "attr", attrs1, &mut row1).unwrap();
906 writer.add_row(row1);
907
908 let attrs2 = Attributes::from(vec![make_kv(
909 "val",
910 OtlpValue::StringValue("20".to_string()),
911 )]);
912 let mut row2 = writer.alloc_one_row();
913 write_attributes(&mut writer, "attr", attrs2, &mut row2).unwrap();
914 writer.add_row(row2);
915
916 let (schema, rows) = writer.into_schema_and_rows();
917 let col_idx = schema
918 .iter()
919 .position(|c| c.column_name == "attr.val")
920 .unwrap();
921 assert_eq!(schema[col_idx].datatype, ColumnDataType::Int64 as i32);
922 assert_eq!(
923 rows[1].values[col_idx].value_data,
924 Some(ValueData::StringValue("20".to_string()))
925 );
926 }
927
928 #[test]
929 fn test_keep_first_seen_schema_until_frontend_reconciliation() {
930 let mut writer = TableData::new(4, 2);
931
932 let attrs1 = Attributes::from(vec![make_kv(
933 "val",
934 OtlpValue::StringValue("10".to_string()),
935 )]);
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(20))]);
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 let col_idx = schema
947 .iter()
948 .position(|c| c.column_name == "attr.val")
949 .unwrap();
950 assert_eq!(schema[col_idx].datatype, ColumnDataType::String as i32);
951 assert_eq!(
952 rows[0].values[col_idx].value_data,
953 Some(ValueData::StringValue("10".to_string()))
954 );
955 assert_eq!(
956 rows[1].values[col_idx].value_data,
957 Some(ValueData::I64Value(20))
958 );
959 }
960
961 #[test]
962 fn test_keep_mixed_string_and_float_values_until_frontend_reconciliation() {
963 let mut writer = TableData::new(4, 2);
964
965 let attrs1 = Attributes::from(vec![make_kv("val", OtlpValue::DoubleValue(1.5))]);
966 let mut row1 = writer.alloc_one_row();
967 write_attributes(&mut writer, "attr", attrs1, &mut row1).unwrap();
968 writer.add_row(row1);
969
970 let attrs2 = Attributes::from(vec![make_kv(
971 "val",
972 OtlpValue::StringValue("1.5".to_string()),
973 )]);
974 let mut row2 = writer.alloc_one_row();
975 write_attributes(&mut writer, "attr", attrs2, &mut row2).unwrap();
976 writer.add_row(row2);
977
978 let (schema, rows) = writer.into_schema_and_rows();
979 let col_idx = schema
980 .iter()
981 .position(|c| c.column_name == "attr.val")
982 .unwrap();
983 assert_eq!(schema[col_idx].datatype, ColumnDataType::Float64 as i32);
984 assert_eq!(
985 rows[1].values[col_idx].value_data,
986 Some(ValueData::StringValue("1.5".to_string()))
987 );
988 }
989
990 #[test]
991 fn test_keep_mixed_string_and_bool_values_until_frontend_reconciliation() {
992 let mut writer = TableData::new(4, 2);
993
994 let attrs1 = Attributes::from(vec![make_kv(
995 "val",
996 OtlpValue::StringValue("true".to_string()),
997 )]);
998 let mut row1 = writer.alloc_one_row();
999 write_attributes(&mut writer, "attr", attrs1, &mut row1).unwrap();
1000 writer.add_row(row1);
1001
1002 let attrs2 = Attributes::from(vec![make_kv("val", OtlpValue::BoolValue(false))]);
1003 let mut row2 = writer.alloc_one_row();
1004 write_attributes(&mut writer, "attr", attrs2, &mut row2).unwrap();
1005 writer.add_row(row2);
1006
1007 let (schema, rows) = writer.into_schema_and_rows();
1008 let col_idx = schema
1009 .iter()
1010 .position(|c| c.column_name == "attr.val")
1011 .unwrap();
1012 assert_eq!(schema[col_idx].datatype, ColumnDataType::String as i32);
1013 assert_eq!(
1014 rows[0].values[col_idx].value_data,
1015 Some(ValueData::StringValue("true".to_string()))
1016 );
1017 assert_eq!(
1018 rows[1].values[col_idx].value_data,
1019 Some(ValueData::BoolValue(false))
1020 );
1021 }
1022
1023 #[test]
1024 fn test_keep_mixed_binary_and_string_values_until_frontend_reconciliation() {
1025 let mut writer = TableData::new(4, 2);
1026
1027 let attrs1 = Attributes::from(vec![make_kv(
1028 "val",
1029 OtlpValue::BytesValue(vec![1_u8, 2, 3]),
1030 )]);
1031 let mut row1 = writer.alloc_one_row();
1032 write_attributes(&mut writer, "attr", attrs1, &mut row1).unwrap();
1033 writer.add_row(row1);
1034
1035 let attrs2 = Attributes::from(vec![make_kv(
1036 "val",
1037 OtlpValue::StringValue("false".to_string()),
1038 )]);
1039 let mut row2 = writer.alloc_one_row();
1040 write_attributes(&mut writer, "attr", attrs2, &mut row2).unwrap();
1041 writer.add_row(row2);
1042
1043 let (schema, rows) = writer.into_schema_and_rows();
1044 let col_idx = schema
1045 .iter()
1046 .position(|c| c.column_name == "attr.val")
1047 .unwrap();
1048 assert_eq!(schema[col_idx].datatype, ColumnDataType::Binary as i32);
1049 assert_eq!(
1050 rows[0].values[col_idx].value_data,
1051 Some(ValueData::BinaryValue(vec![1_u8, 2, 3]))
1052 );
1053 assert_eq!(
1054 rows[1].values[col_idx].value_data,
1055 Some(ValueData::StringValue("false".to_string()))
1056 );
1057 }
1058
1059 #[test]
1060 fn test_keep_mixed_binary_and_json_values_until_frontend_reconciliation() {
1061 let mut writer = TableData::new(4, 2);
1062
1063 let attrs1 = Attributes::from(vec![make_kv(
1064 "val",
1065 OtlpValue::BytesValue(vec![1_u8, 2, 3]),
1066 )]);
1067 let mut row1 = writer.alloc_one_row();
1068 write_attributes(&mut writer, "attr", attrs1, &mut row1).unwrap();
1069 writer.add_row(row1);
1070
1071 let attrs2 = Attributes::from(vec![make_kv(
1072 "val",
1073 OtlpValue::ArrayValue(ArrayValue {
1074 values: vec![AnyValue {
1075 value: Some(OtlpValue::IntValue(1)),
1076 }],
1077 }),
1078 )]);
1079 let mut row2 = writer.alloc_one_row();
1080 write_attributes(&mut writer, "attr", attrs2, &mut row2).unwrap();
1081 writer.add_row(row2);
1082
1083 let (schema, rows) = writer.into_schema_and_rows();
1084 let col_idx = schema
1085 .iter()
1086 .position(|c| c.column_name == "attr.val")
1087 .unwrap();
1088 assert_eq!(schema[col_idx].datatype, ColumnDataType::Binary as i32);
1089 assert!(matches!(
1090 rows[1].values[col_idx].value_data.as_ref(),
1091 Some(ValueData::BinaryValue(_))
1092 ));
1093 }
1094
1095 #[test]
1096 fn test_build_aux_table_requests_deduplicates_services_and_operations() {
1097 let spans = vec![
1098 make_span("svc-a", "trace-a", "span-a"),
1099 make_span("svc-a", "trace-b", "span-b"),
1100 ];
1101 let mut aux_data = TraceAuxData::default();
1102 for span in &spans {
1103 aux_data.observe_span(span);
1104 }
1105
1106 let (requests, total_rows) =
1107 build_aux_table_requests(aux_data, "opentelemetry_traces").unwrap();
1108 assert_eq!(requests.inserts.len(), 2);
1109 assert_eq!(total_rows, 2);
1110 }
1111 }