1use async_channel as mpsc;
2pub use async_kvdb::*;
3use async_memorydb::MenoryDb;
4use core::sync::atomic::{AtomicBool, Ordering::Relaxed};
5use std::collections::HashMap;
6use std::io::Write;
7use std::path::{Path, PathBuf};
8use std::time::{Duration, Instant};
9use std::{fs, io, thread};
10
11fn key2file(key: &Key) -> String {
12 form_urlencoded::Serializer::new(String::new()).append_key_only(key).finish()
13}
14fn file2key(name: &str) -> Option<Key> {
15 let mut parse = form_urlencoded::parse(name.as_bytes());
16 let data = parse.next()?;
17 Some(data.0.into())
18}
19
20fn load_item(path: &Path, key: &Key) -> Option<Value> {
21 let name = key2file(key);
22 let path = path.join(name);
23 let value: Value = fs::read(path).ok()?.into();
24 Some(value)
25}
26fn load_all(path: &Path) -> HashMap<Key, Value> {
27 let Ok(entrys) = fs::read_dir(path) else {
28 return HashMap::new();
29 };
30 entrys
31 .filter_map(|entry| {
32 let file = entry.ok()?.path();
33 if !file.is_file() {
34 return None;
35 }
36 let name = file.file_name()?;
37 let name = name.to_string_lossy().into_owned();
38 let key = file2key(&name)?;
39 let value = fs::read(file).ok()?;
40 Some((key, value.into()))
41 })
42 .collect()
43}
44
45fn fsdb_exec(path: &Path, op: DbOp) -> io::Result<()> {
46 match op {
47 DbOp::Get { key, ch } => {
48 let value = load_item(path, &key);
49 let _ = ch.send(value);
50 }
51 DbOp::GetMany { keys, ch } => {
52 let values = keys
53 .into_iter()
54 .filter_map(|key| {
55 let value = load_item(path, &key)?;
56 Some((key, value))
57 })
58 .collect();
59 let _ = ch.send(values);
60 }
61 DbOp::GetAll { ch } => {
62 let value = load_all(path);
63 let _ = ch.send(value);
64 }
65 DbOp::Insert { key, value } => {
66 let name = key2file(&key);
67 let path = path.join(name);
68 let mut f = fs::File::create(path)?;
69 let _ = f.write_all(&value);
70 #[cfg(feature = "auto_sync")]
71 {
72 let _ = f.sync_data();
73 }
74 }
75 DbOp::InsertMany { data } => {
76 for (key, value) in data {
77 let name = key2file(&key);
78 let path = path.join(name);
79 let mut f = fs::File::create(path)?;
80 let _ = f.write_all(&value);
81 #[cfg(feature = "auto_sync")]
82 {
83 let _ = f.sync_data();
84 }
85 }
86 }
87 DbOp::Delete { key } => {
88 let name = key2file(&key);
89 let path = path.join(name);
90 let _ = fs::remove_file(path);
91 }
92 DbOp::DeleteMany { keys } => {
93 for entry in fs::read_dir(path)? {
94 let file = entry?.path();
95 if file.is_file() {
96 if let Some(name) = file.file_name() {
97 let name = name.to_string_lossy();
98 if keys.iter().any(|k| key2file(k) == name) {
99 let _ = fs::remove_file(file);
100 }
101 }
102 }
103 }
104 }
105 DbOp::DeleteAll => {
106 for entry in fs::read_dir(path)? {
107 let file = entry?.path();
108 if file.is_file() {
109 let _ = fs::remove_file(file);
110 }
111 }
112 }
113 };
114 Ok(())
115}
116
117pub struct FileDb {
118 cached_all: AtomicBool,
119 mem: MenoryDb,
120 write_ch: mpsc::Sender<DbOp>,
121}
122impl FileDb {
123 pub fn new(path: String, interval_ms: u64) -> io::Result<Self> {
127 let interval_ms_check = (interval_ms / 10).clamp(20, 100);
128 let thread_name = format!("db-{}", &path[path.len().saturating_sub(12)..]);
129 let path = PathBuf::from(path);
130 let mem = HashMap::new();
131 fs::create_dir_all(&path)?;
132 let (tx, receiver) = mpsc::unbounded();
133 thread::Builder::new().name(thread_name).spawn(move || {
134 let mut op_merger = DbOpMerger::new();
135 loop {
136 let start_time = Instant::now();
137 if let Ok(op) = receiver.recv_blocking() {
138 op_merger.merge(op);
139 }
140 let remain_time = interval_ms.saturating_sub(start_time.elapsed().as_millis() as u64);
141 for _ in 0..remain_time / interval_ms_check {
142 while let Ok(op) = receiver.try_recv() {
143 op_merger.merge(op);
144 }
145 if !op_merger.need_read() {
146 std::thread::sleep(Duration::from_millis(interval_ms_check))
147 } else {
148 break;
149 }
150 }
151 while let Ok(op) = receiver.try_recv() {
152 op_merger.merge(op);
153 }
154 let ops = op_merger.into_ops();
155 op_merger = DbOpMerger::new();
156 if ops.clear {
157 let _ = fsdb_exec(&path, DbOp::DeleteAll);
158 }
159 if !ops.insert.is_empty() {
161 let _ = fsdb_exec(&path, DbOp::InsertMany { data: ops.insert });
162 }
163 if !ops.delete.is_empty() {
164 let _ = fsdb_exec(&path, DbOp::DeleteMany { keys: ops.delete });
165 }
166 for (key, chs) in ops.get_one {
168 let value = load_item(&path, &key);
169 for ch in chs {
170 let _ = ch.send(value.clone());
171 }
172 }
173 for op in ops.get_many {
174 let _ = fsdb_exec(&path, op);
175 }
176 for ch in ops.get_all {
177 let value = load_all(&path);
178 let _ = ch.send(value.clone());
179 }
180 }
181 })?;
182 Ok(Self {
183 cached_all: AtomicBool::new(false),
184 mem: MenoryDb::new(mem),
185 write_ch: tx,
186 })
187 }
188 pub async fn clear_cache(&self) {
189 self.cached_all.store(false, Relaxed);
190 self.mem.delete_all().await;
191 }
192 pub async fn load_all(&self) {
193 let (tx, rx) = oneshot::async_channel();
194 let _ = self.write_ch.send(DbOp::GetAll { ch: tx }).await;
195 let Ok(data) = rx.await else { return };
196 self.mem.set_many(data).await;
197 self.cached_all.store(true, Relaxed);
198 }
199}
200
201#[async_trait]
202impl Kvdb for FileDb {
203 async fn scan_keys(&self, filter: &Filter) -> Vec<Key> {
204 if !self.cached_all.load(Relaxed) {
205 self.load_all().await;
206 }
207 self.mem.scan_keys(filter).await
208 }
209 async fn get(&self, key: Key) -> Option<Value> {
210 let ret = self.mem.get(key.clone()).await;
211 if ret.is_some() {
212 return ret;
213 }
214 let (tx, rx) = oneshot::async_channel();
215 let _ = self.write_ch.send(DbOp::Get { key: key.clone(), ch: tx }).await;
216 let value = rx.await.ok()?;
217 if let Some(value) = &value {
218 self.mem.set(key, value.clone()).await;
219 } else {
220 self.mem.delete(key).await;
221 }
222 value
223 }
224 async fn get_many(&self, keys: Vec<Key>) -> HashMap<Key, Value> {
225 if !self.cached_all.load(Relaxed) {
226 self.load_all().await;
227 }
228 self.mem.get_many(keys).await
229 }
230 async fn set(&self, key: Key, value: Value) {
231 self.mem.set(key.clone(), value.clone()).await;
232 let _ = self.write_ch.send(DbOp::Insert { key, value }).await;
233 }
234 async fn set_many(&self, data: HashMap<Key, Value>) {
235 self.mem.set_many(data.clone()).await;
236 let _ = self.write_ch.send(DbOp::InsertMany { data }).await;
237 }
238 async fn delete(&self, key: Key) {
239 self.mem.delete(key.clone()).await;
240 let _ = self.write_ch.send(DbOp::Delete { key }).await;
241 }
242 async fn delete_many(&self, keys: Vec<Key>) {
243 self.mem.delete_many(keys.clone()).await;
244 let _ = self.write_ch.send(DbOp::DeleteMany { keys }).await;
245 }
246 async fn delete_all(&self) {
247 self.mem.delete_all().await;
248 let _ = self.write_ch.send(DbOp::DeleteAll).await;
249 }
250}