Skip to main content

vector/sinks/util/
encoding.rs

1use std::io;
2
3use bytes::BytesMut;
4use itertools::Itertools;
5use tokio_util::codec::Encoder as _;
6use vector_lib::{
7    EstimatedJsonEncodedSizeOf,
8    codecs::{Transformer, encoding::Framer, internal_events::EncoderWriteError},
9    config::telemetry,
10    request_metadata::GroupedCountByteSize,
11};
12
13use crate::event::Event;
14
15pub trait Encoder<T> {
16    /// Encodes the input into the provided writer.
17    ///
18    /// # Errors
19    ///
20    /// If an I/O error is encountered while encoding the input, an error variant will be returned.
21    fn encode_input(
22        &self,
23        input: T,
24        writer: &mut dyn io::Write,
25    ) -> io::Result<(usize, GroupedCountByteSize)>;
26}
27
28impl Encoder<Vec<Event>> for (Transformer, vector_lib::codecs::Encoder<Framer>) {
29    fn encode_input(
30        &self,
31        events: Vec<Event>,
32        writer: &mut dyn io::Write,
33    ) -> io::Result<(usize, GroupedCountByteSize)> {
34        let mut encoder = self.1.clone();
35        let mut bytes_written = 0;
36        let mut n_events_pending = events.len();
37        let is_empty = events.is_empty();
38        let batch_prefix = encoder.batch_prefix();
39        write_all(writer, n_events_pending, batch_prefix)?;
40        bytes_written += batch_prefix.len();
41
42        let mut byte_size = telemetry().create_request_count_byte_size();
43
44        for (position, mut event) in events.into_iter().with_position() {
45            self.0.transform(&mut event);
46
47            // Ensure the json size is calculated after any fields have been removed
48            // by the transformer.
49            byte_size.add_event(&event, event.estimated_json_encoded_size_of());
50
51            let mut bytes = BytesMut::new();
52            match (position, encoder.framer()) {
53                (position, Framer::CharacterDelimited(_) | Framer::NewlineDelimited(_))
54                    if position.is_last() =>
55                {
56                    encoder
57                        .serialize(event, &mut bytes)
58                        .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?;
59                }
60                _ => {
61                    encoder
62                        .encode(event, &mut bytes)
63                        .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?;
64                }
65            }
66            write_all(writer, n_events_pending, &bytes)?;
67            bytes_written += bytes.len();
68            n_events_pending -= 1;
69        }
70
71        let batch_suffix = encoder.batch_suffix(is_empty);
72        assert!(n_events_pending == 0);
73        write_all(writer, 0, batch_suffix)?;
74        bytes_written += batch_suffix.len();
75
76        Ok((bytes_written, byte_size))
77    }
78}
79
80impl Encoder<Event> for (Transformer, vector_lib::codecs::Encoder<()>) {
81    fn encode_input(
82        &self,
83        mut event: Event,
84        writer: &mut dyn io::Write,
85    ) -> io::Result<(usize, GroupedCountByteSize)> {
86        let mut encoder = self.1.clone();
87        self.0.transform(&mut event);
88
89        let mut byte_size = telemetry().create_request_count_byte_size();
90        byte_size.add_event(&event, event.estimated_json_encoded_size_of());
91
92        let mut bytes = BytesMut::new();
93        encoder
94            .serialize(event, &mut bytes)
95            .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?;
96        write_all(writer, 1, &bytes)?;
97        Ok((bytes.len(), byte_size))
98    }
99}
100
101#[cfg(feature = "codecs-arrow")]
102impl Encoder<Vec<Event>> for (Transformer, vector_lib::codecs::BatchEncoder) {
103    fn encode_input(
104        &self,
105        events: Vec<Event>,
106        writer: &mut dyn io::Write,
107    ) -> io::Result<(usize, GroupedCountByteSize)> {
108        use tokio_util::codec::Encoder as _;
109        use vector_lib::internal_event::{ComponentEventsDropped, UNINTENTIONAL};
110
111        let mut encoder = self.1.clone();
112        let mut byte_size = telemetry().create_request_count_byte_size();
113        let n_events = events.len();
114        let mut transformed_events = Vec::with_capacity(n_events);
115
116        for mut event in events {
117            self.0.transform(&mut event);
118            byte_size.add_event(&event, event.estimated_json_encoded_size_of());
119            transformed_events.push(event);
120        }
121
122        let mut bytes = BytesMut::new();
123        encoder
124            .encode(transformed_events, &mut bytes)
125            .map_err(|error| {
126                // Codec error paths emit their own internal event
127                // (e.g. SchemaGenerationError, EncoderNullConstraintError,
128                // EncoderRecordBatchError) which logs the error and increments
129                // component_errors_total. We only emit the drop count here to
130                // avoid double-counting.
131                // n_events is the pre-filter count; Parquet filters non-log
132                // events before encoding, but that only happens if a sink is
133                // misconfigured to send non-log events into a log-only encoder,
134                // so the overcount is not a practical concern.
135                emit!(ComponentEventsDropped::<UNINTENTIONAL> {
136                    count: n_events,
137                    reason: "Failed to batch encode events.",
138                });
139                io::Error::new(io::ErrorKind::InvalidData, error)
140            })?;
141
142        write_all(writer, n_events, &bytes)?;
143        Ok((bytes.len(), byte_size))
144    }
145}
146
147impl Encoder<Vec<Event>> for (Transformer, vector_lib::codecs::EncoderKind) {
148    fn encode_input(
149        &self,
150        events: Vec<Event>,
151        writer: &mut dyn io::Write,
152    ) -> io::Result<(usize, GroupedCountByteSize)> {
153        // Delegate to the specific encoder implementation
154        match &self.1 {
155            vector_lib::codecs::EncoderKind::Framed(encoder) => {
156                (self.0.clone(), *encoder.clone()).encode_input(events, writer)
157            }
158            #[cfg(feature = "codecs-arrow")]
159            vector_lib::codecs::EncoderKind::Batch(encoder) => {
160                (self.0.clone(), encoder.clone()).encode_input(events, writer)
161            }
162        }
163    }
164}
165
166/// Write the buffer to the writer. If the operation fails, emit an internal event which complies with the
167/// instrumentation spec- as this necessitates both an Error and EventsDropped event.
168///
169/// # Arguments
170///
171/// * `writer`           - The object implementing io::Write to write data to.
172/// * `n_events_pending` - The number of events that are dropped if this write fails.
173/// * `buf`              - The buffer to write.
174pub fn write_all(
175    writer: &mut dyn io::Write,
176    n_events_pending: usize,
177    buf: &[u8],
178) -> io::Result<()> {
179    writer.write_all(buf).inspect_err(|error| {
180        emit!(EncoderWriteError {
181            error,
182            count: n_events_pending,
183        });
184    })
185}
186
187pub fn as_tracked_write<F, I, E>(inner: &mut dyn io::Write, input: I, f: F) -> io::Result<usize>
188where
189    F: FnOnce(&mut dyn io::Write, I) -> Result<(), E>,
190    E: Into<io::Error> + 'static,
191{
192    struct Tracked<'inner> {
193        count: usize,
194        inner: &'inner mut dyn io::Write,
195    }
196
197    impl io::Write for Tracked<'_> {
198        fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
199            #[allow(clippy::disallowed_methods)] // We pass on the result of `write` to the caller.
200            let n = self.inner.write(buf)?;
201            self.count += n;
202            Ok(n)
203        }
204
205        fn flush(&mut self) -> io::Result<()> {
206            self.inner.flush()
207        }
208    }
209
210    let mut tracked = Tracked { count: 0, inner };
211    f(&mut tracked, input).map_err(|e| e.into())?;
212    Ok(tracked.count)
213}
214
215#[cfg(test)]
216mod tests {
217    use std::{collections::BTreeMap, env, path::PathBuf};
218
219    use bytes::{BufMut, Bytes};
220    use cfg_if::cfg_if;
221    use vector_lib::{
222        codecs::{
223            CharacterDelimitedEncoder, JsonSerializerConfig, LengthDelimitedEncoder,
224            NewlineDelimitedEncoder, TextSerializerConfig,
225            encoding::{ProtobufSerializerConfig, ProtobufSerializerOptions},
226        },
227        event::LogEvent,
228        internal_event::CountByteSize,
229        json_size::JsonSize,
230    };
231    use vrl::value::{KeyString, Value};
232
233    cfg_if! {
234        if #[cfg(feature = "codecs-arrow")] {
235            use arrow::datatypes::{DataType, Field, Schema as ArrowSchema};
236            use vector_lib::codecs::{
237                BatchEncoder,
238                encoding::{ArrowStreamSerializer, ArrowStreamSerializerConfig, BatchSerializer},
239            };
240            use vector_lib::event_test_util::{clear_recorded_events, contains_name_once};
241        }
242    }
243
244    use super::*;
245
246    #[test]
247    fn test_encode_batch_json_empty() {
248        let encoding = (
249            Transformer::default(),
250            vector_lib::codecs::Encoder::<Framer>::new(
251                CharacterDelimitedEncoder::new(b',').into(),
252                JsonSerializerConfig::default().build().into(),
253            ),
254        );
255
256        let mut writer = Vec::new();
257        let (written, json_size) = encoding.encode_input(vec![], &mut writer).unwrap();
258        assert_eq!(written, 2);
259
260        assert_eq!(String::from_utf8(writer).unwrap(), "[]");
261        assert_eq!(
262            CountByteSize(0, JsonSize::zero()),
263            json_size.size().unwrap()
264        );
265    }
266
267    #[test]
268    fn test_encode_batch_json_single() {
269        let encoding = (
270            Transformer::default(),
271            vector_lib::codecs::Encoder::<Framer>::new(
272                CharacterDelimitedEncoder::new(b',').into(),
273                JsonSerializerConfig::default().build().into(),
274            ),
275        );
276
277        let mut writer = Vec::new();
278        let input = vec![Event::Log(LogEvent::from(BTreeMap::from([(
279            KeyString::from("key"),
280            Value::from("value"),
281        )])))];
282
283        let input_json_size = input
284            .iter()
285            .map(|event| event.estimated_json_encoded_size_of())
286            .sum::<JsonSize>();
287
288        let (written, json_size) = encoding.encode_input(input, &mut writer).unwrap();
289        assert_eq!(written, 17);
290
291        assert_eq!(String::from_utf8(writer).unwrap(), r#"[{"key":"value"}]"#);
292        assert_eq!(CountByteSize(1, input_json_size), json_size.size().unwrap());
293    }
294
295    #[test]
296    fn test_encode_batch_json_multiple() {
297        let encoding = (
298            Transformer::default(),
299            vector_lib::codecs::Encoder::<Framer>::new(
300                CharacterDelimitedEncoder::new(b',').into(),
301                JsonSerializerConfig::default().build().into(),
302            ),
303        );
304
305        let input = vec![
306            Event::Log(LogEvent::from(BTreeMap::from([(
307                KeyString::from("key"),
308                Value::from("value1"),
309            )]))),
310            Event::Log(LogEvent::from(BTreeMap::from([(
311                KeyString::from("key"),
312                Value::from("value2"),
313            )]))),
314            Event::Log(LogEvent::from(BTreeMap::from([(
315                KeyString::from("key"),
316                Value::from("value3"),
317            )]))),
318        ];
319
320        let input_json_size = input
321            .iter()
322            .map(|event| event.estimated_json_encoded_size_of())
323            .sum::<JsonSize>();
324
325        let mut writer = Vec::new();
326        let (written, json_size) = encoding.encode_input(input, &mut writer).unwrap();
327        assert_eq!(written, 52);
328
329        assert_eq!(
330            String::from_utf8(writer).unwrap(),
331            r#"[{"key":"value1"},{"key":"value2"},{"key":"value3"}]"#
332        );
333
334        assert_eq!(CountByteSize(3, input_json_size), json_size.size().unwrap());
335    }
336
337    #[test]
338    fn test_encode_batch_ndjson_empty() {
339        let encoding = (
340            Transformer::default(),
341            vector_lib::codecs::Encoder::<Framer>::new(
342                NewlineDelimitedEncoder::default().into(),
343                JsonSerializerConfig::default().build().into(),
344            ),
345        );
346
347        let mut writer = Vec::new();
348        let (written, json_size) = encoding.encode_input(vec![], &mut writer).unwrap();
349        assert_eq!(written, 0);
350
351        assert_eq!(String::from_utf8(writer).unwrap(), "");
352        assert_eq!(
353            CountByteSize(0, JsonSize::zero()),
354            json_size.size().unwrap()
355        );
356    }
357
358    #[test]
359    fn test_encode_batch_ndjson_single() {
360        let encoding = (
361            Transformer::default(),
362            vector_lib::codecs::Encoder::<Framer>::new(
363                NewlineDelimitedEncoder::default().into(),
364                JsonSerializerConfig::default().build().into(),
365            ),
366        );
367
368        let mut writer = Vec::new();
369        let input = vec![Event::Log(LogEvent::from(BTreeMap::from([(
370            KeyString::from("key"),
371            Value::from("value"),
372        )])))];
373        let input_json_size = input
374            .iter()
375            .map(|event| event.estimated_json_encoded_size_of())
376            .sum::<JsonSize>();
377
378        let (written, json_size) = encoding.encode_input(input, &mut writer).unwrap();
379        assert_eq!(written, 16);
380
381        assert_eq!(String::from_utf8(writer).unwrap(), "{\"key\":\"value\"}\n");
382        assert_eq!(CountByteSize(1, input_json_size), json_size.size().unwrap());
383    }
384
385    #[test]
386    fn test_encode_batch_ndjson_multiple() {
387        let encoding = (
388            Transformer::default(),
389            vector_lib::codecs::Encoder::<Framer>::new(
390                NewlineDelimitedEncoder::default().into(),
391                JsonSerializerConfig::default().build().into(),
392            ),
393        );
394
395        let mut writer = Vec::new();
396        let input = vec![
397            Event::Log(LogEvent::from(BTreeMap::from([(
398                KeyString::from("key"),
399                Value::from("value1"),
400            )]))),
401            Event::Log(LogEvent::from(BTreeMap::from([(
402                KeyString::from("key"),
403                Value::from("value2"),
404            )]))),
405            Event::Log(LogEvent::from(BTreeMap::from([(
406                KeyString::from("key"),
407                Value::from("value3"),
408            )]))),
409        ];
410        let input_json_size = input
411            .iter()
412            .map(|event| event.estimated_json_encoded_size_of())
413            .sum::<JsonSize>();
414
415        let (written, json_size) = encoding.encode_input(input, &mut writer).unwrap();
416        assert_eq!(written, 51);
417
418        assert_eq!(
419            String::from_utf8(writer).unwrap(),
420            "{\"key\":\"value1\"}\n{\"key\":\"value2\"}\n{\"key\":\"value3\"}\n"
421        );
422        assert_eq!(CountByteSize(3, input_json_size), json_size.size().unwrap());
423    }
424
425    #[test]
426    fn test_encode_event_json() {
427        let encoding = (
428            Transformer::default(),
429            vector_lib::codecs::Encoder::<()>::new(JsonSerializerConfig::default().build().into()),
430        );
431
432        let mut writer = Vec::new();
433        let input = Event::Log(LogEvent::from(BTreeMap::from([(
434            KeyString::from("key"),
435            Value::from("value"),
436        )])));
437        let input_json_size = input.estimated_json_encoded_size_of();
438
439        let (written, json_size) = encoding.encode_input(input, &mut writer).unwrap();
440        assert_eq!(written, 15);
441
442        assert_eq!(String::from_utf8(writer).unwrap(), r#"{"key":"value"}"#);
443        assert_eq!(CountByteSize(1, input_json_size), json_size.size().unwrap());
444    }
445
446    #[test]
447    fn test_encode_event_text() {
448        let encoding = (
449            Transformer::default(),
450            vector_lib::codecs::Encoder::<()>::new(TextSerializerConfig::default().build().into()),
451        );
452
453        let mut writer = Vec::new();
454        let input = Event::Log(LogEvent::from(BTreeMap::from([(
455            KeyString::from("message"),
456            Value::from("value"),
457        )])));
458        let input_json_size = input.estimated_json_encoded_size_of();
459
460        let (written, json_size) = encoding.encode_input(input, &mut writer).unwrap();
461        assert_eq!(written, 5);
462
463        assert_eq!(String::from_utf8(writer).unwrap(), r"value");
464        assert_eq!(CountByteSize(1, input_json_size), json_size.size().unwrap());
465    }
466
467    fn test_data_dir() -> PathBuf {
468        PathBuf::from(env::var_os("CARGO_MANIFEST_DIR").unwrap()).join("tests/data/protobuf")
469    }
470
471    #[test]
472    fn test_encode_batch_protobuf_single() {
473        let message_raw = std::fs::read(test_data_dir().join("test_proto.pb")).unwrap();
474        let input_proto_size = message_raw.len();
475
476        // default LengthDelimitedCoderOptions.length_field_length is 4
477        let mut buf = BytesMut::with_capacity(64);
478        buf.reserve(4 + input_proto_size);
479        buf.put_uint(input_proto_size as u64, 4);
480        buf.extend_from_slice(&message_raw[..]);
481        let expected_bytes = buf.freeze();
482
483        let config = ProtobufSerializerConfig {
484            protobuf: ProtobufSerializerOptions {
485                desc_file: test_data_dir().join("test_proto.desc"),
486                message_type: "test_proto.User".to_string(),
487                use_json_names: false,
488            },
489        };
490
491        let encoding = (
492            Transformer::default(),
493            vector_lib::codecs::Encoder::<Framer>::new(
494                LengthDelimitedEncoder::default().into(),
495                config.build().unwrap().into(),
496            ),
497        );
498
499        let mut writer = Vec::new();
500        let input = vec![Event::Log(LogEvent::from(BTreeMap::from([
501            (KeyString::from("id"), Value::from("123")),
502            (KeyString::from("name"), Value::from("Alice")),
503            (KeyString::from("age"), Value::from(30)),
504            (
505                KeyString::from("emails"),
506                Value::from(vec!["alice@example.com", "alice@work.com"]),
507            ),
508        ])))];
509
510        let input_json_size = input
511            .iter()
512            .map(|event| event.estimated_json_encoded_size_of())
513            .sum::<JsonSize>();
514
515        let (written, size) = encoding.encode_input(input, &mut writer).unwrap();
516
517        assert_eq!(input_proto_size, 49);
518        assert_eq!(written, input_proto_size + 4);
519        assert_eq!(CountByteSize(1, input_json_size), size.size().unwrap());
520        assert_eq!(Bytes::copy_from_slice(&writer), expected_bytes);
521    }
522
523    #[test]
524    fn test_encode_batch_protobuf_multiple() {
525        let message_raw = std::fs::read(test_data_dir().join("test_proto.pb")).unwrap();
526        let messages = vec![message_raw.clone(), message_raw.clone()];
527        let total_input_proto_size: usize = messages.iter().map(|m| m.len()).sum();
528
529        let mut buf = BytesMut::with_capacity(128);
530        for message in messages {
531            // default LengthDelimitedCoderOptions.length_field_length is 4
532            buf.reserve(4 + message.len());
533            buf.put_uint(message.len() as u64, 4);
534            buf.extend_from_slice(&message[..]);
535        }
536        let expected_bytes = buf.freeze();
537
538        let config = ProtobufSerializerConfig {
539            protobuf: ProtobufSerializerOptions {
540                desc_file: test_data_dir().join("test_proto.desc"),
541                message_type: "test_proto.User".to_string(),
542                use_json_names: false,
543            },
544        };
545
546        let encoding = (
547            Transformer::default(),
548            vector_lib::codecs::Encoder::<Framer>::new(
549                LengthDelimitedEncoder::default().into(),
550                config.build().unwrap().into(),
551            ),
552        );
553
554        let mut writer = Vec::new();
555        let input = vec![
556            Event::Log(LogEvent::from(BTreeMap::from([
557                (KeyString::from("id"), Value::from("123")),
558                (KeyString::from("name"), Value::from("Alice")),
559                (KeyString::from("age"), Value::from(30)),
560                (
561                    KeyString::from("emails"),
562                    Value::from(vec!["alice@example.com", "alice@work.com"]),
563                ),
564            ]))),
565            Event::Log(LogEvent::from(BTreeMap::from([
566                (KeyString::from("id"), Value::from("123")),
567                (KeyString::from("name"), Value::from("Alice")),
568                (KeyString::from("age"), Value::from(30)),
569                (
570                    KeyString::from("emails"),
571                    Value::from(vec!["alice@example.com", "alice@work.com"]),
572                ),
573            ]))),
574        ];
575
576        let input_json_size: JsonSize = input
577            .iter()
578            .map(|event| event.estimated_json_encoded_size_of())
579            .sum();
580
581        let (written, size) = encoding.encode_input(input, &mut writer).unwrap();
582
583        assert_eq!(total_input_proto_size, 49 * 2);
584        assert_eq!(written, total_input_proto_size + 8);
585        assert_eq!(CountByteSize(2, input_json_size), size.size().unwrap());
586        assert_eq!(Bytes::copy_from_slice(&writer), expected_bytes);
587    }
588
589    #[cfg(feature = "codecs-arrow")]
590    #[test]
591    fn test_encode_batch_arrow_emits_record_batch_error_on_type_mismatch() {
592        clear_recorded_events();
593
594        // Schema declares `message` as Int64, but the event below carries a string,
595        // so `build_record_batch` returns `ArrowEncodingError::ArrowJsonDecode`.
596        let schema = ArrowSchema::new(vec![Field::new("message", DataType::Int64, false)]);
597        let serializer = ArrowStreamSerializer::new(ArrowStreamSerializerConfig::new(schema))
598            .expect("failed to build ArrowStreamSerializer");
599        let encoder = BatchEncoder::new(BatchSerializer::Arrow(serializer));
600        let encoding = (Transformer::default(), encoder);
601
602        let event = Event::Log(LogEvent::from(BTreeMap::from([(
603            KeyString::from("message"),
604            Value::from("not_an_integer"),
605        )])));
606
607        let mut writer = Vec::new();
608        let result = encoding.encode_input(vec![event], &mut writer);
609        assert!(
610            result.is_err(),
611            "type mismatch should fail batch encoding, got {result:?}"
612        );
613
614        contains_name_once("EncoderRecordBatchError")
615            .expect("EncoderRecordBatchError should be emitted on ArrowJsonDecode failure");
616        contains_name_once("ComponentEventsDropped")
617            .expect("ComponentEventsDropped should be emitted by the wrapper");
618    }
619}