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
This commit is contained in:
2026-02-09 07:18:13 +01:00
parent b987decb8d
commit 1e882ba6c3
19 changed files with 2028 additions and 215 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.

133
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"
@@ -4280,6 +4351,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 +5490,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 +6107,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 +6381,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

@@ -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

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 {
@@ -76,63 +75,32 @@ impl SsdpClient {
let bind_addr: SocketAddr = "0.0.0.0:0".parse().unwrap();
socket2.bind(&bind_addr.into())?;
// Le client SSDP n'a pas besoin de join_multicast_v4 :
// il envoie des M-SEARCH et reçoit des réponses unicast.
// Le join est seulement utile pour recevoir les NOTIFY multicast,
// ce dont le SsdpServer s'occupe.
//
// On configure uniquement l'interface de sortie multicast sur
// l'IP principale pour que send_to vers 239.255.255.250 fonctionne.
// Sur macOS, join_multicast_v4 sur des interfaces bridge/VM écrase
// IP_MULTICAST_IF et provoque EHOSTUNREACH (errno 65).
let local_ip: std::net::Ipv4Addr = pmoutils::guess_local_ip()
.parse()
.unwrap_or("127.0.0.1".parse().unwrap());
socket2.set_multicast_if_v4(&local_ip)?;
debug!(
"SSDP client: multicast outgoing interface set to {}",
local_ip
);
let socket: UdpSocket = socket2.into();
socket.set_read_timeout(Some(Duration::from_secs(1)))?;
socket.set_multicast_loop_v4(true)?; // utile en dev local
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 {
debug!("SSDP: joined {} on {}", SSDP_MULTICAST_ADDR, ipv4);
}
}
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
);
}
info!("✅ SSDP client ready on {}", addr);
Ok(Self {
socket: Arc::new(socket),
loopback_socket: Arc::new(loopback_socket),
})
}
@@ -150,50 +118,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,
}
}
@@ -78,93 +74,38 @@ impl SsdpServer {
let bind_addr: SocketAddr = format!("0.0.0.0:{}", SSDP_PORT).parse().unwrap();
socket2.bind(&bind_addr.into())?;
// Configurer l'interface de sortie multicast sur l'IP principale.
// Sur macOS, sans cela le kernel peut router les paquets multicast
// via une interface bridge/VM qui n'a pas de route, causant
// EHOSTUNREACH (errno 65).
let local_ip: std::net::Ipv4Addr = pmoutils::guess_local_ip()
.parse()
.unwrap_or("127.0.0.1".parse().unwrap());
socket2.set_multicast_if_v4(&local_ip)?;
debug!(
"SSDP server: multicast outgoing interface set to {}",
local_ip
);
// Convertir en UdpSocket standard
let 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 via INADDR_ANY (l'OS choisit l'interface)
socket.join_multicast_v4(
&SSDP_MULTICAST_ADDR.parse().unwrap(),
&"0.0.0.0".parse().unwrap(),
)?;
socket.set_read_timeout(Some(Duration::from_secs(1)))?;
socket.set_multicast_loop_v4(true)?; // Important pour dev local
socket.set_multicast_loop_v4(true)?; // utile en dev local
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 +131,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 +162,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 +182,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 +195,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 +238,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 +254,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,316 @@
//! 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 tokio::sync::mpsc;
use pmoupnp::actions::{ActionData, ActionError, ActionHandler};
use pmoupnp::{get, set};
use crate::messages::{CommandParams, PlaybackState, ServerMessage, TransportAction};
use crate::state::SharedState;
type ActionFuture =
std::pin::Pin<Box<dyn std::future::Future<Output = Result<ActionData, ActionError>> + Send>>;
// ─── AVTransport Handlers ───────────────────────────────────────────────────
pub fn play_handler(ws: mpsc::UnboundedSender<ServerMessage>, state: SharedState) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture {
let ws = ws.clone();
let state = state.clone();
Box::pin(async move {
let _ = ws.send(ServerMessage::Command {
action: TransportAction::Play,
params: None,
});
state.write().playback_state = PlaybackState::Playing;
Ok(data)
})
})
}
pub fn stop_handler(ws: mpsc::UnboundedSender<ServerMessage>, state: SharedState) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture {
let ws = ws.clone();
let state = state.clone();
Box::pin(async move {
let _ = ws.send(ServerMessage::Command {
action: TransportAction::Stop,
params: None,
});
state.write().playback_state = PlaybackState::Stopped;
Ok(data)
})
})
}
pub fn pause_handler(
ws: mpsc::UnboundedSender<ServerMessage>,
state: SharedState,
) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture {
let ws = ws.clone();
let state = state.clone();
Box::pin(async move {
let _ = ws.send(ServerMessage::Command {
action: TransportAction::Pause,
params: None,
});
state.write().playback_state = PlaybackState::Paused;
Ok(data)
})
})
}
pub fn next_handler(ws: mpsc::UnboundedSender<ServerMessage>) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture {
let ws = ws.clone();
Box::pin(async move {
let _ = ws.send(ServerMessage::Command {
action: TransportAction::Play,
params: None,
});
Ok(data)
})
})
}
pub fn previous_handler(ws: mpsc::UnboundedSender<ServerMessage>) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture {
let ws = ws.clone();
Box::pin(async move {
let _ = ws.send(ServerMessage::Command {
action: TransportAction::Play,
params: None,
});
Ok(data)
})
})
}
pub fn seek_handler(ws: mpsc::UnboundedSender<ServerMessage>) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture {
let ws = ws.clone();
Box::pin(async move {
let target: String = get!(&data, "Target", String);
let _ = 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: mpsc::UnboundedSender<ServerMessage>,
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!(&data, "CurrentURIMetaData", String);
let _ = 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(_state: SharedState) -> ActionHandler {
Arc::new(move |data: ActionData| -> ActionFuture {
Box::pin(async move {
let _uri: String = get!(&data, "NextURI", String);
let _metadata: String = get!(&data, "NextURIMetaData", String);
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();
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());
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", String::new());
set!(&mut data, "NextURIMetaData", String::new());
Ok(data)
})
})
}
// ─── RenderingControl Handlers ──────────────────────────────────────────────
pub fn set_volume_handler(
ws: mpsc::UnboundedSender<ServerMessage>,
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);
let _ = 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: mpsc::UnboundedSender<ServerMessage>,
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);
let _ = 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,89 @@
//! 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,
},
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,
}
#[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 },
Pong,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct BrowserCapabilities {
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,621 @@
//! 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::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),
}
/// 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.
///
/// Le device est construit avec des action handlers qui relaient
/// les commandes SOAP vers le navigateur via le `ws_sender`.
pub fn create_device(
browser_name: &str,
ws_sender: mpsc::UnboundedSender<ServerMessage>,
state: SharedState,
) -> Result<Device, FactoryError> {
let avtransport = Self::build_avtransport(ws_sender.clone(), state.clone())?;
let renderingcontrol = Self::build_renderingcontrol(ws_sender.clone(), state.clone())?;
let connectionmanager = Self::build_connectionmanager()?;
let device = Device::new(
format!("WebRenderer"),
"MediaRenderer".to_string(),
format!("Web Audio - {}", browser_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)
}
/// Construit le service AVTransport avec les handlers WebSocket
fn build_avtransport(
ws: mpsc::UnboundedSender<ServerMessage>,
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(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.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: mpsc::UnboundedSender<ServerMessage>,
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,105 @@
//! Gestionnaire de sessions WebRenderer
use parking_lot::RwLock;
use std::collections::HashMap;
use std::sync::Arc;
use std::time::{Duration, SystemTime};
use tokio::sync::mpsc;
use pmoupnp::devices::DeviceInstance;
use crate::messages::ServerMessage;
use crate::state::SharedState;
/// Session WebSocket liée à un MediaRenderer privé
pub struct WebRendererSession {
pub token: String,
pub udn: String,
pub device_instance: Arc<DeviceInstance>,
pub ws_sender: mpsc::UnboundedSender<ServerMessage>,
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>>>>,
timeout_duration: Duration,
}
impl SessionManager {
pub fn new(timeout_duration: Duration) -> Self {
let manager = Self {
sessions: 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();
self.sessions.write().insert(token.clone(), session);
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
}
}
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,35 @@
//! État partagé du renderer (backend ↔ navigateur)
use parking_lot::RwLock;
use std::sync::Arc;
use crate::messages::PlaybackState;
/// É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 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,
position: None,
duration: None,
volume: 100,
mute: false,
}
}
}
/// Alias pour l'état partagé
pub type SharedState = Arc<RwLock<RendererState>>;

View File

@@ -0,0 +1,424 @@
//! 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 {
ws.on_upgrade(move |socket: WebSocket| handle_socket(socket, state))
}
/// Gestion de la connexion WebSocket
async fn handle_socket(socket: WebSocket, state: WebSocketState) {
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;
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)) => {
match serde_json::from_str::<ClientMessage>(&text) {
Ok(ClientMessage::Init { capabilities }) => {
// Créer le renderer UPnP pour ce navigateur
match create_renderer_for_browser(&capabilities, tx.clone(), &state).await {
Ok(session) => {
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()
),
},
});
session_token = Some(token.clone());
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::Pong) => {}
Err(e) => {
tracing::warn!(error = %e, "Failed to parse client message");
}
}
}
Ok(Message::Close(_)) => break,
Err(e) => {
tracing::error!("WebSocket error: {}", e);
break;
}
_ => {}
}
}
// Cleanup à la déconnexion
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 un DeviceInstance UPnP et l'enregistre pour un navigateur
async fn create_renderer_for_browser(
capabilities: &BrowserCapabilities,
ws_sender: mpsc::UnboundedSender<ServerMessage>,
ws_state: &WebSocketState,
) -> Result<Arc<WebRendererSession>, crate::error::WebRendererError> {
let shared_state: SharedState = Arc::new(RwLock::new(RendererState::default()));
let token = Uuid::new_v4().to_string();
// Construire le Device model avec les handlers WS
let device = WebRendererFactory::create_device(
&capabilities.user_agent,
ws_sender.clone(),
shared_state.clone(),
)
.map_err(|e| crate::error::WebRendererError::DeviceCreationError(e.to_string()))?;
let device = Arc::new(device);
// Enregistrer le device via UpnpServerExt (gère base_url, register_urls, DEVICE_REGISTRY)
#[cfg(feature = "pmoserver")]
let di = {
use pmoupnp::UpnpServerExt;
let server_arc =
pmoserver::get_server().ok_or(crate::error::WebRendererError::ServerNotAvailable)?;
let di = {
let mut server = server_arc.write().await;
server
.register_device(device)
.await
.map_err(|e| crate::error::WebRendererError::RegistrationError(e.to_string()))?
};
// Enregistrer avec le ControlPoint
register_with_control_point(&di, ws_state)?;
di
};
#[cfg(not(feature = "pmoserver"))]
let di = device.create_instance();
let udn = di.udn().to_string();
let session = Arc::new(WebRendererSession {
token,
udn,
device_instance: di,
ws_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();
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.clone()),
udn.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_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()
}