Skip to main content

codecs/decoding/framing/
newline_delimited.rs

1use bytes::{Bytes, BytesMut};
2use tokio_util::codec::Decoder;
3use vector_config::configurable_component;
4
5use super::{BoxedFramingError, CharacterDelimitedDecoder, OversizedAction};
6
7/// Config used to build a `NewlineDelimitedDecoder`.
8#[configurable_component]
9#[derive(Debug, Clone, Default, PartialEq, Eq)]
10pub struct NewlineDelimitedDecoderConfig {
11    /// Options for the newline delimited decoder.
12    #[serde(default, skip_serializing_if = "vector_core::serde::is_default")]
13    pub newline_delimited: NewlineDelimitedDecoderOptions,
14}
15
16/// Options for building a `NewlineDelimitedDecoder`.
17#[configurable_component]
18#[derive(Clone, Debug, Default, PartialEq, Eq)]
19pub struct NewlineDelimitedDecoderOptions {
20    /// The maximum length of the byte buffer.
21    ///
22    /// This length does *not* include the trailing delimiter.
23    ///
24    /// By default, no maximum length is enforced. If events are malformed, this can lead to
25    /// additional resource usage as events continue to be buffered in memory, and can potentially
26    /// lead to memory exhaustion in extreme cases.
27    ///
28    /// If there is a risk of processing malformed data, such as logs with user-controlled input,
29    /// consider setting the maximum length to a reasonably large value as a safety net. This
30    /// prevents processing from being unbounded.
31    #[serde(skip_serializing_if = "vector_core::serde::is_default")]
32    pub max_length: Option<usize>,
33
34    /// The behavior when a line exceeds `max_length`.
35    ///
36    /// When set to `drop` (the default), the entire oversized line is discarded.
37    /// When set to `truncate`, the line is truncated to `max_length` bytes and the
38    /// remainder is discarded up to the next newline.
39    ///
40    /// This option has no effect if `max_length` is not set.
41    #[serde(default, skip_serializing_if = "vector_core::serde::is_default")]
42    pub oversized_action: OversizedAction,
43}
44
45impl NewlineDelimitedDecoderOptions {
46    /// Creates a `NewlineDelimitedDecoderOptions` with a maximum frame length limit.
47    pub const fn new_with_max_length(max_length: usize) -> Self {
48        Self {
49            max_length: Some(max_length),
50            oversized_action: OversizedAction::Drop,
51        }
52    }
53}
54
55impl NewlineDelimitedDecoderConfig {
56    /// Creates a new `NewlineDelimitedDecoderConfig`.
57    pub fn new() -> Self {
58        Default::default()
59    }
60
61    /// Creates a `NewlineDelimitedDecoder` with a maximum frame length limit.
62    pub const fn new_with_max_length(max_length: usize) -> Self {
63        Self {
64            newline_delimited: { NewlineDelimitedDecoderOptions::new_with_max_length(max_length) },
65        }
66    }
67
68    /// Build the `NewlineDelimitedDecoder` from this configuration.
69    pub const fn build(&self) -> NewlineDelimitedDecoder {
70        let oversized_action = self.newline_delimited.oversized_action;
71        if let Some(max_length) = self.newline_delimited.max_length {
72            NewlineDelimitedDecoder::new_with_max_length(max_length)
73                .with_oversized_action(oversized_action)
74        } else {
75            NewlineDelimitedDecoder::new()
76        }
77    }
78}
79
80/// A codec for handling bytes that are delimited by (a) newline(s).
81#[derive(Debug, Clone)]
82pub struct NewlineDelimitedDecoder(CharacterDelimitedDecoder);
83
84impl NewlineDelimitedDecoder {
85    /// Creates a new `NewlineDelimitedDecoder`.
86    pub const fn new() -> Self {
87        Self(CharacterDelimitedDecoder::new(b'\n'))
88    }
89
90    /// Creates a `NewlineDelimitedDecoder` with a maximum frame length limit.
91    ///
92    /// Any frames longer than `max_length` bytes will be discarded entirely.
93    pub const fn new_with_max_length(max_length: usize) -> Self {
94        Self(CharacterDelimitedDecoder::new_with_max_length(
95            b'\n', max_length,
96        ))
97    }
98
99    /// Sets the behavior when a line exceeds `max_length`.
100    pub const fn with_oversized_action(mut self, action: OversizedAction) -> Self {
101        self.0.oversized_action = action;
102        self
103    }
104}
105
106impl Default for NewlineDelimitedDecoder {
107    fn default() -> Self {
108        Self::new()
109    }
110}
111
112impl Decoder for NewlineDelimitedDecoder {
113    type Item = Bytes;
114    type Error = BoxedFramingError;
115
116    fn decode(&mut self, src: &mut BytesMut) -> Result<Option<Self::Item>, Self::Error> {
117        self.0.decode(src)
118    }
119
120    fn decode_eof(&mut self, src: &mut BytesMut) -> Result<Option<Self::Item>, Self::Error> {
121        self.0.decode_eof(src)
122    }
123}
124
125#[cfg(test)]
126mod tests {
127    use super::*;
128
129    #[test]
130    fn decode_bytes_with_newlines() {
131        let mut input = BytesMut::from("foo\nbar\nbaz");
132        let mut decoder = NewlineDelimitedDecoder::new();
133
134        assert_eq!(decoder.decode(&mut input).unwrap().unwrap(), "foo");
135        assert_eq!(decoder.decode(&mut input).unwrap().unwrap(), "bar");
136        assert_eq!(decoder.decode(&mut input).unwrap(), None);
137    }
138
139    #[test]
140    fn decode_bytes_with_newlines_trailing() {
141        let mut input = BytesMut::from("foo\nbar\nbaz\n");
142        let mut decoder = NewlineDelimitedDecoder::new();
143
144        assert_eq!(decoder.decode(&mut input).unwrap().unwrap(), "foo");
145        assert_eq!(decoder.decode(&mut input).unwrap().unwrap(), "bar");
146        assert_eq!(decoder.decode(&mut input).unwrap().unwrap(), "baz");
147        assert_eq!(decoder.decode(&mut input).unwrap(), None);
148    }
149
150    #[test]
151    fn decode_bytes_with_newlines_and_max_length() {
152        let mut input = BytesMut::from("foo\nbarbara\nbaz\n");
153        let mut decoder = NewlineDelimitedDecoder::new_with_max_length(3);
154
155        assert_eq!(decoder.decode(&mut input).unwrap().unwrap(), "foo");
156        assert_eq!(decoder.decode(&mut input).unwrap().unwrap(), "baz");
157        assert_eq!(decoder.decode(&mut input).unwrap(), None);
158    }
159
160    #[test]
161    fn decode_eof_bytes_with_newlines() {
162        let mut input = BytesMut::from("foo\nbar\nbaz");
163        let mut decoder = NewlineDelimitedDecoder::new();
164
165        assert_eq!(decoder.decode_eof(&mut input).unwrap().unwrap(), "foo");
166        assert_eq!(decoder.decode_eof(&mut input).unwrap().unwrap(), "bar");
167        assert_eq!(decoder.decode_eof(&mut input).unwrap().unwrap(), "baz");
168    }
169
170    #[test]
171    fn decode_eof_bytes_with_newlines_trailing() {
172        let mut input = BytesMut::from("foo\nbar\nbaz\n");
173        let mut decoder = NewlineDelimitedDecoder::new();
174
175        assert_eq!(decoder.decode_eof(&mut input).unwrap().unwrap(), "foo");
176        assert_eq!(decoder.decode_eof(&mut input).unwrap().unwrap(), "bar");
177        assert_eq!(decoder.decode_eof(&mut input).unwrap().unwrap(), "baz");
178        assert_eq!(decoder.decode_eof(&mut input).unwrap(), None);
179    }
180
181    #[test]
182    fn decode_eof_bytes_with_newlines_and_max_length() {
183        let mut input = BytesMut::from("foo\nbarbara\nbaz\n");
184        let mut decoder = NewlineDelimitedDecoder::new_with_max_length(3);
185
186        assert_eq!(decoder.decode_eof(&mut input).unwrap().unwrap(), "foo");
187        assert_eq!(decoder.decode_eof(&mut input).unwrap().unwrap(), "baz");
188        assert_eq!(decoder.decode_eof(&mut input).unwrap(), None);
189    }
190}