Skip to main content

pipeline/etl/transform/transformer/greptime/
coerce.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::sync::Arc;
16
17use api::v1::column_data_type_extension::TypeExt;
18use api::v1::column_def::{options_from_fulltext, options_from_inverted, options_from_skipping};
19use api::v1::{ColumnDataTypeExtension, ColumnOptions, JsonTypeExtension};
20use arrow_schema::extension::{
21    EXTENSION_TYPE_METADATA_KEY, EXTENSION_TYPE_NAME_KEY, ExtensionType,
22};
23use datatypes::extension::json::{Json2ExtensionType, JsonMetadata};
24use datatypes::json::JsonSettings;
25use datatypes::schema::{FulltextOptions, SkippingIndexOptions};
26use datatypes::value::Value;
27use greptime_proto::v1::value::ValueData;
28use greptime_proto::v1::{ColumnDataType, ColumnSchema, SemanticType};
29use snafu::{OptionExt, ResultExt, ensure};
30use vrl::value::Value as VrlValue;
31
32use crate::error::{
33    CoerceIncompatibleTypesSnafu, CoerceJsonTypeToSnafu, CoerceStringToTypeSnafu,
34    CoerceTypeToJsonSnafu, CoerceUnsupportedEpochTypeSnafu, ColumnOptionsSnafu, Error,
35    InvalidTimestampSnafu, Result, TransformIndexStateMismatchSnafu,
36    UnsupportedTypeInPipelineSnafu, VrlRegexValueSnafu,
37};
38use crate::etl::transform::index::Index;
39use crate::etl::transform::transformer::greptime::{
40    vrl_value_to_jsonb_value, vrl_value_to_serde_json,
41};
42use crate::etl::transform::{OnFailure, Transform, TransformIndexOptions};
43
44pub(crate) fn coerce_columns(transform: &Transform) -> Result<Vec<ColumnSchema>> {
45    let mut columns = Vec::new();
46
47    for field in transform.fields.iter() {
48        let column_name = field.target_or_input_field().to_string();
49
50        let ext = if matches!(transform.type_, ColumnDataType::Binary) {
51            Some(ColumnDataTypeExtension {
52                type_ext: Some(TypeExt::JsonType(JsonTypeExtension::JsonBinary.into())),
53            })
54        } else {
55            None
56        };
57
58        let semantic_type = coerce_semantic_type(transform) as i32;
59
60        let column = ColumnSchema {
61            column_name,
62            datatype: transform.type_ as i32,
63            semantic_type,
64            datatype_extension: ext,
65            options: coerce_options(transform)?,
66        };
67        columns.push(column);
68    }
69
70    Ok(columns)
71}
72
73fn coerce_semantic_type(transform: &Transform) -> SemanticType {
74    if transform.tag {
75        return SemanticType::Tag;
76    }
77
78    match transform.index {
79        Some(Index::Tag) => SemanticType::Tag,
80        Some(Index::Time) => SemanticType::Timestamp,
81        Some(Index::Fulltext) | Some(Index::Skipping) | Some(Index::Inverted) | None => {
82            SemanticType::Field
83        }
84    }
85}
86
87fn transform_index_label(index: Option<Index>) -> String {
88    index
89        .map(|index| index.to_string())
90        .unwrap_or_else(|| "none".to_string())
91}
92
93fn validate_transform_index_state(transform: &Transform) -> Result<()> {
94    let Some(index_options) = transform.index_options.as_ref() else {
95        return Ok(());
96    };
97
98    let options_index = index_options.index();
99    let index = transform.index;
100    ensure!(
101        index == Some(options_index),
102        TransformIndexStateMismatchSnafu {
103            index: transform_index_label(index),
104            options: options_index.to_string(),
105        }
106    );
107
108    Ok(())
109}
110
111fn build_fulltext_index_options(transform: &Transform) -> Result<FulltextOptions> {
112    match transform.index_options.as_ref() {
113        None => Ok(FulltextOptions {
114            enable: true,
115            ..Default::default()
116        }),
117        Some(TransformIndexOptions::Fulltext(options)) => Ok(options.clone()),
118        Some(options) => TransformIndexStateMismatchSnafu {
119            index: Index::Fulltext.to_string(),
120            options: options.index().to_string(),
121        }
122        .fail(),
123    }
124}
125
126fn build_skipping_index_options(transform: &Transform) -> Result<SkippingIndexOptions> {
127    match transform.index_options.as_ref() {
128        None => Ok(SkippingIndexOptions::default()),
129        Some(TransformIndexOptions::Skipping(options)) => Ok(options.clone()),
130        Some(options) => TransformIndexStateMismatchSnafu {
131            index: Index::Skipping.to_string(),
132            options: options.index().to_string(),
133        }
134        .fail(),
135    }
136}
137
138fn coerce_options(transform: &Transform) -> Result<Option<ColumnOptions>> {
139    validate_transform_index_state(transform)?;
140
141    let mut options = match transform.index {
142        Some(Index::Fulltext) => {
143            let options = build_fulltext_index_options(transform)?;
144            options_from_fulltext(&options).context(ColumnOptionsSnafu)
145        }
146        Some(Index::Skipping) => {
147            let options = build_skipping_index_options(transform)?;
148            options_from_skipping(&options).context(ColumnOptionsSnafu)
149        }
150        Some(Index::Inverted) => Ok(Some(options_from_inverted())),
151        _ => Ok(None),
152    }?;
153
154    if transform.type_ == ColumnDataType::Json {
155        let extension = Json2ExtensionType::new(Arc::new(JsonMetadata::new(
156            transform.json_settings.clone().unwrap_or_default(),
157        )));
158        let options = options.get_or_insert_default();
159        options.options.insert(
160            EXTENSION_TYPE_NAME_KEY.to_string(),
161            Json2ExtensionType::NAME.to_string(),
162        );
163        if let Some(metadata) = extension.serialize_metadata() {
164            options
165                .options
166                .insert(EXTENSION_TYPE_METADATA_KEY.to_string(), metadata);
167        }
168    }
169
170    Ok(options)
171}
172
173pub(crate) fn coerce_value(
174    val: &VrlValue,
175    transform: &Transform,
176    json_settings: Option<&JsonSettings>,
177) -> Result<Option<ValueData>> {
178    match val {
179        VrlValue::Null => Ok(None),
180        VrlValue::Integer(n) => coerce_i64_value(*n, transform),
181        VrlValue::Float(n) => coerce_f64_value(n.into_inner(), transform),
182        VrlValue::Boolean(b) => coerce_bool_value(*b, transform),
183        VrlValue::Bytes(b) => coerce_string_value(String::from_utf8_lossy(b).as_ref(), transform),
184        VrlValue::Timestamp(ts) => match transform.type_ {
185            ColumnDataType::TimestampNanosecond => Ok(Some(ValueData::TimestampNanosecondValue(
186                ts.timestamp_nanos_opt().context(InvalidTimestampSnafu {
187                    input: ts.to_rfc3339(),
188                })?,
189            ))),
190            ColumnDataType::TimestampMicrosecond => Ok(Some(ValueData::TimestampMicrosecondValue(
191                ts.timestamp_micros(),
192            ))),
193            ColumnDataType::TimestampMillisecond => Ok(Some(ValueData::TimestampMillisecondValue(
194                ts.timestamp_millis(),
195            ))),
196            ColumnDataType::TimestampSecond => {
197                Ok(Some(ValueData::TimestampSecondValue(ts.timestamp())))
198            }
199            _ => CoerceIncompatibleTypesSnafu {
200                msg: "Timestamp can only be coerced to another type",
201            }
202            .fail(),
203        },
204        VrlValue::Array(_) | VrlValue::Object(_) => {
205            coerce_json_value(val, transform, json_settings)
206        }
207        VrlValue::Regex(_) => VrlRegexValueSnafu.fail(),
208    }
209}
210
211fn coerce_bool_value(b: bool, transform: &Transform) -> Result<Option<ValueData>> {
212    let val = match transform.type_ {
213        ColumnDataType::Int8 => ValueData::I8Value(b as i32),
214        ColumnDataType::Int16 => ValueData::I16Value(b as i32),
215        ColumnDataType::Int32 => ValueData::I32Value(b as i32),
216        ColumnDataType::Int64 => ValueData::I64Value(b as i64),
217
218        ColumnDataType::Uint8 => ValueData::U8Value(b as u32),
219        ColumnDataType::Uint16 => ValueData::U16Value(b as u32),
220        ColumnDataType::Uint32 => ValueData::U32Value(b as u32),
221        ColumnDataType::Uint64 => ValueData::U64Value(b as u64),
222
223        ColumnDataType::Float32 => ValueData::F32Value(if b { 1.0 } else { 0.0 }),
224        ColumnDataType::Float64 => ValueData::F64Value(if b { 1.0 } else { 0.0 }),
225
226        ColumnDataType::Boolean => ValueData::BoolValue(b),
227        ColumnDataType::String => ValueData::StringValue(b.to_string()),
228
229        ColumnDataType::TimestampNanosecond
230        | ColumnDataType::TimestampMicrosecond
231        | ColumnDataType::TimestampMillisecond
232        | ColumnDataType::TimestampSecond => match transform.on_failure {
233            Some(OnFailure::Ignore) => return Ok(None),
234            Some(OnFailure::Default) => {
235                return CoerceUnsupportedEpochTypeSnafu { ty: "Default" }.fail();
236            }
237            None => {
238                return CoerceUnsupportedEpochTypeSnafu { ty: "Boolean" }.fail();
239            }
240        },
241
242        ColumnDataType::Binary | ColumnDataType::Json => {
243            return CoerceJsonTypeToSnafu {
244                ty: transform.type_.as_str_name(),
245            }
246            .fail();
247        }
248
249        _ => {
250            return UnsupportedTypeInPipelineSnafu {
251                ty: transform.type_.as_str_name(),
252            }
253            .fail();
254        }
255    };
256
257    Ok(Some(val))
258}
259
260fn handle_coercion_failure(transform: &Transform, error: Error) -> Result<Option<ValueData>> {
261    match transform.on_failure {
262        Some(OnFailure::Ignore) => Ok(None),
263        Some(OnFailure::Default) => match transform.get_default() {
264            Some(default) => Ok(Some(default.clone())),
265            None => transform.get_type_matched_default_val().map(Some),
266        },
267        None => Err(error),
268    }
269}
270
271fn integer_out_of_range<T: std::fmt::Display>(
272    value: T,
273    transform: &Transform,
274) -> Result<Option<ValueData>> {
275    handle_coercion_failure(
276        transform,
277        CoerceIncompatibleTypesSnafu {
278            msg: format!(
279                "integer value `{value}` is out of range for {}",
280                transform.type_.as_str_name()
281            ),
282        }
283        .build(),
284    )
285}
286
287fn coerce_i64_value(n: i64, transform: &Transform) -> Result<Option<ValueData>> {
288    let val = match transform.type_ {
289        ColumnDataType::Int8 => match i8::try_from(n) {
290            Ok(value) => ValueData::I8Value(value.into()),
291            Err(_) => return integer_out_of_range(n, transform),
292        },
293        ColumnDataType::Int16 => match i16::try_from(n) {
294            Ok(value) => ValueData::I16Value(value.into()),
295            Err(_) => return integer_out_of_range(n, transform),
296        },
297        ColumnDataType::Int32 => match i32::try_from(n) {
298            Ok(value) => ValueData::I32Value(value),
299            Err(_) => return integer_out_of_range(n, transform),
300        },
301        ColumnDataType::Int64 => ValueData::I64Value(n),
302
303        ColumnDataType::Uint8 => match u8::try_from(n) {
304            Ok(value) => ValueData::U8Value(value.into()),
305            Err(_) => return integer_out_of_range(n, transform),
306        },
307        ColumnDataType::Uint16 => match u16::try_from(n) {
308            Ok(value) => ValueData::U16Value(value.into()),
309            Err(_) => return integer_out_of_range(n, transform),
310        },
311        ColumnDataType::Uint32 => match u32::try_from(n) {
312            Ok(value) => ValueData::U32Value(value),
313            Err(_) => return integer_out_of_range(n, transform),
314        },
315        ColumnDataType::Uint64 => match u64::try_from(n) {
316            Ok(value) => ValueData::U64Value(value),
317            Err(_) => return integer_out_of_range(n, transform),
318        },
319
320        ColumnDataType::Float32 => ValueData::F32Value(n as f32),
321        ColumnDataType::Float64 => ValueData::F64Value(n as f64),
322
323        ColumnDataType::Boolean => ValueData::BoolValue(n != 0),
324        ColumnDataType::String => ValueData::StringValue(n.to_string()),
325
326        ColumnDataType::TimestampNanosecond => ValueData::TimestampNanosecondValue(n),
327        ColumnDataType::TimestampMicrosecond => ValueData::TimestampMicrosecondValue(n),
328        ColumnDataType::TimestampMillisecond => ValueData::TimestampMillisecondValue(n),
329        ColumnDataType::TimestampSecond => ValueData::TimestampSecondValue(n),
330
331        ColumnDataType::Binary | ColumnDataType::Json => {
332            return CoerceJsonTypeToSnafu {
333                ty: transform.type_.as_str_name(),
334            }
335            .fail();
336        }
337
338        _ => return Ok(None),
339    };
340
341    Ok(Some(val))
342}
343
344fn coerce_u64_value(n: u64, transform: &Transform) -> Result<Option<ValueData>> {
345    let val = match transform.type_ {
346        ColumnDataType::Int8 => match i8::try_from(n) {
347            Ok(value) => ValueData::I8Value(value.into()),
348            Err(_) => return integer_out_of_range(n, transform),
349        },
350        ColumnDataType::Int16 => match i16::try_from(n) {
351            Ok(value) => ValueData::I16Value(value.into()),
352            Err(_) => return integer_out_of_range(n, transform),
353        },
354        ColumnDataType::Int32 => match i32::try_from(n) {
355            Ok(value) => ValueData::I32Value(value),
356            Err(_) => return integer_out_of_range(n, transform),
357        },
358        ColumnDataType::Int64 => match i64::try_from(n) {
359            Ok(value) => ValueData::I64Value(value),
360            Err(_) => return integer_out_of_range(n, transform),
361        },
362
363        ColumnDataType::Uint8 => match u8::try_from(n) {
364            Ok(value) => ValueData::U8Value(value.into()),
365            Err(_) => return integer_out_of_range(n, transform),
366        },
367        ColumnDataType::Uint16 => match u16::try_from(n) {
368            Ok(value) => ValueData::U16Value(value.into()),
369            Err(_) => return integer_out_of_range(n, transform),
370        },
371        ColumnDataType::Uint32 => match u32::try_from(n) {
372            Ok(value) => ValueData::U32Value(value),
373            Err(_) => return integer_out_of_range(n, transform),
374        },
375        ColumnDataType::Uint64 => ValueData::U64Value(n),
376
377        ColumnDataType::Float32 => ValueData::F32Value(n as f32),
378        ColumnDataType::Float64 => ValueData::F64Value(n as f64),
379
380        ColumnDataType::Boolean => ValueData::BoolValue(n != 0),
381        ColumnDataType::String => ValueData::StringValue(n.to_string()),
382
383        ColumnDataType::TimestampNanosecond => match i64::try_from(n) {
384            Ok(value) => ValueData::TimestampNanosecondValue(value),
385            Err(_) => return integer_out_of_range(n, transform),
386        },
387        ColumnDataType::TimestampMicrosecond => match i64::try_from(n) {
388            Ok(value) => ValueData::TimestampMicrosecondValue(value),
389            Err(_) => return integer_out_of_range(n, transform),
390        },
391        ColumnDataType::TimestampMillisecond => match i64::try_from(n) {
392            Ok(value) => ValueData::TimestampMillisecondValue(value),
393            Err(_) => return integer_out_of_range(n, transform),
394        },
395        ColumnDataType::TimestampSecond => match i64::try_from(n) {
396            Ok(value) => ValueData::TimestampSecondValue(value),
397            Err(_) => return integer_out_of_range(n, transform),
398        },
399
400        ColumnDataType::Binary | ColumnDataType::Json => {
401            return CoerceJsonTypeToSnafu {
402                ty: transform.type_.as_str_name(),
403            }
404            .fail();
405        }
406
407        _ => return Ok(None),
408    };
409
410    Ok(Some(val))
411}
412
413fn coerce_f64_value(n: f64, transform: &Transform) -> Result<Option<ValueData>> {
414    let val = match transform.type_ {
415        ColumnDataType::Int8 => ValueData::I8Value(n as i32),
416        ColumnDataType::Int16 => ValueData::I16Value(n as i32),
417        ColumnDataType::Int32 => ValueData::I32Value(n as i32),
418        ColumnDataType::Int64 => ValueData::I64Value(n as i64),
419
420        ColumnDataType::Uint8 => ValueData::U8Value(n as u32),
421        ColumnDataType::Uint16 => ValueData::U16Value(n as u32),
422        ColumnDataType::Uint32 => ValueData::U32Value(n as u32),
423        ColumnDataType::Uint64 => ValueData::U64Value(n as u64),
424
425        ColumnDataType::Float32 => ValueData::F32Value(n as f32),
426        ColumnDataType::Float64 => ValueData::F64Value(n),
427
428        ColumnDataType::Boolean => ValueData::BoolValue(n != 0.0),
429        ColumnDataType::String => ValueData::StringValue(n.to_string()),
430
431        ColumnDataType::TimestampNanosecond
432        | ColumnDataType::TimestampMicrosecond
433        | ColumnDataType::TimestampMillisecond
434        | ColumnDataType::TimestampSecond => match transform.on_failure {
435            Some(OnFailure::Ignore) => return Ok(None),
436            Some(OnFailure::Default) => {
437                return CoerceUnsupportedEpochTypeSnafu { ty: "Default" }.fail();
438            }
439            None => {
440                return CoerceUnsupportedEpochTypeSnafu { ty: "Float" }.fail();
441            }
442        },
443
444        ColumnDataType::Binary | ColumnDataType::Json => {
445            return CoerceJsonTypeToSnafu {
446                ty: transform.type_.as_str_name(),
447            }
448            .fail();
449        }
450
451        _ => return Ok(None),
452    };
453
454    Ok(Some(val))
455}
456
457macro_rules! coerce_string_value {
458    ($s:expr, $transform:expr, $type:ident, $parse:ident) => {
459        match $s.parse::<$type>() {
460            Ok(v) => Ok(Some(ValueData::$parse(v.into()))),
461            Err(_) => handle_coercion_failure(
462                $transform,
463                CoerceStringToTypeSnafu {
464                    s: $s,
465                    ty: $transform.type_.as_str_name(),
466                }
467                .build(),
468            ),
469        }
470    };
471}
472
473fn coerce_string_value(s: &str, transform: &Transform) -> Result<Option<ValueData>> {
474    match transform.type_ {
475        ColumnDataType::Int8 => {
476            coerce_string_value!(s, transform, i8, I8Value)
477        }
478        ColumnDataType::Int16 => {
479            coerce_string_value!(s, transform, i16, I16Value)
480        }
481        ColumnDataType::Int32 => {
482            coerce_string_value!(s, transform, i32, I32Value)
483        }
484        ColumnDataType::Int64 => {
485            coerce_string_value!(s, transform, i64, I64Value)
486        }
487
488        ColumnDataType::Uint8 => {
489            coerce_string_value!(s, transform, u8, U8Value)
490        }
491        ColumnDataType::Uint16 => {
492            coerce_string_value!(s, transform, u16, U16Value)
493        }
494        ColumnDataType::Uint32 => {
495            coerce_string_value!(s, transform, u32, U32Value)
496        }
497        ColumnDataType::Uint64 => {
498            coerce_string_value!(s, transform, u64, U64Value)
499        }
500
501        ColumnDataType::Float32 => {
502            coerce_string_value!(s, transform, f32, F32Value)
503        }
504        ColumnDataType::Float64 => {
505            coerce_string_value!(s, transform, f64, F64Value)
506        }
507
508        ColumnDataType::Boolean => {
509            coerce_string_value!(s, transform, bool, BoolValue)
510        }
511
512        ColumnDataType::String => Ok(Some(ValueData::StringValue(s.to_string()))),
513
514        ColumnDataType::TimestampNanosecond
515        | ColumnDataType::TimestampMicrosecond
516        | ColumnDataType::TimestampMillisecond
517        | ColumnDataType::TimestampSecond => match transform.on_failure {
518            Some(OnFailure::Ignore) => Ok(None),
519            Some(OnFailure::Default) => CoerceUnsupportedEpochTypeSnafu { ty: "Default" }.fail(),
520            None => CoerceUnsupportedEpochTypeSnafu { ty: "String" }.fail(),
521        },
522
523        ColumnDataType::Binary | ColumnDataType::Json => CoerceStringToTypeSnafu {
524            s,
525            ty: transform.type_.as_str_name(),
526        }
527        .fail(),
528
529        _ => Ok(None),
530    }
531}
532
533fn coerce_json_value(
534    v: &VrlValue,
535    transform: &Transform,
536    json_settings: Option<&JsonSettings>,
537) -> Result<Option<ValueData>> {
538    let value = match transform.type_ {
539        ColumnDataType::Binary => {
540            let data: jsonb::Value = vrl_value_to_jsonb_value(v);
541            ValueData::BinaryValue(data.to_vec())
542        }
543        ColumnDataType::Json => {
544            let json = vrl_value_to_serde_json(v);
545            let encoded = if let Some(settings) = json_settings.or(transform.json_settings.as_ref())
546            {
547                settings.encode(json)
548            } else {
549                JsonSettings::default().encode(json)
550            };
551            let value = match encoded {
552                Ok(value) => value,
553                Err(error) => return handle_coercion_failure(transform, error.into()),
554            };
555            let Value::Json(value) = value else {
556                unreachable!()
557            };
558            ValueData::JsonValue(api::helper::encode_json_value(*value))
559        }
560        t => {
561            return CoerceTypeToJsonSnafu {
562                ty: t.as_str_name(),
563            }
564            .fail();
565        }
566    };
567    Ok(Some(value))
568}
569
570#[cfg(test)]
571mod tests {
572
573    use datatypes::data_type::ConcreteDataType;
574    use datatypes::json::JsonTypeHint;
575    use datatypes::schema::{FulltextAnalyzer, FulltextBackend, SkippingIndexType};
576    use vrl::prelude::Bytes;
577
578    use super::*;
579    use crate::etl::field::Fields;
580
581    fn transform(type_: ColumnDataType) -> Transform {
582        Transform {
583            fields: Fields::default(),
584            type_,
585            json_settings: None,
586            default: None,
587            index: None,
588            index_options: None,
589            on_failure: None,
590            tag: false,
591        }
592    }
593
594    fn narrow_i64_value(type_: ColumnDataType, value: i64) -> ValueData {
595        match type_ {
596            ColumnDataType::Int8 => ValueData::I8Value(value as i32),
597            ColumnDataType::Int16 => ValueData::I16Value(value as i32),
598            ColumnDataType::Int32 => ValueData::I32Value(value as i32),
599            ColumnDataType::Uint8 => ValueData::U8Value(value as u32),
600            ColumnDataType::Uint16 => ValueData::U16Value(value as u32),
601            ColumnDataType::Uint32 => ValueData::U32Value(value as u32),
602            _ => unreachable!("narrow integer type required"),
603        }
604    }
605
606    fn checked_u64_value(type_: ColumnDataType, value: u64) -> ValueData {
607        match type_ {
608            ColumnDataType::Int8 => ValueData::I8Value(value as i32),
609            ColumnDataType::Int16 => ValueData::I16Value(value as i32),
610            ColumnDataType::Int32 => ValueData::I32Value(value as i32),
611            ColumnDataType::Int64 => ValueData::I64Value(value as i64),
612            ColumnDataType::Uint8 => ValueData::U8Value(value as u32),
613            ColumnDataType::Uint16 => ValueData::U16Value(value as u32),
614            ColumnDataType::Uint32 => ValueData::U32Value(value as u32),
615            ColumnDataType::TimestampNanosecond => {
616                ValueData::TimestampNanosecondValue(value as i64)
617            }
618            ColumnDataType::TimestampMicrosecond => {
619                ValueData::TimestampMicrosecondValue(value as i64)
620            }
621            ColumnDataType::TimestampMillisecond => {
622                ValueData::TimestampMillisecondValue(value as i64)
623            }
624            ColumnDataType::TimestampSecond => ValueData::TimestampSecondValue(value as i64),
625            _ => unreachable!("checked u64 target type required"),
626        }
627    }
628
629    fn assert_out_of_range(
630        result: crate::error::Result<Option<ValueData>>,
631        value: impl std::fmt::Display,
632        type_: ColumnDataType,
633    ) {
634        let error = result.unwrap_err();
635        assert_eq!(
636            error.to_string(),
637            format!(
638                "Failed to coerce value: integer value `{value}` is out of range for {}",
639                type_.as_str_name()
640            )
641        );
642        assert!(matches!(
643            error,
644            crate::error::Error::CoerceIncompatibleTypes { .. }
645        ));
646    }
647
648    #[test]
649    fn test_coerce_i64_narrowing_boundaries_and_overflows() {
650        // target type, inclusive minimum, inclusive maximum
651        let cases = [
652            (ColumnDataType::Int8, i8::MIN as i64, i8::MAX as i64),
653            (ColumnDataType::Int16, i16::MIN as i64, i16::MAX as i64),
654            (ColumnDataType::Int32, i32::MIN as i64, i32::MAX as i64),
655            (ColumnDataType::Uint8, u8::MIN as i64, u8::MAX as i64),
656            (ColumnDataType::Uint16, u16::MIN as i64, u16::MAX as i64),
657            (ColumnDataType::Uint32, u32::MIN as i64, u32::MAX as i64),
658        ];
659
660        for (type_, min, max) in cases {
661            let transform = transform(type_);
662            assert_eq!(
663                coerce_i64_value(min, &transform).unwrap(),
664                Some(narrow_i64_value(type_, min)),
665                "{type_:?} minimum"
666            );
667            assert_eq!(
668                coerce_i64_value(max, &transform).unwrap(),
669                Some(narrow_i64_value(type_, max)),
670                "{type_:?} maximum"
671            );
672            assert_out_of_range(coerce_i64_value(min - 1, &transform), min - 1, type_);
673            assert_out_of_range(coerce_i64_value(max + 1, &transform), max + 1, type_);
674        }
675    }
676
677    #[test]
678    fn test_coerce_i64_to_u64_rejects_negative_value() {
679        assert_out_of_range(
680            coerce_i64_value(-1, &transform(ColumnDataType::Uint64)),
681            -1,
682            ColumnDataType::Uint64,
683        );
684    }
685
686    #[test]
687    fn test_coerce_u64_checked_conversions() {
688        // target type, inclusive maximum
689        let cases = [
690            (ColumnDataType::Int8, i8::MAX as u64),
691            (ColumnDataType::Int16, i16::MAX as u64),
692            (ColumnDataType::Int32, i32::MAX as u64),
693            (ColumnDataType::Int64, i64::MAX as u64),
694            (ColumnDataType::Uint8, u8::MAX as u64),
695            (ColumnDataType::Uint16, u16::MAX as u64),
696            (ColumnDataType::Uint32, u32::MAX as u64),
697            (ColumnDataType::TimestampNanosecond, i64::MAX as u64),
698            (ColumnDataType::TimestampMicrosecond, i64::MAX as u64),
699            (ColumnDataType::TimestampMillisecond, i64::MAX as u64),
700            (ColumnDataType::TimestampSecond, i64::MAX as u64),
701        ];
702
703        for (type_, max) in cases {
704            let transform = transform(type_);
705            assert_eq!(
706                coerce_u64_value(max, &transform).unwrap(),
707                Some(checked_u64_value(type_, max)),
708                "{type_:?} maximum"
709            );
710            assert_out_of_range(coerce_u64_value(max + 1, &transform), max + 1, type_);
711        }
712
713        assert_eq!(
714            coerce_u64_value(u64::MAX, &transform(ColumnDataType::Uint64)).unwrap(),
715            Some(ValueData::U64Value(u64::MAX))
716        );
717        assert_out_of_range(
718            coerce_u64_value(u64::MAX, &transform(ColumnDataType::TimestampNanosecond)),
719            u64::MAX,
720            ColumnDataType::TimestampNanosecond,
721        );
722    }
723
724    #[test]
725    fn test_coerce_string_without_on_failure() {
726        let transform = Transform {
727            fields: Fields::default(),
728            type_: ColumnDataType::Int32,
729            json_settings: None,
730            default: None,
731            index: None,
732            index_options: None,
733            on_failure: None,
734            tag: false,
735        };
736
737        // valid string
738        {
739            let val = VrlValue::Integer(123);
740            let result = coerce_value(&val, &transform, None).unwrap();
741            assert_eq!(result, Some(ValueData::I32Value(123)));
742        }
743
744        // invalid string
745        {
746            let val = VrlValue::Bytes(Bytes::from("hello"));
747            let result = coerce_value(&val, &transform, None);
748            assert!(result.is_err());
749        }
750    }
751
752    #[test]
753    fn test_coerce_string_with_on_failure_ignore() {
754        let transform = Transform {
755            fields: Fields::default(),
756            type_: ColumnDataType::Int32,
757            json_settings: None,
758            default: None,
759            index: None,
760            index_options: None,
761            on_failure: Some(OnFailure::Ignore),
762            tag: false,
763        };
764
765        let val = VrlValue::Bytes(Bytes::from("hello"));
766        let result = coerce_value(&val, &transform, None).unwrap();
767        assert_eq!(result, None);
768    }
769
770    #[test]
771    fn test_coerce_json2_with_on_failure() {
772        let settings = JsonSettings::try_new(
773            vec![JsonTypeHint {
774                path: vec!["age".to_string()],
775                data_type: ConcreteDataType::int64_datatype(),
776                nullable: false,
777                default_constraint: None,
778                inverted_index: false,
779            }],
780            None,
781        )
782        .unwrap();
783        let mut transform = transform(ColumnDataType::Json);
784        transform.json_settings = Some(settings);
785        transform.on_failure = Some(OnFailure::Ignore);
786        let value: VrlValue = serde_json::json!({"age": "42"}).into();
787
788        assert_eq!(coerce_value(&value, &transform, None).unwrap(), None);
789
790        transform.on_failure = Some(OnFailure::Default);
791        assert_eq!(
792            coerce_value(&value, &transform, None).unwrap(),
793            Some(ValueData::JsonValue(Default::default()))
794        );
795    }
796
797    #[test]
798    fn test_coerce_string_with_on_failure_default() {
799        let mut transform = Transform {
800            fields: Fields::default(),
801            type_: ColumnDataType::Int32,
802            json_settings: None,
803            default: None,
804            index: None,
805            index_options: None,
806            on_failure: Some(OnFailure::Default),
807            tag: false,
808        };
809
810        // with no explicit default value
811        {
812            let val = VrlValue::Bytes(Bytes::from("hello"));
813            let result = coerce_value(&val, &transform, None).unwrap();
814            assert_eq!(result, Some(ValueData::I32Value(0)));
815        }
816
817        // with explicit default value
818        {
819            transform.default = Some(ValueData::I32Value(42));
820            let val = VrlValue::Bytes(Bytes::from("hello"));
821            let result = coerce_value(&val, &transform, None).unwrap();
822            assert_eq!(result, Some(ValueData::I32Value(42)));
823        }
824    }
825
826    #[test]
827    fn test_coerce_fulltext_options_with_custom_values() {
828        let transform = Transform {
829            fields: Fields::default(),
830            type_: ColumnDataType::String,
831            json_settings: None,
832            default: None,
833            index: Some(Index::Fulltext),
834            index_options: Some(TransformIndexOptions::Fulltext(
835                FulltextOptions::new_unchecked(
836                    true,
837                    FulltextAnalyzer::Chinese,
838                    true,
839                    FulltextBackend::Tantivy,
840                    10240,
841                    0.01,
842                ),
843            )),
844            on_failure: None,
845            tag: false,
846        };
847
848        let options = coerce_options(&transform).unwrap().unwrap();
849        let fulltext: FulltextOptions =
850            serde_json::from_str(options.options.get("fulltext").unwrap()).unwrap();
851
852        assert!(fulltext.enable);
853        assert_eq!(fulltext.analyzer.to_string(), "Chinese");
854        assert!(fulltext.case_sensitive);
855        assert_eq!(fulltext.backend.to_string(), "tantivy");
856    }
857
858    #[test]
859    fn test_coerce_skipping_options_with_custom_values() {
860        let transform = Transform {
861            fields: Fields::default(),
862            type_: ColumnDataType::Int64,
863            json_settings: None,
864            default: None,
865            index: Some(Index::Skipping),
866            index_options: Some(TransformIndexOptions::Skipping(
867                SkippingIndexOptions::new_unchecked(2048, 0.02, SkippingIndexType::BloomFilter),
868            )),
869            on_failure: None,
870            tag: false,
871        };
872
873        let options = coerce_options(&transform).unwrap().unwrap();
874        let skipping: SkippingIndexOptions =
875            serde_json::from_str(options.options.get("skipping_index").unwrap()).unwrap();
876
877        assert_eq!(skipping.granularity, 2048);
878        assert_eq!(skipping.false_positive_rate(), 0.02);
879        assert_eq!(skipping.index_type.to_string(), "BLOOM");
880    }
881
882    #[test]
883    fn test_coerce_rejects_mismatched_index_options() {
884        let transform = Transform {
885            fields: Fields::default(),
886            type_: ColumnDataType::String,
887            json_settings: None,
888            default: None,
889            index: Some(Index::Fulltext),
890            index_options: Some(TransformIndexOptions::Skipping(
891                SkippingIndexOptions::new_unchecked(2048, 0.02, SkippingIndexType::BloomFilter),
892            )),
893            on_failure: None,
894            tag: false,
895        };
896
897        assert!(coerce_options(&transform).is_err());
898    }
899
900    #[test]
901    fn test_coerce_rejects_index_options_without_index() {
902        let transform = Transform {
903            fields: Fields::default(),
904            type_: ColumnDataType::String,
905            json_settings: None,
906            default: None,
907            index: None,
908            index_options: Some(TransformIndexOptions::Fulltext(
909                FulltextOptions::new_unchecked(
910                    true,
911                    FulltextAnalyzer::Chinese,
912                    true,
913                    FulltextBackend::Tantivy,
914                    10240,
915                    0.01,
916                ),
917            )),
918            on_failure: None,
919            tag: false,
920        };
921
922        assert!(coerce_options(&transform).is_err());
923    }
924}