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        npk1::{Npk1Build, Npk1BuildError, build_npk1_pack},
11    },
12    compact::add_compact_metadata,
13    cutover::{
14        acquire_repack_lock, cutover, file_len, hash_file, object_file_len,
15        preserve_commit_markers, publish_npk1,
16    },
17    staging::{BuildError, RepackSnapshot, RepackStaging, verify_staged},
18};
19use crate::store::{
20    HeddleError, ObjectStore, Result,
21    pack::{
22        ObjectType, PackObjectId, PackReader, RepackContext, RepackError, RepackInventory,
23        RepackOperation, RepackOutcome, StreamingPackBuilder,
24    },
25};
26
27type StagedBuild = (
28    HashSet<PackObjectId>,
29    HashSet<crate::object::ContentHash>,
30    u64,
31    Option<Npk1Build>,
32);
33
34/// Compact native encoder wired behind the generic background repack seam.
35///
36/// It loads trees one at a time while retaining pack-wide dictionaries and
37/// sketches, verifies every typed id and the exact expected object set, then
38/// installs the immutable replacement before retiring source packs under the
39/// pack-manager write lock. The scheduler itself has no native-pack knowledge.
40pub struct FsRepackOperation {
41    store: FsStore,
42    excluded_blobs: HashSet<crate::object::ContentHash>,
43    #[cfg(test)]
44    corrupt_first_object: bool,
45}
46
47impl FsRepackOperation {
48    /// Create a native filesystem repack payload.
49    pub fn new(store: FsStore) -> Self {
50        Self {
51            store,
52            excluded_blobs: HashSet::new(),
53            #[cfg(test)]
54            corrupt_first_object: false,
55        }
56    }
57
58    /// Exclude one blob from the replacement pack.
59    ///
60    /// This is the storage primitive used by purge: cutover retires every
61    /// source pack only after the replacement has been verified without the
62    /// excluded identity.
63    pub fn excluding_blob(mut self, hash: crate::object::ContentHash) -> Self {
64        self.excluded_blobs.insert(hash);
65        self
66    }
67
68    #[cfg(test)]
69    pub(super) fn with_corrupted_output(mut self) -> Self {
70        self.corrupt_first_object = true;
71        self
72    }
73
74    fn inventory(&self) -> Result<RepackInventory> {
75        let loose_blobs = list_hashes_from_dir(&blobs_dir(self.store.root()))?;
76        let loose_trees = list_hashes_from_dir(&trees_dir(self.store.root()))?;
77        let loose_states = list_state_ids_from_dir(&states_dir(self.store.root()))?;
78        let loose_bytes = loose_blobs
79            .iter()
80            .map(|hash| object_file_len(&blobs_dir(self.store.root()), hash))
81            .chain(
82                loose_trees
83                    .iter()
84                    .map(|hash| object_file_len(&trees_dir(self.store.root()), hash)),
85            )
86            .chain(
87                loose_states
88                    .iter()
89                    .map(|id| file_len(&state_path(self.store.root(), id))),
90            )
91            .sum();
92        let manager =
93            self.store.pack_manager().read().map_err(|_| {
94                HeddleError::Config("Failed to acquire pack manager lock".to_string())
95            })?;
96        let paths = manager.pack_file_paths();
97        let generic_pack_count = paths.len();
98        let mut pack_bytes: u64 = paths
99            .iter()
100            .flat_map(|(pack, index)| [*pack, *index])
101            .map(file_len)
102            .sum();
103        let mut ids = manager.list_all_ids()?;
104        drop(manager);
105        let npk1 =
106            self.store.npk1_manager().read().map_err(|_| {
107                HeddleError::Config("Failed to acquire NPK1 manager lock".to_string())
108            })?;
109        let npk1_paths = npk1.file_paths();
110        pack_bytes =
111            pack_bytes.saturating_add(npk1_paths.iter().map(|path| file_len(path)).sum::<u64>());
112        ids.extend(npk1.list_ids()?.into_iter().map(PackObjectId::Hash));
113        let unique = ids.iter().copied().collect::<HashSet<_>>().len() as u64;
114        Ok(RepackInventory {
115            loose_objects: (loose_blobs.len() + loose_trees.len() + loose_states.len()) as u64,
116            loose_bytes,
117            pack_count: (generic_pack_count + npk1_paths.len()) as u64,
118            pack_bytes,
119            duplicate_objects: ids.len() as u64 - unique,
120            packed_objects: ids.len() as u64,
121        })
122    }
123
124    fn execute(&self, context: &RepackContext) -> std::result::Result<RepackOutcome, RepackError> {
125        let packs = packs_dir(self.store.root());
126        fs::create_dir_all(&packs).map_err(RepackError::operation)?;
127        let _operation_lock = acquire_repack_lock(&packs, context)?;
128        context.checkpoint(0)?;
129
130        let snapshot = RepackSnapshot::capture(&self.store).map_err(RepackError::operation)?;
131        if snapshot.ids.is_empty()
132            && snapshot.loose_blobs.is_empty()
133            && snapshot.loose_trees.is_empty()
134            && snapshot.loose_states.is_empty()
135        {
136            return Ok(RepackOutcome::default());
137        }
138
139        let staging = RepackStaging::new(&packs).map_err(RepackError::operation)?;
140        let (expected_generic, expected_trees, logical_bytes, npk1_build) = self
141            .build_staged(&snapshot, &staging, context)
142            .map_err(|error| match error {
143                BuildError::Cancelled(error) => error,
144                BuildError::Store(error) => RepackError::operation(error),
145            })?;
146        verify_staged(&staging, &expected_generic, &expected_trees, context)?;
147        context.checkpoint(0)?;
148
149        // Cutover starts here and is intentionally non-cancellable. The new
150        // immutable files are durable before any source path is retired.
151        let (new_npk1_name, npk1_preexisting) = if npk1_build.is_some() {
152            let (name, preexisting) =
153                publish_npk1(&packs, &staging.npk1).map_err(RepackError::operation)?;
154            (Some(name), preexisting)
155        } else {
156            (None, false)
157        };
158        let new_pack_name = hash_file(&staging.pack).map_err(RepackError::operation)?;
159        let replacement_preexisting = packs.join(format!("{new_pack_name}.pack")).exists()
160            && packs.join(format!("{new_pack_name}.idx")).exists();
161        ObjectStore::install_pack_streaming(&self.store, &staging.pack, &staging.index)
162            .map_err(RepackError::operation)?;
163        preserve_commit_markers(&packs, &new_pack_name, &snapshot.commit_artifact_ids)
164            .map_err(RepackError::operation)?;
165        let cutover = cutover(
166            &self.store,
167            &snapshot,
168            &new_pack_name,
169            new_npk1_name.as_deref(),
170            replacement_preexisting,
171            npk1_preexisting,
172        )
173        .map_err(RepackError::operation)?;
174        // The replacement is authoritative before loose copies are pruned.
175        // `prune_loose_objects_impl` rechecks exact packed content, so a
176        // concurrent loose write is removed only when the durable pack has it.
177        let (_, loose_bytes) = self
178            .store
179            .prune_loose_objects_impl()
180            .map_err(RepackError::operation)?;
181        let reclaimed = cutover
182            .removed_pack_bytes
183            .saturating_add(loose_bytes)
184            .saturating_sub(cutover.replacement_bytes);
185
186        Ok(RepackOutcome {
187            objects_repacked: (expected_generic.len() + expected_trees.len()) as u64,
188            bytes_repacked: logical_bytes,
189            bytes_reclaimed: reclaimed,
190        })
191    }
192
193    fn build_staged(
194        &self,
195        snapshot: &RepackSnapshot,
196        staging: &RepackStaging,
197        context: &RepackContext,
198    ) -> std::result::Result<StagedBuild, BuildError> {
199        let pack_file = OpenOptions::new()
200            .create_new(true)
201            .read(true)
202            .write(true)
203            .open(&staging.pack)?;
204        let mut compression = self.store.compression();
205        compression.max_delta_size = 0;
206        let mut builder = StreamingPackBuilder::new(
207            pack_file,
208            staging.index.clone(),
209            compression,
210            staging.buckets.clone(),
211        )?;
212        let mut expected = HashSet::new();
213        let mut expected_trees = HashSet::new();
214        let mut state_ids = Vec::new();
215        let mut tree_hashes = Vec::new();
216        let mut blob_hashes = Vec::new();
217        #[cfg(test)]
218        let mut corrupt_first = self.corrupt_first_object;
219        #[cfg(not(test))]
220        let mut corrupt_first = false;
221        let mut logical_bytes = 0u64;
222
223        for (pack, index) in &snapshot.old_pack_files {
224            let reader = PackReader::open(pack, index)?;
225            let mut checkpoint_error = None;
226            let visit = reader.visit_objects(|id, object_type, data| {
227                if object_type == ObjectType::Blob
228                    && matches!(id, PackObjectId::Hash(hash) if self.excluded_blobs.contains(&hash))
229                {
230                    if let Err(error) = context.checkpoint(data.len() as u64) {
231                        checkpoint_error = Some(error);
232                        return Err(HeddleError::InvalidObject(
233                            "repack source walk interrupted".to_string(),
234                        ));
235                    }
236                    return Ok(());
237                }
238                match (id, object_type) {
239                    (PackObjectId::Hash(hash), ObjectType::Tree) => {
240                        if expected_trees.insert(hash) {
241                            tree_hashes.push(hash);
242                        }
243                    }
244                    (PackObjectId::Hash(hash), ObjectType::Blob) => {
245                        if !expected.insert(id) {
246                            return Ok(());
247                        }
248                        blob_hashes.push(hash);
249                    }
250                    (PackObjectId::StateId(state_id), ObjectType::State) => {
251                        if !expected.insert(id) {
252                            return Ok(());
253                        }
254                        state_ids.push(state_id);
255                    }
256                    _ => {
257                        if !expected.insert(id) {
258                            return Ok(());
259                        }
260                        let mut data = data.to_vec();
261                        if corrupt_first {
262                            data.push(0xff);
263                            corrupt_first = false;
264                        }
265                        logical_bytes = logical_bytes.saturating_add(data.len() as u64);
266                        builder.add_id(id, object_type, data)?;
267                    }
268                }
269                if let Err(error) = context.checkpoint(data.len() as u64) {
270                    checkpoint_error = Some(error);
271                    return Err(HeddleError::InvalidObject(
272                        "repack source walk interrupted".to_string(),
273                    ));
274                }
275                Ok(())
276            });
277            if let Some(error) = checkpoint_error {
278                return Err(BuildError::Cancelled(error));
279            }
280            visit?;
281        }
282        for hash in &snapshot.npk1_trees {
283            if expected_trees.insert(*hash) {
284                tree_hashes.push(*hash);
285            }
286        }
287        for hash in &snapshot.loose_blobs {
288            let id = PackObjectId::Hash(*hash);
289            if self.excluded_blobs.contains(hash) {
290                continue;
291            }
292            if expected.insert(id) {
293                blob_hashes.push(*hash);
294            }
295        }
296        for hash in &snapshot.loose_trees {
297            if expected_trees.insert(*hash) {
298                tree_hashes.push(*hash);
299            }
300        }
301        for id in &snapshot.loose_states {
302            let object_id = PackObjectId::StateId(*id);
303            if expected.insert(object_id) {
304                state_ids.push(*id);
305            }
306        }
307        let compact = add_compact_metadata(
308            &self.store,
309            &mut builder,
310            &state_ids,
311            &tree_hashes,
312            &blob_hashes,
313            context,
314            &mut corrupt_first,
315        )?;
316        logical_bytes = logical_bytes.saturating_add(compact.logical_bytes);
317        builder.finalize()?;
318        let npk1_build = if compact.tree_order.is_empty() {
319            None
320        } else {
321            let build = build_npk1_pack(
322                &self.store,
323                &compact.tree_order,
324                &compact.tree_parents,
325                &staging.npk1,
326                context,
327            )
328            .map_err(|error| match error {
329                Npk1BuildError::Store(error) => BuildError::Store(error),
330                Npk1BuildError::Cancelled(error) => BuildError::Cancelled(error),
331            })?;
332            logical_bytes = logical_bytes.saturating_add(build.logical_bytes);
333            Some(build)
334        };
335        Ok((expected, expected_trees, logical_bytes, npk1_build))
336    }
337}
338
339impl RepackOperation for FsRepackOperation {
340    fn key(&self) -> String {
341        self.store.root().display().to_string()
342    }
343
344    fn inspect(&self) -> std::result::Result<RepackInventory, RepackError> {
345        self.inventory().map_err(RepackError::operation)
346    }
347
348    fn run(&self, context: &RepackContext) -> std::result::Result<RepackOutcome, RepackError> {
349        self.execute(context)
350    }
351}