wedb_embed 0.1.2

Embedded database engine providing Redis-like APIs, built on fjall / 嵌入式数据库引擎,提供类似 Redis 的接口,底层基于 fjall 开发
Documentation
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,
};

/// 消息流结构操作接口 (Streams)
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<()>;
}