1use 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
31const 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
109pub(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
119fn 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}