use crate::{
api::stream::{
StreamEntry,
conf::{
StreamAddOptions, StreamAutoClaimOptions, StreamClaimOptions, StreamLenOptions,
StreamPendingOptions, StreamRangeOptions, StreamTrimOptions,
},
meta::{
StreamAutoClaimResult, StreamClaimResult, StreamConsumerGroupMeta, StreamConsumerMeta,
StreamGetPendingEntryResult, StreamId, StreamInfo, StreamNack, StreamReadResult,
},
},
error::Result,
traits::DbLike,
};
pub trait Stream: DbLike {
fn xlast_id<K: AsRef<[u8]>>(&self, key: K) -> Result<StreamId>;
fn xadd<K: AsRef<[u8]>, FK: AsRef<[u8]>, FV: AsRef<[u8]>>(
&self,
key: K,
options: impl Into<StreamAddOptions>,
field_vals: &[(FK, FV)],
) -> Result<StreamId>;
fn xadd_simple<K: AsRef<[u8]>, FK: AsRef<[u8]>, FV: AsRef<[u8]>>(
&self,
key: K,
id_opt: Option<StreamId>,
field_vals: &[(FK, FV)],
) -> Result<StreamId>;
fn xlen<K: AsRef<[u8]>>(&self, key: K) -> Result<u64>;
fn xlen_with_options<K: AsRef<[u8]>>(&self, key: K, options: StreamLenOptions) -> Result<u64>;
fn xrange<K: AsRef<[u8]>>(
&self,
key: K,
start: StreamId,
end: StreamId,
count: Option<usize>,
) -> Result<Vec<StreamEntry>>;
fn xrevrange<K: AsRef<[u8]>>(
&self,
key: K,
end: StreamId,
start: StreamId,
count: Option<usize>,
) -> Result<Vec<StreamEntry>>;
fn xrange_with_options<K: AsRef<[u8]>>(
&self,
key: K,
options: StreamRangeOptions,
) -> Result<Vec<StreamEntry>>;
fn xtrim<K: AsRef<[u8]>>(&self, key: K, options: StreamTrimOptions) -> Result<u64>;
fn xdel<K: AsRef<[u8]>>(&self, key: K, ids: &[StreamId]) -> Result<u64>;
fn xgroup_create<K: AsRef<[u8]>>(
&self,
key: K,
group_name: &str,
last_id: &str,
mkstream: bool,
entries_read: Option<i64>,
) -> Result<()>;
fn xgroup_destroy<K: AsRef<[u8]>>(&self, key: K, group_name: &str) -> Result<bool>;
fn xgroup_create_consumer<K: AsRef<[u8]>>(
&self,
key: K,
group_name: &str,
consumer_name: &str,
) -> Result<i32>;
fn xgroup_del_consumer<K: AsRef<[u8]>>(
&self,
key: K,
group_name: &str,
consumer_name: &str,
) -> Result<u64>;
fn xgroup_set_id<K: AsRef<[u8]>>(
&self,
key: K,
group_name: &str,
last_id: &str,
entries_read: Option<i64>,
) -> Result<()>;
fn xack<K: AsRef<[u8]>>(&self, key: K, group_name: &str, entry_ids: &[StreamId]) -> Result<u64>;
fn xclaim<K: AsRef<[u8]>>(
&self,
key: K,
group_name: &str,
consumer_name: &str,
min_idle_time_ms: u64,
entry_ids: &[StreamId],
options: StreamClaimOptions,
) -> Result<StreamClaimResult>;
fn xautoclaim<K: AsRef<[u8]>>(
&self,
key: K,
group_name: &str,
consumer_name: &str,
options: StreamAutoClaimOptions,
) -> Result<StreamAutoClaimResult>;
fn xpending_summary<K: AsRef<[u8]>>(
&self,
key: K,
group_name: &str,
) -> Result<StreamGetPendingEntryResult>;
fn xpending_range<K: AsRef<[u8]>>(
&self,
key: K,
group_name: &str,
options: StreamPendingOptions,
) -> Result<Vec<StreamNack>>;
fn xread<K: AsRef<[u8]>>(
&self,
key: K,
start_id: StreamId,
count: Option<usize>,
) -> Result<Vec<StreamEntry>>;
fn xread_streams(
&self,
streams: &[(&str, StreamId)],
count: Option<usize>,
) -> Result<Vec<StreamReadResult>>;
fn xreadgroup<K: AsRef<[u8]>>(
&self,
key: K,
group_name: &str,
consumer_name: &str,
start_id_str: &str,
count: Option<usize>,
noack: bool,
) -> Result<Vec<StreamEntry>>;
fn xreadgroup_streams(
&self,
group_name: &str,
consumer_name: &str,
streams: &[(&str, &str)],
count: Option<usize>,
noack: bool,
) -> Result<Vec<StreamReadResult>>;
fn xinfo_stream<K: AsRef<[u8]>>(
&self,
key: K,
full: bool,
count: Option<usize>,
) -> Result<StreamInfo>;
fn xinfo_groups<K: AsRef<[u8]>>(&self, key: K) -> Result<Vec<(String, StreamConsumerGroupMeta)>>;
fn xinfo_consumers<K: AsRef<[u8]>>(
&self,
key: K,
group_name: &str,
) -> Result<Vec<(String, StreamConsumerMeta)>>;
fn xsetid<K: AsRef<[u8]>>(
&self,
key: K,
last_generated_id: StreamId,
entries_added: Option<u64>,
max_deleted_id: Option<StreamId>,
) -> Result<()>;
}