1use byteorder::{ByteOrder, LittleEndian, ReadBytesExt};
9use hashbrown::HashMap;
10use radix::RadixNum;
11use rocksdb::backup::{
12 BackupEngine as DBBackupEngine, BackupEngineOptions as DBBackupEngineOptions,
13 RestoreOptions as DBRestoreOptions,
14};
15use rocksdb::{
16 DB, DBCompactionStyle, DBCompressionType, Env as DBEnv, Error as DBError, FlushOptions,
17 WriteBatch, WriteOptions,
18};
19use std::fmt;
20use std::fs;
21use std::io::{self, Cursor};
22use std::path::{Path, PathBuf};
23use std::str;
24use std::sync::{Arc, Mutex, RwLock, RwLockReadGuard, RwLockWriteGuard};
25use std::thread;
26use std::time::{Duration, SystemTime};
27use std::vec::Drain;
28
29use crate::config::ConfigStoreKVDatabase;
30
31use super::generic::{
32 StoreGeneric, StoreGenericActionBuilder, StoreGenericBuilder, StoreGenericPool,
33};
34use super::identifiers::*;
35use super::item::StoreItemPart;
36use super::keyer::{StoreKeyerBuilder, StoreKeyerHasher, StoreKeyerKey, StoreKeyerPrefix};
37
38#[derive(Clone)]
41pub struct StoreKVPool {
42 pool: Arc<RwLock<HashMap<StoreKVKey, Arc<StoreKV>>>>,
43 kv_store_config: Arc<crate::config::ConfigStoreKV>,
44 store_access_lock: Arc<RwLock<()>>,
45 store_acquire_lock: Arc<Mutex<()>>,
46 store_flush_lock: Arc<Mutex<()>>,
47}
48
49pub struct StoreKVBuilder {
50 kv_store_config: Arc<crate::config::ConfigStoreKV>,
51}
52
53pub struct StoreKV {
54 database: DB,
55 last_used: RwLock<SystemTime>,
56 last_flushed: RwLock<SystemTime>,
57 pub lock: RwLock<bool>,
58 kv_store_config: Arc<crate::config::ConfigStoreKV>,
59}
60
61pub struct StoreKVActionBuilder<'build> {
62 pub kv_pool: &'build StoreKVPool,
63}
64
65pub struct StoreKVAction<'a> {
66 store: Option<Arc<StoreKV>>,
67 bucket: StoreItemPart<'a>,
68}
69
70#[derive(PartialEq, Eq, Hash, Clone, Copy)]
71pub struct StoreKVKey {
72 collection_hash: StoreKVAtom,
73}
74
75#[derive(PartialEq)]
76pub enum StoreKVAcquireMode {
77 Any,
78 OpenOnly,
79}
80
81type StoreKVAtom = u32;
82
83const ATOM_HASH_RADIX: usize = 16;
84
85impl StoreKVPool {
86 pub fn new(kv_store_config: Arc<crate::config::ConfigStoreKV>) -> Self {
87 Self {
88 pool: Arc::default(),
89 kv_store_config,
90 store_access_lock: Arc::default(),
91 store_acquire_lock: Arc::default(),
92 store_flush_lock: Arc::default(),
93 }
94 }
95
96 pub fn count(&self) -> usize {
97 self.pool.read().unwrap().len()
98 }
99
100 pub fn lock_read_access<'a>(&'a self) -> RwLockReadGuard<'a, ()> {
101 self.store_access_lock.read().unwrap()
102 }
103
104 pub fn lock_write_access<'a>(&'a self) -> RwLockWriteGuard<'a, ()> {
105 self.store_access_lock.write().unwrap()
106 }
107
108 pub fn acquire(
109 &self,
110 mode: StoreKVAcquireMode,
111 collection: impl AsRef<str>,
112 ) -> Result<Option<Arc<StoreKV>>, ()> {
113 let collection = collection.as_ref();
114 let pool_key = StoreKVKey::from_str(collection);
115
116 let _acquire = self.store_acquire_lock.lock().unwrap();
119
120 let store_pool_read = self.pool.read().unwrap();
122
123 if let Some(store_kv) = store_pool_read.get(&pool_key) {
124 Self::proceed_acquire_cache("kv", collection, pool_key, store_kv).map(Some)
125 } else {
126 tracing::info!(
127 "kv store not in pool for collection: {} {}, opening it",
128 collection,
129 pool_key
130 );
131
132 drop(store_pool_read);
135
136 let can_open_db = if mode == StoreKVAcquireMode::OpenOnly {
138 self.kv_store_config.path(pool_key.collection_hash).exists()
139 } else {
140 true
141 };
142
143 let builder = StoreKVBuilder {
144 kv_store_config: Arc::clone(&self.kv_store_config),
145 };
146
147 if can_open_db {
151 Self::proceed_acquire_open("kv", collection, pool_key, &self.pool, &builder)
152 .map(Some)
153 } else {
154 Ok(None)
155 }
156 }
157 }
158
159 fn close(&self, collection_hash: StoreKVAtom) {
160 tracing::debug!(
161 "closing key-value database for collection: <{:x}>",
162 collection_hash
163 );
164
165 let mut store_pool_write = self.pool.write().unwrap();
166
167 let collection_target = StoreKVKey::from_atom(collection_hash);
168
169 store_pool_write.remove(&collection_target);
170 }
171
172 pub fn janitor(&self) {
173 Self::proceed_janitor(
174 "kv",
175 &self.pool,
176 self.kv_store_config.pool.inactive_after,
177 &self.store_access_lock,
178 )
179 }
180
181 pub fn backup(&self, path: &Path) -> Result<(), io::Error> {
182 tracing::debug!("backing up all kv stores to path: {:?}", path);
183
184 fs::create_dir_all(path)?;
186
187 self.dump_action(
189 "backup",
190 &self.kv_store_config.path,
191 path,
192 &Self::backup_item,
193 )
194 }
195
196 pub fn restore(&self, path: &Path) -> Result<(), io::Error> {
197 tracing::debug!("restoring all kv stores from path: {:?}", path);
198
199 self.dump_action(
201 "restore",
202 path,
203 &self.kv_store_config.path,
204 &Self::restore_item,
205 )
206 }
207
208 pub fn flush(&self, force: bool) {
209 tracing::debug!("scanning for kv store pool items to flush to disk");
210
211 let _flush = self.store_flush_lock.lock().unwrap();
214
215 let mut keys_flush: Vec<StoreKVKey> = Vec::new();
217
218 {
219 let _access = self.store_access_lock.write().unwrap();
222
223 let store_pool_read = self.pool.read().unwrap();
224
225 for (key, store) in &*store_pool_read {
226 let not_flushed_for = store
231 .last_flushed
232 .read()
233 .unwrap()
234 .elapsed()
235 .unwrap_or_else(|err| {
236 tracing::error!(
237 "kv key: {} last flush duration clock issue, zeroing: {}",
238 key,
239 err
240 );
241
242 Duration::from_secs(0)
244 })
245 .as_secs();
246
247 if force || not_flushed_for >= self.kv_store_config.database.flush_after {
248 tracing::info!(
249 "kv key: {} not flushed for: {} seconds, may flush",
250 key,
251 not_flushed_for
252 );
253
254 keys_flush.push(*key);
255 } else {
256 tracing::debug!(
257 "kv key: {} not flushed for: {} seconds, no flush",
258 key,
259 not_flushed_for
260 );
261 }
262 }
263 }
264
265 if keys_flush.is_empty() {
267 tracing::info!("no kv store pool items need to be flushed at the moment");
268
269 return;
270 }
271
272 let mut count_flushed = 0;
274
275 {
276 for key in &keys_flush {
277 {
278 let _access = self.store_access_lock.write().unwrap();
281
282 if let Some(store) = self.pool.read().unwrap().get(key) {
283 tracing::debug!("kv key: {} flush started", key);
284
285 if let Err(err) = store.flush() {
286 tracing::error!("kv key: {} flush failed: {}", key, err);
287 } else {
288 count_flushed += 1;
289
290 tracing::debug!("kv key: {} flush complete", key);
291 }
292
293 *store.last_flushed.write().unwrap() = SystemTime::now();
295 }
296 }
297
298 thread::yield_now();
300 }
301 }
302
303 tracing::info!(
304 "done scanning for kv store pool items to flush to disk (flushed: {})",
305 count_flushed
306 );
307 }
308
309 #[allow(clippy::type_complexity)]
310 fn dump_action(
311 &self,
312 action: &str,
313 read_path: &Path,
314 write_path: &Path,
315 fn_item: &dyn Fn(&Self, &Path, &Path, &str) -> Result<(), io::Error>,
316 ) -> Result<(), io::Error> {
317 for collection in fs::read_dir(read_path)? {
319 let collection = collection?;
320
321 if let (Ok(collection_file_type), Some(collection_name)) =
323 (collection.file_type(), collection.file_name().to_str())
324 {
325 if collection_file_type.is_dir() {
326 tracing::debug!("kv collection ongoing {}: {}", action, collection_name);
327
328 fn_item(self, write_path, &collection.path(), collection_name)?;
329 }
330 }
331 }
332
333 Ok(())
334 }
335
336 fn backup_item(
337 &self,
338 backup_path: &Path,
339 _origin_path: &Path,
340 collection_name: &str,
341 ) -> Result<(), io::Error> {
342 let _access = self.store_access_lock.write().unwrap();
345
346 let kv_backup_path = backup_path.join(collection_name);
348
349 tracing::debug!(
350 "kv collection: {} backing up to path: {:?}",
351 collection_name,
352 kv_backup_path
353 );
354
355 if kv_backup_path.exists() {
357 fs::remove_dir_all(&kv_backup_path)?;
358 }
359
360 fs::create_dir_all(backup_path.join(collection_name))?;
362
363 if let Ok(collection_radix) = RadixNum::from_str(collection_name, ATOM_HASH_RADIX) {
366 if let Ok(collection_hash) = collection_radix.as_decimal() {
367 let origin_kv = StoreKVBuilder {
368 kv_store_config: Arc::clone(&self.kv_store_config),
369 }
370 .open(collection_hash as StoreKVAtom)
371 .map_err(|_| io::Error::other("database open failure"))?;
372
373 let kv_backup_options = DBBackupEngineOptions::new(&kv_backup_path)
375 .map_err(|_| io::Error::other("backup engine options acquire failure"))?;
376 let kv_backup_environment = DBEnv::new()
377 .map_err(|_| io::Error::other("backup engine environment acquire failure"))?;
378
379 let mut kv_backup_engine =
380 DBBackupEngine::open(&kv_backup_options, &kv_backup_environment)
381 .map_err(|_| io::Error::other("backup engine failure"))?;
382
383 kv_backup_engine
385 .create_new_backup(&origin_kv)
386 .map_err(|_| io::Error::other("database backup failure"))?;
387
388 tracing::info!(
389 "kv collection: {} backed up to path: {:?}",
390 collection_name,
391 kv_backup_path
392 );
393 }
394 }
395
396 Ok(())
397 }
398
399 fn restore_item(
400 &self,
401 _backup_path: &Path,
402 origin_path: &Path,
403 collection_name: &str,
404 ) -> Result<(), io::Error> {
405 let _access = self.store_access_lock.write().unwrap();
408
409 tracing::debug!(
410 "kv collection: {} restoring from path: {:?}",
411 collection_name,
412 origin_path
413 );
414
415 if let Ok(collection_radix) = RadixNum::from_str(collection_name, ATOM_HASH_RADIX) {
418 if let Ok(collection_hash) = collection_radix.as_decimal() {
419 self.close(collection_hash as StoreKVAtom);
421
422 let kv_path = self.kv_store_config.path(collection_hash as StoreKVAtom);
424
425 if kv_path.exists() {
427 fs::remove_dir_all(&kv_path)?;
428 }
429
430 fs::create_dir_all(&kv_path)?;
432
433 let kv_backup_options = DBBackupEngineOptions::new(&origin_path)
435 .map_err(|_| io::Error::other("backup engine options acquire failure"))?;
436 let kv_backup_environment = DBEnv::new()
437 .map_err(|_| io::Error::other("backup engine environment acquire failure"))?;
438
439 let mut kv_backup_engine =
440 DBBackupEngine::open(&kv_backup_options, &kv_backup_environment)
441 .map_err(|_| io::Error::other("backup engine failure"))?;
442
443 kv_backup_engine
444 .restore_from_latest_backup(&kv_path, &kv_path, &DBRestoreOptions::default())
445 .map_err(|_| io::Error::other("database restore failure"))?;
446
447 tracing::info!(
448 "kv collection: {} restored to path: {:?} from backup: {:?}",
449 collection_name,
450 kv_path,
451 origin_path
452 );
453 }
454 }
455
456 Ok(())
457 }
458}
459
460impl StoreGenericPool<StoreKVKey, StoreKV, StoreKVBuilder> for StoreKVPool {}
461
462impl StoreKVBuilder {
463 fn open(&self, collection_hash: StoreKVAtom) -> Result<DB, DBError> {
464 tracing::debug!(
465 "opening key-value database for collection: <{:x}>",
466 collection_hash
467 );
468
469 let db_options = self.configure();
471
472 DB::open(&db_options, self.kv_store_config.path(collection_hash))
474 }
475
476 #[rustfmt::skip]
477 fn configure(&self) -> rocksdb::Options {
478 tracing::debug!("configuring key-value database");
479
480 let ConfigStoreKVDatabase {
482 flush_after: _,
483 compress,
484 parallelism,
485 max_open_files,
486 max_flushes,
487 write_ahead_log: _,
488 write_buffer_size,
489 max_write_buffer_number,
490 min_write_buffer_number,
491 min_write_buffer_number_to_merge,
492 block_cache_size,
493 cache_index_and_filter_blocks,
494 compression_type,
495 wal_compression_type,
496 wal_ttl_seconds,
497 wal_size_limit_mb,
498 wal_bytes_per_sync,
499 wal_recovery_mode,
500 compression_level,
501 min_level_to_compress,
502 level_zero_file_num_compaction_trigger,
503 level_zero_slowdown_writes_trigger,
504 level_zero_stop_writes_trigger,
505 max_bytes_for_level_base,
506 max_bytes_for_level_multiplier,
507 target_file_size_base,
508 max_background_jobs,
509 max_subcompactions,
510 stats_dump_period_sec,
511 } = &self.kv_store_config.database;
512
513 let mut db_options = rocksdb::Options::default();
515
516 macro_rules! if_some {
517 ($opts:ident.$set_fn:ident($value:expr)) => {
518 if let Some(value) = $value {
519 $opts.$set_fn(*value);
520 }
521 };
522 }
523
524 db_options.create_if_missing(true);
526 db_options.set_use_fsync(false);
527 db_options.set_compaction_style(DBCompactionStyle::Level);
528
529 if_some!(db_options.set_write_buffer_size(write_buffer_size.map(|n| n * 1024).as_ref()));
531 if_some!(db_options.set_min_write_buffer_number(min_write_buffer_number));
532 if_some!(db_options.set_min_write_buffer_number_to_merge(min_write_buffer_number_to_merge));
533 if_some!(db_options.set_max_write_buffer_number(max_write_buffer_number));
534
535 if_some!(db_options.set_max_open_files(max_open_files));
536
537 if let Some(block_cache_size) = block_cache_size {
541 let cache = rocksdb::Cache::new_lru_cache((*block_cache_size as usize) * 1024 * 1024);
542 let mut block_opts = rocksdb::BlockBasedOptions::default();
543 block_opts.set_block_cache(&cache);
544 if_some!(block_opts.set_cache_index_and_filter_blocks(cache_index_and_filter_blocks));
545 db_options.set_block_based_table_factory(&block_opts);
546 }
547
548 if let Some(compress) = compress {
551 db_options.set_compression_type(if *compress {
552 DBCompressionType::Zstd
553 } else {
554 DBCompressionType::None
555 });
556 }
557 if_some!(db_options.set_compression_type(compression_type));
558 if let Some(compression_level) = compression_level {
559 db_options.set_compression_options(
560 -14,
561 *compression_level,
562 0,
563 0,
564 );
565 }
566
567 if_some!(db_options.set_wal_compression_type(wal_compression_type));
568 if_some!(db_options.set_wal_ttl_seconds(wal_ttl_seconds));
569 if_some!(db_options.set_wal_size_limit_mb(wal_size_limit_mb));
570 if_some!(db_options.set_wal_bytes_per_sync(wal_bytes_per_sync));
571 if_some!(db_options.set_wal_recovery_mode(wal_recovery_mode));
572
573 if_some!(db_options.set_min_level_to_compress(min_level_to_compress));
574
575 if_some!(db_options.set_level_zero_file_num_compaction_trigger(level_zero_file_num_compaction_trigger));
576 if_some!(db_options.set_level_zero_slowdown_writes_trigger(level_zero_slowdown_writes_trigger));
577 if_some!(db_options.set_level_zero_stop_writes_trigger(level_zero_stop_writes_trigger));
578
579 if_some!(db_options.set_max_bytes_for_level_base(max_bytes_for_level_base));
580 if_some!(db_options.set_max_bytes_for_level_multiplier(max_bytes_for_level_multiplier));
581 if_some!(db_options.set_target_file_size_base(target_file_size_base));
582
583 let mut max_background_jobs = *max_background_jobs;
584 if max_background_jobs.is_none() {
585 if let Some(max_flushes) = max_flushes {
586 max_background_jobs = Some((max_subcompactions.unwrap_or(1) + max_flushes) as i32);
587 }
588 }
589 if_some!(db_options.set_max_background_jobs(max_background_jobs.as_ref()));
590 if_some!(db_options.set_max_subcompactions(max_subcompactions));
591
592 if_some!(db_options.set_stats_dump_period_sec(stats_dump_period_sec));
593
594 if_some!(db_options.increase_parallelism(parallelism));
595
596 db_options
597 }
598}
599
600impl crate::config::ConfigStoreKV {
601 fn path(&self, collection_hash: StoreKVAtom) -> PathBuf {
602 self.path.join(format!("{:x}", collection_hash))
603 }
604}
605
606impl StoreGenericBuilder<StoreKVKey, StoreKV> for StoreKVBuilder {
607 fn build(&self, pool_key: StoreKVKey) -> Result<StoreKV, ()> {
608 self.open(pool_key.collection_hash)
609 .map(|db| {
610 let now = SystemTime::now();
611
612 StoreKV {
613 database: db,
614 last_used: RwLock::new(now),
615 last_flushed: RwLock::new(now),
616 lock: RwLock::new(false),
617 kv_store_config: Arc::clone(&self.kv_store_config),
618 }
619 })
620 .map_err(|err| {
621 tracing::error!("failed opening kv: {}", err);
622 })
623 }
624}
625
626impl StoreKV {
627 pub fn get(&self, key: &[u8]) -> Result<Option<Vec<u8>>, DBError> {
628 self.database.get(key)
629 }
630
631 pub fn put(&self, key: &[u8], data: &[u8]) -> Result<(), DBError> {
632 let mut batch = WriteBatch::default();
633
634 batch.put(key, data);
635
636 self.do_write(batch)
637 }
638
639 pub fn delete(&self, key: &[u8]) -> Result<(), DBError> {
640 let mut batch = WriteBatch::default();
641
642 batch.delete(key);
643
644 self.do_write(batch)
645 }
646
647 fn flush(&self) -> Result<(), DBError> {
648 let mut flush_options = FlushOptions::default();
650
651 flush_options.set_wait(true);
652
653 self.database.flush_opt(&flush_options)
655 }
656
657 fn do_write(&self, batch: WriteBatch) -> Result<(), DBError> {
658 let mut write_options = WriteOptions::default();
660
661 if !self.kv_store_config.database.write_ahead_log {
663 tracing::debug!("ignoring wal for kv write");
664
665 write_options.disable_wal(true);
666 } else {
667 tracing::debug!("using wal for kv write");
668
669 write_options.disable_wal(false);
670 }
671
672 self.database.write_opt(batch, &write_options)
674 }
675}
676
677impl StoreGeneric for StoreKV {
678 fn ref_last_used(&self) -> &RwLock<SystemTime> {
679 &self.last_used
680 }
681}
682
683impl<'build> StoreKVActionBuilder<'build> {
684 pub fn access(bucket: StoreItemPart, store: Option<Arc<StoreKV>>) -> StoreKVAction {
685 Self::build(bucket, store)
686 }
687
688 pub fn erase<T: AsRef<str>>(&self, collection: T, bucket: Option<T>) -> Result<u32, ()> {
689 self.dispatch_erase("kv", collection, bucket)
690 }
691
692 fn build(bucket: StoreItemPart, store: Option<Arc<StoreKV>>) -> StoreKVAction {
693 StoreKVAction { store, bucket }
694 }
695}
696
697impl<'build> StoreGenericActionBuilder for StoreKVActionBuilder<'build> {
698 fn proceed_erase_collection(&self, collection_str: &str) -> Result<u32, ()> {
699 let collection_atom = StoreKeyerHasher::to_compact(collection_str);
700 let collection_path = self.kv_pool.kv_store_config.path(collection_atom);
701
702 self.kv_pool.close(collection_atom);
704
705 if collection_path.exists() {
706 tracing::debug!(
707 "kv collection store exists, erasing: {}/* at path: {:?}",
708 collection_str,
709 &collection_path
710 );
711
712 let erase_result = fs::remove_dir_all(&collection_path);
714
715 if erase_result.is_ok() {
716 tracing::debug!("done with kv collection erasure");
717
718 Ok(1)
719 } else {
720 Err(())
721 }
722 } else {
723 tracing::debug!(
724 "kv collection store does not exist, consider already erased: {}/* at path: {:?}",
725 collection_str,
726 &collection_path
727 );
728
729 Ok(0)
730 }
731 }
732
733 fn proceed_erase_bucket(&self, _collection: &str, _bucket: &str) -> Result<u32, ()> {
734 Err(())
737 }
738}
739
740impl<'a> StoreKVAction<'a> {
741 pub fn get_meta_to_value(&self, meta: StoreMetaKey) -> Result<Option<StoreMetaValue>, ()> {
745 if let Some(ref store) = self.store {
746 let store_key = StoreKeyerBuilder::meta_to_value(self.bucket.as_str(), &meta);
747
748 tracing::debug!("store get meta-to-value: {}", store_key);
749
750 match store.get(&store_key.as_bytes()) {
751 Ok(Some(value)) => {
752 tracing::debug!("got meta-to-value: {}", store_key);
753
754 Ok(if let Ok(value) = str::from_utf8(&value) {
755 match meta {
756 StoreMetaKey::IIDIncr => value
757 .parse::<StoreObjectIID>()
758 .ok()
759 .map(StoreMetaValue::IIDIncr)
760 .or(None),
761 }
762 } else {
763 None
764 })
765 }
766 Ok(None) => {
767 tracing::debug!("no meta-to-value found: {}", store_key);
768
769 Ok(None)
770 }
771 Err(err) => {
772 tracing::error!(
773 "error getting meta-to-value: {} with trace: {}",
774 store_key,
775 err
776 );
777
778 Err(())
779 }
780 }
781 } else {
782 Ok(None)
783 }
784 }
785
786 pub fn set_meta_to_value(&self, meta: StoreMetaKey, value: StoreMetaValue) -> Result<(), ()> {
787 if let Some(ref store) = self.store {
788 let store_key = StoreKeyerBuilder::meta_to_value(self.bucket.as_str(), &meta);
789
790 tracing::debug!("store set meta-to-value: {}", store_key);
791
792 let value_string = match value {
793 StoreMetaValue::IIDIncr(iid_incr) => iid_incr.to_string(),
794 };
795
796 store
797 .put(&store_key.as_bytes(), value_string.as_bytes())
798 .or(Err(()))
799 } else {
800 Err(())
801 }
802 }
803
804 pub fn get_term_to_iids(
808 &self,
809 term_hashed: StoreTermHashed,
810 ) -> Result<Option<Vec<StoreObjectIID>>, ()> {
811 if let Some(ref store) = self.store {
812 let store_key = StoreKeyerBuilder::term_to_iids(self.bucket.as_str(), term_hashed);
813
814 tracing::debug!("store get term-to-iids: {}", store_key);
815
816 match store.get(&store_key.as_bytes()) {
817 Ok(Some(value)) => {
818 tracing::debug!(
819 "got term-to-iids: {} with encoded value: {:?}",
820 store_key,
821 &*value
822 );
823
824 Self::decode_u32_list(&*value)
825 .or(Err(()))
826 .map(|value_decoded| {
827 tracing::debug!(
828 "got term-to-iids: {} with decoded value: {:?}",
829 store_key,
830 &value_decoded
831 );
832
833 Some(value_decoded)
834 })
835 }
836 Ok(None) => {
837 tracing::debug!("no term-to-iids found: {}", store_key);
838
839 Ok(None)
840 }
841 Err(err) => {
842 tracing::error!(
843 "error getting term-to-iids: {} with trace: {}",
844 store_key,
845 err
846 );
847
848 Err(())
849 }
850 }
851 } else {
852 Ok(None)
853 }
854 }
855
856 pub fn set_term_to_iids(
857 &self,
858 term_hashed: StoreTermHashed,
859 iids: &[StoreObjectIID],
860 ) -> Result<(), ()> {
861 if let Some(ref store) = self.store {
862 let store_key = StoreKeyerBuilder::term_to_iids(self.bucket.as_str(), term_hashed);
863
864 tracing::debug!("store set term-to-iids: {}", store_key);
865
866 let iids_encoded = Self::encode_u32_list(iids);
868
869 tracing::debug!(
870 "store set term-to-iids: {} with encoded value: {:?}",
871 store_key,
872 iids_encoded
873 );
874
875 store.put(&store_key.as_bytes(), &iids_encoded).or(Err(()))
876 } else {
877 Err(())
878 }
879 }
880
881 pub fn delete_term_to_iids(&self, term_hashed: StoreTermHashed) -> Result<(), ()> {
882 if let Some(ref store) = self.store {
883 let store_key = StoreKeyerBuilder::term_to_iids(self.bucket.as_str(), term_hashed);
884
885 tracing::debug!("store delete term-to-iids: {}", store_key);
886
887 store.delete(&store_key.as_bytes()).or(Err(()))
888 } else {
889 Err(())
890 }
891 }
892
893 pub fn get_oid_to_iid(&self, oid: StoreObjectOID<'a>) -> Result<Option<StoreObjectIID>, ()> {
897 if let Some(ref store) = self.store {
898 let store_key = StoreKeyerBuilder::oid_to_iid(self.bucket.as_str(), oid);
899
900 tracing::debug!("store get oid-to-iid: {}", store_key);
901
902 match store.get(&store_key.as_bytes()) {
903 Ok(Some(value)) => {
904 tracing::debug!(
905 "got oid-to-iid: {} with encoded value: {:?}",
906 store_key,
907 &*value
908 );
909
910 Self::decode_u32(&*value).or(Err(())).map(|value_decoded| {
911 tracing::debug!(
912 "got oid-to-iid: {} with decoded value: {:?}",
913 store_key,
914 &value_decoded
915 );
916
917 Some(value_decoded)
918 })
919 }
920 Ok(None) => {
921 tracing::debug!("no oid-to-iid found: {}", store_key);
922
923 Ok(None)
924 }
925 Err(err) => {
926 tracing::error!(
927 "error getting oid-to-iid: {} with trace: {}",
928 store_key,
929 err
930 );
931
932 Err(())
933 }
934 }
935 } else {
936 Ok(None)
937 }
938 }
939
940 pub fn set_oid_to_iid(&self, oid: StoreObjectOID<'a>, iid: StoreObjectIID) -> Result<(), ()> {
941 if let Some(ref store) = self.store {
942 let store_key = StoreKeyerBuilder::oid_to_iid(self.bucket.as_str(), oid);
943
944 tracing::debug!("store set oid-to-iid: {}", store_key);
945
946 let iid_encoded = Self::encode_u32(iid);
948
949 tracing::debug!(
950 "store set oid-to-iid: {} with encoded value: {:?}",
951 store_key,
952 iid_encoded
953 );
954
955 store.put(&store_key.as_bytes(), &iid_encoded).or(Err(()))
956 } else {
957 Err(())
958 }
959 }
960
961 pub fn delete_oid_to_iid(&self, oid: StoreObjectOID<'a>) -> Result<(), ()> {
962 if let Some(ref store) = self.store {
963 let store_key = StoreKeyerBuilder::oid_to_iid(self.bucket.as_str(), oid);
964
965 tracing::debug!("store delete oid-to-iid: {}", store_key);
966
967 store.delete(&store_key.as_bytes()).or(Err(()))
968 } else {
969 Err(())
970 }
971 }
972
973 pub fn get_iid_to_oid(&self, iid: StoreObjectIID) -> Result<Option<String>, ()> {
977 if let Some(ref store) = self.store {
978 let store_key = StoreKeyerBuilder::iid_to_oid(self.bucket.as_str(), iid);
979
980 tracing::debug!("store get iid-to-oid: {}", store_key);
981
982 match store.get(&store_key.as_bytes()) {
983 Ok(Some(value)) => Ok(str::from_utf8(&value).ok().map(|value| value.to_string())),
984 Ok(None) => Ok(None),
985 Err(_) => Err(()),
986 }
987 } else {
988 Ok(None)
989 }
990 }
991
992 pub fn set_iid_to_oid(&self, iid: StoreObjectIID, oid: StoreObjectOID<'a>) -> Result<(), ()> {
993 if let Some(ref store) = self.store {
994 let store_key = StoreKeyerBuilder::iid_to_oid(self.bucket.as_str(), iid);
995
996 tracing::debug!("store set iid-to-oid: {}", store_key);
997
998 store.put(&store_key.as_bytes(), oid.as_bytes()).or(Err(()))
999 } else {
1000 Err(())
1001 }
1002 }
1003
1004 pub fn delete_iid_to_oid(&self, iid: StoreObjectIID) -> Result<(), ()> {
1005 if let Some(ref store) = self.store {
1006 let store_key = StoreKeyerBuilder::iid_to_oid(self.bucket.as_str(), iid);
1007
1008 tracing::debug!("store delete iid-to-oid: {}", store_key);
1009
1010 store.delete(&store_key.as_bytes()).or(Err(()))
1011 } else {
1012 Err(())
1013 }
1014 }
1015
1016 pub fn get_iid_to_terms(
1020 &self,
1021 iid: StoreObjectIID,
1022 ) -> Result<Option<Vec<StoreTermHashed>>, ()> {
1023 if let Some(ref store) = self.store {
1024 let store_key = StoreKeyerBuilder::iid_to_terms(self.bucket.as_str(), iid);
1025
1026 tracing::debug!("store get iid-to-terms: {}", store_key);
1027
1028 match store.get(&store_key.as_bytes()) {
1029 Ok(Some(value)) => {
1030 tracing::debug!(
1031 "got iid-to-terms: {} with encoded value: {:?}",
1032 store_key,
1033 &*value
1034 );
1035
1036 Self::decode_u32_list(&*value)
1037 .or(Err(()))
1038 .map(|value_decoded| {
1039 tracing::debug!(
1040 "got iid-to-terms: {} with decoded value: {:?}",
1041 store_key,
1042 &value_decoded
1043 );
1044
1045 if !value_decoded.is_empty() {
1046 Some(value_decoded)
1047 } else {
1048 None
1049 }
1050 })
1051 }
1052 Ok(None) => Ok(None),
1053 Err(_) => Err(()),
1054 }
1055 } else {
1056 Ok(None)
1057 }
1058 }
1059
1060 pub fn set_iid_to_terms(
1061 &self,
1062 iid: StoreObjectIID,
1063 terms_hashed: &[StoreTermHashed],
1064 ) -> Result<(), ()> {
1065 if let Some(ref store) = self.store {
1066 let store_key = StoreKeyerBuilder::iid_to_terms(self.bucket.as_str(), iid);
1067
1068 tracing::debug!("store set iid-to-terms: {}", store_key);
1069
1070 let terms_hashed_encoded = Self::encode_u32_list(terms_hashed);
1072
1073 tracing::debug!(
1074 "store set iid-to-terms: {} with encoded value: {:?}",
1075 store_key,
1076 terms_hashed_encoded
1077 );
1078
1079 store
1080 .put(&store_key.as_bytes(), &terms_hashed_encoded)
1081 .or(Err(()))
1082 } else {
1083 Err(())
1084 }
1085 }
1086
1087 pub fn delete_iid_to_terms(&self, iid: StoreObjectIID) -> Result<(), ()> {
1088 if let Some(ref store) = self.store {
1089 let store_key = StoreKeyerBuilder::iid_to_terms(self.bucket.as_str(), iid);
1090
1091 tracing::debug!("store delete iid-to-terms: {}", store_key);
1092
1093 store.delete(&store_key.as_bytes()).or(Err(()))
1094 } else {
1095 Err(())
1096 }
1097 }
1098
1099 pub fn batch_flush_bucket(
1100 &self,
1101 iid: StoreObjectIID,
1102 oid: StoreObjectOID<'a>,
1103 iid_terms_hashed: &[StoreTermHashed],
1104 ) -> Result<u32, ()> {
1105 let mut count = 0;
1106
1107 tracing::debug!(
1108 "store batch flush bucket: {} with hashed terms: {:?}",
1109 iid,
1110 iid_terms_hashed
1111 );
1112
1113 match (
1115 self.delete_oid_to_iid(oid),
1116 self.delete_iid_to_oid(iid),
1117 self.delete_iid_to_terms(iid),
1118 ) {
1119 (Ok(_), Ok(_), Ok(_)) => {
1120 for iid_term in iid_terms_hashed {
1122 if let Ok(Some(mut iid_term_iids)) = self.get_term_to_iids(*iid_term) {
1123 if iid_term_iids.contains(&iid) {
1124 count += 1;
1125
1126 iid_term_iids.retain(|cur_iid| cur_iid != &iid);
1128 }
1129
1130 let is_ok = if iid_term_iids.is_empty() {
1131 self.delete_term_to_iids(*iid_term).is_ok()
1132 } else {
1133 self.set_term_to_iids(*iid_term, &iid_term_iids).is_ok()
1134 };
1135
1136 if !is_ok {
1137 return Err(());
1138 }
1139 }
1140 }
1141
1142 Ok(count)
1143 }
1144 _ => Err(()),
1145 }
1146 }
1147
1148 pub fn batch_truncate_object(
1149 &self,
1150 term_hashed: StoreTermHashed,
1151 term_iids_drain: Drain<StoreObjectIID>,
1152 ) -> Result<u32, ()> {
1153 let mut count = 0;
1154
1155 for term_iid_drain in term_iids_drain {
1156 tracing::debug!("store batch truncate object iid: {}", term_iid_drain);
1157
1158 if let Ok(Some(mut term_iid_drain_terms)) = self.get_iid_to_terms(term_iid_drain) {
1160 count += 1;
1161
1162 term_iid_drain_terms.retain(|cur_term| cur_term != &term_hashed);
1163
1164 if term_iid_drain_terms.is_empty() {
1166 if let Ok(Some(term_iid_drain_oid)) = self.get_iid_to_oid(term_iid_drain) {
1168 if self
1169 .batch_flush_bucket(term_iid_drain, &term_iid_drain_oid, &Vec::new())
1170 .is_err()
1171 {
1172 tracing::error!(
1173 "failed executing store batch truncate object batch-flush-bucket"
1174 );
1175 }
1176 } else {
1177 tracing::error!("failed getting store batch truncate object iid-to-oid");
1178 }
1179 } else {
1180 if self
1182 .set_iid_to_terms(term_iid_drain, &term_iid_drain_terms)
1183 .is_err()
1184 {
1185 tracing::error!("failed setting store batch truncate object iid-to-terms");
1186 }
1187 }
1188 }
1189 }
1190
1191 Ok(count)
1192 }
1193
1194 pub fn batch_erase_bucket(&self) -> Result<u32, ()> {
1195 if let Some(ref store) = self.store {
1196 let (k_meta_to_value, k_term_to_iids, k_oid_to_iid, k_iid_to_oid, k_iid_to_terms) = (
1198 StoreKeyerBuilder::meta_to_value(self.bucket.as_str(), &StoreMetaKey::IIDIncr),
1199 StoreKeyerBuilder::term_to_iids(self.bucket.as_str(), 0),
1200 StoreKeyerBuilder::oid_to_iid(self.bucket.as_str(), ""),
1201 StoreKeyerBuilder::iid_to_oid(self.bucket.as_str(), 0),
1202 StoreKeyerBuilder::iid_to_terms(self.bucket.as_str(), 0),
1203 );
1204
1205 let key_prefixes: [StoreKeyerPrefix; 5] = [
1206 k_meta_to_value.as_prefix(),
1207 k_term_to_iids.as_prefix(),
1208 k_oid_to_iid.as_prefix(),
1209 k_iid_to_oid.as_prefix(),
1210 k_iid_to_terms.as_prefix(),
1211 ];
1212
1213 for key_prefix in &key_prefixes {
1215 tracing::debug!(
1216 "store batch erase bucket: {} for prefix: {:?}",
1217 self.bucket.as_str(),
1218 key_prefix
1219 );
1220
1221 let key_prefix_start: StoreKeyerKey = [
1224 key_prefix[0],
1225 key_prefix[1],
1226 key_prefix[2],
1227 key_prefix[3],
1228 key_prefix[4],
1229 0,
1230 0,
1231 0,
1232 0,
1233 ];
1234 let key_prefix_end: StoreKeyerKey = [
1235 key_prefix[0],
1236 key_prefix[1],
1237 key_prefix[2],
1238 key_prefix[3],
1239 key_prefix[4],
1240 255,
1241 255,
1242 255,
1243 255,
1244 ];
1245
1246 let mut batch = WriteBatch::default();
1248
1249 batch.delete_range(&key_prefix_start, &key_prefix_end);
1250
1251 if let Err(err) = store.do_write(batch) {
1253 tracing::error!(
1254 "failed in store batch erase bucket: {} with error: {}",
1255 self.bucket.as_str(),
1256 err
1257 );
1258 } else {
1259 store.delete(&key_prefix_end).ok();
1263
1264 tracing::debug!(
1265 "succeeded in store batch erase bucket: {}",
1266 self.bucket.as_str()
1267 );
1268 }
1269 }
1270
1271 tracing::info!(
1272 "done processing store batch erase bucket: {}",
1273 self.bucket.as_str()
1274 );
1275
1276 Ok(1)
1277 } else {
1278 Err(())
1279 }
1280 }
1281
1282 fn encode_u32(decoded: u32) -> [u8; 4] {
1283 let mut encoded = [0; 4];
1284
1285 LittleEndian::write_u32(&mut encoded, decoded);
1286
1287 encoded
1288 }
1289
1290 fn decode_u32(encoded: &[u8]) -> Result<u32, ()> {
1291 Cursor::new(encoded).read_u32::<LittleEndian>().or(Err(()))
1292 }
1293
1294 fn encode_u32_list(decoded: &[u32]) -> Vec<u8> {
1295 let mut encoded = Vec::with_capacity(decoded.len() * 4);
1298
1299 for decoded_item in decoded {
1300 encoded.extend(&Self::encode_u32(*decoded_item))
1301 }
1302
1303 encoded
1304 }
1305
1306 fn decode_u32_list(encoded: &[u8]) -> Result<Vec<u32>, ()> {
1307 let mut decoded = Vec::with_capacity(encoded.len() / 4);
1310
1311 for encoded_chunk in encoded.chunks(4) {
1312 if let Ok(decoded_chunk) = Self::decode_u32(encoded_chunk) {
1313 decoded.push(decoded_chunk);
1314 } else {
1315 return Err(());
1316 }
1317 }
1318
1319 Ok(decoded)
1320 }
1321}
1322
1323impl StoreKVKey {
1324 pub fn from_atom(collection_hash: StoreKVAtom) -> StoreKVKey {
1325 StoreKVKey { collection_hash }
1326 }
1327
1328 #[allow(clippy::should_implement_trait)]
1329 pub fn from_str(collection_str: &str) -> StoreKVKey {
1330 StoreKVKey {
1331 collection_hash: StoreKeyerHasher::to_compact(collection_str),
1332 }
1333 }
1334}
1335
1336impl fmt::Display for StoreKVKey {
1337 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
1338 write!(f, "<{:x}>", self.collection_hash)
1339 }
1340}
1341
1342#[cfg(test)]
1343mod tests {
1344 use super::*;
1345
1346 #[test]
1347 fn it_acquires_database() {
1348 let kv_store_config = test_kv_store_config();
1349 let kv_pool = StoreKVPool::new(kv_store_config);
1350
1351 assert!(kv_pool.acquire(StoreKVAcquireMode::Any, "c:test:1").is_ok());
1352 }
1353
1354 #[test]
1355 fn it_janitors_database() {
1356 let kv_store_config = test_kv_store_config();
1357 let kv_pool = StoreKVPool::new(kv_store_config);
1358
1359 kv_pool.janitor();
1360 }
1361
1362 #[test]
1363 fn it_proceeds_primitives() {
1364 let kv_store_config = test_kv_store_config();
1365 let kv_pool = StoreKVPool::new(kv_store_config);
1366
1367 let store = kv_pool
1368 .acquire(StoreKVAcquireMode::Any, "c:test:2")
1369 .unwrap()
1370 .unwrap();
1371
1372 assert!(store.get(&[0]).is_ok());
1373 assert!(store.put(&[0], &[1, 0, 0, 0]).is_ok());
1374 assert!(store.delete(&[0]).is_ok());
1375 }
1376
1377 #[test]
1378 fn it_proceeds_actions() {
1379 let kv_store_config = test_kv_store_config();
1380 let kv_pool = StoreKVPool::new(kv_store_config);
1381
1382 let store = kv_pool
1383 .acquire(StoreKVAcquireMode::Any, "c:test:3")
1384 .unwrap();
1385 let action =
1386 StoreKVActionBuilder::access(StoreItemPart::from_str("b:test:3").unwrap(), store);
1387
1388 assert!(action.get_meta_to_value(StoreMetaKey::IIDIncr).is_ok());
1389 assert!(
1390 action
1391 .set_meta_to_value(StoreMetaKey::IIDIncr, StoreMetaValue::IIDIncr(1))
1392 .is_ok()
1393 );
1394
1395 assert!(action.get_term_to_iids(1).is_ok());
1396 assert!(action.set_term_to_iids(1, &[0, 1, 2]).is_ok());
1397 assert!(action.delete_term_to_iids(1).is_ok());
1398
1399 assert!(action.get_oid_to_iid(&"s".to_string()).is_ok());
1400 assert!(action.set_oid_to_iid(&"s".to_string(), 4).is_ok());
1401 assert!(action.delete_oid_to_iid(&"s".to_string()).is_ok());
1402
1403 assert!(action.get_iid_to_oid(4).is_ok());
1404 assert!(action.set_iid_to_oid(4, &"s".to_string()).is_ok());
1405 assert!(action.delete_iid_to_oid(4).is_ok());
1406
1407 assert!(action.get_iid_to_terms(4).is_ok());
1408 assert!(action.set_iid_to_terms(4, &[45402]).is_ok());
1409 assert!(action.delete_iid_to_terms(4).is_ok());
1410 }
1411
1412 #[test]
1413 fn it_encodes_atom() {
1414 assert_eq!(StoreKVAction::encode_u32(0), [0, 0, 0, 0]);
1415 assert_eq!(StoreKVAction::encode_u32(1), [1, 0, 0, 0]);
1416 assert_eq!(StoreKVAction::encode_u32(45402), [90, 177, 0, 0]);
1417 }
1418
1419 #[test]
1420 fn it_decodes_atom() {
1421 assert_eq!(StoreKVAction::decode_u32(&[0, 0, 0, 0]), Ok(0));
1422 assert_eq!(StoreKVAction::decode_u32(&[1, 0, 0, 0]), Ok(1));
1423 assert_eq!(StoreKVAction::decode_u32(&[90, 177, 0, 0]), Ok(45402));
1424 }
1425
1426 #[test]
1427 fn it_encodes_atom_list() {
1428 assert_eq!(
1429 StoreKVAction::encode_u32_list(&[0, 2, 3]),
1430 [0, 0, 0, 0, 2, 0, 0, 0, 3, 0, 0, 0]
1431 );
1432 assert_eq!(StoreKVAction::encode_u32_list(&[45402]), [90, 177, 0, 0]);
1433 }
1434
1435 #[test]
1436 fn it_decodes_atom_list() {
1437 assert_eq!(
1438 StoreKVAction::decode_u32_list(&[0, 0, 0, 0, 2, 0, 0, 0, 3, 0, 0, 0]),
1439 Ok(vec![0, 2, 3])
1440 );
1441 assert_eq!(
1442 StoreKVAction::decode_u32_list(&[90, 177, 0, 0]),
1443 Ok(vec![45402])
1444 );
1445 }
1446
1447 fn test_kv_store_config() -> Arc<crate::config::ConfigStoreKV> {
1448 Arc::new(
1449 config::Config::builder()
1450 .add_source(config::File::from_str(
1451 crate::config::tests::defaults_toml(),
1452 config::FileFormat::Toml,
1453 ))
1454 .build()
1455 .unwrap()
1456 .get::<crate::config::ConfigStoreKV>("store.kv")
1457 .unwrap(),
1458 )
1459 }
1460}
1461
1462#[cfg(all(feature = "benchmark", test))]
1463mod benches {
1464 extern crate test;
1465
1466 use super::*;
1467 use test::Bencher;
1468
1469 #[bench]
1470 fn bench_encode_atom(b: &mut Bencher) {
1471 b.iter(|| StoreKVAction::encode_u32(0));
1472 }
1473
1474 #[bench]
1475 fn bench_decode_atom(b: &mut Bencher) {
1476 let encoded_atom = [0, 0, 0, 0];
1477
1478 b.iter(|| StoreKVAction::decode_u32(&encoded_atom));
1479 }
1480
1481 #[bench]
1482 fn bench_encode_atom_list(b: &mut Bencher) {
1483 let atom_list = [0, 2, 3];
1484
1485 b.iter(|| StoreKVAction::encode_u32_list(&atom_list));
1486 }
1487
1488 #[bench]
1489 fn bench_decode_atom_list(b: &mut Bencher) {
1490 let encoded_atom_list = [0, 0, 0, 0, 2, 0, 0, 0, 3, 0, 0, 0];
1491
1492 b.iter(|| StoreKVAction::decode_u32_list(&encoded_atom_list));
1493 }
1494}
1495
1496impl fmt::Debug for StoreKVPool {
1499 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1500 use crate::util::fmt::{AsPrettyMutex, AsPrettyRwLock};
1501
1502 let Self {
1504 pool,
1505 store_access_lock,
1506 store_acquire_lock,
1507 store_flush_lock,
1508 kv_store_config: _kv_store_config,
1511 } = self;
1512
1513 f.debug_struct("StoreKVPool")
1514 .field("pool", &AsPrettyRwLock(pool))
1515 .field("store_access_lock", &AsPrettyRwLock(store_access_lock))
1516 .field("store_acquire_lock", &AsPrettyMutex(store_acquire_lock))
1517 .field("store_flush_lock", &AsPrettyMutex(store_flush_lock))
1518 .finish_non_exhaustive()
1519 }
1520}
1521
1522impl fmt::Debug for StoreKVKey {
1523 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1524 fmt::Display::fmt(&self, f)
1525 }
1526}
1527
1528impl fmt::Debug for StoreKV {
1529 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1530 use crate::util::fmt::AsPrettyRwLock;
1531
1532 let Self {
1534 database,
1535 last_used,
1536 last_flushed,
1537 lock,
1538 kv_store_config: _kv_store_config,
1541 } = self;
1542
1543 f.debug_struct("StoreKV")
1544 .field("database", database)
1545 .field("last_used", &AsPrettyRwLock(last_used))
1546 .field("last_flushed", &AsPrettyRwLock(last_flushed))
1547 .field("lock", &AsPrettyRwLock(lock))
1548 .finish_non_exhaustive()
1549 }
1550}