use std::time::Duration;
use crate::auth::AuthConfig;
use crate::metadata::MetadataRecoveryStrategy;
#[non_exhaustive]
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum AcknowledgementMode {
#[default]
Implicit,
Explicit,
}
#[non_exhaustive]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AcknowledgeType {
Accept = 1,
Release = 2,
Reject = 3,
Renew = 4,
}
impl AcknowledgeType {
#[inline]
pub fn to_i8(self) -> i8 {
self as i8
}
}
#[non_exhaustive]
#[derive(Debug, Clone)]
pub struct ShareConsumerConfig {
pub(crate) bootstrap_servers: String,
pub(crate) group_id: String,
pub(crate) client_id: String,
pub(crate) acknowledgement_mode: AcknowledgementMode,
pub(crate) fetch_min_bytes: i32,
pub(crate) fetch_max_bytes: i32,
pub(crate) max_poll_records: i32,
pub(crate) max_buffered_records: i32,
pub(crate) max_records: i32,
pub(crate) batch_size: i32,
pub(crate) fetch_max_wait: Duration,
pub(crate) request_timeout: Duration,
pub(crate) connect_timeout: Duration,
pub(crate) session_timeout: Duration,
pub(crate) heartbeat_interval: Duration,
pub(crate) client_rack: Option<String>,
pub(crate) metadata_max_age: Duration,
pub(crate) metadata_recovery_strategy: MetadataRecoveryStrategy,
pub(crate) metadata_recovery_rebootstrap_trigger: Duration,
pub(crate) metadata_topic_cache_ttl: Option<Duration>,
pub(crate) auth: Option<AuthConfig>,
pub(crate) max_decompressed_size: usize,
pub(crate) transport: crate::network::TransportConfig,
}
impl Default for ShareConsumerConfig {
fn default() -> Self {
Self {
bootstrap_servers: String::new(),
group_id: String::new(),
client_id: "krafka".to_string(),
acknowledgement_mode: AcknowledgementMode::Implicit,
fetch_min_bytes: 1,
fetch_max_bytes: 52_428_800, max_poll_records: 500,
max_buffered_records: 500,
max_records: 5000,
batch_size: 500,
fetch_max_wait: Duration::from_millis(500),
request_timeout: Duration::from_secs(30),
connect_timeout: crate::network::DEFAULT_CONNECT_TIMEOUT,
session_timeout: Duration::from_secs(45),
heartbeat_interval: Duration::from_secs(5),
client_rack: None,
metadata_max_age: Duration::from_secs(300),
metadata_recovery_strategy: MetadataRecoveryStrategy::Rebootstrap,
metadata_recovery_rebootstrap_trigger: Duration::from_secs(300),
metadata_topic_cache_ttl: Some(Duration::from_secs(300)),
auth: None,
max_decompressed_size: crate::protocol::RecordBatch::MAX_DECOMPRESSED_SIZE,
transport: crate::network::TransportConfig::default(),
}
}
}
impl ShareConsumerConfig {
#[inline]
pub fn bootstrap_servers(&self) -> &str {
&self.bootstrap_servers
}
#[inline]
pub fn group_id(&self) -> &str {
&self.group_id
}
#[inline]
pub fn client_id(&self) -> &str {
&self.client_id
}
#[inline]
pub fn acknowledgement_mode(&self) -> AcknowledgementMode {
self.acknowledgement_mode
}
#[inline]
pub fn session_timeout(&self) -> Duration {
self.session_timeout
}
#[inline]
pub fn heartbeat_interval(&self) -> Duration {
self.heartbeat_interval
}
#[inline]
pub fn fetch_min_bytes(&self) -> i32 {
self.fetch_min_bytes
}
#[inline]
pub fn fetch_max_bytes(&self) -> i32 {
self.fetch_max_bytes
}
#[inline]
pub fn max_poll_records(&self) -> i32 {
self.max_poll_records
}
#[inline]
pub fn max_buffered_records(&self) -> i32 {
self.max_buffered_records
}
#[inline]
pub fn max_records(&self) -> i32 {
self.max_records
}
#[inline]
pub fn batch_size(&self) -> i32 {
self.batch_size
}
#[inline]
pub fn fetch_max_wait(&self) -> Duration {
self.fetch_max_wait
}
#[inline]
pub fn request_timeout(&self) -> Duration {
self.request_timeout
}
#[inline]
pub fn connect_timeout(&self) -> Duration {
self.connect_timeout
}
#[inline]
pub fn client_rack(&self) -> Option<&str> {
self.client_rack.as_deref()
}
#[inline]
pub fn metadata_max_age(&self) -> Duration {
self.metadata_max_age
}
#[inline]
pub fn metadata_recovery_strategy(&self) -> MetadataRecoveryStrategy {
self.metadata_recovery_strategy
}
#[inline]
pub fn metadata_recovery_rebootstrap_trigger(&self) -> Duration {
self.metadata_recovery_rebootstrap_trigger
}
#[inline]
pub fn metadata_topic_cache_ttl(&self) -> Option<Duration> {
self.metadata_topic_cache_ttl
}
#[inline]
pub fn auth(&self) -> Option<&AuthConfig> {
self.auth.as_ref()
}
#[inline]
pub fn max_decompressed_size(&self) -> usize {
self.max_decompressed_size
}
}