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#[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 #[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 #[serde(default = "default_queue_capacity")]
33 queue_capacity: NonZeroUsize,
34
35 #[serde(default)]
37 overflow_strategy: OverflowStrategy,
38
39 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#[configurable_component]
60#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
61#[serde(rename_all = "snake_case")]
62pub enum OverflowStrategy {
63 #[default]
69 Block,
70
71 DropNewest,
77
78 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 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 (!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 assert_eq!(Poll::Pending, futures::poll!(out_stream.next()));
281
282 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 assert_eq!(Poll::Pending, futures::poll!(out_stream.next()));
312
313 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 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 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 assert_eq!(Poll::Pending, futures::poll!(out_stream.next()));
353
354 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}