1use 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#[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#[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#[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#[derive(Clone, Debug, PartialEq, Eq)]
156pub struct TreePage {
157 pub entries: Vec<TreeEntry>,
158 pub resume_cursor: TreeResumeCursor,
159}
160
161#[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
176type 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 return Err(TreeStreamError::Malformed(
192 "HRT1 redacted projection must not enter the tree stream".into(),
193 ));
194 }
195 if &magic == TREE_SALTED_MAGIC {
196 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 pub fn next_entry(&mut self) -> Result<Option<TreeEntry>, TreeStreamError> {
338 Ok(self.next_entry_with_position()?.map(|(entry, _)| entry))
339 }
340
341 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 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 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}