1#[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;
34const VERIFIED_LOOSE_BLOB_CACHE_CAPACITY: usize = 65_536;
41pub(super) const RECENT_BLOB_CACHE_MAX_BYTES: usize = 4 * 1024 * 1024;
45pub(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#[derive(Debug)]
86pub(super) struct RecentObjectCache<K, V> {
87 entries: HashMap<K, RecentObjectCacheEntry<V>>,
88 eviction_clock: VecDeque<K>,
89 capacity: usize,
90 byte_budget: Option<usize>,
92 sizer: fn(&V) -> usize,
95 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 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 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 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 pub(super) fn contains(&self, key: &K) -> bool {
154 self.entries.contains_key(key)
155 }
156
157 #[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 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 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
239pub 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 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 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 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 pub fn init(&self) -> Result<()> {
375 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 pub fn root(&self) -> &Path {
389 &self.root
390 }
391
392 pub fn compression(&self) -> CompressionConfig {
394 self.compression
395 }
396
397 pub fn set_compression(&mut self, compression: CompressionConfig) {
399 self.compression = compression;
400 }
401
402 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 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 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 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 pub(super) fn reload_packs_if_stale(&self) -> Result<bool> {
487 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 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 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 #[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 #[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 pub(super) fn write_loose_object_cache(&self, path: &Path, data: &[u8]) -> Result<()> {
622 self.write_reconstructible_cache(path, data)
623 }
624
625 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 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 return true;
676 };
677 *depth -= 1;
678 if *depth == 0 {
679 depths.remove(&self.root);
680 true
681 } else {
682 false
683 }
684 });
685 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}