Skip to main content

pipeline/
etl.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
15#![allow(dead_code)]
16pub mod ctx_req;
17pub mod field;
18pub mod processor;
19pub mod transform;
20pub mod value;
21
22use std::collections::HashMap;
23
24use api::v1::Row;
25use common_time::timestamp::TimeUnit;
26use itertools::Itertools;
27use processor::{Processor, Processors};
28use snafu::{OptionExt, ResultExt, ensure};
29use transform::Transforms;
30use vrl::core::Value as VrlValue;
31use yaml_rust::{Yaml, YamlLoader};
32
33use crate::dispatcher::{Dispatcher, Rule};
34use crate::error::{
35    ArrayElementMustBeObjectSnafu, AutoTransformOneTimestampSnafu, Error,
36    IntermediateKeyIndexSnafu, InvalidVersionNumberSnafu, Result, TransformArrayElementSnafu,
37    YamlLoadSnafu, YamlParseSnafu,
38};
39use crate::etl::processor::ProcessorKind;
40use crate::etl::transform::transformer::greptime::{
41    RowWithTableSuffix, values_to_row, values_to_rows,
42};
43use crate::tablesuffix::TableSuffixTemplate;
44use crate::{
45    ContextOpt, GreptimeTransformer, IdentityTimeIndex, PipelineContext, SchemaInfo,
46    unwrap_or_continue_if_err,
47};
48
49const DESCRIPTION: &str = "description";
50const DOC_VERSION: &str = "version";
51const PROCESSORS: &str = "processors";
52const TRANSFORM: &str = "transform";
53const TRANSFORMS: &str = "transforms";
54const DISPATCHER: &str = "dispatcher";
55const TABLESUFFIX: &str = "table_suffix";
56
57pub enum Content<'a> {
58    Json(&'a str),
59    Yaml(&'a str),
60}
61
62pub fn parse(input: &Content) -> Result<Pipeline> {
63    match input {
64        Content::Yaml(str) => {
65            let docs = YamlLoader::load_from_str(str).context(YamlLoadSnafu)?;
66
67            ensure!(docs.len() == 1, YamlParseSnafu);
68
69            let doc = &docs[0];
70
71            let description = doc[DESCRIPTION].as_str().map(|s| s.to_string());
72
73            let doc_version = (&doc[DOC_VERSION]).try_into()?;
74
75            let processors = if let Some(v) = doc[PROCESSORS].as_vec() {
76                v.try_into()?
77            } else {
78                Processors::default()
79            };
80
81            let transformers = if let Some(v) = doc[TRANSFORMS].as_vec().or(doc[TRANSFORM].as_vec())
82            {
83                v.try_into()?
84            } else {
85                Transforms::default()
86            };
87
88            let transformer = if transformers.is_empty() {
89                // use auto transform
90                // check processors have at least one timestamp-related processor
91                let cnt = processors
92                    .iter()
93                    .filter_map(|p| match p {
94                        ProcessorKind::Date(d) if !d.ignore_missing() => Some(
95                            d.fields
96                                .iter()
97                                .map(|f| (f.target_or_input_field(), TimeUnit::Nanosecond))
98                                .collect_vec(),
99                        ),
100                        ProcessorKind::Epoch(e) if !e.ignore_missing() => Some(
101                            e.fields
102                                .iter()
103                                .map(|f| (f.target_or_input_field(), (&e.resolution).into()))
104                                .collect_vec(),
105                        ),
106                        _ => None,
107                    })
108                    .flatten()
109                    .collect_vec();
110                ensure!(cnt.len() == 1, AutoTransformOneTimestampSnafu);
111
112                let (ts_name, timeunit) = cnt.first().unwrap();
113                TransformerMode::AutoTransform(ts_name.to_string(), *timeunit)
114            } else {
115                TransformerMode::GreptimeTransformer(GreptimeTransformer::new(
116                    transformers,
117                    &doc_version,
118                )?)
119            };
120
121            let dispatcher = if !doc[DISPATCHER].is_badvalue() {
122                Some(Dispatcher::try_from(&doc[DISPATCHER])?)
123            } else {
124                None
125            };
126
127            let tablesuffix = if !doc[TABLESUFFIX].is_badvalue() {
128                Some(TableSuffixTemplate::try_from(&doc[TABLESUFFIX])?)
129            } else {
130                None
131            };
132
133            Ok(Pipeline {
134                doc_version,
135                description,
136                processors,
137                transformer,
138                dispatcher,
139                tablesuffix,
140            })
141        }
142        Content::Json(_) => unimplemented!(),
143    }
144}
145
146#[derive(Debug, Default, Copy, Clone, PartialEq, Eq)]
147pub enum PipelineDocVersion {
148    /// 1. All fields meant to be preserved have to explicitly set in the transform section.
149    /// 2. Or no transform is set, then the auto-transform will be used.
150    #[default]
151    V1,
152
153    /// A combination of transform and auto-transform.
154    /// First it goes through the transform section,
155    /// then use auto-transform to set the rest fields.
156    ///
157    /// This is useful if you only want to set the index field,
158    /// and let the normal fields be auto-inferred.
159    V2,
160}
161
162impl TryFrom<&Yaml> for PipelineDocVersion {
163    type Error = Error;
164
165    fn try_from(value: &Yaml) -> Result<Self> {
166        if value.is_badvalue() || value.is_null() {
167            return Ok(PipelineDocVersion::V1);
168        }
169
170        let version = match value {
171            Yaml::String(s) => s
172                .parse::<i64>()
173                .map_err(|_| InvalidVersionNumberSnafu { version: s.clone() }.build())?,
174            Yaml::Integer(i) => *i,
175            _ => {
176                return InvalidVersionNumberSnafu {
177                    version: value.as_str().unwrap_or_default().to_string(),
178                }
179                .fail();
180            }
181        };
182
183        match version {
184            1 => Ok(PipelineDocVersion::V1),
185            2 => Ok(PipelineDocVersion::V2),
186            _ => InvalidVersionNumberSnafu {
187                version: version.to_string(),
188            }
189            .fail(),
190        }
191    }
192}
193
194#[derive(Debug)]
195pub struct Pipeline {
196    doc_version: PipelineDocVersion,
197    description: Option<String>,
198    processors: processor::Processors,
199    dispatcher: Option<Dispatcher>,
200    transformer: TransformerMode,
201    tablesuffix: Option<TableSuffixTemplate>,
202}
203
204#[derive(Debug, Clone)]
205pub enum TransformerMode {
206    GreptimeTransformer(GreptimeTransformer),
207    AutoTransform(String, TimeUnit),
208}
209
210/// Where the pipeline executed is dispatched to, with context information
211#[derive(Debug, Hash, PartialEq, Eq, Clone, PartialOrd, Ord)]
212pub struct DispatchedTo {
213    pub table_suffix: String,
214    pub pipeline: Option<String>,
215}
216
217impl From<&Rule> for DispatchedTo {
218    fn from(value: &Rule) -> Self {
219        DispatchedTo {
220            table_suffix: value.table_suffix.clone(),
221            pipeline: value.pipeline.clone(),
222        }
223    }
224}
225
226impl DispatchedTo {
227    /// Generate destination table name from input
228    pub fn dispatched_to_table_name(&self, original: &str) -> String {
229        [original, &self.table_suffix].concat()
230    }
231}
232
233/// The result of pipeline execution
234#[derive(Debug)]
235pub enum PipelineExecOutput {
236    Transformed(TransformedOutput),
237    DispatchedTo(DispatchedTo, VrlValue),
238    Filtered,
239}
240
241/// The result after processors and dispatcher rules have run.
242#[derive(Debug)]
243pub enum PipelineProcessOutput {
244    Processed(VrlValue),
245    DispatchedTo(DispatchedTo, VrlValue),
246    Filtered,
247}
248
249/// Output from a successful pipeline transformation.
250///
251/// Rows are grouped by their ContextOpt, with each row having its own optional
252/// table_suffix for routing to different tables when using one-to-many expansion.
253/// This enables true per-row configuration options where different rows can have
254/// different database settings (TTL, merge mode, etc.).
255#[derive(Debug)]
256pub struct TransformedOutput {
257    /// Rows grouped by their ContextOpt, each with optional table suffix
258    pub rows_by_context: HashMap<ContextOpt, Vec<RowWithTableSuffix>>,
259}
260
261impl PipelineExecOutput {
262    // Note: This is a test only function, do not use it in production.
263    pub fn into_transformed(self) -> Option<Vec<RowWithTableSuffix>> {
264        if let Self::Transformed(TransformedOutput { rows_by_context }) = self {
265            // For backward compatibility, merge all rows with a default ContextOpt
266            Some(rows_by_context.into_values().flatten().collect())
267        } else {
268            None
269        }
270    }
271
272    // New method for accessing the HashMap structure directly
273    pub fn into_transformed_hashmap(self) -> Option<HashMap<ContextOpt, Vec<RowWithTableSuffix>>> {
274        if let Self::Transformed(TransformedOutput { rows_by_context }) = self {
275            Some(rows_by_context)
276        } else {
277            None
278        }
279    }
280
281    // Note: This is a test only function, do not use it in production.
282    pub fn into_dispatched(self) -> Option<DispatchedTo> {
283        if let Self::DispatchedTo(d, _) = self {
284            Some(d)
285        } else {
286            None
287        }
288    }
289}
290
291impl Pipeline {
292    fn is_v1(&self) -> bool {
293        self.doc_version == PipelineDocVersion::V1
294    }
295
296    pub fn exec_mut(
297        &self,
298        val: VrlValue,
299        pipeline_ctx: &PipelineContext<'_>,
300        schema_info: &mut SchemaInfo,
301    ) -> Result<PipelineExecOutput> {
302        match self.process_mut(val)? {
303            PipelineProcessOutput::Processed(val) => self
304                .transform_mut(val, pipeline_ctx, schema_info)
305                .map(PipelineExecOutput::Transformed),
306            PipelineProcessOutput::DispatchedTo(dispatched_to, val) => {
307                Ok(PipelineExecOutput::DispatchedTo(dispatched_to, val))
308            }
309            PipelineProcessOutput::Filtered => Ok(PipelineExecOutput::Filtered),
310        }
311    }
312
313    pub fn process_mut(&self, mut val: VrlValue) -> Result<PipelineProcessOutput> {
314        for processor in self.processors.iter() {
315            val = processor.exec_mut(val)?;
316            if val.is_null() {
317                return Ok(PipelineProcessOutput::Filtered);
318            }
319        }
320
321        if let Some(rule) = self.dispatcher.as_ref().and_then(|d| d.exec(&val)) {
322            return Ok(PipelineProcessOutput::DispatchedTo(rule.into(), val));
323        }
324
325        Ok(PipelineProcessOutput::Processed(val))
326    }
327
328    pub fn transform_mut(
329        &self,
330        val: VrlValue,
331        pipeline_ctx: &PipelineContext<'_>,
332        schema_info: &mut SchemaInfo,
333    ) -> Result<TransformedOutput> {
334        let mut val = if val.is_array() {
335            val
336        } else {
337            VrlValue::Array(vec![val])
338        };
339
340        let rows_by_context = match self.transformer() {
341            TransformerMode::GreptimeTransformer(greptime_transformer) => {
342                transform_array_elements_by_ctx(
343                    // SAFETY: by line 326, val must be an array
344                    val.as_array_mut().unwrap(),
345                    greptime_transformer,
346                    self.is_v1(),
347                    schema_info,
348                    pipeline_ctx,
349                    self.tablesuffix.as_ref(),
350                )?
351            }
352            TransformerMode::AutoTransform(ts_name, time_unit) => {
353                let def = crate::PipelineDefinition::GreptimeIdentityPipeline(Some(
354                    IdentityTimeIndex::Epoch(ts_name.clone(), *time_unit, false),
355                ));
356                let n_ctx =
357                    PipelineContext::new(&def, pipeline_ctx.pipeline_param, pipeline_ctx.channel);
358                values_to_rows(
359                    schema_info,
360                    val,
361                    &n_ctx,
362                    None,
363                    true,
364                    self.tablesuffix.as_ref(),
365                )?
366            }
367        };
368
369        Ok(TransformedOutput { rows_by_context })
370    }
371
372    pub fn processors(&self) -> &processor::Processors {
373        &self.processors
374    }
375
376    pub fn transformer(&self) -> &TransformerMode {
377        &self.transformer
378    }
379
380    pub fn resolve_table_suffix(&self, value: &VrlValue) -> Option<String> {
381        ContextOpt::resolve_table_suffix(self.tablesuffix.as_ref(), value)
382    }
383
384    // the method is for test purpose
385    pub fn schemas(&self) -> Option<&Vec<greptime_proto::v1::ColumnSchema>> {
386        match &self.transformer {
387            TransformerMode::GreptimeTransformer(t) => Some(t.schemas()),
388            TransformerMode::AutoTransform(_, _) => None,
389        }
390    }
391
392    pub fn is_variant_table_name(&self) -> bool {
393        // even if the pipeline doesn't have dispatcher or table_suffix,
394        // it can still be a variant because of VRL processor and hint
395        self.dispatcher.is_some() || self.tablesuffix.is_some()
396    }
397}
398
399/// Transforms an array of VRL values into rows grouped by their ContextOpt.
400/// Each element can have its own ContextOpt for per-row configuration.
401fn transform_array_elements_by_ctx(
402    arr: &mut [VrlValue],
403    transformer: &GreptimeTransformer,
404    is_v1: bool,
405    schema_info: &mut SchemaInfo,
406    pipeline_ctx: &PipelineContext<'_>,
407    tablesuffix_template: Option<&TableSuffixTemplate>,
408) -> Result<HashMap<ContextOpt, Vec<RowWithTableSuffix>>> {
409    let skip_error = pipeline_ctx.pipeline_param.skip_error();
410    let mut rows_by_context = HashMap::new();
411
412    for (index, element) in arr.iter_mut().enumerate() {
413        if !element.is_object() {
414            unwrap_or_continue_if_err!(
415                ArrayElementMustBeObjectSnafu {
416                    index,
417                    actual_type: element.kind_str().to_string(),
418                }
419                .fail(),
420                skip_error
421            );
422        }
423
424        let table_suffix = ContextOpt::resolve_table_suffix(tablesuffix_template, element);
425        let values = unwrap_or_continue_if_err!(
426            transformer.transform_mut_with_schema(
427                element,
428                is_v1,
429                schema_info,
430                table_suffix.as_deref(),
431            ),
432            skip_error
433        );
434        if is_v1 {
435            // v1 mode: just use transformer output directly
436            let opt = unwrap_or_continue_if_err!(
437                ContextOpt::from_pipeline_map_to_opt(element),
438                skip_error
439            );
440            rows_by_context
441                .entry(opt)
442                .or_insert_with(Vec::new)
443                .push((Row { values }, table_suffix));
444        } else {
445            // v2 mode: combine with auto-transform for remaining fields
446            let mut value = element.clone();
447            let opt = unwrap_or_continue_if_err!(
448                ContextOpt::from_pipeline_map_to_opt(&mut value),
449                skip_error
450            );
451            let row = unwrap_or_continue_if_err!(
452                values_to_row(schema_info, value, pipeline_ctx, Some(values), false,)
453                    .map_err(Box::new)
454                    .context(TransformArrayElementSnafu { index }),
455                skip_error
456            );
457            rows_by_context
458                .entry(opt)
459                .or_default()
460                .push((row, table_suffix));
461        }
462    }
463
464    Ok(rows_by_context)
465}
466
467pub(crate) fn find_key_index(intermediate_keys: &[String], key: &str, kind: &str) -> Result<usize> {
468    intermediate_keys
469        .iter()
470        .position(|k| k == key)
471        .context(IntermediateKeyIndexSnafu { kind, key })
472}
473
474/// This macro is test only, do not use it in production.
475/// The schema_info cannot be used in auto-transform ts-infer mode for lacking the ts schema.
476///
477/// Usage:
478/// ```ignore
479/// let (pipeline, schema_info, pipeline_def, pipeline_param) = setup_pipeline!(pipeline);
480/// let pipeline_ctx = PipelineContext::new(&pipeline_def, &pipeline_param, Channel::Unknown);
481/// ```
482#[macro_export]
483macro_rules! setup_pipeline {
484    ($pipeline:expr) => {{
485        use std::sync::Arc;
486
487        use $crate::{GreptimePipelineParams, Pipeline, PipelineDefinition, SchemaInfo};
488
489        let pipeline: Arc<Pipeline> = Arc::new($pipeline);
490        let schema = pipeline.schemas().unwrap();
491        let schema_info = SchemaInfo::from_schema_list(schema.clone());
492
493        let pipeline_def = PipelineDefinition::Resolved(pipeline.clone());
494        let pipeline_param = GreptimePipelineParams::default();
495
496        (pipeline, schema_info, pipeline_def, pipeline_param)
497    }};
498}
499
500#[cfg(test)]
501mod tests {
502    use std::collections::BTreeMap;
503    use std::sync::Arc;
504
505    use api::v1::Rows;
506    use datatypes::schema::{FulltextOptions, SkippingIndexOptions};
507    use greptime_proto::v1::value::ValueData;
508    use greptime_proto::v1::{self, ColumnDataType, SemanticType};
509    use vrl::prelude::Bytes;
510    use vrl::value::KeyString;
511
512    use super::*;
513
514    #[test]
515    fn test_pipeline_prepare() {
516        let input_value_str = r#"
517                    {
518                        "my_field": "1,2",
519                        "foo": "bar",
520                        "ts": "1"
521                    }
522                "#;
523        let input_value: serde_json::Value = serde_json::from_str(input_value_str).unwrap();
524
525        let pipeline_yaml = r#"description: 'Pipeline for Apache Tomcat'
526processors:
527    - csv:
528        field: my_field
529        target_fields: field1, field2
530    - epoch:
531        field: ts
532        resolution: ns
533transform:
534    - field: field1
535      type: uint32
536    - field: field2
537      type: uint32
538    - field: ts
539      type: timestamp, ns
540      index: time
541    "#;
542
543        let pipeline: Pipeline = parse(&Content::Yaml(pipeline_yaml)).unwrap();
544        let (pipeline, mut schema_info, pipeline_def, pipeline_param) = setup_pipeline!(pipeline);
545        let pipeline_ctx = PipelineContext::new(
546            &pipeline_def,
547            &pipeline_param,
548            session::context::Channel::Unknown,
549        );
550
551        let payload = input_value.into();
552        let mut result = pipeline
553            .exec_mut(payload, &pipeline_ctx, &mut schema_info)
554            .unwrap()
555            .into_transformed()
556            .unwrap();
557
558        let (row, _table_suffix) = result.swap_remove(0);
559        assert_eq!(row.values[0].value_data, Some(ValueData::U32Value(1)));
560        assert_eq!(row.values[1].value_data, Some(ValueData::U32Value(2)));
561        match &row.values[2].value_data {
562            Some(ValueData::TimestampNanosecondValue(v)) => {
563                assert_ne!(v, &0);
564            }
565            _ => panic!("expect null value"),
566        }
567    }
568
569    #[test]
570    fn test_dissect_pipeline() {
571        let message = r#"129.37.245.88 - meln1ks [01/Aug/2024:14:22:47 +0800] "PATCH /observability/metrics/production HTTP/1.0" 501 33085"#.to_string();
572        let pipeline_str = r#"processors:
573    - dissect:
574        fields:
575          - message
576        patterns:
577          - "%{ip} %{?ignored} %{username} [%{ts}] \"%{method} %{path} %{proto}\" %{status} %{bytes}"
578    - date:
579        fields:
580          - ts
581        formats:
582          - "%d/%b/%Y:%H:%M:%S %z"
583
584transform:
585    - fields:
586        - ip
587        - username
588        - method
589        - path
590        - proto
591      type: string
592    - fields:
593        - status
594      type: uint16
595    - fields:
596        - bytes
597      type: uint32
598    - field: ts
599      type: timestamp, ns
600      index: time"#;
601        let pipeline: Pipeline = parse(&Content::Yaml(pipeline_str)).unwrap();
602        let pipeline = Arc::new(pipeline);
603        let schema = pipeline.schemas().unwrap();
604        let mut schema_info = SchemaInfo::from_schema_list(schema.clone());
605
606        let pipeline_def = crate::PipelineDefinition::Resolved(pipeline.clone());
607        let pipeline_param = crate::GreptimePipelineParams::default();
608        let pipeline_ctx = PipelineContext::new(
609            &pipeline_def,
610            &pipeline_param,
611            session::context::Channel::Unknown,
612        );
613        let payload = VrlValue::Object(BTreeMap::from([(
614            KeyString::from("message"),
615            VrlValue::Bytes(Bytes::from(message)),
616        )]));
617
618        let result = pipeline
619            .exec_mut(payload, &pipeline_ctx, &mut schema_info)
620            .unwrap()
621            .into_transformed()
622            .unwrap();
623
624        assert_eq!(schema_info.schema.len(), result[0].0.values.len());
625        let test = [
626            (
627                ColumnDataType::String as i32,
628                Some(ValueData::StringValue("129.37.245.88".into())),
629            ),
630            (
631                ColumnDataType::String as i32,
632                Some(ValueData::StringValue("meln1ks".into())),
633            ),
634            (
635                ColumnDataType::String as i32,
636                Some(ValueData::StringValue("PATCH".into())),
637            ),
638            (
639                ColumnDataType::String as i32,
640                Some(ValueData::StringValue(
641                    "/observability/metrics/production".into(),
642                )),
643            ),
644            (
645                ColumnDataType::String as i32,
646                Some(ValueData::StringValue("HTTP/1.0".into())),
647            ),
648            (
649                ColumnDataType::Uint16 as i32,
650                Some(ValueData::U16Value(501)),
651            ),
652            (
653                ColumnDataType::Uint32 as i32,
654                Some(ValueData::U32Value(33085)),
655            ),
656            (
657                ColumnDataType::TimestampNanosecond as i32,
658                Some(ValueData::TimestampNanosecondValue(1722493367000000000)),
659            ),
660        ];
661        // manually set schema
662        let schema = pipeline.schemas().unwrap();
663        for i in 0..schema.len() {
664            let schema = &schema[i];
665            let value = &result[0].0.values[i];
666            assert_eq!(schema.datatype, test[i].0);
667            assert_eq!(value.value_data, test[i].1);
668        }
669    }
670
671    #[test]
672    fn test_csv_pipeline() {
673        let input_value_str = r#"
674                    {
675                        "my_field": "1,2",
676                        "foo": "bar",
677                        "ts": "1"
678                    }
679                "#;
680        let input_value: serde_json::Value = serde_json::from_str(input_value_str).unwrap();
681
682        let pipeline_yaml = r#"
683    description: Pipeline for Apache Tomcat
684    processors:
685      - csv:
686          field: my_field
687          target_fields: field1, field2
688      - epoch:
689          field: ts
690          resolution: ns
691    transform:
692      - field: field1
693        type: uint32
694      - field: field2
695        type: uint32
696      - field: ts
697        type: timestamp, ns
698        index: time
699    "#;
700
701        let pipeline: Pipeline = parse(&Content::Yaml(pipeline_yaml)).unwrap();
702        let (pipeline, mut schema_info, pipeline_def, pipeline_param) = setup_pipeline!(pipeline);
703        let pipeline_ctx = PipelineContext::new(
704            &pipeline_def,
705            &pipeline_param,
706            session::context::Channel::Unknown,
707        );
708
709        let payload = input_value.into();
710        let result = pipeline
711            .exec_mut(payload, &pipeline_ctx, &mut schema_info)
712            .unwrap()
713            .into_transformed()
714            .unwrap();
715        assert_eq!(
716            result[0].0.values[0].value_data,
717            Some(ValueData::U32Value(1))
718        );
719        assert_eq!(
720            result[0].0.values[1].value_data,
721            Some(ValueData::U32Value(2))
722        );
723        match &result[0].0.values[2].value_data {
724            Some(ValueData::TimestampNanosecondValue(v)) => {
725                assert_ne!(v, &0);
726            }
727            _ => panic!("expect null value"),
728        }
729    }
730
731    #[test]
732    fn test_date_pipeline() {
733        let input_value_str = r#"
734                {
735                    "my_field": "1,2",
736                    "foo": "bar",
737                    "test_time": "2014-5-17T04:34:56+00:00"
738                }
739            "#;
740        let input_value: serde_json::Value = serde_json::from_str(input_value_str).unwrap();
741
742        let pipeline_yaml = r#"---
743description: Pipeline for Apache Tomcat
744
745processors:
746    - date:
747        field: test_time
748
749transform:
750    - field: test_time
751      type: timestamp, ns
752      index: time
753    "#;
754
755        let pipeline: Pipeline = parse(&Content::Yaml(pipeline_yaml)).unwrap();
756        let pipeline = Arc::new(pipeline);
757        let schema = pipeline.schemas().unwrap();
758        let mut schema_info = SchemaInfo::from_schema_list(schema.clone());
759
760        let pipeline_def = crate::PipelineDefinition::Resolved(pipeline.clone());
761        let pipeline_param = crate::GreptimePipelineParams::default();
762        let pipeline_ctx = PipelineContext::new(
763            &pipeline_def,
764            &pipeline_param,
765            session::context::Channel::Unknown,
766        );
767        let schema = pipeline.schemas().unwrap().clone();
768        let result = input_value.into();
769
770        let rows_with_suffix = pipeline
771            .exec_mut(result, &pipeline_ctx, &mut schema_info)
772            .unwrap()
773            .into_transformed()
774            .unwrap();
775        let output = Rows {
776            schema,
777            rows: rows_with_suffix.into_iter().map(|(r, _)| r).collect(),
778        };
779        let schemas = output.schema;
780
781        assert_eq!(schemas.len(), 1);
782        let schema = schemas[0].clone();
783        assert_eq!("test_time", schema.column_name);
784        assert_eq!(ColumnDataType::TimestampNanosecond as i32, schema.datatype);
785        assert_eq!(SemanticType::Timestamp as i32, schema.semantic_type);
786
787        let row = output.rows[0].clone();
788        assert_eq!(1, row.values.len());
789        let value_data = row.values[0].clone().value_data;
790        assert_eq!(
791            Some(v1::value::ValueData::TimestampNanosecondValue(
792                1400301296000000000
793            )),
794            value_data
795        );
796    }
797
798    #[test]
799    fn test_dispatcher() {
800        let pipeline_yaml = r#"
801---
802description: Pipeline for Apache Tomcat
803
804processors:
805  - epoch:
806      field: ts
807      resolution: ns
808
809dispatcher:
810  field: typename
811  rules:
812    - value: http
813      table_suffix: http_events
814    - value: database
815      table_suffix: db_events
816      pipeline: database_pipeline
817
818transform:
819  - field: typename
820    type: string
821  - field: ts
822    type: timestamp, ns
823    index: time
824"#;
825        let pipeline: Pipeline = parse(&Content::Yaml(pipeline_yaml)).unwrap();
826        let dispatcher = pipeline.dispatcher.expect("expect dispatcher");
827        assert_eq!(dispatcher.field, "typename");
828
829        assert_eq!(dispatcher.rules.len(), 2);
830
831        assert_eq!(
832            dispatcher.rules[0],
833            crate::dispatcher::Rule {
834                value: VrlValue::Bytes(Bytes::from("http")),
835                table_suffix: "http_events".to_string(),
836                pipeline: None
837            }
838        );
839
840        assert_eq!(
841            dispatcher.rules[1],
842            crate::dispatcher::Rule {
843                value: VrlValue::Bytes(Bytes::from("database")),
844                table_suffix: "db_events".to_string(),
845                pipeline: Some("database_pipeline".to_string()),
846            }
847        );
848
849        let bad_yaml1 = r#"
850---
851description: Pipeline for Apache Tomcat
852
853processors:
854  - epoch:
855      field: ts
856      resolution: ns
857
858dispatcher:
859  _field: typename
860  rules:
861    - value: http
862      table_suffix: http_events
863    - value: database
864      table_suffix: db_events
865      pipeline: database_pipeline
866
867transform:
868  - field: typename
869    type: string
870  - field: ts
871    type: timestamp, ns
872    index: time
873"#;
874        let bad_yaml2 = r#"
875---
876description: Pipeline for Apache Tomcat
877
878processors:
879  - epoch:
880      field: ts
881      resolution: ns
882dispatcher:
883  field: typename
884  rules:
885    - value: http
886      _table_suffix: http_events
887    - value: database
888      _table_suffix: db_events
889      pipeline: database_pipeline
890
891transform:
892  - field: typename
893    type: string
894  - field: ts
895    type: timestamp, ns
896    index: time
897"#;
898        let bad_yaml3 = r#"
899---
900description: Pipeline for Apache Tomcat
901
902processors:
903  - epoch:
904      field: ts
905      resolution: ns
906dispatcher:
907  field: typename
908  rules:
909    - _value: http
910      table_suffix: http_events
911    - _value: database
912      table_suffix: db_events
913      pipeline: database_pipeline
914
915transform:
916  - field: typename
917    type: string
918  - field: ts
919    type: timestamp, ns
920    index: time
921"#;
922
923        let r: Result<Pipeline> = parse(&Content::Yaml(bad_yaml1));
924        assert!(r.is_err());
925        let r: Result<Pipeline> = parse(&Content::Yaml(bad_yaml2));
926        assert!(r.is_err());
927        let r: Result<Pipeline> = parse(&Content::Yaml(bad_yaml3));
928        assert!(r.is_err());
929    }
930
931    /// Test one-to-many VRL pipeline expansion.
932    /// A VRL processor can return an array, which results in multiple output rows.
933    #[test]
934    fn test_one_to_many_vrl_expansion() {
935        let pipeline_yaml = r#"
936processors:
937  - epoch:
938      field: timestamp
939      resolution: ms
940  - vrl:
941      source: |
942        events = del(.events)
943        base_host = del(.host)
944        base_ts = del(.timestamp)
945        map_values(array!(events)) -> |event| {
946            {
947                "host": base_host,
948                "event_type": event.type,
949                "event_value": event.value,
950                "timestamp": base_ts
951            }
952        }
953
954transform:
955  - field: host
956    type: string
957  - field: event_type
958    type: string
959  - field: event_value
960    type: int32
961  - field: timestamp
962    type: timestamp, ms
963    index: time
964"#;
965
966        let pipeline: Pipeline = parse(&Content::Yaml(pipeline_yaml)).unwrap();
967        let (pipeline, mut schema_info, pipeline_def, pipeline_param) = setup_pipeline!(pipeline);
968        let pipeline_ctx = PipelineContext::new(
969            &pipeline_def,
970            &pipeline_param,
971            session::context::Channel::Unknown,
972        );
973
974        // Input with 3 events
975        let input_value: serde_json::Value = serde_json::from_str(
976            r#"{
977                "host": "server1",
978                "timestamp": 1716668197217,
979                "events": [
980                    {"type": "cpu", "value": 80},
981                    {"type": "memory", "value": 60},
982                    {"type": "disk", "value": 45}
983                ]
984            }"#,
985        )
986        .unwrap();
987
988        let payload = input_value.into();
989        let result = pipeline
990            .exec_mut(payload, &pipeline_ctx, &mut schema_info)
991            .unwrap()
992            .into_transformed()
993            .unwrap();
994
995        // Should produce 3 rows from 1 input
996        assert_eq!(result.len(), 3);
997
998        // Verify each row has correct structure
999        for (row, _table_suffix) in &result {
1000            assert_eq!(row.values.len(), 4); // host, event_type, event_value, timestamp
1001            // First value should be "server1"
1002            assert_eq!(
1003                row.values[0].value_data,
1004                Some(ValueData::StringValue("server1".to_string()))
1005            );
1006            // Last value should be the timestamp
1007            assert_eq!(
1008                row.values[3].value_data,
1009                Some(ValueData::TimestampMillisecondValue(1716668197217))
1010            );
1011        }
1012
1013        // Verify event types
1014        let event_types: Vec<_> = result
1015            .iter()
1016            .map(|(r, _)| match &r.values[1].value_data {
1017                Some(ValueData::StringValue(s)) => s.clone(),
1018                _ => panic!("expected string"),
1019            })
1020            .collect();
1021        assert!(event_types.contains(&"cpu".to_string()));
1022        assert!(event_types.contains(&"memory".to_string()));
1023        assert!(event_types.contains(&"disk".to_string()));
1024    }
1025
1026    /// Test that single object output still works (backward compatibility)
1027    #[test]
1028    fn test_single_object_output_unchanged() {
1029        let pipeline_yaml = r#"
1030processors:
1031  - epoch:
1032      field: ts
1033      resolution: ms
1034  - vrl:
1035      source: |
1036        .processed = true
1037        .
1038
1039transform:
1040  - field: name
1041    type: string
1042  - field: processed
1043    type: boolean
1044  - field: ts
1045    type: timestamp, ms
1046    index: time
1047"#;
1048
1049        let pipeline: Pipeline = parse(&Content::Yaml(pipeline_yaml)).unwrap();
1050        let (pipeline, mut schema_info, pipeline_def, pipeline_param) = setup_pipeline!(pipeline);
1051        let pipeline_ctx = PipelineContext::new(
1052            &pipeline_def,
1053            &pipeline_param,
1054            session::context::Channel::Unknown,
1055        );
1056
1057        let input_value: serde_json::Value = serde_json::from_str(
1058            r#"{
1059                "name": "test",
1060                "ts": 1716668197217
1061            }"#,
1062        )
1063        .unwrap();
1064
1065        let payload = input_value.into();
1066        let result = pipeline
1067            .exec_mut(payload, &pipeline_ctx, &mut schema_info)
1068            .unwrap()
1069            .into_transformed()
1070            .unwrap();
1071
1072        // Should produce exactly 1 row
1073        assert_eq!(result.len(), 1);
1074        assert_eq!(
1075            result[0].0.values[0].value_data,
1076            Some(ValueData::StringValue("test".to_string()))
1077        );
1078        assert_eq!(
1079            result[0].0.values[1].value_data,
1080            Some(ValueData::BoolValue(true))
1081        );
1082    }
1083
1084    /// Test that empty array produces zero rows
1085    #[test]
1086    fn test_empty_array_produces_zero_rows() {
1087        let pipeline_yaml = r#"
1088processors:
1089  - vrl:
1090      source: |
1091        .events
1092
1093transform:
1094  - field: value
1095    type: int32
1096  - field: greptime_timestamp
1097    type: timestamp, ns
1098    index: time
1099"#;
1100
1101        let pipeline: Pipeline = parse(&Content::Yaml(pipeline_yaml)).unwrap();
1102        let (pipeline, mut schema_info, pipeline_def, pipeline_param) = setup_pipeline!(pipeline);
1103        let pipeline_ctx = PipelineContext::new(
1104            &pipeline_def,
1105            &pipeline_param,
1106            session::context::Channel::Unknown,
1107        );
1108
1109        let input_value: serde_json::Value = serde_json::from_str(r#"{"events": []}"#).unwrap();
1110
1111        let payload = input_value.into();
1112        let result = pipeline
1113            .exec_mut(payload, &pipeline_ctx, &mut schema_info)
1114            .unwrap()
1115            .into_transformed()
1116            .unwrap();
1117
1118        // Empty array should produce zero rows
1119        assert_eq!(result.len(), 0);
1120    }
1121
1122    /// Test that array elements must be objects
1123    #[test]
1124    fn test_array_element_must_be_object() {
1125        let pipeline_yaml = r#"
1126processors:
1127  - vrl:
1128      source: |
1129        .items
1130
1131transform:
1132  - field: value
1133    type: int32
1134  - field: greptime_timestamp
1135    type: timestamp, ns
1136    index: time
1137"#;
1138
1139        let pipeline: Pipeline = parse(&Content::Yaml(pipeline_yaml)).unwrap();
1140        let (pipeline, mut schema_info, pipeline_def, pipeline_param) = setup_pipeline!(pipeline);
1141        let pipeline_ctx = PipelineContext::new(
1142            &pipeline_def,
1143            &pipeline_param,
1144            session::context::Channel::Unknown,
1145        );
1146
1147        // Array with non-object elements should fail
1148        let input_value: serde_json::Value =
1149            serde_json::from_str(r#"{"items": [1, 2, 3]}"#).unwrap();
1150
1151        let payload = input_value.into();
1152        let result = pipeline.exec_mut(payload, &pipeline_ctx, &mut schema_info);
1153
1154        assert!(result.is_err());
1155        let err_msg = result.unwrap_err().to_string();
1156        assert!(
1157            err_msg.contains("must be an object"),
1158            "Expected error about non-object element, got: {}",
1159            err_msg
1160        );
1161    }
1162
1163    /// Test one-to-many with table suffix from VRL hint
1164    #[test]
1165    fn test_one_to_many_with_table_suffix_hint() {
1166        let pipeline_yaml = r#"
1167processors:
1168  - epoch:
1169      field: ts
1170      resolution: ms
1171  - vrl:
1172      source: |
1173        .greptime_table_suffix = "_" + string!(.category)
1174        .
1175
1176transform:
1177  - field: name
1178    type: string
1179  - field: category
1180    type: string
1181  - field: ts
1182    type: timestamp, ms
1183    index: time
1184"#;
1185
1186        let pipeline: Pipeline = parse(&Content::Yaml(pipeline_yaml)).unwrap();
1187        let (pipeline, mut schema_info, pipeline_def, pipeline_param) = setup_pipeline!(pipeline);
1188        let pipeline_ctx = PipelineContext::new(
1189            &pipeline_def,
1190            &pipeline_param,
1191            session::context::Channel::Unknown,
1192        );
1193
1194        let input_value: serde_json::Value = serde_json::from_str(
1195            r#"{
1196                "name": "test",
1197                "category": "metrics",
1198                "ts": 1716668197217
1199            }"#,
1200        )
1201        .unwrap();
1202
1203        let payload = input_value.into();
1204        let result = pipeline
1205            .exec_mut(payload, &pipeline_ctx, &mut schema_info)
1206            .unwrap()
1207            .into_transformed()
1208            .unwrap();
1209
1210        // Should have table suffix extracted per row
1211        assert_eq!(result.len(), 1);
1212        assert_eq!(result[0].1, Some("_metrics".to_string()));
1213    }
1214
1215    /// Test one-to-many with per-row table suffix
1216    #[test]
1217    fn test_one_to_many_per_row_table_suffix() {
1218        let pipeline_yaml = r#"
1219processors:
1220  - epoch:
1221      field: timestamp
1222      resolution: ms
1223  - vrl:
1224      source: |
1225        events = del(.events)
1226        base_ts = del(.timestamp)
1227
1228        map_values(array!(events)) -> |event| {
1229            suffix = "_" + string!(event.category)
1230            {
1231                "name": event.name,
1232                "value": event.value,
1233                "timestamp": base_ts,
1234                "greptime_table_suffix": suffix
1235            }
1236        }
1237
1238transform:
1239  - field: name
1240    type: string
1241  - field: value
1242    type: int32
1243  - field: timestamp
1244    type: timestamp, ms
1245    index: time
1246"#;
1247
1248        let pipeline: Pipeline = parse(&Content::Yaml(pipeline_yaml)).unwrap();
1249        let (pipeline, mut schema_info, pipeline_def, pipeline_param) = setup_pipeline!(pipeline);
1250        let pipeline_ctx = PipelineContext::new(
1251            &pipeline_def,
1252            &pipeline_param,
1253            session::context::Channel::Unknown,
1254        );
1255
1256        // Input with events that should go to different tables
1257        let input_value: serde_json::Value = serde_json::from_str(
1258            r#"{
1259                "timestamp": 1716668197217,
1260                "events": [
1261                    {"name": "cpu_usage", "value": 80, "category": "cpu"},
1262                    {"name": "mem_usage", "value": 60, "category": "memory"},
1263                    {"name": "cpu_temp", "value": 45, "category": "cpu"}
1264                ]
1265            }"#,
1266        )
1267        .unwrap();
1268
1269        let payload = input_value.into();
1270        let result = pipeline
1271            .exec_mut(payload, &pipeline_ctx, &mut schema_info)
1272            .unwrap()
1273            .into_transformed()
1274            .unwrap();
1275
1276        // Should produce 3 rows
1277        assert_eq!(result.len(), 3);
1278
1279        // Collect table suffixes
1280        let table_suffixes: Vec<_> = result.iter().map(|(_, suffix)| suffix.clone()).collect();
1281
1282        // Should have different table suffixes per row
1283        assert!(table_suffixes.contains(&Some("_cpu".to_string())));
1284        assert!(table_suffixes.contains(&Some("_memory".to_string())));
1285
1286        // Count rows per table suffix
1287        let cpu_count = table_suffixes
1288            .iter()
1289            .filter(|s| *s == &Some("_cpu".to_string()))
1290            .count();
1291        let memory_count = table_suffixes
1292            .iter()
1293            .filter(|s| *s == &Some("_memory".to_string()))
1294            .count();
1295        assert_eq!(cpu_count, 2);
1296        assert_eq!(memory_count, 1);
1297    }
1298
1299    /// Test that one-to-many mapping preserves per-row ContextOpt in HashMap
1300    #[test]
1301    fn test_one_to_many_hashmap_contextopt_preservation() {
1302        let pipeline_yaml = r#"
1303processors:
1304  - epoch:
1305      field: timestamp
1306      resolution: ms
1307  - vrl:
1308      source: |
1309        events = del(.events)
1310        base_ts = del(.timestamp)
1311
1312        map_values(array!(events)) -> |event| {
1313            # Set different TTL values per event type
1314            ttl = if event.type == "critical" {
1315                "1h"
1316            } else if event.type == "warning" {
1317                "24h"
1318            } else {
1319                "7d"
1320            }
1321
1322            {
1323                "host": del(.host),
1324                "event_type": event.type,
1325                "event_value": event.value,
1326                "timestamp": base_ts,
1327                "greptime_ttl": ttl
1328            }
1329        }
1330
1331transform:
1332  - field: host
1333    type: string
1334  - field: event_type
1335    type: string
1336  - field: event_value
1337    type: int32
1338  - field: timestamp
1339    type: timestamp, ms
1340    index: time
1341"#;
1342
1343        let pipeline: Pipeline = parse(&Content::Yaml(pipeline_yaml)).unwrap();
1344        let (pipeline, mut schema_info, pipeline_def, pipeline_param) = setup_pipeline!(pipeline);
1345        let pipeline_ctx = PipelineContext::new(
1346            &pipeline_def,
1347            &pipeline_param,
1348            session::context::Channel::Unknown,
1349        );
1350
1351        // Input with events that should have different ContextOpt values
1352        let input_value: serde_json::Value = serde_json::from_str(
1353            r#"{
1354                "host": "server1",
1355                "timestamp": 1716668197217,
1356                "events": [
1357                    {"type": "critical", "value": 100},
1358                    {"type": "warning", "value": 50},
1359                    {"type": "info", "value": 25}
1360                ]
1361            }"#,
1362        )
1363        .unwrap();
1364
1365        let payload = input_value.into();
1366        let result = pipeline
1367            .exec_mut(payload, &pipeline_ctx, &mut schema_info)
1368            .unwrap();
1369
1370        // Extract the HashMap structure
1371        let rows_by_context = result.into_transformed_hashmap().unwrap();
1372
1373        // Should have 3 different ContextOpt groups due to different TTL values
1374        assert_eq!(rows_by_context.len(), 3);
1375
1376        // Verify each ContextOpt group has exactly 1 row and different configurations
1377        let mut context_opts = Vec::new();
1378        for (opt, rows) in &rows_by_context {
1379            assert_eq!(rows.len(), 1); // Each group should have exactly 1 row
1380            context_opts.push(opt.clone());
1381        }
1382
1383        // ContextOpts should be different due to different TTL values
1384        assert_ne!(context_opts[0], context_opts[1]);
1385        assert_ne!(context_opts[1], context_opts[2]);
1386        assert_ne!(context_opts[0], context_opts[2]);
1387
1388        // Verify the rows are correctly structured
1389        for rows in rows_by_context.values() {
1390            for (row, _table_suffix) in rows {
1391                assert_eq!(row.values.len(), 4); // host, event_type, event_value, timestamp
1392            }
1393        }
1394    }
1395
1396    /// Test that single object input still works with HashMap structure
1397    #[test]
1398    fn test_single_object_hashmap_compatibility() {
1399        let pipeline_yaml = r#"
1400processors:
1401  - epoch:
1402      field: ts
1403      resolution: ms
1404  - vrl:
1405      source: |
1406        .processed = true
1407        .
1408
1409transform:
1410  - field: name
1411    type: string
1412  - field: processed
1413    type: boolean
1414  - field: ts
1415    type: timestamp, ms
1416    index: time
1417"#;
1418
1419        let pipeline: Pipeline = parse(&Content::Yaml(pipeline_yaml)).unwrap();
1420        let (pipeline, mut schema_info, pipeline_def, pipeline_param) = setup_pipeline!(pipeline);
1421        let pipeline_ctx = PipelineContext::new(
1422            &pipeline_def,
1423            &pipeline_param,
1424            session::context::Channel::Unknown,
1425        );
1426
1427        let input_value: serde_json::Value = serde_json::from_str(
1428            r#"{
1429                "name": "test",
1430                "ts": 1716668197217
1431            }"#,
1432        )
1433        .unwrap();
1434
1435        let payload = input_value.into();
1436        let result = pipeline
1437            .exec_mut(payload, &pipeline_ctx, &mut schema_info)
1438            .unwrap();
1439
1440        // Extract the HashMap structure
1441        let rows_by_context = result.into_transformed_hashmap().unwrap();
1442
1443        // Single object should produce exactly 1 ContextOpt group
1444        assert_eq!(rows_by_context.len(), 1);
1445
1446        let (_opt, rows) = rows_by_context.into_iter().next().unwrap();
1447        assert_eq!(rows.len(), 1);
1448
1449        // Verify the row structure
1450        let (row, _table_suffix) = &rows[0];
1451        assert_eq!(row.values.len(), 3); // name, processed, timestamp
1452    }
1453
1454    /// Test that empty arrays work correctly with HashMap structure
1455    #[test]
1456    fn test_empty_array_hashmap() {
1457        let pipeline_yaml = r#"
1458processors:
1459  - vrl:
1460      source: |
1461        .events
1462
1463transform:
1464  - field: value
1465    type: int32
1466  - field: greptime_timestamp
1467    type: timestamp, ns
1468    index: time
1469"#;
1470
1471        let pipeline: Pipeline = parse(&Content::Yaml(pipeline_yaml)).unwrap();
1472        let (pipeline, mut schema_info, pipeline_def, pipeline_param) = setup_pipeline!(pipeline);
1473        let pipeline_ctx = PipelineContext::new(
1474            &pipeline_def,
1475            &pipeline_param,
1476            session::context::Channel::Unknown,
1477        );
1478
1479        let input_value: serde_json::Value = serde_json::from_str(r#"{"events": []}"#).unwrap();
1480
1481        let payload = input_value.into();
1482        let result = pipeline
1483            .exec_mut(payload, &pipeline_ctx, &mut schema_info)
1484            .unwrap();
1485
1486        // Extract the HashMap structure
1487        let rows_by_context = result.into_transformed_hashmap().unwrap();
1488
1489        // Empty array should produce empty HashMap
1490        assert_eq!(rows_by_context.len(), 0);
1491    }
1492
1493    #[test]
1494    fn test_pipeline_detailed_index_options_roundtrip() {
1495        let pipeline_yaml = r#"
1496transform:
1497  - field: message
1498    type: string
1499    index:
1500      type: fulltext
1501      options:
1502        analyzer: Chinese
1503        case_sensitive: true
1504        backend: tantivy
1505  - field: trace_id
1506    type: int64
1507    index:
1508      type: skipping
1509      options:
1510        granularity: 2048
1511        false_positive_rate: 0.02
1512        type: BLOOM
1513  - field: ts
1514    type: timestamp, ns
1515    index: time
1516"#;
1517
1518        let pipeline: Pipeline = parse(&Content::Yaml(pipeline_yaml)).unwrap();
1519        let schema = pipeline.schemas().unwrap().clone();
1520
1521        let message = schema
1522            .iter()
1523            .find(|column| column.column_name == "message")
1524            .unwrap();
1525        let trace_id = schema
1526            .iter()
1527            .find(|column| column.column_name == "trace_id")
1528            .unwrap();
1529        let message_options = message.options.clone();
1530        let trace_id_options = trace_id.options.clone();
1531
1532        let fulltext: FulltextOptions = serde_json::from_str(
1533            message
1534                .options
1535                .as_ref()
1536                .unwrap()
1537                .options
1538                .get("fulltext")
1539                .unwrap(),
1540        )
1541        .unwrap();
1542        assert!(fulltext.enable);
1543        assert_eq!(fulltext.analyzer.to_string(), "Chinese");
1544        assert!(fulltext.case_sensitive);
1545        assert_eq!(fulltext.backend.to_string(), "tantivy");
1546
1547        let skipping: SkippingIndexOptions = serde_json::from_str(
1548            trace_id
1549                .options
1550                .as_ref()
1551                .unwrap()
1552                .options
1553                .get("skipping_index")
1554                .unwrap(),
1555        )
1556        .unwrap();
1557        assert_eq!(skipping.granularity, 2048);
1558        assert_eq!(skipping.false_positive_rate(), 0.02);
1559        assert_eq!(skipping.index_type.to_string(), "BLOOM");
1560
1561        let roundtrip_schema = SchemaInfo::from_schema_list(schema)
1562            .column_schemas()
1563            .unwrap();
1564        let roundtrip_message = roundtrip_schema
1565            .iter()
1566            .find(|column| column.column_name == "message")
1567            .unwrap();
1568        let roundtrip_trace_id = roundtrip_schema
1569            .iter()
1570            .find(|column| column.column_name == "trace_id")
1571            .unwrap();
1572
1573        assert_eq!(message_options, roundtrip_message.options);
1574        assert_eq!(trace_id_options, roundtrip_trace_id.options);
1575    }
1576}