Skip to main content

kevy_embedded/
ops_zset_algebra.rs

1//! zset algebra facades (Redis 6.2): `ZINTERSTORE` / `ZUNIONSTORE` /
2//! `ZDIFFSTORE` / `ZINTERCARD`, plus the set-algebra `*STORE` forms
3//! (`SINTERSTORE` / `SUNIONSTORE` / `SDIFFSTORE`) that 1.6.0 shipped
4//! only as reads.
5//!
6//! Locking: sources are read under their own shard locks sequentially,
7//! then `dst` is written under its shard lock — same non-atomic
8//! window as [`Store::copy`] (a concurrent writer between the reads
9//! and the store can be observed). For a fully atomic combination use
10//! the same ops inside [`Store::atomic_all_shards`].
11//!
12//! AOF: the **effect** is logged (`DEL dst` + plain `ZADD`/`SADD` of
13//! the result), never the combination — deterministic on replay and
14//! replica-apply regardless of source state.
15
16use crate::KevyResult;
17
18use kevy_store::{ZAggregate, zdiff, zinter, zintercard, zunion};
19
20/// One source key's scored members (sets contribute score 1.0).
21type ScoredInput = Vec<(Vec<u8>, f64)>;
22
23use crate::store::ensure_writable;
24use crate::store::{Store, commit_write, store_err};
25
26impl Store {
27    fn gather_scored(&self, keys: &[&[u8]]) -> KevyResult<Vec<ScoredInput>> {
28        keys.iter()
29            .map(|k| self.wshard(k).store.zset_or_set_members(k).map_err(store_err))
30            .collect()
31    }
32
33    fn store_zset_result(&self, dst: &[u8], pairs: &[(Vec<u8>, f64)]) -> KevyResult<usize> {
34        ensure_writable(self)?;
35        let mut g = self.wshard(dst);
36        let n = g.store.zstore_result(dst, pairs);
37        commit_write(&mut g, &[b"DEL", dst])?;
38        if !pairs.is_empty() {
39            let score_strs: Vec<Vec<u8>> =
40                pairs.iter().map(|(_, s)| format!("{s}").into_bytes()).collect();
41            let mut argv: Vec<&[u8]> = Vec::with_capacity(2 + pairs.len() * 2);
42            argv.push(b"ZADD");
43            argv.push(dst);
44            for (i, (m, _)) in pairs.iter().enumerate() {
45                argv.push(&score_strs[i]);
46                argv.push(m);
47            }
48            commit_write(&mut g, &argv)?;
49        }
50        Ok(n)
51    }
52
53    /// `ZINTERSTORE dst keys… [WEIGHTS …] [AGGREGATE SUM|MIN|MAX]` —
54    /// returns the stored cardinality.
55    pub fn zinterstore(
56        &self,
57        dst: &[u8],
58        keys: &[&[u8]],
59        weights: Option<&[f64]>,
60        aggregate: ZAggregate,
61    ) -> KevyResult<usize> {
62        let inputs = self.gather_scored(keys)?;
63        self.store_zset_result(dst, &zinter(&inputs, weights, aggregate))
64    }
65
66    /// `ZUNIONSTORE dst keys… [WEIGHTS …] [AGGREGATE …]`.
67    pub fn zunionstore(
68        &self,
69        dst: &[u8],
70        keys: &[&[u8]],
71        weights: Option<&[f64]>,
72        aggregate: ZAggregate,
73    ) -> KevyResult<usize> {
74        let inputs = self.gather_scored(keys)?;
75        self.store_zset_result(dst, &zunion(&inputs, weights, aggregate))
76    }
77
78    /// `ZDIFFSTORE dst keys…` (no weights/aggregate — Redis 6.2).
79    pub fn zdiffstore(&self, dst: &[u8], keys: &[&[u8]]) -> KevyResult<usize> {
80        let inputs = self.gather_scored(keys)?;
81        self.store_zset_result(dst, &zdiff(&inputs))
82    }
83
84    /// `ZINTERCARD keys… [LIMIT n]` — `limit = 0` means unlimited.
85    pub fn zintercard(&self, keys: &[&[u8]], limit: usize) -> KevyResult<usize> {
86        let inputs = self.gather_scored(keys)?;
87        Ok(zintercard(&inputs, limit))
88    }
89
90    /// `SINTERSTORE dst keys…` — set-algebra store form (members only).
91    pub fn sinterstore(&self, dst: &[u8], keys: &[&[u8]]) -> KevyResult<usize> {
92        let members = self.sinter(keys)?;
93        self.store_set_result(dst, &members)
94    }
95
96    /// `SUNIONSTORE dst keys…`.
97    pub fn sunionstore(&self, dst: &[u8], keys: &[&[u8]]) -> KevyResult<usize> {
98        let members = self.sunion(keys)?;
99        self.store_set_result(dst, &members)
100    }
101
102    /// `SDIFFSTORE dst keys…`.
103    pub fn sdiffstore(&self, dst: &[u8], keys: &[&[u8]]) -> KevyResult<usize> {
104        let members = self.sdiff(keys)?;
105        self.store_set_result(dst, &members)
106    }
107
108    fn store_set_result(&self, dst: &[u8], members: &[Vec<u8>]) -> KevyResult<usize> {
109        ensure_writable(self)?;
110        let mut g = self.wshard(dst);
111        g.store.del(&[dst]);
112        commit_write(&mut g, &[b"DEL", dst])?;
113        if members.is_empty() {
114            return Ok(0);
115        }
116        let member_refs: Vec<&[u8]> = members.iter().map(Vec::as_slice).collect();
117        let n = g.store.sadd(dst, &member_refs).map_err(store_err)?;
118        let mut argv: Vec<&[u8]> = Vec::with_capacity(2 + members.len());
119        argv.push(b"SADD");
120        argv.push(dst);
121        argv.extend(members.iter().map(Vec::as_slice));
122        commit_write(&mut g, &argv)?;
123        Ok(n)
124    }
125}