diff --git a/changelog.d/socket_tcp_disconnect_mode.feature.md b/changelog.d/socket_tcp_disconnect_mode.feature.md new file mode 100644 index 0000000000000..e255d1bd692f6 --- /dev/null +++ b/changelog.d/socket_tcp_disconnect_mode.feature.md @@ -0,0 +1,3 @@ +Added a `disconnect_mode` configuration option to the `socket` (TCP mode), `logstash`, `fluent`, `syslog` (TCP mode), and `statsd` (TCP mode) sources. This controls how Vector closes TCP connections on shutdown or when `max_connection_duration_secs` elapses. The `drain` mode (the default, and existing behaviour) maintains the existing graceful shutdown behaviour while the `abort` mode closes connections immediately without waiting for the client to acknowledge the shutdown. This is useful for clients that never read from the socket and therefore cannot detect a graceful shutdown. + +authors: tronboto diff --git a/src/sources/fluent/mod.rs b/src/sources/fluent/mod.rs index 8233d055aa7ec..c055d7357ee22 100644 --- a/src/sources/fluent/mod.rs +++ b/src/sources/fluent/mod.rs @@ -18,7 +18,7 @@ use vector_lib::{ use vrl::value::{Kind, Value, kind::Collection}; use super::util::decompression::{CappedDecoder, max_decompressed_size_bytes}; -use super::util::net::{SocketListenAddr, TcpSource, TcpSourceAck, TcpSourceAcker}; +use super::util::net::{DisconnectMode, SocketListenAddr, TcpSource, TcpSourceAck, TcpSourceAcker}; use crate::{ config::{ DataType, GenerateConfig, Resource, SourceAcknowledgementsConfig, SourceConfig, @@ -190,6 +190,10 @@ pub struct FluentTcpConfig { #[configurable(derived)] #[serde(default, deserialize_with = "bool_or_struct")] acknowledgements: SourceAcknowledgementsConfig, + + #[configurable(derived)] + #[serde(default)] + disconnect_mode: DisconnectMode, } impl FluentTcpConfig { @@ -216,6 +220,7 @@ impl FluentTcpConfig { tls_client_metadata_key, self.receive_buffer_bytes, None, + self.disconnect_mode, cx, self.acknowledgements, self.connection_limit, @@ -279,6 +284,7 @@ impl GenerateConfig for FluentConfig { receive_buffer_bytes: None, acknowledgements: Default::default(), connection_limit: Some(2), + disconnect_mode: DisconnectMode::Drain, }), log_namespace: None, }) @@ -1169,6 +1175,7 @@ mod tests { receive_buffer_bytes: None, acknowledgements: true.into(), connection_limit: None, + disconnect_mode: DisconnectMode::Drain, }), log_namespace: None, } @@ -1242,6 +1249,7 @@ mod tests { receive_buffer_bytes: None, acknowledgements: false.into(), connection_limit: None, + disconnect_mode: DisconnectMode::Drain, }), log_namespace: Some(true), }; @@ -1300,6 +1308,7 @@ mod tests { receive_buffer_bytes: None, acknowledgements: false.into(), connection_limit: None, + disconnect_mode: DisconnectMode::Drain, }), log_namespace: None, }; @@ -1529,6 +1538,7 @@ mod integration_tests { receive_buffer_bytes: None, acknowledgements: false.into(), connection_limit: None, + disconnect_mode: DisconnectMode::Drain, }), log_namespace: None, } diff --git a/src/sources/logstash.rs b/src/sources/logstash.rs index f0ceecb17fe70..11a0564b9cc1c 100644 --- a/src/sources/logstash.rs +++ b/src/sources/logstash.rs @@ -24,7 +24,7 @@ use vrl::value::{KeyString, Kind, kind::Collection}; use super::util::decompression::{ CappedDecoder, max_decompressed_size_bytes, max_zlib_compressed_frame_size_bytes, }; -use super::util::net::{SocketListenAddr, TcpSource, TcpSourceAck, TcpSourceAcker}; +use super::util::net::{DisconnectMode, SocketListenAddr, TcpSource, TcpSourceAck, TcpSourceAcker}; use crate::{ config::{ DataType, GenerateConfig, Resource, SourceAcknowledgementsConfig, SourceConfig, @@ -62,6 +62,10 @@ pub struct LogstashConfig { #[configurable(metadata(docs::type_unit = "connections"))] connection_limit: Option, + #[configurable(derived)] + #[serde(default)] + disconnect_mode: DisconnectMode, + #[configurable(derived)] #[serde(default, deserialize_with = "bool_or_struct")] acknowledgements: SourceAcknowledgementsConfig, @@ -125,6 +129,7 @@ impl Default for LogstashConfig { receive_buffer_bytes: None, acknowledgements: Default::default(), connection_limit: None, + disconnect_mode: DisconnectMode::Drain, log_namespace: None, } } @@ -164,6 +169,7 @@ impl SourceConfig for LogstashConfig { tls_client_metadata_key, self.receive_buffer_bytes, None, + self.disconnect_mode, cx, self.acknowledgements, self.connection_limit, @@ -902,6 +908,7 @@ mod test { receive_buffer_bytes: None, acknowledgements: true.into(), connection_limit: None, + disconnect_mode: DisconnectMode::Drain, log_namespace: None, } .build(SourceContext::new_test(sender, None)) @@ -1667,6 +1674,7 @@ mod integration_tests { receive_buffer_bytes: None, acknowledgements: false.into(), connection_limit: None, + disconnect_mode: DisconnectMode::Drain, log_namespace: None, } .build(SourceContext::new_test(sender, None)) diff --git a/src/sources/socket/mod.rs b/src/sources/socket/mod.rs index ca12b9f25a14d..a1114e40d259f 100644 --- a/src/sources/socket/mod.rs +++ b/src/sources/socket/mod.rs @@ -143,6 +143,7 @@ impl SourceConfig for SocketConfig { tls_client_metadata_key, config.receive_buffer_bytes(), config.max_connection_duration_secs(), + config.disconnect_mode(), cx, false.into(), config.connection_limit, @@ -374,7 +375,7 @@ mod test { event::{Event, LogEvent}, shutdown::{ShutdownSignal, SourceShutdownCoordinator}, sinks::util::tcp::TcpSinkConfig, - sources::util::net::SocketListenAddr, + sources::util::net::{DisconnectMode, SocketListenAddr}, test_util::{ addr::{PortGuard, next_addr, next_addr_any}, collect_n, collect_n_limited, @@ -930,6 +931,38 @@ mod test { } } + #[tokio::test] + async fn tcp_disconnect_mode_abort_on_shutdown() { + let source_id = ComponentKey::from("tcp_disconnect_mode_abort_on_shutdown"); + let (tx, _) = SourceSender::new_test(); + let (guard, addr) = next_addr(); + let (cx, mut shutdown) = SourceContext::new_shutdown(&source_id, tx); + + let mut source_config = TcpConfig::from_address(addr.into()); + source_config.set_disconnect_mode(DisconnectMode::Abort); + let source_task = SocketConfig::from(source_config).build(cx).await.unwrap(); + + drop(tokio::spawn(source_task)); + wait_for_tcp_and_release(guard, addr).await; + + let mut stream = TcpStream::connect(addr) + .await + .expect("stream should be able to connect"); + let mut buffer = [0u8; 10]; + + let deadline = Instant::now() + Duration::from_secs(10); + tokio::spawn(shutdown.shutdown_source(&source_id, deadline)); + + let read_result = tokio::time::timeout(Duration::from_secs(5), stream.read(&mut buffer)) + .await + .expect("timed out waiting for connection to close"); + + match read_result { + Err(e) => assert_eq!(e.kind(), std::io::ErrorKind::ConnectionReset), + Ok(n) => panic!("expected connection reset, got Ok({n})"), + } + } + //////// UDP TESTS //////// async fn send_lines_udp(to: SocketAddr, lines: impl IntoIterator) -> UdpSocket { send_lines_udp_from(bind_unused_udp(), to, lines) diff --git a/src/sources/socket/tcp.rs b/src/sources/socket/tcp.rs index 5e1873e6f9cde..ef69308d26837 100644 --- a/src/sources/socket/tcp.rs +++ b/src/sources/socket/tcp.rs @@ -16,7 +16,7 @@ use crate::{ codecs::Decoder, event::Event, serde::default_decoding, - sources::util::net::{SocketListenAddr, TcpNullAcker, TcpSource}, + sources::util::net::{DisconnectMode, SocketListenAddr, TcpNullAcker, TcpSource}, tcp::TcpKeepaliveConfig, tls::TlsSourceConfig, }; @@ -79,6 +79,10 @@ pub struct TcpConfig { #[configurable(metadata(docs::type_unit = "connections"))] pub connection_limit: Option, + #[configurable(derived)] + #[serde(default)] + disconnect_mode: DisconnectMode, + #[configurable(derived)] pub(super) framing: Option, @@ -115,6 +119,7 @@ impl TcpConfig { framing: None, decoding: default_decoding(), connection_limit: None, + disconnect_mode: DisconnectMode::Drain, log_namespace: None, } } @@ -159,11 +164,20 @@ impl TcpConfig { self.max_connection_duration_secs } + pub const fn disconnect_mode(&self) -> DisconnectMode { + self.disconnect_mode + } + pub const fn set_max_connection_duration_secs(&mut self, val: Option) -> &mut Self { self.max_connection_duration_secs = val; self } + pub const fn set_disconnect_mode(&mut self, val: DisconnectMode) -> &mut Self { + self.disconnect_mode = val; + self + } + pub const fn set_shutdown_timeout_secs(&mut self, val: u64) -> &mut Self { self.shutdown_timeout_secs = Duration::from_secs(val); self diff --git a/src/sources/statsd/mod.rs b/src/sources/statsd/mod.rs index 55650b6eaa0c1..1eb18e9e38ae8 100644 --- a/src/sources/statsd/mod.rs +++ b/src/sources/statsd/mod.rs @@ -21,7 +21,9 @@ use vector_lib::{ }; use self::parser::ParseError; -use super::util::net::{SocketListenAddr, TcpNullAcker, TcpSource, try_bind_udp_socket}; +use super::util::net::{ + DisconnectMode, SocketListenAddr, TcpNullAcker, TcpSource, try_bind_udp_socket, +}; use crate::{ SourceSender, codecs::Decoder, @@ -139,6 +141,10 @@ pub struct TcpConfig { #[configurable(metadata(docs::type_unit = "connections"))] connection_limit: Option, + #[configurable(derived)] + #[serde(default)] + disconnect_mode: DisconnectMode, + /// Whether or not to sanitize incoming statsd key names. When "true", keys are sanitized by: /// - "/" is replaced with "-" /// - All whitespace is replaced with "_" @@ -163,6 +169,7 @@ impl TcpConfig { shutdown_timeout_secs: default_shutdown_timeout_secs(), receive_buffer_bytes: None, connection_limit: None, + disconnect_mode: DisconnectMode::Drain, sanitize: default_sanitize(), convert_to: default_convert_to(), } @@ -223,6 +230,7 @@ impl SourceConfig for StatsdConfig { tls_client_metadata_key, config.receive_buffer_bytes, None, + config.disconnect_mode, cx, false.into(), config.connection_limit, diff --git a/src/sources/syslog.rs b/src/sources/syslog.rs index 912b7b716f67a..97e04c0858d15 100644 --- a/src/sources/syslog.rs +++ b/src/sources/syslog.rs @@ -36,7 +36,9 @@ use crate::{ }, net, shutdown::ShutdownSignal, - sources::util::net::{SocketListenAddr, TcpNullAcker, TcpSource, try_bind_udp_socket}, + sources::util::net::{ + DisconnectMode, SocketListenAddr, TcpNullAcker, TcpSource, try_bind_udp_socket, + }, tcp::TcpKeepaliveConfig, tls::{MaybeTlsSettings, TlsSourceConfig}, }; @@ -100,6 +102,10 @@ pub enum Mode { /// The maximum number of TCP connections that are allowed at any given time. connection_limit: Option, + + #[configurable(derived)] + #[serde(default)] + disconnect_mode: DisconnectMode, }, /// Listen on UDP. @@ -155,6 +161,7 @@ impl Default for SyslogConfig { tls: None, receive_buffer_bytes: None, connection_limit: None, + disconnect_mode: DisconnectMode::Drain, }, host_key: None, max_length: crate::serde::default_max_length(), @@ -188,6 +195,7 @@ impl SourceConfig for SyslogConfig { tls, receive_buffer_bytes, connection_limit, + disconnect_mode, } => { let source = SyslogTcpSource { max_length: self.max_length, @@ -210,6 +218,7 @@ impl SourceConfig for SyslogConfig { tls_client_metadata_key, receive_buffer_bytes, None, + disconnect_mode, cx, false.into(), connection_limit, @@ -1149,6 +1158,7 @@ mod test { tls: None, receive_buffer_bytes: None, connection_limit: None, + disconnect_mode: DisconnectMode::Drain, }); let key = ComponentKey::from("in"); @@ -1360,6 +1370,7 @@ mod test { tls: None, receive_buffer_bytes: None, connection_limit: None, + disconnect_mode: DisconnectMode::Drain, }); let key = ComponentKey::from("in"); diff --git a/src/sources/util/net/mod.rs b/src/sources/util/net/mod.rs index 4ae7920973b7a..7d0b1eb9e8f3f 100644 --- a/src/sources/util/net/mod.rs +++ b/src/sources/util/net/mod.rs @@ -10,8 +10,8 @@ use vector_lib::configurable::configurable_component; #[cfg(feature = "sources-utils-net-tcp")] pub use self::tcp::{ - MAX_IN_FLIGHT_EVENTS_TARGET, TcpNullAcker, TcpSource, TcpSourceAck, TcpSourceAcker, - request_limiter::RequestLimiter, try_bind_tcp_listener, + DisconnectMode, MAX_IN_FLIGHT_EVENTS_TARGET, TcpNullAcker, TcpSource, TcpSourceAck, + TcpSourceAcker, request_limiter::RequestLimiter, try_bind_tcp_listener, }; #[cfg(feature = "sources-utils-net-udp")] pub use self::udp::try_bind_udp_socket; diff --git a/src/sources/util/net/tcp/mod.rs b/src/sources/util/net/tcp/mod.rs index aebce3dc9b6ee..3f9a673b78eb9 100644 --- a/src/sources/util/net/tcp/mod.rs +++ b/src/sources/util/net/tcp/mod.rs @@ -20,6 +20,7 @@ use vector_lib::{ EstimatedJsonEncodedSizeOf, codecs::{ReadyFrames, StreamDecodingError, internal_events::DecoderFramingError}, config::{LegacyKey, LogNamespace, SourceAcknowledgementsConfig}, + configurable::configurable_component, event::{BatchNotifier, BatchStatus, Event}, finalization::AddBatchNotifier, lookup::{OwnedValuePath, path}, @@ -72,6 +73,18 @@ pub async fn try_bind_tcp_listener( .map(|listener| listener.with_allowlist(allowlist)) } +/// Controls how Vector closes TCP connections during shutdown or connection expiry. +#[configurable_component] +#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)] +#[serde(rename_all = "snake_case")] +pub enum DisconnectMode { + /// Gracefully signal the end of the connection and wait for the client to close. + #[default] + Drain, + /// Immediately terminate the connection, without waiting for the client. + Abort, +} + #[derive(Clone, Copy, Eq, PartialEq)] pub enum TcpSourceAck { Ack, @@ -126,6 +139,7 @@ where tls_client_metadata_key: Option, receive_buffer_bytes: Option, max_connection_duration_secs: Option, + disconnect_mode: DisconnectMode, cx: SourceContext, acknowledgements: SourceAcknowledgementsConfig, max_connections: Option, @@ -215,6 +229,7 @@ where keepalive, receive_buffer_bytes, max_connection_duration_secs, + disconnect_mode, source, tripwire, peer_addr, @@ -255,6 +270,7 @@ async fn handle_stream( keepalive: Option, receive_buffer_bytes: Option, max_connection_duration_secs: Option, + disconnect_mode: DisconnectMode, source: T, mut tripwire: BoxFuture<'static, ()>, peer_addr: SocketAddr, @@ -320,13 +336,13 @@ async fn handle_stream( let mut permit = tokio::select! { _ = &mut tripwire => break, Some(_) = &mut connection_close_timeout => { - if close_socket(reader.get_ref().get_ref().get_ref()) { + if close_socket(reader.get_ref().get_ref().get_ref(), disconnect_mode) { break; } None }, _ = &mut shutdown_signal => { - if close_socket(reader.get_ref().get_ref().get_ref()) { + if close_socket(reader.get_ref().get_ref().get_ref(), disconnect_mode) { break; } None @@ -343,7 +359,7 @@ async fn handle_stream( tokio::select! { _ = &mut tripwire => break, _ = &mut shutdown_signal => { - if close_socket(reader.get_ref().get_ref().get_ref()) { + if close_socket(reader.get_ref().get_ref().get_ref(), disconnect_mode) { break; } }, @@ -461,12 +477,26 @@ async fn handle_stream( } } -fn close_socket(socket: &MaybeTlsIncomingStream) -> bool { - debug!("Start graceful shutdown."); - // Close our write part of TCP socket to signal the other side - // that it should stop writing and close the channel. +fn close_socket( + socket: &MaybeTlsIncomingStream, + disconnect_mode: DisconnectMode, +) -> bool { + // Signal to the client that we are done with this connection + // according to the configured disconnect mode. if let Some(stream) = socket.get_ref() { let socket = SockRef::from(stream); + match disconnect_mode { + DisconnectMode::Drain => { + debug!("Closing connection gracefully."); + } + DisconnectMode::Abort => { + debug!("Terminating connection."); + if let Err(error) = socket.set_linger(Some(std::time::Duration::ZERO)) { + warn!(message = "Failed to set SO_LINGER=0 on TCP socket.", %error); + } + return true; + } + } if let Err(error) = socket.shutdown(std::net::Shutdown::Write) { warn!(message = "Failed in signalling to the other side to close the TCP channel.", %error); }