feat(#113): 100% blob storage model — PostgreSQL metadata + DedupService blobs

BREAKING CHANGE: Storage model completely rewritten. All file/folder
metadata now lives in PostgreSQL (storage schema). File content stored
as content-addressable blobs via DedupService. Filesystem directories
are no longer used for user storage.

New components:
- storage.folders / storage.files / storage.trash_items (PG schema)
- FolderDbRepository: virtual folders backed by PG
- FileBlobReadRepository: file reads via PG metadata + dedup blobs
- FileBlobWriteRepository: file writes via PG metadata + dedup blobs
- TrashDbRepository: soft-delete trash using is_trashed flags

Removed legacy FS components (~5500 lines deleted):
- FolderFsRepository, FileFsReadRepository, FileFsWriteRepository
- CompositeFileRepository, ParallelFileProcessor
- IdMappingService, IdMappingOptimizer, FileMetadataCache
- BufferPool, FileSystemUtils, RepositoryErrors
- TrashFsRepository, FolderFsRepositoryTrash

DI rewired: build_app_state() now requires PgPool (no FS fallback).
FileUploadService.new_with_read() and FileRetrievalService.new_with_cache()
constructors added for blob model (no write-behind needed).

Closes #113
This commit is contained in:
Dionisio
2026-02-14 17:54:25 +01:00
parent f25987e553
commit 3c7c16f07e
27 changed files with 1883 additions and 7429 deletions
-504
View File
@@ -1,504 +0,0 @@
use std::cmp::min;
use std::collections::VecDeque;
use std::sync::Arc;
use std::time::{Duration, Instant};
use tokio::sync::{Mutex, Semaphore};
use tracing::debug;
/// Default buffer size in the pool
pub const DEFAULT_BUFFER_SIZE: usize = 64 * 1024; // 64KB
/// Default maximum number of buffers in the pool
pub const DEFAULT_MAX_BUFFERS: usize = 100;
/// Default time-to-live for an inactive buffer (in seconds)
pub const DEFAULT_BUFFER_TTL: u64 = 60;
/// Buffer pooling to optimize read/write operations
pub struct BufferPool {
/// Pool of available buffers
pool: Mutex<VecDeque<PooledBuffer>>,
/// Semaphore to limit the maximum number of buffers
limit: Semaphore,
/// Size of buffers in the pool
buffer_size: usize,
/// Pool statistics
stats: Mutex<BufferPoolStats>,
/// Time-to-live for an inactive buffer
buffer_ttl: Duration,
}
/// Structure for tracking pool statistics
#[derive(Debug, Clone, Default)]
pub struct BufferPoolStats {
/// Total number of get operations
pub gets: usize,
/// Number of pool hits (successful reuse)
pub hits: usize,
/// Number of misses (new buffer creation)
pub misses: usize,
/// Number of returns to the pool
pub returns: usize,
/// Number of TTL evictions
pub evictions: usize,
/// Maximum number of buffers reached
pub max_buffers_reached: usize,
/// Semaphore waits
pub waits: usize,
}
/// Pool buffer with management metadata
struct PooledBuffer {
/// Actual byte buffer
buffer: Vec<u8>,
/// Timestamp of when it was added/returned to the pool
last_used: Instant,
}
/// Borrowed buffer from the pool with automatic cleanup
#[derive(Clone)]
pub struct BorrowedBuffer {
/// Current buffer
buffer: Vec<u8>,
/// Actual used size of the buffer
used_size: usize,
/// Reference to the pool for returning
pool: Arc<BufferPool>,
/// Whether the buffer should be returned to the pool or not
return_to_pool: bool,
}
impl BufferPool {
/// 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)),
limit: Semaphore::new(max_buffers),
buffer_size,
stats: Mutex::new(BufferPoolStats::default()),
buffer_ttl: Duration::from_secs(buffer_ttl_secs),
})
}
/// Creates a pool with default configuration
pub fn default() -> Arc<Self> {
Self::new(DEFAULT_BUFFER_SIZE, DEFAULT_MAX_BUFFERS, DEFAULT_BUFFER_TTL)
}
/// 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 {
// Increment get counter
{
let mut stats = self.stats.lock().await;
stats.gets += 1;
}
// 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 permits available, waiting
{
let mut stats = self.stats.lock().await;
stats.waits += 1;
stats.max_buffers_reached += 1;
}
debug!("Buffer pool: waiting for available buffer");
let permit = self
.limit
.acquire()
.await
.expect("Semaphore should not be closed");
debug!("Buffer pool: acquired buffer after waiting");
permit.forget();
}
};
// 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() {
// Check if the buffer has expired
if pooled_buffer.last_used.elapsed() > self.buffer_ttl {
// Expired buffer, discard and create a new one
let mut stats = self.stats.lock().await;
stats.evictions += 1;
stats.misses += 1;
drop(stats);
debug!("Buffer pool: evicted expired buffer");
// Create new buffer (reusing the permit)
drop(pool_locked); // Release the lock before returning
BorrowedBuffer {
buffer: vec![0; self.buffer_size],
used_size: 0,
pool: pool_arc,
return_to_pool: true,
}
} else {
// Valid buffer, reuse it
let mut stats = self.stats.lock().await;
stats.hits += 1;
drop(stats);
// Release the lock before returning
drop(pool_locked);
// Clear buffer for security
pooled_buffer.buffer.fill(0);
BorrowedBuffer {
buffer: pooled_buffer.buffer,
used_size: 0,
pool: pool_arc,
return_to_pool: true,
}
}
} else {
// No buffers available, create a new one
let mut stats = self.stats.lock().await;
stats.misses += 1;
drop(stats);
// Release the lock before returning
drop(pool_locked);
debug!("Buffer pool: creating new buffer");
BorrowedBuffer {
buffer: vec![0; self.buffer_size],
used_size: 0,
pool: pool_arc,
return_to_pool: true,
}
}
}
/// Returns a buffer to the pool
async fn return_buffer(&self, mut buffer: Vec<u8>) {
// 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
);
// Release the semaphore permit even if we discard the buffer
self.limit.add_permits(1);
return;
}
// Resize to ensure correct capacity
buffer.resize(self.buffer_size, 0);
// Add to the pool
let mut pool_locked = self.pool.lock().await;
pool_locked.push_back(PooledBuffer {
buffer,
last_used: Instant::now(),
});
// Update statistics
let mut stats = self.stats.lock().await;
stats.returns += 1;
// Release the semaphore permit so another caller can acquire a buffer
drop(pool_locked);
drop(stats);
self.limit.add_permits(1);
}
/// 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;
// Count expired
let count_before = pool_locked.len();
// Filter keeping only non-expired
pool_locked.retain(|buffer| buffer.last_used.elapsed() <= self.buffer_ttl);
// Count how many were removed
let removed = count_before - pool_locked.len();
if removed > 0 {
// Update statistics
let mut stats = self.stats.lock().await;
stats.evictions += removed;
debug!("Buffer pool: cleaned {} expired buffers", removed);
}
}
/// Gets current pool statistics
pub async fn get_stats(&self) -> BufferPoolStats {
self.stats.lock().await.clone()
}
/// Starts the periodic cleanup task
pub fn start_cleaner(pool: Arc<Self>) {
tokio::spawn(async move {
let interval = Duration::from_secs(30); // Clean every 30 seconds
loop {
tokio::time::sleep(interval).await;
pool.clean_expired_buffers().await;
// Log statistics periodically
let stats = pool.get_stats().await;
debug!(
"Buffer pool stats: gets={}, hits={}, misses={}, hit_ratio={:.2}%, returns={}, \
evictions={}, max_reached={}, waits={}",
stats.gets,
stats.hits,
stats.misses,
if stats.gets > 0 {
(stats.hits as f64 * 100.0) / stats.gets as f64
} else {
0.0
},
stats.returns,
stats.evictions,
stats.max_buffers_reached,
stats.waits
);
}
});
}
}
impl Clone for BufferPool {
fn clone(&self) -> Self {
Self {
pool: Mutex::new(VecDeque::new()),
limit: Semaphore::new(self.limit.available_permits()),
buffer_size: self.buffer_size,
stats: Mutex::new(BufferPoolStats::default()),
buffer_ttl: self.buffer_ttl,
}
}
}
impl BorrowedBuffer {
/// Accesses the internal buffer
pub fn as_mut_slice(&mut self) -> &mut [u8] {
&mut self.buffer
}
/// Gets a reference to the used data
pub fn as_slice(&self) -> &[u8] {
&self.buffer[..self.used_size]
}
/// Sets how many bytes were actually used
pub fn set_used(&mut self, size: usize) {
self.used_size = min(size, self.buffer.len());
}
/// Converts into a Vec<u8> that includes only the used data
pub fn into_vec(mut self) -> Vec<u8> {
// Mark to not return to pool
self.return_to_pool = false;
// Create a new vector with only the used data
self.buffer[..self.used_size].to_vec()
}
/// 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]);
self.used_size = copy_size;
copy_size
}
/// 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
}
/// Gets the total buffer size
pub fn capacity(&self) -> usize {
self.buffer.len()
}
/// Gets the used buffer size
pub fn used_size(&self) -> usize {
self.used_size
}
}
// When a BorrowedBuffer is dropped, it is returned to the pool
impl Drop for BorrowedBuffer {
fn drop(&mut self) {
if self.return_to_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 the return so that drop doesn't block
// return_buffer will release the semaphore permit
tokio::spawn(async move {
pool.return_buffer(buffer).await;
});
} else {
// Buffer not returned to pool, but we still need to release the semaphore permit
self.pool.limit.add_permits(1);
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn test_buffer_pooling() {
// Create small pool for testing
let pool = BufferPool::new(1024, 5, 60);
// 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");
// Get another buffer
let buffer2 = pool.get_buffer().await;
// Verify stats
let stats = pool.get_stats().await;
assert_eq!(stats.gets, 2);
assert_eq!(stats.hits, 0); // no hits yet
assert_eq!(stats.misses, 2); // all are misses
// Return buffer1 to pool (implicitly via drop)
drop(buffer1);
// Allow the async return to occur
tokio::time::sleep(Duration::from_millis(10)).await;
// Get another buffer (should reuse the returned one)
let buffer3 = pool.get_buffer().await;
// Verify updated stats
let stats = pool.get_stats().await;
assert_eq!(stats.gets, 3);
assert_eq!(stats.hits, 1); // now there should be a hit
assert_eq!(stats.returns, 1); // one buffer returned
// Cleanup
drop(buffer2);
drop(buffer3);
}
#[tokio::test]
async fn test_buffer_operations() {
let pool = BufferPool::new(1024, 10, 60);
// Get buffer
let mut buffer = pool.get_buffer().await;
// Write data
buffer.copy_from_slice(b"Hello, world!");
assert_eq!(buffer.used_size(), 13);
assert_eq!(buffer.as_slice(), b"Hello, world!");
// Convert to vec and verify
let vec = buffer.into_vec(); // This prevents returning to pool
assert_eq!(vec, b"Hello, world!");
// 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);
}
#[tokio::test]
async fn test_pool_limit() {
// Pool with only 3 buffers
let pool = BufferPool::new(1024, 3, 60);
// Get 3 buffers (reaches the limit)
let buffer1 = pool.get_buffer().await;
let buffer2 = pool.get_buffer().await;
let buffer3 = pool.get_buffer().await;
// Verify stats
let stats = pool.get_stats().await;
assert_eq!(stats.gets, 3);
assert_eq!(stats.waits, 0); // no waits yet
// 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
});
// Give time for the task to try to take the buffer
tokio::time::sleep(Duration::from_millis(50)).await;
// Verify there is a wait
let stats = pool.get_stats().await;
assert_eq!(stats.waits, 1);
// Release a buffer
drop(buffer1);
// Give time for the async return and for the waiting task to get its buffer
tokio::time::sleep(Duration::from_millis(50)).await;
// Verify the task was able to continue
assert!(handle.await.unwrap());
// Cleanup
drop(buffer2);
drop(buffer3);
}
#[tokio::test]
async fn test_ttl_expiration() {
// Pool with very short TTL for testing
let pool = BufferPool::new(1024, 5, 1); // 1 second TTL
// Get and return a buffer
let buffer = pool.get_buffer().await;
drop(buffer);
// Allow the async return to occur
tokio::time::sleep(Duration::from_millis(50)).await;
// Verify there is a buffer in the pool
let stats = pool.get_stats().await;
assert_eq!(stats.returns, 1);
// Wait for the TTL to expire
tokio::time::sleep(Duration::from_secs(2)).await;
// Clean expired
pool.clean_expired_buffers().await;
// Get another buffer (should be a miss since the previous one expired)
let _buffer2 = pool.get_buffer().await;
// Verify stats
let stats = pool.get_stats().await;
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)
}
}
@@ -6,14 +6,12 @@ use flate2::read::GzEncoder as GzEncoderRead;
use futures::{Stream, StreamExt};
use std::io;
use std::io::Read;
use std::sync::Arc;
use tracing::error;
use crate::application::ports::compression_ports::{
CompressionLevel as PortCompressionLevel, CompressionPort,
};
use crate::domain::errors::DomainError;
use crate::infrastructure::services::buffer_pool::BufferPool;
/// Compression level for files
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
@@ -73,22 +71,12 @@ pub trait CompressionService: Send + Sync {
}
/// Gzip compression service implementation
pub struct GzipCompressionService {
/// Buffer pool for memory optimization
buffer_pool: Option<Arc<BufferPool>>,
}
pub struct GzipCompressionService;
impl GzipCompressionService {
/// Creates a new service instance
pub fn new() -> Self {
Self { buffer_pool: None }
}
/// 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),
}
Self
}
}
@@ -96,63 +84,6 @@ impl GzipCompressionService {
impl CompressionService for GzipCompressionService {
/// Compresses data in memory using Gzip
async fn compress_data(&self, data: &[u8], level: CompressionLevel) -> io::Result<Vec<u8>> {
// If we have a buffer pool, use a borrowed buffer for compression
if let Some(pool) = &self.buffer_pool {
// Estimate the compression size (approximately 80% of original for typical cases)
let estimated_size = (data.len() as f64 * 0.8) as usize;
// Get a buffer from the pool
let buffer = pool.get_buffer().await;
// Check if the buffer is large enough
if buffer.capacity() >= estimated_size {
// 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();
// 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());
// 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,
};
// 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;
// Verify result
match result {
Ok(Ok(())) => {
// Get the buffer and convert it to Vec<u8>
let buffer = buffer_ptr.lock().await;
let cloned_buffer = buffer.clone();
drop(buffer); // Release the mutex first
return Ok(cloned_buffer.into_vec());
}
Ok(Err(e)) => {
error!("Compression error with buffer pool: {}", e);
// Fall back to standard implementation
}
Err(e) => {
error!("Compression task error with buffer pool: {}", e);
// Fall back to standard implementation
}
}
}
}
// 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 || {
@@ -170,61 +101,7 @@ impl CompressionService for GzipCompressionService {
/// Decompresses data in memory
async fn decompress_data(&self, compressed_data: &[u8]) -> io::Result<Vec<u8>> {
// If we have a buffer pool, use a borrowed buffer for decompression
if let Some(pool) = &self.buffer_pool {
// Estimate the decompression size (approximately 5x of compressed for typical cases)
let estimated_size = compressed_data.len() * 5;
// Get a buffer from the pool
let buffer = pool.get_buffer().await;
// Check if the buffer is large enough
if buffer.capacity() >= estimated_size {
// 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();
// Decompress data
let result = tokio::task::spawn_blocking(move || {
let mut decoder = GzDecoder::new(&data[..]);
// Try to lock the mutex
let mut buffer_guard = match futures::executor::block_on(buffer_clone.lock()) {
buffer => 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;
// Verify result
match result {
Ok(Ok(())) => {
// Get the buffer and convert it to Vec<u8>
let buffer = buffer_ptr.lock().await;
let cloned_buffer = buffer.clone();
drop(buffer); // Release the mutex first
return Ok(cloned_buffer.into_vec());
}
Ok(Err(e)) => {
error!("Decompression error with buffer pool: {}", e);
// Fall back to standard implementation
}
Err(e) => {
error!("Decompression task error with buffer pool: {}", e);
// Fall back to standard implementation
}
}
}
}
// 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
let data = compressed_data.to_vec();
tokio::task::spawn_blocking(move || {
let mut decoder = GzDecoder::new(&data[..]);
let mut decompressed = Vec::new();
@@ -1,749 +0,0 @@
use futures::future::BoxFuture;
use mime_guess::from_path;
use std::collections::{HashMap, VecDeque};
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::{Duration, Instant, UNIX_EPOCH};
use tokio::fs;
use tokio::sync::RwLock;
use tokio::time;
use tracing::debug;
use crate::domain::entities::file::File;
use crate::common::config::AppConfig;
/// Cache entry types
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CacheEntryType {
/// File
File,
/// Directory
Directory,
/// Unknown type
Unknown,
}
/// Cache statistics for monitoring
#[derive(Debug, Clone, Default)]
pub struct CacheStats {
/// Number of cache hits
pub hits: usize,
/// Number of cache misses
pub misses: usize,
/// Number of manual invalidations
pub invalidations: usize,
/// Number of automatic expirations
pub expirations: usize,
/// Number of cache inserts
pub inserts: usize,
/// Total time saved (milliseconds)
pub time_saved_ms: u64,
}
/// Complete cached file metadata
#[derive(Debug, Clone)]
pub struct FileMetadata {
/// Absolute file path
pub path: PathBuf,
/// Whether the file physically exists
pub exists: bool,
/// Entry type (file, directory)
pub entry_type: CacheEntryType,
/// Size in bytes (for files)
pub size: Option<u64>,
/// MIME type (for files)
pub mime_type: Option<String>,
/// Creation timestamp (UNIX epoch seconds)
pub created_at: Option<u64>,
/// Modification timestamp (UNIX epoch seconds)
pub modified_at: Option<u64>,
/// Previous access (used for LRU)
pub last_access: Instant,
/// Cache expiration time
pub expires_at: Instant,
/// Number of accesses to this entry
pub access_count: usize,
}
impl FileMetadata {
/// Creates a new metadata entry
pub fn new(
path: PathBuf,
exists: bool,
entry_type: CacheEntryType,
size: Option<u64>,
mime_type: Option<String>,
created_at: Option<u64>,
modified_at: Option<u64>,
ttl: Duration,
) -> Self {
let now = Instant::now();
Self {
path,
exists,
entry_type,
size,
mime_type,
created_at,
modified_at,
last_access: now,
expires_at: now + ttl,
access_count: 1,
}
}
/// Updates the last access time
pub fn touch(&mut self) {
self.last_access = Instant::now();
self.access_count += 1;
}
/// Checks if the entry has expired
pub fn is_expired(&self) -> bool {
Instant::now() > self.expires_at
}
/// Updates the expiration time with a new TTL
pub fn update_expiry(&mut self, ttl: Duration) {
self.expires_at = Instant::now() + ttl;
}
}
/// Advanced file metadata cache
pub struct FileMetadataCache {
/// Main metadata cache
metadata_cache: RwLock<HashMap<PathBuf, FileMetadata>>,
/// LRU queue for cache management
lru_queue: RwLock<VecDeque<PathBuf>>,
/// Cache usage statistics
stats: RwLock<CacheStats>,
/// Global application configuration
config: AppConfig,
/// Adaptive TTL for popular entries
ttl_multiplier: f64,
/// Popularity threshold for extended TTL
popularity_threshold: usize,
/// Maximum cache size
max_entries: usize,
}
impl FileMetadataCache {
/// 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, // Popular entries have 5x TTL
popularity_threshold: 10, // After 10 accesses it's considered popular
max_entries,
}
}
/// 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());
let mime_type = Some(file.mime_type().to_string());
let created_at = Some(file.created_at());
let modified_at = Some(file.modified_at());
// Use a standard TTL
let ttl = Duration::from_secs(60); // 1 minute
FileMetadata::new(
abs_path,
true,
entry_type,
size,
mime_type,
created_at,
modified_at,
ttl,
)
}
/// Creates a default instance
pub fn default() -> Self {
Self::new(AppConfig::default(), 10_000)
}
/// Creates a cache instance with default configuration
pub fn default_with_config(config: AppConfig) -> Self {
Self::new(config, 50_000) // Larger cache for production system
}
/// 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) {
// Check if expired
if metadata.is_expired() {
// Remove from cache if expired
cache.remove(path);
// Update statistics
let mut stats = self.stats.write().await;
stats.misses += 1;
stats.expirations += 1;
debug!("Cache entry expired for: {}", path.display());
return None;
}
// Update access time
metadata.touch();
// For popular entries, extend TTL
if metadata.access_count >= self.popularity_threshold {
let new_ttl = match metadata.entry_type {
CacheEntryType::File => Duration::from_millis(
(self.config.timeouts.file_operation_ms as f64 * self.ttl_multiplier)
as u64,
),
CacheEntryType::Directory => Duration::from_millis(
(self.config.timeouts.dir_operation_ms as f64 * self.ttl_multiplier) as u64,
),
_ => Duration::from_secs(60), // 1 minute by default
};
metadata.update_expiry(new_ttl);
debug!("Extended TTL for popular entry: {}", path.display());
}
// Calculate approximate time saved
let elapsed = start_time.elapsed().as_millis() as u64;
let estimated_io_time: u64 = 10; // We assume 10ms minimum for IO operation
let time_saved = estimated_io_time.saturating_sub(elapsed);
// Update statistics
let mut stats = self.stats.write().await;
stats.hits += 1;
stats.time_saved_ms += time_saved;
debug!("Cache hit for: {}", path.display());
// Also keep the LRU queue updated
self.update_lru(path.to_path_buf()).await;
// Clone to return
return Some(metadata.clone());
}
// Not found in cache
let mut stats = self.stats.write().await;
stats.misses += 1;
debug!("Cache miss for: {}", path.display());
None
}
/// Updates the LRU queue
async fn update_lru(&self, path: PathBuf) {
let mut lru = self.lru_queue.write().await;
// Remove if already exists
if let Some(pos) = lru.iter().position(|p| p == &path) {
lru.remove(pos);
}
// Add to the end (most recent)
lru.push_back(path);
}
/// 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);
}
None
}
/// 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);
}
None
}
/// 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);
}
None
}
/// 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;
}
None
}
/// 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;
}
None
}
/// Refreshes metadata for a path
pub async fn refresh_metadata(&self, path: &Path) -> Result<FileMetadata, std::io::Error> {
// Perform actual filesystem read
let metadata = fs::metadata(path).await?;
// Determine entry type
let entry_type = if metadata.is_dir() {
CacheEntryType::Directory
} else if metadata.is_file() {
CacheEntryType::File
} else {
CacheEntryType::Unknown
};
// Get size for files
let size = if metadata.is_file() {
Some(metadata.len())
} else {
None
};
// Get MIME type for files
let mime_type = if metadata.is_file() {
Some(from_path(path).first_or_octet_stream().to_string())
} else {
None
};
// Get timestamps
let created_at = metadata
.created()
.map(|time| {
time.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs()
})
.ok();
let modified_at = metadata
.modified()
.map(|time| {
time.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs()
})
.ok();
// 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)
};
// Create metadata entry
let file_metadata = FileMetadata::new(
path.to_path_buf(),
true,
entry_type,
size,
mime_type,
created_at,
modified_at,
ttl,
);
// Update cache
self.update_cache(file_metadata.clone()).await;
Ok(file_metadata)
}
/// Updates the cache with new metadata
pub async fn update_cache(&self, metadata: FileMetadata) {
// Avoid full cache before inserting
self.ensure_capacity().await;
let path = metadata.path.clone();
// Insert into cache
{
let mut cache = self.metadata_cache.write().await;
cache.insert(path.clone(), metadata);
// Update statistics
let mut stats = self.stats.write().await;
stats.inserts += 1;
}
// Update the LRU queue
self.update_lru(path).await;
}
/// Ensures there is space in the cache
async fn ensure_capacity(&self) {
let cache_size = {
let cache = self.metadata_cache.read().await;
cache.len()
};
if cache_size >= self.max_entries {
self.evict_lru_entries(cache_size / 10).await; // Free up 10%
}
}
/// Removes least recently used entries
async fn evict_lru_entries(&self, count: usize) {
let mut paths_to_remove = Vec::with_capacity(count);
// Get entries to remove from the LRU queue
{
let mut lru = self.lru_queue.write().await;
for _ in 0..count {
if let Some(path) = lru.pop_front() {
paths_to_remove.push(path);
} else {
break;
}
}
}
// Remove from the main cache
{
let mut cache = self.metadata_cache.write().await;
for path in paths_to_remove {
cache.remove(&path);
}
}
debug!("Evicted {} LRU entries from cache", count);
}
/// Invalidate a specific cache entry
pub async fn invalidate(&self, path: &Path) {
// Remove from the main cache
{
let mut cache = self.metadata_cache.write().await;
cache.remove(path);
// Update statistics
let mut stats = self.stats.write().await;
stats.invalidations += 1;
}
// Remove from the LRU queue
let path_buf = path.to_path_buf();
{
let mut lru = self.lru_queue.write().await;
if let Some(pos) = lru.iter().position(|p| p == &path_buf) {
lru.remove(pos);
}
}
debug!("Invalidated cache entry for: {}", path.display());
}
/// 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();
// Find all paths that start with the directory
{
let cache = self.metadata_cache.read().await;
for path in cache.keys() {
let path_str = path.to_string_lossy().to_string();
if path_str.starts_with(&dir_str) {
paths_to_remove.push(path.clone());
}
}
}
// Update statistics
{
let mut stats = self.stats.write().await;
stats.invalidations += paths_to_remove.len();
}
// Remove each found path
for path in paths_to_remove {
self.invalidate(&path).await;
}
debug!("Invalidated directory and contents: {}", dir_path.display());
}
/// Get current cache statistics
pub async fn get_stats(&self) -> CacheStats {
let stats = self.stats.read().await;
stats.clone()
}
/// Clears all expired entries from the cache
pub async fn clear_expired(&self) {
let now = Instant::now();
let mut paths_to_remove = Vec::new();
// Find expired entries
{
let cache = self.metadata_cache.read().await;
for (path, metadata) in cache.iter() {
if now > metadata.expires_at {
paths_to_remove.push(path.clone());
}
}
}
// Update statistics
{
let mut stats = self.stats.write().await;
stats.expirations += paths_to_remove.len();
}
// Save the number of entries for logging
let num_paths = paths_to_remove.len();
// Remove expired entries
for path in paths_to_remove {
self.invalidate(&path).await;
}
debug!("Cleared {} expired entries from cache", num_paths);
}
/// 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); // Every minute
loop {
// Wait for the interval
time::sleep(cleanup_interval).await;
// Clean expired entries
cache.clear_expired().await;
// Log statistics
let stats = cache.get_stats().await;
let cache_size = {
let cache_map = cache.metadata_cache.read().await;
cache_map.len()
};
debug!(
"Cache stats: size={}, hits={}, misses={}, hit_ratio={:.2}%, time_saved={}ms",
cache_size,
stats.hits,
stats.misses,
if stats.hits + stats.misses > 0 {
(stats.hits as f64 * 100.0) / (stats.hits + stats.misses) as f64
} else {
0.0
},
stats.time_saved_ms
);
}
})
}
/// 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
}
/// Internal preload implementation with depth tracking
async fn _preload_directory_internal(
&self,
dir_path: &Path,
recursive: bool,
max_depth: usize,
current_depth: usize,
) -> Result<usize, std::io::Error> {
Box::pin(async move {
if current_depth > max_depth {
return Ok(0);
}
// Get directory entries
let mut entries = fs::read_dir(dir_path).await?;
let mut count = 0;
// Process each entry
while let Some(entry) = entries.next_entry().await? {
let path = entry.path();
let metadata = fs::metadata(&path).await?;
// Refresh metadata for this entry
self.refresh_metadata(&path).await?;
count += 1;
// Recursively process subdirectories if needed
if recursive && metadata.is_dir() {
// Box to break recursion
count += self
._preload_directory_internal(&path, recursive, max_depth, current_depth + 1)
.await?;
}
}
Ok(count)
})
.await
}
}
// ─── MetadataCachePort implementation ────────────────────────
use crate::application::ports::cache_ports::{CachedMetadataDto, MetadataCachePort};
use crate::common::errors::DomainError;
use async_trait::async_trait;
#[async_trait]
impl MetadataCachePort for FileMetadataCache {
async fn get_metadata(&self, path: &Path) -> Option<CachedMetadataDto> {
// Delegate to the existing rich get_metadata, then project into the DTO.
let fm = FileMetadataCache::get_metadata(self, path).await?;
Some(CachedMetadataDto {
path: fm.path,
exists: fm.exists,
is_file: fm.entry_type == CacheEntryType::File,
size: fm.size,
mime_type: fm.mime_type,
created_at: fm.created_at,
modified_at: fm.modified_at,
})
}
async fn is_file(&self, path: &Path) -> Option<bool> {
FileMetadataCache::is_file(self, path).await
}
async fn refresh_metadata(&self, path: &Path) -> Result<CachedMetadataDto, DomainError> {
let fm = FileMetadataCache::refresh_metadata(self, path)
.await
.map_err(|e| DomainError::internal_error("MetadataCache", e.to_string()))?;
Ok(CachedMetadataDto {
path: fm.path,
exists: fm.exists,
is_file: fm.entry_type == CacheEntryType::File,
size: fm.size,
mime_type: fm.mime_type,
created_at: fm.created_at,
modified_at: fm.modified_at,
})
}
async fn invalidate(&self, path: &Path) {
FileMetadataCache::invalidate(self, path).await
}
async fn invalidate_directory(&self, dir_path: &Path) {
FileMetadataCache::invalidate_directory(self, dir_path).await
}
}
#[cfg(test)]
mod tests {
use super::*;
use tempfile::tempdir;
use tokio::fs::File;
use tokio::io::AsyncWriteExt;
#[tokio::test]
async fn test_cache_operations() {
// Create temporary directory for tests
let temp_dir = tempdir().unwrap();
let file_path = temp_dir.path().join("test_file.txt");
// 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);
// Create cache
let config = AppConfig::default();
let cache = FileMetadataCache::new(config, 1000);
// Verify initial miss
assert!(cache.exists(&file_path).await.is_none());
// 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
// 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));
// Invalidate and verify it no longer exists in cache
cache.invalidate(&file_path).await;
assert!(cache.exists(&file_path).await.is_none());
// Verify statistics
let stats = cache.get_stats().await;
assert_eq!(stats.inserts, 1);
assert_eq!(stats.invalidations, 1);
assert!(stats.hits > 0);
}
#[tokio::test]
async fn test_directory_operations() {
// 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();
let sub_dir = base_path.join("subdir");
fs::create_dir(&sub_dir).await.unwrap();
let file1 = base_path.join("file1.txt");
let file2 = sub_dir.join("file2.txt");
File::create(&file1).await.unwrap();
File::create(&file2).await.unwrap();
// Create cache
let config = AppConfig::default();
let cache = FileMetadataCache::new(config, 1000);
// 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
// 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));
// Invalidate directory and contents
cache.invalidate_directory(&base_path).await;
// 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());
}
}
@@ -1,326 +0,0 @@
use std::io::Error as IoError;
use std::path::Path;
use tempfile::NamedTempFile;
use tokio::fs::{self, File, OpenOptions};
use tokio::io::AsyncWriteExt;
use tracing::{error, warn};
/// Utility functions for file system operations with proper synchronization
pub struct FileSystemUtils;
impl FileSystemUtils {
/// Writes data to a file with fsync to ensure durability
/// Uses a safe atomic write pattern: write to temp file, fsync, rename
pub async fn atomic_write<P: AsRef<Path>>(path: P, contents: &[u8]) -> Result<(), IoError> {
let path = path.as_ref();
// Ensure parent directory exists
if let Some(parent) = path.parent() {
fs::create_dir_all(parent).await?;
}
// Create a temporary file in the same directory
let dir = path.parent().unwrap_or_else(|| Path::new("."));
let temp_file = match NamedTempFile::new_in(dir) {
Ok(file) => file,
Err(e) => {
error!(
"Failed to create temporary file in {}: {}",
dir.display(),
e
);
return Err(IoError::other(format!(
"Failed to create temporary file: {}",
e
)));
}
};
let temp_path = temp_file.path().to_path_buf();
// Convert to tokio file and write contents
let std_file = temp_file.as_file().try_clone()?;
let mut file = File::from_std(std_file);
file.write_all(contents).await?;
// Ensure data is synced to disk
file.flush().await?;
file.sync_all().await?;
// Rename the temporary file to the target path (atomic operation on most filesystems)
fs::rename(&temp_path, path).await?;
// Sync the directory to ensure the rename is persisted
if let Some(parent) = path.parent() {
match Self::sync_directory(parent).await {
Ok(_) => {}
Err(e) => {
warn!(
"Failed to sync directory {}: {}. File was written but directory entry might not be durable.",
parent.display(),
e
);
}
}
}
Ok(())
}
/// Creates or appends to a file with fsync
pub async fn write_with_sync<P: AsRef<Path>>(
path: P,
contents: &[u8],
append: bool,
) -> Result<(), IoError> {
let path = path.as_ref();
// Ensure parent directory exists
if let Some(parent) = path.parent() {
fs::create_dir_all(parent).await?;
}
// Open file with appropriate options
let mut file = OpenOptions::new()
.write(true)
.create(true)
.truncate(!append)
.append(append)
.open(path)
.await?;
// Write contents
file.write_all(contents).await?;
// Ensure data is synced to disk
file.flush().await?;
file.sync_all().await?;
Ok(())
}
/// Creates directories with fsync
pub async fn create_dir_with_sync<P: AsRef<Path>>(path: P) -> Result<(), IoError> {
let path = path.as_ref();
// Create directory
fs::create_dir_all(path).await?;
// Sync the directory
Self::sync_directory(path).await?;
// Sync parent directory to ensure directory creation is persisted
if let Some(parent) = path.parent() {
match Self::sync_directory(parent).await {
Ok(_) => {}
Err(e) => {
warn!(
"Failed to sync parent directory {}: {}. Directory was created but entry might not be durable.",
parent.display(),
e
);
}
}
}
Ok(())
}
/// Renames a file or directory with proper syncing
pub async fn rename_with_sync<P: AsRef<Path>, Q: AsRef<Path>>(
from: P,
to: Q,
) -> Result<(), IoError> {
let from = from.as_ref();
let to = to.as_ref();
// Ensure parent directory of destination exists
if let Some(parent) = to.parent() {
fs::create_dir_all(parent).await?;
}
// Perform rename
fs::rename(from, to).await?;
// Sync parent directories to ensure rename is persisted
if let Some(from_parent) = from.parent() {
match Self::sync_directory(from_parent).await {
Ok(_) => {}
Err(e) => {
warn!(
"Failed to sync source directory {}: {}. Rename completed but might not be durable.",
from_parent.display(),
e
);
}
}
}
if let Some(to_parent) = to.parent() {
match Self::sync_directory(to_parent).await {
Ok(_) => {}
Err(e) => {
warn!(
"Failed to sync destination directory {}: {}. Rename completed but might not be durable.",
to_parent.display(),
e
);
}
}
}
Ok(())
}
/// Removes a file with directory syncing
pub async fn remove_file_with_sync<P: AsRef<Path>>(path: P) -> Result<(), IoError> {
let path = path.as_ref();
// Remove file
fs::remove_file(path).await?;
// Sync parent directory to ensure removal is persisted
if let Some(parent) = path.parent() {
match Self::sync_directory(parent).await {
Ok(_) => {}
Err(e) => {
warn!(
"Failed to sync directory after file removal {}: {}. File was removed but entry might not be durable.",
parent.display(),
e
);
}
}
}
Ok(())
}
/// Removes a directory with parent directory syncing
pub async fn remove_dir_with_sync<P: AsRef<Path>>(
path: P,
recursive: bool,
) -> Result<(), IoError> {
let path = path.as_ref();
// Remove directory
if recursive {
fs::remove_dir_all(path).await?;
} else {
fs::remove_dir(path).await?;
}
// Sync parent directory to ensure removal is persisted
if let Some(parent) = path.parent() {
match Self::sync_directory(parent).await {
Ok(_) => {}
Err(e) => {
warn!(
"Failed to sync directory after directory removal {}: {}. Directory was removed but entry might not be durable.",
parent.display(),
e
);
}
}
}
Ok(())
}
/// Syncs a directory to ensure its contents are durable
async fn sync_directory<P: AsRef<Path>>(path: P) -> Result<(), IoError> {
let path = path.as_ref();
// Open directory with read permissions
let dir_file = match OpenOptions::new().read(true).open(path).await {
Ok(file) => file,
Err(e) => {
warn!(
"Failed to open directory for syncing {}: {}",
path.display(),
e
);
return Err(e);
}
};
// Sync the directory
dir_file.sync_all().await
}
}
#[cfg(test)]
mod tests {
use super::*;
use tempfile::tempdir;
use tokio::fs;
use tokio::io::AsyncReadExt;
#[tokio::test]
async fn test_atomic_write() {
let temp_dir = tempdir().unwrap();
let file_path = temp_dir.path().join("test.txt");
// Write data atomically
FileSystemUtils::atomic_write(&file_path, b"Hello, world!")
.await
.unwrap();
// Read back the data
let mut file = fs::File::open(&file_path).await.unwrap();
let mut contents = String::new();
file.read_to_string(&mut contents).await.unwrap();
assert_eq!(contents, "Hello, world!");
}
#[tokio::test]
async fn test_write_with_sync() {
let temp_dir = tempdir().unwrap();
let file_path = temp_dir.path().join("test.txt");
// Write data with sync
FileSystemUtils::write_with_sync(&file_path, b"First line\n", false)
.await
.unwrap();
// Append data
FileSystemUtils::write_with_sync(&file_path, b"Second line", true)
.await
.unwrap();
// Read back the data
let mut file = fs::File::open(&file_path).await.unwrap();
let mut contents = String::new();
file.read_to_string(&mut contents).await.unwrap();
assert_eq!(contents, "First line\nSecond line");
}
#[tokio::test]
async fn test_rename_with_sync() {
let temp_dir = tempdir().unwrap();
let source_path = temp_dir.path().join("source.txt");
let dest_path = temp_dir.path().join("dest.txt");
// Create source file
FileSystemUtils::write_with_sync(&source_path, b"Test content", false)
.await
.unwrap();
// Rename file
FileSystemUtils::rename_with_sync(&source_path, &dest_path)
.await
.unwrap();
// Verify source doesn't exist
assert!(!source_path.exists());
// Verify destination exists
let mut file = fs::File::open(&dest_path).await.unwrap();
let mut contents = String::new();
file.read_to_string(&mut contents).await.unwrap();
assert_eq!(contents, "Test content");
}
}
@@ -1,663 +0,0 @@
use async_trait::async_trait;
use std::collections::{HashMap, HashSet};
use std::sync::Arc;
use std::time::{Duration, Instant};
use tokio::sync::{Mutex, RwLock, Semaphore};
use tracing::{debug, error, info, warn};
use crate::application::ports::outbound::IdMappingPort;
use crate::common::errors::DomainError;
use crate::domain::services::path_service::StoragePath;
use crate::infrastructure::services::id_mapping_service::{IdMappingError, IdMappingService};
/// Maximum number of entries in the cache
const MAX_CACHE_SIZE: usize = 10_000;
/// Cache time-to-live (in seconds)
const CACHE_TTL_SECONDS: u64 = 60 * 5; // 5 minutes
/// Optimizer for batch ID mapping operations
pub struct IdMappingOptimizer {
/// Base ID mapping service
base_service: Arc<IdMappingService>,
/// Path to ID cache (path -> id)
path_to_id_cache: RwLock<HashMap<String, (String, Instant)>>,
/// ID to path cache (id -> path)
id_to_path_cache: RwLock<HashMap<String, (String, Instant)>>,
/// Hit counter
stats: RwLock<OptimizerStats>,
/// Semaphore to limit batch operations
batch_limiter: Semaphore,
/// Pending batch queue
pending_batch: Mutex<BatchQueue>,
}
/// Optimizer statistics
#[derive(Debug, Default, Clone)]
pub struct OptimizerStats {
/// Total number of get_path_by_id queries
pub path_by_id_queries: usize,
/// Number of cache hits for get_path_by_id
pub path_by_id_hits: usize,
/// Total number of get_or_create_id queries
pub get_id_queries: usize,
/// Number of cache hits for get_or_create_id
pub get_id_hits: usize,
/// Number of batch operations performed
pub batch_operations: usize,
/// Total number of IDs processed in batch
pub batch_items_processed: usize,
/// Last cache cleanup timestamp
pub last_cleanup: Option<Instant>,
}
/// Queue for batch operations
#[derive(Default)]
struct BatchQueue {
/// Pending paths to get/create ID
path_to_id_requests: HashSet<String>,
/// Pending IDs to get path
id_to_path_requests: HashSet<String>,
}
/// Result of a batch operation
struct BatchResult {
/// Path to ID mapping
path_to_id: HashMap<String, String>,
/// ID to path mapping
id_to_path: HashMap<String, String>,
}
impl IdMappingOptimizer {
/// Creates a new optimizer for the ID mapping service
pub fn new(base_service: Arc<IdMappingService>) -> Self {
Self {
base_service,
path_to_id_cache: RwLock::new(HashMap::with_capacity(1000)),
id_to_path_cache: RwLock::new(HashMap::with_capacity(1000)),
stats: RwLock::new(OptimizerStats::default()),
batch_limiter: Semaphore::new(2), // Limit to 2 concurrent batch operations
pending_batch: Mutex::new(BatchQueue::default()),
}
}
/// Gets optimizer statistics
pub async fn get_stats(&self) -> OptimizerStats {
self.stats.read().await.clone()
}
/// Cleans expired cache entries
pub async fn cleanup_cache(&self) {
let now = Instant::now();
let ttl = Duration::from_secs(CACHE_TTL_SECONDS);
// Clean path_to_id cache
{
let mut cache = self.path_to_id_cache.write().await;
let initial_size = cache.len();
// Retain only non-expired entries
cache.retain(|_, (_, timestamp)| now.duration_since(*timestamp) < ttl);
let removed = initial_size - cache.len();
if removed > 0 {
debug!("Cleaned {} expired entries from path_to_id cache", removed);
}
}
// Clean id_to_path cache
{
let mut cache = self.id_to_path_cache.write().await;
let initial_size = cache.len();
// Retain only non-expired entries
cache.retain(|_, (_, timestamp)| now.duration_since(*timestamp) < ttl);
let removed = initial_size - cache.len();
if removed > 0 {
debug!("Cleaned {} expired entries from id_to_path cache", removed);
}
}
// Update statistics
{
let mut stats = self.stats.write().await;
stats.last_cleanup = Some(now);
}
}
/// Starts periodic cleanup task
pub fn start_cleanup_task(optimizer: Arc<Self>) {
tokio::spawn(async move {
let cleanup_interval = Duration::from_secs(CACHE_TTL_SECONDS / 2);
loop {
tokio::time::sleep(cleanup_interval).await;
optimizer.cleanup_cache().await;
// Log statistics periodically
let stats = optimizer.get_stats().await;
info!(
"ID Mapping Optimizer stats - Path queries: {}, hits: {} ({}%), ID queries: {}, hits: {} ({}%), Batch ops: {}, items: {}",
stats.path_by_id_queries,
stats.path_by_id_hits,
if stats.path_by_id_queries > 0 {
stats.path_by_id_hits as f64 * 100.0 / stats.path_by_id_queries as f64
} else {
0.0
},
stats.get_id_queries,
stats.get_id_hits,
if stats.get_id_queries > 0 {
stats.get_id_hits as f64 * 100.0 / stats.get_id_queries as f64
} else {
0.0
},
stats.batch_operations,
stats.batch_items_processed
);
}
});
}
/// Adds a request to the pending queue for batch processing
async fn queue_path_to_id_request(
&self,
path: &StoragePath,
) -> Result<Option<String>, IdMappingError> {
let path_str = path.to_string();
// Check first in the cache
{
let cache = self.path_to_id_cache.read().await;
if let Some((id, _)) = cache.get(&path_str) {
// Update statistics
{
let mut stats = self.stats.write().await;
stats.get_id_hits += 1;
}
return Ok(Some(id.clone()));
}
}
// 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);
}
// Not found in cache, must be processed in batch
Ok(None)
}
/// Processes pending requests in batch
async fn process_batch(&self) -> Result<BatchResult, IdMappingError> {
// Acquire permit for batch operation
let _permit = self.batch_limiter.acquire().await.unwrap();
// Get pending requests
let (path_requests, id_requests) = {
let mut batch_queue = self.pending_batch.lock().await;
let paths = std::mem::take(&mut batch_queue.path_to_id_requests);
let ids = std::mem::take(&mut batch_queue.id_to_path_requests);
(paths, ids)
};
// Create results
let mut result = BatchResult {
path_to_id: HashMap::with_capacity(path_requests.len()),
id_to_path: HashMap::with_capacity(id_requests.len()),
};
// 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 {
Ok(id) => {
result.path_to_id.insert(path_str.clone(), id.clone());
result.id_to_path.insert(id, path_str);
}
Err(e) => {
error!("Error batch-processing path {}: {}", path_str, e);
// Continue with remaining requests
}
}
}
// Process id->path requests in batch
for id in id_requests {
match self.base_service.get_path_by_id(&id).await {
Ok(path) => {
let path_str = path.to_string();
result.id_to_path.insert(id.clone(), path_str.clone());
result.path_to_id.insert(path_str, id);
}
Err(e) => {
error!("Error batch-processing ID {}: {}", id, e);
// Continue with remaining requests
}
}
}
// Update cache with batch results
{
let mut path_cache = self.path_to_id_cache.write().await;
let mut id_cache = self.id_to_path_cache.write().await;
let now = Instant::now();
for (path, id) in &result.path_to_id {
path_cache.insert(path.clone(), (id.clone(), now));
}
for (id, path) in &result.id_to_path {
id_cache.insert(id.clone(), (path.clone(), now));
}
}
// 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();
}
// 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 {
error!("Error saving ID mapping changes: {}", e);
}
});
Ok(result)
}
/// Forces processing of pending requests if there are enough
async fn trigger_batch_if_needed(&self, min_batch_size: usize) -> Result<(), IdMappingError> {
// 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
};
// Process if necessary
if should_process {
self.process_batch().await?;
}
Ok(())
}
/// Preload a set of paths to get their IDs in batch
pub async fn preload_paths(&self, paths: Vec<StoragePath>) -> Result<(), IdMappingError> {
// Only proceed if there are paths to load
if paths.is_empty() {
return Ok(());
}
// Paths we need to load (those not in cache)
let mut paths_to_load = Vec::new();
// Check cache first
{
let cache = self.path_to_id_cache.read().await;
for path in paths {
let path_str = path.to_string();
if !cache.contains_key(&path_str) {
paths_to_load.push(path_str);
}
}
}
// If all were in cache, finish
if paths_to_load.is_empty() {
return Ok(());
}
// Add paths to queue for batch processing
{
let mut batch_queue = self.pending_batch.lock().await;
for path in paths_to_load {
batch_queue.path_to_id_requests.insert(path);
}
}
// Execute batch processing immediately
self.process_batch().await?;
Ok(())
}
/// Preload a set of IDs to get their paths in batch
pub async fn preload_ids(&self, ids: Vec<String>) -> Result<(), IdMappingError> {
// Only proceed if there are IDs to load
if ids.is_empty() {
return Ok(());
}
// IDs we need to load (those not in cache)
let mut ids_to_load = Vec::new();
// Check cache first
{
let cache = self.id_to_path_cache.read().await;
for id in ids {
if !cache.contains_key(&id) {
ids_to_load.push(id);
}
}
}
// If all were in cache, finish
if ids_to_load.is_empty() {
return Ok(());
}
// Add IDs to queue for batch processing
{
let mut batch_queue = self.pending_batch.lock().await;
for id in ids_to_load {
batch_queue.id_to_path_requests.insert(id);
}
}
// Execute batch processing immediately
self.process_batch().await?;
Ok(())
}
}
#[async_trait]
impl IdMappingPort for IdMappingOptimizer {
async fn get_or_create_id(&self, path: &StoragePath) -> Result<String, DomainError> {
// Update statistics
{
let mut stats = self.stats.write().await;
stats.get_id_queries += 1;
}
let path_str = path.to_string();
// Check cache first
{
let cache = self.path_to_id_cache.read().await;
if let Some((id, _)) = cache.get(&path_str) {
// Update statistics
{
let mut stats = self.stats.write().await;
stats.get_id_hits += 1;
}
return Ok(id.clone());
}
}
// If not in cache, try adding to batch queue first
let queued_result = self.queue_path_to_id_request(path).await?;
if let Some(id) = queued_result {
return Ok(id);
}
// Trigger batch processing if enough items accumulated
self.trigger_batch_if_needed(20).await?;
// Try to get from the base service
let id = self.base_service.get_or_create_id(path).await?;
// Update cache with the new ID
{
let mut path_cache = self.path_to_id_cache.write().await;
let mut id_cache = self.id_to_path_cache.write().await;
let now = Instant::now();
// Control cache size
if path_cache.len() >= MAX_CACHE_SIZE {
warn!(
"Path-to-ID cache size reached limit ({}), clearing oldest entries",
MAX_CACHE_SIZE
);
path_cache.clear();
}
if id_cache.len() >= MAX_CACHE_SIZE {
warn!(
"ID-to-path cache size reached limit ({}), clearing oldest entries",
MAX_CACHE_SIZE
);
id_cache.clear();
}
path_cache.insert(path_str.clone(), (id.clone(), now));
id_cache.insert(id.clone(), (path_str, now));
}
Ok(id)
}
async fn get_path_by_id(&self, id: &str) -> Result<StoragePath, DomainError> {
// Update statistics
{
let mut stats = self.stats.write().await;
stats.path_by_id_queries += 1;
}
// Check first in the cache
{
let cache = self.id_to_path_cache.read().await;
if let Some((path_str, _)) = cache.get(id) {
// Update statistics
{
let mut stats = self.stats.write().await;
stats.path_by_id_hits += 1;
}
return Ok(StoragePath::from_string(path_str));
}
}
// Get from the base service
let path = self.base_service.get_path_by_id(id).await?;
// Update cache
{
let mut id_cache = self.id_to_path_cache.write().await;
let mut path_cache = self.path_to_id_cache.write().await;
let now = Instant::now();
let path_str = path.to_string();
// Control cache size
if id_cache.len() >= MAX_CACHE_SIZE {
warn!(
"ID-to-path cache size reached limit ({}), clearing oldest entries",
MAX_CACHE_SIZE
);
id_cache.clear();
}
if path_cache.len() >= MAX_CACHE_SIZE {
warn!(
"Path-to-ID cache size reached limit ({}), clearing oldest entries",
MAX_CACHE_SIZE
);
path_cache.clear();
}
id_cache.insert(id.to_string(), (path_str.clone(), now));
path_cache.insert(path_str, (id.to_string(), now));
}
Ok(path)
}
async fn update_path(&self, id: &str, new_path: &StoragePath) -> Result<(), DomainError> {
// 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;
// Remove the ID entry
if let Some((old_path, _)) = id_cache.remove(id) {
path_cache.remove(&old_path);
}
}
// Update in the base service
let result = self.base_service.update_path(id, new_path).await?;
// 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;
let now = Instant::now();
let path_str = new_path.to_string();
id_cache.insert(id.to_string(), (path_str.clone(), now));
path_cache.insert(path_str, (id.to_string(), now));
}
Ok(result)
}
async fn remove_id(&self, id: &str) -> Result<(), DomainError> {
// 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;
// Remove the ID entry
if let Some((path, _)) = id_cache.remove(id) {
path_cache.remove(&path);
}
}
// Remove from the base service
self.base_service.remove_id(id).await?;
Ok(())
}
async fn save_changes(&self) -> Result<(), DomainError> {
// Delegate to the base service
self.base_service.save_changes().await?;
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use tempfile::tempdir;
async fn create_test_service() -> (Arc<IdMappingService>, Arc<IdMappingOptimizer>) {
let temp_dir = tempdir().unwrap();
let map_path = temp_dir.path().join("id_map.json");
let base_service = Arc::new(IdMappingService::new(map_path).await.unwrap());
let optimizer = Arc::new(IdMappingOptimizer::new(base_service.clone()));
(base_service, optimizer)
}
#[tokio::test]
async fn test_basic_caching() {
let (_, optimizer) = create_test_service().await;
let path = StoragePath::from_string("/test/file.txt");
// 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");
// 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");
// 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");
}
#[tokio::test]
async fn test_batch_processing() {
let (_, optimizer) = create_test_service().await;
// 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
)));
}
// Preload the paths
optimizer.preload_paths(paths.clone()).await.unwrap();
// 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");
}
// 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"
);
// Verify all subsequent queries are cache hits
assert_eq!(
stats.get_id_hits, 50,
"All subsequente queries should be cache hits"
);
}
#[tokio::test]
async fn test_cache_cleanup() {
let (_, optimizer) = create_test_service().await;
// Create some entries
let path = StoragePath::from_string("/test/cleanup.txt");
let id = optimizer.get_or_create_id(&path).await.unwrap();
// 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");
}
// Run cleanup (should not remove anything yet)
optimizer.cleanup_cache().await;
// 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");
{
let stats = optimizer.get_stats().await;
assert_eq!(stats.get_id_hits, 1, "Should have 1 hit after cleanup");
}
}
}
@@ -1,712 +0,0 @@
use async_trait::async_trait;
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use std::path::PathBuf;
use tokio::fs;
use tokio::sync::{Mutex, RwLock};
use tokio::time;
use uuid::Uuid;
use crate::application::ports::outbound::IdMappingPort;
use crate::common::config::TimeoutConfig;
use crate::common::errors::{DomainError, ErrorKind};
use crate::domain::services::path_service::StoragePath;
/// Specific error for the ID mapping service
#[derive(Debug, thiserror::Error)]
pub enum IdMappingError {
#[error("ID not found: {0}")]
NotFound(String),
#[error("IO error: {0}")]
IoError(#[from] std::io::Error),
#[error("Timeout error: {0}")]
Timeout(String),
#[error("Serialization error: {0}")]
SerializationError(#[from] serde_json::Error),
#[error("Other error: {0}")]
Other(String),
}
// Implement conversion from IdMappingError to DomainError
impl From<IdMappingError> for DomainError {
fn from(err: IdMappingError) -> Self {
match err {
IdMappingError::NotFound(id) => DomainError::not_found("IdMapping", id),
IdMappingError::IoError(e) => DomainError::new(
ErrorKind::InternalError,
"IdMapping",
format!("IO error: {}", e),
)
.with_source(e),
IdMappingError::Timeout(msg) => {
DomainError::timeout("IdMapping", format!("Timeout: {}", msg))
}
IdMappingError::SerializationError(e) => DomainError::new(
ErrorKind::InternalError,
"IdMapping",
format!("Serialization error: {}", e),
)
.with_source(e),
IdMappingError::Other(msg) => DomainError::new(
ErrorKind::InternalError,
"IdMapping",
format!("Other error: {}", msg),
),
}
}
}
/// 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>, // Field for efficient bidirectional lookup
version: u32, // Version to detect changes
}
/// Service to manage mappings between paths and unique IDs
pub struct IdMappingService {
map_path: PathBuf,
id_map: RwLock<IdMap>,
save_mutex: Mutex<()>, // To prevent multiple concurrent saves
timeouts: TimeoutConfig,
pending_save: RwLock<bool>, // Indicates if there are pending changes
}
impl IdMappingService {
/// 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?;
Ok(Self {
map_path,
id_map: RwLock::new(id_map),
save_mutex: Mutex::new(()),
timeouts,
pending_save: RwLock::new(false),
})
}
/// 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 {
Self {
map_path: PathBuf::from("/tmp/dummy_id_map.json"),
id_map: RwLock::new(IdMap::default()),
save_mutex: Mutex::new(()),
timeouts: TimeoutConfig::default(),
pending_save: RwLock::new(false),
}
}
/// Creates an in-memory ID mapping service (for testing - original version)
pub fn new_in_memory() -> Self {
Self {
map_path: PathBuf::from("memory"),
id_map: RwLock::new(IdMap::default()),
save_mutex: Mutex::new(()),
timeouts: TimeoutConfig::default(),
pending_save: RwLock::new(false),
}
}
/// 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() {
// Try to read with timeout to avoid indefinite blocking
let read_result = time::timeout(timeouts.lock_timeout(), fs::read_to_string(map_path))
.await
.map_err(|_| {
DomainError::timeout(
"IdMapping",
format!("Timeout reading 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),
)
})?;
// Parse the JSON
match serde_json::from_str::<IdMap>(&content) {
Ok(mut map) => {
// 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 {
map.id_to_path.insert(id.clone(), path.clone());
rebuild_count += 1;
}
tracing::info!("Rebuilt inverse mapping with {} entries", rebuild_count);
}
tracing::info!(
"Loaded ID map with {} entries (version: {})",
map.path_to_id.len(),
map.version
);
return Ok(map);
}
Err(e) => {
tracing::error!("Error parsing ID map: {}", e);
// 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);
} else {
tracing::info!("Backed up corrupted ID map to {}", backup_path.display());
}
tracing::info!("Creating new empty map after error");
return Ok(IdMap {
path_to_id: HashMap::new(),
id_to_path: HashMap::new(),
version: 1, // Start with version 1
});
}
}
}
// 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, // Start with version 1
};
// Ensure directory exists
if let Some(parent) = map_path.parent()
&& !parent.exists()
&& let Err(e) = fs::create_dir_all(parent).await
{
tracing::error!("Failed to create directory for ID map: {}", e);
}
// Write empty map to file (best-effort: the in-memory map is valid even if disk write fails)
match serde_json::to_string_pretty(&empty_map) {
Ok(json) => {
if let Err(e) = fs::write(map_path, json).await {
tracing::warn!(
"Could not write initial empty ID map (will retry on next save): {}",
e
);
} else {
tracing::info!("Created initial empty ID map at {}", map_path.display());
}
}
Err(e) => {
tracing::error!("Failed to serialize empty ID map: {}", e);
}
}
Ok(empty_map)
}
/// Saves the ID map to disk safely
async fn save_id_map(&self) -> Result<(), DomainError> {
// 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")
})?;
// Create JSON with read lock to minimize lock hold time
let json = {
let mut map = time::timeout(self.timeouts.lock_timeout(), self.id_map.write())
.await
.map_err(|_| {
DomainError::timeout("IdMapping", "Timeout acquiring write lock for ID mapping")
})?;
// Increment version only if there are pending changes to save
let pending = *self.pending_save.read().await;
if pending {
map.version += 1;
tracing::debug!("Incrementing ID map version to {}", map.version);
}
// Use serde with reasonably safe defaults
serde_json::to_string_pretty(&*map).map_err(|e| {
DomainError::internal_error(
"IdMapping",
format!("Failed to serialize ID map to JSON: {}", e),
)
})?
};
// 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
),
)
})?;
// 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
),
)
})?;
// Reset pending flag
{
let mut pending = self.pending_save.write().await;
*pending = false;
}
tracing::info!("Saved ID map successfully to {}", self.map_path.display());
Ok(())
}
/// Generates a unique ID
fn generate_id(&self) -> String {
Uuid::new_v4().to_string()
}
/// Marks changes as pending
async fn mark_pending(&self) {
let mut pending = self.pending_save.write().await;
*pending = true;
}
/// 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();
// First attempt with read lock (more efficient)
{
let read_result =
match time::timeout(self.timeouts.lock_timeout(), self.id_map.read()).await {
Ok(guard) => guard,
Err(_) => {
return Err(IdMappingError::Timeout(
"Timeout acquiring read lock for ID mapping".to_string(),
));
}
};
if let Some(id) = read_result.path_to_id.get(&path_str) {
return Ok(id.clone());
}
}
// If not found, acquire write lock
let write_result =
match time::timeout(self.timeouts.lock_timeout(), self.id_map.write()).await {
Ok(guard) => guard,
Err(_) => {
return Err(IdMappingError::Timeout(
"Timeout acquiring write lock for ID mapping".to_string(),
));
}
};
let mut map = write_result;
// 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());
}
// 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);
// 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);
Ok(id)
}
/// 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(), self.id_map.read()).await {
Ok(guard) => guard,
Err(_) => {
return Err(IdMappingError::Timeout(
"Timeout acquiring read lock for ID lookup".to_string(),
));
}
};
if let Some(path_str) = read_result.id_to_path.get(id) {
return Ok(StoragePath::from_string(path_str));
}
Err(IdMappingError::NotFound(id.to_string()))
}
/// 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(), self.id_map.write()).await {
Ok(guard) => guard,
Err(_) => {
return Err(IdMappingError::Timeout(
"Timeout acquiring write lock for ID update".to_string(),
));
}
};
let mut map = write_result;
// 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);
// 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);
// Mark as pending
drop(map); // Release the write lock before acquiring another
self.mark_pending().await;
tracing::debug!(
"Updated path mapping for ID {}: {} -> {}",
id,
old_path,
new_path.to_string()
);
Ok(())
} else {
Err(IdMappingError::NotFound(id.to_string()))
}
}
/// 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(), self.id_map.write()).await {
Ok(guard) => guard,
Err(_) => {
return Err(IdMappingError::Timeout(
"Timeout acquiring write lock for ID removal".to_string(),
));
}
};
let mut map = write_result;
// Find the path to remove it
if let Some(path) = map.id_to_path.remove(id) {
map.path_to_id.remove(&path);
// Mark as pending
drop(map); // Release the write lock before acquiring another
self.mark_pending().await;
tracing::debug!("Removed ID mapping: {} -> {}", id, path);
Ok(())
} else {
Err(IdMappingError::NotFound(id.to_string()))
}
}
/// Saves pending changes to disk immediately, without debounce
pub async fn save_pending_changes(&self) -> Result<(), IdMappingError> {
// Check if there are pending changes
{
let pending = self.pending_save.read().await;
if !*pending {
return Ok(());
}
}
// 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()
);
// Explicitly verify that the file exists and has size
match std::fs::metadata(&self.map_path) {
Ok(metadata) => {
if metadata.len() > 0 {
tracing::info!(
"Verified saved map file exists with size: {} bytes",
metadata.len()
);
} else {
tracing::warn!(
"Map file exists but has zero size - this might cause issues"
);
}
}
Err(e) => {
tracing::error!("Failed to verify saved map file: {}", e);
// 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::other(format!(
"Failed to verify and retry save: {}",
retry_err
))));
}
tracing::info!("Second save attempt succeeded");
}
}
Ok(())
}
Err(e) => {
tracing::error!(
"Failed to save ID map to {}: {}",
self.map_path.display(),
e
);
// 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(_) => {
tracing::info!("Second save attempt succeeded after initial failure");
Ok(())
}
Err(retry_e) => {
tracing::error!("Second save attempt also failed: {}", retry_e);
Err(IdMappingError::IoError(std::io::Error::other(format!(
"Failed to save ID mappings after retry: {}",
retry_e
))))
}
}
}
}
}
}
#[async_trait]
impl IdMappingPort for IdMappingService {
/// 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
),
)
})
}
/// 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),
)
})
}
/// 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
),
)
})
}
/// 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))
})
}
/// 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),
)
})
}
}
// The extension methods were moved to the IdMappingPort trait as default implementations
// 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)
pub fn new_sync(map_path: PathBuf) -> Self {
// Create a minimal implementation for initialization purposes
Self {
map_path,
id_map: RwLock::new(IdMap::default()),
save_mutex: Mutex::new(()),
timeouts: TimeoutConfig::default(),
pending_save: RwLock::new(false),
}
}
}
impl Clone for IdMappingService {
fn clone(&self) -> Self {
// 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()), // This is not used in the async task
save_mutex: Mutex::new(()), // Neither is this
timeouts: self.timeouts.clone(),
pending_save: RwLock::new(false),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::time::Duration;
use tempfile::tempdir;
#[tokio::test]
async fn test_get_or_create_id() {
let temp_dir = tempdir().unwrap();
let map_path = temp_dir.path().join("id_map.json");
let service = IdMappingService::new(map_path).await.unwrap();
let path = StoragePath::from_string("/test/file.txt");
let id = service.get_or_create_id(&path).await.unwrap();
assert!(!id.is_empty(), "ID should not be empty");
// 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");
}
#[tokio::test]
async fn test_update_path() {
let temp_dir = tempdir().unwrap();
let map_path = temp_dir.path().join("id_map.json");
let service = IdMappingService::new(map_path).await.unwrap();
let old_path = StoragePath::from_string("/test/old.txt");
let id = service.get_or_create_id(&old_path).await.unwrap();
let new_path = StoragePath::from_string("/test/new.txt");
service.update_path(&id, &new_path).await.unwrap();
let retrieved_path = service.get_path_by_id(&id).await.unwrap();
assert_eq!(retrieved_path, new_path, "Path should be updated");
}
#[tokio::test]
async fn test_save_and_load() {
let temp_dir = tempdir().unwrap();
let map_path = temp_dir.path().join("id_map.json");
// Create and populate the service
let service = IdMappingService::new(map_path.clone()).await.unwrap();
let path1 = StoragePath::from_string("/test/file1.txt");
let path2 = StoragePath::from_string("/test/file2.txt");
let id1 = service.get_or_create_id(&path1).await.unwrap();
let id2 = service.get_or_create_id(&path2).await.unwrap();
// Save changes
service.save_pending_changes().await.unwrap();
// Wait to ensure the async save completes
tokio::time::sleep(Duration::from_millis(500)).await;
// Create a new service that should load the same map
let service2 = IdMappingService::new(map_path).await.unwrap();
// 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();
assert_eq!(id1, loaded_id1, "ID1 should be preserved");
assert_eq!(id2, loaded_id2, "ID2 should be preserved");
}
#[tokio::test]
async fn test_concurrent_operations() {
use futures::future::join_all;
let temp_dir = tempdir().unwrap();
let map_path = temp_dir.path().join("id_map.json");
let service = std::sync::Arc::new(IdMappingService::new(map_path).await.unwrap());
// 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));
let service_clone = service.clone();
tasks.push(tokio::spawn(async move {
service_clone.get_or_create_id(&path).await
}));
}
// Wait for all to finish
let results = join_all(tasks).await;
// Verify that all succeeded
for result in results {
assert!(
result.unwrap().is_ok(),
"Concurrent operations should succeed"
);
}
// Save changes
service.save_pending_changes().await.unwrap();
}
}
-5
View File
@@ -1,13 +1,8 @@
pub mod buffer_pool;
pub mod chunked_upload_service;
pub mod compression_service;
pub mod dedup_service;
pub mod file_content_cache;
pub mod file_metadata_cache;
pub mod file_system_i18n_service;
pub mod file_system_utils;
pub mod id_mapping_optimizer;
pub mod id_mapping_service;
pub mod image_transcode_service;
pub mod jwt_service;
pub mod oidc_service;