use std::time::Duration;
use super::{
SaslClientAuthenticator, SaslClientAuthenticatorFactory, SaslConfig, SecurityConfig, TlsConfig,
};
pub const DEFAULT_REQUEST_TIMEOUT: Duration = Duration::from_secs(30);
pub const DEFAULT_SOCKET_CONNECTION_SETUP_TIMEOUT: Duration = Duration::from_secs(10);
pub const DEFAULT_SOCKET_CONNECTION_SETUP_TIMEOUT_MAX: Duration = Duration::from_secs(30);
pub const DEFAULT_METADATA_MAX_AGE: Duration = Duration::from_mins(5);
pub const DEFAULT_METADATA_MAX_IDLE: Duration = Duration::from_mins(5);
pub const DEFAULT_METADATA_REBOOTSTRAP_TRIGGER: Duration = Duration::from_mins(5);
pub const DEFAULT_METADATA_REFRESH_BACKOFF_INITIAL: Duration = Duration::from_millis(100);
pub const DEFAULT_METADATA_REFRESH_BACKOFF_MAX: Duration = Duration::from_secs(1);
pub(super) const DEFAULT_CONNECTIONS_MAX_IDLE: Duration = Duration::from_mins(9);
pub const DEFAULT_MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION: usize = 5;
pub const DEFAULT_BROKER_QUEUE_CAPACITY: usize = 1024;
pub const DEFAULT_RECONNECT_BACKOFF_INITIAL: Duration = Duration::from_millis(100);
pub const DEFAULT_RECONNECT_BACKOFF_MAX: Duration = Duration::from_secs(5);
pub const DEFAULT_BUFFER_POOL_CAPACITY: usize = 0;
pub const DEFAULT_TCP_NODELAY: bool = true;
pub const DEFAULT_REUSE_ADDRESS: bool = true;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MetadataRecoveryStrategy {
Rebootstrap,
None,
}
#[derive(Debug, Clone, PartialEq)]
pub struct ConnectionConfig {
pub socket: SocketConfig,
pub security: SecurityConfig,
pub tls: TlsConfig,
pub sasl: SaslConfig,
pub transport: TransportConfig,
pub request_timeout: Duration,
pub socket_connection_setup_timeout: Duration,
pub socket_connection_setup_timeout_max: Duration,
pub read_buffer_capacity: Option<usize>,
pub metadata_max_age: Duration,
pub metadata_max_idle: Duration,
pub metadata_recovery_strategy: MetadataRecoveryStrategy,
pub metadata_rebootstrap_trigger: Duration,
pub metadata_refresh_backoff_initial: Duration,
pub metadata_refresh_backoff_max: Duration,
pub connections_max_idle: Duration,
pub max_in_flight_requests_per_connection: usize,
pub broker_queue_capacity: usize,
pub reconnect_backoff_initial: Duration,
pub reconnect_backoff_max: Duration,
pub buffer_pool_capacity: usize,
pub allow_auto_topic_creation: bool,
pub use_all_dns_ips: bool,
}
impl Default for ConnectionConfig {
fn default() -> Self {
Self {
socket: SocketConfig::default(),
security: SecurityConfig::default(),
tls: TlsConfig::default(),
sasl: SaslConfig::default(),
transport: TransportConfig::Plaintext,
request_timeout: DEFAULT_REQUEST_TIMEOUT,
socket_connection_setup_timeout: DEFAULT_SOCKET_CONNECTION_SETUP_TIMEOUT,
socket_connection_setup_timeout_max: DEFAULT_SOCKET_CONNECTION_SETUP_TIMEOUT_MAX,
read_buffer_capacity: None,
metadata_max_age: DEFAULT_METADATA_MAX_AGE,
metadata_max_idle: DEFAULT_METADATA_MAX_IDLE,
metadata_recovery_strategy: MetadataRecoveryStrategy::Rebootstrap,
metadata_rebootstrap_trigger: DEFAULT_METADATA_REBOOTSTRAP_TRIGGER,
metadata_refresh_backoff_initial: DEFAULT_METADATA_REFRESH_BACKOFF_INITIAL,
metadata_refresh_backoff_max: DEFAULT_METADATA_REFRESH_BACKOFF_MAX,
connections_max_idle: DEFAULT_CONNECTIONS_MAX_IDLE,
max_in_flight_requests_per_connection: DEFAULT_MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION,
broker_queue_capacity: DEFAULT_BROKER_QUEUE_CAPACITY,
reconnect_backoff_initial: DEFAULT_RECONNECT_BACKOFF_INITIAL,
reconnect_backoff_max: DEFAULT_RECONNECT_BACKOFF_MAX,
buffer_pool_capacity: DEFAULT_BUFFER_POOL_CAPACITY,
allow_auto_topic_creation: false,
use_all_dns_ips: true,
}
}
}
impl From<SocketConfig> for ConnectionConfig {
fn from(socket: SocketConfig) -> Self {
Self {
socket,
..Self::default()
}
}
}
impl ConnectionConfig {
#[must_use]
pub const fn socket(mut self, socket: SocketConfig) -> Self {
self.socket = socket;
self
}
#[must_use]
pub fn sasl(mut self, sasl: SaslConfig) -> Self {
self.sasl = sasl;
self
}
#[must_use]
pub fn sasl_client_authenticator(
mut self,
authenticator: impl SaslClientAuthenticator,
) -> Self {
self.sasl = self.sasl.client_authenticator(authenticator);
self
}
#[must_use]
pub fn sasl_client_authenticator_factory(
mut self,
factory: impl SaslClientAuthenticatorFactory,
) -> Self {
self.sasl = self.sasl.client_authenticator_factory(factory);
self
}
#[must_use]
pub const fn request_timeout(mut self, timeout: Duration) -> Self {
self.request_timeout = timeout;
self
}
#[must_use]
pub const fn socket_connection_setup_timeout(mut self, timeout: Duration) -> Self {
self.socket_connection_setup_timeout = timeout;
self
}
#[must_use]
pub const fn socket_connection_setup_timeout_max(mut self, timeout: Duration) -> Self {
self.socket_connection_setup_timeout_max = timeout;
self
}
#[must_use]
pub const fn read_buffer_capacity(mut self, bytes: usize) -> Self {
self.read_buffer_capacity = Some(bytes);
self
}
#[must_use]
pub const fn metadata_max_age(mut self, timeout: Duration) -> Self {
self.metadata_max_age = timeout;
self
}
#[must_use]
pub const fn metadata_max_idle(mut self, timeout: Duration) -> Self {
self.metadata_max_idle = timeout;
self
}
#[must_use]
pub const fn metadata_recovery_strategy(mut self, strategy: MetadataRecoveryStrategy) -> Self {
self.metadata_recovery_strategy = strategy;
self
}
#[must_use]
pub const fn metadata_rebootstrap_trigger(mut self, timeout: Duration) -> Self {
self.metadata_rebootstrap_trigger = timeout;
self
}
#[must_use]
pub const fn metadata_refresh_backoff_initial(mut self, timeout: Duration) -> Self {
self.metadata_refresh_backoff_initial = timeout;
self
}
#[must_use]
pub const fn metadata_refresh_backoff_max(mut self, timeout: Duration) -> Self {
self.metadata_refresh_backoff_max = timeout;
self
}
#[must_use]
pub const fn connections_max_idle(mut self, timeout: Duration) -> Self {
self.connections_max_idle = timeout;
self
}
#[must_use]
pub const fn max_in_flight_requests_per_connection(mut self, requests: usize) -> Self {
self.max_in_flight_requests_per_connection = if requests == 0 { 1 } else { requests };
self
}
#[must_use]
pub const fn broker_queue_capacity(mut self, requests: usize) -> Self {
self.broker_queue_capacity = if requests == 0 { 1 } else { requests };
self
}
#[must_use]
pub const fn reconnect_backoff_initial(mut self, timeout: Duration) -> Self {
self.reconnect_backoff_initial = timeout;
self
}
#[must_use]
pub const fn reconnect_backoff_max(mut self, timeout: Duration) -> Self {
self.reconnect_backoff_max = timeout;
self
}
#[must_use]
pub const fn buffer_pool_capacity(mut self, buffers: usize) -> Self {
self.buffer_pool_capacity = buffers;
self
}
#[must_use]
pub const fn allow_auto_topic_creation(mut self, allow: bool) -> Self {
self.allow_auto_topic_creation = allow;
self
}
}
#[non_exhaustive]
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum TransportConfig {
Plaintext,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SocketConfig {
pub send_buffer_bytes: Option<usize>,
pub receive_buffer_bytes: Option<usize>,
pub tcp_nodelay: bool,
pub tcp_keepalive: Option<TcpKeepaliveConfig>,
pub tcp_notsent_lowat_bytes: Option<u32>,
pub tcp_quickack: Option<bool>,
pub tcp_user_timeout_ms: Option<Duration>,
pub tcp_congestion: Option<TcpCongestionControl>,
pub reuse_address: bool,
}
impl SocketConfig {
pub const SEND_BUFFER_BYTES_CONFIG: &'static str = "send.buffer.bytes";
pub const RECEIVE_BUFFER_BYTES_CONFIG: &'static str = "receive.buffer.bytes";
#[must_use]
pub const fn tcp_nodelay(mut self, enabled: bool) -> Self {
self.tcp_nodelay = enabled;
self
}
#[must_use]
pub const fn tcp_keepalive(mut self, keepalive: Option<TcpKeepaliveConfig>) -> Self {
self.tcp_keepalive = keepalive;
self
}
#[must_use]
pub const fn send_buffer_bytes(mut self, bytes: usize) -> Self {
self.send_buffer_bytes = Some(bytes);
self
}
#[must_use]
pub const fn receive_buffer_bytes(mut self, bytes: usize) -> Self {
self.receive_buffer_bytes = Some(bytes);
self
}
#[must_use]
pub const fn tcp_notsent_lowat_bytes(mut self, bytes: u32) -> Self {
self.tcp_notsent_lowat_bytes = Some(bytes);
self
}
#[must_use]
pub const fn tcp_quickack(mut self, enabled: bool) -> Self {
self.tcp_quickack = Some(enabled);
self
}
#[must_use]
pub const fn tcp_user_timeout_ms(mut self, timeout: Option<Duration>) -> Self {
self.tcp_user_timeout_ms = timeout;
self
}
#[must_use]
pub const fn tcp_congestion(mut self, congestion: TcpCongestionControl) -> Self {
self.tcp_congestion = Some(congestion);
self
}
#[must_use]
pub const fn reuse_address(mut self, enabled: bool) -> Self {
self.reuse_address = enabled;
self
}
}
impl Default for SocketConfig {
fn default() -> Self {
Self {
send_buffer_bytes: None,
receive_buffer_bytes: None,
tcp_nodelay: DEFAULT_TCP_NODELAY,
tcp_keepalive: None,
tcp_notsent_lowat_bytes: None,
tcp_quickack: None,
tcp_user_timeout_ms: None,
tcp_congestion: None,
reuse_address: DEFAULT_REUSE_ADDRESS,
}
}
}
#[cfg(test)]
mod default_tests {
#![allow(
clippy::expect_used,
clippy::missing_assert_message,
clippy::unwrap_used,
reason = "Unit test fixtures fail fastest with contextual unwrap/expect calls."
)]
use std::time::Duration;
use super::{
ConnectionConfig, DEFAULT_BROKER_QUEUE_CAPACITY, DEFAULT_BUFFER_POOL_CAPACITY,
DEFAULT_CONNECTIONS_MAX_IDLE, DEFAULT_MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION,
DEFAULT_METADATA_MAX_AGE, DEFAULT_RECONNECT_BACKOFF_INITIAL, DEFAULT_RECONNECT_BACKOFF_MAX,
DEFAULT_REQUEST_TIMEOUT, DEFAULT_REUSE_ADDRESS, DEFAULT_SOCKET_CONNECTION_SETUP_TIMEOUT,
DEFAULT_SOCKET_CONNECTION_SETUP_TIMEOUT_MAX, DEFAULT_TCP_NODELAY, SocketConfig,
TcpKeepaliveConfig,
};
#[test]
fn connection_defaults_are_named_kafka_runtime_values() {
let config = ConnectionConfig::default();
assert_eq!(config.request_timeout, DEFAULT_REQUEST_TIMEOUT);
assert_eq!(
config.socket_connection_setup_timeout,
DEFAULT_SOCKET_CONNECTION_SETUP_TIMEOUT
);
assert_eq!(
config.socket_connection_setup_timeout_max,
DEFAULT_SOCKET_CONNECTION_SETUP_TIMEOUT_MAX
);
assert_eq!(config.metadata_max_age, DEFAULT_METADATA_MAX_AGE);
assert_eq!(config.connections_max_idle, DEFAULT_CONNECTIONS_MAX_IDLE);
assert_eq!(
config.max_in_flight_requests_per_connection,
DEFAULT_MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION
);
assert_eq!(config.broker_queue_capacity, DEFAULT_BROKER_QUEUE_CAPACITY);
assert_eq!(
config.reconnect_backoff_initial,
DEFAULT_RECONNECT_BACKOFF_INITIAL
);
assert_eq!(config.reconnect_backoff_max, DEFAULT_RECONNECT_BACKOFF_MAX);
assert_eq!(config.buffer_pool_capacity, DEFAULT_BUFFER_POOL_CAPACITY);
}
#[test]
fn connection_builder_methods_set_every_runtime_field() {
let socket = SocketConfig::default().tcp_nodelay(false);
let config = ConnectionConfig::default()
.socket(socket.clone())
.request_timeout(Duration::from_millis(1))
.socket_connection_setup_timeout(Duration::from_millis(2))
.socket_connection_setup_timeout_max(Duration::from_millis(3))
.read_buffer_capacity(4)
.metadata_max_age(Duration::from_millis(5))
.connections_max_idle(Duration::from_millis(9))
.max_in_flight_requests_per_connection(0)
.broker_queue_capacity(0)
.reconnect_backoff_initial(Duration::from_millis(6))
.reconnect_backoff_max(Duration::from_millis(7))
.buffer_pool_capacity(8);
assert_eq!(ConnectionConfig::from(socket.clone()).socket, socket);
assert_eq!(config.request_timeout, Duration::from_millis(1));
assert_eq!(
config.socket_connection_setup_timeout,
Duration::from_millis(2)
);
assert_eq!(
config.socket_connection_setup_timeout_max,
Duration::from_millis(3)
);
assert_eq!(config.read_buffer_capacity, Some(4));
assert_eq!(config.metadata_max_age, Duration::from_millis(5));
assert_eq!(config.connections_max_idle, Duration::from_millis(9));
assert_eq!(config.max_in_flight_requests_per_connection, 1);
assert_eq!(config.broker_queue_capacity, 1);
assert_eq!(config.reconnect_backoff_initial, Duration::from_millis(6));
assert_eq!(config.reconnect_backoff_max, Duration::from_millis(7));
assert_eq!(config.buffer_pool_capacity, 8);
}
#[test]
fn socket_defaults_and_builders_cover_tcp_options() {
let keepalive = TcpKeepaliveConfig {
idle: Duration::from_secs(1),
interval: Duration::from_secs(2),
};
let socket = SocketConfig::default()
.tcp_nodelay(false)
.tcp_keepalive(Some(keepalive))
.send_buffer_bytes(128)
.receive_buffer_bytes(256)
.tcp_notsent_lowat_bytes(512)
.tcp_quickack(true)
.tcp_user_timeout_ms(Some(Duration::from_secs(3)))
.reuse_address(false);
assert_eq!(SocketConfig::default().tcp_nodelay, DEFAULT_TCP_NODELAY);
assert_eq!(SocketConfig::default().reuse_address, DEFAULT_REUSE_ADDRESS);
assert_eq!(socket.tcp_keepalive, Some(keepalive));
assert_eq!(socket.send_buffer_bytes, Some(128));
assert_eq!(socket.receive_buffer_bytes, Some(256));
assert_eq!(socket.tcp_notsent_lowat_bytes, Some(512));
assert_eq!(socket.tcp_quickack, Some(true));
assert_eq!(socket.tcp_user_timeout_ms, Some(Duration::from_secs(3)));
assert!(!socket.reuse_address);
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct TcpKeepaliveConfig {
pub idle: Duration,
pub interval: Duration,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum TcpCongestionControl {
Bbr,
Cubic,
Reno,
}
impl TcpCongestionControl {
#[cfg(any(target_os = "freebsd", target_os = "linux"))]
pub(crate) const fn as_bytes(self) -> &'static [u8] {
match self {
Self::Bbr => b"bbr",
Self::Cubic => b"cubic",
Self::Reno => b"reno",
}
}
}
#[cfg(test)]
mod tests {
#![allow(
clippy::expect_used,
clippy::missing_assert_message,
clippy::unwrap_used,
reason = "Unit test fixtures fail fastest with contextual unwrap/expect calls."
)]
use std::time::Duration;
use super::{ConnectionConfig, SocketConfig, TcpCongestionControl, TransportConfig};
#[test]
fn socket_config_uses_kafka_style_keys() {
assert_eq!(SocketConfig::SEND_BUFFER_BYTES_CONFIG, "send.buffer.bytes");
assert_eq!(
SocketConfig::RECEIVE_BUFFER_BYTES_CONFIG,
"receive.buffer.bytes"
);
}
#[test]
fn connection_config_layers_socket_and_transport() {
let socket = SocketConfig::default().send_buffer_bytes(65_536);
let config = ConnectionConfig::default()
.socket(socket)
.request_timeout(Duration::from_secs(7))
.socket_connection_setup_timeout(Duration::from_secs(3))
.socket_connection_setup_timeout_max(Duration::from_secs(11))
.connections_max_idle(Duration::from_secs(13))
.read_buffer_capacity(256 * 1024);
assert_eq!(config.socket.send_buffer_bytes, Some(65_536));
assert_eq!(config.transport, TransportConfig::Plaintext);
assert_eq!(config.request_timeout, Duration::from_secs(7));
assert_eq!(
config.socket_connection_setup_timeout,
Duration::from_secs(3)
);
assert_eq!(
config.socket_connection_setup_timeout_max,
Duration::from_secs(11)
);
assert_eq!(config.connections_max_idle, Duration::from_secs(13));
assert_eq!(config.read_buffer_capacity, Some(256 * 1024));
assert_eq!(config.metadata_max_age, Duration::from_mins(5));
assert_eq!(config.max_in_flight_requests_per_connection, 5);
assert_eq!(config.broker_queue_capacity, 1024);
}
#[test]
fn socket_config_exposes_deep_tcp_knobs() {
let config = SocketConfig::default()
.tcp_notsent_lowat_bytes(65_536)
.tcp_quickack(true)
.tcp_user_timeout_ms(Some(Duration::from_secs(15)))
.tcp_congestion(TcpCongestionControl::Bbr);
assert_eq!(config.tcp_notsent_lowat_bytes, Some(65_536));
assert_eq!(config.tcp_quickack, Some(true));
assert_eq!(config.tcp_user_timeout_ms, Some(Duration::from_secs(15)));
assert_eq!(config.tcp_congestion, Some(TcpCongestionControl::Bbr));
}
}