Skip to main content

ruvio_client/
models.rs

1use std::time::{Duration, SystemTime};
2
3/// A sorted-set member and its score.
4#[derive(Debug, Clone, PartialEq)]
5pub struct SortedSetEntry {
6    /// The member.
7    pub member: String,
8    /// The member's score.
9    pub score: f64,
10}
11
12impl SortedSetEntry {
13    /// A member and its score.
14    pub fn new(member: impl Into<String>, score: f64) -> Self {
15        Self {
16            member: member.into(),
17            score,
18        }
19    }
20}
21
22impl<S: Into<String>> From<(S, f64)> for SortedSetEntry {
23    fn from((member, score): (S, f64)) -> Self {
24        Self::new(member, score)
25    }
26}
27
28/// When `HEXPIRE` and its siblings may change a field's TTL. A field without a TTL counts
29/// as an infinite TTL for `IfGreater` and `IfLess`.
30#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
31pub enum FieldExpireCondition {
32    /// Always set the TTL.
33    #[default]
34    Always,
35    /// Only when the field has no TTL (`NX`).
36    IfNoTtl,
37    /// Only when the field already has a TTL (`XX`).
38    IfTtl,
39    /// Only when the new deadline is later (`GT`).
40    IfGreater,
41    /// Only when the new deadline is earlier (`LT`).
42    IfLess,
43}
44
45/// `ZADD` options. The server refuses `NX` with `XX`, `GT`, or `LT`, and `GT` with `LT`.
46#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
47pub struct SortedSetAddOptions {
48    /// Only add new members (`NX`).
49    pub if_not_exists: bool,
50    /// Only update existing members (`XX`).
51    pub if_exists: bool,
52    /// Only update when the new score is greater (`GT`). A missing member is still added.
53    pub greater_than: bool,
54    /// Only update when the new score is lower (`LT`). A missing member is still added.
55    pub less_than: bool,
56    /// Count changed members as well as new ones (`CH`).
57    pub changed: bool,
58}
59
60/// `XREAD` and `XREADGROUP` options.
61#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
62pub struct StreamReadOptions {
63    /// Most entries to return per stream (`COUNT`).
64    pub count: Option<u64>,
65    /// How long to wait for a new entry (`BLOCK`). `Duration::ZERO` waits until an entry arrives.
66    /// The connection sends nothing else while it waits.
67    pub block: Option<Duration>,
68    /// `XREADGROUP` only: advance the group without adding pending entries (`NOACK`).
69    pub no_ack: bool,
70}
71
72/// `XCLAIM` options.
73#[derive(Debug, Clone, Default, PartialEq, Eq)]
74pub struct StreamClaimOptions {
75    /// Set the claimed entries' idle time (`IDLE`). Cannot be combined with [`Self::time`].
76    pub idle: Option<Duration>,
77    /// Set the claimed entries' last delivery time (`TIME`).
78    pub time: Option<SystemTime>,
79    /// Set the delivery count (`RETRYCOUNT`).
80    pub retry_count: Option<u64>,
81    /// Claim ids that are in the stream but not pending (`FORCE`).
82    pub force: bool,
83    /// Move the group's last delivered id (`LASTID`).
84    pub last_id: Option<String>,
85}
86
87/// Optional narrowing for the extended `XPENDING` form.
88#[derive(Debug, Clone, Default, PartialEq, Eq)]
89pub struct StreamPendingFilter {
90    /// Only entries owned by this consumer.
91    pub consumer: Option<String>,
92    /// Only entries idle for at least this long (`IDLE`).
93    pub min_idle: Option<Duration>,
94}
95
96/// One row of the extended `XPENDING` form.
97#[derive(Debug, Clone, PartialEq, Eq)]
98pub struct StreamPendingEntry {
99    /// Entry id.
100    pub id: String,
101    /// Consumer that owns the entry.
102    pub consumer: String,
103    /// Milliseconds since the entry was last delivered.
104    pub idle_millis: i64,
105    /// How many times the entry was delivered.
106    pub delivery_count: i64,
107}
108
109/// One `SCAN` page. A cursor of 0 means the scan is complete.
110#[derive(Debug, Clone, PartialEq, Eq)]
111pub struct ScanPage {
112    /// Cursor for the next call.
113    pub cursor: u64,
114    /// Keys on this page.
115    pub keys: Vec<String>,
116}
117
118/// One `HSCAN` page. A cursor of 0 means the scan is complete.
119#[derive(Debug, Clone, PartialEq, Eq)]
120pub struct HashScanPage {
121    /// Cursor for the next call.
122    pub cursor: u64,
123    /// Field and value pairs on this page.
124    pub fields: Vec<(String, String)>,
125}
126
127/// One stream entry.
128#[derive(Debug, Clone, PartialEq, Eq)]
129pub struct StreamEntry {
130    /// Entry id.
131    pub id: String,
132    /// Field and value pairs in the order they were added.
133    pub fields: Vec<(String, String)>,
134}
135
136/// Entries read from one stream key.
137#[derive(Debug, Clone, PartialEq, Eq)]
138pub struct StreamReadResult {
139    /// Stream key.
140    pub key: String,
141    /// Entries read from the key.
142    pub entries: Vec<StreamEntry>,
143}
144
145/// A Pub/Sub delivery: `message`, `pmessage`, or `smessage`.
146#[derive(Debug, Clone, PartialEq, Eq)]
147pub struct PubSubMessage {
148    /// `message`, `pmessage`, or `smessage`.
149    pub kind: String,
150    /// The matching pattern for `pmessage`.
151    pub pattern: Option<String>,
152    /// Channel the message was published to.
153    pub channel: String,
154    /// Raw payload.
155    pub payload: Vec<u8>,
156}