Merge pull request 'push-mtksvvvquolu' (#74) from push-mtksvvvquolu into main
Some checks failed
Build and Push Docker Image / build (push) Failing after 12m57s

Reviewed-on: #74
This commit was merged in pull request #74.
This commit is contained in:
2026-02-21 20:55:23 +01:00
38 changed files with 3322 additions and 126 deletions

View File

@@ -0,0 +1,97 @@
# WebRenderer UPnP privé par navigateur
## Vue d'ensemble
Le crate `pmowebrenderer` transforme chaque navigateur connecté en un **MediaRenderer UPnP privé**. Quand un navigateur se connecte via WebSocket, le backend Rust crée dynamiquement un device UPnP dédié. Le ControlPoint envoie des commandes SOAP à ce device, et les action handlers les relaient au navigateur via WebSocket. Le navigateur joue l'audio via `<audio>` et renvoie l'état au backend.
## Architecture
```
Browser (Vue.js) Rust Backend ControlPoint
| | |
|-- WS connect --------------->| |
|<-- SessionCreated (token) ---| |
|-- Init (capabilities) ------>| |
| |-- register_device() -------->| (Server)
| | (Device + Services custom) |
| |-- push_renderer() ---------->| (CP registry)
| | |
| |<-- SOAP Play (control_handler)
|<-- Command(Play, uri) -------| (action handler -> WS) |
|-- StateUpdate(Playing) ----->| |
| |-- update StateVarInstance -->| (evented -> SSE)
| | |
|-- WS disconnect ------------>| |
| |-- device_says_byebye() ----->| (CP registry)
```
## Flux de connexion
1. Le navigateur ouvre une WebSocket vers `/api/webrenderer/ws`
2. Il envoie un message `Init` avec ses capabilities (user_agent, formats supportes)
3. Le backend construit un `Device` UPnP avec des `Service` models custom :
- AVTransport (Play, Stop, Pause, Seek, SetURI, GetPositionInfo, etc.)
- RenderingControl (SetVolume, GetVolume, SetMute, GetMute)
- ConnectionManager (GetProtocolInfo)
4. Chaque Action a un handler qui capture le `mpsc::UnboundedSender<ServerMessage>` du WS
5. Le device est enregistre via `Server::register_device()` (routes SOAP + DEVICE_REGISTRY)
6. Un `RendererInfo` est pousse dans le `DeviceRegistry` du ControlPoint via `push_renderer()`
7. Le backend renvoie un `SessionCreated` avec le token et les infos du renderer
## Decision cle : Services dynamiques (zero changement pmoupnp)
Plutot que de modifier pmoupnp pour permettre l'override de handlers post-creation, on **construit des `Service` models dynamiques** pour chaque session WebSocket :
- Les `StateVariable` statics de pmomediarenderer sont reutilisees via `Arc::clone(&*VAR)`
- De nouvelles `Action` sont creees avec `Action::new()`, configurees avec `set_handler()` puis wrappees en `Arc`
- Les handlers capturent le sender WS et le `SharedState` (clone a chaque appel via `Fn` closure)
- Le `Device` est construit avec ces services custom, puis enregistre normalement
Cela reutilise toute l'infrastructure existante sans modification de pmoupnp ni pmomediarenderer.
## Propagation d'etat bidirectionnelle
### SOAP -> Navigateur (commandes)
Les action handlers des services AVTransport/RenderingControl :
1. Lisent les arguments SOAP depuis `ActionData` via la macro `get!()`
2. Envoient un `ServerMessage::Command` ou `SetVolume`/`SetMute` via le canal mpsc
3. Mettent a jour le `SharedState` local
4. Retournent les arguments OUT via `set!()` si necessaire
### Navigateur -> UPnP (etats)
Quand le navigateur envoie `StateUpdate`, `PositionUpdate`, `MetadataUpdate` ou `VolumeUpdate` :
1. Le `SharedState` est mis a jour
2. Les `StateVarInstance` du `DeviceInstance` sont mises a jour via `set_value(StateValue::...)`
3. Les variables evented declenchent les notifications UPnP captees par le watcher du ControlPoint
## Cycle de vie
- **Connexion** : creation du Device, enregistrement aupres du Server et du ControlPoint
- **Session active** : le `SessionManager` gere un timeout de 30 minutes d'inactivite
- **Deconnexion WS** : appel a `device_says_byebye()` sur le registry du ControlPoint pour marquer offline
- **Cleanup automatique** : le `SessionManager` verifie toutes les 60 secondes les sessions expirees
- **Pas de SSDP** : les WebRenderers sont injectes directement, max_age de 86400s
## Structure des fichiers
| Fichier | Role |
|---------|------|
| `handlers.rs` | Action handlers SOAP->WS (play, stop, pause, seek, set_uri, get_*_info, volume, mute) |
| `renderer.rs` | `WebRendererFactory` : construction dynamique Device/Services avec handlers |
| `websocket.rs` | Handler WS : connexion, reception messages, creation device, propagation etat |
| `session.rs` | `SessionManager` : gestion des sessions avec timeout |
| `config.rs` | `WebRendererExt` trait pour `pmoserver::Server` (enregistrement route WS) |
| `messages.rs` | Types de messages WS (ServerMessage, ClientMessage, etc.) |
| `state.rs` | `RendererState` et `SharedState` (etat partage entre handlers et WS) |
| `error.rs` | Types d'erreur du crate |
| `lib.rs` | Exports publics |
## Points d'attention
1. **parking_lot::RwLockWriteGuard non-Send** : les guards de `SharedState` doivent etre dropes avant tout `.await` dans les handlers async. Utiliser des blocs `{ ... }` pour limiter la portee.
2. **Fn vs FnOnce** : les `ActionHandler` sont `Fn` (appeles plusieurs fois). Les closures doivent cloner `ws` et `state` a chaque appel, avant le `async move`.
3. **Routes Axum persistantes** : Axum ne supporte pas la suppression de routes. Les routes SOAP d'un device deconnecte persistent mais les handlers retournent des erreurs naturellement.
4. **Acces au Server** : le WebSocket handler utilise `pmoserver::get_server()` (singleton global) pour enregistrer les devices dynamiquement.

View File

@@ -0,0 +1,138 @@
Parfait. Voici un **schéma fonctionnel minimal** pour un **MediaRenderer UPnP privé par navigateur** avec **token**. Lidée est de rester fidèle à ton backend Rust existant et à la webapp Vue.js. En s'appuyant sur l'architecture de PMOMusic, j'aimerais que tu proposes un plan détaillé pour implémenter un tel système de Média Renderer.
- L'application web se trouve dans: [@webapp](file:///Users/coissac/Sync/maison/Petite_maisons/src/pmomusic/pmoapp/webapp)
- Tu as un prototype de Média Renderer dans: [@pmomediarenderer](file:///Users/coissac/Sync/maison/Petite_maisons/src/pmomusic/pmomediarenderer)
- Le contrôle point est dans : [@pmocontrol](file:///Users/coissac/Sync/maison/Petite_maisons/src/pmomusic/pmocontrol)
- Tu implémenteras ce nouveau système de Média Renderer dans la CRATe pmowebrenderer
Tu mettras une version du plan en Markdown dans le répertoire [@Architecture](file:///Users/coissac/Sync/maison/Petite_maisons/src/pmomusic/Blackboard/Architecture) .
---
## 1. Flow général
```
Browser (Vue.js Control Point)
┌───────────────┐
│ UI / audio │
│ WebSocket │
└───────▲───────┘
│ token
Rust backend (UPnP MediaRenderer)
┌───────────────────────────┐
│ Token → Renderer mapping │
│ Device XML / SOAP endpoints│
│ Play/Pause/Stop → WS → Browser │
└───────────────────────────┘
```
---
## 2. Étapes détaillées
### a) Création du renderer
1. Le navigateur se connecte via WebSocket ou HTTP.
2. Rust génère un token unique pour ce client :
```rust
use uuid::Uuid;
let token = Uuid::new_v4().to_string();
```
3. Rust crée une instance MediaRenderer **privée**, associée à ce token :
* Device description XML : `/renderer/<token>/desc.xml`
* AVTransport SOAP : `/renderer/<token>/avtransport`
* RenderingControl SOAP : `/renderer/<token>/renderingcontrol`
---
### b) Control Point
* La webapp Vue.js reçoit le token et la “déclare” au Control Point :
```js
const renderer = {
token: "abcd-1234-efgh",
name: "Browser Renderer"
};
// Ajout au control point local
controlPoint.addRenderer(renderer);
```
* Toutes les commandes Play/Pause/Stop incluent ce token :
```js
ws.send(JSON.stringify({
token: renderer.token,
action: "play",
uri: "http://localhost:8080/media.mp3"
}));
```
---
### c) Backend Rust : dispatcher les commandes
* Rust reçoit le JSON avec le token.
* Vérifie que le token correspond à un renderer actif.
* Transmet la commande au navigateur via WebSocket (ou HTTP push) :
```rust
match msg.action.as_str() {
"play" => send_ws_to_browser(&token, format!("play:{}", msg.uri)),
"pause" => send_ws_to_browser(&token, "pause".to_string()),
"stop" => send_ws_to_browser(&token, "stop".to_string()),
_ => (),
}
```
* Rust met à jour létat du renderer (AVTransport/RenderingControl) pour le Control Point.
---
### d) Lecture côté navigateur
* Le navigateur reçoit la commande via WebSocket et pilote `<audio>` :
```js
ws.onmessage = (evt) => {
const msg = evt.data;
if(msg.startsWith("play:")) {
audio.src = msg.split(":")[1];
audio.play();
} else if(msg === "pause") {
audio.pause();
} else if(msg === "stop") {
audio.pause();
audio.currentTime = 0;
}
};
```
---
### e) Fermeture / cleanup
* Quand le navigateur se déconnecte :
* Rust supprime le renderer associé au token
* Émet un **byebye virtuel** pour le Control Point (si nécessaire)
* Libère toutes les ressources
---
## 3. Points clés
1. **Token unique** = session privée + sécurité
2. **Pas besoin de SSDP / annonce** : le renderer est dédié à un navigateur connu
3. **Control Point Vue.js** sait exactement quel renderer utiliser
4. **Rust backend** reste seul responsable de limplémentation UPnP
5. **Lecture réelle** = navigateur via `<audio>` ou `<video>`
---
💡 Ce modèle est très proche de ce que font **Mopidy avec Iris**, **Kodi Remote**, ou **Chromecast / local cast** : le device est connu et dédié, pas besoin de découverte réseau.

134
Cargo.lock generated
View File

@@ -19,6 +19,7 @@ dependencies = [
"pmoserver",
"pmosource",
"pmoupnp",
"pmowebrenderer",
"serde_json",
"tokio",
"tracing",
@@ -516,6 +517,7 @@ dependencies = [
"tower 0.5.2",
"tower-layer",
"tower-service",
"tracing",
]
[[package]]
@@ -526,6 +528,7 @@ checksum = "5b098575ebe77cb6d14fc7f32749631a6e44edbef6b796f89b020e99ba20d425"
dependencies = [
"axum-core 0.5.5",
"axum-macros",
"base64 0.22.1",
"bytes",
"form_urlencoded",
"futures-util",
@@ -544,8 +547,10 @@ dependencies = [
"serde_json",
"serde_path_to_error",
"serde_urlencoded",
"sha1",
"sync_wrapper",
"tokio",
"tokio-tungstenite",
"tower 0.5.2",
"tower-layer",
"tower-service",
@@ -570,6 +575,7 @@ dependencies = [
"sync_wrapper",
"tower-layer",
"tower-service",
"tracing",
]
[[package]]
@@ -605,6 +611,30 @@ dependencies = [
"tower-service",
]
[[package]]
name = "axum-extra"
version = "0.9.6"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c794b30c904f0a1c2fb7740f7df7f7972dfaa14ef6f57cb6178dc63e5dca2f04"
dependencies = [
"axum 0.7.9",
"axum-core 0.4.5",
"bytes",
"fastrand",
"futures-util",
"headers",
"http",
"http-body",
"http-body-util",
"mime",
"multer",
"pin-project-lite",
"serde",
"tower 0.5.2",
"tower-layer",
"tower-service",
]
[[package]]
name = "axum-macros"
version = "0.5.0"
@@ -671,7 +701,7 @@ dependencies = [
"portable-atomic",
"portable-atomic-util",
"serde",
"spin",
"spin 0.10.0",
"wasm-bindgen",
"wasm-bindgen-futures",
]
@@ -2089,6 +2119,30 @@ dependencies = [
"num-traits",
]
[[package]]
name = "headers"
version = "0.4.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b3314d5adb5d94bcdf56771f2e50dbbc80bb4bdf88967526706205ac9eff24eb"
dependencies = [
"base64 0.22.1",
"bytes",
"headers-core",
"http",
"httpdate",
"mime",
"sha1",
]
[[package]]
name = "headers-core"
version = "0.3.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "54b4a22553d4242c49fddb9ba998a99962b5cc6f22cb5a3482bec22522403ce4"
dependencies = [
"http",
]
[[package]]
name = "heapless"
version = "0.8.0"
@@ -3041,6 +3095,23 @@ dependencies = [
"pxfm",
]
[[package]]
name = "multer"
version = "3.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "83e87776546dc87511aa5ee218730c92b666d7264ab6ed41f9d215af9cd5224b"
dependencies = [
"bytes",
"encoding_rs",
"futures-util",
"http",
"httparse",
"memchr",
"mime",
"spin 0.9.8",
"version_check",
]
[[package]]
name = "native-tls"
version = "0.2.14"
@@ -4194,6 +4265,7 @@ dependencies = [
"tokio",
"tokio-stream",
"tokio-util",
"tower 0.5.2",
"tracing",
"tracing-subscriber",
"utoipa",
@@ -4280,6 +4352,31 @@ dependencies = [
"xmltree 0.10.3",
]
[[package]]
name = "pmowebrenderer"
version = "0.1.0"
dependencies = [
"async-trait",
"axum 0.8.7",
"axum-extra",
"futures",
"parking_lot",
"pmoconfig",
"pmocontrol",
"pmodidl",
"pmomediarenderer",
"pmoserver",
"pmoupnp",
"pmoutils",
"serde",
"serde_json",
"thiserror 2.0.17",
"tokio",
"tower-http",
"tracing",
"uuid",
]
[[package]]
name = "png"
version = "0.18.0"
@@ -5394,6 +5491,12 @@ dependencies = [
"libsoxr-sys",
]
[[package]]
name = "spin"
version = "0.9.8"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "6980e8d7511241f8acf4aebddbb1ff938df5eebe98691418c4468d0b72a96a67"
[[package]]
name = "spin"
version = "0.10.0"
@@ -6005,6 +6108,18 @@ dependencies = [
"tokio-stream",
]
[[package]]
name = "tokio-tungstenite"
version = "0.28.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d25a406cddcc431a75d3d9afc6a7c0f7428d4891dd973e4d54c56b46127bf857"
dependencies = [
"futures-util",
"log",
"tokio",
"tungstenite",
]
[[package]]
name = "tokio-util"
version = "0.7.17"
@@ -6267,6 +6382,23 @@ version = "0.2.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b"
[[package]]
name = "tungstenite"
version = "0.28.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8628dcc84e5a09eb3d8423d6cb682965dea9133204e8fb3efee74c2a0c259442"
dependencies = [
"bytes",
"data-encoding",
"http",
"httparse",
"log",
"rand 0.9.2",
"sha1",
"thiserror 2.0.17",
"utf-8",
]
[[package]]
name = "typeid"
version = "1.0.3"

View File

@@ -4,6 +4,7 @@ members = [
"PMOMusic",
"pmoupnp",
"pmomediarenderer",
"pmowebrenderer",
"pmomediaserver",
"pmoconfig",
"pmoutils",
@@ -37,7 +38,7 @@ async-trait = "0.1"
# Error handling
anyhow = "1.0"
thiserror = "2.0" # ⚠️ Unifier sur 2.0 (vous avez 1.0 et 2.0)
thiserror = "2.0"
# Logging
tracing = "0.1.41"
@@ -46,7 +47,7 @@ tracing-subscriber = { version = "0.3", features = ["fmt", "env-filter"] }
# HTTP/XML
reqwest = { version = "0.12", default-features = false }
ureq = "3.1"
quick-xml = { version = "0.38", features = ["serialize"] } # ⚠️ Unifier 0.37→0.38
quick-xml = { version = "0.38", features = ["serialize"] }
axum = "0.8.4"
futures = "0.3"
@@ -57,4 +58,4 @@ crossbeam-channel = "0.5"
rand = "0.9"
# Testing
tokio-test = "0.4"
tokio-test = "0.4"

View File

@@ -161,15 +161,17 @@ watch:
ci: fmt-check clippy test doc-build webapp
@echo "$(GREEN)✓ Toutes les vérifications CI passées$(NC)"
## run: Lance le binaire en mode debug
## run: Compile en debug et lance via Terminal.app (requis pour le multicast sur macOS Sequoia+)
run: debug
@echo "$(YELLOW)→ Lancement de l'application...$(NC)"
./target/debug/$(BINARY_NAME)
@echo "$(YELLOW)→ Lancement de l'application via Terminal.app...$(NC)"
@echo "$(BLUE) (Terminal.app est nécessaire pour le multicast sur macOS Sequoia+)$(NC)"
@osascript -e 'tell application "Terminal" to do script "cd \"$(CURDIR)\" && ./target/debug/$(BINARY_NAME) 2>&1 | tee pmomusic.log; exit"'
## run-release: Lance le binaire en mode release
## run-release: Compile en release et lance via Terminal.app (requis pour le multicast sur macOS Sequoia+)
run-release: release
@echo "$(YELLOW)→ Lancement de l'application (release)...$(NC)"
./$(RUST_TARGET)/$(BINARY_NAME)
@echo "$(YELLOW)→ Lancement de l'application (release) via Terminal.app...$(NC)"
@echo "$(BLUE) (Terminal.app est nécessaire pour le multicast sur macOS Sequoia+)$(NC)"
@osascript -e 'tell application "Terminal" to do script "cd \"$(CURDIR)\" && ./$(RUST_TARGET)/$(BINARY_NAME) 2>&1 | tee pmomusic.log; exit"'
## size: Affiche la taille du binaire
size:

View File

@@ -1,6 +1,6 @@
[package]
name = "PMOMusic"
version = "0.3.20"
version = "0.3.21"
edition = "2024"
[dependencies]
@@ -15,6 +15,7 @@ pmoaudiocache = { path = "../pmoaudiocache", features = ["pmoserver"]}
pmoaudio-ext = { path = "../pmoaudio-ext", features = ["all"] }
pmoapp = { path = "../pmoapp", features = ["pmoserver"] }
pmocontrol = { path = "../pmocontrol", features = ["pmoserver"] }
pmowebrenderer = { path = "../pmowebrenderer", features = ["pmoserver"] }
tokio = { workspace = true, features = ["rt-multi-thread", "macros", "sync", "time", "signal"] }
tracing = { workspace = true }

View File

@@ -7,6 +7,7 @@ use pmomediaserver::{
use pmoserver::Server;
use pmosource::MusicSourceExt;
use pmoupnp::UpnpServerExt;
use pmowebrenderer::WebRendererExt;
use tracing::info;
#[tokio::main]
@@ -104,13 +105,22 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
// Enregistrer le Control Point (découverte renderers/serveurs + API REST + SSE)
info!("🎛️ Registering Control Point...");
let _control_point = server
let control_point = server
.write()
.await
.register_control_point(5)
.await
.expect("Failed to register Control Point");
// Enregistrer le WebRenderer (endpoint WebSocket pour renderers navigateur)
info!("🌐 Registering WebRenderer...");
server
.write()
.await
.register_web_renderer(control_point)
.await
.expect("Failed to register WebRenderer");
// Ajouter la webapp via le trait WebAppExt
info!("📡 Registering Web application...");
server
@@ -127,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

@@ -100,14 +100,7 @@ function formatTime(ms: number | null | undefined): string {
}
const currentTime = computed(() => formatTime(state.value?.position_ms));
const totalTime = computed(() => {
const duration = state.value?.duration_ms;
const transport = state.value?.transport_state;
console.log(
`[CurrentTrack] rendererId=${props.rendererId}, duration_ms=${duration}, transport=${transport}, title=${state.value?.current_track?.title}`,
);
return formatTime(duration);
});
const totalTime = computed(() => formatTime(state.value?.duration_ms));
const hasCover = computed(
() => !!metadata.value?.album_art_uri && !imageError.value,
@@ -357,20 +350,14 @@ const swipeOpacity = computed(() => {
@click="openCoverOverlay"
>
<img
v-if="cacheBustedUrl"
ref="coverImageRef"
:style="{
opacity:
cacheBustedUrl && imageLoaded && !imageError ? 1 : 0,
visibility:
cacheBustedUrl && imageLoaded && !imageError
? 'visible'
: 'hidden',
position:
cacheBustedUrl && imageLoaded && !imageError
? 'relative'
: 'absolute',
opacity: imageLoaded && !imageError ? 1 : 0,
visibility: imageLoaded && !imageError ? 'visible' : 'hidden',
position: imageLoaded && !imageError ? 'relative' : 'absolute',
}"
:src="cacheBustedUrl || ''"
:src="cacheBustedUrl"
:alt="metadata?.album || 'Album cover'"
class="cover-image"
@load="handleImageLoad"

View File

@@ -65,9 +65,6 @@ export function useCoverImage(
if (retryCount.value < maxRetries) {
retryCount.value++;
console.log(
`[useCoverImage] Retrying image load (${retryCount.value}/${maxRetries}): ${currentUrl.value}`,
);
setTimeout(() => {
if (!currentUrl.value) return;
@@ -77,7 +74,6 @@ export function useCoverImage(
currentUrl.value,
retryCount.value,
);
console.log(`[useCoverImage] Retry URL: ${cacheBustedUrl.value}`);
}, retryDelay * retryCount.value); // Exponential backoff
} else {
console.error(
@@ -89,9 +85,6 @@ export function useCoverImage(
// Handle successful image load
function handleImageLoad() {
console.log(
`[useCoverImage] Image loaded successfully: ${currentUrl.value}`,
);
imageLoaded.value = true;
imageError.value = false;
retryCount.value = 0;
@@ -119,18 +112,11 @@ export function useCoverImage(
watch(
imageUrl,
(newUri, oldUri) => {
console.log(
`[useCoverImage] URL changed from "${oldUri}" to "${newUri}"`,
);
imageError.value = false;
retryCount.value = 0;
// Si c'est un changement d'URL (pas l'initialisation)
if (oldUri && newUri && oldUri !== newUri) {
console.log(
`[useCoverImage] Changing image, keeping old one visible during load`,
);
isLoadingNewImage.value = true;
// On garde imageLoaded à true pour garder l'ancienne image visible
} else if (!newUri) {
@@ -148,9 +134,6 @@ export function useCoverImage(
if (newUri) {
// Generate cache-busted URL
cacheBustedUrl.value = getCacheBustedUrl(newUri, 0);
console.log(
`[useCoverImage] New cache-busted URL: ${cacheBustedUrl.value}`,
);
} else {
cacheBustedUrl.value = null;
}

View File

@@ -37,9 +37,6 @@ function ensureSSEConnected() {
switch (event.type) {
case 'online':
// Nouveau serveur découvert
console.log(`[useMediaServers] Serveur ${serverId} (${event.friendly_name}) est maintenant en ligne`)
// Ajouter au cache avec les infos disponibles
const server: MediaServerSummary = {
id: serverId,
@@ -55,9 +52,6 @@ function ensureSSEConnected() {
break
case 'offline':
// Serveur déconnecté
console.log(`[useMediaServers] Serveur ${serverId} est maintenant hors ligne`)
// Marquer comme offline dans le cache
const existingServer = serversCache.value.get(serverId)
if (existingServer) {
@@ -77,7 +71,6 @@ function ensureSSEConnected() {
case 'global_updated':
// Invalider tout le cache de ce serveur
console.log(`[useMediaServers] GlobalUpdated pour ${serverId}`)
const globalKeysToDelete: string[] = []
browseCache.value.forEach((_, key) => {
if (key.startsWith(serverId + '/')) {
@@ -89,7 +82,6 @@ function ensureSSEConnected() {
case 'containers_updated':
// Invalider les containers spécifiques
console.log(`[useMediaServers] ContainersUpdated pour ${serverId}:`, event.container_ids)
event.container_ids.forEach(containerId => {
const key = `${serverId}/${containerId}`
browseCache.value.delete(key)

View File

@@ -48,11 +48,6 @@ function ensureSSEConnected() {
// Gérer les événements Online/Offline différemment
if (event.type === "online") {
// Nouveau renderer découvert
console.log(
`[useRenderers] Renderer ${rendererId} (${event.friendly_name}) est maintenant en ligne`,
);
// Ajouter au cache avec les infos disponibles
// Note: on n'a pas toutes les infos (capabilities, protocol) donc on fetch ensuite
const renderer: RendererSummary = {
@@ -86,11 +81,6 @@ function ensureSSEConnected() {
}
if (event.type === "offline") {
// Renderer déconnecté
console.log(
`[useRenderers] Renderer ${rendererId} est maintenant hors ligne`,
);
// Marquer comme offline dans le cache
const renderer = renderersCache.value.get(rendererId);
if (renderer) {
@@ -109,9 +99,6 @@ function ensureSSEConnected() {
snapshotState.lastEventAt.set(rendererId, timestamp);
const snapshot = snapshotState.snapshots.get(rendererId);
console.log(
`[SSE Event] type=${event.type}, rendererId=${rendererId}, snapshot exists=${!!snapshot}, transport=${snapshot?.state?.transport_state}`,
);
// Si pas de snapshot, on doit fetch
if (!snapshot) {
@@ -126,11 +113,6 @@ function ensureSSEConnected() {
break;
case "position_changed":
// Debug: log pour tracer les oscillations
console.log(
`[position_changed] renderer=${rendererId}, duration=${event.track_duration}`,
);
// Mettre à jour position et durée de manière atomique pour garantir la cohérence
// Le backend envoie TOUJOURS les deux valeurs (même si null)
@@ -177,9 +159,6 @@ function ensureSSEConnected() {
...snapshot,
state: newState,
};
console.log(
`[position_changed] Setting new snapshot for ${rendererId}, state ref=${Object.prototype.toString.call(newState)}, transport=${newState.transport_state}`,
);
snapshotState.snapshots.set(rendererId, newSnapshot);
break;

View File

@@ -0,0 +1,600 @@
/**
* 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.
* L'audio est géré via deux éléments <audio> en ping-pong connectés à un
* AudioContext (MediaElementSourceNode) pour le contrôle du volume/mute.
*
* ## Flux gapless
* 1. SetAVTransportURI(N) → slotA.src = N, slotA.load()
* 2. Play → slotA.play()
* 3. SetNextAVTransportURI(N+1) → slotB.src = N+1, slotB.load() (préchargement)
* 4. slotA "ended" → currentSlot = B, slotB.play() immédiat
* → TrackEnded envoyé au backend
* 5. Cycle recommence depuis 3 (slotA redevient le "next")
*/
import { ref, onMounted, onUnmounted, readonly } from "vue";
// ─── Types (miroir de messages.rs) ────────────────────────────────────────────
interface BrowserCapabilities {
instance_id: string;
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" | "set_next_uri";
interface CommandParams {
uri?: string;
metadata?: 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: "session_created"; token: string; renderer_info: RendererInfo }
| StateSyncMessage
| { 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: "track_ended" }
| { 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 ─────────────────────────────────────────
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);
}
// ─── Moteur Audio (HTMLAudioElement + MediaElementSourceNode) ─────────────────
class GaplessEngine {
/** Les deux éléments <audio> fixes (ping-pong) */
private readonly slots: [HTMLAudioElement, HTMLAudioElement];
/** Index du slot actuellement en lecture (0 ou 1) */
private currentSlot: 0 | 1 = 0;
/** URI préchargée dans le slot "next" (l'autre) */
private nextUri: string | null = null;
private volume = 1.0;
private muted = false;
/** Durée de la piste courante (secondes), lue via loadedmetadata */
private _duration = 0;
/**
* Indique qu'un play() a été reçu avant set_uri (race condition).
* setCurrent() le détectera et lancera la lecture automatiquement.
*/
private playPending = false;
onStateChange: (state: PlaybackState) => void = () => {};
onPosition: (pos: number, dur: number) => void = () => {};
onTrackEnded: () => void = () => {};
private positionInterval: ReturnType<typeof setInterval> | null = null;
constructor() {
this.slots = [
this.makeAudioElement(),
this.makeAudioElement(),
];
}
// ── Volume / Mute ────────────────────────────────────────────────────────
setVolume(v: number) {
this.volume = v;
for (const el of this.slots) {
el.volume = this.muted ? 0 : v;
}
}
setMute(m: boolean) {
this.muted = m;
for (const el of this.slots) {
el.volume = m ? 0 : this.volume;
}
}
// ── Transport ────────────────────────────────────────────────────────────
/** Charge la piste courante (sans la jouer). */
setCurrent(uri: string): void {
this.nextUri = null;
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];
// Retirer l'écouteur "ended" de l'autre slot si présent
const otherSlot = (1 - this.currentSlot) as 0 | 1;
this.slots[otherSlot].onended = null;
this.slots[otherSlot].pause();
el.onended = null;
el.pause();
el.src = uri;
el.load();
// Récupérer la durée dès que les métadonnées sont disponibles
el.onloadedmetadata = () => {
this._duration = el.duration || 0;
};
}
/** Précharge la piste suivante dans l'autre slot. */
setNext(uri: string): void {
this.nextUri = uri;
const nextSlot = (1 - this.currentSlot) as 0 | 1;
const el = this.slots[nextSlot];
el.src = uri;
el.preload = "auto";
el.load();
}
async play(): Promise<void> {
const el = this.slots[this.currentSlot];
// Si pas de source : play() est arrivé avant set_uri (race condition).
// On mémorise et setCurrent() déclenchera la lecture dès qu'il sera appelé.
if (!el.src || el.src === window.location.href) {
this.playPending = true;
return;
}
this.playPending = false;
el.onended = () => this.onCurrentEnded();
// Si l'élément n'a pas encore de données, attendre canplay.
if (el.readyState < HTMLMediaElement.HAVE_FUTURE_DATA) {
await new Promise<void>((resolve) => {
const onCanPlay = () => {
el.removeEventListener("canplay", onCanPlay);
resolve();
};
el.addEventListener("canplay", onCanPlay);
});
}
try {
await el.play();
} catch (e) {
console.warn("[GaplessEngine] play() failed (autoplay blocked?):", e);
throw e;
}
this.startPositionTimer();
this.onStateChange("PLAYING");
}
pause(): void {
const el = this.slots[this.currentSlot];
el.pause();
this.stopPositionTimer();
this.onStateChange("PAUSED");
this.sendPosition();
}
stop(): void {
this.playPending = false;
const el = this.slots[this.currentSlot];
el.onended = null;
el.pause();
el.currentTime = 0;
this.nextUri = null;
this.stopPositionTimer();
this.onStateChange("STOPPED");
}
seek(toSeconds: number): void {
const el = this.slots[this.currentSlot];
el.currentTime = toSeconds;
}
destroy(): void {
this.stopPositionTimer();
for (const el of this.slots) {
el.onended = null;
el.pause();
el.src = "";
}
}
// ── Privé ────────────────────────────────────────────────────────────────
private makeAudioElement(): HTMLAudioElement {
const el = document.createElement("audio");
el.preload = "auto";
return el;
}
private onCurrentEnded(): void {
const nextSlot = (1 - this.currentSlot) as 0 | 1;
if (this.nextUri !== null) {
// Le slot suivant est préchargé, on bascule
this.currentSlot = nextSlot;
this.nextUri = null;
this._duration = 0;
const nextEl = this.slots[this.currentSlot];
nextEl.onended = () => this.onCurrentEnded();
// Récupérer la durée de la nouvelle piste courante
if (nextEl.duration && isFinite(nextEl.duration)) {
this._duration = nextEl.duration;
} else {
nextEl.onloadedmetadata = () => {
this._duration = nextEl.duration || 0;
};
}
// Informer le backend (qui fera le swap current←next côté serveur)
this.onTrackEnded();
// Démarrer immédiatement (le préchargement a eu lieu)
nextEl.play().catch((e) =>
console.error("[GaplessEngine] next.play() failed:", e),
);
// L'état reste PLAYING
} else {
// Pas de suivant : fin de lecture
this.stopPositionTimer();
this.onStateChange("STOPPED");
this.onTrackEnded();
}
}
private startPositionTimer(): void {
if (this.positionInterval !== null) return;
this.positionInterval = setInterval(() => this.sendPosition(), 1000);
}
private stopPositionTimer(): void {
if (this.positionInterval !== null) {
clearInterval(this.positionInterval);
this.positionInterval = null;
}
}
private sendPosition(): void {
const el = this.slots[this.currentSlot];
const pos = el.currentTime || 0;
const dur = (el.duration && isFinite(el.duration)) ? el.duration : this._duration;
this.onPosition(pos, dur);
}
}
// ─── Composable ───────────────────────────────────────────────────────────────
export function useWebRenderer() {
const connected = ref(false);
const rendererInfo = ref<RendererInfo | null>(null);
let ws: WebSocket | null = null;
let engine: GaplessEngine | 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));
}
}
// ── Initialisation du moteur ──────────────────────────────────────────────
function initEngine(): GaplessEngine {
const e = new GaplessEngine();
e.onStateChange = (state) => {
send({ type: "state_update", state });
};
e.onPosition = (pos, dur) => {
send({
type: "position_update",
position: secondsToUpnpTime(pos),
duration: secondsToUpnpTime(dur),
});
};
e.onTrackEnded = () => {
send({ type: "track_ended" });
};
return e;
}
// ── Exécution des commandes UPnP ──────────────────────────────────────────
async function execCommand(action: TransportAction, params?: CommandParams) {
if (!engine) return;
switch (action) {
case "set_uri":
if (params?.uri) {
engine.setCurrent(params.uri);
}
break;
case "set_next_uri":
if (params?.uri) {
engine.setNext(params.uri);
}
break;
case "play":
await engine.play().catch((e) =>
console.warn("[WebRenderer] play() failed:", e),
);
break;
case "pause":
engine.pause();
break;
case "stop":
engine.stop();
break;
case "seek":
if (params?.position) {
engine.seek(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;
onConnectedCallback?.();
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":
void execCommand(msg.action, msg.params);
break;
case "set_volume":
engine?.setVolume(msg.volume / 100);
break;
case "set_mute":
engine?.setMute(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`;
ws = new WebSocket(url);
ws.onopen = () => {
send({
type: "init",
capabilities: {
instance_id: getOrCreateInstanceId(),
user_agent: navigator.userAgent,
supported_formats: getSupportedFormats(),
},
});
};
ws.onmessage = handleMessage;
ws.onclose = () => {
connected.value = false;
rendererInfo.value = null;
ws = null;
};
ws.onerror = (err) => {
console.error("[WebRenderer] Erreur WebSocket :", err);
};
}
// ── Déconnexion ───────────────────────────────────────────────────────────
function disconnect() {
engine?.destroy();
if (ws) {
ws.close(1000, "Page unloaded");
ws = null;
}
connected.value = false;
}
// ── Cycle de vie ──────────────────────────────────────────────────────────
onMounted(() => {
engine = initEngine();
connect();
window.addEventListener("beforeunload", disconnect);
});
onUnmounted(() => {
disconnect();
engine = 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

@@ -34,12 +34,7 @@ const uiStore = useUIStore();
sse.onConnectionChange((connected) => {
uiStore.setSSEConnected(connected);
if (connected) {
console.log("[App] SSE connecté");
}
});
// Démarrer la connexion SSE
sse.connect();
console.log("[App] PMOControl initialisé");

View File

@@ -31,13 +31,10 @@ export class PMOControlSSE {
return
}
console.log('[SSE] Connexion à /api/control/events...')
try {
this.eventSource = new EventSource('/api/control/events')
this.eventSource.onopen = () => {
console.log('[SSE] Connexion établie')
this.reconnectAttempts = 0
this.isConnected = true
this.notifyConnectionCallbacks(true)
@@ -75,7 +72,6 @@ export class PMOControlSSE {
}
if (this.eventSource) {
console.log('[SSE] Déconnexion')
this.eventSource.close()
this.eventSource = null
this.isConnected = false
@@ -99,8 +95,6 @@ export class PMOControlSSE {
this.maxReconnectDelay
)
console.log(`[SSE] Reconnexion dans ${delay / 1000}s (tentative ${this.reconnectAttempts})`)
this.reconnectTimer = window.setTimeout(() => {
this.reconnectTimer = null
this.connect()

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,
@@ -343,4 +368,5 @@ const currentTabProps = computed(() => {
opacity: 0;
transform: translateX(-20px);
}
</style>

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

@@ -230,8 +230,11 @@ async fn serve_file_with_streaming<C: CacheConfig + 'static>(
if let Some(generator) = param_generator {
if let Some(data) = generator(cache.clone(), pk.to_string(), param.to_string()).await {
// Le générateur a créé les données, les servir directement
let response =
(StatusCode::OK, [("content-type", content_type)], data).into_response();
let response = (
StatusCode::OK,
[("content-type", content_type), ("access-control-allow-origin", "*")],
data,
).into_response();
if response.status().is_success() {
cache.notify_broadcast(pk, &qualifier).await;
@@ -342,7 +345,8 @@ async fn stream_file_progressive(
let mut response_builder = axum::http::Response::builder()
.status(StatusCode::PARTIAL_CONTENT)
.header(header::CONTENT_TYPE, content_type)
.header(header::ACCEPT_RANGES, "bytes");
.header(header::ACCEPT_RANGES, "bytes")
.header("Access-Control-Allow-Origin", "*");
if let Some(total_size) = expected_size {
response_builder = response_builder.header(
@@ -370,7 +374,8 @@ async fn stream_file_progressive(
let mut response = axum::http::Response::builder()
.status(StatusCode::OK)
.header(header::CONTENT_TYPE, content_type)
.header(header::ACCEPT_RANGES, "bytes");
.header(header::ACCEPT_RANGES, "bytes")
.header("Access-Control-Allow-Origin", "*");
if let Some(size) = expected_size {
response = response.header(header::CONTENT_LENGTH, size);
@@ -414,6 +419,10 @@ async fn serve_complete_file(
parts
.headers
.insert(header::CONTENT_TYPE, content_type.parse().unwrap());
parts.headers.insert(
"Access-Control-Allow-Origin".parse::<axum::http::HeaderName>().unwrap(),
"*".parse().unwrap(),
);
// Convertir ServeFileSystemResponseBody en Body
axum::http::Response::from_parts(parts, Body::new(body))
@@ -515,9 +524,27 @@ pub fn create_file_router_with_generator<C: CacheConfig + 'static>(
let path_with_param = format!("/{}/{}/{{pk}}/{{param}}", cache_name, cache_type);
let path_without_param = format!("/{}/{}/{{pk}}", cache_name, cache_type);
async fn cors_preflight() -> impl IntoResponse {
(
StatusCode::NO_CONTENT,
[
("Access-Control-Allow-Origin", "*"),
("Access-Control-Allow-Methods", "GET, OPTIONS"),
("Access-Control-Allow-Headers", "Range, Content-Type"),
("Access-Control-Max-Age", "86400"),
],
)
}
Router::new()
.route(&path_without_param, get(get_file::<C>))
.route(&path_with_param, get(get_file_with_param::<C>))
.route(
&path_without_param,
get(get_file::<C>).options(cors_preflight),
)
.route(
&path_with_param,
get(get_file_with_param::<C>).options(cors_preflight),
)
.with_state((cache, content_type, param_generator))
}

View File

@@ -878,6 +878,10 @@ impl ControlPoint {
renderer.set_last_metadata(Some(metadata));
renderer.set_playback_source(PlaybackSource::FromQueue);
// Prefetch next track if supported (gapless playback)
self.prefetch_next_track(&renderer, renderer_id);
Ok(())
}
@@ -972,6 +976,51 @@ impl ControlPoint {
}
}
/// Prefetch the next track for a renderer (gapless playback support).
///
/// Called by the WebRenderer backend after a TrackEnded event to trigger
/// SetNextAVTransportURI for the next-next track (N+2).
pub fn prefetch_next_for_renderer(&self, renderer_id: &DeviceId) {
if let Some(renderer) = self.music_renderer_by_id(renderer_id) {
self.prefetch_next_track(&renderer, renderer_id);
}
}
/// Advances the queue index by one for a gapless transition, then prefetches the next-next track.
///
/// Called by the WebRenderer backend after a gapless TrackEnded event. In the gapless case,
/// the browser has already advanced to the next track autonomously, so the backend queue
/// index must be updated to stay in sync. Then we prefetch track N+2 for the next gapless
/// transition.
pub fn advance_queue_and_prefetch(&self, renderer_id: &DeviceId) {
if let Some(renderer) = self.music_renderer_by_id(renderer_id) {
// Advance the queue index to match the browser's gapless transition
match renderer.advance_queue_index() {
Ok(true) => {
debug!(
renderer = renderer_id.0.as_str(),
"advance_queue_and_prefetch: queue index advanced for gapless transition"
);
// Now prefetch the next-next track (N+2)
self.prefetch_next_track(&renderer, renderer_id);
}
Ok(false) => {
debug!(
renderer = renderer_id.0.as_str(),
"advance_queue_and_prefetch: no next track in queue, skipping prefetch"
);
}
Err(err) => {
debug!(
renderer = renderer_id.0.as_str(),
error = %err,
"advance_queue_and_prefetch: failed to advance queue index"
);
}
}
}
}
/// Jumps to a specific index in the queue and starts playback.
pub fn play_queue_index(
&self,
@@ -1019,6 +1068,10 @@ impl ControlPoint {
}
renderer.set_playback_source(PlaybackSource::FromQueue);
// Prefetch next track if supported (gapless playback)
self.prefetch_next_track(&renderer, renderer_id);
Ok(())
}

View File

@@ -879,6 +879,21 @@ impl MusicRenderer {
Ok(())
}
/// Advance the queue index by one without starting playback.
///
/// Used by the WebRenderer gapless path: the browser autonomously transitions
/// to the next track, so only the backend queue pointer needs to be updated
/// to stay in sync.
pub fn advance_queue_index(&self) -> Result<bool, ControlPointError> {
let mut backend = self.lock_backend_for("advance_queue_index");
let advanced = backend.advance()?;
drop(backend);
if advanced {
self.emit_queue_updated();
}
Ok(advanced)
}
/// Play from a specific index in the queue.
pub fn play_from_index(&self, index: usize) -> Result<(), ControlPointError> {
// Reset the has_played flag before starting playback to prevent

View File

@@ -1,4 +1,6 @@
use crate::avtransport::variables::{A_ARG_TYPE_INSTANCE_ID, TRANSPORTSTATE, TRANSPORTSTATUS};
use crate::avtransport::variables::{
A_ARG_TYPE_INSTANCE_ID, TRANSPORTPLAYSPEED, TRANSPORTSTATE, TRANSPORTSTATUS,
};
use pmoupnp::define_action;
define_action! {
@@ -6,5 +8,6 @@ define_action! {
in "InstanceID" => A_ARG_TYPE_INSTANCE_ID,
out "CurrentTransportState" => TRANSPORTSTATE,
out "CurrentTransportStatus" => TRANSPORTSTATUS,
out "CurrentSpeed" => TRANSPORTPLAYSPEED,
}
}

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

@@ -517,13 +517,11 @@ impl RadioParadiseSource {
}
// Normaliser l'albumArtURI : rendre absolu si chemin relatif, sinon fallback par défaut
if let Some(art) = item.album_art.as_mut() {
if art.starts_with('/') {
*art = format!("{}{}", self.base_url, art);
}
} else {
item.album_art = Some(self.default_cover_url());
}
item.album_art = match item.album_art.take() {
Some(art) if art.starts_with('/') => Some(format!("{}{}", self.base_url, art)),
Some(_) => Some(self.default_cover_url()), // URL externe brute → fallback local
None => Some(self.default_cover_url()),
};
}
Ok(items)
@@ -590,13 +588,11 @@ impl RadioParadiseSource {
item.genre = Some("Radio Paradise".to_string());
}
if let Some(art) = item.album_art.as_mut() {
if art.starts_with('/') {
*art = format!("{}{}", self.base_url, art);
}
} else {
item.album_art = Some(self.default_cover_url());
}
item.album_art = match item.album_art.take() {
Some(art) if art.starts_with('/') => Some(format!("{}{}", self.base_url, art)),
Some(_) => Some(self.default_cover_url()),
None => Some(self.default_cover_url()),
};
}
Ok(items)
@@ -921,6 +917,13 @@ impl MusicSource for RadioParadiseSource {
if item.genre.is_none() {
item.genre = Some("Radio Paradise".to_string());
}
item.album_art = match item.album_art.take() {
Some(art) if art.starts_with('/') => {
Some(format!("{}{}", self.base_url, art))
}
Some(_) => Some(self.default_cover_url()),
None => Some(self.default_cover_url()),
};
adjusted.push(item);
}
@@ -987,13 +990,13 @@ impl MusicSource for RadioParadiseSource {
item.genre = Some("Radio Paradise".to_string());
}
if let Some(art) = item.album_art.as_mut() {
if art.starts_with('/') {
*art = format!("{}{}", self.base_url, art);
item.album_art = match item.album_art.take() {
Some(art) if art.starts_with('/') => {
Some(format!("{}{}", self.base_url, art))
}
} else {
item.album_art = Some(self.default_cover_url());
}
Some(_) => Some(self.default_cover_url()),
None => Some(self.default_cover_url()),
};
let expected_id =
format!("radio-paradise:channel:{}:liveplaylist:track:{}", slug, pk);

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

@@ -22,7 +22,7 @@ base64 = "0.22.1"
thiserror = { workspace = true }
anyhow = { workspace = true }
xmltree = "0.11.0"
axum = "0.8.4"
axum = { version = "0.8.4", features = ["ws"] }
tokio = { workspace = true, features = ["rt-multi-thread", "macros", "sync"] }
serde = { workspace = true }
serde_json = { workspace = true }

View File

@@ -79,10 +79,11 @@ impl SsdpClient {
socket.set_read_timeout(Some(Duration::from_secs(1)))?;
socket.set_multicast_loop_v4(true)?; // utile en dev local
let multicast_addr = SSDP_MULTICAST_ADDR.parse().unwrap();
for iface in get_if_addrs::get_if_addrs()? {
if let std::net::IpAddr::V4(ipv4) = iface.ip() {
if !ipv4.is_loopback() {
match socket.join_multicast_v4(&SSDP_MULTICAST_ADDR.parse().unwrap(), &ipv4) {
match socket.join_multicast_v4(&multicast_addr, &ipv4) {
Ok(()) => {
debug!("SSDP: joined {} on {}", SSDP_MULTICAST_ADDR, ipv4);
}
@@ -97,6 +98,20 @@ impl SsdpClient {
}
}
// Sur macOS, chaque join_multicast_v4 écrase IP_MULTICAST_IF.
// On remet explicitement l'interface de sortie sur l'IP principale
// pour que send_to vers 239.255.255.250 utilise la bonne interface.
let local_ip: std::net::Ipv4Addr = pmoutils::guess_local_ip()
.parse()
.unwrap_or("0.0.0.0".parse().unwrap());
let socket2 = Socket::from(socket);
socket2.set_multicast_if_v4(&local_ip)?;
debug!(
"SSDP client: multicast outgoing interface set to {}",
local_ip
);
let socket: UdpSocket = socket2.into();
info!("✅ SSDP client ready on {}", addr);
Ok(Self {

View File

@@ -75,7 +75,7 @@ impl SsdpServer {
socket2.bind(&bind_addr.into())?;
// Convertir en UdpSocket standard
let socket: UdpSocket = socket2.into();
let mut socket: UdpSocket = socket2.into();
// Rejoindre le groupe multicast
socket.join_multicast_v4(
@@ -83,6 +83,22 @@ impl SsdpServer {
&"0.0.0.0".parse().unwrap(),
)?;
// Sur macOS, join_multicast_v4 peut positionner IP_MULTICAST_IF
// sur une interface bridge/VM. On remet explicitement l'interface
// de sortie sur l'IP principale.
let local_ip: std::net::Ipv4Addr = pmoutils::guess_local_ip()
.parse()
.unwrap_or("0.0.0.0".parse().unwrap());
{
let socket2 = Socket::from(socket);
socket2.set_multicast_if_v4(&local_ip)?;
debug!(
"SSDP server: multicast outgoing interface set to {}",
local_ip
);
socket = socket2.into();
}
socket.set_read_timeout(Some(Duration::from_secs(1)))?;
socket.set_multicast_loop_v4(false)?;

38
pmowebrenderer/Cargo.toml Normal file
View File

@@ -0,0 +1,38 @@
[package]
name = "pmowebrenderer"
version = "0.1.0"
edition = "2021"
[dependencies]
pmoupnp = { path = "../pmoupnp" }
pmomediarenderer = { path = "../pmomediarenderer" }
pmoserver = { path = "../pmoserver", optional = true }
pmocontrol = { path = "../pmocontrol", optional = true }
pmoconfig = { path = "../pmoconfig" }
# Async runtime
tokio = { workspace = true, features = ["full"] }
async-trait = { workspace = true }
# WebSocket
axum = { workspace = true, features = ["ws"] }
axum-extra = { version = "0.9", features = ["typed-header"] }
tower-http = { version = "0.6", features = ["fs", "trace"] }
futures = "0.3"
# 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"]

View File

@@ -0,0 +1,57 @@
//! Intégration avec pmoserver — enregistrement des routes WebRenderer
#[cfg(feature = "pmoserver")]
use std::sync::Arc;
#[cfg(feature = "pmoserver")]
use std::time::Duration;
#[cfg(feature = "pmoserver")]
use async_trait::async_trait;
#[cfg(feature = "pmoserver")]
use pmocontrol::ControlPoint;
#[cfg(feature = "pmoserver")]
use crate::error::WebRendererError;
#[cfg(feature = "pmoserver")]
use crate::session::SessionManager;
#[cfg(feature = "pmoserver")]
use crate::websocket::{websocket_handler, WebSocketState};
/// Trait pour étendre pmoserver::Server avec les routes WebRenderer
#[cfg(feature = "pmoserver")]
#[async_trait]
pub trait WebRendererExt {
async fn register_web_renderer(
&mut self,
control_point: Arc<ControlPoint>,
) -> Result<(), WebRendererError>;
}
#[cfg(feature = "pmoserver")]
#[async_trait]
impl WebRendererExt for pmoserver::Server {
async fn register_web_renderer(
&mut self,
control_point: Arc<ControlPoint>,
) -> Result<(), WebRendererError> {
let session_manager = Arc::new(SessionManager::new(Duration::from_secs(30 * 60)));
let ws_state = WebSocketState {
session_manager,
control_point,
};
self.add_any_handler_with_state("/api/webrenderer/ws", websocket_handler, ws_state)
.await;
tracing::info!("WebRenderer WebSocket endpoint registered at /api/webrenderer/ws");
Ok(())
}
}
#[cfg(not(feature = "pmoserver"))]
pub trait WebRendererExt {}
#[cfg(not(feature = "pmoserver"))]
impl WebRendererExt for () {}

View File

@@ -0,0 +1,24 @@
//! Erreurs liées au WebRenderer
use thiserror::Error;
#[derive(Error, Debug)]
pub enum WebRendererError {
#[error("Session not found: {0}")]
SessionNotFound(String),
#[error("Failed to send message to websocket: {0}")]
WebSocketSendError(String),
#[error("Invalid argument: {0}")]
InvalidArgument(String),
#[error("Failed to create device: {0}")]
DeviceCreationError(String),
#[error("Failed to register device: {0}")]
RegistrationError(String),
#[error("Server not available")]
ServerNotAvailable,
}

View File

@@ -0,0 +1,340 @@
//! Action handlers SOAP → WebSocket pour le WebRenderer
//!
//! Chaque handler bridge une action UPnP vers une commande WebSocket
//! envoyée au navigateur, ou lit l'état partagé pour les requêtes GET.
use std::sync::Arc;
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::{SharedSender, SharedState};
type ActionFuture =
std::pin::Pin<Box<dyn std::future::Future<Output = Result<ActionData, ActionError>> + Send>>;
// ─── AVTransport Handlers ───────────────────────────────────────────────────
pub fn play_handler(ws: SharedSender, state: SharedState) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture {
let ws = ws.clone();
let state = state.clone();
Box::pin(async move {
ws.send(ServerMessage::Command {
action: TransportAction::Play,
params: None,
});
state.write().playback_state = PlaybackState::Playing;
Ok(data)
})
})
}
pub fn stop_handler(ws: SharedSender, state: SharedState) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture {
let ws = ws.clone();
let state = state.clone();
Box::pin(async move {
ws.send(ServerMessage::Command {
action: TransportAction::Stop,
params: None,
});
state.write().playback_state = PlaybackState::Stopped;
Ok(data)
})
})
}
pub fn pause_handler(ws: SharedSender, state: SharedState) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture {
let ws = ws.clone();
let state = state.clone();
Box::pin(async move {
ws.send(ServerMessage::Command {
action: TransportAction::Pause,
params: None,
});
state.write().playback_state = PlaybackState::Paused;
Ok(data)
})
})
}
pub fn next_handler(ws: SharedSender) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture {
let ws = ws.clone();
Box::pin(async move {
ws.send(ServerMessage::Command {
action: TransportAction::Play,
params: None,
});
Ok(data)
})
})
}
pub fn previous_handler(ws: SharedSender) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture {
let ws = ws.clone();
Box::pin(async move {
ws.send(ServerMessage::Command {
action: TransportAction::Play,
params: None,
});
Ok(data)
})
})
}
pub fn seek_handler(ws: SharedSender) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture {
let ws = ws.clone();
Box::pin(async move {
let target: String = get!(&data, "Target", String);
ws.send(ServerMessage::Command {
action: TransportAction::Seek,
params: Some(CommandParams {
uri: None,
metadata: None,
position: Some(target),
}),
});
Ok(data)
})
})
}
pub fn set_uri_handler(ws: SharedSender, state: SharedState) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture {
let ws = ws.clone();
let state = state.clone();
Box::pin(async move {
let uri: String = get!(&data, "CurrentURI", String);
let metadata: String = get_value::<String>(&data, "CurrentURIMetaData")
.or_else(|_| {
get_value::<DIDLLite>(&data, "CurrentURIMetaData")
.map(|didl| didl.to_xml())
})
.unwrap_or_default();
ws.send(ServerMessage::Command {
action: TransportAction::SetUri,
params: Some(CommandParams {
uri: Some(uri.clone()),
metadata: Some(metadata.clone()),
position: None,
}),
});
{
let mut s = state.write();
s.current_uri = Some(uri);
s.current_metadata = Some(metadata);
s.playback_state = PlaybackState::Transitioning;
}
Ok(data)
})
})
}
pub fn set_next_uri_handler(ws: SharedSender, state: SharedState) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture {
let ws = ws.clone();
let state = state.clone();
Box::pin(async move {
let uri: String = get!(&data, "NextURI", String);
let metadata: String = get_value::<String>(&data, "NextURIMetaData")
.or_else(|_| {
get_value::<DIDLLite>(&data, "NextURIMetaData")
.map(|didl| didl.to_xml())
})
.unwrap_or_default();
ws.send(ServerMessage::Command {
action: TransportAction::SetNextUri,
params: Some(CommandParams {
uri: Some(uri.clone()),
metadata: Some(metadata.clone()),
position: None,
}),
});
{
let mut s = state.write();
s.next_uri = Some(uri);
s.next_metadata = Some(metadata);
}
Ok(data)
})
})
}
pub fn get_position_info_handler(state: SharedState) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture {
let state = state.clone();
Box::pin(async move {
let mut data = data;
let s = state.read();
set!(
&mut data,
"Track",
if s.current_uri.is_some() { 1u32 } else { 0u32 }
);
set!(
&mut data,
"TrackDuration",
s.duration.clone().unwrap_or_else(|| "00:00:00".to_string())
);
set!(
&mut data,
"TrackURI",
s.current_uri.clone().unwrap_or_default()
);
set!(
&mut data,
"TrackMetaData",
s.current_metadata.clone().unwrap_or_default()
);
set!(
&mut data,
"RelTime",
s.position.clone().unwrap_or_else(|| "00:00:00".to_string())
);
set!(
&mut data,
"AbsTime",
s.position.clone().unwrap_or_else(|| "00:00:00".to_string())
);
Ok(data)
})
})
}
pub fn get_transport_info_handler(state: SharedState) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture {
let state = state.clone();
Box::pin(async move {
let mut data = data;
let s = state.read();
tracing::info!("[WebRenderer] GetTransportInfo handler called, state={:?}", s.playback_state);
let transport_state = match s.playback_state {
PlaybackState::Stopped => "STOPPED",
PlaybackState::Playing => "PLAYING",
PlaybackState::Paused => "PAUSED_PLAYBACK",
PlaybackState::Transitioning => "TRANSITIONING",
};
set!(
&mut data,
"CurrentTransportState",
transport_state.to_string()
);
set!(&mut data, "CurrentTransportStatus", "OK".to_string());
set!(&mut data, "CurrentSpeed", "1".to_string());
Ok(data)
})
})
}
pub fn get_media_info_handler(state: SharedState) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture {
let state = state.clone();
Box::pin(async move {
let mut data = data;
let s = state.read();
set!(
&mut data,
"NrTracks",
if s.current_uri.is_some() { 1u32 } else { 0u32 }
);
set!(
&mut data,
"CurrentURI",
s.current_uri.clone().unwrap_or_default()
);
set!(
&mut data,
"CurrentURIMetaData",
s.current_metadata.clone().unwrap_or_default()
);
set!(
&mut data,
"NextURI",
s.next_uri.clone().unwrap_or_default()
);
set!(
&mut data,
"NextURIMetaData",
s.next_metadata.clone().unwrap_or_default()
);
Ok(data)
})
})
}
// ─── RenderingControl Handlers ──────────────────────────────────────────────
pub fn set_volume_handler(ws: SharedSender, state: SharedState) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture {
let ws = ws.clone();
let state = state.clone();
Box::pin(async move {
let volume: u16 = get!(&data, "DesiredVolume", u16);
ws.send(ServerMessage::SetVolume { volume });
state.write().volume = volume;
Ok(data)
})
})
}
pub fn get_volume_handler(state: SharedState) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture {
let state = state.clone();
Box::pin(async move {
let mut data = data;
let volume = state.read().volume;
set!(&mut data, "CurrentVolume", volume);
Ok(data)
})
})
}
pub fn set_mute_handler(ws: SharedSender, state: SharedState) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture {
let ws = ws.clone();
let state = state.clone();
Box::pin(async move {
let mute: bool = get!(&data, "DesiredMute", bool);
ws.send(ServerMessage::SetMute { mute });
state.write().mute = mute;
Ok(data)
})
})
}
pub fn get_mute_handler(state: SharedState) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture {
let state = state.clone();
Box::pin(async move {
let mut data = data;
let mute = state.read().mute;
set!(&mut data, "CurrentMute", mute);
Ok(data)
})
})
}
// ─── ConnectionManager Handlers ─────────────────────────────────────────────
pub fn get_protocol_info_handler() -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture {
Box::pin(async move {
let mut data = data;
set!(&mut data, "Source", String::new());
set!(
&mut data,
"Sink",
"http-get:*:audio/mpeg:*,http-get:*:audio/mp4:*,http-get:*:audio/ogg:*,http-get:*:audio/flac:*,http-get:*:audio/wav:*,http-get:*:audio/x-flac:*,http-get:*:audio/aac:*,http-get:*:audio/webm:*".to_string()
);
Ok(data)
})
})
}

25
pmowebrenderer/src/lib.rs Normal file
View File

@@ -0,0 +1,25 @@
//! PMO Web Renderer - Transforme un navigateur en MediaRenderer UPnP privé
mod error;
mod handlers;
mod messages;
mod renderer;
mod session;
mod state;
mod websocket;
#[cfg(feature = "pmoserver")]
mod config;
pub use error::WebRendererError;
pub use messages::{
BrowserCapabilities, ClientMessage, CommandParams, PlaybackState, RendererInfo, ServerMessage,
TrackMetadata, TransportAction,
};
pub use renderer::{FactoryError, WebRendererFactory};
pub use session::{SessionManager, WebRendererSession};
pub use state::{RendererState, SharedState};
pub use websocket::{websocket_handler, WebSocketState};
#[cfg(feature = "pmoserver")]
pub use config::WebRendererExt;

View File

@@ -0,0 +1,108 @@
//! Messages WebSocket pour la communication Backend ↔ Navigateur
use serde::{Deserialize, Serialize};
/// Messages envoyés du Backend → Navigateur
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum ServerMessage {
SessionCreated {
token: String,
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 {
action: TransportAction,
#[serde(skip_serializing_if = "Option::is_none")]
params: Option<CommandParams>,
},
SetVolume {
volume: u16,
},
SetMute {
mute: bool,
},
Ping,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum TransportAction {
Play,
Pause,
Stop,
Seek,
SetUri,
SetNextUri,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CommandParams {
#[serde(skip_serializing_if = "Option::is_none")]
pub uri: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub metadata: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub position: Option<String>,
}
/// Messages envoyés du Navigateur → Backend
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum ClientMessage {
Init { capabilities: BrowserCapabilities },
StateUpdate { state: PlaybackState },
PositionUpdate { position: String, duration: String },
MetadataUpdate { metadata: TrackMetadata },
VolumeUpdate { volume: u16, mute: bool },
/// Envoyé quand la piste courante se termine naturellement (gapless).
/// Le backend fait avancer current → next dans l'état partagé.
TrackEnded,
Pong,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
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 supported_formats: Vec<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
pub enum PlaybackState {
Stopped,
Playing,
Paused,
Transitioning,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TrackMetadata {
pub title: Option<String>,
pub artist: Option<String>,
pub album: Option<String>,
pub duration: Option<String>,
pub album_art_uri: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct RendererInfo {
pub udn: String,
pub friendly_name: String,
pub model_name: String,
pub description_url: String,
}

View File

@@ -0,0 +1,650 @@
//! Factory pour créer des instances WebRenderer privées
//!
//! Construit des Device/Service models UPnP dynamiques avec des action handlers
//! qui relaient les commandes SOAP vers le navigateur via WebSocket.
use std::sync::Arc;
use thiserror::Error;
use tokio::sync::mpsc;
use pmoupnp::actions::{Action, Argument};
use pmoupnp::devices::Device;
use pmoupnp::services::Service;
use crate::handlers;
use crate::messages::ServerMessage;
use crate::state::{SharedSender, SharedState};
// ─── Réimport des variables statiques de pmomediarenderer ───────────────────
// Variables AVTransport
use pmomediarenderer::avtransport::variables::{
ABSOLUTETIMEPOSITION, AVTRANSPORTNEXTURI, AVTRANSPORTNEXTURIMETADATA, AVTRANSPORTURI,
AVTRANSPORTURIMETADATA, A_ARG_TYPE_INSTANCE_ID as AVT_INSTANCE_ID, A_ARG_TYPE_PLAY_SPEED,
A_ARG_TYPE_SEEKMODE, CURRENTMEDIADURATION, CURRENTPLAYMODE, CURRENTTRACK, CURRENTTRACKDURATION,
CURRENTTRACKMETADATA, CURRENTTRACKURI, NUMBEROFTRACKS, PLAYBACKSTORAGEMEDIUM,
POSSIBLEPLAYBACKSTORAGEMEDIA, RELATIVETIMEPOSITION, SEEKMODE, TRANSPORTPLAYSPEED,
TRANSPORTSTATE, TRANSPORTSTATUS,
};
// Variables RenderingControl
use pmomediarenderer::renderingcontrol::variables::{
A_ARG_TYPE_CHANNEL, A_ARG_TYPE_INSTANCE_ID as RC_INSTANCE_ID, MUTE, VOLUME,
};
// Variables ConnectionManager
use pmomediarenderer::connectionmanager::variables::{
A_ARG_TYPE_AVTRANSPORTID, A_ARG_TYPE_CONNECTIONID, A_ARG_TYPE_CONNECTIONSTATUS,
A_ARG_TYPE_DIRECTION, A_ARG_TYPE_PROTOCOLINFO, A_ARG_TYPE_RCSID, CURRENTCONNECTIONIDS,
SINKPROTOCOLINFO, SOURCEPROTOCOLINFO,
};
#[derive(Error, Debug)]
pub enum FactoryError {
#[error("Failed to add service to device: {0}")]
ServiceError(String),
#[error("Failed to add action to service: {0}")]
ActionError(String),
#[error("Failed to add variable to service: {0}")]
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;
impl WebRendererFactory {
/// Crée un Device model UPnP complet pour un WebRenderer.
///
/// `device_name` sert de clé pour retrouver l'UDN persistant dans la config.
/// `browser_ua` est le User-Agent complet (pour déterminer le nom affiché).
///
/// 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>,
state: SharedState,
) -> Result<(Device, SharedSender), FactoryError> {
let shared_sender = SharedSender::new(ws_sender);
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 short_name = extract_browser_name(browser_ua);
let device = Device::new(
device_name.to_string(),
"MediaRenderer".to_string(),
format!("Web Audio {}", short_name),
);
device
.add_service(Arc::new(avtransport))
.map_err(|e| FactoryError::ServiceError(format!("{:?}", e)))?;
device
.add_service(Arc::new(renderingcontrol))
.map_err(|e| FactoryError::ServiceError(format!("{:?}", e)))?;
device
.add_service(Arc::new(connectionmanager))
.map_err(|e| FactoryError::ServiceError(format!("{:?}", e)))?;
Ok((device, shared_sender))
}
/// Construit le service AVTransport avec les handlers WebSocket
fn build_avtransport(
ws: SharedSender,
state: SharedState,
) -> Result<Service, FactoryError> {
let mut svc = Service::new("AVTransport".to_string());
let add_var = |svc: &mut Service, var: &Arc<pmoupnp::state_variables::StateVariable>| {
svc.add_variable(Arc::clone(var))
.map_err(|e| FactoryError::VariableError(e.to_string()))
};
// Ajouter toutes les variables d'état
add_var(&mut svc, &AVT_INSTANCE_ID)?;
add_var(&mut svc, &A_ARG_TYPE_PLAY_SPEED)?;
add_var(&mut svc, &A_ARG_TYPE_SEEKMODE)?;
add_var(&mut svc, &ABSOLUTETIMEPOSITION)?;
add_var(&mut svc, &AVTRANSPORTNEXTURI)?;
add_var(&mut svc, &AVTRANSPORTNEXTURIMETADATA)?;
add_var(&mut svc, &AVTRANSPORTURI)?;
add_var(&mut svc, &AVTRANSPORTURIMETADATA)?;
add_var(&mut svc, &CURRENTMEDIADURATION)?;
add_var(&mut svc, &CURRENTPLAYMODE)?;
add_var(&mut svc, &CURRENTTRACK)?;
add_var(&mut svc, &CURRENTTRACKDURATION)?;
add_var(&mut svc, &CURRENTTRACKMETADATA)?;
add_var(&mut svc, &CURRENTTRACKURI)?;
add_var(&mut svc, &NUMBEROFTRACKS)?;
add_var(&mut svc, &PLAYBACKSTORAGEMEDIUM)?;
add_var(&mut svc, &POSSIBLEPLAYBACKSTORAGEMEDIA)?;
add_var(&mut svc, &RELATIVETIMEPOSITION)?;
add_var(&mut svc, &SEEKMODE)?;
add_var(&mut svc, &TRANSPORTPLAYSPEED)?;
add_var(&mut svc, &TRANSPORTSTATE)?;
add_var(&mut svc, &TRANSPORTSTATUS)?;
let add_action = |svc: &mut Service, action: Arc<Action>| {
svc.add_action(action)
.map_err(|e| FactoryError::ActionError(e.to_string()))
};
// Play
let mut play = Action::new("Play".to_string());
play.add_argument(Arc::new(Argument::new_in(
"InstanceID".to_string(),
Arc::clone(&AVT_INSTANCE_ID),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
play.add_argument(Arc::new(Argument::new_in(
"Speed".to_string(),
Arc::clone(&TRANSPORTPLAYSPEED),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
play.set_handler(handlers::play_handler(ws.clone(), state.clone()));
add_action(&mut svc, Arc::new(play))?;
// Stop
let mut stop = Action::new("Stop".to_string());
stop.add_argument(Arc::new(Argument::new_in(
"InstanceID".to_string(),
Arc::clone(&AVT_INSTANCE_ID),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
stop.set_handler(handlers::stop_handler(ws.clone(), state.clone()));
add_action(&mut svc, Arc::new(stop))?;
// Pause
let mut pause = Action::new("Pause".to_string());
pause
.add_argument(Arc::new(Argument::new_in(
"InstanceID".to_string(),
Arc::clone(&AVT_INSTANCE_ID),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
pause.set_handler(handlers::pause_handler(ws.clone(), state.clone()));
add_action(&mut svc, Arc::new(pause))?;
// Next
let mut next = Action::new("Next".to_string());
next.add_argument(Arc::new(Argument::new_in(
"InstanceID".to_string(),
Arc::clone(&AVT_INSTANCE_ID),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
next.set_handler(handlers::next_handler(ws.clone()));
add_action(&mut svc, Arc::new(next))?;
// Previous
let mut previous = Action::new("Previous".to_string());
previous
.add_argument(Arc::new(Argument::new_in(
"InstanceID".to_string(),
Arc::clone(&AVT_INSTANCE_ID),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
previous.set_handler(handlers::previous_handler(ws.clone()));
add_action(&mut svc, Arc::new(previous))?;
// Seek
let mut seek = Action::new("Seek".to_string());
seek.add_argument(Arc::new(Argument::new_in(
"InstanceID".to_string(),
Arc::clone(&AVT_INSTANCE_ID),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
seek.add_argument(Arc::new(Argument::new_in(
"Unit".to_string(),
Arc::clone(&A_ARG_TYPE_SEEKMODE),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
seek.add_argument(Arc::new(Argument::new_in(
"Target".to_string(),
Arc::clone(&SEEKMODE),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
seek.set_handler(handlers::seek_handler(ws.clone()));
add_action(&mut svc, Arc::new(seek))?;
// SetAVTransportURI
let mut set_uri = Action::new("SetAVTransportURI".to_string());
set_uri
.add_argument(Arc::new(Argument::new_in(
"InstanceID".to_string(),
Arc::clone(&AVT_INSTANCE_ID),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
set_uri
.add_argument(Arc::new(Argument::new_in(
"CurrentURI".to_string(),
Arc::clone(&AVTRANSPORTURI),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
set_uri
.add_argument(Arc::new(Argument::new_in(
"CurrentURIMetaData".to_string(),
Arc::clone(&AVTRANSPORTURIMETADATA),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
set_uri.set_handler(handlers::set_uri_handler(ws.clone(), state.clone()));
add_action(&mut svc, Arc::new(set_uri))?;
// SetNextAVTransportURI
let mut set_next_uri = Action::new("SetNextAVTransportURI".to_string());
set_next_uri
.add_argument(Arc::new(Argument::new_in(
"InstanceID".to_string(),
Arc::clone(&AVT_INSTANCE_ID),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
set_next_uri
.add_argument(Arc::new(Argument::new_in(
"NextURI".to_string(),
Arc::clone(&AVTRANSPORTNEXTURI),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
set_next_uri
.add_argument(Arc::new(Argument::new_in(
"NextURIMetaData".to_string(),
Arc::clone(&AVTRANSPORTNEXTURIMETADATA),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
set_next_uri.set_handler(handlers::set_next_uri_handler(ws.clone(), state.clone()));
add_action(&mut svc, Arc::new(set_next_uri))?;
// GetPositionInfo
let mut get_pos = Action::new("GetPositionInfo".to_string());
get_pos
.add_argument(Arc::new(Argument::new_in(
"InstanceID".to_string(),
Arc::clone(&AVT_INSTANCE_ID),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_pos
.add_argument(Arc::new(Argument::new_out(
"Track".to_string(),
Arc::clone(&CURRENTTRACK),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_pos
.add_argument(Arc::new(Argument::new_out(
"TrackDuration".to_string(),
Arc::clone(&CURRENTTRACKDURATION),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_pos
.add_argument(Arc::new(Argument::new_out(
"TrackURI".to_string(),
Arc::clone(&CURRENTTRACKURI),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_pos
.add_argument(Arc::new(Argument::new_out(
"TrackMetaData".to_string(),
Arc::clone(&CURRENTTRACKMETADATA),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_pos
.add_argument(Arc::new(Argument::new_out(
"RelTime".to_string(),
Arc::clone(&RELATIVETIMEPOSITION),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_pos
.add_argument(Arc::new(Argument::new_out(
"AbsTime".to_string(),
Arc::clone(&ABSOLUTETIMEPOSITION),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_pos.set_stateful(false);
get_pos.set_handler(handlers::get_position_info_handler(state.clone()));
add_action(&mut svc, Arc::new(get_pos))?;
// GetTransportInfo
let mut get_info = Action::new("GetTransportInfo".to_string());
get_info
.add_argument(Arc::new(Argument::new_in(
"InstanceID".to_string(),
Arc::clone(&AVT_INSTANCE_ID),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_info
.add_argument(Arc::new(Argument::new_out(
"CurrentTransportState".to_string(),
Arc::clone(&TRANSPORTSTATE),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_info
.add_argument(Arc::new(Argument::new_out(
"CurrentTransportStatus".to_string(),
Arc::clone(&TRANSPORTSTATUS),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_info
.add_argument(Arc::new(Argument::new_out(
"CurrentSpeed".to_string(),
Arc::clone(&TRANSPORTPLAYSPEED),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_info.set_stateful(false);
get_info.set_handler(handlers::get_transport_info_handler(state.clone()));
add_action(&mut svc, Arc::new(get_info))?;
// GetMediaInfo
let mut get_media = Action::new("GetMediaInfo".to_string());
get_media
.add_argument(Arc::new(Argument::new_in(
"InstanceID".to_string(),
Arc::clone(&AVT_INSTANCE_ID),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_media
.add_argument(Arc::new(Argument::new_out(
"NrTracks".to_string(),
Arc::clone(&NUMBEROFTRACKS),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_media
.add_argument(Arc::new(Argument::new_out(
"CurrentURI".to_string(),
Arc::clone(&AVTRANSPORTURI),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_media
.add_argument(Arc::new(Argument::new_out(
"CurrentURIMetaData".to_string(),
Arc::clone(&AVTRANSPORTURIMETADATA),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_media
.add_argument(Arc::new(Argument::new_out(
"NextURI".to_string(),
Arc::clone(&AVTRANSPORTNEXTURI),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_media
.add_argument(Arc::new(Argument::new_out(
"NextURIMetaData".to_string(),
Arc::clone(&AVTRANSPORTNEXTURIMETADATA),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_media.set_stateful(false);
get_media.set_handler(handlers::get_media_info_handler(state.clone()));
add_action(&mut svc, Arc::new(get_media))?;
// GetTransportSettings (passthrough)
let mut get_settings = Action::new("GetTransportSettings".to_string());
get_settings
.add_argument(Arc::new(Argument::new_in(
"InstanceID".to_string(),
Arc::clone(&AVT_INSTANCE_ID),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
add_action(&mut svc, Arc::new(get_settings))?;
// GetDeviceCapabilities (passthrough)
let mut get_caps = Action::new("GetDeviceCapabilities".to_string());
get_caps
.add_argument(Arc::new(Argument::new_in(
"InstanceID".to_string(),
Arc::clone(&AVT_INSTANCE_ID),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
add_action(&mut svc, Arc::new(get_caps))?;
// GetCurrentTransportActions (passthrough)
let mut get_actions = Action::new("GetCurrentTransportActions".to_string());
get_actions
.add_argument(Arc::new(Argument::new_in(
"InstanceID".to_string(),
Arc::clone(&AVT_INSTANCE_ID),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
add_action(&mut svc, Arc::new(get_actions))?;
Ok(svc)
}
/// Construit le service RenderingControl avec les handlers WebSocket
fn build_renderingcontrol(
ws: SharedSender,
state: SharedState,
) -> Result<Service, FactoryError> {
let mut svc = Service::new("RenderingControl".to_string());
svc.add_variable(Arc::clone(&RC_INSTANCE_ID))
.map_err(|e| FactoryError::VariableError(e.to_string()))?;
svc.add_variable(Arc::clone(&A_ARG_TYPE_CHANNEL))
.map_err(|e| FactoryError::VariableError(e.to_string()))?;
svc.add_variable(Arc::clone(&VOLUME))
.map_err(|e| FactoryError::VariableError(e.to_string()))?;
svc.add_variable(Arc::clone(&MUTE))
.map_err(|e| FactoryError::VariableError(e.to_string()))?;
// SetVolume
let mut set_vol = Action::new("SetVolume".to_string());
set_vol
.add_argument(Arc::new(Argument::new_in(
"InstanceID".to_string(),
Arc::clone(&RC_INSTANCE_ID),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
set_vol
.add_argument(Arc::new(Argument::new_in(
"Channel".to_string(),
Arc::clone(&A_ARG_TYPE_CHANNEL),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
set_vol
.add_argument(Arc::new(Argument::new_in(
"DesiredVolume".to_string(),
Arc::clone(&VOLUME),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
set_vol.set_handler(handlers::set_volume_handler(ws.clone(), state.clone()));
svc.add_action(Arc::new(set_vol))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
// GetVolume
let mut get_vol = Action::new("GetVolume".to_string());
get_vol
.add_argument(Arc::new(Argument::new_in(
"InstanceID".to_string(),
Arc::clone(&RC_INSTANCE_ID),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_vol
.add_argument(Arc::new(Argument::new_in(
"Channel".to_string(),
Arc::clone(&A_ARG_TYPE_CHANNEL),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_vol
.add_argument(Arc::new(Argument::new_out(
"CurrentVolume".to_string(),
Arc::clone(&VOLUME),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_vol.set_stateful(false);
get_vol.set_handler(handlers::get_volume_handler(state.clone()));
svc.add_action(Arc::new(get_vol))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
// SetMute
let mut set_mute = Action::new("SetMute".to_string());
set_mute
.add_argument(Arc::new(Argument::new_in(
"InstanceID".to_string(),
Arc::clone(&RC_INSTANCE_ID),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
set_mute
.add_argument(Arc::new(Argument::new_in(
"Channel".to_string(),
Arc::clone(&A_ARG_TYPE_CHANNEL),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
set_mute
.add_argument(Arc::new(Argument::new_in(
"DesiredMute".to_string(),
Arc::clone(&MUTE),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
set_mute.set_handler(handlers::set_mute_handler(ws.clone(), state.clone()));
svc.add_action(Arc::new(set_mute))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
// GetMute
let mut get_mute = Action::new("GetMute".to_string());
get_mute
.add_argument(Arc::new(Argument::new_in(
"InstanceID".to_string(),
Arc::clone(&RC_INSTANCE_ID),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_mute
.add_argument(Arc::new(Argument::new_in(
"Channel".to_string(),
Arc::clone(&A_ARG_TYPE_CHANNEL),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_mute
.add_argument(Arc::new(Argument::new_out(
"CurrentMute".to_string(),
Arc::clone(&MUTE),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_mute.set_stateful(false);
get_mute.set_handler(handlers::get_mute_handler(state.clone()));
svc.add_action(Arc::new(get_mute))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
Ok(svc)
}
/// Construit le service ConnectionManager
fn build_connectionmanager() -> Result<Service, FactoryError> {
let mut svc = Service::new("ConnectionManager".to_string());
svc.add_variable(Arc::clone(&A_ARG_TYPE_CONNECTIONID))
.map_err(|e| FactoryError::VariableError(e.to_string()))?;
svc.add_variable(Arc::clone(&A_ARG_TYPE_CONNECTIONSTATUS))
.map_err(|e| FactoryError::VariableError(e.to_string()))?;
svc.add_variable(Arc::clone(&A_ARG_TYPE_DIRECTION))
.map_err(|e| FactoryError::VariableError(e.to_string()))?;
svc.add_variable(Arc::clone(&A_ARG_TYPE_PROTOCOLINFO))
.map_err(|e| FactoryError::VariableError(e.to_string()))?;
svc.add_variable(Arc::clone(&A_ARG_TYPE_RCSID))
.map_err(|e| FactoryError::VariableError(e.to_string()))?;
svc.add_variable(Arc::clone(&A_ARG_TYPE_AVTRANSPORTID))
.map_err(|e| FactoryError::VariableError(e.to_string()))?;
svc.add_variable(Arc::clone(&CURRENTCONNECTIONIDS))
.map_err(|e| FactoryError::VariableError(e.to_string()))?;
svc.add_variable(Arc::clone(&SINKPROTOCOLINFO))
.map_err(|e| FactoryError::VariableError(e.to_string()))?;
svc.add_variable(Arc::clone(&SOURCEPROTOCOLINFO))
.map_err(|e| FactoryError::VariableError(e.to_string()))?;
// GetProtocolInfo
let mut get_proto = Action::new("GetProtocolInfo".to_string());
get_proto
.add_argument(Arc::new(Argument::new_out(
"Source".to_string(),
Arc::clone(&SOURCEPROTOCOLINFO),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_proto
.add_argument(Arc::new(Argument::new_out(
"Sink".to_string(),
Arc::clone(&SINKPROTOCOLINFO),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_proto.set_stateful(false);
get_proto.set_handler(handlers::get_protocol_info_handler());
svc.add_action(Arc::new(get_proto))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
// GetCurrentConnectionIDs
let mut get_ids = Action::new("GetCurrentConnectionIDs".to_string());
get_ids
.add_argument(Arc::new(Argument::new_out(
"ConnectionIDs".to_string(),
Arc::clone(&CURRENTCONNECTIONIDS),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
svc.add_action(Arc::new(get_ids))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
// GetCurrentConnectionInfo
let mut get_conn = Action::new("GetCurrentConnectionInfo".to_string());
get_conn
.add_argument(Arc::new(Argument::new_in(
"ConnectionID".to_string(),
Arc::clone(&A_ARG_TYPE_CONNECTIONID),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_conn
.add_argument(Arc::new(Argument::new_out(
"RcsID".to_string(),
Arc::clone(&A_ARG_TYPE_RCSID),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_conn
.add_argument(Arc::new(Argument::new_out(
"AVTransportID".to_string(),
Arc::clone(&A_ARG_TYPE_AVTRANSPORTID),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_conn
.add_argument(Arc::new(Argument::new_out(
"ProtocolInfo".to_string(),
Arc::clone(&A_ARG_TYPE_PROTOCOLINFO),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_conn
.add_argument(Arc::new(Argument::new_out(
"PeerConnectionManager".to_string(),
Arc::clone(&A_ARG_TYPE_PROTOCOLINFO),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_conn
.add_argument(Arc::new(Argument::new_out(
"PeerConnectionID".to_string(),
Arc::clone(&A_ARG_TYPE_CONNECTIONID),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_conn
.add_argument(Arc::new(Argument::new_out(
"Direction".to_string(),
Arc::clone(&A_ARG_TYPE_DIRECTION),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
get_conn
.add_argument(Arc::new(Argument::new_out(
"Status".to_string(),
Arc::clone(&A_ARG_TYPE_CONNECTIONSTATUS),
)))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
svc.add_action(Arc::new(get_conn))
.map_err(|e| FactoryError::ActionError(format!("{:?}", e)))?;
Ok(svc)
}
}

View File

@@ -0,0 +1,133 @@
//! Gestionnaire de sessions WebRenderer
use parking_lot::RwLock;
use std::collections::HashMap;
use std::sync::Arc;
use std::time::{Duration, SystemTime};
use pmoupnp::devices::DeviceInstance;
use crate::state::{SharedSender, SharedState};
/// Session WebSocket liée à un MediaRenderer privé
pub struct WebRendererSession {
pub token: String,
pub udn: String,
pub device_instance: Arc<DeviceInstance>,
/// Sender partagé : mis à jour à chaque reconnexion WebSocket.
pub shared_sender: SharedSender,
pub state: SharedState,
pub created_at: SystemTime,
pub last_activity: Arc<RwLock<SystemTime>>,
}
/// Gestionnaire global des sessions
#[derive(Clone)]
pub struct SessionManager {
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,
}
impl SessionManager {
pub fn new(timeout_duration: Duration) -> Self {
let manager = Self {
sessions: Arc::new(RwLock::new(HashMap::new())),
senders: Arc::new(RwLock::new(HashMap::new())),
states: Arc::new(RwLock::new(HashMap::new())),
timeout_duration,
};
manager.spawn_cleanup_task();
manager
}
pub fn add_session(&self, session: Arc<WebRendererSession>) {
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.senders.write().insert(udn.clone(), sender);
self.states.write().insert(udn, state);
tracing::info!(token = %token, "WebRenderer session added");
}
pub fn get_session(&self, token: &str) -> Option<Arc<WebRendererSession>> {
let sessions = self.sessions.read();
if let Some(session) = sessions.get(token) {
*session.last_activity.write() = SystemTime::now();
Some(session.clone())
} else {
None
}
}
/// 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>> {
let session = self.sessions.write().remove(token);
if let Some(ref s) = session {
tracing::info!(token = %s.token, udn = %s.udn, "WebRenderer session removed");
}
session
}
pub fn list_sessions(&self) -> Vec<Arc<WebRendererSession>> {
self.sessions.read().values().cloned().collect()
}
fn spawn_cleanup_task(&self) {
let sessions = self.sessions.clone();
let timeout = self.timeout_duration;
tokio::spawn(async move {
let mut interval = tokio::time::interval(Duration::from_secs(60));
loop {
interval.tick().await;
let now = SystemTime::now();
let mut to_remove = Vec::new();
{
let sessions_guard = sessions.read();
for (token, session) in sessions_guard.iter() {
let last = session.last_activity.read();
if let Ok(elapsed) = now.duration_since(*last) {
if elapsed > timeout {
to_remove.push(token.clone());
}
}
}
}
if !to_remove.is_empty() {
let mut sessions_guard = sessions.write();
for token in to_remove {
sessions_guard.remove(&token);
tracing::info!(token = %token, "Session expired and removed");
}
}
}
});
}
}

View File

@@ -0,0 +1,71 @@
//! État partagé du renderer (backend ↔ navigateur)
use parking_lot::RwLock;
use std::sync::Arc;
use tokio::sync::mpsc;
use crate::messages::{PlaybackState, ServerMessage};
/// État temps-réel du renderer (partagé backend ↔ navigateur)
#[derive(Debug, Clone)]
pub struct RendererState {
pub playback_state: PlaybackState,
pub current_uri: Option<String>,
pub current_metadata: Option<String>,
pub next_uri: Option<String>,
pub next_metadata: Option<String>,
pub position: Option<String>,
pub duration: Option<String>,
pub volume: u16,
pub mute: bool,
}
impl Default for RendererState {
fn default() -> Self {
Self {
playback_state: PlaybackState::Stopped,
current_uri: None,
current_metadata: None,
next_uri: None,
next_metadata: None,
position: None,
duration: None,
volume: 100,
mute: false,
}
}
}
/// Alias pour l'état partagé
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

@@ -0,0 +1,646 @@
//! Handler WebSocket pour les connexions navigateur → WebRenderer
use std::sync::Arc;
use std::time::SystemTime;
use axum::extract::ws::{Message, WebSocket};
use axum::extract::{State, WebSocketUpgrade};
use axum::response::IntoResponse;
use futures::{SinkExt, StreamExt};
use parking_lot::RwLock;
use tokio::sync::mpsc;
use uuid::Uuid;
use pmoupnp::devices::DeviceInstance;
use pmoupnp::variable_types::StateValue;
#[cfg(not(feature = "pmoserver"))]
use pmoupnp::UpnpModel;
use pmoupnp::UpnpTypedInstance;
use crate::messages::*;
use crate::renderer::WebRendererFactory;
use crate::session::{SessionManager, WebRendererSession};
use crate::state::{RendererState, SharedState};
#[cfg(feature = "pmoserver")]
use pmocontrol::model::{RendererCapabilities, RendererProtocol};
#[cfg(feature = "pmoserver")]
use pmocontrol::{ControlPoint, DeviceId};
/// État partagé du serveur WebSocket
#[derive(Clone)]
pub struct WebSocketState {
pub session_manager: Arc<SessionManager>,
#[cfg(feature = "pmoserver")]
pub control_point: Arc<ControlPoint>,
}
/// Handler pour la connexion WebSocket
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
let (tx, mut rx) = mpsc::unbounded_channel::<ServerMessage>();
// Task pour envoyer les messages du canal vers le WebSocket
let send_task = tokio::spawn(async move {
while let Some(msg) = rx.recv().await {
if let Ok(json) = serde_json::to_string(&msg) {
if sink.send(Message::Text(json.into())).await.is_err() {
break;
}
}
}
});
let mut session_token: Option<String> = None;
#[allow(unused_variables, unused_assignments, unused_mut)]
let mut device_udn: Option<String> = None;
// Boucle de réception des messages du navigateur
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 ou reconnecter 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();
let model = session.device_instance.get_model();
// Envoyer la confirmation au navigateur
let _ = tx.send(ServerMessage::SessionCreated {
token: token.clone(),
renderer_info: RendererInfo {
udn: udn.clone(),
friendly_name: model.friendly_name().to_string(),
model_name: model.model_name().to_string(),
description_url: format!(
"{}{}",
session.device_instance.base_url(),
session.device_instance.description_route()
),
},
});
// 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());
#[cfg(feature = "pmoserver")]
{ device_udn = Some(udn.clone()); }
state.session_manager.add_session(session);
tracing::info!(
token = %token,
udn = %udn,
"WebRenderer initialized for browser: {}",
capabilities.user_agent
);
}
Err(e) => {
tracing::error!("Failed to create WebRenderer: {:?}", e);
break;
}
}
}
Ok(ClientMessage::StateUpdate { state: new_state }) => {
if let Some(ref token) = session_token {
if let Some(session) = state.session_manager.get_session(token) {
{
session.state.write().playback_state = new_state.clone();
}
update_transport_state_var(&session.device_instance, &new_state)
.await;
}
}
}
Ok(ClientMessage::PositionUpdate { position, duration }) => {
if let Some(ref token) = session_token {
if let Some(session) = state.session_manager.get_session(token) {
{
let mut s = session.state.write();
s.position = Some(position.clone());
s.duration = Some(duration.clone());
}
update_position_vars(
&session.device_instance,
&position,
&duration,
)
.await;
}
}
}
Ok(ClientMessage::MetadataUpdate { metadata }) => {
if let Some(ref token) = session_token {
if let Some(session) = state.session_manager.get_session(token) {
let didl = build_didl_from_metadata(&metadata);
{
session.state.write().current_metadata = Some(didl.clone());
}
update_metadata_var(&session.device_instance, &didl).await;
}
}
}
Ok(ClientMessage::VolumeUpdate { volume, mute }) => {
if let Some(ref token) = session_token {
if let Some(session) = state.session_manager.get_session(token) {
{
let mut s = session.state.write();
s.volume = volume;
s.mute = mute;
}
update_volume_vars(&session.device_instance, volume, mute).await;
}
}
}
Ok(ClientMessage::TrackEnded) => {
if let Some(ref token) = session_token {
if let Some(session) = state.session_manager.get_session(token) {
let (next_uri, next_metadata, had_next) = {
let mut s = session.state.write();
let uri = s.next_uri.take();
let meta = s.next_metadata.take();
let had_next = uri.is_some();
s.current_uri = uri.clone();
s.current_metadata = meta.clone();
s.next_uri = None;
s.next_metadata = None;
s.position = None;
s.duration = None;
// Si on avait une piste suivante (gapless), on reste en Playing.
// Sinon, on reste en Stopped pour que le watcher déclenche l'auto-advance.
if had_next {
s.playback_state = PlaybackState::Playing;
}
// Si had_next == false, le navigateur a déjà envoyé state_update:STOPPED,
// donc s.playback_state est déjà Stopped. On le laisse tel quel.
(uri, meta, had_next)
};
let new_state = if had_next {
PlaybackState::Playing
} else {
PlaybackState::Stopped
};
update_transport_state_var(
&session.device_instance,
&new_state,
)
.await;
// Mettre à jour AVTransportURI pour que le ControlPoint
// voie la nouvelle piste courante et envoie SetNextAVTransportURI
update_uri_vars(
&session.device_instance,
next_uri.as_deref().unwrap_or(""),
next_metadata.as_deref().unwrap_or(""),
"", // next_uri vide : le ControlPoint le remplira
"",
)
.await;
// Si c'était une transition gapless (on avait une piste suivante),
// avancer l'index de la queue dans le ControlPoint et prefetch la piste N+2.
// Si pas de piste suivante, le watcher verra STOPPED et déclenchera l'auto-advance.
#[cfg(feature = "pmoserver")]
if had_next {
let udn = session.udn.clone();
let cp = state.control_point.clone();
tokio::spawn(async move {
cp.advance_queue_and_prefetch(
&pmocontrol::DeviceId(udn),
);
});
}
tracing::debug!(
uri = ?next_uri,
had_next,
"WebRenderer TrackEnded: advanced to next track"
);
}
}
}
Ok(ClientMessage::Pong) => {}
Err(e) => {
tracing::warn!(error = %e, "Failed to parse client message");
}
}
}
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) => {
tracing::error!("WebSocket error: {}", e);
break;
}
_ => {}
}
}
// Cleanup à la déconnexion
tracing::info!("WebRenderer WebSocket handler exiting (session_token={:?})", session_token);
if let Some(token) = session_token {
state.session_manager.remove_session(&token);
}
// Marquer le renderer comme offline dans le ControlPoint
#[cfg(feature = "pmoserver")]
if let Some(ref udn) = device_udn {
if let Ok(mut registry) = state.control_point.registry().write() {
registry.device_says_byebye(udn);
}
tracing::info!(udn = %udn, "WebRenderer disconnected and marked offline");
}
send_task.abort();
}
/// 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(
capabilities: &BrowserCapabilities,
ws_sender: mpsc::UnboundedSender<ServerMessage>,
ws_state: &WebSocketState,
) -> Result<Arc<WebRendererSession>, crate::error::WebRendererError> {
let token = Uuid::new_v4().to_string();
// Persister l'UDN dérivé de l'instance_id dans la config pour que device_instance.rs
// le retrouve de façon déterministe. La clé ("MediaRenderer", instance_id) est unique
// par onglet/navigateur et stable entre les reloads.
let instance_udn = capabilities.instance_id.clone();
if let Err(e) = pmoconfig::get_config().set_device_udn(
"MediaRenderer",
&instance_udn,
instance_udn.clone(),
) {
tracing::warn!("WebRenderer: failed to persist UDN in config: {:?}", e);
}
// 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)
// Retourne (DeviceInstance, SharedSender effectif, SharedState effective pour cette session)
#[cfg(feature = "pmoserver")]
let (di, shared_sender, shared_state) = {
use pmoupnp::UpnpServerExt;
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 =
pmoserver::get_server().ok_or(crate::error::WebRendererError::ServerNotAvailable)?;
let existing_di = {
let server = server_arc.read().await;
server.get_device(&candidate_udn)
};
if let Some(di) = existing_di {
tracing::info!(udn = %candidate_udn, "WebRenderer: reusing device from registry (session expired)");
// 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");
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"))]
let (di, shared_sender, shared_state) = {
let new_state: SharedState = Arc::new(RwLock::new(RendererState::default()));
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()))?;
(Arc::new(device).create_instance(), new_sender, new_state)
};
let session = Arc::new(WebRendererSession {
token,
udn: full_udn,
device_instance: di,
shared_sender,
state: shared_state,
created_at: SystemTime::now(),
last_activity: Arc::new(RwLock::new(SystemTime::now())),
});
Ok(session)
}
/// Enregistre le DeviceInstance auprès du ControlPoint comme un renderer
#[cfg(feature = "pmoserver")]
fn register_with_control_point(
di: &Arc<DeviceInstance>,
ws_state: &WebSocketState,
) -> Result<(), crate::error::WebRendererError> {
let base_url = di.base_url().to_string();
let udn = di.udn().to_ascii_lowercase();
// Préfixer avec "uuid:" pour correspondre au format SSDP et éviter les doublons
let udn_with_prefix = format!("uuid:{}", udn);
let device_route = di.route();
let model = di.get_model();
let avtransport_control_url = Some(format!(
"{}{}/service/AVTransport/control",
base_url, device_route
));
let rendering_control_url = Some(format!(
"{}{}/service/RenderingControl/control",
base_url, device_route
));
let connection_manager_url = Some(format!(
"{}{}/service/ConnectionManager/control",
base_url, device_route
));
let renderer_info = pmocontrol::RendererInfo::make(
DeviceId(udn_with_prefix.clone()),
udn_with_prefix.clone(),
model.friendly_name().to_string(),
model.model_name().to_string(),
"PMOMusic".to_string(),
RendererProtocol::UpnpAvOnly,
RendererCapabilities {
has_avtransport: true,
has_avtransport_set_next: true,
has_rendering_control: true,
has_connection_manager: true,
..Default::default()
},
format!("{}{}", base_url, di.description_route()),
"PMOMusic WebRenderer/1.0".to_string(),
Some("urn:schemas-upnp-org:service:AVTransport:1".to_string()),
avtransport_control_url,
Some("urn:schemas-upnp-org:service:RenderingControl:1".to_string()),
rendering_control_url,
Some("urn:schemas-upnp-org:service:ConnectionManager:1".to_string()),
connection_manager_url,
None, // oh_playlist_service_type
None, // oh_playlist_control_url
None, // oh_playlist_event_sub_url
None, // oh_info_service_type
None, // oh_info_control_url
None, // oh_info_event_sub_url
None, // oh_time_service_type
None, // oh_time_control_url
None, // oh_time_event_sub_url
None, // oh_volume_service_type
None, // oh_volume_control_url
None, // oh_radio_service_type
None, // oh_radio_control_url
None, // oh_product_service_type
None, // oh_product_control_url
);
if let Ok(mut registry) = ws_state.control_point.registry().write() {
// max_age élevé car pas de SSDP — cleanup à la déconnexion WS
registry.push_renderer(&renderer_info, 86400);
}
tracing::info!(udn = %udn, "WebRenderer registered with ControlPoint");
Ok(())
}
// ─── Mise à jour des StateVarInstance UPnP ──────────────────────────────────
async fn update_transport_state_var(di: &DeviceInstance, state: &PlaybackState) {
let upnp_state = match state {
PlaybackState::Stopped => "STOPPED",
PlaybackState::Playing => "PLAYING",
PlaybackState::Paused => "PAUSED_PLAYBACK",
PlaybackState::Transitioning => "TRANSITIONING",
};
if let Some(service) = di.get_service("AVTransport") {
if let Some(var) = service.get_variable("TransportState") {
let _ = var
.set_value(StateValue::String(upnp_state.to_string()))
.await;
}
}
}
async fn update_position_vars(di: &DeviceInstance, position: &str, duration: &str) {
if let Some(service) = di.get_service("AVTransport") {
if let Some(var) = service.get_variable("RelativeTimePosition") {
let _ = var
.set_value(StateValue::String(position.to_string()))
.await;
}
if let Some(var) = service.get_variable("AbsoluteTimePosition") {
let _ = var
.set_value(StateValue::String(position.to_string()))
.await;
}
if let Some(var) = service.get_variable("CurrentTrackDuration") {
let _ = var
.set_value(StateValue::String(duration.to_string()))
.await;
}
}
}
async fn update_uri_vars(
di: &DeviceInstance,
current_uri: &str,
current_metadata: &str,
next_uri: &str,
next_metadata: &str,
) {
if let Some(service) = di.get_service("AVTransport") {
if let Some(var) = service.get_variable("AVTransportURI") {
let _ = var
.set_value(StateValue::String(current_uri.to_string()))
.await;
}
if let Some(var) = service.get_variable("AVTransportURIMetaData") {
let _ = var
.set_value(StateValue::String(current_metadata.to_string()))
.await;
}
if let Some(var) = service.get_variable("AVTransportNextURI") {
let _ = var
.set_value(StateValue::String(next_uri.to_string()))
.await;
}
if let Some(var) = service.get_variable("AVTransportNextURIMetaData") {
let _ = var
.set_value(StateValue::String(next_metadata.to_string()))
.await;
}
}
}
async fn update_metadata_var(di: &DeviceInstance, didl: &str) {
if let Some(service) = di.get_service("AVTransport") {
if let Some(var) = service.get_variable("CurrentTrackMetaData") {
let _ = var.set_value(StateValue::String(didl.to_string())).await;
}
}
}
async fn update_volume_vars(di: &DeviceInstance, volume: u16, mute: bool) {
if let Some(service) = di.get_service("RenderingControl") {
if let Some(var) = service.get_variable("Volume") {
let _ = var.set_value(StateValue::UI2(volume)).await;
}
if let Some(var) = service.get_variable("Mute") {
let _ = var.set_value(StateValue::Boolean(mute)).await;
}
}
}
fn build_didl_from_metadata(metadata: &TrackMetadata) -> String {
use pmodidl::{DIDLLite, Item, Resource};
use pmoutils::ToXmlElement;
let item = Item {
id: "0".to_string(),
parent_id: "-1".to_string(),
restricted: Some("1".to_string()),
title: metadata
.title
.clone()
.unwrap_or_else(|| "Unknown".to_string()),
creator: None,
class: "object.item.audioItem.musicTrack".to_string(),
artist: metadata.artist.clone(),
album: metadata.album.clone(),
genre: None,
album_art: metadata.album_art_uri.clone(),
album_art_pk: None,
date: None,
original_track_number: None,
resources: vec![Resource {
protocol_info: "http-get:*:audio/*:*".to_string(),
duration: metadata.duration.clone(),
url: "".to_string(),
bits_per_sample: None,
sample_frequency: None,
nr_audio_channels: None,
}],
descriptions: vec![],
};
let didl = DIDLLite {
items: vec![item],
..Default::default()
};
didl.to_xml()
}

View File

@@ -1 +1 @@
0.3.20
0.3.21