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 let pack_manager = SnapshotPackManager::new(packs_dir(&root));
313 let npk1_manager = Npk1Manager::new(packs_dir(&root));
314 Self {
315 root,
316 compression: CompressionConfig::default(),
317 snapshot_delta_search: false,
318 pack_manager: RwLock::new(pack_manager),
319 npk1_manager: RwLock::new(npk1_manager),
320 recent_blobs: RwLock::new(RecentObjectCache::with_byte_budget(
321 RECENT_BLOB_CACHE_CAPACITY,
322 RECENT_BLOB_CACHE_MAX_TOTAL_BYTES,
323 |blob: &Blob| blob.content().len(),
324 )),
325 recent_trees: RwLock::new(RecentObjectCache::with_capacity(RECENT_TREE_CACHE_CAPACITY)),
326 recent_states: RwLock::new(RecentObjectCache::with_capacity(
327 RECENT_TREE_CACHE_CAPACITY,
328 )),
329 external_source: None,
330 loose_object_write_mode: LooseObjectWriteMode::Durable,
331 pending_directory_syncs: Mutex::new(BTreeSet::new()),
332 #[cfg(test)]
333 snapshot_batch_flushes: AtomicUsize::new(0),
334 verified_loose_blobs: RwLock::new(RecentObjectCache::with_capacity(
335 VERIFIED_LOOSE_BLOB_CACHE_CAPACITY,
336 )),
337 }
338 }
339
340 pub fn with_compression(root: impl AsRef<Path>, compression: CompressionConfig) -> Self {
342 let root = root.as_ref().to_path_buf();
343 let pack_manager = SnapshotPackManager::new(packs_dir(&root));
344 let npk1_manager = Npk1Manager::new(packs_dir(&root));
345 Self {
346 root,
347 compression,
348 snapshot_delta_search: false,
349 pack_manager: RwLock::new(pack_manager),
350 npk1_manager: RwLock::new(npk1_manager),
351 recent_blobs: RwLock::new(RecentObjectCache::with_byte_budget(
352 RECENT_BLOB_CACHE_CAPACITY,
353 RECENT_BLOB_CACHE_MAX_TOTAL_BYTES,
354 |blob: &Blob| blob.content().len(),
355 )),
356 recent_trees: RwLock::new(RecentObjectCache::with_capacity(RECENT_TREE_CACHE_CAPACITY)),
357 recent_states: RwLock::new(RecentObjectCache::with_capacity(
358 RECENT_TREE_CACHE_CAPACITY,
359 )),
360 external_source: None,
361 loose_object_write_mode: LooseObjectWriteMode::Durable,
362 pending_directory_syncs: Mutex::new(BTreeSet::new()),
363 #[cfg(test)]
364 snapshot_batch_flushes: AtomicUsize::new(0),
365 verified_loose_blobs: RwLock::new(RecentObjectCache::with_capacity(
366 VERIFIED_LOOSE_BLOB_CACHE_CAPACITY,
367 )),
368 }
369 }
370
371 pub fn init(&self) -> Result<()> {
373 crate::fs_atomic::create_dir_all_durable(&blobs_dir(&self.root))?;
376 crate::fs_atomic::create_dir_all_durable(&trees_dir(&self.root))?;
377 crate::fs_atomic::create_dir_all_durable(&partial_trees_dir(&self.root))?;
378 crate::fs_atomic::create_dir_all_durable(&tree_lineage_dir(&self.root))?;
379 crate::fs_atomic::create_dir_all_durable(&states_dir(&self.root))?;
380 crate::fs_atomic::create_dir_all_durable(&actions_dir(&self.root))?;
381 crate::fs_atomic::create_dir_all_durable(&packs_dir(&self.root))?;
382 Ok(())
383 }
384
385 pub fn root(&self) -> &Path {
387 &self.root
388 }
389
390 pub fn compression(&self) -> CompressionConfig {
392 self.compression
393 }
394
395 pub fn set_compression(&mut self, compression: CompressionConfig) {
397 self.compression = compression;
398 }
399
400 pub fn set_snapshot_delta_search(&mut self, enabled: bool) {
402 self.snapshot_delta_search = enabled;
403 }
404
405 pub fn loose_object_write_mode(&self) -> LooseObjectWriteMode {
406 self.loose_object_write_mode
407 }
408
409 pub fn set_loose_object_write_mode(&mut self, mode: LooseObjectWriteMode) {
410 self.loose_object_write_mode = mode;
411 }
412
413 pub fn set_external_source(&mut self, source: Arc<dyn super::super::ExternalObjectSource>) {
416 self.external_source = Some(source);
417 }
418
419 fn flush_pending_directory_syncs(&self) -> Result<usize> {
420 let pending_dirs = {
421 let mut guard = self.pending_directory_syncs.lock().map_err(|_| {
422 crate::store::HeddleError::Config(
423 "Failed to acquire pending directory sync lock".to_string(),
424 )
425 })?;
426 if guard.is_empty() {
427 return Ok(0);
428 }
429 let dirs = guard.iter().cloned().collect::<Vec<_>>();
430 guard.clear();
431 dirs
432 };
433
434 for (index, dir) in pending_dirs.iter().enumerate() {
435 if let Err(error) = sync_directory(dir) {
436 if let Ok(mut guard) = self.pending_directory_syncs.lock() {
437 guard.extend(pending_dirs[index..].iter().cloned());
438 }
439 return Err(error.into());
440 }
441 }
442
443 Ok(pending_dirs.len())
444 }
445
446 pub fn reload_packs(&self) -> Result<()> {
452 let packs = packs_dir(&self.root);
453 let _ = super::pack_install_journal::recover_pack_install_intents_with_ttl(
454 &packs,
455 Some(super::pack_install_journal::DEFAULT_PACK_INSTALL_INTENT_TTL_SECS),
456 )?;
457 let _ = super::fs_pack::prune_unpaired_pack_files(&packs)?;
459 let mut manager = self.pack_manager.write().map_err(|_| {
460 crate::store::HeddleError::Config("Failed to acquire pack manager lock".to_string())
461 })?;
462 manager.reload()?;
463 drop(manager);
464 let mut npk1 = self.npk1_manager.write().map_err(|_| {
465 crate::store::HeddleError::Config("Failed to acquire NPK1 manager lock".to_string())
466 })?;
467 npk1.reload()
468 }
469
470 pub(super) fn reload_packs_if_stale(&self) -> Result<bool> {
484 let generic_stale = {
486 let manager = self.pack_manager.read().map_err(|_| {
487 crate::store::HeddleError::Config("Failed to acquire pack manager lock".to_string())
488 })?;
489 manager.needs_reload()?
490 };
491 let npk1_stale = {
492 let manager = self.npk1_manager.read().map_err(|_| {
493 crate::store::HeddleError::Config("Failed to acquire NPK1 manager lock".to_string())
494 })?;
495 manager.needs_reload()?
496 };
497 if !generic_stale && !npk1_stale {
498 return Ok(false);
499 }
500 let mut manager = self.pack_manager.write().map_err(|_| {
504 crate::store::HeddleError::Config("Failed to acquire pack manager lock".to_string())
505 })?;
506 let generic_reloaded = manager.reload_if_stale()?;
507 drop(manager);
508 let mut npk1 = self.npk1_manager.write().map_err(|_| {
509 crate::store::HeddleError::Config("Failed to acquire NPK1 manager lock".to_string())
510 })?;
511 let npk1_reloaded = if npk1.needs_reload()? {
512 npk1.reload()?;
513 true
514 } else {
515 false
516 };
517 Ok(generic_reloaded || npk1_reloaded)
518 }
519
520 pub fn pack_manager(&self) -> &RwLock<SnapshotPackManager> {
522 &self.pack_manager
523 }
524
525 pub(super) fn npk1_manager(&self) -> &RwLock<Npk1Manager> {
526 &self.npk1_manager
527 }
528
529 pub fn clear_recent_object_caches(&self) {
530 if let Ok(mut blobs) = self.recent_blobs.write() {
531 *blobs = RecentObjectCache::with_byte_budget(
532 RECENT_BLOB_CACHE_CAPACITY,
533 RECENT_BLOB_CACHE_MAX_TOTAL_BYTES,
534 |blob: &Blob| blob.content().len(),
535 );
536 }
537 if let Ok(mut trees) = self.recent_trees.write() {
538 *trees = RecentObjectCache::with_capacity(RECENT_TREE_CACHE_CAPACITY);
539 }
540 if let Ok(mut states) = self.recent_states.write() {
541 *states = RecentObjectCache::with_capacity(RECENT_TREE_CACHE_CAPACITY);
542 }
543 }
544
545 #[cfg(test)]
553 pub(super) fn evict_recent_blob(&self, hash: &ContentHash) {
554 if let Ok(mut cache) = self.recent_blobs.write() {
555 cache.remove(hash);
556 }
557 }
558
559 pub fn pack_ids(&self) -> Result<Vec<PackObjectId>> {
560 let manager = self.pack_manager.read().map_err(|_| {
561 crate::store::HeddleError::Config("Failed to acquire pack manager lock".to_string())
562 })?;
563 let mut ids = manager.list_all_ids()?;
564 drop(manager);
565 let npk1 = self.npk1_manager.read().map_err(|_| {
566 crate::store::HeddleError::Config("Failed to acquire NPK1 manager lock".to_string())
567 })?;
568 ids.extend(npk1.list_ids()?.into_iter().map(PackObjectId::Hash));
569 ids.sort();
570 ids.dedup();
571 Ok(ids)
572 }
573
574 pub(super) fn write_loose_object_atomic(&self, path: &Path, data: &[u8]) -> Result<()> {
575 let batch_active = SNAPSHOT_WRITE_BATCH_DEPTHS
576 .with(|depths| depths.borrow().get(&self.root).copied().unwrap_or_default() > 0);
577 let configured_mode = if batch_active {
578 LooseObjectWriteMode::BatchDirectorySync
579 } else {
580 self.loose_object_write_mode
581 };
582
583 let mode = match configured_mode {
584 LooseObjectWriteMode::Durable => AtomicWriteMode::Durable,
585 LooseObjectWriteMode::BatchDirectorySync => AtomicWriteMode::BatchDirectorySync,
586 };
587 write_atomic(path, data, mode, Some(&self.pending_directory_syncs))
588 }
589
590 #[allow(dead_code)]
594 pub(super) fn write_pack_atomic(&self, path: &Path, data: &[u8]) -> Result<()> {
595 write_atomic(path, data, AtomicWriteMode::Durable, None)
596 }
597
598 pub(super) fn write_loose_object_cache(&self, path: &Path, data: &[u8]) -> Result<()> {
619 self.write_reconstructible_cache(path, data)
620 }
621
622 pub(super) fn write_reconstructible_cache(&self, path: &Path, data: &[u8]) -> Result<()> {
626 write_atomic(path, data, AtomicWriteMode::NoSync, None)
627 }
628
629 pub(super) fn begin_snapshot_write_batch_impl(&self) -> Result<()> {
630 SNAPSHOT_WRITE_BATCH_DEPTHS.with(|depths| {
631 *depths.borrow_mut().entry(self.root.clone()).or_default() += 1;
632 });
633 Ok(())
634 }
635
636 pub(super) fn flush_snapshot_write_batch_impl(&self) -> Result<()> {
637 let had_batch = SNAPSHOT_WRITE_BATCH_DEPTHS.with(|depths| {
638 let mut depths = depths.borrow_mut();
639 let Some(depth) = depths.get_mut(&self.root) else {
640 return false;
641 };
642 *depth -= 1;
643 if *depth == 0 {
644 depths.remove(&self.root);
645 }
646 true
647 });
648 if !had_batch {
649 return Ok(());
650 }
651
652 #[cfg(test)]
653 self.snapshot_batch_flushes.fetch_add(1, Ordering::Relaxed);
654
655 let _ = self.flush_pending_directory_syncs()?;
662 Ok(())
663 }
664
665 pub(super) fn abort_snapshot_write_batch_impl(&self) {
666 let should_flush = SNAPSHOT_WRITE_BATCH_DEPTHS.with(|depths| {
667 let mut depths = depths.borrow_mut();
668 let Some(depth) = depths.get_mut(&self.root) else {
669 return true;
673 };
674 *depth -= 1;
675 if *depth == 0 {
676 depths.remove(&self.root);
677 true
678 } else {
679 false
680 }
681 });
682 if should_flush {
687 let _ = self.flush_pending_directory_syncs();
688 }
689 }
690
691 #[cfg(test)]
692 pub(super) fn pending_directory_sync_count(&self) -> usize {
693 self.pending_directory_syncs
694 .lock()
695 .map(|pending| pending.len())
696 .unwrap_or(0)
697 }
698
699 #[cfg(test)]
700 pub(super) fn snapshot_batch_flush_count(&self) -> usize {
701 self.snapshot_batch_flushes.load(Ordering::Relaxed)
702 }
703}