1use 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
25pub struct FsRepackOperation {
32 store: FsStore,
33 excluded: HashSet<PackObjectId>,
34 #[cfg(test)]
35 corrupt_first_object: bool,
36}
37
38impl FsRepackOperation {
39 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 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 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 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}