use super::{
BTreeMap, Db, DbOptions, DbStats, Error, FilterStats, LevelFilterStats, LevelStats, Ordering,
Path, ReadVersion, Result, Sequence, Snapshot, Transaction, TransactionOptions,
add_obsolete_blob_stats, lock_poisoned, shutdown_background_workers, table, table_file_bytes,
validate_checkpoint_name,
};
impl Db {
#[must_use]
pub fn snapshot(&self) -> Snapshot {
self.inner
.snapshots
.pinned_snapshot(self.last_committed_sequence())
}
pub fn snapshot_at(&self, version: ReadVersion) -> Result<Snapshot> {
self.ensure_open()?;
let latest = self.last_committed_sequence();
let retained_floor = self.retained_floor_without_active_snapshots(latest);
self.inner
.snapshots
.pinned_retained_snapshot(version.to_sequence(), latest, retained_floor)
}
pub fn create_checkpoint_sync(&self, name: &str) -> Result<ReadVersion> {
self.ensure_open()?;
validate_checkpoint_name(name)?;
if self.inner.options.read_only {
return Err(Error::ReadOnly);
}
let sequence = self.last_committed_sequence();
self.record_checkpoint(name, sequence)?;
Ok(ReadVersion::from_sequence(sequence))
}
pub fn create_checkpoint_at_sync(&self, name: &str, version: ReadVersion) -> Result<()> {
self.ensure_open()?;
validate_checkpoint_name(name)?;
if self.inner.options.read_only {
return Err(Error::ReadOnly);
}
let _pin = self.snapshot_at(version)?;
self.record_checkpoint(name, version.to_sequence())
}
pub async fn create_checkpoint_at(&self, name: &str, version: ReadVersion) -> Result<()> {
self.ensure_open()?;
validate_checkpoint_name(name)?;
if self.inner.options.read_only {
return Err(Error::ReadOnly);
}
let _pin = self.snapshot_at(version)?;
if self.inner.options.storage_mode.is_object_store_persistent() {
self.publish_object_manifest_create_checkpoint(name.to_owned(), version.to_sequence())
.await?;
return Ok(());
}
#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
if self.inner.options.storage_mode.is_browser_persistent() {
return self.create_checkpoint_at_sync(name, version);
}
self.create_checkpoint_at_sync(name, version)
}
pub(in crate::db) fn record_checkpoint(&self, name: &str, sequence: Sequence) -> Result<()> {
if let Some(manifest) = &self.inner.manifest {
manifest
.lock()
.map_err(|_| lock_poisoned("manifest store"))?
.create_checkpoint(name.to_owned(), sequence)?;
} else {
let mut checkpoints = self
.inner
.checkpoints
.lock()
.map_err(|_| lock_poisoned("checkpoint registry"))?;
if checkpoints.insert(name.to_owned(), sequence).is_some() {
return Err(Error::CheckpointAlreadyExists {
name: name.to_owned(),
});
}
}
Ok(())
}
pub fn delete_checkpoint_sync(&self, name: &str) -> Result<()> {
self.ensure_open()?;
validate_checkpoint_name(name)?;
if self.inner.options.read_only {
return Err(Error::ReadOnly);
}
if let Some(manifest) = &self.inner.manifest {
manifest
.lock()
.map_err(|_| lock_poisoned("manifest store"))?
.delete_checkpoint(name.to_owned())?;
} else {
let mut checkpoints = self
.inner
.checkpoints
.lock()
.map_err(|_| lock_poisoned("checkpoint registry"))?;
if checkpoints.remove(name).is_none() {
return Err(Error::CheckpointNotFound {
name: name.to_owned(),
});
}
}
Ok(())
}
pub fn checkpoint_read_version_sync(&self, name: &str) -> Result<ReadVersion> {
self.ensure_open()?;
validate_checkpoint_name(name)?;
let sequence = if let Some(manifest) = &self.inner.manifest {
manifest
.lock()
.map_err(|_| lock_poisoned("manifest store"))?
.checkpoint_sequence(name)
} else {
self.inner
.checkpoints
.lock()
.map_err(|_| lock_poisoned("checkpoint registry"))?
.get(name)
.copied()
}
.ok_or_else(|| Error::CheckpointNotFound {
name: name.to_owned(),
})?;
Ok(ReadVersion::from_sequence(sequence))
}
#[must_use]
pub fn transaction(&self, options: TransactionOptions) -> Transaction {
Transaction::new(self.clone(), self.last_committed_sequence(), options)
}
#[must_use]
pub fn stats(&self) -> DbStats {
let mut stats = self.base_stats();
self.add_wal_stats(&mut stats);
self.add_storage_runtime_stats(&mut stats);
let (blob_read_count, blob_read_bytes) = self.inner.blob_reads.snapshot();
stats.blob_read_count = blob_read_count;
stats.blob_read_bytes = blob_read_bytes;
self.add_scan_and_snapshot_health_stats(&mut stats);
self.add_compaction_level_stats(&mut stats);
self.add_compaction_trigger_stats(&mut stats);
self.add_compaction_skip_stats(&mut stats);
let cache_stats = self.inner.block_cache.stats();
stats.block_cache_hits = cache_stats.hits;
stats.block_cache_misses = cache_stats.misses;
let persistent_path = self.persistent_path();
let mut level_stats = BTreeMap::<u32, LevelStats>::new();
let mut level_filter_stats = BTreeMap::<u32, LevelFilterStats>::new();
let mut live_blob_bytes_by_file = BTreeMap::<u64, u64>::new();
let Ok(buckets) = self.inner.buckets.read() else {
return stats;
};
stats.live_buckets = buckets.len();
for state in buckets.values() {
if let Ok(memtable_bytes) = state.memtable_bytes() {
stats.memtable_bytes = stats.memtable_bytes.saturating_add(memtable_bytes);
}
stats.immutable_memtables = stats
.immutable_memtables
.saturating_add(state.immutable_memtable_count());
let Ok(version) = state.current_version() else {
continue;
};
stats
.read_path
.saturating_add_assign(version.read_path_stats());
for (level_state, tables) in version.level_table_handles() {
let level = level_state.get();
let level_entry = level_stats.entry(level).or_insert(LevelStats {
level,
tables: 0,
bytes: 0,
});
for table in tables {
let properties = table.properties();
let table_bytes = persistent_path.map_or(0, |db_path| {
table_file_bytes(&self.inner.native_storage, db_path, properties.id)
});
let table_filters = table.filter_stats();
stats.filters.saturating_add_assign(table_filters);
let level_filter_entry =
level_filter_stats.entry(level).or_insert(LevelFilterStats {
level,
tables: 0,
filters: FilterStats::default(),
filter_resident_bytes: 0,
});
level_filter_entry.tables += 1;
level_filter_entry
.filters
.saturating_add_assign(table_filters);
level_filter_entry.filter_resident_bytes = level_filter_entry
.filter_resident_bytes
.saturating_add(table.resident_filter_bytes());
stats
.read_path
.saturating_add_assign(table.read_path_stats());
stats.total_tables += 1;
stats.table_bytes = stats.table_bytes.saturating_add(table_bytes);
if properties.level == table::TableLevel::ZERO {
stats.l0_tables += 1;
}
level_entry.tables += 1;
level_entry.bytes = level_entry.bytes.saturating_add(table_bytes);
for reference in &properties.blob_references {
live_blob_bytes_by_file
.entry(reference.file_id)
.and_modify(|bytes| {
*bytes = bytes.saturating_add(reference.referenced_bytes);
})
.or_insert(reference.referenced_bytes);
}
}
}
}
stats.level_tables = level_stats.into_values().collect();
stats.level_filters = level_filter_stats.into_values().collect();
stats.live_blob_files = live_blob_bytes_by_file.len();
stats.live_blob_bytes = live_blob_bytes_by_file.values().copied().sum();
if let Some(db_path) = persistent_path {
add_obsolete_blob_stats(
&self.inner.native_storage,
db_path,
&live_blob_bytes_by_file,
&mut stats,
);
}
stats
}
pub(in crate::db) fn base_stats(&self) -> DbStats {
DbStats {
active_snapshots: self.inner.snapshots.active_count(),
compaction_runs: self.inner.compaction_runs.load(Ordering::Acquire),
compaction_input_tables: self.inner.compaction_input_tables.load(Ordering::Acquire),
compaction_output_tables: self.inner.compaction_output_tables.load(Ordering::Acquire),
compaction_input_bytes: self.inner.compaction_input_bytes.load(Ordering::Acquire),
compaction_output_bytes: self.inner.compaction_output_bytes.load(Ordering::Acquire),
commit_sequences_allocated: self.inner.commit_tracker.last_reserved_sequence().get(),
commit_visible_sequence: self.inner.commit_tracker.visible_sequence().get(),
commit_open_slots: self.inner.commit_tracker.open_slot_count(),
commit_skipped_slots: self.inner.commit_tracker.skipped_slot_count(),
blob_gc_runs: self.inner.blob_gc_runs.load(Ordering::Acquire),
blob_gc_input_bytes: self.inner.blob_gc_input_bytes.load(Ordering::Acquire),
blob_gc_output_bytes: self.inner.blob_gc_output_bytes.load(Ordering::Acquire),
blob_gc_discarded_bytes: self.inner.blob_gc_discarded_bytes.load(Ordering::Acquire),
maintenance_cooperative_yields: self
.inner
.maintenance_cooperative_yields
.load(Ordering::Acquire),
maintenance_budget_exhaustions: self
.inner
.maintenance_budget_exhaustions
.load(Ordering::Acquire),
..DbStats::default()
}
}
pub(in crate::db) fn add_compaction_level_stats(&self, stats: &mut DbStats) {
if let Ok(compaction_levels) = self.inner.compaction_level_stats.lock() {
stats.compaction_levels = compaction_levels.values().cloned().collect();
}
}
pub(in crate::db) fn add_compaction_trigger_stats(&self, stats: &mut DbStats) {
if let Ok(compaction_triggers) = self.inner.compaction_trigger_stats.lock() {
stats.compaction_triggers = compaction_triggers.values().cloned().collect();
}
}
pub(in crate::db) fn add_scan_and_snapshot_health_stats(&self, stats: &mut DbStats) {
let scan_waste = self.inner.scan_waste.snapshot();
stats.scan_internal_records = scan_waste.internal_records;
stats.scan_user_keys = scan_waste.user_keys;
stats.scan_tombstone_hidden_keys = scan_waste.tombstone_hidden_keys;
let visible = self.inner.commit_tracker.visible_sequence();
let oldest_snapshot = self.inner.snapshots.oldest_active_or(visible).min(visible);
stats.oldest_snapshot_seq = oldest_snapshot.get();
stats.oldest_snapshot_lag = visible.get().saturating_sub(oldest_snapshot.get());
}
pub(in crate::db) fn add_compaction_skip_stats(&self, stats: &mut DbStats) {
if let Ok(compaction_skips) = self.inner.compaction_skip_stats.lock() {
stats.compaction_skips = compaction_skips.values().cloned().collect();
}
}
pub(in crate::db) fn add_wal_stats(&self, stats: &mut DbStats) {
if let Some(wal_stats) = self.inner.substrate.wal_stats() {
stats.wal_shards = wal_stats.shards;
stats.wal_open_shards = wal_stats.open_shards;
stats.wal_queue_capacity = wal_stats.queue_capacity;
stats.wal_records_accepted = wal_stats.records_accepted;
stats.wal_bytes_accepted = wal_stats.bytes_accepted;
}
#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
if let Some(wal) = &self.inner.browser_wal {
let wal_stats = wal.stats();
stats.wal_shards = wal_stats.shards;
stats.wal_open_shards = wal_stats.open_shards;
stats.wal_queue_capacity = wal_stats.queue_capacity;
stats.wal_records_accepted = wal_stats.records_accepted;
stats.wal_bytes_accepted = wal_stats.bytes_accepted;
}
}
pub(in crate::db) fn add_storage_runtime_stats(&self, stats: &mut DbStats) {
let storage_stats = self.inner.native_storage.stats();
stats.storage_uses_sync_adapter = storage_stats.uses_blocking_adapter;
stats.storage_uses_platform_io_driver = storage_stats.uses_platform_io_driver;
stats.storage_uses_platform_async_io = storage_stats.uses_platform_async_io;
stats.storage_sync_adapter_tasks = storage_stats.blocking_adapter_tasks;
stats.storage_sync_adapter_queue_capacity = storage_stats.blocking_adapter_queue_capacity;
stats.storage_sync_adapter_queued_tasks = storage_stats.blocking_adapter_queued_tasks;
stats.storage_sync_adapter_submitted_tasks = storage_stats.blocking_adapter_submitted_tasks;
stats.storage_sync_adapter_completed_tasks = storage_stats.blocking_adapter_completed_tasks;
stats.storage_sync_adapter_rejected_tasks = storage_stats.blocking_adapter_rejected_tasks;
stats.storage_sync_adapter_total_runtime_micros =
storage_stats.blocking_adapter_total_runtime_micros;
stats.storage_platform_async_io_tasks = storage_stats.platform_async_io_tasks;
stats.storage_platform_thread_pool_managed_async_tasks =
storage_stats.platform_thread_pool_managed_async_tasks;
stats.storage_platform_sync_fallback_tasks = storage_stats.platform_blocking_fallback_tasks;
stats.storage_platform_io_operations = storage_stats.platform_io_operations;
stats.storage_inline_tasks = storage_stats.inline_tasks;
stats.storage_operations = storage_stats.operations;
}
#[must_use]
pub fn options(&self) -> &DbOptions {
&self.inner.options
}
#[must_use]
pub(crate) fn last_committed_sequence(&self) -> Sequence {
self.inner.commit_tracker.visible_sequence()
}
#[must_use]
pub fn latest_read_version(&self) -> ReadVersion {
ReadVersion::from_sequence(self.last_committed_sequence())
}
#[must_use]
pub fn oldest_retained_read_version(&self) -> ReadVersion {
ReadVersion::from_sequence(self.oldest_retained_sequence())
}
pub(in crate::db) fn oldest_retained_sequence(&self) -> Sequence {
let latest = self.last_committed_sequence();
self.inner
.snapshots
.oldest_active_or(latest)
.min(self.retained_floor_without_active_snapshots(latest))
}
pub(in crate::db) fn retained_floor_without_active_snapshots(
&self,
latest: Sequence,
) -> Sequence {
self.configured_retention_floor(latest)
.min(self.oldest_checkpoint_sequence_or(latest))
}
pub(in crate::db) fn configured_retention_floor(&self, latest: Sequence) -> Sequence {
let keep = self.inner.options.keep_last_read_versions.saturating_sub(1);
Sequence::new(latest.get().saturating_sub(keep))
}
pub(in crate::db) fn oldest_checkpoint_sequence_or(&self, fallback: Sequence) -> Sequence {
if let Some(manifest) = &self.inner.manifest {
return manifest
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.state()
.checkpoints()
.values()
.copied()
.min()
.unwrap_or(fallback);
}
self.inner
.checkpoints
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.values()
.copied()
.min()
.unwrap_or(fallback)
}
pub fn close_sync(&self) {
self.inner.closed.store(true, Ordering::Release);
shutdown_background_workers(
&self.inner.maintenance,
&self.inner.runtime_shutdown,
&self.inner.background_workers,
);
let Ok(()) = self.inner.publish_barrier.close() else {
return;
};
if let Some(db_path) = self.persistent_path().map(Path::to_path_buf) {
let _ = self.cleanup_pending_obsolete_table_files(&db_path);
let _ = self.cleanup_pending_obsolete_blob_files(&db_path);
}
super::super::release_browser_writer_lease(&self.inner);
self.inner.substrate.release_writer_lease();
}
pub(crate) fn ensure_open(&self) -> Result<()> {
if self.inner.closed.load(Ordering::Acquire) {
Err(Error::Closed)
} else {
Ok(())
}
}
}