Implémentation du WebRenderer et améliorations du serveur

Ajout de la fonctionnalité WebRenderer permettant au navigateur de se connecter comme renderer UPnP.

- Implémentation du composant Vue useWebRenderer pour gérer la connexion WebSocket et le contrôle audio
- Ajout d'un proxy WebSocket dans Vite pour rediriger les requêtes /api/webrenderer/ws vers le serveur
- Modification du serveur pour permettre l'enregistrement dynamique de routes et extraction du JoinHandle
- Améliorations du parsing des métadonnées AVTransport pour gérer les types String et DIDLLite
- Ajout de la détection et du nommage des navigateurs dans le WebRenderer
- Mise à jour des dépendances avec tower 0.5.2
- Correction de la gestion des threads lors de l'arrêt du serveur HTTP
This commit is contained in:
2026-02-21 00:15:38 +01:00
parent 7276945b92
commit cddd2afbd8
11 changed files with 443 additions and 12 deletions

1
Cargo.lock generated
View File

@@ -4265,6 +4265,7 @@ dependencies = [
"tokio",
"tokio-stream",
"tokio-util",
"tower 0.5.2",
"tracing",
"tracing-subscriber",
"utoipa",

View File

@@ -137,8 +137,13 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
info!("✅ PMOMusic is ready!");
info!("Press Ctrl+C to stop...");
// Attendre le signal Ctrl+C et l'arrêt du serveur HTTP
server.write().await.wait().await;
// Extraire le join_handle AVANT de libérer le write lock,
// pour pouvoir l'awaiter sans tenir le write lock du serveur global.
// (Tenir le write lock pendant wait() bloquerait register_device() dynamique)
let join_handle = server.write().await.take_join_handle();
if let Some(h) = join_handle {
let _ = h.await;
}
// Le serveur HTTP est arrêté, mais des threads (ControlPoint, etc.) peuvent encore tourner
// Attendre 2 secondes pour laisser le temps aux threads de se terminer

View File

@@ -0,0 +1,332 @@
/**
* Composable pour gérer le WebRenderer navigateur.
*
* Se connecte automatiquement au WebSocket /api/webrenderer/ws au montage
* et se déconnecte proprement au démontage ou à la fermeture de la page.
*
* Le navigateur est ainsi vu comme un renderer UPnP par le ControlPoint.
* Un élément <audio> headless exécute les commandes de transport reçues.
*/
import { ref, onMounted, onUnmounted, readonly } from "vue";
// ─── Types (miroir de messages.rs) ────────────────────────────────────────────
interface BrowserCapabilities {
user_agent: string;
supported_formats: string[];
}
interface RendererInfo {
udn: string;
friendly_name: string;
model_name: string;
description_url: string;
}
type TransportAction = "play" | "pause" | "stop" | "seek" | "set_uri";
interface CommandParams {
uri?: string;
metadata?: string;
position?: string;
}
type ServerMessage =
| { type: "session_created"; token: string; renderer_info: RendererInfo }
| { type: "command"; action: TransportAction; params?: CommandParams }
| { type: "set_volume"; volume: number }
| { type: "set_mute"; mute: boolean }
| { type: "ping" };
type PlaybackState = "PLAYING" | "PAUSED" | "STOPPED" | "TRANSITIONING";
type ClientMessage =
| { type: "init"; capabilities: BrowserCapabilities }
| { type: "state_update"; state: PlaybackState }
| { type: "position_update"; position: string; duration: string }
| { type: "volume_update"; volume: number; mute: boolean }
| { type: "pong" };
// ─── Détection des formats supportés ─────────────────────────────────────────
function getSupportedFormats(): string[] {
const audio = document.createElement("audio");
const formats: Array<[string, string]> = [
["mp3", "audio/mpeg"],
["flac", "audio/flac"],
["ogg", "audio/ogg; codecs=vorbis"],
["opus", "audio/ogg; codecs=opus"],
["aac", "audio/aac"],
["wav", "audio/wav"],
["m4a", 'audio/mp4; codecs="mp4a.40.2"'],
["webm", "audio/webm"],
];
return formats
.filter(([, mime]) => audio.canPlayType(mime) !== "")
.map(([fmt]) => fmt);
}
// ─── Helpers ──────────────────────────────────────────────────────────────────
function secondsToUpnpTime(s: number): string {
const h = Math.floor(s / 3600);
const m = Math.floor((s % 3600) / 60);
const sec = Math.floor(s % 60);
return `${h}:${String(m).padStart(2, "0")}:${String(sec).padStart(2, "0")}`;
}
function upnpTimeToSeconds(t: string): number {
const parts = t.split(":").map(Number);
if (parts.length !== 3) return 0;
const [h, m, s] = parts;
return (h ?? 0) * 3600 + (m ?? 0) * 60 + (s ?? 0);
}
// ─── Composable ───────────────────────────────────────────────────────────────
export function useWebRenderer() {
const connected = ref(false);
const rendererInfo = ref<RendererInfo | null>(null);
let ws: WebSocket | null = null;
let audio: HTMLAudioElement | null = null;
let positionTimer: ReturnType<typeof setInterval> | null = null;
let onConnectedCallback: (() => void) | null = null;
// ── Envoi d'un message au backend ────────────────────────────────────────
function send(msg: ClientMessage) {
if (ws && ws.readyState === WebSocket.OPEN) {
ws.send(JSON.stringify(msg));
}
}
// ── Player audio headless ─────────────────────────────────────────────────
function sendPosition() {
if (!audio) return;
const pos = isFinite(audio.currentTime) ? audio.currentTime : 0;
const dur = isFinite(audio.duration) ? audio.duration : 0;
send({
type: "position_update",
position: secondsToUpnpTime(pos),
duration: secondsToUpnpTime(dur),
});
}
function startPositionTimer() {
if (positionTimer !== null) return;
positionTimer = setInterval(sendPosition, 1000);
}
function stopPositionTimer() {
if (positionTimer !== null) {
clearInterval(positionTimer);
positionTimer = null;
}
}
function createAudio(): HTMLAudioElement {
const el = new Audio();
el.preload = "auto";
el.addEventListener("play", () => {
send({ type: "state_update", state: "PLAYING" });
startPositionTimer();
});
el.addEventListener("pause", () => {
send({ type: "state_update", state: "PAUSED" });
stopPositionTimer();
sendPosition();
});
el.addEventListener("ended", () => {
send({ type: "state_update", state: "STOPPED" });
stopPositionTimer();
});
el.addEventListener("waiting", () => {
send({ type: "state_update", state: "TRANSITIONING" });
});
el.addEventListener("canplay", () => {
if (!el.paused) send({ type: "state_update", state: "PLAYING" });
});
el.addEventListener("volumechange", () => {
const vol = Math.round(el.volume * 100);
send({ type: "volume_update", volume: vol, mute: el.muted });
});
el.addEventListener("error", () => {
console.error("[WebRenderer] Erreur audio :", el.error);
send({ type: "state_update", state: "STOPPED" });
stopPositionTimer();
});
return el;
}
// ── Exécution des commandes UPnP ──────────────────────────────────────────
async function execCommand(action: TransportAction, params?: CommandParams) {
if (!audio) return;
console.debug(`[WebRenderer] Commande: ${action}`, params);
switch (action) {
case "set_uri":
if (params?.uri) {
audio.pause();
audio.src = params.uri;
audio.load();
send({ type: "state_update", state: "STOPPED" });
}
break;
case "play":
if (audio.src) {
try {
await audio.play();
} catch (e) {
console.error("[WebRenderer] play() refusé :", e);
}
}
break;
case "pause":
audio.pause();
break;
case "stop":
audio.pause();
audio.currentTime = 0;
send({ type: "state_update", state: "STOPPED" });
break;
case "seek":
if (params?.position) {
audio.currentTime = upnpTimeToSeconds(params.position);
}
break;
}
}
// ── Gestion des messages entrants ─────────────────────────────────────────
function handleMessage(event: MessageEvent) {
let msg: ServerMessage;
try {
msg = JSON.parse(event.data as string) as ServerMessage;
} catch {
console.warn("[WebRenderer] Message non-JSON reçu :", event.data);
return;
}
switch (msg.type) {
case "session_created":
rendererInfo.value = msg.renderer_info;
connected.value = true;
console.info(
`[WebRenderer] Session créée — UDN: ${msg.renderer_info.udn}`,
);
// Notifier le parent pour qu'il rafraîchisse la liste des renderers
// L'événement SSE peut arriver avant que le subscriber soit prêt
onConnectedCallback?.();
break;
case "command":
void execCommand(msg.action, msg.params);
break;
case "set_volume":
if (audio) {
audio.volume = Math.max(0, Math.min(1, msg.volume / 100));
}
break;
case "set_mute":
if (audio) {
audio.muted = msg.mute;
}
break;
case "ping":
send({ type: "pong" });
break;
}
}
// ── Connexion ─────────────────────────────────────────────────────────────
function connect() {
if (ws) return;
const protocol = location.protocol === "https:" ? "wss:" : "ws:";
const url = `${protocol}//${location.host}/api/webrenderer/ws`;
console.info(`[WebRenderer] Connexion à : ${url}`);
ws = new WebSocket(url);
ws.onopen = () => {
console.info("[WebRenderer] WebSocket ouvert, envoi Init");
send({
type: "init",
capabilities: {
user_agent: navigator.userAgent,
supported_formats: getSupportedFormats(),
},
});
};
ws.onmessage = handleMessage;
ws.onclose = () => {
connected.value = false;
rendererInfo.value = null;
ws = null;
console.info("[WebRenderer] Déconnecté");
};
ws.onerror = (err) => {
console.error("[WebRenderer] Erreur WebSocket :", err);
};
}
// ── Déconnexion ───────────────────────────────────────────────────────────
function disconnect() {
stopPositionTimer();
if (audio) {
audio.pause();
audio.src = "";
}
if (ws) {
ws.close(1000, "Page unloaded");
ws = null;
}
connected.value = false;
}
// ── Cycle de vie ──────────────────────────────────────────────────────────
onMounted(() => {
audio = createAudio();
connect();
window.addEventListener("beforeunload", disconnect);
});
onUnmounted(() => {
disconnect();
audio = null;
window.removeEventListener("beforeunload", disconnect);
});
// ── API publique ──────────────────────────────────────────────────────────
return {
/** true quand la session WebRenderer est établie */
connected: readonly(connected),
/** Infos du renderer UPnP créé pour ce navigateur */
rendererInfo: readonly(rendererInfo),
/** Callback appelé quand la session est créée (pour rafraîchir la liste des renderers) */
onConnected(fn: () => void) {
onConnectedCallback = fn;
},
};
}

View File

@@ -4,6 +4,7 @@ import { useRoute, useRouter } from "vue-router";
import { useTabs } from "@/composables/useTabs";
import { useRenderers } from "@/composables/useRenderers";
import { useMediaServers } from "@/composables/useMediaServers";
import { useWebRenderer } from "@/composables/useWebRenderer";
import { useSwipe } from "@vueuse/core";
// Import des composants
@@ -21,6 +22,11 @@ const { tabs, activeTabId, switchTab, activeTab, syncWithRenderers, isEmpty } =
const { allRenderers, fetchRenderers, getStateById } = useRenderers();
const { allServers, fetchServers } = useMediaServers();
// WebRenderer : ce navigateur s'enregistre automatiquement comme renderer UPnP
const webRenderer = useWebRenderer();
// Rafraîchir la liste des renderers dès que la session WebRenderer est établie
webRenderer.onConnected(() => void fetchRenderers(true));
// État des drawers
const drawerOpen = ref(false);
const rendererDrawerOpen = ref(false);
@@ -97,13 +103,24 @@ function handleRendererSelect(rendererId: string) {
}
}
// Filtre la liste des renderers pour exclure les WebRenderers d'autres navigateurs.
// Seul le WebRenderer créé par ce navigateur (identifié par son UDN) reste visible.
function filterRenderers(renderers: typeof allRenderers.value) {
const myUdn = webRenderer.rendererInfo.value?.udn ?? null;
return renderers.filter((r) => {
if (r.model_name !== "WebRenderer") return true; // renderer classique : toujours visible
if (myUdn === null) return false; // pas encore de session : masquer tous les WebRenderers
return r.id === myUdn; // ne garder que le nôtre
});
}
// Sync route query params avec l'état des tabs
onMounted(async () => {
// Fetch renderers et servers au montage
await Promise.all([fetchRenderers(), fetchServers()]);
// Sync initial des tabs avec les renderers
syncWithRenderers(allRenderers.value);
// Sync initial des tabs avec les renderers (en excluant les WebRenderers étrangers)
syncWithRenderers(filterRenderers(allRenderers.value));
// Restaurer l'onglet actif depuis l'URL
const urlTabId = route.query.tab as string;
@@ -116,11 +133,19 @@ onMounted(async () => {
watch(
() => allRenderers.value,
(newRenderers) => {
syncWithRenderers(newRenderers);
syncWithRenderers(filterRenderers(newRenderers));
},
{ deep: true },
);
// Watch l'UDN du WebRenderer local : quand il s'établit, resync pour faire apparaître notre onglet
watch(
() => webRenderer.rendererInfo.value?.udn,
() => {
syncWithRenderers(filterRenderers(allRenderers.value));
},
);
// Watch les changements d'URL pour changer d'onglet
watch(
() => route.query.tab,

View File

@@ -13,6 +13,11 @@ export default defineConfig({
},
server: {
proxy: {
'/api/webrenderer/ws': {
target: 'ws://localhost:8080',
ws: true,
changeOrigin: true,
},
'/api': {
target: 'http://localhost:8080',
changeOrigin: true,

View File

@@ -31,6 +31,11 @@ fn avtransporturimetadataparser(value: &str) -> Result<Box<dyn Reflect>, StateVa
}
fn avtransporturimetadatamarshal(value: &dyn Reflect) -> Result<String, StateVariableError> {
// Si c'est déjà une String, on la retourne directement
if let Some(s) = value.as_any().downcast_ref::<String>() {
return Ok(s.clone());
}
// Sinon, on essaie de convertir depuis DIDLLite
let didl = value
.downcast_ref::<DIDLLite>()
.ok_or_else(|| StateVariableError::ConversionError("DIDLLite".into()))?;
@@ -45,6 +50,8 @@ pub static AVTRANSPORTURIMETADATA: Lazy<Arc<StateVariable>> =
sv.set_value_parser(Arc::new(avtransporturimetadataparser))
.expect("Failed to set parser");
sv.set_value_marshaler(Arc::new(avtransporturimetadatamarshal))
.expect("Failed to set marshaler");
Arc::new(sv)
});

View File

@@ -8,6 +8,7 @@ pmoconfig = { path = "../pmoconfig", features = ["api"] }
anyhow = { workspace = true }
axum = "0.8.4"
tower = { version = "0.5", features = ["util"] }
tokio = { workspace = true, features = ["rt-multi-thread", "macros", "sync", "time", "signal"] }
tokio-stream = "0.1"
tokio-util = "0.7"

View File

@@ -550,7 +550,6 @@ impl Server {
self.join_handle = Some(tokio::spawn(async move {
let server_future = async {
let r = router.read().await.clone();
let listener = match tokio::net::TcpListener::bind(addr).await {
Ok(l) => l,
Err(e) => {
@@ -559,7 +558,19 @@ impl Server {
}
};
axum::serve(listener, r.into_make_service())
// Utiliser un router dynamique qui relit le router à chaque requête.
// Cela permet d'enregistrer de nouvelles routes après le démarrage du serveur
// (ex: WebRenderer dynamique).
let dynamic_router = axum::Router::new().fallback(move |req: axum::extract::Request| {
let router = router.clone();
async move {
use tower::ServiceExt;
let r = router.read().await.clone();
r.into_service::<axum::body::Body>().oneshot(req).await
}
});
axum::serve(listener, dynamic_router.into_make_service())
.with_graceful_shutdown(async move {
let _ = shutdown_rx.await;
})
@@ -598,6 +609,12 @@ impl Server {
}
}
/// Extrait le JoinHandle du serveur HTTP pour pouvoir l'awaiter
/// sans tenir le write lock du serveur global.
pub fn take_join_handle(&mut self) -> Option<JoinHandle<()>> {
self.join_handle.take()
}
/// Retourne l'URL de base complète du serveur (schéma + hôte + port).
///
/// La valeur configurable peut omettre le schéma ou le port ; cette méthode

View File

@@ -6,8 +6,10 @@
use std::sync::Arc;
use tokio::sync::mpsc;
use pmoupnp::actions::{ActionData, ActionError, ActionHandler};
use pmodidl::DIDLLite;
use pmoupnp::actions::{ActionData, ActionError, ActionHandler, get_value};
use pmoupnp::{get, set};
use pmoutils::ToXmlElement;
use crate::messages::{CommandParams, PlaybackState, ServerMessage, TransportAction};
use crate::state::SharedState;
@@ -118,7 +120,13 @@ pub fn set_uri_handler(
let state = state.clone();
Box::pin(async move {
let uri: String = get!(&data, "CurrentURI", String);
let metadata: String = get!(&data, "CurrentURIMetaData", String);
// CurrentURIMetaData peut être String ou DIDLLite (après parsing)
let metadata: String = get_value::<String>(&data, "CurrentURIMetaData")
.or_else(|_| {
get_value::<DIDLLite>(&data, "CurrentURIMetaData")
.map(|didl| didl.to_xml())
})
.unwrap_or_default();
let _ = ws.send(ServerMessage::Command {
action: TransportAction::SetUri,
params: Some(CommandParams {

View File

@@ -48,6 +48,23 @@ pub enum FactoryError {
VariableError(String),
}
/// Extrait un nom de navigateur court depuis un User-Agent complet.
fn extract_browser_name(ua: &str) -> &str {
if ua.contains("Edg/") || ua.contains("EdgA/") {
"Edge"
} else if ua.contains("OPR/") || ua.contains("Opera") {
"Opera"
} else if ua.contains("Chrome/") {
"Chrome"
} else if ua.contains("Firefox/") {
"Firefox"
} else if ua.contains("Safari/") {
"Safari"
} else {
"Browser"
}
}
/// Factory pour créer des Device UPnP WebRenderer avec des handlers WebSocket
pub struct WebRendererFactory;
@@ -65,10 +82,11 @@ impl WebRendererFactory {
let renderingcontrol = Self::build_renderingcontrol(ws_sender.clone(), state.clone())?;
let connectionmanager = Self::build_connectionmanager()?;
let short_name = extract_browser_name(browser_name);
let device = Device::new(
format!("WebRenderer"),
"WebRenderer".to_string(),
"MediaRenderer".to_string(),
format!("Web Audio - {}", browser_name),
format!("Web Audio {}", short_name),
);
device
.add_service(Arc::new(avtransport))

View File

@@ -40,11 +40,13 @@ pub async fn websocket_handler(
ws: WebSocketUpgrade,
State(state): State<WebSocketState>,
) -> impl IntoResponse {
tracing::info!("WebRenderer WebSocket upgrade request received");
ws.on_upgrade(move |socket: WebSocket| handle_socket(socket, state))
}
/// Gestion de la connexion WebSocket
async fn handle_socket(socket: WebSocket, state: WebSocketState) {
tracing::info!("WebRenderer WebSocket connection established");
let (mut sink, mut stream) = socket.split();
// Canal pour envoyer des messages au navigateur
@@ -68,11 +70,14 @@ async fn handle_socket(socket: WebSocket, state: WebSocketState) {
while let Some(msg_result) = stream.next().await {
match msg_result {
Ok(Message::Text(text)) => {
tracing::info!("WebRenderer received text message: {}", &text);
match serde_json::from_str::<ClientMessage>(&text) {
Ok(ClientMessage::Init { capabilities }) => {
tracing::info!("WebRenderer Init received, creating renderer...");
// Créer le renderer UPnP pour ce navigateur
match create_renderer_for_browser(&capabilities, tx.clone(), &state).await {
Ok(session) => {
tracing::info!("WebRenderer create_renderer_for_browser OK");
let token = session.token.clone();
let udn = session.udn.clone();
@@ -106,7 +111,7 @@ async fn handle_socket(socket: WebSocket, state: WebSocketState) {
);
}
Err(e) => {
tracing::error!("Failed to create WebRenderer: {}", e);
tracing::error!("Failed to create WebRenderer: {:?}", e);
break;
}
}
@@ -204,12 +209,14 @@ async fn create_renderer_for_browser(
let token = Uuid::new_v4().to_string();
// Construire le Device model avec les handlers WS
tracing::info!("WebRenderer: creating device model...");
let device = WebRendererFactory::create_device(
&capabilities.user_agent,
ws_sender.clone(),
shared_state.clone(),
)
.map_err(|e| crate::error::WebRendererError::DeviceCreationError(e.to_string()))?;
tracing::info!("WebRenderer: device model created");
let device = Arc::new(device);
@@ -218,19 +225,24 @@ async fn create_renderer_for_browser(
let di = {
use pmoupnp::UpnpServerExt;
tracing::info!("WebRenderer: getting server arc...");
let server_arc =
pmoserver::get_server().ok_or(crate::error::WebRendererError::ServerNotAvailable)?;
tracing::info!("WebRenderer: got server arc, acquiring write lock...");
let di = {
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
register_with_control_point(&di, ws_state)?;
tracing::info!("WebRenderer: registered with ControlPoint");
di
};