1use std::sync::Arc;
16
17use api::v1::ColumnDataType;
18use api::v1::value::ValueData;
19use chrono::{DateTime, Utc};
20use common_time::Timestamp;
21use common_time::timestamp::TimeUnit;
22use datatypes::timestamp::TimestampNanosecond;
23use itertools::Itertools;
24use session::context::Channel;
25use snafu::{OptionExt, ensure};
26use util::to_pipeline_version;
27use vrl::value::Value as VrlValue;
28
29use crate::error::{
30 CastTypeSnafu, InvalidCustomTimeIndexSnafu, InvalidTimestampSnafu, PipelineMissingSnafu, Result,
31};
32use crate::etl::value::{MS_RESOLUTION, NS_RESOLUTION, S_RESOLUTION, US_RESOLUTION};
33use crate::table::PipelineTable;
34use crate::{GreptimePipelineParams, Pipeline};
35
36mod pipeline_cache;
37pub mod pipeline_operator;
38pub mod table;
39pub mod util;
40
41pub type PipelineVersion = Option<TimestampNanosecond>;
47
48pub type PipelineInfo = (Timestamp, PipelineRef);
50
51pub type PipelineTableRef = Arc<PipelineTable>;
52pub type PipelineRef = Arc<Pipeline>;
53
54#[derive(Default)]
57pub struct SelectInfo {
58 pub keys: Vec<String>,
59}
60
61impl From<String> for SelectInfo {
66 fn from(value: String) -> Self {
67 let mut keys: Vec<String> = value.split(',').map(|s| s.to_string()).sorted().collect();
68 keys.dedup();
69
70 SelectInfo { keys }
71 }
72}
73
74impl SelectInfo {
75 pub fn is_empty(&self) -> bool {
76 self.keys.is_empty()
77 }
78}
79
80pub const GREPTIME_INTERNAL_IDENTITY_PIPELINE_NAME: &str = "greptime_identity";
81pub const GREPTIME_INTERNAL_TRACE_PIPELINE_V0_NAME: &str = "greptime_trace_v0";
82pub const GREPTIME_INTERNAL_TRACE_PIPELINE_V1_NAME: &str = "greptime_trace_v1";
83pub const GREPTIME_INTERNAL_TRACE_PIPELINE_V2_NAME: &str = "greptime_trace_v2";
85
86#[derive(Debug, Clone)]
89pub enum PipelineDefinition {
90 Resolved(Arc<Pipeline>),
91 ByNameAndValue((String, PipelineVersion)),
92 GreptimeIdentityPipeline(Option<IdentityTimeIndex>),
93}
94
95impl PipelineDefinition {
96 pub fn from_name(
97 name: &str,
98 version: PipelineVersion,
99 custom_time_index: Option<(String, bool)>,
100 ) -> Result<Self> {
101 if name == GREPTIME_INTERNAL_IDENTITY_PIPELINE_NAME {
102 Ok(Self::GreptimeIdentityPipeline(
103 custom_time_index
104 .map(|(config, ignore_errors)| {
105 IdentityTimeIndex::from_config(config, ignore_errors)
106 })
107 .transpose()?,
108 ))
109 } else {
110 Ok(Self::ByNameAndValue((name.to_owned(), version)))
111 }
112 }
113
114 pub fn is_identity(&self) -> bool {
115 matches!(self, Self::GreptimeIdentityPipeline(_))
116 }
117
118 pub fn get_custom_ts(&self) -> Option<&IdentityTimeIndex> {
119 if let Self::GreptimeIdentityPipeline(custom_ts) = self {
120 custom_ts.as_ref()
121 } else {
122 None
123 }
124 }
125}
126
127pub struct PipelineContext<'a> {
128 pub pipeline_definition: &'a PipelineDefinition,
129 pub pipeline_param: &'a GreptimePipelineParams,
130 pub channel: Channel,
131}
132
133impl<'a> PipelineContext<'a> {
134 pub fn new(
135 pipeline_definition: &'a PipelineDefinition,
136 pipeline_param: &'a GreptimePipelineParams,
137 channel: Channel,
138 ) -> Self {
139 Self {
140 pipeline_definition,
141 pipeline_param,
142 channel,
143 }
144 }
145}
146pub enum PipelineWay {
147 OtlpLogDirect(Box<SelectInfo>),
148 Pipeline(PipelineDefinition),
149 OtlpTraceDirectV0,
150 OtlpTraceDirectV1,
151 OtlpTraceDirectV2,
152}
153
154impl PipelineWay {
155 pub fn from_name_and_default(
156 name: Option<&str>,
157 version: Option<&str>,
158 default_pipeline: Option<PipelineWay>,
159 ) -> Result<PipelineWay> {
160 if let Some(pipeline_name) = name {
161 if pipeline_name == GREPTIME_INTERNAL_TRACE_PIPELINE_V2_NAME {
162 Ok(PipelineWay::OtlpTraceDirectV2)
163 } else if pipeline_name == GREPTIME_INTERNAL_TRACE_PIPELINE_V1_NAME {
164 Ok(PipelineWay::OtlpTraceDirectV1)
165 } else if pipeline_name == GREPTIME_INTERNAL_TRACE_PIPELINE_V0_NAME {
166 Ok(PipelineWay::OtlpTraceDirectV0)
167 } else {
168 Ok(PipelineWay::Pipeline(PipelineDefinition::from_name(
169 pipeline_name,
170 to_pipeline_version(version)?,
171 None,
172 )?))
173 }
174 } else if let Some(default_pipeline) = default_pipeline {
175 Ok(default_pipeline)
176 } else {
177 PipelineMissingSnafu.fail()
178 }
179 }
180}
181
182const IDENTITY_TS_EPOCH: &str = "epoch";
183const IDENTITY_TS_DATESTR: &str = "datestr";
184
185#[derive(Debug, Clone)]
186pub enum IdentityTimeIndex {
187 Epoch(String, TimeUnit, bool),
188 DateStr(String, String, bool),
189}
190
191impl IdentityTimeIndex {
192 pub fn from_config(config: String, ignore_errors: bool) -> Result<Self> {
193 let parts = config.split(';').collect::<Vec<&str>>();
194 ensure!(
195 parts.len() == 3,
196 InvalidCustomTimeIndexSnafu {
197 config,
198 reason: "config format: '<field>;<type>;<config>'",
199 }
200 );
201
202 let field = parts[0].to_string();
203 match parts[1] {
204 IDENTITY_TS_EPOCH => match parts[2] {
205 NS_RESOLUTION => Ok(IdentityTimeIndex::Epoch(
206 field,
207 TimeUnit::Nanosecond,
208 ignore_errors,
209 )),
210 US_RESOLUTION => Ok(IdentityTimeIndex::Epoch(
211 field,
212 TimeUnit::Microsecond,
213 ignore_errors,
214 )),
215 MS_RESOLUTION => Ok(IdentityTimeIndex::Epoch(
216 field,
217 TimeUnit::Millisecond,
218 ignore_errors,
219 )),
220 S_RESOLUTION => Ok(IdentityTimeIndex::Epoch(
221 field,
222 TimeUnit::Second,
223 ignore_errors,
224 )),
225 _ => InvalidCustomTimeIndexSnafu {
226 config,
227 reason: "epoch type must be one of ns, us, ms, s",
228 }
229 .fail(),
230 },
231 IDENTITY_TS_DATESTR => Ok(IdentityTimeIndex::DateStr(
232 field,
233 parts[2].to_string(),
234 ignore_errors,
235 )),
236 _ => InvalidCustomTimeIndexSnafu {
237 config,
238 reason: "identity time index type must be one of epoch, datestr",
239 }
240 .fail(),
241 }
242 }
243
244 pub fn get_column_name(&self) -> &str {
245 match self {
246 IdentityTimeIndex::Epoch(field, _, _) => field,
247 IdentityTimeIndex::DateStr(field, _, _) => field,
248 }
249 }
250
251 pub fn get_ignore_errors(&self) -> bool {
252 match self {
253 IdentityTimeIndex::Epoch(_, _, ignore_errors) => *ignore_errors,
254 IdentityTimeIndex::DateStr(_, _, ignore_errors) => *ignore_errors,
255 }
256 }
257
258 pub fn get_datatype(&self) -> ColumnDataType {
259 match self {
260 IdentityTimeIndex::Epoch(_, unit, _) => match unit {
261 TimeUnit::Nanosecond => ColumnDataType::TimestampNanosecond,
262 TimeUnit::Microsecond => ColumnDataType::TimestampMicrosecond,
263 TimeUnit::Millisecond => ColumnDataType::TimestampMillisecond,
264 TimeUnit::Second => ColumnDataType::TimestampSecond,
265 },
266 IdentityTimeIndex::DateStr(_, _, _) => ColumnDataType::TimestampNanosecond,
267 }
268 }
269
270 pub fn get_timestamp_value(&self, value: Option<&VrlValue>) -> Result<ValueData> {
271 match self {
272 IdentityTimeIndex::Epoch(_, unit, ignore_errors) => {
273 let v = match value {
274 Some(VrlValue::Integer(v)) => *v,
275 Some(VrlValue::Bytes(s)) => match String::from_utf8_lossy(s).parse::<i64>() {
276 Ok(v) => v,
277 Err(_) => {
278 return if_ignore_errors(
279 *ignore_errors,
280 *unit,
281 format!(
282 "failed to convert {} to number",
283 String::from_utf8_lossy(s)
284 ),
285 );
286 }
287 },
288 Some(VrlValue::Timestamp(timestamp)) => datetime_utc_to_unit(timestamp, unit)?,
289 Some(v) => {
290 return if_ignore_errors(
291 *ignore_errors,
292 *unit,
293 format!("unsupported value type to convert to timestamp: {}", v),
294 );
295 }
296 None => {
297 return if_ignore_errors(
298 *ignore_errors,
299 *unit,
300 "missing field".to_string(),
301 );
302 }
303 };
304 Ok(time_unit_to_value_data(*unit, v))
305 }
306 IdentityTimeIndex::DateStr(_, format, ignore_errors) => {
307 let v = match value {
308 Some(VrlValue::Bytes(s)) => String::from_utf8_lossy(s),
309 Some(v) => {
310 return if_ignore_errors(
311 *ignore_errors,
312 TimeUnit::Nanosecond,
313 format!("unsupported value type to convert to date string: {}", v),
314 );
315 }
316 None => {
317 return if_ignore_errors(
318 *ignore_errors,
319 TimeUnit::Nanosecond,
320 "missing field".to_string(),
321 );
322 }
323 };
324
325 let timestamp = match chrono::DateTime::parse_from_str(&v, format) {
326 Ok(ts) => ts,
327 Err(_) => {
328 return if_ignore_errors(
329 *ignore_errors,
330 TimeUnit::Nanosecond,
331 format!("failed to parse date string: {}, format: {}", v, format),
332 );
333 }
334 };
335
336 Ok(ValueData::TimestampNanosecondValue(
337 timestamp
338 .timestamp_nanos_opt()
339 .context(InvalidTimestampSnafu {
340 input: timestamp.to_rfc3339(),
341 })?,
342 ))
343 }
344 }
345 }
346}
347
348fn datetime_utc_to_unit(timestamp: &DateTime<Utc>, unit: &TimeUnit) -> Result<i64> {
349 let ts = match unit {
350 TimeUnit::Nanosecond => timestamp
351 .timestamp_nanos_opt()
352 .context(InvalidTimestampSnafu {
353 input: timestamp.to_rfc3339(),
354 })?,
355 TimeUnit::Microsecond => timestamp.timestamp_micros(),
356 TimeUnit::Millisecond => timestamp.timestamp_millis(),
357 TimeUnit::Second => timestamp.timestamp(),
358 };
359 Ok(ts)
360}
361
362fn if_ignore_errors(ignore_errors: bool, unit: TimeUnit, msg: String) -> Result<ValueData> {
363 if ignore_errors {
364 Ok(time_unit_to_value_data(
365 unit,
366 Timestamp::current_time(unit).value(),
367 ))
368 } else {
369 CastTypeSnafu { msg }.fail()
370 }
371}
372
373fn time_unit_to_value_data(unit: TimeUnit, v: i64) -> ValueData {
374 match unit {
375 TimeUnit::Nanosecond => ValueData::TimestampNanosecondValue(v),
376 TimeUnit::Microsecond => ValueData::TimestampMicrosecondValue(v),
377 TimeUnit::Millisecond => ValueData::TimestampMillisecondValue(v),
378 TimeUnit::Second => ValueData::TimestampSecondValue(v),
379 }
380}