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}