style: cargo fmt --all
This commit is contained in:
@@ -733,18 +733,16 @@ impl BatchOperationService {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// ── Finalize ─────────────────────────────────────────────────────
|
// ── Finalize ─────────────────────────────────────────────────────
|
||||||
let mut compat_writer = zip.close().await.map_err(|e| {
|
let mut compat_writer = zip
|
||||||
BatchOperationError::Internal(format!("ZIP finalize error: {}", e))
|
.close()
|
||||||
})?;
|
.await
|
||||||
compat_writer.close().await.map_err(|e| {
|
.map_err(|e| BatchOperationError::Internal(format!("ZIP finalize error: {}", e)))?;
|
||||||
BatchOperationError::Internal(format!("ZIP flush error: {}", e))
|
compat_writer
|
||||||
})?;
|
.close()
|
||||||
|
.await
|
||||||
|
.map_err(|e| BatchOperationError::Internal(format!("ZIP flush error: {}", e)))?;
|
||||||
|
|
||||||
let file_size = temp
|
let file_size = temp.as_file().metadata().map(|m| m.len()).unwrap_or(0);
|
||||||
.as_file()
|
|
||||||
.metadata()
|
|
||||||
.map(|m| m.len())
|
|
||||||
.unwrap_or(0);
|
|
||||||
|
|
||||||
info!(
|
info!(
|
||||||
"Batch download ZIP created: {} bytes in {}ms",
|
"Batch download ZIP created: {} bytes in {}ms",
|
||||||
@@ -776,17 +774,18 @@ impl BatchOperationService {
|
|||||||
let mut stream = std::pin::Pin::from(stream);
|
let mut stream = std::pin::Pin::from(stream);
|
||||||
|
|
||||||
while let Some(chunk) = stream.next().await {
|
while let Some(chunk) = stream.next().await {
|
||||||
let bytes = chunk.map_err(|e| {
|
let bytes =
|
||||||
BatchOperationError::Internal(format!("stream read: {}", e))
|
chunk.map_err(|e| BatchOperationError::Internal(format!("stream read: {}", e)))?;
|
||||||
})?;
|
writer
|
||||||
writer.write_all(&bytes).await.map_err(|e| {
|
.write_all(&bytes)
|
||||||
BatchOperationError::Internal(format!("zip chunk write: {}", e))
|
.await
|
||||||
})?;
|
.map_err(|e| BatchOperationError::Internal(format!("zip chunk write: {}", e)))?;
|
||||||
}
|
}
|
||||||
|
|
||||||
writer.close().await.map_err(|e| {
|
writer
|
||||||
BatchOperationError::Internal(format!("zip entry close: {}", e))
|
.close()
|
||||||
})?;
|
.await
|
||||||
|
.map_err(|e| BatchOperationError::Internal(format!("zip entry close: {}", e)))?;
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -194,12 +194,11 @@ impl FileUploadUseCase for FileUploadService {
|
|||||||
};
|
};
|
||||||
|
|
||||||
// Spool to temp file + hash
|
// Spool to temp file + hash
|
||||||
let temp = tempfile::NamedTempFile::new().map_err(|e| {
|
let temp = tempfile::NamedTempFile::new()
|
||||||
DomainError::internal_error("FileUpload", format!("temp file: {e}"))
|
.map_err(|e| DomainError::internal_error("FileUpload", format!("temp file: {e}")))?;
|
||||||
})?;
|
tokio::fs::write(temp.path(), content)
|
||||||
tokio::fs::write(temp.path(), content).await.map_err(|e| {
|
.await
|
||||||
DomainError::internal_error("FileUpload", format!("write temp: {e}"))
|
.map_err(|e| DomainError::internal_error("FileUpload", format!("write temp: {e}")))?;
|
||||||
})?;
|
|
||||||
let hash = hex::encode(Sha256::digest(content));
|
let hash = hex::encode(Sha256::digest(content));
|
||||||
|
|
||||||
let file = self
|
let file = self
|
||||||
@@ -224,12 +223,11 @@ impl FileUploadUseCase for FileUploadService {
|
|||||||
/// then delegates to the streaming update/create path.
|
/// then delegates to the streaming update/create path.
|
||||||
async fn update_file(&self, path: &str, content: &[u8]) -> Result<(), DomainError> {
|
async fn update_file(&self, path: &str, content: &[u8]) -> Result<(), DomainError> {
|
||||||
// Spool to temp file + hash
|
// Spool to temp file + hash
|
||||||
let temp = tempfile::NamedTempFile::new().map_err(|e| {
|
let temp = tempfile::NamedTempFile::new()
|
||||||
DomainError::internal_error("FileUpload", format!("temp file: {e}"))
|
.map_err(|e| DomainError::internal_error("FileUpload", format!("temp file: {e}")))?;
|
||||||
})?;
|
tokio::fs::write(temp.path(), content)
|
||||||
tokio::fs::write(temp.path(), content).await.map_err(|e| {
|
.await
|
||||||
DomainError::internal_error("FileUpload", format!("write temp: {e}"))
|
.map_err(|e| DomainError::internal_error("FileUpload", format!("write temp: {e}")))?;
|
||||||
})?;
|
|
||||||
let hash = hex::encode(Sha256::digest(content));
|
let hash = hex::encode(Sha256::digest(content));
|
||||||
|
|
||||||
self.update_file_streaming(
|
self.update_file_streaming(
|
||||||
|
|||||||
@@ -272,7 +272,11 @@ impl SearchUseCase for SearchService {
|
|||||||
* - Human-readable size formatting
|
* - Human-readable size formatting
|
||||||
* - Pagination
|
* - Pagination
|
||||||
*/
|
*/
|
||||||
async fn search(&self, criteria: SearchCriteriaDto, user_id: &str) -> Result<Arc<SearchResultsDto>> {
|
async fn search(
|
||||||
|
&self,
|
||||||
|
criteria: SearchCriteriaDto,
|
||||||
|
user_id: &str,
|
||||||
|
) -> Result<Arc<SearchResultsDto>> {
|
||||||
let start = Instant::now();
|
let start = Instant::now();
|
||||||
|
|
||||||
// Try to get from cache
|
// Try to get from cache
|
||||||
@@ -369,7 +373,8 @@ impl SearchUseCase for SearchService {
|
|||||||
criteria.sort_by.clone(),
|
criteria.sort_by.clone(),
|
||||||
));
|
));
|
||||||
|
|
||||||
self.store_in_cache(cache_key, Arc::clone(&search_results)).await;
|
self.store_in_cache(cache_key, Arc::clone(&search_results))
|
||||||
|
.await;
|
||||||
return Ok(search_results);
|
return Ok(search_results);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -454,7 +459,8 @@ impl SearchUseCase for SearchService {
|
|||||||
));
|
));
|
||||||
|
|
||||||
// Store in cache — Arc::clone is ~1 ns (atomic increment)
|
// Store in cache — Arc::clone is ~1 ns (atomic increment)
|
||||||
self.store_in_cache(cache_key, Arc::clone(&search_results)).await;
|
self.store_in_cache(cache_key, Arc::clone(&search_results))
|
||||||
|
.await;
|
||||||
|
|
||||||
Ok(search_results)
|
Ok(search_results)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -510,7 +510,13 @@ mod tests {
|
|||||||
&self,
|
&self,
|
||||||
_folder_id: &str,
|
_folder_id: &str,
|
||||||
) -> Result<
|
) -> Result<
|
||||||
std::pin::Pin<Box<dyn futures::Stream<Item = Result<crate::domain::entities::file::File, DomainError>> + Send>>,
|
std::pin::Pin<
|
||||||
|
Box<
|
||||||
|
dyn futures::Stream<
|
||||||
|
Item = Result<crate::domain::entities::file::File, DomainError>,
|
||||||
|
> + Send,
|
||||||
|
>,
|
||||||
|
>,
|
||||||
DomainError,
|
DomainError,
|
||||||
> {
|
> {
|
||||||
Ok(Box::pin(futures::stream::empty()))
|
Ok(Box::pin(futures::stream::empty()))
|
||||||
|
|||||||
@@ -208,7 +208,10 @@ impl FileReadPort for MockFileRepository {
|
|||||||
async fn stream_files_in_subtree(
|
async fn stream_files_in_subtree(
|
||||||
&self,
|
&self,
|
||||||
_folder_id: &str,
|
_folder_id: &str,
|
||||||
) -> std::result::Result<Pin<Box<dyn Stream<Item = std::result::Result<File, DomainError>> + Send>>, DomainError> {
|
) -> std::result::Result<
|
||||||
|
Pin<Box<dyn Stream<Item = std::result::Result<File, DomainError>> + Send>>,
|
||||||
|
DomainError,
|
||||||
|
> {
|
||||||
Ok(Box::pin(futures::stream::empty()))
|
Ok(Box::pin(futures::stream::empty()))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -502,7 +505,10 @@ mod tests {
|
|||||||
// Arrange
|
// Arrange
|
||||||
let trashed_files = Arc::new(Mutex::new(HashMap::new()));
|
let trashed_files = Arc::new(Mutex::new(HashMap::new()));
|
||||||
let trashed_folders = Arc::new(Mutex::new(HashMap::new()));
|
let trashed_folders = Arc::new(Mutex::new(HashMap::new()));
|
||||||
let trash_repo = Arc::new(MockTrashRepository::new(trashed_files.clone(), trashed_folders.clone()));
|
let trash_repo = Arc::new(MockTrashRepository::new(
|
||||||
|
trashed_files.clone(),
|
||||||
|
trashed_folders.clone(),
|
||||||
|
));
|
||||||
let file_repo = Arc::new(MockFileRepository::new(trashed_files));
|
let file_repo = Arc::new(MockFileRepository::new(trashed_files));
|
||||||
let folder_repo = Arc::new(MockFolderRepository::new(trashed_folders));
|
let folder_repo = Arc::new(MockFolderRepository::new(trashed_folders));
|
||||||
|
|
||||||
@@ -573,7 +579,10 @@ mod tests {
|
|||||||
// Arrange
|
// Arrange
|
||||||
let trashed_files = Arc::new(Mutex::new(HashMap::new()));
|
let trashed_files = Arc::new(Mutex::new(HashMap::new()));
|
||||||
let trashed_folders = Arc::new(Mutex::new(HashMap::new()));
|
let trashed_folders = Arc::new(Mutex::new(HashMap::new()));
|
||||||
let trash_repo = Arc::new(MockTrashRepository::new(trashed_files.clone(), trashed_folders.clone()));
|
let trash_repo = Arc::new(MockTrashRepository::new(
|
||||||
|
trashed_files.clone(),
|
||||||
|
trashed_folders.clone(),
|
||||||
|
));
|
||||||
let file_repo = Arc::new(MockFileRepository::new(trashed_files));
|
let file_repo = Arc::new(MockFileRepository::new(trashed_files));
|
||||||
let folder_repo = Arc::new(MockFolderRepository::new(trashed_folders));
|
let folder_repo = Arc::new(MockFolderRepository::new(trashed_folders));
|
||||||
|
|
||||||
@@ -635,7 +644,10 @@ mod tests {
|
|||||||
// Arrange
|
// Arrange
|
||||||
let trashed_files = Arc::new(Mutex::new(HashMap::new()));
|
let trashed_files = Arc::new(Mutex::new(HashMap::new()));
|
||||||
let trashed_folders = Arc::new(Mutex::new(HashMap::new()));
|
let trashed_folders = Arc::new(Mutex::new(HashMap::new()));
|
||||||
let trash_repo = Arc::new(MockTrashRepository::new(trashed_files.clone(), trashed_folders.clone()));
|
let trash_repo = Arc::new(MockTrashRepository::new(
|
||||||
|
trashed_files.clone(),
|
||||||
|
trashed_folders.clone(),
|
||||||
|
));
|
||||||
let file_repo = Arc::new(MockFileRepository::new(trashed_files));
|
let file_repo = Arc::new(MockFileRepository::new(trashed_files));
|
||||||
let folder_repo = Arc::new(MockFolderRepository::new(trashed_folders));
|
let folder_repo = Arc::new(MockFolderRepository::new(trashed_folders));
|
||||||
|
|
||||||
@@ -702,7 +714,10 @@ mod tests {
|
|||||||
// Arrange
|
// Arrange
|
||||||
let trashed_files = Arc::new(Mutex::new(HashMap::new()));
|
let trashed_files = Arc::new(Mutex::new(HashMap::new()));
|
||||||
let trashed_folders = Arc::new(Mutex::new(HashMap::new()));
|
let trashed_folders = Arc::new(Mutex::new(HashMap::new()));
|
||||||
let trash_repo = Arc::new(MockTrashRepository::new(trashed_files.clone(), trashed_folders.clone()));
|
let trash_repo = Arc::new(MockTrashRepository::new(
|
||||||
|
trashed_files.clone(),
|
||||||
|
trashed_folders.clone(),
|
||||||
|
));
|
||||||
let file_repo = Arc::new(MockFileRepository::new(trashed_files));
|
let file_repo = Arc::new(MockFileRepository::new(trashed_files));
|
||||||
let folder_repo = Arc::new(MockFolderRepository::new(trashed_folders));
|
let folder_repo = Arc::new(MockFolderRepository::new(trashed_folders));
|
||||||
|
|
||||||
@@ -768,7 +783,10 @@ mod tests {
|
|||||||
// Arrange
|
// Arrange
|
||||||
let trashed_files = Arc::new(Mutex::new(HashMap::new()));
|
let trashed_files = Arc::new(Mutex::new(HashMap::new()));
|
||||||
let trashed_folders = Arc::new(Mutex::new(HashMap::new()));
|
let trashed_folders = Arc::new(Mutex::new(HashMap::new()));
|
||||||
let trash_repo = Arc::new(MockTrashRepository::new(trashed_files.clone(), trashed_folders.clone()));
|
let trash_repo = Arc::new(MockTrashRepository::new(
|
||||||
|
trashed_files.clone(),
|
||||||
|
trashed_folders.clone(),
|
||||||
|
));
|
||||||
let file_repo = Arc::new(MockFileRepository::new(trashed_files));
|
let file_repo = Arc::new(MockFileRepository::new(trashed_files));
|
||||||
let folder_repo = Arc::new(MockFolderRepository::new(trashed_folders));
|
let folder_repo = Arc::new(MockFolderRepository::new(trashed_folders));
|
||||||
|
|
||||||
|
|||||||
+7
-4
@@ -36,10 +36,10 @@ use crate::application::services::{
|
|||||||
use crate::common::config::AppConfig;
|
use crate::common::config::AppConfig;
|
||||||
use crate::common::errors::DomainError;
|
use crate::common::errors::DomainError;
|
||||||
use crate::domain::services::i18n_service::I18nService;
|
use crate::domain::services::i18n_service::I18nService;
|
||||||
|
use crate::infrastructure::repositories::pg::SharePgRepository;
|
||||||
use crate::infrastructure::repositories::pg::{
|
use crate::infrastructure::repositories::pg::{
|
||||||
FileBlobReadRepository, FileBlobWriteRepository, FolderDbRepository, TrashDbRepository,
|
FileBlobReadRepository, FileBlobWriteRepository, FolderDbRepository, TrashDbRepository,
|
||||||
};
|
};
|
||||||
use crate::infrastructure::repositories::pg::SharePgRepository;
|
|
||||||
use crate::infrastructure::services::file_content_cache::{
|
use crate::infrastructure::services::file_content_cache::{
|
||||||
FileContentCache, FileContentCacheConfig,
|
FileContentCache, FileContentCacheConfig,
|
||||||
};
|
};
|
||||||
@@ -333,11 +333,13 @@ impl AppServiceFactory {
|
|||||||
|
|
||||||
// Build a password hasher for share password verification
|
// Build a password hasher for share password verification
|
||||||
let password_hasher: Arc<dyn crate::application::ports::auth_ports::PasswordHasherPort> =
|
let password_hasher: Arc<dyn crate::application::ports::auth_ports::PasswordHasherPort> =
|
||||||
Arc::new(crate::infrastructure::services::password_hasher::Argon2PasswordHasher::new(
|
Arc::new(
|
||||||
|
crate::infrastructure::services::password_hasher::Argon2PasswordHasher::new(
|
||||||
self.config.auth.hash_memory_cost,
|
self.config.auth.hash_memory_cost,
|
||||||
self.config.auth.hash_time_cost,
|
self.config.auth.hash_time_cost,
|
||||||
self.config.auth.hash_parallelism,
|
self.config.auth.hash_parallelism,
|
||||||
));
|
),
|
||||||
|
);
|
||||||
|
|
||||||
let service = Arc::new(ShareService::new(
|
let service = Arc::new(ShareService::new(
|
||||||
Arc::new(self.config.clone()),
|
Arc::new(self.config.clone()),
|
||||||
@@ -469,7 +471,8 @@ impl AppServiceFactory {
|
|||||||
recent_service = Some(recent.clone());
|
recent_service = Some(recent.clone());
|
||||||
apps.recent_service = Some(recent);
|
apps.recent_service = Some(recent);
|
||||||
|
|
||||||
storage_usage_service = Some(self.create_storage_usage_service(&repos, &pool, &maintenance_pool));
|
storage_usage_service =
|
||||||
|
Some(self.create_storage_usage_service(&repos, &pool, &maintenance_pool));
|
||||||
|
|
||||||
// Auth services
|
// Auth services
|
||||||
if self.config.features.enable_auth {
|
if self.config.features.enable_auth {
|
||||||
|
|||||||
@@ -150,27 +150,29 @@ impl TrashRepository for TrashDbRepository {
|
|||||||
async fn delete_expired_bulk(&self) -> Result<(u64, u64)> {
|
async fn delete_expired_bulk(&self) -> Result<(u64, u64)> {
|
||||||
let cutoff = Utc::now() - chrono::Duration::days(self.retention_days);
|
let cutoff = Utc::now() - chrono::Duration::days(self.retention_days);
|
||||||
|
|
||||||
let mut tx = self.pool.begin().await.map_err(|e| {
|
let mut tx = self
|
||||||
DomainError::internal_error("TrashDb", format!("begin tx: {e}"))
|
.pool
|
||||||
})?;
|
.begin()
|
||||||
|
.await
|
||||||
|
.map_err(|e| DomainError::internal_error("TrashDb", format!("begin tx: {e}")))?;
|
||||||
|
|
||||||
// 1. Bulk-delete expired trashed files.
|
// 1. Bulk-delete expired trashed files.
|
||||||
// The PG trigger `trg_files_decrement_blob_ref` automatically
|
// The PG trigger `trg_files_decrement_blob_ref` automatically
|
||||||
// decrements blob ref_count for every deleted row.
|
// decrements blob ref_count for every deleted row.
|
||||||
let files_deleted = sqlx::query(
|
let files_deleted =
|
||||||
"DELETE FROM storage.files WHERE is_trashed = TRUE AND trashed_at < $1",
|
sqlx::query("DELETE FROM storage.files WHERE is_trashed = TRUE AND trashed_at < $1")
|
||||||
)
|
|
||||||
.bind(cutoff)
|
.bind(cutoff)
|
||||||
.execute(&mut *tx)
|
.execute(&mut *tx)
|
||||||
.await
|
.await
|
||||||
.map_err(|e| DomainError::internal_error("TrashDb", format!("bulk delete files: {e}")))?
|
.map_err(|e| {
|
||||||
|
DomainError::internal_error("TrashDb", format!("bulk delete files: {e}"))
|
||||||
|
})?
|
||||||
.rows_affected();
|
.rows_affected();
|
||||||
|
|
||||||
// 2. Bulk-delete expired trashed folders.
|
// 2. Bulk-delete expired trashed folders.
|
||||||
// FK ON DELETE CASCADE handles descendant folders and their files.
|
// FK ON DELETE CASCADE handles descendant folders and their files.
|
||||||
let folders_deleted = sqlx::query(
|
let folders_deleted =
|
||||||
"DELETE FROM storage.folders WHERE is_trashed = TRUE AND trashed_at < $1",
|
sqlx::query("DELETE FROM storage.folders WHERE is_trashed = TRUE AND trashed_at < $1")
|
||||||
)
|
|
||||||
.bind(cutoff)
|
.bind(cutoff)
|
||||||
.execute(&mut *tx)
|
.execute(&mut *tx)
|
||||||
.await
|
.await
|
||||||
@@ -179,9 +181,9 @@ impl TrashRepository for TrashDbRepository {
|
|||||||
})?
|
})?
|
||||||
.rows_affected();
|
.rows_affected();
|
||||||
|
|
||||||
tx.commit().await.map_err(|e| {
|
tx.commit()
|
||||||
DomainError::internal_error("TrashDb", format!("commit tx: {e}"))
|
.await
|
||||||
})?;
|
.map_err(|e| DomainError::internal_error("TrashDb", format!("commit tx: {e}")))?;
|
||||||
|
|
||||||
Ok((files_deleted, folders_deleted))
|
Ok((files_deleted, folders_deleted))
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -157,8 +157,8 @@ impl UploadSession {
|
|||||||
/// Persist the full session metadata once (on create).
|
/// Persist the full session metadata once (on create).
|
||||||
async fn persist_metadata(&self) -> Result<(), String> {
|
async fn persist_metadata(&self) -> Result<(), String> {
|
||||||
let path = self.temp_dir.join(SESSION_META_FILE);
|
let path = self.temp_dir.join(SESSION_META_FILE);
|
||||||
let json = serde_json::to_vec(self)
|
let json =
|
||||||
.map_err(|e| format!("Failed to serialise session: {e}"))?;
|
serde_json::to_vec(self).map_err(|e| format!("Failed to serialise session: {e}"))?;
|
||||||
// Atomic write: write to .tmp then rename
|
// Atomic write: write to .tmp then rename
|
||||||
let tmp = self.temp_dir.join("session.json.tmp");
|
let tmp = self.temp_dir.join("session.json.tmp");
|
||||||
fs::write(&tmp, &json)
|
fs::write(&tmp, &json)
|
||||||
@@ -213,9 +213,7 @@ impl ChunkedUploadService {
|
|||||||
};
|
};
|
||||||
|
|
||||||
if recovered_count > 0 {
|
if recovered_count > 0 {
|
||||||
tracing::info!(
|
tracing::info!("♻️ Recovered {recovered_count} chunked-upload session(s) from disk");
|
||||||
"♻️ Recovered {recovered_count} chunked-upload session(s) from disk"
|
|
||||||
);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// Start cleanup task
|
// Start cleanup task
|
||||||
@@ -325,10 +323,7 @@ impl ChunkedUploadService {
|
|||||||
// ── Cleanup ──────────────────────────────────────────────────────────
|
// ── Cleanup ──────────────────────────────────────────────────────────
|
||||||
|
|
||||||
/// Background task to clean expired sessions
|
/// Background task to clean expired sessions
|
||||||
async fn cleanup_loop(
|
async fn cleanup_loop(sessions: Arc<DashMap<String, UploadSession>>, temp_base_dir: PathBuf) {
|
||||||
sessions: Arc<DashMap<String, UploadSession>>,
|
|
||||||
temp_base_dir: PathBuf,
|
|
||||||
) {
|
|
||||||
let mut interval = tokio::time::interval(Duration::from_secs(3600)); // Every hour
|
let mut interval = tokio::time::interval(Duration::from_secs(3600)); // Every hour
|
||||||
|
|
||||||
loop {
|
loop {
|
||||||
@@ -463,7 +458,8 @@ impl ChunkedUploadService {
|
|||||||
) -> Result<ChunkUploadResponseDto, String> {
|
) -> Result<ChunkUploadResponseDto, String> {
|
||||||
// Validate session exists and chunk index is valid
|
// Validate session exists and chunk index is valid
|
||||||
let (chunk_path, expected_size) = {
|
let (chunk_path, expected_size) = {
|
||||||
let session = self.sessions
|
let session = self
|
||||||
|
.sessions
|
||||||
.get(upload_id)
|
.get(upload_id)
|
||||||
.ok_or_else(|| format!("Upload session not found: {}", upload_id))?;
|
.ok_or_else(|| format!("Upload session not found: {}", upload_id))?;
|
||||||
|
|
||||||
@@ -500,9 +496,8 @@ impl ChunkedUploadService {
|
|||||||
// worker free for other connections.
|
// worker free for other connections.
|
||||||
if let Some(ref expected_checksum) = checksum {
|
if let Some(ref expected_checksum) = checksum {
|
||||||
let data_clone = data.clone(); // Bytes::clone is O(1) — just an Arc increment
|
let data_clone = data.clone(); // Bytes::clone is O(1) — just an Arc increment
|
||||||
let actual_checksum = tokio::task::spawn_blocking(move || {
|
let actual_checksum =
|
||||||
format!("{:x}", md5::compute(&data_clone))
|
tokio::task::spawn_blocking(move || format!("{:x}", md5::compute(&data_clone)))
|
||||||
})
|
|
||||||
.await
|
.await
|
||||||
.map_err(|e| format!("MD5 checksum task failed: {e}"))?;
|
.map_err(|e| format!("MD5 checksum task failed: {e}"))?;
|
||||||
|
|
||||||
@@ -527,7 +522,8 @@ impl ChunkedUploadService {
|
|||||||
// Disk I/O (persist_progress) is done AFTER the ref is dropped so
|
// Disk I/O (persist_progress) is done AFTER the ref is dropped so
|
||||||
// concurrent uploads to other sessions are never blocked.
|
// concurrent uploads to other sessions are never blocked.
|
||||||
let (bytes_received, progress, is_complete, persist_path, persist_bitmask) = {
|
let (bytes_received, progress, is_complete, persist_path, persist_bitmask) = {
|
||||||
let mut session = self.sessions
|
let mut session = self
|
||||||
|
.sessions
|
||||||
.get_mut(upload_id)
|
.get_mut(upload_id)
|
||||||
.ok_or_else(|| "Session disappeared".to_string())?;
|
.ok_or_else(|| "Session disappeared".to_string())?;
|
||||||
|
|
||||||
@@ -571,11 +567,9 @@ impl ChunkedUploadService {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/// Get upload status
|
/// Get upload status
|
||||||
async fn get_status_inner(
|
async fn get_status_inner(&self, upload_id: &str) -> Result<UploadStatusResponseDto, String> {
|
||||||
&self,
|
let session = self
|
||||||
upload_id: &str,
|
.sessions
|
||||||
) -> Result<UploadStatusResponseDto, String> {
|
|
||||||
let session = self.sessions
|
|
||||||
.get(upload_id)
|
.get(upload_id)
|
||||||
.ok_or_else(|| format!("Upload session not found: {}", 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
|
// Clone the session data and drop the DashMap ref immediately
|
||||||
// so the shard is not held during the expensive assembly step.
|
// so the shard is not held during the expensive assembly step.
|
||||||
let session = {
|
let session = {
|
||||||
let entry = self.sessions
|
let entry = self
|
||||||
|
.sessions
|
||||||
.get(upload_id)
|
.get(upload_id)
|
||||||
.ok_or_else(|| format!("Upload session not found: {}", 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
|
let chunks_meta: Vec<(usize, PathBuf)> = session
|
||||||
.chunks
|
.chunks
|
||||||
.iter()
|
.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();
|
.collect();
|
||||||
let total_size = session.total_size;
|
let total_size = session.total_size;
|
||||||
|
|
||||||
let hash = tokio::task::spawn_blocking(move || -> Result<String, String> {
|
let hash = tokio::task::spawn_blocking(move || -> Result<String, String> {
|
||||||
use std::io::{Read, Write, BufWriter as StdBufWriter};
|
use std::io::{BufWriter as StdBufWriter, Read, Write};
|
||||||
|
|
||||||
let raw_output = std::fs::OpenOptions::new()
|
let raw_output = std::fs::OpenOptions::new()
|
||||||
.create(true)
|
.create(true)
|
||||||
@@ -920,8 +920,7 @@ mod tests {
|
|||||||
};
|
};
|
||||||
|
|
||||||
let json = serde_json::to_vec(&session).expect("serialise");
|
let json = serde_json::to_vec(&session).expect("serialise");
|
||||||
let restored: UploadSession =
|
let restored: UploadSession = serde_json::from_slice(&json).expect("deserialise");
|
||||||
serde_json::from_slice(&json).expect("deserialise");
|
|
||||||
|
|
||||||
assert_eq!(restored.id, session.id);
|
assert_eq!(restored.id, session.id);
|
||||||
assert_eq!(restored.filename, session.filename);
|
assert_eq!(restored.filename, session.filename);
|
||||||
@@ -996,7 +995,9 @@ mod tests {
|
|||||||
|
|
||||||
let recovered = ChunkedUploadService::recover_sessions(&base).await;
|
let recovered = ChunkedUploadService::recover_sessions(&base).await;
|
||||||
assert_eq!(recovered.len(), 1);
|
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.filename, "bigfile.bin");
|
||||||
assert_eq!(session.folder_id, Some("folder-x".into()));
|
assert_eq!(session.folder_id, Some("folder-x".into()));
|
||||||
assert_eq!(session.chunks[0].status, ChunkStatus::Complete);
|
assert_eq!(session.chunks[0].status, ChunkStatus::Complete);
|
||||||
@@ -1051,10 +1052,8 @@ mod tests {
|
|||||||
assert!(status.pending_chunks.is_empty());
|
assert!(status.pending_chunks.is_empty());
|
||||||
|
|
||||||
// 4. Complete (assemble)
|
// 4. Complete (assemble)
|
||||||
let (path, filename, _folder, _ct, size, hash) = service
|
let (path, filename, _folder, _ct, size, hash) =
|
||||||
.complete_upload_inner(&id)
|
service.complete_upload_inner(&id).await.expect("complete");
|
||||||
.await
|
|
||||||
.expect("complete");
|
|
||||||
assert_eq!(filename, "test.txt");
|
assert_eq!(filename, "test.txt");
|
||||||
assert_eq!(size, 1024);
|
assert_eq!(size, 1024);
|
||||||
assert!(!hash.is_empty());
|
assert!(!hash.is_empty());
|
||||||
@@ -1078,7 +1077,13 @@ mod tests {
|
|||||||
let service = ChunkedUploadService::new(base.clone()).await;
|
let service = ChunkedUploadService::new(base.clone()).await;
|
||||||
|
|
||||||
let resp = service
|
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
|
.await
|
||||||
.expect("create");
|
.expect("create");
|
||||||
|
|
||||||
@@ -1156,12 +1161,18 @@ mod tests {
|
|||||||
chunk_size: 512,
|
chunk_size: 512,
|
||||||
chunks: vec![
|
chunks: vec![
|
||||||
ChunkInfo {
|
ChunkInfo {
|
||||||
index: 0, offset: 0, size: 512,
|
index: 0,
|
||||||
status: ChunkStatus::Pending, checksum: None,
|
offset: 0,
|
||||||
|
size: 512,
|
||||||
|
status: ChunkStatus::Pending,
|
||||||
|
checksum: None,
|
||||||
},
|
},
|
||||||
ChunkInfo {
|
ChunkInfo {
|
||||||
index: 1, offset: 512, size: 512,
|
index: 1,
|
||||||
status: ChunkStatus::Pending, checksum: None,
|
offset: 512,
|
||||||
|
size: 512,
|
||||||
|
status: ChunkStatus::Pending,
|
||||||
|
checksum: None,
|
||||||
},
|
},
|
||||||
],
|
],
|
||||||
created_at: Utc::now(),
|
created_at: Utc::now(),
|
||||||
@@ -1172,14 +1183,20 @@ mod tests {
|
|||||||
|
|
||||||
// Write metadata
|
// Write metadata
|
||||||
let json = serde_json::to_vec(&session).unwrap();
|
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
|
// Write progress marking both chunks complete
|
||||||
let bitmask = vec![0b00000011u8]; // bits 0 and 1
|
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"
|
// 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 recovered = ChunkedUploadService::recover_sessions(&base).await;
|
||||||
let s = recovered.get("partial-session").expect("must be recovered");
|
let s = recovered.get("partial-session").expect("must be recovered");
|
||||||
|
|||||||
@@ -769,9 +769,7 @@ impl DedupService {
|
|||||||
.bind(BATCH_SIZE)
|
.bind(BATCH_SIZE)
|
||||||
.fetch_all(self.maintenance_pool.as_ref())
|
.fetch_all(self.maintenance_pool.as_ref())
|
||||||
.await
|
.await
|
||||||
.map_err(|e| {
|
.map_err(|e| DomainError::internal_error("Dedup", format!("GC batch failed: {e}")))?;
|
||||||
DomainError::internal_error("Dedup", format!("GC batch failed: {e}"))
|
|
||||||
})?;
|
|
||||||
|
|
||||||
if batch.is_empty() {
|
if batch.is_empty() {
|
||||||
break;
|
break;
|
||||||
|
|||||||
@@ -192,11 +192,7 @@ impl ThumbnailService {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// 2. Generate thumbnail (CPU-bound, runs in spawn_blocking)
|
// 2. Generate thumbnail (CPU-bound, runs in spawn_blocking)
|
||||||
tracing::info!(
|
tracing::info!("🎨 Generating thumbnail: {} {:?}", file_id_owned, size);
|
||||||
"🎨 Generating thumbnail: {} {:?}",
|
|
||||||
file_id_owned,
|
|
||||||
size
|
|
||||||
);
|
|
||||||
match self.generate_thumbnail(&original_owned, size).await {
|
match self.generate_thumbnail(&original_owned, size).await {
|
||||||
Ok(bytes) => {
|
Ok(bytes) => {
|
||||||
// Save to disk (best-effort — don't fail the request)
|
// Save to disk (best-effort — don't fail the request)
|
||||||
@@ -245,14 +241,17 @@ impl ThumbnailService {
|
|||||||
let max_dim = size.max_dimension();
|
let max_dim = size.max_dimension();
|
||||||
|
|
||||||
// Acquire semaphore permit — bounds peak RAM from concurrent decodes
|
// Acquire semaphore permit — bounds peak RAM from concurrent decodes
|
||||||
let _permit = self.decode_semaphore.acquire().await
|
let _permit = self
|
||||||
|
.decode_semaphore
|
||||||
|
.acquire()
|
||||||
|
.await
|
||||||
.map_err(|_| ThumbnailError::TaskError("Decode semaphore closed".into()))?;
|
.map_err(|_| ThumbnailError::TaskError("Decode semaphore closed".into()))?;
|
||||||
|
|
||||||
// Run image processing in blocking thread pool
|
// Run image processing in blocking thread pool
|
||||||
let result = tokio::task::spawn_blocking(move || -> Result<Vec<u8>, ThumbnailError> {
|
let result = tokio::task::spawn_blocking(move || -> Result<Vec<u8>, ThumbnailError> {
|
||||||
// Single read: load file once into memory, then work from the buffer
|
// Single read: load file once into memory, then work from the buffer
|
||||||
let data = std::fs::read(&path)
|
let data =
|
||||||
.map_err(|e| ThumbnailError::ImageError(e.to_string()))?;
|
std::fs::read(&path).map_err(|e| ThumbnailError::ImageError(e.to_string()))?;
|
||||||
|
|
||||||
// Safety check: read dimensions from in-memory buffer (no 2nd I/O)
|
// Safety check: read dimensions from in-memory buffer (no 2nd I/O)
|
||||||
let (w, h) = image::ImageReader::new(std::io::Cursor::new(&data))
|
let (w, h) = image::ImageReader::new(std::io::Cursor::new(&data))
|
||||||
@@ -318,7 +317,10 @@ impl ThumbnailService {
|
|||||||
let _permit = match self.decode_semaphore.acquire().await {
|
let _permit = match self.decode_semaphore.acquire().await {
|
||||||
Ok(p) => p,
|
Ok(p) => p,
|
||||||
Err(_) => {
|
Err(_) => {
|
||||||
tracing::warn!("Decode semaphore closed, skipping thumbnails for {}", file_id);
|
tracing::warn!(
|
||||||
|
"Decode semaphore closed, skipping thumbnails for {}",
|
||||||
|
file_id
|
||||||
|
);
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
@@ -328,8 +330,8 @@ impl ThumbnailService {
|
|||||||
// Single spawn_blocking: 1 read + 1 decode + 3 resize + 3 encode
|
// Single spawn_blocking: 1 read + 1 decode + 3 resize + 3 encode
|
||||||
let results = tokio::task::spawn_blocking(move || {
|
let results = tokio::task::spawn_blocking(move || {
|
||||||
// Single read: load file once into memory
|
// Single read: load file once into memory
|
||||||
let data = std::fs::read(&path)
|
let data =
|
||||||
.map_err(|e| ThumbnailError::ImageError(e.to_string()))?;
|
std::fs::read(&path).map_err(|e| ThumbnailError::ImageError(e.to_string()))?;
|
||||||
|
|
||||||
// Safety check: read dimensions from in-memory buffer (no 2nd I/O)
|
// Safety check: read dimensions from in-memory buffer (no 2nd I/O)
|
||||||
let (w, h) = image::ImageReader::new(std::io::Cursor::new(&data))
|
let (w, h) = image::ImageReader::new(std::io::Cursor::new(&data))
|
||||||
@@ -372,10 +374,7 @@ impl ThumbnailService {
|
|||||||
|
|
||||||
let mut buf = Vec::new();
|
let mut buf = Vec::new();
|
||||||
thumb
|
thumb
|
||||||
.write_to(
|
.write_to(&mut std::io::Cursor::new(&mut buf), ImageFormat::WebP)
|
||||||
&mut std::io::Cursor::new(&mut buf),
|
|
||||||
ImageFormat::WebP,
|
|
||||||
)
|
|
||||||
.map_err(|e| ThumbnailError::ImageError(e.to_string()))?;
|
.map_err(|e| ThumbnailError::ImageError(e.to_string()))?;
|
||||||
|
|
||||||
Ok((size, Bytes::from(buf)))
|
Ok((size, Bytes::from(buf)))
|
||||||
@@ -388,15 +387,11 @@ impl ThumbnailService {
|
|||||||
let thumbnails = match results {
|
let thumbnails = match results {
|
||||||
Ok(Ok(t)) => t,
|
Ok(Ok(t)) => t,
|
||||||
Ok(Err(e)) => {
|
Ok(Err(e)) => {
|
||||||
tracing::warn!(
|
tracing::warn!("Thumbnail generation failed for {}: {}", file_id, e);
|
||||||
"Thumbnail generation failed for {}: {}", file_id, e
|
|
||||||
);
|
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
tracing::warn!(
|
tracing::warn!("Thumbnail task panicked for {}: {}", file_id, e);
|
||||||
"Thumbnail task panicked for {}: {}", file_id, e
|
|
||||||
);
|
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|||||||
@@ -17,10 +17,7 @@ pub struct TrashCleanupService {
|
|||||||
}
|
}
|
||||||
|
|
||||||
impl TrashCleanupService {
|
impl TrashCleanupService {
|
||||||
pub fn new(
|
pub fn new(trash_repository: Arc<dyn TrashRepository>, cleanup_interval_hours: u64) -> Self {
|
||||||
trash_repository: Arc<dyn TrashRepository>,
|
|
||||||
cleanup_interval_hours: u64,
|
|
||||||
) -> Self {
|
|
||||||
Self {
|
Self {
|
||||||
trash_repository,
|
trash_repository,
|
||||||
cleanup_interval_hours: cleanup_interval_hours.max(1), // Minimum 1 hour
|
cleanup_interval_hours: cleanup_interval_hours.max(1), // Minimum 1 hour
|
||||||
@@ -60,9 +57,7 @@ impl TrashCleanupService {
|
|||||||
|
|
||||||
/// Bulk-delete all expired trash items in a single transaction.
|
/// Bulk-delete all expired trash items in a single transaction.
|
||||||
#[instrument(skip(trash_repository))]
|
#[instrument(skip(trash_repository))]
|
||||||
async fn cleanup_expired_items(
|
async fn cleanup_expired_items(trash_repository: Arc<dyn TrashRepository>) -> Result<()> {
|
||||||
trash_repository: Arc<dyn TrashRepository>,
|
|
||||||
) -> Result<()> {
|
|
||||||
debug!("Starting bulk cleanup of expired trash items");
|
debug!("Starting bulk cleanup of expired trash items");
|
||||||
|
|
||||||
let (files, folders) = trash_repository.delete_expired_bulk().await?;
|
let (files, folders) = trash_repository.delete_expired_bulk().await?;
|
||||||
|
|||||||
@@ -616,10 +616,12 @@ pub async fn download_batch(
|
|||||||
.as_file()
|
.as_file()
|
||||||
.metadata()
|
.metadata()
|
||||||
.map(|m| m.len())
|
.map(|m| m.len())
|
||||||
.map_err(|e| (
|
.map_err(|e| {
|
||||||
|
(
|
||||||
StatusCode::INTERNAL_SERVER_ERROR,
|
StatusCode::INTERNAL_SERVER_ERROR,
|
||||||
format!("Failed to read temp file metadata: {}", e),
|
format!("Failed to read temp file metadata: {}", e),
|
||||||
))?;
|
)
|
||||||
|
})?;
|
||||||
|
|
||||||
// Split into the already-open fd + auto-delete path
|
// Split into the already-open fd + auto-delete path
|
||||||
let (std_file, temp_path) = temp_file.into_parts();
|
let (std_file, temp_path) = temp_file.into_parts();
|
||||||
@@ -644,7 +646,9 @@ pub async fn download_batch(
|
|||||||
|
|
||||||
// Keep TempPath alive in response extensions so the file is only
|
// Keep TempPath alive in response extensions so the file is only
|
||||||
// deleted AFTER the body stream finishes sending.
|
// deleted AFTER the body stream finishes sending.
|
||||||
response.extensions_mut().insert(std::sync::Arc::new(temp_path));
|
response
|
||||||
|
.extensions_mut()
|
||||||
|
.insert(std::sync::Arc::new(temp_path));
|
||||||
|
|
||||||
Ok(response)
|
Ok(response)
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user