use std::collections::HashMap;
use std::sync::{Arc, RwLock, RwLockReadGuard};
use crate::store::StoreItemPart;
use crate::store::kv::KvStoreId;
use crate::util::hash::NoopU32HasherBuilder;
#[macro_use]
mod macros;
mod types;
mod count;
mod flushb;
mod flushc;
mod flusho;
mod list;
mod pop;
mod push;
mod search;
mod suggest;
pub use types::*;
pub struct Executor {
pub app_conf: Arc<crate::Config>,
pub kv_pool: crate::store::kv::KvStorePool,
pub fst_pool: crate::store::fst::FstStorePool,
pub dynamic_conf_store: Arc<DynamicConfigStore>,
}
impl std::fmt::Debug for Executor {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let Self {
kv_pool,
fst_pool,
dynamic_conf_store,
app_conf: _app_conf,
} = self;
f.debug_struct("Executor")
.field("kv_pool", kv_pool)
.field("fst_pool", fst_pool)
.field("dynamic_conf_store", &dynamic_conf_store)
.finish_non_exhaustive()
}
}
#[derive(Default)]
pub struct DynamicConfigStore(RwLock<HashMap<u32, DynamicConfig, NoopU32HasherBuilder>>);
impl DynamicConfigStore {
pub fn insert(&self, collection: StoreItemPart, config: DynamicConfig) {
(self.0.write().unwrap()).insert(collection.into_compact(), config);
}
pub fn get(&self, collection: StoreItemPart) -> Option<DynamicConfig> {
(self.0.read().unwrap())
.get(&collection.into_compact())
.copied()
}
pub fn read<'a>(&'a self) -> DynamicConfigStoreReadGuard<'a> {
DynamicConfigStoreReadGuard(self.0.read().unwrap())
}
}
impl std::fmt::Debug for DynamicConfigStore {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
use crate::util::fmt::AsPrettyRwLock;
f.debug_tuple("DynamicConfStore")
.field(&AsPrettyRwLock(&self.0))
.finish()
}
}
pub struct DynamicConfigStoreReadGuard<'a>(
RwLockReadGuard<'a, HashMap<u32, DynamicConfig, NoopU32HasherBuilder>>,
);
impl<'a> DynamicConfigStoreReadGuard<'a> {
pub fn iter(&self) -> impl Iterator<Item = (&u32, &DynamicConfig)> {
self.0.iter()
}
pub fn get(&self, key: &u32) -> Option<&DynamicConfig> {
self.0.get(key)
}
}
#[derive(Debug, Default, Clone, Copy)]
pub struct DynamicConfig {
pub sonic: DynamicConfigSonic,
pub rocksdb: DynamicConfigRocksDb,
}
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
pub struct DynamicConfigSonic {
pub disable_janitor_tasks: Option<bool>,
pub disable_fst_consolidate_task: Option<bool>,
pub disable_kv_flush_task: Option<bool>,
}
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
pub struct DynamicConfigRocksDb {
pub disable_auto_compactions: Option<bool>,
pub unordered_write: Option<bool>,
pub memtable: Option<RocksDbMemtable>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RocksDbMemtable {
Default,
Vector,
}
impl Executor {
pub fn set_dynamic_conf(
&self,
collection: StoreItemPart,
new_conf: DynamicConfig,
) -> Result<(), Box<dyn std::error::Error>> {
let kv_store_id = KvStoreId::from_part(collection);
tracing::debug!(
?new_conf.rocksdb,
"Re-opening KV store connection for {kv_store_id:?} with new dynamic configuration overrides…"
);
let mut kv_pool_write_guard = self.kv_pool.write().unwrap();
self.kv_pool
.close(kv_store_id, Some(&mut kv_pool_write_guard));
self.kv_pool
.acquire(
true,
collection,
Some(&mut kv_pool_write_guard),
|options| {
let DynamicConfigRocksDb {
disable_auto_compactions,
unordered_write,
memtable,
} = &new_conf.rocksdb;
if let Some(disable_auto_compactions) = disable_auto_compactions {
options.set_disable_auto_compactions(*disable_auto_compactions);
}
if let Some(unordered_write) = unordered_write {
options.set_unordered_write(*unordered_write);
}
match memtable {
Some(RocksDbMemtable::Vector) => {
options.set_memtable_factory(rocksdb::MemtableFactory::Vector);
options.set_allow_concurrent_memtable_write(false);
}
None | Some(RocksDbMemtable::Default) => {}
}
},
)
.map_err(|()| std::io::Error::other("Error re-opening connection"))?;
drop(kv_pool_write_guard);
tracing::info!(
?new_conf.rocksdb,
"KV store connection for {collection:?} successfully re-opened"
);
self.dynamic_conf_store.insert(collection, new_conf);
Ok(())
}
}