Skip to main content

pipeline/etl/
transform.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
15pub mod index;
16pub mod transformer;
17
18use std::collections::HashMap;
19
20use api::helper::ColumnDataTypeWrapper;
21use api::v1::ColumnDataType;
22use api::v1::value::ValueData;
23use chrono::Utc;
24use datatypes::json::{JsonSettings, JsonTypeHint};
25use datatypes::schema::{FulltextOptions, SkippingIndexOptions};
26use snafu::{OptionExt, ResultExt, ensure};
27use sql::parsers::utils::{
28    validate_column_fulltext_create_option, validate_column_skipping_index_create_option,
29};
30
31use crate::error::{
32    Error, FieldMustBeTypeSnafu, InvalidJson2TypeHintSnafu, KeyMustBeStringSnafu,
33    ParseJson2TypeHintPathSnafu, Result, TransformElementMustBeMapSnafu,
34    TransformFieldMustBeSetSnafu, TransformIndexOptionMustBeScalarSnafu, TransformIndexOptionSnafu,
35    TransformIndexOptionUnsupportedSnafu, TransformIndexOptionsUnsupportedSnafu,
36    TransformIndexTypeMismatchSnafu, TransformIndexTypeMustBeSetSnafu,
37    TransformIndexUnsupportedFieldSnafu, TransformOnFailureInvalidValueSnafu,
38    TransformTypeMustBeSetSnafu, UnsupportedTypeInPipelineSnafu,
39};
40use crate::etl::field::Fields;
41use crate::etl::processor::{yaml_bool, yaml_new_field, yaml_new_fields, yaml_string};
42use crate::etl::transform::index::Index;
43use crate::etl::value::{parse_str_type, parse_str_value};
44
45const TRANSFORM_FIELD: &str = "field";
46const TRANSFORM_FIELDS: &str = "fields";
47const TRANSFORM_TYPE: &str = "type";
48const TRANSFORM_INDEX: &str = "index";
49const TRANSFORM_INDEX_TYPE_FIELD: &str = "index.type";
50const TRANSFORM_INDEX_OPTIONS: &str = "options";
51const TRANSFORM_INDEX_OPTIONS_FIELD: &str = "index.options";
52const TRANSFORM_TAG: &str = "tag";
53const TRANSFORM_DEFAULT: &str = "default";
54const TRANSFORM_ON_FAILURE: &str = "on_failure";
55const JSON2_TYPE: &str = "json2";
56const JSON2_TYPE_HINT: &str = "type.json2[]";
57const JSON2_TYPE_HINT_PATH: &str = "path";
58
59pub use transformer::greptime::GreptimeTransformer;
60
61/// On Failure behavior when transform fails
62#[derive(Debug, Clone, Default, Copy)]
63pub enum OnFailure {
64    // Return None if transform fails
65    #[default]
66    Ignore,
67    // Return default value of the field if transform fails
68    // Default value depends on the type of the field, or explicitly set by user
69    Default,
70}
71
72impl std::str::FromStr for OnFailure {
73    type Err = Error;
74
75    fn from_str(s: &str) -> Result<Self> {
76        match s {
77            "ignore" => Ok(OnFailure::Ignore),
78            "default" => Ok(OnFailure::Default),
79            _ => TransformOnFailureInvalidValueSnafu { value: s }.fail(),
80        }
81    }
82}
83
84#[derive(Debug, Default, Clone)]
85pub struct Transforms {
86    pub(crate) transforms: Vec<Transform>,
87}
88
89impl Transforms {
90    pub fn transforms(&self) -> &Vec<Transform> {
91        &self.transforms
92    }
93}
94
95impl std::ops::Deref for Transforms {
96    type Target = Vec<Transform>;
97
98    fn deref(&self) -> &Self::Target {
99        &self.transforms
100    }
101}
102
103impl std::ops::DerefMut for Transforms {
104    fn deref_mut(&mut self) -> &mut Self::Target {
105        &mut self.transforms
106    }
107}
108
109impl TryFrom<&Vec<yaml_rust::Yaml>> for Transforms {
110    type Error = Error;
111
112    fn try_from(docs: &Vec<yaml_rust::Yaml>) -> Result<Self> {
113        let mut transforms = Vec::with_capacity(32);
114        let mut all_output_keys: Vec<String> = Vec::with_capacity(32);
115        let mut all_required_keys = Vec::with_capacity(32);
116
117        for doc in docs {
118            let transform_builder: Transform = doc
119                .as_hash()
120                .context(TransformElementMustBeMapSnafu)?
121                .try_into()?;
122            let mut transform_output_keys = transform_builder
123                .fields
124                .iter()
125                .map(|f| f.target_or_input_field().to_string())
126                .collect();
127            all_output_keys.append(&mut transform_output_keys);
128
129            let mut transform_required_keys = transform_builder
130                .fields
131                .iter()
132                .map(|f| f.input_field().to_string())
133                .collect();
134            all_required_keys.append(&mut transform_required_keys);
135
136            transforms.push(transform_builder);
137        }
138
139        all_required_keys.sort();
140
141        Ok(Transforms { transforms })
142    }
143}
144
145/// only field is required
146#[derive(Debug, Clone)]
147pub struct Transform {
148    pub fields: Fields,
149    pub type_: ColumnDataType,
150    pub(crate) json_settings: Option<JsonSettings>,
151    pub default: Option<ValueData>,
152    pub index: Option<Index>,
153    pub index_options: Option<TransformIndexOptions>,
154    pub tag: bool,
155    pub on_failure: Option<OnFailure>,
156}
157
158#[derive(Debug, Clone, PartialEq, Eq)]
159pub enum TransformIndexOptions {
160    Fulltext(FulltextOptions),
161    Skipping(SkippingIndexOptions),
162}
163
164impl TransformIndexOptions {
165    pub(crate) fn index(&self) -> Index {
166        match self {
167            TransformIndexOptions::Fulltext(_) => Index::Fulltext,
168            TransformIndexOptions::Skipping(_) => Index::Skipping,
169        }
170    }
171
172    #[cfg(test)]
173    pub(crate) fn as_fulltext(&self) -> Option<&FulltextOptions> {
174        match self {
175            TransformIndexOptions::Fulltext(options) => Some(options),
176            TransformIndexOptions::Skipping(_) => None,
177        }
178    }
179
180    #[cfg(test)]
181    pub(crate) fn as_skipping(&self) -> Option<&SkippingIndexOptions> {
182        match self {
183            TransformIndexOptions::Skipping(options) => Some(options),
184            TransformIndexOptions::Fulltext(_) => None,
185        }
186    }
187}
188
189// valid types
190// ColumnDataType::Int8
191// ColumnDataType::Int16
192// ColumnDataType::Int32
193// ColumnDataType::Int64
194// ColumnDataType::Uint8
195// ColumnDataType::Uint16
196// ColumnDataType::Uint32
197// ColumnDataType::Uint64
198// ColumnDataType::Float32
199// ColumnDataType::Float64
200// ColumnDataType::Boolean
201// ColumnDataType::String
202// ColumnDataType::TimestampNanosecond
203// ColumnDataType::TimestampMicrosecond
204// ColumnDataType::TimestampMillisecond
205// ColumnDataType::TimestampSecond
206// ColumnDataType::Binary (JSONB)
207// ColumnDataType::Json (JSON2)
208
209impl Transform {
210    pub(crate) fn get_default(&self) -> Option<&ValueData> {
211        self.default.as_ref()
212    }
213
214    pub(crate) fn get_type_matched_default_val(&self) -> Result<ValueData> {
215        get_default_for_type(&self.type_)
216    }
217
218    pub(crate) fn get_default_value_when_data_is_none(&self) -> Option<ValueData> {
219        if is_timestamp_type(&self.type_) && self.index.is_some_and(|i| i == Index::Time) {
220            let now = Utc::now();
221            match self.type_ {
222                ColumnDataType::TimestampSecond => {
223                    return Some(ValueData::TimestampSecondValue(now.timestamp()));
224                }
225                ColumnDataType::TimestampMillisecond => {
226                    return Some(ValueData::TimestampMillisecondValue(now.timestamp_millis()));
227                }
228                ColumnDataType::TimestampMicrosecond => {
229                    return Some(ValueData::TimestampMicrosecondValue(now.timestamp_micros()));
230                }
231                ColumnDataType::TimestampNanosecond => {
232                    return Some(ValueData::TimestampNanosecondValue(
233                        now.timestamp_nanos_opt()?,
234                    ));
235                }
236                _ => {}
237            }
238        }
239        None
240    }
241
242    pub(crate) fn is_timeindex(&self) -> bool {
243        self.index.is_some_and(|i| i == Index::Time)
244    }
245}
246
247fn is_timestamp_type(ty: &ColumnDataType) -> bool {
248    matches!(
249        ty,
250        ColumnDataType::TimestampSecond
251            | ColumnDataType::TimestampMillisecond
252            | ColumnDataType::TimestampMicrosecond
253            | ColumnDataType::TimestampNanosecond
254    )
255}
256
257fn get_default_for_type(ty: &ColumnDataType) -> Result<ValueData> {
258    let v = match ty {
259        ColumnDataType::Boolean => ValueData::BoolValue(false),
260        ColumnDataType::Int8 => ValueData::I8Value(0),
261        ColumnDataType::Int16 => ValueData::I16Value(0),
262        ColumnDataType::Int32 => ValueData::I32Value(0),
263        ColumnDataType::Int64 => ValueData::I64Value(0),
264        ColumnDataType::Uint8 => ValueData::U8Value(0),
265        ColumnDataType::Uint16 => ValueData::U16Value(0),
266        ColumnDataType::Uint32 => ValueData::U32Value(0),
267        ColumnDataType::Uint64 => ValueData::U64Value(0),
268        ColumnDataType::Float32 => ValueData::F32Value(0.0),
269        ColumnDataType::Float64 => ValueData::F64Value(0.0),
270        ColumnDataType::Binary => ValueData::BinaryValue(jsonb::Value::Null.to_vec()),
271        ColumnDataType::Json => ValueData::JsonValue(Default::default()),
272        ColumnDataType::String => ValueData::StringValue(String::new()),
273
274        ColumnDataType::TimestampSecond => ValueData::TimestampSecondValue(0),
275        ColumnDataType::TimestampMillisecond => ValueData::TimestampMillisecondValue(0),
276        ColumnDataType::TimestampMicrosecond => ValueData::TimestampMicrosecondValue(0),
277        ColumnDataType::TimestampNanosecond => ValueData::TimestampNanosecondValue(0),
278
279        _ => UnsupportedTypeInPipelineSnafu {
280            ty: ty.as_str_name(),
281        }
282        .fail()?,
283    };
284    Ok(v)
285}
286
287fn parse_transform_index(
288    value: &yaml_rust::Yaml,
289) -> Result<(Index, Option<HashMap<String, String>>)> {
290    match value {
291        yaml_rust::Yaml::String(_) => {
292            let index_str = yaml_string(value, TRANSFORM_INDEX)?;
293            Ok((index_str.try_into()?, None))
294        }
295        yaml_rust::Yaml::Hash(hash) => {
296            let mut index = None;
297            let mut index_options = None;
298
299            for (k, v) in hash {
300                let key = k
301                    .as_str()
302                    .with_context(|| KeyMustBeStringSnafu { k: k.clone() })?;
303                match key {
304                    TRANSFORM_TYPE => {
305                        let index_str = yaml_string(v, TRANSFORM_INDEX_TYPE_FIELD)?;
306                        index = Some(index_str.try_into()?);
307                    }
308                    TRANSFORM_INDEX_OPTIONS => {
309                        index_options = Some(parse_transform_index_options(v)?);
310                    }
311                    _ => {
312                        return TransformIndexUnsupportedFieldSnafu {
313                            field: key.to_string(),
314                        }
315                        .fail();
316                    }
317                }
318            }
319
320            Ok((
321                index.context(TransformIndexTypeMustBeSetSnafu)?,
322                index_options,
323            ))
324        }
325        _ => FieldMustBeTypeSnafu {
326            field: TRANSFORM_INDEX,
327            ty: "string or map",
328        }
329        .fail(),
330    }
331}
332
333fn parse_transform_index_options(value: &yaml_rust::Yaml) -> Result<HashMap<String, String>> {
334    let hash = value.as_hash().context(FieldMustBeTypeSnafu {
335        field: TRANSFORM_INDEX_OPTIONS_FIELD,
336        ty: "map",
337    })?;
338    let mut options = HashMap::with_capacity(hash.len());
339
340    for (k, v) in hash {
341        let key = k
342            .as_str()
343            .with_context(|| KeyMustBeStringSnafu { k: k.clone() })?;
344
345        let field = format!("{TRANSFORM_INDEX_OPTIONS_FIELD}.{key}");
346        let value = match v {
347            yaml_rust::Yaml::String(v) => v.clone(),
348            yaml_rust::Yaml::Boolean(v) => v.to_string(),
349            yaml_rust::Yaml::Integer(v) => v.to_string(),
350            yaml_rust::Yaml::Real(v) => v.clone(),
351            _ => {
352                return TransformIndexOptionMustBeScalarSnafu { field }.fail();
353            }
354        };
355        options.insert(key.to_string(), value);
356    }
357
358    Ok(options)
359}
360
361fn lower_typed_transform_index_options<T>(
362    index: Index,
363    index_options: Option<HashMap<String, String>>,
364    validate: fn(&str) -> bool,
365    wrap: fn(T) -> TransformIndexOptions,
366) -> Result<Option<TransformIndexOptions>>
367where
368    T: TryFrom<HashMap<String, String>, Error = datatypes::error::Error>,
369{
370    index_options
371        .map(|opts| {
372            for key in opts.keys() {
373                ensure!(
374                    validate(key),
375                    TransformIndexOptionUnsupportedSnafu {
376                        index: index.to_string(),
377                        key: key.clone(),
378                    }
379                );
380            }
381
382            let options = opts.try_into().context(TransformIndexOptionSnafu {
383                index: index.to_string(),
384            })?;
385
386            Ok(wrap(options))
387        })
388        .transpose()
389}
390
391fn lower_transform_index_options(
392    index: Index,
393    column_type: &ColumnDataType,
394    index_options: Option<HashMap<String, String>>,
395) -> Result<Option<TransformIndexOptions>> {
396    match index {
397        Index::Fulltext => {
398            ensure!(
399                *column_type == ColumnDataType::String,
400                TransformIndexTypeMismatchSnafu {
401                    index: index.to_string(),
402                    expected: ColumnDataType::String.as_str_name().to_string(),
403                    actual: column_type.as_str_name().to_string(),
404                }
405            );
406
407            lower_typed_transform_index_options(
408                index,
409                index_options,
410                validate_column_fulltext_create_option,
411                TransformIndexOptions::Fulltext,
412            )
413        }
414        Index::Skipping => lower_typed_transform_index_options(
415            index,
416            index_options,
417            validate_column_skipping_index_create_option,
418            TransformIndexOptions::Skipping,
419        ),
420        Index::Inverted | Index::Time | Index::Tag => {
421            ensure!(
422                index_options.is_none(),
423                TransformIndexOptionsUnsupportedSnafu {
424                    index: index.to_string(),
425                }
426            );
427            Ok(None)
428        }
429    }
430}
431
432fn parse_transform_type(value: &yaml_rust::Yaml) -> Result<(ColumnDataType, Option<JsonSettings>)> {
433    if let Some(type_name) = value.as_str() {
434        return Ok((parse_str_type(type_name)?, None));
435    }
436
437    let config = value.as_hash().context(FieldMustBeTypeSnafu {
438        field: TRANSFORM_TYPE,
439        ty: "string or map",
440    })?;
441    ensure!(
442        config.len() == 1,
443        InvalidJson2TypeHintSnafu {
444            reason: "transform type map must contain exactly one `json2` field".to_string()
445        }
446    );
447    let (type_name, hints) = config.iter().next().context(InvalidJson2TypeHintSnafu {
448        reason: "transform type map must contain a `json2` field".to_string(),
449    })?;
450    let type_name = type_name.as_str().with_context(|| KeyMustBeStringSnafu {
451        k: type_name.clone(),
452    })?;
453    ensure!(
454        type_name.eq_ignore_ascii_case(JSON2_TYPE),
455        InvalidJson2TypeHintSnafu {
456            reason: format!("unsupported transform type map `{type_name}`")
457        }
458    );
459
460    let hints = hints.as_vec().context(FieldMustBeTypeSnafu {
461        field: JSON2_TYPE,
462        ty: "list",
463    })?;
464    let hints = hints
465        .iter()
466        .map(parse_json2_type_hint)
467        .collect::<Result<Vec<_>>>()?;
468    Ok((
469        ColumnDataType::Json,
470        Some(JsonSettings::try_new(hints, None)?),
471    ))
472}
473
474fn parse_json2_type_hint(value: &yaml_rust::Yaml) -> Result<JsonTypeHint> {
475    let config = value.as_hash().context(FieldMustBeTypeSnafu {
476        field: JSON2_TYPE_HINT,
477        ty: "map",
478    })?;
479    let mut path = None;
480    let mut type_name = None;
481    let mut index = None;
482
483    for (key, value) in config {
484        let key = key
485            .as_str()
486            .with_context(|| KeyMustBeStringSnafu { k: key.clone() })?;
487        match key {
488            JSON2_TYPE_HINT_PATH => path = Some(yaml_string(value, JSON2_TYPE_HINT_PATH)?),
489            TRANSFORM_TYPE => type_name = Some(yaml_string(value, TRANSFORM_TYPE)?),
490            TRANSFORM_INDEX => index = Some(value),
491            _ => {
492                return InvalidJson2TypeHintSnafu {
493                    reason: format!("unsupported field `{key}`"),
494                }
495                .fail();
496            }
497        }
498    }
499
500    let path = path.context(InvalidJson2TypeHintSnafu {
501        reason: "`path` must be set".to_string(),
502    })?;
503    let path = sql::parse_json2_type_hint_path(&path)
504        .with_context(|_| ParseJson2TypeHintPathSnafu { path: path.clone() })?;
505    let type_name = type_name.context(InvalidJson2TypeHintSnafu {
506        reason: "`type` must be set".to_string(),
507    })?;
508    let type_ = parse_str_type(&type_name)?;
509    ensure!(
510        matches!(
511            type_,
512            ColumnDataType::String
513                | ColumnDataType::Int64
514                | ColumnDataType::Uint64
515                | ColumnDataType::Float64
516                | ColumnDataType::Boolean
517        ),
518        InvalidJson2TypeHintSnafu {
519            reason: format!("unsupported type `{type_name}`")
520        }
521    );
522    let data_type = ColumnDataTypeWrapper::new(type_, None).into();
523    let inverted_index = if let Some(value) = index {
524        let (index, options) = parse_transform_index(value)?;
525        ensure!(
526            index == Index::Inverted,
527            InvalidJson2TypeHintSnafu {
528                reason: format!("unsupported index `{index}`")
529            }
530        );
531        lower_transform_index_options(index, &ColumnDataType::Json, options)?;
532        true
533    } else {
534        false
535    };
536
537    Ok(JsonTypeHint {
538        path,
539        data_type,
540        inverted_index,
541    })
542}
543
544impl TryFrom<&yaml_rust::yaml::Hash> for Transform {
545    type Error = Error;
546
547    fn try_from(hash: &yaml_rust::yaml::Hash) -> Result<Self> {
548        let mut fields = Fields::default();
549        let mut default = None;
550        let mut index_value = None;
551        let mut tag = false;
552        let mut on_failure = None;
553
554        let mut type_ = None;
555        let mut json_settings = None;
556
557        for (k, v) in hash {
558            let key = k
559                .as_str()
560                .with_context(|| KeyMustBeStringSnafu { k: k.clone() })?;
561            match key {
562                TRANSFORM_FIELD => {
563                    fields = Fields::one(yaml_new_field(v, TRANSFORM_FIELD)?);
564                }
565
566                TRANSFORM_FIELDS => {
567                    fields = yaml_new_fields(v, TRANSFORM_FIELDS)?;
568                }
569
570                TRANSFORM_TYPE => {
571                    let (parsed_type, parsed_json_settings) = parse_transform_type(v)?;
572                    type_ = Some(parsed_type);
573                    json_settings = parsed_json_settings;
574                }
575
576                TRANSFORM_INDEX => {
577                    index_value = Some(v);
578                }
579
580                TRANSFORM_TAG => {
581                    tag = yaml_bool(v, TRANSFORM_TAG)?;
582                }
583
584                TRANSFORM_DEFAULT => {
585                    default = match v {
586                        yaml_rust::Yaml::Real(r) => Some(r.clone()),
587                        yaml_rust::Yaml::Integer(i) => Some(i.to_string()),
588                        yaml_rust::Yaml::String(s) => Some(s.clone()),
589                        yaml_rust::Yaml::Boolean(b) => Some(b.to_string()),
590                        yaml_rust::Yaml::Array(_)
591                        | yaml_rust::Yaml::Hash(_)
592                        | yaml_rust::Yaml::Alias(_)
593                        | yaml_rust::Yaml::Null
594                        | yaml_rust::Yaml::BadValue => None,
595                    };
596                }
597
598                TRANSFORM_ON_FAILURE => {
599                    let on_failure_str = yaml_string(v, TRANSFORM_ON_FAILURE)?;
600                    on_failure = Some(on_failure_str.parse()?);
601                }
602
603                _ => {}
604            }
605        }
606
607        // ensure fields and type
608        ensure!(!fields.is_empty(), TransformFieldMustBeSetSnafu);
609        let type_ = type_.context(TransformTypeMustBeSetSnafu {
610            fields: format!("{:?}", fields),
611        })?;
612
613        let (index, index_options) = match index_value {
614            Some(value) => {
615                let (index, raw_index_options) = parse_transform_index(value)?;
616                let index_options =
617                    lower_transform_index_options(index, &type_, raw_index_options)?;
618                (Some(index), index_options)
619            }
620            None => (None, None),
621        };
622
623        let final_default = if let Some(default_value) = default {
624            let target = parse_str_value(&type_, &default_value)?;
625            on_failure = Some(OnFailure::Default);
626            Some(target)
627        } else {
628            None
629        };
630
631        let builder = Transform {
632            fields,
633            type_,
634            json_settings,
635            default: final_default,
636            index,
637            index_options,
638            on_failure,
639            tag,
640        };
641
642        Ok(builder)
643    }
644}
645
646#[cfg(test)]
647mod tests {
648    use yaml_rust::YamlLoader;
649
650    use super::*;
651
652    fn parse_transform(yaml: &str) -> Result<Transform> {
653        let docs = YamlLoader::load_from_str(yaml).unwrap();
654        docs[0].as_hash().unwrap().try_into()
655    }
656
657    #[test]
658    fn test_transform_parses_json2_type_hints() {
659        let transform = parse_transform(
660            r#"
661field: payload
662type:
663  json2:
664    - path: "user.id"
665      type: int64
666      index:
667        type: inverted
668    - path: 'attrs."http.status_code"'
669      type: string
670"#,
671        )
672        .unwrap();
673
674        assert_eq!(transform.type_, ColumnDataType::Json);
675        let hints = transform.json_settings.as_ref().unwrap().type_hints();
676        assert_eq!(hints.len(), 2);
677        assert_eq!(hints[0].path, ["user", "id"]);
678        assert_eq!(
679            hints[0].data_type,
680            datatypes::prelude::ConcreteDataType::int64_datatype()
681        );
682        assert!(hints[0].inverted_index);
683        assert_eq!(hints[1].path, ["attrs", "http.status_code"]);
684    }
685
686    #[test]
687    fn test_transform_rejects_json2_nullable_and_default() {
688        for option in ["nullable: false", "default: 7"] {
689            let err = parse_transform(&format!(
690                r#"
691field: payload
692type:
693  json2:
694    - path: score
695      type: float64
696      {option}
697"#,
698            ))
699            .unwrap_err();
700
701            assert!(err.to_string().contains("unsupported field"), "{err}");
702        }
703    }
704
705    #[test]
706    fn test_transform_parses_legacy_string_index() {
707        let transform = parse_transform(
708            r#"
709field: message
710type: string
711index: fulltext
712"#,
713        )
714        .unwrap();
715
716        assert_eq!(transform.index, Some(Index::Fulltext));
717        assert!(transform.index_options.is_none());
718    }
719
720    #[test]
721    fn test_transform_parses_index_object_without_options() {
722        let transform = parse_transform(
723            r#"
724field: message
725type: string
726index:
727  type: inverted
728"#,
729        )
730        .unwrap();
731
732        assert_eq!(transform.index, Some(Index::Inverted));
733        assert!(transform.index_options.is_none());
734    }
735
736    #[test]
737    fn test_transform_parses_index_object_with_scalar_options() {
738        let transform = parse_transform(
739            r#"
740field: message
741type: string
742index:
743  type: fulltext
744  options:
745    analyzer: English
746    case_sensitive: false
747    granularity: 2048
748    false_positive_rate: 0.02
749"#,
750        )
751        .unwrap();
752
753        assert_eq!(transform.index, Some(Index::Fulltext));
754        let options = transform.index_options.as_ref().unwrap();
755        let fulltext = options.as_fulltext().unwrap();
756        assert!(fulltext.enable);
757        assert_eq!(fulltext.analyzer.to_string(), "English");
758        assert!(!fulltext.case_sensitive);
759        assert_eq!(fulltext.granularity, 2048);
760        assert_eq!(fulltext.false_positive_rate(), 0.02);
761    }
762
763    #[test]
764    fn test_transform_rejects_invalid_index_options_type() {
765        let result = parse_transform(
766            r#"
767field: message
768type: string
769index:
770  type: fulltext
771  options: invalid
772"#,
773        );
774
775        assert!(result.is_err());
776    }
777
778    #[test]
779    fn test_transform_rejects_non_scalar_index_option_value() {
780        let result = parse_transform(
781            r#"
782field: message
783type: string
784index:
785  type: fulltext
786  options:
787    analyzer:
788      kind: English
789"#,
790        );
791
792        assert!(result.is_err());
793    }
794
795    #[test]
796    fn test_transform_rejects_unknown_index_field() {
797        let result = parse_transform(
798            r#"
799field: message
800type: string
801index:
802  type: fulltext
803  config: {}
804"#,
805        );
806
807        assert!(result.is_err());
808    }
809
810    #[test]
811    fn test_transform_rejects_unsupported_fulltext_option_key() {
812        let result = parse_transform(
813            r#"
814field: message
815type: string
816index:
817  type: fulltext
818  options:
819    tokenizer: english
820"#,
821        );
822
823        assert!(result.is_err());
824    }
825
826    #[test]
827    fn test_transform_rejects_options_for_inverted_index() {
828        let result = parse_transform(
829            r#"
830field: message
831type: string
832index:
833  type: inverted
834  options:
835    backend: bloom
836"#,
837        );
838
839        assert!(result.is_err());
840    }
841
842    #[test]
843    fn test_transform_rejects_empty_options_for_unsupported_indexes() {
844        for index in ["inverted", "time", "tag"] {
845            let yaml = format!(
846                r#"
847field: message
848type: string
849index:
850  type: {index}
851  options: {{}}
852"#
853            );
854
855            let result = parse_transform(&yaml);
856            assert!(
857                result.is_err(),
858                "expected `{index}` to reject empty options"
859            );
860        }
861    }
862
863    #[test]
864    fn test_transform_rejects_fulltext_index_on_non_string_column() {
865        let result = parse_transform(
866            r#"
867field: count
868type: int64
869index: fulltext
870"#,
871        );
872
873        assert!(result.is_err());
874    }
875
876    #[test]
877    fn test_transform_allows_skipping_index_on_numeric_column() {
878        let transform = parse_transform(
879            r#"
880field: count
881type: int64
882index:
883  type: skipping
884  options:
885    granularity: 2048
886    false_positive_rate: 0.02
887    type: BLOOM
888"#,
889        )
890        .unwrap();
891
892        assert_eq!(transform.index, Some(Index::Skipping));
893        let skipping = transform
894            .index_options
895            .as_ref()
896            .unwrap()
897            .as_skipping()
898            .unwrap();
899        assert_eq!(skipping.granularity, 2048);
900        assert_eq!(skipping.false_positive_rate(), 0.02);
901        assert_eq!(skipping.index_type.to_string(), "BLOOM");
902    }
903}