Files
pmomusic/pmocontrol/src/pmoserver_ext.rs
Eric Coissac ea5936717a Add pmocovers integration for external cover URL proxying
- Introduce optional `pmocover` dependency in pmocontrol
- Add cover URL transformation logic for both REST and SSE endpoints using `pmocovers::proxy_cover_url`/sync
- Refactor `/covers/proxy?...=` handler to accept any external URL and return local cached route
- Implement helper functions `transform_cover_url` (async) & sync variant for consistent cover URL normalization
- Update Cargo.lock to include `pmocovers`
2026-04-04 12:05:10 +02:00

2727 lines
89 KiB
Rust

//! 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::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<ControlPoint>,
}
#[cfg(feature = "pmoserver")]
impl ControlPointState {
pub fn new(control_point: Arc<ControlPoint>) -> 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<RendererSummary>)
),
tag = "control"
)]
async fn list_renderers(State(state): State<ControlPointState>) -> Json<Vec<RendererSummary>> {
// 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::<Vec<_>>()
})
.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<ControlPointState>,
Path(renderer_id): Path<String>,
) -> Result<Json<RendererState>, (StatusCode, Json<ErrorResponse>)> {
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<ControlPointState>,
Path(renderer_id): Path<String>,
headers: HeaderMap,
) -> Result<Json<FullRendererSnapshot>, (StatusCode, Json<ErrorResponse>)> {
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);
}
}
}
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<ControlPointState>,
Path(renderer_id): Path<String>,
headers: HeaderMap,
) -> Result<Json<QueueSnapshot>, (StatusCode, Json<ErrorResponse>)> {
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);
// Transform cover URLs in all queue items
for item in &mut snapshot.queue.items {
if let Some(ref mut metadata) = item.metadata {
if let Some(ref album_art) = metadata.album_art_uri {
if let Some(transformed) = transform_cover_url(Some(album_art), &base_url).await {
metadata.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<AttachedPlaylistInfo>),
(status = 404, description = "Renderer non trouvé", body = ErrorResponse)
),
tag = "control"
)]
async fn get_renderer_binding(
State(state): State<ControlPointState>,
Path(renderer_id): Path<String>,
) -> Result<Json<Option<AttachedPlaylistInfo>>, (StatusCode, Json<ErrorResponse>)> {
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<ControlPointState>,
Path(renderer_id): Path<String>,
) -> Result<Json<SuccessResponse>, (StatusCode, Json<ErrorResponse>)> {
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<ControlPointState>,
Path(renderer_id): Path<String>,
) -> Result<Json<SuccessResponse>, (StatusCode, Json<ErrorResponse>)> {
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<ControlPointState>,
Path(renderer_id): Path<String>,
) -> Result<Json<SuccessResponse>, (StatusCode, Json<ErrorResponse>)> {
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<ControlPointState>,
Path(renderer_id): Path<String>,
) -> Result<Json<SuccessResponse>, (StatusCode, Json<ErrorResponse>)> {
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<ControlPointState>,
Path(renderer_id): Path<String>,
) -> Result<Json<SuccessResponse>, (StatusCode, Json<ErrorResponse>)> {
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<ControlPointState>,
Path(renderer_id): Path<String>,
Json(payload): Json<SeekQueueRequest>,
) -> Result<Json<SuccessResponse>, (StatusCode, Json<ErrorResponse>)> {
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<ControlPointState>,
Path(renderer_id): Path<String>,
Json(payload): Json<SeekRequest>,
) -> Result<Json<SuccessResponse>, (StatusCode, Json<ErrorResponse>)> {
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<ControlPointState>,
Path(renderer_id): Path<String>,
Json(req): Json<VolumeSetRequest>,
) -> Result<Json<SuccessResponse>, (StatusCode, Json<ErrorResponse>)> {
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<ControlPointState>,
Path(renderer_id): Path<String>,
) -> Result<Json<SuccessResponse>, (StatusCode, Json<ErrorResponse>)> {
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::<u16, anyhow::Error>(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<ControlPointState>,
Path(renderer_id): Path<String>,
) -> Result<Json<SuccessResponse>, (StatusCode, Json<ErrorResponse>)> {
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::<u16, anyhow::Error>(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<ControlPointState>,
Path(renderer_id): Path<String>,
) -> Result<Json<SuccessResponse>, (StatusCode, Json<ErrorResponse>)> {
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::<bool, anyhow::Error>(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<ControlPointState>,
Path(renderer_id): Path<String>,
Json(req): Json<SleepTimerRequest>,
) -> Result<Json<SleepTimerState>, (StatusCode, Json<ErrorResponse>)> {
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<ControlPointState>,
Path(renderer_id): Path<String>,
Json(req): Json<SleepTimerRequest>,
) -> Result<Json<SleepTimerState>, (StatusCode, Json<ErrorResponse>)> {
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<ControlPointState>,
Path(renderer_id): Path<String>,
) -> Result<Json<SuccessResponse>, (StatusCode, Json<ErrorResponse>)> {
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<ControlPointState>,
Path(renderer_id): Path<String>,
) -> Result<Json<SleepTimerState>, (StatusCode, Json<ErrorResponse>)> {
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<ControlPointState>,
Path(renderer_id): Path<String>,
) -> Result<Json<StreamState>, (StatusCode, Json<ErrorResponse>)> {
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<ControlPointState>,
Path(renderer_id): Path<String>,
) -> Result<Json<SuccessResponse>, (StatusCode, Json<ErrorResponse>)> {
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<ControlPointState>,
Path(renderer_id): Path<String>,
Json(req): Json<AttachPlaylistRequest>,
) -> Result<Json<SuccessResponse>, (StatusCode, Json<ErrorResponse>)> {
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<ControlPointState>,
Path(renderer_id): Path<String>,
) -> Json<SuccessResponse> {
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<ControlPointState>,
Path(renderer_id): Path<String>,
Json(req): Json<PlayContentRequest>,
) -> Result<Json<SuccessResponse>, (StatusCode, Json<ErrorResponse>)> {
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<ControlPointState>,
Path(renderer_id): Path<String>,
Json(req): Json<PlayContentRequest>,
) -> Result<Json<SuccessResponse>, (StatusCode, Json<ErrorResponse>)> {
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<ControlPointState>,
Path(renderer_id): Path<String>,
Json(req): Json<PlayContentRequest>,
) -> Result<Json<SuccessResponse>, (StatusCode, Json<ErrorResponse>)> {
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<ControlPointState>,
Path(source_renderer_id): Path<String>,
Json(req): Json<TransferQueueRequest>,
) -> Result<Json<SuccessResponse>, (StatusCode, Json<ErrorResponse>)> {
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<MediaServerSummary>)
),
tag = "control"
)]
async fn list_servers(State(state): State<ControlPointState>) -> Json<Vec<MediaServerSummary>> {
// 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::<Vec<_>>()
})
.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<u32>,
}
/// 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<u32>, Query, description = "Index de départ (défaut: 0)"),
("limit" = Option<u32>, 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<ControlPointState>,
Path((server_id, container_id)): Path<(String, String)>,
Query(params): Query<BrowseParams>,
headers: HeaderMap,
) -> Result<Json<BrowseResponse>, (StatusCode, Json<ErrorResponse>)> {
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: None,
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<ErrorResponse>) {
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<String> {
let url = url?;
// Si c'est déjà une route locale de notre cache, la transformer en URL absolue
if url.starts_with("/covers/") {
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) {
return Some(url.to_string());
}
// Pour les autres URLs, utiliser le mechanisme de proxy standard (comme Qobuz)
match pmocovers::proxy_cover_url(url, base_url).await {
Ok(local_url) => 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<Vec<PlaybackItem>> {
// 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<PlaybackItem> = 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,
}
/// GET /control/servers/{server_id}/search?q=<query> - 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<ControlPointState>,
Path(server_id): Path<String>,
Query(params): Query<SearchQuery>,
) -> Result<Json<BrowseResponse>, (StatusCode, Json<ErrorResponse>)> {
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<ContainerEntry> = entries
.into_iter()
.map(|e| ContainerEntry {
id: e.id,
title: e.title,
class: e.class,
is_container: e.is_container,
child_count: None,
artist: e.artist,
album: e.album,
album_art_uri: e.album_art_uri,
})
.collect();
Ok(Json(BrowseResponse {
container_id: "search".to_string(),
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<ControlPoint>) -> 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<Arc<ControlPoint>>;
/// 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<ControlPoint>);
}
#[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>> {
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<ControlPoint>) {
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;
}
}