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, list_state_ids_from_dir},
9        fs_paths::{blobs_dir, packs_dir, state_path, states_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    excluded: HashSet<PackObjectId>,
34    #[cfg(test)]
35    corrupt_first_object: bool,
36}
37
38impl FsRepackOperation {
39    /// Create a native filesystem repack payload.
40    pub fn new(store: FsStore) -> Self {
41        Self {
42            store,
43            excluded: HashSet::new(),
44            #[cfg(test)]
45            corrupt_first_object: false,
46        }
47    }
48
49    /// Exclude one blob from the replacement pack.
50    ///
51    /// This is the storage primitive used by purge: cutover retires every
52    /// source pack only after the replacement has been verified without the
53    /// excluded identity.
54    pub fn excluding_blob(mut self, hash: crate::object::ContentHash) -> Self {
55        self.excluded.insert(PackObjectId::Hash(hash));
56        self
57    }
58
59    #[cfg(test)]
60    pub(super) fn with_corrupted_output(mut self) -> Self {
61        self.corrupt_first_object = true;
62        self
63    }
64
65    fn inventory(&self) -> Result<RepackInventory> {
66        let loose_blobs = list_hashes_from_dir(&blobs_dir(self.store.root()))?;
67        let loose_trees = list_hashes_from_dir(&trees_dir(self.store.root()))?;
68        let loose_states = list_state_ids_from_dir(&states_dir(self.store.root()))?;
69        let loose_bytes = loose_blobs
70            .iter()
71            .map(|hash| object_file_len(&blobs_dir(self.store.root()), hash))
72            .chain(
73                loose_trees
74                    .iter()
75                    .map(|hash| object_file_len(&trees_dir(self.store.root()), hash)),
76            )
77            .chain(
78                loose_states
79                    .iter()
80                    .map(|id| file_len(&state_path(self.store.root(), id))),
81            )
82            .sum();
83        let manager =
84            self.store.pack_manager().read().map_err(|_| {
85                HeddleError::Config("Failed to acquire pack manager lock".to_string())
86            })?;
87        let paths = manager.pack_file_paths();
88        let pack_bytes = paths
89            .iter()
90            .flat_map(|(pack, index)| [*pack, *index])
91            .map(file_len)
92            .sum();
93        let ids = manager.list_all_ids()?;
94        let unique = ids.iter().copied().collect::<HashSet<_>>().len() as u64;
95        Ok(RepackInventory {
96            loose_objects: (loose_blobs.len() + loose_trees.len() + loose_states.len()) as u64,
97            loose_bytes,
98            pack_count: paths.len() as u64,
99            pack_bytes,
100            duplicate_objects: ids.len() as u64 - unique,
101            packed_objects: ids.len() as u64,
102        })
103    }
104
105    fn execute(&self, context: &RepackContext) -> std::result::Result<RepackOutcome, RepackError> {
106        let packs = packs_dir(self.store.root());
107        fs::create_dir_all(&packs).map_err(RepackError::operation)?;
108        let _operation_lock = acquire_repack_lock(&packs, context)?;
109        context.checkpoint(0)?;
110
111        let snapshot = RepackSnapshot::capture(&self.store).map_err(RepackError::operation)?;
112        if snapshot.ids.is_empty()
113            && snapshot.loose_blobs.is_empty()
114            && snapshot.loose_trees.is_empty()
115            && snapshot.loose_states.is_empty()
116        {
117            return Ok(RepackOutcome::default());
118        }
119
120        let staging = RepackStaging::new(&packs).map_err(RepackError::operation)?;
121        let (expected, logical_bytes) =
122            self.build_staged(&snapshot, &staging, context)
123                .map_err(|error| match error {
124                    BuildError::Cancelled(error) => error,
125                    BuildError::Store(error) => RepackError::operation(error),
126                })?;
127        verify_staged(&staging, &expected, context)?;
128        context.checkpoint(0)?;
129
130        // Cutover starts here and is intentionally non-cancellable. The new
131        // pair is durable before any source path is retired.
132        let new_pack_name = hash_file(&staging.pack).map_err(RepackError::operation)?;
133        let replacement_preexisting = packs.join(format!("{new_pack_name}.pack")).exists()
134            && packs.join(format!("{new_pack_name}.idx")).exists();
135        ObjectStore::install_pack_streaming(&self.store, &staging.pack, &staging.index)
136            .map_err(RepackError::operation)?;
137        preserve_commit_markers(&packs, &new_pack_name, &snapshot.commit_artifact_ids)
138            .map_err(RepackError::operation)?;
139        let cutover = cutover(
140            &self.store,
141            &snapshot,
142            &new_pack_name,
143            replacement_preexisting,
144        )
145        .map_err(RepackError::operation)?;
146        // The replacement is authoritative before loose copies are pruned.
147        // `prune_loose_objects_impl` rechecks exact packed content, so a
148        // concurrent loose write is removed only when the durable pack has it.
149        let (_, loose_bytes) = self
150            .store
151            .prune_loose_objects_impl()
152            .map_err(RepackError::operation)?;
153        let reclaimed = cutover
154            .removed_pack_bytes
155            .saturating_add(loose_bytes)
156            .saturating_sub(cutover.replacement_bytes);
157
158        Ok(RepackOutcome {
159            objects_repacked: expected.len() as u64,
160            bytes_repacked: logical_bytes,
161            bytes_reclaimed: reclaimed,
162        })
163    }
164
165    fn build_staged(
166        &self,
167        snapshot: &RepackSnapshot,
168        staging: &RepackStaging,
169        context: &RepackContext,
170    ) -> std::result::Result<(HashSet<PackObjectId>, u64), BuildError> {
171        let pack_file = OpenOptions::new()
172            .create_new(true)
173            .read(true)
174            .write(true)
175            .open(&staging.pack)?;
176        let mut compression = self.store.compression();
177        compression.max_delta_size = 0;
178        let mut builder = StreamingPackBuilder::new(
179            pack_file,
180            staging.index.clone(),
181            compression,
182            staging.buckets.clone(),
183        )?;
184        let mut expected = HashSet::new();
185        let mut state_ids = Vec::new();
186        let mut tree_hashes = Vec::new();
187        let mut blob_hashes = Vec::new();
188        #[cfg(test)]
189        let mut corrupt_first = self.corrupt_first_object;
190        #[cfg(not(test))]
191        let mut corrupt_first = false;
192        let mut logical_bytes = 0u64;
193
194        for (pack, index) in &snapshot.old_pack_files {
195            let reader = PackReader::open(pack, index)?;
196            let mut checkpoint_error = None;
197            let visit = reader.visit_objects(|id, object_type, data| {
198                if self.excluded.contains(&id) {
199                    if let Err(error) = context.checkpoint(data.len() as u64) {
200                        checkpoint_error = Some(error);
201                        return Err(HeddleError::InvalidObject(
202                            "repack source walk interrupted".to_string(),
203                        ));
204                    }
205                    return Ok(());
206                }
207                if !expected.insert(id) {
208                    return Ok(());
209                }
210                match (id, object_type) {
211                    (PackObjectId::Hash(hash), ObjectType::Blob) => blob_hashes.push(hash),
212                    (PackObjectId::Hash(hash), ObjectType::Tree) => tree_hashes.push(hash),
213                    (PackObjectId::StateId(state_id), ObjectType::State) => {
214                        state_ids.push(state_id);
215                    }
216                    _ => {
217                        let mut data = data.to_vec();
218                        if corrupt_first {
219                            data.push(0xff);
220                            corrupt_first = false;
221                        }
222                        logical_bytes = logical_bytes.saturating_add(data.len() as u64);
223                        builder.add_id(id, object_type, data)?;
224                    }
225                }
226                if let Err(error) = context.checkpoint(data.len() as u64) {
227                    checkpoint_error = Some(error);
228                    return Err(HeddleError::InvalidObject(
229                        "repack source walk interrupted".to_string(),
230                    ));
231                }
232                Ok(())
233            });
234            if let Some(error) = checkpoint_error {
235                return Err(BuildError::Cancelled(error));
236            }
237            visit?;
238        }
239        for hash in &snapshot.loose_blobs {
240            let id = PackObjectId::Hash(*hash);
241            if self.excluded.contains(&id) {
242                continue;
243            }
244            if expected.insert(id) {
245                blob_hashes.push(*hash);
246            }
247        }
248        for hash in &snapshot.loose_trees {
249            let id = PackObjectId::Hash(*hash);
250            if expected.insert(id) {
251                tree_hashes.push(*hash);
252            }
253        }
254        for id in &snapshot.loose_states {
255            let object_id = PackObjectId::StateId(*id);
256            if expected.insert(object_id) {
257                state_ids.push(*id);
258            }
259        }
260        logical_bytes = logical_bytes.saturating_add(add_compact_metadata(
261            &self.store,
262            &mut builder,
263            &state_ids,
264            &tree_hashes,
265            &blob_hashes,
266            context,
267            &mut corrupt_first,
268        )?);
269        builder.finalize()?;
270        Ok((expected, logical_bytes))
271    }
272}
273
274impl RepackOperation for FsRepackOperation {
275    fn key(&self) -> String {
276        self.store.root().display().to_string()
277    }
278
279    fn inspect(&self) -> std::result::Result<RepackInventory, RepackError> {
280        self.inventory().map_err(RepackError::operation)
281    }
282
283    fn run(&self, context: &RepackContext) -> std::result::Result<RepackOutcome, RepackError> {
284        self.execute(context)
285    }
286}