use std::fs;
use std::path::Path;
use std::sync::Arc;
use kmp_domain::PortError;
use redb::{Database, TableDefinition};
use super::format_version;
pub(crate) const NODES: TableDefinition<&str, &[u8]> = TableDefinition::new("nodes");
pub(crate) const RELATIONS: TableDefinition<(&str, &str, &str), &[u8]> =
TableDefinition::new("relations_by_source");
pub(crate) const RELATIONS_BY_TARGET: TableDefinition<(&str, &str, &str), ()> =
TableDefinition::new("relations_by_target");
pub(crate) const DETAILS: TableDefinition<&str, &[u8]> = TableDefinition::new("details");
pub(crate) const ANCHORS: TableDefinition<&str, ()> = TableDefinition::new("memory_anchors");
pub(crate) const EVENT_LOG: TableDefinition<u64, &[u8]> = TableDefinition::new("event_log");
pub(crate) const AGGREGATES: TableDefinition<&str, &[u8]> = TableDefinition::new("aggregates");
pub(crate) const IDEMPOTENCY: TableDefinition<&str, &[u8]> = TableDefinition::new("idempotency");
pub(crate) const PROCESSED: TableDefinition<(&str, &str), ()> =
TableDefinition::new("processed_events");
pub(crate) const CHECKPOINTS: TableDefinition<(&str, &str), &[u8]> =
TableDefinition::new("projection_checkpoints");
pub(crate) const SNAPSHOTS: TableDefinition<(&str, &str), &[u8]> =
TableDefinition::new("snapshots");
#[derive(Debug, Clone)]
pub struct EmbeddedKernelStore {
database: Arc<Database>,
}
impl EmbeddedKernelStore {
pub fn open(data_dir: &Path) -> Result<Self, PortError> {
fs::create_dir_all(data_dir).map_err(|error| {
PortError::Unavailable(format!(
"embedded store could not create data dir `{}`: {error}",
data_dir.display()
))
})?;
format_version::check_or_stamp(data_dir)?;
let store_file = format_version::store_file_path(data_dir);
fs::create_dir_all(store_file.parent().expect("store file has a parent")).map_err(
|error| {
PortError::Unavailable(format!(
"embedded store could not create store dir under `{}`: {error}",
data_dir.display()
))
},
)?;
Self::open_store_file(&store_file)
}
pub(crate) fn open_store_file(store_file: &Path) -> Result<Self, PortError> {
let database = Database::create(store_file).map_err(|error| {
PortError::Unavailable(format!(
"embedded store could not open `{}`: {error}",
store_file.display()
))
})?;
let store = Self {
database: Arc::new(database),
};
store.initialize_tables()?;
Ok(store)
}
fn initialize_tables(&self) -> Result<(), PortError> {
let tx = self.begin_write()?;
{
tx.open_table(NODES).map_err(table_error)?;
tx.open_table(RELATIONS).map_err(table_error)?;
tx.open_table(RELATIONS_BY_TARGET).map_err(table_error)?;
tx.open_table(DETAILS).map_err(table_error)?;
tx.open_table(ANCHORS).map_err(table_error)?;
tx.open_table(EVENT_LOG).map_err(table_error)?;
tx.open_table(AGGREGATES).map_err(table_error)?;
tx.open_table(IDEMPOTENCY).map_err(table_error)?;
tx.open_table(PROCESSED).map_err(table_error)?;
tx.open_table(CHECKPOINTS).map_err(table_error)?;
tx.open_table(SNAPSHOTS).map_err(table_error)?;
}
tx.commit().map_err(commit_error)
}
pub(crate) fn begin_write(&self) -> Result<redb::WriteTransaction, PortError> {
self.database.begin_write().map_err(|error| {
PortError::Unavailable(format!("embedded store write transaction failed: {error}"))
})
}
pub(crate) fn begin_read(&self) -> Result<redb::ReadTransaction, PortError> {
use redb::ReadableDatabase;
self.database.begin_read().map_err(|error| {
PortError::Unavailable(format!("embedded store read transaction failed: {error}"))
})
}
pub(crate) async fn run<T, F>(&self, task: F) -> Result<T, PortError>
where
T: Send + 'static,
F: FnOnce(&EmbeddedKernelStore) -> Result<T, PortError> + Send + 'static,
{
let store = self.clone();
tokio::task::spawn_blocking(move || task(&store))
.await
.map_err(|error| {
PortError::Unavailable(format!("embedded store worker failed: {error}"))
})?
}
pub async fn event_log_stats(&self) -> Result<(u64, u64), PortError> {
self.run(|store| {
let tx = store.begin_read()?;
let log = tx.open_table(EVENT_LOG).map_err(table_error)?;
let mut count = 0u64;
let mut last_sequence = 0u64;
for row in redb::ReadableTable::iter(&log).map_err(range_error)? {
let (key, _) = row.map_err(range_error)?;
count += 1;
last_sequence = key.value();
}
Ok((count, last_sequence))
})
.await
}
}
pub(crate) fn aggregate_key(root_node_id: &str, role: &str) -> String {
format!("{root_node_id}\u{1f}{role}")
}
pub(crate) fn table_error(error: redb::TableError) -> PortError {
PortError::Unavailable(format!("embedded store table access failed: {error}"))
}
pub(crate) fn storage_error(error: redb::StorageError) -> PortError {
PortError::Unavailable(format!("embedded store storage access failed: {error}"))
}
pub(crate) fn range_error(error: impl std::fmt::Display) -> PortError {
PortError::Unavailable(format!("embedded store range read failed: {error}"))
}
pub(crate) fn commit_error(error: redb::CommitError) -> PortError {
PortError::Unavailable(format!("embedded store commit failed: {error}"))
}
impl EmbeddedKernelStore {
pub fn compact_data_dir(data_dir: &Path) -> Result<bool, PortError> {
format_version::check_or_stamp(data_dir)?;
let store_file = format_version::store_file_path(data_dir);
let mut database = Database::create(&store_file).map_err(|error| {
PortError::Unavailable(format!(
"embedded store could not open `{}` for compaction: {error}",
store_file.display()
))
})?;
database.compact().map_err(|error| {
PortError::Unavailable(format!(
"embedded store compaction failed for `{}`: {error}",
store_file.display()
))
})
}
}