1use std::fs::File;
49use std::os::unix::fs::FileExt;
50use std::path::Path;
51use std::sync::Arc;
52
53use anyhow::{Result, anyhow};
54use arrow::array::{
55 BooleanArray, BooleanBuilder, FixedSizeBinaryArray, FixedSizeBinaryBuilder, StringArray,
56 StringBuilder, UInt32Array, UInt32Builder, UInt64Array, UInt64Builder,
57};
58use arrow::datatypes::Schema;
59use arrow::ipc::writer::StreamWriter;
60use arrow::record_batch::RecordBatch;
61
62use crate::index::{
63 ChunkLoc, LOOKUP_MODULE, MULTI_INDEX_MAGIC, ManifestEntry, RESERVED_PKG_TYPE, TRIE_MODULE,
64 data_subindex_schema, is_reserved_module, lookup_schema, read_znippy_full_manifest,
65 write_manifest_bytes,
66};
67use crate::index::{
68 CARRIED_RESERVED_MODULES, META_MODULE, ZNIPPY_DELTA_MODULE, read_reserved_section_bytes,
69};
70use crate::meta_index::{
71 MetaTable, build_meta_batch, decode_meta_section, meta_schema,
72};
73use crate::meta_sink::{ArchiveMetaSink, GroupKey};
74
75pub struct ArrowIpcSinkAppend {
81 file: Arc<File>,
82 cursor: u64,
83 entries: Vec<ManifestEntry>,
84 lookup_paths: Vec<String>,
85 lookup_locs: Vec<ChunkLoc>,
86 carried: Vec<(String, ChunkLoc)>,
92 meta: Option<MetaTable>,
103 carried_reserved: Vec<(String, Vec<u8>)>,
117 delta_map: Vec<(String, u32, String)>,
123}
124
125impl ArrowIpcSinkAppend {
126 pub fn new(file: Arc<File>, blob_end_offset: u64) -> Self {
129 Self {
130 file,
131 cursor: blob_end_offset,
132 entries: Vec::new(),
133 lookup_paths: Vec::new(),
134 lookup_locs: Vec::new(),
135 carried: Vec::new(),
136 meta: None,
137 carried_reserved: Vec::new(),
138 delta_map: Vec::new(),
139 }
140 }
141
142 pub fn with_meta(mut self, meta: MetaTable) -> Self {
145 self.meta = Some(meta);
146 self
147 }
148
149 pub fn merge_meta(&mut self, rows: impl IntoIterator<Item = crate::meta_index::MetaEntry>) {
152 self.meta.get_or_insert_with(MetaTable::new).extend(rows);
153 }
154
155 pub fn meta(&self) -> Option<&MetaTable> {
157 self.meta.as_ref()
158 }
159
160 pub fn open_existing(path: &Path) -> Result<Self> {
171 let file = Arc::new(
172 std::fs::OpenOptions::new()
173 .read(true)
174 .write(true)
175 .open(path)
176 .map_err(|e| anyhow!("append: open {} for resume: {e}", path.display()))?,
177 );
178
179 let (entries, _manifest_offset) = read_znippy_full_manifest(path)?;
181 if entries.is_empty() {
182 return Err(anyhow!("append: archive {} has an empty manifest", path.display()));
183 }
184
185 let blob_end = entries
188 .iter()
189 .map(|e| e.index_offset)
190 .min()
191 .ok_or_else(|| anyhow!("append: no sections in manifest"))?;
192
193 let (paths, locs) = recover_rows(path, &entries)?;
197 let carried: Vec<(String, ChunkLoc)> = paths.into_iter().zip(locs).collect();
198
199 let meta = match read_reserved_section_bytes(path, META_MODULE)? {
206 None => None,
207 Some(bytes) => Some(decode_meta_section(&bytes)?.to_table()),
208 };
209
210 let mut carried_reserved = Vec::new();
215 for module in CARRIED_RESERVED_MODULES {
216 if let Some(bytes) = read_reserved_section_bytes(path, module)? {
217 carried_reserved.push(((*module).to_string(), bytes));
218 }
219 }
220
221 let delta_map = read_delta_map(path)?;
226
227 Ok(Self {
228 file,
229 cursor: blob_end,
230 entries: Vec::new(), lookup_paths: Vec::new(),
232 lookup_locs: Vec::new(),
233 carried,
234 meta,
235 carried_reserved,
236 delta_map,
237 })
238 }
239
240 pub fn blob_end(&self) -> u64 {
245 self.cursor
246 }
247
248 pub fn recovered_rows(&self) -> usize {
250 self.carried.len()
251 }
252
253 fn drop_carried_paths(&mut self, replacing: &std::collections::HashSet<&str>) -> usize {
266 if self.carried.is_empty() || replacing.is_empty() {
267 return 0;
268 }
269 let before = self.carried.len();
270 self.carried.retain(|(p, _)| !replacing.contains(p.as_str()));
271 before - self.carried.len()
272 }
273
274 fn emit_carried(&mut self) -> Result<()> {
279 if self.carried.is_empty() {
280 return Ok(());
281 }
282 let carried = std::mem::take(&mut self.carried);
283 let (paths, locs): (Vec<String>, Vec<ChunkLoc>) = carried.into_iter().unzip();
284 let batch = base_batch_from_rows(&paths, &locs)?;
285 self.push_subindex(data_subindex_schema().as_ref(), &[batch], GroupKey {
286 pkg_type: 0,
287 repo: String::new(),
288 module_name: String::new(),
289 })
290 }
291
292 fn accumulate_lookup(&mut self, batch: &RecordBatch) {
295 let cols = (|| {
296 Some((
297 batch.column_by_name("relative_path")?.as_any().downcast_ref::<StringArray>()?,
298 batch.column_by_name("chunk_seq")?.as_any().downcast_ref::<UInt32Array>()?,
299 batch.column_by_name("fdata_offset")?.as_any().downcast_ref::<UInt64Array>()?,
300 batch.column_by_name("compressed")?.as_any().downcast_ref::<BooleanArray>()?,
301 batch.column_by_name("uncompressed_size")?.as_any().downcast_ref::<UInt64Array>()?,
302 batch.column_by_name("blob_offset")?.as_any().downcast_ref::<UInt64Array>()?,
303 batch.column_by_name("blob_size")?.as_any().downcast_ref::<UInt64Array>()?,
304 batch.column_by_name("checksum")?.as_any().downcast_ref::<FixedSizeBinaryArray>()?,
305 ))
306 })();
307 let Some((paths, chunk_seq, fdata, compressed, usz, blob_off, blob_sz, checksum)) = cols
308 else { return; };
309 for i in 0..batch.num_rows() {
310 let mut ck = [0u8; 32];
311 ck.copy_from_slice(checksum.value(i));
312 self.lookup_paths.push(paths.value(i).to_string());
313 self.lookup_locs.push(ChunkLoc {
314 chunk_seq: chunk_seq.value(i),
315 fdata_offset: fdata.value(i),
316 blob_offset: blob_off.value(i),
317 blob_size: blob_sz.value(i),
318 uncompressed_size: usz.value(i),
319 compressed: compressed.value(i),
320 checksum: ck,
321 });
322 }
323 }
324
325 fn write_lookup_and_trie(&mut self) -> Result<()> {
326 let n = self.lookup_paths.len();
327 let mut order: Vec<usize> = (0..n).collect();
328 order.sort_by(|&a, &b| {
329 self.lookup_paths[a].cmp(&self.lookup_paths[b])
330 .then(self.lookup_locs[a].chunk_seq.cmp(&self.lookup_locs[b].chunk_seq))
331 });
332
333 let schema = lookup_schema();
334 let batch = base_batch_permuted(
335 schema.clone(),
336 &self.lookup_paths,
337 &self.lookup_locs,
338 &order,
339 )?;
340 self.push_subindex(&schema, &[batch], GroupKey {
341 pkg_type: RESERVED_PKG_TYPE,
342 repo: String::new(),
343 module_name: LOOKUP_MODULE.to_string(),
344 })?;
345
346 let mut builder = fst::MapBuilder::memory();
347 let mut prev: Option<&str> = None;
348 for (sorted_idx, &orig) in order.iter().enumerate() {
349 let p = self.lookup_paths[orig].as_str();
350 if prev != Some(p) {
351 builder.insert(p.as_bytes(), sorted_idx as u64)
352 .map_err(|e| anyhow!("trie insert: {e}"))?;
353 prev = Some(p);
354 }
355 }
356 let trie_bytes = builder.into_inner().map_err(|e| anyhow!("trie finish: {e}"))?;
357 self.write_raw_section(&trie_bytes, GroupKey {
358 pkg_type: RESERVED_PKG_TYPE,
359 repo: String::new(),
360 module_name: TRIE_MODULE.to_string(),
361 })
362 }
363
364 fn write_meta_subindex(&mut self) -> Result<()> {
371 let Some(table) = self.meta.take() else {
372 return Ok(());
373 };
374 let batch = build_meta_batch(&table)?;
375 let schema = meta_schema();
376 self.push_subindex(schema.as_ref(), &[batch], GroupKey {
377 pkg_type: RESERVED_PKG_TYPE,
378 repo: String::new(),
379 module_name: META_MODULE.to_string(),
380 })
381 }
382
383 fn write_carried_reserved(&mut self) -> Result<()> {
389 for (module, bytes) in std::mem::take(&mut self.carried_reserved) {
390 self.write_raw_section(&bytes, GroupKey {
391 pkg_type: RESERVED_PKG_TYPE,
392 repo: String::new(),
393 module_name: module,
394 })?;
395 }
396 Ok(())
397 }
398
399 fn write_delta_map(&mut self) -> Result<()> {
401 if self.delta_map.is_empty() {
402 return Ok(());
403 }
404 let rows = std::mem::take(&mut self.delta_map);
405 let paths = StringArray::from(rows.iter().map(|r| r.0.as_str()).collect::<Vec<_>>());
406 let seqs = UInt32Array::from(rows.iter().map(|r| r.1).collect::<Vec<_>>());
407 let bases = StringArray::from(rows.iter().map(|r| r.2.as_str()).collect::<Vec<_>>());
408 let schema = crate::index::delta_map_schema();
409 let batch = RecordBatch::try_new(
410 Arc::clone(&schema),
411 vec![Arc::new(paths), Arc::new(seqs), Arc::new(bases)],
412 )
413 .map_err(|e| anyhow!("delta map batch: {e}"))?;
414 self.push_subindex(schema.as_ref(), &[batch], GroupKey {
415 pkg_type: RESERVED_PKG_TYPE,
416 repo: String::new(),
417 module_name: ZNIPPY_DELTA_MODULE.to_string(),
418 })
419 }
420
421 pub fn file(&self) -> &Arc<File> {
423 &self.file
424 }
425
426 pub fn advance_blob_end(&mut self, n: u64) {
428 self.cursor += n;
429 }
430
431 pub fn replace_carried(&mut self, path: &str, locs: Vec<ChunkLoc>) {
436 self.carried.retain(|(p, _)| p != path);
437 for loc in locs {
438 self.carried.push((path.to_string(), loc));
439 }
440 }
441
442 pub fn push_delta_map_row(&mut self, path: String, chunk_seq: u32, base: String) {
444 self.delta_map.retain(|(p, s, _)| !(p == &path && *s == chunk_seq));
445 self.delta_map.push((path, chunk_seq, base));
446 }
447
448 fn write_raw_section(&mut self, bytes: &[u8], key: GroupKey) -> Result<()> {
449 let start = self.cursor;
450 self.file.write_all_at(bytes, start)?;
451 self.cursor += bytes.len() as u64;
452 self.entries.push(ManifestEntry {
453 pkg_type: key.pkg_type,
454 repo: key.repo,
455 module_name: key.module_name,
456 index_offset: start,
457 index_len: bytes.len() as u64,
458 row_count: 0,
459 });
460 Ok(())
461 }
462}
463
464impl ArchiveMetaSink for ArrowIpcSinkAppend {
465 fn push_subindex(
466 &mut self,
467 schema: &Schema,
468 batches: &[RecordBatch],
469 key: GroupKey,
470 ) -> Result<()> {
471 let sub_start = self.cursor;
472 let mut sub_bytes: Vec<u8> = Vec::new();
473 let mut sw = StreamWriter::try_new(&mut sub_bytes, schema)
474 .map_err(|e| anyhow!("sub-index writer: {e}"))?;
475 let mut row_count = 0u64;
476 for batch in batches {
477 row_count += batch.num_rows() as u64;
478 sw.write(batch).map_err(|e| anyhow!("sub-index write: {e}"))?;
479 }
480 sw.finish().map_err(|e| anyhow!("sub-index finish: {e}"))?;
481
482 if !is_reserved_module(&key.module_name) {
489 for batch in batches {
490 self.accumulate_lookup(batch);
491 }
492 }
493
494 let sub_len = sub_bytes.len() as u64;
495 self.file.write_all_at(&sub_bytes, sub_start)?;
496 self.cursor += sub_len;
497
498 self.entries.push(ManifestEntry {
499 pkg_type: key.pkg_type,
500 repo: key.repo,
501 module_name: key.module_name,
502 index_offset: sub_start,
503 index_len: sub_len,
504 row_count,
505 });
506 Ok(())
507 }
508
509 fn finish(mut self: Box<Self>) -> Result<u64> {
510 self.emit_carried()?;
513 self.write_lookup_and_trie()?;
514 self.write_meta_subindex()?;
515 self.write_carried_reserved()?;
516 self.write_delta_map()?;
517
518 let manifest_offset = self.cursor;
519 let manifest_bytes =
520 write_manifest_bytes(&self.entries).map_err(|e| anyhow!("manifest: {e}"))?;
521 self.file.write_all_at(&manifest_bytes, manifest_offset)?;
522
523 let after = manifest_offset + manifest_bytes.len() as u64;
524 self.file.write_all_at(&MULTI_INDEX_MAGIC, after)?;
525 self.file.write_all_at(
526 &manifest_offset.to_le_bytes(),
527 after + MULTI_INDEX_MAGIC.len() as u64,
528 )?;
529 let final_len = after + MULTI_INDEX_MAGIC.len() as u64 + 8;
532 self.file.set_len(final_len)?;
533 self.file.sync_all()?;
534
535 Ok(final_len)
536 }
537}
538
539
540pub fn read_delta_map(path: &Path) -> Result<Vec<(String, u32, String)>> {
552 use arrow::ipc::reader::StreamReader;
553 let Some(bytes) = read_reserved_section_bytes(path, ZNIPPY_DELTA_MODULE)? else {
554 return Ok(Vec::new());
555 };
556 let reader = StreamReader::try_new(std::io::Cursor::new(bytes), None)
557 .map_err(|e| anyhow!("delta map: {e}"))?;
558 let mut out = Vec::new();
559 for batch in reader {
560 let batch = batch.map_err(|e| anyhow!("delta map batch: {e}"))?;
561 let paths = batch
562 .column_by_name("relative_path")
563 .and_then(|c| c.as_any().downcast_ref::<StringArray>())
564 .ok_or_else(|| anyhow!("delta map: missing relative_path"))?;
565 let seqs = batch
566 .column_by_name("chunk_seq")
567 .and_then(|c| c.as_any().downcast_ref::<UInt32Array>())
568 .ok_or_else(|| anyhow!("delta map: missing chunk_seq"))?;
569 let bases = batch
570 .column_by_name("base_path")
571 .and_then(|c| c.as_any().downcast_ref::<StringArray>())
572 .ok_or_else(|| anyhow!("delta map: missing base_path"))?;
573 for r in 0..batch.num_rows() {
574 out.push((paths.value(r).to_string(), seqs.value(r), bases.value(r).to_string()));
575 }
576 }
577 Ok(out)
578}
579
580
581pub fn supersede_as_delta(
617 archive: &Path,
618 superseded: &str,
619 base: &str,
620 chain_depth: usize,
621 compression_level: i32,
622) -> Result<SupersedeOutcome> {
623 if superseded == base {
624 return Err(anyhow!("an entry cannot be a delta against itself: {superseded}"));
625 }
626 if chain_depth + 1 > MAX_GENERATION_CHAIN {
627 return Ok(SupersedeOutcome::ChainTooLong);
628 }
629
630 let (old_bytes, base_bytes) = {
631 let ar = crate::ZnippyArchive::open(archive)?;
632 (
636 ar.extract_file_verified(superseded)?,
637 ar.extract_file_verified(base)?,
638 )
639 };
640
641 let delta = crate::archive::encode_delta_against(&base_bytes, &old_bytes);
642 if (delta.len() as f64) >= crate::archive::DELTA_SIZE_ALPHA * (old_bytes.len() as f64) {
646 return Ok(SupersedeOutcome::NotSmaller {
647 delta_bytes: delta.len() as u64,
648 stored_bytes: old_bytes.len() as u64,
649 });
650 }
651
652 let mut sink = ArrowIpcSinkAppend::open_existing(archive)?;
653 let at = sink.blob_end();
654 let mut ctx = crate::codec::CompressCtx::new(compression_level)?;
659 let frame = ctx.compress(&delta).ok();
660 let (on_disk, compressed): (&[u8], bool) = match frame.as_deref() {
661 Some(f) if f.len() < delta.len() => (f, true),
662 _ => (&delta, false),
663 };
664 sink.file().write_all_at(on_disk, at)?;
665 sink.advance_blob_end(on_disk.len() as u64);
666
667 sink.replace_carried(superseded, vec![ChunkLoc {
669 chunk_seq: 0,
670 fdata_offset: 0,
671 blob_offset: at,
672 blob_size: on_disk.len() as u64,
673 uncompressed_size: old_bytes.len() as u64,
676 compressed,
677 checksum: *blake3::hash(&old_bytes).as_bytes(),
678 }]);
679 sink.push_delta_map_row(superseded.to_string(), 0, base.to_string());
680 Box::new(sink).finish()?;
681
682 Ok(SupersedeOutcome::Delta {
683 stored_bytes: old_bytes.len() as u64,
684 delta_bytes: on_disk.len() as u64,
685 chain_depth: chain_depth + 1,
686 })
687}
688
689#[derive(Debug, Clone, Copy, PartialEq, Eq)]
691pub struct CompactReport {
692 pub bytes_before: u64,
693 pub bytes_after: u64,
694 pub rows: u64,
696 pub delta_rows: u64,
698}
699
700pub fn compact_archive(archive: &Path) -> Result<CompactReport> {
730 let bytes_before = std::fs::metadata(archive)?.len();
731 let src = ArrowIpcSinkAppend::open_existing(archive)?;
732
733 let staged = {
734 let unique = std::time::SystemTime::now()
735 .duration_since(std::time::UNIX_EPOCH)
736 .map(|d| d.as_nanos())
737 .unwrap_or(0);
738 let mut p = archive.as_os_str().to_owned();
739 p.push(format!(".compact-{}-{unique}", std::process::id()));
740 std::path::PathBuf::from(p)
741 };
742 let out = Arc::new(
743 std::fs::OpenOptions::new()
744 .read(true)
745 .write(true)
746 .create(true)
747 .truncate(true)
748 .open(&staged)
749 .map_err(|e| anyhow!("compact: staging {}: {e}", staged.display()))?,
750 );
751
752 let mut sink = ArrowIpcSinkAppend::new(Arc::clone(&out), 0);
753 sink.meta = src.meta.clone();
754 sink.carried_reserved = src.carried_reserved.clone();
755 sink.delta_map = src.delta_map.clone();
756 let delta_rows = sink.delta_map.len() as u64;
757
758 let mut cursor = 0u64;
759 let mut buf: Vec<u8> = Vec::new();
763 for (path, loc) in &src.carried {
764 let n = loc.blob_size as usize;
765 buf.clear();
766 buf.resize(n, 0);
767 src.file
768 .read_exact_at(&mut buf, loc.blob_offset)
769 .map_err(|e| anyhow!("compact: reading {path} at {}: {e}", loc.blob_offset))?;
770 out.write_all_at(&buf, cursor)?;
771 let mut moved = loc.clone();
772 moved.blob_offset = cursor;
773 cursor += loc.blob_size;
774 sink.carried.push((path.clone(), moved));
775 }
776 let rows = sink.carried.len() as u64;
777 sink.cursor = cursor;
778 Box::new(sink).finish()?;
779
780 out.sync_all()?;
781 drop(out);
782 drop(src);
783 std::fs::rename(&staged, archive)?;
784 if let Some(parent) = archive.parent() {
785 if let Ok(f) = std::fs::File::open(parent) {
786 let _ = f.sync_all();
787 }
788 }
789
790 Ok(CompactReport {
791 bytes_before,
792 bytes_after: std::fs::metadata(archive)?.len(),
793 rows,
794 delta_rows,
795 })
796}
797
798#[derive(Debug, PartialEq, Eq)]
801pub enum SupersedeOutcome {
802 Delta { stored_bytes: u64, delta_bytes: u64, chain_depth: usize },
804 NotSmaller { delta_bytes: u64, stored_bytes: u64 },
807 ChainTooLong,
810}
811
812pub const MAX_GENERATION_CHAIN: usize = 32;
837
838pub fn base_batch_from_rows(paths: &[String], locs: &[ChunkLoc]) -> Result<RecordBatch> {
850 let order: Vec<usize> = (0..paths.len()).collect();
851 base_batch_permuted(data_subindex_schema(), paths, locs, &order)
852}
853
854pub(crate) fn base_batch_permuted(
857 schema: Arc<Schema>,
858 paths: &[String],
859 locs: &[ChunkLoc],
860 order: &[usize],
861) -> Result<RecordBatch> {
862 let n = order.len();
863 let mut path_b = StringBuilder::with_capacity(n, n * 16);
864 let mut seq_b = UInt32Builder::with_capacity(n);
865 let mut fdata_b = UInt64Builder::with_capacity(n);
866 let mut comp_b = BooleanBuilder::with_capacity(n);
867 let mut usz_b = UInt64Builder::with_capacity(n);
868 let mut boff_b = UInt64Builder::with_capacity(n);
869 let mut bsz_b = UInt64Builder::with_capacity(n);
870 let mut ck_b = FixedSizeBinaryBuilder::with_capacity(n, 32);
871 for &i in order {
872 let loc = &locs[i];
873 path_b.append_value(&paths[i]);
874 seq_b.append_value(loc.chunk_seq);
875 fdata_b.append_value(loc.fdata_offset);
876 comp_b.append_value(loc.compressed);
877 usz_b.append_value(loc.uncompressed_size);
878 boff_b.append_value(loc.blob_offset);
879 bsz_b.append_value(loc.blob_size);
880 ck_b.append_value(loc.checksum).expect("checksum is 32 bytes");
881 }
882 Ok(RecordBatch::try_new(
883 schema,
884 vec![
885 Arc::new(path_b.finish()),
886 Arc::new(seq_b.finish()),
887 Arc::new(fdata_b.finish()),
888 Arc::new(comp_b.finish()),
889 Arc::new(usz_b.finish()),
890 Arc::new(boff_b.finish()),
891 Arc::new(bsz_b.finish()),
892 Arc::new(ck_b.finish()),
893 ],
894 )?)
895}
896
897pub(crate) fn write_blobs(
903 file: &File,
904 cursor: u64,
905 files: &[(String, Vec<u8>)],
906 ctx: &mut crate::codec::CompressCtx,
907 policy: crate::SkipPolicy,
908) -> Result<(Vec<String>, Vec<ChunkLoc>, u64)> {
909 let mut paths = Vec::with_capacity(files.len());
910 let mut locs = Vec::with_capacity(files.len());
911 let mut cursor = cursor;
912 for (rel, bytes) in files {
913 let checksum = *blake3::hash(bytes).as_bytes();
914 let skip = policy.skip_by_path(std::path::Path::new(rel.as_str()));
924 let frame = if skip { Vec::new() } else { ctx.compress(bytes)? };
925 let (on_disk, compressed): (&[u8], bool) = if !skip && frame.len() < bytes.len() {
926 (&frame, true)
927 } else {
928 (bytes, false)
929 };
930 let blob_offset = cursor;
931 file.write_all_at(on_disk, blob_offset)?;
932 cursor += on_disk.len() as u64;
933 paths.push(rel.clone());
934 locs.push(ChunkLoc {
935 chunk_seq: 0,
936 fdata_offset: 0,
937 blob_offset,
938 blob_size: on_disk.len() as u64,
939 uncompressed_size: bytes.len() as u64,
940 compressed,
941 checksum,
942 });
943 }
944 Ok((paths, locs, cursor))
945}
946
947#[derive(Debug, Clone)]
949pub struct AppendReport {
950 pub rows_before: u64,
953 pub rows_replaced: u64,
956 pub rows_added: u64,
958 pub blob_append_offset: u64,
960 pub blob_bytes_added: u64,
962 pub sealed_total_bytes: u64,
964}
965
966pub fn append_files(
986 archive: &Path,
987 new_files: &[(String, Vec<u8>)],
988 compression_level: i32,
989) -> Result<AppendReport> {
990 append_files_with_meta(archive, new_files, compression_level, None)
991}
992
993pub fn append_files_with_policy(
1000 archive: &Path,
1001 new_files: &[(String, Vec<u8>)],
1002 compression_level: i32,
1003 policy: crate::SkipPolicy,
1004) -> Result<AppendReport> {
1005 let sink = ArrowIpcSinkAppend::open_existing(archive)?;
1006 let rows_before = sink.recovered_rows() as u64;
1007 let blob_append_offset = sink.blob_end();
1008 write_files_into_sink(
1009 sink,
1010 new_files,
1011 compression_level,
1012 rows_before,
1013 blob_append_offset,
1014 policy,
1015 )
1016}
1017
1018pub fn append_files_with_meta(
1025 archive: &Path,
1026 new_files: &[(String, Vec<u8>)],
1027 compression_level: i32,
1028 meta: Option<MetaTable>,
1029) -> Result<AppendReport> {
1030 let mut sink = ArrowIpcSinkAppend::open_existing(archive)?;
1031 if let Some(table) = meta {
1032 sink.merge_meta(table.rows().to_vec());
1033 }
1034 let rows_before = sink.recovered_rows() as u64;
1035 let blob_append_offset = sink.blob_end();
1036 write_files_into_sink(
1037 sink,
1038 new_files,
1039 compression_level,
1040 rows_before,
1041 blob_append_offset,
1042 crate::SkipPolicy::resolve(),
1047 )
1048}
1049
1050pub fn create_archive(
1055 archive: &Path,
1056 files: &[(String, Vec<u8>)],
1057 compression_level: i32,
1058) -> Result<AppendReport> {
1059 create_archive_with_meta(archive, files, compression_level, None)
1060}
1061
1062pub fn create_archive_with_meta(
1070 archive: &Path,
1071 files: &[(String, Vec<u8>)],
1072 compression_level: i32,
1073 meta: Option<MetaTable>,
1074) -> Result<AppendReport> {
1075 let blob_file = Arc::new(
1076 File::create(archive)
1077 .map_err(|e| anyhow!("create archive {}: {e}", archive.display()))?,
1078 );
1079 let mut sink = ArrowIpcSinkAppend::new(blob_file, 0);
1080 sink.meta = meta;
1081 write_files_into_sink(sink, files, compression_level, 0, 0, crate::SkipPolicy::resolve())
1082}
1083
1084pub fn create_archive_to_vec(
1095 files: &[(String, Vec<u8>)],
1096 compression_level: i32,
1097) -> Result<(Vec<u8>, AppendReport)> {
1098 let anon = Arc::new(
1099 tempfile::tempfile().map_err(|e| anyhow!("anonymous archive fd: {e}"))?,
1100 );
1101 let sink = ArrowIpcSinkAppend::new(anon.clone(), 0);
1102 let report =
1103 write_files_into_sink(sink, files, compression_level, 0, 0, crate::SkipPolicy::resolve())?;
1104 let mut bytes = vec![0u8; report.sealed_total_bytes as usize];
1106 anon.read_exact_at(&mut bytes, 0)
1107 .map_err(|e| anyhow!("read back anonymous archive: {e}"))?;
1108 Ok((bytes, report))
1109}
1110
1111fn write_files_into_sink(
1116 mut sink: ArrowIpcSinkAppend,
1117 new_files: &[(String, Vec<u8>)],
1118 compression_level: i32,
1119 rows_before: u64,
1120 blob_append_offset: u64,
1121 policy: crate::SkipPolicy,
1122) -> Result<AppendReport> {
1123 use crate::codec::CompressCtx;
1124
1125 let incoming: std::collections::HashSet<&str> =
1129 new_files.iter().map(|(rel, _)| rel.as_str()).collect();
1130 let rows_replaced = sink.drop_carried_paths(&incoming) as u64;
1131
1132 let blob_file = sink.file.clone();
1136 let mut ctx = CompressCtx::new(compression_level)?;
1137 let (paths, locs, cursor) =
1138 write_blobs(&blob_file, blob_append_offset, new_files, &mut ctx, policy)?;
1139 let blob_bytes_added = cursor - blob_append_offset;
1140 blob_file.sync_all()?;
1141
1142 sink.cursor = cursor;
1145
1146 let batch = base_batch_from_rows(&paths, &locs)?;
1151 let schema = data_subindex_schema();
1152 let rows_added = batch.num_rows() as u64;
1153 sink.push_subindex(
1154 schema.as_ref(),
1155 &[batch],
1156 GroupKey { pkg_type: 0, repo: String::new(), module_name: String::new() },
1157 )?;
1158
1159 let sealed_total_bytes = Box::new(sink).finish()?;
1160
1161 Ok(AppendReport {
1162 rows_before,
1163 rows_replaced,
1164 rows_added,
1165 blob_append_offset,
1166 blob_bytes_added,
1167 sealed_total_bytes,
1168 })
1169}
1170
1171pub(crate) fn recover_rows(
1175 path: &Path,
1176 entries: &[ManifestEntry],
1177) -> Result<(Vec<String>, Vec<ChunkLoc>)> {
1178 use std::io::{Read, Seek, SeekFrom};
1179
1180 let mut file = File::open(path)?;
1181
1182 if let Some(lk) = entries.iter().find(|e| e.module_name == LOOKUP_MODULE) {
1185 file.seek(SeekFrom::Start(lk.index_offset))?;
1186 let mut bytes = vec![0u8; lk.index_len as usize];
1187 file.read_exact(&mut bytes)?;
1188 return decode_base_rows(&bytes);
1189 }
1190
1191 let mut paths = Vec::new();
1193 let mut locs = Vec::new();
1194 for e in entries {
1195 if is_reserved_module(&e.module_name) {
1196 continue;
1197 }
1198 file.seek(SeekFrom::Start(e.index_offset))?;
1199 let mut bytes = vec![0u8; e.index_len as usize];
1200 file.read_exact(&mut bytes)?;
1201 let (mut p, mut l) = decode_base_rows(&bytes)?;
1202 paths.append(&mut p);
1203 locs.append(&mut l);
1204 }
1205 Ok((paths, locs))
1206}
1207
1208#[cfg(all(test, feature = "openzl"))]
1212mod tests {
1213 use super::*;
1214 use crate::codec::CompressCtx;
1215 use crate::meta::{BlobMeta, ChunkMeta};
1216 use crate::{ArrowIpcSink, ZnippyArchive, ZnippyReader};
1217 use std::time::{SystemTime, UNIX_EPOCH};
1218
1219 fn unique_dir(tag: &str) -> std::path::PathBuf {
1220 let ns = SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_nanos();
1221 let d = std::env::temp_dir().join(format!("znippy_append_{tag}_{ns}_{:?}", std::thread::current().id()));
1222 std::fs::create_dir_all(&d).unwrap();
1223 d
1224 }
1225
1226 #[test]
1245 fn an_append_honours_the_skip_policy_instead_of_compressing_everything() {
1246 let dir = unique_dir("skip_policy");
1247 let archive = dir.join("a.znippy");
1248 create_archive(&archive, &[("seed.txt".into(), b"seed".to_vec())], 3).unwrap();
1249
1250 let squishy = vec![b'A'; 256 * 1024];
1252 let name = format!("pack-{}.pack", "0f".repeat(20));
1253 let report =
1254 append_files(&archive, &[(name.clone(), squishy.clone())], 3).unwrap();
1255
1256 assert_eq!(
1257 report.blob_bytes_added,
1258 squishy.len() as u64,
1259 "a `.pack` entry must be stored RAW. {} bytes were written for a {}-byte input, so \
1260 the codec ran over a file the extension table already said was compressed",
1261 report.blob_bytes_added,
1262 squishy.len()
1263 );
1264 assert_eq!(crate::get_file(&archive, &name).unwrap(), squishy);
1266
1267 let report2 =
1271 append_files(&archive, &[("plain.txt".into(), squishy.clone())], 3).unwrap();
1272 assert!(
1273 report2.blob_bytes_added < squishy.len() as u64 / 10,
1274 "an ordinary name must still be compressed; {} bytes for {}",
1275 report2.blob_bytes_added,
1276 squishy.len()
1277 );
1278 assert_eq!(crate::get_file(&archive, "plain.txt").unwrap(), squishy);
1279
1280 let report3 = append_files_with_policy(
1282 &archive,
1283 &[("also-plain.txt".into(), squishy.clone())],
1284 3,
1285 crate::SkipPolicy::already_compressed(),
1286 )
1287 .unwrap();
1288 assert_eq!(
1289 report3.blob_bytes_added,
1290 squishy.len() as u64,
1291 "`already_compressed()` must store raw whatever the name says"
1292 );
1293
1294 std::fs::remove_dir_all(&dir).ok();
1295 }
1296
1297 fn synth(n: usize, salt: u64) -> Vec<(String, Vec<u8>)> {
1300 (0..n)
1301 .map(|i| {
1302 let g = (i.wrapping_mul(2_654_435_761) ^ salt as usize) % 1000;
1303 let p = format!("repo/grp{g:03}/file{:08}_{salt}.bin", i);
1304 let body = format!("payload {i} salt {salt} {}\n", "z".repeat(8 + (i % 40)));
1305 (p, body.into_bytes())
1306 })
1307 .collect()
1308 }
1309
1310 fn write_fresh<S: ArchiveMetaSink + 'static>(
1315 path: &Path,
1316 files: &[(String, Vec<u8>)],
1317 make_sink: impl FnOnce(Arc<File>, u64) -> S,
1318 ) -> u64 {
1319 let file = Arc::new(File::create(path).unwrap());
1320 let mut ctx = CompressCtx::new(3).unwrap();
1321 let mut blobs = Vec::new();
1322 let mut paths = Vec::new();
1323 let mut cursor = 0u64;
1324 for (fi, (rel, bytes)) in files.iter().enumerate() {
1325 let checksum = *blake3::hash(bytes).as_bytes();
1326 let frame = ctx.compress(bytes).unwrap();
1327 let (on_disk, compressed): (&[u8], bool) =
1328 if frame.len() < bytes.len() { (&frame, true) } else { (bytes, false) };
1329 file.write_all_at(on_disk, cursor).unwrap();
1330 let blob_offset = cursor;
1331 cursor += on_disk.len() as u64;
1332 paths.push(rel.clone());
1333 blobs.push(BlobMeta {
1334 blob_offset,
1335 blob_size: on_disk.len() as u64,
1336 chunk_meta: ChunkMeta {
1337 fdata_offset: 0,
1338 file_index: fi as u64,
1339 chunk_seq: 0,
1340 checksum,
1341 compressed,
1342 uncompressed_size: bytes.len() as u64,
1343 compressed_size: on_disk.len() as u64,
1344 },
1345 });
1346 }
1347 let resolver = { let p = paths.clone(); move |fi: u64| p[fi as usize].clone() };
1348 let batch = crate::build_metadata_batch(&blobs, resolver, &[], &[]).unwrap();
1349 let schema = data_subindex_schema();
1352 let mut sink = make_sink(file.clone(), cursor);
1353 sink.push_subindex(
1354 schema.as_ref(),
1355 &[batch],
1356 GroupKey { pkg_type: 0, repo: String::new(), module_name: String::new() },
1357 )
1358 .unwrap();
1359 Box::new(sink).finish().unwrap()
1360 }
1361
1362 #[test]
1367 fn clone_fresh_path_is_byte_identical_to_original() {
1368 let dir = unique_dir("parity");
1369 let files = synth(2_000, 1);
1370
1371 let a = dir.join("a.znippy");
1372 let b = dir.join("b.znippy");
1373 let len_a = write_fresh(&a, &files, ArrowIpcSink::new);
1374 let len_b = write_fresh(&b, &files, ArrowIpcSinkAppend::new);
1375
1376 assert_eq!(len_a, len_b, "clone seal produced a different total length");
1377 let bytes_a = std::fs::read(&a).unwrap();
1378 let bytes_b = std::fs::read(&b).unwrap();
1379 assert_eq!(
1380 bytes_a, bytes_b,
1381 "clone's fresh write path is NOT byte-identical to ArrowIpcSink — parity broken"
1382 );
1383 let _ = std::fs::remove_dir_all(&dir);
1384 }
1385
1386 #[test]
1390 fn native_append_roundtrips_old_and_new_files() {
1391 let dir = unique_dir("resume");
1392 let archive = dir.join("store.znippy");
1393 let orig = synth(1_500, 7);
1394 write_fresh(&archive, &orig, ArrowIpcSink::new);
1395
1396 let added = synth(300, 99);
1397 let report = append_files(&archive, &added, 3).unwrap();
1398 assert_eq!(report.rows_before, orig.len() as u64, "must recover all original rows");
1399 assert_eq!(report.rows_added, added.len() as u64);
1400 assert!(report.blob_bytes_added > 0, "append must write new blob bytes");
1401 assert!(
1402 report.sealed_total_bytes > report.blob_append_offset,
1403 "re-sealed file must be larger than the old blob region"
1404 );
1405
1406 let ar = ZnippyArchive::open(&archive).unwrap();
1409 let mut listed = ar.list_files().unwrap();
1410 listed.sort();
1411 let mut expected: Vec<String> =
1412 orig.iter().chain(added.iter()).map(|(p, _)| p.clone()).collect();
1413 expected.sort();
1414 assert_eq!(listed, expected, "index must list exactly old+new files after append");
1415
1416 for (p, bytes) in orig.iter().chain(added.iter()) {
1417 let got = ar.extract_file(p).unwrap();
1418 assert_eq!(&got, bytes, "byte mismatch after append for {p}");
1419 }
1420
1421 let probe = &added[123].0;
1423 let chunks = crate::locate_file(&archive, probe).unwrap();
1424 assert!(!chunks.is_empty(), "appended file must be locatable via the re-sealed lookup");
1425 let _ = std::fs::remove_dir_all(&dir);
1426 }
1427
1428 #[test]
1433 fn create_archive_seeds_then_grows() {
1434 let dir = unique_dir("create");
1435 let seed = synth(40, 5);
1436
1437 let made = dir.join("made.znippy");
1439 let report = create_archive(&made, &seed, 3).unwrap();
1440 assert_eq!(report.rows_before, 0, "fresh archive has no prior rows");
1441 assert_eq!(report.rows_added, seed.len() as u64);
1442
1443 let ref_path = dir.join("ref.znippy");
1444 write_fresh(&ref_path, &seed, ArrowIpcSinkAppend::new);
1445 assert_eq!(
1446 std::fs::read(&made).unwrap(),
1447 std::fs::read(&ref_path).unwrap(),
1448 "create_archive must be byte-identical to the proven fresh write path"
1449 );
1450
1451 let ar = ZnippyArchive::open(&made).unwrap();
1453 for (p, bytes) in &seed {
1454 assert_eq!(&ar.extract_file(p).unwrap(), bytes, "seed byte mismatch for {p}");
1455 }
1456
1457 let added = synth(15, 88);
1459 let rep2 = append_files(&made, &added, 3).unwrap();
1460 assert_eq!(rep2.rows_before, seed.len() as u64, "append must recover seeded rows");
1461 assert_eq!(rep2.rows_added, added.len() as u64);
1462
1463 let ar2 = ZnippyArchive::open(&made).unwrap();
1464 for (p, bytes) in seed.iter().chain(added.iter()) {
1465 assert_eq!(&ar2.extract_file(p).unwrap(), bytes, "byte mismatch after grow for {p}");
1466 }
1467 let _ = std::fs::remove_dir_all(&dir);
1468 }
1469
1470 #[test]
1475 fn create_archive_to_vec_is_filesystem_free_and_round_trips() {
1476 let files = synth(24, 7);
1477 let (bytes, report) = create_archive_to_vec(&files, 3).unwrap();
1478 assert_eq!(report.rows_before, 0, "fresh in-memory archive has no prior rows");
1479 assert_eq!(report.rows_added, files.len() as u64);
1480 assert_eq!(bytes.len() as u64, report.sealed_total_bytes, "vec len == sealed size");
1481
1482 let dir = unique_dir("tovec");
1483 let ref_path = dir.join("ref.znippy");
1485 create_archive(&ref_path, &files, 3).unwrap();
1486 assert_eq!(
1487 bytes,
1488 std::fs::read(&ref_path).unwrap(),
1489 "in-memory archive must be byte-identical to create_archive's file output"
1490 );
1491 let p = dir.join("from_mem.znippy");
1493 std::fs::write(&p, &bytes).unwrap();
1494 let ar = ZnippyArchive::open(&p).unwrap();
1495 for (name, content) in &files {
1496 assert_eq!(&ar.extract_file(name).unwrap(), content, "byte mismatch for {name}");
1497 }
1498 let _ = std::fs::remove_dir_all(&dir);
1499 }
1500
1501 fn recorded_format_version(path: &Path) -> Option<String> {
1508 use std::io::{Read, Seek, SeekFrom};
1509
1510 use arrow::ipc::reader::StreamReader;
1511
1512 let entries = crate::index::read_znippy_manifest(path).ok()?;
1513 let mut file = File::open(path).ok()?;
1514 for e in &entries {
1515 if is_reserved_module(&e.module_name) {
1516 continue;
1517 }
1518 file.seek(SeekFrom::Start(e.index_offset)).ok()?;
1519 let mut bytes = vec![0u8; e.index_len as usize];
1520 file.read_exact(&mut bytes).ok()?;
1521 let reader = StreamReader::try_new(std::io::Cursor::new(bytes), None).ok()?;
1522 return reader
1523 .schema()
1524 .metadata()
1525 .get(crate::index::FORMAT_VERSION_KEY)
1526 .cloned();
1527 }
1528 None
1529 }
1530
1531 #[test]
1543 fn every_appended_archive_records_the_format_version() {
1544 let dir = unique_dir("fmtver");
1545 let want = crate::index::ZNIPPY_FORMAT_VERSION.to_string();
1546
1547 let made = dir.join("made.znippy");
1549 create_archive(&made, &synth(12, 3), 3).unwrap();
1550 assert_eq!(
1551 recorded_format_version(&made).as_deref(),
1552 Some(want.as_str()),
1553 "create_archive must stamp the on-disk format version"
1554 );
1555
1556 append_files(&made, &synth(7, 91), 3).unwrap();
1559 assert_eq!(
1560 recorded_format_version(&made).as_deref(),
1561 Some(want.as_str()),
1562 "append_files must re-stamp the format version on the re-sealed archive"
1563 );
1564
1565 let (bytes, _) = create_archive_to_vec(&synth(9, 4), 3).unwrap();
1567 let mem = dir.join("mem.znippy");
1568 std::fs::write(&mem, &bytes).unwrap();
1569 assert_eq!(
1570 recorded_format_version(&mem).as_deref(),
1571 Some(want.as_str()),
1572 "create_archive_to_vec must stamp the on-disk format version"
1573 );
1574
1575 let ar = ZnippyArchive::open(&mem).unwrap();
1577 for (p, body) in &synth(9, 4) {
1578 assert_eq!(&ar.extract_file(p).unwrap(), body, "byte mismatch for {p}");
1579 }
1580 let _ = std::fs::remove_dir_all(&dir);
1581 }
1582
1583 #[test]
1593 fn metadata_is_searchable_survives_a_reseal_and_absence_is_reported_as_absence() {
1594 use crate::meta_index::{ArchiveMeta, MetaSearch, MetaTable, MetaValue, read_archive_meta};
1595
1596 let dir = unique_dir("meta");
1597 let files = synth(30, 2);
1598
1599 let plain = dir.join("plain.znippy");
1601 create_archive(&plain, &files, 3).unwrap();
1602 let m = read_archive_meta(&plain).unwrap();
1603 assert_eq!(m, ArchiveMeta::NoMetadata, "an archive with no meta section must say so");
1604 assert!(!m.is_searchable());
1605 assert_eq!(m.find_by_key("build-thing"), MetaSearch::NoMetadata);
1606 assert!(
1607 m.find_by_key("build-thing").hits().is_none(),
1608 "absence must NOT present itself as an empty result set"
1609 );
1610
1611 let empty = dir.join("empty.znippy");
1613 create_archive_with_meta(&empty, &files, 3, Some(MetaTable::new())).unwrap();
1614 let me = read_archive_meta(&empty).unwrap();
1615 assert!(me.is_searchable(), "a sealed empty index WAS searched");
1616 assert!(me.index().is_some_and(|i| i.is_empty()));
1617 assert_eq!(me.find_by_key("build-thing"), MetaSearch::Hits(&[]));
1618
1619 let wasm = b"\0asm\x01\0\0\0".to_vec();
1621 let (p0, p1, p2) = (files[0].0.clone(), files[1].0.clone(), files[2].0.clone());
1622 let mut t = MetaTable::new();
1623 t.insert(p0.clone(), "build-thing", MetaValue::Bytes(wasm.clone()))
1624 .insert(p0.clone(), "build-thing.abi", "wasi-p2")
1625 .insert(p1.clone(), "build-thing", MetaValue::Bytes(wasm.clone()))
1626 .insert(p2.clone(), "coverage", 0.5f64)
1627 .insert_archive("producer", "znippy");
1628
1629 let ar = dir.join("meta.znippy");
1630 create_archive_with_meta(&ar, &files, 3, Some(t)).unwrap();
1631
1632 let m = read_archive_meta(&ar).unwrap();
1633 let idx = m.index().expect("sealed index is present");
1634 assert_eq!(idx.len(), 5);
1635 let hits = m.find_by_key("build-thing").hits().unwrap();
1636 assert_eq!(hits.len(), 2, "exactly the two entries that carry one");
1637 let mut got: Vec<&str> = hits.iter().filter_map(|h| h.path()).collect();
1638 got.sort();
1639 let mut want = vec![p0.as_str(), p1.as_str()];
1640 want.sort();
1641 assert_eq!(got, want, "the search names the right ENTRIES");
1642 assert_eq!(hits[0].value.as_bytes(), Some(&wasm[..]), "and the right VALUE");
1643 assert_eq!(idx.archive_value("producer").and_then(MetaValue::as_str), Some("znippy"));
1644 assert_eq!(
1645 idx.find_by_prefix("build-thing").len(),
1646 3,
1647 "prefix sweeps build-thing + build-thing.abi"
1648 );
1649
1650 let reader = ZnippyArchive::open(&ar).unwrap();
1652 let want_bytes = &files.iter().find(|(p, _)| *p == p0).unwrap().1;
1653 assert_eq!(&reader.extract_file(hits[0].path().unwrap()).unwrap(), want_bytes);
1654
1655 let added = synth(6, 77);
1657 let mut more = MetaTable::new();
1658 more.insert(added[0].0.clone(), "build-thing", MetaValue::Bytes(wasm.clone()));
1659 append_files_with_meta(&ar, &added, 3, Some(more)).unwrap();
1660
1661 let m2 = read_archive_meta(&ar).unwrap();
1662 let hits2 = m2.find_by_key("build-thing").hits().unwrap();
1663 assert_eq!(hits2.len(), 3, "the append merged, it did not replace");
1664 assert_eq!(
1665 m2.index().unwrap().archive_value("producer").and_then(MetaValue::as_str),
1666 Some("znippy"),
1667 "the archive-level row survived the re-seal"
1668 );
1669
1670 append_files(&ar, &synth(3, 91), 3).unwrap();
1672 assert_eq!(
1673 read_archive_meta(&ar).unwrap().find_by_key("build-thing").hits().unwrap().len(),
1674 3,
1675 "an append that says nothing about metadata must not erase it"
1676 );
1677
1678 let ar2 = ZnippyArchive::open(&ar).unwrap();
1681 let n_files = files.len() + added.len() + 3;
1682 assert_eq!(
1683 ar2.file_count(),
1684 n_files,
1685 "metadata rows leaked into the data index"
1686 );
1687 let lookup_bytes = read_reserved_section_bytes(&ar, LOOKUP_MODULE).unwrap().unwrap();
1688 assert_eq!(
1689 decode_base_rows(&lookup_bytes).unwrap().0.len(),
1690 n_files,
1691 "metadata rows leaked into the random-access lookup"
1692 );
1693 for (p, bytes) in files.iter().chain(added.iter()) {
1694 assert_eq!(&ar2.extract_file(p).unwrap(), bytes, "byte mismatch for {p}");
1695 }
1696 let _ = std::fs::remove_dir_all(&dir);
1697 }
1698
1699 #[test]
1724 fn an_append_carries_the_ref_log_forward_and_drops_the_derived_sections() {
1725 use crate::index::{
1726 GUNNAR_GRAPH_MODULE, GUNNAR_REFS_MODULE, GUNNAR_SECRETS_MODULE,
1727 };
1728 use crate::meta_sink::{ReservedSection, ReservedSectionBuilder};
1729
1730 let dir = unique_dir("carried_reserved");
1731 let archive = dir.join("a.znippy");
1732
1733 let refs_bytes = b"refs/heads/main 0123456789abcdef -- push 1".to_vec();
1734 let secrets_bytes = b"\x00ciphertext-only, never plaintext".to_vec();
1735 let graph_bytes = b"a commit graph derived from the objects".to_vec();
1736
1737 {
1740 let f = Arc::new(File::create(&archive).unwrap());
1741 let mut cursor = 0u64;
1742 let mut blobs = Vec::new();
1743 for (i, (name, bytes)) in [("obj/a.bin", b"first".to_vec())].iter().enumerate() {
1744 use std::os::unix::fs::FileExt;
1745 f.write_all_at(bytes, cursor).unwrap();
1746 blobs.push(BlobMeta {
1747 blob_offset: cursor,
1748 blob_size: bytes.len() as u64,
1749 chunk_meta: ChunkMeta {
1750 fdata_offset: 0,
1751 file_index: i as u64,
1752 chunk_seq: 0,
1753 checksum: *blake3::hash(bytes).as_bytes(),
1754 compressed: false,
1755 uncompressed_size: bytes.len() as u64,
1756 compressed_size: bytes.len() as u64,
1757 },
1758 });
1759 cursor += bytes.len() as u64;
1760 let _ = name;
1761 }
1762 let batch = crate::index::build_metadata_batch(
1763 &blobs,
1764 |_fi: u64| "obj/a.bin".to_string(),
1765 &[],
1766 &[],
1767 )
1768 .unwrap();
1769 let mut sink = ArrowIpcSink::new(Arc::clone(&f), cursor);
1772 let (r, s, g) = (refs_bytes.clone(), secrets_bytes.clone(), graph_bytes.clone());
1773 let builder: ReservedSectionBuilder = Box::new(move |_lookup| {
1774 Ok(vec![
1775 ReservedSection::raw(GUNNAR_REFS_MODULE, r),
1776 ReservedSection::raw(GUNNAR_SECRETS_MODULE, s),
1777 ReservedSection::raw(GUNNAR_GRAPH_MODULE, g),
1778 ])
1779 });
1780 sink = sink.with_reserved_builder(builder);
1781 sink.push_subindex(
1782 crate::index::lookup_schema().as_ref(),
1783 &[batch],
1784 GroupKey { pkg_type: 0, repo: String::new(), module_name: String::new() },
1785 )
1786 .unwrap();
1787 Box::new(sink).finish().unwrap();
1788 }
1789
1790 for m in [GUNNAR_REFS_MODULE, GUNNAR_SECRETS_MODULE, GUNNAR_GRAPH_MODULE] {
1792 assert!(
1793 read_reserved_section_bytes(&archive, m).unwrap().is_some(),
1794 "{m} must be present BEFORE the append"
1795 );
1796 }
1797
1798 append_files(&archive, &[("obj/b.bin".to_string(), b"second".to_vec())], 3).unwrap();
1800
1801 assert_eq!(
1803 read_reserved_section_bytes(&archive, GUNNAR_REFS_MODULE).unwrap(),
1804 Some(refs_bytes),
1805 "an object-carrying push must not erase the ref log — this is the 2026-08-04 bug"
1806 );
1807 assert_eq!(
1808 read_reserved_section_bytes(&archive, GUNNAR_SECRETS_MODULE).unwrap(),
1809 Some(secrets_bytes),
1810 );
1811 assert_eq!(
1813 read_reserved_section_bytes(&archive, GUNNAR_GRAPH_MODULE).unwrap(),
1814 None,
1815 "a commit graph derived from the old object set must NOT be carried forward"
1816 );
1817
1818 let reader = ZnippyArchive::open(&archive).unwrap();
1820 assert_eq!(reader.extract_file("obj/a.bin").unwrap(), b"first".to_vec());
1821 assert_eq!(reader.extract_file("obj/b.bin").unwrap(), b"second".to_vec());
1822
1823 let _ = std::fs::remove_dir_all(&dir);
1824 }
1825
1826
1827 #[test]
1835 fn a_superseded_entry_becomes_a_delta_and_reads_back_exactly() {
1836 let dir = unique_dir("supersede");
1837 let archive = dir.join("a.znippy");
1838
1839 let mut st = 0x243f_6a88_85a3_08d3u64;
1842 let gen_n: Vec<u8> = (0..400_000u32)
1843 .map(|_| {
1844 st = st.wrapping_add(0x9e37_79b9_7f4a_7c15);
1845 let mut z = st;
1846 z = (z ^ (z >> 30)).wrapping_mul(0xbf58_476d_1ce4_e5b9);
1847 (z ^ (z >> 27)) as u8
1848 })
1849 .collect();
1850 let mut gen_n1 = gen_n.clone();
1852 gen_n1.extend_from_slice(&gen_n[..20_000]);
1853
1854 create_archive(
1855 &archive,
1856 &[
1857 ("pack-N.pack".to_string(), gen_n.clone()),
1858 ("pack-N1.pack".to_string(), gen_n1.clone()),
1859 ],
1860 3,
1861 )
1862 .unwrap();
1863 let before = std::fs::metadata(&archive).unwrap().len();
1864
1865 let out = supersede_as_delta(&archive, "pack-N.pack", "pack-N1.pack", 0, 3).unwrap();
1866 let (stored, delta_bytes) = match out {
1867 SupersedeOutcome::Delta { stored_bytes, delta_bytes, chain_depth } => {
1868 assert_eq!(chain_depth, 1);
1869 (stored_bytes, delta_bytes)
1870 }
1871 other => panic!("expected a delta, got {other:?}"),
1872 };
1873 assert_eq!(stored, gen_n.len() as u64);
1874 assert!(
1875 delta_bytes * 20 < stored,
1876 "N is a near-prefix of N+1, so the delta must be a small fraction of it; \
1877 got {delta_bytes} against {stored}"
1878 );
1879
1880 let ar = ZnippyArchive::open(&archive).unwrap();
1883 assert_eq!(ar.extract_file("pack-N.pack").unwrap(), gen_n, "the superseded generation");
1884 assert_eq!(ar.extract_file("pack-N1.pack").unwrap(), gen_n1, "the live generation");
1885 assert_eq!(ar.extract_file_verified("pack-N.pack").unwrap(), gen_n);
1886 assert_eq!(ar.file_size("pack-N.pack"), Some(gen_n.len() as u64));
1887
1888 assert!(
1892 read_delta_map(&archive)
1893 .unwrap()
1894 .iter()
1895 .all(|(p, _, _)| p == "pack-N.pack"),
1896 "only the superseded generation may be a delta"
1897 );
1898
1899 let _ = before;
1900 let _ = std::fs::remove_dir_all(&dir);
1901 }
1902
1903 #[test]
1905 fn the_writer_refuses_a_pointless_delta_and_an_over_long_chain() {
1906 let dir = unique_dir("supersede_refuse");
1907 let archive = dir.join("a.znippy");
1908 let mut st = 0x1234_5678_9abc_def0u64;
1909 let noise = |n: usize, st: &mut u64| -> Vec<u8> {
1910 (0..n)
1911 .map(|_| {
1912 *st = st.wrapping_add(0x9e37_79b9_7f4a_7c15);
1913 let mut z = *st;
1914 z = (z ^ (z >> 30)).wrapping_mul(0xbf58_476d_1ce4_e5b9);
1915 (z ^ (z >> 27)) as u8
1916 })
1917 .collect()
1918 };
1919 let a = noise(60_000, &mut st);
1920 let b = noise(60_000, &mut st); create_archive(
1922 &archive,
1923 &[("a.pack".to_string(), a.clone()), ("b.pack".to_string(), b.clone())],
1924 3,
1925 )
1926 .unwrap();
1927
1928 match supersede_as_delta(&archive, "a.pack", "b.pack", 0, 3).unwrap() {
1929 SupersedeOutcome::NotSmaller { .. } => {}
1930 other => panic!("unrelated bytes must not produce a delta, got {other:?}"),
1931 }
1932 let ar = ZnippyArchive::open(&archive).unwrap();
1934 assert_eq!(ar.extract_file("a.pack").unwrap(), a);
1935 assert!(read_delta_map(&archive).unwrap().is_empty());
1936
1937 assert_eq!(
1940 supersede_as_delta(&archive, "a.pack", "b.pack", MAX_GENERATION_CHAIN, 3).unwrap(),
1941 SupersedeOutcome::ChainTooLong
1942 );
1943 assert!(supersede_as_delta(&archive, "a.pack", "a.pack", 0, 3).is_err());
1944
1945 let _ = std::fs::remove_dir_all(&dir);
1946 }
1947
1948 #[test]
1956 fn a_generation_chain_reads_every_generation_back_exactly() {
1957 let dir = unique_dir("genchain");
1958 let archive = dir.join("g.znippy");
1959 let mut st = 0x0f0f_0f0f_dead_beefu64;
1960 let mut gens: Vec<Vec<u8>> = Vec::new();
1961 let mut cur: Vec<u8> = (0..120_000u32)
1962 .map(|_| {
1963 st = st.wrapping_add(0x9e37_79b9_7f4a_7c15);
1964 let mut z = st;
1965 z = (z ^ (z >> 30)).wrapping_mul(0xbf58_476d_1ce4_e5b9);
1966 (z ^ (z >> 27)) as u8
1967 })
1968 .collect();
1969 gens.push(cur.clone());
1970 for g in 1..4 {
1971 cur.extend_from_slice(format!("generation {g} tail ").repeat(200).as_bytes());
1972 gens.push(cur.clone());
1973 }
1974 let files: Vec<(String, Vec<u8>)> = gens
1975 .iter()
1976 .enumerate()
1977 .map(|(i, b)| (format!("pack-{i}.pack"), b.clone()))
1978 .collect();
1979 create_archive(&archive, &files, 3).unwrap();
1980
1981 for i in (0..3).rev() {
1984 let out = supersede_as_delta(
1985 &archive,
1986 &format!("pack-{i}.pack"),
1987 &format!("pack-{}.pack", i + 1),
1988 3 - 1 - i,
1989 3,
1990 )
1991 .unwrap();
1992 assert!(matches!(out, SupersedeOutcome::Delta { .. }), "gen {i}: {out:?}");
1993 }
1994
1995 let ar = ZnippyArchive::open(&archive).unwrap();
1996 for (i, want) in gens.iter().enumerate() {
1997 assert_eq!(
1998 &ar.extract_file(&format!("pack-{i}.pack")).unwrap(),
1999 want,
2000 "generation {i} at chain depth {}",
2001 3 - i
2002 );
2003 }
2004 assert!(
2006 read_delta_map(&archive)
2007 .unwrap()
2008 .iter()
2009 .all(|(p, _, _)| p != "pack-3.pack"),
2010 "the live generation must stay whole"
2011 );
2012 let _ = std::fs::remove_dir_all(&dir);
2013 }
2014
2015
2016 #[test]
2021 #[ignore]
2022 fn perf_real_generation_chain() {
2023 let Ok(d) = std::env::var("ZNIPPY_CHAIN_DIR") else { return };
2024 let dir = unique_dir("realchain");
2025 let archive = dir.join("g.znippy");
2026 let mut files: Vec<(String, Vec<u8>)> = Vec::new();
2027 for i in 0.. {
2028 let p = std::path::Path::new(&d).join(format!("pack-{i}.pack"));
2029 if !p.is_file() {
2030 break;
2031 }
2032 files.push((format!("pack-{i}.pack"), std::fs::read(&p).unwrap()));
2033 }
2034 let n = files.len();
2035 let raw: u64 = files.iter().map(|(_, b)| b.len() as u64).sum();
2036
2037 let t0 = std::time::Instant::now();
2038 create_archive(&archive, &files, 3).unwrap();
2039 let seal_ms = t0.elapsed().as_secs_f64() * 1e3;
2040 let sealed = std::fs::metadata(&archive).unwrap().len();
2041
2042 println!("gen,stored_b,delta_b,ratio_x,supersede_ms");
2044 let mut encode_ms_total = 0.0;
2045 let mut delta_total = 0u64;
2046 for i in (0..n - 1).rev() {
2047 let t = std::time::Instant::now();
2048 let out = supersede_as_delta(
2049 &archive,
2050 &format!("pack-{i}.pack"),
2051 &format!("pack-{}.pack", i + 1),
2052 n - 2 - i,
2053 3,
2054 )
2055 .unwrap();
2056 let ms = t.elapsed().as_secs_f64() * 1e3;
2057 encode_ms_total += ms;
2058 match out {
2059 SupersedeOutcome::Delta { stored_bytes, delta_bytes, .. } => {
2060 delta_total += delta_bytes;
2061 println!(
2062 "{i},{stored_bytes},{delta_bytes},{:.2},{ms:.0}",
2063 stored_bytes as f64 / delta_bytes as f64
2064 );
2065 }
2066 other => println!("{i},-,-,-,{ms:.0} ({other:?})"),
2067 }
2068 }
2069
2070 let ar = ZnippyArchive::open(&archive).unwrap();
2073 let mut live: Vec<(String, Vec<u8>)> = Vec::new();
2074 let t = std::time::Instant::now();
2075 for i in 0..n {
2076 let name = format!("pack-{i}.pack");
2077 let got = ar.extract_file(&name).unwrap();
2078 assert_eq!(got, files[i].1, "generation {i} did not read back exactly");
2079 live.push((name, got));
2080 }
2081 let read_all_ms = t.elapsed().as_secs_f64() * 1e3;
2082
2083 println!("gen,depth,read_ms");
2085 for i in 0..n {
2086 let name = format!("pack-{i}.pack");
2087 let t = std::time::Instant::now();
2088 for _ in 0..5 {
2089 let _ = ar.extract_file(&name).unwrap();
2090 }
2091 println!("{i},{},{:.2}", n - 1 - i, t.elapsed().as_secs_f64() * 1e3 / 5.0);
2092 }
2093
2094 let compacted = dir.join("c.znippy");
2095 let mut packed: Vec<(String, Vec<u8>)> = Vec::new();
2096 for i in 0..n {
2097 packed.push(live[i].clone());
2098 }
2099 let overhead = sealed - raw;
2102 let steady = files[n - 1].1.len() as u64 + delta_total + overhead;
2103 println!(
2104 "SUMMARY generations={n} raw={raw} sealed={sealed} live_whole={} deltas={delta_total} \
2105 steady_state={steady} saving_x={:.2} seal_ms={seal_ms:.0} \
2106 supersede_ms_total={encode_ms_total:.0} read_all_ms={read_all_ms:.0}",
2107 files[n - 1].1.len(),
2108 raw as f64 / steady as f64
2109 );
2110 let _ = (compacted, packed);
2111 let _ = std::fs::remove_dir_all(&dir);
2112 }
2113
2114 #[test]
2115 fn archive_reader_matches_free_functions() {
2116 let dir = unique_dir("reader");
2117 let archive = dir.join("store.znippy");
2118 let files = synth(600, 13);
2119 write_fresh(&archive, &files, ArrowIpcSink::new);
2120
2121 let reader = crate::ArchiveReader::open(&archive).unwrap();
2122 assert_eq!(reader.row_count(), files.len(), "one chunk per synth file");
2123
2124 for (p, _) in files.iter().step_by(37) {
2126 let cached = reader.locate(p);
2127 let free = crate::locate_file(&archive, p).unwrap();
2128 assert!(!cached.is_empty(), "cached reader failed to locate {p}");
2129 assert_eq!(cached, free, "cached vs free locate diverged for {p}");
2130 }
2131
2132 let missing = "repo/does/not/exist.bin";
2134 assert!(reader.locate(missing).is_empty());
2135 assert!(crate::locate_file(&archive, missing).unwrap().is_empty());
2136
2137 assert_eq!(
2139 reader.files_meta(),
2140 crate::get_all_files_meta(&archive).unwrap(),
2141 "cached files_meta diverged from get_all_files_meta"
2142 );
2143
2144 for prefix in ["repo/grp001/", "repo/grp0", "repo/", ""] {
2146 assert_eq!(
2147 reader.files_meta_with_prefix(prefix),
2148 crate::get_files_meta_with_prefix(&archive, prefix).unwrap(),
2149 "cached prefix meta diverged for {prefix:?}"
2150 );
2151 }
2152
2153 let _ = std::fs::remove_dir_all(&dir);
2154 }
2155
2156 #[test]
2174 fn compaction_reclaims_the_superseded_blob_and_changes_no_entry() {
2175 let dir = unique_dir("compact");
2176 let archive = dir.join("c.znippy");
2177
2178 let mut st = 0x5151_2323_abcd_ef01u64;
2179 let base: Vec<u8> = (0..600_000u32)
2180 .map(|_| {
2181 st = st.wrapping_add(0x9e37_79b9_7f4a_7c15);
2182 let mut z = st;
2183 z = (z ^ (z >> 30)).wrapping_mul(0xbf58_476d_1ce4_e5b9);
2184 (z ^ (z >> 27)) as u8
2185 })
2186 .collect();
2187 let mut gens: Vec<Vec<u8>> = vec![base.clone()];
2188 for g in 1..4 {
2189 let mut next = gens[g - 1].clone();
2190 next.extend_from_slice(format!("generation {g} tail ").repeat(300).as_bytes());
2191 gens.push(next);
2192 }
2193 let files: Vec<(String, Vec<u8>)> = gens
2194 .iter()
2195 .enumerate()
2196 .map(|(i, b)| (format!("pack-{i}.pack"), b.clone()))
2197 .collect();
2198 create_archive(&archive, &files, 3).unwrap();
2199 let raw: u64 = gens.iter().map(|g| g.len() as u64).sum();
2200 let whole = std::fs::metadata(&archive).unwrap().len();
2201
2202 for i in (0..3).rev() {
2203 let out = supersede_as_delta(
2204 &archive,
2205 &format!("pack-{i}.pack"),
2206 &format!("pack-{}.pack", i + 1),
2207 3 - 1 - i,
2208 3,
2209 )
2210 .unwrap();
2211 assert!(matches!(out, SupersedeOutcome::Delta { .. }), "gen {i}: {out:?}");
2212 }
2213
2214 let after_supersede = std::fs::metadata(&archive).unwrap().len();
2217 assert!(
2218 after_supersede >= whole,
2219 "supersede shrank the file ({whole} -> {after_supersede}); if that is now true, this \
2220 test and everything built on the dead-payload finding wants revisiting"
2221 );
2222
2223 let report = compact_archive(&archive).unwrap();
2224 let compacted = std::fs::metadata(&archive).unwrap().len();
2225 assert_eq!(report.bytes_after, compacted);
2226 assert_eq!(report.rows, 4, "a compaction must not change the row count");
2227 assert_eq!(report.delta_rows, 3, "the delta map must travel across it");
2228 assert!(
2229 compacted * 2 < raw,
2230 "the compacted archive is {compacted} bytes for {raw} bytes of generations — the dead \
2231 payload was not reclaimed"
2232 );
2233
2234 let ar = ZnippyArchive::open(&archive).unwrap();
2236 for (i, want) in gens.iter().enumerate() {
2237 assert_eq!(
2238 &ar.extract_file(&format!("pack-{i}.pack")).unwrap(),
2239 want,
2240 "generation {i} did not survive the compaction"
2241 );
2242 }
2243 let map = read_delta_map(&archive).unwrap();
2245 assert!(
2246 map.iter().all(|(p, _, _)| p != "pack-3.pack"),
2247 "compaction must not put the live generation behind a link: {map:?}"
2248 );
2249
2250 let again = compact_archive(&archive).unwrap();
2253 assert_eq!(again.rows, 4);
2254 assert_eq!(again.delta_rows, 3);
2255 let ar = ZnippyArchive::open(&archive).unwrap();
2256 for (i, want) in gens.iter().enumerate() {
2257 assert_eq!(&ar.extract_file(&format!("pack-{i}.pack")).unwrap(), want);
2258 }
2259
2260 let _ = std::fs::remove_dir_all(&dir);
2261 }
2262
2263 #[test]
2267 fn compaction_carries_the_reserved_logs() {
2268 use crate::index::GUNNAR_REFS_MODULE;
2269 use crate::meta_sink::{ReservedSection, ReservedSectionBuilder};
2270
2271 let dir = unique_dir("compact_reserved");
2272 let archive = dir.join("r.znippy");
2273 let refs_bytes = b"refs/heads/main 0123456789abcdef".to_vec();
2274
2275 {
2276 use std::os::unix::fs::FileExt;
2277 let f = Arc::new(File::create(&archive).unwrap());
2278 let payload = b"an entry".to_vec();
2279 f.write_all_at(&payload, 0).unwrap();
2280 let blobs = vec![BlobMeta {
2281 blob_offset: 0,
2282 blob_size: payload.len() as u64,
2283 chunk_meta: ChunkMeta {
2284 fdata_offset: 0,
2285 file_index: 0,
2286 chunk_seq: 0,
2287 checksum: *blake3::hash(&payload).as_bytes(),
2288 compressed: false,
2289 uncompressed_size: payload.len() as u64,
2290 compressed_size: payload.len() as u64,
2291 },
2292 }];
2293 let batch =
2294 crate::index::build_metadata_batch(&blobs, |_| "obj/a.bin".to_string(), &[], &[])
2295 .unwrap();
2296 let mut sink = ArrowIpcSink::new(Arc::clone(&f), payload.len() as u64);
2297 let carried = refs_bytes.clone();
2298 let builder: ReservedSectionBuilder =
2299 Box::new(move |_lookup| Ok(vec![ReservedSection::raw(GUNNAR_REFS_MODULE, carried)]));
2300 sink = sink.with_reserved_builder(builder);
2301 sink.push_subindex(
2302 crate::index::lookup_schema().as_ref(),
2303 &[batch],
2304 GroupKey {
2305 pkg_type: 0,
2306 repo: String::new(),
2307 module_name: String::new(),
2308 },
2309 )
2310 .unwrap();
2311 Box::new(sink).finish().unwrap();
2312 }
2313 assert_eq!(
2314 read_reserved_section_bytes(&archive, GUNNAR_REFS_MODULE).unwrap(),
2315 Some(refs_bytes.clone())
2316 );
2317
2318 compact_archive(&archive).unwrap();
2319 assert_eq!(
2320 read_reserved_section_bytes(&archive, GUNNAR_REFS_MODULE).unwrap(),
2321 Some(refs_bytes),
2322 "the compaction dropped __gunnar_refs__"
2323 );
2324 let _ = std::fs::remove_dir_all(&dir);
2325 }
2326}
2327
2328pub(crate) fn decode_base_rows(bytes: &[u8]) -> Result<(Vec<String>, Vec<ChunkLoc>)> {
2332 use arrow::ipc::reader::StreamReader;
2333
2334 let reader = StreamReader::try_new(std::io::Cursor::new(bytes), None)
2335 .map_err(|e| anyhow!("append: lookup ipc reader: {e}"))?;
2336 let mut paths = Vec::new();
2337 let mut locs = Vec::new();
2338 for batch in reader {
2339 let batch = batch.map_err(|e| anyhow!("append: lookup batch decode: {e}"))?;
2340 let get = |n: &str| batch.column_by_name(n)
2341 .ok_or_else(|| anyhow!("append: lookup missing column {n}"));
2342 let p = get("relative_path")?.as_any().downcast_ref::<StringArray>()
2343 .ok_or_else(|| anyhow!("relative_path type"))?;
2344 let seq = get("chunk_seq")?.as_any().downcast_ref::<UInt32Array>()
2345 .ok_or_else(|| anyhow!("chunk_seq type"))?;
2346 let fdata = get("fdata_offset")?.as_any().downcast_ref::<UInt64Array>()
2347 .ok_or_else(|| anyhow!("fdata_offset type"))?;
2348 let comp = get("compressed")?.as_any().downcast_ref::<BooleanArray>()
2349 .ok_or_else(|| anyhow!("compressed type"))?;
2350 let usz = get("uncompressed_size")?.as_any().downcast_ref::<UInt64Array>()
2351 .ok_or_else(|| anyhow!("uncompressed_size type"))?;
2352 let boff = get("blob_offset")?.as_any().downcast_ref::<UInt64Array>()
2353 .ok_or_else(|| anyhow!("blob_offset type"))?;
2354 let bsz = get("blob_size")?.as_any().downcast_ref::<UInt64Array>()
2355 .ok_or_else(|| anyhow!("blob_size type"))?;
2356 let ck = get("checksum")?.as_any().downcast_ref::<FixedSizeBinaryArray>()
2357 .ok_or_else(|| anyhow!("checksum type"))?;
2358 for i in 0..batch.num_rows() {
2359 let mut c = [0u8; 32];
2360 c.copy_from_slice(ck.value(i));
2361 paths.push(p.value(i).to_string());
2362 locs.push(ChunkLoc {
2363 chunk_seq: seq.value(i),
2364 fdata_offset: fdata.value(i),
2365 blob_offset: boff.value(i),
2366 blob_size: bsz.value(i),
2367 uncompressed_size: usz.value(i),
2368 compressed: comp.value(i),
2369 checksum: c,
2370 });
2371 }
2372 }
2373 Ok((paths, locs))
2374}