diff --git a/Cargo.lock b/Cargo.lock index 14bf7770..e31779f6 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -99,6 +99,7 @@ checksum = "d10e4f991a553474232bc0a31799f6d24b034a84c0971d80d2e2f78b2e576e40" dependencies = [ "compression-codecs", "compression-core", + "futures-io", "pin-project-lite", "tokio", ] @@ -147,6 +148,21 @@ dependencies = [ "syn", ] +[[package]] +name = "async_zip" +version = "0.0.18" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0d8c50d65ce1b0e0cb65a785ff615f78860d7754290647d3b983208daa4f85e6" +dependencies = [ + "async-compression", + "crc32fast", + "futures-lite", + "pin-project", + "thiserror", + "tokio", + "tokio-util", +] + [[package]] name = "atoi" version = "2.0.0" @@ -920,6 +936,19 @@ version = "0.3.31" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9e5c1b78ca4aae1ac06c48a526a655760685149f0d465d21f37abfe57ce075c6" +[[package]] +name = "futures-lite" +version = "2.6.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f78e10609fe0e0b3f4157ffab1876319b5b0db102a2c60dc4626306dc46b44ad" +dependencies = [ + "fastrand", + "futures-core", + "futures-io", + "parking", + "pin-project-lite", +] + [[package]] name = "futures-macro" version = "0.3.31" @@ -1793,6 +1822,7 @@ dependencies = [ "argon2", "async-stream", "async-trait", + "async_zip", "axum", "base64", "bytes", @@ -1933,6 +1963,26 @@ version = "2.3.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9b4f627cb1b25917193a259e49bdad08f671f8d9708acfd5fe0a8c1455d87220" +[[package]] +name = "pin-project" +version = "1.1.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "677f1add503faace112b9f1373e43e9e054bfdd22ff1a63c1bc485eaec6a6a8a" +dependencies = [ + "pin-project-internal", +] + +[[package]] +name = "pin-project-internal" +version = "1.1.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6e918e4ff8c4549eb882f14b3a4bc8c8bc93de829416eacf579f1207a8fbf861" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + [[package]] name = "pin-project-lite" version = "0.2.16" @@ -3060,6 +3110,7 @@ checksum = "9ae9cec805b01e8fc3fd2fe289f89149a9b66dd16786abd8b19cfa7b48cb0098" dependencies = [ "bytes", "futures-core", + "futures-io", "futures-sink", "pin-project-lite", "tokio", diff --git a/Cargo.toml b/Cargo.toml index 8d6f0819..10978baa 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -8,7 +8,7 @@ default-run = "oxicloud" [dependencies] axum = { version = "0.8.8", features = ["multipart", "http1", "tokio", "macros"] } tokio = { version = "1.49.0", features = ["full"] } -tokio-util = { version = "0.7.18", features = ["io", "codec"] } +tokio-util = { version = "0.7.18", features = ["io", "codec", "compat"] } tokio-stream = { version = "0.1.18", features = ["fs"] } bytes = "1.11.1" tempfile = "3.25.0" @@ -49,6 +49,7 @@ reqwest = { version = "0.12", default-features = false, features = ["json", "rus base64 = "0.22.1" fs2 = "0.4" rayon = "1.10" +async_zip = { version = "0.0.18", features = ["tokio", "deflate"] } [features] default = [] diff --git a/src/infrastructure/services/zip_service.rs b/src/infrastructure/services/zip_service.rs index b16d9889..9ce6f45d 100644 --- a/src/infrastructure/services/zip_service.rs +++ b/src/infrastructure/services/zip_service.rs @@ -7,13 +7,16 @@ use crate::{ common::errors::{DomainError, ErrorKind, Result}, }; use async_trait::async_trait; +use async_zip::base::write::ZipFileWriter; +use async_zip::{Compression, ZipEntryBuilder}; +use futures::io::AsyncWriteExt as FuturesWriteExt; use futures::StreamExt; -use std::io::Write; use std::sync::Arc; use tempfile::NamedTempFile; use thiserror::Error; +use tokio::io::BufWriter; +use tokio_util::compat::Compat; use tracing::*; -use zip::{ZipWriter, write::SimpleFileOptions}; /// Error related to ZIP file creation #[derive(Debug, Error)] @@ -22,7 +25,7 @@ pub enum ZipError { IoError(#[from] std::io::Error), #[error("ZIP error: {0}")] - ZipError(#[from] zip::result::ZipError), + AsyncZipError(#[from] async_zip::error::ZipError), #[error("Error reading file: {0}")] FileReadError(String), @@ -34,24 +37,21 @@ pub enum ZipError { FolderNotFound(String), } -// Implement From for DomainError to allow the use of ? impl From for DomainError { fn from(err: ZipError) -> Self { DomainError::new(ErrorKind::InternalError, "zip_service", err.to_string()) } } -// Implement From for DomainError directly -impl From for DomainError { - fn from(err: zip::result::ZipError) -> Self { - DomainError::new(ErrorKind::InternalError, "zip_service", err.to_string()) - } -} +/// Type alias for the fully-async ZIP writer backed by a buffered tokio file. +type AsyncZipWriter = ZipFileWriter>>; /// Service for creating ZIP files. /// -/// Writes the ZIP archive to a temporary file on disk so that only one file's -/// stream-chunk (~64 KB) is held in memory at a time, regardless of archive size. +/// Uses `async_zip` for fully-async archive creation. Every write (headers, +/// compressed chunk data, central directory) goes through +/// `tokio::io::BufWriter` → `tokio::fs::File`, so **no Tokio worker is ever +/// blocked** by disk I/O or compression. pub struct ZipService { file_service: Arc, folder_service: Arc, @@ -70,7 +70,7 @@ impl ZipService { } /// Creates a ZIP file backed by a temporary file, containing the contents - /// of a folder and all its subfolders. Returns the `NamedTempFile` so the + /// of a folder and all its subfolders. Returns the `NamedTempFile` so the /// caller can stream it and let the OS clean up on drop. pub async fn create_folder_zip( &self, @@ -91,15 +91,15 @@ impl ZipService { } }; - // Create a temp file to back the ZIP archive (O(1) RAM) + // Create a temp file; open a second async handle for writing. let temp = NamedTempFile::new().map_err(ZipError::IoError)?; - let raw_file = temp.reopen().map_err(ZipError::IoError)?; - let mut zip = ZipWriter::new(raw_file); + let tokio_file = tokio::fs::File::create(temp.path()) + .await + .map_err(ZipError::IoError)?; - // Set compression options - let options = SimpleFileOptions::default() - .compression_method(zip::CompressionMethod::Deflated) - .unix_permissions(0o755); + // 256 KB buffer keeps syscall count low. + let buf_writer = BufWriter::with_capacity(256 * 1024, tokio_file); + let mut zip = ZipFileWriter::with_tokio(buf_writer); // Track processed folders to avoid cycles let mut processed_folders = std::collections::HashSet::new(); @@ -109,25 +109,24 @@ impl ZipService { &mut zip, &folder, folder_name, - &options, &mut processed_folders, ) .await?; - // Finalize the ZIP (flushes central directory) - zip.finish()?; + // Finalize: writes central directory, then flush buffered data to disk. + let mut compat_writer = zip.close().await.map_err(ZipError::AsyncZipError)?; + compat_writer.close().await.map_err(ZipError::IoError)?; Ok(temp) } - /// Iterative BFS over the folder tree. Writes entries directly to the - /// file-backed `ZipWriter` so memory stays flat. + /// Iterative BFS over the folder tree. Writes entries directly to the + /// async `ZipFileWriter` so memory stays flat. async fn process_folder_recursively( &self, - zip: &mut ZipWriter, + zip: &mut AsyncZipWriter, folder: &FolderDto, path: &str, - options: &SimpleFileOptions, processed_folders: &mut std::collections::HashSet, ) -> Result<()> { struct PendingFolder { @@ -148,10 +147,12 @@ impl ZipService { } processed_folders.insert(folder_id.clone()); - // Directory entry + // Directory entry (Stored, zero-length body) let folder_path = format!("{}/", current.path); - match zip.add_directory(&folder_path, *options) { - Ok(_) => debug!("Folder added to ZIP: {}", folder_path), + let dir_entry = + ZipEntryBuilder::new(folder_path.clone().into(), Compression::Stored); + match zip.write_entry_whole(dir_entry, &[]).await { + Ok(()) => debug!("Folder added to ZIP: {}", folder_path), Err(e) => { warn!("Could not add folder to ZIP (may already exist): {}", e); } @@ -171,7 +172,7 @@ impl ZipService { }; for file in files { - self.add_file_to_zip_streamed(zip, &file, &folder_path, options) + self.add_file_to_zip_streamed(zip, &file, &folder_path) .await?; } @@ -200,29 +201,33 @@ impl ZipService { Ok(()) } - /// Streams file content in chunks (~64 KB) into the ZIP entry, keeping - /// peak memory independent of individual file sizes. + /// Streams file content in chunks (~64 KB) into an async ZIP entry, + /// keeping peak memory independent of individual file sizes. async fn add_file_to_zip_streamed( &self, - zip: &mut ZipWriter, + zip: &mut AsyncZipWriter, file: &FileDto, folder_path: &str, - options: &SimpleFileOptions, ) -> Result<()> { let file_path = format!("{}{}", folder_path, file.name); info!("Adding file to ZIP: {}", file_path); let file_id = file.id.to_string(); - // Start the ZIP entry - zip.start_file_from_path(std::path::Path::new(&file_path), *options) - .map_err(ZipError::ZipError)?; + // Open a streaming entry with Deflate compression + let entry = ZipEntryBuilder::new(file_path.clone().into(), Compression::Deflate); + let mut entry_writer = zip + .write_entry_stream(entry) + .await + .map_err(ZipError::AsyncZipError)?; // Stream file contents in chunks instead of loading all into RAM let stream = match self.file_service.get_file_stream(&file_id).await { Ok(s) => s, Err(e) => { error!("Error opening file stream {}: {}", file_id, e); + // Close the partially-opened entry before returning + let _ = entry_writer.close().await; return Err(ZipError::FileReadError(format!( "Error streaming file {}: {}", file_id, e @@ -236,9 +241,18 @@ impl ZipService { while let Some(chunk_result) = stream.next().await { let bytes = chunk_result.map_err(ZipError::IoError)?; - zip.write_all(&bytes).map_err(ZipError::IoError)?; + entry_writer + .write_all(&bytes) + .await + .map_err(ZipError::IoError)?; } + // Finalize the entry (writes data descriptor with CRC + sizes) + entry_writer + .close() + .await + .map_err(ZipError::AsyncZipError)?; + debug!("File added to ZIP: {}", file_path); Ok(()) }