use std::time::Duration;
use ruststream::SubscriptionSource;
use crate::broker::ConnectedPulsarBroker;
use crate::error::PulsarError;
use crate::subscriber::PulsarSubscriber;
use crate::topic::PulsarTopic;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum SubscriptionType {
Exclusive,
#[default]
Shared,
Failover,
KeyShared,
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[must_use]
pub struct DeadLetter {
pub(crate) topic: String,
pub(crate) max_deliveries: usize,
}
impl DeadLetter {
pub fn new(topic: impl Into<String>) -> Self {
Self {
topic: topic.into(),
max_deliveries: 5,
}
}
pub fn max_deliveries(mut self, max: usize) -> Self {
self.max_deliveries = max;
self
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum Topics {
List(Vec<String>),
Pattern(String),
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[must_use]
pub struct PulsarSubscription {
pub(crate) topics: Topics,
pub(crate) subscription: String,
pub(crate) sub_type: SubscriptionType,
pub(crate) dead_letter: Option<DeadLetter>,
pub(crate) ack_timeout: Option<Duration>,
}
impl PulsarSubscription {
pub fn new(topic: impl Into<String>, subscription: impl Into<String>) -> Self {
Self {
topics: Topics::List(vec![topic.into()]),
subscription: subscription.into(),
sub_type: SubscriptionType::default(),
dead_letter: None,
ack_timeout: None,
}
}
pub fn topics<I, S>(topics: I, subscription: impl Into<String>) -> Self
where
I: IntoIterator<Item = S>,
S: Into<String>,
{
Self {
topics: Topics::List(topics.into_iter().map(Into::into).collect()),
..Self::new(String::new(), subscription)
}
}
pub fn pattern(pattern: impl Into<String>, subscription: impl Into<String>) -> Self {
Self {
topics: Topics::Pattern(pattern.into()),
..Self::new(String::new(), subscription)
}
}
pub fn subscription_type(mut self, sub_type: SubscriptionType) -> Self {
self.sub_type = sub_type;
self
}
pub fn dead_letter(mut self, dead_letter: DeadLetter) -> Self {
self.dead_letter = Some(dead_letter);
self
}
pub fn ack_timeout(mut self, timeout: Duration) -> Self {
self.ack_timeout = Some(timeout);
self
}
#[must_use]
pub fn subscription(&self) -> &str {
&self.subscription
}
pub(crate) fn display_topic(&self) -> String {
match &self.topics {
Topics::List(topics) => topics.join(","),
Topics::Pattern(pattern) => pattern.clone(),
}
}
pub(crate) fn validate(&self) -> Result<(), PulsarError> {
if self.subscription.is_empty() {
return Err(PulsarError::Invalid(
"subscription name must be non-empty".into(),
));
}
match &self.topics {
Topics::List(topics) => {
if topics.is_empty() || topics.iter().any(String::is_empty) {
return Err(PulsarError::Invalid("topics must be non-empty".into()));
}
for topic in topics {
let _ = PulsarTopic::parse(topic)?;
}
}
Topics::Pattern(pattern) => {
regex::Regex::new(pattern).map_err(|e| {
PulsarError::Invalid(format!("invalid topic pattern '{pattern}': {e}"))
})?;
}
}
Ok(())
}
}
impl SubscriptionSource<ConnectedPulsarBroker> for PulsarSubscription {
type Subscriber = PulsarSubscriber;
fn name(&self) -> &str {
match &self.topics {
Topics::List(topics) if topics.len() == 1 => &topics[0],
_ => &self.subscription,
}
}
async fn subscribe(
self,
connected: &ConnectedPulsarBroker,
) -> Result<PulsarSubscriber, PulsarError> {
connected.subscribe_descriptor(self).await
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn empty_subscription_is_rejected_before_io() {
assert!(matches!(
PulsarSubscription::new("orders", "").validate(),
Err(PulsarError::Invalid(_))
));
}
#[test]
fn malformed_topics_are_rejected_before_io() {
assert!(matches!(
PulsarSubscription::new("a/b", "workers").validate(),
Err(PulsarError::Invalid(_))
));
}
#[test]
fn malformed_patterns_are_rejected_before_io() {
assert!(matches!(
PulsarSubscription::pattern("orders-(", "workers").validate(),
Err(PulsarError::Invalid(_))
));
}
}