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