diff --git a/.kilo/plans/1775631632458-happy-cactus.md b/.kilo/plans/1775631632458-happy-cactus.md new file mode 100644 index 00000000..314962fd --- /dev/null +++ b/.kilo/plans/1775631632458-happy-cactus.md @@ -0,0 +1,225 @@ +# Audit PMO Control - Optimisation Playlist OpenHome + +## Résumé Exécutif + +L'utilisateur rapporte des lenteurs significatives lors de la manipulation de playlists de ~1000 titres avec les renderers OpenHome. Les renderers Chromecast et UPnP (avec queue interne) ne sont pas affectés. + +## État des Optimisations Deja Implémentées + +Le precedent plan dans `Blackboard/Todo/enorme_playlist.md` a deja été partiellement implémenté: + +| optimisation | Statut | Emplacement | +|-------------|-------|------------| +| MAX_BATCH = 256 pour ReadList | ✅ FAIT | `openhome.rs:1015` | +| Élimination double queue_snapshot() | ✅ FAIT | `replace_queue_with_pivot()` et `replace_queue_standard_lcs()` | +| LCS préfixe/suffixe (lcs_flags_optimized) | ✅ FAIT | `openhome.rs:869` | +| Polling adaptatif (is_active) | ✅ FAIT | `musicrenderer.rs:325-374` | +| Consolidation invalidation caches | ✅ FAIT | `openhome.rs:222-235` | + +## Contraintes Protocolaires Découvertes + +L'action `Insert` **ne supporte PAS l'insertion par lot** - chaque appel prend un seul Uri/Metadata. + +## Nouvelles Optimisations (Contre-propositions Utilisateur) + +### OPT-1: Fast Path pour 99% des cas de sync_queue + +**Observation**: 99% des changements de playlist sont: +- Insertion de nouvelles tracks en **fin de queue** +- Délétion de tracks en **début de queue** +- Rarement des changements nécessitant un vrai alignement LCS + +**Solution**: Ajouter une détection de pattern avant d'appeler LCS: + +```rust +fn smart_sync(&mut self, items: Vec) -> Result<(), ControlPointError> { + let current_ids = self.track_ids()?; + + // Cas 1: Append only (insertion en fin) + if items.starts_with(¤t_ids) { + // Fast path: juste ajouter les nouveaux items + return self.append_only(items.skip(current_ids.len())); + } + + // Cas 2: Delete from beginning + if current_ids.starts_with(&items) { + // Fast path: supprimer de la fin + return self.delete_from_beginning(current_ids.len() - items.len()); + } + + // Cas 3: Full LCS only for complex reorderings + return self.replace_queue_standard_lcs(items); +} +``` + +**Impact**: 99% des sync_queue passent de O(N²) à O(N) + +--- + +### OPT-2: Queue FIFO pour Opérations OpenHome (Thread Background) + +**Concept**: Une file d'attente FIFO des opérations SOAP exécutée dans un thread dédié. + +```rust +pub struct OpenHomeOpQueue { + queue: Arc>>, + worker_handle: Option>, +} + +pub enum OpenHomeOp { + Insert { uri: String, metadata: String, after_id: u32 }, + Delete { track_id: u32 }, + DeleteAll, + SeekId { id: u32 }, + Play, + Pause, + Stop, + SetVolume { volume: u16 }, + // Meta operations + UpdateMetadata { track_id: u32, metadata: String }, +} + +impl OpenHomeOpQueue { + /// Push operation to the FIFO queue + pub fn push(&self, op: OpenHomeOp) { + self.queue.lock().unwrap().push(op); + } + + /// Push with priority (volume, play, stop - need fast response) + pub fn push_first(&self, op: OpenHomeOp) { + self.queue.lock().unwrap().push_front(op); + } + + /// Clear all pending operations (client can flush) + pub fn clear(&self) { + self.queue.lock().unwrap().clear(); + } + + /// Worker thread consumes operations + fn worker_loop(&self) { + loop { + let op = self.queue.lock().unwrap().pop_front(); + match op { + Some(op) => self.execute(op), + None => thread::sleep(Duration::from_millis(10)), + } + } + } +} +``` + +**Benefits**: +- UI non-bloquante (les operations sont lancées et exec en background) +- Batching naturel (plusieurs operations sont executes en sequence) +- Priorité via `push_first()` pour play/stop/volume +- `clear()` permet d'annuler les operations en attente (ex: playlist changée) + +**Implémentation suggérée**: +1. Créer `src/queue/openhome_op_queue.rs` avec la structure +2. Intégrer dans `OpenHomeQueue` ou `OpenHomeRenderer` +3. Thread de worker lancé au démarrage du control point + +--- + +### OPT-3: Métadonnées en Tâche de Fond + +**Observation**: L'appli utilise-t-elle vraiment les métadonnées de la queue OpenHome, ou un cache local? + +Si le control-point maintient son propre cache (plus probable): +- Les mises à jour de métadonnées peuvent être traitées en background +- Pas besoin de sync immédiate des métadonnées + +**Solution**: Queue séparée pour les operations de métadonnées: + +```rust +// Haute priorité (opérations critiques) +let high_priority_queue: OpenHomeOpQueue; + +// Basse priorité (métadonnées) +let metadata_queue: OpenHomeOpQueue; +``` + +**Implémentation**: +1. Séparer les operations critiques (play/stop/seek/volume) de metadata +2. Metadata update traités en background avec délais +3. Le cache local du control-point est mis à jour indépendamment + +--- + +### OPT-4: Connection Pooling HTTP + +Chaque appel SOAP crée une nouvelle connexion. Avec 1000 insertions: +- Overhead TCP: ~10-50ms par appel +- Total: 10-50 secondes overhead réseau + +**Solution**: Agent HTTP static avec connection reuse: + +```rust +// soap_client.rs +static HTTP_AGENT: Lazy = Lazy::new(|| { + Agent::config_builder() + .timeout_global(Some(Duration::from_secs(30))) + .build() +}); +``` + +--- + +## Plan d'Implémentation Proposé + +### Phase 1: Fast Path LCS (Priorité Haute) + +1. Ajouter `detect_sync_pattern()` dans `openhome.rs` +2. Implémenter `append_only()` et `delete_from_beginning()` +3. Tester avec playlists réelles + +### Phase 2: Queue FIFO Opérations (Priorité Haute) + +1. Créer `src/queue/openhome_op_queue.rs` +2. Implémenter `push()`, `push_first()`, `clear()` +3. Thread worker avec loop de consommation +4. Intégrer dans `OpenHomeRenderer` + +### Phase 3: Séparation Métadonnées (Priorité Moyenne) + +1. Créer queue séparée pour metadata +2. Implémenter batch processing + +### Phase 4: Connection Pooling (Priorité Basse) + +1. Modifier `soap_client.rs` pour agent static + +--- + +## Questions pour Clarification + +1. **Cache Métadonnées**: Le control-point utilise-t-il vraiment les métadonnées de la queue OpenHome, ou maintient-il son propre cache qui est alimenté indépendamment? + +2. **Priorité des Opérations**: Pour `push_first()`, quelles opérations nécessitent une réponse rapide? + - Volume (immédiat) + - Play/Pause/Stop (immédiat) + - Seek (rapide) + - Insert (peut être différé) + +3. **Comportement en cas de conflit**: Si le client fait `clear()` et que le worker est en train d'exécuter une opération: + - Annuler l'opération en cours? ( risky - peut laisser le renderer dans un état inconsistent) + - Laisser finir l'opération en cours? (plus sur) + +--- + +## Tests Recommandés + +```bash +# Compiler +cargo build -p pmocontrol + +# Benchmark LCS fast paths +# - Cas: append 100 tracks to 900 = O(N) +# - Cas: delete 100 from 900 = O(N) +# - Cas: reorder = O(N²) avec LCS +``` + +--- + +*Plan mis à jour avec contre-propositions utilisateur* +*Date: 2026-04-08* \ No newline at end of file diff --git a/Blackboard/Todo/enorme_playlist_step2.md b/Blackboard/Todo/enorme_playlist_step2.md new file mode 100644 index 00000000..0200e455 --- /dev/null +++ b/Blackboard/Todo/enorme_playlist_step2.md @@ -0,0 +1,235 @@ +# Optimisation Playlist OpenHome — Étape 2 + +**Contexte**: Suite de `enorme_playlist.md`. Les optimisations de base (MAX_BATCH=256, LCS prefixe/suffixe, polling adaptatif, caches TTL) sont faites. Les lenteurs persistent sur les grandes playlists (~1000 titres). Les renderers Chromecast et UPnP sont moins affectés mais bénéficieront également de certaines optimisations. + +**Contrainte architecturale fondamentale**: La queue OpenHome est la **source de vérité unique**. Un miroir local persistent a été tenté et abandonné — impossible à maintenir en sync quand d'autres control points (BubbleUPnP, Linn, etc.) modifient la queue. Toute optimisation doit respecter cette contrainte. + +**Note sur le LCS**: `lcs_flags_optimized` (openhome.rs:869) gère déjà le cas dominant (préfixe/suffixe communs). Pour un append de 100 tracks à 900 existants, le LCS est O(1) — ce n'est pas le goulot. Le coût réel est les appels SOAP : connexions TCP × (RTT + handshake) et les `ReadList` pour reconstruire le snapshot. + +--- + +## Analyse des Goulots Réels + +Pour un `sync_queue` "append 100 tracks à 900 existants" aujourd'hui : + +| Étape | Appels SOAP | Coût estimé | +|-------|------------|-------------| +| `id_array()` | 1 | ~5ms | +| `read_list()` — 4 batches × 256 | 4 | ~20ms | +| 100 × `insert()` | 100 | ~100 × (RTT + **TCP handshake**) | +| **Total TCP handshakes** | 105 | **105 × 10-50ms = 1-5 secondes** | + +Les deux leviers : (1) éliminer les handshakes TCP, (2) éliminer les `ReadList` quand inutiles. + +--- + +## Plan d'Implémentation + +### Phase 1 — Connection Pooling (1-2h, PRIORITÉ MAXIMALE) + +**Fichier**: `pmocontrol/src/soap_client.rs` + +**Problème**: lignes 65-72 créent un nouvel `ureq::Agent` à chaque appel SOAP = nouvelle connexion TCP à chaque fois. + +**Solution**: Agent statique partagé via `OnceLock`. + +```rust +use std::sync::OnceLock; + +static SOAP_AGENT: OnceLock = OnceLock::new(); + +fn get_soap_agent() -> &'static ureq::Agent { + SOAP_AGENT.get_or_init(|| { + ureq::Agent::config_builder() + .http_status_as_error(false) + .timeout_global(Some(Duration::from_secs(30))) + .build() + .into() + }) +} +``` + +Dans `invoke_upnp_action_with_timeout()`, remplacer la construction de l'agent par : + +```rust +// Cas normal : réutiliser l'agent partagé (keep-alive, connection pooling) +// Cas custom timeout : agent dédié (rare — timeout global suffit en pratique) +let agent_owned; +let agent: &ureq::Agent = if timeout.is_some() { + agent_owned = ureq::Agent::config_builder() + .http_status_as_error(false) + .timeout_global(timeout) + .build() + .into(); + &agent_owned +} else { + get_soap_agent() +}; +``` + +**Vérifier** que `ureq::Agent` maintient bien un pool de connexions HTTP/1.1 keep-alive entre les appels (comportement documenté de ureq v3 — l'agent est conçu pour être réutilisé). + +**Impact**: Élimine les TCP handshakes répétés pour tous les renderers (OpenHome, UPnP, Chromecast). Pour 100 inserts : 99 handshakes économisés × 10-50ms = **1-5 secondes récupérées**. + +--- + +### Phase 2 — Fast Path Session (2-3h, PRIORITÉ HAUTE) + +**Concept**: Sans miroir persistant (abandonné), on peut quand même éviter les `ReadList` dans le cas dominant en maintenant un état **éphémère de session** — valable uniquement entre deux `sync_queue` consécutifs, et invalidé dès qu'on détecte une incohérence. + +**Principe**: Après chaque `sync_queue`, mémoriser : +- la liste d'IDs résultante (déjà dans `track_ids_cache` TTL 1s) +- le `after_id` du dernier insert (pour pouvoir appender sans `id_array`) + +Au prochain `sync_queue`, tenter de détecter le pattern sans `ReadList` : + +```rust +fn try_fast_path(&self, new_items: &[PlaybackItem]) -> FastPathResult { + // Récupérer les IDs actuels (cache ou 1 appel id_array) + let current_ids = self.track_ids()?; + let current_len = current_ids.len(); + let new_len = new_items.len(); + + // Fast path 1: append only + // Condition: new_items a plus d'items, et les current_len premiers de new_items + // ont les mêmes didl_id que les items actuels (vérifiable depuis id_array + metadata_cache local) + if new_len > current_len { + let prefix_matches = self.check_prefix_matches(¤t_ids, &new_items[..current_len]); + if prefix_matches { + return FastPathResult::AppendOnly { items: &new_items[current_len..] }; + } + } + + // Fast path 2: delete from end + if new_len < current_len { + let prefix_matches = self.check_prefix_matches(¤t_ids[..new_len], new_items); + if prefix_matches { + let to_delete = ¤t_ids[new_len..]; + return FastPathResult::DeleteFromEnd { ids: to_delete }; + } + } + + // Cas général: fallback ReadList + LCS + FastPathResult::NeedFullSync +} +``` + +**`check_prefix_matches`** : compare `current_ids[i]` avec `new_items[i]` en utilisant le `metadata_cache` local (déjà en mémoire) pour résoudre les URIs/didl_ids des IDs connus. Si un ID n'est pas en cache → fast path impossible → fallback. + +**Clé**: Cette vérification utilise uniquement le `metadata_cache` local (HashMap en mémoire, nano-secondes) et `id_array()` (déjà caché TTL 1s). Zéro appel `ReadList` dans le cas heureux. Si la vérification échoue (incohérence détectée, cache manquant) → fallback propre vers `ReadList` + LCS, source de vérité OpenHome préservée. + +**Fichier**: `pmocontrol/src/queue/openhome.rs`, ajouter `try_fast_path()` et l'intégrer en début de `sync_queue()`. + +--- + +### Phase 3 — Queue FIFO Async (8-12h, PRIORITÉ MOYENNE) + +**Prérequis**: Phases 1 et 2 complétées. + +**Concept**: Exécuter les opérations SOAP dans un thread dédié pour rendre `sync_queue()` non-bloquant du point de vue de l'appelant. + +**Contrainte technique**: Les `insert()` sont chaînés — chaque appel retourne un `new_id` utilisé comme `after_id` du suivant. Le worker doit maintenir cet état interne. + +**Fichier à créer**: `pmocontrol/src/queue/openhome_op_queue.rs` + +```rust +pub enum OpenHomeOp { + /// Insert séquentiel — after_id géré en interne (last_inserted_id) + InsertAtEnd { uri: String, metadata: String, didl_id: String }, + /// Insert après un ID connu (ex: après le pivot) + InsertAfter { after_id: u32, uri: String, metadata: String, didl_id: String }, + DeleteId { track_id: u32 }, + DeleteAll, + SeekId { id: u32 }, + Play, + Pause, + Stop, + SetVolume { volume: u16 }, +} + +pub struct OpenHomeOpQueue { + sender: mpsc::Sender, + last_error: Arc>>, + completion: Arc<(Mutex, Condvar)>, +} + +impl OpenHomeOpQueue { + pub fn push(&self, op: OpenHomeOp) { ... } + /// Opérations critiques passent devant (play/stop/volume) + pub fn push_priority(&self, op: OpenHomeOp) { ... } + /// Vider la file d'attente (ex: nouvelle playlist demandée avant fin de sync) + pub fn clear_pending(&self) { ... } + /// Attendre que toutes les opérations soient exécutées + pub fn wait_completion(&self) { ... } + /// Récupérer la dernière erreur (non-bloquant) + pub fn take_error(&self) -> Option { ... } +} +``` + +**Comportement en cas d'erreur**: vider la file d'attente, signaler l'erreur via `last_error`, invalider les caches OpenHome (forcer re-sync depuis source de vérité au prochain appel). + +**Comportement en cas de `clear_pending()` pendant exécution**: laisser l'opération en cours se terminer (plus sûr — évite de laisser le renderer dans un état inconsistant), vider le reste. + +**Intégration dans `OpenHomeQueue`**: remplacer les appels directs `playlist_client.insert()` / `playlist_client.delete_id()` par des `op_queue.push()`. Les opérations qui ont besoin d'une réponse synchrone (ex: `current_track()`, `queue_snapshot()`) continuent d'appeler directement le `playlist_client` — mais doivent d'abord attendre la complétion de la file (`wait_completion()`). + +--- + +### Phase 4 — Throttle des replace_item (2-3h, PRIORITÉ BASSE) + +**Contexte**: `replace_item()` (openhome.rs:1315) fait `delete_id` + `insert` pour mettre à jour une piste. Sur un stream radio qui change de morceau, la durée est mise à jour fréquemment → paires SOAP inutiles car le `metadata_cache` local est déjà la source pour l'UI. + +**Solution**: Ne pas envoyer le `delete_id` + `insert` OpenHome si une mise à jour pour ce `track_id` a déjà été envoyée dans les N dernières secondes. Le `metadata_cache` local est mis à jour immédiatement (pour l'UI), et l'opération OpenHome est différée ou ignorée. + +```rust +fn replace_item(&mut self, index: usize, item: PlaybackItem) -> Result<(), ControlPointError> { + let track_id = self.track_ids()?[index]; + + // Toujours mettre à jour le cache local immédiatement (pour l'UI) + self.cache_metadata(track_id, item.metadata.clone()); + + // Throttle: si replace OpenHome récent pour ce track, sauter l'opération SOAP + if self.is_recent_replace(track_id, Duration::from_secs(5)) { + return Ok(()); + } + self.mark_replace_time(track_id); + + // Opération SOAP (delete + insert) + // ... code existant ... +} +``` + +Ajouter `last_replace_times: Mutex>` dans `OpenHomeQueue`. + +--- + +## Ordre d'Implémentation + +``` +Phase 1 → mesurer gain TCP → Phase 2 → mesurer gain ReadList → Phase 3 → Phase 4 +``` + +Tester chaque phase avec des playlists réelles (~1000 tracks) avant de poursuivre. Ne pas combiner les phases pour pouvoir isoler les régressions. + +## Tests + +```bash +# Phase 1 — vérifier connexions TCP réutilisées +# tcpdump -i lo port 60000 -c 200 (ou port du renderer) +# Avant: SYN à chaque appel SOAP +# Après: SYN unique, flux keep-alive + +# Phase 2 — vérifier fast path activé +# RUST_LOG=debug cargo run ... 2>&1 | grep "fast path" +# Cas append 100 → "fast path: AppendOnly, 100 inserts" +# Cas delete 100 → "fast path: DeleteFromEnd, 100 deletes" +# Cas reorder → "fast path: NeedFullSync, falling back to ReadList+LCS" + +# Phase 3 — vérifier non-blocage +# sync_queue() doit retourner en <1ms (les 100 inserts continuent en background) +``` + +--- + +*Contrainte architecturale intégrée: miroir local abandonné (désynchronisation avec autres control points). Toutes les phases respectent OpenHome comme source de vérité unique.* + +*Date: 2026-04-08* diff --git a/pmoapp/webapp/package-lock.json b/pmoapp/webapp/package-lock.json index a122e80a..7904c3eb 100644 --- a/pmoapp/webapp/package-lock.json +++ b/pmoapp/webapp/package-lock.json @@ -1609,9 +1609,9 @@ "license": "MIT" }, "node_modules/vite": { - "version": "7.3.1", - "resolved": "https://registry.npmjs.org/vite/-/vite-7.3.1.tgz", - "integrity": "sha512-w+N7Hifpc3gRjZ63vYBXA56dvvRlNWRczTdmCBBa+CotUzAPf5b7YMdMR/8CQoeYE5LX3W4wj6RYTgonm1b9DA==", + "version": "7.3.2", + "resolved": "https://registry.npmjs.org/vite/-/vite-7.3.2.tgz", + "integrity": "sha512-Bby3NOsna2jsjfLVOHKes8sGwgl4TT0E6vvpYgnAYDIF/tie7MRaFthmKuHx1NSXjiTueXH3do80FMQgvEktRg==", "dev": true, "license": "MIT", "dependencies": { diff --git a/pmocontrol/src/queue/openhome.rs b/pmocontrol/src/queue/openhome.rs index 4007535b..6e0dc3c5 100644 --- a/pmocontrol/src/queue/openhome.rs +++ b/pmocontrol/src/queue/openhome.rs @@ -170,6 +170,8 @@ pub struct OpenHomeQueue { /// Permet de maintenir des métadonnées à jour même si le service OpenHome /// ne permet pas de les modifier directement. metadata_cache: Mutex>>, + /// Cache URI by track ID for fast path matching + uri_by_id: Mutex>, /// Cache for track IDs to avoid redundant IdArray SOAP calls track_ids_cache: Arc>, /// Cache for current track ID to avoid redundant Id SOAP calls @@ -178,6 +180,16 @@ pub struct OpenHomeQueue { read_list_cache: Arc>, } +/// Fast path result for queue sync optimization +enum FastPathResult { + /// New items are an append to the current queue + AppendOnly { new_items: Vec }, + /// Items were deleted from the end + DeleteFromEnd { delete_ids: Vec }, + /// No fast path possible - need full LCS sync + NeedFullSync, +} + impl OpenHomeQueue { pub fn new( renderer_id: DeviceId, @@ -191,6 +203,7 @@ impl OpenHomeQueue { info_client, product_client, metadata_cache: Mutex::new(HashMap::new()), + uri_by_id: Mutex::new(HashMap::new()), track_ids_cache: Arc::new(Mutex::new(TrackIdsCache::new())), current_track_id_cache: Arc::new(Mutex::new(CurrentTrackIdCache::new())), read_list_cache: Arc::new(Mutex::new(ReadListCache::new())), @@ -229,6 +242,117 @@ impl OpenHomeQueue { self.track_ids_cache.lock().unwrap().invalidate(); self.read_list_cache.lock().unwrap().invalidate(); self.current_track_id_cache.lock().unwrap().invalidate(); + self.uri_by_id.lock().unwrap().clear(); + } + + /// Tries to detect a simple append-only or delete-from-end pattern without ReadList. + /// + /// This is a fast path optimization that avoids the expensive ReadList calls when: + /// 1. New items are an append to existing queue (only inserts needed) + /// 2. Items were deleted from the end (only deletes needed) + /// + /// Returns `FastPathResult::NeedFullSync` if the pattern doesn't match or + /// if cache data is insufficient to verify. + fn try_fast_path(&self, new_items: &[PlaybackItem]) -> FastPathResult { + let current_ids = match self.track_ids() { + Ok(ids) => ids, + Err(e) => { + debug!( + renderer = self.renderer_id.0.as_str(), + error = %e, + "try_fast_path: failed to get current IDs, falling back to full sync" + ); + return FastPathResult::NeedFullSync; + } + }; + + let current_len = current_ids.len(); + let new_len = new_items.len(); + + if new_len > current_len { + let prefix_matches = self.check_prefix_matches(¤t_ids, &new_items[..current_len]); + if prefix_matches { + debug!( + renderer = self.renderer_id.0.as_str(), + current_len, new_len, "try_fast_path: detected append-only pattern" + ); + return FastPathResult::AppendOnly { + new_items: new_items[current_len..].to_vec(), + }; + } + } else if new_len < current_len { + let prefix_matches = self.check_prefix_matches(¤t_ids[..new_len], new_items); + if prefix_matches { + let delete_ids: Vec = current_ids[new_len..].to_vec(); + debug!( + renderer = self.renderer_id.0.as_str(), + current_len, + new_len, + delete_count = delete_ids.len(), + "try_fast_path: detected delete-from-end pattern" + ); + return FastPathResult::DeleteFromEnd { delete_ids }; + } + } + + debug!( + renderer = self.renderer_id.0.as_str(), + current_len, new_len, "try_fast_path: no fast path detected, falling back to full sync" + ); + FastPathResult::NeedFullSync + } + + /// Checks if the prefix of the new items matches the current queue items. + /// Uses local metadata cache or falls back to URI comparison. + fn check_prefix_matches(&self, current_ids: &[u32], new_items: &[PlaybackItem]) -> bool { + if current_ids.len() != new_items.len() { + return false; + } + + let metadata_cache = match self.metadata_cache.lock() { + Ok(cache) => cache, + Err(_) => return false, + }; + + let uri_cache = match self.uri_by_id.lock() { + Ok(cache) => cache, + Err(_) => return false, + }; + + for (idx, new_item) in new_items.iter().enumerate() { + let current_id = current_ids[idx]; + + // First check if we have cached metadata for this track_id that matches + if let Some(cached_metadata) = metadata_cache.get(¤t_id) { + if let Some(cached) = cached_metadata { + // Compare using title + artist as identity (for streams) + if let Some(ref new_metadata) = new_item.metadata { + let same_title = cached.title.as_ref() == new_metadata.title.as_ref(); + let same_artist = cached.artist.as_ref() == new_metadata.artist.as_ref(); + if same_title && same_artist { + continue; + } + } + } + } + + // No metadata match - fall back to URI comparison using uri_by_id cache + if new_item.uri.is_empty() { + return false; + } + + // Compare against cached URI for this track_id + if let Some(cached_uri) = uri_cache.get(¤t_id) { + if cached_uri == &new_item.uri { + continue; + } + } + + // No match found - fall back to full sync + return false; + } + + true } /// Met à jour les métadonnées d'un item de la queue à l'index spécifié. @@ -250,7 +374,7 @@ impl OpenHomeQueue { metadata: Option, ) -> Result<(), ControlPointError> { let track_id = self.position_to_id(index)?; - self.cache_metadata(track_id, metadata); + self.cache_metadata(track_id, metadata, ""); Ok(()) } @@ -260,8 +384,19 @@ impl OpenHomeQueue { /// chanson donc toute durée est acceptée. /// Pour les fichiers normaux (non-streams), les métadonnées sont acceptées telles quelles. /// Cette fonction est le SEUL point d'entrée pour modifier le cache. - fn cache_metadata(&self, track_id: u32, new_metadata: Option) { + fn cache_metadata( + &self, + track_id: u32, + new_metadata: Option, + uri: &str, + ) { let mut cache = self.metadata_cache.lock().unwrap(); + let mut uri_cache = self.uri_by_id.lock().unwrap(); + + // Update URI cache + if !uri.is_empty() { + uri_cache.insert(track_id, uri.to_string()); + } // Vérifier s'il y a déjà des métadonnées en cache if let Some(cached_meta) = cache.get(&track_id) { @@ -366,7 +501,7 @@ impl OpenHomeQueue { ); drop(cache); // Libérer le lock avant d'appeler cache_metadata // Mettre en cache pour éviter les oscillations sur les flux radio - self.cache_metadata(entry.id, fresh.clone()); + self.cache_metadata(entry.id, fresh.clone(), entry.uri()); fresh } }; @@ -397,7 +532,7 @@ impl OpenHomeQueue { .insert(after_id, &item.uri, &metadata_xml)?; // Enregistrer les métadonnées dans le cache - self.cache_metadata(new_id, item.metadata); + self.cache_metadata(new_id, item.metadata, &item.uri); Ok(new_id) } @@ -432,7 +567,7 @@ impl OpenHomeQueue { .insert(previous_id, &item.uri, &metadata)?; // Enregistrer les métadonnées dans le cache - self.cache_metadata(new_id, item.metadata); + self.cache_metadata(new_id, item.metadata, &item.uri); previous_id = new_id; } @@ -498,7 +633,7 @@ impl OpenHomeQueue { previous_id = existing_id; // Mettre à jour les métadonnées de l'item existant conservé - self.cache_metadata(existing_id, item.metadata.clone()); + self.cache_metadata(existing_id, item.metadata.clone(), &item.uri); debug!( renderer = self.renderer_id.0.as_str(), @@ -514,7 +649,7 @@ impl OpenHomeQueue { .insert(previous_id, &item.uri, &metadata)?; // Enregistrer les métadonnées du nouvel item - self.cache_metadata(new_id, item.metadata.clone()); + self.cache_metadata(new_id, item.metadata.clone(), &item.uri); debug!( renderer = self.renderer_id.0.as_str(), @@ -590,7 +725,11 @@ impl OpenHomeQueue { let previous_id = pivot_id as u32; // Mettre à jour les métadonnées du pivot - self.cache_metadata(previous_id, new_items[pivot_idx_new].metadata.clone()); + self.cache_metadata( + previous_id, + new_items[pivot_idx_new].metadata.clone(), + &new_items[pivot_idx_new].uri, + ); debug!( renderer = self.renderer_id.0.as_str(), @@ -731,7 +870,7 @@ impl OpenHomeQueue { // Mettre à jour les métadonnées de l'item existant conservé // La fonction cache_metadata gère la protection contre la diminution de durée - self.cache_metadata(existing_id, item.metadata); + self.cache_metadata(existing_id, item.metadata, &item.uri); } else { let metadata = build_metadata_xml(&item); let new_id = self @@ -739,7 +878,7 @@ impl OpenHomeQueue { .insert(previous_id, &item.uri, &metadata)?; // Enregistrer les métadonnées du nouvel item - self.cache_metadata(new_id, item.metadata); + self.cache_metadata(new_id, item.metadata, &item.uri); previous_id = new_id; } @@ -1124,7 +1263,7 @@ impl QueueBackend for OpenHomeQueue { .insert(previous_id, &item.uri, &metadata)?; // Enregistrer les métadonnées dans le cache - self.cache_metadata(new_id, item.metadata); + self.cache_metadata(new_id, item.metadata, &item.uri); previous_id = new_id; } @@ -1167,6 +1306,57 @@ impl QueueBackend for OpenHomeQueue { return Ok(()); } + // Fast path: try to detect simple append-only or delete-from-end patterns + // without the expensive ReadList call that queue_snapshot() would trigger + match self.try_fast_path(&items) { + FastPathResult::AppendOnly { new_items } => { + debug!( + renderer = self.renderer_id.0.as_str(), + new_items_count = new_items.len(), + "sync_queue: fast path - append-only detected" + ); + // Insert new items at the end, using the returned new_id as after_id for next insert + // Start after the last existing track (cached, no SOAP call needed) + let current_ids = self.track_ids()?; + let mut after_id = current_ids.last().copied().unwrap_or(OPENHOME_PLAYLIST_HEAD_ID); + let mut uri_cache = self.uri_by_id.lock().unwrap(); + for item in &new_items { + let metadata = build_metadata_xml(item); + let uri = item.uri.as_str(); + let metadata_xml = metadata.as_str(); + after_id = self.playlist_client.insert(after_id, uri, metadata_xml)?; + // Update URI cache + if !item.uri.is_empty() { + uri_cache.insert(after_id, item.uri.clone()); + } + } + drop(uri_cache); + // Invalidate caches after inserts + self.invalidate_track_caches(); + return Ok(()); + } + FastPathResult::DeleteFromEnd { delete_ids } => { + debug!( + renderer = self.renderer_id.0.as_str(), + delete_count = delete_ids.len(), + "sync_queue: fast path - delete-from-end detected" + ); + // Delete items from the end + for id in delete_ids.iter().rev() { + self.playlist_client.delete_id(*id)?; + } + // Invalidate caches after deletes + self.invalidate_all_caches(); + return Ok(()); + } + FastPathResult::NeedFullSync => { + debug!( + renderer = self.renderer_id.0.as_str(), + "sync_queue: fast path not applicable, proceeding with full sync" + ); + } + } + // Synchronize local state with the actual OpenHome playlist before computing // differences. Without this, any drift between our cache and the renderer // (e.g., manual edits from another control point) would keep the stale items. @@ -1344,7 +1534,8 @@ impl QueueBackend for OpenHomeQueue { // Mettre à jour le cache avec les nouvelles métadonnées self.metadata_cache.lock().unwrap().remove(&track_id); - self.cache_metadata(new_id, item.metadata); + self.uri_by_id.lock().unwrap().remove(&track_id); + self.cache_metadata(new_id, item.metadata, &item.uri); if ci == Some(index) { self.playlist_client.seek_id(new_id)?; diff --git a/pmocontrol/src/soap_client.rs b/pmocontrol/src/soap_client.rs index 580bfbdc..4b91d4f6 100644 --- a/pmocontrol/src/soap_client.rs +++ b/pmocontrol/src/soap_client.rs @@ -1,12 +1,30 @@ +use std::sync::Arc; +use std::sync::OnceLock; use std::time::Duration; use anyhow::{Context, Result}; -use pmoupnp::soap::{SoapEnvelope, build_soap_request, parse_soap_envelope}; +use pmoupnp::soap::{build_soap_request, parse_soap_envelope, SoapEnvelope}; use tracing::{debug, trace, warn}; use ureq::Agent; use crate::errors::ControlPointError; +static SOAP_AGENT: OnceLock> = OnceLock::new(); + +fn get_soap_agent() -> Arc { + SOAP_AGENT + .get_or_init(|| { + Arc::new( + Agent::config_builder() + .http_status_as_error(false) + .timeout_global(Some(Duration::from_secs(30))) + .build() + .into(), + ) + }) + .clone() +} + /// Result of a SOAP call: /// - HTTP status code /// - raw XML body (always) @@ -62,21 +80,20 @@ pub fn invoke_upnp_action_with_timeout( trace!(body = body_xml.as_str(), "SOAP request body"); - let mut builder = Agent::config_builder(); - builder = builder.http_status_as_error(false); - if let Some(duration) = timeout { - builder = builder.timeout_global(Some(duration)); - } + let request = if let Some(duration) = timeout { + let config = Agent::config_builder() + .http_status_as_error(false) + .timeout_global(Some(duration)) + .build(); + let agent: Agent = config.into(); + agent.post(control_url) + } else { + get_soap_agent().post(control_url) + }; - let config = builder.build(); - let agent: Agent = config.into(); - - // 3. SOAPAction header let soap_action_header = format!(r#""{}#{}""#, service_type, action); - // 4. HTTP POST - let mut response = agent - .post(control_url) + let mut response = request .header("Content-Type", r#"text/xml; charset="utf-8""#) .header("SOAPAction", &soap_action_header) .send(body_xml)