diff --git a/pmocontrol/src/control_point.rs b/pmocontrol/src/control_point.rs index b4210ca8..11051e5e 100644 --- a/pmocontrol/src/control_point.rs +++ b/pmocontrol/src/control_point.rs @@ -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" + ); + } } })?; diff --git a/pmocontrol/src/media_server_events.rs b/pmocontrol/src/media_server_events.rs index 8c7e9b47..cd2d011d 100644 --- a/pmocontrol/src/media_server_events.rs +++ b/pmocontrol/src/media_server_events.rs @@ -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); } diff --git a/pmocontrol/src/model.rs b/pmocontrol/src/model.rs index ad315abb..0ad26c2f 100644 --- a/pmocontrol/src/model.rs +++ b/pmocontrol/src/model.rs @@ -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, }, + Online { + id: RendererId, + info: RendererInfo, + }, + Offline { + id: RendererId, + }, } #[derive(Clone, Debug)] @@ -134,4 +141,11 @@ pub enum MediaServerEvent { server_id: ServerId, container_ids: Vec, }, + Online { + server_id: ServerId, + info: MediaServerInfo, + }, + Offline { + server_id: ServerId, + }, } diff --git a/pmocontrol/src/registry.rs b/pmocontrol/src/registry.rs index c8e54eca..affea6f3 100644 --- a/pmocontrol/src/registry.rs +++ b/pmocontrol/src/registry.rs @@ -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 { + 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: diff --git a/pmocontrol/src/sse.rs b/pmocontrol/src/sse.rs index ddb103dd..06439a4c 100644 --- a/pmocontrol/src/sse.rs +++ b/pmocontrol/src/sse.rs @@ -82,6 +82,17 @@ pub enum RendererEventPayload { container_id: Option, timestamp: chrono::DateTime, }, + Online { + renderer_id: String, + friendly_name: String, + model_name: String, + manufacturer: String, + timestamp: chrono::DateTime, + }, + Offline { + renderer_id: String, + timestamp: chrono::DateTime, + }, } /// Payload SSE pour un événement serveur de médias @@ -99,6 +110,17 @@ pub enum MediaServerEventPayload { container_ids: Vec, timestamp: chrono::DateTime, }, + Online { + server_id: String, + friendly_name: String, + model_name: String, + manufacturer: String, + timestamp: chrono::DateTime, + }, + Offline { + server_id: String, + timestamp: chrono::DateTime, + }, } /// 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>) -> 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>) -> 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); diff --git a/pmoupnp/src/ssdp/client.rs b/pmoupnp/src/ssdp/client.rs index 8210b108..8fcb1b8f 100644 --- a/pmoupnp/src/ssdp/client.rs +++ b/pmoupnp/src/ssdp/client.rs @@ -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, }