Skip to main content

async_kvdb/
lib.rs

1#[cfg(feature = "json")]
2mod json;
3#[cfg(feature = "proto")]
4mod proto;
5
6pub use async_trait::async_trait;
7use bytes::Bytes;
8use smol_str::SmolStr;
9use std::collections::HashMap;
10
11#[cfg(feature = "json")]
12pub use json::KvdbJsonExt;
13#[cfg(feature = "proto")]
14pub use proto::KvdbProtoExt;
15
16pub type Key = SmolStr;
17pub type Value = Bytes;
18pub type KeyValue = (Key, Value);
19pub type Filter = dyn Fn(&Key) -> bool + Send + Sync;
20
21pub enum DbOp {
22    Get {
23        key: Key,
24        ch: oneshot::Sender<Option<Value>>,
25    },
26    GetMany {
27        keys: Vec<Key>,
28        ch: oneshot::Sender<HashMap<Key, Value>>,
29    },
30    GetAll {
31        ch: oneshot::Sender<HashMap<Key, Value>>,
32    },
33    Insert {
34        key: Key,
35        value: Value,
36    },
37    InsertMany {
38        data: HashMap<Key, Value>,
39    },
40    Delete {
41        key: Key,
42    },
43    DeleteMany {
44        keys: Vec<Key>,
45    },
46    DeleteAll,
47}
48
49#[derive(Default)]
50pub struct DbOps {
51    pub clear: bool,
52    pub insert: HashMap<Key, Value>,
53    pub delete: Vec<Key>,
54    pub get_one: HashMap<Key, Vec<oneshot::Sender<Option<Value>>>>,
55    pub get_many: Vec<DbOp>,
56    pub get_all: Vec<oneshot::Sender<HashMap<Key, Value>>>,
57}
58
59#[derive(Default)]
60pub struct DbOpMerger {
61    need_clear: bool,
62    ops: HashMap<Key, Option<Value>>,
63    get_one: HashMap<Key, Vec<oneshot::Sender<Option<Value>>>>,
64    get_many: Vec<DbOp>,
65    get_all: Vec<oneshot::Sender<HashMap<Key, Value>>>,
66}
67impl DbOpMerger {
68    pub fn new() -> Self {
69        Self {
70            need_clear: false,
71            ops: HashMap::new(),
72            get_one: HashMap::new(),
73            get_many: Vec::new(),
74            get_all: Vec::new(),
75        }
76    }
77    pub fn need_read(&self) -> bool {
78        !(self.get_one.is_empty() && self.get_many.is_empty() && self.get_all.is_empty())
79    }
80    pub fn need_write(&self) -> bool {
81        self.need_clear || !self.ops.is_empty()
82    }
83    pub fn is_empty(&self) -> bool {
84        self.get_one.is_empty() && self.get_many.is_empty() && self.get_all.is_empty() && !self.need_clear && self.ops.is_empty()
85    }
86    pub fn merge(&mut self, op: DbOp) {
87        match op {
88            DbOp::Get { key, ch } => {
89                let chs = self.get_one.entry(key).or_default();
90                chs.push(ch);
91            }
92            DbOp::GetMany { .. } => {
93                self.get_many.push(op);
94            }
95            DbOp::GetAll { ch } => {
96                self.get_all.push(ch);
97            }
98            DbOp::Insert { key, value } => {
99                self.ops.insert(key, Some(value));
100            }
101            DbOp::InsertMany { data } => {
102                for (key, value) in data {
103                    self.ops.insert(key, Some(value));
104                }
105            }
106            DbOp::Delete { key } => {
107                self.ops.insert(key, None);
108            }
109            DbOp::DeleteMany { keys } => {
110                for key in keys {
111                    self.ops.insert(key, None);
112                }
113            }
114            DbOp::DeleteAll => {
115                self.need_clear = true;
116                self.ops.clear();
117            }
118        }
119    }
120    pub fn into_ops(self) -> DbOps {
121        let mut ops = DbOps::default();
122        for (k, v) in self.ops {
123            if let Some(v) = v {
124                ops.insert.insert(k, v);
125            } else {
126                ops.delete.push(k);
127            }
128        }
129        ops.clear = self.need_clear;
130        ops.get_one = self.get_one;
131        ops.get_many = self.get_many;
132        ops.get_all = self.get_all;
133        ops
134    }
135}
136
137#[async_trait]
138pub trait Kvdb {
139    async fn scan_keys(&self, filter: &Filter) -> Vec<Key>;
140    async fn get(&self, key: Key) -> Option<Value> {
141        let keys = Vec::from([key]);
142        self.get_many(keys).await.into_values().next()
143    }
144    async fn get_many(&self, keys: Vec<Key>) -> HashMap<Key, Value>;
145    async fn set(&self, key: Key, value: Value) {
146        let data = HashMap::from([(key, value)]);
147        self.set_many(data).await
148    }
149    async fn set_many(&self, data: HashMap<Key, Value>);
150    async fn delete(&self, key: Key) {
151        let keys = Vec::from([key]);
152        self.delete_many(keys).await
153    }
154    async fn delete_many(&self, keys: Vec<Key>);
155    async fn delete_all(&self);
156}
157
158#[test]
159fn empty() {
160    let op_merger = DbOpMerger::default();
161    assert!(op_merger.is_empty());
162}