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}