Skip to main content

mkit_core/
batch.rs

1//! Batched-durability object writes.
2//!
3//! [`WriteBatch`] amortises the cost of crash durability across every
4//! object written by one logical command (an `add`, a `commit`, a pack
5//! unpack): objects are staged as barrier-synced temp files and become
6//! durable *and visible together* at [`WriteBatch::commit`], with **one**
7//! full flush per batch instead of two per object.
8//!
9//! # Durability contract
10//!
11//! The invariant the store actually needs is not "every object is
12//! durable the moment it is written" — it is:
13//!
14//! > A ref or index file is only ever written after every object it
15//! > references is durable, and a crash never produces a visible object
16//! > that fails the read-time hash check.
17//!
18//! `WriteBatch` preserves both halves:
19//!
20//! * Staged objects are **invisible** until `commit()` — renames are
21//!   deferred until after the batch's full flush, so another process's
22//!   `contains()` dedup can never observe (and then reference) an
23//!   object whose bytes are not yet durable.
24//! * `commit()` returns only after one full flush, every rename, and a
25//!   deduplicated flush of each touched shard directory. Callers MUST
26//!   order their ref/index writes after `commit()`.
27//! * A dropped (never committed) batch unlinks its temp files and
28//!   leaves the store untouched — aborting is free.
29//!
30//! # How the single flush is enough
31//!
32//! This is git's `core.fsyncMethod=batch` design (bulk-checkin) and
33//! `SQLite`'s macOS sync strategy:
34//!
35//! Every staged file gets a real per-file writeback at commit time —
36//! Apple `fcntl(F_BARRIERFSYNC)`, Linux `fdatasync`, Windows
37//! `FlushFileBuffers` — issued **concurrently** from a scoped-thread
38//! pool so the cost is device latency at queue depth, not
39//! latency×objects. The trailing constant-count full flushes
40//! (`F_FULLFSYNC` on Apple) cover ordering and the dirent updates.
41//! This holds on every filesystem; it does not depend on ext4
42//! ordered-data journaling.
43//!
44//! Workloads that prefer the historical schedule can select
45//! [`SyncPolicy::PerObject`] (config key `durability.objects =
46//! per-object`), which reproduces the old write path exactly.
47
48use std::collections::{HashMap, HashSet};
49use std::fmt;
50use std::fs::{self, File, OpenOptions};
51use std::io::{self, Write};
52use std::path::{Path, PathBuf};
53use std::sync::Mutex;
54
55use tempfile::TempPath;
56
57use crate::hash::Hash;
58use crate::store::{
59    MAX_RAW_OBJECT_SIZE, ObjectSink, ObjectStore, StoreError, StoreResult, sync_parent_dir,
60    temp_file_in,
61};
62
63/// When object writes become durable.
64#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
65pub enum SyncPolicy {
66    /// Historical behaviour: full flush + dir flush per object, object
67    /// visible immediately. O(objects) full flushes.
68    PerObject,
69    /// Stage now, make durable + visible together at
70    /// [`WriteBatch::commit`]. O(1) full flushes per batch. Default.
71    ///
72    /// (There is deliberately no "no flushes" policy: a visible object
73    /// MUST be durable — SPEC-OBJECTS §10.1 — because every writer's
74    /// content-addressed dedup trusts visibility. Ephemeral snapshots
75    /// belong in [`crate::store::EphemeralSink`], which never touches
76    /// the store.)
77    #[default]
78    Batch,
79}
80
81/// Flush/rename primitive seam between the store and the OS.
82///
83/// Production code uses [`RealSyncer`]; unit tests inject a recording
84/// double to assert flush *ordering* and *counts* — the tests in this
85/// module are the proof (and CI regression guard) of the O(1) full
86/// flushes per batch claim. Distinct from [`SyncPolicy`], which is a
87/// production knob deciding *whether* to flush; the syncer decides
88/// *how*.
89pub(crate) trait Syncer: Send + Sync + fmt::Debug {
90    /// Writeback + ordering barrier for one staged file. Must guarantee
91    /// the file's bytes reach the device before any later
92    /// [`Syncer::full`] completes, without forcing a device cache
93    /// flush.
94    fn barrier(&self, file: &File, path: &Path) -> io::Result<()>;
95    /// Full durable flush (device cache included) of `file`.
96    fn full(&self, file: &File, path: &Path) -> io::Result<()>;
97    /// Atomically rename a staged temp file into its final path,
98    /// replacing any existing file.
99    fn rename(&self, tmp: TempPath, final_path: &Path) -> io::Result<()>;
100    /// Durably flush the directory entry updates of `dir` — the legacy
101    /// per-object schedule ([`SyncPolicy::PerObject`] and
102    /// `ObjectStore::write`).
103    fn dir_sync(&self, dir: &Path) -> io::Result<()>;
104    /// Writeback + ordering barrier for the dirent updates of `dir`,
105    /// without forcing a device cache flush. Must guarantee the dirents
106    /// reach the device before a later [`Syncer::device_flush`]
107    /// completes. Batched schedule only; needs the trailing
108    /// `device_flush` to be durable.
109    fn dir_barrier(&self, dir: &Path) -> io::Result<()>;
110    /// Terminal full flush of the device write cache, anchored at the
111    /// store's `objects/` root. Makes everything previously
112    /// barrier-ordered (file data and dirents alike) durable.
113    fn device_flush(&self, objects_root: &Path) -> io::Result<()>;
114}
115
116/// Production [`Syncer`].
117#[derive(Debug)]
118pub(crate) struct RealSyncer;
119
120impl RealSyncer {
121    /// Per-file write barrier: order this file's writeback ahead of the
122    /// batch's terminal device flush, as cheaply as the platform allows.
123    #[cfg(any(target_os = "macos", target_os = "ios"))]
124    fn file_barrier(file: &File) -> io::Result<()> {
125        use std::os::unix::io::AsRawFd;
126        // Apple: `File::sync_data()`/`sync_all()` both map to
127        // `fcntl(F_FULLFSYNC)` — a full device-cache flush. F_BARRIERFSYNC
128        // is the cheaper primitive the module docs assume: it forces this
129        // file's writeback and orders it ahead of later writes WITHOUT a
130        // device flush. The single terminal `device_flush` (one
131        // F_FULLFSYNC) still makes the whole batch durable, so the per-
132        // file step stays a true barrier rather than N full flushes.
133        //
134        // SAFETY: `fcntl(2)` with `F_BARRIERFSYNC` takes only the fd and
135        // the command — it reads/writes no user memory, and the fd is
136        // valid for the borrow of `file`.
137        #[allow(unsafe_code)]
138        let rc = unsafe { libc::fcntl(file.as_raw_fd(), libc::F_BARRIERFSYNC) };
139        if rc == -1 {
140            // Some filesystems reject the fcntl — fall back to the full
141            // flush rather than weaken durability.
142            return file.sync_data();
143        }
144        Ok(())
145    }
146
147    /// Linux: `fdatasync`. Windows: `FlushFileBuffers` (requires the
148    /// write-capable handle the commit path now opens).
149    #[cfg(not(any(target_os = "macos", target_os = "ios")))]
150    fn file_barrier(file: &File) -> io::Result<()> {
151        file.sync_data()
152    }
153}
154
155impl Syncer for RealSyncer {
156    fn barrier(&self, file: &File, _path: &Path) -> io::Result<()> {
157        // Per-file writeback on EVERY platform — Apple: fcntl
158        // F_BARRIERFSYNC (writeback + ordering barrier, no device-cache
159        // flush); Linux: fdatasync; Windows: FlushFileBuffers. This is
160        // what makes a committed batch durable on all filesystems, not
161        // just metadata-journaling ones in ordered-data mode — XFS,
162        // btrfs, ext4 data=writeback, and NTFS get the same guarantee.
163        // The cost is bounded: barriers are issued concurrently from a
164        // thread pool at commit, so wall-clock is latency/queue-depth,
165        // not latency×objects.
166        Self::file_barrier(file)
167    }
168
169    fn full(&self, file: &File, _path: &Path) -> io::Result<()> {
170        file.sync_all()
171    }
172
173    fn rename(&self, tmp: TempPath, final_path: &Path) -> io::Result<()> {
174        // Cross-platform atomic replace: rename(2) on Unix, MoveFileExW
175        // with MOVEFILE_REPLACE_EXISTING on Windows.
176        tmp.persist(final_path).map_err(|e| e.error)?;
177        Ok(())
178    }
179
180    fn dir_sync(&self, dir: &Path) -> io::Result<()> {
181        sync_parent_dir(dir)
182    }
183
184    #[cfg(any(target_os = "macos", target_os = "ios"))]
185    fn dir_barrier(&self, dir: &Path) -> io::Result<()> {
186        // F_BARRIERFSYNC on the directory fd: pushes the dirent
187        // updates toward the device and orders them ahead of the
188        // batch's terminal F_FULLFSYNC, at a fraction of its cost.
189        // Some filesystems reject the fcntl on directories — fall back
190        // to the full dir fsync rather than weaken durability.
191        match File::open(dir) {
192            Ok(d) => d.sync_data().or_else(|_| d.sync_all()),
193            Err(e) if e.kind() == io::ErrorKind::NotFound => Ok(()),
194            Err(e) => Err(e),
195        }
196    }
197
198    #[cfg(not(any(target_os = "macos", target_os = "ios")))]
199    fn dir_barrier(&self, dir: &Path) -> io::Result<()> {
200        // Linux: the directory fsync IS the durability mechanism (the
201        // journal commit orders ordered-mode file data ahead of the
202        // dirents), so the "barrier" must stay a real fsync. Windows:
203        // no directory flush primitive — no-op, the device_flush
204        // covers what the OS exposes.
205        sync_parent_dir(dir)
206    }
207
208    #[cfg(unix)]
209    fn device_flush(&self, objects_root: &Path) -> io::Result<()> {
210        // macOS: sync_all on any fd is F_FULLFSYNC — flushes the whole
211        // device write cache, making every prior barrier durable.
212        // Linux: fsync of the objects root; cheap insurance on top of
213        // the per-dir fsyncs above.
214        match File::open(objects_root) {
215            Ok(d) => d.sync_all(),
216            Err(e) if e.kind() == io::ErrorKind::NotFound => Ok(()),
217            Err(e) => Err(e),
218        }
219    }
220
221    #[cfg(not(unix))]
222    #[allow(clippy::unnecessary_wraps)]
223    fn device_flush(&self, _objects_root: &Path) -> io::Result<()> {
224        Ok(())
225    }
226}
227
228#[derive(Debug, Default)]
229struct BatchState {
230    /// hash → staged-but-invisible temp file, insertion-deduped. The
231    /// final path is derived from the hash at commit (`path_for`) — a
232    /// single source of truth for the layout. Handles are closed
233    /// (`TempPath`) so a large batch does not exhaust the fd limit;
234    /// deletion-on-drop is retained for abort cleanup.
235    staged: HashMap<Hash, TempPath>,
236    /// Shard directories whose dirents this batch must flush at commit.
237    /// Includes dedup hits: an object made visible by another process
238    /// may not have a durable dirent yet, and our commit is about to
239    /// reference it.
240    touched_shards: HashSet<PathBuf>,
241    /// Shard directories this batch has already `create_dir_all`'d —
242    /// at most 256 exist, so memoizing saves ~one mkdir syscall per
243    /// object on large ingests.
244    created_shards: HashSet<PathBuf>,
245}
246
247/// A set of object writes that become durable and visible together.
248/// Created by [`ObjectStore::batch`]. See the module docs for the
249/// durability contract.
250#[derive(Debug)]
251pub struct WriteBatch<'s> {
252    store: &'s ObjectStore,
253    policy: SyncPolicy,
254    // Interior mutability so `&self` writes work and a future parallel
255    // ingest can share one batch across worker threads.
256    inner: Mutex<BatchState>,
257}
258
259impl<'s> WriteBatch<'s> {
260    pub(crate) fn new(store: &'s ObjectStore, policy: SyncPolicy) -> Self {
261        Self {
262            store,
263            policy,
264            inner: Mutex::new(BatchState::default()),
265        }
266    }
267
268    /// Hash `bytes`, dedup against staged and on-disk objects, and
269    /// stage (policy `Batch`) or durably write (policy `PerObject`) the
270    /// object. Returns the BLAKE3 hash either way.
271    pub fn write(&self, bytes: &[u8]) -> StoreResult<Hash> {
272        self.write_parts(&[bytes])
273    }
274
275    /// [`Self::write`] for an object whose bytes are the concatenation
276    /// of `parts`, hashed and written streaming — no concatenated
277    /// buffer is materialised.
278    ///
279    /// # Panics
280    ///
281    /// Panics only if the internal hash-to-path mapping produces a path
282    /// without a parent directory (impossible by construction) or if a
283    /// previous write panicked while holding the batch mutex.
284    pub fn write_parts(&self, parts: &[&[u8]]) -> StoreResult<Hash> {
285        let mut total: usize = 0;
286        for p in parts {
287            total = total
288                .checked_add(p.len())
289                .ok_or(StoreError::ObjectTooLarge)?;
290        }
291        if total > MAX_RAW_OBJECT_SIZE {
292            return Err(StoreError::ObjectTooLarge);
293        }
294        // One id dispatch shared by every part-wise sink: a merkle type
295        // buffers + uses its BMT root, a byte-hashed type streams.
296        self.write_prehashed(crate::object::object_id_from_parts(parts), parts)
297    }
298
299    /// Stage `parts` under the caller-supplied content hash, skipping
300    /// the BLAKE3 pass. `pub(crate)` and reserved for callers that have
301    /// PROVABLY just hashed the same bytes (pack unpack hashes every
302    /// entry to build its report; re-hashing in the batch doubled the
303    /// CPU of every clone/fetch). A wrong hash here would corrupt the
304    /// content addressing — never expose this publicly.
305    pub(crate) fn write_prehashed(&self, h: Hash, parts: &[&[u8]]) -> StoreResult<Hash> {
306        let final_path = self.store.path_for(&h);
307        let shard_dir = final_path
308            .parent()
309            .expect("object path always has a 2-hex parent")
310            .to_path_buf();
311
312        // Short lock: staged-dedup check + mkdir memoization decision.
313        // The file I/O below runs OUTSIDE the lock so concurrent
314        // writers sharing one batch don't convoy on each other's
315        // write_all calls.
316        let need_mkdir = {
317            let st = self.inner.lock().expect("batch state mutex poisoned");
318            if st.staged.contains_key(&h) {
319                return Ok(h);
320            }
321            !st.created_shards.contains(&shard_dir)
322        };
323        if final_path.exists() {
324            // Dedup hit: the object is visible, but if another process
325            // renamed it and has not yet flushed the dirent, it may not
326            // be durable. We are about to reference it, so flush its
327            // shard dir at commit. The `mtime` is left alone (see
328            // `store::refresh_mtime`).
329            self.inner
330                .lock()
331                .expect("batch state mutex poisoned")
332                .touched_shards
333                .insert(shard_dir);
334            return Ok(h);
335        }
336        if need_mkdir {
337            fs::create_dir_all(&shard_dir)?;
338        }
339        let file_name = final_path
340            .file_name()
341            .expect("object path has file name")
342            .to_string_lossy();
343        let mut tmp = temp_file_in(&shard_dir, &file_name)?;
344        for p in parts {
345            tmp.as_file_mut().write_all(p)?;
346        }
347        let syncer = self.store.syncer();
348        match self.policy {
349            SyncPolicy::PerObject => {
350                // Historical write path, immediately durable + visible.
351                syncer.full(tmp.as_file(), tmp.path())?;
352                syncer.rename(tmp.into_temp_path(), &final_path)?;
353                syncer.dir_sync(&shard_dir)?;
354                let mut st = self.inner.lock().expect("batch state mutex poisoned");
355                st.created_shards.insert(shard_dir);
356            }
357            SyncPolicy::Batch => {
358                // No flush here: barriers for every staged file are
359                // issued concurrently at commit() — sequential
360                // per-file barriers would re-serialise the batch on
361                // device latency (measured ~4.6ms per F_BARRIERFSYNC
362                // on Apple SSDs, ~8s for a 100 MiB ingest).
363                let mut st = self.inner.lock().expect("batch state mutex poisoned");
364                // Lost race against a concurrent writer of the same
365                // object within this batch: keep theirs, drop our tmp
366                // (content-addressed — byte-identical by construction).
367                st.staged.entry(h).or_insert_with(|| tmp.into_temp_path());
368                st.touched_shards.insert(shard_dir.clone());
369                st.created_shards.insert(shard_dir);
370            }
371        }
372        Ok(h)
373    }
374
375    /// True when `h` is staged in this batch or already in the store.
376    ///
377    /// # Panics
378    ///
379    /// Panics only if a previous write panicked while holding the batch
380    /// mutex.
381    #[must_use]
382    pub fn contains(&self, h: &Hash) -> bool {
383        self.inner
384            .lock()
385            .expect("batch state mutex poisoned")
386            .staged
387            .contains_key(h)
388            || self.store.contains(h)
389    }
390
391    /// Make every staged object durable and visible: one full flush,
392    /// then all renames, then deduplicated shard-directory flushes.
393    ///
394    /// After `commit()` returns `Ok`, every hash returned by
395    /// [`Self::write`]/[`Self::write_parts`] is durable AND visible.
396    /// Callers MUST call this before reading any object written by this
397    /// batch and before writing any ref/index that references one.
398    ///
399    /// If `commit()` fails partway, already-renamed objects remain
400    /// visible (content-addressing makes re-running the command
401    /// idempotent) and not-yet-renamed temp files are unlinked on drop.
402    ///
403    /// # Panics
404    ///
405    /// Panics only if a previous write panicked while holding the batch
406    /// mutex.
407    pub fn commit(self) -> StoreResult<()> {
408        let st = self.inner.into_inner().expect("batch state mutex poisoned");
409        let syncer = self.store.syncer();
410        match self.policy {
411            // Every write was already made durable and visible.
412            SyncPolicy::PerObject => Ok(()),
413            SyncPolicy::Batch => {
414                let staged: Vec<(Hash, TempPath)> = st.staged.into_iter().collect();
415                // 1. Barrier every staged file, concurrently. Each
416                //    barrier initiates writeback for its file and
417                //    orders it ahead of the full flush below; issuing
418                //    them from worker threads overlaps their device
419                //    latency (queue depth) instead of paying it
420                //    serially per file. All barriers complete (joined)
421                //    before the flush is issued.
422                parallel_io(staged.len(), |i| {
423                    // Write-capable handle: the barrier maps to
424                    // `FlushFileBuffers` on Windows, which rejects a
425                    // read-only handle. `write(true)` opens the existing
426                    // temp without truncating it.
427                    let f = OpenOptions::new().write(true).open(&staged[i].1)?;
428                    syncer.barrier(&f, &staged[i].1)
429                })?;
430                // 2. One full flush — any staged file serves as the
431                //    anchor; the barriers ordered every staged write
432                //    ahead of it (see module docs). Pure-dedup batches
433                //    (nothing staged) skip it: the objects were made
434                //    durable by whoever renamed them into visibility.
435                if let Some((_, tmp)) = staged.first() {
436                    // Write-capable handle for the same Windows reason as
437                    // the barrier above (`sync_all` → `FlushFileBuffers`).
438                    let f = OpenOptions::new().write(true).open(tmp)?;
439                    syncer.full(&f, tmp)?;
440                }
441                // 3. Renames: objects become visible only now, after
442                //    their bytes are durable — another process's dedup
443                //    can never reference a non-durable object. Final
444                //    paths derive from the hashes (single layout rule).
445                for (h, tmp) in staged {
446                    syncer.rename(tmp, &self.store.path_for(&h))?;
447                }
448                // 4. Dirent barriers, once per touched shard dir,
449                //    concurrently (sorted first so the work list is
450                //    deterministic).
451                let mut shards: Vec<PathBuf> = st.touched_shards.into_iter().collect();
452                shards.sort();
453                if !shards.is_empty() {
454                    parallel_io(shards.len(), |i| syncer.dir_barrier(&shards[i]))?;
455                    // 5. Terminal device flush: makes the dirent
456                    //    barriers (and, on platforms where step 2 was
457                    //    a barrier-anchored flush, everything) durable.
458                    syncer.device_flush(self.store.objects_root())?;
459                }
460                Ok(())
461            }
462        }
463    }
464}
465
466/// Worker-pool cap for [`parallel_io`]. Deliberately NOT
467/// `std::thread::available_parallelism()`: a worker here spends nearly
468/// all its time blocked in the kernel on a barrier/fsync-class syscall,
469/// not consuming CPU, so CPU core count is the wrong resource to size
470/// against (the classic pool-sizing formula — thread count scaling with
471/// `1 + wait/compute`, not compute alone — puts this workload's ideal
472/// count far above core count; `tokio::task::spawn_blocking`'s default
473/// pool of 512 threads exists for the identical reason).
474///
475/// 64 is chosen from a benchmark sweep (16/32/64/128 workers, 1000-object
476/// batches) on macOS/APFS (issue #864), repeated with 25-sample mean/
477/// stddev tracking after an initial small-sample pass overstated its
478/// confidence in the 128 result. The 16-vs-64 comparison is robust and
479/// reproduced across three independent runs: 64 is consistently ~20-30%
480/// faster than 16 with visibly tighter variance (stddev in the 10-30ms
481/// range vs 16's 25-40ms, non-overlapping in every run). 32 and 128,
482/// by contrast, showed high run-to-run variance (stddev over 100ms) and
483/// no reliably-ordered result against 64 — plausibly 128 threads hitting
484/// real scheduling contention on a 16-core machine, but not confidently
485/// distinguishable from noise at the sample sizes tested here. Treat
486/// "64" as validated; treat any specific claim about 32 or 128 as
487/// unresolved, not as "128 regresses." fsync concurrency is ultimately
488/// gated by filesystem journal serialization, not raw device queue
489/// depth, so diminishing (and possibly reversing) returns somewhere
490/// past 64 is expected in principle even if this data can't yet pin
491/// down exactly where. Linux/ext4 numbers are not gathered at all; re-tune
492/// if they diverge meaningfully from the macOS data.
493const MAX_SYNC_WORKERS: usize = 64;
494
495/// Run `op(0..count)` across a pool of up to [`MAX_SYNC_WORKERS`] scoped
496/// threads, joining them all before returning. Concurrency overlaps
497/// per-item device latency (~ms per flush primitive on Apple SSDs) that
498/// would otherwise serialise a batch; correctness only needs *all*
499/// items complete before the caller proceeds, which the join
500/// guarantees. Returns the first error observed, if any.
501fn parallel_io(count: usize, op: impl Fn(usize) -> io::Result<()> + Sync) -> io::Result<()> {
502    if count == 0 {
503        return Ok(());
504    }
505    let workers = MAX_SYNC_WORKERS.min(count);
506    if workers == 1 {
507        return (0..count).try_for_each(op);
508    }
509    let next = std::sync::atomic::AtomicUsize::new(0);
510    let mut results: Vec<io::Result<()>> = Vec::new();
511    std::thread::scope(|scope| {
512        let handles: Vec<_> = (0..workers)
513            .map(|_| {
514                let next = &next;
515                let op = &op;
516                scope.spawn(move || -> io::Result<()> {
517                    loop {
518                        let i = next.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
519                        if i >= count {
520                            return Ok(());
521                        }
522                        op(i)?;
523                    }
524                })
525            })
526            .collect();
527        for h in handles {
528            results.push(h.join().expect("parallel io worker panicked"));
529        }
530    });
531    results.into_iter().collect()
532}
533
534impl ObjectSink for WriteBatch<'_> {
535    fn put(&self, bytes: &[u8]) -> StoreResult<Hash> {
536        self.write(bytes)
537    }
538
539    fn put_parts(&self, parts: &[&[u8]]) -> StoreResult<Hash> {
540        self.write_parts(parts)
541    }
542
543    fn has(&self, h: &Hash) -> bool {
544        self.contains(h)
545    }
546}
547
548/// Test doubles shared with other modules' sync-behaviour tests
549/// (`pack.rs` asserts unpack costs one flush; later, `worktree.rs` and
550/// `ops/diff.rs` assert their flush budgets).
551#[cfg(test)]
552pub(crate) mod testing {
553    use super::*;
554
555    /// Every Syncer call, in order. `Rename` records the temp path so
556    /// tests can pair a rename with the `Barrier` that staged it. The
557    /// terminal device flush is recorded as `Full` so "full flush
558    /// count" stays the single number the O(1) tests bound.
559    #[derive(Debug, Clone, PartialEq, Eq)]
560    pub(crate) enum Ev {
561        Barrier(PathBuf),
562        Full(PathBuf),
563        Rename { tmp: PathBuf, dst: PathBuf },
564        DirSync(PathBuf),
565        DirBarrier(PathBuf),
566    }
567
568    /// Recording double: logs ordering, performs renames for real (so
569    /// the store stays functional under test), skips actual flushes
570    /// (they have no observable filesystem effect).
571    #[derive(Debug, Default)]
572    pub(crate) struct RecordingSyncer {
573        events: Mutex<Vec<Ev>>,
574    }
575
576    impl RecordingSyncer {
577        pub(crate) fn events(&self) -> Vec<Ev> {
578            self.events.lock().unwrap().clone()
579        }
580    }
581
582    impl Syncer for RecordingSyncer {
583        fn barrier(&self, _file: &File, path: &Path) -> io::Result<()> {
584            self.events
585                .lock()
586                .unwrap()
587                .push(Ev::Barrier(path.to_path_buf()));
588            Ok(())
589        }
590
591        fn full(&self, _file: &File, path: &Path) -> io::Result<()> {
592            self.events
593                .lock()
594                .unwrap()
595                .push(Ev::Full(path.to_path_buf()));
596            Ok(())
597        }
598
599        fn rename(&self, tmp: TempPath, final_path: &Path) -> io::Result<()> {
600            self.events.lock().unwrap().push(Ev::Rename {
601                tmp: tmp.to_path_buf(),
602                dst: final_path.to_path_buf(),
603            });
604            tmp.persist(final_path).map_err(|e| e.error)?;
605            Ok(())
606        }
607
608        fn dir_sync(&self, dir: &Path) -> io::Result<()> {
609            self.events
610                .lock()
611                .unwrap()
612                .push(Ev::DirSync(dir.to_path_buf()));
613            Ok(())
614        }
615
616        fn dir_barrier(&self, dir: &Path) -> io::Result<()> {
617            self.events
618                .lock()
619                .unwrap()
620                .push(Ev::DirBarrier(dir.to_path_buf()));
621            Ok(())
622        }
623
624        fn device_flush(&self, objects_root: &Path) -> io::Result<()> {
625            self.events
626                .lock()
627                .unwrap()
628                .push(Ev::Full(objects_root.to_path_buf()));
629            Ok(())
630        }
631    }
632}
633
634#[cfg(test)]
635mod tests {
636    use super::testing::{Ev, RecordingSyncer};
637    use super::*;
638    use crate::hash;
639    use proptest::prelude::*;
640    use std::sync::Arc;
641    use tempfile::TempDir;
642
643    fn fresh_store() -> (TempDir, ObjectStore) {
644        let dir = TempDir::new().expect("tempdir");
645        let store =
646            ObjectStore::init(&crate::layout::RepoLayout::single(dir.path())).expect("init");
647        (dir, store)
648    }
649
650    /// Store with a shared `RecordingSyncer` injected.
651    fn recording_store() -> (TempDir, ObjectStore, Arc<RecordingSyncer>) {
652        let (dir, mut store) = fresh_store();
653        let rec = Arc::new(RecordingSyncer::default());
654        store.set_syncer(rec.clone());
655        (dir, store, rec)
656    }
657
658    /// Count object files (62-hex names) under `objects/`, recursively.
659    fn object_file_count(store: &ObjectStore) -> usize {
660        store.iter_object_hashes().unwrap().len()
661    }
662
663    /// Count ALL files under `objects/` including temp files.
664    fn any_file_count(store: &ObjectStore) -> usize {
665        fn walk(dir: &Path, n: &mut usize) {
666            if let Ok(rd) = fs::read_dir(dir) {
667                for e in rd.flatten() {
668                    let p = e.path();
669                    if p.is_dir() {
670                        walk(&p, n);
671                    } else {
672                        *n += 1;
673                    }
674                }
675            }
676        }
677        let mut n = 0;
678        walk(store.objects_root(), &mut n);
679        n
680    }
681
682    // ---- Cycle 1: staging invisibility --------------------------------
683
684    #[test]
685    fn staged_object_is_invisible_until_commit() {
686        let (dir, store) = fresh_store();
687        let batch = store.batch();
688        let h = batch.write(b"staged bytes").unwrap();
689
690        // A second, independent handle must not see the object yet.
691        let other = ObjectStore::open(&crate::layout::RepoLayout::single(dir.path())).unwrap();
692        assert!(
693            !other.contains(&h),
694            "staged object must be invisible before commit"
695        );
696        assert!(other.read(&h).is_err());
697
698        batch.commit().unwrap();
699        assert!(other.contains(&h), "committed object must be visible");
700        assert_eq!(other.read(&h).unwrap(), b"staged bytes");
701    }
702
703    #[test]
704    fn batch_contains_sees_staged_and_disk() {
705        let (_dir, store) = fresh_store();
706        let on_disk = store.write(b"already stored").unwrap();
707        let batch = store.batch();
708        let staged = batch.write(b"only staged").unwrap();
709
710        assert!(batch.contains(&on_disk), "must see on-disk objects");
711        assert!(batch.contains(&staged), "must see its own staged objects");
712        let phony = hash::hash(b"never written");
713        assert!(!batch.contains(&phony));
714    }
715
716    #[test]
717    fn dropped_batch_leaves_no_tmp_files_and_no_objects() {
718        let (_dir, store) = fresh_store();
719        {
720            let batch = store.batch();
721            batch.write(b"abort me 1").unwrap();
722            batch.write(b"abort me 2").unwrap();
723            batch.write(b"abort me 3").unwrap();
724            // dropped without commit
725        }
726        assert_eq!(object_file_count(&store), 0, "no objects may be visible");
727        assert_eq!(
728            any_file_count(&store),
729            0,
730            "no temp files may leak from an aborted batch"
731        );
732    }
733
734    // ---- Cycle 2: flush ordering (the O(1) proof) ----------------------
735
736    fn fifty_distinct_objects() -> Vec<Vec<u8>> {
737        (0u32..50)
738            .map(|i| format!("object #{i}").into_bytes())
739            .collect()
740    }
741
742    #[test]
743    fn batch_commit_full_flush_count_is_constant() {
744        // The O(1) proof: a committed batch costs exactly TWO full
745        // flushes — one covering staged file data (pre-rename), one
746        // terminal device flush covering the dirent updates —
747        // regardless of how many objects it stages.
748        for count in [3usize, 50] {
749            let (_dir, store, rec) = recording_store();
750            let batch = store.batch();
751            for bytes in fifty_distinct_objects().into_iter().take(count) {
752                batch.write(&bytes).unwrap();
753            }
754            // Duplicates must not add flushes either.
755            batch.write(b"object #0").unwrap();
756            batch.commit().unwrap();
757
758            let fulls = rec
759                .events()
760                .iter()
761                .filter(|e| matches!(e, Ev::Full(_)))
762                .count();
763            assert_eq!(
764                fulls, 2,
765                "a {count}-object batch must cost exactly two full flushes"
766            );
767        }
768    }
769
770    #[test]
771    fn every_rename_is_preceded_by_its_barrier_and_the_full_flush() {
772        let (_dir, store, rec) = recording_store();
773        let batch = store.batch();
774        for bytes in fifty_distinct_objects() {
775            batch.write(&bytes).unwrap();
776        }
777        batch.commit().unwrap();
778
779        let evs = rec.events();
780        let full_pos = evs
781            .iter()
782            .position(|e| matches!(e, Ev::Full(_)))
783            .expect("one full flush");
784        // Every barrier must complete before the full flush is issued —
785        // the flush only covers writes the barriers pushed ahead of it.
786        let last_barrier = evs
787            .iter()
788            .rposition(|e| matches!(e, Ev::Barrier(_)))
789            .expect("barriers recorded");
790        assert!(
791            last_barrier < full_pos,
792            "all barriers (last at {last_barrier}) must precede the full flush at {full_pos}"
793        );
794        for (i, ev) in evs.iter().enumerate() {
795            if let Ev::Rename { tmp, .. } = ev {
796                assert!(
797                    i > full_pos,
798                    "rename at {i} must come after the full flush at {full_pos}"
799                );
800                let barrier_pos = evs
801                    .iter()
802                    .position(|e| matches!(e, Ev::Barrier(p) if p == tmp))
803                    .unwrap_or_else(|| panic!("no barrier recorded for {}", tmp.display()));
804                assert!(
805                    barrier_pos < i,
806                    "barrier for {} must precede its rename",
807                    tmp.display()
808                );
809            }
810        }
811    }
812
813    #[test]
814    fn dir_syncs_come_after_all_renames_and_are_deduped() {
815        let (_dir, store, rec) = recording_store();
816        let batch = store.batch();
817        let mut shards = HashSet::new();
818        for bytes in fifty_distinct_objects() {
819            let h = batch.write(&bytes).unwrap();
820            shards.insert(store.path_for(&h).parent().unwrap().to_path_buf());
821        }
822        batch.commit().unwrap();
823
824        let evs = rec.events();
825        let last_rename = evs
826            .iter()
827            .rposition(|e| matches!(e, Ev::Rename { .. }))
828            .expect("renames recorded");
829        let dir_barriers: Vec<(usize, &PathBuf)> = evs
830            .iter()
831            .enumerate()
832            .filter_map(|(i, e)| match e {
833                Ev::DirBarrier(p) => Some((i, p)),
834                _ => None,
835            })
836            .collect();
837
838        let synced: HashSet<PathBuf> = dir_barriers.iter().map(|(_, p)| (*p).clone()).collect();
839        assert_eq!(
840            dir_barriers.len(),
841            synced.len(),
842            "each shard dir must be flushed exactly once"
843        );
844        assert_eq!(synced, shards, "exactly the touched shards are flushed");
845        for (i, p) in &dir_barriers {
846            assert!(
847                *i > last_rename,
848                "dir barrier of {} at {i} must come after the last rename at {last_rename}",
849                p.display()
850            );
851        }
852        // The terminal device flush makes the dirent barriers durable —
853        // it must be the last sync event of the batch.
854        let last_full = evs
855            .iter()
856            .rposition(|e| matches!(e, Ev::Full(_)))
857            .expect("device flush recorded");
858        let last_dir_barrier = dir_barriers.last().expect("dir barriers recorded").0;
859        assert!(
860            last_full > last_dir_barrier,
861            "device flush at {last_full} must follow the last dir barrier at {last_dir_barrier}"
862        );
863    }
864
865    #[test]
866    fn dedup_hit_still_dir_syncs_at_commit() {
867        let (_dir, store, rec) = recording_store();
868        // Pre-store the object (e.g. another process raced us there).
869        let h = store.write(b"already present").unwrap();
870        let shard = store.path_for(&h).parent().unwrap().to_path_buf();
871
872        let batch = store.batch();
873        let h2 = batch.write(b"already present").unwrap();
874        assert_eq!(h, h2);
875        let before = rec.events().len();
876        batch.commit().unwrap();
877
878        let evs = rec.events()[before..].to_vec();
879        assert!(
880            !evs.iter().any(|e| matches!(e, Ev::Rename { .. })),
881            "dedup hit must not stage or rename anything"
882        );
883        assert!(
884            evs.contains(&Ev::DirBarrier(shard)),
885            "commit must still flush the dedup-hit shard dir: its dirent \
886             may not be durable yet and we are about to reference it"
887        );
888        assert!(
889            matches!(evs.last(), Some(Ev::Full(_))),
890            "the dir barrier needs a trailing device flush to be durable"
891        );
892    }
893
894    #[test]
895    fn sync_policy_per_object_matches_legacy_event_pattern() {
896        // Legacy pattern = ObjectStore::write: Full(tmp), Rename, DirSync
897        // per object, in that order, objects visible immediately.
898        let (_dir, store, rec) = recording_store();
899        let legacy_h = store.write(b"legacy path").unwrap();
900        let legacy: Vec<Ev> = rec.events();
901
902        let batch = store.batch_with_policy(SyncPolicy::PerObject);
903        let h = batch.write(b"per object path").unwrap();
904        assert!(
905            store.contains(&h),
906            "PerObject writes must be visible immediately, pre-commit"
907        );
908        let per_object: Vec<Ev> = rec.events()[legacy.len()..].to_vec();
909
910        let kinds = |evs: &[Ev]| -> Vec<u8> {
911            evs.iter()
912                .map(|e| match e {
913                    Ev::Barrier(_) => 0u8,
914                    Ev::Full(_) => 1,
915                    Ev::Rename { .. } => 2,
916                    Ev::DirSync(_) => 3,
917                    Ev::DirBarrier(_) => 4,
918                })
919                .collect()
920        };
921        assert_eq!(
922            kinds(&per_object),
923            kinds(&legacy),
924            "PerObject batch must reproduce the legacy per-write sync pattern"
925        );
926        // commit() of a PerObject batch is a no-op.
927        let before = rec.events().len();
928        batch.commit().unwrap();
929        assert_eq!(rec.events().len(), before);
930        let _ = legacy_h;
931    }
932
933    // ---- Cycle 3: equivalence & limits ---------------------------------
934
935    #[test]
936    fn idempotent_duplicate_writes_in_one_batch_stage_once() {
937        let (_dir, store, rec) = recording_store();
938        let batch = store.batch();
939        let h1 = batch.write(b"twice staged").unwrap();
940        let h2 = batch.write(b"twice staged").unwrap();
941        assert_eq!(h1, h2);
942        batch.commit().unwrap();
943
944        let evs = rec.events();
945        let barriers = evs.iter().filter(|e| matches!(e, Ev::Barrier(_))).count();
946        let renames = evs
947            .iter()
948            .filter(|e| matches!(e, Ev::Rename { .. }))
949            .count();
950        assert_eq!(barriers, 1, "duplicate must not re-stage");
951        assert_eq!(renames, 1, "duplicate must not re-rename");
952    }
953
954    #[test]
955    fn batch_write_rejects_oversize() {
956        // Mirrors store::tests::write_rejects_oversize: a REAL cap+1
957        // body (lazily mapped zero pages via `alloc_zeroed`, so no RSS
958        // cost) must be refused by the batch's size guard before any
959        // hashing or staging.
960        let (_dir, store) = fresh_store();
961        let batch = store.batch();
962        let oversize = vec![0u8; MAX_RAW_OBJECT_SIZE + 1];
963        let err = batch.write(&oversize).unwrap_err();
964        assert!(matches!(err, StoreError::ObjectTooLarge), "got {err:?}");
965        drop(oversize);
966        // The batch stays usable after the rejection.
967        let h = batch.write(&[0u8; 16]).unwrap();
968        batch.commit().unwrap();
969        assert!(store.contains(&h));
970    }
971
972    proptest! {
973        #[test]
974        fn batch_write_hash_equals_store_write_hash(bytes in proptest::collection::vec(any::<u8>(), 0..4096)) {
975            let (_dir, store) = fresh_store();
976            let batch = store.batch();
977            let h_batch = batch.write(&bytes).unwrap();
978            batch.commit().unwrap();
979            let on_disk_via_batch = store.read(&h_batch).unwrap();
980
981            let (_dir2, store2) = fresh_store();
982            let h_store = store2.write(&bytes).unwrap();
983            prop_assert_eq!(h_batch, h_store, "batch and store writes must agree on the hash");
984            prop_assert_eq!(on_disk_via_batch, bytes, "on-disk bytes must round-trip");
985        }
986
987        #[test]
988        fn write_parts_equals_concatenated_write(
989            parts in proptest::collection::vec(proptest::collection::vec(any::<u8>(), 0..512), 0..8)
990        ) {
991            let (_dir, store) = fresh_store();
992            let concatenated: Vec<u8> = parts.iter().flatten().copied().collect();
993
994            let batch = store.batch();
995            let slices: Vec<&[u8]> = parts.iter().map(Vec::as_slice).collect();
996            let h_parts = batch.write_parts(&slices).unwrap();
997            let h_whole = batch.write(&concatenated).unwrap();
998            batch.commit().unwrap();
999
1000            prop_assert_eq!(h_parts, h_whole, "parts and whole must hash identically");
1001            prop_assert_eq!(store.read(&h_parts).unwrap(), concatenated);
1002        }
1003    }
1004}