ruststream_fred/
stream.rs1use std::sync::atomic::{AtomicU64, Ordering};
15use std::time::Duration;
16
17use ruststream::SubscriptionSource;
18
19use crate::broker::ConnectedRedisBroker;
20use crate::deadletter::PoisonPolicy;
21use crate::delay::{DelayConfig, DelayedRetry};
22use crate::{error::RedisError, subscriber::RedisSubscriber};
23
24const DEFAULT_COUNT: u64 = 64;
25const DEFAULT_BLOCK: Duration = Duration::from_secs(5);
26
27fn auto_consumer() -> String {
30 static COUNTER: AtomicU64 = AtomicU64::new(0);
31 let n = COUNTER.fetch_add(1, Ordering::Relaxed);
32 format!("ruststream-{n}")
33}
34
35#[derive(Debug, Clone, Default)]
38pub enum StreamStart {
39 #[default]
41 New,
42 Beginning,
44 Id(String),
46}
47
48impl StreamStart {
49 pub(crate) fn as_id(&self) -> &str {
50 match self {
51 Self::New => "$",
52 Self::Beginning => "0",
53 Self::Id(id) => id,
54 }
55 }
56}
57
58#[derive(Debug, Clone)]
59pub(crate) enum ReadMode {
60 Fresh,
62 Reclaim { min_idle: Duration },
64}
65
66#[derive(Debug, Clone)]
82#[must_use]
83pub struct RedisStream {
84 key: String,
85 group: Option<String>,
86 consumer: Option<String>,
87 count: Option<u64>,
88 block: Option<Duration>,
89 start: StreamStart,
90 mode: ReadMode,
91 dead_letter: Option<String>,
92 max_deliveries: Option<u64>,
93 delayed_retry: Option<DelayedRetry>,
94}
95
96impl RedisStream {
97 pub fn new(key: impl Into<String>) -> Self {
101 Self {
102 key: key.into(),
103 group: None,
104 consumer: None,
105 count: None,
106 block: None,
107 start: StreamStart::New,
108 mode: ReadMode::Fresh,
109 dead_letter: None,
110 max_deliveries: None,
111 delayed_retry: None,
112 }
113 }
114
115 pub fn reclaim(key: impl Into<String>, min_idle: Duration) -> Self {
122 Self {
123 key: key.into(),
124 group: None,
125 consumer: None,
126 count: None,
127 block: None,
128 start: StreamStart::New,
129 mode: ReadMode::Reclaim { min_idle },
130 dead_letter: None,
131 max_deliveries: None,
132 delayed_retry: None,
133 }
134 }
135
136 pub fn group(mut self, group: impl Into<String>) -> Self {
138 self.group = Some(group.into());
139 self
140 }
141
142 pub fn consumer(mut self, consumer: impl Into<String>) -> Self {
144 self.consumer = Some(consumer.into());
145 self
146 }
147
148 pub const fn count(mut self, count: u64) -> Self {
150 self.count = Some(count);
151 self
152 }
153
154 pub const fn block(mut self, block: Duration) -> Self {
158 self.block = Some(block);
159 self
160 }
161
162 pub fn start_id(mut self, start: StreamStart) -> Self {
165 self.start = start;
166 self
167 }
168
169 pub fn dead_letter(mut self, key: impl Into<String>) -> Self {
173 self.dead_letter = Some(key.into());
174 self
175 }
176
177 pub const fn max_deliveries(mut self, max: u64) -> Self {
184 self.max_deliveries = Some(max);
185 self
186 }
187
188 pub fn delayed_retry(mut self, retry: DelayedRetry) -> Self {
198 self.delayed_retry = Some(retry);
199 self
200 }
201
202 #[must_use]
204 pub fn key(&self) -> &str {
205 &self.key
206 }
207
208 pub(crate) fn group_or_err(&self) -> Result<&str, RedisError> {
209 self.group.as_deref().ok_or_else(|| {
210 RedisError::InvalidOptions(format!(
211 "stream subscription on `{}` requires a consumer group: call .group(name)",
212 self.key
213 ))
214 })
215 }
216
217 pub(crate) fn consumer_or_auto(&self) -> String {
218 self.consumer.clone().unwrap_or_else(auto_consumer)
219 }
220
221 pub(crate) fn count_or_default(&self) -> u64 {
222 self.count.unwrap_or(DEFAULT_COUNT)
223 }
224
225 pub(crate) fn block_or_default(&self) -> Duration {
226 self.block.unwrap_or(DEFAULT_BLOCK)
227 }
228
229 pub(crate) const fn start(&self) -> &StreamStart {
230 &self.start
231 }
232
233 pub(crate) fn mode(&self) -> ReadMode {
234 self.mode.clone()
235 }
236
237 pub(crate) fn poison_policy(&self) -> PoisonPolicy {
238 PoisonPolicy {
239 dead_letter: self.dead_letter.clone(),
240 max_deliveries: self.max_deliveries,
241 }
242 }
243
244 pub(crate) fn delay_config(&self) -> Option<DelayConfig> {
245 self.delayed_retry.as_ref().map(DelayConfig::from_retry)
246 }
247}
248
249impl SubscriptionSource<ConnectedRedisBroker> for RedisStream {
250 type Subscriber = RedisSubscriber;
251
252 fn name(&self) -> &str {
253 self.key()
254 }
255
256 async fn subscribe(
257 self,
258 connected: &ConnectedRedisBroker,
259 ) -> Result<Self::Subscriber, RedisError> {
260 connected.subscribe(self).await
261 }
262}
263
264#[cfg(feature = "testing")]
265impl SubscriptionSource<crate::testing::ConnectedRedisTestBroker> for RedisStream {
266 type Subscriber = crate::testing::RedisTestSubscriber;
267
268 fn name(&self) -> &str {
269 self.key()
270 }
271
272 async fn subscribe(
273 self,
274 connected: &crate::testing::ConnectedRedisTestBroker,
275 ) -> Result<Self::Subscriber, RedisError> {
276 connected.subscribe(self.key()).await
277 }
278}
279
280#[cfg(test)]
281mod tests {
282 use super::*;
283
284 #[test]
285 fn group_is_required() {
286 let err = RedisStream::new("orders").group_or_err().unwrap_err();
287 assert!(matches!(err, RedisError::InvalidOptions(msg) if msg.contains("consumer group")));
288 }
289
290 #[test]
291 fn group_set_resolves() {
292 let s = RedisStream::new("orders").group("workers");
293 assert_eq!(s.group_or_err().expect("group set"), "workers");
294 }
295
296 #[test]
297 fn start_maps_to_redis_ids() {
298 assert_eq!(StreamStart::New.as_id(), "$");
299 assert_eq!(StreamStart::Beginning.as_id(), "0");
300 assert_eq!(StreamStart::Id("5-0".into()).as_id(), "5-0");
301 }
302
303 #[test]
304 fn reclaim_carries_min_idle() {
305 let s = RedisStream::reclaim("orders", Duration::from_secs(30)).group("g");
306 assert!(matches!(s.mode(), ReadMode::Reclaim { min_idle } if min_idle.as_secs() == 30));
307 }
308}