Skip to main content

mkit_core/pack/
window.rs

1//! Windowed, resumable, sans-IO pack decoding (SPEC-PACKFILE ยง11).
2//!
3//! Entries are provisional until [`Step::Done`]. On any error the caller must
4//! discard everything staged by this run, including entries from earlier slices.
5//! This reader does not validate objects, resolve deltas, or write to a store.
6//!
7//! Wrong feed ranges, invalid window geometry, bad cursors, and checksum or
8//! pack-id mismatches use [`PackError::PackfileCorrupted`]. Budget or allocation
9//! failures use [`PackError::PackfileTooLarge`]. Framing and trailing-data errors
10//! match [`super::PackEntries`]; zstd uses `ZstdEntryTruncated`,
11//! `DecompressedSizeOverCap`, `DecompressedSizeMismatch`, and `ZstdDecompress`.
12//! The decoded budget includes owned entry buffers, including carried wire bytes
13//! plus decoded output when both are live. Decoder-internal memory is additional
14//! (in particular the ruzstd ring buffer; see the parent module).
15//!
16//! A resumed run re-fetches its current window and checks the checkpoint's
17//! prefix commitment before yielding entries. Completed windows are not re-read.
18//! Without a requested pack id, the initial run first fetches the trailer
19//! window(s), retaining that anchor in its cursor to bind even an unseen suffix.
20
21use super::{
22    DecodeLimits, MAGIC, MAX_ENTRIES, MAX_TOTAL_PAYLOAD, PackEntry, PackError, decode_payload,
23    zstd_claim,
24};
25use crate::hash::Hash;
26use std::borrow::Cow;
27mod cursor;
28mod tree;
29pub use cursor::WindowCursor;
30use tree::Tree;
31
32/// A range which must be supplied in full to [`WindowReader::feed`].
33#[derive(Debug, Clone, Copy, PartialEq, Eq)]
34#[non_exhaustive]
35pub struct WindowRequest {
36    pub offset: u64,
37    pub len: u64,
38}
39
40/// Where the entry a reader last yielded sat in the pack: its complete
41/// frame, type and length prefix included, exactly as `DecodedEntry` reports
42/// it for the buffered decoder.
43#[derive(Debug, Clone, Copy, PartialEq, Eq)]
44#[non_exhaustive]
45pub struct FrameInfo {
46    /// Offset of the complete frame in the pack.
47    pub offset: u64,
48    /// Length of the complete frame, including its five header bytes.
49    pub length: u64,
50    /// Encoded frame type (`0x00`, `0x02`, `0x03` or `0x04`).
51    pub wire_type: u8,
52}
53
54/// The next action for a window-reader driver.
55#[derive(Debug)]
56#[non_exhaustive]
57pub enum Step {
58    /// Fetch exactly this range, then call `feed`.
59    NeedWindow(WindowRequest),
60    /// An owned entry, provisional until `Done`.
61    Entry(PackEntry<'static>),
62    /// Framing, trailer, and optional requested pack id are verified.
63    Done(WindowSummary),
64}
65
66/// Verified pack metadata; raw-only means every wire type was `0x00`.
67#[derive(Debug, Clone, PartialEq, Eq)]
68#[non_exhaustive]
69pub struct WindowSummary {
70    pub version: u32,
71    pub entry_count: u32,
72    pub raw_only: bool,
73    pub first_non_raw: Option<u32>,
74}
75
76#[derive(Debug, Clone, Copy, PartialEq, Eq)]
77enum Phase {
78    Anchor,
79    Header,
80    Boundary,
81    Frame,
82    Payload,
83    Finish,
84    Done,
85    Failed,
86}
87
88/// An I/O-free decoder retaining only the current window and current entry.
89///
90/// **Binding.** `Done` means every entry yielded across the whole run chain (the original run and every resume
91/// through its cursors) is an entry, in order, of the one pack whose bytes hash to the verified pack id: the trailer,
92/// and `expected_pack_id` when set. A cursor from pack A used on a source that yields different bytes for any
93/// not-yet-verified range fails with `PackfileCorrupted`, never `Done`. The reader does not re-read ranges it has
94/// already verified. **Keeping the source immutable across resumes is the caller's job** (informative: WP-4.8 binds
95/// R2 range reads to the object's etag). A source whose already-verified prefix changed after verification can still
96/// reach `Done`, but only with entries of the verified pack.
97#[derive(Debug)]
98pub struct WindowReader {
99    state: WindowCursor,
100    limits: DecodeLimits,
101    phase: Phase,
102    window: Vec<u8>,
103    start: u64,
104    end: u64,
105    before: (Tree, Tree),
106    boundary: Option<WindowCursor>,
107    request: Option<WindowRequest>,
108    header: [u8; 5],
109    header_used: usize,
110    kind: u8,
111    payload_len: u64,
112    carry: Vec<u8>,
113    trailer: [u8; 32],
114    resume_prefix: Option<Hash>,
115    last_frame: Option<FrameInfo>,
116    #[cfg(test)]
117    peak: usize,
118}
119
120fn us(n: u64) -> Result<usize, PackError> {
121    usize::try_from(n).map_err(|_| PackError::PackfileTooLarge)
122}
123fn geometry(pack_len: u64, window: u64) -> Result<(), PackError> {
124    if pack_len < 44 {
125        return Err(PackError::PackfileTooShort);
126    }
127    if !(64 << 10..=64 << 20).contains(&window) || !window.is_power_of_two() {
128        return Err(PackError::PackfileCorrupted);
129    }
130    if pack_len > MAX_TOTAL_PAYLOAD + u64::from(MAX_ENTRIES) * 5 + 44 {
131        return Err(PackError::PackfileTooLarge);
132    }
133    Ok(())
134}
135fn copy_bytes(bytes: &[u8]) -> Result<Vec<u8>, PackError> {
136    let mut out = Vec::new();
137    out.try_reserve_exact(bytes.len())
138        .map_err(|_| PackError::PackfileTooLarge)?;
139    out.extend_from_slice(bytes);
140    Ok(out)
141}
142
143impl WindowReader {
144    /// Start a pack with power-of-two windows between 64 KiB and 64 MiB.
145    ///
146    /// # Errors
147    /// Short pack lengths give `PackfileTooShort`; invalid geometry gives
148    /// `PackfileCorrupted`.
149    pub fn new(
150        pack_len: u64,
151        window_size: u64,
152        limits: DecodeLimits,
153        expected_pack_id: Option<Hash>,
154    ) -> Result<Self, PackError> {
155        geometry(pack_len, window_size)?;
156        Ok(Self::from_state(
157            WindowCursor::initial(pack_len, window_size, expected_pack_id),
158            limits,
159            if expected_pack_id.is_some() {
160                Phase::Header
161            } else {
162                Phase::Anchor
163            },
164        ))
165    }
166
167    /// Restore an entry boundary. The first request contains the next entry.
168    ///
169    /// # Errors
170    /// Inconsistent cursors give `PackfileCorrupted`.
171    pub fn resume(cursor: &WindowCursor, limits: DecodeLimits) -> Result<Self, PackError> {
172        cursor.validate()?;
173        let mut reader = Self::from_state(cursor.clone(), limits, Phase::Boundary);
174        reader.boundary = Some(cursor.clone());
175        reader.resume_prefix = cursor.window_prefix;
176        Ok(reader)
177    }
178
179    fn from_state(state: WindowCursor, limits: DecodeLimits, phase: Phase) -> Self {
180        let start = state.completed * state.window_size; // validated cursor geometry
181        let before = (state.trailer_tree.clone(), state.id_tree.clone());
182        Self {
183            state,
184            limits,
185            phase,
186            window: Vec::new(),
187            start,
188            end: start,
189            before,
190            boundary: None,
191            request: None,
192            header: [0; 5],
193            header_used: 0,
194            kind: 0,
195            payload_len: 0,
196            carry: Vec::new(),
197            trailer: [0; 32],
198            resume_prefix: None,
199            last_frame: None,
200            #[cfg(test)]
201            peak: 0,
202        }
203    }
204
205    fn need(&mut self, offset: u64) -> Result<Step, PackError> {
206        let len = self
207            .state
208            .pack_len
209            .checked_sub(offset)
210            .ok_or(PackError::PackfileCorrupted)?
211            .min(self.state.window_size);
212        if len == 0 {
213            return Err(PackError::UnexpectedEof);
214        }
215        self.window = Vec::new();
216        let request = WindowRequest { offset, len };
217        self.request = Some(request);
218        Ok(Step::NeedWindow(request))
219    }
220
221    /// Advance to a request, an entry, or verified completion.
222    ///
223    /// # Errors
224    /// Malformed packs, resource limits, and failed integrity checks use the
225    /// errors documented at module level. A parsing error makes the reader inert.
226    pub fn step(&mut self) -> Result<Step, PackError> {
227        if let Some(request) = self.request {
228            return Ok(Step::NeedWindow(request));
229        }
230        let result = self.advance();
231        if result.is_err() {
232            self.phase = Phase::Failed;
233        }
234        result
235    }
236
237    fn advance(&mut self) -> Result<Step, PackError> {
238        loop {
239            match self.phase {
240                Phase::Failed => return Err(PackError::PackfileCorrupted),
241                Phase::Done => return Ok(Step::Done(self.summary())),
242                Phase::Anchor => {
243                    if self.end == self.state.pack_len {
244                        self.state.anchor = Some(self.trailer);
245                        self.trailer = [0; 32];
246                        self.start = 0;
247                        self.end = 0;
248                        self.phase = Phase::Header;
249                    } else {
250                        let offset = if self.end == 0 {
251                            (self.state.split() / self.state.window_size) * self.state.window_size
252                        } else {
253                            self.end
254                        };
255                        return self.need(offset);
256                    }
257                }
258                Phase::Header => {
259                    if let Some(step) = self.advance_header()? {
260                        return Ok(step);
261                    }
262                }
263                Phase::Boundary => {
264                    // Resume must fetch the containing window, even at EOF.
265                    if self.window.is_empty() {
266                        return self.need(self.start);
267                    }
268                    if self.state.index == self.state.count {
269                        self.phase = Phase::Finish;
270                        continue;
271                    }
272                    self.header_used = 0;
273                    self.phase = Phase::Frame;
274                }
275                Phase::Frame => {
276                    if let Some(step) = self.advance_frame()? {
277                        return Ok(step);
278                    }
279                }
280                Phase::Payload => {
281                    if let Some(step) = self.advance_payload()? {
282                        return Ok(step);
283                    }
284                }
285                Phase::Finish => {
286                    if self.state.pos != self.state.split() {
287                        return Err(PackError::TrailingData);
288                    }
289                    if self.end < self.state.pack_len {
290                        return self.need(self.end);
291                    }
292                    if self
293                        .state
294                        .anchor
295                        .is_some_and(|anchor| anchor != self.trailer)
296                        || self.state.trailer_tree.root != Some(self.trailer)
297                        || self
298                            .state
299                            .expected
300                            .is_some_and(|id| self.state.id_tree.root != Some(id))
301                    {
302                        return Err(PackError::PackfileCorrupted);
303                    }
304                    self.phase = Phase::Done;
305                }
306            }
307        }
308    }
309
310    fn advance_header(&mut self) -> Result<Option<Step>, PackError> {
311        if self.window.is_empty() {
312            return self.need(0).map(Some);
313        }
314        if &self.window[..4] != MAGIC {
315            return Err(PackError::InvalidMagic);
316        }
317        self.state.version = u32::from_le_bytes(
318            self.window[4..8]
319                .try_into()
320                .map_err(|_| PackError::UnexpectedEof)?,
321        );
322        if !matches!(self.state.version, 1 | 2) {
323            return Err(PackError::UnsupportedVersion(self.state.version));
324        }
325        self.state.count = u32::from_le_bytes(
326            self.window[8..12]
327                .try_into()
328                .map_err(|_| PackError::UnexpectedEof)?,
329        );
330        if self.state.count > MAX_ENTRIES {
331            return Err(PackError::TooManyObjects(self.state.count));
332        }
333        self.phase = Phase::Boundary;
334        self.boundary = Some(self.boundary_state());
335        Ok(None)
336    }
337
338    fn advance_frame(&mut self) -> Result<Option<Step>, PackError> {
339        if self
340            .state
341            .pos
342            .checked_add(
343                u64::try_from(5 - self.header_used).map_err(|_| PackError::PackfileTooLarge)?,
344            )
345            .is_none_or(|end| end > self.state.split())
346        {
347            return Err(PackError::UnexpectedEof);
348        }
349        if self.state.pos == self.end {
350            return self.need(self.end).map(Some);
351        }
352        let begin = us(self
353            .state
354            .pos
355            .checked_sub(self.start)
356            .ok_or(PackError::PackfileCorrupted)?)?;
357        let n = (5 - self.header_used).min(self.window.len() - begin);
358        self.header[self.header_used..self.header_used + n]
359            .copy_from_slice(&self.window[begin..begin + n]);
360        self.header_used += n;
361        self.state.pos = self
362            .state
363            .pos
364            .checked_add(u64::try_from(n).map_err(|_| PackError::PackfileTooLarge)?)
365            .ok_or(PackError::PackfileTooLarge)?;
366        if self.header_used != 5 {
367            return Ok(None);
368        }
369        self.kind = self.header[0];
370        self.payload_len = u64::from(u32::from_le_bytes(
371            self.header[1..]
372                .try_into()
373                .map_err(|_| PackError::UnexpectedEof)?,
374        ));
375        self.state.payload_sum = self
376            .state
377            .payload_sum
378            .checked_add(self.payload_len)
379            .ok_or(PackError::PackfileTooLarge)?;
380        if self
381            .limits
382            .entry_geometry
383            .is_some_and(|(frame, _)| self.payload_len > frame)
384            || self.state.payload_sum > MAX_TOTAL_PAYLOAD
385        {
386            return Err(PackError::PackfileTooLarge);
387        }
388        if self.payload_len
389            > self
390                .state
391                .split()
392                .checked_sub(self.state.pos)
393                .ok_or(PackError::UnexpectedEof)?
394        {
395            return Err(PackError::UnexpectedEof);
396        }
397        match self.kind {
398            0 => {}
399            2 | 4 if self.kind == 2 || self.state.version == 2 => {
400                if self.payload_len < 32 {
401                    return Err(PackError::DeltaEntryTruncated);
402                }
403            }
404            3 if self.state.version == 2 => {}
405            other => return Err(PackError::InvalidEntryType(other)),
406        }
407        if self.kind != 0 && self.state.first_non_raw.is_none() {
408            self.state.first_non_raw = Some(self.state.index);
409        }
410        self.phase = Phase::Payload;
411        Ok(None)
412    }
413
414    fn advance_payload(&mut self) -> Result<Option<Step>, PackError> {
415        let carried = u64::try_from(self.carry.len()).map_err(|_| PackError::PackfileTooLarge)?;
416        let remaining = self
417            .payload_len
418            .checked_sub(carried)
419            .ok_or(PackError::PackfileCorrupted)?;
420        if self.state.pos == self.end && remaining != 0 {
421            return self.need(self.end).map(Some);
422        }
423        let begin = us(self
424            .state
425            .pos
426            .checked_sub(self.start)
427            .ok_or(PackError::PackfileCorrupted)?)?;
428        let available =
429            u64::try_from(self.window.len() - begin).map_err(|_| PackError::PackfileTooLarge)?;
430        if self.carry.is_empty() && remaining <= available {
431            let end = begin
432                .checked_add(us(remaining)?)
433                .ok_or(PackError::PackfileTooLarge)?;
434            let payload = &self.window[begin..end];
435            self.check_budget(payload, 0)?;
436            let entry = own_entry(decode_payload(self.kind, self.state.version, payload)?)?;
437            self.record_peak(entry_len(&entry));
438            self.state.pos = self
439                .state
440                .pos
441                .checked_add(remaining)
442                .ok_or(PackError::PackfileTooLarge)?;
443            self.entry_finished()?;
444            return Ok(Some(Step::Entry(entry)));
445        }
446        if self.carry.is_empty() {
447            if self.payload_len
448                > self
449                    .limits
450                    .entry_geometry
451                    .map_or(self.limits.max_decoded_bytes, |(frame, _)| frame)
452            {
453                return Err(PackError::PackfileTooLarge);
454            }
455            self.carry
456                .try_reserve_exact(us(self.payload_len)?)
457                .map_err(|_| PackError::PackfileTooLarge)?;
458        }
459        let n = remaining.min(available);
460        let end = begin
461            .checked_add(us(n)?)
462            .ok_or(PackError::PackfileTooLarge)?;
463        self.carry.extend_from_slice(&self.window[begin..end]);
464        self.state.pos = self
465            .state
466            .pos
467            .checked_add(n)
468            .ok_or(PackError::PackfileTooLarge)?;
469        self.record_peak(0);
470        // Check a zstd claim as soon as its prefix has arrived.
471        let prefix = if self.kind == 4 { 36 } else { 4 };
472        if matches!(self.kind, 3 | 4) && self.carry.len() >= prefix {
473            self.check_budget(&self.carry, self.payload_len)?;
474        }
475        if u64::try_from(self.carry.len()).map_err(|_| PackError::PackfileTooLarge)?
476            == self.payload_len
477        {
478            self.check_budget(&self.carry, self.payload_len)?;
479            let released = self.release_carried_window()?;
480            let entry = if matches!(self.kind, 0 | 2) {
481                let mut bytes = std::mem::take(&mut self.carry);
482                if self.kind == 0 {
483                    PackEntry::Raw {
484                        bytes: Cow::Owned(bytes),
485                    }
486                } else {
487                    let base = bytes[..32]
488                        .try_into()
489                        .map_err(|_| PackError::DeltaEntryTruncated)?;
490                    bytes.drain(..32);
491                    PackEntry::Delta {
492                        base,
493                        stream: Cow::Owned(bytes),
494                    }
495                }
496            } else {
497                let entry = own_entry(decode_payload(self.kind, self.state.version, &self.carry)?)?;
498                self.record_peak(entry_len(&entry));
499                self.carry = Vec::new();
500                entry
501            };
502            if let Some(cursor) = released {
503                let frame = self.last_frame;
504                *self = Self::resume(&cursor, self.limits)?;
505                self.last_frame = frame;
506            } else {
507                self.entry_finished()?;
508            }
509            return Ok(Some(Step::Entry(entry)));
510        }
511        Ok(None)
512    }
513
514    // Authenticate the post-entry prefix before releasing the carried frame's window.
515    fn release_carried_window(&mut self) -> Result<Option<WindowCursor>, PackError> {
516        if self.limits.entry_geometry.is_none() || !matches!(self.kind, 3 | 4) {
517            return Ok(None);
518        }
519        self.entry_finished()?;
520        let cursor = self.checkpoint().ok_or(PackError::PackfileCorrupted)?;
521        self.window = Vec::new();
522        Ok(Some(cursor))
523    }
524
525    fn check_budget(&self, payload: &[u8], carried: u64) -> Result<(), PackError> {
526        if let Some((frame, stream)) = self.limits.entry_geometry {
527            let claim = match self.kind {
528                3 => zstd_claim(payload)?.0 as u64,
529                4 => zstd_claim(payload.get(32..).ok_or(PackError::DeltaEntryTruncated)?)?.0 as u64,
530                _ => self.payload_len,
531            };
532            let cap = match self.kind {
533                2 => stream.saturating_add(32),
534                4 => stream,
535                _ => self.limits.max_decoded_bytes,
536            };
537            return if self.payload_len <= frame && claim <= cap {
538                Ok(())
539            } else {
540                Err(PackError::PackfileTooLarge)
541            };
542        }
543        let charge = if matches!(self.kind, 3 | 4) {
544            let prefix = if self.kind == 4 { 32 } else { 0 };
545            let (claim, _) = zstd_claim(
546                payload
547                    .get(prefix..)
548                    .ok_or(PackError::DeltaEntryTruncated)?,
549            )?;
550            carried
551                .checked_add(u64::try_from(claim).map_err(|_| PackError::PackfileTooLarge)?)
552                .ok_or(PackError::PackfileTooLarge)?
553        } else {
554            self.payload_len
555        };
556        if charge > self.limits.max_decoded_bytes {
557            Err(PackError::PackfileTooLarge)
558        } else {
559            Ok(())
560        }
561    }
562
563    fn entry_finished(&mut self) -> Result<(), PackError> {
564        let length = self
565            .payload_len
566            .checked_add(5)
567            .ok_or(PackError::PackfileTooLarge)?;
568        self.last_frame = Some(FrameInfo {
569            offset: self
570                .state
571                .pos
572                .checked_sub(length)
573                .ok_or(PackError::PackfileCorrupted)?,
574            length,
575            wire_type: self.kind,
576        });
577        self.state.index = self
578            .state
579            .index
580            .checked_add(1)
581            .ok_or(PackError::PackfileTooLarge)?;
582        self.phase = Phase::Boundary;
583        self.boundary = Some(self.boundary_state());
584        Ok(())
585    }
586
587    /// Supply exactly the outstanding range. Wrong ranges leave it pending.
588    ///
589    /// # Errors
590    /// Wrong ranges give `PackfileCorrupted`; failed allocations give
591    /// `PackfileTooLarge`.
592    pub fn feed(&mut self, offset: u64, bytes: &[u8]) -> Result<(), PackError> {
593        self.validate_feed(offset, bytes)?;
594        let retain = self.retains_window();
595        let buffer = if retain {
596            copy_bytes(bytes)?
597        } else {
598            Vec::new()
599        };
600        self.feed_data(offset, bytes)?;
601        if retain {
602            self.window = buffer;
603            self.record_peak(0);
604        }
605        Ok(())
606    }
607
608    /// Supply a range by transferring its buffer, avoiding two full-window copies.
609    ///
610    /// # Errors
611    /// Wrong ranges give `PackfileCorrupted`; allocation failure gives `PackfileTooLarge`.
612    pub fn feed_owned(&mut self, offset: u64, bytes: Vec<u8>) -> Result<(), PackError> {
613        self.validate_feed(offset, &bytes)?;
614        let retain = self.retains_window();
615        let bytes = if retain && bytes.capacity() > us(self.state.window_size)? {
616            copy_bytes(&bytes)?
617        } else {
618            bytes
619        };
620        self.feed_data(offset, &bytes)?;
621        if retain {
622            self.window = bytes;
623            self.record_peak(0);
624        }
625        Ok(())
626    }
627
628    fn retains_window(&self) -> bool {
629        !(self.phase == Phase::Anchor && self.state.pack_len > self.state.window_size)
630    }
631
632    fn validate_feed(&self, offset: u64, bytes: &[u8]) -> Result<WindowRequest, PackError> {
633        let request = self.request.ok_or(PackError::PackfileCorrupted)?;
634        if request.offset != offset
635            || u64::try_from(bytes.len()).map_err(|_| PackError::PackfileTooLarge)? != request.len
636        {
637            return Err(PackError::PackfileCorrupted);
638        }
639        Ok(request)
640    }
641
642    fn feed_data(&mut self, offset: u64, bytes: &[u8]) -> Result<(), PackError> {
643        let request = self.validate_feed(offset, bytes)?;
644        if let Some(expected) = self.resume_prefix {
645            let prefix_len = self
646                .state
647                .pos
648                .checked_sub(offset)
649                .ok_or(PackError::PackfileCorrupted)?;
650            let prefix = bytes
651                .get(..us(prefix_len)?)
652                .ok_or(PackError::PackfileCorrupted)?;
653            if crate::hash::hash(prefix) != expected {
654                return Err(PackError::PackfileCorrupted);
655            }
656            self.resume_prefix = None;
657        }
658        if self.phase == Phase::Anchor && self.state.pack_len > self.state.window_size {
659            self.end = offset
660                .checked_add(request.len)
661                .ok_or(PackError::PackfileCorrupted)?;
662            self.collect_trailer(offset, bytes)?;
663        } else {
664            self.before = (self.state.trailer_tree.clone(), self.state.id_tree.clone());
665            self.state.trailer_tree.absorb(
666                offset,
667                bytes,
668                self.state.split(),
669                self.state.window_size,
670            )?;
671            if self.state.expected.is_some() {
672                self.state.id_tree.absorb(
673                    offset,
674                    bytes,
675                    self.state.pack_len,
676                    self.state.window_size,
677                )?;
678            }
679            self.start = offset;
680            self.end = offset
681                .checked_add(request.len)
682                .ok_or(PackError::PackfileCorrupted)?;
683            self.collect_trailer(offset, bytes)?;
684            if self.phase == Phase::Anchor {
685                // A single-window pack already supplies the trailer and all
686                // stream bytes. Reuse its buffer and hashes for header parsing.
687                self.state.anchor = Some(self.trailer);
688                self.phase = Phase::Header;
689            }
690        }
691        self.request = None;
692        Ok(())
693    }
694
695    fn collect_trailer(&mut self, offset: u64, bytes: &[u8]) -> Result<(), PackError> {
696        let from = self.state.split().max(offset);
697        if from < self.end {
698            let src = us(from
699                .checked_sub(offset)
700                .ok_or(PackError::PackfileCorrupted)?)?;
701            let dst = us(from
702                .checked_sub(self.state.split())
703                .ok_or(PackError::PackfileCorrupted)?)?;
704            let n = us(self
705                .end
706                .checked_sub(from)
707                .ok_or(PackError::PackfileCorrupted)?)?;
708            self.trailer[dst..dst + n].copy_from_slice(&bytes[src..src + n]);
709        }
710        Ok(())
711    }
712
713    /// The frame of the entry the latest [`Step::Entry`] returned, or `None`
714    /// before the first one. A resumed reader starts with `None`.
715    #[must_use]
716    pub fn last_frame(&self) -> Option<FrameInfo> {
717        self.last_frame
718    }
719
720    /// A compact checkpoint only at an entry boundary, never within an entry.
721    ///
722    /// Hashes the current window prefix lazily on this call. Returns `None` if
723    /// the reader has released the bytes needed for that prefix, including while
724    /// fetching a split trailer or awaiting a resumed window. Keep the previous
725    /// cursor in that case. No bytes are needed when the boundary is exactly at
726    /// a window start. Mid-entry states, including carried straddling entries,
727    /// never produce a checkpoint.
728    #[must_use]
729    pub fn checkpoint(&self) -> Option<WindowCursor> {
730        if !matches!(self.phase, Phase::Boundary | Phase::Finish | Phase::Done) {
731            return None;
732        }
733        let mut cursor = self.boundary.clone()?;
734        if self.window.is_empty()
735            && self.limits.entry_geometry.is_some()
736            && self.resume_prefix.is_some()
737        {
738            return Some(cursor);
739        }
740        let start = cursor.completed.checked_mul(cursor.window_size)?;
741        let prefix_len = cursor.pos.checked_sub(start)?;
742        cursor.window_prefix = if prefix_len == 0 {
743            None
744        } else {
745            if self.start != start {
746                return None;
747            }
748            let end = usize::try_from(prefix_len).ok()?;
749            Some(crate::hash::hash(self.window.get(..end)?))
750        };
751        Some(cursor)
752    }
753
754    fn boundary_state(&self) -> WindowCursor {
755        let mut state = self.state.clone();
756        state.completed = state.pos / state.window_size;
757        if state.completed == self.start / state.window_size {
758            state.trailer_tree = self.before.0.clone();
759            state.id_tree = self.before.1.clone();
760        }
761        state
762    }
763
764    fn summary(&self) -> WindowSummary {
765        WindowSummary {
766            version: self.state.version,
767            entry_count: self.state.count,
768            raw_only: self.state.first_non_raw.is_none(),
769            first_non_raw: self.state.first_non_raw,
770        }
771    }
772
773    #[cfg_attr(not(test), allow(clippy::unused_self))] // counter exists only in tests
774    fn record_peak(&mut self, extra: usize) {
775        #[cfg(test)]
776        {
777            self.peak = self
778                .peak
779                .max(self.window.capacity() + self.carry.capacity() + extra);
780        }
781        #[cfg(not(test))]
782        let _ = extra;
783    }
784}
785
786fn own_entry(entry: PackEntry<'_>) -> Result<PackEntry<'static>, PackError> {
787    fn own(bytes: Cow<'_, [u8]>) -> Result<Cow<'static, [u8]>, PackError> {
788        Ok(Cow::Owned(match bytes {
789            Cow::Owned(bytes) => bytes,
790            Cow::Borrowed(bytes) => copy_bytes(bytes)?,
791        }))
792    }
793    Ok(match entry {
794        PackEntry::Raw { bytes } => PackEntry::Raw { bytes: own(bytes)? },
795        PackEntry::Delta { base, stream } => PackEntry::Delta {
796            base,
797            stream: own(stream)?,
798        },
799    })
800}
801fn entry_len(entry: &PackEntry<'_>) -> usize {
802    let bytes = match entry {
803        PackEntry::Raw { bytes } => bytes,
804        PackEntry::Delta { stream, .. } => stream,
805    };
806    match bytes {
807        Cow::Borrowed(bytes) => bytes.len(),
808        Cow::Owned(bytes) => bytes.capacity(),
809    }
810}
811
812/// A synchronous source serving exact ranges, without retaining earlier windows.
813pub trait WindowSource {
814    /// Read exactly `len` bytes beginning at `offset`.
815    ///
816    /// # Errors
817    /// Return a pack error for a short read or source failure.
818    fn read_window(&mut self, offset: u64, len: u64) -> Result<Vec<u8>, PackError>;
819}
820impl WindowSource for &[u8] {
821    fn read_window(&mut self, offset: u64, len: u64) -> Result<Vec<u8>, PackError> {
822        let end = offset.checked_add(len).ok_or(PackError::UnexpectedEof)?;
823        copy_bytes(
824            self.get(us(offset)?..us(end)?)
825                .ok_or(PackError::UnexpectedEof)?,
826        )
827    }
828}
829
830/// Drive a reader synchronously, dropping each window after feeding it.
831///
832/// # Errors
833/// Propagates reader, source, or sink errors. All previously delivered entries
834/// must be discarded on error, including a final integrity failure.
835pub fn read_all<S: WindowSource>(
836    source: &mut S,
837    pack_len: u64,
838    window_size: u64,
839    limits: DecodeLimits,
840    expected_pack_id: Option<Hash>,
841    mut sink: impl FnMut(PackEntry<'static>) -> Result<(), PackError>,
842) -> Result<WindowSummary, PackError> {
843    let mut reader = WindowReader::new(pack_len, window_size, limits, expected_pack_id)?;
844    loop {
845        match reader.step()? {
846            Step::NeedWindow(request) => {
847                let bytes = source.read_window(request.offset, request.len)?;
848                reader.feed_owned(request.offset, bytes)?;
849            }
850            Step::Entry(entry) => sink(entry)?,
851            Step::Done(summary) => return Ok(summary),
852        }
853    }
854}
855
856#[cfg(test)]
857mod tests;