10 Commits

Author SHA1 Message Date
a502e4fc3e Implémentation de la reconnexion stable avec synchronisation d'état
Cette mise à jour permet une reconnexion stable du renderer Web après un reload de page, en conservant l'état audio (URI, volume, lecture en cours). 

- Ajout d'un identifiant d'instance stable (UUID) dans le navigateur pour retrouver le même renderer UPnP
- Implémentation d'un mécanisme de synchronisation d'état (StateSync) lors des reconnexions
- Mise à jour du système de sender partagé (SharedSender) pour permettre le remplacement des connexions WebSocket
- Correction des handlers UPnP pour utiliser le nouveau type SharedSender
- Amélioration de la gestion des erreurs d'autoplay dans le moteur audio
- Mise à jour du numéro de version vers 0.3.21
2026-02-21 20:54:43 +01:00
32464c93fb Fix race condition in play() and set_uri
Corrige une condition de course où play() pouvait être appelé avant set_uri. Ajoute un flag playPending pour reporter la lecture jusqu'à ce que la source soit définie, et gère correctement les cas où l'élément média n'a pas encore les données nécessaires.
2026-02-21 15:57:52 +01:00
b3f1917441 Optimize album art URL handling and simplify image display logic
This commit refactors the album art URL normalization logic in Rust to use match expressions for cleaner and more concise code. It also simplifies the image display logic in Vue by removing the redundant cacheBustedUrl check, ensuring the image is displayed based solely on load and error states.
2026-02-21 15:48:15 +01:00
86df772354 Remove unnecessary console logs
This commit removes various console.log statements that were used for debugging purposes. These logs clutter the output and provide no value in production. The changes affect multiple files including components and composables related to media control and rendering.
2026-02-21 15:23:54 +01:00
fbe7291b43 Implémentation du playback gapless avec préchargement des pistes
Cette modification implémente un système de playback gapless pour les renderers web, permettant une transition fluide entre les pistes. Le changement inclut l'ajout d'un moteur audio qui utilise deux éléments <audio> en ping-pong connectés à un AudioContext pour le contrôle du volume/mute et le contrôle du flux gapless. Les commandes UPnP ont été mises à jour pour gérer le préchargement des pistes suivantes via SetNextAVTransportURI. Le backend a été modifié pour gérer les transitions gapless et précharger les pistes suivantes. Des ajustements ont également été apportés au serveur pour permettre les requêtes CORS nécessaires à ce fonctionnement.
2026-02-21 14:57:37 +01:00
03775f7574 Add CurrentSpeed to GetTransportInfo
This commit adds the CurrentSpeed field to the GetTransportInfo action, including updates to the action definition, handler, and web renderer. The speed is set to '1' as a default value.
2026-02-21 09:43:39 +01:00
58ec5f96f2 Normaliser le format de l'UDN avec le préfixe 'uuid:'
Ce commit normalise le format de l'UDN en ajoutant le préfixe 'uuid:' pour qu'il corresponde au format SSDP. Cela permet d'éviter les doublons et garantit une compatibilité avec les protocoles réseau. Les modifications affectent la création de l'instance et l'initialisation de la session WebRenderer.
2026-02-21 00:24:24 +01:00
cddd2afbd8 Implémentation du WebRenderer et améliorations du serveur
Ajout de la fonctionnalité WebRenderer permettant au navigateur de se connecter comme renderer UPnP.

- Implémentation du composant Vue useWebRenderer pour gérer la connexion WebSocket et le contrôle audio
- Ajout d'un proxy WebSocket dans Vite pour rediriger les requêtes /api/webrenderer/ws vers le serveur
- Modification du serveur pour permettre l'enregistrement dynamique de routes et extraction du JoinHandle
- Améliorations du parsing des métadonnées AVTransport pour gérer les types String et DIDLLite
- Ajout de la détection et du nommage des navigateurs dans le WebRenderer
- Mise à jour des dépendances avec tower 0.5.2
- Correction de la gestion des threads lors de l'arrêt du serveur HTTP
2026-02-21 00:16:05 +01:00
7276945b92 Fix SSDP multicast on macOS Sequoia+
Corrige les problèmes de multicast SSDP sur macOS Sequoia+ en ajustant les interfaces multicast et en utilisant Terminal.app pour l'exécution des binaires en debug et release.

