Skip to main content

servers/otlp/trace/
v2.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 api::v1::value::ValueData;
16use api::v1::{ColumnDataType, ColumnSchema, RowInsertRequests, SemanticType};
17use snafu::ensure;
18
19use crate::error::{Result, TimestampOverflowSnafu};
20use crate::otlp::trace::span::TraceSpan;
21use crate::otlp::trace::v1::span_duration_nano;
22use crate::otlp::trace::{
23    DURATION_NANO_COLUMN, PARENT_SPAN_ID_COLUMN, RESOURCE_ATTRIBUTES_COLUMN,
24    SCOPE_ATTRIBUTES_COLUMN, SCOPE_NAME_COLUMN, SCOPE_VERSION_COLUMN, SERVICE_NAME_COLUMN,
25    SPAN_ATTRIBUTES_COLUMN, SPAN_EVENTS_COLUMN, SPAN_ID_COLUMN, SPAN_KIND_COLUMN,
26    SPAN_LINKS_COLUMN, SPAN_NAME_COLUMN, SPAN_STATUS_CODE, SPAN_STATUS_MESSAGE_COLUMN,
27    TIMESTAMP_COLUMN, TIMESTAMP_END_COLUMN, TRACE_ID_COLUMN, TRACE_STATE_COLUMN,
28};
29use crate::row_writer::{self, MultiTableData, TableData};
30
31// Preallocate for the fixed v2 schema: 3 timing columns, 3 span/trace IDs,
32// 2 kind/name columns, 2 status columns, 1 trace state, 2 scope name/version
33// columns, 1 service name, 3 JSON2 attribute columns, and 2 JSON event/link columns.
34const APPROXIMATE_COLUMN_COUNT: usize = 19;
35
36struct FixedTraceColumnIndexes {
37    timestamp: usize,
38    timestamp_end: usize,
39    duration_nano: usize,
40    parent_span_id: usize,
41    trace_id: usize,
42    span_id: usize,
43    span_kind: usize,
44    span_name: usize,
45    span_status_code: usize,
46    span_status_message: usize,
47    trace_state: usize,
48    scope_name: usize,
49    scope_version: usize,
50    service_name: usize,
51    span_attributes: usize,
52    scope_attributes: usize,
53    resource_attributes: usize,
54    span_events: usize,
55    span_links: usize,
56}
57
58impl FixedTraceColumnIndexes {
59    fn resolve(writer: &mut TableData) -> Result<Self> {
60        let timestamp = writer.ensure_column(api::v1::helper::time_index_column_schema(
61            TIMESTAMP_COLUMN,
62            ColumnDataType::TimestampNanosecond,
63        ))?;
64        let mut field = |name: &str, datatype: ColumnDataType| {
65            writer.ensure_column(ColumnSchema {
66                column_name: name.to_string(),
67                datatype: datatype as i32,
68                semantic_type: SemanticType::Field as i32,
69                ..Default::default()
70            })
71        };
72        Ok(Self {
73            timestamp,
74            timestamp_end: field(TIMESTAMP_END_COLUMN, ColumnDataType::TimestampNanosecond)?,
75            duration_nano: field(DURATION_NANO_COLUMN, ColumnDataType::Int64)?,
76            parent_span_id: field(PARENT_SPAN_ID_COLUMN, ColumnDataType::String)?,
77            trace_id: field(TRACE_ID_COLUMN, ColumnDataType::String)?,
78            span_id: field(SPAN_ID_COLUMN, ColumnDataType::String)?,
79            span_kind: field(SPAN_KIND_COLUMN, ColumnDataType::String)?,
80            span_name: field(SPAN_NAME_COLUMN, ColumnDataType::String)?,
81            span_status_code: field(SPAN_STATUS_CODE, ColumnDataType::String)?,
82            span_status_message: field(SPAN_STATUS_MESSAGE_COLUMN, ColumnDataType::String)?,
83            trace_state: field(TRACE_STATE_COLUMN, ColumnDataType::String)?,
84            scope_name: field(SCOPE_NAME_COLUMN, ColumnDataType::String)?,
85            scope_version: field(SCOPE_VERSION_COLUMN, ColumnDataType::String)?,
86            service_name: writer.ensure_column(ColumnSchema {
87                column_name: SERVICE_NAME_COLUMN.to_string(),
88                datatype: ColumnDataType::String as i32,
89                semantic_type: SemanticType::Tag as i32,
90                ..Default::default()
91            })?,
92            span_attributes: writer.ensure_column(row_writer::build_json2_column_schema(
93                SPAN_ATTRIBUTES_COLUMN,
94            ))?,
95            scope_attributes: writer.ensure_column(row_writer::build_json2_column_schema(
96                SCOPE_ATTRIBUTES_COLUMN,
97            ))?,
98            resource_attributes: writer.ensure_column(row_writer::build_json2_column_schema(
99                RESOURCE_ATTRIBUTES_COLUMN,
100            ))?,
101            span_events: writer
102                .ensure_column(row_writer::build_json_column_schema(SPAN_EVENTS_COLUMN))?,
103            span_links: writer
104                .ensure_column(row_writer::build_json_column_schema(SPAN_LINKS_COLUMN))?,
105        })
106    }
107}
108
109/// Converts trace spans into row insert requests for the main v2 trace table.
110pub(super) fn v2_to_grpc_main_insert_requests(
111    spans: &[TraceSpan],
112    table_name: &str,
113) -> Result<(RowInsertRequests, usize)> {
114    let mut tables = MultiTableData::default();
115    tables.add_table_data(table_name, build_trace_table_data(spans)?);
116    Ok(tables.into_row_insert_requests())
117}
118
119/// Builds the fixed row-oriented payload for the main v2 trace table.
120fn build_trace_table_data(spans: &[TraceSpan]) -> Result<TableData> {
121    let mut writer = TableData::new(APPROXIMATE_COLUMN_COUNT, spans.len());
122    if spans.is_empty() {
123        return Ok(writer);
124    }
125    let columns = FixedTraceColumnIndexes::resolve(&mut writer)?;
126    for span in spans {
127        write_span_to_row(&mut writer, span, &columns)?;
128    }
129    Ok(writer)
130}
131
132fn write_span_to_row(
133    writer: &mut TableData,
134    span: &TraceSpan,
135    columns: &FixedTraceColumnIndexes,
136) -> Result<()> {
137    ensure!(
138        span.start_in_nanosecond <= i64::MAX as u64,
139        TimestampOverflowSnafu {
140            error: "`span.start_in_nanosecond`",
141        }
142    );
143    ensure!(
144        span.end_in_nanosecond <= i64::MAX as u64,
145        TimestampOverflowSnafu {
146            error: "`span.end_in_nanosecond`",
147        }
148    );
149    let mut row = writer.alloc_one_row();
150
151    let duration = span_duration_nano(span);
152    for (index, value) in [
153        (
154            columns.timestamp,
155            Some(ValueData::TimestampNanosecondValue(
156                span.start_in_nanosecond as i64,
157            )),
158        ),
159        (
160            columns.timestamp_end,
161            Some(ValueData::TimestampNanosecondValue(
162                span.end_in_nanosecond as i64,
163            )),
164        ),
165        (columns.duration_nano, Some(ValueData::I64Value(duration))),
166        (
167            columns.parent_span_id,
168            span.parent_span_id.clone().map(ValueData::StringValue),
169        ),
170        (
171            columns.trace_id,
172            Some(ValueData::StringValue(span.trace_id.clone())),
173        ),
174        (
175            columns.span_id,
176            Some(ValueData::StringValue(span.span_id.clone())),
177        ),
178        (
179            columns.span_kind,
180            Some(ValueData::StringValue(span.span_kind.clone())),
181        ),
182        (
183            columns.span_name,
184            Some(ValueData::StringValue(span.span_name.clone())),
185        ),
186        (
187            columns.span_status_code,
188            Some(ValueData::StringValue(span.span_status_code.clone())),
189        ),
190        (
191            columns.span_status_message,
192            Some(ValueData::StringValue(span.span_status_message.clone())),
193        ),
194        (
195            columns.trace_state,
196            Some(ValueData::StringValue(span.trace_state.clone())),
197        ),
198        (
199            columns.scope_name,
200            Some(ValueData::StringValue(span.scope_name.clone())),
201        ),
202        (
203            columns.scope_version,
204            Some(ValueData::StringValue(span.scope_version.clone())),
205        ),
206        (
207            columns.service_name,
208            span.service_name.clone().map(ValueData::StringValue),
209        ),
210        (
211            columns.span_attributes,
212            Some(row_writer::encode_json2(&span.span_attributes)?),
213        ),
214        (
215            columns.scope_attributes,
216            Some(row_writer::encode_json2(span.scope_attributes.as_ref())?),
217        ),
218        (
219            columns.resource_attributes,
220            Some(row_writer::encode_json2(span.resource_attributes.as_ref())?),
221        ),
222        (
223            columns.span_events,
224            Some(ValueData::BinaryValue(
225                jsonb::Value::from(span.span_events.clone()).to_vec(),
226            )),
227        ),
228        (
229            columns.span_links,
230            Some(ValueData::BinaryValue(
231                jsonb::Value::from(span.span_links.clone()).to_vec(),
232            )),
233        ),
234    ] {
235        row[index].value_data = value;
236    }
237
238    writer.add_row(row);
239    Ok(())
240}