Ajoute un media server à l'application PMOMusic

This commit is contained in:
2025-10-17 09:01:24 +02:00
parent 552d8d90fb
commit 6a87046845
9 changed files with 1167 additions and 26 deletions

View File

@@ -17,6 +17,8 @@ devices:
udn: d7eaad15-7d21-4411-926a-bc1eea0713db
mediarenderer:
udn: f9ef6c21-0ed3-470c-9846-bc1ae85fea62
mediaserver:
udn: 4aa1d843-22ca-4a1d-b383-43a9d197d875
mediaserver:
qobuz:
udn: 28963b75-4c5f-4da7-b10e-ffafd

7
Cargo.lock generated
View File

@@ -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]]

View File

@@ -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"] }

View File

@@ -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::<Webapp>("/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");

View File

@@ -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"] }

View File

@@ -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<String, String> {
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<Container> = 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<dyn MusicSource>,
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<dyn MusicSource>,
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
}
}

View File

@@ -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;

View File

@@ -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<SourceRegistry> = 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<dyn MusicSource>)
///
/// # Examples
///
/// ```ignore
/// let qobuz = Arc::new(QobuzSource::new(credentials));
/// server.register_music_source(qobuz).await;
/// ```
async fn register_music_source(&mut self, source: Arc<dyn MusicSource>);
/// 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<Arc<dyn MusicSource>>;
/// 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<Arc<dyn MusicSource>>;
/// 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<dyn MusicSource>) {
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<Arc<dyn MusicSource>> {
let registry = get_source_registry().await;
registry.get(id).await
}
async fn list_music_sources(&self) -> Vec<Arc<dyn MusicSource>> {
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<Container> {
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<BrowseResult> {
Ok(BrowseResult::Items(vec![]))
}
async fn resolve_uri(&self, object_id: &str) -> Result<String> {
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<Option<Item>> {
Err(pmosource::MusicSourceError::FifoNotSupported)
}
async fn update_id(&self) -> u32 {
0
}
async fn last_change(&self) -> Option<SystemTime> {
None
}
async fn get_items(&self, _offset: usize, _count: usize) -> Result<Vec<Item>> {
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));
}
}

View File

@@ -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<RwLock<...>>` 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<RwLock<HashMap<String, Arc<dyn MusicSource>>>>,
}
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<dyn MusicSource>) {
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<Arc<dyn MusicSource>> {
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<Arc<dyn MusicSource>> {
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<Container> {
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<BrowseResult> {
Ok(BrowseResult::Items(vec![]))
}
async fn resolve_uri(&self, object_id: &str) -> Result<String> {
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<Option<Item>> {
Err(pmosource::MusicSourceError::FifoNotSupported)
}
async fn update_id(&self) -> u32 {
0
}
async fn last_change(&self) -> Option<SystemTime> {
None
}
async fn get_items(&self, _offset: usize, _count: usize) -> Result<Vec<Item>> {
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);
}
}