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