1#![allow(
9 clippy::cast_possible_truncation,
10 clippy::cast_sign_loss,
11 clippy::cast_possible_wrap
12)]
13
14use std::{
15 ffi::OsString,
16 fmt::Display,
17 path::PathBuf,
18 sync::{
19 atomic::{AtomicBool, Ordering},
20 Arc,
21 },
22};
23
24use linera_base::ensure;
25use rocksdb::{BlockBasedOptions, Cache, DBCompactionStyle, SliceTransform, WriteBufferManager};
26use serde::{Deserialize, Serialize};
27use sysinfo::{MemoryRefreshKind, RefreshKind, System};
28use tempfile::TempDir;
29use thiserror::Error;
30
31#[cfg(with_metrics)]
32use crate::metering::MeteredDatabase;
33#[cfg(with_testing)]
34use crate::store::TestKeyValueDatabase;
35use crate::{
36 batch::{Batch, WriteOperation},
37 common::get_upper_bound_option,
38 lru_caching::{LruCachingConfig, LruCachingDatabase},
39 store::{
40 KeyValueDatabase, KeyValueStoreError, ReadableKeyValueStore, WithError,
41 WritableKeyValueStore,
42 },
43};
44
45static ROOT_KEY_DOMAIN: [u8; 1] = [0];
47static STORED_ROOT_KEYS_PREFIX: u8 = 1;
48
49const MAX_VALUE_SIZE: usize = 3 * 1024 * 1024 * 1024 - 400;
52
53const MAX_KEY_SIZE: usize = 8 * 1024 * 1024 - 400;
56
57const WRITE_BUFFER_SIZE: usize = 16 * 1024 * 1024; const MAX_WRITE_BUFFER_NUMBER: i32 = 6;
63
64fn get_available_memory(sys: &System) -> usize {
65 sys.cgroup_limits()
66 .map_or_else(|| sys.total_memory() as usize, |c| c.total_memory as usize)
67}
68
69fn get_available_cpus() -> i32 {
70 std::thread::available_parallelism().map_or(1, |p| p.get() as i32)
71}
72
73const HYPER_CLOCK_CACHE_BLOCK_SIZE: usize = 8 * 1024; type DB = rocksdb::DBWithThreadMode<rocksdb::MultiThreaded>;
77
78#[derive(Clone, Copy, Debug, PartialEq, Eq, Deserialize, Serialize)]
85pub enum RocksDbSpawnMode {
86 SpawnBlocking,
88 BlockInPlace,
90}
91
92impl RocksDbSpawnMode {
93 pub fn get_spawn_mode_from_runtime() -> Self {
95 if tokio::runtime::Handle::current().metrics().num_workers() > 1 {
96 RocksDbSpawnMode::BlockInPlace
97 } else {
98 RocksDbSpawnMode::SpawnBlocking
99 }
100 }
101
102 #[inline]
104 async fn spawn<F, I, O>(&self, f: F, input: I) -> Result<O, RocksDbStoreInternalError>
105 where
106 F: FnOnce(I) -> Result<O, RocksDbStoreInternalError> + Send + 'static,
107 I: Send + 'static,
108 O: Send + 'static,
109 {
110 Ok(match self {
111 RocksDbSpawnMode::BlockInPlace => tokio::task::block_in_place(move || f(input))?,
112 RocksDbSpawnMode::SpawnBlocking => {
113 tokio::task::spawn_blocking(move || f(input)).await??
114 }
115 })
116 }
117}
118
119impl Display for RocksDbSpawnMode {
120 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
121 match &self {
122 RocksDbSpawnMode::SpawnBlocking => write!(f, "spawn_blocking"),
123 RocksDbSpawnMode::BlockInPlace => write!(f, "block_in_place"),
124 }
125 }
126}
127
128fn check_key_size(key: &[u8]) -> Result<(), RocksDbStoreInternalError> {
129 ensure!(
130 key.len() <= MAX_KEY_SIZE,
131 RocksDbStoreInternalError::KeyTooLong
132 );
133 Ok(())
134}
135
136#[derive(Clone)]
137struct RocksDbStoreExecutor {
138 db: Arc<DB>,
139 start_key: Vec<u8>,
140}
141
142impl RocksDbStoreExecutor {
143 fn contains_keys_internal(
144 &self,
145 keys: Vec<Vec<u8>>,
146 ) -> Result<Vec<bool>, RocksDbStoreInternalError> {
147 let size = keys.len();
148 let mut results = vec![false; size];
149 let mut indices = Vec::new();
150 let mut keys_red = Vec::new();
151 for (i, key) in keys.into_iter().enumerate() {
152 check_key_size(&key)?;
153 let mut full_key = self.start_key.to_vec();
154 full_key.extend(key);
155 if self.db.key_may_exist(&full_key) {
156 indices.push(i);
157 keys_red.push(full_key);
158 }
159 }
160 let values_red = self.db.multi_get(keys_red);
161 for (index, value) in indices.into_iter().zip(values_red) {
162 results[index] = value?.is_some();
163 }
164 Ok(results)
165 }
166
167 fn read_multi_values_bytes_internal(
168 &self,
169 keys: Vec<Vec<u8>>,
170 ) -> Result<Vec<Option<Vec<u8>>>, RocksDbStoreInternalError> {
171 for key in &keys {
172 check_key_size(key)?;
173 }
174 let full_keys = keys
175 .into_iter()
176 .map(|key| {
177 let mut full_key = self.start_key.to_vec();
178 full_key.extend(key);
179 full_key
180 })
181 .collect::<Vec<_>>();
182 let entries = self.db.multi_get(&full_keys);
183 Ok(entries.into_iter().collect::<Result<_, _>>()?)
184 }
185
186 fn get_find_prefix_iterator(
187 &self,
188 prefix: &[u8],
189 ) -> rocksdb::DBRawIteratorWithThreadMode<'_, DB> {
190 let mut read_opts = rocksdb::ReadOptions::default();
192 read_opts.set_async_io(true);
194
195 let upper_bound = get_upper_bound_option(prefix);
197 if let Some(upper_bound) = upper_bound {
198 read_opts.set_iterate_upper_bound(upper_bound);
199 }
200
201 let mut iter = self.db.raw_iterator_opt(read_opts);
202 iter.seek(prefix);
203 iter
204 }
205
206 fn find_keys_by_prefix_internal(
207 &self,
208 key_prefix: Vec<u8>,
209 ) -> Result<Vec<Vec<u8>>, RocksDbStoreInternalError> {
210 check_key_size(&key_prefix)?;
211
212 let mut prefix = self.start_key.clone();
213 prefix.extend(key_prefix);
214 let len = prefix.len();
215
216 let mut iter = self.get_find_prefix_iterator(&prefix);
217 let mut keys = Vec::new();
218 while let Some(key) = iter.key() {
219 keys.push(key[len..].to_vec());
220 iter.next();
221 }
222 Ok(keys)
223 }
224
225 #[expect(clippy::type_complexity)]
226 fn find_key_values_by_prefix_internal(
227 &self,
228 key_prefix: Vec<u8>,
229 ) -> Result<Vec<(Vec<u8>, Vec<u8>)>, RocksDbStoreInternalError> {
230 check_key_size(&key_prefix)?;
231 let mut prefix = self.start_key.clone();
232 prefix.extend(key_prefix);
233 let len = prefix.len();
234
235 let mut iter = self.get_find_prefix_iterator(&prefix);
236 let mut key_values = Vec::new();
237 while let Some((key, value)) = iter.item() {
238 let key_value = (key[len..].to_vec(), value.to_vec());
239 key_values.push(key_value);
240 iter.next();
241 }
242 Ok(key_values)
243 }
244
245 fn write_batch_internal(
246 &self,
247 batch: Batch,
248 write_root_key: bool,
249 ) -> Result<(), RocksDbStoreInternalError> {
250 let mut inner_batch = rocksdb::WriteBatchWithTransaction::default();
251 for operation in batch.operations {
252 match operation {
253 WriteOperation::Delete { key } => {
254 check_key_size(&key)?;
255 let mut full_key = self.start_key.to_vec();
256 full_key.extend(key);
257 inner_batch.delete(&full_key)
258 }
259 WriteOperation::Put { key, value } => {
260 check_key_size(&key)?;
261 let mut full_key = self.start_key.to_vec();
262 full_key.extend(key);
263 inner_batch.put(&full_key, value)
264 }
265 WriteOperation::DeletePrefix { key_prefix } => {
266 check_key_size(&key_prefix)?;
267 let mut full_key1 = self.start_key.to_vec();
268 full_key1.extend(&key_prefix);
269 let full_key2 =
270 get_upper_bound_option(&full_key1).expect("the first entry cannot be 255");
271 inner_batch.delete_range(&full_key1, &full_key2);
272 }
273 }
274 }
275 if write_root_key {
276 let mut full_key = self.start_key.to_vec();
277 full_key[0] = STORED_ROOT_KEYS_PREFIX;
278 inner_batch.put(&full_key, vec![]);
279 }
280 self.db.write(inner_batch)?;
281 Ok(())
282 }
283}
284
285#[derive(Clone)]
287pub struct RocksDbStoreInternal {
288 executor: RocksDbStoreExecutor,
289 path_with_guard: PathWithGuard,
290 spawn_mode: RocksDbSpawnMode,
291 root_key_written: Arc<AtomicBool>,
292}
293
294#[derive(Clone)]
296pub struct RocksDbDatabaseInternal {
297 executor: RocksDbStoreExecutor,
298 path_with_guard: PathWithGuard,
299 spawn_mode: RocksDbSpawnMode,
300}
301
302impl WithError for RocksDbDatabaseInternal {
303 type Error = RocksDbStoreInternalError;
304}
305
306#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Deserialize, Serialize, strum::EnumString)]
311#[serde(rename_all = "kebab-case")]
312#[strum(serialize_all = "kebab-case")]
313pub enum RocksDbStatisticsLevel {
314 DisableAll,
316 #[default]
318 ExceptHistogramOrTimers,
319 ExceptTimers,
321 ExceptDetailedTimers,
323 ExceptTimeForMutex,
325 All,
327}
328
329impl RocksDbStatisticsLevel {
330 fn to_rocksdb(self) -> rocksdb::statistics::StatsLevel {
331 use rocksdb::statistics::StatsLevel;
332 match self {
333 Self::DisableAll => StatsLevel::DisableAll,
334 Self::ExceptHistogramOrTimers => StatsLevel::ExceptHistogramOrTimers,
335 Self::ExceptTimers => StatsLevel::ExceptTimers,
336 Self::ExceptDetailedTimers => StatsLevel::ExceptDetailedTimers,
337 Self::ExceptTimeForMutex => StatsLevel::ExceptTimeForMutex,
338 Self::All => StatsLevel::All,
339 }
340 }
341}
342
343#[cfg(test)]
344mod statistics_level_tests {
345 use std::str::FromStr as _;
346
347 use super::RocksDbStatisticsLevel;
348
349 #[test]
350 fn parses_kebab_case_names() {
351 let cases = [
352 ("disable-all", RocksDbStatisticsLevel::DisableAll),
353 (
354 "except-histogram-or-timers",
355 RocksDbStatisticsLevel::ExceptHistogramOrTimers,
356 ),
357 ("except-timers", RocksDbStatisticsLevel::ExceptTimers),
358 (
359 "except-detailed-timers",
360 RocksDbStatisticsLevel::ExceptDetailedTimers,
361 ),
362 (
363 "except-time-for-mutex",
364 RocksDbStatisticsLevel::ExceptTimeForMutex,
365 ),
366 ("all", RocksDbStatisticsLevel::All),
367 ];
368 for (name, expected) in cases {
369 assert_eq!(RocksDbStatisticsLevel::from_str(name), Ok(expected));
370 }
371 assert!(RocksDbStatisticsLevel::from_str("not-a-level").is_err());
372 }
373}
374
375#[derive(Clone, Debug, Deserialize, Serialize)]
377pub struct RocksDbStoreInternalConfig {
378 pub path_with_guard: PathWithGuard,
380 pub spawn_mode: RocksDbSpawnMode,
382 #[serde(default)]
386 pub enable_statistics: bool,
387 #[serde(default)]
389 pub statistics_level: RocksDbStatisticsLevel,
390}
391
392impl RocksDbDatabaseInternal {
393 fn check_namespace(namespace: &str) -> Result<(), RocksDbStoreInternalError> {
394 if !namespace
395 .chars()
396 .all(|character| character.is_ascii_alphanumeric() || character == '_')
397 {
398 return Err(RocksDbStoreInternalError::InvalidNamespace);
399 }
400 Ok(())
401 }
402
403 fn build(
404 config: &RocksDbStoreInternalConfig,
405 namespace: &str,
406 ) -> Result<RocksDbDatabaseInternal, RocksDbStoreInternalError> {
407 let start_key = ROOT_KEY_DOMAIN.to_vec();
408 let temp_store = RocksDbStoreInternal::build(config, namespace, start_key)?;
410 Ok(RocksDbDatabaseInternal {
411 executor: temp_store.executor,
412 path_with_guard: temp_store.path_with_guard,
413 spawn_mode: temp_store.spawn_mode,
414 })
415 }
416}
417
418impl RocksDbStoreInternal {
419 fn build(
420 config: &RocksDbStoreInternalConfig,
421 namespace: &str,
422 start_key: Vec<u8>,
423 ) -> Result<RocksDbStoreInternal, RocksDbStoreInternalError> {
424 RocksDbDatabaseInternal::check_namespace(namespace)?;
425 let mut path_buf = config.path_with_guard.path_buf.clone();
426 let mut path_with_guard = config.path_with_guard.clone();
427 path_buf.push(namespace);
428 path_with_guard.path_buf = path_buf.clone();
429 let spawn_mode = config.spawn_mode;
430 if !std::path::Path::exists(&path_buf) {
431 std::fs::create_dir_all(path_buf.clone())?;
432 }
433 let sys = System::new_with_specifics(
434 RefreshKind::nothing().with_memory(MemoryRefreshKind::nothing().with_ram()),
435 );
436 let num_cpus = get_available_cpus();
437 let total_ram = get_available_memory(&sys);
438
439 let mut options = rocksdb::Options::default();
440 options.create_if_missing(true);
441 options.create_missing_column_families(true);
442
443 options.set_write_buffer_size(WRITE_BUFFER_SIZE);
445 options.set_max_write_buffer_number(MAX_WRITE_BUFFER_NUMBER);
446 options.set_compression_type(rocksdb::DBCompressionType::Lz4);
447 options.set_level_zero_slowdown_writes_trigger(8);
448 options.set_level_zero_stop_writes_trigger(12);
449 options.set_level_zero_file_num_compaction_trigger(2);
450 options.increase_parallelism(num_cpus);
454 options.set_max_background_jobs(num_cpus);
455 options.set_max_subcompactions(num_cpus as u32);
456 options.set_level_compaction_dynamic_level_bytes(true);
457
458 options.set_compaction_style(DBCompactionStyle::Level);
459 options.set_target_file_size_base(2 * WRITE_BUFFER_SIZE as u64);
460
461 let mut block_options = BlockBasedOptions::default();
462 block_options.set_pin_l0_filter_and_index_blocks_in_cache(true);
463 block_options.set_cache_index_and_filter_blocks(true);
464 block_options.set_block_cache(&Cache::new_hyper_clock_cache(
470 total_ram / 4,
471 HYPER_CLOCK_CACHE_BLOCK_SIZE,
472 ));
473
474 let write_buffer_manager =
477 WriteBufferManager::new_write_buffer_manager(total_ram / 4, true);
478 options.set_write_buffer_manager(&write_buffer_manager);
479
480 block_options.set_bloom_filter(10.0, false);
482 block_options.set_whole_key_filtering(false);
483
484 block_options.set_block_size(32 * 1024);
486 block_options.set_format_version(5);
488
489 options.set_block_based_table_factory(&block_options);
490
491 let prefix_extractor = SliceTransform::create_fixed_prefix(8);
494 options.set_prefix_extractor(prefix_extractor);
495
496 options.set_memtable_prefix_bloom_ratio(0.125);
498 options.set_optimize_filters_for_hits(true);
500 options.set_allow_mmap_reads(true);
502 options.set_advise_random_on_open(false);
504
505 if config.enable_statistics {
506 options.enable_statistics();
507 options.set_statistics_level(config.statistics_level.to_rocksdb());
508 }
509
510 let db = Arc::new(DB::open(&options, path_buf)?);
511 #[cfg(with_metrics)]
512 if config.enable_statistics {
513 statistics_metrics::register(Arc::new(options), db.clone());
514 }
515 let executor = RocksDbStoreExecutor { db, start_key };
516 Ok(RocksDbStoreInternal {
517 executor,
518 path_with_guard,
519 spawn_mode,
520 root_key_written: Arc::new(AtomicBool::new(false)),
521 })
522 }
523}
524
525#[cfg(with_metrics)]
531mod statistics_metrics {
532 use std::sync::{Arc, OnceLock};
533
534 use prometheus::{
535 core::{Collector, Desc},
536 proto::MetricFamily,
537 IntGauge,
538 };
539 use rocksdb::{statistics::Ticker, Options};
540
541 use super::DB;
542
543 enum Source {
544 Ticker(Ticker),
545 Property(&'static str),
546 }
547
548 struct Entry {
549 source: Source,
550 gauge: IntGauge,
551 }
552
553 fn definitions() -> Vec<(&'static str, &'static str, Source)> {
554 vec![
555 (
556 "linera_rocksdb_block_cache_hit",
557 "Cumulative RocksDB block cache hits since open",
558 Source::Ticker(Ticker::BlockCacheHit),
559 ),
560 (
561 "linera_rocksdb_block_cache_miss",
562 "Cumulative RocksDB block cache misses since open",
563 Source::Ticker(Ticker::BlockCacheMiss),
564 ),
565 (
566 "linera_rocksdb_compact_read_bytes",
567 "Cumulative bytes read during compaction since open",
568 Source::Ticker(Ticker::CompactReadBytes),
569 ),
570 (
571 "linera_rocksdb_compact_write_bytes",
572 "Cumulative bytes written during compaction since open",
573 Source::Ticker(Ticker::CompactWriteBytes),
574 ),
575 (
576 "linera_rocksdb_flush_write_bytes",
577 "Cumulative bytes written during flushes since open",
578 Source::Ticker(Ticker::FlushWriteBytes),
579 ),
580 (
581 "linera_rocksdb_stall_micros",
582 "Cumulative write-stall time in microseconds since open",
583 Source::Ticker(Ticker::StallMicros),
584 ),
585 (
586 "linera_rocksdb_bytes_written",
587 "Cumulative user bytes written since open",
588 Source::Ticker(Ticker::BytesWritten),
589 ),
590 (
591 "linera_rocksdb_bytes_read",
592 "Cumulative user bytes read since open",
593 Source::Ticker(Ticker::BytesRead),
594 ),
595 (
596 "linera_rocksdb_wal_bytes",
597 "Cumulative bytes written to the write-ahead log since open",
598 Source::Ticker(Ticker::WalFileBytes),
599 ),
600 (
601 "linera_rocksdb_bloom_filter_useful",
602 "Cumulative count of reads avoided by the bloom filter since open",
603 Source::Ticker(Ticker::BloomFilterUseful),
604 ),
605 (
606 "linera_rocksdb_memtable_hit",
607 "Cumulative memtable hits since open",
608 Source::Ticker(Ticker::MemtableHit),
609 ),
610 (
611 "linera_rocksdb_memtable_miss",
612 "Cumulative memtable misses since open",
613 Source::Ticker(Ticker::MemtableMiss),
614 ),
615 (
616 "linera_rocksdb_number_keys_written",
617 "Cumulative number of keys written since open",
618 Source::Ticker(Ticker::NumberKeysWritten),
619 ),
620 (
621 "linera_rocksdb_num_files_at_level0",
622 "Number of files at level 0",
623 Source::Property("rocksdb.num-files-at-level0"),
624 ),
625 (
626 "linera_rocksdb_estimate_pending_compaction_bytes",
627 "Estimated bytes pending compaction",
628 Source::Property("rocksdb.estimate-pending-compaction-bytes"),
629 ),
630 (
631 "linera_rocksdb_num_running_compactions",
632 "Number of currently running compactions",
633 Source::Property("rocksdb.num-running-compactions"),
634 ),
635 (
636 "linera_rocksdb_num_running_flushes",
637 "Number of currently running flushes",
638 Source::Property("rocksdb.num-running-flushes"),
639 ),
640 (
641 "linera_rocksdb_is_write_stopped",
642 "Whether writes are currently stopped (1) or not (0)",
643 Source::Property("rocksdb.is-write-stopped"),
644 ),
645 (
646 "linera_rocksdb_actual_delayed_write_rate",
647 "Current delayed write rate in bytes/s (0 when not delayed)",
648 Source::Property("rocksdb.actual-delayed-write-rate"),
649 ),
650 (
651 "linera_rocksdb_cur_size_all_mem_tables",
652 "Approximate size in bytes of all active and unflushed memtables",
653 Source::Property("rocksdb.cur-size-all-mem-tables"),
654 ),
655 (
656 "linera_rocksdb_num_immutable_mem_table",
657 "Number of immutable memtables not yet flushed",
658 Source::Property("rocksdb.num-immutable-mem-table"),
659 ),
660 (
661 "linera_rocksdb_live_sst_files_size",
662 "Total size in bytes of all live SST files",
663 Source::Property("rocksdb.live-sst-files-size"),
664 ),
665 (
666 "linera_rocksdb_total_sst_files_size",
667 "Total size in bytes of all SST files including obsolete ones",
668 Source::Property("rocksdb.total-sst-files-size"),
669 ),
670 (
671 "linera_rocksdb_estimate_num_keys",
672 "Estimated number of keys in the database",
673 Source::Property("rocksdb.estimate-num-keys"),
674 ),
675 (
676 "linera_rocksdb_block_cache_usage",
677 "Memory in bytes used by the block cache",
678 Source::Property("rocksdb.block-cache-usage"),
679 ),
680 (
681 "linera_rocksdb_block_cache_capacity",
682 "Capacity in bytes of the block cache",
683 Source::Property("rocksdb.block-cache-capacity"),
684 ),
685 ]
686 }
687
688 struct RocksDbStatisticsCollector {
689 options: Arc<Options>,
690 db: Arc<DB>,
691 entries: Vec<Entry>,
692 }
693
694 impl RocksDbStatisticsCollector {
695 fn new(options: Arc<Options>, db: Arc<DB>) -> Self {
696 let entries = definitions()
697 .into_iter()
698 .map(|(name, help, source)| Entry {
699 source,
700 gauge: IntGauge::new(name, help)
701 .expect("RocksDB statistics metric name is valid"),
702 })
703 .collect();
704 Self {
705 options,
706 db,
707 entries,
708 }
709 }
710 }
711
712 impl Collector for RocksDbStatisticsCollector {
713 fn desc(&self) -> Vec<&Desc> {
714 self.entries
715 .iter()
716 .flat_map(|entry| entry.gauge.desc())
717 .collect()
718 }
719
720 fn collect(&self) -> Vec<MetricFamily> {
721 self.entries
722 .iter()
723 .flat_map(|entry| {
724 let value = match &entry.source {
725 Source::Ticker(ticker) => self.options.get_ticker_count(*ticker) as i64,
726 Source::Property(property) => {
727 self.db
728 .property_int_value(*property)
729 .ok()
730 .flatten()
731 .unwrap_or(0) as i64
732 }
733 };
734 entry.gauge.set(value);
735 entry.gauge.collect()
736 })
737 .collect()
738 }
739 }
740
741 pub(super) fn register(options: Arc<Options>, db: Arc<DB>) {
742 static REGISTERED: OnceLock<()> = OnceLock::new();
743 if REGISTERED.set(()).is_err() {
744 tracing::warn!(
745 "RocksDB statistics collector is already registered; skipping additional store"
746 );
747 return;
748 }
749 let collector = RocksDbStatisticsCollector::new(options, db);
750 if let Err(error) = prometheus::register(Box::new(collector)) {
751 tracing::warn!("failed to register the RocksDB statistics collector: {error}");
752 }
753 }
754
755 #[cfg(test)]
756 mod tests {
757 use std::collections::HashSet;
758
759 use super::{definitions, IntGauge};
760
761 #[test]
762 fn definitions_build_unique_valid_gauges() {
763 let definitions = definitions();
764 assert!(!definitions.is_empty());
765 let mut names = HashSet::new();
766 for (name, help, _source) in &definitions {
767 assert!(!help.is_empty(), "metric {name} has empty help text");
768 assert!(names.insert(*name), "duplicate metric name: {name}");
769 IntGauge::new(*name, *help).expect("metric definition should be valid");
770 }
771 }
772 }
773}
774
775impl WithError for RocksDbStoreInternal {
776 type Error = RocksDbStoreInternalError;
777}
778
779impl ReadableKeyValueStore for RocksDbStoreInternal {
780 const MAX_KEY_SIZE: usize = MAX_KEY_SIZE;
781
782 fn root_key(&self) -> Result<Vec<u8>, RocksDbStoreInternalError> {
783 assert!(self.executor.start_key.starts_with(&ROOT_KEY_DOMAIN));
784 let root_key = bcs::from_bytes(&self.executor.start_key[ROOT_KEY_DOMAIN.len()..])?;
785 Ok(root_key)
786 }
787
788 async fn read_value_bytes(
789 &self,
790 key: &[u8],
791 ) -> Result<Option<Vec<u8>>, RocksDbStoreInternalError> {
792 check_key_size(key)?;
793 let db = self.executor.db.clone();
794 let mut full_key = self.executor.start_key.to_vec();
795 full_key.extend(key);
796 self.spawn_mode
797 .spawn(move |x| Ok(db.get(&x)?), full_key)
798 .await
799 }
800
801 async fn contains_key(&self, key: &[u8]) -> Result<bool, RocksDbStoreInternalError> {
802 check_key_size(key)?;
803 let db = self.executor.db.clone();
804 let mut full_key = self.executor.start_key.to_vec();
805 full_key.extend(key);
806 self.spawn_mode
807 .spawn(
808 move |x| {
809 if !db.key_may_exist(&x) {
810 return Ok(false);
811 }
812 Ok(db.get(&x)?.is_some())
813 },
814 full_key,
815 )
816 .await
817 }
818
819 async fn contains_keys(
820 &self,
821 keys: &[Vec<u8>],
822 ) -> Result<Vec<bool>, RocksDbStoreInternalError> {
823 let executor = self.executor.clone();
824 self.spawn_mode
825 .spawn(move |x| executor.contains_keys_internal(x), keys.to_vec())
826 .await
827 }
828
829 async fn read_multi_values_bytes(
830 &self,
831 keys: &[Vec<u8>],
832 ) -> Result<Vec<Option<Vec<u8>>>, RocksDbStoreInternalError> {
833 let executor = self.executor.clone();
834 self.spawn_mode
835 .spawn(
836 move |x| executor.read_multi_values_bytes_internal(x),
837 keys.to_vec(),
838 )
839 .await
840 }
841
842 async fn find_keys_by_prefix(
843 &self,
844 key_prefix: &[u8],
845 ) -> Result<Vec<Vec<u8>>, RocksDbStoreInternalError> {
846 let executor = self.executor.clone();
847 let key_prefix = key_prefix.to_vec();
848 self.spawn_mode
849 .spawn(
850 move |x| executor.find_keys_by_prefix_internal(x),
851 key_prefix,
852 )
853 .await
854 }
855
856 async fn find_key_values_by_prefix(
857 &self,
858 key_prefix: &[u8],
859 ) -> Result<Vec<(Vec<u8>, Vec<u8>)>, RocksDbStoreInternalError> {
860 let executor = self.executor.clone();
861 let key_prefix = key_prefix.to_vec();
862 self.spawn_mode
863 .spawn(
864 move |x| executor.find_key_values_by_prefix_internal(x),
865 key_prefix,
866 )
867 .await
868 }
869}
870
871impl WritableKeyValueStore for RocksDbStoreInternal {
872 const MAX_VALUE_SIZE: usize = MAX_VALUE_SIZE;
873
874 async fn write_batch(&self, batch: Batch) -> Result<(), RocksDbStoreInternalError> {
875 let write_root_key = !self.root_key_written.fetch_or(true, Ordering::SeqCst);
876 let executor = self.executor.clone();
877 self.spawn_mode
878 .spawn(
879 move |x| executor.write_batch_internal(x, write_root_key),
880 batch,
881 )
882 .await
883 }
884
885 async fn clear_journal(&self) -> Result<(), RocksDbStoreInternalError> {
886 Ok(())
887 }
888}
889
890impl KeyValueDatabase for RocksDbDatabaseInternal {
891 type Config = RocksDbStoreInternalConfig;
892 type Store = RocksDbStoreInternal;
893
894 fn get_name() -> String {
895 "rocksdb internal".to_string()
896 }
897
898 async fn connect(
899 config: &Self::Config,
900 namespace: &str,
901 ) -> Result<Self, RocksDbStoreInternalError> {
902 Self::build(config, namespace)
903 }
904
905 fn open_shared(&self, root_key: &[u8]) -> Result<Self::Store, RocksDbStoreInternalError> {
906 let mut start_key = ROOT_KEY_DOMAIN.to_vec();
907 start_key.extend(bcs::to_bytes(root_key)?);
908 let mut executor = self.executor.clone();
909 executor.start_key = start_key;
910 Ok(RocksDbStoreInternal {
911 executor,
912 path_with_guard: self.path_with_guard.clone(),
913 spawn_mode: self.spawn_mode,
914 root_key_written: Arc::new(AtomicBool::new(false)),
915 })
916 }
917
918 fn open_exclusive(&self, root_key: &[u8]) -> Result<Self::Store, RocksDbStoreInternalError> {
919 self.open_shared(root_key)
920 }
921
922 async fn list_all(config: &Self::Config) -> Result<Vec<String>, RocksDbStoreInternalError> {
923 let entries = std::fs::read_dir(config.path_with_guard.path_buf.clone())?;
924 let mut namespaces = Vec::new();
925 for entry in entries {
926 let entry = entry?;
927 if !entry.file_type()?.is_dir() {
928 return Err(RocksDbStoreInternalError::NonDirectoryNamespace);
929 }
930 let namespace = match entry.file_name().into_string() {
931 Err(error) => {
932 return Err(RocksDbStoreInternalError::IntoStringError(error));
933 }
934 Ok(namespace) => namespace,
935 };
936 namespaces.push(namespace);
937 }
938 Ok(namespaces)
939 }
940
941 async fn list_root_keys(&self) -> Result<Vec<Vec<u8>>, RocksDbStoreInternalError> {
942 let mut store = self.open_shared(&[])?;
943 store.executor.start_key = vec![STORED_ROOT_KEYS_PREFIX];
944 let bcs_root_keys = store.find_keys_by_prefix(&[]).await?;
945 let mut root_keys = Vec::new();
946 for bcs_root_key in bcs_root_keys {
947 let root_key = bcs::from_bytes::<Vec<u8>>(&bcs_root_key)?;
948 root_keys.push(root_key);
949 }
950 Ok(root_keys)
951 }
952
953 async fn delete_all(config: &Self::Config) -> Result<(), RocksDbStoreInternalError> {
954 let namespaces = Self::list_all(config).await?;
955 for namespace in namespaces {
956 let mut path_buf = config.path_with_guard.path_buf.clone();
957 path_buf.push(&namespace);
958 std::fs::remove_dir_all(path_buf.as_path())?;
959 }
960 Ok(())
961 }
962
963 async fn exists(
964 config: &Self::Config,
965 namespace: &str,
966 ) -> Result<bool, RocksDbStoreInternalError> {
967 Self::check_namespace(namespace)?;
968 let mut path_buf = config.path_with_guard.path_buf.clone();
969 path_buf.push(namespace);
970 let test = std::path::Path::exists(&path_buf);
971 Ok(test)
972 }
973
974 async fn create(
975 config: &Self::Config,
976 namespace: &str,
977 ) -> Result<(), RocksDbStoreInternalError> {
978 Self::check_namespace(namespace)?;
979 let mut path_buf = config.path_with_guard.path_buf.clone();
980 path_buf.push(namespace);
981 if std::path::Path::exists(&path_buf) {
982 return Err(RocksDbStoreInternalError::StoreAlreadyExists);
983 }
984 std::fs::create_dir_all(path_buf)?;
985 Ok(())
986 }
987
988 async fn delete(
989 config: &Self::Config,
990 namespace: &str,
991 ) -> Result<(), RocksDbStoreInternalError> {
992 Self::check_namespace(namespace)?;
993 let mut path_buf = config.path_with_guard.path_buf.clone();
994 path_buf.push(namespace);
995 let path = path_buf.as_path();
996 std::fs::remove_dir_all(path)?;
997 Ok(())
998 }
999}
1000
1001#[cfg(with_testing)]
1002impl TestKeyValueDatabase for RocksDbDatabaseInternal {
1003 async fn new_test_config() -> Result<RocksDbStoreInternalConfig, RocksDbStoreInternalError> {
1004 let path_with_guard = PathWithGuard::new_testing();
1005 let spawn_mode = RocksDbSpawnMode::get_spawn_mode_from_runtime();
1006 Ok(RocksDbStoreInternalConfig {
1007 path_with_guard,
1008 spawn_mode,
1009 enable_statistics: false,
1010 statistics_level: RocksDbStatisticsLevel::default(),
1011 })
1012 }
1013}
1014
1015#[derive(Error, Debug)]
1017pub enum RocksDbStoreInternalError {
1018 #[error("Store already exists")]
1020 StoreAlreadyExists,
1021
1022 #[error("tokio join error: {0}")]
1024 TokioJoinError(#[from] tokio::task::JoinError),
1025
1026 #[error("RocksDB error: {0}")]
1028 RocksDb(#[from] rocksdb::Error),
1029
1030 #[error("Namespaces should be directories")]
1032 NonDirectoryNamespace,
1033
1034 #[error("error in the conversion from OsString: {0:?}")]
1036 IntoStringError(OsString),
1037
1038 #[error("The key must have at most 8 MiB")]
1040 KeyTooLong,
1041
1042 #[error("Namespace contains forbidden characters")]
1044 InvalidNamespace,
1045
1046 #[error("Filesystem error: {0}")]
1048 FsError(#[from] std::io::Error),
1049
1050 #[error(transparent)]
1052 BcsError(#[from] bcs::Error),
1053}
1054
1055#[derive(Clone, Debug, Deserialize, Serialize)]
1057pub struct PathWithGuard {
1058 pub path_buf: PathBuf,
1060 #[serde(skip)]
1062 _dir_guard: Option<Arc<TempDir>>,
1063}
1064
1065impl PathWithGuard {
1066 pub fn new(path_buf: PathBuf) -> Self {
1068 Self {
1069 path_buf,
1070 _dir_guard: None,
1071 }
1072 }
1073
1074 #[cfg(with_testing)]
1076 fn new_testing() -> PathWithGuard {
1077 let dir = TempDir::new().unwrap();
1078 let path_buf = dir.path().to_path_buf();
1079 let dir_guard = Some(Arc::new(dir));
1080 PathWithGuard {
1081 path_buf,
1082 _dir_guard: dir_guard,
1083 }
1084 }
1085}
1086
1087impl PartialEq for PathWithGuard {
1088 fn eq(&self, other: &Self) -> bool {
1089 self.path_buf == other.path_buf
1090 }
1091}
1092impl Eq for PathWithGuard {}
1093
1094impl KeyValueStoreError for RocksDbStoreInternalError {
1095 const BACKEND: &'static str = "rocks_db";
1096}
1097
1098pub type RocksDbStoreError = RocksDbStoreInternalError;
1100
1101pub type RocksDbStoreConfig = LruCachingConfig<RocksDbStoreInternalConfig>;
1103
1104#[cfg(with_metrics)]
1106pub type RocksDbDatabase =
1107 MeteredDatabase<LruCachingDatabase<MeteredDatabase<RocksDbDatabaseInternal>>>;
1108#[cfg(not(with_metrics))]
1110pub type RocksDbDatabase = LruCachingDatabase<RocksDbDatabaseInternal>;
1111
1112#[cfg(with_testing)]
1113impl crate::backends::DatabaseBackup for RocksDbDatabaseInternal {
1114 fn backup_to(&self, dir: &std::path::Path) -> anyhow::Result<()> {
1115 use rocksdb::{
1116 backup::{BackupEngine, BackupEngineOptions},
1117 Env,
1118 };
1119 let opts = BackupEngineOptions::new(dir)?;
1120 let env = Env::new()?;
1121 let mut engine = BackupEngine::open(&opts, &env)?;
1122 engine.create_new_backup_flush(&*self.executor.db, true)?;
1123 Ok(())
1124 }
1125}