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