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