perf: optimize hot paths — batch concurrency, pagination, sorting, transcoding, folder ops
- batch_operations.rs: replace join_all with buffer_unordered, Arc<str> for shared IDs, remove redundant clones and dead Semaphore - folder_db_repository.rs: use COUNT(*) OVER() for single-query pagination; UPDATE RETURNING for rename/move (eliminates extra SELECTs) - folder_service.rs: remove StorageTransaction wrapper from rename/move — direct repo call (4→2 and 5→3 queries) - search_service.rs: replace sort_by(to_lowercase) with sort_by_cached_key (N vs 2·N·log₂N allocations) - image_transcode_service.rs: dynamic rayon pool sizing via available_parallelism() instead of hardcoded 2 threads - Remove dead transactions module (zero consumers after folder_service refactor)
This commit is contained in:
@@ -2,6 +2,5 @@ pub mod adapters;
|
||||
pub mod dtos;
|
||||
pub mod ports;
|
||||
pub mod services;
|
||||
pub mod transactions;
|
||||
|
||||
// Re-exportaciones para facilitar el acceso a los principales puertos
|
||||
|
||||
@@ -1,13 +1,12 @@
|
||||
use async_zip::base::write::ZipFileWriter;
|
||||
use async_zip::{Compression, ZipEntryBuilder};
|
||||
use futures::io::AsyncWriteExt as FuturesWriteExt;
|
||||
use futures::{Future, StreamExt, future::join_all};
|
||||
use futures::{Future, StreamExt, stream};
|
||||
use std::collections::HashMap;
|
||||
use std::sync::Arc;
|
||||
use tempfile::NamedTempFile;
|
||||
use thiserror::Error;
|
||||
use tokio::io::BufWriter;
|
||||
use tokio::sync::Semaphore;
|
||||
use tracing::info;
|
||||
|
||||
use crate::application::dtos::file_dto::FileDto;
|
||||
@@ -71,7 +70,6 @@ pub struct BatchOperationService {
|
||||
folder_service: Arc<FolderService>,
|
||||
trash_service: Option<Arc<dyn TrashUseCase>>,
|
||||
config: AppConfig,
|
||||
semaphore: Arc<Semaphore>,
|
||||
}
|
||||
|
||||
impl BatchOperationService {
|
||||
@@ -82,16 +80,12 @@ impl BatchOperationService {
|
||||
folder_service: Arc<FolderService>,
|
||||
config: AppConfig,
|
||||
) -> Self {
|
||||
// Limit concurrency based on configuration
|
||||
let max_concurrency = config.concurrency.max_concurrent_files;
|
||||
|
||||
Self {
|
||||
file_retrieval,
|
||||
file_management,
|
||||
folder_service,
|
||||
trash_service: None,
|
||||
config,
|
||||
semaphore: Arc::new(Semaphore::new(max_concurrency)),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -123,6 +117,7 @@ impl BatchOperationService {
|
||||
) -> Result<BatchResult<FileDto>, BatchOperationError> {
|
||||
info!("Starting batch copy of {} files", file_ids.len());
|
||||
let start_time = std::time::Instant::now();
|
||||
let max_concurrent = self.config.concurrency.max_concurrent_files;
|
||||
|
||||
// Create result structure
|
||||
let mut result = BatchResult {
|
||||
@@ -134,31 +129,23 @@ impl BatchOperationService {
|
||||
},
|
||||
};
|
||||
|
||||
// Define the operation to perform for each file
|
||||
let operations = file_ids.into_iter().map(|file_id| {
|
||||
// Arc<str> avoids N heap-clones of the same string
|
||||
let target_folder: Option<Arc<str>> = target_folder_id.map(|s| Arc::from(s.as_str()));
|
||||
|
||||
// buffer_unordered materialises only max_concurrent futures at a time
|
||||
let mut operation_stream = stream::iter(file_ids.into_iter().map(|file_id| {
|
||||
let mgmt = self.file_management.clone();
|
||||
let target_folder = target_folder_id.clone();
|
||||
let semaphore = self.semaphore.clone();
|
||||
let target_folder = target_folder.clone();
|
||||
|
||||
async move {
|
||||
// Acquire semaphore permit
|
||||
let permit = semaphore.acquire().await.unwrap();
|
||||
|
||||
let copy_result = mgmt.copy_file(&file_id, target_folder.clone()).await;
|
||||
|
||||
// Release the permit explicitly (also released on drop)
|
||||
drop(permit);
|
||||
|
||||
// Return the result along with the ID to identify successes/failures
|
||||
let copy_result = mgmt.copy_file(&file_id, target_folder.map(|s| s.to_string())).await;
|
||||
(file_id, copy_result)
|
||||
}
|
||||
});
|
||||
}))
|
||||
.buffer_unordered(max_concurrent);
|
||||
|
||||
// Execute all operations in parallel with concurrency control
|
||||
let operation_results = join_all(operations).await;
|
||||
|
||||
// Process the results
|
||||
for (file_id, operation_result) in operation_results {
|
||||
// Process results as they complete
|
||||
while let Some((file_id, operation_result)) = operation_stream.next().await {
|
||||
match operation_result {
|
||||
Ok(file) => {
|
||||
result.successful.push(file);
|
||||
@@ -195,6 +182,7 @@ impl BatchOperationService {
|
||||
) -> Result<BatchResult<FileDto>, BatchOperationError> {
|
||||
info!("Starting batch move of {} files", file_ids.len());
|
||||
let start_time = std::time::Instant::now();
|
||||
let max_concurrent = self.config.concurrency.max_concurrent_files;
|
||||
|
||||
// Create result structure
|
||||
let mut result = BatchResult {
|
||||
@@ -206,31 +194,20 @@ impl BatchOperationService {
|
||||
},
|
||||
};
|
||||
|
||||
// Define the operation to perform for each file
|
||||
let operations = file_ids.into_iter().map(|file_id| {
|
||||
let target_folder: Option<Arc<str>> = target_folder_id.map(|s| Arc::from(s.as_str()));
|
||||
|
||||
let mut operation_stream = stream::iter(file_ids.into_iter().map(|file_id| {
|
||||
let mgmt = self.file_management.clone();
|
||||
let target_folder = target_folder_id.clone();
|
||||
let semaphore = self.semaphore.clone();
|
||||
let target_folder = target_folder.clone();
|
||||
|
||||
async move {
|
||||
// Acquire semaphore permit
|
||||
let permit = semaphore.acquire().await.unwrap();
|
||||
|
||||
let move_result = mgmt.move_file(&file_id, target_folder.clone()).await;
|
||||
|
||||
// Release the permit explicitly
|
||||
drop(permit);
|
||||
|
||||
// Return the result along with the ID to identify successes/failures
|
||||
let move_result = mgmt.move_file(&file_id, target_folder.map(|s| s.to_string())).await;
|
||||
(file_id, move_result)
|
||||
}
|
||||
});
|
||||
}))
|
||||
.buffer_unordered(max_concurrent);
|
||||
|
||||
// Execute all operations in parallel with concurrency control
|
||||
let operation_results = join_all(operations).await;
|
||||
|
||||
// Process the results
|
||||
for (file_id, operation_result) in operation_results {
|
||||
while let Some((file_id, operation_result)) = operation_stream.next().await {
|
||||
match operation_result {
|
||||
Ok(file) => {
|
||||
result.successful.push(file);
|
||||
@@ -278,30 +255,19 @@ impl BatchOperationService {
|
||||
};
|
||||
|
||||
// Define the operation to perform for each file
|
||||
let operations = file_ids.into_iter().map(|file_id| {
|
||||
let mut operation_stream = stream::iter(file_ids.into_iter().map(|file_id| {
|
||||
let mgmt = self.file_management.clone();
|
||||
let semaphore = self.semaphore.clone();
|
||||
let id_clone = file_id.clone();
|
||||
|
||||
async move {
|
||||
// Acquire semaphore permit
|
||||
let permit = semaphore.acquire().await.unwrap();
|
||||
|
||||
let delete_result = mgmt.delete_file(&file_id).await;
|
||||
|
||||
// Release the permit explicitly
|
||||
drop(permit);
|
||||
|
||||
// Return the result along with the ID
|
||||
(id_clone.clone(), delete_result.map(|_| id_clone))
|
||||
let id_for_result = file_id.clone();
|
||||
(file_id, delete_result.map(|_| id_for_result))
|
||||
}
|
||||
});
|
||||
}))
|
||||
.buffer_unordered(self.config.concurrency.max_concurrent_files);
|
||||
|
||||
// Execute all operations in parallel with concurrency control
|
||||
let operation_results = join_all(operations).await;
|
||||
|
||||
// Process the results
|
||||
for (file_id, operation_result) in operation_results {
|
||||
// Process results as they complete
|
||||
while let Some((file_id, operation_result)) = operation_stream.next().await {
|
||||
match operation_result {
|
||||
Ok(id) => {
|
||||
result.successful.push(id);
|
||||
@@ -349,29 +315,18 @@ impl BatchOperationService {
|
||||
};
|
||||
|
||||
// Define the operation to perform for each file
|
||||
let operations = file_ids.into_iter().map(|file_id| {
|
||||
let mut operation_stream = stream::iter(file_ids.into_iter().map(|file_id| {
|
||||
let retrieval = self.file_retrieval.clone();
|
||||
let semaphore = self.semaphore.clone();
|
||||
|
||||
async move {
|
||||
// Acquire semaphore permit
|
||||
let permit = semaphore.acquire().await.unwrap();
|
||||
|
||||
let get_result = retrieval.get_file(&file_id).await;
|
||||
|
||||
// Release the permit explicitly
|
||||
drop(permit);
|
||||
|
||||
// Return the result along with the ID
|
||||
(file_id, get_result)
|
||||
}
|
||||
});
|
||||
}))
|
||||
.buffer_unordered(self.config.concurrency.max_concurrent_files);
|
||||
|
||||
// Execute all operations in parallel with concurrency control
|
||||
let operation_results = join_all(operations).await;
|
||||
|
||||
// Process the results
|
||||
for (file_id, operation_result) in operation_results {
|
||||
// Process results as they complete
|
||||
while let Some((file_id, operation_result)) = operation_stream.next().await {
|
||||
match operation_result {
|
||||
Ok(file) => {
|
||||
result.successful.push(file);
|
||||
@@ -421,31 +376,22 @@ impl BatchOperationService {
|
||||
};
|
||||
|
||||
// Define the operation to perform for each folder
|
||||
let operations = folder_ids.into_iter().map(|folder_id| {
|
||||
// Arc<str> avoids N heap-clones of the caller string
|
||||
let caller: Arc<str> = Arc::from(caller_id);
|
||||
|
||||
let mut operation_stream = stream::iter(folder_ids.into_iter().map(|folder_id| {
|
||||
let folder_service = self.folder_service.clone();
|
||||
let semaphore = self.semaphore.clone();
|
||||
let id_clone = folder_id.clone();
|
||||
let caller = caller_id.to_string();
|
||||
let caller = caller.clone();
|
||||
|
||||
async move {
|
||||
// Acquire semaphore permit
|
||||
let permit = semaphore.acquire().await.unwrap();
|
||||
|
||||
let delete_result = folder_service.delete_folder(&folder_id, &caller).await;
|
||||
|
||||
// Release the permit explicitly
|
||||
drop(permit);
|
||||
|
||||
// Return the result along with the ID
|
||||
(id_clone.clone(), delete_result.map(|_| id_clone))
|
||||
let id_for_result = folder_id.clone();
|
||||
(folder_id, delete_result.map(|_| id_for_result))
|
||||
}
|
||||
});
|
||||
}))
|
||||
.buffer_unordered(self.config.concurrency.max_concurrent_files);
|
||||
|
||||
// Execute all operations in parallel with concurrency control
|
||||
let operation_results = join_all(operations).await;
|
||||
|
||||
// Process the results
|
||||
for (folder_id, operation_result) in operation_results {
|
||||
while let Some((folder_id, operation_result)) = operation_stream.next().await {
|
||||
match operation_result {
|
||||
Ok(id) => {
|
||||
result.successful.push(id);
|
||||
@@ -497,23 +443,21 @@ impl BatchOperationService {
|
||||
},
|
||||
};
|
||||
|
||||
let operations = file_ids.into_iter().map(|file_id| {
|
||||
let uid: Arc<str> = Arc::from(user_id);
|
||||
|
||||
let mut operation_stream = stream::iter(file_ids.into_iter().map(|file_id| {
|
||||
let trash = trash_service.clone();
|
||||
let semaphore = self.semaphore.clone();
|
||||
let uid = user_id.to_string();
|
||||
let id_clone = file_id.clone();
|
||||
let uid = uid.clone();
|
||||
|
||||
async move {
|
||||
let permit = semaphore.acquire().await.unwrap();
|
||||
let trash_result = trash.move_to_trash(&file_id, "file", &uid).await;
|
||||
drop(permit);
|
||||
(id_clone.clone(), trash_result.map(|_| id_clone))
|
||||
let id_for_result = file_id.clone();
|
||||
(file_id, trash_result.map(|_| id_for_result))
|
||||
}
|
||||
});
|
||||
}))
|
||||
.buffer_unordered(self.config.concurrency.max_concurrent_files);
|
||||
|
||||
let operation_results = join_all(operations).await;
|
||||
|
||||
for (file_id, operation_result) in operation_results {
|
||||
while let Some((file_id, operation_result)) = operation_stream.next().await {
|
||||
match operation_result {
|
||||
Ok(id) => {
|
||||
result.successful.push(id);
|
||||
@@ -564,23 +508,21 @@ impl BatchOperationService {
|
||||
},
|
||||
};
|
||||
|
||||
let operations = folder_ids.into_iter().map(|folder_id| {
|
||||
let uid: Arc<str> = Arc::from(user_id);
|
||||
|
||||
let mut operation_stream = stream::iter(folder_ids.into_iter().map(|folder_id| {
|
||||
let trash = trash_service.clone();
|
||||
let semaphore = self.semaphore.clone();
|
||||
let uid = user_id.to_string();
|
||||
let id_clone = folder_id.clone();
|
||||
let uid = uid.clone();
|
||||
|
||||
async move {
|
||||
let permit = semaphore.acquire().await.unwrap();
|
||||
let trash_result = trash.move_to_trash(&folder_id, "folder", &uid).await;
|
||||
drop(permit);
|
||||
(id_clone.clone(), trash_result.map(|_| id_clone))
|
||||
let id_for_result = folder_id.clone();
|
||||
(folder_id, trash_result.map(|_| id_for_result))
|
||||
}
|
||||
});
|
||||
}))
|
||||
.buffer_unordered(self.config.concurrency.max_concurrent_files);
|
||||
|
||||
let operation_results = join_all(operations).await;
|
||||
|
||||
for (folder_id, operation_result) in operation_results {
|
||||
while let Some((folder_id, operation_result)) = operation_stream.next().await {
|
||||
match operation_result {
|
||||
Ok(id) => {
|
||||
result.successful.push(id);
|
||||
@@ -627,24 +569,23 @@ impl BatchOperationService {
|
||||
},
|
||||
};
|
||||
|
||||
let operations = folder_ids.into_iter().map(|folder_id| {
|
||||
let target: Option<Arc<str>> = target_folder_id.map(|s| Arc::from(s.as_str()));
|
||||
let caller: Arc<str> = Arc::from(caller_id);
|
||||
|
||||
let mut operation_stream = stream::iter(folder_ids.into_iter().map(|folder_id| {
|
||||
let folder_service = self.folder_service.clone();
|
||||
let target = target_folder_id.clone();
|
||||
let semaphore = self.semaphore.clone();
|
||||
let caller = caller_id.to_string();
|
||||
let target = target.clone();
|
||||
let caller = caller.clone();
|
||||
|
||||
async move {
|
||||
let permit = semaphore.acquire().await.unwrap();
|
||||
let dto = MoveFolderDto { parent_id: target };
|
||||
let dto = MoveFolderDto { parent_id: target.map(|s| s.to_string()) };
|
||||
let move_result = folder_service.move_folder(&folder_id, dto, &caller).await;
|
||||
drop(permit);
|
||||
(folder_id, move_result)
|
||||
}
|
||||
});
|
||||
}))
|
||||
.buffer_unordered(self.config.concurrency.max_concurrent_files);
|
||||
|
||||
let operation_results = join_all(operations).await;
|
||||
|
||||
for (folder_id, operation_result) in operation_results {
|
||||
while let Some((folder_id, operation_result)) = operation_stream.next().await {
|
||||
match operation_result {
|
||||
Ok(folder) => {
|
||||
result.successful.push(folder);
|
||||
@@ -873,7 +814,7 @@ impl BatchOperationService {
|
||||
) -> Result<BatchResult<T>, BatchOperationError>
|
||||
where
|
||||
T: Clone + Send + 'static + std::fmt::Debug,
|
||||
F: Fn(T, Arc<Semaphore>) -> Fut + Clone + Send + Sync + 'static,
|
||||
F: Fn(T) -> Fut + Clone + Send + Sync + 'static,
|
||||
Fut: Future<Output = Result<T, DomainError>> + Send + 'static,
|
||||
{
|
||||
info!(
|
||||
@@ -892,33 +833,25 @@ impl BatchOperationService {
|
||||
},
|
||||
};
|
||||
|
||||
// Convert each item to a task
|
||||
let tasks = items.iter().map(|item| {
|
||||
let item_clone = item.clone();
|
||||
// buffer_unordered materialises only max_concurrent futures at a time
|
||||
let mut operation_stream = stream::iter(items.into_iter().map(|item| {
|
||||
let op = operation.clone();
|
||||
let semaphore = self.semaphore.clone();
|
||||
|
||||
async move {
|
||||
// The provided function must handle semaphore acquisition
|
||||
let op_result = op(item_clone.clone(), semaphore).await;
|
||||
|
||||
// Return the result along with the original item for identification
|
||||
(item_clone, op_result)
|
||||
let op_result = op(item.clone()).await;
|
||||
(item, op_result)
|
||||
}
|
||||
});
|
||||
}))
|
||||
.buffer_unordered(self.config.concurrency.max_concurrent_files);
|
||||
|
||||
// Execute all tasks in parallel
|
||||
let operation_results = join_all(tasks).await;
|
||||
|
||||
// Process results
|
||||
for (item, operation_result) in operation_results {
|
||||
// Process results as they complete
|
||||
while let Some((item, operation_result)) = operation_stream.next().await {
|
||||
match operation_result {
|
||||
Ok(result_item) => {
|
||||
result.successful.push(result_item);
|
||||
result.stats.successful += 1;
|
||||
}
|
||||
Err(e) => {
|
||||
// Convert item to string for error reporting
|
||||
result.failed.push((format!("{:?}", item), e.to_string()));
|
||||
result.stats.failed += 1;
|
||||
}
|
||||
@@ -960,34 +893,23 @@ impl BatchOperationService {
|
||||
};
|
||||
|
||||
// Define the operation for each folder
|
||||
let operations = folders.into_iter().map(|(name, parent_id)| {
|
||||
let mut operation_stream = stream::iter(folders.into_iter().map(|(name, parent_id)| {
|
||||
let folder_service = self.folder_service.clone();
|
||||
let semaphore = self.semaphore.clone();
|
||||
|
||||
async move {
|
||||
// Acquire semaphore permit
|
||||
let permit = semaphore.acquire().await.unwrap();
|
||||
|
||||
let dto = crate::application::dtos::folder_dto::CreateFolderDto {
|
||||
name: name.clone(),
|
||||
parent_id: parent_id.clone(),
|
||||
};
|
||||
let create_result = folder_service.create_folder(dto).await;
|
||||
|
||||
// Release the permit explicitly
|
||||
drop(permit);
|
||||
|
||||
// Return the result with an identifier for errors
|
||||
let id = format!("{}:{}", name, parent_id.unwrap_or_default());
|
||||
(id, create_result)
|
||||
}
|
||||
});
|
||||
}))
|
||||
.buffer_unordered(self.config.concurrency.max_concurrent_files);
|
||||
|
||||
// Execute all operations in parallel
|
||||
let operation_results = join_all(operations).await;
|
||||
|
||||
// Process the results
|
||||
for (id, operation_result) in operation_results {
|
||||
// Process results as they complete
|
||||
while let Some((id, operation_result)) = operation_stream.next().await {
|
||||
match operation_result {
|
||||
Ok(folder) => {
|
||||
result.successful.push(folder);
|
||||
@@ -1035,29 +957,18 @@ impl BatchOperationService {
|
||||
};
|
||||
|
||||
// Define the operation for each folder
|
||||
let operations = folder_ids.into_iter().map(|folder_id| {
|
||||
let mut operation_stream = stream::iter(folder_ids.into_iter().map(|folder_id| {
|
||||
let folder_service = self.folder_service.clone();
|
||||
let semaphore = self.semaphore.clone();
|
||||
|
||||
async move {
|
||||
// Acquire semaphore permit
|
||||
let permit = semaphore.acquire().await.unwrap();
|
||||
|
||||
let get_result = folder_service.get_folder(&folder_id).await;
|
||||
|
||||
// Release the permit explicitly
|
||||
drop(permit);
|
||||
|
||||
// Return the result with its ID
|
||||
(folder_id, get_result)
|
||||
}
|
||||
});
|
||||
}))
|
||||
.buffer_unordered(self.config.concurrency.max_concurrent_files);
|
||||
|
||||
// Execute all operations in parallel
|
||||
let operation_results = join_all(operations).await;
|
||||
|
||||
// Process the results
|
||||
for (folder_id, operation_result) in operation_results {
|
||||
// Process results as they complete
|
||||
while let Some((folder_id, operation_result)) = operation_stream.next().await {
|
||||
match operation_result {
|
||||
Ok(folder) => {
|
||||
result.successful.push(folder);
|
||||
@@ -1105,11 +1016,8 @@ mod tests {
|
||||
AppConfig::default(),
|
||||
);
|
||||
|
||||
// Define a generic test operation
|
||||
let operation = |item: i32, semaphore: Arc<Semaphore>| async move {
|
||||
// Acquire and release the semaphore
|
||||
let _permit = semaphore.acquire().await.unwrap();
|
||||
|
||||
// Define a generic test operation (no more semaphore parameter)
|
||||
let operation = |item: i32| async move {
|
||||
if item % 2 == 0 {
|
||||
// Simulate success for even numbers
|
||||
Ok(item * 2)
|
||||
|
||||
@@ -3,7 +3,6 @@ use crate::application::dtos::folder_dto::{
|
||||
};
|
||||
use crate::application::ports::inbound::FolderUseCase;
|
||||
use crate::application::ports::outbound::FolderStoragePort;
|
||||
use crate::application::transactions::storage_transaction::StorageTransaction;
|
||||
use crate::common::errors::{DomainError, ErrorKind};
|
||||
use crate::domain::services::path_service::StoragePath;
|
||||
use async_trait::async_trait;
|
||||
@@ -398,52 +397,11 @@ impl FolderUseCase for FolderService {
|
||||
return Err(DomainError::not_found("Folder", id));
|
||||
}
|
||||
|
||||
// Create transaction for renaming
|
||||
let mut transaction = StorageTransaction::new("rename_folder");
|
||||
|
||||
// Main operation: rename folder
|
||||
// Clone all values to avoid lifetime issues
|
||||
let folder_storage = self.folder_storage.clone();
|
||||
let id_owned = id.to_string();
|
||||
let name_owned = dto.name.clone();
|
||||
|
||||
// Create future with owned values
|
||||
let rename_op = async move {
|
||||
folder_storage.rename_folder(&id_owned, name_owned).await?;
|
||||
Ok(())
|
||||
};
|
||||
let rollback_op = {
|
||||
let original_name = existing_folder.name().to_string();
|
||||
let storage = self.folder_storage.clone();
|
||||
let id_clone = id.to_string();
|
||||
|
||||
async move {
|
||||
// In case of failure, restore the original name
|
||||
storage
|
||||
.rename_folder(&id_clone, original_name)
|
||||
.await
|
||||
.map(|_| ())
|
||||
.map_err(|e| {
|
||||
DomainError::new(
|
||||
ErrorKind::InternalError,
|
||||
"Folder",
|
||||
format!("Failed to rollback folder rename: {}", e),
|
||||
)
|
||||
})
|
||||
}
|
||||
};
|
||||
|
||||
// Add to the transaction
|
||||
transaction.add_operation(rename_op, rollback_op);
|
||||
|
||||
// Execute transaction
|
||||
transaction.commit().await?;
|
||||
|
||||
// Get the renamed folder
|
||||
let folder = self.folder_storage.get_folder(id).await.map_err(|e| {
|
||||
// Rename folder — UPDATE RETURNING gives us the updated row directly
|
||||
let folder = self.folder_storage.rename_folder(id, dto.name).await.map_err(|e| {
|
||||
DomainError::internal_error(
|
||||
"FolderStorage",
|
||||
format!("Failed to get renamed folder with ID: {}: {}", id, e),
|
||||
format!("Failed to rename folder with ID: {}: {}", id, e),
|
||||
)
|
||||
})?;
|
||||
|
||||
@@ -495,55 +453,12 @@ impl FolderUseCase for FolderService {
|
||||
// TODO: Ideally we should verify the entire hierarchy to prevent cycles
|
||||
}
|
||||
|
||||
// Create transaction for moving
|
||||
let mut transaction = StorageTransaction::new("move_folder");
|
||||
|
||||
// Main operation: move folder
|
||||
// Clone all values to avoid lifetime issues
|
||||
let folder_storage = self.folder_storage.clone();
|
||||
let id_owned = id.to_string();
|
||||
// Get parent ID as owned string or None
|
||||
let parent_id_owned = dto.parent_id.as_ref().map(|p| p.to_string());
|
||||
|
||||
// Create future with owned values
|
||||
let move_op = async move {
|
||||
// Convert Option<String> to Option<&str>
|
||||
let parent_ref = parent_id_owned.as_deref();
|
||||
folder_storage.move_folder(&id_owned, parent_ref).await?;
|
||||
Ok(())
|
||||
};
|
||||
let rollback_op = {
|
||||
let original_parent_id = source_folder.parent_id().map(String::from);
|
||||
let storage = self.folder_storage.clone();
|
||||
let id_clone = id.to_string();
|
||||
|
||||
async move {
|
||||
// In case of failure, restore the original location
|
||||
storage
|
||||
.move_folder(&id_clone, original_parent_id.as_deref())
|
||||
.await
|
||||
.map(|_| ())
|
||||
.map_err(|e| {
|
||||
DomainError::new(
|
||||
ErrorKind::InternalError,
|
||||
"Folder",
|
||||
format!("Failed to rollback folder move: {}", e),
|
||||
)
|
||||
})
|
||||
}
|
||||
};
|
||||
|
||||
// Add to the transaction
|
||||
transaction.add_operation(move_op, rollback_op);
|
||||
|
||||
// Execute transaction
|
||||
transaction.commit().await?;
|
||||
|
||||
// Get the moved folder
|
||||
let folder = self.folder_storage.get_folder(id).await.map_err(|e| {
|
||||
// Move folder — UPDATE RETURNING gives us the updated row directly
|
||||
let parent_ref = dto.parent_id.as_deref();
|
||||
let folder = self.folder_storage.move_folder(id, parent_ref).await.map_err(|e| {
|
||||
DomainError::internal_error(
|
||||
"FolderStorage",
|
||||
format!("Failed to get moved folder with ID: {}: {}", id, e),
|
||||
format!("Failed to move folder with ID: {}: {}", id, e),
|
||||
)
|
||||
})?;
|
||||
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
use async_trait::async_trait;
|
||||
use std::cmp::Reverse;
|
||||
use std::sync::Arc;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
@@ -323,15 +324,13 @@ impl SearchUseCase for SearchService {
|
||||
.map(|f| Self::enrich_folder(f, query))
|
||||
.collect();
|
||||
|
||||
// Sort folders
|
||||
// Sort folders (cached_key avoids O(N log N) temporary String allocations)
|
||||
match criteria.sort_by.as_str() {
|
||||
"name" => {
|
||||
enriched_folders
|
||||
.sort_by(|a, b| a.name.to_lowercase().cmp(&b.name.to_lowercase()));
|
||||
enriched_folders.sort_by_cached_key(|f| f.name.to_lowercase());
|
||||
}
|
||||
"name_desc" => {
|
||||
enriched_folders
|
||||
.sort_by(|a, b| b.name.to_lowercase().cmp(&a.name.to_lowercase()));
|
||||
enriched_folders.sort_by_cached_key(|f| Reverse(f.name.to_lowercase()));
|
||||
}
|
||||
"date" => {
|
||||
enriched_folders.sort_by(|a, b| a.modified_at.cmp(&b.modified_at));
|
||||
@@ -411,13 +410,13 @@ impl SearchUseCase for SearchService {
|
||||
.map(|f| Self::enrich_folder(f, query))
|
||||
.collect();
|
||||
|
||||
// ── Sort folders (files already sorted by SQL ORDER BY) ──
|
||||
// ── Sort folders (cached_key avoids O(N log N) temporary String allocations) ──
|
||||
match criteria.sort_by.as_str() {
|
||||
"name" => {
|
||||
enriched_folders.sort_by(|a, b| a.name.to_lowercase().cmp(&b.name.to_lowercase()));
|
||||
enriched_folders.sort_by_cached_key(|f| f.name.to_lowercase());
|
||||
}
|
||||
"name_desc" => {
|
||||
enriched_folders.sort_by(|a, b| b.name.to_lowercase().cmp(&a.name.to_lowercase()));
|
||||
enriched_folders.sort_by_cached_key(|f| Reverse(f.name.to_lowercase()));
|
||||
}
|
||||
"date" => {
|
||||
enriched_folders.sort_by(|a, b| a.modified_at.cmp(&b.modified_at));
|
||||
|
||||
@@ -1 +0,0 @@
|
||||
pub mod storage_transaction;
|
||||
@@ -1,150 +0,0 @@
|
||||
use crate::common::errors::{DomainError, ErrorKind};
|
||||
use std::future::Future;
|
||||
use std::pin::Pin;
|
||||
|
||||
/// Type for async operations and rollbacks
|
||||
type TransactionOp = Pin<Box<dyn Future<Output = Result<(), DomainError>> + Send>>;
|
||||
|
||||
/// Transaction for storage operations
|
||||
/// Allows defining a set of operations and their corresponding rollbacks
|
||||
pub struct StorageTransaction {
|
||||
/// Operations to execute
|
||||
operations: Vec<Box<dyn FnOnce() -> TransactionOp + Send>>,
|
||||
/// Rollback operations to revert changes in case of error
|
||||
rollbacks: Vec<Box<dyn FnOnce() -> TransactionOp + Send>>,
|
||||
/// Transaction name for logging
|
||||
name: String,
|
||||
}
|
||||
|
||||
impl StorageTransaction {
|
||||
/// Creates a new transaction
|
||||
pub fn new(name: &str) -> Self {
|
||||
Self {
|
||||
operations: Vec::new(),
|
||||
rollbacks: Vec::new(),
|
||||
name: name.to_string(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Adds an operation to the transaction with its corresponding rollback
|
||||
pub fn add_operation<F, R>(&mut self, operation: F, rollback: R)
|
||||
where
|
||||
F: Future<Output = Result<(), DomainError>> + Send + 'static,
|
||||
R: Future<Output = Result<(), DomainError>> + Send + 'static,
|
||||
{
|
||||
self.operations.push(Box::new(move || Box::pin(operation)));
|
||||
self.rollbacks.push(Box::new(move || Box::pin(rollback)));
|
||||
}
|
||||
|
||||
/// Adds an operation without rollback (for cleanup or logging)
|
||||
pub fn add_finalizer<F>(&mut self, finalizer: F)
|
||||
where
|
||||
F: Future<Output = Result<(), DomainError>> + Send + 'static,
|
||||
{
|
||||
// The rollback is a no-op
|
||||
let noop = async { Ok(()) };
|
||||
|
||||
self.operations.push(Box::new(move || Box::pin(finalizer)));
|
||||
self.rollbacks.push(Box::new(move || Box::pin(noop)));
|
||||
}
|
||||
|
||||
/// Executes the transaction by applying all operations in order
|
||||
/// If any fails, executes rollbacks in reverse order
|
||||
pub async fn commit(mut self) -> Result<(), DomainError> {
|
||||
tracing::debug!("Starting transaction: {}", self.name);
|
||||
|
||||
let mut completed_ops = Vec::new();
|
||||
|
||||
// Extract operations to avoid ownership issues
|
||||
let operations = std::mem::take(&mut self.operations);
|
||||
let transaction_name = self.name.clone();
|
||||
|
||||
// Execute operations
|
||||
for (i, op) in operations.into_iter().enumerate() {
|
||||
match op().await {
|
||||
Ok(()) => {
|
||||
completed_ops.push(i);
|
||||
tracing::trace!(
|
||||
"Operation {} completed in transaction: {}",
|
||||
i,
|
||||
transaction_name
|
||||
);
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::error!(
|
||||
"Error in operation {} of transaction {}: {}",
|
||||
i,
|
||||
transaction_name,
|
||||
e
|
||||
);
|
||||
|
||||
// Execute rollbacks for completed operations in reverse order
|
||||
self.rollback(completed_ops).await?;
|
||||
|
||||
return Err(DomainError::new(
|
||||
ErrorKind::InternalError,
|
||||
"Transaction",
|
||||
format!("Transaction '{}' failed: {}", transaction_name, e),
|
||||
)
|
||||
.with_source(e));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
tracing::debug!("Transaction completed successfully: {}", transaction_name);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Executes rollbacks for completed operations
|
||||
async fn rollback(mut self, completed_ops: Vec<usize>) -> Result<(), DomainError> {
|
||||
tracing::warn!("Starting rollback for transaction: {}", self.name);
|
||||
|
||||
let mut rollback_errors = Vec::new();
|
||||
|
||||
// Extract rollbacks to avoid ownership issues
|
||||
let mut rollbacks = Vec::new();
|
||||
std::mem::swap(&mut rollbacks, &mut self.rollbacks);
|
||||
|
||||
// Execute rollbacks in reverse order
|
||||
for i in completed_ops.into_iter().rev() {
|
||||
if i < rollbacks.len() {
|
||||
// Take ownership of the rollback (get a mutable reference)
|
||||
if let Some(rb) = rollbacks.get_mut(i) {
|
||||
// Swap with an empty function
|
||||
let rollback = std::mem::replace(rb, Box::new(|| Box::pin(async { Ok(()) })));
|
||||
if let Err(e) = rollback().await {
|
||||
tracing::error!(
|
||||
"Error in rollback of operation {} in transaction {}: {}",
|
||||
i,
|
||||
self.name,
|
||||
e
|
||||
);
|
||||
rollback_errors.push(e);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// If there were errors during rollback, report them
|
||||
if !rollback_errors.is_empty() {
|
||||
tracing::error!(
|
||||
"Errors during transaction rollback {}: {} errors",
|
||||
self.name,
|
||||
rollback_errors.len()
|
||||
);
|
||||
|
||||
return Err(DomainError::new(
|
||||
ErrorKind::InternalError,
|
||||
"Transaction",
|
||||
format!(
|
||||
"Errors during transaction '{}' rollback: {} errors",
|
||||
self.name,
|
||||
rollback_errors.len()
|
||||
),
|
||||
));
|
||||
}
|
||||
|
||||
tracing::info!("Transaction rollback completed: {}", self.name);
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
@@ -258,6 +258,9 @@ impl FolderRepository for FolderDbRepository {
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// Paginated folder listing — single query with `COUNT(*) OVER()` window
|
||||
/// function so the total matching count comes back alongside the data rows,
|
||||
/// eliminating a separate COUNT round-trip.
|
||||
async fn list_folders_paginated(
|
||||
&self,
|
||||
parent_id: Option<&str>,
|
||||
@@ -265,34 +268,14 @@ impl FolderRepository for FolderDbRepository {
|
||||
limit: usize,
|
||||
include_total: bool,
|
||||
) -> Result<(Vec<Folder>, Option<usize>), DomainError> {
|
||||
let total = if include_total {
|
||||
let count: i64 = if let Some(pid) = parent_id {
|
||||
sqlx::query_scalar(
|
||||
"SELECT COUNT(*) FROM storage.folders WHERE parent_id = $1::uuid AND NOT is_trashed",
|
||||
)
|
||||
.bind(pid)
|
||||
.fetch_one(self.pool())
|
||||
.await
|
||||
} else {
|
||||
sqlx::query_scalar(
|
||||
"SELECT COUNT(*) FROM storage.folders WHERE parent_id IS NULL AND NOT is_trashed",
|
||||
)
|
||||
.fetch_one(self.pool())
|
||||
.await
|
||||
}
|
||||
.map_err(|e| DomainError::internal_error("FolderDb", format!("count: {e}")))?;
|
||||
Some(count as usize)
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
let rows: Vec<(String, String, String, Option<String>, String, i64, i64)> =
|
||||
let rows: Vec<(String, String, String, Option<String>, String, i64, i64, i64)> =
|
||||
if let Some(pid) = parent_id {
|
||||
sqlx::query_as(
|
||||
r#"
|
||||
SELECT id::text, name, path, parent_id::text, user_id,
|
||||
EXTRACT(EPOCH FROM created_at)::bigint,
|
||||
EXTRACT(EPOCH FROM updated_at)::bigint
|
||||
EXTRACT(EPOCH FROM updated_at)::bigint,
|
||||
COUNT(*) OVER() AS total_count
|
||||
FROM storage.folders
|
||||
WHERE parent_id = $1::uuid AND NOT is_trashed
|
||||
ORDER BY name
|
||||
@@ -309,7 +292,8 @@ impl FolderRepository for FolderDbRepository {
|
||||
r#"
|
||||
SELECT id::text, name, path, parent_id::text, user_id,
|
||||
EXTRACT(EPOCH FROM created_at)::bigint,
|
||||
EXTRACT(EPOCH FROM updated_at)::bigint
|
||||
EXTRACT(EPOCH FROM updated_at)::bigint,
|
||||
COUNT(*) OVER() AS total_count
|
||||
FROM storage.folders
|
||||
WHERE parent_id IS NULL AND NOT is_trashed
|
||||
ORDER BY name
|
||||
@@ -323,15 +307,24 @@ impl FolderRepository for FolderDbRepository {
|
||||
}
|
||||
.map_err(|e| DomainError::internal_error("FolderDb", format!("paginate: {e}")))?;
|
||||
|
||||
// total_count is identical in every row; 0 when the result set is empty.
|
||||
let total = if include_total {
|
||||
Some(rows.first().map_or(0, |r| r.7) as usize)
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
let folders: Result<Vec<Folder>, DomainError> = rows
|
||||
.into_iter()
|
||||
.map(|(id, name, path, pid, uid, ca, ma)| {
|
||||
.map(|(id, name, path, pid, uid, ca, ma, _total)| {
|
||||
Self::row_to_folder(id, name, path, pid, Some(uid), ca, ma)
|
||||
})
|
||||
.collect();
|
||||
Ok((folders?, total))
|
||||
}
|
||||
|
||||
/// Paginated folder listing filtered by owner — single query with
|
||||
/// `COUNT(*) OVER()` to avoid a separate COUNT round-trip.
|
||||
async fn list_folders_by_owner_paginated(
|
||||
&self,
|
||||
parent_id: Option<&str>,
|
||||
@@ -340,36 +333,14 @@ impl FolderRepository for FolderDbRepository {
|
||||
limit: usize,
|
||||
include_total: bool,
|
||||
) -> Result<(Vec<Folder>, Option<usize>), DomainError> {
|
||||
let total = if include_total {
|
||||
let count: i64 = if let Some(pid) = parent_id {
|
||||
sqlx::query_scalar(
|
||||
"SELECT COUNT(*) FROM storage.folders WHERE parent_id = $1::uuid AND user_id = $2 AND NOT is_trashed",
|
||||
)
|
||||
.bind(pid)
|
||||
.bind(owner_id)
|
||||
.fetch_one(self.pool())
|
||||
.await
|
||||
} else {
|
||||
sqlx::query_scalar(
|
||||
"SELECT COUNT(*) FROM storage.folders WHERE parent_id IS NULL AND user_id = $1 AND NOT is_trashed",
|
||||
)
|
||||
.bind(owner_id)
|
||||
.fetch_one(self.pool())
|
||||
.await
|
||||
}
|
||||
.map_err(|e| DomainError::internal_error("FolderDb", format!("count_by_owner: {e}")))?;
|
||||
Some(count as usize)
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
let rows: Vec<(String, String, String, Option<String>, String, i64, i64)> =
|
||||
let rows: Vec<(String, String, String, Option<String>, String, i64, i64, i64)> =
|
||||
if let Some(pid) = parent_id {
|
||||
sqlx::query_as(
|
||||
r#"
|
||||
SELECT id::text, name, path, parent_id::text, user_id,
|
||||
EXTRACT(EPOCH FROM created_at)::bigint,
|
||||
EXTRACT(EPOCH FROM updated_at)::bigint
|
||||
EXTRACT(EPOCH FROM updated_at)::bigint,
|
||||
COUNT(*) OVER() AS total_count
|
||||
FROM storage.folders
|
||||
WHERE parent_id = $1::uuid AND user_id = $2 AND NOT is_trashed
|
||||
ORDER BY name
|
||||
@@ -387,7 +358,8 @@ impl FolderRepository for FolderDbRepository {
|
||||
r#"
|
||||
SELECT id::text, name, path, parent_id::text, user_id,
|
||||
EXTRACT(EPOCH FROM created_at)::bigint,
|
||||
EXTRACT(EPOCH FROM updated_at)::bigint
|
||||
EXTRACT(EPOCH FROM updated_at)::bigint,
|
||||
COUNT(*) OVER() AS total_count
|
||||
FROM storage.folders
|
||||
WHERE parent_id IS NULL AND user_id = $1 AND NOT is_trashed
|
||||
ORDER BY name
|
||||
@@ -404,9 +376,15 @@ impl FolderRepository for FolderDbRepository {
|
||||
DomainError::internal_error("FolderDb", format!("paginate_by_owner: {e}"))
|
||||
})?;
|
||||
|
||||
let total = if include_total {
|
||||
Some(rows.first().map_or(0, |r| r.7) as usize)
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
let folders: Result<Vec<Folder>, DomainError> = rows
|
||||
.into_iter()
|
||||
.map(|(id, name, path, pid, uid, ca, ma)| {
|
||||
.map(|(id, name, path, pid, uid, ca, ma, _total)| {
|
||||
Self::row_to_folder(id, name, path, pid, Some(uid), ca, ma)
|
||||
})
|
||||
.collect();
|
||||
@@ -417,16 +395,19 @@ impl FolderRepository for FolderDbRepository {
|
||||
// The BEFORE UPDATE trigger recomputes path/lpath for this row;
|
||||
// the AFTER UPDATE cascade trigger then batch-updates all
|
||||
// descendants in a single UPDATE using the GiST lpath index.
|
||||
sqlx::query(
|
||||
let row = sqlx::query_as::<_, (String, String, String, Option<String>, String, i64, i64)>(
|
||||
r#"
|
||||
UPDATE storage.folders
|
||||
SET name = $1, updated_at = NOW()
|
||||
WHERE id = $2::uuid AND NOT is_trashed
|
||||
RETURNING id::text, name, path, parent_id::text, user_id,
|
||||
EXTRACT(EPOCH FROM created_at)::bigint,
|
||||
EXTRACT(EPOCH FROM updated_at)::bigint
|
||||
"#,
|
||||
)
|
||||
.bind(&new_name)
|
||||
.bind(id)
|
||||
.execute(self.pool())
|
||||
.fetch_optional(self.pool())
|
||||
.await
|
||||
.map_err(|e| {
|
||||
if let sqlx::Error::Database(ref db_err) = e
|
||||
@@ -435,9 +416,10 @@ impl FolderRepository for FolderDbRepository {
|
||||
return DomainError::already_exists("Folder", format!("{new_name} already exists"));
|
||||
}
|
||||
DomainError::internal_error("FolderDb", format!("rename: {e}"))
|
||||
})?;
|
||||
})?
|
||||
.ok_or_else(|| DomainError::not_found("Folder", id))?;
|
||||
|
||||
self.get_folder(id).await
|
||||
Self::row_to_folder(row.0, row.1, row.2, row.3, Some(row.4), row.5, row.6)
|
||||
}
|
||||
|
||||
async fn move_folder(
|
||||
@@ -448,20 +430,24 @@ impl FolderRepository for FolderDbRepository {
|
||||
// The BEFORE UPDATE trigger recomputes path/lpath for this row;
|
||||
// the AFTER UPDATE cascade trigger then batch-updates all
|
||||
// descendants in a single UPDATE using the GiST lpath index.
|
||||
sqlx::query(
|
||||
let row = sqlx::query_as::<_, (String, String, String, Option<String>, String, i64, i64)>(
|
||||
r#"
|
||||
UPDATE storage.folders
|
||||
SET parent_id = $1::uuid, updated_at = NOW()
|
||||
WHERE id = $2::uuid AND NOT is_trashed
|
||||
RETURNING id::text, name, path, parent_id::text, user_id,
|
||||
EXTRACT(EPOCH FROM created_at)::bigint,
|
||||
EXTRACT(EPOCH FROM updated_at)::bigint
|
||||
"#,
|
||||
)
|
||||
.bind(new_parent_id)
|
||||
.bind(id)
|
||||
.execute(self.pool())
|
||||
.fetch_optional(self.pool())
|
||||
.await
|
||||
.map_err(|e| DomainError::internal_error("FolderDb", format!("move: {e}")))?;
|
||||
.map_err(|e| DomainError::internal_error("FolderDb", format!("move: {e}")))?
|
||||
.ok_or_else(|| DomainError::not_found("Folder", id))?;
|
||||
|
||||
self.get_folder(id).await
|
||||
Self::row_to_folder(row.0, row.1, row.2, row.3, Some(row.4), row.5, row.6)
|
||||
}
|
||||
|
||||
async fn delete_folder(&self, id: &str) -> Result<(), DomainError> {
|
||||
|
||||
@@ -26,16 +26,28 @@ use crate::domain::errors::{DomainError, ErrorKind};
|
||||
/// Maximum file size for transcoding (5MB - larger files stream directly)
|
||||
pub const MAX_TRANSCODE_SIZE: u64 = 5 * 1024 * 1024;
|
||||
|
||||
/// Number of threads in the dedicated transcoding pool
|
||||
const TRANSCODE_POOL_THREADS: usize = 2;
|
||||
/// Minimum number of threads in the dedicated transcoding pool
|
||||
const MIN_TRANSCODE_THREADS: usize = 2;
|
||||
|
||||
/// Compute the number of transcoding threads: half the available CPUs,
|
||||
/// with a floor of `MIN_TRANSCODE_THREADS`. `available_parallelism()`
|
||||
/// respects cgroup limits (Docker/K8s) and CPU affinity masks.
|
||||
fn transcode_thread_count() -> usize {
|
||||
let cpus = std::thread::available_parallelism()
|
||||
.map(|n| n.get())
|
||||
.unwrap_or(MIN_TRANSCODE_THREADS);
|
||||
(cpus / 2).max(MIN_TRANSCODE_THREADS)
|
||||
}
|
||||
|
||||
/// Dedicated rayon thread pool for CPU-bound image transcoding.
|
||||
/// Isolated from Tokio's blocking pool to prevent starvation of other I/O.
|
||||
/// Thread count scales with available CPUs (half cores, min 2).
|
||||
fn transcode_pool() -> &'static rayon::ThreadPool {
|
||||
static POOL: OnceLock<rayon::ThreadPool> = OnceLock::new();
|
||||
POOL.get_or_init(|| {
|
||||
let threads = transcode_thread_count();
|
||||
rayon::ThreadPoolBuilder::new()
|
||||
.num_threads(TRANSCODE_POOL_THREADS)
|
||||
.num_threads(threads)
|
||||
.thread_name(|idx| format!("transcode-{idx}"))
|
||||
.build()
|
||||
.expect("Failed to create transcode thread pool")
|
||||
@@ -171,7 +183,7 @@ impl ImageTranscodeService {
|
||||
fs::create_dir_all(self.cache_dir.join("webp")).await?;
|
||||
tracing::info!(
|
||||
"🖼️ Image transcode service initialized (rayon pool: {} threads, cache dir: {:?})",
|
||||
TRANSCODE_POOL_THREADS,
|
||||
transcode_thread_count(),
|
||||
self.cache_dir
|
||||
);
|
||||
Ok(())
|
||||
|
||||
Reference in New Issue
Block a user