photon_backend/models/
subscribe_opts.rs1use super::SubscriptionMode;
4
5#[derive(Debug, Clone)]
10pub struct GroupOpts {
11 pub group_id: String,
13 pub shard_count: Option<u32>,
15 pub instance_id: Option<String>,
17}
18
19impl GroupOpts {
20 pub fn new(group_id: impl Into<String>) -> Self {
22 Self {
23 group_id: group_id.into(),
24 shard_count: None,
25 instance_id: None,
26 }
27 }
28
29 #[must_use]
31 pub const fn shards(mut self, count: u32) -> Self {
32 self.shard_count = Some(count);
33 self
34 }
35
36 #[must_use]
38 pub fn instance_id(mut self, id: impl Into<String>) -> Self {
39 self.instance_id = Some(id.into());
40 self
41 }
42}
43
44#[derive(Debug, Clone, Default)]
79pub struct SubscribeOpts {
80 pub subscription_name: Option<String>,
82 pub topic_key_filter: Option<String>,
84 pub mode: SubscriptionMode,
86 pub consumer_group: Option<GroupOpts>,
88}
89
90impl SubscribeOpts {
91 #[must_use]
93 pub const fn default_ephemeral() -> Self {
94 Self {
95 subscription_name: None,
96 topic_key_filter: None,
97 mode: SubscriptionMode::Ephemeral,
98 consumer_group: None,
99 }
100 }
101
102 #[must_use]
104 pub fn broadcast() -> Self {
105 Self {
106 mode: SubscriptionMode::Durable,
107 ..Self::default_ephemeral()
108 }
109 }
110
111 #[must_use]
113 pub fn consumer_group(mut self, opts: GroupOpts) -> Self {
114 self.subscription_name = None;
115 self.topic_key_filter = None;
116 self.mode = SubscriptionMode::Durable;
117 self.consumer_group = Some(opts);
118 self
119 }
120
121 #[must_use]
123 pub fn subscription_name(mut self, name: impl Into<String>) -> Self {
124 self.subscription_name = Some(name.into());
125 self
126 }
127
128 #[must_use]
130 pub fn topic_key_filter(mut self, key: impl Into<String>) -> Self {
131 self.topic_key_filter = Some(key.into());
132 self
133 }
134
135 #[must_use]
137 pub const fn mode(mut self, mode: SubscriptionMode) -> Self {
138 self.mode = mode;
139 self
140 }
141}
142
143#[derive(Debug, Default)]
148pub struct SubscriptionHandle {
149 _private: (),
150}
151
152impl SubscriptionHandle {
153 #[must_use]
155 pub const fn new() -> Self {
156 Self { _private: () }
157 }
158}
159
160#[cfg(test)]
161mod tests {
162 use super::*;
163
164 #[test]
165 fn default_ephemeral_has_no_name_or_filter() {
166 let opts = SubscribeOpts::default_ephemeral();
167 assert!(opts.subscription_name.is_none());
168 assert!(opts.topic_key_filter.is_none());
169 assert!(opts.consumer_group.is_none());
170 assert_eq!(opts.mode, SubscriptionMode::Ephemeral);
171 }
172
173 #[test]
174 fn broadcast_is_durable() {
175 let opts = SubscribeOpts::broadcast().subscription_name("worker-a");
176 assert_eq!(opts.mode, SubscriptionMode::Durable);
177 assert_eq!(opts.subscription_name.as_deref(), Some("worker-a"));
178 }
179
180 #[test]
181 fn builder_setters_apply() {
182 let opts = SubscribeOpts::default_ephemeral()
183 .topic_key_filter("alice")
184 .mode(SubscriptionMode::Durable);
185 assert_eq!(opts.topic_key_filter.as_deref(), Some("alice"));
186 assert_eq!(opts.mode, SubscriptionMode::Durable);
187 }
188
189 #[test]
190 fn consumer_group_clears_name_and_filter() {
191 let opts = SubscribeOpts::broadcast()
192 .subscription_name("worker-a")
193 .topic_key_filter("alice")
194 .consumer_group(GroupOpts::new("group-1").shards(8).instance_id("node-2"));
195
196 assert!(opts.subscription_name.is_none());
197 assert!(opts.topic_key_filter.is_none());
198 assert_eq!(opts.mode, SubscriptionMode::Durable);
199
200 let group = opts.consumer_group.expect("group opts set");
201 assert_eq!(group.group_id, "group-1");
202 assert_eq!(group.shard_count, Some(8));
203 assert_eq!(group.instance_id.as_deref(), Some("node-2"));
204 }
205}