use std::str::FromStr;
use std::time::Duration;
use crate::cnf::Config;
use crate::kvs::config::{SyncMode, parse_duration};
use crate::sys::TOTAL_SYSTEM_MEMORY;
const KIB: u64 = 1024;
const MIB: u64 = 1024 * KIB;
const GIB: u64 = 1024 * MIB;
fn default_compaction_readahead_size() -> usize {
(256 * KIB) as usize
}
fn default_block_cache_size() -> usize {
let mem = *TOTAL_SYSTEM_MEMORY;
mem.saturating_div(2).saturating_sub(GIB).max(16 * MIB) as usize
}
fn default_write_buffer_size() -> usize {
let mem = *TOTAL_SYSTEM_MEMORY;
if mem < GIB {
(32 * MIB) as usize
} else if mem < 16 * GIB {
(64 * MIB) as usize
} else {
(128 * MIB) as usize
}
}
fn default_max_write_buffer_number() -> usize {
let mem = *TOTAL_SYSTEM_MEMORY;
if mem < 4 * GIB {
2
} else if mem < 16 * GIB {
4
} else if mem < 64 * GIB {
8
} else {
32
}
}
fn default_max_concurrent_subcompactions() -> u32 {
let cpu_count = std::thread::available_parallelism().map(|x| x.get()).unwrap_or(1);
(cpu_count as u32).clamp(1, 4)
}
fn default_jobs_count() -> usize {
let cpu_count = std::thread::available_parallelism().map(|x| x.get()).unwrap_or(1);
let memory_limited_jobs = (*TOTAL_SYSTEM_MEMORY / (128 * MIB)).max(2) as usize;
(cpu_count * 2).min(memory_limited_jobs)
}
fn default_target_file_size_base() -> u64 {
let mem = *TOTAL_SYSTEM_MEMORY;
if mem < GIB {
8 * MIB
} else if mem < 4 * GIB {
16 * MIB
} else if mem < 16 * GIB {
32 * MIB
} else {
64 * MIB
}
}
fn default_grouped_commit_max_batch_size() -> usize {
let mem = *TOTAL_SYSTEM_MEMORY;
if mem < GIB {
256
} else if mem < 4 * GIB {
1024
} else {
4096
}
}
fn default_keep_log_file_num() -> usize {
let mem = *TOTAL_SYSTEM_MEMORY;
if mem < GIB {
2
} else if mem < 4 * GIB {
5
} else {
10
}
}
fn default_max_auto_readahead_size() -> usize {
(256 * KIB) as usize
}
fn default_max_open_files() -> usize {
let mem = *TOTAL_SYSTEM_MEMORY;
if mem < GIB {
256
} else if mem < 4 * GIB {
512
} else {
1026
}
}
fn default_enable_blob_files() -> bool {
true
}
fn default_blob_file_size() -> u64 {
let mem = *TOTAL_SYSTEM_MEMORY;
if mem < GIB {
16 * MIB
} else if mem < 4 * GIB {
64 * MIB
} else if mem < 16 * GIB {
128 * MIB
} else {
256 * MIB
}
}
fn default_wal_size_limit() -> u64 {
let mem = *TOTAL_SYSTEM_MEMORY;
if mem < GIB {
32
} else if mem < 16 * GIB {
128
} else {
0
}
}
#[derive(Default, Clone, Debug)]
pub enum BlobCompression {
#[default]
Snappy,
Lz4,
Zstd,
None,
}
impl FromStr for BlobCompression {
type Err = ();
fn from_str(s: &str) -> Result<Self, Self::Err> {
if s.eq_ignore_ascii_case("none") {
Ok(Self::None)
} else if s.eq_ignore_ascii_case("lz4") {
Ok(Self::Lz4)
} else if s.eq_ignore_ascii_case("snappy") {
Ok(Self::Snappy)
} else if s.eq_ignore_ascii_case("zstd") {
Ok(Self::Zstd)
} else {
Err(())
}
}
}
fn default_file_compaction_trigger() -> usize {
2
}
fn default_level0_slowdown_writes_trigger() -> i32 {
8
}
fn default_level0_stop_writes_trigger() -> i32 {
12
}
fn default_periodic_compaction_seconds() -> u64 {
3600
}
fn default_shutdown_wait_for_compact_seconds() -> u64 {
30
}
fn default_runtime_worker_threads() -> usize {
std::cmp::max(4, num_cpus::get())
}
#[derive(Debug, Clone)]
pub struct RocksDbConfig {
pub versioned: bool,
pub retention: Duration,
pub sync_mode: SyncMode,
pub thread_count: usize,
pub jobs_count: usize,
pub max_open_files: usize,
pub block_size: usize,
pub wal_size_limit: u64,
pub target_file_size_base: u64,
pub target_file_size_multiplier: usize,
pub file_compaction_trigger: usize,
pub level0_slowdown_writes_trigger: i32,
pub level0_stop_writes_trigger: i32,
pub periodic_compaction_seconds: u64,
pub compaction_readahead_size: usize,
pub max_concurrent_subcompactions: u32,
pub enable_pipelined_writes: bool,
pub keep_log_file_num: usize,
pub storage_log_level: String,
pub compaction_style: String,
pub universal_size_ratio: i32,
pub universal_min_merge_width: u32,
pub universal_max_merge_width: u32,
pub universal_max_size_amplification_percent: u32,
pub universal_compression_size_percent: i32,
pub universal_stop_style: String,
pub deletion_factory_window_size: usize,
pub deletion_factory_delete_count: usize,
pub deletion_factory_ratio: f64,
pub enable_blob_files: bool,
pub min_blob_size: u64,
pub blob_file_size: u64,
pub blob_compression_type: BlobCompression,
pub enable_blob_gc: bool,
pub blob_gc_age_cutoff: f64,
pub blob_gc_force_threshold: f64,
pub blob_compaction_readahead_size: u64,
pub block_cache_size: usize,
pub write_buffer_size: usize,
pub max_write_buffer_number: usize,
pub min_write_buffer_number_to_merge: usize,
pub sst_max_allowed_space_usage: u64,
pub grouped_commit_timeout: u64,
pub grouped_commit_wait_threshold: usize,
pub grouped_commit_max_batch_size: usize,
pub initial_auto_readahead_size: usize,
pub max_auto_readahead_size: usize,
pub file_reads_for_auto_readahead: u64,
pub prefix_extractor_enabled: bool,
pub whole_key_filtering: bool,
pub memtable_prefix_bloom_ratio: f64,
pub scan_verify_checksums: bool,
pub compact_on_shutdown: bool,
pub shutdown_wait_for_compact_seconds: u64,
pub runtime_worker_threads: usize,
pub runtime_reserve: usize,
}
impl Default for RocksDbConfig {
fn default() -> Self {
let cpu_count = std::thread::available_parallelism().map(|x| x.get()).unwrap_or(1);
Self {
versioned: false,
retention: Duration::ZERO,
sync_mode: SyncMode::Every,
thread_count: cpu_count,
jobs_count: default_jobs_count(),
max_open_files: default_max_open_files(),
block_size: 64 * 1024,
wal_size_limit: default_wal_size_limit(),
target_file_size_base: default_target_file_size_base(),
target_file_size_multiplier: 2,
file_compaction_trigger: default_file_compaction_trigger(),
level0_slowdown_writes_trigger: default_level0_slowdown_writes_trigger(),
level0_stop_writes_trigger: default_level0_stop_writes_trigger(),
periodic_compaction_seconds: default_periodic_compaction_seconds(),
compaction_readahead_size: default_compaction_readahead_size(),
max_concurrent_subcompactions: default_max_concurrent_subcompactions(),
enable_pipelined_writes: true,
keep_log_file_num: default_keep_log_file_num(),
storage_log_level: "warn".to_owned(),
compaction_style: "level".to_owned(),
universal_size_ratio: 1,
universal_min_merge_width: 2,
universal_max_merge_width: u32::MAX,
universal_max_size_amplification_percent: 200,
universal_compression_size_percent: -1,
universal_stop_style: "total".to_owned(),
deletion_factory_window_size: 500,
deletion_factory_delete_count: 25,
deletion_factory_ratio: 0.5,
enable_blob_files: default_enable_blob_files(),
min_blob_size: 4 * 1024,
blob_file_size: default_blob_file_size(),
blob_compression_type: Default::default(),
enable_blob_gc: true,
blob_gc_age_cutoff: 0.5,
blob_gc_force_threshold: 0.5,
blob_compaction_readahead_size: 0,
block_cache_size: default_block_cache_size(),
write_buffer_size: default_write_buffer_size(),
max_write_buffer_number: default_max_write_buffer_number(),
min_write_buffer_number_to_merge: 2,
sst_max_allowed_space_usage: 0,
grouped_commit_timeout: Duration::from_millis(5).as_nanos() as u64,
grouped_commit_wait_threshold: 12,
grouped_commit_max_batch_size: default_grouped_commit_max_batch_size(),
initial_auto_readahead_size: 8 * 1024,
max_auto_readahead_size: default_max_auto_readahead_size(),
file_reads_for_auto_readahead: 2,
prefix_extractor_enabled: true,
whole_key_filtering: true,
memtable_prefix_bloom_ratio: 0.1,
scan_verify_checksums: true,
compact_on_shutdown: false,
shutdown_wait_for_compact_seconds: default_shutdown_wait_for_compact_seconds(),
runtime_worker_threads: default_runtime_worker_threads(),
runtime_reserve: 2,
}
}
}
impl Config for RocksDbConfig {
fn parse(&mut self, map: &crate::cnf::ConfigMap) {
map.parse_key_bool("datastore_versioned", &mut self.versioned)
.parse_key_with("datastore_retention", &mut self.retention, |x| parse_duration(x).ok())
.parse_key("rocksdb_thread_count", &mut self.thread_count)
.parse_key("rocksdb_jobs_count", &mut self.jobs_count)
.parse_key("rocksdb_max_open_files", &mut self.max_open_files)
.parse_key("rocksdb_block_size", &mut self.block_size)
.parse_key("rocksdb_wal_size_limit", &mut self.wal_size_limit)
.parse_key("rocksdb_target_file_size_base", &mut self.target_file_size_base)
.parse_key("rocksdb_target_file_size_multiplier", &mut self.target_file_size_multiplier)
.parse_key("rocksdb_file_compaction_trigger", &mut self.file_compaction_trigger)
.parse_key(
"rocksdb_level0_slowdown_writes_trigger",
&mut self.level0_slowdown_writes_trigger,
)
.parse_key("rocksdb_level0_stop_writes_trigger", &mut self.level0_stop_writes_trigger)
.parse_key("rocksdb_periodic_compaction_seconds", &mut self.periodic_compaction_seconds)
.parse_key("rocksdb_compaction_readahead_size", &mut self.compaction_readahead_size)
.parse_key(
"rocksdb_max_concurrent_subcompactions",
&mut self.max_concurrent_subcompactions,
)
.parse_key_bool("rocksdb_enable_pipelined_writes", &mut self.enable_pipelined_writes)
.parse_key("rocksdb_keep_log_file_num", &mut self.keep_log_file_num)
.parse_key("rocksdb_storage_log_level", &mut self.storage_log_level)
.parse_key("rocksdb_compaction_style", &mut self.compaction_style)
.parse_key("rocksdb_universal_size_ratio", &mut self.universal_size_ratio)
.parse_key("rocksdb_universal_min_merge_width", &mut self.universal_min_merge_width)
.parse_key("rocksdb_universal_max_merge_width", &mut self.universal_max_merge_width)
.parse_key(
"rocksdb_universal_max_size_amplification_percent",
&mut self.universal_max_size_amplification_percent,
)
.parse_key(
"rocksdb_universal_compression_size_percent",
&mut self.universal_compression_size_percent,
)
.parse_key("rocksdb_universal_stop_style", &mut self.universal_stop_style)
.parse_key(
"rocksdb_deletion_factory_window_size",
&mut self.deletion_factory_window_size,
)
.parse_key(
"rocksdb_deletion_factory_delete_count",
&mut self.deletion_factory_delete_count,
)
.parse_key("rocksdb_deletion_factory_ratio", &mut self.deletion_factory_ratio)
.parse_key_bool("rocksdb_enable_blob_files", &mut self.enable_blob_files)
.parse_key("rocksdb_min_blob_size", &mut self.min_blob_size)
.parse_key("rocksdb_blob_file_size", &mut self.blob_file_size)
.parse_key("rocksdb_blob_compression_type", &mut self.blob_compression_type)
.parse_key_bool("rocksdb_enable_blob_gc", &mut self.enable_blob_gc)
.parse_key("rocksdb_blob_gc_age_cutoff", &mut self.blob_gc_age_cutoff)
.parse_key("rocksdb_blob_gc_force_threshold", &mut self.blob_gc_force_threshold)
.parse_key(
"rocksdb_blob_compaction_readahead_size",
&mut self.blob_compaction_readahead_size,
)
.parse_key("rocksdb_block_cache_size", &mut self.block_cache_size)
.parse_key("rocksdb_write_buffer_size", &mut self.write_buffer_size)
.parse_key("rocksdb_max_write_buffer_number", &mut self.max_write_buffer_number)
.parse_key(
"rocksdb_min_write_buffer_number_to_merge",
&mut self.min_write_buffer_number_to_merge,
)
.parse_key("rocksdb_sst_max_allowed_space_usage", &mut self.sst_max_allowed_space_usage)
.parse_key("rocksdb_grouped_commit_timeout", &mut self.grouped_commit_timeout)
.parse_key(
"rocksdb_grouped_commit_wait_threshold",
&mut self.grouped_commit_wait_threshold,
)
.parse_key(
"rocksdb_grouped_commit_max_batch_size",
&mut self.grouped_commit_max_batch_size,
)
.parse_key("rocksdb_initial_auto_readahead_size", &mut self.initial_auto_readahead_size)
.parse_key("rocksdb_max_auto_readahead_size", &mut self.max_auto_readahead_size)
.parse_key(
"rocksdb_file_reads_for_auto_readahead",
&mut self.file_reads_for_auto_readahead,
)
.parse_key_bool("rocksdb_prefix_extractor_enabled", &mut self.prefix_extractor_enabled)
.parse_key_bool("rocksdb_whole_key_filtering", &mut self.whole_key_filtering)
.parse_key("rocksdb_memtable_prefix_bloom_ratio", &mut self.memtable_prefix_bloom_ratio)
.parse_key_bool("rocksdb_scan_verify_checksums", &mut self.scan_verify_checksums)
.parse_key_bool("rocksdb_compact_on_shutdown", &mut self.compact_on_shutdown)
.parse_key(
"rocksdb_shutdown_wait_for_compact_seconds",
&mut self.shutdown_wait_for_compact_seconds,
)
.parse_key("rocksdb_runtime_reserve", &mut self.runtime_reserve)
.parse_key("runtime_worker_threads", &mut self.runtime_worker_threads);
if map.has_key("datastore_sync") {
map.parse_key("datastore_sync", &mut self.sync_mode);
} else {
map.parse_key("datastore_sync_data", &mut self.sync_mode);
}
}
}
#[cfg(test)]
mod test {
use std::time::Duration;
use crate::cnf::ConfigMap;
use crate::kvs::config::SyncMode;
use crate::kvs::rocksdb::RocksDbConfig;
#[test]
fn test_rocksdb_config_defaults() {
let map = ConfigMap::empty();
let config = map.load::<RocksDbConfig>();
assert!(!config.versioned);
assert_eq!(config.retention, Duration::ZERO);
assert_eq!(config.sync_mode, SyncMode::Every);
}
#[test]
fn test_rocksdb_config_compaction_defaults() {
let config = ConfigMap::empty().load::<RocksDbConfig>();
assert_eq!(config.file_compaction_trigger, 2);
assert_eq!(config.level0_slowdown_writes_trigger, 8);
assert_eq!(config.level0_stop_writes_trigger, 12);
assert_eq!(config.periodic_compaction_seconds, 3600);
assert_eq!(config.deletion_factory_window_size, 500);
assert_eq!(config.deletion_factory_delete_count, 25);
assert!((config.deletion_factory_ratio - 0.5).abs() < f64::EPSILON);
assert_eq!(config.compaction_style, "level");
assert_eq!(config.universal_size_ratio, 1);
assert_eq!(config.universal_min_merge_width, 2);
assert_eq!(config.universal_max_merge_width, u32::MAX);
assert_eq!(config.universal_max_size_amplification_percent, 200);
assert_eq!(config.universal_compression_size_percent, -1);
assert_eq!(config.universal_stop_style, "total");
assert!(!config.compact_on_shutdown);
assert_eq!(config.shutdown_wait_for_compact_seconds, 30);
}
#[test]
fn test_rocksdb_config_compaction_overrides() {
let map = ConfigMap::empty()
.with_key_value("rocksdb_file_compaction_trigger", "5")
.with_key_value("rocksdb_level0_slowdown_writes_trigger", "16")
.with_key_value("rocksdb_level0_stop_writes_trigger", "24")
.with_key_value("rocksdb_periodic_compaction_seconds", "0");
let config = map.load::<RocksDbConfig>();
assert_eq!(config.file_compaction_trigger, 5);
assert_eq!(config.level0_slowdown_writes_trigger, 16);
assert_eq!(config.level0_stop_writes_trigger, 24);
assert_eq!(config.periodic_compaction_seconds, 0);
}
#[test]
fn test_rocksdb_config_shutdown_overrides() {
let map = ConfigMap::empty()
.with_key_value("rocksdb_compact_on_shutdown", "true")
.with_key_value("rocksdb_shutdown_wait_for_compact_seconds", "5");
let config = map.load::<RocksDbConfig>();
assert!(config.compact_on_shutdown);
assert_eq!(config.shutdown_wait_for_compact_seconds, 5);
}
#[test]
fn test_rocksdb_config_universal_overrides() {
let map = ConfigMap::empty()
.with_key_value("rocksdb_compaction_style", "universal")
.with_key_value("rocksdb_universal_size_ratio", "5")
.with_key_value("rocksdb_universal_min_merge_width", "3")
.with_key_value("rocksdb_universal_max_merge_width", "16")
.with_key_value("rocksdb_universal_max_size_amplification_percent", "150")
.with_key_value("rocksdb_universal_compression_size_percent", "75")
.with_key_value("rocksdb_universal_stop_style", "similar_size");
let config = map.load::<RocksDbConfig>();
assert_eq!(config.compaction_style, "universal");
assert_eq!(config.universal_size_ratio, 5);
assert_eq!(config.universal_min_merge_width, 3);
assert_eq!(config.universal_max_merge_width, 16);
assert_eq!(config.universal_max_size_amplification_percent, 150);
assert_eq!(config.universal_compression_size_percent, 75);
assert_eq!(config.universal_stop_style, "similar_size");
}
#[test]
fn test_rocksdb_config_sync_every() {
let map =
ConfigMap::from_config_string("sync=every").map_keys(|x| format!("datastore_{x}"));
let config = map.load::<RocksDbConfig>();
assert_eq!(config.sync_mode, SyncMode::Every);
}
#[test]
fn test_rocksdb_config_sync_never() {
let map =
ConfigMap::from_config_string("sync=never").map_keys(|x| format!("datastore_{x}"));
let config = map.load::<RocksDbConfig>();
assert_eq!(config.sync_mode, SyncMode::Never);
}
#[test]
fn test_rocksdb_config_sync_periodic() {
let map =
ConfigMap::from_config_string("sync=200ms").map_keys(|x| format!("datastore_{x}"));
let config = map.load::<RocksDbConfig>();
assert_eq!(config.sync_mode, SyncMode::Interval(Duration::from_millis(200)));
}
#[test]
fn test_rocksdb_config_sync_periodic_seconds() {
let map = ConfigMap::from_config_string("sync=5s").map_keys(|x| format!("datastore_{x}"));
let config = map.load::<RocksDbConfig>();
assert_eq!(config.sync_mode, SyncMode::Interval(Duration::from_secs(5)));
}
#[test]
fn test_rocksdb_config_full_params() {
let map = ConfigMap::from_config_string("versioned=true&retention=30d&sync=every")
.map_keys(|x| format!("datastore_{x}"));
let config = map.load::<RocksDbConfig>();
assert!(config.versioned);
assert_eq!(config.retention, Duration::from_secs(30 * 24 * 60 * 60));
assert_eq!(config.sync_mode, SyncMode::Every);
}
#[test]
fn test_rocksdb_config_runtime_worker_threads_default() {
let config = ConfigMap::empty().load::<RocksDbConfig>();
assert!(
config.runtime_worker_threads >= 4,
"expected default runtime_worker_threads >= 4, got {}",
config.runtime_worker_threads,
);
}
#[test]
fn test_rocksdb_config_runtime_worker_threads_override() {
let config = ConfigMap::empty()
.with_key_value("runtime_worker_threads", "1")
.load::<RocksDbConfig>();
assert_eq!(config.runtime_worker_threads, 1);
}
}