Skip to main content

linera_views/backends/
rocks_db.rs

1// Copyright (c) Zefchain Labs, Inc.
2// SPDX-License-Identifier: Apache-2.0
3
4//! Implements [`crate::store::KeyValueStore`] for the RocksDB database.
5
6// RocksDB's C API uses `i32` and signed sizes; casts at this boundary are
7// by design.
8#![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
45/// The prefixes being used in the system
46static ROOT_KEY_DOMAIN: [u8; 1] = [0];
47static STORED_ROOT_KEYS_PREFIX: u8 = 1;
48
49// The maximum size of values in RocksDB is 3 GiB
50// For offset reasons we decrease by 400
51const MAX_VALUE_SIZE: usize = 3 * 1024 * 1024 * 1024 - 400;
52
53// The maximum size of keys in RocksDB is 8 MiB
54// For offset reasons we decrease by 400
55const MAX_KEY_SIZE: usize = 8 * 1024 * 1024 - 400;
56
57// A small write buffer keeps the memtable flushing even on low-write workloads. A large buffer
58// (e.g. 256 MiB) lets a slowly-filling memtable grow huge and accumulate range tombstones, so every
59// point read has to scan it — which severely amplifies reads during e.g. a long chain
60// synchronization, where blocks are re-executed with little net data written.
61const WRITE_BUFFER_SIZE: usize = 16 * 1024 * 1024; // 16 MiB
62const 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; // 8 KiB
74
75/// The RocksDB client that we use.
76type DB = rocksdb::DBWithThreadMode<rocksdb::MultiThreaded>;
77
78/// The choice of the spawning mode.
79/// `SpawnBlocking` always works and is the safest.
80/// `BlockInPlace` can only be used in multi-threaded environment.
81/// One way to select that is to select BlockInPlace when
82/// `tokio::runtime::Handle::current().metrics().num_workers() > 1`
83/// `BlockInPlace` is documented in <https://docs.rs/tokio/latest/tokio/task/fn.block_in_place.html>
84#[derive(Clone, Copy, Debug, PartialEq, Eq, Deserialize, Serialize)]
85pub enum RocksDbSpawnMode {
86    /// This uses the `spawn_blocking` function of Tokio.
87    SpawnBlocking,
88    /// This uses the `block_in_place` function of Tokio.
89    BlockInPlace,
90}
91
92impl RocksDbSpawnMode {
93    /// Obtains the spawning mode from runtime.
94    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    /// Runs the computation for a function according to the selected policy.
103    #[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        // Configure ReadOptions optimized for SSDs and iterator performance
191        let mut read_opts = rocksdb::ReadOptions::default();
192        // Enable async I/O for better concurrency
193        read_opts.set_async_io(true);
194
195        // Set precise upper bound to minimize key traversal
196        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/// The inner client
286#[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/// Database-level connection to RocksDB for managing namespaces and partitions.
295#[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/// The level of detail collected by RocksDB's internal statistics.
307///
308/// This mirrors [`rocksdb::statistics::StatsLevel`]. The levels are nested: each one
309/// collects a superset of the data collected by the previous one, at increasing cost.
310#[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    /// Collect nothing.
315    DisableAll,
316    /// Collect tickers (counters) only; skip all histograms and timers.
317    #[default]
318    ExceptHistogramOrTimers,
319    /// Collect tickers and histograms, but skip timer statistics.
320    ExceptTimers,
321    /// Collect everything except time spent inside the mutex lock and on compression.
322    ExceptDetailedTimers,
323    /// Collect everything except the counters that require taking time inside the mutex lock.
324    ExceptTimeForMutex,
325    /// Collect everything, including the duration of mutex operations.
326    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/// The initial configuration of the system
376#[derive(Clone, Debug, Deserialize, Serialize)]
377pub struct RocksDbStoreInternalConfig {
378    /// The path to the storage containing the namespaces
379    pub path_with_guard: PathWithGuard,
380    /// The chosen spawn mode
381    pub spawn_mode: RocksDbSpawnMode,
382    /// Whether to enable RocksDB's internal statistics collection and export it as
383    /// Prometheus metrics. Disabled by default to avoid overhead in clients that do not
384    /// scrape metrics; enabled explicitly for the workers.
385    #[serde(default)]
386    pub enable_statistics: bool,
387    /// The level of detail collected when `enable_statistics` is set.
388    #[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        // Create a store to extract its executor and configuration
409        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        // Flush in-memory buffer to disk more often
444        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        // We deliberately give RocksDB one background thread *per* CPU so that
451        // flush + (N-1) compactions can hammer the NVMe at full bandwidth while
452        // still leaving enough CPU time for the foreground application threads.
453        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        // Allocate 1/4 of total RAM for RocksDB block cache, which is a reasonable balance:
465        // - Large enough to significantly improve read performance by caching frequently accessed blocks
466        // - Small enough to leave memory for other system components
467        // - Follows common practice for database caching in server environments
468        // - Prevents excessive memory pressure that could lead to swapping or OOM conditions
469        block_options.set_block_cache(&Cache::new_hyper_clock_cache(
470            total_ram / 4,
471            HYPER_CLOCK_CACHE_BLOCK_SIZE,
472        ));
473
474        // Cap total memtable memory to prevent unbounded growth when multiple column
475        // families are used or many memtables accumulate before flushing.
476        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        // Configure bloom filters for prefix iteration optimization
481        block_options.set_bloom_filter(10.0, false);
482        block_options.set_whole_key_filtering(false);
483
484        // 32KB blocks instead of default 4KB - reduces iterator seeks
485        block_options.set_block_size(32 * 1024);
486        // Use latest format for better compression and performance
487        block_options.set_format_version(5);
488
489        options.set_block_based_table_factory(&block_options);
490
491        // Configure prefix extraction for bloom filter optimization
492        // Use 8 bytes: ROOT_KEY_DOMAIN (1 byte) + BCS variant (1-2 bytes) + identifier start (4-5 bytes)
493        let prefix_extractor = SliceTransform::create_fixed_prefix(8);
494        options.set_prefix_extractor(prefix_extractor);
495
496        // 12.5% of memtable size for bloom filter
497        options.set_memtable_prefix_bloom_ratio(0.125);
498        // Skip bloom filter for memtable when key exists
499        options.set_optimize_filters_for_hits(true);
500        // Use memory-mapped files for faster reads
501        options.set_allow_mmap_reads(true);
502        // Don't use random access pattern since we do prefix scans
503        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/// Exports RocksDB's internal statistics as Prometheus metrics.
526///
527/// The collector reads the values lazily at scrape time: cumulative tickers via
528/// `get_ticker_count` and instantaneous LSM state via `GetIntProperty`. Neither requires
529/// the (more expensive) histogram/timer statistics levels.
530#[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/// The error type for [`RocksDbStoreInternal`]
1016#[derive(Error, Debug)]
1017pub enum RocksDbStoreInternalError {
1018    /// Store already exists
1019    #[error("Store already exists")]
1020    StoreAlreadyExists,
1021
1022    /// Tokio join error in RocksDB.
1023    #[error("tokio join error: {0}")]
1024    TokioJoinError(#[from] tokio::task::JoinError),
1025
1026    /// RocksDB error.
1027    #[error("RocksDB error: {0}")]
1028    RocksDb(#[from] rocksdb::Error),
1029
1030    /// The database contains a file which is not a directory
1031    #[error("Namespaces should be directories")]
1032    NonDirectoryNamespace,
1033
1034    /// Error converting `OsString` to `String`
1035    #[error("error in the conversion from OsString: {0:?}")]
1036    IntoStringError(OsString),
1037
1038    /// The key must have at most 8 MiB
1039    #[error("The key must have at most 8 MiB")]
1040    KeyTooLong,
1041
1042    /// Namespace contains forbidden characters
1043    #[error("Namespace contains forbidden characters")]
1044    InvalidNamespace,
1045
1046    /// Filesystem error
1047    #[error("Filesystem error: {0}")]
1048    FsError(#[from] std::io::Error),
1049
1050    /// BCS serialization error.
1051    #[error(transparent)]
1052    BcsError(#[from] bcs::Error),
1053}
1054
1055/// A path and the guard for the temporary directory if needed
1056#[derive(Clone, Debug, Deserialize, Serialize)]
1057pub struct PathWithGuard {
1058    /// The path to the data
1059    pub path_buf: PathBuf,
1060    /// The guard for the directory if one is needed
1061    #[serde(skip)]
1062    _dir_guard: Option<Arc<TempDir>>,
1063}
1064
1065impl PathWithGuard {
1066    /// Creates a `PathWithGuard` from an existing path.
1067    pub fn new(path_buf: PathBuf) -> Self {
1068        Self {
1069            path_buf,
1070            _dir_guard: None,
1071        }
1072    }
1073
1074    /// Returns the test path for RocksDB without common config.
1075    #[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
1098/// The composed error type for the `RocksDbStore`
1099pub type RocksDbStoreError = RocksDbStoreInternalError;
1100
1101/// The composed config type for the `RocksDbStore`
1102pub type RocksDbStoreConfig = LruCachingConfig<RocksDbStoreInternalConfig>;
1103
1104/// The `RocksDbDatabase` composed type with metrics
1105#[cfg(with_metrics)]
1106pub type RocksDbDatabase =
1107    MeteredDatabase<LruCachingDatabase<MeteredDatabase<RocksDbDatabaseInternal>>>;
1108/// The `RocksDbDatabase` composed type
1109#[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}