adding features

This commit is contained in:
DioCrafts
2025-03-19 19:52:12 +01:00
parent d9bbd575d2
commit 6e055c3043
19 changed files with 1150 additions and 42 deletions
@@ -0,0 +1,184 @@
use std::path::PathBuf;
use std::sync::Arc;
use async_trait::async_trait;
use bytes::Bytes;
use futures::Stream;
use crate::domain::entities::file::File;
use crate::application::ports::storage_ports::FileReadPort;
use crate::common::errors::DomainError;
use crate::domain::repositories::file_repository::FileRepositoryResult;
use crate::infrastructure::repositories::file_metadata_manager::{FileMetadataManager, MetadataError};
use crate::infrastructure::repositories::file_path_resolver::FilePathResolver;
use crate::domain::services::path_service::StoragePath;
use crate::infrastructure::repositories::parallel_file_processor::ParallelFileProcessor;
use crate::common::config::AppConfig;
/// Implementación de repositorio para operaciones de lectura de archivos
pub struct FileFsReadRepository {
root_path: PathBuf,
metadata_manager: Arc<FileMetadataManager>,
path_resolver: Arc<FilePathResolver>,
config: AppConfig,
parallel_processor: Option<Arc<ParallelFileProcessor>>,
}
impl FileFsReadRepository {
/// Crea un nuevo repositorio de lectura de archivos
pub fn new(
root_path: PathBuf,
metadata_manager: Arc<FileMetadataManager>,
path_resolver: Arc<FilePathResolver>,
config: AppConfig,
parallel_processor: Option<Arc<ParallelFileProcessor>>,
) -> Self {
Self {
root_path,
metadata_manager,
path_resolver,
config,
parallel_processor,
}
}
/// Crea una entidad de archivo a partir de metadatos
async fn create_file_entity(
&self,
id: String,
name: String,
storage_path: StoragePath,
size: u64,
mime_type: String,
folder_id: Option<String>,
created_at: Option<u64>,
modified_at: Option<u64>,
) -> FileRepositoryResult<File> {
// If timestamps are provided, use them; otherwise, let File::new create default timestamps
if let (Some(created), Some(modified)) = (created_at, modified_at) {
File::with_timestamps(
id,
name,
storage_path,
size,
mime_type,
folder_id,
created,
modified,
)
.map_err(|e| crate::domain::repositories::file_repository::FileRepositoryError::Other(e.to_string()))
} else {
File::new(
id,
name,
storage_path,
size,
mime_type,
folder_id,
)
.map_err(|e| crate::domain::repositories::file_repository::FileRepositoryError::Other(e.to_string()))
}
}
/// Obtiene un archivo por su ID
async fn get_file_by_id(&self, id: &str) -> FileRepositoryResult<File> {
// Obtener la ruta del archivo usando el resolver de rutas
let storage_path = self.path_resolver.get_path_by_id(id).await?;
// Verificar que el archivo existe físicamente
let abs_path = self.path_resolver.resolve_storage_path(&storage_path);
if !self.metadata_manager.file_exists(&abs_path).await
.map_err(|e| crate::domain::repositories::file_repository::FileRepositoryError::Other(e.to_string()))? {
return Err(crate::domain::repositories::file_repository::FileRepositoryError::NotFound(
format!("File {} not found at {}", id, storage_path.to_string())
));
}
// Obtener metadatos del archivo
let (size, created_at, modified_at) = self.metadata_manager.get_file_metadata(&abs_path).await
.map_err(|e| match e {
MetadataError::IoError(io_err) => crate::domain::repositories::file_repository::FileRepositoryError::IoError(io_err),
MetadataError::Timeout(msg) => crate::domain::repositories::file_repository::FileRepositoryError::Timeout(msg),
MetadataError::Unavailable(msg) => crate::domain::repositories::file_repository::FileRepositoryError::NotFound(msg),
})?;
// Obtener nombre del archivo de la ruta
let name = match storage_path.file_name() {
Some(name) => name,
None => {
return Err(crate::domain::repositories::file_repository::FileRepositoryError::InvalidPath(
storage_path.to_string()
));
}
};
// Determinar ID de carpeta padre
let parent = storage_path.parent();
let folder_id: Option<String> = if parent.is_none() || parent.as_ref().unwrap().is_empty() {
None // Root folder
} else {
None // En implementación real, buscar ID de la carpeta padre
};
// Determinar tipo MIME
let mime_type = mime_guess::from_path(&abs_path)
.first_or_octet_stream()
.to_string();
// Crear entidad de archivo
let file = self.create_file_entity(
id.to_string(),
name,
storage_path,
size,
mime_type,
folder_id,
Some(created_at),
Some(modified_at),
).await?;
Ok(file)
}
}
#[async_trait]
impl FileReadPort for FileFsReadRepository {
async fn get_file(&self, id: &str) -> Result<File, DomainError> {
self.get_file_by_id(id).await
.map_err(|e| match e {
crate::domain::repositories::file_repository::FileRepositoryError::NotFound(msg) => DomainError::not_found("File", msg),
crate::domain::repositories::file_repository::FileRepositoryError::IoError(io_err) => DomainError::internal_error("File", io_err.to_string()),
crate::domain::repositories::file_repository::FileRepositoryError::Timeout(msg) => DomainError::internal_error("File", msg),
_ => DomainError::internal_error("File", e.to_string()),
})
}
async fn list_files(&self, folder_id: Option<&str>) -> Result<Vec<File>, DomainError> {
// Implementación real debe obtener la lista de archivos en una carpeta
// Por ahora, devolvemos lista vacía
Ok(Vec::new())
}
async fn get_file_content(&self, id: &str) -> Result<Vec<u8>, DomainError> {
// Primero obtenemos el archivo para verificar existencia
let file = self.get_file_by_id(id).await
.map_err(|e| match e {
crate::domain::repositories::file_repository::FileRepositoryError::NotFound(msg) => DomainError::not_found("File", msg),
crate::domain::repositories::file_repository::FileRepositoryError::IoError(io_err) => DomainError::internal_error("File", io_err.to_string()),
crate::domain::repositories::file_repository::FileRepositoryError::Timeout(msg) => DomainError::internal_error("File", msg),
_ => DomainError::internal_error("File", e.to_string()),
})?;
// Ruta absoluta del archivo
let abs_path = self.path_resolver.resolve_storage_path(file.storage_path());
// Implementación real debe leer el contenido del archivo
// Por ahora, devolvemos un vector vacío
Ok(Vec::new())
}
async fn get_file_stream(&self, id: &str) -> Result<Box<dyn Stream<Item = Result<Bytes, std::io::Error>> + Send>, DomainError> {
// Implementación real debe devolver un stream de bytes del archivo
// Por ahora, lanzamos un error
Err(DomainError::internal_error("File stream", "Stream functionality not yet implemented"))
}
}
@@ -0,0 +1,132 @@
use std::path::PathBuf;
use std::sync::Arc;
use async_trait::async_trait;
use crate::domain::entities::file::File;
use crate::application::ports::storage_ports::FileWritePort;
use crate::common::errors::DomainError;
use crate::domain::repositories::file_repository::FileRepositoryResult;
use crate::infrastructure::repositories::file_metadata_manager::{FileMetadataManager, MetadataError};
use crate::infrastructure::repositories::file_path_resolver::FilePathResolver;
use crate::domain::services::path_service::StoragePath;
use crate::infrastructure::repositories::parallel_file_processor::ParallelFileProcessor;
use crate::common::config::AppConfig;
use crate::application::services::storage_mediator::StorageMediator;
/// Implementación de repositorio para operaciones de escritura de archivos
pub struct FileFsWriteRepository {
root_path: PathBuf,
metadata_manager: Arc<FileMetadataManager>,
path_resolver: Arc<FilePathResolver>,
storage_mediator: Arc<dyn StorageMediator>,
config: AppConfig,
parallel_processor: Option<Arc<ParallelFileProcessor>>,
}
impl FileFsWriteRepository {
/// Crea un nuevo repositorio de escritura de archivos
pub fn new(
root_path: PathBuf,
metadata_manager: Arc<FileMetadataManager>,
path_resolver: Arc<FilePathResolver>,
storage_mediator: Arc<dyn StorageMediator>,
config: AppConfig,
parallel_processor: Option<Arc<ParallelFileProcessor>>,
) -> Self {
Self {
root_path,
metadata_manager,
path_resolver,
storage_mediator,
config,
parallel_processor,
}
}
/// Crea directorios padres si es necesario
async fn ensure_parent_directory(&self, abs_path: &PathBuf) -> FileRepositoryResult<()> {
if let Some(parent) = abs_path.parent() {
tokio::time::timeout(
self.config.timeouts.dir_timeout(),
tokio::fs::create_dir_all(parent)
).await
.map_err(|_| crate::domain::repositories::file_repository::FileRepositoryError::Timeout(
format!("Timeout creating parent directory: {}", parent.display())
))?
.map_err(crate::domain::repositories::file_repository::FileRepositoryError::IoError)?;
}
Ok(())
}
/// Crea una entidad de archivo a partir de metadatos
async fn create_file_entity(
&self,
id: String,
name: String,
storage_path: StoragePath,
size: u64,
mime_type: String,
folder_id: Option<String>,
created_at: Option<u64>,
modified_at: Option<u64>,
) -> FileRepositoryResult<File> {
// If timestamps are provided, use them; otherwise, let File::new create default timestamps
if let (Some(created), Some(modified)) = (created_at, modified_at) {
File::with_timestamps(
id,
name,
storage_path,
size,
mime_type,
folder_id,
created,
modified,
)
.map_err(|e| crate::domain::repositories::file_repository::FileRepositoryError::Other(e.to_string()))
} else {
File::new(
id,
name,
storage_path,
size,
mime_type,
folder_id,
)
.map_err(|e| crate::domain::repositories::file_repository::FileRepositoryError::Other(e.to_string()))
}
}
/// Elimina un archivo de forma no bloqueante
async fn delete_file_non_blocking(&self, _abs_path: PathBuf) -> FileRepositoryResult<()> {
// Implementación real debe eliminar el archivo
// Por ahora, devolvemos OK
Ok(())
}
}
#[async_trait]
impl FileWritePort for FileFsWriteRepository {
async fn save_file(
&self,
name: String,
folder_id: Option<String>,
content_type: String,
content: Vec<u8>,
) -> Result<File, DomainError> {
// Implementación real debe guardar el archivo en disco
// Por ahora, devolvemos un error
Err(DomainError::internal_error("File save", "Save functionality not yet implemented"))
}
async fn move_file(&self, file_id: &str, target_folder_id: Option<String>) -> Result<File, DomainError> {
// Implementación real debe mover el archivo a otra carpeta
// Por ahora, devolvemos un error
Err(DomainError::internal_error("File move", "Move functionality not yet implemented"))
}
async fn delete_file(&self, id: &str) -> Result<(), DomainError> {
// Implementación real debe eliminar el archivo
// Por ahora, devolvemos un error
Err(DomainError::internal_error("File delete", "Delete functionality not yet implemented"))
}
}
@@ -0,0 +1,158 @@
use std::path::PathBuf;
use std::sync::Arc;
use tokio::time;
use tokio::fs;
use std::time::Duration;
use crate::infrastructure::services::file_metadata_cache::{FileMetadataCache, CacheEntryType, FileMetadata};
use crate::common::config::AppConfig;
use crate::common::errors::DomainError;
/// Gestor de metadatos de archivos que encapsula la lógica de caché
pub struct FileMetadataManager {
metadata_cache: Arc<FileMetadataCache>,
config: AppConfig,
}
#[derive(Debug, thiserror::Error)]
pub enum MetadataError {
#[error("Error de E/S al acceder a los metadatos: {0}")]
IoError(#[from] std::io::Error),
#[error("Timeout al acceder a los metadatos: {0}")]
Timeout(String),
#[error("Metadatos no disponibles: {0}")]
Unavailable(String),
}
impl From<MetadataError> for DomainError {
fn from(err: MetadataError) -> Self {
match err {
MetadataError::IoError(e) => DomainError::internal_error("FileMetadata", e.to_string()),
MetadataError::Timeout(msg) => DomainError::internal_error("FileMetadata", msg),
MetadataError::Unavailable(msg) => DomainError::not_found("FileMetadata", msg),
}
}
}
impl FileMetadataManager {
/// Crea un nuevo gestor de metadatos
pub fn new(metadata_cache: Arc<FileMetadataCache>, config: AppConfig) -> Self {
Self {
metadata_cache,
config,
}
}
/// Comprueba si un archivo existe en la ruta especificada con caché
pub async fn file_exists(&self, abs_path: &PathBuf) -> Result<bool, MetadataError> {
// Intentar obtener del caché avanzado primero
if let Some(is_file) = self.metadata_cache.is_file(&abs_path).await {
tracing::debug!("Metadata cache hit for existence check: {} - path: {}", is_file, abs_path.display());
return Ok(is_file);
}
// Si no está en caché, verificar directamente y actualizar caché
tracing::debug!("Metadata cache miss for existence check: {}", abs_path.display());
// Utilizar timeout para evitar bloqueo
match time::timeout(
self.config.timeouts.file_timeout(),
fs::metadata(&abs_path)
).await {
Ok(Ok(metadata)) => {
let is_file = metadata.is_file();
// Actualizar la caché con información fresca
if let Err(e) = self.metadata_cache.refresh_metadata(&abs_path).await {
tracing::warn!("Failed to update cache for {}: {}", abs_path.display(), e);
}
if is_file {
tracing::debug!("File exists and is accessible: {}", abs_path.display());
Ok(true)
} else {
tracing::warn!("Path exists but is not a file: {}", abs_path.display());
Ok(false)
}
},
Ok(Err(e)) => {
tracing::warn!("File check failed: {} - {}", abs_path.display(), e);
// Añadir a caché como no existente
let entry_type = CacheEntryType::Unknown;
let file_metadata = FileMetadata::new(
abs_path.clone(),
false,
entry_type,
None,
None,
None,
None,
Duration::from_millis(self.config.timeouts.file_operation_ms),
);
self.metadata_cache.update_cache(file_metadata).await;
Ok(false)
},
Err(_) => {
tracing::warn!("Timeout checking file metadata: {}", abs_path.display());
Err(MetadataError::Timeout(format!("Timeout checking file: {}", abs_path.display())))
}
}
}
/// Obtiene metadatos de archivo (tamaño, fechas creación/modificación) con caché
pub async fn get_file_metadata(&self, abs_path: &PathBuf) -> Result<(u64, u64, u64), MetadataError> {
// Intentar obtener de caché primero
if let Some(cached_metadata) = self.metadata_cache.get_metadata(abs_path).await {
if let (Some(size), Some(created_at), Some(modified_at)) =
(cached_metadata.size, cached_metadata.created_at, cached_metadata.modified_at) {
tracing::debug!("Using cached metadata for: {}", abs_path.display());
return Ok((size, created_at, modified_at));
}
}
// Si no está en caché o metadatos incompletos, cargar desde sistema de archivos
let metadata = match time::timeout(
self.config.timeouts.file_timeout(),
fs::metadata(&abs_path)
).await {
Ok(Ok(metadata)) => metadata,
Ok(Err(e)) => return Err(MetadataError::IoError(e)),
Err(_) => return Err(MetadataError::Timeout(
format!("Timeout getting metadata for: {}", abs_path.display())
)),
};
let size = metadata.len();
// Get creation timestamp
let created_at = metadata.created()
.map(|time| time.duration_since(std::time::UNIX_EPOCH).unwrap_or_default().as_secs())
.unwrap_or_else(|_| 0);
// Get modification timestamp
let modified_at = metadata.modified()
.map(|time| time.duration_since(std::time::UNIX_EPOCH).unwrap_or_default().as_secs())
.unwrap_or_else(|_| 0);
// Actualizar caché si es posible
if let Err(e) = self.metadata_cache.refresh_metadata(abs_path).await {
tracing::warn!("Failed to update metadata cache for {}: {}", abs_path.display(), e);
}
Ok((size, created_at, modified_at))
}
/// Invalida la entrada de caché para un archivo
pub async fn invalidate(&self, abs_path: &PathBuf) {
self.metadata_cache.invalidate(abs_path).await;
}
/// Invalida la entrada de caché para un directorio y su contenido
pub async fn invalidate_directory(&self, dir_path: &PathBuf) {
self.metadata_cache.invalidate_directory(dir_path).await;
}
}
@@ -0,0 +1,90 @@
use std::path::PathBuf;
use std::sync::Arc;
use async_trait::async_trait;
use crate::domain::services::path_service::{PathService, StoragePath};
use crate::application::services::storage_mediator::StorageMediator;
use crate::infrastructure::services::id_mapping_service::{IdMappingService, IdMappingError};
use crate::domain::repositories::file_repository::FileRepositoryError;
use crate::common::errors::DomainError;
use crate::application::ports::storage_ports::FilePathResolutionPort;
/// Resuelve rutas de archivos y gestiona el mapeo de IDs a rutas
pub struct FilePathResolver {
path_service: Arc<PathService>,
storage_mediator: Arc<dyn StorageMediator>,
id_mapping_service: Arc<IdMappingService>,
}
impl FilePathResolver {
/// Crea un nuevo resolver de rutas
pub fn new(
path_service: Arc<PathService>,
storage_mediator: Arc<dyn StorageMediator>,
id_mapping_service: Arc<IdMappingService>,
) -> Self {
Self {
path_service,
storage_mediator,
id_mapping_service,
}
}
/// Resuelve una ruta de dominio a una ruta física absoluta
pub fn resolve_storage_path(&self, storage_path: &StoragePath) -> PathBuf {
self.path_service.resolve_path(storage_path)
}
/// Resuelve una ruta PathBuf a una ruta física absoluta (legacy)
pub fn resolve_legacy_path(&self, relative_path: &std::path::Path) -> PathBuf {
self.storage_mediator.resolve_path(relative_path)
}
/// Obtiene la ruta de un archivo por su ID
pub async fn get_path_by_id(&self, id: &str) -> Result<StoragePath, FileRepositoryError> {
self.id_mapping_service.get_path_by_id(id).await
.map_err(FileRepositoryError::from)
}
/// Actualiza la ruta para un ID existente
pub async fn update_path(&self, id: &str, storage_path: &StoragePath) -> Result<(), FileRepositoryError> {
self.id_mapping_service.update_path(id, storage_path).await
.map_err(FileRepositoryError::from)
}
/// Obtiene o crea un ID para una ruta
pub async fn get_or_create_id(&self, storage_path: &StoragePath) -> Result<String, FileRepositoryError> {
self.id_mapping_service.get_or_create_id(storage_path).await
.map_err(FileRepositoryError::from)
}
/// Elimina un ID del mapeo
pub async fn remove_id(&self, id: &str) -> Result<(), FileRepositoryError> {
self.id_mapping_service.remove_id(id).await
.map_err(FileRepositoryError::from)
}
/// Guarda cambios pendientes
pub async fn save_changes(&self) -> Result<(), FileRepositoryError> {
self.id_mapping_service.save_pending_changes().await
.map_err(FileRepositoryError::from)
}
}
// Implementación de FilePathResolutionPort
#[async_trait]
impl FilePathResolutionPort for FilePathResolver {
async fn get_file_path(&self, id: &str) -> Result<StoragePath, DomainError> {
self.get_path_by_id(id).await
.map_err(|e| match e {
FileRepositoryError::NotFound(id) => DomainError::not_found("File", id),
FileRepositoryError::IoError(e) => DomainError::internal_error("FilePath", e.to_string()),
FileRepositoryError::Timeout(msg) => DomainError::internal_error("FilePath", msg),
_ => DomainError::internal_error("FilePath", e.to_string()),
})
}
fn resolve_path(&self, storage_path: &StoragePath) -> PathBuf {
self.resolve_storage_path(storage_path)
}
}
+11
View File
@@ -2,3 +2,14 @@ pub mod file_fs_repository;
pub mod folder_fs_repository;
pub mod parallel_file_processor;
// Nuevos repositorios refactorizados
pub mod file_metadata_manager;
pub mod file_path_resolver;
pub mod file_fs_read_repository;
pub mod file_fs_write_repository;
// Re-exportar para facilitar acceso
pub use file_metadata_manager::FileMetadataManager;
pub use file_path_resolver::FilePathResolver;
pub use file_fs_read_repository::FileFsReadRepository;
pub use file_fs_write_repository::FileFsWriteRepository;