codecs/decoding/framing/
newline_delimited.rs1use bytes::{Bytes, BytesMut};
2use tokio_util::codec::Decoder;
3use vector_config::configurable_component;
4
5use super::{BoxedFramingError, CharacterDelimitedDecoder, OversizedAction};
6
7#[configurable_component]
9#[derive(Debug, Clone, Default, PartialEq, Eq)]
10pub struct NewlineDelimitedDecoderConfig {
11 #[serde(default, skip_serializing_if = "vector_core::serde::is_default")]
13 pub newline_delimited: NewlineDelimitedDecoderOptions,
14}
15
16#[configurable_component]
18#[derive(Clone, Debug, Default, PartialEq, Eq)]
19pub struct NewlineDelimitedDecoderOptions {
20 #[serde(skip_serializing_if = "vector_core::serde::is_default")]
32 pub max_length: Option<usize>,
33
34 #[serde(default, skip_serializing_if = "vector_core::serde::is_default")]
42 pub oversized_action: OversizedAction,
43}
44
45impl NewlineDelimitedDecoderOptions {
46 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 pub fn new() -> Self {
58 Default::default()
59 }
60
61 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 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#[derive(Debug, Clone)]
82pub struct NewlineDelimitedDecoder(CharacterDelimitedDecoder);
83
84impl NewlineDelimitedDecoder {
85 pub const fn new() -> Self {
87 Self(CharacterDelimitedDecoder::new(b'\n'))
88 }
89
90 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 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}