Découverte upnp
This commit is contained in:
@@ -130,8 +130,13 @@ impl ControlPoint {
|
||||
// SsdpClient
|
||||
let client = SsdpClient::new()?; // pmoupnp::ssdp::SsdpClient
|
||||
|
||||
// Clone pour le thread de renouvellement périodique
|
||||
let client_for_renewal = client.clone();
|
||||
|
||||
// Arc utilisé dans le thread
|
||||
let registry_for_thread = Arc::clone(®istry);
|
||||
let event_bus_for_discovery = event_bus.clone();
|
||||
let media_event_bus_for_discovery = media_event_bus.clone();
|
||||
|
||||
// Thread de découverte
|
||||
thread::spawn(move || {
|
||||
@@ -167,12 +172,177 @@ impl ControlPoint {
|
||||
|
||||
if let Ok(mut reg) = registry_for_thread.write() {
|
||||
for update in updates {
|
||||
// Émettre les événements Online/Offline avant d'appliquer l'update
|
||||
match &update {
|
||||
DeviceUpdate::RendererOnline(info) => {
|
||||
event_bus_for_discovery.broadcast(RendererEvent::Online {
|
||||
id: info.id.clone(),
|
||||
info: info.clone(),
|
||||
});
|
||||
}
|
||||
DeviceUpdate::RendererOfflineById(id) => {
|
||||
event_bus_for_discovery.broadcast(RendererEvent::Offline {
|
||||
id: id.clone(),
|
||||
});
|
||||
}
|
||||
DeviceUpdate::RendererOfflineByUdn(udn) => {
|
||||
// Trouver l'ID avant de marquer offline
|
||||
if let Some(renderer) = reg.get_renderer_by_udn(udn) {
|
||||
event_bus_for_discovery.broadcast(RendererEvent::Offline {
|
||||
id: renderer.id.clone(),
|
||||
});
|
||||
}
|
||||
}
|
||||
DeviceUpdate::ServerOnline(info) => {
|
||||
media_event_bus_for_discovery.broadcast(MediaServerEvent::Online {
|
||||
server_id: info.id.clone(),
|
||||
info: info.clone(),
|
||||
});
|
||||
}
|
||||
DeviceUpdate::ServerOfflineById(id) => {
|
||||
media_event_bus_for_discovery.broadcast(MediaServerEvent::Offline {
|
||||
server_id: id.clone(),
|
||||
});
|
||||
}
|
||||
DeviceUpdate::ServerOfflineByUdn(udn) => {
|
||||
// Trouver l'ID avant de marquer offline
|
||||
if let Some(server) = reg.get_server_by_udn(udn) {
|
||||
media_event_bus_for_discovery.broadcast(MediaServerEvent::Offline {
|
||||
server_id: server.id.clone(),
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
reg.apply_update(update);
|
||||
}
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
// Thread de renouvellement périodique des M-SEARCH
|
||||
// Envoie des requêtes de découverte toutes les 60 secondes pour forcer
|
||||
// les nouveaux appareils à se présenter
|
||||
thread::spawn(move || {
|
||||
let search_targets = [
|
||||
"ssdp:all",
|
||||
"urn:schemas-upnp-org:device:MediaRenderer:1",
|
||||
"urn:av-openhome-org:device:MediaRenderer:1",
|
||||
"urn:schemas-upnp-org:device:MediaServer:1",
|
||||
"urn:schemas-wiimu-com:service:PlayQueue:1",
|
||||
];
|
||||
|
||||
loop {
|
||||
// Attendre 60 secondes avant le prochain cycle
|
||||
thread::sleep(Duration::from_secs(60));
|
||||
|
||||
debug!("Sending periodic M-SEARCH for device discovery");
|
||||
|
||||
// Envoyer les M-SEARCH
|
||||
for st in &search_targets {
|
||||
if let Err(e) = client_for_renewal.send_msearch(st, 3) {
|
||||
warn!("Failed to send periodic M-SEARCH for {}: {}", st, e);
|
||||
}
|
||||
thread::sleep(Duration::from_millis(200));
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
// Thread de vérification de présence périodique
|
||||
// Vérifie toutes les 60 secondes que les devices connus sont toujours accessibles
|
||||
let registry_for_presence = Arc::clone(®istry);
|
||||
let event_bus_for_presence = event_bus.clone();
|
||||
let media_event_bus_for_presence = media_event_bus.clone();
|
||||
thread::spawn(move || {
|
||||
use ureq::Agent;
|
||||
|
||||
// HTTP client avec timeout court pour les vérifications de présence
|
||||
let config = Agent::config_builder()
|
||||
.timeout_global(Some(Duration::from_secs(5)))
|
||||
.build();
|
||||
let agent: Agent = config.into();
|
||||
|
||||
loop {
|
||||
// Attendre 60 secondes avant le prochain cycle
|
||||
thread::sleep(Duration::from_secs(60));
|
||||
|
||||
debug!("Starting periodic presence check for devices");
|
||||
|
||||
let mut updates = Vec::new();
|
||||
|
||||
// Lire la liste des devices
|
||||
if let Ok(reg) = registry_for_presence.read() {
|
||||
// Vérifier les renderers
|
||||
for renderer in reg.list_renderers() {
|
||||
if !renderer.online {
|
||||
continue; // Skip déjà offline
|
||||
}
|
||||
|
||||
// Faire un HTTP HEAD pour vérifier la présence
|
||||
match agent.head(&renderer.location).call() {
|
||||
Ok(_) => {
|
||||
// Device répond toujours
|
||||
debug!("Renderer {} ({:?}) is still online",
|
||||
renderer.friendly_name, renderer.id);
|
||||
}
|
||||
Err(e) => {
|
||||
// Device ne répond plus
|
||||
warn!("Renderer {} ({:?}) is no longer responding: {} - marking offline",
|
||||
renderer.friendly_name, renderer.id, e);
|
||||
updates.push(DeviceUpdate::RendererOfflineById(renderer.id));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Vérifier les servers
|
||||
for server in reg.list_servers() {
|
||||
if !server.online {
|
||||
continue; // Skip déjà offline
|
||||
}
|
||||
|
||||
// Faire un HTTP HEAD pour vérifier la présence
|
||||
match agent.head(&server.location).call() {
|
||||
Ok(_) => {
|
||||
// Device répond toujours
|
||||
debug!("Server {} ({:?}) is still online",
|
||||
server.friendly_name, server.id);
|
||||
}
|
||||
Err(e) => {
|
||||
// Device ne répond plus
|
||||
warn!("Server {} ({:?}) is no longer responding: {} - marking offline",
|
||||
server.friendly_name, server.id, e);
|
||||
updates.push(DeviceUpdate::ServerOfflineById(server.id));
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Appliquer les updates et émettre les événements
|
||||
if !updates.is_empty() {
|
||||
if let Ok(mut reg) = registry_for_presence.write() {
|
||||
for update in updates {
|
||||
// Émettre les événements Offline
|
||||
match &update {
|
||||
DeviceUpdate::RendererOfflineById(id) => {
|
||||
event_bus_for_presence.broadcast(RendererEvent::Offline {
|
||||
id: id.clone(),
|
||||
});
|
||||
}
|
||||
DeviceUpdate::ServerOfflineById(id) => {
|
||||
media_event_bus_for_presence.broadcast(MediaServerEvent::Offline {
|
||||
server_id: id.clone(),
|
||||
});
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
|
||||
reg.apply_update(update);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
// Thread de découverte mDNS pour Chromecast
|
||||
let registry_for_mdns = Arc::clone(®istry);
|
||||
thread::spawn(move || {
|
||||
@@ -555,6 +725,19 @@ impl ControlPoint {
|
||||
}
|
||||
}
|
||||
}
|
||||
MediaServerEvent::Online { server_id, info } => {
|
||||
debug!(
|
||||
server = server_id.0.as_str(),
|
||||
friendly_name = info.friendly_name.as_str(),
|
||||
"MediaServer came online"
|
||||
);
|
||||
}
|
||||
MediaServerEvent::Offline { server_id } => {
|
||||
debug!(
|
||||
server = server_id.0.as_str(),
|
||||
"MediaServer went offline"
|
||||
);
|
||||
}
|
||||
}
|
||||
})?;
|
||||
|
||||
|
||||
@@ -541,6 +541,9 @@ impl MediaServerEventWorker {
|
||||
"Broadcasting MediaServerEvent::ContainersUpdated"
|
||||
);
|
||||
}
|
||||
MediaServerEvent::Online { .. } | MediaServerEvent::Offline { .. } => {
|
||||
// Online/Offline events are generated from SSDP discovery, not from notify payloads
|
||||
}
|
||||
}
|
||||
self.bus.broadcast(event);
|
||||
}
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
use crate::capabilities::{PlaybackPositionInfo, PlaybackState};
|
||||
use crate::control_point::PlaylistBinding;
|
||||
use crate::media_server::ServerId;
|
||||
use crate::media_server::{MediaServerInfo, ServerId};
|
||||
|
||||
#[derive(Clone, Debug, PartialEq, Eq, Hash)]
|
||||
pub struct RendererId(pub String);
|
||||
@@ -122,6 +122,13 @@ pub enum RendererEvent {
|
||||
id: RendererId,
|
||||
binding: Option<PlaylistBinding>,
|
||||
},
|
||||
Online {
|
||||
id: RendererId,
|
||||
info: RendererInfo,
|
||||
},
|
||||
Offline {
|
||||
id: RendererId,
|
||||
},
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug)]
|
||||
@@ -134,4 +141,11 @@ pub enum MediaServerEvent {
|
||||
server_id: ServerId,
|
||||
container_ids: Vec<String>,
|
||||
},
|
||||
Online {
|
||||
server_id: ServerId,
|
||||
info: MediaServerInfo,
|
||||
},
|
||||
Offline {
|
||||
server_id: ServerId,
|
||||
},
|
||||
}
|
||||
|
||||
@@ -133,6 +133,15 @@ impl DeviceRegistry {
|
||||
}
|
||||
}
|
||||
|
||||
/// Helper: get a server by UDN (case-insensitive, via udn_index).
|
||||
pub fn get_server_by_udn(&self, udn: &str) -> Option<MediaServerInfo> {
|
||||
let lookup = udn.to_ascii_lowercase();
|
||||
match self.udn_index.get(&lookup) {
|
||||
Some(DeviceKey::Server(id)) => self.servers.get(id).cloned(),
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
|
||||
/// Construct an AvTransportClient for a given renderer id, if possible.
|
||||
///
|
||||
/// Returns:
|
||||
|
||||
@@ -82,6 +82,17 @@ pub enum RendererEventPayload {
|
||||
container_id: Option<String>,
|
||||
timestamp: chrono::DateTime<chrono::Utc>,
|
||||
},
|
||||
Online {
|
||||
renderer_id: String,
|
||||
friendly_name: String,
|
||||
model_name: String,
|
||||
manufacturer: String,
|
||||
timestamp: chrono::DateTime<chrono::Utc>,
|
||||
},
|
||||
Offline {
|
||||
renderer_id: String,
|
||||
timestamp: chrono::DateTime<chrono::Utc>,
|
||||
},
|
||||
}
|
||||
|
||||
/// Payload SSE pour un événement serveur de médias
|
||||
@@ -99,6 +110,17 @@ pub enum MediaServerEventPayload {
|
||||
container_ids: Vec<String>,
|
||||
timestamp: chrono::DateTime<chrono::Utc>,
|
||||
},
|
||||
Online {
|
||||
server_id: String,
|
||||
friendly_name: String,
|
||||
model_name: String,
|
||||
manufacturer: String,
|
||||
timestamp: chrono::DateTime<chrono::Utc>,
|
||||
},
|
||||
Offline {
|
||||
server_id: String,
|
||||
timestamp: chrono::DateTime<chrono::Utc>,
|
||||
},
|
||||
}
|
||||
|
||||
/// Payload SSE unifié pour tous les événements
|
||||
@@ -205,6 +227,21 @@ pub async fn renderer_events_sse(
|
||||
timestamp,
|
||||
}
|
||||
}
|
||||
RendererEvent::Online { id, info } => {
|
||||
RendererEventPayload::Online {
|
||||
renderer_id: id.0,
|
||||
friendly_name: info.friendly_name,
|
||||
model_name: info.model_name,
|
||||
manufacturer: info.manufacturer,
|
||||
timestamp,
|
||||
}
|
||||
}
|
||||
RendererEvent::Offline { id } => {
|
||||
RendererEventPayload::Offline {
|
||||
renderer_id: id.0,
|
||||
timestamp,
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
if let Ok(json) = serde_json::to_string(&payload) {
|
||||
@@ -266,6 +303,21 @@ pub async fn media_server_events_sse(
|
||||
timestamp,
|
||||
}
|
||||
}
|
||||
MediaServerEvent::Online { server_id, info } => {
|
||||
MediaServerEventPayload::Online {
|
||||
server_id: server_id.0,
|
||||
friendly_name: info.friendly_name,
|
||||
model_name: info.model_name,
|
||||
manufacturer: info.manufacturer,
|
||||
timestamp,
|
||||
}
|
||||
}
|
||||
MediaServerEvent::Offline { server_id } => {
|
||||
MediaServerEventPayload::Offline {
|
||||
server_id: server_id.0,
|
||||
timestamp,
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
if let Ok(json) = serde_json::to_string(&payload) {
|
||||
@@ -379,6 +431,21 @@ pub async fn all_events_sse(State(control_point): State<Arc<ControlPoint>>) -> i
|
||||
timestamp,
|
||||
}
|
||||
}
|
||||
RendererEvent::Online { id, info } => {
|
||||
RendererEventPayload::Online {
|
||||
renderer_id: id.0,
|
||||
friendly_name: info.friendly_name,
|
||||
model_name: info.model_name,
|
||||
manufacturer: info.manufacturer,
|
||||
timestamp,
|
||||
}
|
||||
}
|
||||
RendererEvent::Offline { id } => {
|
||||
RendererEventPayload::Offline {
|
||||
renderer_id: id.0,
|
||||
timestamp,
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
let payload = UnifiedEventPayload::Renderer(renderer_payload);
|
||||
@@ -405,6 +472,21 @@ pub async fn all_events_sse(State(control_point): State<Arc<ControlPoint>>) -> i
|
||||
timestamp,
|
||||
}
|
||||
}
|
||||
MediaServerEvent::Online { server_id, info } => {
|
||||
MediaServerEventPayload::Online {
|
||||
server_id: server_id.0,
|
||||
friendly_name: info.friendly_name,
|
||||
model_name: info.model_name,
|
||||
manufacturer: info.manufacturer,
|
||||
timestamp,
|
||||
}
|
||||
}
|
||||
MediaServerEvent::Offline { server_id } => {
|
||||
MediaServerEventPayload::Offline {
|
||||
server_id: server_id.0,
|
||||
timestamp,
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
let payload = UnifiedEventPayload::MediaServer(server_payload);
|
||||
|
||||
@@ -54,6 +54,7 @@ pub enum SsdpEvent {
|
||||
}
|
||||
|
||||
/// Client SSDP pour envoyer des M-SEARCH et écouter les annonces
|
||||
#[derive(Clone)]
|
||||
pub struct SsdpClient {
|
||||
socket: Arc<UdpSocket>,
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user