chore: translate all Spanish comments and log messages to English

This commit is contained in:
Dionisio
2026-02-12 09:41:25 +01:00
parent bc169de8f6
commit d31a413e57
89 changed files with 2480 additions and 2482 deletions
+6 -6
View File
@@ -17,21 +17,21 @@ pub async fn create_auth_services(
pool: Arc<PgPool>,
folder_service: Option<Arc<FolderService>>
) -> Result<AuthServices> {
// Crear servicio de tokens JWT (implementación de TokenServicePort)
// Create JWT token service (TokenServicePort implementation)
let token_service: Arc<dyn TokenServicePort> = Arc::new(JwtTokenService::new(
config.auth.jwt_secret.clone(),
config.auth.access_token_expiry_secs,
config.auth.refresh_token_expiry_secs,
));
// Crear servicio de hashing de contraseñas
// Create password hashing service
let password_hasher = Arc::new(Argon2PasswordHasher::new());
// Crear repositorios PostgreSQL
// Create PostgreSQL repositories
let user_repository = Arc::new(UserPgRepository::new(pool.clone()));
let session_repository = Arc::new(SessionPgRepository::new(pool.clone()));
// Crear servicio de aplicación de autenticación
// Create authentication application service
let mut auth_app_service = AuthApplicationService::new(
user_repository,
session_repository,
@@ -39,7 +39,7 @@ pub async fn create_auth_services(
token_service.clone(),
);
// Configurar servicio de carpetas si está disponible
// Configure folder service if available
if let Some(folder_svc) = folder_service {
auth_app_service = auth_app_service.with_folder_service(folder_svc);
}
@@ -57,7 +57,7 @@ pub async fn create_auth_services(
}
}
// Empaquetar servicio en Arc
// Package service in Arc
let auth_application_service = Arc::new(auth_app_service);
Ok(AuthServices {
+13 -13
View File
@@ -4,7 +4,7 @@ use std::time::Duration;
use crate::common::config::AppConfig;
pub async fn create_database_pool(config: &AppConfig) -> Result<PgPool> {
tracing::info!("Inicializando conexión a PostgreSQL con URL: {}",
tracing::info!("Initializing PostgreSQL connection with URL: {}",
config.database.connection_string.replace("postgres://", "postgres://[user]:[pass]@"));
// Add a more robust connection attempt with retries
@@ -13,9 +13,9 @@ pub async fn create_database_pool(config: &AppConfig) -> Result<PgPool> {
while attempt < MAX_ATTEMPTS {
attempt += 1;
tracing::info!("Intento de conexión a PostgreSQL #{}", attempt);
tracing::info!("PostgreSQL connection attempt #{}", attempt);
// Crear el pool de conexiones con las opciones de configuración
// Create the connection pool with configuration options
match PgPoolOptions::new()
.max_connections(config.database.max_connections)
.min_connections(config.database.min_connections)
@@ -25,10 +25,10 @@ pub async fn create_database_pool(config: &AppConfig) -> Result<PgPool> {
.connect(&config.database.connection_string)
.await {
Ok(pool) => {
// Verificar la conexión
// Verify the connection
match sqlx::query("SELECT 1").execute(&pool).await {
Ok(_) => {
tracing::info!("Conexión a PostgreSQL establecida correctamente");
tracing::info!("PostgreSQL connection established successfully");
// Verify if migrations have been applied
let migration_check = sqlx::query("SELECT EXISTS (SELECT 1 FROM pg_tables WHERE schemaname = 'auth' AND tablename = 'users')")
@@ -39,34 +39,34 @@ pub async fn create_database_pool(config: &AppConfig) -> Result<PgPool> {
Ok(row) => {
let tables_exist: bool = row.get(0);
if !tables_exist {
tracing::warn!("Las tablas de la base de datos no existen. Por favor, ejecuta las migraciones con: cargo run --bin migrate --features migrations");
tracing::warn!("Database tables do not exist. Please run migrations with: cargo run --bin migrate --features migrations");
}
},
Err(_) => {
tracing::warn!("No se pudo verificar el estado de las migraciones. Por favor, ejecuta las migraciones con: cargo run --bin migrate --features migrations");
tracing::warn!("Could not verify migration status. Please run migrations with: cargo run --bin migrate --features migrations");
}
}
return Ok(pool);
},
Err(e) => {
tracing::error!("Error al verificar conexión: {}", e);
tracing::warn!("La base de datos parece no estar configurada. Por favor, ejecuta las migraciones con: cargo run --bin migrate --features migrations");
tracing::error!("Error verifying connection: {}", e);
tracing::warn!("The database appears to not be configured. Please run migrations with: cargo run --bin migrate --features migrations");
if attempt >= MAX_ATTEMPTS {
return Err(anyhow::anyhow!("Error al verificar la conexión a PostgreSQL: {}", e));
return Err(anyhow::anyhow!("Error verifying PostgreSQL connection: {}", e));
}
}
}
},
Err(e) => {
tracing::error!("Error al conectar a PostgreSQL: {}", e);
tracing::error!("Error connecting to PostgreSQL: {}", e);
if attempt >= MAX_ATTEMPTS {
return Err(anyhow::anyhow!("Error en la conexión a PostgreSQL: {}", e));
return Err(anyhow::anyhow!("Error in PostgreSQL connection: {}", e));
}
tokio::time::sleep(Duration::from_secs(1)).await;
}
}
}
Err(anyhow::anyhow!("No se pudo establecer la conexión a PostgreSQL después de {} intentos", MAX_ATTEMPTS))
Err(anyhow::anyhow!("Could not establish PostgreSQL connection after {} attempts", MAX_ATTEMPTS))
}
@@ -10,11 +10,11 @@ use crate::common::errors::DomainError;
use crate::domain::entities::file::File;
use crate::domain::services::path_service::StoragePath;
/// Composite que envuelve `Arc<dyn FileReadPort>` + `Arc<dyn FileWritePort>`
/// y delega cada método al port correspondiente.
/// Composite that wraps `Arc<dyn FileReadPort>` + `Arc<dyn FileWritePort>`
/// and delegates each method to the corresponding port.
///
/// Gracias al blanket impl `impl<T: FileReadPort + FileWritePort> FileStoragePort for T {}`
/// este tipo obtiene `FileStoragePort` automáticamente.
/// Thanks to the blanket impl `impl<T: FileReadPort + FileWritePort> FileStoragePort for T {}`
/// this type automatically gets `FileStoragePort`.
pub struct CompositeFileRepository {
read: Arc<dyn FileReadPort>,
write: Arc<dyn FileWritePort>,
@@ -21,9 +21,9 @@ use crate::infrastructure::services::path_service::PathService;
use crate::domain::services::path_service::StoragePath;
use crate::common::config::AppConfig;
/// Implementación de repositorio para operaciones de **lectura** de archivos.
/// Repository implementation for file **read** operations.
///
/// Implementa `FileReadPort`:
/// Implements `FileReadPort`:
/// get_file, list_files, get_file_content, get_file_stream,
/// get_file_range_stream, get_file_mmap, get_file_path, get_parent_folder_id.
pub struct FileFsReadRepository {
@@ -37,7 +37,7 @@ pub struct FileFsReadRepository {
}
impl FileFsReadRepository {
/// Constructor completo con todas las dependencias de infraestructura.
/// Full constructor with all infrastructure dependencies.
pub fn new(
root_path: PathBuf,
storage_mediator: Arc<dyn StorageMediator>,
@@ -58,7 +58,7 @@ impl FileFsReadRepository {
}
}
/// Stub para pruebas (no realiza I/O real).
/// Stub for testing (does not perform real I/O).
pub fn default_stub() -> Self {
Self {
root_path: PathBuf::from("./storage"),
@@ -75,7 +75,7 @@ impl FileFsReadRepository {
}
}
// ─── helpers internos ────────────────────────────────────
// ─── internal helpers ─────────────────────────────────────
fn resolve_storage_path(&self, storage_path: &StoragePath) -> PathBuf {
self.path_service.resolve_path(storage_path)
@@ -21,9 +21,9 @@ use crate::infrastructure::services::path_service::PathService;
use crate::domain::services::path_service::StoragePath;
use crate::common::config::AppConfig;
/// Implementación de repositorio para operaciones de **escritura** de archivos.
/// Repository implementation for file **write** operations.
///
/// Implementa `FileWritePort`:
/// Implements `FileWritePort`:
/// save_file, save_file_from_stream, move_file, delete_file,
/// update_file_content, register_file_deferred.
pub struct FileFsWriteRepository {
@@ -37,7 +37,7 @@ pub struct FileFsWriteRepository {
}
impl FileFsWriteRepository {
/// Constructor completo con todas las dependencias.
/// Full constructor with all dependencies.
pub fn new(
root_path: PathBuf,
storage_mediator: Arc<dyn StorageMediator>,
@@ -50,7 +50,7 @@ impl FileFsWriteRepository {
Self { root_path, storage_mediator, id_mapping_service, path_service, metadata_cache, config, parallel_processor }
}
/// Stub para pruebas (no realiza I/O real).
/// Stub for testing (does not perform real I/O).
pub fn default_stub() -> Self {
Self {
root_path: PathBuf::from("./storage"),
@@ -5,147 +5,147 @@ use tracing::{debug, error};
use crate::infrastructure::repositories::repository_errors::FolderRepositoryResult;
use crate::infrastructure::repositories::folder_fs_repository::FolderFsRepository;
// Este archivo contiene la implementación de los métodos relacionados con la papelera
// para el repositorio de carpetas FolderFsRepository
// This file contains the implementation of trash-related methods
// for the FolderFsRepository folder repository
// Implementación de métodos de papelera para el repositorio de carpetas
// Implementation of trash methods for the folder repository
impl FolderFsRepository {
// Obtiene la ruta completa a la papelera
// Gets the full path to the trash directory
fn get_trash_dir(&self) -> PathBuf {
self.get_root_path().join(".trash").join("folders")
}
// Crea una ruta única en la papelera para la carpeta
// Creates a unique path in the trash for the folder
async fn create_trash_folder_path(&self, folder_id: &str) -> FolderRepositoryResult<PathBuf> {
let trash_dir = self.get_trash_dir();
// Asegurarse que el directorio de la papelera existe
// Ensure the trash directory exists
if !trash_dir.exists() {
fs::create_dir_all(&trash_dir).await
.map_err(|e| FolderRepositoryError::StorageError(e.to_string()))?;
}
// Crear una ruta única para la carpeta en la papelera
// Create a unique path for the folder in the trash
Ok(trash_dir.join(folder_id))
}
}
// Implementación de los métodos públicos del trait FolderRepository relacionados con la papelera
// Implementation of public FolderRepository trait methods related to trash
// Implementation of internal methods for trash functionality
// These will be enabled when the trash feature is re-enabled
impl FolderFsRepository {
/// Helper method that will be used for trash functionality
pub(crate) async fn _trash_move_to_trash(&self, folder_id: &str) -> FolderRepositoryResult<()> {
debug!("Moviendo carpeta a la papelera: {}", folder_id);
debug!("Moving folder to trash: {}", folder_id);
// Obtener la ruta física de la carpeta
// Get the physical path of the folder
let folder_path = match self.get_mapped_folder_path(folder_id).await {
Ok(path) => path,
Err(e) => {
error!("Error obteniendo ruta de la carpeta {}: {:?}", folder_id, e);
error!("Error getting folder path {}: {:?}", folder_id, e);
return Err(e);
}
};
let folder_path_buf = PathBuf::from(folder_path.to_string());
// Verificamos que la carpeta existe
// Verify the folder exists
if !folder_path_buf.exists() {
return Err(FolderRepositoryError::NotFound(format!("Folder not found: {}", folder_id)));
}
// Crear directorio en la papelera
// Create directory in the trash
let trash_folder_path = self.create_trash_folder_path(folder_id).await?;
// Mover la carpeta físicamente a la papelera
// Physically move the folder to the trash
match fs::rename(&folder_path_buf, &trash_folder_path).await {
Ok(_) => {
debug!("Carpeta movida a papelera: {} -> {}", folder_path_buf.display(), trash_folder_path.display());
debug!("Folder moved to trash: {} -> {}", folder_path_buf.display(), trash_folder_path.display());
// Actualizar el mapeo al nuevo path en la papelera
// Update the mapping to the new path in the trash
if let Err(e) = self.update_mapped_folder_path(folder_id, &trash_folder_path).await {
error!("Error actualizando mapeo de carpeta en papelera: {}", e);
error!("Error updating folder mapping in trash: {}", e);
return Err(e);
}
Ok(())
},
Err(e) => {
error!("Error moviendo carpeta a papelera: {}", e);
error!("Error moving folder to trash: {}", e);
Err(FolderRepositoryError::StorageError(e.to_string()))
}
}
}
/// Restaura una carpeta desde la papelera a su ubicación original
/// Restores a folder from the trash to its original location
pub(crate) async fn _trash_restore_from_trash(&self, folder_id: &str, original_path: &str) -> FolderRepositoryResult<()> {
debug!("Restaurando carpeta {} a {}", folder_id, original_path);
debug!("Restoring folder {} to {}", folder_id, original_path);
// Obtener la ruta actual en la papelera
// Get the current path in the trash
let current_path = match self.get_mapped_folder_path(folder_id).await {
Ok(path) => PathBuf::from(path),
Err(e) => {
error!("Error obteniendo ruta actual de la carpeta {}: {:?}", folder_id, e);
error!("Error getting current folder path {}: {:?}", folder_id, e);
return Err(e);
}
};
// Convertir la ruta original a PathBuf
// Convert the original path to PathBuf
let original_path_buf = PathBuf::from(original_path);
// Asegurar que el directorio padre de destino existe
// Ensure the destination parent directory exists
if let Some(parent) = original_path_buf.parent() {
if !parent.exists() {
fs::create_dir_all(parent).await
.map_err(|e| {
error!("Error creando directorio padre para restauración: {}", e);
error!("Error creating parent directory for restoration: {}", e);
FolderRepositoryError::StorageError(e.to_string())
})?;
}
}
// Mover la carpeta de la papelera a su ubicación original
// Move the folder from the trash to its original location
match fs::rename(&current_path, &original_path_buf).await {
Ok(_) => {
debug!("Carpeta restaurada: {} -> {}", current_path.display(), original_path_buf.display());
debug!("Folder restored: {} -> {}", current_path.display(), original_path_buf.display());
// Actualizar el mapeo a la ruta original
// Update the mapping to the original path
if let Err(e) = self.update_mapped_folder_path(folder_id, &original_path_buf).await {
error!("Error actualizando mapeo de carpeta restaurada: {}", e);
error!("Error updating restored folder mapping: {}", e);
return Err(e);
}
Ok(())
},
Err(e) => {
error!("Error restaurando carpeta: {}", e);
error!("Error restoring folder: {}", e);
Err(FolderRepositoryError::StorageError(e.to_string()))
}
}
}
/// Elimina una carpeta permanentemente (usado por la papelera)
/// Permanently deletes a folder (used by the trash)
pub(crate) async fn _trash_delete_folder_permanently(&self, folder_id: &str) -> FolderRepositoryResult<()> {
debug!("Eliminando carpeta permanentemente: {}", folder_id);
debug!("Permanently deleting folder: {}", folder_id);
// Similar a delete_folder pero sin validaciones adicionales
// Similar to delete_folder but without additional validations
let folder_path = match self.get_mapped_folder_path(folder_id).await {
Ok(path) => PathBuf::from(path),
Err(e) => {
error!("Error obteniendo ruta de la carpeta {}: {:?}", folder_id, e);
error!("Error getting folder path {}: {:?}", folder_id, e);
return Err(e);
}
};
// Eliminar la carpeta recursivamente
// Delete the folder recursively
if folder_path.exists() {
match fs::remove_dir_all(&folder_path).await {
Ok(_) => {
debug!("Carpeta eliminada permanentemente: {}", folder_path.display());
debug!("Folder permanently deleted: {}", folder_path.display());
},
Err(e) => {
error!("Error eliminando carpeta permanentemente: {}", e);
// No reportar error si la carpeta ya no existe
error!("Error permanently deleting folder: {}", e);
// Don't report error if the folder no longer exists
if e.kind() != std::io::ErrorKind::NotFound {
return Err(FolderRepositoryError::StorageError(e.to_string()));
}
@@ -153,16 +153,16 @@ impl FolderFsRepository {
}
}
// Eliminar el mapeo
// Remove the mapping
if let Err(e) = self.remove_mapped_folder_id(folder_id).await {
error!("Error eliminando mapeo de la carpeta: {}", e);
error!("Error removing folder mapping: {}", e);
return Err(e);
}
debug!("Carpeta eliminada permanentemente con éxito: {}", folder_id);
debug!("Folder permanently deleted successfully: {}", folder_id);
Ok(())
}
}
// Re-exportaciones necesarias para el compilador
// Re-exports needed by the compiler
use crate::infrastructure::repositories::repository_errors::FolderRepositoryError;
@@ -20,9 +20,9 @@ impl CalendarEventPgRepository {
#[async_trait]
impl CalendarEventRepository for CalendarEventPgRepository {
async fn create_event(&self, event: CalendarEvent) -> CalendarEventRepositoryResult<CalendarEvent> {
// Este método necesitaría una implementación completa que construya el CalendarEvent
// desde el resultado de la query, utilizando métodos del constructor
// Para esta demostración, vamos a retornar el mismo evento
// This method would need a full implementation that builds the CalendarEvent
// from the query result, using constructor methods
// For this demonstration, we return the same event
sqlx::query(
r#"
@@ -50,7 +50,7 @@ impl CalendarEventRepository for CalendarEventPgRepository {
.await
.map_err(|e| DomainError::database_error(format!("Failed to create calendar event: {}", e)))?;
// Devolvemos el mismo evento en vez de un resultado
// We return the same event instead of a result
Ok(event)
}
@@ -86,8 +86,8 @@ impl CalendarEventRepository for CalendarEventPgRepository {
.await
.map_err(|e| DomainError::database_error(format!("Failed to update calendar event: {}", e)))?;
// En una implementación completa, recuperaríamos el evento actualizado
// Por simplicidad, devolvemos el mismo evento que recibimos
// In a full implementation, we would retrieve the updated event
// For simplicity, we return the same event we received
Ok(event)
}
@@ -176,9 +176,9 @@ impl CalendarEventRepository for CalendarEventPgRepository {
.map_err(|e| DomainError::database_error(format!("Failed to get calendar event by id: {}", e)))?
.ok_or_else(|| DomainError::not_found("Calendar Event", id.to_string()))?;
// En una implementación real, construiríamos un objeto CalendarEvent completo
// Por simplicidad, creamos un objeto con valores predeterminados para
// demostrar el enfoque sin macros
// In a real implementation, we would build a complete CalendarEvent object
// For simplicity, we create an object with default values to
// demonstrate the approach without macros
let event = CalendarEvent::with_id(
row.get("id"),
@@ -32,14 +32,14 @@ impl CalendarRepository for CalendarPgRepository {
.bind(calendar.owner_id())
.bind(calendar.description())
.bind(calendar.color())
.bind(false) // is_public no existe como campo
.bind(false) // is_public doesn't exist as a field
.bind(calendar.created_at())
.bind(calendar.updated_at())
.fetch_one(&*self.pool)
.await
.map_err(|e| DomainError::database_error(format!("Failed to create calendar: {}", e)))?;
// Construir el objeto Calendar utilizando su constructor with_id
// Build the Calendar object using its with_id constructor
let result = Calendar::with_id(
row.get("id"),
row.get("name"),
@@ -66,14 +66,14 @@ impl CalendarRepository for CalendarPgRepository {
.bind(calendar.name())
.bind(calendar.description())
.bind(calendar.color())
.bind(false) // is_public no existe como campo
.bind(false) // is_public doesn't exist as a field
.bind(now)
.bind(calendar.id())
.fetch_one(&*self.pool)
.await
.map_err(|e| DomainError::database_error(format!("Failed to update calendar: {}", e)))?;
// Construir el objeto Calendar utilizando su constructor with_id
// Build the Calendar object using its with_id constructor
let result = Calendar::with_id(
row.get("id"),
row.get("name"),
@@ -8,7 +8,7 @@ use crate::application::dtos::favorites_dto::FavoriteItemDto;
use crate::application::ports::favorites_ports::FavoritesRepositoryPort;
use crate::common::errors::{Result, DomainError, ErrorKind};
/// Implementación PostgreSQL del puerto de persistencia de favoritos.
/// PostgreSQL implementation of the favorites persistence port.
pub struct FavoritesPgRepository {
db_pool: Arc<PgPool>,
}
@@ -8,7 +8,7 @@ use crate::application::dtos::recent_dto::RecentItemDto;
use crate::application::ports::recent_ports::RecentItemsRepositoryPort;
use crate::common::errors::{Result, DomainError, ErrorKind};
/// Implementación PostgreSQL del puerto de persistencia de elementos recientes.
/// PostgreSQL implementation of the recent items persistence port.
pub struct RecentItemsPgRepository {
db_pool: Arc<PgPool>,
}
@@ -10,7 +10,7 @@ use crate::application::ports::auth_ports::SessionStoragePort;
use crate::common::errors::DomainError;
use crate::infrastructure::repositories::pg::transaction_utils::with_transaction;
// Implementar From<sqlx::Error> para SessionRepositoryError para permitir conversiones automáticas
// Implement From<sqlx::Error> for SessionRepositoryError to allow automatic conversions
impl From<sqlx::Error> for SessionRepositoryError {
fn from(err: sqlx::Error) -> Self {
SessionPgRepository::map_sqlx_error(err)
@@ -26,14 +26,14 @@ impl SessionPgRepository {
Self { pool }
}
// Método auxiliar para mapear errores SQL a errores de dominio
// Helper method to map SQL errors to domain errors
pub fn map_sqlx_error(err: sqlx::Error) -> SessionRepositoryError {
match err {
sqlx::Error::RowNotFound => {
SessionRepositoryError::NotFound("Sesión no encontrada".to_string())
SessionRepositoryError::NotFound("Session not found".to_string())
},
_ => SessionRepositoryError::DatabaseError(
format!("Error de base de datos: {}", err)
format!("Database error: {}", err)
),
}
}
@@ -41,9 +41,9 @@ impl SessionPgRepository {
#[async_trait]
impl SessionRepository for SessionPgRepository {
/// Crea una nueva sesión utilizando una transacción
/// Creates a new session using a transaction
async fn create_session(&self, session: Session) -> SessionRepositoryResult<Session> {
// Crear una copia de la sesión para el closure
// Create a copy of the session for the closure
let session_clone = session.clone();
with_transaction(
@@ -51,7 +51,7 @@ impl SessionRepository for SessionPgRepository {
"create_session",
|tx| {
Box::pin(async move {
// Insertar la sesión
// Insert the session
sqlx::query(
r#"
INSERT INTO auth.sessions (
@@ -74,8 +74,8 @@ impl SessionRepository for SessionPgRepository {
.await
.map_err(Self::map_sqlx_error)?;
// Opcionalmente, actualizar el último login del usuario
// dentro de la misma transacción
// Optionally, update the user's last login
// within the same transaction
sqlx::query(
r#"
UPDATE auth.users
@@ -87,12 +87,12 @@ impl SessionRepository for SessionPgRepository {
.execute(&mut **tx)
.await
.map_err(|e| {
// Convertimos el error pero sin interrumpir la creación
// de la sesión si falla la actualización
tracing::warn!("No se pudo actualizar last_login_at para usuario {}: {}",
// Convert the error but without interrupting session
// creation if the update fails
tracing::warn!("Could not update last_login_at for user {}: {}",
session_clone.user_id(), e);
SessionRepositoryError::DatabaseError(format!(
"Sesión creada pero no se pudo actualizar last_login_at: {}", e
"Session created but could not update last_login_at: {}", e
))
})?;
@@ -104,7 +104,7 @@ impl SessionRepository for SessionPgRepository {
Ok(session)
}
/// Obtiene una sesión por ID
/// Gets a session by ID
async fn get_session_by_id(&self, id: &str) -> SessionRepositoryResult<Session> {
let row = sqlx::query(
r#"
@@ -132,7 +132,7 @@ impl SessionRepository for SessionPgRepository {
))
}
/// Obtiene una sesión por token de actualización
/// Gets a session by refresh token
async fn get_session_by_refresh_token(&self, refresh_token: &str) -> SessionRepositoryResult<Session> {
let row = sqlx::query(
r#"
@@ -160,7 +160,7 @@ impl SessionRepository for SessionPgRepository {
))
}
/// Obtiene todas las sesiones de un usuario
/// Gets all sessions for a user
async fn get_sessions_by_user_id(&self, user_id: &str) -> SessionRepositoryResult<Vec<Session>> {
let rows = sqlx::query(
r#"
@@ -195,16 +195,16 @@ impl SessionRepository for SessionPgRepository {
Ok(sessions)
}
/// Revoca una sesión específica utilizando una transacción
/// Revokes a specific session using a transaction
async fn revoke_session(&self, session_id: &str) -> SessionRepositoryResult<()> {
let id = session_id.to_string(); // Clone para uso en closure
let id = session_id.to_string(); // Clone for use in closure
with_transaction(
&self.pool,
"revoke_session",
|tx| {
Box::pin(async move {
// Revocar la sesión
// Revoke the session
let result = sqlx::query(
r#"
UPDATE auth.sessions
@@ -218,14 +218,14 @@ impl SessionRepository for SessionPgRepository {
.await
.map_err(Self::map_sqlx_error)?;
// Si encontramos la sesión, podemos registrar un evento de seguridad
// If we found the session, we can log a security event
if let Some(row) = result {
let user_id: String = row.try_get("user_id").unwrap_or_default();
// Registrar evento de seguridad (en una tabla de seguridad)
// Esto es opcional pero muestra cómo se puede realizar operaciones
// adicionales en la misma transacción
tracing::info!("Sesión con ID {} del usuario {} revocada", id, user_id);
// Log security event (in a security table)
// This is optional but shows how additional operations
// can be performed in the same transaction
tracing::info!("Session with ID {} for user {} revoked", id, user_id);
}
Ok(())
@@ -234,16 +234,16 @@ impl SessionRepository for SessionPgRepository {
).await
}
/// Revoca todas las sesiones de un usuario utilizando una transacción
/// Revokes all sessions for a user using a transaction
async fn revoke_all_user_sessions(&self, user_id: &str) -> SessionRepositoryResult<u64> {
let user_id_clone = user_id.to_string(); // Clone para uso en closure
let user_id_clone = user_id.to_string(); // Clone for use in closure
with_transaction(
&self.pool,
"revoke_all_user_sessions",
|tx| {
Box::pin(async move {
// Revocar todas las sesiones del usuario
// Revoke all sessions for the user
let result = sqlx::query(
r#"
UPDATE auth.sessions
@@ -258,9 +258,9 @@ impl SessionRepository for SessionPgRepository {
let affected = result.rows_affected();
// Registrar evento de seguridad
// Log security event
if affected > 0 {
tracing::info!("Revocadas {} sesiones del usuario {}", affected, user_id_clone);
tracing::info!("Revoked {} sessions for user {}", affected, user_id_clone);
}
Ok(affected)
@@ -269,7 +269,7 @@ impl SessionRepository for SessionPgRepository {
).await
}
/// Elimina sesiones expiradas
/// Deletes expired sessions
async fn delete_expired_sessions(&self) -> SessionRepositoryResult<u64> {
let now = Utc::now();
@@ -288,7 +288,7 @@ impl SessionRepository for SessionPgRepository {
}
}
// Implementación del puerto de almacenamiento para la capa de aplicación
// Implementation of the storage port for the application layer
#[async_trait]
impl SessionStoragePort for SessionPgRepository {
async fn create_session(&self, session: Session) -> Result<Session, DomainError> {
@@ -9,7 +9,7 @@ use crate::application::ports::auth_ports::UserStoragePort;
use crate::common::errors::DomainError;
use crate::infrastructure::repositories::pg::transaction_utils::with_transaction;
// Implementar From<sqlx::Error> para UserRepositoryError para permitir conversiones automáticas
// Implement From<sqlx::Error> for UserRepositoryError to allow automatic conversions
impl From<sqlx::Error> for UserRepositoryError {
fn from(err: sqlx::Error) -> Self {
UserPgRepository::map_sqlx_error(err)
@@ -25,26 +25,26 @@ impl UserPgRepository {
Self { pool }
}
// Método auxiliar para mapear errores SQL a errores de dominio
// Helper method to map SQL errors to domain errors
pub fn map_sqlx_error(err: sqlx::Error) -> UserRepositoryError {
match err {
sqlx::Error::RowNotFound => {
UserRepositoryError::NotFound("Usuario no encontrado".to_string())
UserRepositoryError::NotFound("User not found".to_string())
},
sqlx::Error::Database(db_err) => {
if db_err.code().map_or(false, |code| code == "23505") {
// Código para violación de unicidad en PostgreSQL
// PostgreSQL uniqueness violation code
UserRepositoryError::AlreadyExists(
"Usuario o email ya existe".to_string()
"User or email already exists".to_string()
)
} else {
UserRepositoryError::DatabaseError(
format!("Error de base de datos: {}", db_err)
format!("Database error: {}", db_err)
)
}
},
_ => UserRepositoryError::DatabaseError(
format!("Error de base de datos: {}", err)
format!("Database error: {}", err)
),
}
}
@@ -52,7 +52,7 @@ impl UserPgRepository {
#[async_trait]
impl UserRepository for UserPgRepository {
/// Crea un nuevo usuario utilizando una transacción
/// Creates a new user using a transaction
async fn create_user(&self, user: User) -> UserRepositoryResult<User> {
// Creamos una copia del usuario para el closure
let user_clone = user.clone();
@@ -61,14 +61,14 @@ impl UserRepository for UserPgRepository {
&self.pool,
"create_user",
|tx| {
// Necesitamos mover el closure a un BoxFuture para devolver dentro
// de la llamada with_transaction
// We need to move the closure into a BoxFuture to return inside
// the with_transaction call
Box::pin(async move {
// Usamos los getters para extraer los valores
// Convertimos user.role() a string para pasarlo como texto plano
// Use getters to extract the values
// Convert user.role() to string to pass it as plain text
let role_str = user_clone.role().to_string();
// Modificar el SQL para hacer un cast explícito al tipo auth.userrole
// Modify the SQL to do an explicit cast to the auth.userrole type
let _result = sqlx::query(
r#"
INSERT INTO auth.users (
@@ -87,7 +87,7 @@ impl UserRepository for UserPgRepository {
.bind(user_clone.username())
.bind(user_clone.email())
.bind(user_clone.password_hash())
.bind(&role_str) // Convertir a string pero con cast explícito en SQL
.bind(&role_str) // Convert to string but with explicit cast in SQL
.bind(user_clone.storage_quota_bytes())
.bind(user_clone.storage_used_bytes())
.bind(user_clone.created_at())
@@ -100,18 +100,18 @@ impl UserRepository for UserPgRepository {
.await
.map_err(Self::map_sqlx_error)?;
// Podríamos realizar operaciones adicionales aquí,
// como configurar permisos, roles, etc.
// We could perform additional operations here,
// such as configuring permissions, roles, etc.
Ok(user_clone)
}) as BoxFuture<'_, UserRepositoryResult<User>>
}
).await?;
Ok(user) // Devolvemos el usuario original por simplicidad
Ok(user) // Return the original user for simplicity
}
/// Obtiene un usuario por ID
/// Gets a user by ID
async fn get_user_by_id(&self, id: &str) -> UserRepositoryResult<User> {
let row = sqlx::query(
r#"
@@ -153,7 +153,7 @@ impl UserRepository for UserPgRepository {
))
}
/// Obtiene un usuario por nombre de usuario
/// Gets a user by username
async fn get_user_by_username(&self, username: &str) -> UserRepositoryResult<User> {
let row = sqlx::query(
r#"
@@ -195,7 +195,7 @@ impl UserRepository for UserPgRepository {
))
}
/// Obtiene un usuario por correo electrónico
/// Gets a user by email
async fn get_user_by_email(&self, email: &str) -> UserRepositoryResult<User> {
let row = sqlx::query(
r#"
@@ -237,9 +237,9 @@ impl UserRepository for UserPgRepository {
))
}
/// Actualiza un usuario existente utilizando una transacción
/// Updates an existing user using a transaction
async fn update_user(&self, user: User) -> UserRepositoryResult<User> {
// Creamos una copia del usuario para el closure
// Create a copy of the user for the closure
let user_clone = user.clone();
with_transaction(
@@ -247,7 +247,7 @@ impl UserRepository for UserPgRepository {
"update_user",
|tx| {
Box::pin(async move {
// Actualizar el usuario
// Update the user
sqlx::query(
r#"
UPDATE auth.users
@@ -278,8 +278,8 @@ impl UserRepository for UserPgRepository {
.await
.map_err(Self::map_sqlx_error)?;
// Podríamos realizar operaciones adicionales aquí dentro
// de la misma transacción, como actualizar permisos, etc.
// We could perform additional operations here inside
// the same transaction, such as updating permissions, etc.
Ok(user_clone)
}) as BoxFuture<'_, UserRepositoryResult<User>>
@@ -289,7 +289,7 @@ impl UserRepository for UserPgRepository {
Ok(user)
}
/// Actualiza solo el uso de almacenamiento de un usuario
/// Updates only the storage usage of a user
async fn update_storage_usage(&self, user_id: &str, usage_bytes: i64) -> UserRepositoryResult<()> {
sqlx::query(
r#"
@@ -309,7 +309,7 @@ impl UserRepository for UserPgRepository {
Ok(())
}
/// Actualiza la fecha de último inicio de sesión
/// Updates the last login date
async fn update_last_login(&self, user_id: &str) -> UserRepositoryResult<()> {
sqlx::query(
r#"
@@ -328,7 +328,7 @@ impl UserRepository for UserPgRepository {
Ok(())
}
/// Lista usuarios con paginación
/// Lists users with pagination
async fn list_users(&self, limit: i64, offset: i64) -> UserRepositoryResult<Vec<User>> {
let rows = sqlx::query(
r#"
@@ -378,7 +378,7 @@ impl UserRepository for UserPgRepository {
Ok(users)
}
/// Activa o desactiva un usuario
/// Activates or deactivates a user
async fn set_user_active_status(&self, user_id: &str, active: bool) -> UserRepositoryResult<()> {
sqlx::query(
r#"
@@ -398,7 +398,7 @@ impl UserRepository for UserPgRepository {
Ok(())
}
/// Cambia la contraseña de un usuario
/// Changes a user's password
async fn change_password(&self, user_id: &str, password_hash: &str) -> UserRepositoryResult<()> {
sqlx::query(
r#"
@@ -418,9 +418,9 @@ impl UserRepository for UserPgRepository {
Ok(())
}
/// Cambia el rol de un usuario
/// Changes a user's role
async fn change_role(&self, user_id: &str, role: UserRole) -> UserRepositoryResult<()> {
// Convertir el rol a string para el binding
// Convert the role to string for the binding
let role_str = role.to_string();
sqlx::query(
@@ -441,7 +441,7 @@ impl UserRepository for UserPgRepository {
Ok(())
}
/// Lista usuarios por rol
/// Lists users by role
async fn list_users_by_role(&self, role: &str) -> UserRepositoryResult<Vec<User>> {
let rows = sqlx::query(
r#"
@@ -490,7 +490,7 @@ impl UserRepository for UserPgRepository {
Ok(users)
}
/// Elimina un usuario
/// Deletes a user
async fn delete_user(&self, user_id: &str) -> UserRepositoryResult<()> {
sqlx::query(
r#"
@@ -548,7 +548,7 @@ impl UserRepository for UserPgRepository {
))
}
/// Actualiza la cuota de almacenamiento de un usuario
/// Updates a user's storage quota
async fn update_storage_quota(&self, user_id: &str, quota_bytes: i64) -> UserRepositoryResult<()> {
sqlx::query(
r#"
@@ -568,7 +568,7 @@ impl UserRepository for UserPgRepository {
Ok(())
}
/// Cuenta el número total de usuarios
/// Counts the total number of users
async fn count_users(&self) -> UserRepositoryResult<i64> {
let row = sqlx::query(
"SELECT COUNT(*) as count FROM auth.users"
@@ -581,7 +581,7 @@ impl UserRepository for UserPgRepository {
Ok(count)
}
/// Obtiene estadísticas de almacenamiento agregadas
/// Gets aggregated storage statistics
async fn get_storage_stats(&self) -> UserRepositoryResult<StorageStats> {
let row = sqlx::query(
r#"
@@ -610,7 +610,7 @@ impl UserRepository for UserPgRepository {
}
}
// Implementación del puerto de almacenamiento para la capa de aplicación
// Storage port implementation for the application layer
#[async_trait]
impl UserStoragePort for UserPgRepository {
async fn create_user(&self, user: User) -> Result<User, DomainError> {
@@ -12,7 +12,7 @@ use crate::{
},
};
// Estructura para almacenar en el sistema de archivos
// Structure for storing in the file system
#[derive(Debug, Clone, Serialize, Deserialize)]
struct ShareRecord {
id: String,
@@ -38,12 +38,12 @@ impl ShareFsRepository {
Self { config }
}
/// Obtiene la ruta del archivo JSON donde se almacenan los enlaces compartidos
/// Gets the path to the JSON file where shared links are stored
fn get_shares_path(&self) -> String {
format!("{}/shares.json", self.config.storage_path.display())
}
/// Lee todos los enlaces compartidos del archivo JSON
/// Reads all shared links from the JSON file
async fn read_shares(&self) -> Result<Vec<ShareRecord>, io::Error> {
let path = self.get_shares_path();
let path = Path::new(&path);
@@ -58,12 +58,12 @@ impl ShareFsRepository {
Ok(shares)
}
/// Guarda todos los enlaces compartidos en el archivo JSON
/// Saves all shared links to the JSON file
async fn write_shares(&self, shares: &[ShareRecord]) -> Result<(), io::Error> {
let path = self.get_shares_path();
let json = serde_json::to_string_pretty(shares)?;
// Asegúrate de que el directorio existe
// Make sure the directory exists
let dir = Path::new(&path).parent().unwrap();
if !dir.exists() {
fs::create_dir_all(dir).await?
@@ -72,7 +72,7 @@ impl ShareFsRepository {
fs::write(path, json).await
}
/// Convierte un registro del sistema de archivos a una entidad de dominio
/// Converts a file system record to a domain entity
fn to_entity(&self, record: &ShareRecord) -> Share {
let item_type = ShareItemType::try_from(record.item_type.as_str())
.unwrap_or(ShareItemType::File);
@@ -97,7 +97,7 @@ impl ShareFsRepository {
)
}
/// Convierte una entidad de dominio a un registro para el sistema de archivos
/// Converts a domain entity to a file system record
fn to_record(&self, share: &Share) -> ShareRecord {
ShareRecord {
id: share.id().to_string(),
@@ -122,16 +122,16 @@ impl ShareStoragePort for ShareFsRepository {
let mut shares = self.read_shares().await
.map_err(|e| DomainError::internal_error("Share", e.to_string()))?;
// Verifica si el enlace ya existe
// Check if the link already exists
let existing_index = shares.iter().position(|s| s.id == share.id());
let record = self.to_record(share);
if let Some(index) = existing_index {
// Actualización
// Update
shares[index] = record;
} else {
// Inserción
// Insert
shares.push(record);
}
@@ -190,16 +190,16 @@ impl ShareStoragePort for ShareFsRepository {
let mut shares = self.read_shares().await
.map_err(|e| DomainError::internal_error("Share", e.to_string()))?;
// Busca el índice del enlace a actualizar
// Find the index of the link to update
let index = shares.iter().position(|s| s.id == share.id())
.ok_or_else(|| {
DomainError::not_found("Share", format!("Share with ID {} not found for update", share.id()))
})?;
// Actualiza el registro
// Update the record
shares[index] = self.to_record(share);
// Guarda los cambios
// Save changes
self.write_shares(&shares).await
.map_err(|e| DomainError::internal_error("Share", e.to_string()))?;
@@ -210,16 +210,16 @@ impl ShareStoragePort for ShareFsRepository {
let mut shares = self.read_shares().await
.map_err(|e| DomainError::internal_error("Share", e.to_string()))?;
// Encuentra el índice del enlace a eliminar
// Find the index of the link to delete
let initial_len = shares.len();
shares.retain(|s| s.id != id);
// Si no se eliminó ningún enlace, significa que no existía
// If no link was deleted, it means it didn't exist
if shares.len() == initial_len {
return Err(DomainError::not_found("Share", format!("Share with ID {} not found for deletion", id)));
}
// Guarda los cambios
// Save changes
self.write_shares(&shares).await
.map_err(|e| DomainError::internal_error("Share", e.to_string()))?;
@@ -230,15 +230,15 @@ impl ShareStoragePort for ShareFsRepository {
let shares = self.read_shares().await
.map_err(|e| DomainError::internal_error("Share", e.to_string()))?;
// Filtra los enlaces del usuario
// Filter the user's links
let user_shares: Vec<ShareRecord> = shares.into_iter()
.filter(|s| s.created_by == user_id)
.collect();
// Calcula el total
// Calculate the total
let total = user_shares.len();
// Aplica la paginación
// Apply pagination
let paginated: Vec<Share> = user_shares.iter()
.skip(offset)
.take(limit)
@@ -12,7 +12,7 @@ use crate::domain::entities::trashed_item::{TrashedItem, TrashedItemType};
use crate::domain::repositories::trash_repository::TrashRepository;
use crate::application::ports::outbound::IdMappingPort;
/// Estructura para almacenar elementos en la papelera en formato JSON
/// Structure for storing trash items in JSON format
#[derive(Debug, Serialize, Deserialize)]
struct TrashedItemEntry {
id: String,
@@ -25,7 +25,7 @@ struct TrashedItemEntry {
deletion_date: String,
}
/// Implementación del repositorio de papelera usando el sistema de archivos
/// Trash repository implementation using the file system
pub struct TrashFsRepository {
trash_dir: PathBuf,
trash_index_path: PathBuf,
@@ -45,7 +45,7 @@ impl TrashFsRepository {
}
}
/// Asegura que existe el directorio de papelera
/// Ensures the trash directory exists
async fn ensure_trash_dir(&self) -> Result<()> {
debug!("Checking if trash directory exists: {}", self.trash_dir.display());
if !self.trash_dir.exists() {
@@ -105,7 +105,7 @@ impl TrashFsRepository {
Ok(())
}
/// Obtiene todas las entradas del índice de papelera
/// Gets all entries from the trash index
async fn get_trash_entries(&self) -> Result<Vec<TrashedItemEntry>> {
self.ensure_trash_dir().await?;
@@ -134,7 +134,7 @@ impl TrashFsRepository {
Ok(entries)
}
/// Guarda todas las entradas en el índice de papelera
/// Saves all entries to the trash index
async fn save_trash_entries(&self, entries: Vec<TrashedItemEntry>) -> Result<()> {
self.ensure_trash_dir().await?;
@@ -155,7 +155,7 @@ impl TrashFsRepository {
Ok(())
}
/// Convierte una entrada JSON a entidad TrashedItem
/// Converts a JSON entry to a TrashedItem entity
fn entry_to_trashed_item(&self, entry: TrashedItemEntry) -> Result<TrashedItem> {
let item_type = match entry.item_type.as_str() {
"file" => TrashedItemType::File,
@@ -206,7 +206,7 @@ impl TrashFsRepository {
))
}
/// Convierte una entidad TrashedItem a entrada JSON
/// Converts a TrashedItem entity to a JSON entry
fn trashed_item_to_entry(&self, item: &TrashedItem) -> TrashedItemEntry {
TrashedItemEntry {
id: item.id().to_string(),
@@ -228,9 +228,9 @@ impl TrashFsRepository {
impl TrashRepository for TrashFsRepository {
#[instrument(skip(self))]
async fn add_to_trash(&self, item: &TrashedItem) -> Result<()> {
debug!("Añadiendo elemento a la papelera: id={}, user={}", item.id(), item.user_id());
debug!("Adding item to trash: id={}, user={}", item.id(), item.user_id());
// Aseguramos que existe el directorio de la papelera para este usuario
// Ensure the trash directory exists for this user
let user_trash_dir = self.trash_dir.join("files").join(item.user_id().to_string());
debug!("User trash directory path: {}", user_trash_dir.display());
@@ -268,7 +268,7 @@ impl TrashRepository for TrashFsRepository {
#[instrument(skip(self))]
async fn get_trash_items(&self, user_id: &Uuid) -> Result<Vec<TrashedItem>> {
debug!("Obteniendo elementos en papelera para usuario: {}", user_id);
debug!("Getting trash items for user: {}", user_id);
let entries = self.get_trash_entries().await?;
@@ -290,7 +290,7 @@ impl TrashRepository for TrashFsRepository {
#[instrument(skip(self))]
async fn get_trash_item(&self, id: &Uuid, user_id: &Uuid) -> Result<Option<TrashedItem>> {
debug!("Buscando elemento en papelera: id={}, user={}", id, user_id);
debug!("Looking for item in trash: id={}, user={}", id, user_id);
let entries = self.get_trash_entries().await?;
@@ -311,7 +311,7 @@ impl TrashRepository for TrashFsRepository {
#[instrument(skip(self))]
async fn restore_from_trash(&self, id: &Uuid, user_id: &Uuid) -> Result<()> {
debug!("Restaurando elemento de la papelera: id={}, user={}", id, user_id);
debug!("Restoring item from trash: id={}, user={}", id, user_id);
let mut entries = self.get_trash_entries().await?;
@@ -333,16 +333,16 @@ impl TrashRepository for TrashFsRepository {
#[instrument(skip(self))]
async fn delete_permanently(&self, id: &Uuid, user_id: &Uuid) -> Result<()> {
debug!("Eliminando permanentemente elemento de la papelera: id={}, user={}", id, user_id);
debug!("Permanently deleting item from trash: id={}, user={}", id, user_id);
// Simplemente eliminamos la entrada del índice
// Los archivos físicos se eliminarán a través del repositorio correspondiente
// Simply remove the entry from the index
// Physical files will be deleted through the corresponding repository
self.restore_from_trash(id, user_id).await
}
#[instrument(skip(self))]
async fn clear_trash(&self, user_id: &Uuid) -> Result<()> {
debug!("Limpiando papelera para usuario: {}", user_id);
debug!("Clearing trash for user: {}", user_id);
let mut entries = self.get_trash_entries().await?;
let user_id_str = user_id.to_string();
@@ -355,7 +355,7 @@ impl TrashRepository for TrashFsRepository {
#[instrument(skip(self))]
async fn get_expired_items(&self) -> Result<Vec<TrashedItem>> {
debug!("Buscando elementos de papelera expirados");
debug!("Looking for expired trash items");
let entries = self.get_trash_entries().await?;
let now = Utc::now();
+109 -109
View File
@@ -5,71 +5,71 @@ use tokio::sync::{Mutex, Semaphore};
use std::time::{Duration, Instant};
use tracing::debug;
/// Tamaño por defecto de los buffers en el pool
/// Default buffer size in the pool
pub const DEFAULT_BUFFER_SIZE: usize = 64 * 1024; // 64KB
/// Número máximo por defecto de buffers en el pool
/// Default maximum number of buffers in the pool
pub const DEFAULT_MAX_BUFFERS: usize = 100;
/// Tiempo de vida por defecto de un buffer inactivo (en segundos)
/// Default time-to-live for an inactive buffer (in seconds)
pub const DEFAULT_BUFFER_TTL: u64 = 60;
/// Buffer pooling para optimizar operaciones de lectura/escritura
/// Buffer pooling to optimize read/write operations
pub struct BufferPool {
/// Pool de buffers disponibles
/// Pool of available buffers
pool: Mutex<VecDeque<PooledBuffer>>,
/// Semáforo para limitar el número máximo de buffers
/// Semaphore to limit the maximum number of buffers
limit: Semaphore,
/// Tamaño de los buffers en el pool
/// Size of buffers in the pool
buffer_size: usize,
/// Estadísticas del pool
/// Pool statistics
stats: Mutex<BufferPoolStats>,
/// Tiempo de vida de un buffer inactivo
/// Time-to-live for an inactive buffer
buffer_ttl: Duration,
}
/// Estructura para tracking de estadísticas del pool
/// Structure for tracking pool statistics
#[derive(Debug, Clone, Default)]
pub struct BufferPoolStats {
/// Número total de operaciones de get
/// Total number of get operations
pub gets: usize,
/// Número de hits del pool (reutilización exitosa)
/// Number of pool hits (successful reuse)
pub hits: usize,
/// Número de misses (creación de nuevo buffer)
/// Number of misses (new buffer creation)
pub misses: usize,
/// Número de retornos al pool
/// Number of returns to the pool
pub returns: usize,
/// Número de eviction por TTL
/// Number of TTL evictions
pub evictions: usize,
/// Número máximo de buffers alcanzado
/// Maximum number of buffers reached
pub max_buffers_reached: usize,
/// Esperas por semáforo
/// Semaphore waits
pub waits: usize,
}
/// Buffer del pool con metadatos para gestión
/// Pool buffer with management metadata
struct PooledBuffer {
/// Buffer real de bytes
/// Actual byte buffer
buffer: Vec<u8>,
/// Timestamp de cuándo se añadió/retornó al pool
/// Timestamp of when it was added/returned to the pool
last_used: Instant,
}
/// Buffer prestado del pool con cleanup automático
/// Borrowed buffer from the pool with automatic cleanup
#[derive(Clone)]
pub struct BorrowedBuffer {
/// Buffer actual
/// Current buffer
buffer: Vec<u8>,
/// Tamaño real utilizado del buffer
/// Actual used size of the buffer
used_size: usize,
/// Referencia al pool para retornar
/// Reference to the pool for returning
pool: Arc<BufferPool>,
/// Si el buffer debe o no retornarse al pool
/// Whether the buffer should be returned to the pool or not
return_to_pool: bool,
}
impl BufferPool {
/// Crea un nuevo pool de buffers
/// Creates a new buffer pool
pub fn new(buffer_size: usize, max_buffers: usize, buffer_ttl_secs: u64) -> Arc<Self> {
Arc::new(Self {
pool: Mutex::new(VecDeque::with_capacity(max_buffers)),
@@ -80,7 +80,7 @@ impl BufferPool {
})
}
/// Crea un pool con configuración por defecto
/// Creates a pool with default configuration
pub fn default() -> Arc<Self> {
Self::new(
DEFAULT_BUFFER_SIZE,
@@ -89,25 +89,25 @@ impl BufferPool {
)
}
/// Obtiene un buffer del pool o crea uno nuevo si es necesario.
/// Gets a buffer from the pool or creates a new one if needed.
/// This version takes an Arc<Self> to ensure the BorrowedBuffer keeps a proper
/// reference to the shared pool (not a clone).
#[allow(unused_variables)]
pub async fn get_buffer(self: &Arc<Self>) -> BorrowedBuffer {
// Incrementar contador de gets
// Increment get counter
{
let mut stats = self.stats.lock().await;
stats.gets += 1;
}
// Control de concurrencia
// Concurrency control
// Acquire a semaphore permit. If none available, wait.
// We forget() the permit so it doesn't auto-release on drop.
// Instead, the permit is manually released in return_buffer/Drop via add_permits(1).
match self.limit.try_acquire() {
Ok(permit) => permit.forget(),
Err(_) => {
// No hay permisos disponibles, esperamos
// No permits available, waiting
{
let mut stats = self.stats.lock().await;
stats.waits += 1;
@@ -121,15 +121,15 @@ impl BufferPool {
}
};
// Intentar obtener un buffer existente del pool
// Try to get an existing buffer from the pool
let mut pool_locked = self.pool.lock().await;
let pool_arc = Arc::clone(self);
if let Some(mut pooled_buffer) = pool_locked.pop_front() {
// Verificar si el buffer ha expirado
// Check if the buffer has expired
if pooled_buffer.last_used.elapsed() > self.buffer_ttl {
// Buffer expirado, descartamos y creamos uno nuevo
// Expired buffer, discard and create a new one
let mut stats = self.stats.lock().await;
stats.evictions += 1;
stats.misses += 1;
@@ -137,8 +137,8 @@ impl BufferPool {
debug!("Buffer pool: evicted expired buffer");
// Crear nuevo buffer (reutilizando el permiso)
drop(pool_locked); // Liberar el lock antes de retornar
// Create new buffer (reusing the permit)
drop(pool_locked); // Release the lock before returning
BorrowedBuffer {
buffer: vec![0; self.buffer_size],
@@ -147,15 +147,15 @@ impl BufferPool {
return_to_pool: true,
}
} else {
// Buffer válido, lo reutilizamos
// Valid buffer, reuse it
let mut stats = self.stats.lock().await;
stats.hits += 1;
drop(stats);
// Liberar el lock antes de retornar
// Release the lock before returning
drop(pool_locked);
// Limpiar buffer por seguridad
// Clear buffer for security
pooled_buffer.buffer.fill(0);
BorrowedBuffer {
@@ -166,12 +166,12 @@ impl BufferPool {
}
}
} else {
// No hay buffers disponibles, creamos uno nuevo
// No buffers available, create a new one
let mut stats = self.stats.lock().await;
stats.misses += 1;
drop(stats);
// Liberar el lock antes de retornar
// Release the lock before returning
drop(pool_locked);
debug!("Buffer pool: creating new buffer");
@@ -185,9 +185,9 @@ impl BufferPool {
}
}
/// Retorna un buffer al pool
/// Returns a buffer to the pool
async fn return_buffer(&self, mut buffer: Vec<u8>) {
// Si el buffer es del tamaño incorrecto, lo descartamos
// If the buffer is the wrong size, discard it
if buffer.capacity() != self.buffer_size {
debug!("Buffer pool: discarding buffer of wrong size: {} (expected {})",
buffer.capacity(), self.buffer_size);
@@ -196,10 +196,10 @@ impl BufferPool {
return;
}
// Resize para asegurar capacidad correcta
// Resize to ensure correct capacity
buffer.resize(self.buffer_size, 0);
// Añadir al pool
// Add to the pool
let mut pool_locked = self.pool.lock().await;
pool_locked.push_back(PooledBuffer {
@@ -207,7 +207,7 @@ impl BufferPool {
last_used: Instant::now(),
});
// Actualizar estadísticas
// Update statistics
let mut stats = self.stats.lock().await;
stats.returns += 1;
@@ -217,24 +217,24 @@ impl BufferPool {
self.limit.add_permits(1);
}
/// Limpia buffers expirados del pool
/// Cleans expired buffers from the pool
pub async fn clean_expired_buffers(&self) {
let _now = Instant::now();
let mut pool_locked = self.pool.lock().await;
// Contar expirados
// Count expired
let count_before = pool_locked.len();
// Filtrar manteniendo solo los no expirados
// Filter keeping only non-expired
pool_locked.retain(|buffer| {
buffer.last_used.elapsed() <= self.buffer_ttl
});
// Contar cuántos se eliminaron
// Count how many were removed
let removed = count_before - pool_locked.len();
if removed > 0 {
// Actualizar estadísticas
// Update statistics
let mut stats = self.stats.lock().await;
stats.evictions += removed;
@@ -242,21 +242,21 @@ impl BufferPool {
}
}
/// Obtiene estadísticas actuales del pool
/// Gets current pool statistics
pub async fn get_stats(&self) -> BufferPoolStats {
self.stats.lock().await.clone()
}
/// Inicia la tarea periódica de limpieza
/// Starts the periodic cleanup task
pub fn start_cleaner(pool: Arc<Self>) {
tokio::spawn(async move {
let interval = Duration::from_secs(30); // Limpiar cada 30 segundos
let interval = Duration::from_secs(30); // Clean every 30 seconds
loop {
tokio::time::sleep(interval).await;
pool.clean_expired_buffers().await;
// Loguear estadísticas periódicamente
// Log statistics periodically
let stats = pool.get_stats().await;
debug!("Buffer pool stats: gets={}, hits={}, misses={}, hit_ratio={:.2}%, returns={}, \
evictions={}, max_reached={}, waits={}",
@@ -286,31 +286,31 @@ impl Clone for BufferPool {
}
impl BorrowedBuffer {
/// Accede al buffer interno
/// Accesses the internal buffer
pub fn as_mut_slice(&mut self) -> &mut [u8] {
&mut self.buffer
}
/// Obtiene una referencia a los datos utilizados
/// Gets a reference to the used data
pub fn as_slice(&self) -> &[u8] {
&self.buffer[..self.used_size]
}
/// Establece cuántos bytes se utilizaron realmente
/// Sets how many bytes were actually used
pub fn set_used(&mut self, size: usize) {
self.used_size = min(size, self.buffer.len());
}
/// Convierte en un Vec<u8> que incluye solo los datos utilizados
/// Converts into a Vec<u8> that includes only the used data
pub fn into_vec(mut self) -> Vec<u8> {
// Marcar para no devolver al pool
// Mark to not return to pool
self.return_to_pool = false;
// Crear un nuevo vector solo con los datos utilizados
// Create a new vector with only the used data
self.buffer[..self.used_size].to_vec()
}
/// Copia datos a este buffer y actualiza el tamaño usado
/// Copies data to this buffer and updates the used size
pub fn copy_from_slice(&mut self, data: &[u8]) -> usize {
let copy_size = min(data.len(), self.buffer.len());
self.buffer[..copy_size].copy_from_slice(&data[..copy_size]);
@@ -318,32 +318,32 @@ impl BorrowedBuffer {
copy_size
}
/// Impide que el buffer se devuelva al pool al destruirse
/// Prevents the buffer from being returned to the pool on destruction
pub fn do_not_return(mut self) -> Self {
self.return_to_pool = false;
self
}
/// Obtiene el tamaño total del buffer
/// Gets the total buffer size
pub fn capacity(&self) -> usize {
self.buffer.len()
}
/// Obtiene el tamaño usado del buffer
/// Gets the used buffer size
pub fn used_size(&self) -> usize {
self.used_size
}
}
// Cuando se hace drop de un BorrowedBuffer, lo devuelve al pool
// When a BorrowedBuffer is dropped, it is returned to the pool
impl Drop for BorrowedBuffer {
fn drop(&mut self) {
if self.return_to_pool {
// Tomar posesión del buffer y crear un clone del pool
// Take ownership of the buffer and create a clone of the pool
let buffer = std::mem::take(&mut self.buffer);
let pool = self.pool.clone();
// Spawn del return para que el drop no bloquee
// Spawn the return so that drop doesn't block
// return_buffer will release the semaphore permit
tokio::spawn(async move {
pool.return_buffer(buffer).await;
@@ -361,39 +361,39 @@ mod tests {
#[tokio::test]
async fn test_buffer_pooling() {
// Crear pool pequeño para testing
// Create small pool for testing
let pool = BufferPool::new(1024, 5, 60);
// Obtener un buffer
// Get a buffer
let mut buffer1 = pool.get_buffer().await;
buffer1.copy_from_slice(b"test data");
assert_eq!(buffer1.as_slice(), b"test data");
// Obtener otro buffer
// Get another buffer
let buffer2 = pool.get_buffer().await;
// Verificar stats
// Verify stats
let stats = pool.get_stats().await;
assert_eq!(stats.gets, 2);
assert_eq!(stats.hits, 0); // sin hits todavía
assert_eq!(stats.misses, 2); // todos son misses
assert_eq!(stats.hits, 0); // no hits yet
assert_eq!(stats.misses, 2); // all are misses
// Devolver buffer1 al pool (implícitamente por drop)
// Return buffer1 to pool (implicitly via drop)
drop(buffer1);
// Permitir que el return asíncrono ocurra
// Allow the async return to occur
tokio::time::sleep(Duration::from_millis(10)).await;
// Obtener otro buffer (debería reutilizar el retornado)
// Get another buffer (should reuse the returned one)
let buffer3 = pool.get_buffer().await;
// Verificar stats actualizados
// Verify updated stats
let stats = pool.get_stats().await;
assert_eq!(stats.gets, 3);
assert_eq!(stats.hits, 1); // ahora debería haber un hit
assert_eq!(stats.returns, 1); // un buffer retornado
assert_eq!(stats.hits, 1); // now there should be a hit
assert_eq!(stats.returns, 1); // one buffer returned
// Limpiar
// Cleanup
drop(buffer2);
drop(buffer3);
}
@@ -402,19 +402,19 @@ mod tests {
async fn test_buffer_operations() {
let pool = BufferPool::new(1024, 10, 60);
// Obtener buffer
// Get buffer
let mut buffer = pool.get_buffer().await;
// Escribir datos
// Write data
buffer.copy_from_slice(b"Hello, world!");
assert_eq!(buffer.used_size(), 13);
assert_eq!(buffer.as_slice(), b"Hello, world!");
// Convertir a vec y verificar
let vec = buffer.into_vec(); // Esto impide retornar al pool
// Convert to vec and verify
let vec = buffer.into_vec(); // This prevents returning to pool
assert_eq!(vec, b"Hello, world!");
// Verificar que no se incrementan los returns (buffer no retornado)
// Verify that returns are not incremented (buffer not returned)
tokio::time::sleep(Duration::from_millis(10)).await;
let stats = pool.get_stats().await;
assert_eq!(stats.returns, 0);
@@ -422,76 +422,76 @@ mod tests {
#[tokio::test]
async fn test_pool_limit() {
// Pool con solo 3 buffers
// Pool with only 3 buffers
let pool = BufferPool::new(1024, 3, 60);
// Obtener 3 buffers (alcanza el límite)
// Get 3 buffers (reaches the limit)
let buffer1 = pool.get_buffer().await;
let buffer2 = pool.get_buffer().await;
let buffer3 = pool.get_buffer().await;
// Verificar stats
// Verify stats
let stats = pool.get_stats().await;
assert_eq!(stats.gets, 3);
assert_eq!(stats.waits, 0); // sin esperas todavía
assert_eq!(stats.waits, 0); // no waits yet
// Intentar obtener un 4º buffer en una tarea separada (debería esperar)
// Try to get a 4th buffer in a separate task (should wait)
let pool_clone = pool.clone();
let handle = tokio::spawn(async move {
let _buffer4 = pool_clone.get_buffer().await;
true
});
// Dar tiempo para que la tarea intente tomar el buffer
// Give time for the task to try to take the buffer
tokio::time::sleep(Duration::from_millis(50)).await;
// Verificar que hay una espera
// Verify there is a wait
let stats = pool.get_stats().await;
assert_eq!(stats.waits, 1);
// Liberar un buffer
// Release a buffer
drop(buffer1);
// Dar tiempo para el retorno asíncrono y para que la tarea en espera obtenga su buffer
// Give time for the async return and for the waiting task to get its buffer
tokio::time::sleep(Duration::from_millis(50)).await;
// Verificar que la tarea pudo continuar
// Verify the task was able to continue
assert!(handle.await.unwrap());
// Limpiar
// Cleanup
drop(buffer2);
drop(buffer3);
}
#[tokio::test]
async fn test_ttl_expiration() {
// Pool con TTL muy corto para testing
let pool = BufferPool::new(1024, 5, 1); // 1 segundo TTL
// Pool with very short TTL for testing
let pool = BufferPool::new(1024, 5, 1); // 1 second TTL
// Obtener y devolver un buffer
// Get and return a buffer
let buffer = pool.get_buffer().await;
drop(buffer);
// Permitir que el return asíncrono ocurra
// Allow the async return to occur
tokio::time::sleep(Duration::from_millis(50)).await;
// Verificar que hay un buffer en el pool
// Verify there is a buffer in the pool
let stats = pool.get_stats().await;
assert_eq!(stats.returns, 1);
// Esperar a que expire el TTL
// Wait for the TTL to expire
tokio::time::sleep(Duration::from_secs(2)).await;
// Limpiar expirados
// Clean expired
pool.clean_expired_buffers().await;
// Obtener otro buffer (debería ser un miss ya que el anterior expiró)
// Get another buffer (should be a miss since the previous one expired)
let _buffer2 = pool.get_buffer().await;
// Verificar stats
// Verify stats
let stats = pool.get_stats().await;
assert_eq!(stats.evictions, 1); // un buffer expirado
assert_eq!(stats.hits, 0); // sin hits (el buffer expiró)
assert_eq!(stats.misses, 2); // dos misses (1er y 3er get)
assert_eq!(stats.evictions, 1); // one expired buffer
assert_eq!(stats.hits, 0); // no hits (the buffer expired)
assert_eq!(stats.misses, 2); // two misses (1st and 3rd get)
}
}
@@ -16,16 +16,16 @@ use crate::application::ports::compression_ports::{
use crate::domain::errors::DomainError;
use crate::infrastructure::services::buffer_pool::BufferPool;
/// Nivel de compresión para ficheros
/// Compression level for files
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CompressionLevel {
/// Sin compresión (solo para transferencia)
/// No compression (transfer only)
None = 0,
/// Compresión rápida con menor ratio
/// Fast compression with lower ratio
Fast = 1,
/// Compresión balanceada (por defecto)
/// Balanced compression (default)
Default = 6,
/// Compresión máxima (más lenta)
/// Maximum compression (slower)
Best = 9,
}
@@ -40,49 +40,49 @@ impl From<CompressionLevel> for Compression {
}
}
/// Umbral de tamaño para decidir si se comprime o no
/// Size threshold to decide whether to compress or not
const COMPRESSION_SIZE_THRESHOLD: u64 = 1024 * 50; // 50KB
/// Interfaz para servicios de compresión
/// Interface for compression services
#[async_trait]
pub trait CompressionService: Send + Sync {
/// Comprime datos en memoria
/// Compresses data in memory
async fn compress_data(&self, data: &[u8], level: CompressionLevel) -> io::Result<Vec<u8>>;
/// Descomprime datos en memoria
/// Decompresses data in memory
async fn decompress_data(&self, compressed_data: &[u8]) -> io::Result<Vec<u8>>;
/// Comprime un stream de datos
/// Compresses a data stream
fn compress_stream<S>(&self, stream: S, level: CompressionLevel)
-> impl Stream<Item = io::Result<Bytes>> + Send
where
S: Stream<Item = io::Result<Bytes>> + Send + 'static + Unpin;
/// Descomprime un stream de datos
/// Decompresses a data stream
fn decompress_stream<S>(&self, compressed_stream: S)
-> impl Stream<Item = io::Result<Bytes>> + Send
where
S: Stream<Item = io::Result<Bytes>> + Send + 'static + Unpin;
/// Determina si un archivo debe ser comprimido basado en su tipo MIME y tamaño
/// Determines whether a file should be compressed based on its MIME type and size
fn should_compress(&self, mime_type: &str, size: u64) -> bool;
}
/// Implementación de servicios de compresión usando Gzip
/// Gzip compression service implementation
pub struct GzipCompressionService {
/// Pool de buffers para optimización de memoria
/// Buffer pool for memory optimization
buffer_pool: Option<Arc<BufferPool>>,
}
impl GzipCompressionService {
/// Crea una nueva instancia del servicio
/// Creates a new service instance
pub fn new() -> Self {
Self {
buffer_pool: None,
}
}
/// Crea una nueva instancia del servicio con buffer pool
/// Creates a new service instance with buffer pool
pub fn new_with_buffer_pool(buffer_pool: Arc<BufferPool>) -> Self {
Self {
buffer_pool: Some(buffer_pool),
@@ -92,64 +92,64 @@ impl GzipCompressionService {
#[async_trait]
impl CompressionService for GzipCompressionService {
/// Comprime datos en memoria usando Gzip
/// Compresses data in memory using Gzip
async fn compress_data(&self, data: &[u8], level: CompressionLevel) -> io::Result<Vec<u8>> {
// Si tenemos un buffer pool, usar un buffer prestado para la compresión
// If we have a buffer pool, use a borrowed buffer for compression
if let Some(pool) = &self.buffer_pool {
// Estimar el tamaño de la compresión (aproximadamente 80% del original para casos típicos)
// Estimate the compression size (approximately 80% of original for typical cases)
let estimated_size = (data.len() as f64 * 0.8) as usize;
// Obtener un buffer del pool
// Get a buffer from the pool
let buffer = pool.get_buffer().await;
// Comprobar si el buffer es suficientemente grande
// Check if the buffer is large enough
if buffer.capacity() >= estimated_size {
// Ejecutar la compresión en un worker thread usando el buffer
// Run compression in a worker thread using the buffer
let buffer_ptr = Arc::new(tokio::sync::Mutex::new(buffer));
let buffer_clone = buffer_ptr.clone();
// Comprimir datos
// Clonar los datos para evitar problemas de lifetime
// Compress data
// Clone the data to avoid lifetime issues
let data_owned = data.to_vec();
let result = tokio::task::spawn_blocking(move || {
let mut encoder = GzEncoderRead::new(&data_owned[..], level.into());
// Intentar bloquear el mutex (no debería fallar ya que estamos en un hilo separado)
// Try to lock the mutex (should not fail since we are in a separate thread)
let mut buffer_guard = match futures::executor::block_on(buffer_clone.lock()) {
buffer => buffer,
};
// Leer directamente en el buffer
// Read directly into the buffer
let read_bytes = encoder.read(buffer_guard.as_mut_slice())?;
buffer_guard.set_used(read_bytes);
Ok(()) as io::Result<()>
}).await;
// Verificar resultado
// Verify result
match result {
Ok(Ok(())) => {
// Obtener el buffer y convertirlo a Vec<u8>
// Get the buffer and convert it to Vec<u8>
let buffer = buffer_ptr.lock().await;
let cloned_buffer = buffer.clone();
drop(buffer); // Liberar el mutex primero
drop(buffer); // Release the mutex first
return Ok(cloned_buffer.into_vec());
},
Ok(Err(e)) => {
error!("Error en compresión con buffer pool: {}", e);
// Continuar con implementación estándar
error!("Compression error with buffer pool: {}", e);
// Fall back to standard implementation
},
Err(e) => {
error!("Error en task de compresión con buffer pool: {}", e);
// Continuar con implementación estándar
error!("Compression task error with buffer pool: {}", e);
// Fall back to standard implementation
}
}
}
}
// Implementación estándar si no hay buffer pool o el buffer es insuficiente
// Clonar los datos para evitar problemas de lifetime
// Standard implementation if there is no buffer pool or the buffer is insufficient
// Clone the data to avoid lifetime issues
let data_owned = data.to_vec();
tokio::task::spawn_blocking(move || {
@@ -158,79 +158,79 @@ impl CompressionService for GzipCompressionService {
encoder.read_to_end(&mut compressed)?;
Ok(compressed)
}).await.unwrap_or_else(|e| {
error!("Error en task de compresión: {}", e);
error!("Compression task error: {}", e);
Err(io::Error::new(io::ErrorKind::Other, e.to_string()))
})
}
/// Descomprime datos en memoria
/// Decompresses data in memory
async fn decompress_data(&self, compressed_data: &[u8]) -> io::Result<Vec<u8>> {
// Si tenemos un buffer pool, usar un buffer prestado para la descompresión
// If we have a buffer pool, use a borrowed buffer for decompression
if let Some(pool) = &self.buffer_pool {
// Estimar el tamaño de la descompresión (aproximadamente 5x del comprimido para casos típicos)
// Estimate the decompression size (approximately 5x of compressed for typical cases)
let estimated_size = compressed_data.len() * 5;
// Obtener un buffer del pool
// Get a buffer from the pool
let buffer = pool.get_buffer().await;
// Comprobar si el buffer es suficientemente grande
// Check if the buffer is large enough
if buffer.capacity() >= estimated_size {
// Clonar datos comprimidos para mover al worker
// Clone compressed data to move to the worker
let data = compressed_data.to_vec();
let buffer_ptr = Arc::new(tokio::sync::Mutex::new(buffer));
let buffer_clone = buffer_ptr.clone();
// Descomprimir datos
// Decompress data
let result = tokio::task::spawn_blocking(move || {
let mut decoder = GzDecoder::new(&data[..]);
// Intentar bloquear el mutex
// Try to lock the mutex
let mut buffer_guard = match futures::executor::block_on(buffer_clone.lock()) {
buffer => buffer,
};
// Leer directamente en el buffer
// Read directly into the buffer
let read_bytes = decoder.read(buffer_guard.as_mut_slice())?;
buffer_guard.set_used(read_bytes);
Ok(()) as io::Result<()>
}).await;
// Verificar resultado
// Verify result
match result {
Ok(Ok(())) => {
// Obtener el buffer y convertirlo a Vec<u8>
// Get the buffer and convert it to Vec<u8>
let buffer = buffer_ptr.lock().await;
let cloned_buffer = buffer.clone();
drop(buffer); // Liberar el mutex primero
drop(buffer); // Release the mutex first
return Ok(cloned_buffer.into_vec());
},
Ok(Err(e)) => {
error!("Error en descompresión con buffer pool: {}", e);
// Continuar con implementación estándar
error!("Decompression error with buffer pool: {}", e);
// Fall back to standard implementation
},
Err(e) => {
error!("Error en task de descompresión con buffer pool: {}", e);
// Continuar con implementación estándar
error!("Decompression task error with buffer pool: {}", e);
// Fall back to standard implementation
}
}
}
}
// Implementación estándar si no hay buffer pool o el buffer es insuficiente
let data = compressed_data.to_vec(); // Clonar para mover al worker
// Standard implementation if there is no buffer pool or the buffer is insufficient
let data = compressed_data.to_vec(); // Clone to move to the worker
tokio::task::spawn_blocking(move || {
let mut decoder = GzDecoder::new(&data[..]);
let mut decompressed = Vec::new();
decoder.read_to_end(&mut decompressed)?;
Ok(decompressed)
}).await.unwrap_or_else(|e| {
error!("Error en task de descompresión: {}", e);
error!("Decompression task error: {}", e);
Err(io::Error::new(io::ErrorKind::Other, e.to_string()))
})
}
/// Comprime un stream de bytes
/// Compresses a byte stream
fn compress_stream<S>(&self, stream: S, level: CompressionLevel)
-> impl Stream<Item = io::Result<Bytes>> + Send
where
@@ -271,7 +271,7 @@ impl CompressionService for GzipCompressionService {
})
}
/// Descomprime un stream de bytes
/// Decompresses a byte stream
fn decompress_stream<S>(&self, compressed_stream: S)
-> impl Stream<Item = io::Result<Bytes>> + Send
where
@@ -310,14 +310,14 @@ impl CompressionService for GzipCompressionService {
})
}
/// Determina si un archivo debe ser comprimido basado en su tipo MIME y tamaño
/// Determines whether a file should be compressed based on its MIME type and size
fn should_compress(&self, mime_type: &str, size: u64) -> bool {
// No comprimir archivos muy pequeños (overhead)
// Do not compress very small files (overhead)
if size < COMPRESSION_SIZE_THRESHOLD {
return false;
}
// No comprimir archivos ya comprimidos
// Do not compress already compressed files
if mime_type.starts_with("image/")
&& !mime_type.contains("svg")
&& !mime_type.contains("bmp") {
@@ -345,7 +345,7 @@ impl CompressionService for GzipCompressionService {
return false;
}
// Comprimir archivos de texto, documentos, y otros tipos compresibles
// Compress text files, documents, and other compressible types
true
}
}
@@ -389,19 +389,19 @@ mod tests {
async fn test_compress_decompress_data() {
let service = GzipCompressionService::new();
// Datos de prueba
// Test data
let data = "Hello, world! ".repeat(1000).into_bytes();
// Comprimir
// Compress
let compressed = CompressionService::compress_data(&service, &data, CompressionLevel::Default).await.unwrap();
// Verificar que la compresión reduce el tamaño
// Verify that compression reduces the size
assert!(compressed.len() < data.len());
// Descomprimir
// Decompress
let decompressed = CompressionService::decompress_data(&service, &compressed).await.unwrap();
// Verificar que los datos originales se recuperan correctamente
// Verify that the original data is recovered correctly
assert_eq!(decompressed, data);
}
@@ -409,30 +409,30 @@ mod tests {
async fn test_compress_decompress_stream() {
let service = GzipCompressionService::new();
// Crear datos de prueba
// Create test data
let chunks = vec![
Ok(Bytes::from("Hello, ")),
Ok(Bytes::from("world! ")),
Ok(Bytes::from("This is a test of streaming compression.")),
];
// Convertir a stream
// Convert to stream
let input_stream = futures::stream::iter(chunks);
// Comprimir el stream
// Compress the stream
let compressed_stream = service.compress_stream(input_stream, CompressionLevel::Default);
// Recolectar los bytes comprimidos
// Collect the compressed bytes
let compressed_bytes = compressed_stream
.try_fold(Vec::new(), |mut acc, chunk| async move {
acc.extend_from_slice(&chunk);
Ok(acc)
}).await.unwrap();
// Descomprimir los datos
// Decompress the data
let decompressed = CompressionService::decompress_data(&service, &compressed_bytes).await.unwrap();
// Verificar resultado
// Verify result
let expected = "Hello, world! This is a test of streaming compression.";
assert_eq!(String::from_utf8(decompressed).unwrap(), expected);
}
@@ -441,17 +441,17 @@ mod tests {
fn test_should_compress() {
let service = GzipCompressionService::new();
// Casos que no deberían comprimirse
// Cases that should not be compressed
assert!(!CompressionService::should_compress(&service, "image/jpeg", 100 * 1024));
assert!(!CompressionService::should_compress(&service, "video/mp4", 10 * 1024 * 1024));
assert!(!CompressionService::should_compress(&service, "application/zip", 5 * 1024 * 1024));
// Casos que sí deberían comprimirse
// Cases that should be compressed
assert!(CompressionService::should_compress(&service, "text/html", 100 * 1024));
assert!(CompressionService::should_compress(&service, "application/json", 200 * 1024));
assert!(CompressionService::should_compress(&service, "text/plain", 1024 * 1024));
// Archivos pequeños no deberían comprimirse independientemente del tipo
// Small files should not be compressed regardless of type
assert!(!CompressionService::should_compress(&service, "text/html", 10 * 1024));
}
}
+122 -122
View File
@@ -13,61 +13,61 @@ use crate::domain::entities::file::File;
use crate::common::config::AppConfig;
/// Tipos de entradas en caché
/// Cache entry types
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CacheEntryType {
/// Archivo
/// File
File,
/// Directorio
/// Directory
Directory,
/// Tipo desconocido
/// Unknown type
Unknown,
}
/// Estadísticas de caché para monitoreo
/// Cache statistics for monitoring
#[derive(Debug, Clone, Default)]
pub struct CacheStats {
/// Número de hits en caché
/// Number of cache hits
pub hits: usize,
/// Número de misses en caché
/// Number of cache misses
pub misses: usize,
/// Número de invalidaciones manuales
/// Number of manual invalidations
pub invalidations: usize,
/// Número de expiraciones automáticas
/// Number of automatic expirations
pub expirations: usize,
/// Número de inserciones en caché
/// Number of cache inserts
pub inserts: usize,
/// Tiempo total ahorrado (milisegundos)
/// Total time saved (milliseconds)
pub time_saved_ms: u64,
}
/// Metadatos completos de archivo en caché
/// Complete cached file metadata
#[derive(Debug, Clone)]
pub struct FileMetadata {
/// Ruta absoluta del archivo
/// Absolute file path
pub path: PathBuf,
/// Si el archivo existe físicamente
/// Whether the file physically exists
pub exists: bool,
/// Tipo de entrada (archivo, directorio)
/// Entry type (file, directory)
pub entry_type: CacheEntryType,
/// Tamaño en bytes (para archivos)
/// Size in bytes (for files)
pub size: Option<u64>,
/// Tipo MIME (para archivos)
/// MIME type (for files)
pub mime_type: Option<String>,
/// Timestamp de creación (UNIX epoch seconds)
/// Creation timestamp (UNIX epoch seconds)
pub created_at: Option<u64>,
/// Timestamp de modificación (UNIX epoch seconds)
/// Modification timestamp (UNIX epoch seconds)
pub modified_at: Option<u64>,
/// Acceso previo (usado para LRU)
/// Previous access (used for LRU)
pub last_access: Instant,
/// Tiempo de expiración de la caché
/// Cache expiration time
pub expires_at: Instant,
/// Número de accesos a esta entrada
/// Number of accesses to this entry
pub access_count: usize,
}
impl FileMetadata {
/// Crea una nueva entrada de metadatos
/// Creates a new metadata entry
pub fn new(
path: PathBuf,
exists: bool,
@@ -94,56 +94,56 @@ impl FileMetadata {
}
}
/// Actualiza el tiempo de último acceso
/// Updates the last access time
pub fn touch(&mut self) {
self.last_access = Instant::now();
self.access_count += 1;
}
/// Verifica si la entrada ha expirado
/// Checks if the entry has expired
pub fn is_expired(&self) -> bool {
Instant::now() > self.expires_at
}
/// Actualiza el tiempo de expiración con un nuevo TTL
/// Updates the expiration time with a new TTL
pub fn update_expiry(&mut self, ttl: Duration) {
self.expires_at = Instant::now() + ttl;
}
}
/// Caché avanzada de metadatos de archivos
/// Advanced file metadata cache
pub struct FileMetadataCache {
/// Caché principal de metadatos
/// Main metadata cache
metadata_cache: RwLock<HashMap<PathBuf, FileMetadata>>,
/// Cola LRU para administración de caché
/// LRU queue for cache management
lru_queue: RwLock<VecDeque<PathBuf>>,
/// Estadísticas de uso del caché
/// Cache usage statistics
stats: RwLock<CacheStats>,
/// Configuración global de la aplicación
/// Global application configuration
config: AppConfig,
/// TTL adaptativo para entradas populares
/// Adaptive TTL for popular entries
ttl_multiplier: f64,
/// Umbral de popularidad para TTL extendido
/// Popularity threshold for extended TTL
popularity_threshold: usize,
/// Tamaño máximo de caché
/// Maximum cache size
max_entries: usize,
}
impl FileMetadataCache {
/// Crea una nueva instancia de caché de metadatos
/// Creates a new metadata cache instance
pub fn new(config: AppConfig, max_entries: usize) -> Self {
Self {
metadata_cache: RwLock::new(HashMap::with_capacity(max_entries)),
lru_queue: RwLock::new(VecDeque::with_capacity(max_entries)),
stats: RwLock::new(CacheStats::default()),
config,
ttl_multiplier: 5.0, // Entradas populares tienen 5x TTL
popularity_threshold: 10, // Después de 10 accesos se considera popular
ttl_multiplier: 5.0, // Popular entries have 5x TTL
popularity_threshold: 10, // After 10 accesses it's considered popular
max_entries,
}
}
/// Crea un objeto FileMetadata a partir de un objeto File
/// Creates a FileMetadata object from a File object
pub fn create_metadata_from_file(file: &File, abs_path: PathBuf) -> FileMetadata {
let entry_type = CacheEntryType::File;
let size = Some(file.size());
@@ -151,8 +151,8 @@ impl FileMetadataCache {
let created_at = Some(file.created_at());
let modified_at = Some(file.modified_at());
// Usar un TTL estándar
let ttl = Duration::from_secs(60); // 1 minuto
// Use a standard TTL
let ttl = Duration::from_secs(60); // 1 minute
FileMetadata::new(
abs_path,
@@ -166,28 +166,28 @@ impl FileMetadataCache {
)
}
/// Crea una instancia por defecto
/// Creates a default instance
pub fn default() -> Self {
Self::new(AppConfig::default(), 10_000)
}
/// Crea una instancia de caché con configuración por defecto
/// Creates a cache instance with default configuration
pub fn default_with_config(config: AppConfig) -> Self {
Self::new(config, 50_000) // Caché más grande para sistema en producción
Self::new(config, 50_000) // Larger cache for production system
}
/// Obtiene los metadatos de un archivo si están en caché
/// Gets file metadata if cached
pub async fn get_metadata(&self, path: &Path) -> Option<FileMetadata> {
let start_time = Instant::now();
let mut cache = self.metadata_cache.write().await;
if let Some(metadata) = cache.get_mut(path) {
// Verificar si ha expirado
// Check if expired
if metadata.is_expired() {
// Eliminar de caché si expiró
// Remove from cache if expired
cache.remove(path);
// Actualizar estadísticas
// Update statistics
let mut stats = self.stats.write().await;
stats.misses += 1;
stats.expirations += 1;
@@ -197,10 +197,10 @@ impl FileMetadataCache {
return None;
}
// Actualizar tiempo de acceso
// Update access time
metadata.touch();
// Para entradas populares, extender TTL
// For popular entries, extend TTL
if metadata.access_count >= self.popularity_threshold {
let new_ttl = match metadata.entry_type {
CacheEntryType::File => Duration::from_millis(
@@ -209,33 +209,33 @@ impl FileMetadataCache {
CacheEntryType::Directory => Duration::from_millis(
(self.config.timeouts.dir_operation_ms as f64 * self.ttl_multiplier) as u64
),
_ => Duration::from_secs(60), // 1 minuto por defecto
_ => Duration::from_secs(60), // 1 minute by default
};
metadata.update_expiry(new_ttl);
debug!("Extended TTL for popular entry: {}", path.display());
}
// Calcular tiempo ahorrado aproximado
// Calculate approximate time saved
let elapsed = start_time.elapsed().as_millis() as u64;
let estimated_io_time: u64 = 10; // Asumimos 10ms mínimo para operación de IO
let estimated_io_time: u64 = 10; // We assume 10ms minimum for IO operation
let time_saved = estimated_io_time.saturating_sub(elapsed);
// Actualizar estadísticas
// Update statistics
let mut stats = self.stats.write().await;
stats.hits += 1;
stats.time_saved_ms += time_saved;
debug!("Cache hit for: {}", path.display());
// Mantener también la cola LRU actualizada
// Also keep the LRU queue updated
self.update_lru(path.to_path_buf()).await;
// Clonar para retornar
// Clone to return
return Some(metadata.clone());
}
// No encontrado en caché
// Not found in cache
let mut stats = self.stats.write().await;
stats.misses += 1;
@@ -243,20 +243,20 @@ impl FileMetadataCache {
None
}
/// Actualiza la cola LRU
/// Updates the LRU queue
async fn update_lru(&self, path: PathBuf) {
let mut lru = self.lru_queue.write().await;
// Eliminar si ya existe
// Remove if already exists
if let Some(pos) = lru.iter().position(|p| p == &path) {
lru.remove(pos);
}
// Agregar al final (más reciente)
// Add to the end (most recent)
lru.push_back(path);
}
/// Verifica si un archivo existe
/// Checks if a file exists
pub async fn exists(&self, path: &Path) -> Option<bool> {
if let Some(metadata) = self.get_metadata(path).await {
return Some(metadata.exists);
@@ -265,7 +265,7 @@ impl FileMetadataCache {
None
}
/// Verifica si un path es un directorio
/// Checks if a path is a directory
pub async fn is_dir(&self, path: &Path) -> Option<bool> {
if let Some(metadata) = self.get_metadata(path).await {
return Some(metadata.entry_type == CacheEntryType::Directory);
@@ -274,7 +274,7 @@ impl FileMetadataCache {
None
}
/// Verifica si un path es un archivo
/// Checks if a path is a file
pub async fn is_file(&self, path: &Path) -> Option<bool> {
if let Some(metadata) = self.get_metadata(path).await {
return Some(metadata.entry_type == CacheEntryType::File);
@@ -283,7 +283,7 @@ impl FileMetadataCache {
None
}
/// Obtiene el tamaño de un archivo
/// Gets the size of a file
pub async fn get_size(&self, path: &Path) -> Option<u64> {
if let Some(metadata) = self.get_metadata(path).await {
return metadata.size;
@@ -292,7 +292,7 @@ impl FileMetadataCache {
None
}
/// Obtiene el tipo MIME de un archivo
/// Gets the MIME type of a file
pub async fn get_mime_type(&self, path: &Path) -> Option<String> {
if let Some(metadata) = self.get_metadata(path).await {
return metadata.mime_type;
@@ -301,12 +301,12 @@ impl FileMetadataCache {
None
}
/// Refresca los metadatos de un path
/// Refreshes metadata for a path
pub async fn refresh_metadata(&self, path: &Path) -> Result<FileMetadata, std::io::Error> {
// Realizar lectura real del sistema de archivos
// Perform actual filesystem read
let metadata = fs::metadata(path).await?;
// Determinar tipo de entrada
// Determine entry type
let entry_type = if metadata.is_dir() {
CacheEntryType::Directory
} else if metadata.is_file() {
@@ -315,21 +315,21 @@ impl FileMetadataCache {
CacheEntryType::Unknown
};
// Obtener tamaño para archivos
// Get size for files
let size = if metadata.is_file() {
Some(metadata.len())
} else {
None
};
// Obtener tipo MIME para archivos
// Get MIME type for files
let mime_type = if metadata.is_file() {
Some(from_path(path).first_or_octet_stream().to_string())
} else {
None
};
// Obtener timestamps
// Get timestamps
let created_at = metadata.created()
.map(|time| time.duration_since(UNIX_EPOCH).unwrap_or_default().as_secs())
.ok();
@@ -338,14 +338,14 @@ impl FileMetadataCache {
.map(|time| time.duration_since(UNIX_EPOCH).unwrap_or_default().as_secs())
.ok();
// Determinar TTL apropiado
// Determine appropriate TTL
let ttl = if metadata.is_dir() {
Duration::from_millis(self.config.timeouts.dir_operation_ms)
} else {
Duration::from_millis(self.config.timeouts.file_operation_ms)
};
// Crear entrada de metadatos
// Create metadata entry
let file_metadata = FileMetadata::new(
path.to_path_buf(),
true,
@@ -357,34 +357,34 @@ impl FileMetadataCache {
ttl,
);
// Actualizar caché
// Update cache
self.update_cache(file_metadata.clone()).await;
Ok(file_metadata)
}
/// Actualiza la caché con nuevos metadatos
/// Updates the cache with new metadata
pub async fn update_cache(&self, metadata: FileMetadata) {
// Evitar caché llena antes de insertar
// Avoid full cache before inserting
self.ensure_capacity().await;
let path = metadata.path.clone();
// Insertar en caché
// Insert into cache
{
let mut cache = self.metadata_cache.write().await;
cache.insert(path.clone(), metadata);
// Actualizar estadísticas
// Update statistics
let mut stats = self.stats.write().await;
stats.inserts += 1;
}
// Actualizar la cola LRU
// Update the LRU queue
self.update_lru(path).await;
}
/// Asegura que hay espacio en la caché
/// Ensures there is space in the cache
async fn ensure_capacity(&self) {
let cache_size = {
let cache = self.metadata_cache.read().await;
@@ -392,15 +392,15 @@ impl FileMetadataCache {
};
if cache_size >= self.max_entries {
self.evict_lru_entries(cache_size / 10).await; // Liberar 10%
self.evict_lru_entries(cache_size / 10).await; // Free up 10%
}
}
/// Elimina entradas menos recientemente usadas
/// Removes least recently used entries
async fn evict_lru_entries(&self, count: usize) {
let mut paths_to_remove = Vec::with_capacity(count);
// Obtener entries a eliminar de la cola LRU
// Get entries to remove from the LRU queue
{
let mut lru = self.lru_queue.write().await;
for _ in 0..count {
@@ -412,7 +412,7 @@ impl FileMetadataCache {
}
}
// Eliminar de la caché principal
// Remove from the main cache
{
let mut cache = self.metadata_cache.write().await;
for path in paths_to_remove {
@@ -423,19 +423,19 @@ impl FileMetadataCache {
debug!("Evicted {} LRU entries from cache", count);
}
/// Invalidar una entrada específica de caché
/// Invalidate a specific cache entry
pub async fn invalidate(&self, path: &Path) {
// Eliminar de la caché principal
// Remove from the main cache
{
let mut cache = self.metadata_cache.write().await;
cache.remove(path);
// Actualizar estadísticas
// Update statistics
let mut stats = self.stats.write().await;
stats.invalidations += 1;
}
// Eliminar de la cola LRU
// Remove from the LRU queue
let path_buf = path.to_path_buf();
{
let mut lru = self.lru_queue.write().await;
@@ -447,12 +447,12 @@ impl FileMetadataCache {
debug!("Invalidated cache entry for: {}", path.display());
}
/// Invalidar recursivamente entradas bajo un directorio
/// Recursively invalidate entries under a directory
pub async fn invalidate_directory(&self, dir_path: &Path) {
let dir_str = dir_path.to_string_lossy().to_string();
let mut paths_to_remove = Vec::new();
// Encontrar todos los paths que comienzan con el directorio
// Find all paths that start with the directory
{
let cache = self.metadata_cache.read().await;
for path in cache.keys() {
@@ -463,13 +463,13 @@ impl FileMetadataCache {
}
}
// Actualizar estadísticas
// Update statistics
{
let mut stats = self.stats.write().await;
stats.invalidations += paths_to_remove.len();
}
// Eliminar cada path encontrado
// Remove each found path
for path in paths_to_remove {
self.invalidate(&path).await;
}
@@ -477,18 +477,18 @@ impl FileMetadataCache {
debug!("Invalidated directory and contents: {}", dir_path.display());
}
/// Obtener estadísticas actuales de la caché
/// Get current cache statistics
pub async fn get_stats(&self) -> CacheStats {
let stats = self.stats.read().await;
stats.clone()
}
/// Limpia todas las entradas expiradas de la caché
/// Clears all expired entries from the cache
pub async fn clear_expired(&self) {
let now = Instant::now();
let mut paths_to_remove = Vec::new();
// Encontrar entradas expiradas
// Find expired entries
{
let cache = self.metadata_cache.read().await;
for (path, metadata) in cache.iter() {
@@ -498,16 +498,16 @@ impl FileMetadataCache {
}
}
// Actualizar estadísticas
// Update statistics
{
let mut stats = self.stats.write().await;
stats.expirations += paths_to_remove.len();
}
// Guardar la cantidad de entradas para el logging
// Save the number of entries for logging
let num_paths = paths_to_remove.len();
// Eliminar entradas expiradas
// Remove expired entries
for path in paths_to_remove {
self.invalidate(&path).await;
}
@@ -515,19 +515,19 @@ impl FileMetadataCache {
debug!("Cleared {} expired entries from cache", num_paths);
}
/// Inicia el proceso de limpieza periódica
/// Starts the periodic cleanup process
pub fn start_cleanup_task(cache: Arc<Self>) -> BoxFuture<'static, ()> {
Box::pin(async move {
let cleanup_interval = Duration::from_secs(60); // Cada minuto
let cleanup_interval = Duration::from_secs(60); // Every minute
loop {
// Esperar el intervalo
// Wait for the interval
time::sleep(cleanup_interval).await;
// Limpiar entradas expiradas
// Clean expired entries
cache.clear_expired().await;
// Registrar estadísticas
// Log statistics
let stats = cache.get_stats().await;
let cache_size = {
let cache_map = cache.metadata_cache.read().await;
@@ -550,12 +550,12 @@ impl FileMetadataCache {
})
}
/// Precarga metadatos de directorios completos (útil para inicialización)
/// Preloads metadata for entire directories (useful for initialization)
pub async fn preload_directory(&self, dir_path: &Path, recursive: bool, max_depth: usize) -> Result<usize, std::io::Error> {
self._preload_directory_internal(dir_path, recursive, max_depth, 0).await
}
/// Implementación interna de precarga con seguimiento de profundidad
/// Internal preload implementation with depth tracking
async fn _preload_directory_internal(
&self,
dir_path: &Path,
@@ -568,20 +568,20 @@ impl FileMetadataCache {
return Ok(0);
}
// Obtener entradas del directorio
// Get directory entries
let mut entries = fs::read_dir(dir_path).await?;
let mut count = 0;
// Procesar cada entrada
// Process each entry
while let Some(entry) = entries.next_entry().await? {
let path = entry.path();
let metadata = fs::metadata(&path).await?;
// Refrescar metadatos de esta entrada
// Refresh metadata for this entry
self.refresh_metadata(&path).await?;
count += 1;
// Recursivamente procesar subdirectorios si es necesario
// Recursively process subdirectories if needed
if recursive && metadata.is_dir() {
// Box to break recursion
count += self._preload_directory_internal(
@@ -657,37 +657,37 @@ mod tests {
#[tokio::test]
async fn test_cache_operations() {
// Crear directorio temporal para pruebas
// Create temporary directory for tests
let temp_dir = tempdir().unwrap();
let file_path = temp_dir.path().join("test_file.txt");
// Crear un archivo de prueba
// Create a test file
let mut file = File::create(&file_path).await.unwrap();
file.write_all(b"test content").await.unwrap();
file.flush().await.unwrap();
drop(file);
// Crear caché
// Create cache
let config = AppConfig::default();
let cache = FileMetadataCache::new(config, 1000);
// Verificar miss inicial
// Verify initial miss
assert!(cache.exists(&file_path).await.is_none());
// Refrescar y verificar hit
// Refresh and verify hit
let metadata = cache.refresh_metadata(&file_path).await.unwrap();
assert_eq!(metadata.entry_type, CacheEntryType::File);
assert_eq!(metadata.size, Some(12)); // "test content" = 12 bytes
// Verificar que ahora existe en caché
// Verify it now exists in cache
assert_eq!(cache.exists(&file_path).await, Some(true));
assert_eq!(cache.is_file(&file_path).await, Some(true));
// Invalidar y verificar que ya no existe en caché
// Invalidate and verify it no longer exists in cache
cache.invalidate(&file_path).await;
assert!(cache.exists(&file_path).await.is_none());
// Verificar estadísticas
// Verify statistics
let stats = cache.get_stats().await;
assert_eq!(stats.inserts, 1);
assert_eq!(stats.invalidations, 1);
@@ -696,7 +696,7 @@ mod tests {
#[tokio::test]
async fn test_directory_operations() {
// Crear estructura de directorios para pruebas
// Create directory structure for tests
let temp_dir = tempdir().unwrap();
// Canonicalize to handle macOS /var -> /private/var symlinks
let base_path = temp_dir.path().canonicalize().unwrap();
@@ -709,24 +709,24 @@ mod tests {
File::create(&file1).await.unwrap();
File::create(&file2).await.unwrap();
// Crear caché
// Create cache
let config = AppConfig::default();
let cache = FileMetadataCache::new(config, 1000);
// Precargar directorio recursivamente
// Preload directory recursively
// preload_directory caches the *contents* of the directory, not the root itself
let count = cache.preload_directory(&base_path, true, 2).await.unwrap();
assert_eq!(count, 3); // subdir, file1, file2
// Verificar existencia en caché (solo contenido, no la raíz)
// Verify existence in cache (only contents, not the root)
assert_eq!(cache.is_dir(&sub_dir).await, Some(true));
assert_eq!(cache.is_file(&file1).await, Some(true));
assert_eq!(cache.is_file(&file2).await, Some(true));
// Invalidar directorio y contenido
// Invalidate directory and contents
cache.invalidate_directory(&base_path).await;
// Verificar que nada existe en caché
// Verify nothing exists in cache
assert!(cache.exists(&sub_dir).await.is_none());
assert!(cache.exists(&file1).await.is_none());
assert!(cache.exists(&file2).await.is_none());
@@ -179,7 +179,7 @@ impl 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;
@@ -189,19 +189,19 @@ impl IdMappingOptimizer {
}
}
// Si no está en caché, agregar a la cola de batch
// If not in cache, add to batch queue
{
let mut batch_queue = self.pending_batch.lock().await;
batch_queue.path_to_id_requests.insert(path_str);
}
// No encontrado en caché, debe procesarse en batch
// Not found in cache, must be processed in batch
Ok(None)
}
/// Procesa las solicitudes pendientes en batch
/// Processes pending requests in batch
async fn process_batch(&self) -> Result<BatchResult, IdMappingError> {
// Adquirir permiso para operación batch
// Acquire permit for batch operation
let _permit = self.batch_limiter.acquire().await.unwrap();
// Get pending requests
@@ -214,13 +214,13 @@ impl IdMappingOptimizer {
(paths, ids)
};
// Crear resultados
// Create results
let mut result = BatchResult {
path_to_id: HashMap::with_capacity(path_requests.len()),
id_to_path: HashMap::with_capacity(id_requests.len()),
};
// Procesar solicitudes path->id en batch
// Process path->id requests in batch
for path_str in path_requests {
let path = StoragePath::from_string(&path_str);
match self.base_service.get_or_create_id(&path).await {
@@ -230,12 +230,12 @@ impl IdMappingOptimizer {
},
Err(e) => {
error!("Error batch-processing path {}: {}", path_str, e);
// Continuar con las demás solicitudes
// Continue with remaining requests
}
}
}
// Procesar solicitudes id->path en batch
// Process id->path requests in batch
for id in id_requests {
match self.base_service.get_path_by_id(&id).await {
Ok(path) => {
@@ -245,7 +245,7 @@ impl IdMappingOptimizer {
},
Err(e) => {
error!("Error batch-processing ID {}: {}", id, e);
// Continuar con las demás solicitudes
// Continue with remaining requests
}
}
}
@@ -266,14 +266,14 @@ impl IdMappingOptimizer {
}
}
// Actualizar estadísticas
// Update statistics
{
let mut stats = self.stats.write().await;
stats.batch_operations += 1;
stats.batch_items_processed += result.path_to_id.len() + result.id_to_path.len();
}
// Guardar los cambios al disco en segundo plano
// Save changes to disk in the background
let service_clone = self.base_service.clone();
tokio::spawn(async move {
if let Err(e) = service_clone.save_pending_changes().await {
@@ -284,15 +284,15 @@ impl IdMappingOptimizer {
Ok(result)
}
/// Fuerza el procesamiento de solicitudes pendientes si hay suficientes
/// Forces processing of pending requests if there are enough
async fn trigger_batch_if_needed(&self, min_batch_size: usize) -> Result<(), IdMappingError> {
// Verificar si hay suficientes solicitudes pendientes
// Check if there are enough pending requests
let should_process = {
let batch_queue = self.pending_batch.lock().await;
batch_queue.path_to_id_requests.len() + batch_queue.id_to_path_requests.len() >= min_batch_size
};
// Procesar si es necesario
// Process if necessary
if should_process {
self.process_batch().await?;
}
@@ -300,17 +300,17 @@ impl IdMappingOptimizer {
Ok(())
}
/// Precargar un conjunto de rutas para obtener sus IDs en batch
/// Preload a set of paths to get their IDs in batch
pub async fn preload_paths(&self, paths: Vec<StoragePath>) -> Result<(), IdMappingError> {
// Solo proceder si hay rutas para cargar
// Only proceed if there are paths to load
if paths.is_empty() {
return Ok(());
}
// Rutas que debemos cargar (las que no están en caché)
// Paths we need to load (those not in cache)
let mut paths_to_load = Vec::new();
// Verificar primero el caché
// Check cache first
{
let cache = self.path_to_id_cache.read().await;
for path in paths {
@@ -321,12 +321,12 @@ impl IdMappingOptimizer {
}
}
// Si todos estaban en caché, terminar
// If all were in cache, finish
if paths_to_load.is_empty() {
return Ok(());
}
// Agregar rutas a la cola para procesamiento batch
// Add paths to queue for batch processing
{
let mut batch_queue = self.pending_batch.lock().await;
for path in paths_to_load {
@@ -334,23 +334,23 @@ impl IdMappingOptimizer {
}
}
// Ejecutar procesamiento batch inmediatamente
// Execute batch processing immediately
self.process_batch().await?;
Ok(())
}
/// Precargar un conjunto de IDs para obtener sus rutas en batch
/// Preload a set of IDs to get their paths in batch
pub async fn preload_ids(&self, ids: Vec<String>) -> Result<(), IdMappingError> {
// Solo proceder si hay IDs para cargar
// Only proceed if there are IDs to load
if ids.is_empty() {
return Ok(());
}
// IDs que debemos cargar (los que no están en caché)
// IDs we need to load (those not in cache)
let mut ids_to_load = Vec::new();
// Verificar primero el caché
// Check cache first
{
let cache = self.id_to_path_cache.read().await;
for id in ids {
@@ -360,12 +360,12 @@ impl IdMappingOptimizer {
}
}
// Si todos estaban en caché, terminar
// If all were in cache, finish
if ids_to_load.is_empty() {
return Ok(());
}
// Agregar IDs a la cola para procesamiento batch
// Add IDs to queue for batch processing
{
let mut batch_queue = self.pending_batch.lock().await;
for id in ids_to_load {
@@ -373,7 +373,7 @@ impl IdMappingOptimizer {
}
}
// Ejecutar procesamiento batch inmediatamente
// Execute batch processing immediately
self.process_batch().await?;
Ok(())
@@ -383,7 +383,7 @@ impl IdMappingOptimizer {
#[async_trait]
impl IdMappingPort for IdMappingOptimizer {
async fn get_or_create_id(&self, path: &StoragePath) -> Result<String, DomainError> {
// Actualizar estadísticas
// Update statistics
{
let mut stats = self.stats.write().await;
stats.get_id_queries += 1;
@@ -391,7 +391,7 @@ impl IdMappingPort for IdMappingOptimizer {
let path_str = path.to_string();
// Verificar primero en el caché
// Check cache first
{
let cache = self.path_to_id_cache.read().await;
if let Some((id, _)) = cache.get(&path_str) {
@@ -443,7 +443,7 @@ impl IdMappingPort for IdMappingOptimizer {
}
async fn get_path_by_id(&self, id: &str) -> Result<StoragePath, DomainError> {
// Actualizar estadísticas
// Update statistics
{
let mut stats = self.stats.write().await;
stats.path_by_id_queries += 1;
@@ -493,21 +493,21 @@ impl IdMappingPort for IdMappingOptimizer {
}
async fn update_path(&self, id: &str, new_path: &StoragePath) -> Result<(), DomainError> {
// Invalidar caché para este ID
// Invalidate cache for this ID
{
let mut id_cache = self.id_to_path_cache.write().await;
let mut path_cache = self.path_to_id_cache.write().await;
// Eliminar la entrada del ID
// Remove the ID entry
if let Some((old_path, _)) = id_cache.remove(id) {
path_cache.remove(&old_path);
}
}
// Actualizar en el servicio base
// Update in the base service
let result = self.base_service.update_path(id, new_path).await?;
// Actualizar caché con el nuevo mapeo
// Update cache with new mapping
{
let mut id_cache = self.id_to_path_cache.write().await;
let mut path_cache = self.path_to_id_cache.write().await;
@@ -523,25 +523,25 @@ impl IdMappingPort for IdMappingOptimizer {
}
async fn remove_id(&self, id: &str) -> Result<(), DomainError> {
// Invalidar caché para este ID
// Invalidate cache for this ID
{
let mut id_cache = self.id_to_path_cache.write().await;
let mut path_cache = self.path_to_id_cache.write().await;
// Eliminar la entrada del ID
// Remove the ID entry
if let Some((path, _)) = id_cache.remove(id) {
path_cache.remove(&path);
}
}
// Eliminar en el servicio base
// Remove from the base service
self.base_service.remove_id(id).await?;
Ok(())
}
async fn save_changes(&self) -> Result<(), DomainError> {
// Delegar al servicio base
// Delegate to the base service
self.base_service.save_changes().await?;
Ok(())
@@ -569,15 +569,15 @@ mod tests {
let path = StoragePath::from_string("/test/file.txt");
// Primera llamada debería usar el servicio base
// First call should use the base service
let id = optimizer.get_or_create_id(&path).await.unwrap();
assert!(!id.is_empty(), "ID should not be empty");
// Segunda llamada debería usar caché
// Second call should use cache
let id2 = optimizer.get_or_create_id(&path).await.unwrap();
assert_eq!(id, id2, "Same path should return same ID");
// Verificar estadísticas de caché
// Verify cache statistics
let stats = optimizer.get_stats().await;
assert_eq!(stats.get_id_queries, 2, "Should have 2 queries");
assert_eq!(stats.get_id_hits, 1, "Should have 1 hit");
@@ -587,27 +587,27 @@ mod tests {
async fn test_batch_processing() {
let (_, optimizer) = create_test_service().await;
// Crear un lote de rutas
// Create a batch of paths
let mut paths = Vec::new();
for i in 0..50 {
paths.push(StoragePath::from_string(&format!("/test/batch/file{}.txt", i)));
}
// Precargar las rutas
// Preload the paths
optimizer.preload_paths(paths.clone()).await.unwrap();
// Verificar que todas están en caché
// Verify all are in cache
for path in &paths {
let id = optimizer.get_or_create_id(path).await.unwrap();
assert!(!id.is_empty(), "ID should be available for path");
}
// Verificar estadísticas
// Verify statistics
let stats = optimizer.get_stats().await;
assert_eq!(stats.batch_operations, 1, "Should have 1 batch operation");
assert!(stats.batch_items_processed >= 50, "Should have processed at least 50 items");
// Verificar que todas las consultas posteriores son hits en caché
// Verify all subsequent queries are cache hits
assert_eq!(stats.get_id_hits, 50, "All subsequente queries should be cache hits");
}
@@ -615,21 +615,21 @@ mod tests {
async fn test_cache_cleanup() {
let (_, optimizer) = create_test_service().await;
// Crear algunas entradas
// Create some entries
let path = StoragePath::from_string("/test/cleanup.txt");
let id = optimizer.get_or_create_id(&path).await.unwrap();
// Verificar estadísticas iniciales
// Verify initial statistics
{
let stats = optimizer.get_stats().await;
assert_eq!(stats.get_id_queries, 1, "Should have 1 query");
assert_eq!(stats.get_id_hits, 0, "Should have 0 hits");
}
// Ejecutar limpieza (no debería eliminar nada todavía)
// Run cleanup (should not remove anything yet)
optimizer.cleanup_cache().await;
// Verificar que el caché sigue funcionando
// Verify cache is still working
let id2 = optimizer.get_or_create_id(&path).await.unwrap();
assert_eq!(id, id2, "Cache should still work after cleanup");
@@ -12,7 +12,7 @@ use crate::common::errors::{DomainError, ErrorKind};
use crate::application::ports::outbound::IdMappingPort;
use crate::common::config::TimeoutConfig;
/// Error específico para el servicio de mapeo de IDs
/// Specific error for the ID mapping service
#[derive(Debug, thiserror::Error)]
pub enum IdMappingError {
#[error("ID not found: {0}")]
@@ -31,7 +31,7 @@ pub enum IdMappingError {
Other(String),
}
// Implementar conversión de IdMappingError a DomainError
// Implement conversion from IdMappingError to DomainError
impl From<IdMappingError> for DomainError {
fn from(err: IdMappingError) -> Self {
match err {
@@ -59,25 +59,25 @@ impl From<IdMappingError> for DomainError {
}
}
/// Estructura para almacenar IDs mapeados a sus rutas
/// Structure to store IDs mapped to their paths
#[derive(Serialize, Deserialize, Debug, Default)]
struct IdMap {
path_to_id: HashMap<String, String>,
id_to_path: HashMap<String, String>, // Campo para búsqueda bidireccional eficiente
version: u32, // Versión para detectar cambios
id_to_path: HashMap<String, String>, // Field for efficient bidirectional lookup
version: u32, // Version to detect changes
}
/// Servicio para gestionar mapeos entre rutas y IDs únicos
/// Service to manage mappings between paths and unique IDs
pub struct IdMappingService {
map_path: PathBuf,
id_map: RwLock<IdMap>,
save_mutex: Mutex<()>, // Para evitar múltiples guardados concurrentes
save_mutex: Mutex<()>, // To prevent multiple concurrent saves
timeouts: TimeoutConfig,
pending_save: RwLock<bool>, // Indica si hay cambios pendientes
pending_save: RwLock<bool>, // Indicates if there are pending changes
}
impl IdMappingService {
/// Crea un nuevo servicio de mapeo de IDs
/// Creates a new ID mapping service
pub async fn new(map_path: PathBuf) -> Result<Self, DomainError> {
let timeouts = TimeoutConfig::default();
let id_map = Self::load_id_map(&map_path, &timeouts).await?;
@@ -91,7 +91,7 @@ impl IdMappingService {
})
}
/// Crea un servicio de mapeo de IDs en memoria (para pruebas)
/// Creates an in-memory ID mapping service (for testing)
///
/// Similar functionality as new_in_memory but with a simpler signature for dummy use
pub fn dummy() -> Self {
@@ -104,7 +104,7 @@ impl IdMappingService {
}
}
/// Crea un servicio de mapeo de IDs en memoria (para pruebas - versión original)
/// Creates an in-memory ID mapping service (for testing - original version)
pub fn new_in_memory() -> Self {
Self {
map_path: PathBuf::from("memory"),
@@ -115,10 +115,10 @@ impl IdMappingService {
}
}
/// Carga el mapa de IDs desde disco con manejo robusto de errores
/// Loads the ID map from disk with robust error handling
async fn load_id_map(map_path: &PathBuf, timeouts: &TimeoutConfig) -> Result<IdMap, DomainError> {
if map_path.exists() {
// Intentar leer con timeout para evitar bloqueos indefinidos
// Try to read with timeout to avoid indefinite blocking
let read_result = time::timeout(
timeouts.lock_timeout(),
fs::read_to_string(map_path)
@@ -127,10 +127,10 @@ impl IdMappingService {
let content = read_result.map_err(|e| DomainError::internal_error("IdMapping", format!("Failed to read ID map from {}: {}", map_path.display(), e)))?;
// Parsear el JSON
// Parse the JSON
match serde_json::from_str::<IdMap>(&content) {
Ok(mut map) => {
// Reconstruir el mapa inverso si es necesario
// Rebuild the inverse map if necessary
if map.id_to_path.is_empty() && !map.path_to_id.is_empty() {
let mut rebuild_count = 0;
for (path, id) in &map.path_to_id {
@@ -146,7 +146,7 @@ impl IdMappingService {
},
Err(e) => {
tracing::error!("Error parsing ID map: {}", e);
// Intentar hacer un respaldo del archivo corrupto
// Try to backup the corrupted file
let backup_path = map_path.with_extension("json.bak");
if let Err(copy_err) = tokio::fs::copy(map_path, &backup_path).await {
tracing::error!("Failed to backup corrupted map file: {}", copy_err);
@@ -158,18 +158,18 @@ impl IdMappingService {
return Ok(IdMap {
path_to_id: HashMap::new(),
id_to_path: HashMap::new(),
version: 1, // Iniciar con versión 1
version: 1, // Start with version 1
});
}
}
}
// Devolver un mapa vacío si el archivo no existe y crear el archivo
// Return an empty map if the file doesn't exist and create the file
tracing::info!("No existing ID map found, creating new empty map");
let empty_map = IdMap {
path_to_id: HashMap::new(),
id_to_path: HashMap::new(),
version: 1, // Iniciar con versión 1
version: 1, // Start with version 1
};
// Ensure directory exists
@@ -198,16 +198,16 @@ impl IdMappingService {
Ok(empty_map)
}
/// Guarda el mapa de IDs en disco de manera segura
/// Saves the ID map to disk safely
async fn save_id_map(&self) -> Result<(), DomainError> {
// Adquirir bloqueo exclusivo para salvar
// Acquire exclusive lock for saving
let _lock = time::timeout(
self.timeouts.lock_timeout(),
self.save_mutex.lock()
).await
.map_err(|_| DomainError::timeout("IdMapping", "Timeout acquiring save lock for ID mapping"))?;
// Crear JSON con el lock de lectura para minimizar el tiempo de bloqueo
// Create JSON with read lock to minimize lock hold time
let json = {
let mut map = time::timeout(
self.timeouts.lock_timeout(),
@@ -215,7 +215,7 @@ impl IdMappingService {
).await
.map_err(|_| DomainError::timeout("IdMapping", "Timeout acquiring write lock for ID mapping"))?;
// Incrementar versión sólo si hay cambios por guardar
// Increment version only if there are pending changes to save
let pending = *self.pending_save.read().await;
if pending {
map.version += 1;
@@ -227,16 +227,16 @@ impl IdMappingService {
.map_err(|e| DomainError::internal_error("IdMapping", format!("Failed to serialize ID map to JSON: {}", e)))?
};
// Escribir a un archivo temporal primero para evitar corrupción
// Write to a temporary file first to avoid corruption
let temp_path = self.map_path.with_extension("json.tmp");
fs::write(&temp_path, &json).await
.map_err(|e| DomainError::internal_error("IdMapping", format!("Failed to write temporary ID map to {}: {}", temp_path.display(), e)))?;
// Realizar el rename atómico
// Perform the atomic rename
fs::rename(&temp_path, &self.map_path).await
.map_err(|e| DomainError::internal_error("IdMapping", format!("Failed to rename temporary ID map to {}: {}", self.map_path.display(), e)))?;
// Resetear flag de pendientes
// Reset pending flag
{
let mut pending = self.pending_save.write().await;
*pending = false;
@@ -246,22 +246,22 @@ impl IdMappingService {
Ok(())
}
/// Genera un ID único
/// Generates a unique ID
fn generate_id(&self) -> String {
Uuid::new_v4().to_string()
}
/// Marca cambios como pendientes
/// Marks changes as pending
async fn mark_pending(&self) {
let mut pending = self.pending_save.write().await;
*pending = true;
}
/// Obtiene el ID para una ruta o genera uno nuevo si no existe
/// Gets the ID for a path or generates a new one if it doesn't exist
pub async fn get_or_create_id(&self, path: &StoragePath) -> Result<String, IdMappingError> {
let path_str = path.to_string();
// Primer intento con lock de lectura (más eficiente)
// First attempt with read lock (more efficient)
{
let read_result = match time::timeout(
self.timeouts.lock_timeout(),
@@ -276,7 +276,7 @@ impl IdMappingService {
}
}
// Si no se encuentra, adquirir lock de escritura
// If not found, acquire write lock
let write_result = match time::timeout(
self.timeouts.lock_timeout(),
self.id_map.write()
@@ -287,18 +287,18 @@ impl IdMappingService {
let mut map = write_result;
// Verificar nuevamente (podría haberse agregado mientras esperábamos el lock)
// Check again (it could have been added while we were waiting for the lock)
if let Some(id) = map.path_to_id.get(&path_str) {
return Ok(id.clone());
}
// Generar un nuevo ID y almacenarlo
// Generate a new ID and store it
let id = self.generate_id();
map.path_to_id.insert(path_str.clone(), id.clone());
map.id_to_path.insert(id.clone(), path_str);
// Marcar como pendiente para guardar
drop(map); // Liberar el write lock antes de adquirir otro
// Mark as pending for saving
drop(map); // Release the write lock before acquiring another
self.mark_pending().await;
tracing::debug!("Created new ID mapping: {} -> {}", path.to_string(), id);
@@ -306,7 +306,7 @@ impl IdMappingService {
Ok(id)
}
/// Obtiene una ruta por su ID con manejo de timeout
/// Gets a path by its ID with timeout handling
pub async fn get_path_by_id(&self, id: &str) -> Result<StoragePath, IdMappingError> {
let read_result = match time::timeout(
self.timeouts.lock_timeout(),
@@ -323,7 +323,7 @@ impl IdMappingService {
Err(IdMappingError::NotFound(id.to_string()))
}
/// Actualiza el mapeo de un ID existente a una nueva ruta
/// Updates the mapping of an existing ID to a new path
pub async fn update_path(&self, id: &str, new_path: &StoragePath) -> Result<(), IdMappingError> {
let write_result = match time::timeout(
self.timeouts.lock_timeout(),
@@ -335,17 +335,17 @@ impl IdMappingService {
let mut map = write_result;
// Buscar la ruta anterior para eliminarla
// Find the previous path to remove it
if let Some(old_path) = map.id_to_path.get(id).cloned() {
map.path_to_id.remove(&old_path);
// Registrar la nueva ruta
// Register the new path
let new_path_str = new_path.to_string();
map.path_to_id.insert(new_path_str.clone(), id.to_string());
map.id_to_path.insert(id.to_string(), new_path_str);
// Marcar como pendiente
drop(map); // Liberar el write lock antes de adquirir otro
// Mark as pending
drop(map); // Release the write lock before acquiring another
self.mark_pending().await;
tracing::debug!("Updated path mapping for ID {}: {} -> {}",
@@ -357,7 +357,7 @@ impl IdMappingService {
}
}
/// Elimina un ID del mapa
/// Removes an ID from the map
pub async fn remove_id(&self, id: &str) -> Result<(), IdMappingError> {
let write_result = match time::timeout(
self.timeouts.lock_timeout(),
@@ -369,12 +369,12 @@ impl IdMappingService {
let mut map = write_result;
// Buscar la ruta para eliminarla
// Find the path to remove it
if let Some(path) = map.id_to_path.remove(id) {
map.path_to_id.remove(&path);
// Marcar como pendiente
drop(map); // Liberar el write lock antes de adquirir otro
// Mark as pending
drop(map); // Release the write lock before acquiring another
self.mark_pending().await;
tracing::debug!("Removed ID mapping: {} -> {}", id, path);
@@ -384,9 +384,9 @@ impl IdMappingService {
}
}
/// Guarda cambios pendientes al disco inmediatamente, sin debounce
/// Saves pending changes to disk immediately, without debounce
pub async fn save_pending_changes(&self) -> Result<(), IdMappingError> {
// Verificar si hay cambios pendientes
// Check if there are pending changes
{
let pending = self.pending_save.read().await;
if !*pending {
@@ -394,12 +394,12 @@ impl IdMappingService {
}
}
// Guardar inmediatamente (sin debounce ni spawn)
// Save immediately (without debounce or spawn)
match self.save_id_map().await {
Ok(_) => {
tracing::info!("ID mappings saved successfully to disk at {}", self.map_path.display());
// Verificar explícitamente que el archivo existe y tiene tamaño
// Explicitly verify that the file exists and has size
match std::fs::metadata(&self.map_path) {
Ok(metadata) => {
if metadata.len() > 0 {
@@ -410,7 +410,7 @@ impl IdMappingService {
},
Err(e) => {
tracing::error!("Failed to verify saved map file: {}", e);
// Intentar un segundo guardado si la verificación falla
// Try a second save if verification fails
if let Err(retry_err) = self.save_id_map().await {
tracing::error!("Second save attempt also failed: {}", retry_err);
return Err(IdMappingError::IoError(std::io::Error::new(
@@ -426,7 +426,7 @@ impl IdMappingService {
},
Err(e) => {
tracing::error!("Failed to save ID map to {}: {}", self.map_path.display(), e);
// Intentar un segundo guardado con retraso en caso de error
// Try a second save with delay in case of error
tokio::time::sleep(tokio::time::Duration::from_millis(100)).await;
match self.save_id_map().await {
Ok(_) => {
@@ -448,31 +448,31 @@ impl IdMappingService {
#[async_trait]
impl IdMappingPort for IdMappingService {
/// Obtiene el ID para una ruta o genera uno nuevo si no existe
/// Gets the ID for a path or generates a new one if it doesn't exist
async fn get_or_create_id(&self, path: &StoragePath) -> Result<String, DomainError> {
self.get_or_create_id(path).await
.map_err(|e| DomainError::internal_error("IdMapping", format!("Failed to get or create ID for path: {}: {}", path.to_string(), e)))
}
/// Obtiene una ruta por su ID con manejo de timeout
/// Gets a path by its ID with timeout handling
async fn get_path_by_id(&self, id: &str) -> Result<StoragePath, DomainError> {
self.get_path_by_id(id).await
.map_err(|e| DomainError::internal_error("IdMapping", format!("Failed to get path for ID: {}: {}", id, e)))
}
/// Actualiza el mapeo de un ID existente a una nueva ruta
/// Updates the mapping of an existing ID to a new path
async fn update_path(&self, id: &str, new_path: &StoragePath) -> Result<(), DomainError> {
self.update_path(id, new_path).await
.map_err(|e| DomainError::internal_error("IdMapping", format!("Failed to update path for ID: {} to {}: {}", id, new_path.to_string(), e)))
}
/// Elimina un ID del mapa
/// Removes an ID from the map
async fn remove_id(&self, id: &str) -> Result<(), DomainError> {
self.remove_id(id).await
.map_err(|e| DomainError::internal_error("IdMapping", format!("Failed to remove ID: {}: {}", id, e)))
}
/// Guarda cambios pendientes al disco
/// Saves pending changes to disk
async fn save_changes(&self) -> Result<(), DomainError> {
self.save_pending_changes().await
.map_err(|e| DomainError::internal_error("IdMapping", format!("Failed to save pending ID mapping changes: {}", e)))
@@ -481,7 +481,7 @@ impl IdMappingPort for IdMappingService {
// The extension methods were moved to the IdMappingPort trait as default implementations
// Implementar Clone para poder usar en tokio::spawn
// Implement Clone to allow use in tokio::spawn
/// Synchronous helper for contexts where we can't use async
impl IdMappingService {
/// Create a new service synchronously (only for stubs and initialization)
@@ -499,13 +499,13 @@ impl IdMappingService {
impl Clone for IdMappingService {
fn clone(&self) -> Self {
// No podemos clonar directamente los RwLock/Mutex,
// pero podemos crear nuevas instancias que apunten al mismo Arc interno
// Sin embargo, en este caso simplemente necesitamos la map_path
// We cannot directly clone the RwLock/Mutex,
// but we can create new instances that point to the same internal Arc
// However, in this case we simply need the map_path
Self {
map_path: self.map_path.clone(),
id_map: RwLock::new(IdMap::default()), // Esto no se usa en el task asíncrono
save_mutex: Mutex::new(()), // Esto tampoco
id_map: RwLock::new(IdMap::default()), // This is not used in the async task
save_mutex: Mutex::new(()), // Neither is this
timeouts: self.timeouts.clone(),
pending_save: RwLock::new(false),
}
@@ -530,7 +530,7 @@ mod tests {
assert!(!id.is_empty(), "ID should not be empty");
// Verificar que el mismo ID se devuelve para la misma ruta
// Verify that the same ID is returned for the same path
let id2 = service.get_or_create_id(&path).await.unwrap();
assert_eq!(id, id2, "Same path should return same ID");
}
@@ -557,7 +557,7 @@ mod tests {
let temp_dir = tempdir().unwrap();
let map_path = temp_dir.path().join("id_map.json");
// Crear y poblar el servicio
// Create and populate the service
let service = IdMappingService::new(map_path.clone()).await.unwrap();
let path1 = StoragePath::from_string("/test/file1.txt");
@@ -565,16 +565,16 @@ mod tests {
let id1 = service.get_or_create_id(&path1).await.unwrap();
let id2 = service.get_or_create_id(&path2).await.unwrap();
// Guardar cambios
// Save changes
service.save_pending_changes().await.unwrap();
// Esperar para asegurar que el guardado asíncrono termine
// Wait to ensure the async save completes
tokio::time::sleep(Duration::from_millis(500)).await;
// Crear un nuevo servicio que debería cargar el mismo mapa
// Create a new service that should load the same map
let service2 = IdMappingService::new(map_path).await.unwrap();
// Verificar que los IDs coinciden
// Verify that the IDs match
let loaded_id1 = service2.get_or_create_id(&path1).await.unwrap();
let loaded_id2 = service2.get_or_create_id(&path2).await.unwrap();
@@ -591,7 +591,7 @@ mod tests {
let service = std::sync::Arc::new(IdMappingService::new(map_path).await.unwrap());
// Crear múltiples tareas que intentan acceder simultáneamente
// Create multiple tasks that attempt simultaneous access
let mut tasks = Vec::new();
for i in 0..100 {
let path = StoragePath::from_string(&format!("/test/concurrent/file{}.txt", i));
@@ -602,15 +602,15 @@ mod tests {
}));
}
// Esperar a que todas terminen
// Wait for all to finish
let results = join_all(tasks).await;
// Verificar que todas tuvieron éxito
// Verify that all succeeded
for result in results {
assert!(result.unwrap().is_ok(), "Concurrent operations should succeed");
}
// Guardar cambios
// Save changes
service.save_pending_changes().await.unwrap();
}
}
+3 -3
View File
@@ -110,7 +110,7 @@ impl TokenServicePort for JwtTokenService {
DomainError::new(
ErrorKind::InternalError,
"TokenService",
format!("Error al generar token: {}", e)
format!("Error generating token: {}", e)
)
})
}
@@ -126,12 +126,12 @@ impl TokenServicePort for JwtTokenService {
.map_err(|e| {
match e.kind() {
jsonwebtoken::errors::ErrorKind::ExpiredSignature => {
DomainError::new(ErrorKind::AccessDenied, "TokenService", "Token expirado")
DomainError::new(ErrorKind::AccessDenied, "TokenService", "Token expired")
},
_ => DomainError::new(
ErrorKind::AccessDenied,
"TokenService",
format!("Token inválido: {}", e)
format!("Invalid token: {}", e)
),
}
})?;
+18 -18
View File
@@ -1,9 +1,9 @@
//! PathService - Servicio de infraestructura para manejo de rutas de almacenamiento
//! PathService - Infrastructure service for storage path management
//!
//! Este servicio fue movido desde domain/services porque implementa traits de application
//! (StoragePort, StorageMediator) y tiene dependencias de sistema de archivos (tokio::fs).
//! This service was moved from domain/services because it implements application traits
//! (StoragePort, StorageMediator) and has file system dependencies (tokio::fs).
//!
//! StoragePath (Value Object) permanece en domain/services/path_service.rs
//! StoragePath (Value Object) remains in domain/services/path_service.rs
use std::path::{Path, PathBuf};
use async_trait::async_trait;
@@ -15,18 +15,18 @@ use crate::application::services::storage_mediator::{StorageMediator, StorageMed
use crate::domain::entities::folder::Folder;
use crate::domain::services::path_service::StoragePath;
/// Servicio de infraestructura para manejar operaciones con rutas de almacenamiento
/// Infrastructure service for handling storage path operations
pub struct PathService {
root_path: PathBuf,
}
impl PathService {
/// Crea un nuevo servicio de rutas con una raíz específica
/// Creates a new path service with a specific root
pub fn new(root_path: PathBuf) -> Self {
Self { root_path }
}
/// Convierte una ruta del dominio a una ruta física absoluta
/// Converts a domain path to an absolute physical path
pub fn resolve_path(&self, storage_path: &StoragePath) -> PathBuf {
let mut path = self.root_path.clone();
for segment in storage_path.segments() {
@@ -35,7 +35,7 @@ impl PathService {
path
}
/// Convierte una ruta física a una ruta de dominio
/// Converts a physical path to a domain path
pub fn to_storage_path(&self, physical_path: &Path) -> Option<StoragePath> {
physical_path.strip_prefix(&self.root_path).ok().map(|rel_path| {
let segments: Vec<String> = rel_path
@@ -49,12 +49,12 @@ impl PathService {
})
}
/// Crea una ruta de archivo dentro de una carpeta
/// Creates a file path within a folder
pub fn create_file_path(&self, folder_path: &StoragePath, file_name: &str) -> StoragePath {
folder_path.join(file_name)
}
/// Verifica si una ruta es directamente hija de otra
/// Checks if a path is a direct child of another
pub fn is_direct_child(&self, parent_path: &StoragePath, potential_child: &StoragePath) -> bool {
if let Some(child_parent) = potential_child.parent() {
&child_parent == parent_path
@@ -63,7 +63,7 @@ impl PathService {
}
}
/// Verifica si una ruta está en la raíz
/// Checks if a path is at the root
pub fn is_in_root(&self, path: &StoragePath) -> bool {
path.parent().map_or(true, |p| p.is_empty())
}
@@ -73,9 +73,9 @@ impl PathService {
&self.root_path
}
/// Valida una ruta para asegurar que no contiene componentes peligrosos
/// Validates a path to ensure it doesn't contain dangerous components
pub fn validate_path(&self, path: &StoragePath) -> Result<(), DomainError> {
// Verificar que no haya segmentos vacíos
// Check for empty segments
if path.segments().iter().any(|s| s.is_empty()) {
return Err(DomainError::new(
ErrorKind::InvalidInput,
@@ -84,7 +84,7 @@ impl PathService {
));
}
// Verificar que no haya caracteres peligrosos
// Check for dangerous characters
let dangerous_chars = ['\\', ':', '*', '?', '"', '<', '>', '|'];
for segment in path.segments() {
if segment.contains(&dangerous_chars[..]) {
@@ -95,7 +95,7 @@ impl PathService {
));
}
// Verificar que no empiece con . (oculto en Unix)
// Check that it doesn't start with . (hidden in Unix)
if segment.starts_with('.') && segment != ".well-known" {
return Err(DomainError::new(
ErrorKind::InvalidInput,
@@ -120,13 +120,13 @@ impl StoragePort for PathService {
}
async fn ensure_directory(&self, storage_path: &StoragePath) -> Result<(), DomainError> {
// Primero validar la ruta
// First validate the path
self.validate_path(storage_path)?;
// Resolver a ruta física
// Resolve to physical path
let physical_path = self.resolve_path(storage_path);
// Crear directorios si no existen
// Create directories if they don't exist
if !physical_path.exists() {
fs::create_dir_all(&physical_path).await
.map_err(|e| DomainError::new(
@@ -7,7 +7,7 @@ use crate::common::errors::Result;
use crate::domain::repositories::trash_repository::TrashRepository;
use crate::application::ports::trash_ports::TrashUseCase;
/// Servicio para la limpieza automática de elementos expirados en la papelera
/// Service for automatic cleanup of expired items in the trash
pub struct TrashCleanupService {
trash_service: Arc<dyn TrashUseCase>,
trash_repository: Arc<dyn TrashRepository>,
@@ -23,75 +23,75 @@ impl TrashCleanupService {
Self {
trash_service,
trash_repository,
cleanup_interval_hours: cleanup_interval_hours.max(1), // Mínimo 1 hora
cleanup_interval_hours: cleanup_interval_hours.max(1), // Minimum 1 hour
}
}
/// Inicia el trabajo de limpieza periódica
/// Starts the periodic cleanup job
#[instrument(skip(self))]
pub async fn start_cleanup_job(&self) {
let trash_repository = self.trash_repository.clone();
let trash_service = self.trash_service.clone();
let interval_hours = self.cleanup_interval_hours;
info!("Iniciando trabajo de limpieza de papelera con intervalo de {} horas", interval_hours);
info!("Starting trash cleanup job with interval of {} hours", interval_hours);
tokio::spawn(async move {
let interval_duration = Duration::from_secs(interval_hours * 60 * 60);
let mut interval = time::interval(interval_duration);
// Primera ejecución inmediata
// First immediate execution
Self::cleanup_expired_items(trash_repository.clone(), trash_service.clone()).await
.unwrap_or_else(|e| error!("Error en la limpieza inicial de la papelera: {:?}", e));
.unwrap_or_else(|e| error!("Error in initial trash cleanup: {:?}", e));
loop {
interval.tick().await;
debug!("Ejecutando tarea programada de limpieza de papelera");
debug!("Running scheduled trash cleanup task");
if let Err(e) = Self::cleanup_expired_items(
trash_repository.clone(),
trash_service.clone()
).await {
error!("Error en la limpieza programada de la papelera: {:?}", e);
error!("Error in scheduled trash cleanup: {:?}", e);
}
}
});
}
/// Limpia los elementos expirados en la papelera
/// Cleans up expired items in the trash
#[instrument(skip(trash_repository, trash_service))]
async fn cleanup_expired_items(
trash_repository: Arc<dyn TrashRepository>,
trash_service: Arc<dyn TrashUseCase>,
) -> Result<()> {
debug!("Comenzando limpieza de elementos expirados en la papelera");
debug!("Starting cleanup of expired items in the trash");
// Obtener todos los elementos expirados
// Get all expired items
let expired_items = trash_repository.get_expired_items().await?;
if expired_items.is_empty() {
debug!("No hay elementos expirados para limpiar");
debug!("No expired items to clean up");
return Ok(());
}
info!("Encontrados {} elementos expirados para eliminar", expired_items.len());
info!("Found {} expired items to delete", expired_items.len());
// Eliminar cada elemento expirado
// Delete each expired item
for item in expired_items {
let trash_id = item.id().to_string();
let user_id = item.user_id().to_string();
debug!("Eliminando elemento expirado: id={}, user={}", trash_id, user_id);
debug!("Deleting expired item: id={}, user={}", trash_id, user_id);
// Si falla una eliminación, continuar con las demás
// If a deletion fails, continue with the rest
if let Err(e) = trash_service.delete_permanently(&trash_id, &user_id).await {
error!("Error eliminando elemento expirado {}: {:?}", trash_id, e);
error!("Error deleting expired item {}: {:?}", trash_id, e);
} else {
debug!("Elemento expirado eliminado correctamente: {}", trash_id);
debug!("Expired item deleted successfully: {}", trash_id);
}
}
info!("Limpieza de papelera completada");
info!("Trash cleanup completed");
Ok(())
}
}
+47 -47
View File
@@ -13,47 +13,47 @@ use crate::{
};
use std::sync::Arc;
/// Error relacionado con la creación de archivos ZIP
/// Error related to ZIP file creation
#[derive(Debug, Error)]
pub enum ZipError {
#[error("Error de IO: {0}")]
#[error("IO error: {0}")]
IoError(#[from] std::io::Error),
#[error("Error de ZIP: {0}")]
#[error("ZIP error: {0}")]
ZipError(#[from] zip::result::ZipError),
#[error("Error al leer el archivo: {0}")]
#[error("Error reading file: {0}")]
FileReadError(String),
#[error("Error al obtener contenido de carpeta: {0}")]
#[error("Error getting folder contents: {0}")]
FolderContentsError(String),
#[error("Carpeta no encontrada: {0}")]
#[error("Folder not found: {0}")]
FolderNotFound(String),
}
// Implementar From<ZipError> para DomainError para permitir el uso de ?
// Implement From<ZipError> for DomainError to allow the use of ?
impl From<ZipError> for DomainError {
fn from(err: ZipError) -> Self {
DomainError::new(ErrorKind::InternalError, "zip_service", err.to_string())
}
}
// Implementar From<zip::result::ZipError> para DomainError directamente
// Implement From<zip::result::ZipError> for DomainError directly
impl From<zip::result::ZipError> for DomainError {
fn from(err: zip::result::ZipError) -> Self {
DomainError::new(ErrorKind::InternalError, "zip_service", err.to_string())
}
}
/// Servicio para crear archivos ZIP
/// Service for creating ZIP files
pub struct ZipService {
file_service: Arc<dyn FileRetrievalUseCase>,
folder_service: Arc<dyn FolderUseCase>,
}
impl ZipService {
/// Crea una nueva instancia del servicio ZIP con una referencia al servicio de archivos
/// Creates a new instance of the ZIP service with a reference to the file service
pub fn new(file_service: Arc<dyn FileRetrievalUseCase>, folder_service: Arc<dyn FolderUseCase>) -> Self {
Self {
file_service,
@@ -61,33 +61,33 @@ impl ZipService {
}
}
/// Crea un archivo ZIP con el contenido de una carpeta y todas sus subcarpetas
/// Retorna los bytes del ZIP
/// Creates a ZIP file with the contents of a folder and all its subfolders
/// Returns the ZIP bytes
pub async fn create_folder_zip(&self, folder_id: &str, folder_name: &str) -> Result<Vec<u8>> {
info!("Creando ZIP para carpeta: {} (ID: {})", folder_name, folder_id);
info!("Creating ZIP for folder: {} (ID: {})", folder_name, folder_id);
// Verificar si la carpeta existe
// Verify if the folder exists
let folder = match self.folder_service.get_folder(folder_id).await {
Ok(folder) => folder,
Err(e) => {
error!("Error al obtener carpeta {}: {}", folder_id, e);
error!("Error getting folder {}: {}", folder_id, e);
return Err(ZipError::FolderNotFound(folder_id.to_string()).into());
}
};
// Crear un buffer en memoria para el ZIP
// Create an in-memory buffer for the ZIP
let buf = Cursor::new(Vec::new());
let mut zip = ZipWriter::new(buf);
// Establecer opciones de compresión
// Set compression options
let options = SimpleFileOptions::default()
.compression_method(zip::CompressionMethod::Deflated)
.unix_permissions(0o755);
// Objeto para seguir las carpetas procesadas y evitar ciclos
// Object to track processed folders and avoid cycles
let mut processed_folders = std::collections::HashSet::new();
// Procesamos la carpeta raíz y construimos el ZIP
// Process the root folder and build the ZIP
self.process_folder_recursively(
&mut zip,
&folder,
@@ -96,20 +96,20 @@ impl ZipService {
&mut processed_folders
).await?;
// Finalizar el ZIP y obtener los bytes
// Finalize the ZIP and get the bytes
let mut zip_buf = zip.finish()?;
let mut bytes = Vec::new();
match zip_buf.read_to_end(&mut bytes) {
Ok(_) => Ok(bytes),
Err(e) => {
error!("Error al leer ZIP finalizado: {}", e);
error!("Error reading finalized ZIP: {}", e);
Err(ZipError::IoError(e).into())
}
}
}
// Implementación alternativa para evitar recursión en async
// Alternative implementation to avoid recursion in async
async fn process_folder_recursively(
&self,
zip: &mut ZipWriter<Cursor<Vec<u8>>>,
@@ -118,63 +118,63 @@ impl ZipService {
options: &SimpleFileOptions,
processed_folders: &mut std::collections::HashSet<String>
) -> Result<()> {
// Estructura para representar el trabajo pendiente
// Structure to represent pending work
struct PendingFolder {
folder: FolderDto,
path: String,
}
// Cola de trabajo para procesamiento iterativo
// Work queue for iterative processing
let mut work_queue = vec![PendingFolder {
folder: folder.clone(),
path: path.to_string(),
}];
// Procesar la cola mientras haya elementos
// Process the queue while there are elements
while let Some(current) = work_queue.pop() {
let folder_id = current.folder.id.to_string();
// Evitar ciclos
// Avoid cycles
if processed_folders.contains(&folder_id) {
continue;
}
processed_folders.insert(folder_id.clone());
// Crear la entrada de directorio en el ZIP
// Create the directory entry in the ZIP
let folder_path = format!("{}/", current.path);
match zip.add_directory(&folder_path, *options) {
Ok(_) => debug!("Carpeta agregada al ZIP: {}", folder_path),
Ok(_) => debug!("Folder added to ZIP: {}", folder_path),
Err(e) => {
warn!("No se pudo agregar carpeta al ZIP (puede que ya exista): {}", e);
// Continuamos aunque falle crear el directorio (podría estar duplicado)
warn!("Could not add folder to ZIP (it may already exist): {}", e);
// Continue even if creating the directory fails (it could be a duplicate)
}
}
// Agregar archivos de la carpeta al ZIP
// Add files from the folder to the ZIP
let files = match self.file_service.list_files(Some(&folder_id)).await {
Ok(files) => files,
Err(e) => {
error!("Error al listar archivos en carpeta {}: {}", folder_id, e);
return Err(ZipError::FolderContentsError(format!("Error al listar archivos: {}", e)).into());
error!("Error listing files in folder {}: {}", folder_id, e);
return Err(ZipError::FolderContentsError(format!("Error listing files: {}", e)).into());
}
};
// Agregar cada archivo al ZIP
// Add each file to the ZIP
for file in files {
self.add_file_to_zip(zip, &file, &folder_path, options).await?;
}
// Procesar subcarpetas
// Process subfolders
let subfolders = match self.folder_service.list_folders(Some(&folder_id)).await {
Ok(folders) => folders,
Err(e) => {
error!("Error al listar subcarpetas en {}: {}", folder_id, e);
return Err(ZipError::FolderContentsError(format!("Error al listar subcarpetas: {}", e)).into());
error!("Error listing subfolders in {}: {}", folder_id, e);
return Err(ZipError::FolderContentsError(format!("Error listing subfolders: {}", e)).into());
}
};
// Agregar subcarpetas a la cola
// Add subfolders to the queue
for subfolder in subfolders {
let subfolder_path = format!("{}/{}", current.path, subfolder.name);
work_queue.push(PendingFolder {
@@ -187,7 +187,7 @@ impl ZipService {
Ok(())
}
// Agrega un archivo al ZIP
// Adds a file to the ZIP
async fn add_file_to_zip(
&self,
zip: &mut ZipWriter<Cursor<Vec<u8>>>,
@@ -196,34 +196,34 @@ impl ZipService {
options: &SimpleFileOptions,
) -> Result<()> {
let file_path = format!("{}{}", folder_path, file.name);
info!("Agregando archivo al ZIP: {}", file_path);
info!("Adding file to ZIP: {}", file_path);
// Obtener el contenido del archivo
// Get the file content
let file_id = file.id.to_string();
let content = match self.file_service.get_file_content(&file_id).await {
Ok(content) => content,
Err(e) => {
error!("Error al leer contenido del archivo {}: {}", file_id, e);
return Err(ZipError::FileReadError(format!("Error al leer archivo {}: {}", file_id, e)).into());
error!("Error reading file content {}: {}", file_id, e);
return Err(ZipError::FileReadError(format!("Error reading file {}: {}", file_id, e)).into());
}
};
// Escribir archivo al ZIP
// Write file to the ZIP
match zip.start_file_from_path(std::path::Path::new(&file_path), *options) {
Ok(_) => {
match zip.write_all(&content) {
Ok(_) => {
debug!("Archivo agregado al ZIP: {}", file_path);
debug!("File added to ZIP: {}", file_path);
Ok(())
},
Err(e) => {
error!("Error al escribir contenido del archivo {}: {}", file_path, e);
error!("Error writing file content {}: {}", file_path, e);
Err(ZipError::IoError(e).into())
}
}
},
Err(e) => {
error!("Error al iniciar archivo en ZIP {}: {}", file_path, e);
error!("Error starting file in ZIP {}: {}", file_path, e);
Err(ZipError::ZipError(e).into())
}
}