vole_document/store/
embedded.rs1use 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
29pub const STORE_MAGIC: [u8; 8] = *b"VOLDST\x1A\x00";
31pub const STORE_BACKEND_EMBEDDED: u8 = 0x01;
33pub const STORE_FORMAT_VERSION: u8 = 1;
35
36#[derive(Debug, Clone)]
38pub struct EmbeddedStore {
39 root: PathBuf,
40}
41
42impl EmbeddedStore {
43 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 pub fn root(&self) -> &Path {
58 &self.root
59 }
60
61 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 pub fn is_empty(&self) -> Result<bool> {
77 Ok(self.stats()?.object_count == 0)
78 }
79
80 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 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 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 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
291fn 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
323fn 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}