Skip to main content

ruvio_client/
commands.rs

1// client_coverage: APPEND DEL EVAL EXEC EXPIRE GET HGET HGETALL HSET INCR
2// LPUSH LRANGE MGET MSET MULTI PING SADD SET SMEMBERS STRLEN TTL ZADD ZRANGE
3
4use 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    /// Sends `MGET`. A missing key is `None`.
27    pub fn mget(&mut self, keys: &[&str]) -> Result<Vec<Option<String>>, Error> {
28        self.optional_strings(join("MGET", keys))
29    }
30
31    /// Sends `LIMIT`. `None` means this call did not fit in the window.
32    /// `Some` is how many calls remain after this one.
33    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    /// Sends `ONCE`. `1` is the first call. `0` means the key is already taken.
47    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    /// Native reservation: [acquired], [pending], [conflict], or [completed, binary result].
52    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    /// Publishes a result only for the matching live owner, without renewing retention.
70    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    /// Sends `TAKE`. `None` means the key did not hold enough and was left unchanged.
81    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    /// Sends `GETEX`. Returns the string and sets its TTL to `ttl_ms`.
95    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    /// Sends `LEASE`. `None` means another holder still has the lease.
100    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    /// Sends `SEMAPHORE key max ttl_ms`. `Some(token)` while fewer than `max` holders are
114    /// live; `None` when full. Pass the token to [`Client::release`].
115    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    /// Sends `INCRBY key delta MAX limit`. `None` means the result would pass `limit`, and the
129    /// key was left unchanged.
130    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    /// Sends `RELEASE`. Returns 1 when the token matched the holder of a lease or a semaphore.
144    pub fn release(&mut self, key: &str, token: u64) -> Result<i64, Error> {
145        self.integer(["RELEASE", key, &token.to_string()])
146    }
147
148    /// Sends `SETV`. Omit `version` to create version 1.
149    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    /// Sends `CHANGES START` and returns the cursor to read after.
160    pub fn changes_start(&mut self) -> Result<i64, Error> {
161        self.integer(["CHANGES", "START"])
162    }
163
164    /// Sends `CHANGES` and returns `(sequence, text)` notes after `cursor`.
165    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    /// Sends `MSET`.
192    pub fn mset(&mut self, pairs: &[(&str, &str)]) -> Result<(), Error> {
193        self.ok(with_pairs(args!["MSET"], pairs))
194    }
195
196    /// Sends `GETSET` and returns the previous value.
197    pub fn getset(&mut self, key: &str, value: &str) -> Result<Option<String>, Error> {
198        self.bulk_or_null(["GETSET", key, value])
199    }
200
201    /// Sends `APPEND` and returns the new length.
202    pub fn append(&mut self, key: &str, value: &str) -> Result<i64, Error> {
203        self.integer(["APPEND", key, value])
204    }
205
206    /// Sends `STRLEN`.
207    pub fn strlen(&mut self, key: &str) -> Result<i64, Error> {
208        self.integer(["STRLEN", key])
209    }
210
211    /// Sends `DECR`.
212    pub fn decr(&mut self, key: &str) -> Result<i64, Error> {
213        self.integer(["DECR", key])
214    }
215
216    /// Sends `INCRBY`.
217    pub fn incr_by(&mut self, key: &str, delta: i64) -> Result<i64, Error> {
218        self.integer(args!["INCRBY", key, delta])
219    }
220
221    /// Sends `DECRBY`.
222    pub fn decr_by(&mut self, key: &str, delta: i64) -> Result<i64, Error> {
223        self.integer(args!["DECRBY", key, delta])
224    }
225
226    /// Sends `UNLINK` for one key.
227    pub fn unlink_key(&mut self, key: &str) -> Result<i64, Error> {
228        self.unlink(&[key])
229    }
230
231    /// Sends `UNLINK` and returns how many keys were removed.
232    pub fn unlink(&mut self, keys: &[&str]) -> Result<i64, Error> {
233        self.integer(join("UNLINK", keys))
234    }
235
236    /// Sends `EXISTS` for one key.
237    pub fn exists_key(&mut self, key: &str) -> Result<i64, Error> {
238        self.exists(&[key])
239    }
240
241    /// Sends `EXISTS` and returns how many of the keys exist.
242    pub fn exists(&mut self, keys: &[&str]) -> Result<i64, Error> {
243        self.integer(join("EXISTS", keys))
244    }
245
246    /// Sends `TYPE`: `string`, `list`, `set`, `hash`, `zset`, `stream`, or `none`.
247    pub fn key_type(&mut self, key: &str) -> Result<String, Error> {
248        self.text(["TYPE", key])
249    }
250
251    /// Sends `RENAME`.
252    pub fn rename(&mut self, key: &str, new_key: &str) -> Result<(), Error> {
253        self.ok(["RENAME", key, new_key])
254    }
255
256    /// Sends `SCAN` with an optional `MATCH` pattern and `COUNT` hint.
257    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    /// Sends `DBSIZE` for the selected database.
273    pub fn dbsize(&mut self) -> Result<i64, Error> {
274        self.integer(["DBSIZE"])
275    }
276
277    /// Sends `EXPIREAT`, or `PEXPIREAT` when `when` is not a whole second.
278    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    /// Sends `PEXPIRE`. False when the key does not exist.
289    pub fn pexpire(&mut self, key: &str, ttl: Duration) -> Result<bool, Error> {
290        self.flag(args!["PEXPIRE", key, duration_millis(ttl)])
291    }
292
293    /// Sends `TTL`: seconds left, -1 without an expiry, -2 when the key is missing.
294    pub fn ttl(&mut self, key: &str) -> Result<i64, Error> {
295        self.integer(["TTL", key])
296    }
297
298    /// Sends `PTTL`: milliseconds left, -1 without an expiry, -2 when the key is missing.
299    pub fn pttl(&mut self, key: &str) -> Result<i64, Error> {
300        self.integer(["PTTL", key])
301    }
302
303    /// Sends `PERSIST`. False when the key is missing or has no expiry.
304    pub fn persist(&mut self, key: &str) -> Result<bool, Error> {
305        self.flag(["PERSIST", key])
306    }
307
308    /// Sends `LPUSH` and returns the new length.
309    pub fn lpush(&mut self, key: &str, values: &[&str]) -> Result<i64, Error> {
310        self.integer(with(args!["LPUSH", key], values))
311    }
312
313    /// Sends `RPUSH` and returns the new length.
314    pub fn rpush(&mut self, key: &str, values: &[&str]) -> Result<i64, Error> {
315        self.integer(with(args!["RPUSH", key], values))
316    }
317
318    /// Sends `LPOP` for one element.
319    pub fn lpop(&mut self, key: &str) -> Result<Option<String>, Error> {
320        self.bulk_or_null(["LPOP", key])
321    }
322
323    /// Sends `LPOP` with a count. A missing list returns an empty vector.
324    pub fn lpop_count(&mut self, key: &str, count: u64) -> Result<Vec<String>, Error> {
325        self.strings(args!["LPOP", key, count])
326    }
327
328    /// Sends `RPOP` for one element.
329    pub fn rpop(&mut self, key: &str) -> Result<Option<String>, Error> {
330        self.bulk_or_null(["RPOP", key])
331    }
332
333    /// Sends `RPOP` with a count. A missing list returns an empty vector.
334    pub fn rpop_count(&mut self, key: &str, count: u64) -> Result<Vec<String>, Error> {
335        self.strings(args!["RPOP", key, count])
336    }
337
338    /// Sends `LLEN`.
339    pub fn llen(&mut self, key: &str) -> Result<i64, Error> {
340        self.integer(["LLEN", key])
341    }
342
343    /// Sends `LINDEX`.
344    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    /// Sends `LRANGE`. `stop` is inclusive and negative indexes count from the tail.
349    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    /// Sends `LTRIM`.
354    pub fn ltrim(&mut self, key: &str, start: i64, stop: i64) -> Result<(), Error> {
355        self.ok(args!["LTRIM", key, start, stop])
356    }
357
358    /// Sends `BLPOP` for one key. Returns `(key, value)`, or `None` when
359    /// `timeout` passes. `Duration::ZERO` waits forever.
360    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    /// Sends `BLPOP`. Returns `(key, value)`, or `None` when `timeout` passes.
369    /// `Duration::ZERO` waits forever. The timeout rounds up to whole seconds.
370    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    /// Sends `BRPOP` for one key. Returns `(key, value)`, or `None` when
379    /// `timeout` passes.
380    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    /// Sends `BRPOP`. Returns `(key, value)`, or `None` when `timeout` passes.
389    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    /// Sends `SADD` and returns how many members were new.
398    pub fn sadd(&mut self, key: &str, members: &[&str]) -> Result<i64, Error> {
399        self.integer(with(args!["SADD", key], members))
400    }
401
402    /// Sends `SREM` and returns how many members were removed.
403    pub fn srem(&mut self, key: &str, members: &[&str]) -> Result<i64, Error> {
404        self.integer(with(args!["SREM", key], members))
405    }
406
407    /// Sends `SISMEMBER`.
408    pub fn sismember(&mut self, key: &str, member: &str) -> Result<bool, Error> {
409        self.flag(["SISMEMBER", key, member])
410    }
411
412    /// Sends `SCARD`.
413    pub fn scard(&mut self, key: &str) -> Result<i64, Error> {
414        self.integer(["SCARD", key])
415    }
416
417    /// Sends `SMEMBERS`.
418    pub fn smembers(&mut self, key: &str) -> Result<Vec<String>, Error> {
419        self.strings(["SMEMBERS", key])
420    }
421
422    /// Sends `SINTER`.
423    pub fn sinter(&mut self, keys: &[&str]) -> Result<Vec<String>, Error> {
424        self.strings(join("SINTER", keys))
425    }
426
427    /// Sends `SUNION`.
428    pub fn sunion(&mut self, keys: &[&str]) -> Result<Vec<String>, Error> {
429        self.strings(join("SUNION", keys))
430    }
431
432    /// Sends `SDIFF`.
433    pub fn sdiff(&mut self, keys: &[&str]) -> Result<Vec<String>, Error> {
434        self.strings(join("SDIFF", keys))
435    }
436
437    /// Sends `SINTERSTORE` and returns the size of `destination`.
438    pub fn sinterstore(&mut self, destination: &str, keys: &[&str]) -> Result<i64, Error> {
439        self.integer(with(args!["SINTERSTORE", destination], keys))
440    }
441
442    /// Sends `SUNIONSTORE` and returns the size of `destination`.
443    pub fn sunionstore(&mut self, destination: &str, keys: &[&str]) -> Result<i64, Error> {
444        self.integer(with(args!["SUNIONSTORE", destination], keys))
445    }
446
447    /// Sends `SDIFFSTORE` and returns the size of `destination`.
448    pub fn sdiffstore(&mut self, destination: &str, keys: &[&str]) -> Result<i64, Error> {
449        self.integer(with(args!["SDIFFSTORE", destination], keys))
450    }
451
452    /// Sends `SMOVE`. False when `member` was not in `source`.
453    pub fn smove(&mut self, source: &str, destination: &str, member: &str) -> Result<bool, Error> {
454        self.flag(["SMOVE", source, destination, member])
455    }
456
457    /// Sends `SPOP` for one member.
458    pub fn spop(&mut self, key: &str) -> Result<Option<String>, Error> {
459        self.bulk_or_null(["SPOP", key])
460    }
461
462    /// Sends `SPOP` with a count and returns the removed members.
463    pub fn spop_count(&mut self, key: &str, count: u64) -> Result<Vec<String>, Error> {
464        self.strings(args!["SPOP", key, count])
465    }
466
467    /// Sends `SRANDMEMBER` for one member.
468    pub fn srandmember(&mut self, key: &str) -> Result<Option<String>, Error> {
469        self.bulk_or_null(["SRANDMEMBER", key])
470    }
471
472    /// Sends `SRANDMEMBER` with a count. A positive count returns distinct members;
473    /// a negative count may repeat members.
474    pub fn srandmember_count(&mut self, key: &str, count: i64) -> Result<Vec<String>, Error> {
475        self.strings(args!["SRANDMEMBER", key, count])
476    }
477
478    /// Sends `HSET` and returns how many fields were new.
479    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    /// Sends `HGET`.
484    pub fn hget(&mut self, key: &str, field: &str) -> Result<Option<String>, Error> {
485        self.bulk_or_null(["HGET", key, field])
486    }
487
488    /// Sends `HDEL` and returns how many fields were removed.
489    pub fn hdel(&mut self, key: &str, fields: &[&str]) -> Result<i64, Error> {
490        self.integer(with(args!["HDEL", key], fields))
491    }
492
493    /// Sends `HLEN`.
494    pub fn hlen(&mut self, key: &str) -> Result<i64, Error> {
495        self.integer(["HLEN", key])
496    }
497
498    /// Sends `HGETALL`.
499    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    /// Sends `HMGET`. A missing field is `None`.
506    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    /// Sends `HEXISTS`.
511    pub fn hexists(&mut self, key: &str, field: &str) -> Result<bool, Error> {
512        self.flag(["HEXISTS", key, field])
513    }
514
515    /// Sends `HKEYS`.
516    pub fn hkeys(&mut self, key: &str) -> Result<Vec<String>, Error> {
517        self.strings(["HKEYS", key])
518    }
519
520    /// Sends `HVALS`.
521    pub fn hvals(&mut self, key: &str) -> Result<Vec<String>, Error> {
522        self.strings(["HVALS", key])
523    }
524
525    /// Sends `HINCRBY` and returns the new value.
526    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    /// Sends `HSETNX`. False when the field already existed.
531    pub fn hsetnx(&mut self, key: &str, field: &str, value: &str) -> Result<bool, Error> {
532        self.flag(["HSETNX", key, field, value])
533    }
534
535    /// Sends `HSTRLEN`.
536    pub fn hstrlen(&mut self, key: &str, field: &str) -> Result<i64, Error> {
537        self.integer(["HSTRLEN", key, field])
538    }
539
540    /// Sets a TTL on hash fields. Sends `HEXPIRE` for whole seconds, otherwise `HPEXPIRE`.
541    /// One reply per field: `1` set, `0` condition not met, `2` deleted by a zero TTL,
542    /// `-2` no such field or key.
543    ///
544    /// ```no_run
545    /// # use std::time::Duration;
546    /// # let mut db = ruvio_client::Client::connect("127.0.0.1", 6379)?;
547    /// db.hset("session:42", &[("token", "abc"), ("user", "ada")])?;
548    /// db.hexpire("session:42", Duration::from_secs(900), &["token"])?;
549    /// # Ok::<(), ruvio_client::Error>(())
550    /// ```
551    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    /// Like [`Client::hexpire`] with `NX`, `XX`, `GT`, or `LT`.
561    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    /// Sends `HPEXPIRE`: a TTL in milliseconds on hash fields.
578    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    /// Sets an absolute deadline on hash fields. Sends `HEXPIREAT` for a whole second,
592    /// otherwise `HPEXPIREAT`. A time in the past deletes the field (reply `2`).
593    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    /// Sends `HTTL`: seconds left per field, `-1` without a TTL, `-2` when missing.
610    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    /// Sends `HPTTL`: milliseconds left per field, `-1` without a TTL, `-2` when missing.
619    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    /// Sends `HEXPIRETIME`: the Unix second each field expires at, `-1` or `-2` as `HTTL`.
628    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    /// Sends `HPEXPIRETIME`: the Unix millisecond each field expires at.
637    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    /// Sends `HPERSIST`: `1` removed a TTL, `-1` the field had none, `-2` it is missing.
646    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    /// Sends `HSCAN` with an optional `MATCH` pattern and `COUNT` hint.
655    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    /// Sends `ZADD` with `(member, score)` pairs and returns how many members were new.
678    ///
679    /// ```no_run
680    /// # let mut db = ruvio_client::Client::connect("127.0.0.1", 6379)?;
681    /// db.zadd("board", &[("ada", 100.0), ("bob", 50.0)])?;
682    /// # Ok::<(), ruvio_client::Error>(())
683    /// ```
684    pub fn zadd(&mut self, key: &str, entries: &[(&str, f64)]) -> Result<i64, Error> {
685        self.zadd_with(key, SortedSetAddOptions::default(), entries)
686    }
687
688    /// Sends `ZADD` with `NX`, `XX`, `GT`, `LT`, or `CH`.
689    /// Without [`SortedSetAddOptions::changed`] the reply counts new members.
690    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    /// Sends `ZADD ... INCR`. Returns the new score, or `None` when an option skipped the update.
707    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    /// Sends `ZINCRBY` and returns the new score.
723    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    /// Sends `GEOADD`. Each place is longitude, latitude, then member.
728    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    /// Sends `GEODIST`. `unit` is `m`, `km`, `ft`, or `mi`. `None` asks for meters.
741    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    /// Sends `GEOHASH`. A missing member is `None`.
760    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    /// Sends `GEOPOS`. A missing member is `None`.
768    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    /// Sends `GEOSEARCH ... FROMLONLAT ... BYRADIUS ... ASC`.
809    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    /// Sends `GEOSEARCHSTORE` from a member and a radius. Returns how many members were stored.
831    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    /// Sends `GEORADIUSBYMEMBER ... ASC`.
852    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    /// Sends `GEORADIUS ... ASC`.
863    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    /// Sends `ZRANGE` by rank. `stop` is inclusive.
883    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    /// Sends `ZREVRANGE` by rank. `stop` is inclusive.
888    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    /// Sends `ZRANGE ... WITHSCORES`.
893    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    /// Sends `ZREVRANGE ... WITHSCORES`.
903    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    /// Sends `ZRANGE min max BYLEX`. Bounds are `-`, `+`, `[member`, or `(member`.
913    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    /// Sends `ZRANGEBYSCORE ... WITHSCORES`. `min` and `max` accept `-inf`, `+inf`, and `(` for exclusive bounds.
918    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    /// Sends `ZREM` and returns how many members were removed.
928    pub fn zrem(&mut self, key: &str, members: &[&str]) -> Result<i64, Error> {
929        self.integer(with(args!["ZREM", key], members))
930    }
931
932    /// Sends `ZPOPMIN key count`: removes and returns up to `count` lowest-scored members.
933    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    /// Sends `ZPOPMAX key count`: removes and returns up to `count` highest-scored members.
938    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    /// Sends `BZPOPMIN`. Returns the key and the popped member, or `None` when `timeout`
943    /// passes. `Duration::ZERO` waits forever. The timeout rounds up to whole seconds.
944    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    /// Sends `BZPOPMAX`. Like [`Client::bzpopmin`], highest score first.
953    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    /// Sends `ZREMRANGEBYRANK`. `stop` is inclusive; negative ranks count from the end.
962    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    /// Sends `ZREMRANGEBYSCORE`. Bounds accept `-inf`, `+inf`, and `(` for exclusive.
967    pub fn zremrangebyscore(&mut self, key: &str, min: &str, max: &str) -> Result<i64, Error> {
968        self.integer(["ZREMRANGEBYSCORE", key, min, max])
969    }
970
971    /// Sends `ZREMRANGEBYLEX`. Bounds are `-`, `+`, `[member`, or `(member`.
972    pub fn zremrangebylex(&mut self, key: &str, min: &str, max: &str) -> Result<i64, Error> {
973        self.integer(["ZREMRANGEBYLEX", key, min, max])
974    }
975
976    /// Sends `ZCARD`.
977    pub fn zcard(&mut self, key: &str) -> Result<i64, Error> {
978        self.integer(["ZCARD", key])
979    }
980
981    /// Sends `ZSCORE`.
982    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    /// Sends `ZRANK`: the 0-based rank from the lowest score.
989    pub fn zrank(&mut self, key: &str, member: &str) -> Result<Option<i64>, Error> {
990        self.optional_integer(["ZRANK", key, member])
991    }
992
993    /// Sends `ZREVRANK`: the 0-based rank from the highest score.
994    pub fn zrevrank(&mut self, key: &str, member: &str) -> Result<Option<i64>, Error> {
995        self.optional_integer(["ZREVRANK", key, member])
996    }
997
998    /// Sends `BF.RESERVE`.
999    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    /// Sends `BF.ADD`. False when the item may already have been added.
1004    pub fn bf_add(&mut self, key: &str, item: &str) -> Result<bool, Error> {
1005        self.flag(["BF.ADD", key, item])
1006    }
1007
1008    /// Sends `BF.EXISTS`. True means "maybe added"; false means "never added".
1009    pub fn bf_exists(&mut self, key: &str, item: &str) -> Result<bool, Error> {
1010        self.flag(["BF.EXISTS", key, item])
1011    }
1012
1013    /// Sends `XADD NOMKSTREAM`. A missing key stays absent and the reply is nil.
1014    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    /// Sends `XADD key DELAY ms`. The entry stays hidden from `XREADGROUP` until the delay passes.
1024    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    /// Sends `XADD` and returns the entry id. Use `*` for `id` to let the server pick one. `ms-*` fixes the millisecond.
1038    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    /// Sends `XADD key MAXLEN n`: the stream keeps exactly the newest `max_len` entries.
1043    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    /// Sends `XLEN`.
1057    pub fn xlen(&mut self, key: &str) -> Result<i64, Error> {
1058        self.integer(["XLEN", key])
1059    }
1060
1061    /// Sends `XINFO STREAM`.
1062    pub fn xinfo_stream(&mut self, key: &str) -> Result<RespValue, Error> {
1063        self.run(["XINFO", "STREAM", key])
1064    }
1065
1066    /// Sends `XINFO GROUPS`.
1067    pub fn xinfo_groups(&mut self, key: &str) -> Result<RespValue, Error> {
1068        self.run(["XINFO", "GROUPS", key])
1069    }
1070
1071    /// Sends `XINFO CONSUMERS`.
1072    pub fn xinfo_consumers(&mut self, key: &str, group: &str) -> Result<RespValue, Error> {
1073        self.run(["XINFO", "CONSUMERS", key, group])
1074    }
1075
1076    /// Sends `XSETID`. The id must be `ms-seq` and at least the stream's current top id.
1077    pub fn xsetid(&mut self, key: &str, id: &str) -> Result<(), Error> {
1078        self.ok(["XSETID", key, id])
1079    }
1080
1081    /// Sends `XAUTOCLAIM`. `start` is exclusive (`0-0` scans from the beginning).
1082    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    /// Sends `XRANGE` from `start` to `end` (`-` and `+` for the ends), with an optional `COUNT`.
1101    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    /// Sends `XREVRANGE` from `end` down to `start`, with an optional `COUNT`.
1118    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    /// Sends `XDEL` and returns how many entries were removed.
1135    pub fn xdel(&mut self, key: &str, ids: &[&str]) -> Result<i64, Error> {
1136        self.integer(with(args!["XDEL", key], ids))
1137    }
1138
1139    /// Sends `XTRIM MINID`. Entries below `id` are removed. `~` is trimmed exactly.
1140    pub fn xtrim_minid(&mut self, key: &str, id: &str) -> Result<i64, Error> {
1141        self.integer(args!["XTRIM", key, "MINID", id])
1142    }
1143
1144    /// Sends `XTRIM MAXLEN` and returns how many entries were removed.
1145    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    /// Sends `XREAD` for `(key, id)` pairs. A block that times out returns an empty vector.
1150    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    /// Sends `XGROUP CREATECONSUMER`. Returns 1 when the consumer name is new.
1169    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    /// Sends `XGROUP CREATE`. With `make_stream` it adds `MKSTREAM`.
1179    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    /// Sends `XREADGROUP GROUP` with `COUNT`, `BLOCK`, and `NOACK` from `options`.
1196    /// Use `>` as the id for entries never delivered to the group.
1197    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    /// Sends `XGROUP DESTROY`. False when the group did not exist.
1217    pub fn xgroup_destroy(&mut self, key: &str, group: &str) -> Result<bool, Error> {
1218        self.flag(["XGROUP", "DESTROY", key, group])
1219    }
1220
1221    /// Sends `XGROUP SETID`.
1222    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    /// Sends `XGROUP DELCONSUMER` and returns how many pending entries the consumer had.
1227    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    /// Sends `XACK` and returns how many entries were acknowledged.
1237    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    /// Sends the summary form of `XPENDING`: count, smallest id, largest id, and per-consumer counts.
1242    pub fn xpending_summary(&mut self, key: &str, group: &str) -> Result<RespValue, Error> {
1243        self.run(["XPENDING", key, group])
1244    }
1245
1246    /// Sends the extended `XPENDING` form: ids from `start` to `end` (inclusive), at most `count`,
1247    /// narrowed by `filter`.
1248    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    /// Sends `XCLAIM` with `IDLE`, `TIME`, `RETRYCOUNT`, `FORCE`, or `LASTID` from `options`.
1289    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    /// Sends `XCLAIM ... JUSTID` and returns only the claimed ids.
1304    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    /// Sends `INFO`, or `INFO section`. The `prometheus` section returns Prometheus text.
1320    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    /// Sends `CONFIG GET` and returns `(parameter, value)` pairs.
1328    pub fn config_get(&mut self, parameter: &str) -> Result<Vec<(String, String)>, Error> {
1329        pairs(&self.run(["CONFIG", "GET", parameter])?)
1330    }
1331
1332    /// Sends `SAVE`.
1333    pub fn save(&mut self) -> Result<(), Error> {
1334        self.ok(["SAVE"])
1335    }
1336
1337    /// Sends `BGSAVE`.
1338    pub fn bgsave(&mut self) -> Result<(), Error> {
1339        self.ok(["BGSAVE"])
1340    }
1341
1342    /// Sends `FLUSHDB`. Only the selected database is emptied.
1343    pub fn flushdb(&mut self) -> Result<(), Error> {
1344        self.ok(["FLUSHDB"])
1345    }
1346
1347    /// Sends `FLUSHALL`. Every database is emptied.
1348    pub fn flushall(&mut self) -> Result<(), Error> {
1349        self.ok(["FLUSHALL"])
1350    }
1351
1352    /// Sends `SWAPDB`.
1353    pub fn swapdb(&mut self, first: u32, second: u32) -> Result<(), Error> {
1354        self.ok(args!["SWAPDB", first, second])
1355    }
1356
1357    /// Sends `MOVE`. False when the key is missing or already exists in `database`.
1358    pub fn move_key(&mut self, key: &str, database: u32) -> Result<bool, Error> {
1359        self.flag(args!["MOVE", key, database])
1360    }
1361
1362    /// Sends `MULTI`. Until [`Self::exec`], the server answers `QUEUED`, so queue commands
1363    /// with [`Client::execute`] rather than the typed methods.
1364    pub fn multi(&mut self) -> Result<(), Error> {
1365        self.ok(["MULTI"])
1366    }
1367
1368    /// Sends `EXEC`. `None` when a watched key changed and the transaction was aborted.
1369    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    /// Sends `DISCARD`.
1378    pub fn discard(&mut self) -> Result<(), Error> {
1379        self.ok(["DISCARD"])
1380    }
1381
1382    /// Sends `WATCH` for one key.
1383    pub fn watch_key(&mut self, key: &str) -> Result<(), Error> {
1384        self.watch(&[key])
1385    }
1386
1387    /// Sends `WATCH`.
1388    pub fn watch(&mut self, keys: &[&str]) -> Result<(), Error> {
1389        self.ok(join("WATCH", keys))
1390    }
1391
1392    /// Sends `UNWATCH`.
1393    pub fn unwatch(&mut self) -> Result<(), Error> {
1394        self.ok(["UNWATCH"])
1395    }
1396
1397    /// Sends `EVAL`. The key count is `keys.len()`.
1398    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    /// Sends `EVALSHA`. The key count is `keys.len()`.
1408    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    /// Sends `SCRIPT LOAD` and returns the SHA1 for [`Self::evalsha`].
1418    pub fn script_load(&mut self, script: &str) -> Result<String, Error> {
1419        self.text(["SCRIPT", "LOAD", script])
1420    }
1421
1422    /// Sends `SCRIPT EXISTS`, one flag per hash.
1423    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    /// Sends `SCRIPT FLUSH`.
1435    pub fn script_flush(&mut self) -> Result<(), Error> {
1436        self.ok(["SCRIPT", "FLUSH"])
1437    }
1438
1439    /// Sends `SCRIPT KILL`.
1440    pub fn script_kill(&mut self) -> Result<(), Error> {
1441        self.ok(["SCRIPT", "KILL"])
1442    }
1443
1444    /// Sends `FCALL`. The key count is `keys.len()`.
1445    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    /// Sends `FUNCTION LOAD name body`. The body is plain Lua, not a Redis library shebang,
1455    /// and is lost when the server restarts.
1456    pub fn function_load(&mut self, name: &str, body: &str) -> Result<(), Error> {
1457        self.ok(["FUNCTION", "LOAD", name, body])
1458    }
1459
1460    /// Sends `FUNCTION LIST`.
1461    pub fn function_list(&mut self) -> Result<RespValue, Error> {
1462        self.run(["FUNCTION", "LIST"])
1463    }
1464
1465    /// Sends `FUNCTION DELETE`. False when no function had that name.
1466    pub fn function_delete(&mut self, name: &str) -> Result<bool, Error> {
1467        self.flag(["FUNCTION", "DELETE", name])
1468    }
1469
1470    /// Sends `ACL WHOAMI`.
1471    pub fn acl_whoami(&mut self) -> Result<String, Error> {
1472        self.text(["ACL", "WHOAMI"])
1473    }
1474
1475    /// Sends `ACL USERS`.
1476    pub fn acl_users(&mut self) -> Result<Vec<String>, Error> {
1477        self.strings(["ACL", "USERS"])
1478    }
1479
1480    /// Sends `ACL LIST`.
1481    pub fn acl_list(&mut self) -> Result<Vec<String>, Error> {
1482        self.strings(["ACL", "LIST"])
1483    }
1484
1485    /// Sends `ACL GETUSER`.
1486    pub fn acl_getuser(&mut self, user: &str) -> Result<RespValue, Error> {
1487        self.run(["ACL", "GETUSER", user])
1488    }
1489
1490    /// Sends `ACL CAT`.
1491    pub fn acl_cat(&mut self) -> Result<Vec<String>, Error> {
1492        self.strings(["ACL", "CAT"])
1493    }
1494
1495    /// Sends `ACL SETUSER user rule...`.
1496    pub fn acl_setuser(&mut self, user: &str, rules: &[&str]) -> Result<(), Error> {
1497        self.ok(with(args!["ACL", "SETUSER", user], rules))
1498    }
1499
1500    /// Sends `ACL LOAD`.
1501    pub fn acl_load(&mut self) -> Result<(), Error> {
1502        self.ok(["ACL", "LOAD"])
1503    }
1504
1505    /// Sends `ACL SAVE`.
1506    pub fn acl_save(&mut self) -> Result<(), Error> {
1507        self.ok(["ACL", "SAVE"])
1508    }
1509
1510    /// Sends `CLUSTER SLOTS`.
1511    pub fn cluster_slots(&mut self) -> Result<RespValue, Error> {
1512        self.run(["CLUSTER", "SLOTS"])
1513    }
1514
1515    /// Sends `CLUSTER NODES`.
1516    pub fn cluster_nodes(&mut self) -> Result<String, Error> {
1517        self.text(["CLUSTER", "NODES"])
1518    }
1519
1520    /// Sends `CLUSTER SHARDS`.
1521    pub fn cluster_shards(&mut self) -> Result<RespValue, Error> {
1522        self.run(["CLUSTER", "SHARDS"])
1523    }
1524
1525    /// Sends `CLUSTER INFO`.
1526    pub fn cluster_info(&mut self) -> Result<String, Error> {
1527        self.text(["CLUSTER", "INFO"])
1528    }
1529
1530    /// Sends `CLUSTER MYID`.
1531    pub fn cluster_myid(&mut self) -> Result<String, Error> {
1532        self.text(["CLUSTER", "MYID"])
1533    }
1534
1535    /// Sends `CLUSTER KEYSLOT`.
1536    pub fn cluster_keyslot(&mut self, key: &str) -> Result<i64, Error> {
1537        self.integer(["CLUSTER", "KEYSLOT", key])
1538    }
1539
1540    /// Sends `READONLY`. Ruvio has no replicas; the reply is still `OK` in cluster mode.
1541    pub fn readonly(&mut self) -> Result<(), Error> {
1542        self.ok(["READONLY"])
1543    }
1544
1545    /// Sends `READWRITE`.
1546    pub fn readwrite(&mut self) -> Result<(), Error> {
1547        self.ok(["READWRITE"])
1548    }
1549
1550    /// Sends `CLIENT GETNAME`.
1551    pub fn client_getname(&mut self) -> Result<Option<String>, Error> {
1552        self.bulk_or_null(["CLIENT", "GETNAME"])
1553    }
1554
1555    /// Sends `CLIENT SETNAME`.
1556    pub fn client_setname(&mut self, name: &str) -> Result<(), Error> {
1557        self.ok(["CLIENT", "SETNAME", name])
1558    }
1559
1560    /// Sends `CLIENT TRACKING ON` or `OFF`.
1561    pub fn client_tracking(&mut self, enabled: bool) -> Result<(), Error> {
1562        self.ok(["CLIENT", "TRACKING", if enabled { "ON" } else { "OFF" }])
1563    }
1564
1565    /// Sends `HELLO protocol`. Connecting itself does not send this, and this client
1566    /// only decodes RESP2 replies.
1567    pub fn hello(&mut self, protocol: u8) -> Result<RespValue, Error> {
1568        self.run(args!["HELLO", protocol])
1569    }
1570
1571    /// Sends `PUBLISH` and returns how many subscribers received the message.
1572    pub fn publish(&mut self, channel: &str, message: &str) -> Result<i64, Error> {
1573        self.integer(["PUBLISH", channel, message])
1574    }
1575
1576    /// Sends `SPUBLISH` and returns how many subscribers received the message.
1577    pub fn spublish(&mut self, channel: &str, message: &str) -> Result<i64, Error> {
1578        self.integer(["SPUBLISH", channel, message])
1579    }
1580
1581    /// Sends `SUBSCRIBE` and reads one confirmation per channel.
1582    /// Opens a dedicated Pub/Sub socket so `PING` and `PUBLISH` stay on the
1583    /// command connection. Then read deliveries with [`Self::next_message`].
1584    pub fn subscribe(&mut self, channels: &[&str]) -> Result<Vec<RespValue>, Error> {
1585        self.subscription("SUBSCRIBE", channels, false)
1586    }
1587
1588    /// Sends `UNSUBSCRIBE`. No channels drops every exact subscription on this connection.
1589    pub fn unsubscribe(&mut self, channels: &[&str]) -> Result<Vec<RespValue>, Error> {
1590        self.subscription("UNSUBSCRIBE", channels, true)
1591    }
1592
1593    /// Sends `PSUBSCRIBE` and reads one confirmation per pattern.
1594    pub fn psubscribe(&mut self, patterns: &[&str]) -> Result<Vec<RespValue>, Error> {
1595        self.subscription("PSUBSCRIBE", patterns, false)
1596    }
1597
1598    /// Sends `PUNSUBSCRIBE`. No patterns drops every pattern subscription on this connection.
1599    pub fn punsubscribe(&mut self, patterns: &[&str]) -> Result<Vec<RespValue>, Error> {
1600        self.subscription("PUNSUBSCRIBE", patterns, true)
1601    }
1602
1603    /// Sends `SSUBSCRIBE` and reads one confirmation per channel.
1604    pub fn ssubscribe(&mut self, channels: &[&str]) -> Result<Vec<RespValue>, Error> {
1605        self.subscription("SSUBSCRIBE", channels, false)
1606    }
1607
1608    /// Sends `SUNSUBSCRIBE`. No channels drops every sharded subscription on this connection.
1609    pub fn sunsubscribe(&mut self, channels: &[&str]) -> Result<Vec<RespValue>, Error> {
1610        self.subscription("SUNSUBSCRIBE", channels, true)
1611    }
1612
1613    /// Reads replies until a `message`, `pmessage`, or `smessage` arrives, skipping
1614    /// subscribe confirmations. Blocks until [`Self::set_read_timeout`] passes.
1615    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    /// Calls `on_message` for each delivery. Safe on the command client after
1626    /// `subscribe`; deliveries are read from the dedicated Pub/Sub socket.
1627    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
1716/// Appends an optional `NX`/`XX`/`GT`/`LT`, then `FIELDS numfields field ...`.
1717fn 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}