pub mod conf;
pub mod meta;
pub use conf::{
StreamAddOptions, StreamAutoClaimOptions, StreamClaimOptions, StreamLenOptions,
StreamPendingOptions, StreamRangeOptions, StreamReadOptions, StreamTrimOptions,
StreamTrimStrategy, StreamXGroupCreateOptions, XAdd, XGroup, XRange, XRead, XTrim,
};
pub use meta::{
NextStreamEntryIdStrategy, StreamAutoClaimResult, StreamClaimResult, StreamConsumerGroupMeta,
StreamConsumerMeta, StreamGetPendingEntryResult, StreamId, StreamInfo, StreamMeta, StreamNack,
StreamPelEntry, StreamReadResult, StreamSubkeyType,
};
pub type StreamEntry = (StreamId, Vec<(String, String)>);
use rapidhash::{HashMapExt, RapidHashMap as HashMap, RapidHashSet as HashSet};
use std::collections::{BTreeMap, VecDeque};
use std::str;
use crate::db::WeDb;
use crate::error::{Error, Result};
use crate::key_composer::KeyComposer;
use crate::meta::generate_version;
use crate::stream::meta::parse_hex_u64_16;
const DEFAULT_NS: &str = "default";
const ERR_STREAM_NOT_FOUND: &str = "Stream not found";
const ERR_GROUP_NOT_FOUND: &str = "Consumer group not found";
const ERR_GROUP_BUSY: &str = "BUSYGROUP Consumer Group name already exists";
const ERR_XGROUP_KEY_REQUIRE_EXIST: &str = "The XGROUP subcommand requires the key to exist. Note that for CREATE you may want to use the MKSTREAM option to create an empty stream automatically.";
const ERR_XGROUP_KEY_MUST_EXIST: &str = "The XGROUP subcommand requires the key to exist.";
const ERR_XGROUP_KEY_GROUP_MUST_EXIST: &str =
"The XGROUP subcommand requires the key and group to exist.";
const ERR_XGROUP_GROUP_MUST_EXIST: &str = "The XGROUP subcommand requires the group to exist.";
const ERR_INVALID_START_ID_INTERVAL: &str = "invalid start ID for the interval";
const ERR_INVALID_END_ID_INTERVAL: &str = "invalid end ID for the interval";
const ERR_SET_ID_SMALLER_THAN_TOP: &str =
"The ID specified in XSETID is smaller than the target stream top item";
const ERR_SET_ID_ENTRIES_ADDED_SMALLER: &str =
"The entries_added specified in XSETID is smaller than the target stream length";
const ERR_SET_ID_MAX_DEL_GREATER: &str =
"The ID specified in XSETID is smaller than the provided max_deleted_entry_id";
const ERR_EMPTY_STREAM_ENTRIES_ADDED: &str =
"an empty stream should have non-zero value of ENTRIESADDED";
const ERR_EMPTY_STREAM_MAX_DELETED: &str = "an empty stream should have MAXDELETEDID";
#[inline]
pub(crate) fn parse_stream_id_from_subkey(sub: &[u8]) -> Option<StreamId> {
if sub.len() == 33
&& sub[16] == b':'
&& let (Some(ms), Some(seq)) = (parse_hex_u64_16(&sub[..16]), parse_hex_u64_16(&sub[17..]))
{
return Some(StreamId::new(ms, seq));
}
let id_sub = str::from_utf8(sub).ok()?;
let (ms_hex, seq_hex) = id_sub.split_once(':')?;
let ms = u64::from_str_radix(ms_hex, 16).ok()?;
let seq = u64::from_str_radix(seq_hex, 16).ok()?;
Some(StreamId::new(ms, seq))
}
#[inline]
pub fn encode_stream_entry_pairs<FK: AsRef<[u8]>, FV: AsRef<[u8]>>(fields: &[(FK, FV)]) -> Vec<u8> {
let mut total_len = 4;
for (k, v) in fields {
total_len += 4 + k.as_ref().len() + 4 + v.as_ref().len();
}
let mut buf = Vec::with_capacity(total_len);
buf.extend_from_slice(&(fields.len() as u32).to_be_bytes());
for (k, v) in fields {
let kb = k.as_ref();
let vb = v.as_ref();
buf.extend_from_slice(&(kb.len() as u32).to_be_bytes());
buf.extend_from_slice(kb);
buf.extend_from_slice(&(vb.len() as u32).to_be_bytes());
buf.extend_from_slice(vb);
}
buf
}
#[inline]
pub fn encode_stream_entry_fields(fields: &[(String, String)]) -> Vec<u8> {
encode_stream_entry_pairs(fields)
}
#[inline]
pub fn decode_stream_entry_fields(mut bytes: &[u8]) -> Option<Vec<(String, String)>> {
if bytes.len() < 4 {
return None;
}
let count = u32::from_be_bytes(bytes[..4].try_into().ok()?) as usize;
bytes = &bytes[4..];
let mut fields = Vec::with_capacity(count);
for _ in 0..count {
if bytes.len() < 4 {
return None;
}
let klen = u32::from_be_bytes(bytes[..4].try_into().ok()?) as usize;
bytes = &bytes[4..];
if bytes.len() < klen {
return None;
}
let k = str::from_utf8(&bytes[..klen]).ok()?.to_string();
bytes = &bytes[klen..];
if bytes.len() < 4 {
return None;
}
let vlen = u32::from_be_bytes(bytes[..4].try_into().ok()?) as usize;
bytes = &bytes[4..];
if bytes.len() < vlen {
return None;
}
let v = str::from_utf8(&bytes[..vlen]).ok()?.to_string();
bytes = &bytes[vlen..];
fields.push((k, v));
}
Some(fields)
}
pub fn check_lag_valid(stream_meta: &StreamMeta, group_meta: &mut StreamConsumerGroupMeta) {
let mut valid = false;
if stream_meta.entries_added == 0 {
group_meta.lag = 0;
valid = true;
} else if group_meta.entries_read != -1
&& !stream_range_has_tombstones(stream_meta, group_meta.last_delivered_id)
{
group_meta.lag = stream_meta
.entries_added
.saturating_sub(group_meta.entries_read as u64);
valid = true;
} else {
let entries_read = stream_estimate_distance_from_first_ever_entry(
stream_meta,
group_meta.last_delivered_id,
);
if entries_read != -1 {
group_meta.lag = stream_meta
.entries_added
.saturating_sub(entries_read as u64);
valid = true;
}
}
if !valid {
group_meta.lag = u64::MAX;
}
}
fn stream_range_has_tombstones(meta: &StreamMeta, start_id: StreamId) -> bool {
let end_id = StreamId::max();
if meta.base.size == 0 || meta.max_deleted_entry_id.is_min() {
return false;
}
if meta.first_entry_id > meta.max_deleted_entry_id {
return false;
}
start_id <= meta.max_deleted_entry_id && meta.max_deleted_entry_id <= end_id
}
fn stream_estimate_distance_from_first_ever_entry(meta: &StreamMeta, id: StreamId) -> i64 {
if meta.entries_added == 0 {
return 0;
}
if meta.base.size == 0 && id < meta.last_entry_id {
return meta.entries_added as i64;
}
if id == meta.last_entry_id {
return meta.entries_added as i64;
} else if id > meta.last_entry_id {
return -1;
}
if meta.max_deleted_entry_id.is_min() || meta.max_deleted_entry_id < meta.first_entry_id {
if id < meta.first_entry_id {
return meta.entries_added.saturating_sub(meta.base.size) as i64;
} else if id == meta.first_entry_id {
return (meta.entries_added.saturating_sub(meta.base.size) + 1) as i64;
}
}
-1
}
impl WeDb {
#[inline]
fn get_stream_meta(&self, key_str: &str, now_ms: u64) -> Result<Option<StreamMeta>> {
let kc = KeyComposer::new(DEFAULT_NS);
let meta_k = kc.stream_meta(key_str);
match self.meta_ks.get(meta_k.as_bytes())? {
Some(m_bytes) => {
if let Some(meta) = StreamMeta::decode(&m_bytes)
&& !meta.is_expired(now_ms)
{
return Ok(Some(meta));
}
Ok(None)
}
None => Ok(None),
}
}
pub fn xlast_id<K: AsRef<[u8]>>(&self, key: K) -> Result<StreamId> {
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let now_ms = coarsetime::Clock::now_since_epoch().as_millis();
match self.get_stream_meta(k_str, now_ms)? {
Some(meta) => Ok(meta.last_generated_id),
None => Ok(StreamId::min()),
}
}
fn trim_stream_internal(
&self,
meta: &mut StreamMeta,
key: &str,
options: StreamTrimOptions,
batch: &mut fjall::OwnedWriteBatch,
) -> Result<u64> {
if meta.base.size == 0 {
return Ok(0);
}
if options.strategy == StreamTrimStrategy::MaxLen && meta.base.size <= options.max_len {
return Ok(0);
}
if options.strategy == StreamTrimStrategy::MinId && meta.first_entry_id >= options.min_id {
return Ok(0);
}
let kc = KeyComposer::new(DEFAULT_NS);
let prefix = kc.stream_prefix(key);
let mut delete_cnt = 0u64;
let mut last_deleted_id = StreamId::min();
let mut new_first_id = StreamId::min();
let mut found_next_first = false;
let mut iter = self.data_ks.prefix(&prefix);
while let Some(g) = iter.next() {
let (k, _) = g.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
if let Some(sid) = parse_stream_id_from_subkey(&k[prefix.len()..]) {
if options.strategy == StreamTrimStrategy::MaxLen
&& (meta.base.size.saturating_sub(delete_cnt)) <= options.max_len
{
new_first_id = sid;
found_next_first = true;
break;
}
if options.strategy == StreamTrimStrategy::MinId && sid >= options.min_id {
new_first_id = sid;
found_next_first = true;
break;
}
delete_cnt += 1;
last_deleted_id = sid;
batch.remove(&self.data_ks, &*k);
if let Some(lim) = options.limit
&& delete_cnt as usize >= lim
{
for next_g in iter.by_ref() {
let (nk, _) = next_g.into_inner()?;
if !nk.starts_with(&prefix) {
break;
}
if let Some(nsid) = parse_stream_id_from_subkey(&nk[prefix.len()..]) {
new_first_id = nsid;
found_next_first = true;
break;
}
}
break;
}
}
}
if delete_cnt > 0 {
meta.base.size = meta.base.size.saturating_sub(delete_cnt);
if last_deleted_id > meta.max_deleted_entry_id {
meta.max_deleted_entry_id = last_deleted_id;
}
if meta.base.size == 0 {
meta.first_entry_id.clear();
meta.last_entry_id.clear();
meta.recorded_first_entry_id.clear();
} else if found_next_first {
meta.first_entry_id = new_first_id;
meta.recorded_first_entry_id = new_first_id;
}
}
Ok(delete_cnt)
}
pub fn xadd<K: AsRef<[u8]>, FK: AsRef<[u8]>, FV: AsRef<[u8]>>(
&self,
key: K,
options: impl Into<StreamAddOptions>,
field_vals: &[(FK, FV)],
) -> Result<StreamId> {
let options = options.into();
let kc = KeyComposer::new(DEFAULT_NS);
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.stream_meta(k_str);
let now_ms = coarsetime::Clock::now_since_epoch().as_millis();
let opt_m = self.meta_ks.get(meta_k.as_bytes())?;
let mut meta = match opt_m {
Some(m_bytes) => {
let decoded = StreamMeta::decode(&m_bytes).unwrap_or_else(|| StreamMeta::new(0, 0));
if decoded.is_expired(now_ms) {
if options.nomkstream {
return Err(Error::not_found(ERR_STREAM_NOT_FOUND));
}
StreamMeta::new(0, generate_version())
} else {
decoded
}
}
None => {
if options.nomkstream {
return Err(Error::not_found(ERR_STREAM_NOT_FOUND));
}
StreamMeta::new(0, generate_version())
}
};
let next_entry_id = options
.next_id_strategy
.generate_id(meta.last_generated_id, now_ms)?;
let mut batch = self.db.batch();
let mut should_add = true;
if options.trim_options.strategy != StreamTrimStrategy::None {
let mut trim_opts = options.trim_options;
if trim_opts.strategy == StreamTrimStrategy::MaxLen {
trim_opts.max_len = if trim_opts.max_len > 0 {
trim_opts.max_len - 1
} else {
0
};
}
self.trim_stream_internal(&mut meta, k_str, trim_opts, &mut batch)?;
if trim_opts.strategy == StreamTrimStrategy::MinId && next_entry_id < trim_opts.min_id {
should_add = false;
}
if trim_opts.strategy == StreamTrimStrategy::MaxLen && options.trim_options.max_len == 0
{
should_add = false;
}
}
if should_add {
let encoded_payload = encode_stream_entry_pairs(field_vals);
let item_k = kc.stream_item(k_str, next_entry_id.ms, next_entry_id.seq);
batch.insert(&self.data_ks, item_k.as_bytes(), encoded_payload);
meta.last_generated_id = next_entry_id;
meta.last_entry_id = next_entry_id;
meta.base.size += 1;
if meta.base.size == 1 {
meta.first_entry_id = next_entry_id;
meta.recorded_first_entry_id = next_entry_id;
}
} else {
meta.last_generated_id = next_entry_id;
meta.max_deleted_entry_id = next_entry_id;
}
meta.entries_added += 1;
batch.insert(&self.meta_ks, meta_k.as_bytes(), meta.encode());
batch.commit()?;
Ok(next_entry_id)
}
pub 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> {
let strategy = match id_opt {
Some(id) => NextStreamEntryIdStrategy::FullySpecified(id),
None => NextStreamEntryIdStrategy::Auto,
};
let options = StreamAddOptions {
trim_options: StreamTrimOptions::none(),
next_id_strategy: strategy,
nomkstream: false,
};
self.xadd(key, options, field_vals)
}
pub fn xlen<K: AsRef<[u8]>>(&self, key: K) -> Result<u64> {
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let now_ms = coarsetime::Clock::now_since_epoch().as_millis();
match self.get_stream_meta(k_str, now_ms)? {
Some(meta) => Ok(meta.base.size),
None => Ok(0),
}
}
pub fn xlen_with_options<K: AsRef<[u8]>>(
&self,
key: K,
options: StreamLenOptions,
) -> Result<u64> {
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let now_ms = coarsetime::Clock::now_since_epoch().as_millis();
let meta = match self.get_stream_meta(k_str, now_ms)? {
Some(meta) => meta,
None => return Ok(0),
};
if !options.with_entry_id {
return Ok(meta.base.size);
}
if options.entry_id > meta.last_entry_id {
return Ok(if options.to_first { meta.base.size } else { 0 });
}
if options.entry_id < meta.first_entry_id {
return Ok(if options.to_first { 0 } else { meta.base.size });
}
if (!options.to_first && options.entry_id == meta.first_entry_id)
|| (options.to_first && options.entry_id == meta.last_entry_id)
{
return Ok(meta.base.size.saturating_sub(1));
}
let kc = KeyComposer::new(DEFAULT_NS);
let prefix = kc.stream_prefix(k_str);
let mut count = 0u64;
for g in self.data_ks.prefix(&prefix) {
let (k, _) = g.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
if let Some(sid) = parse_stream_id_from_subkey(&k[prefix.len()..]) {
if options.to_first {
if sid >= options.entry_id {
break;
}
count += 1;
} else if sid > options.entry_id {
count += 1;
}
}
}
Ok(count)
}
pub fn xrange<K: AsRef<[u8]>>(
&self,
key: K,
start: StreamId,
end: StreamId,
count: Option<usize>,
) -> Result<Vec<StreamEntry>> {
let options = StreamRangeOptions {
start,
end,
count,
reverse: false,
exclude_start: false,
exclude_end: false,
};
self.xrange_with_options(key, options)
}
pub fn xrevrange<K: AsRef<[u8]>>(
&self,
key: K,
end: StreamId,
start: StreamId,
count: Option<usize>,
) -> Result<Vec<StreamEntry>> {
let options = StreamRangeOptions {
start: end,
end: start,
count,
reverse: true,
exclude_start: false,
exclude_end: false,
};
self.xrange_with_options(key, options)
}
pub fn xrange_with_options<K: AsRef<[u8]>>(
&self,
key: K,
options: StreamRangeOptions,
) -> Result<Vec<StreamEntry>> {
if options.exclude_start && options.start.is_max() {
return Err(Error::invalid_data(ERR_INVALID_START_ID_INTERVAL));
}
if options.exclude_end && options.end.is_min() {
return Err(Error::invalid_data(ERR_INVALID_END_ID_INTERVAL));
}
if let Some(0) = options.count {
return Ok(Vec::new());
}
if (!options.reverse && options.end < options.start)
|| (options.reverse && options.start < options.end)
{
return Ok(Vec::new());
}
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let now_ms = coarsetime::Clock::now_since_epoch().as_millis();
if self.get_stream_meta(k_str, now_ms)?.is_none() {
return Ok(Vec::new());
}
let kc = KeyComposer::new(DEFAULT_NS);
let prefix = kc.stream_prefix(k_str);
let max_count = options.count.unwrap_or(usize::MAX);
if !options.reverse {
let mut results = Vec::new();
for g in self.data_ks.prefix(&prefix) {
let (k, v) = g.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
if let Some(sid) = parse_stream_id_from_subkey(&k[prefix.len()..]) {
if options.exclude_end {
if sid >= options.end {
break;
}
} else if sid > options.end {
break;
}
let within_start = if options.exclude_start {
sid > options.start
} else {
sid >= options.start
};
let within_end = if options.exclude_end {
sid < options.end
} else {
sid <= options.end
};
if within_start && within_end {
let fields = decode_stream_entry_fields(&v).unwrap_or_default();
results.push((sid, fields));
if results.len() >= max_count {
break;
}
}
}
}
Ok(results)
} else {
let mut window: VecDeque<(StreamId, fjall::Slice)> = VecDeque::new();
for g in self.data_ks.prefix(&prefix) {
let (k, v) = g.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
if let Some(sid) = parse_stream_id_from_subkey(&k[prefix.len()..]) {
if options.exclude_start {
if sid >= options.start {
break;
}
} else if sid > options.start {
break;
}
let within_start = if options.exclude_start {
sid < options.start
} else {
sid <= options.start
};
let within_end = if options.exclude_end {
sid > options.end
} else {
sid >= options.end
};
if within_start && within_end {
if window.len() >= max_count {
window.pop_front();
}
window.push_back((sid, v));
}
}
}
let mut results = Vec::with_capacity(window.len());
for (sid, v) in window.into_iter().rev() {
let fields = decode_stream_entry_fields(&v).unwrap_or_default();
results.push((sid, fields));
}
Ok(results)
}
}
pub fn xtrim<K: AsRef<[u8]>>(&self, key: K, options: StreamTrimOptions) -> Result<u64> {
if options.strategy == StreamTrimStrategy::None {
return Ok(0);
}
let kc = KeyComposer::new(DEFAULT_NS);
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.stream_meta(k_str);
let now_ms = coarsetime::Clock::now_since_epoch().as_millis();
let mut meta = match self.get_stream_meta(k_str, now_ms)? {
Some(meta) => meta,
None => return Ok(0),
};
let mut batch = self.db.batch();
let deleted = self.trim_stream_internal(&mut meta, k_str, options, &mut batch)?;
if deleted > 0 {
batch.insert(&self.meta_ks, meta_k.as_bytes(), meta.encode());
batch.commit()?;
}
Ok(deleted)
}
pub fn xdel<K: AsRef<[u8]>>(&self, key: K, ids: &[StreamId]) -> Result<u64> {
let kc = KeyComposer::new(DEFAULT_NS);
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.stream_meta(k_str);
let now_ms = coarsetime::Clock::now_since_epoch().as_millis();
let mut meta = match self.get_stream_meta(k_str, now_ms)? {
Some(meta) => meta,
None => return Ok(0),
};
if meta.base.size == 0 {
return Ok(0);
}
let mut deleted_cnt = 0u64;
let mut deleted_ids = HashSet::default();
let mut batch = self.db.batch();
for &id in ids {
let item_k = kc.stream_item(k_str, id.ms, id.seq);
if self.data_ks.contains_key(item_k.as_bytes())? {
deleted_cnt += 1;
deleted_ids.insert(id);
batch.remove(&self.data_ks, item_k.as_bytes());
if meta.max_deleted_entry_id < id {
meta.max_deleted_entry_id = id;
}
}
}
if deleted_cnt > 0 {
meta.base.size = meta.base.size.saturating_sub(deleted_cnt);
if meta.base.size == 0 {
meta.first_entry_id.clear();
meta.last_entry_id.clear();
meta.recorded_first_entry_id.clear();
} else {
let need_new_first = deleted_ids.contains(&meta.first_entry_id);
let need_new_last = deleted_ids.contains(&meta.last_entry_id);
if need_new_first || need_new_last {
let prefix = kc.stream_prefix(k_str);
let mut found_first = !need_new_first;
let mut new_last = StreamId::min();
for g in self.data_ks.prefix(&prefix) {
let (k, _) = g.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
if let Some(sid) = parse_stream_id_from_subkey(&k[prefix.len()..])
&& !deleted_ids.contains(&sid)
{
if !found_first {
meta.first_entry_id = sid;
meta.recorded_first_entry_id = sid;
found_first = true;
}
if need_new_last && sid > new_last {
new_last = sid;
}
}
}
if need_new_last && !new_last.is_min() {
meta.last_entry_id = new_last;
}
}
}
batch.insert(&self.meta_ks, meta_k.as_bytes(), meta.encode());
batch.commit()?;
}
Ok(deleted_cnt)
}
pub fn xgroup_create<K: AsRef<[u8]>>(
&self,
key: K,
group_name: &str,
last_id: &str,
mkstream: bool,
entries_read: Option<i64>,
) -> Result<()> {
let kc = KeyComposer::new(DEFAULT_NS);
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.stream_meta(k_str);
let now_ms = coarsetime::Clock::now_since_epoch().as_millis();
let opt_m = self.meta_ks.get(meta_k.as_bytes())?;
let mut meta = match opt_m {
Some(m_bytes) => {
let decoded = StreamMeta::decode(&m_bytes).unwrap_or_else(|| StreamMeta::new(0, 0));
if decoded.is_expired(now_ms) {
if !mkstream {
return Err(Error::invalid_data(ERR_XGROUP_KEY_REQUIRE_EXIST));
}
StreamMeta::new(0, generate_version())
} else {
decoded
}
}
None => {
if !mkstream {
return Err(Error::invalid_data(ERR_XGROUP_KEY_REQUIRE_EXIST));
}
StreamMeta::new(0, generate_version())
}
};
let group_k = kc.stream_group_meta(k_str, group_name);
if self.data_ks.contains_key(group_k.as_bytes())? {
return Err(Error::redis(ERR_GROUP_BUSY));
}
let last_delivered_id = if last_id == "$" {
meta.last_entry_id
} else {
StreamId::parse(last_id)?
};
let group_meta = StreamConsumerGroupMeta {
consumer_number: 0,
pending_number: 0,
last_delivered_id,
entries_read: entries_read.unwrap_or(-1),
lag: 0,
};
let mut batch = self.db.batch();
batch.insert(&self.data_ks, group_k.as_bytes(), group_meta.encode());
meta.group_number += 1;
batch.insert(&self.meta_ks, meta_k.as_bytes(), meta.encode());
batch.commit()?;
Ok(())
}
pub fn xgroup_destroy<K: AsRef<[u8]>>(&self, key: K, group_name: &str) -> Result<bool> {
let kc = KeyComposer::new(DEFAULT_NS);
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.stream_meta(k_str);
let now_ms = coarsetime::Clock::now_since_epoch().as_millis();
let mut meta = match self.get_stream_meta(k_str, now_ms)? {
Some(meta) => meta,
None => {
return Err(Error::invalid_data(ERR_XGROUP_KEY_MUST_EXIST));
}
};
let group_k = kc.stream_group_meta(k_str, group_name);
if !self.data_ks.contains_key(group_k.as_bytes())? {
return Ok(false);
}
let mut batch = self.db.batch();
batch.remove(&self.data_ks, group_k.as_bytes());
let c_prefix = kc.stream_consumer_prefix(k_str, group_name);
for g in self.data_ks.prefix(&c_prefix) {
let (k, _) = g.into_inner()?;
if k.starts_with(&c_prefix) {
batch.remove(&self.data_ks, k);
}
}
let p_prefix = kc.stream_pel_prefix(k_str, group_name);
for g in self.data_ks.prefix(&p_prefix) {
let (k, _) = g.into_inner()?;
if k.starts_with(&p_prefix) {
batch.remove(&self.data_ks, k);
}
}
meta.group_number = meta.group_number.saturating_sub(1);
batch.insert(&self.meta_ks, meta_k.as_bytes(), meta.encode());
batch.commit()?;
Ok(true)
}
pub fn xgroup_create_consumer<K: AsRef<[u8]>>(
&self,
key: K,
group_name: &str,
consumer_name: &str,
) -> Result<i32> {
let kc = KeyComposer::new(DEFAULT_NS);
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let now_ms = coarsetime::Clock::now_since_epoch().as_millis();
if self.get_stream_meta(k_str, now_ms)?.is_none() {
return Err(Error::invalid_data(ERR_XGROUP_KEY_MUST_EXIST));
}
let group_k = kc.stream_group_meta(k_str, group_name);
let group_bytes = match self.data_ks.get(group_k.as_bytes())? {
Some(b) => b,
None => {
return Err(Error::invalid_data(ERR_XGROUP_KEY_GROUP_MUST_EXIST));
}
};
let mut group_meta = StreamConsumerGroupMeta::decode(&group_bytes).unwrap_or_default();
let consumer_k = kc.stream_consumer_meta(k_str, group_name, consumer_name);
if self.data_ks.contains_key(consumer_k.as_bytes())? {
return Ok(0);
}
let consumer_meta = StreamConsumerMeta {
pending_number: 0,
last_attempted_interaction_ms: now_ms,
last_successful_interaction_ms: now_ms,
};
group_meta.consumer_number += 1;
let mut batch = self.db.batch();
batch.insert(&self.data_ks, consumer_k.as_bytes(), consumer_meta.encode());
batch.insert(&self.data_ks, group_k.as_bytes(), group_meta.encode());
batch.commit()?;
Ok(1)
}
pub fn xgroup_del_consumer<K: AsRef<[u8]>>(
&self,
key: K,
group_name: &str,
consumer_name: &str,
) -> Result<u64> {
let kc = KeyComposer::new(DEFAULT_NS);
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let group_k = kc.stream_group_meta(k_str, group_name);
let group_bytes = match self.data_ks.get(group_k.as_bytes())? {
Some(b) => b,
None => return Ok(0),
};
let mut group_meta = StreamConsumerGroupMeta::decode(&group_bytes).unwrap_or_default();
let consumer_k = kc.stream_consumer_meta(k_str, group_name, consumer_name);
let consumer_bytes = match self.data_ks.get(consumer_k.as_bytes())? {
Some(b) => b,
None => return Ok(0),
};
let consumer_meta = StreamConsumerMeta::decode(&consumer_bytes).unwrap_or_default();
let deleted_pel = consumer_meta.pending_number;
let mut batch = self.db.batch();
batch.remove(&self.data_ks, consumer_k.as_bytes());
let p_prefix = kc.stream_pel_prefix(k_str, group_name);
for g in self.data_ks.prefix(&p_prefix) {
let (k, v) = g.into_inner()?;
if k.starts_with(&p_prefix)
&& let Some(pel) = StreamPelEntry::decode(&v)
&& pel.consumer_name == consumer_name
{
batch.remove(&self.data_ks, k);
}
}
group_meta.consumer_number = group_meta.consumer_number.saturating_sub(1);
group_meta.pending_number = group_meta.pending_number.saturating_sub(deleted_pel);
batch.insert(&self.data_ks, group_k.as_bytes(), group_meta.encode());
batch.commit()?;
Ok(deleted_pel)
}
pub fn xgroup_set_id<K: AsRef<[u8]>>(
&self,
key: K,
group_name: &str,
last_id: &str,
entries_read: Option<i64>,
) -> Result<()> {
let kc = KeyComposer::new(DEFAULT_NS);
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let now_ms = coarsetime::Clock::now_since_epoch().as_millis();
let meta = match self.get_stream_meta(k_str, now_ms)? {
Some(meta) => meta,
None => {
return Err(Error::invalid_data(ERR_XGROUP_KEY_MUST_EXIST));
}
};
let group_k = kc.stream_group_meta(k_str, group_name);
let group_bytes = match self.data_ks.get(group_k.as_bytes())? {
Some(b) => b,
None => {
return Err(Error::invalid_data(ERR_XGROUP_GROUP_MUST_EXIST));
}
};
let mut group_meta = StreamConsumerGroupMeta::decode(&group_bytes).unwrap_or_default();
let parsed_id = if last_id == "$" {
meta.last_entry_id
} else {
StreamId::parse(last_id)?
};
group_meta.last_delivered_id = parsed_id;
if let Some(er) = entries_read {
group_meta.entries_read = er;
}
self.data_ks
.insert(group_k.as_bytes(), group_meta.encode())?;
Ok(())
}
pub fn xack<K: AsRef<[u8]>>(
&self,
key: K,
group_name: &str,
entry_ids: &[StreamId],
) -> Result<u64> {
let kc = KeyComposer::new(DEFAULT_NS);
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let group_k = kc.stream_group_meta(k_str, group_name);
let group_bytes = match self.data_ks.get(group_k.as_bytes())? {
Some(b) => b,
None => return Ok(0),
};
let mut group_meta = StreamConsumerGroupMeta::decode(&group_bytes).unwrap_or_default();
let mut acknowledged = 0u64;
let mut consumer_acks: HashMap<String, u64> = HashMap::new();
let mut batch = self.db.batch();
for &id in entry_ids {
let pel_k = kc.stream_pel_item(k_str, group_name, id.ms, id.seq);
if let Some(pel_bytes) = self.data_ks.get(pel_k.as_bytes())? {
if let Some(pel_entry) = StreamPelEntry::decode(&pel_bytes) {
*consumer_acks.entry(pel_entry.consumer_name).or_insert(0) += 1;
}
acknowledged += 1;
batch.remove(&self.data_ks, pel_k.as_bytes());
}
}
if acknowledged > 0 {
group_meta.pending_number = group_meta.pending_number.saturating_sub(acknowledged);
batch.insert(&self.data_ks, group_k.as_bytes(), group_meta.encode());
for (consumer_name, ack_cnt) in consumer_acks {
let consumer_k = kc.stream_consumer_meta(k_str, group_name, &consumer_name);
if let Some(c_bytes) = self.data_ks.get(consumer_k.as_bytes())?
&& let Some(mut c_meta) = StreamConsumerMeta::decode(&c_bytes)
{
c_meta.pending_number = c_meta.pending_number.saturating_sub(ack_cnt);
batch.insert(&self.data_ks, consumer_k.as_bytes(), c_meta.encode());
}
}
batch.commit()?;
}
Ok(acknowledged)
}
pub 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> {
let kc = KeyComposer::new(DEFAULT_NS);
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let now_ms = coarsetime::Clock::now_since_epoch().as_millis();
let _meta = match self.get_stream_meta(k_str, now_ms)? {
Some(meta) => meta,
None => {
return Err(Error::not_found(ERR_STREAM_NOT_FOUND));
}
};
let group_k = kc.stream_group_meta(k_str, group_name);
let group_bytes = match self.data_ks.get(group_k.as_bytes())? {
Some(b) => b,
None => {
return Err(Error::not_found(ERR_GROUP_NOT_FOUND));
}
};
let mut group_meta = StreamConsumerGroupMeta::decode(&group_bytes).unwrap_or_default();
let consumer_k = kc.stream_consumer_meta(k_str, group_name, consumer_name);
let mut consumer_meta = match self.data_ks.get(consumer_k.as_bytes())? {
Some(b) => StreamConsumerMeta::decode(&b).unwrap_or_default(),
None => {
group_meta.consumer_number += 1;
StreamConsumerMeta::default()
}
};
consumer_meta.last_attempted_interaction_ms = now_ms;
consumer_meta.last_successful_interaction_ms = now_ms;
let mut result = StreamClaimResult::default();
let mut batch = self.db.batch();
for &id in entry_ids {
let item_k = kc.stream_item(k_str, id.ms, id.seq);
let raw_item_val = self.data_ks.get(item_k.as_bytes())?;
if raw_item_val.is_none() {
continue;
}
let pel_k = kc.stream_pel_item(k_str, group_name, id.ms, id.seq);
let raw_pel_val = self.data_ks.get(pel_k.as_bytes())?;
let mut pel_entry = if let Some(b) = raw_pel_val {
StreamPelEntry::decode(&b).unwrap_or(StreamPelEntry {
last_delivery_time_ms: 0,
last_delivery_count: 0,
consumer_name: String::new(),
})
} else if options.force {
group_meta.pending_number += 1;
StreamPelEntry {
last_delivery_time_ms: 0,
last_delivery_count: 0,
consumer_name: String::new(),
}
} else {
continue;
};
if now_ms.saturating_sub(pel_entry.last_delivery_time_ms) < min_idle_time_ms {
continue;
}
if options.just_id {
result.ids.push(id);
} else if let Some(val_bytes) = raw_item_val {
let fields = decode_stream_entry_fields(&val_bytes).unwrap_or_default();
result.entries.push((id, fields));
}
if !pel_entry.consumer_name.is_empty() && pel_entry.consumer_name != consumer_name {
let orig_k = kc.stream_consumer_meta(k_str, group_name, &pel_entry.consumer_name);
if let Some(orig_b) = self.data_ks.get(orig_k.as_bytes())?
&& let Some(mut orig_c_meta) = StreamConsumerMeta::decode(&orig_b)
{
orig_c_meta.pending_number = orig_c_meta.pending_number.saturating_sub(1);
batch.insert(&self.data_ks, orig_k.as_bytes(), orig_c_meta.encode());
}
}
if pel_entry.consumer_name != consumer_name {
consumer_meta.pending_number += 1;
pel_entry.consumer_name = consumer_name.to_string();
}
if options.with_time {
pel_entry.last_delivery_time_ms = options.last_delivery_time_ms;
} else {
pel_entry.last_delivery_time_ms = now_ms.saturating_sub(options.idle_time_ms);
}
if pel_entry.last_delivery_time_ms > now_ms {
pel_entry.last_delivery_time_ms = now_ms;
}
if options.with_retry_count {
pel_entry.last_delivery_count = options.last_delivery_count;
} else if !options.just_id {
pel_entry.last_delivery_count += 1;
}
batch.insert(&self.data_ks, pel_k.as_bytes(), pel_entry.encode());
}
if let Some(last_deliv) = options.last_delivered_id
&& last_deliv > group_meta.last_delivered_id
{
group_meta.last_delivered_id = last_deliv;
}
batch.insert(&self.data_ks, consumer_k.as_bytes(), consumer_meta.encode());
batch.insert(&self.data_ks, group_k.as_bytes(), group_meta.encode());
batch.commit()?;
Ok(result)
}
pub fn xautoclaim<K: AsRef<[u8]>>(
&self,
key: K,
group_name: &str,
consumer_name: &str,
options: StreamAutoClaimOptions,
) -> Result<StreamAutoClaimResult> {
if options.exclude_start && options.start_id.is_max() {
return Err(Error::invalid_data(ERR_INVALID_START_ID_INTERVAL));
}
let kc = KeyComposer::new(DEFAULT_NS);
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let now_ms = coarsetime::Clock::now_since_epoch().as_millis();
let _meta = match self.get_stream_meta(k_str, now_ms)? {
Some(meta) => meta,
None => {
return Err(Error::not_found(ERR_STREAM_NOT_FOUND));
}
};
let group_k = kc.stream_group_meta(k_str, group_name);
let group_bytes = match self.data_ks.get(group_k.as_bytes())? {
Some(b) => b,
None => {
return Err(Error::not_found(ERR_GROUP_NOT_FOUND));
}
};
let mut group_meta = StreamConsumerGroupMeta::decode(&group_bytes).unwrap_or_default();
let consumer_k = kc.stream_consumer_meta(k_str, group_name, consumer_name);
let mut consumer_meta = match self.data_ks.get(consumer_k.as_bytes())? {
Some(b) => StreamConsumerMeta::decode(&b).unwrap_or_default(),
None => {
group_meta.consumer_number += 1;
StreamConsumerMeta::default()
}
};
consumer_meta.last_attempted_interaction_ms = now_ms;
consumer_meta.last_successful_interaction_ms = now_ms;
let p_prefix = kc.stream_pel_prefix(k_str, group_name);
let mut count = options.count;
let mut attempts = options.attempts_factors * count;
let mut deleted_entries = Vec::new();
let mut claimed_entries = Vec::new();
let mut next_claim_id = StreamId::min();
let mut has_next = false;
let mut last_scanned_id = options.start_id;
let mut batch = self.db.batch();
let mut claimed_from_others: HashMap<String, u64> = HashMap::new();
let mut deleted_from_consumers: HashMap<String, u64> = HashMap::new();
let mut iter = self.data_ks.prefix(&p_prefix);
for g in iter.by_ref() {
if count == 0 || attempts == 0 {
let (k, _) = g.into_inner()?;
if k.starts_with(&p_prefix)
&& let Some(sid) = parse_stream_id_from_subkey(&k[p_prefix.len()..])
&& sid > last_scanned_id
{
next_claim_id = sid;
has_next = true;
}
break;
}
let (k, v) = g.into_inner()?;
if !k.starts_with(&p_prefix) {
break;
}
if let Some(sid) = parse_stream_id_from_subkey(&k[p_prefix.len()..]) {
if sid < options.start_id {
continue;
}
if options.exclude_start && sid == options.start_id {
continue;
}
last_scanned_id = sid;
attempts -= 1;
if let Some(mut pel_entry) = StreamPelEntry::decode(&v) {
if now_ms.saturating_sub(pel_entry.last_delivery_time_ms)
< options.min_idle_time_ms
{
continue;
}
let item_k = kc.stream_item(k_str, sid.ms, sid.seq);
let raw_item_val = self.data_ks.get(item_k.as_bytes())?;
if raw_item_val.is_none() {
deleted_entries.push(sid);
batch.remove(&self.data_ks, &*k);
*deleted_from_consumers
.entry(pel_entry.consumer_name.clone())
.or_insert(0) += 1;
count -= 1;
continue;
}
let fields = if options.just_id {
Vec::new()
} else if let Some(ref val_slice) = raw_item_val {
decode_stream_entry_fields(val_slice).unwrap_or_default()
} else {
Vec::new()
};
claimed_entries.push((sid, fields));
count -= 1;
if pel_entry.consumer_name != consumer_name {
*claimed_from_others
.entry(pel_entry.consumer_name.clone())
.or_insert(0) += 1;
pel_entry.consumer_name = consumer_name.to_string();
pel_entry.last_delivery_time_ms = now_ms;
if !options.just_id {
pel_entry.last_delivery_count += 1;
}
batch.insert(&self.data_ks, &*k, pel_entry.encode());
}
}
}
}
if !has_next && (count == 0 || attempts == 0) {
for g in iter {
let (k, _) = g.into_inner()?;
if !k.starts_with(&p_prefix) {
break;
}
if let Some(sid) = parse_stream_id_from_subkey(&k[p_prefix.len()..])
&& sid > last_scanned_id
{
next_claim_id = sid;
has_next = true;
break;
}
}
}
if !has_next {
next_claim_id = StreamId::min();
}
let total_claimed: u64 = claimed_from_others.values().sum();
let total_deleted = deleted_entries.len() as u64;
if total_claimed > 0 || total_deleted > 0 {
consumer_meta.pending_number += total_claimed;
if let Some(dec) = deleted_from_consumers.get(consumer_name) {
consumer_meta.pending_number = consumer_meta.pending_number.saturating_sub(*dec);
}
batch.insert(&self.data_ks, consumer_k.as_bytes(), consumer_meta.encode());
for (other_consumer, claimed_cnt) in claimed_from_others {
let other_k = kc.stream_consumer_meta(k_str, group_name, &other_consumer);
if let Some(other_b) = self.data_ks.get(other_k.as_bytes())?
&& let Some(mut other_meta) = StreamConsumerMeta::decode(&other_b)
{
let del_cnt = deleted_from_consumers
.get(&other_consumer)
.copied()
.unwrap_or(0);
other_meta.pending_number = other_meta
.pending_number
.saturating_sub(claimed_cnt + del_cnt);
batch.insert(&self.data_ks, other_k.as_bytes(), other_meta.encode());
}
}
if total_deleted > 0 {
group_meta.pending_number = group_meta.pending_number.saturating_sub(total_deleted);
batch.insert(&self.data_ks, group_k.as_bytes(), group_meta.encode());
}
}
batch.commit()?;
Ok(StreamAutoClaimResult {
next_claim_id,
entries: claimed_entries,
deleted_ids: deleted_entries,
})
}
pub fn xpending_summary<K: AsRef<[u8]>>(
&self,
key: K,
group_name: &str,
) -> Result<StreamGetPendingEntryResult> {
let kc = KeyComposer::new(DEFAULT_NS);
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let p_prefix = kc.stream_pel_prefix(k_str, group_name);
let mut first_id = StreamId::max();
let mut last_id = StreamId::min();
let mut total_pending = 0u64;
let mut consumer_counts: BTreeMap<String, u64> = BTreeMap::new();
for g in self.data_ks.prefix(&p_prefix) {
let (k, v) = g.into_inner()?;
if !k.starts_with(&p_prefix) {
break;
}
if let Some(sid) = parse_stream_id_from_subkey(&k[p_prefix.len()..])
&& let Some(pel) = StreamPelEntry::decode(&v)
{
total_pending += 1;
if sid < first_id {
first_id = sid;
}
if sid > last_id {
last_id = sid;
}
*consumer_counts.entry(pel.consumer_name).or_insert(0) += 1;
}
}
if total_pending == 0 {
return Ok(StreamGetPendingEntryResult {
pending_number: 0,
first_entry_id: StreamId::min(),
last_entry_id: StreamId::min(),
consumer_infos: Vec::new(),
});
}
let consumer_infos = consumer_counts.into_iter().collect();
Ok(StreamGetPendingEntryResult {
pending_number: total_pending,
first_entry_id: first_id,
last_entry_id: last_id,
consumer_infos,
})
}
pub fn xpending_range<K: AsRef<[u8]>>(
&self,
key: K,
group_name: &str,
options: StreamPendingOptions,
) -> Result<Vec<StreamNack>> {
let kc = KeyComposer::new(DEFAULT_NS);
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let p_prefix = kc.stream_pel_prefix(k_str, group_name);
let max_count = options.count.unwrap_or(usize::MAX);
let mut results = Vec::new();
let now_ms = coarsetime::Clock::now_since_epoch().as_millis();
for g in self.data_ks.prefix(&p_prefix) {
let (k, v) = g.into_inner()?;
if !k.starts_with(&p_prefix) {
break;
}
if let Some(sid) = parse_stream_id_from_subkey(&k[p_prefix.len()..]) {
if options.exclude_end {
if sid >= options.end_id {
break;
}
} else if sid > options.end_id {
break;
}
let within_start = if options.exclude_start {
sid > options.start_id
} else {
sid >= options.start_id
};
let within_end = if options.exclude_end {
sid < options.end_id
} else {
sid <= options.end_id
};
if within_start
&& within_end
&& let Some(pel) = StreamPelEntry::decode(&v)
{
if options.with_time
&& now_ms.saturating_sub(pel.last_delivery_time_ms) < options.idle_time
{
continue;
}
if let Some(ref target_c) = options.consumer
&& &pel.consumer_name != target_c
{
continue;
}
results.push(StreamNack {
id: sid,
pel_entry: pel,
});
if results.len() >= max_count {
break;
}
}
}
}
Ok(results)
}
pub fn xread<K: AsRef<[u8]>>(
&self,
key: K,
start_id: StreamId,
count: Option<usize>,
) -> Result<Vec<StreamEntry>> {
let options = StreamRangeOptions {
start: start_id,
end: StreamId::max(),
count,
reverse: false,
exclude_start: true,
exclude_end: false,
};
self.xrange_with_options(key, options)
}
pub fn xread_streams(
&self,
streams: &[(&str, StreamId)],
count: Option<usize>,
) -> Result<Vec<StreamReadResult>> {
let mut results = Vec::with_capacity(streams.len());
for &(stream_name, start_id) in streams {
let entries = self.xread(stream_name, start_id, count)?;
if !entries.is_empty() {
results.push(StreamReadResult {
name: stream_name.to_string(),
entries,
});
}
}
Ok(results)
}
pub 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>> {
let kc = KeyComposer::new(DEFAULT_NS);
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let now_ms = coarsetime::Clock::now_since_epoch().as_millis();
let _meta = match self.get_stream_meta(k_str, now_ms)? {
Some(meta) => meta,
None => {
return Err(Error::not_found(ERR_STREAM_NOT_FOUND));
}
};
let group_k = kc.stream_group_meta(k_str, group_name);
let group_bytes = match self.data_ks.get(group_k.as_bytes())? {
Some(b) => b,
None => {
return Err(Error::not_found(ERR_GROUP_NOT_FOUND));
}
};
let mut group_meta = StreamConsumerGroupMeta::decode(&group_bytes).unwrap_or_default();
let consumer_k = kc.stream_consumer_meta(k_str, group_name, consumer_name);
let mut consumer_meta = match self.data_ks.get(consumer_k.as_bytes())? {
Some(b) => StreamConsumerMeta::decode(&b).unwrap_or_default(),
None => {
group_meta.consumer_number += 1;
StreamConsumerMeta::default()
}
};
consumer_meta.last_attempted_interaction_ms = now_ms;
consumer_meta.last_successful_interaction_ms = now_ms;
let max_count = count.unwrap_or(usize::MAX);
let mut batch = self.db.batch();
if start_id_str == ">" {
let start_id = group_meta.last_delivered_id;
let options = StreamRangeOptions {
start: start_id,
end: StreamId::max(),
count: Some(max_count),
reverse: false,
exclude_start: true,
exclude_end: false,
};
let entries = self.xrange_with_options(key.as_ref(), options)?;
let mut max_id = StreamId::min();
for (sid, _) in &entries {
if *sid > max_id {
max_id = *sid;
}
if !noack {
let pel_k = kc.stream_pel_item(k_str, group_name, sid.ms, sid.seq);
let pel_entry = StreamPelEntry {
last_delivery_time_ms: now_ms,
last_delivery_count: 1,
consumer_name: consumer_name.to_string(),
};
batch.insert(&self.data_ks, pel_k.as_bytes(), pel_entry.encode());
group_meta.entries_read += 1;
group_meta.pending_number += 1;
consumer_meta.pending_number += 1;
}
}
if max_id > group_meta.last_delivered_id {
group_meta.last_delivered_id = max_id;
}
batch.insert(&self.data_ks, group_k.as_bytes(), group_meta.encode());
batch.insert(&self.data_ks, consumer_k.as_bytes(), consumer_meta.encode());
batch.commit()?;
Ok(entries)
} else {
let start_id = StreamId::parse(start_id_str)?;
let p_prefix = kc.stream_pel_prefix(k_str, group_name);
let mut entries = Vec::new();
for g in self.data_ks.prefix(&p_prefix) {
let (k, v) = g.into_inner()?;
if !k.starts_with(&p_prefix) {
break;
}
if let Some(sid) = parse_stream_id_from_subkey(&k[p_prefix.len()..])
&& sid >= start_id
&& let Some(mut pel) = StreamPelEntry::decode(&v)
&& pel.consumer_name == consumer_name
{
let item_k = kc.stream_item(k_str, sid.ms, sid.seq);
if let Some(item_v) = self.data_ks.get(item_k.as_bytes())? {
let fields = decode_stream_entry_fields(&item_v).unwrap_or_default();
entries.push((sid, fields));
pel.last_delivery_count += 1;
pel.last_delivery_time_ms = now_ms;
batch.insert(&self.data_ks, k, pel.encode());
if entries.len() >= max_count {
break;
}
}
}
}
batch.insert(&self.data_ks, group_k.as_bytes(), group_meta.encode());
batch.insert(&self.data_ks, consumer_k.as_bytes(), consumer_meta.encode());
batch.commit()?;
Ok(entries)
}
}
pub fn xreadgroup_streams(
&self,
group_name: &str,
consumer_name: &str,
streams: &[(&str, &str)],
count: Option<usize>,
noack: bool,
) -> Result<Vec<StreamReadResult>> {
let mut results = Vec::with_capacity(streams.len());
for &(stream_name, start_id_str) in streams {
let entries = self.xreadgroup(
stream_name,
group_name,
consumer_name,
start_id_str,
count,
noack,
)?;
if !entries.is_empty() {
results.push(StreamReadResult {
name: stream_name.to_string(),
entries,
});
}
}
Ok(results)
}
pub fn xinfo_stream<K: AsRef<[u8]>>(
&self,
key: K,
full: bool,
count: Option<usize>,
) -> Result<StreamInfo> {
let kc = KeyComposer::new(DEFAULT_NS);
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let now_ms = coarsetime::Clock::now_since_epoch().as_millis();
let meta = match self.get_stream_meta(k_str, now_ms)? {
Some(meta) => meta,
None => {
return Err(Error::not_found(ERR_STREAM_NOT_FOUND));
}
};
let mut first_entry = None;
let mut last_entry = None;
if meta.base.size > 0 && !meta.first_entry_id.is_min() {
let item_k = kc.stream_item(k_str, meta.first_entry_id.ms, meta.first_entry_id.seq);
if let Some(v) = self.data_ks.get(item_k.as_bytes())? {
let fields = decode_stream_entry_fields(&v).unwrap_or_default();
first_entry = Some((meta.first_entry_id, fields));
}
}
if meta.base.size > 0 && !meta.last_entry_id.is_min() {
let item_k = kc.stream_item(k_str, meta.last_entry_id.ms, meta.last_entry_id.seq);
if let Some(v) = self.data_ks.get(item_k.as_bytes())? {
let fields = decode_stream_entry_fields(&v).unwrap_or_default();
last_entry = Some((meta.last_entry_id, fields));
}
}
let entries = if full {
self.xrange(key.as_ref(), StreamId::min(), StreamId::max(), count)?
} else {
Vec::new()
};
Ok(StreamInfo {
size: meta.base.size,
entries_added: meta.entries_added,
last_generated_id: meta.last_generated_id,
max_deleted_entry_id: meta.max_deleted_entry_id,
recorded_first_entry_id: meta.recorded_first_entry_id,
first_entry,
last_entry,
groups: meta.group_number,
entries,
})
}
pub fn xinfo_groups<K: AsRef<[u8]>>(
&self,
key: K,
) -> Result<Vec<(String, StreamConsumerGroupMeta)>> {
let kc = KeyComposer::new(DEFAULT_NS);
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let now_ms = coarsetime::Clock::now_since_epoch().as_millis();
let meta = match self.get_stream_meta(k_str, now_ms)? {
Some(meta) => meta,
None => {
return Err(Error::not_found(ERR_STREAM_NOT_FOUND));
}
};
let g_prefix = kc.stream_group_prefix(k_str);
let mut groups = Vec::new();
for g in self.data_ks.prefix(&g_prefix) {
let (k, v) = g.into_inner()?;
if !k.starts_with(&g_prefix) {
break;
}
let group_name = str::from_utf8(&k[g_prefix.len()..])
.unwrap_or("")
.to_string();
if let Some(mut g_meta) = StreamConsumerGroupMeta::decode(&v) {
check_lag_valid(&meta, &mut g_meta);
groups.push((group_name, g_meta));
}
}
Ok(groups)
}
pub fn xinfo_consumers<K: AsRef<[u8]>>(
&self,
key: K,
group_name: &str,
) -> Result<Vec<(String, StreamConsumerMeta)>> {
let kc = KeyComposer::new(DEFAULT_NS);
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let now_ms = coarsetime::Clock::now_since_epoch().as_millis();
if self.get_stream_meta(k_str, now_ms)?.is_none() {
return Err(Error::not_found(ERR_STREAM_NOT_FOUND));
}
let c_prefix = kc.stream_consumer_prefix(k_str, group_name);
let mut consumers = Vec::new();
for g in self.data_ks.prefix(&c_prefix) {
let (k, v) = g.into_inner()?;
if !k.starts_with(&c_prefix) {
break;
}
let consumer_name = str::from_utf8(&k[c_prefix.len()..])
.unwrap_or("")
.to_string();
if let Some(c_meta) = StreamConsumerMeta::decode(&v) {
consumers.push((consumer_name, c_meta));
}
}
Ok(consumers)
}
pub fn xsetid<K: AsRef<[u8]>>(
&self,
key: K,
last_generated_id: StreamId,
entries_added: Option<u64>,
max_deleted_id: Option<StreamId>,
) -> Result<()> {
if let Some(max_del) = max_deleted_id
&& last_generated_id < max_del
{
return Err(Error::invalid_data(ERR_SET_ID_MAX_DEL_GREATER));
}
let kc = KeyComposer::new(DEFAULT_NS);
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.stream_meta(k_str);
let now_ms = coarsetime::Clock::now_since_epoch().as_millis();
let opt_m = self.meta_ks.get(meta_k.as_bytes())?;
let is_empty = match &opt_m {
Some(m_bytes) => match StreamMeta::decode(m_bytes) {
Some(m) => m.is_expired(now_ms),
None => true,
},
None => true,
};
if is_empty {
if entries_added.is_none() || entries_added == Some(0) {
return Err(Error::invalid_data(ERR_EMPTY_STREAM_ENTRIES_ADDED));
}
if max_deleted_id.is_none() || max_deleted_id == Some(StreamId::min()) {
return Err(Error::invalid_data(ERR_EMPTY_STREAM_MAX_DELETED));
}
}
let mut meta = match opt_m {
Some(m_bytes) => {
let decoded = StreamMeta::decode(&m_bytes).unwrap_or_else(|| StreamMeta::new(0, 0));
if decoded.is_expired(now_ms) {
StreamMeta::new(0, generate_version())
} else {
decoded
}
}
None => StreamMeta::new(0, generate_version()),
};
if meta.base.size > 0 && last_generated_id < meta.last_generated_id {
return Err(Error::invalid_data(ERR_SET_ID_SMALLER_THAN_TOP));
}
if meta.base.size > 0
&& let Some(ea) = entries_added
&& ea < meta.base.size
{
return Err(Error::invalid_data(ERR_SET_ID_ENTRIES_ADDED_SMALLER));
}
meta.last_generated_id = last_generated_id;
if let Some(ea) = entries_added {
meta.entries_added = ea;
}
if let Some(max_del) = max_deleted_id
&& !max_del.is_min()
{
meta.max_deleted_entry_id = max_del;
}
self.meta_ks.insert(meta_k.as_bytes(), meta.encode())?;
Ok(())
}
}