mod scan;
#[cfg(test)]
mod tests;
mod versions;
use std::collections::HashMap;
use std::path::PathBuf;
use std::sync::Arc;
use nodedb_array::schema::ArraySchema;
use nodedb_array::tile::cell_payload::CellPayload;
use nodedb_array::types::coord::value::CoordValue;
use nodedb_wal::crypto::WalEncryptionKey;
use super::manifest::{Manifest, ManifestError, SegmentRef, segment_path};
use super::segment_handle::{SegmentHandle, SegmentHandleError};
use crate::engine::array::memtable::Memtable;
pub type CellVersion = (u64, Vec<CoordValue>, i64, CellPayload);
pub struct ArrayStore {
root: PathBuf,
schema: Arc<ArraySchema>,
schema_hash: u64,
manifest: Manifest,
pub(crate) memtable: Memtable,
pub(crate) segments: HashMap<String, SegmentHandle>,
next_segment_seq: u64,
kek: Option<WalEncryptionKey>,
}
#[derive(Debug, thiserror::Error)]
pub enum ArrayStoreError {
#[error(transparent)]
Manifest(#[from] ManifestError),
#[error(transparent)]
Segment(#[from] SegmentHandleError),
#[error("array store io: {detail}")]
Io { detail: String },
#[error("schema_hash mismatch: store={store:x} new={new:x}")]
SchemaHashMismatch { store: u64, new: u64 },
}
impl ArrayStore {
pub fn open(
root: PathBuf,
schema: Arc<ArraySchema>,
schema_hash: u64,
) -> Result<Self, ArrayStoreError> {
std::fs::create_dir_all(&root).map_err(|e| ArrayStoreError::Io {
detail: format!("mkdir {root:?}: {e}"),
})?;
let manifest = Manifest::load_or_new(&root, schema_hash)?;
if manifest.schema_hash != schema_hash && !manifest.segments.is_empty() {
return Err(ArrayStoreError::SchemaHashMismatch {
store: manifest.schema_hash,
new: schema_hash,
});
}
let mut segments = HashMap::with_capacity(manifest.segments.len());
let mut max_seq: u64 = 0;
for seg in &manifest.segments {
let h = SegmentHandle::open(
&segment_path(&root, &seg.id),
seg.id.clone(),
schema_hash,
None,
)?;
if let Some(seq) = parse_segment_seq(&seg.id) {
max_seq = max_seq.max(seq);
}
segments.insert(seg.id.clone(), h);
}
Ok(Self {
root,
schema,
schema_hash,
manifest,
memtable: Memtable::new(),
segments,
next_segment_seq: max_seq + 1,
kek: None,
})
}
pub fn set_kek(&mut self, kek: WalEncryptionKey) {
self.kek = Some(kek);
}
pub fn kek(&self) -> Option<&WalEncryptionKey> {
self.kek.as_ref()
}
pub fn root(&self) -> &std::path::Path {
&self.root
}
pub fn schema(&self) -> &Arc<ArraySchema> {
&self.schema
}
pub fn schema_hash(&self) -> u64 {
self.schema_hash
}
pub fn manifest(&self) -> &Manifest {
&self.manifest
}
pub fn manifest_mut(&mut self) -> &mut Manifest {
&mut self.manifest
}
pub fn allocate_segment_id(&mut self) -> String {
let seq = self.next_segment_seq;
self.next_segment_seq += 1;
format!("{seq:010}.ndas")
}
pub fn install_segment(&mut self, seg: SegmentRef) -> Result<(), ArrayStoreError> {
let h = SegmentHandle::open(
&segment_path(&self.root, &seg.id),
seg.id.clone(),
self.schema_hash,
self.kek.as_ref(),
)?;
self.segments.insert(seg.id.clone(), h);
self.manifest.append(seg);
Ok(())
}
pub fn replace_segments(
&mut self,
removed: &[String],
added: Vec<SegmentRef>,
) -> Result<(), ArrayStoreError> {
let mut new_handles = Vec::with_capacity(added.len());
for seg in &added {
let h = SegmentHandle::open(
&segment_path(&self.root, &seg.id),
seg.id.clone(),
self.schema_hash,
self.kek.as_ref(),
)?;
new_handles.push(h);
}
self.manifest.replace(removed, added);
for id in removed {
self.segments.remove(id);
}
for h in new_handles {
self.segments.insert(h.id().to_string(), h);
}
Ok(())
}
pub fn persist_manifest(&self) -> Result<(), ArrayStoreError> {
self.manifest.persist(&self.root)?;
Ok(())
}
pub fn unlink_segment(&self, id: &str) -> Result<(), ArrayStoreError> {
let path = segment_path(&self.root, id);
match std::fs::remove_file(&path) {
Ok(()) => Ok(()),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
Err(e) => Err(ArrayStoreError::Io {
detail: format!("unlink {path:?}: {e}"),
}),
}
}
}
fn parse_segment_seq(id: &str) -> Option<u64> {
id.split_once('.').and_then(|(stem, _)| stem.parse().ok())
}