Skip to main content

pipeline/
manager.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::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
41/// Pipeline version. An optional timestamp with nanosecond precision.
42///
43/// If the version is None, it means the latest version of the pipeline.
44/// User can specify the version by providing a timestamp string formatted as iso8601.
45/// When it used in cache key, it will be converted to i64 meaning the number of nanoseconds since the epoch.
46pub type PipelineVersion = Option<TimestampNanosecond>;
47
48/// Pipeline info. A tuple of timestamp and pipeline reference.
49pub type PipelineInfo = (Timestamp, PipelineRef);
50
51pub type PipelineTableRef = Arc<PipelineTable>;
52pub type PipelineRef = Arc<Pipeline>;
53
54/// SelectInfo is used to store the selected keys from OpenTelemetry record attrs
55/// The key is used to uplift value from the attributes and serve as column name in the table
56#[derive(Default)]
57pub struct SelectInfo {
58    pub keys: Vec<String>,
59}
60
61/// Try to convert a string to SelectInfo
62/// The string should be a comma-separated list of keys
63/// example: "key1,key2,key3"
64/// The keys will be sorted and deduplicated
65impl 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";
83/// Built-in OTLP trace pipeline that stores semi-structured fields as JSON2.
84pub const GREPTIME_INTERNAL_TRACE_PIPELINE_V2_NAME: &str = "greptime_trace_v2";
85
86/// Enum for holding information of a pipeline, which is either pipeline itself,
87/// or information that be used to retrieve a pipeline from `PipelineHandler`
88#[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}