pub use crate::db::builder::CloneBuilder;
pub use crate::db::builder::CloneSourceSpec;
use std::collections::BTreeSet;
use crate::checkpoint::{Checkpoint, CheckpointCreateResult};
use crate::compactions_store::CompactionsStore;
use crate::compactor::{Compaction, CompactionSpec, Compactor, CompactorStateView};
use crate::compactor_state::VersionedCompactions;
use crate::compactor_state_protocols::CompactorStateReader;
use crate::config::{CheckpointOptions, GarbageCollectorOptions};
use crate::db::builder::GarbageCollectorBuilder;
use crate::error::SlateDBError;
use crate::manifest::store::{ManifestStore, StoredManifest};
use crate::manifest::VersionedManifest;
use slatedb_common::clock::SystemClock;
use crate::object_stores::{ObjectStoreType, ObjectStores};
use crate::retrying_object_store::RetryingObjectStore;
use crate::seq_tracker::FindOption;
use crate::utils::IdGenerator;
use bytes::Bytes;
use chrono::{DateTime, Utc};
use futures::StreamExt;
use object_store::path::Path;
use object_store::{ObjectStore, ObjectStoreExt};
use rand::RngCore;
use slatedb_common::DbRand;
use std::env;
use std::env::VarError;
use std::ops::{Bound, RangeBounds};
use std::sync::Arc;
use std::time::Duration;
use tokio_util::sync::CancellationToken;
use ulid::Ulid;
use uuid::Uuid;
pub use crate::db::builder::AdminBuilder;
use crate::merge_operator::MergeOperatorType;
use crate::wal::WalAdmin;
use slatedb_txn_obj::TransactionalObject;
pub struct Admin {
pub(crate) path: Path,
pub(crate) object_stores: ObjectStores,
pub(crate) system_clock: Arc<dyn SystemClock>,
pub(crate) rand: Arc<DbRand>,
pub(crate) object_store_max_retries: Option<u32>,
#[cfg(feature = "compaction_filters")]
pub(crate) compaction_filter_supplier:
Option<Arc<dyn crate::compaction_filter::CompactionFilterSupplier>>,
pub(crate) merge_operator: Option<MergeOperatorType>,
pub(crate) wal_admin: Arc<dyn WalAdmin>,
}
impl Admin {
pub async fn read_manifest(
&self,
maybe_id: Option<u64>,
) -> Result<Option<VersionedManifest>, crate::Error> {
let manifest_store = self.manifest_store();
let manifest = if let Some(id) = maybe_id {
manifest_store
.try_read_manifest(id)
.await
.map_err(crate::Error::from)?
.map(|manifest| VersionedManifest::from_manifest(id, manifest))
} else {
manifest_store
.try_read_latest_manifest()
.await
.map_err(crate::Error::from)?
};
Ok(manifest)
}
pub async fn list_manifests<R: RangeBounds<u64>>(
&self,
range: R,
) -> Result<Vec<VersionedManifest>, crate::Error> {
let manifest_store = self.manifest_store();
let manifest_metadata = manifest_store
.list_manifests(range)
.await
.map_err(crate::Error::from)?;
let mut manifests = Vec::with_capacity(manifest_metadata.len());
for metadata in manifest_metadata {
match manifest_store
.try_read_manifest(metadata.id)
.await
.map_err(crate::Error::from)?
{
Some(manifest) => {
manifests.push(VersionedManifest::from_manifest(metadata.id, manifest))
}
None => log::warn!(
"listed manifest missing on read, skipping [id={}]",
metadata.id
),
}
}
Ok(manifests)
}
pub async fn read_compactions(
&self,
maybe_id: Option<u64>,
) -> Result<Option<VersionedCompactions>, crate::Error> {
let compactions_store = self.compactions_store();
let compactions = if let Some(id) = maybe_id {
compactions_store
.try_read_compactions(id)
.await
.map_err(crate::Error::from)?
.map(|compactions| VersionedCompactions::from_compactions(id, compactions))
} else {
compactions_store
.try_read_latest_compactions()
.await
.map_err(crate::Error::from)?
};
Ok(compactions)
}
pub async fn read_compaction(
&self,
compaction_id: Ulid,
maybe_id: Option<u64>,
) -> Result<Option<Compaction>, crate::Error> {
let compactions_store = self.compactions_store();
let compactions = if let Some(compactions_id) = maybe_id {
compactions_store
.try_read_compactions(compactions_id)
.await
.map_err(crate::Error::from)?
} else {
compactions_store
.try_read_latest_compactions()
.await
.map_err(crate::Error::from)?
.map(|compactions| compactions.compactions)
};
let Some(compactions) = compactions else {
return Ok(None);
};
let Some(compaction) = compactions.get(&compaction_id) else {
return Ok(None);
};
Ok(Some(compaction.clone()))
}
pub async fn read_compactor_state_view(&self) -> Result<CompactorStateView, crate::Error> {
let manifest_store = Arc::new(self.manifest_store());
let compactions_store = Arc::new(self.compactions_store());
let reader = CompactorStateReader::new(&manifest_store, &compactions_store);
reader.read_view().await.map_err(crate::Error::from)
}
pub async fn submit_compaction(
&self,
spec: CompactionSpec,
) -> Result<Compaction, crate::Error> {
let compactions_store = Arc::new(self.compactions_store());
let rand = Arc::new(DbRand::new(self.rand.rng().next_u64()));
let compaction_id =
Compactor::submit(spec, compactions_store, rand, self.system_clock.clone()).await?;
let Some(compaction) = self.read_compaction(compaction_id, None).await? else {
return Err(crate::Error::from(SlateDBError::InvalidDBState));
};
Ok(compaction)
}
pub async fn list_compactions<R: RangeBounds<u64>>(
&self,
range: R,
) -> Result<Vec<VersionedCompactions>, crate::Error> {
let compactions_store = self.compactions_store();
let compactions_metadata = compactions_store
.list_compactions(range)
.await
.map_err(crate::Error::from)?;
let mut compactions = Vec::with_capacity(compactions_metadata.len());
for metadata in compactions_metadata {
let stored_compactions = compactions_store
.read_compactions(metadata.id)
.await
.map_err(crate::Error::from)?;
compactions.push(VersionedCompactions::from_compactions(
metadata.id,
stored_compactions,
));
}
Ok(compactions)
}
pub async fn list_checkpoints(
&self,
name_filter: Option<&str>,
) -> Result<Vec<Checkpoint>, crate::Error> {
let manifest_store = self.manifest_store();
let manifest = manifest_store
.read_latest_manifest()
.await
.map_err(crate::Error::from)?
.manifest;
let checkpoints = match name_filter {
Some("") => manifest
.core
.checkpoints
.into_iter()
.filter(|cp| cp.name.as_deref() == Some("") || cp.name.is_none())
.collect(),
Some(name) => manifest
.core
.checkpoints
.into_iter()
.filter(|cp| cp.name.as_deref() == Some(name))
.collect(),
None => manifest.core.checkpoints,
};
Ok(checkpoints)
}
pub async fn run_gc_once(&self, gc_opts: GarbageCollectorOptions) -> Result<(), crate::Error> {
let gc = GarbageCollectorBuilder::new(
self.path.clone(),
self.object_stores.store_of(ObjectStoreType::Main).clone(),
)
.with_system_clock(self.system_clock.clone())
.with_wal_gc(self.wal_admin.garbage_collector(&self.path))
.with_wal_object_store(self.object_stores.store_of(ObjectStoreType::Wal).clone())
.with_options(gc_opts)
.with_seed(self.rand.rng().next_u64())
.build();
gc.run_gc_once().await;
Ok(())
}
pub async fn run_gc(&self, cancellation_token: CancellationToken) -> Result<(), crate::Error> {
self.run_gc_with_options(cancellation_token, GarbageCollectorOptions::default())
.await
}
pub async fn run_gc_with_options(
&self,
cancellation_token: CancellationToken,
gc_opts: GarbageCollectorOptions,
) -> Result<(), crate::Error> {
let gc = GarbageCollectorBuilder::new(
self.path.clone(),
self.object_stores.store_of(ObjectStoreType::Main).clone(),
)
.with_system_clock(self.system_clock.clone())
.with_wal_gc(self.wal_admin.garbage_collector(&self.path))
.with_wal_object_store(self.object_stores.store_of(ObjectStoreType::Wal).clone())
.with_options(gc_opts)
.with_seed(self.rand.rng().next_u64())
.build();
gc.start()?;
tokio::select! {
result = gc.join() => result,
_ = cancellation_token.cancelled() => {
gc.stop().await
}
}
}
pub async fn run_compactor(
&self,
cancellation_token: CancellationToken,
) -> Result<(), crate::Error> {
self.run_compactor_with_options(
cancellation_token,
crate::config::CompactorOptions::default(),
)
.await
}
pub async fn run_compactor_with_options(
&self,
cancellation_token: CancellationToken,
options: crate::config::CompactorOptions,
) -> Result<(), crate::Error> {
#[allow(unused_mut)]
let mut builder = crate::CompactorBuilder::new(
self.path.clone(),
self.object_stores.store_of(ObjectStoreType::Main).clone(),
)
.with_options(options)
.with_system_clock(self.system_clock.clone())
.with_seed(self.rand.rng().next_u64());
#[cfg(feature = "compaction_filters")]
if let Some(supplier) = &self.compaction_filter_supplier {
builder = builder.with_compaction_filter_supplier(supplier.clone());
}
if let Some(merge_operator) = &self.merge_operator {
builder = builder.with_merge_operator(merge_operator.clone());
}
let compactor = builder.build();
compactor.start().await?;
tokio::select! {
result = compactor.join() => result,
_ = cancellation_token.cancelled() => {
compactor.stop().await
}
}
}
pub async fn run_compaction_worker(
&self,
cancellation_token: CancellationToken,
) -> Result<(), crate::Error> {
self.run_compaction_worker_with_options(
cancellation_token,
crate::config::CompactionWorkerOptions::default(),
)
.await
}
pub async fn run_compaction_worker_with_options(
&self,
cancellation_token: CancellationToken,
options: crate::config::CompactionWorkerOptions,
) -> Result<(), crate::Error> {
#[allow(unused_mut)]
let mut builder = crate::CompactionWorkerBuilder::new(
self.path.clone(),
self.object_stores.store_of(ObjectStoreType::Main).clone(),
)
.with_options(options)
.with_system_clock(self.system_clock.clone())
.with_seed(self.rand.rng().next_u64());
#[cfg(feature = "compaction_filters")]
if let Some(supplier) = &self.compaction_filter_supplier {
builder = builder.with_compaction_filter_supplier(supplier.clone());
}
let worker = builder.build().await?;
worker.start()?;
tokio::select! {
result = worker.join() => result,
_ = cancellation_token.cancelled() => {
worker.stop().await
}
}
}
pub async fn create_detached_checkpoint(
&self,
options: &CheckpointOptions,
) -> Result<CheckpointCreateResult, crate::Error> {
let manifest_store = Arc::new(self.manifest_store());
let mut stored_manifest =
StoredManifest::load(manifest_store, self.system_clock.clone()).await?;
let configured_wal_uri = self.object_stores.has_wal_object_store().then(String::new);
stored_manifest
.db_state()
.validate_wal_object_store_uri(configured_wal_uri.as_deref())?;
let checkpoint_id = self.rand.rng().gen_uuid();
let checkpoint = stored_manifest
.write_checkpoint(checkpoint_id, options)
.await?;
Ok(CheckpointCreateResult {
id: checkpoint.id,
manifest_id: checkpoint.manifest_id,
})
}
pub async fn refresh_checkpoint(
&self,
id: Uuid,
lifetime: Option<Duration>,
) -> Result<(), crate::Error> {
let manifest_store = Arc::new(self.manifest_store());
let mut stored_manifest =
StoredManifest::load(manifest_store, self.system_clock.clone()).await?;
stored_manifest
.maybe_apply_update(|stored_manifest| {
let mut dirty = stored_manifest.prepare_dirty()?;
let expire_time = lifetime.map(|l| self.system_clock.now() + l);
let Some(_) = dirty.value.core.checkpoints.iter_mut().find_map(|c| {
if c.id == id {
c.expire_time = expire_time;
return Some(());
}
None
}) else {
return Err(SlateDBError::InvalidDBState);
};
Ok(Some(dirty))
})
.await
.map_err(Into::into)
}
pub async fn delete_checkpoint(&self, id: Uuid) -> Result<(), crate::Error> {
let manifest_store = Arc::new(self.manifest_store());
let mut stored_manifest =
StoredManifest::load(manifest_store, self.system_clock.clone()).await?;
stored_manifest
.maybe_apply_update(|stored_manifest| {
let mut dirty = stored_manifest.prepare_dirty()?;
let checkpoints: Vec<Checkpoint> = dirty
.value
.core
.checkpoints
.iter()
.filter(|c| c.id != id)
.cloned()
.collect();
dirty.value.core.checkpoints = checkpoints;
Ok(Some(dirty))
})
.await
.map_err(Into::into)
}
pub async fn delete_db(&self, confirm: bool) -> Result<Vec<String>, crate::Error> {
let main = self.retrying_store(ObjectStoreType::Main);
if !confirm {
return self.list_prefix(&main).await;
}
let marker = self.path.clone().join(".deleting");
let marker_exists = main.get(&marker).await.map(|_| true).or_else(|e| match e {
object_store::Error::NotFound { .. } => Ok(false),
other => Err(SlateDBError::from(other)),
})?;
let manifest = self.manifest_store().try_read_latest_manifest().await?;
if manifest.is_none()
&& !marker_exists
&& !collect_prefix(&main, &self.path).await?.is_empty()
{
return Err(SlateDBError::InvalidDBState.into());
}
if let Some(manifest) = manifest.as_ref() {
for external_db in manifest.external_dbs() {
let Some(final_checkpoint_id) = external_db.final_checkpoint_id else {
continue;
};
let parent_store = Arc::new(ManifestStore::new(
&Path::from(external_db.path.as_str()),
self.retrying_store(ObjectStoreType::Main),
));
let mut parent =
match StoredManifest::load(parent_store, self.system_clock.clone()).await {
Ok(parent) => parent,
Err(SlateDBError::LatestTransactionalObjectVersionMissing) => continue,
Err(e) => return Err(e.into()),
};
parent.delete_checkpoint(final_checkpoint_id).await?;
}
}
if !marker_exists {
main.put(&marker, Bytes::new().into())
.await
.map_err(SlateDBError::from)?;
}
let mut deleted = self.delete_prefix(&main, Some(&marker)).await?;
deleted.extend(
self.wal_admin
.delete_wal(&self.path, false)
.await
.map_err(SlateDBError::from)?,
);
main.delete(&marker).await.map_err(SlateDBError::from)?;
deleted.push(marker.to_string());
Ok(deleted)
}
async fn list_prefix(&self, main: &Arc<dyn ObjectStore>) -> Result<Vec<String>, crate::Error> {
let paths = collect_prefix(main, &self.path).await?;
let mut paths = paths.iter().map(Path::to_string).collect::<BTreeSet<_>>();
let wal_paths = self
.wal_admin
.delete_wal(&self.path, true)
.await
.map_err(SlateDBError::from)?;
for wp in wal_paths {
paths.insert(wp);
}
Ok(paths.into_iter().collect())
}
async fn delete_prefix(
&self,
store: &Arc<dyn ObjectStore>,
keep: Option<&Path>,
) -> Result<Vec<String>, crate::Error> {
let mut deleted = Vec::new();
for path in collect_prefix(store, &self.path).await? {
if Some(&path) == keep {
continue;
}
store.delete(&path).await.map_err(SlateDBError::from)?;
deleted.push(path);
}
Ok(deleted.into_iter().map(|p| p.to_string()).collect())
}
pub async fn get_timestamp_for_sequence(
&self,
seq: u64,
round_up: bool,
) -> Result<Option<DateTime<Utc>>, crate::Error> {
let manifest_store = self.manifest_store();
let id_manifest = manifest_store.try_read_latest_manifest().await?;
let Some(manifest) = id_manifest else {
return Ok(None);
};
let opt = if round_up {
FindOption::RoundUp
} else {
FindOption::RoundDown
};
Ok(manifest.core().sequence_tracker.find_ts(seq, opt))
}
pub async fn get_sequence_for_timestamp(
&self,
ts: DateTime<Utc>,
round_up: bool,
) -> Result<Option<u64>, crate::Error> {
let manifest_store = self.manifest_store();
let id_manifest = manifest_store.try_read_latest_manifest().await?;
let Some(manifest) = id_manifest else {
return Ok(None);
};
let opt = if round_up {
FindOption::RoundUp
} else {
FindOption::RoundDown
};
Ok(manifest.core().sequence_tracker.find_seq(ts, opt))
}
fn retrying_store(&self, store_type: ObjectStoreType) -> Arc<dyn ObjectStore> {
Arc::new(RetryingObjectStore::new(
self.object_stores.store_of(store_type).clone(),
self.rand.clone(),
self.system_clock.clone(),
self.object_store_max_retries,
))
}
fn manifest_store(&self) -> ManifestStore {
ManifestStore::new(&self.path, self.retrying_store(ObjectStoreType::Main))
}
fn compactions_store(&self) -> CompactionsStore {
CompactionsStore::new(&self.path, self.retrying_store(ObjectStoreType::Main))
}
pub fn create_clone_builder_from_source(
&self,
source: CloneSourceSpec<(Bound<Bytes>, Bound<Bytes>)>,
) -> CloneBuilder<(Bound<Bytes>, Bound<Bytes>)> {
CloneBuilder::new(
self.path.clone(),
source,
self.retrying_store(ObjectStoreType::Main),
)
.with_wal_admin(self.wal_admin.clone())
}
pub fn builder<P: Into<Path>>(path: P, object_store: Arc<dyn ObjectStore>) -> AdminBuilder<P> {
AdminBuilder::new(path, object_store)
}
}
fn get_env_variable(name: &str) -> Result<String, SlateDBError> {
env::var(name).map_err(|e| match e {
VarError::NotPresent => SlateDBError::InvalidEnvironmentVariable {
key: name.to_string(),
value: None,
},
VarError::NotUnicode(not_unicode_value) => SlateDBError::InvalidEnvironmentVariable {
key: name.to_string(),
value: Some(format!("{:?}", not_unicode_value)),
},
})
}
pub fn load_object_store_from_env(
env_file: Option<String>,
) -> Result<Arc<dyn ObjectStore>, crate::Error> {
dotenvy::from_filename(env_file.unwrap_or(String::from(".env"))).ok();
let cloud_provider = get_env_variable("CLOUD_PROVIDER")?;
match cloud_provider.to_lowercase().as_str() {
"local" => load_local(),
"memory" => load_memory(),
#[cfg(feature = "aws")]
"aws" => load_aws(),
#[cfg(feature = "azure")]
"azure" => load_azure(),
#[cfg(feature = "gcp")]
"gcp" => load_gcp(),
invalid_value => Err(SlateDBError::InvalidEnvironmentVariable {
key: "CLOUD_PROVIDER".to_string(),
value: Some(invalid_value.to_string()),
}
.into()),
}
}
pub fn load_local() -> Result<Arc<dyn ObjectStore>, crate::Error> {
let local_path = get_env_variable("LOCAL_PATH")?;
let lfs =
object_store::local::LocalFileSystem::new_with_prefix(local_path).map_err(|error| {
SlateDBError::ObjectStoreError(Arc::new(object_store::Error::Generic {
store: "local",
source: Box::new(error),
}))
})?;
Ok(Arc::new(lfs) as Arc<dyn ObjectStore>)
}
pub fn load_memory() -> Result<Arc<dyn ObjectStore>, crate::Error> {
Ok(Arc::new(object_store::memory::InMemory::new()) as Arc<dyn ObjectStore>)
}
#[cfg(feature = "aws")]
pub fn load_aws() -> Result<Arc<dyn ObjectStore>, crate::Error> {
let builder = object_store::aws::AmazonS3Builder::from_env();
Ok(Arc::new(builder.build().map_err(|error| {
SlateDBError::ObjectStoreError(Arc::new(object_store::Error::Generic {
store: "AmazonS3",
source: Box::new(error),
}))
})?) as Arc<dyn ObjectStore>)
}
#[cfg(feature = "azure")]
pub fn load_azure() -> Result<Arc<dyn ObjectStore>, crate::Error> {
let builder = object_store::azure::MicrosoftAzureBuilder::from_env();
Ok(Arc::new(builder.build().map_err(|error| {
SlateDBError::ObjectStoreError(Arc::new(object_store::Error::Generic {
store: "MicrosoftAzure",
source: Box::new(error),
}))
})?) as Arc<dyn ObjectStore>)
}
#[cfg(feature = "gcp")]
pub fn load_gcp() -> Result<Arc<dyn ObjectStore>, crate::Error> {
let builder = object_store::gcp::GoogleCloudStorageBuilder::from_env();
Ok(Arc::new(builder.build().map_err(|error| {
SlateDBError::ObjectStoreError(Arc::new(object_store::Error::Generic {
store: "GoogleCloudStorage",
source: Box::new(error),
}))
})?) as Arc<dyn ObjectStore>)
}
async fn collect_prefix(
store: &Arc<dyn ObjectStore>,
prefix: &Path,
) -> Result<Vec<Path>, crate::Error> {
let mut listing = store.list(Some(prefix));
let mut paths = Vec::new();
while let Some(meta) = listing
.next()
.await
.transpose()
.map_err(SlateDBError::from)?
{
paths.push(meta.location);
}
Ok(paths)
}
#[cfg(test)]
mod tests {
use crate::admin::{load_object_store_from_env, AdminBuilder};
use crate::compactions_store::{CompactionsStore, StoredCompactions};
use crate::compactor_state::{Compaction, CompactionSpec, CompactionStatus, SourceId};
use crate::config::{
CheckpointOptions, CompactionWorkerOptions, CompactorOptions, GarbageCollectorOptions,
};
use crate::manifest::store::{ManifestStore, StoredManifest};
use crate::manifest::ManifestCore;
use crate::test_utils::{FlakyObjectStore, StringConcatMergeOperator};
use crate::ErrorKind;
use object_store::memory::InMemory;
use object_store::path::Path;
use object_store::ObjectStore;
use slatedb_common::clock::DefaultSystemClock;
use std::sync::Arc;
use tokio_util::sync::CancellationToken;
use ulid::Ulid;
#[test]
fn test_load_object_store_from_env() {
figment::Jail::expect_with(|jail| {
let err = load_object_store_from_env(None).expect_err("expected invalid env error");
assert_eq!(err.kind(), ErrorKind::Invalid);
assert_eq!(
err.to_string(),
"Invalid error: invalid environment variable CLOUD_PROVIDER value `null`"
);
jail.create_file("invalid.env", "CLOUD_PROVIDER=invalid")
.expect("failed to create temp env file");
let err = load_object_store_from_env(Some("invalid.env".to_string()))
.expect_err("expected invalid provider error");
assert_eq!(err.kind(), ErrorKind::Invalid);
assert_eq!(
err.to_string(),
"Invalid error: invalid environment variable CLOUD_PROVIDER value `invalid`"
);
std::env::remove_var("CLOUD_PROVIDER");
jail.create_file("memory.env", "CLOUD_PROVIDER=memory")
.expect("failed to create temp env file");
let r = load_object_store_from_env(Some("memory.env".to_string()));
let store = r.expect("expected memory object store");
assert_eq!(store.to_string(), "InMemory");
Ok(())
});
}
#[test]
fn test_load_local_invalid_path_maps_to_unavailable() {
figment::Jail::expect_with(|jail| {
jail.set_env("LOCAL_PATH", "missing-local-path");
let err = super::load_local().expect_err("expected invalid local-path error");
assert_eq!(err.kind(), ErrorKind::Unavailable);
Ok(())
});
}
#[tokio::test]
async fn test_admin_read_manifest() {
let object_store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let path = Path::from("/tmp/test_admin_read_manifest");
let manifest_store = Arc::new(ManifestStore::new(&path, object_store.clone()));
let mut stored = StoredManifest::create_new_db(
manifest_store,
ManifestCore::new(),
Arc::new(DefaultSystemClock::new()),
)
.await
.unwrap();
let mut dirty = stored.prepare_dirty().unwrap();
dirty.value.core.next_wal_sst_id = 17;
dirty.value.core.last_l0_seq = 9;
dirty.value.writer_epoch = 3;
dirty.value.compactor_epoch = 5;
stored.update(dirty).await.unwrap();
let admin = AdminBuilder::new(path.clone(), object_store).build();
let latest = admin
.read_manifest(None)
.await
.unwrap()
.expect("expected manifest");
assert_eq!(latest.id, 2);
assert_eq!(latest.manifest.writer_epoch, 3);
assert_eq!(latest.manifest.compactor_epoch, 5);
assert_eq!(latest.manifest.core.next_wal_sst_id, 17);
assert_eq!(latest.manifest.core.last_l0_seq, 9);
let first = admin
.read_manifest(Some(1))
.await
.unwrap()
.expect("expected manifest");
assert_eq!(first.id, 1);
assert_eq!(first.manifest.writer_epoch, 0);
assert_eq!(first.manifest.compactor_epoch, 0);
assert_eq!(first.manifest.core.next_wal_sst_id, 1);
assert_eq!(first.manifest.core.last_l0_seq, 0);
}
#[tokio::test(start_paused = true)]
async fn test_admin_run_gc_with_options_stops_on_cancellation() {
let object_store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let path = Path::from("/tmp/test_admin_run_gc_with_options_stops_on_cancellation");
let admin = AdminBuilder::new(path, object_store).build();
let cancellation_token = CancellationToken::new();
cancellation_token.cancel();
let result = admin
.run_gc_with_options(cancellation_token, GarbageCollectorOptions::default())
.await;
assert!(matches!(result, Ok(())));
}
#[tokio::test(start_paused = true)]
async fn test_admin_run_compactor_with_options_stops_on_cancellation() {
let object_store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let path = Path::from("/tmp/test_admin_run_compactor_with_options_stops_on_cancellation");
StoredManifest::create_new_db(
Arc::new(ManifestStore::new(&path, object_store.clone())),
ManifestCore::new(),
Arc::new(DefaultSystemClock::new()),
)
.await
.unwrap();
let admin = AdminBuilder::new(path, object_store).build();
let cancellation_token = CancellationToken::new();
cancellation_token.cancel();
let result = admin
.run_compactor_with_options(
cancellation_token,
CompactorOptions {
worker: None,
..CompactorOptions::default()
},
)
.await;
assert!(matches!(result, Ok(())));
}
#[tokio::test(start_paused = true)]
async fn test_admin_run_compaction_worker_with_options_stops_on_cancellation() {
let object_store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let path =
Path::from("/tmp/test_admin_run_compaction_worker_with_options_stops_on_cancellation");
let admin = AdminBuilder::new(path, object_store).build();
let cancellation_token = CancellationToken::new();
cancellation_token.cancel();
let result = admin
.run_compaction_worker_with_options(
cancellation_token,
CompactionWorkerOptions::default(),
)
.await;
assert!(matches!(result, Ok(())));
}
#[tokio::test]
async fn test_admin_list_manifests() {
let object_store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let path = Path::from("/tmp/test_admin_list_manifests");
let manifest_store = Arc::new(ManifestStore::new(&path, object_store.clone()));
let mut stored = StoredManifest::create_new_db(
manifest_store,
ManifestCore::new(),
Arc::new(DefaultSystemClock::new()),
)
.await
.unwrap();
let mut dirty = stored.prepare_dirty().unwrap();
dirty.value.core.next_wal_sst_id = 5;
dirty.value.core.last_l0_seq = 10;
dirty.value.writer_epoch = 2;
dirty.value.compactor_epoch = 4;
stored.update(dirty).await.unwrap();
let mut dirty = stored.prepare_dirty().unwrap();
dirty.value.core.next_wal_sst_id = 8;
dirty.value.core.last_l0_seq = 20;
dirty.value.writer_epoch = 3;
dirty.value.compactor_epoch = 6;
stored.update(dirty).await.unwrap();
let admin = AdminBuilder::new(path.clone(), object_store).build();
let all = admin.list_manifests(..).await.unwrap();
assert_eq!(
all.iter().map(|manifest| manifest.id).collect::<Vec<_>>(),
vec![1, 2, 3]
);
assert_eq!(
all.iter()
.map(|manifest| manifest.manifest.core.last_l0_seq)
.collect::<Vec<_>>(),
vec![0, 10, 20]
);
assert_eq!(
all.iter()
.map(|manifest| manifest.manifest.writer_epoch)
.collect::<Vec<_>>(),
vec![0, 2, 3]
);
assert_eq!(
all.iter()
.map(|manifest| manifest.manifest.compactor_epoch)
.collect::<Vec<_>>(),
vec![0, 4, 6]
);
let bounded = admin.list_manifests(2..3).await.unwrap();
assert_eq!(
bounded
.iter()
.map(|manifest| manifest.id)
.collect::<Vec<_>>(),
vec![2]
);
let left_bounded = admin.list_manifests(2..).await.unwrap();
assert_eq!(
left_bounded
.iter()
.map(|manifest| manifest.id)
.collect::<Vec<_>>(),
vec![2, 3]
);
let right_bounded = admin.list_manifests(..3).await.unwrap();
assert_eq!(
right_bounded
.iter()
.map(|manifest| manifest.id)
.collect::<Vec<_>>(),
vec![1, 2]
);
}
#[tokio::test]
async fn test_admin_list_manifests_retries_transient_failure() {
let inner: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let flaky = Arc::new(FlakyObjectStore::new(inner, 0).with_list_failures(1, 0));
let path = Path::from("/tmp/test_admin_list_manifests_retries_transient_failure");
let admin = AdminBuilder::new(path, flaky.clone()).build();
let manifests = admin
.list_manifests(..)
.await
.expect("list should succeed after retrying the transient failure");
assert!(manifests.is_empty());
assert_eq!(flaky.list_attempts(), 2);
}
#[tokio::test]
async fn test_admin_create_detached_checkpoint_retries_transient_put() {
let inner: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let path = Path::from("/tmp/test_admin_create_detached_checkpoint_retries_transient_put");
let db = crate::Db::open(path.clone(), inner.clone()).await.unwrap();
db.put(b"key", b"value").await.unwrap();
db.close().await.unwrap();
let flaky = Arc::new(FlakyObjectStore::new(inner, 1));
let admin = AdminBuilder::new(path, flaky.clone()).build();
admin
.create_detached_checkpoint(&CheckpointOptions::default())
.await
.expect("checkpoint should succeed after retrying the transient put");
assert!(flaky.put_attempts() >= 2);
}
#[tokio::test]
async fn test_admin_terminal_object_store_error_maps_to_unavailable() {
let inner: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let path = Path::from("/tmp/test_admin_terminal_object_store_error_maps_to_unavailable");
let db = crate::Db::open(path.clone(), inner.clone()).await.unwrap();
db.put(b"key", b"value").await.unwrap();
db.close().await.unwrap();
let failing = Arc::new(FlakyObjectStore::new(inner, 0).with_put_precondition_always());
let admin = AdminBuilder::new(path, failing.clone()).build();
let err = admin
.create_detached_checkpoint(&CheckpointOptions::default())
.await
.expect_err("expected terminal precondition failure to surface");
assert_eq!(err.kind(), ErrorKind::Unavailable);
assert_eq!(failing.put_attempts(), 1);
}
#[tokio::test]
async fn test_admin_read_compactor_state_view_missing_manifest_maps_to_data() {
let object_store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let path = Path::from("/tmp/test_admin_read_compactor_state_view_missing_manifest");
let admin = AdminBuilder::new(path, object_store).build();
let err = admin
.read_compactor_state_view()
.await
.err()
.expect("expected missing manifest error");
assert_eq!(err.kind(), ErrorKind::Data);
}
#[tokio::test]
async fn test_admin_read_compactions() {
let object_store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let path = Path::from("/tmp/test_admin_read_compactions");
let compactions_store = Arc::new(CompactionsStore::new(&path, object_store.clone()));
let mut stored = StoredCompactions::create(compactions_store.clone(), 7)
.await
.unwrap();
let compaction_id = Ulid::new();
let compaction = Compaction::new(
compaction_id,
CompactionSpec::new(vec![SourceId::SortedRun(3)], 7),
);
let mut dirty = stored.prepare_dirty().unwrap();
dirty.value.insert(compaction);
dirty.value.compactor_epoch = 9;
stored.update(dirty).await.unwrap();
let admin = AdminBuilder::new(path.clone(), object_store).build();
let latest = admin
.read_compactions(None)
.await
.unwrap()
.expect("expected compactions");
let expected_latest = compactions_store.read_compactions(2).await.unwrap();
assert_eq!(latest.id, 2);
assert_eq!(latest.compactions.compactor_epoch, 9);
assert_eq!(latest.compactions, expected_latest);
let first = admin
.read_compactions(Some(1))
.await
.unwrap()
.expect("expected compactions");
let expected_first = compactions_store.read_compactions(1).await.unwrap();
assert_eq!(first.id, 1);
assert_eq!(first.compactions.compactor_epoch, 7);
assert_eq!(first.compactions, expected_first);
}
#[tokio::test]
async fn test_admin_list_compactions() {
let object_store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let path = Path::from("/tmp/test_admin_list_compactions");
let compactions_store = Arc::new(CompactionsStore::new(&path, object_store.clone()));
let mut stored = StoredCompactions::create(compactions_store.clone(), 2)
.await
.unwrap();
let mut dirty = stored.prepare_dirty().unwrap();
dirty.value.insert(Compaction::new(
Ulid::new(),
CompactionSpec::new(vec![SourceId::SortedRun(3)], 7),
));
dirty.value.compactor_epoch = 4;
stored.update(dirty).await.unwrap();
let mut dirty = stored.prepare_dirty().unwrap();
dirty.value.insert(Compaction::new(
Ulid::new(),
CompactionSpec::new(vec![SourceId::SortedRun(5)], 9),
));
dirty.value.compactor_epoch = 6;
stored.update(dirty).await.unwrap();
let admin = AdminBuilder::new(path.clone(), object_store).build();
let listed = admin.list_compactions(..).await.unwrap();
let ids: Vec<u64> = listed.iter().map(|compactions| compactions.id).collect();
assert_eq!(ids, vec![1, 2, 3]);
assert_eq!(
listed
.iter()
.map(|compactions| compactions.compactions.core.recent_compactions().count())
.collect::<Vec<_>>(),
vec![0, 1, 2]
);
assert_eq!(
listed
.iter()
.map(|compactions| compactions.compactions.compactor_epoch)
.collect::<Vec<_>>(),
vec![2, 4, 6]
);
let bounded = admin.list_compactions(2..3).await.unwrap();
assert_eq!(
bounded
.iter()
.map(|compactions| compactions.id)
.collect::<Vec<_>>(),
vec![2]
);
let left_bounded = admin.list_compactions(2..).await.unwrap();
assert_eq!(
left_bounded
.iter()
.map(|compactions| compactions.id)
.collect::<Vec<_>>(),
vec![2, 3]
);
let right_bounded = admin.list_compactions(..3).await.unwrap();
assert_eq!(
right_bounded
.iter()
.map(|compactions| compactions.id)
.collect::<Vec<_>>(),
vec![1, 2]
);
}
#[tokio::test]
async fn test_admin_read_compaction() {
let object_store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let path = Path::from("/tmp/test_admin_read_compaction");
let compactions_store = Arc::new(CompactionsStore::new(&path, object_store.clone()));
let mut stored = StoredCompactions::create(compactions_store.clone(), 0)
.await
.unwrap();
let compaction_id = Ulid::new();
let compaction = Compaction::new(
compaction_id,
CompactionSpec::new(vec![SourceId::SortedRun(3)], 7),
);
let mut dirty = stored.prepare_dirty().unwrap();
dirty.value.insert(compaction);
stored.update(dirty).await.unwrap();
let admin = AdminBuilder::new(path.clone(), object_store).build();
let compaction = admin
.read_compaction(compaction_id, None)
.await
.unwrap()
.expect("expected compaction");
assert_eq!(compaction.id(), compaction_id);
}
#[tokio::test]
async fn test_admin_submit_compaction() {
let object_store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let path = Path::from("/tmp/test_admin_submit_compaction");
let compactions_store = Arc::new(CompactionsStore::new(&path, object_store.clone()));
StoredCompactions::create(compactions_store.clone(), 0)
.await
.unwrap();
let admin = AdminBuilder::new(path.clone(), object_store).build();
let spec = CompactionSpec::new(vec![SourceId::SortedRun(3)], 3);
let compaction = admin.submit_compaction(spec).await.unwrap();
assert_eq!(compaction.spec().destination(), Some(3));
assert_eq!(compaction.spec().sources(), &[SourceId::SortedRun(3)]);
assert_eq!(compaction.status(), CompactionStatus::Submitted);
}
#[cfg(feature = "compaction_filters")]
#[test]
fn test_admin_builder_with_compaction_filter_supplier() {
use crate::compaction_filter::{
CompactionFilter, CompactionFilterDecision, CompactionFilterError,
CompactionFilterSupplier, CompactionJobContext,
};
use crate::types::RowEntry;
struct NoopFilter;
#[async_trait::async_trait]
impl CompactionFilter for NoopFilter {
async fn filter(
&mut self,
_entry: &RowEntry,
) -> Result<CompactionFilterDecision, CompactionFilterError> {
Ok(CompactionFilterDecision::Keep)
}
async fn on_compaction_end(&mut self) -> Result<(), CompactionFilterError> {
Ok(())
}
}
struct NoopFilterSupplier;
#[async_trait::async_trait]
impl CompactionFilterSupplier for NoopFilterSupplier {
async fn create_compaction_filter(
&self,
_context: &CompactionJobContext,
) -> Result<Box<dyn CompactionFilter>, CompactionFilterError> {
Ok(Box::new(NoopFilter))
}
}
let object_store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let admin = AdminBuilder::new("/tmp/test_filter_supplier", object_store)
.with_compaction_filter_supplier(Arc::new(NoopFilterSupplier))
.build();
assert!(admin.compaction_filter_supplier.is_some());
}
#[test]
fn test_admin_builder_with_merge_operator() {
let object_store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let admin = AdminBuilder::new("/tmp/test_merge_operator", object_store)
.with_merge_operator(Arc::new(StringConcatMergeOperator))
.build();
assert!(admin.merge_operator.is_some());
}
#[tokio::test]
async fn test_create_clone_builder() {
use crate::admin::CloneSourceSpec;
use crate::manifest::store::ManifestStore;
use crate::Db;
let object_store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let parent_path = Path::from("/tmp/test_parent");
let clone_path = Path::from("/tmp/test_clone");
let parent_db = Db::open(parent_path.clone(), object_store.clone())
.await
.unwrap();
parent_db.close().await.unwrap();
let admin = AdminBuilder::new(clone_path.clone(), object_store.clone()).build();
let r = admin.create_clone_builder_from_source(CloneSourceSpec::new(parent_path.clone()));
r.build().await.expect("clone should succeed");
let clone_manifest_store = ManifestStore::new(&clone_path, object_store.clone());
let manifest = clone_manifest_store.read_latest_manifest().await;
assert!(manifest.is_ok(), "cloned manifest should exist");
}
#[tokio::test]
async fn test_delete_db_removes_checkpoint_from_parent() {
use crate::admin::CloneSourceSpec;
use crate::config::CheckpointOptions;
use crate::manifest::store::{ManifestStore, StoredManifest};
use crate::Db;
let object_store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let system_clock = Arc::new(DefaultSystemClock::new());
let parent_path = Path::from("/tmp/test_cleanup_parent");
let clone_path = Path::from("/tmp/test_cleanup_clone");
let parent_db = Db::open(parent_path.clone(), object_store.clone())
.await
.unwrap();
parent_db.close().await.unwrap();
let parent_admin = AdminBuilder::new(parent_path.clone(), object_store.clone()).build();
let unrelated = parent_admin
.create_detached_checkpoint(&CheckpointOptions::default())
.await
.unwrap()
.id;
let clone_admin = AdminBuilder::new(clone_path.clone(), object_store.clone()).build();
clone_admin
.create_clone_builder_from_source(CloneSourceSpec::new(parent_path.clone()))
.build()
.await
.expect("clone should succeed");
let clone_ms = Arc::new(ManifestStore::new(&clone_path, object_store.clone()));
let clone_stored = StoredManifest::load(clone_ms, system_clock.clone())
.await
.unwrap();
let pinned = clone_stored.manifest().external_dbs[0]
.final_checkpoint_id
.expect("clone pins a final_checkpoint_id in the parent");
let read_parent_checkpoints = || {
let object_store = object_store.clone();
let system_clock = system_clock.clone();
let parent_path = parent_path.clone();
async move {
let ms = Arc::new(ManifestStore::new(&parent_path, object_store));
let stored = StoredManifest::load(ms, system_clock).await.unwrap();
stored
.manifest()
.core
.checkpoints
.iter()
.map(|c| c.id)
.collect::<Vec<_>>()
}
};
let before = read_parent_checkpoints().await;
assert!(
before.contains(&pinned),
"parent should have pinned checkpoint before cleanup"
);
assert!(
before.contains(&unrelated),
"parent should have unrelated checkpoint"
);
clone_admin
.delete_db(true)
.await
.expect("delete should succeed");
let after = read_parent_checkpoints().await;
assert!(
!after.contains(&pinned),
"pinned checkpoint should be gone from parent"
);
assert!(
after.contains(&unrelated),
"unrelated checkpoint should remain"
);
}
#[tokio::test]
async fn test_delete_db_deletes_clone_after_parent_already_gone() {
use crate::admin::CloneSourceSpec;
use crate::Db;
use futures::StreamExt;
use object_store::ObjectStoreExt;
let object_store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let parent_path = Path::from("/tmp/test_delete_orphan_parent");
let clone_path = Path::from("/tmp/test_delete_orphan_clone");
Db::open(parent_path.clone(), object_store.clone())
.await
.unwrap()
.close()
.await
.unwrap();
let clone_admin = AdminBuilder::new(clone_path.clone(), object_store.clone()).build();
clone_admin
.create_clone_builder_from_source(CloneSourceSpec::new(parent_path.clone()))
.build()
.await
.expect("clone should succeed");
let mut parent_listing = object_store.list(Some(&parent_path));
while let Some(meta) = parent_listing.next().await {
object_store.delete(&meta.unwrap().location).await.unwrap();
}
clone_admin
.delete_db(true)
.await
.expect("clone should delete even with parent already gone");
assert_eq!(
object_store.list(Some(&clone_path)).count().await,
0,
"clone objects should be gone"
);
}
#[tokio::test]
async fn test_delete_db_deletes_own_objects_only_with_confirm() {
use crate::admin::CloneSourceSpec;
use crate::Db;
use futures::StreamExt;
let object_store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let parent_path = Path::from("/tmp/test_delete_confirm_parent");
let clone_path = Path::from("/tmp/test_delete_confirm_clone");
Db::open(parent_path.clone(), object_store.clone())
.await
.unwrap()
.close()
.await
.unwrap();
let clone_admin = AdminBuilder::new(clone_path.clone(), object_store.clone()).build();
clone_admin
.create_clone_builder_from_source(CloneSourceSpec::new(parent_path.clone()))
.build()
.await
.expect("clone should succeed");
let count_under = |prefix: Path| {
let object_store = object_store.clone();
async move { object_store.list(Some(&prefix)).count().await }
};
let initial = count_under(clone_path.clone()).await;
assert!(initial > 0, "clone should have objects");
let would_delete = clone_admin
.delete_db(false)
.await
.expect("dry run should succeed");
assert_eq!(
would_delete.len(),
initial,
"dry run should report every object under the prefix"
);
assert_eq!(
count_under(clone_path.clone()).await,
initial,
"dry run must not delete anything"
);
clone_admin
.delete_db(true)
.await
.expect("delete should succeed");
assert_eq!(
count_under(clone_path.clone()).await,
0,
"clone objects should be gone after confirm"
);
clone_admin
.delete_db(true)
.await
.expect("second delete should be a no-op");
}
#[tokio::test]
async fn test_delete_db_finishes_partial_deletion_via_marker() {
use crate::admin::CloneSourceSpec;
use crate::Db;
use futures::StreamExt;
use object_store::ObjectStoreExt;
let object_store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let parent_path = Path::from("/tmp/test_delete_partial_parent");
let clone_path = Path::from("/tmp/test_delete_partial_clone");
Db::open(parent_path.clone(), object_store.clone())
.await
.unwrap()
.close()
.await
.unwrap();
let clone_admin = AdminBuilder::new(clone_path.clone(), object_store.clone()).build();
clone_admin
.create_clone_builder_from_source(CloneSourceSpec::new(parent_path.clone()))
.build()
.await
.expect("clone should succeed");
object_store
.put(
&clone_path.clone().join(".deleting"),
bytes::Bytes::new().into(),
)
.await
.unwrap();
let manifest_prefix = Path::from("/tmp/test_delete_partial_clone/manifest");
let mut listing = object_store.list(Some(&manifest_prefix));
while let Some(meta) = listing.next().await {
object_store.delete(&meta.unwrap().location).await.unwrap();
}
let count_under = |prefix: Path| {
let object_store = object_store.clone();
async move { object_store.list(Some(&prefix)).count().await }
};
assert!(
count_under(clone_path.clone()).await > 0,
"leftover clone objects should remain after partial deletion"
);
clone_admin
.delete_db(true)
.await
.expect("delete should finish partial deletion");
assert_eq!(
count_under(clone_path.clone()).await,
0,
"leftover objects, including the marker, should be gone"
);
}
#[tokio::test]
async fn test_delete_db_refuses_without_manifest_or_marker() {
use futures::StreamExt;
use object_store::ObjectStoreExt;
let object_store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let dir = Path::from("/tmp/test_delete_fat_finger");
object_store
.put(
&dir.clone().join("important.txt"),
bytes::Bytes::from_static(b"keepme").into(),
)
.await
.unwrap();
let admin = AdminBuilder::new(dir.clone(), object_store.clone()).build();
admin
.delete_db(true)
.await
.expect_err("delete must refuse a dir with no manifest and no marker");
let count = object_store.list(Some(&dir)).count().await;
assert_eq!(count, 1, "the unrelated object must be left untouched");
}
#[cfg(feature = "wal_disable")]
#[tokio::test]
async fn test_create_clone_with_multiple_sources() {
use crate::config::{PutOptions, Settings, WriteOptions};
use crate::manifest::store::ManifestStore;
use crate::{admin::CloneSourceSpec, Db};
let object_store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let grandparent_path1 = Path::from("/tmp/test_grandparent1");
let grandparent_path2 = Path::from("/tmp/test_grandparent2");
let parent_path1 = Path::from("/tmp/test_parent1");
let parent_path2 = Path::from("/tmp/test_parent2");
let clone_path = Path::from("/tmp/test_clone_multi");
let settings = Settings {
wal_enabled: false,
..Settings::default()
};
let write_opts = WriteOptions {
..Default::default()
};
let grandparent_db1 = Db::builder(grandparent_path1.clone(), object_store.clone())
.with_settings(settings.clone())
.build()
.await
.unwrap();
grandparent_db1
.put_with_options(b"a", b"1", &PutOptions::default(), &write_opts)
.await
.unwrap();
grandparent_db1.close().await.unwrap();
let grandparent_db2 = Db::builder(grandparent_path2.clone(), object_store.clone())
.with_settings(settings)
.build()
.await
.unwrap();
grandparent_db2
.put_with_options(b"z", b"2", &PutOptions::default(), &write_opts)
.await
.unwrap();
grandparent_db2.close().await.unwrap();
AdminBuilder::new(parent_path1.clone(), object_store.clone())
.build()
.create_clone_builder_from_source(CloneSourceSpec::new(grandparent_path1.clone()))
.build()
.await
.expect("parent clone 1 should succeed");
AdminBuilder::new(parent_path2.clone(), object_store.clone())
.build()
.create_clone_builder_from_source(CloneSourceSpec::new(grandparent_path2.clone()))
.build()
.await
.expect("parent clone 2 should succeed");
let admin = AdminBuilder::new(clone_path.clone(), object_store.clone()).build();
admin
.create_clone_builder_from_source(CloneSourceSpec::new(parent_path1.clone()))
.with_source(CloneSourceSpec::new(parent_path2.clone()))
.build()
.await
.expect("clone with multiple sources should succeed");
let clone_manifest_store = ManifestStore::new(&clone_path, object_store.clone());
let manifest = clone_manifest_store.read_latest_manifest().await;
assert!(manifest.is_ok(), "cloned manifest should exist");
let manifest_data = manifest.unwrap();
assert_eq!(
manifest_data.manifest.external_dbs.len(),
2,
"clone should have an external database for each parent"
);
}
#[cfg(feature = "wal_disable")]
#[tokio::test]
async fn test_delete_db_removes_checkpoints_from_all_parents() {
use crate::config::{PutOptions, Settings, WriteOptions};
use crate::manifest::store::{ManifestStore, StoredManifest};
use crate::{admin::CloneSourceSpec, Db};
use uuid::Uuid;
let object_store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let system_clock = Arc::new(DefaultSystemClock::new());
let parent_path1 = Path::from("/tmp/test_cleanup_multi_parent1");
let parent_path2 = Path::from("/tmp/test_cleanup_multi_parent2");
let clone_path = Path::from("/tmp/test_cleanup_multi_clone");
let settings = Settings {
wal_enabled: false,
..Settings::default()
};
let write_opts = WriteOptions::default();
for (path, key) in [(&parent_path1, b"a"), (&parent_path2, b"z")] {
let db = Db::builder(path.clone(), object_store.clone())
.with_settings(settings.clone())
.build()
.await
.unwrap();
db.put_with_options(key, b"1", &PutOptions::default(), &write_opts)
.await
.unwrap();
db.close().await.unwrap();
}
let clone_admin = AdminBuilder::new(clone_path.clone(), object_store.clone()).build();
clone_admin
.create_clone_builder_from_source(CloneSourceSpec::new(parent_path1.clone()))
.with_source(CloneSourceSpec::new(parent_path2.clone()))
.build()
.await
.expect("clone with multiple sources should succeed");
let clone_ms = Arc::new(ManifestStore::new(&clone_path, object_store.clone()));
let clone_stored = StoredManifest::load(clone_ms, system_clock.clone())
.await
.unwrap();
let pinned: Vec<(String, Uuid)> = clone_stored
.manifest()
.external_dbs
.iter()
.map(|e| (e.path.clone(), e.final_checkpoint_id.unwrap()))
.collect();
assert_eq!(pinned.len(), 2);
clone_admin
.delete_db(true)
.await
.expect("delete should succeed");
for (parent_path, checkpoint_id) in pinned {
let ms = Arc::new(ManifestStore::new(
&parent_path.into(),
object_store.clone(),
));
let stored = StoredManifest::load(ms, system_clock.clone())
.await
.unwrap();
assert!(
!stored
.manifest()
.core
.checkpoints
.iter()
.any(|c| c.id == checkpoint_id),
"pinned checkpoint should be removed from every parent"
);
}
}
}
#[cfg(test)]
mod gc_tolerant_list_manifests_tests {
use crate::admin::AdminBuilder;
use crate::manifest::store::{ManifestStore, StoredManifest};
use crate::manifest::ManifestCore;
use crate::test_utils::FlakyObjectStore;
use object_store::memory::InMemory;
use object_store::path::Path;
use object_store::ObjectStore;
use slatedb_common::clock::DefaultSystemClock;
use std::sync::Arc;
#[tokio::test]
async fn test_list_manifests_skips_manifest_gced_between_list_and_read() {
let inner: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let flaky = Arc::new(FlakyObjectStore::new(inner, 0));
let store: Arc<dyn ObjectStore> = flaky.clone();
let path = Path::from("/tmp/test_gc_tolerant_list_manifests");
let manifest_store = Arc::new(ManifestStore::new(&path, store.clone()));
let mut sm = StoredManifest::create_new_db(
manifest_store,
ManifestCore::new(),
Arc::new(DefaultSystemClock::new()),
)
.await
.unwrap();
sm.update(sm.prepare_dirty().unwrap()).await.unwrap();
sm.update(sm.prepare_dirty().unwrap()).await.unwrap();
let admin = AdminBuilder::new(path.clone(), store.clone()).build();
let ids: Vec<u64> = admin
.list_manifests(..)
.await
.unwrap()
.iter()
.map(|vm| vm.id())
.collect();
assert_eq!(ids, vec![1, 2, 3]);
flaky.with_get_not_found_failures(1);
let ids: Vec<u64> = admin
.list_manifests(..)
.await
.unwrap()
.iter()
.map(|vm| vm.id())
.collect();
assert_eq!(
ids,
vec![2, 3],
"a manifest GC'd mid-listing should be skipped, not fail the operation"
);
}
}