1use std::io;
22
23use crate::content::list::{
24 MAXIMUM_LIST_SIZE, uncounted_list_entries, uncounted_list_entry,
25 uncounted_list_entry_with_provider,
26};
27use crate::content::provider::SegmentProvider;
28use crate::error::{Error, Result};
29use crate::segment::record::RecordIdentifier;
30
31pub const SMALL_VALUE_LIMIT: u64 = 128;
33
34pub const MEDIUM_VALUE_LIMIT: u64 = (1 << 14) + 128;
36
37pub const BLOCK_SIZE: u64 = 4096;
39
40#[derive(Clone, PartialEq, Eq, Debug)]
42pub enum BinaryValue {
43 Inline {
45 length: u64,
47 record_identifier: RecordIdentifier,
51 },
52 External {
54 blob_identifier: String,
57 },
58}
59
60#[derive(Clone, Copy, Debug)]
61enum BinaryStreamSource {
62 Direct {
63 record_identifier: RecordIdentifier,
64 content_offset: usize,
65 },
66 Blocks {
67 list_identifier: RecordIdentifier,
68 block_count: u64,
69 },
70}
71
72pub struct BinaryStream<
87 'provider,
88 Provider: SegmentProvider + ?Sized = dyn SegmentProvider + 'provider,
89> {
90 provider: &'provider Provider,
91 source: BinaryStreamSource,
92 length: u64,
93 position: u64,
94 resolved_block: Option<(u64, RecordIdentifier)>,
95}
96
97impl<Provider: SegmentProvider + ?Sized> BinaryStream<'_, Provider> {
98 fn resolve_block_identifier(
99 &mut self,
100 list_identifier: RecordIdentifier,
101 block_count: u64,
102 block_index: u64,
103 ) -> Result<RecordIdentifier> {
104 if let Some((resolved_index, identifier)) = self.resolved_block
105 && resolved_index == block_index
106 {
107 return Ok(identifier);
108 }
109 let identifier = uncounted_list_entry_with_provider(
110 self.provider,
111 list_identifier,
112 block_count,
113 block_index,
114 )?;
115 self.resolved_block = Some((block_index, identifier));
116 Ok(identifier)
117 }
118
119 fn current_block_identifier(&mut self) -> Result<Option<RecordIdentifier>> {
120 if self.position == self.length {
121 return Ok(None);
122 }
123 match self.source {
124 BinaryStreamSource::Direct { .. } => Ok(None),
125 BinaryStreamSource::Blocks {
126 list_identifier,
127 block_count,
128 } => self
129 .resolve_block_identifier(list_identifier, block_count, self.position / BLOCK_SIZE)
130 .map(Some),
131 }
132 }
133
134 #[must_use]
136 pub const fn len(&self) -> u64 {
137 self.length
138 }
139
140 #[must_use]
142 pub const fn is_empty(&self) -> bool {
143 self.length == 0
144 }
145
146 #[must_use]
148 pub const fn position(&self) -> u64 {
149 self.position
150 }
151
152 pub fn read_chunk(&mut self, buffer: &mut [u8]) -> Result<usize> {
159 if buffer.is_empty() || self.position == self.length {
160 return Ok(0);
161 }
162
163 let buffer_capacity = u64::try_from(buffer.len()).unwrap_or(u64::MAX);
164 let remaining = self.length - self.position;
165 let maximum_read = remaining.min(buffer_capacity);
166
167 let read_length = match self.source {
168 BinaryStreamSource::Direct {
169 record_identifier,
170 content_offset,
171 } => {
172 let read_length =
173 usize::try_from(maximum_read).map_err(|_| Error::InvalidFormat {
174 details: format!(
175 "binary read length does not fit this platform in record \
176 {record_identifier}"
177 ),
178 })?;
179 let position =
180 usize::try_from(self.position).map_err(|_| Error::InvalidFormat {
181 details: format!(
182 "binary position does not fit this platform in record \
183 {record_identifier}"
184 ),
185 })?;
186 let offset =
187 content_offset
188 .checked_add(position)
189 .ok_or_else(|| Error::InvalidFormat {
190 details: format!(
191 "binary content offset overflows in record {record_identifier}"
192 ),
193 })?;
194 let view = self.provider.segment(record_identifier.segment)?;
195 buffer[..read_length].copy_from_slice(view.read_bytes(
196 record_identifier.record_number,
197 offset,
198 read_length,
199 )?);
200 read_length
201 }
202 BinaryStreamSource::Blocks {
203 list_identifier,
204 block_count,
205 } => {
206 let block_index = self.position / BLOCK_SIZE;
207 let block_offset = self.position % BLOCK_SIZE;
208 let block_remaining = BLOCK_SIZE - block_offset;
209 let read_length =
210 usize::try_from(maximum_read.min(block_remaining)).map_err(|_| {
211 Error::InvalidFormat {
212 details: format!(
213 "binary block read length does not fit this platform for list \
214 {list_identifier}"
215 ),
216 }
217 })?;
218 let block_identifier =
219 self.resolve_block_identifier(list_identifier, block_count, block_index)?;
220 let view = self.provider.segment(block_identifier.segment)?;
221 buffer[..read_length].copy_from_slice(view.read_bytes(
222 block_identifier.record_number,
223 usize::try_from(block_offset).map_err(|_| Error::InvalidFormat {
224 details: format!(
225 "binary block offset does not fit this platform in record \
226 {block_identifier}"
227 ),
228 })?,
229 read_length,
230 )?);
231 read_length
232 }
233 };
234 let read_length_u64 = u64::try_from(read_length).map_err(|_| Error::InvalidFormat {
235 details: "binary read length does not fit the 64-bit value format".to_owned(),
236 })?;
237 self.position += read_length_u64;
238 Ok(read_length)
239 }
240}
241
242impl<Provider: SegmentProvider + ?Sized> io::Read for BinaryStream<'_, Provider> {
243 fn read(&mut self, buffer: &mut [u8]) -> io::Result<usize> {
244 self.read_chunk(buffer).map_err(|error| match error {
245 Error::InputOutput(source) => source,
246 other => io::Error::other(other),
247 })
248 }
249}
250
251pub fn read_value_length(
254 provider: &dyn SegmentProvider,
255 identifier: RecordIdentifier,
256) -> Result<u64> {
257 let view = provider.segment(identifier.segment)?;
258 let head = view.read_u8(identifier.record_number, 0)?;
259 if head & 0x80 == 0 {
260 Ok(u64::from(head))
261 } else if head & 0x40 == 0 {
262 let stored = view.read_u16(identifier.record_number, 0)?;
263 Ok(u64::from(stored & 0x3FFF) + SMALL_VALUE_LIMIT)
264 } else if head & 0x20 == 0 {
265 let stored = view.read_u64(identifier.record_number, 0)?;
266 Ok((stored & 0x1FFF_FFFF_FFFF_FFFF) + MEDIUM_VALUE_LIMIT)
267 } else {
268 Err(Error::InvalidFormat {
269 details: format!(
270 "value record {identifier} starts with external binary marker {head:#04x}; \
271 its length is not stored in the segment"
272 ),
273 })
274 }
275}
276
277pub(crate) fn read_string_stored_length(
280 provider: &dyn SegmentProvider,
281 identifier: RecordIdentifier,
282) -> Result<u64> {
283 let view = provider.segment(identifier.segment)?;
284 let head = view.read_u8(identifier.record_number, 0)?;
285 if head & 0x80 == 0 {
286 return Ok(u64::from(head));
287 }
288 if head & 0x40 == 0 {
289 return Ok(
290 u64::from(view.read_u16(identifier.record_number, 0)? & 0x3fff) + SMALL_VALUE_LIMIT,
291 );
292 }
293 if head & 0xe0 != 0xc0 {
294 return Err(Error::InvalidFormat {
295 details: format!(
296 "record {identifier} starts with binary marker {head:#04x} and is not a string"
297 ),
298 });
299 }
300 let stored = view.read_u64(identifier.record_number, 0)?;
301 let length = (stored & 0x3fff_ffff_ffff_ffff) + MEDIUM_VALUE_LIMIT;
302 if length >= i32::MAX as u64 {
303 return Err(Error::InvalidFormat {
304 details: format!("string of {length} bytes in record {identifier} is too long"),
305 });
306 }
307 Ok(length)
308}
309
310pub(crate) fn read_string_with_stored_byte_budget(
313 provider: &dyn SegmentProvider,
314 identifier: RecordIdentifier,
315 maximum_stored_bytes: u64,
316 consumed_stored_bytes: &mut u64,
317) -> Result<String> {
318 let length = read_string_stored_length(provider, identifier)?;
319 let attempted_stored_bytes = consumed_stored_bytes.saturating_add(length);
320 if attempted_stored_bytes > maximum_stored_bytes {
321 return Err(Error::StringMaterializationBudgetExceeded {
322 maximum_stored_bytes,
323 attempted_stored_bytes,
324 value_identifier: identifier,
325 });
326 }
327 *consumed_stored_bytes = attempted_stored_bytes;
328 read_string(provider, identifier)
329}
330
331pub fn read_string(provider: &dyn SegmentProvider, identifier: RecordIdentifier) -> Result<String> {
336 let view = provider.segment(identifier.segment)?;
337 let head = view.read_u8(identifier.record_number, 0)?;
338 if head & 0x80 == 0 {
339 let length = usize::from(head);
341 let bytes = view.read_bytes(identifier.record_number, 1, length)?;
342 return Ok(String::from_utf8_lossy(bytes).into_owned());
343 }
344 if head & 0x40 == 0 {
345 let stored = view.read_u16(identifier.record_number, 0)?;
347 let length = usize::from(stored & 0x3FFF) + SMALL_VALUE_LIMIT as usize;
348 let bytes = view.read_bytes(identifier.record_number, 2, length)?;
349 return Ok(String::from_utf8_lossy(bytes).into_owned());
350 }
351 if head & 0x20 != 0 {
352 return Err(Error::InvalidFormat {
353 details: format!(
354 "record {identifier} starts with binary marker {head:#04x} and is not a string"
355 ),
356 });
357 }
358 let stored = view.read_u64(identifier.record_number, 0)?;
362 let length = (stored & 0x3FFF_FFFF_FFFF_FFFF) + MEDIUM_VALUE_LIMIT;
363 if length >= i32::MAX as u64 {
364 return Err(Error::InvalidFormat {
365 details: format!("string of {length} bytes in record {identifier} is too long"),
366 });
367 }
368 let list_identifier = view.read_record_identifier(identifier.record_number, 8, 0)?;
369 let bytes = read_block_list(provider, list_identifier, length)?;
370 Ok(String::from_utf8_lossy(&bytes).into_owned())
371}
372
373fn read_block_list(
376 provider: &dyn SegmentProvider,
377 list_identifier: RecordIdentifier,
378 length: u64,
379) -> Result<Vec<u8>> {
380 let block_count = length.div_ceil(BLOCK_SIZE);
381 let block_identifiers = uncounted_list_entries(provider, list_identifier, block_count)?;
382 let mut content = Vec::with_capacity((length as usize).min(1 << 20));
385 let mut remaining = length;
386 for block_identifier in block_identifiers {
387 let block_length = remaining.min(BLOCK_SIZE) as usize;
388 let view = provider.segment(block_identifier.segment)?;
389 content.extend_from_slice(view.read_bytes(
390 block_identifier.record_number,
391 0,
392 block_length,
393 )?);
394 remaining -= block_length as u64;
395 }
396 Ok(content)
397}
398
399pub fn read_binary_value(
402 provider: &dyn SegmentProvider,
403 identifier: RecordIdentifier,
404) -> Result<BinaryValue> {
405 let view = provider.segment(identifier.segment)?;
406 let head = view.read_u8(identifier.record_number, 0)?;
407 if head & 0x80 == 0 || head & 0x40 == 0 || head & 0x20 == 0 {
408 return Ok(BinaryValue::Inline {
410 length: read_value_length(provider, identifier)?,
411 record_identifier: identifier,
412 });
413 }
414 if head & 0x10 == 0 {
415 let stored = view.read_u16(identifier.record_number, 0)?;
418 let length = usize::from(stored & 0x0FFF);
419 let bytes = view.read_bytes(identifier.record_number, 2, length)?;
420 return Ok(BinaryValue::External {
421 blob_identifier: String::from_utf8_lossy(bytes).into_owned(),
422 });
423 }
424 if head & 0x08 == 0 {
425 let string_identifier = view.read_record_identifier(identifier.record_number, 1, 0)?;
428 return Ok(BinaryValue::External {
429 blob_identifier: read_string(provider, string_identifier)?,
430 });
431 }
432 Err(Error::InvalidFormat {
433 details: format!("unexpected value record marker {head:#04x} in record {identifier}"),
434 })
435}
436
437pub fn read_binary_stream<Provider: SegmentProvider + ?Sized>(
449 provider: &Provider,
450 identifier: RecordIdentifier,
451) -> Result<BinaryStream<'_, Provider>> {
452 let view = provider.segment(identifier.segment)?;
453 let head = view.read_u8(identifier.record_number, 0)?;
454 let (length, source) = if head & 0x80 == 0 {
455 (
456 u64::from(head),
457 BinaryStreamSource::Direct {
458 record_identifier: identifier,
459 content_offset: 1,
460 },
461 )
462 } else if head & 0x40 == 0 {
463 let stored = view.read_u16(identifier.record_number, 0)?;
464 (
465 u64::from(stored & 0x3FFF) + SMALL_VALUE_LIMIT,
466 BinaryStreamSource::Direct {
467 record_identifier: identifier,
468 content_offset: 2,
469 },
470 )
471 } else if head & 0x20 == 0 {
472 let stored = view.read_u64(identifier.record_number, 0)?;
473 let length = (stored & 0x1FFF_FFFF_FFFF_FFFF) + MEDIUM_VALUE_LIMIT;
474 let block_count = length.div_ceil(BLOCK_SIZE);
475 if block_count > MAXIMUM_LIST_SIZE {
476 return Err(Error::InvalidFormat {
477 details: format!(
478 "binary of {length} bytes in record {identifier} needs {block_count} blocks, \
479 exceeding the list maximum of {MAXIMUM_LIST_SIZE}"
480 ),
481 });
482 }
483 let list_identifier = view.read_record_identifier(identifier.record_number, 8, 0)?;
484 (
485 length,
486 BinaryStreamSource::Blocks {
487 list_identifier,
488 block_count,
489 },
490 )
491 } else if head & 0x10 == 0 {
492 let stored = view.read_u16(identifier.record_number, 0)?;
493 let length = usize::from(stored & 0x0FFF);
494 let blob_identifier =
495 String::from_utf8_lossy(view.read_bytes(identifier.record_number, 2, length)?)
496 .into_owned();
497 return Err(Error::ExternalBinaryContentUnavailable { blob_identifier });
498 } else if head & 0x08 == 0 {
499 let string_identifier = view.read_record_identifier(identifier.record_number, 1, 0)?;
500 return Err(Error::ExternalBinaryContentUnavailableByRecord {
501 value_identifier: identifier,
502 blob_identifier_record: string_identifier,
503 });
504 } else {
505 return Err(Error::InvalidFormat {
506 details: format!("unexpected value record marker {head:#04x} in record {identifier}"),
507 });
508 };
509
510 Ok(BinaryStream {
511 provider,
512 source,
513 length,
514 position: 0,
515 resolved_block: None,
516 })
517}
518
519pub fn read_binary_content(
528 provider: &dyn SegmentProvider,
529 identifier: RecordIdentifier,
530) -> Result<Vec<u8>> {
531 let mut stream = match read_binary_stream(provider, identifier) {
532 Ok(stream) => stream,
533 Err(Error::ExternalBinaryContentUnavailableByRecord { .. }) => {
534 let BinaryValue::External { blob_identifier } =
535 read_binary_value(provider, identifier)?
536 else {
537 return Err(Error::InvalidFormat {
538 details: format!(
539 "binary value {identifier} changed classification while resolving its \
540 external identifier"
541 ),
542 });
543 };
544 return Err(Error::ExternalBinaryContentUnavailable { blob_identifier });
545 }
546 Err(error) => return Err(error),
547 };
548 let mut content = Vec::with_capacity(
549 usize::try_from(stream.len())
550 .unwrap_or(usize::MAX)
551 .min(1 << 20),
552 );
553 let mut buffer = [0u8; 8192];
554 loop {
555 let read_length = stream.read_chunk(&mut buffer)?;
556 if read_length == 0 {
557 return Ok(content);
558 }
559 content.extend_from_slice(&buffer[..read_length]);
560 }
561}
562
563pub fn verify_binary_content(
571 provider: &dyn SegmentProvider,
572 identifier: RecordIdentifier,
573) -> Result<()> {
574 let mut stream = match read_binary_stream(provider, identifier) {
575 Ok(stream) => stream,
576 Err(Error::ExternalBinaryContentUnavailable { .. }) => return Ok(()),
577 Err(Error::ExternalBinaryContentUnavailableByRecord {
578 blob_identifier_record,
579 ..
580 }) => return verify_string_content(provider, blob_identifier_record),
581 Err(error) => return Err(error),
582 };
583 let mut buffer = [0u8; 8192];
584 loop {
585 if stream.read_chunk(&mut buffer)? == 0 {
586 return Ok(());
587 }
588 }
589}
590
591fn verify_string_content(
592 provider: &dyn SegmentProvider,
593 identifier: RecordIdentifier,
594) -> Result<()> {
595 let length = read_string_stored_length(provider, identifier)?;
596 let view = provider.segment(identifier.segment)?;
597 let head = view.read_u8(identifier.record_number, 0)?;
598 if head & 0x80 == 0 {
599 view.read_bytes(identifier.record_number, 1, length as usize)?;
600 return Ok(());
601 }
602 if head & 0x40 == 0 {
603 view.read_bytes(identifier.record_number, 2, length as usize)?;
604 return Ok(());
605 }
606
607 let list_identifier = view.read_record_identifier(identifier.record_number, 8, 0)?;
608 let block_count = length.div_ceil(BLOCK_SIZE);
609 let mut remaining = length;
610 for block_index in 0..block_count {
611 let block_identifier =
612 uncounted_list_entry(provider, list_identifier, block_count, block_index)?;
613 let block_length = remaining.min(BLOCK_SIZE) as usize;
614 provider.segment(block_identifier.segment)?.read_bytes(
615 block_identifier.record_number,
616 0,
617 block_length,
618 )?;
619 remaining -= block_length as u64;
620 }
621 Ok(())
622}
623
624pub fn inline_binary_contents_equal(
631 provider: &dyn SegmentProvider,
632 first: RecordIdentifier,
633 second: RecordIdentifier,
634 expected_length: u64,
635) -> Result<bool> {
636 if first == second {
637 return Ok(true);
638 }
639 let mut first_stream = read_binary_stream(provider, first)?;
640 let mut second_stream = read_binary_stream(provider, second)?;
641 if first_stream.len() != expected_length || second_stream.len() != expected_length {
642 return Err(Error::InvalidFormat {
643 details: format!(
644 "inline binary comparison expected {expected_length} bytes, but records {first} \
645 and {second} declare {} and {} bytes",
646 first_stream.len(),
647 second_stream.len()
648 ),
649 });
650 }
651
652 let mut first_buffer = [0u8; 8192];
653 let mut second_buffer = [0u8; 8192];
654 loop {
655 let first_block = first_stream.current_block_identifier()?;
656 let second_block = second_stream.current_block_identifier()?;
657 if first_block.is_some() && first_block == second_block {
658 let first_chunk = (first_stream.length - first_stream.position)
663 .min(BLOCK_SIZE - first_stream.position % BLOCK_SIZE);
664 let second_chunk = (second_stream.length - second_stream.position)
665 .min(BLOCK_SIZE - second_stream.position % BLOCK_SIZE);
666 let skipped = first_chunk.min(second_chunk);
667 first_stream.position += skipped;
668 second_stream.position += skipped;
669 continue;
670 }
671 let first_length = first_stream.read_chunk(&mut first_buffer)?;
672 let second_length = second_stream.read_chunk(&mut second_buffer)?;
673 if first_length != second_length
674 || first_buffer[..first_length] != second_buffer[..second_length]
675 {
676 return Ok(false);
677 }
678 if first_length == 0 {
679 return Ok(true);
680 }
681 }
682}
683
684#[cfg(test)]
685mod tests {
686 use std::cell::Cell;
687 use std::io::Read;
688 use std::sync::Arc;
689
690 use super::{
691 BLOCK_SIZE, BinaryStream, BinaryValue, MEDIUM_VALUE_LIMIT, SMALL_VALUE_LIMIT,
692 inline_binary_contents_equal, read_binary_content, read_binary_stream, read_binary_value,
693 read_string, read_value_length, verify_binary_content,
694 };
695 use crate::content::list::uncounted_list_entry;
696 use crate::content::provider::{SegmentProvider, tests::MemorySegmentProvider};
697 use crate::content::template::{Template, read_template};
698 use crate::error::{Error, Result};
699 use crate::segment::identifier::SegmentIdentifier;
700 use crate::segment::parsed_segment::{
701 MAXIMUM_SEGMENT_SIZE,
702 tests::{bulk_segment_identifier, data_segment_identifier, synthetic_data_segment},
703 };
704 use crate::segment::record::RecordIdentifier;
705 use crate::segment::view::SegmentView;
706
707 fn small_string_record(text: &str) -> Vec<u8> {
708 let mut bytes = vec![text.len() as u8];
709 bytes.extend_from_slice(text.as_bytes());
710 bytes
711 }
712
713 fn direct_binary_record(content: &[u8]) -> Vec<u8> {
714 let length = content.len() as u64;
715 let mut record = if length < SMALL_VALUE_LIMIT {
716 vec![length as u8]
717 } else {
718 assert!(length < MEDIUM_VALUE_LIMIT);
719 ((0x8000u16) | (length as u16 - SMALL_VALUE_LIMIT as u16))
720 .to_be_bytes()
721 .to_vec()
722 };
723 record.extend_from_slice(content);
724 record
725 }
726
727 fn local_record_identifier(record_number: u32) -> Vec<u8> {
728 let mut bytes = vec![0, 0];
729 bytes.extend_from_slice(&record_number.to_be_bytes());
730 bytes
731 }
732
733 fn referenced_record_identifier(reference: u16, record_number: u32) -> Vec<u8> {
734 let mut bytes = reference.to_be_bytes().to_vec();
735 bytes.extend_from_slice(&record_number.to_be_bytes());
736 bytes
737 }
738
739 fn repeated_local_identifiers(record_number: u32, count: usize) -> Vec<u8> {
740 let identifier = local_record_identifier(record_number);
741 let mut bytes = Vec::with_capacity(identifier.len() * count);
742 for _ in 0..count {
743 bytes.extend_from_slice(&identifier);
744 }
745 bytes
746 }
747
748 fn long_binary_record(length: u64, list_record_number: u32) -> Vec<u8> {
749 assert!(length >= MEDIUM_VALUE_LIMIT);
750 let mut record = ((length - MEDIUM_VALUE_LIMIT) | (0x3 << 62))
751 .to_be_bytes()
752 .to_vec();
753 record.extend_from_slice(&local_record_identifier(list_record_number));
754 record
755 }
756
757 struct CountingProvider<'provider> {
758 inner: &'provider MemorySegmentProvider,
759 segment_reads: Cell<usize>,
760 }
761
762 impl<'provider> CountingProvider<'provider> {
763 fn new(inner: &'provider MemorySegmentProvider) -> Self {
764 Self {
765 inner,
766 segment_reads: Cell::new(0),
767 }
768 }
769
770 fn segment_reads(&self) -> usize {
771 self.segment_reads.get()
772 }
773
774 fn reset_segment_reads(&self) {
775 self.segment_reads.set(0);
776 }
777 }
778
779 impl SegmentProvider for CountingProvider<'_> {
780 fn segment(&self, identifier: SegmentIdentifier) -> Result<SegmentView<'_>> {
781 self.segment_reads.set(self.segment_reads.get() + 1);
782 self.inner.segment(identifier)
783 }
784
785 fn string(&self, identifier: RecordIdentifier) -> Result<Arc<str>> {
786 read_string(self, identifier).map(Arc::from)
787 }
788
789 fn template(&self, identifier: RecordIdentifier) -> Result<Arc<Template>> {
790 read_template(self, identifier).map(Arc::new)
791 }
792 }
793
794 #[test]
795 fn reads_small_strings() {
796 let segment = data_segment_identifier(1);
797 let mut provider = MemorySegmentProvider::default();
798 provider.insert(
799 segment,
800 synthetic_data_segment(&[], &[(0, 4, small_string_record("jcr:content"))]),
801 );
802 let identifier = RecordIdentifier::new(segment, 0);
803 assert_eq!(
804 read_value_length(&provider, identifier).expect("length"),
805 11
806 );
807 assert_eq!(
808 read_string(&provider, identifier).expect("string"),
809 "jcr:content"
810 );
811 }
812
813 #[test]
814 fn reads_empty_and_boundary_small_strings() {
815 let segment = data_segment_identifier(1);
816 let longest_small = "x".repeat(127);
817 let mut provider = MemorySegmentProvider::default();
818 provider.insert(
819 segment,
820 synthetic_data_segment(
821 &[],
822 &[
823 (0, 4, small_string_record("")),
824 (1, 4, small_string_record(&longest_small)),
825 ],
826 ),
827 );
828 assert_eq!(
829 read_string(&provider, RecordIdentifier::new(segment, 0)).expect("empty"),
830 ""
831 );
832 assert_eq!(
833 read_string(&provider, RecordIdentifier::new(segment, 1)).expect("boundary"),
834 longest_small
835 );
836 }
837
838 #[test]
839 fn reads_medium_strings() {
840 let segment = data_segment_identifier(1);
841 let text = "y".repeat(128);
842 let mut record = ((0x8000u16) | (text.len() as u16 - 128))
843 .to_be_bytes()
844 .to_vec();
845 record.extend_from_slice(text.as_bytes());
846 let mut provider = MemorySegmentProvider::default();
847 provider.insert(segment, synthetic_data_segment(&[], &[(0, 4, record)]));
848 let identifier = RecordIdentifier::new(segment, 0);
849 assert_eq!(
850 read_value_length(&provider, identifier).expect("length"),
851 128
852 );
853 assert_eq!(read_string(&provider, identifier).expect("string"), text);
854 }
855
856 #[test]
857 fn reads_long_strings_from_block_lists() {
858 let segment = data_segment_identifier(1);
859 let text = "z".repeat(20_000);
860 let mut records: Vec<(u32, u8, Vec<u8>)> = Vec::new();
862 for (block_index, chunk) in text.as_bytes().chunks(4096).enumerate() {
863 records.push((1 + block_index as u32, 5, chunk.to_vec()));
864 }
865 let mut bucket = Vec::new();
867 for block_record in 1..=5u32 {
868 bucket.extend_from_slice(&[0, 0]);
869 bucket.extend_from_slice(&block_record.to_be_bytes());
870 }
871 records.push((10, 2, bucket));
872 let mut value = ((text.len() as u64 - 16512) | (0x3 << 62))
874 .to_be_bytes()
875 .to_vec();
876 value.extend_from_slice(&[0, 0]);
877 value.extend_from_slice(&10u32.to_be_bytes());
878 records.push((11, 4, value));
879
880 let mut provider = MemorySegmentProvider::default();
881 provider.insert(segment, synthetic_data_segment(&[], &records));
882 let identifier = RecordIdentifier::new(segment, 11);
883 assert_eq!(
884 read_value_length(&provider, identifier).expect("length"),
885 20_000
886 );
887 assert_eq!(read_string(&provider, identifier).expect("string"), text);
888 assert_eq!(
889 read_binary_content(&provider, identifier).expect("content"),
890 text.as_bytes()
891 );
892 }
893
894 #[test]
895 fn classifies_external_binaries() {
896 let segment = data_segment_identifier(1);
897 let blob_identifier = "datastore-reference-0001";
898 let mut short_external = ((0xE000u16) | blob_identifier.len() as u16)
899 .to_be_bytes()
900 .to_vec();
901 short_external.extend_from_slice(blob_identifier.as_bytes());
902
903 let mut long_external = vec![0xF0u8];
906 long_external.extend_from_slice(&[0, 0]);
907 long_external.extend_from_slice(&1u32.to_be_bytes());
908 let truncated_external = vec![0xE0, 20, b'x'];
909
910 let mut provider = MemorySegmentProvider::default();
911 provider.insert(
912 segment,
913 synthetic_data_segment(
914 &[],
915 &[
916 (0, 8, short_external),
917 (1, 4, small_string_record(blob_identifier)),
918 (2, 8, long_external),
919 (3, 8, truncated_external),
920 ],
921 ),
922 );
923
924 for record_number in [0u32, 2] {
925 let value = read_binary_value(&provider, RecordIdentifier::new(segment, record_number))
926 .expect("binary value");
927 assert_eq!(
928 value,
929 BinaryValue::External {
930 blob_identifier: blob_identifier.to_owned()
931 },
932 "record {record_number}"
933 );
934 }
935
936 match read_binary_content(&provider, RecordIdentifier::new(segment, 0)) {
937 Err(Error::ExternalBinaryContentUnavailable {
938 blob_identifier: reported,
939 }) => {
940 assert_eq!(reported, blob_identifier);
941 }
942 other => panic!("expected external binary error, got {other:?}"),
943 }
944 match read_binary_content(&provider, RecordIdentifier::new(segment, 2)) {
945 Err(Error::ExternalBinaryContentUnavailable {
946 blob_identifier: reported,
947 }) => assert_eq!(reported, blob_identifier),
948 other => panic!("expected legacy long-external error, got {other:?}"),
949 }
950 match read_binary_stream(&provider, RecordIdentifier::new(segment, 0)) {
951 Err(Error::ExternalBinaryContentUnavailable {
952 blob_identifier: reported,
953 }) => assert_eq!(reported, blob_identifier),
954 _ => panic!("expected short external binary stream error"),
955 }
956 match read_binary_stream(&provider, RecordIdentifier::new(segment, 2)) {
957 Err(Error::ExternalBinaryContentUnavailableByRecord {
958 value_identifier,
959 blob_identifier_record,
960 }) => {
961 assert_eq!(value_identifier, RecordIdentifier::new(segment, 2));
962 assert_eq!(blob_identifier_record, RecordIdentifier::new(segment, 1));
963 }
964 _ => panic!("expected bounded long external binary stream error"),
965 }
966 for record_number in [0u32, 2] {
967 verify_binary_content(&provider, RecordIdentifier::new(segment, record_number))
968 .expect("external binaries have no local content to verify");
969 }
970 assert!(matches!(
971 verify_binary_content(&provider, RecordIdentifier::new(segment, 3)),
972 Err(Error::InvalidFormat { .. })
973 ));
974 }
975
976 #[test]
977 fn long_external_stream_does_not_follow_the_identifier_record() {
978 let value_segment = data_segment_identifier(11);
979 let identifier_segment = data_segment_identifier(12);
980 let mut long_external = vec![0xF0u8];
981 long_external.extend_from_slice(&referenced_record_identifier(1, 7));
982
983 let mut inner = MemorySegmentProvider::default();
984 inner.insert(
985 value_segment,
986 synthetic_data_segment(&[identifier_segment], &[(3, 8, long_external)]),
987 );
988 inner.insert(
989 identifier_segment,
990 synthetic_data_segment(
991 &[],
992 &[(7, 4, small_string_record("identifier-must-not-be-read"))],
993 ),
994 );
995 let provider = CountingProvider::new(&inner);
996 let value_identifier = RecordIdentifier::new(value_segment, 3);
997
998 match read_binary_stream(&provider, value_identifier) {
999 Err(Error::ExternalBinaryContentUnavailableByRecord {
1000 value_identifier: reported_value,
1001 blob_identifier_record,
1002 }) => {
1003 assert_eq!(reported_value, value_identifier);
1004 assert_eq!(
1005 blob_identifier_record,
1006 RecordIdentifier::new(identifier_segment, 7)
1007 );
1008 }
1009 _ => panic!("expected bounded long-external error"),
1010 }
1011 assert_eq!(
1012 provider.segment_reads(),
1013 1,
1014 "opening the stream reads only the value record's segment"
1015 );
1016 provider.reset_segment_reads();
1017 verify_binary_content(&provider, value_identifier)
1018 .expect("consistency verification validates the identifier record");
1019 assert!(
1020 provider.segment_reads() > 1,
1021 "verification follows the identifier while the stream opener stays bounded"
1022 );
1023 }
1024
1025 #[test]
1026 fn long_external_verification_reports_missing_or_non_string_identifiers() {
1027 let value_segment = data_segment_identifier(13);
1028 let identifier_segment = data_segment_identifier(14);
1029 let mut long_external = vec![0xF0u8];
1030 long_external.extend_from_slice(&referenced_record_identifier(1, 7));
1031 let value_identifier = RecordIdentifier::new(value_segment, 3);
1032
1033 let mut missing = MemorySegmentProvider::default();
1034 missing.insert(
1035 value_segment,
1036 synthetic_data_segment(&[identifier_segment], &[(3, 8, long_external.clone())]),
1037 );
1038 assert!(matches!(
1039 verify_binary_content(&missing, value_identifier),
1040 Err(Error::SegmentNotFound { segment_identifier })
1041 if segment_identifier == identifier_segment
1042 ));
1043
1044 let mut nested = MemorySegmentProvider::default();
1045 nested.insert(
1046 value_segment,
1047 synthetic_data_segment(&[identifier_segment], &[(3, 8, long_external)]),
1048 );
1049 nested.insert(
1050 identifier_segment,
1051 synthetic_data_segment(&[], &[(7, 8, vec![0xe0, 0])]),
1052 );
1053 assert!(matches!(
1054 verify_binary_content(&nested, value_identifier),
1055 Err(Error::InvalidFormat { .. })
1056 ));
1057 }
1058
1059 #[test]
1060 fn reads_inline_binaries() {
1061 let segment = data_segment_identifier(1);
1062 let content = vec![0x00u8, 0xFF, 0x7F, 0x80];
1063 let record = direct_binary_record(&content);
1064 let mut provider = MemorySegmentProvider::default();
1065 provider.insert(segment, synthetic_data_segment(&[], &[(0, 4, record)]));
1066 let identifier = RecordIdentifier::new(segment, 0);
1067 let value = read_binary_value(&provider, identifier).expect("binary value");
1068 assert_eq!(
1069 value,
1070 BinaryValue::Inline {
1071 length: 4,
1072 record_identifier: identifier
1073 }
1074 );
1075 assert_eq!(
1076 read_binary_content(&provider, identifier).expect("content"),
1077 content
1078 );
1079 }
1080
1081 #[test]
1082 fn streams_small_and_medium_binaries_through_io_read() {
1083 let segment = data_segment_identifier(1);
1084 let empty = Vec::new();
1085 let boundary_small: Vec<u8> = (0..127).map(|index| index as u8).collect();
1086 let first_medium: Vec<u8> = (0..128).map(|index| index as u8).collect();
1087 let boundary_medium: Vec<u8> = (0..16_511).map(|index| (index % 251) as u8).collect();
1088 let mut provider = MemorySegmentProvider::default();
1089 provider.insert(
1090 segment,
1091 synthetic_data_segment(
1092 &[],
1093 &[
1094 (0, 4, direct_binary_record(&empty)),
1095 (1, 4, direct_binary_record(&boundary_small)),
1096 (2, 4, direct_binary_record(&first_medium)),
1097 (3, 4, direct_binary_record(&boundary_medium)),
1098 ],
1099 ),
1100 );
1101
1102 for (record_number, expected) in [
1103 (0, empty.as_slice()),
1104 (1, boundary_small.as_slice()),
1105 (2, first_medium.as_slice()),
1106 (3, boundary_medium.as_slice()),
1107 ] {
1108 let mut stream =
1109 read_binary_stream(&provider, RecordIdentifier::new(segment, record_number))
1110 .expect("stream");
1111 assert_eq!(stream.len(), expected.len() as u64);
1112 assert_eq!(stream.is_empty(), expected.is_empty());
1113 assert_eq!(stream.position(), 0);
1114
1115 let mut actual = Vec::new();
1116 let mut buffer = [0u8; 37];
1117 loop {
1118 let read_length = stream.read(&mut buffer).expect("read");
1119 if read_length == 0 {
1120 break;
1121 }
1122 actual.extend_from_slice(&buffer[..read_length]);
1123 }
1124 assert_eq!(actual, expected);
1125 assert_eq!(stream.position(), expected.len() as u64);
1126 assert_eq!(
1127 read_binary_content(&provider, RecordIdentifier::new(segment, record_number))
1128 .expect("compatibility helper"),
1129 expected
1130 );
1131 }
1132 }
1133
1134 #[test]
1135 fn streams_long_binary_from_partial_bulk_segment_across_block_boundaries() {
1136 let data_segment = data_segment_identifier(1);
1137 let bulk_segment = bulk_segment_identifier(2);
1138 let content: Vec<u8> = (0..20_000).map(|index| (index % 251) as u8).collect();
1139 let first_virtual_offset = (MAXIMUM_SEGMENT_SIZE - content.len()) as u32;
1140
1141 let mut block_list = Vec::new();
1142 for block_offset in (0..content.len()).step_by(BLOCK_SIZE as usize) {
1143 block_list.extend_from_slice(&referenced_record_identifier(
1144 1,
1145 first_virtual_offset + block_offset as u32,
1146 ));
1147 }
1148 let value_record = long_binary_record(content.len() as u64, 10);
1149
1150 let mut provider = MemorySegmentProvider::default();
1151 provider.insert(bulk_segment, content.clone());
1152 provider.insert(
1153 data_segment,
1154 synthetic_data_segment(
1155 &[bulk_segment],
1156 &[(10, 2, block_list), (11, 4, value_record)],
1157 ),
1158 );
1159
1160 let identifier = RecordIdentifier::new(data_segment, 11);
1161 let mut stream = read_binary_stream(&provider, identifier).expect("stream");
1162 let mut first = [0u8; 17];
1163 assert_eq!(stream.read(&mut first).expect("first bytes"), first.len());
1164 assert_eq!(&first, &content[..17]);
1165
1166 let mut through_boundary = [0u8; 5000];
1167 assert_eq!(
1168 stream.read(&mut through_boundary).expect("rest of block"),
1169 BLOCK_SIZE as usize - first.len(),
1170 "a read stops at the current block boundary"
1171 );
1172 assert_eq!(stream.position(), BLOCK_SIZE);
1173
1174 let mut actual = content[..BLOCK_SIZE as usize].to_vec();
1175 stream.read_to_end(&mut actual).expect("remaining blocks");
1176 assert_eq!(actual, content);
1177 assert_eq!(
1178 read_binary_content(&provider, identifier).expect("compatibility helper"),
1179 content
1180 );
1181 }
1182
1183 #[test]
1184 fn large_declared_binary_resolves_only_the_current_list_branch() {
1185 let segment = data_segment_identifier(1);
1186 let block_count = 256u64;
1187 let length = block_count * BLOCK_SIZE;
1188
1189 let mut child_bucket = Vec::new();
1194 for _ in 0..255 {
1195 child_bucket.extend_from_slice(&local_record_identifier(1));
1196 }
1197 let mut top_bucket = local_record_identifier(10);
1198 top_bucket.extend_from_slice(&local_record_identifier(2));
1199 let mut repeated_block = vec![0u8; BLOCK_SIZE as usize];
1200 repeated_block[0] = 0xAB;
1201 repeated_block[1] = 0xCD;
1202 let records = [
1203 (1, 5, repeated_block),
1204 (2, 5, vec![0xEF; BLOCK_SIZE as usize]),
1205 (10, 2, child_bucket),
1206 (11, 2, top_bucket),
1207 (12, 4, long_binary_record(length, 11)),
1208 ];
1209 let mut inner = MemorySegmentProvider::default();
1210 inner.insert(segment, synthetic_data_segment(&[], &records));
1211 let provider = CountingProvider::new(&inner);
1212
1213 let mut stream =
1214 read_binary_stream(&provider, RecordIdentifier::new(segment, 12)).expect("stream");
1215 assert_eq!(stream.len(), length);
1216 assert_eq!(
1217 provider.segment_reads(),
1218 1,
1219 "opening reads only the value record"
1220 );
1221
1222 provider.reset_segment_reads();
1223 let mut byte = [0u8; 1];
1224 assert_eq!(stream.read(&mut byte).expect("first byte"), 1);
1225 assert_eq!(byte, [0xAB]);
1226 assert_eq!(
1227 provider.segment_reads(),
1228 3,
1229 "one top bucket, one child bucket, and the current block"
1230 );
1231
1232 assert_eq!(stream.read(&mut byte).expect("second byte"), 1);
1233 assert_eq!(byte, [0xCD]);
1234 assert_eq!(
1235 provider.segment_reads(),
1236 4,
1237 "the one-entry block cache avoids traversing the list again"
1238 );
1239
1240 let mut last_byte = None;
1245 let mut buffer = [0u8; 8192];
1246 loop {
1247 let read_length = stream.read(&mut buffer).expect("remaining binary");
1248 if read_length == 0 {
1249 break;
1250 }
1251 last_byte = Some(buffer[read_length - 1]);
1252 }
1253 assert_eq!(stream.position(), length);
1254 assert_eq!(last_byte, Some(0xEF));
1255 let block_count = usize::try_from(block_count).expect("fixture block count fits usize");
1256 assert_eq!(
1257 provider.segment_reads(),
1258 4 + 1 + (block_count - 2) * 3 + 2,
1259 "blocks 1 through 254 resolve two buckets and one block; the pass-through final \
1260 child resolves one bucket and one block"
1261 );
1262 }
1263
1264 #[test]
1265 fn canonical_list_resolver_accepts_concrete_and_erased_providers_at_boundaries() {
1266 let segment = data_segment_identifier(1);
1267
1268 let mut first_three_level_root = local_record_identifier(101);
1272 first_three_level_root.extend_from_slice(&local_record_identifier(2));
1273
1274 let maximum_size = crate::content::list::MAXIMUM_LIST_SIZE;
1278 let records = [
1279 (1, 5, vec![0x11]),
1280 (2, 5, vec![0x22]),
1281 (3, 5, vec![0x33]),
1282 (100, 2, first_three_level_root),
1283 (101, 2, repeated_local_identifiers(102, 255)),
1284 (102, 2, repeated_local_identifiers(1, 255)),
1285 (200, 2, repeated_local_identifiers(201, 255)),
1286 (201, 2, repeated_local_identifiers(202, 255)),
1287 (202, 2, repeated_local_identifiers(3, 255)),
1288 ];
1289 let mut provider = MemorySegmentProvider::default();
1290 provider.insert(segment, synthetic_data_segment(&[], &records));
1291
1292 for (list_record, size, index, expected_record) in [
1293 (100, 65_026, 65_024, 1),
1294 (100, 65_026, 65_025, 2),
1295 (200, maximum_size, 0, 3),
1296 (200, maximum_size, maximum_size - 1, 3),
1297 ] {
1298 let list_identifier = RecordIdentifier::new(segment, list_record);
1299 let expected = RecordIdentifier::new(segment, expected_record);
1300 assert_eq!(
1301 uncounted_list_entry(&provider, list_identifier, size, index)
1302 .expect("concrete provider"),
1303 expected
1304 );
1305 let erased: &dyn SegmentProvider = &provider;
1306 assert_eq!(
1307 uncounted_list_entry(erased, list_identifier, size, index)
1308 .expect("erased provider"),
1309 expected,
1310 "generic and erased traversal differ at size {size}, index {index}"
1311 );
1312 }
1313 }
1314
1315 #[test]
1316 fn binary_verification_streams_a_large_list_and_reports_a_truncated_second_block() {
1317 let segment = data_segment_identifier(1);
1318 let block_count = 65_026u64;
1319 let length = block_count * BLOCK_SIZE;
1320
1321 let mut root = local_record_identifier(101);
1327 root.extend_from_slice(&local_record_identifier(900));
1328 let mut leaf = local_record_identifier(800);
1329 leaf.extend_from_slice(&repeated_local_identifiers(900, 254));
1330 let records = [
1331 (10, 4, long_binary_record(length, 100)),
1332 (100, 2, root),
1333 (101, 2, repeated_local_identifiers(102, 255)),
1334 (102, 2, leaf),
1335 (800, 5, vec![0x11; BLOCK_SIZE as usize]),
1336 (900, 5, vec![0x22; BLOCK_SIZE as usize - 4]),
1337 ];
1338 let mut inner = MemorySegmentProvider::default();
1339 inner.insert(segment, synthetic_data_segment(&[], &records));
1340 let provider = CountingProvider::new(&inner);
1341
1342 let error = verify_binary_content(&provider, RecordIdentifier::new(segment, 10))
1343 .expect_err("truncated second block");
1344 assert!(matches!(error, Error::InvalidFormat { .. }));
1345 assert_eq!(
1346 provider.segment_reads(),
1347 9,
1348 "one value head plus three list levels and one block per streamed chunk"
1349 );
1350 }
1351
1352 #[test]
1353 fn binary_comparison_streams_large_lists_across_a_boundary_before_truncation() {
1354 let segment = data_segment_identifier(1);
1355 let block_count = 65_026u64;
1356 let length = block_count * BLOCK_SIZE;
1357
1358 let mut first_root = local_record_identifier(101);
1359 first_root.extend_from_slice(&local_record_identifier(900));
1360 let mut first_leaf = local_record_identifier(800);
1361 first_leaf.extend_from_slice(&repeated_local_identifiers(801, 254));
1362 let mut second_root = local_record_identifier(201);
1363 second_root.extend_from_slice(&local_record_identifier(900));
1364 let mut second_leaf = local_record_identifier(802);
1365 second_leaf.extend_from_slice(&repeated_local_identifiers(900, 254));
1366 let records = [
1367 (10, 4, long_binary_record(length, 100)),
1368 (11, 4, long_binary_record(length, 200)),
1369 (100, 2, first_root),
1370 (101, 2, repeated_local_identifiers(102, 255)),
1371 (102, 2, first_leaf),
1372 (200, 2, second_root),
1373 (201, 2, repeated_local_identifiers(202, 255)),
1374 (202, 2, second_leaf),
1375 (800, 5, vec![0x33; BLOCK_SIZE as usize]),
1376 (801, 5, vec![0x44; BLOCK_SIZE as usize]),
1377 (802, 5, vec![0x33; BLOCK_SIZE as usize]),
1378 (900, 5, vec![0x44; BLOCK_SIZE as usize - 4]),
1379 ];
1380 let mut inner = MemorySegmentProvider::default();
1381 inner.insert(segment, synthetic_data_segment(&[], &records));
1382 let provider = CountingProvider::new(&inner);
1383
1384 let error = inline_binary_contents_equal(
1385 &provider,
1386 RecordIdentifier::new(segment, 10),
1387 RecordIdentifier::new(segment, 11),
1388 length,
1389 )
1390 .expect_err("second stream has a truncated second block");
1391 assert!(matches!(error, Error::InvalidFormat { .. }));
1392 assert_eq!(
1393 provider.segment_reads(),
1394 18,
1395 "two value heads plus two three-level list traversals and block reads per chunk"
1396 );
1397 }
1398
1399 #[test]
1400 fn binary_comparison_checks_the_supplied_length_against_both_records() {
1401 let segment = data_segment_identifier(1);
1402 let records = [
1403 (10, 4, direct_binary_record(b"abc")),
1404 (11, 4, direct_binary_record(b"abc")),
1405 ];
1406 let mut provider = MemorySegmentProvider::default();
1407 provider.insert(segment, synthetic_data_segment(&[], &records));
1408 let first = RecordIdentifier::new(segment, 10);
1409 let second = RecordIdentifier::new(segment, 11);
1410
1411 let error = inline_binary_contents_equal(&provider, first, second, 2)
1412 .expect_err("the caller-provided length is part of the comparison contract");
1413 assert!(matches!(
1414 error,
1415 Error::InvalidFormat { details }
1416 if details.contains("expected 2 bytes")
1417 && details.contains("declare 3 and 3 bytes")
1418 ));
1419 assert!(
1420 inline_binary_contents_equal(&provider, first, second, 3)
1421 .expect("matching declared and expected lengths")
1422 );
1423 }
1424
1425 #[test]
1426 fn binary_comparison_preserves_the_shared_block_identifier_fast_path() {
1427 let segment = data_segment_identifier(1);
1428 let length = MEDIUM_VALUE_LIMIT;
1429 let shared_missing_blocks = repeated_local_identifiers(999, 5);
1430 let records = [
1431 (10, 4, long_binary_record(length, 20)),
1432 (11, 4, long_binary_record(length, 21)),
1433 (20, 2, shared_missing_blocks.clone()),
1434 (21, 2, shared_missing_blocks),
1435 ];
1436 let mut inner = MemorySegmentProvider::default();
1437 inner.insert(segment, synthetic_data_segment(&[], &records));
1438 let provider = CountingProvider::new(&inner);
1439
1440 assert!(
1441 inline_binary_contents_equal(
1442 &provider,
1443 RecordIdentifier::new(segment, 10),
1444 RecordIdentifier::new(segment, 11),
1445 length,
1446 )
1447 .expect("immutable shared block identifiers are content-equal")
1448 );
1449 assert_eq!(
1450 provider.segment_reads(),
1451 12,
1452 "the two heads and two lazy list lookups per block are read, but shared blocks are not"
1453 );
1454 }
1455
1456 #[test]
1457 fn stream_preserves_corrupt_and_missing_record_errors() {
1458 let segment = data_segment_identifier(1);
1459 let truncated_direct = direct_binary_record(&[1, 2, 3, 4]);
1460 let mut provider = MemorySegmentProvider::default();
1461 provider.insert(
1462 segment,
1463 synthetic_data_segment(&[], &[(0, 4, truncated_direct[..3].to_vec())]),
1464 );
1465 let mut stream =
1466 read_binary_stream(&provider, RecordIdentifier::new(segment, 0)).expect("stream");
1467 let mut content = [0u8; 4];
1468 let error = stream
1469 .read_chunk(&mut content)
1470 .expect_err("truncated value");
1471 assert!(matches!(error, Error::InvalidFormat { .. }));
1472 assert_eq!(stream.position(), 0, "failed reads do not advance");
1473 assert!(matches!(
1474 read_binary_content(&provider, RecordIdentifier::new(segment, 0)),
1475 Err(Error::InvalidFormat { .. })
1476 ));
1477
1478 let mut stream =
1479 read_binary_stream(&provider, RecordIdentifier::new(segment, 0)).expect("stream");
1480 let error = stream.read(&mut content).expect_err("io error");
1481 assert!(matches!(
1482 error
1483 .get_ref()
1484 .and_then(|source| source.downcast_ref::<Error>()),
1485 Some(Error::InvalidFormat { .. })
1486 ));
1487 }
1488
1489 #[test]
1490 fn stream_reports_missing_long_block_after_completed_prefix() {
1491 let segment = data_segment_identifier(1);
1492 let length = MEDIUM_VALUE_LIMIT;
1493 let mut block_list = Vec::new();
1494 for record_number in [1u32, 2, 3, 4, 99] {
1495 block_list.extend_from_slice(&local_record_identifier(record_number));
1496 }
1497 let records = [
1498 (1, 5, vec![1; BLOCK_SIZE as usize]),
1499 (2, 5, vec![2; BLOCK_SIZE as usize]),
1500 (3, 5, vec![3; BLOCK_SIZE as usize]),
1501 (4, 5, vec![4; BLOCK_SIZE as usize]),
1502 (10, 2, block_list),
1503 (11, 4, long_binary_record(length, 10)),
1504 ];
1505 let mut provider = MemorySegmentProvider::default();
1506 provider.insert(segment, synthetic_data_segment(&[], &records));
1507 let mut stream =
1508 read_binary_stream(&provider, RecordIdentifier::new(segment, 11)).expect("stream");
1509 let mut block = [0u8; BLOCK_SIZE as usize];
1510 for expected in 1..=4 {
1511 assert_eq!(
1512 stream.read_chunk(&mut block).expect("present block"),
1513 block.len()
1514 );
1515 assert!(block.iter().all(|&byte| byte == expected));
1516 }
1517 let error = stream
1518 .read_chunk(&mut block)
1519 .expect_err("missing fifth block");
1520 assert!(matches!(error, Error::InvalidFormat { .. }));
1521 assert_eq!(stream.position(), 4 * BLOCK_SIZE);
1522 }
1523
1524 #[test]
1525 fn stream_rejects_a_truncated_block_list_bucket() {
1526 let segment = data_segment_identifier(1);
1527 let length = MEDIUM_VALUE_LIMIT;
1528 let mut truncated_block_list = Vec::new();
1529 for record_number in 1..=4u32 {
1530 truncated_block_list.extend_from_slice(&local_record_identifier(record_number));
1531 }
1532 let records = [
1536 (1, 5, vec![1; BLOCK_SIZE as usize]),
1537 (2, 5, vec![2; BLOCK_SIZE as usize]),
1538 (3, 5, vec![3; BLOCK_SIZE as usize]),
1539 (4, 5, vec![4; BLOCK_SIZE as usize]),
1540 (10, 4, long_binary_record(length, 20)),
1541 (20, 2, truncated_block_list),
1542 ];
1543 let mut provider = MemorySegmentProvider::default();
1544 provider.insert(segment, synthetic_data_segment(&[], &records));
1545 let mut stream =
1546 read_binary_stream(&provider, RecordIdentifier::new(segment, 10)).expect("stream");
1547 let mut block = [0u8; BLOCK_SIZE as usize];
1548 for _ in 0..4 {
1549 assert_eq!(
1550 stream.read_chunk(&mut block).expect("present block"),
1551 block.len()
1552 );
1553 }
1554 match stream.read_chunk(&mut block) {
1555 Err(Error::InvalidFormat { details }) => {
1556 assert!(details.contains("record 20"), "unexpected error: {details}");
1557 }
1558 _ => panic!("expected truncated list bucket error"),
1559 }
1560 assert_eq!(stream.position(), 4 * BLOCK_SIZE);
1561 }
1562
1563 #[test]
1564 fn stream_preserves_missing_segment_identity() {
1565 let data_segment = data_segment_identifier(1);
1566 let missing_bulk_segment = bulk_segment_identifier(2);
1567 let length = MEDIUM_VALUE_LIMIT;
1568 let missing_block =
1569 referenced_record_identifier(1, (MAXIMUM_SEGMENT_SIZE - BLOCK_SIZE as usize) as u32);
1570 let mut block_list = Vec::new();
1571 for _ in 0..length.div_ceil(BLOCK_SIZE) {
1572 block_list.extend_from_slice(&missing_block);
1573 }
1574 let mut provider = MemorySegmentProvider::default();
1575 provider.insert(
1576 data_segment,
1577 synthetic_data_segment(
1578 &[missing_bulk_segment],
1579 &[(10, 2, block_list), (11, 4, long_binary_record(length, 10))],
1580 ),
1581 );
1582
1583 let mut stream =
1584 read_binary_stream(&provider, RecordIdentifier::new(data_segment, 11)).expect("stream");
1585 let mut byte = [0u8; 1];
1586 match stream.read_chunk(&mut byte) {
1587 Err(Error::SegmentNotFound { segment_identifier }) => {
1588 assert_eq!(segment_identifier, missing_bulk_segment);
1589 }
1590 _ => panic!("expected missing bulk segment error"),
1591 }
1592 assert_eq!(stream.position(), 0);
1593 }
1594
1595 #[test]
1596 fn stream_rejects_truncated_heads_and_oversized_block_lists() {
1597 let segment = data_segment_identifier(1);
1598 let truncated = ((MEDIUM_VALUE_LIMIT - MEDIUM_VALUE_LIMIT) | (0x3 << 62))
1599 .to_be_bytes()
1600 .to_vec();
1601 let oversized_length = (crate::content::list::MAXIMUM_LIST_SIZE + 1) * BLOCK_SIZE;
1602 let oversized = long_binary_record(oversized_length, 99);
1603 let mut provider = MemorySegmentProvider::default();
1604 provider.insert(
1605 segment,
1606 synthetic_data_segment(&[], &[(0, 4, truncated), (1, 4, oversized)]),
1607 );
1608
1609 assert!(matches!(
1610 read_binary_stream(&provider, RecordIdentifier::new(segment, 0)),
1611 Err(Error::InvalidFormat { .. })
1612 ));
1613 match read_binary_stream(&provider, RecordIdentifier::new(segment, 1)) {
1614 Err(Error::InvalidFormat { details }) => {
1615 assert!(details.contains("exceeding the list maximum"));
1616 }
1617 _ => panic!("expected oversized list error"),
1618 }
1619 }
1620
1621 #[test]
1622 fn binary_stream_is_send_when_its_provider_is_sync() {
1623 fn assert_send<Type: Send>() {}
1624 assert_send::<BinaryStream<'static, MemorySegmentProvider>>();
1625 }
1626
1627 #[test]
1628 fn rejects_invalid_markers() {
1629 let segment = data_segment_identifier(1);
1630 let mut provider = MemorySegmentProvider::default();
1631 provider.insert(
1632 segment,
1633 synthetic_data_segment(&[], &[(0, 4, vec![0xF8, 0, 0, 0])]),
1634 );
1635 let identifier = RecordIdentifier::new(segment, 0);
1636 assert!(read_binary_value(&provider, identifier).is_err());
1637 assert!(read_binary_stream(&provider, identifier).is_err());
1638 assert!(read_string(&provider, identifier).is_err());
1639 assert!(read_value_length(&provider, identifier).is_err());
1640 }
1641}