From 9dee19394709fc2d76dcfea202c6d1e46af584c4 Mon Sep 17 00:00:00 2001 From: Eric Coissac Date: Thu, 9 Apr 2026 11:19:54 +0200 Subject: [PATCH] =?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