On fait la même chose maintenant coté serveur de musique.

This commit is contained in:
2025-12-01 19:55:55 +01:00
parent 7b1d35cebc
commit f434c521e8
14 changed files with 968 additions and 114 deletions

1
Cargo.lock generated
View File

@@ -3164,6 +3164,7 @@ version = "0.1.0"
dependencies = [
"anyhow",
"crossbeam-channel",
"pmodidl",
"pmoupnp",
"quick-xml 0.38.4",
"thiserror 2.0.17",

View File

@@ -5,6 +5,7 @@ edition = "2024"
[dependencies]
pmoupnp = { path = "../pmoupnp" }
pmodidl = { path = "../pmodidl" }
quick-xml = "0.38.4"
thiserror = "2.0.17"
ureq = "3.1.4"

View File

@@ -229,7 +229,9 @@ fn print_backend(prefix: &str, renderer: &MusicRenderer) {
MusicRenderer::Upnp(_) => "UpnpRenderer (UPnP AV / DLNA)",
MusicRenderer::LinkPlay(_) => "LinkPlayRenderer (LinkPlay HTTP)",
MusicRenderer::ArylicTcp(_) => "ArylicTcpRenderer (ARylic TCP Protocol)",
MusicRenderer::HybridUpnpArylic{..} => "Hybrid UpnpArylicRenderer (UPnP AV / DLNA + ARylic TCP Protocol)",
MusicRenderer::HybridUpnpArylic { .. } => {
"Hybrid UpnpArylicRenderer (UPnP AV / DLNA + ARylic TCP Protocol)"
}
};
println!("{prefix}Backend : {backend}");
}

View File