- Dans le Makefile : modifie les cibles 'run' et 'run-release' pour lancer l'application via Terminal.app, nécessaire pour le multicast sur macOS Sequoia+
- Dans pmoupnp/src/ssdp/client.rs : améliore la gestion des interfaces multicast en rejoignant le groupe sur toutes les interfaces non loopback et en configurant explicitement l'interface de sortie
- Dans pmoupnp/src/ssdp/server.rs : corrige la configuration multicast du serveur en maintenant l'interface de sortie correcte et en désactivant le multicast loopback
2026-02-20 18:40:39 +01:00
1e882ba6c3 feat: implémentation du WebRenderer UPnP privé par navigateur
Ajout de la fonctionnalité WebRenderer permettant à chaque navigateur connecté de devenir un MediaRenderer UPnP privé.

- Création du crate pmowebrenderer avec l'architecture complète
- Implémentation des handlers SOAP → WebSocket pour les services AVTransport, RenderingControl et ConnectionManager
- Intégration avec le ControlPoint pour l'enregistrement dynamique des renderers
- Gestion des sessions avec timeout et cleanup automatique
- Support des commandes de transport (play, pause, stop, seek) et du contrôle du volume
- Mise à jour des dépendances dans Cargo.toml et Cargo.lock
- Documentation de l'architecture dans Blackboard/Architecture/webrenderer.md
2026-02-19 16:15:18 +01:00
38 changed files with 3207 additions and 323 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

