use std::collections::{BTreeMap, HashMap};
use std::path::Path;
use frame_core::component::ComponentId;
use haematite::{
BranchHandle, BranchKind, BranchRefStore, BranchRegistry, CheckoutError, CommitDurability,
CommitRequest, ConflictPolicy, DiskStore, Hash, LeafNode, Node, NodeStore, SnapshotRegistry,
checkout, commit_branch, create_branch_with_kind, current_timestamp, fork_from_with_kind,
merge, open_branch, remove_branch,
};
use crate::codec::{
archive_meta_key, component_from_meta_key, decode_merge, decode_meta, encode_merge,
encode_meta, entity_bounds, entity_key, merge_bounds, merge_key, meta_bounds, meta_key,
};
use crate::error::StateError;
use crate::types::{EntityId, MergePolicy, MergeRecord, MetaRecord, MetaState};
pub(crate) const META_BRANCH: &str = "frame-state/meta";
const SHARD: usize = 0;
pub(crate) struct Storage {
pub(crate) nodes: DiskStore,
pub(crate) refs: BranchRefStore,
pub(crate) registry: BranchRegistry,
pub(crate) snapshots: SnapshotRegistry,
pub(crate) meta: BranchHandle,
pub(crate) empty_root: Hash,
pub(crate) active: HashMap<ComponentId, BranchHandle>,
pub(crate) work: HashMap<String, BranchHandle>,
}
impl Storage {
pub(crate) fn open(path: &Path) -> Result<Self, StateError> {
std::fs::create_dir_all(path)?;
let mut nodes = DiskStore::new(path.join("nodes"))?;
let mut refs = BranchRefStore::open(path.join("refs"))?;
let registry = BranchRegistry::new();
let snapshots = SnapshotRegistry::open(path.join("snapshots"))?;
let (meta, empty_root) = if let Some(record) = refs.get(META_BRANCH) {
let root = record
.shards
.first()
.ok_or_else(|| StateError::CorruptRecord {
detail: "meta branch has no shard".to_owned(),
})?
.fork_anchor;
(open_branch(META_BRANCH, &refs, ®istry)?, root)
} else {
let root = nodes.put(&Node::Leaf(LeafNode::new(Vec::new())?))?;
nodes.sync_dirty_dirs()?;
let branch = create_branch_with_kind(
META_BRANCH,
[(SHARD, root)],
BranchKind::Namespace,
&mut refs,
®istry,
current_timestamp(),
)?;
(branch, root)
};
Ok(Self {
nodes,
refs,
registry,
snapshots,
meta,
empty_root,
active: HashMap::new(),
work: HashMap::new(),
})
}
pub(crate) fn meta_record(
&self,
component: ComponentId,
) -> Result<Option<MetaRecord>, StateError> {
self.meta
.get(SHARD, &self.nodes, &meta_key(component))?
.map(|bytes| decode_meta(component, &bytes))
.transpose()
}
pub(crate) fn meta_records(&self) -> Result<Vec<MetaRecord>, StateError> {
let (from, to) = meta_bounds();
checkout(&self.nodes, self.meta.current_root())
.range(&from, &to)
.map(|row| {
let (key, value) = row?;
let component = component_from_meta_key(&key)?;
decode_meta(component, &value)
})
.collect()
}
pub(crate) fn write_meta(&mut self, record: &MetaRecord) -> Result<(), StateError> {
self.meta
.put(SHARD, meta_key(record.component), encode_meta(record)?)?;
self.commit_meta()
}
pub(crate) fn finish_archive_meta(&mut self, record: &MetaRecord) -> Result<(), StateError> {
let MetaState::Archived { generation, .. } = record.state else {
return Err(StateError::CorruptRecord {
detail: "archive evidence requires Archived state".to_owned(),
});
};
let encoded = encode_meta(record)?;
self.meta.put(SHARD, meta_key(record.component), &encoded)?;
self.meta.put(
SHARD,
archive_meta_key(record.component, generation),
encoded,
)?;
self.commit_meta()
}
fn commit_meta(&mut self) -> Result<(), StateError> {
commit_named(&self.meta, &mut self.nodes, &self.registry, &mut self.refs)?;
Ok(())
}
pub(crate) fn create_namespace(&mut self, component: ComponentId) -> Result<(), StateError> {
let name = namespace_name(component);
let branch = create_branch_with_kind(
&name,
[(SHARD, self.empty_root)],
BranchKind::Namespace,
&mut self.refs,
&self.registry,
current_timestamp(),
)?;
self.active.insert(component, branch);
Ok(())
}
pub(crate) fn namespace_exists(&self, component: ComponentId) -> bool {
self.refs.get(&namespace_name(component)).is_some()
}
pub(crate) fn namespace(&mut self, component: ComponentId) -> Result<BranchHandle, StateError> {
if let Some(branch) = self.active.get(&component) {
return Ok(branch.clone());
}
let name = namespace_name(component);
if self.refs.get(&name).is_none() {
return Err(StateError::NotActive { component });
}
let branch = open_branch(&name, &self.refs, &self.registry)?;
self.active.insert(component, branch.clone());
Ok(branch)
}
pub(crate) fn put_entity(
&mut self,
component: ComponentId,
entity: &[u8],
) -> Result<EntityId, StateError> {
let id = EntityId::of(entity);
self.namespace(component)?
.put(SHARD, entity_key(id), entity)?;
Ok(id)
}
pub(crate) fn get_entity(
&mut self,
component: ComponentId,
id: EntityId,
) -> Result<Option<Vec<u8>>, StateError> {
Ok(self
.namespace(component)?
.get(SHARD, &self.nodes, &entity_key(id))?)
}
pub(crate) fn delete_entity(
&mut self,
component: ComponentId,
id: EntityId,
) -> Result<(), StateError> {
self.namespace(component)?.delete(SHARD, entity_key(id))?;
Ok(())
}
pub(crate) fn enumerate_entities(
&mut self,
component: ComponentId,
) -> Result<Vec<(EntityId, Vec<u8>)>, StateError> {
let root = self.namespace(component)?.current_root();
entities_at(&self.nodes, root)
}
pub(crate) fn commit_namespace(&mut self, component: ComponentId) -> Result<Hash, StateError> {
let branch = self.namespace(component)?;
commit_named(&branch, &mut self.nodes, &self.registry, &mut self.refs)?;
Ok(branch.current_root())
}
pub(crate) fn fork_work(
&mut self,
component: ComponentId,
branch_name: &str,
) -> Result<(), StateError> {
let parent = self.namespace(component)?;
let branch = fork_from_with_kind(
&parent,
branch_name,
BranchKind::Work,
&mut self.refs,
&self.registry,
current_timestamp(),
)?;
self.work.insert(branch_name.to_owned(), branch);
Ok(())
}
pub(crate) fn work_branch(&mut self, name: &str) -> Result<BranchHandle, StateError> {
if let Some(branch) = self.work.get(name) {
return Ok(branch.clone());
}
if self.refs.get(name).is_none() {
return Err(StateError::WorkNotFound {
branch: name.to_owned(),
});
}
let branch = open_branch(name, &self.refs, &self.registry)?;
self.work.insert(name.to_owned(), branch.clone());
Ok(branch)
}
pub(crate) fn commit_work(&mut self, name: &str) -> Result<Hash, StateError> {
let branch = self.work_branch(name)?;
commit_named(&branch, &mut self.nodes, &self.registry, &mut self.refs)?;
Ok(branch.current_root())
}
pub(crate) fn merge_work(
&mut self,
component: ComponentId,
work_name: &str,
declared_policy: &MergePolicy,
engine_policy: &ConflictPolicy,
) -> Result<MergeRecord, StateError> {
let source = self.work_branch(work_name)?;
commit_named(&source, &mut self.nodes, &self.registry, &mut self.refs)?;
let destination = self.namespace(component)?;
commit_named(
&destination,
&mut self.nodes,
&self.registry,
&mut self.refs,
)?;
let source_root = source.current_root();
let destination_root_before = destination.current_root();
let report = merge(
&mut self.nodes,
&source,
&destination,
&self.refs,
engine_policy,
)?;
apply_root_diff(
&destination,
&self.nodes,
destination_root_before,
report.merged_root,
)?;
let sequence = self.next_merge_sequence(destination.current_root())?;
let record = MergeRecord {
sequence,
policy: declared_policy.clone(),
source_branch: work_name.to_owned(),
destination_branch: namespace_name(component),
source_root,
destination_root_before,
merged_content_root: report.merged_root,
};
destination.put(SHARD, merge_key(sequence), encode_merge(&record)?)?;
commit_named(
&destination,
&mut self.nodes,
&self.registry,
&mut self.refs,
)?;
let removed = remove_branch(work_name, &mut self.refs)?;
if removed.is_none() {
return Err(StateError::WorkNotFound {
branch: work_name.to_owned(),
});
}
self.work.remove(work_name);
Ok(record)
}
fn next_merge_sequence(&self, root: Hash) -> Result<u64, StateError> {
let (from, to) = merge_bounds();
let count = checkout(&self.nodes, root)
.range(&from, &to)
.collect::<Result<Vec<_>, _>>()?
.len();
u64::try_from(count).map_err(|_| StateError::CorruptRecord {
detail: "merge record count exceeds u64".to_owned(),
})
}
pub(crate) fn merge_records(
&mut self,
component: ComponentId,
) -> Result<Vec<MergeRecord>, StateError> {
let root = self.namespace(component)?.current_root();
let (from, to) = merge_bounds();
checkout(&self.nodes, root)
.range(&from, &to)
.map(|row| {
let (_, value) = row?;
decode_merge(&value)
})
.collect()
}
}
pub(crate) fn namespace_name(component: ComponentId) -> String {
format!("component/{component}")
}
pub(crate) fn work_name(component: ComponentId, name: &str) -> String {
const HEX: &[u8; 16] = b"0123456789abcdef";
let mut encoded = String::with_capacity(name.len().saturating_mul(2));
for byte in name.bytes() {
encoded.push(char::from(HEX[usize::from(byte >> 4)]));
encoded.push(char::from(HEX[usize::from(byte & 0x0f)]));
}
format!("work/{component}/{encoded}")
}
pub(crate) fn archive_name(component: ComponentId, generation: u64) -> String {
format!("archive/{component}/{generation:020}")
}
pub(crate) fn entities_at(
nodes: &DiskStore,
root: Hash,
) -> Result<Vec<(EntityId, Vec<u8>)>, StateError> {
let (from, to) = entity_bounds();
checkout(nodes, root)
.range(&from, &to)
.map(|row| {
let (key, value) =
row.map_err(|error| StateError::Checkout(CheckoutError::Tree(error)))?;
let bytes = <[u8; 32]>::try_from(key.get(1..).unwrap_or_default()).map_err(|_| {
StateError::CorruptRecord {
detail: "entity key is not region byte plus 32-byte digest".to_owned(),
}
})?;
Ok((EntityId::from_bytes(bytes), value))
})
.collect()
}
pub(crate) fn commit_named(
branch: &BranchHandle,
nodes: &mut DiskStore,
registry: &BranchRegistry,
refs: &mut BranchRefStore,
) -> Result<(), StateError> {
commit_branch(
branch,
nodes,
registry,
CommitRequest {
durability: CommitDurability::Durable { refs },
extra_parents: &[],
timestamp: current_timestamp(),
},
)?;
Ok(())
}
fn apply_root_diff(
destination: &BranchHandle,
nodes: &DiskStore,
before: Hash,
after: Hash,
) -> Result<(), StateError> {
let before_entries = all_entries(nodes, before)?;
let after_entries = all_entries(nodes, after)?;
for key in before_entries.keys() {
if !after_entries.contains_key(key) {
destination.delete(SHARD, key)?;
}
}
for (key, value) in after_entries {
if before_entries.get(&key) != Some(&value) {
destination.put(SHARD, key, value)?;
}
}
Ok(())
}
fn all_entries(nodes: &DiskStore, root: Hash) -> Result<BTreeMap<Vec<u8>, Vec<u8>>, StateError> {
checkout(nodes, root)
.range(&[], &[u8::MAX])
.collect::<Result<BTreeMap<_, _>, _>>()
.map_err(StateError::from)
}