Skip to main content

vector/transforms/reduce/
merge_strategy.rs

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