From 08b6c4236d5e1eed33c52df78fde77398e5709b8 Mon Sep 17 00:00:00 2001 From: Eric Coissac Date: Fri, 17 Oct 2025 19:48:53 +0200 Subject: [PATCH] =?UTF-8?q?Ajout=20de=20fonctionnalit=C3=A9=20de=20downloa?= =?UTF-8?q?d=20asynchrone=20au=20pmocache?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- Cargo.lock | 2 + pmocache/Cargo.toml | 6 +- pmocache/DOWNLOAD_MODULE.md | 345 ++++++++++++++ pmocache/examples/README_EXAMPLES.md | 205 ++++++++ pmocache/examples/simple_transformer.rs | 53 +++ pmocache/examples/test_download.rs | 26 ++ .../examples/test_download_transformer.rs | 293 ++++++++++++ pmocache/src/api.rs | 317 +++++++++++++ pmocache/src/cache.rs | 261 +++++++++-- pmocache/src/cache_trait.rs | 23 +- pmocache/src/download.rs | 437 ++++++++++++++++++ pmocache/src/lib.rs | 21 +- pmocache/src/openapi.rs | 62 +++ pmocache/src/pmoserver_ext.rs | 231 +++++++-- 14 files changed, 2194 insertions(+), 88 deletions(-) create mode 100644 pmocache/DOWNLOAD_MODULE.md create mode 100644 pmocache/examples/README_EXAMPLES.md create mode 100644 pmocache/examples/simple_transformer.rs create mode 100644 pmocache/examples/test_download.rs create mode 100644 pmocache/examples/test_download_transformer.rs create mode 100644 pmocache/src/api.rs create mode 100644 pmocache/src/download.rs create mode 100644 pmocache/src/openapi.rs diff --git a/Cargo.lock b/Cargo.lock index f35eb600..0fcc78f7 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2279,12 +2279,14 @@ dependencies = [ "anyhow", "axum", "chrono", + "futures-util", "hex", "reqwest", "rusqlite", "serde", "sha1", "tokio", + "tokio-util", "tracing", "utoipa", ] diff --git a/pmocache/Cargo.toml b/pmocache/Cargo.toml index e5fada70..d827196a 100644 --- a/pmocache/Cargo.toml +++ b/pmocache/Cargo.toml @@ -8,7 +8,8 @@ edition = "2021" rusqlite = { version = "0.37.0", features = ["bundled"] } # HTTP client -reqwest = { version = "0.12", features = ["blocking"] } +reqwest = { version = "0.12", features = ["blocking", "stream"] } +futures-util = "0.3" # Cryptographie sha1 = "0.10" @@ -28,8 +29,9 @@ utoipa = { version = "5.3", optional = true } # Feature pour pmoserver (extension HTTP) axum = { version = "0.8", optional = true } tracing = { version = "0.1", optional = true } +tokio-util = { version = "0.7", features = ["io"], optional = true } [features] default = [] openapi = ["dep:utoipa"] -pmoserver = ["dep:axum", "dep:tracing"] +pmoserver = ["dep:axum", "dep:tracing", "dep:tokio-util"] diff --git a/pmocache/DOWNLOAD_MODULE.md b/pmocache/DOWNLOAD_MODULE.md new file mode 100644 index 00000000..ef793108 --- /dev/null +++ b/pmocache/DOWNLOAD_MODULE.md @@ -0,0 +1,345 @@ +# Module Download + +Module de téléchargement asynchrone avec support de transformation de stream. + +## Vue d'ensemble + +Le module `download` permet de télécharger des fichiers depuis une URL en tâche de fond avec : +- Suivi de la progression en temps réel +- Support de transformations de stream (conversion, compression, etc.) +- API non-bloquante avec attentes conditionnelles +- Gestion d'erreurs robuste + +## API + +### Types principaux + +#### `Download` +Objet représentant un téléchargement en cours, partagé via `Arc`. + +**Méthodes:** +- `filename() -> &Path` - Retourne le chemin du fichier de destination +- `wait_until_min_size(size: u64) -> Result<(), String>` - Attend que le fichier atteigne une taille minimale +- `wait_until_finished() -> Result<(), String>` - Attend la fin complète du téléchargement +- `open() -> io::Result` - Ouvre le fichier pour lecture +- `pos() -> u64` - Position de lecture actuelle +- `set_pos(pos: u64)` - Définit la position de lecture +- `expected_size() -> Option` - Taille attendue du fichier source (via Content-Length) +- `current_size() -> u64` - Taille actuellement téléchargée (source) +- `transformed_size() -> u64` - Taille des données transformées écrites +- `finished() -> bool` - Indique si le téléchargement est terminé +- `error() -> Option` - Retourne l'erreur éventuelle + +#### `StreamTransformer` +Type pour une fonction de transformation de stream. + +```rust +pub type StreamTransformer = Box< + dyn FnOnce( + reqwest::Response, + tokio::fs::File, + Arc, + ) -> Pin> + Send>> + + Send, +>; +``` + +**Paramètres:** +1. `reqwest::Response` - La réponse HTTP avec le stream de données +2. `tokio::fs::File` - Le fichier de destination ouvert en écriture +3. `Arc` - Callback pour mettre à jour la progression (taille transformée) + +**Retour:** +- `Future>` - Future qui se résout quand la transformation est terminée + +### Fonctions + +#### `download(filename, url) -> Arc` +Télécharge un fichier sans transformation. + +```rust +use pmocache::download::download; + +let dl = download("/tmp/file.dat", "https://example.com/file.dat"); +dl.wait_until_finished().await?; +``` + +#### `download_with_transformer(filename, url, transformer) -> Arc` +Télécharge un fichier avec une transformation optionnelle du stream. + +```rust +use pmocache::download::{download_with_transformer, StreamTransformer}; + +let transformer: StreamTransformer = Box::new(|response, mut file, update_progress| { + Box::pin(async move { + // Votre logique de transformation ici + Ok(()) + }) +}); + +let dl = download_with_transformer("/tmp/output.dat", "https://example.com/input.dat", Some(transformer)); +``` + +## Exemples d'utilisation + +### 1. Téléchargement simple + +```rust +use pmocache::download::download; + +#[tokio::main] +async fn main() -> Result<(), Box> { + let dl = download("/tmp/rust.html", "https://www.rust-lang.org/"); + + println!("Téléchargement démarré..."); + + // Attendre au moins 1KB + dl.wait_until_min_size(1024).await?; + println!("Au moins 1KB téléchargés"); + + // Attendre la fin + dl.wait_until_finished().await?; + println!("Terminé! Taille: {} bytes", dl.current_size().await); + + Ok(()) +} +``` + +### 2. Transformation en majuscules + +```rust +use pmocache::download::{download_with_transformer, StreamTransformer}; +use futures_util::StreamExt; +use tokio::io::AsyncWriteExt; + +fn uppercase_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 en majuscules + let uppercase: Vec = chunk + .iter() + .map(|&b| b.to_ascii_uppercase()) + .collect(); + + file.write_all(&uppercase).await.map_err(|e| e.to_string())?; + + total += uppercase.len() as u64; + update_progress(total); + } + + file.flush().await.map_err(|e| e.to_string())?; + Ok(()) + }) + }) +} + +#[tokio::main] +async fn main() { + let transformer = uppercase_transformer(); + let dl = download_with_transformer("/tmp/UPPERCASE.txt", "https://example.com/text.txt", Some(transformer)); + + dl.wait_until_finished().await.unwrap(); + println!("Fichier converti en majuscules!"); +} +``` + +### 3. Compression GZIP à la volée + +```rust +use pmocache::download::{download_with_transformer, StreamTransformer}; +use futures_util::StreamExt; +use tokio::io::AsyncWriteExt; +use async_compression::tokio::write::GzipEncoder; + +fn gzip_transformer() -> StreamTransformer { + Box::new(|response, file, update_progress| { + Box::pin(async move { + let mut encoder = GzipEncoder::new(file); + 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())?; + encoder.write_all(&chunk).await.map_err(|e| e.to_string())?; + + total += chunk.len() as u64; + update_progress(total); + } + + encoder.shutdown().await.map_err(|e| e.to_string())?; + Ok(()) + }) + }) +} +``` + +### 4. Conversion d'image (concept) + +```rust +// Exemple conceptuel de conversion WebP +// (nécessiterait une bibliothèque de traitement d'images) + +fn webp_transformer() -> StreamTransformer { + Box::new(|response, mut file, update_progress| { + Box::pin(async move { + // 1. Télécharger l'image en mémoire + let bytes = response.bytes().await.map_err(|e| e.to_string())?; + + // 2. Décoder l'image source + let img = image::load_from_memory(&bytes) + .map_err(|e| format!("Failed to decode image: {}", e))?; + + // 3. Encoder en WebP + let mut webp_data = Vec::new(); + let encoder = webp::Encoder::from_image(&img) + .map_err(|e| format!("Failed to create WebP encoder: {}", e))?; + let webp = encoder.encode(75.0); // Qualité 75% + webp_data.extend_from_slice(&*webp); + + // 4. Écrire le résultat + file.write_all(&webp_data).await.map_err(|e| e.to_string())?; + file.flush().await.map_err(|e| e.to_string())?; + + update_progress(webp_data.len() as u64); + Ok(()) + }) + }) +} + +// Utilisation +let transformer = webp_transformer(); +let dl = download_with_transformer( + "/tmp/image.webp", + "https://example.com/image.jpg", + Some(transformer) +); +``` + +### 5. Conversion audio (concept) + +```rust +// Exemple conceptuel de conversion MP3 -> FLAC +// (nécessiterait des bibliothèques audio comme symphonia) + +fn mp3_to_flac_transformer() -> StreamTransformer { + Box::new(|response, mut file, update_progress| { + Box::pin(async move { + // 1. Télécharger le MP3 en mémoire + let mp3_bytes = response.bytes().await.map_err(|e| e.to_string())?; + + // 2. Décoder le MP3 + let cursor = std::io::Cursor::new(mp3_bytes); + let mp3_decoder = minimp3::Decoder::new(cursor); + + let mut samples = Vec::new(); + let mut sample_rate = 0; + let mut channels = 0; + + for frame in mp3_decoder { + let frame = frame.map_err(|e| format!("MP3 decode error: {:?}", e))?; + if sample_rate == 0 { + sample_rate = frame.sample_rate; + channels = frame.channels; + } + samples.extend_from_slice(&frame.data); + } + + // 3. Encoder en FLAC + let mut flac_encoder = claxon::FlacEncoder::new( + &mut file, + sample_rate, + channels as u32, + 16, // bits per sample + ).map_err(|e| format!("FLAC encoder error: {:?}", e))?; + + for sample in samples { + flac_encoder.write_sample(sample as i32) + .map_err(|e| format!("FLAC write error: {:?}", e))?; + } + + flac_encoder.finish() + .map_err(|e| format!("FLAC finalize error: {:?}", e))?; + + file.flush().await.map_err(|e| e.to_string())?; + + // Note: on ne peut pas facilement connaître la taille finale avant d'avoir tout encodé + // Pour un suivi précis, il faudrait encoder par chunks + Ok(()) + }) + }) +} +``` + +## Cas d'usage dans PMOMusic + +### 1. Cache audio avec conversion +```rust +// Télécharger du MP3 et le convertir en FLAC pour le cache +let transformer = mp3_to_flac_transformer(); +let dl = download_with_transformer( + cache_path, + audio_url, + Some(transformer) +); +``` + +### 2. Cache d'images avec WebP +```rust +// Télécharger une image et la convertir en WebP +let transformer = webp_transformer(); +let dl = download_with_transformer( + cover_cache_path, + cover_url, + Some(transformer) +); +``` + +### 3. Streaming progressif +```rust +// Commencer à lire le fichier dès qu'on a assez de données +let dl = download(audio_path, stream_url); + +// Attendre au moins 256KB pour commencer la lecture +dl.wait_until_min_size(256 * 1024).await?; + +// Ouvrir le fichier et commencer à lire pendant que le téléchargement continue +let file = dl.open()?; +// ... lecture du fichier +``` + +## Notes d'implémentation + +### Thread safety +- Tous les objets sont thread-safe via `Arc` et `RwLock` +- Le téléchargement s'exécute dans un `tokio::spawn` séparé +- Les callbacks de progression utilisent `Arc` pour être partagés + +### Gestion des erreurs +- Les erreurs sont capturées et stockées dans l'état +- `wait_until_*` retourne l'erreur si elle existe +- Le téléchargement est marqué comme terminé même en cas d'erreur + +### Performance +- Téléchargement par chunks (stream) +- Transformation à la volée sans buffer intermédiaire complet (selon le transformer) +- Mise à jour de la progression asynchrone via spawn + +## Dépendances + +```toml +[dependencies] +reqwest = { version = "0.12", features = ["stream"] } +futures-util = "0.3" +tokio = { version = "1.0", features = ["full"] } + +# Optionnel selon les transformers utilisés +async-compression = "0.4" # Pour GZIP +image = "0.24" # Pour images +webp = "0.2" # Pour WebP +``` diff --git a/pmocache/examples/README_EXAMPLES.md b/pmocache/examples/README_EXAMPLES.md new file mode 100644 index 00000000..05aba976 --- /dev/null +++ b/pmocache/examples/README_EXAMPLES.md @@ -0,0 +1,205 @@ +# Exemples du module Download + +Ce répertoire contient des exemples d'utilisation du module `download` de pmocache. + +## Fichiers + +### `test_download.rs` +Exemple basique de téléchargement sans transformation. + +**Utilisation:** +```bash +cargo run --example test_download +``` + +### `test_download_transformer.rs` +Exemples complets de transformers : +- Transformation en majuscules +- Suppression de header (skip N bytes) +- Numérotation des lignes +- Compression GZIP (commenté, nécessite async-compression) + +**Utilisation:** +```bash +cargo run --example test_download_transformer +``` + +### `simple_transformer.rs` +Exemple de documentation montrant la syntaxe et l'API. + +**Utilisation:** +```bash +cargo run --example simple_transformer +``` + +## Concepts clés + +### 1. Téléchargement simple + +```rust +use pmocache::download::download; + +let dl = download("/tmp/file.dat", "https://example.com/file.dat"); +dl.wait_until_finished().await?; +``` + +### 2. Téléchargement avec transformer + +Un transformer est une fonction qui : +1. Reçoit le stream de réponse HTTP +2. Reçoit un fichier ouvert en écriture +3. Reçoit un callback de progression +4. Traite les données à la volée +5. Écrit le résultat transformé dans le fichier + +```rust +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 + let transformed = your_transformation(&chunk); + + // Écrire le résultat + file.write_all(&transformed).await.map_err(|e| e.to_string())?; + + // Mettre à jour la progression + total += transformed.len() as u64; + update_progress(total); + } + + file.flush().await.map_err(|e| e.to_string())?; + Ok(()) + }) +}); + +let dl = download_with_transformer("/tmp/output.dat", "https://example.com/input.dat", Some(transformer)); +``` + +### 3. Suivi de progression + +```rust +let dl = download("/tmp/file.dat", "https://example.com/file.dat"); + +// Attendre au moins 1MB +dl.wait_until_min_size(1024 * 1024).await?; +println!("Au moins 1MB téléchargés"); + +// Voir la progression +loop { + let current = dl.current_size().await; + let expected = dl.expected_size().await; + + if let Some(total) = expected { + println!("Progression: {}/{} bytes ({:.1}%)", + current, total, 100.0 * current as f64 / total as f64); + } else { + println!("Téléchargés: {} bytes", current); + } + + if dl.finished().await { + break; + } + + tokio::time::sleep(Duration::from_millis(100)).await; +} +``` + +## Cas d'usage pour PMOMusic + +### Conversion d'images pour le cache + +```rust +// Télécharger une couverture d'album et la convertir en WebP +fn webp_transformer() -> StreamTransformer { + Box::new(|response, mut file, update_progress| { + Box::pin(async move { + let bytes = response.bytes().await.map_err(|e| e.to_string())?; + let img = image::load_from_memory(&bytes) + .map_err(|e| format!("Decode error: {}", e))?; + + let encoder = webp::Encoder::from_image(&img) + .map_err(|e| format!("Encode error: {}", e))?; + let webp = encoder.encode(75.0); + + file.write_all(&*webp).await.map_err(|e| e.to_string())?; + file.flush().await.map_err(|e| e.to_string())?; + + update_progress(webp.len() as u64); + Ok(()) + }) + }) +} + +// Utilisation dans pmocovers +let transformer = webp_transformer(); +let dl = download_with_transformer(cache_path, cover_url, Some(transformer)); +``` + +### Conversion audio pour le cache + +```rust +// Télécharger du MP3 et le convertir en FLAC +fn mp3_to_flac_transformer() -> StreamTransformer { + Box::new(|response, mut file, update_progress| { + Box::pin(async move { + let mp3_bytes = response.bytes().await.map_err(|e| e.to_string())?; + + // Décoder MP3 + let decoded = decode_mp3(&mp3_bytes)?; + + // Encoder FLAC + let flac_bytes = encode_flac(&decoded)?; + + file.write_all(&flac_bytes).await.map_err(|e| e.to_string())?; + file.flush().await.map_err(|e| e.to_string())?; + + update_progress(flac_bytes.len() as u64); + Ok(()) + }) + }) +} + +// Utilisation dans pmoaudiocache +let transformer = mp3_to_flac_transformer(); +let dl = download_with_transformer(cache_path, audio_url, Some(transformer)); +``` + +### Streaming progressif + +```rust +// Commencer à lire pendant le téléchargement +let dl = download(audio_path, stream_url); + +// Attendre le buffer minimal (256KB) +dl.wait_until_min_size(256 * 1024).await?; + +// Ouvrir et commencer à lire +let mut file = dl.open()?; +let mut buffer = [0u8; 4096]; + +loop { + // Lire ce qui est disponible + match file.read(&mut buffer) { + Ok(0) if dl.finished().await => break, // EOF + Ok(0) => { + // Pas encore de données, attendre un peu + tokio::time::sleep(Duration::from_millis(10)).await; + } + Ok(n) => { + // Traiter les données lues + process_audio_chunk(&buffer[..n]); + } + Err(e) => return Err(e.into()), + } +} +``` + +## Voir aussi + +- [DOWNLOAD_MODULE.md](../DOWNLOAD_MODULE.md) - Documentation complète du module +- [src/download.rs](../src/download.rs) - Code source diff --git a/pmocache/examples/simple_transformer.rs b/pmocache/examples/simple_transformer.rs new file mode 100644 index 00000000..0ced7f21 --- /dev/null +++ b/pmocache/examples/simple_transformer.rs @@ -0,0 +1,53 @@ +/// Exemple minimal de transformer sans dépendances externes complexes + +// Import direct du type depuis le module +// Note: Cet exemple montre comment utiliser l'API de transformation + +fn main() { + println!("Exemple d'utilisation du module download avec transformers\n"); + + println!("1. Téléchargement simple:"); + println!(" let dl = download(\"/tmp/file.dat\", \"https://example.com/file.dat\");"); + println!(" dl.wait_until_finished().await?;\n"); + + println!("2. Téléchargement avec transformation:"); + println!(" let transformer: StreamTransformer = Box::new(|response, mut file, update_progress| {{"); + println!(" Box::pin(async move {{"); + println!(" let mut stream = response.bytes_stream();"); + println!(" let mut total = 0u64;"); + println!(); + println!(" while let Some(chunk_result) = stream.next().await {{"); + println!(" let chunk = chunk_result.map_err(|e| e.to_string())?;"); + println!(); + println!(" // Transformation ici (ex: compression, conversion)"); + println!(" let transformed = process(chunk);"); + println!(); + println!(" file.write_all(&transformed).await.map_err(|e| e.to_string())?;"); + println!(" total += transformed.len() as u64;"); + println!(" update_progress(total);"); + println!(" }}"); + println!(); + println!(" file.flush().await.map_err(|e| e.to_string())?;"); + println!(" Ok(())"); + println!(" }})"); + println!(" }});\n"); + + println!(" let dl = download_with_transformer(\"/tmp/out.dat\", \"https://example.com/in.dat\", Some(transformer));"); + println!(" dl.wait_until_finished().await?;\n"); + + println!("3. Méthodes disponibles sur Download:"); + println!(" - filename() : Chemin du fichier"); + println!(" - current_size() : Taille téléchargée (source)"); + println!(" - transformed_size() : Taille transformée (destination)"); + println!(" - expected_size() : Taille attendue (Content-Length)"); + println!(" - finished() : Téléchargement terminé?"); + println!(" - error() : Erreur éventuelle"); + println!(" - wait_until_min_size(n) : Attend au moins n bytes"); + println!(" - wait_until_finished() : Attend la fin"); + println!(" - open() : Ouvre le fichier pour lecture"); + println!(" - pos() / set_pos() : Position de lecture\n"); + + println!("Pour des exemples complets, voir:"); + println!(" - examples/test_download_transformer.rs"); + println!(" - DOWNLOAD_MODULE.md"); +} diff --git a/pmocache/examples/test_download.rs b/pmocache/examples/test_download.rs new file mode 100644 index 00000000..79e7d674 --- /dev/null +++ b/pmocache/examples/test_download.rs @@ -0,0 +1,26 @@ +// Simple test pour vérifier la compilation du module download + +#[tokio::main] +async fn main() { + println!("Module download compilé avec succès!"); + + // Test basique (commenté pour ne pas vraiment télécharger) + /* + let dl = download::download("/tmp/test.html", "https://www.rust-lang.org/"); + + println!("Téléchargement démarré..."); + + match dl.wait_until_min_size(100).await { + Ok(_) => println!("Au moins 100 bytes téléchargés"), + Err(e) => eprintln!("Erreur: {}", e), + } + + match dl.wait_until_finished().await { + Ok(_) => { + println!("Téléchargement terminé!"); + println!("Taille finale: {} bytes", dl.current_size().await); + } + Err(e) => eprintln!("Erreur: {}", e), + } + */ +} diff --git a/pmocache/examples/test_download_transformer.rs b/pmocache/examples/test_download_transformer.rs new file mode 100644 index 00000000..a99dea4b --- /dev/null +++ b/pmocache/examples/test_download_transformer.rs @@ -0,0 +1,293 @@ +// Exemple d'utilisation du module download avec transformations + +use pmocache::download::{download_with_transformer, StreamTransformer}; +use futures_util::StreamExt; +use tokio::io::AsyncWriteExt; + +/// Exemple de transformer qui compresse les données en gzip +/// +/// Note: Cette fonction nécessite la dépendance `async-compression` +/// Pour l'utiliser, ajoutez à Cargo.toml: +/// ```toml +/// [dev-dependencies] +/// async-compression = { version = "0.4", features = ["tokio", "gzip"] } +/// ``` +#[allow(dead_code)] +fn create_gzip_transformer() -> StreamTransformer { + // Commenté car nécessite async-compression + // Décommentez si vous ajoutez la dépendance + unimplemented!("Cette fonction nécessite la dépendance async-compression") + + /* + Box::new(|response, mut file, update_progress| { + Box::pin(async move { + use async_compression::tokio::write::GzipEncoder; + + let mut encoder = GzipEncoder::new(&mut file); + let mut stream = response.bytes_stream(); + let mut total_written = 0u64; + + while let Some(chunk_result) = stream.next().await { + let chunk = chunk_result.map_err(|e| format!("Failed to read chunk: {}", e))?; + + encoder + .write_all(&chunk) + .await + .map_err(|e| format!("Failed to write compressed data: {}", e))?; + + total_written += chunk.len() as u64; + update_progress(total_written); + } + + encoder + .shutdown() + .await + .map_err(|e| format!("Failed to finalize compression: {}", e))?; + + Ok(()) + }) + }) + */ +} + +/// Exemple de transformer qui convertit les données en majuscules (exemple simple) +fn create_uppercase_transformer() -> StreamTransformer { + Box::new(|response, mut file, update_progress| { + Box::pin(async move { + let mut stream = response.bytes_stream(); + let mut total_written = 0u64; + + while let Some(chunk_result) = stream.next().await { + let chunk = chunk_result.map_err(|e| format!("Failed to read chunk: {}", e))?; + + // Transformer en majuscules (seulement pour texte ASCII) + let transformed: Vec = chunk + .iter() + .map(|&b| if b.is_ascii_lowercase() { b.to_ascii_uppercase() } else { b }) + .collect(); + + file.write_all(&transformed) + .await + .map_err(|e| format!("Failed to write: {}", e))?; + + total_written += transformed.len() as u64; + update_progress(total_written); + } + + file.flush() + .await + .map_err(|e| format!("Failed to flush: {}", e))?; + + Ok(()) + }) + }) +} + +/// Exemple de transformer qui saute les N premiers bytes (utile pour enlever des headers) +fn create_skip_header_transformer(skip_bytes: usize) -> StreamTransformer { + Box::new(move |response, mut file, update_progress| { + Box::pin(async move { + let mut stream = response.bytes_stream(); + let mut skipped = 0usize; + let mut total_written = 0u64; + + while let Some(chunk_result) = stream.next().await { + let chunk = chunk_result.map_err(|e| format!("Failed to read chunk: {}", e))?; + + let to_write = if skipped < skip_bytes { + let remaining_to_skip = skip_bytes - skipped; + if chunk.len() <= remaining_to_skip { + skipped += chunk.len(); + continue; + } else { + skipped = skip_bytes; + &chunk[remaining_to_skip..] + } + } else { + &chunk[..] + }; + + file.write_all(to_write) + .await + .map_err(|e| format!("Failed to write: {}", e))?; + + total_written += to_write.len() as u64; + update_progress(total_written); + } + + file.flush() + .await + .map_err(|e| format!("Failed to flush: {}", e))?; + + Ok(()) + }) + }) +} + +/// Exemple de transformer qui compte les lignes et ajoute des numéros +fn create_line_number_transformer() -> StreamTransformer { + Box::new(|response, mut file, update_progress| { + Box::pin(async move { + let mut stream = response.bytes_stream(); + let mut line_number = 1u32; + let mut buffer = Vec::new(); + let mut total_written = 0u64; + + while let Some(chunk_result) = stream.next().await { + let chunk = chunk_result.map_err(|e| format!("Failed to read chunk: {}", e))?; + buffer.extend_from_slice(&chunk); + + // Traiter les lignes complètes dans le buffer + while let Some(newline_pos) = buffer.iter().position(|&b| b == b'\n') { + let line = &buffer[..newline_pos]; + + // Écrire le numéro de ligne et la ligne + let numbered_line = format!("{:6}: ", line_number); + file.write_all(numbered_line.as_bytes()) + .await + .map_err(|e| format!("Failed to write: {}", e))?; + + file.write_all(line) + .await + .map_err(|e| format!("Failed to write: {}", e))?; + + file.write_all(b"\n") + .await + .map_err(|e| format!("Failed to write: {}", e))?; + + total_written += numbered_line.len() as u64 + line.len() as u64 + 1; + update_progress(total_written); + + line_number += 1; + buffer.drain(..=newline_pos); + } + } + + // Traiter la dernière ligne si elle n'a pas de newline + if !buffer.is_empty() { + let numbered_line = format!("{:6}: ", line_number); + file.write_all(numbered_line.as_bytes()) + .await + .map_err(|e| format!("Failed to write: {}", e))?; + + file.write_all(&buffer) + .await + .map_err(|e| format!("Failed to write: {}", e))?; + + total_written += numbered_line.len() as u64 + buffer.len() as u64; + update_progress(total_written); + } + + file.flush() + .await + .map_err(|e| format!("Failed to flush: {}", e))?; + + Ok(()) + }) + }) +} + +#[tokio::main] +async fn main() { + println!("=== Exemples de transformers pour le module download ===\n"); + + let temp_dir = std::env::temp_dir(); + + // Exemple 1: Téléchargement avec transformation en majuscules + println!("1. Téléchargement avec transformation en MAJUSCULES"); + let uppercase_file = temp_dir.join("uppercase_example.txt"); + let _ = std::fs::remove_file(&uppercase_file); + + let transformer = create_uppercase_transformer(); + let dl = download_with_transformer( + &uppercase_file, + "https://www.rust-lang.org/", + Some(transformer), + ); + + println!(" Téléchargement démarré..."); + match dl.wait_until_finished().await { + Ok(_) => { + println!(" ✓ Téléchargement terminé!"); + println!(" - Taille source: {} bytes", dl.current_size().await); + println!(" - Taille transformée: {} bytes", dl.transformed_size().await); + } + Err(e) => { + eprintln!(" ✗ Erreur: {}", e); + } + } + + // Exemple 2: Skip header + println!("\n2. Téléchargement en sautant les 100 premiers bytes"); + let skip_file = temp_dir.join("skip_header_example.txt"); + let _ = std::fs::remove_file(&skip_file); + + let transformer = create_skip_header_transformer(100); + let dl = download_with_transformer( + &skip_file, + "https://www.rust-lang.org/", + Some(transformer), + ); + + match dl.wait_until_finished().await { + Ok(_) => { + println!(" ✓ Téléchargement terminé!"); + println!(" - Taille transformée: {} bytes", dl.transformed_size().await); + } + Err(e) => { + eprintln!(" ✗ Erreur: {}", e); + } + } + + // Exemple 3: Numérotation des lignes + println!("\n3. Téléchargement avec numérotation des lignes"); + let numbered_file = temp_dir.join("numbered_example.txt"); + let _ = std::fs::remove_file(&numbered_file); + + let transformer = create_line_number_transformer(); + let dl = download_with_transformer( + &numbered_file, + "https://www.rust-lang.org/", + Some(transformer), + ); + + match dl.wait_until_finished().await { + Ok(_) => { + println!(" ✓ Téléchargement terminé!"); + println!(" - Taille transformée: {} bytes", dl.transformed_size().await); + } + Err(e) => { + eprintln!(" ✗ Erreur: {}", e); + } + } + + println!("\n=== Exemples terminés ==="); + println!("Fichiers créés dans: {:?}", temp_dir); + + // Note: Commenté car nécessite la dépendance async-compression + /* + println!("\n4. Téléchargement avec compression GZIP"); + let gzip_file = temp_dir.join("compressed_example.gz"); + let _ = std::fs::remove_file(&gzip_file); + + let transformer = create_gzip_transformer(); + let dl = download_with_transformer( + &gzip_file, + "https://www.rust-lang.org/", + Some(transformer), + ); + + match dl.wait_until_finished().await { + Ok(_) => { + println!(" ✓ Téléchargement terminé!"); + println!(" - Taille source: {} bytes", dl.current_size().await); + println!(" - Taille compressée: {} bytes", dl.transformed_size().await); + let ratio = 100.0 * dl.transformed_size().await as f64 / dl.current_size().await as f64; + println!(" - Ratio de compression: {:.1}%", ratio); + } + Err(e) => { + eprintln!(" ✗ Erreur: {}", e); + } + } + */ +} diff --git a/pmocache/src/api.rs b/pmocache/src/api.rs new file mode 100644 index 00000000..b18f4651 --- /dev/null +++ b/pmocache/src/api.rs @@ -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, + /// Taille après transformation + pub transformed_size: Option, + /// Taille totale attendue + pub expected_size: Option, + /// Téléchargement terminé + pub finished: bool, + /// Erreur éventuelle + pub error: Option, +} + +/// 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, +} + +/// 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( + State(cache): State>>, +) -> 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( + State(cache): State>>, + Path(pk): Path, +) -> 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( + State(cache): State>>, + Path(pk): Path, +) -> 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( + State(cache): State>>, + Json(req): Json, +) -> 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( + State(cache): State>>, + Path(pk): Path, +) -> 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( + State(cache): State>>, +) -> 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( + State(cache): State>>, +) -> 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(), + } +} diff --git a/pmocache/src/cache.rs b/pmocache/src/cache.rs index 4f8bf234..056a4453 100644 --- a/pmocache/src/cache.rs +++ b/pmocache/src/cache.rs @@ -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`. -/// 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 { /// Répertoire de stockage @@ -53,6 +56,8 @@ pub struct Cache { base_url: String, /// Base de données SQLite pub db: Arc, + /// Map des downloads en cours (pk -> Download) + downloads: Arc>>>, /// Phantom data pour le type de configuration _phantom: std::marker::PhantomData, } @@ -75,12 +80,24 @@ impl Cache { 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 { + 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 Cache { /// /// La clé primaire (pk) du fichier dans le cache pub async fn add_from_url(&self, url: &str, collection: Option<&str>) -> Result { - 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 { + 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 Cache { 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 Cache { 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> { + 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 { + 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 { + 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 { + 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 Cache { /// Implémentation du trait FileCache pour Cache -impl FileCache for Cache { - fn cache_type(&self) -> &str { - C::cache_type() +impl FileCache for Cache { + fn get_cache_dir(&self) -> &Path { + self.cache_dir() + } + + fn get_database(&self) -> Arc { + self.db.clone() + } + + fn get_base_url(&self) -> &str { + &self.base_url } fn validate_data(&self, data: &[u8]) -> Result> { @@ -246,23 +444,12 @@ impl FileCache for Cache { self.add_from_url(url, collection).await } - async fn ensure_from_url(&self, url: &str, collection: Option<&str>) -> Result { - self.ensure_from_url(url, collection).await + async fn add_from_file(&self, path: &str, collection: Option<&str>) -> Result { + self.add_from_file(path, collection).await } - async fn add(&self, url: &str, data: &[u8], collection: Option<&str>) -> Result { - // 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 { + self.ensure_from_url(url, collection).await } async fn get(&self, pk: &str) -> Result { @@ -280,12 +467,4 @@ impl FileCache for Cache { 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() - } } diff --git a/pmocache/src/cache_trait.rs b/pmocache/src/cache_trait.rs index 9e86a99b..a6845115 100644 --- a/pmocache/src/cache_trait.rs +++ b/pmocache/src/cache_trait.rs @@ -12,6 +12,7 @@ pub trait FileCache: Send + Sync { fn get_cache_dir(&self) -> &Path; fn get_database(&self) -> Arc; 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: 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: 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: Send + Sync { /// La clé primaire (pk) du fichier dans le cache async fn add_from_url(&self, url: &str, collection: Option<&str>) -> Result; + /// 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; + /// 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. diff --git a/pmocache/src/download.rs b/pmocache/src/download.rs new file mode 100644 index 00000000..30233818 --- /dev/null +++ b/pmocache/src/download.rs @@ -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, + ) -> Pin> + 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, + /// 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, +} + +/// 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>, +} + +impl Download { + /// Crée une nouvelle instance de Download + fn new(filename: PathBuf) -> Arc { + 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::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 { + 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 { + 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 qui permet de suivre la progression du téléchargement +pub fn download>(filename: P, url: &str) -> Arc { + 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 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>( + filename: P, + url: &str, + transformer: Option, +) -> Arc { + 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>, + transformer: Option, +) -> 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 = 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>, +) -> 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); + } +} diff --git a/pmocache/src/lib.rs b/pmocache/src/lib.rs index 534faa62..3c281e17 100644 --- a/pmocache/src/lib.rs +++ b/pmocache/src/lib.rs @@ -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, +}; diff --git a/pmocache/src/openapi.rs b/pmocache/src/openapi.rs new file mode 100644 index 00000000..127bad0b --- /dev/null +++ b/pmocache/src/openapi.rs @@ -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::, + $crate::api::get_item_info::, + $crate::api::get_download_status::, + $crate::api::add_item::, + $crate::api::delete_item::, + $crate::api::purge_cache::, + $crate::api::consolidate_cache::, + ), + 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; + }; +} diff --git a/pmocache/src/pmoserver_ext.rs b/pmocache/src/pmoserver_ext.rs index db262557..3840adfe 100644 --- a/pmocache/src/pmoserver_ext.rs +++ b/pmocache/src/pmoserver_ext.rs @@ -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( State((cache, content_type)): State<(Arc>, &'static str)>, Path(pk): Path, ) -> 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( State((cache, content_type)): State<(Arc>, &'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( + cache: &Arc>, + 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, + 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( 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( /// /// # 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( /// "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( cache: Arc>, 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::)) - .route("/:pk/:param", get(get_file_with_param::)) + .route(&path_without_param, get(get_file::)) + .route(&path_with_param, get(get_file_with_param::)) .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( + cache: Arc>, +) -> Router { + use crate::api; + + Router::new() + .route( + "/", + get(api::list_items::) + .post(api::add_item::) + .delete(api::purge_cache::), + ) + .route( + "/:pk", + get(api::get_item_info::) + .delete(api::delete_item::), + ) + .route("/:pk/status", get(api::get_download_status::)) + .route("/consolidate", post(api::consolidate_cache::)) + .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( + &mut self, + cache_dir: &str, + limit: usize, + content_type: &'static str, + ) -> anyhow::Result>>; +}