From 6b3851de19e1018058be8416cd85006677d3b4f2 Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 5 Nov 2025 15:49:23 +0000 Subject: [PATCH] feat: Add play_and_cache example for Radio Paradise streaming and playback MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Crée un nouvel exemple complet qui démontre l'utilisation de tout le pipeline: - Téléchargement d'un bloc Radio Paradise - Cache FLAC via FlacCacheSink - Playlist alimentée automatiquement - Lecture en temps réel via PlaylistSource et AudioSink Architecture à deux pipelines : Pipeline 1 (Download & Cache): RadioParadiseStreamSource → FlacCacheSink (avec playlist abonnée) Pipeline 2 (Playback): PlaylistSource (lit la playlist) → AudioSink (joue l'audio) Les deux pipelines s'exécutent en parallèle, permettant la lecture pendant le téléchargement. Modifications: - Ajout pmoaudio-ext avec feature playlist dans pmoparadise - Nouvelle feature "full" combinant pmoaudio + pmoaudio-ext - Logs détaillés à tous les niveaux (DEBUG) Usage: cargo run --example play_and_cache --features full -- --- Cargo.lock | 1 + pmoparadise/Cargo.toml | 5 + pmoparadise/examples/play_and_cache.rs | 261 +++++++++++++++++++++++++ 3 files changed, 267 insertions(+) create mode 100644 pmoparadise/examples/play_and_cache.rs diff --git a/Cargo.lock b/Cargo.lock index 16f3b90d..011554dd 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3070,6 +3070,7 @@ dependencies = [ "futures-util", "hex", "pmoaudio", + "pmoaudio-ext", "pmoaudiocache", "pmoconfig", "pmocovers", diff --git a/pmoparadise/Cargo.toml b/pmoparadise/Cargo.toml index fe7b770c..987e3ab9 100644 --- a/pmoparadise/Cargo.toml +++ b/pmoparadise/Cargo.toml @@ -51,6 +51,9 @@ symphonia = { version = "0.5", features = ["all"] } # Audio decoding - claxon for FLAC streaming claxon = "0.4" +# pmoaudio-ext with playlist support (optional for examples) +pmoaudio-ext = { path = "../pmoaudio-ext", optional = true, features = ["playlist"] } + # Common music source traits pmosource = { path = "../pmosource" } @@ -87,6 +90,8 @@ pmoconfig = ["dep:pmoconfig"] cache = [] # Active le support pmoaudio node (RadioParadiseStreamSource) pmoaudio = ["dep:pmoaudio", "dep:pmoflac", "dep:pmometadata", "dep:futures-util"] +# Active le support complet avec playlist (pour les exemples avancés) +full = ["pmoaudio", "dep:pmoaudio-ext"] [dev-dependencies] # Tests diff --git a/pmoparadise/examples/play_and_cache.rs b/pmoparadise/examples/play_and_cache.rs new file mode 100644 index 00000000..e4f5ae11 --- /dev/null +++ b/pmoparadise/examples/play_and_cache.rs @@ -0,0 +1,261 @@ +//! Télécharge un bloc Radio Paradise, le cache, et le joue en même temps +//! +//! Ce programme démontre l'utilisation complète de la chaîne : +//! 1. RadioParadiseStreamSource - Télécharge et décode un bloc FLAC +//! 2. FlacCacheSink - Cache chaque piste en FLAC et alimente une playlist +//! 3. PlaylistSource - Lit la playlist pendant le téléchargement +//! 4. AudioSink - Joue l'audio sur la sortie standard +//! +//! Architecture : +//! ```text +//! Pipeline 1 (Download & Cache): +//! RadioParadiseStreamSource → FlacCacheSink (avec playlist abonnée) +//! +//! Pipeline 2 (Playback): +//! PlaylistSource (lit la playlist) → AudioSink (joue l'audio) +//! ``` +//! +//! Usage: +//! cargo run --example play_and_cache --features playlist -- +//! +//! Exemple: +//! cargo run --example play_and_cache --features playlist -- 0 # Main Mix +//! cargo run --example play_and_cache --features playlist -- 2 # Rock Mix + +use pmoaudio::{AudioPipelineNode, AudioSink}; +use pmoaudio_ext::{FlacCacheSink, PlaylistSource}; +use pmoaudiocache::AudioCacheConfigExt; +use pmocovers::CoverCacheConfigExt; +use pmoparadise::{RadioParadiseClient, RadioParadiseStreamSource}; +use std::env; +use tokio_util::sync::CancellationToken; + +#[tokio::main] +async fn main() -> Result<(), Box> { + // Initialiser tracing avec beaucoup de logs + tracing_subscriber::fmt() + .with_env_filter( + tracing_subscriber::EnvFilter::from_default_env() + .add_directive(tracing::Level::DEBUG.into()) + .add_directive("pmoaudio=debug".parse()?) + .add_directive("pmoaudio_ext=debug".parse()?) + .add_directive("pmoplaylist=debug".parse()?) + .add_directive("pmoparadise=debug".parse()?) + .add_directive("pmoaudiocache=debug".parse()?) + ) + .init(); + + tracing::info!("=== Radio Paradise Play & Cache ==="); + + // Récupérer les arguments + let args: Vec = env::args().collect(); + if args.len() != 2 { + eprintln!("Usage: {} ", args[0]); + eprintln!(); + eprintln!("Downloads a Radio Paradise block, caches it, and plays it simultaneously."); + eprintln!(); + eprintln!("Channel IDs:"); + eprintln!(" 0 - Main Mix (eclectic, diverse mix)"); + eprintln!(" 1 - Mellow Mix (smooth, chilled music)"); + eprintln!(" 2 - Rock Mix (classic & modern rock)"); + eprintln!(" 3 - World/Etc Mix (global sounds)"); + std::process::exit(1); + } + + let channel_id: u8 = match args[1].parse() { + Ok(id) if id <= 3 => id, + _ => { + eprintln!("Error: channel_id must be a number between 0 and 3"); + std::process::exit(1); + } + }; + + tracing::info!("Channel ID: {}", channel_id); + + // ═══════════════════════════════════════════════════════════════════════════ + // Initialiser les caches et le gestionnaire de playlist + // ═══════════════════════════════════════════════════════════════════════════ + + tracing::info!("Initializing configuration..."); + let mut config = pmoconfig::Config::load_config(&std::env::var("PMO_CONFIG_DIR").unwrap_or_else(|_| "/tmp/pmomusic".to_string()))?; + + tracing::info!("Initializing audio cache..."); + let audio_cache = config.init_audio_cache_configured().await?; + tracing::debug!("Audio cache initialized: {:?}", audio_cache.root()); + + tracing::info!("Initializing cover cache..."); + let cover_cache = config.init_cover_cache_configured().await?; + tracing::debug!("Cover cache initialized: {:?}", cover_cache.root()); + + tracing::info!("Initializing playlist manager..."); + let playlist_manager = pmoplaylist::Manager::init().await; + tracing::debug!("Playlist manager initialized"); + + // ═══════════════════════════════════════════════════════════════════════════ + // Créer la playlist pour ce channel + // ═══════════════════════════════════════════════════════════════════════════ + + let playlist_id = format!("radio-paradise-ch{}", channel_id); + tracing::info!("Creating playlist: {}", playlist_id); + + // Créer la playlist (ou la vider si elle existe) + let mut writer = playlist_manager.create_persistent_playlist(playlist_id.clone()).await?; + writer.set_title(format!("Radio Paradise - Channel {}", channel_id)).await?; + writer.flush().await?; // Vider la playlist si elle existait + tracing::debug!("Playlist created and flushed"); + + // Créer le reader pour la lecture + let reader = playlist_manager.get_read_handle(&playlist_id).await?; + tracing::debug!("Read handle created"); + + // ═══════════════════════════════════════════════════════════════════════════ + // Récupérer les infos du bloc à télécharger + // ═══════════════════════════════════════════════════════════════════════════ + + tracing::info!("Fetching current block metadata..."); + let client = RadioParadiseClient::builder() + .channel(channel_id) + .build() + .await?; + + let block = client.get_block(None).await?; + + tracing::info!("Block Information:"); + tracing::info!(" Event ID: {}", block.event); + tracing::info!(" Songs: {}", block.song_count()); + tracing::info!(" Duration: {:.1} minutes", block.length as f64 / 60000.0); + tracing::info!(""); + + tracing::info!("Tracklist:"); + for (index, song) in block.songs_ordered() { + tracing::info!( + " {:2}. {} - {} ({})", + index + 1, + song.artist, + song.title, + song.album.as_deref().unwrap_or("Unknown Album") + ); + } + tracing::info!(""); + + // ═══════════════════════════════════════════════════════════════════════════ + // Pipeline 1: Téléchargement et cache + // ═══════════════════════════════════════════════════════════════════════════ + + tracing::info!("Creating download pipeline..."); + + // Créer la source Radio Paradise + let mut download_source = RadioParadiseStreamSource::new(client); + download_source.push_block_id(block.event); + tracing::debug!("RadioParadiseStreamSource created with block {}", block.event); + + // Créer le sink de cache FLAC + let mut cache_sink = FlacCacheSink::new(audio_cache.clone(), cover_cache.clone()); + cache_sink.register_playlist(writer); + tracing::debug!("FlacCacheSink created and registered with playlist"); + + // Connecter source → sink + download_source.register(Box::new(cache_sink)); + tracing::info!("Download pipeline connected: RadioParadiseStreamSource → FlacCacheSink"); + + // ═══════════════════════════════════════════════════════════════════════════ + // Pipeline 2: Lecture depuis la playlist + // ═══════════════════════════════════════════════════════════════════════════ + + tracing::info!("Creating playback pipeline..."); + + // Créer la source playlist + let mut playlist_source = PlaylistSource::new(reader, audio_cache.clone()); + tracing::debug!("PlaylistSource created"); + + // Créer le sink audio avec volume à 80% + let audio_sink = AudioSink::with_volume(0.8); + tracing::debug!("AudioSink created with volume 0.8"); + + // Connecter playlist → audio + playlist_source.register(Box::new(audio_sink)); + tracing::info!("Playback pipeline connected: PlaylistSource → AudioSink"); + + // ═══════════════════════════════════════════════════════════════════════════ + // Lancer les deux pipelines en parallèle + // ═══════════════════════════════════════════════════════════════════════════ + + tracing::info!(""); + tracing::info!("========================================"); + tracing::info!("Starting both pipelines..."); + tracing::info!("Pipeline 1: Downloading and caching"); + tracing::info!("Pipeline 2: Playing from playlist"); + tracing::info!("========================================"); + tracing::info!(""); + + let stop_token = CancellationToken::new(); + let stop_token_download = stop_token.clone(); + let stop_token_playback = stop_token.clone(); + + // Gérer Ctrl+C + let stop_token_ctrl_c = stop_token.clone(); + tokio::spawn(async move { + tokio::signal::ctrl_c().await.ok(); + tracing::warn!("Received Ctrl+C, stopping..."); + stop_token_ctrl_c.cancel(); + }); + + let start = std::time::Instant::now(); + + // Lancer les deux pipelines en parallèle + let download_handle = tokio::spawn(async move { + tracing::info!("[DOWNLOAD] Pipeline starting..."); + let result = Box::new(download_source).run(stop_token_download).await; + match &result { + Ok(()) => tracing::info!("[DOWNLOAD] Pipeline completed successfully"), + Err(e) => tracing::error!("[DOWNLOAD] Pipeline error: {}", e), + } + result + }); + + let playback_handle = tokio::spawn(async move { + // Attendre un peu que le premier track soit disponible + tokio::time::sleep(tokio::time::Duration::from_secs(2)).await; + tracing::info!("[PLAYBACK] Pipeline starting..."); + let result = Box::new(playlist_source).run(stop_token_playback).await; + match &result { + Ok(()) => tracing::info!("[PLAYBACK] Pipeline completed successfully"), + Err(e) => tracing::error!("[PLAYBACK] Pipeline error: {}", e), + } + result + }); + + // Attendre les deux pipelines + let (download_result, playback_result) = tokio::join!(download_handle, playback_handle); + + let elapsed = start.elapsed(); + + // Vérifier les résultats + match (download_result, playback_result) { + (Ok(Ok(())), Ok(Ok(()))) => { + tracing::info!(""); + tracing::info!("========================================"); + tracing::info!("✓ Both pipelines completed successfully"); + tracing::info!(" Total time: {:.2}s", elapsed.as_secs_f64()); + tracing::info!("========================================"); + } + (download_res, playback_res) => { + tracing::error!(""); + tracing::error!("========================================"); + if let Err(e) = download_res { + tracing::error!("✗ Download pipeline error: {:?}", e); + } else if let Ok(Err(e)) = download_res { + tracing::error!("✗ Download pipeline error: {}", e); + } + if let Err(e) = playback_res { + tracing::error!("✗ Playback pipeline error: {:?}", e); + } else if let Ok(Err(e)) = playback_res { + tracing::error!("✗ Playback pipeline error: {}", e); + } + tracing::error!("========================================"); + return Err("Pipeline error".into()); + } + } + + Ok(()) +}