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 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 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 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 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
166pub 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)] 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 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 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 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}