use super::{
Arc, BlobLevelMergePolicy, CompactionReservation, Db, DurabilityMode, Error,
HostStorageBackend, KeyRange, LsmCompactionOutput, MaintenanceBudget, MaintenanceOutcome,
NamedCompactionInput, NamedCompactionOutput, NamedFlushInput, ObjectClient, ObjectStoreBackend,
Path, PendingCompactionOutputs, Result, Sequence, StorageMode, StorageObjectDeleteBackend,
StorageObjectId, StorageObjectKind, Table, blob, compaction_trigger_stat_deltas,
is_level_layout_compaction_error, lock_poisoned, referenced_blob_file_ids_from_manifest,
referenced_table_file_ids, should_rewrite_blob_indexes_for_compaction,
sync_storage_directory_after_renames, table,
};
#[cfg(not(all(target_arch = "wasm32", target_os = "unknown")))]
use super::{Ordering, shutdown_background_workers, sync_storage_directory_after_renames_async};
#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
use crate::{ReadVersion, storage::BrowserStorageBackend};
impl Db {
pub fn persist_sync(&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::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()))
}
}
}
#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
pub(in crate::db) fn browser_storage(&self) -> Result<BrowserStorageBackend> {
self.inner
.browser_storage
.clone()
.ok_or_else(|| Error::Corruption {
message: "browser persistent database is missing storage backend".to_owned(),
})
}
#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
pub(in crate::db) fn browser_db_path(&self) -> Result<&Path> {
self.inner
.options
.storage_mode
.browser_path()
.ok_or_else(|| Error::Corruption {
message: "browser persistent database is missing namespace path".to_owned(),
})
}
pub(in crate::db) fn object_storage(&self) -> Result<ObjectStoreBackend> {
self.inner
.object_storage
.clone()
.ok_or_else(|| Error::Corruption {
message: "object-store database is missing storage backend".to_owned(),
})
}
pub(in crate::db) fn object_wal_storage(&self) -> Result<ObjectStoreBackend> {
self.inner
.object_wal_storage
.clone()
.ok_or_else(|| Error::Corruption {
message: "object-store database is missing WAL backend".to_owned(),
})
}
pub(in crate::db) fn object_store_db_path(&self) -> &Path {
&self.inner.object_storage_prefix
}
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 db_path = self.object_store_db_path();
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(
db_path,
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 db_path = self.object_store_db_path();
if !self.has_immutable_memtables()? {
return Ok(MaintenanceOutcome::default());
}
self.run_flush_once_with_budget_object_store_async(db_path, 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 {
sync_storage_directory_after_renames(&self.inner.native_storage, 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 {
sync_storage_directory_after_renames_async(&self.inner.native_storage, db_path).await
}
}
pub(in crate::db) async fn run_flush_once_with_budget_object_store_async(
&self,
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(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,
db_path: &Path,
flush_inputs: &[NamedFlushInput],
) -> Result<()> {
if flush_inputs.is_empty() {
return Ok(());
}
let flush_sequence = flush_inputs
.iter()
.map(|input| input.input.freeze_sequence)
.max()
.expect("non-empty flush input list has a max sequence");
let backend = self.object_storage()?;
let mut written_tables = Vec::with_capacity(flush_inputs.len());
for input in flush_inputs {
let table_path = table::table_path(db_path, input.input.table_id);
let table = table::write_table_with_backend_async(
&backend,
&table_path,
input.input.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)));
}
self.publish_flushed_tables_object_store_async(&written_tables, flush_sequence)
.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(flush_sequence)
.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?;
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?;
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_bucket(
&self,
name: String,
options: crate::BucketOptions,
) -> Result<()> {
let (mut object, _serialize) = self.checkout_object_manifest().await?;
object.create_bucket(name, options).await?;
self.install_object_manifest_after_durable_publish("bucket creation", object)
}
pub(in crate::db) async fn publish_object_manifest_drop_bucket(
&self,
name: String,
) -> Result<()> {
let (mut object, _serialize) = self.checkout_object_manifest().await?;
object.drop_bucket(name).await?;
self.install_object_manifest_after_durable_publish("bucket drop", object)
}
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 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))
}
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);
}
let Some(_flush_guard) = self.inner.maintenance.try_start_flush() else {
return Err(Error::runtime_busy(
"object-store orphan GC cannot run during a flush",
));
};
let backend = self.object_storage()?;
let db_path = self.object_store_db_path();
let (referenced_tables, referenced_blobs) = {
let manifest = self
.inner
.manifest
.as_ref()
.ok_or_else(|| Error::Corruption {
message: "object-store database is missing manifest store".to_owned(),
})?;
let manifest = manifest
.lock()
.map_err(|_| lock_poisoned("manifest store"))?;
(
referenced_table_file_ids(manifest.state()),
referenced_blob_file_ids_from_manifest(manifest.state()),
)
};
let mut deleted = 0_usize;
for table_id in table::list_table_file_ids_with_backend_async(&backend, db_path).await? {
if !referenced_tables.contains(&table_id) {
backend
.delete_object(StorageObjectId::native_file(
StorageObjectKind::Table,
table::table_path(db_path, table_id),
))
.await?;
deleted += 1;
}
}
for file_id in blob::list_blob_file_ids_with_backend_async(&backend, db_path).await? {
if !referenced_blobs.contains(&file_id) {
backend
.delete_object(StorageObjectId::native_file(
StorageObjectKind::Blob,
blob::blob_path(db_path, file_id),
))
.await?;
deleted += 1;
}
}
Ok(deleted)
}
#[allow(clippy::too_many_lines)] pub(in crate::db) async fn run_compaction_once_object_store_async(
&self,
db_path: &Path,
range: &KeyRange,
local_l0_compaction: bool,
budget: MaintenanceBudget,
) -> Result<MaintenanceOutcome> {
let oldest_active_snapshot = self.oldest_retained_sequence();
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(
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,
db_path: &Path,
oldest_active_snapshot: Sequence,
compaction_inputs: &[NamedCompactionInput],
) -> Result<PendingCompactionOutputs> {
let backend = self.object_storage()?;
let mut outputs = Vec::with_capacity(compaction_inputs.len());
let mut written_table_ids = Vec::new();
let mut next_table_id = self.next_table_id()?;
for input in compaction_inputs {
let force_rewrite_trivial =
input.tree.options.blob_level_merge_policy == BlobLevelMergePolicy::Always;
if input.input.trivial_move && !force_rewrite_trivial {
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 payloads = input.tree.build_compaction_table_payloads(
&input.input,
&input.input.compaction_range,
oldest_active_snapshot,
self.inner.options.target_table_bytes,
)?;
let mut table_options = input.input.table_options.clone();
table_options.rewrite_blob_indexes = should_rewrite_blob_indexes_for_compaction(
&input.input,
&payloads,
input.tree.options.blob_level_merge_policy,
);
let mut output_tables = Vec::with_capacity(payloads.len());
for payload in payloads {
let table_id = next_table_id;
next_table_id = next_table_id.next().ok_or_else(|| Error::Corruption {
message: "table id counter overflow".to_owned(),
})?;
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,
&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 Some(wal) = &self.inner.browser_wal else {
return Ok(());
};
let storage = self.browser_storage()?;
wal.persist(&storage, self.browser_db_path()?, 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()))
}
}
}
pub fn flush_sync(&self) -> Result<()> {
self.ensure_open()?;
if self.inner.options.read_only {
return Err(Error::ReadOnly);
}
self.take_background_maintenance_error()?;
if self.inner.options.storage_mode.is_browser_persistent() {
return Err(Error::unsupported_backend(
"browser persistent flush requires async maintenance",
));
}
if self.inner.options.storage_mode.is_object_store_persistent() {
return Err(Error::unsupported_backend(
"object-store flush requires the async API",
));
}
let Some(path) = self.persistent_path() else {
return Ok(());
};
let db_path = path.to_path_buf();
let target_sequence = self.freeze_public_flush_target()?;
let mut should_compact = false;
while self.has_immutable_memtables_at_or_below(target_sequence)? {
self.take_background_maintenance_error()?;
if self.run_flush_once(&db_path, false)? {
should_compact |= self.l0_pressure_exceeded()?;
continue;
}
self.request_background_flush();
self.record_cooperative_maintenance_yield();
self.inner.maintenance.wait_until_flush_idle();
}
if should_compact
|| self.l0_pressure_exceeded()?
|| self.foreground_l0_overlap_pressure_exceeded()?
{
self.run_compaction_barrier(&db_path, &KeyRange::all(), true)?;
}
self.cleanup_pending_obsolete_table_files(&db_path)?;
self.cleanup_pending_obsolete_blob_files(&db_path)?;
self.take_background_maintenance_error()?;
Ok(())
}
#[cfg(not(all(target_arch = "wasm32", target_os = "unknown")))]
pub(in crate::db) async fn flush_native_async(&self) -> Result<()> {
self.ensure_open()?;
if self.inner.options.read_only {
return Err(Error::ReadOnly);
}
self.take_background_maintenance_error()?;
if self.inner.options.storage_mode.is_object_store_persistent() {
return Err(Error::unsupported_backend(
"object-store flush requires the async API",
));
}
let Some(path) = self.persistent_path() else {
return Ok(());
};
let db_path = path.to_path_buf();
let target_sequence = self.freeze_public_flush_target()?;
let mut should_compact = false;
while self.has_immutable_memtables_at_or_below(target_sequence)? {
self.take_background_maintenance_error()?;
let (flush_should_compact, outcome) = self
.run_flush_once_with_budget_native_async(
&db_path,
false,
MaintenanceBudget::unbounded(),
)
.await?;
if outcome.busy {
self.request_background_flush();
self.record_cooperative_maintenance_yield();
self.inner.maintenance.wait_until_flush_idle();
continue;
}
should_compact |= flush_should_compact;
}
if should_compact
|| self.l0_pressure_exceeded()?
|| self.foreground_l0_overlap_pressure_exceeded()?
{
self.run_compaction_barrier_native_async(&db_path, &KeyRange::all(), true)
.await?;
}
self.cleanup_pending_obsolete_table_files_native_async(&db_path)
.await?;
self.cleanup_pending_obsolete_blob_files_native_async(&db_path)
.await?;
self.take_background_maintenance_error()?;
Ok(())
}
#[cfg(not(all(target_arch = "wasm32", target_os = "unknown")))]
pub(in crate::db) async fn close_native_async(&self) -> Result<()> {
self.inner.closed.store(true, Ordering::Release);
shutdown_background_workers(
&self.inner.maintenance,
&self.inner.runtime_shutdown,
&self.inner.background_workers,
);
self.inner.publish_barrier.close()?;
if let Some(db_path) = self.persistent_path().map(Path::to_path_buf) {
self.cleanup_pending_obsolete_table_files_native_async(&db_path)
.await?;
self.cleanup_pending_obsolete_blob_files_native_async(&db_path)
.await?;
}
super::super::release_browser_writer_lease(&self.inner);
self.inner.substrate.release_writer_lease();
Ok(())
}
#[allow(clippy::needless_pass_by_value)]
pub fn compact_range_sync(&self, range: KeyRange) -> Result<()> {
self.take_background_maintenance_error()?;
self.compact_range_internal(range)
}
#[allow(clippy::needless_pass_by_value)]
pub fn compact_range_with_budget_sync(
&self,
range: KeyRange,
budget: MaintenanceBudget,
) -> Result<MaintenanceOutcome> {
self.take_background_maintenance_error()?;
self.compact_range_with_budget_internal(range, budget)
}
#[allow(clippy::needless_pass_by_value)]
pub(in crate::db) fn compact_range_internal(&self, range: KeyRange) -> Result<()> {
self.ensure_open()?;
if self.inner.options.read_only {
return Err(Error::ReadOnly);
}
if self.inner.options.storage_mode.is_browser_persistent() {
return Err(Error::unsupported_backend(
"browser persistent compaction requires async maintenance",
));
}
if self.inner.options.storage_mode.is_object_store_persistent() {
return Err(Error::unsupported_backend(
"object-store compaction requires the async API",
));
}
let Some(path) = self.persistent_path() else {
return Ok(());
};
let db_path = path.to_path_buf();
self.run_compaction_barrier(&db_path, &range, false)?;
self.cleanup_pending_obsolete_table_files(&db_path)?;
Ok(())
}
#[allow(clippy::needless_pass_by_value)]
pub(in crate::db) fn compact_range_with_budget_internal(
&self,
range: KeyRange,
budget: MaintenanceBudget,
) -> Result<MaintenanceOutcome> {
self.ensure_open()?;
if self.inner.options.read_only {
return Err(Error::ReadOnly);
}
if self.inner.options.storage_mode.is_browser_persistent() {
return Err(Error::unsupported_backend(
"browser persistent compaction requires async maintenance",
));
}
if self.inner.options.storage_mode.is_object_store_persistent() {
return Err(Error::unsupported_backend(
"object-store compaction requires the async API",
));
}
let Some(path) = self.persistent_path() else {
return Ok(MaintenanceOutcome::default());
};
let db_path = path.to_path_buf();
self.run_compaction_once_with_budget(&db_path, &range, false, budget)
}
pub fn run_maintenance_with_budget_sync(
&self,
budget: MaintenanceBudget,
) -> Result<MaintenanceOutcome> {
self.take_background_maintenance_error()?;
self.ensure_open()?;
if self.inner.options.read_only {
return Err(Error::ReadOnly);
}
if self.inner.options.storage_mode.is_browser_persistent() {
return Err(Error::unsupported_backend(
"browser persistent maintenance requires async maintenance",
));
}
if self.inner.options.storage_mode.is_object_store_persistent() {
return Err(Error::unsupported_backend(
"object-store maintenance requires the async API",
));
}
let Some(path) = self.persistent_path() else {
return Ok(MaintenanceOutcome::default());
};
let db_path = path.to_path_buf();
let mut outcome = MaintenanceOutcome::default();
let mut should_compact = self.l0_pressure_exceeded()?;
if self.has_immutable_memtables()? {
let (flush_should_compact, flush_outcome) =
self.run_flush_once_with_budget(&db_path, false, budget)?;
should_compact |= flush_should_compact;
outcome.add_assign(flush_outcome);
}
if should_compact {
let compaction_outcome =
self.run_compaction_once_with_budget(&db_path, &KeyRange::all(), true, budget)?;
outcome.add_assign(compaction_outcome);
}
if outcome.made_progress() {
self.cleanup_pending_obsolete_table_files(&db_path)?;
self.cleanup_pending_obsolete_blob_files(&db_path)?;
}
self.take_background_maintenance_error()?;
Ok(outcome)
}
#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
pub(in crate::db) async fn flush_browser_async(&self) -> Result<()> {
self.ensure_open()?;
if self.inner.options.read_only {
return Err(Error::ReadOnly);
}
self.take_background_maintenance_error()?;
let db_path = self.browser_db_path()?;
let target_sequence = self.freeze_public_flush_target()?;
let mut should_compact = false;
while self.has_immutable_memtables_at_or_below(target_sequence)? {
self.take_background_maintenance_error()?;
let (flush_should_compact, outcome) = self
.run_flush_once_with_budget_browser_async(
db_path,
false,
MaintenanceBudget::unbounded(),
)
.await?;
if outcome.busy {
return Err(Error::runtime_busy(
"browser persistent flush is already active",
));
}
should_compact |= flush_should_compact;
if outcome.flushes == 0 {
break;
}
}
if should_compact
|| self.l0_pressure_exceeded()?
|| self.foreground_l0_overlap_pressure_exceeded()?
{
let outcome = self
.run_compaction_once_with_budget_browser_async(
db_path,
&KeyRange::all(),
true,
MaintenanceBudget::unbounded(),
)
.await?;
if outcome.busy {
return Err(Error::runtime_busy(
"browser persistent compaction is already active",
));
}
}
self.cleanup_pending_obsolete_table_files_browser_async(db_path)
.await?;
self.cleanup_pending_obsolete_blob_files_browser_async(db_path)
.await?;
self.take_background_maintenance_error()?;
Ok(())
}
#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
pub(in crate::db) async fn compact_range_browser_async(&self, range: KeyRange) -> Result<()> {
self.take_background_maintenance_error()?;
self.ensure_open()?;
if self.inner.options.read_only {
return Err(Error::ReadOnly);
}
let outcome = self
.run_compaction_once_with_budget_browser_async(
self.browser_db_path()?,
&range,
false,
MaintenanceBudget::unbounded(),
)
.await?;
if outcome.busy {
return Err(Error::runtime_busy(
"browser persistent compaction is already active",
));
}
Ok(())
}
#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
pub(in crate::db) async fn compact_range_with_budget_browser_async(
&self,
range: KeyRange,
budget: MaintenanceBudget,
) -> Result<MaintenanceOutcome> {
self.take_background_maintenance_error()?;
self.ensure_open()?;
if self.inner.options.read_only {
return Err(Error::ReadOnly);
}
self.run_compaction_once_with_budget_browser_async(
self.browser_db_path()?,
&range,
false,
budget,
)
.await
}
#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
pub(in crate::db) async fn run_maintenance_with_budget_browser_async(
&self,
budget: MaintenanceBudget,
) -> Result<MaintenanceOutcome> {
self.take_background_maintenance_error()?;
self.ensure_open()?;
if self.inner.options.read_only {
return Err(Error::ReadOnly);
}
let db_path = self.browser_db_path()?;
let mut outcome = MaintenanceOutcome::default();
let mut should_compact = self.l0_pressure_exceeded()?;
if self.has_immutable_memtables()? {
let (flush_should_compact, flush_outcome) = self
.run_flush_once_with_budget_browser_async(db_path, false, budget)
.await?;
should_compact |= flush_should_compact;
outcome.add_assign(flush_outcome);
}
if should_compact {
let compaction_outcome = self
.run_compaction_once_with_budget_browser_async(
db_path,
&KeyRange::all(),
true,
budget,
)
.await?;
outcome.add_assign(compaction_outcome);
}
if outcome.made_progress() {
self.cleanup_pending_obsolete_table_files_browser_async(db_path)
.await?;
self.cleanup_pending_obsolete_blob_files_browser_async(db_path)
.await?;
}
self.take_background_maintenance_error()?;
Ok(outcome)
}
}