1use std::collections::HashMap;
5
6use std::time::{Duration, SystemTime};
7
8use crate::client::{Client, duration_millis, duration_secs, join, unexpected, unix_millis};
9
10use crate::error::Error;
11use crate::models::{
12 FieldExpireCondition, HashScanPage, PubSubMessage, ScanPage, SortedSetAddOptions,
13 SortedSetEntry, StreamClaimOptions, StreamEntry, StreamPendingEntry, StreamPendingFilter,
14 StreamReadOptions, StreamReadResult,
15};
16
17use crate::value::RespValue;
18
19macro_rules! args {
20 ($($argument:expr),* $(,)?) => {
21 vec![$($argument.to_string()),*]
22 };
23}
24
25impl Client {
26 pub fn mget(&mut self, keys: &[&str]) -> Result<Vec<Option<String>>, Error> {
28 self.optional_strings(join("MGET", keys))
29 }
30
31 pub fn limit(&mut self, key: &str, max: u64, window_ms: u64) -> Result<Option<i64>, Error> {
34 let value = self.run(["LIMIT", key, &max.to_string(), &window_ms.to_string()])?;
35
36 if value.is_null() {
37 return Ok(None);
38 }
39
40 value
41 .as_integer()
42 .map(Some)
43 .ok_or_else(|| unexpected("LIMIT", &value))
44 }
45
46 pub fn once(&mut self, key: &str, ttl_ms: u64) -> Result<i64, Error> {
48 self.integer(["ONCE", key, &ttl_ms.to_string()])
49 }
50
51 pub fn idempotency_begin(
53 &mut self,
54 key: &str,
55 fingerprint: &str,
56 owner: &str,
57 ttl_ms: u64,
58 ) -> Result<RespValue, Error> {
59 self.run([
60 "IDEM",
61 key,
62 "BEGIN",
63 fingerprint,
64 owner,
65 &ttl_ms.to_string(),
66 ])
67 }
68
69 pub fn idempotency_complete(
71 &mut self,
72 key: &str,
73 fingerprint: &str,
74 owner: &str,
75 result: &str,
76 ) -> Result<i64, Error> {
77 self.integer(["IDEM", key, "COMPLETE", fingerprint, owner, result])
78 }
79
80 pub fn take(&mut self, key: &str, count: i64) -> Result<Option<i64>, Error> {
82 let value = self.run(["TAKE", key, &count.to_string()])?;
83
84 if value.is_null() {
85 return Ok(None);
86 }
87
88 value
89 .as_integer()
90 .map(Some)
91 .ok_or_else(|| unexpected("TAKE", &value))
92 }
93
94 pub fn getex(&mut self, key: &str, ttl_ms: u64) -> Result<Option<String>, Error> {
96 self.bulk_or_null(["GETEX", key, &ttl_ms.to_string()])
97 }
98
99 pub fn lease(&mut self, key: &str, ttl_ms: u64) -> Result<Option<i64>, Error> {
101 let value = self.run(["LEASE", key, &ttl_ms.to_string()])?;
102
103 if value.is_null() {
104 return Ok(None);
105 }
106
107 value
108 .as_integer()
109 .map(Some)
110 .ok_or_else(|| unexpected("LEASE", &value))
111 }
112
113 pub fn semaphore(&mut self, key: &str, max: u64, ttl_ms: u64) -> Result<Option<i64>, Error> {
116 let value = self.run(["SEMAPHORE", key, &max.to_string(), &ttl_ms.to_string()])?;
117
118 if value.is_null() {
119 return Ok(None);
120 }
121
122 value
123 .as_integer()
124 .map(Some)
125 .ok_or_else(|| unexpected("SEMAPHORE", &value))
126 }
127
128 pub fn incr_by_max(&mut self, key: &str, delta: i64, limit: i64) -> Result<Option<i64>, Error> {
131 let value = self.run(args!["INCRBY", key, delta, "MAX", limit])?;
132
133 if value.is_null() {
134 return Ok(None);
135 }
136
137 value
138 .as_integer()
139 .map(Some)
140 .ok_or_else(|| unexpected("INCRBY", &value))
141 }
142
143 pub fn release(&mut self, key: &str, token: u64) -> Result<i64, Error> {
145 self.integer(["RELEASE", key, &token.to_string()])
146 }
147
148 pub fn setv(&mut self, key: &str, value: &str, version: Option<u64>) -> Result<i64, Error> {
150 let mut arguments = vec!["SETV".to_owned(), key.to_owned(), value.to_owned()];
151
152 if let Some(version) = version {
153 arguments.push(version.to_string());
154 }
155
156 self.integer(arguments)
157 }
158
159 pub fn changes_start(&mut self) -> Result<i64, Error> {
161 self.integer(["CHANGES", "START"])
162 }
163
164 pub fn changes(&mut self, cursor: u64) -> Result<Vec<(i64, String)>, Error> {
166 let value = self.run(["CHANGES", &cursor.to_string()])?;
167 let Some(rows) = value.as_array() else {
168 return Err(unexpected("CHANGES", &value));
169 };
170 let mut notes = Vec::new();
171
172 for row in rows {
173 let Some(parts) = row.as_array() else {
174 return Err(unexpected("CHANGES", row));
175 };
176 let sequence = parts
177 .first()
178 .and_then(|part| part.as_integer())
179 .ok_or_else(|| unexpected("CHANGES", row))?;
180 let text = parts
181 .get(1)
182 .and_then(|part| part.as_string().ok().flatten())
183 .unwrap_or_default();
184
185 notes.push((sequence, text));
186 }
187
188 Ok(notes)
189 }
190
191 pub fn mset(&mut self, pairs: &[(&str, &str)]) -> Result<(), Error> {
193 self.ok(with_pairs(args!["MSET"], pairs))
194 }
195
196 pub fn getset(&mut self, key: &str, value: &str) -> Result<Option<String>, Error> {
198 self.bulk_or_null(["GETSET", key, value])
199 }
200
201 pub fn append(&mut self, key: &str, value: &str) -> Result<i64, Error> {
203 self.integer(["APPEND", key, value])
204 }
205
206 pub fn strlen(&mut self, key: &str) -> Result<i64, Error> {
208 self.integer(["STRLEN", key])
209 }
210
211 pub fn decr(&mut self, key: &str) -> Result<i64, Error> {
213 self.integer(["DECR", key])
214 }
215
216 pub fn incr_by(&mut self, key: &str, delta: i64) -> Result<i64, Error> {
218 self.integer(args!["INCRBY", key, delta])
219 }
220
221 pub fn decr_by(&mut self, key: &str, delta: i64) -> Result<i64, Error> {
223 self.integer(args!["DECRBY", key, delta])
224 }
225
226 pub fn unlink_key(&mut self, key: &str) -> Result<i64, Error> {
228 self.unlink(&[key])
229 }
230
231 pub fn unlink(&mut self, keys: &[&str]) -> Result<i64, Error> {
233 self.integer(join("UNLINK", keys))
234 }
235
236 pub fn exists_key(&mut self, key: &str) -> Result<i64, Error> {
238 self.exists(&[key])
239 }
240
241 pub fn exists(&mut self, keys: &[&str]) -> Result<i64, Error> {
243 self.integer(join("EXISTS", keys))
244 }
245
246 pub fn key_type(&mut self, key: &str) -> Result<String, Error> {
248 self.text(["TYPE", key])
249 }
250
251 pub fn rename(&mut self, key: &str, new_key: &str) -> Result<(), Error> {
253 self.ok(["RENAME", key, new_key])
254 }
255
256 pub fn scan(
258 &mut self,
259 cursor: u64,
260 pattern: Option<&str>,
261 count: Option<u64>,
262 ) -> Result<ScanPage, Error> {
263 let mut arguments = args!["SCAN", cursor];
264
265 push_scan_options(&mut arguments, pattern, count);
266
267 let (cursor, keys) = parse_scan("SCAN", self.run(arguments)?)?;
268
269 Ok(ScanPage { cursor, keys })
270 }
271
272 pub fn dbsize(&mut self) -> Result<i64, Error> {
274 self.integer(["DBSIZE"])
275 }
276
277 pub fn expire_at(&mut self, key: &str, when: SystemTime) -> Result<bool, Error> {
279 let millis = unix_millis(when)?;
280
281 if millis % 1000 == 0 {
282 return self.flag(args!["EXPIREAT", key, millis / 1000]);
283 }
284
285 self.flag(args!["PEXPIREAT", key, millis])
286 }
287
288 pub fn pexpire(&mut self, key: &str, ttl: Duration) -> Result<bool, Error> {
290 self.flag(args!["PEXPIRE", key, duration_millis(ttl)])
291 }
292
293 pub fn ttl(&mut self, key: &str) -> Result<i64, Error> {
295 self.integer(["TTL", key])
296 }
297
298 pub fn pttl(&mut self, key: &str) -> Result<i64, Error> {
300 self.integer(["PTTL", key])
301 }
302
303 pub fn persist(&mut self, key: &str) -> Result<bool, Error> {
305 self.flag(["PERSIST", key])
306 }
307
308 pub fn lpush(&mut self, key: &str, values: &[&str]) -> Result<i64, Error> {
310 self.integer(with(args!["LPUSH", key], values))
311 }
312
313 pub fn rpush(&mut self, key: &str, values: &[&str]) -> Result<i64, Error> {
315 self.integer(with(args!["RPUSH", key], values))
316 }
317
318 pub fn lpop(&mut self, key: &str) -> Result<Option<String>, Error> {
320 self.bulk_or_null(["LPOP", key])
321 }
322
323 pub fn lpop_count(&mut self, key: &str, count: u64) -> Result<Vec<String>, Error> {
325 self.strings(args!["LPOP", key, count])
326 }
327
328 pub fn rpop(&mut self, key: &str) -> Result<Option<String>, Error> {
330 self.bulk_or_null(["RPOP", key])
331 }
332
333 pub fn rpop_count(&mut self, key: &str, count: u64) -> Result<Vec<String>, Error> {
335 self.strings(args!["RPOP", key, count])
336 }
337
338 pub fn llen(&mut self, key: &str) -> Result<i64, Error> {
340 self.integer(["LLEN", key])
341 }
342
343 pub fn lindex(&mut self, key: &str, index: i64) -> Result<Option<String>, Error> {
345 self.bulk_or_null(args!["LINDEX", key, index])
346 }
347
348 pub fn lrange(&mut self, key: &str, start: i64, stop: i64) -> Result<Vec<String>, Error> {
350 self.strings(args!["LRANGE", key, start, stop])
351 }
352
353 pub fn ltrim(&mut self, key: &str, start: i64, stop: i64) -> Result<(), Error> {
355 self.ok(args!["LTRIM", key, start, stop])
356 }
357
358 pub fn blpop_key(
361 &mut self,
362 key: &str,
363 timeout: Duration,
364 ) -> Result<Option<(String, String)>, Error> {
365 self.blpop(&[key], timeout)
366 }
367
368 pub fn blpop(
371 &mut self,
372 keys: &[&str],
373 timeout: Duration,
374 ) -> Result<Option<(String, String)>, Error> {
375 self.blocking_pop("BLPOP", keys, timeout)
376 }
377
378 pub fn brpop_key(
381 &mut self,
382 key: &str,
383 timeout: Duration,
384 ) -> Result<Option<(String, String)>, Error> {
385 self.brpop(&[key], timeout)
386 }
387
388 pub fn brpop(
390 &mut self,
391 keys: &[&str],
392 timeout: Duration,
393 ) -> Result<Option<(String, String)>, Error> {
394 self.blocking_pop("BRPOP", keys, timeout)
395 }
396
397 pub fn sadd(&mut self, key: &str, members: &[&str]) -> Result<i64, Error> {
399 self.integer(with(args!["SADD", key], members))
400 }
401
402 pub fn srem(&mut self, key: &str, members: &[&str]) -> Result<i64, Error> {
404 self.integer(with(args!["SREM", key], members))
405 }
406
407 pub fn sismember(&mut self, key: &str, member: &str) -> Result<bool, Error> {
409 self.flag(["SISMEMBER", key, member])
410 }
411
412 pub fn scard(&mut self, key: &str) -> Result<i64, Error> {
414 self.integer(["SCARD", key])
415 }
416
417 pub fn smembers(&mut self, key: &str) -> Result<Vec<String>, Error> {
419 self.strings(["SMEMBERS", key])
420 }
421
422 pub fn sinter(&mut self, keys: &[&str]) -> Result<Vec<String>, Error> {
424 self.strings(join("SINTER", keys))
425 }
426
427 pub fn sunion(&mut self, keys: &[&str]) -> Result<Vec<String>, Error> {
429 self.strings(join("SUNION", keys))
430 }
431
432 pub fn sdiff(&mut self, keys: &[&str]) -> Result<Vec<String>, Error> {
434 self.strings(join("SDIFF", keys))
435 }
436
437 pub fn sinterstore(&mut self, destination: &str, keys: &[&str]) -> Result<i64, Error> {
439 self.integer(with(args!["SINTERSTORE", destination], keys))
440 }
441
442 pub fn sunionstore(&mut self, destination: &str, keys: &[&str]) -> Result<i64, Error> {
444 self.integer(with(args!["SUNIONSTORE", destination], keys))
445 }
446
447 pub fn sdiffstore(&mut self, destination: &str, keys: &[&str]) -> Result<i64, Error> {
449 self.integer(with(args!["SDIFFSTORE", destination], keys))
450 }
451
452 pub fn smove(&mut self, source: &str, destination: &str, member: &str) -> Result<bool, Error> {
454 self.flag(["SMOVE", source, destination, member])
455 }
456
457 pub fn spop(&mut self, key: &str) -> Result<Option<String>, Error> {
459 self.bulk_or_null(["SPOP", key])
460 }
461
462 pub fn spop_count(&mut self, key: &str, count: u64) -> Result<Vec<String>, Error> {
464 self.strings(args!["SPOP", key, count])
465 }
466
467 pub fn srandmember(&mut self, key: &str) -> Result<Option<String>, Error> {
469 self.bulk_or_null(["SRANDMEMBER", key])
470 }
471
472 pub fn srandmember_count(&mut self, key: &str, count: i64) -> Result<Vec<String>, Error> {
475 self.strings(args!["SRANDMEMBER", key, count])
476 }
477
478 pub fn hset(&mut self, key: &str, fields: &[(&str, &str)]) -> Result<i64, Error> {
480 self.integer(with_pairs(args!["HSET", key], fields))
481 }
482
483 pub fn hget(&mut self, key: &str, field: &str) -> Result<Option<String>, Error> {
485 self.bulk_or_null(["HGET", key, field])
486 }
487
488 pub fn hdel(&mut self, key: &str, fields: &[&str]) -> Result<i64, Error> {
490 self.integer(with(args!["HDEL", key], fields))
491 }
492
493 pub fn hlen(&mut self, key: &str) -> Result<i64, Error> {
495 self.integer(["HLEN", key])
496 }
497
498 pub fn hgetall(&mut self, key: &str) -> Result<HashMap<String, String>, Error> {
500 let pairs = pairs(&self.run(["HGETALL", key])?)?;
501
502 Ok(pairs.into_iter().collect())
503 }
504
505 pub fn hmget(&mut self, key: &str, fields: &[&str]) -> Result<Vec<Option<String>>, Error> {
507 self.optional_strings(with(args!["HMGET", key], fields))
508 }
509
510 pub fn hexists(&mut self, key: &str, field: &str) -> Result<bool, Error> {
512 self.flag(["HEXISTS", key, field])
513 }
514
515 pub fn hkeys(&mut self, key: &str) -> Result<Vec<String>, Error> {
517 self.strings(["HKEYS", key])
518 }
519
520 pub fn hvals(&mut self, key: &str) -> Result<Vec<String>, Error> {
522 self.strings(["HVALS", key])
523 }
524
525 pub fn hincr_by(&mut self, key: &str, field: &str, delta: i64) -> Result<i64, Error> {
527 self.integer(args!["HINCRBY", key, field, delta])
528 }
529
530 pub fn hsetnx(&mut self, key: &str, field: &str, value: &str) -> Result<bool, Error> {
532 self.flag(["HSETNX", key, field, value])
533 }
534
535 pub fn hstrlen(&mut self, key: &str, field: &str) -> Result<i64, Error> {
537 self.integer(["HSTRLEN", key, field])
538 }
539
540 pub fn hexpire(
552 &mut self,
553 key: &str,
554 ttl: Duration,
555 fields: &[&str],
556 ) -> Result<Vec<i64>, Error> {
557 self.hexpire_with(key, ttl, FieldExpireCondition::Always, fields)
558 }
559
560 pub fn hexpire_with(
562 &mut self,
563 key: &str,
564 ttl: Duration,
565 condition: FieldExpireCondition,
566 fields: &[&str],
567 ) -> Result<Vec<i64>, Error> {
568 let arguments = if ttl.subsec_nanos() == 0 {
569 args!["HEXPIRE", key, duration_secs(ttl)]
570 } else {
571 args!["HPEXPIRE", key, duration_millis(ttl)]
572 };
573
574 self.integers(field_block(arguments, condition, fields))
575 }
576
577 pub fn hpexpire(
579 &mut self,
580 key: &str,
581 ttl: Duration,
582 fields: &[&str],
583 ) -> Result<Vec<i64>, Error> {
584 self.integers(field_block(
585 args!["HPEXPIRE", key, duration_millis(ttl)],
586 FieldExpireCondition::Always,
587 fields,
588 ))
589 }
590
591 pub fn hexpire_at(
594 &mut self,
595 key: &str,
596 when: SystemTime,
597 fields: &[&str],
598 ) -> Result<Vec<i64>, Error> {
599 let millis = unix_millis(when)?;
600 let arguments = if millis % 1000 == 0 {
601 args!["HEXPIREAT", key, millis / 1000]
602 } else {
603 args!["HPEXPIREAT", key, millis]
604 };
605
606 self.integers(field_block(arguments, FieldExpireCondition::Always, fields))
607 }
608
609 pub fn httl(&mut self, key: &str, fields: &[&str]) -> Result<Vec<i64>, Error> {
611 self.integers(field_block(
612 args!["HTTL", key],
613 FieldExpireCondition::Always,
614 fields,
615 ))
616 }
617
618 pub fn hpttl(&mut self, key: &str, fields: &[&str]) -> Result<Vec<i64>, Error> {
620 self.integers(field_block(
621 args!["HPTTL", key],
622 FieldExpireCondition::Always,
623 fields,
624 ))
625 }
626
627 pub fn hexpiretime(&mut self, key: &str, fields: &[&str]) -> Result<Vec<i64>, Error> {
629 self.integers(field_block(
630 args!["HEXPIRETIME", key],
631 FieldExpireCondition::Always,
632 fields,
633 ))
634 }
635
636 pub fn hpexpiretime(&mut self, key: &str, fields: &[&str]) -> Result<Vec<i64>, Error> {
638 self.integers(field_block(
639 args!["HPEXPIRETIME", key],
640 FieldExpireCondition::Always,
641 fields,
642 ))
643 }
644
645 pub fn hpersist(&mut self, key: &str, fields: &[&str]) -> Result<Vec<i64>, Error> {
647 self.integers(field_block(
648 args!["HPERSIST", key],
649 FieldExpireCondition::Always,
650 fields,
651 ))
652 }
653
654 pub fn hscan(
656 &mut self,
657 key: &str,
658 cursor: u64,
659 pattern: Option<&str>,
660 count: Option<u64>,
661 ) -> Result<HashScanPage, Error> {
662 let mut arguments = args!["HSCAN", key, cursor];
663
664 push_scan_options(&mut arguments, pattern, count);
665
666 let (cursor, flat) = parse_scan("HSCAN", self.run(arguments)?)?;
667 let mut fields = Vec::with_capacity(flat.len() / 2);
668 let mut items = flat.into_iter();
669
670 while let (Some(field), Some(value)) = (items.next(), items.next()) {
671 fields.push((field, value));
672 }
673
674 Ok(HashScanPage { cursor, fields })
675 }
676
677 pub fn zadd(&mut self, key: &str, entries: &[(&str, f64)]) -> Result<i64, Error> {
685 self.zadd_with(key, SortedSetAddOptions::default(), entries)
686 }
687
688 pub fn zadd_with(
691 &mut self,
692 key: &str,
693 options: SortedSetAddOptions,
694 entries: &[(&str, f64)],
695 ) -> Result<i64, Error> {
696 let mut arguments = zadd_arguments(key, options);
697
698 for (member, score) in entries {
699 arguments.push(score.to_string());
700 arguments.push((*member).to_owned());
701 }
702
703 self.integer(arguments)
704 }
705
706 pub fn zadd_incr(
708 &mut self,
709 key: &str,
710 member: &str,
711 delta: f64,
712 options: SortedSetAddOptions,
713 ) -> Result<Option<f64>, Error> {
714 let mut arguments = zadd_arguments(key, options);
715
716 arguments.extend(args!["INCR", delta, member]);
717 self.bulk_or_null(arguments)?
718 .map(|text| parse_score(&text))
719 .transpose()
720 }
721
722 pub fn zincr_by(&mut self, key: &str, member: &str, delta: f64) -> Result<f64, Error> {
724 parse_score(&self.text(args!["ZINCRBY", key, delta, member])?)
725 }
726
727 pub fn geo_add(&mut self, key: &str, places: &[(f64, f64, &str)]) -> Result<i64, Error> {
729 let mut arguments = args!["GEOADD", key];
730
731 for (longitude, latitude, member) in places {
732 arguments.push(longitude.to_string());
733 arguments.push(latitude.to_string());
734 arguments.push((*member).to_owned());
735 }
736
737 self.integer(arguments)
738 }
739
740 pub fn geo_dist(
742 &mut self,
743 key: &str,
744 from: &str,
745 to: &str,
746 unit: Option<&str>,
747 ) -> Result<Option<f64>, Error> {
748 let mut arguments = args!["GEODIST", key, from, to];
749
750 if let Some(unit) = unit {
751 arguments.push(unit.to_owned());
752 }
753
754 self.bulk_or_null(arguments)?
755 .map(|text| parse_score(&text))
756 .transpose()
757 }
758
759 pub fn geo_hash(&mut self, key: &str, members: &[&str]) -> Result<Vec<Option<String>>, Error> {
761 let mut arguments = args!["GEOHASH", key];
762
763 arguments.extend(members.iter().map(|member| (*member).to_owned()));
764 self.optional_strings(arguments)
765 }
766
767 pub fn geo_pos(
769 &mut self,
770 key: &str,
771 members: &[&str],
772 ) -> Result<Vec<Option<(f64, f64)>>, Error> {
773 let mut arguments = args!["GEOPOS", key];
774
775 arguments.extend(members.iter().map(|member| (*member).to_owned()));
776
777 let reply = self.run(arguments)?;
778 let Some(items) = reply.as_array() else {
779 return Err(unexpected("GEOPOS", &reply));
780 };
781
782 let mut points = Vec::with_capacity(items.len());
783
784 for item in items {
785 if item.is_null() {
786 points.push(None);
787
788 continue;
789 }
790
791 let Some(pair) = item.as_array() else {
792 return Err(unexpected("GEOPOS", item));
793 };
794
795 if pair.len() != 2 {
796 return Err(unexpected("GEOPOS", item));
797 }
798
799 let longitude = parse_score(&pair[0].as_string()?.unwrap_or_default())?;
800 let latitude = parse_score(&pair[1].as_string()?.unwrap_or_default())?;
801
802 points.push(Some((longitude, latitude)));
803 }
804
805 Ok(points)
806 }
807
808 pub fn geo_search(
810 &mut self,
811 key: &str,
812 longitude: f64,
813 latitude: f64,
814 radius: f64,
815 unit: &str,
816 ) -> Result<Vec<String>, Error> {
817 self.strings(args![
818 "GEOSEARCH",
819 key,
820 "FROMLONLAT",
821 longitude,
822 latitude,
823 "BYRADIUS",
824 radius,
825 unit,
826 "ASC"
827 ])
828 }
829
830 pub fn geo_search_store(
832 &mut self,
833 destination: &str,
834 source: &str,
835 member: &str,
836 radius: f64,
837 unit: &str,
838 ) -> Result<i64, Error> {
839 self.integer(args![
840 "GEOSEARCHSTORE",
841 destination,
842 source,
843 "FROMMEMBER",
844 member,
845 "BYRADIUS",
846 radius,
847 unit
848 ])
849 }
850
851 pub fn geo_radius_by_member(
853 &mut self,
854 key: &str,
855 member: &str,
856 radius: f64,
857 unit: &str,
858 ) -> Result<Vec<String>, Error> {
859 self.strings(args!["GEORADIUSBYMEMBER", key, member, radius, unit, "ASC"])
860 }
861
862 pub fn geo_radius(
864 &mut self,
865 key: &str,
866 longitude: f64,
867 latitude: f64,
868 radius: f64,
869 unit: &str,
870 ) -> Result<Vec<String>, Error> {
871 self.strings(args![
872 "GEORADIUS",
873 key,
874 longitude,
875 latitude,
876 radius,
877 unit,
878 "ASC"
879 ])
880 }
881
882 pub fn zrange(&mut self, key: &str, start: i64, stop: i64) -> Result<Vec<String>, Error> {
884 self.strings(args!["ZRANGE", key, start, stop])
885 }
886
887 pub fn zrevrange(&mut self, key: &str, start: i64, stop: i64) -> Result<Vec<String>, Error> {
889 self.strings(args!["ZREVRANGE", key, start, stop])
890 }
891
892 pub fn zrange_with_scores(
894 &mut self,
895 key: &str,
896 start: i64,
897 stop: i64,
898 ) -> Result<Vec<SortedSetEntry>, Error> {
899 sorted_set(&self.run(args!["ZRANGE", key, start, stop, "WITHSCORES"])?)
900 }
901
902 pub fn zrevrange_with_scores(
904 &mut self,
905 key: &str,
906 start: i64,
907 stop: i64,
908 ) -> Result<Vec<SortedSetEntry>, Error> {
909 sorted_set(&self.run(args!["ZREVRANGE", key, start, stop, "WITHSCORES"])?)
910 }
911
912 pub fn zrange_by_lex(&mut self, key: &str, min: &str, max: &str) -> Result<Vec<String>, Error> {
914 self.strings(["ZRANGE", key, min, max, "BYLEX"])
915 }
916
917 pub fn zrange_by_score_with_scores(
919 &mut self,
920 key: &str,
921 min: &str,
922 max: &str,
923 ) -> Result<Vec<SortedSetEntry>, Error> {
924 sorted_set(&self.run(["ZRANGEBYSCORE", key, min, max, "WITHSCORES"])?)
925 }
926
927 pub fn zrem(&mut self, key: &str, members: &[&str]) -> Result<i64, Error> {
929 self.integer(with(args!["ZREM", key], members))
930 }
931
932 pub fn zpopmin(&mut self, key: &str, count: u64) -> Result<Vec<SortedSetEntry>, Error> {
934 sorted_set(&self.run(args!["ZPOPMIN", key, count])?)
935 }
936
937 pub fn zpopmax(&mut self, key: &str, count: u64) -> Result<Vec<SortedSetEntry>, Error> {
939 sorted_set(&self.run(args!["ZPOPMAX", key, count])?)
940 }
941
942 pub fn bzpopmin(
945 &mut self,
946 keys: &[&str],
947 timeout: Duration,
948 ) -> Result<Option<(String, SortedSetEntry)>, Error> {
949 self.blocking_sorted_pop("BZPOPMIN", keys, timeout)
950 }
951
952 pub fn bzpopmax(
954 &mut self,
955 keys: &[&str],
956 timeout: Duration,
957 ) -> Result<Option<(String, SortedSetEntry)>, Error> {
958 self.blocking_sorted_pop("BZPOPMAX", keys, timeout)
959 }
960
961 pub fn zremrangebyrank(&mut self, key: &str, start: i64, stop: i64) -> Result<i64, Error> {
963 self.integer(args!["ZREMRANGEBYRANK", key, start, stop])
964 }
965
966 pub fn zremrangebyscore(&mut self, key: &str, min: &str, max: &str) -> Result<i64, Error> {
968 self.integer(["ZREMRANGEBYSCORE", key, min, max])
969 }
970
971 pub fn zremrangebylex(&mut self, key: &str, min: &str, max: &str) -> Result<i64, Error> {
973 self.integer(["ZREMRANGEBYLEX", key, min, max])
974 }
975
976 pub fn zcard(&mut self, key: &str) -> Result<i64, Error> {
978 self.integer(["ZCARD", key])
979 }
980
981 pub fn zscore(&mut self, key: &str, member: &str) -> Result<Option<f64>, Error> {
983 self.bulk_or_null(["ZSCORE", key, member])?
984 .map(|text| parse_score(&text))
985 .transpose()
986 }
987
988 pub fn zrank(&mut self, key: &str, member: &str) -> Result<Option<i64>, Error> {
990 self.optional_integer(["ZRANK", key, member])
991 }
992
993 pub fn zrevrank(&mut self, key: &str, member: &str) -> Result<Option<i64>, Error> {
995 self.optional_integer(["ZREVRANK", key, member])
996 }
997
998 pub fn bf_reserve(&mut self, key: &str, error_rate: f64, capacity: u64) -> Result<(), Error> {
1000 self.ok(args!["BF.RESERVE", key, error_rate, capacity])
1001 }
1002
1003 pub fn bf_add(&mut self, key: &str, item: &str) -> Result<bool, Error> {
1005 self.flag(["BF.ADD", key, item])
1006 }
1007
1008 pub fn bf_exists(&mut self, key: &str, item: &str) -> Result<bool, Error> {
1010 self.flag(["BF.EXISTS", key, item])
1011 }
1012
1013 pub fn xadd_nomkstream(
1015 &mut self,
1016 key: &str,
1017 id: &str,
1018 fields: &[(&str, &str)],
1019 ) -> Result<Option<String>, Error> {
1020 self.bulk_or_null(with_pairs(args!["XADD", key, "NOMKSTREAM", id], fields))
1021 }
1022
1023 pub fn xadd_delay(
1025 &mut self,
1026 key: &str,
1027 delay_ms: u64,
1028 id: &str,
1029 fields: &[(&str, &str)],
1030 ) -> Result<String, Error> {
1031 self.text(with_pairs(
1032 args!["XADD", key, "DELAY", delay_ms, id],
1033 fields,
1034 ))
1035 }
1036
1037 pub fn xadd(&mut self, key: &str, id: &str, fields: &[(&str, &str)]) -> Result<String, Error> {
1039 self.text(with_pairs(args!["XADD", key, id], fields))
1040 }
1041
1042 pub fn xadd_maxlen(
1044 &mut self,
1045 key: &str,
1046 max_len: u64,
1047 id: &str,
1048 fields: &[(&str, &str)],
1049 ) -> Result<String, Error> {
1050 self.text(with_pairs(
1051 args!["XADD", key, "MAXLEN", max_len, id],
1052 fields,
1053 ))
1054 }
1055
1056 pub fn xlen(&mut self, key: &str) -> Result<i64, Error> {
1058 self.integer(["XLEN", key])
1059 }
1060
1061 pub fn xinfo_stream(&mut self, key: &str) -> Result<RespValue, Error> {
1063 self.run(["XINFO", "STREAM", key])
1064 }
1065
1066 pub fn xinfo_groups(&mut self, key: &str) -> Result<RespValue, Error> {
1068 self.run(["XINFO", "GROUPS", key])
1069 }
1070
1071 pub fn xinfo_consumers(&mut self, key: &str, group: &str) -> Result<RespValue, Error> {
1073 self.run(["XINFO", "CONSUMERS", key, group])
1074 }
1075
1076 pub fn xsetid(&mut self, key: &str, id: &str) -> Result<(), Error> {
1078 self.ok(["XSETID", key, id])
1079 }
1080
1081 pub fn xautoclaim(
1083 &mut self,
1084 key: &str,
1085 group: &str,
1086 consumer: &str,
1087 min_idle_ms: u64,
1088 start: &str,
1089 ) -> Result<RespValue, Error> {
1090 self.run([
1091 "XAUTOCLAIM",
1092 key,
1093 group,
1094 consumer,
1095 &min_idle_ms.to_string(),
1096 start,
1097 ])
1098 }
1099
1100 pub fn xrange(
1102 &mut self,
1103 key: &str,
1104 start: &str,
1105 end: &str,
1106 count: Option<u64>,
1107 ) -> Result<Vec<StreamEntry>, Error> {
1108 let mut arguments = args!["XRANGE", key, start, end];
1109
1110 if let Some(count) = count {
1111 arguments.extend(args!["COUNT", count]);
1112 }
1113
1114 stream_entries(&self.run(arguments)?)
1115 }
1116
1117 pub fn xrevrange(
1119 &mut self,
1120 key: &str,
1121 end: &str,
1122 start: &str,
1123 count: Option<u64>,
1124 ) -> Result<Vec<StreamEntry>, Error> {
1125 let mut arguments = args!["XREVRANGE", key, end, start];
1126
1127 if let Some(count) = count {
1128 arguments.extend(args!["COUNT", count]);
1129 }
1130
1131 stream_entries(&self.run(arguments)?)
1132 }
1133
1134 pub fn xdel(&mut self, key: &str, ids: &[&str]) -> Result<i64, Error> {
1136 self.integer(with(args!["XDEL", key], ids))
1137 }
1138
1139 pub fn xtrim_minid(&mut self, key: &str, id: &str) -> Result<i64, Error> {
1141 self.integer(args!["XTRIM", key, "MINID", id])
1142 }
1143
1144 pub fn xtrim_maxlen(&mut self, key: &str, max_len: u64) -> Result<i64, Error> {
1146 self.integer(args!["XTRIM", key, "MAXLEN", max_len])
1147 }
1148
1149 pub fn xread(
1151 &mut self,
1152 options: &StreamReadOptions,
1153 streams: &[(&str, &str)],
1154 ) -> Result<Vec<StreamReadResult>, Error> {
1155 if options.no_ack {
1156 return Err(Error::Protocol(
1157 "NOACK applies only to XREADGROUP".to_owned(),
1158 ));
1159 }
1160
1161 let mut arguments = args!["XREAD"];
1162
1163 push_stream_read_options(&mut arguments, options);
1164 push_streams(&mut arguments, streams);
1165 stream_read(&self.run(arguments)?)
1166 }
1167
1168 pub fn xgroup_create_consumer(
1170 &mut self,
1171 key: &str,
1172 group: &str,
1173 consumer: &str,
1174 ) -> Result<i64, Error> {
1175 self.integer(args!["XGROUP", "CREATECONSUMER", key, group, consumer])
1176 }
1177
1178 pub fn xgroup_create(
1180 &mut self,
1181 key: &str,
1182 group: &str,
1183 id: &str,
1184 make_stream: bool,
1185 ) -> Result<(), Error> {
1186 let mut arguments = args!["XGROUP", "CREATE", key, group, id];
1187
1188 if make_stream {
1189 arguments.push("MKSTREAM".to_owned());
1190 }
1191
1192 self.ok(arguments)
1193 }
1194
1195 pub fn xreadgroup(
1198 &mut self,
1199 group: &str,
1200 consumer: &str,
1201 options: &StreamReadOptions,
1202 streams: &[(&str, &str)],
1203 ) -> Result<Vec<StreamReadResult>, Error> {
1204 let mut arguments = args!["XREADGROUP", "GROUP", group, consumer];
1205
1206 push_stream_read_options(&mut arguments, options);
1207
1208 if options.no_ack {
1209 arguments.push("NOACK".to_owned());
1210 }
1211
1212 push_streams(&mut arguments, streams);
1213 stream_read(&self.run(arguments)?)
1214 }
1215
1216 pub fn xgroup_destroy(&mut self, key: &str, group: &str) -> Result<bool, Error> {
1218 self.flag(["XGROUP", "DESTROY", key, group])
1219 }
1220
1221 pub fn xgroup_setid(&mut self, key: &str, group: &str, id: &str) -> Result<(), Error> {
1223 self.ok(["XGROUP", "SETID", key, group, id])
1224 }
1225
1226 pub fn xgroup_delconsumer(
1228 &mut self,
1229 key: &str,
1230 group: &str,
1231 consumer: &str,
1232 ) -> Result<i64, Error> {
1233 self.integer(["XGROUP", "DELCONSUMER", key, group, consumer])
1234 }
1235
1236 pub fn xack(&mut self, key: &str, group: &str, ids: &[&str]) -> Result<i64, Error> {
1238 self.integer(with(args!["XACK", key, group], ids))
1239 }
1240
1241 pub fn xpending_summary(&mut self, key: &str, group: &str) -> Result<RespValue, Error> {
1243 self.run(["XPENDING", key, group])
1244 }
1245
1246 pub fn xpending(
1249 &mut self,
1250 key: &str,
1251 group: &str,
1252 start: &str,
1253 end: &str,
1254 count: u64,
1255 filter: &StreamPendingFilter,
1256 ) -> Result<Vec<StreamPendingEntry>, Error> {
1257 let mut arguments = args!["XPENDING", key, group];
1258
1259 if let Some(idle) = filter.min_idle {
1260 arguments.extend(args!["IDLE", duration_millis(idle)]);
1261 }
1262
1263 arguments.extend(args![start, end, count]);
1264
1265 if let Some(consumer) = &filter.consumer {
1266 arguments.push(consumer.clone());
1267 }
1268
1269 let reply = self.run(arguments)?;
1270 let mut rows = Vec::new();
1271
1272 for row in reply.as_array().unwrap_or_default() {
1273 let Some([id, owner, idle, deliveries]) = row.as_array() else {
1274 return Err(unexpected("XPENDING", row));
1275 };
1276
1277 rows.push(StreamPendingEntry {
1278 id: id.as_string()?.unwrap_or_default(),
1279 consumer: owner.as_string()?.unwrap_or_default(),
1280 idle_millis: idle.as_integer().unwrap_or_default(),
1281 delivery_count: deliveries.as_integer().unwrap_or_default(),
1282 });
1283 }
1284
1285 Ok(rows)
1286 }
1287
1288 pub fn xclaim(
1290 &mut self,
1291 key: &str,
1292 group: &str,
1293 consumer: &str,
1294 min_idle: Duration,
1295 ids: &[&str],
1296 options: &StreamClaimOptions,
1297 ) -> Result<Vec<StreamEntry>, Error> {
1298 let arguments = claim_arguments(key, group, consumer, min_idle, ids, options)?;
1299
1300 stream_entries(&self.run(arguments)?)
1301 }
1302
1303 pub fn xclaim_ids(
1305 &mut self,
1306 key: &str,
1307 group: &str,
1308 consumer: &str,
1309 min_idle: Duration,
1310 ids: &[&str],
1311 options: &StreamClaimOptions,
1312 ) -> Result<Vec<String>, Error> {
1313 let mut arguments = claim_arguments(key, group, consumer, min_idle, ids, options)?;
1314
1315 arguments.push("JUSTID".to_owned());
1316 self.strings(arguments)
1317 }
1318
1319 pub fn info(&mut self, section: Option<&str>) -> Result<String, Error> {
1321 match section {
1322 Some(section) => self.text(["INFO", section]),
1323 None => self.text(["INFO"]),
1324 }
1325 }
1326
1327 pub fn config_get(&mut self, parameter: &str) -> Result<Vec<(String, String)>, Error> {
1329 pairs(&self.run(["CONFIG", "GET", parameter])?)
1330 }
1331
1332 pub fn save(&mut self) -> Result<(), Error> {
1334 self.ok(["SAVE"])
1335 }
1336
1337 pub fn bgsave(&mut self) -> Result<(), Error> {
1339 self.ok(["BGSAVE"])
1340 }
1341
1342 pub fn flushdb(&mut self) -> Result<(), Error> {
1344 self.ok(["FLUSHDB"])
1345 }
1346
1347 pub fn flushall(&mut self) -> Result<(), Error> {
1349 self.ok(["FLUSHALL"])
1350 }
1351
1352 pub fn swapdb(&mut self, first: u32, second: u32) -> Result<(), Error> {
1354 self.ok(args!["SWAPDB", first, second])
1355 }
1356
1357 pub fn move_key(&mut self, key: &str, database: u32) -> Result<bool, Error> {
1359 self.flag(args!["MOVE", key, database])
1360 }
1361
1362 pub fn multi(&mut self) -> Result<(), Error> {
1365 self.ok(["MULTI"])
1366 }
1367
1368 pub fn exec(&mut self) -> Result<Option<Vec<RespValue>>, Error> {
1370 match self.run(["EXEC"])? {
1371 RespValue::Null => Ok(None),
1372 RespValue::Array(items) => Ok(Some(items)),
1373 other => Err(unexpected("EXEC", &other)),
1374 }
1375 }
1376
1377 pub fn discard(&mut self) -> Result<(), Error> {
1379 self.ok(["DISCARD"])
1380 }
1381
1382 pub fn watch_key(&mut self, key: &str) -> Result<(), Error> {
1384 self.watch(&[key])
1385 }
1386
1387 pub fn watch(&mut self, keys: &[&str]) -> Result<(), Error> {
1389 self.ok(join("WATCH", keys))
1390 }
1391
1392 pub fn unwatch(&mut self) -> Result<(), Error> {
1394 self.ok(["UNWATCH"])
1395 }
1396
1397 pub fn eval(
1399 &mut self,
1400 script: &str,
1401 keys: &[&str],
1402 arguments: &[&str],
1403 ) -> Result<RespValue, Error> {
1404 self.run(script_arguments("EVAL", script, keys, arguments))
1405 }
1406
1407 pub fn evalsha(
1409 &mut self,
1410 sha: &str,
1411 keys: &[&str],
1412 arguments: &[&str],
1413 ) -> Result<RespValue, Error> {
1414 self.run(script_arguments("EVALSHA", sha, keys, arguments))
1415 }
1416
1417 pub fn script_load(&mut self, script: &str) -> Result<String, Error> {
1419 self.text(["SCRIPT", "LOAD", script])
1420 }
1421
1422 pub fn script_exists(&mut self, hashes: &[&str]) -> Result<Vec<bool>, Error> {
1424 let reply = self.run(with(args!["SCRIPT", "EXISTS"], hashes))?;
1425
1426 Ok(reply
1427 .as_array()
1428 .unwrap_or_default()
1429 .iter()
1430 .map(|item| item.as_integer().unwrap_or_default() > 0)
1431 .collect())
1432 }
1433
1434 pub fn script_flush(&mut self) -> Result<(), Error> {
1436 self.ok(["SCRIPT", "FLUSH"])
1437 }
1438
1439 pub fn script_kill(&mut self) -> Result<(), Error> {
1441 self.ok(["SCRIPT", "KILL"])
1442 }
1443
1444 pub fn fcall(
1446 &mut self,
1447 name: &str,
1448 keys: &[&str],
1449 arguments: &[&str],
1450 ) -> Result<RespValue, Error> {
1451 self.run(script_arguments("FCALL", name, keys, arguments))
1452 }
1453
1454 pub fn function_load(&mut self, name: &str, body: &str) -> Result<(), Error> {
1457 self.ok(["FUNCTION", "LOAD", name, body])
1458 }
1459
1460 pub fn function_list(&mut self) -> Result<RespValue, Error> {
1462 self.run(["FUNCTION", "LIST"])
1463 }
1464
1465 pub fn function_delete(&mut self, name: &str) -> Result<bool, Error> {
1467 self.flag(["FUNCTION", "DELETE", name])
1468 }
1469
1470 pub fn acl_whoami(&mut self) -> Result<String, Error> {
1472 self.text(["ACL", "WHOAMI"])
1473 }
1474
1475 pub fn acl_users(&mut self) -> Result<Vec<String>, Error> {
1477 self.strings(["ACL", "USERS"])
1478 }
1479
1480 pub fn acl_list(&mut self) -> Result<Vec<String>, Error> {
1482 self.strings(["ACL", "LIST"])
1483 }
1484
1485 pub fn acl_getuser(&mut self, user: &str) -> Result<RespValue, Error> {
1487 self.run(["ACL", "GETUSER", user])
1488 }
1489
1490 pub fn acl_cat(&mut self) -> Result<Vec<String>, Error> {
1492 self.strings(["ACL", "CAT"])
1493 }
1494
1495 pub fn acl_setuser(&mut self, user: &str, rules: &[&str]) -> Result<(), Error> {
1497 self.ok(with(args!["ACL", "SETUSER", user], rules))
1498 }
1499
1500 pub fn acl_load(&mut self) -> Result<(), Error> {
1502 self.ok(["ACL", "LOAD"])
1503 }
1504
1505 pub fn acl_save(&mut self) -> Result<(), Error> {
1507 self.ok(["ACL", "SAVE"])
1508 }
1509
1510 pub fn cluster_slots(&mut self) -> Result<RespValue, Error> {
1512 self.run(["CLUSTER", "SLOTS"])
1513 }
1514
1515 pub fn cluster_nodes(&mut self) -> Result<String, Error> {
1517 self.text(["CLUSTER", "NODES"])
1518 }
1519
1520 pub fn cluster_shards(&mut self) -> Result<RespValue, Error> {
1522 self.run(["CLUSTER", "SHARDS"])
1523 }
1524
1525 pub fn cluster_info(&mut self) -> Result<String, Error> {
1527 self.text(["CLUSTER", "INFO"])
1528 }
1529
1530 pub fn cluster_myid(&mut self) -> Result<String, Error> {
1532 self.text(["CLUSTER", "MYID"])
1533 }
1534
1535 pub fn cluster_keyslot(&mut self, key: &str) -> Result<i64, Error> {
1537 self.integer(["CLUSTER", "KEYSLOT", key])
1538 }
1539
1540 pub fn readonly(&mut self) -> Result<(), Error> {
1542 self.ok(["READONLY"])
1543 }
1544
1545 pub fn readwrite(&mut self) -> Result<(), Error> {
1547 self.ok(["READWRITE"])
1548 }
1549
1550 pub fn client_getname(&mut self) -> Result<Option<String>, Error> {
1552 self.bulk_or_null(["CLIENT", "GETNAME"])
1553 }
1554
1555 pub fn client_setname(&mut self, name: &str) -> Result<(), Error> {
1557 self.ok(["CLIENT", "SETNAME", name])
1558 }
1559
1560 pub fn client_tracking(&mut self, enabled: bool) -> Result<(), Error> {
1562 self.ok(["CLIENT", "TRACKING", if enabled { "ON" } else { "OFF" }])
1563 }
1564
1565 pub fn hello(&mut self, protocol: u8) -> Result<RespValue, Error> {
1568 self.run(args!["HELLO", protocol])
1569 }
1570
1571 pub fn publish(&mut self, channel: &str, message: &str) -> Result<i64, Error> {
1573 self.integer(["PUBLISH", channel, message])
1574 }
1575
1576 pub fn spublish(&mut self, channel: &str, message: &str) -> Result<i64, Error> {
1578 self.integer(["SPUBLISH", channel, message])
1579 }
1580
1581 pub fn subscribe(&mut self, channels: &[&str]) -> Result<Vec<RespValue>, Error> {
1585 self.subscription("SUBSCRIBE", channels, false)
1586 }
1587
1588 pub fn unsubscribe(&mut self, channels: &[&str]) -> Result<Vec<RespValue>, Error> {
1590 self.subscription("UNSUBSCRIBE", channels, true)
1591 }
1592
1593 pub fn psubscribe(&mut self, patterns: &[&str]) -> Result<Vec<RespValue>, Error> {
1595 self.subscription("PSUBSCRIBE", patterns, false)
1596 }
1597
1598 pub fn punsubscribe(&mut self, patterns: &[&str]) -> Result<Vec<RespValue>, Error> {
1600 self.subscription("PUNSUBSCRIBE", patterns, true)
1601 }
1602
1603 pub fn ssubscribe(&mut self, channels: &[&str]) -> Result<Vec<RespValue>, Error> {
1605 self.subscription("SSUBSCRIBE", channels, false)
1606 }
1607
1608 pub fn sunsubscribe(&mut self, channels: &[&str]) -> Result<Vec<RespValue>, Error> {
1610 self.subscription("SUNSUBSCRIBE", channels, true)
1611 }
1612
1613 pub fn next_message(&mut self) -> Result<PubSubMessage, Error> {
1616 loop {
1617 let reply = self.read_message()?;
1618
1619 if let Some(message) = pubsub_message(&reply)? {
1620 return Ok(message);
1621 }
1622 }
1623 }
1624
1625 pub fn listen<F>(&mut self, mut on_message: F) -> Result<(), Error>
1628 where
1629 F: FnMut(PubSubMessage),
1630 {
1631 loop {
1632 on_message(self.next_message()?);
1633 }
1634 }
1635
1636 fn subscription(
1637 &mut self,
1638 command: &str,
1639 names: &[&str],
1640 allow_empty: bool,
1641 ) -> Result<Vec<RespValue>, Error> {
1642 if names.is_empty() && !allow_empty {
1643 return Err(Error::Protocol(format!(
1644 "{command} needs at least one channel or pattern"
1645 )));
1646 }
1647
1648 self.run_replies(join(command, names), names.len().max(1))
1649 }
1650
1651 fn optional_integer<I, S>(&mut self, arguments: I) -> Result<Option<i64>, Error>
1652 where
1653 I: IntoIterator<Item = S>,
1654 S: AsRef<str>,
1655 {
1656 let (command, reply) = self.run_named(arguments)?;
1657
1658 match reply {
1659 RespValue::Null => Ok(None),
1660 RespValue::Integer(value) => Ok(Some(value)),
1661 other => Err(unexpected(&command, &other)),
1662 }
1663 }
1664
1665 fn integers(&mut self, arguments: Vec<String>) -> Result<Vec<i64>, Error> {
1666 let command = arguments.first().cloned().unwrap_or_default();
1667 let reply = self.run(arguments)?;
1668
1669 integer_array(&command, &reply)
1670 }
1671
1672 fn blocking_sorted_pop(
1673 &mut self,
1674 command: &str,
1675 keys: &[&str],
1676 timeout: Duration,
1677 ) -> Result<Option<(String, SortedSetEntry)>, Error> {
1678 let mut arguments = join(command, keys);
1679
1680 arguments.push(duration_secs(timeout).to_string());
1681
1682 match self.run(arguments)? {
1683 RespValue::Null => Ok(None),
1684 RespValue::Array(items) if items.len() == 3 => Ok(Some((
1685 items[0].as_string()?.unwrap_or_default(),
1686 SortedSetEntry::new(
1687 items[1].as_string()?.unwrap_or_default(),
1688 parse_score(&items[2].as_string()?.unwrap_or_default())?,
1689 ),
1690 ))),
1691 other => Err(unexpected(command, &other)),
1692 }
1693 }
1694
1695 fn blocking_pop(
1696 &mut self,
1697 command: &str,
1698 keys: &[&str],
1699 timeout: Duration,
1700 ) -> Result<Option<(String, String)>, Error> {
1701 let mut arguments = join(command, keys);
1702
1703 arguments.push(duration_secs(timeout).to_string());
1704
1705 match self.run(arguments)? {
1706 RespValue::Null => Ok(None),
1707 RespValue::Array(items) if items.len() == 2 => Ok(Some((
1708 items[0].as_string()?.unwrap_or_default(),
1709 items[1].as_string()?.unwrap_or_default(),
1710 ))),
1711 other => Err(unexpected(command, &other)),
1712 }
1713 }
1714}
1715
1716fn field_block(
1718 mut arguments: Vec<String>,
1719 condition: FieldExpireCondition,
1720 fields: &[&str],
1721) -> Vec<String> {
1722 let flag = match condition {
1723 FieldExpireCondition::Always => None,
1724 FieldExpireCondition::IfNoTtl => Some("NX"),
1725 FieldExpireCondition::IfTtl => Some("XX"),
1726 FieldExpireCondition::IfGreater => Some("GT"),
1727 FieldExpireCondition::IfLess => Some("LT"),
1728 };
1729
1730 arguments.extend(flag.map(str::to_owned));
1731 arguments.push("FIELDS".to_owned());
1732 arguments.push(fields.len().to_string());
1733
1734 with(arguments, fields)
1735}
1736
1737fn integer_array(command: &str, reply: &RespValue) -> Result<Vec<i64>, Error> {
1738 let Some(items) = reply.as_array() else {
1739 return Err(unexpected(command, reply));
1740 };
1741
1742 items
1743 .iter()
1744 .map(|item| item.as_integer().ok_or_else(|| unexpected(command, item)))
1745 .collect()
1746}
1747
1748fn with(mut arguments: Vec<String>, rest: &[&str]) -> Vec<String> {
1749 arguments.extend(rest.iter().map(|argument| (*argument).to_owned()));
1750
1751 arguments
1752}
1753
1754fn with_pairs(mut arguments: Vec<String>, pairs: &[(&str, &str)]) -> Vec<String> {
1755 for (left, right) in pairs {
1756 arguments.push((*left).to_owned());
1757 arguments.push((*right).to_owned());
1758 }
1759
1760 arguments
1761}
1762
1763fn push_scan_options(arguments: &mut Vec<String>, pattern: Option<&str>, count: Option<u64>) {
1764 if let Some(pattern) = pattern {
1765 arguments.extend(args!["MATCH", pattern]);
1766 }
1767
1768 if let Some(count) = count {
1769 arguments.extend(args!["COUNT", count]);
1770 }
1771}
1772
1773fn push_stream_read_options(arguments: &mut Vec<String>, options: &StreamReadOptions) {
1774 if let Some(count) = options.count {
1775 arguments.extend(args!["COUNT", count]);
1776 }
1777
1778 if let Some(block) = options.block {
1779 arguments.extend(args!["BLOCK", duration_millis(block)]);
1780 }
1781}
1782
1783fn push_streams(arguments: &mut Vec<String>, streams: &[(&str, &str)]) {
1784 arguments.push("STREAMS".to_owned());
1785 arguments.extend(streams.iter().map(|(key, _)| (*key).to_owned()));
1786 arguments.extend(streams.iter().map(|(_, id)| (*id).to_owned()));
1787}
1788
1789fn zadd_arguments(key: &str, options: SortedSetAddOptions) -> Vec<String> {
1790 let mut arguments = args!["ZADD", key];
1791 let flags = [
1792 (options.if_not_exists, "NX"),
1793 (options.if_exists, "XX"),
1794 (options.greater_than, "GT"),
1795 (options.less_than, "LT"),
1796 (options.changed, "CH"),
1797 ];
1798
1799 for (enabled, flag) in flags {
1800 if enabled {
1801 arguments.push(flag.to_owned());
1802 }
1803 }
1804
1805 arguments
1806}
1807
1808fn claim_arguments(
1809 key: &str,
1810 group: &str,
1811 consumer: &str,
1812 min_idle: Duration,
1813 ids: &[&str],
1814 options: &StreamClaimOptions,
1815) -> Result<Vec<String>, Error> {
1816 if options.idle.is_some() && options.time.is_some() {
1817 return Err(Error::Protocol(
1818 "IDLE and TIME cannot be combined".to_owned(),
1819 ));
1820 }
1821
1822 let mut arguments = with(
1823 args!["XCLAIM", key, group, consumer, duration_millis(min_idle)],
1824 ids,
1825 );
1826
1827 if let Some(idle) = options.idle {
1828 arguments.extend(args!["IDLE", duration_millis(idle)]);
1829 }
1830
1831 if let Some(time) = options.time {
1832 arguments.extend(args!["TIME", unix_millis(time)?]);
1833 }
1834
1835 if let Some(retry_count) = options.retry_count {
1836 arguments.extend(args!["RETRYCOUNT", retry_count]);
1837 }
1838
1839 if options.force {
1840 arguments.push("FORCE".to_owned());
1841 }
1842
1843 if let Some(last_id) = &options.last_id {
1844 arguments.extend(args!["LASTID", last_id]);
1845 }
1846
1847 Ok(arguments)
1848}
1849
1850fn script_arguments(command: &str, script: &str, keys: &[&str], arguments: &[&str]) -> Vec<String> {
1851 let values = with(args![command, script, keys.len()], keys);
1852
1853 with(values, arguments)
1854}
1855
1856fn parse_scan(command: &str, reply: RespValue) -> Result<(u64, Vec<String>), Error> {
1857 let Some([cursor, items]) = reply.as_array() else {
1858 return Err(unexpected(command, &reply));
1859 };
1860
1861 let cursor = cursor
1862 .as_string()?
1863 .unwrap_or_default()
1864 .parse()
1865 .map_err(|_| Error::Protocol(format!("{command} returned a bad cursor")))?;
1866 let mut keys = Vec::new();
1867
1868 for item in items.as_array().unwrap_or_default() {
1869 keys.push(item.as_string()?.unwrap_or_default());
1870 }
1871
1872 Ok((cursor, keys))
1873}
1874
1875fn parse_score(text: &str) -> Result<f64, Error> {
1876 text.parse()
1877 .map_err(|_| Error::Protocol(format!("{text:?} is not a score")))
1878}
1879
1880fn pairs(reply: &RespValue) -> Result<Vec<(String, String)>, Error> {
1881 let items = reply.as_array().unwrap_or_default();
1882 let mut values = Vec::with_capacity(items.len() / 2);
1883
1884 for pair in items.chunks_exact(2) {
1885 values.push((
1886 pair[0].as_string()?.unwrap_or_default(),
1887 pair[1].as_string()?.unwrap_or_default(),
1888 ));
1889 }
1890
1891 Ok(values)
1892}
1893
1894fn sorted_set(reply: &RespValue) -> Result<Vec<SortedSetEntry>, Error> {
1895 pairs(reply)?
1896 .into_iter()
1897 .map(|(member, score)| Ok(SortedSetEntry::new(member, parse_score(&score)?)))
1898 .collect()
1899}
1900
1901fn stream_entries(reply: &RespValue) -> Result<Vec<StreamEntry>, Error> {
1902 let mut entries = Vec::new();
1903
1904 for item in reply.as_array().unwrap_or_default() {
1905 let Some([id, fields]) = item.as_array() else {
1906 continue;
1907 };
1908
1909 entries.push(StreamEntry {
1910 id: id.as_string()?.unwrap_or_default(),
1911 fields: pairs(fields)?,
1912 });
1913 }
1914
1915 Ok(entries)
1916}
1917
1918fn stream_read(reply: &RespValue) -> Result<Vec<StreamReadResult>, Error> {
1919 let mut results = Vec::new();
1920
1921 for item in reply.as_array().unwrap_or_default() {
1922 let Some([key, entries]) = item.as_array() else {
1923 continue;
1924 };
1925
1926 results.push(StreamReadResult {
1927 key: key.as_string()?.unwrap_or_default(),
1928 entries: stream_entries(entries)?,
1929 });
1930 }
1931
1932 Ok(results)
1933}
1934
1935fn pubsub_message(reply: &RespValue) -> Result<Option<PubSubMessage>, Error> {
1936 let items = reply.as_array().unwrap_or_default();
1937 let kind = match items.first() {
1938 Some(first) => first.as_string()?.unwrap_or_default(),
1939 None => return Ok(None),
1940 };
1941
1942 let (pattern, channel, payload) = match (kind.as_str(), items) {
1943 ("message" | "smessage", [_, channel, payload]) => (None, channel, payload),
1944 ("pmessage", [_, pattern, channel, payload]) => (pattern.as_string()?, channel, payload),
1945
1946 _ => return Ok(None),
1947 };
1948
1949 let payload = match payload {
1950 RespValue::Bulk(bytes) => bytes.clone(),
1951 other => other.as_string()?.unwrap_or_default().into_bytes(),
1952 };
1953
1954 Ok(Some(PubSubMessage {
1955 kind,
1956 pattern,
1957 channel: channel.as_string()?.unwrap_or_default(),
1958 payload,
1959 }))
1960}
1961
1962#[cfg(test)]
1963mod tests {
1964 use super::{field_block, integer_array, pubsub_message};
1965 use crate::models::FieldExpireCondition;
1966 use crate::value::RespValue;
1967
1968 #[test]
1969 fn field_block_puts_the_condition_before_fields() {
1970 assert_eq!(
1971 field_block(
1972 vec!["HEXPIRE".into(), "h".into(), "10".into()],
1973 FieldExpireCondition::IfGreater,
1974 &["a", "b"],
1975 ),
1976 ["HEXPIRE", "h", "10", "GT", "FIELDS", "2", "a", "b"]
1977 );
1978 assert_eq!(
1979 field_block(
1980 vec!["HTTL".into(), "h".into()],
1981 FieldExpireCondition::Always,
1982 &["a"]
1983 ),
1984 ["HTTL", "h", "FIELDS", "1", "a"]
1985 );
1986 assert_eq!(
1987 integer_array(
1988 "HTTL",
1989 &RespValue::Array(vec![RespValue::Integer(-2), RespValue::Integer(5)])
1990 )
1991 .unwrap(),
1992 vec![-2, 5]
1993 );
1994 }
1995
1996 #[test]
1997 fn pubsub_message_reads_payload_from_an_array() {
1998 let reply = RespValue::Array(vec![
1999 RespValue::Bulk(b"message".to_vec()),
2000 RespValue::Bulk(b"news".to_vec()),
2001 RespValue::Bulk(b"hello-ruvio".to_vec()),
2002 ]);
2003
2004 let parsed = pubsub_message(&reply).unwrap().unwrap();
2005
2006 assert_eq!(parsed.kind, "message");
2007 assert_eq!(parsed.channel, "news");
2008 assert_eq!(parsed.payload, b"hello-ruvio");
2009 assert!(
2010 pubsub_message(&RespValue::Array(vec![
2011 RespValue::Bulk(b"subscribe".to_vec()),
2012 RespValue::Bulk(b"news".to_vec()),
2013 RespValue::Integer(1),
2014 ]))
2015 .unwrap()
2016 .is_none()
2017 );
2018 }
2019}