#![forbid(unsafe_code)]
use std::path::Path;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Condvar, Mutex};
use crate::core::extent::ChunkId;
use crate::optimizer::foreground::ForegroundPolicy;
use crate::optimizer::policy::OptimizeOptions;
use crate::store::directory;
use crate::store::transaction::CrashHooks;
use crate::store::{NewEntry, Store, StoreConfig, StoreError};
pub mod metrics;
pub use metrics::{
AccountingMetrics, CacheMetrics, DsfbMetrics, EngineMetrics, FormatInfo, GcMetrics,
METRIC_REGISTRY, MetricDef, PhaseMetrics, PhysicalMetrics,
};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
#[repr(i32)]
pub enum ErrorCode {
Ok = 0,
NotFound = 1,
InvalidArgument = 2,
CorruptStore = 3,
IncompatibleFormat = 4,
ResourceLimit = 5,
Io = 6,
Busy = 7,
Unsupported = 8,
Internal = 9,
Closed = 10,
}
impl ErrorCode {
pub const fn as_i32(self) -> i32 {
self as i32
}
pub const fn from_i32(v: i32) -> Option<Self> {
Some(match v {
0 => Self::Ok,
1 => Self::NotFound,
2 => Self::InvalidArgument,
3 => Self::CorruptStore,
4 => Self::IncompatibleFormat,
5 => Self::ResourceLimit,
6 => Self::Io,
7 => Self::Busy,
8 => Self::Unsupported,
9 => Self::Internal,
10 => Self::Closed,
_ => return None,
})
}
pub const fn name(self) -> &'static str {
match self {
Self::Ok => "ok",
Self::NotFound => "not_found",
Self::InvalidArgument => "invalid_argument",
Self::CorruptStore => "corrupt_store",
Self::IncompatibleFormat => "incompatible_format",
Self::ResourceLimit => "resource_limit",
Self::Io => "io",
Self::Busy => "busy",
Self::Unsupported => "unsupported",
Self::Internal => "internal",
Self::Closed => "closed",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct EngineError {
pub code: ErrorCode,
pub message: String,
}
impl EngineError {
pub fn new(code: ErrorCode, message: impl Into<String>) -> Self {
Self {
code,
message: message.into(),
}
}
pub const fn code(&self) -> ErrorCode {
self.code
}
}
impl std::fmt::Display for EngineError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}: {}", self.code.name(), self.message)
}
}
impl std::error::Error for EngineError {}
impl From<StoreError> for EngineError {
fn from(e: StoreError) -> Self {
let (code, message) = match e {
StoreError::Segment(s) => (ErrorCode::CorruptStore, format!("segment: {s:?}")),
StoreError::Superblock(s) => (ErrorCode::CorruptStore, format!("superblock: {s}")),
StoreError::Index(s) => (ErrorCode::CorruptStore, format!("index: {s}")),
StoreError::MissingObject(cid) => (
ErrorCode::CorruptStore,
format!("referenced object {cid} is missing"),
),
StoreError::MissingChunk(cid) => (
ErrorCode::CorruptStore,
format!("chunk descriptor for {cid} is missing"),
),
StoreError::Descriptor(s) => (ErrorCode::CorruptStore, format!("descriptor: {s}")),
StoreError::NotOpen => (ErrorCode::Internal, "store not open".to_string()),
StoreError::Locked => (ErrorCode::Busy, "store is locked by another process".into()),
StoreError::Config(s) => (ErrorCode::InvalidArgument, format!("config: {s}")),
StoreError::Limit(s) => (ErrorCode::ResourceLimit, format!("limit: {s}")),
StoreError::Full(s) => (ErrorCode::ResourceLimit, format!("full: {s}")),
StoreError::Io(s) => (ErrorCode::Io, s),
StoreError::CrashSimulated(s) => (ErrorCode::Internal, format!("crash hook: {s}")),
StoreError::Invariant(s) => (ErrorCode::CorruptStore, format!("invariant: {s}")),
StoreError::ReadOnly => (ErrorCode::Unsupported, "store opened read-only".into()),
StoreError::IncompatibleFormat(err) => (ErrorCode::IncompatibleFormat, err.to_string()),
};
EngineError::new(code, message)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
pub struct BlobId(pub [u8; 32]);
impl BlobId {
pub const fn new(bytes: [u8; 32]) -> Self {
Self(bytes)
}
pub const fn as_bytes(&self) -> &[u8; 32] {
&self.0
}
pub const fn as_chunk_id(&self) -> ChunkId {
ChunkId::new(self.0)
}
pub fn from_hex(s: &str) -> Option<Self> {
ChunkId::from_hex(s).map(|c| Self(*c.as_bytes()))
}
}
impl From<ChunkId> for BlobId {
fn from(c: ChunkId) -> Self {
Self(*c.as_bytes())
}
}
impl std::fmt::Display for BlobId {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
for b in self.0 {
write!(f, "{b:02x}")?;
}
Ok(())
}
}
impl serde::Serialize for BlobId {
fn serialize<S: serde::Serializer>(&self, s: S) -> Result<S::Ok, S::Error> {
s.serialize_str(&self.to_string())
}
}
impl<'de> serde::Deserialize<'de> for BlobId {
fn deserialize<D: serde::Deserializer<'de>>(d: D) -> Result<Self, D::Error> {
let s = String::deserialize(d)?;
Self::from_hex(&s).ok_or_else(|| serde::de::Error::custom("invalid 64-hex blob id"))
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Durability {
Ack,
Durable,
}
#[derive(Debug, Clone, Copy, PartialEq)]
pub struct EngineOpenOptions {
pub io_backend: crate::store::io::IoBackendKind,
pub io_uring_entries: u32,
pub capacity_override: Option<u64>,
pub root_uid: u32,
pub root_gid: u32,
pub foreground: ForegroundPolicy,
pub read_only: bool,
}
impl Default for EngineOpenOptions {
fn default() -> Self {
let c = StoreConfig::default();
Self {
io_backend: c.io_backend,
io_uring_entries: c.io_uring_entries,
capacity_override: None,
root_uid: c.root_uid,
root_gid: c.root_gid,
foreground: c.foreground,
read_only: false,
}
}
}
const ENGINE_DIR_NAME: &[u8] = b".engine";
const TMP_PREFIX: &[u8] = b"blob-";
const TMP_SUFFIX: &[u8] = b"-tmp-";
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct CompactionReport {
pub unreachable_before_bytes: u64,
pub reclaimed_bytes: u64,
pub unreachable_after_bytes: u64,
pub physical_used_after_bytes: u64,
pub live_bytes_after: u64,
}
struct EngineInner {
store: Mutex<Option<Arc<Store>>>,
closed: AtomicBool,
ops_in_flight: AtomicU64,
drain_cv: Condvar,
}
pub struct Engine {
inner: Arc<EngineInner>,
}
impl std::fmt::Debug for Engine {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Engine").finish_non_exhaustive()
}
}
impl Engine {
fn generate_uuid(path: &Path) -> [u8; 16] {
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos())
.unwrap_or(0);
let mut h = blake3::Hasher::new();
h.update(&now.to_le_bytes());
h.update(&std::process::id().to_le_bytes());
h.update(path.to_string_lossy().as_bytes());
let out = h.finalize();
let mut u = [0u8; 16];
u.copy_from_slice(&out.as_bytes()[..16]);
u
}
pub fn create(path: &Path, opts: &EngineOpenOptions) -> Result<Self, EngineError> {
let config = StoreConfig {
io_backend: opts.io_backend,
io_uring_entries: opts.io_uring_entries,
capacity_override: opts.capacity_override,
root_uid: opts.root_uid,
root_gid: opts.root_gid,
foreground: opts.foreground,
..StoreConfig::default()
};
let uuid = Self::generate_uuid(path);
let store = Store::create(path, &config, uuid).map_err(EngineError::from)?;
let hooks = CrashHooks::none();
store
.epoch_create(
1,
ENGINE_DIR_NAME,
NewEntry::dir(0o700, opts.root_uid, opts.root_gid),
&hooks,
)
.map_err(EngineError::from)?;
store
.durability_barrier(&hooks)
.map_err(EngineError::from)?;
let engine = Self {
inner: Arc::new(EngineInner {
store: Mutex::new(Some(Arc::new(store))),
closed: AtomicBool::new(false),
ops_in_flight: AtomicU64::new(0),
drain_cv: Condvar::new(),
}),
};
Ok(engine)
}
pub fn open(path: &Path, opts: &EngineOpenOptions) -> Result<Self, EngineError> {
if !path.join("segments").is_dir() {
return Err(EngineError::new(
ErrorCode::NotFound,
format!("no entropyfs store at {}", path.display()),
));
}
let config = StoreConfig {
io_backend: opts.io_backend,
io_uring_entries: opts.io_uring_entries,
capacity_override: opts.capacity_override,
read_only: opts.read_only,
..StoreConfig::default()
};
let store = Store::open(path, &config).map_err(EngineError::from)?;
let hooks = CrashHooks::none();
let dir_ino = {
let ep = store.epoch();
match store
.dir_lookup_epoch(&ep, 1, ENGINE_DIR_NAME)
.map_err(EngineError::from)?
{
Some(e) if e.d_type == directory::dt::DT_DIR => e.ino,
Some(_) => {
return Err(EngineError::new(
ErrorCode::CorruptStore,
"engine namespace name is occupied by a non-directory entry",
));
}
None if opts.read_only => {
return Err(EngineError::new(
ErrorCode::NotFound,
"store has no engine blob namespace (created read-only; \
nothing to read)",
));
}
None => {
drop(ep);
let _ = store
.epoch_create(
1,
ENGINE_DIR_NAME,
NewEntry::dir(0o700, config.root_uid, config.root_gid),
&hooks,
)
.map_err(EngineError::from)?;
let ep3 = store.epoch();
store
.dir_lookup_epoch(&ep3, 1, ENGINE_DIR_NAME)
.map_err(EngineError::from)?
.ok_or_else(|| {
EngineError::new(ErrorCode::Internal, "namespace create vanished")
})?
.ino
}
}
};
if !opts.read_only {
{
let ep = store.epoch();
let entries = store
.read_dir_epoch(&ep, dir_ino)
.map_err(EngineError::from)?;
let stale: Vec<Vec<u8>> = entries
.into_iter()
.filter(|(name, e)| e.d_type == directory::dt::DT_REG && is_tmp_name(name))
.map(|(name, _)| name)
.collect();
drop(ep);
for name in stale {
store
.epoch_unlink(dir_ino, &name, false, &hooks)
.map_err(EngineError::from)?;
}
}
}
Ok(Self {
inner: Arc::new(EngineInner {
store: Mutex::new(Some(Arc::new(store))),
closed: AtomicBool::new(false),
ops_in_flight: AtomicU64::new(0),
drain_cv: Condvar::new(),
}),
})
}
fn acquire_store(&self) -> Result<Arc<Store>, EngineError> {
if self.inner.closed.load(Ordering::Acquire) {
return Err(EngineError::new(ErrorCode::Closed, "engine is closed"));
}
self.inner.ops_in_flight.fetch_add(1, Ordering::AcqRel);
let guard = self.inner.store.lock().unwrap_or_else(|p| p.into_inner());
let store = match guard.as_ref() {
Some(s) => Arc::clone(s),
None => {
drop(guard);
self.finish_op();
return Err(EngineError::new(ErrorCode::Closed, "engine is closed"));
}
};
drop(guard);
Ok(store)
}
fn finish_op(&self) {
if self.inner.ops_in_flight.fetch_sub(1, Ordering::AcqRel) == 1 {
self.inner.drain_cv.notify_all();
}
}
fn engine_dir_ino(&self, store: &Store) -> Result<u64, EngineError> {
let ep = store.epoch();
let entry = store
.dir_lookup_epoch(&ep, 1, ENGINE_DIR_NAME)
.map_err(EngineError::from)?
.ok_or_else(|| {
EngineError::new(
ErrorCode::CorruptStore,
"engine namespace directory is missing from the store root",
)
})?;
if entry.d_type != directory::dt::DT_DIR {
return Err(EngineError::new(
ErrorCode::CorruptStore,
"engine namespace name is occupied by a non-directory entry",
));
}
Ok(entry.ino)
}
fn blob_name(id: &BlobId) -> Vec<u8> {
id.to_string().into_bytes()
}
fn tmp_name(id: &BlobId, counter: u64) -> Vec<u8> {
let mut v = Vec::with_capacity(64 + 24);
v.extend_from_slice(TMP_PREFIX);
v.extend_from_slice(&id.to_string().into_bytes());
v.extend_from_slice(TMP_SUFFIX);
v.extend_from_slice(format!("{counter:016x}").as_bytes());
v
}
pub fn put_blob(&self, bytes: &[u8]) -> Result<BlobId, EngineError> {
self.put_blob_with(bytes, Durability::Ack)
}
pub fn put_blob_with(
&self,
bytes: &[u8],
durability: Durability,
) -> Result<BlobId, EngineError> {
let store = self.acquire_store()?;
let _op = OpGuard::new(self);
let id = BlobId::from(ChunkId::of(bytes));
crate::perf::trace::span!(
"engine.put_blob",
op = "put_blob",
len = bytes.len() as u64,
id = trace_id(&id),
durable = matches!(durability, Durability::Durable)
);
let hooks = CrashHooks::none();
let dir_ino = self.engine_dir_ino(&store)?;
let final_name = Self::blob_name(&id);
{
let ep = store.epoch();
match store
.dir_lookup_epoch(&ep, dir_ino, &final_name)
.map_err(EngineError::from)?
{
Some(e) if e.d_type == directory::dt::DT_REG => return Ok(id),
Some(_) => {}
None => {}
}
}
let tmp = Self::tmp_name(&id, self.tmp_counter());
let file_ino = store
.epoch_create(
dir_ino,
&tmp,
NewEntry::file(0o600, store.config().root_uid, store.config().root_gid),
&hooks,
)
.map_err(EngineError::from)?;
if !bytes.is_empty() {
store
.epoch_write(
file_ino,
0,
bytes,
OptimizeOptions::default(),
store.config().foreground,
&hooks,
)
.map_err(EngineError::from)?;
}
match store
.epoch_rename(dir_ino, &tmp, dir_ino, &final_name, &hooks)
.map_err(EngineError::from)
{
Ok(_) => {}
Err(e) if e.code == ErrorCode::CorruptStore && e.message.contains("already exists") => {
let ep = store.epoch();
let exists = store
.dir_lookup_epoch(&ep, dir_ino, &final_name)
.map_err(EngineError::from)?
.map(|e| e.d_type == directory::dt::DT_REG)
.unwrap_or(false);
if !exists {
return Err(EngineError::new(
ErrorCode::Internal,
"concurrent put lost its blob and the winner vanished",
));
}
}
Err(e) => return Err(e),
}
if durability == Durability::Durable {
store
.durability_barrier(&hooks)
.map_err(EngineError::from)?;
}
Ok(id)
}
fn tmp_counter(&self) -> u64 {
static COUNTER: AtomicU64 = AtomicU64::new(0);
COUNTER.fetch_add(1, Ordering::Relaxed)
}
pub fn get_blob(&self, id: BlobId) -> Result<Vec<u8>, EngineError> {
let bytes = self.read_blob_range(id, 0, usize::MAX)?;
let got = ChunkId::of(&bytes);
if BlobId::from(got) != id {
return Err(EngineError::new(
ErrorCode::CorruptStore,
format!("blob {id} materialized to bytes hashing to {got} (content mismatch)"),
));
}
Ok(bytes)
}
pub fn read_blob_range(
&self,
id: BlobId,
offset: u64,
len: usize,
) -> Result<Vec<u8>, EngineError> {
let store = self.acquire_store()?;
let _op = OpGuard::new(self);
crate::perf::trace::span!(
"engine.read_blob_range",
op = "read_range",
id = trace_id(&id),
offset = offset,
len = len as u64
);
let dir_ino = self.engine_dir_ino(&store)?;
let name = Self::blob_name(&id);
let ep = store.epoch();
let entry = store
.dir_lookup_epoch(&ep, dir_ino, &name)
.map_err(EngineError::from)?
.ok_or_else(|| EngineError::new(ErrorCode::NotFound, format!("blob {id} not found")))?;
if entry.d_type != directory::dt::DT_REG {
return Err(EngineError::new(
ErrorCode::CorruptStore,
format!("blob name {id} is occupied by a non-file entry"),
));
}
let inode = store
.get_inode_epoch(&ep, entry.ino)
.map_err(EngineError::from)?
.ok_or_else(|| {
EngineError::new(ErrorCode::CorruptStore, format!("blob {id} inode missing"))
})?;
let size = inode.size;
if offset >= size {
return Ok(Vec::new());
}
let avail = size.saturating_sub(offset).min(len as u64);
store
.read_file_epoch(&ep, entry.ino, offset, avail)
.map_err(EngineError::from)
}
pub fn contains(&self, id: BlobId) -> Result<bool, EngineError> {
let store = self.acquire_store()?;
let _op = OpGuard::new(self);
let dir_ino = self.engine_dir_ino(&store)?;
let ep = store.epoch();
Ok(store
.dir_lookup_epoch(&ep, dir_ino, &Self::blob_name(&id))
.map_err(EngineError::from)?
.map(|e| e.d_type == directory::dt::DT_REG)
.unwrap_or(false))
}
pub fn sync(&self) -> Result<(), EngineError> {
let store = self.acquire_store()?;
let _op = OpGuard::new(self);
crate::perf::trace::span!("engine.sync", op = "sync");
store
.durability_barrier(&CrashHooks::none())
.map_err(EngineError::from)
}
pub fn compact(&self) -> Result<CompactionReport, EngineError> {
let store = self.acquire_store()?;
let _op = OpGuard::new(self);
crate::perf::trace::span!("engine.compact", op = "compact");
let hooks = CrashHooks::none();
let dir_ino = self.engine_dir_ino(&store)?;
let stale: Vec<Vec<u8>> = {
let ep = store.epoch();
store
.read_dir_epoch(&ep, dir_ino)
.map_err(EngineError::from)?
.into_iter()
.filter(|(name, e)| e.d_type == directory::dt::DT_REG && is_tmp_name(name))
.map(|(name, _)| name)
.collect()
};
for name in stale {
store
.epoch_unlink(dir_ino, &name, false, &hooks)
.map_err(EngineError::from)?;
}
store
.ensure_epoch_flushed(&hooks)
.map_err(EngineError::from)?;
let unreachable_before =
crate::store::gc::unreachable_bytes(&store).map_err(EngineError::from)?;
let reclaimed =
crate::store::gc::compact_full(&store, &hooks).map_err(EngineError::from)?;
let unreachable_after =
crate::store::gc::unreachable_bytes(&store).map_err(EngineError::from)?;
let physical_used_after_bytes = store.physical_used();
let live_bytes_after = crate::store::physical::physical_report(&store)
.map(|r| r.live_bytes)
.unwrap_or(0);
Ok(CompactionReport {
unreachable_before_bytes: unreachable_before,
reclaimed_bytes: reclaimed,
unreachable_after_bytes: unreachable_after,
physical_used_after_bytes,
live_bytes_after,
})
}
pub fn metrics(&self) -> Result<EngineMetrics, EngineError> {
let store = self.acquire_store()?;
let _op = OpGuard::new(self);
collect_engine_metrics(&store)
}
pub fn close(&self) -> Result<(), EngineError> {
if self.inner.closed.swap(true, Ordering::AcqRel) {
return Err(EngineError::new(ErrorCode::Closed, "engine is closed"));
}
let mut guard = self.inner.store.lock().unwrap_or_else(|p| p.into_inner());
while self.inner.ops_in_flight.load(Ordering::Acquire) != 0 {
guard = self
.inner
.drain_cv
.wait(guard)
.unwrap_or_else(|p| p.into_inner());
}
*guard = None;
Ok(())
}
pub fn is_closed(&self) -> bool {
self.inner.closed.load(Ordering::Acquire)
}
}
impl Drop for Engine {
fn drop(&mut self) {
let _ = self.close();
}
}
struct OpGuard<'a> {
engine: &'a Engine,
}
impl<'a> OpGuard<'a> {
fn new(engine: &'a Engine) -> Self {
Self { engine }
}
}
impl Drop for OpGuard<'_> {
fn drop(&mut self) {
self.engine.finish_op();
}
}
fn is_tmp_name(name: &[u8]) -> bool {
name.starts_with(TMP_PREFIX) && contains_suffix(name, TMP_SUFFIX)
}
fn contains_suffix(name: &[u8], suffix: &[u8]) -> bool {
if name.len() < suffix.len() {
return false;
}
&name[name.len() - suffix.len()..] == suffix
}
fn trace_id(id: &BlobId) -> String {
id.to_string()[..8].to_string()
}
pub fn collect_engine_metrics(store: &Store) -> Result<EngineMetrics, EngineError> {
let root = store.current_root();
let bits = store.feature_bits();
let stats = store.stats();
let capacity = store.physical_capacity();
let used = store.physical_used();
let phys = crate::store::physical::physical_report(store).ok();
let dsfb = store.dsfb_stats();
let (live_b, dead_b, hidden_b, unindexed_b, torn_b, pad_b, fmt_b, unexp_b) = match &phys {
Some(r) => (
r.live_bytes,
r.dead_indexed_bytes,
r.index_hidden_bytes,
r.unindexed_bytes,
r.torn_bytes,
r.zero_padding_bytes,
r.format_overhead_bytes,
r.unexplained(),
),
None => (0, 0, 0, 0, 0, 0, 0, 0),
};
let phases: Vec<PhaseMetrics> = store
.perf()
.snapshot()
.into_iter()
.map(|row| PhaseMetrics {
phase: row.phase.to_string(),
count: row.count,
total_ms: row.total_ms,
p50_us: row.p50_us,
p95_us: row.p95_us,
p99_us: row.p99_us,
})
.collect();
let blob_count = {
let ep = store.epoch();
match store
.dir_lookup_epoch(&ep, 1, ENGINE_DIR_NAME)
.ok()
.flatten()
{
Some(e) if e.d_type == directory::dt::DT_DIR => store
.read_dir_epoch(&ep, e.ino)
.map(|entries| {
entries
.iter()
.filter(|(_, e)| e.d_type == directory::dt::DT_REG)
.count() as u64
})
.unwrap_or(0),
_ => 0,
}
};
Ok(EngineMetrics {
schema_version: 1,
format: FormatInfo {
format_major: root.format_major,
format_minor: root.format_minor,
compat: bits.compat,
ro_compat: bits.ro_compat,
incompat: bits.incompat,
io_backend: store.config().io_backend.name().to_string(),
},
accounting: AccountingMetrics {
logical_bytes: store.logical_bytes().unwrap_or(0),
reachable_bytes: stats.reachable_bytes,
physical_used_bytes: used,
physical_capacity_bytes: capacity,
physical_free_bytes: capacity.saturating_sub(used),
object_count: store.object_index().len() as u64,
data_record_count: stats.data_record_count,
blob_count,
},
physical: PhysicalMetrics {
live_bytes: live_b,
dead_indexed_bytes: dead_b,
index_hidden_bytes: hidden_b,
unindexed_bytes: unindexed_b,
torn_bytes: torn_b,
zero_padding_bytes: pad_b,
format_overhead_bytes: fmt_b,
unexplained_bytes: unexp_b,
},
gc: GcMetrics {
unreachable_bytes: stats.unreachable_bytes,
},
dsfb: DsfbMetrics {
tracked_chunks: dsfb.tracked_chunks,
steps: dsfb.steps,
drift_events: dsfb.drift_events,
slew_events: dsfb.slew_events,
narrowed_searches: dsfb.narrowed_searches,
candidates_evaluated: store.candidates_evaluated(),
},
cache: CacheMetrics {
model_cache_hits: store.model_cache_hits(),
model_cache_misses: store.model_cache_misses(),
},
write_path_phases: phases,
})
}
#[cfg(test)]
mod tests {
use super::*;
use tempfile::TempDir;
fn store_dir() -> TempDir {
tempfile::Builder::new()
.prefix("entropyfs-engine-test")
.tempdir()
.expect("tempdir")
}
#[test]
fn blob_identity_semantics() {
let id = BlobId::from(ChunkId::of(b"hello"));
let hex = id.to_string();
assert_eq!(hex.len(), 64);
let back = BlobId::from_hex(&hex).unwrap();
assert_eq!(back, id);
assert!(BlobId::from_hex("zz").is_none());
assert!(BlobId::from_hex(&"0".repeat(64)).is_some());
}
#[test]
fn put_get_roundtrip_and_dedup() {
let dir = store_dir();
let engine = Engine::create(dir.path(), &EngineOpenOptions::default()).unwrap();
let data = b"the quick brown fox jumps over the lazy dog".repeat(100);
let id = engine.put_blob(&data).unwrap();
let id2 = engine.put_blob(&data).unwrap();
assert_eq!(id, id2);
assert!(engine.contains(id).unwrap());
let out = engine.get_blob(id).unwrap();
assert_eq!(out, data);
let other = engine.put_blob(b"different bytes").unwrap();
assert_ne!(other, id);
assert_eq!(engine.get_blob(other).unwrap(), b"different bytes");
assert!(
!engine
.contains(BlobId::from(ChunkId::of(b"never put")))
.unwrap()
);
engine.close().unwrap();
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn empty_blob() {
let dir = store_dir();
let engine = Engine::create(dir.path(), &EngineOpenOptions::default()).unwrap();
let id = engine.put_blob(b"").unwrap();
assert!(engine.contains(id).unwrap());
assert_eq!(engine.get_blob(id).unwrap(), b"");
let range = engine.read_blob_range(id, 0, 10).unwrap();
assert!(range.is_empty());
engine.close().unwrap();
}
#[test]
fn large_blob_and_range_reads() {
let dir = store_dir();
let engine = Engine::create(dir.path(), &EngineOpenOptions::default()).unwrap();
let mut data = Vec::with_capacity(2 << 20);
for i in 0..(2 << 20) {
data.push((i % 251) as u8);
}
let id = engine.put_blob(&data).unwrap();
assert_eq!(engine.get_blob(id).unwrap(), data);
assert_eq!(engine.read_blob_range(id, 0, 4).unwrap(), &data[..4]);
let mid = 70_000u64;
assert_eq!(
engine.read_blob_range(id, mid, 100).unwrap(),
&data[mid as usize..mid as usize + 100]
);
assert!(
engine
.read_blob_range(id, data.len() as u64, 10)
.unwrap()
.is_empty()
);
let clipped = engine
.read_blob_range(id, data.len() as u64 - 8, 100)
.unwrap();
assert_eq!(clipped.len(), 8);
engine.close().unwrap();
}
#[test]
fn persistence_across_reopen() {
let dir = store_dir();
let id;
{
let engine = Engine::create(dir.path(), &EngineOpenOptions::default()).unwrap();
let data = b"persistent bytes across a reopen".repeat(40);
id = engine.put_blob(&data).unwrap();
engine.sync().unwrap();
engine.close().unwrap();
}
{
let engine = Engine::open(dir.path(), &EngineOpenOptions::default()).unwrap();
assert!(engine.contains(id).unwrap());
let data = b"persistent bytes across a reopen".repeat(40);
assert_eq!(engine.get_blob(id).unwrap(), data);
let id2 = engine.put_blob(b"post-reopen").unwrap();
assert_eq!(engine.get_blob(id2).unwrap(), b"post-reopen");
engine.close().unwrap();
}
}
#[test]
fn open_legacy_store_creates_namespace_lazily() {
let dir = store_dir();
let config = StoreConfig::default();
let store = Store::create(dir.path(), &config, [9u8; 16]).unwrap();
drop(store);
let engine = Engine::open(dir.path(), &EngineOpenOptions::default()).unwrap();
let id = engine.put_blob(b"legacy store adoption").unwrap();
assert_eq!(engine.get_blob(id).unwrap(), b"legacy store adoption");
engine.close().unwrap();
}
#[test]
fn missing_blob_is_not_found() {
let dir = store_dir();
let engine = Engine::create(dir.path(), &EngineOpenOptions::default()).unwrap();
let err = engine
.get_blob(BlobId::from(ChunkId::of(b"nope")))
.unwrap_err();
assert_eq!(err.code(), ErrorCode::NotFound);
engine.close().unwrap();
}
#[test]
fn concurrent_puts_and_gets() {
let dir = store_dir();
let engine = Arc::new(Engine::create(dir.path(), &EngineOpenOptions::default()).unwrap());
let mut handles = Vec::new();
for t in 0..8 {
let engine = Arc::clone(&engine);
handles.push(std::thread::spawn(move || {
for i in 0..8 {
let payload = format!("thread-{t}-blob-{i}:{}", "x".repeat(10_000 + t * 1000));
let id = engine.put_blob(payload.as_bytes()).unwrap();
let out = engine.get_blob(id).unwrap();
assert_eq!(out, payload.as_bytes());
}
}));
}
for h in handles {
h.join().unwrap();
}
let m = engine.metrics().unwrap();
assert_eq!(m.accounting.blob_count, 64);
engine.close().unwrap();
}
#[test]
fn close_drains_and_rejects() {
let dir = store_dir();
let engine = Engine::create(dir.path(), &EngineOpenOptions::default()).unwrap();
let id = engine.put_blob(b"before close").unwrap();
engine.close().unwrap();
assert_eq!(engine.get_blob(id).unwrap_err().code(), ErrorCode::Closed);
assert!(engine.is_closed());
assert_eq!(engine.close().unwrap_err().code(), ErrorCode::Closed);
}
#[test]
fn compact_reclaims_and_preserves() {
let dir = store_dir();
let engine = Engine::create(dir.path(), &EngineOpenOptions::default()).unwrap();
let mut ids = Vec::new();
for i in 0..64 {
let payload = format!("blob-{i}:{}", "y".repeat(50_000 + i));
ids.push(engine.put_blob(payload.as_bytes()).unwrap());
}
for i in 0..32 {
let payload = format!("blob-{i}-v2:{}", "y".repeat(50_000 + i));
ids.push(engine.put_blob(payload.as_bytes()).unwrap());
}
engine.sync().unwrap();
let before = engine.metrics().unwrap();
let report = engine.compact().unwrap();
let after = engine.metrics().unwrap();
assert!(
report.unreachable_after_bytes < report.unreachable_before_bytes,
"compaction must reclaim: before {} after {}",
report.unreachable_before_bytes,
report.unreachable_after_bytes
);
assert!(
report.unreachable_after_bytes <= 4096,
"residual must be bounded (superseded root record), got {}",
report.unreachable_after_bytes
);
assert!(after.accounting.physical_used_bytes <= before.accounting.physical_used_bytes);
for id in &ids[32..] {
let out = engine.get_blob(*id).unwrap();
assert!(!out.is_empty());
}
engine.close().unwrap();
}
}