use std::path::PathBuf;
use std::sync::Arc;
use std::time::Duration;
use pigeonhole_engine::{
CachePriority, CompactionStyle, Compression, EngineOptions, FamilyKind, FamilyOptions,
};
use pigeonhole_format::Durability;
use pigeonhole_io::VfsRef;
use crate::MergeOperator;
pub(crate) const I64_ADD: &str = "pigeonhole.i64_add";
pub fn days(n: u64) -> Duration {
Duration::from_secs(n.saturating_mul(86_400))
}
#[derive(Debug, Clone)]
pub struct Options {
durability: Durability,
shards: usize,
compaction_cores: usize,
pin_threads: bool,
memtable_budget: u64,
block_cache: Option<usize>,
row_cache: usize,
shm_dir: Option<PathBuf>,
create_if_missing: bool,
merge_operators: Vec<Arc<dyn MergeOperator>>,
allow_unregistered_merge_operators: bool,
allow_fuse: bool,
vfs: Option<VfsRef>,
wal_segment_size: Option<u64>,
tablet_changes: bool,
tablet_balance: Option<(Duration, u64, u64)>,
write_stall_timeout: Option<Duration>,
}
impl Default for Options {
fn default() -> Self {
Self {
durability: Durability::GroupSync,
shards: 0,
compaction_cores: 0,
pin_threads: false,
memtable_budget: 64 << 20,
block_cache: None,
row_cache: 0,
shm_dir: None,
create_if_missing: true,
merge_operators: Vec::new(),
allow_unregistered_merge_operators: false,
allow_fuse: false,
vfs: None,
wal_segment_size: None,
tablet_changes: true,
tablet_balance: None,
write_stall_timeout: None,
}
}
}
impl Options {
pub fn durability(mut self, durability: Durability) -> Self {
self.durability = durability;
self
}
pub fn shards(mut self, n: usize) -> Self {
self.shards = n;
self
}
pub fn compaction_cores(mut self, k: usize) -> Self {
self.compaction_cores = k;
self
}
pub fn pin_threads(mut self, yes: bool) -> Self {
self.pin_threads = yes;
self
}
pub fn memtable_budget(mut self, bytes: u64) -> Self {
self.memtable_budget = bytes;
self
}
pub fn write_stall_timeout(mut self, timeout: Duration) -> Self {
self.write_stall_timeout = Some(timeout);
self
}
pub fn block_cache(mut self, bytes: usize) -> Self {
self.block_cache = Some(bytes);
self
}
pub fn row_cache(mut self, bytes: usize) -> Self {
self.row_cache = bytes;
self
}
pub fn shm_dir(mut self, dir: impl Into<PathBuf>) -> Self {
self.shm_dir = Some(dir.into());
self
}
pub fn create_if_missing(mut self, yes: bool) -> Self {
self.create_if_missing = yes;
self
}
pub fn merge_operator(mut self, op: Arc<dyn MergeOperator>) -> Self {
self.merge_operators.push(op);
self
}
pub fn allow_unregistered_merge_operators(mut self, yes: bool) -> Self {
self.allow_unregistered_merge_operators = yes;
self
}
pub fn allow_fuse(mut self, yes: bool) -> Self {
self.allow_fuse = yes;
self
}
pub fn tablet_changes(mut self, yes: bool) -> Self {
self.tablet_changes = yes;
self
}
#[doc(hidden)]
pub fn vfs(mut self, vfs: pigeonhole_io::VfsRef) -> Self {
self.vfs = Some(vfs);
self
}
#[doc(hidden)]
pub fn wal_segment_size(mut self, bytes: u64) -> Self {
self.wal_segment_size = Some(bytes);
self
}
#[doc(hidden)]
pub fn tablet_balance(mut self, interval: Duration, min_writes: u64, split_bytes: u64) -> Self {
self.tablet_balance = Some((interval, min_writes, split_bytes));
self
}
}
#[derive(Debug, Clone, Default)]
pub struct ReaderOptions {
block_cache: Option<usize>,
shm_dir: Option<PathBuf>,
merge_operators: Vec<Arc<dyn MergeOperator>>,
allow_fuse: bool,
vfs: Option<VfsRef>,
}
impl ReaderOptions {
pub fn block_cache(mut self, bytes: usize) -> Self {
self.block_cache = Some(bytes);
self
}
pub fn shm_dir(mut self, dir: impl Into<PathBuf>) -> Self {
self.shm_dir = Some(dir.into());
self
}
pub fn merge_operator(mut self, op: Arc<dyn MergeOperator>) -> Self {
self.merge_operators.push(op);
self
}
pub fn allow_fuse(mut self, yes: bool) -> Self {
self.allow_fuse = yes;
self
}
#[doc(hidden)]
pub fn vfs(mut self, vfs: pigeonhole_io::VfsRef) -> Self {
self.vfs = Some(vfs);
self
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Default)]
pub enum Priority {
Low,
#[default]
Normal,
High,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
pub enum Compaction {
#[default]
Leveled,
Tiered,
FifoByTime,
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub struct Family {
options: FamilyOptions,
}
impl Family {
pub fn counter() -> Self {
Self {
options: FamilyOptions::default()
.merge_operator(I64_ADD.to_owned())
.kind(FamilyKind::Counter),
}
}
pub fn max_versions(mut self, n: u32) -> Self {
self.options.max_versions = n;
self
}
pub fn ttl(mut self, ttl: Duration) -> Self {
self.options.ttl_micros = u64::try_from(ttl.as_micros()).unwrap_or(u64::MAX);
self
}
pub fn bloom_bits(mut self, bits: u8) -> Self {
self.options.bloom_bits = bits;
self
}
pub fn blob_threshold(mut self, bytes: u32) -> Self {
self.options.blob_threshold = bytes;
self
}
pub fn lz4(mut self) -> Self {
self.options.compression = Compression::Lz4;
self
}
pub fn zstd(mut self, level: i8) -> Self {
self.options.compression = Compression::Zstd;
self.options.compression_level = level;
self
}
pub fn uncompressed(mut self) -> Self {
self.options.compression = Compression::None;
self
}
pub fn block_size(mut self, bytes: u32) -> Self {
self.options.block_size = bytes;
self
}
pub fn merge_operator(mut self, name: &str) -> Self {
name.clone_into(&mut self.options.merge_operator);
self
}
pub fn cache_priority(mut self, priority: Priority) -> Self {
self.options.cache_priority = match priority {
Priority::Low => CachePriority::Low,
Priority::Normal => CachePriority::Normal,
Priority::High => CachePriority::High,
};
self
}
pub fn compaction(mut self, strategy: Compaction) -> Self {
self.options.compaction = match strategy {
Compaction::Leveled => CompactionStyle::Leveled,
Compaction::Tiered => CompactionStyle::Tiered,
Compaction::FifoByTime => CompactionStyle::FifoByTime,
};
self
}
}
impl Options {
pub(crate) fn to_engine(&self) -> EngineOptions {
let vfs = self.vfs.clone().unwrap_or_else(default_vfs);
let mut o = EngineOptions::new(vfs);
o.create_if_missing = self.create_if_missing;
o.shards = self.shards;
o.compaction_threads = self.compaction_cores;
o.pin_threads = self.pin_threads;
o.durability = self.durability;
o.memtable_budget = self.memtable_budget;
o.memtable_freeze_bytes = self.memtable_budget / 4;
if let Some(bytes) = self.block_cache {
o.block_cache_bytes = bytes;
}
o.row_cache_bytes = self.row_cache;
o.shm_dir.clone_from(&self.shm_dir);
o.allow_unregistered_merge = self.allow_unregistered_merge_operators;
o.allow_fuse = self.allow_fuse;
if let Some(bytes) = self.wal_segment_size {
o.wal.segment_size = bytes;
}
o.tablet_changes = self.tablet_changes;
if let Some(timeout) = self.write_stall_timeout {
o.write_stall_timeout_nanos = u64::try_from(timeout.as_nanos()).unwrap_or(u64::MAX);
}
if let Some((interval, min_writes, split_bytes)) = self.tablet_balance {
o.balance_interval_nanos = u64::try_from(interval.as_nanos()).unwrap_or(u64::MAX);
o.balance_min_writes = min_writes;
o.tablet_split_bytes = split_bytes;
}
for op in &self.merge_operators {
o.merge_operators.register(Arc::clone(op));
}
o
}
}
impl ReaderOptions {
pub(crate) fn to_engine(&self) -> EngineOptions {
let vfs = self.vfs.clone().unwrap_or_else(default_vfs);
let mut o = EngineOptions::new(vfs);
if let Some(bytes) = self.block_cache {
o.block_cache_bytes = bytes;
}
o.shm_dir.clone_from(&self.shm_dir);
o.allow_fuse = self.allow_fuse;
for op in &self.merge_operators {
o.merge_operators.register(Arc::clone(op));
}
o
}
}
impl Family {
pub(crate) fn to_engine(&self) -> FamilyOptions {
self.options.clone()
}
}
fn default_vfs() -> VfsRef {
pigeonhole_io::pread::PreadVfs::new(0)
}