From f434c521e8f38fd3828f173e3aef04c6a9d7772f Mon Sep 17 00:00:00 2001 From: Eric Coissac Date: Mon, 1 Dec 2025 19:55:55 +0100 Subject: [PATCH] =?UTF-8?q?On=20fait=20la=20m=C3=AAme=20chose=20maintenant?= =?UTF-8?q?=20cot=C3=A9=20serveur=20de=20musique.?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- Cargo.lock | 1 + pmocontrol/Cargo.toml | 1 + pmocontrol/examples/renderer_demo.rs | 4 +- pmocontrol/src/arylic_tcp.rs | 43 ++- pmocontrol/src/control_point.rs | 394 ++++++++++++++++++++++-- pmocontrol/src/discovery.rs | 3 +- pmocontrol/src/lib.rs | 12 +- pmocontrol/src/media_server.rs | 430 +++++++++++++++++++++++++++ pmocontrol/src/model.rs | 26 -- pmocontrol/src/music_renderer.rs | 4 +- pmocontrol/src/playback_queue.rs | 73 +++++ pmocontrol/src/provider.rs | 57 ++-- pmocontrol/src/registry.rs | 13 +- pmocontrol/src/soap_client.rs | 21 +- 14 files changed, 968 insertions(+), 114 deletions(-) create mode 100644 pmocontrol/src/media_server.rs create mode 100644 pmocontrol/src/playback_queue.rs diff --git a/Cargo.lock b/Cargo.lock index 691963c6..22f972dd 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3164,6 +3164,7 @@ version = "0.1.0" dependencies = [ "anyhow", "crossbeam-channel", + "pmodidl", "pmoupnp", "quick-xml 0.38.4", "thiserror 2.0.17", diff --git a/pmocontrol/Cargo.toml b/pmocontrol/Cargo.toml index 2416a266..016aa53a 100644 --- a/pmocontrol/Cargo.toml +++ b/pmocontrol/Cargo.toml @@ -5,6 +5,7 @@ edition = "2024" [dependencies] pmoupnp = { path = "../pmoupnp" } +pmodidl = { path = "../pmodidl" } quick-xml = "0.38.4" thiserror = "2.0.17" ureq = "3.1.4" diff --git a/pmocontrol/examples/renderer_demo.rs b/pmocontrol/examples/renderer_demo.rs index 046d5894..edd8aa92 100644 --- a/pmocontrol/examples/renderer_demo.rs +++ b/pmocontrol/examples/renderer_demo.rs @@ -229,7 +229,9 @@ fn print_backend(prefix: &str, renderer: &MusicRenderer) { MusicRenderer::Upnp(_) => "UpnpRenderer (UPnP AV / DLNA)", MusicRenderer::LinkPlay(_) => "LinkPlayRenderer (LinkPlay HTTP)", MusicRenderer::ArylicTcp(_) => "ArylicTcpRenderer (ARylic TCP Protocol)", - MusicRenderer::HybridUpnpArylic{..} => "Hybrid UpnpArylicRenderer (UPnP AV / DLNA + ARylic TCP Protocol)", + MusicRenderer::HybridUpnpArylic { .. } => { + "Hybrid UpnpArylicRenderer (UPnP AV / DLNA + ARylic TCP Protocol)" + } }; println!("{prefix}Backend : {backend}"); } diff --git a/pmocontrol/src/arylic_tcp.rs b/pmocontrol/src/arylic_tcp.rs index 73648ed1..8ff41791 100644 --- a/pmocontrol/src/arylic_tcp.rs +++ b/pmocontrol/src/arylic_tcp.rs @@ -24,7 +24,6 @@ fn last_command_time() -> &'static Mutex { LAST_COMMAND_TIME.get_or_init(|| Mutex::new(Instant::now())) } - const ARYLIC_TCP_PORT: u16 = 8899; const PACKET_HEADER: [u8; 4] = [0x18, 0x96, 0x18, 0x20]; const RESERVED_BYTES: [u8; 8] = [0; 8]; @@ -387,7 +386,10 @@ fn connect(host: &str, port: u16, timeout: Duration) -> Result { let elapsed = last_time.elapsed(); if elapsed < Duration::from_millis(200) { let wait = Duration::from_millis(200) - elapsed; - debug!("Waiting {:?} before sending command to respect 200ms interval", wait); + debug!( + "Waiting {:?} before sending command to respect 200ms interval", + wait + ); thread::sleep(wait); } *last_time = Instant::now(); @@ -521,9 +523,12 @@ fn send_command_with_mode( let mut stream = connect(host, port, timeout)?; let packet = encode_packet(payload); - stream - .write_all(&packet) - .with_context(|| format!("Failed to write Arylic TCP packet for {}: {}", host, payload))?; + stream.write_all(&packet).with_context(|| { + format!( + "Failed to write Arylic TCP packet for {}: {}", + host, payload + ) + })?; stream.flush().with_context(|| { format!( "Failed to flush Arylic TCP stream for {} (command {})", @@ -540,8 +545,9 @@ fn send_command_with_mode( let _ = stream.shutdown(Shutdown::Write); Ok(None) } - ResponseMode::Required(expected) => read_expected_response(&mut stream, host, payload, expected) - .map(Some), + ResponseMode::Required(expected) => { + read_expected_response(&mut stream, host, payload, expected).map(Some) + } ResponseMode::Optional(expected) => { for _ in 0..MAX_RESPONSE_ATTEMPTS { match read_packet(&mut stream) { @@ -615,7 +621,13 @@ fn send_command_required( payload: &str, expected: &[&str], ) -> Result { - match send_command_with_mode(host, port, timeout, payload, ResponseMode::Required(expected))? { + match send_command_with_mode( + host, + port, + timeout, + payload, + ResponseMode::Required(expected), + )? { Some(s) => Ok(s), None => Err(anyhow!( "Arylic TCP: no response payload for required command {}", @@ -631,14 +643,15 @@ fn send_command_optional( payload: &str, expected: &[&str], ) -> Result> { - send_command_with_mode(host, port, timeout, payload, ResponseMode::Optional(expected)) + send_command_with_mode( + host, + port, + timeout, + payload, + ResponseMode::Optional(expected), + ) } -fn send_command_no_response( - host: &str, - port: u16, - timeout: Duration, - payload: &str, -) -> Result<()> { +fn send_command_no_response(host: &str, port: u16, timeout: Duration, payload: &str) -> Result<()> { send_command_with_mode(host, port, timeout, payload, ResponseMode::None).map(|_| ()) } diff --git a/pmocontrol/src/control_point.rs b/pmocontrol/src/control_point.rs index f241ceb6..258d54d3 100644 --- a/pmocontrol/src/control_point.rs +++ b/pmocontrol/src/control_point.rs @@ -1,19 +1,25 @@ -use std::collections::{HashMap, HashSet}; +use std::collections::HashMap; use std::io; -use std::sync::{Arc, RwLock}; +use std::sync::{Arc, Mutex, RwLock}; use std::thread; use std::time::Duration; +use anyhow::anyhow; use crossbeam_channel::Receiver; use pmoupnp::ssdp::SsdpClient; +use tracing::{debug, error, warn}; use crate::MusicRenderer; use crate::capabilities::{ - PlaybackPosition, PlaybackPositionInfo, PlaybackState, PlaybackStatus, VolumeControl, + PlaybackPosition, PlaybackPositionInfo, PlaybackState, PlaybackStatus, TransportControl, + VolumeControl, }; use crate::discovery::DiscoveryManager; use crate::events::RendererEventBus; +use crate::media_server::{MediaServerInfo, ServerId}; use crate::model::{RendererEvent, RendererId, RendererProtocol}; +use crate::music_renderer::op_not_supported; +use crate::playback_queue::{PlaybackItem, PlaybackQueue}; use crate::provider::HttpXmlDescriptionProvider; use crate::registry::{DeviceRegistry, DeviceRegistryRead, DeviceUpdate}; use crate::upnp_renderer::UpnpRenderer; @@ -25,6 +31,7 @@ use crate::upnp_renderer::UpnpRenderer; pub struct ControlPoint { registry: Arc>, event_bus: RendererEventBus, + runtime: Arc, } impl ControlPoint { @@ -34,6 +41,7 @@ impl ControlPoint { pub fn spawn(timeout_secs: u64) -> io::Result { let registry = Arc::new(RwLock::new(DeviceRegistry::new())); let event_bus = RendererEventBus::new(); + let runtime = Arc::new(RuntimeState::new()); // SsdpClient let client = SsdpClient::new()?; // pmoupnp::ssdp::SsdpClient @@ -83,11 +91,10 @@ impl ControlPoint { let runtime_cp = ControlPoint { registry: Arc::clone(®istry), event_bus: event_bus.clone(), + runtime: Arc::clone(&runtime), }; thread::spawn(move || { - let mut cache: HashMap = HashMap::new(); - loop { let renderers = { let reg = runtime_cp.registry.read().unwrap(); @@ -97,12 +104,9 @@ impl ControlPoint { .collect::>() }; - let mut seen_ids = HashSet::new(); - for renderer in renderers { let info = renderer.info(); - // Ne pas poller les renderers offline if !info.online { continue; } @@ -113,20 +117,12 @@ impl ControlPoint { } let renderer_id = info.id.clone(); - seen_ids.insert(renderer_id.clone()); + let prev_snapshot = runtime_cp.runtime.snapshot_for(&renderer_id); + let mut new_snapshot = prev_snapshot.clone(); + let prev_position = prev_snapshot.position.clone(); - let entry = cache - .entry(renderer_id.clone()) - .or_insert_with(RendererRuntimeSnapshot::default); - - // Keep a snapshot of the previous position to compute logical - // state transitions based on time deltas. - let prev_position = entry.position.clone(); - - // 1) Poll position first, so that the state logic can use the - // freshly updated position when available. if let Ok(position) = renderer.playback_position() { - let has_changed = match entry.position.as_ref() { + let has_changed = match prev_snapshot.position.as_ref() { Some(prev) => !playback_position_equal(prev, &position), None => true, }; @@ -138,19 +134,17 @@ impl ControlPoint { }); } - entry.position = Some(position); + new_snapshot.position = Some(position); } - // 2) Poll raw playback state and compute a logical state that - // compensates for buggy devices (Arylic / LinkPlay). if let Ok(raw_state) = renderer.playback_state() { let logical_state = compute_logical_playback_state( &raw_state, prev_position.as_ref(), - entry.position.as_ref(), + new_snapshot.position.as_ref(), ); - let has_changed = match entry.state.as_ref() { + let has_changed = match prev_snapshot.state.as_ref() { Some(prev) => !playback_state_equal(prev, &logical_state), None => true, }; @@ -160,32 +154,38 @@ impl ControlPoint { id: renderer_id.clone(), state: logical_state.clone(), }); - entry.state = Some(logical_state); } + + new_snapshot.state = Some(logical_state); } if let Ok(volume) = renderer.volume() { - if entry.last_volume != Some(volume) { + if prev_snapshot.last_volume != Some(volume) { runtime_cp.emit_renderer_event(RendererEvent::VolumeChanged { id: renderer_id.clone(), volume, }); - entry.last_volume = Some(volume); } + + new_snapshot.last_volume = Some(volume); } if let Ok(mute) = renderer.mute() { - if entry.last_mute != Some(mute) { + if prev_snapshot.last_mute != Some(mute) { runtime_cp.emit_renderer_event(RendererEvent::MuteChanged { id: renderer_id.clone(), mute, }); - entry.last_mute = Some(mute); } + + new_snapshot.last_mute = Some(mute); } + + runtime_cp + .runtime + .update_snapshot(&renderer_id, new_snapshot); } - cache.retain(|id, _| seen_ids.contains(id)); thread::sleep(Duration::from_secs(1)); } }); @@ -193,6 +193,7 @@ impl ControlPoint { Ok(Self { registry, event_bus, + runtime, }) } @@ -254,6 +255,184 @@ impl ControlPoint { .and_then(|info| MusicRenderer::from_registry_info(info, ®)) } + /// Snapshot list of media servers currently known by the registry. + pub fn list_media_servers(&self) -> Vec { + let reg = self.registry.read().unwrap(); + reg.list_servers() + } + + /// Lookup a media server by id. + pub fn media_server(&self, id: &ServerId) -> Option { + let reg = self.registry.read().unwrap(); + reg.get_server(id) + } + + pub fn clear_queue(&self, renderer_id: &RendererId) -> anyhow::Result<()> { + if !self.runtime.has_entry(renderer_id) { + let err = Self::runtime_entry_missing(renderer_id); + warn!( + renderer = renderer_id.0.as_str(), + "Cannot clear queue: renderer not registered in runtime" + ); + return Err(err); + } + + let removed = self + .runtime + .with_queue_mut(renderer_id, |queue| { + let removed = queue.len(); + queue.clear(); + removed + }) + .ok_or_else(|| Self::runtime_entry_missing(renderer_id))?; + + debug!( + renderer = renderer_id.0.as_str(), + items_removed = removed, + queue_len = 0, + "Cleared playback queue" + ); + Ok(()) + } + + pub fn enqueue_items( + &self, + renderer_id: &RendererId, + items: Vec, + ) -> anyhow::Result<()> { + if !self.runtime.has_entry(renderer_id) { + let err = Self::runtime_entry_missing(renderer_id); + warn!( + renderer = renderer_id.0.as_str(), + "Cannot enqueue items: renderer not registered in runtime" + ); + return Err(err); + } + + let item_count = items.len(); + let new_len = self + .runtime + .with_queue_mut(renderer_id, |queue| { + queue.enqueue_many(items); + queue.len() + }) + .ok_or_else(|| Self::runtime_entry_missing(renderer_id))?; + + debug!( + renderer = renderer_id.0.as_str(), + added = item_count, + queue_len = new_len, + "Enqueued playback items" + ); + Ok(()) + } + + pub fn get_queue_snapshot( + &self, + renderer_id: &RendererId, + ) -> anyhow::Result> { + if !self.runtime.has_entry(renderer_id) { + let err = Self::runtime_entry_missing(renderer_id); + warn!( + renderer = renderer_id.0.as_str(), + "Cannot snapshot queue: renderer not registered in runtime" + ); + return Err(err); + } + + self.runtime + .queue_snapshot(renderer_id) + .ok_or_else(|| Self::runtime_entry_missing(renderer_id)) + } + + pub fn play_next_from_queue(&self, renderer_id: &RendererId) -> anyhow::Result<()> { + if !self.runtime.has_entry(renderer_id) { + let err = Self::runtime_entry_missing(renderer_id); + warn!( + renderer = renderer_id.0.as_str(), + "Cannot advance queue: renderer not registered in runtime" + ); + return Err(err); + } + + let Some((item, remaining_after)) = self.runtime.dequeue_next(renderer_id) else { + debug!( + renderer = renderer_id.0.as_str(), + "play_next_from_queue: queue is empty" + ); + self.runtime + .set_playback_source(renderer_id, PlaybackSource::None); + return Ok(()); + }; + + let queue_before = remaining_after + 1; + debug!( + renderer = renderer_id.0.as_str(), + queue_before, + queue_after = remaining_after, + uri = item.uri.as_str(), + "Dequeued next playback item" + ); + + let renderer = self + .music_renderer_by_id(renderer_id) + .ok_or_else(|| { + warn!( + renderer = renderer_id.0.as_str(), + "Renderer disappeared before queue playback could start" + ); + anyhow!("Renderer {} not found", renderer_id.0) + })?; + + if matches!(renderer.info().protocol, RendererProtocol::OpenHomeOnly) { + self.runtime + .with_queue_mut(renderer_id, |queue| queue.enqueue_front(item)) + .ok_or_else(|| Self::runtime_entry_missing(renderer_id))?; + self.runtime + .set_playback_source(renderer_id, PlaybackSource::None); + return Err(op_not_supported( + "play_next_from_queue", + "OpenHomeOnly renderer", + )); + } + + let playback = (|| -> anyhow::Result<()> { + renderer.play_uri(&item.uri, "")?; + renderer.play()?; + Ok(()) + })(); + + if let Err(err) = playback { + error!( + renderer = renderer_id.0.as_str(), + error = %err, + "Failed to start playback for queued item" + ); + if self + .runtime + .with_queue_mut(renderer_id, |queue| queue.enqueue_front(item)) + .is_none() + { + warn!( + renderer = renderer_id.0.as_str(), + "Failed to requeue item after playback error" + ); + } + self.runtime + .set_playback_source(renderer_id, PlaybackSource::None); + return Err(err); + } + + self.runtime + .set_playback_source(renderer_id, PlaybackSource::FromQueue); + debug!( + renderer = renderer_id.0.as_str(), + queue_len = remaining_after, + "Started playback from queue" + ); + Ok(()) + } + /// Subscribe to renderer events emitted by the control point runtime. /// /// Each subscriber receives all future events independently. @@ -261,13 +440,51 @@ impl ControlPoint { self.event_bus.subscribe() } - #[allow(dead_code)] pub(crate) fn emit_renderer_event(&self, event: RendererEvent) { + self.handle_renderer_event(&event); self.event_bus.broadcast(event); } + + fn handle_renderer_event(&self, event: &RendererEvent) { + if let RendererEvent::StateChanged { id, state } = event { + match state { + PlaybackState::Stopped => { + if self.runtime.is_playing_from_queue(id) { + debug!( + renderer = id.0.as_str(), + "Renderer stopped after queue-driven playback; advancing" + ); + if let Err(err) = self.play_next_from_queue(id) { + error!( + renderer = id.0.as_str(), + error = %err, + "Auto-advance failed; clearing queue playback state" + ); + self.runtime + .set_playback_source(id, PlaybackSource::None); + } + } else { + self.runtime + .set_playback_source(id, PlaybackSource::None); + } + } + PlaybackState::Playing => { + self.runtime.mark_external_if_idle(id); + } + _ => {} + } + } + } + + fn runtime_entry_missing(renderer_id: &RendererId) -> anyhow::Error { + anyhow!( + "Renderer {} not registered in control point runtime", + renderer_id.0 + ) + } } -#[derive(Default)] +#[derive(Clone, Default)] struct RendererRuntimeSnapshot { state: Option, position: Option, @@ -275,6 +492,117 @@ struct RendererRuntimeSnapshot { last_mute: Option, } +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] +enum PlaybackSource { + #[default] + None, + FromQueue, + External, +} + +#[derive(Default)] +struct RendererRuntimeEntry { + snapshot: RendererRuntimeSnapshot, + queue: PlaybackQueue, + playback_source: PlaybackSource, +} + +struct RuntimeState { + entries: Mutex>, +} + +impl RuntimeState { + fn new() -> Self { + Self { + entries: Mutex::new(HashMap::new()), + } + } + + fn snapshot_for(&self, id: &RendererId) -> RendererRuntimeSnapshot { + let entries = self.entries.lock().unwrap(); + entries + .get(id) + .map(|entry| entry.snapshot.clone()) + .unwrap_or_default() + } + + fn update_snapshot(&self, id: &RendererId, snapshot: RendererRuntimeSnapshot) { + self.with_entry(id, |entry| { + entry.snapshot = snapshot; + }); + } + + fn has_entry(&self, id: &RendererId) -> bool { + let entries = self.entries.lock().unwrap(); + entries.contains_key(id) + } + + fn with_queue_mut(&self, id: &RendererId, f: F) -> Option + where + F: FnOnce(&mut PlaybackQueue) -> R, + { + let mut entries = self.entries.lock().unwrap(); + entries + .get_mut(id) + .map(|entry| f(&mut entry.queue)) + } + + fn queue_snapshot(&self, id: &RendererId) -> Option> { + let entries = self.entries.lock().unwrap(); + entries.get(id).map(|entry| entry.queue.snapshot()) + } + + fn dequeue_next(&self, id: &RendererId) -> Option<(PlaybackItem, usize)> { + let mut entries = self.entries.lock().unwrap(); + let entry = entries.get_mut(id)?; + let item = entry.queue.dequeue()?; + let remaining = entry.queue.len(); + Some((item, remaining)) + } + + fn set_playback_source(&self, id: &RendererId, source: PlaybackSource) { + let mut entries = self.entries.lock().unwrap(); + if let Some(entry) = entries.get_mut(id) { + entry.playback_source = source; + } + } + + fn playback_source(&self, id: &RendererId) -> PlaybackSource { + let entries = self.entries.lock().unwrap(); + entries + .get(id) + .map(|entry| entry.playback_source) + .unwrap_or(PlaybackSource::None) + } + + fn is_playing_from_queue(&self, id: &RendererId) -> bool { + matches!( + self.playback_source(id), + PlaybackSource::FromQueue + ) + } + + fn mark_external_if_idle(&self, id: &RendererId) { + let mut entries = self.entries.lock().unwrap(); + if let Some(entry) = entries.get_mut(id) { + if matches!(entry.playback_source, PlaybackSource::None) { + entry.playback_source = PlaybackSource::External; + } + } + } + + fn with_entry(&self, id: &RendererId, f: F) -> R + where + F: FnOnce(&mut RendererRuntimeEntry) -> R, + { + let mut entries = self.entries.lock().unwrap(); + let entry = entries + .entry(id.clone()) + .or_insert_with(RendererRuntimeEntry::default); + f(entry) + } +} + /// Parse "HH:MM:SS" style time strings to seconds. /// /// Returns None for empty or sentinel values such as "NOT_IMPLEMENTED" or "-:--:--". diff --git a/pmocontrol/src/discovery.rs b/pmocontrol/src/discovery.rs index 29934723..013d7198 100644 --- a/pmocontrol/src/discovery.rs +++ b/pmocontrol/src/discovery.rs @@ -3,7 +3,8 @@ use std::time::SystemTime; use pmoupnp::ssdp::SsdpEvent; -use crate::model::{MediaServerInfo, RendererInfo}; +use crate::media_server::MediaServerInfo; +use crate::model::RendererInfo; use crate::registry::DeviceUpdate; /// État connu pour un endpoint UPnP identifié par son UDN. diff --git a/pmocontrol/src/lib.rs b/pmocontrol/src/lib.rs index 129025e1..20f6bb90 100644 --- a/pmocontrol/src/lib.rs +++ b/pmocontrol/src/lib.rs @@ -7,8 +7,10 @@ pub mod connection_manager_client; pub mod control_point; pub mod discovery; pub mod linkplay; +pub mod media_server; pub mod model; pub mod music_renderer; +pub mod playback_queue; pub mod provider; pub mod registry; pub mod rendering_control_client; @@ -24,15 +26,17 @@ pub use capabilities::{ pub use connection_manager_client::{ConnectionInfo, ConnectionManagerClient, ProtocolInfo}; pub use control_point::ControlPoint; pub use linkplay::LinkPlayRenderer; +pub use media_server::{ + MediaBrowser, MediaEntry, MediaResource, MediaServerInfo, MusicServer, ServerId, + UpnpMediaServer, +}; pub use music_renderer::MusicRenderer; +pub use playback_queue::{PlaybackItem, PlaybackQueue}; pub use rendering_control_client::RenderingControlClient; pub use upnp_renderer::UpnpRenderer; pub use discovery::{DeviceDescriptionProvider, DiscoveredEndpoint, DiscoveryManager}; -pub use model::{ - MediaServerCapabilities, MediaServerId, MediaServerInfo, RendererCapabilities, RendererEvent, - RendererId, RendererInfo, RendererProtocol, -}; +pub use model::{RendererCapabilities, RendererEvent, RendererId, RendererInfo, RendererProtocol}; pub use provider::HttpXmlDescriptionProvider; pub use registry::{DeviceRegistry, DeviceRegistryRead, DeviceUpdate}; diff --git a/pmocontrol/src/media_server.rs b/pmocontrol/src/media_server.rs new file mode 100644 index 00000000..a6494371 --- /dev/null +++ b/pmocontrol/src/media_server.rs @@ -0,0 +1,430 @@ +use std::time::{Duration, SystemTime}; + +use anyhow::{Result, anyhow}; +use pmodidl::{self, DIDLLite}; +use pmoupnp::soap::SoapEnvelope; +use pmoupnp::soap::error_codes; +use xmltree::{Element, XMLNode}; + +use crate::soap_client::{SoapCallResult, invoke_upnp_action_with_timeout}; + +/// Unique identifier for a media server registered by the control point. +#[derive(Clone, Debug, Eq, PartialEq, Hash)] +pub struct ServerId(pub String); + +/// Snapshot of a media server discovered through UPnP SSDP. +#[derive(Clone, Debug)] +pub struct MediaServerInfo { + pub id: ServerId, + pub udn: String, + pub friendly_name: String, + pub model_name: String, + pub manufacturer: String, + pub location: String, + pub server_header: String, + pub online: bool, + pub last_seen: SystemTime, + pub max_age: u32, + pub has_content_directory: bool, + pub content_directory_service_type: Option, + pub content_directory_control_url: Option, +} + +/// Simplified view over a DIDL-Lite resource entry. +#[derive(Clone, Debug)] +pub struct MediaResource { + pub uri: String, + pub protocol_info: String, + pub duration: Option, +} + +/// Representation of either a container or an item returned by ContentDirectory. +#[derive(Clone, Debug)] +pub struct MediaEntry { + pub id: String, + pub parent_id: String, + pub title: String, + pub is_container: bool, + pub class: String, + pub resources: Vec, +} + +/// Backend-agnostic media browsing contract. +pub trait MediaBrowser { + fn browse_root(&self) -> Result>; + fn browse_children(&self, object_id: &str, start: u32, count: u32) -> Result>; + fn browse_object(&self, object_id: &str) -> Result; + fn search( + &self, + container_id: &str, + query: &str, + start: u32, + count: u32, + ) -> Result>; +} + +/// Façade over every supported media server backend. +#[derive(Clone, Debug)] +pub enum MusicServer { + Upnp(UpnpMediaServer), +} + +impl MusicServer { + pub fn from_info(info: &MediaServerInfo, timeout: Duration) -> Result { + Ok(MusicServer::Upnp(UpnpMediaServer::new( + info.clone(), + timeout, + ))) + } + + pub fn id(&self) -> &ServerId { + match self { + MusicServer::Upnp(upnp) => upnp.id(), + } + } + + pub fn info(&self) -> &MediaServerInfo { + match self { + MusicServer::Upnp(upnp) => upnp.info(), + } + } +} + +impl MediaBrowser for MusicServer { + fn browse_root(&self) -> Result> { + match self { + MusicServer::Upnp(upnp) => upnp.browse_root(), + } + } + + fn browse_children(&self, object_id: &str, start: u32, count: u32) -> Result> { + match self { + MusicServer::Upnp(upnp) => upnp.browse_children(object_id, start, count), + } + } + + fn browse_object(&self, object_id: &str) -> Result { + match self { + MusicServer::Upnp(upnp) => upnp.browse_object(object_id), + } + } + + fn search( + &self, + container_id: &str, + query: &str, + start: u32, + count: u32, + ) -> Result> { + match self { + MusicServer::Upnp(upnp) => upnp.search(container_id, query, start, count), + } + } +} + +/// Single UPnP ContentDirectory backend implementation. +#[derive(Clone, Debug)] +pub struct UpnpMediaServer { + info: MediaServerInfo, + timeout: Duration, +} + +impl UpnpMediaServer { + pub fn new(info: MediaServerInfo, timeout: Duration) -> Self { + Self { info, timeout } + } + + pub fn id(&self) -> &ServerId { + &self.info.id + } + + pub fn info(&self) -> &MediaServerInfo { + &self.info + } + + fn browse_with_flag( + &self, + object_id: &str, + browse_flag: &str, + start: u32, + count: u32, + ) -> Result> { + let start_str = start.to_string(); + let count_str = count.to_string(); + let args = vec![ + ("ObjectID", object_id.to_string()), + ("BrowseFlag", browse_flag.to_string()), + ("Filter", "*".to_string()), + ("StartingIndex", start_str), + ("RequestedCount", count_str), + ("SortCriteria", String::new()), + ]; + + let response = self.invoke_content_directory("Browse", None, args)?; + let envelope = response + .envelope + .ok_or_else(|| anyhow!("Missing SOAP envelope in Browse response"))?; + let didl_xml = extract_result_payload(&envelope, "BrowseResponse")?; + map_didl_entries(&didl_xml) + } + + fn search_impl( + &self, + container_id: &str, + query: &str, + start: u32, + count: u32, + ) -> Result> { + let start_str = start.to_string(); + let count_str = count.to_string(); + + let args = vec![ + ("ContainerID", container_id.to_string()), + ("SearchCriteria", query.to_string()), + ("Filter", "*".to_string()), + ("StartingIndex", start_str), + ("RequestedCount", count_str), + ("SortCriteria", String::new()), + ]; + + let response = self.invoke_content_directory("Search", Some("search"), args)?; + let envelope = response + .envelope + .ok_or_else(|| anyhow!("Missing SOAP envelope in Search response"))?; + let didl_xml = extract_result_payload(&envelope, "SearchResponse")?; + map_didl_entries(&didl_xml) + } + + fn invoke_content_directory( + &self, + action: &str, + op_name: Option<&str>, + args: Vec<(&'static str, String)>, + ) -> Result { + let op = op_name.unwrap_or(action); + let (control_url, service_type) = self.content_directory_endpoints(op)?; + let borrowed_args: Vec<(&str, &str)> = args.iter().map(|(k, v)| (*k, v.as_str())).collect(); + + let call_result = invoke_upnp_action_with_timeout( + control_url, + service_type, + action, + &borrowed_args, + Some(self.timeout), + )?; + + if !call_result.status.is_success() { + if let Some(env) = &call_result.envelope { + if let Some(err) = parse_upnp_error(env) { + if should_map_to_not_supported(action, op_name, err.error_code) { + return Err(server_op_not_supported(op, "UpnpMediaServer")); + } + + return Err(anyhow!( + "{} failed with UPnP error {}: {}", + action, + err.error_code, + err.error_description + )); + } + } + + return Err(anyhow!( + "{} failed with HTTP status {} and body: {}", + action, + call_result.status, + call_result.raw_body + )); + } + + if let Some(env) = &call_result.envelope { + if let Some(err) = parse_upnp_error(env) { + if should_map_to_not_supported(action, op_name, err.error_code) { + return Err(server_op_not_supported(op, "UpnpMediaServer")); + } + + return Err(anyhow!( + "{} returned UPnP error {}: {}", + action, + err.error_code, + err.error_description + )); + } + } + + Ok(call_result) + } + + fn content_directory_endpoints(&self, op_name: &str) -> Result<(&str, &str)> { + if !self.info.has_content_directory { + return Err(server_op_not_supported(op_name, "UpnpMediaServer")); + } + + let control_url = self + .info + .content_directory_control_url + .as_deref() + .ok_or_else(|| server_op_not_supported(op_name, "UpnpMediaServer"))?; + let service_type = self + .info + .content_directory_service_type + .as_deref() + .ok_or_else(|| server_op_not_supported(op_name, "UpnpMediaServer"))?; + + Ok((control_url, service_type)) + } +} + +impl MediaBrowser for UpnpMediaServer { + fn browse_root(&self) -> Result> { + self.browse_with_flag("0", "BrowseDirectChildren", 0, 0) + } + + fn browse_children(&self, object_id: &str, start: u32, count: u32) -> Result> { + self.browse_with_flag(object_id, "BrowseDirectChildren", start, count) + } + + fn browse_object(&self, object_id: &str) -> Result { + let entries = self.browse_with_flag(object_id, "BrowseMetadata", 0, 1)?; + entries + .into_iter() + .next() + .ok_or_else(|| anyhow!("Object {} was not returned by the server", object_id)) + } + + fn search( + &self, + container_id: &str, + query: &str, + start: u32, + count: u32, + ) -> Result> { + self.search_impl(container_id, query, start, count) + } +} + +fn map_didl_entries(xml: &str) -> Result> { + let trimmed = xml.trim(); + if trimmed.is_empty() { + return Ok(Vec::new()); + } + + let didl: DIDLLite = pmodidl::parse_metadata::(trimmed) + .map_err(|err| anyhow!("Failed to parse DIDL-Lite payload: {}", err))? + .data; + + let mut entries = Vec::new(); + + for container in didl.containers { + entries.push(MediaEntry { + id: container.id, + parent_id: container.parent_id, + title: container.title, + is_container: true, + class: container.class, + resources: Vec::new(), + }); + } + + for item in didl.items { + let resources = item + .resources + .into_iter() + .map(|res| MediaResource { + uri: res.url, + protocol_info: res.protocol_info, + duration: res.duration, + }) + .collect(); + + entries.push(MediaEntry { + id: item.id, + parent_id: item.parent_id, + title: item.title, + is_container: false, + class: item.class, + resources, + }); + } + + Ok(entries) +} + +fn extract_result_payload(envelope: &SoapEnvelope, response_suffix: &str) -> Result { + let response = find_child_with_suffix(&envelope.body.content, response_suffix) + .ok_or_else(|| anyhow!("Missing {} element in SOAP body", response_suffix))?; + let result_elem = find_child_with_suffix(response, "Result") + .ok_or_else(|| anyhow!("Missing Result element in {}", response_suffix))?; + + let payload = result_elem + .get_text() + .map(|t| t.to_string()) + .unwrap_or_default(); + + Ok(payload) +} + +fn find_child_with_suffix<'a>(parent: &'a Element, suffix: &str) -> Option<&'a Element> { + parent.children.iter().find_map(|node| match node { + XMLNode::Element(elem) if elem.name.ends_with(suffix) => Some(elem), + _ => None, + }) +} + +fn parse_upnp_error(envelope: &SoapEnvelope) -> Option { + let fault = find_child_with_suffix(&envelope.body.content, "Fault")?; + let detail = find_child_with_suffix(fault, "detail")?; + let upnp_error = find_child_with_suffix(detail, "UPnPError")?; + + let error_code_elem = upnp_error.children.iter().find_map(|node| match node { + XMLNode::Element(elem) if elem.name.ends_with("errorCode") => Some(elem), + _ => None, + })?; + let error_code_text = error_code_elem.get_text()?.trim().to_string(); + let error_code = error_code_text.parse::().ok()?; + + let error_description = upnp_error + .children + .iter() + .find_map(|node| match node { + XMLNode::Element(elem) if elem.name.ends_with("errorDescription") => { + elem.get_text().map(|t| t.trim().to_string()) + } + _ => None, + }) + .unwrap_or_default(); + + Some(UpnpError { + error_code, + error_description, + }) +} + +fn should_map_to_not_supported(action: &str, op_name: Option<&str>, error_code: u32) -> bool { + if op_name.is_none() { + return false; + } + + let optional_code = error_codes::OPTIONAL_ACTION_NOT_IMPLEMENTED + .parse::() + .unwrap_or(602); + let invalid_action_code = error_codes::INVALID_ACTION.parse::().unwrap_or(401); + + let cd_not_supported = matches!(action, "Search"); + + cd_not_supported && (error_code == optional_code || error_code == invalid_action_code) +} + +fn server_op_not_supported(op: &str, backend: &str) -> anyhow::Error { + anyhow!( + "MusicServer operation '{}' is not supported by backend '{}'", + op, + backend + ) +} + +#[derive(Debug, Clone)] +struct UpnpError { + pub error_code: u32, + pub error_description: String, +} diff --git a/pmocontrol/src/model.rs b/pmocontrol/src/model.rs index 02ae40b9..6bbf7ca4 100644 --- a/pmocontrol/src/model.rs +++ b/pmocontrol/src/model.rs @@ -3,9 +3,6 @@ use crate::capabilities::{PlaybackPositionInfo, PlaybackState}; #[derive(Clone, Debug, PartialEq, Eq, Hash)] pub struct RendererId(pub String); -#[derive(Clone, Debug, PartialEq, Eq, Hash)] -pub struct MediaServerId(pub String); - #[derive(Clone, Debug)] pub enum RendererProtocol { UpnpAvOnly, @@ -53,29 +50,6 @@ pub struct RendererInfo { pub connection_manager_control_url: Option, } -#[derive(Clone, Debug, Default)] -pub struct MediaServerCapabilities { - pub has_content_directory: bool, - pub has_connection_manager: bool, -} - -#[derive(Clone, Debug)] -pub struct MediaServerInfo { - pub id: MediaServerId, - pub udn: String, - pub friendly_name: String, - pub model_name: String, - pub manufacturer: String, - - pub capabilities: MediaServerCapabilities, - - pub location: String, - pub server_header: String, - pub online: bool, - pub last_seen: std::time::SystemTime, - pub max_age: u32, -} - #[derive(Clone, Debug)] pub enum RendererEvent { StateChanged { diff --git a/pmocontrol/src/music_renderer.rs b/pmocontrol/src/music_renderer.rs index d20f0b8f..3eb918a2 100644 --- a/pmocontrol/src/music_renderer.rs +++ b/pmocontrol/src/music_renderer.rs @@ -12,7 +12,7 @@ use crate::{ ArylicTcpRenderer, DeviceRegistry, LinkPlayRenderer, PlaybackPosition, PlaybackState, TransportControl, UpnpRenderer, VolumeControl, }; -use anyhow::{anyhow, Result}; +use anyhow::{Result, anyhow}; use tracing::warn; /// Backend-agnostic façade exposing transport, volume, and status contracts. @@ -33,7 +33,7 @@ pub enum MusicRenderer { } /// Build a standardized error when an operation is not supported by a backend. -fn op_not_supported(op: &str, backend: &str) -> anyhow::Error { +pub(crate) fn op_not_supported(op: &str, backend: &str) -> anyhow::Error { anyhow!( "MusicRenderer operation '{}' is not supported by backend '{}'", op, diff --git a/pmocontrol/src/playback_queue.rs b/pmocontrol/src/playback_queue.rs new file mode 100644 index 00000000..df1ae0b2 --- /dev/null +++ b/pmocontrol/src/playback_queue.rs @@ -0,0 +1,73 @@ +use std::collections::VecDeque; + +use crate::media_server::ServerId; + +#[derive(Clone, Debug)] +pub struct PlaybackItem { + pub uri: String, + pub title: Option, + pub server_id: Option, + pub object_id: Option, +} + +impl PlaybackItem { + pub fn new(uri: impl Into) -> Self { + Self { + uri: uri.into(), + title: None, + server_id: None, + object_id: None, + } + } +} + +#[derive(Clone, Debug, Default)] +pub struct PlaybackQueue { + items: VecDeque, +} + +impl PlaybackQueue { + pub fn new() -> Self { + Self { + items: VecDeque::new(), + } + } + + pub fn len(&self) -> usize { + self.items.len() + } + + pub fn is_empty(&self) -> bool { + self.items.is_empty() + } + + pub fn clear(&mut self) { + self.items.clear(); + } + + pub fn enqueue(&mut self, item: PlaybackItem) { + self.items.push_back(item); + } + + pub fn enqueue_many>(&mut self, items: I) { + for item in items { + self.items.push_back(item); + } + } + + pub fn enqueue_front(&mut self, item: PlaybackItem) { + self.items.push_front(item); + } + + pub fn dequeue(&mut self) -> Option { + self.items.pop_front() + } + + pub fn peek(&self) -> Option<&PlaybackItem> { + self.items.front() + } + + pub fn snapshot(&self) -> Vec { + self.items.iter().cloned().collect() + } +} diff --git a/pmocontrol/src/provider.rs b/pmocontrol/src/provider.rs index 6958e6bb..f511f2b3 100644 --- a/pmocontrol/src/provider.rs +++ b/pmocontrol/src/provider.rs @@ -9,10 +9,8 @@ use crate::arylic_tcp::detect_arylic_tcp; use crate::avtransport_client::AvTransportClient; use crate::discovery::{DeviceDescriptionProvider, DiscoveredEndpoint}; use crate::linkplay::detect_linkplay_http; -use crate::model::{ - MediaServerCapabilities, MediaServerId, MediaServerInfo, RendererCapabilities, RendererId, - RendererInfo, RendererProtocol, -}; +use crate::media_server::{MediaServerInfo, ServerId}; +use crate::model::{RendererCapabilities, RendererId, RendererInfo, RendererProtocol}; use ureq::Agent; @@ -52,6 +50,10 @@ struct ParsedDeviceDescription { // ConnectionManager endpoint (if present in serviceList) connection_manager_service_type: Option, connection_manager_control_url: Option, + + // ContentDirectory endpoint (if present in serviceList) + content_directory_service_type: Option, + content_directory_control_url: Option, } impl ParsedDeviceDescription { @@ -201,6 +203,21 @@ impl HttpXmlDescriptionProvider { ); } } + + if lower + .contains("urn:schemas-upnp-org:service:contentdirectory:") + { + if parsed.content_directory_service_type.is_none() { + parsed.content_directory_service_type = + Some(st.clone()); + parsed.content_directory_control_url = + Some(ctrl.clone()); + debug!( + "Found ContentDirectory service for {}: type={} controlURL={}", + endpoint.udn, st, ctrl + ); + } + } } in_service = false; @@ -340,21 +357,31 @@ impl HttpXmlDescriptionProvider { .as_deref() .unwrap_or_else(|| endpoint.udn.as_str()); let udn = raw_udn.to_ascii_lowercase(); - let caps = detect_server_capabilities(&parsed.service_types); + let has_content_directory = parsed.service_types.iter().any(|st| { + st.to_ascii_lowercase() + .contains("urn:schemas-upnp-org:service:contentdirectory:") + }); let now = SystemTime::now(); + let content_directory_control_url = parsed + .content_directory_control_url + .as_ref() + .map(|ctrl| resolve_control_url(&endpoint.location, ctrl)); + Some(MediaServerInfo { - id: MediaServerId(udn.clone()), + id: ServerId(udn.clone()), udn, friendly_name: parsed.friendly_name.clone().unwrap_or_default(), model_name: parsed.model_name.clone().unwrap_or_default(), manufacturer: parsed.manufacturer.clone().unwrap_or_default(), - capabilities: caps, location: endpoint.location.clone(), server_header: endpoint.server_header.clone(), online: true, last_seen: now, max_age: endpoint.max_age, + has_content_directory, + content_directory_service_type: parsed.content_directory_service_type.clone(), + content_directory_control_url, }) } @@ -442,22 +469,6 @@ fn detect_renderer_protocol(caps: &RendererCapabilities) -> RendererProtocol { } } -fn detect_server_capabilities(service_types: &[String]) -> MediaServerCapabilities { - let mut caps = MediaServerCapabilities::default(); - - for st in service_types { - let lower = st.to_ascii_lowercase(); - if lower.contains("urn:schemas-upnp-org:service:contentdirectory:") { - caps.has_content_directory = true; - } - if lower.contains("urn:schemas-upnp-org:service:connectionmanager:") { - caps.has_connection_manager = true; - } - } - - caps -} - /// Resolve a possibly relative controlURL against the description URL. /// /// - If `control_url` is already absolute (starts with http:// or https://), it is returned as-is. diff --git a/pmocontrol/src/registry.rs b/pmocontrol/src/registry.rs index e9e25f7d..ebcaea2f 100644 --- a/pmocontrol/src/registry.rs +++ b/pmocontrol/src/registry.rs @@ -3,19 +3,20 @@ use std::time::SystemTime; use crate::avtransport_client::AvTransportClient; use crate::connection_manager_client::ConnectionManagerClient; -use crate::model::{MediaServerId, MediaServerInfo, RendererId, RendererInfo}; +use crate::media_server::{MediaServerInfo, ServerId}; +use crate::model::{RendererId, RendererInfo}; use crate::rendering_control_client::RenderingControlClient; #[derive(Clone, Debug)] enum DeviceKey { Renderer(RendererId), - Server(MediaServerId), + Server(ServerId), } #[derive(Debug, Default)] pub struct DeviceRegistry { renderers: HashMap, - servers: HashMap, + servers: HashMap, udn_index: HashMap, } @@ -28,7 +29,7 @@ pub trait DeviceRegistryRead { fn list_servers(&self) -> Vec; fn get_renderer(&self, id: &RendererId) -> Option; - fn get_server(&self, id: &MediaServerId) -> Option; + fn get_server(&self, id: &ServerId) -> Option; } impl DeviceRegistryRead for DeviceRegistry { @@ -44,7 +45,7 @@ impl DeviceRegistryRead for DeviceRegistry { self.renderers.get(id).cloned() } - fn get_server(&self, id: &MediaServerId) -> Option { + fn get_server(&self, id: &ServerId) -> Option { self.servers.get(id).cloned() } } @@ -56,7 +57,7 @@ pub enum DeviceUpdate { RendererOfflineByUdn(String), ServerOnline(MediaServerInfo), - ServerOfflineById(MediaServerId), + ServerOfflineById(ServerId), ServerOfflineByUdn(String), } diff --git a/pmocontrol/src/soap_client.rs b/pmocontrol/src/soap_client.rs index 8906b33d..77c0a5a3 100644 --- a/pmocontrol/src/soap_client.rs +++ b/pmocontrol/src/soap_client.rs @@ -1,3 +1,5 @@ +use std::time::Duration; + use anyhow::{Context, Result}; use pmoupnp::soap::{SoapEnvelope, build_soap_request, parse_soap_envelope}; use ureq::Agent; @@ -24,13 +26,26 @@ pub fn invoke_upnp_action( action: &str, args: &[(&str, &str)], ) -> Result { - // 1. Build SOAP request XML using pmoupnp + invoke_upnp_action_with_timeout(control_url, service_type, action, args, None) +} + +pub fn invoke_upnp_action_with_timeout( + control_url: &str, + service_type: &str, + action: &str, + args: &[(&str, &str)], + timeout: Option, +) -> Result { let body_xml = build_soap_request(service_type, action, args) .context("Failed to build SOAP request body")?; - // 2. Build agent that does NOT treat 4xx/5xx as errors - let config = Agent::config_builder().http_status_as_error(false).build(); + let mut builder = Agent::config_builder(); + builder = builder.http_status_as_error(false); + if let Some(duration) = timeout { + builder = builder.timeout_global(Some(duration)); + } + let config = builder.build(); let agent: Agent = config.into(); // 3. SOAPAction header