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