adaptation de la crate pmocovers
This commit is contained in:
@@ -7,7 +7,7 @@
|
||||
//! - Supprimer des items
|
||||
//! - Purger et consolider le cache
|
||||
|
||||
use crate::{Cache, CacheConfig, CacheEntry};
|
||||
use crate::{Cache, CacheConfig};
|
||||
use axum::{
|
||||
extract::{Path, State},
|
||||
http::StatusCode,
|
||||
|
||||
@@ -11,6 +11,7 @@ use std::collections::HashMap;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::sync::Arc;
|
||||
use tokio::sync::RwLock;
|
||||
use tracing;
|
||||
|
||||
/// Trait pour définir les paramètres du cache
|
||||
pub trait CacheConfig: Send + Sync {
|
||||
@@ -46,7 +47,6 @@ pub trait CacheConfig: Send + Sync {
|
||||
/// Note : Ce type est conçu pour être utilisé derrière un `Arc<Cache>`.
|
||||
/// La synchronisation est gérée par le Mutex interne de la base de données SQLite
|
||||
/// et par le RwLock pour la map des downloads.
|
||||
#[derive(Debug)]
|
||||
pub struct Cache<C: CacheConfig> {
|
||||
/// Répertoire de stockage
|
||||
dir: PathBuf,
|
||||
@@ -58,12 +58,14 @@ pub struct Cache<C: CacheConfig> {
|
||||
pub db: Arc<DB>,
|
||||
/// Map des downloads en cours (pk -> Download)
|
||||
downloads: Arc<RwLock<HashMap<String, Arc<Download>>>>,
|
||||
/// Factory pour créer des transformers (optionnel)
|
||||
transformer_factory: Option<Arc<dyn Fn() -> StreamTransformer + Send + Sync>>,
|
||||
/// Phantom data pour le type de configuration
|
||||
_phantom: std::marker::PhantomData<C>,
|
||||
}
|
||||
|
||||
impl<C: CacheConfig> Cache<C> {
|
||||
/// Crée un nouveau cache
|
||||
/// Crée un nouveau cache sans transformer
|
||||
///
|
||||
/// # Arguments
|
||||
///
|
||||
@@ -71,6 +73,52 @@ impl<C: CacheConfig> Cache<C> {
|
||||
/// * `limit` - Limite de taille du cache (nombre d'éléments)
|
||||
/// * `base_url` - URL de base pour la génération d'URLs
|
||||
pub fn new(dir: &str, limit: usize, base_url: &str) -> Result<Self> {
|
||||
Self::with_transformer(dir, limit, base_url, None)
|
||||
}
|
||||
|
||||
/// Crée un nouveau cache avec un transformer optionnel
|
||||
///
|
||||
/// # Arguments
|
||||
///
|
||||
/// * `dir` - Répertoire de stockage du cache
|
||||
/// * `limit` - Limite de taille du cache (nombre d'éléments)
|
||||
/// * `base_url` - URL de base pour la génération d'URLs
|
||||
/// * `transformer_factory` - Factory pour créer des transformers à chaque téléchargement
|
||||
///
|
||||
/// # Exemple
|
||||
///
|
||||
/// ```rust,no_run
|
||||
/// use pmocache::{Cache, CacheConfig, StreamTransformer};
|
||||
/// use std::sync::Arc;
|
||||
///
|
||||
/// struct MyConfig;
|
||||
/// impl CacheConfig for MyConfig {
|
||||
/// fn file_extension() -> &'static str { "dat" }
|
||||
/// }
|
||||
///
|
||||
/// let transformer_factory = Arc::new(|| {
|
||||
/// // Créer un transformer qui convertit les données
|
||||
/// Box::new(|response, file, progress| {
|
||||
/// Box::pin(async move {
|
||||
/// // Transformation personnalisée
|
||||
/// Ok(())
|
||||
/// })
|
||||
/// }) as StreamTransformer
|
||||
/// });
|
||||
///
|
||||
/// let cache = Cache::<MyConfig>::with_transformer(
|
||||
/// "./cache",
|
||||
/// 1000,
|
||||
/// "http://localhost:8080",
|
||||
/// Some(transformer_factory)
|
||||
/// ).unwrap();
|
||||
/// ```
|
||||
pub fn with_transformer(
|
||||
dir: &str,
|
||||
limit: usize,
|
||||
base_url: &str,
|
||||
transformer_factory: Option<Arc<dyn Fn() -> StreamTransformer + Send + Sync>>,
|
||||
) -> Result<Self> {
|
||||
let directory = PathBuf::from(dir);
|
||||
std::fs::create_dir_all(&directory)?;
|
||||
let db = DB::init(&directory.join("cache.db"), C::table_name())?;
|
||||
@@ -81,18 +129,11 @@ impl<C: CacheConfig> Cache<C> {
|
||||
base_url: base_url.to_string(),
|
||||
db: Arc::new(db),
|
||||
downloads: Arc::new(RwLock::new(HashMap::new())),
|
||||
transformer_factory,
|
||||
_phantom: std::marker::PhantomData,
|
||||
})
|
||||
}
|
||||
|
||||
/// Retourne le transformer pour ce cache
|
||||
///
|
||||
/// Par défaut retourne None (pas de transformation).
|
||||
/// Les caches spécialisés peuvent surcharger cette méthode.
|
||||
fn get_transformer(&self) -> Option<StreamTransformer> {
|
||||
None
|
||||
}
|
||||
|
||||
/// Télécharge un fichier depuis une URL et l'ajoute au cache
|
||||
///
|
||||
/// Utilise le module download pour gérer le téléchargement asynchrone.
|
||||
@@ -120,10 +161,11 @@ impl<C: CacheConfig> Cache<C> {
|
||||
}
|
||||
|
||||
// Lancer le téléchargement avec transformer
|
||||
let transformer = self.transformer_factory.as_ref().map(|f| f());
|
||||
let download = download_with_transformer(
|
||||
&file_path,
|
||||
url,
|
||||
self.get_transformer(),
|
||||
transformer,
|
||||
);
|
||||
|
||||
// Stocker dans la map des downloads en cours
|
||||
@@ -135,6 +177,12 @@ impl<C: CacheConfig> Cache<C> {
|
||||
// Ajouter immédiatement à la DB
|
||||
self.db.add(&pk, url, collection)?;
|
||||
|
||||
// Appliquer la politique d'éviction LRU si nécessaire
|
||||
// Cela garantit que le cache respecte toujours la limite configurée
|
||||
if let Err(e) = self.enforce_limit().await {
|
||||
tracing::warn!("Error enforcing cache limit: {}", e);
|
||||
}
|
||||
|
||||
// Lancer une tâche de nettoyage en background
|
||||
let downloads_clone = self.downloads.clone();
|
||||
let pk_clone = pk.clone();
|
||||
@@ -413,11 +461,78 @@ impl<C: CacheConfig> Cache<C> {
|
||||
&self.base_url
|
||||
}
|
||||
|
||||
/// Construit le chemin complet d'un fichier dans le cache avec le param par défaut
|
||||
///
|
||||
/// Format: `{pk}.{default_param}.{extension}`
|
||||
pub fn file_path(&self, pk: &str) -> PathBuf {
|
||||
self.file_path_with_qualifier(pk, C::default_param())
|
||||
}
|
||||
|
||||
/// Construit le chemin d'un fichier dans le cache avec un qualificatif
|
||||
///
|
||||
/// Format: `{pk}.{qualifier}.{extension}`
|
||||
pub fn file_path_with_qualifier(&self, pk: &str, qualifier: &str) -> PathBuf {
|
||||
self.dir.join(format!("{}.{}.{}", pk, qualifier, C::file_extension()))
|
||||
}
|
||||
|
||||
/// Valide les données avant de les stocker
|
||||
/// Par défaut, accepte toutes les données
|
||||
pub fn validate_data(&self, data: &[u8]) -> Result<Vec<u8>> {
|
||||
Ok(data.to_vec())
|
||||
}
|
||||
|
||||
/// Applique la politique d'éviction LRU (Least Recently Used)
|
||||
///
|
||||
/// Si le nombre d'entrées dépasse la limite configurée, supprime
|
||||
/// les entrées les plus anciennes (moins récemment utilisées).
|
||||
///
|
||||
/// Cette méthode :
|
||||
/// 1. Compte le nombre total d'entrées
|
||||
/// 2. Si > limit, récupère les N entrées les plus anciennes
|
||||
/// 3. Supprime ces entrées de la DB et leurs fichiers du disque
|
||||
///
|
||||
/// # Returns
|
||||
///
|
||||
/// Le nombre d'entrées supprimées
|
||||
pub async fn enforce_limit(&self) -> Result<usize> {
|
||||
let count = self.db.count()?;
|
||||
|
||||
if count <= self.limit {
|
||||
return Ok(0);
|
||||
}
|
||||
|
||||
let to_remove = count - self.limit;
|
||||
let old_entries = self.db.get_oldest(to_remove)?;
|
||||
|
||||
let mut removed = 0;
|
||||
for entry in old_entries {
|
||||
// Supprimer tous les fichiers avec ce pk (toutes variantes)
|
||||
if let Ok(mut dir_entries) = tokio::fs::read_dir(&self.dir).await {
|
||||
while let Ok(Some(dir_entry)) = dir_entries.next_entry().await {
|
||||
if let Some(filename) = dir_entry.file_name().to_str() {
|
||||
// Format: {pk}.{param}.{ext}
|
||||
if filename.starts_with(&entry.pk) && filename.starts_with(&format!("{}.", entry.pk)) {
|
||||
let _ = tokio::fs::remove_file(dir_entry.path()).await;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Supprimer de la base de données
|
||||
if let Err(e) = self.db.delete(&entry.pk) {
|
||||
tracing::warn!("Error deleting entry {} from DB: {}", entry.pk, e);
|
||||
} else {
|
||||
removed += 1;
|
||||
}
|
||||
}
|
||||
|
||||
if removed > 0 {
|
||||
tracing::info!("LRU eviction: removed {} old entries (cache size: {} -> {})",
|
||||
removed, count, count - removed);
|
||||
}
|
||||
|
||||
Ok(removed)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -248,4 +248,54 @@ impl DB {
|
||||
conn.execute(&sql, [pk])?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Compte le nombre total d'entrées dans le cache
|
||||
///
|
||||
/// # Returns
|
||||
///
|
||||
/// Le nombre total d'entrées
|
||||
pub fn count(&self) -> rusqlite::Result<usize> {
|
||||
let conn = self.conn.lock().unwrap();
|
||||
let sql = format!("SELECT COUNT(*) FROM {}", self.table_name);
|
||||
let count: i64 = conn.query_row(&sql, [], |row| row.get(0))?;
|
||||
Ok(count as usize)
|
||||
}
|
||||
|
||||
/// Récupère les N entrées les plus anciennes (LRU - Least Recently Used)
|
||||
///
|
||||
/// Trie par last_used (les plus anciens en premier), puis par hits (les moins utilisés).
|
||||
/// Utile pour implémenter une politique d'éviction LRU.
|
||||
///
|
||||
/// # Arguments
|
||||
///
|
||||
/// * `limit` - Nombre maximum d'entrées à récupérer
|
||||
///
|
||||
/// # Returns
|
||||
///
|
||||
/// Liste des entrées les plus anciennes, triées par last_used ASC
|
||||
pub fn get_oldest(&self, limit: usize) -> rusqlite::Result<Vec<CacheEntry>> {
|
||||
let conn = self.conn.lock().unwrap();
|
||||
let sql = format!(
|
||||
"SELECT pk, source_url, collection, hits, last_used
|
||||
FROM {}
|
||||
ORDER BY last_used ASC, hits ASC
|
||||
LIMIT ?1",
|
||||
self.table_name
|
||||
);
|
||||
|
||||
let mut stmt = conn.prepare(&sql)?;
|
||||
|
||||
let entries = stmt.query_map([limit], |row| {
|
||||
Ok(CacheEntry {
|
||||
pk: row.get(0)?,
|
||||
source_url: row.get(1)?,
|
||||
collection: row.get(2)?,
|
||||
hits: row.get(3)?,
|
||||
last_used: row.get(4)?,
|
||||
})
|
||||
})?
|
||||
.collect::<rusqlite::Result<Vec<_>>>()?;
|
||||
|
||||
Ok(entries)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -40,8 +40,6 @@
|
||||
#[cfg(feature = "pmoserver")]
|
||||
use crate::{Cache, CacheConfig};
|
||||
#[cfg(feature = "pmoserver")]
|
||||
use crate::cache_trait::FileCache;
|
||||
#[cfg(feature = "pmoserver")]
|
||||
use axum::{
|
||||
body::Body,
|
||||
extract::{Path, State},
|
||||
@@ -56,42 +54,83 @@ use std::sync::Arc;
|
||||
use tokio_util::io::ReaderStream;
|
||||
#[cfg(feature = "pmoserver")]
|
||||
use tracing::warn;
|
||||
#[cfg(feature = "pmoserver")]
|
||||
use std::pin::Pin;
|
||||
#[cfg(feature = "pmoserver")]
|
||||
use std::future::Future;
|
||||
|
||||
/// Type pour le callback de génération de param
|
||||
///
|
||||
/// Appelé quand un fichier avec param n'existe pas.
|
||||
/// Permet de générer à la volée (ex: redimensionnement d'images).
|
||||
///
|
||||
/// # Arguments
|
||||
///
|
||||
/// - `cache`: le cache
|
||||
/// - `pk`: clé primaire
|
||||
/// - `param`: paramètre demandé (ex: "256" pour une taille)
|
||||
///
|
||||
/// # Retourne
|
||||
///
|
||||
/// Les données générées ou None si le param n'est pas supporté
|
||||
#[cfg(feature = "pmoserver")]
|
||||
pub type ParamGenerator<C> = Arc<
|
||||
dyn Fn(Arc<Cache<C>>, String, String)
|
||||
-> Pin<Box<dyn Future<Output = Option<Vec<u8>>> + Send>>
|
||||
+ Send + Sync
|
||||
>;
|
||||
|
||||
/// Handler générique pour GET /{cache_name}/{cache_type}/{pk}
|
||||
/// Sert un fichier avec le param par défaut
|
||||
#[cfg(feature = "pmoserver")]
|
||||
async fn get_file<C: CacheConfig + 'static>(
|
||||
State((cache, content_type)): State<(Arc<Cache<C>>, &'static str)>,
|
||||
State((cache, content_type, param_generator)): State<(Arc<Cache<C>>, &'static str, Option<ParamGenerator<C>>)>,
|
||||
Path(pk): Path<String>,
|
||||
) -> Response {
|
||||
// Utiliser le param par défaut
|
||||
let param = C::default_param();
|
||||
serve_file_with_streaming(&cache, &pk, param, content_type).await
|
||||
serve_file_with_streaming(&cache, &pk, param, content_type, param_generator).await
|
||||
}
|
||||
|
||||
/// Handler générique pour GET /{cache_name}/{cache_type}/{pk}/{param}
|
||||
/// Sert un fichier avec un param spécifique
|
||||
#[cfg(feature = "pmoserver")]
|
||||
async fn get_file_with_param<C: CacheConfig + 'static>(
|
||||
State((cache, content_type)): State<(Arc<Cache<C>>, &'static str)>,
|
||||
State((cache, content_type, param_generator)): State<(Arc<Cache<C>>, &'static str, Option<ParamGenerator<C>>)>,
|
||||
Path((pk, param)): Path<(String, String)>,
|
||||
) -> Response {
|
||||
serve_file_with_streaming(&cache, &pk, ¶m, content_type).await
|
||||
serve_file_with_streaming(&cache, &pk, ¶m, content_type, param_generator).await
|
||||
}
|
||||
|
||||
/// Fonction utilitaire pour servir un fichier avec streaming progressif
|
||||
///
|
||||
/// Si le fichier est en cours de téléchargement, il est streamé au fur et à mesure.
|
||||
/// Sinon, le fichier complet est servi normalement.
|
||||
/// Si le fichier n'existe pas et qu'un param_generator est fourni, tente de générer le param.
|
||||
#[cfg(feature = "pmoserver")]
|
||||
async fn serve_file_with_streaming<C: CacheConfig>(
|
||||
cache: &Arc<Cache<C>>,
|
||||
pk: &str,
|
||||
param: &str,
|
||||
content_type: &'static str,
|
||||
param_generator: Option<ParamGenerator<C>>,
|
||||
) -> Response {
|
||||
let file_path = cache.file_path_with_qualifier(pk, param);
|
||||
|
||||
// Si le fichier n'existe pas et qu'on a un générateur, l'utiliser
|
||||
if !file_path.exists() {
|
||||
if let Some(generator) = param_generator {
|
||||
if let Some(data) = generator(cache.clone(), pk.to_string(), param.to_string()).await {
|
||||
// Le générateur a créé les données, les servir directement
|
||||
return (
|
||||
StatusCode::OK,
|
||||
[("content-type", content_type)],
|
||||
data,
|
||||
).into_response();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Mettre à jour les stats d'utilisation
|
||||
if let Err(e) = cache.db.update_hit(pk) {
|
||||
warn!("Error updating hit count for {}: {}", pk, e);
|
||||
@@ -218,18 +257,63 @@ async fn serve_complete_file(
|
||||
pub fn create_file_router<C: CacheConfig + 'static>(
|
||||
cache: Arc<Cache<C>>,
|
||||
content_type: &'static str,
|
||||
) -> Router {
|
||||
create_file_router_with_generator(cache, content_type, None)
|
||||
}
|
||||
|
||||
/// Crée un router pour servir les fichiers d'un cache avec générateur de param
|
||||
///
|
||||
/// Similaire à `create_file_router` mais permet de fournir un générateur
|
||||
/// pour créer des variantes à la volée (ex: redimensionnement d'images).
|
||||
///
|
||||
/// # Arguments
|
||||
///
|
||||
/// * `cache` - Instance du cache
|
||||
/// * `content_type` - Type MIME des fichiers (ex: "image/webp", "audio/flac")
|
||||
/// * `param_generator` - Générateur optionnel pour créer des params à la volée
|
||||
///
|
||||
/// # Exemple
|
||||
///
|
||||
/// ```rust,no_run
|
||||
/// use pmocache::pmoserver_ext::{create_file_router_with_generator, ParamGenerator};
|
||||
/// use std::sync::Arc;
|
||||
///
|
||||
/// # async fn example(cache: std::sync::Arc<pmocache::Cache<CoversConfig>>) {
|
||||
/// let generator: ParamGenerator<CoversConfig> = Arc::new(|cache, pk, param| {
|
||||
/// Box::pin(async move {
|
||||
/// // Générer une variante si param est numérique
|
||||
/// if let Ok(size) = param.parse::<usize>() {
|
||||
/// // Générer et retourner les données
|
||||
/// Some(vec![])
|
||||
/// } else {
|
||||
/// None
|
||||
/// }
|
||||
/// })
|
||||
/// });
|
||||
///
|
||||
/// let router = create_file_router_with_generator(
|
||||
/// cache.clone(),
|
||||
/// "image/webp",
|
||||
/// Some(generator)
|
||||
/// );
|
||||
/// # }
|
||||
/// ```
|
||||
#[cfg(feature = "pmoserver")]
|
||||
pub fn create_file_router_with_generator<C: CacheConfig + 'static>(
|
||||
cache: Arc<Cache<C>>,
|
||||
content_type: &'static str,
|
||||
param_generator: Option<ParamGenerator<C>>,
|
||||
) -> Router {
|
||||
let cache_name = C::cache_name();
|
||||
let cache_type = C::cache_type();
|
||||
|
||||
let path_base = format!("/{}/{}", cache_name, cache_type);
|
||||
let path_with_param = format!("/{}/{}/:pk/:param", cache_name, cache_type);
|
||||
let path_without_param = format!("/{}/{}/:pk", cache_name, cache_type);
|
||||
let path_with_param = format!("/{}/{}/{{pk}}/{{param}}", cache_name, cache_type);
|
||||
let path_without_param = format!("/{}/{}/{{pk}}", cache_name, cache_type);
|
||||
|
||||
Router::new()
|
||||
.route(&path_without_param, get(get_file::<C>))
|
||||
.route(&path_with_param, get(get_file_with_param::<C>))
|
||||
.with_state((cache, content_type))
|
||||
.with_state((cache, content_type, param_generator))
|
||||
}
|
||||
|
||||
/// Crée un router pour l'API REST du cache
|
||||
@@ -261,11 +345,11 @@ pub fn create_api_router<C: CacheConfig + 'static>(
|
||||
.delete(api::purge_cache::<C>),
|
||||
)
|
||||
.route(
|
||||
"/:pk",
|
||||
"/{pk}",
|
||||
get(api::get_item_info::<C>)
|
||||
.delete(api::delete_item::<C>),
|
||||
)
|
||||
.route("/:pk/status", get(api::get_download_status::<C>))
|
||||
.route("/{pk}/status", get(api::get_download_status::<C>))
|
||||
.route("/consolidate", post(api::consolidate_cache::<C>))
|
||||
.with_state(cache)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user