#![allow(clippy::too_many_arguments, clippy::type_complexity)]
use crate::error::Result;
use crate::executor::CommandExecutor;
use crate::value;
use async_trait::async_trait;
use bytes::Bytes;
use redis::{Cmd, ToRedisArgs};
pub type StreamEntry = (String, Vec<(Bytes, Bytes)>);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum StreamTrimStrategy {
MaxLen,
MinId,
}
#[derive(Debug, Clone)]
pub struct StreamTrimOptions {
strategy: StreamTrimStrategy,
exact: bool,
threshold: String,
limit: Option<i64>,
}
impl StreamTrimOptions {
pub fn max_len(exact: bool, threshold: i64, limit: Option<i64>) -> Self {
Self {
strategy: StreamTrimStrategy::MaxLen,
exact,
threshold: threshold.to_string(),
limit,
}
}
pub fn min_id(exact: bool, threshold: impl Into<String>, limit: Option<i64>) -> Self {
Self {
strategy: StreamTrimStrategy::MinId,
exact,
threshold: threshold.into(),
limit,
}
}
pub(crate) fn add_to(&self, cmd: &mut Cmd) {
match self.strategy {
StreamTrimStrategy::MaxLen => cmd.arg("MAXLEN"),
StreamTrimStrategy::MinId => cmd.arg("MINID"),
};
cmd.arg(if self.exact { "=" } else { "~" });
cmd.arg(&self.threshold);
if let Some(l) = self.limit {
cmd.arg("LIMIT").arg(l);
}
}
}
#[derive(Debug, Clone, Default)]
pub struct StreamAddOptions {
pub make_stream: bool,
pub trim: Option<StreamTrimOptions>,
}
impl StreamAddOptions {
pub(crate) fn add_to(&self, cmd: &mut Cmd) {
if !self.make_stream {
cmd.arg("NOMKSTREAM");
}
if let Some(t) = &self.trim {
t.add_to(cmd);
}
}
}
#[derive(Debug, Clone, Copy, Default)]
pub struct StreamReadOptions {
pub block_ms: Option<i64>,
pub count: Option<i64>,
}
impl StreamReadOptions {
pub(crate) fn add_to(&self, cmd: &mut Cmd) {
if let Some(b) = self.block_ms {
cmd.arg("BLOCK").arg(b);
}
if let Some(c) = self.count {
cmd.arg("COUNT").arg(c);
}
}
}
#[derive(Debug, Clone, Copy, Default)]
pub struct StreamReadGroupOptions {
pub block_ms: Option<i64>,
pub count: Option<i64>,
pub no_ack: bool,
}
impl StreamReadGroupOptions {
pub(crate) fn add_to(&self, cmd: &mut Cmd) {
if let Some(b) = self.block_ms {
cmd.arg("BLOCK").arg(b);
}
if let Some(c) = self.count {
cmd.arg("COUNT").arg(c);
}
if self.no_ack {
cmd.arg("NOACK");
}
}
}
#[derive(Debug, Clone, Copy, Default)]
pub struct StreamGroupCreateOptions {
pub make_stream: bool,
pub entries_read: Option<i64>,
}
impl StreamGroupCreateOptions {
pub(crate) fn add_to(&self, cmd: &mut Cmd) {
if self.make_stream {
cmd.arg("MKSTREAM");
}
if let Some(e) = self.entries_read {
cmd.arg("ENTRIESREAD").arg(e);
}
}
}
#[derive(Debug, Clone, Copy, Default)]
pub struct StreamClaimOptions {
pub idle: Option<i64>,
pub idle_unix_time: Option<i64>,
pub retry_count: Option<i64>,
pub is_force: bool,
}
impl StreamClaimOptions {
pub(crate) fn add_to(&self, cmd: &mut Cmd) {
if let Some(i) = self.idle {
cmd.arg("IDLE").arg(i);
}
if let Some(t) = self.idle_unix_time {
cmd.arg("TIME").arg(t);
}
if let Some(r) = self.retry_count {
cmd.arg("RETRYCOUNT").arg(r);
}
if self.is_force {
cmd.arg("FORCE");
}
}
}
pub type PendingConsumer = (Bytes, i64);
#[derive(Debug, Clone, Default)]
pub struct XPendingSummary {
pub count: i64,
pub min_id: Option<Bytes>,
pub max_id: Option<Bytes>,
pub consumers: Vec<PendingConsumer>,
}
#[derive(Debug, Clone)]
pub struct XPendingEntry {
pub id: Bytes,
pub consumer: Bytes,
pub idle_ms: i64,
pub delivery_count: i64,
}
#[async_trait]
pub trait StreamCommands: CommandExecutor {
async fn xadd<K, F, V>(&self, key: K, id: &str, fields: &[(F, V)]) -> Result<Option<String>>
where
K: ToRedisArgs + Send + Sync,
F: ToRedisArgs + Send + Sync,
V: ToRedisArgs + Send + Sync,
{
let mut cmd = Cmd::new();
cmd.arg("XADD").arg(key).arg(id);
for (f, v) in fields {
cmd.arg(f).arg(v);
}
value::to_opt_string(self.execute_command(cmd, None).await?)
}
async fn xlen<K: ToRedisArgs + Send>(&self, key: K) -> Result<i64> {
let mut cmd = Cmd::new();
cmd.arg("XLEN").arg(key);
value::to_i64(self.execute_command(cmd, None).await?)
}
async fn xdel<K: ToRedisArgs + Send>(&self, key: K, ids: &[&str]) -> Result<i64> {
let mut cmd = Cmd::new();
cmd.arg("XDEL").arg(key);
for id in ids {
cmd.arg(*id);
}
value::to_i64(self.execute_command(cmd, None).await?)
}
async fn xtrim_maxlen<K: ToRedisArgs + Send>(
&self,
key: K,
maxlen: i64,
approximate: bool,
) -> Result<i64> {
let mut cmd = Cmd::new();
cmd.arg("XTRIM").arg(key).arg("MAXLEN");
if approximate {
cmd.arg("~");
}
cmd.arg(maxlen);
value::to_i64(self.execute_command(cmd, None).await?)
}
async fn xrange<K: ToRedisArgs + Send>(
&self,
key: K,
start: &str,
end: &str,
) -> Result<Vec<StreamEntry>> {
let mut cmd = Cmd::new();
cmd.arg("XRANGE").arg(key).arg(start).arg(end);
parse_entries(self.execute_command(cmd, None).await?)
}
async fn xrevrange<K: ToRedisArgs + Send>(
&self,
key: K,
end: &str,
start: &str,
) -> Result<Vec<StreamEntry>> {
let mut cmd = Cmd::new();
cmd.arg("XREVRANGE").arg(key).arg(end).arg(start);
parse_entries(self.execute_command(cmd, None).await?)
}
async fn xgroup_create<K: ToRedisArgs + Send>(
&self,
key: K,
group: &str,
id: &str,
mkstream: bool,
) -> Result<()> {
let mut cmd = Cmd::new();
cmd.arg("XGROUP").arg("CREATE").arg(key).arg(group).arg(id);
if mkstream {
cmd.arg("MKSTREAM");
}
value::to_unit(self.execute_command(cmd, None).await?)
}
async fn xgroup_destroy<K: ToRedisArgs + Send>(&self, key: K, group: &str) -> Result<bool> {
let mut cmd = Cmd::new();
cmd.arg("XGROUP").arg("DESTROY").arg(key).arg(group);
value::to_bool(self.execute_command(cmd, None).await?)
}
async fn xack<K: ToRedisArgs + Send>(&self, key: K, group: &str, ids: &[&str]) -> Result<i64> {
let mut cmd = Cmd::new();
cmd.arg("XACK").arg(key).arg(group);
for id in ids {
cmd.arg(*id);
}
value::to_i64(self.execute_command(cmd, None).await?)
}
async fn xadd_options<K, F, V>(
&self,
key: K,
id: &str,
fields: &[(F, V)],
options: &StreamAddOptions,
) -> Result<Option<String>>
where
K: ToRedisArgs + Send + Sync,
F: ToRedisArgs + Send + Sync,
V: ToRedisArgs + Send + Sync,
{
let mut cmd = Cmd::new();
cmd.arg("XADD").arg(key);
options.add_to(&mut cmd);
cmd.arg(id);
for (f, v) in fields {
cmd.arg(f).arg(v);
}
value::to_opt_string(self.execute_command(cmd, None).await?)
}
async fn xtrim_minid<K: ToRedisArgs + Send>(
&self,
key: K,
minid: &str,
approximate: bool,
) -> Result<i64> {
let mut cmd = Cmd::new();
cmd.arg("XTRIM").arg(key).arg("MINID");
if approximate {
cmd.arg("~");
}
cmd.arg(minid);
value::to_i64(self.execute_command(cmd, None).await?)
}
async fn xread<K: ToRedisArgs + Send + Sync>(
&self,
keys_ids: &[(K, &str)],
options: Option<StreamReadOptions>,
) -> Result<Vec<(Bytes, Vec<StreamEntry>)>> {
let mut cmd = Cmd::new();
cmd.arg("XREAD");
if let Some(o) = options {
o.add_to(&mut cmd);
}
cmd.arg("STREAMS");
for (k, _) in keys_ids {
cmd.arg(k);
}
for (_, id) in keys_ids {
cmd.arg(*id);
}
parse_stream_read(self.execute_command(cmd, None).await?)
}
async fn xreadgroup<K: ToRedisArgs + Send + Sync>(
&self,
group: &str,
consumer: &str,
keys_ids: &[(K, &str)],
options: Option<StreamReadGroupOptions>,
) -> Result<Vec<(Bytes, Vec<StreamEntry>)>> {
let mut cmd = Cmd::new();
cmd.arg("XREADGROUP").arg("GROUP").arg(group).arg(consumer);
if let Some(o) = options {
o.add_to(&mut cmd);
}
cmd.arg("STREAMS");
for (k, _) in keys_ids {
cmd.arg(k);
}
for (_, id) in keys_ids {
cmd.arg(*id);
}
parse_stream_read(self.execute_command(cmd, None).await?)
}
async fn xclaim<K: ToRedisArgs + Send>(
&self,
key: K,
group: &str,
consumer: &str,
min_idle_time_ms: i64,
ids: &[&str],
options: Option<StreamClaimOptions>,
) -> Result<Vec<StreamEntry>> {
let mut cmd = Cmd::new();
cmd.arg("XCLAIM")
.arg(key)
.arg(group)
.arg(consumer)
.arg(min_idle_time_ms);
for id in ids {
cmd.arg(*id);
}
if let Some(o) = options {
o.add_to(&mut cmd);
}
parse_entries(self.execute_command(cmd, None).await?)
}
async fn xclaim_justid<K: ToRedisArgs + Send>(
&self,
key: K,
group: &str,
consumer: &str,
min_idle_time_ms: i64,
ids: &[&str],
options: Option<StreamClaimOptions>,
) -> Result<Vec<String>> {
let mut cmd = Cmd::new();
cmd.arg("XCLAIM")
.arg(key)
.arg(group)
.arg(consumer)
.arg(min_idle_time_ms);
for id in ids {
cmd.arg(*id);
}
if let Some(o) = options {
o.add_to(&mut cmd);
}
cmd.arg("JUSTID");
collect_strings(self.execute_command(cmd, None).await?)
}
async fn xautoclaim<K: ToRedisArgs + Send>(
&self,
key: K,
group: &str,
consumer: &str,
min_idle_time_ms: i64,
start: &str,
count: Option<i64>,
) -> Result<(String, Vec<StreamEntry>, Vec<String>)> {
let mut cmd = Cmd::new();
cmd.arg("XAUTOCLAIM")
.arg(key)
.arg(group)
.arg(consumer)
.arg(min_idle_time_ms)
.arg(start);
if let Some(c) = count {
cmd.arg("COUNT").arg(c);
}
parse_autoclaim(self.execute_command(cmd, None).await?)
}
async fn xautoclaim_justid<K: ToRedisArgs + Send>(
&self,
key: K,
group: &str,
consumer: &str,
min_idle_time_ms: i64,
start: &str,
count: Option<i64>,
) -> Result<(String, Vec<String>, Vec<String>)> {
let mut cmd = Cmd::new();
cmd.arg("XAUTOCLAIM")
.arg(key)
.arg(group)
.arg(consumer)
.arg(min_idle_time_ms)
.arg(start);
if let Some(c) = count {
cmd.arg("COUNT").arg(c);
}
cmd.arg("JUSTID");
parse_autoclaim_justid(self.execute_command(cmd, None).await?)
}
async fn xpending<K: ToRedisArgs + Send>(
&self,
key: K,
group: &str,
) -> Result<XPendingSummary> {
let mut cmd = Cmd::new();
cmd.arg("XPENDING").arg(key).arg(group);
parse_xpending_summary(self.execute_command(cmd, None).await?)
}
async fn xpending_range<K: ToRedisArgs + Send>(
&self,
key: K,
group: &str,
start: &str,
end: &str,
count: i64,
min_idle_time_ms: Option<i64>,
consumer: Option<&str>,
) -> Result<Vec<XPendingEntry>> {
let mut cmd = Cmd::new();
cmd.arg("XPENDING").arg(key).arg(group);
if let Some(idle) = min_idle_time_ms {
cmd.arg("IDLE").arg(idle);
}
cmd.arg(start).arg(end).arg(count);
if let Some(c) = consumer {
cmd.arg(c);
}
parse_xpending_range(self.execute_command(cmd, None).await?)
}
async fn xinfo_stream<K: ToRedisArgs + Send>(
&self,
key: K,
) -> Result<Vec<(Bytes, redis::Value)>> {
let mut cmd = Cmd::new();
cmd.arg("XINFO").arg("STREAM").arg(key);
parse_field_value_map(self.execute_command(cmd, None).await?)
}
async fn xinfo_stream_full<K: ToRedisArgs + Send>(
&self,
key: K,
count: Option<i64>,
) -> Result<Vec<(Bytes, redis::Value)>> {
let mut cmd = Cmd::new();
cmd.arg("XINFO").arg("STREAM").arg(key).arg("FULL");
if let Some(c) = count {
cmd.arg("COUNT").arg(c);
}
parse_field_value_map(self.execute_command(cmd, None).await?)
}
async fn xinfo_groups<K: ToRedisArgs + Send>(
&self,
key: K,
) -> Result<Vec<Vec<(Bytes, redis::Value)>>> {
let mut cmd = Cmd::new();
cmd.arg("XINFO").arg("GROUPS").arg(key);
parse_list_of_maps(self.execute_command(cmd, None).await?)
}
async fn xinfo_consumers<K: ToRedisArgs + Send>(
&self,
key: K,
group: &str,
) -> Result<Vec<Vec<(Bytes, redis::Value)>>> {
let mut cmd = Cmd::new();
cmd.arg("XINFO").arg("CONSUMERS").arg(key).arg(group);
parse_list_of_maps(self.execute_command(cmd, None).await?)
}
async fn xsetid<K: ToRedisArgs + Send>(
&self,
key: K,
last_id: &str,
entries_added: Option<i64>,
max_deleted_id: Option<&str>,
) -> Result<()> {
let mut cmd = Cmd::new();
cmd.arg("XSETID").arg(key).arg(last_id);
if let Some(e) = entries_added {
cmd.arg("ENTRIESADDED").arg(e);
}
if let Some(m) = max_deleted_id {
cmd.arg("MAXDELETEDID").arg(m);
}
value::to_unit(self.execute_command(cmd, None).await?)
}
async fn xgroup_create_options<K: ToRedisArgs + Send>(
&self,
key: K,
group: &str,
id: &str,
options: &StreamGroupCreateOptions,
) -> Result<()> {
let mut cmd = Cmd::new();
cmd.arg("XGROUP").arg("CREATE").arg(key).arg(group).arg(id);
options.add_to(&mut cmd);
value::to_unit(self.execute_command(cmd, None).await?)
}
async fn xgroup_create_consumer<K: ToRedisArgs + Send>(
&self,
key: K,
group: &str,
consumer: &str,
) -> Result<bool> {
let mut cmd = Cmd::new();
cmd.arg("XGROUP")
.arg("CREATECONSUMER")
.arg(key)
.arg(group)
.arg(consumer);
value::to_bool(self.execute_command(cmd, None).await?)
}
async fn xgroup_del_consumer<K: ToRedisArgs + Send>(
&self,
key: K,
group: &str,
consumer: &str,
) -> Result<i64> {
let mut cmd = Cmd::new();
cmd.arg("XGROUP")
.arg("DELCONSUMER")
.arg(key)
.arg(group)
.arg(consumer);
value::to_i64(self.execute_command(cmd, None).await?)
}
async fn xgroup_set_id<K: ToRedisArgs + Send>(
&self,
key: K,
group: &str,
id: &str,
entries_read: Option<i64>,
) -> Result<()> {
let mut cmd = Cmd::new();
cmd.arg("XGROUP").arg("SETID").arg(key).arg(group).arg(id);
if let Some(e) = entries_read {
cmd.arg("ENTRIESREAD").arg(e);
}
value::to_unit(self.execute_command(cmd, None).await?)
}
}
fn parse_entries(v: redis::Value) -> Result<Vec<StreamEntry>> {
let pairs: Vec<(redis::Value, redis::Value)> = match v {
redis::Value::Nil => return Ok(Vec::new()),
redis::Value::Map(pairs) => pairs,
redis::Value::Array(items) => {
let mut out = Vec::with_capacity(items.len());
for entry in items {
if let redis::Value::Array(mut parts) = entry
&& parts.len() == 2
{
let fields = parts.pop().unwrap();
let id = parts.pop().unwrap();
out.push((id, fields));
}
}
out
}
other => {
return Err(crate::error::GlideError::Request(format!(
"unexpected stream reply: {other:?}"
)));
}
};
let mut out = Vec::with_capacity(pairs.len());
for (id_val, fields_val) in pairs {
let id = value::to_string(id_val)?;
let fv = parse_fields(fields_val)?;
out.push((id, fv));
}
Ok(out)
}
fn parse_fields(v: redis::Value) -> Result<Vec<(Bytes, Bytes)>> {
let items = match v {
redis::Value::Array(items) => items,
redis::Value::Nil => return Ok(Vec::new()),
other => return Ok(vec![(value::to_bytes(other)?, Bytes::new())]),
};
if items
.iter()
.all(|it| matches!(it, redis::Value::Array(inner) if inner.len() == 2))
{
let mut out = Vec::with_capacity(items.len());
for it in items {
if let redis::Value::Array(mut pair) = it {
let val = value::to_bytes(pair.pop().unwrap())?;
let field = value::to_bytes(pair.pop().unwrap())?;
out.push((field, val));
}
}
return Ok(out);
}
let mut out = Vec::with_capacity(items.len() / 2);
let mut iter = items.into_iter();
while let (Some(f), Some(val)) = (iter.next(), iter.next()) {
out.push((value::to_bytes(f)?, value::to_bytes(val)?));
}
Ok(out)
}
impl<T: CommandExecutor + ?Sized> StreamCommands for T {}
fn collect_strings(v: redis::Value) -> Result<Vec<String>> {
match v {
redis::Value::Nil => Ok(Vec::new()),
redis::Value::Array(items) => items.into_iter().map(value::to_string).collect(),
other => Ok(vec![value::to_string(other)?]),
}
}
fn parse_stream_read(v: redis::Value) -> Result<Vec<(Bytes, Vec<StreamEntry>)>> {
let pairs: Vec<(redis::Value, redis::Value)> = match v {
redis::Value::Nil => return Ok(Vec::new()),
redis::Value::Map(pairs) => pairs,
redis::Value::Array(items) => {
let mut out = Vec::with_capacity(items.len());
for entry in items {
if let redis::Value::Array(mut parts) = entry
&& parts.len() == 2
{
let entries = parts.pop().unwrap();
let key = parts.pop().unwrap();
out.push((key, entries));
}
}
out
}
other => {
return Err(crate::error::GlideError::Request(format!(
"unexpected XREAD reply: {other:?}"
)));
}
};
let mut out = Vec::with_capacity(pairs.len());
for (key_val, entries_val) in pairs {
let key = value::to_bytes(key_val)?;
let entries = parse_entries(entries_val)?;
out.push((key, entries));
}
Ok(out)
}
fn parse_autoclaim(v: redis::Value) -> Result<(String, Vec<StreamEntry>, Vec<String>)> {
match v {
redis::Value::Array(mut items) if items.len() == 2 || items.len() == 3 => {
let deleted = if items.len() == 3 {
collect_strings(items.pop().unwrap())?
} else {
Vec::new()
};
let entries = parse_entries(items.pop().unwrap())?;
let cursor = value::to_string(items.pop().unwrap())?;
Ok((cursor, entries, deleted))
}
other => Err(crate::error::GlideError::Request(format!(
"unexpected XAUTOCLAIM reply: {other:?}"
))),
}
}
fn parse_autoclaim_justid(v: redis::Value) -> Result<(String, Vec<String>, Vec<String>)> {
match v {
redis::Value::Array(mut items) if items.len() == 2 || items.len() == 3 => {
let deleted = if items.len() == 3 {
collect_strings(items.pop().unwrap())?
} else {
Vec::new()
};
let ids = collect_strings(items.pop().unwrap())?;
let cursor = value::to_string(items.pop().unwrap())?;
Ok((cursor, ids, deleted))
}
other => Err(crate::error::GlideError::Request(format!(
"unexpected XAUTOCLAIM JUSTID reply: {other:?}"
))),
}
}
fn parse_xpending_summary(v: redis::Value) -> Result<XPendingSummary> {
let mut items = match v {
redis::Value::Array(items) if items.len() == 4 => items,
redis::Value::Nil => return Ok(XPendingSummary::default()),
other => {
return Err(crate::error::GlideError::Request(format!(
"unexpected XPENDING summary reply: {other:?}"
)));
}
};
let consumers_val = items.pop().unwrap();
let max_val = items.pop().unwrap();
let min_val = items.pop().unwrap();
let count = value::to_i64(items.pop().unwrap())?;
let consumers = match consumers_val {
redis::Value::Nil => Vec::new(),
redis::Value::Array(list) => {
let mut out = Vec::with_capacity(list.len());
for it in list {
if let redis::Value::Array(mut pair) = it
&& pair.len() == 2
{
let cnt = value::to_i64(pair.pop().unwrap())?;
let name = value::to_bytes(pair.pop().unwrap())?;
out.push((name, cnt));
}
}
out
}
_ => Vec::new(),
};
Ok(XPendingSummary {
count,
min_id: value::to_opt_bytes(min_val)?,
max_id: value::to_opt_bytes(max_val)?,
consumers,
})
}
fn parse_xpending_range(v: redis::Value) -> Result<Vec<XPendingEntry>> {
let items = match v {
redis::Value::Nil => return Ok(Vec::new()),
redis::Value::Array(items) => items,
other => {
return Err(crate::error::GlideError::Request(format!(
"unexpected XPENDING range reply: {other:?}"
)));
}
};
let mut out = Vec::with_capacity(items.len());
for it in items {
if let redis::Value::Array(mut parts) = it
&& parts.len() == 4
{
let delivery_count = value::to_i64(parts.pop().unwrap())?;
let idle_ms = value::to_i64(parts.pop().unwrap())?;
let consumer = value::to_bytes(parts.pop().unwrap())?;
let id = value::to_bytes(parts.pop().unwrap())?;
out.push(XPendingEntry {
id,
consumer,
idle_ms,
delivery_count,
});
}
}
Ok(out)
}
fn parse_field_value_map(v: redis::Value) -> Result<Vec<(Bytes, redis::Value)>> {
match v {
redis::Value::Nil => Ok(Vec::new()),
redis::Value::Map(pairs) => pairs
.into_iter()
.map(|(k, val)| Ok((value::to_bytes(k)?, val)))
.collect(),
redis::Value::Array(items) => {
let mut out = Vec::with_capacity(items.len() / 2);
let mut iter = items.into_iter();
while let (Some(k), Some(val)) = (iter.next(), iter.next()) {
out.push((value::to_bytes(k)?, val));
}
Ok(out)
}
other => Err(crate::error::GlideError::Request(format!(
"unexpected XINFO reply: {other:?}"
))),
}
}
fn parse_list_of_maps(v: redis::Value) -> Result<Vec<Vec<(Bytes, redis::Value)>>> {
match v {
redis::Value::Nil => Ok(Vec::new()),
redis::Value::Array(items) => items.into_iter().map(parse_field_value_map).collect(),
other => Err(crate::error::GlideError::Request(format!(
"unexpected XINFO list reply: {other:?}"
))),
}
}
#[cfg(test)]
mod tests {
use super::*;
fn args_of(cmd: &Cmd) -> Vec<String> {
cmd.args_iter()
.filter_map(|a| match a {
redis::Arg::Simple(bytes) => Some(String::from_utf8_lossy(bytes).into_owned()),
redis::Arg::Cursor => None,
})
.collect()
}
#[test]
fn trim_options_maxlen_args() {
let mut cmd = Cmd::new();
StreamTrimOptions::max_len(true, 100, None).add_to(&mut cmd);
assert_eq!(args_of(&cmd), vec!["MAXLEN", "=", "100"]);
let mut cmd = Cmd::new();
StreamTrimOptions::max_len(false, 100, Some(10)).add_to(&mut cmd);
assert_eq!(args_of(&cmd), vec!["MAXLEN", "~", "100", "LIMIT", "10"]);
}
#[test]
fn trim_options_minid_args() {
let mut cmd = Cmd::new();
StreamTrimOptions::min_id(false, "1526985054069-0", None).add_to(&mut cmd);
assert_eq!(args_of(&cmd), vec!["MINID", "~", "1526985054069-0"]);
}
#[test]
fn add_options_args() {
let opts = StreamAddOptions {
make_stream: false,
trim: Some(StreamTrimOptions::max_len(true, 5, None)),
};
let mut cmd = Cmd::new();
opts.add_to(&mut cmd);
assert_eq!(args_of(&cmd), vec!["NOMKSTREAM", "MAXLEN", "=", "5"]);
}
#[test]
fn read_group_options_args() {
let opts = StreamReadGroupOptions {
block_ms: Some(500),
count: Some(10),
no_ack: true,
};
let mut cmd = Cmd::new();
opts.add_to(&mut cmd);
assert_eq!(args_of(&cmd), vec!["BLOCK", "500", "COUNT", "10", "NOACK"]);
}
#[test]
fn claim_options_args() {
let opts = StreamClaimOptions {
idle: Some(100),
idle_unix_time: None,
retry_count: Some(3),
is_force: true,
};
let mut cmd = Cmd::new();
opts.add_to(&mut cmd);
assert_eq!(
args_of(&cmd),
vec!["IDLE", "100", "RETRYCOUNT", "3", "FORCE"]
);
}
}