Optimize UPnP discovery and server response handling
Refactor UPnP discovery to use a thread pool for fetching device descriptions, improving performance and responsiveness. Also, use spawn_blocking in server endpoints to prevent blocking the async runtime when performing sync operations. Additionally, add read timeout to SSDP socket and handle WouldBlock errors gracefully. Refactorisation du contrôle des renderers et amélioration de l'interface utilisateur Refactorisation complète du composant VolumeControl avec gestion d'erreur améliorée et debounce. Migration de l'interface de contrôle des renderers : - Suppression de l'ancienne barre d'onglets en bas - Intégration d'une nouvelle barre d'infos en bas avec les détails du renderer actif - Création d'un nouveau drawer pour la sélection des renderers - Déplacement des informations du renderer de l'onglet vers la barre d'infos Améliorations UI/UX : - Nouvelle interface de contrôle des renderers dans le drawer avec boutons de lecture/pause - Mise à jour des styles et animations pour une meilleure expérience utilisateur - Adaptation responsive pour les appareils mobiles
This commit is contained in:
@@ -1,14 +1,25 @@
|
||||
use crate::{DeviceRegistry, discovery::upnp_provider::ParsedDeviceDescription};
|
||||
use crossbeam_channel::{Sender, bounded};
|
||||
use pmoupnp::ssdp::SsdpEvent;
|
||||
use std::sync::{Arc, Mutex, RwLock};
|
||||
use std::thread;
|
||||
|
||||
use crate::discovery::manager::UDNRegistry;
|
||||
|
||||
/// Gestionnaire des événements SSDP -> DeviceUpdate.
|
||||
/// Task to fetch a device description
|
||||
struct FetchTask {
|
||||
udn: String,
|
||||
location: String,
|
||||
server_header: String,
|
||||
max_age: u32,
|
||||
registry: Arc<RwLock<DeviceRegistry>>,
|
||||
}
|
||||
|
||||
/// Gestionnaire des événements SSDP -> DeviceUpdate.
|
||||
pub struct UpnpDiscoveryManager {
|
||||
device_registry: Arc<RwLock<DeviceRegistry>>,
|
||||
udn_cache: Arc<Mutex<UDNRegistry>>,
|
||||
fetch_sender: Sender<FetchTask>,
|
||||
}
|
||||
|
||||
impl UpnpDiscoveryManager {
|
||||
@@ -16,9 +27,39 @@ impl UpnpDiscoveryManager {
|
||||
device_registry: Arc<RwLock<DeviceRegistry>>,
|
||||
udn_cache: Arc<Mutex<UDNRegistry>>,
|
||||
) -> Self {
|
||||
// Create a bounded channel for fetch tasks (max 10 pending tasks)
|
||||
let (sender, receiver) = bounded::<FetchTask>(10);
|
||||
|
||||
// Spawn a pool of 3 worker threads to process fetch tasks
|
||||
for _ in 0..3 {
|
||||
let receiver = receiver.clone();
|
||||
thread::spawn(move || {
|
||||
while let Ok(task) = receiver.recv() {
|
||||
// Fetch + parse the device description (may take up to 5 seconds)
|
||||
if let Ok(info) = ParsedDeviceDescription::new(
|
||||
&task.udn,
|
||||
&task.location,
|
||||
&task.server_header,
|
||||
5,
|
||||
) {
|
||||
if let Some(renderer_info) = info.build_renderer() {
|
||||
if let Ok(mut reg) = task.registry.write() {
|
||||
reg.push_renderer(&renderer_info, task.max_age);
|
||||
}
|
||||
} else if let Some(server_info) = info.build_server() {
|
||||
if let Ok(mut reg) = task.registry.write() {
|
||||
reg.push_server(&server_info, task.max_age);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
Self {
|
||||
device_registry,
|
||||
udn_cache,
|
||||
fetch_sender: sender,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -48,20 +89,18 @@ impl UpnpDiscoveryManager {
|
||||
UDNRegistry::should_fetch(self.udn_cache.clone(), &udn, max_age as u64);
|
||||
|
||||
if should_fetch {
|
||||
// Fetch + parse the device description
|
||||
if let Ok(info) =
|
||||
ParsedDeviceDescription::new(&udn, &location, &server_header, 5)
|
||||
{
|
||||
if let Some(renderer_info) = info.build_renderer() {
|
||||
if let Ok(mut reg) = self.device_registry.write() {
|
||||
reg.push_renderer(&renderer_info, max_age);
|
||||
}
|
||||
} else if let Some(server_info) = info.build_server() {
|
||||
if let Ok(mut reg) = self.device_registry.write() {
|
||||
reg.push_server(&server_info, max_age);
|
||||
}
|
||||
}
|
||||
}
|
||||
// Send fetch task to worker pool (non-blocking)
|
||||
// If the channel is full, try_send will fail and we skip this fetch
|
||||
let task = FetchTask {
|
||||
udn: udn.clone(),
|
||||
location: location.clone(),
|
||||
server_header: server_header.clone(),
|
||||
max_age,
|
||||
registry: Arc::clone(&self.device_registry),
|
||||
};
|
||||
|
||||
// Use try_send to avoid blocking if the queue is full
|
||||
let _ = self.fetch_sender.try_send(task);
|
||||
} else {
|
||||
// Even if we don't fetch, we MUST update last_seen to prevent timeout
|
||||
// This is critical: SSDP Alive messages arrive more frequently than max_age/2,
|
||||
|
||||
@@ -91,22 +91,29 @@ impl ControlPointState {
|
||||
tag = "control"
|
||||
)]
|
||||
async fn list_renderers(State(state): State<ControlPointState>) -> Json<Vec<RendererSummary>> {
|
||||
let renderers = state.control_point.list_music_renderers();
|
||||
// Use spawn_blocking to avoid blocking the tokio runtime
|
||||
// This is critical because list_music_renderers acquires a RwLock
|
||||
let control_point = state.control_point.clone();
|
||||
let summaries = tokio::task::spawn_blocking(move || {
|
||||
let renderers = control_point.list_music_renderers();
|
||||
|
||||
let summaries: Vec<RendererSummary> = renderers
|
||||
.into_iter()
|
||||
.map(|r| {
|
||||
let info = r.info();
|
||||
RendererSummary {
|
||||
id: r.id().0.clone(),
|
||||
friendly_name: r.friendly_name().to_string(),
|
||||
model_name: r.model_name().to_string(),
|
||||
protocol: protocol_summary(&info.protocol()),
|
||||
capabilities: capability_summary(&info.capabilities()),
|
||||
online: r.is_online(),
|
||||
}
|
||||
})
|
||||
.collect();
|
||||
renderers
|
||||
.into_iter()
|
||||
.map(|r| {
|
||||
let info = r.info();
|
||||
RendererSummary {
|
||||
id: r.id().0.clone(),
|
||||
friendly_name: r.friendly_name().to_string(),
|
||||
model_name: r.model_name().to_string(),
|
||||
protocol: protocol_summary(&info.protocol()),
|
||||
capabilities: capability_summary(&info.capabilities()),
|
||||
online: r.is_online(),
|
||||
}
|
||||
})
|
||||
.collect::<Vec<_>>()
|
||||
})
|
||||
.await
|
||||
.unwrap_or_default();
|
||||
|
||||
Json(summaries)
|
||||
}
|
||||
@@ -156,10 +163,22 @@ async fn get_renderer_full_snapshot(
|
||||
Path(renderer_id): Path<String>,
|
||||
) -> Result<Json<FullRendererSnapshot>, (StatusCode, Json<ErrorResponse>)> {
|
||||
let rid = DeviceId(renderer_id.clone());
|
||||
let snapshot = state
|
||||
.control_point
|
||||
.renderer_full_snapshot(&rid)
|
||||
.map_err(|err| map_snapshot_error(renderer_id, err))?;
|
||||
|
||||
// Use spawn_blocking because renderer_full_snapshot does sync UPnP calls
|
||||
let control_point = state.control_point.clone();
|
||||
let rid_clone = rid.clone();
|
||||
let snapshot =
|
||||
tokio::task::spawn_blocking(move || control_point.renderer_full_snapshot(&rid_clone))
|
||||
.await
|
||||
.map_err(|e| {
|
||||
(
|
||||
StatusCode::INTERNAL_SERVER_ERROR,
|
||||
Json(ErrorResponse {
|
||||
error: format!("Task error: {}", e),
|
||||
}),
|
||||
)
|
||||
})?
|
||||
.map_err(|err| map_snapshot_error(renderer_id, err))?;
|
||||
|
||||
Ok(Json(snapshot))
|
||||
}
|
||||
@@ -1519,17 +1538,23 @@ async fn add_after_current(
|
||||
tag = "control"
|
||||
)]
|
||||
async fn list_servers(State(state): State<ControlPointState>) -> Json<Vec<MediaServerSummary>> {
|
||||
let servers = state.control_point.list_media_servers().unwrap_or_default();
|
||||
// Use spawn_blocking to avoid blocking the tokio runtime
|
||||
let control_point = state.control_point.clone();
|
||||
let summaries = tokio::task::spawn_blocking(move || {
|
||||
let servers = control_point.list_media_servers().unwrap_or_default();
|
||||
|
||||
let summaries: Vec<MediaServerSummary> = servers
|
||||
.into_iter()
|
||||
.map(|s| MediaServerSummary {
|
||||
id: s.id().0.clone(),
|
||||
friendly_name: s.friendly_name().to_string(),
|
||||
model_name: s.model_name().to_string(),
|
||||
online: s.is_online(),
|
||||
})
|
||||
.collect();
|
||||
servers
|
||||
.into_iter()
|
||||
.map(|s| MediaServerSummary {
|
||||
id: s.id().0.clone(),
|
||||
friendly_name: s.friendly_name().to_string(),
|
||||
model_name: s.model_name().to_string(),
|
||||
online: s.is_online(),
|
||||
})
|
||||
.collect::<Vec<_>>()
|
||||
})
|
||||
.await
|
||||
.unwrap_or_default();
|
||||
|
||||
Json(summaries)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user