Skip to main content

servers/otlp/trace/
v0.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::HashSet;
16
17use api::v1::value::ValueData;
18use api::v1::{ColumnDataType, RowInsertRequests};
19use common_catalog::consts::{trace_operations_table_name, trace_services_table_name};
20use common_grpc::precision::Precision;
21use pipeline::{GreptimePipelineParams, PipelineWay};
22use session::context::QueryContextRef;
23
24use crate::error::Result;
25use crate::otlp::trace::span::TraceSpan;
26use crate::otlp::trace::{
27    DURATION_NANO_COLUMN, PARENT_SPAN_ID_COLUMN, SERVICE_NAME_COLUMN, SPAN_ATTRIBUTES_COLUMN,
28    SPAN_EVENTS_COLUMN, SPAN_ID_COLUMN, SPAN_KIND_COLUMN, SPAN_NAME_COLUMN, SPAN_STATUS_CODE,
29    SPAN_STATUS_MESSAGE_COLUMN, TIMESTAMP_COLUMN, TRACE_ID_COLUMN, TRACE_STATE_COLUMN,
30    TraceAuxData,
31};
32use crate::otlp::utils::{make_column_data, make_string_column_data};
33use crate::query_handler::PipelineHandlerRef;
34use crate::row_writer::{self, MultiTableData, TableData};
35
36const APPROXIMATE_COLUMN_COUNT: usize = 24;
37
38// Use a timestamp(2100-01-01 00:00:00) as large as possible.
39const MAX_TIMESTAMP: i64 = 4102444800000000000;
40
41/// Converts trace spans into row insert requests for the main v0 trace table.
42///
43/// Auxiliary service and operation table writes are built separately so the
44/// caller can update them only after the main span write succeeds.
45pub fn v0_to_grpc_main_insert_requests(
46    spans: &[TraceSpan],
47    _pipeline: &PipelineWay,
48    _pipeline_params: &GreptimePipelineParams,
49    table_name: &str,
50    _query_ctx: &QueryContextRef,
51    _pipeline_handler: PipelineHandlerRef,
52) -> Result<(RowInsertRequests, usize)> {
53    let mut multi_table_writer = MultiTableData::default();
54    let trace_writer = build_trace_table_data(spans)?;
55    multi_table_writer.add_table_data(table_name, trace_writer);
56
57    Ok(multi_table_writer.into_row_insert_requests())
58}
59
60/// Builds the row-oriented payload for the main v0 trace table.
61pub fn build_trace_table_data(spans: &[TraceSpan]) -> Result<TableData> {
62    let mut trace_writer = TableData::new(APPROXIMATE_COLUMN_COUNT, spans.len());
63    for span in spans.iter().cloned() {
64        write_span_to_row(&mut trace_writer, span)?;
65    }
66
67    Ok(trace_writer)
68}
69
70/// Builds row insert requests for the v0 trace auxiliary tables.
71pub fn build_aux_table_requests(
72    aux_data: TraceAuxData,
73    table_name: &str,
74) -> Result<(RowInsertRequests, usize)> {
75    let mut multi_table_writer = MultiTableData::default();
76    let mut trace_services_writer = TableData::new(APPROXIMATE_COLUMN_COUNT, 1);
77    let mut trace_operations_writer = TableData::new(APPROXIMATE_COLUMN_COUNT, 1);
78
79    write_trace_services_to_row(&mut trace_services_writer, aux_data.services)?;
80    write_trace_operations_to_row(&mut trace_operations_writer, aux_data.operations)?;
81
82    multi_table_writer.add_table_data(trace_services_table_name(table_name), trace_services_writer);
83    multi_table_writer.add_table_data(
84        trace_operations_table_name(table_name),
85        trace_operations_writer,
86    );
87    Ok(multi_table_writer.into_row_insert_requests())
88}
89
90pub fn write_span_to_row(writer: &mut TableData, span: TraceSpan) -> Result<()> {
91    let mut row = writer.alloc_one_row();
92
93    // write ts
94    row_writer::write_ts_to_nanos(
95        writer,
96        TIMESTAMP_COLUMN,
97        Some(span.start_in_nanosecond as i64),
98        Precision::Nanosecond,
99        &mut row,
100    )?;
101
102    // write fields
103    let fields = vec![
104        make_column_data(
105            "timestamp_end",
106            ColumnDataType::TimestampNanosecond,
107            Some(ValueData::TimestampNanosecondValue(
108                span.end_in_nanosecond as i64,
109            )),
110        ),
111        // The v0 data model is frozen: `duration_nano` stays UInt64. The
112        // signed→unsigned ingest compatibility layer only runs on the v1 path
113        // (see `trace_ingest.rs`), so flipping v0 here would break writes into
114        // every pre-existing v0 table at mito's schema check. New v1 tables
115        // are signed; v0 tables keep the legacy unsigned column indefinitely.
116        make_column_data(
117            DURATION_NANO_COLUMN,
118            ColumnDataType::Uint64,
119            Some(ValueData::U64Value(
120                span.end_in_nanosecond - span.start_in_nanosecond,
121            )),
122        ),
123        make_string_column_data(TRACE_ID_COLUMN, Some(span.trace_id)),
124        make_string_column_data(SPAN_ID_COLUMN, Some(span.span_id)),
125        make_string_column_data(PARENT_SPAN_ID_COLUMN, span.parent_span_id),
126        make_string_column_data(SPAN_KIND_COLUMN, Some(span.span_kind)),
127        make_string_column_data(SPAN_NAME_COLUMN, Some(span.span_name)),
128        make_string_column_data(SPAN_STATUS_CODE, Some(span.span_status_code)),
129        make_string_column_data(SPAN_STATUS_MESSAGE_COLUMN, Some(span.span_status_message)),
130        make_string_column_data(TRACE_STATE_COLUMN, Some(span.trace_state)),
131    ];
132    row_writer::write_fields(writer, fields.into_iter(), &mut row)?;
133
134    if let Some(service_name) = span.service_name {
135        row_writer::write_tag(writer, SERVICE_NAME_COLUMN, service_name, &mut row)?;
136    }
137
138    row_writer::write_json(
139        writer,
140        SPAN_ATTRIBUTES_COLUMN,
141        span.span_attributes.into(),
142        &mut row,
143    )?;
144    row_writer::write_json(
145        writer,
146        SPAN_EVENTS_COLUMN,
147        span.span_events.into(),
148        &mut row,
149    )?;
150    row_writer::write_json(writer, "span_links", span.span_links.into(), &mut row)?;
151
152    // write fields
153    let fields = vec![
154        make_string_column_data("scope_name", Some(span.scope_name)),
155        make_string_column_data("scope_version", Some(span.scope_version)),
156    ];
157    row_writer::write_fields(writer, fields.into_iter(), &mut row)?;
158
159    row_writer::write_json(
160        writer,
161        "scope_attributes",
162        span.scope_attributes.into(),
163        &mut row,
164    )?;
165
166    row_writer::write_json(
167        writer,
168        "resource_attributes",
169        span.resource_attributes.into(),
170        &mut row,
171    )?;
172
173    writer.add_row(row);
174
175    Ok(())
176}
177
178fn write_trace_services_to_row(writer: &mut TableData, services: HashSet<String>) -> Result<()> {
179    for service_name in services {
180        let mut row = writer.alloc_one_row();
181        // Write the timestamp as 0.
182        row_writer::write_ts_to_nanos(
183            writer,
184            TIMESTAMP_COLUMN,
185            Some(MAX_TIMESTAMP),
186            Precision::Nanosecond,
187            &mut row,
188        )?;
189
190        // Write the `service_name` column.
191        row_writer::write_fields(
192            writer,
193            std::iter::once(make_string_column_data(
194                SERVICE_NAME_COLUMN,
195                Some(service_name),
196            )),
197            &mut row,
198        )?;
199        writer.add_row(row);
200    }
201
202    Ok(())
203}
204
205fn write_trace_operations_to_row(
206    writer: &mut TableData,
207    operations: HashSet<(String, String, String)>,
208) -> Result<()> {
209    for (service_name, span_name, span_kind) in operations {
210        let mut row = writer.alloc_one_row();
211        // Write the timestamp as 0.
212        row_writer::write_ts_to_nanos(
213            writer,
214            TIMESTAMP_COLUMN,
215            Some(MAX_TIMESTAMP),
216            Precision::Nanosecond,
217            &mut row,
218        )?;
219
220        // Write the `service_name`, `span_name`, and `span_kind` columns.
221        row_writer::write_fields(
222            writer,
223            vec![
224                make_string_column_data(SERVICE_NAME_COLUMN, Some(service_name)),
225                make_string_column_data(SPAN_NAME_COLUMN, Some(span_name)),
226                make_string_column_data(SPAN_KIND_COLUMN, Some(span_kind)),
227            ]
228            .into_iter(),
229            &mut row,
230        )?;
231        writer.add_row(row);
232    }
233
234    Ok(())
235}
236
237#[cfg(test)]
238mod tests {
239    use api::v1::ColumnDataType;
240    use api::v1::value::ValueData;
241
242    use super::{build_aux_table_requests, build_trace_table_data};
243    use crate::otlp::trace::attributes::Attributes;
244    use crate::otlp::trace::span::{SpanEvents, SpanLinks, TraceSpan};
245    use crate::otlp::trace::{DURATION_NANO_COLUMN, TraceAuxData};
246
247    fn make_span(service_name: &str, trace_id: &str, span_id: &str) -> TraceSpan {
248        TraceSpan {
249            service_name: Some(service_name.to_string()),
250            trace_id: trace_id.to_string(),
251            span_id: span_id.to_string(),
252            parent_span_id: None,
253            resource_attributes: Attributes::from(vec![]),
254            scope_name: "scope".to_string(),
255            scope_version: "v1".to_string(),
256            scope_attributes: Attributes::from(vec![]),
257            trace_state: String::new(),
258            span_name: "op".to_string(),
259            span_kind: "SPAN_KIND_SERVER".to_string(),
260            span_status_code: "STATUS_CODE_UNSET".to_string(),
261            span_status_message: String::new(),
262            span_attributes: Attributes::from(vec![]),
263            span_events: SpanEvents::from(vec![]),
264            span_links: SpanLinks::from(vec![]),
265            start_in_nanosecond: 1,
266            end_in_nanosecond: 2,
267        }
268    }
269
270    #[test]
271    fn test_build_trace_table_data_from_span_subset() {
272        let spans = [
273            make_span("svc-a", "trace-a", "span-a"),
274            make_span("svc-b", "trace-b", "span-b"),
275        ];
276
277        let writer = build_trace_table_data(&spans[..1]).unwrap();
278        let (_, rows) = writer.into_schema_and_rows();
279        assert_eq!(rows.len(), 1);
280    }
281
282    #[test]
283    fn test_v0_duration_nano_stays_uint64() {
284        // The v0 data model is frozen on the unsigned `duration_nano`. This
285        // pins the schema so the unsigned→signed transition of the built-in
286        // models (which flips only v1) cannot silently leak into v0 and break
287        // writes into pre-existing v0 tables.
288        let writer = build_trace_table_data(&[make_span("svc-a", "trace-a", "span-a")]).unwrap();
289        let (schema, rows) = writer.into_schema_and_rows();
290
291        let idx = schema
292            .iter()
293            .position(|c| c.column_name == DURATION_NANO_COLUMN)
294            .unwrap();
295        assert_eq!(schema[idx].datatype, ColumnDataType::Uint64 as i32);
296        assert_eq!(rows[0].values[idx].value_data, Some(ValueData::U64Value(1)));
297    }
298
299    #[test]
300    fn test_build_aux_table_requests_deduplicates_services_and_operations() {
301        let spans = vec![
302            make_span("svc-a", "trace-a", "span-a"),
303            make_span("svc-a", "trace-b", "span-b"),
304        ];
305        let mut aux_data = TraceAuxData::default();
306        for span in &spans {
307            aux_data.observe_span(span);
308        }
309
310        let (requests, total_rows) =
311            build_aux_table_requests(aux_data, "opentelemetry_traces").unwrap();
312        assert_eq!(requests.inserts.len(), 2);
313        assert_eq!(total_rows, 2);
314    }
315}