Implémentation de la reconnexion stable avec synchronisation d'état

Cette mise à jour permet une reconnexion stable du renderer Web après un reload de page, en conservant l'état audio (URI, volume, lecture en cours). 

- Ajout d'un identifiant d'instance stable (UUID) dans le navigateur pour retrouver le même renderer UPnP
- Implémentation d'un mécanisme de synchronisation d'état (StateSync) lors des reconnexions
- Mise à jour du système de sender partagé (SharedSender) pour permettre le remplacement des connexions WebSocket
- Correction des handlers UPnP pour utiliser le nouveau type SharedSender
- Amélioration de la gestion des erreurs d'autoplay dans le moteur audio
- Mise à jour du numéro de version vers 0.3.21
This commit is contained in:
2026-02-21 20:54:22 +01:00
parent 32464c93fb
commit a502e4fc3e
9 changed files with 373 additions and 106 deletions

View File

@@ -1,6 +1,6 @@
[package] [package]
name = "PMOMusic" name = "PMOMusic"
version = "0.3.20" version = "0.3.21"
edition = "2024" edition = "2024"
[dependencies] [dependencies]

View File

@@ -22,6 +22,7 @@ import { ref, onMounted, onUnmounted, readonly } from "vue";
// ─── Types (miroir de messages.rs) ──────────────────────────────────────────── // ─── Types (miroir de messages.rs) ────────────────────────────────────────────
interface BrowserCapabilities { interface BrowserCapabilities {
instance_id: string;
user_agent: string; user_agent: string;
supported_formats: string[]; supported_formats: string[];
} }
@@ -41,8 +42,21 @@ interface CommandParams {
position?: string; position?: string;
} }
interface StateSyncMessage {
type: "state_sync";
current_uri?: string;
current_metadata?: string;
next_uri?: string;
next_metadata?: string;
playback_state: PlaybackState;
position?: string;
volume: number;
mute: boolean;
}
type ServerMessage = type ServerMessage =
| { type: "session_created"; token: string; renderer_info: RendererInfo } | { type: "session_created"; token: string; renderer_info: RendererInfo }
| StateSyncMessage
| { type: "command"; action: TransportAction; params?: CommandParams } | { type: "command"; action: TransportAction; params?: CommandParams }
| { type: "set_volume"; volume: number } | { type: "set_volume"; volume: number }
| { type: "set_mute"; mute: boolean } | { type: "set_mute"; mute: boolean }
@@ -58,6 +72,44 @@ type ClientMessage =
| { type: "track_ended" } | { type: "track_ended" }
| { type: "pong" }; | { type: "pong" };
// ─── Identifiant stable de l'instance navigateur ─────────────────────────────
const INSTANCE_ID_KEY = "pmomusic_webrenderer_instance_id";
/**
* Génère un UUID v4. Utilise crypto.randomUUID() si disponible (HTTPS/localhost),
* sinon fallback sur crypto.getRandomValues() (disponible partout, y compris HTTP).
*/
function generateUUID(): string {
if (typeof crypto.randomUUID === "function") {
return crypto.randomUUID();
}
const bytes = new Uint8Array(16);
crypto.getRandomValues(bytes);
bytes[6] = (bytes[6]! & 0x0f) | 0x40;
bytes[8] = (bytes[8]! & 0x3f) | 0x80;
const hex = Array.from(bytes).map((b) => b.toString(16).padStart(2, "0")).join("");
return `${hex.slice(0, 8)}-${hex.slice(8, 12)}-${hex.slice(12, 16)}-${hex.slice(16, 20)}-${hex.slice(20)}`;
}
/**
* Retourne un UUID stable pour cette instance navigateur.
* Généré une fois, persisté en localStorage, réutilisé entre les reloads.
*/
function getOrCreateInstanceId(): string {
try {
let id = localStorage.getItem(INSTANCE_ID_KEY);
if (!id) {
id = generateUUID();
localStorage.setItem(INSTANCE_ID_KEY, id);
}
return id;
} catch {
// localStorage unavailable (private mode, etc.) → use a session-scoped UUID
return generateUUID();
}
}
// ─── Détection des formats supportés ───────────────────────────────────────── // ─── Détection des formats supportés ─────────────────────────────────────────
function getSupportedFormats(): string[] { function getSupportedFormats(): string[] {
@@ -151,7 +203,41 @@ class GaplessEngine {
setCurrent(uri: string): void { setCurrent(uri: string): void {
this.nextUri = null; this.nextUri = null;
this.onStateChange("TRANSITIONING"); this.onStateChange("TRANSITIONING");
this._loadCurrent(uri);
// Si play() est arrivé avant set_uri (race condition), lancer la lecture maintenant
if (this.playPending) {
this.playPending = false;
this.play().catch((e) =>
console.error("[GaplessEngine] deferred play() failed:", e),
);
}
}
/**
* Restaure l'état audio après un reload sans notifier le serveur de TRANSITIONING.
* Le serveur connaît déjà l'état ; on recharge juste l'audio localement.
* Si shouldPlay=true mais que l'autoplay est bloqué, on notifie PAUSED
* pour que l'interface puisse proposer un bouton Play fonctionnel.
*/
async syncRestore(currentUri: string, nextUri: string | undefined, shouldPlay: boolean): Promise<void> {
this._loadCurrent(currentUri);
if (nextUri) {
this.setNext(nextUri);
}
if (shouldPlay) {
try {
await this.play();
// play() a réussi : onStateChange("PLAYING") a déjà été appelé dans play()
} catch {
// Autoplay bloqué par le navigateur : signaler PAUSED au serveur
// L'audio est chargé, un clic Play suffira à démarrer
this.onStateChange("PAUSED");
}
}
}
private _loadCurrent(uri: string): void {
const el = this.slots[this.currentSlot]; const el = this.slots[this.currentSlot];
// Retirer l'écouteur "ended" de l'autre slot si présent // Retirer l'écouteur "ended" de l'autre slot si présent
const otherSlot = (1 - this.currentSlot) as 0 | 1; const otherSlot = (1 - this.currentSlot) as 0 | 1;
@@ -167,14 +253,6 @@ class GaplessEngine {
el.onloadedmetadata = () => { el.onloadedmetadata = () => {
this._duration = el.duration || 0; this._duration = el.duration || 0;
}; };
// Si play() est arrivé avant set_uri (race condition), lancer la lecture maintenant
if (this.playPending) {
this.playPending = false;
this.play().catch((e) =>
console.error("[GaplessEngine] deferred play() failed:", e),
);
}
} }
/** Précharge la piste suivante dans l'autre slot. */ /** Précharge la piste suivante dans l'autre slot. */
@@ -214,8 +292,8 @@ class GaplessEngine {
try { try {
await el.play(); await el.play();
} catch (e) { } catch (e) {
console.error("[GaplessEngine] play() failed:", e); console.warn("[GaplessEngine] play() failed (autoplay blocked?):", e);
return; throw e;
} }
this.startPositionTimer(); this.startPositionTimer();
@@ -381,7 +459,9 @@ export function useWebRenderer() {
break; break;
case "play": case "play":
await engine.play(); await engine.play().catch((e) =>
console.warn("[WebRenderer] play() failed:", e),
);
break; break;
case "pause": case "pause":
@@ -418,6 +498,17 @@ export function useWebRenderer() {
onConnectedCallback?.(); onConnectedCallback?.();
break; break;
case "state_sync":
if (engine && msg.current_uri) {
engine.setVolume(msg.volume / 100);
engine.setMute(msg.mute);
const shouldPlay = msg.playback_state === "PLAYING" || msg.playback_state === "TRANSITIONING";
// syncRestore recharge l'audio sans notifier le serveur de l'état
// (le serveur connaît déjà l'état ; on évite un aller-retour TRANSITIONING)
engine.syncRestore(msg.current_uri, msg.next_uri, shouldPlay);
}
break;
case "command": case "command":
void execCommand(msg.action, msg.params); void execCommand(msg.action, msg.params);
break; break;
@@ -449,6 +540,7 @@ export function useWebRenderer() {
send({ send({
type: "init", type: "init",
capabilities: { capabilities: {
instance_id: getOrCreateInstanceId(),
user_agent: navigator.userAgent, user_agent: navigator.userAgent,
supported_formats: getSupportedFormats(), supported_formats: getSupportedFormats(),
}, },

View File

@@ -4,7 +4,6 @@
//! envoyée au navigateur, ou lit l'état partagé pour les requêtes GET. //! envoyée au navigateur, ou lit l'état partagé pour les requêtes GET.
use std::sync::Arc; use std::sync::Arc;
use tokio::sync::mpsc;
use pmodidl::DIDLLite; use pmodidl::DIDLLite;
use pmoupnp::actions::{ActionData, ActionError, ActionHandler, get_value}; use pmoupnp::actions::{ActionData, ActionError, ActionHandler, get_value};
@@ -12,19 +11,19 @@ use pmoupnp::{get, set};
use pmoutils::ToXmlElement; use pmoutils::ToXmlElement;
use crate::messages::{CommandParams, PlaybackState, ServerMessage, TransportAction}; use crate::messages::{CommandParams, PlaybackState, ServerMessage, TransportAction};
use crate::state::SharedState; use crate::state::{SharedSender, SharedState};
type ActionFuture = type ActionFuture =
std::pin::Pin<Box<dyn std::future::Future<Output = Result<ActionData, ActionError>> + Send>>; std::pin::Pin<Box<dyn std::future::Future<Output = Result<ActionData, ActionError>> + Send>>;
// ─── AVTransport Handlers ─────────────────────────────────────────────────── // ─── AVTransport Handlers ───────────────────────────────────────────────────
pub fn play_handler(ws: mpsc::UnboundedSender<ServerMessage>, state: SharedState) -> ActionHandler { pub fn play_handler(ws: SharedSender, state: SharedState) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture { Arc::new(move |data: ActionData| -> ActionFuture {
let ws = ws.clone(); let ws = ws.clone();
let state = state.clone(); let state = state.clone();
Box::pin(async move { Box::pin(async move {
let _ = ws.send(ServerMessage::Command { ws.send(ServerMessage::Command {
action: TransportAction::Play, action: TransportAction::Play,
params: None, params: None,
}); });
@@ -34,12 +33,12 @@ pub fn play_handler(ws: mpsc::UnboundedSender<ServerMessage>, state: SharedState
}) })
} }
pub fn stop_handler(ws: mpsc::UnboundedSender<ServerMessage>, state: SharedState) -> ActionHandler { pub fn stop_handler(ws: SharedSender, state: SharedState) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture { Arc::new(move |data: ActionData| -> ActionFuture {
let ws = ws.clone(); let ws = ws.clone();
let state = state.clone(); let state = state.clone();
Box::pin(async move { Box::pin(async move {
let _ = ws.send(ServerMessage::Command { ws.send(ServerMessage::Command {
action: TransportAction::Stop, action: TransportAction::Stop,
params: None, params: None,
}); });
@@ -49,15 +48,12 @@ pub fn stop_handler(ws: mpsc::UnboundedSender<ServerMessage>, state: SharedState
}) })
} }
pub fn pause_handler( pub fn pause_handler(ws: SharedSender, state: SharedState) -> ActionHandler {
ws: mpsc::UnboundedSender<ServerMessage>,
state: SharedState,
) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture { Arc::new(move |data: ActionData| -> ActionFuture {
let ws = ws.clone(); let ws = ws.clone();
let state = state.clone(); let state = state.clone();
Box::pin(async move { Box::pin(async move {
let _ = ws.send(ServerMessage::Command { ws.send(ServerMessage::Command {
action: TransportAction::Pause, action: TransportAction::Pause,
params: None, params: None,
}); });
@@ -67,11 +63,11 @@ pub fn pause_handler(
}) })
} }
pub fn next_handler(ws: mpsc::UnboundedSender<ServerMessage>) -> ActionHandler { pub fn next_handler(ws: SharedSender) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture { Arc::new(move |data: ActionData| -> ActionFuture {
let ws = ws.clone(); let ws = ws.clone();
Box::pin(async move { Box::pin(async move {
let _ = ws.send(ServerMessage::Command { ws.send(ServerMessage::Command {
action: TransportAction::Play, action: TransportAction::Play,
params: None, params: None,
}); });
@@ -80,11 +76,11 @@ pub fn next_handler(ws: mpsc::UnboundedSender<ServerMessage>) -> ActionHandler {
}) })
} }
pub fn previous_handler(ws: mpsc::UnboundedSender<ServerMessage>) -> ActionHandler { pub fn previous_handler(ws: SharedSender) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture { Arc::new(move |data: ActionData| -> ActionFuture {
let ws = ws.clone(); let ws = ws.clone();
Box::pin(async move { Box::pin(async move {
let _ = ws.send(ServerMessage::Command { ws.send(ServerMessage::Command {
action: TransportAction::Play, action: TransportAction::Play,
params: None, params: None,
}); });
@@ -93,12 +89,12 @@ pub fn previous_handler(ws: mpsc::UnboundedSender<ServerMessage>) -> ActionHandl
}) })
} }
pub fn seek_handler(ws: mpsc::UnboundedSender<ServerMessage>) -> ActionHandler { pub fn seek_handler(ws: SharedSender) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture { Arc::new(move |data: ActionData| -> ActionFuture {
let ws = ws.clone(); let ws = ws.clone();
Box::pin(async move { Box::pin(async move {
let target: String = get!(&data, "Target", String); let target: String = get!(&data, "Target", String);
let _ = ws.send(ServerMessage::Command { ws.send(ServerMessage::Command {
action: TransportAction::Seek, action: TransportAction::Seek,
params: Some(CommandParams { params: Some(CommandParams {
uri: None, uri: None,
@@ -111,23 +107,19 @@ pub fn seek_handler(ws: mpsc::UnboundedSender<ServerMessage>) -> ActionHandler {
}) })
} }
pub fn set_uri_handler( pub fn set_uri_handler(ws: SharedSender, state: SharedState) -> ActionHandler {
ws: mpsc::UnboundedSender<ServerMessage>,
state: SharedState,
) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture { Arc::new(move |data: ActionData| -> ActionFuture {
let ws = ws.clone(); let ws = ws.clone();
let state = state.clone(); let state = state.clone();
Box::pin(async move { Box::pin(async move {
let uri: String = get!(&data, "CurrentURI", String); let uri: String = get!(&data, "CurrentURI", String);
// CurrentURIMetaData peut être String ou DIDLLite (après parsing)
let metadata: String = get_value::<String>(&data, "CurrentURIMetaData") let metadata: String = get_value::<String>(&data, "CurrentURIMetaData")
.or_else(|_| { .or_else(|_| {
get_value::<DIDLLite>(&data, "CurrentURIMetaData") get_value::<DIDLLite>(&data, "CurrentURIMetaData")
.map(|didl| didl.to_xml()) .map(|didl| didl.to_xml())
}) })
.unwrap_or_default(); .unwrap_or_default();
let _ = ws.send(ServerMessage::Command { ws.send(ServerMessage::Command {
action: TransportAction::SetUri, action: TransportAction::SetUri,
params: Some(CommandParams { params: Some(CommandParams {
uri: Some(uri.clone()), uri: Some(uri.clone()),
@@ -146,10 +138,7 @@ pub fn set_uri_handler(
}) })
} }
pub fn set_next_uri_handler( pub fn set_next_uri_handler(ws: SharedSender, state: SharedState) -> ActionHandler {
ws: mpsc::UnboundedSender<ServerMessage>,
state: SharedState,
) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture { Arc::new(move |data: ActionData| -> ActionFuture {
let ws = ws.clone(); let ws = ws.clone();
let state = state.clone(); let state = state.clone();
@@ -161,7 +150,7 @@ pub fn set_next_uri_handler(
.map(|didl| didl.to_xml()) .map(|didl| didl.to_xml())
}) })
.unwrap_or_default(); .unwrap_or_default();
let _ = ws.send(ServerMessage::Command { ws.send(ServerMessage::Command {
action: TransportAction::SetNextUri, action: TransportAction::SetNextUri,
params: Some(CommandParams { params: Some(CommandParams {
uri: Some(uri.clone()), uri: Some(uri.clone()),
@@ -283,16 +272,13 @@ pub fn get_media_info_handler(state: SharedState) -> ActionHandler {
// ─── RenderingControl Handlers ────────────────────────────────────────────── // ─── RenderingControl Handlers ──────────────────────────────────────────────
pub fn set_volume_handler( pub fn set_volume_handler(ws: SharedSender, state: SharedState) -> ActionHandler {
ws: mpsc::UnboundedSender<ServerMessage>,
state: SharedState,
) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture { Arc::new(move |data: ActionData| -> ActionFuture {
let ws = ws.clone(); let ws = ws.clone();
let state = state.clone(); let state = state.clone();
Box::pin(async move { Box::pin(async move {
let volume: u16 = get!(&data, "DesiredVolume", u16); let volume: u16 = get!(&data, "DesiredVolume", u16);
let _ = ws.send(ServerMessage::SetVolume { volume }); ws.send(ServerMessage::SetVolume { volume });
state.write().volume = volume; state.write().volume = volume;
Ok(data) Ok(data)
}) })
@@ -311,16 +297,13 @@ pub fn get_volume_handler(state: SharedState) -> ActionHandler {
}) })
} }
pub fn set_mute_handler( pub fn set_mute_handler(ws: SharedSender, state: SharedState) -> ActionHandler {
ws: mpsc::UnboundedSender<ServerMessage>,
state: SharedState,
) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture { Arc::new(move |data: ActionData| -> ActionFuture {
let ws = ws.clone(); let ws = ws.clone();
let state = state.clone(); let state = state.clone();
Box::pin(async move { Box::pin(async move {
let mute: bool = get!(&data, "DesiredMute", bool); let mute: bool = get!(&data, "DesiredMute", bool);
let _ = ws.send(ServerMessage::SetMute { mute }); ws.send(ServerMessage::SetMute { mute });
state.write().mute = mute; state.write().mute = mute;
Ok(data) Ok(data)
}) })

View File

@@ -10,6 +10,18 @@ pub enum ServerMessage {
token: String, token: String,
renderer_info: RendererInfo, renderer_info: RendererInfo,
}, },
/// Envoyé après SessionCreated lors d'une reconnexion pour resynchroniser
/// l'état audio du navigateur (URI courante, état de lecture, etc.).
StateSync {
current_uri: Option<String>,
current_metadata: Option<String>,
next_uri: Option<String>,
next_metadata: Option<String>,
playback_state: PlaybackState,
position: Option<String>,
volume: u16,
mute: bool,
},
Command { Command {
action: TransportAction, action: TransportAction,
#[serde(skip_serializing_if = "Option::is_none")] #[serde(skip_serializing_if = "Option::is_none")]
@@ -62,6 +74,9 @@ pub enum ClientMessage {
#[derive(Debug, Clone, Serialize, Deserialize)] #[derive(Debug, Clone, Serialize, Deserialize)]
pub struct BrowserCapabilities { pub struct BrowserCapabilities {
/// Identifiant stable de l'instance navigateur (UUID stocké en localStorage).
/// Permet de réutiliser le même renderer UPnP après un reload de page.
pub instance_id: String,
pub user_agent: String, pub user_agent: String,
pub supported_formats: Vec<String>, pub supported_formats: Vec<String>,
} }

View File

@@ -13,7 +13,7 @@ use pmoupnp::services::Service;
use crate::handlers; use crate::handlers;
use crate::messages::ServerMessage; use crate::messages::ServerMessage;
use crate::state::SharedState; use crate::state::{SharedSender, SharedState};
// ─── Réimport des variables statiques de pmomediarenderer ─────────────────── // ─── Réimport des variables statiques de pmomediarenderer ───────────────────
// Variables AVTransport // Variables AVTransport
@@ -71,20 +71,25 @@ pub struct WebRendererFactory;
impl WebRendererFactory { impl WebRendererFactory {
/// Crée un Device model UPnP complet pour un WebRenderer. /// Crée un Device model UPnP complet pour un WebRenderer.
/// ///
/// Le device est construit avec des action handlers qui relaient /// `device_name` sert de clé pour retrouver l'UDN persistant dans la config.
/// les commandes SOAP vers le navigateur via le `ws_sender`. /// `browser_ua` est le User-Agent complet (pour déterminer le nom affiché).
pub fn create_device( ///
browser_name: &str, /// Retourne le Device et le `SharedSender` associé. Le `SharedSender` peut être
/// mis à jour à chaque reconnexion WebSocket via `shared_sender.set(new_tx)`.
pub fn create_device_with_name(
device_name: &str,
browser_ua: &str,
ws_sender: mpsc::UnboundedSender<ServerMessage>, ws_sender: mpsc::UnboundedSender<ServerMessage>,
state: SharedState, state: SharedState,
) -> Result<Device, FactoryError> { ) -> Result<(Device, SharedSender), FactoryError> {
let avtransport = Self::build_avtransport(ws_sender.clone(), state.clone())?; let shared_sender = SharedSender::new(ws_sender);
let renderingcontrol = Self::build_renderingcontrol(ws_sender.clone(), state.clone())?; let avtransport = Self::build_avtransport(shared_sender.clone(), state.clone())?;
let renderingcontrol = Self::build_renderingcontrol(shared_sender.clone(), state.clone())?;
let connectionmanager = Self::build_connectionmanager()?; let connectionmanager = Self::build_connectionmanager()?;
let short_name = extract_browser_name(browser_name); let short_name = extract_browser_name(browser_ua);
let device = Device::new( let device = Device::new(
"WebRenderer".to_string(), device_name.to_string(),
"MediaRenderer".to_string(), "MediaRenderer".to_string(),
format!("Web Audio {}", short_name), format!("Web Audio {}", short_name),
); );
@@ -98,12 +103,12 @@ impl WebRendererFactory {
.add_service(Arc::new(connectionmanager)) .add_service(Arc::new(connectionmanager))
.map_err(|e| FactoryError::ServiceError(format!("{:?}", e)))?; .map_err(|e| FactoryError::ServiceError(format!("{:?}", e)))?;
Ok(device) Ok((device, shared_sender))
} }
/// Construit le service AVTransport avec les handlers WebSocket /// Construit le service AVTransport avec les handlers WebSocket
fn build_avtransport( fn build_avtransport(
ws: mpsc::UnboundedSender<ServerMessage>, ws: SharedSender,
state: SharedState, state: SharedState,
) -> Result<Service, FactoryError> { ) -> Result<Service, FactoryError> {
let mut svc = Service::new("AVTransport".to_string()); let mut svc = Service::new("AVTransport".to_string());
@@ -419,7 +424,7 @@ impl WebRendererFactory {
/// Construit le service RenderingControl avec les handlers WebSocket /// Construit le service RenderingControl avec les handlers WebSocket
fn build_renderingcontrol( fn build_renderingcontrol(
ws: mpsc::UnboundedSender<ServerMessage>, ws: SharedSender,
state: SharedState, state: SharedState,
) -> Result<Service, FactoryError> { ) -> Result<Service, FactoryError> {
let mut svc = Service::new("RenderingControl".to_string()); let mut svc = Service::new("RenderingControl".to_string());

View File

@@ -4,19 +4,18 @@ use parking_lot::RwLock;
use std::collections::HashMap; use std::collections::HashMap;
use std::sync::Arc; use std::sync::Arc;
use std::time::{Duration, SystemTime}; use std::time::{Duration, SystemTime};
use tokio::sync::mpsc;
use pmoupnp::devices::DeviceInstance; use pmoupnp::devices::DeviceInstance;
use crate::messages::ServerMessage; use crate::state::{SharedSender, SharedState};
use crate::state::SharedState;
/// Session WebSocket liée à un MediaRenderer privé /// Session WebSocket liée à un MediaRenderer privé
pub struct WebRendererSession { pub struct WebRendererSession {
pub token: String, pub token: String,
pub udn: String, pub udn: String,
pub device_instance: Arc<DeviceInstance>, pub device_instance: Arc<DeviceInstance>,
pub ws_sender: mpsc::UnboundedSender<ServerMessage>, /// Sender partagé : mis à jour à chaque reconnexion WebSocket.
pub shared_sender: SharedSender,
pub state: SharedState, pub state: SharedState,
pub created_at: SystemTime, pub created_at: SystemTime,
pub last_activity: Arc<RwLock<SystemTime>>, pub last_activity: Arc<RwLock<SystemTime>>,
@@ -26,6 +25,12 @@ pub struct WebRendererSession {
#[derive(Clone)] #[derive(Clone)]
pub struct SessionManager { pub struct SessionManager {
sessions: Arc<RwLock<HashMap<String, Arc<WebRendererSession>>>>, sessions: Arc<RwLock<HashMap<String, Arc<WebRendererSession>>>>,
/// Map UDN → SharedSender, persiste même après suppression de la session.
/// Permet de retrouver et mettre à jour le sender à la reconnexion.
senders: Arc<RwLock<HashMap<String, SharedSender>>>,
/// Map UDN → SharedState, persiste même après suppression de la session.
/// Permet de réutiliser l'état partagé avec les handlers du device existant.
states: Arc<RwLock<HashMap<String, SharedState>>>,
timeout_duration: Duration, timeout_duration: Duration,
} }
@@ -33,6 +38,8 @@ impl SessionManager {
pub fn new(timeout_duration: Duration) -> Self { pub fn new(timeout_duration: Duration) -> Self {
let manager = Self { let manager = Self {
sessions: Arc::new(RwLock::new(HashMap::new())), sessions: Arc::new(RwLock::new(HashMap::new())),
senders: Arc::new(RwLock::new(HashMap::new())),
states: Arc::new(RwLock::new(HashMap::new())),
timeout_duration, timeout_duration,
}; };
@@ -42,7 +49,12 @@ impl SessionManager {
pub fn add_session(&self, session: Arc<WebRendererSession>) { pub fn add_session(&self, session: Arc<WebRendererSession>) {
let token = session.token.clone(); let token = session.token.clone();
let udn = session.udn.clone();
let sender = session.shared_sender.clone();
let state = session.state.clone();
self.sessions.write().insert(token.clone(), session); self.sessions.write().insert(token.clone(), session);
self.senders.write().insert(udn.clone(), sender);
self.states.write().insert(udn, state);
tracing::info!(token = %token, "WebRenderer session added"); tracing::info!(token = %token, "WebRenderer session added");
} }
@@ -56,6 +68,22 @@ impl SessionManager {
} }
} }
/// Retrouve une session par UDN du device (indépendant du token WebSocket).
pub fn get_session_by_udn(&self, udn: &str) -> Option<Arc<WebRendererSession>> {
let sessions = self.sessions.read();
sessions.values().find(|s| s.udn == udn).cloned()
}
/// Retrouve le SharedSender par UDN (persiste même après suppression de session).
pub fn get_sender_by_udn(&self, udn: &str) -> Option<SharedSender> {
self.senders.read().get(udn).cloned()
}
/// Retrouve le SharedState par UDN (persiste même après suppression de session).
pub fn get_state_by_udn(&self, udn: &str) -> Option<SharedState> {
self.states.read().get(udn).cloned()
}
pub fn remove_session(&self, token: &str) -> Option<Arc<WebRendererSession>> { pub fn remove_session(&self, token: &str) -> Option<Arc<WebRendererSession>> {
let session = self.sessions.write().remove(token); let session = self.sessions.write().remove(token);
if let Some(ref s) = session { if let Some(ref s) = session {

View File

@@ -2,8 +2,9 @@
use parking_lot::RwLock; use parking_lot::RwLock;
use std::sync::Arc; use std::sync::Arc;
use tokio::sync::mpsc;
use crate::messages::PlaybackState; use crate::messages::{PlaybackState, ServerMessage};
/// État temps-réel du renderer (partagé backend ↔ navigateur) /// État temps-réel du renderer (partagé backend ↔ navigateur)
#[derive(Debug, Clone)] #[derive(Debug, Clone)]
@@ -37,3 +38,34 @@ impl Default for RendererState {
/// Alias pour l'état partagé /// Alias pour l'état partagé
pub type SharedState = Arc<RwLock<RendererState>>; pub type SharedState = Arc<RwLock<RendererState>>;
/// Sender WebSocket partagé et remplaçable entre les reconnexions.
///
/// Les handlers UPnP capturent ce `Arc` à la création du device. À chaque
/// reconnexion WebSocket (reload de page), on remplace le sender interne via
/// `set()`, sans avoir à recréer le device ni ses handlers.
#[derive(Clone)]
pub struct SharedSender(Arc<RwLock<Option<mpsc::UnboundedSender<ServerMessage>>>>);
impl SharedSender {
pub fn new(sender: mpsc::UnboundedSender<ServerMessage>) -> Self {
Self(Arc::new(RwLock::new(Some(sender))))
}
/// Envoie un message au navigateur. Ignore silencieusement si déconnecté.
pub fn send(&self, msg: ServerMessage) {
if let Some(tx) = self.0.read().as_ref() {
let _ = tx.send(msg);
}
}
/// Remplace le sender (appelé à la reconnexion WebSocket).
pub fn set(&self, sender: mpsc::UnboundedSender<ServerMessage>) {
*self.0.write() = Some(sender);
}
/// Retire le sender (appelé à la déconnexion).
pub fn clear(&self) {
*self.0.write() = None;
}
}

View File

@@ -64,6 +64,7 @@ async fn handle_socket(socket: WebSocket, state: WebSocketState) {
}); });
let mut session_token: Option<String> = None; let mut session_token: Option<String> = None;
#[allow(unused_variables, unused_assignments, unused_mut)]
let mut device_udn: Option<String> = None; let mut device_udn: Option<String> = None;
// Boucle de réception des messages du navigateur // Boucle de réception des messages du navigateur
@@ -74,7 +75,7 @@ async fn handle_socket(socket: WebSocket, state: WebSocketState) {
match serde_json::from_str::<ClientMessage>(&text) { match serde_json::from_str::<ClientMessage>(&text) {
Ok(ClientMessage::Init { capabilities }) => { Ok(ClientMessage::Init { capabilities }) => {
tracing::info!("WebRenderer Init received, creating renderer..."); tracing::info!("WebRenderer Init received, creating renderer...");
// Créer le renderer UPnP pour ce navigateur // Créer ou reconnecter le renderer UPnP pour ce navigateur.
match create_renderer_for_browser(&capabilities, tx.clone(), &state).await { match create_renderer_for_browser(&capabilities, tx.clone(), &state).await {
Ok(session) => { Ok(session) => {
tracing::info!("WebRenderer create_renderer_for_browser OK"); tracing::info!("WebRenderer create_renderer_for_browser OK");
@@ -98,8 +99,32 @@ async fn handle_socket(socket: WebSocket, state: WebSocketState) {
}, },
}); });
// Si une URI est déjà chargée (reconnexion en cours de lecture),
// envoyer l'état complet pour que le navigateur puisse reprendre.
{
let s = session.state.read();
if s.current_uri.is_some() {
let _ = tx.send(ServerMessage::StateSync {
current_uri: s.current_uri.clone(),
current_metadata: s.current_metadata.clone(),
next_uri: s.next_uri.clone(),
next_metadata: s.next_metadata.clone(),
playback_state: s.playback_state.clone(),
position: s.position.clone(),
volume: s.volume,
mute: s.mute,
});
tracing::info!(
udn = %udn,
state = ?s.playback_state,
"WebRenderer: sent StateSync to reconnected browser"
);
}
}
session_token = Some(token.clone()); session_token = Some(token.clone());
device_udn = Some(udn.clone()); #[cfg(feature = "pmoserver")]
{ device_udn = Some(udn.clone()); }
state.session_manager.add_session(session); state.session_manager.add_session(session);
@@ -237,7 +262,13 @@ async fn handle_socket(socket: WebSocket, state: WebSocketState) {
} }
} }
} }
Ok(Message::Close(_)) => break, Ok(Message::Binary(b)) => {
tracing::warn!("WebRenderer received binary message ({} bytes)", b.len());
}
Ok(Message::Close(_)) => {
tracing::info!("WebRenderer WebSocket closed by client");
break;
}
Err(e) => { Err(e) => {
tracing::error!("WebSocket error: {}", e); tracing::error!("WebSocket error: {}", e);
break; break;
@@ -247,6 +278,7 @@ async fn handle_socket(socket: WebSocket, state: WebSocketState) {
} }
// Cleanup à la déconnexion // Cleanup à la déconnexion
tracing::info!("WebRenderer WebSocket handler exiting (session_token={:?})", session_token);
if let Some(token) = session_token { if let Some(token) = session_token {
state.session_manager.remove_session(&token); state.session_manager.remove_session(&token);
} }
@@ -263,65 +295,145 @@ async fn handle_socket(socket: WebSocket, state: WebSocketState) {
send_task.abort(); send_task.abort();
} }
/// Crée un DeviceInstance UPnP et l'enregistre pour un navigateur /// Crée ou reconnecte un DeviceInstance UPnP pour un navigateur.
///
/// - Première connexion : crée le device, l'enregistre, crée la session.
/// - Reconnexion (reload) : retrouve la session existante par UDN, met à jour le
/// `SharedSender` avec le nouveau tx WebSocket (les handlers continuent de fonctionner),
/// et crée une nouvelle session avec un nouveau token.
async fn create_renderer_for_browser( async fn create_renderer_for_browser(
capabilities: &BrowserCapabilities, capabilities: &BrowserCapabilities,
ws_sender: mpsc::UnboundedSender<ServerMessage>, ws_sender: mpsc::UnboundedSender<ServerMessage>,
ws_state: &WebSocketState, ws_state: &WebSocketState,
) -> Result<Arc<WebRendererSession>, crate::error::WebRendererError> { ) -> Result<Arc<WebRendererSession>, crate::error::WebRendererError> {
let shared_state: SharedState = Arc::new(RwLock::new(RendererState::default()));
let token = Uuid::new_v4().to_string(); let token = Uuid::new_v4().to_string();
// Construire le Device model avec les handlers WS // Persister l'UDN dérivé de l'instance_id dans la config pour que device_instance.rs
tracing::info!("WebRenderer: creating device model..."); // le retrouve de façon déterministe. La clé ("MediaRenderer", instance_id) est unique
let device = WebRendererFactory::create_device( // par onglet/navigateur et stable entre les reloads.
&capabilities.user_agent, let instance_udn = capabilities.instance_id.clone();
ws_sender.clone(), if let Err(e) = pmoconfig::get_config().set_device_udn(
shared_state.clone(), "MediaRenderer",
) &instance_udn,
.map_err(|e| crate::error::WebRendererError::DeviceCreationError(e.to_string()))?; instance_udn.clone(),
tracing::info!("WebRenderer: device model created"); ) {
tracing::warn!("WebRenderer: failed to persist UDN in config: {:?}", e);
}
let device = Arc::new(device); // UDN normalisé tel que stocké dans le DEVICE_REGISTRY (sans préfixe "uuid:")
let candidate_udn = instance_udn.to_ascii_lowercase();
// UDN avec préfixe "uuid:" pour le ControlPoint et la session
let full_udn = format!("uuid:{}", candidate_udn);
// ── Reconnexion : session existante par UDN ───────────────────────────────
// Si une session avec ce même UDN existe encore dans le SessionManager, on
// met à jour son SharedSender (les handlers UPnP enverront vers le nouveau WS).
if let Some(existing_session) = ws_state.session_manager.get_session_by_udn(&full_udn) {
tracing::info!(udn = %full_udn, "WebRenderer: reconnecting via existing session");
existing_session.shared_sender.set(ws_sender.clone());
#[cfg(feature = "pmoserver")]
register_with_control_point(&existing_session.device_instance, ws_state)?;
// Nouvelle session avec nouveau token, mais même device/state/sender partagés
let session = Arc::new(WebRendererSession {
token,
udn: full_udn,
device_instance: existing_session.device_instance.clone(),
shared_sender: existing_session.shared_sender.clone(),
state: existing_session.state.clone(),
created_at: existing_session.created_at,
last_activity: existing_session.last_activity.clone(),
});
return Ok(session);
}
// ── Première connexion : création complète ────────────────────────────────
// Enregistrer le device via UpnpServerExt (gère base_url, register_urls, DEVICE_REGISTRY) // Enregistrer le device via UpnpServerExt (gère base_url, register_urls, DEVICE_REGISTRY)
// Retourne (DeviceInstance, SharedSender effectif, SharedState effective pour cette session)
#[cfg(feature = "pmoserver")] #[cfg(feature = "pmoserver")]
let di = { let (di, shared_sender, shared_state) = {
use pmoupnp::UpnpServerExt; use pmoupnp::UpnpServerExt;
tracing::info!("WebRenderer: getting server arc..."); tracing::info!("WebRenderer: candidate UDN = {}", candidate_udn);
// Vérifier si un device avec ce même UDN est déjà dans le DEVICE_REGISTRY
// (cas où la session a expiré mais le device est encore enregistré).
let server_arc = let server_arc =
pmoserver::get_server().ok_or(crate::error::WebRendererError::ServerNotAvailable)?; pmoserver::get_server().ok_or(crate::error::WebRendererError::ServerNotAvailable)?;
tracing::info!("WebRenderer: got server arc, acquiring write lock..."); let existing_di = {
let server = server_arc.read().await;
let di = { server.get_device(&candidate_udn)
let mut server = server_arc.write().await;
tracing::info!("WebRenderer: write lock acquired, registering device...");
server
.register_device(device)
.await
.map_err(|e| crate::error::WebRendererError::RegistrationError(e.to_string()))?
}; };
tracing::info!("WebRenderer: device registered, registering with ControlPoint...");
// Enregistrer avec le ControlPoint if let Some(di) = existing_di {
register_with_control_point(&di, ws_state)?; tracing::info!(udn = %candidate_udn, "WebRenderer: reusing device from registry (session expired)");
tracing::info!("WebRenderer: registered with ControlPoint"); // Mettre à jour le SharedSender de ce device (session supprimée mais device toujours dans registry).
// Le SharedSender et le SharedState sont ceux capturés dans les handlers du di existant.
let effective_sender = if let Some(existing_sender) = ws_state.session_manager.get_sender_by_udn(&full_udn) {
existing_sender.set(ws_sender);
tracing::info!(udn = %full_udn, "WebRenderer: updated SharedSender for reused device");
existing_sender
} else {
// Fallback : ne devrait pas arriver mais on crée un sender neuf
tracing::warn!(udn = %full_udn, "WebRenderer: no SharedSender found for reused device");
let new_state: SharedState = Arc::new(RwLock::new(RendererState::default()));
let (_, new_sender) = WebRendererFactory::create_device_with_name(
&instance_udn, &capabilities.user_agent, ws_sender, new_state.clone(),
).map_err(|e| crate::error::WebRendererError::DeviceCreationError(e.to_string()))?;
new_sender
};
let effective_state = ws_state.session_manager.get_state_by_udn(&full_udn)
.unwrap_or_else(|| Arc::new(RwLock::new(RendererState::default())));
register_with_control_point(&di, ws_state)?;
(di, effective_sender, effective_state)
} else {
// Véritablement première connexion : créer device + state + sender
let new_state: SharedState = Arc::new(RwLock::new(RendererState::default()));
tracing::info!("WebRenderer: creating device model...");
let (device, new_sender) = WebRendererFactory::create_device_with_name(
&instance_udn,
&capabilities.user_agent,
ws_sender,
new_state.clone(),
)
.map_err(|e| crate::error::WebRendererError::DeviceCreationError(e.to_string()))?;
tracing::info!("WebRenderer: device model created");
di let device = Arc::new(device);
tracing::info!("WebRenderer: registering new device...");
let di = {
let mut server = server_arc.write().await;
server
.register_device(device)
.await
.map_err(|e| crate::error::WebRendererError::RegistrationError(e.to_string()))?
};
tracing::info!("WebRenderer: device registered");
register_with_control_point(&di, ws_state)?;
(di, new_sender, new_state)
}
}; };
#[cfg(not(feature = "pmoserver"))] #[cfg(not(feature = "pmoserver"))]
let di = device.create_instance(); let (di, shared_sender, shared_state) = {
let new_state: SharedState = Arc::new(RwLock::new(RendererState::default()));
// Normaliser l'UDN avec le préfixe "uuid:" pour correspondre au format SSDP let (device, new_sender) = WebRendererFactory::create_device_with_name(
let udn = format!("uuid:{}", di.udn().to_ascii_lowercase()); &instance_udn,
&capabilities.user_agent,
ws_sender,
new_state.clone(),
)
.map_err(|e| crate::error::WebRendererError::DeviceCreationError(e.to_string()))?;
(Arc::new(device).create_instance(), new_sender, new_state)
};
let session = Arc::new(WebRendererSession { let session = Arc::new(WebRendererSession {
token, token,
udn, udn: full_udn,
device_instance: di, device_instance: di,
ws_sender, shared_sender,
state: shared_state, state: shared_state,
created_at: SystemTime::now(), created_at: SystemTime::now(),
last_activity: Arc::new(RwLock::new(SystemTime::now())), last_activity: Arc::new(RwLock::new(SystemTime::now())),

View File

@@ -1 +1 @@
0.3.20 0.3.21