From 1a7e5dd790775c142760c2e6fe6ebd9b439b7bdd Mon Sep 17 00:00:00 2001 From: Eric Coissac Date: Wed, 17 Dec 2025 21:41:24 +0100 Subject: [PATCH] Debug openhome --- .DS_Store | Bin 18436 -> 16388 bytes .../webapp/src/services/pmocontrol/types.ts | 1 + pmocache/src/api.rs | 8 +- pmocontrol/examples/openhome_tester.rs | 320 ++++++++ pmocontrol/examples/test_attach_playlist.rs | 16 +- pmocontrol/src/control_point.rs | 737 +++++++++++++----- pmocontrol/src/control_point/music_queue.rs | 13 +- .../src/control_point/openhome_queue.rs | 227 +++++- pmocontrol/src/lib.rs | 1 + pmocontrol/src/model.rs | 2 + pmocontrol/src/music_renderer.rs | 114 ++- pmocontrol/src/openhome.rs | 187 +++++ pmocontrol/src/openhome_client.rs | 451 +++++++++-- pmocontrol/src/openhome_playlist.rs | 2 + pmocontrol/src/openhome_renderer.rs | 154 ++-- pmocontrol/src/pmoserver_ext.rs | 9 +- pmocontrol/src/provider.rs | 40 +- pmocontrol/src/soap_client.rs | 159 +++- pmocontrol/src/upnp_renderer.rs | 2 + pmocovers/src/api.rs | 7 +- pmoserver/src/server.rs | 8 +- pmoupnp/src/upnp_server.rs | 4 +- 22 files changed, 2015 insertions(+), 447 deletions(-) create mode 100644 pmocontrol/examples/openhome_tester.rs create mode 100644 pmocontrol/src/openhome.rs diff --git a/.DS_Store b/.DS_Store index 563111b848472da02223eb997491b0132953068d..20e2b30943e4506e0ee53bb9554f2bf8dff51d4b 100644 GIT binary patch delta 247 zcmZpfz}V8jsKUU&!01$(V#vT?0c3IjF^Cr6VPIh4Vqj1J(gxfURkRr!HwGMLo?K`k zJ-NV;pOJNQfQ=IKBD z+$rVF%xJoqM?sxwa*xs^M%Kv=N)nsXRW>qi?owl8+I&Pqj(xH)%OfC`+^k`s%Q&&b zfKhC7jp16x&8p@)teZtN6&U3i7}y|oSpcnNa7A|~(EiuVllfhId4Ns=c?Aq6r}Ok~ NPPF>K2sA^U5daOoF{%In delta 1088 zcmcIjO-vI}5T3VNY|~b0MSfhNR8nYYY2DqnP!pn{(P)ebCoUdJ+oeb-Yz4|sOq2sA z9)(FXUOec{gGe?W@Z!O<-aMO_c=F`M#0&bpttOP<(S3P4^XAR_zIktEHaEoPBM~Da zaj9+@L}O5Wm{Nv@m;6jBdv#ggLb?!D>cE>Nn+p6f-NpA} z;oc@ImElvS+ca%0m8&J&uJXCQ$M^kI0!t&R0xcrC8Z9<+s+T3_ffzGKoxzQcK!UZ1 zxbhfTU>C_i3_0j|$UNmVY^I>jApA7xn2Uc`JT-L6YT`})gP{x%d4qLY#aO2bW!$v| zP?HvD1dww8uY++G7;6V<$d)+}`_9|GPvXWPTy2_%LkaFBjH;V38=x|2MZj5z(FTxp z$g7U0Tfi*fj7>8f&1#Bz!P$(3m3YOqts?blpw9wpzWI^GPdBj-2a?{cl|hl<*Qsvj zRXQcO&699i^_N2E{5{ge9cB_GXvz}Q8YOtBBC*U$ZhmgYq357?@+~s~eR$jK;AdtM zPy$dNB4Js=*vNAz)5-qc5a0F6^?9tvRfBX%bG2dS;dZS>mS2Qx-hN^5iOgx9^CsYi zhS2y&CJr?#Pwco>$sFf*Ud`DvhFYA9p0g?oBg*B}#MVlG&oQ^&X)FMF<#6|*b-t5a~s_#>yJ(State(cache): State>>) -> impl IntoResponse { +pub async fn list_items( + State(cache): State>>, +) -> impl IntoResponse { match cache.db.get_all(true) { Ok(entries) => (StatusCode::OK, Json(entries)).into_response(), Err(e) => ( @@ -389,7 +391,9 @@ pub async fn delete_item( /// Purge complètement le cache /// /// Supprime tous les items et vide la base de données. Opération irréversible. -pub async fn purge_cache(State(cache): State>>) -> impl IntoResponse { +pub async fn purge_cache( + State(cache): State>>, +) -> impl IntoResponse { match cache.purge().await { Ok(_) => ( StatusCode::OK, diff --git a/pmocontrol/examples/openhome_tester.rs b/pmocontrol/examples/openhome_tester.rs new file mode 100644 index 00000000..c78f1195 --- /dev/null +++ b/pmocontrol/examples/openhome_tester.rs @@ -0,0 +1,320 @@ +use std::collections::HashSet; +use std::env; +use std::thread; +use std::time::{Duration, Instant}; + +use anyhow::{Context, Result, anyhow}; +use pmocontrol::{ + DeviceDescriptionProvider, DiscoveredEndpoint, HttpXmlDescriptionProvider, MusicRenderer, + RendererInfo, + control_point::ControlPoint, + openhome_client::{ + OPENHOME_PLAYLIST_HEAD_ID, OhInfoClient, OhPlaylistClient, OhProductClient, OhRadioClient, + OhTimeClient, OhVolumeClient, parse_track_metadata_from_didl, + }, +}; + +fn main() -> Result<()> { + let _ = tracing_subscriber::fmt() + .with_env_filter(tracing_subscriber::EnvFilter::from_default_env()) + .try_init(); + + let args: Vec = env::args().collect(); + + let renderer = match args.len() { + 1 => auto_discover_renderer() + .context("Unable to auto-discover the OpenHome renderer; rerun with explicit UDN and description URL")?, + 3 => { + let udn = args[1].to_ascii_lowercase(); + let description_url = args[2].clone(); + renderer_from_description(&udn, &description_url)? + } + _ => { + eprintln!( + "Usage:\n {0} # auto-discover the single OpenHome renderer\n {0} ", + args[0] + ); + std::process::exit(1); + } + }; + + println!( + "Renderer: {} ({})", + renderer.friendly_name, renderer.model_name + ); + println!("UDN: {}", renderer.udn); + println!( + "OpenHome services -> playlist:{} info:{} time:{} volume:{} radio:{} product:{}", + renderer.oh_playlist_control_url.is_some(), + renderer.oh_info_control_url.is_some(), + renderer.oh_time_control_url.is_some(), + renderer.oh_volume_control_url.is_some(), + renderer.oh_radio_control_url.is_some(), + renderer.oh_product_control_url.is_some(), + ); + + if let Some(result) = test_product(&renderer) { + result?; + } + if let Some(result) = test_playlist(&renderer) { + result?; + } + if let Some(result) = test_info(&renderer) { + result?; + } + if let Some(result) = test_time(&renderer) { + result?; + } + if let Some(result) = test_volume(&renderer) { + result?; + } + if let Some(result) = test_radio(&renderer) { + result?; + } + + Ok(()) +} + +fn auto_discover_renderer() -> Result { + println!("Auto-discovering OpenHome renderer via SSDP…"); + let control_point = ControlPoint::spawn(5).context("Failed to start control point")?; + let deadline = Instant::now() + Duration::from_secs(12); + + while Instant::now() < deadline { + let renderers = control_point.list_music_renderers(); + let mut seen = HashSet::new(); + let mut openhome_infos = Vec::new(); + + for renderer in renderers { + if let MusicRenderer::OpenHome(oh) = renderer { + if seen.insert(oh.id().0.clone()) { + openhome_infos.push(oh.info.clone()); + } + } + } + + match openhome_infos.len() { + 0 => { + thread::sleep(Duration::from_millis(500)); + } + 1 => { + let info = openhome_infos.remove(0); + println!( + "Discovered OpenHome renderer '{}' ({})", + info.friendly_name, info.udn + ); + return Ok(info); + } + _ => { + println!("Detected multiple OpenHome renderers:"); + for info in &openhome_infos { + println!(" - {} ({})", info.friendly_name, info.udn); + } + return Err(anyhow!( + "Multiple OpenHome renderers present; rerun with explicit UDN + description URL." + )); + } + } + } + + Err(anyhow!("No OpenHome renderer discovered on the network.")) +} + +fn renderer_from_description(udn: &str, description_url: &str) -> Result { + let endpoint = DiscoveredEndpoint::new( + udn.to_ascii_lowercase(), + description_url.to_string(), + "openhome-tester/1.0".into(), + 1800, + ); + let provider = HttpXmlDescriptionProvider::new(5); + provider + .build_renderer_info(&endpoint) + .with_context(|| format!("Device {} is not a usable renderer", description_url)) +} + +fn test_product(info: &RendererInfo) -> Option> { + let control_url = info.oh_product_control_url.as_ref()?; + let service_type = info.oh_product_service_type.as_ref()?; + let client = OhProductClient::new(control_url.clone(), service_type.clone()); + + Some((|| { + println!("\n[Product] Listing sources…"); + let sources = client.source_xml()?; + for (idx, source) in sources.iter().enumerate() { + println!( + " #{idx} {} (type={}, visible={})", + source.name, source.source_type, source.visible + ); + } + let current = client.source_index()?; + println!("[Product] Current source index: {current}"); + client.ensure_playlist_source_selected()?; + println!("[Product] Playlist source is now selected"); + Ok(()) + })()) +} + +fn test_playlist(info: &RendererInfo) -> Option> { + let control_url = info.oh_playlist_control_url.as_ref()?; + let service_type = info.oh_playlist_service_type.as_ref()?; + let client = OhPlaylistClient::new(control_url.clone(), service_type.clone()); + + Some((|| { + println!("\n[Playlist] TracksMax={}", client.tracks_max()?); + let ids = client.id_array()?; + println!( + "[Playlist] IdArray contains {} entries (showing up to 5)", + ids.len() + ); + if ids.is_empty() { + println!("[Playlist] Playlist is empty"); + } else { + let preview = &ids[..ids.len().min(5)]; + let entries = client.read_list(preview)?; + for entry in entries { + println!( + " - id={} uri={} title={}", + entry.id, + entry.uri, + parse_track_metadata_from_didl(&entry.metadata_xml) + .and_then(|m| m.title) + .unwrap_or_else(|| "".into()) + ); + } + } + + exercise_playlist_mutations(&client)?; + Ok(()) + })()) +} + +fn exercise_playlist_mutations(client: &OhPlaylistClient) -> Result<()> { + println!("\n[Playlist] Exercising DeleteAll → Insert → DeleteAll"); + client.delete_all()?; + println!(" DeleteAll succeeded"); + std::thread::sleep(Duration::from_millis(200)); + let _ = client.id_array()?; + + let test_uri = env::var("OPENHOME_TEST_URI") + .unwrap_or_else(|_| "http://ice1.somafm.com/groovesalad-128-mp3".into()); + let metadata = build_sample_metadata_xml(&test_uri); + + println!( + " Inserting sample track at head (after_id = {:#x})…", + OPENHOME_PLAYLIST_HEAD_ID + ); + let new_id = match client.insert(OPENHOME_PLAYLIST_HEAD_ID, &test_uri, &metadata) { + Ok(id) => { + println!(" Insert succeeded with id {}", id); + id + } + Err(err) => { + println!( + " Insert with sentinel failed: {}. Retrying with aAfterId=0…", + err + ); + let id = client.insert(0, &test_uri, &metadata)?; + println!(" Insert with aAfterId=0 succeeded with id {}", id); + id + } + }; + + let ids = client.id_array()?; + println!(" Playlist now reports {} id(s): {:?}", ids.len(), ids); + + let entries = client.read_list(&[new_id])?; + if let Some(entry) = entries.first() { + println!(" ReadList confirms id={} uri={}", entry.id, entry.uri); + } + + client.delete_all()?; + println!(" Playlist restored to empty state"); + Ok(()) +} + +fn build_sample_metadata_xml(uri: &str) -> String { + format!( + concat!( + r#""#, + r#""#, + r#"PMO Test Track"#, + r#""#, + "{uri}", + r#""#, + r#"object.item.audioItem.musicTrack"#, + r#""# + ), + uri = uri + ) +} + +fn test_info(info: &RendererInfo) -> Option> { + let control_url = info.oh_info_control_url.as_ref()?; + let service_type = info.oh_info_service_type.as_ref()?; + let client = OhInfoClient::new(control_url.clone(), service_type.clone()); + + Some((|| { + println!("\n[Info] Reading current track…"); + let track = client.track()?; + println!(" URI: {}", track.uri); + if let Some(metadata) = track.metadata() { + println!( + " Metadata: {} – {}", + metadata.artist.as_deref().unwrap_or(""), + metadata.title.as_deref().unwrap_or("") + ); + } + match client.transport_state() { + Ok(state) => println!(" Transport state: {}", state), + Err(err) => println!(" TransportState call failed: {err}"), + } + Ok(()) + })()) +} + +fn test_time(info: &RendererInfo) -> Option> { + let control_url = info.oh_time_control_url.as_ref()?; + let service_type = info.oh_time_service_type.as_ref()?; + let client = OhTimeClient::new(control_url.clone(), service_type.clone()); + + Some((|| { + println!("\n[Time] Querying position…"); + let pos = client.position()?; + println!( + " Tracks={} duration={}s elapsed={}s", + pos.track_count, pos.duration_secs, pos.elapsed_secs + ); + Ok(()) + })()) +} + +fn test_volume(info: &RendererInfo) -> Option> { + let control_url = info.oh_volume_control_url.as_ref()?; + let service_type = info.oh_volume_service_type.as_ref()?; + let client = OhVolumeClient::new(control_url.clone(), service_type.clone()); + + Some((|| { + println!("\n[Volume] Current volume: {}", client.volume()?); + println!("[Volume] Muted: {}", client.mute()?); + Ok(()) + })()) +} + +fn test_radio(info: &RendererInfo) -> Option> { + let control_url = info.oh_radio_control_url.as_ref()?; + let service_type = info.oh_radio_service_type.as_ref()?; + let client = OhRadioClient::new(control_url.clone(), service_type.clone()); + + Some((|| { + println!("\n[Radio] Fetching channel #0 metadata…"); + let channel = client.channel(0)?; + println!(" URI: {}", channel.uri); + if let Some(meta) = channel.metadata_xml { + println!(" Metadata XML: {}", meta); + } + Ok(()) + })()) +} diff --git a/pmocontrol/examples/test_attach_playlist.rs b/pmocontrol/examples/test_attach_playlist.rs index e4cef051..a7797451 100644 --- a/pmocontrol/examples/test_attach_playlist.rs +++ b/pmocontrol/examples/test_attach_playlist.rs @@ -7,14 +7,13 @@ /// /// Usage: /// cargo run --example test_attach_playlist - use anyhow::{Context, Result}; use pmocontrol::{ media_server::{MediaBrowser, MediaEntry, MediaServer, ServerId}, openhome_client::OhPlaylistClient, }; use tracing::{debug, error, info, warn}; -use tracing_subscriber::{fmt, prelude::*, EnvFilter}; +use tracing_subscriber::{EnvFilter, fmt, prelude::*}; #[tokio::main] async fn main() -> Result<()> { @@ -47,7 +46,8 @@ async fn main() -> Result<()> { info!("📡 Step 1: Connecting to MediaServer"); let content_directory_url = format!( "{}/device/{}/service/ContentDirectory/control", - media_server_url, media_server_id.replace("uuid:", "") + media_server_url, + media_server_id.replace("uuid:", "") ); let server = MediaServer::new( @@ -82,7 +82,8 @@ async fn main() -> Result<()> { // Display the first few items info!("📋 First items in playlist:"); for (idx, entry) in entries.iter().take(3).enumerate() { - info!(" [{}] {} - {}", + info!( + " [{}] {} - {}", idx, entry.title, entry.artist.as_deref().unwrap_or("Unknown") @@ -163,8 +164,11 @@ async fn main() -> Result<()> { } info!(""); - info!("✅ Test completed: {}/{} items inserted successfully", - inserted_count, entries.len()); + info!( + "✅ Test completed: {}/{} items inserted successfully", + inserted_count, + entries.len() + ); Ok(()) } diff --git a/pmocontrol/src/control_point.rs b/pmocontrol/src/control_point.rs index 11d8077d..2a2d04e6 100644 --- a/pmocontrol/src/control_point.rs +++ b/pmocontrol/src/control_point.rs @@ -8,20 +8,21 @@ use std::sync::{Arc, Mutex, MutexGuard, RwLock}; use std::thread; use std::time::{Duration, Instant}; -use anyhow::{Context, anyhow}; -use crossbeam_channel::{Receiver, Sender, unbounded}; +use anyhow::{anyhow, Context}; +use crossbeam_channel::{unbounded, Receiver, Sender}; use pmodidl::{DIDLLite, Item as DidlItem, Resource as DidlResource}; use pmoupnp::ssdp::SsdpClient; use quick_xml::se::to_string as to_didl_string; use thiserror::Error; use tracing::{debug, error, info, trace, warn}; -use ureq::{Agent, http}; +use ureq::{http, Agent}; use xmltree::{Element, XMLNode}; pub mod music_queue; pub mod openhome_queue; -use crate::MusicRenderer; +pub const OPENHOME_SNAPSHOT_CACHE_TTL: Duration = Duration::from_secs(2); + use crate::capabilities::{ PlaybackPosition, PlaybackPositionInfo, PlaybackState, PlaybackStatus, TransportControl, VolumeControl, @@ -35,14 +36,17 @@ use crate::media_server_events::spawn_media_server_event_runtime; use crate::model::TrackMetadata; use crate::model::{MediaServerEvent, RendererEvent, RendererId, RendererInfo}; use crate::music_renderer::{ - OpenHomeQueueProvider, RendererRuntimeState, set_openhome_queue_provider, + set_openhome_queue_provider, OpenHomeQueueProvider, RendererRuntimeState, }; #[cfg(feature = "pmoserver")] use crate::openapi::{ CurrentTrackMetadata, FullRendererSnapshot, QueueItem, QueueSnapshotView, RendererBindingView, RendererStateView, }; -use crate::openhome_client::{OhInfoClient, OhPlaylistClient, parse_track_metadata_from_didl}; +use crate::openhome::{ + build_info_client, build_playlist_client, build_product_client, OhServiceKind, +}; +use crate::openhome_client::parse_track_metadata_from_didl; use crate::openhome_playlist::{OpenHomePlaylistSnapshot, OpenHomePlaylistTrack}; use crate::openhome_renderer::{format_seconds, map_openhome_state}; use crate::provider::HttpXmlDescriptionProvider; @@ -50,6 +54,7 @@ use crate::queue_backend::{EnqueueMode, PlaybackItem, QueueBackend}; use crate::queue_interne::InternalQueue; use crate::registry::{DeviceRegistry, DeviceRegistryRead, DeviceUpdate}; use crate::upnp_renderer::UpnpRenderer; +use crate::MusicRenderer; /// Optional attachment between a renderer playback queue and a server-side /// DIDL-Lite playlist container. @@ -391,70 +396,68 @@ impl ControlPoint { thread::Builder::new() .name("cp-media-server-event-worker".into()) - .spawn(move || { - loop { - let event = match media_rx.recv() { - Ok(e) => e, - Err(_) => { - warn!("MediaServerEvent channel closed, worker exiting"); - break; - } - }; + .spawn(move || loop { + let event = match media_rx.recv() { + Ok(e) => e, + Err(_) => { + warn!("MediaServerEvent channel closed, worker exiting"); + break; + } + }; - match event { - MediaServerEvent::GlobalUpdated { - server_id, - system_update_id, - } => { - info!( - server = server_id.0.as_str(), - system_update_id = system_update_id, - "MediaServer global update" - ); - } - MediaServerEvent::ContainersUpdated { - server_id, - container_ids, - } => { - let renderers_to_refresh: Vec = { - let mut bindings = bindings_for_media_worker.lock().unwrap(); - let mut to_refresh = Vec::new(); + match event { + MediaServerEvent::GlobalUpdated { + server_id, + system_update_id, + } => { + info!( + server = server_id.0.as_str(), + system_update_id = system_update_id, + "MediaServer global update" + ); + } + MediaServerEvent::ContainersUpdated { + server_id, + container_ids, + } => { + let renderers_to_refresh: Vec = { + let mut bindings = bindings_for_media_worker.lock().unwrap(); + let mut to_refresh = Vec::new(); - for (renderer_id, binding) in bindings.iter_mut() { - if binding.server_id == server_id - && container_ids.contains(&binding.container_id) - { - binding.pending_refresh = true; - binding.has_seen_update = true; - to_refresh.push(renderer_id.clone()); - } + for (renderer_id, binding) in bindings.iter_mut() { + if binding.server_id == server_id + && container_ids.contains(&binding.container_id) + { + binding.pending_refresh = true; + binding.has_seen_update = true; + to_refresh.push(renderer_id.clone()); } + } - to_refresh - }; + to_refresh + }; - for renderer_id in renderers_to_refresh { - debug!( + for renderer_id in renderers_to_refresh { + debug!( + renderer = renderer_id.0.as_str(), + server = server_id.0.as_str(), + "Triggering queue refresh for bound playlist" + ); + + if let Err(err) = refresh_attached_queue_for( + ®istry_for_media_worker, + &runtime_for_media_worker, + &bindings_for_media_worker, + &renderer_id, + &event_bus_for_media_worker, + None, + ) { + warn!( renderer = renderer_id.0.as_str(), server = server_id.0.as_str(), - "Triggering queue refresh for bound playlist" + error = %err, + "Failed to refresh queue from playlist container" ); - - if let Err(err) = refresh_attached_queue_for( - ®istry_for_media_worker, - &runtime_for_media_worker, - &bindings_for_media_worker, - &renderer_id, - &event_bus_for_media_worker, - None, - ) { - warn!( - renderer = renderer_id.0.as_str(), - server = server_id.0.as_str(), - error = %err, - "Failed to refresh queue from playlist container" - ); - } } } } @@ -782,6 +785,15 @@ impl ControlPoint { renderer.openhome_playlist_snapshot() } + pub fn get_cached_openhome_playlist_snapshot( + &self, + renderer_id: &RendererId, + ttl: Duration, + ) -> anyhow::Result { + let renderer = self.openhome_renderer(renderer_id)?; + self.runtime.openhome_snapshot_cached(&renderer, ttl) + } + pub fn get_openhome_playlist_len(&self, renderer_id: &RendererId) -> anyhow::Result { let renderer = self.openhome_renderer(renderer_id)?; renderer.openhome_playlist_len() @@ -798,12 +810,122 @@ impl ControlPoint { .ok_or_else(|| anyhow!("Renderer {} not found", renderer_id.0))?; let info = renderer.info(); - let (runtime_snapshot, queue_items, current_index) = + // MICRO-PATCH 5: Chemin OpenHome complètement découplé du miroir local + if self.runtime.uses_openhome_playlist(renderer_id) { + // Pour OpenHome: récupérer UNIQUEMENT runtime_snapshot pour volume/state/position + // Ne PAS utiliser queue_items/current_index du runtime (miroir local) + let (runtime_snapshot, _, _) = self.runtime.renderer_snapshot_bundle(renderer_id); + + // OpenHome est la source de vérité pour la queue - pas de fallback au runtime + let snapshot = self.get_cached_openhome_playlist_snapshot( + renderer_id, + OPENHOME_SNAPSHOT_CACHE_TTL, + )?; + + let queue_items: Vec = snapshot + .tracks + .iter() + .map(|track| playback_item_from_openhome_track(renderer_id, track)) + .collect(); + + let queue_len = snapshot.tracks.len(); + + // Pour OpenHome: current_index vient UNIQUEMENT d'OpenHome, pas d'heuristiques runtime + let queue_current_index = snapshot.current_index.or_else(|| { + snapshot.current_id.and_then(|id| { + snapshot.tracks.iter().position(|track| track.id == id) + }) + }); + + debug!( + renderer = renderer_id.0.as_str(), + current_id = ?snapshot.current_id, + current_index = ?queue_current_index, + track_count = snapshot.tracks.len(), + "renderer_full_snapshot: OpenHome snapshot retrieved" + ); + + let queue_view_items: Vec = queue_items + .iter() + .enumerate() + .map(|(index, item)| QueueItem { + index, + uri: item.uri.clone(), + title: item.metadata.as_ref().and_then(|m| m.title.clone()), + artist: item.metadata.as_ref().and_then(|m| m.artist.clone()), + album: item.metadata.as_ref().and_then(|m| m.album.clone()), + album_art_uri: item.metadata.as_ref().and_then(|m| m.album_art_uri.clone()), + server_id: Some(item.media_server_id.0.clone()), + object_id: Some(item.didl_id.clone()), + }) + .collect(); + + let queue_view = QueueSnapshotView { + renderer_id: renderer_id.0.clone(), + items: queue_view_items, + current_index: queue_current_index, + }; + + let binding = self.current_queue_playlist_binding(renderer_id).map( + |(server_id, container_id, has_seen_update)| RendererBindingView { + server_id: server_id.0, + container_id, + has_seen_update, + }, + ); + + let (position_ms, duration_ms) = + convert_runtime_position(runtime_snapshot.position.as_ref()); + let queue_current_metadata = queue_current_index + .and_then(|idx| queue_items.get(idx)) + .map(current_track_from_playback_item); + + // MICRO-PATCH 5: Pour OpenHome, préférer les métadonnées depuis le snapshot OpenHome + // car runtime_snapshot.last_metadata n'est jamais mis à jour pour OpenHome + let current_track = queue_current_metadata.or_else(|| { + runtime_snapshot + .last_metadata + .as_ref() + .map(|meta| CurrentTrackMetadata { + title: meta.title.clone(), + artist: meta.artist.clone(), + album: meta.album.clone(), + album_art_uri: meta.album_art_uri.clone(), + }) + }); + + let state_view = RendererStateView { + id: renderer_id.0.clone(), + friendly_name: info.friendly_name.clone(), + transport_state: runtime_snapshot + .state + .as_ref() + .map(|state| state.as_str().to_string()) + .unwrap_or_else(|| "UNKNOWN".to_string()), + position_ms, + duration_ms, + volume: runtime_snapshot + .last_volume + .and_then(|value| u8::try_from(value).ok()), + mute: runtime_snapshot.last_mute, + queue_len, + attached_playlist: binding.clone(), + current_track, + }; + + return Ok(FullRendererSnapshot { + state: state_view, + queue: queue_view, + binding, + }); + } + + // Chemin non-OpenHome (UPnP AV): continue d'utiliser le runtime + let (runtime_snapshot, queue_items, mut queue_current_index) = self.runtime.renderer_snapshot_bundle(renderer_id); let playback_source = self.runtime.playback_source(renderer_id); let queue_len = queue_items.len(); - let mut queue_current_index = current_index; if queue_current_index.is_none() { if let Some(position) = runtime_snapshot.position.as_ref() { if let Some(uri) = position.track_uri.as_ref() { @@ -1159,7 +1281,11 @@ impl ControlPoint { return Ok(()); } - self.play_next_from_queue(renderer_id) + if self.runtime.uses_openhome_playlist(renderer_id) { + self.play_current_from_queue(renderer_id) + } else { + self.play_next_from_queue(renderer_id) + } } /// Stop playback in response to user action (e.g., Stop button in UI). @@ -1442,63 +1568,119 @@ impl ControlPoint { fn play_current_openhome(&self, renderer_id: &RendererId) -> anyhow::Result<()> { let renderer = self.openhome_renderer(renderer_id)?; - // Pour les renderers OpenHome, vérifier d'abord si la playlist native a des pistes - // Cela couvre le cas où l'utilisateur a ajouté des morceaux directement via l'interface OpenHome - let native_playlist_len = renderer.openhome_playlist_len().unwrap_or(0); + // MICRO-PATCH 5: Pour OpenHome, récupérer les données directement depuis la playlist native + // au lieu du miroir local (qui peut être vide ou obsolète) + match self.get_cached_openhome_playlist_snapshot( + renderer_id, + OPENHOME_SNAPSHOT_CACHE_TTL, + ) { + Ok(snapshot) => { + if snapshot.tracks.is_empty() { + debug!( + renderer = renderer_id.0.as_str(), + "OpenHome playlist is empty" + ); + self.runtime + .set_playback_source(renderer_id, PlaybackSource::None); + return Ok(()); + } - if native_playlist_len > 0 { - // La playlist native a des pistes, on peut simplement appeler play() - // Cela reprendra la lecture à partir du morceau actuel (ou du premier si rien n'est en cours) - renderer.play()?; - self.runtime - .set_playback_source(renderer_id, PlaybackSource::FromQueue); - self.sync_openhome_playlist_for(renderer_id)?; - info!( - renderer = renderer_id.0.as_str(), - playlist_len = native_playlist_len, - "Started OpenHome native playlist playback" - ); - return Ok(()); + // Trouver le track_id courant depuis le snapshot OpenHome + let target_track_id = if let Some(current_id) = snapshot.current_id { + Some(current_id) + } else if let Some(current_idx) = snapshot.current_index { + snapshot.tracks.get(current_idx).map(|track| track.id) + } else { + snapshot.tracks.first().map(|track| track.id) + }; + + if let Some(track_id) = target_track_id { + match renderer.openhome_playlist_play_id(track_id) { + Ok(()) => { + self.runtime + .set_playback_source(renderer_id, PlaybackSource::FromQueue); + self.sync_openhome_playlist_for(renderer_id)?; + info!( + renderer = renderer_id.0.as_str(), + track_id, + playlist_len = snapshot.tracks.len(), + "Started OpenHome playlist playback (current item)" + ); + return Ok(()); + } + Err(err) => { + warn!( + renderer = renderer_id.0.as_str(), + track_id, + error = %err, + "PlayId failed, falling back to Play()" + ); + } + } + } + + // Fallback: appeler Play() sans spécifier de track_id + renderer.play()?; + self.runtime + .set_playback_source(renderer_id, PlaybackSource::FromQueue); + self.sync_openhome_playlist_for(renderer_id)?; + info!( + renderer = renderer_id.0.as_str(), + playlist_len = snapshot.tracks.len(), + "Started OpenHome native playlist playback" + ); + return Ok(()); + } + Err(err) => { + warn!( + renderer = renderer_id.0.as_str(), + error = %err, + "Failed to fetch OpenHome playlist snapshot, playlist might be empty" + ); + debug!( + renderer = renderer_id.0.as_str(), + "OpenHome playlist is empty or unavailable" + ); + self.runtime + .set_playback_source(renderer_id, PlaybackSource::None); + return Ok(()); + } } - - // Fallback: utiliser la queue locale du runtime si la playlist native est vide - // (ce cas se produit quand on a enqueue des items via le control point) - let Some((item, _)) = self.runtime.peek_current(renderer_id) else { - debug!( - renderer = renderer_id.0.as_str(), - "OpenHome playlist is empty or no current item (both native and queue)" - ); - self.runtime - .set_playback_source(renderer_id, PlaybackSource::None); - return Ok(()); - }; - - let track_id = openhome_track_id_from_item(&item) - .ok_or_else(|| anyhow!("Current OpenHome item has no track id"))?; - renderer.openhome_playlist_play_id(track_id)?; - self.runtime - .set_playback_source(renderer_id, PlaybackSource::FromQueue); - self.sync_openhome_playlist_for(renderer_id)?; - info!( - renderer = renderer_id.0.as_str(), - track_id, "Started OpenHome playlist playback (current item from queue)" - ); - Ok(()) } fn play_next_openhome(&self, renderer_id: &RendererId) -> anyhow::Result<()> { let renderer = self.openhome_renderer(renderer_id)?; - let (queue, current_index) = self - .runtime - .queue_full_snapshot(renderer_id) - .ok_or_else(|| Self::runtime_entry_missing(renderer_id))?; - let next_item = match current_index { - Some(idx) => queue.get(idx + 1), - None => queue.first(), + // MICRO-PATCH 5: Récupérer les données directement depuis OpenHome au lieu du miroir local + let snapshot = self.get_cached_openhome_playlist_snapshot( + renderer_id, + OPENHOME_SNAPSHOT_CACHE_TTL, + )?; + + if snapshot.tracks.is_empty() { + debug!( + renderer = renderer_id.0.as_str(), + "OpenHome playlist is empty, cannot play next" + ); + self.runtime + .set_playback_source(renderer_id, PlaybackSource::None); + return Ok(()); + } + + // Déterminer le prochain track_id + let next_track_id = match snapshot.current_index { + Some(idx) => { + // Prendre la piste suivante si elle existe + snapshot + .tracks + .get(idx + 1) + .map(|track| track.id) + .or_else(|| snapshot.tracks.first().map(|track| track.id)) + } + None => snapshot.tracks.first().map(|track| track.id), }; - let Some(item) = next_item else { + let Some(track_id) = next_track_id else { debug!( renderer = renderer_id.0.as_str(), "No OpenHome track available to advance to" @@ -1508,8 +1690,6 @@ impl ControlPoint { return Ok(()); }; - let track_id = openhome_track_id_from_item(item) - .ok_or_else(|| anyhow!("Next OpenHome item has no track id"))?; renderer.openhome_playlist_play_id(track_id)?; self.runtime .set_playback_source(renderer_id, PlaybackSource::FromQueue); @@ -1534,15 +1714,15 @@ fn convert_runtime_position(position: Option<&PlaybackPositionInfo>) -> (Option< } fn build_openhome_queue(info: &RendererInfo) -> Option { - let control_url = info.oh_playlist_control_url.as_ref()?.clone(); - let service_type = info.oh_playlist_service_type.as_ref()?.clone(); - let playlist = OhPlaylistClient::new(control_url, service_type); - let info_client = info - .oh_info_control_url - .as_ref() - .zip(info.oh_info_service_type.as_ref()) - .map(|(url, ty)| OhInfoClient::new(url.clone(), ty.clone())); - Some(OpenHomeQueue::new(info.id.clone(), playlist, info_client)) + let playlist = build_playlist_client(info)?; + let info_client = build_info_client(info); + let product_client = build_product_client(info); + Some(OpenHomeQueue::new( + info.id.clone(), + playlist, + info_client, + product_client, + )) } #[cfg(feature = "pmoserver")] @@ -1586,6 +1766,7 @@ enum PlaylistBackend { struct RendererRuntimeEntry { snapshot: RendererRuntimeSnapshot, pub queue: crate::control_point::music_queue::MusicQueue, + openhome_cache: OpenHomePlaylistCache, playback_source: PlaybackSource, user_stop_requested: bool, playlist_backend: PlaylistBackend, @@ -1596,6 +1777,7 @@ impl Default for RendererRuntimeEntry { Self { snapshot: RendererRuntimeSnapshot::default(), queue: MusicQueue::Internal(InternalQueue::new()), + openhome_cache: OpenHomePlaylistCache::default(), playback_source: PlaybackSource::None, user_stop_requested: false, playlist_backend: PlaylistBackend::PMOQueue, @@ -1603,6 +1785,13 @@ impl Default for RendererRuntimeEntry { } } +#[derive(Debug, Default, Clone)] +struct OpenHomePlaylistCache { + ids: Option>, + snapshot: Option, + last_refresh: Option, +} + struct RuntimeState { entries: Mutex>, } @@ -1667,6 +1856,11 @@ impl OpenHomeQueueProvider for RuntimeOpenHomeQueueProvider { ) -> anyhow::Result> { self.runtime.renderer_state_mut(renderer_id) } + + fn invalidate_openhome_cache(&self, renderer_id: &RendererId) -> anyhow::Result<()> { + self.runtime.invalidate_openhome_cache(renderer_id); + Ok(()) + } } impl RuntimeState { @@ -1846,6 +2040,51 @@ impl RuntimeState { }); } + fn invalidate_openhome_cache(&self, id: &RendererId) { + self.with_entry(id, |entry| { + entry.openhome_cache = OpenHomePlaylistCache::default(); + }); + } + + fn get_openhome_cache(&self, id: &RendererId) -> Option { + let entries = self.entries.lock().unwrap(); + entries.get(id).map(|entry| entry.openhome_cache.clone()) + } + + fn set_openhome_cache(&self, id: &RendererId, cache: OpenHomePlaylistCache) { + self.with_entry(id, |entry| { + entry.openhome_cache = cache; + }); + } + + fn openhome_snapshot_cached( + &self, + renderer: &MusicRenderer, + ttl: Duration, + ) -> anyhow::Result { + let ids = renderer.openhome_playlist_ids()?; + let now = Instant::now(); + + if let Some(cache) = self.get_openhome_cache(renderer.id()) { + if let (Some(cached_ids), Some(snapshot), Some(last_refresh)) = + (cache.ids, cache.snapshot, cache.last_refresh) + { + if cached_ids == ids && now.saturating_duration_since(last_refresh) <= ttl { + return Ok(snapshot); + } + } + } + + let snapshot = renderer.fetch_openhome_playlist_snapshot()?; + let cache = OpenHomePlaylistCache { + ids: Some(ids), + snapshot: Some(snapshot.clone()), + last_refresh: Some(now), + }; + self.set_openhome_cache(renderer.id(), cache); + Ok(snapshot) + } + fn set_playlist_backend(&self, id: &RendererId, backend: PlaylistBackend) { self.with_entry(id, |entry| { entry.playlist_backend = backend; @@ -1969,17 +2208,38 @@ fn refresh_attached_queue_for( // Step 3: Create MusicServer and browse container let music_server = MusicServer::from_info(&server_info, Duration::from_secs(5))?; - let entries = match music_server.browse_children(&container_id, 0, 64) { - Ok(e) => e, - Err(err) => { - warn!( - renderer = renderer_id.0.as_str(), - server = server_id.0.as_str(), - container = container_id.as_str(), - error = %err, - "Failed to browse playlist container for refresh" - ); - return Err(err); + const MAX_BROWSE_ATTEMPTS: usize = 3; + const BROWSE_RETRY_DELAY_MS: u64 = 200; + let mut attempt = 1; + let entries = loop { + match music_server.browse_children(&container_id, 0, 64) { + Ok(e) => break e, + Err(err) => { + if attempt >= MAX_BROWSE_ATTEMPTS { + warn!( + renderer = renderer_id.0.as_str(), + server = server_id.0.as_str(), + container = container_id.as_str(), + attempts = attempt, + error = %err, + "Failed to browse playlist container for refresh" + ); + return Err(err); + } + + debug!( + renderer = renderer_id.0.as_str(), + server = server_id.0.as_str(), + container = container_id.as_str(), + attempt, + error = %err, + "Browse attempt failed, retrying" + ); + thread::sleep(Duration::from_millis( + BROWSE_RETRY_DELAY_MS * attempt as u64, + )); + attempt += 1; + } } }; @@ -1997,6 +2257,7 @@ fn refresh_attached_queue_for( "Refreshed playlist is empty, clearing queue" ); runtime.with_music_queue_mut(renderer_id, |queue| queue.clear_queue())?; + runtime.invalidate_openhome_cache(renderer_id); // Emit QueueUpdated event event_bus.broadcast(RendererEvent::QueueUpdated { @@ -2007,7 +2268,10 @@ fn refresh_attached_queue_for( return Ok(()); } - // Step 5: Intelligent refresh: try to keep current item if it's still in the new list + // Step 5: GENTLE SYNCHRONIZATION + // Use incremental replace_queue() instead of replace_with_attached_playlist() + // This uses LCS algorithm to minimize playlist operations and avoid interrupting playback + // Get the full queue snapshot to access the item currently being played let (full_queue, current_idx) = runtime .queue_full_snapshot(renderer_id) @@ -2016,6 +2280,24 @@ fn refresh_attached_queue_for( // Get the item currently being played (at current_index), not the next one in queue let current_item = current_idx.and_then(|idx| full_queue.get(idx).cloned()); + // Check if renderer is playing - if so, we MUST preserve the current track + let is_playing = { + let renderer_info = { + let reg = registry.read().unwrap(); + reg.get_renderer(renderer_id) + }; + + if let Some(info) = renderer_info { + if let Some(renderer) = MusicRenderer::from_registry_info(info, registry) { + matches!(renderer.playback_state(), Ok(PlaybackState::Playing)) + } else { + false + } + } else { + false + } + }; + let item_found_at = current_item.as_ref().and_then(|current| { let current_uid = current.unique_id(); new_items @@ -2030,6 +2312,7 @@ fn refresh_attached_queue_for( let final_queue_len = runtime.with_music_queue_mut(renderer_id, |queue| { if let Some(idx) = item_found_at { + // Current item is in new playlist - use gentle incremental update queue.replace_queue(new_items.clone(), Some(idx))?; info!( renderer = renderer_id.0.as_str(), @@ -2039,27 +2322,47 @@ fn refresh_attached_queue_for( current_index = idx, upcoming = new_items.len().saturating_sub(idx + 1), current_preserved = true, - "Refreshed queue from playlist container" + is_playing, + "Gentle refresh: current item found in new playlist" ); Ok(new_items.len()) } else if let Some(ref current) = current_item { - let mut combined = Vec::with_capacity(new_items.len() + 1); - combined.push(current.clone()); - combined.extend(new_items.clone()); - queue.replace_queue(combined, Some(0))?; - info!( - renderer = renderer_id.0.as_str(), - server = server_id.0.as_str(), - container = container_id.as_str(), - total_items = new_items.len() + 1, - current_index = 0, - upcoming = new_items.len(), - current_preserved = true, - current_reinserted = true, - "Refreshed queue from playlist container (current item reinserted at start)" - ); - Ok(new_items.len() + 1) + // Current item NOT in new playlist + if is_playing { + // CRITICAL: Preserve current playing track by inserting it at the beginning + let mut combined = Vec::with_capacity(new_items.len() + 1); + combined.push(current.clone()); + combined.extend(new_items.clone()); + queue.replace_queue(combined, Some(0))?; + info!( + renderer = renderer_id.0.as_str(), + server = server_id.0.as_str(), + container = container_id.as_str(), + total_items = new_items.len() + 1, + current_index = 0, + upcoming = new_items.len(), + current_preserved = true, + current_reinserted = true, + is_playing = true, + "Gentle refresh: preserved playing track not in new playlist" + ); + Ok(new_items.len() + 1) + } else { + // Not playing, use new playlist as-is + queue.replace_queue(new_items.clone(), None)?; + info!( + renderer = renderer_id.0.as_str(), + server = server_id.0.as_str(), + container = container_id.as_str(), + total_items = new_items.len(), + current_preserved = false, + is_playing = false, + "Gentle refresh: replaced queue (not playing)" + ); + Ok(new_items.len()) + } } else { + // No current item, use new playlist as-is queue.replace_queue(new_items.clone(), None)?; info!( renderer = renderer_id.0.as_str(), @@ -2067,11 +2370,13 @@ fn refresh_attached_queue_for( container = container_id.as_str(), total_items = new_items.len(), current_preserved = false, - "Refreshed queue from playlist container (no current item)" + is_playing, + "Gentle refresh: no current item" ); Ok(new_items.len()) } })?; + runtime.invalidate_openhome_cache(renderer_id); // Emit QueueUpdated event event_bus.broadcast(RendererEvent::QueueUpdated { @@ -2234,25 +2539,31 @@ fn playback_item_track_metadata(item: &PlaybackItem) -> TrackMetadata { }) } +fn openhome_renderer_from_registry( + registry: &Arc>, + renderer_id: &RendererId, +) -> anyhow::Result { + let info = { + let reg = registry.read().unwrap(); + reg.get_renderer(renderer_id) + .ok_or_else(|| OpenHomeAccessError::RendererNotFound(renderer_id.0.clone()))? + }; + let renderer = MusicRenderer::from_registry_info(info, registry) + .and_then(|r| match r { + MusicRenderer::OpenHome(_) => Some(r), + _ => None, + }) + .ok_or_else(|| OpenHomeAccessError::PlaylistNotSupported(renderer_id.0.clone()))?; + Ok(renderer) +} + fn sync_openhome_playlist( registry: &Arc>, runtime: &Arc, event_bus: &RendererEventBus, renderer_id: &RendererId, ) -> anyhow::Result<()> { - let renderer = { - let info = { - let reg = registry.read().unwrap(); - reg.get_renderer(renderer_id) - .ok_or_else(|| OpenHomeAccessError::RendererNotFound(renderer_id.0.clone()))? - }; - MusicRenderer::from_registry_info(info, registry) - .and_then(|r| match r { - MusicRenderer::OpenHome(_) => Some(r), - _ => None, - }) - .ok_or_else(|| OpenHomeAccessError::PlaylistNotSupported(renderer_id.0.clone()))? - }; + let renderer = openhome_renderer_from_registry(registry, renderer_id)?; let snapshot = renderer.fetch_openhome_playlist_snapshot()?; let playback_items: Vec = snapshot @@ -2261,6 +2572,13 @@ fn sync_openhome_playlist( .map(|track| playback_item_from_openhome_track(renderer_id, track)) .collect(); + debug!( + renderer = renderer_id.0.as_str(), + fetched_items = playback_items.len(), + current_id = ?snapshot.current_id, + "sync_openhome_playlist: fetched snapshot from renderer" + ); + let current_id = snapshot.current_id; let current_index = current_id.and_then(|id| { @@ -2278,6 +2596,12 @@ fn sync_openhome_playlist( }) .unwrap_or(0); + debug!( + renderer = renderer_id.0.as_str(), + queue_len, + "sync_openhome_playlist: updated local queue" + ); + event_bus.broadcast(RendererEvent::QueueUpdated { id: renderer_id.clone(), queue_length: queue_len, @@ -2385,19 +2709,19 @@ impl OpenHomeEventRuntime { } if let Some(url) = info.oh_playlist_event_sub_url.clone() { - let key = OpenHomeSubscriptionKey::new(&info.id, OpenHomeServiceKind::Playlist); + let key = OpenHomeSubscriptionKey::new(&info.id, OhServiceKind::Playlist); active.insert(key.clone()); self.ensure_subscription(key, info.clone(), url); } if let Some(url) = info.oh_info_event_sub_url.clone() { - let key = OpenHomeSubscriptionKey::new(&info.id, OpenHomeServiceKind::Info); + let key = OpenHomeSubscriptionKey::new(&info.id, OhServiceKind::Info); active.insert(key.clone()); self.ensure_subscription(key, info.clone(), url); } if let Some(url) = info.oh_time_event_sub_url.clone() { - let key = OpenHomeSubscriptionKey::new(&info.id, OpenHomeServiceKind::Time); + let key = OpenHomeSubscriptionKey::new(&info.id, OhServiceKind::Time); active.insert(key.clone()); self.ensure_subscription(key, info.clone(), url); } @@ -2508,33 +2832,49 @@ impl OpenHomeEventRuntime { } match entry.service { - OpenHomeServiceKind::Playlist => { + OhServiceKind::Playlist => { if properties .iter() .any(|(name, _)| is_id_array_property(name)) { if self.runtime.uses_openhome_playlist(&entry.renderer.id) { - if let Err(err) = sync_openhome_playlist( - &self.registry, - &self.runtime, - &self.event_bus, - &entry.renderer.id, - ) { - warn!( - renderer = entry.renderer.friendly_name.as_str(), - error = %err, - "Failed to sync OpenHome playlist after IdArray event" - ); + self.runtime + .invalidate_openhome_cache(&entry.renderer.id); + + match openhome_renderer_from_registry(&self.registry, &entry.renderer.id) { + Ok(renderer) => match renderer.openhome_playlist_len() { + Ok(queue_len) => { + self.event_bus.broadcast(RendererEvent::QueueUpdated { + id: entry.renderer.id.clone(), + queue_length: queue_len, + }); + } + Err(err) => { + warn!( + renderer = entry.renderer.friendly_name.as_str(), + error = %err, + "Failed to read OpenHome playlist length after IdArray event" + ); + } + }, + Err(err) => { + warn!( + renderer = entry.renderer.friendly_name.as_str(), + error = %err, + "Failed to build OpenHome renderer after IdArray event" + ); + } } } } } - OpenHomeServiceKind::Info => { + OhServiceKind::Info => { self.handle_info_properties(&entry.renderer.id, properties); } - OpenHomeServiceKind::Time => { + OhServiceKind::Time => { self.handle_time_properties(&entry.renderer.id, properties); } + OhServiceKind::Volume | OhServiceKind::Product => {} } } @@ -2789,11 +3129,11 @@ impl OpenHomeEventRuntime { #[derive(Clone, Debug, PartialEq, Eq, Hash)] struct OpenHomeSubscriptionKey { renderer_id: RendererId, - service: OpenHomeServiceKind, + service: OhServiceKind, } impl OpenHomeSubscriptionKey { - fn new(renderer_id: &RendererId, service: OpenHomeServiceKind) -> Self { + fn new(renderer_id: &RendererId, service: OhServiceKind) -> Self { Self { renderer_id: renderer_id.clone(), service, @@ -2801,26 +3141,9 @@ impl OpenHomeSubscriptionKey { } } -#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)] -enum OpenHomeServiceKind { - Playlist, - Info, - Time, -} - -impl OpenHomeServiceKind { - fn as_str(&self) -> &'static str { - match self { - OpenHomeServiceKind::Playlist => "playlist", - OpenHomeServiceKind::Info => "info", - OpenHomeServiceKind::Time => "time", - } - } -} - struct OpenHomeSubscriptionState { renderer: RendererInfo, - service: OpenHomeServiceKind, + service: OhServiceKind, event_sub_url: String, callback_path: String, sid: Option, @@ -2829,7 +3152,7 @@ struct OpenHomeSubscriptionState { } impl OpenHomeSubscriptionState { - fn new(renderer: RendererInfo, service: OpenHomeServiceKind, event_sub_url: String) -> Self { + fn new(renderer: RendererInfo, service: OhServiceKind, event_sub_url: String) -> Self { Self { callback_path: build_openhome_callback_path(&renderer.id, service), renderer, @@ -2990,7 +3313,7 @@ fn write_openhome_http_response( stream.write_all(response.as_bytes()) } -fn build_openhome_callback_path(id: &RendererId, service: OpenHomeServiceKind) -> String { +fn build_openhome_callback_path(id: &RendererId, service: OhServiceKind) -> String { let mut sanitized = String::new(); for ch in id.0.chars() { if ch.is_ascii_alphanumeric() { @@ -3008,7 +3331,7 @@ fn build_openhome_callback_path(id: &RendererId, service: OpenHomeServiceKind) - fn parse_openhome_propertyset( renderer_id: &RendererId, - service: &OpenHomeServiceKind, + service: &OhServiceKind, body: &[u8], ) -> Vec<(String, String)> { let mut properties = Vec::new(); diff --git a/pmocontrol/src/control_point/music_queue.rs b/pmocontrol/src/control_point/music_queue.rs index db365044..2923b3ea 100644 --- a/pmocontrol/src/control_point/music_queue.rs +++ b/pmocontrol/src/control_point/music_queue.rs @@ -1,4 +1,4 @@ -use anyhow::{Result, anyhow}; +use anyhow::{anyhow, Result}; use crate::control_point::openhome_queue::OpenHomeQueue; use crate::openhome_playlist::OpenHomePlaylistSnapshot; @@ -24,6 +24,17 @@ impl MusicQueue { )), } } + + pub fn replace_with_attached_playlist( + &mut self, + items: Vec, + current_index: Option, + ) -> Result<()> { + match self { + MusicQueue::OpenHome(queue) => queue.replace_entire_playlist(items, current_index), + MusicQueue::Internal(queue) => queue.replace_queue(items, current_index), + } + } } impl Default for MusicQueue { diff --git a/pmocontrol/src/control_point/openhome_queue.rs b/pmocontrol/src/control_point/openhome_queue.rs index 167a12d5..9d19e45d 100644 --- a/pmocontrol/src/control_point/openhome_queue.rs +++ b/pmocontrol/src/control_point/openhome_queue.rs @@ -1,11 +1,13 @@ -use anyhow::{Result, anyhow}; +use anyhow::{anyhow, Result}; use pmodidl::DIDLLite; use quick_xml::escape::escape; +use tracing::debug; use crate::media_server::ServerId; use crate::model::RendererId; use crate::openhome_client::{ - OhInfoClient, OhPlaylistClient, OhTrackEntry, parse_track_metadata_from_didl, + parse_track_metadata_from_didl, OhInfoClient, OhPlaylistClient, OhProductClient, OhTrackEntry, + OPENHOME_PLAYLIST_HEAD_ID, }; use crate::openhome_playlist::{OpenHomePlaylistSnapshot, OpenHomePlaylistTrack}; use crate::queue_backend::{PlaybackItem, QueueBackend, QueueSnapshot}; @@ -16,6 +18,7 @@ pub struct OpenHomeQueue { pub renderer_id: RendererId, pub playlist: OhPlaylistClient, pub info_client: Option, + pub product_client: Option, pub items: Vec, pub current_index: Option, track_ids: Vec, @@ -26,11 +29,13 @@ impl OpenHomeQueue { renderer_id: RendererId, playlist: OhPlaylistClient, info_client: Option, + product_client: Option, ) -> Self { Self { renderer_id, playlist, info_client, + product_client, items: Vec::new(), current_index: None, track_ids: Vec::new(), @@ -43,6 +48,7 @@ impl OpenHomeQueue { /// `OpenHomeRenderer::snapshot_openhome_playlist` but converts entries /// directly into `PlaybackItem`s. pub fn refresh_from_openhome(&mut self) -> Result<()> { + self.ensure_playlist_source_selected()?; let entries = self.playlist.read_all_tracks()?; let mut items = Vec::with_capacity(entries.len()); let mut track_ids = Vec::with_capacity(entries.len()); @@ -56,8 +62,16 @@ impl OpenHomeQueue { .info_client .as_ref() .and_then(|client| client.id().ok()); - let current_index = - current_id.and_then(|id| track_ids.iter().position(|entry_id| *entry_id == id)); + let previous_index = self + .current_index + .and_then(|idx| if idx < track_ids.len() { Some(idx) } else { None }); + let mut current_index = current_id + .and_then(|id| track_ids.iter().position(|entry_id| *entry_id == id)) + .or(previous_index); + + if current_index.is_none() && !track_ids.is_empty() { + current_index = Some(0); + } self.items = items; self.track_ids = track_ids; @@ -85,6 +99,7 @@ impl OpenHomeQueue { current_id: self .current_index .and_then(|idx| self.track_ids.get(idx).copied()), + current_index: self.current_index, tracks, }) } @@ -110,12 +125,14 @@ impl OpenHomeQueue { } }; + self.ensure_playlist_source_selected()?; self.playlist.play_id(id)?; self.current_index = Some(index); Ok(()) } pub fn clear(&mut self) -> Result<()> { + self.ensure_playlist_source_selected()?; self.playlist.delete_all()?; self.items.clear(); self.track_ids.clear(); @@ -123,12 +140,55 @@ impl OpenHomeQueue { Ok(()) } + /// Replace the remote OpenHome playlist entirely with `items`. + /// + /// This is used when attaching a media server playlist: we want to drop any + /// stale entries (even if they were inserted by another control point) and + /// rebuild the renderer playlist from scratch. + pub fn replace_entire_playlist( + &mut self, + items: Vec, + current_index: Option, + ) -> Result<()> { + self.ensure_playlist_source_selected()?; + self.playlist.delete_all()?; + self.items.clear(); + self.track_ids.clear(); + self.current_index = None; + + if items.is_empty() { + return Ok(()); + } + + let mut rebuilt_items = Vec::with_capacity(items.len()); + let mut rebuilt_ids = Vec::with_capacity(items.len()); + let mut previous_id = OPENHOME_PLAYLIST_HEAD_ID; + + for item in items { + let metadata = build_metadata_xml(&item); + let new_id = self.playlist.insert(previous_id, &item.uri, &metadata)?; + previous_id = new_id; + rebuilt_ids.push(new_id); + rebuilt_items.push(self.item_with_openhome_id(item, new_id)); + } + + let normalized = current_index + .filter(|&idx| idx < rebuilt_ids.len()) + .or_else(|| Some(0)); + + self.items = rebuilt_items; + self.track_ids = rebuilt_ids; + self.current_index = normalized; + Ok(()) + } + pub fn add_playback_item( &mut self, item: PlaybackItem, after_id: Option, play: bool, ) -> Result { + self.ensure_playlist_source_selected()?; let metadata_xml = build_metadata_xml(&item); let insert_after = match after_id { Some(id) => id, @@ -161,7 +221,8 @@ impl OpenHomeQueue { } self.track_ids.insert(insert_index, new_id); - self.items.insert(insert_index, item); + let stored_item = self.item_with_openhome_id(item, new_id); + self.items.insert(insert_index, stored_item); self.current_index = if play { Some(insert_index) @@ -173,6 +234,20 @@ impl OpenHomeQueue { Ok(new_id) } + fn ensure_playlist_source_selected(&self) -> Result<()> { + if let Some(product) = &self.product_client { + product.ensure_playlist_source_selected().map_err(|err| { + anyhow!( + "Failed to select OpenHome Playlist source for {}: {}", + self.renderer_id.0, + err + ) + }) + } else { + Ok(()) + } + } + fn playback_item_from_entry(&self, entry: &OhTrackEntry) -> PlaybackItem { let metadata = parse_track_metadata_from_didl(&entry.metadata_xml); let didl_id = didl_id_from_metadata(&entry.metadata_xml) @@ -187,6 +262,12 @@ impl OpenHomeQueue { } } + fn item_with_openhome_id(&self, mut item: PlaybackItem, track_id: u32) -> PlaybackItem { + item.didl_id = format!("openhome:{}", track_id); + item.media_server_id = ServerId(format!("openhome:{}", self.renderer_id.0)); + item + } + fn ensure_track_id(&mut self, index: usize) -> Result { if index >= self.items.len() { return Err(anyhow!("Index out of bounds in OpenHomeQueue: {}", index)); @@ -272,6 +353,41 @@ fn build_metadata_xml(item: &PlaybackItem) -> String { xml } +fn lcs_flags(current: &[PlaybackItem], desired: &[PlaybackItem]) -> (Vec, Vec) { + let m = current.len(); + let n = desired.len(); + let mut dp = vec![vec![0u32; n + 1]; m + 1]; + + for i in 0..m { + for j in 0..n { + if current[i].uri == desired[j].uri { + dp[i + 1][j + 1] = dp[i][j] + 1; + } else { + dp[i + 1][j + 1] = dp[i + 1][j].max(dp[i][j + 1]); + } + } + } + + let mut keep_current = vec![false; m]; + let mut keep_desired = vec![false; n]; + let (mut i, mut j) = (m, n); + + while i > 0 && j > 0 { + if current[i - 1].uri == desired[j - 1].uri { + keep_current[i - 1] = true; + keep_desired[j - 1] = true; + i -= 1; + j -= 1; + } else if dp[i - 1][j] >= dp[i][j - 1] { + i -= 1; + } else { + j -= 1; + } + } + + (keep_current, keep_desired) +} + impl QueueBackend for OpenHomeQueue { fn queue_snapshot(&self) -> Result { Ok(QueueSnapshot { @@ -284,6 +400,7 @@ impl QueueBackend for OpenHomeQueue { let normalized = index.filter(|&i| i < self.items.len()); if let Some(idx) = normalized { let track_id = self.ensure_track_id(idx)?; + self.ensure_playlist_source_selected()?; self.playlist.play_id(track_id)?; } self.current_index = normalized; @@ -295,27 +412,90 @@ impl QueueBackend for OpenHomeQueue { items: Vec, current_index: Option, ) -> Result<()> { - self.playlist.delete_all()?; - - let mut previous_id = 0u32; - let mut inserted_ids = Vec::with_capacity(items.len()); - - for item in &items { - let metadata = build_metadata_xml(item); - let new_id = self.playlist.insert(previous_id, &item.uri, &metadata)?; - previous_id = new_id; - inserted_ids.push(new_id); + self.ensure_playlist_source_selected()?; + if items.is_empty() { + self.playlist.delete_all()?; + self.items.clear(); + self.track_ids.clear(); + self.current_index = None; + return Ok(()); } - let normalized = current_index.filter(|&i| i < inserted_ids.len()); - if let Some(idx) = normalized { - if let Some(track_id) = inserted_ids.get(idx).copied() { - self.playlist.play_id(track_id)?; + // Synchronize local state with the actual OpenHome playlist before computing + // differences. Without this, any drift between our cache and the renderer + // (e.g., manual edits from another control point) would keep the stale items. + self.refresh_from_openhome()?; + debug!( + renderer = self.renderer_id.0.as_str(), + actual_items = self.items.len(), + "OpenHome playlist state refreshed before replace_queue" + ); + + let (keep_current, keep_desired) = lcs_flags(&self.items, &items); + + let items_to_keep = keep_current.iter().filter(|&&k| k).count(); + let items_to_delete = keep_current.iter().filter(|&&k| !k).count(); + let items_to_add = keep_desired.iter().filter(|&&k| !k).count(); + + debug!( + renderer = self.renderer_id.0.as_str(), + keep = items_to_keep, + delete = items_to_delete, + add = items_to_add, + "LCS computed: minimizing OpenHome playlist operations" + ); + + for idx in (0..self.track_ids.len()).rev() { + if !keep_current[idx] { + let track_id = self.track_ids[idx]; + self.playlist.delete_id(track_id)?; + self.track_ids.remove(idx); + self.items.remove(idx); } } - self.items = items; - self.track_ids = inserted_ids; + let remaining_ids = self.track_ids.clone(); + let mut remaining_idx = 0usize; + let mut previous_id = OPENHOME_PLAYLIST_HEAD_ID; + let mut rebuilt_items = Vec::with_capacity(items.len()); + let mut rebuilt_ids = Vec::with_capacity(items.len()); + + for (idx, item) in items.into_iter().enumerate() { + if keep_desired[idx] { + if remaining_idx >= remaining_ids.len() { + return Err(anyhow!( + "OpenHome playlist refresh bookkeeping mismatch (kept entries underflow)" + )); + } + let existing_id = remaining_ids[remaining_idx]; + remaining_idx += 1; + previous_id = existing_id; + rebuilt_ids.push(existing_id); + rebuilt_items.push(self.item_with_openhome_id(item, existing_id)); + } else { + let metadata = build_metadata_xml(&item); + let new_id = self.playlist.insert(previous_id, &item.uri, &metadata)?; + previous_id = new_id; + rebuilt_ids.push(new_id); + rebuilt_items.push(self.item_with_openhome_id(item, new_id)); + } + } + + if remaining_idx != remaining_ids.len() { + return Err(anyhow!( + "OpenHome playlist refresh bookkeeping mismatch (kept entries overflow)" + )); + } + + let previous_index = self + .current_index + .and_then(|idx| if idx < rebuilt_ids.len() { Some(idx) } else { None }); + let normalized = current_index + .filter(|&i| i < rebuilt_ids.len()) + .or(previous_index) + .or_else(|| if rebuilt_ids.is_empty() { None } else { Some(0) }); + self.items = rebuilt_items; + self.track_ids = rebuilt_ids; self.current_index = normalized; Ok(()) } @@ -329,9 +509,10 @@ impl QueueBackend for OpenHomeQueue { return Ok(()); } + self.ensure_playlist_source_selected()?; let track_id = self.ensure_track_id(index)?; let before_id = if index == 0 { - 0 + OPENHOME_PLAYLIST_HEAD_ID } else { self.ensure_track_id(index - 1)? }; @@ -344,7 +525,7 @@ impl QueueBackend for OpenHomeQueue { self.playlist.play_id(new_id)?; } - self.items[index] = item; + self.items[index] = self.item_with_openhome_id(item, new_id); self.track_ids[index] = new_id; Ok(()) } diff --git a/pmocontrol/src/lib.rs b/pmocontrol/src/lib.rs index 36f12a98..37c2c951 100644 --- a/pmocontrol/src/lib.rs +++ b/pmocontrol/src/lib.rs @@ -11,6 +11,7 @@ pub mod linkplay; pub mod media_server; pub mod model; pub mod music_renderer; +pub mod openhome; pub mod openhome_client; pub mod openhome_playlist; pub mod openhome_renderer; diff --git a/pmocontrol/src/model.rs b/pmocontrol/src/model.rs index 8b4e5baa..f9ce9e15 100644 --- a/pmocontrol/src/model.rs +++ b/pmocontrol/src/model.rs @@ -85,6 +85,8 @@ pub struct RendererInfo { pub oh_volume_control_url: Option, pub oh_radio_service_type: Option, pub oh_radio_control_url: Option, + pub oh_product_service_type: Option, + pub oh_product_control_url: Option, } #[derive(Clone, Debug)] diff --git a/pmocontrol/src/music_renderer.rs b/pmocontrol/src/music_renderer.rs index 6c3a4bcb..bbb1f276 100644 --- a/pmocontrol/src/music_renderer.rs +++ b/pmocontrol/src/music_renderer.rs @@ -63,6 +63,7 @@ pub trait OpenHomeQueueProvider: Send + Sync + 'static { &'a self, renderer_id: &RendererId, ) -> Result>; + fn invalidate_openhome_cache(&self, renderer_id: &RendererId) -> Result<()>; } static OPENHOME_QUEUE_PROVIDER: OnceLock> = OnceLock::new(); @@ -187,18 +188,7 @@ impl MusicRenderer { } pub fn openhome_playlist_snapshot(&self) -> Result { - if let Some(provider) = OPENHOME_QUEUE_PROVIDER.get() { - let state = provider.renderer_state(self.id())?; - match &state.queue { - MusicQueue::OpenHome(queue) => queue.openhome_playlist_snapshot(), - _ => Err(op_not_supported( - "openhome_playlist_snapshot", - self.unsupported_backend_name(), - )), - } - } else { - self.fetch_openhome_playlist_snapshot() - } + self.fetch_openhome_playlist_snapshot() } pub(crate) fn fetch_openhome_playlist_snapshot(&self) -> Result { @@ -212,46 +202,41 @@ impl MusicRenderer { } pub fn openhome_playlist_len(&self) -> Result { - if let Some(provider) = OPENHOME_QUEUE_PROVIDER.get() { - let state = provider.renderer_state(self.id())?; - match &state.queue { - MusicQueue::OpenHome(queue) => Ok(queue.len()), - _ => Err(op_not_supported( - "openhome_playlist_len", - self.unsupported_backend_name(), - )), - } - } else { - Ok(self.fetch_openhome_playlist_snapshot()?.tracks.len()) + match self { + MusicRenderer::OpenHome(renderer) => renderer.openhome_playlist_len(), + _ => Err(op_not_supported( + "openhome_playlist_len", + self.unsupported_backend_name(), + )), } } pub fn openhome_playlist_ids(&self) -> Result> { - if let Some(provider) = OPENHOME_QUEUE_PROVIDER.get() { - let state = provider.renderer_state(self.id())?; - match &state.queue { - MusicQueue::OpenHome(queue) => Ok(queue.openhome_track_ids()), - _ => Err(op_not_supported( - "openhome_playlist_ids", - self.unsupported_backend_name(), - )), - } - } else { - let snapshot = self.fetch_openhome_playlist_snapshot()?; - Ok(snapshot.tracks.into_iter().map(|track| track.id).collect()) + match self { + MusicRenderer::OpenHome(renderer) => renderer.openhome_playlist_ids(), + _ => Err(op_not_supported( + "openhome_playlist_ids", + self.unsupported_backend_name(), + )), } } pub fn openhome_playlist_clear(&self) -> Result<()> { if let Some(provider) = OPENHOME_QUEUE_PROVIDER.get() { - let mut state = provider.renderer_state_mut(self.id())?; - match &mut *state.queue { - MusicQueue::OpenHome(queue) => queue.clear(), - _ => Err(op_not_supported( - "openhome_playlist_clear", - self.unsupported_backend_name(), - )), + let result = { + let mut state = provider.renderer_state_mut(self.id())?; + match &mut *state.queue { + MusicQueue::OpenHome(queue) => queue.clear(), + _ => Err(op_not_supported( + "openhome_playlist_clear", + self.unsupported_backend_name(), + )), + } + }; + if result.is_ok() { + provider.invalidate_openhome_cache(self.id())?; } + result } else { self.fetch_openhome_playlist_clear() } @@ -275,17 +260,24 @@ impl MusicRenderer { play: bool, ) -> Result { if let Some(provider) = OPENHOME_QUEUE_PROVIDER.get() { - let mut state = provider.renderer_state_mut(self.id())?; - match &mut *state.queue { - MusicQueue::OpenHome(queue) => { - let playback_item = Self::playback_item_from_params(self.id(), uri, metadata)?; - queue.add_playback_item(playback_item, after_id, play) + let result = { + let mut state = provider.renderer_state_mut(self.id())?; + match &mut *state.queue { + MusicQueue::OpenHome(queue) => { + let playback_item = + Self::playback_item_from_params(self.id(), uri, metadata)?; + queue.add_playback_item(playback_item, after_id, play) + } + _ => Err(op_not_supported( + "openhome_playlist_add_track", + self.unsupported_backend_name(), + )), } - _ => Err(op_not_supported( - "openhome_playlist_add_track", - self.unsupported_backend_name(), - )), + }; + if result.is_ok() { + provider.invalidate_openhome_cache(self.id())?; } + result } else { self.fetch_openhome_playlist_add_track(uri, metadata, after_id, play) } @@ -293,14 +285,20 @@ impl MusicRenderer { pub fn openhome_playlist_play_id(&self, id: u32) -> Result<()> { if let Some(provider) = OPENHOME_QUEUE_PROVIDER.get() { - let mut state = provider.renderer_state_mut(self.id())?; - match &mut *state.queue { - MusicQueue::OpenHome(queue) => queue.select_track_id(id), - _ => Err(op_not_supported( - "openhome_playlist_play_id", - self.unsupported_backend_name(), - )), + let result = { + let mut state = provider.renderer_state_mut(self.id())?; + match &mut *state.queue { + MusicQueue::OpenHome(queue) => queue.select_track_id(id), + _ => Err(op_not_supported( + "openhome_playlist_play_id", + self.unsupported_backend_name(), + )), + } + }; + if result.is_ok() { + provider.invalidate_openhome_cache(self.id())?; } + result } else { self.fetch_openhome_playlist_play_id(id) } diff --git a/pmocontrol/src/openhome.rs b/pmocontrol/src/openhome.rs new file mode 100644 index 00000000..81d4f6fb --- /dev/null +++ b/pmocontrol/src/openhome.rs @@ -0,0 +1,187 @@ +use crate::model::RendererInfo; +use crate::openhome_client::{ + OhInfoClient, OhPlaylistClient, OhProductClient, OhRadioClient, OhTimeClient, OhVolumeClient, +}; + +#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)] +pub enum OhServiceKind { + Playlist, + Info, + Time, + Volume, + Product, +} + +impl OhServiceKind { + pub fn as_str(&self) -> &'static str { + match self { + OhServiceKind::Playlist => "playlist", + OhServiceKind::Info => "info", + OhServiceKind::Time => "time", + OhServiceKind::Volume => "volume", + OhServiceKind::Product => "product", + } + } +} + +pub struct OhServiceEndpoint<'a> { + pub control_url: &'a str, + pub service_type: &'a str, +} + +fn endpoint_for<'a>(info: &'a RendererInfo, kind: OhServiceKind) -> Option> { + let (control_url, service_type) = match kind { + OhServiceKind::Playlist => ( + info.oh_playlist_control_url.as_deref()?, + info.oh_playlist_service_type.as_deref()?, + ), + OhServiceKind::Info => ( + info.oh_info_control_url.as_deref()?, + info.oh_info_service_type.as_deref()?, + ), + OhServiceKind::Time => ( + info.oh_time_control_url.as_deref()?, + info.oh_time_service_type.as_deref()?, + ), + OhServiceKind::Volume => ( + info.oh_volume_control_url.as_deref()?, + info.oh_volume_service_type.as_deref()?, + ), + OhServiceKind::Product => ( + info.oh_product_control_url.as_deref()?, + info.oh_product_service_type.as_deref()?, + ), + }; + + Some(OhServiceEndpoint { + control_url, + service_type, + }) +} + +pub fn control_url_for<'a>(info: &'a RendererInfo, kind: OhServiceKind) -> Option<&'a str> { + endpoint_for(info, kind).map(|endpoint| endpoint.control_url) +} + +pub fn service_type_for<'a>(info: &'a RendererInfo, kind: OhServiceKind) -> Option<&'a str> { + endpoint_for(info, kind).map(|endpoint| endpoint.service_type) +} + +pub fn build_playlist_client(info: &RendererInfo) -> Option { + let endpoint = endpoint_for(info, OhServiceKind::Playlist)?; + Some(OhPlaylistClient::new( + endpoint.control_url.to_string(), + endpoint.service_type.to_string(), + )) +} + +pub fn build_info_client(info: &RendererInfo) -> Option { + let endpoint = endpoint_for(info, OhServiceKind::Info)?; + Some(OhInfoClient::new( + endpoint.control_url.to_string(), + endpoint.service_type.to_string(), + )) +} + +pub fn build_time_client(info: &RendererInfo) -> Option { + let endpoint = endpoint_for(info, OhServiceKind::Time)?; + Some(OhTimeClient::new( + endpoint.control_url.to_string(), + endpoint.service_type.to_string(), + )) +} + +pub fn build_volume_client(info: &RendererInfo) -> Option { + let endpoint = endpoint_for(info, OhServiceKind::Volume)?; + Some(OhVolumeClient::new( + endpoint.control_url.to_string(), + endpoint.service_type.to_string(), + )) +} + +pub fn build_product_client(info: &RendererInfo) -> Option { + let endpoint = endpoint_for(info, OhServiceKind::Product)?; + Some(OhProductClient::new( + endpoint.control_url.to_string(), + endpoint.service_type.to_string(), + )) +} + +pub fn build_radio_client(info: &RendererInfo) -> Option { + let control_url = info.oh_radio_control_url.as_ref()?; + let service_type = info.oh_radio_service_type.as_ref()?; + Some(OhRadioClient::new( + control_url.clone(), + service_type.clone(), + )) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::model::{RendererCapabilities, RendererId, RendererInfo, RendererProtocol}; + + fn sample_renderer_info() -> RendererInfo { + RendererInfo { + id: RendererId("renderer".into()), + udn: "renderer".into(), + friendly_name: "Renderer".into(), + model_name: "Model".into(), + manufacturer: "Maker".into(), + protocol: RendererProtocol::OpenHomeOnly, + capabilities: RendererCapabilities::default(), + location: "http://host:1234/description.xml".into(), + server_header: "test".into(), + online: true, + last_seen: std::time::SystemTime::now(), + max_age: 1800, + avtransport_service_type: None, + avtransport_control_url: None, + rendering_control_service_type: None, + rendering_control_control_url: None, + connection_manager_service_type: None, + connection_manager_control_url: None, + oh_playlist_service_type: Some("urn:av-openhome-org:service:Playlist:1".into()), + oh_playlist_control_url: Some("http://host/oh/playlist".into()), + oh_playlist_event_sub_url: Some("http://host/events/playlist".into()), + oh_info_service_type: Some("urn:av-openhome-org:service:Info:1".into()), + oh_info_control_url: Some("http://host/oh/info".into()), + oh_info_event_sub_url: Some("http://host/events/info".into()), + oh_time_service_type: Some("urn:av-openhome-org:service:Time:1".into()), + oh_time_control_url: Some("http://host/oh/time".into()), + oh_time_event_sub_url: Some("http://host/events/time".into()), + oh_volume_service_type: Some("urn:av-openhome-org:service:Volume:1".into()), + oh_volume_control_url: Some("http://host/oh/volume".into()), + oh_radio_service_type: None, + oh_radio_control_url: None, + oh_product_service_type: Some("urn:av-openhome-org:service:Product:1".into()), + oh_product_control_url: Some("http://host/oh/product".into()), + } + } + + #[test] + fn selects_correct_playlist_endpoint() { + let info = sample_renderer_info(); + let endpoint = endpoint_for(&info, OhServiceKind::Playlist).unwrap(); + assert_eq!(endpoint.control_url, "http://host/oh/playlist"); + assert_eq!( + endpoint.service_type, + "urn:av-openhome-org:service:Playlist:1" + ); + } + + #[test] + fn returns_none_when_service_missing() { + let mut info = sample_renderer_info(); + info.oh_playlist_control_url = None; + assert!(endpoint_for(&info, OhServiceKind::Playlist).is_none()); + } + + #[test] + fn info_and_playlist_use_different_urls() { + let info = sample_renderer_info(); + let playlist = control_url_for(&info, OhServiceKind::Playlist).unwrap(); + let info_url = control_url_for(&info, OhServiceKind::Info).unwrap(); + assert_ne!(playlist, info_url); + } +} diff --git a/pmocontrol/src/openhome_client.rs b/pmocontrol/src/openhome_client.rs index e9bb8fee..324a6f2e 100644 --- a/pmocontrol/src/openhome_client.rs +++ b/pmocontrol/src/openhome_client.rs @@ -1,9 +1,15 @@ use crate::model::TrackMetadata; -use crate::soap_client::{SoapCallResult, invoke_upnp_action}; -use anyhow::{Result, anyhow}; +use crate::soap_client::{invoke_upnp_action, parse_upnp_error, SoapCallResult}; +use anyhow::{anyhow, Result}; use pmoupnp::soap::SoapEnvelope; +use tracing::{debug, info, warn}; use xmltree::{Element, XMLNode}; +/// Value used by OpenHome renderers to indicate "insert at the head". +/// Several implementations, including upmpdcli, expect zero rather than the +/// historical 0xFFFFFFFF sentinel. +pub const OPENHOME_PLAYLIST_HEAD_ID: u32 = 0; + #[derive(Debug, Clone)] pub struct OhTrackEntry { pub id: u32, @@ -38,6 +44,13 @@ pub struct OhRadioChannel { pub metadata_xml: Option, } +#[derive(Debug, Clone)] +pub struct OhProductSource { + pub name: String, + pub source_type: String, + pub visible: bool, +} + #[derive(Debug, Clone)] pub struct OhPlaylistClient { pub control_url: String, @@ -62,7 +75,7 @@ impl OhPlaylistClient { .map(|id| id.to_string()) .collect::>() .join(","); - let args = [("aIdList", id_list_csv.as_str())]; + let args = [("IdList", id_list_csv.as_str())]; let call_result = invoke_upnp_action(&self.control_url, &self.service_type, "ReadList", &args)?; @@ -71,16 +84,24 @@ impl OhPlaylistClient { let response = find_child_with_suffix(&envelope.body.content, "ReadListResponse") .ok_or_else(|| anyhow!("Missing ReadListResponse element in SOAP body"))?; - let track_list_xml = extract_child_text(response, "aTrackList")?; - parse_track_list(&track_list_xml) + let track_list_b64 = + extract_child_text_any(response, &["aTrackList", "TrackList", "aValue", "Value"])?; + let track_list_sample: String = track_list_b64.chars().take(256).collect(); + debug!( + control_url = self.control_url.as_str(), + track_list_len = track_list_b64.len(), + track_list_sample = %track_list_sample, + "OpenHome ReadList returned raw TrackList content" + ); + parse_track_list(&track_list_b64) } pub fn insert(&self, after_id: u32, uri: &str, metadata: &str) -> Result { let after_id_str = after_id.to_string(); let args = [ - ("aAfterId", after_id_str.as_str()), - ("aUri", uri), - ("aMetadata", metadata), + ("AfterId", after_id_str.as_str()), + ("Uri", uri), + ("Metadata", metadata), ]; let call_result = @@ -89,7 +110,8 @@ impl OhPlaylistClient { let envelope = ensure_success("Insert", &call_result)?; let response = find_child_with_suffix(&envelope.body.content, "InsertResponse") .ok_or_else(|| anyhow!("Missing InsertResponse element in SOAP body"))?; - let new_id_text = extract_child_text(response, "aNewId")?; + let new_id_text = + extract_child_text_any(response, &["aNewId", "NewId", "aValue", "Value"])?; let new_id = new_id_text .parse::() .map_err(|_| anyhow!("Invalid aNewId value: {}", new_id_text))?; @@ -99,7 +121,7 @@ impl OhPlaylistClient { pub fn play_id(&self, id: u32) -> Result<()> { let id_str = id.to_string(); - let args = [("aId", id_str.as_str())]; + let args = [("Id", id_str.as_str())]; let call_result = invoke_upnp_action(&self.control_url, &self.service_type, "PlayId", &args)?; @@ -135,7 +157,7 @@ impl OhPlaylistClient { pub fn seek_second_absolute(&self, second: u32) -> Result<()> { let second_str = second.to_string(); - let args = [("aSecond", second_str.as_str())]; + let args = [("Second", second_str.as_str())]; let call_result = invoke_upnp_action( &self.control_url, &self.service_type, @@ -148,7 +170,7 @@ impl OhPlaylistClient { pub fn delete_id(&self, id: u32) -> Result<()> { let id_str = id.to_string(); - let args = [("aId", id_str.as_str())]; + let args = [("Id", id_str.as_str())]; let call_result = invoke_upnp_action(&self.control_url, &self.service_type, "DeleteId", &args)?; @@ -169,7 +191,7 @@ impl OhPlaylistClient { let envelope = ensure_success("TracksMax", &call_result)?; let response = find_child_with_suffix(&envelope.body.content, "TracksMaxResponse") .ok_or_else(|| anyhow!("Missing TracksMaxResponse element in SOAP body"))?; - let value_text = extract_child_text(response, "aValue")?; + let value_text = extract_child_text_any(response, &["aValue", "Value"])?; let value = value_text .parse::() .map_err(|_| anyhow!("Invalid TracksMax value: {}", value_text))?; @@ -185,7 +207,10 @@ impl OhPlaylistClient { .ok_or_else(|| anyhow!("Missing IdArrayResponse element in SOAP body"))?; // Try to extract the array element. If missing, assume empty playlist. - let array_text = match extract_child_text_any(response, &["aArray", "aIdArray"]) { + let array_text = match extract_child_text_any( + response, + &["aArray", "Array", "aIdArray", "IdArray", "aValue", "Value"], + ) { Ok(text) => text, Err(_) => { // Element not found - playlist is likely empty @@ -215,16 +240,49 @@ impl OhPlaylistClient { pub fn read_all_tracks(&self) -> Result> { let ids = self.id_array()?; + debug!( + control_url = self.control_url.as_str(), + id_count = ids.len(), + "OpenHome Playlist IdArray returned" + ); + if ids.is_empty() { + info!( + control_url = self.control_url.as_str(), + "OpenHome Playlist is empty (no track IDs)" + ); return Ok(Vec::new()); } const MAX_BATCH: usize = 64; let mut entries = Vec::with_capacity(ids.len()); for chunk in ids.chunks(MAX_BATCH) { - let mut batch = self.read_list(chunk)?; - entries.append(&mut batch); + match self.read_list(chunk) { + Ok(mut batch) => entries.append(&mut batch), + Err(err) if chunk.len() > 1 && is_invalid_entry_id_error(&err) => { + debug!( + control_url = self.control_url.as_str(), + requested = chunk.len(), + "ReadList chunk failed with invalid entry ids, falling back to per-id requests" + ); + for id in chunk { + match self.read_list(&[*id]) { + Ok(mut single) => entries.append(&mut single), + Err(inner_err) => return Err(inner_err), + } + } + } + Err(err) => return Err(err), + } } + + info!( + control_url = self.control_url.as_str(), + track_count = entries.len(), + expected_count = ids.len(), + "OpenHome Playlist tracks read" + ); + Ok(entries) } } @@ -244,15 +302,14 @@ impl OhInfoClient { } pub fn track(&self) -> Result { - use tracing::debug; - let call_result = invoke_upnp_action(&self.control_url, &self.service_type, "Track", &[])?; let envelope = ensure_success("Track", &call_result)?; let response = find_child_with_suffix(&envelope.body.content, "TrackResponse") .ok_or_else(|| anyhow!("Missing TrackResponse element in SOAP body"))?; - let uri = extract_child_text(response, "aUri")?; + let uri = extract_child_text_any(response, &["aUri", "Uri", "aValue", "Value"]) + .unwrap_or_default(); let metadata_xml = extract_child_text_optional(response, "aMetadata") .unwrap_or(None) .filter(|s| !s.is_empty()); @@ -278,7 +335,8 @@ impl OhInfoClient { let response = find_child_with_suffix(&envelope.body.content, "NextResponse") .ok_or_else(|| anyhow!("Missing NextResponse element in SOAP body"))?; - let uri = extract_child_text(response, "aUri")?; + let uri = extract_child_text_any(response, &["aUri", "Uri", "aValue", "Value"]) + .unwrap_or_default(); let metadata_xml = extract_child_text_optional(response, "aMetadata") .unwrap_or(None) .filter(|s| !s.is_empty()); @@ -306,7 +364,7 @@ impl OhInfoClient { let envelope = ensure_success("TransportState", &call_result)?; let response = find_child_with_suffix(&envelope.body.content, "TransportStateResponse") .ok_or_else(|| anyhow!("Missing TransportStateResponse element in SOAP body"))?; - let state = extract_child_text(response, "aState")?; + let state = extract_child_text_any(response, &["aState", "State", "aValue", "Value"])?; Ok(state) } @@ -337,15 +395,18 @@ impl OhTimeClient { let response = find_child_with_suffix(&envelope.body.content, "TimeResponse") .ok_or_else(|| anyhow!("Missing TimeResponse element in SOAP body"))?; - let track_count = extract_child_text(response, "aTrackCount")? - .parse::() - .map_err(|_| anyhow!("Invalid aTrackCount value in Time response"))?; - let duration_secs = extract_child_text(response, "aDuration")? - .parse::() - .map_err(|_| anyhow!("Invalid aDuration value in Time response"))?; - let elapsed_secs = extract_child_text(response, "aSeconds")? - .parse::() - .map_err(|_| anyhow!("Invalid aSeconds value in Time response"))?; + let track_count = + extract_child_text_any(response, &["aTrackCount", "TrackCount", "aValue", "Value"])? + .parse::() + .map_err(|_| anyhow!("Invalid aTrackCount value in Time response"))?; + let duration_secs = + extract_child_text_any(response, &["aDuration", "Duration", "aValue", "Value"])? + .parse::() + .map_err(|_| anyhow!("Invalid aDuration value in Time response"))?; + let elapsed_secs = + extract_child_text_any(response, &["aSeconds", "Seconds", "aValue", "Value"])? + .parse::() + .map_err(|_| anyhow!("Invalid aSeconds value in Time response"))?; Ok(OhTimePosition { track_count, @@ -374,7 +435,7 @@ impl OhVolumeClient { let envelope = ensure_success("Volume", &call_result)?; let response = find_child_with_suffix(&envelope.body.content, "VolumeResponse") .ok_or_else(|| anyhow!("Missing VolumeResponse element in SOAP body"))?; - let value = extract_child_text(response, "aVolume")?; + let value = extract_child_text_any(response, &["aVolume", "Volume", "aValue", "Value"])?; let parsed = value .parse::() .map_err(|_| anyhow!("Invalid volume value: {}", value))?; @@ -383,7 +444,7 @@ impl OhVolumeClient { pub fn set_volume(&self, vol: u16) -> Result<()> { let vol_str = vol.to_string(); - let args = [("aVolume", vol_str.as_str())]; + let args = [("Value", vol_str.as_str())]; let call_result = invoke_upnp_action(&self.control_url, &self.service_type, "SetVolume", &args)?; handle_action_response("SetVolume", &call_result) @@ -394,13 +455,13 @@ impl OhVolumeClient { let envelope = ensure_success("Mute", &call_result)?; let response = find_child_with_suffix(&envelope.body.content, "MuteResponse") .ok_or_else(|| anyhow!("Missing MuteResponse element in SOAP body"))?; - let value = extract_child_text(response, "aMute")?; + let value = extract_child_text_any(response, &["aMute", "Mute", "aValue", "Value"])?; parse_bool(&value) } pub fn set_mute(&self, mute: bool) -> Result<()> { let mute_str = if mute { "1" } else { "0" }; - let args = [("aMute", mute_str)]; + let args = [("Mute", mute_str)]; let call_result = invoke_upnp_action(&self.control_url, &self.service_type, "SetMute", &args)?; handle_action_response("SetMute", &call_result) @@ -423,7 +484,7 @@ impl OhRadioClient { pub fn play_channel(&self, id: u32) -> Result<()> { let id_str = id.to_string(); - let args = [("aId", id_str.as_str())]; + let args = [("Id", id_str.as_str())]; let call_result = invoke_upnp_action(&self.control_url, &self.service_type, "PlayChannel", &args)?; handle_action_response("PlayChannel", &call_result) @@ -431,7 +492,7 @@ impl OhRadioClient { pub fn channel(&self, id: u32) -> Result { let id_str = id.to_string(); - let args = [("aId", id_str.as_str())]; + let args = [("Id", id_str.as_str())]; let call_result = invoke_upnp_action(&self.control_url, &self.service_type, "Channel", &args)?; @@ -439,7 +500,8 @@ impl OhRadioClient { let response = find_child_with_suffix(&envelope.body.content, "ChannelResponse") .ok_or_else(|| anyhow!("Missing ChannelResponse element in SOAP body"))?; - let uri = extract_child_text(response, "aUri")?; + let uri = extract_child_text_any(response, &["aUri", "Uri", "aValue", "Value"]) + .unwrap_or_default(); let metadata_xml = extract_child_text_optional(response, "aMetadata") .unwrap_or(None) .filter(|s| !s.is_empty()); @@ -448,9 +510,130 @@ impl OhRadioClient { } } -pub fn parse_track_metadata_from_didl(xml: &str) -> Option { - use tracing::debug; +#[derive(Debug, Clone)] +pub struct OhProductClient { + pub control_url: String, + pub service_type: String, +} +impl OhProductClient { + pub fn new(control_url: String, service_type: String) -> Self { + Self { + control_url, + service_type, + } + } + + pub fn source_xml(&self) -> Result> { + let call_result = + invoke_upnp_action(&self.control_url, &self.service_type, "SourceXml", &[])?; + let envelope = ensure_success("SourceXml", &call_result)?; + let response = find_child_with_suffix(&envelope.body.content, "SourceXmlResponse") + .ok_or_else(|| anyhow!("Missing SourceXmlResponse element in SOAP body"))?; + let xml = extract_child_text_any(response, &["aSourceXml", "aXml", "aValue", "Value"])?; + parse_product_source_list(&xml) + } + + pub fn source_index(&self) -> Result { + let call_result = + invoke_upnp_action(&self.control_url, &self.service_type, "SourceIndex", &[])?; + let envelope = ensure_success("SourceIndex", &call_result)?; + let response = find_child_with_suffix(&envelope.body.content, "SourceIndexResponse") + .ok_or_else(|| anyhow!("Missing SourceIndexResponse element in SOAP body"))?; + let value = extract_child_text_any(response, &["aIndex", "Index", "aValue", "Value"])?; + value + .parse::() + .map_err(|_| anyhow!("Invalid Product.SourceIndex value: {}", value)) + } + + pub fn set_source_index(&self, index: u32) -> Result<()> { + let value = index.to_string(); + let args = [("Index", value.as_str())]; + let call_result = invoke_upnp_action( + &self.control_url, + &self.service_type, + "SetSourceIndex", + &args, + )?; + handle_action_response("SetSourceIndex", &call_result) + } + + pub fn ensure_playlist_source_selected(&self) -> Result<()> { + let sources = self.source_xml()?; + + // Log all available sources for diagnostics + info!( + control_url = self.control_url.as_str(), + source_count = sources.len(), + "OpenHome Product sources available" + ); + for (idx, source) in sources.iter().enumerate() { + debug!( + control_url = self.control_url.as_str(), + index = idx, + name = source.name.as_str(), + source_type = source.source_type.as_str(), + visible = source.visible, + "OpenHome source" + ); + } + + let playlist_index = sources + .iter() + .position(|source| source.source_type.eq_ignore_ascii_case("playlist")) + .ok_or_else(|| { + warn!( + control_url = self.control_url.as_str(), + available_types = ?sources.iter().map(|s| s.source_type.as_str()).collect::>(), + "OpenHome Product source list does not expose a Playlist entry" + ); + anyhow!("OpenHome Product source list does not expose a Playlist entry") + })?; + let playlist_index = playlist_index as u32; + let current_index = self.source_index()?; + + // Log current source state + let current_source = sources.get(current_index as usize); + info!( + control_url = self.control_url.as_str(), + current_index, + current_source_name = current_source.map(|s| s.name.as_str()).unwrap_or("unknown"), + current_source_type = current_source.map(|s| s.source_type.as_str()).unwrap_or("unknown"), + playlist_index, + needs_switch = current_index != playlist_index, + "OpenHome source state" + ); + + if current_index != playlist_index { + info!( + control_url = self.control_url.as_str(), + from_index = current_index, + to_index = playlist_index, + "Switching OpenHome Product source to Playlist" + ); + self.set_source_index(playlist_index)?; + + // Verify the switch was successful + let new_index = self.source_index()?; + if new_index == playlist_index { + info!( + control_url = self.control_url.as_str(), + "Successfully switched to Playlist source" + ); + } else { + warn!( + control_url = self.control_url.as_str(), + expected = playlist_index, + actual = new_index, + "Source switch may have failed - index mismatch" + ); + } + } + Ok(()) + } +} + +pub fn parse_track_metadata_from_didl(xml: &str) -> Option { if xml.trim().is_empty() { return None; } @@ -477,11 +660,28 @@ pub fn parse_track_metadata_from_didl(xml: &str) -> Option { }) } -fn parse_track_list(xml: &str) -> Result> { - if xml.trim().is_empty() { +fn parse_track_list(payload: &str) -> Result> { + let trimmed = payload.trim(); + if trimmed.is_empty() { return Ok(Vec::new()); } + let (xml, was_base64) = if trimmed.starts_with('<') { + (trimmed.to_string(), false) + } else { + let bytes = decode_base64(trimmed)?; + let decoded = String::from_utf8(bytes) + .map_err(|err| anyhow!("TrackList payload not valid UTF-8 after base64 decode: {err}"))?; + (decoded, true) + }; + let xml_sample: String = xml.chars().take(256).collect(); + debug!( + raw_base64 = was_base64, + decoded_len = xml.len(), + decoded_sample = %xml_sample, + "Decoded OpenHome TrackList payload" + ); + let mut reader = std::io::Cursor::new(xml.as_bytes()); let root = Element::parse(&mut reader) .map_err(|err| anyhow!("Failed to parse OpenHome TrackList XML: {}", err))?; @@ -499,12 +699,19 @@ fn parse_track_list(xml: &str) -> Result> { } fn parse_track_entry(elem: &Element) -> Result { - let id_text = extract_child_text(elem, "Id")?; + let id_text = extract_child_text_local(elem, "Id")?; + if id_text.contains(',') { + debug!( + raw_entry = %elem.name, + raw_id = id_text.as_str(), + "Unexpected multi-value Id element in OpenHome TrackList entry" + ); + } let id = id_text .parse::() .map_err(|_| anyhow!("Invalid OpenHome Entry Id: {}", id_text))?; - let uri = extract_child_text(elem, "Uri")?; - let metadata_xml = extract_child_text_optional(elem, "Metadata")?.unwrap_or_default(); + let uri = extract_child_text_local(elem, "Uri")?; + let metadata_xml = extract_child_text_optional_local(elem, "Metadata")?.unwrap_or_default(); Ok(OhTrackEntry { id, @@ -513,6 +720,48 @@ fn parse_track_entry(elem: &Element) -> Result { }) } +fn parse_product_source_list(xml: &str) -> Result> { + if xml.trim().is_empty() { + return Ok(Vec::new()); + } + + let mut reader = std::io::Cursor::new(xml.as_bytes()); + let root = Element::parse(&mut reader) + .map_err(|err| anyhow!("Failed to parse OpenHome SourceXml payload: {}", err))?; + + let mut sources = Vec::new(); + for node in &root.children { + if let XMLNode::Element(elem) = node { + if elem.name.ends_with("Source") { + let name = extract_child_text(elem, "Name")?; + let source_type = extract_child_text(elem, "Type")?; + let visible = extract_child_text_optional(elem, "Visible")? + .map(|v| parse_visible_flag(&v)) + .unwrap_or(true); + + sources.push(OhProductSource { + name, + source_type, + visible, + }); + } + } + } + + Ok(sources) +} + +fn parse_visible_flag(value: &str) -> bool { + let trimmed = value.trim(); + if trimmed.eq_ignore_ascii_case("true") { + return true; + } + if trimmed.eq_ignore_ascii_case("false") { + return false; + } + trimmed == "1" +} + fn ensure_success<'a>(action: &str, call_result: &'a SoapCallResult) -> Result<&'a SoapEnvelope> { if !call_result.status.is_success() { if let Some(env) = &call_result.envelope { @@ -555,42 +804,6 @@ fn handle_action_response(action: &str, call_result: &SoapCallResult) -> Result< Ok(()) } -#[derive(Debug, Clone)] -struct UpnpError { - pub error_code: u32, - pub error_description: String, -} - -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 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), @@ -598,6 +811,17 @@ fn find_child_with_suffix<'a>(parent: &'a Element, suffix: &str) -> Option<&'a E }) } +fn find_child_with_local_name<'a>(parent: &'a Element, local: &str) -> Option<&'a Element> { + parent.children.iter().find_map(|node| { + if let XMLNode::Element(elem) = node { + if elem.name == local || elem.name.ends_with(&format!(":{}", local)) { + return Some(elem); + } + } + None + }) +} + fn extract_child_text(parent: &Element, suffix: &str) -> Result { let child = find_child_with_suffix(parent, suffix) .ok_or_else(|| anyhow!("Missing {suffix} element in response"))?; @@ -634,6 +858,28 @@ fn extract_child_text_any(parent: &Element, suffixes: &[&str]) -> Result )) } +fn extract_child_text_local(parent: &Element, local: &str) -> Result { + let child = find_child_with_local_name(parent, local) + .ok_or_else(|| anyhow!("Missing {local} element in response"))?; + let text = child + .get_text() + .map(|t| t.trim().to_string()) + .ok_or_else(|| anyhow!("{local} element missing text in response"))?; + Ok(text) +} + +fn extract_child_text_optional_local(parent: &Element, local: &str) -> Result> { + if let Some(child) = find_child_with_local_name(parent, local) { + let text = child + .get_text() + .map(|t| t.trim().to_string()) + .unwrap_or_default(); + Ok(Some(text)) + } else { + Ok(None) + } +} + fn parse_bool(value: &str) -> Result { match value.trim() { "0" => Ok(false), @@ -678,3 +924,48 @@ pub(crate) fn decode_base64(input: &str) -> Result> { Ok(output) } + +fn is_invalid_entry_id_error(err: &anyhow::Error) -> bool { + let msg = format!("{err}"); + msg.contains("Invalid OpenHome Entry Id") +} + +#[cfg(test)] +mod tests { + use super::*; + use std::io::Cursor; + + #[test] + fn parse_insert_response_accepts_newid_without_prefix() { + let xml = r#"1337"#; + let mut cursor = Cursor::new(xml.as_bytes()); + let response = Element::parse(&mut cursor).expect("valid xml"); + let value = + extract_child_text_any(&response, &["aNewId", "NewId", "aValue", "Value"]).unwrap(); + assert_eq!(value, "1337"); + } + + #[test] + fn parse_readlist_response_accepts_tracklist_without_prefix() { + let xml = r#"PGVudHJ5PjwvZW50cnk+"#; + let mut cursor = Cursor::new(xml.as_bytes()); + let response = Element::parse(&mut cursor).expect("valid xml"); + let value = + extract_child_text_any(&response, &["aTrackList", "TrackList", "aValue", "Value"]) + .expect("tracklist"); + assert_eq!(value, "PGVudHJ5PjwvZW50cnk+"); + } + + #[test] + fn parse_idarray_response_accepts_array_without_prefix() { + let xml = r#"1AAAAAQAAAAI="#; + let mut cursor = Cursor::new(xml.as_bytes()); + let response = Element::parse(&mut cursor).expect("valid xml"); + let value = extract_child_text_any( + &response, + &["aArray", "Array", "aIdArray", "IdArray", "aValue", "Value"], + ) + .expect("array content"); + assert_eq!(value, "AAAAAQAAAAI="); + } +} diff --git a/pmocontrol/src/openhome_playlist.rs b/pmocontrol/src/openhome_playlist.rs index 7ff368da..b9f85f7a 100644 --- a/pmocontrol/src/openhome_playlist.rs +++ b/pmocontrol/src/openhome_playlist.rs @@ -6,6 +6,8 @@ pub struct OpenHomePlaylistSnapshot { pub renderer_id: String, /// ID courant dans la playlist (si connu). pub current_id: Option, + /// Position courante dans la playlist (si connue). + pub current_index: Option, /// Tracks présents dans la playlist native. pub tracks: Vec, } diff --git a/pmocontrol/src/openhome_renderer.rs b/pmocontrol/src/openhome_renderer.rs index e92195d8..3f325018 100644 --- a/pmocontrol/src/openhome_renderer.rs +++ b/pmocontrol/src/openhome_renderer.rs @@ -4,12 +4,16 @@ use crate::capabilities::{ }; use crate::model::{RendererId, RendererInfo, RendererProtocol}; use crate::music_renderer::op_not_supported; +use crate::openhome::{ + build_info_client, build_playlist_client, build_product_client, build_radio_client, + build_time_client, build_volume_client, +}; use crate::openhome_client::{ - OhInfoClient, OhPlaylistClient, OhRadioClient, OhTimeClient, OhTrackEntry, OhVolumeClient, - parse_track_metadata_from_didl, + parse_track_metadata_from_didl, OhInfoClient, OhPlaylistClient, OhProductClient, OhRadioClient, + OhTimeClient, OhTrackEntry, OhVolumeClient, OPENHOME_PLAYLIST_HEAD_ID, }; use crate::openhome_playlist::{OpenHomePlaylistSnapshot, OpenHomePlaylistTrack}; -use anyhow::{Result, anyhow}; +use anyhow::{anyhow, Result}; use tracing::debug; #[derive(Clone, Debug)] @@ -19,6 +23,7 @@ pub struct OpenHomeRenderer { info_client: Option, time_client: Option, volume_client: Option, + product_client: Option, #[allow(dead_code)] radio_client: Option, } @@ -30,6 +35,7 @@ impl OpenHomeRenderer { info_client: build_info_client(&info), time_client: build_time_client(&info), volume_client: build_volume_client(&info), + product_client: build_product_client(&info), radio_client: build_radio_client(&info), info, } @@ -68,9 +74,12 @@ impl OpenHomeRenderer { } fn playlist_client_for(&self, op: &str) -> Result<&OhPlaylistClient> { - self.playlist + let playlist = self + .playlist .as_ref() - .ok_or_else(|| op_not_supported(op, "OpenHome Playlist")) + .ok_or_else(|| op_not_supported(op, "OpenHome Playlist"))?; + self.ensure_playlist_source_selected()?; + Ok(playlist) } fn info_client_for(&self, op: &str) -> Result<&OhInfoClient> { @@ -91,19 +100,99 @@ impl OpenHomeRenderer { .ok_or_else(|| op_not_supported(op, "OpenHome Volume")) } + fn ensure_playlist_source_selected(&self) -> Result<()> { + if let Some(product) = &self.product_client { + product.ensure_playlist_source_selected().map_err(|err| { + anyhow!( + "Failed to select OpenHome Playlist source for {}: {}", + self.info.id.0, + err + ) + }) + } else { + Ok(()) + } + } + pub(crate) fn snapshot_openhome_playlist(&self) -> Result { let playlist = self.playlist_client_for("snapshot_openhome_playlist")?; let entries = playlist.read_all_tracks()?; - let current_id = self + + // Essayer d'obtenir current_id depuis Info.Id() + let mut current_id = self .info_client .as_ref() - .and_then(|client| client.id().ok()); + .and_then(|client| { + match client.id() { + Ok(id) => { + debug!( + renderer = self.info.id.0.as_str(), + current_id = id, + "OpenHome Info service returned current_id" + ); + Some(id) + } + Err(err) => { + debug!( + renderer = self.info.id.0.as_str(), + error = %err, + "OpenHome Info.Id() failed, will try Info.Track()" + ); + None + } + } + }); + + // Fallback: Si Info.Id() échoue, essayer Info.Track() et matcher l'URI + if current_id.is_none() { + if let Some(client) = self.info_client.as_ref() { + match client.track() { + Ok(track_info) => { + debug!( + renderer = self.info.id.0.as_str(), + track_uri = track_info.uri.as_str(), + "OpenHome Info.Track() returned, searching by URI" + ); + // Trouver l'entry qui matche cet URI + current_id = entries.iter() + .find(|entry| entry.uri == track_info.uri) + .map(|entry| { + debug!( + renderer = self.info.id.0.as_str(), + found_id = entry.id, + "Found current_id by matching URI" + ); + entry.id + }); + } + Err(err) => { + debug!( + renderer = self.info.id.0.as_str(), + error = %err, + "OpenHome Info.Track() also failed" + ); + } + } + } + } + + let current_index = + current_id.and_then(|id| entries.iter().position(|entry| entry.id == id)); + + debug!( + renderer = self.info.id.0.as_str(), + current_id = ?current_id, + current_index = ?current_index, + track_count = entries.len(), + "snapshot_openhome_playlist completed" + ); let tracks = entries.iter().map(convert_oh_track_entry).collect(); Ok(OpenHomePlaylistSnapshot { renderer_id: self.info.id.0.clone(), current_id, + current_index, tracks, }) } @@ -138,7 +227,11 @@ impl OpenHomeRenderer { let playlist = self.playlist_client_for("add_track_openhome")?; let insert_after = match after_id { Some(id) => id, - None => playlist.id_array()?.last().copied().unwrap_or(0), + None => playlist + .id_array()? + .last() + .copied() + .unwrap_or(OPENHOME_PLAYLIST_HEAD_ID), }; let new_id = playlist.insert(insert_after, uri, metadata)?; @@ -166,8 +259,10 @@ impl TransportControl for OpenHomeRenderer { ); } - let new_id = playlist.insert(0, uri, meta)?; - playlist.play_id(new_id) + // Reuse the same insertion logic as the queue path so that we honor + // renderer expectations (IdArray sequencing, etc.). + self.add_track_openhome(uri, meta, None, true)?; + Ok(()) } fn play(&self) -> Result<()> { @@ -311,42 +406,3 @@ fn convert_oh_track_entry(entry: &OhTrackEntry) -> OpenHomePlaylistTrack { album_art_uri: metadata.and_then(|m| m.album_art_uri), } } - -fn build_playlist_client(info: &RendererInfo) -> Option { - let control_url = info.oh_playlist_control_url.as_ref()?; - let service_type = info.oh_playlist_service_type.as_ref()?; - Some(OhPlaylistClient::new( - control_url.clone(), - service_type.clone(), - )) -} - -fn build_info_client(info: &RendererInfo) -> Option { - let control_url = info.oh_info_control_url.as_ref()?; - let service_type = info.oh_info_service_type.as_ref()?; - Some(OhInfoClient::new(control_url.clone(), service_type.clone())) -} - -fn build_time_client(info: &RendererInfo) -> Option { - let control_url = info.oh_time_control_url.as_ref()?; - let service_type = info.oh_time_service_type.as_ref()?; - Some(OhTimeClient::new(control_url.clone(), service_type.clone())) -} - -fn build_volume_client(info: &RendererInfo) -> Option { - let control_url = info.oh_volume_control_url.as_ref()?; - let service_type = info.oh_volume_service_type.as_ref()?; - Some(OhVolumeClient::new( - control_url.clone(), - service_type.clone(), - )) -} - -fn build_radio_client(info: &RendererInfo) -> Option { - let control_url = info.oh_radio_control_url.as_ref()?; - let service_type = info.oh_radio_service_type.as_ref()?; - Some(OhRadioClient::new( - control_url.clone(), - service_type.clone(), - )) -} diff --git a/pmocontrol/src/pmoserver_ext.rs b/pmocontrol/src/pmoserver_ext.rs index 310763c3..137c6fe6 100644 --- a/pmocontrol/src/pmoserver_ext.rs +++ b/pmocontrol/src/pmoserver_ext.rs @@ -4,7 +4,9 @@ //! et naviguer dans les serveurs de médias. #[cfg(feature = "pmoserver")] -use crate::control_point::{ControlPoint, OpenHomeAccessError}; +use crate::control_point::{ + ControlPoint, OpenHomeAccessError, OPENHOME_SNAPSHOT_CACHE_TTL, +}; #[cfg(feature = "pmoserver")] use crate::media_server::{MediaBrowser, MediaEntry, MusicServer, ServerId}; #[cfg(feature = "pmoserver")] @@ -1099,7 +1101,10 @@ async fn get_openhome_playlist( let rid_for_task = rid.clone(); let fetch_task = tokio::task::spawn_blocking(move || { - control_point.get_openhome_playlist_snapshot(&rid_for_task) + control_point.get_cached_openhome_playlist_snapshot( + &rid_for_task, + OPENHOME_SNAPSHOT_CACHE_TTL, + ) }); let snapshot = fetch_task diff --git a/pmocontrol/src/provider.rs b/pmocontrol/src/provider.rs index 98f8f589..29bd3aee 100644 --- a/pmocontrol/src/provider.rs +++ b/pmocontrol/src/provider.rs @@ -1,7 +1,7 @@ use std::io::BufReader; use std::time::{Duration, SystemTime}; -use quick_xml::{Error as XmlError, Reader, events::Event}; +use quick_xml::{events::Event, Error as XmlError, Reader}; use thiserror::Error; use tracing::{debug, warn}; @@ -69,6 +69,8 @@ struct ParsedDeviceDescription { oh_volume_control_url: Option, oh_radio_service_type: Option, oh_radio_control_url: Option, + oh_product_service_type: Option, + oh_product_control_url: Option, } impl ParsedDeviceDescription { @@ -301,6 +303,17 @@ impl HttpXmlDescriptionProvider { ); } } + + if lower.contains("urn:av-openhome-org:service:product:") { + if parsed.oh_product_service_type.is_none() { + parsed.oh_product_service_type = Some(st.clone()); + parsed.oh_product_control_url = Some(ctrl.clone()); + debug!( + "Found OpenHome Product for {}: type={} controlURL={}", + endpoint.udn, st, ctrl + ); + } + } } in_service = false; @@ -463,6 +476,11 @@ impl HttpXmlDescriptionProvider { .oh_radio_control_url .as_ref() .map(|ctrl| resolve_control_url(&endpoint.location, ctrl)), + oh_product_service_type: parsed.oh_product_service_type.clone(), + oh_product_control_url: parsed + .oh_product_control_url + .as_ref() + .map(|ctrl| resolve_control_url(&endpoint.location, ctrl)), }) } @@ -620,6 +638,26 @@ pub(crate) fn resolve_control_url(description_url: &str, control_url: &str) -> S control_url.to_string() } +#[cfg(test)] +mod tests { + use super::resolve_control_url; + + #[test] + fn resolves_relative_path_against_description() { + let base = "http://192.0.2.10:49152/device.xml"; + let control = "/upnp/control/playlist"; + let resolved = resolve_control_url(base, control); + assert_eq!(resolved, "http://192.0.2.10:49152/upnp/control/playlist"); + } + + #[test] + fn leaves_absolute_url_untouched() { + let url = "http://renderer.local:1400/MediaRenderer/Control"; + let resolved = resolve_control_url("http://example.invalid/device.xml", url); + assert_eq!(resolved, url); + } +} + impl DeviceDescriptionProvider for HttpXmlDescriptionProvider { fn build_renderer_info(&self, endpoint: &DiscoveredEndpoint) -> Option { match self.fetch_and_parse(endpoint) { diff --git a/pmocontrol/src/soap_client.rs b/pmocontrol/src/soap_client.rs index e20d0003..d289ccca 100644 --- a/pmocontrol/src/soap_client.rs +++ b/pmocontrol/src/soap_client.rs @@ -1,8 +1,8 @@ use std::time::Duration; use anyhow::{Context, Result}; -use pmoupnp::soap::{SoapEnvelope, build_soap_request, parse_soap_envelope}; -use tracing::{debug, trace}; +use pmoupnp::soap::{build_soap_request, parse_soap_envelope, SoapEnvelope}; +use tracing::{debug, trace, warn}; use ureq::Agent; /// Result of a SOAP call: @@ -15,6 +15,14 @@ pub struct SoapCallResult { pub envelope: Option, } +pub fn build_soap_body( + action: &str, + service_type: &str, + args: &[(&str, &str)], +) -> Result { + build_soap_request(service_type, action, args) +} + /// Invoke a UPnP SOAP action on a control URL. /// /// - `control_url`: full HTTP URL of the service control endpoint @@ -37,13 +45,15 @@ pub fn invoke_upnp_action_with_timeout( args: &[(&str, &str)], timeout: Option, ) -> Result { - let body_xml = build_soap_request(service_type, action, args) - .context("Failed to build SOAP request body")?; + let body_xml = + build_soap_body(action, service_type, args).context("Failed to build SOAP request body")?; + let arg_log = summarize_args_for_log(args); debug!( url = control_url, action = action, service_type = service_type, + args = ?arg_log, "Sending SOAP request" ); @@ -82,14 +92,145 @@ pub fn invoke_upnp_action_with_timeout( .context("Failed to read SOAP response body")?; // 6. Try to parse SOAP envelope; non-fatal on failure - let envelope = match parse_soap_envelope(raw_body.as_bytes()) { - Ok(env) => Some(env), - Err(_) => None, - }; + let parsed_envelope = parse_soap_envelope(raw_body.as_bytes()).ok(); + + if !status.is_success() { + if is_oh_info_invalid_action(service_type, action, parsed_envelope.as_ref()) { + debug!( + url = control_url, + action = action, + service_type = service_type, + status = status.as_u16(), + "OpenHome Info action not supported (Invalid Action)" + ); + } else { + warn!( + url = control_url, + action = action, + service_type = service_type, + status = status.as_u16(), + body_snippet = %response_snippet(&raw_body), + "SOAP call returned non-success status" + ); + } + } Ok(SoapCallResult { status, raw_body, - envelope, + envelope: parsed_envelope, }) } + +fn summarize_args_for_log<'a>(args: &'a [(&'a str, &'a str)]) -> Vec { + args.iter() + .map(|(name, value)| format!("{}:{}B {}", name, value.len(), preview_value(value))) + .collect() +} + +fn preview_value(value: &str) -> String { + const MAX_PREVIEW: usize = 96; + if value.len() <= MAX_PREVIEW { + value.to_string() + } else { + format!("{}…", &value[..MAX_PREVIEW]) + } +} + +fn response_snippet(body: &str) -> String { + const MAX_LEN: usize = 256; + let trimmed = body.trim(); + if trimmed.len() <= MAX_LEN { + trimmed.to_string() + } else { + format!("{}…", &trimmed[..MAX_LEN]) + } +} + +fn is_oh_info_invalid_action( + service_type: &str, + action: &str, + envelope: Option<&SoapEnvelope>, +) -> bool { + if service_type != "urn:av-openhome-org:service:Info:1" { + return false; + } + if action != "Id" && action != "TransportState" { + return false; + } + let Some(env) = envelope else { + return false; + }; + match parse_upnp_error(env) { + Some(err) if err.error_code == 401 => true, + _ => false, + } +} + +#[derive(Debug, Clone)] +pub struct UpnpError { + pub error_code: u32, + pub error_description: String, +} + +pub 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 { + xmltree::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 { + xmltree::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 find_child_with_suffix<'a>( + parent: &'a xmltree::Element, + suffix: &str, +) -> Option<&'a xmltree::Element> { + parent.children.iter().find_map(|node| match node { + xmltree::XMLNode::Element(elem) if elem.name.ends_with(suffix) => Some(elem), + _ => None, + }) +} + +#[cfg(test)] +mod tests { + use super::build_soap_body; + + #[test] + fn build_body_preserves_openhome_argument_names() { + let args = [ + ("AfterId", "0"), + ("Uri", "http://example.test/audio.flac"), + ("Metadata", ""), + ]; + let xml = + build_soap_body("Insert", "urn:av-openhome-org:service:Playlist:1", &args).unwrap(); + assert!(xml.contains("0")); + assert!(xml.contains("http://example.test/audio.flac")); + assert!(xml.contains("")); + assert!(xml.contains("")); + assert!(xml.contains("DIDL-Lite")); + assert!(!xml.contains("")); + } +} diff --git a/pmocontrol/src/upnp_renderer.rs b/pmocontrol/src/upnp_renderer.rs index 7d87bffc..df8e880e 100644 --- a/pmocontrol/src/upnp_renderer.rs +++ b/pmocontrol/src/upnp_renderer.rs @@ -203,6 +203,8 @@ mod tests { oh_volume_control_url: None, oh_radio_service_type: None, oh_radio_control_url: None, + oh_product_service_type: None, + oh_product_control_url: None, } } diff --git a/pmocovers/src/api.rs b/pmocovers/src/api.rs index 7e227876..060f8c2f 100644 --- a/pmocovers/src/api.rs +++ b/pmocovers/src/api.rs @@ -2,12 +2,7 @@ use crate::cache; use crate::Cache; -use axum::{ - extract::State, - http::StatusCode, - response::IntoResponse, - Json, -}; +use axum::{extract::State, http::StatusCode, response::IntoResponse, Json}; use pmocache::api::{AddItemRequest, AddItemResponse, ErrorResponse}; use std::sync::Arc; diff --git a/pmoserver/src/server.rs b/pmoserver/src/server.rs index f7fc0d6c..e04f4247 100644 --- a/pmoserver/src/server.rs +++ b/pmoserver/src/server.rs @@ -551,7 +551,13 @@ impl Server { self.join_handle = Some(tokio::spawn(async move { let server_future = async { let r = router.read().await.clone(); - let listener = tokio::net::TcpListener::bind(addr).await.unwrap(); + let listener = match tokio::net::TcpListener::bind(addr).await { + Ok(l) => l, + Err(e) => { + error!("Failed to bind to {}: {}", addr, e); + panic!("Cannot start server: {}", e); + } + }; axum::serve(listener, r.into_make_service()) .with_graceful_shutdown(async move { diff --git a/pmoupnp/src/upnp_server.rs b/pmoupnp/src/upnp_server.rs index a29d1179..e76ccf40 100644 --- a/pmoupnp/src/upnp_server.rs +++ b/pmoupnp/src/upnp_server.rs @@ -25,15 +25,15 @@ use std::sync::RwLock; use pmoserver::Server; use utoipa::OpenApi; +use crate::UpnpModel; use crate::devices::errors::DeviceError; use crate::devices::{Device, DeviceInstance, DeviceRegistry}; use crate::ssdp::SsdpServer; use crate::upnp_api::UpnpApiExt; -use crate::UpnpModel; use pmoaudiocache::Cache as AudioCache; use pmocovers::Cache as CoverCache; -use pmoutils::{find_process_using_port, TransportProtocol}; +use pmoutils::{TransportProtocol, find_process_using_port}; /// Registre de devices global et thread-safe. ///