diff --git a/Cargo.lock b/Cargo.lock index 9b32a950..79b650f0 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1756,7 +1756,7 @@ checksum = "42f5e15c9953c5e4ccceeb2e7382a716482c34515315f7b03532b8b4e8393d2d" [[package]] name = "oxicloud" -version = "0.4.2" +version = "0.5.0" dependencies = [ "anyhow", "argon2", diff --git a/src/application/services/batch_operations.rs b/src/application/services/batch_operations.rs index ebbee241..590e1087 100644 --- a/src/application/services/batch_operations.rs +++ b/src/application/services/batch_operations.rs @@ -733,18 +733,16 @@ impl BatchOperationService { } // ── Finalize ───────────────────────────────────────────────────── - let mut compat_writer = zip.close().await.map_err(|e| { - BatchOperationError::Internal(format!("ZIP finalize error: {}", e)) - })?; - compat_writer.close().await.map_err(|e| { - BatchOperationError::Internal(format!("ZIP flush error: {}", e)) - })?; + let mut compat_writer = zip + .close() + .await + .map_err(|e| BatchOperationError::Internal(format!("ZIP finalize error: {}", e)))?; + compat_writer + .close() + .await + .map_err(|e| BatchOperationError::Internal(format!("ZIP flush error: {}", e)))?; - let file_size = temp - .as_file() - .metadata() - .map(|m| m.len()) - .unwrap_or(0); + let file_size = temp.as_file().metadata().map(|m| m.len()).unwrap_or(0); info!( "Batch download ZIP created: {} bytes in {}ms", @@ -776,17 +774,18 @@ impl BatchOperationService { let mut stream = std::pin::Pin::from(stream); while let Some(chunk) = stream.next().await { - let bytes = chunk.map_err(|e| { - BatchOperationError::Internal(format!("stream read: {}", e)) - })?; - writer.write_all(&bytes).await.map_err(|e| { - BatchOperationError::Internal(format!("zip chunk write: {}", e)) - })?; + let bytes = + chunk.map_err(|e| BatchOperationError::Internal(format!("stream read: {}", e)))?; + writer + .write_all(&bytes) + .await + .map_err(|e| BatchOperationError::Internal(format!("zip chunk write: {}", e)))?; } - writer.close().await.map_err(|e| { - BatchOperationError::Internal(format!("zip entry close: {}", e)) - })?; + writer + .close() + .await + .map_err(|e| BatchOperationError::Internal(format!("zip entry close: {}", e)))?; Ok(()) } diff --git a/src/infrastructure/services/chunked_upload_service.rs b/src/infrastructure/services/chunked_upload_service.rs index 44e00cbe..dd215802 100644 --- a/src/infrastructure/services/chunked_upload_service.rs +++ b/src/infrastructure/services/chunked_upload_service.rs @@ -157,8 +157,8 @@ impl UploadSession { /// Persist the full session metadata once (on create). async fn persist_metadata(&self) -> Result<(), String> { let path = self.temp_dir.join(SESSION_META_FILE); - let json = serde_json::to_vec(self) - .map_err(|e| format!("Failed to serialise session: {e}"))?; + let json = + serde_json::to_vec(self).map_err(|e| format!("Failed to serialise session: {e}"))?; // Atomic write: write to .tmp then rename let tmp = self.temp_dir.join("session.json.tmp"); fs::write(&tmp, &json) @@ -213,9 +213,7 @@ impl ChunkedUploadService { }; if recovered_count > 0 { - tracing::info!( - "♻️ Recovered {recovered_count} chunked-upload session(s) from disk" - ); + tracing::info!("♻️ Recovered {recovered_count} chunked-upload session(s) from disk"); } // Start cleanup task @@ -325,10 +323,7 @@ impl ChunkedUploadService { // ── Cleanup ────────────────────────────────────────────────────────── /// Background task to clean expired sessions - async fn cleanup_loop( - sessions: Arc>, - temp_base_dir: PathBuf, - ) { + async fn cleanup_loop(sessions: Arc>, temp_base_dir: PathBuf) { let mut interval = tokio::time::interval(Duration::from_secs(3600)); // Every hour loop { @@ -463,7 +458,8 @@ impl ChunkedUploadService { ) -> Result { // Validate session exists and chunk index is valid let (chunk_path, expected_size) = { - let session = self.sessions + let session = self + .sessions .get(upload_id) .ok_or_else(|| format!("Upload session not found: {}", upload_id))?; @@ -500,11 +496,10 @@ impl ChunkedUploadService { // worker free for other connections. if let Some(ref expected_checksum) = checksum { let data_clone = data.clone(); // Bytes::clone is O(1) — just an Arc increment - let actual_checksum = tokio::task::spawn_blocking(move || { - format!("{:x}", md5::compute(&data_clone)) - }) - .await - .map_err(|e| format!("MD5 checksum task failed: {e}"))?; + let actual_checksum = + tokio::task::spawn_blocking(move || format!("{:x}", md5::compute(&data_clone))) + .await + .map_err(|e| format!("MD5 checksum task failed: {e}"))?; if actual_checksum != *expected_checksum { return Err(format!( @@ -527,7 +522,8 @@ impl ChunkedUploadService { // Disk I/O (persist_progress) is done AFTER the ref is dropped so // concurrent uploads to other sessions are never blocked. let (bytes_received, progress, is_complete, persist_path, persist_bitmask) = { - let mut session = self.sessions + let mut session = self + .sessions .get_mut(upload_id) .ok_or_else(|| "Session disappeared".to_string())?; @@ -571,11 +567,9 @@ impl ChunkedUploadService { } /// Get upload status - async fn get_status_inner( - &self, - upload_id: &str, - ) -> Result { - let session = self.sessions + async fn get_status_inner(&self, upload_id: &str) -> Result { + let session = self + .sessions .get(upload_id) .ok_or_else(|| format!("Upload session not found: {}", upload_id))?; @@ -613,7 +607,8 @@ impl ChunkedUploadService { // Clone the session data and drop the DashMap ref immediately // so the shard is not held during the expensive assembly step. let session = { - let entry = self.sessions + let entry = self + .sessions .get(upload_id) .ok_or_else(|| format!("Upload session not found: {}", upload_id))?; @@ -640,12 +635,17 @@ impl ChunkedUploadService { let chunks_meta: Vec<(usize, PathBuf)> = session .chunks .iter() - .map(|c| (c.index, session.temp_dir.join(format!("chunk_{:06}", c.index)))) + .map(|c| { + ( + c.index, + session.temp_dir.join(format!("chunk_{:06}", c.index)), + ) + }) .collect(); let total_size = session.total_size; let hash = tokio::task::spawn_blocking(move || -> Result { - use std::io::{Read, Write, BufWriter as StdBufWriter}; + use std::io::{BufWriter as StdBufWriter, Read, Write}; let raw_output = std::fs::OpenOptions::new() .create(true) @@ -920,8 +920,7 @@ mod tests { }; let json = serde_json::to_vec(&session).expect("serialise"); - let restored: UploadSession = - serde_json::from_slice(&json).expect("deserialise"); + let restored: UploadSession = serde_json::from_slice(&json).expect("deserialise"); assert_eq!(restored.id, session.id); assert_eq!(restored.filename, session.filename); @@ -996,7 +995,9 @@ mod tests { let recovered = ChunkedUploadService::recover_sessions(&base).await; assert_eq!(recovered.len(), 1); - let session = recovered.get(&upload_id).expect("session must be recovered"); + let session = recovered + .get(&upload_id) + .expect("session must be recovered"); assert_eq!(session.filename, "bigfile.bin"); assert_eq!(session.folder_id, Some("folder-x".into())); assert_eq!(session.chunks[0].status, ChunkStatus::Complete); @@ -1051,10 +1052,8 @@ mod tests { assert!(status.pending_chunks.is_empty()); // 4. Complete (assemble) - let (path, filename, _folder, _ct, size, hash) = service - .complete_upload_inner(&id) - .await - .expect("complete"); + let (path, filename, _folder, _ct, size, hash) = + service.complete_upload_inner(&id).await.expect("complete"); assert_eq!(filename, "test.txt"); assert_eq!(size, 1024); assert!(!hash.is_empty()); @@ -1078,7 +1077,13 @@ mod tests { let service = ChunkedUploadService::new(base.clone()).await; let resp = service - .create_session_inner("x.bin".into(), None, "application/octet-stream".into(), 512, Some(512)) + .create_session_inner( + "x.bin".into(), + None, + "application/octet-stream".into(), + 512, + Some(512), + ) .await .expect("create"); @@ -1156,12 +1161,18 @@ mod tests { chunk_size: 512, chunks: vec![ ChunkInfo { - index: 0, offset: 0, size: 512, - status: ChunkStatus::Pending, checksum: None, + index: 0, + offset: 0, + size: 512, + status: ChunkStatus::Pending, + checksum: None, }, ChunkInfo { - index: 1, offset: 512, size: 512, - status: ChunkStatus::Pending, checksum: None, + index: 1, + offset: 512, + size: 512, + status: ChunkStatus::Pending, + checksum: None, }, ], created_at: Utc::now(), @@ -1172,14 +1183,20 @@ mod tests { // Write metadata let json = serde_json::to_vec(&session).unwrap(); - fs::write(session_dir.join(SESSION_META_FILE), &json).await.unwrap(); + fs::write(session_dir.join(SESSION_META_FILE), &json) + .await + .unwrap(); // Write progress marking both chunks complete let bitmask = vec![0b00000011u8]; // bits 0 and 1 - fs::write(session_dir.join(PROGRESS_FILE), &bitmask).await.unwrap(); + fs::write(session_dir.join(PROGRESS_FILE), &bitmask) + .await + .unwrap(); // But only create chunk_000000 on disk — chunk_000001 is "missing" - fs::write(session_dir.join("chunk_000000"), &[0u8; 512]).await.unwrap(); + fs::write(session_dir.join("chunk_000000"), &[0u8; 512]) + .await + .unwrap(); let recovered = ChunkedUploadService::recover_sessions(&base).await; let s = recovered.get("partial-session").expect("must be recovered");