1use std::collections::HashSet;
2
3use bytes::{Bytes, BytesMut};
4use chrono::{DateTime, Utc};
5use dyn_clone::DynClone;
6use ordered_float::NotNan;
7use vector_lib::configurable::configurable_component;
8use vrl::compiler::value::VrlValueArithmetic;
9use vrl::path::{OwnedSegment, OwnedTargetPath};
10
11use crate::event::{LogEvent, Value};
12
13#[configurable_component]
15#[derive(Clone, Debug, PartialEq)]
16#[cfg_attr(feature = "proptest", derive(proptest_derive::Arbitrary))]
17#[serde(rename_all = "snake_case")]
18pub enum MergeStrategy {
19 Discard,
21
22 Retain,
26
27 Sum,
29
30 Max,
32
33 Min,
35
36 Array,
38
39 Concat,
41
42 ConcatNewline,
44
45 ConcatRaw,
47
48 ShortestArray,
50
51 LongestArray,
53
54 FlatUnique,
56}
57
58#[derive(Debug, Clone)]
59struct DiscardMerger {
60 v: Value,
61}
62
63impl DiscardMerger {
64 const fn new(v: Value) -> Self {
65 Self { v }
66 }
67}
68
69impl ReduceValueMerger for DiscardMerger {
70 fn add(&mut self, _v: Value) -> Result<(), String> {
71 Ok(())
72 }
73
74 fn insert_into(
75 self: Box<Self>,
76 path: &OwnedTargetPath,
77 v: &mut LogEvent,
78 ) -> Result<(), String> {
79 v.insert(path, self.v);
80 Ok(())
81 }
82}
83
84#[derive(Debug, Clone)]
85struct RetainMerger {
86 v: Value,
87}
88
89impl RetainMerger {
90 #[allow(clippy::missing_const_for_fn)] fn new(v: Value) -> Self {
92 Self { v }
93 }
94}
95
96impl ReduceValueMerger for RetainMerger {
97 fn add(&mut self, v: Value) -> Result<(), String> {
98 if Value::Null != v {
99 self.v = v;
100 }
101 Ok(())
102 }
103
104 fn insert_into(
105 self: Box<Self>,
106 path: &OwnedTargetPath,
107 v: &mut LogEvent,
108 ) -> Result<(), String> {
109 v.insert(path, self.v);
110 Ok(())
111 }
112}
113
114#[derive(Debug, Clone)]
115struct ConcatMerger {
116 v: BytesMut,
117 join_by: Option<Vec<u8>>,
118}
119
120impl ConcatMerger {
121 fn new(v: Bytes, join_by: Option<char>) -> Self {
122 let join_by = join_by.map(|c| c.to_string().into_bytes());
124
125 Self {
126 v: BytesMut::from(&v[..]),
127 join_by,
128 }
129 }
130}
131
132impl ReduceValueMerger for ConcatMerger {
133 fn add(&mut self, v: Value) -> Result<(), String> {
134 if let Value::Bytes(b) = v {
135 if let Some(buf) = self.join_by.as_ref() {
136 self.v.extend(&buf[..]);
137 }
138 self.v.extend_from_slice(&b);
139 Ok(())
140 } else {
141 Err(format!(
142 "expected string value, found: '{}'",
143 v.to_string_lossy()
144 ))
145 }
146 }
147
148 fn insert_into(
149 self: Box<Self>,
150 path: &OwnedTargetPath,
151 v: &mut LogEvent,
152 ) -> Result<(), String> {
153 v.insert(path, Value::Bytes(self.v.into()));
154 Ok(())
155 }
156}
157
158#[derive(Debug, Clone)]
159struct ConcatArrayMerger {
160 v: Vec<Value>,
161}
162
163impl ConcatArrayMerger {
164 const fn new(v: Vec<Value>) -> Self {
165 Self { v }
166 }
167}
168
169impl ReduceValueMerger for ConcatArrayMerger {
170 fn add(&mut self, v: Value) -> Result<(), String> {
171 if let Value::Array(a) = v {
172 self.v.extend_from_slice(&a);
173 } else {
174 self.v.push(v);
175 }
176 Ok(())
177 }
178
179 fn insert_into(
180 self: Box<Self>,
181 path: &OwnedTargetPath,
182 v: &mut LogEvent,
183 ) -> Result<(), String> {
184 v.insert(path, Value::Array(self.v));
185 Ok(())
186 }
187}
188
189#[derive(Debug, Clone)]
190struct ArrayMerger {
191 v: Vec<Value>,
192}
193
194impl ArrayMerger {
195 fn new(v: Value) -> Self {
196 Self { v: vec![v] }
197 }
198}
199
200impl ReduceValueMerger for ArrayMerger {
201 fn add(&mut self, v: Value) -> Result<(), String> {
202 self.v.push(v);
203 Ok(())
204 }
205
206 fn insert_into(
207 self: Box<Self>,
208 path: &OwnedTargetPath,
209 v: &mut LogEvent,
210 ) -> Result<(), String> {
211 v.insert(path, Value::Array(self.v));
212 Ok(())
213 }
214}
215
216#[derive(Debug, Clone)]
217struct LongestArrayMerger {
218 v: Vec<Value>,
219}
220
221impl LongestArrayMerger {
222 const fn new(v: Vec<Value>) -> Self {
223 Self { v }
224 }
225}
226
227impl ReduceValueMerger for LongestArrayMerger {
228 fn add(&mut self, v: Value) -> Result<(), String> {
229 if let Value::Array(a) = v {
230 if a.len() > self.v.len() {
231 self.v = a;
232 }
233 Ok(())
234 } else {
235 Err(format!(
236 "expected array value, found: '{}'",
237 v.to_string_lossy()
238 ))
239 }
240 }
241
242 fn insert_into(
243 self: Box<Self>,
244 path: &OwnedTargetPath,
245 v: &mut LogEvent,
246 ) -> Result<(), String> {
247 v.insert(path, Value::Array(self.v));
248 Ok(())
249 }
250}
251
252#[derive(Debug, Clone)]
253struct ShortestArrayMerger {
254 v: Vec<Value>,
255}
256
257impl ShortestArrayMerger {
258 const fn new(v: Vec<Value>) -> Self {
259 Self { v }
260 }
261}
262
263impl ReduceValueMerger for ShortestArrayMerger {
264 fn add(&mut self, v: Value) -> Result<(), String> {
265 if let Value::Array(a) = v {
266 if a.len() < self.v.len() {
267 self.v = a;
268 }
269 Ok(())
270 } else {
271 Err(format!(
272 "expected array value, found: '{}'",
273 v.to_string_lossy()
274 ))
275 }
276 }
277
278 fn insert_into(
279 self: Box<Self>,
280 path: &OwnedTargetPath,
281 v: &mut LogEvent,
282 ) -> Result<(), String> {
283 v.insert(path, Value::Array(self.v));
284 Ok(())
285 }
286}
287
288#[derive(Debug, Clone)]
289struct FlatUniqueMerger {
290 v: HashSet<Value>,
291}
292
293#[allow(clippy::mutable_key_type)] fn insert_value(h: &mut HashSet<Value>, v: Value) {
295 match v {
296 Value::Object(m) => {
297 for (_, v) in m {
298 h.insert(v);
299 }
300 }
301 Value::Array(vec) => {
302 for v in vec {
303 h.insert(v);
304 }
305 }
306 _ => {
307 h.insert(v);
308 }
309 }
310}
311
312impl FlatUniqueMerger {
313 #[allow(clippy::mutable_key_type)] fn new(v: Value) -> Self {
315 let mut h = HashSet::default();
316 insert_value(&mut h, v);
317 Self { v: h }
318 }
319}
320
321impl ReduceValueMerger for FlatUniqueMerger {
322 fn add(&mut self, v: Value) -> Result<(), String> {
323 insert_value(&mut self.v, v);
324 Ok(())
325 }
326
327 fn insert_into(
328 self: Box<Self>,
329 path: &OwnedTargetPath,
330 v: &mut LogEvent,
331 ) -> Result<(), String> {
332 v.insert(path, Value::Array(self.v.into_iter().collect()));
333 Ok(())
334 }
335}
336
337#[derive(Debug, Clone)]
338struct TimestampWindowMerger {
339 started: DateTime<Utc>,
340 latest: DateTime<Utc>,
341}
342
343impl TimestampWindowMerger {
344 const fn new(v: DateTime<Utc>) -> Self {
345 Self {
346 started: v,
347 latest: v,
348 }
349 }
350}
351
352impl ReduceValueMerger for TimestampWindowMerger {
353 fn add(&mut self, v: Value) -> Result<(), String> {
354 if let Value::Timestamp(ts) = v {
355 self.latest = ts
356 } else {
357 return Err(format!(
358 "expected timestamp value, found: {}",
359 v.to_string_lossy()
360 ));
361 }
362 Ok(())
363 }
364
365 fn insert_into(
366 self: Box<Self>,
367 path: &OwnedTargetPath,
368 v: &mut LogEvent,
369 ) -> Result<(), String> {
370 let mut end_path = path.clone();
374 if let Some(OwnedSegment::Field(last)) = end_path.path.segments.last_mut() {
375 *last = format!("{last}_end").into();
376 }
377 v.insert(&end_path, Value::Timestamp(self.latest));
378 v.insert(path, Value::Timestamp(self.started));
379 Ok(())
380 }
381}
382
383#[derive(Debug, Clone)]
384enum NumberMergerValue {
385 Int(i64),
386 Float(NotNan<f64>),
387}
388
389impl From<i64> for NumberMergerValue {
390 fn from(v: i64) -> Self {
391 NumberMergerValue::Int(v)
392 }
393}
394
395impl From<NotNan<f64>> for NumberMergerValue {
396 fn from(v: NotNan<f64>) -> Self {
397 NumberMergerValue::Float(v)
398 }
399}
400
401#[derive(Debug, Clone)]
402struct AddNumbersMerger {
403 v: NumberMergerValue,
404}
405
406impl AddNumbersMerger {
407 const fn new(v: NumberMergerValue) -> Self {
408 Self { v }
409 }
410}
411
412impl ReduceValueMerger for AddNumbersMerger {
413 fn add(&mut self, v: Value) -> Result<(), String> {
414 let current = match &self.v {
417 NumberMergerValue::Int(value) => Value::Integer(*value),
418 NumberMergerValue::Float(value) => Value::Float(*value),
419 };
420 let sum = current.try_add(v).map_err(|error| error.to_string())?;
421 self.v = match sum {
422 Value::Integer(value) => NumberMergerValue::Int(value),
423 Value::Float(value) => NumberMergerValue::Float(value),
424 _ => return Err("summing numeric values produced a non-numeric value".to_owned()),
425 };
426
427 Ok(())
428 }
429
430 fn insert_into(
431 self: Box<Self>,
432 path: &OwnedTargetPath,
433 v: &mut LogEvent,
434 ) -> Result<(), String> {
435 match self.v {
436 NumberMergerValue::Float(f) => v.insert(path, Value::Float(f)),
437 NumberMergerValue::Int(i) => v.insert(path, Value::Integer(i)),
438 };
439 Ok(())
440 }
441}
442
443#[derive(Debug, Clone)]
444struct MaxNumberMerger {
445 v: NumberMergerValue,
446}
447
448impl MaxNumberMerger {
449 const fn new(v: NumberMergerValue) -> Self {
450 Self { v }
451 }
452}
453
454impl ReduceValueMerger for MaxNumberMerger {
455 fn add(&mut self, v: Value) -> Result<(), String> {
456 match v {
459 Value::Integer(i) => {
460 match self.v {
461 NumberMergerValue::Int(i2) => {
462 if i > i2 {
463 self.v = NumberMergerValue::Int(i);
464 }
465 }
466 NumberMergerValue::Float(f2) => {
467 let f = NotNan::new(i as f64).unwrap();
468 if f > f2 {
469 self.v = NumberMergerValue::Float(f);
470 }
471 }
472 };
473 }
474 Value::Float(f) => {
475 let f2 = match self.v {
476 NumberMergerValue::Int(i2) => NotNan::new(i2 as f64).unwrap(),
477 NumberMergerValue::Float(f2) => f2,
478 };
479 if f > f2 {
480 self.v = NumberMergerValue::Float(f);
481 }
482 }
483 _ => {
484 return Err(format!(
485 "expected numeric value, found: '{}'",
486 v.to_string_lossy()
487 ));
488 }
489 }
490 Ok(())
491 }
492
493 fn insert_into(
494 self: Box<Self>,
495 path: &OwnedTargetPath,
496 v: &mut LogEvent,
497 ) -> Result<(), String> {
498 match self.v {
499 NumberMergerValue::Float(f) => v.insert(path, Value::Float(f)),
500 NumberMergerValue::Int(i) => v.insert(path, Value::Integer(i)),
501 };
502 Ok(())
503 }
504}
505
506#[derive(Debug, Clone)]
507struct MinNumberMerger {
508 v: NumberMergerValue,
509}
510
511impl MinNumberMerger {
512 const fn new(v: NumberMergerValue) -> Self {
513 Self { v }
514 }
515}
516
517impl ReduceValueMerger for MinNumberMerger {
518 fn add(&mut self, v: Value) -> Result<(), String> {
519 match v {
522 Value::Integer(i) => {
523 match self.v {
524 NumberMergerValue::Int(i2) => {
525 if i < i2 {
526 self.v = NumberMergerValue::Int(i);
527 }
528 }
529 NumberMergerValue::Float(f2) => {
530 let f = NotNan::new(i as f64).unwrap();
531 if f < f2 {
532 self.v = NumberMergerValue::Float(f);
533 }
534 }
535 };
536 }
537 Value::Float(f) => {
538 let f2 = match self.v {
539 NumberMergerValue::Int(i2) => NotNan::new(i2 as f64).unwrap(),
540 NumberMergerValue::Float(f2) => f2,
541 };
542 if f < f2 {
543 self.v = NumberMergerValue::Float(f);
544 }
545 }
546 _ => {
547 return Err(format!(
548 "expected numeric value, found: '{}'",
549 v.to_string_lossy()
550 ));
551 }
552 }
553 Ok(())
554 }
555
556 fn insert_into(
557 self: Box<Self>,
558 path: &OwnedTargetPath,
559 v: &mut LogEvent,
560 ) -> Result<(), String> {
561 match self.v {
562 NumberMergerValue::Float(f) => v.insert(path, Value::Float(f)),
563 NumberMergerValue::Int(i) => v.insert(path, Value::Integer(i)),
564 };
565 Ok(())
566 }
567}
568
569pub trait ReduceValueMerger: std::fmt::Debug + Send + Sync + DynClone {
570 fn add(&mut self, v: Value) -> Result<(), String>;
571 fn insert_into(self: Box<Self>, path: &OwnedTargetPath, v: &mut LogEvent)
572 -> Result<(), String>;
573}
574
575dyn_clone::clone_trait_object!(ReduceValueMerger);
576
577impl From<Value> for Box<dyn ReduceValueMerger> {
578 fn from(v: Value) -> Self {
579 match v {
580 Value::Integer(i) => Box::new(AddNumbersMerger::new(i.into())),
581 Value::Float(f) => Box::new(AddNumbersMerger::new(f.into())),
582 Value::Timestamp(ts) => Box::new(TimestampWindowMerger::new(ts)),
583 Value::Object(_) => Box::new(DiscardMerger::new(v)),
584 Value::Null => Box::new(DiscardMerger::new(v)),
585 Value::Boolean(_) => Box::new(DiscardMerger::new(v)),
586 Value::Bytes(_) => Box::new(DiscardMerger::new(v)),
587 Value::Regex(_) => Box::new(DiscardMerger::new(v)),
588 Value::Array(_) => Box::new(DiscardMerger::new(v)),
589 }
590 }
591}
592
593pub(crate) fn get_value_merger(
594 v: Value,
595 m: &MergeStrategy,
596) -> Result<Box<dyn ReduceValueMerger>, String> {
597 match m {
598 MergeStrategy::Sum => match v {
599 Value::Integer(i) => Ok(Box::new(AddNumbersMerger::new(i.into()))),
600 Value::Float(f) => Ok(Box::new(AddNumbersMerger::new(f.into()))),
601 _ => Err(format!(
602 "expected number value, found: '{}'",
603 v.to_string_lossy()
604 )),
605 },
606 MergeStrategy::Max => match v {
607 Value::Integer(i) => Ok(Box::new(MaxNumberMerger::new(i.into()))),
608 Value::Float(f) => Ok(Box::new(MaxNumberMerger::new(f.into()))),
609 _ => Err(format!(
610 "expected number value, found: '{}'",
611 v.to_string_lossy()
612 )),
613 },
614 MergeStrategy::Min => match v {
615 Value::Integer(i) => Ok(Box::new(MinNumberMerger::new(i.into()))),
616 Value::Float(f) => Ok(Box::new(MinNumberMerger::new(f.into()))),
617 _ => Err(format!(
618 "expected number value, found: '{}'",
619 v.to_string_lossy()
620 )),
621 },
622 MergeStrategy::Concat => match v {
623 Value::Bytes(b) => Ok(Box::new(ConcatMerger::new(b, Some(' ')))),
624 Value::Array(a) => Ok(Box::new(ConcatArrayMerger::new(a))),
625 _ => Err(format!(
626 "expected string or array value, found: '{}'",
627 v.to_string_lossy()
628 )),
629 },
630 MergeStrategy::ConcatNewline => match v {
631 Value::Bytes(b) => Ok(Box::new(ConcatMerger::new(b, Some('\n')))),
632 _ => Err(format!(
633 "expected string value, found: '{}'",
634 v.to_string_lossy()
635 )),
636 },
637 MergeStrategy::ConcatRaw => match v {
638 Value::Bytes(b) => Ok(Box::new(ConcatMerger::new(b, None))),
639 _ => Err(format!(
640 "expected string value, found: '{}'",
641 v.to_string_lossy()
642 )),
643 },
644 MergeStrategy::Array => Ok(Box::new(ArrayMerger::new(v))),
645 MergeStrategy::ShortestArray => match v {
646 Value::Array(a) => Ok(Box::new(ShortestArrayMerger::new(a))),
647 _ => Err(format!(
648 "expected array value, found: '{}'",
649 v.to_string_lossy()
650 )),
651 },
652 MergeStrategy::LongestArray => match v {
653 Value::Array(a) => Ok(Box::new(LongestArrayMerger::new(a))),
654 _ => Err(format!(
655 "expected array value, found: '{}'",
656 v.to_string_lossy()
657 )),
658 },
659 MergeStrategy::Discard => Ok(Box::new(DiscardMerger::new(v))),
660 MergeStrategy::Retain => Ok(Box::new(RetainMerger::new(v))),
661 MergeStrategy::FlatUnique => Ok(Box::new(FlatUniqueMerger::new(v))),
662 }
663}
664
665#[cfg(test)]
666mod test {
667 use serde_json::json;
668 use vrl::owned_event_path;
669
670 use super::*;
671 use crate::event::LogEvent;
672
673 #[test]
674 fn initial_values() {
675 assert!(get_value_merger("foo".into(), &MergeStrategy::Discard).is_ok());
676 assert!(get_value_merger("foo".into(), &MergeStrategy::Retain).is_ok());
677 assert!(get_value_merger("foo".into(), &MergeStrategy::Sum).is_err());
678 assert!(get_value_merger("foo".into(), &MergeStrategy::Max).is_err());
679 assert!(get_value_merger("foo".into(), &MergeStrategy::Min).is_err());
680 assert!(get_value_merger("foo".into(), &MergeStrategy::Array).is_ok());
681 assert!(get_value_merger("foo".into(), &MergeStrategy::LongestArray).is_err());
682 assert!(get_value_merger("foo".into(), &MergeStrategy::ShortestArray).is_err());
683 assert!(get_value_merger("foo".into(), &MergeStrategy::Concat).is_ok());
684 assert!(get_value_merger("foo".into(), &MergeStrategy::ConcatNewline).is_ok());
685 assert!(get_value_merger("foo".into(), &MergeStrategy::ConcatRaw).is_ok());
686 assert!(get_value_merger("foo".into(), &MergeStrategy::FlatUnique).is_ok());
687
688 assert!(get_value_merger(42.into(), &MergeStrategy::Discard).is_ok());
689 assert!(get_value_merger(42.into(), &MergeStrategy::Retain).is_ok());
690 assert!(get_value_merger(42.into(), &MergeStrategy::Sum).is_ok());
691 assert!(get_value_merger(42.into(), &MergeStrategy::Min).is_ok());
692 assert!(get_value_merger(42.into(), &MergeStrategy::Max).is_ok());
693 assert!(get_value_merger(42.into(), &MergeStrategy::Array).is_ok());
694 assert!(get_value_merger(42.into(), &MergeStrategy::LongestArray).is_err());
695 assert!(get_value_merger(42.into(), &MergeStrategy::ShortestArray).is_err());
696 assert!(get_value_merger(42.into(), &MergeStrategy::Concat).is_err());
697 assert!(get_value_merger(42.into(), &MergeStrategy::ConcatNewline).is_err());
698 assert!(get_value_merger(42.into(), &MergeStrategy::ConcatRaw).is_err());
699 assert!(get_value_merger(42.into(), &MergeStrategy::FlatUnique).is_ok());
700
701 assert!(get_value_merger(42.into(), &MergeStrategy::Discard).is_ok());
702 assert!(get_value_merger(42.into(), &MergeStrategy::Retain).is_ok());
703 assert!(get_value_merger(4.2.into(), &MergeStrategy::Sum).is_ok());
704 assert!(get_value_merger(4.2.into(), &MergeStrategy::Min).is_ok());
705 assert!(get_value_merger(4.2.into(), &MergeStrategy::Max).is_ok());
706 assert!(get_value_merger(4.2.into(), &MergeStrategy::Array).is_ok());
707 assert!(get_value_merger(4.2.into(), &MergeStrategy::LongestArray).is_err());
708 assert!(get_value_merger(4.2.into(), &MergeStrategy::ShortestArray).is_err());
709 assert!(get_value_merger(4.2.into(), &MergeStrategy::Concat).is_err());
710 assert!(get_value_merger(4.2.into(), &MergeStrategy::ConcatNewline).is_err());
711 assert!(get_value_merger(4.2.into(), &MergeStrategy::ConcatRaw).is_err());
712 assert!(get_value_merger(4.2.into(), &MergeStrategy::FlatUnique).is_ok());
713
714 assert!(get_value_merger(true.into(), &MergeStrategy::Discard).is_ok());
715 assert!(get_value_merger(true.into(), &MergeStrategy::Retain).is_ok());
716 assert!(get_value_merger(true.into(), &MergeStrategy::Sum).is_err());
717 assert!(get_value_merger(true.into(), &MergeStrategy::Max).is_err());
718 assert!(get_value_merger(true.into(), &MergeStrategy::Min).is_err());
719 assert!(get_value_merger(true.into(), &MergeStrategy::Array).is_ok());
720 assert!(get_value_merger(true.into(), &MergeStrategy::LongestArray).is_err());
721 assert!(get_value_merger(true.into(), &MergeStrategy::ShortestArray).is_err());
722 assert!(get_value_merger(true.into(), &MergeStrategy::Concat).is_err());
723 assert!(get_value_merger(true.into(), &MergeStrategy::ConcatNewline).is_err());
724 assert!(get_value_merger(true.into(), &MergeStrategy::ConcatRaw).is_err());
725 assert!(get_value_merger(true.into(), &MergeStrategy::FlatUnique).is_ok());
726
727 assert!(get_value_merger(Utc::now().into(), &MergeStrategy::Discard).is_ok());
728 assert!(get_value_merger(Utc::now().into(), &MergeStrategy::Retain).is_ok());
729 assert!(get_value_merger(Utc::now().into(), &MergeStrategy::Sum).is_err());
730 assert!(get_value_merger(Utc::now().into(), &MergeStrategy::Max).is_err());
731 assert!(get_value_merger(Utc::now().into(), &MergeStrategy::Min).is_err());
732 assert!(get_value_merger(Utc::now().into(), &MergeStrategy::Array).is_ok());
733 assert!(get_value_merger(Utc::now().into(), &MergeStrategy::LongestArray).is_err());
734 assert!(get_value_merger(Utc::now().into(), &MergeStrategy::ShortestArray).is_err());
735 assert!(get_value_merger(Utc::now().into(), &MergeStrategy::Concat).is_err());
736 assert!(get_value_merger(Utc::now().into(), &MergeStrategy::ConcatNewline).is_err());
737 assert!(get_value_merger(Utc::now().into(), &MergeStrategy::ConcatRaw).is_err());
738 assert!(get_value_merger(Utc::now().into(), &MergeStrategy::Discard).is_ok());
739 assert!(get_value_merger(Utc::now().into(), &MergeStrategy::FlatUnique).is_ok());
740
741 assert!(get_value_merger(json!([]).into(), &MergeStrategy::Discard).is_ok());
742 assert!(get_value_merger(json!([]).into(), &MergeStrategy::Retain).is_ok());
743 assert!(get_value_merger(json!([]).into(), &MergeStrategy::Sum).is_err());
744 assert!(get_value_merger(json!([]).into(), &MergeStrategy::Max).is_err());
745 assert!(get_value_merger(json!([]).into(), &MergeStrategy::Min).is_err());
746 assert!(get_value_merger(json!([]).into(), &MergeStrategy::Array).is_ok());
747 assert!(get_value_merger(json!([]).into(), &MergeStrategy::LongestArray).is_ok());
748 assert!(get_value_merger(json!([]).into(), &MergeStrategy::ShortestArray).is_ok());
749 assert!(get_value_merger(json!([]).into(), &MergeStrategy::Concat).is_ok());
750 assert!(get_value_merger(json!([]).into(), &MergeStrategy::ConcatNewline).is_err());
751 assert!(get_value_merger(json!([]).into(), &MergeStrategy::ConcatRaw).is_err());
752 assert!(get_value_merger(json!([]).into(), &MergeStrategy::FlatUnique).is_ok());
753
754 assert!(get_value_merger(json!({}).into(), &MergeStrategy::Discard).is_ok());
755 assert!(get_value_merger(json!({}).into(), &MergeStrategy::Retain).is_ok());
756 assert!(get_value_merger(json!({}).into(), &MergeStrategy::Sum).is_err());
757 assert!(get_value_merger(json!({}).into(), &MergeStrategy::Max).is_err());
758 assert!(get_value_merger(json!({}).into(), &MergeStrategy::Min).is_err());
759 assert!(get_value_merger(json!({}).into(), &MergeStrategy::Array).is_ok());
760 assert!(get_value_merger(json!({}).into(), &MergeStrategy::LongestArray).is_err());
761 assert!(get_value_merger(json!({}).into(), &MergeStrategy::ShortestArray).is_err());
762 assert!(get_value_merger(json!({}).into(), &MergeStrategy::Concat).is_err());
763 assert!(get_value_merger(json!({}).into(), &MergeStrategy::ConcatNewline).is_err());
764 assert!(get_value_merger(json!({}).into(), &MergeStrategy::ConcatRaw).is_err());
765 assert!(get_value_merger(json!({}).into(), &MergeStrategy::FlatUnique).is_ok());
766
767 assert!(get_value_merger(json!(null).into(), &MergeStrategy::Discard).is_ok());
768 assert!(get_value_merger(json!(null).into(), &MergeStrategy::Retain).is_ok());
769 assert!(get_value_merger(json!(null).into(), &MergeStrategy::Sum).is_err());
770 assert!(get_value_merger(json!(null).into(), &MergeStrategy::Max).is_err());
771 assert!(get_value_merger(json!(null).into(), &MergeStrategy::Min).is_err());
772 assert!(get_value_merger(json!(null).into(), &MergeStrategy::Array).is_ok());
773 assert!(get_value_merger(json!(null).into(), &MergeStrategy::LongestArray).is_err());
774 assert!(get_value_merger(json!(null).into(), &MergeStrategy::ShortestArray).is_err());
775 assert!(get_value_merger(json!(null).into(), &MergeStrategy::Concat).is_err());
776 assert!(get_value_merger(json!(null).into(), &MergeStrategy::ConcatNewline).is_err());
777 assert!(get_value_merger(json!(null).into(), &MergeStrategy::ConcatRaw).is_err());
778 assert!(get_value_merger(json!(null).into(), &MergeStrategy::FlatUnique).is_ok());
779 }
780
781 #[test]
782 fn merging_values() {
783 assert_eq!(
784 merge("foo".into(), "bar".into(), &MergeStrategy::Discard),
785 Ok("foo".into())
786 );
787 assert_eq!(
788 merge("foo".into(), "bar".into(), &MergeStrategy::Retain),
789 Ok("bar".into())
790 );
791 assert_eq!(
792 merge("foo".into(), "bar".into(), &MergeStrategy::Array),
793 Ok(json!(["foo", "bar"]).into())
794 );
795 assert_eq!(
796 merge("foo".into(), "bar".into(), &MergeStrategy::Concat),
797 Ok("foo bar".into())
798 );
799 assert_eq!(
800 merge("foo".into(), "bar".into(), &MergeStrategy::ConcatNewline),
801 Ok("foo\nbar".into())
802 );
803 assert_eq!(
804 merge("foo".into(), "bar".into(), &MergeStrategy::ConcatRaw),
805 Ok("foobar".into())
806 );
807 assert!(merge("foo".into(), 42.into(), &MergeStrategy::Concat).is_err());
808 assert!(merge("foo".into(), 4.2.into(), &MergeStrategy::Concat).is_err());
809 assert!(merge("foo".into(), true.into(), &MergeStrategy::Concat).is_err());
810 assert!(merge("foo".into(), Utc::now().into(), &MergeStrategy::Concat).is_err());
811 assert!(merge("foo".into(), json!({}).into(), &MergeStrategy::Concat).is_err());
812 assert!(merge("foo".into(), json!([]).into(), &MergeStrategy::Concat).is_err());
813 assert!(merge("foo".into(), json!(null).into(), &MergeStrategy::Concat).is_err());
814
815 assert_eq!(
816 merge("foo".into(), "bar".into(), &MergeStrategy::ConcatNewline),
817 Ok("foo\nbar".into())
818 );
819
820 assert_eq!(
821 merge(21.into(), 21.into(), &MergeStrategy::Sum),
822 Ok(42.into())
823 );
824 assert_eq!(
825 merge(21.into(), 21.0.into(), &MergeStrategy::Sum),
826 Ok(42.0.into())
827 );
828 assert_eq!(
829 merge(21.0.into(), 21.into(), &MergeStrategy::Sum),
830 Ok(42.0.into())
831 );
832 assert_eq!(
833 merge(41.into(), 42.into(), &MergeStrategy::Max),
834 Ok(42.into())
835 );
836 assert_eq!(
837 merge(42.into(), 41.into(), &MergeStrategy::Max),
838 Ok(42.into())
839 );
840 assert_eq!(
841 merge(42.into(), 43.into(), &MergeStrategy::Min),
842 Ok(42.into())
843 );
844 assert_eq!(
845 merge(43.into(), 42.into(), &MergeStrategy::Min),
846 Ok(42.into())
847 );
848
849 assert_eq!(
850 merge(2.1.into(), 2.1.into(), &MergeStrategy::Sum),
851 Ok(4.2.into())
852 );
853 assert_eq!(
854 merge(4.1.into(), 4.2.into(), &MergeStrategy::Max),
855 Ok(4.2.into())
856 );
857 assert_eq!(
858 merge(4.2.into(), 4.1.into(), &MergeStrategy::Max),
859 Ok(4.2.into())
860 );
861 assert_eq!(
862 merge(4.2.into(), 4.3.into(), &MergeStrategy::Min),
863 Ok(4.2.into())
864 );
865 assert_eq!(
866 merge(4.3.into(), 4.2.into(), &MergeStrategy::Min),
867 Ok(4.2.into())
868 );
869
870 assert_eq!(
871 merge(
872 json!([4_i64]).into(),
873 json!([2_i64]).into(),
874 &MergeStrategy::Concat
875 ),
876 Ok(json!([4_i64, 2_i64]).into())
877 );
878 assert_eq!(
879 merge(json!([]).into(), 42_i64.into(), &MergeStrategy::Concat),
880 Ok(json!([42_i64]).into())
881 );
882
883 assert_eq!(
884 merge(
885 json!([34_i64]).into(),
886 json!([42_i64, 43_i64]).into(),
887 &MergeStrategy::ShortestArray
888 ),
889 Ok(json!([34_i64]).into())
890 );
891 assert_eq!(
892 merge(
893 json!([34_i64]).into(),
894 json!([42_i64, 43_i64]).into(),
895 &MergeStrategy::LongestArray
896 ),
897 Ok(json!([42_i64, 43_i64]).into())
898 );
899
900 let v = merge(34_i64.into(), 43_i64.into(), &MergeStrategy::FlatUnique).unwrap();
901 match v.clone() {
902 Value::Array(v) => {
903 let v: Vec<_> = v
904 .into_iter()
905 .map(|i| {
906 if let Value::Integer(i) = i {
907 i
908 } else {
909 panic!("Bad value");
910 }
911 })
912 .collect();
913 assert_eq!(v.iter().filter(|i| **i == 34i64).count(), 1);
914 assert_eq!(v.iter().filter(|i| **i == 43i64).count(), 1);
915 }
916 _ => {
917 panic!("Not array");
918 }
919 }
920 let v = merge(v, 34_i32.into(), &MergeStrategy::FlatUnique).unwrap();
921 if let Value::Array(v) = v {
922 let v: Vec<_> = v
923 .into_iter()
924 .map(|i| {
925 if let Value::Integer(i) = i {
926 i
927 } else {
928 panic!("Bad value");
929 }
930 })
931 .collect();
932 assert_eq!(v.iter().filter(|i| **i == 34i64).count(), 1);
933 assert_eq!(v.iter().filter(|i| **i == 43i64).count(), 1);
934 } else {
935 panic!("Not array");
936 }
937 }
938
939 #[test]
940 fn sum_returns_nan_errors() {
941 assert_eq!(
942 merge(
943 f64::INFINITY.into(),
944 f64::NEG_INFINITY.into(),
945 &MergeStrategy::Sum,
946 ),
947 Err("operation would produce NaN".to_owned())
948 );
949 }
950
951 fn merge(initial: Value, additional: Value, strategy: &MergeStrategy) -> Result<Value, String> {
952 let mut merger = get_value_merger(initial, strategy)?;
953 merger.add(additional)?;
954 let mut output = LogEvent::default();
955 let out_path = owned_event_path!("out");
956 merger.insert_into(&out_path, &mut output)?;
957 Ok(output.remove(&out_path).unwrap())
958 }
959
960 #[test]
965 fn timestamp_window_merger_handles_quoted_path() {
966 let t1 = Utc::now();
967 let t2 = t1 + chrono::Duration::seconds(5);
968
969 let mut merger: Box<dyn ReduceValueMerger> = Value::Timestamp(t1).into();
971 merger.add(Value::Timestamp(t2)).unwrap();
972
973 let mut out = LogEvent::default();
974 let path = vrl::path::parse_target_path("\"foo.bar\"").unwrap();
975 merger.insert_into(&path, &mut out).expect("insert_into");
976
977 assert_eq!(out.get(&path), Some(&Value::Timestamp(t1)));
978 let end_path = vrl::path::parse_target_path("\"foo.bar_end\"").unwrap();
979 assert_eq!(out.get(&end_path), Some(&Value::Timestamp(t2)));
980 }
981}