Skip to main content

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}