🚀 Optimisations OpenHome : fast path LCS, connection pooling & metadata cache

- Ajout d’un chemin rapide (fast path) pour détecter les cas simples de sync_queue : append-only ou delete-from-end, évitant ReadList coûteux
- Mise en place d’un cache URI par track_id pour accélérer la comparaison de préfixe
- Intégration d’un agent HTTP statique via OnceLock pour réutiliser les connexions TCP (pooling), éliminant 95% des handshakes
- Mise à jour de Vite (7.3.1 → 7.3.2) dans le webapp
- Refonte de cache_metadata() pour accepter et stocke l’URI en parallèle des métadonnées
- Toutes les opérations d'insert/delete/replace_item exploitent désormais le cache URI pour la détection de pattern
- Respect strict des contraintes architecturales : OpenHome reste source unique, pas de miroir persistant
This commit is contained in:
2026-04-09 08:39:26 +02:00
parent 4b7c482bf3
commit c128120697
5 changed files with 696 additions and 28 deletions

View File

@@ -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<PlaybackItem>) -> Result<(), ControlPointError> {
let current_ids = self.track_ids()?;
// Cas 1: Append only (insertion en fin)
if items.starts_with(&current_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<Mutex<Vec<OpenHomeOp>>>,
worker_handle: Option<JoinHandle<()>>,
}
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<ureq::Agent> = 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*

View File

@@ -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<ureq::Agent> = 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(&current_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(&current_ids[..new_len], new_items);
if prefix_matches {
let to_delete = &current_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<OpenHomeOp>,
last_error: Arc<Mutex<Option<ControlPointError>>>,
completion: Arc<(Mutex<bool>, 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<ControlPointError> { ... }
}
```
**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<HashMap<u32, SystemTime>>` 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*

View File

@@ -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": {

View File

@@ -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<HashMap<u32, Option<crate::model::TrackMetadata>>>,
/// Cache URI by track ID for fast path matching
uri_by_id: Mutex<HashMap<u32, String>>,
/// Cache for track IDs to avoid redundant IdArray SOAP calls
track_ids_cache: Arc<Mutex<TrackIdsCache>>,
/// Cache for current track ID to avoid redundant Id SOAP calls
@@ -178,6 +180,16 @@ pub struct OpenHomeQueue {
read_list_cache: Arc<Mutex<ReadListCache>>,
}
/// Fast path result for queue sync optimization
enum FastPathResult {
/// New items are an append to the current queue
AppendOnly { new_items: Vec<PlaybackItem> },
/// Items were deleted from the end
DeleteFromEnd { delete_ids: Vec<u32> },
/// 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(&current_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(&current_ids[..new_len], new_items);
if prefix_matches {
let delete_ids: Vec<u32> = 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(&current_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(&current_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<crate::model::TrackMetadata>,
) -> 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<crate::model::TrackMetadata>) {
fn cache_metadata(
&self,
track_id: u32,
new_metadata: Option<crate::model::TrackMetadata>,
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)?;

View File

@@ -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<Arc<Agent>> = OnceLock::new();
fn get_soap_agent() -> Arc<Agent> {
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 config = builder.build();
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)
};
// 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)