kevy_embedded/
ops_zset_algebra.rs1use crate::KevyResult;
17
18use kevy_store::{ZAggregate, zdiff, zinter, zintercard, zunion};
19
20type 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 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 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 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 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 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 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 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}