Merge pull request 'debugage massif de la queue openhome' (#25) from push-tkukkpplqtql into main
All checks were successful
Build and Push Docker Image / build (push) Successful in 28m25s
All checks were successful
Build and Push Docker Image / build (push) Successful in 28m25s
Reviewed-on: #25
This commit was merged in pull request #25.
This commit is contained in:
@@ -36,9 +36,36 @@ function ensureSSEConnected() {
|
||||
const serverId = event.server_id
|
||||
|
||||
switch (event.type) {
|
||||
case 'global_updated':
|
||||
// Invalider tout le cache de ce serveur
|
||||
console.log(`[useMediaServers] GlobalUpdated pour ${serverId}`)
|
||||
case 'online':
|
||||
// Nouveau serveur découvert
|
||||
console.log(`[useMediaServers] Serveur ${serverId} (${event.friendly_name}) est maintenant en ligne`)
|
||||
|
||||
// Ajouter au cache avec les infos disponibles
|
||||
const server: MediaServerSummary = {
|
||||
id: serverId,
|
||||
friendly_name: event.friendly_name,
|
||||
model_name: event.model_name,
|
||||
online: true,
|
||||
}
|
||||
serversCache.value.set(serverId, server)
|
||||
|
||||
// Fetch la liste complète pour avoir les bonnes infos
|
||||
// On ne le fait pas ici car on pourrait déclencher trop de requêtes
|
||||
// La liste se mettra à jour au prochain refresh automatique
|
||||
break
|
||||
|
||||
case 'offline':
|
||||
// Serveur déconnecté
|
||||
console.log(`[useMediaServers] Serveur ${serverId} est maintenant hors ligne`)
|
||||
|
||||
// Marquer comme offline dans le cache
|
||||
const existingServer = serversCache.value.get(serverId)
|
||||
if (existingServer) {
|
||||
existingServer.online = false
|
||||
serversCache.value.set(serverId, existingServer)
|
||||
}
|
||||
|
||||
// Invalider tout le cache browse de ce serveur
|
||||
const keysToDelete: string[] = []
|
||||
browseCache.value.forEach((_, key) => {
|
||||
if (key.startsWith(serverId + '/')) {
|
||||
@@ -48,6 +75,18 @@ function ensureSSEConnected() {
|
||||
keysToDelete.forEach(key => browseCache.value.delete(key))
|
||||
break
|
||||
|
||||
case 'global_updated':
|
||||
// Invalider tout le cache de ce serveur
|
||||
console.log(`[useMediaServers] GlobalUpdated pour ${serverId}`)
|
||||
const globalKeysToDelete: string[] = []
|
||||
browseCache.value.forEach((_, key) => {
|
||||
if (key.startsWith(serverId + '/')) {
|
||||
globalKeysToDelete.push(key)
|
||||
}
|
||||
})
|
||||
globalKeysToDelete.forEach(key => browseCache.value.delete(key))
|
||||
break
|
||||
|
||||
case 'containers_updated':
|
||||
// Invalider les containers spécifiques
|
||||
console.log(`[useMediaServers] ContainersUpdated pour ${serverId}:`, event.container_ids)
|
||||
|
||||
@@ -45,6 +45,63 @@ function ensureSSEConnected() {
|
||||
sse.onRendererEvent((event) => {
|
||||
const rendererId = event.renderer_id
|
||||
const timestamp = Date.parse(event.timestamp ?? '') || Date.now()
|
||||
|
||||
// Gérer les événements Online/Offline différemment
|
||||
if (event.type === 'online') {
|
||||
// Nouveau renderer découvert
|
||||
console.log(`[useRenderers] Renderer ${rendererId} (${event.friendly_name}) est maintenant en ligne`)
|
||||
|
||||
// Ajouter au cache avec les infos disponibles
|
||||
// Note: on n'a pas toutes les infos (capabilities, protocol) donc on fetch ensuite
|
||||
const renderer: RendererSummary = {
|
||||
id: rendererId,
|
||||
friendly_name: event.friendly_name,
|
||||
model_name: event.model_name,
|
||||
protocol: 'upnp', // Valeur par défaut, sera mise à jour par le fetch
|
||||
capabilities: {
|
||||
has_avtransport: false,
|
||||
has_avtransport_set_next: false,
|
||||
has_rendering_control: false,
|
||||
has_connection_manager: false,
|
||||
has_linkplay_http: false,
|
||||
has_arylic_tcp: false,
|
||||
has_oh_playlist: false,
|
||||
has_oh_volume: false,
|
||||
has_oh_info: false,
|
||||
has_oh_time: false,
|
||||
has_oh_radio: false,
|
||||
},
|
||||
online: true,
|
||||
}
|
||||
renderersCache.value.set(rendererId, renderer)
|
||||
|
||||
// Fetch la liste complète pour avoir les bonnes infos
|
||||
void fetchRenderers(true)
|
||||
|
||||
// Fetch le snapshot complet pour ce renderer
|
||||
void fetchRendererSnapshot(rendererId, { force: true })
|
||||
return
|
||||
}
|
||||
|
||||
if (event.type === 'offline') {
|
||||
// Renderer déconnecté
|
||||
console.log(`[useRenderers] Renderer ${rendererId} est maintenant hors ligne`)
|
||||
|
||||
// Marquer comme offline dans le cache
|
||||
const renderer = renderersCache.value.get(rendererId)
|
||||
if (renderer) {
|
||||
renderer.online = false
|
||||
renderersCache.value.set(rendererId, renderer)
|
||||
}
|
||||
|
||||
// Supprimer le snapshot (il n'est plus valide)
|
||||
snapshotState.snapshots.delete(rendererId)
|
||||
snapshotState.lastSnapshotAt.delete(rendererId)
|
||||
snapshotState.lastEventAt.delete(rendererId)
|
||||
return
|
||||
}
|
||||
|
||||
// Pour les autres événements, comportement existant
|
||||
snapshotState.lastEventAt.set(rendererId, timestamp)
|
||||
const lastSnapshot = snapshotState.lastSnapshotAt.get(rendererId) ?? 0
|
||||
if (!snapshotState.snapshots.has(rendererId) || timestamp > lastSnapshot) {
|
||||
|
||||
@@ -176,10 +176,14 @@ export type RendererEventPayload =
|
||||
| { type: 'metadata_changed'; renderer_id: string; title: string | null; artist: string | null; album: string | null; album_art_uri: string | null; timestamp: string }
|
||||
| { type: 'queue_updated'; renderer_id: string; queue_length: number; timestamp: string }
|
||||
| { type: 'binding_changed'; renderer_id: string; server_id: string | null; container_id: string | null; timestamp: string }
|
||||
| { type: 'online'; renderer_id: string; friendly_name: string; model_name: string; manufacturer: string; timestamp: string }
|
||||
| { type: 'offline'; renderer_id: string; timestamp: string }
|
||||
|
||||
export type MediaServerEventPayload =
|
||||
| { type: 'global_updated'; server_id: string; system_update_id: number | null; timestamp: string }
|
||||
| { type: 'containers_updated'; server_id: string; container_ids: string[]; timestamp: string }
|
||||
| { type: 'online'; server_id: string; friendly_name: string; model_name: string; manufacturer: string; timestamp: string }
|
||||
| { type: 'offline'; server_id: string; timestamp: string }
|
||||
|
||||
export type UnifiedEventPayload =
|
||||
| { category: 'renderer' } & RendererEventPayload
|
||||
|
||||
@@ -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"
|
||||
);
|
||||
}
|
||||
}
|
||||
})?;
|
||||
|
||||
@@ -1512,6 +1695,36 @@ impl ControlPoint {
|
||||
container_id: &str,
|
||||
auto_play: bool,
|
||||
) -> anyhow::Result<()> {
|
||||
// CRITICAL: When attaching a new playlist to a renderer, we must UNCONDITIONALLY
|
||||
// clear the RENDERER queue first (but NOT the local queue cache, which will be
|
||||
// replaced by refresh_attached_queue_for() using replace_entire_playlist()).
|
||||
//
|
||||
// Attach workflow: Clear renderer → Fill with new playlist
|
||||
// Update workflow: Gentle sync (preserve current item, use LCS)
|
||||
info!(
|
||||
renderer = renderer_id.0.as_str(),
|
||||
server = server_id.0.as_str(),
|
||||
container = container_id,
|
||||
"Attaching new playlist: clearing renderer queue"
|
||||
);
|
||||
|
||||
// Clear the renderer's queue for OpenHome renderers
|
||||
// We also sync the local cache to reflect the empty state, which will trigger
|
||||
// refresh_attached_queue_for() to use replace_entire_playlist() instead of gentle sync
|
||||
if self.runtime.uses_openhome_playlist(renderer_id) {
|
||||
let renderer = self.openhome_renderer(renderer_id)?;
|
||||
renderer.openhome_playlist_clear()?;
|
||||
// Sync local cache to reflect the empty renderer state
|
||||
self.sync_openhome_playlist_for(renderer_id)?;
|
||||
debug!(
|
||||
renderer = renderer_id.0.as_str(),
|
||||
"Cleared OpenHome renderer playlist and synced local cache"
|
||||
);
|
||||
} else {
|
||||
// For non-OpenHome renderers, use the standard clear_queue
|
||||
self.clear_queue(renderer_id)?;
|
||||
}
|
||||
|
||||
let binding = PlaylistBinding {
|
||||
server_id: server_id.clone(),
|
||||
container_id: container_id.to_string(),
|
||||
@@ -2425,6 +2638,9 @@ fn refresh_attached_queue_for(
|
||||
let current_item = current_idx.and_then(|idx| full_queue.get(idx).cloned());
|
||||
|
||||
// Check if renderer is playing - if so, we MUST preserve the current track
|
||||
// We use multiple signals to determine playback state:
|
||||
// 1. Direct playback_state() query (may fail on some renderers like upmpdcli)
|
||||
// 2. Presence of current_idx (if we have a current track, likely playing)
|
||||
let is_playing = {
|
||||
let renderer_info = {
|
||||
let reg = registry.read().unwrap();
|
||||
@@ -2433,7 +2649,20 @@ fn refresh_attached_queue_for(
|
||||
|
||||
if let Some(info) = renderer_info {
|
||||
if let Some(renderer) = MusicRenderer::from_registry_info(info, registry) {
|
||||
matches!(renderer.playback_state(), Ok(PlaybackState::Playing))
|
||||
// Try direct query first
|
||||
if matches!(renderer.playback_state(), Ok(PlaybackState::Playing)) {
|
||||
true
|
||||
} else if current_idx.is_some() && !full_queue.is_empty() {
|
||||
// Fallback: if we have a current index and non-empty queue,
|
||||
// assume playback is happening (handles renderers where playback_state() fails)
|
||||
debug!(
|
||||
renderer = renderer_id.0.as_str(),
|
||||
"playback_state() failed or not Playing, but current_idx is set - assuming playback"
|
||||
);
|
||||
true
|
||||
} else {
|
||||
false
|
||||
}
|
||||
} else {
|
||||
false
|
||||
}
|
||||
|
||||
@@ -58,20 +58,51 @@ impl OpenHomeQueue {
|
||||
track_ids.push(entry.id);
|
||||
}
|
||||
|
||||
let current_id = self
|
||||
.info_client
|
||||
.as_ref()
|
||||
.and_then(|client| client.id().ok());
|
||||
let previous_index = self
|
||||
.current_index
|
||||
.and_then(|idx| if idx < track_ids.len() { Some(idx) } else { None });
|
||||
let mut current_index = current_id
|
||||
.and_then(|id| track_ids.iter().position(|entry_id| *entry_id == id))
|
||||
.or(previous_index);
|
||||
// Try multiple methods to determine the currently playing track, from most to least reliable:
|
||||
// 1. Info.Id() - Direct ID query (fastest, but fails if track no longer in playlist)
|
||||
// 2. Info.Track() - Returns URI, which we can search for (works even if track removed)
|
||||
// 3. None - No current track can be determined
|
||||
let current_id = self.info_client.as_ref().and_then(|client| {
|
||||
// Try Info.Id() first
|
||||
if let Ok(id) = client.id() {
|
||||
debug!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
track_id = id,
|
||||
"Detected current track via Info.Id()"
|
||||
);
|
||||
return Some(id);
|
||||
}
|
||||
|
||||
if current_index.is_none() && !track_ids.is_empty() {
|
||||
current_index = Some(0);
|
||||
}
|
||||
// If Id() fails, try Track() to get the URI and search for it
|
||||
if let Ok(track_info) = client.track() {
|
||||
debug!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
track_uri = track_info.uri.as_str(),
|
||||
"Info.Id() failed, searching for current track by URI from Info.Track()"
|
||||
);
|
||||
return entries
|
||||
.iter()
|
||||
.find(|entry| entry.uri == track_info.uri)
|
||||
.map(|entry| {
|
||||
debug!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
found_id = entry.id,
|
||||
found_uri = entry.uri.as_str(),
|
||||
"Found current track ID by matching URI"
|
||||
);
|
||||
entry.id
|
||||
});
|
||||
}
|
||||
|
||||
debug!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
"Both Info.Id() and Info.Track() failed, cannot determine current track"
|
||||
);
|
||||
None
|
||||
});
|
||||
|
||||
let current_index = current_id
|
||||
.and_then(|id| track_ids.iter().position(|entry_id| *entry_id == id));
|
||||
|
||||
self.items = items;
|
||||
self.track_ids = track_ids;
|
||||
@@ -283,6 +314,389 @@ impl OpenHomeQueue {
|
||||
.copied()
|
||||
.ok_or_else(|| anyhow!("Failed to resolve OpenHome track id at index {}", index))
|
||||
}
|
||||
|
||||
/// CASE 1: Replace queue while preserving the currently playing item as first.
|
||||
/// The currently playing item is NOT in the new playlist, so we keep it as the first
|
||||
/// item and append the entire new playlist after it.
|
||||
fn replace_queue_preserve_current(
|
||||
&mut self,
|
||||
new_items: Vec<PlaybackItem>,
|
||||
playing_idx: usize,
|
||||
playing_id: u32,
|
||||
) -> Result<()> {
|
||||
// Re-read the current playlist state to get fresh IDs
|
||||
// This minimizes race conditions where IDs become invalid between our last refresh
|
||||
// and now (due to UPnP events from the server)
|
||||
let current_entries = self.playlist.read_all_tracks()?;
|
||||
let current_ids: Vec<u32> = current_entries.iter().map(|e| e.id).collect();
|
||||
|
||||
debug!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
fresh_id_count = current_ids.len(),
|
||||
cached_id_count = self.track_ids.len(),
|
||||
"Re-read playlist before deletions to avoid stale ID errors"
|
||||
);
|
||||
|
||||
// Find the playing track in the fresh list
|
||||
let fresh_playing_idx = current_ids.iter().position(|&id| id == playing_id);
|
||||
|
||||
if fresh_playing_idx.is_none() {
|
||||
debug!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
playing_id,
|
||||
"Playing track not found in fresh playlist - renderer state may have changed, aborting modification"
|
||||
);
|
||||
// The playing track is gone - don't try to manipulate the playlist
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
// Delete everything except the currently playing item (using fresh IDs)
|
||||
for &track_id in current_ids.iter().rev() {
|
||||
if track_id != playing_id {
|
||||
self.playlist.delete_id_if_exists(track_id)?;
|
||||
}
|
||||
}
|
||||
|
||||
// Rebuild: [currently_playing, new_items...]
|
||||
let mut rebuilt_items = Vec::with_capacity(1 + new_items.len());
|
||||
let mut rebuilt_ids = Vec::with_capacity(1 + new_items.len());
|
||||
|
||||
rebuilt_items.push(self.items[playing_idx].clone());
|
||||
rebuilt_ids.push(playing_id);
|
||||
|
||||
let mut previous_id = playing_id;
|
||||
for item in new_items {
|
||||
let metadata = build_metadata_xml(&item);
|
||||
let new_id = self.playlist.insert(previous_id, &item.uri, &metadata)?;
|
||||
previous_id = new_id;
|
||||
rebuilt_ids.push(new_id);
|
||||
rebuilt_items.push(self.item_with_openhome_id(item, new_id));
|
||||
}
|
||||
|
||||
self.items = rebuilt_items;
|
||||
self.track_ids = rebuilt_ids;
|
||||
self.current_index = Some(0); // Currently playing is now at index 0
|
||||
|
||||
debug!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
"Gentle sync completed: preserved playing track as first item (not in new playlist)"
|
||||
);
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// CASE 2: Replace queue with double-LCS (before and after the pivot).
|
||||
/// The currently playing item IS in the new playlist, so we use it as a pivot
|
||||
/// and apply LCS separately to the portions before and after it.
|
||||
fn replace_queue_with_pivot(
|
||||
&mut self,
|
||||
new_items: Vec<PlaybackItem>,
|
||||
pivot_idx_new: usize,
|
||||
pivot_id: u32,
|
||||
) -> Result<()> {
|
||||
debug!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
pivot_id,
|
||||
pivot_idx_new,
|
||||
new_playlist_len = new_items.len(),
|
||||
"Starting replace_queue_with_pivot - will re-read from OpenHome"
|
||||
);
|
||||
|
||||
// Re-read the current playlist state from OpenHome (the ONLY source of truth)
|
||||
// This is CRITICAL to avoid deleting IDs that no longer exist, which can
|
||||
// put the renderer (upmpdcli) into a degraded state where Info.TransportState()
|
||||
// starts returning HTTP 500 errors.
|
||||
let current_entries = self.playlist.read_all_tracks()?;
|
||||
|
||||
// Convert entries to PlaybackItems - this is the REAL current state
|
||||
let mut fresh_items = Vec::with_capacity(current_entries.len());
|
||||
let mut fresh_ids = Vec::with_capacity(current_entries.len());
|
||||
for entry in ¤t_entries {
|
||||
fresh_items.push(self.playback_item_from_entry(entry));
|
||||
fresh_ids.push(entry.id);
|
||||
}
|
||||
|
||||
debug!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
fresh_count = fresh_items.len(),
|
||||
"Re-read playlist from OpenHome (source of truth)"
|
||||
);
|
||||
|
||||
// Find the pivot in the fresh list
|
||||
let fresh_pivot_idx = fresh_ids.iter().position(|&id| id == pivot_id);
|
||||
|
||||
if fresh_pivot_idx.is_none() {
|
||||
debug!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
pivot_id,
|
||||
"Pivot track not found in fresh playlist - renderer state changed, aborting"
|
||||
);
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let fresh_pivot_idx = fresh_pivot_idx.unwrap();
|
||||
|
||||
// Split fresh data at the pivot - use ONLY fresh data, ignore cache
|
||||
let old_before: Vec<PlaybackItem> = fresh_items[..fresh_pivot_idx].to_vec();
|
||||
let old_after: Vec<PlaybackItem> = fresh_items[fresh_pivot_idx + 1..].to_vec();
|
||||
let old_ids_before: Vec<u32> = fresh_ids[..fresh_pivot_idx].to_vec();
|
||||
let old_ids_after: Vec<u32> = fresh_ids[fresh_pivot_idx + 1..].to_vec();
|
||||
|
||||
let new_before = &new_items[..pivot_idx_new];
|
||||
let new_after = &new_items[pivot_idx_new + 1..];
|
||||
|
||||
// LCS on the AFTER part (using fresh data from OpenHome)
|
||||
let (keep_old_after, keep_new_after) = lcs_flags(&old_after, new_after);
|
||||
|
||||
// LCS on the BEFORE part (using fresh data from OpenHome)
|
||||
let (keep_old_before, keep_new_before) = lcs_flags(&old_before, new_before);
|
||||
|
||||
// Delete items marked for deletion in AFTER part (reverse order)
|
||||
for (idx, &track_id) in old_ids_after.iter().enumerate().rev() {
|
||||
if !keep_old_after[idx] {
|
||||
debug!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
track_id,
|
||||
position = "AFTER pivot",
|
||||
"RENDERER OP: DeleteId({})",
|
||||
track_id
|
||||
);
|
||||
self.playlist.delete_id_if_exists(track_id)?;
|
||||
}
|
||||
}
|
||||
|
||||
// Delete items marked for deletion in BEFORE part (reverse order)
|
||||
for (idx, &track_id) in old_ids_before.iter().enumerate().rev() {
|
||||
if !keep_old_before[idx] {
|
||||
debug!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
track_id,
|
||||
position = "BEFORE pivot",
|
||||
"RENDERER OP: DeleteId({})",
|
||||
track_id
|
||||
);
|
||||
self.playlist.delete_id_if_exists(track_id)?;
|
||||
}
|
||||
}
|
||||
|
||||
// Rebuild the playlist: [BEFORE, PIVOT, AFTER]
|
||||
let mut rebuilt_items = Vec::with_capacity(new_items.len());
|
||||
let mut rebuilt_ids = Vec::with_capacity(new_items.len());
|
||||
|
||||
// Collect IDs of kept items in BEFORE part (in order)
|
||||
let remaining_before: Vec<u32> = old_ids_before
|
||||
.iter()
|
||||
.enumerate()
|
||||
.filter_map(|(idx, &id)| if keep_old_before[idx] { Some(id) } else { None })
|
||||
.collect();
|
||||
|
||||
let mut remaining_before_idx = 0;
|
||||
let mut previous_id = OPENHOME_PLAYLIST_HEAD_ID;
|
||||
|
||||
// Rebuild BEFORE part
|
||||
for (idx, item) in new_before.iter().enumerate() {
|
||||
if keep_new_before[idx] {
|
||||
let existing_id = remaining_before[remaining_before_idx];
|
||||
remaining_before_idx += 1;
|
||||
previous_id = existing_id;
|
||||
rebuilt_ids.push(existing_id);
|
||||
rebuilt_items.push(self.item_with_openhome_id(item.clone(), existing_id));
|
||||
debug!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
track_id = existing_id,
|
||||
position = "BEFORE pivot",
|
||||
"KEPT existing track ID {}",
|
||||
existing_id
|
||||
);
|
||||
} else {
|
||||
let metadata = build_metadata_xml(item);
|
||||
let new_id = self.playlist.insert(previous_id, &item.uri, &metadata)?;
|
||||
debug!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
after_id = previous_id,
|
||||
new_id,
|
||||
position = "BEFORE pivot",
|
||||
"RENDERER OP: Insert(after={}) -> new_id={}",
|
||||
previous_id,
|
||||
new_id
|
||||
);
|
||||
previous_id = new_id;
|
||||
rebuilt_ids.push(new_id);
|
||||
rebuilt_items.push(self.item_with_openhome_id(item.clone(), new_id));
|
||||
}
|
||||
}
|
||||
|
||||
// Add PIVOT (keeps its ID!)
|
||||
rebuilt_ids.push(pivot_id);
|
||||
rebuilt_items.push(self.item_with_openhome_id(new_items[pivot_idx_new].clone(), pivot_id));
|
||||
previous_id = pivot_id;
|
||||
debug!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
pivot_id,
|
||||
pivot_idx_new,
|
||||
"PIVOT preserved with ID {} at index {}",
|
||||
pivot_id,
|
||||
pivot_idx_new
|
||||
);
|
||||
|
||||
// Collect IDs of kept items in AFTER part (in order)
|
||||
let remaining_after: Vec<u32> = old_ids_after
|
||||
.iter()
|
||||
.enumerate()
|
||||
.filter_map(|(idx, &id)| if keep_old_after[idx] { Some(id) } else { None })
|
||||
.collect();
|
||||
|
||||
let mut remaining_after_idx = 0;
|
||||
|
||||
// Rebuild AFTER part
|
||||
for (idx, item) in new_after.iter().enumerate() {
|
||||
if keep_new_after[idx] {
|
||||
let existing_id = remaining_after[remaining_after_idx];
|
||||
remaining_after_idx += 1;
|
||||
previous_id = existing_id;
|
||||
rebuilt_ids.push(existing_id);
|
||||
rebuilt_items.push(self.item_with_openhome_id(item.clone(), existing_id));
|
||||
debug!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
track_id = existing_id,
|
||||
position = "AFTER pivot",
|
||||
"KEPT existing track ID {}",
|
||||
existing_id
|
||||
);
|
||||
} else {
|
||||
let metadata = build_metadata_xml(item);
|
||||
let new_id = self.playlist.insert(previous_id, &item.uri, &metadata)?;
|
||||
debug!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
after_id = previous_id,
|
||||
new_id,
|
||||
position = "AFTER pivot",
|
||||
"RENDERER OP: Insert(after={}) -> new_id={}",
|
||||
previous_id,
|
||||
new_id
|
||||
);
|
||||
previous_id = new_id;
|
||||
rebuilt_ids.push(new_id);
|
||||
rebuilt_items.push(self.item_with_openhome_id(item.clone(), new_id));
|
||||
}
|
||||
}
|
||||
|
||||
self.items = rebuilt_items;
|
||||
self.track_ids = rebuilt_ids;
|
||||
self.current_index = Some(pivot_idx_new); // Pivot is at its new position
|
||||
|
||||
// VERIFICATION: Check that pivot ID is preserved
|
||||
let final_pivot_id = self.track_ids.get(pivot_idx_new).copied();
|
||||
if final_pivot_id != Some(pivot_id) {
|
||||
return Err(anyhow!(
|
||||
"CRITICAL BUG: Pivot ID changed from {} to {:?} during replace_queue_with_pivot!",
|
||||
pivot_id,
|
||||
final_pivot_id
|
||||
));
|
||||
}
|
||||
|
||||
debug!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
pivot_idx = pivot_idx_new,
|
||||
pivot_id,
|
||||
final_playlist_len = self.track_ids.len(),
|
||||
pivot_verified = true,
|
||||
"Gentle sync completed: double-LCS with pivot (playing track preserved)"
|
||||
);
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Standard LCS-based replacement (used when no currently playing item).
|
||||
fn replace_queue_standard_lcs(
|
||||
&mut self,
|
||||
items: Vec<PlaybackItem>,
|
||||
current_index: Option<usize>,
|
||||
) -> Result<()> {
|
||||
let (keep_current, keep_desired) = lcs_flags(&self.items, &items);
|
||||
|
||||
let items_to_keep = keep_current.iter().filter(|&&k| k).count();
|
||||
let items_to_delete = keep_current.iter().filter(|&&k| !k).count();
|
||||
let items_to_add = keep_desired.iter().filter(|&&k| !k).count();
|
||||
|
||||
debug!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
keep = items_to_keep,
|
||||
delete = items_to_delete,
|
||||
add = items_to_add,
|
||||
"LCS computed: minimizing OpenHome playlist operations"
|
||||
);
|
||||
|
||||
// If we're replacing everything (keep=0), use delete_all() instead of
|
||||
// individual delete_id() calls. This is much more robust for live playlists
|
||||
// where track IDs can become invalid between refresh and deletion.
|
||||
if items_to_keep == 0 && items_to_delete > 0 {
|
||||
debug!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
"Using delete_all() for complete replacement (more robust for live playlists)"
|
||||
);
|
||||
self.playlist.delete_all()?;
|
||||
self.track_ids.clear();
|
||||
self.items.clear();
|
||||
} else {
|
||||
// Selective deletion when keeping some items
|
||||
for idx in (0..self.track_ids.len()).rev() {
|
||||
if !keep_current[idx] {
|
||||
let track_id = self.track_ids[idx];
|
||||
// Use delete_id_if_exists() to handle cases where another control point
|
||||
// may have already modified the playlist
|
||||
self.playlist.delete_id_if_exists(track_id)?;
|
||||
self.track_ids.remove(idx);
|
||||
self.items.remove(idx);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
let remaining_ids = self.track_ids.clone();
|
||||
let mut remaining_idx = 0usize;
|
||||
let mut previous_id = OPENHOME_PLAYLIST_HEAD_ID;
|
||||
let mut rebuilt_items = Vec::with_capacity(items.len());
|
||||
let mut rebuilt_ids = Vec::with_capacity(items.len());
|
||||
|
||||
for (idx, item) in items.into_iter().enumerate() {
|
||||
if keep_desired[idx] {
|
||||
if remaining_idx >= remaining_ids.len() {
|
||||
return Err(anyhow!(
|
||||
"OpenHome playlist refresh bookkeeping mismatch (kept entries underflow)"
|
||||
));
|
||||
}
|
||||
let existing_id = remaining_ids[remaining_idx];
|
||||
remaining_idx += 1;
|
||||
previous_id = existing_id;
|
||||
rebuilt_ids.push(existing_id);
|
||||
rebuilt_items.push(self.item_with_openhome_id(item, existing_id));
|
||||
} else {
|
||||
let metadata = build_metadata_xml(&item);
|
||||
let new_id = self.playlist.insert(previous_id, &item.uri, &metadata)?;
|
||||
previous_id = new_id;
|
||||
rebuilt_ids.push(new_id);
|
||||
rebuilt_items.push(self.item_with_openhome_id(item, new_id));
|
||||
}
|
||||
}
|
||||
|
||||
if remaining_idx != remaining_ids.len() {
|
||||
return Err(anyhow!(
|
||||
"OpenHome playlist refresh bookkeeping mismatch (kept entries overflow)"
|
||||
));
|
||||
}
|
||||
|
||||
let previous_index = self
|
||||
.current_index
|
||||
.and_then(|idx| if idx < rebuilt_ids.len() { Some(idx) } else { None });
|
||||
let normalized = current_index
|
||||
.filter(|&i| i < rebuilt_ids.len())
|
||||
.or(previous_index)
|
||||
.or_else(|| if rebuilt_ids.is_empty() { None } else { Some(0) });
|
||||
self.items = rebuilt_items;
|
||||
self.track_ids = rebuilt_ids;
|
||||
self.current_index = normalized;
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
pub fn didl_id_from_metadata(xml: &str) -> Option<String> {
|
||||
@@ -425,92 +839,86 @@ impl QueueBackend for OpenHomeQueue {
|
||||
// differences. Without this, any drift between our cache and the renderer
|
||||
// (e.g., manual edits from another control point) would keep the stale items.
|
||||
self.refresh_from_openhome()?;
|
||||
|
||||
// Try to get the currently playing track ID from the renderer.
|
||||
// Note: Some OpenHome renderers (like upmpdcli) don't reliably support Info.Id(),
|
||||
// so we fall back to using our internal current_index pointer.
|
||||
let currently_playing_id_from_renderer = self
|
||||
.info_client
|
||||
.as_ref()
|
||||
.and_then(|client| client.id().ok());
|
||||
|
||||
// Find the currently playing item in our local state.
|
||||
// Priority: 1) Renderer-reported ID, 2) Our internal current_index
|
||||
let playing_info = if let Some(id) = currently_playing_id_from_renderer {
|
||||
// CASE: Renderer explicitly reported the playing track ID
|
||||
self.track_ids
|
||||
.iter()
|
||||
.position(|&tid| tid == id)
|
||||
.map(|idx| (idx, id, self.items[idx].uri.clone()))
|
||||
} else if let Some(idx) = self.current_index {
|
||||
// CASE: Use our internal pointer (fallback for renderers without Info.Id() support)
|
||||
if idx < self.track_ids.len() && idx < self.items.len() {
|
||||
let id = self.track_ids[idx];
|
||||
let uri = self.items[idx].uri.clone();
|
||||
debug!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
current_index = idx,
|
||||
track_id = id,
|
||||
"Using internal current_index as fallback (renderer didn't report playing ID)"
|
||||
);
|
||||
Some((idx, id, uri))
|
||||
} else {
|
||||
None
|
||||
}
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
debug!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
actual_items = self.items.len(),
|
||||
currently_playing_id_from_renderer = ?currently_playing_id_from_renderer,
|
||||
playing_info_detected = playing_info.is_some(),
|
||||
"OpenHome playlist state refreshed before replace_queue"
|
||||
);
|
||||
|
||||
let (keep_current, keep_desired) = lcs_flags(&self.items, &items);
|
||||
if let Some((playing_idx, playing_id, playing_uri)) = playing_info {
|
||||
// Find if the currently playing item is in the new playlist (by URI)
|
||||
let new_playing_idx = items.iter().position(|item| item.uri == playing_uri);
|
||||
|
||||
let items_to_keep = keep_current.iter().filter(|&&k| k).count();
|
||||
let items_to_delete = keep_current.iter().filter(|&&k| !k).count();
|
||||
let items_to_add = keep_desired.iter().filter(|&&k| !k).count();
|
||||
if let Some(pivot_idx) = new_playing_idx {
|
||||
// CASE 2: Currently playing item IS in the new playlist
|
||||
// Use gentle double-LCS strategy: preserve the pivot and sync before/after separately
|
||||
debug!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
playing_idx,
|
||||
pivot_idx,
|
||||
"Gentle sync: currently playing item found in new playlist at index {}",
|
||||
pivot_idx
|
||||
);
|
||||
|
||||
debug!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
keep = items_to_keep,
|
||||
delete = items_to_delete,
|
||||
add = items_to_add,
|
||||
"LCS computed: minimizing OpenHome playlist operations"
|
||||
);
|
||||
self.replace_queue_with_pivot(items, pivot_idx, playing_id)?;
|
||||
} else {
|
||||
// CASE 1: Currently playing item NOT in the new playlist
|
||||
// Keep it as first item and append the new playlist after it
|
||||
debug!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
playing_idx,
|
||||
"Gentle sync: currently playing item not in new playlist, preserving as first item"
|
||||
);
|
||||
|
||||
// If we're replacing everything (keep=0), use delete_all() instead of
|
||||
// individual delete_id() calls. This is much more robust for live playlists
|
||||
// where track IDs can become invalid between refresh and deletion.
|
||||
if items_to_keep == 0 && items_to_delete > 0 {
|
||||
self.replace_queue_preserve_current(items, playing_idx, playing_id)?;
|
||||
}
|
||||
} else {
|
||||
// No currently playing item or can't determine it - use standard LCS
|
||||
debug!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
"Using delete_all() for complete replacement (more robust for live playlists)"
|
||||
"No currently playing item, using standard LCS sync"
|
||||
);
|
||||
self.playlist.delete_all()?;
|
||||
self.track_ids.clear();
|
||||
self.items.clear();
|
||||
} else {
|
||||
// Selective deletion when keeping some items
|
||||
for idx in (0..self.track_ids.len()).rev() {
|
||||
if !keep_current[idx] {
|
||||
let track_id = self.track_ids[idx];
|
||||
self.playlist.delete_id(track_id)?;
|
||||
self.track_ids.remove(idx);
|
||||
self.items.remove(idx);
|
||||
}
|
||||
}
|
||||
self.replace_queue_standard_lcs(items, current_index)?;
|
||||
}
|
||||
|
||||
let remaining_ids = self.track_ids.clone();
|
||||
let mut remaining_idx = 0usize;
|
||||
let mut previous_id = OPENHOME_PLAYLIST_HEAD_ID;
|
||||
let mut rebuilt_items = Vec::with_capacity(items.len());
|
||||
let mut rebuilt_ids = Vec::with_capacity(items.len());
|
||||
|
||||
for (idx, item) in items.into_iter().enumerate() {
|
||||
if keep_desired[idx] {
|
||||
if remaining_idx >= remaining_ids.len() {
|
||||
return Err(anyhow!(
|
||||
"OpenHome playlist refresh bookkeeping mismatch (kept entries underflow)"
|
||||
));
|
||||
}
|
||||
let existing_id = remaining_ids[remaining_idx];
|
||||
remaining_idx += 1;
|
||||
previous_id = existing_id;
|
||||
rebuilt_ids.push(existing_id);
|
||||
rebuilt_items.push(self.item_with_openhome_id(item, existing_id));
|
||||
} else {
|
||||
let metadata = build_metadata_xml(&item);
|
||||
let new_id = self.playlist.insert(previous_id, &item.uri, &metadata)?;
|
||||
previous_id = new_id;
|
||||
rebuilt_ids.push(new_id);
|
||||
rebuilt_items.push(self.item_with_openhome_id(item, new_id));
|
||||
}
|
||||
}
|
||||
|
||||
if remaining_idx != remaining_ids.len() {
|
||||
return Err(anyhow!(
|
||||
"OpenHome playlist refresh bookkeeping mismatch (kept entries overflow)"
|
||||
));
|
||||
}
|
||||
|
||||
let previous_index = self
|
||||
.current_index
|
||||
.and_then(|idx| if idx < rebuilt_ids.len() { Some(idx) } else { None });
|
||||
let normalized = current_index
|
||||
.filter(|&i| i < rebuilt_ids.len())
|
||||
.or(previous_index)
|
||||
.or_else(|| if rebuilt_ids.is_empty() { None } else { Some(0) });
|
||||
self.items = rebuilt_items;
|
||||
self.track_ids = rebuilt_ids;
|
||||
self.current_index = normalized;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -531,7 +939,9 @@ impl QueueBackend for OpenHomeQueue {
|
||||
self.ensure_track_id(index - 1)?
|
||||
};
|
||||
|
||||
self.playlist.delete_id(track_id)?;
|
||||
// Use delete_id_if_exists() to handle cases where another control point
|
||||
// may have already modified the playlist
|
||||
self.playlist.delete_id_if_exists(track_id)?;
|
||||
let metadata = build_metadata_xml(&item);
|
||||
let new_id = self.playlist.insert(before_id, &item.uri, &metadata)?;
|
||||
|
||||
|
||||
@@ -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,
|
||||
},
|
||||
}
|
||||
|
||||
@@ -121,12 +121,12 @@ impl OhPlaylistClient {
|
||||
|
||||
pub fn play_id(&self, id: u32) -> Result<()> {
|
||||
let id_str = id.to_string();
|
||||
let args = [("Id", id_str.as_str())];
|
||||
let args = [("Value", id_str.as_str())];
|
||||
|
||||
let call_result =
|
||||
invoke_upnp_action(&self.control_url, &self.service_type, "PlayId", &args)?;
|
||||
invoke_upnp_action(&self.control_url, &self.service_type, "SeekId", &args)?;
|
||||
|
||||
handle_action_response("PlayId", &call_result)
|
||||
handle_action_response("SeekId", &call_result)
|
||||
}
|
||||
|
||||
pub fn play(&self) -> Result<()> {
|
||||
@@ -170,7 +170,7 @@ impl OhPlaylistClient {
|
||||
|
||||
pub fn delete_id(&self, id: u32) -> Result<()> {
|
||||
let id_str = id.to_string();
|
||||
let args = [("Id", id_str.as_str())];
|
||||
let args = [("Value", id_str.as_str())];
|
||||
|
||||
let call_result =
|
||||
invoke_upnp_action(&self.control_url, &self.service_type, "DeleteId", &args)?;
|
||||
@@ -178,12 +178,59 @@ impl OhPlaylistClient {
|
||||
handle_action_response("DeleteId", &call_result)
|
||||
}
|
||||
|
||||
/// Attempts to delete an OpenHome playlist entry by ID.
|
||||
/// Unlike delete_id(), this function silently ignores errors related to invalid/missing IDs,
|
||||
/// which is useful in multi-control-point scenarios where playlist state may have changed.
|
||||
///
|
||||
/// Returns:
|
||||
/// - Ok(true) if the ID was successfully deleted
|
||||
/// - Ok(false) if the ID didn't exist (logged as warning)
|
||||
/// - Err(_) for other errors (network issues, etc.)
|
||||
pub fn delete_id_if_exists(&self, id: u32) -> Result<bool> {
|
||||
match self.delete_id(id) {
|
||||
Ok(()) => Ok(true),
|
||||
Err(err) => {
|
||||
// Check if this is an error about an invalid/missing ID
|
||||
// OpenHome servers may return different error messages/codes for this case
|
||||
let err_msg = format!("{err}");
|
||||
if err_msg.contains("Invalid")
|
||||
|| err_msg.contains("invalid")
|
||||
|| err_msg.contains("not found")
|
||||
|| err_msg.contains("does not exist")
|
||||
|| err_msg.contains("unknown")
|
||||
|| err_msg.contains("500") // HTTP 500 = renderer in inconsistent state
|
||||
|| err_msg.contains("Action Failed") // UPnP error 501
|
||||
{
|
||||
warn!(
|
||||
control_url = self.control_url.as_str(),
|
||||
id,
|
||||
"DeleteId silently ignored - ID does not exist or renderer in inconsistent state (likely modified by events)"
|
||||
);
|
||||
Ok(false)
|
||||
} else {
|
||||
// Re-throw other errors (network issues, etc.)
|
||||
Err(err)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub fn delete_all(&self) -> Result<()> {
|
||||
let call_result =
|
||||
invoke_upnp_action(&self.control_url, &self.service_type, "DeleteAll", &[])?;
|
||||
handle_action_response("DeleteAll", &call_result)
|
||||
}
|
||||
|
||||
pub fn current_id(&self) -> Result<String> {
|
||||
let call_result =
|
||||
invoke_upnp_action(&self.control_url, &self.service_type, "Id", &[])?;
|
||||
let envelope = ensure_success("Id", &call_result)?;
|
||||
let response = find_child_with_suffix(&envelope.body.content, "IdResponse")
|
||||
.ok_or_else(|| anyhow!("Missing IdResponse element in SOAP body"))?;
|
||||
let value: String = extract_child_text_any(response, &["aValue", "Value"])?;
|
||||
Ok(value)
|
||||
}
|
||||
|
||||
pub fn tracks_max(&self) -> Result<u32> {
|
||||
let call_result =
|
||||
invoke_upnp_action(&self.control_url, &self.service_type, "TracksMax", &[])?;
|
||||
@@ -191,7 +238,7 @@ impl OhPlaylistClient {
|
||||
let envelope = ensure_success("TracksMax", &call_result)?;
|
||||
let response = find_child_with_suffix(&envelope.body.content, "TracksMaxResponse")
|
||||
.ok_or_else(|| anyhow!("Missing TracksMaxResponse element in SOAP body"))?;
|
||||
let value_text = extract_child_text_any(response, &["aValue", "Value"])?;
|
||||
let value_text: String = extract_child_text_any(response, &["aValue", "Value"])?;
|
||||
let value = value_text
|
||||
.parse::<u32>()
|
||||
.map_err(|_| anyhow!("Invalid TracksMax value: {}", value_text))?;
|
||||
@@ -700,16 +747,26 @@ fn parse_track_list(payload: &str) -> Result<Vec<OhTrackEntry>> {
|
||||
|
||||
fn parse_track_entry(elem: &Element) -> Result<OhTrackEntry> {
|
||||
let id_text = extract_child_text_local(elem, "Id")?;
|
||||
|
||||
// Some renderers (like upmpdcli) return a comma-separated list of IDs in the Id element
|
||||
// when using ReadList. This is a non-standard compact format that we cannot parse properly
|
||||
// because we need to fetch each track individually to get its Uri and Metadata.
|
||||
// Return an error to force the fallback to individual Read() calls.
|
||||
if id_text.contains(',') {
|
||||
debug!(
|
||||
raw_entry = %elem.name,
|
||||
raw_id = id_text.as_str(),
|
||||
"Unexpected multi-value Id element in OpenHome TrackList entry"
|
||||
"Multi-value Id element detected - forcing fallback to individual reads"
|
||||
);
|
||||
return Err(anyhow!(
|
||||
"Renderer returned comma-separated IDs in single Entry - need individual reads"
|
||||
));
|
||||
}
|
||||
|
||||
let id = id_text
|
||||
.parse::<u32>()
|
||||
.map_err(|_| anyhow!("Invalid OpenHome Entry Id: {}", id_text))?;
|
||||
|
||||
let uri = extract_child_text_local(elem, "Uri")?;
|
||||
let metadata_xml = extract_child_text_optional_local(elem, "Metadata")?.unwrap_or_default();
|
||||
|
||||
@@ -927,7 +984,7 @@ pub(crate) fn decode_base64(input: &str) -> Result<Vec<u8>> {
|
||||
|
||||
fn is_invalid_entry_id_error(err: &anyhow::Error) -> bool {
|
||||
let msg = format!("{err}");
|
||||
msg.contains("Invalid OpenHome Entry Id")
|
||||
msg.contains("Invalid OpenHome Entry Id") || msg.contains("comma-separated IDs")
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -341,13 +341,6 @@ fn spawn_playlist_event_handler(manager: Arc<ParadiseChannelManager>) {
|
||||
e
|
||||
);
|
||||
}
|
||||
if let Err(e) = append_track_to_history(descriptor, &cache_pk).await {
|
||||
tracing::warn!(
|
||||
"Failed to update history for channel {}: {}",
|
||||
descriptor.display_name,
|
||||
e
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -361,22 +354,3 @@ fn channel_from_live_playlist(playlist_id: &str) -> Option<&'static ChannelDescr
|
||||
.iter()
|
||||
.find(|descriptor| descriptor.slug == slug)
|
||||
}
|
||||
|
||||
async fn append_track_to_history(descriptor: &ChannelDescriptor, cache_pk: &str) -> Result<()> {
|
||||
let playlist_id = format!("radio-paradise-history-{}", descriptor.slug);
|
||||
let manager = pmoplaylist::PlaylistManager();
|
||||
let handle = manager
|
||||
.get_persistent_write_handle(playlist_id.clone())
|
||||
.await
|
||||
.with_context(|| format!("Failed to get history playlist {}", playlist_id))?;
|
||||
|
||||
if handle.contains_pk(cache_pk).await? {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
handle
|
||||
.push(cache_pk.to_string())
|
||||
.await
|
||||
.with_context(|| format!("Failed to append {} to {}", cache_pk, playlist_id))?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -272,14 +272,25 @@ impl RadioParadisePlaylistFeeder {
|
||||
// 5. Sauvegarder les métadonnées
|
||||
self.save_metadata(&pk, song, &block).await?;
|
||||
|
||||
// 6. Calculer le TTL
|
||||
// 6. Vérifier si la chanson est déjà dans la playlist
|
||||
if self.playlist_handle.contains_pk(&pk).await? {
|
||||
tracing::debug!(
|
||||
"RadioParadisePlaylistFeeder: Skipping duplicate song {} - {} (pk={})",
|
||||
idx,
|
||||
song.title,
|
||||
pk
|
||||
);
|
||||
continue;
|
||||
}
|
||||
|
||||
// 7. Calculer le TTL
|
||||
let sched_end = song
|
||||
.sched_end_time_ms()
|
||||
.ok_or_else(|| anyhow::anyhow!("Cannot calculate TTL without sched_time_millis"))?;
|
||||
let ttl_ms = sched_end.saturating_sub(now_ms);
|
||||
let ttl = Duration::from_millis(ttl_ms);
|
||||
|
||||
// 7. Push dans la playlist avec TTL
|
||||
// 8. Push dans la playlist avec TTL
|
||||
self.playlist_handle.push_with_ttl(pk.clone(), ttl).await?;
|
||||
|
||||
tracing::info!(
|
||||
|
||||
@@ -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