2025-04-02 01:22:05 +02:00
|
|
|
use crate::{
|
|
|
|
|
application::dtos::file_dto::FileDto,
|
|
|
|
|
application::dtos::folder_dto::FolderDto,
|
2026-02-08 13:40:23 +01:00
|
|
|
application::ports::file_ports::FileRetrievalUseCase,
|
2026-02-14 01:29:34 +01:00
|
|
|
application::ports::inbound::FolderUseCase,
|
2026-02-08 13:40:23 +01:00
|
|
|
application::ports::zip_ports::ZipPort,
|
2026-02-14 01:29:34 +01:00
|
|
|
common::errors::{DomainError, ErrorKind, Result},
|
2025-04-02 01:22:05 +02:00
|
|
|
};
|
2026-02-14 01:29:34 +01:00
|
|
|
use async_trait::async_trait;
|
2026-02-22 22:29:07 +01:00
|
|
|
use futures::StreamExt;
|
|
|
|
|
use std::io::Write;
|
2025-04-02 01:22:05 +02:00
|
|
|
use std::sync::Arc;
|
2026-02-22 22:29:07 +01:00
|
|
|
use tempfile::NamedTempFile;
|
2026-02-14 01:29:34 +01:00
|
|
|
use thiserror::Error;
|
|
|
|
|
use tracing::*;
|
|
|
|
|
use zip::{ZipWriter, write::SimpleFileOptions};
|
2025-04-02 01:22:05 +02:00
|
|
|
|
2026-02-12 09:41:25 +01:00
|
|
|
/// Error related to ZIP file creation
|
2025-04-02 01:22:05 +02:00
|
|
|
#[derive(Debug, Error)]
|
|
|
|
|
pub enum ZipError {
|
2026-02-12 09:41:25 +01:00
|
|
|
#[error("IO error: {0}")]
|
2025-04-02 01:22:05 +02:00
|
|
|
IoError(#[from] std::io::Error),
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-12 09:41:25 +01:00
|
|
|
#[error("ZIP error: {0}")]
|
2025-04-02 01:22:05 +02:00
|
|
|
ZipError(#[from] zip::result::ZipError),
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-12 09:41:25 +01:00
|
|
|
#[error("Error reading file: {0}")]
|
2025-04-02 01:22:05 +02:00
|
|
|
FileReadError(String),
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-12 09:41:25 +01:00
|
|
|
#[error("Error getting folder contents: {0}")]
|
2025-04-02 01:22:05 +02:00
|
|
|
FolderContentsError(String),
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-12 09:41:25 +01:00
|
|
|
#[error("Folder not found: {0}")]
|
2025-04-02 01:22:05 +02:00
|
|
|
FolderNotFound(String),
|
|
|
|
|
}
|
|
|
|
|
|
2026-02-12 09:41:25 +01:00
|
|
|
// Implement From<ZipError> for DomainError to allow the use of ?
|
2025-04-02 01:22:05 +02:00
|
|
|
impl From<ZipError> for DomainError {
|
|
|
|
|
fn from(err: ZipError) -> Self {
|
|
|
|
|
DomainError::new(ErrorKind::InternalError, "zip_service", err.to_string())
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-02-12 09:41:25 +01:00
|
|
|
// Implement From<zip::result::ZipError> for DomainError directly
|
2025-04-02 01:22:05 +02:00
|
|
|
impl From<zip::result::ZipError> for DomainError {
|
|
|
|
|
fn from(err: zip::result::ZipError) -> Self {
|
|
|
|
|
DomainError::new(ErrorKind::InternalError, "zip_service", err.to_string())
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-02-22 22:29:07 +01:00
|
|
|
/// 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.
|
2025-04-02 01:22:05 +02:00
|
|
|
pub struct ZipService {
|
2026-02-08 13:40:23 +01:00
|
|
|
file_service: Arc<dyn FileRetrievalUseCase>,
|
2025-04-02 01:22:05 +02:00
|
|
|
folder_service: Arc<dyn FolderUseCase>,
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
impl ZipService {
|
2026-02-22 22:29:07 +01:00
|
|
|
/// Creates a new instance of the ZIP service
|
2026-02-14 01:29:34 +01:00
|
|
|
pub fn new(
|
|
|
|
|
file_service: Arc<dyn FileRetrievalUseCase>,
|
|
|
|
|
folder_service: Arc<dyn FolderUseCase>,
|
|
|
|
|
) -> Self {
|
2025-04-02 01:22:05 +02:00
|
|
|
Self {
|
|
|
|
|
file_service,
|
|
|
|
|
folder_service,
|
|
|
|
|
}
|
|
|
|
|
}
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-22 22:29:07 +01:00
|
|
|
/// Creates a ZIP file backed by a temporary file, containing the contents
|
|
|
|
|
/// 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,
|
|
|
|
|
folder_id: &str,
|
|
|
|
|
folder_name: &str,
|
|
|
|
|
) -> Result<NamedTempFile> {
|
2026-02-14 01:29:34 +01:00
|
|
|
info!(
|
|
|
|
|
"Creating ZIP for folder: {} (ID: {})",
|
|
|
|
|
folder_name, folder_id
|
|
|
|
|
);
|
|
|
|
|
|
2026-02-22 22:29:07 +01:00
|
|
|
// Verify the folder exists
|
2025-04-02 01:22:05 +02:00
|
|
|
let folder = match self.folder_service.get_folder(folder_id).await {
|
|
|
|
|
Ok(folder) => folder,
|
|
|
|
|
Err(e) => {
|
2026-02-12 09:41:25 +01:00
|
|
|
error!("Error getting folder {}: {}", folder_id, e);
|
2025-04-02 01:22:05 +02:00
|
|
|
return Err(ZipError::FolderNotFound(folder_id.to_string()).into());
|
|
|
|
|
}
|
|
|
|
|
};
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-22 22:29:07 +01:00
|
|
|
// Create a temp file to back the ZIP archive (O(1) RAM)
|
|
|
|
|
let temp = NamedTempFile::new().map_err(ZipError::IoError)?;
|
|
|
|
|
let raw_file = temp.reopen().map_err(ZipError::IoError)?;
|
|
|
|
|
let mut zip = ZipWriter::new(raw_file);
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-12 09:41:25 +01:00
|
|
|
// Set compression options
|
2025-04-02 01:22:05 +02:00
|
|
|
let options = SimpleFileOptions::default()
|
|
|
|
|
.compression_method(zip::CompressionMethod::Deflated)
|
|
|
|
|
.unix_permissions(0o755);
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-22 22:29:07 +01:00
|
|
|
// Track processed folders to avoid cycles
|
2025-04-02 01:22:05 +02:00
|
|
|
let mut processed_folders = std::collections::HashSet::new();
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-22 22:29:07 +01:00
|
|
|
// Build the ZIP iteratively
|
2025-04-02 01:22:05 +02:00
|
|
|
self.process_folder_recursively(
|
|
|
|
|
&mut zip,
|
|
|
|
|
&folder,
|
|
|
|
|
folder_name,
|
|
|
|
|
&options,
|
2026-02-14 01:29:34 +01:00
|
|
|
&mut processed_folders,
|
|
|
|
|
)
|
|
|
|
|
.await?;
|
|
|
|
|
|
2026-02-22 22:29:07 +01:00
|
|
|
// Finalize the ZIP (flushes central directory)
|
|
|
|
|
zip.finish()?;
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-22 22:29:07 +01:00
|
|
|
Ok(temp)
|
2025-04-02 01:22:05 +02:00
|
|
|
}
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-22 22:29:07 +01:00
|
|
|
/// Iterative BFS over the folder tree. Writes entries directly to the
|
|
|
|
|
/// file-backed `ZipWriter` so memory stays flat.
|
2025-04-02 01:22:05 +02:00
|
|
|
async fn process_folder_recursively(
|
|
|
|
|
&self,
|
2026-02-22 22:29:07 +01:00
|
|
|
zip: &mut ZipWriter<std::fs::File>,
|
2025-04-02 01:22:05 +02:00
|
|
|
folder: &FolderDto,
|
|
|
|
|
path: &str,
|
|
|
|
|
options: &SimpleFileOptions,
|
2026-02-14 01:29:34 +01:00
|
|
|
processed_folders: &mut std::collections::HashSet<String>,
|
2025-04-02 01:22:05 +02:00
|
|
|
) -> Result<()> {
|
|
|
|
|
struct PendingFolder {
|
|
|
|
|
folder: FolderDto,
|
|
|
|
|
path: String,
|
|
|
|
|
}
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2025-04-02 01:22:05 +02:00
|
|
|
let mut work_queue = vec![PendingFolder {
|
|
|
|
|
folder: folder.clone(),
|
|
|
|
|
path: path.to_string(),
|
|
|
|
|
}];
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2025-04-02 01:22:05 +02:00
|
|
|
while let Some(current) = work_queue.pop() {
|
|
|
|
|
let folder_id = current.folder.id.to_string();
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2025-04-02 01:22:05 +02:00
|
|
|
if processed_folders.contains(&folder_id) {
|
|
|
|
|
continue;
|
|
|
|
|
}
|
|
|
|
|
processed_folders.insert(folder_id.clone());
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-22 22:29:07 +01:00
|
|
|
// Directory entry
|
2025-04-02 01:22:05 +02:00
|
|
|
let folder_path = format!("{}/", current.path);
|
|
|
|
|
match zip.add_directory(&folder_path, *options) {
|
2026-02-12 09:41:25 +01:00
|
|
|
Ok(_) => debug!("Folder added to ZIP: {}", folder_path),
|
2025-04-02 01:22:05 +02:00
|
|
|
Err(e) => {
|
2026-02-22 22:29:07 +01:00
|
|
|
warn!("Could not add folder to ZIP (may already exist): {}", e);
|
2025-04-02 01:22:05 +02:00
|
|
|
}
|
|
|
|
|
}
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-22 22:29:07 +01:00
|
|
|
// Files in this folder
|
2025-04-02 01:22:05 +02:00
|
|
|
let files = match self.file_service.list_files(Some(&folder_id)).await {
|
|
|
|
|
Ok(files) => files,
|
|
|
|
|
Err(e) => {
|
2026-02-12 09:41:25 +01:00
|
|
|
error!("Error listing files in folder {}: {}", folder_id, e);
|
2026-02-14 01:29:34 +01:00
|
|
|
return Err(ZipError::FolderContentsError(format!(
|
|
|
|
|
"Error listing files: {}",
|
|
|
|
|
e
|
|
|
|
|
))
|
|
|
|
|
.into());
|
2025-04-02 01:22:05 +02:00
|
|
|
}
|
|
|
|
|
};
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2025-04-02 01:22:05 +02:00
|
|
|
for file in files {
|
2026-02-22 22:29:07 +01:00
|
|
|
self.add_file_to_zip_streamed(zip, &file, &folder_path, options)
|
2026-02-14 01:29:34 +01:00
|
|
|
.await?;
|
2025-04-02 01:22:05 +02:00
|
|
|
}
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-22 22:29:07 +01:00
|
|
|
// Subfolders
|
2025-04-02 01:22:05 +02:00
|
|
|
let subfolders = match self.folder_service.list_folders(Some(&folder_id)).await {
|
|
|
|
|
Ok(folders) => folders,
|
|
|
|
|
Err(e) => {
|
2026-02-12 09:41:25 +01:00
|
|
|
error!("Error listing subfolders in {}: {}", folder_id, e);
|
2026-02-14 01:29:34 +01:00
|
|
|
return Err(ZipError::FolderContentsError(format!(
|
|
|
|
|
"Error listing subfolders: {}",
|
|
|
|
|
e
|
|
|
|
|
))
|
|
|
|
|
.into());
|
2025-04-02 01:22:05 +02:00
|
|
|
}
|
|
|
|
|
};
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2025-04-02 01:22:05 +02:00
|
|
|
for subfolder in subfolders {
|
|
|
|
|
let subfolder_path = format!("{}/{}", current.path, subfolder.name);
|
|
|
|
|
work_queue.push(PendingFolder {
|
|
|
|
|
folder: subfolder,
|
|
|
|
|
path: subfolder_path,
|
|
|
|
|
});
|
|
|
|
|
}
|
|
|
|
|
}
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2025-04-02 01:22:05 +02:00
|
|
|
Ok(())
|
|
|
|
|
}
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-22 22:29:07 +01:00
|
|
|
/// Streams file content in chunks (~64 KB) into the ZIP entry, keeping
|
|
|
|
|
/// peak memory independent of individual file sizes.
|
|
|
|
|
async fn add_file_to_zip_streamed(
|
2025-04-02 01:22:05 +02:00
|
|
|
&self,
|
2026-02-22 22:29:07 +01:00
|
|
|
zip: &mut ZipWriter<std::fs::File>,
|
2025-04-02 01:22:05 +02:00
|
|
|
file: &FileDto,
|
|
|
|
|
folder_path: &str,
|
|
|
|
|
options: &SimpleFileOptions,
|
|
|
|
|
) -> Result<()> {
|
|
|
|
|
let file_path = format!("{}{}", folder_path, file.name);
|
2026-02-12 09:41:25 +01:00
|
|
|
info!("Adding file to ZIP: {}", file_path);
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2025-04-02 01:22:05 +02:00
|
|
|
let file_id = file.id.to_string();
|
2026-02-22 22:29:07 +01:00
|
|
|
|
|
|
|
|
// Start the ZIP entry
|
|
|
|
|
zip.start_file_from_path(std::path::Path::new(&file_path), *options)
|
|
|
|
|
.map_err(ZipError::ZipError)?;
|
|
|
|
|
|
|
|
|
|
// 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,
|
2025-04-02 01:22:05 +02:00
|
|
|
Err(e) => {
|
2026-02-22 22:29:07 +01:00
|
|
|
error!("Error opening file stream {}: {}", file_id, e);
|
2026-02-14 01:29:34 +01:00
|
|
|
return Err(ZipError::FileReadError(format!(
|
2026-02-22 22:29:07 +01:00
|
|
|
"Error streaming file {}: {}",
|
2026-02-14 01:29:34 +01:00
|
|
|
file_id, e
|
|
|
|
|
))
|
|
|
|
|
.into());
|
2025-04-02 01:22:05 +02:00
|
|
|
}
|
|
|
|
|
};
|
2026-02-14 01:29:34 +01:00
|
|
|
|
2026-02-22 22:29:07 +01:00
|
|
|
// Pin the stream so StreamExt::next() can be called
|
|
|
|
|
let mut stream = std::pin::Pin::from(stream);
|
|
|
|
|
|
|
|
|
|
while let Some(chunk_result) = stream.next().await {
|
|
|
|
|
let bytes = chunk_result.map_err(ZipError::IoError)?;
|
|
|
|
|
zip.write_all(&bytes).map_err(ZipError::IoError)?;
|
2025-04-02 01:22:05 +02:00
|
|
|
}
|
2026-02-22 22:29:07 +01:00
|
|
|
|
|
|
|
|
debug!("File added to ZIP: {}", file_path);
|
|
|
|
|
Ok(())
|
2025-04-02 01:22:05 +02:00
|
|
|
}
|
2026-02-08 13:40:23 +01:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// ─── Port implementation ─────────────────────────────────────────────────────
|
|
|
|
|
|
|
|
|
|
#[async_trait]
|
|
|
|
|
impl ZipPort for ZipService {
|
|
|
|
|
async fn create_folder_zip(
|
|
|
|
|
&self,
|
|
|
|
|
folder_id: &str,
|
|
|
|
|
folder_name: &str,
|
2026-02-22 22:29:07 +01:00
|
|
|
) -> std::result::Result<NamedTempFile, DomainError> {
|
2026-02-08 13:40:23 +01:00
|
|
|
self.create_folder_zip(folder_id, folder_name).await
|
|
|
|
|
}
|
2026-02-14 01:29:34 +01:00
|
|
|
}
|