Skip to main content

heddle_pack/store/pack/
manager.rs

1// SPDX-License-Identifier: Apache-2.0
2//! Pack file manager for coordinating multiple pack files.
3
4use std::{
5    collections::HashMap,
6    fs,
7    path::{Path, PathBuf},
8    sync::{OnceLock, RwLock},
9    time::SystemTime,
10};
11
12use tracing::{debug, instrument, trace};
13
14use crate::{
15    object::ContentHash,
16    store::{
17        Result,
18        pack::{ObjectType, PackObjectId, PackReadTier, PackReader},
19    },
20};
21
22/// Format-only coordinator for loaded pack and index files.
23///
24/// Object-domain indexes belong in a wrapper owned by the consuming crate.
25pub struct PackManager {
26    packs_dir: PathBuf,
27    packs: Vec<CachedPack>,
28    scratch_root: PathBuf,
29    object_locations: RwLock<ObjectLocationIndex>,
30    eager_object_locations: bool,
31}
32
33#[derive(Default)]
34struct ObjectLocationIndex {
35    locations: HashMap<PackObjectId, ObjectLocation>,
36    complete: bool,
37}
38
39#[derive(Clone, Copy)]
40struct ObjectLocation {
41    pack_index: usize,
42    tier: PackReadTier,
43}
44
45struct CachedPack {
46    pack_path: PathBuf,
47    index_path: PathBuf,
48    reader: OnceLock<Option<PackReader<'static>>>,
49    scratch_root: PathBuf,
50}
51
52impl CachedPack {
53    fn discovered(pack_path: PathBuf, index_path: PathBuf, scratch_root: PathBuf) -> Self {
54        Self {
55            pack_path,
56            index_path,
57            reader: OnceLock::new(),
58            scratch_root,
59        }
60    }
61
62    fn validated(
63        pack_path: PathBuf,
64        index_path: PathBuf,
65        reader: PackReader<'static>,
66        scratch_root: PathBuf,
67    ) -> Self {
68        Self {
69            pack_path,
70            index_path,
71            reader: OnceLock::from(Some(reader)),
72            scratch_root,
73        }
74    }
75
76    fn reader(&self) -> Option<&PackReader<'static>> {
77        self.reader
78            .get_or_init(|| {
79                match PackReader::open_lazy(&self.pack_path, &self.index_path, &self.scratch_root) {
80                    Ok(reader) => Some(reader),
81                    Err(error) => {
82                        debug!(pack = ?self.pack_path, %error, "Failed to open pack");
83                        None
84                    }
85                }
86            })
87            .as_ref()
88    }
89
90    fn verified_reader(&self) -> Option<&PackReader<'static>> {
91        self.reader
92            .get_or_init(|| {
93                match PackReader::open(&self.pack_path, &self.index_path, &self.scratch_root) {
94                    Ok(reader) => Some(reader),
95                    Err(error) => {
96                        debug!(pack = ?self.pack_path, %error, "Failed to open pack");
97                        None
98                    }
99                }
100            })
101            .as_ref()
102    }
103}
104
105impl PackManager {
106    pub fn new(packs_dir: PathBuf, scratch_root: PathBuf) -> Self {
107        Self::new_with_scratch(packs_dir, scratch_root, force_eager_pack_index())
108    }
109
110    #[cfg(test)]
111    fn new_with_index_mode(packs_dir: PathBuf, eager_object_locations: bool) -> Self {
112        let scratch_root = packs_dir.join("tmp");
113        Self::new_with_scratch(packs_dir, scratch_root, eager_object_locations)
114    }
115
116    fn new_with_scratch(
117        packs_dir: PathBuf,
118        scratch_root: PathBuf,
119        eager_object_locations: bool,
120    ) -> Self {
121        let packs = Self::load_packs(&packs_dir, &scratch_root).unwrap_or_default();
122        let object_locations = Self::initial_object_locations(&packs, eager_object_locations);
123        Self {
124            packs_dir,
125            packs,
126            scratch_root,
127            object_locations: RwLock::new(object_locations),
128            eager_object_locations,
129        }
130    }
131
132    fn discover_pack_paths(packs_dir: &Path) -> Result<Vec<(PathBuf, PathBuf)>> {
133        let mut packs = Vec::new();
134
135        if !packs_dir.exists() {
136            return Ok(packs);
137        }
138
139        for entry in fs::read_dir(packs_dir)? {
140            let entry = entry?;
141            let path = entry.path();
142
143            if path.extension().map(|e| e == "pack").unwrap_or(false) {
144                let index_path = path.with_extension("idx");
145                if index_path.exists() {
146                    packs.push((path, index_path));
147                }
148            }
149        }
150
151        // Pack names are content hashes, so lexical order says nothing about
152        // recency. Keep the oldest pack first: point lookups walk this vector
153        // backwards and current snapshot trees overwhelmingly live in the
154        // newest incremental pack. The path is a deterministic tie-breaker
155        // for filesystems with coarse timestamp precision.
156        packs.sort_by(|left, right| {
157            pack_modified(&left.0)
158                .cmp(&pack_modified(&right.0))
159                .then_with(|| left.0.cmp(&right.0))
160        });
161
162        debug!(count = packs.len(), "Discovered pack files");
163        Ok(packs)
164    }
165
166    fn load_packs(packs_dir: &Path, scratch_root: &Path) -> Result<Vec<CachedPack>> {
167        Ok(Self::discover_pack_paths(packs_dir)?
168            .into_iter()
169            .map(|(pack_path, index_path)| {
170                CachedPack::discovered(pack_path, index_path, scratch_root.to_path_buf())
171            })
172            .collect())
173    }
174
175    pub fn reload(&mut self) -> Result<()> {
176        self.packs = Self::load_packs(&self.packs_dir, &self.scratch_root)?;
177        self.reset_object_locations();
178        Ok(())
179    }
180
181    fn initial_object_locations(packs: &[CachedPack], eager: bool) -> ObjectLocationIndex {
182        if !eager {
183            return ObjectLocationIndex::default();
184        }
185        let mut locations = HashMap::new();
186        for (pack_index, pack) in packs.iter().enumerate() {
187            let Some(reader) = pack.verified_reader() else {
188                continue;
189            };
190            let Ok(objects) = reader.indexed_read_tiers() else {
191                continue;
192            };
193            for (id, tier) in objects {
194                remember_location(&mut locations, id, pack_index, tier);
195            }
196        }
197        ObjectLocationIndex {
198            locations,
199            complete: true,
200        }
201    }
202
203    fn reset_object_locations(&mut self) {
204        self.object_locations = RwLock::new(Self::initial_object_locations(
205            &self.packs,
206            self.eager_object_locations,
207        ));
208    }
209
210    fn object_location(&self, id: &PackObjectId) -> Result<Option<usize>> {
211        {
212            let index = self
213                .object_locations
214                .read()
215                .unwrap_or_else(std::sync::PoisonError::into_inner);
216            if let Some(location) = index.locations.get(id) {
217                return Ok(Some(location.pack_index));
218            }
219            if index.complete {
220                return Ok(None);
221            }
222        }
223
224        let mut index = self
225            .object_locations
226            .write()
227            .unwrap_or_else(std::sync::PoisonError::into_inner);
228        if !index.complete {
229            for (pack_index, pack) in self.packs.iter().enumerate() {
230                let Some(reader) = pack.reader() else {
231                    continue;
232                };
233                let Ok(objects) = reader.indexed_read_tiers() else {
234                    continue;
235                };
236                for (object_id, tier) in objects {
237                    remember_location(&mut index.locations, object_id, pack_index, tier);
238                }
239            }
240            index.complete = true;
241        }
242        Ok(index.locations.get(id).map(|location| location.pack_index))
243    }
244
245    /// Locate one object through each pack's sorted index without building the
246    /// cross-pack map. Newer packs are checked first because incremental
247    /// snapshots usually place the current state/tree chain in the newest pack.
248    fn point_object_location(&self, id: &PackObjectId) -> Result<Option<usize>> {
249        {
250            let index = self
251                .object_locations
252                .read()
253                .unwrap_or_else(std::sync::PoisonError::into_inner);
254            if let Some(location) = index.locations.get(id) {
255                return Ok(Some(location.pack_index));
256            }
257            if index.complete {
258                return Ok(None);
259            }
260        }
261        for (pack_index, pack) in self.packs.iter().enumerate().rev() {
262            let Some(reader) = pack.reader() else {
263                continue;
264            };
265            if reader.contains_object(id)? {
266                return Ok(Some(pack_index));
267            }
268        }
269        Ok(None)
270    }
271
272    /// Add a complete pack/index pair to the in-memory format index.
273    pub fn add_pack(&mut self, pack_path: PathBuf, index_path: PathBuf) -> Result<()> {
274        if self.packs.iter().any(|pack| pack.pack_path == pack_path) {
275            return Ok(());
276        }
277        let reader = PackReader::open(&pack_path, &index_path, &self.scratch_root)?;
278        let pack_index = self.packs.len();
279        let cached =
280            CachedPack::validated(pack_path, index_path, reader, self.scratch_root.clone());
281        self.packs.push(cached);
282        let mut index = self
283            .object_locations
284            .write()
285            .unwrap_or_else(std::sync::PoisonError::into_inner);
286        if index.complete {
287            let objects = self.packs[pack_index]
288                .reader()
289                .ok_or_else(|| {
290                    crate::store::StoreError::InvalidObject("new pack reader unavailable".into())
291                })?
292                .indexed_read_tiers()?;
293            for (id, tier) in objects {
294                remember_location(&mut index.locations, id, pack_index, tier);
295            }
296        }
297        Ok(())
298    }
299
300    /// Check whether the immutable pack set on disk differs from this snapshot.
301    ///
302    /// Comparing only counts misses the decisive repack transition (`one old`
303    /// → `one replacement`). Exact path comparison lets another `FsStore`
304    /// recover after an atomic cutover even when cardinality is unchanged.
305    /// Half-installed packs remain filtered by `discover_pack_paths`.
306    pub fn needs_reload(&self) -> Result<bool> {
307        let discovered = Self::discover_pack_paths(&self.packs_dir)?;
308        Ok(discovered.len() != self.packs.len()
309            || discovered
310                .iter()
311                .zip(&self.packs)
312                .any(|((pack, index), cached)| {
313                    *pack != cached.pack_path || *index != cached.index_path
314                }))
315    }
316
317    /// Reload the pack list when the immutable pack set changed on disk.
318    ///
319    /// Catches the multi-instance case: two `FsStore`s back the same
320    /// shared object dir (typical for lightweight thread worktrees,
321    /// where the worktree's repo opens its own store but points at
322    /// the main repo's `.heddle/`). When the worktree's store installs
323    /// a new pack, the main repo's already-open `pack_manager`
324    /// doesn't know about it; without this `get_blob`/`has_blob`
325    /// from the main repo would surface "object not found".
326    pub fn reload_if_stale(&mut self) -> Result<bool> {
327        if !self.needs_reload()? {
328            return Ok(false);
329        }
330        debug!("PackManager: pack set changed under us, reloading");
331        self.reload()?;
332        Ok(true)
333    }
334
335    pub fn get_object(&self, id: &PackObjectId) -> Result<Option<(ObjectType, Vec<u8>)>> {
336        let Some(pack_index) = self.point_object_location(id)? else {
337            trace!("Object not found in any pack");
338            return Ok(None);
339        };
340        let Some(reader) = self.packs[pack_index].reader() else {
341            return Ok(None);
342        };
343        let object = reader.get_object(id)?;
344        if object.is_some() {
345            trace!("Found object in pack");
346        }
347        Ok(object)
348    }
349
350    /// Return the physical tier that will serve `id`.
351    ///
352    /// When both layouts contain the same immutable object, lookup always
353    /// selects the hot random-access record before a solid frame.
354    pub fn object_read_tier(&self, id: &PackObjectId) -> Result<Option<PackReadTier>> {
355        let _ = self.object_location(id)?;
356        let index = self
357            .object_locations
358            .read()
359            .unwrap_or_else(std::sync::PoisonError::into_inner);
360        Ok(index.locations.get(id).map(|location| location.tier))
361    }
362
363    /// Read `id` from one specific discovered pack without building the
364    /// cross-pack location index. The record identity remains validated by
365    /// [`PackReader`]; object-domain consumers must validate decoded content.
366    pub fn get_object_from_pack(
367        &self,
368        pack_path: &Path,
369        id: &PackObjectId,
370    ) -> Result<Option<(ObjectType, Vec<u8>)>> {
371        let Some(pack) = self.packs.iter().find(|pack| pack.pack_path == pack_path) else {
372            return Ok(None);
373        };
374        let Some(reader) = pack.reader() else {
375            return Ok(None);
376        };
377        reader.get_object(id)
378    }
379
380    /// List the identities in one specific pack without building the
381    /// cross-pack location index.
382    pub fn list_ids_from_pack(&self, pack_path: &Path) -> Result<Vec<PackObjectId>> {
383        let Some(pack) = self.packs.iter().find(|pack| pack.pack_path == pack_path) else {
384            return Ok(Vec::new());
385        };
386        let Some(reader) = pack.reader() else {
387            return Ok(Vec::new());
388        };
389        reader.list_ids()
390    }
391
392    #[instrument(skip(self), fields(hash = %hash.short()))]
393    pub fn get_hashed_object(&self, hash: &ContentHash) -> Result<Option<(ObjectType, Vec<u8>)>> {
394        self.get_object(&PackObjectId::Hash(*hash))
395    }
396
397    /// Look up the logical object type without decoding the object payload.
398    pub fn get_hashed_object_type(&self, hash: &ContentHash) -> Result<Option<ObjectType>> {
399        let id = PackObjectId::Hash(*hash);
400        let Some(pack_index) = self.point_object_location(&id)? else {
401            return Ok(None);
402        };
403        let Some(reader) = self.packs[pack_index].reader() else {
404            return Ok(None);
405        };
406        reader.get_hashed_object_type(hash)
407    }
408
409    /// Zero-copy variant of `get_hashed_object`. Returns
410    /// [`bytes::Bytes`] views into the underlying pack mmap when
411    /// the entry is non-delta and stored uncompressed; falls back
412    /// to the standard decompress-into-Vec path otherwise.
413    pub fn get_hashed_object_bytes(
414        &self,
415        hash: &ContentHash,
416    ) -> Result<Option<(ObjectType, bytes::Bytes)>> {
417        let id = PackObjectId::Hash(*hash);
418        let Some(pack_index) = self.point_object_location(&id)? else {
419            return Ok(None);
420        };
421        let Some(reader) = self.packs[pack_index].reader() else {
422            return Ok(None);
423        };
424        reader.get_object_bytes(&id)
425    }
426
427    pub fn has_object(&self, hash: &ContentHash) -> bool {
428        self.point_object_location(&PackObjectId::Hash(*hash))
429            .is_ok_and(|location| location.is_some())
430    }
431
432    /// Look up the uncompressed size of `hash` across all loaded
433    /// packs without decompressing the payload. Returns `Ok(None)`
434    /// when the object isn't in any loaded pack.
435    pub fn get_hashed_object_size(&self, hash: &ContentHash) -> Result<Option<u64>> {
436        let id = PackObjectId::Hash(*hash);
437        let Some(pack_index) = self.point_object_location(&id)? else {
438            return Ok(None);
439        };
440        let Some(reader) = self.packs[pack_index].reader() else {
441            return Ok(None);
442        };
443        reader.get_hashed_object_size(hash)
444    }
445
446    pub fn has_object_id(&self, id: &PackObjectId) -> bool {
447        self.point_object_location(id)
448            .is_ok_and(|location| location.is_some())
449    }
450
451    /// List all object hashes across all packs.
452    pub fn list_all_hashes(&self) -> Result<Vec<ContentHash>> {
453        let mut hashes = Vec::new();
454        for pack in &self.packs {
455            if let Some(reader) = pack.reader() {
456                hashes.extend(reader.list_hashes()?);
457            }
458        }
459        Ok(hashes)
460    }
461
462    pub fn list_all_ids(&self) -> Result<Vec<PackObjectId>> {
463        let mut ids = Vec::new();
464        for pack in &self.packs {
465            if let Some(reader) = pack.reader() {
466                ids.extend(reader.list_ids()?);
467            }
468        }
469        Ok(ids)
470    }
471
472    /// Return paths of all pack files (for deletion during aggressive repack).
473    pub fn pack_file_paths(&self) -> Vec<(&Path, &Path)> {
474        self.packs
475            .iter()
476            .map(|pack| (pack.pack_path.as_path(), pack.index_path.as_path()))
477            .collect()
478    }
479
480    pub fn pack_count(&self) -> usize {
481        self.packs.len()
482    }
483
484    pub fn packs_dir(&self) -> &Path {
485        &self.packs_dir
486    }
487}
488
489fn pack_modified(path: &Path) -> SystemTime {
490    fs::metadata(path)
491        .and_then(|metadata| metadata.modified())
492        .unwrap_or(SystemTime::UNIX_EPOCH)
493}
494
495fn remember_location(
496    locations: &mut HashMap<PackObjectId, ObjectLocation>,
497    id: PackObjectId,
498    pack_index: usize,
499    tier: PackReadTier,
500) {
501    let candidate = ObjectLocation { pack_index, tier };
502    match locations.get_mut(&id) {
503        Some(existing)
504            if existing.tier == PackReadTier::SolidFrame && tier == PackReadTier::Hot =>
505        {
506            *existing = candidate;
507        }
508        Some(_) => {}
509        None => {
510            locations.insert(id, candidate);
511        }
512    }
513}
514
515fn force_eager_pack_index() -> bool {
516    std::env::var("HEDDLE_PERF_FORCE_EAGER_PACK_INDEX")
517        .is_ok_and(|value| matches!(value.as_str(), "1" | "true" | "yes"))
518}
519
520#[cfg(test)]
521#[path = "manager_tests.rs"]
522mod tests;