Skip to main content

objects/store/fs/
fs_store.rs

1// SPDX-License-Identifier: Apache-2.0
2//! Core FsStore structure.
3
4#[cfg(test)]
5use std::sync::atomic::AtomicUsize;
6use std::{
7    collections::{BTreeSet, HashMap, VecDeque},
8    hash::Hash,
9    path::{Path, PathBuf},
10    sync::{
11        Arc, Mutex, RwLock,
12        atomic::{AtomicBool, Ordering},
13    },
14};
15
16use heddle_format::compression::CompressionConfig;
17
18use super::{
19    fs_io::{AtomicWriteMode, write_atomic},
20    fs_paths::{
21        actions_dir, blobs_dir, packs_dir, partial_trees_dir, states_dir, tree_lineage_dir,
22        trees_dir,
23    },
24    npk1::Npk1Manager,
25};
26use crate::{
27    fs_atomic::sync_directory,
28    object::{Blob, ContentHash, State, StateId, Tree},
29    store::{Result, SnapshotPackManager, pack::PackObjectId},
30};
31
32const RECENT_BLOB_CACHE_CAPACITY: usize = 2_048;
33const RECENT_TREE_CACHE_CAPACITY: usize = 1_024;
34/// Soft cap on the in-process loose-blob verification cache. Each
35/// entry is one `ContentHash` (~32 bytes) so this is ≈2 MB of memory
36/// for the upper bound, and clock eviction is bounded by hash
37/// hits rather than store size. 65k entries covers the typical hot
38/// working set for million-blob monorepos; a daemon that materialises
39/// dozens of unrelated trees won't drift toward unbounded growth.
40const VERIFIED_LOOSE_BLOB_CACHE_CAPACITY: usize = 65_536;
41/// Blobs larger than this are not stored in `recent_blobs` so a single
42/// multi-MB read cannot thrash the hot working set. 4 MiB matches the
43/// typical "large file" boundary used elsewhere in the object path.
44pub(super) const RECENT_BLOB_CACHE_MAX_BYTES: usize = 4 * 1024 * 1024;
45/// Total-byte budget for `recent_blobs`. Without it, populate-on-read
46/// could retain `RECENT_BLOB_CACHE_CAPACITY` (2048) × the 4 MiB
47/// per-entry gate ≈ 8 GiB of deep-cloned blob bytes for a read-only
48/// workload (mount / `heddled`) that streams many cold blobs. 256 MiB
49/// caps the resident blob-cache footprint while still holding a deep
50/// hot working set of small objects (the common case).
51pub(super) const RECENT_BLOB_CACHE_MAX_TOTAL_BYTES: usize = 256 * 1024 * 1024;
52
53thread_local! {
54    static SNAPSHOT_WRITE_BATCH_DEPTHS: std::cell::RefCell<HashMap<PathBuf, usize>> =
55        std::cell::RefCell::new(HashMap::new());
56}
57
58#[derive(Clone, Copy, Debug, Eq, PartialEq)]
59pub enum LooseObjectWriteMode {
60    Durable,
61    BatchDirectorySync,
62}
63
64/// Bounded in-process object cache with second-chance clock eviction.
65///
66/// Two independent caps are enforced on every [`insert`](Self::insert):
67///
68/// * `capacity` — the maximum entry *count*.
69/// * `byte_budget` — a soft cap on the cumulative *bytes* of the
70///   cached values, sized by the per-entry `sizer` closure. `None`
71///   disables the byte cap (caches whose values are effectively
72///   fixed-size, e.g. the `()`-valued verified-loose cache).
73///
74/// The byte budget is what keeps populate-on-read bounded: a read-only
75/// workload (mount / `heddled`) that streams many multi-MB blobs
76/// through `get_blob` can otherwise retain `capacity × max-entry-bytes`
77/// of deep-cloned `Vec`s. With the budget, inserting a new large blob
78/// advances the second-chance clock until the total fits.
79///
80/// [`get`](Self::get) marks the entry recently used through an atomic bit, so
81/// cache hits need only a shared map lock. Eviction advances a clock queue and
82/// gives marked entries one additional chance before removal. Both hits and
83/// amortized eviction stay O(1), including for the 65k-entry verification
84/// cache.
85#[derive(Debug)]
86pub(super) struct RecentObjectCache<K, V> {
87    entries: HashMap<K, RecentObjectCacheEntry<V>>,
88    eviction_clock: VecDeque<K>,
89    capacity: usize,
90    /// Soft cap on cumulative cached bytes; `None` = count-only.
91    byte_budget: Option<usize>,
92    /// `sizer(value)` in bytes. Only consulted when `byte_budget`
93    /// is `Some`.
94    sizer: fn(&V) -> usize,
95    /// Running sum of `sizer(v)` over all `entries`.
96    cached_bytes: usize,
97}
98
99#[derive(Debug)]
100struct RecentObjectCacheEntry<V> {
101    value: V,
102    recently_accessed: AtomicBool,
103}
104
105impl<K, V> RecentObjectCache<K, V>
106where
107    K: Copy + Eq + Hash,
108{
109    /// Count-capped cache with no byte budget. Used for caches whose
110    /// values are effectively fixed-size (e.g. the verified-loose
111    /// marker cache).
112    pub(super) fn with_capacity(capacity: usize) -> Self {
113        Self {
114            entries: HashMap::new(),
115            eviction_clock: VecDeque::new(),
116            capacity,
117            byte_budget: None,
118            sizer: |_| 0,
119            cached_bytes: 0,
120        }
121    }
122
123    /// Cache capped by *both* entry count and cumulative bytes.
124    /// `sizer` reports each value's heap-ish footprint; the cache
125    /// advances the second-chance clock until both caps hold.
126    pub(super) fn with_byte_budget(
127        capacity: usize,
128        byte_budget: usize,
129        sizer: fn(&V) -> usize,
130    ) -> Self {
131        Self {
132            entries: HashMap::new(),
133            eviction_clock: VecDeque::new(),
134            capacity,
135            byte_budget: Some(byte_budget),
136            sizer,
137            cached_bytes: 0,
138        }
139    }
140
141    /// Lookup with lock-free second-chance promotion inside an already-held
142    /// shared map lock. Only insertion and eviction require exclusive access.
143    pub(super) fn get(&self, key: &K) -> Option<&V> {
144        let entry = self.entries.get(key)?;
145        entry.recently_accessed.store(true, Ordering::Relaxed);
146        Some(&entry.value)
147    }
148
149    /// Presence check without promotion. Cheap enough to run under a
150    /// read lock — used both by verified-loose probes and by `has_*`
151    /// existence checks that must not serialize concurrent readers on
152    /// the exclusive write lock a promoting `get` would need.
153    pub(super) fn contains(&self, key: &K) -> bool {
154        self.entries.contains_key(key)
155    }
156
157    /// Drop `key` from the cache entirely. Returns the evicted value if
158    /// present. Targeted counterpart to the redaction-`purge` cache
159    /// drop: a purged blob's bytes must not linger in `recent_blobs`
160    /// where a long-lived process would keep serving (or reporting
161    /// present) the destroyed content. The production purge path drops
162    /// the whole cache via `clear_recent_caches` (it crosses the
163    /// generic `ObjectStore` seam); this per-key variant backs the
164    /// store-level `evict_recent_blob` used in tests.
165    #[cfg(test)]
166    pub(super) fn remove(&mut self, key: &K) -> Option<V> {
167        let removed = self.entries.remove(key)?.value;
168        self.cached_bytes = self.cached_bytes.saturating_sub((self.sizer)(&removed));
169        Some(removed)
170    }
171
172    pub(super) fn insert(&mut self, key: K, value: V) {
173        if self.capacity == 0 {
174            return;
175        }
176        let new_bytes = self.byte_budget.map(|_| (self.sizer)(&value)).unwrap_or(0);
177        let entry = RecentObjectCacheEntry {
178            value,
179            recently_accessed: AtomicBool::new(false),
180        };
181        if let Some(old) = self.entries.insert(key, entry) {
182            self.cached_bytes = self.cached_bytes.saturating_sub(
183                self.byte_budget
184                    .map(|_| (self.sizer)(&old.value))
185                    .unwrap_or(0),
186            );
187        } else {
188            self.eviction_clock.push_back(key);
189        }
190        self.cached_bytes += new_bytes;
191        self.evict_to_fit(key);
192    }
193
194    /// Advance the second-chance clock until both the count cap and the byte
195    /// budget hold. A recently read entry is marked cold and moved to the back
196    /// once before it can be evicted. The freshly inserted entry starts at the
197    /// back, so it is not the first target (a single entry larger than the
198    /// whole budget is kept — the budget is a soft cap, not a hard per-entry
199    /// gate; the per-entry `RECENT_BLOB_CACHE_MAX_BYTES` gate already bounds
200    /// the largest thing that reaches here).
201    fn evict_to_fit(&mut self, admitted_key: K) {
202        loop {
203            let over_count = self.entries.len() > self.capacity;
204            let over_bytes = self
205                .byte_budget
206                .is_some_and(|budget| self.cached_bytes > budget && self.entries.len() > 1);
207            if !over_count && !over_bytes {
208                break;
209            }
210            let Some(candidate) = self.eviction_clock.pop_front() else {
211                break;
212            };
213            let Some(entry) = self.entries.get(&candidate) else {
214                continue;
215            };
216            // The value that triggered this pass has not had an opportunity to
217            // serve a read yet. Keep it for this pass when another victim
218            // exists; otherwise an all-hot full cache would cycle through the
219            // residents and evict the new external object immediately.
220            if candidate == admitted_key && self.entries.len() > 1 {
221                self.eviction_clock.push_back(candidate);
222                continue;
223            }
224            if entry.recently_accessed.swap(false, Ordering::Relaxed) {
225                self.eviction_clock.push_back(candidate);
226                continue;
227            }
228            if let Some(evicted) = self.entries.remove(&candidate) {
229                self.cached_bytes = self.cached_bytes.saturating_sub(
230                    self.byte_budget
231                        .map(|_| (self.sizer)(&evicted.value))
232                        .unwrap_or(0),
233                );
234            }
235        }
236    }
237}
238
239/// Filesystem-based storage for Heddle objects.
240///
241/// Layout:
242/// ```text
243/// .heddle/
244///   objects/
245///     blobs/
246///       ab/
247///         cdef1234...
248///     trees/
249///       ab/
250///         cdef1234...
251///     states/
252///       <state_id>.state
253///   actions/
254///     <action_id>.action
255///   packs/
256///     <hash>.pack
257///     <hash>.idx
258/// ```
259pub struct FsStore {
260    pub(super) root: PathBuf,
261    pub(super) compression: CompressionConfig,
262    pub(super) snapshot_delta_search: bool,
263    pack_manager: RwLock<SnapshotPackManager>,
264    npk1_manager: RwLock<Npk1Manager>,
265    pub(super) recent_blobs: RwLock<RecentObjectCache<ContentHash, Blob>>,
266    pub(super) recent_trees: RwLock<RecentObjectCache<ContentHash, Tree>>,
267    pub(super) recent_states: RwLock<RecentObjectCache<StateId, State>>,
268    pub(super) external_source: Option<Arc<dyn super::super::ExternalObjectSource>>,
269    loose_object_write_mode: LooseObjectWriteMode,
270    pending_directory_syncs: Mutex<BTreeSet<PathBuf>>,
271    #[cfg(test)]
272    snapshot_batch_flushes: AtomicUsize,
273    /// In-process trust cache for loose-blob cache mirrors. A hash
274    /// enters this bounded clock cache when this process either (a) wrote the blob
275    /// itself via `promote_to_loose_uncompressed` or (b) successfully
276    /// hash-verified it on first read. Bytes-on-disk for any entry
277    /// in this cache can be trusted without a re-hash by subsequent
278    /// `loose_blob_path` calls within the same process.
279    ///
280    /// Capped at [`VERIFIED_LOOSE_BLOB_CACHE_CAPACITY`] entries so a
281    /// long-lived process (`heddled`) materialising many unrelated
282    /// trees doesn't drift into unbounded memory growth. Second-chance
283    /// eviction; an evicted hash pays one extra BLAKE3 on its next
284    /// read (cost-of-evict ≈ working-set-size BLAKE3 ops). Stored as
285    /// `RecentObjectCache<…, ()>` to share the clock-eviction
286    /// machinery with the other on-store caches; the unit value is
287    /// a marker that the corresponding loose mirror was verified.
288    ///
289    /// Pairs with `AtomicWriteMode::NoSync` on the write side: a
290    /// crashed promote leaves a torn cache-mirror file, but its
291    /// hash won't match on the next process's first-read verify,
292    /// so the reader falls through to a fresh promote off the pack.
293    pub(super) verified_loose_blobs: RwLock<RecentObjectCache<ContentHash, ()>>,
294}
295
296impl Clone for FsStore {
297    fn clone(&self) -> Self {
298        let mut cloned = Self::with_compression(&self.root, self.compression);
299        cloned.snapshot_delta_search = self.snapshot_delta_search;
300        cloned.loose_object_write_mode = self.loose_object_write_mode;
301        cloned.external_source = self.external_source.clone();
302        cloned
303    }
304}
305
306impl FsStore {
307    /// Create a new filesystem store rooted at the given path.
308    ///
309    /// The path should be the `.heddle` directory.
310    pub fn new(root: impl AsRef<Path>) -> Self {
311        let root = root.as_ref().to_path_buf();
312        crate::store::pack::sweep_scratch(&root.join("tmp"));
313        let pack_manager = SnapshotPackManager::new(packs_dir(&root), root.join("tmp"));
314        let npk1_manager = Npk1Manager::new(packs_dir(&root));
315        Self {
316            root,
317            compression: CompressionConfig::default(),
318            snapshot_delta_search: false,
319            pack_manager: RwLock::new(pack_manager),
320            npk1_manager: RwLock::new(npk1_manager),
321            recent_blobs: RwLock::new(RecentObjectCache::with_byte_budget(
322                RECENT_BLOB_CACHE_CAPACITY,
323                RECENT_BLOB_CACHE_MAX_TOTAL_BYTES,
324                |blob: &Blob| blob.content().len(),
325            )),
326            recent_trees: RwLock::new(RecentObjectCache::with_capacity(RECENT_TREE_CACHE_CAPACITY)),
327            recent_states: RwLock::new(RecentObjectCache::with_capacity(
328                RECENT_TREE_CACHE_CAPACITY,
329            )),
330            external_source: None,
331            loose_object_write_mode: LooseObjectWriteMode::Durable,
332            pending_directory_syncs: Mutex::new(BTreeSet::new()),
333            #[cfg(test)]
334            snapshot_batch_flushes: AtomicUsize::new(0),
335            verified_loose_blobs: RwLock::new(RecentObjectCache::with_capacity(
336                VERIFIED_LOOSE_BLOB_CACHE_CAPACITY,
337            )),
338        }
339    }
340
341    /// Create a new filesystem store with custom compression settings.
342    pub fn with_compression(root: impl AsRef<Path>, compression: CompressionConfig) -> Self {
343        let root = root.as_ref().to_path_buf();
344        crate::store::pack::sweep_scratch(&root.join("tmp"));
345        let pack_manager = SnapshotPackManager::new(packs_dir(&root), root.join("tmp"));
346        let npk1_manager = Npk1Manager::new(packs_dir(&root));
347        Self {
348            root,
349            compression,
350            snapshot_delta_search: false,
351            pack_manager: RwLock::new(pack_manager),
352            npk1_manager: RwLock::new(npk1_manager),
353            recent_blobs: RwLock::new(RecentObjectCache::with_byte_budget(
354                RECENT_BLOB_CACHE_CAPACITY,
355                RECENT_BLOB_CACHE_MAX_TOTAL_BYTES,
356                |blob: &Blob| blob.content().len(),
357            )),
358            recent_trees: RwLock::new(RecentObjectCache::with_capacity(RECENT_TREE_CACHE_CAPACITY)),
359            recent_states: RwLock::new(RecentObjectCache::with_capacity(
360                RECENT_TREE_CACHE_CAPACITY,
361            )),
362            external_source: None,
363            loose_object_write_mode: LooseObjectWriteMode::Durable,
364            pending_directory_syncs: Mutex::new(BTreeSet::new()),
365            #[cfg(test)]
366            snapshot_batch_flushes: AtomicUsize::new(0),
367            verified_loose_blobs: RwLock::new(RecentObjectCache::with_capacity(
368                VERIFIED_LOOSE_BLOB_CACHE_CAPACITY,
369            )),
370        }
371    }
372
373    /// Initialize the directory structure.
374    pub fn init(&self) -> Result<()> {
375        // Durable create so the object-store layout dirs survive crash
376        // between mkdir and first object write (L6 residual migration).
377        crate::fs_atomic::create_dir_all_durable(&blobs_dir(&self.root))?;
378        crate::fs_atomic::create_dir_all_durable(&trees_dir(&self.root))?;
379        crate::fs_atomic::create_dir_all_durable(&partial_trees_dir(&self.root))?;
380        crate::fs_atomic::create_dir_all_durable(&tree_lineage_dir(&self.root))?;
381        crate::fs_atomic::create_dir_all_durable(&states_dir(&self.root))?;
382        crate::fs_atomic::create_dir_all_durable(&actions_dir(&self.root))?;
383        crate::fs_atomic::create_dir_all_durable(&packs_dir(&self.root))?;
384        Ok(())
385    }
386
387    /// Get the root path.
388    pub fn root(&self) -> &Path {
389        &self.root
390    }
391
392    /// Get the compression configuration.
393    pub fn compression(&self) -> CompressionConfig {
394        self.compression
395    }
396
397    /// Set the compression configuration.
398    pub fn set_compression(&mut self, compression: CompressionConfig) {
399        self.compression = compression;
400    }
401
402    /// Enable or disable sliding-window delta search for snapshot packs.
403    pub fn set_snapshot_delta_search(&mut self, enabled: bool) {
404        self.snapshot_delta_search = enabled;
405    }
406
407    pub fn loose_object_write_mode(&self) -> LooseObjectWriteMode {
408        self.loose_object_write_mode
409    }
410
411    pub fn set_loose_object_write_mode(&mut self, mode: LooseObjectWriteMode) {
412        self.loose_object_write_mode = mode;
413    }
414
415    /// Configure a read-through source for objects not present in the native
416    /// store. Writes always remain native.
417    pub fn set_external_source(&mut self, source: Arc<dyn super::super::ExternalObjectSource>) {
418        self.external_source = Some(source);
419    }
420
421    fn flush_pending_directory_syncs(&self) -> Result<usize> {
422        let pending_dirs = {
423            let mut guard = self.pending_directory_syncs.lock().map_err(|_| {
424                crate::store::HeddleError::Config(
425                    "Failed to acquire pending directory sync lock".to_string(),
426                )
427            })?;
428            if guard.is_empty() {
429                return Ok(0);
430            }
431            let dirs = guard.iter().cloned().collect::<Vec<_>>();
432            guard.clear();
433            dirs
434        };
435
436        for (index, dir) in pending_dirs.iter().enumerate() {
437            if let Err(error) = sync_directory(dir) {
438                if let Ok(mut guard) = self.pending_directory_syncs.lock() {
439                    guard.extend(pending_dirs[index..].iter().cloned());
440                }
441                return Err(error.into());
442            }
443        }
444
445        Ok(pending_dirs.len())
446    }
447
448    /// Reload pack files from disk.
449    ///
450    /// Runs L8 install-intent recovery first so crash windows between pack
451    /// and index publish are finished or aborted before packs are loaded.
452    /// Uses the default intent TTL so abandoned staging is swept.
453    pub fn reload_packs(&self) -> Result<()> {
454        crate::store::pack::sweep_scratch(&self.root.join("tmp"));
455        let packs = packs_dir(&self.root);
456        let _ = super::pack_install_journal::recover_pack_install_intents_with_ttl(
457            &packs,
458            Some(super::pack_install_journal::DEFAULT_PACK_INSTALL_INTENT_TTL_SECS),
459        )?;
460        // Option D backstop: remove any legacy unpaired packs without intent.
461        let _ = super::fs_pack::prune_unpaired_pack_files(&packs)?;
462        let mut manager = self.pack_manager.write().map_err(|_| {
463            crate::store::HeddleError::Config("Failed to acquire pack manager lock".to_string())
464        })?;
465        manager.reload()?;
466        drop(manager);
467        let mut npk1 = self.npk1_manager.write().map_err(|_| {
468            crate::store::HeddleError::Config("Failed to acquire NPK1 manager lock".to_string())
469        })?;
470        npk1.reload()
471    }
472
473    /// Reload pack files only if the immutable pack set changed on disk.
474    /// Cheap discovery when nothing changed; full reload when a sibling
475    /// `FsStore` installed a pack or atomically replaced a generation.
476    ///
477    /// Returns `true` when a reload happened. Used by `get_*` and
478    /// `has_*` paths after an in-memory miss to recover from the
479    /// "two FsStores backing the same `.heddle/` directory" case
480    /// (typical for lightweight thread worktrees).
481    ///
482    /// Double-checked locking: the read-lock fast path means a
483    /// thundering herd of concurrent misses doesn't serialize on
484    /// the write lock; only the first thread that observes a stale
485    /// view escalates and does the reload.
486    pub(super) fn reload_packs_if_stale(&self) -> Result<bool> {
487        // Fast path: read-lock and bail out if the disk snapshot still matches.
488        let generic_stale = {
489            let manager = self.pack_manager.read().map_err(|_| {
490                crate::store::HeddleError::Config("Failed to acquire pack manager lock".to_string())
491            })?;
492            manager.needs_reload()?
493        };
494        let npk1_stale = {
495            let manager = self.npk1_manager.read().map_err(|_| {
496                crate::store::HeddleError::Config("Failed to acquire NPK1 manager lock".to_string())
497            })?;
498            manager.needs_reload()?
499        };
500        if !generic_stale && !npk1_stale {
501            return Ok(false);
502        }
503        // Slow path: take the write lock and re-check (another
504        // thread may have already reloaded between our drop and
505        // re-acquire).
506        let mut manager = self.pack_manager.write().map_err(|_| {
507            crate::store::HeddleError::Config("Failed to acquire pack manager lock".to_string())
508        })?;
509        let generic_reloaded = manager.reload_if_stale()?;
510        drop(manager);
511        let mut npk1 = self.npk1_manager.write().map_err(|_| {
512            crate::store::HeddleError::Config("Failed to acquire NPK1 manager lock".to_string())
513        })?;
514        let npk1_reloaded = if npk1.needs_reload()? {
515            npk1.reload()?;
516            true
517        } else {
518            false
519        };
520        Ok(generic_reloaded || npk1_reloaded)
521    }
522
523    /// Get the pack manager for pack operations.
524    pub fn pack_manager(&self) -> &RwLock<SnapshotPackManager> {
525        &self.pack_manager
526    }
527
528    pub(super) fn npk1_manager(&self) -> &RwLock<Npk1Manager> {
529        &self.npk1_manager
530    }
531
532    pub fn clear_recent_object_caches(&self) {
533        if let Ok(mut blobs) = self.recent_blobs.write() {
534            *blobs = RecentObjectCache::with_byte_budget(
535                RECENT_BLOB_CACHE_CAPACITY,
536                RECENT_BLOB_CACHE_MAX_TOTAL_BYTES,
537                |blob: &Blob| blob.content().len(),
538            );
539        }
540        if let Ok(mut trees) = self.recent_trees.write() {
541            *trees = RecentObjectCache::with_capacity(RECENT_TREE_CACHE_CAPACITY);
542        }
543        if let Ok(mut states) = self.recent_states.write() {
544            *states = RecentObjectCache::with_capacity(RECENT_TREE_CACHE_CAPACITY);
545        }
546    }
547
548    /// Drop a single blob hash from the in-process `recent_blobs`
549    /// cache. Targeted counterpart to the redaction-`purge` cache drop:
550    /// after the loose bytes are physically deleted, a long-lived
551    /// process must not keep serving (or reporting present) the purged
552    /// content from cache. Idempotent — a miss is a no-op. Test-only:
553    /// the production purge path crosses the generic `ObjectStore` seam
554    /// and drops the whole cache via `clear_recent_caches`.
555    #[cfg(test)]
556    pub(super) fn evict_recent_blob(&self, hash: &ContentHash) {
557        if let Ok(mut cache) = self.recent_blobs.write() {
558            cache.remove(hash);
559        }
560    }
561
562    pub fn pack_ids(&self) -> Result<Vec<PackObjectId>> {
563        let manager = self.pack_manager.read().map_err(|_| {
564            crate::store::HeddleError::Config("Failed to acquire pack manager lock".to_string())
565        })?;
566        let mut ids = manager.list_all_ids()?;
567        drop(manager);
568        let npk1 = self.npk1_manager.read().map_err(|_| {
569            crate::store::HeddleError::Config("Failed to acquire NPK1 manager lock".to_string())
570        })?;
571        ids.extend(npk1.list_ids()?.into_iter().map(PackObjectId::Hash));
572        ids.sort();
573        ids.dedup();
574        Ok(ids)
575    }
576
577    pub(super) fn write_loose_object_atomic(&self, path: &Path, data: &[u8]) -> Result<()> {
578        let batch_active = SNAPSHOT_WRITE_BATCH_DEPTHS
579            .with(|depths| depths.borrow().get(&self.root).copied().unwrap_or_default() > 0);
580        let configured_mode = if batch_active {
581            LooseObjectWriteMode::BatchDirectorySync
582        } else {
583            self.loose_object_write_mode
584        };
585
586        let mode = match configured_mode {
587            LooseObjectWriteMode::Durable => AtomicWriteMode::Durable,
588            LooseObjectWriteMode::BatchDirectorySync => AtomicWriteMode::BatchDirectorySync,
589        };
590        write_atomic(path, data, mode, Some(&self.pending_directory_syncs))
591    }
592
593    /// Durable atomic write for pack/index bytes when not going through the
594    /// L8 journal (tests / rare call sites). Prefer
595    /// [`super::pack_install_journal::install_pack_bytes_journaled`].
596    #[allow(dead_code)]
597    pub(super) fn write_pack_atomic(&self, path: &Path, data: &[u8]) -> Result<()> {
598        write_atomic(path, data, AtomicWriteMode::Durable, None)
599    }
600
601    /// Atomic write tuned for *cache-mirror* loose objects: no fsync
602    /// at any level. The authoritative copy lives in a pack; if a
603    /// crash leaves the cache mirror torn, the read-side hash check
604    /// catches it and `promote_to_loose_uncompressed` rebuilds it
605    /// from the pack on the next access.
606    ///
607    /// On macOS APFS, `sync_data` alone costs ~5 ms per call (it
608    /// behaves like `F_FULLFSYNC` for tiny writes), and the parent
609    /// directory fsync is ~3-10 ms on top. For 1k blobs, that's
610    /// 5-15 seconds of pure fsync wallclock — the dominant cost in
611    /// the cold materialize path. Dropping both pays back ~30× on
612    /// raw create+rename throughput (measured: 200/s with sync_data
613    /// vs 5500/s without).
614    ///
615    /// Safety contract: this is only valid for files whose authority
616    /// lives elsewhere. Used by `promote_to_loose_uncompressed`; the
617    /// matching `loose_blob_path` reader hash-verifies before
618    /// trusting the bytes. Do *not* use for `put_blob` / `put_tree`
619    /// / `put_state` — those are the authoritative copy and must
620    /// survive a crash.
621    pub(super) fn write_loose_object_cache(&self, path: &Path, data: &[u8]) -> Result<()> {
622        self.write_reconstructible_cache(path, data)
623    }
624
625    /// Atomically publish reconstructible cache bytes without a durability
626    /// barrier. The caller must be able to rebuild the file from an
627    /// authoritative object after a crash.
628    pub(super) fn write_reconstructible_cache(&self, path: &Path, data: &[u8]) -> Result<()> {
629        write_atomic(path, data, AtomicWriteMode::NoSync, None)
630    }
631
632    pub(super) fn begin_snapshot_write_batch_impl(&self) -> Result<()> {
633        SNAPSHOT_WRITE_BATCH_DEPTHS.with(|depths| {
634            *depths.borrow_mut().entry(self.root.clone()).or_default() += 1;
635        });
636        Ok(())
637    }
638
639    pub(super) fn flush_snapshot_write_batch_impl(&self) -> Result<()> {
640        let had_batch = SNAPSHOT_WRITE_BATCH_DEPTHS.with(|depths| {
641            let mut depths = depths.borrow_mut();
642            let Some(depth) = depths.get_mut(&self.root) else {
643                return false;
644            };
645            *depth -= 1;
646            if *depth == 0 {
647                depths.remove(&self.root);
648            }
649            true
650        });
651        if !had_batch {
652            return Ok(());
653        }
654
655        #[cfg(test)]
656        self.snapshot_batch_flushes.fetch_add(1, Ordering::Relaxed);
657
658        // Batches may overlap across snapshot preparers. Each successful
659        // preparer must establish durability for its own writes before it can
660        // publish an oplog edge, even while another batch remains active.
661        // Draining the shared set is safe: entries taken by another flush are
662        // already durable, and every write from this batch was queued before
663        // this call acquired the set.
664        let _ = self.flush_pending_directory_syncs()?;
665        Ok(())
666    }
667
668    pub(super) fn abort_snapshot_write_batch_impl(&self) {
669        let should_flush = SNAPSHOT_WRITE_BATCH_DEPTHS.with(|depths| {
670            let mut depths = depths.borrow_mut();
671            let Some(depth) = depths.get_mut(&self.root) else {
672                // A preceding flush may have removed the thread-local batch
673                // before its directory sync failed. Preserve abort's
674                // conservative retry of those pending syncs.
675                return true;
676            };
677            *depth -= 1;
678            if *depth == 0 {
679                depths.remove(&self.root);
680                true
681            } else {
682                false
683            }
684        });
685        // Immutable objects staged by a failed snapshot are harmless orphans.
686        // Never clear another concurrent preparation's pending directory syncs;
687        // when this was the last batch, conservatively make every staged rename
688        // durable before returning.
689        if should_flush {
690            let _ = self.flush_pending_directory_syncs();
691        }
692    }
693
694    #[cfg(test)]
695    pub(super) fn pending_directory_sync_count(&self) -> usize {
696        self.pending_directory_syncs
697            .lock()
698            .map(|pending| pending.len())
699            .unwrap_or(0)
700    }
701
702    #[cfg(test)]
703    pub(super) fn snapshot_batch_flush_count(&self) -> usize {
704        self.snapshot_batch_flushes.load(Ordering::Relaxed)
705    }
706}