use std::path::{Path, PathBuf};
use crate::branch::{BranchRefStore, Timestamp, current_timestamp};
use crate::ids::ShardId;
use crate::shard::router::{SHARD_STORE_DIR, SHARD_WAL_FILE};
use crate::store::{DiskStore, NodeStore};
use crate::tree::{Hash, LeafNode, Node, TreePolicy, batch_mutate_owned};
use crate::wal::{DurableWal, FsyncPolicy, WalRecovery};
use super::super::DatabaseError;
use super::super::vacuum::CANONICAL_REFS_DIR;
fn shard_paths(data_dir: &Path, shard_id: usize) -> (PathBuf, PathBuf, PathBuf) {
let shard_dir = data_dir.join(format!("shard-{shard_id}"));
let store_dir = shard_dir.join(SHARD_STORE_DIR);
let wal_path = shard_dir.join(SHARD_WAL_FILE);
(shard_dir, store_dir, wal_path)
}
fn fail(context: &str, error: impl std::fmt::Display) -> DatabaseError {
DatabaseError::MigrationFailed(format!("{context}: {error}"))
}
pub(super) enum ShardRebuild {
Empty,
Unchanged,
Rebuilt,
}
fn collect_entries<S: NodeStore + ?Sized>(
store: &S,
root: Hash,
out: &mut Vec<(Vec<u8>, Vec<u8>)>,
) -> Result<(), DatabaseError> {
let node = store
.get(&root)
.map_err(|error| fail("stream shard node", error))?
.ok_or_else(|| {
DatabaseError::MigrationFailed(format!("node {root:?} missing while streaming a root"))
})?;
match &*node {
Node::Leaf(leaf) => out.extend(leaf.entries().iter().cloned()),
Node::Internal(internal) => {
for (_separator, child) in internal.children() {
collect_entries(store, *child, out)?;
}
}
}
Ok(())
}
fn empty_root<S: NodeStore + ?Sized>(store: &mut S) -> Result<Hash, DatabaseError> {
let leaf = LeafNode::new(Vec::new()).map_err(|error| fail("empty leaf", error))?;
store
.put(&Node::Leaf(leaf))
.map_err(|error| fail("store empty leaf", error))
}
fn rebuild_root<S: NodeStore + ?Sized>(
store: &mut S,
source_root: Hash,
policy: TreePolicy,
) -> Result<Hash, DatabaseError> {
let mut entries = Vec::new();
collect_entries(store, source_root, &mut entries)?;
let baseline = empty_root(store)?;
if entries.is_empty() {
return Ok(baseline);
}
let batch: Vec<(Vec<u8>, Option<Vec<u8>>)> = entries
.into_iter()
.map(|(key, value)| (key, Some(value)))
.collect();
batch_mutate_owned(store, baseline, batch, policy).map_err(|error| fail("rebuild root", error))
}
pub(super) fn rebuild_shard(
data_dir: &Path,
shard_id: usize,
policy: TreePolicy,
) -> Result<ShardRebuild, DatabaseError> {
let (shard_dir, store_dir, wal_path) = shard_paths(data_dir, shard_id);
if !shard_dir.exists() {
return Ok(ShardRebuild::Empty);
}
let mut store = DiskStore::new(&store_dir).map_err(|error| fail("open shard store", error))?;
let recovered = WalRecovery::recover_path(&wal_path, &store)
.map_err(|error| fail("recover shard wal", error))?;
let Some(source_root) = recovered.committed_root() else {
return Ok(ShardRebuild::Empty);
};
let new_root = rebuild_root(&mut store, source_root, policy)?;
if new_root == source_root {
return Ok(ShardRebuild::Unchanged);
}
store
.sync_dirty_dirs()
.map_err(|error| fail("sync shard store dirs", error))?;
let mut wal = DurableWal::new(&wal_path, FsyncPolicy::CommitOnly)
.map_err(|error| fail("open shard wal", error))?;
wal.commit(new_root)
.map_err(|error| fail("commit shard root", error))?;
Ok(ShardRebuild::Rebuilt)
}
fn branches_dir(data_dir: &Path) -> Option<PathBuf> {
let dir = data_dir.join(CANONICAL_REFS_DIR);
dir.is_dir().then_some(dir)
}
pub(super) fn branch_names(data_dir: &Path) -> Result<Vec<String>, DatabaseError> {
let Some(dir) = branches_dir(data_dir) else {
return Ok(Vec::new());
};
let refs = BranchRefStore::open(&dir).map_err(|error| fail("open branch refstore", error))?;
Ok(refs.list().map(|record| record.name.clone()).collect())
}
pub(super) enum BranchRebuild {
Unchanged,
Rebuilt,
}
pub(super) fn rebuild_one_branch(
data_dir: &Path,
name: &str,
policy: TreePolicy,
) -> Result<BranchRebuild, DatabaseError> {
let Some(dir) = branches_dir(data_dir) else {
return Ok(BranchRebuild::Unchanged);
};
let mut refs =
BranchRefStore::open(&dir).map_err(|error| fail("open branch refstore", error))?;
let Some(record) = refs.get(name) else {
return Ok(BranchRebuild::Unchanged);
};
let created = record.created;
let seq = record.seq;
let parents = record.parents.clone();
let shards: Vec<(ShardId, Hash)> = record
.shards
.iter()
.map(|shard| (shard.shard_id, shard.head))
.collect();
let mut new_heads: Vec<(ShardId, Hash)> = Vec::with_capacity(shards.len());
let mut changed = false;
for (shard_id, head) in shards {
let (_shard_dir, store_dir, _wal_path) = shard_paths(data_dir, shard_id);
let mut store = DiskStore::new(&store_dir)
.map_err(|error| fail("open shard store for branch", error))?;
let new_head = rebuild_root(&mut store, head, policy)?;
if new_head != head {
changed = true;
store
.sync_dirty_dirs()
.map_err(|error| fail("sync branch head store dirs", error))?;
}
new_heads.push((shard_id, new_head));
}
if !changed {
return Ok(BranchRebuild::Unchanged);
}
let timestamp: Timestamp = current_timestamp();
refs.advance(name, created, seq, &new_heads, parents, timestamp)
.map_err(|error| fail("advance branch head", error))?;
Ok(BranchRebuild::Rebuilt)
}