//! End-to-end queue demo that prefers the PMOMusic media server and exercises //! the ControlPoint playback queue API. use std::collections::VecDeque; use std::env; use std::process; use std::thread; use std::time::Duration; use anyhow::{Context, Result}; use pmocontrol::model::TrackMetadata; use pmocontrol::{ ControlPoint, DeviceRegistryRead, MediaBrowser, MediaEntry, MediaServerEvent, MediaServerInfo, MusicRenderer, MusicServer, PlaybackItem, PlaybackPosition, PlaybackPositionInfo, RendererInfo, }; const DEFAULT_TIMEOUT_SECS: u64 = 5; const DEFAULT_DISCOVERY_SECS: u64 = 5; const DEFAULT_MAX_TRACKS: usize = 3; const MONITOR_DURATION_SECS: u64 = 600; const MONITOR_POLL_SECS: u64 = 5; const MAX_BROWSE_DEPTH: usize = 2; fn main() -> Result<()> { let _ = tracing_subscriber::fmt::try_init(); let config = CliConfig::parse_from_env().unwrap_or_else(|err| { eprintln!("Error parsing arguments: {err}"); print_usage_and_exit(); }); if config.max_tracks == 0 { eprintln!("--max-tracks must be >= 1"); process::exit(1); } println!( "Starting queue_pmomusic_demo with timeout={}s discovery={}s max_tracks={}", config.timeout_secs, config.discovery_secs, config.max_tracks ); // ControlPoint::spawn starts the HttpXmlDescriptionProvider + DiscoveryManager combo. let control_point = ControlPoint::spawn(config.timeout_secs).context("Failed to start control point")?; println!( "Discovery running for {} seconds before selecting devices...", config.discovery_secs ); thread::sleep(Duration::from_secs(config.discovery_secs)); let registry = control_point.registry(); let (renderer, server_info) = { let reg = registry.read().expect("registry poisoned"); let renderer_candidates: Vec = reg .list_renderers() .into_iter() .filter(|info| !is_pmomusic_renderer(info)) .collect(); let renderer = pick_renderer(renderer_candidates) .unwrap_or_else(|| no_renderer_and_exit("No suitable renderer found after discovery.")); let server = pick_media_server(reg.list_servers()) .unwrap_or_else(|| no_server_and_exit("No media server with ContentDirectory.")); (renderer, server) }; println!( "Selected renderer \"{}\" (protocol={:?}, id={})", renderer.friendly_name, renderer.protocol, renderer.id.0 ); println!( "Selected media server \"{}\" at {} (id={})", server_info.friendly_name, server_info.location, server_info.id.0 ); let renderer_instance = MusicRenderer::from_registry_info(renderer.clone(), ®istry) .expect("Selected renderer is not usable by MusicRenderer façade"); let supports_set_next = renderer_instance .as_upnp() .map(|upnp| upnp.supports_set_next()) .unwrap_or(false); println!( "Renderer \"{}\": AVTransport present = {}, SetNextAVTransportURI supported = {}", renderer.friendly_name, renderer.capabilities.has_avtransport, supports_set_next ); let timeout = Duration::from_secs(config.timeout_secs); let server = MusicServer::from_info(&server_info, timeout).context("Failed to init MusicServer")?; let root_entries = server .browse_root() .context("Failed to browse ContentDirectory root")?; println!("Root returned {} entries", root_entries.len()); // Try to find a playlist container first let (playback_items, bound_container_id) = collect_playable_items_with_binding(&server, &root_entries, config.max_tracks) .context("Failed to derive playable items from ContentDirectory root/children")?; if playback_items.is_empty() { println!("No playable tracks were found on the selected server."); process::exit(1); } println!( "Discovered {} playable items; enqueuing…", playback_items.len() ); if let Some(ref container_id) = bound_container_id { println!( "Found playlist container '{}' to bind queue to", container_id ); } else { println!("No playlist container found; queue will not be bound to server"); } let mut planned_queue: VecDeque; let renderer_id = renderer.id.clone(); control_point .clear_queue(&renderer_id) .context("Failed to clear playback queue")?; control_point .enqueue_items(&renderer_id, playback_items) .context("Failed to enqueue playback items")?; // Attach queue to playlist container if we found one if let Some(container_id) = bound_container_id { control_point .attach_queue_to_playlist(&renderer_id, server_info.id.clone(), container_id.clone()) .context("Failed to attach queue to playlist container")?; println!( "✓ Queue attached to playlist container {} on server {}", container_id, server_info.friendly_name ); } let snapshot = control_point .get_queue_snapshot(&renderer_id) .context("Failed to snapshot queue after enqueue")?; print_queue_snapshot(&snapshot); planned_queue = snapshot.clone().into(); control_point .play_next_from_queue(&renderer_id) .context("Failed to start playback from queue")?; let mut current_track = planned_queue.pop_front(); let remaining = control_point .get_queue_snapshot(&renderer_id) .context("Failed to snapshot queue after play_next_from_queue")?; planned_queue = remaining.clone().into(); println!( "Playback started on \"{}\"; {} tracks remaining in queue.", renderer.friendly_name, remaining.len() ); println!( "Monitoring queue auto-advance for {} seconds (poll every {}s)…", MONITOR_DURATION_SECS, MONITOR_POLL_SECS ); // Subscribe to media server events to observe playlist updates let media_event_rx = control_point.subscribe_media_server_events(); let poll_count = MONITOR_DURATION_SECS / MONITOR_POLL_SECS; for tick in 0..poll_count { thread::sleep(Duration::from_secs(MONITOR_POLL_SECS)); // Drain any MediaServerEvent that arrived since last poll loop { match media_event_rx.try_recv() { Ok(MediaServerEvent::GlobalUpdated { server_id, system_update_id, }) => { println!( " 📢 MediaServer {} global update (SystemUpdateID={:?})", server_id.0, system_update_id ); } Ok(MediaServerEvent::ContainersUpdated { server_id, container_ids, }) => { println!( " 📢 MediaServer {} containers updated: {:?}", server_id.0, container_ids ); // Check if our bound container was updated if let Some((bound_server, bound_container, _)) = control_point.current_queue_playlist_binding(&renderer_id) { if bound_server == server_id && container_ids.contains(&bound_container) { println!( " 🔄 Bound playlist container '{}' was updated, queue will refresh automatically", bound_container ); } } } Err(_) => break, // No more events, continue with normal monitoring } } let snapshot = control_point .get_queue_snapshot(&renderer_id) .context("Queue snapshot failed during monitoring loop")?; let new_plan: VecDeque = snapshot.clone().into(); if planned_queue.len() > new_plan.len() { let removed = planned_queue.len() - new_plan.len(); for _ in 0..removed { current_track = planned_queue.pop_front(); } } planned_queue = new_plan; let playback_info = control_point .music_renderer_by_id(&renderer_id) .and_then(|renderer| renderer.playback_position().ok()); let title = current_track_title(current_track.as_ref()); if let Some(info) = playback_info { println!( "[tick {tick}] Queue length = {} | now playing: {} [{}]", snapshot.len(), title, format_playback_position(&info) ); } else { println!( "[tick {tick}] Queue length = {} | now playing: {} [position unavailable]", snapshot.len(), title ); } } println!("Monitoring finished, exiting."); Ok(()) } #[derive(Debug)] struct CliConfig { timeout_secs: u64, discovery_secs: u64, max_tracks: usize, } impl CliConfig { fn parse_from_env() -> Result { let mut timeout_secs = DEFAULT_TIMEOUT_SECS; let mut discovery_secs = DEFAULT_DISCOVERY_SECS; let mut max_tracks = DEFAULT_MAX_TRACKS; let mut args = env::args().skip(1); while let Some(arg) = args.next() { match arg.as_str() { "--timeout-secs" => { let value = args .next() .ok_or_else(|| "--timeout-secs requires a value".to_string())?; timeout_secs = value.parse().map_err(|err| { format!("Invalid value for --timeout-secs ({value}): {err}") })?; } "--discovery-secs" => { let value = args .next() .ok_or_else(|| "--discovery-secs requires a value".to_string())?; discovery_secs = value.parse().map_err(|err| { format!("Invalid value for --discovery-secs ({value}): {err}") })?; } "--max-tracks" => { let value = args .next() .ok_or_else(|| "--max-tracks requires a value".to_string())?; max_tracks = value.parse().map_err(|err| { format!("Invalid value for --max-tracks ({value}): {err}") })?; } "--help" | "-h" => { print_usage_and_exit(); } unknown => { return Err(format!("Unknown argument: {unknown}")); } } } Ok(Self { timeout_secs, discovery_secs, max_tracks, }) } } fn pick_renderer(renderers: Vec) -> Option { let mut candidates: Vec = renderers; if candidates.is_empty() { return None; } println!("Renderer candidates:"); for (idx, info) in candidates.iter().enumerate() { println!( " [{}] {} | model={} | location={} | online={}", idx, info.friendly_name, info.model_name, info.location, info.online ); } let selected = candidates.remove(0); println!( "Automatically selecting renderer index 0: {}", selected.friendly_name ); Some(selected) } fn pick_media_server(servers: Vec) -> Option { let mut candidates: Vec = servers .into_iter() .filter(|info| info.has_content_directory) .filter(|info| info.content_directory_control_url.is_some()) .collect(); if candidates.is_empty() { return None; } if let Some(idx) = candidates.iter().position(is_pmomusic_server) { let server = candidates.remove(idx); println!( "Preferring PMOMusic server \"{}\" (server header: {}).", server.friendly_name, server.server_header ); Some(server) } else { println!("No PMOMusic server discovered; falling back to first ContentDirectory server."); Some(candidates.remove(0)) } } fn is_pmomusic_server(info: &MediaServerInfo) -> bool { let name = info.friendly_name.to_ascii_lowercase(); let header = info.server_header.to_ascii_lowercase(); name.contains("pmomusic") || header.contains("pmomusic") } fn is_pmomusic_renderer(info: &RendererInfo) -> bool { info.friendly_name .to_ascii_lowercase() .contains("pmomusic audio renderer") } fn collect_playable_items_with_binding( server: &MusicServer, entries: &[MediaEntry], max_tracks: usize, ) -> Result<(Vec, Option)> { // First, try to find a playlist container let playlist_container = entries.iter().find(|entry| { entry.is_container && entry .class .to_ascii_lowercase() .contains("object.container.playlistcontainer") }); if let Some(playlist) = playlist_container { println!( "Found playlist container: '{}' (id: {}, class: {})", playlist.title, playlist.id, playlist.class ); // Browse the playlist container match server.browse_children(&playlist.id, 0, max_tracks as u32) { Ok(children) => { let mut items = Vec::new(); for entry in &children { if let Some(item) = playback_item_from_entry(server, entry) { items.push(item); if items.len() >= max_tracks { break; } } } if !items.is_empty() { return Ok((items, Some(playlist.id.clone()))); } println!( "Playlist container '{}' is empty, falling back to general browse", playlist.title ); } Err(err) => { println!( "Failed to browse playlist container '{}': {}, falling back", playlist.title, err ); } } } else { println!("No playlist container found in root entries, using fallback"); } // Fallback: collect from any container/item let mut items = Vec::new(); for entry in entries { gather_items_from_entry(server, entry, max_tracks, 0, &mut items)?; if items.len() >= max_tracks { break; } } Ok((items, None)) } fn gather_items_from_entry( server: &MusicServer, entry: &MediaEntry, max_tracks: usize, depth: usize, out: &mut Vec, ) -> Result<()> { if out.len() >= max_tracks { return Ok(()); } if entry.is_container { if depth >= MAX_BROWSE_DEPTH { return Ok(()); } match server.browse_children(&entry.id, 0, 50) { Ok(children) => { for child in children { gather_items_from_entry(server, &child, max_tracks, depth + 1, out)?; if out.len() >= max_tracks { break; } } } Err(err) => { tracing::warn!( container_id = entry.id.as_str(), error = %err, "Failed to browse child container" ); } } return Ok(()); } if let Some(item) = playback_item_from_entry(server, entry) { out.push(item); } Ok(()) } fn playback_item_from_entry(server: &MusicServer, entry: &MediaEntry) -> Option { if entry.title.to_ascii_lowercase().contains("live stream") { return None; } let resource = entry.resources.iter().find(|res| res.is_audio())?; let metadata = TrackMetadata { title: Some(entry.title.clone()), artist: entry.artist.clone(), album: entry.album.clone(), genre: entry.genre.clone(), album_art_uri: entry.album_art_uri.clone(), date: entry.date.clone(), track_number: entry.track_number.clone(), creator: entry.creator.clone(), }; Some(PlaybackItem { media_server_id: server.id().clone(), didl_id: entry.id.clone(), uri: resource.uri.clone(), metadata: Some(metadata), }) } fn print_queue_snapshot(items: &[PlaybackItem]) { println!("Current queue snapshot ({} items):", items.len()); for (idx, item) in items.iter().enumerate() { let label = item .metadata .as_ref() .and_then(|meta| meta.title.as_deref()) .unwrap_or_else(|| item.uri.as_str()); println!(" [{}] {} -> {}", idx, label, item.uri); } if items.is_empty() { println!(" "); } } fn current_track_title(item: Option<&PlaybackItem>) -> String { match item { Some(track) => track .metadata .as_ref() .and_then(|meta| meta.title.as_deref()) .unwrap_or_else(|| track.uri.as_str()) .to_string(), None => "".to_string(), } } fn format_playback_position(info: &PlaybackPositionInfo) -> String { let rel = info.rel_time.as_deref().unwrap_or("-"); let dur = info.track_duration.as_deref().unwrap_or("-"); format!("{rel} / {dur}") } fn no_renderer_and_exit(message: &str) -> ! { println!("{message}"); process::exit(1); } fn no_server_and_exit(message: &str) -> ! { println!("{message}"); process::exit(1); } fn print_usage_and_exit() -> ! { println!( "Usage: cargo run -p pmocontrol --example queue_pmomusic_demo -- [--timeout-secs N] [--discovery-secs N] [--max-tracks N]" ); process::exit(1); }