Mise en place de la synchronisation de la playlist du point de contrôle avec celle du serveur de média.

This commit is contained in:
2025-12-02 14:41:53 +01:00
parent aeb9332387
commit f99ab2a41e
9 changed files with 19861 additions and 59 deletions

1
package.json Normal file
View File

@@ -0,0 +1 @@
{}

View File

@@ -15,6 +15,7 @@ fn main() -> std::io::Result<()> {
"urn:schemas-upnp-org:device:MediaRenderer:1",
"urn:av-openhome-org:device:MediaRenderer:1",
"urn:schemas-upnp-org:device:MediaServer:1",
"urn:schemas-wiimu-com:service:PlayQueue:1",
];
for st in &search_targets {
if let Err(err) = client.send_msearch(st, 3) {

View File

@@ -0,0 +1,749 @@
//! Full interactive control point CLI demo.
//!
//! This example demonstrates all capabilities of the ControlPoint:
//! - Interactive device selection (renderer + media server)
//! - ContentDirectory navigation with back/forward support
//! - Queue construction and playlist binding
//! - Full playback control (play/pause/stop/next/volume/mute)
//! - Real-time event monitoring (renderer and media server events)
//! - Live playlist observation and auto-refresh
use std::io::{self, BufRead, Write};
use std::process;
use std::sync::Arc;
use std::thread;
use std::time::Duration;
use anyhow::{anyhow, Context, Result};
use pmocontrol::{
ControlPoint, DeviceRegistryRead, MediaBrowser, MediaEntry, MediaResource, MediaServerEvent,
MediaServerInfo, MusicServer, PlaybackItem, PlaybackPosition, PlaybackStatus, RendererEvent,
RendererInfo, TransportControl, VolumeControl,
};
const DEFAULT_TIMEOUT_SECS: u64 = 5;
const DEFAULT_DISCOVERY_SECS: u64 = 15;
fn main() -> Result<()> {
let _ = tracing_subscriber::fmt::try_init();
println!("=== Full Control Point Interactive Demo ===");
println!("Starting control point with timeout={}s", DEFAULT_TIMEOUT_SECS);
let control_point = ControlPoint::spawn(DEFAULT_TIMEOUT_SECS)
.context("Failed to start control point")?;
println!("Discovery running for {} seconds...", DEFAULT_DISCOVERY_SECS);
thread::sleep(Duration::from_secs(DEFAULT_DISCOVERY_SECS));
// Step 1: List and select renderer
let registry = control_point.registry();
let renderer_info = {
let reg = registry.read().expect("registry poisoned");
let renderers = reg.list_renderers();
println!("\n=== Available Media Renderers ===");
if renderers.is_empty() {
eprintln!("No renderers discovered. Exiting.");
process::exit(1);
}
print_renderers(&renderers);
select_renderer(&renderers)?
};
println!("\n✓ Selected renderer: {} (id={})",
renderer_info.friendly_name, renderer_info.id.0);
// Step 2: List and select media server
let server_info = {
let reg = registry.read().expect("registry poisoned");
let servers: Vec<MediaServerInfo> = reg.list_servers()
.into_iter()
.filter(|s| s.has_content_directory && s.content_directory_control_url.is_some())
.collect();
println!("\n=== Available Media Servers ===");
if servers.is_empty() {
eprintln!("No media servers with ContentDirectory discovered. Exiting.");
process::exit(1);
}
print_servers(&servers);
select_server(&servers)?
};
println!("\n✓ Selected server: {} (id={})",
server_info.friendly_name, server_info.id.0);
let timeout = Duration::from_secs(DEFAULT_TIMEOUT_SECS);
let server = MusicServer::from_info(&server_info, timeout)
.context("Failed to initialize MusicServer")?;
// Step 3: Navigate ContentDirectory and select a container for queue
println!("\n=== ContentDirectory Navigation ===");
println!("Commands: [number]=enter container, 'b'=back, 's'=select current as queue source, 'q'=quit navigation");
let (selected_items, selected_container_id) = navigate_and_select(&server)?;
if selected_items.is_empty() {
println!("No playable items selected. Exiting.");
process::exit(1);
}
println!("\n✓ Selected {} items for playback queue", selected_items.len());
// Step 4: Build queue
let renderer_id = renderer_info.id.clone();
println!("\n=== Building Playback Queue ===");
control_point.clear_queue(&renderer_id)
.context("Failed to clear queue")?;
control_point.enqueue_items(&renderer_id, selected_items)
.context("Failed to enqueue items")?;
println!("✓ Queue built with {} items",
control_point.get_queue_snapshot(&renderer_id)?.len());
// Step 5: Ask about playlist binding
if let Some(container_id) = selected_container_id {
println!("\nAttach this queue to playlist container '{}' for auto-refresh? (y/n)", container_id);
if read_yes_no()? {
control_point.attach_queue_to_playlist(
&renderer_id,
server_info.id.clone(),
container_id.clone(),
);
println!("✓ Queue attached to playlist container '{}'", container_id);
} else {
println!("Queue will not be bound to playlist (no auto-refresh)");
}
}
// Step 6: Start playback
println!("\n=== Starting Playback ===");
control_point.play_next_from_queue(&renderer_id)
.context("Failed to start playback")?;
println!("✓ Playback started");
// Step 7: Spawn event monitoring threads
let control_point_arc = Arc::new(control_point);
spawn_renderer_event_thread(Arc::clone(&control_point_arc), renderer_id.clone());
spawn_media_server_event_thread(Arc::clone(&control_point_arc), renderer_id.clone());
// Step 8: Interactive control loop
println!("\n=== Interactive Control ===");
print_help();
run_control_loop(Arc::clone(&control_point_arc), renderer_id)?;
println!("\nExiting. Goodbye!");
Ok(())
}
/// Print list of renderers with index.
fn print_renderers(renderers: &[RendererInfo]) {
for (idx, info) in renderers.iter().enumerate() {
println!(" [{}] {} | model={} | protocol={:?} | online={} | id={}",
idx, info.friendly_name, info.model_name, info.protocol, info.online, info.id.0);
}
}
/// Print list of servers with index.
fn print_servers(servers: &[MediaServerInfo]) {
for (idx, info) in servers.iter().enumerate() {
println!(" [{}] {} | model={} | manufacturer={} | online={} | id={}",
idx, info.friendly_name, info.model_name, info.manufacturer, info.online, info.id.0);
}
}
/// Interactive renderer selection.
fn select_renderer(renderers: &[RendererInfo]) -> Result<RendererInfo> {
loop {
print!("\nSelect renderer (index 0-{}): ", renderers.len() - 1);
io::stdout().flush()?;
let mut input = String::new();
io::stdin().read_line(&mut input)?;
match input.trim().parse::<usize>() {
Ok(idx) if idx < renderers.len() => {
return Ok(renderers[idx].clone());
}
_ => {
println!("Invalid selection. Please enter a number between 0 and {}", renderers.len() - 1);
}
}
}
}
/// Interactive server selection.
fn select_server(servers: &[MediaServerInfo]) -> Result<MediaServerInfo> {
loop {
print!("\nSelect media server (index 0-{}): ", servers.len() - 1);
io::stdout().flush()?;
let mut input = String::new();
io::stdin().read_line(&mut input)?;
match input.trim().parse::<usize>() {
Ok(idx) if idx < servers.len() => {
return Ok(servers[idx].clone());
}
_ => {
println!("Invalid selection. Please enter a number between 0 and {}", servers.len() - 1);
}
}
}
}
/// Navigation state for ContentDirectory browsing.
struct NavigationState {
/// Stack of (container_id, container_title) for back navigation.
path_stack: Vec<(String, String)>,
/// Current container ID being browsed.
current_container_id: String,
/// Current container title.
current_container_title: String,
}
impl NavigationState {
fn new(root_id: String, root_title: String) -> Self {
Self {
path_stack: Vec::new(),
current_container_id: root_id,
current_container_title: root_title,
}
}
fn enter_container(&mut self, container_id: String, container_title: String) {
self.path_stack.push((
self.current_container_id.clone(),
self.current_container_title.clone(),
));
self.current_container_id = container_id;
self.current_container_title = container_title;
}
fn go_back(&mut self) -> bool {
if let Some((parent_id, parent_title)) = self.path_stack.pop() {
self.current_container_id = parent_id;
self.current_container_title = parent_title;
true
} else {
false
}
}
}
/// Navigate ContentDirectory and let user select a container for queue.
fn navigate_and_select(server: &MusicServer) -> Result<(Vec<PlaybackItem>, Option<String>)> {
let root_entries = server.browse_root()
.context("Failed to browse root")?;
let mut nav_state = NavigationState::new("0".to_string(), "Root".to_string());
loop {
println!("\n--- Browsing: {} (id: {}) ---",
nav_state.current_container_title, nav_state.current_container_id);
let entries = if nav_state.current_container_id == "0" {
root_entries.clone()
} else {
server.browse_children(&nav_state.current_container_id, 0, 100)
.context("Failed to browse container")?
};
if entries.is_empty() {
println!("(empty container)");
} else {
print_entries(&entries);
}
print!("\nCommand [0-{} to enter container / b=back / s=select / q=quit]: ",
entries.len().saturating_sub(1));
io::stdout().flush()?;
let mut input = String::new();
io::stdin().read_line(&mut input)?;
let input = input.trim();
match input {
"q" => {
println!("Navigation cancelled.");
return Ok((Vec::new(), None));
}
"b" => {
if !nav_state.go_back() {
println!("Already at root.");
}
}
"s" => {
// Select current container
println!("Collecting items from '{}'...", nav_state.current_container_title);
let items = collect_playable_items(server, &entries)?;
if items.is_empty() {
println!("No playable items found in this container.");
continue;
}
return Ok((items, Some(nav_state.current_container_id.clone())));
}
_ => {
// Try to parse as index
match input.parse::<usize>() {
Ok(idx) if idx < entries.len() => {
let entry = &entries[idx];
if entry.is_container {
nav_state.enter_container(entry.id.clone(), entry.title.clone());
} else {
println!("Entry {} is an audio item, not a container. Use 's' to select the current container for playback.", idx);
}
}
_ => {
println!("Invalid command. Enter a number (0-{}), 'b' to go back, 's' to select, or 'q' to quit.",
entries.len().saturating_sub(1));
}
}
}
}
}
}
/// Print MediaEntry list.
fn print_entries(entries: &[MediaEntry]) {
for (idx, entry) in entries.iter().enumerate() {
if entry.is_container {
println!(" [{}] \u{1F4C1} {} (id: {}, class: {})",
idx, entry.title, entry.id, entry.class);
} else {
let has_audio = entry.resources.iter().any(is_audio_resource);
let icon = if has_audio { "\u{266B}" } else { "\u{1F4C4}" };
println!(" [{}] {} {} (id: {}, class: {})",
idx, icon, entry.title, entry.id, entry.class);
}
}
}
/// Collect playable items from MediaEntry list (including nested containers).
fn collect_playable_items(server: &MusicServer, entries: &[MediaEntry]) -> Result<Vec<PlaybackItem>> {
let mut items = Vec::new();
for entry in entries {
if entry.is_container {
// Recursively collect from container
match server.browse_children(&entry.id, 0, 100) {
Ok(children) => {
items.extend(collect_playable_items(server, &children)?);
}
Err(err) => {
eprintln!("Warning: failed to browse container '{}': {}", entry.title, err);
}
}
} else {
if let Some(item) = playback_item_from_entry(server, entry) {
items.push(item);
}
}
}
Ok(items)
}
/// Convert MediaEntry to PlaybackItem.
fn playback_item_from_entry(server: &MusicServer, entry: &MediaEntry) -> Option<PlaybackItem> {
let resource = entry.resources.iter().find(|res| is_audio_resource(res))?;
let mut item = PlaybackItem::new(resource.uri.clone());
item.title = Some(entry.title.clone());
item.server_id = Some(server.id().clone());
item.object_id = Some(entry.id.clone());
Some(item)
}
/// Check if MediaResource is audio.
fn is_audio_resource(res: &MediaResource) -> bool {
let lower = res.protocol_info.to_ascii_lowercase();
if lower.contains("audio/") {
return true;
}
lower
.split(':')
.nth(2)
.map(|mime| mime.starts_with("audio/"))
.unwrap_or(false)
}
/// Read yes/no from stdin.
fn read_yes_no() -> Result<bool> {
loop {
print!("> ");
io::stdout().flush()?;
let mut input = String::new();
io::stdin().read_line(&mut input)?;
match input.trim().to_lowercase().as_str() {
"y" | "yes" => return Ok(true),
"n" | "no" => return Ok(false),
_ => println!("Please enter 'y' or 'n'"),
}
}
}
/// Spawn thread to display renderer events in real-time.
fn spawn_renderer_event_thread(
control_point: Arc<ControlPoint>,
renderer_id: pmocontrol::model::RendererId,
) {
let event_rx = control_point.subscribe_events();
thread::spawn(move || {
loop {
match event_rx.recv() {
Ok(event) => {
match event {
RendererEvent::StateChanged { id, state } => {
if id == renderer_id {
println!("\n[EVENT] Renderer state changed: {:?}", state);
print!("> ");
io::stdout().flush().ok();
}
}
RendererEvent::PositionChanged { id, position } => {
if id == renderer_id {
let rel = position.rel_time.as_deref().unwrap_or("-");
let dur = position.track_duration.as_deref().unwrap_or("-");
println!("\n[EVENT] Position: {} / {}", rel, dur);
print!("> ");
io::stdout().flush().ok();
}
}
RendererEvent::VolumeChanged { id, volume } => {
if id == renderer_id {
println!("\n[EVENT] Volume changed: {}", volume);
print!("> ");
io::stdout().flush().ok();
}
}
RendererEvent::MuteChanged { id, mute } => {
if id == renderer_id {
println!("\n[EVENT] Mute changed: {}", mute);
print!("> ");
io::stdout().flush().ok();
}
}
}
}
Err(_) => {
eprintln!("\n[EVENT] Renderer event channel closed");
break;
}
}
}
});
}
/// Spawn thread to display media server events in real-time.
fn spawn_media_server_event_thread(
control_point: Arc<ControlPoint>,
renderer_id: pmocontrol::model::RendererId,
) {
let media_event_rx = control_point.subscribe_media_server_events();
thread::spawn(move || {
loop {
match media_event_rx.recv() {
Ok(event) => {
match event {
MediaServerEvent::GlobalUpdated { server_id, system_update_id } => {
println!("\n[MEDIA EVENT] Server {} global update (SystemUpdateID={})",
server_id.0, system_update_id.unwrap_or(0));
print!("> ");
io::stdout().flush().ok();
}
MediaServerEvent::ContainersUpdated { server_id, container_ids } => {
println!("\n[MEDIA EVENT] Server {} containers updated: {:?}",
server_id.0, container_ids);
// Check if 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!("[MEDIA EVENT] → Bound playlist '{}' was updated!", bound_container);
}
}
print!("> ");
io::stdout().flush().ok();
}
}
}
Err(_) => {
eprintln!("\n[MEDIA EVENT] Media server event channel closed");
break;
}
}
}
});
}
/// Print help message for control commands.
fn print_help() {
println!("Available commands:");
println!(" p - Pause");
println!(" r - Resume/Play");
println!(" s - Stop");
println!(" n - Next track");
println!(" + - Volume +5");
println!(" - - Volume -5");
println!(" m - Toggle mute");
println!(" i - Show renderer info (state, position, volume)");
println!(" k - Show current queue");
println!(" b - Show playlist binding");
println!(" h - Show this help");
println!(" q - Quit gracefully");
println!(" Q - Quit immediately");
}
/// Main interactive control loop.
fn run_control_loop(
control_point: Arc<ControlPoint>,
renderer_id: pmocontrol::model::RendererId,
) -> Result<()> {
let stdin = io::stdin();
let mut reader = stdin.lock();
loop {
print!("> ");
io::stdout().flush()?;
let mut line = String::new();
reader.read_line(&mut line)?;
let cmd = line.trim();
if cmd.is_empty() {
continue;
}
match cmd {
"q" => {
println!("Quitting...");
break;
}
"Q" => {
println!("Quitting immediately!");
process::exit(0);
}
"h" => {
print_help();
}
"p" => {
if let Err(err) = pause_renderer(&control_point, &renderer_id) {
eprintln!("Pause failed: {}", err);
} else {
println!("✓ Paused");
}
}
"r" => {
if let Err(err) = resume_renderer(&control_point, &renderer_id) {
eprintln!("Resume failed: {}", err);
} else {
println!("✓ Resumed");
}
}
"s" => {
if let Err(err) = stop_renderer(&control_point, &renderer_id) {
eprintln!("Stop failed: {}", err);
} else {
println!("✓ Stopped");
}
}
"n" => {
if let Err(err) = control_point.play_next_from_queue(&renderer_id) {
eprintln!("Next track failed: {}", err);
} else {
println!("✓ Playing next track");
}
}
"+" => {
if let Err(err) = adjust_volume(&control_point, &renderer_id, 5) {
eprintln!("Volume adjustment failed: {}", err);
} else {
println!("✓ Volume +5");
}
}
"-" => {
if let Err(err) = adjust_volume(&control_point, &renderer_id, -5) {
eprintln!("Volume adjustment failed: {}", err);
} else {
println!("✓ Volume -5");
}
}
"m" => {
if let Err(err) = toggle_mute(&control_point, &renderer_id) {
eprintln!("Mute toggle failed: {}", err);
} else {
println!("✓ Mute toggled");
}
}
"i" => {
if let Err(err) = show_renderer_info(&control_point, &renderer_id) {
eprintln!("Failed to get renderer info: {}", err);
}
}
"k" => {
if let Err(err) = show_queue(&control_point, &renderer_id) {
eprintln!("Failed to get queue: {}", err);
}
}
"b" => {
show_binding(&control_point, &renderer_id);
}
_ => {
println!("Unknown command '{}'. Type 'h' for help.", cmd);
}
}
}
Ok(())
}
/// Pause the renderer.
fn pause_renderer(
control_point: &ControlPoint,
renderer_id: &pmocontrol::model::RendererId,
) -> Result<()> {
let renderer = control_point.music_renderer_by_id(renderer_id)
.ok_or_else(|| anyhow!("Renderer not found"))?;
renderer.pause()?;
Ok(())
}
/// Resume/play the renderer.
fn resume_renderer(
control_point: &ControlPoint,
renderer_id: &pmocontrol::model::RendererId,
) -> Result<()> {
let renderer = control_point.music_renderer_by_id(renderer_id)
.ok_or_else(|| anyhow!("Renderer not found"))?;
renderer.play()?;
Ok(())
}
/// Stop the renderer.
fn stop_renderer(
control_point: &ControlPoint,
renderer_id: &pmocontrol::model::RendererId,
) -> Result<()> {
let renderer = control_point.music_renderer_by_id(renderer_id)
.ok_or_else(|| anyhow!("Renderer not found"))?;
renderer.stop()?;
Ok(())
}
/// Adjust volume by delta (-100 to +100).
fn adjust_volume(
control_point: &ControlPoint,
renderer_id: &pmocontrol::model::RendererId,
delta: i32,
) -> Result<()> {
let renderer = control_point.music_renderer_by_id(renderer_id)
.ok_or_else(|| anyhow!("Renderer not found"))?;
let current = renderer.volume()?;
let new_volume = (current as i32 + delta).clamp(0, 100) as u16;
renderer.set_volume(new_volume)?;
Ok(())
}
/// Toggle mute.
fn toggle_mute(
control_point: &ControlPoint,
renderer_id: &pmocontrol::model::RendererId,
) -> Result<()> {
let renderer = control_point.music_renderer_by_id(renderer_id)
.ok_or_else(|| anyhow!("Renderer not found"))?;
let current_mute = renderer.mute()?;
renderer.set_mute(!current_mute)?;
Ok(())
}
/// Show renderer info (state, position, volume, mute).
fn show_renderer_info(
control_point: &ControlPoint,
renderer_id: &pmocontrol::model::RendererId,
) -> Result<()> {
let renderer = control_point.music_renderer_by_id(renderer_id)
.ok_or_else(|| anyhow!("Renderer not found"))?;
println!("\n=== Renderer Info ===");
println!("Name: {}", renderer.info().friendly_name);
match renderer.playback_state() {
Ok(state) => println!("State: {:?}", state),
Err(err) => println!("State: <error: {}>", err),
}
match renderer.playback_position() {
Ok(pos) => {
let rel = pos.rel_time.as_deref().unwrap_or("-");
let dur = pos.track_duration.as_deref().unwrap_or("-");
println!("Position: {} / {}", rel, dur);
}
Err(err) => println!("Position: <error: {}>", err),
}
match renderer.volume() {
Ok(vol) => println!("Volume: {}", vol),
Err(err) => println!("Volume: <error: {}>", err),
}
match renderer.mute() {
Ok(mute) => println!("Mute: {}", mute),
Err(err) => println!("Mute: <error: {}>", err),
}
Ok(())
}
/// Show current queue.
fn show_queue(
control_point: &ControlPoint,
renderer_id: &pmocontrol::model::RendererId,
) -> Result<()> {
let queue = control_point.get_queue_snapshot(renderer_id)?;
println!("\n=== Queue ({} items) ===", queue.len());
if queue.is_empty() {
println!(" <empty>");
} else {
for (idx, item) in queue.iter().enumerate() {
let title = item.title.as_deref().unwrap_or("<no title>");
println!(" [{}] {}", idx, title);
}
}
Ok(())
}
/// Show playlist binding info.
fn show_binding(
control_point: &ControlPoint,
renderer_id: &pmocontrol::model::RendererId,
) {
match control_point.current_queue_playlist_binding(renderer_id) {
Some((server_id, container_id, has_seen_update)) => {
println!("\n=== Playlist Binding ===");
println!("Server ID: {}", server_id.0);
println!("Container ID: {}", container_id);
println!("Has seen update: {}", has_seen_update);
}
None => {
println!("\n=== Playlist Binding ===");
println!(" <no binding>");
}
}
}

