diff --git a/.DS_Store b/.DS_Store index b4d7ee26..8c67cf49 100644 Binary files a/.DS_Store and b/.DS_Store differ diff --git a/Blackboard/ToThinkAbout/plant_webrenderer_server_side_streaming.md b/Blackboard/ToThinkAbout/plant_webrenderer_server_side_streaming.md new file mode 100644 index 00000000..6f5709a9 --- /dev/null +++ b/Blackboard/ToThinkAbout/plant_webrenderer_server_side_streaming.md @@ -0,0 +1,1193 @@ +# Plan d'implémentation détaillé : WebRenderer streaming audio côté serveur + +## Synthèse de l'exploration du code + +### Architecture actuelle + +L'architecture actuelle repose sur : + +1. `pmowebrenderer/src/websocket.rs` - Handler WebSocket : crée ou reconnecte un device UPnP par navigateur, envoie des commandes JSON au navigateur +2. `pmowebrenderer/src/handlers.rs` - Action handlers UPnP qui forwardent vers le WebSocket via `SharedSender` +3. `pmowebrenderer/src/session.rs` - `SessionManager` : HashMap token → session, persistance UDN/sender entre reconnexions +4. `pmowebrenderer/src/state.rs` - `RendererState` + `SharedSender` (remplaçable à chaque reconnexion WS) +5. `pmowebrenderer/src/renderer.rs` - `WebRendererFactory` : construit les services UPnP (AVTransport, RenderingControl, ConnectionManager) +6. `pmowebrenderer/src/config.rs` - Trait `WebRendererExt` : enregistre la route WS dans pmoserver +7. Frontend : `useWebRenderer.ts` - composable Vue.js gérant WebSocket + `GaplessEngine` (deux `HTMLAudioElement` en ping-pong) + +### Infrastructure réutilisable confirmée + +- `StreamingFlacSink` + `StreamHandle` dans `pmoaudio-ext/src/sinks/streaming_flac_sink.rs` : broadcast multi-clients, gapless, backpressure, header FLAC caché pour late-joiners +- Pattern HTTP streaming dans `pmomediaserver/src/paradise_streaming.rs` : `Body::from_stream(ReaderStream::new(stream))` avec headers corrects +- `PlaylistSource` dans `pmoaudio-ext/src/sources/playlist_source.rs` : modèle pour `source_loader.rs`, incluant décodage FLAC via `pmoflac::decode_audio_stream`, gestion cache progressif +- `pmoserver::Server` : API `add_handler_with_state`, `add_post_handler_with_state`, `add_any_handler_with_state`, `add_router` + +### Dépendances existantes de pmowebrenderer + +Le crate actuel ne dépend pas de `pmoaudio-ext`, `pmoaudio`, `pmoflac` ou `pmoaudiocache`. Il faudra les ajouter. + +--- + +## Ordre d'implémentation + +Les étapes sont organisées de façon à avoir un système compilable et testable à chaque jalon, en allant des fondations vers l'intégration. + +--- + +## Etape 1 : Préparer le `Cargo.toml` de pmowebrenderer + +**Fichier concerné :** `pmowebrenderer/Cargo.toml` + +**Modifications :** +```toml +[dependencies] +# Existants conservés +pmoupnp = { path = "../pmoupnp" } +pmomediarenderer = { path = "../pmomediarenderer" } +pmoserver = { path = "../pmoserver", optional = true } +pmocontrol = { path = "../pmocontrol", optional = true } +pmoconfig = { path = "../pmoconfig" } + +# Nouveaux : pipeline audio +pmoaudio-ext = { path = "../pmoaudio-ext", features = ["http-stream", "playlist"], optional = true } +pmoaudio = { path = "../pmoaudio", optional = true } +pmoflac = { path = "../pmoflac", optional = true } +pmoaudiocache = { path = "../pmoaudiocache", optional = true } +pmometadata = { path = "../pmometadata", optional = true } + +# Async runtime +tokio = { workspace = true, features = ["full"] } +async-trait = { workspace = true } + +# HTTP streaming (on retire "ws" de axum si WebSocket supprimé en phase finale) +axum = { workspace = true, features = ["ws"] } # garder "ws" pendant la migration +axum-extra = { version = "0.9", features = ["typed-header"] } +tower-http = { version = "0.6", features = ["fs", "trace"] } +futures = "0.3" +bytes = "1.0" +tokio-util = { version = "0.7", features = ["io"] } + +# Serialization +serde = { workspace = true } +serde_json = { workspace = true } + +# Utilities +uuid = { workspace = true, features = ["v4", "serde"] } +parking_lot = "0.12" +thiserror = { workspace = true } +tracing = { workspace = true } + +pmodidl = { path = "../pmodidl" } +pmoutils = { path = "../pmoutils" } + +[features] +default = [] +pmoserver = ["dep:pmoserver", "dep:pmocontrol"] +audio-pipeline = ["dep:pmoaudio-ext", "dep:pmoaudio", "dep:pmoflac", "dep:pmoaudiocache", "dep:pmometadata"] +``` + +**Note :** Ajouter `pmoaudio-ext` dans le workspace `Cargo.toml` (il n'y est pas encore). + +**Piège :** `pmoaudio-ext` n'est pas dans `[workspace.members]` du `Cargo.toml` racine. Il faut l'y ajouter. Vérifier aussi que `pmoaudio-ext` avec feature `playlist` n'introduit pas de dépendance circulaire via `pmoplaylist`. + +--- + +## Etape 2 : Créer `pipeline.rs` — canal de contrôle et état par instance + +**Fichier :** `pmowebrenderer/src/pipeline.rs` (nouveau) + +Ce module définit les types partagés entre tous les handlers sans logique de pipeline encore. + +```rust +use tokio::sync::mpsc; + +/// Commandes envoyées au pipeline audio de l'instance +#[derive(Debug)] +pub enum PipelineControl { + LoadUri(String), + LoadNextUri(String), + Play, + Pause, + Stop, + Seek(f64), // secondes + SetVolume(u16), // 0-100 + SetMute(bool), +} + +/// Handle vers le pipeline audio d'une instance WebRenderer serveur +#[derive(Clone)] +pub struct PipelineHandle { + /// Canal de contrôle vers la task pipeline + pub control_tx: mpsc::Sender, + /// Token d'annulation pour stopper le pipeline + pub stop_token: tokio_util::sync::CancellationToken, +} + +impl PipelineHandle { + pub async fn send(&self, cmd: PipelineControl) { + let _ = self.control_tx.send(cmd).await; + } +} +``` + +**Points d'attention :** +- `PipelineControl` doit être `Send` (pas de `Arc>` dans les variantes). +- Le `CancellationToken` de `tokio_util` est déjà utilisé dans `pmoaudio-ext`, importer depuis là. + +--- + +## Etape 3 : Modifier `state.rs` — ajouter stream handle et pipeline handle + +**Fichier :** `pmowebrenderer/src/state.rs` (modifier) + +Remplacer `SharedSender` (qui envoie vers le WS) par un `SharedStreamHandle` et un `PipelineHandle`. + +```rust +use parking_lot::RwLock; +use std::sync::Arc; + +#[cfg(feature = "audio-pipeline")] +use pmoaudio_ext::sinks::streaming_flac_sink::StreamHandle; + +use crate::messages::PlaybackState; +#[cfg(feature = "audio-pipeline")] +use crate::pipeline::PipelineHandle; + +/// État temps-réel du renderer (partagé backend ↔ pipeline) +#[derive(Debug, Clone)] +pub struct RendererState { + pub playback_state: PlaybackState, + pub current_uri: Option, + pub current_metadata: Option, + pub next_uri: Option, + pub next_metadata: Option, + pub position: Option, // mis à jour par la task de suivi position + pub duration: Option, + pub volume: u16, + pub mute: bool, +} + +impl Default for RendererState { /* identique à l'existant */ } + +pub type SharedState = Arc>; + +/// Handle vers le flux FLAC HTTP d'une instance (remplace SharedSender) +#[cfg(feature = "audio-pipeline")] +#[derive(Clone)] +pub struct SharedStreamHandle(Arc>>); + +#[cfg(feature = "audio-pipeline")] +impl SharedStreamHandle { + pub fn new(handle: StreamHandle) -> Self { + Self(Arc::new(RwLock::new(Some(handle)))) + } + + pub fn get(&self) -> Option { + self.0.read().clone() + } + + pub fn clear(&self) { + *self.0.write() = None; + } +} + +// Pour la compatibilité avec les handlers qui utilisent encore SharedSender +// pendant la phase de migration, on conserve SharedSender dans le module +// mais on la rend conditionnelle au feature "pmoserver" seul (sans audio-pipeline). +``` + +**Points d'attention :** +- Conserver `SharedSender` derrière une feature flag pendant la migration. Les handlers UPnP existants (`handlers.rs`) l'utilisent encore. +- La `StreamHandle` de `pmoaudio-ext` est `Clone` — on peut la partager sans `Arc>` en réalité, mais un wrapper permet de la remplacer à la reconnexion. + +--- + +## Etape 4 : Créer `source_loader.rs` — ouverture d'URI arbitraires + +**Fichier :** `pmowebrenderer/src/source_loader.rs` (nouveau) + +Ce module sait ouvrir une URI quelconque (URL HTTP, chemin fichier local, chemin Samba) et la transformer en `AsyncRead` de bytes PCM via `pmoflac::decode_audio_stream`. + +```rust +use pmoflac::decode_audio_stream; +use std::path::Path; +use tokio::io::AsyncRead; + +pub enum SourceKind { + LocalFile(std::path::PathBuf), + HttpUrl(String), +} + +pub fn classify_uri(uri: &str) -> SourceKind { + if uri.starts_with("http://") || uri.starts_with("https://") { + SourceKind::HttpUrl(uri.to_string()) + } else { + // Chemin fichier (absolu ou Samba monté) + SourceKind::LocalFile(std::path::PathBuf::from(uri)) + } +} + +/// Ouvre une source audio quelconque et retourne le stream PCM décodé. +/// Retourne (stream, sample_rate, bits_per_sample, channels) +pub async fn open_uri( + uri: &str, +) -> Result<(impl AsyncRead + Send + Unpin, pmoflac::StreamInfo), SourceError> { + match classify_uri(uri) { + SourceKind::LocalFile(path) => { + let file = tokio::fs::File::open(&path).await + .map_err(|e| SourceError::Io(e.to_string()))?; + let stream = decode_audio_stream(file).await + .map_err(|e| SourceError::Decode(e.to_string()))?; + let info = stream.info().clone(); + Ok((stream, info)) + } + SourceKind::HttpUrl(url) => { + // Utiliser reqwest pour streamer l'URL HTTP externe + // Même approche que pmoaudiocache qui télécharge depuis URLs Qobuz + let response = reqwest::get(&url).await + .map_err(|e| SourceError::Http(e.to_string()))?; + let byte_stream = response.bytes_stream(); + // Convertir en AsyncRead + use tokio_util::io::StreamReader; + use futures::TryStreamExt; + let reader = StreamReader::new( + byte_stream.map_err(|e| std::io::Error::new(std::io::ErrorKind::Other, e)) + ); + let stream = decode_audio_stream(reader).await + .map_err(|e| SourceError::Decode(e.to_string()))?; + let info = stream.info().clone(); + Ok((stream, info)) + } + } +} + +#[derive(thiserror::Error, Debug)] +pub enum SourceError { + #[error("IO error: {0}")] + Io(String), + #[error("Decode error: {0}")] + Decode(String), + #[error("HTTP error: {0}")] + Http(String), +} +``` + +**Points d'attention :** +- `decode_audio_stream` prend un `AsyncRead + Send + Unpin`. `tokio::fs::File` et `StreamReader` satisfont ces contraintes. +- Pour les URLs Qobuz (URLs signées avec expiration), la source sera ouverte immédiatement par le handler `SetAVTransportURI` — pas de retry automatique. Si l'URL expire pendant la lecture, le pipeline s'arrêtera proprement via `StopReason`. +- Les chemins Samba supposent que le partage est monté localement sur le serveur. Aucun traitement spécial nécessaire, `tokio::fs::File::open` suffit. +- Ajouter `reqwest` en dépendance de pmowebrenderer avec feature `stream`. + +--- + +## Etape 5 : Créer `pipeline.rs` — logique complète du pipeline audio + +**Fichier :** `pmowebrenderer/src/pipeline.rs` (compléter l'étape 2) + +La task pipeline tourne en background pour chaque instance. Elle reçoit des `PipelineControl`, ouvre les sources via `source_loader`, alimente la `StreamingFlacSink`. + +```rust +use pmoaudio::{ + AudioChunk, AudioChunkData, AudioError, AudioSegment, _AudioSegment, + pipeline::{AudioPipelineNode, PipelineHandle as AudioPipelineHandle}, +}; +use pmoaudio_ext::sinks::streaming_flac_sink::{StreamHandle, StreamingFlacSink}; +use pmoflac::EncoderOptions; +use tokio::sync::mpsc; +use tokio_util::sync::CancellationToken; + +pub struct InstancePipeline { + pub stream_handle: StreamHandle, + pub control_tx: mpsc::Sender, + pub stop_token: CancellationToken, +} + +impl InstancePipeline { + /// Crée et démarre un pipeline pour une instance WebRenderer. + /// Retourne immédiatement ; la task tourne en background. + pub fn start() -> Self { + let stop_token = CancellationToken::new(); + let (control_tx, mut control_rx) = mpsc::channel::(16); + let (sink, stream_handle) = StreamingFlacSink::new( + EncoderOptions::default(), + 24, // bits_per_sample — 24 bits pour haute qualité + ); + + // La task de pipeline gère le cycle de vie + let stop_token_clone = stop_token.clone(); + let sink_arc: Arc>>> = + Arc::new(Mutex::new(None)); + + tokio::spawn(pipeline_task( + control_rx, + stop_token_clone, + sink, + )); + + Self { + stream_handle, + control_tx, + stop_token, + } + } +} + +async fn pipeline_task( + mut control_rx: mpsc::Receiver, + stop_token: CancellationToken, + sink: StreamingFlacSink, +) { + // État interne de la task + let mut current_pipeline_handle: Option = None; + let mut pending_next_uri: Option = None; + + // La StreamingFlacSink est démarrée une fois, tourne en continu. + // Le canal PCM (sink.get_tx()) reçoit les AudioSegments. + // Pour l'alimenter, on lance une task source séparée par piste. + + let sink_tx = sink.get_tx().expect("StreamingFlacSink must have a tx"); + let sink_stop = stop_token.clone(); + tokio::spawn(async move { + Box::new(sink).run(sink_stop).await.ok(); + }); + + loop { + tokio::select! { + _ = stop_token.cancelled() => { + // Stopper la source courante + if let Some(h) = current_pipeline_handle.take() { + h.stop(); + } + break; + } + + cmd = control_rx.recv() => { + match cmd { + None => break, + Some(PipelineControl::LoadUri(uri)) => { + // Stopper la source courante si elle existe + if let Some(h) = current_pipeline_handle.take() { + h.stop(); + } + // Lancer une nouvelle source + current_pipeline_handle = Some( + spawn_source_task(uri, sink_tx.clone(), stop_token.clone()).await + ); + } + Some(PipelineControl::LoadNextUri(uri)) => { + pending_next_uri = Some(uri); + } + Some(PipelineControl::Play) => { + // Pipeline serveur : Play ne fait rien de spécial + // La source alimente automatiquement dès LoadUri + } + Some(PipelineControl::Stop) => { + if let Some(h) = current_pipeline_handle.take() { + h.stop(); + } + // Envoyer EndOfStream au sink + let _ = sink_tx.send(Arc::new(AudioSegment::new_end_of_stream(0, 0.0))).await; + } + Some(PipelineControl::Seek(pos_sec)) => { + // Pour seek : stopper la source, relancer depuis la position + if let Some(current_uri) = get_current_uri() { + if let Some(h) = current_pipeline_handle.take() { + h.stop(); + } + current_pipeline_handle = Some( + spawn_source_task_from(current_uri, pos_sec, sink_tx.clone(), stop_token.clone()).await + ); + } + } + Some(PipelineControl::SetVolume(vol)) => { + // TODO : Volume DSP côté serveur (hors scope initial) + } + Some(PipelineControl::SetMute(mute)) => { + // TODO + } + _ => {} + } + } + } + } +} +``` + +**Points d'attention critiques :** + +1. **Gapless** : Le gapless avec `StreamingFlacSink` en mode `restart_encoder_on_track_boundary: false` signifie que le flux FLAC est continu. La frontière de piste est gérée par `SyncMarker::TrackBoundary`. Pour enchaîner deux sources indépendantes, il faut envoyer un `SyncMarker::TrackBoundary` entre les deux flux PCM vers le sink. La `StreamingFlacSink` ne restarte pas l'encodeur — le navigateur ne voit pas d'interruption. + +2. **Architecture source** : La source doit envoyer des `Arc` directement au `sink_tx` (le canal d'entrée de `StreamingFlacSink`). Ce n'est pas la même architecture que `Node` qui utilise les `children`. Ici, on alimente directement le sink via son `Sender>`. C'est exact car `sink.get_tx()` retourne le `Sender` d'entrée du `Node`. + +3. **Position tracking** : Sans retour du navigateur (le navigateur ne connaît que le flux FLAC), la position doit être trackée côté serveur. La `StreamingFlacSink` maintient un `current_timestamp` accessible via `StreamHandle` (indirectement via les métadonnées). Une task périodique lit ce timestamp et met à jour `RendererState.position`. + +4. **Seek** : Un seek implique de relancer la source depuis une position donnée. Pour les fichiers FLAC locaux, `pmoflac::decode_audio_stream` ne supporte pas le seek natif sur un stream. Il faudra rouvrir le fichier et lire en avançant les frames. Approche pragmatique : seek = stop + reopen + skip samples (coûteux mais simple). Le navigateur rebuffère ~1s comme indiqué dans l'archi. + +--- + +## Etape 6 : Créer `register.rs` — handlers POST /register et DELETE /{id} + +**Fichier :** `pmowebrenderer/src/register.rs` (nouveau) + +Ce module remplace `websocket.rs`. La logique de création du device UPnP est ici, adaptée du code de `create_renderer_for_browser` dans `websocket.rs`. + +```rust +use axum::{ + extract::{Path, State}, + http::StatusCode, + response::IntoResponse, + Json, +}; +use serde::{Deserialize, Serialize}; +use std::sync::Arc; + +use crate::registry::RendererRegistry; +use crate::pipeline::InstancePipeline; + +#[derive(Debug, Deserialize)] +pub struct RegisterRequest { + pub instance_id: String, + pub user_agent: String, +} + +#[derive(Debug, Serialize)] +pub struct RegisterResponse { + pub stream_url: String, +} + +/// POST /api/webrenderer/register +pub async fn register_handler( + State(registry): State>, + Json(req): Json, +) -> impl IntoResponse { + match registry.register_or_reconnect(&req.instance_id, &req.user_agent).await { + Ok(stream_url) => (StatusCode::OK, Json(RegisterResponse { stream_url })).into_response(), + Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()).into_response(), + } +} + +/// DELETE /api/webrenderer/{id} +pub async fn unregister_handler( + State(registry): State>, + Path(instance_id): Path, +) -> impl IntoResponse { + registry.unregister(&instance_id).await; + StatusCode::NO_CONTENT +} +``` + +**Points d'attention :** +- La logique de `create_renderer_for_browser` dans `websocket.rs` gère 4 cas (reconnexion session active, device dans registry mais session expirée, première connexion, etc.). Cette logique doit être reprise quasi à l'identique dans `RendererRegistry::register_or_reconnect`. +- La `stream_url` retournée est `/api/webrenderer/{instance_id}/stream` — URL relative, correcte en local et via proxy. +- Plus de `SharedSender` dans `WebRendererSession`. À la place : `stream_handle: SharedStreamHandle` + `pipeline: PipelineHandle`. + +--- + +## Etape 7 : Créer `stream.rs` — handler HTTP GET /stream + +**Fichier :** `pmowebrenderer/src/stream.rs` (nouveau) + +Ce module suit exactement le pattern de `pmomediaserver/src/paradise_streaming.rs`. + +```rust +use axum::{ + body::Body, + extract::{Path, State}, + http::{ + StatusCode, + header::{ACCEPT_RANGES, CACHE_CONTROL, CONNECTION, CONTENT_TYPE}, + }, + response::{IntoResponse, Response}, +}; +use std::sync::Arc; +use tokio_util::io::ReaderStream; + +use crate::registry::RendererRegistry; + +/// GET /api/webrenderer/{id}/stream +/// +/// Retourne le flux FLAC continu de l'instance. +/// La déconnexion du client est détectée automatiquement par la coupure du flux. +pub async fn stream_handler( + State(registry): State>, + Path(instance_id): Path, +) -> impl IntoResponse { + let handle = match registry.get_stream_handle(&instance_id).await { + Some(h) => h, + None => return StatusCode::NOT_FOUND.into_response(), + }; + + // Abonner ce client au flux FLAC + let flac_stream = handle.subscribe_flac(); + + Response::builder() + .status(StatusCode::OK) + .header(CONTENT_TYPE, "audio/flac") + .header(CACHE_CONTROL, "no-store, no-transform") + .header(CONNECTION, "keep-alive") + .header(ACCEPT_RANGES, "none") + .body(Body::from_stream(ReaderStream::new(flac_stream))) + .unwrap() + .into_response() +} +``` + +**Points d'attention :** +- `ReaderStream` de `tokio-util` convertit un `AsyncRead` en `Stream>`. `FlacClientStream` implémente `AsyncRead`. Le pattern est identique à `paradise_streaming.rs`. +- La détection de déconnexion se fait via `FlacClientStream::Drop` qui décrémente le compteur de clients dans `StreamHandle::client_disconnected()`. La `StreamingFlacSink` peut être configurée avec `set_auto_stop(true)` pour stopper le pipeline si plus aucun client n'écoute. +- Un client qui se reconnecte (`reload` de page) reçoit d'abord le header FLAC caché dans `SharedStreamHandleInner::header`, puis la suite du flux courant via `register_client()`. Ce mécanisme est déjà dans `SharedStreamHandleInner` — vérifier que `subscribe_flac()` envoie bien le header en premier (c'est le cas dans l'implémentation actuelle via `header_cache`). + +--- + +## Etape 8 : Créer `registry.rs` — remplace `session.rs` + +**Fichier :** `pmowebrenderer/src/registry.rs` (nouveau) + +Le `RendererRegistry` remplace `SessionManager`. La session est maintenant liée au flux FLAC, pas au WebSocket. + +```rust +use parking_lot::RwLock; +use std::collections::HashMap; +use std::sync::Arc; +use std::time::{Duration, SystemTime}; + +use pmoupnp::devices::DeviceInstance; + +#[cfg(feature = "audio-pipeline")] +use crate::pipeline::{InstancePipeline, PipelineHandle}; +#[cfg(feature = "audio-pipeline")] +use crate::state::SharedStreamHandle; +use crate::state::SharedState; + +/// Instance WebRenderer côté serveur +pub struct WebRendererInstance { + pub instance_id: String, // UUID stable du navigateur (localStorage) + pub udn: String, // "uuid:{instance_id}" + pub device_instance: Arc, + pub state: SharedState, + #[cfg(feature = "audio-pipeline")] + pub stream_handle: SharedStreamHandle, + #[cfg(feature = "audio-pipeline")] + pub pipeline: PipelineHandle, + pub created_at: SystemTime, + pub last_stream_connect: Arc>>, +} + +/// Registre global des instances WebRenderer actives +pub struct RendererRegistry { + instances: Arc>>>, + /// Map UDN → instance, pour retrouver par UDN depuis les handlers UPnP + by_udn: Arc>>>, + #[cfg(feature = "pmoserver")] + control_point: Arc, +} + +impl RendererRegistry { + /// Enregistre ou reconnecte une instance. + /// Retourne l'URL de stream relative. + pub async fn register_or_reconnect( + &self, + instance_id: &str, + user_agent: &str, + ) -> Result { + // Vérifier si une instance existe déjà pour cet instance_id + { + let instances = self.instances.read(); + if let Some(existing) = instances.get(instance_id) { + // Reconnexion : l'instance existe, le pipeline tourne toujours + // Réannoncer auprès du ControlPoint + #[cfg(feature = "pmoserver")] + self.register_with_control_point(&existing.device_instance)?; + + tracing::info!(instance_id, "WebRenderer: reconnected"); + return Ok(format!("/api/webrenderer/{}/stream", instance_id)); + } + } + + // Première connexion : créer device UPnP + pipeline + // (Reprend la logique de create_renderer_for_browser dans websocket.rs) + let instance = self.create_instance(instance_id, user_agent).await?; + let stream_url = format!("/api/webrenderer/{}/stream", instance_id); + + let instance = Arc::new(instance); + { + let mut instances = self.instances.write(); + instances.insert(instance_id.to_string(), instance.clone()); + } + { + let mut by_udn = self.by_udn.write(); + by_udn.insert(instance.udn.clone(), instance.clone()); + } + + tracing::info!(instance_id, "WebRenderer: registered new instance"); + Ok(stream_url) + } + + /// Retourne le StreamHandle pour l'endpoint /stream + pub async fn get_stream_handle(&self, instance_id: &str) -> Option { + self.instances.read() + .get(instance_id) + .and_then(|i| i.stream_handle.get()) + } + + /// Retourne le PipelineHandle pour les handlers UPnP + pub fn get_pipeline_by_udn(&self, udn: &str) -> Option { + self.by_udn.read() + .get(udn) + .map(|i| i.pipeline.clone()) + } + + pub fn get_state_by_udn(&self, udn: &str) -> Option { + self.by_udn.read() + .get(udn) + .map(|i| i.state.clone()) + } + + pub async fn unregister(&self, instance_id: &str) { + if let Some(instance) = self.instances.write().remove(instance_id) { + self.by_udn.write().remove(&instance.udn); + // Stopper le pipeline + instance.pipeline.stop_token.cancel(); + // Annoncer SSDP byebye + #[cfg(feature = "pmoserver")] + if let Ok(mut registry) = self.control_point.registry().write() { + registry.device_says_byebye(&instance.udn); + } + tracing::info!(instance_id, "WebRenderer: unregistered"); + } + } +} +``` + +**Points d'attention :** +- La méthode `create_instance` reprend quasi mot pour mot la logique de `create_renderer_for_browser` dans `websocket.rs` (vérifier registry DEVICE_REGISTRY, créer device model, `server.register_device`, `register_with_control_point`), mais au lieu de créer `SharedSender`, elle crée `InstancePipeline::start()`. +- Le cleanup "session expirée" de l'ancien `SessionManager` est remplacé par la détection de déconnexion du flux FLAC : quand `FlacClientStream` est droppé, `client_disconnected()` est appelé. Si `auto_stop` est activé, le pipeline s'arrête. +- Cleanup proactif via une task de polling : si le dernier client FLAC s'est déconnecté depuis plus de 5 minutes et que le ControlPoint indique que le renderer n'est plus en usage, on peut unregister. + +--- + +## Etape 9 : Modifier `handlers.rs` — brancher sur le pipeline au lieu du WS + +**Fichier :** `pmowebrenderer/src/handlers.rs` (modifier) + +Les handlers UPnP ne font plus `ws.send(...)` mais `pipeline.send(PipelineControl::...)`. + +```rust +// Nouveau play_handler +pub fn play_handler(pipeline: PipelineHandle, state: SharedState) -> ActionHandler { + Arc::new(move |data: ActionData| -> ActionFuture { + let pipeline = pipeline.clone(); + let state = state.clone(); + Box::pin(async move { + pipeline.send(PipelineControl::Play).await; + state.write().playback_state = PlaybackState::Playing; + Ok(data) + }) + }) +} + +// Nouveau set_uri_handler — la différence clé +pub fn set_uri_handler(pipeline: PipelineHandle, state: SharedState) -> ActionHandler { + Arc::new(move |data: ActionData| -> ActionFuture { + let pipeline = pipeline.clone(); + let state = state.clone(); + Box::pin(async move { + let uri: String = get!(&data, "CurrentURI", String); + let metadata: String = /* ... identique à l'existant ... */; + + // Envoyer au pipeline serveur (plus de WebSocket) + pipeline.send(PipelineControl::LoadUri(uri.clone())).await; + + { + let mut s = state.write(); + s.current_uri = Some(uri); + s.current_metadata = Some(metadata); + s.playback_state = PlaybackState::Transitioning; + } + Ok(data) + }) + }) +} + +// set_next_uri_handler +pub fn set_next_uri_handler(pipeline: PipelineHandle, state: SharedState) -> ActionHandler { + Arc::new(move |data: ActionData| -> ActionFuture { + let pipeline = pipeline.clone(); + let state = state.clone(); + Box::pin(async move { + let uri: String = get!(&data, "NextURI", String); + let metadata: String = /* ... */; + pipeline.send(PipelineControl::LoadNextUri(uri.clone())).await; + { + let mut s = state.write(); + s.next_uri = Some(uri); + s.next_metadata = Some(metadata); + } + Ok(data) + }) + }) +} + +// seek_handler +pub fn seek_handler(pipeline: PipelineHandle) -> ActionHandler { + Arc::new(move |data: ActionData| -> ActionFuture { + let pipeline = pipeline.clone(); + Box::pin(async move { + let target: String = get!(&data, "Target", String); + // Convertir "H:MM:SS" en secondes + let pos_sec = upnp_time_to_seconds(&target); + pipeline.send(PipelineControl::Seek(pos_sec)).await; + Ok(data) + }) + }) +} + +// Volume et Mute : côté serveur ou navigateur ? +// Phase initiale : volume = 100 fixe côté serveur, navigateur gère avec l'élément