1use std::{
2 collections::BTreeMap,
3 convert::TryFrom,
4 io,
5 net::SocketAddr,
6 num::{NonZeroU64, NonZeroUsize},
7 time::Duration,
8};
9
10use bytes::{Buf, Bytes, BytesMut};
11use smallvec::{SmallVec, smallvec};
12use snafu::{ResultExt, Snafu};
13use tokio_util::codec::Decoder;
14use vector_lib::{
15 codecs::{BytesDeserializerConfig, StreamDecodingError},
16 config::{LegacyKey, LogNamespace},
17 configurable::configurable_component,
18 ipallowlist::IpAllowlistConfig,
19 lookup::{OwnedValuePath, event_path, metadata_path, owned_value_path, path},
20 schema::Definition,
21};
22use vrl::value::{KeyString, Kind, kind::Collection};
23
24use super::util::decompression::{
25 CappedDecoder, max_decompressed_size_bytes, max_zlib_compressed_frame_size_bytes,
26};
27use super::util::net::{SocketListenAddr, TcpSource, TcpSourceAck, TcpSourceAcker};
28use crate::{
29 config::{
30 DataType, GenerateConfig, Resource, SourceAcknowledgementsConfig, SourceConfig,
31 SourceContext, SourceOutput, log_schema,
32 },
33 event::{Event, LogEvent, Value},
34 serde::bool_or_struct,
35 tcp::TcpKeepaliveConfig,
36 tls::{MaybeTlsSettings, TlsSourceConfig},
37 types,
38};
39
40#[configurable_component(source("logstash", "Collect logs from a Logstash agent."))]
42#[derive(Clone, Debug)]
43pub struct LogstashConfig {
44 #[configurable(derived)]
45 address: SocketListenAddr,
46
47 #[configurable(derived)]
48 keepalive: Option<TcpKeepaliveConfig>,
49
50 #[configurable(derived)]
51 pub permit_origin: Option<IpAllowlistConfig>,
52
53 #[configurable(derived)]
54 tls: Option<TlsSourceConfig>,
55
56 #[configurable(metadata(docs::type_unit = "bytes"))]
58 #[configurable(metadata(docs::examples = 65536))]
59 receive_buffer_bytes: Option<usize>,
60
61 #[configurable(metadata(docs::type_unit = "connections"))]
63 connection_limit: Option<u32>,
64
65 #[configurable(metadata(docs::type_unit = "seconds"))]
71 tls_handshake_timeout_secs: Option<NonZeroU64>,
72
73 #[configurable(derived)]
74 #[serde(default, deserialize_with = "bool_or_struct")]
75 acknowledgements: SourceAcknowledgementsConfig,
76
77 #[configurable(metadata(docs::hidden))]
79 #[serde(default)]
80 log_namespace: Option<bool>,
81}
82
83impl LogstashConfig {
84 fn schema_definition(&self, log_namespace: LogNamespace) -> Definition {
86 let host_key = log_schema()
88 .host_key()
89 .cloned()
90 .map(LegacyKey::InsertIfEmpty);
91
92 let tls_client_metadata_path = self
93 .tls
94 .as_ref()
95 .and_then(|tls| tls.client_metadata_key.as_ref())
96 .and_then(|k| k.path.clone())
97 .map(LegacyKey::Overwrite);
98
99 BytesDeserializerConfig
100 .schema_definition(log_namespace)
101 .with_standard_vector_source_metadata()
102 .with_source_metadata(
103 LogstashConfig::NAME,
104 None,
105 &owned_value_path!("timestamp"),
106 Kind::timestamp().or_undefined(),
107 Some("timestamp"),
108 )
109 .with_source_metadata(
110 LogstashConfig::NAME,
111 host_key,
112 &owned_value_path!("host"),
113 Kind::bytes(),
114 Some("host"),
115 )
116 .with_source_metadata(
117 Self::NAME,
118 tls_client_metadata_path,
119 &owned_value_path!("tls_client_metadata"),
120 Kind::object(Collection::empty().with_unknown(Kind::bytes())).or_undefined(),
121 None,
122 )
123 }
124}
125
126impl Default for LogstashConfig {
127 fn default() -> Self {
128 Self {
129 address: SocketListenAddr::SocketAddr("0.0.0.0:5044".parse().unwrap()),
130 keepalive: None,
131 permit_origin: None,
132 tls: None,
133 receive_buffer_bytes: None,
134 acknowledgements: Default::default(),
135 connection_limit: None,
136 tls_handshake_timeout_secs: None,
137 log_namespace: None,
138 }
139 }
140}
141
142impl GenerateConfig for LogstashConfig {
143 fn generate_config() -> serde_json::Value {
144 serde_json::to_value(LogstashConfig::default()).unwrap()
145 }
146}
147
148#[async_trait::async_trait]
149#[typetag::serde(name = "logstash")]
150impl SourceConfig for LogstashConfig {
151 async fn build(&self, cx: SourceContext) -> crate::Result<super::Source> {
152 let log_namespace = cx.log_namespace(self.log_namespace);
153 let source = LogstashSource {
154 timestamp_converter: types::Conversion::Timestamp(cx.globals.timezone()),
155 legacy_host_key_path: log_schema().host_key().cloned(),
156 log_namespace,
157 };
158 let shutdown_secs = Duration::from_secs(30);
159 let tls_config = self.tls.as_ref().map(|tls| tls.tls_config.clone());
160 let tls_client_metadata_key = self
161 .tls
162 .as_ref()
163 .and_then(|tls| tls.client_metadata_key.clone())
164 .and_then(|k| k.path);
165
166 let tls = MaybeTlsSettings::from_config(tls_config.as_ref(), true)?;
167 source.run(
168 self.address,
169 self.keepalive,
170 shutdown_secs,
171 tls,
172 None, tls_client_metadata_key,
174 self.receive_buffer_bytes,
175 None,
176 self.tls_handshake_timeout_secs,
177 cx,
178 self.acknowledgements,
179 self.connection_limit,
180 self.permit_origin.clone().map(Into::into),
181 LogstashConfig::NAME,
182 log_namespace,
183 )
184 }
185
186 fn outputs(&self, global_log_namespace: LogNamespace) -> Vec<SourceOutput> {
187 vec![SourceOutput::new_maybe_logs(
190 DataType::Log,
191 self.schema_definition(global_log_namespace.merge(self.log_namespace)),
192 )]
193 }
194
195 fn resources(&self) -> Vec<Resource> {
196 vec![self.address.as_tcp_resource()]
197 }
198
199 fn can_acknowledge(&self) -> bool {
200 true
201 }
202}
203
204#[derive(Debug, Clone)]
205struct LogstashSource {
206 timestamp_converter: types::Conversion,
207 log_namespace: LogNamespace,
208 legacy_host_key_path: Option<OwnedValuePath>,
209}
210
211impl TcpSource for LogstashSource {
212 type Error = DecodeError;
213 type Item = LogstashEventFrame;
214 type Decoder = LogstashDecoder;
215 type Acker = LogstashAcker;
216
217 fn decoder(&self) -> Self::Decoder {
218 LogstashDecoder::new()
219 }
220
221 fn handle_events(&self, events: &mut [Event], host: SocketAddr) {
222 let now = chrono::Utc::now();
223 for event in events {
224 let log = event.as_mut_log();
225
226 self.log_namespace.insert_vector_metadata(
227 log,
228 log_schema().source_type_key(),
229 path!("source_type"),
230 Bytes::from_static(LogstashConfig::NAME.as_bytes()),
231 );
232
233 let log_timestamp = log.get(event_path!("@timestamp")).and_then(|timestamp| {
234 self.timestamp_converter
235 .convert::<Value>(timestamp.coerce_to_bytes())
236 .ok()
237 });
238
239 match self.log_namespace {
244 LogNamespace::Vector => {
245 if let Some(timestamp) = log_timestamp {
246 log.insert(metadata_path!(LogstashConfig::NAME, "timestamp"), timestamp);
247 }
248 log.insert(metadata_path!("vector", "ingest_timestamp"), now);
249 }
250 LogNamespace::Legacy => {
251 if let Some(timestamp_key) = log_schema().timestamp_key_target_path() {
252 log.insert(
253 timestamp_key,
254 log_timestamp.unwrap_or_else(|| Value::from(now)),
255 );
256 }
257 }
258 }
259
260 self.log_namespace.insert_source_metadata(
261 LogstashConfig::NAME,
262 log,
263 self.legacy_host_key_path
264 .as_ref()
265 .map(LegacyKey::InsertIfEmpty),
266 path!("host"),
267 host.ip().to_string(),
268 );
269 }
270 }
271
272 fn build_acker(&self, frames: &[Self::Item]) -> Self::Acker {
273 LogstashAcker::new(frames)
274 }
275}
276
277struct LogstashAcker {
278 acknowledgements: SmallVec<[(LogstashProtocolVersion, u32); 1]>,
297}
298
299impl LogstashAcker {
300 fn new(frames: &[LogstashEventFrame]) -> Self {
301 let acknowledgements = frames
302 .iter()
303 .filter(|frame| frame.window_end)
305 .map(|frame| (frame.protocol, frame.sequence_number))
306 .collect();
307
308 Self { acknowledgements }
309 }
310}
311
312impl TcpSourceAcker for LogstashAcker {
313 fn build_ack(self, ack: TcpSourceAck) -> Option<Bytes> {
315 match ack {
316 TcpSourceAck::Ack if !self.acknowledgements.is_empty() => {
317 let mut bytes: Vec<u8> = Vec::with_capacity(self.acknowledgements.len() * 6);
318 for (protocol_version, sequence_number) in self.acknowledgements {
319 bytes.push(protocol_version.into());
320 bytes.push(LogstashFrameType::Ack.into());
321 bytes.extend(sequence_number.to_be_bytes().iter());
322 }
323 Some(Bytes::from(bytes))
324 }
325 _ => None,
326 }
327 }
328}
329
330#[derive(Debug)]
331enum LogstashDecoderReadState {
332 ReadProtocol,
333 ReadType(LogstashProtocolVersion),
334 ReadFrame(LogstashProtocolVersion, LogstashFrameType),
335 PendingDecompressed {
338 buf: BytesMut,
339 decoder: Box<LogstashDecoder>,
340 },
341}
342
343#[derive(Debug)]
344struct LogstashDecoder {
345 state: LogstashDecoderReadState,
346 window_events_remaining: Option<NonZeroUsize>,
350 nested: bool,
357 max_frame_size: usize,
362}
363
364impl LogstashDecoder {
365 fn new() -> Self {
366 Self {
367 state: LogstashDecoderReadState::ReadProtocol,
368 window_events_remaining: None,
369 nested: false,
370 max_frame_size: max_decompressed_size_bytes(),
371 }
372 }
373
374 fn new_nested(window_events_remaining: Option<NonZeroUsize>) -> Self {
375 Self {
376 state: LogstashDecoderReadState::ReadProtocol,
377 window_events_remaining,
378 nested: true,
379 max_frame_size: max_decompressed_size_bytes(),
380 }
381 }
382
383 const fn annotate_frame(&mut self, frame: &mut LogstashEventFrame) {
394 match self.window_events_remaining {
395 Some(remaining) if remaining.get() == 1 => {
396 frame.window_end = true;
397 self.window_events_remaining = None;
398 }
399 Some(remaining) => {
400 frame.window_end = false;
401 self.window_events_remaining = NonZeroUsize::new(remaining.get() - 1); }
403 None => {
404 frame.window_end = true;
407 }
408 }
409 }
410}
411
412#[derive(Debug, Snafu)]
413pub enum DecodeError {
414 #[snafu(display("i/o error: {}", source))]
415 IO { source: io::Error },
416 #[snafu(display("Unknown logstash protocol version: {}", version))]
417 UnknownProtocolVersion { version: char },
418 #[snafu(display("Unknown logstash protocol message type: {}", frame_type))]
419 UnknownFrameType { frame_type: char },
420 #[snafu(display("Failed to decode JSON frame: {}", source))]
421 JsonFrameFailedDecode { source: serde_json::Error },
422 #[snafu(display("Failed to decompress compressed frame: {}", source))]
423 DecompressionFailed { source: io::Error },
424 #[snafu(display(
425 "Received a WindowSize frame before the current window completed ({remaining} events still expected)"
426 ))]
427 PrematureWindowSize { remaining: usize },
428 #[snafu(display("Compressed frame contains a nested compressed frame"))]
429 NestedCompressedFrame,
430 #[snafu(display(
431 "logstash frame exceeds maximum size before decoding: {size} bytes buffered, limit is {max} bytes"
432 ))]
433 FrameTooLarge { size: usize, max: usize },
434}
435
436impl StreamDecodingError for DecodeError {
437 fn can_continue(&self) -> bool {
438 false
444 }
445}
446
447impl From<io::Error> for DecodeError {
448 fn from(source: io::Error) -> Self {
449 DecodeError::IO { source }
450 }
451}
452
453#[derive(Debug, Clone, Copy)]
454enum LogstashProtocolVersion {
455 V1, V2, }
458
459impl From<LogstashProtocolVersion> for u8 {
460 fn from(frame_type: LogstashProtocolVersion) -> u8 {
461 use LogstashProtocolVersion::*;
462
463 match frame_type {
464 V1 => b'1',
465 V2 => b'2',
466 }
467 }
468}
469
470impl TryFrom<u8> for LogstashProtocolVersion {
471 type Error = DecodeError;
472
473 fn try_from(frame_type: u8) -> Result<LogstashProtocolVersion, DecodeError> {
474 use LogstashProtocolVersion::*;
475
476 match frame_type {
477 b'1' => Ok(V1),
478 b'2' => Ok(V2),
479 version => Err(DecodeError::UnknownProtocolVersion {
480 version: version as char,
481 }),
482 }
483 }
484}
485
486#[derive(Debug, Clone, Copy)]
487enum LogstashFrameType {
488 Ack, WindowSize, Data, Json, Compressed, }
494
495impl From<LogstashFrameType> for u8 {
496 fn from(frame_type: LogstashFrameType) -> u8 {
497 use LogstashFrameType::*;
498
499 match frame_type {
500 Ack => b'A',
501 WindowSize => b'W',
502 Data => b'D',
503 Json => b'J',
504 Compressed => b'C',
505 }
506 }
507}
508
509impl TryFrom<u8> for LogstashFrameType {
510 type Error = DecodeError;
511
512 fn try_from(frame_type: u8) -> Result<LogstashFrameType, DecodeError> {
513 use LogstashFrameType::*;
514
515 match frame_type {
516 b'A' => Ok(Ack),
517 b'W' => Ok(WindowSize),
518 b'D' => Ok(Data),
519 b'J' => Ok(Json),
520 b'C' => Ok(Compressed),
521 frame_type => Err(DecodeError::UnknownFrameType {
522 frame_type: frame_type as char,
523 }),
524 }
525 }
526}
527
528#[derive(Debug)]
530struct LogstashEventFrame {
531 protocol: LogstashProtocolVersion,
532 sequence_number: u32,
533 fields: BTreeMap<KeyString, serde_json::Value>,
534 window_end: bool,
537}
538
539impl Decoder for LogstashDecoder {
542 type Item = (LogstashEventFrame, usize);
543 type Error = DecodeError;
544
545 fn decode(&mut self, src: &mut BytesMut) -> Result<Option<Self::Item>, Self::Error> {
546 loop {
553 self.state = match self.state {
557 LogstashDecoderReadState::PendingDecompressed {
560 ref mut buf,
561 ref mut decoder,
562 } => match decoder.decode(buf)? {
563 Some(frame) => return Ok(Some(frame)),
564 None => {
565 self.window_events_remaining = decoder.window_events_remaining;
568 LogstashDecoderReadState::ReadProtocol
569 }
570 },
571 LogstashDecoderReadState::ReadProtocol => {
572 if src.remaining() < 1 {
573 return Ok(None);
574 }
575
576 use LogstashProtocolVersion::*;
577
578 match LogstashProtocolVersion::try_from(src.get_u8())? {
579 V1 => LogstashDecoderReadState::ReadType(V1),
580 V2 => LogstashDecoderReadState::ReadType(V2),
581 }
582 }
583 LogstashDecoderReadState::ReadType(protocol) => {
584 if src.remaining() < 1 {
585 return Ok(None);
586 }
587
588 use LogstashFrameType::*;
589
590 match LogstashFrameType::try_from(src.get_u8())? {
591 WindowSize => LogstashDecoderReadState::ReadFrame(protocol, WindowSize),
592 Data => LogstashDecoderReadState::ReadFrame(protocol, Data),
593 Json => LogstashDecoderReadState::ReadFrame(protocol, Json),
594 Compressed => LogstashDecoderReadState::ReadFrame(protocol, Compressed),
595 Ack => LogstashDecoderReadState::ReadFrame(protocol, Ack),
596 }
597 }
598 LogstashDecoderReadState::ReadFrame(_protocol, LogstashFrameType::WindowSize) => {
620 if let Some(remaining) = self.window_events_remaining {
630 return Err(DecodeError::PrematureWindowSize {
631 remaining: remaining.get(),
632 });
633 }
634
635 if src.remaining() < 4 {
636 return Ok(None);
637 }
638
639 let window_size = src.get_u32() as usize;
640 self.window_events_remaining = NonZeroUsize::new(window_size);
641
642 LogstashDecoderReadState::ReadProtocol
643 }
644 LogstashDecoderReadState::ReadFrame(_protocol, LogstashFrameType::Ack) => {
648 if src.remaining() < 4 {
649 return Ok(None);
650 }
651
652 let _sequence_number = src.get_u32();
653
654 LogstashDecoderReadState::ReadProtocol
655 }
656 LogstashDecoderReadState::ReadFrame(protocol, LogstashFrameType::Data) => {
658 let Some((mut frame, byte_size)) = decode_data_frame(protocol, src) else {
659 if src.remaining() > self.max_frame_size {
663 return Err(DecodeError::FrameTooLarge {
664 size: src.remaining(),
665 max: self.max_frame_size,
666 });
667 }
668 return Ok(None);
669 };
670 self.annotate_frame(&mut frame);
671
672 self.state = LogstashDecoderReadState::ReadProtocol;
673 return Ok(Some((frame, byte_size)));
674 }
675 LogstashDecoderReadState::ReadFrame(protocol, LogstashFrameType::Json) => {
677 let Some((mut frame, byte_size)) =
678 decode_json_frame(protocol, src, self.max_frame_size)?
679 else {
680 return Ok(None);
681 };
682 self.annotate_frame(&mut frame);
683
684 self.state = LogstashDecoderReadState::ReadProtocol;
685 return Ok(Some((frame, byte_size)));
686 }
687 LogstashDecoderReadState::ReadFrame(_protocol, LogstashFrameType::Compressed) => {
695 if self.nested {
696 return Err(DecodeError::NestedCompressedFrame);
697 }
698
699 let Some(buf) = decode_compressed_frame(src)? else {
700 return Ok(None);
701 };
702
703 LogstashDecoderReadState::PendingDecompressed {
704 buf,
705 decoder: Box::new(LogstashDecoder::new_nested(
706 self.window_events_remaining,
707 )),
708 }
709 }
710 };
711 }
712 }
713}
714
715fn decode_data_frame(
717 protocol: LogstashProtocolVersion,
718 src: &mut BytesMut,
719) -> Option<(LogstashEventFrame, usize)> {
720 let mut rest = src.as_ref();
721
722 if rest.remaining() < 8 {
723 return None;
724 }
725 let sequence_number = rest.get_u32();
726 let pair_count = rest.get_u32();
727 if pair_count == 0 {
728 return None; }
730
731 let mut fields = BTreeMap::<KeyString, serde_json::Value>::new();
732 for _ in 0..pair_count {
733 let (key, value, right) = decode_pair(rest)?;
734 rest = right;
735
736 fields.insert(
737 String::from_utf8_lossy(key).into(),
738 String::from_utf8_lossy(value).into(),
739 );
740 }
741
742 let byte_size = bytes_remaining(src, rest);
743 src.advance(byte_size);
744
745 Some((
746 LogstashEventFrame {
747 protocol,
748 sequence_number,
749 fields,
750 window_end: false,
751 },
752 byte_size,
753 ))
754}
755
756fn decode_pair(mut rest: &[u8]) -> Option<(&[u8], &[u8], &[u8])> {
757 if rest.remaining() < 4 {
758 return None;
759 }
760 let key_length = rest.get_u32() as usize;
761
762 if rest.remaining() < key_length {
763 return None;
764 }
765 let (key, right) = rest.split_at(key_length);
766 rest = right;
767
768 if rest.remaining() < 4 {
769 return None;
770 }
771 let value_length = rest.get_u32() as usize;
772 if rest.remaining() < value_length {
773 return None;
774 }
775 let (value, right) = rest.split_at(value_length);
776 Some((key, value, right))
777}
778
779fn decode_json_frame(
780 protocol: LogstashProtocolVersion,
781 src: &mut BytesMut,
782 max_frame_size: usize,
783) -> Result<Option<(LogstashEventFrame, usize)>, DecodeError> {
784 let mut rest = src.as_ref();
785
786 if rest.remaining() < 8 {
787 return Ok(None);
788 }
789 let sequence_number = rest.get_u32();
790 let payload_size = rest.get_u32() as usize;
791
792 if payload_size > max_frame_size {
796 return Err(DecodeError::FrameTooLarge {
797 size: payload_size,
798 max: max_frame_size,
799 });
800 }
801
802 if rest.remaining() < payload_size {
803 return Ok(None);
804 }
805
806 let (slice, right) = rest.split_at(payload_size);
807 rest = right;
808
809 let fields: BTreeMap<KeyString, serde_json::Value> =
810 serde_json::from_slice(slice).context(JsonFrameFailedDecodeSnafu {})?;
811
812 let byte_size = bytes_remaining(src, rest);
813 src.advance(byte_size);
814
815 Ok(Some((
816 LogstashEventFrame {
817 protocol,
818 sequence_number,
819 fields,
820 window_end: false,
821 },
822 byte_size,
823 )))
824}
825
826fn decode_compressed_frame(src: &mut BytesMut) -> Result<Option<BytesMut>, DecodeError> {
828 let mut rest = src.as_ref();
829
830 if rest.remaining() < 4 {
831 return Ok(None);
832 }
833 let payload_size = rest.get_u32() as usize;
834 let limit = max_decompressed_size_bytes();
835
836 let compressed_limit = max_zlib_compressed_frame_size_bytes();
841 if payload_size > compressed_limit {
842 return Err(DecodeError::DecompressionFailed {
843 source: io::Error::other(format!(
844 "compressed frame payload size {payload_size} exceeds limit of {compressed_limit} bytes"
845 )),
846 });
847 }
848
849 if rest.remaining() < payload_size {
850 return Ok(None);
851 }
852
853 let (slice, right) = rest.split_at(payload_size);
854 rest = right;
855
856 let res = CappedDecoder::zlib_with_limit(io::Cursor::new(slice), limit)
857 .decompress()
858 .map(|v| BytesMut::from(v.as_slice()))
859 .context(DecompressionFailedSnafu);
860
861 let byte_size = bytes_remaining(src, rest);
862 src.advance(byte_size);
863
864 Ok(Some(res?))
865}
866
867fn bytes_remaining(src: &BytesMut, rest: &[u8]) -> usize {
868 let remaining = rest.remaining();
869 src.remaining() - remaining
870}
871
872impl From<LogstashEventFrame> for Event {
873 fn from(frame: LogstashEventFrame) -> Self {
874 Event::Log(LogEvent::from(
875 frame
876 .fields
877 .into_iter()
878 .map(|(key, value)| (key, Value::from(value)))
879 .collect::<BTreeMap<_, _>>(),
880 ))
881 }
882}
883
884impl From<LogstashEventFrame> for SmallVec<[Event; 1]> {
885 fn from(frame: LogstashEventFrame) -> Self {
886 smallvec![frame.into()]
887 }
888}
889
890#[cfg(test)]
891mod test {
892 use std::io::Write;
893
894 use bytes::BufMut;
895 use flate2::{Compression, write::ZlibEncoder};
896 use futures::{Stream, StreamExt, stream};
897 use rand::{RngExt, rng};
898 use tokio::io::{AsyncReadExt, AsyncWriteExt};
899 use vector_lib::codecs::ReadyFrames;
900 use vector_lib::lookup::OwnedTargetPath;
901 use vrl::event_path;
902 use vrl::value::kind::Collection;
903
904 use super::*;
905 use crate::{
906 SourceSender,
907 event::EventStatus,
908 test_util::{
909 addr::next_addr,
910 components::{SOCKET_PUSH_SOURCE_TAGS, assert_source_compliance},
911 spawn_collect_n, wait_for_tcp,
912 },
913 };
914
915 #[test]
916 fn generate_config() {
917 crate::test_util::test_generate_config::<LogstashConfig>();
918 }
919
920 #[tokio::test]
921 async fn test_delivered() {
922 test_protocol(EventStatus::Delivered, true).await;
923 }
924
925 #[tokio::test]
926 async fn test_failed() {
927 test_protocol(EventStatus::Rejected, false).await;
928 }
929
930 async fn start_logstash(
931 status: EventStatus,
932 ) -> (SocketAddr, impl Stream<Item = Event> + Unpin) {
933 let (sender, recv) = SourceSender::new_test_finalize(status);
934 let (_guard, address) = next_addr();
935 let source = LogstashConfig {
936 address: address.into(),
937 tls: None,
938 permit_origin: None,
939 keepalive: None,
940 receive_buffer_bytes: None,
941 acknowledgements: true.into(),
942 connection_limit: None,
943 tls_handshake_timeout_secs: None,
944 log_namespace: None,
945 }
946 .build(SourceContext::new_test(sender, None))
947 .await
948 .unwrap();
949 tokio::spawn(source);
950 wait_for_tcp(address).await;
951 (address, recv)
952 }
953
954 async fn test_protocol(status: EventStatus, sends_ack: bool) {
955 let events = assert_source_compliance(&SOCKET_PUSH_SOURCE_TAGS, async {
956 let (address, recv) = start_logstash(status).await;
957 spawn_collect_n(
958 send_req(address, &[("message", "Hello, world!")], sends_ack),
959 recv,
960 1,
961 )
962 .await
963 })
964 .await;
965
966 assert_eq!(events.len(), 1);
967 let log = events[0].as_log();
968 assert_eq!(
969 log.get(event_path!("message")).unwrap().to_string_lossy(),
970 "Hello, world!".to_string()
971 );
972 assert_eq!(
973 log.get(event_path!("source_type"))
974 .unwrap()
975 .to_string_lossy(),
976 "logstash".to_string()
977 );
978 assert!(log.get(event_path!("host")).is_some());
979 assert!(log.get(event_path!("timestamp")).is_some());
980 }
981
982 fn push_req(req: &mut BytesMut, seq: u32, pairs: &[(&str, &str)]) {
983 req.put_u8(b'2');
984 req.put_u8(b'D');
985 req.put_u32(seq);
986 req.put_u32(pairs.len() as u32);
987 for (key, value) in pairs {
988 req.put_u32(key.len() as u32);
989 req.put(key.as_bytes());
990 req.put_u32(value.len() as u32);
991 req.put(value.as_bytes());
992 }
993 }
994
995 fn encode_req(seq: u32, pairs: &[(&str, &str)]) -> Bytes {
996 let mut req = BytesMut::new();
997 push_req(&mut req, seq, pairs);
998 req.into()
999 }
1000
1001 fn push_window_size(req: &mut BytesMut, size: u32) {
1002 req.put_u8(b'2');
1003 req.put_u8(b'W');
1004 req.put_u32(size);
1005 }
1006
1007 fn push_compressed(req: &mut BytesMut, inner: &[u8]) {
1008 let mut encoder = ZlibEncoder::new(Vec::new(), Compression::default());
1009 encoder.write_all(inner).unwrap();
1010 let compressed = encoder.finish().unwrap();
1011
1012 req.put_u8(b'2');
1013 req.put_u8(b'C');
1014 req.put_u32(compressed.len() as u32);
1015 req.put(compressed.as_slice());
1016 }
1017
1018 fn decode_frames(mut src: BytesMut) -> Vec<(LogstashEventFrame, usize)> {
1019 let mut decoder = LogstashDecoder::new();
1020 let mut frames = Vec::new();
1021
1022 while let Some(frame) = decoder.decode(&mut src).unwrap() {
1023 frames.push(frame);
1024 }
1025
1026 assert_eq!(src.len(), 0);
1027 frames
1028 }
1029
1030 fn decode_acknowledgements(mut ack: Bytes) -> Vec<u32> {
1031 let mut acknowledgements = Vec::new();
1032
1033 while !ack.is_empty() {
1034 assert!(
1035 ack.len() >= 6,
1036 "ack stream ended with {} trailing bytes",
1037 ack.len()
1038 );
1039 assert_eq!(ack.get_u8(), b'2');
1040 assert_eq!(ack.get_u8(), b'A');
1041 acknowledgements.push(ack.get_u32());
1042 }
1043
1044 acknowledgements
1045 }
1046
1047 fn decoded_sequence_numbers(decoded: &[(LogstashEventFrame, usize)]) -> Vec<u32> {
1048 decoded
1049 .iter()
1050 .map(|(frame, _)| frame.sequence_number)
1051 .collect::<Vec<_>>()
1052 }
1053
1054 fn assert_decoded_sequences(
1055 decoded: &[(LogstashEventFrame, usize)],
1056 expected_sequences: &[u32],
1057 ) {
1058 assert_eq!(decoded_sequence_numbers(decoded), expected_sequences);
1059 }
1060
1061 async fn assert_acknowledgements_for_ready_frames(
1062 decoded: Vec<(LogstashEventFrame, usize)>,
1063 expected_sequences: &[u32],
1064 expected_acknowledgements: &[u32],
1065 ) {
1066 assert_decoded_sequences(&decoded, expected_sequences);
1067
1068 let stream = stream::iter(decoded.into_iter().map(Ok::<_, DecodeError>));
1069 let mut ready = ReadyFrames::with_capacity(stream, 16);
1070 let (frames, _) = ready.next().await.unwrap().unwrap();
1071
1072 let acknowledgements = LogstashAcker::new(&frames)
1076 .build_ack(TcpSourceAck::Ack)
1077 .map_or_else(Vec::new, decode_acknowledgements);
1078
1079 assert!(ready.next().await.is_none());
1080 assert_eq!(acknowledgements, expected_acknowledgements);
1081 }
1082
1083 fn decode_frames_and_assert_sequences(
1084 src: BytesMut,
1085 expected_sequences: &[u32],
1086 ) -> Vec<(LogstashEventFrame, usize)> {
1087 let decoded = decode_frames(src);
1088 assert_decoded_sequences(&decoded, expected_sequences);
1089 decoded
1090 }
1091
1092 fn decode_frames_with_decoder(
1093 decoder: &mut LogstashDecoder,
1094 mut src: BytesMut,
1095 ) -> Vec<(LogstashEventFrame, usize)> {
1096 let mut frames = Vec::new();
1097
1098 while let Some(frame) = decoder.decode(&mut src).unwrap() {
1099 frames.push(frame);
1100 }
1101
1102 assert_eq!(src.len(), 0);
1103 frames
1104 }
1105
1106 fn decode_frames_with_decoder_and_assert_sequences(
1107 decoder: &mut LogstashDecoder,
1108 src: BytesMut,
1109 expected_sequences: &[u32],
1110 ) -> Vec<(LogstashEventFrame, usize)> {
1111 let decoded = decode_frames_with_decoder(decoder, src);
1112 assert_decoded_sequences(&decoded, expected_sequences);
1113 decoded
1114 }
1115
1116 #[test]
1117 fn v1_decoder_does_not_panic() {
1118 let seq = rng().random_range(1..u32::MAX);
1119 let req = encode_req(seq, &[("message", "Hello, World!")]);
1120 for i in 0..req.len() - 1 {
1121 assert!(
1122 decode_data_frame(LogstashProtocolVersion::V1, &mut BytesMut::from(&req[..i]))
1123 .is_none()
1124 );
1125 }
1126 }
1127
1128 #[test]
1135 fn malformed_json_frame_is_a_fatal_decode_error() {
1136 let mut decoder = LogstashDecoder::new();
1137 let mut src = BytesMut::new();
1138 src.put_u8(b'2');
1139 src.put_u8(b'J');
1140 src.put_u32(1); let bad = b"{ not valid json ";
1142 src.put_u32(bad.len() as u32); src.put(&bad[..]);
1144
1145 let err = decoder.decode(&mut src).unwrap_err();
1146 assert!(matches!(err, DecodeError::JsonFrameFailedDecode { .. }));
1147 assert!(
1148 !err.can_continue(),
1149 "a malformed JSON frame must be fatal so the connection closes",
1150 );
1151 }
1152
1153 #[test]
1154 fn malformed_compressed_frame_is_a_fatal_decode_error() {
1155 let mut decoder = LogstashDecoder::new();
1156 let mut src = BytesMut::new();
1157 src.put_u8(b'2');
1158 src.put_u8(b'C');
1159 let garbage = b"this is not a zlib stream";
1160 src.put_u32(garbage.len() as u32); src.put(&garbage[..]);
1162
1163 let err = decoder.decode(&mut src).unwrap_err();
1164 assert!(matches!(err, DecodeError::DecompressionFailed { .. }));
1165 assert!(!err.can_continue());
1166 }
1167
1168 #[test]
1169 fn premature_window_size_frame_is_a_fatal_decode_error() {
1170 let mut decoder = LogstashDecoder::new();
1183 let mut src = BytesMut::new();
1184 push_window_size(&mut src, 2);
1185 push_req(&mut src, 1, &[("message", "only one of two")]);
1186 push_window_size(&mut src, 5); assert!(decoder.decode(&mut src).unwrap().is_some());
1190
1191 let err = decoder.decode(&mut src).unwrap_err();
1194 assert!(matches!(err, DecodeError::PrematureWindowSize { .. }));
1195 assert!(
1196 !err.can_continue(),
1197 "a premature WindowSize must be fatal so the connection closes",
1198 );
1199 }
1200
1201 #[test]
1202 fn premature_window_size_inside_compressed_payload_is_fatal() {
1203 let mut inner = BytesMut::new();
1207 push_window_size(&mut inner, 2);
1208 push_req(&mut inner, 1, &[("message", "only one of two")]);
1209 push_window_size(&mut inner, 5); let mut req = BytesMut::new();
1212 push_compressed(&mut req, &inner);
1213
1214 let mut decoder = LogstashDecoder::new();
1215 assert!(decoder.decode(&mut req).unwrap().is_some());
1217 let err = decoder.decode(&mut req).unwrap_err();
1219 assert!(matches!(err, DecodeError::PrematureWindowSize { .. }));
1220 assert!(
1221 !err.can_continue(),
1222 "a premature WindowSize inside a compressed frame must be fatal",
1223 );
1224 }
1225
1226 #[test]
1227 fn compressed_frame_expansion_is_incremental() {
1228 let mut inner = BytesMut::new();
1231 for seq in 1..=1000 {
1232 push_req(&mut inner, seq, &[("m", "")]);
1233 }
1234
1235 let mut req = BytesMut::new();
1236 push_compressed(&mut req, &inner);
1237
1238 let mut decoder = LogstashDecoder::new();
1239 let first = decoder.decode(&mut req).unwrap().unwrap();
1240 assert_eq!(first.0.sequence_number, 1);
1241
1242 assert!(
1245 matches!(
1246 decoder.state,
1247 LogstashDecoderReadState::PendingDecompressed { .. }
1248 ),
1249 "compressed expansion must be incremental, got {:?}",
1250 decoder.state
1251 );
1252
1253 let second = decoder.decode(&mut req).unwrap().unwrap();
1255 assert_eq!(second.0.sequence_number, 2);
1256 }
1257
1258 #[test]
1259 fn nested_compressed_frame_is_a_fatal_decode_error() {
1260 let mut inner = BytesMut::new();
1261 push_req(&mut inner, 1, &[("message", "should never be reached")]);
1262
1263 let mut middle = BytesMut::new();
1264 push_compressed(&mut middle, &inner);
1265
1266 let mut req = BytesMut::new();
1267 push_compressed(&mut req, &middle);
1268
1269 let mut decoder = LogstashDecoder::new();
1270 let err = decoder.decode(&mut req).unwrap_err();
1271 assert!(matches!(err, DecodeError::NestedCompressedFrame));
1272 assert!(!err.can_continue());
1273 }
1274
1275 #[test]
1276 fn oversized_frames_are_rejected() {
1277 let mut json = BytesMut::new();
1281 json.put_u8(b'2');
1282 json.put_u8(b'J');
1283 json.put_u32(1); json.put_u32(100); let mut data = BytesMut::new();
1287 data.put_u8(b'2');
1288 data.put_u8(b'D');
1289 data.put_u32(1); data.put_u32(u32::MAX); data.put(&[0u8; 8][..]); for (frame_type, mut src, expected_size) in [("json", json, 100), ("data", data, 16)] {
1294 let mut decoder = LogstashDecoder::new();
1295 decoder.max_frame_size = 8;
1296
1297 let err = decoder.decode(&mut src).unwrap_err();
1298 assert!(
1299 matches!(err, DecodeError::FrameTooLarge { size, max: 8 } if size == expected_size),
1300 "{frame_type} frame: unexpected error {err:?}",
1301 );
1302 assert!(
1303 !err.can_continue(),
1304 "{frame_type} frame: an oversized frame must be fatal so the connection closes",
1305 );
1306 }
1307 }
1308
1309 #[test]
1310 fn frames_within_the_cap_still_decode() {
1311 let mut decoder = LogstashDecoder::new();
1312 decoder.max_frame_size = 100;
1313
1314 let mut req = BytesMut::new();
1315 push_req(&mut req, 1, &[("message", "hello")]);
1316
1317 let decoded = decode_frames_with_decoder(&mut decoder, req);
1318 assert_decoded_sequences(&decoded, &[1]);
1319 }
1320
1321 #[test]
1322 fn fragmented_input_is_assembled_across_decode_calls() {
1323 let mut decoder = LogstashDecoder::new();
1326 let full = encode_req(7, &[("message", "hello")]);
1327
1328 let mut src = BytesMut::new();
1329 let mut decoded = Vec::new();
1330 for byte in full.iter() {
1331 src.put_u8(*byte);
1332 if let Some(frame) = decoder.decode(&mut src).unwrap() {
1333 decoded.push(frame);
1334 }
1335 }
1336 assert_decoded_sequences(&decoded, &[7]);
1337 }
1338
1339 #[tokio::test]
1340 async fn malformed_frame_closes_connection_without_ack() {
1341 let (address, _recv) = start_logstash(EventStatus::Delivered).await;
1342
1343 let mut socket = tokio::net::TcpStream::connect(address).await.unwrap();
1344
1345 let mut req = BytesMut::new();
1347 req.put_u8(b'2');
1348 req.put_u8(b'J');
1349 req.put_u32(1); let bad = b"{ not valid json ";
1351 req.put_u32(bad.len() as u32); req.put(&bad[..]);
1353 socket.write_all(&req).await.unwrap();
1354
1355 let mut output = BytesMut::new();
1358 let result = socket.read_buf(&mut output).await;
1359 assert!(
1360 matches!(result, Ok(0)) || result.is_err(),
1361 "expected the connection to close; read returned {result:?} with {output:?}",
1362 );
1363 assert!(
1364 output.is_empty(),
1365 "no ACK should be sent for a malformed frame, got {output:?}",
1366 );
1367 }
1368
1369 #[tokio::test]
1370 async fn distinct_windows_do_not_share_an_ack_domain() {
1371 let mut req = BytesMut::new();
1372 push_window_size(&mut req, 1);
1373 push_req(&mut req, 1, &[("message", "first window")]);
1374 push_window_size(&mut req, 2);
1375 push_req(&mut req, 1, &[("message", "second window first")]);
1376 push_req(&mut req, 2, &[("message", "second window second")]);
1377
1378 let decoded = decode_frames_and_assert_sequences(req, &[1, 1, 2]);
1379 assert_acknowledgements_for_ready_frames(decoded, &[1, 1, 2], &[1, 2]).await;
1380 }
1381
1382 #[tokio::test]
1383 async fn distinct_windows_with_monotonic_sequences_ack_the_first_window() {
1384 let mut req = BytesMut::new();
1385 push_window_size(&mut req, 2);
1386 push_req(&mut req, 1, &[("message", "first window first")]);
1387 push_req(&mut req, 2, &[("message", "first window second")]);
1388 push_window_size(&mut req, 2);
1389 push_req(&mut req, 3, &[("message", "second window first")]);
1390 push_req(&mut req, 4, &[("message", "second window second")]);
1391
1392 let decoded = decode_frames_and_assert_sequences(req, &[1, 2, 3, 4]);
1393 assert_acknowledgements_for_ready_frames(decoded, &[1, 2, 3, 4], &[2, 4]).await;
1394 }
1395
1396 #[tokio::test]
1397 async fn incomplete_window_is_not_acked() {
1398 let mut req = BytesMut::new();
1406 push_window_size(&mut req, 4);
1407 push_req(&mut req, 1, &[("message", "only event in partial window")]);
1408
1409 let decoded = decode_frames_and_assert_sequences(req, &[1]);
1410 assert_acknowledgements_for_ready_frames(decoded, &[1], &[]).await;
1411 }
1412
1413 #[tokio::test]
1414 async fn window_split_across_compressed_frames_acks_once_on_completion() {
1415 let mut first = BytesMut::new();
1424 push_req(&mut first, 1, &[("message", "w4 first")]);
1425 push_req(&mut first, 2, &[("message", "w4 second")]);
1426
1427 let mut second = BytesMut::new();
1428 push_req(&mut second, 3, &[("message", "w4 third")]);
1429 push_req(&mut second, 4, &[("message", "w4 fourth")]);
1430
1431 let mut req = BytesMut::new();
1435 push_window_size(&mut req, 4);
1436 push_compressed(&mut req, &first);
1437 push_compressed(&mut req, &second);
1438
1439 let decoded = decode_frames_and_assert_sequences(req, &[1, 2, 3, 4]);
1440 assert_acknowledgements_for_ready_frames(decoded, &[1, 2, 3, 4], &[4]).await;
1441 }
1442
1443 #[tokio::test]
1444 async fn window_larger_than_ready_frames_capacity_in_one_compressed_frame_acks_once() {
1445 const WINDOW: u32 = 5;
1446 const CAPACITY: usize = 2;
1447
1448 let mut inner = BytesMut::new();
1449 for seq in 1..=WINDOW {
1450 push_req(&mut inner, seq, &[("message", "event in oversized window")]);
1451 }
1452
1453 const SMALL_WINDOW: u32 = 2;
1460 let mut small_inner = BytesMut::new();
1461 for seq in 1..=SMALL_WINDOW {
1462 push_req(
1463 &mut small_inner,
1464 seq,
1465 &[("message", "event in small window")],
1466 );
1467 }
1468
1469 let mut req = BytesMut::new();
1472 push_window_size(&mut req, WINDOW);
1473 push_compressed(&mut req, &inner);
1474 push_window_size(&mut req, SMALL_WINDOW);
1475 push_compressed(&mut req, &small_inner);
1476
1477 let decoded = decode_frames_and_assert_sequences(req, &[1, 2, 3, 4, 5, 1, 2]);
1478
1479 let stream = stream::iter(decoded.into_iter().map(Ok::<_, DecodeError>));
1480 let mut ready = ReadyFrames::with_capacity(stream, CAPACITY);
1481 let mut acknowledgements = Vec::new();
1482
1483 while let Some(result) = ready.next().await {
1484 let (frames, _byte_size) = result.unwrap();
1485 let acks = LogstashAcker::new(&frames)
1486 .build_ack(TcpSourceAck::Ack)
1487 .map_or_else(Vec::new, decode_acknowledgements);
1488 acknowledgements.push(acks);
1489 }
1490
1491 assert_eq!(acknowledgements, vec![vec![], vec![], vec![5], vec![2]]);
1499
1500 let all_acks: Vec<u32> = acknowledgements.into_iter().flatten().collect();
1502 assert_eq!(
1503 all_acks,
1504 vec![5, 2],
1505 "each window must be ACKed exactly once"
1506 );
1507 }
1508
1509 #[tokio::test]
1510 async fn complete_window_then_incomplete_window_acks_only_the_complete_one() {
1511 let mut req = BytesMut::new();
1517 push_window_size(&mut req, 2);
1518 push_req(&mut req, 1, &[("message", "complete window first")]);
1519 push_req(&mut req, 2, &[("message", "complete window second")]);
1520 push_window_size(&mut req, 1000);
1521 push_req(&mut req, 1, &[("message", "partial window first")]);
1522 push_req(&mut req, 2, &[("message", "partial window second")]);
1523 push_req(&mut req, 3, &[("message", "partial window third")]);
1524
1525 let decoded = decode_frames_and_assert_sequences(req, &[1, 2, 1, 2, 3]);
1526 assert_acknowledgements_for_ready_frames(decoded, &[1, 2, 1, 2, 3], &[2]).await;
1527 }
1528
1529 #[tokio::test]
1530 async fn compressed_frames_preserve_inner_window_boundaries() {
1531 let mut inner = BytesMut::new();
1532 push_window_size(&mut inner, 2);
1533 push_req(&mut inner, 1, &[("message", "compressed first")]);
1534 push_req(&mut inner, 2, &[("message", "compressed second")]);
1535
1536 let mut req = BytesMut::new();
1537 push_compressed(&mut req, &inner);
1538
1539 let decoded = decode_frames_and_assert_sequences(req, &[1, 2]);
1540 assert_acknowledgements_for_ready_frames(decoded, &[1, 2], &[2]).await;
1541 }
1542
1543 #[tokio::test]
1544 async fn single_window_split_across_ready_frames_acks_only_on_completion() {
1545 let mut req = BytesMut::new();
1550 push_window_size(&mut req, 4);
1551 push_req(&mut req, 1, &[("message", "first")]);
1552 push_req(&mut req, 2, &[("message", "second")]);
1553 push_req(&mut req, 3, &[("message", "third")]);
1554 push_req(&mut req, 4, &[("message", "fourth")]);
1555
1556 let decoded = decode_frames_and_assert_sequences(req, &[1, 2, 3, 4]);
1557
1558 let stream = stream::iter(decoded.into_iter().map(Ok::<_, DecodeError>));
1559 let mut ready = ReadyFrames::with_capacity(stream, 2);
1560 let mut acknowledgements = Vec::new();
1561
1562 while let Some(result) = ready.next().await {
1563 let (frames, _byte_size) = result.unwrap();
1564 let acks = LogstashAcker::new(&frames)
1565 .build_ack(TcpSourceAck::Ack)
1566 .map_or_else(Vec::new, decode_acknowledgements);
1567 acknowledgements.push(acks);
1568 }
1569
1570 assert_eq!(acknowledgements, vec![vec![], vec![4]]);
1573 }
1574
1575 #[tokio::test]
1576 async fn fresh_window_after_completed_window_is_accepted() {
1577 let mut decoder = LogstashDecoder::new();
1585
1586 let mut first_batch = BytesMut::new();
1587 push_window_size(&mut first_batch, 1);
1588 push_req(&mut first_batch, 1, &[("message", "first window")]);
1589 let decoded =
1590 decode_frames_with_decoder_and_assert_sequences(&mut decoder, first_batch, &[1]);
1591 assert_acknowledgements_for_ready_frames(decoded, &[1], &[1]).await;
1592
1593 let mut second_batch = BytesMut::new();
1594 push_window_size(&mut second_batch, 1);
1595 push_req(
1596 &mut second_batch,
1597 1,
1598 &[("message", "fresh window after completion")],
1599 );
1600 let decoded =
1601 decode_frames_with_decoder_and_assert_sequences(&mut decoder, second_batch, &[1]);
1602 assert_acknowledgements_for_ready_frames(decoded, &[1], &[1]).await;
1603 }
1604
1605 async fn send_req(address: SocketAddr, pairs: &[(&str, &str)], sends_ack: bool) {
1606 let seq = rng().random_range(1..u32::MAX);
1607 let mut socket = tokio::net::TcpStream::connect(address).await.unwrap();
1608
1609 let req = encode_req(seq, pairs);
1610 socket.write_all(&req).await.unwrap();
1611
1612 let mut output = BytesMut::new();
1613 socket.read_buf(&mut output).await.unwrap();
1614
1615 if sends_ack {
1616 assert_eq!(output.get_u8(), b'2');
1617 assert_eq!(output.get_u8(), b'A');
1618 assert_eq!(output.get_u32(), seq);
1619 }
1620 assert_eq!(output.len(), 0);
1621 }
1622
1623 #[test]
1624 fn output_schema_definition_vector_namespace() {
1625 let config = LogstashConfig {
1626 log_namespace: Some(true),
1627 ..Default::default()
1628 };
1629
1630 let definitions = config
1631 .outputs(LogNamespace::Vector)
1632 .remove(0)
1633 .schema_definition(true);
1634
1635 let expected_definition =
1636 Definition::new_with_default_metadata(Kind::bytes(), [LogNamespace::Vector])
1637 .with_meaning(OwnedTargetPath::event_root(), "message")
1638 .with_metadata_field(
1639 &owned_value_path!("vector", "source_type"),
1640 Kind::bytes(),
1641 None,
1642 )
1643 .with_metadata_field(
1644 &owned_value_path!("vector", "ingest_timestamp"),
1645 Kind::timestamp(),
1646 None,
1647 )
1648 .with_metadata_field(
1649 &owned_value_path!(LogstashConfig::NAME, "timestamp"),
1650 Kind::timestamp().or_undefined(),
1651 Some("timestamp"),
1652 )
1653 .with_metadata_field(
1654 &owned_value_path!(LogstashConfig::NAME, "host"),
1655 Kind::bytes(),
1656 Some("host"),
1657 )
1658 .with_metadata_field(
1659 &owned_value_path!(LogstashConfig::NAME, "tls_client_metadata"),
1660 Kind::object(Collection::empty().with_unknown(Kind::bytes())).or_undefined(),
1661 None,
1662 );
1663
1664 assert_eq!(definitions, Some(expected_definition))
1665 }
1666
1667 #[test]
1668 fn output_schema_definition_legacy_namespace() {
1669 let config = LogstashConfig::default();
1670
1671 let definitions = config
1672 .outputs(LogNamespace::Legacy)
1673 .remove(0)
1674 .schema_definition(true);
1675
1676 let expected_definition = Definition::new_with_default_metadata(
1677 Kind::object(Collection::empty()),
1678 [LogNamespace::Legacy],
1679 )
1680 .with_event_field(
1681 &owned_value_path!("message"),
1682 Kind::bytes(),
1683 Some("message"),
1684 )
1685 .with_event_field(&owned_value_path!("source_type"), Kind::bytes(), None)
1686 .with_event_field(&owned_value_path!("timestamp"), Kind::timestamp(), None)
1687 .with_event_field(&owned_value_path!("host"), Kind::bytes(), Some("host"));
1688
1689 assert_eq!(definitions, Some(expected_definition))
1690 }
1691}
1692
1693#[cfg(all(test, feature = "logstash-integration-tests"))]
1694mod integration_tests {
1695 use std::time::Duration;
1696
1697 use futures::Stream;
1698 use tokio::time::timeout;
1699 use vrl::event_path;
1700
1701 use super::*;
1702 use crate::{
1703 SourceSender,
1704 config::SourceContext,
1705 event::EventStatus,
1706 test_util::{
1707 collect_n,
1708 components::{SOCKET_PUSH_SOURCE_TAGS, assert_source_compliance},
1709 wait_for_tcp,
1710 },
1711 tls::{TlsConfig, TlsEnableableConfig},
1712 };
1713
1714 fn heartbeat_address() -> String {
1715 std::env::var("HEARTBEAT_ADDRESS")
1716 .expect("Address of Beats Heartbeat service must be specified.")
1717 }
1718
1719 #[tokio::test]
1720 async fn beats_heartbeat() {
1721 let events = assert_source_compliance(&SOCKET_PUSH_SOURCE_TAGS, async {
1722 let out = source(heartbeat_address(), None).await;
1723
1724 timeout(Duration::from_secs(60), collect_n(out, 1))
1725 .await
1726 .unwrap()
1727 })
1728 .await;
1729
1730 assert!(!events.is_empty());
1731
1732 let log = events[0].as_log();
1733 assert_eq!(
1734 log.get(event_path!("@metadata", "beat")),
1735 Some(String::from("heartbeat").into()).as_ref()
1736 );
1737 assert_eq!(
1738 log.get(event_path!("summary", "up")),
1739 Some(1.into()).as_ref()
1740 );
1741 assert!(log.get(event_path!("timestamp")).is_some());
1742 assert!(log.get(event_path!("host")).is_some());
1743 }
1744
1745 fn logstash_address() -> String {
1746 std::env::var("LOGSTASH_ADDRESS")
1747 .expect("Listen address for `logstash` source must be specified.")
1748 }
1749
1750 #[tokio::test]
1751 async fn logstash() {
1752 let events = assert_source_compliance(&SOCKET_PUSH_SOURCE_TAGS, async {
1753 let out = source(
1754 logstash_address(),
1755 Some(TlsEnableableConfig {
1756 enabled: Some(true),
1757 options: TlsConfig {
1758 crt_file: Some(
1759 "tests/integration/shared/data/host.docker.internal.crt".into(),
1760 ),
1761 key_file: Some(
1762 "tests/integration/shared/data/host.docker.internal.key".into(),
1763 ),
1764 ..Default::default()
1765 },
1766 }),
1767 )
1768 .await;
1769
1770 timeout(Duration::from_secs(60), collect_n(out, 1))
1771 .await
1772 .unwrap()
1773 })
1774 .await;
1775
1776 assert!(!events.is_empty());
1777
1778 let log = events[0].as_log();
1779 assert!(
1780 log.get(event_path!("line"))
1781 .unwrap()
1782 .to_string_lossy()
1783 .contains("Hello World")
1784 );
1785 assert!(log.get(event_path!("host")).is_some());
1786 }
1787
1788 async fn source(
1789 address: String,
1790 tls: Option<TlsEnableableConfig>,
1791 ) -> impl Stream<Item = Event> + Unpin {
1792 let (sender, recv) = SourceSender::new_test_finalize(EventStatus::Delivered);
1793 let address: SocketAddr = address.parse().unwrap();
1794 let tls_config = TlsSourceConfig {
1795 client_metadata_key: None,
1796 tls_config: tls.unwrap_or_default(),
1797 };
1798 tokio::spawn(async move {
1799 LogstashConfig {
1800 address: address.into(),
1801 tls: Some(tls_config),
1802 keepalive: None,
1803 permit_origin: None,
1804 receive_buffer_bytes: None,
1805 acknowledgements: false.into(),
1806 connection_limit: None,
1807 tls_handshake_timeout_secs: None,
1808 log_namespace: None,
1809 }
1810 .build(SourceContext::new_test(sender, None))
1811 .await
1812 .unwrap()
1813 .await
1814 .unwrap()
1815 });
1816 wait_for_tcp(address).await;
1817 recv
1818 }
1819}