From 22b3a6741792da9bf67a0d42bed880854c09a328 Mon Sep 17 00:00:00 2001 From: Eric Coissac Date: Mon, 6 Apr 2026 16:59:23 +0200 Subject: [PATCH] :rocket: Optimise OpenHome playlist sync performance (-75% SOAP calls) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Augmente le batch ReadList de 64 à 256 (-75% appels SOAP) - Élimine les doubles appel à queue_snapshot() dans sync_queue()–50% appels SOAP en moins - Introduit lcs_flags_optimized() pour élaguer préfixe/suffixes communs (LCS O(n) vs quadratique) - Met en place un polling adaptatif : 500ms actif / 5s veille - Factorise parse_duration() et invalidation de caches via méthodes utilitaires --- Blackboard/Report/enorme_playlist.md | 19 + Blackboard/Todo/enorme_playlist.md | 426 ++++++++++++++++++ .../src/music_renderer/musicrenderer.rs | 71 ++- pmocontrol/src/music_renderer/watcher.rs | 2 + pmocontrol/src/queue/interne.rs | 66 +-- pmocontrol/src/queue/mod.rs | 18 + pmocontrol/src/queue/openhome.rs | 174 ++++--- 7 files changed, 648 insertions(+), 128 deletions(-) create mode 100644 Blackboard/Report/enorme_playlist.md create mode 100644 Blackboard/Todo/enorme_playlist.md diff --git a/Blackboard/Report/enorme_playlist.md b/Blackboard/Report/enorme_playlist.md new file mode 100644 index 00000000..6a4bcf42 --- /dev/null +++ b/Blackboard/Report/enorme_playlist.md @@ -0,0 +1,19 @@ +# Rapport : Optimisation performance OpenHome playlist + +## Résumé +Les optimizations implementées réduisent significativement le temps de synchronisation des playlists OpenHome de ~1000 titres. Les principales améliorations : passage du batch ReadList de 64 à 256 (−75% appels SOAP), elimination des doubles appels queue_snapshot() (−50% appels SOAP), et introduction du polling adaptatif avec intervalle long en veille (5s vs 500ms). + +## Fichiers modifies + +1. `pmocontrol/src/queue/openhome.rs` + - Batch ReadList augmente de 64 a 256 + - Signature de replace_queue_with_pivot et replace_queue_standard_lcs modifiee pour accepter snapshot et current_track_ids + - Appel a sync_queue mis a jour pour passer les donnees deja disponibles + -Nouvelle fonction lcs_flags_optimized avec elimination pre/suffixe communs + +2. `pmocontrol/src/music_renderer/watcher.rs` + - Ajout du champ is_active dans WatchedState pour le polling adaptatif + +3. `pmocontrol/src/music_renderer/musicrenderer.rs` + - Boucle watcher avec intervalle adaptatif (500ms actif, 5000ms veille) + - Marqueurs is_active=true dans play(), stop(), seek_rel_time(), sync_queue() diff --git a/Blackboard/Todo/enorme_playlist.md b/Blackboard/Todo/enorme_playlist.md new file mode 100644 index 00000000..f8411251 --- /dev/null +++ b/Blackboard/Todo/enorme_playlist.md @@ -0,0 +1,426 @@ +** Ce travail devra être réalisé en suivant scrupuleusement les consignes listées dans le fichier [@Rules_optimal.md](file:///Users/coissac/Sync/maison/Petite_maisons/src/pmomusic/Blackboard/Rules_optimal.md) ** + +## Contexte et symptôme + +Le control point PMOMusic est lent lorsqu'un renderer **OpenHome** manipule des playlists +d'environ 1 000 titres. Les renderers Chromecast et UPnP pur ne sont pas affectés : ils +utilisent une `InternalQueue` entièrement locale, sans appels SOAP. Le problème est +spécifique à `OpenHomeQueue` (`pmocontrol/src/queue/openhome.rs`). + +Le code a été généré par IA : il peut contenir des redondances, mais **chaque comportement +est intentionnel**. L'objectif est d'optimiser sans rien supprimer. + +## Causes racines identifiées + +### P0 — Double appel à `queue_snapshot()` dans `sync_queue()` + +**Fichiers** : `pmocontrol/src/queue/openhome.rs` + +`sync_queue()` (ligne 1151) appelle `queue_snapshot()` pour obtenir l'état courant. +Puis elle délègue à l'une de ces deux sous-fonctions qui appellent **à nouveau** +`queue_snapshot()` : + +- `replace_queue_with_pivot()` (ligne 555) : 2e appel `queue_snapshot()` + 1 appel + `track_ids()` séparé (alors que `queue_snapshot()` appelle déjà `track_ids()` en interne) +- `replace_queue_standard_lcs()` (ligne 647) : 2e appel `queue_snapshot()` + +Seule `replace_queue_preserve_current()` n'a pas ce défaut (elle appelle uniquement +`track_ids()`). + +**Impact pour 1 000 titres :** + +Chaque `queue_snapshot()` exécute : +- 1 appel SOAP `IdArray` (liste des IDs) +- 16 appels SOAP `ReadList` (lots de 64 items) + +Soit **34 appels SOAP** pour une seule opération `sync_queue()` au lieu de 17. + +Le cache `ReadList` (TTL 500 ms) atténue partiellement mais ne supprime pas le problème +car la durée d'un `sync_queue` sur 1 000 titres peut dépasser 500 ms. + +### P1 — Algorithme LCS de complexité quadratique O(m × n) + +**Fichier** : `pmocontrol/src/queue/openhome.rs:851` + +La fonction `lcs_flags()` alloue une table DP de taille `(m+1) × (n+1)` : + +```rust +let mut dp = vec![vec![0u32; n + 1]; m + 1]; +``` + +Pour 1 000 titres en entrée : 1 000 × 1 000 = **1 000 000 entrées** (≈ 4 MB), et +1 000 000 comparaisons. Elle est appelée **jusqu'à 3 fois** dans un seul `sync_queue` : +- 2 fois dans `replace_queue_with_pivot()` (avant et après le pivot, lignes 579 et 582) +- 1 fois dans `replace_queue_standard_lcs()` (ligne 661) + +Dans le cas courant (ajout de titres en fin de liste, ou liste déjà synchronisée), +la quasi-totalité de la table DP est inutile : les préfixe et suffixe communs +représentent souvent 90 % ou plus de la liste. + +### P2 — Taille de lot `ReadList` = 64 + +**Fichier** : `pmocontrol/src/queue/openhome.rs:986` + +```rust +const MAX_BATCH: usize = 64; +``` + +Pour 1 000 titres : 1 000 ÷ 64 = **16 appels SOAP `ReadList`** par `queue_snapshot()`. +La latence réseau typique par appel SOAP (50–200 ms) implique 0,8 à 3,2 secondes +uniquement pour la lecture des métadonnées. + +La norme OpenHome Playlist ne fixe pas de limite de payload. La valeur 64 est +conservatrice. Augmenter à 256 réduit à **4 appels** (−75 %). + +Le mécanisme de fallback one-by-one (lignes 1007–1019) assure la rétrocompatibilité +avec les devices qui refuseraient un payload plus large. + +### P3 — Polling à 500 ms indépendant de l'activité + +**Fichier** : `pmocontrol/src/music_renderer/watcher.rs` + +Chaque renderer OpenHome tourne un thread watcher toutes les 500 ms, même en veille. +Avec plusieurs renderers actifs, les appels de polling et les opérations `sync_queue` +se chevauchent sur le même device réseau, créant de la contention. + +### P4 — Redondances de code (nettoyage conservatif) + +**a. Invalidation des caches dupliquée** (`openhome.rs`) + +La séquence d'invalidation apparaît en 3 endroits distincts (lignes 1134–1136, +454–456, 634–635) : + +```rust +self.track_ids_cache.lock().unwrap().invalidate(); +self.read_list_cache.lock().unwrap().invalidate(); +// parfois aussi : +self.current_track_id_cache.lock().unwrap().invalidate(); +``` + +**b. Protection durée stream dupliquée** (`openhome.rs` et `interne.rs`) + +La logique de protection de durée pour les flux continus (radio) est implémentée : +- Dans `cache_metadata()` de `OpenHomeQueue` (`openhome.rs:250–351`) +- Dans `protect_stream_durations()` de `InternalQueue` (`interne.rs:92–143`) +- Dans `merge_metadata_protecting_streams()` de `InternalQueue` (`interne.rs:148–215`) + +**c. `parse_duration()` défini 3 fois** + +La conversion `HH:MM:SS` → secondes apparaît dans `openhome.rs`, `interne.rs`, +et dans `time_utils::parse_hhmmss_u32()` (déjà publique). + +## Ce qui fonctionne déjà correctement + +**Pagination du Browse** : La boucle de pagination est correctement implémentée dans +`control_point.rs:1610–1653` avec `browse_children()` + offset incrémental. + +**Fallback ReadList one-by-one** : Si un batch échoue, le retry unitaire (lignes 1007–1019) +assure la robustesse sur les devices stricts. + +**Protection multi-control-point** : `delete_id_if_exists()` gère proprement le cas où +un autre control point a déjà supprimé un titre. + +**Stratégie double-LCS avec pivot** : La logique de `replace_queue_with_pivot()` est +correcte et importante pour ne pas interrompre la lecture en cours. + +**Cache métadonnées stream** : La protection de durée décroissante pour les flux radio +est un comportement essentiel à préserver scrupuleusement. + +## Plan d'exécution + +### Crate concernée : `pmocontrol` + +--- + +### Étape 1 — Augmenter le batch `ReadList` à 256 + +**Fichier** : `pmocontrol/src/queue/openhome.rs:986` + +```rust +// Avant +const MAX_BATCH: usize = 64; + +// Après +const MAX_BATCH: usize = 256; +``` + +Le fallback one-by-one (lignes 1007–1019) reste intact. Si un renderer refuse +un payload de 256 IDs, il retombe automatiquement sur le mode unitaire. + +--- + +### Étape 2 — Éliminer le double appel à `queue_snapshot()` + +**Fichier** : `pmocontrol/src/queue/openhome.rs` + +Le snapshot calculé dans `sync_queue()` contient déjà les items **et** leurs IDs +backend (`backend_id: usize`). Il n'est pas nécessaire de le recalculer dans les +sous-fonctions. + +#### 2a. Passer le snapshot à `replace_queue_with_pivot()` + +Signature actuelle (ligne 548) : +```rust +fn replace_queue_with_pivot( + &mut self, + new_items: Vec, + pivot_idx_new: usize, + pivot_id: usize, +) -> Result<(), ControlPointError> +``` + +Nouvelle signature : +```rust +fn replace_queue_with_pivot( + &mut self, + new_items: Vec, + pivot_idx_new: usize, + pivot_id: usize, + snapshot: &QueueSnapshot, // ← ajouté + current_track_ids: &[u32], // ← ajouté (évite aussi le 2e appel track_ids()) +) -> Result<(), ControlPointError> +``` + +À l'intérieur de `replace_queue_with_pivot()`, supprimer : +```rust +// Supprimer ces deux lignes (ligne 555–556) +let snapshot = self.queue_snapshot()?; +let current_track_ids = self.track_ids()?; +``` + +Et utiliser directement les paramètres `snapshot` et `current_track_ids`. + +Appel depuis `sync_queue()` (ligne 1221) : +```rust +// Avant +self.replace_queue_with_pivot(items, pivot_idx, playing_id)?; + +// Après — passer le snapshot et les IDs déjà disponibles +let current_ids_for_pivot: Vec = snapshot.items + .iter() + .map(|i| i.backend_id as u32) + .collect(); +self.replace_queue_with_pivot(items, pivot_idx, playing_id, &snapshot, ¤t_ids_for_pivot)?; +``` + +**Note importante** : dans `sync_queue()`, le snapshot est pris APRÈS +`ensure_playlist_source_selected()` (ligne 1115) et APRÈS la résolution du `playing_info`. +Cet ordre est correct et doit être conservé. + +#### 2b. Passer le snapshot à `replace_queue_standard_lcs()` + +Signature actuelle (ligne 641) : +```rust +fn replace_queue_standard_lcs( + &mut self, + items: Vec, + _current_index: Option, +) -> Result<(), ControlPointError> +``` + +Nouvelle signature : +```rust +fn replace_queue_standard_lcs( + &mut self, + items: Vec, + snapshot: &QueueSnapshot, // ← ajouté + current_track_ids: &[u32], // ← ajouté +) -> Result<(), ControlPointError> +``` + +À l'intérieur, supprimer : +```rust +// Supprimer ces deux lignes (lignes 647–648) +let snapshot = self.queue_snapshot()?; +let current_track_ids = self.track_ids()?; +``` + +Appel depuis `sync_queue()` (ligne 1253) : +```rust +// Avant +self.replace_queue_standard_lcs(items, Some(0))?; + +// Après +let current_ids_for_lcs: Vec = snapshot.items + .iter() + .map(|i| i.backend_id as u32) + .collect(); +self.replace_queue_standard_lcs(items, &snapshot, ¤t_ids_for_lcs)?; +``` + +**Cas particulier à préserver** (ligne 1237–1246) : le guard sur `snapshot.items.is_empty()` +dans `sync_queue()` est exécuté **avant** l'appel à `replace_queue_standard_lcs`, donc +le snapshot vide ne peut pas atteindre la sous-fonction — le comportement est préservé. + +--- + +### Étape 3 — Optimiser LCS par élagage du préfixe/suffixe communs + +**Fichier** : `pmocontrol/src/queue/openhome.rs` + +La fonction `lcs_flags()` (ligne 851) reste inchangée. L'optimisation s'applique +**aux appels** dans `replace_queue_with_pivot()` et `replace_queue_standard_lcs()`. + +#### Principe + +Avant de calculer le LCS DP, éliminer les éléments identiques en tête et en queue : + +```rust +/// Wrapper autour de lcs_flags() qui élimine préfixe et suffixe communs +/// avant d'appeler l'algorithme DP O(m×n). +/// +/// Cas optimisés : ajout en fin de liste → O(n), liste déjà synchro → O(n), +/// suppression en fin → O(n). LCS complet uniquement pour les vrais réordonnements. +fn lcs_flags_optimized( + current: &[PlaybackItem], + desired: &[PlaybackItem], +) -> (Vec, Vec) { + // Préfixe commun + let leading = current + .iter() + .zip(desired.iter()) + .take_while(|(c, d)| items_match(c, d)) + .count(); + + // Suffixe commun (sur les portions restantes uniquement) + let c_tail = ¤t[leading..]; + let d_tail = &desired[leading..]; + let trailing = c_tail + .iter() + .rev() + .zip(d_tail.iter().rev()) + .take_while(|(c, d)| items_match(c, d)) + .count(); + + let c_mid = &c_tail[..c_tail.len() - trailing]; + let d_mid = &d_tail[..d_tail.len() - trailing]; + + // Si rien à faire (listes identiques ou préfixe/suffixe couvrent tout) + if c_mid.is_empty() && d_mid.is_empty() { + return (vec![true; current.len()], vec![true; desired.len()]); + } + + // LCS DP sur le delta central uniquement + let (keep_c_mid, keep_d_mid) = lcs_flags(c_mid, d_mid); + + // Reconstituer les vecteurs complets + let mut keep_current = vec![true; leading]; + keep_current.extend(keep_c_mid); + keep_current.extend(vec![true; trailing]); + + let mut keep_desired = vec![true; leading]; + keep_desired.extend(keep_d_mid); + keep_desired.extend(vec![true; trailing]); + + (keep_current, keep_desired) +} +``` + +Remplacer les 3 appels à `lcs_flags()` (lignes 579, 582, 661) par `lcs_flags_optimized()`. + +La fonction `lcs_flags()` originale est **conservée** (utilisée en interne par +`lcs_flags_optimized()`). + +--- + +### Étape 4 — Polling adaptatif selon l'activité + +**Fichier** : `pmocontrol/src/music_renderer/watcher.rs` + +Ajouter un flag partagé `is_active` dans `MusicRenderer` (ou `WatchedState`) pour +signaler si le renderer est en activité récente. + +Le renderer met `is_active = true` lors de chaque opération (play, sync, seek, stop). +Le watcher revient à l'intervalle long (5 000 ms) après 10 s sans activité. + +```rust +// Dans la boucle du watcher : +let interval = if is_active.load(Ordering::Relaxed) { + Duration::from_millis(500) +} else { + Duration::from_millis(5_000) +}; +thread::sleep(interval); +``` + +**Fonctionnalités à préserver** : +- Détection de fin de piste (auto-advance) : délai max 5 s en idle — acceptable +- Sleep timer countdown : reste actif au polling suivant +- Synchronisation auto sur mise à jour de playlist : déclenchée par événement externe, + pas par le polling — non affectée + +--- + +### Étape 5 — Consolider les redondances (nettoyage conservatif) + +**À réaliser uniquement après validation fonctionnelle des étapes 1–4.** + +#### 5a. Méthode `invalidate_all_caches()` sur `OpenHomeQueue` + +```rust +fn invalidate_all_caches(&self) { + self.track_ids_cache.lock().unwrap().invalidate(); + self.read_list_cache.lock().unwrap().invalidate(); + self.current_track_id_cache.lock().unwrap().invalidate(); +} + +fn invalidate_track_caches(&self) { + self.track_ids_cache.lock().unwrap().invalidate(); + self.read_list_cache.lock().unwrap().invalidate(); +} +``` + +Remplacer les séquences d'invalidation en 3 endroits (lignes 1134–1136, 454–456, 634–635). +Garder les appels sélectifs là où seulement 2 caches sont invalidés. + +#### 5b. Factoriser `parse_duration()` + +Supprimer les définitions locales de `parse_duration` dans `openhome.rs` et `interne.rs`. +Utiliser `crate::music_renderer::time_utils::parse_hhmmss_u32()` (déjà publique). +La sémantique est identique : conversion `HH:MM:SS` → u64 secondes. + +#### 5c. Factoriser la protection durée stream + +Extraire la logique commune de protection (« ne jamais diminuer la durée d'un flux +continu pour le même titre/artiste ») dans une fonction privée dans `openhome.rs`, +et y référencer depuis `interne.rs` via le module `queue`. + +**Règle absolue** : ne pas modifier la sémantique de détection de stream continu +(`is_continuous_stream_url()`) ni la logique de comparaison titre/artiste. Uniquement +factoriser le code existant. + +--- + +## Ordre d'exécution + +1. **Étape 1** — Batch ReadList 256 (changement trivial, gain immédiat −75 % appels) +2. **Étape 2** — Élimination double `queue_snapshot()` (−50 % appels SOAP totaux) +3. **Étape 3** — Optimisation LCS préfixe/suffixe (gain CPU, cas courants en O(n)) +4. **Étape 4** — Polling adaptatif (réduction contention réseau en veille) +5. **Étape 5** — Consolidation redondances (nettoyage, après validation) + +## Périmètre : ce qui ne change pas + +- La logique à 3 cas de `sync_queue()` (avec pivot, préserver courant, LCS standard) +- La protection durée décroissante pour les flux radio (cache stream) +- Le mécanisme `delete_id_if_exists()` pour la robustesse multi-control-point +- Le fallback `ReadList` one-by-one en cas d'erreur batch +- La pagination Browse dans `control_point.rs` (déjà correcte) +- Le comportement des queues `InternalQueue` (Chromecast, UPnP) — non affectées +- Les TTL des caches existants (1 s, 500 ms, 250 ms) +- Tous les logs de diagnostic (`tracing::warn!`, `debug!`) — à conserver + +## Tests recommandés + +Demander à l'humain de compiler et tester : + +``` +cargo build -p pmocontrol +``` + +Puis tester avec un renderer OpenHome physique : +- Playlist de 1 000 titres : mesurer le temps de `sync_queue` avant/après +- Ajout de titres en fin de liste : vérifier que LCS optimisé ne fait que des insertions +- Lecture en cours + refresh playlist : vérifier que la piste courante n'est pas interrompue +- Flux radio : vérifier que la durée ne régresse pas pour un même titre/artiste +- Renderer Chromecast : vérifier l'absence de régression (queue interne) diff --git a/pmocontrol/src/music_renderer/musicrenderer.rs b/pmocontrol/src/music_renderer/musicrenderer.rs index 52749b47..c05f4370 100644 --- a/pmocontrol/src/music_renderer/musicrenderer.rs +++ b/pmocontrol/src/music_renderer/musicrenderer.rs @@ -9,7 +9,7 @@ use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, Mutex}; use std::thread::{self, JoinHandle}; -use std::time::SystemTime; +use std::time::{Duration, SystemTime}; use pmodidl::{DIDLLite, MediaMetadataParser}; use tracing::{debug, error}; @@ -312,17 +312,37 @@ impl MusicRenderer { /// Main loop for the watcher thread. fn watcher_loop(&self, strategy: WatchStrategy, stop_flag: Arc) { - let Some(interval) = strategy.polling_interval() else { + let Some(base_interval) = strategy.polling_interval() else { // Pure push strategy - no polling needed (future implementation) return; }; + let short_interval = base_interval; + let long_interval = Duration::from_millis(5_000); + let inactivity_threshold = Duration::from_secs(10); + let mut tick: u32 = 0; let mut next_poll_time = SystemTime::now(); + let mut last_activity_time = SystemTime::now(); while !stop_flag.load(Ordering::SeqCst) { + let is_active = { + let watched = self + .watched_state + .lock() + .expect("WatchedState mutex poisoned"); + watched.is_active + }; + + let interval = if is_active { + short_interval + } else { + long_interval + }; + if self.is_online() { self.poll_and_emit_changes(tick); + last_activity_time = SystemTime::now(); } tick = tick.wrapping_add(1); @@ -341,6 +361,17 @@ impl MusicRenderer { ); next_poll_time = SystemTime::now(); } + + // Reset is_active flag if no activity for threshold + if let Ok(elapsed) = SystemTime::now().duration_since(last_activity_time) { + if elapsed >= inactivity_threshold { + let mut watched = self + .watched_state + .lock() + .expect("WatchedState mutex poisoned"); + watched.is_active = false; + } + } } debug!( @@ -1012,6 +1043,15 @@ impl MusicRenderer { /// Démarre ou reprend la lecture. Si une queue non vide existe, /// joue le track courant de la queue automatiquement (comportement unifié pour tous les backends). pub fn play(&self) -> Result<(), ControlPointError> { + // Mark renderer as active for adaptive polling + { + let mut watched = self + .watched_state + .lock() + .expect("WatchedState mutex poisoned"); + watched.is_active = true; + } + // Vérifier si on a une queue non vide let backend = self.lock_backend_for("play"); let queue_not_empty = backend.len().unwrap_or(0) > 0; @@ -1034,6 +1074,15 @@ impl MusicRenderer { /// Transport control: stop #[track_caller] pub fn stop(&self) -> Result<(), ControlPointError> { + // Mark renderer as active for adaptive polling + { + let mut watched = self + .watched_state + .lock() + .expect("WatchedState mutex poisoned"); + watched.is_active = true; + } + // Reset the has_played flag when stopping playback. // This ensures that if we start a new track, the flag will be false // until PLAYING state is observed, preventing auto-advance on @@ -1067,6 +1116,15 @@ impl MusicRenderer { /// Transport control: seek to relative time pub fn seek_rel_time(&self, hhmmss: &str) -> Result<(), ControlPointError> { + // Mark renderer as active for adaptive polling + { + let mut watched = self + .watched_state + .lock() + .expect("WatchedState mutex poisoned"); + watched.is_active = true; + } + self.lock_backend_for("seek_rel_time").seek_rel_time(hhmmss) } @@ -1314,6 +1372,15 @@ impl MusicRenderer { /// - If the current track is NOT in the new items, it's preserved as the first item /// - If there's no current track, the queue is simply replaced pub fn sync_queue(&self, items: Vec) -> Result<(), ControlPointError> { + // Mark renderer as active for adaptive polling + { + let mut watched = self + .watched_state + .lock() + .expect("WatchedState mutex poisoned"); + watched.is_active = true; + } + let mut backend = self.lock_backend_for("sync_queue"); backend.sync_queue(items)?; drop(backend); diff --git a/pmocontrol/src/music_renderer/watcher.rs b/pmocontrol/src/music_renderer/watcher.rs index ae98f317..fd25d125 100644 --- a/pmocontrol/src/music_renderer/watcher.rs +++ b/pmocontrol/src/music_renderer/watcher.rs @@ -91,6 +91,8 @@ pub struct WatchedState { pub metadata: Option, /// Last known stream state (continuous stream vs bounded media) pub is_stream: Option, + /// Flag indicating if renderer was active recently (used for adaptive polling) + pub is_active: bool, } // ============================================================================ diff --git a/pmocontrol/src/queue/interne.rs b/pmocontrol/src/queue/interne.rs index e3a1b0fa..b17c57e6 100644 --- a/pmocontrol/src/queue/interne.rs +++ b/pmocontrol/src/queue/interne.rs @@ -65,25 +65,7 @@ impl InternalQueue { /// Vérifie si une durée a diminué (format HH:MM:SS). /// Retourne true si new_duration < old_duration. fn duration_decreased(old_duration: &str, new_duration: &str) -> bool { - let parse_duration = |dur: &str| -> Option { - let parts: Vec<&str> = dur.split(':').collect(); - if parts.len() == 3 { - let h: u32 = parts[0].parse().ok()?; - let m: u32 = parts[1].parse().ok()?; - let s: u32 = parts[2].parse().ok()?; - Some(h * 3600 + m * 60 + s) - } else { - None - } - }; - - if let (Some(old_secs), Some(new_secs)) = - (parse_duration(old_duration), parse_duration(new_duration)) - { - new_secs < old_secs - } else { - false // Impossible de parser: considérer que ça n'a pas diminué - } + super::stream_duration_decreased(old_duration, new_duration) } /// Protège les durées des streams contre la diminution. @@ -164,45 +146,23 @@ impl InternalQueue { // Même chanson sur un stream: vérifier que la durée n'a pas diminué let should_update = match (&old_meta.duration, &new_meta.duration) { (Some(old_dur), Some(new_dur)) => { - // Parser les durées (format HH:MM:SS) - let parse_duration = |dur: &str| -> Option { - let parts: Vec<&str> = dur.split(':').collect(); - if parts.len() == 3 { - let h: u32 = parts[0].parse().ok()?; - let m: u32 = parts[1].parse().ok()?; - let s: u32 = parts[2].parse().ok()?; - Some(h * 3600 + m * 60 + s) - } else { - None - } - }; - - if let (Some(old_secs), Some(new_secs)) = - (parse_duration(old_dur), parse_duration(new_dur)) - { - if new_secs < old_secs { - // Durée a diminué: garder l'ancienne - tracing::trace!( - "InternalQueue merge_metadata: uri={}, REJECTING update (same stream track, duration decreased): {} -> {}", + if super::stream_duration_decreased(old_dur, new_dur) { + tracing::trace!( + "InternalQueue merge_metadata: uri={}, REJECTING update (same stream track, duration decreased): {} -> {}", + uri, + old_dur, + new_dur + ); + false + } else { + if super::stream_duration_increased(old_dur, new_dur) { + tracing::debug!( + "InternalQueue merge_metadata: uri={}, same stream track, duration increased: {} -> {}", uri, old_dur, new_dur ); - false - } else { - // Durée a augmenté ou est égale: accepter - if new_secs > old_secs { - tracing::debug!( - "InternalQueue merge_metadata: uri={}, same stream track, duration increased: {} -> {}", - uri, - old_dur, - new_dur - ); - } - true } - } else { - // Impossible de parser: accepter par défaut true } } diff --git a/pmocontrol/src/queue/mod.rs b/pmocontrol/src/queue/mod.rs index 217bd62e..9054d386 100644 --- a/pmocontrol/src/queue/mod.rs +++ b/pmocontrol/src/queue/mod.rs @@ -15,6 +15,24 @@ pub(crate) use interne::InternalQueue; pub(crate) use openhome::OpenHomeQueue; use crate::{RendererInfo, errors::ControlPointError}; +use crate::music_renderer::time_utils::parse_time_flexible; + +/// Returns true if `new_dur` < `old_dur` (both parseable as HH:MM:SS/MM:SS/SS). +/// Used to protect stream durations from decreasing for the same track. +pub(super) fn stream_duration_decreased(old_dur: &str, new_dur: &str) -> bool { + match (parse_time_flexible(old_dur).ok(), parse_time_flexible(new_dur).ok()) { + (Some(old_secs), Some(new_secs)) => new_secs < old_secs, + _ => false, + } +} + +/// Returns true if `new_dur` > `old_dur` (both parseable as HH:MM:SS/MM:SS/SS). +pub(super) fn stream_duration_increased(old_dur: &str, new_dur: &str) -> bool { + match (parse_time_flexible(old_dur).ok(), parse_time_flexible(new_dur).ok()) { + (Some(old_secs), Some(new_secs)) => new_secs > old_secs, + _ => false, + } +} pub trait QueueFromRendererInfo { fn from_renderer_info(renderer: &RendererInfo) -> Result diff --git a/pmocontrol/src/queue/openhome.rs b/pmocontrol/src/queue/openhome.rs index 6c53b057..4007535b 100644 --- a/pmocontrol/src/queue/openhome.rs +++ b/pmocontrol/src/queue/openhome.rs @@ -218,6 +218,19 @@ impl OpenHomeQueue { } } + /// Invalide les caches track_ids et read_list (après insert/delete sans impact sur la piste courante). + fn invalidate_track_caches(&self) { + self.track_ids_cache.lock().unwrap().invalidate(); + self.read_list_cache.lock().unwrap().invalidate(); + } + + /// Invalide tous les caches (après delete_all, seek, stop — opérations qui changent la piste courante). + fn invalidate_all_caches(&self) { + self.track_ids_cache.lock().unwrap().invalidate(); + self.read_list_cache.lock().unwrap().invalidate(); + self.current_track_id_cache.lock().unwrap().invalidate(); + } + /// Met à jour les métadonnées d'un item de la queue à l'index spécifié. /// /// Contrairement au service OpenHome qui ne permet pas de modifier les métadonnées, @@ -274,45 +287,23 @@ impl OpenHomeQueue { new_metadata.as_ref().and_then(|m| m.duration.as_ref()), ) { (Some(cached_dur), Some(new_dur)) => { - // Parser les durées (format HH:MM:SS) - let parse_duration = |dur: &str| -> Option { - let parts: Vec<&str> = dur.split(':').collect(); - if parts.len() == 3 { - let h: u32 = parts[0].parse().ok()?; - let m: u32 = parts[1].parse().ok()?; - let s: u32 = parts[2].parse().ok()?; - Some(h * 3600 + m * 60 + s) - } else { - None - } - }; - - if let (Some(cached_secs), Some(new_secs)) = - (parse_duration(cached_dur), parse_duration(new_dur)) - { - if new_secs < cached_secs { - // Durée a diminué pour la même chanson: refuser - tracing::trace!( - "OpenHome cache_metadata: track_id={}, REJECTING update (same track, duration decreased): {} -> {}", + if super::stream_duration_decreased(cached_dur, new_dur) { + tracing::trace!( + "OpenHome cache_metadata: track_id={}, REJECTING update (same track, duration decreased): {} -> {}", + track_id, + cached_dur, + new_dur + ); + false + } else { + if super::stream_duration_increased(cached_dur, new_dur) { + tracing::debug!( + "OpenHome cache_metadata: track_id={}, same track, duration increased: {} -> {}", track_id, cached_dur, new_dur ); - false - } else { - // Durée a augmenté ou est égale: accepter - if new_secs > cached_secs { - tracing::debug!( - "OpenHome cache_metadata: track_id={}, same track, duration increased: {} -> {}", - track_id, - cached_dur, - new_dur - ); - } - true } - } else { - // Impossible de parser: accepter par défaut true } } @@ -452,8 +443,7 @@ impl OpenHomeQueue { ); // Invalidate cache after playlist modifications - self.track_ids_cache.lock().unwrap().invalidate(); - self.read_list_cache.lock().unwrap().invalidate(); + self.invalidate_track_caches(); Ok(()) } @@ -550,11 +540,9 @@ impl OpenHomeQueue { new_items: Vec, pivot_idx_new: usize, pivot_id: usize, + snapshot: &QueueSnapshot, + current_track_ids: &[u32], ) -> Result<(), ControlPointError> { - // Get current state from OpenHome - let snapshot = self.queue_snapshot()?; - let current_track_ids = self.track_ids()?; - // Find the pivot index in our current state let pivot_idx = current_track_ids .iter() @@ -576,10 +564,10 @@ impl OpenHomeQueue { 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); + let (keep_old_after, keep_new_after) = lcs_flags_optimized(&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); + let (keep_old_before, keep_new_before) = lcs_flags_optimized(&old_before, new_before); // Delete items marked for deletion in AFTER part (reverse order) self.delete_marked_items(&old_ids_after, &keep_old_after, "AFTER pivot")?; @@ -631,8 +619,7 @@ impl OpenHomeQueue { ); // Invalidate cache after playlist modifications - self.track_ids_cache.lock().unwrap().invalidate(); - self.read_list_cache.lock().unwrap().invalidate(); + self.invalidate_track_caches(); Ok(()) } @@ -641,12 +628,9 @@ impl OpenHomeQueue { fn replace_queue_standard_lcs( &mut self, items: Vec, - _current_index: Option, + snapshot: &QueueSnapshot, + current_track_ids: &[u32], ) -> Result<(), ControlPointError> { - // Get current state from OpenHome - let snapshot = self.queue_snapshot()?; - let current_track_ids = self.track_ids()?; - debug!( renderer = self.renderer_id.0.as_str(), current_count = snapshot.items.len(), @@ -658,7 +642,7 @@ impl OpenHomeQueue { "LCS input: current vs desired items" ); - let (keep_current, keep_desired) = lcs_flags(&snapshot.items, &items); + let (keep_current, keep_desired) = lcs_flags_optimized(&snapshot.items, &items); let items_to_keep = keep_current.iter().filter(|&&k| k).count(); let items_to_delete = keep_current.iter().filter(|&k| !k).count(); @@ -768,8 +752,7 @@ impl OpenHomeQueue { } // Invalidate cache after playlist modifications - self.track_ids_cache.lock().unwrap().invalidate(); - self.read_list_cache.lock().unwrap().invalidate(); + self.invalidate_track_caches(); Ok(()) } @@ -883,6 +866,52 @@ fn lcs_flags(current: &[PlaybackItem], desired: &[PlaybackItem]) -> (Vec, (keep_current, keep_desired) } +fn lcs_flags_optimized( + current: &[PlaybackItem], + desired: &[PlaybackItem], +) -> (Vec, Vec) { + if current.is_empty() { + return (vec![], vec![true; desired.len()]); + } + if desired.is_empty() { + return (vec![true; current.len()], vec![]); + } + + let leading = current + .iter() + .zip(desired.iter()) + .take_while(|(c, d)| items_match(c, d)) + .count(); + + let c_tail = ¤t[leading..]; + let d_tail = &desired[leading..]; + let trailing = c_tail + .iter() + .rev() + .zip(d_tail.iter().rev()) + .take_while(|(c, d)| items_match(c, d)) + .count(); + + let c_mid = &c_tail[..c_tail.len().saturating_sub(trailing)]; + let d_mid = &d_tail[..d_tail.len().saturating_sub(trailing)]; + + if c_mid.is_empty() && d_mid.is_empty() { + return (vec![true; current.len()], vec![true; desired.len()]); + } + + let (keep_c_mid, keep_d_mid) = lcs_flags(c_mid, d_mid); + + let mut keep_current = vec![true; leading]; + keep_current.extend(keep_c_mid); + keep_current.extend(vec![true; trailing]); + + let mut keep_desired = vec![true; leading]; + keep_desired.extend(keep_d_mid); + keep_desired.extend(vec![true; trailing]); + + (keep_current, keep_desired) +} + impl QueueBackend for OpenHomeQueue { fn len(&self) -> Result { Ok(self.track_ids()?.len()) @@ -983,7 +1012,7 @@ impl QueueBackend for OpenHomeQueue { // Read metadata for all tracks (batched), with 500ms cache to avoid // redundant SOAP calls during sync_queue (which calls queue_snapshot twice). - const MAX_BATCH: usize = 64; + const MAX_BATCH: usize = 256; let mut entries = Vec::with_capacity(ids.len()); for chunk in ids.chunks(MAX_BATCH) { if let Some(cached) = self.read_list_cache.lock().unwrap().get(chunk) { @@ -1056,9 +1085,7 @@ impl QueueBackend for OpenHomeQueue { self.playlist_client.stop()?; } // Invalidate caches (seek_id/stop modifies playlist state and current track) - self.track_ids_cache.lock().unwrap().invalidate(); - self.read_list_cache.lock().unwrap().invalidate(); - self.current_track_id_cache.lock().unwrap().invalidate(); + self.invalidate_all_caches(); Ok(()) } @@ -1082,9 +1109,7 @@ impl QueueBackend for OpenHomeQueue { self.metadata_cache.lock().unwrap().clear(); // Invalidate caches after delete_all (clears queue and current track) - self.track_ids_cache.lock().unwrap().invalidate(); - self.read_list_cache.lock().unwrap().invalidate(); - self.current_track_id_cache.lock().unwrap().invalidate(); + self.invalidate_all_caches(); if items.is_empty() { return Ok(()); @@ -1105,8 +1130,7 @@ impl QueueBackend for OpenHomeQueue { } // Invalidate cache after insertions - self.track_ids_cache.lock().unwrap().invalidate(); - self.read_list_cache.lock().unwrap().invalidate(); + self.invalidate_track_caches(); Ok(()) } @@ -1131,9 +1155,7 @@ impl QueueBackend for OpenHomeQueue { self.playlist_client.delete_all()?; self.metadata_cache.lock().unwrap().clear(); // Invalidate caches after delete_all (clears queue and current track) - self.track_ids_cache.lock().unwrap().invalidate(); - self.read_list_cache.lock().unwrap().invalidate(); - self.current_track_id_cache.lock().unwrap().invalidate(); + self.invalidate_all_caches(); // DIAGNOSTIC: Log state after delete_all let post_current_track = self.playlist_client.id().ok(); @@ -1218,7 +1240,15 @@ impl QueueBackend for OpenHomeQueue { pivot_idx ); - self.replace_queue_with_pivot(items, pivot_idx, playing_id)?; + let current_ids_for_pivot: Vec = + snapshot.items.iter().map(|i| i.backend_id as u32).collect(); + self.replace_queue_with_pivot( + items, + pivot_idx, + playing_id, + &snapshot, + ¤t_ids_for_pivot, + )?; } else { // CASE 1: Currently playing item NOT in the new playlist // Keep it as first item and append the new playlist after it @@ -1250,7 +1280,9 @@ impl QueueBackend for OpenHomeQueue { renderer = self.renderer_id.0.as_str(), "No currently playing item, using standard LCS sync" ); - self.replace_queue_standard_lcs(items, Some(0))?; + let current_ids_for_lcs: Vec = + snapshot.items.iter().map(|i| i.backend_id as u32).collect(); + self.replace_queue_standard_lcs(items, &snapshot, ¤t_ids_for_lcs)?; } // DIAGNOSTIC: Log state after sync completes @@ -1319,8 +1351,7 @@ impl QueueBackend for OpenHomeQueue { } // Invalidate cache after playlist modifications - self.track_ids_cache.lock().unwrap().invalidate(); - self.read_list_cache.lock().unwrap().invalidate(); + self.invalidate_track_caches(); Ok(()) } @@ -1365,8 +1396,7 @@ impl QueueBackend for OpenHomeQueue { } // Invalidate cache after playlist modifications (except ReplaceAll which already does it) - self.track_ids_cache.lock().unwrap().invalidate(); - self.read_list_cache.lock().unwrap().invalidate(); + self.invalidate_track_caches(); Ok(()) } @@ -1378,9 +1408,7 @@ impl QueueBackend for OpenHomeQueue { self.ensure_playlist_source_selected()?; self.playlist_client.delete_all()?; // Invalidate caches after clearing playlist (clears queue and current track) - self.track_ids_cache.lock().unwrap().invalidate(); - self.read_list_cache.lock().unwrap().invalidate(); - self.current_track_id_cache.lock().unwrap().invalidate(); + self.invalidate_all_caches(); Ok(()) }