View File

@@ -0,0 +1,581 @@
//! Live PMOMusic demo: binds a renderer queue to a dynamic "Live Playlist" container
//! and monitors ContentDirectory updates over an extended period (~30 minutes).
use std::collections::{VecDeque, HashSet};
use std::env;
use std::process;
use std::thread;
use std::time::Duration;
use anyhow::{Context, Result};
use pmocontrol::{
ControlPoint, DeviceRegistryRead, MediaBrowser, MediaEntry, MediaResource, MediaServerEvent,
MediaServerInfo, MusicRenderer, MusicServer, PlaybackItem, PlaybackPosition,
PlaybackPositionInfo, RendererInfo, RendererProtocol,
};
const DEFAULT_TIMEOUT_SECS: u64 = 5;
const DEFAULT_DISCOVERY_SECS: u64 = 5;
const DEFAULT_MAX_INITIAL_TRACKS: usize = 15;
const MONITOR_DURATION_SECS: u64 = 1800; // ~30 minutes
const MONITOR_POLL_SECS: u64 = 5;
const MAX_BROWSE_DEPTH: usize = 4;
const MAX_CONTAINERS_TO_EXPLORE: usize = 100;
const FALLBACK_MIN_TRACKS: usize = 5;
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();
});
println!(
"Starting live_pmomusic_demo with timeout={}s discovery={}s max_initial_tracks={}",
config.timeout_secs, config.discovery_secs, config.max_initial_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<RendererInfo> = 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_pmomusic_server(reg.list_servers())
.unwrap_or_else(|| no_server_and_exit("No PMOMusic media server found."));
(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(), &registry)
.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")?;
println!("Searching for a Live Playlist container in ContentDirectory...");
let live_playlist_container = find_live_playlist_container(&server)
.context("Failed to search for Live Playlist container")?;
let live_playlist_container = match live_playlist_container {
Some(container) => container,
None => {
println!("No Live Playlist container found on server \"{}\". Exiting.", server_info.friendly_name);
process::exit(1);
}
};
println!(
"✓ Found Live Playlist: '{}' (id: {}, class: {})",
live_playlist_container.title, live_playlist_container.id, live_playlist_container.class
);
// Build initial queue from the live playlist container
let playback_items = collect_playable_items_from_container(
&server,
&live_playlist_container.id,
config.max_initial_tracks,
)
.context("Failed to collect playable items from Live Playlist container")?;
if playback_items.is_empty() {
println!("Live Playlist container '{}' contains no playable tracks.", live_playlist_container.title);
process::exit(1);
}
println!(
"Discovered {} playable items from Live Playlist; enqueuing…",
playback_items.len()
);
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 live playlist container
control_point.attach_queue_to_playlist(
&renderer_id,
server_info.id.clone(),
live_playlist_container.id.clone(),
);
println!(
"✓ Queue attached to Live Playlist container '{}' (id: {}) on server '{}'",
live_playlist_container.title, live_playlist_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);
let mut planned_queue: VecDeque<PlaybackItem> = 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 Live Playlist queue for {} seconds (poll every {}s)…",
MONITOR_DURATION_SECS, MONITOR_POLL_SECS
);
println!("This will observe ContentDirectory updates and auto-advance behavior.");
// 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 Live Playlist container '{}' was updated, queue refresh triggered automatically",
bound_container
);
// Take a fresh snapshot to observe changes
if let Ok(fresh_snapshot) = control_point.get_queue_snapshot(&renderer_id) {
println!(
" → Queue length after refresh: {} items",
fresh_snapshot.len()
);
if !fresh_snapshot.is_empty() {
println!(" → First item: {}",
fresh_snapshot[0].title.as_deref().unwrap_or("<no title>")
);
}
}
}
}
}
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<PlaybackItem> = snapshot.clone().into();
// Detect queue changes (auto-advance)
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 ({}s elapsed), exiting.", MONITOR_DURATION_SECS);
Ok(())
}
#[derive(Debug)]
struct CliConfig {
timeout_secs: u64,
discovery_secs: u64,
max_initial_tracks: usize,
}
impl CliConfig {
fn parse_from_env() -> Result<Self, String> {
let mut timeout_secs = DEFAULT_TIMEOUT_SECS;
let mut discovery_secs = DEFAULT_DISCOVERY_SECS;
let mut max_initial_tracks = DEFAULT_MAX_INITIAL_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-initial-tracks" => {
let value = args
.next()
.ok_or_else(|| "--max-initial-tracks requires a value".to_string())?;
max_initial_tracks = value.parse().map_err(|err| {
format!("Invalid value for --max-initial-tracks ({value}): {err}")
})?;
}
"--help" | "-h" => {
print_usage_and_exit();
}
unknown => {
return Err(format!("Unknown argument: {unknown}"));
}
}
}
Ok(Self {
timeout_secs,
discovery_secs,
max_initial_tracks,
})
}
}
fn pick_renderer(renderers: Vec<RendererInfo>) -> Option<RendererInfo> {
let mut candidates: Vec<RendererInfo> = renderers
.into_iter()
.filter(|info| match info.protocol {
RendererProtocol::OpenHomeOnly => false,
RendererProtocol::UpnpAvOnly | RendererProtocol::Hybrid => true,
})
.collect();
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_pmomusic_server(servers: Vec<MediaServerInfo>) -> Option<MediaServerInfo> {
let mut candidates: Vec<MediaServerInfo> = 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;
}
// Prioritize PMOMusic servers
if let Some(idx) = candidates.iter().position(is_pmomusic_server) {
let server = candidates.remove(idx);
println!(
"✓ Selected PMOMusic server \"{}\" (model: {}, manufacturer: {}).",
server.friendly_name, server.model_name, server.manufacturer
);
return Some(server);
}
// Fallback to first ContentDirectory server if no PMOMusic found
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 model = info.model_name.to_ascii_lowercase();
let manufacturer = info.manufacturer.to_ascii_lowercase();
name.contains("pmomusic") || model.contains("pmomusic") || manufacturer.contains("pmomusic")
}
fn is_pmomusic_renderer(info: &RendererInfo) -> bool {
info.friendly_name
.to_ascii_lowercase()
.contains("pmomusic audio renderer")
}
/// Search for a container whose title contains "Live Playlist" using BFS.
fn find_live_playlist_container(server: &MusicServer) -> Result<Option<MediaEntry>> {
let root_entries = server
.browse_root()
.context("Failed to browse ContentDirectory root")?;
println!("Root returned {} entries, starting BFS search for Live Playlist...", root_entries.len());
// BFS queue: (entry, depth)
let mut queue: VecDeque<(MediaEntry, usize)> = VecDeque::new();
let mut visited: HashSet<String> = HashSet::new();
let mut containers_explored = 0;
// Initialize with root entries
for entry in root_entries {
if entry.is_container {
queue.push_back((entry, 0));
}
}
while let Some((container, depth)) = queue.pop_front() {
// Check exploration limits
if depth > MAX_BROWSE_DEPTH {
continue;
}
if containers_explored >= MAX_CONTAINERS_TO_EXPLORE {
println!("Reached max containers to explore ({}), stopping search.", MAX_CONTAINERS_TO_EXPLORE);
break;
}
// Skip already visited
if visited.contains(&container.id) {
continue;
}
visited.insert(container.id.clone());
containers_explored += 1;
let title_lower = container.title.to_ascii_lowercase();
// Check if this container is a Live Playlist
if title_lower.contains("live playlist") {
println!(
"✓ Found Live Playlist at depth {}: '{}' (id: {})",
depth, container.title, container.id
);
return Ok(Some(container));
}
// Browse children and add containers to queue
match server.browse_children(&container.id, 0, 100) {
Ok(children) => {
for child in children {
if child.is_container && !visited.contains(&child.id) {
queue.push_back((child, depth + 1));
}
}
}
Err(err) => {
tracing::warn!(
container_id = container.id.as_str(),
error = %err,
"Failed to browse container during Live Playlist search"
);
}
}
}
println!("No Live Playlist container found after exploring {} containers.", containers_explored);
Ok(None)
}
/// Collect playable items from a specific container.
/// Retries with fewer items if the initial browse times out.
fn collect_playable_items_from_container(
server: &MusicServer,
container_id: &str,
max_tracks: usize,
) -> Result<Vec<PlaybackItem>> {
println!(
"Attempting to browse Live Playlist container (requesting {} items)...",
max_tracks
);
// First attempt with requested count
let children = match server.browse_children(container_id, 0, max_tracks as u32) {
Ok(children) => children,
Err(err) => {
let err_str = err.to_string().to_lowercase();
if err_str.contains("timeout") && max_tracks > FALLBACK_MIN_TRACKS {
println!(
"Browse timed out, retrying with fewer items ({})...",
FALLBACK_MIN_TRACKS
);
// Fallback: try with minimal number of items
server
.browse_children(container_id, 0, FALLBACK_MIN_TRACKS as u32)
.context("Failed to browse Live Playlist container even with minimal item count")?
} else {
return Err(err).context("Failed to browse Live Playlist container children");
}
}
};
println!("Browse returned {} entries from Live Playlist", children.len());
let mut items = Vec::new();
for entry in &children {
if !entry.is_container {
if let Some(item) = playback_item_from_entry(server, entry) {
items.push(item);
if items.len() >= max_tracks {
break;
}
}
}
}
println!("Extracted {} playable items from Live Playlist", items.len());
Ok(items)
}
fn playback_item_from_entry(server: &MusicServer, entry: &MediaEntry) -> Option<PlaybackItem> {
// Skip live streams (we're looking for regular tracks in a live playlist)
if entry.title.to_ascii_lowercase().contains("live stream") {
return None;
}
let resource = entry.resources.iter().find(|res| is_audio_resource(res))?;
let mut item = PlaybackItem::new(resource.uri.clone());
item.title = Some(entry.title.clone());
item.server_id = Some(server.id().clone());
item.object_id = Some(entry.id.clone());
Some(item)
}
fn is_audio_resource(res: &MediaResource) -> bool {
let lower = res.protocol_info.to_ascii_lowercase();
if lower.contains("audio/") {
return true;
}
lower
.split(':')
.nth(2)
.map(|mime| mime.starts_with("audio/"))
.unwrap_or(false)
}
fn print_queue_snapshot(items: &[PlaybackItem]) {
println!("Current queue snapshot ({} items):", items.len());
for (idx, item) in items.iter().take(10).enumerate() {
let label = item.title.as_deref().unwrap_or_else(|| item.uri.as_str());
println!(" [{}] {}", idx, label);
}
if items.len() > 10 {
println!(" ... and {} more items", items.len() - 10);
}
if items.is_empty() {
println!(" <queue is empty>");
}
}
fn current_track_title(item: Option<&PlaybackItem>) -> String {
match item {
Some(track) => track
.title
.as_deref()
.unwrap_or_else(|| track.uri.as_str())
.to_string(),
None => "<unknown>".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 live_pmomusic_demo -- [--timeout-secs N] [--discovery-secs N] [--max-initial-tracks N]"
);
process::exit(1);
}

