1use std::{
58 collections::HashSet,
59 fs::{File, OpenOptions},
60 io::{self, BufWriter, Cursor, Read, Seek, SeekFrom, Write},
61 path::PathBuf,
62};
63
64use heddle_format::compression::CompressionConfig;
65
66use super::{ObjectType, PackObjectId, PackStats, pack_container_spec, write_container_header};
67
68#[cfg(feature = "zstd")]
74const CSIZE_PLACEHOLDER_LEN: usize = 10;
75use crate::{
76 object::ContentHash,
77 store::{Result, StoreError},
78};
79
80const BUCKETS_PER_VARIANT: usize = 256;
86const TOTAL_BUCKETS: usize = BUCKETS_PER_VARIANT * 3;
88const MAX_OPEN_BUCKET_WRITERS: usize = 32;
92
93const HASH_VARIANT: usize = 0;
97const CHANGEID_VARIANT: usize = 1;
98const ANNOTATED_TAG_VARIANT: usize = 2;
99
100pub trait SyncData {
106 fn sync_data_for_durability(&mut self) -> io::Result<()>;
107}
108
109impl SyncData for File {
110 fn sync_data_for_durability(&mut self) -> io::Result<()> {
111 self.sync_all()
112 }
113}
114
115impl SyncData for Cursor<Vec<u8>> {
116 fn sync_data_for_durability(&mut self) -> io::Result<()> {
117 Ok(())
118 }
119}
120
121pub struct StreamingPackBuilder<W: Write + Read + Seek> {
124 pack_writer: Option<BufWriter<W>>,
130 header_offset: u64,
134 pack_position: u64,
137 record_count: u64,
138 object_count: u64,
139 declared_object_count: Option<u64>,
140 total_uncompressed: u64,
141 total_compressed: u64,
142 #[cfg_attr(not(feature = "zstd"), allow(dead_code))]
147 compression: CompressionConfig,
148 bucket_dir: PathBuf,
151 scratch_lease: Option<super::ScratchLease>,
152 bucket_writers: Vec<Option<BucketWriter>>,
156 open_bucket_writers: usize,
157 bucket_access_tick: u64,
158 bucket_paths: Vec<PathBuf>,
159 index_path: PathBuf,
163 finalized: bool,
166 durable: bool,
167}
168
169struct BucketWriter {
170 writer: BufWriter<File>,
171 last_used: u64,
172}
173
174#[cfg(feature = "zstd")]
175struct CountingWriter<'a, W: Write> {
176 inner: &'a mut W,
177 written: u64,
178}
179
180#[cfg(feature = "zstd")]
181impl<'a, W: Write> CountingWriter<'a, W> {
182 fn new(inner: &'a mut W) -> Self {
183 Self { inner, written: 0 }
184 }
185}
186
187#[cfg(feature = "zstd")]
188impl<W: Write> Write for CountingWriter<'_, W> {
189 fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
190 let written = self.inner.write(buf)?;
191 self.written = self.written.saturating_add(written as u64);
192 Ok(written)
193 }
194
195 fn flush(&mut self) -> std::io::Result<()> {
196 self.inner.flush()
197 }
198}
199
200impl<W: Write + Read + Seek + SyncData> StreamingPackBuilder<W> {
201 #[cfg(feature = "fs")]
217 pub fn new(
218 pack_writer: W,
219 index_path: PathBuf,
220 compression: CompressionConfig,
221 bucket_dir: PathBuf,
222 ) -> Result<Self> {
223 Self::new_inner(pack_writer, index_path, compression, bucket_dir, None, true)
224 }
225
226 #[cfg(feature = "fs")]
233 pub fn new_with_object_count(
234 pack_writer: W,
235 index_path: PathBuf,
236 compression: CompressionConfig,
237 bucket_dir: PathBuf,
238 object_count: u64,
239 ) -> Result<Self> {
240 Self::new_inner(
241 pack_writer,
242 index_path,
243 compression,
244 bucket_dir,
245 Some(object_count),
246 true,
247 )
248 }
249
250 pub fn new_with_object_count_ephemeral(
257 pack_writer: W,
258 index_path: PathBuf,
259 compression: CompressionConfig,
260 bucket_dir: PathBuf,
261 object_count: u64,
262 ) -> Result<Self> {
263 Self::new_inner(
264 pack_writer,
265 index_path,
266 compression,
267 bucket_dir,
268 Some(object_count),
269 false,
270 )
271 }
272
273 fn new_inner(
274 mut pack_writer: W,
275 index_path: PathBuf,
276 compression: CompressionConfig,
277 bucket_dir: PathBuf,
278 declared_object_count: Option<u64>,
279 durable: bool,
280 ) -> Result<Self> {
281 #[cfg(feature = "fs")]
282 if durable {
283 heddle_fs_prims::fs_atomic::create_dir_all_durable(&bucket_dir)
284 .map_err(StoreError::from)?;
285 } else {
286 std::fs::create_dir_all(&bucket_dir).map_err(StoreError::from)?;
287 }
288 #[cfg(not(feature = "fs"))]
289 {
290 debug_assert!(!durable);
291 std::fs::create_dir_all(&bucket_dir).map_err(StoreError::from)?;
292 }
293 let scratch_lease = super::ScratchLease::acquire(&bucket_dir)?;
294 let header_offset = pack_writer.stream_position().map_err(StoreError::from)?;
295
296 let mut header_bytes = Vec::with_capacity(16);
301 write_container_header(
302 &mut header_bytes,
303 pack_container_spec(),
304 declared_object_count.unwrap_or(0),
305 );
306 pack_writer
307 .write_all(&header_bytes)
308 .map_err(StoreError::from)?;
309
310 let bucket_paths: Vec<PathBuf> = (0..TOTAL_BUCKETS)
311 .map(|i| {
312 let variant = match i / BUCKETS_PER_VARIANT {
313 HASH_VARIANT => 'h',
314 CHANGEID_VARIANT => 's',
315 ANNOTATED_TAG_VARIANT => 't',
316 _ => unreachable!("bucket variant is bounded by TOTAL_BUCKETS"),
317 };
318 let prefix = i % BUCKETS_PER_VARIANT;
319 bucket_dir.join(format!("bucket-{variant}-{prefix:02x}"))
320 })
321 .collect();
322 for path in &bucket_paths {
323 let _ = std::fs::remove_file(path);
324 }
325
326 Ok(Self {
327 pack_writer: Some(BufWriter::new(pack_writer)),
328 header_offset,
329 pack_position: header_offset + header_bytes.len() as u64,
330 record_count: 0,
331 object_count: 0,
332 declared_object_count,
333 total_uncompressed: 0,
334 total_compressed: 0,
335 compression,
336 bucket_dir,
337 scratch_lease: Some(scratch_lease),
338 bucket_writers: (0..TOTAL_BUCKETS).map(|_| None).collect(),
339 open_bucket_writers: 0,
340 bucket_access_tick: 0,
341 bucket_paths,
342 index_path,
343 finalized: false,
344 durable,
345 })
346 }
347
348 pub fn flush_pack(&mut self) -> Result<()> {
353 if let Some(writer) = self.pack_writer.as_mut() {
354 writer.flush().map_err(StoreError::from)?;
355 }
356 Ok(())
357 }
358
359 #[cfg(feature = "source-transfer")]
360 pub(super) fn scratch_root(&self) -> &std::path::Path {
361 &self.bucket_dir
362 }
363
364 pub fn add(&mut self, hash: ContentHash, obj_type: ObjectType, data: Vec<u8>) -> Result<()> {
366 self.add_id(PackObjectId::Hash(hash), obj_type, data)
367 }
368
369 pub fn add_id(
392 &mut self,
393 id: PackObjectId,
394 obj_type: ObjectType,
395 data: impl AsRef<[u8]>,
396 ) -> Result<()> {
397 let data = data.as_ref();
398 let pw = self
403 .pack_writer
404 .as_mut()
405 .ok_or_else(|| StoreError::InvalidObject("pack builder is finalized".into()))?;
406 let entry_start = self.pack_position;
407 let offset = entry_start
408 .checked_sub(self.header_offset)
409 .ok_or_else(|| StoreError::InvalidObject("pack position precedes its header".into()))?;
410
411 self.total_uncompressed = self
412 .total_uncompressed
413 .checked_add(data.len() as u64)
414 .ok_or_else(|| StoreError::InvalidObject("pack decoded size overflow".into()))?;
415
416 let mut header_buf = Vec::with_capacity(40);
419 id.encode_tagged(&mut header_buf);
420 let encoded_type = if obj_type == ObjectType::AnnotatedTag {
421 ObjectType::Blob
422 } else {
423 obj_type
424 };
425 super::varint::encode_type_and_size(encoded_type, data.len() as u64, &mut header_buf);
426 pw.write_all(&header_buf).map_err(StoreError::from)?;
427 self.pack_position = self
428 .pack_position
429 .checked_add(header_buf.len() as u64)
430 .ok_or_else(|| {
431 StoreError::InvalidObject("streaming pack position overflow".to_string())
432 })?;
433 #[cfg(feature = "zstd")]
437 let csize_pos = entry_start + header_buf.len() as u64;
438
439 let want_compress: bool;
452 #[cfg(feature = "zstd")]
453 {
454 want_compress = self.compression.enabled && data.len() >= self.compression.min_size;
455 }
456 #[cfg(not(feature = "zstd"))]
457 {
458 want_compress = false;
459 }
460 if !want_compress {
461 let mut csize_buf = Vec::with_capacity(10);
464 super::varint::encode_varint(data.len() as u64, &mut csize_buf);
465 pw.write_all(&csize_buf).map_err(StoreError::from)?;
466 self.pack_position = self
467 .pack_position
468 .checked_add(csize_buf.len() as u64)
469 .ok_or_else(|| {
470 StoreError::InvalidObject("streaming pack position overflow".to_string())
471 })?;
472 pw.write_all(data).map_err(StoreError::from)?;
473 self.pack_position = self
474 .pack_position
475 .checked_add(data.len() as u64)
476 .ok_or_else(|| {
477 StoreError::InvalidObject("streaming pack position overflow".to_string())
478 })?;
479 self.total_compressed += data.len() as u64;
480 } else {
481 #[cfg(feature = "zstd")]
482 {
483 pw.write_all(&[0u8; CSIZE_PLACEHOLDER_LEN])
486 .map_err(StoreError::from)?;
487 self.pack_position = self
488 .pack_position
489 .checked_add(CSIZE_PLACEHOLDER_LEN as u64)
490 .ok_or_else(|| {
491 StoreError::InvalidObject("streaming pack position overflow".to_string())
492 })?;
493 let body_start = self.pack_position;
494 let compressed_size;
495 {
496 let mut counting = CountingWriter::new(&mut *pw);
497 let mut enc =
498 zstd::stream::write::Encoder::new(&mut counting, self.compression.level)
499 .map_err(StoreError::from)?;
500 enc.set_pledged_src_size(Some(data.len() as u64))
505 .map_err(StoreError::from)?;
506 enc.write_all(data).map_err(StoreError::from)?;
507 enc.finish().map_err(StoreError::from)?;
508 compressed_size = counting.written;
509 }
510 self.pack_position =
511 self.pack_position
512 .checked_add(compressed_size)
513 .ok_or_else(|| {
514 StoreError::InvalidObject(
515 "streaming pack position overflow".to_string(),
516 )
517 })?;
518 let body_end = body_start.checked_add(compressed_size).ok_or_else(|| {
519 StoreError::InvalidObject("streaming pack position overflow".to_string())
520 })?;
521 self.total_compressed += compressed_size;
522
523 let mut csize_bytes = [0u8; CSIZE_PLACEHOLDER_LEN];
528 encode_varint_padded_to_10(compressed_size, &mut csize_bytes);
529 pw.flush().map_err(StoreError::from)?;
530 let inner = pw.get_mut();
531 inner
532 .seek(SeekFrom::Start(csize_pos))
533 .map_err(StoreError::from)?;
534 inner.write_all(&csize_bytes).map_err(StoreError::from)?;
535 if compressed_size == data.len() as u64 {
540 inner
541 .seek(SeekFrom::Start(body_start))
542 .map_err(StoreError::from)?;
543 inner.write_all(data).map_err(StoreError::from)?;
544 }
545 inner
546 .seek(SeekFrom::Start(body_end))
547 .map_err(StoreError::from)?;
548 }
549 #[cfg(not(feature = "zstd"))]
550 {
551 unreachable!("compression branch reached without `zstd` feature");
554 }
555 }
556
557 self.add_index_entry(id, offset)?;
558 self.record_count += 1;
559 self.object_count += 1;
560 Ok(())
561 }
562
563 pub fn add_shared_frame(
570 &mut self,
571 ids: &[PackObjectId],
572 obj_type: ObjectType,
573 uncompressed_size: usize,
574 stored_data: &[u8],
575 ) -> Result<()> {
576 if ids.is_empty() {
577 return Err(StoreError::InvalidObject(
578 "compact frame must contain at least one object".to_string(),
579 ));
580 }
581 if !matches!(
582 obj_type,
583 ObjectType::Blob | ObjectType::Tree | ObjectType::State
584 ) {
585 return Err(StoreError::InvalidObject(
586 "shared compact frames may contain only blobs, trees, or states".to_string(),
587 ));
588 }
589 let unique = ids.iter().copied().collect::<HashSet<_>>();
590 if unique.len() != ids.len() {
591 return Err(StoreError::InvalidObject(
592 "compact frame contains duplicate object ids".to_string(),
593 ));
594 }
595 if ids.iter().any(|id| {
596 !matches!(
597 (obj_type, id),
598 (ObjectType::Blob | ObjectType::Tree, PackObjectId::Hash(_))
599 | (ObjectType::State, PackObjectId::StateId(_))
600 )
601 }) {
602 return Err(StoreError::InvalidObject(
603 "compact frame id kind does not match its object type".to_string(),
604 ));
605 }
606
607 let entry_start = self.pack_position;
608 let offset = entry_start
609 .checked_sub(self.header_offset)
610 .expect("header offset should precede compact frame");
611 let mut header = Vec::with_capacity(48);
612 ids[0].encode_tagged(&mut header);
613 super::varint::encode_type_and_size(obj_type, uncompressed_size as u64, &mut header);
614 super::varint::encode_varint(stored_data.len() as u64, &mut header);
615 let writer = self
616 .pack_writer
617 .as_mut()
618 .expect("add_shared_frame called after finalize");
619 writer.write_all(&header).map_err(StoreError::from)?;
620 writer.write_all(stored_data).map_err(StoreError::from)?;
621 self.pack_position = self
622 .pack_position
623 .checked_add((header.len() + stored_data.len()) as u64)
624 .ok_or_else(|| {
625 StoreError::InvalidObject("streaming pack position overflow".to_string())
626 })?;
627 self.total_uncompressed = self
628 .total_uncompressed
629 .saturating_add(uncompressed_size as u64);
630 self.total_compressed = self
631 .total_compressed
632 .saturating_add(stored_data.len() as u64);
633 self.record_count = self
634 .record_count
635 .checked_add(1)
636 .ok_or_else(|| StoreError::InvalidObject("pack record count overflow".to_string()))?;
637 for id in ids {
638 self.add_index_entry(*id, offset)?;
639 }
640 self.object_count = self
641 .object_count
642 .checked_add(ids.len() as u64)
643 .ok_or_else(|| StoreError::InvalidObject("pack object count overflow".to_string()))?;
644 Ok(())
645 }
646
647 fn add_index_entry(&mut self, id: PackObjectId, offset: u64) -> Result<()> {
648 let bucket_idx = bucket_index_for(&id);
649 let bucket = self.get_or_open_bucket(bucket_idx)?;
650 let mut entry = Vec::with_capacity(33 + 8);
651 id.encode_tagged(&mut entry);
652 entry.extend_from_slice(&offset.to_be_bytes());
653 bucket.write_all(&entry).map_err(StoreError::from)
654 }
655
656 fn get_or_open_bucket(&mut self, idx: usize) -> Result<&mut BufWriter<File>> {
657 self.bucket_access_tick = self.bucket_access_tick.wrapping_add(1);
658 let last_used = self.bucket_access_tick;
659 if self.bucket_writers[idx].is_none() {
660 if self.open_bucket_writers >= MAX_OPEN_BUCKET_WRITERS {
661 self.evict_lru_bucket()?;
662 }
663 let path = &self.bucket_paths[idx];
664 let f = OpenOptions::new()
665 .create(true)
666 .append(true)
667 .open(path)
668 .map_err(StoreError::from)?;
669 self.bucket_writers[idx] = Some(BucketWriter {
670 writer: BufWriter::new(f),
671 last_used,
672 });
673 self.open_bucket_writers += 1;
674 } else if let Some(bucket) = self.bucket_writers[idx].as_mut() {
675 bucket.last_used = last_used;
676 }
677 Ok(&mut self.bucket_writers[idx]
678 .as_mut()
679 .ok_or_else(|| StoreError::InvalidObject("index bucket writer unavailable".into()))?
680 .writer)
681 }
682
683 fn evict_lru_bucket(&mut self) -> Result<()> {
684 let Some((idx, _)) = self
685 .bucket_writers
686 .iter()
687 .enumerate()
688 .filter_map(|(idx, bucket)| bucket.as_ref().map(|bucket| (idx, bucket.last_used)))
689 .min_by_key(|(_, last_used)| *last_used)
690 else {
691 return Ok(());
692 };
693
694 if let Some(mut bucket) = self.bucket_writers[idx].take() {
695 bucket.writer.flush().map_err(StoreError::from)?;
696 self.open_bucket_writers -= 1;
697 }
698 Ok(())
699 }
700
701 pub fn finalize(mut self) -> Result<(W, PackStats)> {
712 for bucket in self.bucket_writers.iter_mut().flatten() {
715 bucket.writer.flush().map_err(StoreError::from)?;
716 }
717 for slot in self.bucket_writers.iter_mut() {
720 *slot = None;
721 }
722 self.open_bucket_writers = 0;
723
724 let bw = self
728 .pack_writer
729 .take()
730 .ok_or_else(|| StoreError::InvalidObject("pack builder is finalized".into()))?;
731 let mut writer = bw
732 .into_inner()
733 .map_err(|e| StoreError::from(std::io::Error::other(e.to_string())))?;
734 if let Some(expected) = self.declared_object_count {
735 if expected != self.record_count {
736 return Err(StoreError::InvalidObject(format!(
737 "streaming pack declared {expected} record(s) but added {}",
738 self.record_count
739 )));
740 }
741 } else {
742 writer
743 .seek(SeekFrom::Start(self.header_offset))
744 .map_err(StoreError::from)?;
745 let mut header_bytes = Vec::with_capacity(16);
746 write_container_header(&mut header_bytes, pack_container_spec(), self.record_count);
747 writer.write_all(&header_bytes).map_err(StoreError::from)?;
748 }
749
750 writer
755 .seek(SeekFrom::Start(self.header_offset))
756 .map_err(StoreError::from)?;
757 let mut hasher = blake3::Hasher::new();
758 let mut buf = vec![0u8; 64 * 1024];
759 loop {
760 let n = writer.read(&mut buf).map_err(StoreError::from)?;
761 if n == 0 {
762 break;
763 }
764 hasher.update(&buf[..n]);
765 }
766 let checksum = hasher.finalize();
767
768 writer.seek(SeekFrom::End(0)).map_err(StoreError::from)?;
770 writer
771 .write_all(checksum.as_bytes())
772 .map_err(StoreError::from)?;
773 writer.flush().map_err(StoreError::from)?;
774 if self.durable {
776 writer
777 .sync_data_for_durability()
778 .map_err(StoreError::from)?;
779 }
780
781 let idx_file = File::create(&self.index_path).map_err(StoreError::from)?;
784 let mut idx_writer = BufWriter::new(idx_file);
785 write_index_header(&mut idx_writer, self.object_count)?;
786 let mut entries_written: u64 = 0;
787 for path in self.bucket_paths.iter() {
788 if !path.exists() {
789 continue;
790 }
791 let run = super::disk_sort::sort_records::<41>(File::open(path)?, &self.bucket_dir)?;
792 let mut input = std::io::BufReader::new(run.as_file());
793 while let Some(record) = super::disk_sort::read_record::<41>(&mut input)? {
794 let (id, _) = PackObjectId::decode_tagged(&record)?;
795 let offset = u64::from_be_bytes(record[33..].try_into().map_err(|_| {
796 StoreError::InvalidObject("index bucket offset truncated".into())
797 })?);
798 write_index_entry(&mut idx_writer, id, offset)?;
799 entries_written += 1;
800 }
801 }
802 idx_writer.flush().map_err(StoreError::from)?;
803 let _idx_file = idx_writer
805 .into_inner()
806 .map_err(|e| StoreError::from(std::io::Error::other(e.to_string())))?;
807 #[cfg(feature = "fs")]
808 if self.durable {
809 _idx_file.sync_all().map_err(StoreError::from)?;
810 if let Some(parent) = self.index_path.parent() {
811 heddle_fs_prims::fs_atomic::sync_directory(parent).map_err(StoreError::from)?;
812 }
813 }
814 debug_assert_eq!(
815 entries_written, self.object_count,
816 "streaming index entry count drifted from add() count"
817 );
818
819 for path in self.bucket_paths.iter() {
824 let _ = std::fs::remove_file(path);
825 }
826 self.scratch_lease.take();
827 let _ = std::fs::remove_file(self.bucket_dir.join(super::scratch::LEASE_FILE));
828 let _ = std::fs::remove_dir(&self.bucket_dir);
829 self.finalized = true;
830
831 let stats = PackStats {
832 object_count: self.object_count,
833 total_uncompressed: self.total_uncompressed,
834 total_compressed: self.total_compressed,
835 delta_count: 0,
836 compression_ratio: if self.total_uncompressed == 0 {
837 0.0
838 } else {
839 self.total_compressed as f64 / self.total_uncompressed as f64
840 },
841 };
842
843 Ok((writer, stats))
844 }
845}
846
847fn write_index_header<W: Write>(out: &mut W, count: u64) -> Result<()> {
852 super::pack_index::index_header().write_to(out, count)
853}
854
855fn write_index_entry<W: Write>(out: &mut W, id: PackObjectId, offset: u64) -> Result<()> {
858 let buf = super::pack_index::encode_index_entry(id, offset);
859 out.write_all(&buf).map_err(StoreError::from)
860}
861
862#[cfg(feature = "zstd")]
875fn encode_varint_padded_to_10(value: u64, out: &mut [u8; 10]) {
876 let mut v = value;
877 for slot in out.iter_mut().take(9) {
878 *slot = 0x80 | ((v & 0x7F) as u8);
879 v >>= 7;
880 }
881 out[9] = (v & 0x7F) as u8;
882}
883
884impl<W: Write + Read + Seek> Drop for StreamingPackBuilder<W> {
885 fn drop(&mut self) {
886 if self.finalized {
887 return;
888 }
889 for path in self.bucket_paths.iter() {
892 let _ = std::fs::remove_file(path);
893 }
894 self.scratch_lease.take();
895 let _ = std::fs::remove_file(self.bucket_dir.join(super::scratch::LEASE_FILE));
896 let _ = std::fs::remove_dir(&self.bucket_dir);
897 }
898}
899
900fn bucket_index_for(id: &PackObjectId) -> usize {
904 match id {
905 PackObjectId::Hash(h) => HASH_VARIANT * BUCKETS_PER_VARIANT + h.as_bytes()[0] as usize,
906 PackObjectId::StateId(c) => {
907 CHANGEID_VARIANT * BUCKETS_PER_VARIANT + c.as_bytes()[0] as usize
908 }
909 PackObjectId::AnnotatedTag(hash) => {
910 ANNOTATED_TAG_VARIANT * BUCKETS_PER_VARIANT + hash.as_bytes()[0] as usize
911 }
912 }
913}
914
915#[cfg(test)]
918mod tests {
919 use std::io::Cursor;
920
921 use super::*;
922 use crate::{
923 object::StateId,
924 store::pack::{PackReader, PackStats},
925 };
926
927 fn deterministic_hash(seed: u8) -> ContentHash {
928 let mut bytes = [0u8; 32];
932 bytes[0] = seed;
933 for (i, b) in bytes.iter_mut().enumerate().skip(1) {
934 *b = seed.wrapping_mul(31).wrapping_add(i as u8);
935 }
936 ContentHash::from_bytes(bytes)
937 }
938
939 fn deterministic_state_id(seed: u8) -> StateId {
940 let mut bytes = [0u8; 32];
941 bytes[0] = seed;
942 for (i, b) in bytes.iter_mut().enumerate().skip(1) {
943 *b = seed.wrapping_add(i as u8 * 7);
944 }
945 StateId::from_bytes(bytes)
946 }
947
948 fn fresh_builder(
953 tmp: &tempfile::TempDir,
954 ) -> (StreamingPackBuilder<Cursor<Vec<u8>>>, PathBuf, PathBuf) {
955 let bucket_dir = tmp.path().join("buckets");
956 let index_path = tmp.path().join("test.idx");
957 let cursor = Cursor::new(Vec::<u8>::new());
958 let b = StreamingPackBuilder::new(
959 cursor,
960 index_path.clone(),
961 CompressionConfig::default(),
962 bucket_dir.clone(),
963 )
964 .unwrap();
965 (b, bucket_dir, index_path)
966 }
967
968 fn finalize_cursor(
973 b: StreamingPackBuilder<Cursor<Vec<u8>>>,
974 index_path: &std::path::Path,
975 ) -> (Vec<u8>, Vec<u8>, PackStats) {
976 let (cursor, stats) = b.finalize().unwrap();
977 let index_bytes = std::fs::read(index_path).unwrap();
978 (cursor.into_inner(), index_bytes, stats)
979 }
980
981 #[test]
982 fn empty_pack_finalizes_to_valid_zero_count_pack() {
983 let tmp = tempfile::TempDir::new().unwrap();
984 let (b, bucket_dir, idx_path) = fresh_builder(&tmp);
985 let (pack_data, index_data, stats) = finalize_cursor(b, &idx_path);
986
987 assert_eq!(stats.object_count, 0);
988 let reader = PackReader::from_bytes(pack_data, index_data, &std::env::temp_dir()).unwrap();
990 assert!(reader.list_ids().unwrap().is_empty());
991 assert!(
993 !bucket_dir.exists(),
994 "bucket dir should be cleaned on successful finalize"
995 );
996 }
997
998 #[test]
999 fn single_blob_with_hash_id_round_trips() {
1000 let tmp = tempfile::TempDir::new().unwrap();
1001 let (mut b, _, idx_path) = fresh_builder(&tmp);
1002 let hash = deterministic_hash(0x42);
1003 let payload = b"hello, streaming pack".to_vec();
1004 b.add(hash, ObjectType::Blob, payload.clone()).unwrap();
1005 let (pack_data, index_data, stats) = finalize_cursor(b, &idx_path);
1006
1007 assert_eq!(stats.object_count, 1);
1008 let reader = PackReader::from_bytes(pack_data, index_data, &std::env::temp_dir()).unwrap();
1009 let id = PackObjectId::Hash(hash);
1010 assert!(reader.has_object(&id).unwrap());
1011 let (got_type, got_data) = reader.get_object(&id).unwrap().unwrap();
1012 assert_eq!(got_type, ObjectType::Blob);
1013 assert_eq!(got_data, payload);
1014 }
1015
1016 #[test]
1017 fn single_state_with_change_id_round_trips() {
1018 let tmp = tempfile::TempDir::new().unwrap();
1019 let (mut b, _, idx_path) = fresh_builder(&tmp);
1020 let cid = deterministic_state_id(0xa5);
1021 let payload = b"serialized-state-bytes".to_vec();
1022 b.add_id(
1023 PackObjectId::StateId(cid),
1024 ObjectType::State,
1025 payload.clone(),
1026 )
1027 .unwrap();
1028 let (pack_data, index_data, stats) = finalize_cursor(b, &idx_path);
1029
1030 assert_eq!(stats.object_count, 1);
1031 let reader = PackReader::from_bytes(pack_data, index_data, &std::env::temp_dir()).unwrap();
1032 let id = PackObjectId::StateId(cid);
1033 let (ty, data) = reader.get_object(&id).unwrap().unwrap();
1034 assert_eq!(ty, ObjectType::State);
1035 assert_eq!(data, payload);
1036 }
1037
1038 #[test]
1039 fn shared_compact_tree_frame_reconstructs_each_indexed_object() {
1040 use crate::object::{Tree, TreeEntry};
1041
1042 let tmp = tempfile::TempDir::new().unwrap();
1043 let (mut builder, _, index_path) = fresh_builder(&tmp);
1044 let blob = deterministic_hash(0x33);
1045 let trees = vec![
1046 Tree::from_entries(vec![TreeEntry::file("a", blob, false).unwrap()]),
1047 Tree::from_entries(vec![TreeEntry::file("b", blob, true).unwrap()]),
1048 ];
1049 let ids = trees
1050 .iter()
1051 .map(|tree| PackObjectId::Hash(tree.hash()))
1052 .collect::<Vec<_>>();
1053 let frame = heddle_object_model::compact::encode_tree_frame(&trees).unwrap();
1054 builder
1055 .add_shared_frame(&ids, ObjectType::Tree, frame.len(), &frame)
1056 .unwrap();
1057 let (pack, index, stats) = finalize_cursor(builder, &index_path);
1058 let reader = PackReader::from_bytes(pack, index, &std::env::temp_dir()).unwrap();
1059
1060 assert_eq!(stats.object_count, 2);
1061 assert_eq!(
1062 reader.encoded_payload_bytes(ObjectType::Tree).unwrap(),
1063 frame.len() as u64
1064 );
1065 let hosted = ids
1066 .iter()
1067 .zip(&trees)
1068 .map(|(id, tree)| {
1069 (
1070 *id,
1071 ObjectType::Tree,
1072 tree.encode_canonical().unwrap().len() as u64,
1073 )
1074 })
1075 .collect::<Vec<_>>();
1076 assert!(
1077 reader
1078 .copy_hosted_encoded_subset(&hosted)
1079 .unwrap()
1080 .is_none(),
1081 "repository-local compact frames must use the hosted fallback"
1082 );
1083 for (id, tree) in ids.iter().zip(&trees) {
1084 let (object_type, bytes) = reader.get_object(id).unwrap().unwrap();
1085 assert_eq!(object_type, ObjectType::Tree);
1086 assert_eq!(bytes, tree.encode_canonical().unwrap());
1087 }
1088 }
1089
1090 #[test]
1091 fn compact_tree_extraction_rejects_an_index_alias_with_the_wrong_typed_hash() {
1092 use crate::object::{Tree, TreeEntry};
1093
1094 let tmp = tempfile::TempDir::new().unwrap();
1095 let (mut builder, _, index_path) = fresh_builder(&tmp);
1096 let blob = deterministic_hash(0x34);
1097 let trees = vec![
1098 Tree::from_entries(vec![TreeEntry::file("a", blob, false).unwrap()]),
1099 Tree::from_entries(vec![TreeEntry::file("b", blob, true).unwrap()]),
1100 ];
1101 let wrong_hash = ContentHash::compute_typed("tree", b"not the second tree");
1102 assert_ne!(wrong_hash, trees[1].hash());
1103 let ids = vec![
1104 PackObjectId::Hash(trees[0].hash()),
1105 PackObjectId::Hash(wrong_hash),
1106 ];
1107 let frame = heddle_object_model::compact::encode_tree_frame(&trees).unwrap();
1108 builder
1109 .add_shared_frame(&ids, ObjectType::Tree, frame.len(), &frame)
1110 .unwrap();
1111 let (pack, index, _) = finalize_cursor(builder, &index_path);
1112 let reader = PackReader::from_bytes(pack, index, &std::env::temp_dir()).unwrap();
1113
1114 let error = reader.get_object(&ids[1]).unwrap_err();
1115 assert!(
1116 error
1117 .to_string()
1118 .contains("does not contain indexed object"),
1119 "extraction must derive and verify the tree's typed hash: {error}"
1120 );
1121 }
1122
1123 #[test]
1124 fn shared_lineage_blob_frame_reconstructs_each_indexed_object() {
1125 let tmp = tempfile::TempDir::new().unwrap();
1126 let (mut builder, _, index_path) = fresh_builder(&tmp);
1127 let bodies = [b"newest version".as_slice(), b"older version".as_slice()];
1128 let ids = bodies
1129 .iter()
1130 .map(|body| PackObjectId::Hash(ContentHash::compute_typed("blob", body)))
1131 .collect::<Vec<_>>();
1132 let frame = heddle_object_model::compact::encode_blob_frame(&bodies).unwrap();
1133 builder
1134 .add_shared_frame(&ids, ObjectType::Blob, frame.len(), &frame)
1135 .unwrap();
1136 let (pack, index, stats) = finalize_cursor(builder, &index_path);
1137 let reader = PackReader::from_bytes(pack, index, &std::env::temp_dir()).unwrap();
1138
1139 assert_eq!(stats.object_count, bodies.len() as u64);
1140 assert_eq!(
1141 reader.encoded_payload_bytes(ObjectType::Blob).unwrap(),
1142 frame.len() as u64
1143 );
1144 for (id, expected) in ids.iter().zip(bodies) {
1145 let (object_type, actual) = reader.get_object(id).unwrap().unwrap();
1146 assert_eq!(object_type, ObjectType::Blob);
1147 assert_eq!(actual, expected);
1148 let PackObjectId::Hash(hash) = id else {
1149 unreachable!("blob ids are hashes")
1150 };
1151 assert_eq!(
1152 reader.get_hashed_object_type(hash).unwrap(),
1153 Some(ObjectType::Blob)
1154 );
1155 assert_eq!(
1156 reader.get_hashed_object_size(hash).unwrap(),
1157 Some(expected.len() as u64)
1158 );
1159 }
1160 }
1161
1162 #[test]
1163 fn ordinary_blob_starting_with_frame_magic_remains_ordinary() {
1164 let tmp = tempfile::TempDir::new().unwrap();
1165 let (mut builder, _, index_path) = fresh_builder(&tmp);
1166 let body = b"HCB2 arbitrary user content".to_vec();
1167 let hash = ContentHash::compute_typed("blob", &body);
1168 builder.add(hash, ObjectType::Blob, body.clone()).unwrap();
1169 let (pack, index, _) = finalize_cursor(builder, &index_path);
1170 let reader = PackReader::from_bytes(pack, index, &std::env::temp_dir()).unwrap();
1171
1172 assert_eq!(
1173 reader.get_object(&PackObjectId::Hash(hash)).unwrap(),
1174 Some((ObjectType::Blob, body))
1175 );
1176 }
1177
1178 #[test]
1179 fn corrupt_compact_frame_byte_invalidates_every_contained_object() {
1180 use crate::object::{Tree, TreeEntry};
1181
1182 let tmp = tempfile::TempDir::new().unwrap();
1183 let (mut builder, _, index_path) = fresh_builder(&tmp);
1184 let blob = deterministic_hash(0x44);
1185 let trees = vec![
1186 Tree::from_entries(vec![TreeEntry::file("a", blob, false).unwrap()]),
1187 Tree::from_entries(vec![TreeEntry::file("b", blob, true).unwrap()]),
1188 ];
1189 let ids = trees
1190 .iter()
1191 .map(|tree| PackObjectId::Hash(tree.hash()))
1192 .collect::<Vec<_>>();
1193 let mut frame = heddle_object_model::compact::encode_tree_frame(&trees).unwrap();
1194 let corrupt_at = frame.len() / 2;
1195 frame[corrupt_at] ^= 0x01;
1196 builder
1197 .add_shared_frame(&ids, ObjectType::Tree, frame.len(), &frame)
1198 .unwrap();
1199 let (pack, index, _) = finalize_cursor(builder, &index_path);
1200 let reader = PackReader::from_bytes(pack, index, &std::env::temp_dir()).unwrap();
1201
1202 for id in ids {
1203 let error = reader.get_object(&id).unwrap_err();
1204 assert!(
1205 error
1206 .to_string()
1207 .contains("compact frame checksum mismatch"),
1208 "unexpected error for {id:?}: {error}"
1209 );
1210 }
1211 }
1212
1213 #[test]
1214 fn corrupt_blob_frame_byte_invalidates_every_contained_object() {
1215 let tmp = tempfile::TempDir::new().unwrap();
1216 let (mut builder, _, index_path) = fresh_builder(&tmp);
1217 let bodies = [b"newest version".as_slice(), b"older version".as_slice()];
1218 let ids = bodies
1219 .iter()
1220 .map(|body| PackObjectId::Hash(ContentHash::compute_typed("blob", body)))
1221 .collect::<Vec<_>>();
1222 let mut frame = heddle_object_model::compact::encode_blob_frame(&bodies).unwrap();
1223 let corrupt_at = frame.len() / 2;
1224 frame[corrupt_at] ^= 0x01;
1225 builder
1226 .add_shared_frame(&ids, ObjectType::Blob, frame.len(), &frame)
1227 .unwrap();
1228 let (pack, index, _) = finalize_cursor(builder, &index_path);
1229 let reader = PackReader::from_bytes(pack, index, &std::env::temp_dir()).unwrap();
1230
1231 for id in ids {
1232 let error = reader.get_object(&id).unwrap_err();
1233 assert!(
1234 error
1235 .to_string()
1236 .contains("compact frame checksum mismatch"),
1237 "unexpected error for {id:?}: {error}"
1238 );
1239 }
1240 }
1241
1242 #[test]
1243 fn mixed_hash_and_changeid_ids_all_retrievable() {
1244 let tmp = tempfile::TempDir::new().unwrap();
1245 let (mut b, _, idx_path) = fresh_builder(&tmp);
1246 let blob_hash = deterministic_hash(0x10);
1247 let tree_hash = deterministic_hash(0x20);
1248 let state_cid = deterministic_state_id(0x80);
1249
1250 b.add(blob_hash, ObjectType::Blob, b"blob-bytes".to_vec())
1251 .unwrap();
1252 b.add(tree_hash, ObjectType::Tree, b"serialized-tree".to_vec())
1253 .unwrap();
1254 b.add_id(
1255 PackObjectId::StateId(state_cid),
1256 ObjectType::State,
1257 b"serialized-state",
1258 )
1259 .unwrap();
1260
1261 let (pack_data, index_data, stats) = finalize_cursor(b, &idx_path);
1262 assert_eq!(stats.object_count, 3);
1263 let reader = PackReader::from_bytes(pack_data, index_data, &std::env::temp_dir()).unwrap();
1264 assert_eq!(
1265 reader
1266 .get_object(&PackObjectId::Hash(blob_hash))
1267 .unwrap()
1268 .unwrap()
1269 .1,
1270 b"blob-bytes".to_vec()
1271 );
1272 assert_eq!(
1273 reader
1274 .get_object(&PackObjectId::Hash(tree_hash))
1275 .unwrap()
1276 .unwrap()
1277 .1,
1278 b"serialized-tree".to_vec()
1279 );
1280 assert_eq!(
1281 reader
1282 .get_object(&PackObjectId::StateId(state_cid))
1283 .unwrap()
1284 .unwrap()
1285 .1,
1286 b"serialized-state".to_vec()
1287 );
1288 }
1289
1290 #[test]
1291 fn ten_thousand_objects_round_trip_correctly() {
1292 let tmp = tempfile::TempDir::new().unwrap();
1296 let (mut b, _, idx_path) = fresh_builder(&tmp);
1297 let mut hashes = Vec::with_capacity(10_000);
1298 for i in 0..10_000u32 {
1299 let h = blake3::hash(&i.to_le_bytes());
1302 let hash = ContentHash::from_bytes(*h.as_bytes());
1303 hashes.push(hash);
1304 b.add(hash, ObjectType::Blob, format!("payload-{i}").into_bytes())
1305 .unwrap();
1306 }
1307 let (pack_data, index_data, stats) = finalize_cursor(b, &idx_path);
1308 assert_eq!(stats.object_count, 10_000);
1309
1310 let reader = PackReader::from_bytes(pack_data, index_data, &std::env::temp_dir()).unwrap();
1311 assert_eq!(reader.list_ids().unwrap().len(), 10_000);
1312 for i in [0, 1, 99, 1234, 5_000, 9_999] {
1314 let id = PackObjectId::Hash(hashes[i]);
1315 let (_ty, data) = reader.get_object(&id).unwrap().unwrap();
1316 assert_eq!(data, format!("payload-{i}").into_bytes());
1317 }
1318 }
1319
1320 #[test]
1321 fn bucket_writers_are_lru_capped_below_fd_limit() {
1322 let tmp = tempfile::TempDir::new().unwrap();
1323 let (mut b, _bucket_dir, idx_path) = fresh_builder(&tmp);
1324 let mut ids = Vec::new();
1325
1326 for i in 0..BUCKETS_PER_VARIANT {
1327 let hash = deterministic_hash(i as u8);
1328 ids.push(PackObjectId::Hash(hash));
1329 b.add(hash, ObjectType::Blob, format!("hash-{i}").into_bytes())
1330 .unwrap();
1331 assert!(
1332 b.open_bucket_writers <= MAX_OPEN_BUCKET_WRITERS,
1333 "open bucket writers should stay capped"
1334 );
1335 }
1336
1337 for i in 0..BUCKETS_PER_VARIANT {
1338 let cid = deterministic_state_id(i as u8);
1339 ids.push(PackObjectId::StateId(cid));
1340 b.add_id(
1341 PackObjectId::StateId(cid),
1342 ObjectType::State,
1343 format!("state-{i}").into_bytes(),
1344 )
1345 .unwrap();
1346 assert!(
1347 b.open_bucket_writers <= MAX_OPEN_BUCKET_WRITERS,
1348 "open bucket writers should stay capped"
1349 );
1350 }
1351
1352 for i in 0..BUCKETS_PER_VARIANT {
1353 let hash = deterministic_hash(i as u8);
1354 let id = PackObjectId::AnnotatedTag(hash);
1355 ids.push(id);
1356 b.add_id(
1357 id,
1358 ObjectType::AnnotatedTag,
1359 format!("tag-{i}").into_bytes(),
1360 )
1361 .unwrap();
1362 assert!(
1363 b.open_bucket_writers <= MAX_OPEN_BUCKET_WRITERS,
1364 "open bucket writers should stay capped"
1365 );
1366 }
1367
1368 let (pack_data, index_data, stats) = finalize_cursor(b, &idx_path);
1369 assert_eq!(stats.object_count, TOTAL_BUCKETS as u64);
1370 let reader = PackReader::from_bytes(pack_data, index_data, &std::env::temp_dir()).unwrap();
1371 for id in ids {
1372 assert!(reader.has_object(&id).unwrap(), "missing id {id:?}");
1373 }
1374 }
1375
1376 #[test]
1377 fn index_id_sort_order_matches_packbuilder_output() {
1378 use crate::store::pack::PackBuilder;
1386 let payloads: Vec<(PackObjectId, ObjectType, Vec<u8>)> = (0..200u32)
1387 .map(|i| {
1388 let h = blake3::hash(&i.to_le_bytes());
1389 (
1390 PackObjectId::Hash(ContentHash::from_bytes(*h.as_bytes())),
1391 if i % 3 == 0 {
1392 ObjectType::Tree
1393 } else {
1394 ObjectType::Blob
1395 },
1396 format!("body-{i}").into_bytes(),
1397 )
1398 })
1399 .collect();
1400
1401 let compression = CompressionConfig {
1404 max_delta_size: 0,
1405 ..CompressionConfig::default()
1406 };
1407 let mut classic = PackBuilder::new(compression);
1408 for (id, ty, data) in payloads.iter() {
1409 classic.add_id(*id, *ty, data.clone());
1410 }
1411 let (classic_pack, classic_index, _) = classic.build().unwrap();
1412 let classic_reader =
1413 PackReader::from_bytes(classic_pack, classic_index, &std::env::temp_dir()).unwrap();
1414
1415 let tmp = tempfile::TempDir::new().unwrap();
1416 let bucket_dir = tmp.path().join("buckets");
1417 let idx_path = tmp.path().join("test.idx");
1418 let cursor = Cursor::new(Vec::<u8>::new());
1419 let mut streaming =
1420 StreamingPackBuilder::new(cursor, idx_path.clone(), compression, bucket_dir).unwrap();
1421 for (id, ty, data) in payloads.iter() {
1422 streaming.add_id(*id, *ty, data.clone()).unwrap();
1423 }
1424 let (streaming_pack, streaming_index, _) = finalize_cursor(streaming, &idx_path);
1425 let streaming_reader =
1426 PackReader::from_bytes(streaming_pack, streaming_index, &std::env::temp_dir()).unwrap();
1427
1428 assert_eq!(
1431 streaming_reader.list_ids().unwrap(),
1432 classic_reader.list_ids().unwrap(),
1433 "streaming and classic indices should report the same id sequence"
1434 );
1435 for (id, _ty, want) in payloads.iter().take(10).chain(payloads.iter().skip(190)) {
1438 let (_, got) = streaming_reader.get_object(id).unwrap().unwrap();
1439 assert_eq!(&got, want);
1440 let (_, classic_got) = classic_reader.get_object(id).unwrap().unwrap();
1441 assert_eq!(got, classic_got);
1442 }
1443 }
1444
1445 #[test]
1446 fn corrupted_pack_fails_checksum_verification() {
1447 let tmp = tempfile::TempDir::new().unwrap();
1448 let (mut b, _, idx_path) = fresh_builder(&tmp);
1449 b.add(
1450 deterministic_hash(0x01),
1451 ObjectType::Blob,
1452 b"some bytes".to_vec(),
1453 )
1454 .unwrap();
1455 let (mut pack_data, index_data, _) = finalize_cursor(b, &idx_path);
1456 let body_byte = 18; pack_data[body_byte] ^= 0xff;
1459 let result = PackReader::from_bytes(pack_data, index_data, &std::env::temp_dir());
1460 assert!(
1461 result.is_err(),
1462 "PackReader should reject pack with mutated body"
1463 );
1464 }
1465
1466 #[test]
1467 fn pack_count_in_header_matches_index_entry_count() {
1468 let tmp = tempfile::TempDir::new().unwrap();
1469 let (mut b, _, idx_path) = fresh_builder(&tmp);
1470 for i in 0..7u8 {
1471 b.add(
1472 deterministic_hash(i),
1473 ObjectType::Blob,
1474 format!("p{i}").into_bytes(),
1475 )
1476 .unwrap();
1477 }
1478 let (pack_data, index_data, _) = finalize_cursor(b, &idx_path);
1479 let count = u64::from_be_bytes(pack_data[8..16].try_into().unwrap());
1481 assert_eq!(count, 7);
1482 let reader = PackReader::from_bytes(pack_data, index_data, &std::env::temp_dir()).unwrap();
1483 assert_eq!(reader.list_ids().unwrap().len(), 7);
1484 }
1485
1486 #[test]
1487 fn declared_pack_count_is_written_before_finalize() {
1488 let tmp = tempfile::TempDir::new().unwrap();
1489 let bucket_dir = tmp.path().join("buckets");
1490 let idx_path = tmp.path().join("test.idx");
1491 let cursor = Cursor::new(Vec::<u8>::new());
1492 let mut b = StreamingPackBuilder::new_with_object_count(
1493 cursor,
1494 idx_path.clone(),
1495 CompressionConfig::default(),
1496 bucket_dir,
1497 2,
1498 )
1499 .unwrap();
1500
1501 b.flush_pack().unwrap();
1502 let initial = b.pack_writer.as_ref().unwrap().get_ref().get_ref().clone();
1503 assert_eq!(u64::from_be_bytes(initial[8..16].try_into().unwrap()), 2);
1504
1505 let hash = deterministic_hash(0x40);
1506 b.add(hash, ObjectType::Blob, b"known-count-entry".to_vec())
1507 .unwrap();
1508 b.flush_pack().unwrap();
1509 let after_add = b.pack_writer.as_ref().unwrap().get_ref().get_ref().clone();
1510 assert_eq!(u64::from_be_bytes(after_add[8..16].try_into().unwrap()), 2);
1511
1512 let second_hash = deterministic_hash(0x41);
1513 b.add(second_hash, ObjectType::Blob, b"second-entry".to_vec())
1514 .unwrap();
1515 let (pack_data, index_data, stats) = finalize_cursor(b, &idx_path);
1516
1517 assert_eq!(stats.object_count, 2);
1518 assert_eq!(u64::from_be_bytes(pack_data[8..16].try_into().unwrap()), 2);
1519 let reader = PackReader::from_bytes(pack_data, index_data, &std::env::temp_dir()).unwrap();
1520 assert!(reader.has_object(&PackObjectId::Hash(hash)).unwrap());
1521 assert!(reader.has_object(&PackObjectId::Hash(second_hash)).unwrap());
1522 }
1523
1524 #[test]
1525 fn declared_pack_count_mismatch_fails_finalize() {
1526 let tmp = tempfile::TempDir::new().unwrap();
1527 let bucket_dir = tmp.path().join("buckets");
1528 let idx_path = tmp.path().join("test.idx");
1529 let cursor = Cursor::new(Vec::<u8>::new());
1530 let mut b = StreamingPackBuilder::new_with_object_count(
1531 cursor,
1532 idx_path,
1533 CompressionConfig::default(),
1534 bucket_dir,
1535 2,
1536 )
1537 .unwrap();
1538
1539 b.add(
1540 deterministic_hash(0x50),
1541 ObjectType::Blob,
1542 b"only-entry".to_vec(),
1543 )
1544 .unwrap();
1545 let error = b.finalize().unwrap_err();
1546
1547 assert!(
1548 error
1549 .to_string()
1550 .contains("streaming pack declared 2 record(s) but added 1")
1551 );
1552 }
1553
1554 #[test]
1555 fn bucket_files_are_cleaned_on_successful_finalize() {
1556 let tmp = tempfile::TempDir::new().unwrap();
1557 let bucket_dir = tmp.path().join("buckets");
1558 let idx_path = tmp.path().join("test.idx");
1559 let cursor = Cursor::new(Vec::<u8>::new());
1560 let mut b = StreamingPackBuilder::new(
1561 cursor,
1562 idx_path.clone(),
1563 CompressionConfig::default(),
1564 bucket_dir.clone(),
1565 )
1566 .unwrap();
1567 for i in 0..50u8 {
1568 b.add(deterministic_hash(i), ObjectType::Blob, vec![i; 32])
1569 .unwrap();
1570 }
1571 assert!(bucket_dir.exists());
1573 let bucket_count = std::fs::read_dir(&bucket_dir).unwrap().count();
1574 assert!(bucket_count > 0, "bucket dir should hold some files");
1575 let _ = finalize_cursor(b, &idx_path);
1576 assert!(
1577 !bucket_dir.exists(),
1578 "bucket dir should be removed on finalize"
1579 );
1580 }
1581
1582 #[test]
1583 fn bucket_files_are_cleaned_on_drop_without_finalize() {
1584 let tmp = tempfile::TempDir::new().unwrap();
1585 let bucket_dir = tmp.path().join("buckets");
1586 let idx_path = tmp.path().join("test.idx");
1587 {
1588 let cursor = Cursor::new(Vec::<u8>::new());
1589 let mut b = StreamingPackBuilder::new(
1590 cursor,
1591 idx_path.clone(),
1592 CompressionConfig::default(),
1593 bucket_dir.clone(),
1594 )
1595 .unwrap();
1596 for i in 0..10u8 {
1597 b.add(deterministic_hash(i), ObjectType::Blob, vec![0; 32])
1598 .unwrap();
1599 }
1600 assert!(bucket_dir.exists());
1601 }
1603 assert!(
1604 !idx_path.exists(),
1605 "no index file should have been created without finalize"
1606 );
1607 assert!(
1608 !bucket_dir.exists(),
1609 "bucket dir should be removed on Drop when finalize never ran"
1610 );
1611 }
1612
1613 #[test]
1614 fn large_blob_streams_to_disk_without_double_buffering() {
1615 let tmp = tempfile::TempDir::new().unwrap();
1620 let bucket_dir = tmp.path().join("buckets");
1621 let pack_path = tmp.path().join("pack.dat");
1622 let idx_path = tmp.path().join("pack.idx");
1623 let file = std::fs::OpenOptions::new()
1624 .read(true)
1625 .write(true)
1626 .create(true)
1627 .truncate(true)
1628 .open(&pack_path)
1629 .unwrap();
1630 let mut b = StreamingPackBuilder::new(
1631 file,
1632 idx_path.clone(),
1633 CompressionConfig::default(),
1634 bucket_dir,
1635 )
1636 .unwrap();
1637 let payload: Vec<u8> = (0..4 * 1024 * 1024u32).map(|i| (i & 0xff) as u8).collect();
1638 let hash = deterministic_hash(0xff);
1639 b.add(hash, ObjectType::Blob, payload.clone()).unwrap();
1640 let (_, stats) = b.finalize().unwrap();
1641 let index_data = std::fs::read(&idx_path).unwrap();
1642 assert_eq!(stats.object_count, 1);
1643 let pack_bytes = std::fs::read(&pack_path).unwrap();
1644 let reader = PackReader::from_bytes(pack_bytes, index_data, &std::env::temp_dir()).unwrap();
1647 let (_ty, got) = reader
1648 .get_object(&PackObjectId::Hash(hash))
1649 .unwrap()
1650 .unwrap();
1651 assert_eq!(got, payload);
1652 }
1653
1654 #[test]
1655 fn bucket_distribution_for_random_hashes_is_roughly_uniform() {
1656 let tmp = tempfile::TempDir::new().unwrap();
1663 let bucket_dir = tmp.path().join("buckets");
1664 let idx_path = tmp.path().join("test.idx");
1665 let cursor = Cursor::new(Vec::<u8>::new());
1666 let mut b = StreamingPackBuilder::new(
1667 cursor,
1668 idx_path.clone(),
1669 CompressionConfig::default(),
1670 bucket_dir.clone(),
1671 )
1672 .unwrap();
1673 for i in 0..1024u32 {
1674 let h = blake3::hash(&i.to_le_bytes());
1675 let hash = ContentHash::from_bytes(*h.as_bytes());
1676 b.add(hash, ObjectType::Blob, b"x".to_vec()).unwrap();
1677 }
1678 b.pack_writer.as_mut().unwrap().flush().unwrap();
1680 let mut max_entries = 0usize;
1681 let entry_size = 33 + 8;
1684 for path in b.bucket_paths.iter() {
1685 if path.exists() {
1686 let size = std::fs::metadata(path).unwrap().len() as usize;
1687 let entries = size / entry_size;
1688 if entries > max_entries {
1689 max_entries = entries;
1690 }
1691 }
1692 }
1693 assert!(
1696 max_entries <= 16,
1697 "max bucket has {max_entries} entries; uniform expected ~4"
1698 );
1699 let _ = finalize_cursor(b, &idx_path);
1700 }
1701
1702 #[test]
1703 fn finalize_returns_correct_stats() {
1704 let tmp = tempfile::TempDir::new().unwrap();
1705 let (mut b, _, idx_path) = fresh_builder(&tmp);
1706 let payload = vec![0xabu8; 1024];
1707 for i in 0..5u8 {
1708 b.add(deterministic_hash(i), ObjectType::Blob, payload.clone())
1709 .unwrap();
1710 }
1711 let (_, _, stats) = finalize_cursor(b, &idx_path);
1712 assert_eq!(stats.object_count, 5);
1713 assert_eq!(stats.total_uncompressed, 5 * 1024);
1714 assert!(stats.total_compressed > 0);
1715 assert!(stats.compression_ratio > 0.0);
1716 assert_eq!(stats.delta_count, 0, "streaming builder never deltas");
1717 }
1718
1719 #[cfg(feature = "zstd")]
1720 #[test]
1721 fn streaming_compression_roundtrips_through_zstd_frame() {
1722 let tmp = tempfile::TempDir::new().unwrap();
1730 let (mut b, _, idx_path) = fresh_builder(&tmp);
1731 let payload = vec![0u8; 64 * 1024];
1734 let hash = deterministic_hash(0x77);
1735 b.add(hash, ObjectType::Blob, payload.clone()).unwrap();
1736 let (pack_data, index_data, stats) = finalize_cursor(b, &idx_path);
1737 assert!(
1738 stats.total_compressed < stats.total_uncompressed,
1739 "expected compression ratio < 1.0, got {}/{}",
1740 stats.total_compressed,
1741 stats.total_uncompressed
1742 );
1743 let reader = PackReader::from_bytes(pack_data, index_data, &std::env::temp_dir()).unwrap();
1744 let (_ty, got) = reader
1745 .get_object(&PackObjectId::Hash(hash))
1746 .unwrap()
1747 .unwrap();
1748 assert_eq!(got, payload);
1749 }
1750
1751 #[cfg(feature = "zstd")]
1752 #[test]
1753 fn equal_length_zstd_frame_is_stored_as_raw_payload() {
1754 fn encode_like_streaming_writer(data: &[u8]) -> Vec<u8> {
1758 let mut compressed = Vec::new();
1759 let mut encoder = zstd::stream::write::Encoder::new(&mut compressed, 3).unwrap();
1760 encoder
1761 .set_pledged_src_size(Some(data.len() as u64))
1762 .unwrap();
1763 std::io::Write::write_all(&mut encoder, data).unwrap();
1764 encoder.finish().unwrap();
1765 compressed
1766 }
1767
1768 let mut random = Vec::with_capacity(1024);
1769 for seed in 0u64..32 {
1770 random.extend_from_slice(blake3::hash(&seed.to_le_bytes()).as_bytes());
1771 }
1772 let payload = (256..=random.len())
1773 .find_map(|logical_len| {
1774 (0..=logical_len).find_map(|zero_prefix| {
1775 let mut candidate = random[..logical_len].to_vec();
1776 candidate[..zero_prefix].fill(0);
1777 let compressed = encode_like_streaming_writer(&candidate);
1778 (compressed.len() == candidate.len()).then_some(candidate)
1779 })
1780 })
1781 .expect("fixture search must find an equal-length zstd frame");
1782 let compressed = encode_like_streaming_writer(&payload);
1783 assert_eq!(compressed.len(), payload.len());
1784 assert_eq!(compressed.first(), Some(&0x28), "fixture must be zstd");
1785
1786 let tmp = tempfile::TempDir::new().unwrap();
1787 let (mut builder, _, index_path) = fresh_builder(&tmp);
1788 let id = PackObjectId::Hash(deterministic_hash(0x78));
1789 builder
1790 .add_id(id, ObjectType::StateAttachment, payload.clone())
1791 .unwrap();
1792 let (pack, index, stats) = finalize_cursor(builder, &index_path);
1793
1794 assert_eq!(stats.total_compressed, stats.total_uncompressed);
1795 let reader = PackReader::from_bytes(pack, index, &std::env::temp_dir()).unwrap();
1796 assert_eq!(
1797 reader.get_object(&id).unwrap(),
1798 Some((ObjectType::StateAttachment, payload))
1799 );
1800 }
1801
1802 #[cfg(feature = "zstd")]
1803 #[test]
1804 fn padded_varint_decodes_to_original_value_for_canonical_decoder() {
1805 let cases: &[u64] = &[0, 1, 127, 128, 4096, 1_000_000, 1_000_000_000_000, u64::MAX];
1811 for &value in cases {
1812 let mut buf = [0u8; 10];
1813 super::encode_varint_padded_to_10(value, &mut buf);
1814 let (decoded, consumed) = super::super::varint::decode_varint(&buf)
1815 .expect("padded varint should always decode");
1816 assert_eq!(decoded, value, "varint roundtrip failed for {value}");
1817 assert_eq!(
1818 consumed, 10,
1819 "padded encoding should consume all 10 bytes for {value}"
1820 );
1821 }
1822 }
1823
1824 #[cfg(feature = "zstd")]
1825 #[test]
1826 fn streaming_path_does_not_buffer_compressed_payload_in_memory() {
1827 let tmp = tempfile::TempDir::new().unwrap();
1840 let bucket_dir = tmp.path().join("buckets");
1841 let pack_path = tmp.path().join("pack.dat");
1842 let idx_path = tmp.path().join("pack.idx");
1843 let file = std::fs::OpenOptions::new()
1844 .read(true)
1845 .write(true)
1846 .create(true)
1847 .truncate(true)
1848 .open(&pack_path)
1849 .unwrap();
1850 let mut b = StreamingPackBuilder::new(
1851 file,
1852 idx_path.clone(),
1853 CompressionConfig::default(),
1854 bucket_dir,
1855 )
1856 .unwrap();
1857 let payload = vec![0xa5u8; 8 * 1024 * 1024];
1858 let hash = deterministic_hash(0x66);
1859 b.add(hash, ObjectType::Blob, payload.clone()).unwrap();
1860 let mid_size = std::fs::metadata(&pack_path).unwrap().len();
1864 assert!(
1865 mid_size > 16 + 40,
1866 "pack file should hold real entry data after add; size={mid_size}"
1867 );
1868 let (_, _) = b.finalize().unwrap();
1869 let pack_bytes = std::fs::read(&pack_path).unwrap();
1870 let index_bytes = std::fs::read(&idx_path).unwrap();
1871 let reader =
1872 PackReader::from_bytes(pack_bytes, index_bytes, &std::env::temp_dir()).unwrap();
1873 let (_ty, got) = reader
1874 .get_object(&PackObjectId::Hash(hash))
1875 .unwrap()
1876 .unwrap();
1877 assert_eq!(got, payload);
1878 }
1879
1880 #[test]
1881 fn list_ids_returns_all_added_ids_sorted() {
1882 let tmp = tempfile::TempDir::new().unwrap();
1883 let (mut b, _, idx_path) = fresh_builder(&tmp);
1884 let mut added: Vec<PackObjectId> = Vec::new();
1885 for seed in [0x05u8, 0xa0, 0x12, 0x9f, 0x33] {
1887 let id = PackObjectId::Hash(deterministic_hash(seed));
1888 b.add_id(id, ObjectType::Blob, vec![seed; 4]).unwrap();
1889 added.push(id);
1890 }
1891 for seed in [0x80u8, 0x10, 0xff] {
1892 let id = PackObjectId::StateId(deterministic_state_id(seed));
1893 b.add_id(id, ObjectType::State, vec![seed; 4]).unwrap();
1894 added.push(id);
1895 }
1896 let (pack_data, index_data, _) = finalize_cursor(b, &idx_path);
1897 let reader = PackReader::from_bytes(pack_data, index_data, &std::env::temp_dir()).unwrap();
1898 let mut got = reader.list_ids().unwrap();
1899 let mut sorted = got.clone();
1902 sorted.sort();
1903 assert_eq!(got, sorted, "list_ids must come back sorted");
1904 added.sort();
1906 got.sort();
1907 assert_eq!(got, added);
1908 }
1909}