ruvio-client 0.2.9

RESP2 client for the Ruvio key-value server
Documentation
use std::time::{Duration, SystemTime};

/// A sorted-set member and its score.
#[derive(Debug, Clone, PartialEq)]
pub struct SortedSetEntry {
    /// The member.
    pub member: String,
    /// The member's score.
    pub score: f64,
}

impl SortedSetEntry {
    /// A member and its score.
    pub fn new(member: impl Into<String>, score: f64) -> Self {
        Self {
            member: member.into(),
            score,
        }
    }
}

impl<S: Into<String>> From<(S, f64)> for SortedSetEntry {
    fn from((member, score): (S, f64)) -> Self {
        Self::new(member, score)
    }
}

/// When `HEXPIRE` and its siblings may change a field's TTL. A field without a TTL counts
/// as an infinite TTL for `IfGreater` and `IfLess`.
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub enum FieldExpireCondition {
    /// Always set the TTL.
    #[default]
    Always,
    /// Only when the field has no TTL (`NX`).
    IfNoTtl,
    /// Only when the field already has a TTL (`XX`).
    IfTtl,
    /// Only when the new deadline is later (`GT`).
    IfGreater,
    /// Only when the new deadline is earlier (`LT`).
    IfLess,
}

/// `ZADD` options. The server refuses `NX` with `XX`, `GT`, or `LT`, and `GT` with `LT`.
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct SortedSetAddOptions {
    /// Only add new members (`NX`).
    pub if_not_exists: bool,
    /// Only update existing members (`XX`).
    pub if_exists: bool,
    /// Only update when the new score is greater (`GT`). A missing member is still added.
    pub greater_than: bool,
    /// Only update when the new score is lower (`LT`). A missing member is still added.
    pub less_than: bool,
    /// Count changed members as well as new ones (`CH`).
    pub changed: bool,
}

/// `XREAD` and `XREADGROUP` options.
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct StreamReadOptions {
    /// Most entries to return per stream (`COUNT`).
    pub count: Option<u64>,
    /// How long to wait for a new entry (`BLOCK`). `Duration::ZERO` waits until an entry arrives.
    /// The connection sends nothing else while it waits.
    pub block: Option<Duration>,
    /// `XREADGROUP` only: advance the group without adding pending entries (`NOACK`).
    pub no_ack: bool,
}

/// `XCLAIM` options.
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct StreamClaimOptions {
    /// Set the claimed entries' idle time (`IDLE`). Cannot be combined with [`Self::time`].
    pub idle: Option<Duration>,
    /// Set the claimed entries' last delivery time (`TIME`).
    pub time: Option<SystemTime>,
    /// Set the delivery count (`RETRYCOUNT`).
    pub retry_count: Option<u64>,
    /// Claim ids that are in the stream but not pending (`FORCE`).
    pub force: bool,
    /// Move the group's last delivered id (`LASTID`).
    pub last_id: Option<String>,
}

/// Optional narrowing for the extended `XPENDING` form.
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct StreamPendingFilter {
    /// Only entries owned by this consumer.
    pub consumer: Option<String>,
    /// Only entries idle for at least this long (`IDLE`).
    pub min_idle: Option<Duration>,
}

/// One row of the extended `XPENDING` form.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct StreamPendingEntry {
    /// Entry id.
    pub id: String,
    /// Consumer that owns the entry.
    pub consumer: String,
    /// Milliseconds since the entry was last delivered.
    pub idle_millis: i64,
    /// How many times the entry was delivered.
    pub delivery_count: i64,
}

/// One `SCAN` page. A cursor of 0 means the scan is complete.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ScanPage {
    /// Cursor for the next call.
    pub cursor: u64,
    /// Keys on this page.
    pub keys: Vec<String>,
}

/// One `HSCAN` page. A cursor of 0 means the scan is complete.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct HashScanPage {
    /// Cursor for the next call.
    pub cursor: u64,
    /// Field and value pairs on this page.
    pub fields: Vec<(String, String)>,
}

/// One stream entry.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct StreamEntry {
    /// Entry id.
    pub id: String,
    /// Field and value pairs in the order they were added.
    pub fields: Vec<(String, String)>,
}

/// Entries read from one stream key.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct StreamReadResult {
    /// Stream key.
    pub key: String,
    /// Entries read from the key.
    pub entries: Vec<StreamEntry>,
}

/// A Pub/Sub delivery: `message`, `pmessage`, or `smessage`.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PubSubMessage {
    /// `message`, `pmessage`, or `smessage`.
    pub kind: String,
    /// The matching pattern for `pmessage`.
    pub pattern: Option<String>,
    /// Channel the message was published to.
    pub channel: String,
    /// Raw payload.
    pub payload: Vec<u8>,
}