mkit_core/store.rs
1//! Local content-addressed object store.
2//!
3//! Layout (under the [`crate::layout::RepoLayout`] common dir passed to
4//! [`ObjectStore::open`] / [`ObjectStore::init`]):
5//!
6//! ```text
7//! .mkit/
8//! objects/
9//! <2-hex>/<62-hex> # raw canonical object bytes, BLAKE3-named
10//! ```
11//!
12//! Writes are atomic: bytes are first written to a sibling temp file
13//! (`<name>.tmp.<pid>.<rand>`), made durable, then renamed into place.
14//! A crash mid-write leaves only the temp file behind and never
15//! produces a visible object that fails the read-time hash check.
16//!
17//! Durability comes in three shapes (see [`crate::batch`] for the full
18//! contract): [`ObjectStore::write`] flushes per object
19//! ([`SyncPolicy::PerObject`]); [`ObjectStore::batch`] defers visibility
20//! and amortises durability to **one** full flush per batch — the write
21//! path used by every multi-object command (add, commit, pack unpack);
22//! and [`ObjectStore::bulk_writer`] fsyncs each object's contents before
23//! the rename and batches the dir fsyncs at commit (git import). In all
24//! three an object is never visible before its bytes are durable, so a
25//! ref or index written after the write/commit returns can never
26//! reference a non-durable object.
27//!
28//! Reads always verify integrity by recomputing BLAKE3 over the bytes
29//! and comparing against the requested hash; mismatch returns
30//! [`StoreError::HashMismatch`]. The one opt-in exception is
31//! [`ObjectStore::read_unverified`] (and the [`DisplaySource`] adapter
32//! built on it), reserved for display-only rendering where a corrupt
33//! object should surface as a bad render, never as durable state — see
34//! its doc for the full policy (#625).
35//!
36//! See `docs/specs/SPEC-OBJECTS.md` §10 for the path-layout rule.
37
38use std::fs::{self, File};
39use std::io::{self, Read, Write};
40use std::path::{Path, PathBuf};
41use std::process;
42use std::sync::Arc;
43use std::sync::atomic::{AtomicU64, Ordering};
44
45use tempfile::NamedTempFile;
46
47use crate::batch::{RealSyncer, Syncer};
48pub use crate::batch::{SyncPolicy, WriteBatch};
49use crate::hash::{self, Hash, object_path, to_hex};
50use crate::layout::RepoLayout;
51use crate::object::{MkitError, Object, object_id_from_bytes, verified_id_and_object};
52use crate::serialize;
53
54mod memory;
55mod source;
56pub use memory::MemorySource;
57pub use source::{DisplaySource, EphemeralSink, ObjectSource};
58
59/// Top-level repository directory name.
60pub const MKIT_DIR: &str = ".mkit";
61/// Subdirectory under `.mkit/` that holds raw object files.
62pub const OBJECTS_DIR: &str = "objects";
63/// File under `.mkit/` declaring the object-addressing format. Written at
64/// [`ObjectStore::init`] and required (and matched) at [`ObjectStore::open`]
65/// so an older flat-hash-addressed repository is rejected loudly rather
66/// than silently mis-read once merkle addressing is in effect.
67pub const FORMAT_FILE: &str = "format";
68/// The only supported object-addressing format value (see
69/// `docs/specs/SPEC-MERKLE-OBJECTS.md`): Tree/ChunkedBlob keyed by BMT root.
70pub const FORMAT_VALUE: &str = "bmt-v1";
71/// Hard cap on raw object size, enforced on both [`ObjectStore::write`]
72/// and [`ObjectStore::read`].
73pub const MAX_RAW_OBJECT_SIZE: usize = 1024 * 1024 * 1024; // 1 GiB
74
75/// Hard cap on tree-walk recursion depth.
76///
77/// Every core tree-walker (`index::from_tree`, `ops::diff`, `ops::merge`,
78/// `ops::restore`) recurses one native stack frame per directory level.
79/// Content-addressing prevents true cycles (a child's hash can never equal
80/// its parent's), but a crafted untrusted repo can still nest a few thousand
81/// single-entry trees in a tiny pack and overflow the stack — a denial of
82/// service reachable from `clone`/`checkout`/`merge`/`diff`. Walkers thread a
83/// depth counter and abort with a typed `TreeTooDeep` error once
84/// `depth > MAX_TREE_DEPTH`.
85///
86/// 128 is far beyond any legitimately-deep source tree (Git's own working
87/// trees rarely exceed a couple dozen levels) while remaining comfortably
88/// within a default 8 MiB thread stack. Mirrors the `MAX_REF_DEPTH` cap on
89/// the ref-tree walker in `refs.rs`.
90pub const MAX_TREE_DEPTH: usize = 128;
91
92/// Errors raised by the [`ObjectStore`] surface. Distinct from
93/// [`MkitError`] so callers can pattern-match on filesystem failures
94/// without losing the structured-decode-error variants.
95#[derive(Debug, thiserror::Error)]
96pub enum StoreError {
97 #[error("path is not an mkit repository (missing .mkit/objects)")]
98 NotAMkitRepository,
99 #[error(".mkit already exists in this directory")]
100 AlreadyInitialized,
101 #[error(
102 "repository object-addressing format is {found:?}, expected \"{}\" — this repository predates merkle object addressing and is not readable by this mkit (pre-1.0: no migration). See docs/specs/SPEC-MERKLE-OBJECTS.md.",
103 FORMAT_VALUE
104 )]
105 IncompatibleRepoFormat { found: Option<String> },
106 #[error("object {0} not found")]
107 ObjectNotFound(String),
108 #[error("object exceeds {} byte cap", MAX_RAW_OBJECT_SIZE)]
109 ObjectTooLarge,
110 #[error("tree nesting exceeds {} levels", MAX_TREE_DEPTH)]
111 TreeTooDeep,
112 #[error("on-disk bytes hash to {actual}, expected {expected}")]
113 HashMismatch { expected: String, actual: String },
114 #[error(transparent)]
115 Io(#[from] io::Error),
116 #[error(transparent)]
117 Decode(#[from] MkitError),
118 /// [`crate::transfer::plan_pack_with`]'s `encode_deltas` callback
119 /// returned a different number of results than the candidate slice it
120 /// was given — a contract violation by the caller (see that
121 /// function's docs), never a normal runtime condition. Reported as a
122 /// typed error rather than a panic: the built-in caller
123 /// (`mkit-cli`'s push path) can never trigger this since rayon's
124 /// `IndexedParallelIterator::collect` always preserves length, but a
125 /// future or third-party `encode_deltas` with a bug should fail the
126 /// `push` cleanly rather than crash the process.
127 #[error("encode_deltas callback returned {actual} results for {expected} candidates")]
128 DeltaBatchLengthMismatch { expected: usize, actual: usize },
129}
130
131/// Deferred-fsync writer returned by [`ObjectStore::bulk_writer`].
132/// See that method for the crash-safety contract.
133#[derive(Debug)]
134pub struct BulkWriter<'a> {
135 store: &'a ObjectStore,
136 dirs: std::collections::HashSet<PathBuf>,
137 files: std::collections::HashSet<PathBuf>,
138}
139
140impl BulkWriter<'_> {
141 /// Write one object durable-before-visible: contents are fsynced
142 /// before the rename publishes the object, and only the shard-dir
143 /// fsync (rename durability) is deferred to commit. This upholds the
144 /// store's global invariant (an object is never visible before its
145 /// bytes are durable), so another process's `contains()` dedup can
146 /// never reference a half-written object. An existing path is
147 /// byte-verified: matched objects are left in place with a refreshed
148 /// `mtime` (MKIT-55; still fsynced at commit, see below), torn ones
149 /// are rewritten.
150 ///
151 /// # Panics
152 /// Never in practice: object paths always have a 2-hex shard
153 /// parent by construction.
154 pub fn write(&mut self, bytes: &[u8]) -> StoreResult<Hash> {
155 if bytes.len() > MAX_RAW_OBJECT_SIZE {
156 return Err(StoreError::ObjectTooLarge);
157 }
158 let h = object_id_from_bytes(bytes);
159 let final_path = self.store.path_for(&h);
160 // VERIFY-skip an existing object: content addressing makes
161 // byte-equality the exact durability test — a good file (from
162 // a previously committed session, possibly referenced by
163 // native history) must NOT be replaced with an unsynced inode
164 // a power loss could tear, while a torn file from a crashed
165 // session fails the comparison and is healed by the rewrite.
166 let shard_dir = final_path
167 .parent()
168 .expect("object path always has a 2-hex parent");
169 // A match whose mtime cannot be refreshed (MKIT-55, see
170 // `refresh_mtime`) is rewritten like a torn file.
171 if let Ok(existing) = fs::read(&final_path)
172 && existing == bytes
173 && refresh_mtime(&final_path).is_ok()
174 {
175 // The bytes may have matched out of the PAGE CACHE of a
176 // crashed session's unsynced write — byte equality is not a
177 // durability test, so this pre-existing file still needs a
178 // content fsync at commit.
179 self.dirs.insert(shard_dir.to_path_buf());
180 self.files.insert(final_path);
181 return Ok(h);
182 }
183 // `self.dirs` already records every shard this session has
184 // touched — a shard present there was `create_dir_all`'d (or
185 // observed via the byte-match branch above, which also implies
186 // it exists) by an earlier `write` call in this same session,
187 // so a repeat `create_dir_all` here is a guaranteed-redundant
188 // mkdir + is_dir stat (std's fallback path on `AlreadyExists`).
189 // At most 256 shards exist, so memoizing saves ~2 syscalls per
190 // object on large bulk ingests (git-import) — the same
191 // `created_shards` memoization `WriteBatch::write_prehashed`
192 // already applies to the batched-add path.
193 if !self.dirs.contains(shard_dir) {
194 fs::create_dir_all(shard_dir)?;
195 }
196 // Content fsync happens BEFORE the rename, so the object is
197 // durable the instant it becomes visible. Only the rename's
198 // durability (the shard dir fsync) is deferred to commit.
199 crate::atomic::write_content_synced(&final_path, bytes)?;
200 self.dirs.insert(shard_dir.to_path_buf());
201 Ok(h)
202 }
203
204 /// Make the session durable: fsync the CONTENTS of any pre-existing
205 /// objects this batch re-used (newly written objects were already
206 /// content-fsynced before their rename in [`write`](Self::write)),
207 /// then every touched shard directory (renames become durable).
208 pub fn commit(self) -> StoreResult<()> {
209 for file in &self.files {
210 fs::OpenOptions::new().write(true).open(file)?.sync_all()?;
211 }
212 for dir in &self.dirs {
213 crate::atomic::sync_dir(dir)?;
214 }
215 Ok(())
216 }
217}
218
219/// Result alias used throughout this module.
220pub type StoreResult<T> = Result<T, StoreError>;
221
222// Tiny per-process counter for unique temp-file names. We use this
223// instead of pulling in `rand` because the temp name only needs to be
224// unique within the process; the atomic `rename` enforces global
225// correctness even if two processes collide on a name.
226static TEMP_SEQ: AtomicU64 = AtomicU64::new(0);
227
228/// Local content-addressed object store backed by the filesystem.
229#[derive(Debug, Clone)]
230pub struct ObjectStore {
231 /// Absolute path to `<root>/.mkit/objects`.
232 objects_root: PathBuf,
233 /// Flush/rename primitive seam. Production code always uses
234 /// [`RealSyncer`]; unit tests inject a recording double to assert
235 /// flush ordering and counts (the O(1)-flushes-per-batch contract).
236 syncer: Arc<dyn Syncer>,
237 /// Policy handed out by [`Self::batch`]. Defaults to
238 /// [`SyncPolicy::Batch`]; the CLI maps the repo config key
239 /// `durability.objects = per-object` onto it for deployments that
240 /// want the strict historical schedule.
241 default_policy: SyncPolicy,
242 /// Test-only counter of physical reads (`read_raw` calls) this store
243 /// has performed. There is no read-side equivalent of `syncer` to
244 /// inject a recording double against, so this is a plain counter
245 /// instead — used to prove callers (e.g. the pack reader's delta-base
246 /// resolution, issue #643) hit the store at most once per distinct
247 /// object rather than once per reference to it.
248 #[cfg(test)]
249 read_calls: Arc<AtomicU64>,
250}
251
252impl ObjectStore {
253 /// Open the existing repository described by `layout`. Returns
254 /// [`StoreError::NotAMkitRepository`] if the layout's `objects/`
255 /// directory does not exist. The object store is common-dir
256 /// (shared) state — see [`crate::layout`].
257 pub fn open(layout: &RepoLayout) -> StoreResult<Self> {
258 let objects_root = layout.objects_dir();
259 if !objects_root.is_dir() {
260 return Err(StoreError::NotAMkitRepository);
261 }
262 // Reject a repository that does not declare merkle object
263 // addressing — a pre-merkle repo would mis-read every Tree /
264 // ChunkedBlob (and thus every Commit) under the new id scheme.
265 match fs::read_to_string(layout.format_file()) {
266 Ok(s) if s.trim() == FORMAT_VALUE => {}
267 Ok(s) => {
268 return Err(StoreError::IncompatibleRepoFormat {
269 found: Some(s.trim().to_owned()),
270 });
271 }
272 Err(e) if e.kind() == io::ErrorKind::NotFound => {
273 return Err(StoreError::IncompatibleRepoFormat { found: None });
274 }
275 Err(e) => return Err(StoreError::Io(e)),
276 }
277 Ok(Self {
278 objects_root,
279 syncer: Arc::new(RealSyncer),
280 default_policy: SyncPolicy::Batch,
281 #[cfg(test)]
282 read_calls: Arc::new(AtomicU64::new(0)),
283 })
284 }
285
286 /// Initialise a fresh common dir for the repository described by
287 /// `layout`. Returns [`StoreError::AlreadyInitialized`] if the
288 /// common dir already exists.
289 pub fn init(layout: &RepoLayout) -> StoreResult<Self> {
290 if layout.common_dir().exists() {
291 return Err(StoreError::AlreadyInitialized);
292 }
293 let objects_root = layout.objects_dir();
294 fs::create_dir_all(&objects_root)?;
295 // Declare the object-addressing format so a future open by an
296 // incompatible (older) mkit fails loudly (SPEC-MERKLE-OBJECTS §7).
297 fs::write(layout.format_file(), format!("{FORMAT_VALUE}\n").as_bytes())?;
298 Ok(Self {
299 objects_root,
300 syncer: Arc::new(RealSyncer),
301 default_policy: SyncPolicy::Batch,
302 #[cfg(test)]
303 read_calls: Arc::new(AtomicU64::new(0)),
304 })
305 }
306
307 /// Select the [`SyncPolicy`] that [`Self::batch`] hands out. The
308 /// CLI wires the repo config key `durability.objects` here so
309 /// deployments on filesystems where they prefer the strict
310 /// per-object schedule can opt into it (SPEC-OBJECTS §10.1).
311 pub fn set_sync_policy(&mut self, policy: SyncPolicy) {
312 self.default_policy = policy;
313 }
314
315 /// Replace the flush/rename primitives. Test-only seam — see the
316 /// `syncer` field. Not exposed publicly so the production sync
317 /// strategy cannot be silently weakened by downstream code.
318 #[cfg(test)]
319 pub(crate) fn set_syncer(&mut self, syncer: Arc<dyn Syncer>) {
320 self.syncer = syncer;
321 }
322
323 /// The active flush/rename primitives, shared with [`WriteBatch`].
324 pub(crate) fn syncer(&self) -> &Arc<dyn Syncer> {
325 &self.syncer
326 }
327
328 /// Test-only: number of physical filesystem reads this store has
329 /// performed (via [`Self::read`] or [`Self::read_unverified`]) since
330 /// creation. See the `read_calls` field doc for why this is a plain
331 /// counter rather than an injectable double like [`Self::set_syncer`].
332 #[cfg(test)]
333 pub(crate) fn read_call_count(&self) -> u64 {
334 self.read_calls.load(Ordering::Relaxed)
335 }
336
337 /// Start a batched write with the store's configured policy
338 /// (default [`SyncPolicy::Batch`]): objects staged by the batch
339 /// become durable and visible together at [`WriteBatch::commit`],
340 /// with O(1) full flushes per batch instead of per object.
341 #[must_use]
342 pub fn batch(&self) -> WriteBatch<'_> {
343 self.batch_with_policy(self.default_policy)
344 }
345
346 /// Start a batched write with an explicit [`SyncPolicy`].
347 #[must_use]
348 pub fn batch_with_policy(&self, policy: SyncPolicy) -> WriteBatch<'_> {
349 WriteBatch::new(self, policy)
350 }
351
352 /// Returns `true` when `root` contains a `.mkit/objects` directory.
353 #[must_use]
354 pub fn is_repo_root(root: &Path) -> bool {
355 root.join(MKIT_DIR).join(OBJECTS_DIR).is_dir()
356 }
357
358 /// Absolute path to the `objects/` directory.
359 #[must_use]
360 pub fn objects_root(&self) -> &Path {
361 &self.objects_root
362 }
363
364 /// Compute the on-disk path for `hash`, joined under `objects/`.
365 /// `pub(crate)` so [`WriteBatch`] shares the single layout rule.
366 pub(crate) fn path_for(&self, h: &Hash) -> PathBuf {
367 let p = object_path(h);
368 // Both halves are ASCII hex by construction in `object_path`.
369 let dir = std::str::from_utf8(&p.dir).expect("ascii hex");
370 let file = std::str::from_utf8(&p.file).expect("ascii hex");
371 self.objects_root.join(dir).join(file)
372 }
373
374 /// Returns `true` when the object `h` is present in the store. Does
375 /// **not** verify integrity — use [`Self::read`] for that.
376 #[must_use]
377 pub fn contains(&self, h: &Hash) -> bool {
378 self.path_for(h).is_file()
379 }
380
381 /// Write `bytes` to the store, returning their BLAKE3 hash. Atomic:
382 /// writes to a sibling temp file, `fsync`s, then renames into place.
383 /// Idempotent — re-writing the same bytes is a no-op (the temp file
384 /// is unlinked on the early-return path).
385 ///
386 /// # Panics
387 ///
388 /// Panics only if the internal hash-to-path mapping produces a path
389 /// without a parent directory, which is impossible by construction.
390 pub fn write(&self, bytes: &[u8]) -> StoreResult<Hash> {
391 if bytes.len() > MAX_RAW_OBJECT_SIZE {
392 return Err(StoreError::ObjectTooLarge);
393 }
394 let h = object_id_from_bytes(bytes);
395 let final_path = self.path_for(&h);
396 if final_path.exists() {
397 // Dedup hit: the object is visible, but if another process
398 // renamed it and has not yet flushed the dirent, it may not
399 // be durable — and our caller is about to reference it.
400 // The `mtime` is left alone: callers hold locks gc also
401 // takes, so gc's grace window does not apply to them (see
402 // `refresh_mtime`).
403 // Flush its shard dir before returning (SPEC-OBJECTS §10.1
404 // dedup rule, mirroring WriteBatch's touched_shards).
405 self.syncer()
406 .dir_sync(final_path.parent().expect("object path has parent"))?;
407 return Ok(h);
408 }
409 let shard_dir = final_path
410 .parent()
411 .expect("object path always has a 2-hex parent");
412 fs::create_dir_all(shard_dir)?;
413 write_atomic(&final_path, bytes, &**self.syncer())?;
414 Ok(h)
415 }
416
417 /// Begin a bulk-write session: each new object is written
418 /// durable-before-visible (contents fsynced, then renamed), and
419 /// [`BulkWriter::commit`] batches the directory fsyncs (rename
420 /// durability) plus a content fsync of any re-used pre-existing
421 /// objects once at the end, instead of fsyncing a dir per write.
422 ///
423 /// Because content is durable before the rename, the session upholds
424 /// the store's global invariant even under concurrent readers: an
425 /// object another process can `contains()`-dedup against is always
426 /// backed by durable bytes.
427 ///
428 /// Crash-safety contract (deliberately weaker than [`Self::write`]
429 /// only for the rename/dirent half, for callers whose whole
430 /// operation is idempotent — e.g. the deterministic git-import,
431 /// which re-runs from a retained source mirror): after a crash
432 /// BEFORE `commit`, a just-renamed object's dirent may be lost (the
433 /// shard dir is not yet fsynced), so objects may be missing — never
434 /// torn, since contents were fsynced first. Existing paths are
435 /// VERIFIED (byte compare)
436 /// rather than blindly rewritten or blindly trusted: a matching
437 /// file is left in place with only its `mtime` refreshed (it may be
438 /// durable and referenced by native history — replacing it with an
439 /// unsynced inode would put it at risk), a torn one is healed by
440 /// rewrite, and reads always
441 /// BLAKE3-verify. Callers MUST gate bulk sessions behind their
442 /// own crash marker and re-run on detection.
443 #[must_use]
444 pub fn bulk_writer(&self) -> BulkWriter<'_> {
445 BulkWriter {
446 store: self,
447 dirs: std::collections::HashSet::new(),
448 files: std::collections::HashSet::new(),
449 }
450 }
451
452 /// Open `h`'s on-disk file, cap its size at [`MAX_RAW_OBJECT_SIZE`],
453 /// and read it fully into a pre-sized buffer — the shared body of
454 /// [`Self::read`] and [`Self::read_unverified`]; neither hashes nor
455 /// verifies, that's each caller's job.
456 fn read_raw(&self, h: &Hash) -> StoreResult<Vec<u8>> {
457 #[cfg(test)]
458 self.read_calls.fetch_add(1, Ordering::Relaxed);
459 let path = self.path_for(h);
460 let mut file = File::open(&path).map_err(|e| match e.kind() {
461 io::ErrorKind::NotFound => StoreError::ObjectNotFound(to_hex(h)),
462 _ => StoreError::Io(e),
463 })?;
464 let meta = file.metadata()?;
465 let size = meta.len();
466 if u128::from(size) > MAX_RAW_OBJECT_SIZE as u128 {
467 return Err(StoreError::ObjectTooLarge);
468 }
469 // We've already bounded `size` by `MAX_RAW_OBJECT_SIZE` (1 GiB),
470 // which fits in `usize` on every platform we support (32-bit
471 // included). Pre-size the buffer to avoid the doubling
472 // re-allocations of `read_to_end` for large objects.
473 let cap = usize::try_from(size).map_err(|_| StoreError::ObjectTooLarge)?;
474 let mut bytes = Vec::with_capacity(cap);
475 file.read_to_end(&mut bytes)?;
476 Ok(bytes)
477 }
478
479 /// Reads raw bytes for a verifier that must classify a hash mismatch
480 /// under the requested id. Unlike [`Self::read`], this deliberately
481 /// leaves identity verification to the caller; it is crate-private so
482 /// only verification code can use the raw bytes as data.
483 ///
484 /// # Errors
485 ///
486 /// Returns the same store errors as [`Self::read_raw`].
487 pub(crate) fn read_raw_for_verification(&self, h: &Hash) -> StoreResult<Vec<u8>> {
488 self.read_raw(h)
489 }
490
491 /// Read raw bytes for `h`. Verifies that BLAKE3 of the on-disk
492 /// bytes equals `h` and returns [`StoreError::HashMismatch`] on
493 /// failure (the bytes are still discarded so callers cannot
494 /// accidentally use corrupt data).
495 pub fn read(&self, h: &Hash) -> StoreResult<Vec<u8>> {
496 let bytes = self.read_raw(h)?;
497 check_hash(h, &object_id_from_bytes(&bytes))?;
498 Ok(bytes)
499 }
500
501 /// Read a verified pack base using the caller's budgeted allocator.
502 /// The same open handle supplies the size and bytes; growth is rejected
503 /// rather than letting `read_to_end` allocate outside that budget.
504 pub(crate) fn read_with_allocator<E: From<StoreError>>(
505 &self,
506 h: &Hash,
507 allocate: impl FnOnce(usize) -> Result<Vec<u8>, E>,
508 ) -> Result<Vec<u8>, E> {
509 #[cfg(test)]
510 self.read_calls.fetch_add(1, Ordering::Relaxed);
511 let mut file = File::open(self.path_for(h)).map_err(|e| {
512 E::from(match e.kind() {
513 io::ErrorKind::NotFound => StoreError::ObjectNotFound(to_hex(h)),
514 _ => StoreError::Io(e),
515 })
516 })?;
517 let size = file.metadata().map_err(StoreError::Io)?.len();
518 if size > MAX_RAW_OBJECT_SIZE as u64 {
519 return Err(StoreError::ObjectTooLarge.into());
520 }
521 let len = usize::try_from(size).map_err(|_| StoreError::ObjectTooLarge)?;
522 let mut bytes = allocate(len)?;
523 bytes.resize(len, 0);
524 file.read_exact(&mut bytes).map_err(StoreError::Io)?;
525 let mut extra = [0];
526 if file.read(&mut extra).map_err(StoreError::Io)? != 0 {
527 return Err(StoreError::ObjectTooLarge.into());
528 }
529 check_hash(h, &object_id_from_bytes(&bytes))?;
530 Ok(bytes)
531 }
532
533 /// Read raw bytes for `h` WITHOUT the BLAKE3 integrity check that
534 /// [`Self::read`] performs — a [`StoreError::HashMismatch`] can
535 /// therefore never come out of this path.
536 ///
537 /// # Policy — display-only paths ONLY (#625)
538 /// This exists for hot read-only rendering (`diff`/`show`/commit
539 /// summaries formatting a diffstat or patch preview). It is the same
540 /// "cheap read" tradeoff [`Self::object_type`] already makes for shape
541 /// checks (see its doc and [`Self::verify_object_type`], which is the
542 /// *verifying* twin used by tree-publication paths) — extended here to
543 /// the full object body.
544 ///
545 /// NEVER use this for publication (commit/merge/rebase tree writes),
546 /// dedup, fetch/apply, or any path whose output is consumed as data
547 /// rather than displayed and discarded. mkit never feeds unverified
548 /// bytes into state that mkit itself writes, or into output that
549 /// mkit's own tooling round-trips (e.g. `format-patch` piped into
550 /// `git am`). On corruption, callers of `read` get a loud, typed error
551 /// before anything downstream can act on bad bytes; callers of
552 /// `read_unverified` get whatever a decoder makes of the corrupt bytes
553 /// — a decode error, or in the worst case garbage in the *rendered
554 /// output* — but never durable state, because display paths write
555 /// nothing back to the store.
556 ///
557 /// Prefer wrapping the source with [`DisplaySource`] at the render
558 /// call site over calling this directly; that keeps the verify/
559 /// no-verify choice at the boundary instead of threading a flag
560 /// through render code.
561 pub fn read_unverified(&self, h: &Hash) -> StoreResult<Vec<u8>> {
562 self.read_raw(h)
563 }
564
565 /// The object's type tag, from its 6-byte prologue — without
566 /// reading or hash-verifying the body. Backs cheap shape checks
567 /// (e.g. "is this staged hash blob-like?") that previously paid a
568 /// full read+BLAKE3 of every staged blob per status/commit; the
569 /// real read path still integrity-verifies at use time.
570 ///
571 /// # Errors
572 /// [`StoreError::ObjectNotFound`] if absent; [`StoreError::Decode`]
573 /// for a short file, bad magic/version, or unknown tag.
574 pub fn object_type(&self, h: &Hash) -> StoreResult<crate::object::ObjectType> {
575 let path = self.path_for(h);
576 let mut file = File::open(&path).map_err(|e| match e.kind() {
577 io::ErrorKind::NotFound => StoreError::ObjectNotFound(to_hex(h)),
578 _ => StoreError::Io(e),
579 })?;
580 let mut prologue = [0u8; 6];
581 file.read_exact(&mut prologue)
582 .map_err(|_| StoreError::Decode(MkitError::EmptyData))?;
583 if prologue[1..5] != crate::object::MAGIC {
584 return Err(StoreError::Decode(MkitError::InvalidMagic));
585 }
586 if prologue[5] != crate::object::SCHEMA_VERSION {
587 return Err(StoreError::Decode(MkitError::UnsupportedObjectVersion));
588 }
589 crate::object::ObjectType::from_u8(prologue[0]).map_err(StoreError::Decode)
590 }
591
592 /// Convenience: read raw bytes and decode into a typed [`Object`].
593 ///
594 /// A merkelized type ([`Object::Tree`] / [`Object::ChunkedBlob`]) is
595 /// already decoded once by `verified_id_and_object` to compute its
596 /// BMT-root id for the integrity check below — reused directly here
597 /// instead of a second `deserialize` of the same bytes (the two
598 /// decodes this function used to pay for every tree/chunked-blob
599 /// read, once inside [`Self::read`]'s verification and once here).
600 /// Every other type was only ever decoded once (verification there is
601 /// a plain `BLAKE3(bytes)`, no decode), so this path is unchanged for
602 /// them.
603 pub fn read_object(&self, h: &Hash) -> StoreResult<Object> {
604 let bytes = self.read_raw(h)?;
605 let (actual, decoded) = verified_id_and_object(&bytes);
606 check_hash(h, &actual)?;
607 match decoded {
608 Some(obj) => Ok(obj),
609 None => Ok(serialize::deserialize(&bytes)?),
610 }
611 }
612
613 /// Hash-verifying variant of [`object_type`](Self::object_type): reads
614 /// the whole object and confirms its content hashes to `h` before
615 /// returning the declared type. This is the integrity guard for
616 /// tree-publication paths (commit, merge, rebase, …) — `object_type`
617 /// alone reads only the 6-byte prologue, so a staged object corrupted
618 /// after `add` would otherwise be published into a durable tree and
619 /// only fail at later read time. Use `object_type` on hot read-only
620 /// paths (status/diff snapshots) where nothing durable is published.
621 pub fn verify_object_type(&self, h: &Hash) -> StoreResult<crate::object::ObjectType> {
622 let bytes = self.read(h)?; // read() re-hashes and rejects on mismatch
623 if bytes.len() < 6 {
624 return Err(StoreError::Decode(MkitError::EmptyData));
625 }
626 if bytes[1..5] != crate::object::MAGIC {
627 return Err(StoreError::Decode(MkitError::InvalidMagic));
628 }
629 if bytes[5] != crate::object::SCHEMA_VERSION {
630 return Err(StoreError::Decode(MkitError::UnsupportedObjectVersion));
631 }
632 crate::object::ObjectType::from_u8(bytes[0]).map_err(StoreError::Decode)
633 }
634
635 /// Enumerate every object hash currently in the store by walking
636 /// `objects/<2-hex>/<62-hex>`. Entries whose names are not the
637 /// expected hex shape (stray files, atomic-write temp files, unknown
638 /// dirs) are skipped — they are not objects. Used by `gc` to find
639 /// prune candidates.
640 ///
641 /// # Errors
642 /// [`StoreError::Io`] if a directory cannot be read (gc must then
643 /// fail closed rather than prune against a partial enumeration).
644 pub fn iter_object_hashes(&self) -> StoreResult<Vec<Hash>> {
645 let mut out = Vec::new();
646 let shards = match fs::read_dir(&self.objects_root) {
647 Ok(rd) => rd,
648 Err(e) if e.kind() == io::ErrorKind::NotFound => return Ok(out),
649 Err(e) => return Err(StoreError::Io(e)),
650 };
651 for shard in shards {
652 let shard = shard?;
653 if !shard.file_type()?.is_dir() {
654 continue;
655 }
656 let dir_name = shard.file_name();
657 let Some(dir) = dir_name.to_str() else {
658 continue;
659 };
660 if dir.len() != 2 || !dir.bytes().all(|b| b.is_ascii_hexdigit()) {
661 continue;
662 }
663 for entry in fs::read_dir(shard.path())? {
664 let entry = entry?;
665 if !entry.file_type()?.is_file() {
666 continue;
667 }
668 let file_name = entry.file_name();
669 let Some(file) = file_name.to_str() else {
670 continue;
671 };
672 if file.len() != 62 {
673 continue;
674 }
675 // Reassemble the 64-hex id and parse it; skips non-objects.
676 if let Ok(h) = hash::from_hex(&format!("{dir}{file}")) {
677 out.push(h);
678 }
679 }
680 }
681 Ok(out)
682 }
683
684 /// Filesystem metadata for object `h` (size + mtime), for gc's
685 /// grace-window check.
686 ///
687 /// # Errors
688 /// [`StoreError::ObjectNotFound`] if absent, else [`StoreError::Io`].
689 pub fn object_metadata(&self, h: &Hash) -> StoreResult<fs::Metadata> {
690 fs::metadata(self.path_for(h)).map_err(|e| match e.kind() {
691 io::ErrorKind::NotFound => StoreError::ObjectNotFound(to_hex(h)),
692 _ => StoreError::Io(e),
693 })
694 }
695
696 /// Delete object `h` from the store. Idempotent: a missing object is
697 /// not an error (it may have been pruned already).
698 ///
699 /// # Errors
700 /// [`StoreError::Io`] on a filesystem failure other than not-found.
701 pub fn remove_object(&self, h: &Hash) -> StoreResult<()> {
702 match fs::remove_file(self.path_for(h)) {
703 Ok(()) => Ok(()),
704 Err(e) if e.kind() == io::ErrorKind::NotFound => Ok(()),
705 Err(e) => Err(StoreError::Io(e)),
706 }
707 }
708}
709
710/// Shared by [`ObjectStore::read`] and [`ObjectStore::read_object`]: both
711/// compute an `actual` id from the on-disk bytes their own way (plain
712/// `object_id_from_bytes`, or the decode-reusing [`verified_id_and_object`])
713/// and then need the exact same requested-vs-actual check and
714/// [`StoreError::HashMismatch`] shape — kept in one place so the two can't
715/// drift apart.
716fn check_hash(expected: &Hash, actual: &Hash) -> StoreResult<()> {
717 if actual != expected {
718 return Err(StoreError::HashMismatch {
719 expected: to_hex(expected),
720 actual: to_hex(actual),
721 });
722 }
723 Ok(())
724}
725
726/// Create a uniquely-named sibling temp file in `parent` for the object
727/// file `file_name`. The `.{name}.tmp.{pid}.{seq}` shape is what
728/// [`ObjectStore::iter_object_hashes`] and the stale-temp tests rely on
729/// to skip non-objects.
730pub(crate) fn temp_file_in(parent: &Path, file_name: &str) -> io::Result<NamedTempFile> {
731 let pid = process::id();
732 let seq = TEMP_SEQ.fetch_add(1, Ordering::Relaxed);
733 let tmp_name = format!(".{file_name}.tmp.{pid}.{seq}");
734 NamedTempFile::with_prefix_in(tmp_name, parent)
735}
736
737/// Atomically write `bytes` to `final_path` with per-object durability
738/// ([`SyncPolicy::PerObject`]): temp file, full flush, rename into
739/// place, parent-dir flush. On Unix, `rename(2)` is atomic with respect
740/// to concurrent readers and replaces the destination. On Windows,
741/// [`NamedTempFile::persist`] uses `MOVEFILE_REPLACE_EXISTING`
742/// semantics so the replace-existing path works there too.
743///
744/// After a successful rename we flush the parent directory on Unix to
745/// commit the dirent update — without this, the rename can survive a
746/// power loss only in the page cache and the file appears missing on
747/// reboot.
748///
749/// All flush/rename primitives go through `syncer` so unit tests can
750/// assert ordering; [`RealSyncer`] preserves the historical behaviour
751/// exactly.
752fn write_atomic(final_path: &Path, bytes: &[u8], syncer: &dyn Syncer) -> io::Result<()> {
753 let parent = final_path.parent().expect("write_atomic: path has parent");
754 let file_name = final_path
755 .file_name()
756 .expect("write_atomic: path has file name")
757 .to_string_lossy();
758
759 let mut tmp = temp_file_in(parent, &file_name)?;
760 tmp.as_file_mut().write_all(bytes)?;
761 syncer.full(tmp.as_file(), tmp.path())?;
762
763 // NamedTempFile::persist uses a cross-platform atomic replace:
764 // rename(2) on Unix, MoveFileExW with MOVEFILE_REPLACE_EXISTING on Windows.
765 syncer.rename(tmp.into_temp_path(), final_path)?;
766
767 syncer.dir_sync(parent)?;
768 Ok(())
769}
770
771/// Write target shared by [`ObjectStore`] (per-object durability) and
772/// [`WriteBatch`] (batched durability), so ingest code can be written
773/// once against either sink.
774pub trait ObjectSink {
775 /// Store `bytes` as one object, returning its BLAKE3 hash.
776 fn put(&self, bytes: &[u8]) -> StoreResult<Hash>;
777 /// Store the concatenation of `parts` as one object without the
778 /// caller having to materialise the concatenated buffer.
779 fn put_parts(&self, parts: &[&[u8]]) -> StoreResult<Hash>;
780 /// True when the object is already present (or staged) in this sink.
781 fn has(&self, h: &Hash) -> bool;
782}
783
784impl ObjectSink for ObjectStore {
785 fn put(&self, bytes: &[u8]) -> StoreResult<Hash> {
786 self.write(bytes)
787 }
788
789 fn put_parts(&self, parts: &[&[u8]]) -> StoreResult<Hash> {
790 // Per-object path is the legacy/rare sink; hot ingest paths use
791 // WriteBatch, whose put_parts is copy-free.
792 let total: usize = parts.iter().map(|p| p.len()).sum();
793 let mut bytes = Vec::with_capacity(total);
794 for p in parts {
795 bytes.extend_from_slice(p);
796 }
797 self.write(&bytes)
798 }
799
800 fn has(&self, h: &Hash) -> bool {
801 self.contains(h)
802 }
803}
804
805/// Reset an existing object file's `mtime` to now, on a dedup hit.
806///
807/// gc's grace window is measured from the object file's on-disk
808/// `mtime` (SPEC-GC "Concurrent writers and the grace window"). A
809/// writer outside gc's lock set that dedups against an OLD unreachable
810/// object and then publishes a ref over it would otherwise lose the
811/// object to a gc run in between (MKIT-55). The only such writer into
812/// this store is git import, through [`BulkWriter`], so only
813/// [`BulkWriter::write`] calls this. [`ObjectStore::write`] and
814/// [`crate::batch::WriteBatch`] callers hold `worktree.lock`, which gc
815/// also takes, and refreshing there would dirty the inode right before
816/// their per-hit shard-dir fsync, turning a ~10 µs dedup hit into a
817/// journal commit (~3 ms on APFS). The caller treats an `Err` as "not
818/// refreshed" and falls back to rewriting the object, which gives it a
819/// fresh inode and `mtime`; it never reports a dedup hit whose `mtime`
820/// was not refreshed.
821///
822/// - Permissions: on Unix, setting an explicit time needs ownership,
823/// not write access, so a read-only fd works on a read-only object
824/// file; a non-owner gets `EPERM` and so rewrites. Windows needs a
825/// write-capable handle for `SetFileTime`.
826/// - Durability: the new `mtime` is not fsynced. It only has to hold
827/// between this dedup hit and the writer's ref publish, in the same
828/// boot, where gc reads it from the in-memory inode. A crash ends the
829/// writer: a ref it already published durably makes the object a gc
830/// root, and an unpublished one never needed it, so a reverted
831/// `mtime` after a crash cannot lose a reachable object.
832/// - Cost: one `open` + `futimens` (plus `close`) per dedup hit.
833///
834/// This does not close gc's own window between reading an object's
835/// `mtime` and unlinking it (`ops::gc::run_gc`): a refresh that lands
836/// inside it still loses the object.
837pub(crate) fn refresh_mtime(path: &Path) -> io::Result<()> {
838 #[cfg(test)]
839 if FAIL_REFRESH.get() {
840 return Err(io::ErrorKind::NotFound.into());
841 }
842 let mut opts = fs::OpenOptions::new();
843 if cfg!(unix) {
844 opts.read(true);
845 } else {
846 opts.write(true);
847 }
848 opts.open(path)?.set_modified(std::time::SystemTime::now())
849}
850
851#[cfg(test)]
852thread_local! {
853 /// Test seam: make [`refresh_mtime`] fail on this thread, as a
854 /// concurrent gc unlink (`NotFound`) would, to exercise the
855 /// callers' rewrite fallback.
856 static FAIL_REFRESH: std::cell::Cell<bool> = const { std::cell::Cell::new(false) };
857}
858
859/// On Unix, fsync the directory holding the just-renamed file so the
860/// dirent update is durable. No-op on non-Unix (Windows does not expose
861/// a stable directory-fsync primitive via `std::fs`).
862#[cfg(unix)]
863pub(crate) fn sync_parent_dir(parent: &Path) -> io::Result<()> {
864 match File::open(parent) {
865 Ok(dir) => dir.sync_all(),
866 // If the dir disappeared under us (race with external cleanup),
867 // the durability invariant is moot — propagate silently.
868 Err(e) if e.kind() == io::ErrorKind::NotFound => Ok(()),
869 Err(e) => Err(e),
870 }
871}
872
873#[cfg(not(unix))]
874#[allow(clippy::unnecessary_wraps)]
875pub(crate) fn sync_parent_dir(_parent: &Path) -> io::Result<()> {
876 Ok(())
877}
878
879#[cfg(test)]
880mod tests {
881 use super::*;
882 use crate::object::Blob;
883 use std::fs::OpenOptions;
884 use std::io::Seek;
885 use tempfile::TempDir;
886
887 fn fresh_store() -> (TempDir, ObjectStore) {
888 let dir = TempDir::new().expect("tempdir");
889 let store = ObjectStore::init(&RepoLayout::single(dir.path())).expect("init");
890 (dir, store)
891 }
892
893 #[test]
894 fn init_creates_layout() {
895 let dir = TempDir::new().unwrap();
896 let _ = ObjectStore::init(&RepoLayout::single(dir.path())).unwrap();
897 assert!(dir.path().join(MKIT_DIR).is_dir());
898 assert!(dir.path().join(MKIT_DIR).join(OBJECTS_DIR).is_dir());
899 }
900
901 #[test]
902 fn init_rejects_already_initialized() {
903 let dir = TempDir::new().unwrap();
904 let layout = RepoLayout::single(dir.path());
905 ObjectStore::init(&layout).unwrap();
906 let err = ObjectStore::init(&layout).unwrap_err();
907 assert!(matches!(err, StoreError::AlreadyInitialized));
908 }
909
910 #[test]
911 fn open_rejects_non_repo() {
912 let dir = TempDir::new().unwrap();
913 let err = ObjectStore::open(&RepoLayout::single(dir.path())).unwrap_err();
914 assert!(matches!(err, StoreError::NotAMkitRepository));
915 }
916
917 #[test]
918 fn init_writes_format_marker_and_open_accepts_it() {
919 let dir = TempDir::new().unwrap();
920 let layout = RepoLayout::single(dir.path());
921 ObjectStore::init(&layout).unwrap();
922 let marker = fs::read_to_string(dir.path().join(MKIT_DIR).join(FORMAT_FILE)).unwrap();
923 assert_eq!(marker.trim(), FORMAT_VALUE);
924 // A freshly-init'd repo opens cleanly.
925 ObjectStore::open(&layout).unwrap();
926 }
927
928 #[test]
929 fn open_rejects_repo_missing_or_wrong_format_marker() {
930 // A pre-merkle repo (objects dir present, no format marker) must be
931 // rejected loudly rather than silently mis-read.
932 let dir = TempDir::new().unwrap();
933 let layout = RepoLayout::single(dir.path());
934 ObjectStore::init(&layout).unwrap();
935 let marker_path = dir.path().join(MKIT_DIR).join(FORMAT_FILE);
936
937 fs::remove_file(&marker_path).unwrap();
938 assert!(matches!(
939 ObjectStore::open(&layout),
940 Err(StoreError::IncompatibleRepoFormat { found: None })
941 ));
942
943 // A repo declaring some other (e.g. future) format is also rejected.
944 fs::write(&marker_path, b"flat-v0\n").unwrap();
945 assert!(matches!(
946 ObjectStore::open(&layout),
947 Err(StoreError::IncompatibleRepoFormat { found: Some(_) })
948 ));
949 }
950
951 #[test]
952 fn is_repo_root_predicate() {
953 let dir = TempDir::new().unwrap();
954 assert!(!ObjectStore::is_repo_root(dir.path()));
955 ObjectStore::init(&RepoLayout::single(dir.path())).unwrap();
956 assert!(ObjectStore::is_repo_root(dir.path()));
957 }
958
959 #[test]
960 fn write_then_read_roundtrip() {
961 let (_dir, store) = fresh_store();
962 let bytes = b"hello world".to_vec();
963 let h = store.write(&bytes).unwrap();
964 assert!(store.contains(&h));
965 let got = store.read(&h).unwrap();
966 assert_eq!(got, bytes);
967 }
968
969 #[test]
970 fn read_object_deserialises() {
971 let (_dir, store) = fresh_store();
972 let obj = Object::Blob(Blob {
973 data: b"object bytes".to_vec(),
974 });
975 let bytes = serialize::serialize(&obj).unwrap();
976 let h = store.write(&bytes).unwrap();
977 let parsed = store.read_object(&h).unwrap();
978 assert_eq!(parsed, obj);
979 }
980
981 /// `read_object` on a merkelized type (the path `verified_id_and_object`
982 /// added a decoded-object fast path for) round-trips correctly — the
983 /// non-merkle `read_object_deserialises` case above can't exercise the
984 /// `Some(obj)` branch at all.
985 #[test]
986 fn read_object_deserialises_tree() {
987 use crate::object::{EntryMode, Tree, TreeEntry};
988 let (_dir, store) = fresh_store();
989 let obj = Object::Tree(Tree {
990 entries: vec![
991 TreeEntry {
992 name: b"a.txt".to_vec(),
993 mode: EntryMode::Blob,
994 object_hash: crate::hash::hash(b"a"),
995 },
996 TreeEntry {
997 name: b"b.txt".to_vec(),
998 mode: EntryMode::Blob,
999 object_hash: crate::hash::hash(b"b"),
1000 },
1001 ],
1002 });
1003 let bytes = serialize::serialize(&obj).unwrap();
1004 let h = store.write(&bytes).unwrap();
1005 let parsed = store.read_object(&h).unwrap();
1006 assert_eq!(parsed, obj);
1007 // Every `S: ObjectSource + ?Sized` caller (diff/blame/worktree
1008 // blob loading) resolves `read_object` through this trait impl,
1009 // not the inherent method above — pin that it agrees.
1010 let via_trait: &dyn ObjectSource = &store;
1011 assert_eq!(via_trait.read_object(&h).unwrap(), obj);
1012 }
1013
1014 /// Regression test for the double-decode fix: corrupting a stored
1015 /// `Tree`'s bytes must still surface `HashMismatch` through
1016 /// `read_object`'s merkle (`Some(obj)`-reusing) branch, not just
1017 /// through `read`'s. Also exercises `read_object` via the
1018 /// `ObjectSource` trait object — the actual call shape every
1019 /// `S: ObjectSource + ?Sized` caller (`diff::load_tree`,
1020 /// `LoadedBlob::load`) uses — to pin `impl ObjectSource for
1021 /// ObjectStore`'s `read_object` override against regressing back to
1022 /// the trait's double-decoding default.
1023 #[test]
1024 fn read_object_detects_tree_corruption_via_store_and_trait() {
1025 use crate::object::{EntryMode, Tree, TreeEntry};
1026 let (_dir, store) = fresh_store();
1027 let obj = Object::Tree(Tree {
1028 entries: vec![TreeEntry {
1029 name: b"a.txt".to_vec(),
1030 mode: EntryMode::Blob,
1031 object_hash: crate::hash::hash(b"a"),
1032 }],
1033 });
1034 let bytes = serialize::serialize(&obj).unwrap();
1035 let h = store.write(&bytes).unwrap();
1036
1037 let path = store.path_for(&h);
1038 let mut f = OpenOptions::new()
1039 .read(true)
1040 .write(true)
1041 .open(&path)
1042 .unwrap();
1043 f.seek(io::SeekFrom::Start(0)).unwrap();
1044 f.write_all(&[bytes[0] ^ 0xFF]).unwrap();
1045 f.sync_all().unwrap();
1046
1047 assert!(matches!(
1048 store.read_object(&h).unwrap_err(),
1049 StoreError::HashMismatch { .. }
1050 ));
1051 let via_trait: &dyn ObjectSource = &store;
1052 assert!(matches!(
1053 via_trait.read_object(&h).unwrap_err(),
1054 StoreError::HashMismatch { .. }
1055 ));
1056 }
1057
1058 #[test]
1059 fn write_is_idempotent() {
1060 let (_dir, store) = fresh_store();
1061 let bytes = b"duplicate".to_vec();
1062 let h1 = store.write(&bytes).unwrap();
1063 let h2 = store.write(&bytes).unwrap();
1064 assert_eq!(h1, h2);
1065 // Second write must not have produced any stray temp files.
1066 let shard = store.path_for(&h1);
1067 let parent = shard.parent().unwrap();
1068 let entries: Vec<_> = fs::read_dir(parent).unwrap().collect();
1069 assert_eq!(
1070 entries.len(),
1071 1,
1072 "shard dir should contain exactly the final object, no temp leaks"
1073 );
1074 }
1075
1076 #[test]
1077 fn read_missing_returns_not_found() {
1078 let (_dir, store) = fresh_store();
1079 let phony = hash::hash(b"never written");
1080 let err = store.read(&phony).unwrap_err();
1081 assert!(matches!(err, StoreError::ObjectNotFound(_)));
1082 assert!(!store.contains(&phony));
1083 }
1084
1085 #[test]
1086 fn write_rejects_oversize() {
1087 let (_dir, store) = fresh_store();
1088 // A REAL cap+1 body: `vec![0u8; n]` is `alloc_zeroed`, so the
1089 // 1 GiB+1 buffer costs lazily mapped zero pages, not RSS — and
1090 // the guard rejects on `len()` before anything reads the bytes.
1091 let oversize = vec![0u8; MAX_RAW_OBJECT_SIZE + 1];
1092 let err = store.write(&oversize).unwrap_err();
1093 assert!(matches!(err, StoreError::ObjectTooLarge), "got {err:?}");
1094 drop(oversize);
1095 // A realistic small write still works after the rejection.
1096 let h = store.write(&[0u8; 16]).unwrap();
1097 assert!(store.contains(&h));
1098 }
1099
1100 #[test]
1101 fn read_rejects_oversize_on_disk() {
1102 // Construct an oversize on-disk blob by hand and confirm `read`
1103 // refuses it. We use a sentinel size = MAX + 1 file padded with
1104 // zeros; this allocates ~1 GiB of disk, which is unfriendly in
1105 // unit tests, so instead we monkey-patch via a smaller-cap copy
1106 // of the read path: we synthesise a too-large file by truncating
1107 // a real one and verifying the comparison logic at the boundary.
1108 //
1109 // We exercise `MAX_RAW_OBJECT_SIZE` indirectly by writing a
1110 // small object and then *replacing* the on-disk file with one
1111 // whose `metadata().len()` exceeds the cap. We use sparse
1112 // truncation so no real disk is consumed.
1113 let (_dir, store) = fresh_store();
1114 let h = store.write(b"seed").unwrap();
1115 let path = store.path_for(&h);
1116 let f = OpenOptions::new().write(true).open(&path).unwrap();
1117 // Sparse extend to cap+1 bytes; allocates effectively no blocks.
1118 f.set_len(MAX_RAW_OBJECT_SIZE as u64 + 1).unwrap();
1119 drop(f);
1120 let err = store.read(&h).unwrap_err();
1121 assert!(matches!(err, StoreError::ObjectTooLarge));
1122 }
1123
1124 #[test]
1125 fn read_detects_corruption() {
1126 let (_dir, store) = fresh_store();
1127 let bytes = b"trustworthy".to_vec();
1128 let h = store.write(&bytes).unwrap();
1129 // Flip a single byte in the on-disk file and expect HashMismatch.
1130 let path = store.path_for(&h);
1131 {
1132 let mut f = OpenOptions::new()
1133 .read(true)
1134 .write(true)
1135 .open(&path)
1136 .unwrap();
1137 f.seek(io::SeekFrom::Start(0)).unwrap();
1138 f.write_all(&[bytes[0] ^ 0xFF]).unwrap();
1139 f.sync_all().unwrap();
1140 }
1141 let err = store.read(&h).unwrap_err();
1142 match err {
1143 StoreError::HashMismatch { expected, actual } => {
1144 assert_eq!(expected, to_hex(&h));
1145 assert_ne!(actual, expected, "actual must differ once corrupted");
1146 }
1147 other => panic!("expected HashMismatch, got {other:?}"),
1148 }
1149 }
1150
1151 /// Flip the first byte of the on-disk object file for `h`, in place —
1152 /// the "the bytes on disk no longer match `h`" setup for the
1153 /// `read`-vs-`read_unverified` corruption split test below. (The
1154 /// `DisplaySource` corruption tests carry their own copy in
1155 /// `source.rs`.)
1156 fn corrupt_first_byte(store: &ObjectStore, h: &Hash, first_byte: u8) {
1157 let path = store.path_for(h);
1158 let mut f = OpenOptions::new()
1159 .read(true)
1160 .write(true)
1161 .open(&path)
1162 .unwrap();
1163 f.seek(io::SeekFrom::Start(0)).unwrap();
1164 f.write_all(&[first_byte ^ 0xFF]).unwrap();
1165 f.sync_all().unwrap();
1166 }
1167
1168 #[test]
1169 fn read_unverified_returns_bytes_read_rejects_corruption() {
1170 // The core split this whole feature exists for: `read` must still
1171 // fail loudly on a corrupted object, while `read_unverified`
1172 // returns the (corrupt) bytes as-is — display paths only ever get
1173 // a bad render, never silently-wrong durable state.
1174 let (_dir, store) = fresh_store();
1175 let bytes = b"trustworthy".to_vec();
1176 let h = store.write(&bytes).unwrap();
1177 corrupt_first_byte(&store, &h, bytes[0]);
1178
1179 assert!(matches!(
1180 store.read(&h).unwrap_err(),
1181 StoreError::HashMismatch { .. }
1182 ));
1183
1184 let mut corrupted = bytes.clone();
1185 corrupted[0] ^= 0xFF;
1186 assert_eq!(store.read_unverified(&h).unwrap(), corrupted);
1187 }
1188
1189 #[test]
1190 fn path_layout_is_2_then_62_hex() {
1191 let (_dir, store) = fresh_store();
1192 let bytes = b"layout test".to_vec();
1193 let h = store.write(&bytes).unwrap();
1194 let hex = to_hex(&h);
1195 let path = store.path_for(&h);
1196 let parent = path.parent().unwrap();
1197 let parent_name = parent.file_name().unwrap().to_str().unwrap();
1198 let file_name = path.file_name().unwrap().to_str().unwrap();
1199 assert_eq!(parent_name.len(), 2);
1200 assert_eq!(file_name.len(), 62);
1201 assert_eq!(parent_name, &hex[..2]);
1202 assert_eq!(file_name, &hex[2..]);
1203 assert!(path.is_file(), "object file must exist at expected path");
1204 }
1205
1206 #[test]
1207 fn temp_file_left_behind_does_not_satisfy_contains() {
1208 // Simulate a crash mid-write: drop a stale `.tmp.*` file in the
1209 // shard dir without ever renaming. `contains()` must report the
1210 // object as absent, and the shard dir must not contain a real
1211 // entry that passes the hash check.
1212 let (_dir, store) = fresh_store();
1213 let target = hash::hash(b"never finalised");
1214 let final_path = store.path_for(&target);
1215 let shard = final_path.parent().unwrap();
1216 fs::create_dir_all(shard).unwrap();
1217 let stale = shard.join(format!(
1218 ".{}.tmp.0.0",
1219 final_path.file_name().unwrap().to_string_lossy()
1220 ));
1221 fs::write(&stale, b"partial").unwrap();
1222 assert!(stale.is_file());
1223 assert!(!store.contains(&target));
1224 // No file in the shard dir should hash to `target`.
1225 for entry in fs::read_dir(shard).unwrap() {
1226 let p = entry.unwrap().path();
1227 let bytes = fs::read(&p).unwrap();
1228 assert_ne!(
1229 hash::hash(&bytes),
1230 target,
1231 "stale temp file must not satisfy the target hash"
1232 );
1233 }
1234 }
1235}
1236
1237#[cfg(test)]
1238mod bulk_writer_tests {
1239 use super::*;
1240
1241 #[test]
1242 fn bulk_writer_round_trips_and_rewrites() {
1243 let td = tempfile::tempdir().unwrap();
1244 let store = ObjectStore::init(&RepoLayout::single(td.path())).unwrap();
1245 let obj = crate::serialize::serialize(&crate::object::Object::Blob(crate::object::Blob {
1246 data: b"bulk".to_vec(),
1247 }))
1248 .unwrap();
1249 let mut bw = store.bulk_writer();
1250 let h1 = bw.write(&obj).unwrap();
1251 // No existence short-circuit: rewriting is fine and heals
1252 // torn files on idempotent re-runs.
1253 let h2 = bw.write(&obj).unwrap();
1254 assert_eq!(h1, h2);
1255 bw.commit().unwrap();
1256 assert_eq!(store.read(&h1).unwrap(), obj);
1257 // Interoperates with the normal (fsynced) writer.
1258 assert_eq!(store.write(&obj).unwrap(), h1);
1259 }
1260}
1261
1262/// MKIT-55: a dedup hit must refresh the existing object's `mtime`, so
1263/// a writer outside gc's lock set (git import) that is about to publish
1264/// a ref over an old unreachable object does not lose it to a gc with a
1265/// grace window (SPEC-GC "Concurrent writers and the grace window").
1266#[cfg(test)]
1267mod dedup_mtime_tests {
1268 use super::*;
1269 use crate::layout::RepoLayout;
1270 use std::time::{Duration, SystemTime, UNIX_EPOCH};
1271
1272 const DAY: u64 = 24 * 60 * 60;
1273 /// gc's default grace window (SPEC-GC "Status").
1274 const GRACE: u64 = 14 * DAY;
1275
1276 fn repo() -> (tempfile::TempDir, ObjectStore) {
1277 let td = tempfile::tempdir().unwrap();
1278 let store = ObjectStore::init(&RepoLayout::single(td.path())).unwrap();
1279 crate::refs::init(&RepoLayout::single(td.path())).unwrap();
1280 (td, store)
1281 }
1282
1283 fn mtime(store: &ObjectStore, h: &Hash) -> SystemTime {
1284 fs::metadata(store.path_for(h)).unwrap().modified().unwrap()
1285 }
1286
1287 /// Age the object file past the grace window, as an object written
1288 /// long ago and since orphaned (e.g. a superseded commit whose
1289 /// recovery entry expired) would be.
1290 fn backdate(store: &ObjectStore, h: &Hash) -> SystemTime {
1291 let old = SystemTime::now() - Duration::from_secs(30 * DAY);
1292 File::open(store.path_for(h))
1293 .unwrap()
1294 .set_modified(old)
1295 .unwrap();
1296 assert!(mtime(store, h) <= old + Duration::from_secs(1));
1297 old
1298 }
1299
1300 fn now_secs() -> u64 {
1301 SystemTime::now()
1302 .duration_since(UNIX_EPOCH)
1303 .unwrap()
1304 .as_secs()
1305 }
1306
1307 /// Write an object, age it, dedup-write it again through `rewrite`,
1308 /// then run a gc with the default grace: the re-written object must
1309 /// have a fresh `mtime` and survive, while an equally old orphan
1310 /// nobody re-wrote is pruned (so the gc really sweeps).
1311 fn assert_dedup_refreshes(rewrite: impl Fn(&ObjectStore, &[u8]) -> Hash) {
1312 let (td, store) = repo();
1313 let bytes = b"object the import is about to publish".to_vec();
1314 let h = store.write(&bytes).unwrap();
1315 let control = store.write(b"old orphan nobody writes again").unwrap();
1316 backdate(&store, &h);
1317 backdate(&store, &control);
1318
1319 let before = SystemTime::now() - Duration::from_secs(1);
1320 assert_eq!(rewrite(&store, &bytes), h);
1321 assert!(
1322 mtime(&store, &h) >= before,
1323 "dedup hit must refresh the object's mtime"
1324 );
1325
1326 let report = crate::ops::gc::run_gc(
1327 &store,
1328 &RepoLayout::single(td.path()),
1329 now_secs(),
1330 GRACE,
1331 false,
1332 )
1333 .unwrap();
1334 assert!(store.contains(&h), "re-written object pruned: {report:?}");
1335 assert_eq!(store.read(&h).unwrap(), bytes);
1336 assert!(!store.contains(&control), "control orphan kept: {report:?}");
1337 }
1338
1339 /// A failed refresh must not be reported as a dedup hit:
1340 /// `BulkWriter` falls back to rewriting the object, which gives it
1341 /// a fresh `mtime` too.
1342 #[test]
1343 fn failed_refresh_falls_back_to_rewrite() {
1344 fn failing<T>(f: impl FnOnce() -> T) -> T {
1345 FAIL_REFRESH.set(true);
1346 let out = f();
1347 FAIL_REFRESH.set(false);
1348 out
1349 }
1350 assert_dedup_refreshes(|s, b| {
1351 let mut bw = s.bulk_writer();
1352 let h = failing(|| bw.write(b).unwrap());
1353 bw.commit().unwrap();
1354 h
1355 });
1356 }
1357
1358 /// Writable object files only: `BulkWriter::commit` re-opens a
1359 /// re-used object for writing to fsync it, which a read-only object
1360 /// file refuses (independent of the mtime refresh).
1361 #[test]
1362 fn bulk_writer_dedup_refreshes_mtime() {
1363 assert_dedup_refreshes(|s, b| {
1364 let mut bw = s.bulk_writer();
1365 let h = bw.write(b).unwrap();
1366 bw.commit().unwrap();
1367 h
1368 });
1369 }
1370}