use std::{
path::PathBuf,
sync::{Arc, atomic::Ordering::Relaxed},
};
use wdev::Device;
use super::{
database_manager_base::DatabaseManagerBase,
garnet_database::GarnetDatabase,
i_database_manager::{HybridLogStats, IDatabaseManager},
};
use crate::storage::functions::functions_state::FunctionsState;
pub struct SingleDatabaseManager<D: Device> {
pub base: DatabaseManagerBase<D>,
pub db: Arc<GarnetDatabase<D>>,
}
impl<D: Device> SingleDatabaseManager<D> {
pub fn new(checkpoint_dir: PathBuf, db: Arc<GarnetDatabase<D>>) -> Self {
Self {
base: DatabaseManagerBase::new(checkpoint_dir),
db,
}
}
pub fn try_get_or_add_database(&self) -> wkv::Result<(Arc<GarnetDatabase<D>>, bool)> {
self.base.try_get_or_add_database(&self.db)
}
pub async fn recover_checkpoint(
&self,
replica_recover: bool,
recover_from_token: Option<u128>,
) -> wkv::Result<()> {
let _ = replica_recover;
self
.base
.recover_database_checkpoint_async(&self.db, recover_from_token)
.await
.map(|_: Option<_>| ())
}
pub fn try_pause_checkpoints(&self) -> bool {
self.base.try_pause_checkpoints(&self.db)
}
pub fn resume_checkpoints(&self) {
self.base.resume_checkpoints(&self.db);
}
pub async fn take_checkpoint(&self, background: bool) -> wkv::Result<bool> {
let _ = background;
self.base.take_database_checkpoint_async(&self.db).await
}
pub async fn take_checkpoint_helper(&self, entry_ms: u64) -> wkv::Result<bool> {
self
.base
.take_checkpoint_helper_async(&self.db, entry_ms)
.await
}
pub async fn take_on_demand_checkpoint(&self, entry_ms: u64) -> wkv::Result<()> {
self
.base
.take_on_demand_checkpoint_async(&self.db, entry_ms)
.await
}
pub async fn task_checkpoint_based_on_aof_size_limit(&self, limit: u64) -> wkv::Result<()> {
self.base.checkpoint_if_aof_exceeds(&self.db, limit).await?;
Ok(())
}
pub fn commit_to_aof(&self) -> wkv::Result<()> {
self.base.commit_aof(&self.db)
}
pub fn wait_for_commit_to_aof(&self) -> wkv::Result<bool> {
Ok(
self
.db
.aof
.as_ref()
.is_none_or(|aof| aof.flushed_until_address() >= aof.committed_until_address()),
)
}
pub async fn recover_aof(&self) -> wkv::Result<u64> {
self.base.recover_database_aof_async(&self.db).await
}
pub async fn replay_aof(&self, until: u64) -> wkv::Result<u64> {
self.base.replay_database_aof(&self.db, 0, until).await
}
pub fn grow_indexes_if_needed(&self) -> wkv::Result<bool> {
self.base.grow_index_if_needed_async(&self.db)
}
pub async fn execute_object_collection(&self) -> wkv::Result<usize> {
let session = self.db.store.new_session()?;
session.set_active_db(self.db.id.max(0) as u64);
let batch = session.enter_batch();
let storage = crate::storage::session::storage_session::StorageSession::new(batch);
storage.object_collect(|_, _| true).await
}
pub fn start_size_trackers(&self) {
}
pub fn reset_revivification_stats(&self) {
}
pub fn enqueue_commit(&self, until: u64) {
self.db.last_save_store_tail_address.store(until, Relaxed);
}
pub fn get_databases_snapshot(&self) -> Vec<Arc<GarnetDatabase<D>>> {
vec![Arc::clone(&self.db)]
}
pub async fn flush_database(&self) -> wkv::Result<()> {
self.base.reset_database(&self.db).await
}
pub async fn flush_all_databases(&self) -> wkv::Result<()> {
self.flush_database().await
}
pub async fn try_swap_databases(&self, _db_id1: i64, _db_id2: i64) -> bool {
false
}
pub fn create_functions_state(&self) -> FunctionsState {
FunctionsState::new()
}
pub fn try_pause_checkpoints_continuous(&self) -> bool {
self.try_pause_checkpoints()
}
pub async fn collect_hybrid_log_stats(&self) -> wkv::Result<Vec<(i64, HybridLogStats)>> {
self.base.collect_hybrid_log_stats(&self.db).await
}
pub fn safe_flush_aof(&self) -> wkv::Result<bool> {
Ok(match &self.db.aof {
Some(aof) => {
self
.db
.last_save_store_tail_address
.store(aof.flushed_until_address(), Relaxed);
true
}
None => false,
})
}
pub fn recover_vector_sets(&self) -> wkv::Result<u64> {
Ok(0)
}
}
impl<D: Device> IDatabaseManager<D> for SingleDatabaseManager<D> {
async fn try_get_or_add_database(
&self,
_db_id: i64,
) -> wkv::Result<(Arc<GarnetDatabase<D>>, bool)> {
SingleDatabaseManager::try_get_or_add_database(self)
}
fn try_get_database(&self, _db_id: i64) -> Option<Arc<GarnetDatabase<D>>> {
Some(Arc::clone(&self.db))
}
fn try_pause_checkpoints(&self, _db_id: i64) -> bool {
SingleDatabaseManager::try_pause_checkpoints(self)
}
fn resume_checkpoints(&self, _db_id: i64) {
SingleDatabaseManager::resume_checkpoints(self);
}
async fn recover_checkpoint_async(
&self,
replica_recover: bool,
recover_from_token: Option<u128>,
) -> wkv::Result<()> {
SingleDatabaseManager::recover_checkpoint(self, replica_recover, recover_from_token).await
}
async fn take_checkpoint_async(&self, background: bool, _db_id: i64) -> wkv::Result<bool> {
SingleDatabaseManager::take_checkpoint(self, background).await
}
async fn take_on_demand_checkpoint_async(&self, entry_ms: u64, _db_id: i64) -> wkv::Result<()> {
SingleDatabaseManager::take_on_demand_checkpoint(self, entry_ms).await
}
async fn task_checkpoint_based_on_aof_size_limit_async(
&self,
aof_size_limit: u64,
) -> wkv::Result<()> {
SingleDatabaseManager::task_checkpoint_based_on_aof_size_limit(self, aof_size_limit).await
}
fn commit_to_aof_async(&self, _db_id: i64) -> wkv::Result<()> {
SingleDatabaseManager::commit_to_aof(self)
}
fn wait_for_commit_to_aof_async(&self, _db_id: i64) -> wkv::Result<bool> {
SingleDatabaseManager::wait_for_commit_to_aof(self)
}
async fn recover_aof_async(&self) -> wkv::Result<u64> {
SingleDatabaseManager::recover_aof(self).await
}
async fn replay_aof(&self, until: u64) -> wkv::Result<u64> {
SingleDatabaseManager::replay_aof(self, until).await
}
fn grow_indexes_if_needed_async(&self) -> wkv::Result<bool> {
SingleDatabaseManager::grow_indexes_if_needed(self)
}
async fn execute_object_collection(&self, _db_id: i64) -> wkv::Result<usize> {
SingleDatabaseManager::execute_object_collection(self).await
}
fn start_size_trackers(&self) {
SingleDatabaseManager::start_size_trackers(self);
}
fn reset_revivification_stats(&self) {
SingleDatabaseManager::reset_revivification_stats(self);
}
fn enqueue_commit(&self, _db_id: i64, until: u64) {
SingleDatabaseManager::enqueue_commit(self, until);
}
fn get_databases_snapshot(&self) -> Vec<Arc<GarnetDatabase<D>>> {
SingleDatabaseManager::get_databases_snapshot(self)
}
async fn flush_database(&self, _db_id: i64) -> wkv::Result<()> {
SingleDatabaseManager::flush_database(self).await
}
async fn flush_all_databases(&self) -> wkv::Result<()> {
SingleDatabaseManager::flush_all_databases(self).await
}
async fn try_swap_databases(&self, db_id1: i64, db_id2: i64) -> bool {
SingleDatabaseManager::try_swap_databases(self, db_id1, db_id2).await
}
fn create_functions_state(&self, _db_id: i64) -> FunctionsState {
SingleDatabaseManager::create_functions_state(self)
}
async fn collect_hybrid_log_stats(&self) -> wkv::Result<Vec<(i64, HybridLogStats)>> {
SingleDatabaseManager::collect_hybrid_log_stats(self).await
}
fn recover_vector_sets(&self, _db_id: i64) -> wkv::Result<u64> {
SingleDatabaseManager::recover_vector_sets(self)
}
}