@@ -5,7 +5,7 @@ Parfait. Voici un **schéma fonctionnel minimal** pour un **MediaRenderer UPnP p
- 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.
Tu mettras une version du plan en Markdown dans le répertoire [@Architecture](file:///Users/coissac/Sync/maison/Petite_maisons/src/pmomusic/Blackboard/Architecture) .
---
@@ -136,4 +136,3 @@ ws.onmessage = (evt) => {
---
💡 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

@@ -57,7 +57,6 @@ pub enum SsdpEvent {
#[derive(Clone)]
pub struct SsdpClient {
socket: Arc<UdpSocket>,
loopback_socket: Arc<UdpSocket>,
}
impl SsdpClient {
@@ -80,59 +79,43 @@ 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() {
// Join multicast on ALL interfaces, including loopback for local development
match socket.join_multicast_v4(&SSDP_MULTICAST_ADDR.parse().unwrap(), &ipv4) {
Ok(()) => {
if ipv4.is_loopback() {
debug!(
"SSDP: joined {} on {} (localhost - dev mode)",
SSDP_MULTICAST_ADDR, ipv4
);
} else {
if !ipv4.is_loopback() {
match socket.join_multicast_v4(&multicast_addr, &ipv4) {
Ok(()) => {
debug!("SSDP: joined {} on {}", SSDP_MULTICAST_ADDR, ipv4);
}
}
Err(e) => {
warn!(
"SSDP: failed to join {} on {}: {}",
SSDP_MULTICAST_ADDR, ipv4, e
);
Err(e) => {
warn!(
"SSDP: failed to join {} on {}: {}",
SSDP_MULTICAST_ADDR, ipv4, e
);
}
}
}
}
}
// Create a second socket for loopback multicast (when network blocks multicast)
let loopback_socket2 = Socket::new(Domain::IPV4, Type::DGRAM, Some(Protocol::UDP))?;
loopback_socket2.set_reuse_address(true)?;
// Bind on any address with ephemeral port (like the main socket)
let loopback_bind: SocketAddr = "0.0.0.0:0".parse().unwrap();
loopback_socket2.bind(&loopback_bind.into())?;
let loopback_socket: UdpSocket = loopback_socket2.into();
loopback_socket.set_multicast_loop_v4(true)?;
// Join multicast specifically on loopback interface (127.0.0.1)
let loopback_addr: std::net::Ipv4Addr = "127.0.0.1".parse().unwrap();
if let Err(e) =
loopback_socket.join_multicast_v4(&SSDP_MULTICAST_ADDR.parse().unwrap(), &loopback_addr)
{
warn!("Failed to join multicast on loopback socket: {}", e);
} else {
debug!(
"SSDP loopback socket: joined {} on 127.0.0.1",
SSDP_MULTICAST_ADDR
);
}
// 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 {
socket: Arc::new(socket),
loopback_socket: Arc::new(loopback_socket),
})
}
@@ -150,50 +133,19 @@ impl SsdpClient {
SSDP_MULTICAST_ADDR, SSDP_PORT, mx, st
);
let multicast_addr: SocketAddr = format!("{}:{}", SSDP_MULTICAST_ADDR, SSDP_PORT)
let addr: SocketAddr = format!("{}:{}", SSDP_MULTICAST_ADDR, SSDP_PORT)
.parse()
.unwrap();
// Try multicast first
match self.socket.send_to(msg.as_bytes(), multicast_addr) {
match self.socket.send_to(msg.as_bytes(), addr) {
Ok(_) => {
info!("📤 M-SEARCH sent to multicast (ST={}, MX={})", st, mx);
info!("📤 M-SEARCH sent (ST={}, MX={})", st, mx);
debug!(
"📨 M-SEARCH payload\n<details>\n\n```\n{}\n```\n</details>\n",
msg
);
Ok(())
}
Err(e) if e.raw_os_error() == Some(65) || e.raw_os_error() == Some(101) => {
// Error 65 (EHOSTUNREACH) on macOS, 101 (ENETUNREACH) on Linux
// Network blocks multicast (e.g., eduroam) → send to localhost unicast
warn!(
"⚠️ Multicast blocked on network ({}), sending to localhost for local devices",
e
);
// Send to localhost unicast (not multicast) - servers on 127.0.0.1:1900 will receive
let localhost_addr: SocketAddr =
format!("127.0.0.1:{}", SSDP_PORT).parse().unwrap();
match self.loopback_socket.send_to(msg.as_bytes(), localhost_addr) {
Ok(_) => {
info!(
"📤 M-SEARCH sent to localhost unicast (ST={}, MX={})",
st, mx
);
debug!(
"📨 M-SEARCH localhost payload\n<details>\n\n```\n{}\n```\n</details>\n",
msg
);
Ok(())
}
Err(loopback_err) => {
warn!("❌ Failed to send M-SEARCH to localhost: {}", loopback_err);
Err(loopback_err)
}
}
}
Err(e) => {
warn!("❌ Failed to send M-SEARCH: {}", e);
Err(e)

View File

@@ -15,9 +15,6 @@ pub struct SsdpServer {
/// Socket UDP pour SSDP
socket: Option<Arc<UdpSocket>>,
/// Socket dédié pour loopback (quand le réseau bloque le multicast)
loopback_socket: Option<Arc<UdpSocket>>,
}
impl SsdpServer {
@@ -26,7 +23,6 @@ impl SsdpServer {
Self {
devices: Arc::new(RwLock::new(HashMap::new())),
socket: None,
loopback_socket: None,
}
}
@@ -79,92 +75,40 @@ 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 sur toutes les interfaces (y compris loopback)
// Ceci est essentiel pour le développement local avec plusieurs instances
for iface in get_if_addrs::get_if_addrs()? {
if let std::net::IpAddr::V4(ipv4) = iface.ip() {
match socket.join_multicast_v4(&SSDP_MULTICAST_ADDR.parse().unwrap(), &ipv4) {
Ok(()) => {
if ipv4.is_loopback() {
debug!(
"SSDP server: joined {} on {} (localhost - dev mode)",
SSDP_MULTICAST_ADDR, ipv4
);
} else {
debug!("SSDP server: joined {} on {}", SSDP_MULTICAST_ADDR, ipv4);
}
}
Err(e) => {
warn!(
"SSDP server: failed to join {} on {}: {}",
SSDP_MULTICAST_ADDR, ipv4, e
);
}
}
}
// Rejoindre le groupe multicast
socket.join_multicast_v4(
&SSDP_MULTICAST_ADDR.parse().unwrap(),
&"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(true)?; // Important pour dev local
socket.set_multicast_loop_v4(false)?;
let socket = Arc::new(socket);
self.socket = Some(socket.clone());
// Create a second socket for loopback multicast (when network blocks multicast)
let loopback_socket2 = Socket::new(Domain::IPV4, Type::DGRAM, Some(Protocol::UDP))?;
loopback_socket2.set_reuse_address(true)?;
#[cfg(unix)]
{
use std::os::unix::io::AsRawFd;
let fd = loopback_socket2.as_raw_fd();
let optval: libc::c_int = 1;
unsafe {
let result = libc::setsockopt(
fd,
libc::SOL_SOCKET,
libc::SO_REUSEPORT,
&optval as *const _ as *const libc::c_void,
std::mem::size_of_val(&optval) as libc::socklen_t,
);
if result != 0 {
return Err(std::io::Error::last_os_error());
}
}
}
// Bind on any address with ephemeral port (not on port 1900 to avoid conflicts)
let loopback_bind: SocketAddr = "0.0.0.0:0".parse().unwrap();
loopback_socket2.bind(&loopback_bind.into())?;
let loopback_socket: UdpSocket = loopback_socket2.into();
loopback_socket.set_multicast_loop_v4(true)?;
// Join multicast specifically on loopback interface (127.0.0.1)
let loopback_addr: std::net::Ipv4Addr = "127.0.0.1".parse().unwrap();
if let Err(e) =
loopback_socket.join_multicast_v4(&SSDP_MULTICAST_ADDR.parse().unwrap(), &loopback_addr)
{
warn!(
"SSDP server: failed to join multicast on loopback socket: {}",
e
);
} else {
debug!(
"SSDP loopback server socket: joined {} on 127.0.0.1",
SSDP_MULTICAST_ADDR
);
}
let loopback_socket = Arc::new(loopback_socket);
self.loopback_socket = Some(loopback_socket.clone());
info!("✅ SSDP server started on {}", addr);
// Lancer les goroutines d'annonces périodiques et d'écoute M-SEARCH
self.start_periodic_announcements(socket.clone(), loopback_socket.clone());
self.start_periodic_announcements(socket.clone());
self.start_msearch_listener(socket.clone());
Ok(())
@@ -190,10 +134,9 @@ impl SsdpServer {
// Envoyer alive pour tous les NTs
if let Some(ref socket) = self.socket {
let loopback_socket = self.loopback_socket.as_ref();
let nts = device.get_notification_types();
for nt in nts.iter() {
Self::send_alive(socket, loopback_socket, &device, nt, false);
Self::send_alive(socket, &device, nt, false);
// Petit délai pour éviter de saturer le buffer UDP sur macOS
std::thread::sleep(Duration::from_millis(5));
}
@@ -222,13 +165,7 @@ impl SsdpServer {
}
/// Envoie un NOTIFY alive
fn send_alive(
socket: &UdpSocket,
loopback_socket: Option<&Arc<UdpSocket>>,
device: &SsdpDevice,
nt: &str,
is_periodic: bool,
) {
fn send_alive(socket: &UdpSocket, device: &SsdpDevice, nt: &str, is_periodic: bool) {
let usn = if nt.starts_with("uuid:") {
format!("{}", nt)
} else {
@@ -248,11 +185,11 @@ impl SsdpServer {
SSDP_MULTICAST_ADDR, SSDP_PORT, MAX_AGE, device.location, nt, device.server, usn
);
let multicast_addr: SocketAddr = format!("{}:{}", SSDP_MULTICAST_ADDR, SSDP_PORT)
let addr: SocketAddr = format!("{}:{}", SSDP_MULTICAST_ADDR, SSDP_PORT)
.parse()
.unwrap();
match socket.send_to(msg.as_bytes(), multicast_addr) {
match socket.send_to(msg.as_bytes(), addr) {
Ok(_) => {
let label = if is_periodic { " (periodic)" } else { "" };
info!("✅ NOTIFY alive{}: {} (NT={})", label, usn, nt);
@@ -261,34 +198,7 @@ impl SsdpServer {
label, msg
);
}
Err(e) if e.raw_os_error() == Some(65) || e.raw_os_error() == Some(101) => {
// Multicast bloqué sur le réseau, envoyer en unicast sur localhost
if let Some(loopback_sock) = loopback_socket {
// Send to localhost unicast - local clients will receive
let localhost_addr: SocketAddr =
format!("127.0.0.1:{}", SSDP_PORT).parse().unwrap();
match loopback_sock.send_to(msg.as_bytes(), localhost_addr) {
Ok(_) => {
let label = if is_periodic { " (periodic)" } else { "" };
info!("✅ NOTIFY alive{} to localhost: {} (NT={})", label, usn, nt);
}
Err(loopback_err) => {
let label = if is_periodic { "periodic " } else { "" };
warn!(
"❌ Failed to send {}NOTIFY alive to localhost for {}: {}",
label, usn, loopback_err
);
}
}
} else {
let label = if is_periodic { "periodic " } else { "" };
warn!(
"❌ Failed to send {}NOTIFY alive for {} (no loopback socket available): {}",
label, usn, e
);
}
}
Err(e) => {
let label = if is_periodic { "periodic " } else { "" };
warn!("❌ Failed to send {}NOTIFY alive for {}: {}", label, usn, e);
@@ -331,11 +241,7 @@ impl SsdpServer {
}
/// Démarre les annonces périodiques (toutes les MAX_AGE/2 secondes)
fn start_periodic_announcements(
&self,
socket: Arc<UdpSocket>,
loopback_socket: Arc<UdpSocket>,
) {
fn start_periodic_announcements(&self, socket: Arc<UdpSocket>) {
let devices = Arc::clone(&self.devices);
let period = Duration::from_secs((MAX_AGE / 2) as u64);
@@ -351,7 +257,7 @@ impl SsdpServer {
};
for device in &devices_snapshot {
for nt in device.get_notification_types() {
Self::send_alive(&socket, Some(&loopback_socket), device, nt, true);
Self::send_alive(&socket, device, nt, true);
}
}
}

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