use std::fs::File;
use std::os::unix::fs::FileExt;
use std::path::Path;
use std::sync::Arc;
use anyhow::{Result, anyhow};
use arrow::array::{
BooleanArray, BooleanBuilder, FixedSizeBinaryArray, FixedSizeBinaryBuilder, StringArray,
StringBuilder, UInt32Array, UInt32Builder, UInt64Array, UInt64Builder,
};
use arrow::datatypes::Schema;
use arrow::ipc::writer::StreamWriter;
use arrow::record_batch::RecordBatch;
use crate::index::{
ChunkLoc, LOOKUP_MODULE, MULTI_INDEX_MAGIC, ManifestEntry, RESERVED_PKG_TYPE, TRIE_MODULE,
data_subindex_schema, is_reserved_module, lookup_schema, read_znippy_full_manifest,
write_manifest_bytes,
};
use crate::index::{
CARRIED_RESERVED_MODULES, META_MODULE, ZNIPPY_DELTA_MODULE, read_reserved_section_bytes,
};
use crate::meta_index::{
MetaTable, build_meta_batch, decode_meta_section, meta_schema,
};
use crate::meta_sink::{ArchiveMetaSink, GroupKey};
pub struct ArrowIpcSinkAppend {
file: Arc<File>,
cursor: u64,
entries: Vec<ManifestEntry>,
lookup_paths: Vec<String>,
lookup_locs: Vec<ChunkLoc>,
carried: Vec<(String, ChunkLoc)>,
meta: Option<MetaTable>,
carried_reserved: Vec<(String, Vec<u8>)>,
delta_map: Vec<(String, u32, String)>,
}
impl ArrowIpcSinkAppend {
pub fn new(file: Arc<File>, blob_end_offset: u64) -> Self {
Self {
file,
cursor: blob_end_offset,
entries: Vec::new(),
lookup_paths: Vec::new(),
lookup_locs: Vec::new(),
carried: Vec::new(),
meta: None,
carried_reserved: Vec::new(),
delta_map: Vec::new(),
}
}
pub fn with_meta(mut self, meta: MetaTable) -> Self {
self.meta = Some(meta);
self
}
pub fn merge_meta(&mut self, rows: impl IntoIterator<Item = crate::meta_index::MetaEntry>) {
self.meta.get_or_insert_with(MetaTable::new).extend(rows);
}
pub fn meta(&self) -> Option<&MetaTable> {
self.meta.as_ref()
}
pub fn open_existing(path: &Path) -> Result<Self> {
let file = Arc::new(
std::fs::OpenOptions::new()
.read(true)
.write(true)
.open(path)
.map_err(|e| anyhow!("append: open {} for resume: {e}", path.display()))?,
);
let (entries, _manifest_offset) = read_znippy_full_manifest(path)?;
if entries.is_empty() {
return Err(anyhow!("append: archive {} has an empty manifest", path.display()));
}
let blob_end = entries
.iter()
.map(|e| e.index_offset)
.min()
.ok_or_else(|| anyhow!("append: no sections in manifest"))?;
let (paths, locs) = recover_rows(path, &entries)?;
let carried: Vec<(String, ChunkLoc)> = paths.into_iter().zip(locs).collect();
let meta = match read_reserved_section_bytes(path, META_MODULE)? {
None => None,
Some(bytes) => Some(decode_meta_section(&bytes)?.to_table()),
};
let mut carried_reserved = Vec::new();
for module in CARRIED_RESERVED_MODULES {
if let Some(bytes) = read_reserved_section_bytes(path, module)? {
carried_reserved.push(((*module).to_string(), bytes));
}
}
let delta_map = read_delta_map(path)?;
Ok(Self {
file,
cursor: blob_end,
entries: Vec::new(), lookup_paths: Vec::new(),
lookup_locs: Vec::new(),
carried,
meta,
carried_reserved,
delta_map,
})
}
pub fn blob_end(&self) -> u64 {
self.cursor
}
pub fn recovered_rows(&self) -> usize {
self.carried.len()
}
fn drop_carried_paths(&mut self, replacing: &std::collections::HashSet<&str>) -> usize {
if self.carried.is_empty() || replacing.is_empty() {
return 0;
}
let before = self.carried.len();
self.carried.retain(|(p, _)| !replacing.contains(p.as_str()));
before - self.carried.len()
}
fn emit_carried(&mut self) -> Result<()> {
if self.carried.is_empty() {
return Ok(());
}
let carried = std::mem::take(&mut self.carried);
let (paths, locs): (Vec<String>, Vec<ChunkLoc>) = carried.into_iter().unzip();
let batch = base_batch_from_rows(&paths, &locs)?;
self.push_subindex(data_subindex_schema().as_ref(), &[batch], GroupKey {
pkg_type: 0,
repo: String::new(),
module_name: String::new(),
})
}
fn accumulate_lookup(&mut self, batch: &RecordBatch) {
let cols = (|| {
Some((
batch.column_by_name("relative_path")?.as_any().downcast_ref::<StringArray>()?,
batch.column_by_name("chunk_seq")?.as_any().downcast_ref::<UInt32Array>()?,
batch.column_by_name("fdata_offset")?.as_any().downcast_ref::<UInt64Array>()?,
batch.column_by_name("compressed")?.as_any().downcast_ref::<BooleanArray>()?,
batch.column_by_name("uncompressed_size")?.as_any().downcast_ref::<UInt64Array>()?,
batch.column_by_name("blob_offset")?.as_any().downcast_ref::<UInt64Array>()?,
batch.column_by_name("blob_size")?.as_any().downcast_ref::<UInt64Array>()?,
batch.column_by_name("checksum")?.as_any().downcast_ref::<FixedSizeBinaryArray>()?,
))
})();
let Some((paths, chunk_seq, fdata, compressed, usz, blob_off, blob_sz, checksum)) = cols
else { return; };
for i in 0..batch.num_rows() {
let mut ck = [0u8; 32];
ck.copy_from_slice(checksum.value(i));
self.lookup_paths.push(paths.value(i).to_string());
self.lookup_locs.push(ChunkLoc {
chunk_seq: chunk_seq.value(i),
fdata_offset: fdata.value(i),
blob_offset: blob_off.value(i),
blob_size: blob_sz.value(i),
uncompressed_size: usz.value(i),
compressed: compressed.value(i),
checksum: ck,
});
}
}
fn write_lookup_and_trie(&mut self) -> Result<()> {
let n = self.lookup_paths.len();
let mut order: Vec<usize> = (0..n).collect();
order.sort_by(|&a, &b| {
self.lookup_paths[a].cmp(&self.lookup_paths[b])
.then(self.lookup_locs[a].chunk_seq.cmp(&self.lookup_locs[b].chunk_seq))
});
let schema = lookup_schema();
let batch = base_batch_permuted(
schema.clone(),
&self.lookup_paths,
&self.lookup_locs,
&order,
)?;
self.push_subindex(&schema, &[batch], GroupKey {
pkg_type: RESERVED_PKG_TYPE,
repo: String::new(),
module_name: LOOKUP_MODULE.to_string(),
})?;
let mut builder = fst::MapBuilder::memory();
let mut prev: Option<&str> = None;
for (sorted_idx, &orig) in order.iter().enumerate() {
let p = self.lookup_paths[orig].as_str();
if prev != Some(p) {
builder.insert(p.as_bytes(), sorted_idx as u64)
.map_err(|e| anyhow!("trie insert: {e}"))?;
prev = Some(p);
}
}
let trie_bytes = builder.into_inner().map_err(|e| anyhow!("trie finish: {e}"))?;
self.write_raw_section(&trie_bytes, GroupKey {
pkg_type: RESERVED_PKG_TYPE,
repo: String::new(),
module_name: TRIE_MODULE.to_string(),
})
}
fn write_meta_subindex(&mut self) -> Result<()> {
let Some(table) = self.meta.take() else {
return Ok(());
};
let batch = build_meta_batch(&table)?;
let schema = meta_schema();
self.push_subindex(schema.as_ref(), &[batch], GroupKey {
pkg_type: RESERVED_PKG_TYPE,
repo: String::new(),
module_name: META_MODULE.to_string(),
})
}
fn write_carried_reserved(&mut self) -> Result<()> {
for (module, bytes) in std::mem::take(&mut self.carried_reserved) {
self.write_raw_section(&bytes, GroupKey {
pkg_type: RESERVED_PKG_TYPE,
repo: String::new(),
module_name: module,
})?;
}
Ok(())
}
fn write_delta_map(&mut self) -> Result<()> {
if self.delta_map.is_empty() {
return Ok(());
}
let rows = std::mem::take(&mut self.delta_map);
let paths = StringArray::from(rows.iter().map(|r| r.0.as_str()).collect::<Vec<_>>());
let seqs = UInt32Array::from(rows.iter().map(|r| r.1).collect::<Vec<_>>());
let bases = StringArray::from(rows.iter().map(|r| r.2.as_str()).collect::<Vec<_>>());
let schema = crate::index::delta_map_schema();
let batch = RecordBatch::try_new(
Arc::clone(&schema),
vec![Arc::new(paths), Arc::new(seqs), Arc::new(bases)],
)
.map_err(|e| anyhow!("delta map batch: {e}"))?;
self.push_subindex(schema.as_ref(), &[batch], GroupKey {
pkg_type: RESERVED_PKG_TYPE,
repo: String::new(),
module_name: ZNIPPY_DELTA_MODULE.to_string(),
})
}
pub fn file(&self) -> &Arc<File> {
&self.file
}
pub fn advance_blob_end(&mut self, n: u64) {
self.cursor += n;
}
pub fn replace_carried(&mut self, path: &str, locs: Vec<ChunkLoc>) {
self.carried.retain(|(p, _)| p != path);
for loc in locs {
self.carried.push((path.to_string(), loc));
}
}
pub fn push_delta_map_row(&mut self, path: String, chunk_seq: u32, base: String) {
self.delta_map.retain(|(p, s, _)| !(p == &path && *s == chunk_seq));
self.delta_map.push((path, chunk_seq, base));
}
fn write_raw_section(&mut self, bytes: &[u8], key: GroupKey) -> Result<()> {
let start = self.cursor;
self.file.write_all_at(bytes, start)?;
self.cursor += bytes.len() as u64;
self.entries.push(ManifestEntry {
pkg_type: key.pkg_type,
repo: key.repo,
module_name: key.module_name,
index_offset: start,
index_len: bytes.len() as u64,
row_count: 0,
});
Ok(())
}
}
impl ArchiveMetaSink for ArrowIpcSinkAppend {
fn push_subindex(
&mut self,
schema: &Schema,
batches: &[RecordBatch],
key: GroupKey,
) -> Result<()> {
let sub_start = self.cursor;
let mut sub_bytes: Vec<u8> = Vec::new();
let mut sw = StreamWriter::try_new(&mut sub_bytes, schema)
.map_err(|e| anyhow!("sub-index writer: {e}"))?;
let mut row_count = 0u64;
for batch in batches {
row_count += batch.num_rows() as u64;
sw.write(batch).map_err(|e| anyhow!("sub-index write: {e}"))?;
}
sw.finish().map_err(|e| anyhow!("sub-index finish: {e}"))?;
if !is_reserved_module(&key.module_name) {
for batch in batches {
self.accumulate_lookup(batch);
}
}
let sub_len = sub_bytes.len() as u64;
self.file.write_all_at(&sub_bytes, sub_start)?;
self.cursor += sub_len;
self.entries.push(ManifestEntry {
pkg_type: key.pkg_type,
repo: key.repo,
module_name: key.module_name,
index_offset: sub_start,
index_len: sub_len,
row_count,
});
Ok(())
}
fn finish(mut self: Box<Self>) -> Result<u64> {
self.emit_carried()?;
self.write_lookup_and_trie()?;
self.write_meta_subindex()?;
self.write_carried_reserved()?;
self.write_delta_map()?;
let manifest_offset = self.cursor;
let manifest_bytes =
write_manifest_bytes(&self.entries).map_err(|e| anyhow!("manifest: {e}"))?;
self.file.write_all_at(&manifest_bytes, manifest_offset)?;
let after = manifest_offset + manifest_bytes.len() as u64;
self.file.write_all_at(&MULTI_INDEX_MAGIC, after)?;
self.file.write_all_at(
&manifest_offset.to_le_bytes(),
after + MULTI_INDEX_MAGIC.len() as u64,
)?;
let final_len = after + MULTI_INDEX_MAGIC.len() as u64 + 8;
self.file.set_len(final_len)?;
self.file.sync_all()?;
Ok(final_len)
}
}
pub fn read_delta_map(path: &Path) -> Result<Vec<(String, u32, String)>> {
use arrow::ipc::reader::StreamReader;
let Some(bytes) = read_reserved_section_bytes(path, ZNIPPY_DELTA_MODULE)? else {
return Ok(Vec::new());
};
let reader = StreamReader::try_new(std::io::Cursor::new(bytes), None)
.map_err(|e| anyhow!("delta map: {e}"))?;
let mut out = Vec::new();
for batch in reader {
let batch = batch.map_err(|e| anyhow!("delta map batch: {e}"))?;
let paths = batch
.column_by_name("relative_path")
.and_then(|c| c.as_any().downcast_ref::<StringArray>())
.ok_or_else(|| anyhow!("delta map: missing relative_path"))?;
let seqs = batch
.column_by_name("chunk_seq")
.and_then(|c| c.as_any().downcast_ref::<UInt32Array>())
.ok_or_else(|| anyhow!("delta map: missing chunk_seq"))?;
let bases = batch
.column_by_name("base_path")
.and_then(|c| c.as_any().downcast_ref::<StringArray>())
.ok_or_else(|| anyhow!("delta map: missing base_path"))?;
for r in 0..batch.num_rows() {
out.push((paths.value(r).to_string(), seqs.value(r), bases.value(r).to_string()));
}
}
Ok(out)
}
pub fn supersede_as_delta(
archive: &Path,
superseded: &str,
base: &str,
chain_depth: usize,
compression_level: i32,
) -> Result<SupersedeOutcome> {
if superseded == base {
return Err(anyhow!("an entry cannot be a delta against itself: {superseded}"));
}
if chain_depth + 1 > MAX_GENERATION_CHAIN {
return Ok(SupersedeOutcome::ChainTooLong);
}
let (old_bytes, base_bytes) = {
let ar = crate::ZnippyArchive::open(archive)?;
(
ar.extract_file_verified(superseded)?,
ar.extract_file_verified(base)?,
)
};
let delta = crate::archive::encode_delta_against(&base_bytes, &old_bytes);
if (delta.len() as f64) >= crate::archive::DELTA_SIZE_ALPHA * (old_bytes.len() as f64) {
return Ok(SupersedeOutcome::NotSmaller {
delta_bytes: delta.len() as u64,
stored_bytes: old_bytes.len() as u64,
});
}
let mut sink = ArrowIpcSinkAppend::open_existing(archive)?;
let at = sink.blob_end();
let mut ctx = crate::codec::CompressCtx::new(compression_level)?;
let frame = ctx.compress(&delta).ok();
let (on_disk, compressed): (&[u8], bool) = match frame.as_deref() {
Some(f) if f.len() < delta.len() => (f, true),
_ => (&delta, false),
};
sink.file().write_all_at(on_disk, at)?;
sink.advance_blob_end(on_disk.len() as u64);
sink.replace_carried(superseded, vec![ChunkLoc {
chunk_seq: 0,
fdata_offset: 0,
blob_offset: at,
blob_size: on_disk.len() as u64,
uncompressed_size: old_bytes.len() as u64,
compressed,
checksum: *blake3::hash(&old_bytes).as_bytes(),
}]);
sink.push_delta_map_row(superseded.to_string(), 0, base.to_string());
Box::new(sink).finish()?;
Ok(SupersedeOutcome::Delta {
stored_bytes: old_bytes.len() as u64,
delta_bytes: on_disk.len() as u64,
chain_depth: chain_depth + 1,
})
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct CompactReport {
pub bytes_before: u64,
pub bytes_after: u64,
pub rows: u64,
pub delta_rows: u64,
}
pub fn compact_archive(archive: &Path) -> Result<CompactReport> {
let bytes_before = std::fs::metadata(archive)?.len();
let src = ArrowIpcSinkAppend::open_existing(archive)?;
let staged = {
let unique = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos())
.unwrap_or(0);
let mut p = archive.as_os_str().to_owned();
p.push(format!(".compact-{}-{unique}", std::process::id()));
std::path::PathBuf::from(p)
};
let out = Arc::new(
std::fs::OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(true)
.open(&staged)
.map_err(|e| anyhow!("compact: staging {}: {e}", staged.display()))?,
);
let mut sink = ArrowIpcSinkAppend::new(Arc::clone(&out), 0);
sink.meta = src.meta.clone();
sink.carried_reserved = src.carried_reserved.clone();
sink.delta_map = src.delta_map.clone();
let delta_rows = sink.delta_map.len() as u64;
let mut cursor = 0u64;
let mut buf: Vec<u8> = Vec::new();
for (path, loc) in &src.carried {
let n = loc.blob_size as usize;
buf.clear();
buf.resize(n, 0);
src.file
.read_exact_at(&mut buf, loc.blob_offset)
.map_err(|e| anyhow!("compact: reading {path} at {}: {e}", loc.blob_offset))?;
out.write_all_at(&buf, cursor)?;
let mut moved = loc.clone();
moved.blob_offset = cursor;
cursor += loc.blob_size;
sink.carried.push((path.clone(), moved));
}
let rows = sink.carried.len() as u64;
sink.cursor = cursor;
Box::new(sink).finish()?;
out.sync_all()?;
drop(out);
drop(src);
std::fs::rename(&staged, archive)?;
if let Some(parent) = archive.parent() {
if let Ok(f) = std::fs::File::open(parent) {
let _ = f.sync_all();
}
}
Ok(CompactReport {
bytes_before,
bytes_after: std::fs::metadata(archive)?.len(),
rows,
delta_rows,
})
}
#[derive(Debug, PartialEq, Eq)]
pub enum SupersedeOutcome {
Delta { stored_bytes: u64, delta_bytes: u64, chain_depth: usize },
NotSmaller { delta_bytes: u64, stored_bytes: u64 },
ChainTooLong,
}
pub const MAX_GENERATION_CHAIN: usize = 32;
pub fn base_batch_from_rows(paths: &[String], locs: &[ChunkLoc]) -> Result<RecordBatch> {
let order: Vec<usize> = (0..paths.len()).collect();
base_batch_permuted(data_subindex_schema(), paths, locs, &order)
}
pub(crate) fn base_batch_permuted(
schema: Arc<Schema>,
paths: &[String],
locs: &[ChunkLoc],
order: &[usize],
) -> Result<RecordBatch> {
let n = order.len();
let mut path_b = StringBuilder::with_capacity(n, n * 16);
let mut seq_b = UInt32Builder::with_capacity(n);
let mut fdata_b = UInt64Builder::with_capacity(n);
let mut comp_b = BooleanBuilder::with_capacity(n);
let mut usz_b = UInt64Builder::with_capacity(n);
let mut boff_b = UInt64Builder::with_capacity(n);
let mut bsz_b = UInt64Builder::with_capacity(n);
let mut ck_b = FixedSizeBinaryBuilder::with_capacity(n, 32);
for &i in order {
let loc = &locs[i];
path_b.append_value(&paths[i]);
seq_b.append_value(loc.chunk_seq);
fdata_b.append_value(loc.fdata_offset);
comp_b.append_value(loc.compressed);
usz_b.append_value(loc.uncompressed_size);
boff_b.append_value(loc.blob_offset);
bsz_b.append_value(loc.blob_size);
ck_b.append_value(loc.checksum).expect("checksum is 32 bytes");
}
Ok(RecordBatch::try_new(
schema,
vec![
Arc::new(path_b.finish()),
Arc::new(seq_b.finish()),
Arc::new(fdata_b.finish()),
Arc::new(comp_b.finish()),
Arc::new(usz_b.finish()),
Arc::new(boff_b.finish()),
Arc::new(bsz_b.finish()),
Arc::new(ck_b.finish()),
],
)?)
}
pub(crate) fn write_blobs(
file: &File,
cursor: u64,
files: &[(String, Vec<u8>)],
ctx: &mut crate::codec::CompressCtx,
policy: crate::SkipPolicy,
) -> Result<(Vec<String>, Vec<ChunkLoc>, u64)> {
let mut paths = Vec::with_capacity(files.len());
let mut locs = Vec::with_capacity(files.len());
let mut cursor = cursor;
for (rel, bytes) in files {
let checksum = *blake3::hash(bytes).as_bytes();
let skip = policy.skip_by_path(std::path::Path::new(rel.as_str()));
let frame = if skip { Vec::new() } else { ctx.compress(bytes)? };
let (on_disk, compressed): (&[u8], bool) = if !skip && frame.len() < bytes.len() {
(&frame, true)
} else {
(bytes, false)
};
let blob_offset = cursor;
file.write_all_at(on_disk, blob_offset)?;
cursor += on_disk.len() as u64;
paths.push(rel.clone());
locs.push(ChunkLoc {
chunk_seq: 0,
fdata_offset: 0,
blob_offset,
blob_size: on_disk.len() as u64,
uncompressed_size: bytes.len() as u64,
compressed,
checksum,
});
}
Ok((paths, locs, cursor))
}
#[derive(Debug, Clone)]
pub struct AppendReport {
pub rows_before: u64,
pub rows_replaced: u64,
pub rows_added: u64,
pub blob_append_offset: u64,
pub blob_bytes_added: u64,
pub sealed_total_bytes: u64,
}
pub fn append_files(
archive: &Path,
new_files: &[(String, Vec<u8>)],
compression_level: i32,
) -> Result<AppendReport> {
append_files_with_meta(archive, new_files, compression_level, None)
}
pub fn append_files_with_policy(
archive: &Path,
new_files: &[(String, Vec<u8>)],
compression_level: i32,
policy: crate::SkipPolicy,
) -> Result<AppendReport> {
let sink = ArrowIpcSinkAppend::open_existing(archive)?;
let rows_before = sink.recovered_rows() as u64;
let blob_append_offset = sink.blob_end();
write_files_into_sink(
sink,
new_files,
compression_level,
rows_before,
blob_append_offset,
policy,
)
}
pub fn append_files_with_meta(
archive: &Path,
new_files: &[(String, Vec<u8>)],
compression_level: i32,
meta: Option<MetaTable>,
) -> Result<AppendReport> {
let mut sink = ArrowIpcSinkAppend::open_existing(archive)?;
if let Some(table) = meta {
sink.merge_meta(table.rows().to_vec());
}
let rows_before = sink.recovered_rows() as u64;
let blob_append_offset = sink.blob_end();
write_files_into_sink(
sink,
new_files,
compression_level,
rows_before,
blob_append_offset,
crate::SkipPolicy::resolve(),
)
}
pub fn create_archive(
archive: &Path,
files: &[(String, Vec<u8>)],
compression_level: i32,
) -> Result<AppendReport> {
create_archive_with_meta(archive, files, compression_level, None)
}
pub fn create_archive_with_meta(
archive: &Path,
files: &[(String, Vec<u8>)],
compression_level: i32,
meta: Option<MetaTable>,
) -> Result<AppendReport> {
let blob_file = Arc::new(
File::create(archive)
.map_err(|e| anyhow!("create archive {}: {e}", archive.display()))?,
);
let mut sink = ArrowIpcSinkAppend::new(blob_file, 0);
sink.meta = meta;
write_files_into_sink(sink, files, compression_level, 0, 0, crate::SkipPolicy::resolve())
}
pub fn create_archive_to_vec(
files: &[(String, Vec<u8>)],
compression_level: i32,
) -> Result<(Vec<u8>, AppendReport)> {
let anon = Arc::new(
tempfile::tempfile().map_err(|e| anyhow!("anonymous archive fd: {e}"))?,
);
let sink = ArrowIpcSinkAppend::new(anon.clone(), 0);
let report =
write_files_into_sink(sink, files, compression_level, 0, 0, crate::SkipPolicy::resolve())?;
let mut bytes = vec![0u8; report.sealed_total_bytes as usize];
anon.read_exact_at(&mut bytes, 0)
.map_err(|e| anyhow!("read back anonymous archive: {e}"))?;
Ok((bytes, report))
}
fn write_files_into_sink(
mut sink: ArrowIpcSinkAppend,
new_files: &[(String, Vec<u8>)],
compression_level: i32,
rows_before: u64,
blob_append_offset: u64,
policy: crate::SkipPolicy,
) -> Result<AppendReport> {
use crate::codec::CompressCtx;
let incoming: std::collections::HashSet<&str> =
new_files.iter().map(|(rel, _)| rel.as_str()).collect();
let rows_replaced = sink.drop_carried_paths(&incoming) as u64;
let blob_file = sink.file.clone();
let mut ctx = CompressCtx::new(compression_level)?;
let (paths, locs, cursor) =
write_blobs(&blob_file, blob_append_offset, new_files, &mut ctx, policy)?;
let blob_bytes_added = cursor - blob_append_offset;
blob_file.sync_all()?;
sink.cursor = cursor;
let batch = base_batch_from_rows(&paths, &locs)?;
let schema = data_subindex_schema();
let rows_added = batch.num_rows() as u64;
sink.push_subindex(
schema.as_ref(),
&[batch],
GroupKey { pkg_type: 0, repo: String::new(), module_name: String::new() },
)?;
let sealed_total_bytes = Box::new(sink).finish()?;
Ok(AppendReport {
rows_before,
rows_replaced,
rows_added,
blob_append_offset,
blob_bytes_added,
sealed_total_bytes,
})
}
pub(crate) fn recover_rows(
path: &Path,
entries: &[ManifestEntry],
) -> Result<(Vec<String>, Vec<ChunkLoc>)> {
use std::io::{Read, Seek, SeekFrom};
let mut file = File::open(path)?;
if let Some(lk) = entries.iter().find(|e| e.module_name == LOOKUP_MODULE) {
file.seek(SeekFrom::Start(lk.index_offset))?;
let mut bytes = vec![0u8; lk.index_len as usize];
file.read_exact(&mut bytes)?;
return decode_base_rows(&bytes);
}
let mut paths = Vec::new();
let mut locs = Vec::new();
for e in entries {
if is_reserved_module(&e.module_name) {
continue;
}
file.seek(SeekFrom::Start(e.index_offset))?;
let mut bytes = vec![0u8; e.index_len as usize];
file.read_exact(&mut bytes)?;
let (mut p, mut l) = decode_base_rows(&bytes)?;
paths.append(&mut p);
locs.append(&mut l);
}
Ok((paths, locs))
}
#[cfg(all(test, feature = "openzl"))]
mod tests {
use super::*;
use crate::codec::CompressCtx;
use crate::meta::{BlobMeta, ChunkMeta};
use crate::{ArrowIpcSink, ZnippyArchive, ZnippyReader};
use std::time::{SystemTime, UNIX_EPOCH};
fn unique_dir(tag: &str) -> std::path::PathBuf {
let ns = SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_nanos();
let d = std::env::temp_dir().join(format!("znippy_append_{tag}_{ns}_{:?}", std::thread::current().id()));
std::fs::create_dir_all(&d).unwrap();
d
}
#[test]
fn an_append_honours_the_skip_policy_instead_of_compressing_everything() {
let dir = unique_dir("skip_policy");
let archive = dir.join("a.znippy");
create_archive(&archive, &[("seed.txt".into(), b"seed".to_vec())], 3).unwrap();
let squishy = vec![b'A'; 256 * 1024];
let name = format!("pack-{}.pack", "0f".repeat(20));
let report =
append_files(&archive, &[(name.clone(), squishy.clone())], 3).unwrap();
assert_eq!(
report.blob_bytes_added,
squishy.len() as u64,
"a `.pack` entry must be stored RAW. {} bytes were written for a {}-byte input, so \
the codec ran over a file the extension table already said was compressed",
report.blob_bytes_added,
squishy.len()
);
assert_eq!(crate::get_file(&archive, &name).unwrap(), squishy);
let report2 =
append_files(&archive, &[("plain.txt".into(), squishy.clone())], 3).unwrap();
assert!(
report2.blob_bytes_added < squishy.len() as u64 / 10,
"an ordinary name must still be compressed; {} bytes for {}",
report2.blob_bytes_added,
squishy.len()
);
assert_eq!(crate::get_file(&archive, "plain.txt").unwrap(), squishy);
let report3 = append_files_with_policy(
&archive,
&[("also-plain.txt".into(), squishy.clone())],
3,
crate::SkipPolicy::already_compressed(),
)
.unwrap();
assert_eq!(
report3.blob_bytes_added,
squishy.len() as u64,
"`already_compressed()` must store raw whatever the name says"
);
std::fs::remove_dir_all(&dir).ok();
}
fn synth(n: usize, salt: u64) -> Vec<(String, Vec<u8>)> {
(0..n)
.map(|i| {
let g = (i.wrapping_mul(2_654_435_761) ^ salt as usize) % 1000;
let p = format!("repo/grp{g:03}/file{:08}_{salt}.bin", i);
let body = format!("payload {i} salt {salt} {}\n", "z".repeat(8 + (i % 40)));
(p, body.into_bytes())
})
.collect()
}
fn write_fresh<S: ArchiveMetaSink + 'static>(
path: &Path,
files: &[(String, Vec<u8>)],
make_sink: impl FnOnce(Arc<File>, u64) -> S,
) -> u64 {
let file = Arc::new(File::create(path).unwrap());
let mut ctx = CompressCtx::new(3).unwrap();
let mut blobs = Vec::new();
let mut paths = Vec::new();
let mut cursor = 0u64;
for (fi, (rel, bytes)) in files.iter().enumerate() {
let checksum = *blake3::hash(bytes).as_bytes();
let frame = ctx.compress(bytes).unwrap();
let (on_disk, compressed): (&[u8], bool) =
if frame.len() < bytes.len() { (&frame, true) } else { (bytes, false) };
file.write_all_at(on_disk, cursor).unwrap();
let blob_offset = cursor;
cursor += on_disk.len() as u64;
paths.push(rel.clone());
blobs.push(BlobMeta {
blob_offset,
blob_size: on_disk.len() as u64,
chunk_meta: ChunkMeta {
fdata_offset: 0,
file_index: fi as u64,
chunk_seq: 0,
checksum,
compressed,
uncompressed_size: bytes.len() as u64,
compressed_size: on_disk.len() as u64,
},
});
}
let resolver = { let p = paths.clone(); move |fi: u64| p[fi as usize].clone() };
let batch = crate::build_metadata_batch(&blobs, resolver, &[], &[]).unwrap();
let schema = data_subindex_schema();
let mut sink = make_sink(file.clone(), cursor);
sink.push_subindex(
schema.as_ref(),
&[batch],
GroupKey { pkg_type: 0, repo: String::new(), module_name: String::new() },
)
.unwrap();
Box::new(sink).finish().unwrap()
}
#[test]
fn clone_fresh_path_is_byte_identical_to_original() {
let dir = unique_dir("parity");
let files = synth(2_000, 1);
let a = dir.join("a.znippy");
let b = dir.join("b.znippy");
let len_a = write_fresh(&a, &files, ArrowIpcSink::new);
let len_b = write_fresh(&b, &files, ArrowIpcSinkAppend::new);
assert_eq!(len_a, len_b, "clone seal produced a different total length");
let bytes_a = std::fs::read(&a).unwrap();
let bytes_b = std::fs::read(&b).unwrap();
assert_eq!(
bytes_a, bytes_b,
"clone's fresh write path is NOT byte-identical to ArrowIpcSink — parity broken"
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn native_append_roundtrips_old_and_new_files() {
let dir = unique_dir("resume");
let archive = dir.join("store.znippy");
let orig = synth(1_500, 7);
write_fresh(&archive, &orig, ArrowIpcSink::new);
let added = synth(300, 99);
let report = append_files(&archive, &added, 3).unwrap();
assert_eq!(report.rows_before, orig.len() as u64, "must recover all original rows");
assert_eq!(report.rows_added, added.len() as u64);
assert!(report.blob_bytes_added > 0, "append must write new blob bytes");
assert!(
report.sealed_total_bytes > report.blob_append_offset,
"re-sealed file must be larger than the old blob region"
);
let ar = ZnippyArchive::open(&archive).unwrap();
let mut listed = ar.list_files().unwrap();
listed.sort();
let mut expected: Vec<String> =
orig.iter().chain(added.iter()).map(|(p, _)| p.clone()).collect();
expected.sort();
assert_eq!(listed, expected, "index must list exactly old+new files after append");
for (p, bytes) in orig.iter().chain(added.iter()) {
let got = ar.extract_file(p).unwrap();
assert_eq!(&got, bytes, "byte mismatch after append for {p}");
}
let probe = &added[123].0;
let chunks = crate::locate_file(&archive, probe).unwrap();
assert!(!chunks.is_empty(), "appended file must be locatable via the re-sealed lookup");
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn create_archive_seeds_then_grows() {
let dir = unique_dir("create");
let seed = synth(40, 5);
let made = dir.join("made.znippy");
let report = create_archive(&made, &seed, 3).unwrap();
assert_eq!(report.rows_before, 0, "fresh archive has no prior rows");
assert_eq!(report.rows_added, seed.len() as u64);
let ref_path = dir.join("ref.znippy");
write_fresh(&ref_path, &seed, ArrowIpcSinkAppend::new);
assert_eq!(
std::fs::read(&made).unwrap(),
std::fs::read(&ref_path).unwrap(),
"create_archive must be byte-identical to the proven fresh write path"
);
let ar = ZnippyArchive::open(&made).unwrap();
for (p, bytes) in &seed {
assert_eq!(&ar.extract_file(p).unwrap(), bytes, "seed byte mismatch for {p}");
}
let added = synth(15, 88);
let rep2 = append_files(&made, &added, 3).unwrap();
assert_eq!(rep2.rows_before, seed.len() as u64, "append must recover seeded rows");
assert_eq!(rep2.rows_added, added.len() as u64);
let ar2 = ZnippyArchive::open(&made).unwrap();
for (p, bytes) in seed.iter().chain(added.iter()) {
assert_eq!(&ar2.extract_file(p).unwrap(), bytes, "byte mismatch after grow for {p}");
}
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn create_archive_to_vec_is_filesystem_free_and_round_trips() {
let files = synth(24, 7);
let (bytes, report) = create_archive_to_vec(&files, 3).unwrap();
assert_eq!(report.rows_before, 0, "fresh in-memory archive has no prior rows");
assert_eq!(report.rows_added, files.len() as u64);
assert_eq!(bytes.len() as u64, report.sealed_total_bytes, "vec len == sealed size");
let dir = unique_dir("tovec");
let ref_path = dir.join("ref.znippy");
create_archive(&ref_path, &files, 3).unwrap();
assert_eq!(
bytes,
std::fs::read(&ref_path).unwrap(),
"in-memory archive must be byte-identical to create_archive's file output"
);
let p = dir.join("from_mem.znippy");
std::fs::write(&p, &bytes).unwrap();
let ar = ZnippyArchive::open(&p).unwrap();
for (name, content) in &files {
assert_eq!(&ar.extract_file(name).unwrap(), content, "byte mismatch for {name}");
}
let _ = std::fs::remove_dir_all(&dir);
}
fn recorded_format_version(path: &Path) -> Option<String> {
use std::io::{Read, Seek, SeekFrom};
use arrow::ipc::reader::StreamReader;
let entries = crate::index::read_znippy_manifest(path).ok()?;
let mut file = File::open(path).ok()?;
for e in &entries {
if is_reserved_module(&e.module_name) {
continue;
}
file.seek(SeekFrom::Start(e.index_offset)).ok()?;
let mut bytes = vec![0u8; e.index_len as usize];
file.read_exact(&mut bytes).ok()?;
let reader = StreamReader::try_new(std::io::Cursor::new(bytes), None).ok()?;
return reader
.schema()
.metadata()
.get(crate::index::FORMAT_VERSION_KEY)
.cloned();
}
None
}
#[test]
fn every_appended_archive_records_the_format_version() {
let dir = unique_dir("fmtver");
let want = crate::index::ZNIPPY_FORMAT_VERSION.to_string();
let made = dir.join("made.znippy");
create_archive(&made, &synth(12, 3), 3).unwrap();
assert_eq!(
recorded_format_version(&made).as_deref(),
Some(want.as_str()),
"create_archive must stamp the on-disk format version"
);
append_files(&made, &synth(7, 91), 3).unwrap();
assert_eq!(
recorded_format_version(&made).as_deref(),
Some(want.as_str()),
"append_files must re-stamp the format version on the re-sealed archive"
);
let (bytes, _) = create_archive_to_vec(&synth(9, 4), 3).unwrap();
let mem = dir.join("mem.znippy");
std::fs::write(&mem, &bytes).unwrap();
assert_eq!(
recorded_format_version(&mem).as_deref(),
Some(want.as_str()),
"create_archive_to_vec must stamp the on-disk format version"
);
let ar = ZnippyArchive::open(&mem).unwrap();
for (p, body) in &synth(9, 4) {
assert_eq!(&ar.extract_file(p).unwrap(), body, "byte mismatch for {p}");
}
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn metadata_is_searchable_survives_a_reseal_and_absence_is_reported_as_absence() {
use crate::meta_index::{ArchiveMeta, MetaSearch, MetaTable, MetaValue, read_archive_meta};
let dir = unique_dir("meta");
let files = synth(30, 2);
let plain = dir.join("plain.znippy");
create_archive(&plain, &files, 3).unwrap();
let m = read_archive_meta(&plain).unwrap();
assert_eq!(m, ArchiveMeta::NoMetadata, "an archive with no meta section must say so");
assert!(!m.is_searchable());
assert_eq!(m.find_by_key("build-thing"), MetaSearch::NoMetadata);
assert!(
m.find_by_key("build-thing").hits().is_none(),
"absence must NOT present itself as an empty result set"
);
let empty = dir.join("empty.znippy");
create_archive_with_meta(&empty, &files, 3, Some(MetaTable::new())).unwrap();
let me = read_archive_meta(&empty).unwrap();
assert!(me.is_searchable(), "a sealed empty index WAS searched");
assert!(me.index().is_some_and(|i| i.is_empty()));
assert_eq!(me.find_by_key("build-thing"), MetaSearch::Hits(&[]));
let wasm = b"\0asm\x01\0\0\0".to_vec();
let (p0, p1, p2) = (files[0].0.clone(), files[1].0.clone(), files[2].0.clone());
let mut t = MetaTable::new();
t.insert(p0.clone(), "build-thing", MetaValue::Bytes(wasm.clone()))
.insert(p0.clone(), "build-thing.abi", "wasi-p2")
.insert(p1.clone(), "build-thing", MetaValue::Bytes(wasm.clone()))
.insert(p2.clone(), "coverage", 0.5f64)
.insert_archive("producer", "znippy");
let ar = dir.join("meta.znippy");
create_archive_with_meta(&ar, &files, 3, Some(t)).unwrap();
let m = read_archive_meta(&ar).unwrap();
let idx = m.index().expect("sealed index is present");
assert_eq!(idx.len(), 5);
let hits = m.find_by_key("build-thing").hits().unwrap();
assert_eq!(hits.len(), 2, "exactly the two entries that carry one");
let mut got: Vec<&str> = hits.iter().filter_map(|h| h.path()).collect();
got.sort();
let mut want = vec![p0.as_str(), p1.as_str()];
want.sort();
assert_eq!(got, want, "the search names the right ENTRIES");
assert_eq!(hits[0].value.as_bytes(), Some(&wasm[..]), "and the right VALUE");
assert_eq!(idx.archive_value("producer").and_then(MetaValue::as_str), Some("znippy"));
assert_eq!(
idx.find_by_prefix("build-thing").len(),
3,
"prefix sweeps build-thing + build-thing.abi"
);
let reader = ZnippyArchive::open(&ar).unwrap();
let want_bytes = &files.iter().find(|(p, _)| *p == p0).unwrap().1;
assert_eq!(&reader.extract_file(hits[0].path().unwrap()).unwrap(), want_bytes);
let added = synth(6, 77);
let mut more = MetaTable::new();
more.insert(added[0].0.clone(), "build-thing", MetaValue::Bytes(wasm.clone()));
append_files_with_meta(&ar, &added, 3, Some(more)).unwrap();
let m2 = read_archive_meta(&ar).unwrap();
let hits2 = m2.find_by_key("build-thing").hits().unwrap();
assert_eq!(hits2.len(), 3, "the append merged, it did not replace");
assert_eq!(
m2.index().unwrap().archive_value("producer").and_then(MetaValue::as_str),
Some("znippy"),
"the archive-level row survived the re-seal"
);
append_files(&ar, &synth(3, 91), 3).unwrap();
assert_eq!(
read_archive_meta(&ar).unwrap().find_by_key("build-thing").hits().unwrap().len(),
3,
"an append that says nothing about metadata must not erase it"
);
let ar2 = ZnippyArchive::open(&ar).unwrap();
let n_files = files.len() + added.len() + 3;
assert_eq!(
ar2.file_count(),
n_files,
"metadata rows leaked into the data index"
);
let lookup_bytes = read_reserved_section_bytes(&ar, LOOKUP_MODULE).unwrap().unwrap();
assert_eq!(
decode_base_rows(&lookup_bytes).unwrap().0.len(),
n_files,
"metadata rows leaked into the random-access lookup"
);
for (p, bytes) in files.iter().chain(added.iter()) {
assert_eq!(&ar2.extract_file(p).unwrap(), bytes, "byte mismatch for {p}");
}
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn an_append_carries_the_ref_log_forward_and_drops_the_derived_sections() {
use crate::index::{
GUNNAR_GRAPH_MODULE, GUNNAR_REFS_MODULE, GUNNAR_SECRETS_MODULE,
};
use crate::meta_sink::{ReservedSection, ReservedSectionBuilder};
let dir = unique_dir("carried_reserved");
let archive = dir.join("a.znippy");
let refs_bytes = b"refs/heads/main 0123456789abcdef -- push 1".to_vec();
let secrets_bytes = b"\x00ciphertext-only, never plaintext".to_vec();
let graph_bytes = b"a commit graph derived from the objects".to_vec();
{
let f = Arc::new(File::create(&archive).unwrap());
let mut cursor = 0u64;
let mut blobs = Vec::new();
for (i, (name, bytes)) in [("obj/a.bin", b"first".to_vec())].iter().enumerate() {
use std::os::unix::fs::FileExt;
f.write_all_at(bytes, cursor).unwrap();
blobs.push(BlobMeta {
blob_offset: cursor,
blob_size: bytes.len() as u64,
chunk_meta: ChunkMeta {
fdata_offset: 0,
file_index: i as u64,
chunk_seq: 0,
checksum: *blake3::hash(bytes).as_bytes(),
compressed: false,
uncompressed_size: bytes.len() as u64,
compressed_size: bytes.len() as u64,
},
});
cursor += bytes.len() as u64;
let _ = name;
}
let batch = crate::index::build_metadata_batch(
&blobs,
|_fi: u64| "obj/a.bin".to_string(),
&[],
&[],
)
.unwrap();
let mut sink = ArrowIpcSink::new(Arc::clone(&f), cursor);
let (r, s, g) = (refs_bytes.clone(), secrets_bytes.clone(), graph_bytes.clone());
let builder: ReservedSectionBuilder = Box::new(move |_lookup| {
Ok(vec![
ReservedSection::raw(GUNNAR_REFS_MODULE, r),
ReservedSection::raw(GUNNAR_SECRETS_MODULE, s),
ReservedSection::raw(GUNNAR_GRAPH_MODULE, g),
])
});
sink = sink.with_reserved_builder(builder);
sink.push_subindex(
crate::index::lookup_schema().as_ref(),
&[batch],
GroupKey { pkg_type: 0, repo: String::new(), module_name: String::new() },
)
.unwrap();
Box::new(sink).finish().unwrap();
}
for m in [GUNNAR_REFS_MODULE, GUNNAR_SECRETS_MODULE, GUNNAR_GRAPH_MODULE] {
assert!(
read_reserved_section_bytes(&archive, m).unwrap().is_some(),
"{m} must be present BEFORE the append"
);
}
append_files(&archive, &[("obj/b.bin".to_string(), b"second".to_vec())], 3).unwrap();
assert_eq!(
read_reserved_section_bytes(&archive, GUNNAR_REFS_MODULE).unwrap(),
Some(refs_bytes),
"an object-carrying push must not erase the ref log — this is the 2026-08-04 bug"
);
assert_eq!(
read_reserved_section_bytes(&archive, GUNNAR_SECRETS_MODULE).unwrap(),
Some(secrets_bytes),
);
assert_eq!(
read_reserved_section_bytes(&archive, GUNNAR_GRAPH_MODULE).unwrap(),
None,
"a commit graph derived from the old object set must NOT be carried forward"
);
let reader = ZnippyArchive::open(&archive).unwrap();
assert_eq!(reader.extract_file("obj/a.bin").unwrap(), b"first".to_vec());
assert_eq!(reader.extract_file("obj/b.bin").unwrap(), b"second".to_vec());
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn a_superseded_entry_becomes_a_delta_and_reads_back_exactly() {
let dir = unique_dir("supersede");
let archive = dir.join("a.znippy");
let mut st = 0x243f_6a88_85a3_08d3u64;
let gen_n: Vec<u8> = (0..400_000u32)
.map(|_| {
st = st.wrapping_add(0x9e37_79b9_7f4a_7c15);
let mut z = st;
z = (z ^ (z >> 30)).wrapping_mul(0xbf58_476d_1ce4_e5b9);
(z ^ (z >> 27)) as u8
})
.collect();
let mut gen_n1 = gen_n.clone();
gen_n1.extend_from_slice(&gen_n[..20_000]);
create_archive(
&archive,
&[
("pack-N.pack".to_string(), gen_n.clone()),
("pack-N1.pack".to_string(), gen_n1.clone()),
],
3,
)
.unwrap();
let before = std::fs::metadata(&archive).unwrap().len();
let out = supersede_as_delta(&archive, "pack-N.pack", "pack-N1.pack", 0, 3).unwrap();
let (stored, delta_bytes) = match out {
SupersedeOutcome::Delta { stored_bytes, delta_bytes, chain_depth } => {
assert_eq!(chain_depth, 1);
(stored_bytes, delta_bytes)
}
other => panic!("expected a delta, got {other:?}"),
};
assert_eq!(stored, gen_n.len() as u64);
assert!(
delta_bytes * 20 < stored,
"N is a near-prefix of N+1, so the delta must be a small fraction of it; \
got {delta_bytes} against {stored}"
);
let ar = ZnippyArchive::open(&archive).unwrap();
assert_eq!(ar.extract_file("pack-N.pack").unwrap(), gen_n, "the superseded generation");
assert_eq!(ar.extract_file("pack-N1.pack").unwrap(), gen_n1, "the live generation");
assert_eq!(ar.extract_file_verified("pack-N.pack").unwrap(), gen_n);
assert_eq!(ar.file_size("pack-N.pack"), Some(gen_n.len() as u64));
assert!(
read_delta_map(&archive)
.unwrap()
.iter()
.all(|(p, _, _)| p == "pack-N.pack"),
"only the superseded generation may be a delta"
);
let _ = before;
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn the_writer_refuses_a_pointless_delta_and_an_over_long_chain() {
let dir = unique_dir("supersede_refuse");
let archive = dir.join("a.znippy");
let mut st = 0x1234_5678_9abc_def0u64;
let noise = |n: usize, st: &mut u64| -> Vec<u8> {
(0..n)
.map(|_| {
*st = st.wrapping_add(0x9e37_79b9_7f4a_7c15);
let mut z = *st;
z = (z ^ (z >> 30)).wrapping_mul(0xbf58_476d_1ce4_e5b9);
(z ^ (z >> 27)) as u8
})
.collect()
};
let a = noise(60_000, &mut st);
let b = noise(60_000, &mut st); create_archive(
&archive,
&[("a.pack".to_string(), a.clone()), ("b.pack".to_string(), b.clone())],
3,
)
.unwrap();
match supersede_as_delta(&archive, "a.pack", "b.pack", 0, 3).unwrap() {
SupersedeOutcome::NotSmaller { .. } => {}
other => panic!("unrelated bytes must not produce a delta, got {other:?}"),
}
let ar = ZnippyArchive::open(&archive).unwrap();
assert_eq!(ar.extract_file("a.pack").unwrap(), a);
assert!(read_delta_map(&archive).unwrap().is_empty());
assert_eq!(
supersede_as_delta(&archive, "a.pack", "b.pack", MAX_GENERATION_CHAIN, 3).unwrap(),
SupersedeOutcome::ChainTooLong
);
assert!(supersede_as_delta(&archive, "a.pack", "a.pack", 0, 3).is_err());
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn a_generation_chain_reads_every_generation_back_exactly() {
let dir = unique_dir("genchain");
let archive = dir.join("g.znippy");
let mut st = 0x0f0f_0f0f_dead_beefu64;
let mut gens: Vec<Vec<u8>> = Vec::new();
let mut cur: Vec<u8> = (0..120_000u32)
.map(|_| {
st = st.wrapping_add(0x9e37_79b9_7f4a_7c15);
let mut z = st;
z = (z ^ (z >> 30)).wrapping_mul(0xbf58_476d_1ce4_e5b9);
(z ^ (z >> 27)) as u8
})
.collect();
gens.push(cur.clone());
for g in 1..4 {
cur.extend_from_slice(format!("generation {g} tail ").repeat(200).as_bytes());
gens.push(cur.clone());
}
let files: Vec<(String, Vec<u8>)> = gens
.iter()
.enumerate()
.map(|(i, b)| (format!("pack-{i}.pack"), b.clone()))
.collect();
create_archive(&archive, &files, 3).unwrap();
for i in (0..3).rev() {
let out = supersede_as_delta(
&archive,
&format!("pack-{i}.pack"),
&format!("pack-{}.pack", i + 1),
3 - 1 - i,
3,
)
.unwrap();
assert!(matches!(out, SupersedeOutcome::Delta { .. }), "gen {i}: {out:?}");
}
let ar = ZnippyArchive::open(&archive).unwrap();
for (i, want) in gens.iter().enumerate() {
assert_eq!(
&ar.extract_file(&format!("pack-{i}.pack")).unwrap(),
want,
"generation {i} at chain depth {}",
3 - i
);
}
assert!(
read_delta_map(&archive)
.unwrap()
.iter()
.all(|(p, _, _)| p != "pack-3.pack"),
"the live generation must stay whole"
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
#[ignore]
fn perf_real_generation_chain() {
let Ok(d) = std::env::var("ZNIPPY_CHAIN_DIR") else { return };
let dir = unique_dir("realchain");
let archive = dir.join("g.znippy");
let mut files: Vec<(String, Vec<u8>)> = Vec::new();
for i in 0.. {
let p = std::path::Path::new(&d).join(format!("pack-{i}.pack"));
if !p.is_file() {
break;
}
files.push((format!("pack-{i}.pack"), std::fs::read(&p).unwrap()));
}
let n = files.len();
let raw: u64 = files.iter().map(|(_, b)| b.len() as u64).sum();
let t0 = std::time::Instant::now();
create_archive(&archive, &files, 3).unwrap();
let seal_ms = t0.elapsed().as_secs_f64() * 1e3;
let sealed = std::fs::metadata(&archive).unwrap().len();
println!("gen,stored_b,delta_b,ratio_x,supersede_ms");
let mut encode_ms_total = 0.0;
let mut delta_total = 0u64;
for i in (0..n - 1).rev() {
let t = std::time::Instant::now();
let out = supersede_as_delta(
&archive,
&format!("pack-{i}.pack"),
&format!("pack-{}.pack", i + 1),
n - 2 - i,
3,
)
.unwrap();
let ms = t.elapsed().as_secs_f64() * 1e3;
encode_ms_total += ms;
match out {
SupersedeOutcome::Delta { stored_bytes, delta_bytes, .. } => {
delta_total += delta_bytes;
println!(
"{i},{stored_bytes},{delta_bytes},{:.2},{ms:.0}",
stored_bytes as f64 / delta_bytes as f64
);
}
other => println!("{i},-,-,-,{ms:.0} ({other:?})"),
}
}
let ar = ZnippyArchive::open(&archive).unwrap();
let mut live: Vec<(String, Vec<u8>)> = Vec::new();
let t = std::time::Instant::now();
for i in 0..n {
let name = format!("pack-{i}.pack");
let got = ar.extract_file(&name).unwrap();
assert_eq!(got, files[i].1, "generation {i} did not read back exactly");
live.push((name, got));
}
let read_all_ms = t.elapsed().as_secs_f64() * 1e3;
println!("gen,depth,read_ms");
for i in 0..n {
let name = format!("pack-{i}.pack");
let t = std::time::Instant::now();
for _ in 0..5 {
let _ = ar.extract_file(&name).unwrap();
}
println!("{i},{},{:.2}", n - 1 - i, t.elapsed().as_secs_f64() * 1e3 / 5.0);
}
let compacted = dir.join("c.znippy");
let mut packed: Vec<(String, Vec<u8>)> = Vec::new();
for i in 0..n {
packed.push(live[i].clone());
}
let overhead = sealed - raw;
let steady = files[n - 1].1.len() as u64 + delta_total + overhead;
println!(
"SUMMARY generations={n} raw={raw} sealed={sealed} live_whole={} deltas={delta_total} \
steady_state={steady} saving_x={:.2} seal_ms={seal_ms:.0} \
supersede_ms_total={encode_ms_total:.0} read_all_ms={read_all_ms:.0}",
files[n - 1].1.len(),
raw as f64 / steady as f64
);
let _ = (compacted, packed);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn archive_reader_matches_free_functions() {
let dir = unique_dir("reader");
let archive = dir.join("store.znippy");
let files = synth(600, 13);
write_fresh(&archive, &files, ArrowIpcSink::new);
let reader = crate::ArchiveReader::open(&archive).unwrap();
assert_eq!(reader.row_count(), files.len(), "one chunk per synth file");
for (p, _) in files.iter().step_by(37) {
let cached = reader.locate(p);
let free = crate::locate_file(&archive, p).unwrap();
assert!(!cached.is_empty(), "cached reader failed to locate {p}");
assert_eq!(cached, free, "cached vs free locate diverged for {p}");
}
let missing = "repo/does/not/exist.bin";
assert!(reader.locate(missing).is_empty());
assert!(crate::locate_file(&archive, missing).unwrap().is_empty());
assert_eq!(
reader.files_meta(),
crate::get_all_files_meta(&archive).unwrap(),
"cached files_meta diverged from get_all_files_meta"
);
for prefix in ["repo/grp001/", "repo/grp0", "repo/", ""] {
assert_eq!(
reader.files_meta_with_prefix(prefix),
crate::get_files_meta_with_prefix(&archive, prefix).unwrap(),
"cached prefix meta diverged for {prefix:?}"
);
}
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn compaction_reclaims_the_superseded_blob_and_changes_no_entry() {
let dir = unique_dir("compact");
let archive = dir.join("c.znippy");
let mut st = 0x5151_2323_abcd_ef01u64;
let base: Vec<u8> = (0..600_000u32)
.map(|_| {
st = st.wrapping_add(0x9e37_79b9_7f4a_7c15);
let mut z = st;
z = (z ^ (z >> 30)).wrapping_mul(0xbf58_476d_1ce4_e5b9);
(z ^ (z >> 27)) as u8
})
.collect();
let mut gens: Vec<Vec<u8>> = vec![base.clone()];
for g in 1..4 {
let mut next = gens[g - 1].clone();
next.extend_from_slice(format!("generation {g} tail ").repeat(300).as_bytes());
gens.push(next);
}
let files: Vec<(String, Vec<u8>)> = gens
.iter()
.enumerate()
.map(|(i, b)| (format!("pack-{i}.pack"), b.clone()))
.collect();
create_archive(&archive, &files, 3).unwrap();
let raw: u64 = gens.iter().map(|g| g.len() as u64).sum();
let whole = std::fs::metadata(&archive).unwrap().len();
for i in (0..3).rev() {
let out = supersede_as_delta(
&archive,
&format!("pack-{i}.pack"),
&format!("pack-{}.pack", i + 1),
3 - 1 - i,
3,
)
.unwrap();
assert!(matches!(out, SupersedeOutcome::Delta { .. }), "gen {i}: {out:?}");
}
let after_supersede = std::fs::metadata(&archive).unwrap().len();
assert!(
after_supersede >= whole,
"supersede shrank the file ({whole} -> {after_supersede}); if that is now true, this \
test and everything built on the dead-payload finding wants revisiting"
);
let report = compact_archive(&archive).unwrap();
let compacted = std::fs::metadata(&archive).unwrap().len();
assert_eq!(report.bytes_after, compacted);
assert_eq!(report.rows, 4, "a compaction must not change the row count");
assert_eq!(report.delta_rows, 3, "the delta map must travel across it");
assert!(
compacted * 2 < raw,
"the compacted archive is {compacted} bytes for {raw} bytes of generations — the dead \
payload was not reclaimed"
);
let ar = ZnippyArchive::open(&archive).unwrap();
for (i, want) in gens.iter().enumerate() {
assert_eq!(
&ar.extract_file(&format!("pack-{i}.pack")).unwrap(),
want,
"generation {i} did not survive the compaction"
);
}
let map = read_delta_map(&archive).unwrap();
assert!(
map.iter().all(|(p, _, _)| p != "pack-3.pack"),
"compaction must not put the live generation behind a link: {map:?}"
);
let again = compact_archive(&archive).unwrap();
assert_eq!(again.rows, 4);
assert_eq!(again.delta_rows, 3);
let ar = ZnippyArchive::open(&archive).unwrap();
for (i, want) in gens.iter().enumerate() {
assert_eq!(&ar.extract_file(&format!("pack-{i}.pack")).unwrap(), want);
}
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn compaction_carries_the_reserved_logs() {
use crate::index::GUNNAR_REFS_MODULE;
use crate::meta_sink::{ReservedSection, ReservedSectionBuilder};
let dir = unique_dir("compact_reserved");
let archive = dir.join("r.znippy");
let refs_bytes = b"refs/heads/main 0123456789abcdef".to_vec();
{
use std::os::unix::fs::FileExt;
let f = Arc::new(File::create(&archive).unwrap());
let payload = b"an entry".to_vec();
f.write_all_at(&payload, 0).unwrap();
let blobs = vec![BlobMeta {
blob_offset: 0,
blob_size: payload.len() as u64,
chunk_meta: ChunkMeta {
fdata_offset: 0,
file_index: 0,
chunk_seq: 0,
checksum: *blake3::hash(&payload).as_bytes(),
compressed: false,
uncompressed_size: payload.len() as u64,
compressed_size: payload.len() as u64,
},
}];
let batch =
crate::index::build_metadata_batch(&blobs, |_| "obj/a.bin".to_string(), &[], &[])
.unwrap();
let mut sink = ArrowIpcSink::new(Arc::clone(&f), payload.len() as u64);
let carried = refs_bytes.clone();
let builder: ReservedSectionBuilder =
Box::new(move |_lookup| Ok(vec![ReservedSection::raw(GUNNAR_REFS_MODULE, carried)]));
sink = sink.with_reserved_builder(builder);
sink.push_subindex(
crate::index::lookup_schema().as_ref(),
&[batch],
GroupKey {
pkg_type: 0,
repo: String::new(),
module_name: String::new(),
},
)
.unwrap();
Box::new(sink).finish().unwrap();
}
assert_eq!(
read_reserved_section_bytes(&archive, GUNNAR_REFS_MODULE).unwrap(),
Some(refs_bytes.clone())
);
compact_archive(&archive).unwrap();
assert_eq!(
read_reserved_section_bytes(&archive, GUNNAR_REFS_MODULE).unwrap(),
Some(refs_bytes),
"the compaction dropped __gunnar_refs__"
);
let _ = std::fs::remove_dir_all(&dir);
}
}
pub(crate) fn decode_base_rows(bytes: &[u8]) -> Result<(Vec<String>, Vec<ChunkLoc>)> {
use arrow::ipc::reader::StreamReader;
let reader = StreamReader::try_new(std::io::Cursor::new(bytes), None)
.map_err(|e| anyhow!("append: lookup ipc reader: {e}"))?;
let mut paths = Vec::new();
let mut locs = Vec::new();
for batch in reader {
let batch = batch.map_err(|e| anyhow!("append: lookup batch decode: {e}"))?;
let get = |n: &str| batch.column_by_name(n)
.ok_or_else(|| anyhow!("append: lookup missing column {n}"));
let p = get("relative_path")?.as_any().downcast_ref::<StringArray>()
.ok_or_else(|| anyhow!("relative_path type"))?;
let seq = get("chunk_seq")?.as_any().downcast_ref::<UInt32Array>()
.ok_or_else(|| anyhow!("chunk_seq type"))?;
let fdata = get("fdata_offset")?.as_any().downcast_ref::<UInt64Array>()
.ok_or_else(|| anyhow!("fdata_offset type"))?;
let comp = get("compressed")?.as_any().downcast_ref::<BooleanArray>()
.ok_or_else(|| anyhow!("compressed type"))?;
let usz = get("uncompressed_size")?.as_any().downcast_ref::<UInt64Array>()
.ok_or_else(|| anyhow!("uncompressed_size type"))?;
let boff = get("blob_offset")?.as_any().downcast_ref::<UInt64Array>()
.ok_or_else(|| anyhow!("blob_offset type"))?;
let bsz = get("blob_size")?.as_any().downcast_ref::<UInt64Array>()
.ok_or_else(|| anyhow!("blob_size type"))?;
let ck = get("checksum")?.as_any().downcast_ref::<FixedSizeBinaryArray>()
.ok_or_else(|| anyhow!("checksum type"))?;
for i in 0..batch.num_rows() {
let mut c = [0u8; 32];
c.copy_from_slice(ck.value(i));
paths.push(p.value(i).to_string());
locs.push(ChunkLoc {
chunk_seq: seq.value(i),
fdata_offset: fdata.value(i),
blob_offset: boff.value(i),
blob_size: bsz.value(i),
uncompressed_size: usz.value(i),
compressed: comp.value(i),
checksum: c,
});
}
}
Ok((paths, locs))
}