1use 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 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 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 {
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 {
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 {
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 {
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}