Skip to main content

kevy_embedded/
ops_p2.rs

1//! Hash field reads, sorted-set range queries, list slice access, and
2//! the atomic single-call string helpers `getset` / `getdel`.
3//!
4//! Every method is a thin facade over the corresponding
5//! `kevy_store::Store` method, with `commit_write` AOF logging on the
6//! write paths.
7
8use crate::KevyResult;
9
10use kevy_store::ScoreBound;
11
12use crate::store::ensure_writable;
13use crate::store::{Store, commit_write, store_err};
14
15impl Store {
16    // ---- hash mass-getters --------------------------------------------
17
18    /// `HGETALL key` — every `(field, value)` pair in `key`'s hash, in
19    /// arbitrary order. Empty when `key` is absent. Errors on wrong type.
20    pub fn hgetall(&self, key: &[u8]) -> KevyResult<Vec<(Vec<u8>, Vec<u8>)>> {
21        let flat = self.wshard(key).store.hgetall(key).map_err(store_err)?;
22        // kevy-store returns [f0, v0, f1, v1, ...] — pair them up.
23        let mut out = Vec::with_capacity(flat.len() / 2);
24        let mut it = flat.into_iter();
25        while let (Some(f), Some(v)) = (it.next(), it.next()) {
26            out.push((f, v));
27        }
28        Ok(out)
29    }
30
31    /// `HEXISTS key field` — `true` when `field` is present.
32    pub fn hexists(&self, key: &[u8], field: &[u8]) -> KevyResult<bool> {
33        self.wshard(key).store.hexists(key, field).map_err(store_err)
34    }
35
36    /// `HLEN key` — number of fields; 0 when absent.
37    pub fn hlen(&self, key: &[u8]) -> KevyResult<usize> {
38        self.wshard(key).store.hlen(key).map_err(store_err)
39    }
40
41    /// `HKEYS key` — every field name in `key`'s hash.
42    pub fn hkeys(&self, key: &[u8]) -> KevyResult<Vec<Vec<u8>>> {
43        self.wshard(key).store.hkeys(key).map_err(store_err)
44    }
45
46    /// `HVALS key` — every value in `key`'s hash.
47    pub fn hvals(&self, key: &[u8]) -> KevyResult<Vec<Vec<u8>>> {
48        self.wshard(key).store.hvals(key).map_err(store_err)
49    }
50
51    /// `HMGET key field [field ...]` — read multiple fields in one
52    /// call. `None` per requested field that is absent.
53    pub fn hmget(&self, key: &[u8], fields: &[&[u8]]) -> KevyResult<Vec<Option<Vec<u8>>>> {
54        self.wshard(key).store.hmget(key, fields).map_err(store_err)
55    }
56
57    /// `HINCRBY key field delta` — atomic integer increment of a hash
58    /// field. Returns the post-increment value.
59    pub fn hincrby(&self, key: &[u8], field: &[u8], delta: i64) -> KevyResult<i64> {
60        ensure_writable(self)?;
61        let mut g = self.wshard(key);
62        let new_val = g.store.hincrby(key, field, delta).map_err(store_err)?;
63        let delta_str = format!("{delta}");
64        commit_write(&mut g, &[b"HINCRBY", key, field, delta_str.as_bytes()])?;
65        Ok(new_val)
66    }
67
68    // ---- zset mass-readers + atomic incr -----------------------------
69
70    /// `ZRANGE key start stop WITHSCORES` — members in ascending score
71    /// order between rank `start..=stop` (Redis-style inclusive
72    /// indexing; negatives count from the tail). Returns `(member,
73    /// score)` pairs.
74    pub fn zrange(&self, key: &[u8], start: i64, stop: i64) -> KevyResult<Vec<(Vec<u8>, f64)>> {
75        self.wshard(key).store.zrange(key, start, stop).map_err(store_err)
76    }
77
78    /// `ZREVRANGE key start stop WITHSCORES` — `zrange` with the order
79    /// reversed (highest score first). The `start..=stop` indexing is
80    /// against the reversed list, matching Redis semantics.
81    pub fn zrevrange(&self, key: &[u8], start: i64, stop: i64) -> KevyResult<Vec<(Vec<u8>, f64)>> {
82        self.wshard(key).store.zrevrange(key, start, stop).map_err(store_err)
83    }
84
85    /// `ZRANGEBYSCORE` — score-range read. `min` / `max` are
86    /// inclusive; pass `f64::NEG_INFINITY` / `f64::INFINITY` for open
87    /// bounds. Returns `(member, score)` pairs in ascending score
88    /// order. Exclusive bounds are `ZRANGEBYSCORE (` syntax in Redis;
89    /// expose via the dedicated [`Self::zrange_by_score_excl`].
90    pub fn zrange_by_score(
91        &self,
92        key: &[u8],
93        min: f64,
94        max: f64,
95    ) -> KevyResult<Vec<(Vec<u8>, f64)>> {
96        self.wshard(key)
97            .store
98            .zrange_by_score(
99                key,
100                ScoreBound { value: min, exclusive: false },
101                ScoreBound { value: max, exclusive: false },
102            )
103            .map_err(store_err)
104    }
105
106    /// Same as [`Self::zrange_by_score`] but with explicit
107    /// inclusive/exclusive control on each bound (`(min` / `(max` in
108    /// Redis syntax).
109    pub fn zrange_by_score_excl(
110        &self,
111        key: &[u8],
112        min: ScoreBound,
113        max: ScoreBound,
114    ) -> KevyResult<Vec<(Vec<u8>, f64)>> {
115        self.wshard(key).store.zrange_by_score(key, min, max).map_err(store_err)
116    }
117
118    /// `ZRANGEBYSCORE key min max LIMIT offset count` — score-range
119    /// read with pagination (closes the embedded LIMIT gap —
120    /// the server parser always had it).
121    pub fn zrange_by_score_limit(
122        &self,
123        key: &[u8],
124        min: f64,
125        max: f64,
126        offset: usize,
127        count: usize,
128    ) -> KevyResult<Vec<(Vec<u8>, f64)>> {
129        let all = self.zrange_by_score(key, min, max)?;
130        Ok(all.into_iter().skip(offset).take(count).collect())
131    }
132
133    /// `ZREVRANGEBYSCORE key max min LIMIT offset count` — descending
134    /// score-range read with pagination.
135    pub fn zrevrange_by_score_limit(
136        &self,
137        key: &[u8],
138        max: f64,
139        min: f64,
140        offset: usize,
141        count: usize,
142    ) -> KevyResult<Vec<(Vec<u8>, f64)>> {
143        let mut all = self.zrange_by_score(key, min, max)?;
144        all.reverse();
145        Ok(all.into_iter().skip(offset).take(count).collect())
146    }
147
148    /// `zpopmin_below` — pop up to `count` lowest members with
149    /// score strictly `< below` (delayed-job "pop what's due").
150    /// AOF logs the effect (`ZREM` of the popped members).
151    pub fn zpopmin_below(
152        &self,
153        key: &[u8],
154        below: f64,
155        count: usize,
156    ) -> KevyResult<Vec<(Vec<u8>, f64)>> {
157        ensure_writable(self)?;
158        let mut g = self.wshard(key);
159        let items = g.store.zpopmin_below(key, below, count).map_err(store_err)?;
160        if !items.is_empty() {
161            let mut argv: Vec<&[u8]> = Vec::with_capacity(2 + items.len());
162            argv.push(b"ZREM");
163            argv.push(key);
164            argv.extend(items.iter().map(|(m, _)| m.as_slice()));
165            commit_write(&mut g, &argv)?;
166        }
167        Ok(items)
168    }
169
170    /// `ZINCRBY key delta member` — atomic float increment of a member's
171    /// score. Returns the post-increment score.
172    pub fn zincrby(&self, key: &[u8], delta: f64, member: &[u8]) -> KevyResult<f64> {
173        ensure_writable(self)?;
174        let mut g = self.wshard(key);
175        let new_score = g.store.zincrby(key, delta, member).map_err(store_err)?;
176        let delta_str = format!("{delta}");
177        commit_write(&mut g, &[b"ZINCRBY", key, delta_str.as_bytes(), member])?;
178        Ok(new_score)
179    }
180
181    // ---- list slice + index ops --------------------------------------
182
183    /// `LRANGE key start stop` — list slice. Negative indices count
184    /// from the tail. Empty when absent.
185    pub fn lrange(&self, key: &[u8], start: i64, stop: i64) -> KevyResult<Vec<Vec<u8>>> {
186        self.wshard(key).store.lrange(key, start, stop).map_err(store_err)
187    }
188
189    /// `LINDEX key idx` — element at index `idx`; `None` out of range.
190    pub fn lindex(&self, key: &[u8], idx: i64) -> KevyResult<Option<Vec<u8>>> {
191        self.wshard(key).store.lindex(key, idx).map_err(store_err)
192    }
193
194    /// `LREM key count value` — remove up to `|count|` occurrences of
195    /// `value`. `count > 0` from head, `count < 0` from tail,
196    /// `count == 0` all. Returns the count actually removed.
197    pub fn lrem(&self, key: &[u8], count: i64, value: &[u8]) -> KevyResult<usize> {
198        ensure_writable(self)?;
199        let mut g = self.wshard(key);
200        let removed = g.store.lrem(key, count, value).map_err(store_err)?;
201        if removed > 0 {
202            let count_str = format!("{count}");
203            commit_write(&mut g, &[b"LREM", key, count_str.as_bytes(), value])?;
204        }
205        Ok(removed)
206    }
207
208    // ---- string single-call atomic patterns --------------------------
209
210    /// `GETSET key new` — set `key` to `new`, return the previous
211    /// value (or `None` when `key` was absent).
212    pub fn getset(&self, key: &[u8], new: &[u8]) -> KevyResult<Option<Vec<u8>>> {
213        ensure_writable(self)?;
214        let mut g = self.wshard(key);
215        let prev = g.store.getset(key, new.to_vec()).map_err(store_err)?;
216        commit_write(&mut g, &[b"SET", key, new])?;
217        Ok(prev)
218    }
219
220    /// `GETDEL key` — delete `key`, return the previous value
221    /// (`None` when `key` was absent).
222    pub fn getdel(&self, key: &[u8]) -> KevyResult<Option<Vec<u8>>> {
223        ensure_writable(self)?;
224        let mut g = self.wshard(key);
225        let prev = g.store.getdel(key).map_err(store_err)?;
226        if prev.is_some() {
227            commit_write(&mut g, &[b"DEL", key])?;
228        }
229        Ok(prev)
230    }
231}