Compare commits
6 Commits
6285c8ccc8
...
9ad7d7454d
| Author | SHA1 | Date | |
|---|---|---|---|
| 9ad7d7454d | |||
| 39d7c56a99 | |||
| 82197c8103 | |||
| 8ddee7d94f | |||
| 4b793cec59 | |||
| 9e447023a8 |
323
Blackboard/Todo/refactoring_pmocontrol.md
Normal file
323
Blackboard/Todo/refactoring_pmocontrol.md
Normal file
@@ -0,0 +1,323 @@
|
||||
# Refactoring Plan - Crate `pmocontrol`
|
||||
|
||||
## Contexte
|
||||
|
||||
La crate `pmocontrol` implémente un control point UPnP multiprotocole pour contrôler des renderers audio (UPnP/DLNA, OpenHome, LinkPlay, Arylic TCP, Chromecast). Une refactorisation récente avait pour objectif de monter la logique vers les couches abstraites, mais des duplications et des problèmes de conception subsistent.
|
||||
|
||||
---
|
||||
|
||||
## Structure analysée
|
||||
|
||||
- `music_renderer/` : Implémentations concrètes + façade `MusicRenderer`
|
||||
- `queue/` : Gestion abstraite et concrète des files de lecture
|
||||
- `discovery/` : Découverte SSDP et gestion des appareils
|
||||
- `upnp_clients/` : Clients SOAP pour services UPnP
|
||||
- `control_point.rs` : Point de contrôle principal
|
||||
|
||||
Backends : `UpnpRenderer`, `OpenHomeRenderer`, `LinkPlayRenderer`, `ArylicTcpRenderer`, `ChromecastRenderer`, `HybridUpnpArylicRenderer`
|
||||
|
||||
---
|
||||
|
||||
## CATÉGORIE P0 : BUGS LOGIQUES (à corriger immédiatement)
|
||||
|
||||
### BUG-1 : `sync_queue` dans UpnpRenderer ignore le cancel_token
|
||||
|
||||
**Fichier :** `src/music_renderer/upnp_renderer.rs` (lignes ~354-364)
|
||||
|
||||
**Description :** Le paramètre `cancel_token` est reçu comme `_cancel_token` (ignoré) et remplacé par un `Arc::new(AtomicBool::new(false))` fraîchement créé. Les demandes d'annulation de synchronisation de queue sont silencieusement ignorées pour le backend UPnP.
|
||||
|
||||
**Correction :**
|
||||
```rust
|
||||
fn sync_queue(
|
||||
&mut self,
|
||||
items: Vec<PlaybackItem>,
|
||||
cancel_token: &Arc<AtomicBool>, // utiliser le param, pas _cancel_token
|
||||
on_ready: Option<Box<dyn FnOnce() + Send>>,
|
||||
) -> Result<(), ControlPointError> {
|
||||
self.queue
|
||||
.lock()
|
||||
.unwrap()
|
||||
.sync_queue(items, cancel_token, on_ready) // passer le vrai token
|
||||
}
|
||||
```
|
||||
|
||||
**Tâche :** Vérifier également les autres backends (OpenHome, LinkPlay, Arylic, Chromecast) s'ils propagent correctement le cancel_token.
|
||||
|
||||
---
|
||||
|
||||
## CATÉGORIE P1 : DUPLICATIONS MAJEURES (à traiter en priorité)
|
||||
|
||||
### DUP-1 : Implémentation de `QueueBackend` répétée dans les 5+ renderers
|
||||
|
||||
**Fichiers :**
|
||||
- `src/music_renderer/upnp_renderer.rs` (~310-381)
|
||||
- `src/music_renderer/arylic_tcp.rs` (~368-450)
|
||||
- `src/music_renderer/linkplay_renderer.rs` (~246-330)
|
||||
- `src/music_renderer/chromecast_renderer.rs` (~866+)
|
||||
- `src/music_renderer/openhome_renderer.rs` (~629+)
|
||||
|
||||
**Description :** Chaque renderer implémente `QueueBackend` de manière identique : chaque méthode verrouille `self.queue` et délègue à la file sous-jacente. ~150+ lignes de boilerplate.
|
||||
|
||||
**Approche recommandée — Trait délégateur :**
|
||||
```rust
|
||||
// Dans queue/mod.rs ou music_renderer/mod.rs
|
||||
pub trait HasQueue {
|
||||
fn queue(&self) -> &Arc<Mutex<MusicQueue>>;
|
||||
}
|
||||
|
||||
// Impl automatique pour QueueBackend si le type implémente HasQueue
|
||||
impl<T: HasQueue> QueueBackend for T {
|
||||
fn len(&self) -> Result<usize, ControlPointError> {
|
||||
self.queue().lock().unwrap().len()
|
||||
}
|
||||
fn track_ids(&self) -> Result<Vec<u32>, ControlPointError> {
|
||||
self.queue().lock().unwrap().track_ids()
|
||||
}
|
||||
// ... toutes les méthodes déléguantes
|
||||
}
|
||||
|
||||
// Dans chaque renderer : une seule ligne
|
||||
impl HasQueue for UpnpRenderer {
|
||||
fn queue(&self) -> &Arc<Mutex<MusicQueue>> { &self.queue }
|
||||
}
|
||||
```
|
||||
|
||||
**Tâche :** Définir le trait `HasQueue`, implémenter `QueueBackend for T where T: HasQueue`, supprimer les implémentations manuelles dans chaque renderer.
|
||||
|
||||
---
|
||||
|
||||
### DUP-2 : Logique commune de `play_from_queue` dupliquée dans 4+ renderers
|
||||
|
||||
**Fichiers :**
|
||||
- `src/music_renderer/upnp_renderer.rs` (~184-256)
|
||||
- `src/music_renderer/linkplay_renderer.rs` (~189-212)
|
||||
- `src/music_renderer/arylic_tcp.rs` (~311-334)
|
||||
- `src/music_renderer/openhome_renderer.rs` (~partie similaire)
|
||||
|
||||
**Description :** Les 10-12 premières lignes de `play_from_queue` sont identiques dans tous les renderers : verrouillage de queue, gestion de l'index courant, fallback sur index 0 si non défini, récupération de l'item. Seule la partie terminale (play effectif sur le backend) diffère.
|
||||
|
||||
**Approche recommandée — Méthode par défaut dans un trait :**
|
||||
```rust
|
||||
pub trait QueueTransportControl: HasQueue + HasContinuousStream {
|
||||
// Primitive spécifique au backend
|
||||
fn play_item(&self, item: &PlaybackItem) -> Result<(), ControlPointError>;
|
||||
|
||||
// Implémentation commune par défaut
|
||||
fn play_from_queue(&self) -> Result<(), ControlPointError> {
|
||||
let mut queue = self.queue().lock().unwrap();
|
||||
let current_index = match queue.current_index()? {
|
||||
Some(idx) => idx,
|
||||
None => {
|
||||
if queue.len()? > 0 {
|
||||
queue.set_index(Some(0))?;
|
||||
0
|
||||
} else {
|
||||
return Err(ControlPointError::QueueError("Queue is empty".into()));
|
||||
}
|
||||
}
|
||||
};
|
||||
let item = queue.get_item(current_index)?
|
||||
.ok_or_else(|| ControlPointError::QueueError("Current item not found".into()))?;
|
||||
drop(queue);
|
||||
|
||||
let is_stream = is_continuous_stream_url(&item.uri);
|
||||
*self.continuous_stream().lock().unwrap() = is_stream;
|
||||
self.play_item(&item)
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
**Tâche :** Créer `QueueTransportControl` avec une méthode par défaut, implémenter `play_item` dans chaque renderer, supprimer la logique commune dupliquée.
|
||||
|
||||
---
|
||||
|
||||
### DUP-3 : Initialisation redondante des champs partagés dans tous les renderers
|
||||
|
||||
**Fichiers :** Constructeurs dans tous les fichiers renderer
|
||||
|
||||
**Description :** Chaque renderer répète la même construction :
|
||||
```rust
|
||||
let queue = Arc::new(Mutex::new(MusicQueue::from_renderer_info(info)?));
|
||||
// ...
|
||||
continuous_stream: Arc::new(Mutex::new(false)),
|
||||
```
|
||||
|
||||
**Approche recommandée :**
|
||||
```rust
|
||||
pub struct SharedRendererState {
|
||||
pub queue: Arc<Mutex<MusicQueue>>,
|
||||
pub continuous_stream: Arc<Mutex<bool>>,
|
||||
}
|
||||
|
||||
impl SharedRendererState {
|
||||
pub fn from_renderer_info(info: &RendererInfo) -> Result<Self, ControlPointError> {
|
||||
Ok(Self {
|
||||
queue: Arc::new(Mutex::new(MusicQueue::from_renderer_info(info)?)),
|
||||
continuous_stream: Arc::new(Mutex::new(false)),
|
||||
})
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
**Tâche :** Créer `SharedRendererState`, l'utiliser dans tous les constructeurs de renderers.
|
||||
|
||||
---
|
||||
|
||||
### DUP-4 : `parse_didl_duration` implémentée deux fois différemment
|
||||
|
||||
**Fichiers :**
|
||||
- `src/music_renderer/upnp_renderer.rs` (~383-420) : parsing manuel par string search (fragile)
|
||||
- `src/music_renderer/musicrenderer.rs` (~2089-2117) : via parser DIDL-Lite structuré (robuste)
|
||||
|
||||
**Description :** Deux implémentations divergentes. L'une risque de mal parser du DIDL là où l'autre réussit.
|
||||
|
||||
**Tâche :** Conserver uniquement la version via `DIDLLite::parse`, l'exporter depuis `music_renderer/mod.rs`, supprimer la version par string search dans `upnp_renderer.rs`.
|
||||
|
||||
---
|
||||
|
||||
## CATÉGORIE P2 : ALGORITHMES COMPLEXES ET ABSTRACTIONS MAL PLACÉES
|
||||
|
||||
### ALGO-1 : `schedule_sync` dans `music_queue.rs` — logique intriquée
|
||||
|
||||
**Fichier :** `src/queue/music_queue.rs` (~97-200)
|
||||
|
||||
**Description :** La méthode crée un thread worker avec :
|
||||
- Des `AtomicBool` pour synchronisation (sync_in_progress, sync_pending, sync_cancel_token)
|
||||
- Une boucle infinie interne qui re-tente si un nouveau job arrive
|
||||
- Un Guard RAII basé sur `Drop` pour le cleanup
|
||||
- Des closures capturées mêlant synchronisation et logique métier
|
||||
|
||||
Difficile à tester, à observer de l'extérieur, pas de timeout.
|
||||
|
||||
**Tâche :**
|
||||
1. Extraire la logique du worker dans une fonction `sync_worker_loop` avec signature claire
|
||||
2. Documenter le protocole de synchronisation avec les AtomicBool
|
||||
3. Ajouter une stratégie de timeout ou de sortie en cas de blocage
|
||||
|
||||
---
|
||||
|
||||
### ALGO-2 : Enum dispatch sprawl dans `MusicRendererBackend`
|
||||
|
||||
**Fichier :** `src/music_renderer/musicrenderer.rs` (~2134-2495)
|
||||
|
||||
**Description :** L'enum a 6 variantes. Chaque trait implémenté pour l'enum (`TransportControl`, `PlaybackStatus`, `PlaybackPosition`, `RendererBackend`, `QueueBackend`, etc.) contient un `match` sur les 6 variantes. Estimation : 200+ lignes de boilerplate purement mécanique. Ajouter une 7e variante requiert des mises à jour dans 25+ endroits.
|
||||
|
||||
**Approche recommandée — Macro de dispatch :**
|
||||
```rust
|
||||
macro_rules! dispatch {
|
||||
($self:expr, $method:ident($($arg:expr),*)) => {
|
||||
match $self {
|
||||
MusicRendererBackend::Upnp(b) => b.$method($($arg),*),
|
||||
MusicRendererBackend::OpenHome(b) => b.$method($($arg),*),
|
||||
MusicRendererBackend::LinkPlay(b) => b.$method($($arg),*),
|
||||
MusicRendererBackend::ArylicTcp(b) => b.$method($($arg),*),
|
||||
MusicRendererBackend::Chromecast(b) => b.$method($($arg),*),
|
||||
MusicRendererBackend::HybridUpnpArylic { upnp, .. } => upnp.$method($($arg),*),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl TransportControl for MusicRendererBackend {
|
||||
fn play_uri(&self, uri: &str, meta: &str) -> Result<(), ControlPointError> {
|
||||
dispatch!(self, play_uri(uri, meta))
|
||||
}
|
||||
// ...
|
||||
}
|
||||
```
|
||||
|
||||
**Tâche :** Définir la macro `dispatch!`, remplacer les match statements redondants, valider les cas où HybridUpnpArylic a une logique spéciale.
|
||||
|
||||
---
|
||||
|
||||
### ALGO-3 : Logique de protection des durées de streams dupliquée dans 3 endroits
|
||||
|
||||
**Fichiers :**
|
||||
- `src/queue/interne.rs` (~74-126) : `protect_stream_durations`
|
||||
- `src/queue/openhome.rs` (~400+) : logique similaire pour playlists OpenHome
|
||||
- `src/music_renderer/musicrenderer.rs` (~488-537) : dans `poll_and_emit_changes`
|
||||
|
||||
**Description :** La logique "refuser la diminution de durée pour un stream continu" est réimplémentée trois fois. Si la définition de "diminution acceptable" change, il faut modifier 3 fichiers.
|
||||
|
||||
**Tâche :** Créer `music_renderer/stream_utils.rs` (ou équivalent) avec une fonction `protect_stream_duration(old, new, is_stream) -> Option<String>` et l'utiliser dans les 3 endroits.
|
||||
|
||||
---
|
||||
|
||||
### ALGO-4 : Détection de flux continu fragmentée
|
||||
|
||||
**Fichiers :** `stream_detection.rs`, `musicrenderer.rs`, `queue/interne.rs`, `queue/openhome.rs`
|
||||
|
||||
**Description :** La détection "est-ce un stream continu?" passe par plusieurs chemins non unifiés :
|
||||
1. `TrackMetadata::is_continuous_stream`
|
||||
2. Appel `is_continuous_stream_url(uri)` (réseau)
|
||||
3. Absence de durée dans les métadonnées
|
||||
|
||||
Un stream peut être marqué continu dans une couche mais pas l'autre.
|
||||
|
||||
**Tâche :** Créer une fonction canonique unique :
|
||||
```rust
|
||||
pub fn is_continuous_stream(metadata: Option<&TrackMetadata>, uri: &str) -> bool {
|
||||
metadata.map(|m| m.is_continuous_stream).unwrap_or(false)
|
||||
|| is_continuous_stream_url(uri)
|
||||
}
|
||||
```
|
||||
Faire passer tous les codepaths par cette fonction.
|
||||
|
||||
---
|
||||
|
||||
## CATÉGORIE P3 : BONNES PRATIQUES (amélioration continue)
|
||||
|
||||
### BP-1 : `.unwrap()` sur mutex locks (>50 occurrences)
|
||||
|
||||
**Problème :** Si un mutex est empoisonné (panique dans une autre tâche), `.unwrap()` propage la panique. Aucun code ne gère ce cas.
|
||||
|
||||
**Tâche :** Remplacer `.unwrap()` par `.expect("message contextuel")` à court terme. À long terme, envisager `parking_lot::Mutex` (pas de concept de poison).
|
||||
|
||||
---
|
||||
|
||||
### BP-2 : Absence de gestion d'erreur dans les threads watcher et sync
|
||||
|
||||
**Fichiers :** `musicrenderer.rs` (watcher_loop), `music_queue.rs` (schedule_sync)
|
||||
|
||||
**Tâche :** Ajouter `error!` logs dans les threads et décider explicitement de la politique de redémarrage (continuer vs arrêter).
|
||||
|
||||
---
|
||||
|
||||
### BP-3 : Champs `pub` au lieu de `pub(crate)` dans `PlaylistBinding`
|
||||
|
||||
**Fichier :** `src/music_renderer/musicrenderer.rs` (struct `PlaylistBinding`)
|
||||
|
||||
**Tâche :** Rendre les champs `pub` → `pub(crate)` ou privés avec accesseurs.
|
||||
|
||||
---
|
||||
|
||||
### BP-4 : Documentation manquante sur les contrats des traits
|
||||
|
||||
**Fichiers :** `src/music_renderer/capabilities.rs`, `src/queue/backend.rs`
|
||||
|
||||
**Tâche :** Ajouter des doc-comments sur les traits clés (`TransportControl`, `PlaybackStatus`, `QueueBackend`) décrivant les invariants, les pré/post-conditions, et le comportement attendu.
|
||||
|
||||
---
|
||||
|
||||
## PLAN D'EXÉCUTION
|
||||
|
||||
### Phase 1 — Bugs (immédiat)
|
||||
- [ ] **BUG-1** : Corriger le cancel_token ignoré dans `upnp_renderer.rs::sync_queue`
|
||||
- [ ] Vérifier les autres renderers pour le même bug
|
||||
|
||||
### Phase 2 — Éliminer les duplications majeures (1-2 semaines)
|
||||
- [ ] **DUP-1** : Trait `HasQueue` + impl automatique de `QueueBackend`
|
||||
- [ ] **DUP-4** : Unifier `parse_didl_duration` sur la version DIDL-Lite
|
||||
- [ ] **DUP-3** : Créer `SharedRendererState` pour l'init commune
|
||||
- [ ] **DUP-2** : Trait `QueueTransportControl` avec `play_from_queue` par défaut
|
||||
|
||||
### Phase 3 — Simplifier les algorithmes (2-4 semaines)
|
||||
- [ ] **ALGO-2** : Macro `dispatch!` pour `MusicRendererBackend`
|
||||
- [ ] **ALGO-3** : Centraliser la protection des durées de stream
|
||||
- [ ] **ALGO-4** : Unifier la détection de flux continu
|
||||
- [ ] **ALGO-1** : Refactoriser `schedule_sync` (extraire `sync_worker_loop`)
|
||||
|
||||
### Phase 4 — Qualité continue
|
||||
- [ ] **BP-1** : Remplacer les `.unwrap()` critiques
|
||||
- [ ] **BP-2** : Gestion d'erreur dans les threads
|
||||
- [ ] **BP-3** : Visibilité des champs `PlaylistBinding`
|
||||
- [ ] **BP-4** : Documentation des traits
|
||||
2
Cargo.lock
generated
2
Cargo.lock
generated
@@ -4,7 +4,7 @@ version = 4
|
||||
|
||||
[[package]]
|
||||
name = "PMOMusic"
|
||||
version = "0.3.47"
|
||||
version = "0.3.48"
|
||||
dependencies = [
|
||||
"axum 0.8.7",
|
||||
"console-subscriber",
|
||||
|
||||
@@ -19,7 +19,7 @@ impl RendererEventBus {
|
||||
pub(crate) fn subscribe(&self) -> Receiver<RendererEvent> {
|
||||
let (tx, rx) = unbounded::<RendererEvent>();
|
||||
{
|
||||
let mut subscribers = self.subscribers.lock().unwrap();
|
||||
let mut subscribers = self.subscribers.lock().expect("event subscribers mutex poisoned");
|
||||
subscribers.push(tx);
|
||||
}
|
||||
rx
|
||||
@@ -27,7 +27,7 @@ impl RendererEventBus {
|
||||
|
||||
#[allow(dead_code)]
|
||||
pub(crate) fn broadcast(&self, event: RendererEvent) {
|
||||
let mut subscribers = self.subscribers.lock().unwrap();
|
||||
let mut subscribers = self.subscribers.lock().expect("event subscribers mutex poisoned");
|
||||
subscribers.retain(|tx| tx.send(event.clone()).is_ok());
|
||||
}
|
||||
}
|
||||
@@ -47,14 +47,14 @@ impl MediaServerEventBus {
|
||||
pub fn subscribe(&self) -> Receiver<MediaServerEvent> {
|
||||
let (tx, rx) = unbounded::<MediaServerEvent>();
|
||||
{
|
||||
let mut subscribers = self.subscribers.lock().unwrap();
|
||||
let mut subscribers = self.subscribers.lock().expect("event subscribers mutex poisoned");
|
||||
subscribers.push(tx);
|
||||
}
|
||||
rx
|
||||
}
|
||||
|
||||
pub(crate) fn broadcast(&self, event: MediaServerEvent) {
|
||||
let mut subscribers = self.subscribers.lock().unwrap();
|
||||
let mut subscribers = self.subscribers.lock().expect("event subscribers mutex poisoned");
|
||||
subscribers.retain(|tx| tx.send(event.clone()).is_ok());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -12,14 +12,14 @@ use crate::errors::ControlPointError;
|
||||
use crate::linkplay_client::extract_linkplay_host;
|
||||
use crate::model::{PlaybackState, RendererInfo};
|
||||
use crate::music_renderer::capabilities::{
|
||||
PlaybackPosition, PlaybackPositionInfo, PlaybackStatus, QueueTransportControl, RendererBackend,
|
||||
TransportControl, VolumeControl,
|
||||
HasContinuousStream, PlaybackPosition, PlaybackPositionInfo, PlaybackStatus,
|
||||
QueueTransportControl, 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::HasQueue;
|
||||
use crate::music_renderer::RendererFromMediaRendererInfo;
|
||||
use crate::queue::MusicQueue;
|
||||
use crate::queue::{EnqueueMode, PlaybackItem, QueueBackend, QueueSnapshot};
|
||||
use crate::queue::{EnqueueMode, MusicQueue, PlaybackItem, QueueBackend, QueueSnapshot};
|
||||
use crate::DeviceIdentity;
|
||||
|
||||
/// Raw response from Arylic MCU+PINFGET command
|
||||
@@ -137,7 +137,7 @@ impl RendererFromMediaRendererInfo for ArylicTcpRenderer {
|
||||
impl ArylicTcpRenderer {
|
||||
/// Returns true if currently playing a continuous stream (radio without duration)
|
||||
pub fn is_continuous_stream(&self) -> bool {
|
||||
*self.continuous_stream.lock().unwrap()
|
||||
*self.continuous_stream.lock().expect("continuous_stream mutex poisoned")
|
||||
}
|
||||
|
||||
/// Create an ArylicTcpRenderer with a shared queue (for HybridUpnpArylic)
|
||||
@@ -268,7 +268,7 @@ impl PlaybackPosition for ArylicTcpRenderer {
|
||||
|
||||
// Récupérer les métadonnées depuis la queue (avec protection contre diminution de durée)
|
||||
// Normalement current_index est toujours Some() si la queue n'est pas vide (règle métier)
|
||||
let mut queue_guard = self.queue.lock().unwrap();
|
||||
let mut queue_guard = self.queue.lock().expect("queue mutex poisoned");
|
||||
let queue_item = queue_guard.peek_current().ok().flatten();
|
||||
|
||||
if let Some((current_item, _)) = queue_item {
|
||||
@@ -301,140 +301,23 @@ impl PlaybackPosition for ArylicTcpRenderer {
|
||||
}
|
||||
}
|
||||
|
||||
impl RendererBackend for ArylicTcpRenderer {
|
||||
|
||||
impl QueueTransportControl for ArylicTcpRenderer {
|
||||
fn play_item(&self, item: &PlaybackItem) -> Result<(), ControlPointError> {
|
||||
self.play_uri(&item.uri, "")
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
impl HasQueue for ArylicTcpRenderer {
|
||||
fn queue(&self) -> &Arc<Mutex<MusicQueue>> {
|
||||
&self.queue
|
||||
}
|
||||
}
|
||||
|
||||
impl QueueTransportControl for ArylicTcpRenderer {
|
||||
fn play_from_queue(&self) -> Result<(), ControlPointError> {
|
||||
let mut queue = self.queue.lock().unwrap();
|
||||
|
||||
let current_index = match queue.current_index()? {
|
||||
Some(idx) => idx,
|
||||
None => {
|
||||
if queue.len()? > 0 {
|
||||
queue.set_index(Some(0))?;
|
||||
0
|
||||
} else {
|
||||
return Err(ControlPointError::QueueError("Queue is empty".into()));
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
let item = queue
|
||||
.get_item(current_index)?
|
||||
.ok_or_else(|| ControlPointError::QueueError("Current item not found".into()))?;
|
||||
|
||||
let uri = item.uri.clone();
|
||||
drop(queue);
|
||||
|
||||
self.play_uri(&uri, "")
|
||||
}
|
||||
|
||||
fn play_next(&self) -> Result<(), ControlPointError> {
|
||||
{
|
||||
let mut queue = self.queue.lock().unwrap();
|
||||
if !queue.advance()? {
|
||||
return Err(ControlPointError::QueueError("No next track".into()));
|
||||
}
|
||||
}
|
||||
|
||||
self.play_from_queue()
|
||||
}
|
||||
|
||||
fn play_previous(&self) -> Result<(), ControlPointError> {
|
||||
{
|
||||
let mut queue = self.queue.lock().unwrap();
|
||||
if !queue.rewind()? {
|
||||
return Err(ControlPointError::QueueError("No previous track".into()));
|
||||
}
|
||||
}
|
||||
|
||||
self.play_from_queue()
|
||||
}
|
||||
|
||||
fn play_from_index(&self, index: usize) -> Result<(), ControlPointError> {
|
||||
{
|
||||
let mut queue = self.queue.lock().unwrap();
|
||||
queue.set_index(Some(index))?;
|
||||
}
|
||||
|
||||
self.play_from_queue()
|
||||
}
|
||||
}
|
||||
|
||||
impl QueueBackend for ArylicTcpRenderer {
|
||||
fn len(&self) -> Result<usize, ControlPointError> {
|
||||
self.queue.lock().unwrap().len()
|
||||
}
|
||||
|
||||
fn track_ids(&self) -> Result<Vec<u32>, ControlPointError> {
|
||||
self.queue.lock().unwrap().track_ids()
|
||||
}
|
||||
|
||||
fn id_to_position(&self, id: u32) -> Result<usize, ControlPointError> {
|
||||
self.queue.lock().unwrap().id_to_position(id)
|
||||
}
|
||||
|
||||
fn position_to_id(&self, id: usize) -> Result<u32, ControlPointError> {
|
||||
self.queue.lock().unwrap().position_to_id(id)
|
||||
}
|
||||
|
||||
fn current_track(&self) -> Result<Option<u32>, ControlPointError> {
|
||||
self.queue.lock().unwrap().current_track()
|
||||
}
|
||||
|
||||
fn current_index(&self) -> Result<Option<usize>, ControlPointError> {
|
||||
self.queue.lock().unwrap().current_index()
|
||||
}
|
||||
|
||||
fn queue_snapshot(&self) -> Result<QueueSnapshot, ControlPointError> {
|
||||
self.queue.lock().unwrap().queue_snapshot()
|
||||
}
|
||||
|
||||
fn set_index(&mut self, index: Option<usize>) -> Result<(), ControlPointError> {
|
||||
self.queue.lock().unwrap().set_index(index)
|
||||
}
|
||||
|
||||
fn replace_queue(
|
||||
&mut self,
|
||||
items: Vec<PlaybackItem>,
|
||||
current_index: Option<usize>,
|
||||
) -> Result<(), ControlPointError> {
|
||||
self.queue
|
||||
.lock()
|
||||
.unwrap()
|
||||
.replace_queue(items, current_index)
|
||||
}
|
||||
|
||||
fn sync_queue(
|
||||
&mut self,
|
||||
items: Vec<PlaybackItem>,
|
||||
_cancel_token: &Arc<AtomicBool>,
|
||||
on_ready: Option<Box<dyn FnOnce() + Send>>,
|
||||
) -> Result<(), ControlPointError> {
|
||||
self.queue
|
||||
.lock()
|
||||
.unwrap()
|
||||
.sync_queue(items, &Arc::new(AtomicBool::new(false)), on_ready)
|
||||
}
|
||||
|
||||
fn get_item(&self, index: usize) -> Result<Option<PlaybackItem>, ControlPointError> {
|
||||
self.queue.lock().unwrap().get_item(index)
|
||||
}
|
||||
|
||||
fn replace_item(&mut self, index: usize, item: PlaybackItem) -> Result<(), ControlPointError> {
|
||||
self.queue.lock().unwrap().replace_item(index, item)
|
||||
}
|
||||
|
||||
fn enqueue_items(
|
||||
&mut self,
|
||||
items: Vec<PlaybackItem>,
|
||||
mode: EnqueueMode,
|
||||
) -> Result<(), ControlPointError> {
|
||||
self.queue.lock().unwrap().enqueue_items(items, mode)
|
||||
impl HasContinuousStream for ArylicTcpRenderer {
|
||||
fn continuous_stream(&self) -> &Arc<Mutex<bool>> {
|
||||
&self.continuous_stream
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -1,36 +1,96 @@
|
||||
// pmocontrol/src/capabilities.rs
|
||||
use anyhow::Result;
|
||||
use std::sync::{Arc, Mutex};
|
||||
|
||||
use crate::queue::MusicQueue;
|
||||
use crate::{errors::ControlPointError, model::PlaybackState};
|
||||
use crate::queue::{MusicQueue, QueueBackend};
|
||||
use crate::{errors::ControlPointError, model::PlaybackState, PlaybackItem};
|
||||
|
||||
/// Backend-specific operations for renderers.
|
||||
/// Marker trait for renderer backends that own a `MusicQueue`.
|
||||
///
|
||||
/// This trait provides access to backend-specific resources like the queue.
|
||||
pub trait RendererBackend {
|
||||
/// Returns a reference to the queue associated with this backend.
|
||||
/// Implementing this trait automatically provides the full `QueueBackend`
|
||||
/// blanket implementation (see `queue/backend.rs`). Backends only need to
|
||||
/// return a reference to their `Arc<Mutex<MusicQueue>>` field.
|
||||
pub trait HasQueue {
|
||||
fn queue(&self) -> &Arc<Mutex<MusicQueue>>;
|
||||
}
|
||||
|
||||
/// Marker trait for renderer backends that track stream continuity.
|
||||
///
|
||||
/// The flag is `true` while the renderer is playing a continuous stream
|
||||
/// (e.g. an internet radio station) and `false` for bounded media files.
|
||||
/// It is used by the watcher to decide whether auto-advance should be
|
||||
/// suppressed when playback stops.
|
||||
pub trait HasContinuousStream {
|
||||
fn continuous_stream(&self) -> &Arc<Mutex<bool>>;
|
||||
}
|
||||
|
||||
/// Queue-aware transport control operations.
|
||||
///
|
||||
/// These operations combine queue management with transport control,
|
||||
/// allowing navigation (next/previous) and track selection from the queue.
|
||||
#[allow(dead_code)]
|
||||
pub trait QueueTransportControl {
|
||||
/// Play the next track from the queue.
|
||||
fn play_next(&self) -> Result<(), ControlPointError>;
|
||||
|
||||
/// Play the previous track from the queue.
|
||||
#[allow(dead_code)]
|
||||
fn play_previous(&self) -> Result<(), ControlPointError>;
|
||||
pub trait QueueTransportControl: HasQueue + HasContinuousStream {
|
||||
/// Play a specific item from the queue (backend-specific implementation).
|
||||
fn play_item(&self, item: &PlaybackItem) -> Result<(), ControlPointError>;
|
||||
|
||||
/// Play from the queue at the current index (or initialize to 0 if not set).
|
||||
fn play_from_queue(&self) -> Result<(), ControlPointError>;
|
||||
/// This is the default implementation that handles queue navigation.
|
||||
fn play_from_queue(&self) -> Result<(), ControlPointError> {
|
||||
let mut queue = self.queue().lock().expect("queue mutex poisoned");
|
||||
|
||||
let current_index = match queue.current_index()? {
|
||||
Some(idx) => idx,
|
||||
None => {
|
||||
if queue.len()? > 0 {
|
||||
queue.set_index(Some(0))?;
|
||||
0
|
||||
} else {
|
||||
return Err(ControlPointError::QueueError("Queue is empty".into()));
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
let item = queue
|
||||
.get_item(current_index)?
|
||||
.ok_or_else(|| ControlPointError::QueueError("Current item not found".into()))?;
|
||||
|
||||
drop(queue);
|
||||
|
||||
let is_stream = crate::music_renderer::is_continuous_stream(item.metadata.as_ref(), &item.uri);
|
||||
*self.continuous_stream().lock().expect("continuous_stream mutex poisoned") = is_stream;
|
||||
|
||||
self.play_item(&item)
|
||||
}
|
||||
|
||||
/// Play the next track from the queue.
|
||||
fn play_next(&self) -> Result<(), ControlPointError> {
|
||||
{
|
||||
let mut queue = self.queue().lock().expect("queue mutex poisoned");
|
||||
if !queue.advance()? {
|
||||
return Err(ControlPointError::QueueError("No next track".into()));
|
||||
}
|
||||
}
|
||||
self.play_from_queue()
|
||||
}
|
||||
|
||||
/// Play the previous track from the queue.
|
||||
fn play_previous(&self) -> Result<(), ControlPointError> {
|
||||
{
|
||||
let mut queue = self.queue().lock().expect("queue mutex poisoned");
|
||||
if !queue.rewind()? {
|
||||
return Err(ControlPointError::QueueError("No previous track".into()));
|
||||
}
|
||||
}
|
||||
self.play_from_queue()
|
||||
}
|
||||
|
||||
/// Play from a specific index in the queue.
|
||||
fn play_from_index(&self, index: usize) -> Result<(), ControlPointError>;
|
||||
fn play_from_index(&self, index: usize) -> Result<(), ControlPointError> {
|
||||
{
|
||||
let mut queue = self.queue().lock().expect("queue mutex poisoned");
|
||||
queue.set_index(Some(index))?;
|
||||
}
|
||||
self.play_from_queue()
|
||||
}
|
||||
}
|
||||
|
||||
/// Logical playback position across backends.
|
||||
@@ -47,52 +107,97 @@ pub struct PlaybackPositionInfo {
|
||||
pub track_metadata: Option<String>, // DIDL-Lite XML from GetPositionInfo
|
||||
pub track_uri: Option<String>, // Current track URI
|
||||
}
|
||||
/// Provides the current playback position and track metadata.
|
||||
///
|
||||
/// All time fields use the format `"HH:MM:SS"` (or `None` when unavailable).
|
||||
/// `track_metadata` carries a raw DIDL-Lite XML fragment returned by the device;
|
||||
/// callers that only need structured metadata should use `extract_track_metadata`
|
||||
/// from the watcher module instead.
|
||||
pub trait PlaybackPosition {
|
||||
/// Returns the current playback position information.
|
||||
///
|
||||
/// Returns `Err` if the renderer is unreachable or the query fails.
|
||||
fn playback_position(&self) -> Result<PlaybackPositionInfo, ControlPointError>;
|
||||
}
|
||||
|
||||
/// Generic abstraction for playback status (transport state).
|
||||
/// Generic abstraction for the current transport state.
|
||||
///
|
||||
/// For UPnP AV, this is backed by AVTransport::GetTransportInfo.
|
||||
/// For OpenHome, a future implementation will adapt from OH Info/Time.
|
||||
/// # Implementations
|
||||
///
|
||||
/// - **UPnP AV**: backed by `AVTransport::GetTransportInfo`.
|
||||
/// - **OpenHome**: adapted from OH `Info` / `Time` services.
|
||||
/// - **LinkPlay / Arylic**: mapped from the vendor status response.
|
||||
///
|
||||
/// # Postconditions
|
||||
///
|
||||
/// The returned `PlaybackState` must be one of the canonical values defined by
|
||||
/// the `PlaybackState` enum. Backend-specific states that have no canonical
|
||||
/// equivalent should be mapped to the closest approximation (e.g. "BUFFERING"
|
||||
/// → `PlaybackState::Transitioning`).
|
||||
pub trait PlaybackStatus {
|
||||
/// Returns the current transport state of the renderer.
|
||||
fn playback_state(&self) -> Result<PlaybackState, ControlPointError>;
|
||||
}
|
||||
|
||||
/// Abstraction générique des capacités de transport (lecture / pause / stop / seek)
|
||||
/// indépendamment du protocole sous-jacent (UPnP AV, OpenHome, ...).
|
||||
/// Generic transport control abstraction (play / pause / stop / seek),
|
||||
/// independent of the underlying protocol (UPnP AV, OpenHome, …).
|
||||
///
|
||||
/// # Invariants
|
||||
///
|
||||
/// - `play_uri` sets the active resource and begins playback atomically from the
|
||||
/// caller's perspective. Implementations may split this into two protocol steps
|
||||
/// (e.g. `SetAVTransportURI` + `Play` for UPnP AV) but the caller should not
|
||||
/// need to know.
|
||||
/// - `play` / `pause` / `stop` operate on whatever resource is currently loaded;
|
||||
/// they do not change the queue pointer.
|
||||
/// - `seek_rel_time` uses the format `"HH:MM:SS"`. Backends that do not support
|
||||
/// seeking should return `ControlPointError::NotSupported`.
|
||||
///
|
||||
/// # Relation to `QueueTransportControl`
|
||||
///
|
||||
/// `TransportControl` knows nothing about the queue. `QueueTransportControl`
|
||||
/// extends it with queue-aware navigation (`play_next`, `play_previous`, …).
|
||||
pub trait TransportControl {
|
||||
/// Set la ressource à lire (URI + métadonnées) et/ou commence la lecture.
|
||||
/// Load a resource (URI + DIDL-Lite metadata) and begin playback.
|
||||
///
|
||||
/// Selon l'implémentation, cette méthode peut soit :
|
||||
/// - faire un "Set...URI" + "Play" (cas UPnP AV),
|
||||
/// - ou configurer la file de lecture (cas OpenHome, etc.).
|
||||
/// Depending on the backend this may execute as a single atomic operation or as
|
||||
/// two sequential commands (set resource, then play).
|
||||
fn play_uri(&self, uri: &str, meta: &str) -> Result<(), ControlPointError>;
|
||||
|
||||
/// Démarre ou reprend la lecture.
|
||||
/// Start or resume playback of the currently loaded resource.
|
||||
fn play(&self) -> Result<(), ControlPointError>;
|
||||
|
||||
/// Met la lecture en pause.
|
||||
/// Pause the current playback.
|
||||
fn pause(&self) -> Result<(), ControlPointError>;
|
||||
|
||||
/// Arrête la lecture.
|
||||
/// Stop the current playback and release the loaded resource.
|
||||
fn stop(&self) -> Result<(), ControlPointError>;
|
||||
|
||||
/// Seek à un temps relatif (HH:MM:SS) si supporté.
|
||||
/// Seek to a relative time position expressed as `"HH:MM:SS"`.
|
||||
///
|
||||
/// Returns `ControlPointError::NotSupported` when the backend does not
|
||||
/// implement seeking.
|
||||
fn seek_rel_time(&self, hhmmss: &str) -> Result<(), ControlPointError>;
|
||||
}
|
||||
|
||||
/// Abstraction générique des capacités de contrôle de volume / mute.
|
||||
/// Generic volume and mute control abstraction.
|
||||
///
|
||||
/// # Volume scale
|
||||
///
|
||||
/// Volume values are expressed on the native scale of each renderer.
|
||||
/// UPnP AV and OpenHome renderers typically use 0–100. Callers should
|
||||
/// not assume any particular scale; use the values returned by `volume()`
|
||||
/// as the baseline for relative adjustments.
|
||||
pub trait VolumeControl {
|
||||
/// Retourne le volume logique courant (échelle dépendante du renderer).
|
||||
/// Returns the current logical volume (renderer-specific scale).
|
||||
fn volume(&self) -> Result<u16, ControlPointError>;
|
||||
|
||||
/// Définit le volume logique (échelle dépendante du renderer).
|
||||
/// Sets the logical volume (renderer-specific scale).
|
||||
fn set_volume(&self, v: u16) -> Result<(), ControlPointError>;
|
||||
|
||||
/// Indique si le renderer est muet (mute activé).
|
||||
/// Returns `true` when the renderer is muted.
|
||||
fn mute(&self) -> Result<bool, ControlPointError>;
|
||||
|
||||
/// Active ou désactive le mute.
|
||||
/// Enables (`true`) or disables (`false`) mute.
|
||||
fn set_mute(&self, m: bool) -> Result<(), ControlPointError>;
|
||||
}
|
||||
|
||||
@@ -23,11 +23,12 @@ use crate::discovery::chromecast_discovery::{
|
||||
use crate::errors::ControlPointError;
|
||||
use crate::model::{PlaybackState, RendererInfo};
|
||||
use crate::music_renderer::capabilities::{
|
||||
PlaybackPosition, PlaybackPositionInfo, PlaybackStatus, QueueTransportControl, RendererBackend,
|
||||
TransportControl, VolumeControl,
|
||||
HasContinuousStream, PlaybackPosition, PlaybackPositionInfo, PlaybackStatus,
|
||||
QueueTransportControl, TransportControl, VolumeControl,
|
||||
};
|
||||
use crate::music_renderer::musicrenderer::MusicRendererBackend;
|
||||
use crate::music_renderer::time_utils::{format_hhmmss_f64, parse_hhmmss_strict};
|
||||
use crate::music_renderer::HasQueue;
|
||||
use crate::music_renderer::RendererFromMediaRendererInfo;
|
||||
use crate::queue::{EnqueueMode, MusicQueue, PlaybackItem, QueueBackend, QueueSnapshot};
|
||||
use crate::DeviceIdentity;
|
||||
@@ -171,7 +172,7 @@ impl RendererFromMediaRendererInfo for ChromecastRenderer {
|
||||
impl ChromecastRenderer {
|
||||
/// Returns true if currently playing a continuous stream (radio without duration)
|
||||
pub fn is_continuous_stream(&self) -> bool {
|
||||
*self.continuous_stream.lock().unwrap()
|
||||
*self.continuous_stream.lock().expect("continuous_stream mutex poisoned")
|
||||
}
|
||||
|
||||
/// Connect to the device with retry on connection failures.
|
||||
@@ -204,7 +205,7 @@ impl TransportControl for ChromecastRenderer {
|
||||
|
||||
// Détecte si l'URL est un flux continu
|
||||
let is_stream = crate::music_renderer::is_continuous_stream_url(uri);
|
||||
*self.continuous_stream.lock().unwrap() = is_stream;
|
||||
*self.continuous_stream.lock().expect("continuous_stream mutex poisoned") = is_stream;
|
||||
tracing::debug!(
|
||||
"ChromecastRenderer play_uri: URI={}, continuous_stream={}",
|
||||
uri,
|
||||
@@ -799,139 +800,22 @@ impl VolumeControl for ChromecastRenderer {
|
||||
}
|
||||
}
|
||||
|
||||
impl RendererBackend for ChromecastRenderer {
|
||||
|
||||
impl QueueTransportControl for ChromecastRenderer {
|
||||
fn play_item(&self, item: &PlaybackItem) -> Result<(), ControlPointError> {
|
||||
self.play_uri(&item.uri, "")
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
impl HasQueue for ChromecastRenderer {
|
||||
fn queue(&self) -> &Arc<Mutex<MusicQueue>> {
|
||||
&self.queue
|
||||
}
|
||||
}
|
||||
|
||||
impl QueueTransportControl for ChromecastRenderer {
|
||||
fn play_from_queue(&self) -> Result<(), ControlPointError> {
|
||||
let mut queue = self.queue.lock().unwrap();
|
||||
|
||||
let current_index = match queue.current_index()? {
|
||||
Some(idx) => idx,
|
||||
None => {
|
||||
if queue.len()? > 0 {
|
||||
queue.set_index(Some(0))?;
|
||||
0
|
||||
} else {
|
||||
return Err(ControlPointError::QueueError("Queue is empty".into()));
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
let item = queue
|
||||
.get_item(current_index)?
|
||||
.ok_or_else(|| ControlPointError::QueueError("Current item not found".into()))?;
|
||||
|
||||
let uri = item.uri.clone();
|
||||
drop(queue);
|
||||
|
||||
self.play_uri(&uri, "")
|
||||
}
|
||||
|
||||
fn play_next(&self) -> Result<(), ControlPointError> {
|
||||
{
|
||||
let mut queue = self.queue.lock().unwrap();
|
||||
if !queue.advance()? {
|
||||
return Err(ControlPointError::QueueError("No next track".into()));
|
||||
}
|
||||
}
|
||||
|
||||
self.play_from_queue()
|
||||
}
|
||||
|
||||
fn play_previous(&self) -> Result<(), ControlPointError> {
|
||||
{
|
||||
let mut queue = self.queue.lock().unwrap();
|
||||
if !queue.rewind()? {
|
||||
return Err(ControlPointError::QueueError("No previous track".into()));
|
||||
}
|
||||
}
|
||||
|
||||
self.play_from_queue()
|
||||
}
|
||||
|
||||
fn play_from_index(&self, index: usize) -> Result<(), ControlPointError> {
|
||||
{
|
||||
let mut queue = self.queue.lock().unwrap();
|
||||
queue.set_index(Some(index))?;
|
||||
}
|
||||
|
||||
self.play_from_queue()
|
||||
}
|
||||
}
|
||||
|
||||
impl QueueBackend for ChromecastRenderer {
|
||||
fn len(&self) -> Result<usize, ControlPointError> {
|
||||
self.queue.lock().unwrap().len()
|
||||
}
|
||||
|
||||
fn track_ids(&self) -> Result<Vec<u32>, ControlPointError> {
|
||||
self.queue.lock().unwrap().track_ids()
|
||||
}
|
||||
|
||||
fn id_to_position(&self, id: u32) -> Result<usize, ControlPointError> {
|
||||
self.queue.lock().unwrap().id_to_position(id)
|
||||
}
|
||||
|
||||
fn position_to_id(&self, id: usize) -> Result<u32, ControlPointError> {
|
||||
self.queue.lock().unwrap().position_to_id(id)
|
||||
}
|
||||
|
||||
fn current_track(&self) -> Result<Option<u32>, ControlPointError> {
|
||||
self.queue.lock().unwrap().current_track()
|
||||
}
|
||||
|
||||
fn current_index(&self) -> Result<Option<usize>, ControlPointError> {
|
||||
self.queue.lock().unwrap().current_index()
|
||||
}
|
||||
|
||||
fn queue_snapshot(&self) -> Result<QueueSnapshot, ControlPointError> {
|
||||
self.queue.lock().unwrap().queue_snapshot()
|
||||
}
|
||||
|
||||
fn set_index(&mut self, index: Option<usize>) -> Result<(), ControlPointError> {
|
||||
self.queue.lock().unwrap().set_index(index)
|
||||
}
|
||||
|
||||
fn replace_queue(
|
||||
&mut self,
|
||||
items: Vec<PlaybackItem>,
|
||||
current_index: Option<usize>,
|
||||
) -> Result<(), ControlPointError> {
|
||||
self.queue
|
||||
.lock()
|
||||
.unwrap()
|
||||
.replace_queue(items, current_index)
|
||||
}
|
||||
|
||||
fn sync_queue(
|
||||
&mut self,
|
||||
items: Vec<PlaybackItem>,
|
||||
_cancel_token: &Arc<AtomicBool>,
|
||||
on_ready: Option<Box<dyn FnOnce() + Send>>,
|
||||
) -> Result<(), ControlPointError> {
|
||||
self.queue
|
||||
.lock()
|
||||
.unwrap()
|
||||
.sync_queue(items, &Arc::new(AtomicBool::new(false)), on_ready)
|
||||
}
|
||||
|
||||
fn get_item(&self, index: usize) -> Result<Option<PlaybackItem>, ControlPointError> {
|
||||
self.queue.lock().unwrap().get_item(index)
|
||||
}
|
||||
|
||||
fn replace_item(&mut self, index: usize, item: PlaybackItem) -> Result<(), ControlPointError> {
|
||||
self.queue.lock().unwrap().replace_item(index, item)
|
||||
}
|
||||
|
||||
fn enqueue_items(
|
||||
&mut self,
|
||||
items: Vec<PlaybackItem>,
|
||||
mode: EnqueueMode,
|
||||
) -> Result<(), ControlPointError> {
|
||||
self.queue.lock().unwrap().enqueue_items(items, mode)
|
||||
impl HasContinuousStream for ChromecastRenderer {
|
||||
fn continuous_stream(&self) -> &Arc<Mutex<bool>> {
|
||||
&self.continuous_stream
|
||||
}
|
||||
}
|
||||
|
||||
@@ -10,14 +10,14 @@ use crate::linkplay_client::{
|
||||
};
|
||||
use crate::model::{PlaybackState, RendererInfo};
|
||||
use crate::music_renderer::capabilities::{
|
||||
PlaybackPosition, PlaybackPositionInfo, PlaybackStatus, QueueTransportControl, RendererBackend,
|
||||
TransportControl, VolumeControl,
|
||||
HasContinuousStream, PlaybackPosition, PlaybackPositionInfo, PlaybackStatus,
|
||||
QueueTransportControl, TransportControl, VolumeControl,
|
||||
};
|
||||
use crate::music_renderer::musicrenderer::MusicRendererBackend;
|
||||
use crate::music_renderer::time_utils::parse_hhmmss_strict;
|
||||
use crate::music_renderer::HasQueue;
|
||||
use crate::music_renderer::RendererFromMediaRendererInfo;
|
||||
use crate::queue::MusicQueue;
|
||||
use crate::queue::{EnqueueMode, PlaybackItem, QueueBackend, QueueSnapshot};
|
||||
use crate::queue::{EnqueueMode, MusicQueue, PlaybackItem, QueueBackend};
|
||||
use crate::DeviceIdentity;
|
||||
|
||||
const DEFAULT_HTTP_TIMEOUT_SECS: u64 = 3;
|
||||
@@ -91,7 +91,7 @@ impl RendererFromMediaRendererInfo for LinkPlayRenderer {
|
||||
impl LinkPlayRenderer {
|
||||
/// Returns true if currently playing a continuous stream (radio without duration)
|
||||
pub fn is_continuous_stream(&self) -> bool {
|
||||
*self.continuous_stream.lock().unwrap()
|
||||
*self.continuous_stream.lock().expect("continuous_stream mutex poisoned")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -99,7 +99,7 @@ impl TransportControl for LinkPlayRenderer {
|
||||
fn play_uri(&self, uri: &str, _meta: &str) -> Result<(), ControlPointError> {
|
||||
// Détecte si l'URL est un flux continu
|
||||
let is_stream = crate::music_renderer::is_continuous_stream_url(uri);
|
||||
*self.continuous_stream.lock().unwrap() = is_stream;
|
||||
*self.continuous_stream.lock().expect("continuous_stream mutex poisoned") = is_stream;
|
||||
tracing::debug!(
|
||||
"LinkPlayRenderer play_uri: URI={}, continuous_stream={}",
|
||||
uri,
|
||||
@@ -158,7 +158,7 @@ impl PlaybackPosition for LinkPlayRenderer {
|
||||
let mut position_info = self.fetch_status()?.position_info();
|
||||
|
||||
// Use queue metadata instead of direct status metadata to benefit from duration protection
|
||||
let mut queue_guard = self.queue.lock().unwrap();
|
||||
let mut queue_guard = self.queue.lock().expect("queue mutex poisoned");
|
||||
let queue_item = queue_guard.peek_current().ok().flatten();
|
||||
|
||||
if let Some((current_item, _)) = queue_item {
|
||||
@@ -179,139 +179,22 @@ impl PlaybackPosition for LinkPlayRenderer {
|
||||
}
|
||||
}
|
||||
|
||||
impl RendererBackend for LinkPlayRenderer {
|
||||
|
||||
impl QueueTransportControl for LinkPlayRenderer {
|
||||
fn play_item(&self, item: &PlaybackItem) -> Result<(), ControlPointError> {
|
||||
self.play_uri(&item.uri, "")
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
impl HasQueue for LinkPlayRenderer {
|
||||
fn queue(&self) -> &Arc<Mutex<MusicQueue>> {
|
||||
&self.queue
|
||||
}
|
||||
}
|
||||
|
||||
impl QueueTransportControl for LinkPlayRenderer {
|
||||
fn play_from_queue(&self) -> Result<(), ControlPointError> {
|
||||
let mut queue = self.queue.lock().unwrap();
|
||||
|
||||
let current_index = match queue.current_index()? {
|
||||
Some(idx) => idx,
|
||||
None => {
|
||||
if queue.len()? > 0 {
|
||||
queue.set_index(Some(0))?;
|
||||
0
|
||||
} else {
|
||||
return Err(ControlPointError::QueueError("Queue is empty".into()));
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
let item = queue
|
||||
.get_item(current_index)?
|
||||
.ok_or_else(|| ControlPointError::QueueError("Current item not found".into()))?;
|
||||
|
||||
let uri = item.uri.clone();
|
||||
drop(queue);
|
||||
|
||||
self.play_uri(&uri, "")
|
||||
}
|
||||
|
||||
fn play_next(&self) -> Result<(), ControlPointError> {
|
||||
{
|
||||
let mut queue = self.queue.lock().unwrap();
|
||||
if !queue.advance()? {
|
||||
return Err(ControlPointError::QueueError("No next track".into()));
|
||||
}
|
||||
}
|
||||
|
||||
self.play_from_queue()
|
||||
}
|
||||
|
||||
fn play_previous(&self) -> Result<(), ControlPointError> {
|
||||
{
|
||||
let mut queue = self.queue.lock().unwrap();
|
||||
if !queue.rewind()? {
|
||||
return Err(ControlPointError::QueueError("No previous track".into()));
|
||||
}
|
||||
}
|
||||
|
||||
self.play_from_queue()
|
||||
}
|
||||
|
||||
fn play_from_index(&self, index: usize) -> Result<(), ControlPointError> {
|
||||
{
|
||||
let mut queue = self.queue.lock().unwrap();
|
||||
queue.set_index(Some(index))?;
|
||||
}
|
||||
|
||||
self.play_from_queue()
|
||||
}
|
||||
}
|
||||
|
||||
impl QueueBackend for LinkPlayRenderer {
|
||||
fn len(&self) -> Result<usize, ControlPointError> {
|
||||
self.queue.lock().unwrap().len()
|
||||
}
|
||||
|
||||
fn track_ids(&self) -> Result<Vec<u32>, ControlPointError> {
|
||||
self.queue.lock().unwrap().track_ids()
|
||||
}
|
||||
|
||||
fn id_to_position(&self, id: u32) -> Result<usize, ControlPointError> {
|
||||
self.queue.lock().unwrap().id_to_position(id)
|
||||
}
|
||||
|
||||
fn position_to_id(&self, id: usize) -> Result<u32, ControlPointError> {
|
||||
self.queue.lock().unwrap().position_to_id(id)
|
||||
}
|
||||
|
||||
fn current_track(&self) -> Result<Option<u32>, ControlPointError> {
|
||||
self.queue.lock().unwrap().current_track()
|
||||
}
|
||||
|
||||
fn current_index(&self) -> Result<Option<usize>, ControlPointError> {
|
||||
self.queue.lock().unwrap().current_index()
|
||||
}
|
||||
|
||||
fn queue_snapshot(&self) -> Result<QueueSnapshot, ControlPointError> {
|
||||
self.queue.lock().unwrap().queue_snapshot()
|
||||
}
|
||||
|
||||
fn set_index(&mut self, index: Option<usize>) -> Result<(), ControlPointError> {
|
||||
self.queue.lock().unwrap().set_index(index)
|
||||
}
|
||||
|
||||
fn replace_queue(
|
||||
&mut self,
|
||||
items: Vec<PlaybackItem>,
|
||||
current_index: Option<usize>,
|
||||
) -> Result<(), ControlPointError> {
|
||||
self.queue
|
||||
.lock()
|
||||
.unwrap()
|
||||
.replace_queue(items, current_index)
|
||||
}
|
||||
|
||||
fn sync_queue(
|
||||
&mut self,
|
||||
items: Vec<PlaybackItem>,
|
||||
_cancel_token: &Arc<AtomicBool>,
|
||||
on_ready: Option<Box<dyn FnOnce() + Send>>,
|
||||
) -> Result<(), ControlPointError> {
|
||||
self.queue
|
||||
.lock()
|
||||
.unwrap()
|
||||
.sync_queue(items, &Arc::new(AtomicBool::new(false)), on_ready)
|
||||
}
|
||||
|
||||
fn get_item(&self, index: usize) -> Result<Option<PlaybackItem>, ControlPointError> {
|
||||
self.queue.lock().unwrap().get_item(index)
|
||||
}
|
||||
|
||||
fn replace_item(&mut self, index: usize, item: PlaybackItem) -> Result<(), ControlPointError> {
|
||||
self.queue.lock().unwrap().replace_item(index, item)
|
||||
}
|
||||
|
||||
fn enqueue_items(
|
||||
&mut self,
|
||||
items: Vec<PlaybackItem>,
|
||||
mode: EnqueueMode,
|
||||
) -> Result<(), ControlPointError> {
|
||||
self.queue.lock().unwrap().enqueue_items(items, mode)
|
||||
impl HasContinuousStream for LinkPlayRenderer {
|
||||
fn continuous_stream(&self) -> &Arc<Mutex<bool>> {
|
||||
&self.continuous_stream
|
||||
}
|
||||
}
|
||||
|
||||
@@ -6,7 +6,7 @@ mod upnp_renderer;
|
||||
mod openhome;
|
||||
mod openhome_renderer;
|
||||
|
||||
mod capabilities;
|
||||
pub mod capabilities;
|
||||
mod chromecast_renderer;
|
||||
|
||||
mod musicrenderer;
|
||||
@@ -18,13 +18,14 @@ pub mod watcher;
|
||||
use std::sync::{Arc, Mutex};
|
||||
|
||||
pub use crate::music_renderer::capabilities::{
|
||||
PlaybackPosition, PlaybackPositionInfo, PlaybackStatus,
|
||||
HasContinuousStream, HasQueue, PlaybackPosition, PlaybackPositionInfo, PlaybackStatus,
|
||||
QueueTransportControl, TransportControl, VolumeControl,
|
||||
};
|
||||
pub use crate::music_renderer::musicrenderer::{MusicRenderer, PlaylistBinding};
|
||||
pub use crate::music_renderer::sleep_timer::SleepTimer;
|
||||
pub use crate::music_renderer::stream_detection::is_continuous_stream_url;
|
||||
pub use crate::music_renderer::stream_detection::{is_continuous_stream, is_continuous_stream_url};
|
||||
use crate::{
|
||||
RendererInfo, errors::ControlPointError, music_renderer::musicrenderer::MusicRendererBackend,
|
||||
errors::ControlPointError, music_renderer::musicrenderer::MusicRendererBackend, RendererInfo,
|
||||
};
|
||||
|
||||
pub trait RendererFromMediaRendererInfo {
|
||||
|
||||
@@ -19,10 +19,6 @@ use crate::events::RendererEventBus;
|
||||
use crate::model::RendererEvent;
|
||||
use crate::model::{PlaybackSource, PlaybackState, RendererInfo, RendererProtocol, TrackMetadata};
|
||||
use crate::music_renderer::arylic_tcp::ArylicTcpRenderer;
|
||||
use crate::music_renderer::capabilities::{
|
||||
PlaybackPosition, PlaybackPositionInfo, PlaybackStatus, QueueTransportControl, RendererBackend,
|
||||
TransportControl, VolumeControl,
|
||||
};
|
||||
use crate::music_renderer::chromecast_renderer::ChromecastRenderer;
|
||||
use crate::music_renderer::linkplay_renderer::LinkPlayRenderer;
|
||||
use crate::music_renderer::openhome_renderer::OpenHomeRenderer;
|
||||
@@ -33,6 +29,10 @@ use crate::music_renderer::watcher::{
|
||||
WatchedState,
|
||||
};
|
||||
use crate::music_renderer::RendererFromMediaRendererInfo;
|
||||
use crate::music_renderer::{
|
||||
HasContinuousStream, HasQueue, PlaybackPosition, PlaybackPositionInfo, PlaybackStatus,
|
||||
QueueTransportControl, TransportControl, VolumeControl,
|
||||
};
|
||||
use crate::online::DeviceConnectionState;
|
||||
use crate::queue::{EnqueueMode, MusicQueue, PlaybackItem, QueueBackend, QueueSnapshot};
|
||||
use crate::{DeviceId, DeviceIdentity, DeviceOnline};
|
||||
@@ -46,9 +46,9 @@ use tracing::warn;
|
||||
#[derive(Clone, Debug)]
|
||||
pub struct PlaylistBinding {
|
||||
/// MediaServer that owns the playlist container.
|
||||
pub server_id: DeviceId,
|
||||
pub(crate) server_id: DeviceId,
|
||||
/// DIDL-Lite object id of the playlist container.
|
||||
pub container_id: String,
|
||||
pub(crate) container_id: String,
|
||||
/// True once at least one ContainerUpdateIDs notification has been seen.
|
||||
pub(crate) has_seen_update: bool,
|
||||
/// Flag used internally to signal that the queue should be refreshed
|
||||
@@ -58,6 +58,18 @@ pub struct PlaylistBinding {
|
||||
pub(crate) auto_play_on_refresh: bool,
|
||||
}
|
||||
|
||||
impl PlaylistBinding {
|
||||
/// Returns the ID of the MediaServer that owns this playlist container.
|
||||
pub fn server_id(&self) -> &DeviceId {
|
||||
&self.server_id
|
||||
}
|
||||
|
||||
/// Returns the DIDL-Lite object ID of the playlist container.
|
||||
pub fn container_id(&self) -> &str {
|
||||
&self.container_id
|
||||
}
|
||||
}
|
||||
|
||||
/// Backend-agnostic façade exposing transport, volume, and status contracts.
|
||||
#[derive(Clone, Debug)]
|
||||
pub enum MusicRendererBackend {
|
||||
@@ -311,6 +323,14 @@ impl MusicRenderer {
|
||||
}
|
||||
|
||||
/// Main loop for the watcher thread.
|
||||
///
|
||||
/// # Error policy
|
||||
///
|
||||
/// The watcher thread runs for the lifetime of the renderer and never restarts
|
||||
/// automatically. Network/device errors during polling are logged at `debug` level
|
||||
/// and ignored — they are transient and expected when a device is temporarily
|
||||
/// unreachable. Panics inside `poll_and_emit_changes` are caught and logged at
|
||||
/// `error` level so they do not kill the watcher thread.
|
||||
fn watcher_loop(&self, strategy: WatchStrategy, stop_flag: Arc<AtomicBool>) {
|
||||
let Some(base_interval) = strategy.polling_interval() else {
|
||||
// Pure push strategy - no polling needed (future implementation)
|
||||
@@ -341,7 +361,16 @@ impl MusicRenderer {
|
||||
};
|
||||
|
||||
if self.is_online() {
|
||||
// Wrap in catch_unwind so a panic in poll logic does not terminate the watcher.
|
||||
let poll_result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
|
||||
self.poll_and_emit_changes(tick);
|
||||
}));
|
||||
if let Err(_panic) = poll_result {
|
||||
error!(
|
||||
renderer = self.info.friendly_name(),
|
||||
"poll_and_emit_changes panicked; watcher continues"
|
||||
);
|
||||
}
|
||||
last_activity_time = SystemTime::now();
|
||||
}
|
||||
|
||||
@@ -393,14 +422,19 @@ impl MusicRenderer {
|
||||
|
||||
// Step 2: Do all network calls WITHOUT holding any locks
|
||||
// This prevents blocking other threads that need to read watched_state
|
||||
let position = self.playback_position().ok();
|
||||
let raw_state = self.playback_state().ok();
|
||||
// Errors are logged at trace level — device temporarily unreachable is expected.
|
||||
let position = self.playback_position()
|
||||
.inspect_err(|e| tracing::trace!(renderer = self.info.friendly_name(), error = %e, "playback_position failed"))
|
||||
.ok();
|
||||
let raw_state = self.playback_state()
|
||||
.inspect_err(|e| tracing::trace!(renderer = self.info.friendly_name(), error = %e, "playback_state failed"))
|
||||
.ok();
|
||||
|
||||
// Poll volume and mute every other tick (1 second at 500ms interval)
|
||||
let (volume, mute, is_stream) = if tick % 2 == 0 {
|
||||
(
|
||||
self.volume().ok(),
|
||||
self.mute().ok(),
|
||||
self.volume().inspect_err(|e| tracing::trace!(renderer = self.info.friendly_name(), error = %e, "volume poll failed")).ok(),
|
||||
self.mute().inspect_err(|e| tracing::trace!(renderer = self.info.friendly_name(), error = %e, "mute poll failed")).ok(),
|
||||
Some(self.is_playing_a_stream()),
|
||||
)
|
||||
} else {
|
||||
@@ -490,30 +524,11 @@ impl MusicRenderer {
|
||||
if let Some(stream_flag) = is_stream {
|
||||
if stream_flag {
|
||||
if let Some(ref new_duration) = position.track_duration {
|
||||
let mut state = self.state.lock().unwrap();
|
||||
|
||||
// Parse durations to compare (HH:MM:SS format)
|
||||
let parse_duration = |dur_str: &str| -> Option<u32> {
|
||||
let parts: Vec<&str> = dur_str.split(':').collect();
|
||||
if parts.len() == 3 {
|
||||
let h: u32 = parts[0].parse().ok()?;
|
||||
let m: u32 = parts[1].parse().ok()?;
|
||||
let s: u32 = parts[2].parse().ok()?;
|
||||
Some(h * 3600 + m * 60 + s)
|
||||
} else {
|
||||
None
|
||||
}
|
||||
};
|
||||
let mut state = self.state.lock().expect("RendererState mutex poisoned");
|
||||
|
||||
match &state.current_track_duration {
|
||||
Some(stored_duration) => {
|
||||
// Compare new duration with stored one
|
||||
if let (Some(stored_secs), Some(new_secs)) = (
|
||||
parse_duration(stored_duration),
|
||||
parse_duration(new_duration),
|
||||
) {
|
||||
if new_secs > stored_secs {
|
||||
// Duration increased: update stored value and use new one
|
||||
if crate::queue::stream_duration_increased(stored_duration, new_duration) {
|
||||
tracing::debug!(
|
||||
"MusicRenderer [{}]: Stream duration increased: {} -> {}",
|
||||
self.info.friendly_name(),
|
||||
@@ -521,12 +536,11 @@ impl MusicRenderer {
|
||||
new_duration
|
||||
);
|
||||
state.current_track_duration = Some(new_duration.clone());
|
||||
} else {
|
||||
// Duration decreased or equal: keep stored value
|
||||
} else if crate::queue::stream_duration_decreased(stored_duration, new_duration) {
|
||||
// Duration decreased: keep stored value
|
||||
position.track_duration = Some(stored_duration.clone());
|
||||
}
|
||||
}
|
||||
}
|
||||
None => {
|
||||
// First time: store the duration
|
||||
state.current_track_duration = Some(new_duration.clone());
|
||||
@@ -667,7 +681,7 @@ impl MusicRenderer {
|
||||
match state {
|
||||
PlaybackState::Stopped => {
|
||||
{
|
||||
let s = self.state.lock().unwrap();
|
||||
let s = self.state.lock().expect("RendererState mutex poisoned");
|
||||
tracing::debug!(
|
||||
renderer = self.info.friendly_name(),
|
||||
has_played = s.has_played_since_track_start,
|
||||
@@ -700,7 +714,7 @@ impl MusicRenderer {
|
||||
// On autorise l'auto-avance si:
|
||||
// 1. On a bien vu PLAYING OU
|
||||
// 2. Le titre a été lancé depuis plus de 20 secondes
|
||||
let track_start = self.state.lock().unwrap().track_start_time;
|
||||
let track_start = self.state.lock().expect("RendererState mutex poisoned").track_start_time;
|
||||
let elapsed = track_start
|
||||
.and_then(|t| t.elapsed().ok())
|
||||
.unwrap_or_default();
|
||||
@@ -757,7 +771,7 @@ impl MusicRenderer {
|
||||
PlaybackState::NoMedia => {
|
||||
// Handle end of track (Chromecast returns NoMedia when track ends)
|
||||
// This is equivalent to Stopped for auto-advance purposes
|
||||
let s = self.state.lock().unwrap();
|
||||
let s = self.state.lock().expect("RendererState mutex poisoned");
|
||||
let playback_source = s.playback_source;
|
||||
let has_played = s.has_played_since_track_start;
|
||||
let user_stop = s.user_stop_requested;
|
||||
@@ -824,7 +838,7 @@ impl MusicRenderer {
|
||||
self.set_has_played_flag();
|
||||
}
|
||||
PlaybackState::Transitioning => {
|
||||
let s = self.state.lock().unwrap();
|
||||
let s = self.state.lock().expect("RendererState mutex poisoned");
|
||||
tracing::trace!(
|
||||
renderer = self.info.friendly_name(),
|
||||
has_played = s.has_played_since_track_start,
|
||||
@@ -998,7 +1012,7 @@ impl MusicRenderer {
|
||||
/// Get a clone of the queue Arc (for async sync operations).
|
||||
pub fn queue(&self) -> Arc<Mutex<MusicQueue>> {
|
||||
let backend = self.lock_backend_for("queue");
|
||||
crate::music_renderer::capabilities::RendererBackend::queue(&*backend).clone()
|
||||
crate::music_renderer::capabilities::HasQueue::queue(&*backend).clone()
|
||||
}
|
||||
|
||||
/// Get the current queue item without advancing.
|
||||
@@ -1666,13 +1680,13 @@ impl MusicRenderer {
|
||||
|
||||
/// Gets the last known track metadata.
|
||||
pub fn last_metadata(&self) -> Option<TrackMetadata> {
|
||||
self.state.lock().unwrap().last_metadata.clone()
|
||||
self.state.lock().expect("RendererState mutex poisoned").last_metadata.clone()
|
||||
}
|
||||
|
||||
/// Sets the last known track metadata.
|
||||
/// Updates track_start_time and resets current_track_duration only if the metadata actually changes.
|
||||
pub fn set_last_metadata(&self, metadata: Option<TrackMetadata>) {
|
||||
let mut state = self.state.lock().unwrap();
|
||||
let mut state = self.state.lock().expect("RendererState mutex poisoned");
|
||||
let metadata_changed = state.last_metadata != metadata;
|
||||
if metadata_changed {
|
||||
// Pour les flux continus: utiliser dc:date comme track_start_time réel de diffusion.
|
||||
@@ -1693,23 +1707,23 @@ impl MusicRenderer {
|
||||
|
||||
/// Gets the timestamp when the current track started playing.
|
||||
pub fn track_start_time(&self) -> Option<SystemTime> {
|
||||
self.state.lock().unwrap().track_start_time
|
||||
self.state.lock().expect("RendererState mutex poisoned").track_start_time
|
||||
}
|
||||
|
||||
/// Gets the current playback source.
|
||||
pub fn playback_source(&self) -> PlaybackSource {
|
||||
self.state.lock().unwrap().playback_source
|
||||
self.state.lock().expect("RendererState mutex poisoned").playback_source
|
||||
}
|
||||
|
||||
/// Sets the playback source.
|
||||
pub fn set_playback_source(&self, source: PlaybackSource) {
|
||||
self.state.lock().unwrap().playback_source = source;
|
||||
self.state.lock().expect("RendererState mutex poisoned").playback_source = source;
|
||||
}
|
||||
|
||||
/// Checks if currently playing from queue.
|
||||
pub fn is_playing_from_queue(&self) -> bool {
|
||||
matches!(
|
||||
self.state.lock().unwrap().playback_source,
|
||||
self.state.lock().expect("RendererState mutex poisoned").playback_source,
|
||||
PlaybackSource::FromQueue
|
||||
)
|
||||
}
|
||||
@@ -1720,7 +1734,7 @@ impl MusicRenderer {
|
||||
/// Does NOT change None -> External because that would break queue playback
|
||||
/// (the control_point will set it to FromQueue after play_from_queue succeeds).
|
||||
pub fn mark_external_if_idle(&self) {
|
||||
let mut state = self.state.lock().unwrap();
|
||||
let mut state = self.state.lock().expect("RendererState mutex poisoned");
|
||||
if matches!(state.playback_source, PlaybackSource::External) {
|
||||
// Keep External if we were already playing externally
|
||||
} else {
|
||||
@@ -1731,12 +1745,12 @@ impl MusicRenderer {
|
||||
|
||||
/// Marks that the user requested a stop (to prevent auto-advance).
|
||||
pub fn mark_user_stop_requested(&self) {
|
||||
self.state.lock().unwrap().user_stop_requested = true;
|
||||
self.state.lock().expect("RendererState mutex poisoned").user_stop_requested = true;
|
||||
}
|
||||
|
||||
/// Checks and clears the user stop requested flag.
|
||||
pub fn check_and_clear_user_stop_requested(&self) -> bool {
|
||||
let mut state = self.state.lock().unwrap();
|
||||
let mut state = self.state.lock().expect("RendererState mutex poisoned");
|
||||
let was_requested = state.user_stop_requested;
|
||||
state.user_stop_requested = false;
|
||||
was_requested
|
||||
@@ -1747,21 +1761,21 @@ impl MusicRenderer {
|
||||
/// Sets the has_played_since_track_start flag to true.
|
||||
/// Called when PLAYING state is detected.
|
||||
fn set_has_played_flag(&self) {
|
||||
self.state.lock().unwrap().has_played_since_track_start = true;
|
||||
self.state.lock().expect("RendererState mutex poisoned").has_played_since_track_start = true;
|
||||
}
|
||||
|
||||
/// Clears the has_played_since_track_start flag.
|
||||
/// Called when stopping playback or starting a new track.
|
||||
/// This is public so that ControlPoint can reset it when jumping to a new track.
|
||||
pub fn clear_has_played_flag(&self) {
|
||||
self.state.lock().unwrap().has_played_since_track_start = false;
|
||||
self.state.lock().expect("RendererState mutex poisoned").has_played_since_track_start = false;
|
||||
}
|
||||
|
||||
/// Checks and clears the has_played_since_track_start flag.
|
||||
/// Returns true if PLAYING was seen since last track start, false otherwise.
|
||||
/// Used to determine if auto-advance should be allowed.
|
||||
fn check_and_clear_has_played_flag(&self) -> bool {
|
||||
let mut state = self.state.lock().unwrap();
|
||||
let mut state = self.state.lock().expect("RendererState mutex poisoned");
|
||||
let has_played = state.has_played_since_track_start;
|
||||
state.has_played_since_track_start = false;
|
||||
has_played
|
||||
@@ -1777,7 +1791,7 @@ impl MusicRenderer {
|
||||
/// # Errors
|
||||
/// Returns an error if the duration is invalid (0 or > 7200 seconds).
|
||||
pub fn start_sleep_timer(&self, duration_seconds: u32) -> Result<u32, ControlPointError> {
|
||||
let mut state = self.state.lock().unwrap();
|
||||
let mut state = self.state.lock().expect("RendererState mutex poisoned");
|
||||
state
|
||||
.sleep_timer
|
||||
.start(duration_seconds)
|
||||
@@ -1793,7 +1807,7 @@ impl MusicRenderer {
|
||||
/// # Errors
|
||||
/// Returns an error if the duration is invalid (0 or > 7200 seconds).
|
||||
pub fn update_sleep_timer(&self, duration_seconds: u32) -> Result<u32, ControlPointError> {
|
||||
let mut state = self.state.lock().unwrap();
|
||||
let mut state = self.state.lock().expect("RendererState mutex poisoned");
|
||||
state
|
||||
.sleep_timer
|
||||
.update(duration_seconds)
|
||||
@@ -1804,32 +1818,32 @@ impl MusicRenderer {
|
||||
|
||||
/// Cancels the sleep timer.
|
||||
pub fn cancel_sleep_timer(&self) {
|
||||
self.state.lock().unwrap().sleep_timer.cancel();
|
||||
self.state.lock().expect("RendererState mutex poisoned").sleep_timer.cancel();
|
||||
}
|
||||
|
||||
/// Returns the remaining seconds of the sleep timer, or None if no timer is active.
|
||||
pub fn sleep_timer_remaining(&self) -> Option<u32> {
|
||||
self.state.lock().unwrap().sleep_timer.remaining_seconds()
|
||||
self.state.lock().expect("RendererState mutex poisoned").sleep_timer.remaining_seconds()
|
||||
}
|
||||
|
||||
/// Returns the configured duration of the sleep timer in seconds.
|
||||
pub fn sleep_timer_duration(&self) -> u32 {
|
||||
self.state.lock().unwrap().sleep_timer.duration_seconds()
|
||||
self.state.lock().expect("RendererState mutex poisoned").sleep_timer.duration_seconds()
|
||||
}
|
||||
|
||||
/// Returns true if the sleep timer is active.
|
||||
pub fn is_sleep_timer_active(&self) -> bool {
|
||||
self.state.lock().unwrap().sleep_timer.is_active()
|
||||
self.state.lock().expect("RendererState mutex poisoned").sleep_timer.is_active()
|
||||
}
|
||||
|
||||
/// Returns true if the sleep timer has expired.
|
||||
pub fn is_sleep_timer_expired(&self) -> bool {
|
||||
self.state.lock().unwrap().sleep_timer.is_expired()
|
||||
self.state.lock().expect("RendererState mutex poisoned").sleep_timer.is_expired()
|
||||
}
|
||||
|
||||
/// Gets the sleep timer state as a tuple (is_active, duration_seconds, remaining_seconds).
|
||||
pub fn sleep_timer_state(&self) -> (bool, u32, Option<u32>) {
|
||||
let state = self.state.lock().unwrap();
|
||||
let state = self.state.lock().expect("RendererState mutex poisoned");
|
||||
(
|
||||
state.sleep_timer.is_active(),
|
||||
state.sleep_timer.duration_seconds(),
|
||||
@@ -2091,7 +2105,7 @@ impl RendererFromMediaRendererInfo for MusicRendererBackend {
|
||||
/// Extracts the duration attribute from the <res> element in DIDL metadata.
|
||||
/// This is used as a fallback when the renderer doesn't provide track_duration
|
||||
/// in GetPositionInfo or similar calls.
|
||||
fn parse_didl_duration(didl_xml: &str) -> Option<String> {
|
||||
pub(crate) fn parse_didl_duration(didl_xml: &str) -> Option<String> {
|
||||
// Parse DIDL-Lite XML properly using pmodidl
|
||||
let didl = match DIDLLite::parse(didl_xml) {
|
||||
Ok(d) => d,
|
||||
@@ -2129,6 +2143,37 @@ fn parse_rfc3339_to_system_time(s: &str) -> Option<SystemTime> {
|
||||
Some(std::time::UNIX_EPOCH + std::time::Duration::from_secs(secs as u64))
|
||||
}
|
||||
|
||||
|
||||
/// Dispatch a method call to the inner backend, using the UPnP field for
|
||||
/// HybridUpnpArylic.
|
||||
macro_rules! dispatch_upnp {
|
||||
($self:expr, $method:ident($($arg:expr),*)) => {
|
||||
match $self {
|
||||
MusicRendererBackend::Upnp(r) => r.$method($($arg),*),
|
||||
MusicRendererBackend::OpenHome(r) => r.$method($($arg),*),
|
||||
MusicRendererBackend::LinkPlay(r) => r.$method($($arg),*),
|
||||
MusicRendererBackend::ArylicTcp(r) => r.$method($($arg),*),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.$method($($arg),*),
|
||||
MusicRendererBackend::HybridUpnpArylic { upnp, .. } => upnp.$method($($arg),*),
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
/// Dispatch a method call to the inner backend, using the Arylic field for
|
||||
/// HybridUpnpArylic.
|
||||
macro_rules! dispatch_arylic {
|
||||
($self:expr, $method:ident($($arg:expr),*)) => {
|
||||
match $self {
|
||||
MusicRendererBackend::Upnp(r) => r.$method($($arg),*),
|
||||
MusicRendererBackend::OpenHome(r) => r.$method($($arg),*),
|
||||
MusicRendererBackend::LinkPlay(r) => r.$method($($arg),*),
|
||||
MusicRendererBackend::ArylicTcp(r) => r.$method($($arg),*),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.$method($($arg),*),
|
||||
MusicRendererBackend::HybridUpnpArylic { arylic, .. } => arylic.$method($($arg),*),
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
/// Transport control façade that dispatches to whichever backend can fulfill
|
||||
/// the request, returning a standardized error if the backend lacks support.
|
||||
impl TransportControl for MusicRendererBackend {
|
||||
@@ -2145,38 +2190,9 @@ impl TransportControl for MusicRendererBackend {
|
||||
}
|
||||
}
|
||||
|
||||
fn play(&self) -> Result<(), ControlPointError> {
|
||||
match self {
|
||||
MusicRendererBackend::Upnp(upnp) => upnp.play(),
|
||||
MusicRendererBackend::OpenHome(oh) => oh.play(),
|
||||
MusicRendererBackend::LinkPlay(lp) => lp.play(),
|
||||
MusicRendererBackend::ArylicTcp(ary) => ary.play(),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.play(),
|
||||
MusicRendererBackend::HybridUpnpArylic { arylic, .. } => arylic.play(),
|
||||
}
|
||||
}
|
||||
|
||||
fn pause(&self) -> Result<(), ControlPointError> {
|
||||
match self {
|
||||
MusicRendererBackend::Upnp(upnp) => upnp.pause(),
|
||||
MusicRendererBackend::OpenHome(oh) => oh.pause(),
|
||||
MusicRendererBackend::LinkPlay(lp) => lp.pause(),
|
||||
MusicRendererBackend::ArylicTcp(ary) => ary.pause(),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.pause(),
|
||||
MusicRendererBackend::HybridUpnpArylic { arylic, .. } => arylic.pause(),
|
||||
}
|
||||
}
|
||||
|
||||
fn stop(&self) -> Result<(), ControlPointError> {
|
||||
match self {
|
||||
MusicRendererBackend::Upnp(upnp) => upnp.stop(),
|
||||
MusicRendererBackend::OpenHome(oh) => oh.stop(),
|
||||
MusicRendererBackend::LinkPlay(lp) => lp.stop(),
|
||||
MusicRendererBackend::ArylicTcp(ary) => ary.stop(),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.stop(),
|
||||
MusicRendererBackend::HybridUpnpArylic { arylic, .. } => arylic.stop(),
|
||||
}
|
||||
}
|
||||
fn play(&self) -> Result<(), ControlPointError> { dispatch_arylic!(self, play()) }
|
||||
fn pause(&self) -> Result<(), ControlPointError> { dispatch_arylic!(self, pause()) }
|
||||
fn stop(&self) -> Result<(), ControlPointError> { dispatch_arylic!(self, stop()) }
|
||||
|
||||
fn seek_rel_time(&self, hhmmss: &str) -> Result<(), ControlPointError> {
|
||||
match self {
|
||||
@@ -2197,49 +2213,10 @@ impl TransportControl for MusicRendererBackend {
|
||||
/// Hybrid backends may read via Arylic TCP and write via UPnP, but callers
|
||||
/// always depend on a single [`VolumeControl`] entry point.
|
||||
impl VolumeControl for MusicRendererBackend {
|
||||
fn volume(&self) -> Result<u16, ControlPointError> {
|
||||
match self {
|
||||
MusicRendererBackend::HybridUpnpArylic { arylic, .. } => arylic.volume(),
|
||||
MusicRendererBackend::ArylicTcp(ary) => ary.volume(),
|
||||
MusicRendererBackend::OpenHome(oh) => oh.volume(),
|
||||
MusicRendererBackend::Upnp(upnp) => upnp.volume(),
|
||||
MusicRendererBackend::LinkPlay(lp) => lp.volume(),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.volume(),
|
||||
}
|
||||
}
|
||||
|
||||
fn set_volume(&self, vol: u16) -> Result<(), ControlPointError> {
|
||||
match self {
|
||||
MusicRendererBackend::HybridUpnpArylic { upnp, .. } => upnp.set_volume(vol),
|
||||
MusicRendererBackend::ArylicTcp(ary) => ary.set_volume(vol),
|
||||
MusicRendererBackend::OpenHome(oh) => oh.set_volume(vol),
|
||||
MusicRendererBackend::Upnp(upnp) => upnp.set_volume(vol),
|
||||
MusicRendererBackend::LinkPlay(lp) => lp.set_volume(vol),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.set_volume(vol),
|
||||
}
|
||||
}
|
||||
|
||||
fn mute(&self) -> Result<bool, ControlPointError> {
|
||||
match self {
|
||||
MusicRendererBackend::HybridUpnpArylic { arylic, .. } => arylic.mute(),
|
||||
MusicRendererBackend::OpenHome(r) => r.mute(),
|
||||
MusicRendererBackend::Upnp(r) => r.mute(),
|
||||
MusicRendererBackend::LinkPlay(r) => r.mute(),
|
||||
MusicRendererBackend::ArylicTcp(r) => r.mute(),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.mute(),
|
||||
}
|
||||
}
|
||||
|
||||
fn set_mute(&self, m: bool) -> Result<(), ControlPointError> {
|
||||
match self {
|
||||
MusicRendererBackend::HybridUpnpArylic { arylic, .. } => arylic.set_mute(m),
|
||||
MusicRendererBackend::OpenHome(r) => r.set_mute(m),
|
||||
MusicRendererBackend::Upnp(r) => r.set_mute(m),
|
||||
MusicRendererBackend::LinkPlay(r) => r.set_mute(m),
|
||||
MusicRendererBackend::ArylicTcp(r) => r.set_mute(m),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.set_mute(m),
|
||||
}
|
||||
}
|
||||
fn volume(&self) -> Result<u16, ControlPointError> { dispatch_arylic!(self, volume()) }
|
||||
fn set_volume(&self, vol: u16) -> Result<(), ControlPointError> { dispatch_upnp!(self, set_volume(vol)) }
|
||||
fn mute(&self) -> Result<bool, ControlPointError> { dispatch_arylic!(self, mute()) }
|
||||
fn set_mute(&self, m: bool) -> Result<(), ControlPointError> { dispatch_arylic!(self, set_mute(m)) }
|
||||
}
|
||||
|
||||
/// Playback-state queries sourced from the backend best suited for the job.
|
||||
@@ -2263,234 +2240,28 @@ impl PlaybackStatus for MusicRendererBackend {
|
||||
/// regardless of the backend providing the raw transport data.
|
||||
impl PlaybackPosition for MusicRendererBackend {
|
||||
fn playback_position(&self) -> Result<PlaybackPositionInfo, ControlPointError> {
|
||||
match self {
|
||||
MusicRendererBackend::Upnp(r) => r.playback_position(),
|
||||
MusicRendererBackend::OpenHome(r) => r.playback_position(),
|
||||
MusicRendererBackend::LinkPlay(r) => r.playback_position(),
|
||||
MusicRendererBackend::ArylicTcp(r) => r.playback_position(),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.playback_position(),
|
||||
MusicRendererBackend::HybridUpnpArylic { arylic, .. } => arylic.playback_position(),
|
||||
}
|
||||
dispatch_arylic!(self, playback_position())
|
||||
}
|
||||
}
|
||||
|
||||
impl RendererBackend for MusicRendererBackend {
|
||||
fn queue(&self) -> &Arc<Mutex<MusicQueue>> {
|
||||
match self {
|
||||
MusicRendererBackend::Upnp(r) => r.queue(),
|
||||
MusicRendererBackend::OpenHome(r) => r.queue(),
|
||||
MusicRendererBackend::LinkPlay(r) => r.queue(),
|
||||
MusicRendererBackend::ArylicTcp(r) => r.queue(),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.queue(),
|
||||
MusicRendererBackend::HybridUpnpArylic { upnp, .. } => upnp.queue(),
|
||||
}
|
||||
impl HasQueue for MusicRendererBackend {
|
||||
fn queue(&self) -> &Arc<Mutex<MusicQueue>> { dispatch_upnp!(self, queue()) }
|
||||
}
|
||||
|
||||
impl HasContinuousStream for MusicRendererBackend {
|
||||
fn continuous_stream(&self) -> &Arc<Mutex<bool>> {
|
||||
dispatch_arylic!(self, continuous_stream())
|
||||
}
|
||||
}
|
||||
|
||||
impl QueueTransportControl for MusicRendererBackend {
|
||||
fn play_item(&self, item: &PlaybackItem) -> Result<(), ControlPointError> {
|
||||
dispatch_upnp!(self, play_item(item))
|
||||
}
|
||||
fn play_from_queue(&self) -> Result<(), ControlPointError> {
|
||||
match self {
|
||||
MusicRendererBackend::Upnp(r) => r.play_from_queue(),
|
||||
MusicRendererBackend::OpenHome(r) => r.play_from_queue(),
|
||||
MusicRendererBackend::LinkPlay(r) => r.play_from_queue(),
|
||||
MusicRendererBackend::ArylicTcp(r) => r.play_from_queue(),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.play_from_queue(),
|
||||
MusicRendererBackend::HybridUpnpArylic { upnp, .. } => upnp.play_from_queue(),
|
||||
dispatch_upnp!(self, play_from_queue())
|
||||
}
|
||||
}
|
||||
|
||||
fn play_next(&self) -> Result<(), ControlPointError> {
|
||||
match self {
|
||||
MusicRendererBackend::Upnp(r) => r.play_next(),
|
||||
MusicRendererBackend::OpenHome(r) => r.play_next(),
|
||||
MusicRendererBackend::LinkPlay(r) => r.play_next(),
|
||||
MusicRendererBackend::ArylicTcp(r) => r.play_next(),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.play_next(),
|
||||
MusicRendererBackend::HybridUpnpArylic { upnp, .. } => upnp.play_next(),
|
||||
}
|
||||
}
|
||||
|
||||
fn play_previous(&self) -> Result<(), ControlPointError> {
|
||||
match self {
|
||||
MusicRendererBackend::Upnp(r) => r.play_previous(),
|
||||
MusicRendererBackend::OpenHome(r) => r.play_previous(),
|
||||
MusicRendererBackend::LinkPlay(r) => r.play_previous(),
|
||||
MusicRendererBackend::ArylicTcp(r) => r.play_previous(),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.play_previous(),
|
||||
MusicRendererBackend::HybridUpnpArylic { upnp, .. } => upnp.play_previous(),
|
||||
}
|
||||
}
|
||||
|
||||
fn play_from_index(&self, index: usize) -> Result<(), ControlPointError> {
|
||||
match self {
|
||||
MusicRendererBackend::Upnp(r) => r.play_from_index(index),
|
||||
MusicRendererBackend::OpenHome(r) => r.play_from_index(index),
|
||||
MusicRendererBackend::LinkPlay(r) => r.play_from_index(index),
|
||||
MusicRendererBackend::ArylicTcp(r) => r.play_from_index(index),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.play_from_index(index),
|
||||
MusicRendererBackend::HybridUpnpArylic { upnp, .. } => upnp.play_from_index(index),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl QueueBackend for MusicRendererBackend {
|
||||
fn len(&self) -> Result<usize, ControlPointError> {
|
||||
match self {
|
||||
MusicRendererBackend::Upnp(r) => r.len(),
|
||||
MusicRendererBackend::OpenHome(r) => r.len(),
|
||||
MusicRendererBackend::LinkPlay(r) => r.len(),
|
||||
MusicRendererBackend::ArylicTcp(r) => r.len(),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.len(),
|
||||
MusicRendererBackend::HybridUpnpArylic { upnp, .. } => upnp.len(),
|
||||
}
|
||||
}
|
||||
|
||||
fn track_ids(&self) -> Result<Vec<u32>, ControlPointError> {
|
||||
match self {
|
||||
MusicRendererBackend::Upnp(r) => r.track_ids(),
|
||||
MusicRendererBackend::OpenHome(r) => r.track_ids(),
|
||||
MusicRendererBackend::LinkPlay(r) => r.track_ids(),
|
||||
MusicRendererBackend::ArylicTcp(r) => r.track_ids(),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.track_ids(),
|
||||
MusicRendererBackend::HybridUpnpArylic { upnp, .. } => upnp.track_ids(),
|
||||
}
|
||||
}
|
||||
|
||||
fn id_to_position(&self, id: u32) -> Result<usize, ControlPointError> {
|
||||
match self {
|
||||
MusicRendererBackend::Upnp(r) => r.id_to_position(id),
|
||||
MusicRendererBackend::OpenHome(r) => r.id_to_position(id),
|
||||
MusicRendererBackend::LinkPlay(r) => r.id_to_position(id),
|
||||
MusicRendererBackend::ArylicTcp(r) => r.id_to_position(id),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.id_to_position(id),
|
||||
MusicRendererBackend::HybridUpnpArylic { upnp, .. } => upnp.id_to_position(id),
|
||||
}
|
||||
}
|
||||
|
||||
fn position_to_id(&self, id: usize) -> Result<u32, ControlPointError> {
|
||||
match self {
|
||||
MusicRendererBackend::Upnp(r) => r.position_to_id(id),
|
||||
MusicRendererBackend::OpenHome(r) => r.position_to_id(id),
|
||||
MusicRendererBackend::LinkPlay(r) => r.position_to_id(id),
|
||||
MusicRendererBackend::ArylicTcp(r) => r.position_to_id(id),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.position_to_id(id),
|
||||
MusicRendererBackend::HybridUpnpArylic { upnp, .. } => upnp.position_to_id(id),
|
||||
}
|
||||
}
|
||||
|
||||
fn current_track(&self) -> Result<Option<u32>, ControlPointError> {
|
||||
match self {
|
||||
MusicRendererBackend::Upnp(r) => r.current_track(),
|
||||
MusicRendererBackend::OpenHome(r) => r.current_track(),
|
||||
MusicRendererBackend::LinkPlay(r) => r.current_track(),
|
||||
MusicRendererBackend::ArylicTcp(r) => r.current_track(),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.current_track(),
|
||||
MusicRendererBackend::HybridUpnpArylic { upnp, .. } => upnp.current_track(),
|
||||
}
|
||||
}
|
||||
|
||||
fn current_index(&self) -> Result<Option<usize>, ControlPointError> {
|
||||
match self {
|
||||
MusicRendererBackend::Upnp(r) => r.current_index(),
|
||||
MusicRendererBackend::OpenHome(r) => r.current_index(),
|
||||
MusicRendererBackend::LinkPlay(r) => r.current_index(),
|
||||
MusicRendererBackend::ArylicTcp(r) => r.current_index(),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.current_index(),
|
||||
MusicRendererBackend::HybridUpnpArylic { upnp, .. } => upnp.current_index(),
|
||||
}
|
||||
}
|
||||
|
||||
fn queue_snapshot(&self) -> Result<QueueSnapshot, ControlPointError> {
|
||||
match self {
|
||||
MusicRendererBackend::Upnp(r) => r.queue_snapshot(),
|
||||
MusicRendererBackend::OpenHome(r) => r.queue_snapshot(),
|
||||
MusicRendererBackend::LinkPlay(r) => r.queue_snapshot(),
|
||||
MusicRendererBackend::ArylicTcp(r) => r.queue_snapshot(),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.queue_snapshot(),
|
||||
MusicRendererBackend::HybridUpnpArylic { upnp, .. } => upnp.queue_snapshot(),
|
||||
}
|
||||
}
|
||||
|
||||
fn set_index(&mut self, index: Option<usize>) -> Result<(), ControlPointError> {
|
||||
match self {
|
||||
MusicRendererBackend::Upnp(r) => r.set_index(index),
|
||||
MusicRendererBackend::OpenHome(r) => r.set_index(index),
|
||||
MusicRendererBackend::LinkPlay(r) => r.set_index(index),
|
||||
MusicRendererBackend::ArylicTcp(r) => r.set_index(index),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.set_index(index),
|
||||
MusicRendererBackend::HybridUpnpArylic { upnp, .. } => upnp.set_index(index),
|
||||
}
|
||||
}
|
||||
|
||||
fn replace_queue(
|
||||
&mut self,
|
||||
items: Vec<PlaybackItem>,
|
||||
current_index: Option<usize>,
|
||||
) -> Result<(), ControlPointError> {
|
||||
match self {
|
||||
MusicRendererBackend::Upnp(r) => r.replace_queue(items, current_index),
|
||||
MusicRendererBackend::OpenHome(r) => r.replace_queue(items, current_index),
|
||||
MusicRendererBackend::LinkPlay(r) => r.replace_queue(items, current_index),
|
||||
MusicRendererBackend::ArylicTcp(r) => r.replace_queue(items, current_index),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.replace_queue(items, current_index),
|
||||
MusicRendererBackend::HybridUpnpArylic { upnp, .. } => {
|
||||
upnp.replace_queue(items, current_index)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn sync_queue(
|
||||
&mut self,
|
||||
items: Vec<PlaybackItem>,
|
||||
cancel_token: &Arc<AtomicBool>,
|
||||
on_ready: Option<Box<dyn FnOnce() + Send>>,
|
||||
) -> Result<(), ControlPointError> {
|
||||
match self {
|
||||
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)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn get_item(&self, index: usize) -> Result<Option<PlaybackItem>, ControlPointError> {
|
||||
match self {
|
||||
MusicRendererBackend::Upnp(r) => r.get_item(index),
|
||||
MusicRendererBackend::OpenHome(r) => r.get_item(index),
|
||||
MusicRendererBackend::LinkPlay(r) => r.get_item(index),
|
||||
MusicRendererBackend::ArylicTcp(r) => r.get_item(index),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.get_item(index),
|
||||
MusicRendererBackend::HybridUpnpArylic { upnp, .. } => upnp.get_item(index),
|
||||
}
|
||||
}
|
||||
|
||||
fn replace_item(&mut self, index: usize, item: PlaybackItem) -> Result<(), ControlPointError> {
|
||||
match self {
|
||||
MusicRendererBackend::Upnp(r) => r.replace_item(index, item),
|
||||
MusicRendererBackend::OpenHome(r) => r.replace_item(index, item),
|
||||
MusicRendererBackend::LinkPlay(r) => r.replace_item(index, item),
|
||||
MusicRendererBackend::ArylicTcp(r) => r.replace_item(index, item),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.replace_item(index, item),
|
||||
MusicRendererBackend::HybridUpnpArylic { upnp, .. } => upnp.replace_item(index, item),
|
||||
}
|
||||
}
|
||||
|
||||
fn enqueue_items(
|
||||
&mut self,
|
||||
items: Vec<PlaybackItem>,
|
||||
mode: EnqueueMode,
|
||||
) -> Result<(), ControlPointError> {
|
||||
match self {
|
||||
MusicRendererBackend::Upnp(r) => r.enqueue_items(items, mode),
|
||||
MusicRendererBackend::OpenHome(r) => r.enqueue_items(items, mode),
|
||||
MusicRendererBackend::LinkPlay(r) => r.enqueue_items(items, mode),
|
||||
MusicRendererBackend::ArylicTcp(r) => r.enqueue_items(items, mode),
|
||||
MusicRendererBackend::Chromecast(cc) => cc.enqueue_items(items, mode),
|
||||
MusicRendererBackend::HybridUpnpArylic { upnp, .. } => upnp.enqueue_items(items, mode),
|
||||
}
|
||||
dispatch_upnp!(self, play_from_index(index))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2,8 +2,8 @@ use std::sync::{atomic::AtomicBool, Arc, Mutex};
|
||||
use std::time::SystemTime;
|
||||
|
||||
use crate::music_renderer::capabilities::{
|
||||
PlaybackPosition, PlaybackPositionInfo, PlaybackStatus, QueueTransportControl, RendererBackend,
|
||||
TransportControl, VolumeControl,
|
||||
HasContinuousStream, PlaybackPosition, PlaybackPositionInfo, PlaybackStatus,
|
||||
QueueTransportControl, TransportControl, VolumeControl,
|
||||
};
|
||||
use crate::music_renderer::time_utils::{format_hhmmss_u32, parse_time_flexible};
|
||||
use crate::DeviceIdentity;
|
||||
@@ -15,6 +15,7 @@ use crate::music_renderer::openhome::{
|
||||
build_info_client, build_playlist_client, build_product_client, build_radio_client,
|
||||
build_time_client, build_volume_client,
|
||||
};
|
||||
use crate::music_renderer::HasQueue;
|
||||
use crate::music_renderer::RendererFromMediaRendererInfo;
|
||||
use crate::queue::{EnqueueMode, MusicQueue, PlaybackItem, QueueBackend, QueueSnapshot};
|
||||
use crate::upnp_clients::{
|
||||
@@ -88,7 +89,7 @@ impl OpenHomeRenderer {
|
||||
|
||||
/// Returns true if currently playing a continuous stream (radio without duration)
|
||||
pub fn is_continuous_stream(&self) -> bool {
|
||||
*self.continuous_stream.lock().unwrap()
|
||||
*self.continuous_stream.lock().expect("continuous_stream mutex poisoned")
|
||||
}
|
||||
|
||||
pub fn has_playlist(&self) -> bool {
|
||||
@@ -175,14 +176,14 @@ impl OpenHomeRenderer {
|
||||
/// Plus rapide que snapshot_openhome_playlist() pour juste connaître le nombre de pistes.
|
||||
pub(crate) fn openhome_playlist_len(&self) -> Result<usize, ControlPointError> {
|
||||
// Use queue.len() which uses cached track_ids() internally
|
||||
let queue = self.queue.lock().unwrap();
|
||||
let queue = self.queue.lock().expect("queue mutex poisoned");
|
||||
queue.len()
|
||||
}
|
||||
|
||||
/// 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<Vec<u32>, ControlPointError> {
|
||||
let queue = self.queue.lock().unwrap();
|
||||
let queue = self.queue.lock().expect("queue mutex poisoned");
|
||||
if let Some(oh_queue) = queue.as_openhome() {
|
||||
oh_queue.track_ids()
|
||||
} else {
|
||||
@@ -208,7 +209,7 @@ impl OpenHomeRenderer {
|
||||
let insert_after = match after_id {
|
||||
Some(id) => id,
|
||||
None => {
|
||||
let queue = self.queue.lock().unwrap();
|
||||
let queue = self.queue.lock().expect("queue mutex poisoned");
|
||||
if let Some(oh_queue) = queue.as_openhome() {
|
||||
oh_queue
|
||||
.track_ids()?
|
||||
@@ -266,11 +267,6 @@ impl RendererFromMediaRendererInfo for OpenHomeRenderer {
|
||||
}
|
||||
}
|
||||
|
||||
impl RendererBackend for OpenHomeRenderer {
|
||||
fn queue(&self) -> &Arc<Mutex<MusicQueue>> {
|
||||
&self.queue
|
||||
}
|
||||
}
|
||||
|
||||
impl TransportControl for OpenHomeRenderer {
|
||||
fn play_uri(&self, uri: &str, meta: &str) -> Result<(), ControlPointError> {
|
||||
@@ -369,7 +365,7 @@ impl PlaybackPosition for OpenHomeRenderer {
|
||||
|
||||
// Check cache first
|
||||
{
|
||||
let mut cache = self.position_cache.lock().unwrap();
|
||||
let mut cache = self.position_cache.lock().expect("position_cache mutex poisoned");
|
||||
|
||||
// Track calls for warning detection
|
||||
cache.calls_in_last_second.push(now);
|
||||
@@ -410,7 +406,7 @@ impl PlaybackPosition for OpenHomeRenderer {
|
||||
let mut track_metadata_xml = None;
|
||||
|
||||
// Get track ID from queue (uses cached data)
|
||||
let queue_guard_for_id = self.queue.lock().unwrap();
|
||||
let queue_guard_for_id = self.queue.lock().expect("queue mutex poisoned");
|
||||
if let Some(oh_queue) = queue_guard_for_id.as_openhome() {
|
||||
match oh_queue.current_track() {
|
||||
Ok(id_opt) => track_id = id_opt,
|
||||
@@ -423,7 +419,7 @@ impl PlaybackPosition for OpenHomeRenderer {
|
||||
drop(queue_guard_for_id);
|
||||
|
||||
// Use queue API to get current item with cached metadata
|
||||
let mut queue_guard = self.queue.lock().unwrap();
|
||||
let mut queue_guard = self.queue.lock().expect("queue mutex poisoned");
|
||||
if let Ok(Some((current_item, _))) = queue_guard.peek_current() {
|
||||
// Use metadata from queue cache (updated via OpenHome events)
|
||||
track_uri = Some(current_item.uri.clone());
|
||||
@@ -440,7 +436,7 @@ impl PlaybackPosition for OpenHomeRenderer {
|
||||
}
|
||||
|
||||
// Check if the URI has changed to detect track changes
|
||||
let mut cached_uri = self.current_track_uri.lock().unwrap();
|
||||
let mut cached_uri = self.current_track_uri.lock().expect("current_track_uri mutex poisoned");
|
||||
let uri_changed = cached_uri.as_ref() != Some(¤t_item.uri);
|
||||
|
||||
if uri_changed {
|
||||
@@ -452,7 +448,7 @@ impl PlaybackPosition for OpenHomeRenderer {
|
||||
|
||||
// Détecte si la nouvelle URL est un flux continu
|
||||
let is_stream = crate::music_renderer::is_continuous_stream_url(¤t_item.uri);
|
||||
*self.continuous_stream.lock().unwrap() = is_stream;
|
||||
*self.continuous_stream.lock().expect("continuous_stream mutex poisoned") = is_stream;
|
||||
tracing::debug!("OpenHome URI changed, continuous_stream={}", is_stream);
|
||||
|
||||
*cached_uri = Some(current_item.uri.clone());
|
||||
@@ -488,7 +484,7 @@ impl PlaybackPosition for OpenHomeRenderer {
|
||||
|
||||
// Update cache with fresh data
|
||||
{
|
||||
let mut cache = self.position_cache.lock().unwrap();
|
||||
let mut cache = self.position_cache.lock().expect("position_cache mutex poisoned");
|
||||
cache.last_position = Some(position_info.clone());
|
||||
cache.last_update = Some(now);
|
||||
}
|
||||
@@ -496,28 +492,6 @@ impl PlaybackPosition for OpenHomeRenderer {
|
||||
Ok(position_info)
|
||||
}
|
||||
}
|
||||
|
||||
/// Parse duration from DIDL-Lite metadata XML (OpenHome version)
|
||||
#[allow(dead_code)]
|
||||
fn parse_didl_duration_openhome(didl: &str) -> Option<String> {
|
||||
// Search for duration attribute in <res> element
|
||||
let res_start = didl.find("<res ")?;
|
||||
let after_res = &didl[res_start..];
|
||||
let tag_close = after_res.find('>')?;
|
||||
let tag_attrs = &after_res[..tag_close];
|
||||
|
||||
if let Some(duration_start) = tag_attrs.find("duration=\"") {
|
||||
let duration_offset = duration_start + "duration=\"".len();
|
||||
if let Some(duration_end) = tag_attrs[duration_offset..].find('"') {
|
||||
let duration = &tag_attrs[duration_offset..duration_offset + duration_end];
|
||||
return Some(duration.to_string());
|
||||
}
|
||||
}
|
||||
|
||||
tracing::debug!("OpenHome: No duration found in DIDL metadata");
|
||||
None
|
||||
}
|
||||
|
||||
pub(crate) fn map_openhome_state(raw: &str) -> PlaybackState {
|
||||
match raw.trim().to_ascii_uppercase().as_str() {
|
||||
"PLAYING" => PlaybackState::Playing,
|
||||
@@ -529,6 +503,11 @@ pub(crate) fn map_openhome_state(raw: &str) -> PlaybackState {
|
||||
}
|
||||
|
||||
impl QueueTransportControl for OpenHomeRenderer {
|
||||
fn play_item(&self, _item: &PlaybackItem) -> Result<(), ControlPointError> {
|
||||
let playlist = self.playlist_client_for("play_item")?;
|
||||
playlist.play()
|
||||
}
|
||||
|
||||
fn play_from_queue(&self) -> Result<(), ControlPointError> {
|
||||
{
|
||||
let queue = self
|
||||
@@ -553,50 +532,6 @@ impl QueueTransportControl for OpenHomeRenderer {
|
||||
playlist.play()
|
||||
}
|
||||
|
||||
fn play_next(&self) -> Result<(), ControlPointError> {
|
||||
{
|
||||
let mut queue = self
|
||||
.queue
|
||||
.lock()
|
||||
.map_err(|_| ControlPointError::QueueError("Queue mutex poisoned".into()))?;
|
||||
let len = queue.len().unwrap_or(0);
|
||||
let current = queue.current_index().ok().flatten();
|
||||
let current_track_id = queue.current_track().ok().flatten();
|
||||
let all_ids = queue.track_ids().ok().unwrap_or_default();
|
||||
tracing::trace!(
|
||||
queue_len = len,
|
||||
current_index = ?current,
|
||||
current_track_id = ?current_track_id,
|
||||
all_track_ids = ?all_ids,
|
||||
"OpenHome play_next: advancing queue"
|
||||
);
|
||||
if !queue.advance()? {
|
||||
tracing::trace!(
|
||||
queue_len = len,
|
||||
current_index = ?current,
|
||||
"OpenHome play_next: advance() returned false — no next track"
|
||||
);
|
||||
return Err(ControlPointError::QueueError("No next track".into()));
|
||||
}
|
||||
}
|
||||
|
||||
self.play_from_queue()
|
||||
}
|
||||
|
||||
fn play_previous(&self) -> Result<(), ControlPointError> {
|
||||
{
|
||||
let mut queue = self
|
||||
.queue
|
||||
.lock()
|
||||
.map_err(|_| ControlPointError::QueueError("Queue mutex poisoned".into()))?;
|
||||
if !queue.rewind()? {
|
||||
return Err(ControlPointError::QueueError("No previous track".into()));
|
||||
}
|
||||
}
|
||||
|
||||
self.play_from_queue()
|
||||
}
|
||||
|
||||
fn play_from_index(&self, index: usize) -> Result<(), ControlPointError> {
|
||||
// For OpenHome, we need to convert index to track_id
|
||||
let track_id = {
|
||||
@@ -626,90 +561,37 @@ impl QueueTransportControl for OpenHomeRenderer {
|
||||
}
|
||||
}
|
||||
|
||||
impl QueueBackend for OpenHomeRenderer {
|
||||
fn len(&self) -> Result<usize, ControlPointError> {
|
||||
self.queue
|
||||
.lock()
|
||||
.map_err(|_| ControlPointError::QueueError("Queue mutex poisoned".into()))?
|
||||
.len()
|
||||
impl HasQueue for OpenHomeRenderer {
|
||||
fn queue(&self) -> &Arc<Mutex<MusicQueue>> {
|
||||
&self.queue
|
||||
}
|
||||
}
|
||||
|
||||
fn track_ids(&self) -> Result<Vec<u32>, ControlPointError> {
|
||||
self.queue
|
||||
.lock()
|
||||
.map_err(|_| ControlPointError::QueueError("Queue mutex poisoned".into()))?
|
||||
.track_ids()
|
||||
impl HasContinuousStream for OpenHomeRenderer {
|
||||
fn continuous_stream(&self) -> &Arc<Mutex<bool>> {
|
||||
&self.continuous_stream
|
||||
}
|
||||
}
|
||||
|
||||
fn id_to_position(&self, id: u32) -> Result<usize, ControlPointError> {
|
||||
self.queue
|
||||
.lock()
|
||||
.map_err(|_| ControlPointError::QueueError("Queue mutex poisoned".into()))?
|
||||
.id_to_position(id)
|
||||
}
|
||||
|
||||
fn position_to_id(&self, id: usize) -> Result<u32, ControlPointError> {
|
||||
self.queue
|
||||
.lock()
|
||||
.map_err(|_| ControlPointError::QueueError("Queue mutex poisoned".into()))?
|
||||
.position_to_id(id)
|
||||
}
|
||||
|
||||
fn current_track(&self) -> Result<Option<u32>, ControlPointError> {
|
||||
self.queue
|
||||
.lock()
|
||||
.map_err(|_| ControlPointError::QueueError("Queue mutex poisoned".into()))?
|
||||
.current_track()
|
||||
}
|
||||
|
||||
fn current_index(&self) -> Result<Option<usize>, ControlPointError> {
|
||||
self.queue
|
||||
.lock()
|
||||
.map_err(|_| ControlPointError::QueueError("Queue mutex poisoned".into()))?
|
||||
.current_index()
|
||||
}
|
||||
|
||||
fn queue_snapshot(&self) -> Result<QueueSnapshot, ControlPointError> {
|
||||
self.queue
|
||||
.lock()
|
||||
.map_err(|_| ControlPointError::QueueError("Queue mutex poisoned".into()))?
|
||||
.queue_snapshot()
|
||||
}
|
||||
|
||||
fn set_index(&mut self, index: Option<usize>) -> Result<(), ControlPointError> {
|
||||
self.queue
|
||||
.lock()
|
||||
.map_err(|_| ControlPointError::QueueError("Queue mutex poisoned".into()))?
|
||||
.set_index(index)
|
||||
}
|
||||
|
||||
fn replace_queue(
|
||||
impl OpenHomeRenderer {
|
||||
pub fn replace_queue_with_background(
|
||||
&mut self,
|
||||
items: Vec<PlaybackItem>,
|
||||
current_index: Option<usize>,
|
||||
) -> Result<(), ControlPointError> {
|
||||
// ✅ CORRECTION BUG PRODUCTION: On ne charge PAS toutes les métadonnées
|
||||
// dans le thread principal. OpenHome sur 1000 titres inondait la base SQLite
|
||||
// et bloquait TOUS les autres threads (mutex >500ms).
|
||||
//
|
||||
// On fait juste l'insertion minimaliste maintenant. Le préchargement
|
||||
// des métadonnées est délégué à un thread background.
|
||||
self.queue
|
||||
.lock()
|
||||
.map_err(|_| ControlPointError::QueueError("Mutex poisoned".into()))?
|
||||
.replace_queue(items, current_index)?;
|
||||
|
||||
// Background worker: charge les métadonnées petit à petit sans bloquer personne
|
||||
let queue = self.queue.clone();
|
||||
std::thread::spawn(move || {
|
||||
debug!("🔄 OpenHome: préchargement métadonnées queue en background");
|
||||
if let Ok(mut queue) = queue.lock() {
|
||||
// On ne fait que les 10 prochains titres maintenant, le reste on s'en fout
|
||||
if let Ok(Some(idx)) = queue.current_index() {
|
||||
let end = std::cmp::min(idx + 10, queue.len().unwrap_or(0));
|
||||
for i in idx..end {
|
||||
let _ = queue.get_item(i);
|
||||
// Petit délai pour ne pas noyer la base de données
|
||||
std::thread::sleep(std::time::Duration::from_millis(5));
|
||||
}
|
||||
}
|
||||
@@ -719,41 +601,4 @@ impl QueueBackend for OpenHomeRenderer {
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn sync_queue(
|
||||
&mut self,
|
||||
items: Vec<PlaybackItem>,
|
||||
cancel_token: &Arc<AtomicBool>,
|
||||
on_ready: Option<Box<dyn FnOnce() + Send>>,
|
||||
) -> Result<(), ControlPointError> {
|
||||
self.queue
|
||||
.lock()
|
||||
.map_err(|_| ControlPointError::QueueError("Queue mutex poisoned".into()))?
|
||||
.sync_queue(items, cancel_token, on_ready)
|
||||
}
|
||||
|
||||
fn get_item(&self, index: usize) -> Result<Option<PlaybackItem>, ControlPointError> {
|
||||
self.queue
|
||||
.lock()
|
||||
.map_err(|_| ControlPointError::QueueError("Queue mutex poisoned".into()))?
|
||||
.get_item(index)
|
||||
}
|
||||
|
||||
fn replace_item(&mut self, index: usize, item: PlaybackItem) -> Result<(), ControlPointError> {
|
||||
self.queue
|
||||
.lock()
|
||||
.map_err(|_| ControlPointError::QueueError("Queue mutex poisoned".into()))?
|
||||
.replace_item(index, item)
|
||||
}
|
||||
|
||||
fn enqueue_items(
|
||||
&mut self,
|
||||
items: Vec<PlaybackItem>,
|
||||
mode: EnqueueMode,
|
||||
) -> Result<(), ControlPointError> {
|
||||
self.queue
|
||||
.lock()
|
||||
.map_err(|_| ControlPointError::QueueError("Queue mutex poisoned".into()))?
|
||||
.enqueue_items(items, mode)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -53,7 +53,7 @@ pub fn is_continuous_stream_url(url: &str) -> bool {
|
||||
|
||||
// Check cache first
|
||||
{
|
||||
let cache = STREAM_CACHE.lock().unwrap();
|
||||
let cache = STREAM_CACHE.lock().expect("stream cache mutex poisoned");
|
||||
if let Some(&cached_result) = cache.get(url) {
|
||||
trace!("Cache hit for {}: is_stream={}", url, cached_result);
|
||||
return cached_result;
|
||||
@@ -62,7 +62,7 @@ pub fn is_continuous_stream_url(url: &str) -> bool {
|
||||
|
||||
// Check if already being verified
|
||||
{
|
||||
let mut pending = PENDING_CHECKS.lock().unwrap();
|
||||
let mut pending = PENDING_CHECKS.lock().expect("pending checks mutex poisoned");
|
||||
if pending.contains(url) {
|
||||
debug!(
|
||||
"Stream detection already in progress for {}, returning false temporarily",
|
||||
@@ -93,13 +93,13 @@ pub fn is_continuous_stream_url(url: &str) -> bool {
|
||||
|
||||
// Store in cache
|
||||
{
|
||||
let mut cache = STREAM_CACHE.lock().unwrap();
|
||||
let mut cache = STREAM_CACHE.lock().expect("stream cache mutex poisoned");
|
||||
cache.insert(url_owned.clone(), result);
|
||||
}
|
||||
|
||||
// Remove from pending
|
||||
{
|
||||
let mut pending = PENDING_CHECKS.lock().unwrap();
|
||||
let mut pending = PENDING_CHECKS.lock().expect("pending checks mutex poisoned");
|
||||
pending.remove(&url_owned);
|
||||
}
|
||||
|
||||
@@ -185,12 +185,27 @@ fn check_stream_headers(url: &str) -> Result<bool, String> {
|
||||
|
||||
trace!(
|
||||
"Stream detection for {}: content-length={}, chunked={}, streaming_mime={}, is_stream={}",
|
||||
url, has_content_length, is_chunked, is_streaming_mime, is_stream
|
||||
url,
|
||||
has_content_length,
|
||||
is_chunked,
|
||||
is_streaming_mime,
|
||||
is_stream
|
||||
);
|
||||
|
||||
Ok(is_stream)
|
||||
}
|
||||
|
||||
/// Canonical check: returns `true` if this item should be treated as a continuous stream.
|
||||
///
|
||||
/// Checks `metadata.is_continuous_stream` first (already computed at ingest time),
|
||||
/// then falls back to the URL-based HTTP detection.
|
||||
///
|
||||
/// Use this function everywhere transport-layer code needs to decide whether playback is
|
||||
/// a continuous stream (radio) vs bounded media (file/album track).
|
||||
pub fn is_continuous_stream(metadata: Option<&crate::model::TrackMetadata>, uri: &str) -> bool {
|
||||
metadata.map(|m| m.is_continuous_stream).unwrap_or(false) || is_continuous_stream_url(uri)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
@@ -3,10 +3,13 @@ use std::sync::{atomic::AtomicBool, Arc, Mutex};
|
||||
use crate::errors::ControlPointError;
|
||||
use crate::model::PlaybackState;
|
||||
use crate::music_renderer::capabilities::{
|
||||
PlaybackPosition, PlaybackPositionInfo, PlaybackStatus, QueueTransportControl, RendererBackend,
|
||||
TransportControl, VolumeControl,
|
||||
HasContinuousStream, PlaybackPosition, PlaybackPositionInfo, PlaybackStatus,
|
||||
QueueTransportControl, TransportControl, VolumeControl,
|
||||
};
|
||||
use crate::music_renderer::musicrenderer::{build_didl_lite_metadata, MusicRendererBackend};
|
||||
use crate::music_renderer::musicrenderer::{
|
||||
build_didl_lite_metadata, parse_didl_duration, MusicRendererBackend,
|
||||
};
|
||||
use crate::music_renderer::HasQueue;
|
||||
use crate::music_renderer::RendererFromMediaRendererInfo;
|
||||
use crate::queue::{EnqueueMode, MusicQueue, PlaybackItem, QueueBackend, QueueSnapshot};
|
||||
use crate::upnp_clients::{
|
||||
@@ -114,7 +117,7 @@ impl UpnpRenderer {
|
||||
|
||||
/// Returns true if currently playing a continuous stream (radio without duration)
|
||||
pub fn is_continuous_stream(&self) -> bool {
|
||||
*self.continuous_stream.lock().unwrap()
|
||||
*self.continuous_stream.lock().expect("continuous_stream mutex poisoned")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -174,80 +177,27 @@ impl RendererFromMediaRendererInfo for UpnpRenderer {
|
||||
}
|
||||
}
|
||||
|
||||
impl RendererBackend for UpnpRenderer {
|
||||
fn queue(&self) -> &Arc<Mutex<MusicQueue>> {
|
||||
&self.queue
|
||||
}
|
||||
}
|
||||
|
||||
impl QueueTransportControl for UpnpRenderer {
|
||||
fn play_from_queue(&self) -> Result<(), ControlPointError> {
|
||||
let mut queue = self.queue.lock().unwrap();
|
||||
|
||||
// Get or initialize current index
|
||||
let current_index = match queue.current_index()? {
|
||||
Some(idx) => idx,
|
||||
None => {
|
||||
if queue.len()? > 0 {
|
||||
queue.set_index(Some(0))?;
|
||||
0
|
||||
} else {
|
||||
return Err(ControlPointError::QueueError("Queue is empty".into()));
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
// Get the item
|
||||
let item = queue
|
||||
.get_item(current_index)?
|
||||
.ok_or_else(|| ControlPointError::QueueError("Current item not found".into()))?;
|
||||
|
||||
drop(queue);
|
||||
|
||||
// Build metadata - handle optional TrackMetadata
|
||||
fn play_item(&self, item: &PlaybackItem) -> Result<(), ControlPointError> {
|
||||
let metadata = if let Some(ref track_metadata) = item.metadata {
|
||||
build_didl_lite_metadata(track_metadata, &item.uri, &item.protocol_info)
|
||||
} else {
|
||||
// Fallback to minimal DIDL-Lite if no metadata
|
||||
format!(
|
||||
r#"<DIDL-Lite xmlns="urn:schemas-upnp-org:metadata-1-0/DIDL-Lite/"><item id="0" parentID="-1" restricted="1"><res protocolInfo="{}">{}</res></item></DIDL-Lite>"#,
|
||||
item.protocol_info, item.uri
|
||||
)
|
||||
};
|
||||
|
||||
// Détecte si l'URL est un flux continu en interrogeant le serveur HTTP
|
||||
let is_stream = crate::music_renderer::is_continuous_stream_url(&item.uri);
|
||||
*self.continuous_stream.lock().unwrap() = is_stream;
|
||||
|
||||
// Log current queue state for debugging
|
||||
let queue_state = {
|
||||
let queue = self.queue.lock().unwrap();
|
||||
let idx = queue.current_index().unwrap_or(None);
|
||||
let len = queue.len().unwrap_or(0);
|
||||
let uri = item.uri.clone();
|
||||
let title = item.metadata.as_ref().and_then(|m| m.title.clone());
|
||||
(idx, len, uri, title)
|
||||
};
|
||||
tracing::debug!(
|
||||
"UpnpRenderer play_from_queue: index={:?}/{}, uri={}, title={:?}, continuous_stream={}",
|
||||
queue_state.0,
|
||||
queue_state.1,
|
||||
queue_state.2,
|
||||
queue_state.3,
|
||||
is_stream
|
||||
);
|
||||
|
||||
// Parse et cache la durée du DIDL (fallback pour certains amplis)
|
||||
let duration = parse_didl_duration(&metadata);
|
||||
if let Some(ref dur) = duration {
|
||||
tracing::debug!("Caching duration from queue DIDL: {}", dur);
|
||||
*self.cached_duration.lock().unwrap() = Some(dur.clone());
|
||||
*self.cached_duration.lock().expect("cached_duration mutex poisoned") = Some(dur.clone());
|
||||
} else {
|
||||
tracing::debug!("No duration to cache from queue DIDL");
|
||||
*self.cached_duration.lock().unwrap() = None;
|
||||
*self.cached_duration.lock().expect("cached_duration mutex poisoned") = None;
|
||||
}
|
||||
|
||||
// UPNP: SetAVTransportURI + Play
|
||||
let avt = self.avtransport()?;
|
||||
avt.set_av_transport_uri(&item.uri, &metadata)?;
|
||||
avt.play(0, "1")?;
|
||||
@@ -255,160 +205,19 @@ impl QueueTransportControl for UpnpRenderer {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn play_next(&self) -> Result<(), ControlPointError> {
|
||||
let current_idx = {
|
||||
let queue = self.queue.lock().unwrap();
|
||||
let idx = queue.current_index().unwrap_or(None);
|
||||
let len = queue.len().unwrap_or(0);
|
||||
tracing::debug!(
|
||||
current_index = ?idx,
|
||||
queue_len = len,
|
||||
"play_next: attempting to advance"
|
||||
);
|
||||
idx
|
||||
};
|
||||
{
|
||||
let mut queue = self.queue.lock().unwrap();
|
||||
if !queue.advance()? {
|
||||
return Err(ControlPointError::QueueError("No next track".into()));
|
||||
}
|
||||
let new_idx = queue.current_index().unwrap_or(None);
|
||||
tracing::debug!(
|
||||
previous_index = ?current_idx,
|
||||
new_index = ?new_idx,
|
||||
"play_next: advanced"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
self.play_from_queue()
|
||||
}
|
||||
|
||||
fn play_previous(&self) -> Result<(), ControlPointError> {
|
||||
{
|
||||
let mut queue = self.queue.lock().unwrap();
|
||||
if !queue.rewind()? {
|
||||
return Err(ControlPointError::QueueError("No previous track".into()));
|
||||
}
|
||||
}
|
||||
|
||||
self.play_from_queue()
|
||||
}
|
||||
|
||||
fn play_from_index(&self, index: usize) -> Result<(), ControlPointError> {
|
||||
{
|
||||
let mut queue = self.queue.lock().unwrap();
|
||||
queue.set_index(Some(index))?;
|
||||
}
|
||||
// CORRECTIF: Quand on change l'index manuellement (shuffle, sélection d'un titre)
|
||||
// on logue pour être sûr que c'est bien appelé
|
||||
tracing::debug!(index = index, "✅ SHUFFLE / SEEK: play_from_index appelé");
|
||||
|
||||
self.play_from_queue()
|
||||
impl HasQueue for UpnpRenderer {
|
||||
fn queue(&self) -> &Arc<Mutex<MusicQueue>> {
|
||||
&self.queue
|
||||
}
|
||||
}
|
||||
|
||||
impl QueueBackend for UpnpRenderer {
|
||||
fn len(&self) -> Result<usize, ControlPointError> {
|
||||
self.queue.lock().unwrap().len()
|
||||
}
|
||||
|
||||
fn track_ids(&self) -> Result<Vec<u32>, ControlPointError> {
|
||||
self.queue.lock().unwrap().track_ids()
|
||||
}
|
||||
|
||||
fn id_to_position(&self, id: u32) -> Result<usize, ControlPointError> {
|
||||
self.queue.lock().unwrap().id_to_position(id)
|
||||
}
|
||||
|
||||
fn position_to_id(&self, id: usize) -> Result<u32, ControlPointError> {
|
||||
self.queue.lock().unwrap().position_to_id(id)
|
||||
}
|
||||
|
||||
fn current_track(&self) -> Result<Option<u32>, ControlPointError> {
|
||||
self.queue.lock().unwrap().current_track()
|
||||
}
|
||||
|
||||
fn current_index(&self) -> Result<Option<usize>, ControlPointError> {
|
||||
self.queue.lock().unwrap().current_index()
|
||||
}
|
||||
|
||||
fn queue_snapshot(&self) -> Result<QueueSnapshot, ControlPointError> {
|
||||
self.queue.lock().unwrap().queue_snapshot()
|
||||
}
|
||||
|
||||
fn set_index(&mut self, index: Option<usize>) -> Result<(), ControlPointError> {
|
||||
self.queue.lock().unwrap().set_index(index)
|
||||
}
|
||||
|
||||
fn replace_queue(
|
||||
&mut self,
|
||||
items: Vec<PlaybackItem>,
|
||||
current_index: Option<usize>,
|
||||
) -> Result<(), ControlPointError> {
|
||||
self.queue
|
||||
.lock()
|
||||
.unwrap()
|
||||
.replace_queue(items, current_index)
|
||||
}
|
||||
|
||||
fn sync_queue(
|
||||
&mut self,
|
||||
items: Vec<PlaybackItem>,
|
||||
_cancel_token: &Arc<AtomicBool>,
|
||||
on_ready: Option<Box<dyn FnOnce() + Send>>,
|
||||
) -> Result<(), ControlPointError> {
|
||||
self.queue
|
||||
.lock()
|
||||
.unwrap()
|
||||
.sync_queue(items, &Arc::new(AtomicBool::new(false)), on_ready)
|
||||
}
|
||||
|
||||
fn get_item(&self, index: usize) -> Result<Option<PlaybackItem>, ControlPointError> {
|
||||
self.queue.lock().unwrap().get_item(index)
|
||||
}
|
||||
|
||||
fn replace_item(&mut self, index: usize, item: PlaybackItem) -> Result<(), ControlPointError> {
|
||||
self.queue.lock().unwrap().replace_item(index, item)
|
||||
}
|
||||
|
||||
fn enqueue_items(
|
||||
&mut self,
|
||||
items: Vec<PlaybackItem>,
|
||||
mode: EnqueueMode,
|
||||
) -> Result<(), ControlPointError> {
|
||||
self.queue.lock().unwrap().enqueue_items(items, mode)
|
||||
impl HasContinuousStream for UpnpRenderer {
|
||||
fn continuous_stream(&self) -> &Arc<Mutex<bool>> {
|
||||
&self.continuous_stream
|
||||
}
|
||||
}
|
||||
|
||||
/// Parse le DIDL-Lite pour extraire la durée du premier élément <res>
|
||||
fn parse_didl_duration(didl: &str) -> Option<String> {
|
||||
// Recherche de l'élément <res> (avec ou sans espace après)
|
||||
let res_start = didl
|
||||
.find("<res ")
|
||||
.or_else(|| didl.find("<res>"))
|
||||
.or_else(|| didl.find("<res\n"))
|
||||
.or_else(|| didl.find("<res\t"))?;
|
||||
|
||||
let after_res = &didl[res_start..];
|
||||
|
||||
// Recherche de l'attribut duration dans cet élément <res>
|
||||
// Il doit être avant la fermeture du tag (avant '>')
|
||||
if let Some(tag_close) = after_res.find('>') {
|
||||
let tag_attrs = &after_res[..tag_close];
|
||||
|
||||
if let Some(duration_start) = tag_attrs.find("duration=\"") {
|
||||
let duration_offset = duration_start + "duration=\"".len();
|
||||
if let Some(duration_end) = tag_attrs[duration_offset..].find('"') {
|
||||
let duration = &tag_attrs[duration_offset..duration_offset + duration_end];
|
||||
return Some(duration.to_string());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
tracing::warn!("No duration attribute found in DIDL <res> element");
|
||||
None
|
||||
}
|
||||
|
||||
/// Implémentation UPnP AV de `TransportControl` pour [`UpnpRenderer`].
|
||||
///
|
||||
/// Cette impl se base sur AVTransport (InstanceID = 0).
|
||||
@@ -424,7 +233,7 @@ impl TransportControl for UpnpRenderer {
|
||||
|
||||
// Détecte si l'URL est un flux continu en interrogeant le serveur HTTP
|
||||
let is_stream = crate::music_renderer::is_continuous_stream_url(uri);
|
||||
*self.continuous_stream.lock().unwrap() = is_stream;
|
||||
*self.continuous_stream.lock().expect("continuous_stream mutex poisoned") = is_stream;
|
||||
tracing::debug!(
|
||||
"UpnpRenderer play_uri: URI={}, continuous_stream={}",
|
||||
uri,
|
||||
@@ -435,10 +244,10 @@ impl TransportControl for UpnpRenderer {
|
||||
let duration = parse_didl_duration(meta);
|
||||
if let Some(ref dur) = duration {
|
||||
tracing::debug!("Caching duration from DIDL: {}", dur);
|
||||
*self.cached_duration.lock().unwrap() = Some(dur.clone());
|
||||
*self.cached_duration.lock().expect("cached_duration mutex poisoned") = Some(dur.clone());
|
||||
} else {
|
||||
tracing::debug!("No duration to cache from DIDL");
|
||||
*self.cached_duration.lock().unwrap() = None;
|
||||
*self.cached_duration.lock().expect("cached_duration mutex poisoned") = None;
|
||||
}
|
||||
|
||||
let avt = self.avtransport()?;
|
||||
@@ -532,7 +341,7 @@ impl PlaybackPosition for UpnpRenderer {
|
||||
let mut track_metadata_xml = None;
|
||||
let mut track_uri = raw.track_uri.clone();
|
||||
|
||||
let mut queue_guard = self.queue.lock().unwrap();
|
||||
let mut queue_guard = self.queue.lock().expect("queue mutex poisoned");
|
||||
|
||||
// Récupérer l'item courant de la queue
|
||||
// Normalement current_index est toujours Some() si la queue n'est pas vide (règle métier)
|
||||
|
||||
@@ -28,8 +28,85 @@
|
||||
//! - This identity is used by the sync helpers to preserve the current
|
||||
//! track across queue rebuilds when the MediaServer content changes.
|
||||
|
||||
use crate::music_renderer::HasQueue;
|
||||
use crate::queue::MusicQueue;
|
||||
use crate::{errors::ControlPointError, PlaybackItem, QueueSnapshot};
|
||||
use std::sync::{atomic::AtomicBool, Arc};
|
||||
use std::sync::{atomic::AtomicBool, Arc, Mutex};
|
||||
|
||||
/// Blanket implementation of QueueBackend for types that have a queue.
|
||||
/// All methods simply delegate to the underlying MusicQueue.
|
||||
impl<T: HasQueue> QueueBackend for T {
|
||||
fn len(&self) -> Result<usize, ControlPointError> {
|
||||
self.queue().lock().expect("queue mutex poisoned").len()
|
||||
}
|
||||
|
||||
fn track_ids(&self) -> Result<Vec<u32>, ControlPointError> {
|
||||
self.queue().lock().expect("queue mutex poisoned").track_ids()
|
||||
}
|
||||
|
||||
fn id_to_position(&self, id: u32) -> Result<usize, ControlPointError> {
|
||||
self.queue().lock().expect("queue mutex poisoned").id_to_position(id)
|
||||
}
|
||||
|
||||
fn position_to_id(&self, id: usize) -> Result<u32, ControlPointError> {
|
||||
self.queue().lock().expect("queue mutex poisoned").position_to_id(id)
|
||||
}
|
||||
|
||||
fn current_track(&self) -> Result<Option<u32>, ControlPointError> {
|
||||
self.queue().lock().expect("queue mutex poisoned").current_track()
|
||||
}
|
||||
|
||||
fn current_index(&self) -> Result<Option<usize>, ControlPointError> {
|
||||
self.queue().lock().expect("queue mutex poisoned").current_index()
|
||||
}
|
||||
|
||||
fn queue_snapshot(&self) -> Result<QueueSnapshot, ControlPointError> {
|
||||
self.queue().lock().expect("queue mutex poisoned").queue_snapshot()
|
||||
}
|
||||
|
||||
fn set_index(&mut self, index: Option<usize>) -> Result<(), ControlPointError> {
|
||||
self.queue().lock().expect("queue mutex poisoned").set_index(index)
|
||||
}
|
||||
|
||||
fn replace_queue(
|
||||
&mut self,
|
||||
items: Vec<PlaybackItem>,
|
||||
current_index: Option<usize>,
|
||||
) -> Result<(), ControlPointError> {
|
||||
self.queue()
|
||||
.lock()
|
||||
.unwrap()
|
||||
.replace_queue(items, current_index)
|
||||
}
|
||||
|
||||
fn sync_queue(
|
||||
&mut self,
|
||||
items: Vec<PlaybackItem>,
|
||||
cancel_token: &Arc<AtomicBool>,
|
||||
on_ready: Option<Box<dyn FnOnce() + Send>>,
|
||||
) -> Result<(), ControlPointError> {
|
||||
self.queue()
|
||||
.lock()
|
||||
.unwrap()
|
||||
.sync_queue(items, cancel_token, on_ready)
|
||||
}
|
||||
|
||||
fn get_item(&self, index: usize) -> Result<Option<PlaybackItem>, ControlPointError> {
|
||||
self.queue().lock().expect("queue mutex poisoned").get_item(index)
|
||||
}
|
||||
|
||||
fn replace_item(&mut self, index: usize, item: PlaybackItem) -> Result<(), ControlPointError> {
|
||||
self.queue().lock().expect("queue mutex poisoned").replace_item(index, item)
|
||||
}
|
||||
|
||||
fn enqueue_items(
|
||||
&mut self,
|
||||
items: Vec<PlaybackItem>,
|
||||
mode: EnqueueMode,
|
||||
) -> Result<(), ControlPointError> {
|
||||
self.queue().lock().expect("queue mutex poisoned").enqueue_items(items, mode)
|
||||
}
|
||||
}
|
||||
|
||||
/// High-level enqueue mode.
|
||||
///
|
||||
|
||||
@@ -19,7 +19,7 @@ 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 {
|
||||
pub(crate) fn stream_duration_decreased(old_dur: &str, new_dur: &str) -> bool {
|
||||
match (
|
||||
parse_time_flexible(old_dur).ok(),
|
||||
parse_time_flexible(new_dur).ok(),
|
||||
@@ -30,7 +30,7 @@ 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 {
|
||||
pub(crate) fn stream_duration_increased(old_dur: &str, new_dur: &str) -> bool {
|
||||
match (
|
||||
parse_time_flexible(old_dur).ok(),
|
||||
parse_time_flexible(new_dur).ok(),
|
||||
|
||||
@@ -106,7 +106,7 @@ impl MusicQueue {
|
||||
on_complete: Box<dyn Fn(usize) + Send + 'static>,
|
||||
) -> SyncScheduleOutcome {
|
||||
let (sync_in_progress, sync_pending, sync_cancel_token) = {
|
||||
let q = queue_arc.lock().unwrap();
|
||||
let q = queue_arc.lock().expect("MusicQueue mutex poisoned");
|
||||
(
|
||||
Arc::clone(&q.sync_in_progress),
|
||||
Arc::clone(&q.sync_pending),
|
||||
@@ -129,6 +129,14 @@ impl MusicQueue {
|
||||
thread::Builder::new()
|
||||
.name(thread_name)
|
||||
.spawn(move || {
|
||||
// Protocol for the three AtomicBools:
|
||||
// sync_in_progress : set to true before spawn, cleared on Drop via Guard.
|
||||
// sync_pending : set to true by a concurrent caller that arrives while
|
||||
// a sync is already running. The worker re-fetches items
|
||||
// and loops when it detects this flag on exit.
|
||||
// sync_cancel_token: set to true when a new sync request interrupts an
|
||||
// in-progress one. Passed into QueueBackend::sync_queue
|
||||
// so it can abort early.
|
||||
struct Guard(Arc<AtomicBool>);
|
||||
impl Drop for Guard {
|
||||
fn drop(&mut self) {
|
||||
@@ -136,12 +144,44 @@ impl MusicQueue {
|
||||
}
|
||||
}
|
||||
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);
|
||||
|
||||
// Error policy: sync errors are logged by sync_worker_loop (warn level) and
|
||||
// the thread exits normally. No restart — a new sync can be scheduled via
|
||||
// schedule_sync. The Guard Drop clears sync_in_progress unconditionally,
|
||||
// even on panic, keeping the AtomicBool protocol consistent.
|
||||
tracing::debug!(thread = %std::thread::current().name().unwrap_or("?"), "queue-sync thread started");
|
||||
Self::sync_worker_loop(
|
||||
queue_arc,
|
||||
items,
|
||||
pending_items_fn,
|
||||
on_ready,
|
||||
on_complete,
|
||||
sync_pending,
|
||||
sync_cancel_token,
|
||||
);
|
||||
tracing::debug!(thread = %std::thread::current().name().unwrap_or("?"), "queue-sync thread done");
|
||||
})
|
||||
.expect("Failed to spawn queue-sync thread");
|
||||
|
||||
SyncScheduleOutcome::Scheduled
|
||||
}
|
||||
|
||||
/// Inner loop executed by the sync worker thread.
|
||||
///
|
||||
/// Runs at least once with `initial_items`. If a new sync request arrives while the
|
||||
/// loop is running (`sync_pending` becomes true), it re-fetches items via
|
||||
/// `pending_items_fn` and iterates again, allowing the latest playlist state to win.
|
||||
fn sync_worker_loop(
|
||||
queue_arc: Arc<Mutex<MusicQueue>>,
|
||||
initial_items: Vec<PlaybackItem>,
|
||||
pending_items_fn: Box<dyn Fn() -> Result<Vec<PlaybackItem>, ControlPointError> + Send>,
|
||||
initial_on_ready: Option<Box<dyn FnOnce() + Send>>,
|
||||
on_complete: Box<dyn Fn(usize) + Send>,
|
||||
sync_pending: Arc<AtomicBool>,
|
||||
sync_cancel_token: Arc<AtomicBool>,
|
||||
) {
|
||||
let mut current_items = initial_items;
|
||||
let mut current_on_ready = Some(initial_on_ready);
|
||||
let mut on_complete = Some(on_complete);
|
||||
|
||||
loop {
|
||||
sync_pending.store(false, SeqCst);
|
||||
@@ -168,7 +208,7 @@ impl MusicQueue {
|
||||
);
|
||||
|
||||
let result = {
|
||||
let mut q = queue_arc.lock().unwrap();
|
||||
let mut q = queue_arc.lock().expect("MusicQueue mutex poisoned");
|
||||
<MusicQueue as QueueBackend>::sync_queue(
|
||||
&mut q,
|
||||
current_items,
|
||||
@@ -178,8 +218,6 @@ impl MusicQueue {
|
||||
};
|
||||
// Queue lock is released here.
|
||||
// Now safe to call on_ready (which may re-lock the queue).
|
||||
// If on_ready was triggered by the proxy, consume and call it.
|
||||
// If not (cancelled before first insert), keep it to pass to retry.
|
||||
let carry_on_ready = if on_ready_triggered.load(SeqCst) {
|
||||
tracing::debug!(
|
||||
thread = %std::thread::current().name().unwrap_or("?"),
|
||||
@@ -215,7 +253,7 @@ impl MusicQueue {
|
||||
"queue-sync: completed successfully"
|
||||
);
|
||||
if let Some(cb) = on_complete.take() {
|
||||
let queue_len = queue_arc.lock().unwrap().len().unwrap_or(0);
|
||||
let queue_len = queue_arc.lock().expect("MusicQueue mutex poisoned").len().unwrap_or(0);
|
||||
cb(queue_len);
|
||||
}
|
||||
}
|
||||
@@ -236,12 +274,6 @@ impl MusicQueue {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
tracing::debug!(thread = %std::thread::current().name().unwrap_or("?"), "queue-sync thread done");
|
||||
})
|
||||
.expect("Failed to spawn queue-sync thread");
|
||||
|
||||
SyncScheduleOutcome::Scheduled
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -233,16 +233,16 @@ impl OpenHomeQueue {
|
||||
|
||||
/// Invalide les caches track_ids et read_list (après insert/delete sans impact sur la piste courante).
|
||||
fn invalidate_track_caches(&self) {
|
||||
self.track_ids_cache.lock().unwrap().invalidate();
|
||||
self.read_list_cache.lock().unwrap().invalidate();
|
||||
self.track_ids_cache.lock().expect("track_ids_cache mutex poisoned").invalidate();
|
||||
self.read_list_cache.lock().expect("read_list_cache mutex poisoned").invalidate();
|
||||
}
|
||||
|
||||
/// Invalide tous les caches (après delete_all, seek, stop — opérations qui changent la piste courante).
|
||||
fn invalidate_all_caches(&self) {
|
||||
self.track_ids_cache.lock().unwrap().invalidate();
|
||||
self.read_list_cache.lock().unwrap().invalidate();
|
||||
self.current_track_id_cache.lock().unwrap().invalidate();
|
||||
self.uri_by_id.lock().unwrap().clear();
|
||||
self.track_ids_cache.lock().expect("track_ids_cache mutex poisoned").invalidate();
|
||||
self.read_list_cache.lock().expect("read_list_cache mutex poisoned").invalidate();
|
||||
self.current_track_id_cache.lock().expect("current_track_id_cache mutex poisoned").invalidate();
|
||||
self.uri_by_id.lock().expect("uri_by_id mutex poisoned").clear();
|
||||
}
|
||||
|
||||
/// Tries to detect a simple append-only or delete-from-end pattern without ReadList.
|
||||
@@ -390,8 +390,8 @@ impl OpenHomeQueue {
|
||||
new_metadata: Option<crate::model::TrackMetadata>,
|
||||
uri: &str,
|
||||
) {
|
||||
let mut cache = self.metadata_cache.lock().unwrap();
|
||||
let mut uri_cache = self.uri_by_id.lock().unwrap();
|
||||
let mut cache = self.metadata_cache.lock().expect("metadata_cache mutex poisoned");
|
||||
let mut uri_cache = self.uri_by_id.lock().expect("uri_by_id mutex poisoned");
|
||||
|
||||
// Update URI cache
|
||||
if !uri.is_empty() {
|
||||
@@ -481,7 +481,7 @@ impl OpenHomeQueue {
|
||||
// Le cache contient les métadonnées stables mises lors de l'insertion
|
||||
// Les métadonnées de l'entry (venant de ReadList) changent pour les streams
|
||||
let metadata = {
|
||||
let cache = self.metadata_cache.lock().unwrap();
|
||||
let cache = self.metadata_cache.lock().expect("metadata_cache mutex poisoned");
|
||||
if let Some(cached_meta) = cache.get(&entry.id) {
|
||||
// Utiliser les métadonnées stables du cache
|
||||
tracing::trace!(
|
||||
@@ -557,7 +557,7 @@ impl OpenHomeQueue {
|
||||
}
|
||||
if track_id as usize != playing_id {
|
||||
self.playlist_client.delete_id_if_exists(track_id)?;
|
||||
self.metadata_cache.lock().unwrap().remove(&track_id);
|
||||
self.metadata_cache.lock().expect("metadata_cache mutex poisoned").remove(&track_id);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -616,7 +616,7 @@ impl OpenHomeQueue {
|
||||
track_id
|
||||
);
|
||||
self.playlist_client.delete_id_if_exists(track_id)?;
|
||||
self.metadata_cache.lock().unwrap().remove(&track_id);
|
||||
self.metadata_cache.lock().expect("metadata_cache mutex poisoned").remove(&track_id);
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
@@ -858,7 +858,7 @@ impl OpenHomeQueue {
|
||||
"Using delete_all() for complete replacement (safe - no current track or not in new playlist)"
|
||||
);
|
||||
self.playlist_client.delete_all()?;
|
||||
self.metadata_cache.lock().unwrap().clear();
|
||||
self.metadata_cache.lock().expect("metadata_cache mutex poisoned").clear();
|
||||
}
|
||||
} else {
|
||||
for idx in (0..current_track_ids.len()).rev() {
|
||||
@@ -868,7 +868,7 @@ impl OpenHomeQueue {
|
||||
if !keep_current[idx] {
|
||||
let track_id = current_track_ids[idx];
|
||||
self.playlist_client.delete_id_if_exists(track_id)?;
|
||||
self.metadata_cache.lock().unwrap().remove(&track_id);
|
||||
self.metadata_cache.lock().expect("metadata_cache mutex poisoned").remove(&track_id);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1099,7 +1099,7 @@ impl QueueBackend for OpenHomeQueue {
|
||||
self.ensure_playlist_source_selected()?;
|
||||
|
||||
// Lock the cache for the entire operation to prevent race conditions
|
||||
let mut cache = self.track_ids_cache.lock().unwrap();
|
||||
let mut cache = self.track_ids_cache.lock().expect("track_ids_cache mutex poisoned");
|
||||
|
||||
// Check if cache is valid
|
||||
if let Some(cached_ids) = cache.get() {
|
||||
@@ -1147,7 +1147,7 @@ impl QueueBackend for OpenHomeQueue {
|
||||
|
||||
fn current_track(&self) -> Result<Option<u32>, ControlPointError> {
|
||||
// Hold lock during entire operation to prevent race conditions
|
||||
let mut cache = self.current_track_id_cache.lock().unwrap();
|
||||
let mut cache = self.current_track_id_cache.lock().expect("current_track_id_cache mutex poisoned");
|
||||
|
||||
// Return cached value if valid
|
||||
if let Some(cached_id) = cache.get() {
|
||||
@@ -1192,7 +1192,7 @@ impl QueueBackend for OpenHomeQueue {
|
||||
const MAX_BATCH: usize = 256;
|
||||
let mut entries = Vec::with_capacity(ids.len());
|
||||
for chunk in ids.chunks(MAX_BATCH) {
|
||||
if let Some(cached) = self.read_list_cache.lock().unwrap().get(chunk) {
|
||||
if let Some(cached) = self.read_list_cache.lock().expect("read_list_cache mutex poisoned").get(chunk) {
|
||||
trace!(
|
||||
renderer = self.renderer_id.0.as_str(),
|
||||
"ReadList cache hit for {} IDs",
|
||||
@@ -1283,7 +1283,7 @@ impl QueueBackend for OpenHomeQueue {
|
||||
|
||||
self.ensure_playlist_source_selected()?;
|
||||
self.playlist_client.delete_all()?;
|
||||
self.metadata_cache.lock().unwrap().clear();
|
||||
self.metadata_cache.lock().expect("metadata_cache mutex poisoned").clear();
|
||||
|
||||
// Invalidate caches after delete_all (clears queue and current track)
|
||||
self.invalidate_all_caches();
|
||||
@@ -1340,7 +1340,7 @@ impl QueueBackend for OpenHomeQueue {
|
||||
"sync_queue: Empty playlist - clearing queue with delete_all"
|
||||
);
|
||||
self.playlist_client.delete_all()?;
|
||||
self.metadata_cache.lock().unwrap().clear();
|
||||
self.metadata_cache.lock().expect("metadata_cache mutex poisoned").clear();
|
||||
self.invalidate_all_caches();
|
||||
|
||||
let post_current_track = self.playlist_client.id().ok();
|
||||
@@ -1364,7 +1364,7 @@ impl QueueBackend for OpenHomeQueue {
|
||||
.last()
|
||||
.copied()
|
||||
.unwrap_or(OPENHOME_PLAYLIST_HEAD_ID);
|
||||
let mut uri_cache = self.uri_by_id.lock().unwrap();
|
||||
let mut uri_cache = self.uri_by_id.lock().expect("uri_by_id mutex poisoned");
|
||||
for item in &new_items {
|
||||
if cancel_token.load(SeqCst) {
|
||||
return Err(ControlPointError::SyncCancelled);
|
||||
@@ -1583,8 +1583,8 @@ impl QueueBackend for OpenHomeQueue {
|
||||
.insert(before_id, &item.uri, &metadata)?;
|
||||
|
||||
// Mettre à jour le cache avec les nouvelles métadonnées
|
||||
self.metadata_cache.lock().unwrap().remove(&track_id);
|
||||
self.uri_by_id.lock().unwrap().remove(&track_id);
|
||||
self.metadata_cache.lock().expect("metadata_cache mutex poisoned").remove(&track_id);
|
||||
self.uri_by_id.lock().expect("uri_by_id mutex poisoned").remove(&track_id);
|
||||
self.cache_metadata(new_id, item.metadata, &item.uri);
|
||||
|
||||
if ci == Some(index) {
|
||||
|
||||
@@ -817,7 +817,7 @@ impl OhProductClient {
|
||||
|
||||
pub fn source_xml(&self) -> Result<Vec<OhProductSource>> {
|
||||
// Lock the cache for the entire operation to prevent race conditions
|
||||
let mut cache = self.source_xml_cache.lock().unwrap();
|
||||
let mut cache = self.source_xml_cache.lock().expect("source_xml_cache mutex poisoned");
|
||||
|
||||
// Check if cache is valid
|
||||
if let Some(cached_sources) = cache.get() {
|
||||
@@ -841,7 +841,7 @@ impl OhProductClient {
|
||||
|
||||
pub fn source_index(&self) -> Result<u32, ControlPointError> {
|
||||
// Lock the cache for the entire operation to prevent race conditions
|
||||
let mut cache = self.source_index_cache.lock().unwrap();
|
||||
let mut cache = self.source_index_cache.lock().expect("source_index_cache mutex poisoned");
|
||||
|
||||
// Check if cache is valid
|
||||
if let Some(cached_index) = cache.get() {
|
||||
@@ -881,7 +881,7 @@ impl OhProductClient {
|
||||
|
||||
// Invalidate cache after write operation
|
||||
if result.is_ok() {
|
||||
let mut cache = self.source_index_cache.lock().unwrap();
|
||||
let mut cache = self.source_index_cache.lock().expect("source_index_cache mutex poisoned");
|
||||
cache.invalidate();
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user