//! Extension pmoserver pour le Control Point //! //! Ce module fournit une API REST pour contrôler les renderers UPnP //! et naviguer dans les serveurs de médias. #[cfg(feature = "pmoserver")] use crate::control_point::ControlPoint; #[cfg(feature = "pmoserver")] use crate::media_server::{MediaBrowser, playback_item_from_entry}; #[cfg(feature = "pmoserver")] use crate::MediaEntry; #[cfg(feature = "pmoserver")] use crate::model::{RendererCapabilities, RendererProtocol}; #[cfg(feature = "pmoserver")] use crate::openapi::{ AttachPlaylistRequest, AttachedPlaylistInfo, BrowseResponse, ContainerEntry, ErrorResponse, FullRendererSnapshot, MediaServerSummary, PlayContentRequest, QueueSnapshot, RendererCapabilitiesSummary, RendererProtocolSummary, RendererState, RendererSummary, SeekQueueRequest, SeekRequest, SleepTimerRequest, SleepTimerState, StreamState, SuccessResponse, TransferQueueRequest, VolumeSetRequest, }; #[cfg(feature = "pmoserver")] use crate::queue::PlaybackItem; #[cfg(feature = "pmoserver")] use crate::{DeviceId, DeviceIdentity, DeviceOnline}; #[cfg(feature = "pmoserver")] use pmocovers; #[cfg(feature = "pmoserver")] use async_trait::async_trait; #[cfg(feature = "pmoserver")] use axum::{ Json, Router, extract::{Path, Query, State}, http::{StatusCode, header::HeaderMap}, routing::{get, post}, }; #[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_DEFAULT_LIMIT: u32 = 50; #[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); /// État partagé pour l'API ControlPoint #[cfg(feature = "pmoserver")] #[derive(Clone)] pub struct ControlPointState { control_point: Arc, } #[cfg(feature = "pmoserver")] impl ControlPointState { pub fn new(control_point: Arc) -> Self { Self { control_point } } } // ============================================================================ // HANDLERS - RENDERERS // ============================================================================ /// GET /control/renderers - Liste tous les renderers #[cfg(feature = "pmoserver")] #[utoipa::path( get, path = "/renderers", responses( (status = 200, description = "Liste des renderers", body = Vec) ), tag = "control" )] async fn list_renderers(State(state): State) -> Json> { // Use spawn_blocking to avoid blocking the tokio runtime // This is critical because list_music_renderers acquires a RwLock let control_point = state.control_point.clone(); let summaries = tokio::task::spawn_blocking(move || { let renderers = control_point.list_music_renderers(); renderers .into_iter() .map(|r| { let info = r.info(); // Extraire le host depuis la location URL (ex: "http://192.168.1.10:8080/...") let server_host = info.location() .split("://").nth(1) .and_then(|s| s.split('/').next()) .and_then(|s| s.split(':').next()) .map(|h| h.to_string()) .filter(|h| !h.is_empty()); RendererSummary { id: r.id().0.clone(), friendly_name: r.friendly_name().to_string(), model_name: r.model_name().to_string(), protocol: protocol_summary(&info.protocol()), capabilities: capability_summary(&info.capabilities()), online: r.is_online(), server_host, } }) .collect::>() }) .await .unwrap_or_default(); Json(summaries) } /// GET /control/renderers/{renderer_id} - Récupère l'état d'un renderer #[cfg(feature = "pmoserver")] #[utoipa::path( get, path = "/renderers/{renderer_id}", params( ("renderer_id" = String, Path, description = "ID unique du renderer") ), responses( (status = 200, description = "État du renderer", body = RendererState), (status = 404, description = "Renderer non trouvé", body = ErrorResponse) ), tag = "control" )] async fn get_renderer_state( State(state): State, Path(renderer_id): Path, ) -> Result, (StatusCode, Json)> { let rid = DeviceId(renderer_id.clone()); let snapshot = state .control_point .renderer_full_snapshot(&rid) .map_err(|err| map_snapshot_error(renderer_id, err))?; Ok(Json(snapshot.state)) } #[cfg(feature = "pmoserver")] #[utoipa::path( get, path = "/renderers/{renderer_id}/full", params( ("renderer_id" = String, Path, description = "ID unique du renderer") ), responses( (status = 200, description = "Snapshot complet du renderer", body = FullRendererSnapshot), (status = 404, description = "Renderer non trouvé", body = ErrorResponse) ), tag = "control" )] async fn get_renderer_full_snapshot( State(state): State, Path(renderer_id): Path, headers: HeaderMap, ) -> Result, (StatusCode, Json)> { let rid = DeviceId(renderer_id.clone()); // Use spawn_blocking because renderer_full_snapshot does sync UPnP calls let control_point = state.control_point.clone(); let rid_clone = rid.clone(); let mut snapshot = tokio::task::spawn_blocking(move || control_point.renderer_full_snapshot(&rid_clone)) .await .map_err(|e| { ( StatusCode::INTERNAL_SERVER_ERROR, Json(ErrorResponse { error: format!("Task error: {}", e), }), ) })? .map_err(|err| map_snapshot_error(renderer_id, err))?; // Get base_url from request headers for transforming cover URLs let base_url_str = pmoserver::get_base_url_from_request(&headers); let base_url = pmoserver::BaseUrl(base_url_str); // Transform cover URLs in current_track if let Some(ref mut current_track) = snapshot.state.current_track { if let Some(ref album_art) = current_track.album_art_uri { if let Some(transformed) = transform_cover_url(Some(album_art), &base_url).await { current_track.album_art_uri = Some(transformed); } } } // Transform cover URLs in queue items for item in &mut snapshot.queue.items { if let Some(ref album_art) = item.album_art_uri { if let Some(transformed) = transform_cover_url(Some(album_art), &base_url).await { item.album_art_uri = Some(transformed); } } } Ok(Json(snapshot)) } /// GET /control/renderers/{renderer_id}/queue - Récupère la queue d'un renderer #[cfg(feature = "pmoserver")] #[utoipa::path( get, path = "/renderers/{renderer_id}/queue", params( ("renderer_id" = String, Path, description = "ID unique du renderer") ), responses( ( status = 200, description = "Playlist complète du renderer (avec index courant)", body = QueueSnapshot ), (status = 404, description = "Renderer non trouvé", body = ErrorResponse) ), tag = "control" )] async fn get_renderer_queue( State(state): State, Path(renderer_id): Path, headers: HeaderMap, ) -> Result, (StatusCode, Json)> { let rid = DeviceId(renderer_id.clone()); let mut snapshot = state .control_point .renderer_full_snapshot(&rid) .map_err(|err| map_snapshot_error(renderer_id, err))?; // Get base_url from request headers let base_url_str = pmoserver::get_base_url_from_request(&headers); let base_url = pmoserver::BaseUrl(base_url_str.clone()); // Transform cover URLs in all queue items using async version for item in &mut snapshot.queue.items { if let Some(ref album_art) = item.album_art_uri { if let Some(transformed) = transform_cover_url(Some(album_art), &base_url).await { item.album_art_uri = Some(transformed); } } } Ok(Json(snapshot.queue)) } /// GET /control/renderers/{renderer_id}/binding - Récupère le binding playlist #[cfg(feature = "pmoserver")] #[utoipa::path( get, path = "/renderers/{renderer_id}/binding", params( ("renderer_id" = String, Path, description = "ID unique du renderer") ), responses( (status = 200, description = "Binding playlist", body = Option), (status = 404, description = "Renderer non trouvé", body = ErrorResponse) ), tag = "control" )] async fn get_renderer_binding( State(state): State, Path(renderer_id): Path, ) -> Result>, (StatusCode, Json)> { let rid = DeviceId(renderer_id.clone()); let snapshot = state .control_point .renderer_full_snapshot(&rid) .map_err(|err| map_snapshot_error(renderer_id, err))?; Ok(Json(snapshot.binding)) } // ============================================================================ // HANDLERS - TRANSPORT CONTROLS // ============================================================================ /// POST /control/renderers/{renderer_id}/play - Démarre la lecture #[cfg(feature = "pmoserver")] #[utoipa::path( post, path = "/renderers/{renderer_id}/play", params( ("renderer_id" = String, Path, description = "ID unique du renderer") ), responses( (status = 200, description = "Lecture démarrée", body = SuccessResponse), (status = 404, description = "Renderer non trouvé", body = ErrorResponse), (status = 500, description = "Erreur lors de l'exécution", body = ErrorResponse) ), tag = "control" )] async fn play_renderer( State(state): State, Path(renderer_id): Path, ) -> Result, (StatusCode, Json)> { let rid = DeviceId(renderer_id.clone()); let renderer = state .control_point .music_renderer_by_id(&rid) .ok_or_else(|| { ( StatusCode::NOT_FOUND, Json(ErrorResponse { error: format!("Renderer {} not found", renderer_id), }), ) })?; 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(), })) } /// POST /control/renderers/{renderer_id}/pause - Met en pause #[cfg(feature = "pmoserver")] #[utoipa::path( post, path = "/renderers/{renderer_id}/pause", params( ("renderer_id" = String, Path, description = "ID unique du renderer") ), responses( (status = 200, description = "Lecture en pause", body = SuccessResponse), (status = 404, description = "Renderer non trouvé", body = ErrorResponse), (status = 500, description = "Erreur lors de l'exécution", body = ErrorResponse) ), tag = "control" )] async fn pause_renderer( State(state): State, Path(renderer_id): Path, ) -> Result, (StatusCode, Json)> { let rid = DeviceId(renderer_id.clone()); let renderer = state .control_point .music_renderer_by_id(&rid) .ok_or_else(|| { ( StatusCode::NOT_FOUND, Json(ErrorResponse { error: format!("Renderer {} not found", renderer_id), }), ) })?; 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(), })) } /// POST /control/renderers/{renderer_id}/stop - Arrête la lecture #[cfg(feature = "pmoserver")] #[utoipa::path( post, path = "/renderers/{renderer_id}/stop", params( ("renderer_id" = String, Path, description = "ID unique du renderer") ), responses( (status = 200, description = "Lecture arrêtée", body = SuccessResponse), (status = 404, description = "Renderer non trouvé", body = ErrorResponse), (status = 500, description = "Erreur lors de l'exécution", body = ErrorResponse) ), tag = "control" )] async fn stop_renderer( State(state): State, Path(renderer_id): Path, ) -> Result, (StatusCode, Json)> { let rid = DeviceId(renderer_id.clone()); state .control_point .music_renderer_by_id(&rid) .ok_or_else(|| { ( StatusCode::NOT_FOUND, Json(ErrorResponse { error: format!("Renderer {} not found", renderer_id), }), ) })?; let control_point = Arc::clone(&state.control_point); let rid_for_task = rid.clone(); let stop_task = tokio::task::spawn_blocking(move || control_point.user_stop(&rid_for_task)); 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(), })) } /// POST /control/renderers/{renderer_id}/resume - Reprend la lecture depuis la queue #[cfg(feature = "pmoserver")] #[utoipa::path( post, path = "/renderers/{renderer_id}/resume", params( ("renderer_id" = String, Path, description = "ID unique du renderer") ), responses( (status = 200, description = "Lecture reprise", body = SuccessResponse), (status = 404, description = "Renderer non trouvé", body = ErrorResponse), (status = 500, description = "Erreur lors de l'exécution", body = ErrorResponse) ), tag = "control" )] async fn resume_renderer( State(state): State, Path(renderer_id): Path, ) -> Result, (StatusCode, Json)> { let rid = DeviceId(renderer_id.clone()); state .control_point .music_renderer_by_id(&rid) .ok_or_else(|| { ( StatusCode::NOT_FOUND, Json(ErrorResponse { error: format!("Renderer {} not found", renderer_id), }), ) })?; let control_point = Arc::clone(&state.control_point); let rid_for_task = rid.clone(); let resume_task = tokio::task::spawn_blocking(move || control_point.play_current_from_queue(&rid_for_task)); time::timeout(TRANSPORT_COMMAND_TIMEOUT, resume_task) .await .map_err(|_| { warn!( "Resume command for renderer {} exceeded {:?}", renderer_id, TRANSPORT_COMMAND_TIMEOUT ); ( StatusCode::GATEWAY_TIMEOUT, Json(ErrorResponse { error: format!( "Resume command timed out after {}s", TRANSPORT_COMMAND_TIMEOUT.as_secs() ), }), ) })? .map_err(|e| { warn!("Task join error during resume: {}", e); ( StatusCode::INTERNAL_SERVER_ERROR, Json(ErrorResponse { error: format!("Internal task error: {}", e), }), ) })? .map_err(|e| { warn!( "Failed to resume playback for renderer {}: {}", renderer_id, e ); ( StatusCode::INTERNAL_SERVER_ERROR, Json(ErrorResponse { error: format!("Failed to resume playback: {}", e), }), ) })?; Ok(Json(SuccessResponse { message: "Playback resumed".to_string(), })) } /// POST /control/renderers/{renderer_id}/next - Passe au morceau suivant #[cfg(feature = "pmoserver")] #[utoipa::path( post, path = "/renderers/{renderer_id}/next", params( ("renderer_id" = String, Path, description = "ID unique du renderer") ), responses( (status = 200, description = "Piste suivante lancée", body = SuccessResponse), (status = 404, description = "Renderer non trouvé", body = ErrorResponse), (status = 500, description = "Erreur lors de l'exécution", body = ErrorResponse) ), tag = "control" )] async fn next_renderer( State(state): State, Path(renderer_id): Path, ) -> Result, (StatusCode, Json)> { let rid = DeviceId(renderer_id.clone()); state .control_point .music_renderer_by_id(&rid) .ok_or_else(|| { ( StatusCode::NOT_FOUND, Json(ErrorResponse { error: format!("Renderer {} not found", renderer_id), }), ) })?; let control_point = Arc::clone(&state.control_point); let rid_for_task = rid.clone(); let next_task = tokio::task::spawn_blocking(move || control_point.play_next_from_queue(&rid_for_task)); time::timeout(TRANSPORT_COMMAND_TIMEOUT, next_task) .await .map_err(|_| { warn!( "Next command for renderer {} exceeded {:?}", renderer_id, TRANSPORT_COMMAND_TIMEOUT ); ( StatusCode::GATEWAY_TIMEOUT, Json(ErrorResponse { error: format!( "Next command timed out after {}s", TRANSPORT_COMMAND_TIMEOUT.as_secs() ), }), ) })? .map_err(|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 skip to next track for renderer {}: {}", renderer_id, e ); ( StatusCode::INTERNAL_SERVER_ERROR, Json(ErrorResponse { error: format!("Failed to skip to next track: {}", e), }), ) })?; Ok(Json(SuccessResponse { message: "Skipped to next track".to_string(), })) } /// POST /control/renderers/{renderer_id}/queue/seek - Saute à un index spécifique dans la queue #[cfg(feature = "pmoserver")] #[utoipa::path( post, path = "/renderers/{renderer_id}/queue/seek", tag = "control", request_body = SeekQueueRequest, responses( (status = 200, description = "Lecture démarrée à l'index spécifié", body = SuccessResponse), (status = 404, description = "Renderer non trouvé", body = ErrorResponse), (status = 400, description = "Index invalide", body = ErrorResponse), (status = 500, description = "Erreur interne", body = ErrorResponse), ) )] async fn seek_queue_index( State(state): State, Path(renderer_id): Path, Json(payload): Json, ) -> Result, (StatusCode, Json)> { let rid = DeviceId(renderer_id.clone()); state .control_point .music_renderer_by_id(&rid) .ok_or_else(|| { ( StatusCode::NOT_FOUND, Json(ErrorResponse { error: format!("Renderer {} not found", renderer_id), }), ) })?; let control_point = Arc::clone(&state.control_point); let rid_for_task = rid.clone(); let index = payload.index; // Launch the command in background and return immediately // The UI will be updated via SSE events when the state changes tokio::task::spawn(async move { let rid_for_log = rid_for_task.clone(); let result = tokio::task::spawn_blocking(move || { control_point.play_queue_index(&rid_for_task, index) }) .await; match result { Ok(Ok(())) => { debug!( "Successfully started playback at index {} for renderer {}", index, rid_for_log.0 ); } Ok(Err(e)) => { warn!( "Failed to seek to index {} for renderer {}: {}", index, rid_for_log.0, e ); } Err(e) => { warn!( "Task join error during queue seek for renderer {}: {}", rid_for_log.0, e ); } } }); Ok(Json(SuccessResponse { message: format!("Playing item at index {}", index), })) } /// POST /control/renderers/{renderer_id}/seek - Seek à une position spécifique (en secondes) #[cfg(feature = "pmoserver")] #[utoipa::path( post, path = "/control/renderers/{renderer_id}/seek", tag = "control", params( ("renderer_id" = String, Path, description = "ID unique du renderer") ), request_body = SeekRequest, responses( (status = 200, description = "Seek effectué avec succès", body = SuccessResponse), (status = 404, description = "Renderer non trouvé", body = ErrorResponse), (status = 504, description = "Timeout lors du seek", body = ErrorResponse), (status = 500, description = "Erreur interne", body = ErrorResponse), ) )] async fn seek_renderer( State(state): State, Path(renderer_id): Path, Json(payload): Json, ) -> Result, (StatusCode, Json)> { let rid = DeviceId(renderer_id.clone()); let renderer = state .control_point .music_renderer_by_id(&rid) .ok_or_else(|| { ( StatusCode::NOT_FOUND, Json(ErrorResponse { error: format!("Renderer {} not found", renderer_id), }), ) })?; let seconds = payload.seconds; let seek_task = tokio::task::spawn_blocking(move || renderer.seek(seconds)); time::timeout(TRANSPORT_COMMAND_TIMEOUT, seek_task) .await .map_err(|_| { warn!( "Seek command for renderer {} exceeded {:?}", renderer_id, TRANSPORT_COMMAND_TIMEOUT ); ( StatusCode::GATEWAY_TIMEOUT, Json(ErrorResponse { error: format!( "Seek command timed out after {}s", TRANSPORT_COMMAND_TIMEOUT.as_secs() ), }), ) })? .map_err(|e| { warn!("Task join error during seek: {}", e); ( StatusCode::INTERNAL_SERVER_ERROR, Json(ErrorResponse { error: format!("Internal task error: {}", e), }), ) })? .map_err(|e| { warn!( "Failed to seek to {} seconds for renderer {}: {}", seconds, renderer_id, e ); ( StatusCode::INTERNAL_SERVER_ERROR, Json(ErrorResponse { error: format!("Failed to seek to {} seconds: {}", seconds, e), }), ) })?; Ok(Json(SuccessResponse { message: format!("Seeked to {} seconds", seconds), })) } /// POST /control/renderers/{renderer_id}/volume/set - Définit le volume #[cfg(feature = "pmoserver")] #[utoipa::path( post, path = "/renderers/{renderer_id}/volume/set", params( ("renderer_id" = String, Path, description = "ID unique du renderer") ), request_body = VolumeSetRequest, responses( (status = 200, description = "Volume défini", body = SuccessResponse), (status = 400, description = "Requête invalide", body = ErrorResponse), (status = 404, description = "Renderer non trouvé", body = ErrorResponse), (status = 500, description = "Erreur lors de l'exécution", body = ErrorResponse) ), tag = "control" )] async fn set_renderer_volume( State(state): State, Path(renderer_id): Path, Json(req): Json, ) -> Result, (StatusCode, Json)> { let rid = DeviceId(renderer_id.clone()); let renderer = state .control_point .music_renderer_by_id(&rid) .ok_or_else(|| { ( StatusCode::NOT_FOUND, Json(ErrorResponse { error: format!("Renderer {} not found", renderer_id), }), ) })?; 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 {}", volume), })) } /// POST /control/renderers/{renderer_id}/volume/up - Augmente le volume #[cfg(feature = "pmoserver")] #[utoipa::path( post, path = "/renderers/{renderer_id}/volume/up", params( ("renderer_id" = String, Path, description = "ID unique du renderer") ), responses( (status = 200, description = "Volume augmenté", body = SuccessResponse), (status = 404, description = "Renderer non trouvé", body = ErrorResponse), (status = 500, description = "Erreur lors de l'exécution", body = ErrorResponse) ), tag = "control" )] async fn volume_up_renderer( State(state): State, Path(renderer_id): Path, ) -> Result, (StatusCode, Json)> { let rid = DeviceId(renderer_id.clone()); let renderer = state .control_point .music_renderer_by_id(&rid) .ok_or_else(|| { ( StatusCode::NOT_FOUND, Json(ErrorResponse { error: format!("Renderer {} not found", renderer_id), }), ) })?; 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::(new_volume) }); 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), })) } /// POST /control/renderers/{renderer_id}/volume/down - Diminue le volume #[cfg(feature = "pmoserver")] #[utoipa::path( post, path = "/renderers/{renderer_id}/volume/down", params( ("renderer_id" = String, Path, description = "ID unique du renderer") ), responses( (status = 200, description = "Volume diminué", body = SuccessResponse), (status = 404, description = "Renderer non trouvé", body = ErrorResponse), (status = 500, description = "Erreur lors de l'exécution", body = ErrorResponse) ), tag = "control" )] async fn volume_down_renderer( State(state): State, Path(renderer_id): Path, ) -> Result, (StatusCode, Json)> { let rid = DeviceId(renderer_id.clone()); let renderer = state .control_point .music_renderer_by_id(&rid) .ok_or_else(|| { ( StatusCode::NOT_FOUND, Json(ErrorResponse { error: format!("Renderer {} not found", renderer_id), }), ) })?; 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::(new_volume) }); 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), })) } /// POST /control/renderers/{renderer_id}/mute/toggle - Bascule le mute #[cfg(feature = "pmoserver")] #[utoipa::path( post, path = "/renderers/{renderer_id}/mute/toggle", params( ("renderer_id" = String, Path, description = "ID unique du renderer") ), responses( (status = 200, description = "Mute basculé", body = SuccessResponse), (status = 404, description = "Renderer non trouvé", body = ErrorResponse), (status = 500, description = "Erreur lors de l'exécution", body = ErrorResponse) ), tag = "control" )] async fn toggle_mute_renderer( State(state): State, Path(renderer_id): Path, ) -> Result, (StatusCode, Json)> { let rid = DeviceId(renderer_id.clone()); let renderer = state .control_point .music_renderer_by_id(&rid) .ok_or_else(|| { ( StatusCode::NOT_FOUND, Json(ErrorResponse { error: format!("Renderer {} not found", renderer_id), }), ) })?; 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::(new_mute) }); 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" }), })) } // ============================================================================ // HANDLERS - SLEEP TIMER // ============================================================================ /// POST /control/renderers/{renderer_id}/timer/start - Démarre le sleep timer #[cfg(feature = "pmoserver")] #[utoipa::path( post, path = "/renderers/{renderer_id}/timer/start", params( ("renderer_id" = String, Path, description = "ID unique du renderer") ), request_body = SleepTimerRequest, responses( (status = 200, description = "Timer démarré", body = SleepTimerState), (status = 404, description = "Renderer non trouvé", body = ErrorResponse), (status = 400, description = "Durée invalide", body = ErrorResponse), (status = 500, description = "Erreur lors de l'exécution", body = ErrorResponse) ), tag = "control" )] async fn start_sleep_timer( State(state): State, Path(renderer_id): Path, Json(req): Json, ) -> Result, (StatusCode, Json)> { let rid = DeviceId(renderer_id.clone()); let renderer = state .control_point .music_renderer_by_id(&rid) .ok_or_else(|| { ( StatusCode::NOT_FOUND, Json(ErrorResponse { error: format!("Renderer {} not found", renderer_id), }), ) })?; // Start the timer let remaining = renderer .start_sleep_timer(req.duration_seconds) .map_err(|e| { warn!( "Failed to start sleep timer for renderer {}: {}", renderer_id, e ); ( StatusCode::BAD_REQUEST, Json(ErrorResponse { error: format!("Failed to start timer: {}", e), }), ) })?; // Note: TimerStarted event will be emitted by the timer watchdog thread Ok(Json(SleepTimerState { active: true, duration_seconds: req.duration_seconds, remaining_seconds: Some(remaining), })) } /// POST /control/renderers/{renderer_id}/timer/update - Met à jour la durée du timer #[cfg(feature = "pmoserver")] #[utoipa::path( post, path = "/renderers/{renderer_id}/timer/update", params( ("renderer_id" = String, Path, description = "ID unique du renderer") ), request_body = SleepTimerRequest, responses( (status = 200, description = "Timer mis à jour", body = SleepTimerState), (status = 404, description = "Renderer non trouvé", body = ErrorResponse), (status = 400, description = "Durée invalide", body = ErrorResponse), (status = 500, description = "Erreur lors de l'exécution", body = ErrorResponse) ), tag = "control" )] async fn update_sleep_timer( State(state): State, Path(renderer_id): Path, Json(req): Json, ) -> Result, (StatusCode, Json)> { let rid = DeviceId(renderer_id.clone()); let renderer = state .control_point .music_renderer_by_id(&rid) .ok_or_else(|| { ( StatusCode::NOT_FOUND, Json(ErrorResponse { error: format!("Renderer {} not found", renderer_id), }), ) })?; // Update the timer let remaining = renderer .update_sleep_timer(req.duration_seconds) .map_err(|e| { warn!( "Failed to update sleep timer for renderer {}: {}", renderer_id, e ); ( StatusCode::BAD_REQUEST, Json(ErrorResponse { error: format!("Failed to update timer: {}", e), }), ) })?; // Note: TimerUpdated event will be emitted by the timer watchdog thread Ok(Json(SleepTimerState { active: true, duration_seconds: req.duration_seconds, remaining_seconds: Some(remaining), })) } /// POST /control/renderers/{renderer_id}/timer/cancel - Annule le sleep timer #[cfg(feature = "pmoserver")] #[utoipa::path( post, path = "/renderers/{renderer_id}/timer/cancel", params( ("renderer_id" = String, Path, description = "ID unique du renderer") ), responses( (status = 200, description = "Timer annulé", body = SuccessResponse), (status = 404, description = "Renderer non trouvé", body = ErrorResponse) ), tag = "control" )] async fn cancel_sleep_timer( State(state): State, Path(renderer_id): Path, ) -> Result, (StatusCode, Json)> { let rid = DeviceId(renderer_id.clone()); let renderer = state .control_point .music_renderer_by_id(&rid) .ok_or_else(|| { ( StatusCode::NOT_FOUND, Json(ErrorResponse { error: format!("Renderer {} not found", renderer_id), }), ) })?; // Cancel the timer renderer.cancel_sleep_timer(); // Note: TimerCancelled event will be emitted by the timer watchdog thread Ok(Json(SuccessResponse { message: "Sleep timer cancelled".to_string(), })) } /// GET /control/renderers/{renderer_id}/timer - Récupère l'état du sleep timer #[cfg(feature = "pmoserver")] #[utoipa::path( get, path = "/renderers/{renderer_id}/timer", params( ("renderer_id" = String, Path, description = "ID unique du renderer") ), responses( (status = 200, description = "État du timer", body = SleepTimerState), (status = 404, description = "Renderer non trouvé", body = ErrorResponse) ), tag = "control" )] async fn get_sleep_timer_state( State(state): State, Path(renderer_id): Path, ) -> Result, (StatusCode, Json)> { let rid = DeviceId(renderer_id.clone()); let renderer = state .control_point .music_renderer_by_id(&rid) .ok_or_else(|| { ( StatusCode::NOT_FOUND, Json(ErrorResponse { error: format!("Renderer {} not found", renderer_id), }), ) })?; let (active, duration_seconds, remaining_seconds) = renderer.sleep_timer_state(); Ok(Json(SleepTimerState { active, duration_seconds, remaining_seconds, })) } // ============================================================================ // HANDLERS - STREAM STATE // ============================================================================ /// GET /control/renderers/{renderer_id}/stream-state - Récupère l'état stream du renderer #[cfg(feature = "pmoserver")] #[utoipa::path( get, path = "/renderers/{renderer_id}/stream-state", params( ("renderer_id" = String, Path, description = "ID unique du renderer") ), responses( (status = 200, description = "État du flux", body = StreamState), (status = 404, description = "Renderer non trouvé", body = ErrorResponse) ), tag = "control" )] async fn get_stream_state( State(state): State, Path(renderer_id): Path, ) -> Result, (StatusCode, Json)> { let rid = DeviceId(renderer_id.clone()); let renderer = state .control_point .music_renderer_by_id(&rid) .ok_or_else(|| { ( StatusCode::NOT_FOUND, Json(ErrorResponse { error: format!("Renderer {} not found", renderer_id), }), ) })?; let is_stream = renderer.is_playing_a_stream(); let is_playing = renderer .playback_state() .map(|s| matches!(s, crate::model::PlaybackState::Playing)) .unwrap_or(false); Ok(Json(StreamState { is_stream, is_playing, })) } // ============================================================================ // HANDLERS - QUEUE SHUFFLE // ============================================================================ /// POST /control/renderers/{renderer_id}/queue/shuffle - Mélange la queue de lecture #[cfg(feature = "pmoserver")] #[utoipa::path( post, path = "/renderers/{renderer_id}/queue/shuffle", params( ("renderer_id" = String, Path, description = "ID unique du renderer") ), responses( (status = 200, description = "Queue mélangée et lecture démarrée", body = SuccessResponse), (status = 404, description = "Renderer non trouvé", body = ErrorResponse), (status = 400, description = "Queue vide", body = ErrorResponse), (status = 504, description = "Timeout de la commande", body = ErrorResponse), (status = 500, description = "Erreur lors de l'exécution", body = ErrorResponse) ), tag = "control" )] async fn shuffle_queue( State(state): State, Path(renderer_id): Path, ) -> Result, (StatusCode, Json)> { let rid = DeviceId(renderer_id.clone()); // Verify renderer exists state .control_point .music_renderer_by_id(&rid) .ok_or_else(|| { ( StatusCode::NOT_FOUND, Json(ErrorResponse { error: format!("Renderer {} not found", renderer_id), }), ) })?; let control_point = Arc::clone(&state.control_point); let rid_for_task = rid.clone(); let shuffle_task = tokio::task::spawn_blocking(move || control_point.shuffle_queue(&rid_for_task)); time::timeout(QUEUE_COMMAND_TIMEOUT, shuffle_task) .await .map_err(|_| { warn!( "Shuffle command for renderer {} exceeded {:?}", renderer_id, QUEUE_COMMAND_TIMEOUT ); ( StatusCode::GATEWAY_TIMEOUT, Json(ErrorResponse { error: format!( "Shuffle command timed out after {}s", QUEUE_COMMAND_TIMEOUT.as_secs() ), }), ) })? .map_err(|e| { warn!("Task join error during shuffle: {}", e); ( StatusCode::INTERNAL_SERVER_ERROR, Json(ErrorResponse { error: format!("Internal task error: {}", e), }), ) })? .map_err(|e| { warn!( "Failed to shuffle queue for renderer {}: {}", renderer_id, e ); ( StatusCode::BAD_REQUEST, Json(ErrorResponse { error: format!("Failed to shuffle queue: {}", e), }), ) })?; debug!( renderer = renderer_id.as_str(), "Queue shuffled via HTTP API" ); Ok(Json(SuccessResponse { message: "Queue shuffled and playback started".to_string(), })) } // ============================================================================ // HANDLERS - BINDING PLAYLIST // ============================================================================ /// POST /control/renderers/{renderer_id}/binding/attach - Attache une playlist #[cfg(feature = "pmoserver")] #[utoipa::path( post, path = "/renderers/{renderer_id}/binding/attach", params( ("renderer_id" = String, Path, description = "ID unique du renderer") ), request_body = AttachPlaylistRequest, responses( (status = 200, description = "Playlist attachée", body = SuccessResponse), (status = 404, description = "Renderer non trouvé", body = ErrorResponse) ), tag = "control" )] async fn attach_playlist_binding( State(state): State, Path(renderer_id): Path, Json(req): Json, ) -> Result, (StatusCode, Json)> { let rid = DeviceId(renderer_id.clone()); let sid = DeviceId(req.server_id.clone()); let container_id = req.container_id.clone(); let control_point = Arc::clone(&state.control_point); // Spawn blocking task and wait for completion with timeout let attach_task = tokio::task::spawn_blocking(move || { control_point.attach_queue_to_playlist_with_options(&rid, sid, container_id, req.auto_play) }); 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), }), ) })? .map_err(|e| { warn!( renderer = renderer_id.as_str(), server = req.server_id.as_str(), container = req.container_id.as_str(), error = %e, "Failed to attach playlist" ); ( StatusCode::INTERNAL_SERVER_ERROR, Json(ErrorResponse { error: format!("Failed to attach playlist: {}", e), }), ) })?; debug!( renderer = renderer_id.as_str(), server = req.server_id.as_str(), container = req.container_id.as_str(), auto_play = req.auto_play, "Playlist attached via HTTP API" ); Ok(Json(SuccessResponse { message: format!("Playlist {} attached to renderer", req.container_id), })) } /// POST /control/renderers/{renderer_id}/binding/detach - Détache la playlist #[cfg(feature = "pmoserver")] #[utoipa::path( post, path = "/renderers/{renderer_id}/binding/detach", params( ("renderer_id" = String, Path, description = "ID unique du renderer") ), responses( (status = 200, description = "Playlist détachée", body = SuccessResponse) ), tag = "control" )] async fn detach_playlist_binding( State(state): State, Path(renderer_id): Path, ) -> Json { let rid = DeviceId(renderer_id.clone()); state.control_point.detach_queue_playlist(&rid); debug!( renderer = renderer_id.as_str(), "Playlist detached via HTTP API" ); Json(SuccessResponse { message: "Playlist detached".to_string(), }) } // ============================================================================ // HANDLERS - QUEUE CONTENT // ============================================================================ /// POST /control/renderers/{renderer_id}/queue/play - Lire du contenu immédiatement #[cfg(feature = "pmoserver")] #[utoipa::path( post, path = "/renderers/{renderer_id}/queue/play", params( ("renderer_id" = String, Path, description = "ID unique du renderer") ), request_body = PlayContentRequest, responses( (status = 200, description = "Contenu en cours de lecture", body = SuccessResponse), (status = 404, description = "Renderer ou serveur non trouvé", body = ErrorResponse), (status = 500, description = "Erreur lors de l'exécution", body = ErrorResponse) ), tag = "control" )] async fn play_content( State(state): State, Path(renderer_id): Path, Json(req): Json, ) -> Result, (StatusCode, Json)> { let rid = DeviceId(renderer_id.clone()); let sid = DeviceId(req.server_id.clone()); let object_id = req.object_id.clone(); let object_id_for_log = object_id.clone(); // Get renderer to verify it exists state .control_point .music_renderer_by_id(&rid) .ok_or_else(|| { ( StatusCode::NOT_FOUND, Json(ErrorResponse { error: format!("Renderer {} not found", renderer_id), }), ) })?; let control_point = Arc::clone(&state.control_point); let rid_for_log = rid.clone(); let object_id_for_log = object_id.clone(); let object_id_for_debug = object_id_for_log.clone(); // Launch the command in background and return immediately // The UI will be updated via SSE events when playback starts tokio::task::spawn(async move { let result = tokio::task::spawn_blocking(move || { debug!( renderer = rid.0.as_str(), server = sid.0.as_str(), object = object_id.as_str(), "play_content: fetching playback items" ); // Fetch playback items from server let items = fetch_playback_items(&control_point, &sid, &object_id)?; debug!( renderer = rid.0.as_str(), item_count = items.len(), "play_content: fetched items" ); if items.is_empty() { return Err(anyhow::anyhow!("No playable content found")); } if items.len() > 1 { debug!( renderer = rid.0.as_str(), server = sid.0.as_str(), object = object_id.as_str(), item_count = items.len(), "Auto-binding playlist to renderer queue (auto_play = true)" ); control_point.attach_queue_to_playlist_with_options( &rid, sid.clone(), object_id.clone(), true, )?; return Ok(()); } // Clear queue control_point.clear_queue(&rid)?; // Enqueue items control_point.enqueue_items(&rid, items)?; // Start playback // Pour les renderers OpenHome, play_current_from_queue() va gérer automatiquement // la lecture depuis la playlist native si elle existe control_point.play_current_from_queue(&rid)?; Ok::<(), anyhow::Error>(()) }) .await; match result { Ok(Ok(())) => { debug!( "Successfully started playing content {} on renderer {}", object_id_for_log, rid_for_log.0 ); } Ok(Err(e)) => { warn!( "Failed to play content on renderer {}: {}", rid_for_log.0, e ); } Err(e) => { warn!( "Task join error during play content for renderer {}: {}", rid_for_log.0, e ); } } }); debug!( renderer = renderer_id.as_str(), server = req.server_id.as_str(), object = object_id_for_debug.as_str(), "Content playing via HTTP API" ); Ok(Json(SuccessResponse { message: "Content playing".to_string(), })) } /// POST /control/renderers/{renderer_id}/queue/add - Ajouter du contenu à la queue #[cfg(feature = "pmoserver")] #[utoipa::path( post, path = "/renderers/{renderer_id}/queue/add", params( ("renderer_id" = String, Path, description = "ID unique du renderer") ), request_body = PlayContentRequest, responses( (status = 200, description = "Contenu ajouté à la queue", body = SuccessResponse), (status = 404, description = "Renderer ou serveur non trouvé", body = ErrorResponse), (status = 500, description = "Erreur lors de l'exécution", body = ErrorResponse) ), tag = "control" )] async fn add_to_queue( State(state): State, Path(renderer_id): Path, Json(req): Json, ) -> Result, (StatusCode, Json)> { let rid = DeviceId(renderer_id.clone()); let sid = DeviceId(req.server_id.clone()); let object_id = req.object_id.clone(); let object_id_for_log = object_id.clone(); // Verify renderer exists state .control_point .music_renderer_by_id(&rid) .ok_or_else(|| { ( StatusCode::NOT_FOUND, Json(ErrorResponse { error: format!("Renderer {} not found", renderer_id), }), ) })?; let control_point = Arc::clone(&state.control_point); let rid_for_log = rid.clone(); let object_id_for_log = object_id.clone(); let object_id_for_debug = object_id_for_log.clone(); // Launch the command in background and return immediately // The UI will be updated via SSE events when the queue changes tokio::task::spawn(async move { let result = tokio::task::spawn_blocking(move || { // Fetch playback items from server let items = fetch_playback_items(&control_point, &sid, &object_id)?; if items.is_empty() { return Err(anyhow::anyhow!("No playable content found")); } // Enqueue items control_point.enqueue_items(&rid, items)?; Ok::<(), anyhow::Error>(()) }) .await; match result { Ok(Ok(())) => { debug!( "Successfully added content {} to queue for renderer {}", object_id_for_log, rid_for_log.0 ); } Ok(Err(e)) => { warn!( "Failed to add content to queue for renderer {}: {}", rid_for_log.0, e ); } Err(e) => { warn!( "Task join error during add to queue for renderer {}: {}", rid_for_log.0, e ); } } }); debug!( renderer = renderer_id.as_str(), server = req.server_id.as_str(), object = object_id_for_debug.as_str(), "Content added to queue via HTTP API" ); Ok(Json(SuccessResponse { message: "Content added to queue".to_string(), })) } /// POST /control/renderers/{renderer_id}/queue/add-after - Ajouter du contenu après le morceau actuel #[cfg(feature = "pmoserver")] #[utoipa::path( post, path = "/renderers/{renderer_id}/queue/add-after", params( ("renderer_id" = String, Path, description = "ID unique du renderer") ), request_body = PlayContentRequest, responses( (status = 200, description = "Contenu ajouté après le morceau actuel", body = SuccessResponse), (status = 404, description = "Renderer ou serveur non trouvé", body = ErrorResponse), (status = 504, description = "Timeout de la commande", body = ErrorResponse), (status = 500, description = "Erreur lors de l'exécution", body = ErrorResponse) ), tag = "control" )] async fn add_after_current( State(state): State, Path(renderer_id): Path, Json(req): Json, ) -> Result, (StatusCode, Json)> { let rid = DeviceId(renderer_id.clone()); let sid = DeviceId(req.server_id.clone()); let object_id = req.object_id.clone(); let object_id_for_log = object_id.clone(); // Verify renderer exists state .control_point .music_renderer_by_id(&rid) .ok_or_else(|| { ( StatusCode::NOT_FOUND, Json(ErrorResponse { error: format!("Renderer {} not found", renderer_id), }), ) })?; let control_point = Arc::clone(&state.control_point); let rid_for_log = rid.clone(); let object_id_for_log = object_id.clone(); let object_id_for_debug = object_id_for_log.clone(); // Launch the command in background and return immediately // The UI will be updated via SSE events when the queue changes tokio::task::spawn(async move { let result = tokio::task::spawn_blocking(move || { // Fetch playback items from server let items = fetch_playback_items(&control_point, &sid, &object_id)?; if items.is_empty() { return Err(anyhow::anyhow!("No playable content found")); } // Insert items after current using the new method control_point.enqueue_items_with_mode( &rid, items, crate::queue::EnqueueMode::InsertAfterCurrent, )?; Ok::<(), anyhow::Error>(()) }) .await; match result { Ok(Ok(())) => { debug!( "Successfully added content {} after current for renderer {}", object_id_for_log, rid_for_log.0 ); } Ok(Err(e)) => { warn!( "Failed to add content after current for renderer {}: {}", rid_for_log.0, e ); } Err(e) => { warn!( "Task join error during add after current for renderer {}: {}", rid_for_log.0, e ); } } }); debug!( renderer = renderer_id.as_str(), server = req.server_id.as_str(), object = object_id_for_debug.as_str(), "Content added after current via HTTP API" ); Ok(Json(SuccessResponse { message: "Content added after current track".to_string(), })) } /// POST /control/renderers/{renderer_id}/queue/transfer - Transfère la queue vers un autre renderer #[cfg(feature = "pmoserver")] #[utoipa::path( post, path = "/renderers/{renderer_id}/queue/transfer", params( ("renderer_id" = String, Path, description = "ID du renderer source") ), request_body = TransferQueueRequest, responses( (status = 200, description = "Queue transférée avec succès", body = SuccessResponse), (status = 404, description = "Renderer non trouvé", body = ErrorResponse), (status = 500, description = "Erreur lors du transfert", body = ErrorResponse) ), tag = "control" )] async fn transfer_queue( State(state): State, Path(source_renderer_id): Path, Json(req): Json, ) -> Result, (StatusCode, Json)> { let source_id = DeviceId(source_renderer_id.clone()); let dest_id = DeviceId(req.destination_renderer_id.clone()); debug!( source = source_renderer_id.as_str(), dest = req.destination_renderer_id.as_str(), "Transferring queue between renderers via HTTP API" ); let control_point = state.control_point.clone(); tokio::task::spawn_blocking(move || control_point.transfer_queue(&source_id, &dest_id)) .await .map_err(|e| { warn!( source = source_renderer_id.as_str(), dest = req.destination_renderer_id.as_str(), error = ?e, "Failed to spawn transfer_queue task" ); ( StatusCode::INTERNAL_SERVER_ERROR, Json(ErrorResponse { error: format!("Failed to spawn transfer task: {}", e), }), ) })? .map_err(|e| { warn!( source = source_renderer_id.as_str(), dest = req.destination_renderer_id.as_str(), error = ?e, "Failed to transfer queue" ); ( StatusCode::INTERNAL_SERVER_ERROR, Json(ErrorResponse { error: format!("Failed to transfer queue: {}", e), }), ) })?; debug!( source = source_renderer_id.as_str(), dest = req.destination_renderer_id.as_str(), "Queue transferred successfully via HTTP API" ); Ok(Json(SuccessResponse { message: format!( "Queue transferred from {} to {}", source_renderer_id, req.destination_renderer_id ), })) } // ============================================================================ // HANDLERS - MEDIA SERVERS // ============================================================================ /// GET /control/servers - Liste tous les serveurs de médias #[cfg(feature = "pmoserver")] #[utoipa::path( get, path = "/servers", responses( (status = 200, description = "Liste des serveurs de médias", body = Vec) ), tag = "control" )] async fn list_servers(State(state): State) -> Json> { // Use spawn_blocking to avoid blocking the tokio runtime let control_point = state.control_point.clone(); let summaries = tokio::task::spawn_blocking(move || { let servers = control_point.list_media_servers().unwrap_or_default(); servers .into_iter() .map(|s| MediaServerSummary { id: s.id().0.clone(), friendly_name: s.friendly_name().to_string(), model_name: s.model_name().to_string(), online: s.is_online(), }) .collect::>() }) .await .unwrap_or_default(); Json(summaries) } /// Paramètres de pagination pour le browse #[cfg(feature = "pmoserver")] #[derive(Debug, serde::Deserialize)] struct BrowseParams { #[serde(default)] offset: u32, limit: Option, } /// GET /control/servers/{server_id}/containers/{container_id} - Browse un container #[cfg(feature = "pmoserver")] #[utoipa::path( get, path = "/servers/{server_id}/containers/{container_id}", params( ("server_id" = String, Path, description = "ID unique du serveur"), ("container_id" = String, Path, description = "ID du container (use '0' for root)"), ("offset" = Option, Query, description = "Index de départ (défaut: 0)"), ("limit" = Option, Query, description = "Nombre max d'items (défaut: 50)"), ), responses( (status = 200, description = "Contenu du container", body = BrowseResponse), (status = 404, description = "Serveur non trouvé", body = ErrorResponse), (status = 500, description = "Erreur lors du browse", body = ErrorResponse) ), tag = "control" )] async fn browse_container( State(state): State, Path((server_id, container_id)): Path<(String, String)>, Query(params): Query, headers: HeaderMap, ) -> Result, (StatusCode, Json)> { let base_url_str = pmoserver::get_base_url_from_request(&headers); let base_url = pmoserver::BaseUrl(base_url_str); let sid = DeviceId(server_id.clone()); let server = state.control_point.media_server(&sid).ok_or_else(|| { ( StatusCode::NOT_FOUND, Json(ErrorResponse { error: format!("Server {} not found", server_id), }), ) })?; if !server.is_online() { return Err(( StatusCode::SERVICE_UNAVAILABLE, Json(ErrorResponse { error: format!("Server {} is offline", server_id), }), )); } if !server.has_content_directory() { return Err(( StatusCode::NOT_IMPLEMENTED, Json(ErrorResponse { error: format!("Server {} does not support ContentDirectory", server_id), }), )); } let offset = params.offset; let limit = params.limit.unwrap_or(BROWSE_DEFAULT_LIMIT); // Use spawn_blocking to avoid blocking the async runtime with synchronous SOAP calls let container_id_clone = container_id.clone(); let server_clone = server.clone(); let browse_task = tokio::task::spawn_blocking(move || { server_clone.browse_children_paged(&container_id_clone, offset, limit) }); let page = 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 {}: {}", container_id, server_id, e ); ( StatusCode::INTERNAL_SERVER_ERROR, Json(ErrorResponse { error: format!("Failed to browse container: {}", e), }), ) })?; // Transform cover URLs let mut container_entries = Vec::with_capacity(page.entries.len()); for e in page.entries { let album_art_uri = transform_cover_url(e.album_art_uri.as_deref(), &base_url).await; container_entries.push(ContainerEntry { id: e.id, title: e.title, class: e.class, is_container: e.is_container, child_count: e.child_count, artist: e.artist, album: e.album, album_art_uri, }); } Ok(Json(BrowseResponse { container_id, entries: container_entries, total_count: page.total_count, offset, })) } // ============================================================================ // HELPERS // ============================================================================ #[cfg(feature = "pmoserver")] fn map_snapshot_error( renderer_id: String, err: anyhow::Error, ) -> (StatusCode, Json) { warn!( renderer = renderer_id.as_str(), error = %err, "Failed to build renderer snapshot" ); ( StatusCode::NOT_FOUND, Json(ErrorResponse { error: format!("Renderer {} not found", renderer_id), }), ) } /// Transforme une URL de cover pour qu'elle soit accessible depuis le client /// /// Si l'URL est une route locale de notre cache (/covers/...), on la transforme en URL absolue. /// Sinon, on utilise pmocovers::proxy_cover_url() pour mettre en cache et retourner notre URL. /// C'est le même mécanisme que PMO Cache utilise déjà pour Qobuz. async fn transform_cover_url(url: Option<&str>, base_url: &pmoserver::BaseUrl) -> Option { let url = url?; // Si c'est déjà une route locale de notre cache, la transformer en URL absolue if url.starts_with("/covers/") { debug!(url = %url, "Already local cover route"); return Some(base_url.url_for(url)); } // Si c'est une URL de notre instance, la retourner directement if url.starts_with(&base_url.0) { debug!(url = %url, "Already our instance URL"); return Some(url.to_string()); } // Pour les autres URLs, utiliser le mechanisme de proxy standard (comme Qobuz) debug!(url = %url, "Proxyfying cover URL via pmocovers"); match pmocovers::proxy_cover_url(url, base_url).await { Ok(local_url) => { debug!(result = %local_url, "Proxified successfully"); Some(local_url) }, Err(e) => { tracing::warn!("Failed to proxy cover URL {}: {}", url, e); Some(url.to_string()) } } } /// Helper to fetch playback items from a media server object (container or item). /// /// This function browses the server to get the entries and converts them to PlaybackItem. /// For containers, it browses children. For items, it browses metadata. #[cfg(feature = "pmoserver")] fn fetch_playback_items( control_point: &ControlPoint, server_id: &DeviceId, object_id: &str, ) -> anyhow::Result> { // Get server from registry let server = control_point .media_server(server_id) .ok_or_else(|| anyhow::anyhow!("Server {} not found", server_id.0))?; if !server.is_online() { return Err(anyhow::anyhow!("Server {} is offline", server_id.0)); } if !server.has_content_directory() { return Err(anyhow::anyhow!( "Server {} does not support ContentDirectory", server_id.0 )); } // First, get metadata for the object to determine if it's a container or item let object_metadata = server.browse_object(object_id)?; let entries = if object_metadata.is_container { // For containers, browse all children with pagination let mut all_entries = Vec::new(); let mut offset = 0u32; loop { let page = server.browse_children(object_id, offset, BROWSE_DEFAULT_LIMIT)?; let fetched = page.len() as u32; all_entries.extend(page); if fetched < BROWSE_DEFAULT_LIMIT { break; } offset += fetched; } all_entries } else { // For items, use the object itself vec![object_metadata] }; debug!( server_id = server_id.0.as_str(), object_id = object_id, total_entries = entries.len(), containers = entries.iter().filter(|e| e.is_container).count(), items_count = entries.iter().filter(|e| !e.is_container).count(), "Browse returned entries" ); // Convert to PlaybackItem let items: Vec = entries .iter() .filter_map(|entry| playback_item_from_entry(server.clone(), entry)) .collect(); if items.is_empty() && !entries.is_empty() { warn!( server_id = server_id.0.as_str(), object_id = object_id, total_entries = entries.len(), "No playable items found - all entries were filtered out" ); } Ok(items) } #[cfg(feature = "pmoserver")] fn protocol_summary(protocol: &RendererProtocol) -> RendererProtocolSummary { match protocol { RendererProtocol::UpnpAvOnly => RendererProtocolSummary::Upnp, RendererProtocol::OpenHomeOnly => RendererProtocolSummary::Openhome, RendererProtocol::OpenHomeHybrid => RendererProtocolSummary::Hybrid, RendererProtocol::ChromecastOnly => RendererProtocolSummary::Chromecast, } } #[cfg(feature = "pmoserver")] fn capability_summary(caps: &RendererCapabilities) -> RendererCapabilitiesSummary { RendererCapabilitiesSummary { has_avtransport: caps.has_avtransport, has_avtransport_set_next: caps.has_avtransport_set_next, has_rendering_control: caps.has_rendering_control, has_connection_manager: caps.has_connection_manager, has_linkplay_http: caps.has_linkplay_http, has_arylic_tcp: caps.has_arylic_tcp, has_oh_playlist: caps.has_oh_playlist, has_oh_volume: caps.has_oh_volume, has_oh_info: caps.has_oh_info, has_oh_time: caps.has_oh_time, has_oh_radio: caps.has_oh_radio, has_chromecast: caps.has_chromecast, } } /// Paramètres de recherche #[cfg(feature = "pmoserver")] #[derive(Debug, serde::Deserialize)] struct SearchQuery { q: String, } #[cfg(feature = "pmoserver")] fn search_result_container_id(entries: &[ContainerEntry]) -> String { entries .iter() .find(|entry| entry.is_container) .map(|entry| { let parts: Vec<&str> = entry.id.splitn(5, ':').collect(); if parts.len() == 5 && parts[0] == "qobuz" && parts[1] == "search" { // Résultat de recherche Qobuz : reconstruire le container virtuel parent // ex. "qobuz:search:catalog:albums:Beethoven" → "qobuz:search:catalog:all:Beethoven" format!("qobuz:search:{}:all:{}", parts[2], parts[4]) } else { // Résultat d'une autre source (ex. UrlSource) : retourner l'ID tel quel entry.id.clone() } }) .unwrap_or_else(|| "search".to_string()) } /// GET /control/servers/{server_id}/search?q= - Recherche dans un serveur #[cfg(feature = "pmoserver")] #[utoipa::path( get, path = "/servers/{server_id}/search", params( ("server_id" = String, Path, description = "ID unique du serveur"), ("q" = String, Query, description = "Requête de recherche"), ), responses( (status = 200, description = "Résultats de recherche", body = BrowseResponse), (status = 404, description = "Serveur non trouvé", body = ErrorResponse), (status = 500, description = "Erreur lors de la recherche", body = ErrorResponse) ), tag = "control" )] async fn search_server( State(state): State, Path(server_id): Path, Query(params): Query, ) -> Result, (StatusCode, Json)> { let sid = DeviceId(server_id.clone()); let server = state.control_point.media_server(&sid).ok_or_else(|| { ( StatusCode::NOT_FOUND, Json(ErrorResponse { error: format!("Server {} not found", server_id), }), ) })?; if !server.is_online() { return Err(( StatusCode::SERVICE_UNAVAILABLE, Json(ErrorResponse { error: format!("Server {} is offline", server_id), }), )); } if !server.has_content_directory() { return Err(( StatusCode::NOT_IMPLEMENTED, Json(ErrorResponse { error: format!("Server {} does not support ContentDirectory", server_id), }), )); } debug!(server_id = %server_id, query = %params.q, "Search request"); let query = params.q.clone(); let server_clone = server.clone(); let search_task = tokio::task::spawn_blocking(move || { server_clone.search("0", &query, 0, 200) }); let entries = time::timeout(BROWSE_REQUEST_TIMEOUT, search_task) .await .map_err(|_| { warn!( "Search request on server {} exceeded {:?}", server_id, BROWSE_REQUEST_TIMEOUT ); ( StatusCode::GATEWAY_TIMEOUT, Json(ErrorResponse { error: format!( "Search request timed out after {}s", BROWSE_REQUEST_TIMEOUT.as_secs() ), }), ) })? .map_err(|e| { warn!("Task join error during search: {}", e); ( StatusCode::INTERNAL_SERVER_ERROR, Json(ErrorResponse { error: format!("Internal task error: {}", e), }), ) })? .map_err(|e| { warn!("Failed to search on server {}: {}", server_id, e); ( StatusCode::INTERNAL_SERVER_ERROR, Json(ErrorResponse { error: format!("Failed to search: {}", e), }), ) })?; let total_count = entries.len() as u32; debug!(server_id = %server_id, count = total_count, "Search results"); let container_entries: Vec = entries .into_iter() .map(|e| ContainerEntry { id: e.id, title: e.title, class: e.class, is_container: e.is_container, child_count: e.child_count, artist: e.artist, album: e.album, album_art_uri: e.album_art_uri, }) .collect(); let container_id = search_result_container_id(&container_entries); Ok(Json(BrowseResponse { container_id, entries: container_entries, total_count, offset: 0, })) } // ============================================================================ // ROUTER & TRAIT // ============================================================================ /// Crée le router pour l'API Control Point #[cfg(feature = "pmoserver")] pub fn create_api_router(state: ControlPointState, control_point: Arc) -> Router { Router::new() // Renderers .route("/renderers", get(list_renderers)) .route("/renderers/{renderer_id}", get(get_renderer_state)) .route( "/renderers/{renderer_id}/full", get(get_renderer_full_snapshot), ) .route("/renderers/{renderer_id}/queue", get(get_renderer_queue)) .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}/resume", post(resume_renderer)) .route("/renderers/{renderer_id}/next", post(next_renderer)) .route("/renderers/{renderer_id}/seek", post(seek_renderer)) // Queue control .route( "/renderers/{renderer_id}/queue/seek", post(seek_queue_index), ) .route( "/renderers/{renderer_id}/queue/shuffle", post(shuffle_queue), ) // 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), ) // Sleep timer .route("/renderers/{renderer_id}/timer", get(get_sleep_timer_state)) .route( "/renderers/{renderer_id}/timer/start", post(start_sleep_timer), ) .route( "/renderers/{renderer_id}/timer/update", post(update_sleep_timer), ) .route( "/renderers/{renderer_id}/timer/cancel", post(cancel_sleep_timer), ) // Stream state .route( "/renderers/{renderer_id}/stream-state", get(get_stream_state), ) // Playlist binding .route( "/renderers/{renderer_id}/binding/attach", post(attach_playlist_binding), ) .route( "/renderers/{renderer_id}/binding/detach", post(detach_playlist_binding), ) // Queue content .route("/renderers/{renderer_id}/queue/play", post(play_content)) .route("/renderers/{renderer_id}/queue/add", post(add_to_queue)) .route( "/renderers/{renderer_id}/queue/add-after", post(add_after_current), ) .route( "/renderers/{renderer_id}/queue/transfer", post(transfer_queue), ) // Servers .route("/servers", get(list_servers)) .route( "/servers/{server_id}/containers/{container_id}", get(browse_container), ) .route("/servers/{server_id}/search", get(search_server)) .with_state(state) // SSE events - merge the SSE router .merge(crate::sse::create_sse_router(control_point)) } /// Trait d'extension pour pmoserver::Server /// /// Permet d'initialiser le ControlPoint avec routes HTTP complètes #[cfg(feature = "pmoserver")] #[async_trait] pub trait ControlPointExt { /// Enregistre et initialise le Control Point avec son API complète /// /// Cette fonction de haut niveau : /// 1. Lance le runtime du ControlPoint (découverte SSDP, polling renderers, etc.) /// 2. Enregistre toutes les routes HTTP REST /// 3. Enregistre tous les endpoints SSE pour les événements /// 4. Génère la documentation OpenAPI /// /// # Routes créées /// /// - API REST: `/api/control/*` /// - `/renderers` - Liste et état des renderers /// - `/servers` - Liste et navigation des serveurs de médias /// - Contrôles de transport, volume, queue, binding /// - SSE Events: `/api/control/events/*` /// - `/events` - Tous les événements (renderers + serveurs) /// - `/events/renderers` - Événements renderers uniquement /// - `/events/servers` - Événements serveurs uniquement /// - Swagger: `/swagger-ui/control` /// /// # Arguments /// /// * `timeout_secs` - Timeout HTTP pour les requêtes UPnP (recommandé: 5 secondes) /// /// # Returns /// /// Retourne l'instance du ControlPoint dans un Arc pour permettre /// d'interagir avec depuis l'application. /// /// # Errors /// /// Retourne une erreur si le runtime SSDP ne peut pas être démarré. /// /// # Examples /// /// ```ignore /// use pmocontrol::ControlPointExt; /// use pmoserver::Server; /// /// let server = Server::create_upnp_server().await?; /// /// // Enregistrer le Control Point avec timeout de 5 secondes /// let control_point = server /// .write() /// .await /// .register_control_point(5) /// .await?; /// /// // Le Control Point est maintenant actif et ses routes HTTP/SSE sont enregistrées /// // 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>; /// Initialise l'API Control Point (bas niveau) /// /// Cette méthode est appelée automatiquement par `register_control_point()`. /// Utilisez `register_control_point()` pour la plupart des cas d'usage. /// /// # Routes créées /// /// - API REST: `/api/control/*` /// - SSE Events: `/api/control/events/*` /// - Swagger: `/swagger-ui/control` /// /// # Arguments /// /// * `control_point` - Instance du ControlPoint déjà créée async fn init_control_point(&mut self, control_point: Arc); } #[cfg(feature = "pmoserver")] #[async_trait] impl ControlPointExt for pmoserver::Server { async fn register_control_point( &mut self, timeout_secs: u64, ) -> std::io::Result> { use tracing::info; info!("🎛️ Initializing Control Point..."); // 1. Lancer le runtime du ControlPoint let control_point = ControlPoint::spawn(timeout_secs)?; let control_point = Arc::new(control_point); info!("✅ Control Point runtime started"); info!(" - SSDP discovery active"); info!(" - Renderer polling active (1s interval)"); info!(" - MediaServer event subscriptions active"); // 2. Enregistrer les routes HTTP REST et SSE self.init_control_point(control_point.clone()).await; info!("✅ Control Point API registered:"); info!(" - REST API: /api/control/*"); info!(" - SSE Events: /api/control/events/*"); info!(" - OpenAPI docs: /swagger-ui/control"); Ok(control_point) } async fn init_control_point(&mut self, control_point: Arc) { let state = ControlPointState::new(control_point.clone()); // Créer le router API (inclut REST et SSE) let api_router = create_api_router(state, control_point); // L'enregistrer avec OpenAPI self.add_openapi(api_router, crate::openapi::ApiDoc::openapi(), "control") .await; } }