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}