Ajout de fonctionnalité de download asynchrone au pmocache
This commit is contained in:
317
pmocache/src/api.rs
Normal file
317
pmocache/src/api.rs
Normal file
@@ -0,0 +1,317 @@
|
||||
//! API REST générique pour la gestion du cache
|
||||
//!
|
||||
//! Ce module expose une API REST documentée avec OpenAPI/Swagger pour :
|
||||
//! - Lister les items en cache
|
||||
//! - Ajouter des items depuis une URL
|
||||
//! - Consulter le status des downloads en cours
|
||||
//! - Supprimer des items
|
||||
//! - Purger et consolider le cache
|
||||
|
||||
use crate::{Cache, CacheConfig, CacheEntry};
|
||||
use axum::{
|
||||
extract::{Path, State},
|
||||
http::StatusCode,
|
||||
response::IntoResponse,
|
||||
Json,
|
||||
};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::sync::Arc;
|
||||
|
||||
#[cfg(feature = "openapi")]
|
||||
use utoipa::ToSchema;
|
||||
|
||||
/// Statut d'un téléchargement
|
||||
#[derive(Debug, Serialize, Deserialize)]
|
||||
#[cfg_attr(feature = "openapi", derive(ToSchema))]
|
||||
pub struct DownloadStatus {
|
||||
/// Clé primaire de l'item
|
||||
#[cfg_attr(feature = "openapi", schema(example = "1a2b3c4d5e6f7a8b"))]
|
||||
pub pk: String,
|
||||
/// Téléchargement en cours
|
||||
pub in_progress: bool,
|
||||
/// Taille actuelle téléchargée (source)
|
||||
pub current_size: Option<u64>,
|
||||
/// Taille après transformation
|
||||
pub transformed_size: Option<u64>,
|
||||
/// Taille totale attendue
|
||||
pub expected_size: Option<u64>,
|
||||
/// Téléchargement terminé
|
||||
pub finished: bool,
|
||||
/// Erreur éventuelle
|
||||
pub error: Option<String>,
|
||||
}
|
||||
|
||||
/// Requête pour ajouter un item au cache
|
||||
#[derive(Debug, Serialize, Deserialize)]
|
||||
#[cfg_attr(feature = "openapi", derive(ToSchema))]
|
||||
pub struct AddItemRequest {
|
||||
/// URL de la source
|
||||
#[cfg_attr(feature = "openapi", schema(example = "https://example.com/file.dat"))]
|
||||
pub url: String,
|
||||
/// Collection optionnelle
|
||||
#[cfg_attr(feature = "openapi", schema(example = "album:the_wall"))]
|
||||
pub collection: Option<String>,
|
||||
}
|
||||
|
||||
/// Réponse après ajout d'un item
|
||||
#[derive(Debug, Serialize, Deserialize)]
|
||||
#[cfg_attr(feature = "openapi", derive(ToSchema))]
|
||||
pub struct AddItemResponse {
|
||||
/// Clé primaire (pk) de l'item ajouté
|
||||
#[cfg_attr(feature = "openapi", schema(example = "1a2b3c4d5e6f7a8b"))]
|
||||
pub pk: String,
|
||||
/// URL source de l'item
|
||||
#[cfg_attr(feature = "openapi", schema(example = "https://example.com/file.dat"))]
|
||||
pub url: String,
|
||||
/// Message de succès
|
||||
#[cfg_attr(feature = "openapi", schema(example = "Item added successfully"))]
|
||||
pub message: String,
|
||||
}
|
||||
|
||||
/// Réponse de suppression d'un item
|
||||
#[derive(Debug, Serialize, Deserialize)]
|
||||
#[cfg_attr(feature = "openapi", derive(ToSchema))]
|
||||
pub struct DeleteItemResponse {
|
||||
/// Message de succès
|
||||
#[cfg_attr(feature = "openapi", schema(example = "Item deleted successfully"))]
|
||||
pub message: String,
|
||||
}
|
||||
|
||||
/// Réponse d'erreur générique
|
||||
#[derive(Debug, Serialize, Deserialize)]
|
||||
#[cfg_attr(feature = "openapi", derive(ToSchema))]
|
||||
pub struct ErrorResponse {
|
||||
/// Code d'erreur
|
||||
#[cfg_attr(feature = "openapi", schema(example = "NOT_FOUND"))]
|
||||
pub error: String,
|
||||
/// Message descriptif
|
||||
#[cfg_attr(feature = "openapi", schema(example = "Item not found in cache"))]
|
||||
pub message: String,
|
||||
}
|
||||
|
||||
/// Liste tous les items en cache avec leurs statistiques
|
||||
///
|
||||
/// Retourne la liste complète des entrées du cache triées par nombre d'accès décroissant.
|
||||
pub async fn list_items<C: CacheConfig>(
|
||||
State(cache): State<Arc<Cache<C>>>,
|
||||
) -> impl IntoResponse {
|
||||
match cache.db.get_all() {
|
||||
Ok(entries) => (StatusCode::OK, Json(entries)).into_response(),
|
||||
Err(e) => (
|
||||
StatusCode::INTERNAL_SERVER_ERROR,
|
||||
Json(ErrorResponse {
|
||||
error: "DATABASE_ERROR".to_string(),
|
||||
message: format!("Cannot retrieve cache entries: {}", e),
|
||||
}),
|
||||
)
|
||||
.into_response(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Récupère les informations d'un item spécifique
|
||||
///
|
||||
/// Retourne les métadonnées d'un item identifié par sa clé (pk).
|
||||
pub async fn get_item_info<C: CacheConfig>(
|
||||
State(cache): State<Arc<Cache<C>>>,
|
||||
Path(pk): Path<String>,
|
||||
) -> impl IntoResponse {
|
||||
match cache.db.get(&pk) {
|
||||
Ok(entry) => (StatusCode::OK, Json(entry)).into_response(),
|
||||
Err(_) => (
|
||||
StatusCode::NOT_FOUND,
|
||||
Json(ErrorResponse {
|
||||
error: "NOT_FOUND".to_string(),
|
||||
message: format!("Item with pk '{}' not found in cache", pk),
|
||||
}),
|
||||
)
|
||||
.into_response(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Récupère le statut du téléchargement d'un item
|
||||
///
|
||||
/// Retourne le statut actuel du téléchargement (progression, tailles, erreurs).
|
||||
/// Si le téléchargement est terminé, retourne les informations du fichier.
|
||||
pub async fn get_download_status<C: CacheConfig>(
|
||||
State(cache): State<Arc<Cache<C>>>,
|
||||
Path(pk): Path<String>,
|
||||
) -> impl IntoResponse {
|
||||
// Vérifier que l'item existe dans la DB
|
||||
if cache.db.get(&pk).is_err() {
|
||||
return (
|
||||
StatusCode::NOT_FOUND,
|
||||
Json(ErrorResponse {
|
||||
error: "NOT_FOUND".to_string(),
|
||||
message: format!("Item with pk '{}' not found in cache", pk),
|
||||
}),
|
||||
)
|
||||
.into_response();
|
||||
}
|
||||
|
||||
let in_progress = cache.get_download(&pk).await.is_some();
|
||||
let current_size = cache.current_size(&pk).await;
|
||||
let transformed_size = cache.transformed_size(&pk).await;
|
||||
let expected_size = cache.expected_size(&pk).await;
|
||||
let finished = cache.is_finished(&pk).await;
|
||||
|
||||
let error = if let Some(download) = cache.get_download(&pk).await {
|
||||
download.error().await
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
let status = DownloadStatus {
|
||||
pk,
|
||||
in_progress,
|
||||
current_size,
|
||||
transformed_size,
|
||||
expected_size,
|
||||
finished,
|
||||
error,
|
||||
};
|
||||
|
||||
(StatusCode::OK, Json(status)).into_response()
|
||||
}
|
||||
|
||||
/// Ajoute un item au cache depuis une URL
|
||||
///
|
||||
/// Télécharge l'item depuis l'URL fournie et l'ajoute au cache.
|
||||
/// Si l'item existe déjà, il est mis à jour.
|
||||
pub async fn add_item<C: CacheConfig>(
|
||||
State(cache): State<Arc<Cache<C>>>,
|
||||
Json(req): Json<AddItemRequest>,
|
||||
) -> impl IntoResponse {
|
||||
if req.url.is_empty() {
|
||||
return (
|
||||
StatusCode::BAD_REQUEST,
|
||||
Json(ErrorResponse {
|
||||
error: "INVALID_REQUEST".to_string(),
|
||||
message: "URL cannot be empty".to_string(),
|
||||
}),
|
||||
)
|
||||
.into_response();
|
||||
}
|
||||
|
||||
match cache.add_from_url(&req.url, req.collection.as_deref()).await {
|
||||
Ok(pk) => (
|
||||
StatusCode::CREATED,
|
||||
Json(AddItemResponse {
|
||||
pk,
|
||||
url: req.url,
|
||||
message: "Item added successfully".to_string(),
|
||||
}),
|
||||
)
|
||||
.into_response(),
|
||||
Err(e) => (
|
||||
StatusCode::INTERNAL_SERVER_ERROR,
|
||||
Json(ErrorResponse {
|
||||
error: "PROCESSING_ERROR".to_string(),
|
||||
message: format!("Cannot add item: {}", e),
|
||||
}),
|
||||
)
|
||||
.into_response(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Supprime un item du cache
|
||||
///
|
||||
/// Supprime l'item et toutes ses variantes du disque et de la base de données.
|
||||
pub async fn delete_item<C: CacheConfig>(
|
||||
State(cache): State<Arc<Cache<C>>>,
|
||||
Path(pk): Path<String>,
|
||||
) -> impl IntoResponse {
|
||||
// Vérifier que l'item existe
|
||||
if cache.db.get(&pk).is_err() {
|
||||
return (
|
||||
StatusCode::NOT_FOUND,
|
||||
Json(ErrorResponse {
|
||||
error: "NOT_FOUND".to_string(),
|
||||
message: format!("Item with pk '{}' not found in cache", pk),
|
||||
}),
|
||||
)
|
||||
.into_response();
|
||||
}
|
||||
|
||||
// Supprimer tous les fichiers avec ce pk (toutes variantes)
|
||||
let cache_dir = cache.cache_dir();
|
||||
if let Ok(mut entries) = tokio::fs::read_dir(cache_dir).await {
|
||||
while let Ok(Some(entry)) = entries.next_entry().await {
|
||||
if let Some(filename) = entry.file_name().to_str() {
|
||||
// Format: {pk}.{param}.{ext}
|
||||
if filename.starts_with(&pk) && filename.starts_with(&format!("{}.", pk)) {
|
||||
let _ = tokio::fs::remove_file(entry.path()).await;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Supprimer de la base de données
|
||||
match cache.db.delete(&pk) {
|
||||
Ok(_) => (
|
||||
StatusCode::OK,
|
||||
Json(DeleteItemResponse {
|
||||
message: format!("Item '{}' deleted successfully", pk),
|
||||
}),
|
||||
)
|
||||
.into_response(),
|
||||
Err(e) => (
|
||||
StatusCode::INTERNAL_SERVER_ERROR,
|
||||
Json(ErrorResponse {
|
||||
error: "DATABASE_ERROR".to_string(),
|
||||
message: format!("Cannot delete from database: {}", e),
|
||||
}),
|
||||
)
|
||||
.into_response(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Purge complètement le cache
|
||||
///
|
||||
/// Supprime tous les items et vide la base de données. Opération irréversible.
|
||||
pub async fn purge_cache<C: CacheConfig>(
|
||||
State(cache): State<Arc<Cache<C>>>,
|
||||
) -> impl IntoResponse {
|
||||
match cache.purge().await {
|
||||
Ok(_) => (
|
||||
StatusCode::OK,
|
||||
Json(DeleteItemResponse {
|
||||
message: "Cache purged successfully".to_string(),
|
||||
}),
|
||||
)
|
||||
.into_response(),
|
||||
Err(e) => (
|
||||
StatusCode::INTERNAL_SERVER_ERROR,
|
||||
Json(ErrorResponse {
|
||||
error: "PURGE_ERROR".to_string(),
|
||||
message: format!("Cannot purge cache: {}", e),
|
||||
}),
|
||||
)
|
||||
.into_response(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Consolide le cache
|
||||
///
|
||||
/// Re-télécharge les items manquants et supprime les fichiers orphelins.
|
||||
/// Utile pour réparer un cache corrompu.
|
||||
pub async fn consolidate_cache<C: CacheConfig>(
|
||||
State(cache): State<Arc<Cache<C>>>,
|
||||
) -> impl IntoResponse {
|
||||
match cache.consolidate().await {
|
||||
Ok(_) => (
|
||||
StatusCode::OK,
|
||||
Json(DeleteItemResponse {
|
||||
message: "Cache consolidated successfully".to_string(),
|
||||
}),
|
||||
)
|
||||
.into_response(),
|
||||
Err(e) => (
|
||||
StatusCode::INTERNAL_SERVER_ERROR,
|
||||
Json(ErrorResponse {
|
||||
error: "CONSOLIDATE_ERROR".to_string(),
|
||||
message: format!("Cannot consolidate cache: {}", e),
|
||||
}),
|
||||
)
|
||||
.into_response(),
|
||||
}
|
||||
}
|
||||
@@ -3,12 +3,14 @@
|
||||
//! Ce module fournit une interface générique pour gérer un cache de fichiers
|
||||
//! avec métadonnées dans une base de données SQLite.
|
||||
|
||||
use crate::cache_trait::FileCache;
|
||||
use crate::cache_trait::{FileCache, pk_from_url};
|
||||
use crate::db::DB;
|
||||
use crate::download::{Download, download_with_transformer, StreamTransformer};
|
||||
use anyhow::{anyhow, Result};
|
||||
use sha1::{Digest, Sha1};
|
||||
use std::collections::HashMap;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::sync::Arc;
|
||||
use tokio::sync::RwLock;
|
||||
|
||||
/// Trait pour définir les paramètres du cache
|
||||
pub trait CacheConfig: Send + Sync {
|
||||
@@ -42,7 +44,8 @@ pub trait CacheConfig: Send + Sync {
|
||||
/// * `C` - Configuration du cache (implémente `CacheConfig`)
|
||||
///
|
||||
/// 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.
|
||||
/// 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
|
||||
@@ -53,6 +56,8 @@ pub struct Cache<C: CacheConfig> {
|
||||
base_url: String,
|
||||
/// Base de données SQLite
|
||||
pub db: Arc<DB>,
|
||||
/// Map des downloads en cours (pk -> Download)
|
||||
downloads: Arc<RwLock<HashMap<String, Arc<Download>>>>,
|
||||
/// Phantom data pour le type de configuration
|
||||
_phantom: std::marker::PhantomData<C>,
|
||||
}
|
||||
@@ -75,12 +80,24 @@ impl<C: CacheConfig> Cache<C> {
|
||||
limit,
|
||||
base_url: base_url.to_string(),
|
||||
db: Arc::new(db),
|
||||
downloads: Arc::new(RwLock::new(HashMap::new())),
|
||||
_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.
|
||||
/// Le download est tracké dans la map jusqu'à sa fin.
|
||||
///
|
||||
/// # Arguments
|
||||
///
|
||||
/// * `url` - URL du fichier à télécharger
|
||||
@@ -90,13 +107,63 @@ impl<C: CacheConfig> Cache<C> {
|
||||
///
|
||||
/// La clé primaire (pk) du fichier dans le cache
|
||||
pub async fn add_from_url(&self, url: &str, collection: Option<&str>) -> Result<String> {
|
||||
let response = reqwest::get(url).await?;
|
||||
if !response.status().is_success() {
|
||||
return Err(anyhow!("Bad status: {}", response.status()));
|
||||
let pk = pk_from_url(url);
|
||||
let file_path = self.file_path(&pk);
|
||||
|
||||
// Vérifier si déjà en cours de téléchargement
|
||||
{
|
||||
let downloads = self.downloads.read().await;
|
||||
if downloads.contains_key(&pk) {
|
||||
// Download déjà en cours, retourner la clé
|
||||
return Ok(pk);
|
||||
}
|
||||
}
|
||||
|
||||
let data = response.bytes().await?;
|
||||
self.add(url, &data, collection).await
|
||||
// Lancer le téléchargement avec transformer
|
||||
let download = download_with_transformer(
|
||||
&file_path,
|
||||
url,
|
||||
self.get_transformer(),
|
||||
);
|
||||
|
||||
// Stocker dans la map des downloads en cours
|
||||
{
|
||||
let mut downloads = self.downloads.write().await;
|
||||
downloads.insert(pk.clone(), download.clone());
|
||||
}
|
||||
|
||||
// Ajouter immédiatement à la DB
|
||||
self.db.add(&pk, url, collection)?;
|
||||
|
||||
// Lancer une tâche de nettoyage en background
|
||||
let downloads_clone = self.downloads.clone();
|
||||
let pk_clone = pk.clone();
|
||||
tokio::spawn(async move {
|
||||
// Attendre la fin du téléchargement
|
||||
let _ = download.wait_until_finished().await;
|
||||
// Retirer de la map
|
||||
downloads_clone.write().await.remove(&pk_clone);
|
||||
});
|
||||
|
||||
Ok(pk)
|
||||
}
|
||||
|
||||
/// Ajoute un fichier local au cache
|
||||
///
|
||||
/// Le fichier est copié dans le cache via une URL file://
|
||||
///
|
||||
/// # Arguments
|
||||
///
|
||||
/// * `path` - Chemin du fichier local
|
||||
/// * `collection` - Collection optionnelle à laquelle appartient le fichier
|
||||
///
|
||||
/// # Returns
|
||||
///
|
||||
/// La clé primaire (pk) du fichier dans le cache
|
||||
pub async fn add_from_file(&self, path: &str, collection: Option<&str>) -> Result<String> {
|
||||
let canonical_path = std::fs::canonicalize(path)?;
|
||||
let file_url = format!("file://{}", canonical_path.display());
|
||||
self.add_from_url(&file_url, collection).await
|
||||
}
|
||||
|
||||
/// S'assure qu'un fichier est présent dans le cache
|
||||
@@ -180,13 +247,11 @@ impl<C: CacheConfig> Cache<C> {
|
||||
for entry in entries {
|
||||
let file_path = self.file_path(&entry.pk);
|
||||
if !file_path.exists() {
|
||||
match reqwest::get(&entry.source_url).await {
|
||||
Ok(response) if response.status().is_success() => {
|
||||
let data = response.bytes().await?;
|
||||
self.add(&entry.source_url, &data, entry.collection.as_deref())
|
||||
.await?;
|
||||
}
|
||||
_ => {
|
||||
// Re-télécharger le fichier manquant
|
||||
match self.add_from_url(&entry.source_url, entry.collection.as_deref()).await {
|
||||
Ok(_) => {},
|
||||
Err(_) => {
|
||||
// Si le téléchargement échoue, supprimer l'entrée DB
|
||||
self.db.delete(&entry.pk)?;
|
||||
}
|
||||
}
|
||||
@@ -213,6 +278,131 @@ impl<C: CacheConfig> Cache<C> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Récupère l'objet Download pour un pk donné (si en cours)
|
||||
///
|
||||
/// # Arguments
|
||||
///
|
||||
/// * `pk` - Clé primaire du fichier
|
||||
///
|
||||
/// # Returns
|
||||
///
|
||||
/// Some(Download) si le téléchargement est en cours, None sinon
|
||||
pub async fn get_download(&self, pk: &str) -> Option<Arc<Download>> {
|
||||
let downloads = self.downloads.read().await;
|
||||
downloads.get(pk).cloned()
|
||||
}
|
||||
|
||||
/// Retourne la taille actuelle téléchargée (source)
|
||||
///
|
||||
/// Si le download est en cours, retourne la taille téléchargée.
|
||||
/// Sinon, retourne la taille du fichier sur disque.
|
||||
///
|
||||
/// # Arguments
|
||||
///
|
||||
/// * `pk` - Clé primaire du fichier
|
||||
pub async fn current_size(&self, pk: &str) -> Option<u64> {
|
||||
if let Some(download) = self.get_download(pk).await {
|
||||
Some(download.current_size().await)
|
||||
} else {
|
||||
// Fichier terminé, lire la taille du fichier
|
||||
let file_path = self.file_path(pk);
|
||||
if file_path.exists() {
|
||||
std::fs::metadata(file_path).ok().map(|m| m.len())
|
||||
} else {
|
||||
None
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Retourne la taille des données transformées
|
||||
///
|
||||
/// Si le download est en cours, retourne la taille transformée.
|
||||
/// Sinon, retourne la taille du fichier sur disque.
|
||||
///
|
||||
/// # Arguments
|
||||
///
|
||||
/// * `pk` - Clé primaire du fichier
|
||||
pub async fn transformed_size(&self, pk: &str) -> Option<u64> {
|
||||
if let Some(download) = self.get_download(pk).await {
|
||||
Some(download.transformed_size().await)
|
||||
} else {
|
||||
// Fichier terminé, lire la taille du fichier
|
||||
let file_path = self.file_path(pk);
|
||||
if file_path.exists() {
|
||||
std::fs::metadata(file_path).ok().map(|m| m.len())
|
||||
} else {
|
||||
None
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Retourne la taille attendue du fichier (si disponible)
|
||||
///
|
||||
/// # Arguments
|
||||
///
|
||||
/// * `pk` - Clé primaire du fichier
|
||||
pub async fn expected_size(&self, pk: &str) -> Option<u64> {
|
||||
if let Some(download) = self.get_download(pk).await {
|
||||
download.expected_size().await
|
||||
} else {
|
||||
// Fichier terminé, la taille finale est la taille du fichier
|
||||
self.transformed_size(pk).await
|
||||
}
|
||||
}
|
||||
|
||||
/// Indique si le téléchargement est terminé
|
||||
///
|
||||
/// # Arguments
|
||||
///
|
||||
/// * `pk` - Clé primaire du fichier
|
||||
pub async fn is_finished(&self, pk: &str) -> bool {
|
||||
if let Some(download) = self.get_download(pk).await {
|
||||
download.finished().await
|
||||
} else {
|
||||
// Pas dans la map = terminé (ou n'existe pas)
|
||||
self.file_path(pk).exists()
|
||||
}
|
||||
}
|
||||
|
||||
/// Attend qu'un fichier atteigne au moins une taille minimale
|
||||
///
|
||||
/// # Arguments
|
||||
///
|
||||
/// * `pk` - Clé primaire du fichier
|
||||
/// * `min_size` - Taille minimale attendue en bytes
|
||||
pub async fn wait_until_min_size(&self, pk: &str, min_size: u64) -> Result<()> {
|
||||
if let Some(download) = self.get_download(pk).await {
|
||||
download.wait_until_min_size(min_size).await
|
||||
.map_err(|e| anyhow!("Download error: {}", e))
|
||||
} else {
|
||||
// Déjà terminé ou n'existe pas
|
||||
if self.file_path(pk).exists() {
|
||||
Ok(())
|
||||
} else {
|
||||
Err(anyhow!("File not found"))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Attend que le téléchargement soit complètement terminé
|
||||
///
|
||||
/// # Arguments
|
||||
///
|
||||
/// * `pk` - Clé primaire du fichier
|
||||
pub async fn wait_until_finished(&self, pk: &str) -> Result<()> {
|
||||
if let Some(download) = self.get_download(pk).await {
|
||||
download.wait_until_finished().await
|
||||
.map_err(|e| anyhow!("Download error: {}", e))
|
||||
} else {
|
||||
// Déjà terminé ou n'existe pas
|
||||
if self.file_path(pk).exists() {
|
||||
Ok(())
|
||||
} else {
|
||||
Err(anyhow!("File not found"))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Retourne le répertoire du cache
|
||||
pub fn cache_dir(&self) -> &Path {
|
||||
&self.dir
|
||||
@@ -232,9 +422,17 @@ impl<C: CacheConfig> Cache<C> {
|
||||
|
||||
|
||||
/// Implémentation du trait FileCache pour Cache
|
||||
impl<C: CacheConfig> FileCache for Cache<C> {
|
||||
fn cache_type(&self) -> &str {
|
||||
C::cache_type()
|
||||
impl<C: CacheConfig> FileCache<C> for Cache<C> {
|
||||
fn get_cache_dir(&self) -> &Path {
|
||||
self.cache_dir()
|
||||
}
|
||||
|
||||
fn get_database(&self) -> Arc<DB> {
|
||||
self.db.clone()
|
||||
}
|
||||
|
||||
fn get_base_url(&self) -> &str {
|
||||
&self.base_url
|
||||
}
|
||||
|
||||
fn validate_data(&self, data: &[u8]) -> Result<Vec<u8>> {
|
||||
@@ -246,23 +444,12 @@ impl<C: CacheConfig> FileCache for Cache<C> {
|
||||
self.add_from_url(url, collection).await
|
||||
}
|
||||
|
||||
async fn ensure_from_url(&self, url: &str, collection: Option<&str>) -> Result<String> {
|
||||
self.ensure_from_url(url, collection).await
|
||||
async fn add_from_file(&self, path: &str, collection: Option<&str>) -> Result<String> {
|
||||
self.add_from_file(path, collection).await
|
||||
}
|
||||
|
||||
async fn add(&self, url: &str, data: &[u8], collection: Option<&str>) -> Result<String> {
|
||||
// Valider les données avant de les ajouter
|
||||
let validated_data = self.validate_data(data)?;
|
||||
|
||||
let pk = pk_from_url(url);
|
||||
let file_path = self.file_path(&pk);
|
||||
|
||||
if !file_path.exists() {
|
||||
tokio::fs::write(&file_path, &validated_data).await?;
|
||||
}
|
||||
|
||||
self.db.add(&pk, url, collection)?;
|
||||
Ok(pk)
|
||||
async fn ensure_from_url(&self, url: &str, collection: Option<&str>) -> Result<String> {
|
||||
self.ensure_from_url(url, collection).await
|
||||
}
|
||||
|
||||
async fn get(&self, pk: &str) -> Result<PathBuf> {
|
||||
@@ -280,12 +467,4 @@ impl<C: CacheConfig> FileCache for Cache<C> {
|
||||
async fn consolidate(&self) -> Result<()> {
|
||||
self.consolidate().await
|
||||
}
|
||||
|
||||
fn get_cache_dir(&self) -> String {
|
||||
self.cache_dir()
|
||||
}
|
||||
|
||||
fn get_base_url(&self) -> &str {
|
||||
self.get_base_url()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -12,6 +12,7 @@ pub trait FileCache<C: CacheConfig>: Send + Sync {
|
||||
fn get_cache_dir(&self) -> &Path;
|
||||
fn get_database(&self) -> Arc<DB>;
|
||||
fn get_base_url(&self) -> &str;
|
||||
|
||||
/// Valide les données avant de les stocker dans le cache
|
||||
///
|
||||
/// Cette méthode peut être surchargée pour vérifier le type MIME,
|
||||
@@ -34,12 +35,12 @@ pub trait FileCache<C: CacheConfig>: Send + Sync {
|
||||
C::cache_type()
|
||||
}
|
||||
|
||||
/// Retourne le type de cache
|
||||
/// Retourne le nom du cache
|
||||
fn cache_name(&self) -> &'static str {
|
||||
C::cache_name()
|
||||
}
|
||||
|
||||
/// Retourne le type de cache
|
||||
/// Retourne le paramètre par défaut
|
||||
fn default_param(&self) -> &'static str {
|
||||
C::default_param()
|
||||
}
|
||||
@@ -54,9 +55,7 @@ pub trait FileCache<C: CacheConfig>: Send + Sync {
|
||||
C::table_name()
|
||||
}
|
||||
|
||||
|
||||
|
||||
/// Construit le chemin complet d'un fichier dans le cache
|
||||
/// Construit le chemin complet d'un fichier dans le cache
|
||||
///
|
||||
/// Format: `{pk}.{qualificatif}.{extension}`
|
||||
/// Pour le fichier original: `{pk}.orig.{extension}`
|
||||
@@ -83,6 +82,20 @@ pub trait FileCache<C: CacheConfig>: Send + Sync {
|
||||
/// La clé primaire (pk) du fichier dans le cache
|
||||
async fn add_from_url(&self, url: &str, collection: Option<&str>) -> Result<String>;
|
||||
|
||||
/// Ajoute un fichier local au cache
|
||||
///
|
||||
/// Le fichier est copié dans le cache via une URL file://
|
||||
///
|
||||
/// # Arguments
|
||||
///
|
||||
/// * `path` - Chemin du fichier local
|
||||
/// * `collection` - Collection optionnelle à laquelle appartient le fichier
|
||||
///
|
||||
/// # Returns
|
||||
///
|
||||
/// La clé primaire (pk) du fichier dans le cache
|
||||
async fn add_from_file(&self, path: &str, collection: Option<&str>) -> Result<String>;
|
||||
|
||||
/// S'assure qu'un fichier est présent dans le cache
|
||||
///
|
||||
/// Si le fichier existe déjà, retourne sa clé. Sinon, le télécharge.
|
||||
|
||||
437
pmocache/src/download.rs
Normal file
437
pmocache/src/download.rs
Normal file
@@ -0,0 +1,437 @@
|
||||
use std::fs::File;
|
||||
use std::io;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::pin::Pin;
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
use tokio::sync::RwLock;
|
||||
use futures_util::Future;
|
||||
|
||||
/// Type pour une fonction de transformation de stream
|
||||
///
|
||||
/// La fonction reçoit:
|
||||
/// - Le stream de bytes téléchargés
|
||||
/// - Un writer pour écrire les données transformées
|
||||
/// - Un callback pour mettre à jour la progression
|
||||
///
|
||||
/// Elle retourne un Future qui se résout en Result
|
||||
pub type StreamTransformer = Box<
|
||||
dyn FnOnce(
|
||||
reqwest::Response,
|
||||
tokio::fs::File,
|
||||
Arc<dyn Fn(u64) + Send + Sync>,
|
||||
) -> Pin<Box<dyn Future<Output = Result<(), String>> + Send>>
|
||||
+ Send,
|
||||
>;
|
||||
|
||||
/// État interne du téléchargement
|
||||
#[derive(Debug, Clone)]
|
||||
struct DownloadState {
|
||||
/// Taille actuelle téléchargée (du stream source)
|
||||
current_size: u64,
|
||||
/// Taille attendue du fichier source (si connue)
|
||||
expected_size: Option<u64>,
|
||||
/// Taille des données transformées écrites
|
||||
transformed_size: u64,
|
||||
/// Indique si le téléchargement est terminé
|
||||
finished: bool,
|
||||
/// Position de lecture actuelle
|
||||
read_position: u64,
|
||||
/// Erreur éventuelle lors du téléchargement
|
||||
error: Option<String>,
|
||||
}
|
||||
|
||||
/// Objet représentant un téléchargement en cours
|
||||
#[derive(Debug)]
|
||||
pub struct Download {
|
||||
/// Nom du fichier de destination
|
||||
filename: PathBuf,
|
||||
/// État partagé entre le téléchargement et les lectures
|
||||
state: Arc<RwLock<DownloadState>>,
|
||||
}
|
||||
|
||||
impl Download {
|
||||
/// Crée une nouvelle instance de Download
|
||||
fn new(filename: PathBuf) -> Arc<Self> {
|
||||
Arc::new(Self {
|
||||
filename,
|
||||
state: Arc::new(RwLock::new(DownloadState {
|
||||
current_size: 0,
|
||||
expected_size: None,
|
||||
transformed_size: 0,
|
||||
finished: false,
|
||||
read_position: 0,
|
||||
error: None,
|
||||
})),
|
||||
})
|
||||
}
|
||||
|
||||
/// Retourne le nom du fichier
|
||||
pub fn filename(&self) -> &Path {
|
||||
&self.filename
|
||||
}
|
||||
|
||||
/// Attend que le fichier atteigne au moins la taille spécifiée ou soit complètement téléchargé
|
||||
pub async fn wait_until_min_size(&self, min_size: u64) -> Result<(), String> {
|
||||
loop {
|
||||
let state = self.state.read().await;
|
||||
|
||||
// Vérifier s'il y a eu une erreur
|
||||
if let Some(ref error) = state.error {
|
||||
return Err(error.clone());
|
||||
}
|
||||
|
||||
// Vérifier si la condition est remplie
|
||||
if state.transformed_size >= min_size || state.finished {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
drop(state); // Libérer le lock avant de dormir
|
||||
tokio::time::sleep(Duration::from_millis(50)).await;
|
||||
}
|
||||
}
|
||||
|
||||
/// Attend que le téléchargement soit complètement terminé
|
||||
pub async fn wait_until_finished(&self) -> Result<(), String> {
|
||||
loop {
|
||||
let state = self.state.read().await;
|
||||
|
||||
// Vérifier s'il y a eu une erreur
|
||||
if let Some(ref error) = state.error {
|
||||
return Err(error.clone());
|
||||
}
|
||||
|
||||
if state.finished {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
drop(state);
|
||||
tokio::time::sleep(Duration::from_millis(50)).await;
|
||||
}
|
||||
}
|
||||
|
||||
/// Ouvre le fichier pour lecture
|
||||
pub fn open(&self) -> io::Result<File> {
|
||||
File::open(&self.filename)
|
||||
}
|
||||
|
||||
/// Retourne la position actuelle de lecture
|
||||
pub async fn pos(&self) -> u64 {
|
||||
let state = self.state.read().await;
|
||||
state.read_position
|
||||
}
|
||||
|
||||
/// Met à jour la position de lecture
|
||||
pub async fn set_pos(&self, pos: u64) {
|
||||
let mut state = self.state.write().await;
|
||||
state.read_position = pos;
|
||||
}
|
||||
|
||||
/// Retourne la taille attendue du fichier (si disponible)
|
||||
pub async fn expected_size(&self) -> Option<u64> {
|
||||
let state = self.state.read().await;
|
||||
state.expected_size
|
||||
}
|
||||
|
||||
/// Retourne la taille actuellement téléchargée (du stream source)
|
||||
pub async fn current_size(&self) -> u64 {
|
||||
let state = self.state.read().await;
|
||||
state.current_size
|
||||
}
|
||||
|
||||
/// Retourne la taille des données transformées écrites sur disque
|
||||
pub async fn transformed_size(&self) -> u64 {
|
||||
let state = self.state.read().await;
|
||||
state.transformed_size
|
||||
}
|
||||
|
||||
/// Indique si le téléchargement est terminé
|
||||
pub async fn finished(&self) -> bool {
|
||||
let state = self.state.read().await;
|
||||
state.finished
|
||||
}
|
||||
|
||||
/// Retourne l'erreur éventuelle
|
||||
pub async fn error(&self) -> Option<String> {
|
||||
let state = self.state.read().await;
|
||||
state.error.clone()
|
||||
}
|
||||
}
|
||||
|
||||
/// Lance le téléchargement d'une URL dans un fichier
|
||||
///
|
||||
/// # Arguments
|
||||
/// * `filename` - Chemin du fichier de destination
|
||||
/// * `url` - URL à télécharger
|
||||
///
|
||||
/// # Returns
|
||||
/// Un Arc<Download> qui permet de suivre la progression du téléchargement
|
||||
pub fn download<P: AsRef<Path>>(filename: P, url: &str) -> Arc<Download> {
|
||||
download_with_transformer(filename, url, None)
|
||||
}
|
||||
|
||||
/// Lance le téléchargement d'une URL avec transformation du stream
|
||||
///
|
||||
/// # Arguments
|
||||
/// * `filename` - Chemin du fichier de destination
|
||||
/// * `url` - URL à télécharger
|
||||
/// * `transformer` - Fonction optionnelle pour transformer le stream avant sauvegarde
|
||||
///
|
||||
/// # Returns
|
||||
/// Un Arc<Download> qui permet de suivre la progression du téléchargement
|
||||
///
|
||||
/// # Exemple
|
||||
///
|
||||
/// ```rust,no_run
|
||||
/// use pmocache::download::{download_with_transformer, StreamTransformer};
|
||||
/// use futures_util::StreamExt;
|
||||
/// use tokio::io::AsyncWriteExt;
|
||||
///
|
||||
/// // Transformer qui convertit en majuscules (exemple simple)
|
||||
/// let transformer: StreamTransformer = Box::new(|response, mut file, update_progress| {
|
||||
/// Box::pin(async move {
|
||||
/// let mut stream = response.bytes_stream();
|
||||
/// let mut total = 0u64;
|
||||
///
|
||||
/// while let Some(chunk_result) = stream.next().await {
|
||||
/// let chunk = chunk_result.map_err(|e| e.to_string())?;
|
||||
///
|
||||
/// // Transformer les données (ex: conversion, décompression, etc.)
|
||||
/// let transformed = chunk.to_vec(); // Votre transformation ici
|
||||
///
|
||||
/// file.write_all(&transformed).await.map_err(|e| e.to_string())?;
|
||||
///
|
||||
/// total += chunk.len() as u64;
|
||||
/// update_progress(total);
|
||||
/// }
|
||||
///
|
||||
/// file.flush().await.map_err(|e| e.to_string())?;
|
||||
/// Ok(())
|
||||
/// })
|
||||
/// });
|
||||
///
|
||||
/// let dl = download_with_transformer("/tmp/output.txt", "https://example.com/data", Some(transformer));
|
||||
/// ```
|
||||
pub fn download_with_transformer<P: AsRef<Path>>(
|
||||
filename: P,
|
||||
url: &str,
|
||||
transformer: Option<StreamTransformer>,
|
||||
) -> Arc<Download> {
|
||||
let filename = filename.as_ref().to_path_buf();
|
||||
let url = url.to_string();
|
||||
|
||||
let download = Download::new(filename.clone());
|
||||
let state = Arc::clone(&download.state);
|
||||
|
||||
// Lancer le téléchargement en tâche de fond
|
||||
tokio::spawn(async move {
|
||||
if let Err(e) = download_impl(filename, url, state, transformer).await {
|
||||
// L'erreur a déjà été enregistrée dans download_impl
|
||||
eprintln!("Download error: {}", e);
|
||||
}
|
||||
});
|
||||
|
||||
download
|
||||
}
|
||||
|
||||
/// Implémentation du téléchargement
|
||||
async fn download_impl(
|
||||
filename: PathBuf,
|
||||
url: String,
|
||||
state: Arc<RwLock<DownloadState>>,
|
||||
transformer: Option<StreamTransformer>,
|
||||
) -> Result<(), String> {
|
||||
// Créer le client HTTP
|
||||
let client = reqwest::Client::builder()
|
||||
.timeout(Duration::from_secs(300))
|
||||
.build()
|
||||
.map_err(|e| e.to_string())?;
|
||||
|
||||
// Lancer la requête
|
||||
let response = client
|
||||
.get(&url)
|
||||
.send()
|
||||
.await
|
||||
.map_err(|e| {
|
||||
let error = format!("Failed to fetch URL: {}", e);
|
||||
tokio::task::block_in_place(|| {
|
||||
tokio::runtime::Handle::current().block_on(async {
|
||||
let mut s = state.write().await;
|
||||
s.error = Some(error.clone());
|
||||
});
|
||||
});
|
||||
error
|
||||
})?;
|
||||
|
||||
// Vérifier le statut
|
||||
if !response.status().is_success() {
|
||||
let error = format!("HTTP error: {}", response.status());
|
||||
let mut s = state.write().await;
|
||||
s.error = Some(error.clone());
|
||||
s.finished = true;
|
||||
return Err(error);
|
||||
}
|
||||
|
||||
// Récupérer la taille attendue si disponible
|
||||
if let Some(content_length) = response.content_length() {
|
||||
let mut s = state.write().await;
|
||||
s.expected_size = Some(content_length);
|
||||
}
|
||||
|
||||
// Créer le fichier de destination
|
||||
let file = tokio::fs::File::create(&filename)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
let error = format!("Failed to create file: {}", e);
|
||||
tokio::task::block_in_place(|| {
|
||||
tokio::runtime::Handle::current().block_on(async {
|
||||
let mut s = state.write().await;
|
||||
s.error = Some(error.clone());
|
||||
s.finished = true;
|
||||
});
|
||||
});
|
||||
error
|
||||
})?;
|
||||
|
||||
// Si un transformer est fourni, l'utiliser
|
||||
if let Some(transformer) = transformer {
|
||||
// Créer un callback pour mettre à jour la progression
|
||||
let state_clone = Arc::clone(&state);
|
||||
let progress_callback: Arc<dyn Fn(u64) + Send + Sync> = Arc::new(move |transformed_bytes| {
|
||||
let state = Arc::clone(&state_clone);
|
||||
tokio::spawn(async move {
|
||||
let mut s = state.write().await;
|
||||
s.transformed_size = transformed_bytes;
|
||||
});
|
||||
});
|
||||
|
||||
// Appeler le transformer
|
||||
match transformer(response, file, progress_callback).await {
|
||||
Ok(_) => {
|
||||
let mut s = state.write().await;
|
||||
s.finished = true;
|
||||
Ok(())
|
||||
}
|
||||
Err(e) => {
|
||||
let mut s = state.write().await;
|
||||
s.error = Some(e.clone());
|
||||
s.finished = true;
|
||||
Err(e)
|
||||
}
|
||||
}
|
||||
} else {
|
||||
// Comportement par défaut : téléchargement direct sans transformation
|
||||
default_download(response, file, state).await
|
||||
}
|
||||
}
|
||||
|
||||
/// Téléchargement par défaut sans transformation
|
||||
async fn default_download(
|
||||
response: reqwest::Response,
|
||||
mut file: tokio::fs::File,
|
||||
state: Arc<RwLock<DownloadState>>,
|
||||
) -> Result<(), String> {
|
||||
use tokio::io::AsyncWriteExt;
|
||||
use futures_util::StreamExt;
|
||||
|
||||
let mut stream = response.bytes_stream();
|
||||
|
||||
while let Some(chunk_result) = stream.next().await {
|
||||
match chunk_result {
|
||||
Ok(chunk) => {
|
||||
// Écrire le chunk dans le fichier
|
||||
if let Err(e) = file.write_all(&chunk).await {
|
||||
let error = format!("Failed to write to file: {}", e);
|
||||
let mut s = state.write().await;
|
||||
s.error = Some(error.clone());
|
||||
s.finished = true;
|
||||
return Err(error);
|
||||
}
|
||||
|
||||
// Mettre à jour les tailles (identiques sans transformation)
|
||||
let mut s = state.write().await;
|
||||
let chunk_len = chunk.len() as u64;
|
||||
s.current_size += chunk_len;
|
||||
s.transformed_size += chunk_len;
|
||||
}
|
||||
Err(e) => {
|
||||
let error = format!("Failed to read chunk: {}", e);
|
||||
let mut s = state.write().await;
|
||||
s.error = Some(error.clone());
|
||||
s.finished = true;
|
||||
return Err(error);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Fermer le fichier
|
||||
if let Err(e) = file.flush().await {
|
||||
let error = format!("Failed to flush file: {}", e);
|
||||
let mut s = state.write().await;
|
||||
s.error = Some(error.clone());
|
||||
s.finished = true;
|
||||
return Err(error);
|
||||
}
|
||||
|
||||
// Marquer comme terminé
|
||||
let mut s = state.write().await;
|
||||
s.finished = true;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use std::path::PathBuf;
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_download_basic() {
|
||||
let temp_dir = std::env::temp_dir();
|
||||
let filename = temp_dir.join("test_download.txt");
|
||||
|
||||
// Nettoyer si le fichier existe
|
||||
let _ = std::fs::remove_file(&filename);
|
||||
|
||||
// Télécharger un petit fichier de test
|
||||
let dl = download(&filename, "https://www.rust-lang.org/");
|
||||
|
||||
// Attendre la fin du téléchargement
|
||||
match dl.wait_until_finished().await {
|
||||
Ok(_) => {
|
||||
assert!(dl.finished().await);
|
||||
assert!(filename.exists());
|
||||
assert!(dl.current_size().await > 0);
|
||||
}
|
||||
Err(e) => {
|
||||
eprintln!("Download failed: {}", e);
|
||||
}
|
||||
}
|
||||
|
||||
// Nettoyer
|
||||
let _ = std::fs::remove_file(&filename);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_wait_until_min_size() {
|
||||
let temp_dir = std::env::temp_dir();
|
||||
let filename = temp_dir.join("test_download_min_size.txt");
|
||||
|
||||
let _ = std::fs::remove_file(&filename);
|
||||
|
||||
let dl = download(&filename, "https://www.rust-lang.org/");
|
||||
|
||||
// Attendre au moins 100 bytes
|
||||
match dl.wait_until_min_size(100).await {
|
||||
Ok(_) => {
|
||||
let size = dl.current_size().await;
|
||||
assert!(size >= 100 || dl.finished().await);
|
||||
}
|
||||
Err(e) => {
|
||||
eprintln!("Download failed: {}", e);
|
||||
}
|
||||
}
|
||||
|
||||
let _ = std::fs::remove_file(&filename);
|
||||
}
|
||||
}
|
||||
@@ -125,10 +125,27 @@
|
||||
pub mod db;
|
||||
pub mod cache;
|
||||
pub mod cache_trait;
|
||||
pub mod download;
|
||||
|
||||
#[cfg(feature = "pmoserver")]
|
||||
pub mod pmoserver_ext;
|
||||
|
||||
#[cfg(feature = "pmoserver")]
|
||||
pub mod api;
|
||||
|
||||
#[cfg(feature = "openapi")]
|
||||
pub mod openapi;
|
||||
|
||||
pub use db::{DB, CacheEntry};
|
||||
pub use cache::{Cache, CacheConfig, pk_from_url};
|
||||
pub use cache_trait::FileCache;
|
||||
pub use cache::{Cache, CacheConfig};
|
||||
pub use cache_trait::{FileCache, pk_from_url};
|
||||
pub use download::{Download, download, download_with_transformer, StreamTransformer};
|
||||
|
||||
#[cfg(feature = "pmoserver")]
|
||||
pub use pmoserver_ext::{create_file_router, create_api_router, GenericCacheExt};
|
||||
|
||||
#[cfg(all(feature = "pmoserver", feature = "openapi"))]
|
||||
pub use api::{
|
||||
DownloadStatus, AddItemRequest, AddItemResponse,
|
||||
DeleteItemResponse, ErrorResponse,
|
||||
};
|
||||
|
||||
62
pmocache/src/openapi.rs
Normal file
62
pmocache/src/openapi.rs
Normal file
@@ -0,0 +1,62 @@
|
||||
//! Génération de documentation OpenAPI pour l'API du cache générique
|
||||
//!
|
||||
//! Ce module fournit une macro pour créer dynamiquement la documentation OpenAPI
|
||||
//! selon le type de cache (images, audio, etc.).
|
||||
|
||||
/// Macro pour créer une documentation OpenAPI pour un type de cache
|
||||
///
|
||||
/// # Exemple
|
||||
///
|
||||
/// ```rust,ignore
|
||||
/// use pmocache::create_cache_openapi;
|
||||
///
|
||||
/// // Génère une struct OpenApi pour le cache de couvertures
|
||||
/// create_cache_openapi!(
|
||||
/// CoversApiDoc,
|
||||
/// "covers",
|
||||
/// "Covers",
|
||||
/// "Gestion du cache d'images de couvertures"
|
||||
/// );
|
||||
/// ```
|
||||
#[macro_export]
|
||||
macro_rules! create_cache_openapi {
|
||||
($doc_name:ident, $cache_name:expr, $cache_title:expr, $cache_description:expr) => {
|
||||
#[derive(utoipa::OpenApi)]
|
||||
#[openapi(
|
||||
paths(
|
||||
$crate::api::list_items::<Self>,
|
||||
$crate::api::get_item_info::<Self>,
|
||||
$crate::api::get_download_status::<Self>,
|
||||
$crate::api::add_item::<Self>,
|
||||
$crate::api::delete_item::<Self>,
|
||||
$crate::api::purge_cache::<Self>,
|
||||
$crate::api::consolidate_cache::<Self>,
|
||||
),
|
||||
components(
|
||||
schemas(
|
||||
$crate::db::CacheEntry,
|
||||
$crate::api::DownloadStatus,
|
||||
$crate::api::AddItemRequest,
|
||||
$crate::api::AddItemResponse,
|
||||
$crate::api::DeleteItemResponse,
|
||||
$crate::api::ErrorResponse,
|
||||
)
|
||||
),
|
||||
tags(
|
||||
(name = $cache_name, description = concat!("Gestion du cache de ", $cache_title))
|
||||
),
|
||||
info(
|
||||
title = concat!("PMO", $cache_title, " API"),
|
||||
version = "0.1.0",
|
||||
description = $cache_description,
|
||||
contact(
|
||||
name = "PMOMusic",
|
||||
),
|
||||
license(
|
||||
name = "MIT",
|
||||
),
|
||||
)
|
||||
)]
|
||||
pub struct $doc_name;
|
||||
};
|
||||
}
|
||||
@@ -1,16 +1,22 @@
|
||||
//! Extension pmoserver pour servir les fichiers du cache via HTTP
|
||||
//!
|
||||
//! Ce module fournit des handlers génériques pour servir les fichiers
|
||||
//! d'un cache via des routes HTTP structurées.
|
||||
//! d'un cache via des routes HTTP structurées, avec support du streaming progressif.
|
||||
//!
|
||||
//! ## Routes générées
|
||||
//!
|
||||
//! Format: `/{name}/{type}/{pk}[/{param}]`
|
||||
//! Format: `/{cache_name}/{cache_type}/{pk}[/{param}]`
|
||||
//!
|
||||
//! Exemples:
|
||||
//! - `/covers/images/abc123` - Image avec param par défaut (orig)
|
||||
//! - `/covers/images/abc123/thumb` - Image avec param spécifique
|
||||
//! - `/audio/tracks/def456/stream` - Piste audio
|
||||
//! - `/covers/images/abc123/256` - Image redimensionnée 256x256
|
||||
//! - `/audio/tracks/def456` - Piste audio par défaut
|
||||
//! - `/audio/tracks/def456/stream` - Piste audio streamable
|
||||
//!
|
||||
//! ## Streaming progressif
|
||||
//!
|
||||
//! Les fichiers en cours de téléchargement sont automatiquement streamés
|
||||
//! au fur et à mesure de leur disponibilité.
|
||||
//!
|
||||
//! ## Utilisation
|
||||
//!
|
||||
@@ -25,71 +31,142 @@
|
||||
//! "image/webp" // Content-Type
|
||||
//! );
|
||||
//!
|
||||
//! // Le router peut être monté sur n'importe quel chemin
|
||||
//! // Exemple: /covers/images -> GET /covers/images/{pk}
|
||||
//! // -> GET /covers/images/{pk}/{param}
|
||||
//! // Le router sera monté à la racine avec les routes complètes
|
||||
//! // Exemple: GET /covers/images/{pk}
|
||||
//! // GET /covers/images/{pk}/{param}
|
||||
//! # }
|
||||
//! ```
|
||||
|
||||
#[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},
|
||||
http::StatusCode,
|
||||
response::{IntoResponse, Response},
|
||||
routing::get,
|
||||
routing::{get, post},
|
||||
Router,
|
||||
};
|
||||
#[cfg(feature = "pmoserver")]
|
||||
use std::sync::Arc;
|
||||
#[cfg(feature = "pmoserver")]
|
||||
use tokio_util::io::ReaderStream;
|
||||
#[cfg(feature = "pmoserver")]
|
||||
use tracing::warn;
|
||||
|
||||
/// Handler générique pour GET /{pk}
|
||||
/// 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)>,
|
||||
Path(pk): Path<String>,
|
||||
) -> Response {
|
||||
match cache.get(&pk).await {
|
||||
Ok(file_path) => match tokio::fs::read(&file_path).await {
|
||||
Ok(data) => (
|
||||
StatusCode::OK,
|
||||
[("content-type", content_type)],
|
||||
data,
|
||||
)
|
||||
.into_response(),
|
||||
Err(_) => (StatusCode::NOT_FOUND, "File not found").into_response(),
|
||||
},
|
||||
Err(e) => {
|
||||
warn!("Error getting file {}: {}", pk, e);
|
||||
(StatusCode::NOT_FOUND, "Item not found").into_response()
|
||||
}
|
||||
}
|
||||
// Utiliser le param par défaut
|
||||
let param = C::default_param();
|
||||
serve_file_with_streaming(&cache, &pk, param, content_type).await
|
||||
}
|
||||
|
||||
/// Handler générique pour GET /{pk}/{param}
|
||||
/// 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)>,
|
||||
Path((pk, param)): Path<(String, String)>,
|
||||
) -> Response {
|
||||
let file_path = cache.file_path_with_qualifier(&pk, ¶m);
|
||||
serve_file_with_streaming(&cache, &pk, ¶m, content_type).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.
|
||||
#[cfg(feature = "pmoserver")]
|
||||
async fn serve_file_with_streaming<C: CacheConfig>(
|
||||
cache: &Arc<Cache<C>>,
|
||||
pk: &str,
|
||||
param: &str,
|
||||
content_type: &'static str,
|
||||
) -> Response {
|
||||
let file_path = cache.file_path_with_qualifier(pk, param);
|
||||
|
||||
// Mettre à jour les stats d'utilisation
|
||||
if let Err(e) = cache.db.update_hit(pk) {
|
||||
warn!("Error updating hit count for {}: {}", pk, e);
|
||||
}
|
||||
|
||||
// Vérifier si le download est en cours
|
||||
if let Some(download) = cache.get_download(pk).await {
|
||||
// Le fichier est en cours de téléchargement
|
||||
if !download.finished().await {
|
||||
// Streaming progressif
|
||||
return stream_file_progressive(file_path, download, content_type).await;
|
||||
}
|
||||
}
|
||||
|
||||
// Fichier terminé ou pas de download en cours, servir normalement
|
||||
serve_complete_file(file_path, content_type).await
|
||||
}
|
||||
|
||||
/// Stream un fichier en cours de téléchargement de manière progressive
|
||||
#[cfg(feature = "pmoserver")]
|
||||
async fn stream_file_progressive(
|
||||
file_path: std::path::PathBuf,
|
||||
download: Arc<crate::download::Download>,
|
||||
content_type: &'static str,
|
||||
) -> Response {
|
||||
// Attendre qu'au moins 64 KB soient disponibles avant de commencer
|
||||
const MIN_SIZE_TO_START: u64 = 64 * 1024;
|
||||
|
||||
if let Err(e) = download.wait_until_min_size(MIN_SIZE_TO_START).await {
|
||||
warn!("Error waiting for download to start: {}", e);
|
||||
if let Some(error_msg) = download.error().await {
|
||||
return (
|
||||
StatusCode::INTERNAL_SERVER_ERROR,
|
||||
format!("Download error: {}", error_msg),
|
||||
)
|
||||
.into_response();
|
||||
}
|
||||
return (StatusCode::NOT_FOUND, "File not available").into_response();
|
||||
}
|
||||
|
||||
// Ouvrir le fichier en lecture
|
||||
let file = match tokio::fs::File::open(&file_path).await {
|
||||
Ok(f) => f,
|
||||
Err(e) => {
|
||||
warn!("Error opening file {:?}: {}", file_path, e);
|
||||
return (StatusCode::NOT_FOUND, "File not found").into_response();
|
||||
}
|
||||
};
|
||||
|
||||
// Créer un stream à partir du fichier
|
||||
let stream = ReaderStream::new(file);
|
||||
let body = Body::from_stream(stream);
|
||||
|
||||
(
|
||||
StatusCode::OK,
|
||||
[
|
||||
("content-type", content_type),
|
||||
("transfer-encoding", "chunked"),
|
||||
],
|
||||
body,
|
||||
)
|
||||
.into_response()
|
||||
}
|
||||
|
||||
/// Sert un fichier complet déjà téléchargé
|
||||
#[cfg(feature = "pmoserver")]
|
||||
async fn serve_complete_file(
|
||||
file_path: std::path::PathBuf,
|
||||
content_type: &'static str,
|
||||
) -> Response {
|
||||
if !file_path.exists() {
|
||||
warn!("File not found: {:?}", file_path);
|
||||
return (StatusCode::NOT_FOUND, "File not found").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);
|
||||
}
|
||||
|
||||
match tokio::fs::read(&file_path).await {
|
||||
Ok(data) => (
|
||||
StatusCode::OK,
|
||||
@@ -97,12 +174,17 @@ async fn get_file_with_param<C: CacheConfig + 'static>(
|
||||
data,
|
||||
)
|
||||
.into_response(),
|
||||
Err(_) => (StatusCode::NOT_FOUND, "File not found").into_response(),
|
||||
Err(e) => {
|
||||
warn!("Error reading file {:?}: {}", file_path, e);
|
||||
(StatusCode::INTERNAL_SERVER_ERROR, "Error reading file").into_response()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Crée un router pour servir les fichiers d'un cache
|
||||
///
|
||||
/// Crée un router avec les routes complètes incluant cache_name et cache_type.
|
||||
///
|
||||
/// # Arguments
|
||||
///
|
||||
/// * `cache` - Instance du cache
|
||||
@@ -110,8 +192,8 @@ async fn get_file_with_param<C: CacheConfig + 'static>(
|
||||
///
|
||||
/// # Routes créées
|
||||
///
|
||||
/// - `GET /{pk}` - Fichier avec param par défaut
|
||||
/// - `GET /{pk}/{param}` - Fichier avec param spécifique
|
||||
/// - `GET /{cache_name}/{cache_type}/{pk}` - Fichier avec param par défaut
|
||||
/// - `GET /{cache_name}/{cache_type}/{pk}/{param}` - Fichier avec param spécifique
|
||||
///
|
||||
/// # Exemple
|
||||
///
|
||||
@@ -126,8 +208,10 @@ async fn get_file_with_param<C: CacheConfig + 'static>(
|
||||
/// "image/webp"
|
||||
/// );
|
||||
///
|
||||
/// // Monter le router sur /covers/images
|
||||
/// server.add_router("/covers/images", router).await;
|
||||
/// // Le router sera monté à la racine avec les routes complètes:
|
||||
/// // GET /covers/images/{pk}
|
||||
/// // GET /covers/images/{pk}/{param}
|
||||
/// server.add_router("/", router).await;
|
||||
/// # }
|
||||
/// ```
|
||||
#[cfg(feature = "pmoserver")]
|
||||
@@ -135,8 +219,79 @@ pub fn create_file_router<C: CacheConfig + 'static>(
|
||||
cache: Arc<Cache<C>>,
|
||||
content_type: &'static str,
|
||||
) -> 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);
|
||||
|
||||
Router::new()
|
||||
.route("/:pk", get(get_file::<C>))
|
||||
.route("/:pk/:param", get(get_file_with_param::<C>))
|
||||
.route(&path_without_param, get(get_file::<C>))
|
||||
.route(&path_with_param, get(get_file_with_param::<C>))
|
||||
.with_state((cache, content_type))
|
||||
}
|
||||
|
||||
/// Crée un router pour l'API REST du cache
|
||||
///
|
||||
/// # Arguments
|
||||
///
|
||||
/// * `cache` - Instance du cache
|
||||
///
|
||||
/// # Routes créées
|
||||
///
|
||||
/// - `GET /` - Liste des items
|
||||
/// - `POST /` - Ajouter un item
|
||||
/// - `DELETE /` - Purger le cache
|
||||
/// - `GET /{pk}` - Info d'un item
|
||||
/// - `GET /{pk}/status` - Status du download
|
||||
/// - `DELETE /{pk}` - Supprimer un item
|
||||
/// - `POST /consolidate` - Consolider le cache
|
||||
#[cfg(feature = "pmoserver")]
|
||||
pub fn create_api_router<C: CacheConfig + 'static>(
|
||||
cache: Arc<Cache<C>>,
|
||||
) -> Router {
|
||||
use crate::api;
|
||||
|
||||
Router::new()
|
||||
.route(
|
||||
"/",
|
||||
get(api::list_items::<C>)
|
||||
.post(api::add_item::<C>)
|
||||
.delete(api::purge_cache::<C>),
|
||||
)
|
||||
.route(
|
||||
"/:pk",
|
||||
get(api::get_item_info::<C>)
|
||||
.delete(api::delete_item::<C>),
|
||||
)
|
||||
.route("/:pk/status", get(api::get_download_status::<C>))
|
||||
.route("/consolidate", post(api::consolidate_cache::<C>))
|
||||
.with_state(cache)
|
||||
}
|
||||
|
||||
/// Trait d'extension pour pmoserver::Server
|
||||
///
|
||||
/// Permet d'initialiser un cache générique avec routes HTTP complètes
|
||||
#[cfg(feature = "pmoserver")]
|
||||
pub trait GenericCacheExt {
|
||||
/// Initialise un cache générique avec routes complètes
|
||||
///
|
||||
/// # Arguments
|
||||
///
|
||||
/// * `cache_dir` - Répertoire de stockage du cache
|
||||
/// * `limit` - Limite de taille du cache (nombre d'éléments)
|
||||
/// * `content_type` - Type MIME des fichiers (ex: "image/webp", "audio/flac")
|
||||
///
|
||||
/// # Routes créées
|
||||
///
|
||||
/// - Fichiers: `/{cache_name}/{cache_type}/{pk}[/{param}]`
|
||||
/// - API: `/api/{cache_name}/*`
|
||||
/// - Swagger: `/swagger-ui/{cache_name}`
|
||||
async fn init_generic_cache<C: CacheConfig + 'static>(
|
||||
&mut self,
|
||||
cache_dir: &str,
|
||||
limit: usize,
|
||||
content_type: &'static str,
|
||||
) -> anyhow::Result<Arc<Cache<C>>>;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user