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 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
34pub 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 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 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 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 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}