use ruststream::SubscriptionSource;
use crate::broker::KafkaBroker;
use crate::error::KafkaError;
use crate::subscriber::KafkaSubscriber;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
#[non_exhaustive]
pub enum StartOffset {
#[default]
Committed,
Earliest,
Latest,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
#[non_exhaustive]
pub enum Commit {
#[default]
Auto,
Tracked,
}
#[derive(Debug, Clone)]
pub struct KafkaTopic {
topic: String,
group: Option<String>,
start: StartOffset,
commit: Commit,
config: Vec<(String, String)>,
}
impl KafkaTopic {
#[must_use]
pub fn new(topic: impl Into<String>) -> Self {
Self {
topic: topic.into(),
group: None,
start: StartOffset::default(),
commit: Commit::default(),
config: Vec::new(),
}
}
#[must_use]
pub fn group(mut self, group: impl Into<String>) -> Self {
self.group = Some(group.into());
self
}
#[must_use]
pub fn start(mut self, start: StartOffset) -> Self {
self.start = start;
self
}
#[must_use]
pub fn commit(mut self, commit: Commit) -> Self {
self.commit = commit;
self
}
#[must_use]
pub fn config(mut self, key: impl Into<String>, value: impl Into<String>) -> Self {
self.config.push((key.into(), value.into()));
self
}
#[must_use]
pub fn topic(&self) -> &str {
&self.topic
}
pub(crate) fn group_or<'a>(&'a self, fallback: Option<&'a str>) -> Option<&'a str> {
self.group.as_deref().or(fallback)
}
pub(crate) fn start_offset(&self) -> StartOffset {
self.start
}
pub(crate) fn commit_mode(&self) -> Commit {
self.commit
}
pub(crate) fn config_entries(&self) -> &[(String, String)] {
&self.config
}
}
impl SubscriptionSource<KafkaBroker> for KafkaTopic {
type Subscriber = KafkaSubscriber;
fn name(&self) -> &str {
&self.topic
}
async fn subscribe(self, broker: &KafkaBroker) -> Result<Self::Subscriber, KafkaError> {
broker.subscribe(self).await
}
}
#[cfg(feature = "testing")]
impl SubscriptionSource<crate::testing::KafkaTestBroker> for KafkaTopic {
type Subscriber = crate::testing::KafkaTestSubscriber;
fn name(&self) -> &str {
&self.topic
}
async fn subscribe(
self,
broker: &crate::testing::KafkaTestBroker,
) -> Result<Self::Subscriber, KafkaError> {
broker.subscribe(&self.topic).await
}
}