/src/gitoxide/gix-pack/src/cache/delta/traverse/resolve.rs
Line | Count | Source |
1 | | use std::sync::atomic::{AtomicBool, Ordering}; |
2 | | |
3 | | use gix_features::{ |
4 | | progress::Progress, |
5 | | threading::{self, OwnShared}, |
6 | | }; |
7 | | |
8 | | use crate::{ |
9 | | cache::delta::{ |
10 | | Item, |
11 | | traverse::{Context, Error, util::ItemSliceSync}, |
12 | | }, |
13 | | data, |
14 | | data::EntryRange, |
15 | | }; |
16 | | |
17 | | mod node { |
18 | | use crate::cache::delta::{Item, traverse::util::ItemSliceSync}; |
19 | | |
20 | | /// A node in a delta tree, with exclusive access to its item data. |
21 | | pub(crate) struct Node<'a, T: Send> { |
22 | | // SAFETY INVARIANT: see Node::new(). That function is the only one used |
23 | | // to create or modify these fields. |
24 | | item: &'a mut Item<T>, |
25 | | child_items: &'a ItemSliceSync<'a, Item<T>>, |
26 | | } |
27 | | |
28 | | impl<'a, T: Send> Node<'a, T> { |
29 | | /// SAFETY: `item.children` must uniquely reference elements in `child_items` that no other live item does. |
30 | | /// All child items must uphold the same invariant. |
31 | | #[expect(unsafe_code)] |
32 | 0 | pub(super) unsafe fn new(item: &'a mut Item<T>, child_items: &'a ItemSliceSync<'a, Item<T>>) -> Self { |
33 | 0 | Node { item, child_items } |
34 | 0 | } |
35 | | |
36 | | /// Return the pack byte range used to resolve this entry's header and compressed data. |
37 | 0 | pub fn entry_slice(&self) -> crate::data::EntryRange { |
38 | 0 | self.item.offset..self.item.next_offset |
39 | 0 | } |
40 | | |
41 | | /// Return the data associated with this node. |
42 | 0 | pub fn data(&mut self) -> &mut T { |
43 | 0 | &mut self.item.data |
44 | 0 | } |
45 | | |
46 | | /// Return true if this node is a base for other deltas. |
47 | 0 | pub fn has_children(&self) -> bool { |
48 | 0 | !self.item.children().is_empty() |
49 | 0 | } |
50 | | |
51 | 0 | pub fn add_children(&mut self, children: impl IntoIterator<Item = u32>) { |
52 | 0 | self.item.extend_children(children); |
53 | 0 | } |
54 | | |
55 | | /// Transform this node into an iterator over its children. |
56 | 0 | pub fn into_child_iter(self) -> impl Iterator<Item = Node<'a, T>> + 'a { |
57 | 0 | let children = self.child_items; |
58 | | #[expect(unsafe_code)] |
59 | 0 | self.item.children().iter().map(move |&index| { |
60 | | // SAFETY: Tree guarantees that each child index belongs to exactly one parent. |
61 | 0 | let item = unsafe { children.get_mut(index as usize) }; |
62 | | // SAFETY: The child inherits the same uniqueness guarantee. |
63 | 0 | unsafe { Node::new(item, children) } |
64 | 0 | }) |
65 | 0 | } |
66 | | } |
67 | | } |
68 | | |
69 | | use node::Node; |
70 | | |
71 | 0 | fn attach_ref_delta_children<T: Send>( |
72 | 0 | node: &mut Node<'_, T>, |
73 | 0 | entry: &data::Entry, |
74 | 0 | decompressed: &[u8], |
75 | 0 | ref_delta_children: Option<&super::SharedRefDeltaChildren>, |
76 | 0 | object_hash: gix_hash::Kind, |
77 | 0 | ) -> Result<(), Error> { |
78 | 0 | let Some(ref_delta_children) = ref_delta_children else { |
79 | 0 | return Ok(()); |
80 | | }; |
81 | | // Avoid hashing every remaining object once all pending ref-deltas have found their bases. |
82 | 0 | if threading::lock(ref_delta_children).is_empty() { |
83 | 0 | return Ok(()); |
84 | 0 | } |
85 | | |
86 | 0 | let kind = entry.header.as_kind().expect("a fully resolved object has a base kind"); |
87 | 0 | let id = gix_object::compute_hash(object_hash, kind, decompressed)?; |
88 | 0 | if let Some(children) = threading::lock(ref_delta_children).remove(&id) { |
89 | 0 | node.add_children(children); |
90 | 0 | } |
91 | 0 | Ok(()) |
92 | 0 | } |
93 | | |
94 | | /// A parsed entry and its decompressed bytes, ready to serve as a delta base. |
95 | | struct ResolvedBase { |
96 | | /// The pack entry, with a delta header replaced by its resolved object kind. |
97 | | entry: data::Entry, |
98 | | /// The pack offset immediately after the entry. |
99 | | entry_end: u64, |
100 | | /// The fully resolved object bytes. |
101 | | bytes: Vec<u8>, |
102 | | } |
103 | | |
104 | | /// A resolved base shared by sibling work items. |
105 | | /// |
106 | | /// [`OwnShared`] uses an `Arc` for parallel builds and an `Rc` otherwise. Once all siblings have released their clones, |
107 | | /// the task holding the sole reference can use [`OwnShared::try_unwrap()`] to recover the base and reuse its `Vec` |
108 | | /// allocation as scratch space. |
109 | | type SharedResolvedBase = OwnShared<ResolvedBase>; |
110 | | |
111 | | /// A schedulable delta-tree node. |
112 | | /// |
113 | | /// Work items can move between workers because each [`Node`] grants exclusive access to one item, while siblings |
114 | | /// share their parent only through an immutable [`SharedResolvedBase`]. |
115 | | struct WorkItem<'a, T: Send> { |
116 | | /// The traversal level, with roots at level `0`. |
117 | | level: u16, |
118 | | /// The exclusive handle to the tree item being resolved. |
119 | | node: Node<'a, T>, |
120 | | /// The resolved parent's entry and bytes needed to apply this node's delta, or `None` for roots. |
121 | | parent: Option<SharedResolvedBase>, |
122 | | } |
123 | | |
124 | | /// Resolve all delta trees from a shared, lock-free work-stealing pool. |
125 | | /// It's `unsafe` as there is safety-constraints on `items` and `child_items`. |
126 | | /// |
127 | | /// SAFETY: `items` and `child_items` must originate from the same [`crate::cache::delta::Tree`]. |
128 | | #[expect(clippy::too_many_arguments, unsafe_code)] |
129 | | #[deny(unsafe_op_in_unsafe_fn)] |
130 | 0 | pub(super) unsafe fn all<T, F, MBFN, E, R>( |
131 | 0 | items: &mut [Item<T>], |
132 | 0 | child_items: &ItemSliceSync<'_, Item<T>>, |
133 | 0 | thread_limit: Option<usize>, |
134 | 0 | num_objects: usize, |
135 | 0 | objects: gix_features::progress::StepShared, |
136 | 0 | size: gix_features::progress::StepShared, |
137 | 0 | progress: &dyn Progress, |
138 | 0 | resolve: F, |
139 | 0 | resolve_data: &R, |
140 | 0 | modify_base: MBFN, |
141 | 0 | ref_delta_children: Option<super::SharedRefDeltaChildren>, |
142 | 0 | object_hash: gix_hash::Kind, |
143 | 0 | alloc_limit_bytes: Option<usize>, |
144 | 0 | should_interrupt: &AtomicBool, |
145 | 0 | ) -> Result<(), Error> |
146 | 0 | where |
147 | 0 | T: Send, |
148 | 0 | R: Send + Sync, |
149 | 0 | F: for<'r> Fn(EntryRange, &'r R) -> Option<&'r [u8]> + Send + Clone, |
150 | 0 | MBFN: FnMut(&mut T, &dyn Progress, Context<'_>) -> Result<(), E> + Send + Clone, |
151 | 0 | E: std::error::Error + Send + Sync + 'static, |
152 | | { |
153 | 0 | let work = items |
154 | 0 | .iter_mut() |
155 | 0 | .map(|item| { |
156 | | // SAFETY: Required from the caller, and each root item is unique. |
157 | | #[expect(unsafe_code)] |
158 | 0 | let node = unsafe { Node::new(item, child_items) }; |
159 | 0 | WorkItem { |
160 | 0 | level: 0, |
161 | 0 | node, |
162 | 0 | parent: None, |
163 | 0 | } |
164 | 0 | }) |
165 | 0 | .collect::<Vec<_>>(); |
166 | | |
167 | | #[cfg(feature = "parallel")] |
168 | | { |
169 | | resolve_parallel( |
170 | | gix_features::parallel::num_threads(thread_limit).min(num_objects), |
171 | | work, |
172 | | objects, |
173 | | size, |
174 | | progress, |
175 | | resolve, |
176 | | resolve_data, |
177 | | modify_base, |
178 | | ref_delta_children, |
179 | | object_hash, |
180 | | alloc_limit_bytes, |
181 | | should_interrupt, |
182 | | ) |
183 | | } |
184 | | #[cfg(not(feature = "parallel"))] |
185 | | { |
186 | 0 | let _ = thread_limit; |
187 | 0 | let _ = num_objects; |
188 | 0 | resolve_serial( |
189 | 0 | work, |
190 | 0 | objects, |
191 | 0 | size, |
192 | 0 | progress, |
193 | 0 | resolve, |
194 | 0 | resolve_data, |
195 | 0 | modify_base, |
196 | 0 | ref_delta_children, |
197 | 0 | object_hash, |
198 | 0 | alloc_limit_bytes, |
199 | 0 | should_interrupt, |
200 | | ) |
201 | | } |
202 | 0 | } |
203 | | |
204 | | /// Resolve all work on the current thread using `work` as a LIFO stack. |
205 | | /// |
206 | | /// `work` initially contains only roots. Resolving a node adds a [`WorkItem`] for each child to the same stack consumed by |
207 | | /// the loop, so processing continues through dynamically scheduled descendants, including ref-delta children attached |
208 | | /// during resolution, until both they and the remaining roots are exhausted. LIFO order makes children of the current |
209 | | /// node run before roots that were already waiting. |
210 | | #[cfg(not(feature = "parallel"))] |
211 | | #[expect(clippy::too_many_arguments)] |
212 | 0 | fn resolve_serial<T, F, MBFN, E, R>( |
213 | 0 | mut work: Vec<WorkItem<'_, T>>, |
214 | 0 | objects: gix_features::progress::StepShared, |
215 | 0 | size: gix_features::progress::StepShared, |
216 | 0 | progress: &dyn Progress, |
217 | 0 | resolve: F, |
218 | 0 | resolve_data: &R, |
219 | 0 | mut modify_base: MBFN, |
220 | 0 | ref_delta_children: Option<super::SharedRefDeltaChildren>, |
221 | 0 | object_hash: gix_hash::Kind, |
222 | 0 | alloc_limit_bytes: Option<usize>, |
223 | 0 | should_interrupt: &AtomicBool, |
224 | 0 | ) -> Result<(), Error> |
225 | 0 | where |
226 | 0 | T: Send, |
227 | 0 | R: Send + Sync, |
228 | 0 | F: for<'r> Fn(EntryRange, &'r R) -> Option<&'r [u8]> + Send + Clone, |
229 | 0 | MBFN: FnMut(&mut T, &dyn Progress, Context<'_>) -> Result<(), E> + Send + Clone, |
230 | 0 | E: std::error::Error + Send + Sync + 'static, |
231 | | { |
232 | 0 | let mut delta_bytes = Vec::new(); |
233 | 0 | let mut fully_resolved_delta_bytes = Vec::new(); |
234 | 0 | let mut inflate = gix_zlib::Inflate::default(); |
235 | 0 | while let Some(task) = work.pop() { |
236 | 0 | if should_interrupt.load(Ordering::Relaxed) { |
237 | 0 | return Err(Error::Interrupted); |
238 | 0 | } |
239 | 0 | resolve_task( |
240 | 0 | task, |
241 | 0 | &mut delta_bytes, |
242 | 0 | &mut fully_resolved_delta_bytes, |
243 | 0 | &mut inflate, |
244 | 0 | progress, |
245 | 0 | &resolve, |
246 | 0 | resolve_data, |
247 | 0 | &mut modify_base, |
248 | 0 | ref_delta_children.as_ref(), |
249 | 0 | object_hash, |
250 | 0 | alloc_limit_bytes, |
251 | 0 | &objects, |
252 | 0 | &size, |
253 | 0 | |child| work.push(child), |
254 | 0 | )?; |
255 | | } |
256 | 0 | Ok(()) |
257 | 0 | } |
258 | | |
259 | | /// Resolve work in parallel with per-worker LIFO queues and work stealing. |
260 | | /// |
261 | | /// Roots start in a shared queue. Each worker prefers children on its own queue, then steals from peer queues, and finally |
262 | | /// takes another root. This favors completing active trees before starting more roots. |
263 | | /// |
264 | | /// Newly discovered children are scheduled on the current worker, where idle peers can steal them and help with an active |
265 | | /// tree. A worker that temporarily finds no task yields while other work is queued or in progress, since it may expose |
266 | | /// more descendants, and exits only when no work remains. |
267 | | #[cfg(feature = "parallel")] |
268 | | #[expect(clippy::too_many_arguments)] |
269 | | fn resolve_parallel<T, F, MBFN, E, R>( |
270 | | num_threads: usize, |
271 | | work: Vec<WorkItem<'_, T>>, |
272 | | objects: gix_features::progress::StepShared, |
273 | | size: gix_features::progress::StepShared, |
274 | | progress: &dyn Progress, |
275 | | resolve: F, |
276 | | resolve_data: &R, |
277 | | modify_base: MBFN, |
278 | | ref_delta_children: Option<super::SharedRefDeltaChildren>, |
279 | | object_hash: gix_hash::Kind, |
280 | | alloc_limit_bytes: Option<usize>, |
281 | | should_interrupt: &AtomicBool, |
282 | | ) -> Result<(), Error> |
283 | | where |
284 | | T: Send, |
285 | | R: Send + Sync, |
286 | | F: for<'r> Fn(EntryRange, &'r R) -> Option<&'r [u8]> + Send + Clone, |
287 | | MBFN: FnMut(&mut T, &dyn Progress, Context<'_>) -> Result<(), E> + Send + Clone, |
288 | | E: std::error::Error + Send + Sync + 'static, |
289 | | { |
290 | | use std::sync::atomic::AtomicUsize; |
291 | | |
292 | | if num_threads == 0 { |
293 | | return Ok(()); |
294 | | } |
295 | | let roots = crossbeam_deque::Injector::new(); |
296 | | let remaining = AtomicUsize::new(work.len()); |
297 | | for task in work { |
298 | | roots.push(task); |
299 | | } |
300 | | let workers: Vec<_> = (0..num_threads).map(|_| crossbeam_deque::Worker::new_lifo()).collect(); |
301 | | let stealers: Vec<_> = workers.iter().map(crossbeam_deque::Worker::stealer).collect(); |
302 | | let abort = AtomicBool::new(false); |
303 | | |
304 | | gix_features::parallel::threads(|scope| { |
305 | | let mut handles = Vec::with_capacity(num_threads); |
306 | | for (tid, worker) in workers.into_iter().enumerate() { |
307 | | let result = gix_features::parallel::build_thread() |
308 | | .name(format!("gix-pack.traverse_deltas.{tid}")) |
309 | | .spawn_scoped(scope, { |
310 | | let stealers = &stealers; |
311 | | let roots = &roots; |
312 | | let remaining = &remaining; |
313 | | let abort = &abort; |
314 | | let objects = &objects; |
315 | | let size = &size; |
316 | | let resolve = resolve.clone(); |
317 | | let mut modify_base = modify_base.clone(); |
318 | | let ref_delta_children = ref_delta_children.clone(); |
319 | | move || { |
320 | | // Make sure we never deadlock because a panicking worker can't update `remaining` anymore. |
321 | | let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { |
322 | | let mut delta_bytes = Vec::new(); |
323 | | let mut fully_resolved_delta_bytes = Vec::new(); |
324 | | let mut inflate = gix_zlib::Inflate::default(); |
325 | | loop { |
326 | | if abort.load(Ordering::Relaxed) { |
327 | | return Ok(()); |
328 | | } |
329 | | if should_interrupt.load(Ordering::Relaxed) { |
330 | | abort.store(true, Ordering::Relaxed); |
331 | | return Err(Error::Interrupted); |
332 | | } |
333 | | let Some(task) = steal(&worker, stealers, roots) else { |
334 | | if remaining.load(Ordering::Acquire) == 0 { |
335 | | return Ok(()); |
336 | | } |
337 | | std::thread::yield_now(); |
338 | | continue; |
339 | | }; |
340 | | |
341 | | let task_result = resolve_task( |
342 | | task, |
343 | | &mut delta_bytes, |
344 | | &mut fully_resolved_delta_bytes, |
345 | | &mut inflate, |
346 | | progress, |
347 | | &resolve, |
348 | | resolve_data, |
349 | | &mut modify_base, |
350 | | ref_delta_children.as_ref(), |
351 | | object_hash, |
352 | | alloc_limit_bytes, |
353 | | objects, |
354 | | size, |
355 | | |child| { |
356 | | remaining.fetch_add(1, Ordering::Release); |
357 | | worker.push(child); |
358 | | }, |
359 | | ); |
360 | | remaining.fetch_sub(1, Ordering::AcqRel); |
361 | | if let Err(err) = task_result { |
362 | | abort.store(true, Ordering::Relaxed); |
363 | | return Err(err); |
364 | | } |
365 | | } |
366 | | })); |
367 | | if result.is_err() { |
368 | | abort.store(true, Ordering::Relaxed); |
369 | | } |
370 | | result.unwrap_or_else(|payload| std::panic::resume_unwind(payload)) |
371 | | } |
372 | | }); |
373 | | match result { |
374 | | Ok(handle) => handles.push(handle), |
375 | | Err(err) => { |
376 | | abort.store(true, Ordering::Relaxed); |
377 | | for handle in handles { |
378 | | if let Err(payload) = handle.join() { |
379 | | std::panic::resume_unwind(payload); |
380 | | } |
381 | | } |
382 | | return Err(Error::SpawnThread(err)); |
383 | | } |
384 | | } |
385 | | } |
386 | | |
387 | | let mut error = None; |
388 | | for handle in handles { |
389 | | match handle.join() { |
390 | | Ok(Err(err)) if error.is_none() => error = Some(err), |
391 | | Ok(_) => {} |
392 | | Err(payload) => std::panic::resume_unwind(payload), |
393 | | } |
394 | | } |
395 | | error.map_or(Ok(()), Err) |
396 | | }) |
397 | | } |
398 | | |
399 | | /// Take local work or steal it from another worker or the shared root queue, in that order. |
400 | | /// |
401 | | /// Work in peer queues belongs to already active trees and retains shared base buffers. Preferring it advances those |
402 | | /// trees toward completion and releases their buffers before another root starts a new tree. |
403 | | #[cfg(feature = "parallel")] |
404 | | fn steal<T>( |
405 | | worker: &crossbeam_deque::Worker<T>, |
406 | | stealers: &[crossbeam_deque::Stealer<T>], |
407 | | roots: &crossbeam_deque::Injector<T>, |
408 | | ) -> Option<T> { |
409 | | if let Some(task) = worker.pop() { |
410 | | return Some(task); |
411 | | } |
412 | | loop { |
413 | | let mut retry = false; |
414 | | for stealer in stealers { |
415 | | match stealer.steal() { |
416 | | crossbeam_deque::Steal::Success(task) => return Some(task), |
417 | | crossbeam_deque::Steal::Retry => retry = true, |
418 | | crossbeam_deque::Steal::Empty => {} |
419 | | } |
420 | | } |
421 | | match roots.steal() { |
422 | | crossbeam_deque::Steal::Success(task) => return Some(task), |
423 | | crossbeam_deque::Steal::Retry => retry = true, |
424 | | crossbeam_deque::Steal::Empty => {} |
425 | | } |
426 | | if !retry { |
427 | | return None; |
428 | | } |
429 | | } |
430 | | } |
431 | | |
432 | | /// Resolve one work item using scratch buffers owned by the calling worker. |
433 | | /// |
434 | | /// `delta_bytes` holds inflated delta instructions, while `fully_resolved_delta_bytes` receives the object produced by |
435 | | /// applying them. Passing both by mutable reference preserves their allocations across tasks handled by the same worker. |
436 | | /// |
437 | | /// A resolved delta that has children moves its output buffer into [`SharedResolvedBase`]. Once the last sibling releases |
438 | | /// that base, [`OwnShared::try_unwrap()`] can recover its allocation; the largest available delta buffer is retained as |
439 | | /// `fully_resolved_delta_bytes` for a later task. Root objects deliberately use a task-local buffer instead. |
440 | | /// |
441 | | /// `push` separates discovering children from scheduling them. This function calls it with each child [`WorkItem`]; |
442 | | /// serial traversal pushes that item onto its `Vec`, while parallel traversal pushes it onto the current worker's deque |
443 | | /// and updates the shared count of unfinished work. |
444 | | #[expect(clippy::too_many_arguments)] |
445 | 0 | fn resolve_task<'a, T, F, MBFN, E, R>( |
446 | 0 | WorkItem { |
447 | 0 | level, |
448 | 0 | mut node, |
449 | 0 | parent, |
450 | 0 | }: WorkItem<'a, T>, |
451 | 0 | delta_bytes: &mut Vec<u8>, |
452 | 0 | fully_resolved_delta_bytes: &mut Vec<u8>, |
453 | 0 | inflate: &mut gix_zlib::Inflate, |
454 | 0 | progress: &dyn Progress, |
455 | 0 | resolve: &F, |
456 | 0 | resolve_data: &R, |
457 | 0 | modify_base: &mut MBFN, |
458 | 0 | ref_delta_children: Option<&super::SharedRefDeltaChildren>, |
459 | 0 | object_hash: gix_hash::Kind, |
460 | 0 | alloc_limit_bytes: Option<usize>, |
461 | 0 | objects: &gix_features::progress::StepShared, |
462 | 0 | size: &gix_features::progress::StepShared, |
463 | 0 | mut push: impl FnMut(WorkItem<'a, T>), |
464 | 0 | ) -> Result<(), Error> |
465 | 0 | where |
466 | 0 | T: Send, |
467 | 0 | R: Send + Sync, |
468 | 0 | F: for<'r> Fn(EntryRange, &'r R) -> Option<&'r [u8]> + Send, |
469 | 0 | MBFN: FnMut(&mut T, &dyn Progress, Context<'_>) -> Result<(), E> + Send, |
470 | 0 | E: std::error::Error + Send + Sync + 'static, |
471 | | { |
472 | 0 | let is_root = parent.is_none(); |
473 | | // Root buffers either become shared bases or are dropped after inspection. Keeping leaf-root allocations out of |
474 | | // worker scratch avoids retaining an occasionally huge capacity for the worker's lifetime. |
475 | 0 | let mut root_bytes = Vec::new(); |
476 | 0 | let (entry, entry_end) = if let Some(parent) = parent.as_ref() { |
477 | 0 | let (mut entry, entry_end) = decompress_from_resolver( |
478 | 0 | node.entry_slice(), |
479 | 0 | delta_bytes, |
480 | 0 | inflate, |
481 | 0 | resolve, |
482 | 0 | resolve_data, |
483 | 0 | object_hash, |
484 | 0 | alloc_limit_bytes, |
485 | 0 | )?; |
486 | 0 | let (base_size, consumed) = data::delta::decode_header_size(delta_bytes)?; |
487 | 0 | let base_size = decoded_size_limited(base_size, alloc_limit_bytes)?; |
488 | 0 | if parent.bytes.len() != base_size { |
489 | 0 | return Err(data::delta::apply::Error::Corrupt { |
490 | 0 | message: "delta base size does not match base object size", |
491 | 0 | } |
492 | 0 | .into()); |
493 | 0 | } |
494 | 0 | let (result_size, result_header_size) = data::delta::decode_header_size(&delta_bytes[consumed..])?; |
495 | 0 | let result_size = decoded_size_limited(result_size, alloc_limit_bytes)?; |
496 | 0 | resize_with_limit(fully_resolved_delta_bytes, result_size, alloc_limit_bytes)?; |
497 | 0 | data::delta::apply( |
498 | 0 | &parent.bytes, |
499 | 0 | fully_resolved_delta_bytes, |
500 | 0 | &delta_bytes[consumed + result_header_size..], |
501 | 0 | )?; |
502 | 0 | entry.header = parent.entry.header; |
503 | 0 | (entry, entry_end) |
504 | | } else { |
505 | 0 | decompress_from_resolver( |
506 | 0 | node.entry_slice(), |
507 | 0 | &mut root_bytes, |
508 | 0 | inflate, |
509 | 0 | resolve, |
510 | 0 | resolve_data, |
511 | 0 | object_hash, |
512 | 0 | alloc_limit_bytes, |
513 | 0 | )? |
514 | | }; |
515 | | |
516 | 0 | let resolved = ResolvedBase { |
517 | 0 | entry, |
518 | 0 | entry_end, |
519 | 0 | bytes: if is_root { |
520 | 0 | root_bytes |
521 | | } else { |
522 | 0 | std::mem::take(fully_resolved_delta_bytes) |
523 | | }, |
524 | | }; |
525 | 0 | attach_ref_delta_children( |
526 | 0 | &mut node, |
527 | 0 | &resolved.entry, |
528 | 0 | &resolved.bytes, |
529 | 0 | ref_delta_children, |
530 | 0 | object_hash, |
531 | 0 | )?; |
532 | 0 | let has_children = node.has_children(); |
533 | 0 | inspect(&mut node, level, &resolved, progress, modify_base, objects, size)?; |
534 | 0 | let mut reusable = if has_children { |
535 | 0 | let resolved = OwnShared::new(resolved); |
536 | 0 | for child in node.into_child_iter() { |
537 | 0 | push(WorkItem { |
538 | 0 | level: level + 1, |
539 | 0 | node: child, |
540 | 0 | parent: Some(OwnShared::clone(&resolved)), |
541 | 0 | }); |
542 | 0 | } |
543 | 0 | None |
544 | 0 | } else if is_root { |
545 | 0 | None |
546 | | } else { |
547 | 0 | Some(resolved.bytes) |
548 | | }; |
549 | | |
550 | | // This might be a leaf, while its base buffer now is also exclusively available, |
551 | | // and if so, keep the larger buffer. |
552 | 0 | if let Some(parent) = parent { |
553 | 0 | if let Ok(parent) = OwnShared::try_unwrap(parent) { |
554 | 0 | if reusable |
555 | 0 | .as_ref() |
556 | 0 | .is_none_or(|reusable| parent.bytes.capacity() > reusable.capacity()) |
557 | 0 | { |
558 | 0 | reusable = Some(parent.bytes); |
559 | 0 | } |
560 | 0 | } |
561 | 0 | } |
562 | 0 | if let Some(reusable) = reusable { |
563 | 0 | *fully_resolved_delta_bytes = reusable; |
564 | 0 | } |
565 | 0 | fully_resolved_delta_bytes.clear(); |
566 | 0 | Ok(()) |
567 | 0 | } |
568 | | |
569 | | /// Hand one fully resolved node to the caller's inspector and record its progress. |
570 | | /// |
571 | | /// `modify_base` receives mutable access to the node's associated data and a [`Context`] containing the parsed entry, |
572 | | /// its end offset, the resolved object bytes, and its delta-tree level. Only a successful inspection increments the |
573 | | /// object counter and the total number of resolved bytes; inspector errors abort traversal as [`Error::Inspect`]. |
574 | 0 | fn inspect<T, MBFN, E>( |
575 | 0 | node: &mut Node<'_, T>, |
576 | 0 | level: u16, |
577 | 0 | resolved: &ResolvedBase, |
578 | 0 | progress: &dyn Progress, |
579 | 0 | modify_base: &mut MBFN, |
580 | 0 | objects: &gix_features::progress::StepShared, |
581 | 0 | size: &gix_features::progress::StepShared, |
582 | 0 | ) -> Result<(), Error> |
583 | 0 | where |
584 | 0 | T: Send, |
585 | 0 | MBFN: FnMut(&mut T, &dyn Progress, Context<'_>) -> Result<(), E> + Send, |
586 | 0 | E: std::error::Error + Send + Sync + 'static, |
587 | | { |
588 | 0 | modify_base( |
589 | 0 | node.data(), |
590 | 0 | progress, |
591 | 0 | Context { |
592 | 0 | entry: &resolved.entry, |
593 | 0 | entry_end: resolved.entry_end, |
594 | 0 | decompressed: &resolved.bytes, |
595 | 0 | level, |
596 | 0 | }, |
597 | 0 | ) |
598 | 0 | .map_err(|err| Box::new(err) as Box<dyn std::error::Error + Send + Sync>)?; |
599 | 0 | objects.fetch_add(1, Ordering::Relaxed); |
600 | 0 | size.fetch_add(resolved.bytes.len(), Ordering::Relaxed); |
601 | 0 | Ok(()) |
602 | 0 | } |
603 | | |
604 | | /// Resolve and decompress one pack entry into `out`, returning its parsed metadata and end offset. |
605 | | /// |
606 | | /// `resolve` keeps traversal independent of pack storage by borrowing the bytes for `slice` from caller-owned |
607 | | /// `resolve_data`, such as an in-memory buffer or mapped file. |
608 | 0 | fn decompress_from_resolver<F, R>( |
609 | 0 | slice: EntryRange, |
610 | 0 | out: &mut Vec<u8>, |
611 | 0 | inflate: &mut gix_zlib::Inflate, |
612 | 0 | resolve: &F, |
613 | 0 | resolve_data: &R, |
614 | 0 | object_hash: gix_hash::Kind, |
615 | 0 | alloc_limit_bytes: Option<usize>, |
616 | 0 | ) -> Result<(data::Entry, u64), Error> |
617 | 0 | where |
618 | 0 | F: for<'r> Fn(EntryRange, &'r R) -> Option<&'r [u8]> + Send, |
619 | | { |
620 | 0 | let bytes = resolve(slice.clone(), resolve_data).ok_or(Error::ResolveFailed { |
621 | 0 | pack_offset: slice.start, |
622 | 0 | })?; |
623 | 0 | let entry = data::Entry::from_bytes(bytes, slice.start, object_hash)?; |
624 | 0 | let compressed = &bytes[entry.header_size()..]; |
625 | 0 | let decompressed_len = decoded_size_limited(entry.decompressed_size, alloc_limit_bytes)?; |
626 | 0 | decompress_all_at_once_with(inflate, compressed, decompressed_len, out, alloc_limit_bytes)?; |
627 | 0 | Ok((entry, slice.end)) |
628 | 0 | } |
629 | | |
630 | 0 | fn decompress_all_at_once_with( |
631 | 0 | inflate: &mut gix_zlib::Inflate, |
632 | 0 | b: &[u8], |
633 | 0 | decompressed_len: usize, |
634 | 0 | out: &mut Vec<u8>, |
635 | 0 | alloc_limit_bytes: Option<usize>, |
636 | 0 | ) -> Result<(), Error> { |
637 | 0 | resize_with_limit(out, decompressed_len, alloc_limit_bytes)?; |
638 | 0 | inflate.reset(); |
639 | 0 | inflate.once(b, out).map_err(|err| Error::ZlibInflate { |
640 | 0 | source: err, |
641 | | message: "Failed to decompress entry", |
642 | 0 | })?; |
643 | 0 | Ok(()) |
644 | 0 | } |
645 | | |
646 | 0 | fn decoded_size_limited(size: u64, alloc_limit_bytes: Option<usize>) -> Result<usize, Error> { |
647 | 0 | let size: usize = size.try_into().map_err(|_| Error::OutOfMemory)?; |
648 | 0 | if alloc_limit_bytes.is_some_and(|limit| size > limit) { |
649 | 0 | return Err(Error::OutOfMemory); |
650 | 0 | } |
651 | 0 | Ok(size) |
652 | 0 | } |
653 | | |
654 | 0 | fn resize_with_limit(out: &mut Vec<u8>, len: usize, alloc_limit_bytes: Option<usize>) -> Result<(), Error> { |
655 | 0 | if alloc_limit_bytes.is_some_and(|limit| len > limit) { |
656 | 0 | return Err(Error::OutOfMemory); |
657 | 0 | } |
658 | 0 | out.try_reserve(len.saturating_sub(out.len()))?; |
659 | 0 | out.resize(len, 0); |
660 | 0 | Ok(()) |
661 | 0 | } |
662 | | |
663 | | #[cfg(test)] |
664 | | mod tests { |
665 | | use std::{ |
666 | | io::Write, |
667 | | sync::atomic::{AtomicBool, AtomicUsize, Ordering}, |
668 | | time::Duration, |
669 | | }; |
670 | | |
671 | | use gix_features::progress; |
672 | | |
673 | | use crate::{ |
674 | | cache::delta::{Tree, traverse}, |
675 | | data, |
676 | | }; |
677 | | |
678 | | #[test] |
679 | | fn traversal_resolves_children_lazily() { |
680 | | let mut pack = Vec::new(); |
681 | | let root_offset = append_entry(&mut pack, data::entry::Header::Blob, 1, b"A"); |
682 | | let first_child = append_delta(&mut pack, root_offset, b'B'); |
683 | | let second_child = append_delta(&mut pack, root_offset, b'C'); |
684 | | let first_leaf = append_delta(&mut pack, first_child, b'D'); |
685 | | let second_leaf = append_delta(&mut pack, second_child, b'E'); |
686 | | |
687 | | let mut tree = Tree::with_capacity(5).expect("capacity is small"); |
688 | | tree.add_root(root_offset, ()).expect("offsets are increasing"); |
689 | | tree.add_child(root_offset, first_child, ()) |
690 | | .expect("offsets are increasing"); |
691 | | tree.add_child(root_offset, second_child, ()) |
692 | | .expect("offsets are increasing"); |
693 | | tree.add_child(first_child, first_leaf, ()) |
694 | | .expect("offsets are increasing"); |
695 | | tree.add_child(second_child, second_leaf, ()) |
696 | | .expect("offsets are increasing"); |
697 | | |
698 | | let resolve_calls = AtomicUsize::new(0); |
699 | | let calls_at_first_child = AtomicUsize::new(usize::MAX); |
700 | | traverse( |
701 | | tree, |
702 | | &pack, |
703 | | Some(1), |
704 | | None, |
705 | | |slice, pack| { |
706 | | resolve_calls.fetch_add(1, Ordering::Relaxed); |
707 | | pack.get(slice.start as usize..slice.end as usize) |
708 | | }, |
709 | | |(), _progress, context| { |
710 | | if context.level == 1 { |
711 | | calls_at_first_child.fetch_min(resolve_calls.load(Ordering::Relaxed), Ordering::Relaxed); |
712 | | } |
713 | | Ok::<_, std::io::Error>(()) |
714 | | }, |
715 | | ) |
716 | | .expect("valid delta tree"); |
717 | | |
718 | | assert_eq!( |
719 | | calls_at_first_child.load(Ordering::Relaxed), |
720 | | 2, |
721 | | "the first child must be inspected before its siblings are materialized" |
722 | | ); |
723 | | } |
724 | | |
725 | | #[test] |
726 | | fn traversal_parallelizes_children_of_one_root() { |
727 | | let mut pack = Vec::new(); |
728 | | let root_offset = append_entry(&mut pack, data::entry::Header::Blob, 1, b"A"); |
729 | | let child_offsets: Vec<_> = (b'B'..=b'I') |
730 | | .map(|byte| append_delta(&mut pack, root_offset, byte)) |
731 | | .collect(); |
732 | | |
733 | | let mut tree = Tree::with_capacity(1 + child_offsets.len()).expect("capacity is small"); |
734 | | tree.add_root(root_offset, ()).expect("offsets are increasing"); |
735 | | for child_offset in child_offsets { |
736 | | tree.add_child(root_offset, child_offset, ()) |
737 | | .expect("offsets are increasing"); |
738 | | } |
739 | | |
740 | | let active = AtomicUsize::new(0); |
741 | | let max_active = AtomicUsize::new(0); |
742 | | traverse( |
743 | | tree, |
744 | | &pack, |
745 | | Some(2), |
746 | | None, |
747 | | |slice, pack| pack.get(slice.start as usize..slice.end as usize), |
748 | | |(), _progress, context| { |
749 | | if context.level > 0 { |
750 | | let now_active = active.fetch_add(1, Ordering::Relaxed) + 1; |
751 | | max_active.fetch_max(now_active, Ordering::Relaxed); |
752 | | std::thread::sleep(Duration::from_millis(20)); |
753 | | active.fetch_sub(1, Ordering::Relaxed); |
754 | | } |
755 | | Ok::<_, std::io::Error>(()) |
756 | | }, |
757 | | ) |
758 | | .expect("valid delta tree"); |
759 | | |
760 | | let expected = if cfg!(feature = "parallel") { 1 } else { 0 }; |
761 | | assert!( |
762 | | max_active.load(Ordering::Relaxed) > expected, |
763 | | "idle workers must help with children of the last remaining root (if in parallel mode)" |
764 | | ); |
765 | | } |
766 | | |
767 | | #[test] |
768 | | fn traversal_rejects_declared_decompressed_size_over_alloc_limit() { |
769 | | let mut pack = Vec::new(); |
770 | | let root_offset = append_entry(&mut pack, data::entry::Header::Blob, 1, b""); |
771 | | let mut tree = Tree::with_capacity(1).expect("capacity is small"); |
772 | | tree.add_root(root_offset, ()).expect("offsets are increasing"); |
773 | | |
774 | | let err = traverse_with_limit(tree, &pack).expect_err("entry size exceeds the allocation cap"); |
775 | | |
776 | | assert!( |
777 | | matches!(err, traverse::Error::OutOfMemory), |
778 | | "declared decompressed sizes above the cap must be rejected before allocation" |
779 | | ); |
780 | | } |
781 | | |
782 | | #[test] |
783 | | fn traversal_rejects_delta_base_size_over_alloc_limit() { |
784 | | let mut pack = Vec::new(); |
785 | | let root_offset = append_entry(&mut pack, data::entry::Header::Blob, 0, b""); |
786 | | |
787 | | let delta = [1, 0]; |
788 | | let child_offset = pack.len() as u64; |
789 | | append_entry( |
790 | | &mut pack, |
791 | | data::entry::Header::OfsDelta { |
792 | | base_distance: child_offset - root_offset, |
793 | | }, |
794 | | delta.len() as u64, |
795 | | &delta, |
796 | | ); |
797 | | |
798 | | let mut tree = Tree::with_capacity(2).expect("capacity is small"); |
799 | | tree.add_root(root_offset, ()).expect("offsets are increasing"); |
800 | | tree.add_child(root_offset, child_offset, ()) |
801 | | .expect("offsets are increasing"); |
802 | | |
803 | | let err = traverse_with_limit(tree, &pack).expect_err("delta base size exceeds the allocation cap"); |
804 | | |
805 | | assert!( |
806 | | matches!(err, traverse::Error::OutOfMemory), |
807 | | "delta base sizes above the cap must be rejected before comparing them with the decoded base" |
808 | | ); |
809 | | } |
810 | | |
811 | | #[test] |
812 | | fn traversal_rejects_delta_result_size_over_alloc_limit() { |
813 | | let mut pack = Vec::new(); |
814 | | let root_offset = append_entry(&mut pack, data::entry::Header::Blob, 0, b""); |
815 | | |
816 | | let delta = [0, 1, 1, b'A']; |
817 | | let child_offset = pack.len() as u64; |
818 | | append_entry( |
819 | | &mut pack, |
820 | | data::entry::Header::OfsDelta { |
821 | | base_distance: child_offset - root_offset, |
822 | | }, |
823 | | delta.len() as u64, |
824 | | &delta, |
825 | | ); |
826 | | |
827 | | let mut tree = Tree::with_capacity(2).expect("capacity is small"); |
828 | | tree.add_root(root_offset, ()).expect("offsets are increasing"); |
829 | | tree.add_child(root_offset, child_offset, ()) |
830 | | .expect("offsets are increasing"); |
831 | | |
832 | | let err = traverse_with_limit(tree, &pack).expect_err("delta result size exceeds the allocation cap"); |
833 | | |
834 | | assert!( |
835 | | matches!(err, traverse::Error::OutOfMemory), |
836 | | "delta result sizes above the cap must be rejected before resizing the output buffer" |
837 | | ); |
838 | | } |
839 | | |
840 | | fn traverse_with_limit(tree: Tree<()>, pack: &Vec<u8>) -> Result<(), traverse::Error> { |
841 | | traverse( |
842 | | tree, |
843 | | pack, |
844 | | Some(1), |
845 | | Some(0), |
846 | | |slice, pack| pack.get(slice.start as usize..slice.end as usize), |
847 | | |(), _progress, _context| Ok::<_, std::io::Error>(()), |
848 | | ) |
849 | | } |
850 | | |
851 | | fn traverse<F, MBFN>( |
852 | | tree: Tree<()>, |
853 | | pack: &Vec<u8>, |
854 | | thread_limit: Option<usize>, |
855 | | alloc_limit_bytes: Option<usize>, |
856 | | resolve: F, |
857 | | inspect: MBFN, |
858 | | ) -> Result<(), traverse::Error> |
859 | | where |
860 | | F: for<'r> Fn(data::EntryRange, &'r Vec<u8>) -> Option<&'r [u8]> + Send + Clone, |
861 | | MBFN: |
862 | | FnMut(&mut (), &dyn progress::Progress, traverse::Context<'_>) -> Result<(), std::io::Error> + Send + Clone, |
863 | | { |
864 | | let should_interrupt = AtomicBool::new(false); |
865 | | let mut size_progress = progress::Discard; |
866 | | tree.traverse( |
867 | | resolve, |
868 | | pack, |
869 | | pack.len() as u64, |
870 | | inspect, |
871 | | traverse::Options { |
872 | | object_progress: Box::new(progress::Discard), |
873 | | size_progress: &mut size_progress, |
874 | | thread_limit, |
875 | | should_interrupt: &should_interrupt, |
876 | | object_hash: gix_hash::Kind::Sha1, |
877 | | alloc_limit_bytes, |
878 | | }, |
879 | | ) |
880 | | .map(|_| ()) |
881 | | } |
882 | | |
883 | | fn append_delta(pack: &mut Vec<u8>, base_offset: data::Offset, byte: u8) -> data::Offset { |
884 | | let delta = [1, 1, 1, byte]; |
885 | | let offset = pack.len() as data::Offset; |
886 | | append_entry( |
887 | | pack, |
888 | | data::entry::Header::OfsDelta { |
889 | | base_distance: offset - base_offset, |
890 | | }, |
891 | | delta.len() as u64, |
892 | | &delta, |
893 | | ) |
894 | | } |
895 | | |
896 | | fn append_entry( |
897 | | pack: &mut Vec<u8>, |
898 | | header: data::entry::Header, |
899 | | decompressed_size: u64, |
900 | | payload: &[u8], |
901 | | ) -> data::Offset { |
902 | | let offset = pack.len() as data::Offset; |
903 | | header |
904 | | .write_to(decompressed_size, pack) |
905 | | .expect("writing an entry header to memory succeeds"); |
906 | | pack.extend(deflate(payload)); |
907 | | offset |
908 | | } |
909 | | |
910 | | fn deflate(bytes: &[u8]) -> Vec<u8> { |
911 | | let mut out = gix_zlib::stream::deflate::Write::new(Vec::new(), gix_zlib::Compression::BEST_SPEED); |
912 | | out.write_all(bytes).expect("writing to deflater succeeds"); |
913 | | out.flush().expect("flushing deflater succeeds"); |
914 | | out.into_inner() |
915 | | } |
916 | | } |