1use core::fmt;
19use std::sync::atomic::{AtomicBool, AtomicIsize, Ordering};
20use std::sync::Arc;
21
22use anyhow::{anyhow, Context, Result};
23use async_trait::async_trait;
24use serde::de::DeserializeOwned;
25use serde::{Deserialize, Serialize};
26use serde_json::Value;
27use tokio::runtime::Handle;
28use tokio::sync::mpsc;
29use tokio::sync::oneshot;
30use tokio::task::spawn_blocking;
31
32use redb::{
33 Database, ReadableDatabase, ReadableTable, ReadableTableMetadata, TableDefinition,
34 WriteTransaction,
35};
36
37use crate::storage::{AsyncIterator, IterItem, Key, List, Map, StorageDB};
38#[cfg(feature = "ttl")]
39use crate::timestamp_millis;
40use crate::{StorageList, StorageMap, TimestampMillis};
41
42const KV_TABLE: TableDefinition<&[u8], &[u8]> = TableDefinition::new("__kv_table");
48const MAP_TABLE: TableDefinition<&[u8], &[u8]> = TableDefinition::new("__map_table");
49const LIST_TABLE: TableDefinition<&[u8], &[u8]> = TableDefinition::new("__list_table");
50const EXPIRE_KEYS_TABLE: TableDefinition<&[u8], &[u8]> = TableDefinition::new("__expire_key_table");
51const KEY_EXPIRE_TABLE: TableDefinition<&[u8], &[u8]> = TableDefinition::new("__key_expire_table");
52
53#[allow(dead_code)]
55const SEPARATOR: &[u8] = b"@";
56const MAP_KEY_SEPARATOR: &[u8] = b"@__item@";
58const MAP_KEY_COUNT_SUFFIX: &[u8] = b"@__count@";
60const LIST_KEY_COUNT_SUFFIX: &[u8] = b"@__count@";
62const LIST_KEY_CONTENT_SUFFIX: &[u8] = b"@__content@";
64
65const MAP_NAME_PREFIX: &[u8] = b"__rmqtt_map@";
67const LIST_NAME_PREFIX: &[u8] = b"__rmqtt_list@";
68const KEY_PREFIX: &[u8] = b"__rmqtt@";
69#[allow(dead_code)]
70const KEY_PREFIX_LEN: &[u8] = b"__rmqtt_len@";
71const COUNTER_PREFIX: &[u8] = b"__counter@";
72
73const CHANNEL_CAPACITY: usize = 10_000;
75
76type IterResultSender = oneshot::Sender<Result<Vec<(Vec<u8>, Vec<u8>)>>>;
82
83type BatchInsertItem = (Vec<u8>, Vec<u8>, oneshot::Sender<Result<()>>);
85
86#[allow(dead_code)]
92fn make_map_prefix(name: &[u8]) -> Vec<u8> {
93 let mut v = Vec::with_capacity(MAP_NAME_PREFIX.len() + name.len() + 1);
94 v.extend_from_slice(MAP_NAME_PREFIX);
95 v.extend_from_slice(name);
96 v.extend_from_slice(SEPARATOR);
97 v
98}
99
100fn make_map_item_prefix(name: &[u8]) -> Vec<u8> {
102 let mut v = Vec::with_capacity(MAP_NAME_PREFIX.len() + name.len() + MAP_KEY_SEPARATOR.len());
103 v.extend_from_slice(MAP_NAME_PREFIX);
104 v.extend_from_slice(name);
105 v.extend_from_slice(MAP_KEY_SEPARATOR);
106 v
107}
108
109fn make_map_item_key(name: &[u8], key: &[u8]) -> Vec<u8> {
111 let mut v = make_map_item_prefix(name);
112 v.extend_from_slice(key);
113 v
114}
115
116fn make_map_count_key(name: &[u8]) -> Vec<u8> {
118 let mut v = Vec::with_capacity(MAP_NAME_PREFIX.len() + name.len() + MAP_KEY_COUNT_SUFFIX.len());
119 v.extend_from_slice(MAP_NAME_PREFIX);
120 v.extend_from_slice(name);
121 v.extend_from_slice(MAP_KEY_COUNT_SUFFIX);
122 v
123}
124
125fn is_map_count_key(key: &[u8]) -> bool {
127 key.starts_with(MAP_NAME_PREFIX) && key.ends_with(MAP_KEY_COUNT_SUFFIX)
128}
129
130fn map_count_key_to_name(key: &[u8]) -> &[u8] {
132 let start = MAP_NAME_PREFIX.len();
133 let end = key.len() - MAP_KEY_COUNT_SUFFIX.len();
134 &key[start..end]
135}
136
137#[allow(dead_code)]
139fn map_item_key_to_name(key: &[u8]) -> Option<&[u8]> {
140 if let Some(pos) = key
141 .windows(MAP_KEY_SEPARATOR.len())
142 .position(|w| w == MAP_KEY_SEPARATOR)
143 {
144 if key.starts_with(MAP_NAME_PREFIX) {
145 return Some(&key[MAP_NAME_PREFIX.len()..pos]);
146 }
147 }
148 None
149}
150
151#[allow(dead_code)]
153fn make_list_prefix(name: &[u8]) -> Vec<u8> {
154 let mut v = Vec::with_capacity(LIST_NAME_PREFIX.len() + name.len());
155 v.extend_from_slice(LIST_NAME_PREFIX);
156 v.extend_from_slice(name);
157 v
158}
159
160fn make_list_count_key(name: &[u8]) -> Vec<u8> {
162 let mut v =
163 Vec::with_capacity(LIST_NAME_PREFIX.len() + name.len() + LIST_KEY_COUNT_SUFFIX.len());
164 v.extend_from_slice(LIST_NAME_PREFIX);
165 v.extend_from_slice(name);
166 v.extend_from_slice(LIST_KEY_COUNT_SUFFIX);
167 v
168}
169
170fn is_list_count_key(key: &[u8]) -> bool {
172 key.starts_with(LIST_NAME_PREFIX) && key.ends_with(LIST_KEY_COUNT_SUFFIX)
173}
174
175fn list_count_key_to_name(key: &[u8]) -> &[u8] {
177 let start = LIST_NAME_PREFIX.len();
178 let end = key.len() - LIST_KEY_COUNT_SUFFIX.len();
179 &key[start..end]
180}
181
182fn make_list_content_key(name: &[u8], idx: usize) -> Vec<u8> {
184 let idx_bytes = idx.to_be_bytes();
185 let mut v =
186 Vec::with_capacity(LIST_NAME_PREFIX.len() + name.len() + LIST_KEY_CONTENT_SUFFIX.len() + 8);
187 v.extend_from_slice(LIST_NAME_PREFIX);
188 v.extend_from_slice(name);
189 v.extend_from_slice(LIST_KEY_CONTENT_SUFFIX);
190 v.extend_from_slice(&idx_bytes);
191 v
192}
193
194fn make_list_content_prefix(name: &[u8]) -> Vec<u8> {
196 let mut v =
197 Vec::with_capacity(LIST_NAME_PREFIX.len() + name.len() + LIST_KEY_CONTENT_SUFFIX.len());
198 v.extend_from_slice(LIST_NAME_PREFIX);
199 v.extend_from_slice(name);
200 v.extend_from_slice(LIST_KEY_CONTENT_SUFFIX);
201 v
202}
203
204#[allow(dead_code)]
206fn make_counter_key(key: &[u8]) -> Vec<u8> {
207 let mut v = Vec::with_capacity(COUNTER_PREFIX.len() + key.len());
208 v.extend_from_slice(COUNTER_PREFIX);
209 v.extend_from_slice(key);
210 v
211}
212
213#[derive(Clone)]
218struct Pattern(Arc<Vec<PatternChar>>);
219
220impl std::ops::Deref for Pattern {
221 type Target = Vec<PatternChar>;
222 fn deref(&self) -> &Self::Target {
223 &self.0
224 }
225}
226
227#[derive(Clone, Debug)]
228enum PatternChar {
229 Exact(u8),
231 Wildcard,
233 AnyChar,
235}
236
237impl From<&[u8]> for Pattern {
238 fn from(pattern: &[u8]) -> Self {
239 Pattern::parse(pattern)
240 }
241}
242
243impl Pattern {
244 fn parse(pattern: &[u8]) -> Self {
245 let mut parsed = Vec::new();
246 let mut chars = pattern.iter().copied().peekable();
247 while let Some(c) = chars.next() {
248 match c {
249 b'*' => {
250 if !matches!(parsed.last(), Some(PatternChar::Wildcard)) {
252 parsed.push(PatternChar::Wildcard);
253 }
254 }
255 b'?' => parsed.push(PatternChar::AnyChar),
256 b'\\' => {
257 if let Some(next) = chars.next() {
258 parsed.push(PatternChar::Exact(next));
259 }
260 }
261 _ => parsed.push(PatternChar::Exact(c)),
262 }
263 }
264 Pattern(Arc::new(parsed))
265 }
266}
267
268fn is_match<P: Into<Pattern>>(pattern: P, text: &[u8]) -> bool {
269 let pattern = pattern.into();
270 let text_chars = text;
271 let pattern_len = pattern.len();
272 let text_len = text_chars.len();
273
274 let mut dp = vec![vec![false; pattern_len + 1]; text_len + 1];
276 dp[0][0] = true;
277
278 for j in 1..=pattern_len {
279 if matches!(pattern[j - 1], PatternChar::Wildcard) {
280 dp[0][j] = dp[0][j - 1];
281 }
282 }
283
284 for i in 1..=text_len {
285 for j in 1..=pattern_len {
286 match &pattern[j - 1] {
287 PatternChar::Exact(c) => {
288 if text_chars[i - 1] == *c {
289 dp[i][j] = dp[i - 1][j - 1];
290 }
291 }
292 PatternChar::AnyChar => {
293 dp[i][j] = dp[i - 1][j - 1];
294 }
295 PatternChar::Wildcard => {
296 dp[i][j] = dp[i][j - 1] || dp[i - 1][j];
297 }
298 }
299 }
300 }
301
302 dp[text_len][pattern_len]
303}
304
305pub type RedbCleanupFun = fn(&RedbStorageDB);
311
312fn def_cleanup(_db: &RedbStorageDB) {
314 #[cfg(feature = "ttl")]
315 {
316 let db = _db.clone();
317 std::thread::spawn(move || {
318 let limit = 5000;
319 loop {
320 std::thread::sleep(std::time::Duration::from_secs(60));
321 let mut total_cleanups = 0;
322 let now = std::time::Instant::now();
323 loop {
324 let count = db.cleanup(limit);
325 total_cleanups += count;
326 if count > 0 {
327 log::debug!(
328 "def_cleanup: {}, total cleanups: {}, active_count(): {}, cost time: {:?}",
329 count,
330 total_cleanups,
331 db.active_count(),
332 now.elapsed()
333 );
334 }
335 if count < limit {
336 break;
337 }
338 if db.active_count() > 50 {
339 std::thread::sleep(std::time::Duration::from_millis(500));
340 } else {
341 std::thread::sleep(std::time::Duration::from_millis(0));
342 }
343 }
344 if now.elapsed().as_secs() > 3 {
345 log::info!(
346 "total cleanups: {}, cost time: {:?}",
347 total_cleanups,
348 now.elapsed()
349 );
350 }
351 }
352 });
353 }
354}
355
356#[derive(Debug, Clone, Serialize, Deserialize)]
362pub struct RedbConfig {
363 pub path: String,
365 #[serde(default = "RedbConfig::default_cache_size")]
367 pub cache_size: usize,
368 #[serde(skip, default = "RedbConfig::cleanup_f_default")]
370 pub cleanup_f: RedbCleanupFun,
371}
372
373impl RedbConfig {
374 fn default_cache_size() -> usize {
375 1024 * 1024 * 1024 }
377
378 #[inline]
380 fn cleanup_f_default() -> RedbCleanupFun {
381 def_cleanup
382 }
383}
384
385impl Default for RedbConfig {
386 fn default() -> Self {
387 RedbConfig {
388 path: String::default(),
389 cache_size: Self::default_cache_size(),
390 cleanup_f: def_cleanup,
391 }
392 }
393}
394
395pub struct RedbStorageDB {
401 db: Arc<Database>,
402 cmd_tx: mpsc::Sender<Command>,
403 active_count: Arc<AtomicIsize>,
404}
405
406impl Clone for RedbStorageDB {
407 fn clone(&self) -> Self {
408 RedbStorageDB {
409 db: self.db.clone(),
410 cmd_tx: self.cmd_tx.clone(),
411 active_count: self.active_count.clone(),
412 }
413 }
414}
415
416impl fmt::Debug for RedbStorageDB {
417 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
418 f.debug_struct("RedbStorageDB").finish()
419 }
420}
421
422impl RedbStorageDB {
423 pub fn open(cfg: &RedbConfig) -> Result<Self> {
425 if let Some(parent) = std::path::Path::new(&cfg.path).parent() {
427 if !parent.as_os_str().is_empty() {
428 std::fs::create_dir_all(parent)?;
429 }
430 }
431
432 let mut builder = redb::Builder::new();
433 builder.set_cache_size(cfg.cache_size);
434
435 let db = if std::path::Path::new(&cfg.path).exists() {
437 builder.open(&cfg.path)?
438 } else {
439 builder.create(&cfg.path)?
440 };
441
442 {
445 let txn = db.begin_write()?;
446 let _ = txn.open_table(KV_TABLE);
447 let _ = txn.open_table(MAP_TABLE);
448 let _ = txn.open_table(LIST_TABLE);
449 let _ = txn.open_table(EXPIRE_KEYS_TABLE);
450 let _ = txn.open_table(KEY_EXPIRE_TABLE);
451 txn.commit()?;
452 }
453
454 let db = Arc::new(db);
455 let active_count = Arc::new(AtomicIsize::new(0));
456 let (cmd_tx, cmd_rx) = mpsc::channel::<Command>(CHANNEL_CAPACITY);
457
458 Self::start_background_thread(db.clone(), cmd_rx, active_count.clone());
459
460 let storage = RedbStorageDB {
461 db,
462 cmd_tx,
463 active_count,
464 };
465
466 (cfg.cleanup_f)(&storage);
467
468 Ok(storage)
469 }
470
471 #[cfg(feature = "ttl")]
473 #[inline]
474 pub fn cleanup(&self, limit: usize) -> usize {
475 Self::exec_cleanup(&self.db, limit).unwrap_or(0)
476 }
477
478 #[inline]
480 pub fn active_count(&self) -> isize {
481 self.active_count.load(Ordering::Relaxed)
482 }
483
484 async fn send_cmd<R, F>(&self, make_cmd: F) -> Result<R>
486 where
487 R: Send + 'static,
488 F: FnOnce(oneshot::Sender<Result<R>>) -> Command,
489 {
490 let (tx, rx) = oneshot::channel();
491 let cmd = make_cmd(tx);
492 self.cmd_tx
493 .send(cmd)
494 .await
495 .map_err(|_| anyhow!("redb command channel closed"))?;
496 rx.await
497 .map_err(|_| anyhow!("redb command response channel dropped"))?
498 }
499
500 fn start_background_thread(
501 db: Arc<Database>,
502 mut rx: mpsc::Receiver<Command>,
503 active_count: Arc<AtomicIsize>,
504 ) {
505 spawn_blocking(move || {
506 let _handle = Handle::current();
507 'outer: while let Some(cmd) = rx.blocking_recv() {
508 active_count.fetch_add(1, Ordering::Release);
509 match cmd {
510 Command::DBInsert { key, val, tx } => {
512 let result =
513 Self::exec_db_insert(&db, &key, &val).context("redb: db_insert");
514 let _ = tx.send(result);
515 }
516 Command::DBGet { key, tx } => {
517 let result = Self::exec_db_get(&db, &key).context("redb: db_get");
518 let _ = tx.send(result);
519 }
520 Command::DBRemove { key, tx } => {
521 let result = Self::exec_db_remove(&db, &key).context("redb: db_remove");
522 let _ = tx.send(result);
523 }
524 Command::DBContainsKey { key, tx } => {
525 let result =
526 Self::exec_db_contains_key(&db, &key).context("redb: db_contains_key");
527 let _ = tx.send(result);
528 }
529 Command::DBBatchInsert { key_vals, tx } => {
530 let result = Self::exec_db_batch_insert(&db, key_vals)
531 .context("redb: db_batch_insert");
532 let _ = tx.send(result);
533 }
534 Command::DBBatchRemove { keys, tx } => {
535 let result =
536 Self::exec_db_batch_remove(&db, keys).context("redb: db_batch_remove");
537 let _ = tx.send(result);
538 }
539 Command::DBCounterIncr { key, increment, tx } => {
540 let result = Self::exec_counter_incr(&db, &key, increment)
541 .context("redb: counter_incr");
542 let _ = tx.send(result);
543 }
544 Command::DBCounterDecr { key, decrement, tx } => {
545 let result = Self::exec_counter_decr(&db, &key, decrement)
546 .context("redb: counter_decr");
547 let _ = tx.send(result);
548 }
549 Command::DBCounterGet { key, tx } => {
550 let result = Self::exec_counter_get(&db, &key).context("redb: counter_get");
551 let _ = tx.send(result);
552 }
553 Command::DBCounterSet { key, val, tx } => {
554 let result =
555 Self::exec_counter_set(&db, &key, val).context("redb: counter_set");
556 let _ = tx.send(result);
557 }
558 Command::DBLen { tx } => {
559 let result = Self::exec_db_len(&db).context("redb: db_len");
560 let _ = tx.send(result);
561 }
562 Command::DBSize { tx } => {
563 let result = Self::exec_db_size(&db).context("redb: db_size");
564 let _ = tx.send(result);
565 }
566 Command::DBInfo { tx } => {
567 let result = Self::exec_db_info(&db).context("redb: db_info");
568 let _ = tx.send(result);
569 }
570
571 Command::DBMapGet { key, tx } => {
573 let result = Self::exec_db_map_contains_key(&db, &key)
574 .context("redb: db_map_contains_key");
575 let _ = tx.send(result);
576 }
577 Command::DBListGet { key, tx } => {
578 let result = Self::exec_db_list_contains_key(&db, &key)
579 .context("redb: db_list_contains_key");
580 let _ = tx.send(result);
581 }
582
583 Command::DBMapIter { tx } => {
585 let result = Self::exec_db_map_iter(&db).context("redb: db_map_iter");
586 let _ = tx.send(result);
587 }
588 Command::DBListIter { tx } => {
589 let result = Self::exec_db_list_iter(&db).context("redb: db_list_iter");
590 let _ = tx.send(result);
591 }
592 Command::DBScanIter { pattern, tx } => {
593 let result = Self::exec_db_scan(&db, &pattern).context("redb: db_scan");
594 let _ = tx.send(result);
595 }
596
597 #[cfg(feature = "ttl")]
599 Command::DBExpireAt { key, at, tx } => {
600 let result =
601 Self::exec_db_expire_at(&db, &key, at).context("redb: db_expire_at");
602 let _ = tx.send(result);
603 }
604 #[cfg(feature = "ttl")]
605 Command::DBExpire { key, dur, tx } => {
606 let result =
607 Self::exec_db_expire(&db, &key, dur).context("redb: db_expire");
608 let _ = tx.send(result);
609 }
610 #[cfg(feature = "ttl")]
611 Command::DBTtl { key, tx } => {
612 let result = Self::exec_db_ttl(&db, &key).context("redb: db_ttl");
613 let _ = tx.send(result);
614 }
615 #[cfg(feature = "ttl")]
616 Command::DBCleanup { tx } => {
617 let result = Self::exec_cleanup(&db, 1000).context("redb: cleanup");
618 let _ = tx.send(result);
619 }
620
621 Command::MapInsert { map, key, val, tx } => {
623 let result = Self::exec_map_insert(&db, &map.name, &key, &val)
624 .context("redb: map_insert");
625 let _ = tx.send(result);
626 }
627 Command::MapGet { map, key, tx } => {
628 let result =
629 Self::exec_map_get(&db, &map.name, &key).context("redb: map_get");
630 let _ = tx.send(result);
631 }
632 Command::MapRemove { map, key, tx } => {
633 let result =
634 Self::exec_map_remove(&db, &map.name, &key).context("redb: map_remove");
635 let _ = tx.send(result);
636 }
637 Command::MapContainsKey { map, key, tx } => {
638 let result = Self::exec_map_contains_key(&db, &map.name, &key)
639 .context("redb: map_contains_key");
640 let _ = tx.send(result);
641 }
642 #[cfg(feature = "map_len")]
643 Command::MapLen { map, tx } => {
644 let result = Self::exec_map_len(&db, &map.name).context("redb: map_len");
645 let _ = tx.send(result);
646 }
647 Command::MapIsEmpty { map, tx } => {
648 let result =
649 Self::exec_map_is_empty(&db, &map.name).context("redb: map_is_empty");
650 let _ = tx.send(result);
651 }
652 Command::MapClear { map, tx } => {
653 let result =
654 Self::exec_map_clear(&db, &map.name).context("redb: map_clear");
655 let _ = tx.send(result);
656 }
657 Command::MapRemoveAndFetch { map, key, tx } => {
658 let result = Self::exec_map_remove_and_fetch(&db, &map.name, &key)
659 .context("redb: map_remove_and_fetch");
660 let _ = tx.send(result);
661 }
662 Command::MapRemoveWithPrefix { map, prefix, tx } => {
663 let result = Self::exec_map_remove_with_prefix(&db, &map.name, &prefix)
664 .context("redb: map_remove_with_prefix");
665 let _ = tx.send(result);
666 }
667 Command::MapBatchInsert { map, key_vals, tx } => {
668 let result = Self::exec_map_batch_insert(&db, &map.name, key_vals)
669 .context("redb: map_batch_insert");
670 let _ = tx.send(result);
671 }
672 Command::MapBatchRemove { map, keys, tx } => {
673 let result = Self::exec_map_batch_remove(&db, &map.name, keys)
674 .context("redb: map_batch_remove");
675 let _ = tx.send(result);
676 }
677 Command::MapIter { map, tx } => {
678 let result = Self::exec_map_iter(&db, &map.name).context("redb: map_iter");
679 let _ = tx.send(result);
680 }
681 Command::MapKeyIter { map, tx } => {
682 let result =
683 Self::exec_map_key_iter(&db, &map.name).context("redb: map_key_iter");
684 let _ = tx.send(result);
685 }
686 Command::MapPrefixIter { map, prefix, tx } => {
687 let result = Self::exec_map_prefix_iter(&db, &map.name, &prefix)
688 .context("redb: map_prefix_iter");
689 let _ = tx.send(result);
690 }
691 #[cfg(feature = "ttl")]
692 Command::MapExpireAt { map, at, tx } => {
693 let result = Self::exec_map_expire_at(&db, &map.name, at)
694 .context("redb: map_expire_at");
695 let _ = tx.send(result);
696 }
697 #[cfg(feature = "ttl")]
698 Command::MapExpire { map, dur, tx } => {
699 let result =
700 Self::exec_map_expire(&db, &map.name, dur).context("redb: map_expire");
701 let _ = tx.send(result);
702 }
703 #[cfg(feature = "ttl")]
704 Command::MapTTL { map, tx } => {
705 let result = Self::exec_map_ttl(&db, &map.name).context("redb: map_ttl");
706 let _ = tx.send(result);
707 }
708 Command::MapIsExpired { map, tx } => {
709 let result = Self::exec_map_is_expired(&db, &map.name)
710 .context("redb: map_is_expired");
711 let _ = tx.send(result);
712 }
713
714 Command::ListPush { list, val, tx } => {
716 let result =
717 Self::exec_list_push(&db, &list.name, &val).context("redb: list_push");
718 let _ = tx.send(result);
719 }
720 Command::ListPushs { list, vals, tx } => {
721 let result = Self::exec_list_pushs(&db, &list.name, vals)
722 .context("redb: list_pushs");
723 let _ = tx.send(result);
724 }
725 Command::ListPushLimit {
726 list,
727 val,
728 limit,
729 pop_front_if_limited,
730 tx,
731 } => {
732 let result = Self::exec_list_push_limit(
733 &db,
734 &list.name,
735 &val,
736 limit,
737 pop_front_if_limited,
738 )
739 .context("redb: list_push_limit");
740 let _ = tx.send(result);
741 }
742 Command::ListPop { list, tx } => {
743 let result = Self::exec_list_pop(&db, &list.name).context("redb: list_pop");
744 let _ = tx.send(result);
745 }
746 Command::ListAll { list, tx } => {
747 let result = Self::exec_list_all(&db, &list.name).context("redb: list_all");
748 let _ = tx.send(result);
749 }
750 Command::ListGetIndex { list, idx, tx } => {
751 let result = Self::exec_list_get_index(&db, &list.name, idx)
752 .context("redb: list_get_index");
753 let _ = tx.send(result);
754 }
755 Command::ListLen { list, tx } => {
756 let result = Self::exec_list_len(&db, &list.name).context("redb: list_len");
757 let _ = tx.send(result);
758 }
759 Command::ListIsEmpty { list, tx } => {
760 let result = Self::exec_list_is_empty(&db, &list.name)
761 .context("redb: list_is_empty");
762 let _ = tx.send(result);
763 }
764 Command::ListClear { list, tx } => {
765 let result =
766 Self::exec_list_clear(&db, &list.name).context("redb: list_clear");
767 let _ = tx.send(result);
768 }
769 Command::ListIter { list, tx } => {
770 let result =
771 Self::exec_list_iter(&db, &list.name).context("redb: list_iter");
772 let _ = tx.send(result);
773 }
774 #[cfg(feature = "ttl")]
775 Command::ListExpireAt { list, at, tx } => {
776 let result = Self::exec_list_expire_at(&db, &list.name, at)
777 .context("redb: list_expire_at");
778 let _ = tx.send(result);
779 }
780 #[cfg(feature = "ttl")]
781 Command::ListExpire { list, dur, tx } => {
782 let result = Self::exec_list_expire(&db, &list.name, dur)
783 .context("redb: list_expire");
784 let _ = tx.send(result);
785 }
786 #[cfg(feature = "ttl")]
787 Command::ListTTL { list, tx } => {
788 let result = Self::exec_list_ttl(&db, &list.name).context("redb: list_ttl");
789 let _ = tx.send(result);
790 }
791 Command::ListIsExpired { list, tx } => {
792 let result = Self::exec_list_is_expired(&db, &list.name)
793 .context("redb: list_is_expired");
794 let _ = tx.send(result);
795 }
796
797 #[cfg(feature = "circuit-breaker")]
799 Command::DBInsertRaw { key, val, tx } => {
800 let result =
801 Self::exec_db_insert(&db, &key, &val).context("redb: db_insert_raw");
802 let _ = tx.send(result);
803 }
804 #[cfg(feature = "circuit-breaker")]
805 Command::DBGetRaw { key, tx } => {
806 let result = Self::exec_db_get(&db, &key).context("redb: db_get_raw");
807 let _ = tx.send(result);
808 }
809 #[cfg(feature = "circuit-breaker")]
810 Command::DBRemoveRaw { key, tx } => {
811 let result = Self::exec_db_remove(&db, &key).context("redb: db_remove_raw");
812 let _ = tx.send(result);
813 }
814
815 Command::Shutdown => break 'outer,
816 }
817 active_count.fetch_sub(1, Ordering::Release);
818 }
819 });
820 }
821
822 fn exec_db_insert(db: &Database, key: &[u8], val: &[u8]) -> Result<()> {
827 let txn = db.begin_write()?;
828 {
829 let mut table = txn.open_table(KV_TABLE)?;
830 table.insert(key, val)?;
831 #[cfg(feature = "ttl")]
833 Self::clear_key_ttl_in_txn(&txn, key)?;
834 }
835 txn.commit()?;
836 Ok(())
837 }
838
839 #[allow(dead_code)]
842 fn exec_db_insert_batch(db: &Database, batch: &mut Vec<BatchInsertItem>) {
843 let txn = match db.begin_write() {
844 Ok(t) => t,
845 Err(e) => {
846 for (_, _, tx) in batch.drain(..) {
847 let _ = tx.send(Err(anyhow::format_err!("{:?}", e)));
848 }
849 return;
850 }
851 };
852 let mut kv_table = match txn.open_table(KV_TABLE) {
853 Ok(t) => t,
854 Err(e) => {
855 for (_, _, tx) in batch.drain(..) {
856 let _ = tx.send(Err(anyhow::format_err!("{:?}", e)));
857 }
858 return;
859 }
860 };
861 #[cfg(feature = "ttl")]
862 let mut ke_table = match txn.open_table(KEY_EXPIRE_TABLE) {
863 Ok(t) => t,
864 Err(e) => {
865 for (_, _, tx) in batch.drain(..) {
866 let _ = tx.send(Err(anyhow::format_err!("{:?}", e)));
867 }
868 return;
869 }
870 };
871 for (key, val, tx) in batch.drain(..) {
872 if let Err(e) = kv_table.insert(key.as_slice(), val.as_slice()) {
873 let _ = tx.send(Err(anyhow::format_err!("{:?}", e)));
874 continue;
875 }
876 #[cfg(feature = "ttl")]
877 if let Ok(encoded) = postcard::to_stdvec(&TimestampMillis::MAX) {
878 if let Err(e) = ke_table.insert(key.as_slice(), encoded.as_slice()) {
879 let _ = tx.send(Err(anyhow::format_err!("{:?}", e)));
880 continue;
881 }
882 }
883 let _ = tx.send(Ok(()));
884 }
885 drop(kv_table);
886 #[cfg(feature = "ttl")]
887 {
888 drop(ke_table);
889 }
890 if let Err(e) = txn.commit() {
891 eprintln!("redb batch commit failed: {:?}", e);
893 }
894 }
895
896 fn exec_db_get(db: &Database, key: &[u8]) -> Result<Option<Vec<u8>>> {
897 let txn = db.begin_read()?;
898 let table = txn.open_table(KV_TABLE)?;
899 #[cfg(feature = "ttl")]
900 {
901 if let Some(expire_at_bytes) = txn.open_table(KEY_EXPIRE_TABLE)?.get(key)? {
902 let expire_at =
903 i64::from_be_bytes(expire_at_bytes.value().try_into().unwrap_or([0; 8]));
904 if expire_at <= timestamp_millis() {
905 return Ok(None);
906 }
907 }
908 }
909 match table.get(key)? {
910 Some(guard) => {
911 let val = guard.value().to_vec();
912 Ok(Some(val))
913 }
914 None => Ok(None),
915 }
916 }
917
918 fn exec_db_remove(db: &Database, key: &[u8]) -> Result<()> {
919 let txn = db.begin_write()?;
920 {
921 let mut table = txn.open_table(KV_TABLE)?;
922 table.remove(key)?;
923 Self::clear_key_ttl_in_txn(&txn, key)?;
924 }
925 txn.commit()?;
926 Ok(())
927 }
928
929 fn exec_db_contains_key(db: &Database, key: &[u8]) -> Result<bool> {
930 let txn = db.begin_read()?;
931 #[cfg(feature = "ttl")]
932 {
933 if let Ok(kv) = txn.open_table(KV_TABLE)?.get(key) {
934 if kv.is_some() {
935 if let Some(expire_at_bytes) = txn.open_table(KEY_EXPIRE_TABLE)?.get(key)? {
936 let expire_at = i64::from_be_bytes(
937 expire_at_bytes.value().try_into().unwrap_or([0; 8]),
938 );
939 if expire_at <= timestamp_millis() {
940 return Ok(false);
941 }
942 }
943 return Ok(true);
944 }
945 }
946 Ok(false)
947 }
948 #[cfg(not(feature = "ttl"))]
949 {
950 let table = txn.open_table(KV_TABLE)?;
951 Ok(table.get(key)?.is_some())
952 }
953 }
954
955 fn exec_db_batch_insert(db: &Database, key_vals: Vec<(Vec<u8>, Vec<u8>)>) -> Result<()> {
956 let txn = db.begin_write()?;
957 {
958 let mut table = txn.open_table(KV_TABLE)?;
959 for (key, val) in key_vals {
960 table.insert(key.as_slice(), val.as_slice())?;
961 #[cfg(feature = "ttl")]
963 Self::clear_key_ttl_in_txn(&txn, &key)?;
964 }
965 }
966 txn.commit()?;
967 Ok(())
968 }
969
970 fn exec_db_batch_remove(db: &Database, keys: Vec<Vec<u8>>) -> Result<()> {
971 let txn = db.begin_write()?;
972 {
973 let mut table = txn.open_table(KV_TABLE)?;
974 for key in keys {
975 table.remove(key.as_slice())?;
976 Self::clear_key_ttl_in_txn(&txn, &key)?;
977 }
978 }
979 txn.commit()?;
980 Ok(())
981 }
982
983 fn exec_counter_incr(db: &Database, key: &[u8], increment: isize) -> Result<()> {
984 let txn = db.begin_write()?;
985 {
986 let mut table = txn.open_table(KV_TABLE)?;
987 let current: isize = table
988 .get(key)?
989 .and_then(|g| {
990 <[u8; 8]>::try_from(g.value())
991 .ok()
992 .map(isize::from_be_bytes)
993 })
994 .unwrap_or(0);
995 let new = current + increment;
996 table.insert(key, &new.to_be_bytes()[..])?;
997 }
998 txn.commit()?;
999 Ok(())
1000 }
1001
1002 fn exec_counter_decr(db: &Database, key: &[u8], decrement: isize) -> Result<()> {
1003 Self::exec_counter_incr(db, key, -decrement)
1004 }
1005
1006 fn exec_counter_get(db: &Database, key: &[u8]) -> Result<Option<isize>> {
1007 let txn = db.begin_read()?;
1008 let table = txn.open_table(KV_TABLE)?;
1009 match table.get(key)? {
1010 Some(g) => {
1011 let arr: [u8; 8] = <[u8; 8]>::try_from(g.value()).unwrap_or([0; 8]);
1012 let val = isize::from_be_bytes(arr);
1013 Ok(Some(val))
1014 }
1015 None => Ok(None),
1016 }
1017 }
1018
1019 fn exec_counter_set(db: &Database, key: &[u8], val: isize) -> Result<()> {
1020 let txn = db.begin_write()?;
1021 {
1022 let mut table = txn.open_table(KV_TABLE)?;
1023 table.insert(key, &val.to_be_bytes()[..])?;
1024 }
1025 txn.commit()?;
1026 Ok(())
1027 }
1028
1029 fn exec_db_len(db: &Database) -> Result<usize> {
1030 #[cfg(feature = "len")]
1031 {
1032 let txn = db.begin_read()?;
1033 let kv_table = txn.open_table(KV_TABLE)?;
1034 let mut count = 0usize;
1035 for result in kv_table.iter()? {
1036 let (key, _) = result?;
1037 let k = key.value();
1038 if k.starts_with(COUNTER_PREFIX) || k.starts_with(KEY_PREFIX) {
1039 continue;
1040 }
1041 count += 1;
1042 }
1043 Ok(count)
1044 }
1045 #[cfg(not(feature = "len"))]
1046 {
1047 let _ = db;
1048 Ok(0)
1049 }
1050 }
1051
1052 fn exec_db_size(db: &Database) -> Result<usize> {
1053 let txn = db.begin_read()?;
1054 let kv_table = txn.open_table(KV_TABLE)?;
1055 let map_table = txn.open_table(MAP_TABLE)?;
1056 let list_table = txn.open_table(LIST_TABLE)?;
1057 let kv_count = kv_table.len()? as usize;
1058 let map_count = map_table.len()? as usize;
1059 let list_count = list_table.len()? as usize;
1060 Ok(kv_count + map_count + list_count)
1061 }
1062
1063 fn exec_db_info(_db: &Database) -> Result<Value> {
1064 Ok(serde_json::json!({
1065 "storage_engine": "Redb",
1066 }))
1067 }
1068
1069 fn exec_db_map_contains_key(db: &Database, name: &[u8]) -> Result<bool> {
1070 let txn = db.begin_read()?;
1071 #[cfg(feature = "ttl")]
1072 if let Ok(table) = txn.open_table(KEY_EXPIRE_TABLE) {
1073 if let Some(guard) = table.get(name)? {
1074 let expire_at: TimestampMillis = postcard::from_bytes(guard.value())?;
1075 if expire_at <= timestamp_millis() {
1076 return Ok(false);
1077 }
1078 }
1079 }
1080 let table = txn.open_table(MAP_TABLE)?;
1081 let count_key = make_map_count_key(name);
1082 Ok(table.get(count_key.as_slice())?.is_some())
1083 }
1084
1085 fn exec_db_list_contains_key(db: &Database, name: &[u8]) -> Result<bool> {
1086 let txn = db.begin_read()?;
1087 #[cfg(feature = "ttl")]
1088 if let Ok(table) = txn.open_table(KEY_EXPIRE_TABLE) {
1089 if let Some(guard) = table.get(name)? {
1090 let expire_at: TimestampMillis = postcard::from_bytes(guard.value())?;
1091 if expire_at <= timestamp_millis() {
1092 return Ok(false);
1093 }
1094 }
1095 }
1096 let table = txn.open_table(LIST_TABLE)?;
1097 let count_key = make_list_count_key(name);
1098 Ok(table.get(count_key.as_slice())?.is_some())
1099 }
1100
1101 fn exec_db_map_iter(db: &Database) -> Result<Vec<(Vec<u8>, Vec<u8>)>> {
1102 let txn = db.begin_read()?;
1103 let table = txn.open_table(MAP_TABLE)?;
1104 let mut results = Vec::new();
1105 for result in table.iter()? {
1106 let (key, val) = result?;
1107 let k = key.value().to_vec();
1108 if is_map_count_key(&k) {
1109 results.push((k, val.value().to_vec()));
1110 }
1111 }
1112 Ok(results)
1113 }
1114
1115 fn exec_db_list_iter(db: &Database) -> Result<Vec<(Vec<u8>, Vec<u8>)>> {
1116 let txn = db.begin_read()?;
1117 let table = txn.open_table(LIST_TABLE)?;
1118 let mut results = Vec::new();
1119 for result in table.iter()? {
1120 let (key, val) = result?;
1121 let k = key.value().to_vec();
1122 if is_list_count_key(&k) {
1123 results.push((k, val.value().to_vec()));
1124 }
1125 }
1126 Ok(results)
1127 }
1128
1129 fn exec_db_scan(db: &Database, pattern: &[u8]) -> Result<Vec<Vec<u8>>> {
1130 let txn = db.begin_read()?;
1131 let table = txn.open_table(KV_TABLE)?;
1132 let pattern = Pattern::from(pattern);
1133 let mut results = Vec::new();
1134 for result in table.iter()? {
1135 let (key, _) = result?;
1136 let k = key.value();
1137 if k.starts_with(COUNTER_PREFIX) {
1138 continue;
1139 }
1140 if is_match(pattern.clone(), k) {
1141 results.push(k.to_vec());
1142 }
1143 }
1144 let map_table = txn.open_table(MAP_TABLE)?;
1146 for result in map_table.iter()? {
1147 let (key, _) = result?;
1148 let k = key.value();
1149 if is_map_count_key(k) && is_match(pattern.clone(), k) {
1150 results.push(k.to_vec());
1151 }
1152 }
1153 let list_table = txn.open_table(LIST_TABLE)?;
1154 for result in list_table.iter()? {
1155 let (key, _) = result?;
1156 let k = key.value();
1157 if is_list_count_key(k) && is_match(pattern.clone(), k) {
1158 results.push(k.to_vec());
1159 }
1160 }
1161 Ok(results)
1162 }
1163
1164 #[cfg(feature = "ttl")]
1167 fn clear_key_ttl_in_txn(txn: &WriteTransaction, key: &[u8]) -> Result<()> {
1168 let mut ke_tbl = txn.open_table(KEY_EXPIRE_TABLE)?;
1170 let old_expire_bytes: Option<Vec<u8>> = ke_tbl.get(key)?.map(|g| g.value().to_vec());
1171 if let Some(ref expire_bytes) = old_expire_bytes {
1172 let mut ek = Vec::with_capacity(8 + key.len());
1173 ek.extend_from_slice(expire_bytes);
1174 ek.extend_from_slice(key);
1175 let mut ek_tbl = txn.open_table(EXPIRE_KEYS_TABLE)?;
1177 ek_tbl.remove(ek.as_slice())?;
1178 ke_tbl.remove(key)?;
1179 }
1180 Ok(())
1181 }
1182
1183 #[cfg(not(feature = "ttl"))]
1184 fn clear_key_ttl_in_txn(_txn: &WriteTransaction, _key: &[u8]) -> Result<()> {
1185 Ok(())
1186 }
1187
1188 #[cfg(feature = "ttl")]
1189 fn set_key_ttl_in_txn(
1190 txn: &WriteTransaction,
1191 key: &[u8],
1192 expire_at: TimestampMillis,
1193 ) -> Result<bool> {
1194 let at_bytes = expire_at.to_be_bytes();
1196 let mut ek = Vec::with_capacity(8 + key.len());
1197 ek.extend_from_slice(&at_bytes);
1198 ek.extend_from_slice(key);
1199
1200 let mut ke_tbl = txn.open_table(KEY_EXPIRE_TABLE)?;
1202 let old_expire_data: Option<Vec<u8>> = ke_tbl.get(key)?.map(|g| g.value().to_vec());
1203 if let Some(ref old_bytes) = old_expire_data {
1204 let mut old_ek = Vec::with_capacity(8 + key.len());
1205 old_ek.extend_from_slice(old_bytes);
1206 old_ek.extend_from_slice(key);
1207 let mut ek_tbl = txn.open_table(EXPIRE_KEYS_TABLE)?;
1209 ek_tbl.remove(old_ek.as_slice())?;
1210 }
1211
1212 ke_tbl.insert(key, &at_bytes[..])?;
1214 let mut ek_tbl = txn.open_table(EXPIRE_KEYS_TABLE)?;
1215 ek_tbl.insert(ek.as_slice(), &b"kv"[..])?;
1216
1217 Ok(true)
1218 }
1219
1220 #[cfg(feature = "ttl")]
1221 fn exec_db_expire_at(db: &Database, key: &[u8], at: TimestampMillis) -> Result<bool> {
1222 let txn = db.begin_read()?;
1223 let exists = {
1224 let kv = txn.open_table(KV_TABLE)?.get(key)?.is_some()
1225 || txn
1226 .open_table(MAP_TABLE)?
1227 .get(make_map_count_key(key).as_slice())?
1228 .is_some()
1229 || txn
1230 .open_table(LIST_TABLE)?
1231 .get(make_list_count_key(key).as_slice())?
1232 .is_some();
1233 kv
1234 };
1235 drop(txn);
1236 if !exists {
1237 return Ok(false);
1238 }
1239 let txn = db.begin_write()?;
1240 {
1241 Self::set_key_ttl_in_txn(&txn, key, at)?;
1242 }
1243 txn.commit()?;
1244 Ok(true)
1245 }
1246
1247 #[cfg(feature = "ttl")]
1248 fn exec_db_expire(db: &Database, key: &[u8], dur: TimestampMillis) -> Result<bool> {
1249 let at = timestamp_millis() + dur;
1250 Self::exec_db_expire_at(db, key, at)
1251 }
1252
1253 #[cfg(feature = "ttl")]
1254 fn exec_db_ttl(db: &Database, key: &[u8]) -> Result<Option<TimestampMillis>> {
1255 let txn = db.begin_read()?;
1256 let ke_table = txn.open_table(KEY_EXPIRE_TABLE)?;
1257 match ke_table.get(key)? {
1258 Some(guard) => {
1259 let expire_at = i64::from_be_bytes(guard.value().try_into().unwrap_or([0; 8]));
1260 let now = timestamp_millis();
1261 if expire_at <= now {
1262 Ok(None)
1263 } else {
1264 Ok(Some(expire_at - now))
1265 }
1266 }
1267 None => {
1268 let kv = txn.open_table(KV_TABLE)?;
1271 if kv.get(key)?.is_some() {
1272 return Ok(Some(TimestampMillis::MAX));
1273 }
1274 let mk = make_map_count_key(key);
1275 if txn.open_table(MAP_TABLE)?.get(mk.as_slice())?.is_some() {
1276 return Ok(Some(TimestampMillis::MAX));
1277 }
1278 let lk = make_list_count_key(key);
1279 if txn.open_table(LIST_TABLE)?.get(lk.as_slice())?.is_some() {
1280 return Ok(Some(TimestampMillis::MAX));
1281 }
1282 Ok(None)
1283 }
1284 }
1285 }
1286
1287 #[cfg(feature = "ttl")]
1288 fn exec_cleanup(db: &Database, limit: usize) -> Result<usize> {
1289 let now = timestamp_millis();
1290 let mut count = 0usize;
1291
1292 let to_remove: Vec<Vec<u8>> = {
1295 let txn = db.begin_read()?;
1296 let expire_keys = txn.open_table(EXPIRE_KEYS_TABLE)?;
1297 expire_keys
1298 .range::<&[u8]>(..)?
1299 .filter_map(|r| {
1300 r.ok().and_then(|(k, _)| {
1301 let key_bytes = k.value();
1302 if key_bytes.len() < 8 {
1303 return None;
1304 }
1305 let expire_at = u64::from_be_bytes(
1306 <[u8; 8]>::try_from(&key_bytes[..8]).unwrap_or([0; 8]),
1307 ) as TimestampMillis;
1308 if expire_at <= now {
1309 Some(key_bytes.to_vec())
1311 } else {
1312 None
1313 }
1314 })
1315 })
1316 .take(limit)
1317 .collect()
1318 };
1320
1321 if to_remove.is_empty() {
1322 return Ok(0);
1323 }
1324
1325 let txn = db.begin_write()?;
1327 {
1328 let mut expire_keys = txn.open_table(EXPIRE_KEYS_TABLE)?;
1329 let mut key_expire = txn.open_table(KEY_EXPIRE_TABLE)?;
1330 let mut kv_table = txn.open_table(KV_TABLE)?;
1331 let mut map_table = txn.open_table(MAP_TABLE)?;
1332 let mut list_table = txn.open_table(LIST_TABLE)?;
1333
1334 for expired_key in &to_remove {
1335 let (_, actual_key) = expired_key.split_at(8);
1336 let _ = kv_table.remove(actual_key)?;
1337 let _ = map_table.remove(actual_key)?;
1338 let _ = list_table.remove(actual_key)?;
1339 let _ = key_expire.remove(actual_key)?;
1340 let _ = expire_keys.remove(expired_key.as_slice())?;
1341 count += 1;
1342 }
1343 }
1344 txn.commit()?;
1345 Ok(count)
1346 }
1347
1348 #[cfg(feature = "ttl")]
1351 fn is_container_expired(db: &Database, name: &[u8]) -> Result<bool> {
1352 let txn = db.begin_read()?;
1353 let expired = match txn.open_table(KEY_EXPIRE_TABLE)?.get(name)? {
1354 Some(guard) => {
1355 let expire_at = i64::from_be_bytes(guard.value().try_into().unwrap_or([0; 8]));
1356 expire_at <= timestamp_millis()
1357 }
1358 None => {
1359 false
1364 }
1365 };
1366 Ok(expired)
1367 }
1368
1369 fn exec_map_insert(db: &Database, name: &[u8], key: &[u8], val: &[u8]) -> Result<()> {
1372 let txn = db.begin_write()?;
1373 {
1374 let mut table = txn.open_table(MAP_TABLE)?;
1375 let item_key = make_map_item_key(name, key);
1376 let is_new = table.insert(item_key.as_slice(), val)?.is_none();
1377 #[cfg(feature = "ttl")]
1379 {
1380 let mut ke = txn.open_table(KEY_EXPIRE_TABLE)?;
1381 let expired = ke.get(name)?.is_none_or(|g| {
1382 i64::from_be_bytes(g.value().try_into().unwrap_or([0; 8])) <= timestamp_millis()
1383 });
1384 if expired {
1385 let old_bytes = ke.get(name)?.map(|g| g.value().to_vec());
1387 if let Some(ref bytes) = old_bytes {
1388 let mut ek = Vec::with_capacity(8 + name.len());
1389 ek.extend_from_slice(bytes);
1390 ek.extend_from_slice(name);
1391 let mut ek_tbl = txn.open_table(EXPIRE_KEYS_TABLE)?;
1392 ek_tbl.remove(ek.as_slice())?;
1393 ke.remove(name)?;
1394 }
1395 }
1396 }
1397 #[cfg(feature = "map_len")]
1398 if is_new {
1399 let count_key = make_map_count_key(name);
1400 let old = table
1401 .get(count_key.as_slice())?
1402 .map(|g| isize::from_be_bytes(g.value().try_into().unwrap_or([0; 8])))
1403 .unwrap_or(0);
1404 let new_count = old + 1;
1405 table.insert(count_key.as_slice(), &new_count.to_be_bytes()[..])?;
1406 }
1407 }
1408 txn.commit()?;
1409 Ok(())
1410 }
1411
1412 fn exec_map_get(db: &Database, name: &[u8], key: &[u8]) -> Result<Option<Vec<u8>>> {
1413 let txn = db.begin_read()?;
1414 #[cfg(feature = "ttl")]
1415 {
1416 if let Some(expire_at_bytes) = txn.open_table(KEY_EXPIRE_TABLE)?.get(name)? {
1417 let expire_at =
1418 i64::from_be_bytes(expire_at_bytes.value().try_into().unwrap_or([0; 8]));
1419 if expire_at <= timestamp_millis() {
1420 return Ok(None);
1421 }
1422 }
1423 }
1424 let table = txn.open_table(MAP_TABLE)?;
1425 let item_key = make_map_item_key(name, key);
1426 match table.get(item_key.as_slice())? {
1427 Some(guard) => Ok(Some(guard.value().to_vec())),
1428 None => Ok(None),
1429 }
1430 }
1431
1432 fn exec_map_remove(db: &Database, name: &[u8], key: &[u8]) -> Result<()> {
1433 let txn = db.begin_write()?;
1434 {
1435 let mut table = txn.open_table(MAP_TABLE)?;
1436 let item_key = make_map_item_key(name, key);
1437 let existed = table.remove(item_key.as_slice())?.is_some();
1438 #[cfg(feature = "map_len")]
1439 if existed {
1440 let count_key = make_map_count_key(name);
1441 let old = table
1442 .get(count_key.as_slice())?
1443 .map(|g| isize::from_be_bytes(g.value().try_into().unwrap_or([1; 8])))
1444 .unwrap_or(1);
1445 let new_count = (old - 1).max(0);
1446 table.insert(count_key.as_slice(), &new_count.to_be_bytes()[..])?;
1447 }
1448 }
1449 txn.commit()?;
1450 Ok(())
1451 }
1452
1453 fn exec_map_contains_key(db: &Database, name: &[u8], key: &[u8]) -> Result<bool> {
1454 let txn = db.begin_read()?;
1455 #[cfg(feature = "ttl")]
1456 if let Ok(table) = txn.open_table(KEY_EXPIRE_TABLE) {
1457 if let Some(guard) = table.get(name)? {
1458 let expire_at: TimestampMillis = postcard::from_bytes(guard.value())?;
1459 if expire_at <= timestamp_millis() {
1460 return Ok(false);
1461 }
1462 }
1463 }
1464 let table = txn.open_table(MAP_TABLE)?;
1465 let item_key = make_map_item_key(name, key);
1466 Ok(table.get(item_key.as_slice())?.is_some())
1467 }
1468
1469 #[cfg(feature = "map_len")]
1470 fn exec_map_len(db: &Database, name: &[u8]) -> Result<usize> {
1471 let txn = db.begin_read()?;
1472 #[cfg(feature = "ttl")]
1473 if let Ok(table) = txn.open_table(KEY_EXPIRE_TABLE) {
1474 if let Some(guard) = table.get(name)? {
1475 let expire_at = i64::from_be_bytes(guard.value().try_into().unwrap_or([0; 8]));
1476 if expire_at <= timestamp_millis() {
1477 return Ok(0);
1478 }
1479 }
1480 }
1481 let table = txn.open_table(MAP_TABLE)?;
1482 let count_key = make_map_count_key(name);
1483 match table.get(count_key.as_slice())? {
1484 Some(g) => {
1485 let c = isize::from_be_bytes(g.value().try_into().unwrap_or([0; 8]));
1486 Ok(c.max(0) as usize)
1487 }
1488 None => Ok(0),
1489 }
1490 }
1491
1492 fn exec_map_is_empty(db: &Database, name: &[u8]) -> Result<bool> {
1493 let txn = db.begin_read()?;
1494 #[cfg(feature = "ttl")]
1495 if let Ok(table) = txn.open_table(KEY_EXPIRE_TABLE) {
1496 if let Some(guard) = table.get(name)? {
1497 let expire_at: TimestampMillis = postcard::from_bytes(guard.value())?;
1498 if expire_at <= timestamp_millis() {
1499 return Ok(true);
1500 }
1501 }
1502 }
1503 let table = txn.open_table(MAP_TABLE)?;
1504 let item_prefix = make_map_item_prefix(name);
1505 let mut iter = table.range(item_prefix.as_slice()..)?;
1506 Ok(iter.next().is_none())
1507 }
1508
1509 fn exec_map_clear(db: &Database, name: &[u8]) -> Result<()> {
1510 let txn = db.begin_write()?;
1511 {
1512 let mut table = txn.open_table(MAP_TABLE)?;
1513 let item_prefix = make_map_item_prefix(name);
1514 {
1515 let extract = table.extract_from_if(item_prefix.as_slice().., |_, _| true)?;
1516 let _drained: Vec<_> = extract.collect::<Result<Vec<_>, _>>()?;
1517 }
1518 let count_key = make_map_count_key(name);
1519 let _ = table.remove(count_key.as_slice())?;
1520 }
1521 txn.commit()?;
1522 Ok(())
1523 }
1524
1525 fn exec_map_remove_and_fetch(
1526 db: &Database,
1527 name: &[u8],
1528 key: &[u8],
1529 ) -> Result<Option<Vec<u8>>> {
1530 #[cfg(feature = "ttl")]
1531 if Self::is_container_expired(db, name)? {
1532 return Ok(None);
1533 }
1534 let txn = db.begin_write()?;
1535 let (result, existed);
1536 {
1537 let mut table = txn.open_table(MAP_TABLE)?;
1538 let item_key = make_map_item_key(name, key);
1539 let removed = table.remove(item_key.as_slice())?;
1540 result = match removed.map(|g| g.value().to_vec()) {
1541 Some(v) => Ok(Some(v)),
1542 None => Ok(None),
1543 };
1544 existed = result.as_ref().ok().and_then(|r| r.as_ref()).is_some();
1545 }
1546 #[cfg(feature = "map_len")]
1547 if existed {
1548 let mut table = txn.open_table(MAP_TABLE)?;
1549 let count_key = make_map_count_key(name);
1550 let old = table
1551 .get(count_key.as_slice())?
1552 .map(|g| isize::from_be_bytes(g.value().try_into().unwrap_or([1; 8])))
1553 .unwrap_or(1);
1554 let new_count = (old - 1).max(0);
1555 table.insert(count_key.as_slice(), &new_count.to_be_bytes()[..])?;
1556 }
1557 txn.commit()?;
1558 result
1559 }
1560
1561 fn exec_map_remove_with_prefix(db: &Database, name: &[u8], prefix: &[u8]) -> Result<()> {
1562 let txn = db.begin_write()?;
1563 {
1564 let mut table = txn.open_table(MAP_TABLE)?;
1565 let item_prefix = make_map_item_key(name, prefix);
1566 let keys_to_remove: Vec<Vec<u8>> = table
1567 .range(item_prefix.as_slice()..)?
1568 .filter_map(|r| {
1569 r.ok().and_then(|(k, _)| {
1570 let k = k.value().to_vec();
1571 if k.starts_with(item_prefix.as_slice()) {
1572 Some(k)
1573 } else {
1574 None
1575 }
1576 })
1577 })
1578 .collect();
1579 for k in keys_to_remove {
1580 let _ = table.remove(k.as_slice())?;
1581 }
1582 #[cfg(feature = "map_len")]
1583 {
1584 let count_key = make_map_count_key(name);
1585 let item_prefix2 = make_map_item_prefix(name);
1586 let mut remaining = 0usize;
1587 for r in table.range(item_prefix2.as_slice()..)? {
1588 let (k, _) = r?;
1589 if !k.value().starts_with(item_prefix2.as_slice()) {
1590 break;
1591 }
1592 if !k.value().ends_with(b"@__count@") {
1593 remaining += 1;
1594 }
1595 }
1596 let encoded = postcard::to_stdvec(&remaining)?;
1597 table.insert(count_key.as_slice(), encoded.as_slice())?;
1598 }
1599 }
1600 txn.commit()?;
1601 Ok(())
1602 }
1603
1604 fn exec_map_batch_insert(
1605 db: &Database,
1606 name: &[u8],
1607 key_vals: Vec<(Vec<u8>, Vec<u8>)>,
1608 ) -> Result<()> {
1609 let txn = db.begin_write()?;
1610 {
1611 let mut table = txn.open_table(MAP_TABLE)?;
1612 #[cfg(feature = "map_len")]
1613 let mut new_count: usize = 0;
1614 for (key, val) in key_vals {
1615 let item_key = make_map_item_key(name, &key);
1616 let is_new = table.insert(item_key.as_slice(), val.as_slice())?.is_none();
1617 #[cfg(feature = "map_len")]
1618 if is_new {
1619 new_count += 1;
1620 }
1621 }
1622 #[cfg(feature = "map_len")]
1623 if new_count > 0 {
1624 let count_key = make_map_count_key(name);
1625 let old = table
1626 .get(count_key.as_slice())?
1627 .map(|g| isize::from_be_bytes(g.value().try_into().unwrap_or([0; 8])))
1628 .unwrap_or(0);
1629 let total = old + new_count as isize;
1630 table.insert(count_key.as_slice(), &total.to_be_bytes()[..])?;
1631 }
1632 }
1633 txn.commit()?;
1634 Ok(())
1635 }
1636
1637 fn exec_map_batch_remove(db: &Database, name: &[u8], keys: Vec<Vec<u8>>) -> Result<()> {
1638 let txn = db.begin_write()?;
1639 {
1640 let mut table = txn.open_table(MAP_TABLE)?;
1641 #[cfg(feature = "map_len")]
1642 let mut removed_count: usize = 0;
1643 for key in keys {
1644 let item_key = make_map_item_key(name, &key);
1645 let existed = table.remove(item_key.as_slice())?.is_some();
1646 #[cfg(feature = "map_len")]
1647 if existed {
1648 removed_count += 1;
1649 }
1650 }
1651 #[cfg(feature = "map_len")]
1652 if removed_count > 0 {
1653 let count_key = make_map_count_key(name);
1654 let old = table
1655 .get(count_key.as_slice())?
1656 .map(|g| isize::from_be_bytes(g.value().try_into().unwrap_or([0; 8])))
1657 .unwrap_or(0);
1658 let total = (old - removed_count as isize).max(0);
1659 table.insert(count_key.as_slice(), &total.to_be_bytes()[..])?;
1660 }
1661 }
1662 txn.commit()?;
1663 Ok(())
1664 }
1665
1666 fn exec_map_iter(db: &Database, name: &[u8]) -> Result<Vec<(Vec<u8>, Vec<u8>)>> {
1667 #[cfg(feature = "ttl")]
1668 if Self::is_container_expired(db, name)? {
1669 return Ok(Vec::new());
1670 }
1671 let txn = db.begin_read()?;
1672 let table = txn.open_table(MAP_TABLE)?;
1673 let item_prefix = make_map_item_prefix(name);
1674 let mut results = Vec::new();
1675 for result in table.range(item_prefix.as_slice()..)? {
1676 let (key, val) = result?;
1677 let k = key.value();
1678 if !k.starts_with(item_prefix.as_slice()) {
1680 break;
1681 }
1682 let inner_key = k[item_prefix.len()..].to_vec();
1683 results.push((inner_key, val.value().to_vec()));
1684 }
1685 Ok(results)
1686 }
1687
1688 fn exec_map_key_iter(db: &Database, name: &[u8]) -> Result<Vec<Vec<u8>>> {
1689 #[cfg(feature = "ttl")]
1690 if Self::is_container_expired(db, name)? {
1691 return Ok(Vec::new());
1692 }
1693 let txn = db.begin_read()?;
1694 let table = txn.open_table(MAP_TABLE)?;
1695 let item_prefix = make_map_item_prefix(name);
1696 let mut results = Vec::new();
1697 for result in table.range(item_prefix.as_slice()..)? {
1698 let (key, _) = result?;
1699 let k = key.value();
1700 if !k.starts_with(item_prefix.as_slice()) {
1701 break;
1702 }
1703 results.push(k[item_prefix.len()..].to_vec());
1704 }
1705 Ok(results)
1706 }
1707
1708 fn exec_map_prefix_iter(
1709 db: &Database,
1710 name: &[u8],
1711 prefix: &[u8],
1712 ) -> Result<Vec<(Vec<u8>, Vec<u8>)>> {
1713 #[cfg(feature = "ttl")]
1714 if Self::is_container_expired(db, name)? {
1715 return Ok(Vec::new());
1716 }
1717 let txn = db.begin_read()?;
1718 let table = txn.open_table(MAP_TABLE)?;
1719 let start_key = make_map_item_key(name, prefix);
1720 let prefix_key = make_map_item_prefix(name);
1721 let mut results = Vec::new();
1722 for result in table.range(start_key.as_slice()..)? {
1723 let (key, val) = result?;
1724 let k = key.value();
1725 if !k.starts_with(prefix_key.as_slice()) {
1726 break;
1727 }
1728 if !k.starts_with(start_key.as_slice()) {
1729 break;
1730 }
1731 let inner_key = k[prefix_key.len()..].to_vec();
1732 results.push((inner_key, val.value().to_vec()));
1733 }
1734 Ok(results)
1735 }
1736
1737 #[cfg(feature = "ttl")]
1738 fn exec_map_expire_at(db: &Database, name: &[u8], at: TimestampMillis) -> Result<bool> {
1739 let txn = db.begin_write()?;
1740 {
1741 let map_table = txn.open_table(MAP_TABLE)?;
1742 let count_key = make_map_count_key(name);
1743 if map_table.get(count_key.as_slice())?.is_none() {
1744 return Ok(false);
1745 }
1746 Self::set_key_ttl_in_txn(&txn, name, at)?;
1747 }
1748 txn.commit()?;
1749 Ok(true)
1750 }
1751
1752 #[cfg(feature = "ttl")]
1753 fn exec_map_expire(db: &Database, name: &[u8], dur: TimestampMillis) -> Result<bool> {
1754 let at = timestamp_millis() + dur;
1755 Self::exec_map_expire_at(db, name, at)
1756 }
1757
1758 #[cfg(feature = "ttl")]
1759 fn exec_map_ttl(db: &Database, name: &[u8]) -> Result<Option<TimestampMillis>> {
1760 Self::exec_db_ttl(db, name)
1761 }
1762
1763 fn exec_map_is_expired(db: &Database, name: &[u8]) -> Result<bool> {
1764 #[cfg(feature = "ttl")]
1765 {
1766 match Self::exec_db_ttl(db, name)? {
1767 Some(ttl) if ttl > 0 => Ok(false),
1768 _ => Ok(true),
1769 }
1770 }
1771 #[cfg(not(feature = "ttl"))]
1772 {
1773 let _ = (db, name);
1774 Ok(false)
1775 }
1776 }
1777
1778 fn exec_list_push(db: &Database, name: &[u8], val: &[u8]) -> Result<()> {
1781 #[cfg(feature = "ttl")]
1783 {
1784 let check_txn = db.begin_read()?;
1785 let expired = match check_txn.open_table(KEY_EXPIRE_TABLE)?.get(name)? {
1786 Some(guard) => {
1787 let expire_at = i64::from_be_bytes(guard.value().try_into().unwrap_or([0; 8]));
1788 expire_at <= timestamp_millis()
1789 }
1790 None => true,
1791 };
1792 drop(check_txn);
1793 if expired {
1794 let write_txn = db.begin_write()?;
1795 Self::clear_key_ttl_in_txn(&write_txn, name)?;
1796 write_txn.commit()?;
1797 }
1798 }
1799 let txn = db.begin_write()?;
1800 {
1801 let mut table = txn.open_table(LIST_TABLE)?;
1802 let count_key = make_list_count_key(name);
1803 let (start, end) = table
1804 .get(count_key.as_slice())?
1805 .map(|g| postcard::from_bytes::<(usize, usize)>(g.value()).unwrap_or((0, 0)))
1806 .unwrap_or((0, 0));
1807 let content_key = make_list_content_key(name, end);
1808 table.insert(content_key.as_slice(), val)?;
1809 let encoded = postcard::to_stdvec(&(start, end + 1))?;
1810 table.insert(count_key.as_slice(), encoded.as_slice())?;
1811 }
1812 txn.commit()?;
1813 Ok(())
1814 }
1815
1816 fn exec_list_pushs(db: &Database, name: &[u8], vals: Vec<Vec<u8>>) -> Result<()> {
1817 let txn = db.begin_write()?;
1818 {
1819 let mut table = txn.open_table(LIST_TABLE)?;
1820 let count_key = make_list_count_key(name);
1821 let (start, mut end) = table
1822 .get(count_key.as_slice())?
1823 .map(|g| postcard::from_bytes::<(usize, usize)>(g.value()).unwrap_or((0, 0)))
1824 .unwrap_or((0, 0));
1825 for val in &vals {
1826 let content_key = make_list_content_key(name, end);
1827 table.insert(content_key.as_slice(), val.as_slice())?;
1828 end += 1;
1829 }
1830 let encoded = postcard::to_stdvec(&(start, end))?;
1831 table.insert(count_key.as_slice(), encoded.as_slice())?;
1832 }
1833 txn.commit()?;
1834 Ok(())
1835 }
1836
1837 fn exec_list_push_limit(
1838 db: &Database,
1839 name: &[u8],
1840 val: &[u8],
1841 limit: usize,
1842 pop_front_if_limited: bool,
1843 ) -> Result<Option<Vec<u8>>> {
1844 let txn = db.begin_write()?;
1845 let result = {
1846 let mut table = txn.open_table(LIST_TABLE)?;
1847 let count_key = make_list_count_key(name);
1848 let (mut start, mut end) = table
1849 .get(count_key.as_slice())?
1850 .map(|g| postcard::from_bytes::<(usize, usize)>(g.value()).unwrap_or((0, 0)))
1851 .unwrap_or((0, 0));
1852
1853 let mut popped: Option<Vec<u8>> = None;
1854 if end - start >= limit {
1855 if pop_front_if_limited {
1856 let front_key = make_list_content_key(name, start);
1857 popped = table
1858 .remove(front_key.as_slice())?
1859 .map(|g| g.value().to_vec());
1860 start += 1;
1861 } else {
1862 return Err(anyhow!(
1863 "list '{}' is full, limit reached: {}",
1864 String::from_utf8_lossy(name),
1865 limit
1866 ));
1867 }
1868 }
1869
1870 let content_key = make_list_content_key(name, end);
1871 table.insert(content_key.as_slice(), val)?;
1872 end += 1;
1873
1874 let encoded = postcard::to_stdvec(&(start, end))?;
1875 table.insert(count_key.as_slice(), encoded.as_slice())?;
1876
1877 popped
1878 };
1879 txn.commit()?;
1880 Ok(result)
1881 }
1882
1883 fn exec_list_pop(db: &Database, name: &[u8]) -> Result<Option<Vec<u8>>> {
1884 #[cfg(feature = "ttl")]
1885 if Self::is_container_expired(db, name)? {
1886 return Ok(None);
1887 }
1888 let txn = db.begin_write()?;
1889 let result = {
1890 let mut table = txn.open_table(LIST_TABLE)?;
1891 let count_key = make_list_count_key(name);
1892 let (start, end) = table
1893 .get(count_key.as_slice())?
1894 .map(|g| postcard::from_bytes::<(usize, usize)>(g.value()).unwrap_or((0, 0)))
1895 .unwrap_or((0, 0));
1896
1897 if start >= end {
1898 None
1899 } else {
1900 let front_key = make_list_content_key(name, start);
1901 let val = table
1902 .remove(front_key.as_slice())?
1903 .map(|g| g.value().to_vec());
1904 let encoded = postcard::to_stdvec(&(start + 1, end))?;
1905 table.insert(count_key.as_slice(), encoded.as_slice())?;
1906 val
1907 }
1908 };
1909 txn.commit()?;
1910 Ok(result)
1911 }
1912
1913 fn exec_list_all(db: &Database, name: &[u8]) -> Result<Vec<Vec<u8>>> {
1914 let txn = db.begin_read()?;
1915 #[cfg(feature = "ttl")]
1916 if let Ok(table) = txn.open_table(KEY_EXPIRE_TABLE) {
1917 if let Some(guard) = table.get(name)? {
1918 let expire_at: TimestampMillis = postcard::from_bytes(guard.value())?;
1919 if expire_at <= timestamp_millis() {
1920 return Ok(Vec::new());
1921 }
1922 }
1923 }
1924 let table = txn.open_table(LIST_TABLE)?;
1925 let count_key = make_list_count_key(name);
1926 let (start, end) = table
1927 .get(count_key.as_slice())?
1928 .map(|g| postcard::from_bytes::<(usize, usize)>(g.value()).unwrap_or((0, 0)))
1929 .unwrap_or((0, 0));
1930
1931 let mut results = Vec::with_capacity(end.saturating_sub(start));
1932 for i in start..end {
1933 let ck = make_list_content_key(name, i);
1934 if let Some(val) = table.get(ck.as_slice())? {
1935 results.push(val.value().to_vec());
1936 }
1937 }
1938 Ok(results)
1939 }
1940
1941 fn exec_list_get_index(db: &Database, name: &[u8], idx: usize) -> Result<Option<Vec<u8>>> {
1942 #[cfg(feature = "ttl")]
1943 if Self::is_container_expired(db, name)? {
1944 return Ok(None);
1945 }
1946 let txn = db.begin_read()?;
1947 let table = txn.open_table(LIST_TABLE)?;
1948 let count_key = make_list_count_key(name);
1949 let (start, end) = table
1950 .get(count_key.as_slice())?
1951 .map(|g| postcard::from_bytes::<(usize, usize)>(g.value()).unwrap_or((0, 0)))
1952 .unwrap_or((0, 0));
1953
1954 let actual_idx = start + idx;
1955 if actual_idx >= end {
1956 return Ok(None);
1957 }
1958 let ck = make_list_content_key(name, actual_idx);
1959 match table.get(ck.as_slice())? {
1960 Some(guard) => Ok(Some(guard.value().to_vec())),
1961 None => Ok(None),
1962 }
1963 }
1964
1965 fn exec_list_len(db: &Database, name: &[u8]) -> Result<usize> {
1966 let txn = db.begin_read()?;
1967 #[cfg(feature = "ttl")]
1968 if let Ok(table) = txn.open_table(KEY_EXPIRE_TABLE) {
1969 if let Some(guard) = table.get(name)? {
1970 let expire_at: TimestampMillis = postcard::from_bytes(guard.value())?;
1971 if expire_at <= timestamp_millis() {
1972 return Ok(0);
1973 }
1974 }
1975 }
1976 let table = txn.open_table(LIST_TABLE)?;
1977 let count_key = make_list_count_key(name);
1978 match table.get(count_key.as_slice())? {
1979 Some(g) => {
1980 let (start, end) = postcard::from_bytes::<(usize, usize)>(g.value())?;
1981 Ok(end - start)
1982 }
1983 None => Ok(0),
1984 }
1985 }
1986
1987 fn exec_list_is_empty(db: &Database, name: &[u8]) -> Result<bool> {
1988 Ok(Self::exec_list_len(db, name)? == 0)
1989 }
1990
1991 fn exec_list_clear(db: &Database, name: &[u8]) -> Result<()> {
1992 let txn = db.begin_write()?;
1993 {
1994 let mut table = txn.open_table(LIST_TABLE)?;
1995 let content_prefix = make_list_content_prefix(name);
1996 let count_key = make_list_count_key(name);
1997
1998 {
2000 let extract = table.extract_from_if(content_prefix.as_slice().., |_, _| true)?;
2001 let _drained: Vec<_> = extract.collect::<Result<Vec<_>, _>>()?;
2002 }
2003 let _ = table.remove(count_key.as_slice())?;
2004 }
2005 txn.commit()?;
2006 Ok(())
2007 }
2008
2009 fn exec_list_iter(db: &Database, name: &[u8]) -> Result<Vec<Vec<u8>>> {
2010 #[cfg(feature = "ttl")]
2011 if Self::is_container_expired(db, name)? {
2012 return Ok(Vec::new());
2013 }
2014 let txn = db.begin_read()?;
2015 let table = txn.open_table(LIST_TABLE)?;
2016 let count_key = make_list_count_key(name);
2017 let (start, end) = table
2018 .get(count_key.as_slice())?
2019 .map(|g| postcard::from_bytes::<(usize, usize)>(g.value()).unwrap_or((0, 0)))
2020 .unwrap_or((0, 0));
2021
2022 let mut results = Vec::with_capacity(end.saturating_sub(start));
2023 for i in start..end {
2024 let ck = make_list_content_key(name, i);
2025 if let Some(val) = table.get(ck.as_slice())? {
2026 results.push(val.value().to_vec());
2027 }
2028 }
2029 Ok(results)
2030 }
2031
2032 #[cfg(feature = "ttl")]
2033 fn exec_list_expire_at(db: &Database, name: &[u8], at: TimestampMillis) -> Result<bool> {
2034 let txn = db.begin_write()?;
2035 {
2036 let list_table = txn.open_table(LIST_TABLE)?;
2037 let count_key = make_list_count_key(name);
2038 if list_table.get(count_key.as_slice())?.is_none() {
2039 return Ok(false);
2040 }
2041 Self::set_key_ttl_in_txn(&txn, name, at)?;
2042 }
2043 txn.commit()?;
2044 Ok(true)
2045 }
2046
2047 #[cfg(feature = "ttl")]
2048 fn exec_list_expire(db: &Database, name: &[u8], dur: TimestampMillis) -> Result<bool> {
2049 let at = timestamp_millis() + dur;
2050 Self::exec_list_expire_at(db, name, at)
2051 }
2052
2053 #[cfg(feature = "ttl")]
2054 fn exec_list_ttl(db: &Database, name: &[u8]) -> Result<Option<TimestampMillis>> {
2055 Self::exec_db_ttl(db, name)
2056 }
2057
2058 fn exec_list_is_expired(db: &Database, name: &[u8]) -> Result<bool> {
2059 #[cfg(feature = "ttl")]
2060 {
2061 match Self::exec_db_ttl(db, name)? {
2062 Some(ttl) if ttl > 0 => Ok(false),
2063 _ => Ok(true),
2064 }
2065 }
2066 #[cfg(not(feature = "ttl"))]
2067 {
2068 let _ = (db, name);
2069 Ok(false)
2070 }
2071 }
2072}
2073
2074enum Command {
2079 DBInsert {
2081 key: Vec<u8>,
2082 val: Vec<u8>,
2083 tx: oneshot::Sender<Result<()>>,
2084 },
2085 DBGet {
2086 key: Vec<u8>,
2087 tx: oneshot::Sender<Result<Option<Vec<u8>>>>,
2088 },
2089 DBRemove {
2090 key: Vec<u8>,
2091 tx: oneshot::Sender<Result<()>>,
2092 },
2093 DBContainsKey {
2094 key: Vec<u8>,
2095 tx: oneshot::Sender<Result<bool>>,
2096 },
2097 DBBatchInsert {
2098 key_vals: Vec<(Vec<u8>, Vec<u8>)>,
2099 tx: oneshot::Sender<Result<()>>,
2100 },
2101 DBBatchRemove {
2102 keys: Vec<Vec<u8>>,
2103 tx: oneshot::Sender<Result<()>>,
2104 },
2105 DBCounterIncr {
2106 key: Vec<u8>,
2107 increment: isize,
2108 tx: oneshot::Sender<Result<()>>,
2109 },
2110 DBCounterDecr {
2111 key: Vec<u8>,
2112 decrement: isize,
2113 tx: oneshot::Sender<Result<()>>,
2114 },
2115 DBCounterGet {
2116 key: Vec<u8>,
2117 tx: oneshot::Sender<Result<Option<isize>>>,
2118 },
2119 DBCounterSet {
2120 key: Vec<u8>,
2121 val: isize,
2122 tx: oneshot::Sender<Result<()>>,
2123 },
2124 DBLen {
2125 tx: oneshot::Sender<Result<usize>>,
2126 },
2127 DBSize {
2128 tx: oneshot::Sender<Result<usize>>,
2129 },
2130 DBInfo {
2131 tx: oneshot::Sender<Result<Value>>,
2132 },
2133 DBMapGet {
2135 key: Vec<u8>,
2136 tx: oneshot::Sender<Result<bool>>,
2137 },
2138 DBListGet {
2139 key: Vec<u8>,
2140 tx: oneshot::Sender<Result<bool>>,
2141 },
2142 DBMapIter {
2144 tx: IterResultSender,
2145 },
2146 DBListIter {
2147 tx: IterResultSender,
2148 },
2149 DBScanIter {
2150 pattern: Vec<u8>,
2151 tx: oneshot::Sender<Result<Vec<Vec<u8>>>>,
2152 },
2153 #[cfg(feature = "ttl")]
2155 DBExpireAt {
2156 key: Vec<u8>,
2157 at: TimestampMillis,
2158 tx: oneshot::Sender<Result<bool>>,
2159 },
2160 #[cfg(feature = "ttl")]
2161 DBExpire {
2162 key: Vec<u8>,
2163 dur: TimestampMillis,
2164 tx: oneshot::Sender<Result<bool>>,
2165 },
2166 #[cfg(feature = "ttl")]
2167 DBTtl {
2168 key: Vec<u8>,
2169 tx: oneshot::Sender<Result<Option<TimestampMillis>>>,
2170 },
2171 #[cfg(feature = "ttl")]
2172 #[allow(dead_code)]
2173 DBCleanup {
2174 tx: oneshot::Sender<Result<usize>>,
2175 },
2176 MapInsert {
2178 map: RedbStorageMap,
2179 key: Vec<u8>,
2180 val: Vec<u8>,
2181 tx: oneshot::Sender<Result<()>>,
2182 },
2183 MapGet {
2184 map: RedbStorageMap,
2185 key: Vec<u8>,
2186 tx: oneshot::Sender<Result<Option<Vec<u8>>>>,
2187 },
2188 MapRemove {
2189 map: RedbStorageMap,
2190 key: Vec<u8>,
2191 tx: oneshot::Sender<Result<()>>,
2192 },
2193 MapContainsKey {
2194 map: RedbStorageMap,
2195 key: Vec<u8>,
2196 tx: oneshot::Sender<Result<bool>>,
2197 },
2198 #[cfg(feature = "map_len")]
2199 MapLen {
2200 map: RedbStorageMap,
2201 tx: oneshot::Sender<Result<usize>>,
2202 },
2203 MapIsEmpty {
2204 map: RedbStorageMap,
2205 tx: oneshot::Sender<Result<bool>>,
2206 },
2207 MapClear {
2208 map: RedbStorageMap,
2209 tx: oneshot::Sender<Result<()>>,
2210 },
2211 MapRemoveAndFetch {
2212 map: RedbStorageMap,
2213 key: Vec<u8>,
2214 tx: oneshot::Sender<Result<Option<Vec<u8>>>>,
2215 },
2216 MapRemoveWithPrefix {
2217 map: RedbStorageMap,
2218 prefix: Vec<u8>,
2219 tx: oneshot::Sender<Result<()>>,
2220 },
2221 MapBatchInsert {
2222 map: RedbStorageMap,
2223 key_vals: Vec<(Vec<u8>, Vec<u8>)>,
2224 tx: oneshot::Sender<Result<()>>,
2225 },
2226 MapBatchRemove {
2227 map: RedbStorageMap,
2228 keys: Vec<Vec<u8>>,
2229 tx: oneshot::Sender<Result<()>>,
2230 },
2231 MapIter {
2232 map: RedbStorageMap,
2233 tx: IterResultSender,
2234 },
2235 MapKeyIter {
2236 map: RedbStorageMap,
2237 tx: oneshot::Sender<Result<Vec<Vec<u8>>>>,
2238 },
2239 MapPrefixIter {
2240 map: RedbStorageMap,
2241 prefix: Vec<u8>,
2242 tx: IterResultSender,
2243 },
2244 #[cfg(feature = "ttl")]
2245 MapExpireAt {
2246 map: RedbStorageMap,
2247 at: TimestampMillis,
2248 tx: oneshot::Sender<Result<bool>>,
2249 },
2250 #[cfg(feature = "ttl")]
2251 MapExpire {
2252 map: RedbStorageMap,
2253 dur: TimestampMillis,
2254 tx: oneshot::Sender<Result<bool>>,
2255 },
2256 #[cfg(feature = "ttl")]
2257 MapTTL {
2258 map: RedbStorageMap,
2259 tx: oneshot::Sender<Result<Option<TimestampMillis>>>,
2260 },
2261 #[allow(dead_code)]
2262 MapIsExpired {
2263 map: RedbStorageMap,
2264 tx: oneshot::Sender<Result<bool>>,
2265 },
2266 ListPush {
2268 list: RedbStorageList,
2269 val: Vec<u8>,
2270 tx: oneshot::Sender<Result<()>>,
2271 },
2272 ListPushs {
2273 list: RedbStorageList,
2274 vals: Vec<Vec<u8>>,
2275 tx: oneshot::Sender<Result<()>>,
2276 },
2277 ListPushLimit {
2278 list: RedbStorageList,
2279 val: Vec<u8>,
2280 limit: usize,
2281 pop_front_if_limited: bool,
2282 tx: oneshot::Sender<Result<Option<Vec<u8>>>>,
2283 },
2284 ListPop {
2285 list: RedbStorageList,
2286 tx: oneshot::Sender<Result<Option<Vec<u8>>>>,
2287 },
2288 ListAll {
2289 list: RedbStorageList,
2290 tx: oneshot::Sender<Result<Vec<Vec<u8>>>>,
2291 },
2292 ListGetIndex {
2293 list: RedbStorageList,
2294 idx: usize,
2295 tx: oneshot::Sender<Result<Option<Vec<u8>>>>,
2296 },
2297 ListLen {
2298 list: RedbStorageList,
2299 tx: oneshot::Sender<Result<usize>>,
2300 },
2301 ListIsEmpty {
2302 list: RedbStorageList,
2303 tx: oneshot::Sender<Result<bool>>,
2304 },
2305 ListClear {
2306 list: RedbStorageList,
2307 tx: oneshot::Sender<Result<()>>,
2308 },
2309 ListIter {
2310 list: RedbStorageList,
2311 tx: oneshot::Sender<Result<Vec<Vec<u8>>>>,
2312 },
2313 #[cfg(feature = "ttl")]
2314 ListExpireAt {
2315 list: RedbStorageList,
2316 at: TimestampMillis,
2317 tx: oneshot::Sender<Result<bool>>,
2318 },
2319 #[cfg(feature = "ttl")]
2320 ListExpire {
2321 list: RedbStorageList,
2322 dur: TimestampMillis,
2323 tx: oneshot::Sender<Result<bool>>,
2324 },
2325 #[cfg(feature = "ttl")]
2326 ListTTL {
2327 list: RedbStorageList,
2328 tx: oneshot::Sender<Result<Option<TimestampMillis>>>,
2329 },
2330 #[allow(dead_code)]
2331 ListIsExpired {
2332 list: RedbStorageList,
2333 tx: oneshot::Sender<Result<bool>>,
2334 },
2335 #[cfg(feature = "circuit-breaker")]
2337 DBInsertRaw {
2338 key: Vec<u8>,
2339 val: Vec<u8>,
2340 tx: oneshot::Sender<Result<()>>,
2341 },
2342 #[cfg(feature = "circuit-breaker")]
2343 DBGetRaw {
2344 key: Vec<u8>,
2345 tx: oneshot::Sender<Result<Option<Vec<u8>>>>,
2346 },
2347 #[cfg(feature = "circuit-breaker")]
2348 #[allow(dead_code)]
2349 DBRemoveRaw {
2350 key: Vec<u8>,
2351 tx: oneshot::Sender<Result<()>>,
2352 },
2353 #[allow(dead_code)]
2355 Shutdown,
2356}
2357
2358#[derive(Clone)]
2364pub struct RedbStorageMap {
2365 db: RedbStorageDB,
2366 name: Arc<Key>,
2367 #[allow(dead_code)]
2368 expire: Option<TimestampMillis>,
2369 empty: Arc<AtomicBool>,
2370}
2371
2372impl fmt::Debug for RedbStorageMap {
2373 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2374 f.debug_struct("RedbStorageMap")
2375 .field("name", &String::from_utf8_lossy(&self.name))
2376 .finish()
2377 }
2378}
2379
2380impl RedbStorageMap {
2381 fn new(db: RedbStorageDB, name: Key, expire: Option<TimestampMillis>) -> Self {
2382 RedbStorageMap {
2383 db,
2384 name: Arc::new(name),
2385 expire,
2386 empty: Arc::new(AtomicBool::new(true)),
2387 }
2388 }
2389}
2390
2391#[async_trait]
2392impl Map for RedbStorageMap {
2393 fn name(&self) -> &[u8] {
2394 &self.name
2395 }
2396
2397 async fn insert<K, V>(&self, key: K, val: &V) -> Result<()>
2398 where
2399 K: AsRef<[u8]> + Sync + Send,
2400 V: Serialize + Sync + Send + ?Sized,
2401 {
2402 let key_bytes = key.as_ref().to_vec();
2403 let val_bytes = postcard::to_stdvec(val)?;
2404 self.db
2405 .send_cmd(|tx| Command::MapInsert {
2406 map: self.clone(),
2407 key: key_bytes,
2408 val: val_bytes,
2409 tx,
2410 })
2411 .await
2412 }
2413
2414 async fn get<K, V>(&self, key: K) -> Result<Option<V>>
2415 where
2416 K: AsRef<[u8]> + Sync + Send,
2417 V: DeserializeOwned + Sync + Send,
2418 {
2419 let key_bytes = key.as_ref().to_vec();
2420 let result: Option<Vec<u8>> = self
2421 .db
2422 .send_cmd(|tx| Command::MapGet {
2423 map: self.clone(),
2424 key: key_bytes,
2425 tx,
2426 })
2427 .await?;
2428 match result {
2429 Some(bytes) => Ok(Some(postcard::from_bytes(&bytes)?)),
2430 None => Ok(None),
2431 }
2432 }
2433
2434 async fn remove<K>(&self, key: K) -> Result<()>
2435 where
2436 K: AsRef<[u8]> + Sync + Send,
2437 {
2438 let key_bytes = key.as_ref().to_vec();
2439 self.db
2440 .send_cmd(|tx| Command::MapRemove {
2441 map: self.clone(),
2442 key: key_bytes,
2443 tx,
2444 })
2445 .await
2446 }
2447
2448 async fn contains_key<K: AsRef<[u8]> + Sync + Send>(&self, key: K) -> Result<bool> {
2449 let key_bytes = key.as_ref().to_vec();
2450 self.db
2451 .send_cmd(|tx| Command::MapContainsKey {
2452 map: self.clone(),
2453 key: key_bytes,
2454 tx,
2455 })
2456 .await
2457 }
2458
2459 #[cfg(feature = "map_len")]
2460 async fn len(&self) -> Result<usize> {
2461 self.db
2462 .send_cmd(|tx| Command::MapLen {
2463 map: self.clone(),
2464 tx,
2465 })
2466 .await
2467 }
2468
2469 async fn is_empty(&self) -> Result<bool> {
2470 if !self.empty.load(Ordering::Acquire) {
2471 return Ok(false);
2472 }
2473 self.db
2474 .send_cmd(|tx| Command::MapIsEmpty {
2475 map: self.clone(),
2476 tx,
2477 })
2478 .await
2479 }
2480
2481 async fn clear(&self) -> Result<()> {
2482 let result: Result<()> = self
2483 .db
2484 .send_cmd(|tx| Command::MapClear {
2485 map: self.clone(),
2486 tx,
2487 })
2488 .await;
2489 if result.is_ok() {
2490 self.empty.store(true, Ordering::Release);
2491 }
2492 result
2493 }
2494
2495 async fn remove_and_fetch<K, V>(&self, key: K) -> Result<Option<V>>
2496 where
2497 K: AsRef<[u8]> + Sync + Send,
2498 V: DeserializeOwned + Sync + Send,
2499 {
2500 let key_bytes = key.as_ref().to_vec();
2501 let result: Option<Vec<u8>> = self
2502 .db
2503 .send_cmd(|tx| Command::MapRemoveAndFetch {
2504 map: self.clone(),
2505 key: key_bytes,
2506 tx,
2507 })
2508 .await?;
2509 match result {
2510 Some(bytes) => Ok(Some(postcard::from_bytes(&bytes)?)),
2511 None => Ok(None),
2512 }
2513 }
2514
2515 async fn remove_with_prefix<K>(&self, prefix: K) -> Result<()>
2516 where
2517 K: AsRef<[u8]> + Sync + Send,
2518 {
2519 let prefix_bytes = prefix.as_ref().to_vec();
2520 self.db
2521 .send_cmd(|tx| Command::MapRemoveWithPrefix {
2522 map: self.clone(),
2523 prefix: prefix_bytes,
2524 tx,
2525 })
2526 .await
2527 }
2528
2529 async fn batch_insert<V>(&self, key_vals: Vec<(Key, V)>) -> Result<()>
2530 where
2531 V: Serialize + Sync + Send,
2532 {
2533 let mut kvs = Vec::with_capacity(key_vals.len());
2534 for (k, v) in key_vals {
2535 let val_bytes = postcard::to_stdvec(&v)?;
2536 kvs.push((k, val_bytes));
2537 }
2538 self.db
2539 .send_cmd(|tx| Command::MapBatchInsert {
2540 map: self.clone(),
2541 key_vals: kvs,
2542 tx,
2543 })
2544 .await
2545 }
2546
2547 async fn batch_remove(&self, keys: Vec<Key>) -> Result<()> {
2548 self.db
2549 .send_cmd(|tx| Command::MapBatchRemove {
2550 map: self.clone(),
2551 keys,
2552 tx,
2553 })
2554 .await
2555 }
2556
2557 async fn iter<'a, V>(
2558 &'a mut self,
2559 ) -> Result<Box<dyn AsyncIterator<Item = IterItem<V>> + Send + 'a>>
2560 where
2561 V: DeserializeOwned + Sync + Send + 'a + 'static,
2562 {
2563 let entries: Vec<(Vec<u8>, Vec<u8>)> = self
2564 .db
2565 .send_cmd(|tx| Command::MapIter {
2566 map: self.clone(),
2567 tx,
2568 })
2569 .await?;
2570 let iter = RedbIter {
2571 entries,
2572 pos: 0,
2573 _phantom: std::marker::PhantomData,
2574 };
2575 Ok(Box::new(iter))
2576 }
2577
2578 async fn key_iter<'a>(
2579 &'a mut self,
2580 ) -> Result<Box<dyn AsyncIterator<Item = Result<Key>> + Send + 'a>> {
2581 let keys: Vec<Vec<u8>> = self
2582 .db
2583 .send_cmd(|tx| Command::MapKeyIter {
2584 map: self.clone(),
2585 tx,
2586 })
2587 .await?;
2588 let iter = RedbKeyIter { keys, pos: 0 };
2589 Ok(Box::new(iter))
2590 }
2591
2592 async fn prefix_iter<'a, P, V>(
2593 &'a mut self,
2594 prefix: P,
2595 ) -> Result<Box<dyn AsyncIterator<Item = IterItem<V>> + Send + 'a>>
2596 where
2597 P: AsRef<[u8]> + Send + Sync,
2598 V: DeserializeOwned + Sync + Send + 'a + 'static,
2599 {
2600 let prefix_bytes = prefix.as_ref().to_vec();
2601 let entries: Vec<(Vec<u8>, Vec<u8>)> = self
2602 .db
2603 .send_cmd(|tx| Command::MapPrefixIter {
2604 map: self.clone(),
2605 prefix: prefix_bytes,
2606 tx,
2607 })
2608 .await?;
2609 let iter = RedbPrefixIter {
2610 entries,
2611 pos: 0,
2612 _phantom: std::marker::PhantomData,
2613 };
2614 Ok(Box::new(iter))
2615 }
2616
2617 #[cfg(feature = "ttl")]
2618 async fn expire_at(&self, at: TimestampMillis) -> Result<bool> {
2619 self.db
2620 .send_cmd(|tx| Command::MapExpireAt {
2621 map: self.clone(),
2622 at,
2623 tx,
2624 })
2625 .await
2626 }
2627
2628 #[cfg(feature = "ttl")]
2629 async fn expire(&self, dur: TimestampMillis) -> Result<bool> {
2630 self.db
2631 .send_cmd(|tx| Command::MapExpire {
2632 map: self.clone(),
2633 dur,
2634 tx,
2635 })
2636 .await
2637 }
2638
2639 #[cfg(feature = "ttl")]
2640 async fn ttl(&self) -> Result<Option<TimestampMillis>> {
2641 self.db
2642 .send_cmd(|tx| Command::MapTTL {
2643 map: self.clone(),
2644 tx,
2645 })
2646 .await
2647 }
2648}
2649
2650#[derive(Clone)]
2656pub struct RedbStorageList {
2657 db: RedbStorageDB,
2658 name: Arc<Key>,
2659 #[allow(dead_code)]
2660 expire: Option<TimestampMillis>,
2661 empty: Arc<AtomicBool>,
2662}
2663
2664impl fmt::Debug for RedbStorageList {
2665 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2666 f.debug_struct("RedbStorageList")
2667 .field("name", &String::from_utf8_lossy(&self.name))
2668 .finish()
2669 }
2670}
2671
2672impl RedbStorageList {
2673 fn new(db: RedbStorageDB, name: Key, expire: Option<TimestampMillis>) -> Self {
2674 RedbStorageList {
2675 db,
2676 name: Arc::new(name),
2677 expire,
2678 empty: Arc::new(AtomicBool::new(true)),
2679 }
2680 }
2681}
2682
2683#[async_trait]
2684impl List for RedbStorageList {
2685 fn name(&self) -> &[u8] {
2686 &self.name
2687 }
2688
2689 async fn push<V>(&self, val: &V) -> Result<()>
2690 where
2691 V: Serialize + Sync + Send,
2692 {
2693 let val_bytes = postcard::to_stdvec(val)?;
2694 self.db
2695 .send_cmd(|tx| Command::ListPush {
2696 list: self.clone(),
2697 val: val_bytes,
2698 tx,
2699 })
2700 .await
2701 }
2702
2703 async fn pushs<V>(&self, vals: Vec<V>) -> Result<()>
2704 where
2705 V: Serialize + Sync + Send,
2706 {
2707 let mut val_bytes = Vec::with_capacity(vals.len());
2708 for v in vals {
2709 val_bytes.push(postcard::to_stdvec(&v)?);
2710 }
2711 self.db
2712 .send_cmd(|tx| Command::ListPushs {
2713 list: self.clone(),
2714 vals: val_bytes,
2715 tx,
2716 })
2717 .await
2718 }
2719
2720 async fn push_limit<V>(
2721 &self,
2722 val: &V,
2723 limit: usize,
2724 pop_front_if_limited: bool,
2725 ) -> Result<Option<V>>
2726 where
2727 V: Serialize + Sync + Send,
2728 V: DeserializeOwned,
2729 {
2730 let val_bytes = postcard::to_stdvec(val)?;
2731 let result: Option<Vec<u8>> = self
2732 .db
2733 .send_cmd(|tx| Command::ListPushLimit {
2734 list: self.clone(),
2735 val: val_bytes,
2736 limit,
2737 pop_front_if_limited,
2738 tx,
2739 })
2740 .await?;
2741 match result {
2742 Some(bytes) => Ok(Some(postcard::from_bytes(&bytes)?)),
2743 None => Ok(None),
2744 }
2745 }
2746
2747 async fn pop<V>(&self) -> Result<Option<V>>
2748 where
2749 V: DeserializeOwned + Sync + Send,
2750 {
2751 let result: Option<Vec<u8>> = self
2752 .db
2753 .send_cmd(|tx| Command::ListPop {
2754 list: self.clone(),
2755 tx,
2756 })
2757 .await?;
2758 match result {
2759 Some(bytes) => Ok(Some(postcard::from_bytes(&bytes)?)),
2760 None => Ok(None),
2761 }
2762 }
2763
2764 async fn all<V>(&self) -> Result<Vec<V>>
2765 where
2766 V: DeserializeOwned + Sync + Send,
2767 {
2768 let result: Vec<Vec<u8>> = self
2769 .db
2770 .send_cmd(|tx| Command::ListAll {
2771 list: self.clone(),
2772 tx,
2773 })
2774 .await?;
2775 let mut vals = Vec::with_capacity(result.len());
2776 for bytes in result {
2777 vals.push(postcard::from_bytes(&bytes)?);
2778 }
2779 Ok(vals)
2780 }
2781
2782 async fn get_index<V>(&self, idx: usize) -> Result<Option<V>>
2783 where
2784 V: DeserializeOwned + Sync + Send,
2785 {
2786 let result: Option<Vec<u8>> = self
2787 .db
2788 .send_cmd(|tx| Command::ListGetIndex {
2789 list: self.clone(),
2790 idx,
2791 tx,
2792 })
2793 .await?;
2794 match result {
2795 Some(bytes) => Ok(Some(postcard::from_bytes(&bytes)?)),
2796 None => Ok(None),
2797 }
2798 }
2799
2800 async fn len(&self) -> Result<usize> {
2801 self.db
2802 .send_cmd(|tx| Command::ListLen {
2803 list: self.clone(),
2804 tx,
2805 })
2806 .await
2807 }
2808
2809 async fn is_empty(&self) -> Result<bool> {
2810 if !self.empty.load(Ordering::Acquire) {
2811 return Ok(false);
2812 }
2813 self.db
2814 .send_cmd(|tx| Command::ListIsEmpty {
2815 list: self.clone(),
2816 tx,
2817 })
2818 .await
2819 }
2820
2821 async fn clear(&self) -> Result<()> {
2822 let result: Result<()> = self
2823 .db
2824 .send_cmd(|tx| Command::ListClear {
2825 list: self.clone(),
2826 tx,
2827 })
2828 .await;
2829 if result.is_ok() {
2830 self.empty.store(true, Ordering::Release);
2831 }
2832 result
2833 }
2834
2835 async fn iter<'a, V>(
2836 &'a mut self,
2837 ) -> Result<Box<dyn AsyncIterator<Item = Result<V>> + Send + 'a>>
2838 where
2839 V: DeserializeOwned + Sync + Send + 'a + 'static,
2840 {
2841 let entries: Vec<Vec<u8>> = self
2842 .db
2843 .send_cmd(|tx| Command::ListIter {
2844 list: self.clone(),
2845 tx,
2846 })
2847 .await?;
2848 let iter = RedbListValIter {
2849 entries,
2850 pos: 0,
2851 _phantom: std::marker::PhantomData,
2852 };
2853 Ok(Box::new(iter))
2854 }
2855
2856 #[cfg(feature = "ttl")]
2857 async fn expire_at(&self, at: TimestampMillis) -> Result<bool> {
2858 self.db
2859 .send_cmd(|tx| Command::ListExpireAt {
2860 list: self.clone(),
2861 at,
2862 tx,
2863 })
2864 .await
2865 }
2866
2867 #[cfg(feature = "ttl")]
2868 async fn expire(&self, dur: TimestampMillis) -> Result<bool> {
2869 self.db
2870 .send_cmd(|tx| Command::ListExpire {
2871 list: self.clone(),
2872 dur,
2873 tx,
2874 })
2875 .await
2876 }
2877
2878 #[cfg(feature = "ttl")]
2879 async fn ttl(&self) -> Result<Option<TimestampMillis>> {
2880 self.db
2881 .send_cmd(|tx| Command::ListTTL {
2882 list: self.clone(),
2883 tx,
2884 })
2885 .await
2886 }
2887}
2888
2889#[async_trait]
2894impl StorageDB for RedbStorageDB {
2895 type MapType = RedbStorageMap;
2896 type ListType = RedbStorageList;
2897
2898 async fn map<N: AsRef<[u8]> + Sync + Send>(
2899 &self,
2900 name: N,
2901 expire: Option<TimestampMillis>,
2902 ) -> Self::MapType {
2903 let name_bytes = name.as_ref().to_vec();
2904 let map = RedbStorageMap::new(self.clone(), name_bytes, expire);
2905 #[cfg(feature = "ttl")]
2906 if let Some(expire_ms) = expire {
2907 if let Err(e) = self
2908 .send_cmd(|tx| Command::DBExpireAt {
2909 key: map.name.to_vec(),
2910 at: timestamp_millis() + expire_ms,
2911 tx,
2912 })
2913 .await
2914 {
2915 log::warn!(
2916 "redb map '{}' expire_at failed: {e}",
2917 String::from_utf8_lossy(&map.name)
2918 );
2919 }
2920 }
2921 map
2922 }
2923
2924 async fn map_remove<K>(&self, name: K) -> Result<()>
2925 where
2926 K: AsRef<[u8]> + Sync + Send,
2927 {
2928 let name_bytes = name.as_ref().to_vec();
2930 let _ = self
2932 .send_cmd(|tx| Command::MapClear {
2933 map: RedbStorageMap::new(self.clone(), name_bytes.clone(), None),
2934 tx,
2935 })
2936 .await;
2937 Ok(())
2938 }
2939
2940 async fn map_contains_key<K: AsRef<[u8]> + Sync + Send>(&self, key: K) -> Result<bool> {
2941 let key_bytes = key.as_ref().to_vec();
2942 self.send_cmd(|tx| Command::DBMapGet { key: key_bytes, tx })
2943 .await
2944 }
2945
2946 async fn list<N: AsRef<[u8]> + Sync + Send>(
2947 &self,
2948 name: N,
2949 expire: Option<TimestampMillis>,
2950 ) -> Self::ListType {
2951 let name_bytes = name.as_ref().to_vec();
2952 let list = RedbStorageList::new(self.clone(), name_bytes, expire);
2953 #[cfg(feature = "ttl")]
2954 if let Some(expire_ms) = expire {
2955 if let Err(e) = self
2956 .send_cmd(|tx| Command::DBExpireAt {
2957 key: list.name.to_vec(),
2958 at: timestamp_millis() + expire_ms,
2959 tx,
2960 })
2961 .await
2962 {
2963 log::warn!(
2964 "redb list '{}' expire_at failed: {e}",
2965 String::from_utf8_lossy(&list.name)
2966 );
2967 }
2968 }
2969 list
2970 }
2971
2972 async fn list_remove<K>(&self, name: K) -> Result<()>
2973 where
2974 K: AsRef<[u8]> + Sync + Send,
2975 {
2976 let name_bytes = name.as_ref().to_vec();
2977 let _ = self
2978 .send_cmd(|tx| Command::ListClear {
2979 list: RedbStorageList::new(self.clone(), name_bytes.clone(), None),
2980 tx,
2981 })
2982 .await;
2983 Ok(())
2984 }
2985
2986 async fn list_contains_key<K: AsRef<[u8]> + Sync + Send>(&self, key: K) -> Result<bool> {
2987 let key_bytes = key.as_ref().to_vec();
2988 self.send_cmd(|tx| Command::DBListGet { key: key_bytes, tx })
2989 .await
2990 }
2991
2992 async fn insert<K, V>(&self, key: K, val: &V) -> Result<()>
2993 where
2994 K: AsRef<[u8]> + Sync + Send,
2995 V: Serialize + Sync + Send + ?Sized,
2996 {
2997 let key_bytes = key.as_ref().to_vec();
2998 let val_bytes = postcard::to_stdvec(val)?;
2999 self.send_cmd(|tx| Command::DBInsert {
3000 key: key_bytes,
3001 val: val_bytes,
3002 tx,
3003 })
3004 .await
3005 }
3006
3007 async fn get<K, V>(&self, key: K) -> Result<Option<V>>
3008 where
3009 K: AsRef<[u8]> + Sync + Send,
3010 V: DeserializeOwned + Sync + Send,
3011 {
3012 let key_bytes = key.as_ref().to_vec();
3013 let result: Option<Vec<u8>> = self
3014 .send_cmd(|tx| Command::DBGet { key: key_bytes, tx })
3015 .await?;
3016 match result {
3017 Some(bytes) => Ok(Some(postcard::from_bytes(&bytes)?)),
3018 None => Ok(None),
3019 }
3020 }
3021
3022 async fn remove<K>(&self, key: K) -> Result<()>
3023 where
3024 K: AsRef<[u8]> + Sync + Send,
3025 {
3026 let key_bytes = key.as_ref().to_vec();
3027 self.send_cmd(|tx| Command::DBRemove { key: key_bytes, tx })
3028 .await
3029 }
3030
3031 async fn contains_key<K: AsRef<[u8]> + Sync + Send>(&self, key: K) -> Result<bool> {
3032 let key_bytes = key.as_ref().to_vec();
3033 self.send_cmd(|tx| Command::DBContainsKey { key: key_bytes, tx })
3034 .await
3035 }
3036
3037 async fn batch_insert<V>(&self, key_vals: Vec<(Key, V)>) -> Result<()>
3038 where
3039 V: Serialize + Sync + Send,
3040 {
3041 let mut kvs = Vec::with_capacity(key_vals.len());
3042 for (k, v) in key_vals {
3043 let val_bytes = postcard::to_stdvec(&v)?;
3044 kvs.push((k, val_bytes));
3045 }
3046 self.send_cmd(|tx| Command::DBBatchInsert { key_vals: kvs, tx })
3047 .await
3048 }
3049
3050 async fn batch_remove(&self, keys: Vec<Key>) -> Result<()> {
3051 self.send_cmd(|tx| Command::DBBatchRemove { keys, tx })
3052 .await
3053 }
3054
3055 async fn counter_incr<K>(&self, key: K, increment: isize) -> Result<()>
3056 where
3057 K: AsRef<[u8]> + Sync + Send,
3058 {
3059 let key_bytes = key.as_ref().to_vec();
3060 self.send_cmd(|tx| Command::DBCounterIncr {
3061 key: key_bytes,
3062 increment,
3063 tx,
3064 })
3065 .await
3066 }
3067
3068 async fn counter_decr<K>(&self, key: K, decrement: isize) -> Result<()>
3069 where
3070 K: AsRef<[u8]> + Sync + Send,
3071 {
3072 let key_bytes = key.as_ref().to_vec();
3073 self.send_cmd(|tx| Command::DBCounterDecr {
3074 key: key_bytes,
3075 decrement,
3076 tx,
3077 })
3078 .await
3079 }
3080
3081 async fn counter_get<K>(&self, key: K) -> Result<Option<isize>>
3082 where
3083 K: AsRef<[u8]> + Sync + Send,
3084 {
3085 let key_bytes = key.as_ref().to_vec();
3086 self.send_cmd(|tx| Command::DBCounterGet { key: key_bytes, tx })
3087 .await
3088 }
3089
3090 async fn counter_set<K>(&self, key: K, val: isize) -> Result<()>
3091 where
3092 K: AsRef<[u8]> + Sync + Send,
3093 {
3094 let key_bytes = key.as_ref().to_vec();
3095 self.send_cmd(|tx| Command::DBCounterSet {
3096 key: key_bytes,
3097 val,
3098 tx,
3099 })
3100 .await
3101 }
3102
3103 #[cfg(feature = "len")]
3104 async fn len(&self) -> Result<usize> {
3105 self.send_cmd(|tx| Command::DBLen { tx }).await
3106 }
3107
3108 async fn db_size(&self) -> Result<usize> {
3109 self.send_cmd(|tx| Command::DBSize { tx }).await
3110 }
3111
3112 #[cfg(feature = "ttl")]
3113 async fn expire_at<K>(&self, key: K, at: TimestampMillis) -> Result<bool>
3114 where
3115 K: AsRef<[u8]> + Sync + Send,
3116 {
3117 let key_bytes = key.as_ref().to_vec();
3118 self.send_cmd(|tx| Command::DBExpireAt {
3119 key: key_bytes,
3120 at,
3121 tx,
3122 })
3123 .await
3124 }
3125
3126 #[cfg(feature = "ttl")]
3127 async fn expire<K>(&self, key: K, dur: TimestampMillis) -> Result<bool>
3128 where
3129 K: AsRef<[u8]> + Sync + Send,
3130 {
3131 let key_bytes = key.as_ref().to_vec();
3132 self.send_cmd(|tx| Command::DBExpire {
3133 key: key_bytes,
3134 dur,
3135 tx,
3136 })
3137 .await
3138 }
3139
3140 #[cfg(feature = "ttl")]
3141 async fn ttl<K>(&self, key: K) -> Result<Option<TimestampMillis>>
3142 where
3143 K: AsRef<[u8]> + Sync + Send,
3144 {
3145 let key_bytes = key.as_ref().to_vec();
3146 self.send_cmd(|tx| Command::DBTtl { key: key_bytes, tx })
3147 .await
3148 }
3149
3150 async fn map_iter<'a>(
3151 &'a mut self,
3152 ) -> Result<Box<dyn AsyncIterator<Item = Result<StorageMap>> + Send + 'a>> {
3153 let entries: Vec<(Vec<u8>, Vec<u8>)> =
3154 self.send_cmd(|tx| Command::DBMapIter { tx }).await?;
3155 let map = RedbMapIter {
3156 db: self.clone(),
3157 entries,
3158 pos: 0,
3159 };
3160 Ok(Box::new(map))
3161 }
3162
3163 async fn list_iter<'a>(
3164 &'a mut self,
3165 ) -> Result<Box<dyn AsyncIterator<Item = Result<StorageList>> + Send + 'a>> {
3166 let entries: Vec<(Vec<u8>, Vec<u8>)> =
3167 self.send_cmd(|tx| Command::DBListIter { tx }).await?;
3168 let list = RedbListIter {
3169 db: self.clone(),
3170 entries,
3171 pos: 0,
3172 };
3173 Ok(Box::new(list))
3174 }
3175
3176 async fn scan<'a, P>(
3177 &'a mut self,
3178 pattern: P,
3179 ) -> Result<Box<dyn AsyncIterator<Item = Result<Key>> + Send + 'a>>
3180 where
3181 P: AsRef<[u8]> + Sync + Send,
3182 {
3183 let pattern_bytes = pattern.as_ref().to_vec();
3184 let keys: Vec<Vec<u8>> = self
3185 .send_cmd(|tx| Command::DBScanIter {
3186 pattern: pattern_bytes.clone(),
3187 tx,
3188 })
3189 .await?;
3190 let pattern = Pattern::from(pattern_bytes.as_slice());
3191 let iter = AsyncDbKeyIter {
3192 keys,
3193 pos: 0,
3194 pattern,
3195 };
3196 Ok(Box::new(iter))
3197 }
3198
3199 async fn info(&self) -> Result<Value> {
3200 self.send_cmd(|tx| Command::DBInfo { tx }).await
3201 }
3202}
3203
3204#[cfg(feature = "circuit-breaker")]
3209impl RedbStorageDB {
3210 pub(crate) async fn insert_raw(&self, key: &[u8], val: &[u8]) -> Result<()> {
3211 let (tx, rx) = oneshot::channel();
3212 self.cmd_tx
3213 .send(Command::DBInsertRaw {
3214 key: key.to_vec(),
3215 val: val.to_vec(),
3216 tx,
3217 })
3218 .await
3219 .map_err(|_| anyhow!("redb command channel closed"))?;
3220 rx.await
3221 .map_err(|_| anyhow!("redb command response channel dropped"))??;
3222 Ok(())
3223 }
3224
3225 pub(crate) async fn get_raw(&self, key: &[u8]) -> Result<Option<Vec<u8>>> {
3226 let (tx, rx) = oneshot::channel();
3227 self.cmd_tx
3228 .send(Command::DBGetRaw {
3229 key: key.to_vec(),
3230 tx,
3231 })
3232 .await
3233 .map_err(|_| anyhow!("redb command channel closed"))?;
3234 rx.await
3235 .map_err(|_| anyhow!("redb command response channel dropped"))?
3236 }
3237
3238 #[allow(dead_code)]
3239 pub(crate) async fn remove_raw(&self, key: &[u8]) -> Result<()> {
3240 let (tx, rx) = oneshot::channel();
3241 self.cmd_tx
3242 .send(Command::DBRemoveRaw {
3243 key: key.to_vec(),
3244 tx,
3245 })
3246 .await
3247 .map_err(|_| anyhow!("redb command channel closed"))?;
3248 rx.await
3249 .map_err(|_| anyhow!("redb command response channel dropped"))??;
3250 Ok(())
3251 }
3252
3253 pub(crate) async fn batch_insert_raw(&self, key_vals: Vec<(Vec<u8>, Vec<u8>)>) -> Result<()> {
3254 self.send_cmd(|tx| Command::DBBatchInsert { key_vals, tx })
3255 .await
3256 }
3257}
3258
3259#[cfg(feature = "circuit-breaker")]
3260impl RedbStorageMap {
3261 pub(crate) async fn insert_raw(&self, key: &[u8], val: &[u8]) -> Result<()> {
3262 self.db
3263 .send_cmd(|tx| Command::MapInsert {
3264 map: self.clone(),
3265 key: key.to_vec(),
3266 val: val.to_vec(),
3267 tx,
3268 })
3269 .await
3270 }
3271
3272 pub(crate) async fn get_raw(&self, key: &[u8]) -> Result<Option<Vec<u8>>> {
3273 self.db
3274 .send_cmd(|tx| Command::MapGet {
3275 map: self.clone(),
3276 key: key.to_vec(),
3277 tx,
3278 })
3279 .await
3280 }
3281
3282 pub(crate) async fn remove_and_fetch_raw(&self, key: &[u8]) -> Result<Option<Vec<u8>>> {
3283 self.db
3284 .send_cmd(|tx| Command::MapRemoveAndFetch {
3285 map: self.clone(),
3286 key: key.to_vec(),
3287 tx,
3288 })
3289 .await
3290 }
3291
3292 pub(crate) async fn batch_insert_raw(&self, key_vals: Vec<(Vec<u8>, Vec<u8>)>) -> Result<()> {
3293 self.db
3294 .send_cmd(|tx| Command::MapBatchInsert {
3295 map: self.clone(),
3296 key_vals,
3297 tx,
3298 })
3299 .await
3300 }
3301}
3302
3303#[cfg(feature = "circuit-breaker")]
3304impl RedbStorageList {
3305 pub(crate) async fn push_raw(&self, val: &[u8]) -> Result<()> {
3306 self.db
3307 .send_cmd(|tx| Command::ListPush {
3308 list: self.clone(),
3309 val: val.to_vec(),
3310 tx,
3311 })
3312 .await
3313 }
3314
3315 pub(crate) async fn pushs_raw(&self, vals: Vec<Vec<u8>>) -> Result<()> {
3316 self.db
3317 .send_cmd(|tx| Command::ListPushs {
3318 list: self.clone(),
3319 vals,
3320 tx,
3321 })
3322 .await
3323 }
3324
3325 pub(crate) async fn pop_raw(&self) -> Result<Option<Vec<u8>>> {
3326 self.db
3327 .send_cmd(|tx| Command::ListPop {
3328 list: self.clone(),
3329 tx,
3330 })
3331 .await
3332 }
3333
3334 pub(crate) async fn all_raw(&self) -> Result<Vec<Vec<u8>>> {
3335 self.db
3336 .send_cmd(|tx| Command::ListAll {
3337 list: self.clone(),
3338 tx,
3339 })
3340 .await
3341 }
3342
3343 pub(crate) async fn get_index_raw(&self, idx: usize) -> Result<Option<Vec<u8>>> {
3344 self.db
3345 .send_cmd(|tx| Command::ListGetIndex {
3346 list: self.clone(),
3347 idx,
3348 tx,
3349 })
3350 .await
3351 }
3352
3353 pub(crate) async fn push_limit_raw(
3354 &self,
3355 val: &[u8],
3356 limit: usize,
3357 pop_front_if_limited: bool,
3358 ) -> Result<Option<Vec<u8>>> {
3359 self.db
3360 .send_cmd(|tx| Command::ListPushLimit {
3361 list: self.clone(),
3362 val: val.to_vec(),
3363 limit,
3364 pop_front_if_limited,
3365 tx,
3366 })
3367 .await
3368 }
3369}
3370
3371pub struct RedbIter<V> {
3377 entries: Vec<(Vec<u8>, Vec<u8>)>,
3378 pos: usize,
3379 _phantom: std::marker::PhantomData<V>,
3380}
3381
3382#[async_trait]
3383impl<V: DeserializeOwned + Send> AsyncIterator for RedbIter<V> {
3384 type Item = IterItem<V>;
3385 async fn next(&mut self) -> Option<Self::Item> {
3386 if self.pos >= self.entries.len() {
3387 return None;
3388 }
3389 let (key, val_bytes) = self.entries[self.pos].clone();
3390 self.pos += 1;
3391 match postcard::from_bytes(&val_bytes) {
3392 Ok(val) => Some(Ok((key, val))),
3393 Err(e) => Some(Err(e.into())),
3394 }
3395 }
3396}
3397
3398pub struct RedbKeyIter {
3400 keys: Vec<Vec<u8>>,
3401 pos: usize,
3402}
3403
3404#[async_trait]
3405impl AsyncIterator for RedbKeyIter {
3406 type Item = Result<Key>;
3407 async fn next(&mut self) -> Option<Self::Item> {
3408 if self.pos >= self.keys.len() {
3409 return None;
3410 }
3411 let key = self.keys[self.pos].clone();
3412 self.pos += 1;
3413 Some(Ok(key))
3414 }
3415}
3416
3417pub struct RedbPrefixIter<V> {
3419 entries: Vec<(Vec<u8>, Vec<u8>)>,
3420 pos: usize,
3421 _phantom: std::marker::PhantomData<V>,
3422}
3423
3424#[async_trait]
3425impl<V: DeserializeOwned + Send> AsyncIterator for RedbPrefixIter<V> {
3426 type Item = IterItem<V>;
3427 async fn next(&mut self) -> Option<Self::Item> {
3428 if self.pos >= self.entries.len() {
3429 return None;
3430 }
3431 let (key, val_bytes) = self.entries[self.pos].clone();
3432 self.pos += 1;
3433 match postcard::from_bytes(&val_bytes) {
3434 Ok(val) => Some(Ok((key, val))),
3435 Err(e) => Some(Err(e.into())),
3436 }
3437 }
3438}
3439
3440pub struct RedbMapIter {
3442 db: RedbStorageDB,
3443 entries: Vec<(Vec<u8>, Vec<u8>)>,
3444 pos: usize,
3445}
3446
3447#[async_trait]
3448impl AsyncIterator for RedbMapIter {
3449 type Item = Result<StorageMap>;
3450 async fn next(&mut self) -> Option<Self::Item> {
3451 if self.pos >= self.entries.len() {
3452 return None;
3453 }
3454 let (key, _) = self.entries[self.pos].clone();
3455 self.pos += 1;
3456 let map_name = map_count_key_to_name(&key).to_vec();
3457 let map = RedbStorageMap::new(self.db.clone(), map_name, None);
3458 Some(Ok(StorageMap::Redb(map)))
3459 }
3460}
3461
3462pub struct RedbListIter {
3464 db: RedbStorageDB,
3465 entries: Vec<(Vec<u8>, Vec<u8>)>,
3466 pos: usize,
3467}
3468
3469#[async_trait]
3470impl AsyncIterator for RedbListIter {
3471 type Item = Result<StorageList>;
3472 async fn next(&mut self) -> Option<Self::Item> {
3473 if self.pos >= self.entries.len() {
3474 return None;
3475 }
3476 let (key, _) = self.entries[self.pos].clone();
3477 self.pos += 1;
3478 let list_name = list_count_key_to_name(&key).to_vec();
3479 let list = RedbStorageList::new(self.db.clone(), list_name, None);
3480 Some(Ok(StorageList::Redb(list)))
3481 }
3482}
3483
3484pub struct AsyncDbKeyIter {
3486 keys: Vec<Vec<u8>>,
3487 pos: usize,
3488 pattern: Pattern,
3489}
3490
3491#[async_trait]
3492impl AsyncIterator for AsyncDbKeyIter {
3493 type Item = Result<Key>;
3494 async fn next(&mut self) -> Option<Self::Item> {
3495 while self.pos < self.keys.len() {
3496 let key = self.keys[self.pos].clone();
3497 self.pos += 1;
3498 if is_match(self.pattern.clone(), &key) {
3499 if key.starts_with(COUNTER_PREFIX) || key.starts_with(KEY_PREFIX) {
3501 continue;
3502 }
3503 return Some(Ok(key));
3504 }
3505 }
3506 None
3507 }
3508}
3509
3510pub struct RedbListValIter<V> {
3512 entries: Vec<Vec<u8>>,
3513 pos: usize,
3514 _phantom: std::marker::PhantomData<V>,
3515}
3516
3517#[async_trait]
3518impl<V: DeserializeOwned + Send> AsyncIterator for RedbListValIter<V> {
3519 type Item = Result<V>;
3520 async fn next(&mut self) -> Option<Self::Item> {
3521 if self.pos >= self.entries.len() {
3522 return None;
3523 }
3524 let bytes = self.entries[self.pos].clone();
3525 self.pos += 1;
3526 match postcard::from_bytes(&bytes) {
3527 Ok(val) => Some(Ok(val)),
3528 Err(e) => Some(Err(e.into())),
3529 }
3530 }
3531}