diff --git a/Cargo.lock b/Cargo.lock
index 0fb09e78e..82f66641e 100644
--- a/Cargo.lock
+++ b/Cargo.lock
@@ -10489,15 +10489,10 @@ dependencies = [
name = "rustfs-zip"
version = "1.0.0-rc.2"
dependencies = [
- "astral-tokio-tar",
"async-compression",
- "criterion",
"hotpath",
- "tempfile",
"thiserror 2.0.20",
"tokio",
- "tokio-stream",
- "zip",
]
[[package]]
diff --git a/crates/zip/Cargo.toml b/crates/zip/Cargo.toml
index 54cdd9547..01ec5d84e 100644
--- a/crates/zip/Cargo.toml
+++ b/crates/zip/Cargo.toml
@@ -20,7 +20,7 @@ repository.workspace = true
rust-version.workspace = true
version.workspace = true
homepage.workspace = true
-description = "ZIP file handling for RustFS, providing support for reading and writing ZIP archives."
+description = "Archive format detection and async stream decoders for RustFS."
keywords = ["zip", "compression", "rustfs", "Minio"]
categories = ["web-programming", "development-tools", "compression"]
documentation = "https://docs.rs/rustfs-zip/latest/rustfs_zip/"
@@ -28,10 +28,6 @@ documentation = "https://docs.rs/rustfs-zip/latest/rustfs_zip/"
[lib]
doctest = false
-[[bench]]
-name = "zip_benchmark"
-harness = false
-
[features]
default = []
hotpath = ["hotpath/hotpath", "hotpath/tokio"]
@@ -48,16 +44,8 @@ async-compression = { workspace = true, features = [
"zstd",
"xz",
] }
-tokio = { workspace = true, features = ["fs", "io-util", "macros", "rt-multi-thread"] }
-tokio-stream = { workspace = true }
-astral-tokio-tar = { workspace = true }
+tokio = { workspace = true, features = ["io-util", "macros", "rt"] }
thiserror = { workspace = true }
-zip = { workspace = true }
-
-[dev-dependencies]
-criterion = { workspace = true, features = ["html_reports"] }
-tempfile = { workspace = true }
-
[lints]
workspace = true
diff --git a/crates/zip/README.md b/crates/zip/README.md
index 78e105f5f..a9b7865d5 100644
--- a/crates/zip/README.md
+++ b/crates/zip/README.md
@@ -1,9 +1,9 @@
[](https://rustfs.com)
-# RustFS Zip - Archive And Compression Primitives
+# RustFS Zip - Archive Format Detection And Stream Decoding
- High-performance compression and archiving for RustFS object storage
+ Archive format detection and async stream decoders for RustFS object storage
@@ -17,53 +17,23 @@
## 📖 Overview
-**RustFS Zip** provides archive and compression primitives for the [RustFS](https://rustfs.com) distributed object storage system. Today it is primarily used by RustFS archive extract flows to:
+**RustFS Zip** provides the archive primitives used by the [RustFS](https://rustfs.com) archive extract flow:
-- identify archive/compression formats by extension
-- stream tar and tar+compression inputs through async decoders
-- provide small ZIP read/write helpers for local archive workflows
+- identify a compression format from an archive extension
+- wrap an async reader in the matching stream decoder
+- carry the shared default archive guardrails
## Current Features
-- A clearer type model with:
- - `CompressionCodec` for stream codecs
- - `ArchiveKind` for container families
- - `ArchiveFormat` for concrete archive/container combinations
-- Async stream codecs for `gzip`, `bzip2`, `zlib`, `xz`, and `zstd`
-- Tar archive iteration over async readers through `read_archive_entries()` / `extract_tar_entries()`
-- Archive guardrails through `ArchiveLimits` for entry count, entry size, total unpacked size, and path length
-- In-memory compression helpers for payload round-trip workflows
-- Blocking ZIP create/extract helpers for local archive files
-- ZIP helper metadata via `ZipEntry`, including:
- - `compression_method`
- - `archive_kind`
- - `format`
- - `unix_mode`
-- ZIP helper options via `ZipWriteOptions`, including:
- - `compression_level`
- - `create_directory_entries`
-
-## Compatibility
-
-- `CompressionFormat` is retained as a compatibility layer for existing callers
-- New code should prefer `ArchiveFormat`, `ArchiveKind`, and `CompressionCodec` when expressing archive semantics
-
-## ZIP Helper Scope
-
-The file-based ZIP helper APIs are best suited for:
-
-- local archive import/export flows
-- admin-side packaging helpers
-- test fixtures and tooling
-
-They are not intended to be a remote streaming ZIP access engine.
+- `CompressionFormat::from_extension()` for extension-based format detection, including tar-family suffixes such as `tgz`, `tbz2`, `txz`, and `tzst`
+- `CompressionFormat::get_decoder()` for async stream decoding of `gzip`, `bzip2`, `zlib`, `xz`, and `zstd`, plus a pass-through reader for plain `tar`
+- `ArchiveLimits` with the default entry count, entry size, total unpacked size, and path length guardrails
## Current Boundaries
-- ZIP is supported via file-based helper APIs, not the tar-family async stream APIs
-- Tar-family stream APIs are intended for `tar`, `tar.gz`, `tar.bz2`, `tar.xz`, `tar.zst`, and similar compressed tar flows
-- Default archive guardrails are intentionally conservative and do not replace higher-level RustFS object-path validation
-- This crate does not currently implement a general-purpose parallel archive engine
+- ZIP has no stream decoder: `get_decoder()` rejects `CompressionFormat::Zip`, because ZIP needs central-directory semantics that a forward-only stream cannot provide
+- This crate detects formats and hands back decoders; archive iteration, entry writing, and extraction to disk belong to the caller
+- `ArchiveLimits` carries the values only; enforcement and the resulting protocol error belong to the caller
- Archive extraction safety policy remains the responsibility of the RustFS caller for object-store flows
## 📚 Documentation
diff --git a/crates/zip/benches/zip_benchmark.rs b/crates/zip/benches/zip_benchmark.rs
deleted file mode 100644
index c7047918b..000000000
--- a/crates/zip/benches/zip_benchmark.rs
+++ /dev/null
@@ -1,416 +0,0 @@
-// Copyright 2024 RustFS Team
-//
-// Licensed under the Apache License, Version 2.0 (the "License");
-// you may not use this file except in compliance with the License.
-// You may obtain a copy of the License at
-//
-// http://www.apache.org/licenses/LICENSE-2.0
-//
-// Unless required by applicable law or agreed to in writing, software
-// distributed under the License is distributed on an "AS IS" BASIS,
-// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
-// See the License for the specific language governing permissions and
-// limitations under the License.
-
-use criterion::{BenchmarkId, Criterion, Throughput, criterion_group, criterion_main};
-use rustfs_zip::{
- ArchiveLimits, CompressionFormat, CompressionLevel, ZipWriteOptions, create_zip_with_options, extract_tar_entries,
- extract_zip_to_path_with_limits, extract_zip_with_limits,
-};
-use std::hint::black_box;
-use std::sync::Arc;
-use std::sync::atomic::{AtomicUsize, Ordering};
-use tempfile::tempdir;
-use tokio::runtime::Builder;
-use tokio_tar::{Builder as TarBuilder, Header};
-use zip::ZipArchive;
-
-fn build_runtime() -> tokio::runtime::Runtime {
- Builder::new_current_thread()
- .enable_all()
- .build()
- .expect("build tokio runtime for rustfs-zip benchmarks")
-}
-
-async fn build_tar_payload(entry_count: usize, payload_size: usize) -> Vec {
- let sink = tokio::io::duplex(64 * 1024);
- let (writer, mut reader) = sink;
- let write_task = tokio::spawn(async move {
- let mut builder = TarBuilder::new(writer);
- let payload = vec![b'a'; payload_size];
- for index in 0..entry_count {
- let mut header = Header::new_gnu();
- header.set_size(payload.len() as u64);
- header.set_mode(0o644);
- header.set_cksum();
- builder
- .append_data(&mut header, format!("entry-{index}.txt"), &payload[..])
- .await
- .expect("append tar benchmark entry");
- }
- builder.finish().await.expect("finish tar benchmark archive");
- });
-
- let mut output = Vec::new();
- tokio::io::copy(&mut reader, &mut output)
- .await
- .expect("read tar benchmark archive");
- write_task.await.expect("join tar writer task");
- output
-}
-
-async fn build_compressed_tar_payload(format: CompressionFormat, entry_count: usize, payload_size: usize) -> Vec {
- let tar_payload = build_tar_payload(entry_count, payload_size).await;
- rustfs_zip::Compressor::new(format)
- .compress(&tar_payload)
- .await
- .expect("compress tar benchmark payload")
-}
-
-fn bench_tar_family_extract(c: &mut Criterion) {
- let runtime = build_runtime();
- let mut group = c.benchmark_group("zip_tar_family_extract");
-
- for (name, format, entry_count, payload_size) in [
- ("tar_gzip_small_many", CompressionFormat::Gzip, 64usize, 256usize),
- ("tar_zstd_medium", CompressionFormat::Zstd, 16usize, 16 * 1024usize),
- ] {
- let payload = runtime.block_on(build_compressed_tar_payload(format, entry_count, payload_size));
- group.throughput(Throughput::Bytes(payload.len() as u64));
- group.bench_with_input(BenchmarkId::new(name, payload.len()), &payload, |b, payload| {
- b.iter(|| {
- runtime.block_on(async {
- let seen = Arc::new(AtomicUsize::new(0));
- let seen_ref = Arc::clone(&seen);
- extract_tar_entries(std::io::Cursor::new(payload.clone()), format, move |_entry| {
- let seen_ref = Arc::clone(&seen_ref);
- async move {
- seen_ref.fetch_add(1, Ordering::Relaxed);
- Ok(())
- }
- })
- .await
- .expect("extract tar benchmark payload");
- black_box(seen.load(Ordering::Relaxed));
- });
- });
- });
- }
-
- group.finish();
-}
-
-fn bench_zip_helper_round_trip(c: &mut Criterion) {
- let runtime = build_runtime();
- let mut group = c.benchmark_group("zip_helper_round_trip");
-
- let zip_matrix = [
- ("stored_flat_32x128", CompressionLevel::Fastest, 32usize, 128usize, "flat"),
- ("stored_nested_32x256", CompressionLevel::Fastest, 32usize, 256usize, "nested"),
- ("stored_flat_256x128", CompressionLevel::Fastest, 256usize, 128usize, "flat"),
- ("deflated_flat_32x1k", CompressionLevel::Best, 32usize, 1024usize, "flat"),
- ("deflated_nested_256x1k", CompressionLevel::Best, 256usize, 1024usize, "nested"),
- ("deflated_deep_1024x4k", CompressionLevel::Best, 1024usize, 4 * 1024usize, "deep"),
- ];
-
- for (name, compression_level, file_count, payload_size, layout) in zip_matrix {
- let files = (0..file_count)
- .map(|index| {
- let path = match layout {
- "flat" => format!("file-{index}.txt"),
- "nested" => format!("batch-{}/file-{index}.txt", index % 8),
- "deep" => format!("lvl1/lvl2-{}/lvl3-{}/file-{index}.txt", index % 16, index % 32),
- _ => format!("file-{index}.txt"),
- };
- (path, vec![b'b'; payload_size])
- })
- .collect::>();
- let total_bytes = (file_count * payload_size) as u64;
- group.throughput(Throughput::Bytes(total_bytes));
-
- group.bench_with_input(BenchmarkId::new(name, total_bytes), &files, |b, files| {
- b.iter(|| {
- let temp = tempdir().expect("create benchmark tempdir");
- let zip_path = temp.path().join("archive.zip");
- let extract_path = temp.path().join("extract");
- runtime.block_on(async {
- create_zip_with_options(
- &zip_path,
- files.clone(),
- ZipWriteOptions {
- compression_level,
- create_directory_entries: true,
- },
- )
- .await
- .expect("create zip benchmark archive");
-
- let entries = extract_zip_with_limits(&zip_path, &extract_path, ArchiveLimits::default())
- .await
- .expect("extract zip benchmark archive");
- black_box(entries.len());
- });
- });
- });
- }
-
- group.finish();
-}
-
-fn bench_zip_helper_hotspot_breakdown(c: &mut Criterion) {
- let runtime = build_runtime();
- let mut group = c.benchmark_group("zip_helper_hotspot_breakdown");
- let files = (0..32)
- .map(|index| (format!("batch/file-{index}.txt"), vec![b'c'; 256]))
- .collect::>();
- let total_bytes = (32 * 256) as u64;
- group.throughput(Throughput::Bytes(total_bytes));
-
- group.bench_function("fs_setup_cleanup_only", |b| {
- b.iter(|| {
- let temp = tempdir().expect("create benchmark tempdir");
- let zip_path = temp.path().join("archive.zip");
- let extract_path = temp.path().join("extract");
- black_box((zip_path, extract_path));
- });
- });
-
- group.bench_function("zip_create_only_stored_small", |b| {
- b.iter(|| {
- let temp = tempdir().expect("create benchmark tempdir");
- let zip_path = temp.path().join("archive.zip");
- runtime.block_on(async {
- create_zip_with_options(
- &zip_path,
- files.clone(),
- ZipWriteOptions {
- compression_level: CompressionLevel::Fastest,
- create_directory_entries: true,
- },
- )
- .await
- .expect("create zip benchmark archive");
- });
- });
- });
-
- let payload_for_extract = {
- let temp = tempdir().expect("create benchmark tempdir");
- let zip_path = temp.path().join("archive.zip");
- runtime.block_on(async {
- create_zip_with_options(
- &zip_path,
- files.clone(),
- ZipWriteOptions {
- compression_level: CompressionLevel::Fastest,
- create_directory_entries: true,
- },
- )
- .await
- .expect("prepare zip benchmark extract payload");
- });
- std::fs::read(&zip_path).expect("read benchmark zip payload")
- };
-
- group.bench_function("zip_extract_only_stored_small", |b| {
- b.iter(|| {
- let temp = tempdir().expect("create benchmark tempdir");
- let zip_path = temp.path().join("archive.zip");
- let extract_path = temp.path().join("extract");
- std::fs::write(&zip_path, &payload_for_extract).expect("write benchmark zip payload");
- runtime.block_on(async {
- let entries = extract_zip_with_limits(&zip_path, &extract_path, ArchiveLimits::default())
- .await
- .expect("extract zip benchmark archive");
- black_box(entries.len());
- });
- });
- });
-
- group.bench_function("zip_extract_only_stored_small_summary_only", |b| {
- b.iter(|| {
- let temp = tempdir().expect("create benchmark tempdir");
- let zip_path = temp.path().join("archive.zip");
- let extract_path = temp.path().join("extract");
- std::fs::write(&zip_path, &payload_for_extract).expect("write benchmark zip payload");
- runtime.block_on(async {
- let summary = extract_zip_to_path_with_limits(&zip_path, &extract_path, ArchiveLimits::default())
- .await
- .expect("extract zip benchmark summary path");
- black_box(summary.entry_count);
- });
- });
- });
-
- group.bench_function("zip_reader_only_stored_small", |b| {
- b.iter(|| {
- let cursor = std::io::Cursor::new(payload_for_extract.clone());
- let mut archive = ZipArchive::new(cursor).expect("open zip archive for reader-only benchmark");
- let mut total_bytes = 0usize;
- for index in 0..archive.len() {
- let mut zip_file = archive.by_index(index).expect("access zip entry by index");
- let enclosed_name = zip_file
- .enclosed_name()
- .expect("resolve enclosed zip entry name")
- .to_string_lossy()
- .replace('\\', "/");
- let size = zip_file.size();
- assert!(!enclosed_name.is_empty(), "zip reader-only benchmark expects non-empty names");
- assert!(
- size <= ArchiveLimits::default().max_entry_size,
- "zip reader-only benchmark expects small entries"
- );
- if !zip_file.is_dir() {
- let mut sink = [0_u8; 256];
- let bytes_read =
- std::io::Read::read(&mut zip_file, &mut sink).expect("read zip entry payload for reader-only benchmark");
- total_bytes += bytes_read;
- }
- }
- black_box(total_bytes);
- });
- });
-
- group.bench_function("file_write_only_stored_small", |b| {
- b.iter(|| {
- let temp = tempdir().expect("create benchmark tempdir");
- let extract_path = temp.path().join("extract");
- std::fs::create_dir_all(&extract_path).expect("create extract dir for file-write-only benchmark");
- let mut total_bytes = 0usize;
- for index in 0..32 {
- let path = extract_path.join(format!("file-{index}.txt"));
- std::fs::write(&path, [b'c'; 256]).expect("write small file for file-write-only benchmark");
- total_bytes += 256;
- }
- black_box(total_bytes);
- });
- });
-
- group.finish();
-}
-
-fn build_object_archive_files(
- metadata_count: usize,
- metadata_size: usize,
- payload_count: usize,
- payload_size: usize,
-) -> Vec<(String, Vec)> {
- let mut files = Vec::with_capacity(metadata_count * 2 + payload_count);
-
- for index in 0..metadata_count {
- let key_prefix = format!(
- "bucket-a/shard-{}/tenant-{}/dataset-{}/object-{index:04}",
- index % 8,
- index % 16,
- index % 32
- );
- files.push((
- format!("{key_prefix}/meta.json"),
- format!(
- "{{\"key\":\"object-{index:04}\",\"etag\":\"{:032x}\",\"size\":{},\"content_type\":\"application/octet-stream\"}}",
- index,
- payload_size
- )
- .into_bytes(),
- ));
- files.push((format!("{key_prefix}/tags.txt"), vec![b'm'; metadata_size]));
- }
-
- for index in 0..payload_count {
- let payload_prefix = format!(
- "bucket-a/shard-{}/tenant-{}/dataset-{}/object-{index:04}",
- index % 8,
- index % 16,
- index % 32
- );
- files.push((format!("{payload_prefix}/part-00000.bin"), vec![b'p'; payload_size]));
- }
-
- files
-}
-
-fn bench_zip_object_archive_extract(c: &mut Criterion) {
- let runtime = build_runtime();
- let mut group = c.benchmark_group("zip_object_archive_extract");
-
- for (name, compression_level, metadata_count, metadata_size, payload_count, payload_size) in [
- (
- "stored_metadata_heavy_384m_24p",
- CompressionLevel::Fastest,
- 384usize,
- 192usize,
- 24usize,
- 32 * 1024usize,
- ),
- (
- "deflated_mixed_192m_32p",
- CompressionLevel::Best,
- 192usize,
- 256usize,
- 32usize,
- 64 * 1024usize,
- ),
- ] {
- let files = build_object_archive_files(metadata_count, metadata_size, payload_count, payload_size);
- let total_bytes = files.iter().map(|(_, payload)| payload.len() as u64).sum::();
- let payload = {
- let temp = tempdir().expect("create benchmark tempdir");
- let zip_path = temp.path().join("object-archive.zip");
- runtime.block_on(async {
- create_zip_with_options(
- &zip_path,
- files.clone(),
- ZipWriteOptions {
- compression_level,
- create_directory_entries: true,
- },
- )
- .await
- .expect("create object archive benchmark payload");
- });
- std::fs::read(&zip_path).expect("read object archive benchmark payload")
- };
-
- group.throughput(Throughput::Bytes(total_bytes));
- group.bench_function(BenchmarkId::new("extract_full", name), |b| {
- b.iter(|| {
- let temp = tempdir().expect("create benchmark tempdir");
- let zip_path = temp.path().join("archive.zip");
- let extract_path = temp.path().join("extract");
- std::fs::write(&zip_path, &payload).expect("write object archive benchmark payload");
- runtime.block_on(async {
- let entries = extract_zip_with_limits(&zip_path, &extract_path, ArchiveLimits::default())
- .await
- .expect("extract object archive benchmark payload");
- black_box(entries.len());
- });
- });
- });
-
- group.bench_function(BenchmarkId::new("extract_summary_only", name), |b| {
- b.iter(|| {
- let temp = tempdir().expect("create benchmark tempdir");
- let zip_path = temp.path().join("archive.zip");
- let extract_path = temp.path().join("extract");
- std::fs::write(&zip_path, &payload).expect("write object archive benchmark payload");
- runtime.block_on(async {
- let summary = extract_zip_to_path_with_limits(&zip_path, &extract_path, ArchiveLimits::default())
- .await
- .expect("extract object archive benchmark summary");
- black_box(summary.file_count);
- });
- });
- });
- }
-
- group.finish();
-}
-
-criterion_group!(
- benches,
- bench_tar_family_extract,
- bench_zip_helper_round_trip,
- bench_zip_helper_hotspot_breakdown,
- bench_zip_object_archive_extract
-);
-criterion_main!(benches);
diff --git a/crates/zip/src/lib.rs b/crates/zip/src/lib.rs
index d61968fbb..82aa8867a 100644
--- a/crates/zip/src/lib.rs
+++ b/crates/zip/src/lib.rs
@@ -13,21 +13,8 @@
// limitations under the License.
use async_compression::tokio::bufread::{BzDecoder, GzipDecoder, XzDecoder, ZlibDecoder, ZstdDecoder};
-use async_compression::tokio::write::{BzEncoder, GzipEncoder, XzEncoder, ZlibEncoder, ZstdEncoder};
-use std::collections::HashSet;
-use std::future::Future;
-use std::io::{Read, Write};
-use std::path::{Component, Path, PathBuf};
-use std::pin::Pin;
-use std::sync::{Arc, Mutex};
-use std::task::{Context, Poll};
use thiserror::Error;
-use tokio::fs::File;
-use tokio::io::{self, AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt, BufReader, BufWriter};
-use tokio::task::spawn_blocking;
-use tokio_stream::StreamExt;
-use tokio_tar::Archive;
-use zip::{CompressionMethod, ZipArchive, ZipWriter, write::SimpleFileOptions};
+use tokio::io::{AsyncRead, BufReader};
pub type Result = std::result::Result;
@@ -38,51 +25,6 @@ pub enum ZipError {
format: CompressionFormat,
operation: &'static str,
},
- #[error("invalid compression level {0}: value exceeds i32::MAX")]
- InvalidCompressionLevel(u32),
- #[error("unsafe archive entry path: {0}")]
- UnsafeEntryPath(String),
- #[error("archive entry path length {length} exceeds limit {limit}: {path}")]
- EntryPathTooLong { path: String, length: usize, limit: usize },
- #[error("archive entry count {count} exceeds limit {limit}")]
- EntryCountLimitExceeded { count: usize, limit: usize },
- #[error("archive entry '{path}' size {size} exceeds limit {limit}")]
- EntrySizeLimitExceeded { path: String, size: u64, limit: u64 },
- #[error("archive total unpacked size {size} exceeds limit {limit}")]
- TotalUnpackedSizeLimitExceeded { size: u64, limit: u64 },
- #[error(transparent)]
- Io(#[from] io::Error),
- #[error(transparent)]
- Zip(#[from] zip::result::ZipError),
- #[error(transparent)]
- Join(#[from] tokio::task::JoinError),
-}
-
-#[derive(Debug, PartialEq, Eq, Clone, Copy)]
-pub enum CompressionCodec {
- Gzip,
- Bzip2,
- Xz,
- Zlib,
- Zstd,
-}
-
-#[derive(Debug, PartialEq, Eq, Clone, Copy)]
-pub enum ArchiveKind {
- Tar,
- Zip,
-}
-
-#[derive(Debug, PartialEq, Eq, Clone, Copy)]
-pub enum ArchiveFormat {
- Tar,
- TarGzip,
- TarBzip2,
- TarXz,
- TarZlib,
- TarZstd,
- Zip,
- Unknown,
}
#[derive(Debug, PartialEq, Eq, Clone, Copy)]
@@ -97,35 +39,9 @@ pub enum CompressionFormat {
Unknown,
}
-#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
-pub enum CompressionLevel {
- Fastest,
- Best,
- #[default]
- Default,
- Level(u32),
-}
-
-#[derive(Debug, Clone, PartialEq, Eq)]
-pub struct ZipEntry {
- pub name: String,
- pub size: u64,
- pub compressed_size: u64,
- pub is_dir: bool,
- pub compression_method: String,
- pub archive_kind: ArchiveKind,
- pub format: ArchiveFormat,
- pub unix_mode: Option,
-}
-
-#[derive(Debug, Clone, PartialEq, Eq, Default)]
-pub struct ZipExtractSummary {
- pub entry_count: usize,
- pub directory_count: usize,
- pub file_count: usize,
- pub total_unpacked_size: u64,
-}
-
+/// Archive guardrails. The values are carried here so every archive caller
+/// shares one default policy; enforcement belongs to the caller, which maps a
+/// breach onto its own protocol error.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct ArchiveLimits {
pub max_entries: usize,
@@ -147,103 +63,23 @@ impl Default for ArchiveLimits {
}
}
-#[derive(Debug, Clone, Copy, PartialEq, Eq)]
-pub struct ZipWriteOptions {
- pub compression_level: CompressionLevel,
- pub create_directory_entries: bool,
-}
-
-impl Default for ZipWriteOptions {
- fn default() -> Self {
- Self {
- compression_level: CompressionLevel::Default,
- create_directory_entries: false,
- }
- }
-}
-
-const SMALL_ZIP_EXTRACT_FAST_PATH_LIMIT: u64 = 8 * 1024;
-
-#[derive(Clone, Default)]
-struct SharedBuffer {
- inner: Arc>>,
-}
-
-impl SharedBuffer {
- fn into_vec(self) -> Vec {
- self.inner.lock().expect("shared in-memory writer lock poisoned").clone()
- }
-}
-
-impl AsyncWrite for SharedBuffer {
- fn poll_write(self: Pin<&mut Self>, _cx: &mut Context<'_>, buf: &[u8]) -> Poll> {
- let mut inner = self
- .inner
- .lock()
- .map_err(|_| io::Error::other("shared in-memory writer lock poisoned"))?;
- inner.extend_from_slice(buf);
- Poll::Ready(Ok(buf.len()))
- }
-
- fn poll_flush(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll> {
- Poll::Ready(Ok(()))
- }
-
- fn poll_shutdown(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll> {
- Poll::Ready(Ok(()))
- }
-}
-
impl CompressionFormat {
+ /// Map an archive extension onto the stream codec needed to read it.
+ /// Tar-family suffixes (`tgz`, `tbz2`, `txz`, `tzst`, ...) resolve to their
+ /// codec because the tar container itself is read from the decoded stream.
pub fn from_extension(ext: &str) -> Self {
- Self::from_archive_format(ArchiveFormat::from_extension(ext))
- }
-
- pub fn from_archive_format(format: ArchiveFormat) -> Self {
- match format {
- ArchiveFormat::TarGzip => CompressionFormat::Gzip,
- ArchiveFormat::TarBzip2 => CompressionFormat::Bzip2,
- ArchiveFormat::TarXz => CompressionFormat::Xz,
- ArchiveFormat::TarZlib => CompressionFormat::Zlib,
- ArchiveFormat::TarZstd => CompressionFormat::Zstd,
- ArchiveFormat::Tar => CompressionFormat::Tar,
- ArchiveFormat::Zip => CompressionFormat::Zip,
- ArchiveFormat::Unknown => CompressionFormat::Unknown,
+ match ext.to_ascii_lowercase().as_str() {
+ "gz" | "gzip" | "tgz" => CompressionFormat::Gzip,
+ "bz2" | "bzip2" | "tbz" | "tbz2" => CompressionFormat::Bzip2,
+ "xz" | "txz" => CompressionFormat::Xz,
+ "zlib" | "zz" => CompressionFormat::Zlib,
+ "zst" | "zstd" | "tzst" => CompressionFormat::Zstd,
+ "tar" => CompressionFormat::Tar,
+ "zip" => CompressionFormat::Zip,
+ _ => CompressionFormat::Unknown,
}
}
- pub fn archive_format_from_path>(path: P) -> ArchiveFormat {
- ArchiveFormat::from_path(path)
- }
-
- pub fn archive_kind(&self) -> Option {
- match self {
- CompressionFormat::Tar => Some(ArchiveKind::Tar),
- CompressionFormat::Zip => Some(ArchiveKind::Zip),
- CompressionFormat::Gzip
- | CompressionFormat::Bzip2
- | CompressionFormat::Xz
- | CompressionFormat::Zlib
- | CompressionFormat::Zstd
- | CompressionFormat::Unknown => None,
- }
- }
-
- pub fn compression_codec(&self) -> Option {
- match self {
- CompressionFormat::Gzip => Some(CompressionCodec::Gzip),
- CompressionFormat::Bzip2 => Some(CompressionCodec::Bzip2),
- CompressionFormat::Xz => Some(CompressionCodec::Xz),
- CompressionFormat::Zlib => Some(CompressionCodec::Zlib),
- CompressionFormat::Zstd => Some(CompressionCodec::Zstd),
- CompressionFormat::Tar | CompressionFormat::Zip | CompressionFormat::Unknown => None,
- }
- }
-
- pub fn from_path>(path: P) -> Self {
- Self::from_archive_format(ArchiveFormat::from_path(path))
- }
-
pub fn extension(&self) -> &'static str {
match self {
CompressionFormat::Gzip => "gz",
@@ -257,10 +93,6 @@ impl CompressionFormat {
}
}
- pub fn is_supported(&self) -> bool {
- !matches!(self, CompressionFormat::Unknown)
- }
-
pub fn get_decoder(&self, input: R) -> Result>
where
R: AsyncRead + Send + Unpin + 'static,
@@ -290,625 +122,14 @@ impl CompressionFormat {
Ok(decoder)
}
-
- fn convert_level(level: CompressionLevel) -> Result {
- match level {
- CompressionLevel::Fastest => Ok(async_compression::Level::Fastest),
- CompressionLevel::Best => Ok(async_compression::Level::Best),
- CompressionLevel::Default => Ok(async_compression::Level::Default),
- CompressionLevel::Level(n) => {
- let level = i32::try_from(n).map_err(|_| ZipError::InvalidCompressionLevel(n))?;
- Ok(async_compression::Level::Precise(level))
- }
- }
- }
-
- pub fn get_encoder(&self, output: W, level: CompressionLevel) -> Result>
- where
- W: AsyncWrite + Send + Unpin + 'static,
- {
- let writer = BufWriter::new(output);
-
- let encoder: Box = match self {
- CompressionFormat::Gzip => Box::new(GzipEncoder::with_quality(writer, Self::convert_level(level)?)),
- CompressionFormat::Bzip2 => Box::new(BzEncoder::with_quality(writer, Self::convert_level(level)?)),
- CompressionFormat::Zlib => Box::new(ZlibEncoder::with_quality(writer, Self::convert_level(level)?)),
- CompressionFormat::Xz => Box::new(XzEncoder::with_quality(writer, Self::convert_level(level)?)),
- CompressionFormat::Zstd => Box::new(ZstdEncoder::with_quality(writer, Self::convert_level(level)?)),
- CompressionFormat::Tar => Box::new(writer),
- CompressionFormat::Zip => {
- return Err(ZipError::UnsupportedFormat {
- format: *self,
- operation: "stream encoding",
- });
- }
- CompressionFormat::Unknown => {
- return Err(ZipError::UnsupportedFormat {
- format: *self,
- operation: "encoding",
- });
- }
- };
-
- Ok(encoder)
- }
-}
-
-impl ArchiveFormat {
- pub fn from_extension(ext: &str) -> Self {
- match ext.to_ascii_lowercase().as_str() {
- "gz" | "gzip" | "tgz" => ArchiveFormat::TarGzip,
- "bz2" | "bzip2" | "tbz" | "tbz2" => ArchiveFormat::TarBzip2,
- "xz" | "txz" => ArchiveFormat::TarXz,
- "zlib" | "zz" => ArchiveFormat::TarZlib,
- "zst" | "zstd" | "tzst" => ArchiveFormat::TarZstd,
- "tar" => ArchiveFormat::Tar,
- "zip" => ArchiveFormat::Zip,
- _ => ArchiveFormat::Unknown,
- }
- }
-
- pub fn from_path>(path: P) -> Self {
- let path = path.as_ref();
- let lower_name = path.file_name().and_then(|name| name.to_str()).map(str::to_ascii_lowercase);
-
- if let Some(name) = lower_name {
- if name.ends_with(".tar.gz") || name.ends_with(".tgz") {
- return ArchiveFormat::TarGzip;
- }
- if name.ends_with(".tar.bz2") || name.ends_with(".tbz") || name.ends_with(".tbz2") {
- return ArchiveFormat::TarBzip2;
- }
- if name.ends_with(".tar.xz") || name.ends_with(".txz") {
- return ArchiveFormat::TarXz;
- }
- if name.ends_with(".tar.zst") || name.ends_with(".tzst") {
- return ArchiveFormat::TarZstd;
- }
- if name.ends_with(".tar.zlib") {
- return ArchiveFormat::TarZlib;
- }
- }
-
- path.extension()
- .and_then(|s| s.to_str())
- .map(Self::from_extension)
- .unwrap_or(ArchiveFormat::Unknown)
- }
-
- pub fn archive_kind(&self) -> Option {
- match self {
- ArchiveFormat::Tar
- | ArchiveFormat::TarGzip
- | ArchiveFormat::TarBzip2
- | ArchiveFormat::TarXz
- | ArchiveFormat::TarZlib
- | ArchiveFormat::TarZstd => Some(ArchiveKind::Tar),
- ArchiveFormat::Zip => Some(ArchiveKind::Zip),
- ArchiveFormat::Unknown => None,
- }
- }
-
- pub fn compression_codec(&self) -> Option {
- match self {
- ArchiveFormat::TarGzip => Some(CompressionCodec::Gzip),
- ArchiveFormat::TarBzip2 => Some(CompressionCodec::Bzip2),
- ArchiveFormat::TarXz => Some(CompressionCodec::Xz),
- ArchiveFormat::TarZlib => Some(CompressionCodec::Zlib),
- ArchiveFormat::TarZstd => Some(CompressionCodec::Zstd),
- ArchiveFormat::Tar | ArchiveFormat::Zip | ArchiveFormat::Unknown => None,
- }
- }
-
- pub fn extension(&self) -> &'static str {
- match self {
- ArchiveFormat::Tar => "tar",
- ArchiveFormat::TarGzip => "tar.gz",
- ArchiveFormat::TarBzip2 => "tar.bz2",
- ArchiveFormat::TarXz => "tar.xz",
- ArchiveFormat::TarZlib => "tar.zlib",
- ArchiveFormat::TarZstd => "tar.zst",
- ArchiveFormat::Zip => "zip",
- ArchiveFormat::Unknown => "",
- }
- }
-}
-
-/// Read entries from a tar-family archive stream.
-///
-/// Supported formats are:
-/// - `CompressionFormat::Tar`
-/// - `CompressionFormat::Gzip`
-/// - `CompressionFormat::Bzip2`
-/// - `CompressionFormat::Xz`
-/// - `CompressionFormat::Zlib`
-/// - `CompressionFormat::Zstd`
-///
-/// `CompressionFormat::Zip` is intentionally not supported here because ZIP
-/// requires central-directory semantics and is handled through file-based
-/// helper APIs.
-pub async fn read_archive_entries(input: R, format: CompressionFormat, callback: F) -> Result<()>
-where
- R: AsyncRead + Send + Unpin + 'static,
- F: FnMut(tokio_tar::Entry>>) -> Fut + Send + 'static,
- Fut: Future