fix auth errors and add primigenial paper trash
This commit is contained in:
@@ -115,9 +115,9 @@ impl IdMappingService {
|
||||
timeouts.lock_timeout(),
|
||||
fs::read_to_string(map_path)
|
||||
).await
|
||||
.with_context(|| format!("Timeout reading ID map from {}", map_path.display()))?;
|
||||
.map_err(|_| DomainError::timeout("IdMapping", format!("Timeout reading ID map from {}", map_path.display())))?;
|
||||
|
||||
let content = read_result.with_context(|| format!("Failed to read ID map from {}", map_path.display()))?;
|
||||
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
|
||||
match serde_json::from_str::<IdMap>(&content) {
|
||||
@@ -197,7 +197,7 @@ impl IdMappingService {
|
||||
self.timeouts.lock_timeout(),
|
||||
self.save_mutex.lock()
|
||||
).await
|
||||
.with_context(|| "Timeout acquiring save lock for ID mapping")?;
|
||||
.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
|
||||
let json = {
|
||||
@@ -205,7 +205,7 @@ impl IdMappingService {
|
||||
self.timeouts.lock_timeout(),
|
||||
self.id_map.write()
|
||||
).await
|
||||
.with_context(|| "Timeout acquiring write lock for ID mapping")?;
|
||||
.map_err(|_| DomainError::timeout("IdMapping", "Timeout acquiring write lock for ID mapping"))?;
|
||||
|
||||
// Incrementar versión sólo si hay cambios por guardar
|
||||
let pending = *self.pending_save.read().await;
|
||||
@@ -216,17 +216,17 @@ impl IdMappingService {
|
||||
|
||||
// Use serde with reasonably safe defaults
|
||||
serde_json::to_string_pretty(&*map)
|
||||
.with_context(|| "Failed to serialize ID map to JSON")?
|
||||
.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
|
||||
let temp_path = self.map_path.with_extension("json.tmp");
|
||||
fs::write(&temp_path, &json).await
|
||||
.with_context(|| format!("Failed to write temporary ID map to {}", temp_path.display()))?;
|
||||
.map_err(|e| DomainError::internal_error("IdMapping", format!("Failed to write temporary ID map to {}: {}", temp_path.display(), e)))?;
|
||||
|
||||
// Realizar el rename atómico
|
||||
fs::rename(&temp_path, &self.map_path).await
|
||||
.with_context(|| format!("Failed to rename temporary ID map to {}", self.map_path.display()))?;
|
||||
.map_err(|e| DomainError::internal_error("IdMapping", format!("Failed to rename temporary ID map to {}: {}", self.map_path.display(), e)))?;
|
||||
|
||||
// Resetear flag de pendientes
|
||||
{
|
||||
@@ -408,34 +408,36 @@ impl IdMappingPort for IdMappingService {
|
||||
/// Obtiene el ID para una ruta o genera uno nuevo si no existe
|
||||
async fn get_or_create_id(&self, path: &StoragePath) -> Result<String, DomainError> {
|
||||
self.get_or_create_id(path).await
|
||||
.with_context(|| format!("Failed to get or create ID for path: {}", path.to_string()))
|
||||
.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
|
||||
async fn get_path_by_id(&self, id: &str) -> Result<StoragePath, DomainError> {
|
||||
self.get_path_by_id(id).await
|
||||
.with_context(|| format!("Failed to get path for ID: {}", id))
|
||||
.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
|
||||
async fn update_path(&self, id: &str, new_path: &StoragePath) -> Result<(), DomainError> {
|
||||
self.update_path(id, new_path).await
|
||||
.with_context(|| format!("Failed to update path for ID: {} to {}", id, new_path.to_string()))
|
||||
.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
|
||||
async fn remove_id(&self, id: &str) -> Result<(), DomainError> {
|
||||
self.remove_id(id).await
|
||||
.with_context(|| format!("Failed to remove ID: {}", id))
|
||||
.map_err(|e| DomainError::internal_error("IdMapping", format!("Failed to remove ID: {}: {}", id, e)))
|
||||
}
|
||||
|
||||
/// Guarda cambios pendientes al disco
|
||||
async fn save_changes(&self) -> Result<(), DomainError> {
|
||||
self.save_pending_changes().await
|
||||
.with_context(|| "Failed to save pending ID mapping changes")
|
||||
.map_err(|e| DomainError::internal_error("IdMapping", format!("Failed to save pending ID mapping changes: {}", e)))
|
||||
}
|
||||
}
|
||||
|
||||
// The extension methods were moved to the IdMappingPort trait as default implementations
|
||||
|
||||
// Implementar Clone para poder usar en tokio::spawn
|
||||
/// Synchronous helper for contexts where we can't use async
|
||||
impl IdMappingService {
|
||||
|
||||
@@ -4,4 +4,5 @@ pub mod id_mapping_optimizer;
|
||||
pub mod cache_manager;
|
||||
pub mod file_metadata_cache;
|
||||
pub mod compression_service;
|
||||
pub mod buffer_pool;
|
||||
pub mod buffer_pool;
|
||||
pub mod trash_cleanup_service;
|
||||
@@ -0,0 +1,97 @@
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
use tokio::time;
|
||||
use tracing::{debug, error, info, instrument};
|
||||
|
||||
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
|
||||
pub struct TrashCleanupService {
|
||||
trash_service: Arc<dyn TrashUseCase>,
|
||||
trash_repository: Arc<dyn TrashRepository>,
|
||||
cleanup_interval_hours: u64,
|
||||
}
|
||||
|
||||
impl TrashCleanupService {
|
||||
pub fn new(
|
||||
trash_service: Arc<dyn TrashUseCase>,
|
||||
trash_repository: Arc<dyn TrashRepository>,
|
||||
cleanup_interval_hours: u64,
|
||||
) -> Self {
|
||||
Self {
|
||||
trash_service,
|
||||
trash_repository,
|
||||
cleanup_interval_hours: cleanup_interval_hours.max(1), // Mínimo 1 hora
|
||||
}
|
||||
}
|
||||
|
||||
/// Inicia el trabajo de limpieza periódica
|
||||
#[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);
|
||||
|
||||
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
|
||||
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));
|
||||
|
||||
loop {
|
||||
interval.tick().await;
|
||||
debug!("Ejecutando tarea programada de limpieza de papelera");
|
||||
|
||||
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);
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
/// Limpia los elementos expirados en la papelera
|
||||
#[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");
|
||||
|
||||
// Obtener todos los elementos expirados
|
||||
let expired_items = trash_repository.get_expired_items().await?;
|
||||
|
||||
if expired_items.is_empty() {
|
||||
debug!("No hay elementos expirados para limpiar");
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
info!("Encontrados {} elementos expirados para eliminar", expired_items.len());
|
||||
|
||||
// Eliminar cada elemento expirado
|
||||
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);
|
||||
|
||||
// Si falla una eliminación, continuar con las demás
|
||||
if let Err(e) = trash_service.delete_permanently(&trash_id, &user_id).await {
|
||||
error!("Error eliminando elemento expirado {}: {:?}", trash_id, e);
|
||||
} else {
|
||||
debug!("Elemento expirado eliminado correctamente: {}", trash_id);
|
||||
}
|
||||
}
|
||||
|
||||
info!("Limpieza de papelera completada");
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user