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