use std::{path::Path, sync::Arc};
#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
use crate::ReadVersion;
use crate::{
db::{
CompactionReservation, DatabaseStorageRef, Db, LsmCompactionOutput, MaintenanceBudget,
MaintenanceOutcome, NamedCompactionInput, NamedCompactionOutput, NamedFlushInput,
PendingCompactionOutputs, lock_poisoned,
},
error::{Error, Result},
object_store::{ObjectClient, ObjectStoreBackend},
options::{DurabilityMode, HostStorageBackend, StorageMode},
table::{self, Table},
types::{KeyRange, Sequence},
};
#[cfg(not(all(target_arch = "wasm32", target_os = "unknown")))]
use crate::db::engine::sync_storage_directory_after_renames_async;
use crate::db::engine::{
compaction_trigger_stat_deltas, is_level_layout_compaction_error,
sync_storage_directory_after_renames,
};
impl Db {
pub fn persist_sync(&self, mode: DurabilityMode) -> Result<()> {
self.ensure_open()?;
if (self.inner.options.storage_mode.is_wasi_persistent()
|| self.inner.options.storage_mode.is_object_store_persistent())
&& matches!(
mode,
DurabilityMode::SyncData | DurabilityMode::SyncAll | DurabilityMode::SyncAllStrict
)
{
return Err(Error::unsupported_durability(mode));
}
match &self.inner.options.storage_mode {
StorageMode::InMemory => Ok(()),
StorageMode::HostPersistent {
backend: HostStorageBackend::ObjectStore,
} => self.inner.substrate.persist_wal(mode),
StorageMode::Persistent { .. }
| StorageMode::HostPersistent {
backend: HostStorageBackend::Wasi { .. },
} => {
self.inner.substrate.persist_wal(mode)?;
Ok(())
}
StorageMode::HostPersistent { backend } => {
Err(Error::unsupported_backend(backend.as_str()))
}
}
}
pub(in crate::db) async fn flush_object_store_async(&self) -> Result<()> {
self.ensure_open()?;
if self.inner.options.read_only {
return Err(Error::ReadOnly);
}
let DatabaseStorageRef::ObjectStore(resources) = self.inner.storage.resources() else {
return Err(Error::unsupported_backend(
"object-store flush requires an object-store database",
));
};
let target_sequence = self.freeze_public_flush_target()?;
while self.has_immutable_memtables_at_or_below(target_sequence)? {
let outcome = self
.run_flush_once_with_budget_object_store_async(
resources.objects,
resources.prefix,
MaintenanceBudget::unbounded(),
)
.await?;
if outcome.busy {
return Err(Error::runtime_busy("object-store flush is already active"));
}
if outcome.flushes == 0 {
break;
}
}
Ok(())
}
pub(in crate::db) async fn flush_object_store_with_budget_async(
&self,
budget: MaintenanceBudget,
) -> Result<MaintenanceOutcome> {
self.ensure_open()?;
if self.inner.options.read_only {
return Err(Error::ReadOnly);
}
let DatabaseStorageRef::ObjectStore(resources) = self.inner.storage.resources() else {
return Err(Error::unsupported_backend(
"object-store flush requires an object-store database",
));
};
if !self.has_immutable_memtables()? {
return Ok(MaintenanceOutcome::default());
}
self.run_flush_once_with_budget_object_store_async(
resources.objects,
resources.prefix,
budget,
)
.await
}
pub(in crate::db) fn filesystem_publish_durability(&self) -> DurabilityMode {
if self.inner.options.storage_mode.is_wasi_persistent() {
DurabilityMode::Flush
} else {
DurabilityMode::SyncAll
}
}
pub(in crate::db) fn sync_filesystem_directory_after_renames(
&self,
db_path: &Path,
) -> Result<()> {
if self.inner.options.storage_mode.is_wasi_persistent() {
Ok(())
} else {
let DatabaseStorageRef::Filesystem(resources) = self.inner.storage.resources() else {
return Err(Error::unsupported_backend(
"directory synchronization requires filesystem storage",
));
};
sync_storage_directory_after_renames(resources.files, db_path)
}
}
#[cfg(not(all(target_arch = "wasm32", target_os = "unknown")))]
pub(in crate::db) async fn sync_filesystem_directory_after_renames_async(
&self,
db_path: &Path,
) -> Result<()> {
if self.inner.options.storage_mode.is_wasi_persistent() {
Ok(())
} else {
let DatabaseStorageRef::Filesystem(resources) = self.inner.storage.resources() else {
return Err(Error::unsupported_backend(
"directory synchronization requires filesystem storage",
));
};
sync_storage_directory_after_renames_async(resources.files, db_path).await
}
}
pub(in crate::db) async fn run_flush_once_with_budget_object_store_async(
&self,
backend: &ObjectStoreBackend,
db_path: &Path,
budget: MaintenanceBudget,
) -> Result<MaintenanceOutcome> {
let Some(_flush_guard) = self.inner.maintenance.try_start_flush() else {
return Ok(MaintenanceOutcome::busy_outcome());
};
let (flush_inputs, budget_exhausted) = self.collect_flush_inputs_with_budget(budget)?;
let flush_count = flush_inputs.len();
self.write_flush_inputs_object_store_async(backend, db_path, &flush_inputs)
.await?;
let outcome = MaintenanceOutcome {
flushes: flush_count,
budget_exhausted: budget_exhausted && flush_count != 0,
..MaintenanceOutcome::default()
};
if outcome.budget_exhausted {
self.record_maintenance_budget_exhaustion();
}
Ok(outcome)
}
pub(in crate::db) async fn write_flush_inputs_object_store_async(
&self,
backend: &ObjectStoreBackend,
db_path: &Path,
flush_inputs: &[NamedFlushInput],
) -> Result<()> {
if flush_inputs.is_empty() {
return Ok(());
}
let mut file_ids = self
.reserve_file_ids_object_store_async(flush_inputs.len())
.await?;
let mut written_tables = Vec::with_capacity(flush_inputs.len());
for input in flush_inputs {
let table_id = file_ids.next_table_id()?;
let table_path = table::table_path(db_path, table_id);
let table = table::write_table_with_backend_async(
backend,
&table_path,
table_id,
input.input.table_level,
&input.input.table_options,
&input.input.point_records,
&input.input.range_tombstones,
DurabilityMode::Flush,
)
.await?;
written_tables.push((input.bucket.clone(), Arc::new(table)));
}
let replay_floor = self.replay_floor_after_flush_serialized(flush_inputs)?;
self.publish_flushed_tables_object_store_async(&written_tables, replay_floor)
.await?;
Self::install_flushed_tables(flush_inputs, written_tables)
.map_err(|error| self.close_after_durable_publish_error("flush", &error))?;
self.inner
.substrate
.rewrite_wal_after_replay_floor_async(replay_floor)
.await?;
Ok(())
}
pub(in crate::db) async fn publish_flushed_tables_object_store_async(
&self,
tables: &[(String, Arc<Table>)],
flush_sequence: Sequence,
) -> Result<()> {
let edits = tables
.iter()
.map(|(bucket, table)| (bucket.clone(), table.properties().clone()))
.collect::<Vec<_>>();
self.publish_object_manifest_add_tables(edits, flush_sequence)
.await
}
#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
pub(in crate::db) async fn create_checkpoint_browser_async(
&self,
name: &str,
) -> Result<ReadVersion> {
let _manifest_publish = self.inner.browser_manifest_async_lock.lock().await;
let sequence = self.last_committed_sequence();
let manifest = self
.inner
.manifest
.as_ref()
.ok_or_else(|| Error::Corruption {
message: "browser persistent database is missing manifest store".to_owned(),
})?;
let prepared_publish = {
let manifest = manifest
.lock()
.map_err(|_| lock_poisoned("manifest store"))?;
manifest.prepare_create_checkpoint_publish(name.to_owned(), sequence)?
};
prepared_publish.publish_async().await.map_err(|error| {
self.close_after_manifest_durability_failure("checkpoint creation", error)
})?;
self.install_prepared_manifest_after_durable_publish(
"checkpoint creation",
manifest,
prepared_publish,
)?;
Ok(ReadVersion::from_sequence(sequence))
}
#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
pub(in crate::db) async fn delete_checkpoint_browser_async(&self, name: &str) -> Result<()> {
let _manifest_publish = self.inner.browser_manifest_async_lock.lock().await;
let manifest = self
.inner
.manifest
.as_ref()
.ok_or_else(|| Error::Corruption {
message: "browser persistent database is missing manifest store".to_owned(),
})?;
let prepared_publish = {
let manifest = manifest
.lock()
.map_err(|_| lock_poisoned("manifest store"))?;
manifest.prepare_delete_checkpoint_publish(name.to_owned())?
};
prepared_publish.publish_async().await.map_err(|error| {
self.close_after_manifest_durability_failure("checkpoint deletion", error)
})?;
self.install_prepared_manifest_after_durable_publish(
"checkpoint deletion",
manifest,
prepared_publish,
)
}
pub(in crate::db) async fn checkout_object_manifest(
&self,
) -> Result<(
crate::manifest::ObjectManifestStore<Arc<dyn ObjectClient>>,
futures::lock::MutexGuard<'_, ()>,
)> {
let manifest = self
.inner
.manifest
.as_ref()
.ok_or_else(|| Error::Corruption {
message: "object-store database is missing manifest store".to_owned(),
})?;
let serialize = self.inner.object_manifest_async_lock.lock().await;
let object = manifest
.lock()
.map_err(|_| lock_poisoned("manifest store"))?
.clone_object_manifest()?;
Ok((object, serialize))
}
pub(in crate::db) async fn publish_object_manifest_create_checkpoint(
&self,
name: String,
sequence: Sequence,
) -> Result<()> {
let (mut object, _serialize) = self.checkout_object_manifest().await?;
object.create_checkpoint(name, sequence).await?;
self.install_object_manifest_after_durable_publish("checkpoint creation", object)
}
pub(in crate::db) async fn reserve_file_ids_object_store_async(
&self,
count: usize,
) -> Result<crate::manifest::FileIdReservation> {
let (mut object, _serialize) = self.checkout_object_manifest().await?;
let reservation = object.reserve_file_ids(count).await?;
self.install_object_manifest_after_durable_publish("file-id reservation", object)?;
Ok(reservation)
}
pub(in crate::db) async fn publish_object_manifest_delete_checkpoint(
&self,
name: String,
) -> Result<()> {
let (mut object, _serialize) = self.checkout_object_manifest().await?;
object.delete_checkpoint(name).await?;
self.install_object_manifest_after_durable_publish("checkpoint deletion", object)
}
pub(in crate::db) async fn publish_object_manifest_add_tables(
&self,
edits: Vec<(String, table::TableProperties)>,
flush_sequence: Sequence,
) -> Result<()> {
let (mut object, _serialize) = self.checkout_object_manifest().await?;
object.add_tables(edits, flush_sequence).await?;
self.install_object_manifest_after_durable_publish("flush", object)
}
pub(in crate::db) async fn publish_object_manifest_replace_tables(
&self,
edits: Vec<(String, Vec<table::TableId>, Vec<table::TableProperties>)>,
obsolete_blob_ids: Vec<u64>,
pending_deletion_sequence: Sequence,
) -> Result<()> {
let (mut object, _serialize) = self.checkout_object_manifest().await?;
object
.replace_tables_batch_and_mark_blob_deletions(
edits,
obsolete_blob_ids,
pending_deletion_sequence,
)
.await?;
self.install_object_manifest_after_durable_publish("compaction", object)
}
pub(in crate::db) fn install_object_manifest(
&self,
object: crate::manifest::ObjectManifestStore<Arc<dyn ObjectClient>>,
) -> Result<()> {
self.inner
.manifest
.as_ref()
.ok_or_else(|| Error::Corruption {
message: "object-store database is missing manifest store".to_owned(),
})?
.lock()
.map_err(|_| lock_poisoned("manifest store"))?
.install_object_manifest(object)
}
pub(in crate::db) fn install_object_manifest_after_durable_publish(
&self,
operation: &'static str,
object: crate::manifest::ObjectManifestStore<Arc<dyn ObjectClient>>,
) -> Result<()> {
self.install_object_manifest(object)
.map_err(|error| self.close_after_durable_publish_error(operation, &error))
}
#[allow(clippy::unused_async)] pub(in crate::db) async fn cleanup_object_store_orphans_async(&self) -> Result<usize> {
self.ensure_open()?;
if self.inner.options.read_only {
return Err(Error::ReadOnly);
}
Ok(0)
}
#[allow(clippy::too_many_lines)] pub(in crate::db) async fn run_compaction_once_object_store_async(
&self,
backend: &ObjectStoreBackend,
db_path: &Path,
range: &KeyRange,
local_l0_compaction: bool,
budget: MaintenanceBudget,
) -> Result<MaintenanceOutcome> {
let latest = self.last_committed_sequence();
let retained_floor = self.retained_floor_without_active_snapshots(latest);
let (oldest_active_snapshot, _snapshot_guard) =
self.inner.snapshots.begin_compaction(retained_floor);
let compaction_inputs =
self.collect_compaction_inputs(range, oldest_active_snapshot, local_l0_compaction)?;
if compaction_inputs.is_empty() {
return Ok(MaintenanceOutcome::default());
}
let reservations = compaction_inputs
.iter()
.map(|input| CompactionReservation {
bucket: input.bucket.clone(),
range: input.input.compaction_range.clone(),
})
.collect::<Vec<_>>();
let Some(compaction_guard) = self.inner.maintenance.reserve_compactions(reservations)
else {
return Ok(MaintenanceOutcome::busy_outcome());
};
let mut compaction_inputs = compaction_inputs
.into_iter()
.filter(|input| compaction_guard.contains(&input.bucket, &input.input.compaction_range))
.collect::<Vec<_>>();
if compaction_inputs.is_empty() {
return Ok(MaintenanceOutcome::busy_outcome());
}
let limit = budget.compaction_input_limit();
let budget_exhausted = compaction_inputs.len() > limit;
compaction_inputs.truncate(limit);
if compaction_inputs.is_empty() {
return Ok(MaintenanceOutcome::default());
}
if !Self::compaction_inputs_are_current(&compaction_inputs)? {
return Ok(MaintenanceOutcome::busy_outcome());
}
let PendingCompactionOutputs {
outputs: written_tables,
written_table_ids: _,
} = self
.build_compaction_outputs_object_store_async(
backend,
db_path,
oldest_active_snapshot,
&compaction_inputs,
)
.await?;
let output_tables = written_tables
.iter()
.flat_map(|output| output.output.tables.iter().cloned())
.collect::<Vec<_>>();
let input_tables = compaction_inputs
.iter()
.flat_map(|input| input.input.input_tables.iter().cloned())
.collect::<Vec<_>>();
let trigger_stats = compaction_trigger_stat_deltas(&compaction_inputs, &written_tables);
let obsolete_blob_ids =
self.obsolete_blob_ids_for_compaction(&compaction_inputs, &written_tables)?;
match self.validate_compacted_tables(&written_tables) {
Ok(()) => {}
Err(error) if is_level_layout_compaction_error(&error) => {
return Ok(MaintenanceOutcome::default());
}
Err(error) => return Err(error),
}
let edits = written_tables
.iter()
.map(|output| {
(
output.bucket.clone(),
output.output.input_table_ids.clone(),
output
.output
.tables
.iter()
.map(|table| table.properties().clone())
.collect::<Vec<_>>(),
)
})
.collect::<Vec<_>>();
let pending_deletion_sequence = self.last_committed_sequence();
self.publish_object_manifest_replace_tables(
edits,
obsolete_blob_ids,
pending_deletion_sequence,
)
.await?;
self.install_compacted_tables(written_tables)
.map_err(|error| self.close_after_durable_publish_error("compaction", &error))?;
self.record_compaction_stats_from_tables(
compaction_inputs.len(),
&input_tables,
&output_tables,
&trigger_stats,
);
let outcome = MaintenanceOutcome {
compactions: compaction_inputs.len(),
budget_exhausted,
..MaintenanceOutcome::default()
};
if outcome.budget_exhausted {
self.record_maintenance_budget_exhaustion();
}
Ok(outcome)
}
pub(in crate::db) async fn build_compaction_outputs_object_store_async(
&self,
backend: &ObjectStoreBackend,
db_path: &Path,
oldest_active_snapshot: Sequence,
compaction_inputs: &[NamedCompactionInput],
) -> Result<PendingCompactionOutputs> {
let rewrites =
self.prepare_compaction_rewrites(oldest_active_snapshot, compaction_inputs)?;
let output_table_count = rewrites
.iter()
.flatten()
.try_fold(0usize, |count, rewrite| {
count
.checked_add(rewrite.payloads.len())
.ok_or_else(|| Error::Corruption {
message: "compaction output table count overflow".to_owned(),
})
})?;
let mut file_ids = self
.reserve_file_ids_object_store_async(output_table_count)
.await?;
let mut outputs = Vec::with_capacity(compaction_inputs.len());
let mut written_table_ids = Vec::new();
for (input, rewrite) in compaction_inputs.iter().zip(rewrites) {
let Some(rewrite) = rewrite else {
outputs.push(NamedCompactionOutput {
bucket: input.bucket.clone(),
trigger: Some(input.input.trigger),
output: LsmCompactionOutput {
input_table_ids: input.input.input_table_ids.clone(),
tables: vec![input.input.moved_table()?],
},
});
continue;
};
let mut output_tables = Vec::with_capacity(rewrite.payloads.len());
for payload in rewrite.payloads {
let table_id = file_ids.next_table_id()?;
let table_path = table::table_path(db_path, table_id);
written_table_ids.push(table_id);
let table = table::write_table_with_backend_async(
backend,
&table_path,
table_id,
input.input.table_level,
&rewrite.table_options,
&payload.point_records,
&payload.range_tombstones,
DurabilityMode::Flush,
)
.await?;
output_tables.push(Arc::new(table));
}
outputs.push(NamedCompactionOutput {
bucket: input.bucket.clone(),
trigger: Some(input.input.trigger),
output: LsmCompactionOutput {
input_table_ids: input.input.input_table_ids.clone(),
tables: output_tables,
},
});
}
Ok(PendingCompactionOutputs {
outputs,
written_table_ids,
})
}
#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
pub(in crate::db) async fn run_owned_browser_task<T>(
label: &'static str,
task: impl std::future::Future<Output = Result<T>> + 'static,
) -> Result<T>
where
T: 'static,
{
let (sender, receiver) = futures::channel::oneshot::channel();
wasm_bindgen_futures::spawn_local(async move {
let _ = sender.send(task.await);
});
receiver.await.map_err(|_| Error::runtime_busy(label))?
}
#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
pub(in crate::db) async fn persist_browser_async(&self, mode: DurabilityMode) -> Result<()> {
self.ensure_open()?;
if matches!(
mode,
DurabilityMode::SyncData | DurabilityMode::SyncAll | DurabilityMode::SyncAllStrict
) {
return Err(Error::unsupported_durability(mode));
}
let DatabaseStorageRef::Browser(resources) = self.inner.storage.resources() else {
return Err(Error::unsupported_backend(
"browser WAL persistence requires browser storage",
));
};
let Some(wal) = resources.wal else {
return Ok(());
};
wal.persist(resources.files, resources.root, mode).await
}
#[cfg(not(all(target_arch = "wasm32", target_os = "unknown")))]
pub(in crate::db) async fn persist_native_async(&self, mode: DurabilityMode) -> Result<()> {
self.ensure_open()?;
if self.inner.options.storage_mode.is_wasi_persistent()
&& matches!(
mode,
DurabilityMode::SyncData | DurabilityMode::SyncAll | DurabilityMode::SyncAllStrict
)
{
return Err(Error::unsupported_durability(mode));
}
match &self.inner.options.storage_mode {
StorageMode::InMemory => Ok(()),
StorageMode::Persistent { .. }
| StorageMode::HostPersistent {
backend: HostStorageBackend::Wasi { .. } | HostStorageBackend::ObjectStore,
} => self.inner.substrate.persist_wal_async(mode).await,
StorageMode::HostPersistent { backend } => {
Err(Error::unsupported_backend(backend.as_str()))
}
}
}
}