From 4b7c482bf3d38f86e160680af39c31370fe52a04 Mon Sep 17 00:00:00 2001 From: Eric Coissac Date: Wed, 8 Apr 2026 09:01:05 +0200 Subject: [PATCH 1/3] [webapp] Add debug logging to search functionality - Bump Cargo.lock version from v0.3.41 to v0.3.52 - Add console.log statements in ServerDrawer.vue and useMediaServers.ts to trace search flow: handleSearch call, server ID/query passed in API request - Include logging for number of entries returned by search --- Cargo.lock | 2 +- pmoapp/webapp/src/components/unified/ServerDrawer.vue | 1 + pmoapp/webapp/src/composables/useMediaServers.ts | 3 +++ 3 files changed, 5 insertions(+), 1 deletion(-) diff --git a/Cargo.lock b/Cargo.lock index 843135d5..664d3596 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4,7 +4,7 @@ version = 4 [[package]] name = "PMOMusic" -version = "0.3.39" +version = "0.3.40" dependencies = [ "axum 0.8.7", "console-subscriber", diff --git a/pmoapp/webapp/src/components/unified/ServerDrawer.vue b/pmoapp/webapp/src/components/unified/ServerDrawer.vue index 19a7a0df..0f2783e0 100644 --- a/pmoapp/webapp/src/components/unified/ServerDrawer.vue +++ b/pmoapp/webapp/src/components/unified/ServerDrawer.vue @@ -51,6 +51,7 @@ const searchInput = ref(''); const isSearchMode = computed(() => searchQuery.value !== ''); async function handleSearch() { + console.log('[ServerDrawer] handleSearch called, currentServer:', currentServer.value?.id, 'searchInput:', searchInput.value); if (!currentServer.value || !searchInput.value.trim()) return; await searchServer(currentServer.value.id, searchInput.value.trim()); } diff --git a/pmoapp/webapp/src/composables/useMediaServers.ts b/pmoapp/webapp/src/composables/useMediaServers.ts index dd71f314..a9f366a5 100644 --- a/pmoapp/webapp/src/composables/useMediaServers.ts +++ b/pmoapp/webapp/src/composables/useMediaServers.ts @@ -216,6 +216,7 @@ export function useMediaServers() { // Recherche dans un serveur async function searchServer(serverId: string, query: string) { + console.log(`[useMediaServers] searchServer called: serverId=${serverId}, query=${query}`); if (!query.trim()) { searchResults.value = null searchQuery.value = '' @@ -226,8 +227,10 @@ export function useMediaServers() { loading.value = true error.value = null searchQuery.value = query + console.log(`[useMediaServers] Calling API searchServer for server ${serverId}`); const data = await api.searchServer(serverId, query) + console.log(`[useMediaServers] Search returned ${data.entries.length} entries, total=${data.total_count}`); searchResults.value = { container_id: 'search', entries: data.entries, From c128120697fb67059c83a29e67388d9105e285eb Mon Sep 17 00:00:00 2001 From: Eric Coissac Date: Thu, 9 Apr 2026 08:39:26 +0200 Subject: [PATCH 2/3] :rocket: Optimisations OpenHome : fast path LCS, connection pooling & metadata cache MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 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 --- .kilo/plans/1775631632458-happy-cactus.md | 225 +++++++++++++++++++++ Blackboard/Todo/enorme_playlist_step2.md | 235 ++++++++++++++++++++++ pmoapp/webapp/package-lock.json | 6 +- pmocontrol/src/queue/openhome.rs | 215 ++++++++++++++++++-- pmocontrol/src/soap_client.rs | 43 ++-- 5 files changed, 696 insertions(+), 28 deletions(-) create mode 100644 .kilo/plans/1775631632458-happy-cactus.md create mode 100644 Blackboard/Todo/enorme_playlist_step2.md 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) From 9dee19394709fc2d76dcfea202c6d1e46af584c4 Mon Sep 17 00:00:00 2001 From: Eric Coissac Date: Thu, 9 Apr 2026 11:19:54 +0200 Subject: [PATCH 3/3] =?UTF-8?q?:arrow=5Fup:=20version=20to=20v0.3.41=20?= =?UTF-8?q?=E2=80=94=20async=20queue=20sync=20refactoring?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Add `SyncCancelled` error variant for non-fatal cancellations - Refactor queue sync to async via `MusicQueue::schedule_sync()` - Extract browse+conversion into internal helper - Update all `QueueBackend::sync_queue()` signatures to accept cancel token and on_ready callback - Implement early-start logic (on pivot preservation or first insert) - Add `QueueReadyToPlay` and ‘ QueueSyncCancelled➔ SSE events - Replace blocking `refresh_attached_queue_for()` with non-blocking async dispatch in control_point.rs - Bump version to v0.3.41 --- Blackboard/Todo/enorme_playlist_step3.md | 642 ++++++++++++++++++ Cargo.lock | 2 +- PMOMusic/Cargo.toml | 2 +- pmocontrol/src/control_point.rs | 254 ++++++- pmocontrol/src/errors.rs | 2 + pmocontrol/src/model.rs | 6 + pmocontrol/src/music_renderer/arylic_tcp.rs | 22 +- .../src/music_renderer/chromecast_renderer.rs | 13 +- .../src/music_renderer/linkplay_renderer.rs | 20 +- .../src/music_renderer/musicrenderer.rs | 42 +- .../src/music_renderer/openhome_renderer.rs | 19 +- .../src/music_renderer/upnp_renderer.rs | 18 +- pmocontrol/src/queue/backend.rs | 10 +- pmocontrol/src/queue/interne.rs | 66 +- pmocontrol/src/queue/mod.rs | 14 +- pmocontrol/src/queue/music_queue.rs | 356 ++++++++-- pmocontrol/src/queue/openhome.rs | 180 +++-- pmocontrol/src/queue/snapshot.rs | 2 +- pmocontrol/src/registry.rs | 8 + pmocontrol/src/sse.rs | 16 + version.txt | 2 +- 21 files changed, 1455 insertions(+), 241 deletions(-) create mode 100644 Blackboard/Todo/enorme_playlist_step3.md diff --git a/Blackboard/Todo/enorme_playlist_step3.md b/Blackboard/Todo/enorme_playlist_step3.md new file mode 100644 index 00000000..5246fae1 --- /dev/null +++ b/Blackboard/Todo/enorme_playlist_step3.md @@ -0,0 +1,642 @@ +** 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) ** + +# Async Queue Refresh — Étape 3 + +**Contexte**: `refresh_attached_queue_for()` dans `control_point.rs` est appelé de 3 endroits +et bloque son thread pendant toute la synchronisation (browse media server + 100+ opérations +SOAP/mémoire). L'objectif est de factoriser le mécanisme async au niveau de la couche queue +(`MusicQueue`), qui est déjà l'abstraction agnostique du backend. `control_point.rs` ne doit +plus connaître les threads ni les tokens d'annulation. + +**Principe architectural**: la couche queue sait *comment* syncer (mémoire ou SOAP) et donc +aussi *comment* annuler et *quand* signaler que la lecture peut démarrer. `control_point.rs` +sait seulement *quoi* syncer (browse + conversion PlaybackItem). Les deux responsabilités +restent séparées. + +--- + +## Vue d'ensemble des changements + +``` +AVANT + control_point.rs + refresh_attached_queue_for() + → browse() + → sync_queue(items) ← bloquant 1-5s + +APRÈS + control_point.rs + do_queue_refresh_work() ← interne, fait le browse + conversion + MusicQueue (couche queue) + schedule_sync(items, callbacks) ← non-bloquant, retourne immédiatement + → thread "queue-sync-{renderer_id}" + → QueueBackend::sync_queue(items, cancel_token, on_ready) +``` + +--- + +## Fichiers à modifier / créer + +| Fichier | Action | +|---------|--------| +| `pmocontrol/src/errors.rs` | Ajouter variante `SyncCancelled` | +| `pmocontrol/src/queue/backend.rs` | Modifier signature `sync_queue()` | +| `pmocontrol/src/queue/interne.rs` | Adapter signature `sync_queue()` | +| `pmocontrol/src/queue/openhome.rs` | Adapter + points de vérification cancel + on_ready | +| `pmocontrol/src/queue/music_queue.rs` | Ajouter champs async + méthode `schedule_sync()` | +| `pmocontrol/src/queue/mod.rs` | Exporter `SyncScheduleOutcome` | +| `pmocontrol/src/model.rs` | Ajouter événements `QueueReadyToPlay`, `QueueSyncCancelled` | +| `pmocontrol/src/sse.rs` | Sérialiser les deux nouveaux événements | +| `pmocontrol/src/control_point.rs` | Remplacer les 3 call sites bloquants | + +--- + +## Étape 1 — Nouvelle variante d'erreur (`errors.rs`) + +**Fichier**: `pmocontrol/src/errors.rs` + +Ajouter après la ligne 51 (`ControlPoint`) : + +```rust +#[error("Queue sync cancelled (superseded by a newer request)")] +SyncCancelled, +``` + +Cette variante est retournée par `sync_queue()` quand le `cancel_token` passe à `true`. +Elle est **non-fatale** — le coordinator la traite comme un comportement normal, pas une +erreur à logger en `warn!`. + +--- + +## Étape 2 — Modifier la signature de `sync_queue()` dans le trait (`backend.rs`) + +**Fichier**: `pmocontrol/src/queue/backend.rs`, ligne 110 + +```rust +// AVANT +fn sync_queue(&mut self, items: Vec) -> Result<(), ControlPointError>; + +// APRÈS +use std::sync::{Arc, atomic::AtomicBool}; + +fn sync_queue( + &mut self, + items: Vec, + cancel_token: &Arc, + on_ready: Option>, +) -> Result<(), ControlPointError>; +``` + +**Sémantique des paramètres** : +- `cancel_token` : si `true` au moment d'une opération, retourner `Err(SyncCancelled)` immédiatement +- `on_ready` : callback one-shot appelé quand la lecture peut démarrer (voir logique ci-dessous) + +--- + +## Étape 3 — Adapter `InternalQueue::sync_queue()` (`interne.rs`) + +**Fichier**: `pmocontrol/src/queue/interne.rs` + +Trouver la méthode `sync_queue()` et adapter la signature. Le corps reste identique, +avec deux ajouts : + +**1. Vérification du pivot (early start)** : si une piste est en cours de lecture +(un `current_index` est défini dans le snapshot courant), appeler `on_ready` immédiatement +avant toute opération — la piste courante sera préservée. + +**2. Si pas de pivot** (queue vide ou aucune piste en cours) : appeler `on_ready` après +avoir inséré le premier item. + +**3. Vérification cancel** : après chaque item inséré/supprimé (en pratique `InternalQueue` +est rapide mais le principe doit être cohérent) : + +```rust +fn sync_queue( + &mut self, + items: Vec, + cancel_token: &Arc, + mut on_ready: Option>, +) -> Result<(), ControlPointError> { + use std::sync::atomic::Ordering::SeqCst; + + // Early start si pivot présent + let has_current = self.current_index()?.is_some(); + if has_current { + if let Some(f) = on_ready.take() { f(); } + } + + // ... logique existante de sync_queue() ... + // Dans la boucle d'insertions, après le 1er insert : + if on_ready.is_some() { + if let Some(f) = on_ready.take() { f(); } + } + // Après chaque opération : + if cancel_token.load(SeqCst) { + return Err(ControlPointError::SyncCancelled); + } + + Ok(()) +} +``` + +--- + +## Étape 4 — Adapter `OpenHomeQueue::sync_queue()` (`openhome.rs`) + +**Fichier**: `pmocontrol/src/queue/openhome.rs` + +### 4.1 Signature (ligne ~1277) + +```rust +fn sync_queue( + &mut self, + items: Vec, + cancel_token: &Arc, + mut on_ready: Option>, +) -> Result<(), ControlPointError> +``` + +### 4.2 Early start — logique pivot + +Au début de `sync_queue()`, **avant** toute opération SOAP, détecter si un pivot est présent : + +```rust +// Après la récupération du snapshot (ligne ~1323), avant les branches if/else : +let has_pivot = playing_info.is_some(); +if has_pivot { + // Le pivot sera préservé — on peut démarrer la lecture immédiatement + if let Some(f) = on_ready.take() { f(); } +} +``` + +Si pas de pivot (nouvelle playlist via `delete_all` + inserts depuis 0) : appeler `on_ready` +après le **1er insert réussi** dans `replace_queue()` et dans `replace_queue_standard_lcs()`. + +### 4.3 Points de vérification cancel + +Ajouter `if cancel_token.load(SeqCst) { return Err(SyncCancelled); }` aux endroits suivants : + +- Dans `delete_marked_items()` (ligne ~458) : après chaque `delete_id_if_exists()` +- Dans `rebuild_playlist_section()` (ligne ~494) : après chaque `insert()` +- Dans `replace_queue_preserve_current()` (ligne ~419) : après chaque `delete_id_if_exists()` et `insert()` +- Dans `replace_queue_standard_lcs()` (ligne ~693-745) : après chaque delete et insert +- Dans `replace_queue()` (ligne ~1118) : après chaque insert dans la boucle +- Dans le fast path `AppendOnly` (ligne ~1318) : après chaque insert +- Dans le fast path `DeleteFromEnd` (ligne ~1336) : après chaque delete + +### 4.4 Propagation du cancel_token aux helpers + +Les méthodes helper privées qui font des boucles doivent recevoir le token : + +```rust +fn delete_marked_items( + &mut self, + old_ids: &[u32], + keep_flags: &[bool], + position_label: &str, + cancel_token: &Arc, +) -> Result<(), ControlPointError> + +fn rebuild_playlist_section( + &mut self, + // ... params existants ... + cancel_token: &Arc, + on_ready: &mut Option>, +) -> Result +``` + +--- + +## Étape 5 — Adapter `MusicQueue` (dispatch enum) (`music_queue.rs`) + +**Fichier**: `pmocontrol/src/queue/music_queue.rs` + +### 5.1 Adapter le dispatch `sync_queue()` (ligne ~98) + +```rust +fn sync_queue( + &mut self, + items: Vec, + cancel_token: &Arc, + on_ready: Option>, +) -> Result<(), ControlPointError> { + match self { + MusicQueue::Internal(q) => q.sync_queue(items, cancel_token, on_ready), + MusicQueue::OpenHome(q) => q.sync_queue(items, cancel_token, on_ready), + } +} +``` + +### 5.2 Ajouter l'état async dans `MusicQueue` + +`MusicQueue` passe de simple enum de dispatch à une struct qui **encapsule** l'enum backend +et l'état de synchronisation async : + +```rust +// AVANT +pub enum MusicQueue { + Internal(InternalQueue), + OpenHome(OpenHomeQueue), +} + +// APRÈS +pub struct MusicQueue { + backend: MusicQueueBackend, + // État async de synchronisation + sync_in_progress: Arc, + sync_pending: Arc, + sync_cancel_token: Arc, +} + +// L'enum devient privée +enum MusicQueueBackend { + Internal(InternalQueue), + OpenHome(OpenHomeQueue), +} +``` + +**Note**: le changement de `enum` en `struct` implique de mettre à jour toutes les +utilisations de `MusicQueue::Internal(...)` et `MusicQueue::OpenHome(...)` dans le reste +du code (essentiellement `music_queue.rs` lui-même et `mod.rs`). Les callers externes +utilisent `MusicQueue` via `QueueBackend` et `QueueFromRendererInfo` — ils ne sont pas +impactés si l'API publique est préservée. + +### 5.3 Enum résultat et méthode `schedule_sync()` + +```rust +pub enum SyncScheduleOutcome { + /// Thread spawné, sync en cours. + Scheduled, + /// Sync déjà en cours — annulée et nouvelle sync programmée en pending. + AlreadyRunning, +} +``` + +```rust +impl MusicQueue { + /// Lance une synchronisation asynchrone de la queue. + /// + /// - Si aucune sync n'est en cours : spawne un thread, retourne `Scheduled`. + /// - Si une sync est en cours : l'annule, note une sync pending, retourne `AlreadyRunning`. + /// Le thread en cours finira l'opération courante, détectera le cancel, puis + /// relancera la sync avec les nouveaux items via `pending_items_fn`. + /// + /// `pending_items_fn` : closure appelée dans le worker pour re-fetcher les items + /// en cas de pending. Elle doit être Send + 'static car elle s'exécute dans un thread. + /// + /// `on_ready` : appelé dès que la lecture peut démarrer (pivot préservé ou 1er insert). + /// + /// `renderer_id` : utilisé uniquement pour nommer le thread de travail. + pub fn schedule_sync( + &self, // &self car l'état async est derrière Arc + renderer_id: &str, + items: Vec, + pending_items_fn: Box Result, ControlPointError> + Send + 'static>, + on_ready: Option>, + ) -> SyncScheduleOutcome +``` + +**Problème de `&mut self` vs `&self`** : `sync_queue()` dans le trait prend `&mut self` +car les backends mutent leur état. Mais `schedule_sync()` veut spawner un thread qui détient +le backend. Solution : le backend est déjà derrière `Arc>` dans +`MusicRenderer` (ligne 124 de musicrenderer.rs — c'est le `queue` field). Le thread worker +clone cet `Arc` et acquiert le lock pour appeler `sync_queue()`. + +**Signature révisée** : + +```rust +/// Doit être appelé avec un Arc> pour permettre le spawn du thread worker. +pub fn schedule_sync( + queue_arc: &Arc>, + renderer_id: &str, + items: Vec, + pending_items_fn: Box Result, ControlPointError> + Send + 'static>, + on_ready: Option>, +) -> SyncScheduleOutcome { + use std::sync::atomic::Ordering::SeqCst; + + let (sync_in_progress, sync_pending, sync_cancel_token) = { + let q = queue_arc.lock().unwrap(); + ( + Arc::clone(&q.sync_in_progress), + Arc::clone(&q.sync_pending), + Arc::clone(&q.sync_cancel_token), + ) + }; + + if sync_in_progress.swap(true, SeqCst) { + // Sync en cours : annuler et noter pending + sync_cancel_token.store(true, SeqCst); + sync_pending.store(true, SeqCst); + return SyncScheduleOutcome::AlreadyRunning; + } + + // Pas de sync en cours : initialiser et spawner + sync_cancel_token.store(false, SeqCst); + sync_pending.store(false, SeqCst); + + let queue_arc = Arc::clone(queue_arc); + let thread_name = format!("queue-sync-{}", renderer_id); + + thread::Builder::new() + .name(thread_name) + .spawn(move || { + // Guard: libère in_progress à la sortie même en cas de panic + struct Guard(Arc); + impl Drop for Guard { + fn drop(&mut self) { self.0.store(false, SeqCst); } + } + let _guard = Guard(Arc::clone(&sync_in_progress)); + + let mut current_items = items; + let mut current_on_ready = Some(on_ready); + + loop { + sync_pending.store(false, SeqCst); + sync_cancel_token.store(false, SeqCst); + + let result = { + let mut q = queue_arc.lock().unwrap(); + q.backend.sync_queue( + current_items, + &sync_cancel_token, + current_on_ready.take().flatten(), + ) + }; + + match result { + Err(ControlPointError::SyncCancelled) => { + // Annulé normalement — vérifier si pending + } + Err(e) => { + warn!("queue-sync error: {}", e); + } + Ok(()) => {} + } + + if !sync_pending.load(SeqCst) { + break; // Pas de nouvelle sync en attente → terminer + } + + // Nouvelle sync demandée pendant l'exécution → re-fetcher et relancer + match pending_items_fn() { + Ok(new_items) => { + current_items = new_items; + current_on_ready = Some(None); // pas de on_ready pour les re-syncs + } + Err(e) => { + warn!("queue-sync pending re-fetch error: {}", e); + break; + } + } + } + // _guard libère sync_in_progress = false + }) + .expect("Failed to spawn queue-sync thread"); + + SyncScheduleOutcome::Scheduled +} +``` + +--- + +## Étape 6 — Nouveaux événements SSE + +### 6.1 `model.rs` + +Localiser l'enum `RendererEvent` et ajouter : + +```rust +/// Émis dès que la queue peut être lue (pivot préservé ou 1er track inséré). +QueueReadyToPlay { + id: DeviceId, +}, +/// Émis quand une sync est annulée car une nouvelle a été demandée. +QueueSyncCancelled { + id: DeviceId, +}, +``` + +### 6.2 `sse.rs` + +Localiser le match de sérialisation des `RendererEvent` et ajouter les deux variantes. +Suivre le pattern existant de `QueueRefreshing` (ligne ~152) : + +```rust +RendererEvent::QueueReadyToPlay { id } => { + // sérialiser avec type = "queue_ready_to_play" +} +RendererEvent::QueueSyncCancelled { id } => { + // sérialiser avec type = "queue_sync_cancelled" +} +``` + +--- + +## Étape 7 — Modifier `control_point.rs` + +### 7.1 Extraire la logique de browse + +Renommer `refresh_attached_queue_for()` en deux fonctions : + +**`fetch_queue_items_for()`** (nouvelle, interne) : fait le browse + conversion, retourne +`Vec`. C'est la `pending_items_fn` passée à `schedule_sync()`. + +**`schedule_queue_refresh_for()`** (remplace l'ancienne) : appelle `fetch_queue_items_for()`, +puis `MusicQueue::schedule_sync()`. + +```rust +fn fetch_queue_items_for( + registry: &Arc>, + renderer_id: &DeviceId, +) -> Result, ControlPointError> { + // Browse media server + conversion PlaybackItem + // (logique actuellement dans refresh_attached_queue_for() lignes ~1593-1670) +} + +fn schedule_queue_refresh_for( + registry: &Arc>, + renderer_id: &DeviceId, + event_bus: &RendererEventBus, + auto_play_cb: Option Result<(), ControlPointError> + Send + 'static>>, +) -> SyncScheduleOutcome { + let items = match fetch_queue_items_for(registry, renderer_id) { + Ok(items) => items, + Err(e) => { warn!(...); return SyncScheduleOutcome::Scheduled; /* ou erreur */ } + }; + + // on_ready : déclenche auto_play si demandé + émet QueueReadyToPlay SSE + let rid = renderer_id.clone(); + let bus = event_bus.clone(); + let on_ready: Option> = Some(Box::new(move || { + bus.broadcast(RendererEvent::QueueReadyToPlay { id: rid.clone() }); + if let Some(cb) = auto_play_cb { + if let Err(e) = cb(&rid) { + warn!("auto-play callback failed: {}", e); + } + } + })); + + // pending_items_fn : re-fetcher depuis le media server si pending + let registry2 = Arc::clone(registry); + let rid2 = renderer_id.clone(); + let pending_fn = Box::new(move || fetch_queue_items_for(®istry2, &rid2)); + + // Récupérer l'Arc> du renderer depuis le registry + let queue_arc = { + let reg = registry.read().unwrap(); + reg.get_renderer_queue_arc(renderer_id)? // méthode à ajouter dans DeviceRegistry + }; + + // Émettre QueueRefreshing avant de lancer + event_bus.broadcast(RendererEvent::QueueRefreshing { id: renderer_id.clone() }); + + let outcome = MusicQueue::schedule_sync(&queue_arc, &renderer_id.0, items, pending_fn, on_ready); + + if matches!(outcome, SyncScheduleOutcome::AlreadyRunning) { + event_bus.broadcast(RendererEvent::QueueSyncCancelled { id: renderer_id.clone() }); + } + + outcome +} +``` + +### 7.2 Call site 1 — thread "cp-media-server-event-worker" (l.270) + +```rust +// AVANT +let _ = refresh_attached_queue_for(®istry, &renderer_id, &event_bus, None); +// APRÈS +schedule_queue_refresh_for(®istry, &renderer_id, &event_bus, None); +``` + +### 7.3 Call site 2 — thread "cp-playlist-periodic-refresh" (l.340) + +```rust +// AVANT +let _ = refresh_attached_queue_for(®istry_for_periodic, &renderer_id, &event_bus_for_periodic, None); +// APRÈS +schedule_queue_refresh_for(®istry_for_periodic, &renderer_id, &event_bus_for_periodic, None); +``` + +### 7.4 Call site 3 — `attach_queue_to_playlist_async()` (l.1264) + +```rust +pub async fn attach_queue_to_playlist_async( + &self, renderer_id: &DeviceId, server_id: &DeviceId, container_id: &str, auto_play: bool, +) -> Result<(), ControlPointError> { + // 1. Enregistrer la liaison (synchrone, ~1ms) + self.registry.write().unwrap() + .set_playlist_binding(renderer_id, server_id, container_id); + + // 2. Construire le callback auto-play si demandé + let cp = self.clone(); + let rid = renderer_id.clone(); + let auto_play_cb: Option Result<(), ControlPointError> + Send + 'static>> = + if auto_play { + Some(Box::new(move |id| cp.play_current_from_queue(id))) + } else { + None + }; + + // 3. Lancer le refresh async — retourne immédiatement + schedule_queue_refresh_for(&self.registry, renderer_id, &self.event_bus, auto_play_cb); + + Ok(()) // La webapp sera notifiée via SSE (QueueRefreshing → QueueReadyToPlay → QueueUpdated) +} +``` + +### 7.5 Fin de sync — émettre `QueueUpdated` + +L'événement `QueueUpdated` (avec `queue_length`) est actuellement émis à la ligne ~1718 +dans `refresh_attached_queue_for()`. Il doit être émis à la fin du worker thread dans +`MusicQueue::schedule_sync()`, après le `sync_queue()` réussi. + +Passer un `on_complete` callback à `schedule_sync()` (en plus de `on_ready`) : + +```rust +pub fn schedule_sync( + queue_arc: &Arc>, + renderer_id: &str, + items: Vec, + pending_items_fn: Box Result, ControlPointError> + Send + 'static>, + on_ready: Option>, + on_complete: Box, // NOUVEAU — reçoit queue_length +) -> SyncScheduleOutcome +``` + +Dans le worker, après `Ok(())` du `sync_queue()` : + +```rust +Ok(()) => { + let queue_len = queue_arc.lock().unwrap().len().unwrap_or(0); + on_complete(queue_len); +} +``` + +Dans `schedule_queue_refresh_for()` : + +```rust +let bus3 = event_bus.clone(); +let rid3 = renderer_id.clone(); +let on_complete = Box::new(move |queue_len: usize| { + bus3.broadcast(RendererEvent::QueueUpdated { + id: rid3.clone(), + queue_length: queue_len, + }); +}); +``` + +--- + +## Étape 8 — Accès à `Arc>` depuis le registry + +`schedule_queue_refresh_for()` a besoin d'accéder à l'`Arc>` du renderer. +Localiser dans `DeviceRegistry` comment les renderers et leurs queues sont stockés, et ajouter +une méthode : + +```rust +pub fn get_renderer_queue_arc( + &self, + renderer_id: &DeviceId, +) -> Result>, ControlPointError> +``` + +(Ou équivalent selon la structure réelle du registry.) + +--- + +## Ordre d'implémentation + +``` +Étape 1 (errors.rs) ← 5 min +Étape 2 (backend.rs) ← 10 min, casse la compilation → à faire avant les autres +Étape 3 (interne.rs) ← 30 min +Étape 4 (openhome.rs) ← 1-2h (nombreux points de vérification cancel) +Étape 5 (music_queue.rs) ← 2-3h (changement struct + schedule_sync) +Étape 6 (model.rs + sse.rs) ← 30 min +Étape 7 (control_point.rs) ← 1h +Étape 8 (registry) ← 30 min selon la structure +``` + +Après l'étape 2, `cargo build` cassera jusqu'à l'étape 5 incluse — c'est attendu. +Faire les étapes 3, 4, 5 dans la même session sans interrompre. + +## Tests + +```bash +cargo build -p pmocontrol + +# Vérifier les scénarios : +# 1. Attach playlist → QueueRefreshing SSE immédiat, QueueReadyToPlay après 1er insert, +# QueueUpdated après fin complète +# 2. Attach 2e playlist pendant sync en cours → QueueSyncCancelled + nouveau refresh repart +# 3. Piste en cours de lecture pendant sync → QueueReadyToPlay immédiat (pivot préservé) +# 4. Refresh périodique (60s) ne bloque plus son thread + +RUST_LOG=debug cargo run 2>&1 | grep -E "queue.sync|SyncCancelled|Scheduled|AlreadyRunning|on_ready|on_complete" +``` + +--- + +*Date: 2026-04-09* diff --git a/Cargo.lock b/Cargo.lock index 664d3596..f874fddc 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4,7 +4,7 @@ version = 4 [[package]] name = "PMOMusic" -version = "0.3.40" +version = "0.3.41" dependencies = [ "axum 0.8.7", "console-subscriber", diff --git a/PMOMusic/Cargo.toml b/PMOMusic/Cargo.toml index f80e2624..bc890477 100644 --- a/PMOMusic/Cargo.toml +++ b/PMOMusic/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "PMOMusic" -version = "0.3.40" +version = "0.3.41" edition = "2024" [dependencies] diff --git a/pmocontrol/src/control_point.rs b/pmocontrol/src/control_point.rs index 5b8b3068..01220e08 100644 --- a/pmocontrol/src/control_point.rs +++ b/pmocontrol/src/control_point.rs @@ -26,7 +26,7 @@ use crate::openapi::{ CurrentTrackMetadata, FullRendererSnapshot, QueueItem, QueueSnapshotView, RendererBindingView, RendererStateView, }; -use crate::queue::{EnqueueMode, PlaybackItem, QueueSnapshot}; +use crate::queue::{EnqueueMode, PlaybackItem, QueueSnapshot, SyncScheduleOutcome}; use crate::registry::DeviceRegistry; /// Control point minimal : @@ -267,19 +267,12 @@ impl ControlPoint { "Triggering queue refresh for bound playlist" ); - if let Err(err) = refresh_attached_queue_for( + let _ = schedule_queue_refresh_for( ®istry_for_media_worker, &renderer_id, &event_bus_for_media_worker, None, - ) { - warn!( - renderer = renderer_id.0.as_str(), - server = server_id.0.as_str(), - error = %err, - "Failed to refresh queue from playlist container" - ); - } + ); } } MediaServerEvent::Online { server_id, info } => { @@ -337,18 +330,12 @@ impl ControlPoint { "Periodic refresh triggered for bound playlist" ); - if let Err(err) = refresh_attached_queue_for( + let _ = schedule_queue_refresh_for( ®istry_for_periodic, &renderer_id, &event_bus_for_periodic, None, - ) { - warn!( - renderer = renderer_id.0.as_str(), - error = %err, - "Periodic refresh failed for bound playlist" - ); - } + ); } } })?; @@ -1182,19 +1169,25 @@ impl ControlPoint { binding.auto_play_on_refresh = auto_play; renderer.set_playlist_binding(Some(binding)); - let mut auto_start_cb = |rid: &DeviceId| self.play_current_from_queue(rid); - let callback: Option<&mut dyn FnMut(&DeviceId) -> Result<(), ControlPointError>> = + let callback: Option Result<(), ControlPointError> + Send + 'static>> = if auto_play { - Some(&mut auto_start_cb) + let reg = Arc::clone(&self.registry); + Some(Box::new(move |rid: &DeviceId| { + let renderer = reg.read().unwrap().get_renderer(rid).ok_or_else(|| + ControlPointError::ControlPoint(format!("Renderer {} not found", rid.0)) + )?; + renderer.play_current_from_queue() + })) } else { None }; - return refresh_attached_queue_for( + let _ = schedule_queue_refresh_for( &self.registry, renderer_id, &self.event_bus, callback, ); + return Ok(()); } // CRITICAL: When attaching a new playlist to a renderer, we must UNCONDITIONALLY @@ -1247,21 +1240,25 @@ impl ControlPoint { // Note: BindingChanged event is emitted automatically by MusicRenderer::set_playlist_binding() // For initial attach with auto_play, force playback start (don't check if idle) - let mut auto_start_cb = |rid: &DeviceId| { - debug!( - renderer = rid.0.as_str(), - "Attach callback: forcing playback start (not checking if idle)" - ); - self.play_current_from_queue(rid) - }; - let callback: Option<&mut dyn FnMut(&DeviceId) -> Result<(), ControlPointError>> = + let callback: Option Result<(), ControlPointError> + Send + 'static>> = if auto_play { - Some(&mut auto_start_cb) + let reg = Arc::clone(&self.registry); + Some(Box::new(move |rid: &DeviceId| { + debug!( + renderer = rid.0.as_str(), + "Attach callback: forcing playback start (not checking if idle)" + ); + let renderer = reg.read().unwrap().get_renderer(rid).ok_or_else(|| + ControlPointError::ControlPoint(format!("Renderer {} not found", rid.0)) + )?; + renderer.play_current_from_queue() + })) } else { None }; - refresh_attached_queue_for(&self.registry, renderer_id, &self.event_bus, callback) + let _ = schedule_queue_refresh_for(&self.registry, renderer_id, &self.event_bus, callback); + Ok(()) } /// Detach a renderer's queue from its associated playlist container. @@ -1760,6 +1757,199 @@ fn refresh_attached_queue_for( Ok(()) } +fn fetch_queue_items_for( + music_server: &Arc, + container_id: &str, +) -> Result, ControlPointError> { + const MAX_BROWSE_ATTEMPTS: usize = 3; + const BROWSE_RETRY_DELAY_MS: u64 = 200; + const BROWSE_PAGE_SIZE: u32 = 64; + + let entries = { + let mut all_entries = Vec::new(); + let mut offset = 0u32; + loop { + let mut attempt = 1; + let page = loop { + match music_server.browse_children(container_id, offset, BROWSE_PAGE_SIZE) { + Ok(e) => break e, + Err(err) => { + if attempt >= MAX_BROWSE_ATTEMPTS { + return Err(err); + } + thread::sleep(Duration::from_millis( + BROWSE_RETRY_DELAY_MS * attempt as u64, + )); + attempt += 1; + } + } + }; + let fetched = page.len() as u32; + all_entries.extend(page); + if fetched < BROWSE_PAGE_SIZE { + break; + } + offset += fetched; + } + all_entries + }; + + let new_items: Vec = entries + .iter() + .filter_map(|entry| playback_item_from_entry(music_server.clone(), entry)) + .collect(); + + Ok(new_items) +} + +fn schedule_queue_refresh_for( + registry: &Arc>, + renderer_id: &DeviceId, + event_bus: &RendererEventBus, + after_refresh: Option Result<(), ControlPointError> + Send + 'static>>, +) -> SyncScheduleOutcome { + // Step 1: Get renderer from registry + let renderer = { + let reg = registry.read().unwrap(); + reg.get_renderer(renderer_id) + }; + + let renderer = match renderer { + Some(r) => r, + None => { + debug!( + renderer = renderer_id.0.as_str(), + "schedule_queue_refresh_for: renderer not found" + ); + return SyncScheduleOutcome::Scheduled; + } + }; + + // Check if there's a binding + let (server_id, container_id) = { + let binding = match renderer.get_playlist_binding() { + Some(b) => b, + None => { + debug!( + renderer = renderer_id.0.as_str(), + "schedule_queue_refresh_for: no binding present" + ); + return SyncScheduleOutcome::Scheduled; + } + }; + (binding.server_id.clone(), binding.container_id.clone()) + }; + + // Step 2: Get server from registry + let music_server = { + let reg = registry.read().unwrap(); + reg.get_server(&server_id) + }; + + let music_server = match music_server { + Some(s) => s, + None => { + warn!( + renderer = renderer_id.0.as_str(), + server = server_id.0.as_str(), + "schedule_queue_refresh_for: server not found in registry" + ); + return SyncScheduleOutcome::Scheduled; + } + }; + + if !music_server.is_online() { + debug!( + renderer = renderer_id.0.as_str(), + server = server_id.0.as_str(), + "schedule_queue_refresh_for: server offline, skipping refresh" + ); + return SyncScheduleOutcome::Scheduled; + } + + // Emit QueueRefreshing + event_bus.broadcast(RendererEvent::QueueRefreshing { + id: renderer_id.clone(), + }); + + // Get the queue Arc + let queue_arc = { + let reg = registry.read().unwrap(); + match reg.get_renderer_queue_arc(renderer_id) { + Some(q) => q, + None => { + warn!("schedule_queue_refresh_for: could not get queue arc"); + return SyncScheduleOutcome::Scheduled; + } + } + }; + + // Build the pending_items_fn for re-fetching + let registry_clone = Arc::clone(registry); + let server_id_clone = server_id.clone(); + let container_id_clone = container_id.clone(); + let pending_fn = Box::new(move || { + let music_server = { + let reg = registry_clone.read().unwrap(); + reg.get_server(&server_id_clone) + }; + match music_server { + Some(s) => fetch_queue_items_for(&s, &container_id_clone), + None => Err(ControlPointError::MediaServerError("Server not found".to_string())), + } + }); + + // Build on_ready callback + let rid = renderer_id.clone(); + let bus = event_bus.clone(); + let after_cb = after_refresh; + let on_ready: Option> = Some(Box::new(move || { + bus.broadcast(RendererEvent::QueueReadyToPlay { id: rid.clone() }); + if let Some(cb) = after_cb { + if let Err(e) = cb(&rid) { + warn!("auto-play callback failed: {}", e); + } + } + })); + + // Build on_complete callback + let rid2 = renderer_id.clone(); + let bus2 = event_bus.clone(); + let on_complete = Box::new(move |queue_len: usize| { + bus2.broadcast(RendererEvent::QueueUpdated { + id: rid2.clone(), + queue_length: queue_len, + }); + }); + + // First fetch + let items = match fetch_queue_items_for(&music_server, &container_id) { + Ok(items) => items, + Err(e) => { + warn!("schedule_queue_refresh_for: failed to fetch items: {}", e); + return SyncScheduleOutcome::Scheduled; + } + }; + + // Call schedule_sync + let outcome = crate::queue::MusicQueue::schedule_sync( + &queue_arc, + &renderer_id.0, + items, + pending_fn, + on_ready, + on_complete, + ); + + if matches!(outcome, SyncScheduleOutcome::AlreadyRunning) { + event_bus.broadcast(RendererEvent::QueueSyncCancelled { + id: renderer_id.clone(), + }); + } + + outcome +} + fn didl_item_from_playback_item(item: &PlaybackItem) -> DidlItem { let metadata = item.metadata.as_ref(); let title = metadata diff --git a/pmocontrol/src/errors.rs b/pmocontrol/src/errors.rs index 29cc36a4..2d0ef137 100644 --- a/pmocontrol/src/errors.rs +++ b/pmocontrol/src/errors.rs @@ -49,6 +49,8 @@ pub enum ControlPointError { SnapshotError(String), #[error("Error on ControlPoint: {0}")] ControlPoint(String), + #[error("Queue sync cancelled (superseded by a newer request)")] + SyncCancelled, } impl ControlPointError { diff --git a/pmocontrol/src/model.rs b/pmocontrol/src/model.rs index 9686d3a9..1dcf1d38 100644 --- a/pmocontrol/src/model.rs +++ b/pmocontrol/src/model.rs @@ -437,6 +437,12 @@ pub enum RendererEvent { QueueRefreshing { id: DeviceId, }, + QueueReadyToPlay { + id: DeviceId, + }, + QueueSyncCancelled { + id: DeviceId, + }, BindingChanged { id: DeviceId, binding: Option, diff --git a/pmocontrol/src/music_renderer/arylic_tcp.rs b/pmocontrol/src/music_renderer/arylic_tcp.rs index 38a9ac45..82752bfa 100644 --- a/pmocontrol/src/music_renderer/arylic_tcp.rs +++ b/pmocontrol/src/music_renderer/arylic_tcp.rs @@ -1,26 +1,26 @@ -use std::sync::{Arc, Mutex}; +use std::sync::{atomic::AtomicBool, Arc, Mutex}; use std::time::Duration; use serde::Deserialize; use tracing::debug; -use crate::DeviceIdentity; use crate::arylic_client::{ - ARYLIC_TCP_PORT, DEFAULT_TIMEOUT_SECS, send_command_no_response, send_command_optional, - send_command_required, + send_command_no_response, send_command_optional, send_command_required, ARYLIC_TCP_PORT, + DEFAULT_TIMEOUT_SECS, }; use crate::errors::ControlPointError; use crate::linkplay_client::extract_linkplay_host; use crate::model::{PlaybackState, RendererInfo}; -use crate::music_renderer::RendererFromMediaRendererInfo; use crate::music_renderer::capabilities::{ PlaybackPosition, PlaybackPositionInfo, PlaybackStatus, QueueTransportControl, RendererBackend, TransportControl, VolumeControl, }; use crate::music_renderer::musicrenderer::MusicRendererBackend; use crate::music_renderer::time_utils::{format_hhmmss, ms_to_seconds, parse_hhmmss_strict}; +use crate::music_renderer::RendererFromMediaRendererInfo; use crate::queue::MusicQueue; use crate::queue::{EnqueueMode, PlaybackItem, QueueBackend, QueueSnapshot}; +use crate::DeviceIdentity; /// Raw response from Arylic MCU+PINFGET command #[derive(Debug, Deserialize)] @@ -409,8 +409,16 @@ impl QueueBackend for ArylicTcpRenderer { .replace_queue(items, current_index) } - fn sync_queue(&mut self, items: Vec) -> Result<(), ControlPointError> { - self.queue.lock().unwrap().sync_queue(items) + fn sync_queue( + &mut self, + items: Vec, + _cancel_token: &Arc, + on_ready: Option>, + ) -> Result<(), ControlPointError> { + self.queue + .lock() + .unwrap() + .sync_queue(items, &Arc::new(AtomicBool::new(false)), on_ready) } fn get_item(&self, index: usize) -> Result, ControlPointError> { diff --git a/pmocontrol/src/music_renderer/chromecast_renderer.rs b/pmocontrol/src/music_renderer/chromecast_renderer.rs index e00cddf2..9d6fa61f 100644 --- a/pmocontrol/src/music_renderer/chromecast_renderer.rs +++ b/pmocontrol/src/music_renderer/chromecast_renderer.rs @@ -11,6 +11,7 @@ //! are wrapped in sync calls using smol::block_on for compatibility with //! the existing sync trait interfaces. +use std::sync::atomic::AtomicBool; use std::sync::{Arc, Mutex, Once}; use std::thread::JoinHandle; @@ -906,8 +907,16 @@ impl QueueBackend for ChromecastRenderer { .replace_queue(items, current_index) } - fn sync_queue(&mut self, items: Vec) -> Result<(), ControlPointError> { - self.queue.lock().unwrap().sync_queue(items) + fn sync_queue( + &mut self, + items: Vec, + _cancel_token: &Arc, + on_ready: Option>, + ) -> Result<(), ControlPointError> { + self.queue + .lock() + .unwrap() + .sync_queue(items, &Arc::new(AtomicBool::new(false)), on_ready) } fn get_item(&self, index: usize) -> Result, ControlPointError> { diff --git a/pmocontrol/src/music_renderer/linkplay_renderer.rs b/pmocontrol/src/music_renderer/linkplay_renderer.rs index dc0e3b49..689cac7e 100644 --- a/pmocontrol/src/music_renderer/linkplay_renderer.rs +++ b/pmocontrol/src/music_renderer/linkplay_renderer.rs @@ -1,24 +1,24 @@ use std::fmt; -use std::sync::{Arc, Mutex}; +use std::sync::{atomic::AtomicBool, Arc, Mutex}; use std::time::Duration; use ureq::Agent; -use crate::DeviceIdentity; use crate::errors::ControlPointError; use crate::linkplay_client::{ - LinkPlayStatus, build_agent, extract_linkplay_host, fetch_status_for_host, percent_encode, + build_agent, extract_linkplay_host, fetch_status_for_host, percent_encode, LinkPlayStatus, }; use crate::model::{PlaybackState, RendererInfo}; -use crate::music_renderer::RendererFromMediaRendererInfo; use crate::music_renderer::capabilities::{ PlaybackPosition, PlaybackPositionInfo, PlaybackStatus, QueueTransportControl, RendererBackend, TransportControl, VolumeControl, }; use crate::music_renderer::musicrenderer::MusicRendererBackend; use crate::music_renderer::time_utils::parse_hhmmss_strict; +use crate::music_renderer::RendererFromMediaRendererInfo; use crate::queue::MusicQueue; use crate::queue::{EnqueueMode, PlaybackItem, QueueBackend, QueueSnapshot}; +use crate::DeviceIdentity; const DEFAULT_HTTP_TIMEOUT_SECS: u64 = 3; @@ -287,8 +287,16 @@ impl QueueBackend for LinkPlayRenderer { .replace_queue(items, current_index) } - fn sync_queue(&mut self, items: Vec) -> Result<(), ControlPointError> { - self.queue.lock().unwrap().sync_queue(items) + fn sync_queue( + &mut self, + items: Vec, + _cancel_token: &Arc, + on_ready: Option>, + ) -> Result<(), ControlPointError> { + self.queue + .lock() + .unwrap() + .sync_queue(items, &Arc::new(AtomicBool::new(false)), on_ready) } fn get_item(&self, index: usize) -> Result, ControlPointError> { diff --git a/pmocontrol/src/music_renderer/musicrenderer.rs b/pmocontrol/src/music_renderer/musicrenderer.rs index c05f4370..f3a8badd 100644 --- a/pmocontrol/src/music_renderer/musicrenderer.rs +++ b/pmocontrol/src/music_renderer/musicrenderer.rs @@ -976,6 +976,12 @@ impl MusicRenderer { Ok(snapshot) } + /// Get a clone of the queue Arc (for async sync operations). + pub fn queue(&self) -> Arc> { + let backend = self.lock_backend_for("queue"); + crate::music_renderer::capabilities::RendererBackend::queue(&*backend).clone() + } + /// Get the current queue item without advancing. /// Returns the item and count of remaining items after current. pub fn peek_current(&self) -> Result, ControlPointError> { @@ -1372,7 +1378,16 @@ 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 + self.sync_queue_with_callback(items, None, None) + } + + /// Synchronize the queue with new items with optional cancel token and on_ready callback. + pub fn sync_queue_with_callback( + &self, + items: Vec, + cancel_token: Option<&Arc>, + on_ready: Option>, + ) -> Result<(), ControlPointError> { { let mut watched = self .watched_state @@ -1382,7 +1397,9 @@ impl MusicRenderer { } let mut backend = self.lock_backend_for("sync_queue"); - backend.sync_queue(items)?; + let dummy_token = Arc::new(AtomicBool::new(false)); + let token = cancel_token.unwrap_or(&dummy_token); + backend.sync_queue(items, token, on_ready)?; drop(backend); self.emit_queue_updated(); Ok(()) @@ -2283,14 +2300,21 @@ impl QueueBackend for MusicRendererBackend { } } - fn sync_queue(&mut self, items: Vec) -> Result<(), ControlPointError> { + fn sync_queue( + &mut self, + items: Vec, + cancel_token: &Arc, + on_ready: Option>, + ) -> Result<(), ControlPointError> { match self { - MusicRendererBackend::Upnp(r) => r.sync_queue(items), - MusicRendererBackend::OpenHome(r) => r.sync_queue(items), - MusicRendererBackend::LinkPlay(r) => r.sync_queue(items), - MusicRendererBackend::ArylicTcp(r) => r.sync_queue(items), - MusicRendererBackend::Chromecast(cc) => cc.sync_queue(items), - MusicRendererBackend::HybridUpnpArylic { upnp, .. } => upnp.sync_queue(items), + MusicRendererBackend::Upnp(r) => r.sync_queue(items, cancel_token, on_ready), + MusicRendererBackend::OpenHome(r) => r.sync_queue(items, cancel_token, on_ready), + MusicRendererBackend::LinkPlay(r) => r.sync_queue(items, cancel_token, on_ready), + MusicRendererBackend::ArylicTcp(r) => r.sync_queue(items, cancel_token, on_ready), + MusicRendererBackend::Chromecast(cc) => cc.sync_queue(items, cancel_token, on_ready), + MusicRendererBackend::HybridUpnpArylic { upnp, .. } => { + upnp.sync_queue(items, cancel_token, on_ready) + } } } diff --git a/pmocontrol/src/music_renderer/openhome_renderer.rs b/pmocontrol/src/music_renderer/openhome_renderer.rs index b83ea713..c843c4e6 100644 --- a/pmocontrol/src/music_renderer/openhome_renderer.rs +++ b/pmocontrol/src/music_renderer/openhome_renderer.rs @@ -1,4 +1,4 @@ -use std::sync::{Arc, Mutex}; +use std::sync::{atomic::AtomicBool, Arc, Mutex}; use std::time::SystemTime; use crate::music_renderer::capabilities::{ @@ -182,9 +182,8 @@ impl OpenHomeRenderer { /// Retourne les IDs des pistes de la playlist OpenHome. /// Plus rapide que snapshot_openhome_playlist() car ne récupère pas les métadonnées. pub(crate) fn openhome_playlist_ids(&self) -> Result, ControlPointError> { - // Use the queue's cached track_ids() instead of direct id_array() call let queue = self.queue.lock().unwrap(); - if let MusicQueue::OpenHome(oh_queue) = &*queue { + if let Some(oh_queue) = queue.as_openhome() { oh_queue.track_ids() } else { Err(ControlPointError::QueueError( @@ -209,9 +208,8 @@ impl OpenHomeRenderer { let insert_after = match after_id { Some(id) => id, None => { - // Use the queue's cached track_ids() instead of direct id_array() call let queue = self.queue.lock().unwrap(); - if let MusicQueue::OpenHome(oh_queue) = &*queue { + if let Some(oh_queue) = queue.as_openhome() { oh_queue .track_ids()? .last() @@ -413,7 +411,7 @@ impl PlaybackPosition for OpenHomeRenderer { // Get track ID from queue (uses cached data) let queue_guard_for_id = self.queue.lock().unwrap(); - if let MusicQueue::OpenHome(oh_queue) = &*queue_guard_for_id { + if let Some(oh_queue) = queue_guard_for_id.as_openhome() { match oh_queue.current_track() { Ok(id_opt) => track_id = id_opt, Err(err) => debug!( @@ -696,11 +694,16 @@ impl QueueBackend for OpenHomeRenderer { .replace_queue(items, current_index) } - fn sync_queue(&mut self, items: Vec) -> Result<(), ControlPointError> { + fn sync_queue( + &mut self, + items: Vec, + cancel_token: &Arc, + on_ready: Option>, + ) -> Result<(), ControlPointError> { self.queue .lock() .map_err(|_| ControlPointError::QueueError("Queue mutex poisoned".into()))? - .sync_queue(items) + .sync_queue(items, cancel_token, on_ready) } fn get_item(&self, index: usize) -> Result, ControlPointError> { diff --git a/pmocontrol/src/music_renderer/upnp_renderer.rs b/pmocontrol/src/music_renderer/upnp_renderer.rs index 5577727b..11578c5a 100644 --- a/pmocontrol/src/music_renderer/upnp_renderer.rs +++ b/pmocontrol/src/music_renderer/upnp_renderer.rs @@ -1,13 +1,13 @@ -use std::sync::{Arc, Mutex}; +use std::sync::{atomic::AtomicBool, Arc, Mutex}; use crate::errors::ControlPointError; use crate::model::PlaybackState; -use crate::music_renderer::RendererFromMediaRendererInfo; use crate::music_renderer::capabilities::{ PlaybackPosition, PlaybackPositionInfo, PlaybackStatus, QueueTransportControl, RendererBackend, TransportControl, VolumeControl, }; -use crate::music_renderer::musicrenderer::{MusicRendererBackend, build_didl_lite_metadata}; +use crate::music_renderer::musicrenderer::{build_didl_lite_metadata, MusicRendererBackend}; +use crate::music_renderer::RendererFromMediaRendererInfo; use crate::queue::{EnqueueMode, MusicQueue, PlaybackItem, QueueBackend, QueueSnapshot}; use crate::upnp_clients::{ AvTransportClient, ConnectionInfo, ConnectionManagerClient, PositionInfo, ProtocolInfo, @@ -318,8 +318,16 @@ impl QueueBackend for UpnpRenderer { .replace_queue(items, current_index) } - fn sync_queue(&mut self, items: Vec) -> Result<(), ControlPointError> { - self.queue.lock().unwrap().sync_queue(items) + fn sync_queue( + &mut self, + items: Vec, + _cancel_token: &Arc, + on_ready: Option>, + ) -> Result<(), ControlPointError> { + self.queue + .lock() + .unwrap() + .sync_queue(items, &Arc::new(AtomicBool::new(false)), on_ready) } fn get_item(&self, index: usize) -> Result, ControlPointError> { diff --git a/pmocontrol/src/queue/backend.rs b/pmocontrol/src/queue/backend.rs index 4d77c76e..c70d343f 100644 --- a/pmocontrol/src/queue/backend.rs +++ b/pmocontrol/src/queue/backend.rs @@ -28,7 +28,8 @@ //! - This identity is used by the sync helpers to preserve the current //! track across queue rebuilds when the MediaServer content changes. -use crate::{PlaybackItem, QueueSnapshot, errors::ControlPointError}; +use crate::{errors::ControlPointError, PlaybackItem, QueueSnapshot}; +use std::sync::{atomic::AtomicBool, Arc}; /// High-level enqueue mode. /// @@ -107,7 +108,12 @@ pub trait QueueBackend { /// corresponds to the old current index track. /// If the old current index track is absent from the new queue, /// it is kept as the first item and the new items are appended after it. - fn sync_queue(&mut self, items: Vec) -> Result<(), ControlPointError>; + fn sync_queue( + &mut self, + items: Vec, + cancel_token: &Arc, + on_ready: Option>, + ) -> Result<(), ControlPointError>; /// Returns the item at `index`, if it exists. fn get_item(&self, index: usize) -> Result, ControlPointError>; diff --git a/pmocontrol/src/queue/interne.rs b/pmocontrol/src/queue/interne.rs index b17c57e6..c98b8afc 100644 --- a/pmocontrol/src/queue/interne.rs +++ b/pmocontrol/src/queue/interne.rs @@ -20,6 +20,7 @@ use crate::{ queue::{MusicQueue, PlaybackItem, QueueBackend, QueueFromRendererInfo, QueueSnapshot}, DeviceId, DeviceIdentity, RendererInfo, }; +use std::sync::{atomic::AtomicBool, Arc}; /// Internal/local queue implementation. /// @@ -269,25 +270,43 @@ impl QueueBackend for InternalQueue { Ok(()) } - fn sync_queue(&mut self, items: Vec) -> Result<(), ControlPointError> { - use tracing::debug; + fn sync_queue( + &mut self, + items: Vec, + cancel_token: &Arc, + mut on_ready: Option>, + ) -> Result<(), ControlPointError> { + use std::sync::atomic::Ordering::SeqCst; - if items.is_empty() { - return self.replace_queue(Vec::new(), None); + if cancel_token.load(SeqCst) { + return Err(ControlPointError::SyncCancelled); + } + + let has_current = self.current_index().ok().flatten().is_some(); + if has_current { + if let Some(f) = on_ready.take() { + f(); + } + } + + if items.is_empty() { + let _ = self.replace_queue(Vec::new(), None); + return Ok(()); } - // Protéger les durées des streams contre la diminution let updated_items = self.protect_stream_durations(items); - // Récupérer l'item actuel let current = self.current_index.and_then(|idx| { self.items .get(idx) .map(|item| (idx, item.uri.clone(), item.didl_id.clone())) }); + if cancel_token.load(SeqCst) { + return Err(ControlPointError::SyncCancelled); + } + if let Some((_current_idx, current_uri, current_didl_id)) = current { - // Chercher l'item actuel dans la nouvelle liste (par URI d'abord, puis par didl_id) let new_idx = updated_items .iter() .position(|item| item.uri == current_uri) @@ -298,34 +317,25 @@ impl QueueBackend for InternalQueue { }); if let Some(new_idx) = new_idx { - // Item trouvé dans la nouvelle liste - debug!( - renderer = self.renderer_id.0.as_str(), - current_uri = current_uri.as_str(), - new_idx, - "sync_queue: current item found in new playlist" - ); - self.replace_queue(updated_items, Some(new_idx)) + self.replace_queue(updated_items, Some(new_idx))?; } else { - // Item pas trouvé - cela ne devrait pas arriver si la playlist n'a pas changé - // Loguer pour diagnostic - debug!( - renderer = self.renderer_id.0.as_str(), - current_uri = current_uri.as_str(), - current_didl_id = current_didl_id.as_str(), - new_items_count = updated_items.len(), - "sync_queue: current item NOT found in new playlist, preserving as first item" - ); let current_item = self.items[self.current_index.unwrap()].clone(); let mut new_items = Vec::with_capacity(updated_items.len() + 1); new_items.push(current_item); new_items.extend(updated_items); - self.replace_queue(new_items, Some(0)) + self.replace_queue(new_items, Some(0))?; } } else { - // Pas d'item actuel - self.replace_queue(updated_items, None) + self.replace_queue(updated_items, None)?; } + + if !has_current { + if let Some(f) = on_ready.take() { + f(); + } + } + + Ok(()) } fn enqueue_items( @@ -502,6 +512,6 @@ impl QueueFromRendererInfo for InternalQueue { } fn to_backend(self) -> MusicQueue { - MusicQueue::Internal(self) + MusicQueue::from_internal(self) } } diff --git a/pmocontrol/src/queue/mod.rs b/pmocontrol/src/queue/mod.rs index 9054d386..2dfd125c 100644 --- a/pmocontrol/src/queue/mod.rs +++ b/pmocontrol/src/queue/mod.rs @@ -7,20 +7,23 @@ mod snapshot; use std::sync::{Arc, Mutex}; pub use backend::{EnqueueMode, QueueBackend}; -pub use music_queue::MusicQueue; +pub use music_queue::{MusicQueue, SyncScheduleOutcome}; pub use snapshot::{PlaybackItem, QueueSnapshot}; // Internal queue implementations - not part of the public API pub(crate) use interne::InternalQueue; pub(crate) use openhome::OpenHomeQueue; -use crate::{RendererInfo, errors::ControlPointError}; use crate::music_renderer::time_utils::parse_time_flexible; +use crate::{errors::ControlPointError, RendererInfo}; /// 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()) { + match ( + parse_time_flexible(old_dur).ok(), + parse_time_flexible(new_dur).ok(), + ) { (Some(old_secs), Some(new_secs)) => new_secs < old_secs, _ => false, } @@ -28,7 +31,10 @@ pub(super) fn stream_duration_decreased(old_dur: &str, new_dur: &str) -> bool { /// 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()) { + match ( + parse_time_flexible(old_dur).ok(), + parse_time_flexible(new_dur).ok(), + ) { (Some(old_secs), Some(new_secs)) => new_secs > old_secs, _ => false, } diff --git a/pmocontrol/src/queue/music_queue.rs b/pmocontrol/src/queue/music_queue.rs index 44b3d4aa..139c77e5 100644 --- a/pmocontrol/src/queue/music_queue.rs +++ b/pmocontrol/src/queue/music_queue.rs @@ -3,84 +3,297 @@ use crate::queue::{ EnqueueMode, InternalQueue, OpenHomeQueue, QueueBackend, QueueFromRendererInfo, }; use crate::{PlaybackItem, QueueSnapshot, RendererInfo}; +use std::sync::{ + atomic::{AtomicBool, Ordering::SeqCst}, + Arc, Mutex, +}; +use std::thread; + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum SyncScheduleOutcome { + Scheduled, + AlreadyRunning, +} #[derive(Debug)] -pub enum MusicQueue { +enum MusicQueueBackend { Internal(InternalQueue), OpenHome(OpenHomeQueue), } +#[derive(Debug)] +pub struct MusicQueue { + backend: MusicQueueBackend, + sync_in_progress: Arc, + sync_pending: Arc, + sync_cancel_token: Arc, +} + impl MusicQueue { /// Creates a queue appropriate for the given renderer. /// This is the factory method used by QueueFromRendererInfo trait. pub fn from_renderer_info(info: &RendererInfo) -> Result { - if info.capabilities().has_oh_playlist() { - Ok(MusicQueue::OpenHome(OpenHomeQueue::from_renderer_info( - info, - )?)) + let backend = if info.capabilities().has_oh_playlist() { + MusicQueueBackend::OpenHome(OpenHomeQueue::from_renderer_info(info)?) } else { - Ok(MusicQueue::Internal(InternalQueue::from_renderer_info( - info, - )?)) + MusicQueueBackend::Internal(InternalQueue::from_renderer_info(info)?) + }; + + Ok(MusicQueue { + backend, + sync_in_progress: Arc::new(AtomicBool::new(false)), + sync_pending: Arc::new(AtomicBool::new(false)), + sync_cancel_token: Arc::new(AtomicBool::new(false)), + }) + } + + /// Creates a MusicQueue from an InternalQueue backend. + pub fn from_internal(queue: InternalQueue) -> MusicQueue { + MusicQueue { + backend: MusicQueueBackend::Internal(queue), + sync_in_progress: Arc::new(AtomicBool::new(false)), + sync_pending: Arc::new(AtomicBool::new(false)), + sync_cancel_token: Arc::new(AtomicBool::new(false)), } } + + /// Creates a MusicQueue from an OpenHomeQueue backend. + pub fn from_openhome(queue: OpenHomeQueue) -> MusicQueue { + MusicQueue { + backend: MusicQueueBackend::OpenHome(queue), + sync_in_progress: Arc::new(AtomicBool::new(false)), + sync_pending: Arc::new(AtomicBool::new(false)), + sync_cancel_token: Arc::new(AtomicBool::new(false)), + } + } + + /// Returns true if this is an OpenHome backend. + pub fn is_openhome(&self) -> bool { + matches!(self.backend, MusicQueueBackend::OpenHome(_)) + } + + /// Returns a reference to the OpenHome queue if this is an OpenHome backend. + pub fn as_openhome(&self) -> Option<&OpenHomeQueue> { + match &self.backend { + MusicQueueBackend::OpenHome(q) => Some(q), + MusicQueueBackend::Internal(_) => None, + } + } + + /// Returns a mutable reference to the OpenHome queue if this is an OpenHome backend. + pub fn as_openhome_mut(&mut self) -> Option<&mut OpenHomeQueue> { + match &mut self.backend { + MusicQueueBackend::OpenHome(q) => Some(q), + MusicQueueBackend::Internal(_) => None, + } + } + + pub fn len(&self) -> Result { + match &self.backend { + MusicQueueBackend::Internal(q) => q.len(), + MusicQueueBackend::OpenHome(q) => q.len(), + } + } + + pub fn schedule_sync( + queue_arc: &Arc>, + renderer_id: &str, + items: Vec, + pending_items_fn: Box< + dyn Fn() -> Result, ControlPointError> + Send + 'static, + >, + on_ready: Option>, + on_complete: Box, + ) -> SyncScheduleOutcome { + let (sync_in_progress, sync_pending, sync_cancel_token) = { + let q = queue_arc.lock().unwrap(); + ( + Arc::clone(&q.sync_in_progress), + Arc::clone(&q.sync_pending), + Arc::clone(&q.sync_cancel_token), + ) + }; + + if sync_in_progress.swap(true, SeqCst) { + sync_cancel_token.store(true, SeqCst); + sync_pending.store(true, SeqCst); + return SyncScheduleOutcome::AlreadyRunning; + } + + sync_cancel_token.store(false, SeqCst); + sync_pending.store(false, SeqCst); + + let queue_arc = Arc::clone(queue_arc); + let thread_name = format!("queue-sync-{}", renderer_id); + + thread::Builder::new() + .name(thread_name) + .spawn(move || { + struct Guard(Arc); + impl Drop for Guard { + fn drop(&mut self) { + self.0.store(false, SeqCst); + } + } + let _guard = Guard(Arc::clone(&sync_in_progress)); + + let mut current_items = items; + let mut current_on_ready = Some(on_ready); + let mut on_complete = Some(on_complete); + + tracing::debug!(thread = %std::thread::current().name().unwrap_or("?"), "queue-sync thread started"); + + loop { + sync_pending.store(false, SeqCst); + sync_cancel_token.store(false, SeqCst); + + // Extract the real on_ready BEFORE locking the queue. + // on_ready may call play_from_queue() which re-locks the queue, + // so we must NOT call it while holding queue_arc. + let real_on_ready = current_on_ready.take().flatten(); + let on_ready_triggered = Arc::new(AtomicBool::new(false)); + let proxy_on_ready: Option> = + real_on_ready.as_ref().map(|_| { + let flag = Arc::clone(&on_ready_triggered); + Box::new(move || { + flag.store(true, SeqCst); + }) as Box + }); + + tracing::debug!( + thread = %std::thread::current().name().unwrap_or("?"), + items = current_items.len(), + has_on_ready = real_on_ready.is_some(), + "queue-sync: calling sync_queue" + ); + + let result = { + let mut q = queue_arc.lock().unwrap(); + ::sync_queue( + &mut q, + current_items, + &sync_cancel_token, + proxy_on_ready, + ) + }; + // Queue lock is released here. + // Now safe to call on_ready (which may re-lock the queue). + if on_ready_triggered.load(SeqCst) { + tracing::debug!( + thread = %std::thread::current().name().unwrap_or("?"), + "queue-sync: on_ready triggered, calling callback" + ); + if let Some(f) = real_on_ready { + f(); + } + } else if real_on_ready.is_some() { + tracing::debug!( + thread = %std::thread::current().name().unwrap_or("?"), + "queue-sync: on_ready not triggered (cancelled or skipped)" + ); + } + + match result { + Err(ControlPointError::SyncCancelled) => { + tracing::debug!( + thread = %std::thread::current().name().unwrap_or("?"), + "queue-sync: cancelled" + ); + } + Err(e) => { + tracing::warn!("queue-sync error: {}", e); + } + Ok(()) => { + tracing::debug!( + thread = %std::thread::current().name().unwrap_or("?"), + "queue-sync: completed successfully" + ); + if let Some(cb) = on_complete.take() { + let queue_len = queue_arc.lock().unwrap().len().unwrap_or(0); + cb(queue_len); + } + } + } + + if !sync_pending.load(SeqCst) { + break; + } + + match pending_items_fn() { + Ok(new_items) => { + current_items = new_items; + current_on_ready = Some(None); + } + Err(e) => { + tracing::warn!("queue-sync pending re-fetch error: {}", e); + break; + } + } + } + + tracing::debug!(thread = %std::thread::current().name().unwrap_or("?"), "queue-sync thread done"); + }) + .expect("Failed to spawn queue-sync thread"); + + SyncScheduleOutcome::Scheduled + } } impl QueueBackend for MusicQueue { // Primitives fn len(&self) -> Result { - match self { - MusicQueue::Internal(q) => q.len(), - MusicQueue::OpenHome(q) => q.len(), + match &self.backend { + MusicQueueBackend::Internal(q) => q.len(), + MusicQueueBackend::OpenHome(q) => q.len(), } } fn track_ids(&self) -> Result, ControlPointError> { - match self { - MusicQueue::Internal(q) => q.track_ids(), - MusicQueue::OpenHome(q) => q.track_ids(), + match &self.backend { + MusicQueueBackend::Internal(q) => q.track_ids(), + MusicQueueBackend::OpenHome(q) => q.track_ids(), } } fn id_to_position(&self, id: u32) -> Result { - match self { - MusicQueue::Internal(q) => q.id_to_position(id), - MusicQueue::OpenHome(q) => q.id_to_position(id), + match &self.backend { + MusicQueueBackend::Internal(q) => q.id_to_position(id), + MusicQueueBackend::OpenHome(q) => q.id_to_position(id), } } fn position_to_id(&self, id: usize) -> Result { - match self { - MusicQueue::Internal(q) => q.position_to_id(id), - MusicQueue::OpenHome(q) => q.position_to_id(id), + match &self.backend { + MusicQueueBackend::Internal(q) => q.position_to_id(id), + MusicQueueBackend::OpenHome(q) => q.position_to_id(id), } } fn current_track(&self) -> Result, ControlPointError> { - match self { - MusicQueue::Internal(q) => q.current_track(), - MusicQueue::OpenHome(q) => q.current_track(), + match &self.backend { + MusicQueueBackend::Internal(q) => q.current_track(), + MusicQueueBackend::OpenHome(q) => q.current_track(), } } fn current_index(&self) -> Result, ControlPointError> { - match self { - MusicQueue::Internal(q) => q.current_index(), - MusicQueue::OpenHome(q) => q.current_index(), + match &self.backend { + MusicQueueBackend::Internal(q) => q.current_index(), + MusicQueueBackend::OpenHome(q) => q.current_index(), } } fn queue_snapshot(&self) -> Result { - match self { - MusicQueue::Internal(q) => q.queue_snapshot(), - MusicQueue::OpenHome(q) => q.queue_snapshot(), + match &self.backend { + MusicQueueBackend::Internal(q) => q.queue_snapshot(), + MusicQueueBackend::OpenHome(q) => q.queue_snapshot(), } } fn set_index(&mut self, index: Option) -> Result<(), ControlPointError> { - match self { - MusicQueue::Internal(q) => q.set_index(index), - MusicQueue::OpenHome(q) => q.set_index(index), + match &mut self.backend { + MusicQueueBackend::Internal(q) => q.set_index(index), + MusicQueueBackend::OpenHome(q) => q.set_index(index), } } @@ -89,30 +302,35 @@ impl QueueBackend for MusicQueue { items: Vec, current_index: Option, ) -> Result<(), ControlPointError> { - match self { - MusicQueue::Internal(q) => q.replace_queue(items, current_index), - MusicQueue::OpenHome(q) => q.replace_queue(items, current_index), + match &mut self.backend { + MusicQueueBackend::Internal(q) => q.replace_queue(items, current_index), + MusicQueueBackend::OpenHome(q) => q.replace_queue(items, current_index), } } - fn sync_queue(&mut self, items: Vec) -> Result<(), ControlPointError> { - match self { - MusicQueue::Internal(q) => q.sync_queue(items), - MusicQueue::OpenHome(q) => q.sync_queue(items), + fn sync_queue( + &mut self, + items: Vec, + cancel_token: &Arc, + on_ready: Option>, + ) -> Result<(), ControlPointError> { + match &mut self.backend { + MusicQueueBackend::Internal(q) => q.sync_queue(items, cancel_token, on_ready), + MusicQueueBackend::OpenHome(q) => q.sync_queue(items, cancel_token, on_ready), } } fn get_item(&self, index: usize) -> Result, ControlPointError> { - match self { - MusicQueue::Internal(q) => q.get_item(index), - MusicQueue::OpenHome(q) => q.get_item(index), + match &self.backend { + MusicQueueBackend::Internal(q) => q.get_item(index), + MusicQueueBackend::OpenHome(q) => q.get_item(index), } } fn replace_item(&mut self, index: usize, item: PlaybackItem) -> Result<(), ControlPointError> { - match self { - MusicQueue::Internal(q) => q.replace_item(index, item), - MusicQueue::OpenHome(q) => q.replace_item(index, item), + match &mut self.backend { + MusicQueueBackend::Internal(q) => q.replace_item(index, item), + MusicQueueBackend::OpenHome(q) => q.replace_item(index, item), } } @@ -121,59 +339,59 @@ impl QueueBackend for MusicQueue { items: Vec, mode: EnqueueMode, ) -> Result<(), ControlPointError> { - match self { - MusicQueue::Internal(q) => q.enqueue_items(items, mode), - MusicQueue::OpenHome(q) => q.enqueue_items(items, mode), + match &mut self.backend { + MusicQueueBackend::Internal(q) => q.enqueue_items(items, mode), + MusicQueueBackend::OpenHome(q) => q.enqueue_items(items, mode), } } // Optimized helpers fn clear_queue(&mut self) -> Result<(), ControlPointError> { - match self { - MusicQueue::Internal(q) => q.clear_queue(), - MusicQueue::OpenHome(q) => q.clear_queue(), + match &mut self.backend { + MusicQueueBackend::Internal(q) => q.clear_queue(), + MusicQueueBackend::OpenHome(q) => q.clear_queue(), } } fn is_empty(&self) -> Result { - match self { - MusicQueue::Internal(q) => q.is_empty(), - MusicQueue::OpenHome(q) => q.is_empty(), + match &self.backend { + MusicQueueBackend::Internal(q) => q.is_empty(), + MusicQueueBackend::OpenHome(q) => q.is_empty(), } } fn upcoming_len(&self) -> Result { - match self { - MusicQueue::Internal(q) => q.upcoming_len(), - MusicQueue::OpenHome(q) => q.upcoming_len(), + match &self.backend { + MusicQueueBackend::Internal(q) => q.upcoming_len(), + MusicQueueBackend::OpenHome(q) => q.upcoming_len(), } } fn upcoming_items(&self) -> Result, ControlPointError> { - match self { - MusicQueue::Internal(q) => q.upcoming_items(), - MusicQueue::OpenHome(q) => q.upcoming_items(), + match &self.backend { + MusicQueueBackend::Internal(q) => q.upcoming_items(), + MusicQueueBackend::OpenHome(q) => q.upcoming_items(), } } fn peek_current(&mut self) -> Result, ControlPointError> { - match self { - MusicQueue::Internal(q) => q.peek_current(), - MusicQueue::OpenHome(q) => q.peek_current(), + match &mut self.backend { + MusicQueueBackend::Internal(q) => q.peek_current(), + MusicQueueBackend::OpenHome(q) => q.peek_current(), } } fn dequeue_next(&mut self) -> Result, ControlPointError> { - match self { - MusicQueue::Internal(q) => q.dequeue_next(), - MusicQueue::OpenHome(q) => q.dequeue_next(), + match &mut self.backend { + MusicQueueBackend::Internal(q) => q.dequeue_next(), + MusicQueueBackend::OpenHome(q) => q.dequeue_next(), } } fn append_or_init_index(&mut self, items: Vec) -> Result<(), ControlPointError> { - match self { - MusicQueue::Internal(q) => q.append_or_init_index(items), - MusicQueue::OpenHome(q) => q.append_or_init_index(items), + match &mut self.backend { + MusicQueueBackend::Internal(q) => q.append_or_init_index(items), + MusicQueueBackend::OpenHome(q) => q.append_or_init_index(items), } } } diff --git a/pmocontrol/src/queue/openhome.rs b/pmocontrol/src/queue/openhome.rs index 6e0dc3c5..69ae1eb1 100644 --- a/pmocontrol/src/queue/openhome.rs +++ b/pmocontrol/src/queue/openhome.rs @@ -1,5 +1,5 @@ use std::collections::HashMap; -use std::sync::{Arc, Mutex}; +use std::sync::{atomic::AtomicBool, Arc, Mutex}; use std::time::SystemTime; use std::usize; @@ -544,32 +544,44 @@ impl OpenHomeQueue { &mut self, new_items: Vec, playing_id: usize, + cancel_token: &Arc, + on_ready: &mut Option>, ) -> Result<(), ControlPointError> { - // Get current track IDs from OpenHome + use std::sync::atomic::Ordering::SeqCst; + let current_track_ids = self.track_ids()?; - // Delete everything except the currently playing item - // Using delete_id_if_exists() to handle cases where another control point - // may have already modified the playlist for &track_id in current_track_ids.iter().rev() { + if cancel_token.load(SeqCst) { + return Err(ControlPointError::SyncCancelled); + } if track_id as usize != playing_id { self.playlist_client.delete_id_if_exists(track_id)?; self.metadata_cache.lock().unwrap().remove(&track_id); } } - // Insert new items after the currently playing track let mut previous_id = playing_id as u32; + let mut first_insert_done = false; for item in new_items { + if cancel_token.load(SeqCst) { + return Err(ControlPointError::SyncCancelled); + } let metadata = build_metadata_xml(&item); let new_id = self .playlist_client .insert(previous_id, &item.uri, &metadata)?; - // Enregistrer les métadonnées dans le cache self.cache_metadata(new_id, item.metadata, &item.uri); previous_id = new_id; + + if !first_insert_done && on_ready.is_some() { + first_insert_done = true; + if let Some(f) = on_ready.take() { + f(); + } + } } debug!( @@ -577,7 +589,6 @@ impl OpenHomeQueue { "Gentle sync completed: preserved playing track as first item (not in new playlist)" ); - // Invalidate cache after playlist modifications self.invalidate_track_caches(); Ok(()) @@ -589,8 +600,13 @@ impl OpenHomeQueue { old_ids: &[u32], keep_flags: &[bool], position_label: &str, + cancel_token: &Arc, ) -> Result<(), ControlPointError> { + use std::sync::atomic::Ordering::SeqCst; for (idx, &track_id) in old_ids.iter().enumerate().rev() { + if cancel_token.load(SeqCst) { + return Err(ControlPointError::SyncCancelled); + } if !keep_flags[idx] { debug!( renderer = self.renderer_id.0.as_str(), @@ -615,8 +631,11 @@ impl OpenHomeQueue { keep_old_flags: &[bool], mut previous_id: u32, position_label: &str, + cancel_token: &Arc, + on_ready: &mut Option>, ) -> Result { - // Collect IDs of kept items (in order) + use std::sync::atomic::Ordering::SeqCst; + let remaining_ids: Vec = old_ids .iter() .enumerate() @@ -624,15 +643,17 @@ impl OpenHomeQueue { .collect(); let mut remaining_idx = 0; + let mut first_insert_done = false; - // Rebuild section for (idx, item) in new_items.iter().enumerate() { + if cancel_token.load(SeqCst) { + return Err(ControlPointError::SyncCancelled); + } if keep_new_flags[idx] { let existing_id = remaining_ids[remaining_idx]; remaining_idx += 1; previous_id = existing_id; - // Mettre à jour les métadonnées de l'item existant conservé self.cache_metadata(existing_id, item.metadata.clone(), &item.uri); debug!( @@ -648,7 +669,6 @@ impl OpenHomeQueue { .playlist_client .insert(previous_id, &item.uri, &metadata)?; - // Enregistrer les métadonnées du nouvel item self.cache_metadata(new_id, item.metadata.clone(), &item.uri); debug!( @@ -661,6 +681,13 @@ impl OpenHomeQueue { new_id ); previous_id = new_id; + + if !first_insert_done && on_ready.is_some() { + first_insert_done = true; + if let Some(f) = on_ready.take() { + f(); + } + } } } @@ -677,6 +704,8 @@ impl OpenHomeQueue { pivot_id: usize, snapshot: &QueueSnapshot, current_track_ids: &[u32], + cancel_token: &Arc, + on_ready: &mut Option>, ) -> Result<(), ControlPointError> { // Find the pivot index in our current state let pivot_idx = current_track_ids @@ -705,10 +734,15 @@ impl OpenHomeQueue { 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")?; + self.delete_marked_items(&old_ids_after, &keep_old_after, "AFTER pivot", cancel_token)?; // Delete items marked for deletion in BEFORE part (reverse order) - self.delete_marked_items(&old_ids_before, &keep_old_before, "BEFORE pivot")?; + self.delete_marked_items( + &old_ids_before, + &keep_old_before, + "BEFORE pivot", + cancel_token, + )?; // Rebuild the playlist: [BEFORE, PIVOT, AFTER] // Rebuild BEFORE part (we don't need the returned previous_id) @@ -719,6 +753,8 @@ impl OpenHomeQueue { &keep_old_before, OPENHOME_PLAYLIST_HEAD_ID, "BEFORE pivot", + cancel_token, + on_ready, )?; // PIVOT keeps its ID and position - it's the anchor point @@ -748,6 +784,8 @@ impl OpenHomeQueue { &keep_old_after, previous_id, "AFTER pivot", + cancel_token, + on_ready, )?; debug!( @@ -769,7 +807,11 @@ impl OpenHomeQueue { items: Vec, snapshot: &QueueSnapshot, current_track_ids: &[u32], + cancel_token: &Arc, + on_ready: &mut Option>, ) -> Result<(), ControlPointError> { + use std::sync::atomic::Ordering::SeqCst; + debug!( renderer = self.renderer_id.0.as_str(), current_count = snapshot.items.len(), @@ -795,32 +837,22 @@ impl OpenHomeQueue { "LCS computed: minimizing OpenHome playlist operations" ); - // Get current playing track ID BEFORE any modifications let current_track_id = self.playlist_client.id().ok().filter(|&id| id != 0); - // Check if the currently playing track is in the new playlist - // If so, we should NOT use delete_all() - we must preserve it let current_track_in_new_playlist = current_track_id.and_then(|current_id| { items .iter() .position(|item| item.backend_id as u32 == current_id) }); - // If we're replacing everything (keep=0), use delete_all() BUT only if - // there's no currently playing track, OR if the current track is not in the new playlist. - // If current track IS in new playlist, we must preserve it using insert/delete operations. if items_to_keep == 0 && items_to_delete > 0 { if current_track_in_new_playlist.is_some() { - // Current track is in new playlist - use insert/delete instead of delete_all - // to preserve playback debug!( renderer = self.renderer_id.0.as_str(), current_track_in_playlist = true, "Preserving currently playing track - using insert/delete instead of delete_all" ); - // Fall through to selective deletion below } else { - // No current track or not in new playlist - safe to use delete_all debug!( renderer = self.renderer_id.0.as_str(), "Using delete_all() for complete replacement (safe - no current track or not in new playlist)" @@ -829,19 +861,18 @@ impl OpenHomeQueue { self.metadata_cache.lock().unwrap().clear(); } } else { - // Selective deletion when keeping some items for idx in (0..current_track_ids.len()).rev() { + if cancel_token.load(SeqCst) { + return Err(ControlPointError::SyncCancelled); + } if !keep_current[idx] { let track_id = current_track_ids[idx]; - // Use delete_id_if_exists() to handle cases where another control point - // may have already modified the playlist self.playlist_client.delete_id_if_exists(track_id)?; self.metadata_cache.lock().unwrap().remove(&track_id); } } } - // Rebuild by inserting new items let remaining_ids: Vec = current_track_ids .iter() .enumerate() @@ -856,8 +887,12 @@ impl OpenHomeQueue { let mut remaining_idx = 0usize; let mut previous_id = OPENHOME_PLAYLIST_HEAD_ID; + let mut first_insert_done = false; for (idx, item) in items.into_iter().enumerate() { + if cancel_token.load(SeqCst) { + return Err(ControlPointError::SyncCancelled); + } if keep_desired[idx] { if remaining_idx >= remaining_ids.len() { return Err(ControlPointError::OpenHomeError(format!( @@ -868,8 +903,6 @@ impl OpenHomeQueue { remaining_idx += 1; previous_id = existing_id; - // 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, &item.uri); } else { let metadata = build_metadata_xml(&item); @@ -877,10 +910,16 @@ impl OpenHomeQueue { .playlist_client .insert(previous_id, &item.uri, &metadata)?; - // Enregistrer les métadonnées du nouvel item self.cache_metadata(new_id, item.metadata, &item.uri); previous_id = new_id; + + if !first_insert_done && on_ready.is_some() { + first_insert_done = true; + if let Some(f) = on_ready.take() { + f(); + } + } } } @@ -890,7 +929,6 @@ impl OpenHomeQueue { ))); } - // Invalidate cache after playlist modifications self.invalidate_track_caches(); Ok(()) @@ -1274,10 +1312,20 @@ impl QueueBackend for OpenHomeQueue { Ok(()) } - fn sync_queue(&mut self, items: Vec) -> Result<(), ControlPointError> { + fn sync_queue( + &mut self, + items: Vec, + cancel_token: &Arc, + mut on_ready: Option>, + ) -> Result<(), ControlPointError> { + use std::sync::atomic::Ordering::SeqCst; + self.ensure_playlist_source_selected()?; - // DIAGNOSTIC: Log current track state before any modifications + if cancel_token.load(SeqCst) { + return Err(ControlPointError::SyncCancelled); + } + let pre_current_track = self.playlist_client.id().ok(); tracing::warn!( renderer = self.renderer_id.0.as_str(), @@ -1293,10 +1341,8 @@ 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.invalidate_all_caches(); - // DIAGNOSTIC: Log state after delete_all let post_current_track = self.playlist_client.id().ok(); tracing::warn!( renderer = self.renderer_id.0.as_str(), @@ -1306,8 +1352,6 @@ 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!( @@ -1315,23 +1359,25 @@ impl QueueBackend for OpenHomeQueue { 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 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 { + if cancel_token.load(SeqCst) { + return Err(ControlPointError::SyncCancelled); + } 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(()); } @@ -1341,11 +1387,12 @@ impl QueueBackend for OpenHomeQueue { 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() { + if cancel_token.load(SeqCst) { + return Err(ControlPointError::SyncCancelled); + } self.playlist_client.delete_id(*id)?; } - // Invalidate caches after deletes self.invalidate_all_caches(); return Ok(()); } @@ -1357,9 +1404,6 @@ impl QueueBackend for OpenHomeQueue { } } - // 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. let snapshot = self.queue_snapshot()?; debug!( @@ -1370,9 +1414,6 @@ impl QueueBackend for OpenHomeQueue { "sync_queue: snapshot vs new items comparison" ); - // Note: current_index may point to an index that doesn't exist in items - // if the OpenHome renderer is in an inconsistent state (e.g., IdArray returns - // IDs but ReadList returns empty TrackList). We must bounds-check here. let playing_info = snapshot.current_index.and_then(|idx| { if idx < snapshot.items.len() { Some(( @@ -1399,8 +1440,14 @@ impl QueueBackend for OpenHomeQueue { "OpenHome playlist state" ); + let has_pivot = playing_info.is_some(); + if has_pivot { + if let Some(f) = on_ready.take() { + f(); + } + } + if let Some((playing_idx, playing_id, playing_uri, playing_didl_id)) = playing_info { - // Find if the currently playing item is in the new playlist (by URI first, then by didl_id) let new_playing_idx = items .iter() .position(|item| item.uri == playing_uri) @@ -1420,8 +1467,6 @@ impl QueueBackend for OpenHomeQueue { ); if let Some(pivot_idx) = new_playing_idx { - // CASE 2: Currently playing item IS in the new playlist - // Use gentle double-LCS strategy: preserve the pivot and sync before/after separately debug!( renderer = self.renderer_id.0.as_str(), playing_idx, @@ -1438,22 +1483,24 @@ impl QueueBackend for OpenHomeQueue { playing_id, &snapshot, ¤t_ids_for_pivot, + cancel_token, + &mut on_ready, )?; } else { - // CASE 1: Currently playing item NOT in the new playlist - // Keep it as first item and append the new playlist after it debug!( renderer = self.renderer_id.0.as_str(), playing_idx, "Gentle sync: currently playing item not in new playlist, preserving as first item" ); - self.replace_queue_preserve_current(items, playing_id)?; + self.replace_queue_preserve_current( + items, + playing_id, + cancel_token, + &mut on_ready, + )?; } } else { - // No currently playing item or can't determine it - use standard LCS - // BUT first check if this is because the OpenHome device returned empty playlist - // This could cause the queue to be cleared incorrectly if snapshot.items.is_empty() && !items.is_empty() { tracing::warn!( renderer = self.renderer_id.0.as_str(), @@ -1461,8 +1508,6 @@ impl QueueBackend for OpenHomeQueue { new_items = items.len(), "OpenHome playlist appears empty - possible stale cache or device issue, NOT clearing queue" ); - // Don't call replace_queue_standard_lcs with empty snapshot - it would clear our queue - // Instead, just add the new items without deleting existing ones return self.enqueue_items(items, crate::queue::EnqueueMode::AppendToEnd); } @@ -1472,10 +1517,15 @@ impl QueueBackend for OpenHomeQueue { ); 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)?; + self.replace_queue_standard_lcs( + items, + &snapshot, + ¤t_ids_for_lcs, + cancel_token, + &mut on_ready, + )?; } - // DIAGNOSTIC: Log state after sync completes let post_current_track = self.playlist_client.id().ok(); let post_ids = self.track_ids(); tracing::warn!( @@ -1694,6 +1744,6 @@ impl QueueFromRendererInfo for OpenHomeQueue { } fn to_backend(self) -> MusicQueue { - MusicQueue::OpenHome(self) + MusicQueue::from_openhome(self) } } diff --git a/pmocontrol/src/queue/snapshot.rs b/pmocontrol/src/queue/snapshot.rs index 7026f93d..58557d46 100644 --- a/pmocontrol/src/queue/snapshot.rs +++ b/pmocontrol/src/queue/snapshot.rs @@ -1,4 +1,4 @@ -use crate::{DeviceId, model::TrackMetadata}; +use crate::{model::TrackMetadata, DeviceId}; /// Canonical representation of a track in a renderer queue. /// diff --git a/pmocontrol/src/registry.rs b/pmocontrol/src/registry.rs index 2e8a086d..fbf3c907 100644 --- a/pmocontrol/src/registry.rs +++ b/pmocontrol/src/registry.rs @@ -136,6 +136,14 @@ impl DeviceRegistry { self.devices.get(id)?.as_music_renderer().ok() } + pub fn get_renderer_queue_arc( + &self, + id: &DeviceId, + ) -> Option>> { + let renderer = self.get_renderer(id)?; + Some(renderer.queue()) + } + pub fn get_server(&self, id: &DeviceId) -> Option> { self.devices.get(id)?.as_music_server().ok() } diff --git a/pmocontrol/src/sse.rs b/pmocontrol/src/sse.rs index 1d29580f..0f68289c 100644 --- a/pmocontrol/src/sse.rs +++ b/pmocontrol/src/sse.rs @@ -153,6 +153,14 @@ pub enum RendererEventPayload { renderer_id: String, timestamp: chrono::DateTime, }, + QueueReadyToPlay { + renderer_id: String, + timestamp: chrono::DateTime, + }, + QueueSyncCancelled { + renderer_id: String, + timestamp: chrono::DateTime, + }, BindingChanged { renderer_id: String, server_id: Option, @@ -293,6 +301,14 @@ async fn renderer_event_to_payload( renderer_id: id.0, timestamp, }, + RendererEvent::QueueReadyToPlay { id } => RendererEventPayload::QueueReadyToPlay { + renderer_id: id.0, + timestamp, + }, + RendererEvent::QueueSyncCancelled { id } => RendererEventPayload::QueueSyncCancelled { + renderer_id: id.0, + timestamp, + }, RendererEvent::BindingChanged { id, binding } => RendererEventPayload::BindingChanged { renderer_id: id.0, server_id: binding.as_ref().map(|b| b.server_id.0.clone()), diff --git a/version.txt b/version.txt index a02af241..b463c010 100644 --- a/version.txt +++ b/version.txt @@ -1 +1 @@ -0.3.40 +0.3.41