diff --git a/Cargo.lock b/Cargo.lock index 6b56e562..38c60dd8 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -6,7 +6,8 @@ version = 4 name = "PMOMusic" version = "0.1.0" dependencies = [ - "axum", + "axum 0.8.6", + "console-subscriber", "pmoapp", "pmoaudiocache", "pmoconfig", @@ -197,6 +198,33 @@ dependencies = [ "arrayvec", ] +[[package]] +name = "axum" +version = "0.7.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "edca88bc138befd0323b20752846e6587272d3b03b0343c8ea28a6f819e6e71f" +dependencies = [ + "async-trait", + "axum-core 0.4.5", + "bytes", + "futures-util", + "http", + "http-body", + "http-body-util", + "itoa", + "matchit 0.7.3", + "memchr", + "mime", + "percent-encoding", + "pin-project-lite", + "rustversion", + "serde", + "sync_wrapper", + "tower 0.5.2", + "tower-layer", + "tower-service", +] + [[package]] name = "axum" version = "0.8.6" @@ -213,7 +241,7 @@ dependencies = [ "hyper", "hyper-util", "itoa", - "matchit", + "matchit 0.8.4", "memchr", "mime", "percent-encoding", @@ -224,7 +252,7 @@ dependencies = [ "serde_urlencoded", "sync_wrapper", "tokio", - "tower", + "tower 0.5.2", "tower-layer", "tower-service", "tracing", @@ -314,6 +342,12 @@ dependencies = [ "windows-link 0.2.0", ] +[[package]] +name = "base64" +version = "0.21.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9d297deb1925b89f2ccc13d7635fa0714f12c87adce1c75356b39ca9b7178567" + [[package]] name = "base64" version = "0.22.1" @@ -392,7 +426,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d65b08cc2fb06b872143c506ee7e14decb3ade19bf25af9ddd39dfc1799ccc85" dependencies = [ "bevy_macro_utils", - "indexmap", + "indexmap 2.11.4", "proc-macro2", "quote", "syn 2.0.106", @@ -573,6 +607,45 @@ dependencies = [ "crossbeam-utils", ] +[[package]] +name = "console-api" +version = "0.8.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8030735ecb0d128428b64cd379809817e620a40e5001c54465b99ec5feec2857" +dependencies = [ + "futures-core", + "prost", + "prost-types", + "tonic", + "tracing-core", +] + +[[package]] +name = "console-subscriber" +version = "0.4.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6539aa9c6a4cd31f4b1c040f860a1eac9aa80e7df6b05d506a6e7179936d6a01" +dependencies = [ + "console-api", + "crossbeam-channel", + "crossbeam-utils", + "futures-task", + "hdrhistogram", + "humantime", + "hyper-util", + "prost", + "prost-types", + "serde", + "serde_json", + "thread_local", + "tokio", + "tokio-stream", + "tonic", + "tracing", + "tracing-core", + "tracing-subscriber", +] + [[package]] name = "cookie" version = "0.18.1" @@ -1298,7 +1371,7 @@ dependencies = [ "futures-core", "futures-sink", "http", - "indexmap", + "indexmap 2.11.4", "slab", "tokio", "tokio-util", @@ -1324,6 +1397,12 @@ dependencies = [ "byteorder", ] +[[package]] +name = "hashbrown" +version = "0.12.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8a9ee70c43aaf417c914396645a0fa852624801b24ebb7ae78fe8272889ac888" + [[package]] name = "hashbrown" version = "0.15.5" @@ -1352,6 +1431,19 @@ dependencies = [ "hashbrown 0.15.5", ] +[[package]] +name = "hdrhistogram" +version = "7.5.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "765c9198f173dd59ce26ff9f95ef0aafd0a0fe01fb9d72841bc5066a4c06511d" +dependencies = [ + "base64 0.21.7", + "byteorder", + "flate2", + "nom", + "num-traits", +] + [[package]] name = "heapless" version = "0.8.0" @@ -1439,6 +1531,12 @@ version = "1.0.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "df3b46402a9d5adb4c86a0cf463f42e19994e3ee891101b1841f30a545cb49a9" +[[package]] +name = "humantime" +version = "2.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "135b12329e5e3ce057a9f972339ea52bc954fe1e9358ef27f95e89716fbc5424" + [[package]] name = "hyper" version = "1.7.0" @@ -1478,6 +1576,19 @@ dependencies = [ "tower-service", ] +[[package]] +name = "hyper-timeout" +version = "0.5.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2b90d566bffbce6a75bd8b09a05aa8c2cb1fabb6cb348f8840c9e4c90a0d83b0" +dependencies = [ + "hyper", + "hyper-util", + "pin-project-lite", + "tokio", + "tower-service", +] + [[package]] name = "hyper-tls" version = "0.6.0" @@ -1500,7 +1611,7 @@ version = "0.1.17" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3c6995591a8f1380fcb4ba966a252a4b29188d51d2b89e3a252f5305be65aea8" dependencies = [ - "base64", + "base64 0.22.1", "bytes", "futures-channel", "futures-core", @@ -1512,7 +1623,7 @@ dependencies = [ "libc", "percent-encoding", "pin-project-lite", - "socket2", + "socket2 0.6.0", "system-configuration", "tokio", "tower-service", @@ -1691,6 +1802,16 @@ version = "1.12.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e7c5cedc30da3a610cac6b4ba17597bdf7152cf974e8aab3afb3d54455e371c8" +[[package]] +name = "indexmap" +version = "1.9.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bd070e393353796e801d209ad339e89596eb4c8d430d18ede6a1cced8fafbd99" +dependencies = [ + "autocfg", + "hashbrown 0.12.3", +] + [[package]] name = "indexmap" version = "2.11.4" @@ -1931,6 +2052,12 @@ dependencies = [ "regex-automata", ] +[[package]] +name = "matchit" +version = "0.7.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0e7465ac9959cc2b1404e8e2367b43684a6d13790fe23056cc8c6c5a6b7bcb94" + [[package]] name = "matchit" version = "0.8.4" @@ -2336,6 +2463,26 @@ version = "2.3.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9b4f627cb1b25917193a259e49bdad08f671f8d9708acfd5fe0a8c1455d87220" +[[package]] +name = "pin-project" +version = "1.1.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "677f1add503faace112b9f1373e43e9e054bfdd22ff1a63c1bc485eaec6a6a8a" +dependencies = [ + "pin-project-internal", +] + +[[package]] +name = "pin-project-internal" +version = "1.1.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6e918e4ff8c4549eb882f14b3a4bc8c8bc93de829416eacf579f1207a8fbf861" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.106", +] + [[package]] name = "pin-project-lite" version = "0.2.16" @@ -2360,8 +2507,8 @@ version = "1.8.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "740ebea15c5d1428f910cd1a5f52cebf8d25006245ed8ade92702f4943d91e07" dependencies = [ - "base64", - "indexmap", + "base64 0.22.1", + "indexmap 2.11.4", "quick-xml 0.38.3", "serde", "time", @@ -2389,7 +2536,7 @@ name = "pmoaudiocache" version = "0.1.0" dependencies = [ "anyhow", - "axum", + "axum 0.8.6", "chrono", "claxon", "flacenc 0.4.0", @@ -2416,7 +2563,7 @@ name = "pmocache" version = "0.1.0" dependencies = [ "anyhow", - "axum", + "axum 0.8.6", "bytes", "chrono", "futures-util", @@ -2440,7 +2587,7 @@ name = "pmoconfig" version = "0.1.0" dependencies = [ "anyhow", - "axum", + "axum 0.8.6", "dirs", "lazy_static", "log", @@ -2459,7 +2606,7 @@ name = "pmocovers" version = "0.1.0" dependencies = [ "anyhow", - "axum", + "axum 0.8.6", "image", "pmocache", "pmoconfig", @@ -2501,7 +2648,7 @@ name = "pmomediaserver" version = "0.1.0" dependencies = [ "async-trait", - "axum", + "axum 0.8.6", "bevy_reflect", "once_cell", "pmoconfig", @@ -2527,7 +2674,7 @@ dependencies = [ "anyhow", "async-stream", "async-trait", - "axum", + "axum 0.8.6", "bytes", "chrono", "claxon", @@ -2574,7 +2721,7 @@ name = "pmoqobuz" version = "0.1.0" dependencies = [ "anyhow", - "axum", + "axum 0.8.6", "chrono", "hex", "mockito", @@ -2603,7 +2750,7 @@ version = "0.1.0" dependencies = [ "anyhow", "async-stream", - "axum", + "axum 0.8.6", "axum-embed", "axum-server", "futures", @@ -2627,7 +2774,7 @@ version = "0.1.0" dependencies = [ "anyhow", "async-trait", - "axum", + "axum 0.8.6", "lazy_static", "pmoaudiocache", "pmoconfig", @@ -2649,8 +2796,8 @@ name = "pmoupnp" version = "0.1.0" dependencies = [ "anyhow", - "axum", - "base64", + "axum 0.8.6", + "base64 0.22.1", "bevy_reflect", "bevy_reflect_derive", "chrono", @@ -2769,6 +2916,38 @@ dependencies = [ "syn 2.0.106", ] +[[package]] +name = "prost" +version = "0.13.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2796faa41db3ec313a31f7624d9286acf277b52de526150b7e69f3debf891ee5" +dependencies = [ + "bytes", + "prost-derive", +] + +[[package]] +name = "prost-derive" +version = "0.13.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8a56d757972c98b346a9b766e3f02746cde6dd1cd1d1d563472929fdd74bec4d" +dependencies = [ + "anyhow", + "itertools", + "proc-macro2", + "quote", + "syn 2.0.106", +] + +[[package]] +name = "prost-types" +version = "0.13.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "52c2c1bf36ddb1a1c396b3601a3cec27c2462e45f07c386894ec3ccf5332bd16" +dependencies = [ + "prost", +] + [[package]] name = "psl-types" version = "2.0.11" @@ -3028,7 +3207,7 @@ version = "0.12.23" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d429f34c8092b2d42c7c93cec323bb4adeb7c67698f70839adec842ec10c7ceb" dependencies = [ - "base64", + "base64 0.22.1", "bytes", "cookie", "cookie_store", @@ -3058,7 +3237,7 @@ dependencies = [ "tokio", "tokio-native-tls", "tokio-util", - "tower", + "tower 0.5.2", "tower-http", "tower-service", "url", @@ -3349,7 +3528,7 @@ version = "0.9.34+deprecated" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6a8b1a1a2ebf674015cc02edccce75287f1a0130d394307b36743c2f5d504b47" dependencies = [ - "indexmap", + "indexmap 2.11.4", "itoa", "ryu", "serde", @@ -3444,6 +3623,16 @@ dependencies = [ "serde", ] +[[package]] +name = "socket2" +version = "0.5.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e22376abed350d73dd1cd119b57ffccad95b4e585a7cda43e286245ce23c0678" +dependencies = [ + "libc", + "windows-sys 0.52.0", +] + [[package]] name = "socket2" version = "0.6.0" @@ -3905,8 +4094,9 @@ dependencies = [ "pin-project-lite", "signal-hook-registry", "slab", - "socket2", + "socket2 0.6.0", "tokio-macros", + "tracing", "windows-sys 0.59.0", ] @@ -4014,7 +4204,7 @@ version = "0.22.27" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "41fe8c660ae4257887cf66394862d21dbca4a6ddd26f04a3560410406a2f819a" dependencies = [ - "indexmap", + "indexmap 2.11.4", "serde", "serde_spanned", "toml_datetime 0.6.11", @@ -4027,7 +4217,7 @@ version = "0.23.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f3effe7c0e86fdff4f69cdd2ccc1b96f933e24811c5441d44904e8683e27184b" dependencies = [ - "indexmap", + "indexmap 2.11.4", "toml_datetime 0.7.2", "toml_parser", "winnow", @@ -4042,6 +4232,56 @@ dependencies = [ "winnow", ] +[[package]] +name = "tonic" +version = "0.12.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "877c5b330756d856ffcc4553ab34a5684481ade925ecc54bcd1bf02b1d0d4d52" +dependencies = [ + "async-stream", + "async-trait", + "axum 0.7.9", + "base64 0.22.1", + "bytes", + "h2", + "http", + "http-body", + "http-body-util", + "hyper", + "hyper-timeout", + "hyper-util", + "percent-encoding", + "pin-project", + "prost", + "socket2 0.5.10", + "tokio", + "tokio-stream", + "tower 0.4.13", + "tower-layer", + "tower-service", + "tracing", +] + +[[package]] +name = "tower" +version = "0.4.13" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b8fa9be0de6cf49e536ce1851f987bd21a43b771b09473c3549a6c853db37c1c" +dependencies = [ + "futures-core", + "futures-util", + "indexmap 1.9.3", + "pin-project", + "pin-project-lite", + "rand 0.8.5", + "slab", + "tokio", + "tokio-util", + "tower-layer", + "tower-service", + "tracing", +] + [[package]] name = "tower" version = "0.5.2" @@ -4071,7 +4311,7 @@ dependencies = [ "http-body", "iri-string", "pin-project-lite", - "tower", + "tower 0.5.2", "tower-layer", "tower-service", ] @@ -4226,7 +4466,7 @@ version = "5.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2fcc29c80c21c31608227e0912b2d7fddba57ad76b606890627ba8ee7964e993" dependencies = [ - "indexmap", + "indexmap 2.11.4", "serde", "serde_json", "utoipa-gen", @@ -4250,8 +4490,8 @@ version = "9.0.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d047458f1b5b65237c2f6dc6db136945667f40a7668627b3490b9513a3d43a55" dependencies = [ - "axum", - "base64", + "axum 0.8.6", + "base64 0.22.1", "mime_guess", "regex", "rust-embed", @@ -4753,7 +4993,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "08db1edfb05d9b3c1542e521aea074442088292f00b5f28e435c714a98f85031" dependencies = [ "assert-json-diff", - "base64", + "base64 0.22.1", "deadpool", "futures", "http", @@ -4909,7 +5149,7 @@ dependencies = [ "arbitrary", "crc32fast", "flate2", - "indexmap", + "indexmap 2.11.4", "memchr", "zopfli", ] diff --git a/PMOMusic/Cargo.toml b/PMOMusic/Cargo.toml index e175068c..cda6d03a 100644 --- a/PMOMusic/Cargo.toml +++ b/PMOMusic/Cargo.toml @@ -20,3 +20,4 @@ tracing-subscriber = "0.3.20" axum = "0.8.4" serde_json = "1.0.145" utoipa = "5.4" +console-subscriber = "0.4.1" diff --git a/PMOMusic/src/main.rs b/PMOMusic/src/main.rs index f2c4afdd..f17b1220 100644 --- a/PMOMusic/src/main.rs +++ b/PMOMusic/src/main.rs @@ -9,8 +9,13 @@ use tracing::info; #[tokio::main] async fn main() -> Result<(), Box> { // ========== PHASE 1 : Infrastructure UPnP ========== -let server = Server::create_upnp_server().await?; // Routes personnalisées de l'application - server.write().await + // #[cfg(tokio_unstable)] + // console_subscriber::init(); + + let server = Server::create_upnp_server().await?; // Routes personnalisées de l'application + server + .write() + .await .add_route("/info", || async { serde_json::json!({"version": "1.0.0"}) }) @@ -18,7 +23,9 @@ let server = Server::create_upnp_server().await?; // Routes personnalisées d // Initialiser le système de gestion des sources musicales avec API REST info!("📡 Initializing music sources management system..."); - server.write().await + server + .write() + .await .init_music_sources() .await .expect("Failed to initialize music sources API"); @@ -48,7 +55,9 @@ let server = Server::create_upnp_server().await?; // Routes personnalisées d // Enregistrer les devices UPnP (HTTP + SSDP automatique) info!("📡 Registering UPnP devices..."); - let renderer_instance = server.write().await + let renderer_instance = server + .write() + .await .register_device(MEDIA_RENDERER.clone()) .await .expect("Failed to register MediaRenderer"); @@ -59,7 +68,9 @@ let server = Server::create_upnp_server().await?; // Routes personnalisées d renderer_instance.description_route() ); - let server_instance = server.write().await + let server_instance = server + .write() + .await .register_device(MEDIA_SERVER.clone()) .await .expect("Failed to register MediaServer"); @@ -72,7 +83,11 @@ let server = Server::create_upnp_server().await?; // Routes personnalisées d // Ajouter la webapp via le trait WebAppExt info!("📡 Registering Web application..."); - server.write().await.add_webapp_with_redirect::("/app").await; + server + .write() + .await + .add_webapp_with_redirect::("/app") + .await; // ========== PHASE 3 : Démarrage du serveur ========== diff --git a/pmocache/src/db.rs b/pmocache/src/db.rs index 24b3a92b..657a11ae 100644 --- a/pmocache/src/db.rs +++ b/pmocache/src/db.rs @@ -79,9 +79,8 @@ impl<'a> std::ops::DerefMut for ConnGuard<'a> { } } - impl DB { - fn lock_conn(&self, ctx: &'static str) -> ConnGuard<'_> { + fn lock_conn(&self, ctx: &'static str) -> ConnGuard<'_> { trace!("DB mutex → acquiring ({ctx})"); let start = Instant::now(); let guard = self.conn.lock().unwrap(); diff --git a/pmoconfig/src/lib.rs b/pmoconfig/src/lib.rs index 5f6699a8..e7be4059 100644 --- a/pmoconfig/src/lib.rs +++ b/pmoconfig/src/lib.rs @@ -464,7 +464,7 @@ impl Config { _ => { self.set_managed_dir(path, default.to_string())?; default.to_string() - }, + } }; self.resolve_and_create_dir(&dir_path) } diff --git a/pmoparadise/src/paradise/history.rs b/pmoparadise/src/paradise/history.rs index 279d6819..c276a810 100644 --- a/pmoparadise/src/paradise/history.rs +++ b/pmoparadise/src/paradise/history.rs @@ -94,7 +94,6 @@ impl SqliteHistoryBackend { } } - #[async_trait] impl HistoryBackend for SqliteHistoryBackend { async fn append(&self, entry: HistoryEntry) -> anyhow::Result<()> { diff --git a/pmoparadise/src/paradise/mod.rs b/pmoparadise/src/paradise/mod.rs index a8071a58..f2860671 100644 --- a/pmoparadise/src/paradise/mod.rs +++ b/pmoparadise/src/paradise/mod.rs @@ -22,7 +22,7 @@ pub use channel::{ max_channel_id, ChannelDescriptor, ParadiseChannel, ParadiseChannelKind, ParadiseClientStream, ALL_CHANNELS, }; -pub use constants::*; // Export all constants +pub use constants::*; // Export all constants pub use history::{create_history_backend, HistoryBackend, HistoryEntry}; pub use playlist::PlaylistEntry; pub use worker::{load_rp_metadata, ParadiseWorker, RadioParadiseMetadata, WorkerCommand}; diff --git a/pmoparadise/src/paradise/worker.rs b/pmoparadise/src/paradise/worker.rs index fbf7ab2d..5823204e 100644 --- a/pmoparadise/src/paradise/worker.rs +++ b/pmoparadise/src/paradise/worker.rs @@ -56,8 +56,14 @@ impl ParadiseWorker { let join_handle = tokio::spawn(async move { info!(channel = descriptor.slug, "Starting Radio Paradise worker"); - let mut state = - WorkerState::new(descriptor, client, history_max_tracks, playlist, history, cache_manager); + let mut state = WorkerState::new( + descriptor, + client, + history_max_tracks, + playlist, + history, + cache_manager, + ); loop { if let Some(task) = state.scheduled_task.as_mut() { @@ -155,7 +161,24 @@ struct WorkerState { shutdown: bool, } +#[derive(Clone)] +struct SongTaskContext { + cache_manager: Arc, + playlist: SharedPlaylist, + descriptor_id: u8, + slug: &'static str, +} + impl WorkerState { + fn song_task_context(&self) -> SongTaskContext { + SongTaskContext { + cache_manager: Arc::clone(&self.cache_manager), + playlist: self.playlist.clone(), + descriptor_id: self.descriptor.id, + slug: self.descriptor.slug, + } + } + fn new( descriptor: ChannelDescriptor, client: RadioParadiseClient, @@ -452,27 +475,19 @@ impl WorkerState { track_samples.len() ); - // Encode and cache the song - let entry = self - .process_song_from_pcm( - &block, - song_index, - song, - track_samples, - sample_rate, - channels as usize, - bits_per_sample, - ) - .await?; - - self.playlist.push_active(entry).await; - - info!( - channel = self.descriptor.slug, - song_index = song_index, - "🎵 Song '{}' available after {}ms (streaming mode)", - song.title, - current_position_ms + let context = self.song_task_context(); + spawn_song_processing( + context, + block.clone(), + song_index, + song.clone(), + track_samples, + sample_rate, + channels as usize, + bits_per_sample, + self.active_clients, + song.duration, + current_position_ms, ); current_song_idx += 1; @@ -508,19 +523,20 @@ impl WorkerState { track_samples.len() ); - let entry = self - .process_song_from_pcm( - &block, - song_index, - song, - track_samples, - sample_rate, - channels as usize, - bits_per_sample, - ) - .await?; - - self.playlist.push_active(entry).await; + let context = self.song_task_context(); + spawn_song_processing( + context, + block.clone(), + song_index, + song.clone(), + track_samples, + sample_rate, + channels as usize, + bits_per_sample, + self.active_clients, + song.duration, + song_start_ms, + ); } } @@ -567,188 +583,21 @@ impl WorkerState { .ok_or_else(|| anyhow!("Sample slice out of bounds"))?; let track_samples = slice.to_vec(); - let flac_bytes = encode_samples_to_flac( + encode_song_to_cache( + Arc::clone(&self.cache_manager), + self.descriptor.id, + self.descriptor.slug, + block.clone(), + *song_index, + song.clone(), track_samples, - decoded.channels, decoded.sample_rate, + decoded.channels, decoded.bits_per_sample, + self.active_clients, + duration_ms, ) .await - .context("Failed to encode song to FLAC")?; - - let track_id = self.compute_track_id(block, *song_index); - let placeholder_uri = format!("{}#{}", block.url, song_index); - - let mut metadata = TrackMetadata { - original_uri: placeholder_uri.clone(), - cached_audio_pk: None, - cached_cover_pk: None, - }; - - if let Some(cover_pk) = self.cache_cover(block, song).await? { - metadata.cached_cover_pk = Some(cover_pk); - } - - let flac_len = flac_bytes.len() as u64; - let reader = StreamReader::new(stream::iter(vec![Ok::<_, std::io::Error>(Bytes::from( - flac_bytes, - ))])); - - let audio_pk = self - .cache_manager - .cache_audio_from_reader(&track_id, reader, Some(flac_len)) - .await - .map_err(|e| anyhow!("Cache audio error: {e}"))?; - - // Wait for the file to be completely written before continuing - self.cache_manager - .wait_audio_ready(&audio_pk) - .await - .map_err(|e| anyhow!("Wait audio ready error: {e}"))?; - - metadata.cached_audio_pk = Some(audio_pk.clone()); - self.cache_manager - .update_metadata(track_id.clone(), metadata) - .await; - - let file_path = self.cache_manager.audio_file_path(&audio_pk).await; - - let entry = Arc::new(PlaylistEntry::new( - track_id, - self.descriptor.id, - Arc::new(song.clone()), - Utc::now(), - duration_ms, - Some(audio_pk), - file_path, - self.active_clients, - )); - - Ok(entry) - } - - /// Process a song from pre-decoded PCM samples (for streaming mode) - /// - /// This is a variant of `process_song()` that takes PCM samples directly - /// instead of slicing from a DecodedBlock. Used for progressive streaming - /// where samples are decoded on-the-fly. - /// - /// # Arguments - /// - /// * `block` - The block metadata - /// * `song_index` - Index of the song in the block - /// * `song` - The song metadata - /// * `track_samples` - Pre-decoded and sliced PCM samples (interleaved i32) - /// * `sample_rate` - Sample rate (e.g., 44100) - /// * `channels` - Number of channels (e.g., 2 for stereo) - /// * `bits_per_sample` - Bits per sample (e.g., 16) - async fn process_song_from_pcm( - &self, - block: &Block, - song_index: usize, - song: &Song, - track_samples: Vec, - sample_rate: u32, - channels: usize, - bits_per_sample: u32, - ) -> Result> { - let duration_ms = song.duration; - - // Encode PCM to FLAC - let flac_bytes = - encode_samples_to_flac(track_samples, channels, sample_rate, bits_per_sample) - .await - .context("Failed to encode song to FLAC")?; - - let track_id = self.compute_track_id(block, song_index); - let placeholder_uri = format!("{}#{}", block.url, song_index); - - let mut metadata = TrackMetadata { - original_uri: placeholder_uri.clone(), - cached_audio_pk: None, - cached_cover_pk: None, - }; - - // Cache cover art - if let Some(cover_pk) = self.cache_cover(block, song).await? { - metadata.cached_cover_pk = Some(cover_pk); - } - - // Cache audio - let flac_len = flac_bytes.len() as u64; - let reader = StreamReader::new(stream::iter(vec![Ok::<_, std::io::Error>(Bytes::from( - flac_bytes, - ))])); - - let audio_pk = self - .cache_manager - .cache_audio_from_reader(&track_id, reader, Some(flac_len)) - .await - .map_err(|e| anyhow!("Cache audio error: {e}"))?; - - // Wait for the file to be completely written before continuing - self.cache_manager - .wait_audio_ready(&audio_pk) - .await - .map_err(|e| anyhow!("Wait audio ready error: {e}"))?; - - metadata.cached_audio_pk = Some(audio_pk.clone()); - self.cache_manager - .update_metadata(track_id.clone(), metadata.clone()) - .await; - - // Stocker les métadonnées Radio Paradise dans le cache - if let Err(e) = Self::store_rp_metadata( - &self.cache_manager, - &audio_pk, - &track_id, - self.descriptor.id, - song, - duration_ms, - block.event, - metadata.cached_cover_pk.as_deref(), - ) - .await - { - warn!( - channel = self.descriptor.slug, - "Failed to store RP metadata: {e:?}" - ); - } - - let file_path = self.cache_manager.audio_file_path(&audio_pk).await; - - let entry = Arc::new(PlaylistEntry::new( - track_id, - self.descriptor.id, - Arc::new(song.clone()), - Utc::now(), - duration_ms, - Some(audio_pk), - file_path, - self.active_clients, - )); - - Ok(entry) - } - - async fn cache_cover(&self, block: &Block, song: &Song) -> Result> { - if let Some(ref cover_path) = song.cover { - if let Some(cover_url) = block.cover_url(cover_path) { - match self.cache_manager.cache_cover(&cover_url).await { - Ok(pk) => return Ok(Some(pk)), - Err(err) => { - warn!(channel = self.descriptor.slug, "Cover cache error: {err}"); - } - } - } else { - warn!( - channel = self.descriptor.slug, - "Unable to resolve cover URL for {}", cover_path - ); - } - } - Ok(None) } /// Stocke les métadonnées Radio Paradise pour un fichier audio caché @@ -756,48 +605,8 @@ impl WorkerState { /// Cette fonction persiste toutes les métadonnées RP dans la base de données /// du cache audio, permettant leur récupération future sans dépendance aux /// données en mémoire. - async fn store_rp_metadata( - cache_manager: &SourceCacheManager, - audio_pk: &str, - track_id: &str, - channel_id: u8, - song: &Song, - duration_ms: u64, - event: u64, - cover_pk: Option<&str>, - ) -> Result<()> { - use serde_json::json; - - // Métadonnées basiques de la chanson - cache_manager.set_audio_metadata(audio_pk, "rp_title", json!(song.title))?; - cache_manager.set_audio_metadata(audio_pk, "rp_artist", json!(song.artist))?; - cache_manager.set_audio_metadata(audio_pk, "rp_album", json!(song.album))?; - cache_manager.set_audio_metadata(audio_pk, "rp_year", json!(song.year))?; - - // Informations temporelles - cache_manager.set_audio_metadata(audio_pk, "rp_duration_ms", json!(duration_ms))?; - cache_manager.set_audio_metadata(audio_pk, "rp_elapsed_ms", json!(song.elapsed))?; - - // Identifiants Radio Paradise - cache_manager.set_audio_metadata(audio_pk, "rp_track_id", json!(track_id))?; - cache_manager.set_audio_metadata(audio_pk, "rp_channel_id", json!(channel_id))?; - cache_manager.set_audio_metadata(audio_pk, "rp_event", json!(event))?; - - // Métadonnées supplémentaires - cache_manager.set_audio_metadata(audio_pk, "rp_rating", json!(song.rating))?; - cache_manager.set_audio_metadata(audio_pk, "rp_cover_url", json!(song.cover))?; - cache_manager.set_audio_metadata(audio_pk, "rp_cover_pk", json!(cover_pk))?; - - Ok(()) - } - fn compute_track_id(&self, block: &Block, song_index: usize) -> String { - // Use deterministic ID based on block event and song index - // This allows checking if a song is cached before downloading the block - format!( - "rp:{}:event_{}_song_{}", - self.descriptor.id, block.event, song_index - ) + compute_track_id_for_descriptor(self.descriptor.id, block, song_index) } async fn maybe_schedule_poll(&mut self) { @@ -1072,7 +881,7 @@ fn ms_to_frames(ms: u64, sample_rate: u32) -> usize { fn decode_block_audio(data: Vec) -> anyhow::Result { use symphonia::core::audio::SampleBuffer; - use symphonia::core::codecs::{DecoderOptions, CODEC_TYPE_NULL}; + use symphonia::core::codecs::{CODEC_TYPE_NULL, DecoderOptions}; use symphonia::core::errors::Error as SymphoniaError; use symphonia::core::formats::FormatOptions; use symphonia::core::io::MediaSourceStream; @@ -1174,6 +983,213 @@ fn decode_block_audio(data: Vec) -> anyhow::Result { }) } +fn compute_track_id_for_descriptor(descriptor_id: u8, block: &Block, song_index: usize) -> String { + format!( + "rp:{}:event_{}_song_{}", + descriptor_id, block.event, song_index + ) +} + +async fn store_rp_metadata( + cache_manager: &SourceCacheManager, + audio_pk: &str, + track_id: &str, + channel_id: u8, + song: &Song, + duration_ms: u64, + event: u64, + cover_pk: Option<&str>, +) -> Result<()> { + use serde_json::json; + + cache_manager.set_audio_metadata(audio_pk, "rp_title", json!(song.title))?; + cache_manager.set_audio_metadata(audio_pk, "rp_artist", json!(song.artist))?; + cache_manager.set_audio_metadata(audio_pk, "rp_album", json!(song.album))?; + cache_manager.set_audio_metadata(audio_pk, "rp_year", json!(song.year))?; + + cache_manager.set_audio_metadata(audio_pk, "rp_duration_ms", json!(duration_ms))?; + cache_manager.set_audio_metadata(audio_pk, "rp_elapsed_ms", json!(song.elapsed))?; + + cache_manager.set_audio_metadata(audio_pk, "rp_track_id", json!(track_id))?; + cache_manager.set_audio_metadata(audio_pk, "rp_channel_id", json!(channel_id))?; + cache_manager.set_audio_metadata(audio_pk, "rp_event", json!(event))?; + + cache_manager.set_audio_metadata(audio_pk, "rp_rating", json!(song.rating))?; + cache_manager.set_audio_metadata(audio_pk, "rp_cover_url", json!(song.cover))?; + cache_manager.set_audio_metadata(audio_pk, "rp_cover_pk", json!(cover_pk))?; + + Ok(()) +} + +async fn cache_cover_for_song( + cache_manager: &SourceCacheManager, + slug: &'static str, + block: &Block, + song: &Song, +) -> Result> { + if let Some(ref cover_path) = song.cover { + if let Some(cover_url) = block.cover_url(cover_path) { + match cache_manager.cache_cover(&cover_url).await { + Ok(pk) => return Ok(Some(pk)), + Err(err) => { + warn!(channel = slug, "Cover cache error: {err}"); + } + } + } else { + warn!( + channel = slug, + "Unable to resolve cover URL for {}", cover_path + ); + } + } + Ok(None) +} + +async fn encode_song_to_cache( + cache_manager: Arc, + descriptor_id: u8, + slug: &'static str, + block: Block, + song_index: usize, + song: Song, + track_samples: Vec, + sample_rate: u32, + channels: usize, + bits_per_sample: u32, + active_clients: usize, + duration_ms: u64, +) -> Result> { + let flac_bytes = encode_samples_to_flac(track_samples, channels, sample_rate, bits_per_sample) + .await + .context("Failed to encode song to FLAC")?; + + let track_id = compute_track_id_for_descriptor(descriptor_id, &block, song_index); + let placeholder_uri = format!("{}#{}", block.url, song_index); + + let mut metadata = TrackMetadata { + original_uri: placeholder_uri.clone(), + cached_audio_pk: None, + cached_cover_pk: None, + }; + + if let Some(cover_pk) = + cache_cover_for_song(cache_manager.as_ref(), slug, &block, &song).await? + { + metadata.cached_cover_pk = Some(cover_pk); + } + + let flac_len = flac_bytes.len() as u64; + let reader = StreamReader::new(stream::iter(vec![Ok::<_, std::io::Error>(Bytes::from( + flac_bytes, + ))])); + + let audio_pk = cache_manager + .cache_audio_from_reader(&track_id, reader, Some(flac_len)) + .await + .map_err(|e| anyhow!("Cache audio error: {e}"))?; + + cache_manager + .wait_audio_ready(&audio_pk) + .await + .map_err(|e| anyhow!("Wait audio ready error: {e}"))?; + + metadata.cached_audio_pk = Some(audio_pk.clone()); + cache_manager + .update_metadata(track_id.clone(), metadata.clone()) + .await; + + if let Err(e) = store_rp_metadata( + cache_manager.as_ref(), + &audio_pk, + &track_id, + descriptor_id, + &song, + duration_ms, + block.event, + metadata.cached_cover_pk.as_deref(), + ) + .await + { + warn!(channel = slug, "Failed to store RP metadata: {e:?}"); + } + + let file_path = cache_manager.audio_file_path(&audio_pk).await; + + let entry = Arc::new(PlaylistEntry::new( + track_id, + descriptor_id, + Arc::new(song.clone()), + Utc::now(), + duration_ms, + Some(audio_pk), + file_path, + active_clients, + )); + + Ok(entry) +} + +fn spawn_song_processing( + context: SongTaskContext, + block: Block, + song_index: usize, + song: Song, + track_samples: Vec, + sample_rate: u32, + channels: usize, + bits_per_sample: u32, + active_clients: usize, + duration_ms: u64, + position_ms: u64, +) { + tokio::spawn(async move { + let SongTaskContext { + cache_manager, + playlist, + descriptor_id, + slug, + } = context; + + let song_title = song.title.clone(); + + match encode_song_to_cache( + cache_manager, + descriptor_id, + slug, + block, + song_index, + song, + track_samples, + sample_rate, + channels, + bits_per_sample, + active_clients, + duration_ms, + ) + .await + { + Ok(entry) => { + playlist.push_active(entry).await; + info!( + channel = slug, + song_index = song_index, + "🎵 Song '{}' available after {}ms (streaming mode)", + song_title, + position_ms + ); + } + Err(err) => { + warn!( + channel = slug, + song_index = song_index, + "Failed to process song '{}' asynchronously: {err:?}", + song_title + ); + } + } + }); +} + async fn encode_samples_to_flac( samples: Vec, channels: usize, @@ -1185,6 +1201,9 @@ async fn encode_samples_to_flac( use flacenc::component::BitRepr; use flacenc::error::Verify; + // Note: Claxon retourne les samples dans leur résolution native + // Un fichier FLAC 16 bits retourne des samples i32 avec des valeurs dans la plage i16 + // Pas besoin de normalisation supplémentaire let config = flacenc::config::Encoder::default() .into_verified() .map_err(|e| anyhow!("FLAC config error: {e:?}"))?; diff --git a/pmoparadise/src/pmoserver_ext.rs b/pmoparadise/src/pmoserver_ext.rs index 7ddf4f17..f21d7f16 100644 --- a/pmoparadise/src/pmoserver_ext.rs +++ b/pmoparadise/src/pmoserver_ext.rs @@ -555,7 +555,7 @@ async fn get_channel_status( last_change, history_entries: history_len, history_max_tracks: channel.history_max_tracks(), - configured: true, // All channels are always available + configured: true, // All channels are always available cache_collection_id: cache_stats.collection_id, cache_total_tracks: cache_stats.total_tracks, cache_cached_tracks: cache_stats.cached_tracks, diff --git a/pmoparadise/src/source.rs b/pmoparadise/src/source.rs index 17b4b2a4..99c69616 100644 --- a/pmoparadise/src/source.rs +++ b/pmoparadise/src/source.rs @@ -104,7 +104,10 @@ impl RadioParadiseSource { path.push("paradise"); std::fs::create_dir_all(&path).ok(); path.push("history.db"); - (path.to_string_lossy().to_string(), HISTORY_DEFAULT_MAX_TRACKS) + ( + path.to_string_lossy().to_string(), + HISTORY_DEFAULT_MAX_TRACKS, + ) }; let history_backend = create_history_backend(&database_path).map_err(|e| { @@ -175,7 +178,10 @@ impl RadioParadiseSource { path.push("paradise"); std::fs::create_dir_all(&path).ok(); path.push("history.db"); - (path.to_string_lossy().to_string(), HISTORY_DEFAULT_MAX_TRACKS) + ( + path.to_string_lossy().to_string(), + HISTORY_DEFAULT_MAX_TRACKS, + ) }; let history_backend: Arc = diff --git a/pmoparadise/src/streaming.rs b/pmoparadise/src/streaming.rs index 9d45cea7..b38e9be5 100644 --- a/pmoparadise/src/streaming.rs +++ b/pmoparadise/src/streaming.rs @@ -6,7 +6,7 @@ use std::pin::Pin; use std::sync::mpsc::{sync_channel, Receiver, RecvError, SyncSender}; use std::time::{Duration, Instant}; -const CHANNEL_BUFFER_SIZE: usize = 16; +const CHANNEL_BUFFER_SIZE: usize = 64; // Augmenté de 16 à 64 pour réduire les warnings "buffer plein" pub const CHUNK_SIZE_FRAMES: usize = 4096; pub struct ChannelReader { @@ -67,11 +67,11 @@ impl Read for ChannelReader { let to_copy = available.min(buf.len()); buf[..to_copy].copy_from_slice(&chunk[self.position..self.position + to_copy]); self.position += to_copy; - tracing::trace!( - "ChannelReader copied {} bytes (elapsed {:?})", - to_copy, - start.elapsed() - ); + tracing::trace!( + "ChannelReader copied {} bytes (elapsed {:?})", + to_copy, + start.elapsed() + ); return Ok(to_copy); } } @@ -91,10 +91,7 @@ impl Read for ChannelReader { return Err(io::Error::new(io::ErrorKind::Other, e)); } Err(RecvError) => { - tracing::trace!( - "ChannelReader stream closed after {:?}", - start.elapsed() - ); + tracing::trace!("ChannelReader stream closed after {:?}", start.elapsed()); return Ok(0); } } @@ -179,12 +176,23 @@ impl StreamingPCMDecoder { Err(e) => return Err(anyhow::anyhow!("FLAC decode error: {}", e)), }; - let samples: Vec = frame.into_buffer(); - if samples.is_empty() { + let planar_samples: Vec = frame.into_buffer(); + if planar_samples.is_empty() { self.done = true; return Ok(None); } + // IMPORTANT: Claxon retourne les samples en format PLANAR (tous les L, puis tous les R) + // Mais nous avons besoin du format INTERLEAVED (L, R, L, R, ...) pour l'encodage + let block_size = planar_samples.len() / self.channels as usize; + let mut samples = Vec::with_capacity(planar_samples.len()); + + for i in 0..block_size { + for ch in 0..self.channels as usize { + samples.push(planar_samples[ch * block_size + i]); + } + } + let position_ms = { let frames = self.total_samples_decoded / self.channels as u64; (frames * 1000) / self.sample_rate as u64 diff --git a/pmoserver/src/logs/mod.rs b/pmoserver/src/logs/mod.rs index e0bae9cd..919a9fa2 100644 --- a/pmoserver/src/logs/mod.rs +++ b/pmoserver/src/logs/mod.rs @@ -299,16 +299,25 @@ pub fn init_logging() -> LogState { }; if enable_console { - subscriber - .with( - tracing_subscriber::fmt::layer() - .with_target(true) - .with_level(true) - .with_ansi(true), - ) - .init(); + let subscriber = subscriber.with( + tracing_subscriber::fmt::layer() + .with_target(true) + .with_level(true) + .with_ansi(true), + ); + if let Err(e) = subscriber.try_init() { + eprintln!( + "⚠️ tracing subscriber already initialised, skipping console layer: {}", + e + ); + } } else { - subscriber.init(); + if let Err(e) = subscriber.try_init() { + eprintln!( + "⚠️ tracing subscriber already initialised, skipping default registry: {}", + e + ); + } } log_state diff --git a/pmosource/src/cache.rs b/pmosource/src/cache.rs index eb6ab469..09078879 100644 --- a/pmosource/src/cache.rs +++ b/pmosource/src/cache.rs @@ -289,12 +289,7 @@ impl SourceCacheManager { /// cache_manager.set_audio_metadata(audio_pk, "rating", json!(8.5)).unwrap(); /// # } /// ``` - pub fn set_audio_metadata( - &self, - audio_pk: &str, - key: &str, - value: JsonValue, - ) -> Result<()> { + pub fn set_audio_metadata(&self, audio_pk: &str, key: &str, value: JsonValue) -> Result<()> { self.audio_cache .db .set_a_metadata(audio_pk, key, value)