use anyhow::{Result, Context};
use std::path::PathBuf;
use std::sync::Arc;
use traits::{ConnectorTrait, RemoteAsset};
use znippy_compress::stream_packer::{ArchiveEntry, StreamCompressor};
const MAX_TAR_ENTRY_BYTES: u64 = 8 * 1024 * 1024 * 1024;
const UPLOAD_CONCURRENCY: usize = 8;
enum UploadOutcome {
Ok(String),
Failed(String, String),
}
enum TarJob {
Upload { path: String, data: Vec<u8> },
Fatal(anyhow::Error),
}
fn tally_upload(outcome: UploadOutcome, pushed: &mut usize, failed: &mut usize) {
match outcome {
UploadOutcome::Ok(p) => {
*pushed += 1;
println!(" ✓ {p}");
}
UploadOutcome::Failed(p, e) => {
*failed += 1;
println!(" ✗ {p}: {e}");
}
}
}
fn read_capped<R: std::io::Read>(reader: &mut R, max: u64, what: &str) -> Result<Vec<u8>> {
let mut data = Vec::new();
let mut limited = std::io::Read::take(reader, max.saturating_add(1));
std::io::Read::read_to_end(&mut limited, &mut data)
.with_context(|| format!("Failed to read entry: {what}"))?;
if data.len() as u64 > max {
anyhow::bail!(
"tar entry {what} exceeds per-entry cap of {max} bytes (decompression bomb?)"
);
}
Ok(data)
}
pub fn pkg_type_from_format(format: &str) -> i8 {
use traits::ArtifactFormat;
ArtifactFormat::from_format_str(format)
.map(|f| f.znippy_type_id())
.unwrap_or(0)
}
fn safe_join(base: &std::path::Path, rel: &str) -> Result<PathBuf> {
use std::path::Component;
let rel_path = std::path::Path::new(rel);
let mut out = base.to_path_buf();
for comp in rel_path.components() {
match comp {
Component::Normal(c) => out.push(c),
Component::CurDir => {}
Component::ParentDir | Component::RootDir | Component::Prefix(_) => {
anyhow::bail!("refusing unsafe asset path {rel:?} (escapes {})", base.display());
}
}
}
Ok(out)
}
fn is_safe_member_path(rel: &str) -> bool {
use std::path::Component;
std::path::Path::new(rel).components().all(|c| {
matches!(c, Component::Normal(_) | Component::CurDir)
})
}
pub async fn feed_repo_to_sender(
connector: Arc<dyn ConnectorTrait>,
assets: Vec<RemoteAsset>,
sender: crossbeam_channel::Sender<ArchiveEntry>,
pkg_type: i8,
repo_name: String,
) -> Result<(usize, usize)> {
let total = assets.len();
let mut downloaded = 0usize;
let mut failed = 0usize;
for (i, asset) in assets.iter().enumerate() {
print!(" [{}/{}] {} ... ", i + 1, total, asset.path);
match connector.download_asset(asset).await {
Ok(bytes) => {
let n = bytes.len();
sender.send(ArchiveEntry {
relative_path: asset.path.clone(),
data: bytes,
pkg_type: Some(pkg_type),
repo: Some(repo_name.clone()),
}).context("compressor channel closed")?;
downloaded += 1;
println!("✓ ({} bytes)", n);
}
Err(e) => {
failed += 1;
println!("✗ {}", e);
}
}
}
Ok((downloaded, failed))
}
pub async fn dump_all_repos_to_znippy(
connector: Arc<dyn ConnectorTrait>,
repos: Vec<traits::RemoteRepository>,
output: PathBuf,
sequential: bool,
) -> Result<()> {
use znippy_compress::stream_packer::compress_stream;
println!("📦 Creating znippy archive: {}", output.display());
let compressor: StreamCompressor = compress_stream(&output, false)?;
let total_downloaded;
let total_failed;
if sequential {
let sender = compressor.sender().clone();
let mut dl = 0usize;
let mut fail = 0usize;
for repo in &repos {
let pkg_type = pkg_type_from_format(&repo.format);
println!("\n→ {} (format: {}, pkg_type: {})", repo.name, repo.format, pkg_type);
let assets = connector.list_assets(&repo.name).await?;
println!(" {} assets", assets.len());
let (d, f) = feed_repo_to_sender(
connector.clone(), assets, sender.clone(), pkg_type, repo.name.clone(),
).await?;
dl += d;
fail += f;
}
drop(sender);
total_downloaded = dl;
total_failed = fail;
} else {
let mut tasks = Vec::new();
for repo in &repos {
let pkg_type = pkg_type_from_format(&repo.format);
println!("\n→ {} (format: {}, pkg_type: {})", repo.name, repo.format, pkg_type);
let assets = connector.list_assets(&repo.name).await?;
println!(" {} assets queued", assets.len());
let conn = connector.clone();
let sender = compressor.sender().clone();
let repo_name = repo.name.clone();
tasks.push(tokio::spawn(async move {
feed_repo_to_sender(conn, assets, sender, pkg_type, repo_name).await
}));
}
let mut dl = 0usize;
let mut fail = 0usize;
for task in tasks {
let (d, f) = task.await
.context("repo task panicked")?
.context("repo download failed")?;
dl += d;
fail += f;
}
total_downloaded = dl;
total_failed = fail;
}
let report = compressor.finish()?;
println!("\n✅ Archive complete: {} downloaded, {} failed ({} bytes compressed)",
total_downloaded, total_failed, report.compressed_bytes);
Ok(())
}
pub async fn transfer_to_znippy(
connector: &dyn ConnectorTrait,
assets: Vec<RemoteAsset>,
output: PathBuf,
) -> Result<()> {
use znippy_compress::stream_packer::compress_stream;
println!("📦 Creating znippy archive: {}", output.display());
let compressor = compress_stream(&output, false)?;
let sender = compressor.sender().clone();
let total = assets.len();
let mut downloaded = 0usize;
for (i, asset) in assets.iter().enumerate() {
print!(" [{}/{}] {} ... ", i + 1, total, asset.path);
match connector.download_asset(asset).await {
Ok(bytes) => {
let n = bytes.len();
sender.send(ArchiveEntry {
relative_path: asset.path.clone(),
data: bytes,
pkg_type: None,
repo: None,
}).context("Failed to send to compressor")?;
downloaded += 1;
println!("✓ ({} bytes)", n);
}
Err(e) => {
println!("✗ {}", e);
}
}
}
drop(sender);
let report = compressor.finish()?;
println!("\n✅ Znippy archive complete: {} assets packed ({} bytes compressed)",
downloaded, report.compressed_bytes);
Ok(())
}
fn pkg_type_for_file(path: &std::path::Path) -> Option<i8> {
let name = path.file_name()?.to_str()?;
let fmt = match name.rsplit_once('.').map(|(_, e)| e)? {
"crate" => "rust",
"jar" | "pom" | "war" | "ear" => "maven",
"whl" => "python",
"nupkg" => "nuget",
"gem" => "gem",
_ => return None,
};
match pkg_type_from_format(fmt) {
0 => None,
t => Some(t),
}
}
pub async fn pack_directory_to_znippy(
directory: PathBuf,
output: PathBuf,
format: Option<&str>,
) -> Result<()> {
use znippy_compress::stream_packer::compress_stream;
if !directory.is_dir() {
anyhow::bail!("not a directory: {}", directory.display());
}
let forced = format.map(pkg_type_from_format).filter(|t| *t != 0);
let mut files = Vec::new();
walk_all_files(&directory, &mut files)?;
files.sort();
if files.is_empty() {
anyhow::bail!("no files found under {}", directory.display());
}
println!("📦 Packing {} files from {} → {}", files.len(), directory.display(), output.display());
let compressor = compress_stream(&output, false)?;
let sender = compressor.sender().clone();
let total = files.len();
let mut packed = 0usize;
for (i, path) in files.iter().enumerate() {
let rel = path
.strip_prefix(&directory)
.unwrap_or(path)
.to_string_lossy()
.replace('\\', "/");
let pkg_type = forced.or_else(|| pkg_type_for_file(path));
let data = std::fs::read(path)
.with_context(|| format!("Failed to read {}", path.display()))?;
let n = data.len();
sender
.send(ArchiveEntry { relative_path: rel.clone(), data, pkg_type, repo: None })
.context("Failed to send to compressor")?;
packed += 1;
println!(" [{}/{}] {} ({} bytes, pkg_type={:?})", i + 1, total, rel, n, pkg_type);
}
drop(sender);
let report = compressor.finish()?;
println!(
"\n✅ Znippy archive complete: {} files packed ({} bytes compressed)",
packed, report.compressed_bytes
);
Ok(())
}
fn walk_all_files(dir: &std::path::Path, files: &mut Vec<PathBuf>) -> Result<()> {
if !dir.is_dir() {
return Ok(());
}
for entry in std::fs::read_dir(dir)? {
let entry = entry?;
let path = entry.path();
if path.is_dir() {
walk_all_files(&path, files)?;
} else if path.is_file() {
files.push(path);
}
}
Ok(())
}
pub async fn transfer_to_tar_zstd(
connector: &dyn ConnectorTrait,
assets: Vec<RemoteAsset>,
output: PathBuf,
) -> Result<()> {
use std::io::BufWriter;
println!("📦 Creating tar.zst archive: {}", output.display());
let file = std::fs::File::create(&output)
.with_context(|| format!("Failed to create output file: {}", output.display()))?;
let zstd_encoder = zstd::Encoder::new(BufWriter::new(file), 3)
.context("Failed to create zstd encoder")?;
let mut tar_builder = tar::Builder::new(zstd_encoder);
let total = assets.len();
let mut downloaded = 0usize;
let mut failed = 0usize;
for (i, asset) in assets.iter().enumerate() {
print!(" [{}/{}] {} ... ", i + 1, total, asset.path);
match connector.download_asset(asset).await {
Ok(bytes) => {
let mut header = tar::Header::new_gnu();
header.set_size(bytes.len() as u64);
header.set_mode(0o644);
header.set_cksum();
tar_builder.append_data(&mut header, &asset.path, bytes.as_slice())
.with_context(|| format!("Failed to append {} to tar", asset.path))?;
downloaded += 1;
println!("✓ ({} bytes)", bytes.len());
}
Err(e) => {
failed += 1;
println!("✗ {}", e);
}
}
}
let zstd_encoder = tar_builder.into_inner()
.context("Failed to finish tar archive")?;
zstd_encoder.finish()
.context("Failed to finish zstd compression")?;
println!("\n✅ tar.zst archive complete: {} assets packed, {} failed", downloaded, failed);
Ok(())
}
pub async fn transfer_to_directory(
connector: &dyn ConnectorTrait,
assets: Vec<RemoteAsset>,
output: PathBuf,
) -> Result<()> {
println!("📁 Downloading to directory: {}", output.display());
std::fs::create_dir_all(&output)
.with_context(|| format!("Failed to create output directory: {}", output.display()))?;
let total = assets.len();
let mut downloaded = 0usize;
let mut failed = 0usize;
for (i, asset) in assets.iter().enumerate() {
let dest = match safe_join(&output, &asset.path) {
Ok(d) => d,
Err(e) => {
failed += 1;
println!(" [{}/{}] {} ... ✗ {}", i + 1, total, asset.path, e);
continue;
}
};
if let Some(parent) = dest.parent() {
std::fs::create_dir_all(parent)?;
}
print!(" [{}/{}] {} ... ", i + 1, total, asset.path);
match connector.download_asset(asset).await {
Ok(bytes) => {
std::fs::write(&dest, &bytes)
.with_context(|| format!("Failed to write {}", dest.display()))?;
downloaded += 1;
println!("✓ ({} bytes)", bytes.len());
}
Err(e) => {
failed += 1;
println!("✗ {}", e);
}
}
}
println!("\n✅ Done! {} downloaded, {} failed", downloaded, failed);
Ok(())
}
pub async fn push_directory_to_connector(
connector: &dyn ConnectorTrait,
directory: PathBuf,
repository: &str,
) -> Result<()> {
println!("📦 Pushing from directory: {}", directory.display());
let crate_files = find_crate_files(&directory)?;
println!(" Found {} .crate files to push", crate_files.len());
if crate_files.is_empty() {
println!(" Nothing to push.");
return Ok(());
}
let mut pushed = 0usize;
let mut failed = 0usize;
let jobs = crate_files.iter().map(move |crate_path| async move {
let upload_path = extract_crate_relative_path(crate_path)
.to_string_lossy()
.into_owned();
let data = std::fs::read(crate_path)
.with_context(|| format!("Failed to read {}", crate_path.display()))?;
anyhow::Ok(match connector.upload_asset(repository, &upload_path, &data).await {
Ok(()) => UploadOutcome::Ok(upload_path),
Err(e) => UploadOutcome::Failed(upload_path, e.to_string()),
})
});
znippy_zoomies::gatling::io::run(jobs, UPLOAD_CONCURRENCY, |outcome| {
tally_upload(outcome, &mut pushed, &mut failed);
Ok(())
})
.await?;
println!("\n✅ Push complete: {} succeeded, {} failed", pushed, failed);
Ok(())
}
pub async fn push_znippy_to_connector(
connector: &dyn ConnectorTrait,
archive: PathBuf,
repository: &str,
) -> Result<()> {
use znippy_common::{ZnippyArchive, ZnippyReader};
println!("📦 Reading znippy archive: {}", archive.display());
let znippy = ZnippyArchive::open(&archive)?;
let files = znippy.list_files()?;
println!(" Found {} files in archive", files.len());
let mut pushed = 0usize;
let mut failed = 0usize;
let znippy = &znippy;
let jobs = files.iter().map(move |path| async move {
if !is_safe_member_path(path) {
return anyhow::Ok(UploadOutcome::Failed(
path.clone(),
"unsafe member path (absolute or '..') — skipped".to_string(),
));
}
anyhow::Ok(match znippy.extract_file(path) {
Ok(data) => match connector.upload_asset(repository, path, &data).await {
Ok(()) => UploadOutcome::Ok(path.clone()),
Err(e) => UploadOutcome::Failed(path.clone(), e.to_string()),
},
Err(e) => UploadOutcome::Failed(path.clone(), format!("extract: {e}")),
})
});
znippy_zoomies::gatling::io::run(jobs, UPLOAD_CONCURRENCY, |outcome| {
tally_upload(outcome, &mut pushed, &mut failed);
Ok(())
})
.await?;
println!("\n✅ Push complete: {} succeeded, {} failed", pushed, failed);
Ok(())
}
pub async fn push_tar_zstd_to_connector(
connector: &dyn ConnectorTrait,
archive: PathBuf,
repository: &str,
) -> Result<()> {
use std::io::BufReader;
println!("📦 Reading tar.zst archive: {}", archive.display());
let file = std::fs::File::open(&archive)
.with_context(|| format!("Failed to open archive: {}", archive.display()))?;
let zstd_decoder = zstd::Decoder::new(BufReader::new(file))
.context("Failed to create zstd decoder")?;
let mut tar_archive = tar::Archive::new(zstd_decoder);
let mut entries = tar_archive.entries().context("Failed to read tar entries")?;
let mut pushed = 0usize;
let mut failed = 0usize;
let jobs = std::iter::from_fn(move || loop {
let mut entry = match entries.next()? {
Ok(e) => e,
Err(e) => {
return Some(TarJob::Fatal(
anyhow::Error::new(e).context("Failed to read tar entry"),
));
}
};
if entry.header().entry_type() != tar::EntryType::Regular {
continue;
}
let path = match entry.path() {
Ok(p) => p.to_string_lossy().into_owned(),
Err(e) => {
return Some(TarJob::Fatal(
anyhow::Error::new(e).context("Failed to read entry path"),
));
}
};
return Some(match read_capped(&mut entry, MAX_TAR_ENTRY_BYTES, &path) {
Ok(data) => TarJob::Upload { path, data },
Err(e) => TarJob::Fatal(e),
});
})
.map(move |job| async move {
match job {
TarJob::Fatal(e) => Err(e),
TarJob::Upload { path, data } => {
if !is_safe_member_path(&path) {
return anyhow::Ok(UploadOutcome::Failed(
path,
"unsafe member path (absolute or '..') — skipped".to_string(),
));
}
anyhow::Ok(
match connector.upload_asset(repository, &path, &data).await {
Ok(()) => UploadOutcome::Ok(path),
Err(e) => UploadOutcome::Failed(path, e.to_string()),
},
)
}
}
});
znippy_zoomies::gatling::io::run(jobs, UPLOAD_CONCURRENCY, |outcome| {
tally_upload(outcome, &mut pushed, &mut failed);
Ok(())
})
.await?;
println!("\n✅ Push complete: {} succeeded, {} failed", pushed, failed);
Ok(())
}
pub async fn unpack_znippy_to_directory(archive: PathBuf, output: PathBuf) -> Result<()> {
use znippy_common::{ZnippyArchive, ZnippyReader};
println!("📦 Reading znippy archive: {}", archive.display());
let znippy = ZnippyArchive::open(&archive)?;
let files = znippy.list_files()?;
println!(" Found {} files in archive → {}", files.len(), output.display());
std::fs::create_dir_all(&output)
.with_context(|| format!("Failed to create output directory: {}", output.display()))?;
let total = files.len();
let mut written = 0usize;
let mut failed = 0usize;
for (i, path) in files.iter().enumerate() {
let dest = match safe_join(&output, path) {
Ok(d) => d,
Err(e) => {
failed += 1;
println!(" [{}/{}] {} ... ✗ {}", i + 1, total, path, e);
continue;
}
};
match znippy.extract_file(path) {
Ok(data) => {
if let Some(parent) = dest.parent() {
std::fs::create_dir_all(parent)?;
}
std::fs::write(&dest, &data)
.with_context(|| format!("Failed to write {}", dest.display()))?;
written += 1;
println!(" [{}/{}] {} ✓ ({} bytes)", i + 1, total, path, data.len());
}
Err(e) => {
failed += 1;
println!(" [{}/{}] {} ... ✗ extract: {}", i + 1, total, path, e);
}
}
}
println!("\n✅ Unpacked znippy: {} files written, {} failed", written, failed);
if failed > 0 {
anyhow::bail!("{failed} member(s) failed to extract from {}", archive.display());
}
Ok(())
}
pub async fn unpack_tar_zstd_to_directory(archive: PathBuf, output: PathBuf) -> Result<()> {
use std::io::BufReader;
println!("📦 Reading tar.zst archive: {} → {}", archive.display(), output.display());
std::fs::create_dir_all(&output)
.with_context(|| format!("Failed to create output directory: {}", output.display()))?;
let file = std::fs::File::open(&archive)
.with_context(|| format!("Failed to open archive: {}", archive.display()))?;
let zstd_decoder = zstd::Decoder::new(BufReader::new(file))
.context("Failed to create zstd decoder")?;
let mut tar_archive = tar::Archive::new(zstd_decoder);
let mut written = 0usize;
let mut failed = 0usize;
for entry in tar_archive.entries().context("Failed to read tar entries")? {
let mut entry = entry.context("Failed to read tar entry")?;
if entry.header().entry_type() != tar::EntryType::Regular {
continue;
}
let path = entry
.path()
.context("Failed to read entry path")?
.to_string_lossy()
.into_owned();
let dest = match safe_join(&output, &path) {
Ok(d) => d,
Err(e) => {
failed += 1;
println!(" {} ... ✗ {}", path, e);
continue;
}
};
let data = read_capped(&mut entry, MAX_TAR_ENTRY_BYTES, &path)?;
if let Some(parent) = dest.parent() {
std::fs::create_dir_all(parent)?;
}
std::fs::write(&dest, &data)
.with_context(|| format!("Failed to write {}", dest.display()))?;
written += 1;
println!(" {} ✓ ({} bytes)", path, data.len());
}
println!("\n✅ Unpacked tar.zst: {} files written, {} skipped", written, failed);
Ok(())
}
pub async fn pack_directory_to_tar_zstd(directory: PathBuf, output: PathBuf) -> Result<()> {
use std::io::BufWriter;
if !directory.is_dir() {
anyhow::bail!("not a directory: {}", directory.display());
}
let mut files = Vec::new();
walk_all_files(&directory, &mut files)?;
files.sort();
if files.is_empty() {
anyhow::bail!("no files found under {}", directory.display());
}
println!("📦 Packing {} files from {} → {}", files.len(), directory.display(), output.display());
let file = std::fs::File::create(&output)
.with_context(|| format!("Failed to create output file: {}", output.display()))?;
let zstd_encoder = zstd::Encoder::new(BufWriter::new(file), 3)
.context("Failed to create zstd encoder")?;
let mut tar_builder = tar::Builder::new(zstd_encoder);
let total = files.len();
for (i, path) in files.iter().enumerate() {
let rel = path
.strip_prefix(&directory)
.unwrap_or(path)
.to_string_lossy()
.replace('\\', "/");
let data = std::fs::read(path)
.with_context(|| format!("Failed to read {}", path.display()))?;
let mut header = tar::Header::new_gnu();
header.set_size(data.len() as u64);
header.set_mode(0o644);
header.set_cksum();
tar_builder
.append_data(&mut header, &rel, data.as_slice())
.with_context(|| format!("Failed to append {} to tar", rel))?;
println!(" [{}/{}] {} ({} bytes)", i + 1, total, rel, data.len());
}
let zstd_encoder = tar_builder.into_inner().context("Failed to finish tar archive")?;
zstd_encoder.finish().context("Failed to finish zstd compression")?;
println!("\n✅ tar.zst archive complete: {} files packed", total);
Ok(())
}
pub async fn convert_znippy_to_tar_zstd(input: PathBuf, output: PathBuf) -> Result<()> {
let staging = tempfile::tempdir().context("create staging dir for znippy→tar.zst")?;
unpack_znippy_to_directory(input, staging.path().to_path_buf()).await?;
pack_directory_to_tar_zstd(staging.path().to_path_buf(), output).await
}
pub async fn convert_tar_zstd_to_znippy(
input: PathBuf,
output: PathBuf,
format: Option<&str>,
) -> Result<()> {
let staging = tempfile::tempdir().context("create staging dir for tar.zst→znippy")?;
unpack_tar_zstd_to_directory(input, staging.path().to_path_buf()).await?;
pack_directory_to_znippy(staging.path().to_path_buf(), output, format).await
}
fn extract_crate_relative_path(path: &std::path::Path) -> PathBuf {
let components: Vec<_> = path.components().collect();
for (i, comp) in components.iter().enumerate() {
if comp.as_os_str() == "crates" {
return components[i..].iter().collect();
}
}
PathBuf::from(path.file_name().unwrap_or_default())
}
fn find_crate_files(dir: &std::path::Path) -> Result<Vec<PathBuf>> {
let mut files = Vec::new();
walk_dir_recursive(dir, &mut files)?;
files.sort();
Ok(files)
}
fn walk_dir_recursive(dir: &std::path::Path, files: &mut Vec<PathBuf>) -> Result<()> {
if !dir.is_dir() {
return Ok(());
}
for entry in std::fs::read_dir(dir)? {
let entry = entry?;
let path = entry.path();
if path.is_dir() {
walk_dir_recursive(&path, files)?;
} else if path.extension().and_then(|e| e.to_str()) == Some("crate") {
files.push(path);
}
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn safe_join_allows_normal_relative_paths() {
let base = std::path::Path::new("/srv/out");
let j = safe_join(base, "crates/se/serde/serde-1.0.0.crate").unwrap();
assert_eq!(j, base.join("crates/se/serde/serde-1.0.0.crate"));
assert_eq!(safe_join(base, "./a/b").unwrap(), base.join("a/b"));
}
#[test]
fn safe_join_rejects_hostile_upstream_paths() {
let base = std::path::Path::new("/srv/out");
for hostile in [
"../../../../etc/cron.d/holger",
"ok/../../../../etc/passwd",
"/etc/passwd",
"/etc/cron.d/x",
] {
assert!(
safe_join(base, hostile).is_err(),
"must reject hostile asset path {hostile:?}"
);
}
assert!(safe_join(base, "a/../b").is_err()); }
#[test]
fn pkg_type_for_file_agrees_with_znippy_plugin_extensions() {
use znippy_common::plugin::ArchiveTypePlugin;
use znippy_common::plugins::cargo_native::CargoPlugin;
use znippy_common::plugins::gem_native::GemPlugin;
use znippy_common::plugins::skeletons;
let cargo = CargoPlugin::new();
let gem = GemPlugin;
let nuget = skeletons::NugetPlugin;
let cases: Vec<(&str, i8, Vec<String>)> = vec![
("serde-1.0.0.crate", cargo.type_id(), cargo.meta().extensions),
("rails-7.1.0.gem", gem.type_id(), gem.meta().extensions),
("Newtonsoft.Json.13.0.3.nupkg", nuget.type_id(), nuget.meta().extensions),
];
for (filename, znippy_id, exts) in &cases {
assert!(
exts.iter().any(|e| filename.ends_with(e.as_str())),
"{filename}: not covered by znippy handler claims {exts:?}"
);
assert_eq!(
pkg_type_for_file(std::path::Path::new(filename)),
Some(*znippy_id),
"{filename}: holger pkg_type drifted from znippy handler type_id {znippy_id}"
);
}
assert_eq!(
pkg_type_for_file(std::path::Path::new("guava-33.0.jar")),
Some(traits::ArtifactFormat::Maven3.znippy_type_id()),
);
assert_eq!(
pkg_type_for_file(std::path::Path::new("numpy-1.26-cp311.whl")),
Some(traits::ArtifactFormat::Pip.znippy_type_id()),
);
for ambiguous in ["pkg-1.0.tar.gz", "pkg-1.0.tgz", "pkg-1.0.zip", "sym-1.0.snupkg"] {
assert_eq!(
pkg_type_for_file(std::path::Path::new(ambiguous)),
None,
"{ambiguous} must stay untyped (curated adapter policy)"
);
}
}
#[test]
fn read_capped_accepts_within_limit_and_rejects_bomb() {
let mut ok = std::io::Cursor::new(vec![7u8; 100]);
let out = read_capped(&mut ok, 128, "small").unwrap();
assert_eq!(out.len(), 100);
let mut at = std::io::Cursor::new(vec![1u8; 128]);
assert_eq!(read_capped(&mut at, 128, "edge").unwrap().len(), 128);
let mut bomb = std::io::Cursor::new(vec![0u8; 1000]);
let err = read_capped(&mut bomb, 128, "bomb").expect_err("must reject over-cap entry");
assert!(err.to_string().contains("exceeds per-entry cap"), "got: {err}");
}
fn seed_fixture_tree(root: &std::path::Path) -> Vec<(String, Vec<u8>)> {
let fixtures = vec![
("crates/se/serde/serde-1.0.0.crate".to_string(), b"serde crate bytes \x00\x01\x02".to_vec()),
("crates/to/tokio/tokio-1.2.3.crate".to_string(), vec![0xABu8; 4096]),
("readme.txt".to_string(), b"plain top-level file".to_vec()),
];
for (rel, data) in &fixtures {
let dest = root.join(rel);
std::fs::create_dir_all(dest.parent().unwrap()).unwrap();
std::fs::write(&dest, data).unwrap();
}
fixtures
}
fn assert_tree_recovered(dir: &std::path::Path, fixtures: &[(String, Vec<u8>)]) {
for (rel, data) in fixtures {
let got = std::fs::read(dir.join(rel))
.unwrap_or_else(|e| panic!("recovered file {rel} missing: {e}"));
assert_eq!(&got, data, "file {rel} must round-trip byte-for-byte");
}
}
#[tokio::test]
async fn znippy_directory_round_trip_is_byte_exact() {
let work = tempfile::tempdir().unwrap();
let src = work.path().join("src");
std::fs::create_dir_all(&src).unwrap();
let fixtures = seed_fixture_tree(&src);
let archive = work.path().join("out.znippy");
pack_directory_to_znippy(src, archive.clone(), None)
.await
.expect("pack directory → znippy");
assert!(archive.exists(), "znippy archive was written");
let out = work.path().join("out");
unpack_znippy_to_directory(archive, out.clone())
.await
.expect("unpack znippy → directory");
assert_tree_recovered(&out, &fixtures);
}
#[tokio::test]
async fn tar_zstd_directory_round_trip_is_byte_exact() {
let work = tempfile::tempdir().unwrap();
let src = work.path().join("src");
std::fs::create_dir_all(&src).unwrap();
let fixtures = seed_fixture_tree(&src);
let archive = work.path().join("out.tar.zst");
pack_directory_to_tar_zstd(src, archive.clone())
.await
.expect("pack directory → tar.zst");
assert!(archive.exists(), "tar.zst archive was written");
let out = work.path().join("out");
unpack_tar_zstd_to_directory(archive, out.clone())
.await
.expect("unpack tar.zst → directory");
assert_tree_recovered(&out, &fixtures);
}
#[tokio::test]
async fn cross_archive_conversions_are_byte_exact() {
let work = tempfile::tempdir().unwrap();
let src = work.path().join("src");
std::fs::create_dir_all(&src).unwrap();
let fixtures = seed_fixture_tree(&src);
let znippy = work.path().join("a.znippy");
let tarzst = work.path().join("a.tar.zst");
pack_directory_to_znippy(src.clone(), znippy.clone(), None).await.unwrap();
pack_directory_to_tar_zstd(src, tarzst.clone()).await.unwrap();
let z2t = work.path().join("z2t.tar.zst");
convert_znippy_to_tar_zstd(znippy, z2t.clone()).await.expect("znippy → tar.zst");
let out1 = work.path().join("out1");
unpack_tar_zstd_to_directory(z2t, out1.clone()).await.unwrap();
assert_tree_recovered(&out1, &fixtures);
let t2z = work.path().join("t2z.znippy");
convert_tar_zstd_to_znippy(tarzst, t2z.clone(), None).await.expect("tar.zst → znippy");
let out2 = work.path().join("out2");
unpack_znippy_to_directory(t2z, out2.clone()).await.unwrap();
assert_tree_recovered(&out2, &fixtures);
}
}