Coverage Report

Created: 2026-08-05 07:37

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/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
}