use super::SubscriptionMode;
#[derive(Debug, Clone)]
pub struct GroupOpts {
pub group_id: String,
pub shard_count: Option<u32>,
pub instance_id: Option<String>,
}
impl GroupOpts {
pub fn new(group_id: impl Into<String>) -> Self {
Self {
group_id: group_id.into(),
shard_count: None,
instance_id: None,
}
}
#[must_use]
pub const fn shards(mut self, count: u32) -> Self {
self.shard_count = Some(count);
self
}
#[must_use]
pub fn instance_id(mut self, id: impl Into<String>) -> Self {
self.instance_id = Some(id.into());
self
}
}
#[derive(Debug, Clone, Default)]
pub struct SubscribeOpts {
pub subscription_name: Option<String>,
pub topic_key_filter: Option<String>,
pub mode: SubscriptionMode,
pub consumer_group: Option<GroupOpts>,
}
impl SubscribeOpts {
#[must_use]
pub const fn default_ephemeral() -> Self {
Self {
subscription_name: None,
topic_key_filter: None,
mode: SubscriptionMode::Ephemeral,
consumer_group: None,
}
}
#[must_use]
pub fn broadcast() -> Self {
Self {
mode: SubscriptionMode::Durable,
..Self::default_ephemeral()
}
}
#[must_use]
pub fn consumer_group(mut self, opts: GroupOpts) -> Self {
self.subscription_name = None;
self.topic_key_filter = None;
self.mode = SubscriptionMode::Durable;
self.consumer_group = Some(opts);
self
}
#[must_use]
pub fn subscription_name(mut self, name: impl Into<String>) -> Self {
self.subscription_name = Some(name.into());
self
}
#[must_use]
pub fn topic_key_filter(mut self, key: impl Into<String>) -> Self {
self.topic_key_filter = Some(key.into());
self
}
#[must_use]
pub const fn mode(mut self, mode: SubscriptionMode) -> Self {
self.mode = mode;
self
}
}
#[derive(Debug, Default)]
pub struct SubscriptionHandle {
_private: (),
}
impl SubscriptionHandle {
#[must_use]
pub const fn new() -> Self {
Self { _private: () }
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn default_ephemeral_has_no_name_or_filter() {
let opts = SubscribeOpts::default_ephemeral();
assert!(opts.subscription_name.is_none());
assert!(opts.topic_key_filter.is_none());
assert!(opts.consumer_group.is_none());
assert_eq!(opts.mode, SubscriptionMode::Ephemeral);
}
#[test]
fn broadcast_is_durable() {
let opts = SubscribeOpts::broadcast().subscription_name("worker-a");
assert_eq!(opts.mode, SubscriptionMode::Durable);
assert_eq!(opts.subscription_name.as_deref(), Some("worker-a"));
}
#[test]
fn builder_setters_apply() {
let opts = SubscribeOpts::default_ephemeral()
.topic_key_filter("alice")
.mode(SubscriptionMode::Durable);
assert_eq!(opts.topic_key_filter.as_deref(), Some("alice"));
assert_eq!(opts.mode, SubscriptionMode::Durable);
}
#[test]
fn consumer_group_clears_name_and_filter() {
let opts = SubscribeOpts::broadcast()
.subscription_name("worker-a")
.topic_key_filter("alice")
.consumer_group(GroupOpts::new("group-1").shards(8).instance_id("node-2"));
assert!(opts.subscription_name.is_none());
assert!(opts.topic_key_filter.is_none());
assert_eq!(opts.mode, SubscriptionMode::Durable);
let group = opts.consumer_group.expect("group opts set");
assert_eq!(group.group_id, "group-1");
assert_eq!(group.shard_count, Some(8));
assert_eq!(group.instance_id.as_deref(), Some("node-2"));
}
}