Skip to main content

vector/transforms/reduce/
transform.rs

1use std::{
2    collections::{HashMap, hash_map::Entry},
3    pin::Pin,
4    time::{Duration, Instant},
5};
6
7use futures::Stream;
8use indexmap::IndexMap;
9use vector_lib::stream::expiration_map::{Emitter, map_with_expiration};
10use vector_vrl_metrics::MetricsStorage;
11use vrl::{
12    path::{OwnedTargetPath, parse_target_path},
13    prelude::KeyString,
14};
15
16use crate::{
17    conditions::Condition,
18    event::{Event, EventMetadata, LogEvent, discriminant::Discriminant},
19    internal_events::{ReduceAddEventError, ReduceStaleEventFlushed},
20    transforms::{
21        TaskTransform,
22        reduce::{
23            config::ReduceConfig,
24            merge_strategy::{MergeStrategy, ReduceValueMerger, get_value_merger},
25        },
26    },
27};
28
29#[derive(Clone, Debug)]
30struct ReduceState {
31    events: usize,
32    fields: HashMap<OwnedTargetPath, Box<dyn ReduceValueMerger>>,
33    stale_since: Instant,
34    creation: Instant,
35    metadata: EventMetadata,
36}
37
38fn is_covered_by_strategy(
39    path: &OwnedTargetPath,
40    strategies: &IndexMap<OwnedTargetPath, MergeStrategy>,
41) -> bool {
42    let mut current = OwnedTargetPath::event_root();
43    for component in &path.path.segments {
44        current = current.with_field_appended(&component.to_string());
45        if strategies.contains_key(&current) {
46            return true;
47        }
48    }
49    false
50}
51
52impl ReduceState {
53    fn new() -> Self {
54        Self {
55            events: 0,
56            stale_since: Instant::now(),
57            creation: Instant::now(),
58            fields: HashMap::new(),
59            metadata: EventMetadata::default(),
60        }
61    }
62
63    fn add_event(&mut self, e: LogEvent, strategies: &IndexMap<OwnedTargetPath, MergeStrategy>) {
64        self.metadata.merge(e.metadata().clone());
65
66        for (path, strategy) in strategies {
67            if let Some(value) = e.get(path) {
68                match self.fields.entry(path.clone()) {
69                    Entry::Vacant(entry) => match get_value_merger(value.clone(), strategy) {
70                        Ok(m) => {
71                            entry.insert(m);
72                        }
73                        Err(error) => {
74                            warn!(message = "Failed to create value merger.", %error, %path);
75                        }
76                    },
77                    Entry::Occupied(mut entry) => {
78                        if let Err(error) = entry.get_mut().add(value.clone()) {
79                            warn!(message = "Failed to merge value.", %error);
80                        }
81                    }
82                }
83            }
84        }
85
86        if let Some(fields_iter) = e.all_event_fields_skip_array_elements() {
87            for (path, value) in fields_iter {
88                // This should not return an error, unless there is a bug in the event fields iterator.
89                let parsed_path = match parse_target_path(&path) {
90                    Ok(path) => path,
91                    Err(error) => {
92                        emit!(ReduceAddEventError { error, path });
93                        continue;
94                    }
95                };
96                if is_covered_by_strategy(&parsed_path, strategies) {
97                    continue;
98                }
99
100                let maybe_strategy = strategies.get(&parsed_path);
101                match self.fields.entry(parsed_path) {
102                    Entry::Vacant(entry) => {
103                        if let Some(strategy) = maybe_strategy {
104                            match get_value_merger(value.clone(), strategy) {
105                                Ok(m) => {
106                                    entry.insert(m);
107                                }
108                                Err(error) => {
109                                    warn!(message = "Failed to merge value.", %error);
110                                }
111                            }
112                        } else {
113                            entry.insert(value.clone().into());
114                        }
115                    }
116                    Entry::Occupied(mut entry) => {
117                        if let Err(error) = entry.get_mut().add(value.clone()) {
118                            warn!(message = "Failed to merge value.", %error);
119                        }
120                    }
121                }
122            }
123        }
124        // else the event root is not an object (see https://github.com/vectordotdev/vector/issues/18219)
125
126        self.events += 1;
127        self.stale_since = Instant::now();
128    }
129
130    fn flush(mut self) -> LogEvent {
131        let mut event = LogEvent::new_with_metadata(self.metadata);
132        for (path, v) in self.fields.drain() {
133            if let Err(error) = v.insert_into(&path, &mut event) {
134                warn!(message = "Failed to merge values for field.", %error);
135            }
136        }
137        self.events = 0;
138        event
139    }
140}
141
142#[derive(Clone, Debug)]
143pub struct Reduce {
144    expire_after: Duration,
145    flush_period: Duration,
146    end_every_period: Option<Duration>,
147    group_by: Vec<String>,
148    merge_strategies: IndexMap<OwnedTargetPath, MergeStrategy>,
149    reduce_merge_states: HashMap<Discriminant, ReduceState>,
150    ends_when: Option<Condition>,
151    starts_when: Option<Condition>,
152    max_events: Option<usize>,
153}
154
155fn validate_merge_strategies(strategies: IndexMap<KeyString, MergeStrategy>) -> crate::Result<()> {
156    for (path, _) in &strategies {
157        let contains_index = parse_target_path(path)
158            .map_err(|_| format!("Could not parse path: `{path}`"))?
159            .path
160            .segments
161            .iter()
162            .any(|segment| segment.is_index());
163        if contains_index {
164            return Err(format!(
165                "Merge strategies with indexes are currently not supported. Path: `{path}`"
166            )
167            .into());
168        }
169    }
170
171    Ok(())
172}
173
174impl Reduce {
175    pub fn new(
176        config: &ReduceConfig,
177        enrichment_tables: &vector_lib::enrichment::TableRegistry,
178        metrics_storage: &MetricsStorage,
179    ) -> crate::Result<Self> {
180        if config.ends_when.is_some() && config.starts_when.is_some() {
181            return Err("only one of `ends_when` and `starts_when` can be provided".into());
182        }
183
184        let ends_when = config
185            .ends_when
186            .as_ref()
187            .map(|c| c.build(enrichment_tables, metrics_storage))
188            .transpose()?;
189        let starts_when = config
190            .starts_when
191            .as_ref()
192            .map(|c| c.build(enrichment_tables, metrics_storage))
193            .transpose()?;
194        let group_by = config.group_by.clone().into_iter().collect();
195        let max_events = config.max_events.map(|max| max.into());
196
197        validate_merge_strategies(config.merge_strategies.clone())?;
198
199        Ok(Reduce {
200            expire_after: config.expire_after_ms,
201            flush_period: config.flush_period_ms,
202            end_every_period: config.end_every_period_ms,
203            group_by,
204            merge_strategies: config
205                .merge_strategies
206                .iter()
207                .filter_map(|(path, strategy)| {
208                    // TODO Invalid paths are ignored to preserve backwards compatibility.
209                    //      Merge strategy paths should ideally be [`lookup_v2::ConfigTargetPath`]
210                    //      which means an invalid path would result in an configuration error.
211                    let parsed_path = parse_target_path(path).ok();
212                    if parsed_path.is_none() {
213                        warn!(message = "Ignoring strategy with invalid path.", %path);
214                    }
215                    parsed_path.map(|path| (path, strategy.clone()))
216                })
217                .collect(),
218            reduce_merge_states: HashMap::new(),
219            ends_when,
220            starts_when,
221            max_events,
222        })
223    }
224
225    fn flush_into(&mut self, emitter: &mut Emitter<Event>) {
226        let mut flush_discriminants = Vec::new();
227        let now = Instant::now();
228        for (k, t) in &self.reduce_merge_states {
229            if let Some(period) = self.end_every_period
230                && (now - t.creation) >= period
231            {
232                flush_discriminants.push(k.clone());
233            }
234
235            if (now - t.stale_since) >= self.expire_after {
236                flush_discriminants.push(k.clone());
237            }
238        }
239        for k in &flush_discriminants {
240            if let Some(t) = self.reduce_merge_states.remove(k) {
241                emit!(ReduceStaleEventFlushed);
242                emitter.emit(Event::from(t.flush()));
243            }
244        }
245    }
246
247    fn flush_all_into(&mut self, emitter: &mut Emitter<Event>) {
248        self.reduce_merge_states
249            .drain()
250            .for_each(|(_, s)| emitter.emit(Event::from(s.flush())));
251    }
252
253    fn push_or_new_reduce_state(&mut self, event: LogEvent, discriminant: Discriminant) {
254        match self.reduce_merge_states.entry(discriminant) {
255            Entry::Vacant(entry) => {
256                let mut state = ReduceState::new();
257                state.add_event(event, &self.merge_strategies);
258                entry.insert(state);
259            }
260            Entry::Occupied(mut entry) => {
261                entry.get_mut().add_event(event, &self.merge_strategies);
262            }
263        };
264    }
265
266    pub fn transform_one(&mut self, emitter: &mut Emitter<Event>, event: Event) {
267        let (starts_here, event) = match &self.starts_when {
268            Some(condition) => condition.check(event),
269            None => (false, event),
270        };
271
272        let (mut ends_here, event) = match &self.ends_when {
273            Some(condition) => condition.check(event),
274            None => (false, event),
275        };
276
277        let event = event.into_log();
278        let discriminant = Discriminant::from_log_event(&event, &self.group_by);
279
280        if let Some(max_events) = self.max_events {
281            if max_events == 1 {
282                ends_here = true;
283            } else if let Some(entry) = self.reduce_merge_states.get(&discriminant) {
284                // The current event will finish this set
285                if entry.events + 1 == max_events {
286                    ends_here = true;
287                }
288            }
289        }
290
291        if starts_here {
292            if let Some(state) = self.reduce_merge_states.remove(&discriminant) {
293                emitter.emit(state.flush().into());
294            }
295
296            self.push_or_new_reduce_state(event, discriminant)
297        } else if ends_here {
298            emitter.emit(match self.reduce_merge_states.remove(&discriminant) {
299                Some(mut state) => {
300                    state.add_event(event, &self.merge_strategies);
301                    state.flush().into()
302                }
303                None => {
304                    let mut state = ReduceState::new();
305                    state.add_event(event, &self.merge_strategies);
306                    state.flush().into()
307                }
308            });
309        } else {
310            self.push_or_new_reduce_state(event, discriminant)
311        }
312    }
313}
314
315impl TaskTransform<Event> for Reduce {
316    fn transform(
317        self: Box<Self>,
318        input_rx: Pin<Box<dyn Stream<Item = Event> + Send>>,
319    ) -> Pin<Box<dyn Stream<Item = Event> + Send>>
320    where
321        Self: 'static,
322    {
323        let transform_fn = move |me: &mut Box<Reduce>, event, emitter: &mut Emitter<Event>| {
324            me.transform_one(emitter, event);
325        };
326
327        construct_output_stream(self, input_rx, transform_fn)
328    }
329}
330
331pub fn construct_output_stream(
332    reduce: Box<Reduce>,
333    input_rx: Pin<Box<dyn Stream<Item = Event> + Send>>,
334    mut transform_fn: impl FnMut(&mut Box<Reduce>, Event, &mut Emitter<Event>) + Send + Sync + 'static,
335) -> Pin<Box<dyn Stream<Item = Event> + Send>>
336where
337    Reduce: 'static,
338{
339    let flush_period = reduce.flush_period;
340    Box::pin(map_with_expiration(
341        reduce,
342        input_rx,
343        flush_period,
344        move |me, event, emitter| {
345            transform_fn(me, event, emitter);
346        },
347        |me, emitter| {
348            me.flush_into(emitter);
349        },
350        |me, emitter| {
351            me.flush_all_into(emitter);
352        },
353    ))
354}
355
356#[cfg(test)]
357mod test {
358    use std::sync::Arc;
359
360    use indoc::indoc;
361    use serde_json::json;
362    use tokio::sync::mpsc;
363    use tokio_stream::wrappers::ReceiverStream;
364    use vector_lib::{enrichment::TableRegistry, lookup::owned_value_path};
365    use vrl::event_path;
366    use vrl::value::Kind;
367
368    use super::*;
369    use crate::{
370        config::{OutputId, TransformConfig, schema, schema::Definition},
371        event::{LogEvent, Value},
372        test_util::components::assert_transform_compliance,
373        transforms::test::create_topology,
374    };
375
376    #[tokio::test]
377    async fn reduce_from_condition() {
378        let reduce_config = serde_yaml::from_str::<ReduceConfig>(indoc! {"
379            group_by:
380              - request_id
381            ends_when:
382              type: vrl
383              source: exists(.test_end)
384        "})
385        .unwrap();
386
387        assert_transform_compliance(async move {
388            let input_definition = schema::Definition::default_legacy_namespace()
389                .with_event_field(&owned_value_path!("counter"), Kind::integer(), None)
390                .with_event_field(&owned_value_path!("request_id"), Kind::bytes(), None)
391                .with_event_field(
392                    &owned_value_path!("test_end"),
393                    Kind::bytes().or_undefined(),
394                    None,
395                )
396                .with_event_field(
397                    &owned_value_path!("extra_field"),
398                    Kind::bytes().or_undefined(),
399                    None,
400                );
401            let schema_definitions = reduce_config
402                .outputs(&Default::default(), &[("test".into(), input_definition)])
403                .first()
404                .unwrap()
405                .schema_definitions(true)
406                .clone();
407
408            let new_schema_definition = reduce_config.outputs(
409                &Default::default(),
410                &[(OutputId::from("in"), Definition::default_legacy_namespace())],
411            )[0]
412            .clone()
413            .log_schema_definitions
414            .get(&OutputId::from("in"))
415            .unwrap()
416            .clone();
417
418            let (tx, rx) = mpsc::channel(1);
419            let (topology, mut out) = create_topology(ReceiverStream::new(rx), reduce_config).await;
420
421            let mut e_1 = LogEvent::from("test message 1");
422            e_1.insert(event_path!("counter"), 1);
423            e_1.insert(event_path!("request_id"), "1");
424            let mut metadata_1 = e_1.metadata().clone();
425            metadata_1.set_upstream_id(Arc::new(OutputId::from("transform")));
426            metadata_1.set_schema_definition(&Arc::new(new_schema_definition.clone()));
427
428            let mut e_2 = LogEvent::from("test message 2");
429            e_2.insert(event_path!("counter"), 2);
430            e_2.insert(event_path!("request_id"), "2");
431            let mut metadata_2 = e_2.metadata().clone();
432            metadata_2.set_upstream_id(Arc::new(OutputId::from("transform")));
433            metadata_2.set_schema_definition(&Arc::new(new_schema_definition.clone()));
434
435            let mut e_3 = LogEvent::from("test message 3");
436            e_3.insert(event_path!("counter"), 3);
437            e_3.insert(event_path!("request_id"), "1");
438
439            let mut e_4 = LogEvent::from("test message 4");
440            e_4.insert(event_path!("counter"), 4);
441            e_4.insert(event_path!("request_id"), "1");
442            e_4.insert(event_path!("test_end"), "yep");
443
444            let mut e_5 = LogEvent::from("test message 5");
445            e_5.insert(event_path!("counter"), 5);
446            e_5.insert(event_path!("request_id"), "2");
447            e_5.insert(event_path!("extra_field"), "value1");
448            e_5.insert(event_path!("test_end"), "yep");
449
450            for event in [e_1.into(), e_2.into(), e_3.into(), e_4.into(), e_5.into()] {
451                tx.send(event).await.unwrap();
452            }
453
454            let output_1 = out.recv().await.unwrap().into_log();
455            assert_eq!(output_1["message"], "test message 1".into());
456            assert_eq!(output_1["counter"], Value::from(8));
457            assert_eq!(output_1.metadata(), &metadata_1);
458            schema_definitions
459                .values()
460                .for_each(|definition| definition.assert_valid_for_event(&output_1.clone().into()));
461
462            let output_2 = out.recv().await.unwrap().into_log();
463            assert_eq!(output_2["message"], "test message 2".into());
464            assert_eq!(output_2["extra_field"], "value1".into());
465            assert_eq!(output_2["counter"], Value::from(7));
466            assert_eq!(output_2.metadata(), &metadata_2);
467            schema_definitions
468                .values()
469                .for_each(|definition| definition.assert_valid_for_event(&output_2.clone().into()));
470
471            drop(tx);
472            topology.stop().await;
473            assert_eq!(out.recv().await, None);
474        })
475        .await;
476    }
477
478    #[tokio::test]
479    async fn reduce_merge_strategies() {
480        let reduce_config = serde_yaml::from_str::<ReduceConfig>(indoc! {"
481            group_by:
482              - request_id
483            merge_strategies:
484              foo: concat
485              bar: array
486              baz: max
487            ends_when:
488              type: vrl
489              source: exists(.test_end)
490        "})
491        .unwrap();
492
493        assert_transform_compliance(async move {
494            let (tx, rx) = mpsc::channel(1);
495
496            let new_schema_definition = reduce_config.outputs(
497                &Default::default(),
498                &[(OutputId::from("in"), Definition::default_legacy_namespace())],
499            )[0]
500            .clone()
501            .log_schema_definitions
502            .get(&OutputId::from("in"))
503            .unwrap()
504            .clone();
505
506            let (topology, mut out) = create_topology(ReceiverStream::new(rx), reduce_config).await;
507
508            let mut e_1 = LogEvent::from("test message 1");
509            e_1.insert(event_path!("foo"), "first foo");
510            e_1.insert(event_path!("bar"), "first bar");
511            e_1.insert(event_path!("baz"), 2);
512            e_1.insert(event_path!("request_id"), "1");
513            let mut metadata = e_1.metadata().clone();
514            metadata.set_upstream_id(Arc::new(OutputId::from("transform")));
515            metadata.set_schema_definition(&Arc::new(new_schema_definition.clone()));
516            tx.send(e_1.into()).await.unwrap();
517
518            let mut e_2 = LogEvent::from("test message 2");
519            e_2.insert(event_path!("foo"), "second foo");
520            e_2.insert(event_path!("bar"), 2);
521            e_2.insert(event_path!("baz"), "not number");
522            e_2.insert(event_path!("request_id"), "1");
523            tx.send(e_2.into()).await.unwrap();
524
525            let mut e_3 = LogEvent::from("test message 3");
526            e_3.insert(event_path!("foo"), 10);
527            e_3.insert(event_path!("bar"), "third bar");
528            e_3.insert(event_path!("baz"), 3);
529            e_3.insert(event_path!("request_id"), "1");
530            e_3.insert(event_path!("test_end"), "yep");
531            tx.send(e_3.into()).await.unwrap();
532
533            let output_1 = out.recv().await.unwrap().into_log();
534            assert_eq!(output_1["message"], "test message 1".into());
535            assert_eq!(output_1["foo"], "first foo second foo".into());
536            assert_eq!(
537                output_1["bar"],
538                Value::Array(vec!["first bar".into(), 2.into(), "third bar".into()]),
539            );
540            assert_eq!(output_1["baz"], 3.into());
541            assert_eq!(output_1.metadata(), &metadata);
542
543            drop(tx);
544            topology.stop().await;
545            assert_eq!(out.recv().await, None);
546        })
547        .await;
548    }
549
550    #[tokio::test]
551    async fn missing_group_by() {
552        let reduce_config = serde_yaml::from_str::<ReduceConfig>(indoc! {"
553            group_by:
554              - request_id
555            ends_when:
556              type: vrl
557              source: exists(.test_end)
558        "})
559        .unwrap();
560
561        assert_transform_compliance(async move {
562            let (tx, rx) = mpsc::channel(1);
563            let new_schema_definition = reduce_config.outputs(
564                &Default::default(),
565                &[(OutputId::from("in"), Definition::default_legacy_namespace())],
566            )[0]
567            .clone()
568            .log_schema_definitions
569            .get(&OutputId::from("in"))
570            .unwrap()
571            .clone();
572
573            let (topology, mut out) = create_topology(ReceiverStream::new(rx), reduce_config).await;
574
575            let mut e_1 = LogEvent::from("test message 1");
576            e_1.insert(event_path!("counter"), 1);
577            e_1.insert(event_path!("request_id"), "1");
578            let mut metadata_1 = e_1.metadata().clone();
579            metadata_1.set_upstream_id(Arc::new(OutputId::from("transform")));
580            metadata_1.set_schema_definition(&Arc::new(new_schema_definition.clone()));
581            tx.send(e_1.into()).await.unwrap();
582
583            let mut e_2 = LogEvent::from("test message 2");
584            e_2.insert(event_path!("counter"), 2);
585            let mut metadata_2 = e_2.metadata().clone();
586            metadata_2.set_upstream_id(Arc::new(OutputId::from("transform")));
587            metadata_2.set_schema_definition(&Arc::new(new_schema_definition));
588            tx.send(e_2.into()).await.unwrap();
589
590            let mut e_3 = LogEvent::from("test message 3");
591            e_3.insert(event_path!("counter"), 3);
592            e_3.insert(event_path!("request_id"), "1");
593            tx.send(e_3.into()).await.unwrap();
594
595            let mut e_4 = LogEvent::from("test message 4");
596            e_4.insert(event_path!("counter"), 4);
597            e_4.insert(event_path!("request_id"), "1");
598            e_4.insert(event_path!("test_end"), "yep");
599            tx.send(e_4.into()).await.unwrap();
600
601            let mut e_5 = LogEvent::from("test message 5");
602            e_5.insert(event_path!("counter"), 5);
603            e_5.insert(event_path!("extra_field"), "value1");
604            e_5.insert(event_path!("test_end"), "yep");
605            tx.send(e_5.into()).await.unwrap();
606
607            let output_1 = out.recv().await.unwrap().into_log();
608            assert_eq!(output_1["message"], "test message 1".into());
609            assert_eq!(output_1["counter"], Value::from(8));
610            assert_eq!(output_1.metadata(), &metadata_1);
611
612            let output_2 = out.recv().await.unwrap().into_log();
613            assert_eq!(output_2["message"], "test message 2".into());
614            assert_eq!(output_2["extra_field"], "value1".into());
615            assert_eq!(output_2["counter"], Value::from(7));
616            assert_eq!(output_2.metadata(), &metadata_2);
617
618            drop(tx);
619            topology.stop().await;
620            assert_eq!(out.recv().await, None);
621        })
622        .await;
623    }
624
625    #[tokio::test]
626    async fn max_events_0() {
627        let reduce_config = serde_yaml::from_str::<ReduceConfig>(indoc! {"
628            group_by:
629              - id
630            merge_strategies:
631              id: retain
632              message: array
633            max_events: 0
634        "});
635
636        match reduce_config {
637            Ok(_conf) => unreachable!("max_events=0 should be rejected."),
638            Err(err) => assert!(
639                err.to_string()
640                    .contains("invalid value: integer `0`, expected a nonzero usize")
641            ),
642        }
643    }
644
645    #[tokio::test]
646    async fn max_events_1() {
647        let reduce_config = serde_yaml::from_str::<ReduceConfig>(indoc! {"
648            group_by:
649              - id
650            merge_strategies:
651              id: retain
652              message: array
653            max_events: 1
654        "})
655        .unwrap();
656        assert_transform_compliance(async move {
657            let (tx, rx) = mpsc::channel(1);
658            let (topology, mut out) = create_topology(ReceiverStream::new(rx), reduce_config).await;
659
660            let mut e_1 = LogEvent::from("test 1");
661            e_1.insert(event_path!("id"), "1");
662
663            let mut e_2 = LogEvent::from("test 2");
664            e_2.insert(event_path!("id"), "1");
665
666            let mut e_3 = LogEvent::from("test 3");
667            e_3.insert(event_path!("id"), "1");
668
669            for event in [e_1.into(), e_2.into(), e_3.into()] {
670                tx.send(event).await.unwrap();
671            }
672
673            let output_1 = out.recv().await.unwrap().into_log();
674            assert_eq!(output_1["message"], vec!["test 1"].into());
675            let output_2 = out.recv().await.unwrap().into_log();
676            assert_eq!(output_2["message"], vec!["test 2"].into());
677
678            let output_3 = out.recv().await.unwrap().into_log();
679            assert_eq!(output_3["message"], vec!["test 3"].into());
680
681            drop(tx);
682            topology.stop().await;
683            assert_eq!(out.recv().await, None);
684        })
685        .await;
686    }
687
688    #[tokio::test]
689    async fn max_events() {
690        let reduce_config = serde_yaml::from_str::<ReduceConfig>(indoc! {"
691            group_by:
692              - id
693            merge_strategies:
694              id: retain
695              message: array
696            max_events: 3
697        "})
698        .unwrap();
699
700        assert_transform_compliance(async move {
701            let (tx, rx) = mpsc::channel(1);
702            let (topology, mut out) = create_topology(ReceiverStream::new(rx), reduce_config).await;
703
704            let mut e_1 = LogEvent::from("test 1");
705            e_1.insert(event_path!("id"), "1");
706
707            let mut e_2 = LogEvent::from("test 2");
708            e_2.insert(event_path!("id"), "1");
709
710            let mut e_3 = LogEvent::from("test 3");
711            e_3.insert(event_path!("id"), "1");
712
713            let mut e_4 = LogEvent::from("test 4");
714            e_4.insert(event_path!("id"), "1");
715
716            let mut e_5 = LogEvent::from("test 5");
717            e_5.insert(event_path!("id"), "1");
718
719            let mut e_6 = LogEvent::from("test 6");
720            e_6.insert(event_path!("id"), "1");
721
722            for event in [
723                e_1.into(),
724                e_2.into(),
725                e_3.into(),
726                e_4.into(),
727                e_5.into(),
728                e_6.into(),
729            ] {
730                tx.send(event).await.unwrap();
731            }
732
733            let output_1 = out.recv().await.unwrap().into_log();
734            assert_eq!(
735                output_1["message"],
736                vec!["test 1", "test 2", "test 3"].into()
737            );
738
739            let output_2 = out.recv().await.unwrap().into_log();
740            assert_eq!(
741                output_2["message"],
742                vec!["test 4", "test 5", "test 6"].into()
743            );
744
745            drop(tx);
746            topology.stop().await;
747            assert_eq!(out.recv().await, None);
748        })
749        .await
750    }
751
752    #[tokio::test]
753    async fn arrays() {
754        let reduce_config = serde_yaml::from_str::<ReduceConfig>(indoc! {"
755            group_by:
756              - request_id
757            merge_strategies:
758              foo: array
759              bar: concat
760            ends_when:
761              type: vrl
762              source: exists(.test_end)
763        "})
764        .unwrap();
765
766        assert_transform_compliance(async move {
767            let (tx, rx) = mpsc::channel(1);
768
769            let new_schema_definition = reduce_config.outputs(
770                &Default::default(),
771                &[(OutputId::from("in"), Definition::default_legacy_namespace())],
772            )[0]
773            .clone()
774            .log_schema_definitions
775            .get(&OutputId::from("in"))
776            .unwrap()
777            .clone();
778
779            let (topology, mut out) = create_topology(ReceiverStream::new(rx), reduce_config).await;
780
781            let mut e_1 = LogEvent::from("test message 1");
782            e_1.insert(event_path!("foo"), json!([1, 3]));
783            e_1.insert(event_path!("bar"), json!([1, 3]));
784            e_1.insert(event_path!("request_id"), "1");
785            let mut metadata_1 = e_1.metadata().clone();
786            metadata_1.set_upstream_id(Arc::new(OutputId::from("transform")));
787            metadata_1.set_schema_definition(&Arc::new(new_schema_definition.clone()));
788
789            tx.send(e_1.into()).await.unwrap();
790
791            let mut e_2 = LogEvent::from("test message 2");
792            e_2.insert(event_path!("foo"), json!([2, 4]));
793            e_2.insert(event_path!("bar"), json!([2, 4]));
794            e_2.insert(event_path!("request_id"), "2");
795            let mut metadata_2 = e_2.metadata().clone();
796            metadata_2.set_upstream_id(Arc::new(OutputId::from("transform")));
797            metadata_2.set_schema_definition(&Arc::new(new_schema_definition));
798            tx.send(e_2.into()).await.unwrap();
799
800            let mut e_3 = LogEvent::from("test message 3");
801            e_3.insert(event_path!("foo"), json!([5, 7]));
802            e_3.insert(event_path!("bar"), json!([5, 7]));
803            e_3.insert(event_path!("request_id"), "1");
804            tx.send(e_3.into()).await.unwrap();
805
806            let mut e_4 = LogEvent::from("test message 4");
807            e_4.insert(event_path!("foo"), json!("done"));
808            e_4.insert(event_path!("bar"), json!("done"));
809            e_4.insert(event_path!("request_id"), "1");
810            e_4.insert(event_path!("test_end"), "yep");
811            tx.send(e_4.into()).await.unwrap();
812
813            let mut e_5 = LogEvent::from("test message 5");
814            e_5.insert(event_path!("foo"), json!([6, 8]));
815            e_5.insert(event_path!("bar"), json!([6, 8]));
816            e_5.insert(event_path!("request_id"), "2");
817            tx.send(e_5.into()).await.unwrap();
818
819            let mut e_6 = LogEvent::from("test message 6");
820            e_6.insert(event_path!("foo"), json!("done"));
821            e_6.insert(event_path!("bar"), json!("done"));
822            e_6.insert(event_path!("request_id"), "2");
823            e_6.insert(event_path!("test_end"), "yep");
824            tx.send(e_6.into()).await.unwrap();
825
826            let output_1 = out.recv().await.unwrap().into_log();
827            assert_eq!(output_1["foo"], json!([[1, 3], [5, 7], "done"]).into());
828            assert_eq!(output_1["bar"], json!([1, 3, 5, 7, "done"]).into());
829            assert_eq!(output_1.metadata(), &metadata_1);
830
831            let output_2 = out.recv().await.unwrap().into_log();
832            assert_eq!(output_2["foo"], json!([[2, 4], [6, 8], "done"]).into());
833            assert_eq!(output_2["bar"], json!([2, 4, 6, 8, "done"]).into());
834            assert_eq!(output_2.metadata(), &metadata_2);
835
836            drop(tx);
837            topology.stop().await;
838            assert_eq!(out.recv().await, None);
839        })
840        .await;
841    }
842
843    #[tokio::test]
844    async fn strategy_path_with_nested_fields() {
845        let reduce_config = serde_yaml::from_str::<ReduceConfig>(indoc! {"
846            group_by:
847              - id
848            merge_strategies:
849              id: discard
850              message.a.b: array
851            ends_when:
852              type: vrl
853              source: exists(.test_end)
854        "})
855        .unwrap();
856
857        assert_transform_compliance(async move {
858            let (tx, rx) = mpsc::channel(1);
859
860            let (topology, mut out) = create_topology(ReceiverStream::new(rx), reduce_config).await;
861
862            let e_1 = LogEvent::from(Value::from(btreemap! {
863                "id" => 777,
864                "message" => btreemap! {
865                    "a" => btreemap! {
866                        "b" => vec![1,2],
867                        "num" => 1,
868                    },
869                },
870                "arr" => vec![btreemap! { "a" => 1 }, btreemap! { "b" => 1 }]
871            }));
872            let mut metadata_1 = e_1.metadata().clone();
873            metadata_1.set_upstream_id(Arc::new(OutputId::from("reduce")));
874
875            tx.send(e_1.into()).await.unwrap();
876
877            let e_2 = LogEvent::from(Value::from(btreemap! {
878                "id" => 777,
879                "message" => btreemap! {
880                        "a" => btreemap! {
881                            "b" => vec![3,4],
882                            "num" => 2,
883                        },
884                },
885                 "arr" => vec![btreemap! { "a" => 2 }, btreemap! { "b" => 2 }],
886                "test_end" => "done",
887            }));
888            tx.send(e_2.into()).await.unwrap();
889
890            let mut output = out.recv().await.unwrap().into_log();
891
892            // Remove timestamp fields which were automatically added.
893            output.remove_timestamp();
894            output.remove(event_path!("timestamp_end"));
895
896            assert_eq!(
897                *output.value(),
898                btreemap! {
899                    "id" => 777,
900                    "message" => btreemap! {
901                        "a" => btreemap! {
902                            "b" => vec![vec![1, 2], vec![3,4]],
903                            "num" => 3,
904                        },
905                    },
906                    "arr" => vec![btreemap! { "a" => 1 }, btreemap! { "b" => 1 }],
907                    "test_end" => "done",
908                }
909                .into()
910            );
911
912            drop(tx);
913            topology.stop().await;
914            assert_eq!(out.recv().await, None);
915        })
916        .await;
917    }
918
919    #[test]
920    fn invalid_merge_strategies_containing_indexes() {
921        let config = serde_yaml::from_str::<ReduceConfig>(indoc! {"
922            group_by:
923              - id
924            merge_strategies:
925              id: discard
926              'nested.msg[0]': array
927        "})
928        .unwrap();
929        let error = Reduce::new(
930            &config,
931            &TableRegistry::default(),
932            &MetricsStorage::default(),
933        )
934        .unwrap_err();
935        assert_eq!(
936            error.to_string(),
937            "Merge strategies with indexes are currently not supported. Path: `nested.msg[0]`"
938        );
939    }
940
941    #[tokio::test]
942    async fn merge_objects_in_array() {
943        let config = serde_yaml::from_str::<ReduceConfig>(indoc! {r#"
944            group_by:
945              - id
946            merge_strategies:
947              events: array
948              '"a-b"': retain
949              another: discard
950            ends_when:
951              type: vrl
952              source: exists(.test_end)
953        "#})
954        .unwrap();
955
956        assert_transform_compliance(async move {
957            let (tx, rx) = mpsc::channel(1);
958
959            let (topology, mut out) = create_topology(ReceiverStream::new(rx), config).await;
960
961            let v_1 = Value::from(btreemap! {
962                "attrs" => btreemap! {
963                    "nested.msg" => "foo",
964                },
965                "sev" => 2,
966            });
967            let mut e_1 = LogEvent::from(Value::from(
968                btreemap! {"id" => 777, "another" => btreemap!{ "a" => 1}},
969            ));
970            e_1.insert(event_path!("events"), v_1.clone());
971            e_1.insert(&vrl::path::parse_target_path("\"a-b\"").unwrap(), 2);
972            tx.send(e_1.into()).await.unwrap();
973
974            let v_2 = Value::from(btreemap! {
975                "attrs" => btreemap! {
976                    "nested.msg" => "bar",
977                },
978                "sev" => 3,
979            });
980            let mut e_2 = LogEvent::from(Value::from(
981                btreemap! {"id" => 777, "test_end" => "done", "another" => btreemap!{ "b" => 2}},
982            ));
983            e_2.insert(event_path!("events"), v_2.clone());
984            e_2.insert(&vrl::path::parse_target_path("\"a-b\"").unwrap(), 2);
985            tx.send(e_2.into()).await.unwrap();
986
987            let output = out.recv().await.unwrap().into_log();
988            let expected_value = Value::from(btreemap! {
989                "id" => 1554,
990                "events" => vec![v_1, v_2],
991                "another" => btreemap!{ "a" => 1},
992                "a-b" => 2,
993                "test_end" => "done"
994            });
995            assert_eq!(*output.value(), expected_value);
996
997            drop(tx);
998            topology.stop().await;
999            assert_eq!(out.recv().await, None);
1000        })
1001        .await
1002    }
1003
1004    #[tokio::test]
1005    async fn merged_quoted_path() {
1006        let config = serde_yaml::from_str::<ReduceConfig>(indoc! {"
1007            ends_when:
1008              type: vrl
1009              source: exists(.test_end)
1010        "})
1011        .unwrap();
1012
1013        assert_transform_compliance(async move {
1014            let (tx, rx) = mpsc::channel(1);
1015
1016            let (topology, mut out) = create_topology(ReceiverStream::new(rx), config).await;
1017
1018            let e_1 = LogEvent::from(Value::from(btreemap! {"a b" => 1}));
1019            tx.send(e_1.into()).await.unwrap();
1020
1021            let e_2 = LogEvent::from(Value::from(btreemap! {"a b" => 2, "test_end" => "done"}));
1022            tx.send(e_2.into()).await.unwrap();
1023
1024            let output = out.recv().await.unwrap().into_log();
1025            let expected_value = Value::from(btreemap! {
1026                "a b" => 3,
1027                "test_end" => "done"
1028            });
1029            assert_eq!(*output.value(), expected_value);
1030
1031            drop(tx);
1032            topology.stop().await;
1033            assert_eq!(out.recv().await, None);
1034        })
1035        .await
1036    }
1037}