View File

@@ -9,9 +9,9 @@ use std::time::Duration;
use anyhow::{Context, Result};
use pmocontrol::{
ControlPoint, DeviceRegistryRead, MediaBrowser, MediaEntry, MediaResource, MediaServerInfo,
MusicRenderer, MusicServer, PlaybackItem, PlaybackPosition, PlaybackPositionInfo, RendererInfo,
RendererProtocol,
ControlPoint, DeviceRegistryRead, MediaBrowser, MediaEntry, MediaResource, MediaServerEvent,
MediaServerInfo, MusicRenderer, MusicServer, PlaybackItem, PlaybackPosition,
PlaybackPositionInfo, RendererInfo, RendererProtocol,
};
const DEFAULT_TIMEOUT_SECS: u64 = 5;
@@ -94,8 +94,10 @@ fn main() -> Result<()> {
.context("Failed to browse ContentDirectory root")?;
println!("Root returned {} entries", root_entries.len());
let playback_items = collect_playable_items(&server, &root_entries, config.max_tracks)
.context("Failed to derive playable items from ContentDirectory root/children")?;
// 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.");
@@ -107,6 +109,15 @@ fn main() -> Result<()> {
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<PlaybackItem>;
let renderer_id = renderer.id.clone();
control_point
@@ -116,6 +127,19 @@ fn main() -> Result<()> {
.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(),
);
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")?;
@@ -140,9 +164,51 @@ fn main() -> Result<()> {
"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")?;
@@ -302,11 +368,60 @@ fn is_pmomusic_renderer(info: &RendererInfo) -> bool {
.contains("pmomusic audio renderer")
}
fn collect_playable_items(
fn collect_playable_items_with_binding(
server: &MusicServer,
entries: &[MediaEntry],
max_tracks: usize,
) -> Result<Vec<PlaybackItem>> {
) -> Result<(Vec<PlaybackItem>, Option<String>)> {
// 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)?;
@@ -314,7 +429,7 @@ fn collect_playable_items(
break;
}
}
Ok(items)
Ok((items, None))
}
fn gather_items_from_entry(

View File

@@ -7,7 +7,7 @@ use std::time::Duration;
use anyhow::anyhow;
use crossbeam_channel::Receiver;
use pmoupnp::ssdp::SsdpClient;
use tracing::{debug, error, warn};
use tracing::{debug, error, info, warn};
use crate::MusicRenderer;
use crate::capabilities::{
@@ -16,7 +16,9 @@ use crate::capabilities::{
};
use crate::discovery::DiscoveryManager;
use crate::events::{MediaServerEventBus, RendererEventBus};
use crate::media_server::{MediaServerInfo, ServerId};
use crate::media_server::{
MediaBrowser, MediaEntry, MediaResource, MediaServerInfo, MusicServer, ServerId,
};
use crate::media_server_events::spawn_media_server_event_runtime;
use crate::model::{MediaServerEvent, RendererEvent, RendererId, RendererProtocol};
use crate::music_renderer::op_not_supported;
@@ -25,6 +27,25 @@ use crate::provider::HttpXmlDescriptionProvider;
use crate::registry::{DeviceRegistry, DeviceRegistryRead, DeviceUpdate};
use crate::upnp_renderer::UpnpRenderer;
/// Optional attachment between a renderer playback queue and a server-side
/// DIDL-Lite playlist container.
///
/// When a queue is bound to a playlist, the control point will automatically
/// refresh it whenever the server notifies us of changes to that container.
/// User-driven mutations (clear, enqueue, etc.) break the binding automatically.
#[derive(Clone, Debug)]
struct PlaylistBinding {
/// MediaServer that owns the playlist container.
server_id: ServerId,
/// DIDL-Lite object id of the playlist container.
container_id: String,
/// True once at least one ContainerUpdateIDs notification has been seen.
has_seen_update: bool,
/// Flag used internally to signal that the queue should be refreshed
/// from the server container.
pending_refresh: bool,
}
/// Control point minimal :
/// - lance un SsdpClient dans un thread,
/// - passe les SsdpEvent au DiscoveryManager,
@@ -34,6 +55,12 @@ pub struct ControlPoint {
event_bus: RendererEventBus,
media_event_bus: MediaServerEventBus,
runtime: Arc<RuntimeState>,
/// Optional attachment between a renderer playback queue and a
/// server-side DIDL-Lite playlist container.
///
/// Key : RendererId
/// Value : PlaylistBinding
playlist_bindings: Arc<Mutex<HashMap<RendererId, PlaylistBinding>>>,
}
impl ControlPoint {
@@ -45,6 +72,7 @@ impl ControlPoint {
let event_bus = RendererEventBus::new();
let media_event_bus = MediaServerEventBus::new();
let runtime = Arc::new(RuntimeState::new());
let playlist_bindings = Arc::new(Mutex::new(HashMap::new()));
// SsdpClient
let client = SsdpClient::new()?; // pmoupnp::ssdp::SsdpClient
@@ -61,10 +89,11 @@ impl ControlPoint {
// ACTIVE DISCOVERY : envoyer quelques M-SEARCH au démarrage
// pour forcer les devices à répondre rapidement.
let search_targets = [
"ssdp:all",
"urn:schemas-upnp-org:device:MediaRenderer:1",
"urn:av-openhome-org:device:MediaRenderer:1",
"urn:schemas-upnp-org:device:MediaServer:1",
"ssdp:all",
"urn:schemas-upnp-org:device:MediaRenderer:1",
"urn:av-openhome-org:device:MediaRenderer:1",
"urn:schemas-upnp-org:device:MediaServer:1",
"urn:schemas-wiimu-com:service:PlayQueue:1", // <-- AJOUTER
];
for st in &search_targets {
@@ -96,9 +125,11 @@ impl ControlPoint {
event_bus: event_bus.clone(),
media_event_bus: media_event_bus.clone(),
runtime: Arc::clone(&runtime),
playlist_bindings: Arc::clone(&playlist_bindings),
};
thread::spawn(move || {
let mut tick: u32 = 0;
loop {
let infos = {
let reg = runtime_cp.registry.read().unwrap();
@@ -126,6 +157,7 @@ impl ControlPoint {
let mut new_snapshot = prev_snapshot.clone();
let prev_position = prev_snapshot.position.clone();
// Poll position every tick (1s) for smooth UI progress
if let Ok(position) = renderer.playback_position() {
let has_changed = match prev_snapshot.position.as_ref() {
Some(prev) => !playback_position_equal(prev, &position),
@@ -142,6 +174,7 @@ impl ControlPoint {
new_snapshot.position = Some(position);
}
// Poll state every tick to ensure responsive playback control
if let Ok(raw_state) = renderer.playback_state() {
let logical_state = compute_logical_playback_state(
&raw_state,
@@ -154,7 +187,9 @@ impl ControlPoint {
None => true,
};
if has_changed {
// Emit event only for non-transient states to reduce noise
// and avoid overwhelming the renderer during track changes
if has_changed && !matches!(logical_state, PlaybackState::Transitioning) {
runtime_cp.emit_renderer_event(RendererEvent::StateChanged {
id: renderer_id.clone(),
state: logical_state.clone(),
@@ -164,26 +199,30 @@ impl ControlPoint {
new_snapshot.state = Some(logical_state);
}
if let Ok(volume) = renderer.volume() {
if prev_snapshot.last_volume != Some(volume) {
runtime_cp.emit_renderer_event(RendererEvent::VolumeChanged {
id: renderer_id.clone(),
volume,
});
// Poll volume and mute less frequently (every 3 seconds)
// to reduce SOAP overhead without impacting UI responsiveness
if tick % 3 == 0 {
if let Ok(volume) = renderer.volume() {
if prev_snapshot.last_volume != Some(volume) {
runtime_cp.emit_renderer_event(RendererEvent::VolumeChanged {
id: renderer_id.clone(),
volume,
});
}
new_snapshot.last_volume = Some(volume);
}
new_snapshot.last_volume = Some(volume);
}
if let Ok(mute) = renderer.mute() {
if prev_snapshot.last_mute != Some(mute) {
runtime_cp.emit_renderer_event(RendererEvent::MuteChanged {
id: renderer_id.clone(),
mute,
});
}
if let Ok(mute) = renderer.mute() {
if prev_snapshot.last_mute != Some(mute) {
runtime_cp.emit_renderer_event(RendererEvent::MuteChanged {
id: renderer_id.clone(),
mute,
});
new_snapshot.last_mute = Some(mute);
}
new_snapshot.last_mute = Some(mute);
}
runtime_cp
@@ -191,6 +230,8 @@ impl ControlPoint {
.update_snapshot(&renderer_id, new_snapshot);
}
tick = tick.wrapping_add(1);
// Keep 1 second polling for smooth position updates
thread::sleep(Duration::from_secs(1));
}
});
@@ -201,11 +242,138 @@ impl ControlPoint {
timeout_secs,
)?;
// Worker thread to process MediaServerEvent and trigger queue refreshes
// for renderers bound to updated playlist containers
let registry_for_media_worker = Arc::clone(&registry);
let runtime_for_media_worker = Arc::clone(&runtime);
let bindings_for_media_worker = Arc::clone(&playlist_bindings);
let media_rx = media_event_bus.subscribe();
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;
}
};
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<RendererId> = {
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());
}
}
to_refresh
};
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(
&registry_for_media_worker,
&runtime_for_media_worker,
&bindings_for_media_worker,
&renderer_id,
) {
warn!(
renderer = renderer_id.0.as_str(),
server = server_id.0.as_str(),
error = %err,
"Failed to refresh queue from playlist container"
);
}
}
}
}
}
})?;
// Periodic refresh worker for bound playlists
// Every 60 seconds, trigger a refresh for all renderers with active bindings
let registry_for_periodic = Arc::clone(&registry);
let runtime_for_periodic = Arc::clone(&runtime);
let bindings_for_periodic = Arc::clone(&playlist_bindings);
thread::Builder::new()
.name("cp-playlist-periodic-refresh".into())
.spawn(move || {
loop {
// Sleep for 60 seconds between refresh cycles
thread::sleep(Duration::from_secs(60));
// Collect all renderers with active bindings and mark them for refresh
let renderers_to_refresh: Vec<RendererId> = {
let mut bindings = bindings_for_periodic.lock().unwrap();
let mut to_refresh = Vec::new();
for (renderer_id, binding) in bindings.iter_mut() {
binding.pending_refresh = true;
to_refresh.push(renderer_id.clone());
}
to_refresh
};
// Trigger refresh for each bound renderer (outside of lock)
for renderer_id in renderers_to_refresh {
debug!(
renderer = renderer_id.0.as_str(),
"Periodic refresh triggered for bound playlist"
);
if let Err(err) = refresh_attached_queue_for(
&registry_for_periodic,
&runtime_for_periodic,
&bindings_for_periodic,
&renderer_id,
) {
warn!(
renderer = renderer_id.0.as_str(),
error = %err,
"Periodic refresh failed for bound playlist"
);
}
}
}
})?;
Ok(Self {
registry,
event_bus,
media_event_bus,
runtime,
playlist_bindings,
})
}
@@ -308,6 +476,9 @@ impl ControlPoint {
return Err(err);
}
// User-driven mutation: detach any playlist binding
self.detach_binding_on_user_mutation(renderer_id, "clear_queue");
let removed = self
.runtime
.with_queue_mut(renderer_id, |queue| {
@@ -340,6 +511,9 @@ impl ControlPoint {
return Err(err);
}
// User-driven mutation: detach any playlist binding
self.detach_binding_on_user_mutation(renderer_id, "enqueue_items");
let item_count = items.len();
let new_len = self
.runtime
@@ -501,6 +675,92 @@ impl ControlPoint {
self.media_event_bus.subscribe()
}
/// Attach a renderer's playback queue to a server-side playlist container.
///
/// When attached, the queue will be automatically refreshed from the
/// container whenever the server notifies us of changes via ContentDirectory
/// events. The binding is broken if the user explicitly mutates the queue
/// through methods like `clear_queue` or `enqueue_items`.
pub fn attach_queue_to_playlist(
&self,
renderer_id: &RendererId,
server_id: ServerId,
container_id: String,
) {
let mut bindings = self.playlist_bindings.lock().unwrap();
bindings.insert(
renderer_id.clone(),
PlaylistBinding {
server_id: server_id.clone(),
container_id: container_id.clone(),
has_seen_update: false,
pending_refresh: false,
},
);
info!(
renderer = renderer_id.0.as_str(),
server = server_id.0.as_str(),
container = container_id.as_str(),
"Queue attached to playlist container"
);
}
/// Detach a renderer's queue from its associated playlist container.
///
/// After calling this, the queue will no longer be automatically refreshed
/// from the server. If no binding existed, this is a no-op.
pub fn detach_queue_playlist(&self, renderer_id: &RendererId) {
let mut bindings = self.playlist_bindings.lock().unwrap();
if let Some(binding) = bindings.remove(renderer_id) {
info!(
renderer = renderer_id.0.as_str(),
server = binding.server_id.0.as_str(),
container = binding.container_id.as_str(),
"Queue detached from playlist container"
);
} else {
debug!(
renderer = renderer_id.0.as_str(),
"detach_queue_playlist: no binding to remove"
);
}
}
/// Query the current playlist binding for a renderer's queue, if any.
///
/// Returns `(server_id, container_id, has_seen_update)` if the queue is
/// bound to a server playlist container, or `None` otherwise.
pub fn current_queue_playlist_binding(
&self,
renderer_id: &RendererId,
) -> Option<(ServerId, String, bool)> {
let bindings = self.playlist_bindings.lock().unwrap();
bindings.get(renderer_id).map(|binding| {
(
binding.server_id.clone(),
binding.container_id.clone(),
binding.has_seen_update,
)
})
}
/// Internal helper to detach the playlist binding on user-driven mutations.
///
/// This is called by public queue mutation methods (clear, enqueue, etc.)
/// to ensure that explicit user actions break the automatic refresh binding.
fn detach_binding_on_user_mutation(&self, renderer_id: &RendererId, reason: &str) {
let mut bindings = self.playlist_bindings.lock().unwrap();
if let Some(binding) = bindings.remove(renderer_id) {
info!(
renderer = renderer_id.0.as_str(),
server = binding.server_id.0.as_str(),
container = binding.container_id.as_str(),
reason = reason,
"Playlist binding auto-detached due to user mutation"
);
}
}
pub(crate) fn emit_renderer_event(&self, event: RendererEvent) {
self.handle_renderer_event(&event);
self.event_bus.broadcast(event);
@@ -657,6 +917,202 @@ impl RuntimeState {
}
}
/// Internal helper to refresh a renderer's playback queue from its bound
/// playlist container.
///
/// This function is called automatically when a ContentDirectory event indicates
/// that the bound container has been updated. It attempts to preserve the
/// currently playing item when possible.
fn refresh_attached_queue_for(
registry: &Arc<RwLock<DeviceRegistry>>,
runtime: &Arc<RuntimeState>,
bindings: &Arc<Mutex<HashMap<RendererId, PlaylistBinding>>>,
renderer_id: &RendererId,
) -> anyhow::Result<()> {
// Step 1: Check binding and mark refresh as in-progress
let (server_id, container_id) = {
let mut bindings_lock = bindings.lock().unwrap();
let binding = match bindings_lock.get_mut(renderer_id) {
Some(b) => b,
None => {
debug!(
renderer = renderer_id.0.as_str(),
"refresh_attached_queue_for: no binding present"
);
return Ok(());
}
};
if !binding.pending_refresh {
debug!(
renderer = renderer_id.0.as_str(),
"refresh_attached_queue_for: pending_refresh is false, nothing to do"
);
return Ok(());
}
// Mark as processed
binding.pending_refresh = false;
(binding.server_id.clone(), binding.container_id.clone())
};
// Step 2: Fetch MediaServerInfo from registry
let server_info = {
let reg = registry.read().unwrap();
reg.get_server(&server_id)
};
let server_info = match server_info {
Some(info) => info,
None => {
warn!(
renderer = renderer_id.0.as_str(),
server = server_id.0.as_str(),
"refresh_attached_queue_for: server not found in registry"
);
return Ok(());
}
};
if !server_info.online {
debug!(
renderer = renderer_id.0.as_str(),
server = server_id.0.as_str(),
"refresh_attached_queue_for: server offline, skipping refresh"
);
return Ok(());
}
if !server_info.has_content_directory {
debug!(
renderer = renderer_id.0.as_str(),
server = server_id.0.as_str(),
"refresh_attached_queue_for: server has no ContentDirectory"
);
return Ok(());
}
// 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);
}
};
// Step 4: Convert MediaEntry to PlaybackItem
let new_items: Vec<PlaybackItem> = entries
.iter()
.filter_map(|entry| playback_item_from_entry(&music_server, entry))
.collect();
if new_items.is_empty() {
debug!(
renderer = renderer_id.0.as_str(),
server = server_id.0.as_str(),
container = container_id.as_str(),
"Refreshed playlist is empty, clearing queue"
);
runtime.with_queue_mut(renderer_id, |queue| queue.clear());
return Ok(());
}
// Step 5: Intelligent refresh: try to keep current item if it's still in the new list
let current_item = runtime.queue_snapshot(renderer_id).and_then(|q| q.first().cloned());
let item_found_at = current_item.as_ref().and_then(|current| {
new_items.iter().position(|new_item| {
// Match by object_id if both have it
if let (Some(current_obj), Some(new_obj)) = (&current.object_id, &new_item.object_id)
{
return current_obj == new_obj;
}
// Fallback: match by URI
current.uri == new_item.uri
})
});
runtime.with_queue_mut(renderer_id, |queue| {
queue.clear();
if let Some(idx) = item_found_at {
// Current item found: reconstruct queue starting from that item
for i in idx..new_items.len() {
queue.enqueue(new_items[i].clone());
}
info!(
renderer = renderer_id.0.as_str(),
server = server_id.0.as_str(),
container = container_id.as_str(),
queue_len = new_items.len() - idx,
current_preserved = true,
"Refreshed queue from playlist container"
);
} else {
// Current item not found: replace with full new list
for item in new_items.iter() {
queue.enqueue(item.clone());
}
info!(
renderer = renderer_id.0.as_str(),
server = server_id.0.as_str(),
container = container_id.as_str(),
queue_len = new_items.len(),
current_preserved = false,
"Refreshed queue from playlist container (current item not found)"
);
}
});
Ok(())
}
/// Helper to detect if a MediaResource is audio content.
fn is_audio_resource(res: &MediaResource) -> bool {
let lower = res.protocol_info.to_ascii_lowercase();
if lower.contains("audio/") {
return true;
}
// Check MIME type in protocolInfo (format: protocol:network:contentFormat:additionalInfo)
lower
.split(':')
.nth(2)
.map(|mime| mime.starts_with("audio/"))
.unwrap_or(false)
}
/// Helper to convert a MediaEntry to a PlaybackItem.
fn playback_item_from_entry(server: &MusicServer, entry: &MediaEntry) -> Option<PlaybackItem> {
// Ignore containers
if entry.is_container {
return None;
}
// Skip "live stream" entries (heuristic from example)
if entry.title.to_ascii_lowercase().contains("live stream") {
return None;
}
// Find an audio resource
let resource = entry.resources.iter().find(|res| is_audio_resource(res))?;
let mut item = PlaybackItem::new(resource.uri.clone());
item.title = Some(entry.title.clone());
item.server_id = Some(server.id().clone());
item.object_id = Some(entry.id.clone());
Some(item)
}
/// Parse "HH:MM:SS" style time strings to seconds.
///
/// Returns None for empty or sentinel values such as "NOT_IMPLEMENTED" or "-:--:--".

17811
pmocontrol/ssdp_capture.txt Normal file

File diff suppressed because it is too large Load Diff

View File

@@ -25,7 +25,7 @@ use std::collections::HashMap;
use std::net::{SocketAddr, UdpSocket};
use std::sync::Arc;
use std::time::Duration;
use tracing::{debug, info, warn};
use tracing::{debug, info, trace, warn};
/// Événements SSDP intéressants pour un control point
#[derive(Debug, Clone)]
@@ -149,10 +149,7 @@ impl SsdpClient {
Ok((n, from)) => {
let data = String::from_utf8_lossy(&buf[..n]);
if let Some(event) = parse_message(&data, from) {
debug!(
"📥 SSDP datagram from {}\n<details>\n\n```\n{}\n```\n</details>\n",
from, data
);
debug!("📥 SSDP event from {}: {:?}", from, event);
on_event(event);
}
}
@@ -174,7 +171,7 @@ fn parse_message(data: &str, from: SocketAddr) -> Option<SsdpEvent> {
let upper = first_line.to_ascii_uppercase();
let headers = parse_headers(lines);
if upper.starts_with("NOTIFY ") {
let result = if upper.starts_with("NOTIFY ") {
handle_notify(&headers, from)
} else if upper.starts_with("HTTP/") && upper.contains(" 200 ") {
handle_search_response(&headers, from)
@@ -182,19 +179,50 @@ fn parse_message(data: &str, from: SocketAddr) -> Option<SsdpEvent> {
// Another control point querying us; we are not a device, so we ignore.
None
} else {
trace!("Unknown SSDP message type from {}: {}", from, first_line);
None
};
if result.is_none() {
trace!(
"SSDP message from {} could not be parsed:\n{}",
from,
data
);
}
result
}
fn handle_notify(headers: &HashMap<String, String>, from: SocketAddr) -> Option<SsdpEvent> {
// Critical headers: NTS, NT, USN (required by UPnP spec)
let nts = headers.get("NTS")?.to_ascii_lowercase();
let nt = headers.get("NT")?.to_string();
let usn = headers.get("USN")?.to_string();
if nts == "ssdp:alive" {
let location = headers.get("LOCATION")?.to_string();
let server = headers.get("SERVER")?.to_string();
// LOCATION is required for alive notifications
let location = match headers.get("LOCATION") {
Some(loc) => loc.to_string(),
None => {
trace!(
"NOTIFY ssdp:alive from {} missing LOCATION header, ignoring",
from
);
return None;
}
};
// Non-critical headers: SERVER, CACHE-CONTROL
let server = headers
.get("SERVER")
.map(|s| s.to_string())
.unwrap_or_else(|| {
trace!("NOTIFY from {} has no SERVER header, using 'Unknown'", from);
"Unknown".to_string()
});
let max_age = parse_max_age(headers.get("CACHE-CONTROL"));
Some(SsdpEvent::Alive {
usn,
nt,
@@ -206,6 +234,7 @@ fn handle_notify(headers: &HashMap<String, String>, from: SocketAddr) -> Option<
} else if nts == "ssdp:byebye" {
Some(SsdpEvent::ByeBye { usn, nt, from })
} else {
trace!("Unknown NTS value from {}: {}", from, nts);
None
}
}
@@ -214,10 +243,46 @@ fn handle_search_response(
headers: &HashMap<String, String>,
from: SocketAddr,
) -> Option<SsdpEvent> {
let st = headers.get("ST")?.to_string();
let usn = headers.get("USN")?.to_string();
let location = headers.get("LOCATION")?.to_string();
let server = headers.get("SERVER")?.to_string();
// Critical headers: ST, USN, LOCATION (required by UPnP spec)
let st = match headers.get("ST") {
Some(s) => s.to_string(),
None => {
trace!("M-SEARCH response from {} missing ST header, ignoring", from);
return None;
}
};
let usn = match headers.get("USN") {
Some(u) => u.to_string(),
None => {
trace!(
"M-SEARCH response from {} missing USN header, ignoring",
from
);
return None;
}
};
let location = match headers.get("LOCATION") {
Some(loc) => loc.to_string(),
None => {
trace!(
"M-SEARCH response from {} missing LOCATION header, ignoring",
from
);
return None;
}
};
// Non-critical headers: SERVER, CACHE-CONTROL
let server = headers
.get("SERVER")
.map(|s| s.to_string())
.unwrap_or_else(|| {
trace!(
"M-SEARCH response from {} has no SERVER header, using 'Unknown'",
from
);
"Unknown".to_string()
});
let max_age = parse_max_age(headers.get("CACHE-CONTROL"));
Some(SsdpEvent::SearchResponse {
@@ -237,11 +302,29 @@ where
let mut headers = HashMap::new();
for line in lines {
let line = line.trim();
// Empty line marks end of headers
if line.is_empty() {
break;
}
if let Some((name, value)) = line.split_once(':') {
headers.insert(name.trim().to_ascii_uppercase(), value.trim().to_string());
// Split on first ':' only (values may contain ':')
if let Some(colon_pos) = line.find(':') {
let (name, value_with_colon) = line.split_at(colon_pos);
let value = &value_with_colon[1..]; // Skip the ':'
let name = name.trim().to_ascii_uppercase();
let value = value.trim().to_string();
// Skip empty header names or values
if !name.is_empty() && !value.is_empty() {
headers.insert(name, value);
} else {
trace!("Skipping malformed header: '{}'", line);
}
} else {
// Line without ':' is invalid, skip it
trace!("Skipping line without colon: '{}'", line);
}
}
headers
@@ -249,22 +332,27 @@ where
fn parse_max_age(value: Option<&String>) -> u32 {
if let Some(v) = value {
for part in v.split(',') {
let part = part.trim();
if let Some(rest) = part.strip_prefix("max-age=") {
if let Ok(age) = rest.trim().parse::<u32>() {
return age;
}
} else {
let lower = part.to_ascii_lowercase();
if let Some(idx) = lower.find("max-age=") {
let raw = &part[idx + 8..];
if let Ok(age) = raw.trim().parse::<u32>() {
return age;
}
}
// Try case-insensitive search for "max-age=<number>"
let lower = v.to_ascii_lowercase();
if let Some(idx) = lower.find("max-age") {
// Extract everything after "max-age"
let after_key = &v[idx + 7..];
// Skip any whitespace and '=' characters
let after_eq = after_key.trim_start().trim_start_matches('=').trim_start();
// Try to parse the first sequence of digits
let digits: String = after_eq
.chars()
.take_while(|c| c.is_ascii_digit())
.collect();
if let Ok(age) = digits.parse::<u32>() {
return age;
}
}
trace!(
"Could not parse max-age from CACHE-CONTROL: '{}', using default {}",
v,
MAX_AGE
);
}
MAX_AGE
}

0
src/index.js Normal file
View File