From 75878be5640bfa6ee2a80d94d18d15d9b8214ca3 Mon Sep 17 00:00:00 2001 From: Eric Coissac Date: Wed, 3 Dec 2025 23:17:41 +0100 Subject: [PATCH] Deboguage pmo_remote_control --- Cargo.lock | 15 +- pmocontrol/Cargo.toml | 6 + .../examples/full_control_point_demo.rs | 30 +- pmocontrol/examples/pmo_remote_control.rs | 1847 +++++++++++++++++ .../examples/pmomusic_integration_example.rs | 4 +- pmocontrol/src/control_point.rs | 204 +- pmocontrol/src/media_server.rs | 13 +- pmocontrol/src/playback_queue.rs | 12 + pmocontrol/src/pmoserver_ext.rs | 734 +++++-- pmocontrol/src/sse.rs | 14 +- pmodidl/src/lib.rs | 10 +- 11 files changed, 2620 insertions(+), 269 deletions(-) create mode 100644 pmocontrol/examples/pmo_remote_control.rs diff --git a/Cargo.lock b/Cargo.lock index ac7b4d44..dddee6f7 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3252,6 +3252,7 @@ dependencies = [ "chrono", "crossbeam-channel", "crossterm", + "percent-encoding", "pmodidl", "pmoserver", "pmoupnp", @@ -3263,6 +3264,7 @@ dependencies = [ "tokio", "tokio-stream", "tracing", + "tracing-log 0.1.4", "tracing-subscriber", "ureq", "utoipa", @@ -5232,6 +5234,17 @@ dependencies = [ "valuable", ] +[[package]] +name = "tracing-log" +version = "0.1.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f751112709b4e791d8ce53e32c4ed2d353565a795ce84da2285393f41557bdf2" +dependencies = [ + "log", + "once_cell", + "tracing-core", +] + [[package]] name = "tracing-log" version = "0.2.0" @@ -5258,7 +5271,7 @@ dependencies = [ "thread_local", "tracing", "tracing-core", - "tracing-log", + "tracing-log 0.2.0", ] [[package]] diff --git a/pmocontrol/Cargo.toml b/pmocontrol/Cargo.toml index 6551e1b5..8cf8b0e4 100644 --- a/pmocontrol/Cargo.toml +++ b/pmocontrol/Cargo.toml @@ -11,6 +11,7 @@ thiserror = "2.0.17" ureq = "3.1.4" tracing = "0.1.41" tracing-subscriber = "0.3" +tracing-log = "0.1" anyhow = "1.0" xmltree = "0.11.0" crossbeam-channel = "0.5" @@ -29,6 +30,11 @@ tokio-stream = { version = "0.1", features = ["sync"], optional = true } async-stream = { version = "0.3", optional = true } chrono = { version = "0.4", features = ["serde"], optional = true } +[dev-dependencies] +serde = { version = "1.0", features = ["derive"] } +serde_json = "1.0" +percent-encoding = "2.3" + [features] default = [] # Active l'API REST pmoserver diff --git a/pmocontrol/examples/full_control_point_demo.rs b/pmocontrol/examples/full_control_point_demo.rs index f2f0ad0b..5af38635 100644 --- a/pmocontrol/examples/full_control_point_demo.rs +++ b/pmocontrol/examples/full_control_point_demo.rs @@ -6,29 +6,29 @@ use std::collections::HashMap; use std::io::{self, Stdout}; use std::process; -use std::sync::mpsc::{self, Sender}; use std::sync::Arc; +use std::sync::mpsc::{self, Sender}; use std::thread; use std::time::{Duration, Instant}; -use anyhow::{anyhow, Context, Result}; +use anyhow::{Context, Result, anyhow}; use crossterm::event::{self, DisableMouseCapture, EnableMouseCapture, Event, KeyCode, KeyEvent}; use crossterm::execute; use crossterm::terminal::{ - disable_raw_mode, enable_raw_mode, EnterAlternateScreen, LeaveAlternateScreen, + EnterAlternateScreen, LeaveAlternateScreen, disable_raw_mode, enable_raw_mode, }; use pmocontrol::model::TrackMetadata; use pmocontrol::{ ControlPoint, DeviceRegistryRead, MediaBrowser, MediaEntry, MediaResource, MediaServerEvent, - MediaServerInfo, MusicServer, PlaybackItem, PlaybackPosition, PlaybackStatus, - PlaybackPositionInfo, RendererEvent, RendererInfo, TransportControl, VolumeControl, + MediaServerInfo, MusicServer, PlaybackItem, PlaybackPosition, PlaybackPositionInfo, + PlaybackStatus, RendererEvent, RendererInfo, TransportControl, VolumeControl, }; +use ratatui::Terminal; use ratatui::backend::CrosstermBackend; use ratatui::layout::{Alignment, Constraint, Direction, Layout, Rect}; use ratatui::style::{Color, Modifier, Style}; use ratatui::text::{Line, Span}; use ratatui::widgets::{Block, Borders, Clear, Gauge, List, ListItem, ListState, Paragraph}; -use ratatui::Terminal; const DEFAULT_TIMEOUT_SECS: u64 = 5; const DEFAULT_DISCOVERY_SECS: u64 = 15; @@ -725,10 +725,10 @@ impl App { self.control_point.clear_queue(&renderer_id)?; self.control_point.enqueue_items(&renderer_id, items)?; - let (full_queue, current_index) = - self.control_point - .get_full_queue_snapshot(&renderer_id) - .context("failed to snapshot queue after enqueue")?; + let (full_queue, current_index) = self + .control_point + .get_full_queue_snapshot(&renderer_id) + .context("failed to snapshot queue after enqueue")?; self.queue_snapshot = full_queue; self.queue_current_index = current_index; self.record_queue_metadata(); @@ -948,19 +948,15 @@ impl App { } fn refresh_queue_snapshot(&mut self, renderer_id: &pmocontrol::model::RendererId) { - if let Ok((queue, current_index)) = - self.control_point.get_full_queue_snapshot(renderer_id) + if let Ok((queue, current_index)) = self.control_point.get_full_queue_snapshot(renderer_id) { self.queue_snapshot = queue; self.queue_current_index = current_index; self.record_queue_metadata(); if let Some(idx) = self.queue_current_index { if let Some(item) = self.queue_snapshot.get(idx).cloned() { - let needs_update = self - .ui_state - .current_track_uri - .as_deref() - != Some(item.uri.as_str()); + let needs_update = + self.ui_state.current_track_uri.as_deref() != Some(item.uri.as_str()); if needs_update { self.apply_item_as_current(&item); } diff --git a/pmocontrol/examples/pmo_remote_control.rs b/pmocontrol/examples/pmo_remote_control.rs new file mode 100644 index 00000000..d0afb9b5 --- /dev/null +++ b/pmocontrol/examples/pmo_remote_control.rs @@ -0,0 +1,1847 @@ +//! Remote REST-based control point demo using Ratatui. +//! +//! This example mirrors the UX of `full_control_point_demo.rs` but drives a +//! remote PMOMusic server exclusively through the `/api/control` REST API. + +use std::env; +use std::fs::{File, OpenOptions}; +use std::io::{self, Stdout, Write}; +use std::process; +use std::sync::mpsc::TryRecvError; +use std::sync::{mpsc, Arc, Mutex}; +use std::thread; +use std::time::{Duration, Instant}; + +use anyhow::{anyhow, bail, Context, Result}; +use crossterm::event::{self, DisableMouseCapture, EnableMouseCapture, Event, KeyCode, KeyEvent}; +use crossterm::execute; +use crossterm::terminal::{ + disable_raw_mode, enable_raw_mode, EnterAlternateScreen, LeaveAlternateScreen, +}; +use percent_encoding::{utf8_percent_encode, NON_ALPHANUMERIC}; +use ratatui::backend::CrosstermBackend; +use ratatui::layout::{Alignment, Constraint, Direction, Layout, Rect}; +use ratatui::style::{Color, Modifier, Style}; +use ratatui::text::{Line, Span}; +use ratatui::widgets::{Block, Borders, Clear, Gauge, List, ListItem, ListState, Paragraph}; +use ratatui::Terminal; +use serde::de::{DeserializeOwned, Deserializer}; +use serde::{Deserialize, Serialize}; +use tracing::info; +use tracing_subscriber::fmt::writer::BoxMakeWriter; +use tracing_subscriber::EnvFilter; +use ureq::http; +use ureq::{Agent, Body}; + +const DEFAULT_BASE_URL: &str = "http://localhost:8080/api/control"; +const DEFAULT_TIMEOUT: Duration = Duration::from_secs(15); +const TICK_RATE: Duration = Duration::from_millis(200); +const ROOT_CONTAINERS: &[&str] = &["0", "0$"]; + +fn main() -> Result<()> { + // Install panic handler to restore terminal even on panic + std::panic::set_hook(Box::new(|panic_info| { + // Force terminal restoration + let _ = disable_raw_mode(); + let _ = execute!(io::stdout(), LeaveAlternateScreen, DisableMouseCapture); + eprintln!("\n\n❌ Application panicked: {:?}", panic_info); + eprintln!("Terminal has been restored. You can now close this window safely."); + })); + + init_tracing(); + let options = resolve_options()?; + info!( + base_url = %options.base_url, + timeout_ms = options.timeout.as_millis(), + "Démarrage du client REST" + ); + println!("PMO Remote Control demo"); + println!("Using Control API at {}", options.base_url); + + let client = RestClient::new(&options.base_url, options.timeout)?; + let renderers = client + .list_renderers() + .context("Impossible de récupérer la liste des renderers")?; + if renderers.is_empty() { + eprintln!( + "Aucun renderer disponible via {}. Lancement annulé.", + options.base_url + ); + return Ok(()); + } + + let app = App::new(client, renderers); + if let Err(err) = run_app(app) { + eprintln!("Application fermée avec erreur: {err}"); + } + + println!("\nAu revoir !"); + Ok(()) +} + +struct AppOptions { + base_url: String, + timeout: Duration, +} + +fn resolve_options() -> Result { + let mut args = env::args().skip(1); + let mut cli_base: Option = None; + let mut cli_timeout: Option = None; + while let Some(arg) = args.next() { + match arg.as_str() { + "--base-url" => { + let value = args + .next() + .ok_or_else(|| anyhow!("--base-url requiert une valeur"))?; + cli_base = Some(value); + } + "--timeout-ms" => { + let value = args + .next() + .ok_or_else(|| anyhow!("--timeout-ms requiert une valeur"))?; + let millis: u64 = value + .parse() + .with_context(|| format!("Valeur invalide pour --timeout-ms: {value}"))?; + cli_timeout = Some(millis); + } + "--help" | "-h" => { + print_usage(); + process::exit(0); + } + other => bail!("Argument inconnu: {other}. Utilise --help pour l'aide."), + } + } + let base = cli_base + .or_else(|| env::var("PMO_REMOTE_BASE_URL").ok()) + .unwrap_or_else(|| DEFAULT_BASE_URL.to_string()); + let timeout_ms = cli_timeout + .or_else(|| { + env::var("PMO_REMOTE_TIMEOUT_MS") + .ok() + .and_then(|v| v.parse().ok()) + }) + .unwrap_or_else(|| DEFAULT_TIMEOUT.as_millis() as u64); + let timeout = Duration::from_millis(timeout_ms.max(1)); + Ok(AppOptions { + base_url: base, + timeout, + }) +} + +fn print_usage() { + println!( + "Usage: cargo run -p pmocontrol --example pmo_remote_control [-- --base-url --timeout-ms ]" + ); + println!("Variables d'environnement:"); + println!( + " PMO_REMOTE_BASE_URL Override la base de l'API REST (par défaut {DEFAULT_BASE_URL})" + ); + println!( + " PMO_REMOTE_TIMEOUT_MS Timeout HTTP global en millisecondes (par défaut {})", + DEFAULT_TIMEOUT.as_millis() + ); + println!( + " PMO_REMOTE_LOG_FILE Écrit les logs tracing dans ce fichier (append) au lieu de stderr" + ); + println!( + " RUST_LOG Active le filtrage tracing/log (ex: pmocontrol=debug,ureq=debug)" + ); +} + +fn init_tracing() { + let _ = tracing_log::LogTracer::init(); + let writer = log_writer(); + let env_filter = EnvFilter::try_from_default_env() + .or_else(|_| EnvFilter::try_new("info")) + .unwrap_or_else(|_| EnvFilter::new("info")); + let _ = tracing_subscriber::fmt() + .with_env_filter(env_filter) + .with_writer(writer) + .try_init(); +} + +fn log_writer() -> BoxMakeWriter { + if let Ok(path) = env::var("PMO_REMOTE_LOG_FILE") { + match OpenOptions::new().create(true).append(true).open(&path) { + Ok(file) => { + let shared = SharedLogWriter::new(file); + let writer = BoxMakeWriter::new(move || shared.clone()); + return writer; + } + Err(err) => { + eprintln!( + "Impossible d'ouvrir {path} pour les logs tracing: {err}. Retour à stderr" + ); + } + } + } + BoxMakeWriter::new(io::stderr) +} + +#[derive(Clone)] +struct SharedLogWriter { + inner: Arc>, +} + +impl SharedLogWriter { + fn new(file: File) -> Self { + Self { + inner: Arc::new(Mutex::new(file)), + } + } +} + +impl Write for SharedLogWriter { + fn write(&mut self, buf: &[u8]) -> io::Result { + let mut guard = self + .inner + .lock() + .map_err(|err| io::Error::new(io::ErrorKind::Other, err.to_string()))?; + guard.write(buf) + } + + fn flush(&mut self) -> io::Result<()> { + let mut guard = self + .inner + .lock() + .map_err(|err| io::Error::new(io::ErrorKind::Other, err.to_string()))?; + guard.flush() + } +} + +struct App { + client: RestClient, + renderers: Vec, + renderer_index: usize, + selected_renderer: Option, + servers: Vec, + server_index: usize, + selected_server: Option, + browser: Option, + mode: Mode, + ui_state: UiState, + queue_snapshot: Vec, + queue_current_index: Option, + binding_info: Option, + show_queue_overlay: bool, + show_help_overlay: bool, + pending_binding: Option, + status_line: String, + binding_worker: Option>, +} + +#[derive(Clone, Copy, PartialEq, Eq)] +enum Mode { + SelectRenderer, + SelectServer, + Browse, + BindingPrompt, + Control, +} + +#[derive(Clone)] +struct UiState { + renderer_name: String, + server_name: Option, + transport_state: Option, + progress: Option, + volume: Option, + mute: Option, + metadata: Option, + last_status: Option, +} + +#[derive(Clone)] +struct PlaybackProgress { + position_ms: Option, + duration_ms: Option, +} + +#[derive(Clone)] +struct TrackMetadata { + title: Option, + artist: Option, + album: Option, + album_art_uri: Option, +} + +struct PendingBinding { + container_id: String, + container_title: String, +} + +enum BindingWorkerMessage { + Success { container_title: String }, + Failure { error: String }, +} + +struct BrowserState { + server_id: String, + nav_state: NavigationState, + entries: Vec, + selected_index: usize, +} + +struct NavigationState { + path_stack: Vec<(String, String)>, + current_container_id: String, + current_container_title: String, +} + +impl App { + fn new(client: RestClient, renderers: Vec) -> Self { + Self { + client, + renderers, + renderer_index: 0, + selected_renderer: None, + servers: Vec::new(), + server_index: 0, + selected_server: None, + browser: None, + mode: Mode::SelectRenderer, + ui_state: UiState::placeholder(), + queue_snapshot: Vec::new(), + queue_current_index: None, + binding_info: None, + show_queue_overlay: false, + show_help_overlay: false, + pending_binding: None, + status_line: "Sélectionne un renderer avec ↑/↓ puis Entrée".to_string(), + binding_worker: None, + } + } + + fn draw(&self, f: &mut ratatui::Frame<'_>) { + match self.mode { + Mode::SelectRenderer => self.draw_renderer_selection(f), + Mode::SelectServer => self.draw_server_selection(f), + Mode::Browse => self.draw_browser(f), + _ => self.draw_control_screen(f), + }; + + if self.show_queue_overlay { + self.draw_queue_overlay(f); + } + if self.show_help_overlay { + self.draw_help_overlay(f); + } + if matches!(self.mode, Mode::BindingPrompt) { + self.draw_binding_prompt(f); + } + self.draw_status_line(f); + } + + fn draw_renderer_selection(&self, f: &mut ratatui::Frame<'_>) { + let area = f.size(); + let block = Block::default() + .borders(Borders::ALL) + .title("Sélection du renderer"); + let items: Vec = self + .renderers + .iter() + .map(|info| { + let status = if info.online { + "[en ligne]" + } else { + "[hors ligne]" + }; + let text = format!("{status} {} | {}", info.friendly_name, info.model_name); + ListItem::new(text) + }) + .collect(); + let list = List::new(items) + .block(block) + .highlight_style( + Style::default() + .fg(Color::Yellow) + .add_modifier(Modifier::BOLD), + ) + .highlight_symbol("▶ "); + let mut state = ListState::default(); + if !self.renderers.is_empty() { + state.select(Some( + self.renderer_index + .min(self.renderers.len().saturating_sub(1)), + )); + } + f.render_stateful_widget(list, area, &mut state); + } + + fn draw_server_selection(&self, f: &mut ratatui::Frame<'_>) { + let area = f.size(); + let block = Block::default() + .borders(Borders::ALL) + .title("Sélection du serveur"); + let items: Vec = self + .servers + .iter() + .map(|info| { + let status = if info.online { + "[en ligne]" + } else { + "[hors ligne]" + }; + let text = format!("{status} {} | {}", info.friendly_name, info.model_name); + ListItem::new(text) + }) + .collect(); + let list = List::new(items) + .block(block) + .highlight_style( + Style::default() + .fg(Color::Cyan) + .add_modifier(Modifier::BOLD), + ) + .highlight_symbol("▶ "); + let mut state = ListState::default(); + if !self.servers.is_empty() { + state.select(Some( + self.server_index.min(self.servers.len().saturating_sub(1)), + )); + } + f.render_stateful_widget(list, area, &mut state); + } + + fn draw_browser(&self, f: &mut ratatui::Frame<'_>) { + let area = f.size(); + let Some(browser) = &self.browser else { + return; + }; + + let block = Block::default().borders(Borders::ALL).title(Span::styled( + format!( + "Navigation: {} (id: {})", + browser.nav_state.current_container_title, browser.nav_state.current_container_id + ), + Style::default() + .fg(Color::White) + .add_modifier(Modifier::BOLD), + )); + + let items: Vec = browser + .entries + .iter() + .map(|entry| { + let icon = if entry.is_container { "📁" } else { "♪" }; + let text = format!("{icon} {}", entry.title); + ListItem::new(text) + }) + .collect(); + + let list = List::new(items) + .block(block) + .highlight_style( + Style::default() + .fg(Color::Green) + .add_modifier(Modifier::BOLD), + ) + .highlight_symbol("▶ "); + let mut state = ListState::default(); + if !browser.entries.is_empty() { + state.select(Some( + browser + .selected_index + .min(browser.entries.len().saturating_sub(1)), + )); + } + f.render_stateful_widget(list, area, &mut state); + } + + fn draw_control_screen(&self, f: &mut ratatui::Frame<'_>) { + let chunks = Layout::default() + .direction(Direction::Vertical) + .constraints([ + Constraint::Length(5), + Constraint::Min(8), + Constraint::Length(3), + ]) + .split(f.size()); + + self.draw_header(f, chunks[0]); + self.draw_playback_panel(f, chunks[1]); + self.draw_help_strip(f, chunks[2]); + } + + fn draw_header(&self, f: &mut ratatui::Frame<'_>, area: Rect) { + let ui = &self.ui_state; + let renderer = &ui.renderer_name; + let server = ui.server_name.as_deref().unwrap_or(""); + let state = ui + .transport_state + .as_deref() + .unwrap_or("État inconnu") + .to_string(); + let volume = ui + .volume + .map(|v| v.to_string()) + .unwrap_or_else(|| "--".to_string()); + let mute = match ui.mute { + Some(true) => "ON", + Some(false) => "OFF", + None => "??", + }; + + let binding = self + .binding_info + .as_ref() + .map(|info| format!("{} / {}", info.server_id, info.container_id)) + .unwrap_or_else(|| "".to_string()); + + let text = vec![ + Line::from(vec![Span::styled( + format!("Renderer : {renderer}"), + Style::default().fg(Color::Yellow), + )]), + Line::from(vec![Span::raw(format!("Serveur : {server}"))]), + Line::from(vec![Span::raw(format!( + "État : {state} Volume {volume} | Mute {mute}" + ))]), + Line::from(vec![Span::raw(format!("Playlist : {binding}"))]), + ]; + + let paragraph = Paragraph::new(text) + .block(Block::default().borders(Borders::ALL).title("Statut")) + .alignment(Alignment::Left); + f.render_widget(paragraph, area); + } + + fn draw_playback_panel(&self, f: &mut ratatui::Frame<'_>, area: Rect) { + let chunks = Layout::default() + .direction(Direction::Vertical) + .constraints([Constraint::Min(6), Constraint::Length(3)]) + .split(area); + + let meta_lines = render_metadata_block(self.ui_state.metadata.as_ref()); + let paragraph = Paragraph::new(meta_lines) + .block( + Block::default() + .borders(Borders::ALL) + .title("Lecture en cours"), + ) + .wrap(ratatui::widgets::Wrap { trim: true }); + f.render_widget(paragraph, chunks[0]); + + let gauge = match self.ui_state.progress.as_ref() { + Some(progress) => build_progress_gauge(progress), + None => Gauge::default() + .block(Block::default().borders(Borders::ALL).title("Progression")) + .label("en attente...") + .ratio(0.0), + }; + f.render_widget(gauge, chunks[1]); + } + + fn draw_help_strip(&self, f: &mut ratatui::Frame<'_>, area: Rect) { + let lines = vec![ + Line::from("Commandes: h=Aide | R=Renderer | S=Serveur | B=Browse"), + Line::from(" Espace=Play/Pause s=Stop n=Next m=Mute +/-=Volume k=Queue"), + Line::from(" b=Binding dans le browser | q=Quit | ESC ferme les overlays"), + ]; + let paragraph = + Paragraph::new(lines).block(Block::default().borders(Borders::ALL).title("Raccourcis")); + f.render_widget(paragraph, area); + } + + fn draw_help_overlay(&self, f: &mut ratatui::Frame<'_>) { + let area = centered_rect(70, 60, f.size()); + let lines = vec![ + Line::from("Raccourcis disponibles:"), + Line::from(" R / S : re-sélectionner renderer / serveur"), + Line::from(" B : revenir au navigateur depuis l'écran principal"), + Line::from(" ↑/↓ ou +/- : ajuster le volume (via REST)"), + Line::from(" k / h : toggle queue ou aide"), + Line::from( + " Browser : Entrée=ouvrir, ←/Backspace ou r=retour, s ou b=sélectionner", + ), + Line::from(" Binding prompt : y=confirmer, n=annuler"), + Line::from(" ESC : fermer overlay courant"), + ]; + let block = Block::default() + .title("Aide détaillée (h pour fermer)") + .borders(Borders::ALL) + .style(Style::default().bg(Color::Black)); + f.render_widget(Clear, area); + f.render_widget(Paragraph::new(lines).block(block), area); + } + + fn draw_queue_overlay(&self, f: &mut ratatui::Frame<'_>) { + let area = centered_rect(80, 70, f.size()); + let mut lines = Vec::new(); + lines.push(Line::from("Playlist actuelle (▶ = en cours):")); + if self.queue_snapshot.is_empty() { + lines.push(Line::from(" ")); + } else { + for (idx, item) in self.queue_snapshot.iter().enumerate() { + let title = item.title.as_deref().unwrap_or(""); + let artist = item.artist.as_deref().unwrap_or(""); + let prefix = match self.queue_current_index { + Some(current) if current == idx => "▶", + _ => " ", + }; + let line = if artist.is_empty() { + format!("{prefix} [{idx}] {title}") + } else { + format!("{prefix} [{idx}] {artist} - {title}") + }; + lines.push(Line::from(line)); + } + } + let block = Block::default() + .title("Playlist (fermer avec k ou Esc)") + .borders(Borders::ALL) + .style(Style::default().bg(Color::Black)); + f.render_widget(Clear, area); + f.render_widget(Paragraph::new(lines).block(block), area); + } + + fn draw_binding_prompt(&self, f: &mut ratatui::Frame<'_>) { + let area = centered_rect(60, 40, f.size()); + let Some(binding) = &self.pending_binding else { + return; + }; + let renderer = self.ui_state.renderer_name.clone(); + let server = self + .selected_server + .as_ref() + .map(|s| s.friendly_name.clone()) + .unwrap_or_else(|| "".to_string()); + let lines = vec![ + Line::from("Attacher ce conteneur comme playlist distante ?"), + Line::from(format!("Renderer: {renderer}")), + Line::from(format!("Serveur : {server}")), + Line::from(format!( + "Container: {} ({})", + binding.container_title, binding.container_id + )), + Line::from("y = oui | n = non"), + ]; + let block = Block::default() + .title("Binding playlist") + .borders(Borders::ALL) + .style(Style::default().bg(Color::Black)); + f.render_widget(Clear, area); + f.render_widget(Paragraph::new(lines).block(block), area); + } + + fn draw_status_line(&self, f: &mut ratatui::Frame<'_>) { + let area = Rect { + x: 0, + y: f.size().height.saturating_sub(1), + width: f.size().width, + height: 1, + }; + let status = self + .ui_state + .last_status + .clone() + .unwrap_or_else(|| self.status_line.clone()); + let paragraph = Paragraph::new(status).style(Style::default().fg(Color::Gray)); + f.render_widget(paragraph, area); + } + + fn handle_key(&mut self, key: KeyEvent) -> Result { + match self.mode { + Mode::SelectRenderer => self.handle_renderer_key(key), + Mode::SelectServer => self.handle_server_key(key), + Mode::Browse => self.handle_browse_key(key), + Mode::BindingPrompt => self.handle_binding_key(key), + Mode::Control => self.handle_control_key(key), + } + } + + fn handle_renderer_key(&mut self, key: KeyEvent) -> Result { + match key.code { + KeyCode::Char('q') => return Ok(true), + KeyCode::Esc => return Ok(true), + KeyCode::Up => { + if self.renderer_index > 0 { + self.renderer_index -= 1; + } + } + KeyCode::Down => { + if self.renderer_index + 1 < self.renderers.len() { + self.renderer_index += 1; + } + } + KeyCode::Enter => { + let info = self.renderers[self.renderer_index].clone(); + let previous_renderer_id = self + .selected_renderer + .as_ref() + .map(|renderer| renderer.id.clone()); + let mut stop_error: Option = None; + if let Some(prev_id) = previous_renderer_id { + if prev_id != info.id { + if let Err(err) = self.client.stop(&prev_id) { + stop_error = Some(format!( + "Renderer sélectionné mais arrêt de l'ancien impossible: {err}" + )); + } + } + } + self.selected_renderer = Some(info.clone()); + self.ui_state = UiState::new(info.friendly_name.clone()); + self.queue_snapshot.clear(); + self.queue_current_index = None; + self.binding_info = None; + self.selected_server = None; + self.browser = None; + self.show_help_overlay = false; + self.show_queue_overlay = false; + self.pending_binding = None; + self.mode = Mode::SelectServer; + self.status_line = "Sélectionne un serveur avec ↑/↓ puis Entrée".to_string(); + self.load_servers(); + self.refresh_renderer_state(); + self.refresh_queue(); + self.refresh_binding_info(); + if let Some(message) = stop_error { + self.ui_state.set_status(message); + } else { + self.ui_state.set_status("Renderer sélectionné"); + } + } + _ => {} + } + Ok(false) + } + + fn handle_server_key(&mut self, key: KeyEvent) -> Result { + match key.code { + KeyCode::Char('q') => return Ok(true), + KeyCode::Esc => { + self.mode = Mode::SelectRenderer; + self.status_line = "Sélectionne un renderer avec ↑/↓ puis Entrée".to_string(); + } + KeyCode::Up => { + if self.server_index > 0 { + self.server_index -= 1; + } + } + KeyCode::Down => { + if self.server_index + 1 < self.servers.len() { + self.server_index += 1; + } + } + KeyCode::Enter => { + if self.servers.is_empty() { + self.ui_state + .set_status("Aucun serveur disponible. Vérifie le backend."); + return Ok(false); + } + let info = self.servers[self.server_index].clone(); + self.selected_server = Some(info.clone()); + self.ui_state.server_name = Some(info.friendly_name.clone()); + match self.load_browser_for_server(&info) { + Ok(_) => { + self.mode = Mode::Browse; + self.status_line = + "Navigue avec ↑/↓, Entrée pour ouvrir, ←/Backspace ou r pour remonter, s pour sélectionner" + .to_string(); + self.ui_state.set_status("Serveur sélectionné"); + } + Err(err) => { + self.ui_state + .set_status(format!("Navigation impossible: {err}")); + } + } + } + _ => {} + } + Ok(false) + } + + fn handle_browse_key(&mut self, key: KeyEvent) -> Result { + if key.code == KeyCode::Char('q') { + return Ok(true); + } + if self.browser.is_none() { + return Ok(false); + } + + match key.code { + KeyCode::Up => { + if let Some(browser) = self.browser.as_mut() { + if browser.selected_index > 0 { + browser.selected_index -= 1; + } + } + } + KeyCode::Down => { + if let Some(browser) = self.browser.as_mut() { + if browser.selected_index + 1 < browser.entries.len() { + browser.selected_index += 1; + } + } + } + KeyCode::Enter => { + let entry = self + .browser + .as_ref() + .and_then(|b| b.current_entry().cloned()); + if let Some(entry) = entry { + if entry.is_container { + self.enter_container(entry)?; + } else { + self.ui_state + .set_status("Sélectionne un dossier pour le binding."); + } + } + } + KeyCode::Char('s') | KeyCode::Char('b') => { + let entry = self + .browser + .as_ref() + .and_then(|b| b.current_entry().cloned()); + if let Some(entry) = entry { + if entry.is_container { + self.pending_binding = Some(PendingBinding { + container_id: entry.id, + container_title: entry.title, + }); + self.mode = Mode::BindingPrompt; + self.ui_state.set_status("Confirme le binding (y/n)"); + } else { + self.ui_state + .set_status("Impossible de binder un item individuel."); + } + } + } + KeyCode::Left | KeyCode::Backspace => { + self.navigate_browser_up(); + } + KeyCode::Char('r') | KeyCode::Char('R') => { + self.navigate_browser_up(); + } + KeyCode::Char('h') => { + self.show_help_overlay = !self.show_help_overlay; + } + KeyCode::Esc => { + self.mode = Mode::SelectServer; + self.status_line = "Sélectionne un serveur avec ↑/↓ puis Entrée".to_string(); + self.show_help_overlay = false; + self.show_queue_overlay = false; + } + _ => {} + } + Ok(false) + } + + fn handle_binding_key(&mut self, key: KeyEvent) -> Result { + match key.code { + KeyCode::Char('y') => { + self.show_help_overlay = false; + self.show_queue_overlay = false; + self.attach_binding(true)?; + } + KeyCode::Char('n') | KeyCode::Esc => { + self.show_help_overlay = false; + self.show_queue_overlay = false; + self.attach_binding(false)?; + } + KeyCode::Char('q') => return Ok(true), + _ => {} + } + Ok(false) + } + + fn handle_control_key(&mut self, key: KeyEvent) -> Result { + match key.code { + KeyCode::Char('q') => return Ok(true), + KeyCode::Char('R') => { + self.open_renderer_menu(); + } + KeyCode::Char('S') => { + self.open_server_menu(); + } + KeyCode::Char('B') => { + if self.browser.is_some() { + self.mode = Mode::Browse; + self.status_line = + "Navigue avec ↑/↓, Entrée pour ouvrir, ←/Backspace ou r pour remonter, s pour sélectionner" + .to_string(); + } else { + self.ui_state + .set_status("Pas de navigateur actif. Reprends la sélection serveur."); + } + } + KeyCode::Char('h') => { + self.show_help_overlay = !self.show_help_overlay; + } + KeyCode::Char('k') => { + self.show_queue_overlay = !self.show_queue_overlay; + } + KeyCode::Esc => { + self.show_queue_overlay = false; + self.show_help_overlay = false; + } + KeyCode::Char(' ') => { + self.toggle_play_pause()?; + } + KeyCode::Char('p') => { + self.pause_renderer()?; + } + KeyCode::Char('s') => { + self.stop_renderer()?; + } + KeyCode::Char('n') => { + self.play_next()?; + } + KeyCode::Char('+') | KeyCode::Char('=') | KeyCode::Up => { + self.volume_up()?; + } + KeyCode::Char('-') | KeyCode::Down => { + self.volume_down()?; + } + KeyCode::Char('m') => { + self.toggle_mute()?; + } + _ => {} + } + Ok(false) + } + + fn load_servers(&mut self) { + match self.client.list_servers() { + Ok(list) => { + self.servers = list; + self.server_index = 0; + if self.servers.is_empty() { + self.ui_state + .set_status("Aucun serveur disponible. Vérifie PMOMusic."); + } + } + Err(err) => { + self.ui_state + .set_status(format!("Erreur REST serveurs: {err}")); + } + } + } + + fn load_browser_for_server(&mut self, info: &MediaServerSummaryClient) -> Result<()> { + let (root_id, entries) = self.fetch_root_entries(&info.id)?; + self.browser = Some(BrowserState::new( + info.id.clone(), + root_id, + entries, + info.friendly_name.clone(), + )); + Ok(()) + } + + fn fetch_root_entries(&self, server_id: &str) -> Result<(String, Vec)> { + let mut last_err: Option = None; + for &candidate in ROOT_CONTAINERS { + match self.client.browse_container(server_id, candidate) { + Ok(resp) => return Ok((resp.container_id, resp.entries)), + Err(err) => last_err = Some(err), + } + } + Err(last_err.unwrap_or_else(|| { + anyhow!("Impossible de parcourir la racine pour le serveur {server_id}") + })) + } + + fn enter_container(&mut self, entry: ContainerEntryClient) -> Result<()> { + let Some(browser) = self.browser.as_mut() else { + return Ok(()); + }; + let container_id = entry.id.clone(); + let container_title = entry.title.clone(); + match self + .client + .browse_container(&browser.server_id, &container_id) + { + Ok(resp) => { + browser + .nav_state + .enter_container(container_id, container_title); + browser.entries = resp.entries; + browser.selected_index = 0; + } + Err(err) => { + self.ui_state + .set_status(format!("Impossible d'ouvrir: {err}")); + } + } + Ok(()) + } + + fn navigate_browser_up(&mut self) { + let (server_id, container_id) = { + let Some(browser) = self.browser.as_mut() else { + return; + }; + if browser.nav_state.go_back() { + ( + browser.server_id.clone(), + browser.nav_state.current_container_id.clone(), + ) + } else { + self.ui_state.set_status("Déjà à la racine."); + return; + } + }; + match self.client.browse_container(&server_id, &container_id) { + Ok(resp) => { + if let Some(browser) = self.browser.as_mut() { + browser.entries = resp.entries; + browser.selected_index = 0; + } + } + Err(err) => { + self.ui_state + .set_status(format!("Retour impossible: {err}")); + } + } + } + + fn attach_binding(&mut self, attach: bool) -> Result<()> { + if attach { + if self.binding_worker.is_some() { + self.ui_state + .set_status("Binding déjà en cours. Patiente quelques secondes..."); + return Ok(()); + } + let Some(renderer) = self.selected_renderer.as_ref() else { + self.ui_state.set_status("Choisis un renderer en premier."); + return Ok(()); + }; + let Some(server) = self.selected_server.as_ref() else { + self.ui_state.set_status("Choisis un serveur en premier."); + return Ok(()); + }; + let Some(binding) = self.pending_binding.take() else { + return Ok(()); + }; + let client = self.client.clone(); + let renderer_id = renderer.id.clone(); + let server_id = server.id.clone(); + let container_id = binding.container_id.clone(); + let container_title = binding.container_title.clone(); + let (tx, rx) = mpsc::channel(); + self.binding_worker = Some(rx); + self.mode = Mode::Browse; + self.status_line = + "Binding en cours... patiente pendant la préparation de la playlist".to_string(); + self.show_help_overlay = false; + self.show_queue_overlay = false; + self.ui_state + .set_status(format!("Association en cours: {}", container_title)); + thread::spawn(move || { + let outcome = (|| -> Result<()> { + client.attach_playlist(&renderer_id, &server_id, &container_id)?; + // Préparer la lecture en sélectionnant le premier item + client.next(&renderer_id)?; + Ok(()) + })(); + let message = match outcome { + Ok(_) => BindingWorkerMessage::Success { container_title }, + Err(err) => BindingWorkerMessage::Failure { + error: err.to_string(), + }, + }; + let _ = tx.send(message); + }); + } else { + self.pending_binding = None; + self.mode = Mode::Browse; + self.status_line = + "Navigue avec ↑/↓, Entrée pour ouvrir, ←/Backspace ou r pour remonter, s pour sélectionner".to_string(); + self.ui_state.set_status("Binding annulé"); + } + Ok(()) + } + + fn toggle_play_pause(&mut self) -> Result<()> { + let Some(renderer) = self.selected_renderer.as_ref() else { + return Ok(()); + }; + let current = self + .ui_state + .transport_state + .as_deref() + .unwrap_or("") + .to_string(); + if current.eq_ignore_ascii_case("PLAYING") { + self.client.pause(&renderer.id)?; + self.ui_state.set_status("Pause envoyée"); + } else { + self.client.play(&renderer.id)?; + self.ui_state.set_status("Lecture envoyée"); + } + Ok(()) + } + + fn pause_renderer(&mut self) -> Result<()> { + if let Some(renderer) = self.selected_renderer.as_ref() { + self.client.pause(&renderer.id)?; + self.ui_state.set_status("Pause envoyée"); + } + Ok(()) + } + + fn stop_renderer(&mut self) -> Result<()> { + if let Some(renderer) = self.selected_renderer.as_ref() { + self.client.stop(&renderer.id)?; + self.ui_state.set_status("Stop envoyé"); + } + Ok(()) + } + + fn play_next(&mut self) -> Result<()> { + if let Some(renderer) = self.selected_renderer.as_ref() { + self.client.next(&renderer.id)?; + self.ui_state.set_status("Piste suivante demandée"); + } + Ok(()) + } + + fn volume_up(&mut self) -> Result<()> { + if let Some(renderer) = self.selected_renderer.as_ref() { + self.client.volume_up(&renderer.id)?; + self.ui_state.set_status("Volume +"); + } + Ok(()) + } + + fn volume_down(&mut self) -> Result<()> { + if let Some(renderer) = self.selected_renderer.as_ref() { + self.client.volume_down(&renderer.id)?; + self.ui_state.set_status("Volume -"); + } + Ok(()) + } + + fn toggle_mute(&mut self) -> Result<()> { + if let Some(renderer) = self.selected_renderer.as_ref() { + self.client.toggle_mute(&renderer.id)?; + self.ui_state.set_status("Mute togglé"); + } + Ok(()) + } + + fn open_renderer_menu(&mut self) { + match self.client.list_renderers() { + Ok(list) => { + self.renderers = list; + self.renderer_index = 0; + self.mode = Mode::SelectRenderer; + self.status_line = "Sélectionne un renderer avec ↑/↓ puis Entrée".to_string(); + self.show_help_overlay = false; + self.show_queue_overlay = false; + } + Err(err) => { + self.ui_state + .set_status(format!("Impossible de rafraîchir les renderers: {err}")); + } + } + } + + fn open_server_menu(&mut self) { + if self.selected_renderer.is_none() { + self.ui_state.set_status("Sélectionne d'abord un renderer."); + return; + } + self.load_servers(); + self.mode = Mode::SelectServer; + self.status_line = "Sélectionne un serveur avec ↑/↓ puis Entrée".to_string(); + self.show_help_overlay = false; + self.show_queue_overlay = false; + } + + fn refresh_renderer_state(&mut self) { + let Some(renderer) = self.selected_renderer.as_ref() else { + return; + }; + match self.client.get_renderer_state(&renderer.id) { + Ok(state) => { + self.ui_state.transport_state = Some(state.transport_state.clone()); + self.ui_state.volume = state.volume; + self.ui_state.mute = state.mute; + self.ui_state.progress = Some(PlaybackProgress { + position_ms: state.position_ms, + duration_ms: state.duration_ms, + }); + if let Some(info) = state.attached_playlist { + self.binding_info = Some(info); + } + } + Err(err) => { + self.ui_state + .set_status(format!("Erreur REST renderer: {err}")); + } + } + } + + fn refresh_queue(&mut self) { + let Some(renderer) = self.selected_renderer.as_ref() else { + return; + }; + let current_signature = self.capture_current_queue_signature(); + match self.client.get_renderer_queue(&renderer.id) { + Ok(snapshot) => { + let next_items = snapshot.items; + let mut next_index = snapshot.current_index; + if !Self::is_valid_queue_index(next_index, &next_items) { + next_index = current_signature + .and_then(|sig| Self::find_queue_index_by_signature(&next_items, &sig)); + } + self.queue_snapshot = next_items; + self.queue_current_index = next_index; + self.update_current_track_metadata(); + } + Err(err) => { + self.ui_state + .set_status(format!("Erreur REST queue: {err}")); + } + } + } + + fn capture_current_queue_signature(&self) -> Option { + let idx = self.queue_current_index?; + let item = self.queue_snapshot.get(idx)?; + Some(QueueItemSignature::from_item(item)) + } + + fn is_valid_queue_index(index: Option, items: &[QueueItemClient]) -> bool { + match index { + Some(idx) => idx < items.len(), + None => false, + } + } + + fn find_queue_index_by_signature( + items: &[QueueItemClient], + signature: &QueueItemSignature, + ) -> Option { + items.iter().enumerate().find_map(|(idx, item)| { + if signature.matches(item) { + Some(idx) + } else { + None + } + }) + } + + fn refresh_binding_info(&mut self) { + let Some(renderer) = self.selected_renderer.as_ref() else { + return; + }; + match self.client.get_renderer_binding(&renderer.id) { + Ok(binding) => { + self.binding_info = binding; + } + Err(err) => { + self.ui_state + .set_status(format!("Erreur REST binding: {err}")); + } + } + } + + fn update_current_track_metadata(&mut self) { + if let Some(idx) = self.queue_current_index { + if let Some(item) = self.queue_snapshot.get(idx) { + self.ui_state.metadata = Some(TrackMetadata { + title: item.title.clone(), + artist: item.artist.clone(), + album: item.album.clone(), + album_art_uri: item.album_art_uri.clone(), + }); + return; + } + } + self.ui_state.metadata = None; + } + + fn on_tick(&mut self) { + self.poll_binding_worker(); + if self.mode != Mode::Control { + return; + } + self.refresh_renderer_state(); + self.refresh_queue(); + self.refresh_binding_info(); + } + + fn poll_binding_worker(&mut self) { + let Some(receiver) = self.binding_worker.as_ref() else { + return; + }; + match receiver.try_recv() { + Ok(BindingWorkerMessage::Success { container_title }) => { + self.binding_worker = None; + self.finish_binding_success(container_title); + } + Ok(BindingWorkerMessage::Failure { error }) => { + self.binding_worker = None; + self.finish_binding_failure(error); + } + Err(TryRecvError::Empty) => {} + Err(TryRecvError::Disconnected) => { + self.binding_worker = None; + self.finish_binding_failure("Worker binding interrompu (canal fermé)".to_string()); + } + } + } + + fn finish_binding_success(&mut self, container_title: String) { + self.mode = Mode::Control; + self.status_line = "Espace=Play/Pause, n=Next, +/- volume, k=Queue, B=Browse".to_string(); + self.show_help_overlay = false; + self.show_queue_overlay = false; + self.ui_state + .set_status(format!("Lecture lancée depuis {container_title}")); + self.refresh_binding_info(); + self.refresh_queue(); + self.refresh_renderer_state(); + } + + fn finish_binding_failure(&mut self, error: String) { + self.mode = Mode::Browse; + self.status_line = + "Navigue avec ↑/↓, Entrée pour ouvrir, ←/Backspace ou r pour remonter, s pour sélectionner".to_string(); + self.show_help_overlay = false; + self.show_queue_overlay = false; + self.ui_state + .set_status(format!("Binding impossible: {error}")); + } +} + +impl UiState { + fn new(renderer_name: String) -> Self { + Self { + renderer_name, + server_name: None, + transport_state: None, + progress: None, + volume: None, + mute: None, + metadata: None, + last_status: Some("Interface initialisée.".to_string()), + } + } + + fn placeholder() -> Self { + Self::new("".to_string()) + } + + fn set_status>(&mut self, status: S) { + self.last_status = Some(status.into()); + } +} + +impl BrowserState { + fn new( + server_id: String, + root_container_id: String, + entries: Vec, + friendly_name: String, + ) -> Self { + Self { + server_id, + nav_state: NavigationState::new(root_container_id, friendly_name), + entries, + selected_index: 0, + } + } + + fn current_entry(&self) -> Option<&ContainerEntryClient> { + self.entries.get(self.selected_index) + } +} + +impl NavigationState { + fn new(root_id: String, root_title: String) -> Self { + Self { + path_stack: Vec::new(), + current_container_id: root_id, + current_container_title: root_title, + } + } + + fn enter_container(&mut self, container_id: String, container_title: String) { + self.path_stack.push(( + self.current_container_id.clone(), + self.current_container_title.clone(), + )); + self.current_container_id = container_id; + self.current_container_title = container_title; + } + + fn go_back(&mut self) -> bool { + if let Some((parent_id, parent_title)) = self.path_stack.pop() { + self.current_container_id = parent_id; + self.current_container_title = parent_title; + true + } else { + false + } + } +} + +fn run_app(mut app: App) -> Result<()> { + let terminal = setup_terminal()?; + let mut guard = TerminalGuard { terminal }; + let mut last_tick = Instant::now(); + + let result = (|| -> Result<()> { + loop { + guard.terminal.draw(|f| app.draw(f))?; + + let timeout = TICK_RATE + .checked_sub(last_tick.elapsed()) + .unwrap_or_else(|| Duration::from_secs(0)); + + if event::poll(timeout)? { + if let Event::Key(key) = event::read()? { + if app.handle_key(key)? { + break; + } + } + } + + if last_tick.elapsed() >= TICK_RATE { + app.on_tick(); + last_tick = Instant::now(); + } + } + Ok(()) + })(); + + // Restauration garantie via Drop de TerminalGuard + // On tente un stop avec timeout court + if let Some(renderer) = app.selected_renderer.as_ref() { + // Utiliser un timeout très court pour ne pas bloquer le shutdown + let _ = app.client.stop(&renderer.id); + } + + result +} + +fn setup_terminal() -> Result>> { + enable_raw_mode()?; + let mut stdout = io::stdout(); + execute!(stdout, EnterAlternateScreen, EnableMouseCapture)?; + let backend = CrosstermBackend::new(stdout); + let terminal = Terminal::new(backend)?; + Ok(terminal) +} + +fn render_metadata_block(metadata: Option<&TrackMetadata>) -> Vec> { + let mut lines = Vec::new(); + if let Some(meta) = metadata { + let title = meta + .title + .clone() + .unwrap_or_else(|| "".to_string()); + lines.push(Line::from(format!("Titre : {title}"))); + if let Some(artist) = meta.artist.as_deref() { + lines.push(Line::from(format!("Artiste: {artist}"))); + } + if let Some(album) = meta.album.as_deref() { + lines.push(Line::from(format!("Album : {album}"))); + } + if let Some(art) = meta.album_art_uri.as_deref() { + lines.push(Line::from(format!("Cover : {art}"))); + } + } else { + lines.push(Line::from("(En attente des métadonnées...)")); + } + lines +} + +fn build_progress_gauge(progress: &PlaybackProgress) -> Gauge<'static> { + let ratio = match (progress.position_ms, progress.duration_ms) { + (Some(pos), Some(dur)) if dur > 0 => (pos as f64 / dur as f64).clamp(0.0, 1.0), + _ => 0.0, + }; + let label = format!( + "{} / {}", + progress + .position_ms + .and_then(format_time_ms) + .unwrap_or_else(|| "--:--".to_string()), + progress + .duration_ms + .and_then(format_time_ms) + .unwrap_or_else(|| "--:--".to_string()) + ); + Gauge::default() + .block(Block::default().borders(Borders::ALL).title("Progression")) + .gauge_style(Style::default().fg(Color::Magenta)) + .ratio(ratio) + .label(label) +} + +fn format_time_ms(ms: u64) -> Option { + let total_seconds = ms / 1000; + let hours = total_seconds / 3600; + let minutes = (total_seconds % 3600) / 60; + let seconds = total_seconds % 60; + if hours > 0 { + Some(format!("{hours:02}:{minutes:02}:{seconds:02}")) + } else { + Some(format!("{minutes:02}:{seconds:02}")) + } +} + +fn centered_rect(percent_x: u16, percent_y: u16, r: Rect) -> Rect { + let popup_layout = Layout::default() + .direction(Direction::Vertical) + .constraints( + [ + Constraint::Percentage((100 - percent_y) / 2), + Constraint::Percentage(percent_y), + Constraint::Percentage((100 - percent_y) / 2), + ] + .as_ref(), + ) + .split(r); + + Layout::default() + .direction(Direction::Horizontal) + .constraints( + [ + Constraint::Percentage((100 - percent_x) / 2), + Constraint::Percentage(percent_x), + Constraint::Percentage((100 - percent_x) / 2), + ] + .as_ref(), + ) + .split(popup_layout[1])[1] +} + +/// RAII guard pour garantir la restauration du terminal même en cas d'erreur ou de panic +struct TerminalGuard { + terminal: Terminal>, +} + +impl Drop for TerminalGuard { + fn drop(&mut self) { + // Force la restauration du terminal, même si les appels échouent + let _ = disable_raw_mode(); + let _ = execute!( + self.terminal.backend_mut(), + LeaveAlternateScreen, + DisableMouseCapture + ); + let _ = self.terminal.show_cursor(); + } +} + +// ============================================================================ +// REST DTOs +// ============================================================================ + +#[allow(dead_code)] +#[derive(Debug, Clone, Deserialize)] +struct RendererSummaryClient { + id: String, + friendly_name: String, + model_name: String, + protocol: String, + online: bool, +} + +#[allow(dead_code)] +#[derive(Debug, Clone, Deserialize)] +struct RendererStateClient { + id: String, + friendly_name: String, + transport_state: String, + position_ms: Option, + duration_ms: Option, + volume: Option, + mute: Option, + queue_len: usize, + attached_playlist: Option, +} + +#[allow(dead_code)] +#[derive(Debug, Clone, Deserialize)] +struct AttachedPlaylistInfoClient { + server_id: String, + container_id: String, + has_seen_update: bool, +} + +#[allow(dead_code)] +#[derive(Debug, Clone, Deserialize)] +struct QueueItemClient { + index: usize, + uri: String, + title: Option, + artist: Option, + album: Option, + server_id: Option, + object_id: Option, + album_art_uri: Option, +} + +#[allow(dead_code)] +#[derive(Debug, Clone, Deserialize)] +struct QueueSnapshotClient { + renderer_id: String, + #[serde(default, deserialize_with = "deserialize_nullable_vec")] + items: Vec, + current_index: Option, +} + +#[derive(Debug, Clone)] +struct QueueItemSignature { + server_id: Option, + object_id: Option, + uri: String, +} + +impl QueueItemSignature { + fn from_item(item: &QueueItemClient) -> Self { + Self { + server_id: item.server_id.clone(), + object_id: item.object_id.clone(), + uri: item.uri.clone(), + } + } + + fn matches(&self, other: &QueueItemClient) -> bool { + if let (Some(sig_obj), Some(other_obj)) = (&self.object_id, &other.object_id) { + if sig_obj == other_obj { + if let (Some(sig_server), Some(other_server)) = (&self.server_id, &other.server_id) + { + return sig_server == other_server; + } + return true; + } + } + self.uri == other.uri + } +} + +#[allow(dead_code)] +#[derive(Debug, Clone, Deserialize)] +struct MediaServerSummaryClient { + id: String, + friendly_name: String, + model_name: String, + online: bool, +} + +#[allow(dead_code)] +#[derive(Debug, Clone, Deserialize)] +struct ContainerEntryClient { + id: String, + title: String, + class: String, + is_container: bool, + child_count: Option, + artist: Option, + album: Option, + album_art_uri: Option, +} + +#[allow(dead_code)] +#[derive(Debug, Clone, Deserialize)] +struct BrowseResponseClient { + container_id: String, + #[serde(default, deserialize_with = "deserialize_nullable_vec")] + entries: Vec, +} + +#[allow(dead_code)] +#[derive(Debug, Serialize)] +struct VolumeSetRequest { + volume: u8, +} + +#[derive(Debug, Serialize)] +struct AttachPlaylistRequest<'a> { + server_id: &'a str, + container_id: &'a str, +} + +fn deserialize_nullable_vec<'de, D, T>(deserializer: D) -> Result, D::Error> +where + D: Deserializer<'de>, + T: Deserialize<'de>, +{ + let opt = Option::>::deserialize(deserializer)?; + Ok(opt.unwrap_or_default()) +} + +// ============================================================================ +// REST CLIENT +// ============================================================================ + +struct RestClient { + base_url: String, + agent: Agent, +} + +impl Clone for RestClient { + fn clone(&self) -> Self { + Self { + base_url: self.base_url.clone(), + agent: self.agent.clone(), + } + } +} + +impl RestClient { + fn new(base_url: &str, timeout: Duration) -> Result { + let mut builder = Agent::config_builder(); + builder = builder.timeout_global(Some(timeout)); + builder = builder.http_status_as_error(false); + let config = builder.build(); + let agent: Agent = config.into(); + Ok(Self { + base_url: base_url.trim_end_matches('/').to_string(), + agent, + }) + } + + fn list_renderers(&self) -> Result> { + self.get_json(&["renderers"]) + } + + fn get_renderer_state(&self, id: &str) -> Result { + self.get_json(&["renderers", id]) + } + + fn get_renderer_queue(&self, id: &str) -> Result { + self.get_json(&["renderers", id, "queue"]) + } + + fn get_renderer_binding(&self, id: &str) -> Result> { + self.get_json(&["renderers", id, "binding"]) + } + + fn play(&self, id: &str) -> Result<()> { + self.post_empty(&["renderers", id, "play"]) + } + + fn pause(&self, id: &str) -> Result<()> { + self.post_empty(&["renderers", id, "pause"]) + } + + fn stop(&self, id: &str) -> Result<()> { + self.post_empty(&["renderers", id, "stop"]) + } + + fn next(&self, id: &str) -> Result<()> { + self.post_empty(&["renderers", id, "next"]) + } + + #[allow(dead_code)] + fn set_volume(&self, id: &str, volume: u8) -> Result<()> { + let payload = VolumeSetRequest { volume }; + self.post_json(&["renderers", id, "volume", "set"], &payload) + } + + fn volume_up(&self, id: &str) -> Result<()> { + self.post_empty(&["renderers", id, "volume", "up"]) + } + + fn volume_down(&self, id: &str) -> Result<()> { + self.post_empty(&["renderers", id, "volume", "down"]) + } + + fn toggle_mute(&self, id: &str) -> Result<()> { + self.post_empty(&["renderers", id, "mute", "toggle"]) + } + + fn attach_playlist(&self, id: &str, server_id: &str, container_id: &str) -> Result<()> { + let payload = AttachPlaylistRequest { + server_id, + container_id, + }; + self.post_json(&["renderers", id, "binding", "attach"], &payload) + } + + #[allow(dead_code)] + fn detach_playlist(&self, id: &str) -> Result<()> { + self.post_empty(&["renderers", id, "binding", "detach"]) + } + + fn list_servers(&self) -> Result> { + self.get_json(&["servers"]) + } + + fn browse_container( + &self, + server_id: &str, + container_id: &str, + ) -> Result { + self.get_json(&["servers", server_id, "containers", container_id]) + } + + fn get_json(&self, segments: &[&str]) -> Result + where + T: DeserializeOwned, + { + let url = self.build_url_segments(segments); + let response = self.agent.get(&url).call(); + let mut response = Self::handle_response(response)?; + let text = response + .body_mut() + .read_to_string() + .with_context(|| format!("Échec de lecture JSON depuis {url}"))?; + let value = serde_json::from_str(&text) + .with_context(|| format!("Échec de parsing JSON depuis {url}"))?; + Ok(value) + } + + fn post_empty(&self, segments: &[&str]) -> Result<()> { + let url = self.build_url_segments(segments); + let response = self.agent.post(&url).send_empty(); + Self::handle_response(response)?; + Ok(()) + } + + fn post_json(&self, segments: &[&str], payload: &T) -> Result<()> + where + T: Serialize, + { + let url = self.build_url_segments(segments); + let body = serde_json::to_vec(payload)?; + let response = self + .agent + .post(&url) + .header("content-type", "application/json") + .send(body); + Self::handle_response(response)?; + Ok(()) + } + + fn build_url_segments(&self, segments: &[&str]) -> String { + let mut url = self.base_url.clone(); + for segment in segments { + url.push('/'); + url.push_str(&utf8_percent_encode(segment, NON_ALPHANUMERIC).to_string()); + } + url + } + + fn handle_response( + response: Result, ureq::Error>, + ) -> Result> { + match response { + Ok(resp) => { + if resp.status().is_success() { + Ok(resp) + } else { + let mut resp = resp; + let status = resp.status(); + let body = resp + .body_mut() + .read_to_string() + .unwrap_or_else(|_| "".into()); + Err(anyhow!("HTTP {}: {}", status, body)) + } + } + Err(err) => Err(anyhow!(err)), + } + } +} diff --git a/pmocontrol/examples/pmomusic_integration_example.rs b/pmocontrol/examples/pmomusic_integration_example.rs index c840e523..abcf96c0 100644 --- a/pmocontrol/examples/pmomusic_integration_example.rs +++ b/pmocontrol/examples/pmomusic_integration_example.rs @@ -5,9 +5,7 @@ #[cfg(not(feature = "pmoserver"))] fn main() { - eprintln!( - "This example requires the 'pmoserver' feature. Re-run with `--features pmoserver`." - ); + eprintln!("This example requires the 'pmoserver' feature. Re-run with `--features pmoserver`."); } #[cfg(feature = "pmoserver")] diff --git a/pmocontrol/src/control_point.rs b/pmocontrol/src/control_point.rs index 490a0cc2..6e7e17ce 100644 --- a/pmocontrol/src/control_point.rs +++ b/pmocontrol/src/control_point.rs @@ -45,6 +45,8 @@ struct PlaylistBinding { /// Flag used internally to signal that the queue should be refreshed /// from the server container. pending_refresh: bool, + /// Whether the next refresh should auto-start playback if the renderer is idle. + auto_play_on_refresh: bool, } /// Control point minimal : @@ -326,6 +328,7 @@ impl ControlPoint { &bindings_for_media_worker, &renderer_id, &event_bus_for_media_worker, + None, ) { warn!( renderer = renderer_id.0.as_str(), @@ -380,6 +383,7 @@ impl ControlPoint { &bindings_for_periodic, &renderer_id, &event_bus_for_periodic, + None, ) { warn!( renderer = renderer_id.0.as_str(), @@ -718,6 +722,29 @@ impl ControlPoint { Ok(()) } + fn start_queue_playback_if_idle(&self, renderer_id: &RendererId) -> anyhow::Result<()> { + let snapshot = self.runtime.snapshot_for(renderer_id); + let renderer_playing = matches!(snapshot.state, Some(PlaybackState::Playing)); + if renderer_playing || self.runtime.is_playing_from_queue(renderer_id) { + return Ok(()); + } + + let has_items = self + .runtime + .queue_snapshot(renderer_id) + .map(|items| !items.is_empty()) + .unwrap_or(false); + if !has_items { + debug!( + renderer = renderer_id.0.as_str(), + "start_queue_playback_if_idle: queue is empty" + ); + return Ok(()); + } + + self.play_next_from_queue(renderer_id) + } + /// Subscribe to renderer events emitted by the control point runtime. /// /// Each subscriber receives all future events independently. @@ -747,22 +774,43 @@ impl ControlPoint { server_id: ServerId, container_id: String, ) { - let mut bindings = self.playlist_bindings.lock().unwrap(); - bindings.insert( - renderer_id.clone(), - PlaylistBinding { - server_id: server_id.clone(), - container_id: container_id.clone(), - has_seen_update: false, - pending_refresh: false, - }, - ); - info!( - renderer = renderer_id.0.as_str(), - server = server_id.0.as_str(), - container = container_id.as_str(), - "Queue attached to playlist container" - ); + { + let mut bindings = self.playlist_bindings.lock().unwrap(); + bindings.insert( + renderer_id.clone(), + PlaylistBinding { + server_id: server_id.clone(), + container_id: container_id.clone(), + has_seen_update: false, + pending_refresh: true, + auto_play_on_refresh: true, + }, + ); + info!( + renderer = renderer_id.0.as_str(), + server = server_id.0.as_str(), + container = container_id.as_str(), + "Queue attached to playlist container" + ); + } // Drop bindings lock here before calling refresh_attached_queue_for + + let mut auto_start_cb = |rid: &RendererId| self.start_queue_playback_if_idle(rid); + if let Err(err) = refresh_attached_queue_for( + &self.registry, + &self.runtime, + &self.playlist_bindings, + renderer_id, + &self.event_bus, + Some(&mut auto_start_cb), + ) { + warn!( + renderer = renderer_id.0.as_str(), + server = server_id.0.as_str(), + container = container_id.as_str(), + error = %err, + "Initial playlist refresh after attachment failed" + ); + } } /// Detach a renderer's queue from its associated playlist container. @@ -995,9 +1043,10 @@ fn refresh_attached_queue_for( bindings: &Arc>>, renderer_id: &RendererId, event_bus: &RendererEventBus, + mut after_refresh: Option<&mut dyn FnMut(&RendererId) -> anyhow::Result<()>>, ) -> anyhow::Result<()> { // Step 1: Check binding and mark refresh as in-progress - let (server_id, container_id) = { + let (server_id, container_id, auto_play) = { let mut bindings_lock = bindings.lock().unwrap(); let binding = match bindings_lock.get_mut(renderer_id) { Some(b) => b, @@ -1020,7 +1069,13 @@ fn refresh_attached_queue_for( // Mark as processed binding.pending_refresh = false; - (binding.server_id.clone(), binding.container_id.clone()) + let auto_play = binding.auto_play_on_refresh; + binding.auto_play_on_refresh = false; + ( + binding.server_id.clone(), + binding.container_id.clone(), + auto_play, + ) }; // Step 2: Fetch MediaServerInfo from registry @@ -1101,9 +1156,13 @@ fn refresh_attached_queue_for( } // Step 5: Intelligent refresh: try to keep current item if it's still in the new list - let current_item = runtime - .queue_snapshot(renderer_id) - .and_then(|q| q.first().cloned()); + // Get the full queue snapshot to access the item currently being played + let (full_queue, current_idx) = runtime + .queue_full_snapshot(renderer_id) + .unwrap_or((vec![], None)); + + // Get the item currently being played (at current_index), not the next one in queue + let current_item = current_idx.and_then(|idx| full_queue.get(idx).cloned()); let item_found_at = current_item.as_ref().and_then(|current| { new_items.iter().position(|new_item| { @@ -1116,39 +1175,66 @@ fn refresh_attached_queue_for( }) }); - let final_queue_len = runtime.with_queue_mut(renderer_id, |queue| { - queue.clear(); + let final_queue_len = runtime + .with_queue_mut(renderer_id, |queue| { + queue.clear(); - if let Some(idx) = item_found_at { - // Current item found: reconstruct queue starting from that item - for i in idx..new_items.len() { - queue.enqueue(new_items[i].clone()); + if let Some(idx) = item_found_at { + // Current item found: load the ENTIRE new playlist and position at that item + // This preserves items before the current track (as "already played") + for item in new_items.iter() { + queue.enqueue(item.clone()); + } + queue.set_current_index(Some(idx)); + info!( + renderer = renderer_id.0.as_str(), + server = server_id.0.as_str(), + container = container_id.as_str(), + total_items = new_items.len(), + current_index = idx, + upcoming = new_items.len().saturating_sub(idx + 1), + current_preserved = true, + "Refreshed queue from playlist container" + ); + new_items.len() + } else if let Some(ref current) = current_item { + // Current item NOT found: insert it at the beginning, then add new items + // This preserves the currently playing track and prevents it from being lost + queue.enqueue(current.clone()); + for item in new_items.iter() { + queue.enqueue(item.clone()); + } + queue.set_current_index(Some(0)); + info!( + renderer = renderer_id.0.as_str(), + server = server_id.0.as_str(), + container = container_id.as_str(), + total_items = new_items.len() + 1, + current_index = 0, + upcoming = new_items.len(), + current_preserved = true, + current_reinserted = true, + "Refreshed queue from playlist container (current item reinserted at start)" + ); + new_items.len() + 1 + } else { + // No current item: replace with full new list + for item in new_items.iter() { + queue.enqueue(item.clone()); + } + queue.set_current_index(None); + info!( + renderer = renderer_id.0.as_str(), + server = server_id.0.as_str(), + container = container_id.as_str(), + total_items = new_items.len(), + current_preserved = false, + "Refreshed queue from playlist container (no current item)" + ); + new_items.len() } - info!( - renderer = renderer_id.0.as_str(), - server = server_id.0.as_str(), - container = container_id.as_str(), - queue_len = new_items.len() - idx, - current_preserved = true, - "Refreshed queue from playlist container" - ); - new_items.len() - idx - } else { - // Current item not found: replace with full new list - for item in new_items.iter() { - queue.enqueue(item.clone()); - } - info!( - renderer = renderer_id.0.as_str(), - server = server_id.0.as_str(), - container = container_id.as_str(), - queue_len = new_items.len(), - current_preserved = false, - "Refreshed queue from playlist container (current item not found)" - ); - new_items.len() - } - }).unwrap_or(0); + }) + .unwrap_or(0); // Emit QueueUpdated event event_bus.broadcast(RendererEvent::QueueUpdated { @@ -1156,6 +1242,20 @@ fn refresh_attached_queue_for( queue_length: final_queue_len, }); + if auto_play { + if let Some(callback) = after_refresh.as_deref_mut() { + if let Err(err) = callback(renderer_id) { + warn!( + renderer = renderer_id.0.as_str(), + server = server_id.0.as_str(), + container = container_id.as_str(), + error = %err, + "Failed to auto-start playback after playlist refresh" + ); + } + } + } + Ok(()) } diff --git a/pmocontrol/src/media_server.rs b/pmocontrol/src/media_server.rs index ea1b338b..edb6c63d 100644 --- a/pmocontrol/src/media_server.rs +++ b/pmocontrol/src/media_server.rs @@ -344,10 +344,15 @@ fn map_didl_entries(xml: &str) -> Result> { let resources = item .resources .into_iter() - .map(|res| MediaResource { - uri: res.url, - protocol_info: res.protocol_info, - duration: res.duration, + .filter_map(|res| { + if res.url.trim().is_empty() { + return None; + } + Some(MediaResource { + uri: res.url, + protocol_info: res.protocol_info, + duration: res.duration, + }) }) .collect(); diff --git a/pmocontrol/src/playback_queue.rs b/pmocontrol/src/playback_queue.rs index 0aecb890..16049017 100644 --- a/pmocontrol/src/playback_queue.rs +++ b/pmocontrol/src/playback_queue.rs @@ -129,4 +129,16 @@ impl PlaybackQueue { pub fn full_snapshot(&self) -> (Vec, Option) { (self.items.clone(), self.current_index) } + + pub fn set_current_index(&mut self, index: Option) { + if let Some(idx) = index { + if idx < self.items.len() { + self.current_index = Some(idx); + } else { + self.current_index = None; + } + } else { + self.current_index = None; + } + } } diff --git a/pmocontrol/src/pmoserver_ext.rs b/pmocontrol/src/pmoserver_ext.rs index 742f11ec..78ccb0f0 100644 --- a/pmocontrol/src/pmoserver_ext.rs +++ b/pmocontrol/src/pmoserver_ext.rs @@ -11,7 +11,7 @@ use crate::media_server::{MediaBrowser, MusicServer, ServerId}; use crate::model::{RendererId, RendererProtocol}; #[cfg(feature = "pmoserver")] use crate::openapi::{ - AttachedPlaylistInfo, AttachPlaylistRequest, BrowseResponse, ContainerEntry, ErrorResponse, + AttachPlaylistRequest, AttachedPlaylistInfo, BrowseResponse, ContainerEntry, ErrorResponse, MediaServerSummary, QueueItem, QueueSnapshot, RendererState, RendererSummary, SuccessResponse, VolumeSetRequest, }; @@ -22,20 +22,49 @@ use crate::{PlaybackPosition, PlaybackStatus, TransportControl, VolumeControl}; use async_trait::async_trait; #[cfg(feature = "pmoserver")] use axum::{ + Json, Router, extract::{Path, State}, http::StatusCode, routing::{get, post}, - Json, Router, }; #[cfg(feature = "pmoserver")] use std::sync::Arc; #[cfg(feature = "pmoserver")] use std::time::Duration; #[cfg(feature = "pmoserver")] +use tokio::time; +#[cfg(feature = "pmoserver")] use tracing::{debug, warn}; #[cfg(feature = "pmoserver")] use utoipa::OpenApi; +#[cfg(feature = "pmoserver")] +const BROWSE_PAGE_SIZE: u32 = 100; +#[cfg(feature = "pmoserver")] +const MEDIA_SERVER_SOAP_TIMEOUT: Duration = Duration::from_secs(15); +#[cfg(feature = "pmoserver")] +const BROWSE_REQUEST_TIMEOUT: Duration = Duration::from_secs(20); + +// Timeouts for simple commands (play/pause/stop) +#[cfg(feature = "pmoserver")] +const TRANSPORT_COMMAND_TIMEOUT: Duration = Duration::from_secs(5); + +// Timeouts for volume/mute commands (faster than transport) +#[cfg(feature = "pmoserver")] +const VOLUME_COMMAND_TIMEOUT: Duration = Duration::from_secs(3); + +// Timeout for queue operations +#[cfg(feature = "pmoserver")] +const QUEUE_COMMAND_TIMEOUT: Duration = Duration::from_secs(10); + +// Timeout for attach playlist (includes browse + cache + queue update) +#[cfg(feature = "pmoserver")] +const ATTACH_PLAYLIST_TIMEOUT: Duration = Duration::from_secs(60); + +// Timeout for renderer state queries (multiple SOAP calls) +#[cfg(feature = "pmoserver")] +const STATE_QUERY_TIMEOUT: Duration = Duration::from_secs(8); + /// État partagé pour l'API ControlPoint #[cfg(feature = "pmoserver")] #[derive(Clone)] @@ -64,9 +93,7 @@ impl ControlPointState { ), tag = "control" )] -async fn list_renderers( - State(state): State, -) -> Json> { +async fn list_renderers(State(state): State) -> Json> { let renderers = state.control_point.list_music_renderers(); let summaries: Vec = renderers @@ -119,30 +146,64 @@ async fn get_renderer_state( })?; let info = renderer.info(); + let renderer_clone = renderer.clone(); - // État de transport - let transport_state = renderer - .playback_state() - .ok() - .map(state_to_string) - .unwrap_or_else(|| "UNKNOWN".to_string()); + // Spawn blocking task for all SOAP calls to avoid blocking Tokio runtime + let state_task = tokio::task::spawn_blocking(move || { + // État de transport + let transport_state = renderer_clone + .playback_state() + .ok() + .map(state_to_string) + .unwrap_or_else(|| "UNKNOWN".to_string()); - // Position et durée - let (position_ms, duration_ms) = renderer - .playback_position() - .ok() - .and_then(|pos| { - let position = parse_hms_to_ms(pos.rel_time.as_deref()); - let duration = parse_hms_to_ms(pos.track_duration.as_deref()); - Some((position, duration)) - }) - .unwrap_or((None, None)); + // Position et durée + let (position_ms, duration_ms) = renderer_clone + .playback_position() + .ok() + .and_then(|pos| { + let position = parse_hms_to_ms(pos.rel_time.as_deref()); + let duration = parse_hms_to_ms(pos.track_duration.as_deref()); + Some((position, duration)) + }) + .unwrap_or((None, None)); - // Volume et mute - let volume = renderer.volume().ok().and_then(|v| u8::try_from(v).ok()); - let mute = renderer.mute().ok(); + // Volume et mute + let volume = renderer_clone.volume().ok().and_then(|v| u8::try_from(v).ok()); + let mute = renderer_clone.mute().ok(); - // Queue + (transport_state, position_ms, duration_ms, volume, mute) + }); + + let (transport_state, position_ms, duration_ms, volume, mute) = + time::timeout(STATE_QUERY_TIMEOUT, state_task) + .await + .map_err(|_| { + warn!( + "State query for renderer {} exceeded {:?}", + renderer_id, STATE_QUERY_TIMEOUT + ); + ( + StatusCode::GATEWAY_TIMEOUT, + Json(ErrorResponse { + error: format!( + "State query timed out after {}s", + STATE_QUERY_TIMEOUT.as_secs() + ), + }), + ) + })? + .map_err(|e| { + warn!("Task join error during state query: {}", e); + ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(ErrorResponse { + error: format!("Internal task error: {}", e), + }), + ) + })?; + + // Queue (non-blocking, local data) let queue_len = state .control_point .get_queue_snapshot(&rid) @@ -150,15 +211,17 @@ async fn get_renderer_state( .map(|q| q.len()) .unwrap_or(0); - // Playlist binding + // Playlist binding (non-blocking, local data) let attached_playlist = state .control_point .current_queue_playlist_binding(&rid) - .map(|(server_id, container_id, has_seen_update)| AttachedPlaylistInfo { - server_id: server_id.0, - container_id, - has_seen_update, - }); + .map( + |(server_id, container_id, has_seen_update)| AttachedPlaylistInfo { + server_id: server_id.0, + container_id, + has_seen_update, + }, + ); Ok(Json(RendererState { id: info.id.0.clone(), @@ -197,17 +260,18 @@ async fn get_renderer_queue( ) -> Result, (StatusCode, Json)> { let rid = RendererId(renderer_id.clone()); - let (items, current_index) = state - .control_point - .get_full_queue_snapshot(&rid) - .map_err(|e| { - ( - StatusCode::NOT_FOUND, - Json(ErrorResponse { - error: format!("Failed to get queue: {}", e), - }), - ) - })?; + let (items, current_index) = + state + .control_point + .get_full_queue_snapshot(&rid) + .map_err(|e| { + ( + StatusCode::NOT_FOUND, + Json(ErrorResponse { + error: format!("Failed to get queue: {}", e), + }), + ) + })?; let queue_items: Vec = items .into_iter() @@ -253,11 +317,13 @@ async fn get_renderer_binding( let binding = state .control_point .current_queue_playlist_binding(&rid) - .map(|(server_id, container_id, has_seen_update)| AttachedPlaylistInfo { - server_id: server_id.0, - container_id, - has_seen_update, - }); + .map( + |(server_id, container_id, has_seen_update)| AttachedPlaylistInfo { + server_id: server_id.0, + container_id, + has_seen_update, + }, + ); Ok(Json(binding)) } @@ -299,15 +365,44 @@ async fn play_renderer( ) })?; - renderer.play().map_err(|e| { - warn!("Failed to play renderer {}: {}", renderer_id, e); - ( - StatusCode::INTERNAL_SERVER_ERROR, - Json(ErrorResponse { - error: format!("Failed to play: {}", e), - }), - ) - })?; + let renderer_clone = renderer.clone(); + let play_task = tokio::task::spawn_blocking(move || renderer_clone.play()); + + time::timeout(TRANSPORT_COMMAND_TIMEOUT, play_task) + .await + .map_err(|_| { + warn!( + "Play command for renderer {} exceeded {:?}", + renderer_id, TRANSPORT_COMMAND_TIMEOUT + ); + ( + StatusCode::GATEWAY_TIMEOUT, + Json(ErrorResponse { + error: format!( + "Play command timed out after {}s", + TRANSPORT_COMMAND_TIMEOUT.as_secs() + ), + }), + ) + })? + .map_err(|e| { + warn!("Task join error during play: {}", e); + ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(ErrorResponse { + error: format!("Internal task error: {}", e), + }), + ) + })? + .map_err(|e| { + warn!("Failed to play renderer {}: {}", renderer_id, e); + ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(ErrorResponse { + error: format!("Failed to play: {}", e), + }), + ) + })?; Ok(Json(SuccessResponse { message: "Playback started".to_string(), @@ -347,15 +442,44 @@ async fn pause_renderer( ) })?; - renderer.pause().map_err(|e| { - warn!("Failed to pause renderer {}: {}", renderer_id, e); - ( - StatusCode::INTERNAL_SERVER_ERROR, - Json(ErrorResponse { - error: format!("Failed to pause: {}", e), - }), - ) - })?; + let renderer_clone = renderer.clone(); + let pause_task = tokio::task::spawn_blocking(move || renderer_clone.pause()); + + time::timeout(TRANSPORT_COMMAND_TIMEOUT, pause_task) + .await + .map_err(|_| { + warn!( + "Pause command for renderer {} exceeded {:?}", + renderer_id, TRANSPORT_COMMAND_TIMEOUT + ); + ( + StatusCode::GATEWAY_TIMEOUT, + Json(ErrorResponse { + error: format!( + "Pause command timed out after {}s", + TRANSPORT_COMMAND_TIMEOUT.as_secs() + ), + }), + ) + })? + .map_err(|e| { + warn!("Task join error during pause: {}", e); + ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(ErrorResponse { + error: format!("Internal task error: {}", e), + }), + ) + })? + .map_err(|e| { + warn!("Failed to pause renderer {}: {}", renderer_id, e); + ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(ErrorResponse { + error: format!("Failed to pause: {}", e), + }), + ) + })?; Ok(Json(SuccessResponse { message: "Playback paused".to_string(), @@ -395,15 +519,44 @@ async fn stop_renderer( ) })?; - renderer.stop().map_err(|e| { - warn!("Failed to stop renderer {}: {}", renderer_id, e); - ( - StatusCode::INTERNAL_SERVER_ERROR, - Json(ErrorResponse { - error: format!("Failed to stop: {}", e), - }), - ) - })?; + let renderer_clone = renderer.clone(); + let stop_task = tokio::task::spawn_blocking(move || renderer_clone.stop()); + + time::timeout(TRANSPORT_COMMAND_TIMEOUT, stop_task) + .await + .map_err(|_| { + warn!( + "Stop command for renderer {} exceeded {:?}", + renderer_id, TRANSPORT_COMMAND_TIMEOUT + ); + ( + StatusCode::GATEWAY_TIMEOUT, + Json(ErrorResponse { + error: format!( + "Stop command timed out after {}s", + TRANSPORT_COMMAND_TIMEOUT.as_secs() + ), + }), + ) + })? + .map_err(|e| { + warn!("Task join error during stop: {}", e); + ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(ErrorResponse { + error: format!("Internal task error: {}", e), + }), + ) + })? + .map_err(|e| { + warn!("Failed to stop renderer {}: {}", renderer_id, e); + ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(ErrorResponse { + error: format!("Failed to stop: {}", e), + }), + ) + })?; Ok(Json(SuccessResponse { message: "Playback stopped".to_string(), @@ -430,12 +583,44 @@ async fn next_renderer( Path(renderer_id): Path, ) -> Result, (StatusCode, Json)> { let rid = RendererId(renderer_id.clone()); + let control_point = state.control_point.clone(); + let rid_clone = rid.clone(); - state - .control_point - .play_next_from_queue(&rid) + let next_task = tokio::task::spawn_blocking(move || { + control_point.play_next_from_queue(&rid_clone) + }); + + time::timeout(QUEUE_COMMAND_TIMEOUT, next_task) + .await + .map_err(|_| { + warn!( + "Next command for renderer {} exceeded {:?}", + renderer_id, QUEUE_COMMAND_TIMEOUT + ); + ( + StatusCode::GATEWAY_TIMEOUT, + Json(ErrorResponse { + error: format!( + "Next command timed out after {}s", + QUEUE_COMMAND_TIMEOUT.as_secs() + ), + }), + ) + })? .map_err(|e| { - warn!("Failed to advance queue for renderer {}: {}", renderer_id, e); + warn!("Task join error during next: {}", e); + ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(ErrorResponse { + error: format!("Internal task error: {}", e), + }), + ) + })? + .map_err(|e| { + warn!( + "Failed to advance queue for renderer {}: {}", + renderer_id, e + ); ( StatusCode::INTERNAL_SERVER_ERROR, Json(ErrorResponse { @@ -489,18 +674,50 @@ async fn set_renderer_volume( ) })?; - renderer.set_volume(req.volume as u16).map_err(|e| { - warn!("Failed to set volume for renderer {}: {}", renderer_id, e); - ( - StatusCode::INTERNAL_SERVER_ERROR, - Json(ErrorResponse { - error: format!("Failed to set volume: {}", e), - }), - ) - })?; + let renderer_clone = renderer.clone(); + let volume = req.volume; + let volume_task = tokio::task::spawn_blocking(move || { + renderer_clone.set_volume(volume as u16) + }); + + time::timeout(VOLUME_COMMAND_TIMEOUT, volume_task) + .await + .map_err(|_| { + warn!( + "Set volume command for renderer {} exceeded {:?}", + renderer_id, VOLUME_COMMAND_TIMEOUT + ); + ( + StatusCode::GATEWAY_TIMEOUT, + Json(ErrorResponse { + error: format!( + "Set volume command timed out after {}s", + VOLUME_COMMAND_TIMEOUT.as_secs() + ), + }), + ) + })? + .map_err(|e| { + warn!("Task join error during set volume: {}", e); + ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(ErrorResponse { + error: format!("Internal task error: {}", e), + }), + ) + })? + .map_err(|e| { + warn!("Failed to set volume for renderer {}: {}", renderer_id, e); + ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(ErrorResponse { + error: format!("Failed to set volume: {}", e), + }), + ) + })?; Ok(Json(SuccessResponse { - message: format!("Volume set to {}", req.volume), + message: format!("Volume set to {}", volume), })) } @@ -537,25 +754,52 @@ async fn volume_up_renderer( ) })?; - let current = renderer.volume().map_err(|e| { - ( - StatusCode::INTERNAL_SERVER_ERROR, - Json(ErrorResponse { - error: format!("Failed to get current volume: {}", e), - }), - ) - })?; + let renderer_clone = renderer.clone(); + let volume_task = tokio::task::spawn_blocking(move || { + let current = renderer_clone.volume()?; + let new_volume = (current + 5).min(100); + renderer_clone.set_volume(new_volume)?; + Ok::(new_volume) + }); - let new_volume = (current + 5).min(100); - renderer.set_volume(new_volume).map_err(|e| { - warn!("Failed to increase volume for renderer {}: {}", renderer_id, e); - ( - StatusCode::INTERNAL_SERVER_ERROR, - Json(ErrorResponse { - error: format!("Failed to increase volume: {}", e), - }), - ) - })?; + let new_volume = time::timeout(VOLUME_COMMAND_TIMEOUT, volume_task) + .await + .map_err(|_| { + warn!( + "Volume up command for renderer {} exceeded {:?}", + renderer_id, VOLUME_COMMAND_TIMEOUT + ); + ( + StatusCode::GATEWAY_TIMEOUT, + Json(ErrorResponse { + error: format!( + "Volume up command timed out after {}s", + VOLUME_COMMAND_TIMEOUT.as_secs() + ), + }), + ) + })? + .map_err(|e| { + warn!("Task join error during volume up: {}", e); + ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(ErrorResponse { + error: format!("Internal task error: {}", e), + }), + ) + })? + .map_err(|e| { + warn!( + "Failed to increase volume for renderer {}: {}", + renderer_id, e + ); + ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(ErrorResponse { + error: format!("Failed to increase volume: {}", e), + }), + ) + })?; Ok(Json(SuccessResponse { message: format!("Volume increased to {}", new_volume), @@ -595,25 +839,52 @@ async fn volume_down_renderer( ) })?; - let current = renderer.volume().map_err(|e| { - ( - StatusCode::INTERNAL_SERVER_ERROR, - Json(ErrorResponse { - error: format!("Failed to get current volume: {}", e), - }), - ) - })?; + let renderer_clone = renderer.clone(); + let volume_task = tokio::task::spawn_blocking(move || { + let current = renderer_clone.volume()?; + let new_volume = current.saturating_sub(5); + renderer_clone.set_volume(new_volume)?; + Ok::(new_volume) + }); - let new_volume = current.saturating_sub(5); - renderer.set_volume(new_volume).map_err(|e| { - warn!("Failed to decrease volume for renderer {}: {}", renderer_id, e); - ( - StatusCode::INTERNAL_SERVER_ERROR, - Json(ErrorResponse { - error: format!("Failed to decrease volume: {}", e), - }), - ) - })?; + let new_volume = time::timeout(VOLUME_COMMAND_TIMEOUT, volume_task) + .await + .map_err(|_| { + warn!( + "Volume down command for renderer {} exceeded {:?}", + renderer_id, VOLUME_COMMAND_TIMEOUT + ); + ( + StatusCode::GATEWAY_TIMEOUT, + Json(ErrorResponse { + error: format!( + "Volume down command timed out after {}s", + VOLUME_COMMAND_TIMEOUT.as_secs() + ), + }), + ) + })? + .map_err(|e| { + warn!("Task join error during volume down: {}", e); + ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(ErrorResponse { + error: format!("Internal task error: {}", e), + }), + ) + })? + .map_err(|e| { + warn!( + "Failed to decrease volume for renderer {}: {}", + renderer_id, e + ); + ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(ErrorResponse { + error: format!("Failed to decrease volume: {}", e), + }), + ) + })?; Ok(Json(SuccessResponse { message: format!("Volume decreased to {}", new_volume), @@ -653,25 +924,49 @@ async fn toggle_mute_renderer( ) })?; - let current_mute = renderer.mute().map_err(|e| { - ( - StatusCode::INTERNAL_SERVER_ERROR, - Json(ErrorResponse { - error: format!("Failed to get current mute state: {}", e), - }), - ) - })?; + let renderer_clone = renderer.clone(); + let mute_task = tokio::task::spawn_blocking(move || { + let current_mute = renderer_clone.mute()?; + let new_mute = !current_mute; + renderer_clone.set_mute(new_mute)?; + Ok::(new_mute) + }); - let new_mute = !current_mute; - renderer.set_mute(new_mute).map_err(|e| { - warn!("Failed to toggle mute for renderer {}: {}", renderer_id, e); - ( - StatusCode::INTERNAL_SERVER_ERROR, - Json(ErrorResponse { - error: format!("Failed to toggle mute: {}", e), - }), - ) - })?; + let new_mute = time::timeout(VOLUME_COMMAND_TIMEOUT, mute_task) + .await + .map_err(|_| { + warn!( + "Toggle mute command for renderer {} exceeded {:?}", + renderer_id, VOLUME_COMMAND_TIMEOUT + ); + ( + StatusCode::GATEWAY_TIMEOUT, + Json(ErrorResponse { + error: format!( + "Toggle mute command timed out after {}s", + VOLUME_COMMAND_TIMEOUT.as_secs() + ), + }), + ) + })? + .map_err(|e| { + warn!("Task join error during toggle mute: {}", e); + ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(ErrorResponse { + error: format!("Internal task error: {}", e), + }), + ) + })? + .map_err(|e| { + warn!("Failed to toggle mute for renderer {}: {}", renderer_id, e); + ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(ErrorResponse { + error: format!("Failed to toggle mute: {}", e), + }), + ) + })?; Ok(Json(SuccessResponse { message: format!("Mute {}", if new_mute { "enabled" } else { "disabled" }), @@ -704,10 +999,40 @@ async fn attach_playlist_binding( ) -> Result, (StatusCode, Json)> { let rid = RendererId(renderer_id.clone()); let sid = ServerId(req.server_id.clone()); + let container_id = req.container_id.clone(); + let control_point = Arc::clone(&state.control_point); - state - .control_point - .attach_queue_to_playlist(&rid, sid, req.container_id.clone()); + // Spawn blocking task and wait for completion with timeout + let attach_task = tokio::task::spawn_blocking(move || { + control_point.attach_queue_to_playlist(&rid, sid, container_id); + }); + + time::timeout(ATTACH_PLAYLIST_TIMEOUT, attach_task) + .await + .map_err(|_| { + warn!( + "Attach playlist for renderer {} exceeded {:?}", + renderer_id, ATTACH_PLAYLIST_TIMEOUT + ); + ( + StatusCode::GATEWAY_TIMEOUT, + Json(ErrorResponse { + error: format!( + "Attach playlist timed out after {}s", + ATTACH_PLAYLIST_TIMEOUT.as_secs() + ), + }), + ) + })? + .map_err(|e| { + warn!("Task join error during attach playlist: {}", e); + ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(ErrorResponse { + error: format!("Internal task error: {}", e), + }), + ) + })?; debug!( renderer = renderer_id.as_str(), @@ -717,10 +1042,7 @@ async fn attach_playlist_binding( ); Ok(Json(SuccessResponse { - message: format!( - "Playlist {} attached to renderer", - req.container_id - ), + message: format!("Playlist {} attached to renderer", req.container_id), })) } @@ -769,9 +1091,7 @@ async fn detach_playlist_binding( ), tag = "control" )] -async fn list_servers( - State(state): State, -) -> Json> { +async fn list_servers(State(state): State) -> Json> { let servers = state.control_point.list_media_servers(); let summaries: Vec = servers @@ -809,17 +1129,14 @@ async fn browse_container( ) -> Result, (StatusCode, Json)> { let sid = ServerId(server_id.clone()); - let server_info = state - .control_point - .media_server(&sid) - .ok_or_else(|| { - ( - StatusCode::NOT_FOUND, - Json(ErrorResponse { - error: format!("Server {} not found", server_id), - }), - ) - })?; + let server_info = state.control_point.media_server(&sid).ok_or_else(|| { + ( + StatusCode::NOT_FOUND, + Json(ErrorResponse { + error: format!("Server {} not found", server_id), + }), + ) + })?; if !server_info.online { return Err(( @@ -839,8 +1156,8 @@ async fn browse_container( )); } - let music_server = MusicServer::from_info(&server_info, Duration::from_secs(10)).map_err( - |e| { + let music_server = + MusicServer::from_info(&server_info, MEDIA_SERVER_SOAP_TIMEOUT).map_err(|e| { warn!("Failed to create MusicServer for {}: {}", server_id, e); ( StatusCode::INTERNAL_SERVER_ERROR, @@ -848,11 +1165,40 @@ async fn browse_container( error: format!("Failed to initialize server: {}", e), }), ) - }, - )?; + })?; - let entries = music_server - .browse_children(&container_id, 0, 100) + // Use spawn_blocking to avoid blocking the async runtime with synchronous SOAP calls + let container_id_clone = container_id.clone(); + let browse_task = tokio::task::spawn_blocking(move || { + music_server.browse_children(&container_id_clone, 0, BROWSE_PAGE_SIZE) + }); + + let entries = time::timeout(BROWSE_REQUEST_TIMEOUT, browse_task) + .await + .map_err(|_| { + warn!( + "Browse request for container {} on server {} exceeded {:?}", + container_id, server_id, BROWSE_REQUEST_TIMEOUT + ); + ( + StatusCode::GATEWAY_TIMEOUT, + Json(ErrorResponse { + error: format!( + "Browse request timed out after {}s", + BROWSE_REQUEST_TIMEOUT.as_secs() + ), + }), + ) + })? + .map_err(|e| { + warn!("Task join error during browse: {}", e); + ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(ErrorResponse { + error: format!("Internal task error: {}", e), + }), + ) + })? .map_err(|e| { warn!( "Failed to browse container {} on server {}: {}", @@ -939,23 +1285,47 @@ pub fn create_api_router(state: ControlPointState, control_point: Arc std::io::Result>; + async fn register_control_point( + &mut self, + timeout_secs: u64, + ) -> std::io::Result>; /// Initialise l'API Control Point (bas niveau) /// @@ -1041,7 +1414,10 @@ pub trait ControlPointExt { #[cfg(feature = "pmoserver")] #[async_trait] impl ControlPointExt for pmoserver::Server { - async fn register_control_point(&mut self, timeout_secs: u64) -> std::io::Result> { + async fn register_control_point( + &mut self, + timeout_secs: u64, + ) -> std::io::Result> { use tracing::info; info!("🎛️ Initializing Control Point..."); diff --git a/pmocontrol/src/sse.rs b/pmocontrol/src/sse.rs index 67b7394f..62991225 100644 --- a/pmocontrol/src/sse.rs +++ b/pmocontrol/src/sse.rs @@ -10,20 +10,20 @@ //! - GET /api/control/events/servers - Événements serveurs uniquement //! - GET /api/control/events - Tous les événements (agrégés) +#[cfg(feature = "pmoserver")] +use crate::PlaybackState; #[cfg(feature = "pmoserver")] use crate::control_point::ControlPoint; #[cfg(feature = "pmoserver")] use crate::model::{MediaServerEvent, RendererEvent}; #[cfg(feature = "pmoserver")] -use crate::PlaybackState; -#[cfg(feature = "pmoserver")] use async_stream::stream; #[cfg(feature = "pmoserver")] use axum::{ - extract::State, - response::sse::{Event, KeepAlive, Sse}, - response::IntoResponse, Router, + extract::State, + response::IntoResponse, + response::sse::{Event, KeepAlive, Sse}, }; #[cfg(feature = "pmoserver")] use serde::Serialize; @@ -276,9 +276,7 @@ pub async fn media_server_events_sse( ), tag = "control" )] -pub async fn all_events_sse( - State(control_point): State>, -) -> impl IntoResponse { +pub async fn all_events_sse(State(control_point): State>) -> impl IntoResponse { // Convert crossbeam channels to tokio channels for async compatibility let (renderer_tx, mut renderer_rx_tokio) = tokio::sync::mpsc::unbounded_channel(); let (server_tx, mut server_rx_tokio) = tokio::sync::mpsc::unbounded_channel(); diff --git a/pmodidl/src/lib.rs b/pmodidl/src/lib.rs index eb901736..deeb8281 100644 --- a/pmodidl/src/lib.rs +++ b/pmodidl/src/lib.rs @@ -118,7 +118,7 @@ pub struct Container { #[serde(rename = "@id")] pub id: String, - #[serde(rename = "@parentID")] + #[serde(rename = "@parentID", default)] pub parent_id: String, #[serde(rename = "@restricted", skip_serializing_if = "Option::is_none")] @@ -133,7 +133,7 @@ pub struct Container { #[serde(rename = "dc:title", alias = "title")] pub title: String, - #[serde(rename = "upnp:class", alias = "class")] + #[serde(rename = "upnp:class", alias = "class", default)] pub class: String, #[serde(rename = "container", default)] @@ -149,7 +149,7 @@ pub struct Item { #[serde(rename = "@id")] pub id: String, - #[serde(rename = "@parentID")] + #[serde(rename = "@parentID", default)] pub parent_id: String, #[serde(rename = "@restricted", skip_serializing_if = "Option::is_none")] @@ -165,7 +165,7 @@ pub struct Item { )] pub creator: Option, - #[serde(rename = "upnp:class", alias = "class")] + #[serde(rename = "upnp:class", alias = "class", default)] pub class: String, #[serde( @@ -238,7 +238,7 @@ pub struct Resource { #[serde(rename = "@duration", skip_serializing_if = "Option::is_none")] pub duration: Option, - #[serde(rename = "$text")] + #[serde(rename = "$text", default)] pub url: String, }