Deboguage pmo_remote_control

This commit is contained in:
2025-12-03 23:17:41 +01:00
parent 1432f2507c
commit 75878be564
11 changed files with 2620 additions and 269 deletions

15
Cargo.lock generated
View File

@@ -3252,6 +3252,7 @@ dependencies = [
"chrono",
"crossbeam-channel",
"crossterm",
"percent-encoding",
"pmodidl",
"pmoserver",
"pmoupnp",
@@ -3263,6 +3264,7 @@ dependencies = [
"tokio",
"tokio-stream",
"tracing",
"tracing-log 0.1.4",
"tracing-subscriber",
"ureq",
"utoipa",
@@ -5232,6 +5234,17 @@ dependencies = [
"valuable",
]
[[package]]
name = "tracing-log"
version = "0.1.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "f751112709b4e791d8ce53e32c4ed2d353565a795ce84da2285393f41557bdf2"
dependencies = [
"log",
"once_cell",
"tracing-core",
]
[[package]]
name = "tracing-log"
version = "0.2.0"
@@ -5258,7 +5271,7 @@ dependencies = [
"thread_local",
"tracing",
"tracing-core",
"tracing-log",
"tracing-log 0.2.0",
]
[[package]]

View File

@@ -11,6 +11,7 @@ thiserror = "2.0.17"
ureq = "3.1.4"
tracing = "0.1.41"
tracing-subscriber = "0.3"
tracing-log = "0.1"
anyhow = "1.0"
xmltree = "0.11.0"
crossbeam-channel = "0.5"
@@ -29,6 +30,11 @@ tokio-stream = { version = "0.1", features = ["sync"], optional = true }
async-stream = { version = "0.3", optional = true }
chrono = { version = "0.4", features = ["serde"], optional = true }
[dev-dependencies]
serde = { version = "1.0", features = ["derive"] }
serde_json = "1.0"
percent-encoding = "2.3"
[features]
default = []
# Active l'API REST pmoserver

View File

@@ -6,29 +6,29 @@
use std::collections::HashMap;
use std::io::{self, Stdout};
use std::process;
use std::sync::mpsc::{self, Sender};
use std::sync::Arc;
use std::sync::mpsc::{self, Sender};
use std::thread;
use std::time::{Duration, Instant};
use anyhow::{anyhow, Context, Result};
use anyhow::{Context, Result, anyhow};
use crossterm::event::{self, DisableMouseCapture, EnableMouseCapture, Event, KeyCode, KeyEvent};
use crossterm::execute;
use crossterm::terminal::{
disable_raw_mode, enable_raw_mode, EnterAlternateScreen, LeaveAlternateScreen,
EnterAlternateScreen, LeaveAlternateScreen, disable_raw_mode, enable_raw_mode,
};
use pmocontrol::model::TrackMetadata;
use pmocontrol::{
ControlPoint, DeviceRegistryRead, MediaBrowser, MediaEntry, MediaResource, MediaServerEvent,
MediaServerInfo, MusicServer, PlaybackItem, PlaybackPosition, PlaybackStatus,
PlaybackPositionInfo, RendererEvent, RendererInfo, TransportControl, VolumeControl,
MediaServerInfo, MusicServer, PlaybackItem, PlaybackPosition, PlaybackPositionInfo,
PlaybackStatus, RendererEvent, RendererInfo, TransportControl, VolumeControl,
};
use ratatui::Terminal;
use ratatui::backend::CrosstermBackend;
use ratatui::layout::{Alignment, Constraint, Direction, Layout, Rect};
use ratatui::style::{Color, Modifier, Style};
use ratatui::text::{Line, Span};
use ratatui::widgets::{Block, Borders, Clear, Gauge, List, ListItem, ListState, Paragraph};
use ratatui::Terminal;
const DEFAULT_TIMEOUT_SECS: u64 = 5;
const DEFAULT_DISCOVERY_SECS: u64 = 15;
@@ -725,10 +725,10 @@ impl App {
self.control_point.clear_queue(&renderer_id)?;
self.control_point.enqueue_items(&renderer_id, items)?;
let (full_queue, current_index) =
self.control_point
.get_full_queue_snapshot(&renderer_id)
.context("failed to snapshot queue after enqueue")?;
let (full_queue, current_index) = self
.control_point
.get_full_queue_snapshot(&renderer_id)
.context("failed to snapshot queue after enqueue")?;
self.queue_snapshot = full_queue;
self.queue_current_index = current_index;
self.record_queue_metadata();
@@ -948,19 +948,15 @@ impl App {
}
fn refresh_queue_snapshot(&mut self, renderer_id: &pmocontrol::model::RendererId) {
if let Ok((queue, current_index)) =
self.control_point.get_full_queue_snapshot(renderer_id)
if let Ok((queue, current_index)) = self.control_point.get_full_queue_snapshot(renderer_id)
{
self.queue_snapshot = queue;
self.queue_current_index = current_index;
self.record_queue_metadata();
if let Some(idx) = self.queue_current_index {
if let Some(item) = self.queue_snapshot.get(idx).cloned() {
let needs_update = self
.ui_state
.current_track_uri
.as_deref()
!= Some(item.uri.as_str());
let needs_update =
self.ui_state.current_track_uri.as_deref() != Some(item.uri.as_str());
if needs_update {
self.apply_item_as_current(&item);
}

File diff suppressed because it is too large Load Diff

View File

@@ -5,9 +5,7 @@
#[cfg(not(feature = "pmoserver"))]
fn main() {
eprintln!(
"This example requires the 'pmoserver' feature. Re-run with `--features pmoserver`."
);
eprintln!("This example requires the 'pmoserver' feature. Re-run with `--features pmoserver`.");
}
#[cfg(feature = "pmoserver")]

View File

@@ -45,6 +45,8 @@ struct PlaylistBinding {
/// Flag used internally to signal that the queue should be refreshed
/// from the server container.
pending_refresh: bool,
/// Whether the next refresh should auto-start playback if the renderer is idle.
auto_play_on_refresh: bool,
}
/// Control point minimal :
@@ -326,6 +328,7 @@ impl ControlPoint {
&bindings_for_media_worker,
&renderer_id,
&event_bus_for_media_worker,
None,
) {
warn!(
renderer = renderer_id.0.as_str(),
@@ -380,6 +383,7 @@ impl ControlPoint {
&bindings_for_periodic,
&renderer_id,
&event_bus_for_periodic,
None,
) {
warn!(
renderer = renderer_id.0.as_str(),
@@ -718,6 +722,29 @@ impl ControlPoint {
Ok(())
}
fn start_queue_playback_if_idle(&self, renderer_id: &RendererId) -> anyhow::Result<()> {
let snapshot = self.runtime.snapshot_for(renderer_id);
let renderer_playing = matches!(snapshot.state, Some(PlaybackState::Playing));
if renderer_playing || self.runtime.is_playing_from_queue(renderer_id) {
return Ok(());
}
let has_items = self
.runtime
.queue_snapshot(renderer_id)
.map(|items| !items.is_empty())
.unwrap_or(false);
if !has_items {
debug!(
renderer = renderer_id.0.as_str(),
"start_queue_playback_if_idle: queue is empty"
);
return Ok(());
}
self.play_next_from_queue(renderer_id)
}
/// Subscribe to renderer events emitted by the control point runtime.
///
/// Each subscriber receives all future events independently.
@@ -747,22 +774,43 @@ impl ControlPoint {
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"
);
{
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: true,
auto_play_on_refresh: true,
},
);
info!(
renderer = renderer_id.0.as_str(),
server = server_id.0.as_str(),
container = container_id.as_str(),
"Queue attached to playlist container"
);
} // Drop bindings lock here before calling refresh_attached_queue_for
let mut auto_start_cb = |rid: &RendererId| self.start_queue_playback_if_idle(rid);
if let Err(err) = refresh_attached_queue_for(
&self.registry,
&self.runtime,
&self.playlist_bindings,
renderer_id,
&self.event_bus,
Some(&mut auto_start_cb),
) {
warn!(
renderer = renderer_id.0.as_str(),
server = server_id.0.as_str(),
container = container_id.as_str(),
error = %err,
"Initial playlist refresh after attachment failed"
);
}
}
/// Detach a renderer's queue from its associated playlist container.
@@ -995,9 +1043,10 @@ fn refresh_attached_queue_for(
bindings: &Arc<Mutex<HashMap<RendererId, PlaylistBinding>>>,
renderer_id: &RendererId,
event_bus: &RendererEventBus,
mut after_refresh: Option<&mut dyn FnMut(&RendererId) -> anyhow::Result<()>>,
) -> anyhow::Result<()> {
// Step 1: Check binding and mark refresh as in-progress
let (server_id, container_id) = {
let (server_id, container_id, auto_play) = {
let mut bindings_lock = bindings.lock().unwrap();
let binding = match bindings_lock.get_mut(renderer_id) {
Some(b) => b,
@@ -1020,7 +1069,13 @@ fn refresh_attached_queue_for(
// Mark as processed
binding.pending_refresh = false;
(binding.server_id.clone(), binding.container_id.clone())
let auto_play = binding.auto_play_on_refresh;
binding.auto_play_on_refresh = false;
(
binding.server_id.clone(),
binding.container_id.clone(),
auto_play,
)
};
// Step 2: Fetch MediaServerInfo from registry
@@ -1101,9 +1156,13 @@ fn refresh_attached_queue_for(
}
// 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());
// Get the full queue snapshot to access the item currently being played
let (full_queue, current_idx) = runtime
.queue_full_snapshot(renderer_id)
.unwrap_or((vec![], None));
// Get the item currently being played (at current_index), not the next one in queue
let current_item = current_idx.and_then(|idx| full_queue.get(idx).cloned());
let item_found_at = current_item.as_ref().and_then(|current| {
new_items.iter().position(|new_item| {
@@ -1116,39 +1175,66 @@ fn refresh_attached_queue_for(
})
});
let final_queue_len = runtime.with_queue_mut(renderer_id, |queue| {
queue.clear();
let final_queue_len = 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());
if let Some(idx) = item_found_at {
// Current item found: load the ENTIRE new playlist and position at that item
// This preserves items before the current track (as "already played")
for item in new_items.iter() {
queue.enqueue(item.clone());
}
queue.set_current_index(Some(idx));
info!(
renderer = renderer_id.0.as_str(),
server = server_id.0.as_str(),
container = container_id.as_str(),
total_items = new_items.len(),
current_index = idx,
upcoming = new_items.len().saturating_sub(idx + 1),
current_preserved = true,
"Refreshed queue from playlist container"
);
new_items.len()
} else if let Some(ref current) = current_item {
// Current item NOT found: insert it at the beginning, then add new items
// This preserves the currently playing track and prevents it from being lost
queue.enqueue(current.clone());
for item in new_items.iter() {
queue.enqueue(item.clone());
}
queue.set_current_index(Some(0));
info!(
renderer = renderer_id.0.as_str(),
server = server_id.0.as_str(),
container = container_id.as_str(),
total_items = new_items.len() + 1,
current_index = 0,
upcoming = new_items.len(),
current_preserved = true,
current_reinserted = true,
"Refreshed queue from playlist container (current item reinserted at start)"
);
new_items.len() + 1
} else {
// No current item: replace with full new list
for item in new_items.iter() {
queue.enqueue(item.clone());
}
queue.set_current_index(None);
info!(
renderer = renderer_id.0.as_str(),
server = server_id.0.as_str(),
container = container_id.as_str(),
total_items = new_items.len(),
current_preserved = false,
"Refreshed queue from playlist container (no current item)"
);
new_items.len()
}
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"
);
new_items.len() - idx
} 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)"
);
new_items.len()
}
}).unwrap_or(0);
})
.unwrap_or(0);
// Emit QueueUpdated event
event_bus.broadcast(RendererEvent::QueueUpdated {
@@ -1156,6 +1242,20 @@ fn refresh_attached_queue_for(
queue_length: final_queue_len,
});
if auto_play {
if let Some(callback) = after_refresh.as_deref_mut() {
if let Err(err) = callback(renderer_id) {
warn!(
renderer = renderer_id.0.as_str(),
server = server_id.0.as_str(),
container = container_id.as_str(),
error = %err,
"Failed to auto-start playback after playlist refresh"
);
}
}
}
Ok(())
}

View File

@@ -344,10 +344,15 @@ fn map_didl_entries(xml: &str) -> Result<Vec<MediaEntry>> {
let resources = item
.resources
.into_iter()
.map(|res| MediaResource {
uri: res.url,
protocol_info: res.protocol_info,
duration: res.duration,
.filter_map(|res| {
if res.url.trim().is_empty() {
return None;
}
Some(MediaResource {
uri: res.url,
protocol_info: res.protocol_info,
duration: res.duration,
})
})
.collect();

View File

@@ -129,4 +129,16 @@ impl PlaybackQueue {
pub fn full_snapshot(&self) -> (Vec<PlaybackItem>, Option<usize>) {
(self.items.clone(), self.current_index)
}
pub fn set_current_index(&mut self, index: Option<usize>) {
if let Some(idx) = index {
if idx < self.items.len() {
self.current_index = Some(idx);
} else {
self.current_index = None;
}
} else {
self.current_index = None;
}
}
}

View File

@@ -11,7 +11,7 @@ use crate::media_server::{MediaBrowser, MusicServer, ServerId};
use crate::model::{RendererId, RendererProtocol};
#[cfg(feature = "pmoserver")]
use crate::openapi::{
AttachedPlaylistInfo, AttachPlaylistRequest, BrowseResponse, ContainerEntry, ErrorResponse,
AttachPlaylistRequest, AttachedPlaylistInfo, BrowseResponse, ContainerEntry, ErrorResponse,
MediaServerSummary, QueueItem, QueueSnapshot, RendererState, RendererSummary, SuccessResponse,
VolumeSetRequest,
};
@@ -22,20 +22,49 @@ use crate::{PlaybackPosition, PlaybackStatus, TransportControl, VolumeControl};
use async_trait::async_trait;
#[cfg(feature = "pmoserver")]
use axum::{
Json, Router,
extract::{Path, State},
http::StatusCode,
routing::{get, post},
Json, Router,
};
#[cfg(feature = "pmoserver")]
use std::sync::Arc;
#[cfg(feature = "pmoserver")]
use std::time::Duration;
#[cfg(feature = "pmoserver")]
use tokio::time;
#[cfg(feature = "pmoserver")]
use tracing::{debug, warn};
#[cfg(feature = "pmoserver")]
use utoipa::OpenApi;
#[cfg(feature = "pmoserver")]
const BROWSE_PAGE_SIZE: u32 = 100;
#[cfg(feature = "pmoserver")]
const MEDIA_SERVER_SOAP_TIMEOUT: Duration = Duration::from_secs(15);
#[cfg(feature = "pmoserver")]
const BROWSE_REQUEST_TIMEOUT: Duration = Duration::from_secs(20);
// Timeouts for simple commands (play/pause/stop)
#[cfg(feature = "pmoserver")]
const TRANSPORT_COMMAND_TIMEOUT: Duration = Duration::from_secs(5);
// Timeouts for volume/mute commands (faster than transport)
#[cfg(feature = "pmoserver")]
const VOLUME_COMMAND_TIMEOUT: Duration = Duration::from_secs(3);
// Timeout for queue operations
#[cfg(feature = "pmoserver")]
const QUEUE_COMMAND_TIMEOUT: Duration = Duration::from_secs(10);
// Timeout for attach playlist (includes browse + cache + queue update)
#[cfg(feature = "pmoserver")]
const ATTACH_PLAYLIST_TIMEOUT: Duration = Duration::from_secs(60);
// Timeout for renderer state queries (multiple SOAP calls)
#[cfg(feature = "pmoserver")]
const STATE_QUERY_TIMEOUT: Duration = Duration::from_secs(8);
/// État partagé pour l'API ControlPoint
#[cfg(feature = "pmoserver")]
#[derive(Clone)]
@@ -64,9 +93,7 @@ impl ControlPointState {
),
tag = "control"
)]
async fn list_renderers(
State(state): State<ControlPointState>,
) -> Json<Vec<RendererSummary>> {
async fn list_renderers(State(state): State<ControlPointState>) -> Json<Vec<RendererSummary>> {
let renderers = state.control_point.list_music_renderers();
let summaries: Vec<RendererSummary> = renderers
@@ -119,30 +146,64 @@ async fn get_renderer_state(
})?;
let info = renderer.info();
let renderer_clone = renderer.clone();
// État de transport
let transport_state = renderer
.playback_state()
.ok()
.map(state_to_string)
.unwrap_or_else(|| "UNKNOWN".to_string());
// Spawn blocking task for all SOAP calls to avoid blocking Tokio runtime
let state_task = tokio::task::spawn_blocking(move || {
// État de transport
let transport_state = renderer_clone
.playback_state()
.ok()
.map(state_to_string)
.unwrap_or_else(|| "UNKNOWN".to_string());
// Position et durée
let (position_ms, duration_ms) = renderer
.playback_position()
.ok()
.and_then(|pos| {
let position = parse_hms_to_ms(pos.rel_time.as_deref());
let duration = parse_hms_to_ms(pos.track_duration.as_deref());
Some((position, duration))
})
.unwrap_or((None, None));
// Position et durée
let (position_ms, duration_ms) = renderer_clone
.playback_position()
.ok()
.and_then(|pos| {
let position = parse_hms_to_ms(pos.rel_time.as_deref());
let duration = parse_hms_to_ms(pos.track_duration.as_deref());
Some((position, duration))
})
.unwrap_or((None, None));
// Volume et mute
let volume = renderer.volume().ok().and_then(|v| u8::try_from(v).ok());
let mute = renderer.mute().ok();
// Volume et mute
let volume = renderer_clone.volume().ok().and_then(|v| u8::try_from(v).ok());
let mute = renderer_clone.mute().ok();
// Queue
(transport_state, position_ms, duration_ms, volume, mute)
});
let (transport_state, position_ms, duration_ms, volume, mute) =
time::timeout(STATE_QUERY_TIMEOUT, state_task)
.await
.map_err(|_| {
warn!(
"State query for renderer {} exceeded {:?}",
renderer_id, STATE_QUERY_TIMEOUT
);
(
StatusCode::GATEWAY_TIMEOUT,
Json(ErrorResponse {
error: format!(
"State query timed out after {}s",
STATE_QUERY_TIMEOUT.as_secs()
),
}),
)
})?
.map_err(|e| {
warn!("Task join error during state query: {}", e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Internal task error: {}", e),
}),
)
})?;
// Queue (non-blocking, local data)
let queue_len = state
.control_point
.get_queue_snapshot(&rid)
@@ -150,15 +211,17 @@ async fn get_renderer_state(
.map(|q| q.len())
.unwrap_or(0);
// Playlist binding
// Playlist binding (non-blocking, local data)
let attached_playlist = state
.control_point
.current_queue_playlist_binding(&rid)
.map(|(server_id, container_id, has_seen_update)| AttachedPlaylistInfo {
server_id: server_id.0,
container_id,
has_seen_update,
});
.map(
|(server_id, container_id, has_seen_update)| AttachedPlaylistInfo {
server_id: server_id.0,
container_id,
has_seen_update,
},
);
Ok(Json(RendererState {
id: info.id.0.clone(),
@@ -197,17 +260,18 @@ async fn get_renderer_queue(
) -> Result<Json<QueueSnapshot>, (StatusCode, Json<ErrorResponse>)> {
let rid = RendererId(renderer_id.clone());
let (items, current_index) = state
.control_point
.get_full_queue_snapshot(&rid)
.map_err(|e| {
(
StatusCode::NOT_FOUND,
Json(ErrorResponse {
error: format!("Failed to get queue: {}", e),
}),
)
})?;
let (items, current_index) =
state
.control_point
.get_full_queue_snapshot(&rid)
.map_err(|e| {
(
StatusCode::NOT_FOUND,
Json(ErrorResponse {
error: format!("Failed to get queue: {}", e),
}),
)
})?;
let queue_items: Vec<QueueItem> = items
.into_iter()
@@ -253,11 +317,13 @@ async fn get_renderer_binding(
let binding = state
.control_point
.current_queue_playlist_binding(&rid)
.map(|(server_id, container_id, has_seen_update)| AttachedPlaylistInfo {
server_id: server_id.0,
container_id,
has_seen_update,
});
.map(
|(server_id, container_id, has_seen_update)| AttachedPlaylistInfo {
server_id: server_id.0,
container_id,
has_seen_update,
},
);
Ok(Json(binding))
}
@@ -299,15 +365,44 @@ async fn play_renderer(
)
})?;
renderer.play().map_err(|e| {
warn!("Failed to play renderer {}: {}", renderer_id, e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Failed to play: {}", e),
}),
)
})?;
let renderer_clone = renderer.clone();
let play_task = tokio::task::spawn_blocking(move || renderer_clone.play());
time::timeout(TRANSPORT_COMMAND_TIMEOUT, play_task)
.await
.map_err(|_| {
warn!(
"Play command for renderer {} exceeded {:?}",
renderer_id, TRANSPORT_COMMAND_TIMEOUT
);
(
StatusCode::GATEWAY_TIMEOUT,
Json(ErrorResponse {
error: format!(
"Play command timed out after {}s",
TRANSPORT_COMMAND_TIMEOUT.as_secs()
),
}),
)
})?
.map_err(|e| {
warn!("Task join error during play: {}", e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Internal task error: {}", e),
}),
)
})?
.map_err(|e| {
warn!("Failed to play renderer {}: {}", renderer_id, e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Failed to play: {}", e),
}),
)
})?;
Ok(Json(SuccessResponse {
message: "Playback started".to_string(),
@@ -347,15 +442,44 @@ async fn pause_renderer(
)
})?;
renderer.pause().map_err(|e| {
warn!("Failed to pause renderer {}: {}", renderer_id, e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Failed to pause: {}", e),
}),
)
})?;
let renderer_clone = renderer.clone();
let pause_task = tokio::task::spawn_blocking(move || renderer_clone.pause());
time::timeout(TRANSPORT_COMMAND_TIMEOUT, pause_task)
.await
.map_err(|_| {
warn!(
"Pause command for renderer {} exceeded {:?}",
renderer_id, TRANSPORT_COMMAND_TIMEOUT
);
(
StatusCode::GATEWAY_TIMEOUT,
Json(ErrorResponse {
error: format!(
"Pause command timed out after {}s",
TRANSPORT_COMMAND_TIMEOUT.as_secs()
),
}),
)
})?
.map_err(|e| {
warn!("Task join error during pause: {}", e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Internal task error: {}", e),
}),
)
})?
.map_err(|e| {
warn!("Failed to pause renderer {}: {}", renderer_id, e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Failed to pause: {}", e),
}),
)
})?;
Ok(Json(SuccessResponse {
message: "Playback paused".to_string(),
@@ -395,15 +519,44 @@ async fn stop_renderer(
)
})?;
renderer.stop().map_err(|e| {
warn!("Failed to stop renderer {}: {}", renderer_id, e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Failed to stop: {}", e),
}),
)
})?;
let renderer_clone = renderer.clone();
let stop_task = tokio::task::spawn_blocking(move || renderer_clone.stop());
time::timeout(TRANSPORT_COMMAND_TIMEOUT, stop_task)
.await
.map_err(|_| {
warn!(
"Stop command for renderer {} exceeded {:?}",
renderer_id, TRANSPORT_COMMAND_TIMEOUT
);
(
StatusCode::GATEWAY_TIMEOUT,
Json(ErrorResponse {
error: format!(
"Stop command timed out after {}s",
TRANSPORT_COMMAND_TIMEOUT.as_secs()
),
}),
)
})?
.map_err(|e| {
warn!("Task join error during stop: {}", e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Internal task error: {}", e),
}),
)
})?
.map_err(|e| {
warn!("Failed to stop renderer {}: {}", renderer_id, e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Failed to stop: {}", e),
}),
)
})?;
Ok(Json(SuccessResponse {
message: "Playback stopped".to_string(),
@@ -430,12 +583,44 @@ async fn next_renderer(
Path(renderer_id): Path<String>,
) -> Result<Json<SuccessResponse>, (StatusCode, Json<ErrorResponse>)> {
let rid = RendererId(renderer_id.clone());
let control_point = state.control_point.clone();
let rid_clone = rid.clone();
state
.control_point
.play_next_from_queue(&rid)
let next_task = tokio::task::spawn_blocking(move || {
control_point.play_next_from_queue(&rid_clone)
});
time::timeout(QUEUE_COMMAND_TIMEOUT, next_task)
.await
.map_err(|_| {
warn!(
"Next command for renderer {} exceeded {:?}",
renderer_id, QUEUE_COMMAND_TIMEOUT
);
(
StatusCode::GATEWAY_TIMEOUT,
Json(ErrorResponse {
error: format!(
"Next command timed out after {}s",
QUEUE_COMMAND_TIMEOUT.as_secs()
),
}),
)
})?
.map_err(|e| {
warn!("Failed to advance queue for renderer {}: {}", renderer_id, e);
warn!("Task join error during next: {}", e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Internal task error: {}", e),
}),
)
})?
.map_err(|e| {
warn!(
"Failed to advance queue for renderer {}: {}",
renderer_id, e
);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
@@ -489,18 +674,50 @@ async fn set_renderer_volume(
)
})?;
renderer.set_volume(req.volume as u16).map_err(|e| {
warn!("Failed to set volume for renderer {}: {}", renderer_id, e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Failed to set volume: {}", e),
}),
)
})?;
let renderer_clone = renderer.clone();
let volume = req.volume;
let volume_task = tokio::task::spawn_blocking(move || {
renderer_clone.set_volume(volume as u16)
});
time::timeout(VOLUME_COMMAND_TIMEOUT, volume_task)
.await
.map_err(|_| {
warn!(
"Set volume command for renderer {} exceeded {:?}",
renderer_id, VOLUME_COMMAND_TIMEOUT
);
(
StatusCode::GATEWAY_TIMEOUT,
Json(ErrorResponse {
error: format!(
"Set volume command timed out after {}s",
VOLUME_COMMAND_TIMEOUT.as_secs()
),
}),
)
})?
.map_err(|e| {
warn!("Task join error during set volume: {}", e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Internal task error: {}", e),
}),
)
})?
.map_err(|e| {
warn!("Failed to set volume for renderer {}: {}", renderer_id, e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Failed to set volume: {}", e),
}),
)
})?;
Ok(Json(SuccessResponse {
message: format!("Volume set to {}", req.volume),
message: format!("Volume set to {}", volume),
}))
}
@@ -537,25 +754,52 @@ async fn volume_up_renderer(
)
})?;
let current = renderer.volume().map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Failed to get current volume: {}", e),
}),
)
})?;
let renderer_clone = renderer.clone();
let volume_task = tokio::task::spawn_blocking(move || {
let current = renderer_clone.volume()?;
let new_volume = (current + 5).min(100);
renderer_clone.set_volume(new_volume)?;
Ok::<u16, anyhow::Error>(new_volume)
});
let new_volume = (current + 5).min(100);
renderer.set_volume(new_volume).map_err(|e| {
warn!("Failed to increase volume for renderer {}: {}", renderer_id, e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Failed to increase volume: {}", e),
}),
)
})?;
let new_volume = time::timeout(VOLUME_COMMAND_TIMEOUT, volume_task)
.await
.map_err(|_| {
warn!(
"Volume up command for renderer {} exceeded {:?}",
renderer_id, VOLUME_COMMAND_TIMEOUT
);
(
StatusCode::GATEWAY_TIMEOUT,
Json(ErrorResponse {
error: format!(
"Volume up command timed out after {}s",
VOLUME_COMMAND_TIMEOUT.as_secs()
),
}),
)
})?
.map_err(|e| {
warn!("Task join error during volume up: {}", e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Internal task error: {}", e),
}),
)
})?
.map_err(|e| {
warn!(
"Failed to increase volume for renderer {}: {}",
renderer_id, e
);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Failed to increase volume: {}", e),
}),
)
})?;
Ok(Json(SuccessResponse {
message: format!("Volume increased to {}", new_volume),
@@ -595,25 +839,52 @@ async fn volume_down_renderer(
)
})?;
let current = renderer.volume().map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Failed to get current volume: {}", e),
}),
)
})?;
let renderer_clone = renderer.clone();
let volume_task = tokio::task::spawn_blocking(move || {
let current = renderer_clone.volume()?;
let new_volume = current.saturating_sub(5);
renderer_clone.set_volume(new_volume)?;
Ok::<u16, anyhow::Error>(new_volume)
});
let new_volume = current.saturating_sub(5);
renderer.set_volume(new_volume).map_err(|e| {
warn!("Failed to decrease volume for renderer {}: {}", renderer_id, e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Failed to decrease volume: {}", e),
}),
)
})?;
let new_volume = time::timeout(VOLUME_COMMAND_TIMEOUT, volume_task)
.await
.map_err(|_| {
warn!(
"Volume down command for renderer {} exceeded {:?}",
renderer_id, VOLUME_COMMAND_TIMEOUT
);
(
StatusCode::GATEWAY_TIMEOUT,
Json(ErrorResponse {
error: format!(
"Volume down command timed out after {}s",
VOLUME_COMMAND_TIMEOUT.as_secs()
),
}),
)
})?
.map_err(|e| {
warn!("Task join error during volume down: {}", e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Internal task error: {}", e),
}),
)
})?
.map_err(|e| {
warn!(
"Failed to decrease volume for renderer {}: {}",
renderer_id, e
);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Failed to decrease volume: {}", e),
}),
)
})?;
Ok(Json(SuccessResponse {
message: format!("Volume decreased to {}", new_volume),
@@ -653,25 +924,49 @@ async fn toggle_mute_renderer(
)
})?;
let current_mute = renderer.mute().map_err(|e| {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Failed to get current mute state: {}", e),
}),
)
})?;
let renderer_clone = renderer.clone();
let mute_task = tokio::task::spawn_blocking(move || {
let current_mute = renderer_clone.mute()?;
let new_mute = !current_mute;
renderer_clone.set_mute(new_mute)?;
Ok::<bool, anyhow::Error>(new_mute)
});
let new_mute = !current_mute;
renderer.set_mute(new_mute).map_err(|e| {
warn!("Failed to toggle mute for renderer {}: {}", renderer_id, e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Failed to toggle mute: {}", e),
}),
)
})?;
let new_mute = time::timeout(VOLUME_COMMAND_TIMEOUT, mute_task)
.await
.map_err(|_| {
warn!(
"Toggle mute command for renderer {} exceeded {:?}",
renderer_id, VOLUME_COMMAND_TIMEOUT
);
(
StatusCode::GATEWAY_TIMEOUT,
Json(ErrorResponse {
error: format!(
"Toggle mute command timed out after {}s",
VOLUME_COMMAND_TIMEOUT.as_secs()
),
}),
)
})?
.map_err(|e| {
warn!("Task join error during toggle mute: {}", e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Internal task error: {}", e),
}),
)
})?
.map_err(|e| {
warn!("Failed to toggle mute for renderer {}: {}", renderer_id, e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Failed to toggle mute: {}", e),
}),
)
})?;
Ok(Json(SuccessResponse {
message: format!("Mute {}", if new_mute { "enabled" } else { "disabled" }),
@@ -704,10 +999,40 @@ async fn attach_playlist_binding(
) -> Result<Json<SuccessResponse>, (StatusCode, Json<ErrorResponse>)> {
let rid = RendererId(renderer_id.clone());
let sid = ServerId(req.server_id.clone());
let container_id = req.container_id.clone();
let control_point = Arc::clone(&state.control_point);
state
.control_point
.attach_queue_to_playlist(&rid, sid, req.container_id.clone());
// Spawn blocking task and wait for completion with timeout
let attach_task = tokio::task::spawn_blocking(move || {
control_point.attach_queue_to_playlist(&rid, sid, container_id);
});
time::timeout(ATTACH_PLAYLIST_TIMEOUT, attach_task)
.await
.map_err(|_| {
warn!(
"Attach playlist for renderer {} exceeded {:?}",
renderer_id, ATTACH_PLAYLIST_TIMEOUT
);
(
StatusCode::GATEWAY_TIMEOUT,
Json(ErrorResponse {
error: format!(
"Attach playlist timed out after {}s",
ATTACH_PLAYLIST_TIMEOUT.as_secs()
),
}),
)
})?
.map_err(|e| {
warn!("Task join error during attach playlist: {}", e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Internal task error: {}", e),
}),
)
})?;
debug!(
renderer = renderer_id.as_str(),
@@ -717,10 +1042,7 @@ async fn attach_playlist_binding(
);
Ok(Json(SuccessResponse {
message: format!(
"Playlist {} attached to renderer",
req.container_id
),
message: format!("Playlist {} attached to renderer", req.container_id),
}))
}
@@ -769,9 +1091,7 @@ async fn detach_playlist_binding(
),
tag = "control"
)]
async fn list_servers(
State(state): State<ControlPointState>,
) -> Json<Vec<MediaServerSummary>> {
async fn list_servers(State(state): State<ControlPointState>) -> Json<Vec<MediaServerSummary>> {
let servers = state.control_point.list_media_servers();
let summaries: Vec<MediaServerSummary> = servers
@@ -809,17 +1129,14 @@ async fn browse_container(
) -> Result<Json<BrowseResponse>, (StatusCode, Json<ErrorResponse>)> {
let sid = ServerId(server_id.clone());
let server_info = state
.control_point
.media_server(&sid)
.ok_or_else(|| {
(
StatusCode::NOT_FOUND,
Json(ErrorResponse {
error: format!("Server {} not found", server_id),
}),
)
})?;
let server_info = state.control_point.media_server(&sid).ok_or_else(|| {
(
StatusCode::NOT_FOUND,
Json(ErrorResponse {
error: format!("Server {} not found", server_id),
}),
)
})?;
if !server_info.online {
return Err((
@@ -839,8 +1156,8 @@ async fn browse_container(
));
}
let music_server = MusicServer::from_info(&server_info, Duration::from_secs(10)).map_err(
|e| {
let music_server =
MusicServer::from_info(&server_info, MEDIA_SERVER_SOAP_TIMEOUT).map_err(|e| {
warn!("Failed to create MusicServer for {}: {}", server_id, e);
(
StatusCode::INTERNAL_SERVER_ERROR,
@@ -848,11 +1165,40 @@ async fn browse_container(
error: format!("Failed to initialize server: {}", e),
}),
)
},
)?;
})?;
let entries = music_server
.browse_children(&container_id, 0, 100)
// Use spawn_blocking to avoid blocking the async runtime with synchronous SOAP calls
let container_id_clone = container_id.clone();
let browse_task = tokio::task::spawn_blocking(move || {
music_server.browse_children(&container_id_clone, 0, BROWSE_PAGE_SIZE)
});
let entries = time::timeout(BROWSE_REQUEST_TIMEOUT, browse_task)
.await
.map_err(|_| {
warn!(
"Browse request for container {} on server {} exceeded {:?}",
container_id, server_id, BROWSE_REQUEST_TIMEOUT
);
(
StatusCode::GATEWAY_TIMEOUT,
Json(ErrorResponse {
error: format!(
"Browse request timed out after {}s",
BROWSE_REQUEST_TIMEOUT.as_secs()
),
}),
)
})?
.map_err(|e| {
warn!("Task join error during browse: {}", e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Internal task error: {}", e),
}),
)
})?
.map_err(|e| {
warn!(
"Failed to browse container {} on server {}: {}",
@@ -939,23 +1285,47 @@ pub fn create_api_router(state: ControlPointState, control_point: Arc<ControlPoi
.route("/renderers", get(list_renderers))
.route("/renderers/{renderer_id}", get(get_renderer_state))
.route("/renderers/{renderer_id}/queue", get(get_renderer_queue))
.route("/renderers/{renderer_id}/binding", get(get_renderer_binding))
.route(
"/renderers/{renderer_id}/binding",
get(get_renderer_binding),
)
// Transport control
.route("/renderers/{renderer_id}/play", post(play_renderer))
.route("/renderers/{renderer_id}/pause", post(pause_renderer))
.route("/renderers/{renderer_id}/stop", post(stop_renderer))
.route("/renderers/{renderer_id}/next", post(next_renderer))
// Volume control
.route("/renderers/{renderer_id}/volume/set", post(set_renderer_volume))
.route("/renderers/{renderer_id}/volume/up", post(volume_up_renderer))
.route("/renderers/{renderer_id}/volume/down", post(volume_down_renderer))
.route("/renderers/{renderer_id}/mute/toggle", post(toggle_mute_renderer))
.route(
"/renderers/{renderer_id}/volume/set",
post(set_renderer_volume),
)
.route(
"/renderers/{renderer_id}/volume/up",
post(volume_up_renderer),
)
.route(
"/renderers/{renderer_id}/volume/down",
post(volume_down_renderer),
)
.route(
"/renderers/{renderer_id}/mute/toggle",
post(toggle_mute_renderer),
)
// Playlist binding
.route("/renderers/{renderer_id}/binding/attach", post(attach_playlist_binding))
.route("/renderers/{renderer_id}/binding/detach", post(detach_playlist_binding))
.route(
"/renderers/{renderer_id}/binding/attach",
post(attach_playlist_binding),
)
.route(
"/renderers/{renderer_id}/binding/detach",
post(detach_playlist_binding),
)
// Servers
.route("/servers", get(list_servers))
.route("/servers/{server_id}/containers/{container_id}", get(browse_container))
.route(
"/servers/{server_id}/containers/{container_id}",
get(browse_container),
)
.with_state(state)
// SSE events - merge the SSE router
.merge(crate::sse::create_sse_router(control_point))
@@ -1019,7 +1389,10 @@ pub trait ControlPointExt {
/// // On peut l'utiliser directement si besoin
/// let renderers = control_point.list_music_renderers();
/// ```
async fn register_control_point(&mut self, timeout_secs: u64) -> std::io::Result<Arc<ControlPoint>>;
async fn register_control_point(
&mut self,
timeout_secs: u64,
) -> std::io::Result<Arc<ControlPoint>>;
/// Initialise l'API Control Point (bas niveau)
///
@@ -1041,7 +1414,10 @@ pub trait ControlPointExt {
#[cfg(feature = "pmoserver")]
#[async_trait]
impl ControlPointExt for pmoserver::Server {
async fn register_control_point(&mut self, timeout_secs: u64) -> std::io::Result<Arc<ControlPoint>> {
async fn register_control_point(
&mut self,
timeout_secs: u64,
) -> std::io::Result<Arc<ControlPoint>> {
use tracing::info;
info!("🎛️ Initializing Control Point...");

View File

@@ -10,20 +10,20 @@
//! - GET /api/control/events/servers - Événements serveurs uniquement
//! - GET /api/control/events - Tous les événements (agrégés)
#[cfg(feature = "pmoserver")]
use crate::PlaybackState;
#[cfg(feature = "pmoserver")]
use crate::control_point::ControlPoint;
#[cfg(feature = "pmoserver")]
use crate::model::{MediaServerEvent, RendererEvent};
#[cfg(feature = "pmoserver")]
use crate::PlaybackState;
#[cfg(feature = "pmoserver")]
use async_stream::stream;
#[cfg(feature = "pmoserver")]
use axum::{
extract::State,
response::sse::{Event, KeepAlive, Sse},
response::IntoResponse,
Router,
extract::State,
response::IntoResponse,
response::sse::{Event, KeepAlive, Sse},
};
#[cfg(feature = "pmoserver")]
use serde::Serialize;
@@ -276,9 +276,7 @@ pub async fn media_server_events_sse(
),
tag = "control"
)]
pub async fn all_events_sse(
State(control_point): State<Arc<ControlPoint>>,
) -> impl IntoResponse {
pub async fn all_events_sse(State(control_point): State<Arc<ControlPoint>>) -> impl IntoResponse {
// Convert crossbeam channels to tokio channels for async compatibility
let (renderer_tx, mut renderer_rx_tokio) = tokio::sync::mpsc::unbounded_channel();
let (server_tx, mut server_rx_tokio) = tokio::sync::mpsc::unbounded_channel();

View File

@@ -118,7 +118,7 @@ pub struct Container {
#[serde(rename = "@id")]
pub id: String,
#[serde(rename = "@parentID")]
#[serde(rename = "@parentID", default)]
pub parent_id: String,
#[serde(rename = "@restricted", skip_serializing_if = "Option::is_none")]
@@ -133,7 +133,7 @@ pub struct Container {
#[serde(rename = "dc:title", alias = "title")]
pub title: String,
#[serde(rename = "upnp:class", alias = "class")]
#[serde(rename = "upnp:class", alias = "class", default)]
pub class: String,
#[serde(rename = "container", default)]
@@ -149,7 +149,7 @@ pub struct Item {
#[serde(rename = "@id")]
pub id: String,
#[serde(rename = "@parentID")]
#[serde(rename = "@parentID", default)]
pub parent_id: String,
#[serde(rename = "@restricted", skip_serializing_if = "Option::is_none")]
@@ -165,7 +165,7 @@ pub struct Item {
)]
pub creator: Option<String>,
#[serde(rename = "upnp:class", alias = "class")]
#[serde(rename = "upnp:class", alias = "class", default)]
pub class: String,
#[serde(
@@ -238,7 +238,7 @@ pub struct Resource {
#[serde(rename = "@duration", skip_serializing_if = "Option::is_none")]
pub duration: Option<String>,
#[serde(rename = "$text")]
#[serde(rename = "$text", default)]
pub url: String,
}