1use 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#[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#[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#[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#[derive(Clone, Debug, PartialEq, Eq)]
155pub struct TreePage {
156 pub entries: Vec<TreeEntry>,
157 pub resume_cursor: TreeResumeCursor,
158}
159
160#[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 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 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}