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