Skip to main content

vole_document/store/
embedded.rs

1//! The `EmbeddedStore` reference backend: a local content-addressed directory.
2//!
3//! Storage is raw (no compression) so its accounting is transparent and
4//! backend-independent. Layout:
5//!
6//! ```text
7//! STORE_ROOT/
8//!   STORE                    # magic + backend tag + format version (10 bytes)
9//!   objects/<aa>/<bb>/<64-hex-id>       # object payload, raw
10//!   objects/<aa>/<bb>/<64-hex-id>.tmp-<pid>-<nanos>   # in-flight put
11//! ```
12//!
13//! ## Atomic put invariant
14//!
15//! The final name exists ⇒ the object is complete and byte-exact. A put writes
16//! a temp sibling, `write_all` + `sync_all`, then renames onto the final name
17//! (atomic on one filesystem). A crash leaves at most a `.tmp-*` file, swept on
18//! the next [`EmbeddedStore::open`]; only `.tmp-*` names are swept, foreign
19//! entries are left untouched.
20
21use std::fs;
22use std::io::{Read, Seek, SeekFrom, Write};
23use std::path::{Path, PathBuf};
24use std::time::{SystemTime, UNIX_EPOCH};
25
26use crate::error::{Error, Result};
27use crate::store::{Id, ObjectStore, StoreStats};
28
29/// Store marker magic (`VOLDST` + 0x1A text-EOF + NUL).
30pub const STORE_MAGIC: [u8; 8] = *b"VOLDST\x1A\x00";
31/// Backend tag for the raw embedded store.
32pub const STORE_BACKEND_EMBEDDED: u8 = 0x01;
33/// Current on-disk store format version.
34pub const STORE_FORMAT_VERSION: u8 = 1;
35
36/// A content-addressed directory store.
37#[derive(Debug, Clone)]
38pub struct EmbeddedStore {
39    root: PathBuf,
40}
41
42impl EmbeddedStore {
43    /// Open (creating if necessary) a store rooted at `root`.
44    ///
45    /// Creates the directory tree, writes/validates the `STORE` marker, and
46    /// sweeps any leftover `.tmp-*` files from a crashed put.
47    pub fn open(root: impl AsRef<Path>) -> Result<Self> {
48        let root = root.as_ref().to_path_buf();
49        fs::create_dir_all(root.join("objects"))?;
50        let store = EmbeddedStore { root };
51        store.ensure_marker()?;
52        store.sweep_tmp()?;
53        Ok(store)
54    }
55
56    /// The store root directory.
57    pub fn root(&self) -> &Path {
58        &self.root
59    }
60
61    /// The raw stored length of `id`, or a typed error if absent.
62    pub fn len(&self, id: &Id) -> Result<u64> {
63        match fs::metadata(self.object_path(id)) {
64            Ok(m) if m.is_file() => Ok(m.len()),
65            Ok(_) => Err(Error::missing_external_object(format!(
66                "object {id} is not a file"
67            ))),
68            Err(e) if e.kind() == std::io::ErrorKind::NotFound => Err(
69                Error::missing_external_object(format!("object {id} is not present in the store")),
70            ),
71            Err(e) => Err(Error::from(e)),
72        }
73    }
74
75    /// Whether the store holds no objects.
76    pub fn is_empty(&self) -> Result<bool> {
77        Ok(self.stats()?.object_count == 0)
78    }
79
80    /// Store statistics over every present object.
81    pub fn stats(&self) -> Result<StoreStats> {
82        let listed = self.list()?;
83        let total_bytes: u64 = listed.iter().map(|(_, len)| *len).sum();
84        Ok(StoreStats {
85            object_count: listed.len() as u64,
86            total_bytes,
87            stored_bytes: total_bytes,
88        })
89    }
90
91    fn marker_path(&self) -> PathBuf {
92        self.root.join("STORE")
93    }
94
95    fn objects_dir(&self) -> PathBuf {
96        self.root.join("objects")
97    }
98
99    fn object_path(&self, id: &Id) -> PathBuf {
100        let hex = id.to_hex();
101        self.objects_dir()
102            .join(&hex[0..2])
103            .join(&hex[2..4])
104            .join(&hex)
105    }
106
107    fn ensure_marker(&self) -> Result<()> {
108        let path = self.marker_path();
109        if let Ok(bytes) = fs::read(&path) {
110            if bytes.len() != 10 || bytes[0..8] != STORE_MAGIC {
111                return Err(Error::invalid_container(
112                    "store marker has a bad magic or length",
113                ));
114            }
115            if bytes[8] != STORE_BACKEND_EMBEDDED {
116                return Err(Error::unsupported_feature(format!(
117                    "store backend tag {:#04x} is not the embedded backend",
118                    bytes[8]
119                )));
120            }
121            if bytes[9] > STORE_FORMAT_VERSION {
122                return Err(Error::unsupported_version(format!(
123                    "store format version {} is newer than this build ({STORE_FORMAT_VERSION})",
124                    bytes[9]
125                )));
126            }
127            return Ok(());
128        }
129        let mut marker = Vec::with_capacity(10);
130        marker.extend_from_slice(&STORE_MAGIC);
131        marker.push(STORE_BACKEND_EMBEDDED);
132        marker.push(STORE_FORMAT_VERSION);
133        write_atomic(&path, &marker)
134    }
135
136    /// Remove leftover `.tmp-*` files from a crashed put. Foreign entries (any
137    /// name not containing the `.tmp-` infix) are left untouched.
138    fn sweep_tmp(&self) -> Result<()> {
139        let objects = self.objects_dir();
140        for aa in read_dir_dirs(&objects)? {
141            for bb in read_dir_dirs(&aa)? {
142                for entry in fs::read_dir(&bb)? {
143                    let entry = entry?;
144                    if !entry.file_type()?.is_file() {
145                        continue;
146                    }
147                    let name = entry.file_name();
148                    let name = name.to_string_lossy();
149                    if name.contains(".tmp-") {
150                        let _ = fs::remove_file(entry.path());
151                    }
152                }
153            }
154        }
155        Ok(())
156    }
157}
158
159impl ObjectStore for EmbeddedStore {
160    fn put(&mut self, bytes: &[u8]) -> Result<Id> {
161        let id = Id::of(bytes);
162        let path = self.object_path(&id);
163        if path.is_file() {
164            return Ok(id);
165        }
166        let parent = path
167            .parent()
168            .ok_or_else(|| Error::internal_invariant("object path has no parent"))?;
169        fs::create_dir_all(parent)?;
170
171        let hex = id.to_hex();
172        let nanos = SystemTime::now()
173            .duration_since(UNIX_EPOCH)
174            .map(|d| d.as_nanos())
175            .unwrap_or(0);
176        let tmp = parent.join(format!(".{hex}.tmp-{}-{nanos}", std::process::id()));
177
178        let write_result = (|| -> Result<()> {
179            let mut f = fs::File::create(&tmp)?;
180            f.write_all(bytes)?;
181            f.sync_all()?;
182            Ok(())
183        })();
184        if let Err(e) = write_result {
185            let _ = fs::remove_file(&tmp);
186            return Err(e);
187        }
188        if let Err(e) = fs::rename(&tmp, &path) {
189            let _ = fs::remove_file(&tmp);
190            return Err(Error::from(e));
191        }
192        // Best-effort directory durability (Linux).
193        if let Ok(dir) = fs::File::open(parent) {
194            let _ = dir.sync_all();
195        }
196        Ok(id)
197    }
198
199    fn get(&self, id: &Id) -> Result<Vec<u8>> {
200        let bytes = match fs::read(self.object_path(id)) {
201            Ok(b) => b,
202            Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
203                return Err(Error::missing_external_object(format!(
204                    "object {id} is not present in the store"
205                )));
206            }
207            Err(e) => return Err(Error::from(e)),
208        };
209        if Id::of(&bytes) != *id {
210            return Err(Error::integrity_mismatch(format!(
211                "stored object {id} does not hash to its content id"
212            )));
213        }
214        Ok(bytes)
215    }
216
217    fn get_range(&self, id: &Id, offset: u64, len: u64) -> Result<Vec<u8>> {
218        let mut f = match fs::File::open(self.object_path(id)) {
219            Ok(f) => f,
220            Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
221                return Err(Error::missing_external_object(format!(
222                    "object {id} is not present in the store"
223                )));
224            }
225            Err(e) => return Err(Error::from(e)),
226        };
227        let file_len = f.metadata()?.len();
228        let end = offset
229            .checked_add(len)
230            .ok_or_else(|| Error::integrity_mismatch("object range offset+len overflow"))?;
231        if end > file_len {
232            return Err(Error::integrity_mismatch(format!(
233                "object {id} has {file_len} bytes; range [{offset}, {end}) is out of bounds"
234            )));
235        }
236        f.seek(SeekFrom::Start(offset))?;
237        let mut out = vec![0u8; len as usize];
238        f.read_exact(&mut out)?;
239        Ok(out)
240    }
241
242    fn contains(&self, id: &Id) -> Result<bool> {
243        Ok(self.object_path(id).is_file())
244    }
245
246    fn list(&self) -> Result<Vec<(Id, u64)>> {
247        let mut out: Vec<(Id, u64)> = Vec::new();
248        let objects = self.objects_dir();
249        for aa in read_dir_dirs(&objects)? {
250            for bb in read_dir_dirs(&aa)? {
251                for entry in fs::read_dir(&bb)? {
252                    let entry = entry?;
253                    if !entry.file_type()?.is_file() {
254                        continue;
255                    }
256                    let name = entry.file_name();
257                    let Some(name) = name.to_str() else {
258                        continue;
259                    };
260                    let Ok(id) = Id::from_hex(name) else {
261                        continue;
262                    };
263                    out.push((id, entry.metadata()?.len()));
264                }
265            }
266        }
267        out.sort_unstable_by_key(|(id, _)| *id);
268        Ok(out)
269    }
270
271    fn remove(&self, id: &Id) -> Result<u64> {
272        let path = self.object_path(id);
273        let len = match fs::metadata(&path) {
274            Ok(m) if m.is_file() => m.len(),
275            Ok(_) => return Ok(0),
276            Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(0),
277            Err(e) => return Err(Error::from(e)),
278        };
279        fs::remove_file(&path)?;
280        // Best-effort empty-shard cleanup; a non-empty shard is left in place.
281        if let Some(bb) = path.parent() {
282            let _ = fs::remove_dir(bb);
283            if let Some(aa) = bb.parent() {
284                let _ = fs::remove_dir(aa);
285            }
286        }
287        Ok(len)
288    }
289}
290
291/// Write `bytes` to `path` atomically: temp sibling, `sync_all`, rename, fsync
292/// the parent directory. Used for the small `STORE` marker.
293fn write_atomic(path: &Path, bytes: &[u8]) -> Result<()> {
294    let dir = match path.parent() {
295        Some(p) if !p.as_os_str().is_empty() => p.to_path_buf(),
296        _ => PathBuf::from("."),
297    };
298    let name = path
299        .file_name()
300        .map(|s| s.to_string_lossy().into_owned())
301        .unwrap_or_else(|| "STORE".to_string());
302    let nanos = SystemTime::now()
303        .duration_since(UNIX_EPOCH)
304        .map(|d| d.as_nanos())
305        .unwrap_or(0);
306    let tmp = dir.join(format!(".{name}.tmp-{}-{nanos}", std::process::id()));
307
308    {
309        let mut f = fs::File::create(&tmp)?;
310        f.write_all(bytes)?;
311        f.sync_all()?;
312    }
313    if let Err(e) = fs::rename(&tmp, path) {
314        let _ = fs::remove_file(&tmp);
315        return Err(Error::from(e));
316    }
317    if let Ok(d) = fs::File::open(&dir) {
318        let _ = d.sync_all();
319    }
320    Ok(())
321}
322
323/// Subdirectories of `dir` (missing `dir` yields an empty list), sorted for
324/// deterministic traversal.
325fn read_dir_dirs(dir: &Path) -> Result<Vec<PathBuf>> {
326    let mut out: Vec<PathBuf> = Vec::new();
327    let entries = match fs::read_dir(dir) {
328        Ok(e) => e,
329        Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(out),
330        Err(e) => return Err(Error::from(e)),
331    };
332    for entry in entries {
333        let entry = entry?;
334        if entry.file_type()?.is_dir() {
335            out.push(entry.path());
336        }
337    }
338    out.sort_unstable();
339    Ok(out)
340}