1#[cfg(feature = "streams")]
4use crate::{
5 FromRedisValue, RedisWrite, ToRedisArgs, Value,
6 errors::{ParsingError, invalid_type_error},
7 types::HashMap,
8};
9use crate::{from_redis_value, from_redis_value_ref, types::ToSingleRedisArg};
10
11#[derive(PartialEq, Eq, Clone, Debug, Copy)]
17#[non_exhaustive]
18pub enum StreamMaxlen {
19 Equals(usize),
21 Approx(usize),
23}
24
25impl ToRedisArgs for StreamMaxlen {
26 fn write_redis_args<W>(&self, out: &mut W)
27 where
28 W: ?Sized + RedisWrite,
29 {
30 let (ch, val) = match *self {
31 Self::Equals(v) => ("=", v),
32 Self::Approx(v) => ("~", v),
33 };
34 out.write_arg(b"MAXLEN");
35 out.write_arg(ch.as_bytes());
36 val.write_redis_args(out);
37 }
38}
39
40#[derive(Debug)]
43#[non_exhaustive]
44pub enum StreamTrimmingMode {
45 Exact,
47 Approx,
49}
50
51impl ToRedisArgs for StreamTrimmingMode {
52 fn write_redis_args<W>(&self, out: &mut W)
53 where
54 W: ?Sized + RedisWrite,
55 {
56 match self {
57 Self::Exact => out.write_arg(b"="),
58 Self::Approx => out.write_arg(b"~"),
59 }
60 }
61}
62
63#[derive(Debug)]
67#[non_exhaustive]
68pub enum StreamTrimStrategy {
69 MaxLen(StreamTrimmingMode, usize, Option<usize>),
71 MinId(StreamTrimmingMode, String, Option<usize>),
73}
74
75impl StreamTrimStrategy {
76 pub fn maxlen(trim: StreamTrimmingMode, max_entries: usize) -> Self {
78 Self::MaxLen(trim, max_entries, None)
79 }
80
81 pub fn minid(trim: StreamTrimmingMode, stream_id: impl Into<String>) -> Self {
83 Self::MinId(trim, stream_id.into(), None)
84 }
85
86 pub fn limit(self, limit: usize) -> Self {
88 match self {
89 Self::MaxLen(m, t, _) => Self::MaxLen(m, t, Some(limit)),
90 Self::MinId(m, t, _) => Self::MinId(m, t, Some(limit)),
91 }
92 }
93}
94
95impl ToRedisArgs for StreamTrimStrategy {
96 fn write_redis_args<W>(&self, out: &mut W)
97 where
98 W: ?Sized + RedisWrite,
99 {
100 let limit = match self {
101 Self::MaxLen(m, t, limit) => {
102 out.write_arg(b"MAXLEN");
103 m.write_redis_args(out);
104 t.write_redis_args(out);
105 limit
106 }
107 Self::MinId(m, t, limit) => {
108 out.write_arg(b"MINID");
109 m.write_redis_args(out);
110 t.write_redis_args(out);
111 limit
112 }
113 };
114 if let Some(limit) = limit {
115 out.write_arg(b"LIMIT");
116 limit.write_redis_args(out);
117 }
118 }
119}
120
121#[derive(Debug)]
126pub struct StreamTrimOptions {
127 strategy: StreamTrimStrategy,
128 deletion_policy: Option<StreamDeletionPolicy>,
129}
130
131impl StreamTrimOptions {
132 pub fn maxlen(mode: StreamTrimmingMode, max_entries: usize) -> Self {
134 Self {
135 strategy: StreamTrimStrategy::maxlen(mode, max_entries),
136 deletion_policy: None,
137 }
138 }
139
140 pub fn minid(mode: StreamTrimmingMode, stream_id: impl Into<String>) -> Self {
142 Self {
143 strategy: StreamTrimStrategy::minid(mode, stream_id),
144 deletion_policy: None,
145 }
146 }
147
148 pub fn limit(mut self, limit: usize) -> Self {
150 self.strategy = self.strategy.limit(limit);
151 self
152 }
153
154 pub fn set_deletion_policy(mut self, deletion_policy: StreamDeletionPolicy) -> Self {
156 self.deletion_policy = Some(deletion_policy);
157 self
158 }
159}
160
161impl ToRedisArgs for StreamTrimOptions {
162 fn write_redis_args<W>(&self, out: &mut W)
163 where
164 W: ?Sized + RedisWrite,
165 {
166 self.strategy.write_redis_args(out);
167 if let Some(deletion_policy) = self.deletion_policy.as_ref() {
168 deletion_policy.write_redis_args(out);
169 }
170 }
171}
172
173#[derive(Debug, Clone)]
178pub enum StreamIdempotencyMode {
179 Manual {
183 producer_id: String,
185 idempotent_id: String,
187 },
188 Automatic {
192 producer_id: String,
194 },
195}
196
197impl ToRedisArgs for StreamIdempotencyMode {
198 fn write_redis_args<W>(&self, out: &mut W)
199 where
200 W: ?Sized + RedisWrite,
201 {
202 match self {
203 Self::Manual {
204 producer_id,
205 idempotent_id,
206 } => {
207 out.write_arg(b"IDMP");
208 out.write_arg(producer_id.as_bytes());
209 out.write_arg(idempotent_id.as_bytes());
210 }
211 Self::Automatic { producer_id } => {
212 out.write_arg(b"IDMPAUTO");
213 out.write_arg(producer_id.as_bytes());
214 }
215 }
216 }
217}
218
219pub const IDMP_DURATION_MIN: u32 = 1;
221pub const IDMP_DURATION_MAX: u32 = 86400;
223pub const IDMP_MAXSIZE_MIN: u16 = 1;
225pub const IDMP_MAXSIZE_MAX: u16 = 10000;
227
228#[derive(Debug)]
257pub struct StreamConfigOptions {
258 idmp_duration: Option<u32>,
262 idmp_maxsize: Option<u16>,
266}
267
268impl StreamConfigOptions {
269 pub fn with_idempotency_seconds(seconds: u32) -> Result<Self, String> {
278 if !(IDMP_DURATION_MIN..=IDMP_DURATION_MAX).contains(&seconds) {
279 return Err(format!(
280 "IDMP-DURATION must be between {IDMP_DURATION_MIN} and {IDMP_DURATION_MAX} seconds, got: {seconds}"
281 ));
282 }
283 Ok(Self {
284 idmp_duration: Some(seconds),
285 idmp_maxsize: None,
286 })
287 }
288
289 pub fn with_idempotency_maxsize(size: u16) -> Result<Self, String> {
298 if !(IDMP_MAXSIZE_MIN..=IDMP_MAXSIZE_MAX).contains(&size) {
299 return Err(format!(
300 "IDMP-MAXSIZE must be between {IDMP_MAXSIZE_MIN} and {IDMP_MAXSIZE_MAX} entries, got: {size}"
301 ));
302 }
303 Ok(Self {
304 idmp_duration: None,
305 idmp_maxsize: Some(size),
306 })
307 }
308
309 pub fn idempotency_seconds(mut self, seconds: u32) -> Result<Self, String> {
315 if !(IDMP_DURATION_MIN..=IDMP_DURATION_MAX).contains(&seconds) {
316 return Err(format!(
317 "IDMP-DURATION must be between {IDMP_DURATION_MIN} and {IDMP_DURATION_MAX} seconds, got: {seconds}"
318 ));
319 }
320 self.idmp_duration = Some(seconds);
321 Ok(self)
322 }
323
324 pub fn idempotency_maxsize(mut self, size: u16) -> Result<Self, String> {
330 if !(IDMP_MAXSIZE_MIN..=IDMP_MAXSIZE_MAX).contains(&size) {
331 return Err(format!(
332 "IDMP-MAXSIZE must be between {IDMP_MAXSIZE_MIN} and {IDMP_MAXSIZE_MAX} entries, got: {size}"
333 ));
334 }
335 self.idmp_maxsize = Some(size);
336 Ok(self)
337 }
338}
339
340impl ToRedisArgs for StreamConfigOptions {
341 fn write_redis_args<W>(&self, out: &mut W)
342 where
343 W: ?Sized + RedisWrite,
344 {
345 if let Some(duration) = self.idmp_duration {
346 out.write_arg(b"IDMP-DURATION");
347 out.write_arg(duration.to_string().as_bytes());
348 }
349 if let Some(maxsize) = self.idmp_maxsize {
350 out.write_arg(b"IDMP-MAXSIZE");
351 out.write_arg(maxsize.to_string().as_bytes());
352 }
353 }
354}
355
356#[derive(Default, Debug)]
361pub struct StreamAddOptions {
362 nomkstream: bool,
363 trim: Option<StreamTrimStrategy>,
364 deletion_policy: Option<StreamDeletionPolicy>,
365 idempotency: Option<StreamIdempotencyMode>,
366}
367
368impl StreamAddOptions {
369 pub fn nomkstream(mut self) -> Self {
371 self.nomkstream = true;
372 self
373 }
374
375 pub fn trim(mut self, trim: StreamTrimStrategy) -> Self {
377 self.trim = Some(trim);
378 self
379 }
380
381 pub fn set_deletion_policy(mut self, deletion_policy: StreamDeletionPolicy) -> Self {
383 self.deletion_policy = Some(deletion_policy);
384 self
385 }
386
387 pub fn idmp(
409 mut self,
410 producer_id: impl Into<String>,
411 idempotent_id: impl Into<String>,
412 ) -> Self {
413 self.idempotency = Some(StreamIdempotencyMode::Manual {
414 producer_id: producer_id.into(),
415 idempotent_id: idempotent_id.into(),
416 });
417 self
418 }
419
420 pub fn idmpauto(mut self, producer_id: impl Into<String>) -> Self {
441 self.idempotency = Some(StreamIdempotencyMode::Automatic {
442 producer_id: producer_id.into(),
443 });
444 self
445 }
446}
447
448impl ToRedisArgs for StreamAddOptions {
449 fn write_redis_args<W>(&self, out: &mut W)
450 where
451 W: ?Sized + RedisWrite,
452 {
453 if self.nomkstream {
454 out.write_arg(b"NOMKSTREAM");
455 }
456 if let Some(deletion_policy) = self.deletion_policy.as_ref() {
457 deletion_policy.write_redis_args(out);
458 }
459 if let Some(idempotency) = self.idempotency.as_ref() {
460 idempotency.write_redis_args(out);
461 }
462 if let Some(strategy) = self.trim.as_ref() {
463 strategy.write_redis_args(out);
464 }
465 }
466}
467
468#[derive(Default, Debug)]
473pub struct StreamAutoClaimOptions {
474 count: Option<usize>,
475 justid: bool,
476}
477
478impl StreamAutoClaimOptions {
479 pub fn count(mut self, n: usize) -> Self {
481 self.count = Some(n);
482 self
483 }
484
485 pub fn with_justid(mut self) -> Self {
488 self.justid = true;
489 self
490 }
491}
492
493impl ToRedisArgs for StreamAutoClaimOptions {
494 fn write_redis_args<W>(&self, out: &mut W)
495 where
496 W: ?Sized + RedisWrite,
497 {
498 if let Some(ref count) = self.count {
499 out.write_arg(b"COUNT");
500 out.write_arg(format!("{count}").as_bytes());
501 }
502 if self.justid {
503 out.write_arg(b"JUSTID");
504 }
505 }
506}
507
508#[derive(Default, Debug)]
513pub struct StreamClaimOptions {
514 idle: Option<usize>,
516 time: Option<usize>,
518 retry: Option<usize>,
520 force: bool,
522 justid: bool,
525 lastid: Option<String>,
527}
528
529impl StreamClaimOptions {
530 pub fn idle(mut self, ms: usize) -> Self {
532 self.idle = Some(ms);
533 self
534 }
535
536 pub fn time(mut self, ms_time: usize) -> Self {
538 self.time = Some(ms_time);
539 self
540 }
541
542 pub fn retry(mut self, count: usize) -> Self {
544 self.retry = Some(count);
545 self
546 }
547
548 pub fn with_force(mut self) -> Self {
550 self.force = true;
551 self
552 }
553
554 pub fn with_justid(mut self) -> Self {
557 self.justid = true;
558 self
559 }
560
561 pub fn with_lastid(mut self, lastid: impl Into<String>) -> Self {
563 self.lastid = Some(lastid.into());
564 self
565 }
566}
567
568impl ToRedisArgs for StreamClaimOptions {
569 fn write_redis_args<W>(&self, out: &mut W)
570 where
571 W: ?Sized + RedisWrite,
572 {
573 if let Some(ref ms) = self.idle {
574 out.write_arg(b"IDLE");
575 out.write_arg(format!("{ms}").as_bytes());
576 }
577 if let Some(ref ms_time) = self.time {
578 out.write_arg(b"TIME");
579 out.write_arg(format!("{ms_time}").as_bytes());
580 }
581 if let Some(ref count) = self.retry {
582 out.write_arg(b"RETRYCOUNT");
583 out.write_arg(format!("{count}").as_bytes());
584 }
585 if self.force {
586 out.write_arg(b"FORCE");
587 }
588 if self.justid {
589 out.write_arg(b"JUSTID");
590 }
591 if let Some(ref lastid) = self.lastid {
592 out.write_arg(b"LASTID");
593 lastid.write_redis_args(out);
594 }
595 }
596}
597
598type SRGroup = Option<(Vec<Vec<u8>>, Vec<Vec<u8>>)>;
602#[derive(Default, Debug)]
607pub struct StreamReadOptions {
608 block: Option<usize>,
610 count: Option<usize>,
612 noack: Option<bool>,
614 group: SRGroup,
617 claim: Option<usize>,
620}
621
622impl StreamReadOptions {
623 pub fn read_only(&self) -> bool {
626 self.group.is_none()
627 }
628
629 pub fn noack(mut self) -> Self {
633 self.noack = Some(true);
634 self
635 }
636
637 pub fn block(mut self, ms: usize) -> Self {
639 self.block = Some(ms);
640 self
641 }
642
643 pub fn count(mut self, n: usize) -> Self {
645 self.count = Some(n);
646 self
647 }
648
649 pub fn group<GN: ToRedisArgs, CN: ToRedisArgs>(
651 mut self,
652 group_name: GN,
653 consumer_name: CN,
654 ) -> Self {
655 self.group = Some((
656 ToRedisArgs::to_redis_args(&group_name),
657 ToRedisArgs::to_redis_args(&consumer_name),
658 ));
659 self
660 }
661
662 pub fn claim(mut self, min_idle_time: usize) -> Self {
664 self.claim = Some(min_idle_time);
665 self
666 }
667}
668
669impl ToRedisArgs for StreamReadOptions {
670 fn write_redis_args<W>(&self, out: &mut W)
671 where
672 W: ?Sized + RedisWrite,
673 {
674 if let Some(ref group) = self.group {
675 out.write_arg(b"GROUP");
676 for i in &group.0 {
677 out.write_arg(i);
678 }
679 for i in &group.1 {
680 out.write_arg(i);
681 }
682 }
683
684 if let Some(ref ms) = self.block {
685 out.write_arg(b"BLOCK");
686 out.write_arg(format!("{ms}").as_bytes());
687 }
688
689 if let Some(ref n) = self.count {
690 out.write_arg(b"COUNT");
691 out.write_arg(format!("{n}").as_bytes());
692 }
693
694 if self.group.is_some() {
695 if self.noack == Some(true) {
697 out.write_arg(b"NOACK");
698 }
699 if let Some(ref min_idle_time) = self.claim {
701 out.write_arg(b"CLAIM");
702 out.write_arg(format!("{min_idle_time}").as_bytes());
703 }
704 }
705 }
706}
707
708#[derive(Default, Debug, Clone)]
713pub struct StreamAutoClaimReply {
714 pub next_stream_id: String,
716 pub claimed: Vec<StreamId>,
718 pub deleted_ids: Vec<String>,
720 pub invalid_entries: bool,
724}
725
726#[derive(Default, Debug, Clone)]
732pub struct StreamReadReply {
733 pub keys: Vec<StreamKey>,
735}
736
737#[derive(Default, Debug, Clone)]
749pub struct StreamRangeReply {
750 pub ids: Vec<StreamId>,
752}
753
754#[derive(Default, Debug, Clone)]
761pub struct StreamClaimReply {
762 pub ids: Vec<StreamId>,
764}
765
766#[derive(Debug, Clone, Default)]
774#[non_exhaustive]
775pub enum StreamPendingReply {
776 #[default]
778 Empty,
779 Data(StreamPendingData),
781}
782
783impl StreamPendingReply {
784 pub fn count(&self) -> usize {
786 match self {
787 Self::Empty => 0,
788 Self::Data(x) => x.count,
789 }
790 }
791}
792
793#[derive(Default, Debug, Clone)]
797pub struct StreamPendingData {
798 pub count: usize,
800 pub start_id: String,
802 pub end_id: String,
804 pub consumers: Vec<StreamInfoConsumer>,
808}
809
810#[derive(Default, Debug, Clone)]
820pub struct StreamPendingCountReply {
821 pub ids: Vec<StreamPendingId>,
825}
826
827#[derive(Default, Debug, Clone)]
840pub struct StreamInfoStreamReply {
841 pub last_generated_id: String,
844 pub radix_tree_keys: usize,
847 pub groups: usize,
849 pub length: usize,
851 pub first_entry: StreamId,
853 pub last_entry: StreamId,
855}
856
857#[derive(Default, Debug, Clone)]
886#[non_exhaustive]
887pub struct StreamInfoStreamReplyWithIdempotency {
888 pub base: StreamInfoStreamReply,
890 pub idmp_duration: u32,
892 pub idmp_maxsize: u16,
894 pub pids_tracked: usize,
896 pub iids_tracked: usize,
898 pub iids_added: usize,
900 pub iids_duplicates: usize,
902}
903
904#[derive(Default, Debug, Clone)]
910pub struct StreamInfoConsumersReply {
911 pub consumers: Vec<StreamInfoConsumer>,
913}
914
915#[derive(Default, Debug, Clone)]
923pub struct StreamInfoGroupsReply {
924 pub groups: Vec<StreamInfoGroup>,
926}
927
928#[derive(Default, Debug, Clone)]
933pub struct StreamInfoConsumer {
934 pub name: String,
936 pub pending: usize,
938 pub idle: usize,
940}
941
942#[derive(Default, Debug, Clone)]
947pub struct StreamInfoGroup {
948 pub name: String,
950 pub consumers: usize,
952 pub pending: usize,
954 pub last_delivered_id: String,
956 pub entries_read: Option<usize>,
959 pub lag: Option<usize>,
962}
963
964#[derive(Default, Debug, Clone)]
968pub struct StreamPendingId {
969 pub id: String,
971 pub consumer: String,
975 pub last_delivered_ms: usize,
978 pub times_delivered: usize,
980}
981
982#[derive(Default, Debug, Clone)]
984pub struct StreamKey {
985 pub key: String,
987 pub ids: Vec<StreamId>,
989}
990
991#[derive(Default, Debug, Clone, PartialEq)]
994pub struct StreamId {
995 pub id: String,
997 pub map: HashMap<String, Value>,
999 pub milliseconds_elapsed_from_delivery: Option<usize>,
1001 pub delivered_count: Option<usize>,
1003}
1004
1005impl StreamId {
1006 fn from_array_value(v: Value) -> Result<Self, ParsingError> {
1008 let mut stream_id = Self::default();
1009 if let Value::Array(mut values) = v {
1010 if let Some(v) = values.first_mut() {
1011 stream_id.id = from_redis_value(std::mem::take(v))?;
1012 }
1013 if let Some(v) = values.first_mut() {
1014 stream_id.map = from_redis_value(std::mem::take(v))?;
1015 }
1016 }
1017
1018 Ok(stream_id)
1019 }
1020
1021 pub fn get<T: FromRedisValue>(&self, key: &str) -> Option<T> {
1024 match self.map.get(key) {
1025 Some(x) => from_redis_value_ref(x).ok(),
1026 None => None,
1027 }
1028 }
1029
1030 pub fn contains_key(&self, key: &str) -> bool {
1032 self.map.contains_key(key)
1033 }
1034
1035 pub fn len(&self) -> usize {
1037 self.map.len()
1038 }
1039
1040 pub fn is_empty(&self) -> bool {
1042 self.len() == 0
1043 }
1044}
1045
1046type SACRows = Vec<HashMap<String, HashMap<String, Value>>>;
1047
1048impl FromRedisValue for StreamAutoClaimReply {
1049 fn from_redis_value(v: Value) -> Result<Self, ParsingError> {
1050 let Value::Array(mut items) = v else {
1051 invalid_type_error!("Not a array response", v);
1052 };
1053
1054 if items.len() > 3 || items.len() < 2 {
1055 invalid_type_error!("Incorrect number of items", &items);
1056 }
1057
1058 let deleted_ids = if items.len() == 3 {
1059 from_redis_value(items.pop().unwrap())?
1060 } else {
1061 Vec::new()
1062 };
1063 let claimed = items.pop().unwrap();
1065 let next_stream_id = from_redis_value(items.pop().unwrap())?;
1066
1067 let Value::Array(arr) = &claimed else {
1068 invalid_type_error!("Incorrect type", claimed)
1069 };
1070 let Some(entry) = arr.iter().find(|val| !matches!(val, Value::Nil)) else {
1071 return Ok(Self {
1072 next_stream_id,
1073 claimed: Vec::new(),
1074 deleted_ids,
1075 invalid_entries: !arr.is_empty(),
1076 });
1077 };
1078 let (claimed, invalid_entries) = match entry {
1079 Value::BulkString(_) => {
1080 let claimed_count = arr.len();
1082 let ids: Vec<Option<String>> = from_redis_value(claimed)?;
1083
1084 let claimed: Vec<_> = ids
1085 .into_iter()
1086 .filter_map(|id| {
1087 id.map(|id| StreamId {
1088 id,
1089 ..Default::default()
1090 })
1091 })
1092 .collect();
1093 let invalid_entries = claimed.len() < claimed_count;
1095 (claimed, invalid_entries)
1096 }
1097 Value::Array(_) => {
1098 let claimed_count = arr.len();
1100 let rows: SACRows = from_redis_value(claimed)?;
1101
1102 let claimed: Vec<_> = rows
1103 .into_iter()
1104 .flat_map(|row| {
1105 row.into_iter().map(|(id, map)| StreamId {
1106 id,
1107 map,
1108 milliseconds_elapsed_from_delivery: None,
1109 delivered_count: None,
1110 })
1111 })
1112 .collect();
1113 let invalid_entries = claimed.len() < claimed_count;
1115 (claimed, invalid_entries)
1116 }
1117 _ => invalid_type_error!("Incorrect type", claimed),
1118 };
1119
1120 Ok(Self {
1121 next_stream_id,
1122 claimed,
1123 deleted_ids,
1124 invalid_entries,
1125 })
1126 }
1127}
1128
1129type SRRows = Vec<HashMap<String, Vec<HashMap<String, HashMap<String, Value>>>>>;
1130type SRClaimRows =
1131 Vec<HashMap<String, Vec<(String, HashMap<String, Value>, Option<usize>, Option<usize>)>>>;
1132
1133impl FromRedisValue for StreamReadReply {
1134 fn from_redis_value(v: Value) -> Result<Self, ParsingError> {
1135 if let Ok(rows) = from_redis_value::<SRRows>(v.clone()) {
1137 return Ok(Self::from_standard_rows(rows));
1138 }
1139
1140 if let Ok(rows) = from_redis_value::<SRClaimRows>(v.clone()) {
1143 return Ok(Self::from_claim_rows(rows));
1144 }
1145
1146 invalid_type_error!("Could not parse StreamReadReply in any known format", v)
1147 }
1148}
1149
1150impl StreamReadReply {
1151 fn from_standard_rows(rows: SRRows) -> Self {
1152 let keys = rows
1153 .into_iter()
1154 .flat_map(|row| {
1155 row.into_iter().map(|(key, entries)| StreamKey {
1156 key,
1157 ids: entries
1158 .into_iter()
1159 .flat_map(|id_row| {
1160 id_row.into_iter().map(|(id, map)| StreamId {
1161 id,
1162 map,
1163 milliseconds_elapsed_from_delivery: None,
1164 delivered_count: None,
1165 })
1166 })
1167 .collect(),
1168 })
1169 })
1170 .collect();
1171 Self { keys }
1172 }
1173
1174 fn from_claim_rows(rows: SRClaimRows) -> Self {
1175 let keys = rows
1176 .into_iter()
1177 .flat_map(|row| {
1178 row.into_iter().map(|(key, entries)| StreamKey {
1179 key,
1180 ids: entries
1181 .into_iter()
1182 .map(
1183 |(id, map, milliseconds_elapsed_from_delivery, delivered_count)| {
1184 StreamId {
1185 id,
1186 map,
1187 milliseconds_elapsed_from_delivery,
1188 delivered_count,
1189 }
1190 },
1191 )
1192 .collect(),
1193 })
1194 })
1195 .collect();
1196 Self { keys }
1197 }
1198}
1199
1200impl FromRedisValue for StreamRangeReply {
1201 fn from_redis_value(v: Value) -> Result<Self, ParsingError> {
1202 let rows: Vec<HashMap<String, HashMap<String, Value>>> = from_redis_value(v)?;
1203 let ids: Vec<StreamId> = rows
1204 .into_iter()
1205 .flat_map(|row| {
1206 row.into_iter().map(|(id, map)| StreamId {
1207 id,
1208 map,
1209 milliseconds_elapsed_from_delivery: None,
1210 delivered_count: None,
1211 })
1212 })
1213 .collect();
1214 Ok(Self { ids })
1215 }
1216}
1217
1218impl FromRedisValue for StreamClaimReply {
1219 fn from_redis_value(v: Value) -> Result<Self, ParsingError> {
1220 let rows: Vec<HashMap<String, HashMap<String, Value>>> = from_redis_value(v)?;
1221 let ids: Vec<StreamId> = rows
1222 .into_iter()
1223 .flat_map(|row| {
1224 row.into_iter().map(|(id, map)| StreamId {
1225 id,
1226 map,
1227 milliseconds_elapsed_from_delivery: None,
1228 delivered_count: None,
1229 })
1230 })
1231 .collect();
1232 Ok(Self { ids })
1233 }
1234}
1235
1236type SPRInner = (
1237 usize,
1238 Option<String>,
1239 Option<String>,
1240 Vec<Option<(String, String)>>,
1241);
1242impl FromRedisValue for StreamPendingReply {
1243 fn from_redis_value(v: Value) -> Result<Self, ParsingError> {
1244 let (count, start, end, consumer_data): SPRInner = from_redis_value(v)?;
1245
1246 if count == 0 {
1247 Ok(Self::Empty)
1248 } else {
1249 let mut result = StreamPendingData::default();
1250
1251 let start_id = start.ok_or_else(|| {
1252 ParsingError::from(arcstr::literal!(
1253 "IllegalState: Non-zero pending expects start id"
1254 ))
1255 })?;
1256
1257 let end_id = end.ok_or_else(|| {
1258 ParsingError::from(arcstr::literal!(
1259 "IllegalState: Non-zero pending expects end id"
1260 ))
1261 })?;
1262
1263 result.count = count;
1264 result.start_id = start_id;
1265 result.end_id = end_id;
1266
1267 result.consumers = consumer_data
1268 .into_iter()
1269 .flatten()
1270 .map(|(name, pending)| StreamInfoConsumer {
1271 name,
1272 pending: pending.parse().unwrap_or_default(),
1273 ..Default::default()
1274 })
1275 .collect();
1276
1277 Ok(Self::Data(result))
1278 }
1279 }
1280}
1281
1282impl FromRedisValue for StreamPendingCountReply {
1283 fn from_redis_value(v: Value) -> Result<Self, ParsingError> {
1284 let mut reply = Self::default();
1285 match v {
1286 Value::Array(outer_tuple) => {
1287 for outer in outer_tuple {
1288 match outer {
1289 Value::Array(inner_tuple) => match &inner_tuple[..] {
1290 [
1291 Value::BulkString(id_bytes),
1292 Value::BulkString(consumer_bytes),
1293 Value::Int(last_delivered_ms_u64),
1294 Value::Int(times_delivered_u64),
1295 ] => {
1296 let id = String::from_utf8(id_bytes.to_vec())?;
1297 let consumer = String::from_utf8(consumer_bytes.to_vec())?;
1298 let last_delivered_ms = *last_delivered_ms_u64 as usize;
1299 let times_delivered = *times_delivered_u64 as usize;
1300 reply.ids.push(StreamPendingId {
1301 id,
1302 consumer,
1303 last_delivered_ms,
1304 times_delivered,
1305 });
1306 }
1307 _ => fail!(ParsingError::from(arcstr::literal!(
1308 "Cannot parse redis data (3)"
1309 ))),
1310 },
1311 _ => fail!(ParsingError::from(arcstr::literal!(
1312 "Cannot parse redis data (2)"
1313 ))),
1314 }
1315 }
1316 }
1317 _ => fail!(ParsingError::from(arcstr::literal!(
1318 "Cannot parse redis data (1)"
1319 ))),
1320 }
1321 Ok(reply)
1322 }
1323}
1324
1325impl FromRedisValue for StreamInfoStreamReply {
1326 fn from_redis_value(v: Value) -> Result<Self, ParsingError> {
1327 let mut map: HashMap<String, Value> = from_redis_value(v)?;
1328 let mut reply = Self::default();
1329 if let Some(v) = map.remove("last-generated-id") {
1330 reply.last_generated_id = from_redis_value(v)?;
1331 }
1332 if let Some(v) = map.remove("radix-tree-nodes") {
1333 reply.radix_tree_keys = from_redis_value(v)?;
1334 }
1335 if let Some(v) = map.remove("groups") {
1336 reply.groups = from_redis_value(v)?;
1337 }
1338 if let Some(v) = map.remove("length") {
1339 reply.length = from_redis_value(v)?;
1340 }
1341 if let Some(v) = map.remove("first-entry") {
1342 reply.first_entry = StreamId::from_array_value(v)?;
1343 }
1344 if let Some(v) = map.remove("last-entry") {
1345 reply.last_entry = StreamId::from_array_value(v)?;
1346 }
1347 Ok(reply)
1348 }
1349}
1350
1351impl FromRedisValue for StreamInfoStreamReplyWithIdempotency {
1352 fn from_redis_value(v: Value) -> Result<Self, ParsingError> {
1353 let mut map: HashMap<String, Value> = from_redis_value(v)?;
1354
1355 let mut base = StreamInfoStreamReply::default();
1357 if let Some(v) = map.remove("last-generated-id") {
1358 base.last_generated_id = from_redis_value(v)?;
1359 }
1360 if let Some(v) = map.remove("radix-tree-nodes") {
1361 base.radix_tree_keys = from_redis_value(v)?;
1362 }
1363 if let Some(v) = map.remove("groups") {
1364 base.groups = from_redis_value(v)?;
1365 }
1366 if let Some(v) = map.remove("length") {
1367 base.length = from_redis_value(v)?;
1368 }
1369 if let Some(v) = map.remove("first-entry") {
1370 base.first_entry = StreamId::from_array_value(v)?;
1371 }
1372 if let Some(v) = map.remove("last-entry") {
1373 base.last_entry = StreamId::from_array_value(v)?;
1374 }
1375
1376 let mut reply = Self {
1378 base,
1379 ..Default::default()
1380 };
1381
1382 if let Some(v) = map.remove("idmp-duration") {
1383 reply.idmp_duration = from_redis_value(v)?;
1384 }
1385 if let Some(v) = map.remove("idmp-maxsize") {
1386 reply.idmp_maxsize = from_redis_value(v)?;
1387 }
1388 if let Some(v) = map.remove("pids-tracked") {
1389 reply.pids_tracked = from_redis_value(v)?;
1390 }
1391 if let Some(v) = map.remove("iids-tracked") {
1392 reply.iids_tracked = from_redis_value(v)?;
1393 }
1394 if let Some(v) = map.remove("iids-added") {
1395 reply.iids_added = from_redis_value(v)?;
1396 }
1397 if let Some(v) = map.remove("iids-duplicates") {
1398 reply.iids_duplicates = from_redis_value(v)?;
1399 }
1400
1401 Ok(reply)
1402 }
1403}
1404
1405impl FromRedisValue for StreamInfoConsumersReply {
1406 fn from_redis_value(v: Value) -> Result<Self, ParsingError> {
1407 let consumers: Vec<HashMap<String, Value>> = from_redis_value(v)?;
1408 let mut reply = Self::default();
1409 for mut map in consumers {
1410 let mut c = StreamInfoConsumer::default();
1411 if let Some(v) = map.remove("name") {
1412 c.name = from_redis_value(v)?;
1413 }
1414 if let Some(v) = map.remove("pending") {
1415 c.pending = from_redis_value(v)?;
1416 }
1417 if let Some(v) = map.remove("idle") {
1418 c.idle = from_redis_value(v)?;
1419 }
1420 reply.consumers.push(c);
1421 }
1422
1423 Ok(reply)
1424 }
1425}
1426
1427impl FromRedisValue for StreamInfoGroupsReply {
1428 fn from_redis_value(v: Value) -> Result<Self, ParsingError> {
1429 let groups: Vec<HashMap<String, Value>> = from_redis_value(v)?;
1430 let mut reply = Self::default();
1431 for mut map in groups {
1432 let mut g = StreamInfoGroup::default();
1433 if let Some(v) = map.remove("name") {
1434 g.name = from_redis_value(v)?;
1435 }
1436 if let Some(v) = map.remove("pending") {
1437 g.pending = from_redis_value(v)?;
1438 }
1439 if let Some(v) = map.remove("consumers") {
1440 g.consumers = from_redis_value(v)?;
1441 }
1442 if let Some(v) = map.remove("last-delivered-id") {
1443 g.last_delivered_id = from_redis_value(v)?;
1444 }
1445 if let Some(v) = map.remove("entries-read") {
1446 g.entries_read = if let Value::Nil = v {
1447 None
1448 } else {
1449 Some(from_redis_value(v)?)
1450 };
1451 }
1452 if let Some(v) = map.remove("lag") {
1453 g.lag = if let Value::Nil = v {
1454 None
1455 } else {
1456 Some(from_redis_value(v)?)
1457 };
1458 }
1459 reply.groups.push(g);
1460 }
1461 Ok(reply)
1462 }
1463}
1464
1465#[derive(Debug, Clone, Default)]
1467#[non_exhaustive]
1468pub enum StreamDeletionPolicy {
1469 #[default]
1471 KeepRef,
1472 DelRef,
1474 Acked,
1476}
1477
1478impl ToRedisArgs for StreamDeletionPolicy {
1479 fn write_redis_args<W>(&self, out: &mut W)
1480 where
1481 W: ?Sized + RedisWrite,
1482 {
1483 match self {
1484 Self::KeepRef => out.write_arg(b"KEEPREF"),
1485 Self::DelRef => out.write_arg(b"DELREF"),
1486 Self::Acked => out.write_arg(b"ACKED"),
1487 }
1488 }
1489}
1490impl ToSingleRedisArg for StreamDeletionPolicy {}
1491
1492#[cfg(feature = "streams")]
1494#[cfg_attr(docsrs, doc(cfg(feature = "streams")))]
1495#[derive(Debug, PartialEq, Eq)]
1496#[non_exhaustive]
1497pub enum XDelExStatusCode {
1498 IdNotFound = -1,
1500 Deleted = 1,
1502 NotDeletedUnacknowledgedOrStillReferenced = 2,
1505}
1506
1507#[cfg(feature = "streams")]
1508impl FromRedisValue for XDelExStatusCode {
1509 fn from_redis_value(v: Value) -> Result<Self, ParsingError> {
1510 match v {
1511 Value::Int(code) => match code {
1512 -1 => Ok(Self::IdNotFound),
1513 1 => Ok(Self::Deleted),
1514 2 => Ok(Self::NotDeletedUnacknowledgedOrStillReferenced),
1515 _ => Err(format!("Invalid XDelExStatusCode status code: {code}").into()),
1516 },
1517 _ => Err(arcstr::literal!("Response type not XAckDelStatusCode compatible").into()),
1518 }
1519 }
1520}
1521
1522#[cfg(feature = "streams")]
1530#[cfg_attr(docsrs, doc(cfg(feature = "streams")))]
1531#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1532#[non_exhaustive]
1533pub enum StreamNackMode {
1534 Silent,
1538 Fail,
1541 Fatal,
1546}
1547
1548#[cfg(feature = "streams")]
1549impl ToRedisArgs for StreamNackMode {
1550 fn write_redis_args<W>(&self, out: &mut W)
1551 where
1552 W: ?Sized + RedisWrite,
1553 {
1554 match self {
1555 Self::Silent => out.write_arg(b"SILENT"),
1556 Self::Fail => out.write_arg(b"FAIL"),
1557 Self::Fatal => out.write_arg(b"FATAL"),
1558 }
1559 }
1560}
1561
1562#[cfg(feature = "streams")]
1563impl ToSingleRedisArg for StreamNackMode {}
1564
1565#[cfg(feature = "streams")]
1578#[cfg_attr(docsrs, doc(cfg(feature = "streams")))]
1579#[derive(Debug)]
1580pub struct StreamNackOptions {
1581 mode: StreamNackMode,
1583}
1584
1585#[cfg(feature = "streams")]
1586impl StreamNackOptions {
1587 pub fn new(mode: StreamNackMode) -> Self {
1589 Self { mode }
1590 }
1591}
1592
1593#[cfg(feature = "streams")]
1594impl ToRedisArgs for StreamNackOptions {
1595 fn write_redis_args<W>(&self, out: &mut W)
1596 where
1597 W: ?Sized + RedisWrite,
1598 {
1599 self.mode.write_redis_args(out);
1600 }
1601}
1602
1603#[cfg(feature = "streams")]
1605#[cfg_attr(docsrs, doc(cfg(feature = "streams")))]
1606#[derive(Debug, PartialEq, Eq)]
1607#[non_exhaustive]
1608pub enum XAckDelStatusCode {
1609 IdNotFound = -1,
1611 AcknowledgedAndDeleted = 1,
1613 AcknowledgedNotDeletedStillReferenced = 2,
1615}
1616
1617#[cfg(feature = "streams")]
1618impl FromRedisValue for XAckDelStatusCode {
1619 fn from_redis_value(v: Value) -> Result<Self, ParsingError> {
1620 match v {
1621 Value::Int(code) => match code {
1622 -1 => Ok(Self::IdNotFound),
1623 1 => Ok(Self::AcknowledgedAndDeleted),
1624 2 => Ok(Self::AcknowledgedNotDeletedStillReferenced),
1625 _ => Err(arcstr::literal!("Invalid XAckDelStatusCode status code: {code}").into()),
1626 },
1627 _ => Err(arcstr::literal!("Response type not XAckDelStatusCode compatible").into()),
1628 }
1629 }
1630}
1631
1632#[cfg(test)]
1633mod tests {
1634 use super::*;
1635
1636 fn assert_command_eq(object: impl ToRedisArgs, expected: &[u8]) {
1637 let mut out: Vec<Vec<u8>> = Vec::new();
1638
1639 object.write_redis_args(&mut out);
1640
1641 let mut cmd: Vec<u8> = Vec::new();
1642
1643 out.iter_mut().for_each(|item| {
1644 cmd.append(item);
1645 cmd.push(b' ');
1646 });
1647
1648 cmd.pop();
1649
1650 assert_eq!(cmd, expected);
1651 }
1652
1653 mod stream_auto_claim_reply {
1654 use super::*;
1655 use crate::Value;
1656
1657 #[test]
1658 fn short_response() {
1659 let value = Value::Array(vec![Value::BulkString("1713465536578-0".into())]);
1660
1661 StreamAutoClaimReply::from_redis_value(value).unwrap_err();
1662 }
1663
1664 #[test]
1665 fn parses_none_claimed_response() {
1666 let value = Value::Array(vec![
1667 Value::BulkString("0-0".into()),
1668 Value::Array(vec![]),
1669 Value::Array(vec![]),
1670 ]);
1671
1672 let reply: StreamAutoClaimReply = FromRedisValue::from_redis_value(value).unwrap();
1673
1674 assert_eq!(reply.next_stream_id.as_str(), "0-0");
1675 assert_eq!(reply.claimed.len(), 0);
1676 assert_eq!(reply.deleted_ids.len(), 0);
1677 }
1678
1679 #[test]
1680 fn parses_response() {
1681 let value = Value::Array(vec![
1682 Value::BulkString("1713465536578-0".into()),
1683 Value::Array(vec![
1684 Value::Array(vec![
1685 Value::BulkString("1713465533411-0".into()),
1686 Value::Array(vec![
1688 Value::BulkString("name".into()),
1689 Value::BulkString("test".into()),
1690 Value::BulkString("other".into()),
1691 Value::BulkString("whaterver".into()),
1692 ]),
1693 ]),
1694 Value::Array(vec![
1695 Value::BulkString("1713465536069-0".into()),
1696 Value::Array(vec![
1697 Value::BulkString("name".into()),
1698 Value::BulkString("another test".into()),
1699 Value::BulkString("other".into()),
1700 Value::BulkString("something".into()),
1701 ]),
1702 ]),
1703 ]),
1704 Value::Array(vec![Value::BulkString("123456789-0".into())]),
1705 ]);
1706
1707 let reply: StreamAutoClaimReply = FromRedisValue::from_redis_value(value).unwrap();
1708
1709 assert_eq!(reply.next_stream_id.as_str(), "1713465536578-0");
1710 assert_eq!(reply.claimed.len(), 2);
1711 assert_eq!(reply.claimed[0].id.as_str(), "1713465533411-0");
1712 assert!(
1713 matches!(reply.claimed[0].map.get("name"), Some(Value::BulkString(v)) if v == "test".as_bytes())
1714 );
1715 assert_eq!(reply.claimed[1].id.as_str(), "1713465536069-0");
1716 assert_eq!(reply.deleted_ids.len(), 1);
1717 assert!(reply.deleted_ids.contains(&"123456789-0".to_string()));
1718 }
1719
1720 #[test]
1721 fn parses_v6_response() {
1722 let value = Value::Array(vec![
1723 Value::BulkString("1713465536578-0".into()),
1724 Value::Array(vec![
1725 Value::Array(vec![
1726 Value::BulkString("1713465533411-0".into()),
1727 Value::Array(vec![
1728 Value::BulkString("name".into()),
1729 Value::BulkString("test".into()),
1730 Value::BulkString("other".into()),
1731 Value::BulkString("whaterver".into()),
1732 ]),
1733 ]),
1734 Value::Array(vec![
1735 Value::BulkString("1713465536069-0".into()),
1736 Value::Array(vec![
1737 Value::BulkString("name".into()),
1738 Value::BulkString("another test".into()),
1739 Value::BulkString("other".into()),
1740 Value::BulkString("something".into()),
1741 ]),
1742 ]),
1743 ]),
1744 ]);
1746
1747 let reply: StreamAutoClaimReply = FromRedisValue::from_redis_value(value).unwrap();
1748
1749 assert_eq!(reply.next_stream_id.as_str(), "1713465536578-0");
1750 assert_eq!(reply.claimed.len(), 2);
1751 let ids: Vec<_> = reply.claimed.iter().map(|e| e.id.as_str()).collect();
1752 assert!(ids.contains(&"1713465533411-0"));
1753 assert!(ids.contains(&"1713465536069-0"));
1754 assert_eq!(reply.deleted_ids.len(), 0);
1755 }
1756
1757 #[test]
1758 fn parses_justid_response() {
1759 let value = Value::Array(vec![
1760 Value::BulkString("1713465536578-0".into()),
1761 Value::Array(vec![
1762 Value::BulkString("1713465533411-0".into()),
1763 Value::BulkString("1713465536069-0".into()),
1764 ]),
1765 Value::Array(vec![Value::BulkString("123456789-0".into())]),
1766 ]);
1767
1768 let reply: StreamAutoClaimReply = FromRedisValue::from_redis_value(value).unwrap();
1769
1770 assert_eq!(reply.next_stream_id.as_str(), "1713465536578-0");
1771 assert_eq!(reply.claimed.len(), 2);
1772 let ids: Vec<_> = reply.claimed.iter().map(|e| e.id.as_str()).collect();
1773 assert!(ids.contains(&"1713465533411-0"));
1774 assert!(ids.contains(&"1713465536069-0"));
1775 assert_eq!(reply.deleted_ids.len(), 1);
1776 assert!(reply.deleted_ids.contains(&"123456789-0".to_string()));
1777 }
1778
1779 #[test]
1780 fn parses_v6_justid_response() {
1781 let value = Value::Array(vec![
1782 Value::BulkString("1713465536578-0".into()),
1783 Value::Array(vec![
1784 Value::BulkString("1713465533411-0".into()),
1785 Value::BulkString("1713465536069-0".into()),
1786 ]),
1787 ]);
1789
1790 let reply: StreamAutoClaimReply = FromRedisValue::from_redis_value(value).unwrap();
1791
1792 assert_eq!(reply.next_stream_id.as_str(), "1713465536578-0");
1793 assert_eq!(reply.claimed.len(), 2);
1794 let ids: Vec<_> = reply.claimed.iter().map(|e| e.id.as_str()).collect();
1795 assert!(ids.contains(&"1713465533411-0"));
1796 assert!(ids.contains(&"1713465536069-0"));
1797 assert_eq!(reply.deleted_ids.len(), 0);
1798 }
1799 }
1800
1801 mod stream_trim_options {
1802 use super::*;
1803
1804 #[test]
1805 fn maxlen_trim() {
1806 let options = StreamTrimOptions::maxlen(StreamTrimmingMode::Approx, 10);
1807
1808 assert_command_eq(options, b"MAXLEN ~ 10");
1809 }
1810
1811 #[test]
1812 fn maxlen_exact_trim() {
1813 let options = StreamTrimOptions::maxlen(StreamTrimmingMode::Exact, 10);
1814
1815 assert_command_eq(options, b"MAXLEN = 10");
1816 }
1817
1818 #[test]
1819 fn maxlen_trim_limit() {
1820 let options = StreamTrimOptions::maxlen(StreamTrimmingMode::Approx, 10).limit(5);
1821
1822 assert_command_eq(options, b"MAXLEN ~ 10 LIMIT 5");
1823 }
1824 #[test]
1825 fn minid_trim_limit() {
1826 let options = StreamTrimOptions::minid(StreamTrimmingMode::Exact, "123456-7").limit(5);
1827
1828 assert_command_eq(options, b"MINID = 123456-7 LIMIT 5");
1829 }
1830 }
1831
1832 mod stream_add_options {
1833 use super::*;
1834
1835 #[test]
1836 fn the_default() {
1837 let options = StreamAddOptions::default();
1838
1839 assert_command_eq(options, b"");
1840 }
1841
1842 #[test]
1843 fn with_maxlen_trim() {
1844 let options = StreamAddOptions::default()
1845 .trim(StreamTrimStrategy::maxlen(StreamTrimmingMode::Exact, 10));
1846
1847 assert_command_eq(options, b"MAXLEN = 10");
1848 }
1849
1850 #[test]
1851 fn with_nomkstream() {
1852 let options = StreamAddOptions::default().nomkstream();
1853
1854 assert_command_eq(options, b"NOMKSTREAM");
1855 }
1856
1857 #[test]
1858 fn with_nomkstream_and_maxlen_trim() {
1859 let options = StreamAddOptions::default()
1860 .nomkstream()
1861 .trim(StreamTrimStrategy::maxlen(StreamTrimmingMode::Exact, 10));
1862
1863 assert_command_eq(options, b"NOMKSTREAM MAXLEN = 10");
1864 }
1865
1866 #[test]
1867 fn with_idmp_manual_mode() {
1868 let options = StreamAddOptions::default().idmp("producer-1", "iid-1");
1869
1870 assert_command_eq(options, b"IDMP producer-1 iid-1");
1871 }
1872
1873 #[test]
1874 fn with_idmpauto_automatic_mode() {
1875 let options = StreamAddOptions::default().idmpauto("producer-1");
1876
1877 assert_command_eq(options, b"IDMPAUTO producer-1");
1878 }
1879
1880 #[test]
1881 fn with_nomkstream_and_idmp() {
1882 let options = StreamAddOptions::default()
1883 .nomkstream()
1884 .idmp("producer-1", "iid-1");
1885
1886 assert_command_eq(options, b"NOMKSTREAM IDMP producer-1 iid-1");
1887 }
1888
1889 #[test]
1890 fn with_trim_and_idmp() {
1891 let options = StreamAddOptions::default()
1892 .trim(StreamTrimStrategy::maxlen(StreamTrimmingMode::Exact, 100))
1893 .idmp("producer-1", "iid-1");
1894
1895 assert_command_eq(options, b"IDMP producer-1 iid-1 MAXLEN = 100");
1896 }
1897
1898 #[test]
1899 fn with_all_options_and_idmp() {
1900 let options = StreamAddOptions::default()
1901 .nomkstream()
1902 .trim(StreamTrimStrategy::maxlen(StreamTrimmingMode::Approx, 100))
1903 .idmp("producer-1", "iid-1")
1904 .set_deletion_policy(StreamDeletionPolicy::KeepRef);
1905
1906 assert_command_eq(
1907 options,
1908 b"NOMKSTREAM KEEPREF IDMP producer-1 iid-1 MAXLEN ~ 100",
1909 );
1910 }
1911
1912 #[test]
1913 fn with_all_options_and_idmpauto() {
1914 let options = StreamAddOptions::default()
1915 .nomkstream()
1916 .trim(StreamTrimStrategy::minid(
1917 StreamTrimmingMode::Exact,
1918 "123456-0",
1919 ))
1920 .idmpauto("producer-2")
1921 .set_deletion_policy(StreamDeletionPolicy::DelRef);
1922
1923 assert_command_eq(
1924 options,
1925 b"NOMKSTREAM DELREF IDMPAUTO producer-2 MINID = 123456-0",
1926 );
1927 }
1928 }
1929
1930 mod stream_config_options {
1931 use super::*;
1932
1933 const IDMP_CUSTOM_DURATION: u32 = 300;
1934 const IDMP_CUSTOM_MAXSIZE: u16 = 1000;
1935
1936 #[test]
1937 fn with_idempotency_seconds_only() {
1938 let options =
1939 StreamConfigOptions::with_idempotency_seconds(IDMP_CUSTOM_DURATION).unwrap();
1940 assert_command_eq(
1941 options,
1942 format!("IDMP-DURATION {IDMP_CUSTOM_DURATION}").as_bytes(),
1943 );
1944 }
1945
1946 #[test]
1947 fn with_idempotency_maxsize_only() {
1948 let options =
1949 StreamConfigOptions::with_idempotency_maxsize(IDMP_CUSTOM_MAXSIZE).unwrap();
1950 assert_command_eq(
1951 options,
1952 format!("IDMP-MAXSIZE {IDMP_CUSTOM_MAXSIZE}").as_bytes(),
1953 );
1954 }
1955
1956 #[test]
1957 fn with_both_options_starting_with_idempotency_seconds() {
1958 let options = StreamConfigOptions::with_idempotency_seconds(IDMP_CUSTOM_DURATION)
1959 .unwrap()
1960 .idempotency_maxsize(IDMP_CUSTOM_MAXSIZE)
1961 .unwrap();
1962 assert_command_eq(
1963 options,
1964 format!("IDMP-DURATION {IDMP_CUSTOM_DURATION} IDMP-MAXSIZE {IDMP_CUSTOM_MAXSIZE}")
1965 .as_bytes(),
1966 );
1967 }
1968
1969 #[test]
1970 fn with_both_options_starting_with_idempotency_maxsize() {
1971 let options = StreamConfigOptions::with_idempotency_maxsize(IDMP_CUSTOM_MAXSIZE)
1972 .unwrap()
1973 .idempotency_seconds(IDMP_CUSTOM_DURATION)
1974 .unwrap();
1975 assert_command_eq(
1976 options,
1977 format!("IDMP-DURATION {IDMP_CUSTOM_DURATION} IDMP-MAXSIZE {IDMP_CUSTOM_MAXSIZE}")
1978 .as_bytes(),
1979 );
1980 }
1981
1982 #[test]
1983 fn with_max_values() {
1984 let options = StreamConfigOptions::with_idempotency_seconds(IDMP_DURATION_MAX)
1985 .unwrap()
1986 .idempotency_maxsize(IDMP_MAXSIZE_MAX)
1987 .unwrap();
1988 assert_command_eq(
1989 options,
1990 format!("IDMP-DURATION {IDMP_DURATION_MAX} IDMP-MAXSIZE {IDMP_MAXSIZE_MAX}")
1991 .as_bytes(),
1992 );
1993 }
1994
1995 #[test]
1996 fn with_min_values() {
1997 let options = StreamConfigOptions::with_idempotency_seconds(IDMP_DURATION_MIN)
1998 .unwrap()
1999 .idempotency_maxsize(IDMP_MAXSIZE_MIN)
2000 .unwrap();
2001 assert_command_eq(
2002 options,
2003 format!("IDMP-DURATION {IDMP_DURATION_MIN} IDMP-MAXSIZE {IDMP_MAXSIZE_MIN}")
2004 .as_bytes(),
2005 );
2006 }
2007
2008 #[test]
2009 fn error_idempotency_seconds_too_low() {
2010 let result = StreamConfigOptions::with_idempotency_seconds(IDMP_DURATION_MIN - 1);
2011 assert!(result.is_err());
2012 assert!(result.unwrap_err().contains(&format!(
2013 "IDMP-DURATION must be between {IDMP_DURATION_MIN} and {IDMP_DURATION_MAX}"
2014 )));
2015 }
2016
2017 #[test]
2018 fn error_idempotency_seconds_too_high() {
2019 let result = StreamConfigOptions::with_idempotency_seconds(IDMP_DURATION_MAX + 1);
2020 assert!(result.is_err());
2021 assert!(result.unwrap_err().contains(&format!(
2022 "IDMP-DURATION must be between {IDMP_DURATION_MIN} and {IDMP_DURATION_MAX}"
2023 )));
2024 }
2025
2026 #[test]
2027 fn error_idempotency_maxsize_too_low() {
2028 let result = StreamConfigOptions::with_idempotency_maxsize(IDMP_MAXSIZE_MIN - 1);
2029 assert!(result.is_err());
2030 assert!(result.unwrap_err().contains(&format!(
2031 "IDMP-MAXSIZE must be between {IDMP_MAXSIZE_MIN} and {IDMP_MAXSIZE_MAX}"
2032 )));
2033 }
2034
2035 #[test]
2036 fn error_idempotency_maxsize_too_high() {
2037 let result = StreamConfigOptions::with_idempotency_maxsize(IDMP_MAXSIZE_MAX + 1);
2038 assert!(result.is_err());
2039 assert!(result.unwrap_err().contains(&format!(
2040 "IDMP-MAXSIZE must be between {IDMP_MAXSIZE_MIN} and {IDMP_MAXSIZE_MAX}"
2041 )));
2042 }
2043
2044 #[test]
2045 fn error_setter_idempotency_seconds_too_low() {
2046 let result = StreamConfigOptions::with_idempotency_maxsize(IDMP_CUSTOM_MAXSIZE)
2047 .unwrap()
2048 .idempotency_seconds(IDMP_DURATION_MIN - 1);
2049 assert!(result.is_err());
2050 assert!(result.unwrap_err().contains(&format!(
2051 "IDMP-DURATION must be between {IDMP_DURATION_MIN} and {IDMP_DURATION_MAX}"
2052 )));
2053 }
2054
2055 #[test]
2056 fn error_setter_idempotency_maxsize_too_high() {
2057 let result = StreamConfigOptions::with_idempotency_seconds(IDMP_CUSTOM_DURATION)
2058 .unwrap()
2059 .idempotency_maxsize(IDMP_MAXSIZE_MAX + 1);
2060 assert!(result.is_err());
2061 assert!(result.unwrap_err().contains(&format!(
2062 "IDMP-MAXSIZE must be between {IDMP_MAXSIZE_MIN} and {IDMP_MAXSIZE_MAX}"
2063 )));
2064 }
2065 }
2066}