Skip to main content

heddle_object_model/object/
tree_stream.rs

1// SPDX-License-Identifier: Apache-2.0
2//! Streaming Tree entry reader with persistable resume cursors.
3
4use serde::{Deserialize, Serialize};
5use sley::{ObjectFormat as GitObjectFormat, ObjectId as GitObjectId};
6
7use super::{
8    ContentHash, EntryType, FileMode, SpoolId, StateId, Tree, TreeEntry, TreeError,
9    tree::git_format_from_tag,
10    tree_canonical::{
11        TREE_BLOCK_ENCODING_VERSION, TREE_BLOCK_INDEX_LEN, TREE_BLOCK_PREAMBLE_LEN,
12        TREE_ENCODING_VERSION, TREE_HEADER_LEN, TREE_LEAN_ENCODING_VERSION, TREE_LEAN_MAGIC,
13        TreeBlockHeader, TreeBlockIndex, TreeHeader, decode_block_header, decode_block_index,
14        decode_block_payload, decode_entry_frame, decode_header,
15    },
16    tree_source::{TreeBodyIntegrity, TreeByteSource},
17};
18
19/// Failure while streaming or range-resuming a canonical tree.
20#[derive(Debug, thiserror::Error)]
21pub enum TreeStreamError {
22    #[error("invalid tree entry: {0}")]
23    Invalid(#[from] TreeError),
24    #[error(
25        "unsupported tree encoding version {found} (this binary supports {TREE_ENCODING_VERSION}, {TREE_BLOCK_ENCODING_VERSION}, and {TREE_LEAN_ENCODING_VERSION})"
26    )]
27    UnsupportedVersion { found: u8 },
28    #[error("tree resume cursor does not match this object: {0}")]
29    CursorMismatch(String),
30    #[error("truncated tree frame at byte {offset}")]
31    TruncatedFrame { offset: u64 },
32    #[error("tree payload has {extra} trailing byte(s) after declared end")]
33    TrailingBytes { extra: u64 },
34    #[error("tree ended after {decoded} of {expected} declared entries")]
35    UnexpectedEof { expected: u64, decoded: u64 },
36    #[error("tree entry exceeds page byte limit ({decoded_bytes} > {max_decoded_bytes})")]
37    OversizedEntry {
38        decoded_bytes: usize,
39        max_decoded_bytes: usize,
40    },
41    #[error("tree page limits must be nonzero")]
42    InvalidPageLimits,
43    #[error("ranged tree resume requires a verified-placement object source")]
44    UnverifiedRange,
45    #[error("malformed tree encoding: {0}")]
46    Malformed(String),
47    #[error("tree compression failed: {0}")]
48    Compression(String),
49    #[error("tree I/O error: {0}")]
50    Io(#[from] std::io::Error),
51    #[error("decoded tree hash {found} does not match {expected}")]
52    HashMismatch {
53        expected: ContentHash,
54        found: ContentHash,
55    },
56}
57
58/// Caller-sized page budget. Zero limits fail closed.
59///
60/// Fields stay private so callers cannot bypass [`Self::new`].
61#[derive(Clone, Copy, Debug, PartialEq, Eq)]
62pub struct TreePageLimits {
63    max_entries: usize,
64    max_decoded_bytes: usize,
65}
66
67impl TreePageLimits {
68    pub fn new(max_entries: usize, max_decoded_bytes: usize) -> Result<Self, TreeStreamError> {
69        if max_entries == 0 || max_decoded_bytes == 0 {
70            return Err(TreeStreamError::InvalidPageLimits);
71        }
72        Ok(Self {
73            max_entries,
74            max_decoded_bytes,
75        })
76    }
77
78    pub fn max_entries(&self) -> usize {
79        self.max_entries
80    }
81
82    pub fn max_decoded_bytes(&self) -> usize {
83        self.max_decoded_bytes
84    }
85}
86
87/// Persistable entry-boundary cursor bound to a tree id and encoding version.
88#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
89pub struct TreeResumeCursor {
90    pub(crate) tree_id: ContentHash,
91    pub(crate) encoding_version: u8,
92    pub(crate) ordinal: u64,
93    pub(crate) byte_offset: u64,
94    pub(crate) prev_name: Option<String>,
95}
96
97impl TreeResumeCursor {
98    pub fn start(tree_id: ContentHash) -> Self {
99        Self::start_for_version(tree_id, TREE_ENCODING_VERSION)
100    }
101
102    fn start_for_version(tree_id: ContentHash, encoding_version: u8) -> Self {
103        Self {
104            tree_id,
105            encoding_version,
106            ordinal: 0,
107            byte_offset: TREE_HEADER_LEN as u64,
108            prev_name: None,
109        }
110    }
111
112    pub fn tree_id(&self) -> ContentHash {
113        self.tree_id
114    }
115
116    pub fn encoding_version(&self) -> u8 {
117        self.encoding_version
118    }
119
120    pub fn ordinal(&self) -> u64 {
121        self.ordinal
122    }
123
124    pub fn byte_offset(&self) -> u64 {
125        self.byte_offset
126    }
127
128    pub fn prev_name(&self) -> Option<&str> {
129        self.prev_name.as_deref()
130    }
131}
132
133#[derive(Debug)]
134enum TreeReaderLayout {
135    Raw,
136    Lean {
137        body_start: u64,
138    },
139    Blocked {
140        header: TreeBlockHeader,
141        cache: Option<DecodedBlock>,
142        validated_blocks: Vec<bool>,
143    },
144}
145
146#[derive(Debug)]
147struct DecodedBlock {
148    block: usize,
149    raw: Vec<u8>,
150    frames: Vec<(usize, usize)>,
151}
152
153/// One bounded page of decoded entries plus the cursor after the last entry.
154#[derive(Clone, Debug, PartialEq, Eq)]
155pub struct TreePage {
156    pub entries: Vec<TreeEntry>,
157    pub resume_cursor: TreeResumeCursor,
158}
159
160/// Incremental HTR4 reader. Does not materialize a full `Vec<TreeEntry>`.
161#[derive(Debug)]
162pub struct TreeEntryReader<S: TreeByteSource> {
163    source: S,
164    header: TreeHeader,
165    layout: TreeReaderLayout,
166    cursor: TreeResumeCursor,
167    hasher: Option<blake3::Hasher>,
168    lean_hash_entries: Option<Vec<TreeEntry>>,
169    decoded_logical_len: u64,
170    started_at_zero: bool,
171    pending: Option<(TreeEntry, usize)>,
172    finished: bool,
173}
174
175impl<S: TreeByteSource> TreeEntryReader<S> {
176    pub fn open(
177        mut source: S,
178        expected_id: ContentHash,
179        resume: Option<&TreeResumeCursor>,
180    ) -> Result<Self, TreeStreamError> {
181        let mut magic = [0u8; 4];
182        source.read_exact_at(0, &mut magic)?;
183        let (header, layout) = if &magic == TREE_LEAN_MAGIC {
184            let (entry_count, body_start) = read_source_varint(&mut source, 4)?;
185            let entry_count = u64::try_from(entry_count)
186                .map_err(|_| TreeStreamError::Malformed("tree entry count exceeds u64".into()))?;
187            let payload_len = source
188                .len()
189                .checked_sub(body_start)
190                .ok_or(TreeStreamError::TruncatedFrame { offset: body_start })?;
191            (
192                TreeHeader {
193                    version: TREE_LEAN_ENCODING_VERSION,
194                    tree_id: expected_id,
195                    entry_count,
196                    payload_len,
197                    logical_len: 0,
198                },
199                TreeReaderLayout::Lean { body_start },
200            )
201        } else {
202            let mut header_buf = [0u8; TREE_HEADER_LEN];
203            source.read_exact_at(0, &mut header_buf)?;
204            let header = decode_header(&header_buf)?;
205            if header.tree_id != expected_id {
206                return Err(TreeStreamError::HashMismatch {
207                    expected: expected_id,
208                    found: header.tree_id,
209                });
210            }
211            let expected_len = TREE_HEADER_LEN as u64 + header.payload_len;
212            if source.len() < expected_len {
213                return Err(TreeStreamError::TruncatedFrame {
214                    offset: source.len(),
215                });
216            }
217            if source.len() > expected_len {
218                return Err(TreeStreamError::TrailingBytes {
219                    extra: source.len() - expected_len,
220                });
221            }
222            let layout = if header.version == TREE_BLOCK_ENCODING_VERSION {
223                let mut preamble = [0u8; TREE_BLOCK_PREAMBLE_LEN];
224                source.read_exact_at(TREE_HEADER_LEN as u64, &mut preamble)?;
225                let block_header = decode_block_header(&header, &preamble, source.len())?;
226                TreeReaderLayout::Blocked {
227                    header: block_header,
228                    cache: None,
229                    validated_blocks: vec![false; block_header.block_count],
230                }
231            } else {
232                TreeReaderLayout::Raw
233            };
234            (header, layout)
235        };
236        let mut cursor = resume
237            .cloned()
238            .unwrap_or_else(|| TreeResumeCursor::start_for_version(expected_id, header.version));
239        if resume.is_none()
240            && let TreeReaderLayout::Lean { body_start } = &layout
241        {
242            cursor.byte_offset = *body_start;
243        }
244        validate_cursor(&header, &layout, &cursor)?;
245        if cursor.ordinal > 0 && source.integrity() != TreeBodyIntegrity::VerifiedPlacement {
246            return Err(TreeStreamError::UnverifiedRange);
247        }
248        let is_lean = matches!(layout, TreeReaderLayout::Lean { .. });
249        let hasher = (cursor.ordinal == 0 && !is_lean)
250            .then(|| ContentHash::typed_hasher("tree", header.logical_len));
251        let lean_hash_entries = (cursor.ordinal == 0 && is_lean).then(Vec::new);
252        let started_at_zero = cursor.ordinal == 0;
253        let mut reader = Self {
254            source,
255            header,
256            layout,
257            cursor,
258            hasher,
259            lean_hash_entries,
260            decoded_logical_len: 0,
261            started_at_zero,
262            pending: None,
263            finished: false,
264        };
265        reader.arm_pending_at_cursor()?;
266        Ok(reader)
267    }
268
269    pub fn header(&self) -> &TreeHeader {
270        &self.header
271    }
272
273    pub fn bytes_read(&self) -> u64 {
274        self.source.bytes_read()
275    }
276
277    pub fn next_page(
278        &mut self,
279        limits: TreePageLimits,
280    ) -> Result<Option<TreePage>, TreeStreamError> {
281        if limits.max_entries() == 0 || limits.max_decoded_bytes() == 0 {
282            return Err(TreeStreamError::InvalidPageLimits);
283        }
284        if self.cursor.ordinal == self.header.entry_count {
285            return Ok(None);
286        }
287        let mut entries = Vec::new();
288        let mut decoded_bytes = 0usize;
289        while entries.len() < limits.max_entries() && self.cursor.ordinal < self.header.entry_count
290        {
291            let (entry, consumed) = self.take_next_entry()?;
292            let size = entry.decoded_size();
293            if size > limits.max_decoded_bytes() {
294                return Err(TreeStreamError::OversizedEntry {
295                    decoded_bytes: size,
296                    max_decoded_bytes: limits.max_decoded_bytes(),
297                });
298            }
299            if !entries.is_empty()
300                && decoded_bytes.saturating_add(size) > limits.max_decoded_bytes()
301            {
302                self.pending = Some((entry, consumed));
303                break;
304            }
305            self.commit_entry(&entry, consumed)?;
306            decoded_bytes += size;
307            entries.push(entry);
308        }
309        Ok(Some(TreePage {
310            entries,
311            resume_cursor: self.cursor.clone(),
312        }))
313    }
314
315    /// Yield one decoded entry. Used by full-object collect without a page ceiling.
316    pub fn next_entry(&mut self) -> Result<Option<TreeEntry>, TreeStreamError> {
317        if self.cursor.ordinal == self.header.entry_count {
318            return Ok(None);
319        }
320        let (entry, consumed) = self.take_next_entry()?;
321        self.commit_entry(&entry, consumed)?;
322        Ok(Some(entry))
323    }
324
325    pub fn finish_and_verify(&mut self) -> Result<(), TreeStreamError> {
326        if self.cursor.ordinal != self.header.entry_count {
327            return Err(TreeStreamError::UnexpectedEof {
328                expected: self.header.entry_count,
329                decoded: self.cursor.ordinal,
330            });
331        }
332        let payload_end = self.logical_payload_end();
333        if self.cursor.byte_offset != payload_end {
334            return Err(TreeStreamError::TrailingBytes {
335                extra: payload_end.abs_diff(self.cursor.byte_offset),
336            });
337        }
338        if let Some(hasher) = self.hasher.take() {
339            if self.decoded_logical_len != self.header.logical_len {
340                return Err(TreeStreamError::Malformed(
341                    "declared logical length does not match entries".into(),
342                ));
343            }
344            let found = ContentHash::from_bytes(hasher.finalize().into());
345            if found != self.header.tree_id {
346                return Err(TreeStreamError::HashMismatch {
347                    expected: self.header.tree_id,
348                    found,
349                });
350            }
351        } else if let Some(entries) = self.lean_hash_entries.take() {
352            let tree = Tree::try_from_decoded_entries(entries)?;
353            let found = tree.hash();
354            if found != self.header.tree_id {
355                return Err(TreeStreamError::HashMismatch {
356                    expected: self.header.tree_id,
357                    found,
358                });
359            }
360        } else if self.source.integrity() != TreeBodyIntegrity::VerifiedPlacement {
361            return Err(TreeStreamError::UnverifiedRange);
362        }
363        self.validate_complete_block_layout()?;
364        self.finished = true;
365        Ok(())
366    }
367
368    fn take_next_entry(&mut self) -> Result<(TreeEntry, usize), TreeStreamError> {
369        if let Some(pending) = self.pending.take() {
370            return Ok(pending);
371        }
372        self.read_entry_at(self.cursor.byte_offset)
373    }
374
375    fn read_entry_at(&mut self, offset: u64) -> Result<(TreeEntry, usize), TreeStreamError> {
376        if matches!(self.layout, TreeReaderLayout::Lean { .. }) {
377            return self.read_lean_entry(offset);
378        }
379        if matches!(self.layout, TreeReaderLayout::Blocked { .. }) {
380            return self.read_blocked_entry();
381        }
382        let payload_end = TREE_HEADER_LEN as u64 + self.header.payload_len;
383        let mut len_buf = [0u8; 4];
384        self.source.read_exact_at(offset, &mut len_buf)?;
385        let frame_len = u64::from(u32::from_le_bytes(len_buf));
386        let frame_start = offset
387            .checked_add(4)
388            .ok_or(TreeStreamError::TruncatedFrame { offset })?;
389        let frame_end = frame_start
390            .checked_add(frame_len)
391            .ok_or(TreeStreamError::TruncatedFrame { offset })?;
392        if frame_end > payload_end || frame_end > self.source.len() {
393            return Err(TreeStreamError::TruncatedFrame { offset });
394        }
395        let frame_len =
396            usize::try_from(frame_len).map_err(|_| TreeStreamError::TruncatedFrame { offset })?;
397        let mut frame = vec![0u8; frame_len];
398        self.source.read_exact_at(frame_start, &mut frame)?;
399        let entry = decode_entry_frame(&frame)?;
400        Ok((entry, 4 + frame_len))
401    }
402
403    fn logical_payload_end(&self) -> u64 {
404        match &self.layout {
405            TreeReaderLayout::Raw => TREE_HEADER_LEN as u64 + self.header.payload_len,
406            TreeReaderLayout::Lean { .. } => self.source.len(),
407            TreeReaderLayout::Blocked { header, .. } => {
408                TREE_HEADER_LEN as u64 + header.raw_payload_len
409            }
410        }
411    }
412
413    fn read_lean_entry(&mut self, offset: u64) -> Result<(TreeEntry, usize), TreeStreamError> {
414        let mut cursor = offset;
415        let tag = read_source_byte(&mut self.source, &mut cursor)?;
416        let mode = FileMode::from_byte(tag >> 3).ok_or_else(|| {
417            TreeStreamError::Malformed(format!("invalid compact entry mode {}", tag >> 3))
418        })?;
419        let kind = EntryType::from_byte(tag & 0x07).ok_or_else(|| {
420            TreeStreamError::Malformed(format!("invalid compact entry kind {}", tag & 0x07))
421        })?;
422        let prefix = read_source_varint_at(&mut self.source, &mut cursor)?;
423        let suffix_len = read_source_varint_at(&mut self.source, &mut cursor)?;
424        let previous = self.cursor.prev_name.as_deref().unwrap_or("");
425        if prefix > previous.len() {
426            return Err(TreeStreamError::Malformed(
427                "compact name prefix exceeds predecessor".into(),
428            ));
429        }
430        let mut suffix = vec![0u8; suffix_len];
431        self.source.read_exact_at(cursor, &mut suffix)?;
432        cursor = cursor
433            .checked_add(suffix_len as u64)
434            .ok_or_else(|| TreeStreamError::Malformed("compact name offset overflow".into()))?;
435        let mut name = previous.as_bytes()[..prefix].to_vec();
436        name.extend_from_slice(&suffix);
437        let name = String::from_utf8(name)
438            .map_err(|_| TreeStreamError::Malformed("compact entry name is not UTF-8".into()))?;
439        let entry = match kind {
440            EntryType::Blob | EntryType::Tree | EntryType::Symlink => {
441                let mut bytes = [0u8; 32];
442                self.source.read_exact_at(cursor, &mut bytes)?;
443                cursor = cursor.checked_add(32).ok_or_else(|| {
444                    TreeStreamError::Malformed("compact hash offset overflow".into())
445                })?;
446                let hash = ContentHash::from_bytes(bytes);
447                match kind {
448                    EntryType::Blob => TreeEntry::file(name, hash, mode == FileMode::Executable)?,
449                    EntryType::Tree => TreeEntry::directory(name, hash)?,
450                    EntryType::Symlink => TreeEntry::symlink(name, hash)?,
451                    EntryType::Gitlink | EntryType::Spoollink => {
452                        return Err(TreeStreamError::Malformed(
453                            "invalid compact content-addressed kind".into(),
454                        ));
455                    }
456                }
457            }
458            EntryType::Gitlink => {
459                let format = git_format_from_tag(read_source_byte(&mut self.source, &mut cursor)?)?;
460                let oid_len = match format {
461                    GitObjectFormat::Sha1 => 20,
462                    GitObjectFormat::Sha256 => 32,
463                };
464                let mut oid = vec![0u8; oid_len];
465                self.source.read_exact_at(cursor, &mut oid)?;
466                cursor = cursor.checked_add(oid_len as u64).ok_or_else(|| {
467                    TreeStreamError::Malformed("compact gitlink offset overflow".into())
468                })?;
469                let target = GitObjectId::from_raw(format, &oid).map_err(|error| {
470                    TreeStreamError::Malformed(format!("invalid compact gitlink: {error}"))
471                })?;
472                TreeEntry::gitlink(name, target)?
473            }
474            EntryType::Spoollink => {
475                let spool_len = read_source_varint_at(&mut self.source, &mut cursor)?;
476                let mut spool = vec![0u8; spool_len];
477                self.source.read_exact_at(cursor, &mut spool)?;
478                cursor = cursor.checked_add(spool_len as u64).ok_or_else(|| {
479                    TreeStreamError::Malformed("compact spool offset overflow".into())
480                })?;
481                let spool = std::str::from_utf8(&spool).map_err(|_| {
482                    TreeStreamError::Malformed("compact spool id is not UTF-8".into())
483                })?;
484                let spool_id = SpoolId::parse(spool).map_err(|error| {
485                    TreeStreamError::Malformed(format!("invalid compact spool id: {error}"))
486                })?;
487                let mut state = [0u8; 32];
488                self.source.read_exact_at(cursor, &mut state)?;
489                cursor = cursor.checked_add(32).ok_or_else(|| {
490                    TreeStreamError::Malformed("compact state offset overflow".into())
491                })?;
492                TreeEntry::spoollink(name, spool_id, StateId::from_bytes(state))?
493            }
494        };
495        if entry.mode() != mode {
496            return Err(TreeStreamError::Malformed(format!(
497                "compact entry kind/mode mismatch for {}: {kind:?}/{mode:?}",
498                entry.name()
499            )));
500        }
501        let consumed = usize::try_from(cursor - offset)
502            .map_err(|_| TreeStreamError::Malformed("compact entry length exceeds usize".into()))?;
503        Ok((entry, consumed))
504    }
505
506    fn read_blocked_entry(&mut self) -> Result<(TreeEntry, usize), TreeStreamError> {
507        let block_header = match &self.layout {
508            TreeReaderLayout::Blocked { header, .. } => *header,
509            TreeReaderLayout::Raw | TreeReaderLayout::Lean { .. } => {
510                return Err(TreeStreamError::Malformed(
511                    "raw tree entered blocked reader".into(),
512                ));
513            }
514        };
515        let ordinal = usize::try_from(self.cursor.ordinal)
516            .map_err(|_| TreeStreamError::Malformed("tree ordinal exceeds usize".into()))?;
517        let block = ordinal / block_header.block_entries;
518        let within_block = ordinal % block_header.block_entries;
519        let index = self.read_block_index(block, &block_header)?;
520        let block_logical_start = (TREE_HEADER_LEN as u64)
521            .checked_add(index.raw_offset)
522            .ok_or_else(|| TreeStreamError::Malformed("tree block raw offset overflow".into()))?;
523
524        if within_block == 0 {
525            if self.cursor.byte_offset != block_logical_start {
526                return Err(TreeStreamError::CursorMismatch(
527                    "cursor byte offset is not the block restart boundary".into(),
528                ));
529            }
530            return self.read_block_anchor(index);
531        }
532
533        let needs_load = match &self.layout {
534            TreeReaderLayout::Blocked { cache, .. } => {
535                cache.as_ref().is_none_or(|cached| cached.block != block)
536            }
537            TreeReaderLayout::Raw | TreeReaderLayout::Lean { .. } => true,
538        };
539        if needs_load {
540            let decoded = self.load_block(index, block, &block_header)?;
541            if let TreeReaderLayout::Blocked {
542                cache,
543                validated_blocks,
544                ..
545            } = &mut self.layout
546            {
547                *cache = Some(decoded);
548                if let Some(validated) = validated_blocks.get_mut(block) {
549                    *validated = true;
550                }
551            }
552        }
553        let cached = match &self.layout {
554            TreeReaderLayout::Blocked {
555                cache: Some(cached),
556                ..
557            } => cached,
558            _ => {
559                return Err(TreeStreamError::Malformed(
560                    "tree block cache was not populated".into(),
561                ));
562            }
563        };
564        let (frame_start, frame_end) =
565            cached
566                .frames
567                .get(within_block)
568                .copied()
569                .ok_or(TreeStreamError::UnexpectedEof {
570                    expected: self.cursor.ordinal + 1,
571                    decoded: self.cursor.ordinal,
572                })?;
573        let frame =
574            cached
575                .raw
576                .get(frame_start + 4..frame_end)
577                .ok_or(TreeStreamError::TruncatedFrame {
578                    offset: frame_start as u64,
579                })?;
580        let consumed = frame_end - frame_start;
581        let expected_offset = block_logical_start
582            .checked_add(frame_start as u64)
583            .ok_or_else(|| TreeStreamError::Malformed("tree cursor offset overflow".into()))?;
584        if self.cursor.byte_offset != expected_offset {
585            return Err(TreeStreamError::CursorMismatch(
586                "cursor byte offset is not an entry boundary".into(),
587            ));
588        }
589        Ok((decode_entry_frame(frame)?, consumed))
590    }
591
592    fn read_block_index(
593        &mut self,
594        block: usize,
595        block_header: &TreeBlockHeader,
596    ) -> Result<TreeBlockIndex, TreeStreamError> {
597        let relative = block
598            .checked_mul(TREE_BLOCK_INDEX_LEN)
599            .ok_or_else(|| TreeStreamError::Malformed("tree block index overflow".into()))?;
600        let offset = TREE_HEADER_LEN
601            .checked_add(TREE_BLOCK_PREAMBLE_LEN)
602            .and_then(|start| start.checked_add(relative))
603            .ok_or_else(|| TreeStreamError::Malformed("tree block index offset overflow".into()))?;
604        let mut bytes = [0u8; TREE_BLOCK_INDEX_LEN];
605        self.source.read_exact_at(offset as u64, &mut bytes)?;
606        decode_block_index(&bytes, block, block_header, self.source.len())
607    }
608
609    fn read_block_anchor(
610        &mut self,
611        index: TreeBlockIndex,
612    ) -> Result<(TreeEntry, usize), TreeStreamError> {
613        let mut len_bytes = [0u8; 4];
614        self.source
615            .read_exact_at(index.stored_offset, &mut len_bytes)?;
616        let frame_len = u32::from_le_bytes(len_bytes) as usize;
617        let consumed = 4usize
618            .checked_add(frame_len)
619            .ok_or(TreeStreamError::TruncatedFrame {
620                offset: index.stored_offset,
621            })?;
622        if consumed > index.stored_len || consumed > index.raw_len {
623            return Err(TreeStreamError::TruncatedFrame {
624                offset: index.stored_offset,
625            });
626        }
627        let mut frame = vec![0u8; frame_len];
628        self.source
629            .read_exact_at(index.stored_offset + 4, &mut frame)?;
630        Ok((decode_entry_frame(&frame)?, consumed))
631    }
632
633    fn load_block(
634        &mut self,
635        index: TreeBlockIndex,
636        block: usize,
637        block_header: &TreeBlockHeader,
638    ) -> Result<DecodedBlock, TreeStreamError> {
639        let mut stored = vec![0u8; index.stored_len];
640        self.source
641            .read_exact_at(index.stored_offset, &mut stored)?;
642        let raw = decode_block_payload(&stored, index.raw_len)?;
643        let first_entry = block
644            .checked_mul(block_header.block_entries)
645            .ok_or_else(|| TreeStreamError::Malformed("tree block ordinal overflow".into()))?;
646        let remaining = usize::try_from(self.header.entry_count)
647            .map_err(|_| TreeStreamError::Malformed("tree entry count exceeds usize".into()))?
648            .checked_sub(first_entry)
649            .ok_or_else(|| TreeStreamError::Malformed("tree block starts past entries".into()))?;
650        let expected_entries = remaining.min(block_header.block_entries);
651        let frames = decode_block_frames(&raw, expected_entries)?;
652        Ok(DecodedBlock { block, raw, frames })
653    }
654
655    fn validate_complete_block_layout(&mut self) -> Result<(), TreeStreamError> {
656        let block_header = match &self.layout {
657            TreeReaderLayout::Raw | TreeReaderLayout::Lean { .. } => return Ok(()),
658            TreeReaderLayout::Blocked { header, .. } => *header,
659        };
660        let mut expected_raw_offset = 0u64;
661        let mut expected_stored_offset = block_header.index_end;
662        for block in 0..block_header.block_count {
663            let index = self.read_block_index(block, &block_header)?;
664            if index.raw_offset != expected_raw_offset
665                || index.stored_offset != expected_stored_offset
666            {
667                return Err(TreeStreamError::Malformed(
668                    "tree blocks are not contiguous".into(),
669                ));
670            }
671            expected_raw_offset = expected_raw_offset
672                .checked_add(index.raw_len as u64)
673                .ok_or_else(|| TreeStreamError::Malformed("raw tree length overflow".into()))?;
674            expected_stored_offset = expected_stored_offset
675                .checked_add(index.stored_len as u64)
676                .ok_or_else(|| TreeStreamError::Malformed("stored tree length overflow".into()))?;
677            let needs_payload_validation = self.started_at_zero
678                && match &self.layout {
679                    TreeReaderLayout::Blocked {
680                        validated_blocks, ..
681                    } => validated_blocks
682                        .get(block)
683                        .is_none_or(|validated| !validated),
684                    TreeReaderLayout::Raw | TreeReaderLayout::Lean { .. } => false,
685                };
686            if needs_payload_validation {
687                let _ = self.load_block(index, block, &block_header)?;
688                if let TreeReaderLayout::Blocked {
689                    validated_blocks, ..
690                } = &mut self.layout
691                    && let Some(validated) = validated_blocks.get_mut(block)
692                {
693                    *validated = true;
694                }
695            }
696        }
697        if expected_raw_offset != block_header.raw_payload_len
698            || expected_stored_offset != self.source.len()
699        {
700            return Err(TreeStreamError::Malformed(
701                "tree block lengths do not match the header".into(),
702            ));
703        }
704        Ok(())
705    }
706
707    fn commit_entry(&mut self, entry: &TreeEntry, consumed: usize) -> Result<(), TreeStreamError> {
708        if let Some(previous) = self.cursor.prev_name.as_deref()
709            && previous >= entry.name()
710        {
711            return Err(TreeError::InvalidStructure(
712                "entries must be strictly sorted by name".into(),
713            )
714            .into());
715        }
716        if let Some(hasher) = &mut self.hasher {
717            entry.update_hasher(hasher);
718        }
719        if let Some(entries) = &mut self.lean_hash_entries {
720            entries.push(entry.clone());
721        }
722        self.decoded_logical_len = self
723            .decoded_logical_len
724            .checked_add(entry.encoded_len() as u64)
725            .ok_or_else(|| TreeStreamError::Malformed("logical length overflow".into()))?;
726        self.cursor.ordinal += 1;
727        self.cursor.byte_offset += consumed as u64;
728        self.cursor.prev_name = Some(entry.name().to_string());
729        Ok(())
730    }
731
732    fn arm_pending_at_cursor(&mut self) -> Result<(), TreeStreamError> {
733        if self.cursor.ordinal == 0 || self.cursor.ordinal == self.header.entry_count {
734            return Ok(());
735        }
736        let (entry, consumed) = self.read_entry_at(self.cursor.byte_offset)?;
737        if let Some(previous) = self.cursor.prev_name.as_deref()
738            && previous >= entry.name()
739        {
740            return Err(TreeStreamError::CursorMismatch(
741                "cursor previous name is not a valid predecessor".into(),
742            ));
743        }
744        self.pending = Some((entry, consumed));
745        Ok(())
746    }
747}
748
749#[cfg(test)]
750#[path = "tree_stream_proptests.rs"]
751mod tree_stream_proptests;
752#[cfg(test)]
753#[path = "tree_stream_tests.rs"]
754mod tree_stream_tests;
755
756fn decode_block_frames(
757    raw: &[u8],
758    expected_entries: usize,
759) -> Result<Vec<(usize, usize)>, TreeStreamError> {
760    let mut frames = Vec::with_capacity(expected_entries);
761    let mut offset = 0usize;
762    for _ in 0..expected_entries {
763        let len_bytes = raw
764            .get(offset..offset + 4)
765            .ok_or(TreeStreamError::TruncatedFrame {
766                offset: offset as u64,
767            })?;
768        let frame_len = u32::from_le_bytes(
769            len_bytes
770                .try_into()
771                .map_err(|_| TreeStreamError::Malformed("invalid block frame length".into()))?,
772        ) as usize;
773        let frame_start = offset
774            .checked_add(4)
775            .ok_or(TreeStreamError::TruncatedFrame {
776                offset: offset as u64,
777            })?;
778        let frame_end =
779            frame_start
780                .checked_add(frame_len)
781                .ok_or(TreeStreamError::TruncatedFrame {
782                    offset: offset as u64,
783                })?;
784        if frame_end > raw.len() {
785            return Err(TreeStreamError::TruncatedFrame {
786                offset: offset as u64,
787            });
788        }
789        frames.push((offset, frame_end));
790        offset = frame_end;
791    }
792    if offset != raw.len() {
793        return Err(TreeStreamError::TrailingBytes {
794            extra: raw.len().abs_diff(offset) as u64,
795        });
796    }
797    Ok(frames)
798}
799
800fn validate_cursor(
801    header: &TreeHeader,
802    layout: &TreeReaderLayout,
803    cursor: &TreeResumeCursor,
804) -> Result<(), TreeStreamError> {
805    if cursor.encoding_version != header.version {
806        return Err(TreeStreamError::CursorMismatch(format!(
807            "encoding version {} is not {}",
808            cursor.encoding_version, header.version
809        )));
810    }
811    if cursor.tree_id != header.tree_id {
812        return Err(TreeStreamError::CursorMismatch(
813            "cursor tree id does not match the opened object".into(),
814        ));
815    }
816    let payload_end = match layout {
817        TreeReaderLayout::Raw => TREE_HEADER_LEN as u64 + header.payload_len,
818        TreeReaderLayout::Lean { body_start } => body_start + header.payload_len,
819        TreeReaderLayout::Blocked { header, .. } => TREE_HEADER_LEN as u64 + header.raw_payload_len,
820    };
821    let payload_start = match layout {
822        TreeReaderLayout::Raw | TreeReaderLayout::Blocked { .. } => TREE_HEADER_LEN as u64,
823        TreeReaderLayout::Lean { body_start } => *body_start,
824    };
825    if cursor.ordinal > header.entry_count {
826        return Err(TreeStreamError::CursorMismatch(
827            "cursor ordinal is past the declared entry count".into(),
828        ));
829    }
830    if cursor.ordinal == 0 {
831        if cursor.byte_offset != payload_start || cursor.prev_name.is_some() {
832            return Err(TreeStreamError::CursorMismatch(
833                "start cursor must be the first entry boundary".into(),
834            ));
835        }
836        return Ok(());
837    }
838    if cursor.ordinal == header.entry_count {
839        if cursor.byte_offset != payload_end {
840            return Err(TreeStreamError::CursorMismatch(
841                "end cursor is not the declared payload end".into(),
842            ));
843        }
844        return Ok(());
845    }
846    if cursor.byte_offset < payload_start || cursor.byte_offset >= payload_end {
847        return Err(TreeStreamError::CursorMismatch(
848            "cursor byte offset is not inside the payload".into(),
849        ));
850    }
851    Ok(())
852}
853
854fn read_source_byte<S: TreeByteSource>(
855    source: &mut S,
856    offset: &mut u64,
857) -> Result<u8, TreeStreamError> {
858    let mut byte = [0u8; 1];
859    source.read_exact_at(*offset, &mut byte)?;
860    *offset = offset
861        .checked_add(1)
862        .ok_or_else(|| TreeStreamError::Malformed("tree byte offset overflow".into()))?;
863    Ok(byte[0])
864}
865
866fn read_source_varint<S: TreeByteSource>(
867    source: &mut S,
868    offset: u64,
869) -> Result<(usize, u64), TreeStreamError> {
870    let mut cursor = offset;
871    let value = read_source_varint_at(source, &mut cursor)?;
872    Ok((value, cursor))
873}
874
875fn read_source_varint_at<S: TreeByteSource>(
876    source: &mut S,
877    offset: &mut u64,
878) -> Result<usize, TreeStreamError> {
879    let mut value = 0usize;
880    for shift in (0..usize::BITS).step_by(7) {
881        let byte = read_source_byte(source, offset)?;
882        value |= ((byte & 0x7f) as usize)
883            .checked_shl(shift)
884            .ok_or_else(|| TreeStreamError::Malformed("varint overflow".into()))?;
885        if byte & 0x80 == 0 {
886            return Ok(value);
887        }
888    }
889    Err(TreeStreamError::Malformed("varint overflow".into()))
890}
891
892impl Tree {
893    /// Decode HTR4 through the streaming reader and collect the eager `Tree`.
894    pub fn decode_canonical_streamed(data: &[u8]) -> Result<Self, TreeStreamError> {
895        let header = decode_header(data)?;
896        let mut reader = TreeEntryReader::open(
897            super::tree_source::BytesTreeSource::sequential_verify(bytes::Bytes::copy_from_slice(
898                data,
899            )),
900            header.tree_id,
901            None,
902        )?;
903        let mut entries = Vec::new();
904        while let Some(entry) = reader.next_entry()? {
905            entries.push(entry);
906        }
907        reader.finish_and_verify()?;
908        let tree = Tree::try_from_decoded_entries(entries).map_err(TreeStreamError::from)?;
909        let found = tree.hash();
910        if found != header.tree_id {
911            return Err(TreeStreamError::HashMismatch {
912                expected: header.tree_id,
913                found,
914            });
915        }
916        Ok(tree)
917    }
918}