Skip to main content

objects/store/fs/repack/
operation.rs

1// SPDX-License-Identifier: Apache-2.0
2
3use std::{collections::HashSet, fs, fs::OpenOptions};
4
5use super::{
6    super::{
7        FsStore,
8        fs_io::list_hashes_from_dir,
9        fs_paths::{blobs_dir, packs_dir, trees_dir},
10    },
11    compact::add_compact_metadata,
12    cutover::{
13        acquire_repack_lock, cutover, file_len, hash_file, object_file_len, preserve_commit_markers,
14    },
15    staging::{BuildError, RepackSnapshot, RepackStaging, verify_staged},
16};
17use crate::store::{
18    HeddleError, ObjectStore, Result,
19    pack::{
20        ObjectType, PackObjectId, PackReader, RepackContext, RepackError, RepackInventory,
21        RepackOperation, RepackOutcome, StreamingPackBuilder,
22    },
23};
24
25/// Compact native encoder wired behind the generic background repack seam.
26///
27/// It streams one object at a time (bounded memory), verifies every typed id
28/// and the exact expected object set, then installs the immutable replacement
29/// before retiring source packs under the pack-manager write lock. The
30/// scheduler itself has no native-pack knowledge.
31pub struct FsRepackOperation {
32    store: FsStore,
33    #[cfg(test)]
34    corrupt_first_object: bool,
35}
36
37impl FsRepackOperation {
38    /// Create a native filesystem repack payload.
39    pub fn new(store: FsStore) -> Self {
40        Self {
41            store,
42            #[cfg(test)]
43            corrupt_first_object: false,
44        }
45    }
46
47    #[cfg(test)]
48    pub(super) fn with_corrupted_output(mut self) -> Self {
49        self.corrupt_first_object = true;
50        self
51    }
52
53    fn inventory(&self) -> Result<RepackInventory> {
54        let loose_blobs = list_hashes_from_dir(&blobs_dir(self.store.root()))?;
55        let loose_trees = list_hashes_from_dir(&trees_dir(self.store.root()))?;
56        let loose_bytes = loose_blobs
57            .iter()
58            .map(|hash| object_file_len(&blobs_dir(self.store.root()), hash))
59            .chain(
60                loose_trees
61                    .iter()
62                    .map(|hash| object_file_len(&trees_dir(self.store.root()), hash)),
63            )
64            .sum();
65        let manager =
66            self.store.pack_manager().read().map_err(|_| {
67                HeddleError::Config("Failed to acquire pack manager lock".to_string())
68            })?;
69        let paths = manager.pack_file_paths();
70        let pack_bytes = paths
71            .iter()
72            .flat_map(|(pack, index)| [*pack, *index])
73            .map(file_len)
74            .sum();
75        let ids = manager.list_all_ids()?;
76        let unique = ids.iter().copied().collect::<HashSet<_>>().len() as u64;
77        Ok(RepackInventory {
78            loose_objects: (loose_blobs.len() + loose_trees.len()) as u64,
79            loose_bytes,
80            pack_count: paths.len() as u64,
81            pack_bytes,
82            duplicate_objects: ids.len() as u64 - unique,
83            packed_objects: ids.len() as u64,
84        })
85    }
86
87    fn execute(&self, context: &RepackContext) -> std::result::Result<RepackOutcome, RepackError> {
88        let packs = packs_dir(self.store.root());
89        fs::create_dir_all(&packs).map_err(RepackError::operation)?;
90        let _operation_lock = acquire_repack_lock(&packs, context)?;
91        context.checkpoint(0)?;
92
93        let snapshot = RepackSnapshot::capture(&self.store).map_err(RepackError::operation)?;
94        if snapshot.ids.is_empty()
95            && snapshot.loose_blobs.is_empty()
96            && snapshot.loose_trees.is_empty()
97        {
98            return Ok(RepackOutcome::default());
99        }
100
101        let staging = RepackStaging::new(&packs).map_err(RepackError::operation)?;
102        let (expected, logical_bytes) =
103            self.build_staged(&snapshot, &staging, context)
104                .map_err(|error| match error {
105                    BuildError::Cancelled(error) => error,
106                    BuildError::Store(error) => RepackError::operation(error),
107                })?;
108        verify_staged(&staging, &expected, context)?;
109        context.checkpoint(0)?;
110
111        // Cutover starts here and is intentionally non-cancellable. The new
112        // pair is durable before any source path is retired.
113        let new_pack_name = hash_file(&staging.pack).map_err(RepackError::operation)?;
114        let replacement_preexisting = packs.join(format!("{new_pack_name}.pack")).exists()
115            && packs.join(format!("{new_pack_name}.idx")).exists();
116        ObjectStore::install_pack_streaming(&self.store, &staging.pack, &staging.index)
117            .map_err(RepackError::operation)?;
118        preserve_commit_markers(&packs, &new_pack_name, &snapshot.commit_artifact_ids)
119            .map_err(RepackError::operation)?;
120        let cutover = cutover(
121            &self.store,
122            &snapshot,
123            &new_pack_name,
124            replacement_preexisting,
125        )
126        .map_err(RepackError::operation)?;
127        // The replacement is authoritative before loose copies are pruned.
128        // `prune_loose_objects_impl` rechecks exact packed content, so a
129        // concurrent loose write is removed only when the durable pack has it.
130        let (_, loose_bytes) = self
131            .store
132            .prune_loose_objects_impl()
133            .map_err(RepackError::operation)?;
134        let reclaimed = cutover
135            .removed_pack_bytes
136            .saturating_add(loose_bytes)
137            .saturating_sub(cutover.replacement_bytes);
138
139        Ok(RepackOutcome {
140            objects_repacked: expected.len() as u64,
141            bytes_repacked: logical_bytes,
142            bytes_reclaimed: reclaimed,
143        })
144    }
145
146    fn build_staged(
147        &self,
148        snapshot: &RepackSnapshot,
149        staging: &RepackStaging,
150        context: &RepackContext,
151    ) -> std::result::Result<(HashSet<PackObjectId>, u64), BuildError> {
152        let pack_file = OpenOptions::new()
153            .create_new(true)
154            .read(true)
155            .write(true)
156            .open(&staging.pack)?;
157        let mut compression = self.store.compression();
158        compression.max_delta_size = 0;
159        let mut builder = StreamingPackBuilder::new(
160            pack_file,
161            staging.index.clone(),
162            compression,
163            staging.buckets.clone(),
164        )?;
165        let mut expected = HashSet::new();
166        let mut state_ids = Vec::new();
167        let mut tree_hashes = Vec::new();
168        let mut blob_hashes = Vec::new();
169        #[cfg(test)]
170        let mut corrupt_first = self.corrupt_first_object;
171        #[cfg(not(test))]
172        let mut corrupt_first = false;
173        let mut logical_bytes = 0u64;
174
175        for (pack, index) in &snapshot.old_pack_files {
176            let reader = PackReader::open(pack, index)?;
177            let mut checkpoint_error = None;
178            let visit = reader.visit_objects(|id, object_type, data| {
179                if !expected.insert(id) {
180                    return Ok(());
181                }
182                match (id, object_type) {
183                    (PackObjectId::Hash(hash), ObjectType::Blob) => blob_hashes.push(hash),
184                    (PackObjectId::Hash(hash), ObjectType::Tree) => tree_hashes.push(hash),
185                    (PackObjectId::StateId(state_id), ObjectType::State) => {
186                        state_ids.push(state_id);
187                    }
188                    _ => {
189                        let mut data = data.to_vec();
190                        if corrupt_first {
191                            data.push(0xff);
192                            corrupt_first = false;
193                        }
194                        logical_bytes = logical_bytes.saturating_add(data.len() as u64);
195                        builder.add_id(id, object_type, data)?;
196                    }
197                }
198                if let Err(error) = context.checkpoint(data.len() as u64) {
199                    checkpoint_error = Some(error);
200                    return Err(HeddleError::InvalidObject(
201                        "repack source walk interrupted".to_string(),
202                    ));
203                }
204                Ok(())
205            });
206            if let Some(error) = checkpoint_error {
207                return Err(BuildError::Cancelled(error));
208            }
209            visit?;
210        }
211        for hash in &snapshot.loose_blobs {
212            let id = PackObjectId::Hash(*hash);
213            if expected.insert(id) {
214                blob_hashes.push(*hash);
215            }
216        }
217        for hash in &snapshot.loose_trees {
218            let id = PackObjectId::Hash(*hash);
219            if expected.insert(id) {
220                tree_hashes.push(*hash);
221            }
222        }
223        logical_bytes = logical_bytes.saturating_add(add_compact_metadata(
224            &self.store,
225            &mut builder,
226            &state_ids,
227            &tree_hashes,
228            &blob_hashes,
229            context,
230            &mut corrupt_first,
231        )?);
232        builder.finalize()?;
233        Ok((expected, logical_bytes))
234    }
235}
236
237impl RepackOperation for FsRepackOperation {
238    fn key(&self) -> String {
239        self.store.root().display().to_string()
240    }
241
242    fn inspect(&self) -> std::result::Result<RepackInventory, RepackError> {
243        self.inventory().map_err(RepackError::operation)
244    }
245
246    fn run(&self, context: &RepackContext) -> std::result::Result<RepackOutcome, RepackError> {
247        self.execute(context)
248    }
249}