Skip to main content

objects/store/fs/
fs_pack.rs

1// SPDX-License-Identifier: Apache-2.0
2//! Pack and prune operations for FsStore.
3
4use std::{
5    fs,
6    num::NonZeroUsize,
7    path::{Path, PathBuf},
8    sync::Arc,
9};
10
11use super::{
12    FsStore,
13    fs_impl::validate_state_serialized,
14    fs_io::{list_hashes_from_dir, list_state_ids_from_dir, read_file_bytes},
15    fs_paths::{blobs_dir, hash_path, packs_dir, state_path, states_dir, trees_dir},
16};
17use crate::{
18    object::{ContentHash, State, StateAttachment, StateAttachmentId},
19    store::{
20        FsRepackOperation, HeddleError, ObjectStore, RepackPolicy, RepackResourceLimits,
21        RepackSchedule, RepackScheduler, Result, SnapshotCommitArtifact, SnapshotCommitDescriptor,
22        TreeWrite, codec,
23        pack::{ObjectType as PackObjectType, PackBuilder, PackObjectId, PackReader},
24        snapshot_commit::snapshot_commit_marker_path,
25    },
26};
27
28/// Paths of `*.pack` files in `packs_dir` that have no matching `*.idx`.
29///
30/// L8 residual: crash between durable pack and index publish can leave an
31/// unpaired pack that [`FsStore::reload_packs`] ignores. Listing supports
32/// optional GC (design: `docs/program/L8_PACK_INSTALL_JOURNAL.md` Option D).
33/// Does not delete anything.
34pub(crate) fn list_unpaired_pack_files(packs_dir: &Path) -> std::io::Result<Vec<PathBuf>> {
35    if !packs_dir.exists() {
36        return Ok(Vec::new());
37    }
38    let mut unpaired = Vec::new();
39    for entry in fs::read_dir(packs_dir)? {
40        let entry = entry?;
41        let path = entry.path();
42        if path.extension().and_then(|e| e.to_str()) != Some("pack") {
43            continue;
44        }
45        let idx = path.with_extension("idx");
46        if !idx.exists() {
47            unpaired.push(path);
48        }
49    }
50    unpaired.sort();
51    Ok(unpaired)
52}
53
54/// Remove unpaired `*.pack` files (no matching `*.idx`) under `packs_dir`.
55///
56/// Safe for correctness: loaders never open unpaired packs. Bounds L8 disk
57/// leak. Returns `(removed_count, bytes_freed)`.
58pub(crate) fn prune_unpaired_pack_files(packs_dir: &Path) -> std::io::Result<(u64, u64)> {
59    let mut removed = 0u64;
60    let mut bytes_freed = 0u64;
61    for path in list_unpaired_pack_files(packs_dir)? {
62        let bytes = fs::metadata(&path).map(|m| m.len()).unwrap_or(0);
63        match fs::remove_file(&path) {
64            Ok(()) => {
65                removed += 1;
66                bytes_freed = bytes_freed.saturating_add(bytes);
67            }
68            Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
69            Err(e) => return Err(e),
70        }
71    }
72    Ok((removed, bytes_freed))
73}
74
75fn remove_file_ignore_missing(path: &std::path::Path) -> Result<()> {
76    match fs::remove_file(path) {
77        Ok(()) => Ok(()),
78        Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
79        Err(e) => Err(HeddleError::from(e)),
80    }
81}
82
83fn remove_file_counted(path: &Path) -> Result<Option<u64>> {
84    let metadata = match fs::metadata(path) {
85        Ok(metadata) => metadata,
86        Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(None),
87        Err(error) => return Err(HeddleError::from(error)),
88    };
89    match fs::remove_file(path) {
90        Ok(()) => Ok(Some(metadata.len())),
91        Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(None),
92        Err(error) => Err(HeddleError::from(error)),
93    }
94}
95
96impl FsStore {
97    /// Rewrite all packs without `hash`, remove its loose copy, and verify that
98    /// neither local nor external object lookup can still serve the bytes.
99    pub fn remove_blob_everywhere(&self, hash: &ContentHash) -> Result<bool> {
100        let was_present = ObjectStore::has_blob_locally(self, hash)?;
101        if was_present {
102            let scheduler = RepackScheduler::new(
103                RepackPolicy::default(),
104                RepackResourceLimits::new(NonZeroUsize::MIN),
105            );
106            let operation = Arc::new(FsRepackOperation::new(self.clone()).excluding_blob(*hash));
107            let RepackSchedule::Started(handle) = scheduler
108                .repack_now(operation)
109                .map_err(|error| HeddleError::InvalidObject(error.to_string()))?
110            else {
111                return Err(HeddleError::InvalidObject(
112                    "exclusive purge repack did not start".to_string(),
113                ));
114            };
115            handle
116                .wait()
117                .map_err(|error| HeddleError::InvalidObject(error.to_string()))?;
118
119            // The repack operation owns a clone of this store, so its atomic
120            // cutover updates that clone's in-memory pack manager. Reload the
121            // caller's manager from the newly published generation before
122            // checking whether the purged object is still reachable.
123            self.reload_packs()?;
124
125            let loose = hash_path(&blobs_dir(&self.root), hash);
126            match fs::remove_file(&loose) {
127                Ok(()) => {
128                    if let Some(parent) = loose.parent() {
129                        crate::fs_atomic::sync_directory(parent)?;
130                    }
131                }
132                Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
133                Err(error) => return Err(error.into()),
134            }
135        }
136        self.clear_recent_object_caches();
137        if ObjectStore::has_blob_locally(self, hash)?
138            || ObjectStore::get_blob(self, hash)?.is_some()
139            || ObjectStore::get_blob_bytes(self, hash)?.is_some()
140        {
141            return Err(HeddleError::InvalidObject(format!(
142                "purged blob {} remains readable after pack rewrite",
143                hash.short()
144            )));
145        }
146        Ok(was_present)
147    }
148
149    /// Install a structured snapshot closure and its commit artifact through
150    /// the filesystem store's single durable pack barrier.
151    #[doc(hidden)]
152    pub fn put_committed_snapshot_objects_packed(
153        &self,
154        blobs: Vec<(ContentHash, Vec<u8>)>,
155        trees: Vec<TreeWrite>,
156        tree: &TreeWrite,
157        state: &State,
158        attachments: Vec<StateAttachment>,
159        artifact: SnapshotCommitArtifact,
160    ) -> Result<SnapshotCommitDescriptor> {
161        self.put_snapshot_objects_packed_impl(
162            blobs,
163            trees,
164            tree,
165            state,
166            attachments,
167            Some(artifact),
168        )?
169        .ok_or_else(|| {
170            HeddleError::InvalidObject(
171                "committed snapshot pack did not expose its artifact descriptor".to_string(),
172            )
173        })
174    }
175
176    /// Install blobs, root tree, state, and immutable authored attachments
177    /// through one pack publication. Ordinary callers treat the pack as
178    /// pre-oplog staging; committed structured snapshots add their local trust
179    /// marker in the same directory barrier and make the pack authoritative.
180    pub(super) fn put_snapshot_objects_packed_impl(
181        &self,
182        blobs: Vec<(ContentHash, Vec<u8>)>,
183        trees: Vec<TreeWrite>,
184        tree: &TreeWrite,
185        state: &State,
186        attachments: Vec<StateAttachment>,
187        commit_artifact: Option<SnapshotCommitArtifact>,
188    ) -> Result<Option<SnapshotCommitDescriptor>> {
189        // A committed snapshot artifact is installed only after exact-once and
190        // isolation validation. Its freshly-authored StateId cannot be a retry
191        // (dedup returns before this callback), so avoid an expected-negative
192        // pack-directory rescan on every native capture.
193        let state_was_present = if commit_artifact.is_some() {
194            false
195        } else {
196            <Self as ObjectStore>::has_state(self, &state.id())?
197        };
198        let mut compression = self.compression;
199        if !self.snapshot_delta_search {
200            compression.max_delta_size = 0;
201        }
202        let mut builder = PackBuilder::new(compression);
203
204        for (hash, data) in blobs {
205            if commit_artifact.is_none() && ObjectStore::has_blob_locally(self, &hash)? {
206                continue;
207            }
208            builder.add(hash, PackObjectType::Blob, data);
209        }
210
211        let tree_hash = tree.tree.hash();
212        let mut staged_trees = Vec::with_capacity(trees.len() + 1);
213        let mut staged_encodings = Vec::with_capacity(trees.len() + 1);
214        let mut seen_trees = std::collections::HashSet::with_capacity(trees.len() + 1);
215        for authored_tree in trees {
216            let authored_hash = authored_tree.tree.hash();
217            let reuses_materialized = self
218                .try_get_tree_serialized_once(&authored_hash)?
219                .is_some_and(|body| !crate::object::is_delta_tree(&body));
220            if seen_trees.insert(authored_hash)
221                && !reuses_materialized
222                && (commit_artifact.is_some()
223                    || !ObjectStore::has_tree_locally(self, &authored_hash)?)
224            {
225                let encoded = self.encode_tree_write(&authored_tree)?;
226                if matches!(
227                    encoded.kind,
228                    crate::store::codec::TreeEncodingKind::Delta { anchor, .. }
229                        if anchor == authored_hash
230                ) {
231                    return Err(HeddleError::InvalidObject(
232                        "HDC1 result id must differ from its anchor id".to_string(),
233                    ));
234                }
235                builder.add(authored_hash, PackObjectType::Tree, encoded.data);
236                staged_encodings.push((authored_hash, encoded.kind));
237                staged_trees.push((authored_hash, authored_tree.tree));
238            }
239        }
240        let reuses_materialized = self
241            .try_get_tree_serialized_once(&tree_hash)?
242            .is_some_and(|body| !crate::object::is_delta_tree(&body));
243        if !reuses_materialized
244            && (commit_artifact.is_some() || !ObjectStore::has_tree_locally(self, &tree_hash)?)
245            && seen_trees.insert(tree_hash)
246        {
247            let encoded = self.encode_tree_write(tree)?;
248            if matches!(
249                encoded.kind,
250                crate::store::codec::TreeEncodingKind::Delta { anchor, .. }
251                    if anchor == tree_hash
252            ) {
253                return Err(HeddleError::InvalidObject(
254                    "HDC1 result id must differ from its anchor id".to_string(),
255                ));
256            }
257            builder.add(tree_hash, PackObjectType::Tree, encoded.data);
258            staged_encodings.push((tree_hash, encoded.kind));
259            staged_trees.push((tree_hash, tree.tree.clone()));
260        }
261
262        let state_id = state.id();
263        builder.add_id(
264            PackObjectId::StateId(state_id),
265            PackObjectType::State,
266            state.encode_current_msgpack()?,
267        );
268        let attachment_ids = attachments
269            .iter()
270            .map(|attachment| {
271                if attachment.state_id != state_id {
272                    return Err(HeddleError::InvalidObject(
273                        "snapshot attachment targets a different state".to_string(),
274                    ));
275                }
276                let id = attachment.id();
277                builder.add(
278                    *id.as_hash(),
279                    PackObjectType::StateAttachment,
280                    attachment.encode_current_msgpack()?,
281                );
282                Ok(id)
283            })
284            .collect::<Result<Vec<StateAttachmentId>>>()?;
285        let artifact_id = commit_artifact.as_ref().map(SnapshotCommitArtifact::id);
286        let artifact_bytes = commit_artifact
287            .as_ref()
288            .map(rmp_serde::to_vec_named)
289            .transpose()?;
290        if let Some(artifact) = &commit_artifact {
291            artifact.validate()?;
292            let bytes = artifact_bytes.as_ref().ok_or_else(|| {
293                HeddleError::InvalidObject(
294                    "snapshot commit artifact bytes were not encoded".to_string(),
295                )
296            })?;
297            builder.add(artifact.id(), PackObjectType::SnapshotCommit, bytes.clone());
298        }
299
300        let (pack_data, index_data, _stats, retained_objects) =
301            builder.build_retaining_objects()?;
302        let packs = packs_dir(&self.root);
303        let installed_pack_name = if commit_artifact.is_some() {
304            let (Some(artifact_id), Some(artifact_bytes)) = (artifact_id, artifact_bytes) else {
305                return Err(HeddleError::InvalidObject(
306                    "snapshot commit artifact metadata is incomplete".to_string(),
307                ));
308            };
309            super::pack_install_journal::install_committed_snapshot_pack_bytes(
310                &packs,
311                pack_data,
312                index_data,
313                artifact_id,
314                artifact_bytes,
315            )?
316        } else {
317            super::pack_install_journal::install_snapshot_pack_bytes(&packs, pack_data, index_data)?
318        };
319        {
320            let mut manager = self.pack_manager().write().map_err(|_| {
321                HeddleError::Config("Failed to acquire pack manager lock".to_string())
322            })?;
323            manager.add_pack(
324                packs.join(format!("{installed_pack_name}.pack")),
325                packs.join(format!("{installed_pack_name}.idx")),
326            )?;
327        }
328        for (hash, kind) in staged_encodings {
329            self.remember_tree_encoding(hash, kind)?;
330        }
331        self.materialize_packed_attachment_index(&state_id, &attachment_ids, state_was_present)?;
332
333        if let Ok(mut cache) = self.recent_blobs.write() {
334            for (id, object_type, data) in retained_objects {
335                if let (PackObjectId::Hash(hash), PackObjectType::Blob) = (id, object_type) {
336                    cache.insert(hash, crate::object::Blob::new(data));
337                }
338            }
339        }
340        if let Ok(mut cache) = self.recent_trees.write() {
341            for (hash, authored_tree) in staged_trees {
342                cache.insert(hash, authored_tree);
343            }
344        }
345        if let Ok(mut cache) = self.recent_states.write() {
346            let mut cached = state.clone();
347            cached.state_id = state_id;
348            cache.insert(state_id, cached);
349        }
350        let descriptor = if let Some(artifact) = commit_artifact {
351            let pack_path = packs.join(format!("{installed_pack_name}.pack"));
352            let index_path = packs.join(format!("{installed_pack_name}.idx"));
353            let object_ids =
354                PackReader::open(&pack_path, &index_path, &self.root.join("tmp"))?.list_ids()?;
355            Some(SnapshotCommitDescriptor {
356                artifact,
357                pack_name: installed_pack_name,
358                pack_path,
359                object_ids,
360            })
361        } else {
362            None
363        };
364        Ok(descriptor)
365    }
366
367    /// Bulk-install many blobs as a single packfile. Two fsyncs total
368    /// (one for `.pack`, one for `.idx`) regardless of blob count —
369    /// vs. N×fsync if each blob were written loose. Used by the
370    /// snapshot hot path; called at the end of the tree walk with
371    /// every new blob accumulated in memory.
372    ///
373    /// Skips blobs already in the store (whether loose or packed) so
374    /// re-snapshotting an unchanged worktree doesn't churn the pack
375    /// directory. With every blob already known, this is a no-op.
376    pub(super) fn put_blobs_packed_impl(&self, blobs: Vec<(ContentHash, Vec<u8>)>) -> Result<()> {
377        if blobs.is_empty() {
378            return Ok(());
379        }
380        // Snapshot-time pack: skip the sliding-window delta search.
381        // It's a CPU win on similar-content files (the GC packer
382        // benefits) but for a single snapshot the inputs are
383        // unrelated content (random binaries, small text, etc.) and
384        // every pair-wise delta estimate runs across the full
385        // payloads — for 16×4MB blobs that's tens of seconds of
386        // hashing for ~zero compression benefit. GC's
387        // `pack_objects_impl` keeps the full delta search; this
388        // path only optimizes durability + write throughput.
389        let mut compression = self.compression;
390        if !self.snapshot_delta_search {
391            compression.max_delta_size = 0;
392        }
393        let mut builder = PackBuilder::new(compression);
394        let mut added = 0usize;
395        for (hash, data) in blobs {
396            if ObjectStore::has_blob_locally(self, &hash)? {
397                continue;
398            }
399            builder.add(hash, PackObjectType::Blob, data);
400            added += 1;
401        }
402        if added == 0 {
403            return Ok(());
404        }
405        let (pack_data, index_data, _stats, retained_objects) =
406            builder.build_retaining_objects()?;
407
408        // A generic install clears recent-object caches because received packs
409        // can shadow loose objects. This locally-built pack returns ownership
410        // of its original inputs after encoding, so repopulating the cache does
411        // not require a payload-sized staging or `Blob::from_slice` copy.
412        self.install_pack_files(&pack_data, &index_data)?;
413        if let Ok(mut cache) = self.recent_blobs.write() {
414            for (id, object_type, data) in retained_objects {
415                if let (PackObjectId::Hash(hash), PackObjectType::Blob) = (id, object_type) {
416                    cache.insert(hash, crate::object::Blob::new(data));
417                }
418            }
419        }
420        Ok(())
421    }
422
423    /// Consolidate the object store into a single pack.
424    ///
425    /// GC must *shrink* the set of places a reader has to look, not grow
426    /// it. The naive "pack the loose objects into a fresh pack" strategy
427    /// regressed read performance badly: every `maintenance gc` minted a
428    /// brand-new pack *alongside* the existing pack(s) and (by default)
429    /// left the now-redundant loose copies in place. The result was an
430    /// object store with strictly MORE sources to search — loose objects
431    /// plus an ever-growing fleet of packs — and `PackManager::get_object`
432    /// probes every pack linearly, so each extra pack roughly doubled the
433    /// cost of the object lookups that `status`/`diff`/verification do.
434    ///
435    /// This implementation does a true repack: it folds every object
436    /// already living in a pack *together with* the loose blobs and trees
437    /// into one new consolidated pack, installs it, and then deletes the
438    /// superseded packs. Combined with the caller's
439    /// `prune_loose_objects`, the store ends a GC with exactly one pack
440    /// and no loose duplicates — strictly fewer read sources than it
441    /// started with. Running GC again over an already-consolidated store
442    /// is a no-op (nothing loose, one pack already covers everything).
443    ///
444    pub(super) fn pack_objects_impl(&self, delta_search: bool) -> Result<(u64, u64)> {
445        // Serialize every source-pack-retiring path with background repack,
446        // including callers in another process. Ordinary immutable pack
447        // installs remain concurrent and are preserved at scheduler cutover.
448        let _repack_lock = super::repack::acquire_repack_lock_blocking(&packs_dir(&self.root))?;
449        let loose_blobs = list_hashes_from_dir(&blobs_dir(&self.root))?;
450        let loose_trees = list_hashes_from_dir(&trees_dir(&self.root))?;
451
452        // Snapshot what the existing packs already hold, plus the file
453        // paths we'll retire once the consolidated pack is installed.
454        let (existing_ids, old_pack_files, commit_artifact_ids) = {
455            let manager = self.pack_manager().read().map_err(|_| {
456                HeddleError::Config("Failed to acquire pack manager lock".to_string())
457            })?;
458            let ids = manager.list_all_ids()?;
459            let commit_artifact_ids = manager
460                .snapshot_commit_descriptors()?
461                .into_iter()
462                .map(|descriptor| descriptor.artifact.id())
463                .collect::<Vec<_>>();
464            let files: Vec<(std::path::PathBuf, std::path::PathBuf)> = manager
465                .pack_file_paths()
466                .into_iter()
467                .map(|(pack, index)| (pack.to_path_buf(), index.to_path_buf()))
468                .collect();
469            (ids, files, commit_artifact_ids)
470        };
471
472        // Nothing loose and at most one pack already — the store is
473        // already consolidated; don't churn a fresh identical pack.
474        if loose_blobs.is_empty() && loose_trees.is_empty() && old_pack_files.len() <= 1 {
475            return Ok((0, 0));
476        }
477
478        // Consolidation packs every object that's already packed plus the
479        // loose ones. The default path skips the sliding-window delta search
480        // to keep foreground GC latency bounded: it searches the full payloads
481        // of every object and can turn a seconds-long consolidation into
482        // minutes. The caller resolves the repository's GC policy and the
483        // `--aggressive` override into the `delta_search` argument. This
484        // mirrors the snapshot hot path, whose policy is held by the store.
485        let mut compression = self.compression;
486        if !delta_search {
487            compression.max_delta_size = 0;
488        }
489        let mut builder = PackBuilder::new(compression);
490        let loose_tree_set: std::collections::HashSet<ContentHash> =
491            loose_trees.iter().copied().collect();
492        let mut seen: std::collections::HashSet<crate::store::pack::PackObjectId> =
493            std::collections::HashSet::new();
494
495        // 1. Carry forward everything already in a pack so the old packs
496        //    can be retired. `get_object` resolves the body + type for
497        //    any id (blob/tree/state/action), and `add_id` preserves
498        //    content-addressed state objects.
499        for id in existing_ids {
500            if !seen.insert(id) {
501                continue;
502            }
503            let obj_type = {
504                let manager = self.pack_manager().read().map_err(|_| {
505                    HeddleError::Config("Failed to acquire pack manager lock".to_string())
506                })?;
507                manager.get_object(&id)?
508            };
509            if let Some((obj_type, mut data)) = obj_type {
510                if let crate::store::pack::PackObjectId::Hash(hash) = id
511                    && obj_type == PackObjectType::Tree
512                    && loose_tree_set.contains(&hash)
513                    && let Some(loose_data) = ObjectStore::get_tree_serialized(self, &hash)?
514                {
515                    data = loose_data;
516                }
517                builder.add_id(id, obj_type, data);
518            }
519        }
520
521        // 2. Fold in the loose blobs and trees. Skip any whose hash is
522        //    already covered by a carried-forward pack entry.
523        for hash in &loose_blobs {
524            let id = crate::store::pack::PackObjectId::Hash(*hash);
525            if seen.contains(&id) {
526                continue;
527            }
528            if let Some(blob) = ObjectStore::get_blob(self, hash)? {
529                seen.insert(id);
530                builder.add(*hash, PackObjectType::Blob, blob.content().to_vec());
531            }
532        }
533        for hash in &loose_trees {
534            let id = crate::store::pack::PackObjectId::Hash(*hash);
535            if seen.contains(&id) {
536                continue;
537            }
538            if let Some(tree) = ObjectStore::get_tree(self, hash)? {
539                let data = tree.encode_canonical()?;
540                seen.insert(id);
541                builder.add(*hash, PackObjectType::Tree, data);
542            }
543        }
544
545        if seen.is_empty() {
546            return Ok((0, 0));
547        }
548
549        let (pack_data, index_data, stats) = builder.build()?;
550        let new_pack_name = blake3::hash(&pack_data).to_hex();
551        if commit_artifact_ids.is_empty() {
552            self.install_pack_files(&pack_data, &index_data)?;
553        } else {
554            super::pack_install_journal::install_snapshot_pack_bytes_with_commit_markers(
555                &packs_dir(&self.root),
556                pack_data,
557                index_data,
558                &commit_artifact_ids,
559            )?;
560            self.reload_packs()?;
561        }
562        // GC packs *replace* loose objects (followed by
563        // `prune_loose_objects`). Bust the recent-objects caches so
564        // a subsequent get_* doesn't return a stale `Blob`/`Tree`
565        // pointing at a path we're about to delete. The snapshot hot
566        // path doesn't go through here — it calls
567        // `install_pack_files` directly via `put_blobs_packed_impl`,
568        // which keeps its caches warm.
569        self.clear_recent_object_caches();
570
571        // Retire the superseded packs now that the consolidated pack is
572        // durably installed and every object they held has been carried
573        // forward. The consolidated pack is content-addressed, so if it
574        // happened to hash-collide with an old pack (a store that was
575        // already a single consolidated pack) that file is excluded here.
576        // Stack hex digest; compare as &str — no format!/String intermediate.
577        for (pack_path, index_path) in &old_pack_files {
578            let is_new_pack = pack_path
579                .file_stem()
580                .and_then(|stem| stem.to_str())
581                .map(|stem| stem == new_pack_name.as_str())
582                .unwrap_or(false);
583            if is_new_pack {
584                continue;
585            }
586            remove_file_ignore_missing(pack_path)?;
587            remove_file_ignore_missing(index_path)?;
588            for artifact_id in &commit_artifact_ids {
589                remove_file_ignore_missing(&snapshot_commit_marker_path(pack_path, artifact_id))?;
590            }
591        }
592        // Retiring source packs requires a full reload of the pack list.
593        self.reload_packs()?;
594        self.clear_recent_object_caches();
595
596        let saved = stats.total_uncompressed - stats.total_compressed;
597        Ok((stats.object_count, saved))
598    }
599
600    pub(super) fn install_pack_files(&self, pack_data: &[u8], index_data: &[u8]) -> Result<()> {
601        let packs = packs_dir(&self.root);
602        // L8 A+: durable staging + intent journal for in-memory pack install
603        // (same crash-safety as install_pack_files_streaming).
604        // Design: docs/program/L8_PACK_INSTALL_JOURNAL.md
605        let _pack_name = super::pack_install_journal::install_pack_bytes_journaled(
606            &packs, pack_data, index_data,
607        )?;
608        // Pack manager picks up the new files. We do *not* clear the
609        // recent-object caches here — every caller that follows this
610        // with a destructive prune is responsible for clearing them
611        // explicitly. Snapshot installs rely on cache stickiness to
612        // keep tight snapshot loops fast (see
613        // `put_blobs_packed_impl`).
614        self.reload_packs()?;
615        Ok(())
616    }
617
618    /// Move a pack and its index already on disk into the store's
619    /// pack directory, computing the pack's content-hash by streaming
620    /// the file (constant memory regardless of pack size). Pairs with
621    /// `StreamingPackBuilder`: pack data, the index, *and* this
622    /// installation step never load the full pack or index into
623    /// memory.
624    ///
625    /// Sources are staged then published via the L8 A+ install journal
626    /// ([`super::pack_install_journal`]): durable staging + intent, then
627    /// pack/index publish with crash recovery on reload.
628    pub(super) fn install_pack_files_streaming(
629        &self,
630        src_pack_path: &std::path::Path,
631        src_index_path: &std::path::Path,
632    ) -> Result<()> {
633        use std::io::Read;
634
635        let packs = packs_dir(&self.root);
636        crate::fs_atomic::create_dir_all_durable(&packs)?;
637
638        // Stream-hash the pack file to derive its name. 64 KiB chunks
639        // keep the hasher's working set tiny.
640        let mut hasher = blake3::Hasher::new();
641        let mut file = fs::File::open(src_pack_path)?;
642        let mut buf = vec![0u8; 64 * 1024];
643        loop {
644            let n = file.read(&mut buf)?;
645            if n == 0 {
646                break;
647            }
648            hasher.update(&buf[..n]);
649        }
650        drop(file);
651        // Native digest for potential callers; hex String only for the journal
652        // path/name boundary (filenames + intent JSON).
653        let pack_hash = hasher.finalize();
654        let pack_name = pack_hash.to_hex().to_string();
655
656        // L8 A+: durable staging + intent journal, then pack/index publish.
657        // Recovery on reload finishes or aborts incomplete installs.
658        // Design: docs/program/L8_PACK_INSTALL_JOURNAL.md
659        super::pack_install_journal::install_pack_files_journaled(
660            &packs,
661            src_pack_path,
662            src_index_path,
663            &pack_name,
664        )?;
665
666        self.clear_recent_object_caches();
667        self.reload_packs()?;
668        Ok(())
669    }
670
671    /// Remove L8 orphan packs (`.pack` without `.idx`) from this store.
672    pub fn prune_unpaired_packs(&self) -> Result<(u64, u64)> {
673        let packs = packs_dir(&self.root);
674        Ok(prune_unpaired_pack_files(&packs)?)
675    }
676
677    pub(super) fn prune_loose_objects_impl(&self) -> Result<(u64, u64)> {
678        let mut removed = 0u64;
679        let mut bytes_freed = 0u64;
680
681        let blobs = list_hashes_from_dir(&blobs_dir(&self.root))?;
682        let trees = list_hashes_from_dir(&trees_dir(&self.root))?;
683        let states = list_state_ids_from_dir(&states_dir(&self.root))?;
684
685        for hash in &blobs {
686            let packed = self
687                .pack_manager()
688                .read()
689                .map_err(|_| {
690                    HeddleError::Config("Failed to acquire pack manager lock".to_string())
691                })?
692                .get_hashed_object(hash)?
693                .is_some();
694            if packed {
695                let path = hash_path(&blobs_dir(&self.root), hash);
696                if let Some(bytes) = remove_file_counted(&path)? {
697                    bytes_freed = bytes_freed.saturating_add(bytes);
698                    removed += 1;
699                }
700            }
701        }
702
703        for hash in &trees {
704            let path = hash_path(&trees_dir(&self.root), hash);
705            let Some(loose_data) = read_file_bytes(&path)? else {
706                continue;
707            };
708            let loose_body = codec::decode_tree_body(loose_data.as_slice())?;
709            let loose_is_delta = crate::object::is_delta_tree(&loose_body);
710            let loose_tree = self.decode_tree_storage_body(*hash, &loose_body)?;
711            let found = loose_tree.hash();
712            if found != *hash {
713                return Err(HeddleError::Corruption {
714                    expected: *hash,
715                    found,
716                });
717            }
718            let npk1_tree = self
719                .npk1_manager()
720                .read()
721                .map_err(|_| {
722                    HeddleError::Config("Failed to acquire NPK1 manager lock".to_string())
723                })?
724                .get_tree(hash)?;
725            if npk1_tree.as_ref() == Some(&loose_tree) {
726                if let Some(bytes) = remove_file_counted(&path)? {
727                    bytes_freed = bytes_freed.saturating_add(bytes);
728                    removed += 1;
729                }
730                continue;
731            }
732            // Keep main's HDC1 hot-tier safety rule unless an identical,
733            // materialized NPK1 tree has already made the delta redundant.
734            if loose_is_delta {
735                continue;
736            }
737            let packed = self
738                .pack_manager()
739                .read()
740                .map_err(|_| {
741                    HeddleError::Config("Failed to acquire pack manager lock".to_string())
742                })?
743                .get_hashed_object(hash)?;
744            let Some((obj_type, packed_data)) = packed else {
745                continue;
746            };
747            if obj_type != PackObjectType::Tree {
748                continue;
749            }
750            // A loose current tree can intentionally shadow an older packed
751            // schema at the same semantic hash. Preserve that migration copy
752            // until consolidation replaces the legacy body.
753            if crate::object::is_delta_tree(&packed_data) {
754                continue;
755            }
756            let Ok(packed_tree) = codec::decode_tree_serialized_with_key(&packed_data, *hash, None)
757            else {
758                continue;
759            };
760            let packed_found = packed_tree.hash();
761            if packed_found != *hash {
762                return Err(HeddleError::Corruption {
763                    expected: *hash,
764                    found: packed_found,
765                });
766            }
767            if packed_tree == loose_tree
768                && let Some(bytes) = remove_file_counted(&path)?
769            {
770                bytes_freed = bytes_freed.saturating_add(bytes);
771                removed += 1;
772            }
773        }
774
775        for id in &states {
776            let packed = self
777                .pack_manager()
778                .read()
779                .map_err(|_| {
780                    HeddleError::Config("Failed to acquire pack manager lock".to_string())
781                })?
782                .get_object(&PackObjectId::StateId(*id))?;
783            let Some((obj_type, packed_data)) = packed else {
784                continue;
785            };
786            if obj_type != PackObjectType::State {
787                continue;
788            }
789            let path = state_path(&self.root, id);
790            let Some(loose_data) = read_file_bytes(&path)? else {
791                continue;
792            };
793            let loose_state = codec::decode_state(loose_data.as_slice())?;
794            let packed_state = validate_state_serialized(&packed_data, *id)?;
795            if !loose_state.accepts_stored_id(id) {
796                return Err(HeddleError::InvalidObject(format!(
797                    "loose state id mismatch while pruning: expected {id}, computed {}",
798                    loose_state.id()
799                )));
800            }
801            if packed_state == loose_state
802                && let Some(bytes) = remove_file_counted(&path)?
803            {
804                bytes_freed = bytes_freed.saturating_add(bytes);
805                removed += 1;
806            }
807        }
808
809        Ok((removed, bytes_freed))
810    }
811}
812
813#[cfg(test)]
814mod unpaired_pack_tests {
815    use std::fs;
816
817    use super::{list_unpaired_pack_files, prune_unpaired_pack_files};
818
819    #[test]
820    fn list_and_prune_unpaired_packs() {
821        let dir = tempfile::tempdir().unwrap();
822        let packs = dir.path();
823        fs::write(packs.join("aaa.pack"), b"pack-only").unwrap();
824        fs::write(packs.join("bbb.pack"), b"paired-pack").unwrap();
825        fs::write(packs.join("bbb.idx"), b"paired-idx").unwrap();
826        fs::write(packs.join("ccc.idx"), b"index-only").unwrap();
827
828        let listed = list_unpaired_pack_files(packs).unwrap();
829        assert_eq!(listed.len(), 1);
830        assert!(listed[0].ends_with("aaa.pack"));
831
832        let (removed, bytes) = prune_unpaired_pack_files(packs).unwrap();
833        assert_eq!(removed, 1);
834        assert_eq!(bytes, b"pack-only".len() as u64);
835        assert!(!packs.join("aaa.pack").exists());
836        assert!(packs.join("bbb.pack").exists());
837        assert!(packs.join("bbb.idx").exists());
838        assert!(packs.join("ccc.idx").exists());
839        assert!(list_unpaired_pack_files(packs).unwrap().is_empty());
840    }
841
842    #[test]
843    fn missing_packs_dir_is_empty() {
844        let dir = tempfile::tempdir().unwrap();
845        let missing = dir.path().join("nope");
846        assert!(list_unpaired_pack_files(&missing).unwrap().is_empty());
847        assert_eq!(prune_unpaired_pack_files(&missing).unwrap(), (0, 0));
848    }
849}