Skip to main content

servers/otlp/trace/
v1.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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
45// Use a timestamp(2100-01-01 00:00:00) as large as possible.
46const 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/// Distinguishes raw bytes from JSONB values that both use a binary protobuf value.
98#[derive(Debug, Clone, Copy, PartialEq, Eq)]
99pub enum TraceBinaryType {
100    /// Raw OTLP bytes with no datatype extension.
101    Binary,
102    /// An OTLP array or key-value list encoded with the JSONB extension.
103    Json,
104}
105
106impl TraceBinaryType {
107    /// Applies this logical binary type to a column schema.
108    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/// Per-column observations collected while building one trace chunk.
120#[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    /// Records a row once while preserving encounter order.
129    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    /// Returns whether raw binary and JSONB values share this column.
136    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
146/// Sparse column metadata retained for span-level fallback.
147pub struct TraceRetryColumn {
148    /// Row indexes that contain this dynamic column.
149    pub present_rows: Vec<usize>,
150    /// Row-specific binary kinds needed to restore the logical schema.
151    pub binary_types: Vec<(usize, TraceBinaryType)>,
152}
153
154/// Retry metadata keyed by dynamic trace column name.
155pub type TraceRetryColumns = HashMap<String, TraceRetryColumn>;
156
157/// Schema observations for dynamic columns in one converted trace chunk.
158#[derive(Default)]
159pub struct TraceBatchSchema {
160    columns: HashMap<String, TraceBatchColumnSchema>,
161}
162
163impl TraceBatchSchema {
164    /// Records a scalar value type and the row where the column is present.
165    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    /// Records the logical kind of a binary protobuf value for one row.
174    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    /// Records a sparse column occurrence that has no dynamic value type.
190    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    /// Returns the distinct value types in first-seen order for a column.
198    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    /// Returns whether any column mixes raw binary and JSONB values.
205    pub fn has_incompatible_logical_types(&self) -> bool {
206        self.columns
207            .values()
208            .any(TraceBatchColumnSchema::has_incompatible_logical_types)
209    }
210
211    /// Returns whether the named column mixes raw binary and JSONB values.
212    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    /// Converts batch observations into sparse metadata for per-span retries.
219    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
246/// Converts trace spans into row insert requests for the main v1 trace table.
247///
248/// Auxiliary service and operation table writes are built separately so the
249/// caller can update them only after the main span write succeeds.
250pub 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
262/// Converts owned spans into unpadded main-table rows and schema observations.
263pub 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
269/// Builds the main-table request without collecting batch schema observations.
270fn 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
281/// Builds the row-oriented payload for the main v1 trace table.
282pub fn build_trace_table_data(spans: &[TraceSpan]) -> Result<TableData> {
283    build_trace_table_data_from_iter(spans.iter().cloned(), None)
284}
285
286/// Builds trace rows while collecting dynamic column observations.
287fn 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
295/// Shared row builder with optional batch schema observation.
296fn 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
318/// Builds row insert requests for the v1 trace auxiliary tables.
319pub 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
345/// Writes one span and optionally records its dynamic columns for reconciliation.
346fn 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        // Write the timestamp as 0.
471        row_writer::write_ts_to_nanos(
472            writer,
473            TIMESTAMP_COLUMN,
474            Some(MAX_TIMESTAMP),
475            Precision::Nanosecond,
476            &mut row,
477        )?;
478
479        // Write the `service_name` column.
480        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        // Write the timestamp as 0.
498        row_writer::write_ts_to_nanos(
499            writer,
500            TIMESTAMP_COLUMN,
501            Some(MAX_TIMESTAMP),
502            Precision::Nanosecond,
503            &mut row,
504        )?;
505
506        // Write the `service_name`, `span_name`, and `span_kind` columns as tags.
507        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
534/// Writes flattened attributes without coercion and optionally records their actual types.
535fn 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        // skip resource_attributes.service.name because its already copied to
546        // top level as `SERVICE_NAME_COLUMN`
547        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                // Keep the raw request value here. Mixed trace types are reconciled later
558                // in the frontend once we can also see the existing table schema.
559                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                // Do not coerce or promote types while building the request-local rows.
571                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                // Preserving the original value avoids order-dependent behavior inside one batch.
583                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    // Conversion matrix coverage lives in the shared coercion helper tests.
1112}