use super::{
BranchId, Commit, CommitId, NO_COMMIT,
diff::{DiffEntry, diff_roots},
handle::{Branch, BranchMut},
};
use crate::basic::persistent_btree::{EMPTY_ROOT, NodeId, PersistentBTree};
use crate::common::ende::{KeyEnDeOrdered, ValueEnDe};
use crate::common::error::{Result, VsdbError};
use crate::{Mapx, MapxOrd, Orphan};
use serde::{Deserialize, Serialize};
use std::{
collections::{BinaryHeap, HashMap, HashSet},
marker::PhantomData,
time::{SystemTime, UNIX_EPOCH},
};
#[derive(Clone, Debug, Serialize, Deserialize)]
pub(crate) struct BranchState {
pub(crate) name: String,
pub(crate) head: CommitId,
pub(crate) dirty_root: NodeId,
}
#[derive(Clone, Debug)]
pub struct VerMap<K, V> {
pub(crate) tree: PersistentBTree,
pub(crate) commits: MapxOrd<u64, Commit>,
pub(crate) branches: MapxOrd<u64, BranchState>,
pub(crate) branch_names: Mapx<String, u64>,
pub(crate) next_commit: Orphan<u64>,
pub(crate) next_branch: Orphan<u64>,
pub(crate) main_branch: Orphan<u64>,
pub(crate) gc_dirty: Orphan<bool>,
_phantom: PhantomData<(K, V)>,
}
impl<K, V> Serialize for VerMap<K, V> {
fn serialize<S>(&self, serializer: S) -> std::result::Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
crate::common::serialize_typed_handle_meta::<Self, S>(
&(
&self.tree,
&self.commits,
&self.branches,
&self.branch_names,
&self.next_commit,
&self.next_branch,
&self.main_branch,
&self.gc_dirty,
),
serializer,
)
}
}
type VerMapPayload = (
PersistentBTree,
MapxOrd<u64, Commit>,
MapxOrd<u64, BranchState>,
Mapx<String, u64>,
Orphan<u64>,
Orphan<u64>,
Orphan<u64>,
Orphan<bool>,
);
impl<'de, K, V> Deserialize<'de> for VerMap<K, V> {
fn deserialize<D>(deserializer: D) -> std::result::Result<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
let (
tree,
commits,
branches,
branch_names,
next_commit,
next_branch,
main_branch,
gc_dirty,
) = crate::common::deserialize_typed_handle_meta::<Self, VerMapPayload, D>(
deserializer,
)?;
let mut m = VerMap {
tree,
commits,
branches,
branch_names,
next_commit,
next_branch,
main_branch,
gc_dirty,
_phantom: PhantomData,
};
m.repair_commit_ref_counts_if_needed();
m.rebuild_tree_ref_counts();
Ok(m)
}
}
impl<K, V> Default for VerMap<K, V>
where
K: KeyEnDeOrdered,
V: ValueEnDe,
{
fn default() -> Self {
Self::new()
}
}
impl<K, V> VerMap<K, V>
where
K: KeyEnDeOrdered,
V: ValueEnDe,
{
pub fn new() -> Self {
Self::new_with_main("main")
}
pub fn instance_id(&self) -> u64 {
self.tree.instance_id()
}
pub fn save_meta(&self) -> Result<u64> {
let id = self.instance_id();
crate::common::save_instance_meta(id, self)?;
Ok(id)
}
pub fn from_meta(instance_id: u64) -> Result<Self> {
crate::common::load_instance_meta(instance_id)
}
pub fn new_with_main(name: &str) -> Self {
let mut branches: MapxOrd<u64, BranchState> = MapxOrd::new();
let mut branch_names: Mapx<String, u64> = Mapx::new();
let initial_id: BranchId = 1;
let main = BranchState {
name: name.into(),
head: NO_COMMIT,
dirty_root: EMPTY_ROOT,
};
branches.insert(&initial_id, &main);
branch_names.insert(&name.to_string(), &initial_id);
Self {
tree: PersistentBTree::new(),
commits: MapxOrd::new(),
branches,
branch_names,
next_commit: Orphan::new(1), next_branch: Orphan::new(initial_id + 1),
main_branch: Orphan::new(initial_id),
gc_dirty: Orphan::new(false),
_phantom: PhantomData,
}
}
pub(crate) fn get_branch(&self, id: BranchId) -> Result<BranchState> {
self.branches
.get(&id)
.ok_or(VsdbError::BranchNotFound { branch_id: id })
}
pub(crate) fn get_commit_inner(&self, id: CommitId) -> Result<Commit> {
self.commits
.get(&id)
.ok_or(VsdbError::CommitNotFound { commit_id: id })
}
fn branch_name_exists(&self, name: &str) -> bool {
self.branches.iter().any(|(_, state)| state.name == name)
}
pub fn main_branch(&self) -> BranchId {
self.main_branch.get_value()
}
pub fn set_main_branch(&mut self, branch: BranchId) -> Result<()> {
self.get_branch(branch)?;
*self.main_branch.get_mut() = branch;
Ok(())
}
pub fn create_branch(
&mut self,
name: &str,
source_branch: BranchId,
) -> Result<BranchId> {
if self.branch_name_exists(name) {
return Err(VsdbError::BranchAlreadyExists {
name: name.to_string(),
});
}
let src = self.get_branch(source_branch)?;
*self.gc_dirty.get_mut() = true;
let id = self.next_branch.get_value();
*self.next_branch.get_mut() = id + 1;
let state = BranchState {
name: name.into(),
head: src.head,
dirty_root: src.dirty_root,
};
self.branches.insert(&id, &state);
self.branch_names.insert(&name.to_string(), &id);
self.increment_ref(src.head);
self.tree.acquire_node(src.dirty_root);
*self.gc_dirty.get_mut() = false;
Ok(id)
}
pub fn delete_branch(&mut self, branch: BranchId) -> Result<()> {
if branch == self.main_branch.get_value() {
return Err(VsdbError::CannotDeleteMainBranch);
}
let state = self.get_branch(branch)?;
let dead_head = state.head;
let dead_dirty = state.dirty_root;
*self.gc_dirty.get_mut() = true;
self.branch_names.remove(&state.name);
self.branches.remove(&branch);
self.tree.release_node(dead_dirty);
self.decrement_ref(dead_head);
*self.gc_dirty.get_mut() = false;
Ok(())
}
pub fn list_branches(&self) -> Vec<(BranchId, String)> {
self.branches.iter().map(|(id, s)| (id, s.name)).collect()
}
pub fn branch_id(&self, name: &str) -> Option<BranchId> {
self.branch_names.get(&name.to_string())
}
pub fn branch_name(&self, branch: BranchId) -> Option<String> {
self.branches.get(&branch).map(|s| s.name)
}
pub fn has_uncommitted(&self, branch: BranchId) -> Result<bool> {
let state = self.get_branch(branch)?;
if state.head == NO_COMMIT {
Ok(state.dirty_root != EMPTY_ROOT)
} else {
let head_root = self.get_commit_inner(state.head)?.root;
Ok(state.dirty_root != head_root)
}
}
pub fn insert(&mut self, branch: BranchId, key: &K, value: &V) -> Result<()> {
let mut state = self.get_branch(branch)?;
let old_root = state.dirty_root;
state.dirty_root = self.tree.insert(old_root, &key.to_bytes(), &value.encode());
self.tree.acquire_node(state.dirty_root);
self.branches.insert(&branch, &state);
self.tree.release_node(old_root);
Ok(())
}
pub fn remove(&mut self, branch: BranchId, key: &K) -> Result<()> {
let mut state = self.get_branch(branch)?;
let old_root = state.dirty_root;
state.dirty_root = self.tree.remove(old_root, &key.to_bytes());
self.tree.acquire_node(state.dirty_root);
self.branches.insert(&branch, &state);
self.tree.release_node(old_root);
Ok(())
}
pub fn commit(&mut self, branch: BranchId) -> Result<CommitId> {
let state = self.get_branch(branch)?;
*self.gc_dirty.get_mut() = true;
let id = self.next_commit.get_value();
*self.next_commit.get_mut() = id + 1;
let parents = if state.head == NO_COMMIT {
vec![]
} else {
vec![state.head]
};
let commit = Commit {
id,
root: state.dirty_root,
parents,
timestamp_us: now_us(),
ref_count: 1,
};
self.commits.insert(&id, &commit);
self.tree.acquire_node(state.dirty_root);
let new_state = BranchState { head: id, ..state };
self.branches.insert(&branch, &new_state);
*self.gc_dirty.get_mut() = false;
Ok(id)
}
pub fn discard(&mut self, branch: BranchId) -> Result<()> {
let state = self.get_branch(branch)?;
let old_dirty = state.dirty_root;
let root = if state.head == NO_COMMIT {
EMPTY_ROOT
} else {
self.get_commit_inner(state.head)?.root
};
let new_state = BranchState {
dirty_root: root,
..state
};
self.tree.acquire_node(root);
self.branches.insert(&branch, &new_state);
self.tree.release_node(old_dirty);
Ok(())
}
pub fn rollback_to(&mut self, branch: BranchId, target: CommitId) -> Result<()> {
let state = self.get_branch(branch)?;
let _ = self.get_commit_inner(target)?;
if state.head == NO_COMMIT {
return Err(VsdbError::Other {
detail: "target commit is not an ancestor of this branch's head".into(),
});
}
if target != state.head {
let mut queue = vec![state.head];
let mut visited = HashSet::new();
let mut found = false;
while let Some(cur) = queue.pop() {
if cur == NO_COMMIT || !visited.insert(cur) {
continue;
}
if cur == target {
found = true;
break;
}
if let Some(c) = self.commits.get(&cur) {
queue.extend_from_slice(&c.parents);
}
}
if !found {
return Err(VsdbError::Other {
detail: "target commit is not an ancestor of this branch's head"
.into(),
});
}
} else if self.has_uncommitted(branch)? {
return Err(VsdbError::UncommittedChanges { branch_id: branch });
} else {
return Ok(());
}
*self.gc_dirty.get_mut() = true;
let commit = self.get_commit_inner(target)?;
let old_head = state.head;
let old_dirty = state.dirty_root;
let new_state = BranchState {
name: state.name,
head: target,
dirty_root: commit.root,
};
self.branches.insert(&branch, &new_state);
self.tree.acquire_node(commit.root);
self.tree.release_node(old_dirty);
self.increment_ref(target);
self.decrement_ref(old_head);
*self.gc_dirty.get_mut() = false;
Ok(())
}
pub fn merge(&mut self, source: BranchId, target: BranchId) -> Result<CommitId> {
if source == target {
return Err(VsdbError::Other {
detail: "cannot merge a branch into itself".into(),
});
}
if self.has_uncommitted(source)? {
return Err(VsdbError::UncommittedChanges { branch_id: source });
}
if self.has_uncommitted(target)? {
return Err(VsdbError::UncommittedChanges { branch_id: target });
}
let src = self.get_branch(source)?;
let tgt = self.get_branch(target)?;
if src.head == NO_COMMIT {
return Err(VsdbError::Other {
detail: format!("source branch {source} has no commits"),
});
}
if src.head == tgt.head {
return Ok(tgt.head);
}
*self.gc_dirty.get_mut() = true;
if tgt.head == NO_COMMIT {
let src_commit = self.get_commit_inner(src.head)?;
let new_state = BranchState {
head: src.head,
dirty_root: src_commit.root,
..tgt
};
self.branches.insert(&target, &new_state);
self.increment_ref(src.head);
self.tree.acquire_node(src_commit.root);
self.tree.release_node(tgt.dirty_root);
*self.gc_dirty.get_mut() = false;
return Ok(src.head);
}
let src_commit = self.get_commit_inner(src.head)?;
let tgt_commit = self.get_commit_inner(tgt.head)?;
let ancestor_roots: Vec<NodeId> = self
.find_merge_bases(src.head, tgt.head)
.into_iter()
.map(|aid| self.get_commit_inner(aid).map(|c| c.root))
.collect::<Result<_>>()?;
let ancestor_roots = if ancestor_roots.is_empty() {
vec![EMPTY_ROOT]
} else {
ancestor_roots
};
let merged_root = super::merge::three_way_merge_many_bases(
&mut self.tree,
&ancestor_roots,
src_commit.root,
tgt_commit.root,
);
let id = self.next_commit.get_value();
*self.next_commit.get_mut() = id + 1;
let commit = Commit {
id,
root: merged_root,
parents: vec![tgt.head, src.head],
timestamp_us: now_us(),
ref_count: 1,
};
self.commits.insert(&id, &commit);
let new_state = BranchState {
head: id,
dirty_root: merged_root,
..tgt
};
self.branches.insert(&target, &new_state);
self.tree.acquire_node(merged_root); self.tree.acquire_node(merged_root); self.tree.release_node(tgt.dirty_root);
self.increment_ref(src.head);
*self.gc_dirty.get_mut() = false;
Ok(id)
}
fn find_merge_bases(&self, a: CommitId, b: CommitId) -> Vec<CommitId> {
const FROM_A: u8 = 0b001;
const FROM_B: u8 = 0b010;
const BOTH: u8 = FROM_A | FROM_B;
const STALE: u8 = 0b100;
if self.commits.get(&a).is_none() || self.commits.get(&b).is_none() {
return vec![];
}
fn mark(
flags: &mut HashMap<CommitId, u8>,
heap: &mut BinaryHeap<CommitId>,
id: CommitId,
flag: u8,
) {
if id == NO_COMMIT {
return;
}
let slot = flags.entry(id).or_insert(0);
if *slot == 0 {
heap.push(id);
}
*slot |= flag;
}
let mut flags: HashMap<CommitId, u8> = HashMap::new();
let mut heap: BinaryHeap<CommitId> = BinaryHeap::new();
let mut bases = Vec::new();
mark(&mut flags, &mut heap, a, FROM_A);
mark(&mut flags, &mut heap, b, FROM_B);
while heap.iter().any(|id| flags[id] & STALE == 0) {
let id = heap.pop().unwrap();
let mut f = flags[&id];
if f & BOTH == BOTH {
if f & STALE == 0 {
bases.push(id);
}
f |= STALE;
}
if let Some(c) = self.commits.get(&id) {
for &parent in &c.parents {
mark(&mut flags, &mut heap, parent, f);
}
}
}
bases
}
fn find_common_ancestor(&self, a: CommitId, b: CommitId) -> Option<CommitId> {
self.find_merge_bases(a, b).into_iter().max()
}
pub fn fork_point(&self, a: CommitId, b: CommitId) -> Option<CommitId> {
self.find_common_ancestor(a, b)
}
pub fn commit_distance(&self, from: CommitId, ancestor: CommitId) -> Option<u64> {
self.commits.get(&from)?;
self.commits.get(&ancestor)?;
let mut cur = from;
let mut count = 0u64;
while cur != ancestor {
if cur == NO_COMMIT {
return None;
}
let c = self.commits.get(&cur)?;
cur = c.parents.first().copied().unwrap_or(NO_COMMIT);
count += 1;
}
Some(count)
}
pub fn get_commit(&self, commit_id: CommitId) -> Option<Commit> {
self.commits.get(&commit_id)
}
pub fn head_commit(&self, branch: BranchId) -> Result<Option<Commit>> {
let state = self.get_branch(branch)?;
if state.head == NO_COMMIT {
Ok(None)
} else {
Ok(self.commits.get(&state.head))
}
}
pub fn log(&self, branch: BranchId) -> Result<Vec<Commit>> {
let state = self.get_branch(branch)?;
let mut result = Vec::new();
let mut cur = state.head;
while cur != NO_COMMIT {
if let Some(c) = self.commits.get(&cur) {
cur = c.parents.first().copied().unwrap_or(NO_COMMIT);
result.push(c);
} else {
break;
}
}
Ok(result)
}
pub fn diff_commits(&self, from: CommitId, to: CommitId) -> Result<Vec<DiffEntry>> {
let from_commit = self.get_commit_inner(from)?;
let to_commit = self.get_commit_inner(to)?;
Ok(diff_roots(&self.tree, from_commit.root, to_commit.root))
}
pub fn diff_uncommitted(&self, branch: BranchId) -> Result<Vec<DiffEntry>> {
let state = self.get_branch(branch)?;
let head_root = if state.head == NO_COMMIT {
EMPTY_ROOT
} else {
self.get_commit_inner(state.head)?.root
};
Ok(diff_roots(&self.tree, head_root, state.dirty_root))
}
pub fn gc(&mut self) {
if self.gc_dirty.get_value()
|| self.commits.iter().any(|(_, c)| c.ref_count == 0)
{
self.rebuild_ref_counts();
}
let mut live_roots: Vec<NodeId> =
self.commits.iter().map(|(_, c)| c.root).collect();
for (_, s) in self.branches.iter() {
if s.dirty_root != EMPTY_ROOT {
live_roots.push(s.dirty_root);
}
}
self.tree.gc(&live_roots);
}
pub fn branch(&self, id: BranchId) -> Result<Branch<'_, K, V>> {
self.get_branch(id)?;
Ok(Branch { map: self, id })
}
pub fn branch_mut(&mut self, id: BranchId) -> Result<BranchMut<'_, K, V>> {
self.get_branch(id)?;
Ok(BranchMut { map: self, id })
}
pub fn main(&self) -> Branch<'_, K, V> {
Branch {
map: self,
id: self.main_branch(),
}
}
pub fn main_mut(&mut self) -> BranchMut<'_, K, V> {
let id = self.main_branch();
BranchMut { map: self, id }
}
fn increment_ref(&mut self, commit_id: CommitId) {
if commit_id == NO_COMMIT {
return;
}
if let Some(mut c) = self.commits.get(&commit_id) {
c.ref_count += 1;
self.commits.insert(&commit_id, &c);
}
}
fn decrement_ref(&mut self, commit_id: CommitId) {
if commit_id == NO_COMMIT {
return;
}
let already_dirty = self.gc_dirty.get_value();
*self.gc_dirty.get_mut() = true;
let mut work = vec![commit_id];
while let Some(id) = work.pop() {
if id == NO_COMMIT {
continue;
}
let Some(mut c) = self.commits.get(&id) else {
continue; };
c.ref_count = c.ref_count.saturating_sub(1);
if c.ref_count == 0 {
let parents = c.parents.clone();
self.tree.release_node(c.root);
self.commits.remove(&id);
work.extend(parents);
} else {
self.commits.insert(&id, &c);
}
}
if !already_dirty {
*self.gc_dirty.get_mut() = false;
}
}
}
fn now_us() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_micros() as u64)
.unwrap_or(0)
}