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(
82        &self,
83        key: &[u8],
84        start: i64,
85        stop: i64,
86    ) -> KevyResult<Vec<(Vec<u8>, f64)>> {
87        let mut all = self
88            .wshard(key)
89            .store
90            .zrange(key, 0, -1)
91            .map_err(store_err)?;
92        all.reverse();
93        let n = all.len() as i64;
94        if n == 0 {
95            return Ok(Vec::new());
96        }
97        let clamp = |x: i64| -> usize {
98            let v = if x < 0 { (n + x).max(0) } else { x.min(n - 1) };
99            v as usize
100        };
101        let s = clamp(start);
102        let e = clamp(stop);
103        if s > e {
104            return Ok(Vec::new());
105        }
106        Ok(all.into_iter().skip(s).take(e - s + 1).collect())
107    }
108
109    /// `ZRANGEBYSCORE` — score-range read. `min` / `max` are
110    /// inclusive; pass `f64::NEG_INFINITY` / `f64::INFINITY` for open
111    /// bounds. Returns `(member, score)` pairs in ascending score
112    /// order. Exclusive bounds are `ZRANGEBYSCORE (` syntax in Redis;
113    /// expose via the dedicated [`Self::zrange_by_score_excl`].
114    pub fn zrange_by_score(
115        &self,
116        key: &[u8],
117        min: f64,
118        max: f64,
119    ) -> KevyResult<Vec<(Vec<u8>, f64)>> {
120        self.wshard(key)
121            .store
122            .zrange_by_score(
123                key,
124                ScoreBound { value: min, exclusive: false },
125                ScoreBound { value: max, exclusive: false },
126            )
127            .map_err(store_err)
128    }
129
130    /// Same as [`Self::zrange_by_score`] but with explicit
131    /// inclusive/exclusive control on each bound (`(min` / `(max` in
132    /// Redis syntax).
133    pub fn zrange_by_score_excl(
134        &self,
135        key: &[u8],
136        min: ScoreBound,
137        max: ScoreBound,
138    ) -> KevyResult<Vec<(Vec<u8>, f64)>> {
139        self.wshard(key)
140            .store
141            .zrange_by_score(key, min, max)
142            .map_err(store_err)
143    }
144
145    /// `ZRANGEBYSCORE key min max LIMIT offset count` — score-range
146    /// read with pagination (closes the embedded LIMIT gap —
147    /// the server parser always had it).
148    pub fn zrange_by_score_limit(
149        &self,
150        key: &[u8],
151        min: f64,
152        max: f64,
153        offset: usize,
154        count: usize,
155    ) -> KevyResult<Vec<(Vec<u8>, f64)>> {
156        let all = self.zrange_by_score(key, min, max)?;
157        Ok(all.into_iter().skip(offset).take(count).collect())
158    }
159
160    /// `ZREVRANGEBYSCORE key max min LIMIT offset count` — descending
161    /// score-range read with pagination.
162    pub fn zrevrange_by_score_limit(
163        &self,
164        key: &[u8],
165        max: f64,
166        min: f64,
167        offset: usize,
168        count: usize,
169    ) -> KevyResult<Vec<(Vec<u8>, f64)>> {
170        let mut all = self.zrange_by_score(key, min, max)?;
171        all.reverse();
172        Ok(all.into_iter().skip(offset).take(count).collect())
173    }
174
175    /// `zpopmin_below` — pop up to `count` lowest members with
176    /// score strictly `< below` (delayed-job "pop what's due").
177    /// AOF logs the effect (`ZREM` of the popped members).
178    pub fn zpopmin_below(
179        &self,
180        key: &[u8],
181        below: f64,
182        count: usize,
183    ) -> KevyResult<Vec<(Vec<u8>, f64)>> {
184        ensure_writable(self)?;
185        let mut g = self.wshard(key);
186        let items = g.store.zpopmin_below(key, below, count).map_err(store_err)?;
187        if !items.is_empty() {
188            let mut argv: Vec<&[u8]> = Vec::with_capacity(2 + items.len());
189            argv.push(b"ZREM");
190            argv.push(key);
191            argv.extend(items.iter().map(|(m, _)| m.as_slice()));
192            commit_write(&mut g, &argv)?;
193        }
194        Ok(items)
195    }
196
197    /// `ZINCRBY key delta member` — atomic float increment of a member's
198    /// score. Returns the post-increment score.
199    pub fn zincrby(&self, key: &[u8], delta: f64, member: &[u8]) -> KevyResult<f64> {
200        ensure_writable(self)?;
201        let mut g = self.wshard(key);
202        let new_score = g.store.zincrby(key, delta, member).map_err(store_err)?;
203        let delta_str = format!("{delta}");
204        commit_write(&mut g, &[b"ZINCRBY", key, delta_str.as_bytes(), member])?;
205        Ok(new_score)
206    }
207
208    // ---- list slice + index ops --------------------------------------
209
210    /// `LRANGE key start stop` — list slice. Negative indices count
211    /// from the tail. Empty when absent.
212    pub fn lrange(&self, key: &[u8], start: i64, stop: i64) -> KevyResult<Vec<Vec<u8>>> {
213        self.wshard(key).store.lrange(key, start, stop).map_err(store_err)
214    }
215
216    /// `LINDEX key idx` — element at index `idx`; `None` out of range.
217    pub fn lindex(&self, key: &[u8], idx: i64) -> KevyResult<Option<Vec<u8>>> {
218        self.wshard(key).store.lindex(key, idx).map_err(store_err)
219    }
220
221    /// `LREM key count value` — remove up to `|count|` occurrences of
222    /// `value`. `count > 0` from head, `count < 0` from tail,
223    /// `count == 0` all. Returns the count actually removed.
224    pub fn lrem(&self, key: &[u8], count: i64, value: &[u8]) -> KevyResult<usize> {
225        ensure_writable(self)?;
226        let mut g = self.wshard(key);
227        let removed = g.store.lrem(key, count, value).map_err(store_err)?;
228        if removed > 0 {
229            let count_str = format!("{count}");
230            commit_write(&mut g, &[b"LREM", key, count_str.as_bytes(), value])?;
231        }
232        Ok(removed)
233    }
234
235    // ---- string single-call atomic patterns --------------------------
236
237    /// `GETSET key new` — set `key` to `new`, return the previous
238    /// value (or `None` when `key` was absent).
239    pub fn getset(&self, key: &[u8], new: &[u8]) -> KevyResult<Option<Vec<u8>>> {
240        ensure_writable(self)?;
241        let mut g = self.wshard(key);
242        let prev = g.store.getset(key, new.to_vec()).map_err(store_err)?;
243        commit_write(&mut g, &[b"SET", key, new])?;
244        Ok(prev)
245    }
246
247    /// `GETDEL key` — delete `key`, return the previous value
248    /// (`None` when `key` was absent).
249    pub fn getdel(&self, key: &[u8]) -> KevyResult<Option<Vec<u8>>> {
250        ensure_writable(self)?;
251        let mut g = self.wshard(key);
252        let prev = g.store.getdel(key).map_err(store_err)?;
253        if prev.is_some() {
254            commit_write(&mut g, &[b"DEL", key])?;
255        }
256        Ok(prev)
257    }
258}