@@ -24,7 +24,6 @@ fn last_command_time() -> &'static Mutex<Instant> {
LAST_COMMAND_TIME.get_or_init(|| Mutex::new(Instant::now()))
}
const ARYLIC_TCP_PORT: u16 = 8899;
const PACKET_HEADER: [u8; 4] = [0x18, 0x96, 0x18, 0x20];
const RESERVED_BYTES: [u8; 8] = [0; 8];
@@ -387,7 +386,10 @@ fn connect(host: &str, port: u16, timeout: Duration) -> Result<TcpStream> {
let elapsed = last_time.elapsed();
if elapsed < Duration::from_millis(200) {
let wait = Duration::from_millis(200) - elapsed;
debug!("Waiting {:?} before sending command to respect 200ms interval", wait);
debug!(
"Waiting {:?} before sending command to respect 200ms interval",
wait
);
thread::sleep(wait);
}
*last_time = Instant::now();
@@ -521,9 +523,12 @@ fn send_command_with_mode(
let mut stream = connect(host, port, timeout)?;
let packet = encode_packet(payload);
stream
.write_all(&packet)
.with_context(|| format!("Failed to write Arylic TCP packet for {}: {}", host, payload))?;
stream.write_all(&packet).with_context(|| {
format!(
"Failed to write Arylic TCP packet for {}: {}",
host, payload
)
})?;
stream.flush().with_context(|| {
format!(
"Failed to flush Arylic TCP stream for {} (command {})",
@@ -540,8 +545,9 @@ fn send_command_with_mode(
let _ = stream.shutdown(Shutdown::Write);
Ok(None)
}
ResponseMode::Required(expected) => read_expected_response(&mut stream, host, payload, expected)
.map(Some),
ResponseMode::Required(expected) => {
read_expected_response(&mut stream, host, payload, expected).map(Some)
}
ResponseMode::Optional(expected) => {
for _ in 0..MAX_RESPONSE_ATTEMPTS {
match read_packet(&mut stream) {
@@ -615,7 +621,13 @@ fn send_command_required(
payload: &str,
expected: &[&str],
) -> Result<String> {
match send_command_with_mode(host, port, timeout, payload, ResponseMode::Required(expected))? {
match send_command_with_mode(
host,
port,
timeout,
payload,
ResponseMode::Required(expected),
)? {
Some(s) => Ok(s),
None => Err(anyhow!(
"Arylic TCP: no response payload for required command {}",
@@ -631,14 +643,15 @@ fn send_command_optional(
payload: &str,
expected: &[&str],
) -> Result<Option<String>> {
send_command_with_mode(host, port, timeout, payload, ResponseMode::Optional(expected))
send_command_with_mode(
host,
port,
timeout,
payload,
ResponseMode::Optional(expected),
)
}
fn send_command_no_response(
host: &str,
port: u16,
timeout: Duration,
payload: &str,
) -> Result<()> {
fn send_command_no_response(host: &str, port: u16, timeout: Duration, payload: &str) -> Result<()> {
send_command_with_mode(host, port, timeout, payload, ResponseMode::None).map(|_| ())
}

View File

@@ -1,19 +1,25 @@
use std::collections::{HashMap, HashSet};
use std::collections::HashMap;
use std::io;
use std::sync::{Arc, RwLock};
use std::sync::{Arc, Mutex, RwLock};
use std::thread;
use std::time::Duration;
use anyhow::anyhow;
use crossbeam_channel::Receiver;
use pmoupnp::ssdp::SsdpClient;
use tracing::{debug, error, warn};
use crate::MusicRenderer;
use crate::capabilities::{
PlaybackPosition, PlaybackPositionInfo, PlaybackState, PlaybackStatus, VolumeControl,
PlaybackPosition, PlaybackPositionInfo, PlaybackState, PlaybackStatus, TransportControl,
VolumeControl,
};
use crate::discovery::DiscoveryManager;
use crate::events::RendererEventBus;
use crate::media_server::{MediaServerInfo, ServerId};
use crate::model::{RendererEvent, RendererId, RendererProtocol};
use crate::music_renderer::op_not_supported;
use crate::playback_queue::{PlaybackItem, PlaybackQueue};
use crate::provider::HttpXmlDescriptionProvider;
use crate::registry::{DeviceRegistry, DeviceRegistryRead, DeviceUpdate};
use crate::upnp_renderer::UpnpRenderer;
@@ -25,6 +31,7 @@ use crate::upnp_renderer::UpnpRenderer;
pub struct ControlPoint {
registry: Arc<RwLock<DeviceRegistry>>,
event_bus: RendererEventBus,
runtime: Arc<RuntimeState>,
}
impl ControlPoint {
@@ -34,6 +41,7 @@ impl ControlPoint {
pub fn spawn(timeout_secs: u64) -> io::Result<Self> {
let registry = Arc::new(RwLock::new(DeviceRegistry::new()));
let event_bus = RendererEventBus::new();
let runtime = Arc::new(RuntimeState::new());
// SsdpClient
let client = SsdpClient::new()?; // pmoupnp::ssdp::SsdpClient
@@ -83,11 +91,10 @@ impl ControlPoint {
let runtime_cp = ControlPoint {
registry: Arc::clone(&registry),
event_bus: event_bus.clone(),
runtime: Arc::clone(&runtime),
};
thread::spawn(move || {
let mut cache: HashMap<RendererId, RendererRuntimeSnapshot> = HashMap::new();
loop {
let renderers = {
let reg = runtime_cp.registry.read().unwrap();
@@ -97,12 +104,9 @@ impl ControlPoint {
.collect::<Vec<_>>()
};
let mut seen_ids = HashSet::new();
for renderer in renderers {
let info = renderer.info();
// Ne pas poller les renderers offline
if !info.online {
continue;
}
@@ -113,20 +117,12 @@ impl ControlPoint {
}
let renderer_id = info.id.clone();
seen_ids.insert(renderer_id.clone());
let prev_snapshot = runtime_cp.runtime.snapshot_for(&renderer_id);
let mut new_snapshot = prev_snapshot.clone();
let prev_position = prev_snapshot.position.clone();
let entry = cache
.entry(renderer_id.clone())
.or_insert_with(RendererRuntimeSnapshot::default);
// Keep a snapshot of the previous position to compute logical
// state transitions based on time deltas.
let prev_position = entry.position.clone();
// 1) Poll position first, so that the state logic can use the
// freshly updated position when available.
if let Ok(position) = renderer.playback_position() {
let has_changed = match entry.position.as_ref() {
let has_changed = match prev_snapshot.position.as_ref() {
Some(prev) => !playback_position_equal(prev, &position),
None => true,
};
@@ -138,19 +134,17 @@ impl ControlPoint {
});
}
entry.position = Some(position);
new_snapshot.position = Some(position);
}
// 2) Poll raw playback state and compute a logical state that
// compensates for buggy devices (Arylic / LinkPlay).
if let Ok(raw_state) = renderer.playback_state() {
let logical_state = compute_logical_playback_state(
&raw_state,
prev_position.as_ref(),
entry.position.as_ref(),
new_snapshot.position.as_ref(),
);
let has_changed = match entry.state.as_ref() {
let has_changed = match prev_snapshot.state.as_ref() {
Some(prev) => !playback_state_equal(prev, &logical_state),
None => true,
};
@@ -160,32 +154,38 @@ impl ControlPoint {
id: renderer_id.clone(),
state: logical_state.clone(),
});
entry.state = Some(logical_state);
}
new_snapshot.state = Some(logical_state);
}
if let Ok(volume) = renderer.volume() {
if entry.last_volume != Some(volume) {
if prev_snapshot.last_volume != Some(volume) {
runtime_cp.emit_renderer_event(RendererEvent::VolumeChanged {
id: renderer_id.clone(),
volume,
});
entry.last_volume = Some(volume);
}
new_snapshot.last_volume = Some(volume);
}
if let Ok(mute) = renderer.mute() {
if entry.last_mute != Some(mute) {
if prev_snapshot.last_mute != Some(mute) {
runtime_cp.emit_renderer_event(RendererEvent::MuteChanged {
id: renderer_id.clone(),
mute,
});
entry.last_mute = Some(mute);
}
new_snapshot.last_mute = Some(mute);
}
runtime_cp
.runtime
.update_snapshot(&renderer_id, new_snapshot);
}
cache.retain(|id, _| seen_ids.contains(id));
thread::sleep(Duration::from_secs(1));
}
});
@@ -193,6 +193,7 @@ impl ControlPoint {
Ok(Self {
registry,
event_bus,
runtime,
})
}
@@ -254,6 +255,184 @@ impl ControlPoint {
.and_then(|info| MusicRenderer::from_registry_info(info, &reg))
}
/// Snapshot list of media servers currently known by the registry.
pub fn list_media_servers(&self) -> Vec<MediaServerInfo> {
let reg = self.registry.read().unwrap();
reg.list_servers()
}
/// Lookup a media server by id.
pub fn media_server(&self, id: &ServerId) -> Option<MediaServerInfo> {
let reg = self.registry.read().unwrap();
reg.get_server(id)
}
pub fn clear_queue(&self, renderer_id: &RendererId) -> anyhow::Result<()> {
if !self.runtime.has_entry(renderer_id) {
let err = Self::runtime_entry_missing(renderer_id);
warn!(
renderer = renderer_id.0.as_str(),
"Cannot clear queue: renderer not registered in runtime"
);
return Err(err);
}
let removed = self
.runtime
.with_queue_mut(renderer_id, |queue| {
let removed = queue.len();
queue.clear();
removed
})
.ok_or_else(|| Self::runtime_entry_missing(renderer_id))?;
debug!(
renderer = renderer_id.0.as_str(),
items_removed = removed,
queue_len = 0,
"Cleared playback queue"
);
Ok(())
}
pub fn enqueue_items(
&self,
renderer_id: &RendererId,
items: Vec<PlaybackItem>,
) -> anyhow::Result<()> {
if !self.runtime.has_entry(renderer_id) {
let err = Self::runtime_entry_missing(renderer_id);
warn!(
renderer = renderer_id.0.as_str(),
"Cannot enqueue items: renderer not registered in runtime"
);
return Err(err);
}
let item_count = items.len();
let new_len = self
.runtime
.with_queue_mut(renderer_id, |queue| {
queue.enqueue_many(items);
queue.len()
})
.ok_or_else(|| Self::runtime_entry_missing(renderer_id))?;
debug!(
renderer = renderer_id.0.as_str(),
added = item_count,
queue_len = new_len,
"Enqueued playback items"
);
Ok(())
}
pub fn get_queue_snapshot(
&self,
renderer_id: &RendererId,
) -> anyhow::Result<Vec<PlaybackItem>> {
if !self.runtime.has_entry(renderer_id) {
let err = Self::runtime_entry_missing(renderer_id);
warn!(
renderer = renderer_id.0.as_str(),
"Cannot snapshot queue: renderer not registered in runtime"
);
return Err(err);
}
self.runtime
.queue_snapshot(renderer_id)
.ok_or_else(|| Self::runtime_entry_missing(renderer_id))
}
pub fn play_next_from_queue(&self, renderer_id: &RendererId) -> anyhow::Result<()> {
if !self.runtime.has_entry(renderer_id) {
let err = Self::runtime_entry_missing(renderer_id);
warn!(
renderer = renderer_id.0.as_str(),
"Cannot advance queue: renderer not registered in runtime"
);
return Err(err);
}
let Some((item, remaining_after)) = self.runtime.dequeue_next(renderer_id) else {
debug!(
renderer = renderer_id.0.as_str(),
"play_next_from_queue: queue is empty"
);
self.runtime
.set_playback_source(renderer_id, PlaybackSource::None);
return Ok(());
};
let queue_before = remaining_after + 1;
debug!(
renderer = renderer_id.0.as_str(),
queue_before,
queue_after = remaining_after,
uri = item.uri.as_str(),
"Dequeued next playback item"
);
let renderer = self
.music_renderer_by_id(renderer_id)
.ok_or_else(|| {
warn!(
renderer = renderer_id.0.as_str(),
"Renderer disappeared before queue playback could start"
);
anyhow!("Renderer {} not found", renderer_id.0)
})?;
if matches!(renderer.info().protocol, RendererProtocol::OpenHomeOnly) {
self.runtime
.with_queue_mut(renderer_id, |queue| queue.enqueue_front(item))
.ok_or_else(|| Self::runtime_entry_missing(renderer_id))?;
self.runtime
.set_playback_source(renderer_id, PlaybackSource::None);
return Err(op_not_supported(
"play_next_from_queue",
"OpenHomeOnly renderer",
));
}
let playback = (|| -> anyhow::Result<()> {
renderer.play_uri(&item.uri, "")?;
renderer.play()?;
Ok(())
})();
if let Err(err) = playback {
error!(
renderer = renderer_id.0.as_str(),
error = %err,
"Failed to start playback for queued item"
);
if self
.runtime
.with_queue_mut(renderer_id, |queue| queue.enqueue_front(item))
.is_none()
{
warn!(
renderer = renderer_id.0.as_str(),
"Failed to requeue item after playback error"
);
}
self.runtime
.set_playback_source(renderer_id, PlaybackSource::None);
return Err(err);
}
self.runtime
.set_playback_source(renderer_id, PlaybackSource::FromQueue);
debug!(
renderer = renderer_id.0.as_str(),
queue_len = remaining_after,
"Started playback from queue"
);
Ok(())
}
/// Subscribe to renderer events emitted by the control point runtime.
///
/// Each subscriber receives all future events independently.
@@ -261,13 +440,51 @@ impl ControlPoint {
self.event_bus.subscribe()
}
#[allow(dead_code)]
pub(crate) fn emit_renderer_event(&self, event: RendererEvent) {
self.handle_renderer_event(&event);
self.event_bus.broadcast(event);
}
fn handle_renderer_event(&self, event: &RendererEvent) {
if let RendererEvent::StateChanged { id, state } = event {
match state {
PlaybackState::Stopped => {
if self.runtime.is_playing_from_queue(id) {
debug!(
renderer = id.0.as_str(),
"Renderer stopped after queue-driven playback; advancing"
);
if let Err(err) = self.play_next_from_queue(id) {
error!(
renderer = id.0.as_str(),
error = %err,
"Auto-advance failed; clearing queue playback state"
);
self.runtime
.set_playback_source(id, PlaybackSource::None);
}
} else {
self.runtime
.set_playback_source(id, PlaybackSource::None);
}
}
PlaybackState::Playing => {
self.runtime.mark_external_if_idle(id);
}
_ => {}
}
}
}
fn runtime_entry_missing(renderer_id: &RendererId) -> anyhow::Error {
anyhow!(
"Renderer {} not registered in control point runtime",
renderer_id.0
)
}
}
#[derive(Default)]
#[derive(Clone, Default)]
struct RendererRuntimeSnapshot {
state: Option<PlaybackState>,
position: Option<PlaybackPositionInfo>,
@@ -275,6 +492,117 @@ struct RendererRuntimeSnapshot {
last_mute: Option<bool>,
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
enum PlaybackSource {
#[default]
None,
FromQueue,
External,
}
#[derive(Default)]
struct RendererRuntimeEntry {
snapshot: RendererRuntimeSnapshot,
queue: PlaybackQueue,
playback_source: PlaybackSource,
}
struct RuntimeState {
entries: Mutex<HashMap<RendererId, RendererRuntimeEntry>>,
}
impl RuntimeState {
fn new() -> Self {
Self {
entries: Mutex::new(HashMap::new()),
}
}
fn snapshot_for(&self, id: &RendererId) -> RendererRuntimeSnapshot {
let entries = self.entries.lock().unwrap();
entries
.get(id)
.map(|entry| entry.snapshot.clone())
.unwrap_or_default()
}
fn update_snapshot(&self, id: &RendererId, snapshot: RendererRuntimeSnapshot) {
self.with_entry(id, |entry| {
entry.snapshot = snapshot;
});
}
fn has_entry(&self, id: &RendererId) -> bool {
let entries = self.entries.lock().unwrap();
entries.contains_key(id)
}
fn with_queue_mut<F, R>(&self, id: &RendererId, f: F) -> Option<R>
where
F: FnOnce(&mut PlaybackQueue) -> R,
{
let mut entries = self.entries.lock().unwrap();
entries
.get_mut(id)
.map(|entry| f(&mut entry.queue))
}
fn queue_snapshot(&self, id: &RendererId) -> Option<Vec<PlaybackItem>> {
let entries = self.entries.lock().unwrap();
entries.get(id).map(|entry| entry.queue.snapshot())
}
fn dequeue_next(&self, id: &RendererId) -> Option<(PlaybackItem, usize)> {
let mut entries = self.entries.lock().unwrap();
let entry = entries.get_mut(id)?;
let item = entry.queue.dequeue()?;
let remaining = entry.queue.len();
Some((item, remaining))
}
fn set_playback_source(&self, id: &RendererId, source: PlaybackSource) {
let mut entries = self.entries.lock().unwrap();
if let Some(entry) = entries.get_mut(id) {
entry.playback_source = source;
}
}
fn playback_source(&self, id: &RendererId) -> PlaybackSource {
let entries = self.entries.lock().unwrap();
entries
.get(id)
.map(|entry| entry.playback_source)
.unwrap_or(PlaybackSource::None)
}
fn is_playing_from_queue(&self, id: &RendererId) -> bool {
matches!(
self.playback_source(id),
PlaybackSource::FromQueue
)
}
fn mark_external_if_idle(&self, id: &RendererId) {
let mut entries = self.entries.lock().unwrap();
if let Some(entry) = entries.get_mut(id) {
if matches!(entry.playback_source, PlaybackSource::None) {
entry.playback_source = PlaybackSource::External;
}
}
}
fn with_entry<F, R>(&self, id: &RendererId, f: F) -> R
where
F: FnOnce(&mut RendererRuntimeEntry) -> R,
{
let mut entries = self.entries.lock().unwrap();
let entry = entries
.entry(id.clone())
.or_insert_with(RendererRuntimeEntry::default);
f(entry)
}
}
/// Parse "HH:MM:SS" style time strings to seconds.
///
/// Returns None for empty or sentinel values such as "NOT_IMPLEMENTED" or "-:--:--".

View File

@@ -3,7 +3,8 @@ use std::time::SystemTime;
use pmoupnp::ssdp::SsdpEvent;
use crate::model::{MediaServerInfo, RendererInfo};
use crate::media_server::MediaServerInfo;
use crate::model::RendererInfo;
use crate::registry::DeviceUpdate;
/// État connu pour un endpoint UPnP identifié par son UDN.

View File

@@ -7,8 +7,10 @@ pub mod connection_manager_client;
pub mod control_point;
pub mod discovery;
pub mod linkplay;
pub mod media_server;
pub mod model;
pub mod music_renderer;
pub mod playback_queue;
pub mod provider;
pub mod registry;
pub mod rendering_control_client;
@@ -24,15 +26,17 @@ pub use capabilities::{
pub use connection_manager_client::{ConnectionInfo, ConnectionManagerClient, ProtocolInfo};
pub use control_point::ControlPoint;
pub use linkplay::LinkPlayRenderer;
pub use media_server::{
MediaBrowser, MediaEntry, MediaResource, MediaServerInfo, MusicServer, ServerId,
UpnpMediaServer,
};
pub use music_renderer::MusicRenderer;
pub use playback_queue::{PlaybackItem, PlaybackQueue};
pub use rendering_control_client::RenderingControlClient;
pub use upnp_renderer::UpnpRenderer;
pub use discovery::{DeviceDescriptionProvider, DiscoveredEndpoint, DiscoveryManager};
pub use model::{
MediaServerCapabilities, MediaServerId, MediaServerInfo, RendererCapabilities, RendererEvent,
RendererId, RendererInfo, RendererProtocol,
};
pub use model::{RendererCapabilities, RendererEvent, RendererId, RendererInfo, RendererProtocol};
pub use provider::HttpXmlDescriptionProvider;
pub use registry::{DeviceRegistry, DeviceRegistryRead, DeviceUpdate};

View File

@@ -0,0 +1,430 @@
use std::time::{Duration, SystemTime};
use anyhow::{Result, anyhow};
use pmodidl::{self, DIDLLite};
use pmoupnp::soap::SoapEnvelope;
use pmoupnp::soap::error_codes;
use xmltree::{Element, XMLNode};
use crate::soap_client::{SoapCallResult, invoke_upnp_action_with_timeout};
/// Unique identifier for a media server registered by the control point.
#[derive(Clone, Debug, Eq, PartialEq, Hash)]
pub struct ServerId(pub String);
/// Snapshot of a media server discovered through UPnP SSDP.
#[derive(Clone, Debug)]
pub struct MediaServerInfo {
pub id: ServerId,
pub udn: String,
pub friendly_name: String,
pub model_name: String,
pub manufacturer: String,
pub location: String,
pub server_header: String,
pub online: bool,
pub last_seen: SystemTime,
pub max_age: u32,
pub has_content_directory: bool,
pub content_directory_service_type: Option<String>,
pub content_directory_control_url: Option<String>,
}
/// Simplified view over a DIDL-Lite resource entry.
#[derive(Clone, Debug)]
pub struct MediaResource {
pub uri: String,
pub protocol_info: String,
pub duration: Option<String>,
}
/// Representation of either a container or an item returned by ContentDirectory.
#[derive(Clone, Debug)]
pub struct MediaEntry {
pub id: String,
pub parent_id: String,
pub title: String,
pub is_container: bool,
pub class: String,
pub resources: Vec<MediaResource>,
}
/// Backend-agnostic media browsing contract.
pub trait MediaBrowser {
fn browse_root(&self) -> Result<Vec<MediaEntry>>;
fn browse_children(&self, object_id: &str, start: u32, count: u32) -> Result<Vec<MediaEntry>>;
fn browse_object(&self, object_id: &str) -> Result<MediaEntry>;
fn search(
&self,
container_id: &str,
query: &str,
start: u32,
count: u32,
) -> Result<Vec<MediaEntry>>;
}
/// Façade over every supported media server backend.
#[derive(Clone, Debug)]
pub enum MusicServer {
Upnp(UpnpMediaServer),
}
impl MusicServer {
pub fn from_info(info: &MediaServerInfo, timeout: Duration) -> Result<Self> {
Ok(MusicServer::Upnp(UpnpMediaServer::new(
info.clone(),
timeout,
)))
}
pub fn id(&self) -> &ServerId {
match self {
MusicServer::Upnp(upnp) => upnp.id(),
}
}
pub fn info(&self) -> &MediaServerInfo {
match self {
MusicServer::Upnp(upnp) => upnp.info(),
}
}
}
impl MediaBrowser for MusicServer {
fn browse_root(&self) -> Result<Vec<MediaEntry>> {
match self {
MusicServer::Upnp(upnp) => upnp.browse_root(),
}
}
fn browse_children(&self, object_id: &str, start: u32, count: u32) -> Result<Vec<MediaEntry>> {
match self {
MusicServer::Upnp(upnp) => upnp.browse_children(object_id, start, count),
}
}
fn browse_object(&self, object_id: &str) -> Result<MediaEntry> {
match self {
MusicServer::Upnp(upnp) => upnp.browse_object(object_id),
}
}
fn search(
&self,
container_id: &str,
query: &str,
start: u32,
count: u32,
) -> Result<Vec<MediaEntry>> {
match self {
MusicServer::Upnp(upnp) => upnp.search(container_id, query, start, count),
}
}
}
/// Single UPnP ContentDirectory backend implementation.
#[derive(Clone, Debug)]
pub struct UpnpMediaServer {
info: MediaServerInfo,
timeout: Duration,
}
impl UpnpMediaServer {
pub fn new(info: MediaServerInfo, timeout: Duration) -> Self {
Self { info, timeout }
}
pub fn id(&self) -> &ServerId {
&self.info.id
}
pub fn info(&self) -> &MediaServerInfo {
&self.info
}
fn browse_with_flag(
&self,
object_id: &str,
browse_flag: &str,
start: u32,
count: u32,
) -> Result<Vec<MediaEntry>> {
let start_str = start.to_string();
let count_str = count.to_string();
let args = vec![
("ObjectID", object_id.to_string()),
("BrowseFlag", browse_flag.to_string()),
("Filter", "*".to_string()),
("StartingIndex", start_str),
("RequestedCount", count_str),
("SortCriteria", String::new()),
];
let response = self.invoke_content_directory("Browse", None, args)?;
let envelope = response
.envelope
.ok_or_else(|| anyhow!("Missing SOAP envelope in Browse response"))?;
let didl_xml = extract_result_payload(&envelope, "BrowseResponse")?;
map_didl_entries(&didl_xml)
}
fn search_impl(
&self,
container_id: &str,
query: &str,
start: u32,
count: u32,
) -> Result<Vec<MediaEntry>> {
let start_str = start.to_string();
let count_str = count.to_string();
let args = vec![
("ContainerID", container_id.to_string()),
("SearchCriteria", query.to_string()),
("Filter", "*".to_string()),
("StartingIndex", start_str),
("RequestedCount", count_str),
("SortCriteria", String::new()),
];
let response = self.invoke_content_directory("Search", Some("search"), args)?;
let envelope = response
.envelope
.ok_or_else(|| anyhow!("Missing SOAP envelope in Search response"))?;
let didl_xml = extract_result_payload(&envelope, "SearchResponse")?;
map_didl_entries(&didl_xml)
}
fn invoke_content_directory(
&self,
action: &str,
op_name: Option<&str>,
args: Vec<(&'static str, String)>,
) -> Result<SoapCallResult> {
let op = op_name.unwrap_or(action);
let (control_url, service_type) = self.content_directory_endpoints(op)?;
let borrowed_args: Vec<(&str, &str)> = args.iter().map(|(k, v)| (*k, v.as_str())).collect();
let call_result = invoke_upnp_action_with_timeout(
control_url,
service_type,
action,
&borrowed_args,
Some(self.timeout),
)?;
if !call_result.status.is_success() {
if let Some(env) = &call_result.envelope {
if let Some(err) = parse_upnp_error(env) {
if should_map_to_not_supported(action, op_name, err.error_code) {
return Err(server_op_not_supported(op, "UpnpMediaServer"));
}
return Err(anyhow!(
"{} failed with UPnP error {}: {}",
action,
err.error_code,
err.error_description
));
}
}
return Err(anyhow!(
"{} failed with HTTP status {} and body: {}",
action,
call_result.status,
call_result.raw_body
));
}
if let Some(env) = &call_result.envelope {
if let Some(err) = parse_upnp_error(env) {
if should_map_to_not_supported(action, op_name, err.error_code) {
return Err(server_op_not_supported(op, "UpnpMediaServer"));
}
return Err(anyhow!(
"{} returned UPnP error {}: {}",
action,
err.error_code,
err.error_description
));
}
}
Ok(call_result)
}
fn content_directory_endpoints(&self, op_name: &str) -> Result<(&str, &str)> {
if !self.info.has_content_directory {
return Err(server_op_not_supported(op_name, "UpnpMediaServer"));
}
let control_url = self
.info
.content_directory_control_url
.as_deref()
.ok_or_else(|| server_op_not_supported(op_name, "UpnpMediaServer"))?;
let service_type = self
.info
.content_directory_service_type
.as_deref()
.ok_or_else(|| server_op_not_supported(op_name, "UpnpMediaServer"))?;
Ok((control_url, service_type))
}
}
impl MediaBrowser for UpnpMediaServer {
fn browse_root(&self) -> Result<Vec<MediaEntry>> {
self.browse_with_flag("0", "BrowseDirectChildren", 0, 0)
}
fn browse_children(&self, object_id: &str, start: u32, count: u32) -> Result<Vec<MediaEntry>> {
self.browse_with_flag(object_id, "BrowseDirectChildren", start, count)
}
fn browse_object(&self, object_id: &str) -> Result<MediaEntry> {
let entries = self.browse_with_flag(object_id, "BrowseMetadata", 0, 1)?;
entries
.into_iter()
.next()
.ok_or_else(|| anyhow!("Object {} was not returned by the server", object_id))
}
fn search(
&self,
container_id: &str,
query: &str,
start: u32,
count: u32,
) -> Result<Vec<MediaEntry>> {
self.search_impl(container_id, query, start, count)
}
}
fn map_didl_entries(xml: &str) -> Result<Vec<MediaEntry>> {
let trimmed = xml.trim();
if trimmed.is_empty() {
return Ok(Vec::new());
}
let didl: DIDLLite = pmodidl::parse_metadata::<DIDLLite>(trimmed)
.map_err(|err| anyhow!("Failed to parse DIDL-Lite payload: {}", err))?
.data;
let mut entries = Vec::new();
for container in didl.containers {
entries.push(MediaEntry {
id: container.id,
parent_id: container.parent_id,
title: container.title,
is_container: true,
class: container.class,
resources: Vec::new(),
});
}
for item in didl.items {
let resources = item
.resources
.into_iter()
.map(|res| MediaResource {
uri: res.url,
protocol_info: res.protocol_info,
duration: res.duration,
})
.collect();
entries.push(MediaEntry {
id: item.id,
parent_id: item.parent_id,
title: item.title,
is_container: false,
class: item.class,
resources,
});
}
Ok(entries)
}
fn extract_result_payload(envelope: &SoapEnvelope, response_suffix: &str) -> Result<String> {
let response = find_child_with_suffix(&envelope.body.content, response_suffix)
.ok_or_else(|| anyhow!("Missing {} element in SOAP body", response_suffix))?;
let result_elem = find_child_with_suffix(response, "Result")
.ok_or_else(|| anyhow!("Missing Result element in {}", response_suffix))?;
let payload = result_elem
.get_text()
.map(|t| t.to_string())
.unwrap_or_default();
Ok(payload)
}
fn find_child_with_suffix<'a>(parent: &'a Element, suffix: &str) -> Option<&'a Element> {
parent.children.iter().find_map(|node| match node {
XMLNode::Element(elem) if elem.name.ends_with(suffix) => Some(elem),
_ => None,
})
}
fn parse_upnp_error(envelope: &SoapEnvelope) -> Option<UpnpError> {
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::<u32>().ok()?;
let error_description = upnp_error
.children
.iter()
.find_map(|node| match node {
XMLNode::Element(elem) if elem.name.ends_with("errorDescription") => {
elem.get_text().map(|t| t.trim().to_string())
}
_ => None,
})
.unwrap_or_default();
Some(UpnpError {
error_code,
error_description,
})
}
fn should_map_to_not_supported(action: &str, op_name: Option<&str>, error_code: u32) -> bool {
if op_name.is_none() {
return false;
}
let optional_code = error_codes::OPTIONAL_ACTION_NOT_IMPLEMENTED
.parse::<u32>()
.unwrap_or(602);
let invalid_action_code = error_codes::INVALID_ACTION.parse::<u32>().unwrap_or(401);
let cd_not_supported = matches!(action, "Search");
cd_not_supported && (error_code == optional_code || error_code == invalid_action_code)
}
fn server_op_not_supported(op: &str, backend: &str) -> anyhow::Error {
anyhow!(
"MusicServer operation '{}' is not supported by backend '{}'",
op,
backend
)
}
#[derive(Debug, Clone)]
struct UpnpError {
pub error_code: u32,
pub error_description: String,
}

View File

@@ -3,9 +3,6 @@ use crate::capabilities::{PlaybackPositionInfo, PlaybackState};
#[derive(Clone, Debug, PartialEq, Eq, Hash)]
pub struct RendererId(pub String);
#[derive(Clone, Debug, PartialEq, Eq, Hash)]
pub struct MediaServerId(pub String);
#[derive(Clone, Debug)]
pub enum RendererProtocol {
UpnpAvOnly,
@@ -53,29 +50,6 @@ pub struct RendererInfo {
pub connection_manager_control_url: Option<String>,
}
#[derive(Clone, Debug, Default)]
pub struct MediaServerCapabilities {
pub has_content_directory: bool,
pub has_connection_manager: bool,
}
#[derive(Clone, Debug)]
pub struct MediaServerInfo {
pub id: MediaServerId,
pub udn: String,
pub friendly_name: String,
pub model_name: String,
pub manufacturer: String,
pub capabilities: MediaServerCapabilities,
pub location: String,
pub server_header: String,
pub online: bool,
pub last_seen: std::time::SystemTime,
pub max_age: u32,
}
#[derive(Clone, Debug)]
pub enum RendererEvent {
StateChanged {

View File

@@ -12,7 +12,7 @@ use crate::{
ArylicTcpRenderer, DeviceRegistry, LinkPlayRenderer, PlaybackPosition, PlaybackState,
TransportControl, UpnpRenderer, VolumeControl,
};
use anyhow::{anyhow, Result};
use anyhow::{Result, anyhow};
use tracing::warn;
/// Backend-agnostic façade exposing transport, volume, and status contracts.
@@ -33,7 +33,7 @@ pub enum MusicRenderer {
}
/// Build a standardized error when an operation is not supported by a backend.
fn op_not_supported(op: &str, backend: &str) -> anyhow::Error {
pub(crate) fn op_not_supported(op: &str, backend: &str) -> anyhow::Error {
anyhow!(
"MusicRenderer operation '{}' is not supported by backend '{}'",
op,

View File

@@ -0,0 +1,73 @@
use std::collections::VecDeque;
use crate::media_server::ServerId;
#[derive(Clone, Debug)]
pub struct PlaybackItem {
pub uri: String,
pub title: Option<String>,
pub server_id: Option<ServerId>,
pub object_id: Option<String>,
}
impl PlaybackItem {
pub fn new(uri: impl Into<String>) -> Self {
Self {
uri: uri.into(),
title: None,
server_id: None,
object_id: None,
}
}
}
#[derive(Clone, Debug, Default)]
pub struct PlaybackQueue {
items: VecDeque<PlaybackItem>,
}
impl PlaybackQueue {
pub fn new() -> Self {
Self {
items: VecDeque::new(),
}
}
pub fn len(&self) -> usize {
self.items.len()
}
pub fn is_empty(&self) -> bool {
self.items.is_empty()
}
pub fn clear(&mut self) {
self.items.clear();
}
pub fn enqueue(&mut self, item: PlaybackItem) {
self.items.push_back(item);
}
pub fn enqueue_many<I: IntoIterator<Item = PlaybackItem>>(&mut self, items: I) {
for item in items {
self.items.push_back(item);
}
}
pub fn enqueue_front(&mut self, item: PlaybackItem) {
self.items.push_front(item);
}
pub fn dequeue(&mut self) -> Option<PlaybackItem> {
self.items.pop_front()
}
pub fn peek(&self) -> Option<&PlaybackItem> {
self.items.front()
}
pub fn snapshot(&self) -> Vec<PlaybackItem> {
self.items.iter().cloned().collect()
}
}

View File

@@ -9,10 +9,8 @@ use crate::arylic_tcp::detect_arylic_tcp;
use crate::avtransport_client::AvTransportClient;
use crate::discovery::{DeviceDescriptionProvider, DiscoveredEndpoint};
use crate::linkplay::detect_linkplay_http;
use crate::model::{
MediaServerCapabilities, MediaServerId, MediaServerInfo, RendererCapabilities, RendererId,
RendererInfo, RendererProtocol,
};
use crate::media_server::{MediaServerInfo, ServerId};
use crate::model::{RendererCapabilities, RendererId, RendererInfo, RendererProtocol};
use ureq::Agent;
@@ -52,6 +50,10 @@ struct ParsedDeviceDescription {
// ConnectionManager endpoint (if present in serviceList)
connection_manager_service_type: Option<String>,
connection_manager_control_url: Option<String>,
// ContentDirectory endpoint (if present in serviceList)
content_directory_service_type: Option<String>,
content_directory_control_url: Option<String>,
}
impl ParsedDeviceDescription {
@@ -201,6 +203,21 @@ impl HttpXmlDescriptionProvider {
);
}
}
if lower
.contains("urn:schemas-upnp-org:service:contentdirectory:")
{
if parsed.content_directory_service_type.is_none() {
parsed.content_directory_service_type =
Some(st.clone());
parsed.content_directory_control_url =
Some(ctrl.clone());
debug!(
"Found ContentDirectory service for {}: type={} controlURL={}",
endpoint.udn, st, ctrl
);
}
}
}
in_service = false;
@@ -340,21 +357,31 @@ impl HttpXmlDescriptionProvider {
.as_deref()
.unwrap_or_else(|| endpoint.udn.as_str());
let udn = raw_udn.to_ascii_lowercase();
let caps = detect_server_capabilities(&parsed.service_types);
let has_content_directory = parsed.service_types.iter().any(|st| {
st.to_ascii_lowercase()
.contains("urn:schemas-upnp-org:service:contentdirectory:")
});
let now = SystemTime::now();
let content_directory_control_url = parsed
.content_directory_control_url
.as_ref()
.map(|ctrl| resolve_control_url(&endpoint.location, ctrl));
Some(MediaServerInfo {
id: MediaServerId(udn.clone()),
id: ServerId(udn.clone()),
udn,
friendly_name: parsed.friendly_name.clone().unwrap_or_default(),
model_name: parsed.model_name.clone().unwrap_or_default(),
manufacturer: parsed.manufacturer.clone().unwrap_or_default(),
capabilities: caps,
location: endpoint.location.clone(),
server_header: endpoint.server_header.clone(),
online: true,
last_seen: now,
max_age: endpoint.max_age,
has_content_directory,
content_directory_service_type: parsed.content_directory_service_type.clone(),
content_directory_control_url,
})
}
@@ -442,22 +469,6 @@ fn detect_renderer_protocol(caps: &RendererCapabilities) -> RendererProtocol {
}
}
fn detect_server_capabilities(service_types: &[String]) -> MediaServerCapabilities {
let mut caps = MediaServerCapabilities::default();
for st in service_types {
let lower = st.to_ascii_lowercase();
if lower.contains("urn:schemas-upnp-org:service:contentdirectory:") {
caps.has_content_directory = true;
}
if lower.contains("urn:schemas-upnp-org:service:connectionmanager:") {
caps.has_connection_manager = true;
}
}
caps
}
/// Resolve a possibly relative controlURL against the description URL.
///
/// - If `control_url` is already absolute (starts with http:// or https://), it is returned as-is.

View File

@@ -3,19 +3,20 @@ use std::time::SystemTime;
use crate::avtransport_client::AvTransportClient;
use crate::connection_manager_client::ConnectionManagerClient;
use crate::model::{MediaServerId, MediaServerInfo, RendererId, RendererInfo};
use crate::media_server::{MediaServerInfo, ServerId};
use crate::model::{RendererId, RendererInfo};
use crate::rendering_control_client::RenderingControlClient;
#[derive(Clone, Debug)]
enum DeviceKey {
Renderer(RendererId),
Server(MediaServerId),
Server(ServerId),
}
#[derive(Debug, Default)]
pub struct DeviceRegistry {
renderers: HashMap<RendererId, RendererInfo>,
servers: HashMap<MediaServerId, MediaServerInfo>,
servers: HashMap<ServerId, MediaServerInfo>,
udn_index: HashMap<String, DeviceKey>,
}
@@ -28,7 +29,7 @@ pub trait DeviceRegistryRead {
fn list_servers(&self) -> Vec<MediaServerInfo>;
fn get_renderer(&self, id: &RendererId) -> Option<RendererInfo>;
fn get_server(&self, id: &MediaServerId) -> Option<MediaServerInfo>;
fn get_server(&self, id: &ServerId) -> Option<MediaServerInfo>;
}
impl DeviceRegistryRead for DeviceRegistry {
@@ -44,7 +45,7 @@ impl DeviceRegistryRead for DeviceRegistry {
self.renderers.get(id).cloned()
}
fn get_server(&self, id: &MediaServerId) -> Option<MediaServerInfo> {
fn get_server(&self, id: &ServerId) -> Option<MediaServerInfo> {
self.servers.get(id).cloned()
}
}
@@ -56,7 +57,7 @@ pub enum DeviceUpdate {
RendererOfflineByUdn(String),
ServerOnline(MediaServerInfo),
ServerOfflineById(MediaServerId),
ServerOfflineById(ServerId),
ServerOfflineByUdn(String),
}

View File

@@ -1,3 +1,5 @@
use std::time::Duration;
use anyhow::{Context, Result};
use pmoupnp::soap::{SoapEnvelope, build_soap_request, parse_soap_envelope};
use ureq::Agent;
@@ -24,13 +26,26 @@ pub fn invoke_upnp_action(
action: &str,
args: &[(&str, &str)],
) -> Result<SoapCallResult> {
// 1. Build SOAP request XML using pmoupnp
invoke_upnp_action_with_timeout(control_url, service_type, action, args, None)
}
pub fn invoke_upnp_action_with_timeout(
control_url: &str,
service_type: &str,
action: &str,
args: &[(&str, &str)],
timeout: Option<Duration>,
) -> Result<SoapCallResult> {
let body_xml = build_soap_request(service_type, action, args)
.context("Failed to build SOAP request body")?;
// 2. Build agent that does NOT treat 4xx/5xx as errors
let config = Agent::config_builder().http_status_as_error(false).build();
let mut builder = Agent::config_builder();
builder = builder.http_status_as_error(false);
if let Some(duration) = timeout {
builder = builder.timeout_global(Some(duration));
}
let config = builder.build();
let agent: Agent = config.into();
// 3. SOAPAction header