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(¤t) {
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 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 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 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 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 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}