diff --git a/Cargo.lock b/Cargo.lock index c8c8780d..4a5052d5 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -11,6 +11,7 @@ dependencies = [ "pmoconfig", "pmocovers", "pmomediarenderer", + "pmomediaserver", "pmoserver", "pmoupnp", "serde_json", @@ -2348,10 +2349,16 @@ dependencies = [ name = "pmomediaserver" version = "0.1.0" dependencies = [ + "async-trait", "bevy_reflect", "once_cell", "pmodidl", + "pmoserver", + "pmosource", "pmoupnp", + "quick-xml 0.38.3", + "tokio", + "tracing", ] [[package]] diff --git a/PMOMusic/Cargo.toml b/PMOMusic/Cargo.toml index 9473a309..b3326db2 100644 --- a/PMOMusic/Cargo.toml +++ b/PMOMusic/Cargo.toml @@ -7,6 +7,7 @@ edition = "2024" pmoconfig = { path = "../pmoconfig" } pmoupnp = { path = "../pmoupnp"} pmomediarenderer = { path = "../pmomediarenderer" } +pmomediaserver = { path = "../pmomediaserver"} pmoserver = { path = "../pmoserver" } pmocovers = { path = "../pmocovers", features = ["pmoserver"] } pmoapp = { path = "../pmoapp", features = ["pmoserver"] } diff --git a/PMOMusic/src/main.rs b/PMOMusic/src/main.rs index d6402186..86515898 100644 --- a/PMOMusic/src/main.rs +++ b/PMOMusic/src/main.rs @@ -1,15 +1,9 @@ -use pmoupnp::{ - ssdp::SsdpServer, - upnp_api::UpnpApiExt, - UpnpServer, -}; -use pmomediarenderer::MEDIA_RENDERER; -use pmoserver::{ - logs::LoggingOptions, - ServerBuilder -}; -use pmoapp::{Webapp, WebAppExt}; +use pmoapp::{WebAppExt, Webapp}; use pmocovers::CoverCacheExt; +use pmomediarenderer::MEDIA_RENDERER; +use pmomediaserver::{MEDIA_SERVER, MediaServerExt}; +use pmoserver::{ServerBuilder, logs::LoggingOptions}; +use pmoupnp::{UpnpServer, ssdp::SsdpServer, upnp_api::UpnpApiExt}; use tracing::info; #[tokio::main] @@ -20,17 +14,13 @@ async fn main() { // Initialiser le logging et enregistrer les routes de logs server.init_logging().await; - info!("📡 Registering the cover cache..."); - let cache = server.init_cover_cache_configured() - .await - .expect("Cannot initialise the image cache"); - - info!("✅ Cover cache ready at {}", - cache.cache_dir(), - ); - + let cache = server + .init_cover_cache_configured() + .await + .expect("Cannot initialise the image cache"); + info!("✅ Cover cache ready at {}", cache.cache_dir(),); // Routes de base server @@ -39,7 +29,6 @@ async fn main() { }) .await; - // Ajouter la webapp via le trait WebAppExt info!("📡 Registering Web application..."); server.add_webapp_with_redirect::("/app").await; @@ -48,23 +37,58 @@ async fn main() { server.register_upnp_api().await; info!("📡 Registering MediaRenderer..."); - let renderer_instance = server.register_device(MEDIA_RENDERER.clone()) + let renderer_instance = server + .register_device(MEDIA_RENDERER.clone()) .await .expect("Failed to register MediaRenderer routes"); - info!("✅ MediaRenderer ready at {}{}", + info!( + "✅ MediaRenderer ready at {}{}", renderer_instance.base_url(), renderer_instance.description_route() ); + // TODO: Enregistrer les sources musicales + // Exemple d'utilisation du MediaServerExt: + // + // use std::sync::Arc; + // + // info!("📡 Registering music sources..."); + // + // // Exemple: Enregistrer une source Qobuz + // // let qobuz = Arc::new(QobuzSource::new(credentials)); + // // server.register_music_source(qobuz).await; + // + // // Exemple: Enregistrer une source Radio Paradise + // // let radio = Arc::new(RadioParadiseSource::new()); + // // server.register_music_source(radio).await; + // + // // Lister toutes les sources enregistrées + // let sources = server.list_music_sources().await; + // info!("✅ {} music source(s) registered", sources.len()); + // for source in sources { + // info!(" - {} ({})", source.name(), source.id()); + // } + + info!("📡 Registering MediaServer..."); + let server_instance = server + .register_device(MEDIA_SERVER.clone()) + .await + .expect("Failed to register MediaServer routes"); + + info!( + "✅ MediaServer ready at {}{}", + server_instance.base_url(), + server_instance.description_route() + ); + // Créer et démarrer le serveur SSDP info!("📡 Starting SSDP discovery..."); let mut ssdp_server = SsdpServer::new(); ssdp_server.start().expect("Failed to start SSDP server"); // Créer et enregistrer le device SSDP pour le MediaRenderer - let ssdp_device = renderer_instance - .to_ssdp_device("PMOMusic", "1.0"); + let ssdp_device = renderer_instance.to_ssdp_device("PMOMusic", "1.0"); ssdp_server.add_device(ssdp_device); info!("✅ SSDP announcements sent for MediaRenderer"); diff --git a/pmomediaserver/Cargo.toml b/pmomediaserver/Cargo.toml index 48103974..ffcc50ef 100644 --- a/pmomediaserver/Cargo.toml +++ b/pmomediaserver/Cargo.toml @@ -6,6 +6,12 @@ edition = "2024" [dependencies] pmoupnp = { path = "../pmoupnp" } pmodidl = { path = "../pmodidl" } +pmosource = { path = "../pmosource" } +pmoserver = { path = "../pmoserver" } once_cell = "1.20" bevy_reflect = "0.17.1" +tokio = { version = "1", features = ["sync"] } +async-trait = "0.1" +tracing = "0.1" +quick-xml = { version = "0.38.3", features = ["serialize"] } diff --git a/pmomediaserver/src/content_handler.rs b/pmomediaserver/src/content_handler.rs new file mode 100644 index 00000000..3144ce5a --- /dev/null +++ b/pmomediaserver/src/content_handler.rs @@ -0,0 +1,426 @@ +//! # ContentDirectory Handler - Gestionnaire du service ContentDirectory +//! +//! Ce module implémente la logique métier du service ContentDirectory en intégrant +//! les sources musicales enregistrées dans le registre. +//! +//! ## Fonctionnalités +//! +//! - **Navigation multi-sources** : Combine toutes les sources dans une hiérarchie +//! - **Browse** : Parcours des containers et items +//! - **Search** : Recherche dans les sources qui le supportent +//! - **Update ID** : Suivi des changements pour les notifications UPnP + +use crate::server_ext::get_source_registry; +use pmodidl::{Container, DIDLLite}; +use pmosource::{BrowseResult, MusicSource}; +use std::sync::Arc; + +/// Convertit des containers et items en XML DIDL-Lite +fn to_didl_lite(containers: &[Container], items: &[pmodidl::Item]) -> Result { + let didl = DIDLLite { + xmlns: "urn:schemas-upnp-org:metadata-1-0/DIDL-Lite/".to_string(), + xmlns_upnp: Some("urn:schemas-upnp-org:metadata-1-0/upnp/".to_string()), + xmlns_dc: Some("http://purl.org/dc/elements/1.1/".to_string()), + xmlns_dlna: Some("urn:schemas-dlna-org:metadata-1-0/".to_string()), + xmlns_pv: None, + xmlns_sec: None, + containers: containers.to_vec(), + items: items.to_vec(), + }; + + quick_xml::se::to_string(&didl) + .map_err(|e| format!("Failed to serialize DIDL-Lite: {}", e)) +} + +/// Handler pour le service ContentDirectory +/// +/// Ce handler gère toutes les opérations du ContentDirectory en utilisant +/// les sources musicales enregistrées dans le registre global. +pub struct ContentHandler; + +impl ContentHandler { + /// Crée un nouveau ContentHandler + pub fn new() -> Self { + Self + } + + /// Browse un container ou récupère les métadonnées d'un objet + /// + /// # Arguments + /// + /// * `object_id` - L'ID de l'objet à parcourir ("0" pour la racine) + /// * `browse_flag` - "BrowseMetadata" ou "BrowseDirectChildren" + /// * `starting_index` - Index de départ pour la pagination + /// * `requested_count` - Nombre d'éléments demandés (0 = tous) + /// + /// # Returns + /// + /// Un tuple contenant: + /// - Le résultat DIDL-Lite XML + /// - Le nombre d'éléments retournés + /// - Le nombre total d'éléments + /// - L'update ID + /// + /// # Examples + /// + /// ```ignore + /// let handler = ContentHandler::new(); + /// let (didl, returned, total, update_id) = + /// handler.browse("0", "BrowseDirectChildren", 0, 0).await?; + /// ``` + pub async fn browse( + &self, + object_id: &str, + browse_flag: &str, + starting_index: u32, + requested_count: u32, + ) -> Result<(String, u32, u32, u32), String> { + tracing::debug!( + object_id = %object_id, + browse_flag = %browse_flag, + starting_index = %starting_index, + requested_count = %requested_count, + "ContentDirectory::Browse" + ); + + match browse_flag { + "BrowseMetadata" => self.browse_metadata(object_id).await, + "BrowseDirectChildren" => { + self.browse_direct_children(object_id, starting_index, requested_count) + .await + } + _ => Err(format!("Invalid BrowseFlag: {}", browse_flag)), + } + } + + /// Browse les métadonnées d'un objet spécifique + async fn browse_metadata(&self, object_id: &str) -> Result<(String, u32, u32, u32), String> { + if object_id == "0" { + // Retourner le container racine + let root = self.build_root_container().await; + let didl = to_didl_lite(&[root], &[])?; + Ok((didl, 1, 1, 0)) + } else { + // Essayer de trouver l'objet dans les sources + let registry = get_source_registry().await; + + // Vérifier si c'est un container racine d'une source + if let Some(source) = registry.get(object_id).await { + let container = source + .root_container() + .await + .map_err(|e| format!("Failed to get root container: {}", e))?; + let didl = to_didl_lite(&[container], &[])?; + let update_id = source.update_id().await; + return Ok((didl, 1, 1, update_id)); + } + + // Sinon, chercher dans les sources + for source in registry.list_all().await { + if let Ok(result) = source.browse(object_id).await { + // L'objet a été trouvé, retourner ses métadonnées + match result { + BrowseResult::Containers(containers) => { + if let Some(container) = containers.first() { + let didl = to_didl_lite(&[container.clone()], &[])?; + let update_id = source.update_id().await; + return Ok((didl, 1, 1, update_id)); + } + } + BrowseResult::Items(items) => { + if let Some(item) = items.first() { + let didl = to_didl_lite(&[], &[item.clone()])?; + let update_id = source.update_id().await; + return Ok((didl, 1, 1, update_id)); + } + } + BrowseResult::Mixed { containers, items } => { + if let Some(container) = containers.first() { + let didl = to_didl_lite(&[container.clone()], &[])?; + let update_id = source.update_id().await; + return Ok((didl, 1, 1, update_id)); + } else if let Some(item) = items.first() { + let didl = to_didl_lite(&[], &[item.clone()])?; + let update_id = source.update_id().await; + return Ok((didl, 1, 1, update_id)); + } + } + } + } + } + + Err(format!("Object not found: {}", object_id)) + } + } + + /// Browse les enfants directs d'un container + async fn browse_direct_children( + &self, + object_id: &str, + starting_index: u32, + requested_count: u32, + ) -> Result<(String, u32, u32, u32), String> { + if object_id == "0" { + // Retourner toutes les sources comme enfants de la racine + return self.browse_root(starting_index, requested_count).await; + } + + let registry = get_source_registry().await; + + // Vérifier si c'est le container racine d'une source + if let Some(source) = registry.get(object_id).await { + return self + .browse_source_root(source, starting_index, requested_count) + .await; + } + + // Sinon, chercher dans les sources + for source in registry.list_all().await { + if let Ok(result) = source.browse(object_id).await { + return self + .browse_result_to_didl(result, source, starting_index, requested_count) + .await; + } + } + + Err(format!("Container not found: {}", object_id)) + } + + /// Browse la racine (liste toutes les sources) + async fn browse_root( + &self, + starting_index: u32, + requested_count: u32, + ) -> Result<(String, u32, u32, u32), String> { + let registry = get_source_registry().await; + let sources = registry.list_all().await; + + let mut containers = Vec::new(); + for source in sources.iter() { + let container = source + .root_container() + .await + .map_err(|e| format!("Failed to get root container: {}", e))?; + containers.push(container); + } + + // Appliquer la pagination + let total = containers.len(); + let start = starting_index as usize; + let count = if requested_count == 0 { + total - start + } else { + requested_count as usize + }; + + let paginated: Vec = containers + .into_iter() + .skip(start) + .take(count) + .collect(); + + let returned = paginated.len(); + let didl = to_didl_lite(&paginated, &[])?; + + Ok((didl, returned as u32, total as u32, 0)) + } + + /// Browse le container racine d'une source spécifique + async fn browse_source_root( + &self, + source: Arc, + starting_index: u32, + requested_count: u32, + ) -> Result<(String, u32, u32, u32), String> { + let result = source + .browse(source.id()) + .await + .map_err(|e| format!("Browse failed: {}", e))?; + + self.browse_result_to_didl(result, source, starting_index, requested_count) + .await + } + + /// Convertit un BrowseResult en DIDL-Lite XML avec pagination + async fn browse_result_to_didl( + &self, + result: BrowseResult, + source: Arc, + starting_index: u32, + requested_count: u32, + ) -> Result<(String, u32, u32, u32), String> { + let (mut containers, mut items) = match result { + BrowseResult::Containers(c) => (c, vec![]), + BrowseResult::Items(i) => (vec![], i), + BrowseResult::Mixed { containers, items } => (containers, items), + }; + + // Calculer le total avant pagination + let total = (containers.len() + items.len()) as u32; + + // Appliquer la pagination + let start = starting_index as usize; + let count = if requested_count == 0 { + total as usize - start + } else { + requested_count as usize + }; + + // Pagination sur les containers d'abord, puis les items + let total_containers = containers.len(); + if start < total_containers { + // On commence dans les containers + containers = containers.into_iter().skip(start).collect(); + let remaining = count.saturating_sub(containers.len()); + containers.truncate(count); + + if remaining > 0 && !items.is_empty() { + items.truncate(remaining); + } else { + items.clear(); + } + } else { + // On commence dans les items + containers.clear(); + let item_start = start - total_containers; + items = items + .into_iter() + .skip(item_start) + .take(count) + .collect(); + } + + let returned = (containers.len() + items.len()) as u32; + let didl = to_didl_lite(&containers, &items)?; + let update_id = source.update_id().await; + + Ok((didl, returned, total, update_id)) + } + + /// Construit le container racine du MediaServer + async fn build_root_container(&self) -> Container { + let registry = get_source_registry().await; + let child_count = registry.count().await; + + Container { + id: "0".to_string(), + parent_id: "-1".to_string(), + restricted: Some("1".to_string()), + child_count: Some(child_count.to_string()), + title: "PMOMusic".to_string(), + class: "object.container".to_string(), + containers: vec![], + items: vec![], + } + } + + /// Recherche dans toutes les sources qui supportent la recherche + /// + /// # Arguments + /// + /// * `container_id` - ID du container dans lequel rechercher ("0" = partout) + /// * `search_criteria` - Critères de recherche UPnP + /// + /// # Returns + /// + /// Les mêmes informations que browse() + pub async fn search( + &self, + container_id: &str, + search_criteria: &str, + ) -> Result<(String, u32, u32, u32), String> { + tracing::debug!( + container_id = %container_id, + search_criteria = %search_criteria, + "ContentDirectory::Search" + ); + + let registry = get_source_registry().await; + let mut all_containers = Vec::new(); + let mut all_items = Vec::new(); + + // Rechercher dans toutes les sources qui supportent la recherche + for source in registry.list_all().await { + if source.capabilities().supports_search { + if let Ok(result) = source.search(search_criteria).await { + match result { + BrowseResult::Containers(c) => all_containers.extend(c), + BrowseResult::Items(i) => all_items.extend(i), + BrowseResult::Mixed { containers, items } => { + all_containers.extend(containers); + all_items.extend(items); + } + } + } + } + } + + let total = (all_containers.len() + all_items.len()) as u32; + let didl = to_didl_lite(&all_containers, &all_items)?; + + Ok((didl, total, total, 0)) + } + + /// Retourne les capacités de recherche + pub async fn get_search_capabilities(&self) -> String { + // Capacités de recherche de base UPnP + "dc:title,dc:creator,upnp:artist,upnp:album,upnp:genre".to_string() + } + + /// Retourne les capacités de tri + pub async fn get_sort_capabilities(&self) -> String { + // Capacités de tri de base UPnP + "dc:title,dc:date,upnp:artist,upnp:album".to_string() + } + + /// Retourne le system update ID global + pub async fn get_system_update_id(&self) -> u32 { + let registry = get_source_registry().await; + let sources = registry.list_all().await; + + // Combiner les update IDs de toutes les sources + let mut combined_id = 0u32; + for source in sources { + combined_id = combined_id.wrapping_add(source.update_id().await); + } + + combined_id + } +} + +impl Default for ContentHandler { + fn default() -> Self { + Self::new() + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[tokio::test] + async fn test_content_handler_creation() { + let handler = ContentHandler::new(); + let capabilities = handler.get_search_capabilities().await; + assert!(capabilities.contains("dc:title")); + } + + #[tokio::test] + async fn test_browse_root_empty() { + let handler = ContentHandler::new(); + let result = handler.browse("0", "BrowseDirectChildren", 0, 0).await; + assert!(result.is_ok()); + + let (didl, returned, total, _) = result.unwrap(); + assert_eq!(returned, 0); + assert_eq!(total, 0); + assert!(didl.contains("DIDL-Lite")); + } + + #[tokio::test] + async fn test_get_system_update_id() { + let handler = ContentHandler::new(); + let update_id = handler.get_system_update_id().await; + assert_eq!(update_id, 0); // No sources registered + } +} diff --git a/pmomediaserver/src/lib.rs b/pmomediaserver/src/lib.rs index bea3bdca..33858d61 100644 --- a/pmomediaserver/src/lib.rs +++ b/pmomediaserver/src/lib.rs @@ -16,7 +16,7 @@ //! - Type : `urn:schemas-upnp-org:device:MediaServer:1` //! - Services : ContentDirectory:1, ConnectionManager:1 //! -//! # Utilisation +//! # Utilisation de base //! //! ```ignore //! use pmomediaserver::MEDIA_SERVER; @@ -25,9 +25,38 @@ //! let server = MEDIA_SERVER.clone(); //! let instance = server.create_instance(); //! ``` +//! +//! # Gestion des sources musicales +//! +//! Le MediaServer peut diffuser plusieurs sources musicales (Qobuz, Radio Paradise, etc.) +//! via le trait `MediaServerExt` : +//! +//! ```ignore +//! use pmomediaserver::server_ext::MediaServerExt; +//! use pmoserver::ServerBuilder; +//! use std::sync::Arc; +//! +//! let mut server = ServerBuilder::new_configured().build(); +//! +//! // Enregistrer une source musicale +//! let qobuz = Arc::new(QobuzSource::new(credentials)); +//! server.register_music_source(qobuz).await; +//! +//! // Lister toutes les sources +//! let sources = server.list_music_sources().await; +//! for source in sources { +//! println!("Source: {} ({})", source.name(), source.id()); +//! } +//! ``` pub mod contentdirectory; pub mod connectionmanager; pub mod device; +pub mod source_registry; +pub mod server_ext; +pub mod content_handler; pub use device::MEDIA_SERVER; +pub use source_registry::SourceRegistry; +pub use server_ext::{MediaServerExt, get_source_registry}; +pub use content_handler::ContentHandler; diff --git a/pmomediaserver/src/server_ext.rs b/pmomediaserver/src/server_ext.rs new file mode 100644 index 00000000..d0d00737 --- /dev/null +++ b/pmomediaserver/src/server_ext.rs @@ -0,0 +1,278 @@ +//! # Extension trait pour le serveur MediaServer +//! +//! Ce module fournit un trait d'extension pour `pmoserver::Server` permettant +//! d'enregistrer facilement des sources musicales et de configurer le MediaServer. + +use crate::source_registry::SourceRegistry; +use async_trait::async_trait; +use pmosource::MusicSource; +use pmoserver::Server; +use std::sync::Arc; +use tokio::sync::OnceCell; + +/// Extension pour le registre de sources au niveau global +/// +/// Ce registre est partagé par toutes les instances du serveur et permet +/// d'accéder aux sources musicales depuis n'importe où dans l'application. +static GLOBAL_REGISTRY: OnceCell = OnceCell::const_new(); + +/// Initialise le registre global +/// +/// Cette fonction est appelée automatiquement lors de la première utilisation. +async fn init_global_registry() -> &'static SourceRegistry { + GLOBAL_REGISTRY + .get_or_init(|| async { SourceRegistry::new() }) + .await +} + +/// Récupère le registre global de sources +/// +/// # Examples +/// +/// ```ignore +/// use pmomediaserver::server_ext::get_source_registry; +/// +/// let registry = get_source_registry().await; +/// if let Some(source) = registry.get("qobuz").await { +/// // Utiliser la source +/// } +/// ``` +pub async fn get_source_registry() -> &'static SourceRegistry { + init_global_registry().await +} + +/// Trait d'extension pour le serveur permettant l'enregistrement de sources musicales +/// +/// Ce trait ajoute des méthodes pratiques à `Server` pour enregistrer des sources +/// musicales et les rendre disponibles via le MediaServer. +/// +/// # Examples +/// +/// ```ignore +/// use pmomediaserver::server_ext::MediaServerExt; +/// use pmoserver::ServerBuilder; +/// +/// let mut server = ServerBuilder::new_configured().build(); +/// +/// // Enregistrer une source +/// let qobuz = Arc::new(QobuzSource::new()); +/// server.register_music_source(qobuz).await; +/// +/// // Lister toutes les sources +/// let sources = server.list_music_sources().await; +/// ``` +#[async_trait] +pub trait MediaServerExt { + /// Enregistre une source musicale dans le MediaServer + /// + /// La source devient immédiatement disponible via le service ContentDirectory + /// et peut être parcourue par les clients UPnP. + /// + /// # Arguments + /// + /// * `source` - La source musicale à enregistrer (Arc) + /// + /// # Examples + /// + /// ```ignore + /// let qobuz = Arc::new(QobuzSource::new(credentials)); + /// server.register_music_source(qobuz).await; + /// ``` + async fn register_music_source(&mut self, source: Arc); + + /// Récupère une source musicale par son ID + /// + /// # Arguments + /// + /// * `id` - L'ID unique de la source + /// + /// # Returns + /// + /// Un `Arc` vers la source si elle existe, ou `None`. + /// + /// # Examples + /// + /// ```ignore + /// if let Some(source) = server.get_music_source("qobuz").await { + /// println!("Found: {}", source.name()); + /// } + /// ``` + async fn get_music_source(&self, id: &str) -> Option>; + + /// Liste toutes les sources musicales enregistrées + /// + /// # Returns + /// + /// Un vecteur contenant toutes les sources enregistrées. + /// + /// # Examples + /// + /// ```ignore + /// let sources = server.list_music_sources().await; + /// for source in sources { + /// println!("- {} ({})", source.name(), source.id()); + /// } + /// ``` + async fn list_music_sources(&self) -> Vec>; + + /// Compte le nombre de sources musicales enregistrées + /// + /// # Returns + /// + /// Le nombre total de sources. + /// + /// # Examples + /// + /// ```ignore + /// let count = server.count_music_sources().await; + /// println!("Total sources: {}", count); + /// ``` + async fn count_music_sources(&self) -> usize; + + /// Supprime une source musicale du registre + /// + /// # Arguments + /// + /// * `id` - L'ID de la source à supprimer + /// + /// # Returns + /// + /// `true` si la source a été supprimée, `false` si elle n'existait pas. + /// + /// # Examples + /// + /// ```ignore + /// if server.remove_music_source("old-radio").await { + /// println!("Source removed"); + /// } + /// ``` + async fn remove_music_source(&mut self, id: &str) -> bool; +} + +#[async_trait] +impl MediaServerExt for Server { + async fn register_music_source(&mut self, source: Arc) { + let registry = get_source_registry().await; + + tracing::info!( + source_id = %source.id(), + source_name = %source.name(), + "Registering music source to MediaServer" + ); + + registry.register(source).await; + } + + async fn get_music_source(&self, id: &str) -> Option> { + let registry = get_source_registry().await; + registry.get(id).await + } + + async fn list_music_sources(&self) -> Vec> { + let registry = get_source_registry().await; + registry.list_all().await + } + + async fn count_music_sources(&self) -> usize { + let registry = get_source_registry().await; + registry.count().await + } + + async fn remove_music_source(&mut self, id: &str) -> bool { + let registry = get_source_registry().await; + registry.remove(id).await + } +} + +#[cfg(test)] +mod tests { + use super::*; + use pmosource::{MusicSource, Result, BrowseResult}; + use pmodidl::{Container, Item}; + use std::time::SystemTime; + + #[derive(Debug)] + struct DummySource { + id: String, + name: String, + } + + impl DummySource { + fn new(id: &str, name: &str) -> Self { + Self { + id: id.to_string(), + name: name.to_string(), + } + } + } + + #[async_trait::async_trait] + impl MusicSource for DummySource { + fn name(&self) -> &str { + &self.name + } + + fn id(&self) -> &str { + &self.id + } + + fn default_image(&self) -> &[u8] { + &[] + } + + async fn root_container(&self) -> Result { + Ok(Container { + id: self.id.clone(), + parent_id: "0".to_string(), + restricted: Some("1".to_string()), + child_count: Some("0".to_string()), + title: self.name.clone(), + class: "object.container".to_string(), + containers: vec![], + items: vec![], + }) + } + + async fn browse(&self, _object_id: &str) -> Result { + Ok(BrowseResult::Items(vec![])) + } + + async fn resolve_uri(&self, object_id: &str) -> Result { + Ok(format!("http://example.com/{}", object_id)) + } + + fn supports_fifo(&self) -> bool { + false + } + + async fn append_track(&self, _track: Item) -> Result<()> { + Err(pmosource::MusicSourceError::FifoNotSupported) + } + + async fn remove_oldest(&self) -> Result> { + Err(pmosource::MusicSourceError::FifoNotSupported) + } + + async fn update_id(&self) -> u32 { + 0 + } + + async fn last_change(&self) -> Option { + None + } + + async fn get_items(&self, _offset: usize, _count: usize) -> Result> { + Ok(vec![]) + } + } + + #[tokio::test] + async fn test_global_registry_singleton() { + // Vérifier que le registre global est bien un singleton + let registry1 = get_source_registry().await; + let registry2 = get_source_registry().await; + + // Les deux références devraient pointer vers le même registre + assert!(std::ptr::eq(registry1, registry2)); + } +} diff --git a/pmomediaserver/src/source_registry.rs b/pmomediaserver/src/source_registry.rs new file mode 100644 index 00000000..ac99d145 --- /dev/null +++ b/pmomediaserver/src/source_registry.rs @@ -0,0 +1,368 @@ +//! # Source Registry - Gestionnaire de sources musicales +//! +//! Ce module fournit un registre centralisé pour gérer les différentes sources musicales +//! (MusicSource) qui peuvent être diffusées par le MediaServer. +//! +//! ## Fonctionnalités +//! +//! - **Enregistrement de sources** : Ajout de sources musicales au registre +//! - **Accès aux sources** : Récupération des sources enregistrées par ID +//! - **Navigation multi-sources** : Combine les sources dans une hiérarchie unique +//! - **Thread-safe** : Utilise Arc et RwLock pour un accès concurrent + +use pmosource::MusicSource; +use std::collections::HashMap; +use std::sync::Arc; +use tokio::sync::RwLock; + +/// Registre des sources musicales +/// +/// Ce registre maintient une liste de toutes les sources musicales enregistrées +/// et permet de les récupérer par leur ID unique. +/// +/// # Thread Safety +/// +/// Le registre utilise `Arc>` pour permettre un accès concurrent sécurisé. +/// Plusieurs lecteurs peuvent accéder simultanément aux sources, mais l'enregistrement +/// de nouvelles sources nécessite un verrou exclusif. +/// +/// # Examples +/// +/// ```ignore +/// use pmomediaserver::source_registry::SourceRegistry; +/// +/// let registry = SourceRegistry::new(); +/// +/// // Enregistrer une source +/// let source = Arc::new(MyMusicSource::new()); +/// registry.register(source).await; +/// +/// // Récupérer une source +/// if let Some(source) = registry.get("my-source-id").await { +/// let root = source.root_container().await?; +/// } +/// ``` +#[derive(Clone)] +pub struct SourceRegistry { + sources: Arc>>>, +} + +impl SourceRegistry { + /// Crée un nouveau registre vide + /// + /// # Examples + /// + /// ``` + /// use pmomediaserver::source_registry::SourceRegistry; + /// + /// let registry = SourceRegistry::new(); + /// ``` + pub fn new() -> Self { + Self { + sources: Arc::new(RwLock::new(HashMap::new())), + } + } + + /// Enregistre une nouvelle source musicale + /// + /// La source est identifiée par son ID unique (retourné par `source.id()`). + /// Si une source avec le même ID existe déjà, elle sera remplacée. + /// + /// # Arguments + /// + /// * `source` - La source musicale à enregistrer (doit implémenter `MusicSource`) + /// + /// # Examples + /// + /// ```ignore + /// let source = Arc::new(RadioParadise::new()); + /// registry.register(source).await; + /// ``` + pub async fn register(&self, source: Arc) { + let id = source.id().to_string(); + let mut sources = self.sources.write().await; + + tracing::info!( + source_id = %id, + source_name = %source.name(), + "Registering music source" + ); + + sources.insert(id, source); + } + + /// Récupère une source par son ID + /// + /// # Arguments + /// + /// * `id` - L'ID unique de la source + /// + /// # Returns + /// + /// Un `Arc` vers la source si elle existe, ou `None` si aucune source avec cet ID + /// n'est enregistrée. + /// + /// # Examples + /// + /// ```ignore + /// if let Some(source) = registry.get("radio-paradise").await { + /// println!("Found: {}", source.name()); + /// } + /// ``` + pub async fn get(&self, id: &str) -> Option> { + let sources = self.sources.read().await; + sources.get(id).cloned() + } + + /// Liste toutes les sources enregistrées + /// + /// # Returns + /// + /// Un vecteur contenant des clones de toutes les sources enregistrées. + /// + /// # Examples + /// + /// ```ignore + /// let all_sources = registry.list_all().await; + /// for source in all_sources { + /// println!("Source: {} ({})", source.name(), source.id()); + /// } + /// ``` + pub async fn list_all(&self) -> Vec> { + let sources = self.sources.read().await; + sources.values().cloned().collect() + } + + /// Compte le nombre de sources enregistrées + /// + /// # Returns + /// + /// Le nombre total de sources dans le registre. + /// + /// # Examples + /// + /// ```ignore + /// let count = registry.count().await; + /// println!("Total sources: {}", count); + /// ``` + pub async fn count(&self) -> usize { + let sources = self.sources.read().await; + sources.len() + } + + /// Supprime une source du registre + /// + /// # Arguments + /// + /// * `id` - L'ID de la source à supprimer + /// + /// # Returns + /// + /// `true` si la source a été supprimée, `false` si elle n'existait pas. + /// + /// # Examples + /// + /// ```ignore + /// if registry.remove("old-source").await { + /// println!("Source removed"); + /// } + /// ``` + pub async fn remove(&self, id: &str) -> bool { + let mut sources = self.sources.write().await; + + if sources.remove(id).is_some() { + tracing::info!(source_id = %id, "Removed music source"); + true + } else { + false + } + } + + /// Vérifie si une source est enregistrée + /// + /// # Arguments + /// + /// * `id` - L'ID de la source à vérifier + /// + /// # Returns + /// + /// `true` si la source existe, `false` sinon. + /// + /// # Examples + /// + /// ```ignore + /// if registry.contains("qobuz").await { + /// // La source Qobuz est disponible + /// } + /// ``` + pub async fn contains(&self, id: &str) -> bool { + let sources = self.sources.read().await; + sources.contains_key(id) + } +} + +impl Default for SourceRegistry { + fn default() -> Self { + Self::new() + } +} + +#[cfg(test)] +mod tests { + use super::*; + use pmosource::{MusicSource, Result, BrowseResult}; + use pmodidl::{Container, Item}; + use std::time::SystemTime; + + #[derive(Debug)] + struct TestSource { + id: String, + name: String, + } + + impl TestSource { + fn new(id: &str, name: &str) -> Self { + Self { + id: id.to_string(), + name: name.to_string(), + } + } + } + + #[async_trait::async_trait] + impl MusicSource for TestSource { + fn name(&self) -> &str { + &self.name + } + + fn id(&self) -> &str { + &self.id + } + + fn default_image(&self) -> &[u8] { + &[] + } + + async fn root_container(&self) -> Result { + Ok(Container { + id: self.id.clone(), + parent_id: "0".to_string(), + restricted: Some("1".to_string()), + child_count: Some("0".to_string()), + title: self.name.clone(), + class: "object.container".to_string(), + containers: vec![], + items: vec![], + }) + } + + async fn browse(&self, _object_id: &str) -> Result { + Ok(BrowseResult::Items(vec![])) + } + + async fn resolve_uri(&self, object_id: &str) -> Result { + Ok(format!("http://example.com/{}", object_id)) + } + + fn supports_fifo(&self) -> bool { + false + } + + async fn append_track(&self, _track: Item) -> Result<()> { + Err(pmosource::MusicSourceError::FifoNotSupported) + } + + async fn remove_oldest(&self) -> Result> { + Err(pmosource::MusicSourceError::FifoNotSupported) + } + + async fn update_id(&self) -> u32 { + 0 + } + + async fn last_change(&self) -> Option { + None + } + + async fn get_items(&self, _offset: usize, _count: usize) -> Result> { + Ok(vec![]) + } + } + + #[tokio::test] + async fn test_register_and_get() { + let registry = SourceRegistry::new(); + let source = Arc::new(TestSource::new("test-1", "Test Source 1")); + + registry.register(source.clone()).await; + + let retrieved = registry.get("test-1").await; + assert!(retrieved.is_some()); + + let retrieved = retrieved.unwrap(); + assert_eq!(retrieved.id(), "test-1"); + assert_eq!(retrieved.name(), "Test Source 1"); + } + + #[tokio::test] + async fn test_list_all() { + let registry = SourceRegistry::new(); + + registry.register(Arc::new(TestSource::new("test-1", "Test 1"))).await; + registry.register(Arc::new(TestSource::new("test-2", "Test 2"))).await; + registry.register(Arc::new(TestSource::new("test-3", "Test 3"))).await; + + let sources = registry.list_all().await; + assert_eq!(sources.len(), 3); + } + + #[tokio::test] + async fn test_count() { + let registry = SourceRegistry::new(); + assert_eq!(registry.count().await, 0); + + registry.register(Arc::new(TestSource::new("test-1", "Test 1"))).await; + assert_eq!(registry.count().await, 1); + + registry.register(Arc::new(TestSource::new("test-2", "Test 2"))).await; + assert_eq!(registry.count().await, 2); + } + + #[tokio::test] + async fn test_remove() { + let registry = SourceRegistry::new(); + registry.register(Arc::new(TestSource::new("test-1", "Test 1"))).await; + + assert!(registry.contains("test-1").await); + assert!(registry.remove("test-1").await); + assert!(!registry.contains("test-1").await); + assert!(!registry.remove("test-1").await); + } + + #[tokio::test] + async fn test_contains() { + let registry = SourceRegistry::new(); + + assert!(!registry.contains("test-1").await); + + registry.register(Arc::new(TestSource::new("test-1", "Test 1"))).await; + + assert!(registry.contains("test-1").await); + assert!(!registry.contains("test-2").await); + } + + #[tokio::test] + async fn test_replace_source() { + let registry = SourceRegistry::new(); + + registry.register(Arc::new(TestSource::new("test-1", "Old Name"))).await; + let old = registry.get("test-1").await.unwrap(); + assert_eq!(old.name(), "Old Name"); + + registry.register(Arc::new(TestSource::new("test-1", "New Name"))).await; + let new = registry.get("test-1").await.unwrap(); + assert_eq!(new.name(), "New Name"); + + assert_eq!(registry.count().await, 1); + } +}