use std::collections::HashMap;
use std::time::{Duration, SystemTime};
use crate::client::{Client, duration_millis, duration_secs, join, unexpected, unix_millis};
use crate::error::Error;
use crate::models::{
FieldExpireCondition, HashScanPage, PubSubMessage, ScanPage, SortedSetAddOptions,
SortedSetEntry, StreamClaimOptions, StreamEntry, StreamPendingEntry, StreamPendingFilter,
StreamReadOptions, StreamReadResult,
};
use crate::value::RespValue;
macro_rules! args {
($($argument:expr),* $(,)?) => {
vec![$($argument.to_string()),*]
};
}
impl Client {
pub fn mget(&mut self, keys: &[&str]) -> Result<Vec<Option<String>>, Error> {
self.optional_strings(join("MGET", keys))
}
pub fn limit(&mut self, key: &str, max: u64, window_ms: u64) -> Result<Option<i64>, Error> {
let value = self.run(["LIMIT", key, &max.to_string(), &window_ms.to_string()])?;
if value.is_null() {
return Ok(None);
}
value
.as_integer()
.map(Some)
.ok_or_else(|| unexpected("LIMIT", &value))
}
pub fn once(&mut self, key: &str, ttl_ms: u64) -> Result<i64, Error> {
self.integer(["ONCE", key, &ttl_ms.to_string()])
}
pub fn idempotency_begin(
&mut self,
key: &str,
fingerprint: &str,
owner: &str,
ttl_ms: u64,
) -> Result<RespValue, Error> {
self.run([
"IDEM",
key,
"BEGIN",
fingerprint,
owner,
&ttl_ms.to_string(),
])
}
pub fn idempotency_complete(
&mut self,
key: &str,
fingerprint: &str,
owner: &str,
result: &str,
) -> Result<i64, Error> {
self.integer(["IDEM", key, "COMPLETE", fingerprint, owner, result])
}
pub fn take(&mut self, key: &str, count: i64) -> Result<Option<i64>, Error> {
let value = self.run(["TAKE", key, &count.to_string()])?;
if value.is_null() {
return Ok(None);
}
value
.as_integer()
.map(Some)
.ok_or_else(|| unexpected("TAKE", &value))
}
pub fn getex(&mut self, key: &str, ttl_ms: u64) -> Result<Option<String>, Error> {
self.bulk_or_null(["GETEX", key, &ttl_ms.to_string()])
}
pub fn lease(&mut self, key: &str, ttl_ms: u64) -> Result<Option<i64>, Error> {
let value = self.run(["LEASE", key, &ttl_ms.to_string()])?;
if value.is_null() {
return Ok(None);
}
value
.as_integer()
.map(Some)
.ok_or_else(|| unexpected("LEASE", &value))
}
pub fn semaphore(&mut self, key: &str, max: u64, ttl_ms: u64) -> Result<Option<i64>, Error> {
let value = self.run(["SEMAPHORE", key, &max.to_string(), &ttl_ms.to_string()])?;
if value.is_null() {
return Ok(None);
}
value
.as_integer()
.map(Some)
.ok_or_else(|| unexpected("SEMAPHORE", &value))
}
pub fn incr_by_max(&mut self, key: &str, delta: i64, limit: i64) -> Result<Option<i64>, Error> {
let value = self.run(args!["INCRBY", key, delta, "MAX", limit])?;
if value.is_null() {
return Ok(None);
}
value
.as_integer()
.map(Some)
.ok_or_else(|| unexpected("INCRBY", &value))
}
pub fn release(&mut self, key: &str, token: u64) -> Result<i64, Error> {
self.integer(["RELEASE", key, &token.to_string()])
}
pub fn setv(&mut self, key: &str, value: &str, version: Option<u64>) -> Result<i64, Error> {
let mut arguments = vec!["SETV".to_owned(), key.to_owned(), value.to_owned()];
if let Some(version) = version {
arguments.push(version.to_string());
}
self.integer(arguments)
}
pub fn changes_start(&mut self) -> Result<i64, Error> {
self.integer(["CHANGES", "START"])
}
pub fn changes(&mut self, cursor: u64) -> Result<Vec<(i64, String)>, Error> {
let value = self.run(["CHANGES", &cursor.to_string()])?;
let Some(rows) = value.as_array() else {
return Err(unexpected("CHANGES", &value));
};
let mut notes = Vec::new();
for row in rows {
let Some(parts) = row.as_array() else {
return Err(unexpected("CHANGES", row));
};
let sequence = parts
.first()
.and_then(|part| part.as_integer())
.ok_or_else(|| unexpected("CHANGES", row))?;
let text = parts
.get(1)
.and_then(|part| part.as_string().ok().flatten())
.unwrap_or_default();
notes.push((sequence, text));
}
Ok(notes)
}
pub fn mset(&mut self, pairs: &[(&str, &str)]) -> Result<(), Error> {
self.ok(with_pairs(args!["MSET"], pairs))
}
pub fn getset(&mut self, key: &str, value: &str) -> Result<Option<String>, Error> {
self.bulk_or_null(["GETSET", key, value])
}
pub fn append(&mut self, key: &str, value: &str) -> Result<i64, Error> {
self.integer(["APPEND", key, value])
}
pub fn strlen(&mut self, key: &str) -> Result<i64, Error> {
self.integer(["STRLEN", key])
}
pub fn decr(&mut self, key: &str) -> Result<i64, Error> {
self.integer(["DECR", key])
}
pub fn incr_by(&mut self, key: &str, delta: i64) -> Result<i64, Error> {
self.integer(args!["INCRBY", key, delta])
}
pub fn decr_by(&mut self, key: &str, delta: i64) -> Result<i64, Error> {
self.integer(args!["DECRBY", key, delta])
}
pub fn unlink_key(&mut self, key: &str) -> Result<i64, Error> {
self.unlink(&[key])
}
pub fn unlink(&mut self, keys: &[&str]) -> Result<i64, Error> {
self.integer(join("UNLINK", keys))
}
pub fn exists_key(&mut self, key: &str) -> Result<i64, Error> {
self.exists(&[key])
}
pub fn exists(&mut self, keys: &[&str]) -> Result<i64, Error> {
self.integer(join("EXISTS", keys))
}
pub fn key_type(&mut self, key: &str) -> Result<String, Error> {
self.text(["TYPE", key])
}
pub fn rename(&mut self, key: &str, new_key: &str) -> Result<(), Error> {
self.ok(["RENAME", key, new_key])
}
pub fn scan(
&mut self,
cursor: u64,
pattern: Option<&str>,
count: Option<u64>,
) -> Result<ScanPage, Error> {
let mut arguments = args!["SCAN", cursor];
push_scan_options(&mut arguments, pattern, count);
let (cursor, keys) = parse_scan("SCAN", self.run(arguments)?)?;
Ok(ScanPage { cursor, keys })
}
pub fn dbsize(&mut self) -> Result<i64, Error> {
self.integer(["DBSIZE"])
}
pub fn expire_at(&mut self, key: &str, when: SystemTime) -> Result<bool, Error> {
let millis = unix_millis(when)?;
if millis % 1000 == 0 {
return self.flag(args!["EXPIREAT", key, millis / 1000]);
}
self.flag(args!["PEXPIREAT", key, millis])
}
pub fn pexpire(&mut self, key: &str, ttl: Duration) -> Result<bool, Error> {
self.flag(args!["PEXPIRE", key, duration_millis(ttl)])
}
pub fn ttl(&mut self, key: &str) -> Result<i64, Error> {
self.integer(["TTL", key])
}
pub fn pttl(&mut self, key: &str) -> Result<i64, Error> {
self.integer(["PTTL", key])
}
pub fn persist(&mut self, key: &str) -> Result<bool, Error> {
self.flag(["PERSIST", key])
}
pub fn lpush(&mut self, key: &str, values: &[&str]) -> Result<i64, Error> {
self.integer(with(args!["LPUSH", key], values))
}
pub fn rpush(&mut self, key: &str, values: &[&str]) -> Result<i64, Error> {
self.integer(with(args!["RPUSH", key], values))
}
pub fn lpop(&mut self, key: &str) -> Result<Option<String>, Error> {
self.bulk_or_null(["LPOP", key])
}
pub fn lpop_count(&mut self, key: &str, count: u64) -> Result<Vec<String>, Error> {
self.strings(args!["LPOP", key, count])
}
pub fn rpop(&mut self, key: &str) -> Result<Option<String>, Error> {
self.bulk_or_null(["RPOP", key])
}
pub fn rpop_count(&mut self, key: &str, count: u64) -> Result<Vec<String>, Error> {
self.strings(args!["RPOP", key, count])
}
pub fn llen(&mut self, key: &str) -> Result<i64, Error> {
self.integer(["LLEN", key])
}
pub fn lindex(&mut self, key: &str, index: i64) -> Result<Option<String>, Error> {
self.bulk_or_null(args!["LINDEX", key, index])
}
pub fn lrange(&mut self, key: &str, start: i64, stop: i64) -> Result<Vec<String>, Error> {
self.strings(args!["LRANGE", key, start, stop])
}
pub fn ltrim(&mut self, key: &str, start: i64, stop: i64) -> Result<(), Error> {
self.ok(args!["LTRIM", key, start, stop])
}
pub fn blpop_key(
&mut self,
key: &str,
timeout: Duration,
) -> Result<Option<(String, String)>, Error> {
self.blpop(&[key], timeout)
}
pub fn blpop(
&mut self,
keys: &[&str],
timeout: Duration,
) -> Result<Option<(String, String)>, Error> {
self.blocking_pop("BLPOP", keys, timeout)
}
pub fn brpop_key(
&mut self,
key: &str,
timeout: Duration,
) -> Result<Option<(String, String)>, Error> {
self.brpop(&[key], timeout)
}
pub fn brpop(
&mut self,
keys: &[&str],
timeout: Duration,
) -> Result<Option<(String, String)>, Error> {
self.blocking_pop("BRPOP", keys, timeout)
}
pub fn sadd(&mut self, key: &str, members: &[&str]) -> Result<i64, Error> {
self.integer(with(args!["SADD", key], members))
}
pub fn srem(&mut self, key: &str, members: &[&str]) -> Result<i64, Error> {
self.integer(with(args!["SREM", key], members))
}
pub fn sismember(&mut self, key: &str, member: &str) -> Result<bool, Error> {
self.flag(["SISMEMBER", key, member])
}
pub fn scard(&mut self, key: &str) -> Result<i64, Error> {
self.integer(["SCARD", key])
}
pub fn smembers(&mut self, key: &str) -> Result<Vec<String>, Error> {
self.strings(["SMEMBERS", key])
}
pub fn sinter(&mut self, keys: &[&str]) -> Result<Vec<String>, Error> {
self.strings(join("SINTER", keys))
}
pub fn sunion(&mut self, keys: &[&str]) -> Result<Vec<String>, Error> {
self.strings(join("SUNION", keys))
}
pub fn sdiff(&mut self, keys: &[&str]) -> Result<Vec<String>, Error> {
self.strings(join("SDIFF", keys))
}
pub fn sinterstore(&mut self, destination: &str, keys: &[&str]) -> Result<i64, Error> {
self.integer(with(args!["SINTERSTORE", destination], keys))
}
pub fn sunionstore(&mut self, destination: &str, keys: &[&str]) -> Result<i64, Error> {
self.integer(with(args!["SUNIONSTORE", destination], keys))
}
pub fn sdiffstore(&mut self, destination: &str, keys: &[&str]) -> Result<i64, Error> {
self.integer(with(args!["SDIFFSTORE", destination], keys))
}
pub fn smove(&mut self, source: &str, destination: &str, member: &str) -> Result<bool, Error> {
self.flag(["SMOVE", source, destination, member])
}
pub fn spop(&mut self, key: &str) -> Result<Option<String>, Error> {
self.bulk_or_null(["SPOP", key])
}
pub fn spop_count(&mut self, key: &str, count: u64) -> Result<Vec<String>, Error> {
self.strings(args!["SPOP", key, count])
}
pub fn srandmember(&mut self, key: &str) -> Result<Option<String>, Error> {
self.bulk_or_null(["SRANDMEMBER", key])
}
pub fn srandmember_count(&mut self, key: &str, count: i64) -> Result<Vec<String>, Error> {
self.strings(args!["SRANDMEMBER", key, count])
}
pub fn hset(&mut self, key: &str, fields: &[(&str, &str)]) -> Result<i64, Error> {
self.integer(with_pairs(args!["HSET", key], fields))
}
pub fn hget(&mut self, key: &str, field: &str) -> Result<Option<String>, Error> {
self.bulk_or_null(["HGET", key, field])
}
pub fn hdel(&mut self, key: &str, fields: &[&str]) -> Result<i64, Error> {
self.integer(with(args!["HDEL", key], fields))
}
pub fn hlen(&mut self, key: &str) -> Result<i64, Error> {
self.integer(["HLEN", key])
}
pub fn hgetall(&mut self, key: &str) -> Result<HashMap<String, String>, Error> {
let pairs = pairs(&self.run(["HGETALL", key])?)?;
Ok(pairs.into_iter().collect())
}
pub fn hmget(&mut self, key: &str, fields: &[&str]) -> Result<Vec<Option<String>>, Error> {
self.optional_strings(with(args!["HMGET", key], fields))
}
pub fn hexists(&mut self, key: &str, field: &str) -> Result<bool, Error> {
self.flag(["HEXISTS", key, field])
}
pub fn hkeys(&mut self, key: &str) -> Result<Vec<String>, Error> {
self.strings(["HKEYS", key])
}
pub fn hvals(&mut self, key: &str) -> Result<Vec<String>, Error> {
self.strings(["HVALS", key])
}
pub fn hincr_by(&mut self, key: &str, field: &str, delta: i64) -> Result<i64, Error> {
self.integer(args!["HINCRBY", key, field, delta])
}
pub fn hsetnx(&mut self, key: &str, field: &str, value: &str) -> Result<bool, Error> {
self.flag(["HSETNX", key, field, value])
}
pub fn hstrlen(&mut self, key: &str, field: &str) -> Result<i64, Error> {
self.integer(["HSTRLEN", key, field])
}
pub fn hexpire(
&mut self,
key: &str,
ttl: Duration,
fields: &[&str],
) -> Result<Vec<i64>, Error> {
self.hexpire_with(key, ttl, FieldExpireCondition::Always, fields)
}
pub fn hexpire_with(
&mut self,
key: &str,
ttl: Duration,
condition: FieldExpireCondition,
fields: &[&str],
) -> Result<Vec<i64>, Error> {
let arguments = if ttl.subsec_nanos() == 0 {
args!["HEXPIRE", key, duration_secs(ttl)]
} else {
args!["HPEXPIRE", key, duration_millis(ttl)]
};
self.integers(field_block(arguments, condition, fields))
}
pub fn hpexpire(
&mut self,
key: &str,
ttl: Duration,
fields: &[&str],
) -> Result<Vec<i64>, Error> {
self.integers(field_block(
args!["HPEXPIRE", key, duration_millis(ttl)],
FieldExpireCondition::Always,
fields,
))
}
pub fn hexpire_at(
&mut self,
key: &str,
when: SystemTime,
fields: &[&str],
) -> Result<Vec<i64>, Error> {
let millis = unix_millis(when)?;
let arguments = if millis % 1000 == 0 {
args!["HEXPIREAT", key, millis / 1000]
} else {
args!["HPEXPIREAT", key, millis]
};
self.integers(field_block(arguments, FieldExpireCondition::Always, fields))
}
pub fn httl(&mut self, key: &str, fields: &[&str]) -> Result<Vec<i64>, Error> {
self.integers(field_block(
args!["HTTL", key],
FieldExpireCondition::Always,
fields,
))
}
pub fn hpttl(&mut self, key: &str, fields: &[&str]) -> Result<Vec<i64>, Error> {
self.integers(field_block(
args!["HPTTL", key],
FieldExpireCondition::Always,
fields,
))
}
pub fn hexpiretime(&mut self, key: &str, fields: &[&str]) -> Result<Vec<i64>, Error> {
self.integers(field_block(
args!["HEXPIRETIME", key],
FieldExpireCondition::Always,
fields,
))
}
pub fn hpexpiretime(&mut self, key: &str, fields: &[&str]) -> Result<Vec<i64>, Error> {
self.integers(field_block(
args!["HPEXPIRETIME", key],
FieldExpireCondition::Always,
fields,
))
}
pub fn hpersist(&mut self, key: &str, fields: &[&str]) -> Result<Vec<i64>, Error> {
self.integers(field_block(
args!["HPERSIST", key],
FieldExpireCondition::Always,
fields,
))
}
pub fn hscan(
&mut self,
key: &str,
cursor: u64,
pattern: Option<&str>,
count: Option<u64>,
) -> Result<HashScanPage, Error> {
let mut arguments = args!["HSCAN", key, cursor];
push_scan_options(&mut arguments, pattern, count);
let (cursor, flat) = parse_scan("HSCAN", self.run(arguments)?)?;
let mut fields = Vec::with_capacity(flat.len() / 2);
let mut items = flat.into_iter();
while let (Some(field), Some(value)) = (items.next(), items.next()) {
fields.push((field, value));
}
Ok(HashScanPage { cursor, fields })
}
pub fn zadd(&mut self, key: &str, entries: &[(&str, f64)]) -> Result<i64, Error> {
self.zadd_with(key, SortedSetAddOptions::default(), entries)
}
pub fn zadd_with(
&mut self,
key: &str,
options: SortedSetAddOptions,
entries: &[(&str, f64)],
) -> Result<i64, Error> {
let mut arguments = zadd_arguments(key, options);
for (member, score) in entries {
arguments.push(score.to_string());
arguments.push((*member).to_owned());
}
self.integer(arguments)
}
pub fn zadd_incr(
&mut self,
key: &str,
member: &str,
delta: f64,
options: SortedSetAddOptions,
) -> Result<Option<f64>, Error> {
let mut arguments = zadd_arguments(key, options);
arguments.extend(args!["INCR", delta, member]);
self.bulk_or_null(arguments)?
.map(|text| parse_score(&text))
.transpose()
}
pub fn zincr_by(&mut self, key: &str, member: &str, delta: f64) -> Result<f64, Error> {
parse_score(&self.text(args!["ZINCRBY", key, delta, member])?)
}
pub fn geo_add(&mut self, key: &str, places: &[(f64, f64, &str)]) -> Result<i64, Error> {
let mut arguments = args!["GEOADD", key];
for (longitude, latitude, member) in places {
arguments.push(longitude.to_string());
arguments.push(latitude.to_string());
arguments.push((*member).to_owned());
}
self.integer(arguments)
}
pub fn geo_dist(
&mut self,
key: &str,
from: &str,
to: &str,
unit: Option<&str>,
) -> Result<Option<f64>, Error> {
let mut arguments = args!["GEODIST", key, from, to];
if let Some(unit) = unit {
arguments.push(unit.to_owned());
}
self.bulk_or_null(arguments)?
.map(|text| parse_score(&text))
.transpose()
}
pub fn geo_hash(&mut self, key: &str, members: &[&str]) -> Result<Vec<Option<String>>, Error> {
let mut arguments = args!["GEOHASH", key];
arguments.extend(members.iter().map(|member| (*member).to_owned()));
self.optional_strings(arguments)
}
pub fn geo_pos(
&mut self,
key: &str,
members: &[&str],
) -> Result<Vec<Option<(f64, f64)>>, Error> {
let mut arguments = args!["GEOPOS", key];
arguments.extend(members.iter().map(|member| (*member).to_owned()));
let reply = self.run(arguments)?;
let Some(items) = reply.as_array() else {
return Err(unexpected("GEOPOS", &reply));
};
let mut points = Vec::with_capacity(items.len());
for item in items {
if item.is_null() {
points.push(None);
continue;
}
let Some(pair) = item.as_array() else {
return Err(unexpected("GEOPOS", item));
};
if pair.len() != 2 {
return Err(unexpected("GEOPOS", item));
}
let longitude = parse_score(&pair[0].as_string()?.unwrap_or_default())?;
let latitude = parse_score(&pair[1].as_string()?.unwrap_or_default())?;
points.push(Some((longitude, latitude)));
}
Ok(points)
}
pub fn geo_search(
&mut self,
key: &str,
longitude: f64,
latitude: f64,
radius: f64,
unit: &str,
) -> Result<Vec<String>, Error> {
self.strings(args![
"GEOSEARCH",
key,
"FROMLONLAT",
longitude,
latitude,
"BYRADIUS",
radius,
unit,
"ASC"
])
}
pub fn geo_search_store(
&mut self,
destination: &str,
source: &str,
member: &str,
radius: f64,
unit: &str,
) -> Result<i64, Error> {
self.integer(args![
"GEOSEARCHSTORE",
destination,
source,
"FROMMEMBER",
member,
"BYRADIUS",
radius,
unit
])
}
pub fn geo_radius_by_member(
&mut self,
key: &str,
member: &str,
radius: f64,
unit: &str,
) -> Result<Vec<String>, Error> {
self.strings(args!["GEORADIUSBYMEMBER", key, member, radius, unit, "ASC"])
}
pub fn geo_radius(
&mut self,
key: &str,
longitude: f64,
latitude: f64,
radius: f64,
unit: &str,
) -> Result<Vec<String>, Error> {
self.strings(args![
"GEORADIUS",
key,
longitude,
latitude,
radius,
unit,
"ASC"
])
}
pub fn zrange(&mut self, key: &str, start: i64, stop: i64) -> Result<Vec<String>, Error> {
self.strings(args!["ZRANGE", key, start, stop])
}
pub fn zrevrange(&mut self, key: &str, start: i64, stop: i64) -> Result<Vec<String>, Error> {
self.strings(args!["ZREVRANGE", key, start, stop])
}
pub fn zrange_with_scores(
&mut self,
key: &str,
start: i64,
stop: i64,
) -> Result<Vec<SortedSetEntry>, Error> {
sorted_set(&self.run(args!["ZRANGE", key, start, stop, "WITHSCORES"])?)
}
pub fn zrevrange_with_scores(
&mut self,
key: &str,
start: i64,
stop: i64,
) -> Result<Vec<SortedSetEntry>, Error> {
sorted_set(&self.run(args!["ZREVRANGE", key, start, stop, "WITHSCORES"])?)
}
pub fn zrange_by_lex(&mut self, key: &str, min: &str, max: &str) -> Result<Vec<String>, Error> {
self.strings(["ZRANGE", key, min, max, "BYLEX"])
}
pub fn zrange_by_score_with_scores(
&mut self,
key: &str,
min: &str,
max: &str,
) -> Result<Vec<SortedSetEntry>, Error> {
sorted_set(&self.run(["ZRANGEBYSCORE", key, min, max, "WITHSCORES"])?)
}
pub fn zrem(&mut self, key: &str, members: &[&str]) -> Result<i64, Error> {
self.integer(with(args!["ZREM", key], members))
}
pub fn zpopmin(&mut self, key: &str, count: u64) -> Result<Vec<SortedSetEntry>, Error> {
sorted_set(&self.run(args!["ZPOPMIN", key, count])?)
}
pub fn zpopmax(&mut self, key: &str, count: u64) -> Result<Vec<SortedSetEntry>, Error> {
sorted_set(&self.run(args!["ZPOPMAX", key, count])?)
}
pub fn bzpopmin(
&mut self,
keys: &[&str],
timeout: Duration,
) -> Result<Option<(String, SortedSetEntry)>, Error> {
self.blocking_sorted_pop("BZPOPMIN", keys, timeout)
}
pub fn bzpopmax(
&mut self,
keys: &[&str],
timeout: Duration,
) -> Result<Option<(String, SortedSetEntry)>, Error> {
self.blocking_sorted_pop("BZPOPMAX", keys, timeout)
}
pub fn zremrangebyrank(&mut self, key: &str, start: i64, stop: i64) -> Result<i64, Error> {
self.integer(args!["ZREMRANGEBYRANK", key, start, stop])
}
pub fn zremrangebyscore(&mut self, key: &str, min: &str, max: &str) -> Result<i64, Error> {
self.integer(["ZREMRANGEBYSCORE", key, min, max])
}
pub fn zremrangebylex(&mut self, key: &str, min: &str, max: &str) -> Result<i64, Error> {
self.integer(["ZREMRANGEBYLEX", key, min, max])
}
pub fn zcard(&mut self, key: &str) -> Result<i64, Error> {
self.integer(["ZCARD", key])
}
pub fn zscore(&mut self, key: &str, member: &str) -> Result<Option<f64>, Error> {
self.bulk_or_null(["ZSCORE", key, member])?
.map(|text| parse_score(&text))
.transpose()
}
pub fn zrank(&mut self, key: &str, member: &str) -> Result<Option<i64>, Error> {
self.optional_integer(["ZRANK", key, member])
}
pub fn zrevrank(&mut self, key: &str, member: &str) -> Result<Option<i64>, Error> {
self.optional_integer(["ZREVRANK", key, member])
}
pub fn bf_reserve(&mut self, key: &str, error_rate: f64, capacity: u64) -> Result<(), Error> {
self.ok(args!["BF.RESERVE", key, error_rate, capacity])
}
pub fn bf_add(&mut self, key: &str, item: &str) -> Result<bool, Error> {
self.flag(["BF.ADD", key, item])
}
pub fn bf_exists(&mut self, key: &str, item: &str) -> Result<bool, Error> {
self.flag(["BF.EXISTS", key, item])
}
pub fn xadd_nomkstream(
&mut self,
key: &str,
id: &str,
fields: &[(&str, &str)],
) -> Result<Option<String>, Error> {
self.bulk_or_null(with_pairs(args!["XADD", key, "NOMKSTREAM", id], fields))
}
pub fn xadd_delay(
&mut self,
key: &str,
delay_ms: u64,
id: &str,
fields: &[(&str, &str)],
) -> Result<String, Error> {
self.text(with_pairs(
args!["XADD", key, "DELAY", delay_ms, id],
fields,
))
}
pub fn xadd(&mut self, key: &str, id: &str, fields: &[(&str, &str)]) -> Result<String, Error> {
self.text(with_pairs(args!["XADD", key, id], fields))
}
pub fn xadd_maxlen(
&mut self,
key: &str,
max_len: u64,
id: &str,
fields: &[(&str, &str)],
) -> Result<String, Error> {
self.text(with_pairs(
args!["XADD", key, "MAXLEN", max_len, id],
fields,
))
}
pub fn xlen(&mut self, key: &str) -> Result<i64, Error> {
self.integer(["XLEN", key])
}
pub fn xinfo_stream(&mut self, key: &str) -> Result<RespValue, Error> {
self.run(["XINFO", "STREAM", key])
}
pub fn xinfo_groups(&mut self, key: &str) -> Result<RespValue, Error> {
self.run(["XINFO", "GROUPS", key])
}
pub fn xinfo_consumers(&mut self, key: &str, group: &str) -> Result<RespValue, Error> {
self.run(["XINFO", "CONSUMERS", key, group])
}
pub fn xsetid(&mut self, key: &str, id: &str) -> Result<(), Error> {
self.ok(["XSETID", key, id])
}
pub fn xautoclaim(
&mut self,
key: &str,
group: &str,
consumer: &str,
min_idle_ms: u64,
start: &str,
) -> Result<RespValue, Error> {
self.run([
"XAUTOCLAIM",
key,
group,
consumer,
&min_idle_ms.to_string(),
start,
])
}
pub fn xrange(
&mut self,
key: &str,
start: &str,
end: &str,
count: Option<u64>,
) -> Result<Vec<StreamEntry>, Error> {
let mut arguments = args!["XRANGE", key, start, end];
if let Some(count) = count {
arguments.extend(args!["COUNT", count]);
}
stream_entries(&self.run(arguments)?)
}
pub fn xrevrange(
&mut self,
key: &str,
end: &str,
start: &str,
count: Option<u64>,
) -> Result<Vec<StreamEntry>, Error> {
let mut arguments = args!["XREVRANGE", key, end, start];
if let Some(count) = count {
arguments.extend(args!["COUNT", count]);
}
stream_entries(&self.run(arguments)?)
}
pub fn xdel(&mut self, key: &str, ids: &[&str]) -> Result<i64, Error> {
self.integer(with(args!["XDEL", key], ids))
}
pub fn xtrim_minid(&mut self, key: &str, id: &str) -> Result<i64, Error> {
self.integer(args!["XTRIM", key, "MINID", id])
}
pub fn xtrim_maxlen(&mut self, key: &str, max_len: u64) -> Result<i64, Error> {
self.integer(args!["XTRIM", key, "MAXLEN", max_len])
}
pub fn xread(
&mut self,
options: &StreamReadOptions,
streams: &[(&str, &str)],
) -> Result<Vec<StreamReadResult>, Error> {
if options.no_ack {
return Err(Error::Protocol(
"NOACK applies only to XREADGROUP".to_owned(),
));
}
let mut arguments = args!["XREAD"];
push_stream_read_options(&mut arguments, options);
push_streams(&mut arguments, streams);
stream_read(&self.run(arguments)?)
}
pub fn xgroup_create_consumer(
&mut self,
key: &str,
group: &str,
consumer: &str,
) -> Result<i64, Error> {
self.integer(args!["XGROUP", "CREATECONSUMER", key, group, consumer])
}
pub fn xgroup_create(
&mut self,
key: &str,
group: &str,
id: &str,
make_stream: bool,
) -> Result<(), Error> {
let mut arguments = args!["XGROUP", "CREATE", key, group, id];
if make_stream {
arguments.push("MKSTREAM".to_owned());
}
self.ok(arguments)
}
pub fn xreadgroup(
&mut self,
group: &str,
consumer: &str,
options: &StreamReadOptions,
streams: &[(&str, &str)],
) -> Result<Vec<StreamReadResult>, Error> {
let mut arguments = args!["XREADGROUP", "GROUP", group, consumer];
push_stream_read_options(&mut arguments, options);
if options.no_ack {
arguments.push("NOACK".to_owned());
}
push_streams(&mut arguments, streams);
stream_read(&self.run(arguments)?)
}
pub fn xgroup_destroy(&mut self, key: &str, group: &str) -> Result<bool, Error> {
self.flag(["XGROUP", "DESTROY", key, group])
}
pub fn xgroup_setid(&mut self, key: &str, group: &str, id: &str) -> Result<(), Error> {
self.ok(["XGROUP", "SETID", key, group, id])
}
pub fn xgroup_delconsumer(
&mut self,
key: &str,
group: &str,
consumer: &str,
) -> Result<i64, Error> {
self.integer(["XGROUP", "DELCONSUMER", key, group, consumer])
}
pub fn xack(&mut self, key: &str, group: &str, ids: &[&str]) -> Result<i64, Error> {
self.integer(with(args!["XACK", key, group], ids))
}
pub fn xpending_summary(&mut self, key: &str, group: &str) -> Result<RespValue, Error> {
self.run(["XPENDING", key, group])
}
pub fn xpending(
&mut self,
key: &str,
group: &str,
start: &str,
end: &str,
count: u64,
filter: &StreamPendingFilter,
) -> Result<Vec<StreamPendingEntry>, Error> {
let mut arguments = args!["XPENDING", key, group];
if let Some(idle) = filter.min_idle {
arguments.extend(args!["IDLE", duration_millis(idle)]);
}
arguments.extend(args![start, end, count]);
if let Some(consumer) = &filter.consumer {
arguments.push(consumer.clone());
}
let reply = self.run(arguments)?;
let mut rows = Vec::new();
for row in reply.as_array().unwrap_or_default() {
let Some([id, owner, idle, deliveries]) = row.as_array() else {
return Err(unexpected("XPENDING", row));
};
rows.push(StreamPendingEntry {
id: id.as_string()?.unwrap_or_default(),
consumer: owner.as_string()?.unwrap_or_default(),
idle_millis: idle.as_integer().unwrap_or_default(),
delivery_count: deliveries.as_integer().unwrap_or_default(),
});
}
Ok(rows)
}
pub fn xclaim(
&mut self,
key: &str,
group: &str,
consumer: &str,
min_idle: Duration,
ids: &[&str],
options: &StreamClaimOptions,
) -> Result<Vec<StreamEntry>, Error> {
let arguments = claim_arguments(key, group, consumer, min_idle, ids, options)?;
stream_entries(&self.run(arguments)?)
}
pub fn xclaim_ids(
&mut self,
key: &str,
group: &str,
consumer: &str,
min_idle: Duration,
ids: &[&str],
options: &StreamClaimOptions,
) -> Result<Vec<String>, Error> {
let mut arguments = claim_arguments(key, group, consumer, min_idle, ids, options)?;
arguments.push("JUSTID".to_owned());
self.strings(arguments)
}
pub fn info(&mut self, section: Option<&str>) -> Result<String, Error> {
match section {
Some(section) => self.text(["INFO", section]),
None => self.text(["INFO"]),
}
}
pub fn config_get(&mut self, parameter: &str) -> Result<Vec<(String, String)>, Error> {
pairs(&self.run(["CONFIG", "GET", parameter])?)
}
pub fn save(&mut self) -> Result<(), Error> {
self.ok(["SAVE"])
}
pub fn bgsave(&mut self) -> Result<(), Error> {
self.ok(["BGSAVE"])
}
pub fn flushdb(&mut self) -> Result<(), Error> {
self.ok(["FLUSHDB"])
}
pub fn flushall(&mut self) -> Result<(), Error> {
self.ok(["FLUSHALL"])
}
pub fn swapdb(&mut self, first: u32, second: u32) -> Result<(), Error> {
self.ok(args!["SWAPDB", first, second])
}
pub fn move_key(&mut self, key: &str, database: u32) -> Result<bool, Error> {
self.flag(args!["MOVE", key, database])
}
pub fn multi(&mut self) -> Result<(), Error> {
self.ok(["MULTI"])
}
pub fn exec(&mut self) -> Result<Option<Vec<RespValue>>, Error> {
match self.run(["EXEC"])? {
RespValue::Null => Ok(None),
RespValue::Array(items) => Ok(Some(items)),
other => Err(unexpected("EXEC", &other)),
}
}
pub fn discard(&mut self) -> Result<(), Error> {
self.ok(["DISCARD"])
}
pub fn watch_key(&mut self, key: &str) -> Result<(), Error> {
self.watch(&[key])
}
pub fn watch(&mut self, keys: &[&str]) -> Result<(), Error> {
self.ok(join("WATCH", keys))
}
pub fn unwatch(&mut self) -> Result<(), Error> {
self.ok(["UNWATCH"])
}
pub fn eval(
&mut self,
script: &str,
keys: &[&str],
arguments: &[&str],
) -> Result<RespValue, Error> {
self.run(script_arguments("EVAL", script, keys, arguments))
}
pub fn evalsha(
&mut self,
sha: &str,
keys: &[&str],
arguments: &[&str],
) -> Result<RespValue, Error> {
self.run(script_arguments("EVALSHA", sha, keys, arguments))
}
pub fn script_load(&mut self, script: &str) -> Result<String, Error> {
self.text(["SCRIPT", "LOAD", script])
}
pub fn script_exists(&mut self, hashes: &[&str]) -> Result<Vec<bool>, Error> {
let reply = self.run(with(args!["SCRIPT", "EXISTS"], hashes))?;
Ok(reply
.as_array()
.unwrap_or_default()
.iter()
.map(|item| item.as_integer().unwrap_or_default() > 0)
.collect())
}
pub fn script_flush(&mut self) -> Result<(), Error> {
self.ok(["SCRIPT", "FLUSH"])
}
pub fn script_kill(&mut self) -> Result<(), Error> {
self.ok(["SCRIPT", "KILL"])
}
pub fn fcall(
&mut self,
name: &str,
keys: &[&str],
arguments: &[&str],
) -> Result<RespValue, Error> {
self.run(script_arguments("FCALL", name, keys, arguments))
}
pub fn function_load(&mut self, name: &str, body: &str) -> Result<(), Error> {
self.ok(["FUNCTION", "LOAD", name, body])
}
pub fn function_list(&mut self) -> Result<RespValue, Error> {
self.run(["FUNCTION", "LIST"])
}
pub fn function_delete(&mut self, name: &str) -> Result<bool, Error> {
self.flag(["FUNCTION", "DELETE", name])
}
pub fn acl_whoami(&mut self) -> Result<String, Error> {
self.text(["ACL", "WHOAMI"])
}
pub fn acl_users(&mut self) -> Result<Vec<String>, Error> {
self.strings(["ACL", "USERS"])
}
pub fn acl_list(&mut self) -> Result<Vec<String>, Error> {
self.strings(["ACL", "LIST"])
}
pub fn acl_getuser(&mut self, user: &str) -> Result<RespValue, Error> {
self.run(["ACL", "GETUSER", user])
}
pub fn acl_cat(&mut self) -> Result<Vec<String>, Error> {
self.strings(["ACL", "CAT"])
}
pub fn acl_setuser(&mut self, user: &str, rules: &[&str]) -> Result<(), Error> {
self.ok(with(args!["ACL", "SETUSER", user], rules))
}
pub fn acl_load(&mut self) -> Result<(), Error> {
self.ok(["ACL", "LOAD"])
}
pub fn acl_save(&mut self) -> Result<(), Error> {
self.ok(["ACL", "SAVE"])
}
pub fn cluster_slots(&mut self) -> Result<RespValue, Error> {
self.run(["CLUSTER", "SLOTS"])
}
pub fn cluster_nodes(&mut self) -> Result<String, Error> {
self.text(["CLUSTER", "NODES"])
}
pub fn cluster_shards(&mut self) -> Result<RespValue, Error> {
self.run(["CLUSTER", "SHARDS"])
}
pub fn cluster_info(&mut self) -> Result<String, Error> {
self.text(["CLUSTER", "INFO"])
}
pub fn cluster_myid(&mut self) -> Result<String, Error> {
self.text(["CLUSTER", "MYID"])
}
pub fn cluster_keyslot(&mut self, key: &str) -> Result<i64, Error> {
self.integer(["CLUSTER", "KEYSLOT", key])
}
pub fn readonly(&mut self) -> Result<(), Error> {
self.ok(["READONLY"])
}
pub fn readwrite(&mut self) -> Result<(), Error> {
self.ok(["READWRITE"])
}
pub fn client_getname(&mut self) -> Result<Option<String>, Error> {
self.bulk_or_null(["CLIENT", "GETNAME"])
}
pub fn client_setname(&mut self, name: &str) -> Result<(), Error> {
self.ok(["CLIENT", "SETNAME", name])
}
pub fn client_tracking(&mut self, enabled: bool) -> Result<(), Error> {
self.ok(["CLIENT", "TRACKING", if enabled { "ON" } else { "OFF" }])
}
pub fn hello(&mut self, protocol: u8) -> Result<RespValue, Error> {
self.run(args!["HELLO", protocol])
}
pub fn publish(&mut self, channel: &str, message: &str) -> Result<i64, Error> {
self.integer(["PUBLISH", channel, message])
}
pub fn spublish(&mut self, channel: &str, message: &str) -> Result<i64, Error> {
self.integer(["SPUBLISH", channel, message])
}
pub fn subscribe(&mut self, channels: &[&str]) -> Result<Vec<RespValue>, Error> {
self.subscription("SUBSCRIBE", channels, false)
}
pub fn unsubscribe(&mut self, channels: &[&str]) -> Result<Vec<RespValue>, Error> {
self.subscription("UNSUBSCRIBE", channels, true)
}
pub fn psubscribe(&mut self, patterns: &[&str]) -> Result<Vec<RespValue>, Error> {
self.subscription("PSUBSCRIBE", patterns, false)
}
pub fn punsubscribe(&mut self, patterns: &[&str]) -> Result<Vec<RespValue>, Error> {
self.subscription("PUNSUBSCRIBE", patterns, true)
}
pub fn ssubscribe(&mut self, channels: &[&str]) -> Result<Vec<RespValue>, Error> {
self.subscription("SSUBSCRIBE", channels, false)
}
pub fn sunsubscribe(&mut self, channels: &[&str]) -> Result<Vec<RespValue>, Error> {
self.subscription("SUNSUBSCRIBE", channels, true)
}
pub fn next_message(&mut self) -> Result<PubSubMessage, Error> {
loop {
let reply = self.read_message()?;
if let Some(message) = pubsub_message(&reply)? {
return Ok(message);
}
}
}
pub fn listen<F>(&mut self, mut on_message: F) -> Result<(), Error>
where
F: FnMut(PubSubMessage),
{
loop {
on_message(self.next_message()?);
}
}
fn subscription(
&mut self,
command: &str,
names: &[&str],
allow_empty: bool,
) -> Result<Vec<RespValue>, Error> {
if names.is_empty() && !allow_empty {
return Err(Error::Protocol(format!(
"{command} needs at least one channel or pattern"
)));
}
self.run_replies(join(command, names), names.len().max(1))
}
fn optional_integer<I, S>(&mut self, arguments: I) -> Result<Option<i64>, Error>
where
I: IntoIterator<Item = S>,
S: AsRef<str>,
{
let (command, reply) = self.run_named(arguments)?;
match reply {
RespValue::Null => Ok(None),
RespValue::Integer(value) => Ok(Some(value)),
other => Err(unexpected(&command, &other)),
}
}
fn integers(&mut self, arguments: Vec<String>) -> Result<Vec<i64>, Error> {
let command = arguments.first().cloned().unwrap_or_default();
let reply = self.run(arguments)?;
integer_array(&command, &reply)
}
fn blocking_sorted_pop(
&mut self,
command: &str,
keys: &[&str],
timeout: Duration,
) -> Result<Option<(String, SortedSetEntry)>, Error> {
let mut arguments = join(command, keys);
arguments.push(duration_secs(timeout).to_string());
match self.run(arguments)? {
RespValue::Null => Ok(None),
RespValue::Array(items) if items.len() == 3 => Ok(Some((
items[0].as_string()?.unwrap_or_default(),
SortedSetEntry::new(
items[1].as_string()?.unwrap_or_default(),
parse_score(&items[2].as_string()?.unwrap_or_default())?,
),
))),
other => Err(unexpected(command, &other)),
}
}
fn blocking_pop(
&mut self,
command: &str,
keys: &[&str],
timeout: Duration,
) -> Result<Option<(String, String)>, Error> {
let mut arguments = join(command, keys);
arguments.push(duration_secs(timeout).to_string());
match self.run(arguments)? {
RespValue::Null => Ok(None),
RespValue::Array(items) if items.len() == 2 => Ok(Some((
items[0].as_string()?.unwrap_or_default(),
items[1].as_string()?.unwrap_or_default(),
))),
other => Err(unexpected(command, &other)),
}
}
}
fn field_block(
mut arguments: Vec<String>,
condition: FieldExpireCondition,
fields: &[&str],
) -> Vec<String> {
let flag = match condition {
FieldExpireCondition::Always => None,
FieldExpireCondition::IfNoTtl => Some("NX"),
FieldExpireCondition::IfTtl => Some("XX"),
FieldExpireCondition::IfGreater => Some("GT"),
FieldExpireCondition::IfLess => Some("LT"),
};
arguments.extend(flag.map(str::to_owned));
arguments.push("FIELDS".to_owned());
arguments.push(fields.len().to_string());
with(arguments, fields)
}
fn integer_array(command: &str, reply: &RespValue) -> Result<Vec<i64>, Error> {
let Some(items) = reply.as_array() else {
return Err(unexpected(command, reply));
};
items
.iter()
.map(|item| item.as_integer().ok_or_else(|| unexpected(command, item)))
.collect()
}
fn with(mut arguments: Vec<String>, rest: &[&str]) -> Vec<String> {
arguments.extend(rest.iter().map(|argument| (*argument).to_owned()));
arguments
}
fn with_pairs(mut arguments: Vec<String>, pairs: &[(&str, &str)]) -> Vec<String> {
for (left, right) in pairs {
arguments.push((*left).to_owned());
arguments.push((*right).to_owned());
}
arguments
}
fn push_scan_options(arguments: &mut Vec<String>, pattern: Option<&str>, count: Option<u64>) {
if let Some(pattern) = pattern {
arguments.extend(args!["MATCH", pattern]);
}
if let Some(count) = count {
arguments.extend(args!["COUNT", count]);
}
}
fn push_stream_read_options(arguments: &mut Vec<String>, options: &StreamReadOptions) {
if let Some(count) = options.count {
arguments.extend(args!["COUNT", count]);
}
if let Some(block) = options.block {
arguments.extend(args!["BLOCK", duration_millis(block)]);
}
}
fn push_streams(arguments: &mut Vec<String>, streams: &[(&str, &str)]) {
arguments.push("STREAMS".to_owned());
arguments.extend(streams.iter().map(|(key, _)| (*key).to_owned()));
arguments.extend(streams.iter().map(|(_, id)| (*id).to_owned()));
}
fn zadd_arguments(key: &str, options: SortedSetAddOptions) -> Vec<String> {
let mut arguments = args!["ZADD", key];
let flags = [
(options.if_not_exists, "NX"),
(options.if_exists, "XX"),
(options.greater_than, "GT"),
(options.less_than, "LT"),
(options.changed, "CH"),
];
for (enabled, flag) in flags {
if enabled {
arguments.push(flag.to_owned());
}
}
arguments
}
fn claim_arguments(
key: &str,
group: &str,
consumer: &str,
min_idle: Duration,
ids: &[&str],
options: &StreamClaimOptions,
) -> Result<Vec<String>, Error> {
if options.idle.is_some() && options.time.is_some() {
return Err(Error::Protocol(
"IDLE and TIME cannot be combined".to_owned(),
));
}
let mut arguments = with(
args!["XCLAIM", key, group, consumer, duration_millis(min_idle)],
ids,
);
if let Some(idle) = options.idle {
arguments.extend(args!["IDLE", duration_millis(idle)]);
}
if let Some(time) = options.time {
arguments.extend(args!["TIME", unix_millis(time)?]);
}
if let Some(retry_count) = options.retry_count {
arguments.extend(args!["RETRYCOUNT", retry_count]);
}
if options.force {
arguments.push("FORCE".to_owned());
}
if let Some(last_id) = &options.last_id {
arguments.extend(args!["LASTID", last_id]);
}
Ok(arguments)
}
fn script_arguments(command: &str, script: &str, keys: &[&str], arguments: &[&str]) -> Vec<String> {
let values = with(args![command, script, keys.len()], keys);
with(values, arguments)
}
fn parse_scan(command: &str, reply: RespValue) -> Result<(u64, Vec<String>), Error> {
let Some([cursor, items]) = reply.as_array() else {
return Err(unexpected(command, &reply));
};
let cursor = cursor
.as_string()?
.unwrap_or_default()
.parse()
.map_err(|_| Error::Protocol(format!("{command} returned a bad cursor")))?;
let mut keys = Vec::new();
for item in items.as_array().unwrap_or_default() {
keys.push(item.as_string()?.unwrap_or_default());
}
Ok((cursor, keys))
}
fn parse_score(text: &str) -> Result<f64, Error> {
text.parse()
.map_err(|_| Error::Protocol(format!("{text:?} is not a score")))
}
fn pairs(reply: &RespValue) -> Result<Vec<(String, String)>, Error> {
let items = reply.as_array().unwrap_or_default();
let mut values = Vec::with_capacity(items.len() / 2);
for pair in items.chunks_exact(2) {
values.push((
pair[0].as_string()?.unwrap_or_default(),
pair[1].as_string()?.unwrap_or_default(),
));
}
Ok(values)
}
fn sorted_set(reply: &RespValue) -> Result<Vec<SortedSetEntry>, Error> {
pairs(reply)?
.into_iter()
.map(|(member, score)| Ok(SortedSetEntry::new(member, parse_score(&score)?)))
.collect()
}
fn stream_entries(reply: &RespValue) -> Result<Vec<StreamEntry>, Error> {
let mut entries = Vec::new();
for item in reply.as_array().unwrap_or_default() {
let Some([id, fields]) = item.as_array() else {
continue;
};
entries.push(StreamEntry {
id: id.as_string()?.unwrap_or_default(),
fields: pairs(fields)?,
});
}
Ok(entries)
}
fn stream_read(reply: &RespValue) -> Result<Vec<StreamReadResult>, Error> {
let mut results = Vec::new();
for item in reply.as_array().unwrap_or_default() {
let Some([key, entries]) = item.as_array() else {
continue;
};
results.push(StreamReadResult {
key: key.as_string()?.unwrap_or_default(),
entries: stream_entries(entries)?,
});
}
Ok(results)
}
fn pubsub_message(reply: &RespValue) -> Result<Option<PubSubMessage>, Error> {
let items = reply.as_array().unwrap_or_default();
let kind = match items.first() {
Some(first) => first.as_string()?.unwrap_or_default(),
None => return Ok(None),
};
let (pattern, channel, payload) = match (kind.as_str(), items) {
("message" | "smessage", [_, channel, payload]) => (None, channel, payload),
("pmessage", [_, pattern, channel, payload]) => (pattern.as_string()?, channel, payload),
_ => return Ok(None),
};
let payload = match payload {
RespValue::Bulk(bytes) => bytes.clone(),
other => other.as_string()?.unwrap_or_default().into_bytes(),
};
Ok(Some(PubSubMessage {
kind,
pattern,
channel: channel.as_string()?.unwrap_or_default(),
payload,
}))
}
#[cfg(test)]
mod tests {
use super::{field_block, integer_array, pubsub_message};
use crate::models::FieldExpireCondition;
use crate::value::RespValue;
#[test]
fn field_block_puts_the_condition_before_fields() {
assert_eq!(
field_block(
vec!["HEXPIRE".into(), "h".into(), "10".into()],
FieldExpireCondition::IfGreater,
&["a", "b"],
),
["HEXPIRE", "h", "10", "GT", "FIELDS", "2", "a", "b"]
);
assert_eq!(
field_block(
vec!["HTTL".into(), "h".into()],
FieldExpireCondition::Always,
&["a"]
),
["HTTL", "h", "FIELDS", "1", "a"]
);
assert_eq!(
integer_array(
"HTTL",
&RespValue::Array(vec![RespValue::Integer(-2), RespValue::Integer(5)])
)
.unwrap(),
vec![-2, 5]
);
}
#[test]
fn pubsub_message_reads_payload_from_an_array() {
let reply = RespValue::Array(vec![
RespValue::Bulk(b"message".to_vec()),
RespValue::Bulk(b"news".to_vec()),
RespValue::Bulk(b"hello-ruvio".to_vec()),
]);
let parsed = pubsub_message(&reply).unwrap().unwrap();
assert_eq!(parsed.kind, "message");
assert_eq!(parsed.channel, "news");
assert_eq!(parsed.payload, b"hello-ruvio");
assert!(
pubsub_message(&RespValue::Array(vec![
RespValue::Bulk(b"subscribe".to_vec()),
RespValue::Bulk(b"news".to_vec()),
RespValue::Integer(1),
]))
.unwrap()
.is_none()
);
}
}