Skip to main content

vector/transforms/
delay.rs

1use std::{num::NonZeroUsize, pin::Pin, time::Duration};
2
3use async_stream::stream;
4use futures::{Stream, StreamExt};
5use serde_with::serde_as;
6use snafu::Snafu;
7use tokio_util::time::DelayQueue;
8use vector_lib::configurable::configurable_component;
9use vector_lib::internal_event::INTENTIONAL;
10use vector_lib::{config::clone_input_definitions, internal_event::ComponentEventsDropped};
11
12use crate::{
13    conditions::{AnyCondition, Condition},
14    config::{DataType, Input, OutputId, TransformConfig, TransformContext, TransformOutput},
15    event::Event,
16    schema,
17    transforms::{TaskTransform, Transform},
18};
19
20/// Configuration for the `delay` transform.
21#[serde_as]
22#[configurable_component(transform("delay", "Slow down events passing through a topology."))]
23#[derive(Clone, Debug)]
24#[serde(deny_unknown_fields)]
25pub struct DelayConfig {
26    /// Time to delay each event, in milliseconds.
27    #[serde_as(as = "serde_with::DurationMilliSeconds<u64>")]
28    #[configurable(metadata(docs::human_name = "Delay in milliseconds", docs::example = 200))]
29    delay_ms: Duration,
30
31    /// Limit for number of items in the delay queue.
32    #[serde(default = "default_queue_capacity")]
33    queue_capacity: NonZeroUsize,
34
35    /// Strategy to handle full queue capacity.
36    #[serde(default)]
37    overflow_strategy: OverflowStrategy,
38
39    /// Delay events in provided delay periods until the condition is met.
40    condition: Option<AnyCondition>,
41}
42
43const fn default_queue_capacity() -> NonZeroUsize {
44    NonZeroUsize::new(500).expect("static non-zero number")
45}
46
47impl Default for DelayConfig {
48    fn default() -> Self {
49        Self {
50            delay_ms: Default::default(),
51            queue_capacity: default_queue_capacity(),
52            overflow_strategy: Default::default(),
53            condition: Default::default(),
54        }
55    }
56}
57
58/// Event handling behavior when delay queue is full.
59#[configurable_component]
60#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
61#[serde(rename_all = "snake_case")]
62pub enum OverflowStrategy {
63    /// Wait for free space in the queue.
64    ///
65    /// This applies backpressure up the topology, signalling that sources should slow down
66    /// the acceptance/consumption of events. This may cause the system to degenerate if this
67    /// component blocks for too long.
68    #[default]
69    Block,
70
71    /// Drops the event instead of waiting for free space in the queue.
72    ///
73    /// The event will be intentionally dropped. This mode is typically used when performance is the
74    /// highest priority, and it is preferable to temporarily lose events rather than cause a
75    /// slowdown in the acceptance/consumption of events.
76    DropNewest,
77
78    /// Forward the event without any delay to next component.
79    Forward,
80}
81
82impl_generate_config_from_default!(DelayConfig);
83
84#[async_trait::async_trait]
85#[typetag::serde(name = "delay")]
86impl TransformConfig for DelayConfig {
87    async fn build(&self, context: &TransformContext) -> crate::Result<Transform> {
88        if self.delay_ms.as_millis() == 0 {
89            return Err(Box::new(BuildError::ZeroDelayDuration));
90        }
91        Ok(Transform::event_task(Delay::new(self, context)?))
92    }
93
94    fn input(&self) -> Input {
95        Input::all()
96    }
97
98    fn outputs(
99        &self,
100        _: &TransformContext,
101        input_definitions: &[(OutputId, schema::Definition)],
102    ) -> Vec<TransformOutput> {
103        // The event is not modified, so the definition is passed through as-is
104        vec![TransformOutput::new(
105            DataType::all_bits(),
106            clone_input_definitions(input_definitions),
107        )]
108    }
109
110    fn validate(&self, _: &TransformContext) -> Result<(), Vec<String>> {
111        if self.delay_ms.as_millis() == 0 {
112            Err(vec!["delay must not be zero".to_string()])
113        } else {
114            Ok(())
115        }
116    }
117
118    fn validate_env(&self, context: &TransformContext) -> Result<(), Vec<String>> {
119        self.condition
120            .as_ref()
121            .map(|c| {
122                c.validate(&context.enrichment_tables, &context.metrics_storage)
123                    .map_err(|e| vec![format!("condition: {e}")])
124            })
125            .unwrap_or(Ok(()))
126    }
127}
128
129pub struct Delay {
130    delay: Duration,
131    queue: DelayQueue<Event>,
132    queue_capacity: NonZeroUsize,
133    overflow_strategy: OverflowStrategy,
134    condition: Option<Condition>,
135}
136
137impl Delay {
138    pub fn new(config: &DelayConfig, context: &TransformContext) -> crate::Result<Self> {
139        Ok(Self {
140            delay: config.delay_ms,
141            queue: DelayQueue::with_capacity(config.queue_capacity.get()),
142            queue_capacity: config.queue_capacity,
143            overflow_strategy: config.overflow_strategy,
144            condition: config
145                .condition
146                .as_ref()
147                .map(|c| c.build(&context.enrichment_tables, &context.metrics_storage))
148                .transpose()?,
149        })
150    }
151
152    fn check_condition(&self, event: Event, first: bool) -> (bool, Event) {
153        if let Some(condition) = self.condition.as_ref() {
154            condition.check(event)
155        } else {
156            // If this is the first check, we need to ensure at least one delay is
157            // done if no condition is configured
158            (!first, event)
159        }
160    }
161}
162
163impl TaskTransform<Event> for Delay {
164    fn transform(
165        mut self: Box<Self>,
166        mut input_rx: Pin<Box<dyn Stream<Item = Event> + Send>>,
167    ) -> Pin<Box<dyn Stream<Item = Event> + Send>>
168    where
169        Self: 'static,
170    {
171        Box::pin(stream! {
172            let mut done = false;
173            loop {
174                if done && self.queue.is_empty() {
175                    break;
176                }
177                tokio::select! {
178                    biased;
179
180                    Some(res) = self.queue.next() => {
181                        let event = res.into_inner();
182                        let (result, event) = self.check_condition(event, false);
183                        if result {
184                            yield event;
185                        } else {
186                            self.queue.insert(event, self.delay);
187                        }
188                        if done && self.queue.is_empty() {
189                            break;
190                        }
191                    },
192
193                    maybe_event = input_rx.next(), if !done => {
194                        match maybe_event {
195                            None => {
196                                done = true;
197                            }
198                            Some(event) => {
199                                let (result, event) = self.check_condition(event, true);
200                                if result {
201                                    yield event
202                                } else {
203                                    if self.queue_capacity.get() <= self.queue.len() {
204                                        match self.overflow_strategy {
205                                            OverflowStrategy::Block => {
206                                                while self.queue_capacity.get() <= self.queue.len() && let Some(next) = self.queue.next().await {
207                                                    let event = next.into_inner();
208                                                    let (result, event) = self.check_condition(event, false);
209                                                    if result {
210                                                        yield event;
211                                                    } else {
212                                                        self.queue.insert(event, self.delay);
213                                                    }
214                                                }
215                                            },
216                                            OverflowStrategy::DropNewest => {
217                                                emit!(ComponentEventsDropped::<INTENTIONAL> {
218                                                    count: 1,
219                                                    reason: "Queue is full and overflow strategy is drop_newest",
220                                                });
221                                                continue;
222                                            }
223                                            OverflowStrategy::Forward => {
224                                                yield event;
225                                                continue;
226                                            }
227                                        }
228                                    }
229                                    self.queue.insert(event, self.delay);
230                                }
231                            }
232                        }
233                    },
234                }
235            }
236        })
237    }
238}
239
240#[derive(Debug, Snafu)]
241pub enum BuildError {
242    #[snafu(display("The delay duration must not be zero"))]
243    ZeroDelayDuration,
244}
245
246#[cfg(test)]
247mod tests {
248    use indoc::indoc;
249    use std::task::Poll;
250
251    use futures::SinkExt;
252    use vector_lib::event::TraceEvent;
253
254    use super::*;
255    use crate::event::LogEvent;
256
257    #[test]
258    fn generate_config() {
259        crate::test_util::test_generate_config::<DelayConfig>();
260    }
261
262    #[tokio::test]
263    async fn delay_events() {
264        let config = serde_yaml::from_str::<DelayConfig>(indoc! {"
265            delay_ms: 200
266        "})
267        .unwrap();
268
269        let delay =
270            Transform::event_task(Delay::new(&config, &TransformContext::default()).unwrap());
271
272        let delay = delay.into_task();
273
274        let (mut tx, rx) = futures::channel::mpsc::channel(10);
275        let mut out_stream = delay.transform_events(Box::pin(rx));
276
277        tx.send(LogEvent::default().into()).await.unwrap();
278
279        // We should be pending, because we are now waiting for the delay
280        assert_eq!(Poll::Pending, futures::poll!(out_stream.next()));
281
282        // Wait long enough for delay to end
283        tokio::time::sleep(Duration::from_secs_f64(0.3)).await;
284
285        if !matches!(futures::poll!(out_stream.next()), Poll::Ready(Some(_event))) {
286            panic!("Unexpectedly received None or Pending in output stream");
287        }
288    }
289
290    #[tokio::test]
291    async fn delay_events_at_capacity_drop_newest() {
292        let config = serde_yaml::from_str::<DelayConfig>(indoc! {"
293            delay_ms: 200
294            queue_capacity: 1
295            overflow_strategy: drop_newest
296        "})
297        .unwrap();
298
299        let delay =
300            Transform::event_task(Delay::new(&config, &TransformContext::default()).unwrap());
301
302        let delay = delay.into_task();
303
304        let (mut tx, rx) = futures::channel::mpsc::channel(10);
305        let mut out_stream = delay.transform_events(Box::pin(rx));
306
307        tx.send(LogEvent::default().into()).await.unwrap();
308        tx.send(TraceEvent::default().into()).await.unwrap();
309
310        // We should be pending, because we are now waiting for the delay
311        assert_eq!(Poll::Pending, futures::poll!(out_stream.next()));
312
313        // Wait long enough for delay to end
314        tokio::time::sleep(Duration::from_secs_f64(0.3)).await;
315
316        let Poll::Ready(Some(event)) = futures::poll!(out_stream.next()) else {
317            panic!("Unexpectedly received None or Pending in output stream");
318        };
319        assert!(event.try_into_log().is_some());
320
321        // We should be pending, because trace event should have been dropped
322        assert_eq!(Poll::Pending, futures::poll!(out_stream.next()));
323    }
324
325    #[tokio::test]
326    async fn delay_events_at_capacity_pass() {
327        let config = serde_yaml::from_str::<DelayConfig>(indoc! {"
328            delay_ms: 200
329            queue_capacity: 1
330            overflow_strategy: forward
331        "})
332        .unwrap();
333
334        let delay =
335            Transform::event_task(Delay::new(&config, &TransformContext::default()).unwrap());
336
337        let delay = delay.into_task();
338
339        let (mut tx, rx) = futures::channel::mpsc::channel(10);
340        let mut out_stream = delay.transform_events(Box::pin(rx));
341
342        tx.send(LogEvent::default().into()).await.unwrap();
343        tx.send(TraceEvent::default().into()).await.unwrap();
344
345        // First event should be trace, because it is passed right away before delay
346        let Poll::Ready(Some(event)) = futures::poll!(out_stream.next()) else {
347            panic!("Unexpectedly received None or Pending in output stream");
348        };
349        assert!(event.try_into_trace().is_some());
350
351        // We should be pending, because we are now waiting for the delay
352        assert_eq!(Poll::Pending, futures::poll!(out_stream.next()));
353
354        // Wait long enough for delay to end
355        tokio::time::sleep(Duration::from_secs_f64(0.3)).await;
356
357        let Poll::Ready(Some(event)) = futures::poll!(out_stream.next()) else {
358            panic!("Unexpectedly received None or Pending in output stream");
359        };
360        assert!(event.try_into_log().is_some());
361    }
362}