adding comments, moving technical documentation, improve Dockerfile, delete unuseful files
This commit is contained in:
@@ -1,6 +1,6 @@
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
use crate::domain::entities::share::{Share, ShareItemType, SharePermissions};
|
||||
use crate::domain::entities::share::{Share, SharePermissions};
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct ShareDto {
|
||||
|
||||
@@ -6,18 +6,18 @@ use crate::application::ports::file_ports::FileManagementUseCase;
|
||||
use crate::application::ports::storage_ports::FileWritePort;
|
||||
use crate::common::errors::DomainError;
|
||||
|
||||
/// Servicio para operaciones de gestión de archivos
|
||||
/// Service for file management operations
|
||||
pub struct FileManagementService {
|
||||
file_repository: Arc<dyn FileWritePort>,
|
||||
}
|
||||
|
||||
impl FileManagementService {
|
||||
/// Crea un nuevo servicio de gestión de archivos
|
||||
/// Creates a new file management service
|
||||
pub fn new(file_repository: Arc<dyn FileWritePort>) -> Self {
|
||||
Self { file_repository }
|
||||
}
|
||||
|
||||
/// Crea un stub para pruebas
|
||||
/// Creates a stub for testing
|
||||
pub fn default_stub() -> Self {
|
||||
Self {
|
||||
file_repository: Arc::new(crate::infrastructure::repositories::FileFsWriteRepository::default_stub())
|
||||
@@ -28,15 +28,15 @@ impl FileManagementService {
|
||||
#[async_trait]
|
||||
impl FileManagementUseCase for FileManagementService {
|
||||
async fn move_file(&self, file_id: &str, folder_id: Option<String>) -> Result<FileDto, DomainError> {
|
||||
tracing::info!("Moviendo archivo con ID: {} a carpeta: {:?}", file_id, folder_id);
|
||||
tracing::info!("Moving file with ID: {} to folder: {:?}", file_id, folder_id);
|
||||
|
||||
let moved_file = self.file_repository.move_file(file_id, folder_id).await
|
||||
.map_err(|e| {
|
||||
tracing::error!("Error al mover archivo (ID: {}): {}", file_id, e);
|
||||
tracing::error!("Error moving file (ID: {}): {}", file_id, e);
|
||||
e
|
||||
})?;
|
||||
|
||||
tracing::info!("Archivo movido exitosamente: {} (ID: {}) a carpeta: {:?}",
|
||||
tracing::info!("File moved successfully: {} (ID: {}) to folder: {:?}",
|
||||
moved_file.name(), moved_file.id(), moved_file.folder_id());
|
||||
|
||||
Ok(FileDto::from(moved_file))
|
||||
|
||||
@@ -19,23 +19,23 @@ use bytes::Bytes;
|
||||
#[derive(Debug, Error)]
|
||||
pub enum FileServiceError {
|
||||
/// Returned when a requested file cannot be found
|
||||
#[error("Archivo no encontrado: {0}")]
|
||||
#[error("File not found: {0}")]
|
||||
NotFound(String),
|
||||
|
||||
/// Returned when a file operation conflicts with existing files
|
||||
#[error("Archivo ya existe: {0}")]
|
||||
#[error("File already exists: {0}")]
|
||||
Conflict(String),
|
||||
|
||||
/// Returned when file access fails due to permissions or I/O issues
|
||||
#[error("Error de acceso al archivo: {0}")]
|
||||
#[error("File access error: {0}")]
|
||||
AccessError(String),
|
||||
|
||||
/// Returned when a file path is invalid
|
||||
#[error("Ruta de archivo inválida: {0}")]
|
||||
#[error("Invalid file path: {0}")]
|
||||
InvalidPath(String),
|
||||
|
||||
/// Generic internal error for unexpected failures
|
||||
#[error("Error interno: {0}")]
|
||||
#[error("Internal error: {0}")]
|
||||
InternalError(String),
|
||||
}
|
||||
|
||||
@@ -52,7 +52,7 @@ impl From<FileRepositoryError> for FileServiceError {
|
||||
FileRepositoryError::AlreadyExists(path) => FileServiceError::Conflict(path),
|
||||
FileRepositoryError::InvalidPath(path) => FileServiceError::InvalidPath(path),
|
||||
FileRepositoryError::IoError(e) => FileServiceError::AccessError(e.to_string()),
|
||||
FileRepositoryError::Timeout(msg) => FileServiceError::AccessError(format!("Operación expiró: {}", msg)),
|
||||
FileRepositoryError::Timeout(msg) => FileServiceError::AccessError(format!("Operation timed out: {}", msg)),
|
||||
_ => FileServiceError::InternalError(err.to_string()),
|
||||
}
|
||||
}
|
||||
@@ -213,16 +213,16 @@ impl FileService {
|
||||
|
||||
/// Moves a file to a new folder using filesystem operations directly
|
||||
pub async fn move_file(&self, file_id: &str, folder_id: Option<String>) -> FileServiceResult<FileDto> {
|
||||
tracing::info!("Moviendo archivo con ID: {} a carpeta: {:?}", file_id, folder_id);
|
||||
tracing::info!("Moving file with ID: {} to folder: {:?}", file_id, folder_id);
|
||||
|
||||
// Usar la implementación eficiente del repositorio que utiliza rename
|
||||
// Use the efficient repository implementation that uses rename
|
||||
let moved_file = self.file_repository.move_file(file_id, folder_id).await
|
||||
.map_err(|e| {
|
||||
tracing::error!("Error al mover archivo (ID: {}): {}", file_id, e);
|
||||
tracing::error!("Error moving file (ID: {}): {}", file_id, e);
|
||||
FileServiceError::from(e)
|
||||
})?;
|
||||
|
||||
tracing::info!("Archivo movido exitosamente: {} (ID: {}) a carpeta: {:?}",
|
||||
tracing::info!("File moved successfully: {} (ID: {}) to folder: {:?}",
|
||||
moved_file.name(), moved_file.id(), moved_file.folder_id());
|
||||
|
||||
Ok(FileDto::from(moved_file))
|
||||
|
||||
@@ -7,7 +7,7 @@ use crate::{
|
||||
application::{
|
||||
dtos::{
|
||||
pagination::PaginatedResponseDto,
|
||||
share_dto::{CreateShareDto, ShareDto, SharePermissionsDto, UpdateShareDto},
|
||||
share_dto::{CreateShareDto, ShareDto, UpdateShareDto},
|
||||
},
|
||||
ports::{
|
||||
outbound::{FileStoragePort, FolderStoragePort},
|
||||
|
||||
@@ -53,9 +53,9 @@ impl TrashService {
|
||||
}
|
||||
}
|
||||
|
||||
/// Convierte una entidad TrashedItem a un DTO
|
||||
/// Converts a TrashedItem entity to a DTO
|
||||
fn to_dto(&self, item: TrashedItem) -> TrashedItemDto {
|
||||
// Calcular days_until_deletion antes de mover item.original_path
|
||||
// Calculate days_until_deletion before moving item.original_path
|
||||
let days_until_deletion = item.days_until_deletion();
|
||||
|
||||
TrashedItemDto {
|
||||
@@ -72,12 +72,12 @@ impl TrashService {
|
||||
}
|
||||
}
|
||||
|
||||
/// Valida los permisos del usuario sobre un elemento
|
||||
/// Validates user permissions over an item
|
||||
#[instrument(skip(self))]
|
||||
async fn validate_user_ownership(&self, _item_id: &str, _user_id: &str) -> Result<()> {
|
||||
// Aquí implementaríamos la validación de permisos
|
||||
// Por ahora, simplemente devolvemos Ok ya que no tenemos una implementación completa
|
||||
// de permisos por usuario
|
||||
// Here we would implement permission validation
|
||||
// For now, we simply return Ok since we don't have a complete
|
||||
// implementation of user permissions
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
@@ -86,7 +86,7 @@ impl TrashService {
|
||||
impl TrashUseCase for TrashService {
|
||||
#[instrument(skip(self))]
|
||||
async fn get_trash_items(&self, user_id: &str) -> Result<Vec<TrashedItemDto>> {
|
||||
debug!("Obteniendo elementos en papelera para usuario: {}", user_id);
|
||||
debug!("Getting trash items for user: {}", user_id);
|
||||
|
||||
let user_uuid = Uuid::parse_str(user_id)
|
||||
.map_err(|e| DomainError::validation_error("User", format!("Invalid user ID: {}", e)))?;
|
||||
@@ -102,52 +102,52 @@ impl TrashUseCase for TrashService {
|
||||
|
||||
#[instrument(skip(self))]
|
||||
async fn move_to_trash(&self, item_id: &str, item_type: &str, user_id: &str) -> Result<()> {
|
||||
info!("Moviendo a papelera: tipo={}, id={}, usuario={}", item_type, item_id, user_id);
|
||||
info!("Moving to trash: type={}, id={}, user={}", item_type, item_id, user_id);
|
||||
debug!("User UUID validation: {}", user_id);
|
||||
|
||||
// Validate user ownership
|
||||
debug!("Validando permisos de usuario");
|
||||
debug!("Validating user permissions");
|
||||
self.validate_user_ownership(item_id, user_id).await?;
|
||||
debug!("Permisos de usuario validados");
|
||||
debug!("User permissions validated");
|
||||
|
||||
// Parse UUIDs with detailed error handling
|
||||
debug!("Validando UUID del item: {}", item_id);
|
||||
debug!("Validating item UUID: {}", item_id);
|
||||
let item_uuid = match Uuid::parse_str(item_id) {
|
||||
Ok(uuid) => {
|
||||
debug!("UUID del item válido: {}", uuid);
|
||||
debug!("Valid item UUID: {}", uuid);
|
||||
uuid
|
||||
},
|
||||
Err(e) => {
|
||||
error!("UUID del item inválido: {} - Error: {}", item_id, e);
|
||||
error!("Invalid item UUID: {} - Error: {}", item_id, e);
|
||||
return Err(DomainError::validation_error("Item", format!("Invalid item ID: {}", e)));
|
||||
}
|
||||
};
|
||||
|
||||
debug!("Validando UUID del usuario: {}", user_id);
|
||||
debug!("Validating user UUID: {}", user_id);
|
||||
let user_uuid = match Uuid::parse_str(user_id) {
|
||||
Ok(uuid) => {
|
||||
debug!("UUID del usuario válido: {}", uuid);
|
||||
debug!("Valid user UUID: {}", uuid);
|
||||
uuid
|
||||
},
|
||||
Err(e) => {
|
||||
error!("UUID del usuario inválido: {} - Error: {}", user_id, e);
|
||||
error!("Invalid user UUID: {} - Error: {}", user_id, e);
|
||||
return Err(DomainError::validation_error("User", format!("Invalid user ID: {}", e)));
|
||||
}
|
||||
};
|
||||
|
||||
match item_type {
|
||||
"file" => {
|
||||
info!("Procesando archivo para mover a papelera: {}", item_id);
|
||||
info!("Processing file to move to trash: {}", item_id);
|
||||
|
||||
// Obtener el archivo para verificar que existe y capturar sus datos
|
||||
debug!("Obteniendo datos del archivo: {}", item_id);
|
||||
// Get the file to verify it exists and capture its data
|
||||
debug!("Getting file data: {}", item_id);
|
||||
let file = match self.file_repository.get_file_by_id(item_id).await {
|
||||
Ok(file) => {
|
||||
debug!("Archivo encontrado: {} ({})", file.name(), item_id);
|
||||
debug!("File found: {} ({})", file.name(), item_id);
|
||||
file
|
||||
},
|
||||
Err(e) => {
|
||||
error!("Error al obtener archivo: {} - {}", item_id, e);
|
||||
error!("Error getting file: {} - {}", item_id, e);
|
||||
return Err(DomainError::new(
|
||||
ErrorKind::NotFound,
|
||||
"File",
|
||||
@@ -157,10 +157,10 @@ impl TrashUseCase for TrashService {
|
||||
};
|
||||
|
||||
let original_path = file.storage_path().to_string();
|
||||
debug!("Ruta original del archivo: {}", original_path);
|
||||
debug!("Original file path: {}", original_path);
|
||||
|
||||
// Crear el elemento de papelera
|
||||
debug!("Creando objeto TrashedItem para el archivo");
|
||||
// Create the trash item
|
||||
debug!("Creating TrashedItem object for the file");
|
||||
let trashed_item = TrashedItem::new(
|
||||
item_uuid,
|
||||
user_uuid,
|
||||
@@ -169,28 +169,28 @@ impl TrashUseCase for TrashService {
|
||||
original_path,
|
||||
self.retention_days,
|
||||
);
|
||||
debug!("TrashedItem creado con éxito: {} -> {}", file.name(), trashed_item.id);
|
||||
debug!("TrashedItem created successfully: {} -> {}", file.name(), trashed_item.id);
|
||||
|
||||
// Primero añadimos a la papelera para registrar el elemento
|
||||
info!("Añadiendo archivo {} a índice de papelera", item_id);
|
||||
// First add to trash index to register the item
|
||||
info!("Adding file {} to trash index", item_id);
|
||||
match self.trash_repository.add_to_trash(&trashed_item).await {
|
||||
Ok(_) => {
|
||||
debug!("Archivo añadido al índice de papelera con éxito");
|
||||
debug!("File added to trash index successfully");
|
||||
},
|
||||
Err(e) => {
|
||||
error!("Error al añadir archivo al índice de papelera: {}", e);
|
||||
error!("Error adding file to trash index: {}", e);
|
||||
return Err(DomainError::internal_error("TrashRepository", format!("Failed to add file to trash: {}", e)));
|
||||
}
|
||||
};
|
||||
|
||||
// Luego movemos el archivo físicamente a la papelera
|
||||
info!("Moviendo archivo físicamente a la papelera: {}", item_id);
|
||||
// Then physically move the file to trash
|
||||
info!("Physically moving file to trash: {}", item_id);
|
||||
match self.file_repository.move_to_trash(item_id).await {
|
||||
Ok(_) => {
|
||||
debug!("Archivo movido físicamente a papelera con éxito: {}", item_id);
|
||||
debug!("File physically moved to trash successfully: {}", item_id);
|
||||
},
|
||||
Err(e) => {
|
||||
error!("Error al mover archivo físicamente a papelera: {} - {}", item_id, e);
|
||||
error!("Error physically moving file to trash: {} - {}", item_id, e);
|
||||
return Err(DomainError::new(
|
||||
ErrorKind::InternalError,
|
||||
"File",
|
||||
@@ -199,11 +199,11 @@ impl TrashUseCase for TrashService {
|
||||
}
|
||||
}
|
||||
|
||||
info!("Archivo movido a papelera completamente: {}", item_id);
|
||||
info!("File completely moved to trash: {}", item_id);
|
||||
Ok(())
|
||||
},
|
||||
"folder" => {
|
||||
// Obtener la carpeta para verificar que existe y capturar sus datos
|
||||
// Get the folder to verify it exists and capture its data
|
||||
let folder = self.folder_repository.get_folder_by_id(item_id).await
|
||||
.map_err(|e| DomainError::new(
|
||||
ErrorKind::NotFound,
|
||||
@@ -213,7 +213,7 @@ impl TrashUseCase for TrashService {
|
||||
|
||||
let original_path = folder.storage_path().to_string();
|
||||
|
||||
// Crear el elemento de papelera
|
||||
// Create the trash item
|
||||
let trashed_item = TrashedItem::new(
|
||||
item_uuid,
|
||||
user_uuid,
|
||||
@@ -223,7 +223,7 @@ impl TrashUseCase for TrashService {
|
||||
self.retention_days,
|
||||
);
|
||||
|
||||
// Primero añadimos a la papelera para registrar el elemento
|
||||
// First add to trash index to register the item
|
||||
debug!("Adding folder {} to trash repository", item_id);
|
||||
match self.trash_repository.add_to_trash(&trashed_item).await {
|
||||
Ok(_) => debug!("Successfully added folder to trash repository"),
|
||||
@@ -233,7 +233,7 @@ impl TrashUseCase for TrashService {
|
||||
}
|
||||
};
|
||||
|
||||
// Luego movemos la carpeta físicamente a la papelera
|
||||
// Then physically move the folder to trash
|
||||
self.folder_repository.move_to_trash(item_id).await
|
||||
.map_err(|e| DomainError::new(
|
||||
ErrorKind::InternalError,
|
||||
@@ -241,7 +241,7 @@ impl TrashUseCase for TrashService {
|
||||
format!("Error moving folder {} to trash: {}", item_id, e)
|
||||
))?;
|
||||
|
||||
debug!("Carpeta movida a papelera: {}", item_id);
|
||||
debug!("Folder moved to trash: {}", item_id);
|
||||
Ok(())
|
||||
},
|
||||
_ => Err(DomainError::validation_error("Item", format!("Invalid item type: {}", item_type))),
|
||||
@@ -250,7 +250,7 @@ impl TrashUseCase for TrashService {
|
||||
|
||||
#[instrument(skip(self))]
|
||||
async fn restore_item(&self, trash_id: &str, user_id: &str) -> Result<()> {
|
||||
info!("Restaurando elemento {} para usuario {}", trash_id, user_id);
|
||||
info!("Restoring item {} for user {}", trash_id, user_id);
|
||||
|
||||
let trash_uuid = match Uuid::parse_str(trash_id) {
|
||||
Ok(id) => {
|
||||
@@ -283,10 +283,10 @@ impl TrashUseCase for TrashService {
|
||||
info!("Found item in trash: ID={}, Type={:?}, OriginalID={}",
|
||||
trash_id, item.item_type, item.original_id);
|
||||
|
||||
// Restaurar según tipo
|
||||
// Restore based on type
|
||||
match item.item_type {
|
||||
TrashedItemType::File => {
|
||||
// Restaurar el archivo a su ubicación original
|
||||
// Restore the file to its original location
|
||||
let file_id = item.original_id.to_string();
|
||||
let original_path = item.original_path.clone();
|
||||
|
||||
@@ -313,7 +313,7 @@ impl TrashUseCase for TrashService {
|
||||
}
|
||||
},
|
||||
TrashedItemType::Folder => {
|
||||
// Restaurar la carpeta a su ubicación original
|
||||
// Restore the folder to its original location
|
||||
let folder_id = item.original_id.to_string();
|
||||
let original_path = item.original_path.clone();
|
||||
|
||||
@@ -408,7 +408,7 @@ impl TrashUseCase for TrashService {
|
||||
info!("Found item in trash: ID={}, Type={:?}, OriginalID={}",
|
||||
trash_id, item.item_type, item.original_id);
|
||||
|
||||
// Eliminar permanentemente según tipo
|
||||
// Permanently delete based on type
|
||||
match item.item_type {
|
||||
TrashedItemType::File => {
|
||||
// Eliminar el archivo permanentemente
|
||||
@@ -463,7 +463,7 @@ impl TrashUseCase for TrashService {
|
||||
}
|
||||
}
|
||||
|
||||
// Eliminar el item de la papelera siempre, para mantener consistencia
|
||||
// Always remove the item from trash index to maintain consistency
|
||||
info!("Removing entry from trash index: {}", trash_id);
|
||||
match self.trash_repository.delete_permanently(&trash_uuid, &user_uuid).await {
|
||||
Ok(_) => {
|
||||
@@ -497,38 +497,38 @@ impl TrashUseCase for TrashService {
|
||||
|
||||
#[instrument(skip(self))]
|
||||
async fn empty_trash(&self, user_id: &str) -> Result<()> {
|
||||
info!("Vaciando papelera para usuario {}", user_id);
|
||||
info!("Emptying trash for user {}", user_id);
|
||||
|
||||
let user_uuid = Uuid::parse_str(user_id)
|
||||
.map_err(|e| DomainError::validation_error("User", format!("Invalid user ID: {}", e)))?;
|
||||
|
||||
// Obtener todos los elementos en la papelera del usuario
|
||||
// Get all items in the user's trash
|
||||
let items = self.trash_repository.get_trash_items(&user_uuid).await?;
|
||||
|
||||
// Eliminar permanentemente cada elemento
|
||||
// Permanently delete each item
|
||||
for item in items {
|
||||
match item.item_type {
|
||||
TrashedItemType::File => {
|
||||
// Eliminar el archivo permanentemente
|
||||
// Permanently delete the file
|
||||
let file_id = item.original_id.to_string();
|
||||
if let Err(e) = self.file_repository.delete_file_permanently(&file_id).await {
|
||||
error!("Error al eliminar archivo {} permanentemente: {}", file_id, e);
|
||||
error!("Error permanently deleting file {}: {}", file_id, e);
|
||||
}
|
||||
},
|
||||
TrashedItemType::Folder => {
|
||||
// Eliminar la carpeta permanentemente
|
||||
// Permanently delete the folder
|
||||
let folder_id = item.original_id.to_string();
|
||||
if let Err(e) = self.folder_repository.delete_folder_permanently(&folder_id).await {
|
||||
error!("Error al eliminar carpeta {} permanentemente: {}", folder_id, e);
|
||||
error!("Error permanently deleting folder {}: {}", folder_id, e);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Limpiar todos los registros de la papelera para este usuario
|
||||
// Clear all trash records for this user
|
||||
self.trash_repository.clear_trash(&user_uuid).await?;
|
||||
|
||||
info!("Papelera vaciada completamente para usuario {}", user_id);
|
||||
info!("Trash completely emptied for user {}", user_id);
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
+24
-24
@@ -10,11 +10,11 @@ use crate::domain::services::path_service::StoragePath;
|
||||
#[derive(Debug, thiserror::Error)]
|
||||
pub enum FileError {
|
||||
/// Occurs when a file name contains invalid characters or is empty.
|
||||
#[error("Nombre de archivo inválido: {0}")]
|
||||
#[error("Invalid file name: {0}")]
|
||||
InvalidFileName(String),
|
||||
|
||||
/// Occurs when validation fails for any file entity attribute.
|
||||
#[error("Error en la validación: {0}")]
|
||||
#[error("Validation error: {0}")]
|
||||
#[allow(dead_code)]
|
||||
ValidationError(String),
|
||||
}
|
||||
@@ -68,7 +68,7 @@ pub struct File {
|
||||
modified_at: u64,
|
||||
}
|
||||
|
||||
// Ya no necesitamos este módulo, ahora usamos un String directamente
|
||||
// We no longer need this module, now we use a String directly
|
||||
|
||||
impl Default for File {
|
||||
fn default() -> Self {
|
||||
@@ -87,7 +87,7 @@ impl Default for File {
|
||||
}
|
||||
|
||||
impl File {
|
||||
/// Crea un nuevo archivo con validación
|
||||
/// Creates a new file with validation
|
||||
pub fn new(
|
||||
id: String,
|
||||
name: String,
|
||||
@@ -96,7 +96,7 @@ impl File {
|
||||
mime_type: String,
|
||||
folder_id: Option<String>,
|
||||
) -> FileResult<Self> {
|
||||
// Validar nombre de archivo
|
||||
// Validate file name
|
||||
if name.is_empty() || name.contains('/') || name.contains('\\') {
|
||||
return Err(FileError::InvalidFileName(name));
|
||||
}
|
||||
@@ -106,7 +106,7 @@ impl File {
|
||||
.unwrap_or_default()
|
||||
.as_secs();
|
||||
|
||||
// Almacenamos el string de la ruta para compatibilidad con serialización
|
||||
// Store the path string for serialization compatibility
|
||||
let path_string = storage_path.to_string();
|
||||
|
||||
Ok(Self {
|
||||
@@ -122,7 +122,7 @@ impl File {
|
||||
})
|
||||
}
|
||||
|
||||
/// Crea un archivo con timestamps específicos (para reconstrucción)
|
||||
/// Creates a file with specific timestamps (for reconstruction)
|
||||
pub fn with_timestamps(
|
||||
id: String,
|
||||
name: String,
|
||||
@@ -133,12 +133,12 @@ impl File {
|
||||
created_at: u64,
|
||||
modified_at: u64,
|
||||
) -> FileResult<Self> {
|
||||
// Validar nombre de archivo
|
||||
// Validate file name
|
||||
if name.is_empty() || name.contains('/') || name.contains('\\') {
|
||||
return Err(FileError::InvalidFileName(name));
|
||||
}
|
||||
|
||||
// Almacenamos el string de la ruta para compatibilidad con serialización
|
||||
// Store the path string for serialization compatibility
|
||||
let path_string = storage_path.to_string();
|
||||
|
||||
Ok(Self {
|
||||
@@ -191,8 +191,8 @@ impl File {
|
||||
self.modified_at
|
||||
}
|
||||
|
||||
/// Crea una nueva instancia de File desde un DTO
|
||||
/// Esta función es principalmente para conversiones en los batch handlers
|
||||
/// Creates a new File instance from a DTO
|
||||
/// This function is primarily for conversions in batch handlers
|
||||
pub fn from_dto(
|
||||
id: String,
|
||||
name: String,
|
||||
@@ -203,10 +203,10 @@ impl File {
|
||||
created_at: u64,
|
||||
modified_at: u64,
|
||||
) -> Self {
|
||||
// Crear storage_path desde el string
|
||||
// Create storage_path from string
|
||||
let storage_path = StoragePath::from_string(&path);
|
||||
|
||||
// Crear directamente sin validación para evitar errores en conversiones DTO
|
||||
// Create directly without validation to avoid errors in DTO conversions
|
||||
Self {
|
||||
id,
|
||||
name,
|
||||
@@ -220,24 +220,24 @@ impl File {
|
||||
}
|
||||
}
|
||||
|
||||
// Métodos para crear nuevas versiones del archivo (inmutable)
|
||||
// Methods to create new versions of the file (immutable)
|
||||
|
||||
/// Crea una nueva versión del archivo con nombre actualizado
|
||||
/// Creates a new version of the file with updated name
|
||||
#[allow(dead_code)]
|
||||
pub fn with_name(&self, new_name: String) -> FileResult<Self> {
|
||||
// Validar nombre de archivo
|
||||
// Validate file name
|
||||
if new_name.is_empty() || new_name.contains('/') || new_name.contains('\\') {
|
||||
return Err(FileError::InvalidFileName(new_name));
|
||||
}
|
||||
|
||||
// Actualizar ruta basada en el nombre
|
||||
// Update path based on name
|
||||
let parent_path = self.storage_path.parent();
|
||||
let new_storage_path = match parent_path {
|
||||
Some(parent) => parent.join(&new_name),
|
||||
None => StoragePath::from_string(&new_name),
|
||||
};
|
||||
|
||||
// Actualizar representación en string
|
||||
// Update string representation
|
||||
let new_path_string = new_storage_path.to_string();
|
||||
|
||||
let now = std::time::SystemTime::now()
|
||||
@@ -258,15 +258,15 @@ impl File {
|
||||
})
|
||||
}
|
||||
|
||||
/// Crea una nueva versión del archivo con carpeta actualizada
|
||||
/// Creates a new version of the file with updated folder
|
||||
pub fn with_folder(&self, folder_id: Option<String>, folder_path: Option<StoragePath>) -> FileResult<Self> {
|
||||
// Necesitamos una ruta de carpeta para actualizar la ruta del archivo
|
||||
// We need a folder path to update the file path
|
||||
let new_storage_path = match folder_path {
|
||||
Some(path) => path.join(&self.name),
|
||||
None => StoragePath::from_string(&self.name), // Raíz
|
||||
None => StoragePath::from_string(&self.name), // Root
|
||||
};
|
||||
|
||||
// Actualizar representación en string
|
||||
// Update string representation
|
||||
let new_path_string = new_storage_path.to_string();
|
||||
|
||||
let now = std::time::SystemTime::now()
|
||||
@@ -287,7 +287,7 @@ impl File {
|
||||
})
|
||||
}
|
||||
|
||||
/// Crea una nueva versión del archivo con tamaño actualizado
|
||||
/// Creates a new version of the file with updated size
|
||||
#[allow(dead_code)]
|
||||
pub fn with_size(&self, new_size: u64) -> Self {
|
||||
let now = std::time::SystemTime::now()
|
||||
@@ -333,7 +333,7 @@ mod tests {
|
||||
let storage_path = StoragePath::from_string("/test/invalid/file.txt");
|
||||
let file = File::new(
|
||||
"123".to_string(),
|
||||
"file/with/slash.txt".to_string(), // Nombre inválido
|
||||
"file/with/slash.txt".to_string(), // Invalid name
|
||||
storage_path,
|
||||
100,
|
||||
"text/plain".to_string(),
|
||||
|
||||
@@ -1,18 +1,18 @@
|
||||
use serde::{Serialize, Deserialize};
|
||||
use crate::domain::services::path_service::StoragePath;
|
||||
|
||||
/// Error en la creación o manipulación de entidades de carpeta
|
||||
/// Error in the creation or manipulation of folder entities
|
||||
#[derive(Debug, thiserror::Error)]
|
||||
pub enum FolderError {
|
||||
#[error("Nombre de carpeta inválido: {0}")]
|
||||
#[error("Invalid folder name: {0}")]
|
||||
InvalidFolderName(String),
|
||||
|
||||
#[error("Error en la validación: {0}")]
|
||||
#[error("Validation error: {0}")]
|
||||
#[allow(dead_code)]
|
||||
ValidationError(String),
|
||||
}
|
||||
|
||||
/// Tipo de resultado para operaciones con entidades de carpeta
|
||||
/// Result type for folder entity operations
|
||||
pub type FolderResult<T> = Result<T, FolderError>;
|
||||
|
||||
/// Represents a folder entity in the domain
|
||||
@@ -42,7 +42,7 @@ pub struct Folder {
|
||||
modified_at: u64,
|
||||
}
|
||||
|
||||
// Ya no necesitamos este módulo, ahora usamos un String directamente
|
||||
// We no longer need this module, now we use a String directly
|
||||
|
||||
impl Default for Folder {
|
||||
fn default() -> Self {
|
||||
@@ -66,7 +66,7 @@ impl Folder {
|
||||
storage_path: StoragePath,
|
||||
parent_id: Option<String>,
|
||||
) -> FolderResult<Self> {
|
||||
// Validar nombre de carpeta
|
||||
// Validate folder name
|
||||
if name.is_empty() || name.contains('/') || name.contains('\\') {
|
||||
return Err(FolderError::InvalidFolderName(name));
|
||||
}
|
||||
@@ -76,7 +76,7 @@ impl Folder {
|
||||
.unwrap_or_default()
|
||||
.as_secs();
|
||||
|
||||
// Almacenamos el string de la ruta para compatibilidad con serialización
|
||||
// Store the path string for serialization compatibility
|
||||
let path_string = storage_path.to_string();
|
||||
|
||||
Ok(Self {
|
||||
@@ -99,12 +99,12 @@ impl Folder {
|
||||
created_at: u64,
|
||||
modified_at: u64,
|
||||
) -> FolderResult<Self> {
|
||||
// Validar nombre de carpeta
|
||||
// Validate folder name
|
||||
if name.is_empty() || name.contains('/') || name.contains('\\') {
|
||||
return Err(FolderError::InvalidFolderName(name));
|
||||
}
|
||||
|
||||
// Almacenamos el string de la ruta para compatibilidad con serialización
|
||||
// Store the path string for serialization compatibility
|
||||
let path_string = storage_path.to_string();
|
||||
|
||||
Ok(Self {
|
||||
@@ -147,8 +147,8 @@ impl Folder {
|
||||
self.modified_at
|
||||
}
|
||||
|
||||
/// Crea una nueva instancia de Folder desde un DTO
|
||||
/// Esta función es principalmente para conversiones en los batch handlers
|
||||
/// Creates a new Folder instance from a DTO
|
||||
/// This function is primarily for conversions in batch handlers
|
||||
pub fn from_dto(
|
||||
id: String,
|
||||
name: String,
|
||||
@@ -157,10 +157,10 @@ impl Folder {
|
||||
created_at: u64,
|
||||
modified_at: u64,
|
||||
) -> Self {
|
||||
// Crear storage_path desde el string
|
||||
// Create storage_path from the string
|
||||
let storage_path = StoragePath::from_string(&path);
|
||||
|
||||
// Crear directamente sin validación para evitar errores en conversiones DTO
|
||||
// Create directly without validation to avoid errors in DTO conversions
|
||||
Self {
|
||||
id,
|
||||
name,
|
||||
@@ -172,23 +172,23 @@ impl Folder {
|
||||
}
|
||||
}
|
||||
|
||||
// Métodos para crear nuevas versiones de la carpeta (inmutable)
|
||||
// Methods to create new versions of the folder (immutable)
|
||||
|
||||
/// Creates a new version of the folder with updated name
|
||||
pub fn with_name(&self, new_name: String) -> FolderResult<Self> {
|
||||
// Validar nombre de carpeta
|
||||
// Validate folder name
|
||||
if new_name.is_empty() || new_name.contains('/') || new_name.contains('\\') {
|
||||
return Err(FolderError::InvalidFolderName(new_name));
|
||||
}
|
||||
|
||||
// Actualizar ruta basada en el nombre
|
||||
// Update path based on the name
|
||||
let parent_path = self.storage_path.parent();
|
||||
let new_storage_path = match parent_path {
|
||||
Some(parent) => parent.join(&new_name),
|
||||
None => StoragePath::from_string(&new_name),
|
||||
};
|
||||
|
||||
// Actualizar representación en string
|
||||
// Update string representation
|
||||
let new_path_string = new_storage_path.to_string();
|
||||
|
||||
let now = std::time::SystemTime::now()
|
||||
@@ -209,13 +209,13 @@ impl Folder {
|
||||
|
||||
/// Creates a new version of the folder with updated parent
|
||||
pub fn with_parent(&self, parent_id: Option<String>, parent_path: Option<StoragePath>) -> FolderResult<Self> {
|
||||
// Necesitamos una ruta de carpeta para actualizar la ruta
|
||||
// We need a folder path to update the path
|
||||
let new_storage_path = match parent_path {
|
||||
Some(path) => path.join(&self.name),
|
||||
None => StoragePath::from_string(&self.name), // Raíz
|
||||
None => StoragePath::from_string(&self.name), // Root
|
||||
};
|
||||
|
||||
// Actualizar representación en string
|
||||
// Update string representation
|
||||
let new_path_string = new_storage_path.to_string();
|
||||
|
||||
let now = std::time::SystemTime::now()
|
||||
@@ -276,7 +276,7 @@ mod tests {
|
||||
let storage_path = StoragePath::from_string("/test/invalid/folder");
|
||||
let folder = Folder::new(
|
||||
"123".to_string(),
|
||||
"folder/with/slash".to_string(), // Nombre inválido
|
||||
"folder/with/slash".to_string(), // Invalid name
|
||||
storage_path,
|
||||
None,
|
||||
);
|
||||
@@ -302,6 +302,6 @@ mod tests {
|
||||
assert!(renamed.is_ok());
|
||||
let renamed = renamed.unwrap();
|
||||
assert_eq!(renamed.name(), "new_name");
|
||||
assert_eq!(renamed.id(), "123"); // El ID no cambia
|
||||
assert_eq!(renamed.id(), "123"); // The ID doesn't change
|
||||
}
|
||||
}
|
||||
@@ -1,4 +1,3 @@
|
||||
use std::sync::Arc;
|
||||
|
||||
use async_trait::async_trait;
|
||||
use thiserror::Error;
|
||||
|
||||
@@ -45,8 +45,8 @@ use crate::infrastructure::repositories::parallel_file_processor::ParallelFilePr
|
||||
* filesystem-specific details.
|
||||
*/
|
||||
|
||||
// Usar constantes de la configuración centralizada en lugar de valores fijos
|
||||
// Esto se reemplaza con self.config.concurrency.max_concurrent_files más adelante
|
||||
// Use constants from centralized configuration instead of fixed values
|
||||
// This is replaced with self.config.concurrency.max_concurrent_files later
|
||||
|
||||
/// Filesystem implementation of the FileRepository interface
|
||||
pub struct FileFsRepository {
|
||||
@@ -130,16 +130,16 @@ impl FileFsRepository {
|
||||
async fn file_exists_at_storage_path(&self, storage_path: &StoragePath) -> FileRepositoryResult<bool> {
|
||||
let abs_path = self.resolve_storage_path(storage_path);
|
||||
|
||||
// Intentar obtener del caché avanzado primero
|
||||
// Try to get from advanced cache first
|
||||
if let Some(is_file) = self.metadata_cache.is_file(&abs_path).await {
|
||||
tracing::debug!("Metadata cache hit for existence check: {} - path: {}", is_file, abs_path.display());
|
||||
return Ok(is_file);
|
||||
}
|
||||
|
||||
// Si no está en caché, verificar directamente y actualizar caché
|
||||
// If not in cache, verify directly and update cache
|
||||
tracing::debug!("Metadata cache miss for existence check: {}", abs_path.display());
|
||||
|
||||
// Utilizar timeout para evitar bloqueo
|
||||
// Use timeout to avoid blocking
|
||||
match time::timeout(
|
||||
self.config.timeouts.file_timeout(),
|
||||
fs::metadata(&abs_path)
|
||||
@@ -147,7 +147,7 @@ impl FileFsRepository {
|
||||
Ok(Ok(metadata)) => {
|
||||
let is_file = metadata.is_file();
|
||||
|
||||
// Actualizar la caché con información fresca
|
||||
// Update cache with fresh information
|
||||
if let Err(e) = self.metadata_cache.refresh_metadata(&abs_path).await {
|
||||
tracing::warn!("Failed to update cache for {}: {}", abs_path.display(), e);
|
||||
}
|
||||
@@ -163,7 +163,7 @@ impl FileFsRepository {
|
||||
Ok(Err(e)) => {
|
||||
tracing::warn!("File check failed: {} - {}", abs_path.display(), e);
|
||||
|
||||
// Añadir a caché como no existente
|
||||
// Add to cache as non-existent
|
||||
let entry_type = CacheEntryType::Unknown;
|
||||
let file_metadata = crate::infrastructure::services::file_metadata_cache::FileMetadata::new(
|
||||
abs_path.clone(),
|
||||
@@ -191,13 +191,13 @@ impl FileFsRepository {
|
||||
pub async fn file_exists(&self, path: &std::path::Path) -> FileRepositoryResult<bool> {
|
||||
let abs_path = self.resolve_legacy_path(path);
|
||||
|
||||
// Intentar obtener del caché avanzado primero
|
||||
// Try to get from advanced cache first
|
||||
if let Some(is_file) = self.metadata_cache.is_file(&abs_path).await {
|
||||
tracing::debug!("Metadata cache hit for legacy existence check: {} - path: {}", is_file, abs_path.display());
|
||||
return Ok(is_file);
|
||||
}
|
||||
|
||||
// Si no está en caché, verificar directamente
|
||||
// If not in cache, verify directly
|
||||
tracing::info!("Checking if file exists: {} - path: {}", abs_path.exists(), abs_path.display());
|
||||
|
||||
match time::timeout(
|
||||
@@ -207,7 +207,7 @@ impl FileFsRepository {
|
||||
Ok(Ok(metadata)) => {
|
||||
let is_file = metadata.is_file();
|
||||
|
||||
// Actualizar la caché con información fresca
|
||||
// Update cache with fresh information
|
||||
if let Err(e) = self.metadata_cache.refresh_metadata(&abs_path).await {
|
||||
tracing::warn!("Failed to update cache for {}: {}", abs_path.display(), e);
|
||||
}
|
||||
@@ -271,7 +271,7 @@ impl FileFsRepository {
|
||||
|
||||
/// Extracts file metadata from a physical path with timeout and cache
|
||||
async fn get_file_metadata(&self, abs_path: &PathBuf) -> FileRepositoryResult<(u64, u64, u64)> {
|
||||
// Intentar obtener de caché primero
|
||||
// Try to get from cache first
|
||||
if let Some(cached_metadata) = self.metadata_cache.get_metadata(abs_path).await {
|
||||
if let (Some(size), Some(created_at), Some(modified_at)) =
|
||||
(cached_metadata.size, cached_metadata.created_at, cached_metadata.modified_at) {
|
||||
@@ -280,7 +280,7 @@ impl FileFsRepository {
|
||||
}
|
||||
}
|
||||
|
||||
// Si no está en caché o metadatos incompletos, cargar desde sistema de archivos
|
||||
// If not in cache or incomplete metadata, load from filesystem
|
||||
let metadata = match time::timeout(
|
||||
self.config.timeouts.file_timeout(),
|
||||
fs::metadata(&abs_path)
|
||||
@@ -304,7 +304,7 @@ impl FileFsRepository {
|
||||
.map(|time| time.duration_since(std::time::UNIX_EPOCH).unwrap_or_default().as_secs())
|
||||
.unwrap_or_else(|_| 0);
|
||||
|
||||
// Actualizar caché si es posible
|
||||
// Update cache if possible
|
||||
if let Err(e) = self.metadata_cache.refresh_metadata(abs_path).await {
|
||||
tracing::warn!("Failed to update metadata cache for {}: {}", abs_path.display(), e);
|
||||
}
|
||||
@@ -340,7 +340,7 @@ impl FileFsRepository {
|
||||
.map_err(|_| FileRepositoryError::Timeout(format!("Timeout checking file size: {}", abs_path.display())))?
|
||||
.map_err(FileRepositoryError::IoError)?;
|
||||
|
||||
// Utiliza el método del ResourceConfig para determinar si es un archivo grande
|
||||
// Use the ResourceConfig method to determine if it's a large file
|
||||
Ok(self.config.resources.is_large_file(metadata.len()))
|
||||
}
|
||||
|
||||
@@ -395,7 +395,7 @@ impl FileRepositoryError {
|
||||
}
|
||||
}
|
||||
|
||||
// Los errores ya están definidos por la interfaz FileRepositoryError
|
||||
// Errors are already defined by the FileRepositoryError interface
|
||||
|
||||
// Enable cloning for concurrent operations
|
||||
impl Clone for FileFsRepository {
|
||||
|
||||
@@ -5,19 +5,19 @@ use tracing::{debug, error, instrument};
|
||||
use crate::domain::repositories::file_repository::FileRepositoryResult;
|
||||
use crate::infrastructure::repositories::file_fs_repository::FileFsRepository;
|
||||
|
||||
// Este archivo contiene la implementación de los métodos relacionados con la papelera
|
||||
// para el repositorio de archivos FileFsRepository
|
||||
// This file contains the implementation of trash-related methods
|
||||
// for the FileFsRepository file repository
|
||||
|
||||
// Implementación de métodos de papelera para el repositorio de archivos
|
||||
// Implementation of trash methods for the file repository
|
||||
impl FileFsRepository {
|
||||
// Obtiene la ruta completa a la papelera
|
||||
// Gets the complete path to the trash directory
|
||||
fn get_trash_dir(&self) -> PathBuf {
|
||||
let trash_dir = self.get_root_path().join(".trash").join("files");
|
||||
debug!("Base trash directory: {}", trash_dir.display());
|
||||
trash_dir
|
||||
}
|
||||
|
||||
// Obtiene la ruta de la papelera para un usuario específico (si se proporciona)
|
||||
// Gets the trash directory path for a specific user (if provided)
|
||||
fn get_user_trash_dir(&self, user_id: Option<&str>) -> PathBuf {
|
||||
let base_trash_dir = self.get_trash_dir();
|
||||
|
||||
@@ -33,7 +33,7 @@ impl FileFsRepository {
|
||||
}
|
||||
}
|
||||
|
||||
// Crea una ruta única en la papelera para el archivo
|
||||
// Creates a unique path in the trash for the file
|
||||
async fn create_trash_file_path(&self, file_id: &str) -> FileRepositoryResult<PathBuf> {
|
||||
debug!("Creating trash file path for file ID: {}", file_id);
|
||||
|
||||
@@ -62,7 +62,7 @@ impl FileFsRepository {
|
||||
}
|
||||
}
|
||||
|
||||
// Implementación de los métodos públicos del trait FileRepository relacionados con la papelera
|
||||
// Implementation of the public methods of the FileRepository trait related to trash
|
||||
// Note: The FileRepository trait implementation has been moved to file_fs_repository.rs
|
||||
// to avoid duplicate implementations
|
||||
|
||||
@@ -70,101 +70,101 @@ impl FileFsRepository {
|
||||
impl FileFsRepository {
|
||||
/// Helper method that will be used for trash functionality
|
||||
pub(crate) async fn _trash_move_to_trash(&self, file_id: &str) -> FileRepositoryResult<()> {
|
||||
debug!("Moviendo archivo a la papelera: {}", file_id);
|
||||
debug!("Moving file to trash: {}", file_id);
|
||||
|
||||
// Obtener la ruta física del archivo
|
||||
// Creamos un método independiente para acceder al servicio de mapeo de IDs
|
||||
debug!("Obteniendo ruta del archivo con ID: {}", file_id);
|
||||
// Get the physical path of the file
|
||||
// We create an independent method to access the ID mapping service
|
||||
debug!("Getting file path with ID: {}", file_id);
|
||||
let file_path = match self.id_mapping_service().get_file_path(file_id).await {
|
||||
Ok(path) => {
|
||||
debug!("Ruta del archivo obtenida: {}", path.display());
|
||||
debug!("File path obtained: {}", path.display());
|
||||
path
|
||||
},
|
||||
Err(e) => {
|
||||
error!("Error obteniendo ruta del archivo {}: {:?}", file_id, e);
|
||||
error!("Error getting file path {}: {:?}", file_id, e);
|
||||
return Err(FileRepositoryError::IdMappingError(format!("Failed to get file path: {}", e)));
|
||||
}
|
||||
};
|
||||
|
||||
// Verificamos que el archivo existe
|
||||
debug!("Verificando que el archivo existe: {}", file_path.display());
|
||||
// Verify that the file exists
|
||||
debug!("Verifying that the file exists: {}", file_path.display());
|
||||
if !self.file_exists(&file_path).await? {
|
||||
error!("Archivo no encontrado en la ruta especificada: {}", file_path.display());
|
||||
error!("File not found at the specified path: {}", file_path.display());
|
||||
return Err(FileRepositoryError::NotFound(format!("File not found: {}", file_id)));
|
||||
}
|
||||
debug!("Archivo encontrado, continuando con la operación");
|
||||
debug!("File found, continuing with the operation");
|
||||
|
||||
// Crear directorio en la papelera si no existe
|
||||
debug!("Creando path para archivo en papelera");
|
||||
// Create directory in trash if it doesn't exist
|
||||
debug!("Creating path for file in trash");
|
||||
let trash_file_path = self.create_trash_file_path(file_id).await?;
|
||||
debug!("Path en papelera: {}", trash_file_path.display());
|
||||
debug!("Path in trash: {}", trash_file_path.display());
|
||||
|
||||
// Mover el archivo físicamente a la papelera (no actualiza mappings)
|
||||
debug!("Moviendo archivo físicamente a papelera: {} -> {}", file_path.display(), trash_file_path.display());
|
||||
// Physically move the file to trash (doesn't update mappings)
|
||||
debug!("Physically moving file to trash: {} -> {}", file_path.display(), trash_file_path.display());
|
||||
match fs::rename(&file_path, &trash_file_path).await {
|
||||
Ok(_) => {
|
||||
debug!("Archivo movido a papelera exitosamente: {} -> {}", file_path.display(), trash_file_path.display());
|
||||
debug!("File successfully moved to trash: {} -> {}", file_path.display(), trash_file_path.display());
|
||||
|
||||
// Invalidar la caché del archivo original
|
||||
debug!("Invalidando caché para: {}", file_path.display());
|
||||
// Invalidate the cache for the original file
|
||||
debug!("Invalidating cache for: {}", file_path.display());
|
||||
self.metadata_cache().invalidate(&file_path).await;
|
||||
|
||||
// Actualizar el mapeo al nuevo path en la papelera
|
||||
debug!("Actualizando mapeo de ID a nuevo path en papelera");
|
||||
// Update the mapping to the new path in trash
|
||||
debug!("Updating ID mapping to new path in trash");
|
||||
if let Err(e) = self.id_mapping_service().update_file_path(file_id, &trash_file_path).await {
|
||||
error!("Error actualizando mapeo de archivo en papelera: {}", e);
|
||||
error!("Error updating file mapping in trash: {}", e);
|
||||
return Err(FileRepositoryError::MappingError(format!("Failed to update mapping: {}", e)));
|
||||
}
|
||||
debug!("Mapeo actualizado exitosamente");
|
||||
debug!("Mapping successfully updated");
|
||||
|
||||
debug!("Operación de mover a papelera completada con éxito para el archivo: {}", file_id);
|
||||
debug!("Move to trash operation completed successfully for file: {}", file_id);
|
||||
Ok(())
|
||||
},
|
||||
Err(e) => {
|
||||
error!("Error moviendo archivo a papelera: {} -> {}: {}",
|
||||
error!("Error moving file to trash: {} -> {}: {}",
|
||||
file_path.display(), trash_file_path.display(), e);
|
||||
Err(FileRepositoryError::IoError(e))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Restaura un archivo desde la papelera a su ubicación original
|
||||
/// Restores a file from trash to its original location
|
||||
#[instrument(skip(self))]
|
||||
pub(crate) async fn _trash_restore_from_trash(&self, file_id: &str, original_path: &str) -> FileRepositoryResult<()> {
|
||||
debug!("Restaurando archivo {} a {}", file_id, original_path);
|
||||
debug!("Restoring file {} to {}", file_id, original_path);
|
||||
|
||||
// Try to get the current path from the ID mapping service
|
||||
let current_path_result = self.id_mapping_service().get_file_path(file_id).await;
|
||||
|
||||
match current_path_result {
|
||||
Ok(current_path) => {
|
||||
debug!("Ruta actual en papelera: {}", current_path.display());
|
||||
debug!("Current path in trash: {}", current_path.display());
|
||||
|
||||
// Check if the file exists in the trash
|
||||
let file_exists = match fs::metadata(¤t_path).await {
|
||||
Ok(_) => {
|
||||
debug!("Archivo existe en papelera");
|
||||
debug!("File exists in trash");
|
||||
true
|
||||
},
|
||||
Err(e) => {
|
||||
debug!("Archivo no existe en papelera: {} - {}", current_path.display(), e);
|
||||
debug!("File does not exist in trash: {} - {}", current_path.display(), e);
|
||||
false
|
||||
}
|
||||
};
|
||||
|
||||
if !file_exists {
|
||||
error!("El archivo no existe físicamente en la papelera: {}", current_path.display());
|
||||
error!("The file does not physically exist in the trash: {}", current_path.display());
|
||||
return Err(FileRepositoryError::NotFound(format!("File not found in trash: {}", file_id)));
|
||||
}
|
||||
|
||||
// Parse the original path to a PathBuf
|
||||
let original_path_buf = PathBuf::from(original_path);
|
||||
debug!("Ruta original para restauración: {}", original_path_buf.display());
|
||||
debug!("Original path for restoration: {}", original_path_buf.display());
|
||||
|
||||
// Check if a file already exists at the destination
|
||||
let target_exists = fs::metadata(&original_path_buf).await.is_ok();
|
||||
if target_exists {
|
||||
debug!("Ya existe un archivo en la ruta de destino, generando ruta alternativa");
|
||||
debug!("A file already exists at the destination path, generating alternative path");
|
||||
|
||||
// Generate a unique path by adding a suffix
|
||||
// Extract filename and extension
|
||||
@@ -187,16 +187,16 @@ impl FileFsRepository {
|
||||
|
||||
// Create the alternative path
|
||||
let alternative_path = parent_dir.join(new_name);
|
||||
debug!("Ruta alternativa para restauración: {}", alternative_path.display());
|
||||
debug!("Alternative path for restoration: {}", alternative_path.display());
|
||||
|
||||
// Ensure the parent directory exists
|
||||
if let Some(parent) = alternative_path.parent() {
|
||||
if !parent.exists() {
|
||||
debug!("Creando directorio padre para restauración: {}", parent.display());
|
||||
debug!("Creating parent directory for restoration: {}", parent.display());
|
||||
match fs::create_dir_all(parent).await {
|
||||
Ok(_) => debug!("Directorio padre creado exitosamente"),
|
||||
Ok(_) => debug!("Parent directory created successfully"),
|
||||
Err(e) => {
|
||||
error!("Error creando directorio padre: {} - {}", parent.display(), e);
|
||||
error!("Error creating parent directory: {} - {}", parent.display(), e);
|
||||
return Err(FileRepositoryError::IoError(e));
|
||||
}
|
||||
}
|
||||
@@ -204,30 +204,30 @@ impl FileFsRepository {
|
||||
}
|
||||
|
||||
// Move the file from trash to the alternative location
|
||||
debug!("Moviendo archivo de papelera a ubicación alternativa: {} -> {}",
|
||||
debug!("Moving file from trash to alternative location: {} -> {}",
|
||||
current_path.display(), alternative_path.display());
|
||||
match fs::rename(¤t_path, &alternative_path).await {
|
||||
Ok(_) => {
|
||||
debug!("Archivo restaurado exitosamente a ubicación alternativa");
|
||||
debug!("File successfully restored to alternative location");
|
||||
|
||||
// Invalidate cache entries
|
||||
debug!("Invalidando caché para archivo en papelera");
|
||||
debug!("Invalidating cache for file in trash");
|
||||
self.metadata_cache().invalidate(¤t_path).await;
|
||||
|
||||
// Update the ID mapping
|
||||
debug!("Actualizando mapeo de ID a nueva ubicación");
|
||||
debug!("Updating ID mapping to new location");
|
||||
if let Err(e) = self.id_mapping_service().update_file_path(file_id, &alternative_path).await {
|
||||
error!("Error actualizando mapeo de archivo restaurado: {}", e);
|
||||
error!("Error updating mapping of restored file: {}", e);
|
||||
return Err(FileRepositoryError::MappingError(
|
||||
format!("Failed to update mapping: {}", e)
|
||||
));
|
||||
}
|
||||
|
||||
debug!("Restauración a ubicación alternativa completada con éxito");
|
||||
debug!("Restoration to alternative location completed successfully");
|
||||
Ok(())
|
||||
},
|
||||
Err(e) => {
|
||||
error!("Error restaurando archivo a ubicación alternativa: {}", e);
|
||||
error!("Error restoring file to alternative location: {}", e);
|
||||
Err(FileRepositoryError::IoError(e))
|
||||
}
|
||||
}
|
||||
@@ -235,11 +235,11 @@ impl FileFsRepository {
|
||||
// Ensure the parent directory exists
|
||||
if let Some(parent) = original_path_buf.parent() {
|
||||
if !parent.exists() {
|
||||
debug!("Creando directorio padre para restauración: {}", parent.display());
|
||||
debug!("Creating parent directory for restoration: {}", parent.display());
|
||||
match fs::create_dir_all(parent).await {
|
||||
Ok(_) => debug!("Directorio padre creado exitosamente"),
|
||||
Ok(_) => debug!("Parent directory created successfully"),
|
||||
Err(e) => {
|
||||
error!("Error creando directorio padre: {} - {}", parent.display(), e);
|
||||
error!("Error creating parent directory: {} - {}", parent.display(), e);
|
||||
return Err(FileRepositoryError::IoError(e));
|
||||
}
|
||||
}
|
||||
@@ -247,41 +247,41 @@ impl FileFsRepository {
|
||||
}
|
||||
|
||||
// Move the file from trash to its original location
|
||||
debug!("Moviendo archivo de papelera a ubicación original: {} -> {}",
|
||||
debug!("Moving file from trash to original location: {} -> {}",
|
||||
current_path.display(), original_path_buf.display());
|
||||
match fs::rename(¤t_path, &original_path_buf).await {
|
||||
Ok(_) => {
|
||||
debug!("Archivo restaurado exitosamente a ubicación original");
|
||||
debug!("File successfully restored to original location");
|
||||
|
||||
// Invalidate cache entries
|
||||
debug!("Invalidando caché para archivo en papelera");
|
||||
debug!("Invalidating cache for file in trash");
|
||||
self.metadata_cache().invalidate(¤t_path).await;
|
||||
|
||||
// Update the ID mapping
|
||||
debug!("Actualizando mapeo de ID a ubicación original");
|
||||
debug!("Updating ID mapping to original location");
|
||||
if let Err(e) = self.id_mapping_service().update_file_path(file_id, &original_path_buf).await {
|
||||
error!("Error actualizando mapeo de archivo restaurado: {}", e);
|
||||
error!("Error updating mapping of restored file: {}", e);
|
||||
return Err(FileRepositoryError::MappingError(
|
||||
format!("Failed to update mapping: {}", e)
|
||||
));
|
||||
}
|
||||
|
||||
debug!("Restauración a ubicación original completada con éxito");
|
||||
debug!("Restoration to original location completed successfully");
|
||||
Ok(())
|
||||
},
|
||||
Err(e) => {
|
||||
error!("Error restaurando archivo a ubicación original: {}", e);
|
||||
error!("Error restoring file to original location: {}", e);
|
||||
Err(FileRepositoryError::IoError(e))
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
Err(e) => {
|
||||
error!("Error obteniendo ruta actual del archivo {}: {:?}", file_id, e);
|
||||
error!("Error getting current path of file {}: {:?}", file_id, e);
|
||||
|
||||
// Check if the error is because the ID was not found
|
||||
if format!("{}", e).contains("not found") {
|
||||
debug!("ID no encontrado en mapeo, archivo ya no existe en papelera");
|
||||
debug!("ID not found in mapping, file no longer exists in trash");
|
||||
return Err(FileRepositoryError::NotFound(format!("File not found in trash: {}", file_id)));
|
||||
}
|
||||
|
||||
@@ -292,48 +292,48 @@ impl FileFsRepository {
|
||||
}
|
||||
}
|
||||
|
||||
/// Elimina un archivo permanentemente (usado por la papelera)
|
||||
/// Permanently deletes a file (used by trash)
|
||||
#[instrument(skip(self))]
|
||||
pub(crate) async fn _trash_delete_file_permanently(&self, file_id: &str) -> FileRepositoryResult<()> {
|
||||
debug!("Eliminando archivo permanentemente: {}", file_id);
|
||||
debug!("Permanently deleting file: {}", file_id);
|
||||
|
||||
// Get the file path using the ID mapping service
|
||||
let file_path_result = self.id_mapping_service().get_file_path(file_id).await;
|
||||
|
||||
match file_path_result {
|
||||
Ok(file_path) => {
|
||||
debug!("Encontrada ruta para archivo: {} -> {}", file_id, file_path.display());
|
||||
debug!("Found path for file: {} -> {}", file_id, file_path.display());
|
||||
|
||||
// Check if the file physically exists before attempting to delete
|
||||
let file_exists = fs::metadata(&file_path).await.is_ok();
|
||||
|
||||
if file_exists {
|
||||
debug!("Archivo existe físicamente, eliminando: {}", file_path.display());
|
||||
debug!("File exists physically, deleting: {}", file_path.display());
|
||||
|
||||
// Delete the file physically
|
||||
if let Err(e) = fs::remove_file(&file_path).await {
|
||||
error!("Error eliminando archivo permanentemente: {} - {}", file_path.display(), e);
|
||||
error!("Error permanently deleting file: {} - {}", file_path.display(), e);
|
||||
// Don't report error if the file already doesn't exist
|
||||
if e.kind() != std::io::ErrorKind::NotFound {
|
||||
return Err(FileRepositoryError::IoError(e));
|
||||
}
|
||||
} else {
|
||||
debug!("Archivo eliminado físicamente con éxito");
|
||||
debug!("File physically deleted successfully");
|
||||
}
|
||||
|
||||
// Invalidate cache for this file
|
||||
debug!("Invalidando caché para el archivo: {}", file_path.display());
|
||||
debug!("Invalidating cache for file: {}", file_path.display());
|
||||
self.metadata_cache().invalidate(&file_path).await;
|
||||
} else {
|
||||
debug!("Archivo no existe físicamente, solo limpiando mapeos: {}", file_path.display());
|
||||
debug!("File does not exist physically, only cleaning mappings: {}", file_path.display());
|
||||
}
|
||||
|
||||
// Always remove the ID mapping regardless of whether the file exists
|
||||
debug!("Eliminando mapeo de ID: {}", file_id);
|
||||
debug!("Removing ID mapping: {}", file_id);
|
||||
match self.id_mapping_service().remove_id(file_id).await {
|
||||
Ok(_) => debug!("Mapeo de ID eliminado con éxito"),
|
||||
Ok(_) => debug!("ID mapping successfully removed"),
|
||||
Err(e) => {
|
||||
error!("Error eliminando mapeo del archivo: {}", e);
|
||||
error!("Error removing file mapping: {}", e);
|
||||
// Only return error for critical mapping errors, otherwise continue
|
||||
if format!("{}", e).contains("not found") {
|
||||
debug!("ID mapping not found, ignoring this error for deletion");
|
||||
@@ -343,16 +343,16 @@ impl FileFsRepository {
|
||||
}
|
||||
};
|
||||
|
||||
debug!("Archivo eliminado permanentemente con éxito: {}", file_id);
|
||||
debug!("File permanently deleted successfully: {}", file_id);
|
||||
Ok(())
|
||||
},
|
||||
Err(e) => {
|
||||
// This could happen if the file is already deleted or wasn't properly indexed
|
||||
error!("Error obteniendo ruta del archivo {}: {:?}", file_id, e);
|
||||
error!("Error getting file path {}: {:?}", file_id, e);
|
||||
|
||||
// Check if the error is because the ID was not found
|
||||
if format!("{}", e).contains("not found") {
|
||||
debug!("ID no encontrado en mapeo, considerando borrado exitoso: {}", file_id);
|
||||
debug!("ID not found in mapping, considering deletion successful: {}", file_id);
|
||||
// In this case, we consider the file already deleted
|
||||
return Ok(());
|
||||
}
|
||||
@@ -363,5 +363,5 @@ impl FileFsRepository {
|
||||
}
|
||||
}
|
||||
|
||||
// Re-exportaciones necesarias para el compilador
|
||||
// Re-exports needed for the compiler
|
||||
use crate::domain::repositories::file_repository::FileRepositoryError;
|
||||
@@ -17,7 +17,7 @@ use crate::application::services::storage_mediator::StorageMediator;
|
||||
use crate::application::ports::outbound::FolderStoragePort;
|
||||
use crate::common::errors::DomainError;
|
||||
|
||||
// Para poder usar streams en la función list_folders
|
||||
// To be able to use streams in the list_folders function
|
||||
use tokio_stream;
|
||||
|
||||
/// Filesystem implementation of the FolderRepository interface
|
||||
@@ -79,7 +79,7 @@ impl FolderFsRepository {
|
||||
async fn count_directory_items(&self, directory_path: &Path) -> FolderRepositoryResult<usize> {
|
||||
use tokio::fs::read_dir;
|
||||
|
||||
// Timeout para evitar bloqueos
|
||||
// Timeout to avoid blocking
|
||||
let read_dir_timeout = Duration::from_secs(30);
|
||||
let read_dir_result = timeout(
|
||||
read_dir_timeout,
|
||||
@@ -91,7 +91,7 @@ impl FolderFsRepository {
|
||||
let mut entries = result.map_err(FolderRepositoryError::IoError)?;
|
||||
let mut count = 0;
|
||||
|
||||
// Contar entradas manualmente
|
||||
// Count entries manually
|
||||
while let Ok(Some(_)) = entries.next_entry().await {
|
||||
count += 1;
|
||||
}
|
||||
@@ -261,10 +261,10 @@ impl From<FolderRepositoryError> for DomainError {
|
||||
}
|
||||
}
|
||||
|
||||
// Implementar Clone para poder usar en procesamiento concurrente
|
||||
// Implement Clone to use in concurrent processing
|
||||
impl Clone for FolderFsRepository {
|
||||
fn clone(&self) -> Self {
|
||||
// Clonamos los Arc, lo que solo incrementa el contador de referencias
|
||||
// Clone the Arcs, which only increments the reference counter
|
||||
Self {
|
||||
root_path: self.root_path.clone(),
|
||||
storage_mediator: self.storage_mediator.clone(),
|
||||
|
||||
@@ -13,18 +13,18 @@ use crate::common::config::AppConfig;
|
||||
use crate::domain::repositories::file_repository::FileRepositoryError;
|
||||
use crate::infrastructure::services::buffer_pool::BufferPool;
|
||||
|
||||
/// Estructura para el rango de bytes a procesar
|
||||
/// Structure for the byte range to process
|
||||
#[derive(Debug, Clone, Copy)]
|
||||
pub struct ChunkRange {
|
||||
/// Índice del chunk
|
||||
/// Chunk index
|
||||
pub index: usize,
|
||||
/// Posición de inicio en bytes
|
||||
/// Start position in bytes
|
||||
pub start: u64,
|
||||
/// Tamaño del chunk en bytes
|
||||
/// Chunk size in bytes
|
||||
pub size: usize,
|
||||
}
|
||||
|
||||
/// Buffer pooling específico para BytesMut
|
||||
/// Specific buffer pooling for BytesMut
|
||||
pub struct BytesBufferPool {
|
||||
buffers: Mutex<Vec<BytesMut>>,
|
||||
buffer_size: usize,
|
||||
@@ -40,53 +40,53 @@ impl BytesBufferPool {
|
||||
}
|
||||
}
|
||||
|
||||
/// Obtener un buffer del pool o crear uno nuevo
|
||||
/// Get a buffer from the pool or create a new one
|
||||
pub async fn get_buffer(&self) -> BytesMut {
|
||||
let mut buffers = self.buffers.lock().await;
|
||||
|
||||
if let Some(mut buffer) = buffers.pop() {
|
||||
// Reutilizar buffer existente
|
||||
buffer.clear(); // Mantener capacidad, limpiar contenido
|
||||
// Reuse existing buffer
|
||||
buffer.clear(); // Keep capacity, clear content
|
||||
buffer
|
||||
} else {
|
||||
// Crear nuevo buffer si el pool está vacío
|
||||
// Create new buffer if the pool is empty
|
||||
BytesMut::with_capacity(self.buffer_size)
|
||||
}
|
||||
}
|
||||
|
||||
/// Devolver un buffer al pool para reutilización
|
||||
/// Return a buffer to the pool for reuse
|
||||
pub async fn return_buffer(&self, mut buffer: BytesMut) {
|
||||
// Restablece el buffer para reutilización
|
||||
// Reset the buffer for reuse
|
||||
buffer.clear();
|
||||
|
||||
let mut buffers = self.buffers.lock().await;
|
||||
|
||||
// Solo mantener hasta max_buffers
|
||||
// Only keep up to max_buffers
|
||||
if buffers.len() < self.max_buffers {
|
||||
buffers.push(buffer);
|
||||
}
|
||||
// Si ya tenemos suficientes buffers, este se descartará
|
||||
// If we already have enough buffers, this one will be discarded
|
||||
}
|
||||
}
|
||||
|
||||
/// Procesador paralelo de archivos para operaciones IO intensivas
|
||||
/// Parallel file processor for IO-intensive operations
|
||||
pub struct ParallelFileProcessor {
|
||||
/// Configuración de la aplicación
|
||||
/// Application configuration
|
||||
config: AppConfig,
|
||||
/// Semáforo para limitar concurrencia global
|
||||
/// Semaphore to limit global concurrency
|
||||
concurrency_limiter: Arc<Semaphore>,
|
||||
/// Pool de buffers para optimizar memoria
|
||||
/// Buffer pool to optimize memory
|
||||
buffer_pool: Option<Arc<BufferPool>>,
|
||||
/// Pool de buffers BytesMut para operaciones zero-copy
|
||||
/// BytesMut buffer pool for zero-copy operations
|
||||
bytes_pool: Arc<BytesBufferPool>,
|
||||
}
|
||||
|
||||
impl ParallelFileProcessor {
|
||||
/// Crea una nueva instancia del procesador
|
||||
/// Creates a new processor instance
|
||||
pub fn new(config: AppConfig) -> Self {
|
||||
let concurrency_limiter = Arc::new(Semaphore::new(config.concurrency.max_concurrent_io));
|
||||
|
||||
// Crear pool de BytesMut para operaciones eficientes
|
||||
// Create BytesMut pool for efficient operations
|
||||
let chunk_size = config.resources.chunk_size_bytes;
|
||||
let max_chunks = config.concurrency.max_parallel_chunks;
|
||||
let bytes_pool = Arc::new(BytesBufferPool::new(chunk_size, max_chunks * 2));
|
||||
@@ -99,11 +99,11 @@ impl ParallelFileProcessor {
|
||||
}
|
||||
}
|
||||
|
||||
/// Crea una nueva instancia del procesador con un pool de buffers
|
||||
/// Creates a new processor instance with a buffer pool
|
||||
pub fn new_with_buffer_pool(config: AppConfig, buffer_pool: Arc<BufferPool>) -> Self {
|
||||
let concurrency_limiter = Arc::new(Semaphore::new(config.concurrency.max_concurrent_io));
|
||||
|
||||
// Crear pool de BytesMut para operaciones eficientes
|
||||
// Create BytesMut pool for efficient operations
|
||||
let chunk_size = config.resources.chunk_size_bytes;
|
||||
let max_chunks = config.concurrency.max_parallel_chunks;
|
||||
let bytes_pool = Arc::new(BytesBufferPool::new(chunk_size, max_chunks * 2));
|
||||
@@ -116,15 +116,15 @@ impl ParallelFileProcessor {
|
||||
}
|
||||
}
|
||||
|
||||
/// Divide un archivo en chunks para procesamiento paralelo
|
||||
/// Divides a file into chunks for parallel processing
|
||||
pub fn calculate_chunks(&self, file_size: u64) -> Vec<ChunkRange> {
|
||||
// Determinar si el archivo necesita procesamiento paralelo
|
||||
// Determine if the file needs parallel processing
|
||||
let needs_parallel = self.config.resources.needs_parallel_processing(
|
||||
file_size, &self.config.concurrency
|
||||
);
|
||||
|
||||
if !needs_parallel {
|
||||
// Para archivos pequeños, usar un solo chunk
|
||||
// For small files, use a single chunk
|
||||
return vec![ChunkRange {
|
||||
index: 0,
|
||||
start: 0,
|
||||
@@ -132,21 +132,21 @@ impl ParallelFileProcessor {
|
||||
}];
|
||||
}
|
||||
|
||||
// Calcular número óptimo de chunks
|
||||
// Calculate optimal number of chunks
|
||||
let chunk_count = self.config.resources.calculate_optimal_chunks(
|
||||
file_size, &self.config.concurrency
|
||||
);
|
||||
|
||||
// Calcular tamaño de cada chunk
|
||||
// Calculate size of each chunk
|
||||
let chunk_size = self.config.resources.calculate_chunk_size(file_size, chunk_count);
|
||||
|
||||
// Crear los rangos de chunks
|
||||
// Create chunk ranges
|
||||
let mut chunks = Vec::with_capacity(chunk_count);
|
||||
|
||||
let mut start = 0;
|
||||
for i in 0..chunk_count {
|
||||
let current_chunk_size = if i == chunk_count - 1 {
|
||||
// Último chunk puede ser más pequeño
|
||||
// Last chunk might be smaller
|
||||
(file_size - start) as usize
|
||||
} else {
|
||||
chunk_size
|
||||
@@ -167,16 +167,16 @@ impl ParallelFileProcessor {
|
||||
chunks
|
||||
}
|
||||
|
||||
/// Lee un archivo en paralelo y devuelve el contenido completo
|
||||
/// Implementación optimizada usando BytesMut para reducir copias de memoria
|
||||
/// Reads a file in parallel and returns the complete content
|
||||
/// Optimized implementation using BytesMut to reduce memory copies
|
||||
pub async fn read_file_parallel(&self, file_path: &PathBuf) -> Result<Vec<u8>, FileRepositoryError> {
|
||||
// Obtener tamaño del archivo
|
||||
// Get file size
|
||||
let metadata = tokio::fs::metadata(file_path).await
|
||||
.map_err(FileRepositoryError::IoError)?;
|
||||
|
||||
let file_size = metadata.len();
|
||||
|
||||
// Verificar si el archivo es demasiado grande para memoria
|
||||
// Check if the file is too large for memory
|
||||
if !self.config.resources.can_load_in_memory(file_size) {
|
||||
return Err(FileRepositoryError::Other(
|
||||
format!("File too large to load in memory: {} MB (max: {} MB)",
|
||||
@@ -185,19 +185,19 @@ impl ParallelFileProcessor {
|
||||
));
|
||||
}
|
||||
|
||||
// Calcular chunks
|
||||
// Calculate chunks
|
||||
let chunks = self.calculate_chunks(file_size);
|
||||
|
||||
if chunks.len() == 1 {
|
||||
// Para un solo chunk, usar lectura simple con buffer pool si está disponible
|
||||
// For a single chunk, use simple reading with buffer pool if available
|
||||
info!("Reading file with size {}MB as a single chunk", file_size / (1024 * 1024));
|
||||
|
||||
if let Some(pool) = &self.buffer_pool {
|
||||
// Usar buffer del pool para lectura eficiente
|
||||
// Use buffer from the pool for efficient reading
|
||||
debug!("Using buffer pool for single chunk read");
|
||||
let mut buffer = pool.get_buffer().await;
|
||||
|
||||
// Si el buffer es demasiado pequeño, revertir a la implementación estándar
|
||||
// If the buffer is too small, revert to standard implementation
|
||||
if buffer.capacity() < file_size as usize {
|
||||
debug!("Buffer from pool too small ({}), using standard read", buffer.capacity());
|
||||
let content = tokio::fs::read(file_path).await
|
||||
@@ -206,7 +206,7 @@ impl ParallelFileProcessor {
|
||||
return Ok(content);
|
||||
}
|
||||
|
||||
// Usar el buffer de memoria del pool
|
||||
// Use memory buffer from the pool
|
||||
let mut file = File::open(file_path).await
|
||||
.map_err(FileRepositoryError::IoError)?;
|
||||
|
||||
@@ -215,11 +215,11 @@ impl ParallelFileProcessor {
|
||||
|
||||
buffer.set_used(read_size);
|
||||
|
||||
// Convertir en Vec<u8>
|
||||
// Convert to Vec<u8>
|
||||
let content = buffer.into_vec();
|
||||
return Ok(content);
|
||||
} else {
|
||||
// Implementación estándar sin pool
|
||||
// Standard implementation without pool
|
||||
let content = tokio::fs::read(file_path).await
|
||||
.map_err(FileRepositoryError::IoError)?;
|
||||
|
||||
@@ -227,51 +227,51 @@ impl ParallelFileProcessor {
|
||||
}
|
||||
}
|
||||
|
||||
// Para múltiples chunks, usar lectura paralela
|
||||
// For multiple chunks, use parallel reading
|
||||
info!("Reading file with size {}MB in {} parallel chunks using BytesMut",
|
||||
file_size / (1024 * 1024), chunks.len());
|
||||
|
||||
// Crear buffer de resultado final (pre-allocated)
|
||||
// Create final result buffer (pre-allocated)
|
||||
let mut result = BytesMut::with_capacity(file_size as usize);
|
||||
result.resize(file_size as usize, 0);
|
||||
let result_mutex = Arc::new(Mutex::new(result));
|
||||
|
||||
// Crear tareas para cada chunk
|
||||
// Create tasks for each chunk
|
||||
let mut tasks = Vec::with_capacity(chunks.len());
|
||||
|
||||
// Abrir archivo una sola vez y compartirlo
|
||||
// Open file once and share it
|
||||
let file = Arc::new(File::open(file_path).await
|
||||
.map_err(FileRepositoryError::IoError)?);
|
||||
|
||||
// Referencia al pool de BytesMut
|
||||
// Reference to BytesMut pool
|
||||
let bytes_pool = self.bytes_pool.clone();
|
||||
|
||||
// Procesar chunks en paralelo
|
||||
// Process chunks in parallel
|
||||
for chunk in chunks {
|
||||
let file_clone = file.clone();
|
||||
let result_clone = result_mutex.clone();
|
||||
let semaphore_clone = self.concurrency_limiter.clone();
|
||||
let bytes_pool_clone = bytes_pool.clone();
|
||||
|
||||
// Spawn task para este chunk - no hay necesidad de copiar los datos originales
|
||||
// Spawn task for this chunk - no need to copy the original data
|
||||
let task = task::spawn(async move {
|
||||
// Adquirir permiso del semáforo
|
||||
// Acquire semaphore permit
|
||||
let _permit = semaphore_clone.acquire().await.unwrap();
|
||||
|
||||
// Obtener un buffer reusable del pool de BytesMut
|
||||
// Get a reusable buffer from the BytesMut pool
|
||||
let mut chunk_buffer = bytes_pool_clone.get_buffer().await;
|
||||
|
||||
// Asegurar que tenga suficiente capacidad
|
||||
// Ensure it has sufficient capacity
|
||||
if chunk_buffer.capacity() < chunk.size {
|
||||
chunk_buffer = BytesMut::with_capacity(chunk.size);
|
||||
}
|
||||
// Resize al tamaño exacto necesario
|
||||
// Resize to the exact size needed
|
||||
chunk_buffer.resize(chunk.size, 0);
|
||||
|
||||
// Crear un descriptor de archivo duplicado para uso independiente
|
||||
// Create a duplicate file descriptor for independent use
|
||||
let mut file_handle = file_clone.try_clone().await?;
|
||||
|
||||
// Posicionar y leer directamente en el BytesMut
|
||||
// Position and read directly into the BytesMut
|
||||
file_handle.seek(SeekFrom::Start(chunk.start)).await?;
|
||||
let bytes_read = file_handle.read_exact(&mut chunk_buffer[..chunk.size]).await?;
|
||||
|
||||
@@ -282,18 +282,18 @@ impl ParallelFileProcessor {
|
||||
));
|
||||
}
|
||||
|
||||
// Escribir en resultado final
|
||||
// Write to final result
|
||||
let mut result_lock = result_clone.lock().await;
|
||||
let start_pos = chunk.start as usize;
|
||||
let end_pos = start_pos + chunk.size;
|
||||
|
||||
// Usar copy_from_slice para copiar desde BytesMut al buffer de resultado
|
||||
// Use copy_from_slice to copy from BytesMut to result buffer
|
||||
result_lock[start_pos..end_pos].copy_from_slice(&chunk_buffer[..chunk.size]);
|
||||
|
||||
// Devolver el buffer al pool para su reutilización
|
||||
// Return the buffer to the pool for reuse
|
||||
bytes_pool_clone.return_buffer(chunk_buffer).await;
|
||||
|
||||
// Registrar progreso
|
||||
// Log progress
|
||||
debug!("Chunk {} processed: {} bytes from offset {}",
|
||||
chunk.index, chunk.size, chunk.start);
|
||||
|
||||
@@ -303,10 +303,10 @@ impl ParallelFileProcessor {
|
||||
tasks.push(task);
|
||||
}
|
||||
|
||||
// Esperar a que todas las tareas terminen
|
||||
// Wait for all tasks to complete
|
||||
let results = join_all(tasks).await;
|
||||
|
||||
// Verificar errores
|
||||
// Check for errors
|
||||
for (i, task_result) in results.into_iter().enumerate() {
|
||||
match task_result {
|
||||
Ok(Ok(())) => {},
|
||||
@@ -321,7 +321,7 @@ impl ParallelFileProcessor {
|
||||
}
|
||||
}
|
||||
|
||||
// Obtener el resultado final y convertir a Vec<u8>
|
||||
// Get the final result and convert to Vec<u8>
|
||||
let result_buffer = result_mutex.lock().await;
|
||||
let result_vec = result_buffer.to_vec();
|
||||
|
||||
@@ -329,8 +329,8 @@ impl ParallelFileProcessor {
|
||||
Ok(result_vec)
|
||||
}
|
||||
|
||||
/// Escribe un archivo en paralelo desde un buffer
|
||||
/// Implementación optimizada usando BytesMut/Bytes para reducir copias de memoria
|
||||
/// Writes a file in parallel from a buffer
|
||||
/// Optimized implementation using BytesMut/Bytes to reduce memory copies
|
||||
pub async fn write_file_parallel(
|
||||
&self,
|
||||
file_path: &PathBuf,
|
||||
@@ -338,56 +338,56 @@ impl ParallelFileProcessor {
|
||||
) -> Result<(), FileRepositoryError> {
|
||||
let file_size = content.len() as u64;
|
||||
|
||||
// Calcular chunks
|
||||
// Calculate chunks
|
||||
let chunks = self.calculate_chunks(file_size);
|
||||
|
||||
if chunks.len() == 1 {
|
||||
// Para un solo chunk, usar escritura simple
|
||||
// For a single chunk, use simple writing
|
||||
info!("Writing file with size {}MB as a single chunk", file_size / (1024 * 1024));
|
||||
|
||||
// Implementación estándar (el buffer pooling no ofrece ventajas para escritura simple)
|
||||
// Standard implementation (buffer pooling offers no advantages for simple writing)
|
||||
tokio::fs::write(file_path, content).await
|
||||
.map_err(FileRepositoryError::IoError)?;
|
||||
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
// Para múltiples chunks, usar escritura paralela
|
||||
// For multiple chunks, use parallel writing
|
||||
info!("Writing file with size {}MB in {} parallel chunks using Bytes",
|
||||
file_size / (1024 * 1024), chunks.len());
|
||||
|
||||
// Crear archivo (no usamos Mutex para reducir contención)
|
||||
// Create file (we don't use Mutex to reduce contention)
|
||||
let file = File::create(file_path).await
|
||||
.map_err(FileRepositoryError::IoError)?;
|
||||
|
||||
// Convertir contenido a Bytes (un solo paso de copia)
|
||||
// Convert content to Bytes (single copy step)
|
||||
let content_bytes = Bytes::copy_from_slice(content);
|
||||
|
||||
// Crear tareas para cada chunk
|
||||
// Create tasks for each chunk
|
||||
let mut tasks = Vec::with_capacity(chunks.len());
|
||||
|
||||
// Procesar chunks en paralelo
|
||||
// Process chunks in parallel
|
||||
for chunk in chunks {
|
||||
let file_clone = file.try_clone().await
|
||||
.map_err(FileRepositoryError::IoError)?;
|
||||
let semaphore_clone = self.concurrency_limiter.clone();
|
||||
|
||||
// Crear slice de Bytes (no copia datos, solo referencia)
|
||||
// Create Bytes slice (doesn't copy data, only references)
|
||||
let start_idx = chunk.start as usize;
|
||||
let end_idx = start_idx + chunk.size;
|
||||
let chunk_data = content_bytes.slice(start_idx..end_idx);
|
||||
|
||||
// Crear y lanzar tarea
|
||||
// Create and launch task
|
||||
let task = task::spawn(async move {
|
||||
// Adquirir permiso del semáforo
|
||||
// Acquire semaphore permit
|
||||
let _permit = semaphore_clone.acquire().await.unwrap();
|
||||
|
||||
// Posicionar y escribir
|
||||
// Position and write
|
||||
let mut file_handle = file_clone;
|
||||
file_handle.seek(SeekFrom::Start(chunk.start)).await?;
|
||||
file_handle.write_all(&chunk_data).await?;
|
||||
|
||||
// Registrar progreso
|
||||
// Log progress
|
||||
debug!("Chunk {} written: {} bytes at offset {}",
|
||||
chunk.index, chunk.size, chunk.start);
|
||||
|
||||
@@ -397,10 +397,10 @@ impl ParallelFileProcessor {
|
||||
tasks.push(task);
|
||||
}
|
||||
|
||||
// Esperar a que todas las tareas terminen
|
||||
// Wait for all tasks to complete
|
||||
let results = join_all(tasks).await;
|
||||
|
||||
// Verificar errores
|
||||
// Check for errors
|
||||
for (i, task_result) in results.into_iter().enumerate() {
|
||||
match task_result {
|
||||
Ok(Ok(())) => {},
|
||||
@@ -415,7 +415,7 @@ impl ParallelFileProcessor {
|
||||
}
|
||||
}
|
||||
|
||||
// Garantizar que todo se ha escrito correctamente
|
||||
// Ensure everything has been written correctly
|
||||
let mut file_handle = file;
|
||||
file_handle.flush().await.map_err(FileRepositoryError::IoError)?;
|
||||
|
||||
@@ -423,17 +423,17 @@ impl ParallelFileProcessor {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Escribe un chunk en un archivo en una posición específica
|
||||
/// Writes a chunk to a file at a specific position
|
||||
#[allow(dead_code)]
|
||||
async fn write_chunk_optimized(
|
||||
file: &mut File,
|
||||
offset: u64,
|
||||
data: Bytes
|
||||
) -> Result<(), std::io::Error> {
|
||||
// Preparar la escritura en la posición correcta
|
||||
// Prepare writing at the correct position
|
||||
file.seek(SeekFrom::Start(offset)).await?;
|
||||
|
||||
// Escribir datos sin copias adicionales
|
||||
// Write data without additional copies
|
||||
file.write_all(&data).await?;
|
||||
|
||||
Ok(())
|
||||
@@ -447,59 +447,59 @@ mod tests {
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_parallel_read_write() {
|
||||
// Crear configuración con umbral bajo para testing
|
||||
// Create configuration with low threshold for testing
|
||||
let mut config = AppConfig::default();
|
||||
config.concurrency.min_size_for_parallel_chunks_mb = 1; // 1MB para testing
|
||||
config.concurrency.min_size_for_parallel_chunks_mb = 1; // 1MB for testing
|
||||
config.concurrency.max_parallel_chunks = 4;
|
||||
|
||||
let processor = ParallelFileProcessor::new(config);
|
||||
|
||||
// Crear directorio temporal
|
||||
// Create temporary directory
|
||||
let temp_dir = tempdir().unwrap();
|
||||
let file_path = temp_dir.path().join("test_file.bin");
|
||||
|
||||
// Crear datos de prueba (2MB)
|
||||
// Create test data (2MB)
|
||||
let size = 2 * 1024 * 1024;
|
||||
let mut test_data = Vec::with_capacity(size);
|
||||
for i in 0..size {
|
||||
test_data.push((i % 256) as u8);
|
||||
}
|
||||
|
||||
// Escribir archivo en paralelo
|
||||
// Write file in parallel
|
||||
processor.write_file_parallel(&file_path, &test_data).await.unwrap();
|
||||
|
||||
// Leer archivo en paralelo
|
||||
// Read file in parallel
|
||||
let read_data = processor.read_file_parallel(&file_path).await.unwrap();
|
||||
|
||||
// Verificar que los datos son idénticos
|
||||
// Verify that the data is identical
|
||||
assert_eq!(test_data.len(), read_data.len());
|
||||
assert_eq!(test_data, read_data);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_bytesmut_pool() {
|
||||
// Crear pool
|
||||
// Create pool
|
||||
let pool = BytesBufferPool::new(1024, 5);
|
||||
|
||||
// Obtener buffer
|
||||
// Get buffer
|
||||
let mut buffer1 = pool.get_buffer().await;
|
||||
buffer1.put_slice(b"test data");
|
||||
assert_eq!(&buffer1[..9], b"test data");
|
||||
|
||||
// Devolver buffer al pool
|
||||
// Return buffer to the pool
|
||||
pool.return_buffer(buffer1).await;
|
||||
|
||||
// Obtener otro buffer (debería ser el mismo)
|
||||
// Get another buffer (should be the same one)
|
||||
let buffer2 = pool.get_buffer().await;
|
||||
assert_eq!(buffer2.capacity(), 1024);
|
||||
|
||||
// El buffer debería estar vacío (clear)
|
||||
// The buffer should be empty (cleared)
|
||||
assert_eq!(buffer2.len(), 0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_chunk_calculation() {
|
||||
// Crear configuración de prueba
|
||||
// Create test configuration
|
||||
let mut config = AppConfig::default();
|
||||
config.concurrency.min_size_for_parallel_chunks_mb = 100; // 100MB
|
||||
config.concurrency.max_parallel_chunks = 4;
|
||||
@@ -507,18 +507,18 @@ mod tests {
|
||||
|
||||
let processor = ParallelFileProcessor::new(config);
|
||||
|
||||
// Archivo pequeño (10MB)
|
||||
// Small file (10MB)
|
||||
let small_file_size = 10 * 1024 * 1024;
|
||||
let chunks = processor.calculate_chunks(small_file_size);
|
||||
assert_eq!(chunks.len(), 1);
|
||||
assert_eq!(chunks[0].size as u64, small_file_size);
|
||||
|
||||
// Archivo grande (300MB)
|
||||
// Large file (300MB)
|
||||
let large_file_size = 300 * 1024 * 1024;
|
||||
let chunks = processor.calculate_chunks(large_file_size);
|
||||
assert_eq!(chunks.len(), 4); // Limitado a max_parallel_chunks
|
||||
assert_eq!(chunks.len(), 4); // Limited to max_parallel_chunks
|
||||
|
||||
// Verificar que todos los chunks suman el tamaño total
|
||||
// Verify that all chunks add up to the total size
|
||||
let total_size: u64 = chunks.iter().map(|c| c.size as u64).sum();
|
||||
assert_eq!(total_size, large_file_size);
|
||||
}
|
||||
|
||||
@@ -10,60 +10,60 @@ use crate::infrastructure::services::id_mapping_service::{IdMappingService, IdMa
|
||||
use crate::common::errors::DomainError;
|
||||
use crate::application::ports::outbound::IdMappingPort;
|
||||
|
||||
/// Tamaño máximo de entradas en el caché
|
||||
/// Maximum number of entries in the cache
|
||||
const MAX_CACHE_SIZE: usize = 10_000;
|
||||
|
||||
/// Tiempo de vida del caché (en segundos)
|
||||
const CACHE_TTL_SECONDS: u64 = 60 * 5; // 5 minutos
|
||||
/// Cache time-to-live (in seconds)
|
||||
const CACHE_TTL_SECONDS: u64 = 60 * 5; // 5 minutes
|
||||
|
||||
/// Optimizador para operaciones masivas de mapeo de IDs
|
||||
/// Optimizer for batch ID mapping operations
|
||||
pub struct IdMappingOptimizer {
|
||||
/// Servicio base de mapeo de IDs
|
||||
/// Base ID mapping service
|
||||
base_service: Arc<IdMappingService>,
|
||||
|
||||
/// Caché de ID por ruta (path -> id)
|
||||
/// Path to ID cache (path -> id)
|
||||
path_to_id_cache: RwLock<HashMap<String, (String, Instant)>>,
|
||||
|
||||
/// Caché de ruta por ID (id -> path)
|
||||
/// ID to path cache (id -> path)
|
||||
id_to_path_cache: RwLock<HashMap<String, (String, Instant)>>,
|
||||
|
||||
/// Contador de hits
|
||||
/// Hit counter
|
||||
stats: RwLock<OptimizerStats>,
|
||||
|
||||
/// Semáforo para limitar operaciones de batch
|
||||
/// Semaphore to limit batch operations
|
||||
batch_limiter: Semaphore,
|
||||
|
||||
/// Cola de batch pendientes
|
||||
/// Pending batch queue
|
||||
pending_batch: Mutex<BatchQueue>,
|
||||
}
|
||||
|
||||
/// Estadísticas del optimizador
|
||||
/// Optimizer statistics
|
||||
#[derive(Debug, Default, Clone)]
|
||||
pub struct OptimizerStats {
|
||||
/// Número total de consultas get_path_by_id
|
||||
/// Total number of get_path_by_id queries
|
||||
pub path_by_id_queries: usize,
|
||||
/// Número de hits en caché get_path_by_id
|
||||
/// Number of cache hits for get_path_by_id
|
||||
pub path_by_id_hits: usize,
|
||||
|
||||
/// Número total de consultas get_or_create_id
|
||||
/// Total number of get_or_create_id queries
|
||||
pub get_id_queries: usize,
|
||||
/// Número de hits en caché get_or_create_id
|
||||
/// Number of cache hits for get_or_create_id
|
||||
pub get_id_hits: usize,
|
||||
|
||||
/// Número de batch realizados
|
||||
/// Number of batch operations performed
|
||||
pub batch_operations: usize,
|
||||
/// Número total de IDs procesados en batch
|
||||
/// Total number of IDs processed in batch
|
||||
pub batch_items_processed: usize,
|
||||
|
||||
/// Último momento de limpieza de caché
|
||||
/// Last cache cleanup timestamp
|
||||
pub last_cleanup: Option<Instant>,
|
||||
}
|
||||
|
||||
/// Cola para operaciones batch
|
||||
/// Queue for batch operations
|
||||
struct BatchQueue {
|
||||
/// Rutas pendientes para obtener/crear ID
|
||||
/// Pending paths to get/create ID
|
||||
path_to_id_requests: HashSet<String>,
|
||||
/// IDs pendientes para obtener ruta
|
||||
/// Pending IDs to get path
|
||||
id_to_path_requests: HashSet<String>,
|
||||
}
|
||||
|
||||
@@ -76,43 +76,43 @@ impl Default for BatchQueue {
|
||||
}
|
||||
}
|
||||
|
||||
/// Resultado de una operación batch
|
||||
/// Result of a batch operation
|
||||
struct BatchResult {
|
||||
/// Mapeo de ruta a ID
|
||||
/// Path to ID mapping
|
||||
path_to_id: HashMap<String, String>,
|
||||
/// Mapeo de ID a ruta
|
||||
/// ID to path mapping
|
||||
id_to_path: HashMap<String, String>,
|
||||
}
|
||||
|
||||
impl IdMappingOptimizer {
|
||||
/// Crea un nuevo optimizador para el servicio de mapeo de IDs
|
||||
/// Creates a new optimizer for the ID mapping service
|
||||
pub fn new(base_service: Arc<IdMappingService>) -> Self {
|
||||
Self {
|
||||
base_service,
|
||||
path_to_id_cache: RwLock::new(HashMap::with_capacity(1000)),
|
||||
id_to_path_cache: RwLock::new(HashMap::with_capacity(1000)),
|
||||
stats: RwLock::new(OptimizerStats::default()),
|
||||
batch_limiter: Semaphore::new(2), // Limitar a 2 operaciones batch concurrentes
|
||||
batch_limiter: Semaphore::new(2), // Limit to 2 concurrent batch operations
|
||||
pending_batch: Mutex::new(BatchQueue::default()),
|
||||
}
|
||||
}
|
||||
|
||||
/// Obtiene estadísticas del optimizador
|
||||
/// Gets optimizer statistics
|
||||
pub async fn get_stats(&self) -> OptimizerStats {
|
||||
self.stats.read().await.clone()
|
||||
}
|
||||
|
||||
/// Limpia entradas expiradas del caché
|
||||
/// Cleans expired cache entries
|
||||
pub async fn cleanup_cache(&self) {
|
||||
let now = Instant::now();
|
||||
let ttl = Duration::from_secs(CACHE_TTL_SECONDS);
|
||||
|
||||
// Limpiar caché path_to_id
|
||||
// Clean path_to_id cache
|
||||
{
|
||||
let mut cache = self.path_to_id_cache.write().await;
|
||||
let initial_size = cache.len();
|
||||
|
||||
// Retener solo entradas no expiradas
|
||||
// Retain only non-expired entries
|
||||
cache.retain(|_, (_, timestamp)| {
|
||||
now.duration_since(*timestamp) < ttl
|
||||
});
|
||||
@@ -123,12 +123,12 @@ impl IdMappingOptimizer {
|
||||
}
|
||||
}
|
||||
|
||||
// Limpiar caché id_to_path
|
||||
// Clean id_to_path cache
|
||||
{
|
||||
let mut cache = self.id_to_path_cache.write().await;
|
||||
let initial_size = cache.len();
|
||||
|
||||
// Retener solo entradas no expiradas
|
||||
// Retain only non-expired entries
|
||||
cache.retain(|_, (_, timestamp)| {
|
||||
now.duration_since(*timestamp) < ttl
|
||||
});
|
||||
@@ -139,14 +139,14 @@ impl IdMappingOptimizer {
|
||||
}
|
||||
}
|
||||
|
||||
// Actualizar estadísticas
|
||||
// Update statistics
|
||||
{
|
||||
let mut stats = self.stats.write().await;
|
||||
stats.last_cleanup = Some(now);
|
||||
}
|
||||
}
|
||||
|
||||
/// Inicia tarea de limpieza periódica
|
||||
/// Starts periodic cleanup task
|
||||
pub fn start_cleanup_task(optimizer: Arc<Self>) {
|
||||
tokio::spawn(async move {
|
||||
let cleanup_interval = Duration::from_secs(CACHE_TTL_SECONDS / 2);
|
||||
@@ -155,7 +155,7 @@ impl IdMappingOptimizer {
|
||||
tokio::time::sleep(cleanup_interval).await;
|
||||
optimizer.cleanup_cache().await;
|
||||
|
||||
// Loguear estadísticas periódicamente
|
||||
// Log statistics periodically
|
||||
let stats = optimizer.get_stats().await;
|
||||
info!("ID Mapping Optimizer stats - Path queries: {}, hits: {} ({}%), ID queries: {}, hits: {} ({}%), Batch ops: {}, items: {}",
|
||||
stats.path_by_id_queries,
|
||||
@@ -171,11 +171,11 @@ impl IdMappingOptimizer {
|
||||
});
|
||||
}
|
||||
|
||||
/// Agrega una solicitud a la cola pendiente para procesamiento batch
|
||||
/// Adds a request to the pending queue for batch processing
|
||||
async fn queue_path_to_id_request(&self, path: &StoragePath) -> Result<Option<String>, IdMappingError> {
|
||||
let path_str = path.to_string();
|
||||
|
||||
// Verificar primero en el caché
|
||||
// Check first in the cache
|
||||
{
|
||||
let cache = self.path_to_id_cache.read().await;
|
||||
if let Some((id, _)) = cache.get(&path_str) {
|
||||
@@ -204,7 +204,7 @@ impl IdMappingOptimizer {
|
||||
// Adquirir permiso para operación batch
|
||||
let _permit = self.batch_limiter.acquire().await.unwrap();
|
||||
|
||||
// Obtener las solicitudes pendientes
|
||||
// Get pending requests
|
||||
let (path_requests, id_requests) = {
|
||||
let mut batch_queue = self.pending_batch.lock().await;
|
||||
|
||||
@@ -250,7 +250,7 @@ impl IdMappingOptimizer {
|
||||
}
|
||||
}
|
||||
|
||||
// Actualizar caché con los resultados del batch
|
||||
// Update cache with batch results
|
||||
{
|
||||
let mut path_cache = self.path_to_id_cache.write().await;
|
||||
let mut id_cache = self.id_to_path_cache.write().await;
|
||||
@@ -397,7 +397,7 @@ impl IdMappingPort for IdMappingOptimizer {
|
||||
{
|
||||
let cache = self.path_to_id_cache.read().await;
|
||||
if let Some((id, _)) = cache.get(&path_str) {
|
||||
// Actualizar estadísticas
|
||||
// Update statistics
|
||||
{
|
||||
let mut stats = self.stats.write().await;
|
||||
stats.get_id_hits += 1;
|
||||
@@ -407,7 +407,7 @@ impl IdMappingPort for IdMappingOptimizer {
|
||||
}
|
||||
}
|
||||
|
||||
// Si no está en caché, intentar agregar a cola de batch primero
|
||||
// If not in cache, try adding to batch queue first
|
||||
let queued_result = self.queue_path_to_id_request(path).await?;
|
||||
if let Some(id) = queued_result {
|
||||
return Ok(id);
|
||||
@@ -416,17 +416,17 @@ impl IdMappingPort for IdMappingOptimizer {
|
||||
// Trigger batch processing if enough items accumulated
|
||||
self.trigger_batch_if_needed(20).await?;
|
||||
|
||||
// Intentar obtener del servicio base
|
||||
// Try to get from the base service
|
||||
let id = self.base_service.get_or_create_id(path).await?;
|
||||
|
||||
// Actualizar caché con el nuevo ID
|
||||
// Update cache with the new ID
|
||||
{
|
||||
let mut path_cache = self.path_to_id_cache.write().await;
|
||||
let mut id_cache = self.id_to_path_cache.write().await;
|
||||
|
||||
let now = Instant::now();
|
||||
|
||||
// Controlar tamaño del caché
|
||||
// Control cache size
|
||||
if path_cache.len() >= MAX_CACHE_SIZE {
|
||||
warn!("Path-to-ID cache size reached limit ({}), clearing oldest entries", MAX_CACHE_SIZE);
|
||||
path_cache.clear();
|
||||
@@ -451,11 +451,11 @@ impl IdMappingPort for IdMappingOptimizer {
|
||||
stats.path_by_id_queries += 1;
|
||||
}
|
||||
|
||||
// Verificar primero en el caché
|
||||
// Check first in the cache
|
||||
{
|
||||
let cache = self.id_to_path_cache.read().await;
|
||||
if let Some((path_str, _)) = cache.get(id) {
|
||||
// Actualizar estadísticas
|
||||
// Update statistics
|
||||
{
|
||||
let mut stats = self.stats.write().await;
|
||||
stats.path_by_id_hits += 1;
|
||||
@@ -465,10 +465,10 @@ impl IdMappingPort for IdMappingOptimizer {
|
||||
}
|
||||
}
|
||||
|
||||
// Obtener del servicio base
|
||||
// Get from the base service
|
||||
let path = self.base_service.get_path_by_id(id).await?;
|
||||
|
||||
// Actualizar caché
|
||||
// Update cache
|
||||
{
|
||||
let mut id_cache = self.id_to_path_cache.write().await;
|
||||
let mut path_cache = self.path_to_id_cache.write().await;
|
||||
@@ -476,7 +476,7 @@ impl IdMappingPort for IdMappingOptimizer {
|
||||
let now = Instant::now();
|
||||
let path_str = path.to_string();
|
||||
|
||||
// Controlar tamaño del caché
|
||||
// Control cache size
|
||||
if id_cache.len() >= MAX_CACHE_SIZE {
|
||||
warn!("ID-to-path cache size reached limit ({}), clearing oldest entries", MAX_CACHE_SIZE);
|
||||
id_cache.clear();
|
||||
|
||||
@@ -14,7 +14,7 @@ use crate::{
|
||||
dtos::share_dto::{CreateShareDto, UpdateShareDto},
|
||||
ports::share_ports::ShareUseCase
|
||||
},
|
||||
common::errors::{DomainError, ErrorKind},
|
||||
common::errors::ErrorKind,
|
||||
};
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
|
||||
Reference in New Issue
Block a user