20 Commits

Author SHA1 Message Date
000d378789 Merge pull request 'push-mvtuzlstzutl' (#35) from push-mvtuzlstzutl into main
All checks were successful
Build and Push Docker Image / build (push) Successful in 29m38s
Reviewed-on: #35
2025-12-28 16:31:41 +01:00
02a44e9e75 Debug playlist lecture 2025-12-28 16:31:02 +01:00
2f77913caa Enrichissement de la source Qobuz 2025-12-28 15:00:06 +01:00
2353cd736d Merge pull request 'Debug dans pmoaudiocache pour les piste non complètement téléchargées' (#34) from push-wzwlxvklsnuw into main
All checks were successful
Build and Push Docker Image / build (push) Successful in 28m25s
Reviewed-on: #34
2025-12-28 14:20:57 +01:00
572f6e4301 Debug dans pmoaudiocache pour les piste non complètement téléchargées 2025-12-28 14:20:14 +01:00
5cb8cc2f31 Merge pull request 'Debug qobuz, pas encore parfait mais presque.' (#33) from push-ryzlqtmkoyxw into main
Some checks failed
Build and Push Docker Image / build (push) Has been cancelled
Reviewed-on: #33
2025-12-28 13:53:53 +01:00
2b69350162 Debug qobuz pas encore parfait mais mieux 2025-12-28 13:52:31 +01:00
0c64ed1a9a Autorise la configuration des noms UPNP. 2025-12-28 12:37:38 +01:00
876b1e0147 update de qobuz pour utiliser les sources lazy 2025-12-28 12:37:10 +01:00
b915481595 Merge pull request 'Fin provisoire du débugage de PMO contrôle' (#32) from push-nlwpnmvzzmvs into main
All checks were successful
Build and Push Docker Image / build (push) Successful in 28m17s
Reviewed-on: #32
2025-12-28 09:45:43 +01:00
9fe4819640 Fin provisoire du débugage de PMO contrôle 2025-12-28 09:36:55 +01:00
7976ce096e Merge pull request 'push-wmwyloupzyyo' (#31) from push-wmwyloupzyyo into main
All checks were successful
Build and Push Docker Image / build (push) Successful in 29m40s
Reviewed-on: #31
2025-12-27 22:58:58 +01:00
0cbd0c9b30 Remise au propre des abstractions de PMOcontrol. 2025-12-27 22:42:13 +01:00
0abe66f7fb On corrige le bug des serveur de musique autre que le PMO serveur. 2025-12-27 22:06:20 +01:00
ed46ee8109 Merge pull request 'push-xwqwmkpwtxqk' (#30) from push-xwqwmkpwtxqk into main
All checks were successful
Build and Push Docker Image / build (push) Successful in 28m32s
Reviewed-on: #30
2025-12-27 18:06:34 +01:00
dec48a3130 Lenteur détection des serveurs. 2025-12-27 18:05:03 +01:00
46625c20ec Petit débugage de la user interface. 2025-12-27 17:02:29 +01:00
125b5bbc2f Simplification du code OpenHome environ tous les workarounds. 2025-12-27 17:02:29 +01:00
0a50ca38ac Merge pull request 'Amélioration de l'interface utilisateur.' (#29) from push-lxpmrvpuzxkw into main
All checks were successful
Build and Push Docker Image / build (push) Successful in 28m27s
Reviewed-on: #29
2025-12-27 15:28:05 +01:00
40cce9f894 Rends les items de la queue de lecture cliquables. 2025-12-27 15:12:57 +01:00
45 changed files with 3575 additions and 775 deletions

3
.claude-env Normal file
View File

@@ -0,0 +1,3 @@
# Configuration PATH pour Claude Code
# Ce fichier sera lu automatiquement pour configurer l'environnement
export PATH="/Users/coissac/mamba/condabin:/opt/homebrew/lib/ruby/gems/3.4.0/bin:/opt/homebrew/opt/ruby/bin:/Users/coissac/go/bin:/Users/coissac/.cargo/bin:/Users/coissac/.modular/pkg/packages.modular.com_mojo/bin:/Applications/quarto/bin:/Users/coissac/.vscode-oss/extensions/vadimcn.vscode-lldb-1.12.0/bin:/Library/Frameworks/Python.framework/Versions/3.12/bin:/opt/homebrew/bin:/opt/homebrew/sbin:/usr/local/bin:/System/Cryptexes/App/usr/bin:/usr/bin:/bin:/usr/sbin:/sbin:/var/run/com.apple.security.cryptexd/codex.system/bootstrap/usr/local/bin:/var/run/com.apple.security.cryptexd/codex.system/bootstrap/usr/bin:/var/run/com.apple.security.cryptexd/codex.system/bootstrap/usr/appleinternal/bin:/opt/pmk/env/global/bin:/opt/X11/bin:/usr/local/dorado/dorado-0.9.5-osx-arm64/bin:/usr/local/go/bin:/usr/local/src/last-main/bin:/Users/coissac/travail/__MOI__/GO/obitools4/build:/opt/podman/bin:/Applications/quarto/bin:/Users/coissac/.cargo/bin:/Users/coissac/.vscode-oss/extensions/vadimcn.vscode-lldb-1.12.0/bin:/Users/coissac/.vscode-oss/extensions/ms-python.debugpy-2025.14.1-darwin-arm64/bundled/scripts/noConfigScripts:/Users/coissac/.orbstack/bin"

31
.claude/CLAUDE.md Normal file
View File

@@ -0,0 +1,31 @@
# PMOMusic Project Configuration
## Version Control
Ce projet utilise **Jujutsu (jj)** pour le contrôle de version, PAS git.
- Utiliser les commandes `jj` au lieu des commandes `git`
- Bookmark principal : `main`
- Ne jamais suggérer de commandes git
## Environnement
Le PATH et les variables d'environnement sont configurés dans `.claude-env` à la racine du projet.
## Configuration de l'application
- Fichier de configuration principal : `.pmomusic/config.yaml`
- Configuration UPNP personnalisable pour différencier les instances en développement
## Développement
Pendant le développement, plusieurs serveurs PMOMusic peuvent tourner en parallèle. Utiliser la configuration UPNP dans `.pmomusic/config.yaml` pour différencier les instances :
```yaml
host:
upnp:
manufacturer: "PMOMusic-Dev1"
udn_prefix: "pmomusic-dev1"
model_name_prefix: "PMOMusic-Dev1"
friendly_name_prefix: "PMOMusic-Dev1"
```
## Architecture
- Projet Rust multi-crates avec workspaces
- Crates principales : pmoupnp, pmomediaserver, pmomediarenderer, pmoconfig
- Pattern d'extension de configuration via traits (voir pmocache/src/config_ext.rs)

1
Cargo.lock generated
View File

@@ -4021,6 +4021,7 @@ dependencies = [
"reqwest",
"serde",
"serde_json",
"serde_yaml",
"socket2 0.5.10",
"thiserror 2.0.17",
"tokio",

View File

@@ -6,10 +6,19 @@ defineProps<{
item: QueueItem
isCurrent: boolean
}>()
const emit = defineEmits<{
click: [item: QueueItem]
}>()
function handleClick(item: QueueItem) {
console.log('[QueueItem] Click detected on item:', item.index, item.title)
emit('click', item)
}
</script>
<template>
<div :class="['queue-item', { current: isCurrent }]">
<div :class="['queue-item', { current: isCurrent }]" @click="handleClick(item)">
<!-- Indicateur piste en cours -->
<div class="current-indicator" v-if="isCurrent">
<Play :size="16" fill="currentColor" />
@@ -51,6 +60,7 @@ defineProps<{
border-radius: var(--radius-md);
transition: background-color var(--transition-fast);
position: relative;
cursor: pointer;
}
.queue-item:hover {

View File

@@ -3,17 +3,26 @@ import { computed, ref, watch, nextTick, toRef } from 'vue'
import { useRenderer } from '@/composables/useRenderers'
import QueueItem from './QueueItem.vue'
import { Link } from 'lucide-vue-next'
import type { QueueItem as QueueItemType } from '@/services/pmocontrol/types'
const props = defineProps<{
rendererId: string
}>()
const emit = defineEmits<{
clickItem: [item: QueueItemType]
}>()
const { queue, binding } = useRenderer(toRef(props, 'rendererId'))
const isAttached = computed(() => !!binding.value)
const queueContainer = ref<HTMLElement | null>(null)
function handleItemClick(item: QueueItemType) {
emit('clickItem', item)
}
// Auto-scroll vers la piste courante lors de l'ouverture
watch(() => queue.value?.current_index, async (currentIndex) => {
if (currentIndex !== null && currentIndex !== undefined && queueContainer.value) {
@@ -53,6 +62,7 @@ watch(() => queue.value?.current_index, async (currentIndex) => {
:key="item.index"
:item="item"
:is-current="item.index === queue.current_index"
@click="handleItemClick"
/>
</div>

View File

@@ -10,6 +10,8 @@ import StatusBadge from '@/components/pmocontrol/StatusBadge.vue'
import { ChevronUp, ChevronDown, Link } from 'lucide-vue-next'
import { useRenderers } from '@/composables/useRenderers'
import { useUIStore } from '@/stores/ui'
import { api } from '@/services/pmocontrol/api'
import type { QueueItem } from '@/services/pmocontrol/types'
const props = defineProps<{
rendererId: string
@@ -48,6 +50,21 @@ async function handleDetachPlaylist() {
uiStore.notifyError(`Erreur: ${error instanceof Error ? error.message : 'Erreur inconnue'}`)
}
}
// Gérer le clic sur un item de la queue
async function handleQueueItemClick(item: QueueItem) {
try {
await api.seekQueueIndex(props.rendererId, item.index)
console.log('[RendererTabContent] Jumped to queue index:', item.index, item.title)
// Force un refetch immédiat pour synchroniser la cover affichée
// sans attendre l'événement SSE qui peut avoir un délai
await refresh(true)
} catch (error) {
console.error('[RendererTabContent] Error seeking to queue index:', error)
uiStore.notifyError(`Erreur: ${error instanceof Error ? error.message : 'Impossible de sauter à cet item'}`)
}
}
</script>
<template>
@@ -100,7 +117,7 @@ async function handleDetachPlaylist() {
<!-- Colonne droite: Queue (desktop et landscape uniquement) -->
<div v-if="!isMobilePortrait" class="queue-column">
<QueueViewer :renderer-id="rendererId" class="queue-viewer" />
<QueueViewer :renderer-id="rendererId" class="queue-viewer" @click-item="handleQueueItemClick" />
</div>
<!-- Drawer queue (mobile portrait uniquement) -->
@@ -114,7 +131,7 @@ async function handleDetachPlaylist() {
<!-- Contenu du drawer -->
<div class="queue-drawer-content">
<QueueViewer :renderer-id="rendererId" />
<QueueViewer :renderer-id="rendererId" @click-item="handleQueueItemClick" />
</div>
</div>

View File

@@ -147,6 +147,17 @@ class PMOControlAPI {
})
}
/**
* Saute à un index spécifique dans la queue
* POST /api/control/renderers/{id}/queue/seek
*/
async seekQueueIndex(id: string, index: number): Promise<SuccessResponse> {
return this.request<SuccessResponse>(`/renderers/${encodeURIComponent(id)}/queue/seek`, {
method: 'POST',
body: JSON.stringify({ index }),
})
}
// ============================================================================
// CONTRÔLE VOLUME
// ============================================================================

View File

@@ -124,6 +124,24 @@ async fn serve_finalized_pk<C: CacheConfig + 'static>(
warn!("Error updating hit count for {}: {}", pk, e);
}
// ROUTE 1 : Fichier en cours de téléchargement
// Utilise le streaming progressif avec Content-Length pour permettre
// la lecture pendant le téléchargement tout en préservant la durée/position
if let Some(download) = cache.get_download(pk).await {
if !download.finished().await {
let response = stream_file_progressive(file_path, download, content_type).await;
if response.status().is_success() {
cache.notify_broadcast(pk, &qualifier).await;
}
return response;
}
}
// ROUTE 2 : Fichier complètement téléchargé
// Utilise l'ancien système éprouvé qui garantit un passage correct
// de toutes les informations (Content-Length automatique, etc.)
let response = serve_complete_file(file_path, content_type).await;
if response.status().is_success() {
@@ -256,6 +274,8 @@ async fn stream_file_progressive(
download: Arc<crate::download::Download>,
content_type: &'static str,
) -> Response {
use axum::http::header;
// Attendre qu'au moins 64 KB soient disponibles avant de commencer
const MIN_SIZE_TO_START: u64 = 64 * 1024;
@@ -284,15 +304,24 @@ async fn stream_file_progressive(
let stream = ReaderStream::new(file);
let body = Body::from_stream(stream);
(
StatusCode::OK,
[
("content-type", content_type),
("transfer-encoding", "chunked"),
],
body,
)
.into_response()
// Récupérer la taille attendue du fichier si disponible
let expected_size = download.expected_size().await;
// Construire la réponse avec Content-Length si connu
let mut response = axum::http::Response::builder()
.status(StatusCode::OK)
.header(header::CONTENT_TYPE, content_type);
if let Some(size) = expected_size {
// Si on connaît la taille finale, l'envoyer au renderer
// pour qu'il puisse calculer la durée et afficher la position
response = response.header(header::CONTENT_LENGTH, size);
} else {
// Sinon, utiliser chunked encoding
response = response.header(header::TRANSFER_ENCODING, "chunked");
}
response.body(body).unwrap()
}
/// Sert un fichier complet déjà téléchargé

View File

@@ -1,5 +1,10 @@
host:
http_port: "8080"
upnp:
manufacturer: "PMOMusic"
udn_prefix: "pmomusic"
model_name_prefix: "PMOMusic"
friendly_name_prefix: "PMOMusic"
cover_cache:
directory: "cache_covers"
size: 2000

View File

@@ -0,0 +1,120 @@
//! Simple Chromecast info retrieval test
//!
//! This is the most minimal test - just connects and gets device status.
//!
//! Usage:
//! cargo run --example chromecast_info -- <chromecast_ip>
use std::env;
use std::sync::Once;
use rust_cast::CastDevice;
const DEFAULT_DESTINATION_ID: &str = "receiver-0";
const DEFAULT_PORT: u16 = 8009;
/// Ensures the Rustls CryptoProvider is initialized exactly once.
fn ensure_crypto_provider_initialized() {
static INIT: Once = Once::new();
INIT.call_once(|| {
let _ = rustls::crypto::CryptoProvider::install_default(
rustls::crypto::aws_lc_rs::default_provider()
);
});
}
fn main() {
// Parse command line arguments
let args: Vec<String> = env::args().collect();
if args.len() < 2 {
eprintln!("Usage: {} <chromecast_ip>", args[0]);
eprintln!("\nExample:");
eprintln!(" {} 192.168.1.100", args[0]);
std::process::exit(1);
}
let chromecast_ip = &args[1];
println!("═══════════════════════════════════════════════════════");
println!(" Chromecast Info Test");
println!("═══════════════════════════════════════════════════════");
println!();
// Initialize crypto provider
ensure_crypto_provider_initialized();
// Connect to the device
println!("→ Connecting to {}:{}...", chromecast_ip, DEFAULT_PORT);
let cast_device = match CastDevice::connect_without_host_verification(chromecast_ip, DEFAULT_PORT) {
Ok(device) => {
println!(" ✓ Connected");
device
}
Err(e) => {
eprintln!(" ✗ Failed: {}", e);
std::process::exit(1);
}
};
// Connect to receiver channel
println!();
println!("→ Connecting to receiver channel...");
if let Err(e) = cast_device.connection.connect(DEFAULT_DESTINATION_ID.to_string()) {
eprintln!(" ✗ Failed: {}", e);
std::process::exit(1);
}
println!(" ✓ Channel connected");
// Send initial ping
println!();
println!("→ Sending initial ping...");
if let Err(e) = cast_device.heartbeat.ping() {
eprintln!(" ✗ Failed: {}", e);
std::process::exit(1);
}
println!(" ✓ Ping sent");
// Get receiver status
println!();
println!("→ Getting receiver status...");
match cast_device.receiver.get_status() {
Ok(status) => {
println!(" ✓ Status retrieved\n");
// Volume
if let Some(level) = status.volume.level {
println!(" Volume: {:.0}%", level * 100.0);
}
if let Some(muted) = status.volume.muted {
println!(" Muted: {}", muted);
}
// Applications
println!("\n Running applications: {}", status.applications.len());
for (i, app) in status.applications.iter().enumerate() {
println!("\n App #{}:", i + 1);
println!(" Display Name: {}", app.display_name);
println!(" App ID: {}", app.app_id);
println!(" Session ID: {}", app.session_id);
println!(" Transport ID: {}", app.transport_id);
println!(" Status: {}", app.status_text);
println!(" Namespaces: {}", app.namespaces.len());
}
if status.applications.is_empty() {
println!(" (no apps currently running)");
}
}
Err(e) => {
eprintln!(" ✗ Failed: {}", e);
std::process::exit(1);
}
}
println!();
println!("═══════════════════════════════════════════════════════");
println!(" Test completed successfully!");
println!("═══════════════════════════════════════════════════════");
}

View File

@@ -0,0 +1,284 @@
//! Simple Chromecast playback test
//!
//! This example tests basic Chromecast functionality by:
//! 1. Connecting to a Chromecast device
//! 2. Launching the DefaultMediaReceiver app
//! 3. Loading and playing a test media URL
//! 4. Maintaining the heartbeat loop
//!
//! Usage:
//! cargo run --example chromecast_playback_test -- <chromecast_ip> [media_url]
//!
//! Example:
//! cargo run --example chromecast_playback_test -- 192.168.1.100
use std::env;
use std::sync::Once;
use rust_cast::{
CastDevice, ChannelMessage,
channels::{
heartbeat::HeartbeatResponse,
media::{Media, StreamType},
receiver::CastDeviceApp,
},
};
const DEFAULT_DESTINATION_ID: &str = "receiver-0";
const DEFAULT_PORT: u16 = 8009;
// Test media URLs (public domain audio files)
const TEST_MEDIA_URL: &str = "https://www.soundhelix.com/examples/mp3/SoundHelix-Song-1.mp3";
/// Ensures the Rustls CryptoProvider is initialized exactly once.
fn ensure_crypto_provider_initialized() {
static INIT: Once = Once::new();
INIT.call_once(|| {
let _ = rustls::crypto::CryptoProvider::install_default(
rustls::crypto::aws_lc_rs::default_provider()
);
println!("✓ Rustls CryptoProvider initialized");
});
}
/// Detects content type from URL path
fn detect_content_type(url: &str) -> String {
// Detect from URL path - check if path contains /flac/, /mp3/, etc.
if url.contains("/flac/") || url.contains(".flac") {
println!(" ✓ Detected FLAC from URL path");
return "audio/flac".to_string();
}
if url.contains("/mp3/") || url.contains(".mp3") {
println!(" ✓ Detected MP3 from URL path");
return "audio/mpeg".to_string();
}
if url.contains(".m4a") || url.contains(".mp4") || url.contains(".aac") {
println!(" ✓ Detected AAC/M4A from URL path");
return "audio/mp4".to_string();
}
if url.contains("/ogg/") || url.contains(".ogg") {
println!(" ✓ Detected OGG from URL path");
return "audio/ogg".to_string();
}
if url.contains(".opus") {
println!(" ✓ Detected Opus from URL path");
return "audio/opus".to_string();
}
if url.contains(".wav") {
println!(" ✓ Detected WAV from URL path");
return "audio/wav".to_string();
}
// Fallback
println!(" ⚠ Could not detect type, using audio/mpeg as fallback");
"audio/mpeg".to_string()
}
fn main() {
// Setup logging
tracing_subscriber::fmt()
.with_max_level(tracing::Level::DEBUG)
.init();
// Parse command line arguments
let args: Vec<String> = env::args().collect();
if args.len() < 2 {
eprintln!("Usage: {} <chromecast_ip> [media_url]", args[0]);
eprintln!("\nExample:");
eprintln!(" {} 192.168.1.100", args[0]);
eprintln!("\nIf no media URL is provided, will use: {}", TEST_MEDIA_URL);
std::process::exit(1);
}
let chromecast_ip = &args[1];
let media_url = if args.len() > 2 {
&args[2]
} else {
TEST_MEDIA_URL
};
println!("╔════════════════════════════════════════════════════════════╗");
println!("║ Chromecast Playback Test (rust_cast) ║");
println!("╚════════════════════════════════════════════════════════════╝");
println!();
println!("Target: {}:{}", chromecast_ip, DEFAULT_PORT);
println!("Media URL: {}", media_url);
println!();
println!("Detecting media content type...");
let media_type = detect_content_type(media_url);
println!();
// Initialize crypto provider
ensure_crypto_provider_initialized();
// Step 1: Connect to the device
println!("──────────────────────────────────────────────────────────");
println!("STEP 1: Connecting to Chromecast...");
let cast_device = match CastDevice::connect_without_host_verification(chromecast_ip, DEFAULT_PORT) {
Ok(device) => {
println!("✓ Connected to Chromecast");
device
}
Err(e) => {
eprintln!("✗ Failed to connect: {}", e);
std::process::exit(1);
}
};
// Step 2: Connect to the default receiver channel
println!();
println!("STEP 2: Connecting to receiver channel...");
if let Err(e) = cast_device.connection.connect(DEFAULT_DESTINATION_ID.to_string()) {
eprintln!("✗ Failed to connect channel: {}", e);
std::process::exit(1);
}
println!("✓ Channel connected");
// Step 3: Send initial ping (CRITICAL per rust_caster.rs)
println!();
println!("STEP 3: Sending initial heartbeat ping...");
if let Err(e) = cast_device.heartbeat.ping() {
eprintln!("✗ Failed to send initial ping: {}", e);
std::process::exit(1);
}
println!("✓ Initial ping sent");
// Step 4: Get receiver status
println!();
println!("STEP 4: Getting receiver status...");
let status = match cast_device.receiver.get_status() {
Ok(status) => {
println!("✓ Receiver status obtained");
println!(" - Volume: {:.0}%", status.volume.level.unwrap_or(0.5) * 100.0);
println!(" - Muted: {}", status.volume.muted.unwrap_or(false));
println!(" - Running apps: {}", status.applications.len());
status
}
Err(e) => {
eprintln!("✗ Failed to get status: {}", e);
std::process::exit(1);
}
};
// Step 5: Launch DefaultMediaReceiver
println!();
println!("STEP 5: Launching DefaultMediaReceiver app...");
let app = match cast_device.receiver.launch_app(&CastDeviceApp::DefaultMediaReceiver) {
Ok(app) => {
println!("✓ App launched successfully");
println!(" - App ID: {}", app.app_id);
println!(" - Display Name: {}", app.display_name);
println!(" - Session ID: {}", app.session_id);
println!(" - Transport ID: {}", app.transport_id);
app
}
Err(e) => {
eprintln!("✗ Failed to launch app: {}", e);
std::process::exit(1);
}
};
// Step 6: Connect to the app's transport
println!();
println!("STEP 6: Connecting to app transport...");
if let Err(e) = cast_device.connection.connect(app.transport_id.as_str()) {
eprintln!("✗ Failed to connect to app transport: {}", e);
std::process::exit(1);
}
println!("✓ Connected to app transport: {}", app.transport_id);
// Step 7: Load the media
println!();
println!("STEP 7: Loading media...");
println!(" Content-Type: {}", media_type);
let media = Media {
content_id: media_url.to_string(),
content_type: media_type,
stream_type: StreamType::Buffered,
duration: None,
metadata: None,
};
match cast_device.media.load(
app.transport_id.as_str(),
app.session_id.as_str(),
&media,
) {
Ok(status) => {
println!("✓ Media loaded successfully!");
println!(" - Media status entries: {}", status.entries.len());
if let Some(entry) = status.entries.first() {
println!(" - Player state: {:?}", entry.player_state);
println!(" - Media session ID: {}", entry.media_session_id);
if let Some(ref media) = entry.media {
println!(" - Content ID: {}", media.content_id);
println!(" - Stream type: {:?}", media.stream_type);
}
}
}
Err(e) => {
eprintln!("✗ Failed to load media: {}", e);
std::process::exit(1);
}
}
// Step 8: Enter heartbeat loop
println!();
println!("──────────────────────────────────────────────────────────");
println!("STEP 8: Entering heartbeat loop (Ctrl+C to exit)...");
println!("──────────────────────────────────────────────────────────");
println!();
let mut heartbeat_count = 0;
let mut media_message_count = 0;
loop {
match cast_device.receive() {
Ok(ChannelMessage::Heartbeat(response)) => {
if let HeartbeatResponse::Ping = response {
heartbeat_count += 1;
println!("[Heartbeat #{:3}] Received Ping, sending Pong...", heartbeat_count);
if let Err(e) = cast_device.heartbeat.pong() {
eprintln!("✗ Failed to send pong: {}", e);
break;
}
} else {
println!("[Heartbeat] {:?}", response);
}
}
Ok(ChannelMessage::Media(response)) => {
media_message_count += 1;
println!("[Media #{:3}] {:?}", media_message_count, response);
}
Ok(ChannelMessage::Receiver(response)) => {
println!("[Receiver] {:?}", response);
}
Ok(ChannelMessage::Connection(response)) => {
println!("[Connection] {:?}", response);
}
Ok(ChannelMessage::Raw(response)) => {
println!("[Raw] Unsupported message type: {:?}", response);
}
Err(e) => {
eprintln!("✗ Error receiving message: {}", e);
break;
}
}
}
println!();
println!("──────────────────────────────────────────────────────────");
println!("Test completed.");
println!("Total heartbeats: {}", heartbeat_count);
println!("Total media messages: {}", media_message_count);
}

View File

@@ -0,0 +1,409 @@
# Analyse : Rendre rust-cast asynchrone vs autres options
**Date :** 2025-12-27
**Question :** Est-il plus simple d'intégrer rust-cast et le modifier pour le rendre asynchrone ?
---
## 1. Analyse de la codebase rust-cast
### Taille et complexité
```bash
Total : ~5800 lignes de code Rust
Structure modulaire :
├── src/lib.rs (~570 lignes)
├── src/message_manager.rs (~300 lignes)
├── src/channels/
│ ├── media.rs (~800 lignes)
│ ├── receiver.rs (~400 lignes)
│ ├── heartbeat.rs (~100 lignes)
│ └── connection.rs (~100 lignes)
├── src/cast/
│ ├── cast_channel.rs (généré par protobuf)
│ └── proxies.rs (~500 lignes)
└── src/errors.rs, utils.rs (~200 lignes)
```
**Conclusion :** Codebase de taille **modeste et bien structurée**.
---
## 2. Points bloquants identifiés
### 2.1 I/O synchrone bloquant
Tous les I/O passent par `MessageManager<S>``S: Read + Write` :
```rust
// message_manager.rs:246-253
fn read(&self) -> Result<CastMessage, Error> {
let mut buffer: [u8; 4] = [0; 4];
let reader = &mut *self.stream.borrow_mut();
reader.read_exact(&mut buffer)?; // ❌ BLOQUANT
let length = utils::read_u32_from_buffer(&buffer)?;
let mut buffer: Vec<u8> = Vec::with_capacity(length as usize);
let mut limited_reader = reader.take(u64::from(length));
limited_reader.read_to_end(&mut buffer)?; // ❌ BLOQUANT
...
}
```
```rust
// message_manager.rs:138-141
pub fn send(&self, message: CastMessage) -> Result<(), Error> {
...
let writer = &mut *self.stream.borrow_mut();
writer.write_all(&message_length_buffer)?; // ❌ BLOQUANT
writer.write_all(&message_content_buffer)?; // ❌ BLOQUANT
...
}
```
### 2.2 Connexion TLS
```rust
// lib.rs:125
let stream = StreamOwned::new(
conn,
TcpStream::connect((host.as_ref(), port))? // ❌ BLOQUANT
);
```
**Total : 5 points bloquants critiques** (connect, read_exact, read_to_end, 2x write_all)
---
## 3. Effort pour rendre rust-cast asynchrone
### 3.1 Modifications requises
#### A. Remplacer la stack réseau
**Avant (sync) :**
```rust
use std::net::TcpStream;
use rustls::{ClientConnection, StreamOwned};
type TlsStream = StreamOwned<ClientConnection, TcpStream>;
```
**Après (async) :**
```rust
use async_io::Async;
use std::net::TcpStream;
use async_rustls::{TlsConnector, client::TlsStream};
// OU avec tokio :
use tokio::net::TcpStream;
use tokio_rustls::{TlsConnector, client::TlsStream};
```
⚠️ **PROBLÈME :** `rustls::StreamOwned` n'existe pas en version async native. Il faut utiliser :
- `async-rustls` (pour async-std/smol)
- `tokio-rustls` (pour tokio)
Ces crates ont une **API différente** de `rustls::StreamOwned`.
#### B. Modifier `MessageManager`
```diff
- pub struct MessageManager<S> where S: Write + Read {
+ pub struct MessageManager<S> where S: AsyncWrite + AsyncRead + Unpin {
- pub fn send(&self, message: CastMessage) -> Result<(), Error> {
+ pub async fn send(&self, message: CastMessage) -> Result<(), Error> {
...
- writer.write_all(&message_length_buffer)?;
+ writer.write_all(&message_length_buffer).await?;
}
- pub fn receive(&self) -> Result<CastMessage, Error> {
+ pub async fn receive(&self) -> Result<CastMessage, Error> {
...
}
- fn read(&self) -> Result<CastMessage, Error> {
+ async fn read(&self) -> Result<CastMessage, Error> {
- reader.read_exact(&mut buffer)?;
+ reader.read_exact(&mut buffer).await?;
- limited_reader.read_to_end(&mut buffer)?;
+ limited_reader.read_to_end(&mut buffer).await?;
}
}
```
#### C. Propager `async` dans tous les channels
**Avant :**
```rust
// channels/media.rs
impl<'a, S> MediaChannel<'a, S> where S: Write + Read {
pub fn play(&self, ...) -> Result<(), Error> {
self.message_manager.send(...)?;
self.message_manager.receive_find_map(...)
}
}
```
**Après :**
```rust
impl<'a, S> MediaChannel<'a, S> where S: AsyncWrite + AsyncRead + Unpin {
pub async fn play(&self, ...) -> Result<(), Error> {
self.message_manager.send(...).await?;
self.message_manager.receive_find_map(...).await
}
}
```
**Impact :** TOUS les channels (media, receiver, heartbeat, connection) deviennent `async`.
#### D. Modifier `CastDevice`
```diff
impl<'a> CastDevice<'a> {
- pub fn connect<S>(host: S, port: u16) -> Result<CastDevice<'a>, Error>
+ pub async fn connect<S>(host: S, port: u16) -> Result<CastDevice<'a>, Error>
{
...
- let stream = TcpStream::connect((host.as_ref(), port))?;
+ let stream = TcpStream::connect((host.as_ref(), port)).await?;
...
}
- pub fn receive(&self) -> Result<ChannelMessage, Error> {
+ pub async fn receive(&self) -> Result<ChannelMessage, Error> {
- let cast_message = self.message_manager.receive()?;
+ let cast_message = self.message_manager.receive().await?;
...
}
}
```
### 3.2 Estimation de l'effort
| Tâche | Fichiers touchés | Complexité | Temps estimé |
|-------|------------------|------------|--------------|
| Choisir stack async (smol vs tokio) | - | Faible | 1h |
| Migrer vers async-rustls/tokio-rustls | lib.rs | **MOYENNE** | 4-6h |
| Rendre MessageManager async | message_manager.rs | **MOYENNE** | 4-6h |
| Rendre tous les channels async | 4 fichiers | **MOYENNE-ÉLEVÉE** | 8-12h |
| Mettre à jour CastDevice | lib.rs | Moyenne | 2-4h |
| Tests et debug | Tous | **ÉLEVÉE** | 8-16h |
| **TOTAL** | **~10 fichiers** | **ÉLEVÉE** | **27-45 heures** |
⚠️ **RISQUES :**
- API `async-rustls` différente de `rustls::StreamOwned` → peut nécessiter refactoring profond
- Gestion des locks async (`Mutex``async_lock::Mutex` ou `tokio::sync::Mutex`)
- Bugs subtils liés à la concurrence async
- Tests nécessaires pour valider la stabilité
---
## 4. Comparaison des 4 options
### Option 1 : ✅ **Rester avec rust-cast sync et corriger le TLS**
**Effort :** FAIBLE (2-8 heures)
**Actions :**
- Investiguer les erreurs TLS prématurées
- Ajouter retry logic sur les reconnexions
- Améliorer la gestion d'erreur dans [chromecast_renderer.rs](pmocontrol/src/chromecast_renderer.rs:86-104)
- Peut-être ajuster les timeouts de lecture
**Avantages :**
- ✅ Garde l'API sync compatible avec PMOMusic
- ✅ Risque minimal
- ✅ Solution rapide
**Inconvénients :**
- ⚠️ Ne résout peut-être pas tous les problèmes TLS
---
### Option 2 : 🔧 **Forker rust-cast et moderniser le TLS (reste sync)**
**Effort :** MOYEN (8-16 heures)
**Actions :**
- Forker rust-cast sur GitHub/GitLab
- Améliorer la gestion TLS (retry, reconnexion automatique)
- Ajouter logs détaillés
- Corriger les bugs TLS identifiés
- Maintenir un fork privé
**Avantages :**
- ✅ Garde l'API sync
- ✅ Contrôle total sur les correctifs
- ✅ Peut merger les améliorations de upstream
**Inconvénients :**
- ⚠️ Maintenance du fork à long terme
- ⚠️ Doit suivre les mises à jour de rustls
---
### Option 3 : 🔄 **Rendre rust-cast asynchrone**
**Effort :** ÉLEVÉ (27-45 heures)
**Actions :**
- Migrer vers async-rustls ou tokio-rustls
- Rendre tout le code async (MessageManager, channels, CastDevice)
- Adapter PMOMusic pour wrapper les appels async
**Avantages :**
- ✅ Architecture moderne
- ✅ Potentiellement meilleure performance pour gérer plusieurs devices
- ✅ Résout probablement les problèmes TLS via stack moderne
**Inconvénients :**
- ❌ Effort très élevé
- ❌ Risque de bugs subtils
- ❌ PMOMusic doit wrapper tous les appels avec `smol::block_on()`
- ❌ Overhead de conversion sync→async→sync
**⚠️ PARADOXE :** Rendre rust-cast async pour ensuite le wrapper en sync dans PMOMusic = **surcharge inutile**
---
### Option 4 : ❌ **Migrer vers cast-sender (déjà async)**
**Effort :** TRÈS ÉLEVÉ (40-80 heures)
**Problèmes critiques :**
- ❌ API incomplète (pas de get_status, pas de seek)
- ❌ Nécessite architecture stateful complexe
- ❌ Documentation insuffisante (23%)
**Voir :** [cast-sender-evaluation.md](cast-sender-evaluation.md)
---
## 5. Analyse détaillée : Async est-il vraiment utile ?
### 5.1 Cas d'usage PMOMusic
**Architecture actuelle :**
- 1 thread par Chromecast actif (pour le heartbeat)
- Opérations de contrôle (play, pause, volume) : sporadiques
- Pas de gestion massive de connexions simultanées
**Bénéfice de async :**
-**FAIBLE** : PMOMusic n'a pas besoin de gérer 100+ connexions simultanées
-**OVERHEAD** : Wrapping sync→async→sync ajoute de la complexité
### 5.2 Vraie cause des problèmes TLS ?
Les problèmes de "fermeture TLS prématurée" sont probablement dus à :
- Timeout réseau trop court
- Gestion d'erreur insuffisante lors des reconnexions
- Bugs spécifiques de certaines versions de rustls
**Async ne résout PAS directement ces problèmes !**
---
## 6. Recommandation finale
### 🏆 **Option recommandée : Option 1 (Corriger rust-cast sync)**
**Raisons :**
1. **Effort minimal** : 2-8 heures vs 27-45h pour async
2. **Risque minimal** : Garde l'architecture validée
3. **Compatibilité** : Pas de changement dans PMOMusic
4. **Pragmatique** : Résout le problème réel (TLS) sans over-engineering
**Plan d'action concret :**
```rust
// Améliorer la fonction connect_to_device
fn connect_to_device(host: &str, port: u16) -> Result<CastDevice> {
const MAX_RETRIES: u32 = 3;
const RETRY_DELAY_MS: u64 = 1000;
for attempt in 1..=MAX_RETRIES {
match try_connect(host, port) {
Ok(device) => return Ok(device),
Err(e) if attempt < MAX_RETRIES => {
tracing::warn!(
"Connection attempt {} failed: {}. Retrying in {}ms...",
attempt, e, RETRY_DELAY_MS
);
std::thread::sleep(Duration::from_millis(RETRY_DELAY_MS));
}
Err(e) => return Err(e),
}
}
unreachable!()
}
// Ajouter timeout configurable pour les read operations
// Ajouter meilleure gestion d'erreur dans le heartbeat loop
```
---
### 🥈 **Alternative : Option 2 (Fork rust-cast)**
Si l'Option 1 ne suffit pas après investigation, forker permet :
- Corrections TLS plus profondes
- Ajout de fonctionnalités manquantes
- Contrôle total
**Pas besoin de rendre async !**
---
### 🚫 **Options déconseillées :**
-**Option 3** (Async rust-cast) : Effort 5-10x supérieur pour bénéfice marginal
-**Option 4** (cast-sender) : API incomplète, effort encore plus élevé
---
## 7. Conclusion
**NON, rendre rust-cast asynchrone n'est PAS plus simple.**
**Comparaison des efforts :**
| Option | Effort (heures) | Complexité | Risque |
|--------|----------------|------------|--------|
| 1. Corriger rust-cast sync | 2-8 | Faible | Minimal |
| 2. Forker rust-cast | 8-16 | Moyenne | Faible |
| 3. **Async rust-cast** | **27-45** | **Élevée** | **Élevé** |
| 4. Migrer cast-sender | 40-80 | Très élevée | Très élevé |
**Le ratio effort/bénéfice de l'option async est défavorable :**
- **5-10x plus d'effort** que corriger le code sync
- **Bénéfice minimal** pour l'architecture actuelle de PMOMusic
- **Risques élevés** de bugs de concurrence async
**Recommandation :** Commencer par l'**Option 1**, investiguer les vrais problèmes TLS, et envisager l'**Option 2** (fork) uniquement si nécessaire. Éviter absolument l'**Option 3** (async) sauf changement radical d'architecture de PMOMusic.
---
## Annexe : Si vous vouliez quand même faire async...
### Stack recommandée
**Pour PMOMusic (déjà avec smol) :**
```toml
[dependencies]
async-io = "2.3"
async-rustls = "0.4"
futures-lite = "2.1"
```
**Points d'attention :**
- Remplacer tous les `Mutex` par `async_lock::Mutex`
- Gérer correctement le `Unpin` trait pour les streams
- Tester intensivement la gestion des erreurs async
- Prévoir 2-3 semaines de développement + tests
**Mais encore une fois : le jeu n'en vaut pas la chandelle !**

View File

@@ -94,6 +94,12 @@ fn connect_to_device<'a>(host: &'a str, port: u16) -> Result<CastDevice<'a>> {
.connect(DEFAULT_DESTINATION_ID.to_string())
.map_err(|e| anyhow!("Failed to connect channel: {}", e))?;
// Send initial gre to establish heartbeat communication
// This is critical per rust_caster.rs example
device.heartbeat
.ping()
.map_err(|e| anyhow!("Failed to send initial heartbeat ping: {}", e))?;
Ok(device)
}

View File

@@ -14,7 +14,7 @@ use pmodidl::{DIDLLite, Item as DidlItem, Resource as DidlResource};
use pmoupnp::ssdp::SsdpClient;
use quick_xml::se::to_string as to_didl_string;
use thiserror::Error;
use tracing::{debug, error, info, trace, warn};
use tracing::{debug, error, info, warn};
use ureq::{http, Agent};
use xmltree::{Element, XMLNode};
@@ -31,7 +31,9 @@ use crate::control_point::music_queue::MusicQueue;
use crate::control_point::openhome_queue::OpenHomeQueue;
use crate::discovery::DiscoveryManager;
use crate::events::{MediaServerEventBus, RendererEventBus};
use crate::media_server::{MediaBrowser, MediaEntry, MediaServerInfo, MusicServer, ServerId};
use crate::media_server::{
playback_item_from_entry, MediaBrowser, MediaServerInfo, MusicServer, ServerId,
};
use crate::media_server_events::spawn_media_server_event_runtime;
use crate::model::TrackMetadata;
use crate::model::{MediaServerEvent, RendererEvent, RendererId, RendererInfo};
@@ -50,7 +52,7 @@ use crate::openhome_client::parse_track_metadata_from_didl;
use crate::openhome_playlist::{OpenHomePlaylistSnapshot, OpenHomePlaylistTrack};
use crate::openhome_renderer::{format_seconds, map_openhome_state};
use crate::provider::HttpXmlDescriptionProvider;
use crate::queue_backend::{EnqueueMode, PlaybackItem, QueueBackend};
use crate::queue_backend::{EnqueueMode, PlaybackItem, QueueBackend, QueueSnapshot};
use crate::queue_interne::InternalQueue;
use crate::registry::{DeviceRegistry, DeviceRegistryRead, DeviceUpdate};
use crate::upnp_renderer::UpnpRenderer;
@@ -939,17 +941,15 @@ impl ControlPoint {
// User-driven mutation: detach any playlist binding
self.detach_playlist_binding(renderer_id, "clear_queue");
if self.runtime.uses_openhome_playlist(renderer_id) {
let renderer = self.openhome_renderer(renderer_id)?;
renderer.openhome_playlist_clear()?;
self.sync_openhome_playlist_for(renderer_id)?;
debug!(
renderer = renderer_id.0.as_str(),
"Cleared OpenHome playlist"
);
return Ok(());
}
// Clear the queue on the backend (backend-agnostic)
let renderer = self.music_renderer_by_id(renderer_id)
.ok_or_else(|| anyhow!("Renderer {} not found", renderer_id.0))?;
renderer.clear_queue()?;
// Sync backend state to local cache (for OpenHome, this updates the cache)
renderer.sync_queue_state()?;
// Clear the local queue state
let removed = self.runtime.with_music_queue_mut(renderer_id, |queue| {
let removed = queue.upcoming_len()?;
queue.clear_queue()?;
@@ -996,11 +996,7 @@ impl ControlPoint {
// User-driven mutation: detach any playlist binding
self.detach_playlist_binding(renderer_id, "enqueue_items");
if self.runtime.uses_openhome_playlist(renderer_id) {
self.enqueue_items_openhome(renderer_id, items)?;
return Ok(());
}
// Enqueue items using QueueBackend abstraction (works for both backends)
let item_count = items.len();
let new_len = self.runtime.with_music_queue_mut(renderer_id, |queue| {
queue.enqueue_items(items, EnqueueMode::AppendToEnd)?;
@@ -1078,20 +1074,25 @@ impl ControlPoint {
self.runtime.current_track_metadata(renderer_id)
}
/// Force a resynchronization of the OpenHome playlist cache for a renderer.
/// Gets the backend queue snapshot for renderers with persistent queues.
///
/// This is used by external APIs after mutating the native playlist so that
/// the local queue mirrors the renderer state.
pub fn refresh_openhome_playlist(&self, renderer_id: &RendererId) -> anyhow::Result<()> {
self.sync_openhome_playlist_for(renderer_id)
}
pub fn get_openhome_playlist_snapshot(
/// Returns the queue snapshot if the renderer has a backend queue (e.g., OpenHome),
/// or None if it doesn't (e.g., AVTransport).
pub fn get_renderer_queue_snapshot(
&self,
renderer_id: &RendererId,
) -> anyhow::Result<OpenHomePlaylistSnapshot> {
let renderer = self.openhome_renderer(renderer_id)?;
renderer.openhome_playlist_snapshot()
) -> anyhow::Result<Option<QueueSnapshot>> {
let renderer = self.music_renderer_by_id(renderer_id)
.ok_or_else(|| anyhow!("Renderer {} not found", renderer_id.0))?;
renderer.queue_snapshot()
}
/// Gets the length of the backend queue for renderers with persistent queues.
///
/// Returns the queue length if the renderer has a backend queue, or 0 if it doesn't.
pub fn get_renderer_queue_length(&self, renderer_id: &RendererId) -> anyhow::Result<usize> {
let snapshot = self.get_renderer_queue_snapshot(renderer_id)?;
Ok(snapshot.map(|s| s.len()).unwrap_or(0))
}
pub fn get_cached_openhome_playlist_snapshot(
@@ -1103,11 +1104,6 @@ impl ControlPoint {
self.runtime.openhome_snapshot_cached(&renderer, ttl)
}
pub fn get_openhome_playlist_len(&self, renderer_id: &RendererId) -> anyhow::Result<usize> {
let renderer = self.openhome_renderer(renderer_id)?;
renderer.openhome_playlist_len()
}
/// Build a fully consistent snapshot for UI consumers (state + queue + binding).
#[cfg(feature = "pmoserver")]
pub fn renderer_full_snapshot(
@@ -1119,123 +1115,31 @@ impl ControlPoint {
.ok_or_else(|| anyhow!("Renderer {} not found", renderer_id.0))?;
let info = renderer.info();
// MICRO-PATCH 5: Chemin OpenHome complètement découplé du miroir local
if self.runtime.uses_openhome_playlist(renderer_id) {
// Pour OpenHome: récupérer UNIQUEMENT runtime_snapshot pour volume/state/position
// Ne PAS utiliser queue_items/current_index du runtime (miroir local)
let (runtime_snapshot, _, _) = self.runtime.renderer_snapshot_bundle(renderer_id);
// Get runtime snapshot (needed for all backends: volume, state, position)
let (runtime_snapshot, runtime_queue_items, runtime_current_index) =
self.runtime.renderer_snapshot_bundle(renderer_id);
// OpenHome est la source de vérité pour la queue - pas de fallback au runtime
let snapshot = self.get_cached_openhome_playlist_snapshot(
renderer_id,
OPENHOME_SNAPSHOT_CACHE_TTL,
)?;
let queue_items: Vec<PlaybackItem> = snapshot
.tracks
.iter()
.map(|track| playback_item_from_openhome_track(renderer_id, track))
.collect();
let queue_len = snapshot.tracks.len();
// Pour OpenHome: current_index vient UNIQUEMENT d'OpenHome, pas d'heuristiques runtime
let queue_current_index = snapshot.current_index.or_else(|| {
snapshot.current_id.and_then(|id| {
snapshot.tracks.iter().position(|track| track.id == id)
})
});
// Try to get queue from backend (returns Some for OpenHome, None for others)
let backend_queue = renderer.queue_snapshot()?;
// Use backend queue if available (OpenHome), otherwise use runtime queue (AVTransport)
let (queue_items, mut queue_current_index, from_backend) = if let Some(snapshot) = backend_queue {
debug!(
renderer = renderer_id.0.as_str(),
current_id = ?snapshot.current_id,
current_index = ?queue_current_index,
track_count = snapshot.tracks.len(),
"renderer_full_snapshot: OpenHome snapshot retrieved"
current_index = ?snapshot.current_index,
track_count = snapshot.items.len(),
"renderer_full_snapshot: using backend queue (OpenHome)"
);
(snapshot.items, snapshot.current_index, true)
} else {
(runtime_queue_items, runtime_current_index, false)
};
let queue_view_items: Vec<QueueItem> = queue_items
.iter()
.enumerate()
.map(|(index, item)| QueueItem {
index,
uri: item.uri.clone(),
title: item.metadata.as_ref().and_then(|m| m.title.clone()),
artist: item.metadata.as_ref().and_then(|m| m.artist.clone()),
album: item.metadata.as_ref().and_then(|m| m.album.clone()),
album_art_uri: item.metadata.as_ref().and_then(|m| m.album_art_uri.clone()),
server_id: Some(item.media_server_id.0.clone()),
object_id: Some(item.didl_id.clone()),
})
.collect();
let queue_view = QueueSnapshotView {
renderer_id: renderer_id.0.clone(),
items: queue_view_items,
current_index: queue_current_index,
};
let binding = self.current_queue_playlist_binding(renderer_id).map(
|(server_id, container_id, has_seen_update)| RendererBindingView {
server_id: server_id.0,
container_id,
has_seen_update,
},
);
let (position_ms, duration_ms) =
convert_runtime_position(runtime_snapshot.position.as_ref());
let queue_current_metadata = queue_current_index
.and_then(|idx| queue_items.get(idx))
.map(current_track_from_playback_item);
// MICRO-PATCH 5: Pour OpenHome, préférer les métadonnées depuis le snapshot OpenHome
// car runtime_snapshot.last_metadata n'est jamais mis à jour pour OpenHome
let current_track = queue_current_metadata.or_else(|| {
runtime_snapshot
.last_metadata
.as_ref()
.map(|meta| CurrentTrackMetadata {
title: meta.title.clone(),
artist: meta.artist.clone(),
album: meta.album.clone(),
album_art_uri: meta.album_art_uri.clone(),
})
});
let state_view = RendererStateView {
id: renderer_id.0.clone(),
friendly_name: info.friendly_name.clone(),
transport_state: runtime_snapshot
.state
.as_ref()
.map(|state| state.as_str().to_string())
.unwrap_or_else(|| "UNKNOWN".to_string()),
position_ms,
duration_ms,
volume: runtime_snapshot
.last_volume
.and_then(|value| u8::try_from(value).ok()),
mute: runtime_snapshot.last_mute,
queue_len,
attached_playlist: binding.clone(),
current_track,
};
return Ok(FullRendererSnapshot {
state: state_view,
queue: queue_view,
binding,
});
}
// Chemin non-OpenHome (UPnP AV): continue d'utiliser le runtime
let (runtime_snapshot, queue_items, mut queue_current_index) =
self.runtime.renderer_snapshot_bundle(renderer_id);
let playback_source = self.runtime.playback_source(renderer_id);
let queue_len = queue_items.len();
if queue_current_index.is_none() {
// For non-backend queues (AVTransport), try heuristics to determine current_index
if !from_backend && queue_current_index.is_none() {
if let Some(position) = runtime_snapshot.position.as_ref() {
if let Some(uri) = position.track_uri.as_ref() {
if let Some(idx) = queue_items.iter().position(|item| item.uri == *uri) {
@@ -1248,18 +1152,19 @@ impl ControlPoint {
}
}
}
}
if queue_current_index.is_none()
&& matches!(playback_source, PlaybackSource::FromQueue)
&& runtime_snapshot
.state
.as_ref()
.map(|state| matches!(state, PlaybackState::Playing | PlaybackState::Paused))
.unwrap_or(false)
&& !queue_items.is_empty()
{
queue_current_index = Some(0);
// Final fallback: if playing from queue and no index, assume first track
if queue_current_index.is_none()
&& matches!(playback_source, PlaybackSource::FromQueue)
&& runtime_snapshot
.state
.as_ref()
.map(|state| matches!(state, PlaybackState::Playing | PlaybackState::Paused))
.unwrap_or(false)
&& !queue_items.is_empty()
{
queue_current_index = Some(0);
}
}
let queue_view_items: Vec<QueueItem> = queue_items
@@ -1297,16 +1202,32 @@ impl ControlPoint {
.and_then(|idx| queue_items.get(idx))
.map(current_track_from_playback_item);
let current_track = runtime_snapshot
.last_metadata
.as_ref()
.map(|meta| CurrentTrackMetadata {
title: meta.title.clone(),
artist: meta.artist.clone(),
album: meta.album.clone(),
album_art_uri: meta.album_art_uri.clone(),
// For backend queues (OpenHome), prefer queue metadata since runtime isn't updated
// For runtime queues (AVTransport), prefer runtime metadata which is fresher
let current_track = if from_backend {
queue_current_metadata.or_else(|| {
runtime_snapshot
.last_metadata
.as_ref()
.map(|meta| CurrentTrackMetadata {
title: meta.title.clone(),
artist: meta.artist.clone(),
album: meta.album.clone(),
album_art_uri: meta.album_art_uri.clone(),
})
})
.or(queue_current_metadata);
} else {
runtime_snapshot
.last_metadata
.as_ref()
.map(|meta| CurrentTrackMetadata {
title: meta.title.clone(),
artist: meta.artist.clone(),
album: meta.album.clone(),
album_art_uri: meta.album_art_uri.clone(),
})
.or(queue_current_metadata)
};
let state_view = RendererStateView {
id: renderer_id.0.clone(),
@@ -1334,35 +1255,53 @@ impl ControlPoint {
})
}
pub fn clear_openhome_playlist(&self, renderer_id: &RendererId) -> anyhow::Result<()> {
let renderer = self.openhome_renderer(renderer_id)?;
renderer.openhome_playlist_clear()?;
self.sync_openhome_playlist_for(renderer_id)
/// Clears the renderer's backend queue.
///
/// For renderers with persistent queues (OpenHome), this clears the queue on the renderer.
/// For other renderers, this returns an error.
pub fn clear_renderer_queue(&self, renderer_id: &RendererId) -> anyhow::Result<()> {
let renderer = self.music_renderer_by_id(renderer_id)
.ok_or_else(|| anyhow!("Renderer {} not found", renderer_id.0))?;
renderer.clear_queue()?;
renderer.sync_queue_state()
}
pub fn add_openhome_track(
/// Adds a track to the renderer's backend queue.
///
/// For renderers with persistent queues (OpenHome), this adds the track to the queue.
/// For other renderers, this returns an error.
///
/// Returns the backend-specific track ID if applicable.
pub fn add_track_to_renderer(
&self,
renderer_id: &RendererId,
uri: &str,
metadata: &str,
after_id: Option<u32>,
play: bool,
) -> anyhow::Result<()> {
let renderer = self.openhome_renderer(renderer_id)?;
renderer.openhome_playlist_add_track(uri, metadata, after_id, play)?;
self.sync_openhome_playlist_for(renderer_id)
) -> anyhow::Result<Option<u32>> {
let renderer = self.music_renderer_by_id(renderer_id)
.ok_or_else(|| anyhow!("Renderer {} not found", renderer_id.0))?;
let track_id = renderer.add_track_to_queue(uri, metadata, after_id, play)?;
renderer.sync_queue_state()?;
Ok(track_id)
}
pub fn play_openhome_track_id(
/// Selects and plays a specific track from the renderer's backend queue.
///
/// For renderers with persistent queues (OpenHome), this uses the track ID.
/// For other renderers, this returns an error.
pub fn select_renderer_track(
&self,
renderer_id: &RendererId,
track_id: u32,
) -> anyhow::Result<()> {
let renderer = self.openhome_renderer(renderer_id)?;
renderer.openhome_playlist_play_id(track_id)?;
let renderer = self.music_renderer_by_id(renderer_id)
.ok_or_else(|| anyhow!("Renderer {} not found", renderer_id.0))?;
renderer.select_queue_track(track_id)?;
self.runtime
.set_playback_source(renderer_id, PlaybackSource::FromQueue);
self.sync_openhome_playlist_for(renderer_id)
renderer.sync_queue_state()
}
/// Plays the current queue item without advancing the index.
@@ -1381,10 +1320,18 @@ impl ControlPoint {
return Err(err);
}
if self.runtime.uses_openhome_playlist(renderer_id) {
return self.play_current_openhome(renderer_id);
let renderer = self.music_renderer_by_id(renderer_id).ok_or_else(|| {
anyhow!("Renderer {} not found", renderer_id.0)
})?;
// If renderer has a backend queue (OpenHome), use backend playback
if renderer.queue_snapshot()?.is_some() {
renderer.play_current_from_backend_queue()?;
self.runtime.set_playback_source(renderer_id, PlaybackSource::FromQueue);
return Ok(());
}
// Otherwise use local queue playback (AVTransport)
let Some((item, remaining)) = self.runtime.peek_current(renderer_id) else {
debug!(
renderer = renderer_id.0.as_str(),
@@ -1402,13 +1349,9 @@ impl ControlPoint {
"Playing current playback item from queue"
);
let renderer = self.music_renderer_by_id(renderer_id).ok_or_else(|| {
warn!(
renderer = renderer_id.0.as_str(),
"Renderer disappeared before queue playback could start"
);
anyhow!("Renderer {} not found", renderer_id.0)
})?;
// Temporarily disable auto-advance to prevent race condition
// when renderer sends Stopped event during SetAVTransportURI
self.runtime.set_playback_source(renderer_id, PlaybackSource::None);
let playback = (|| -> anyhow::Result<()> {
let didl_metadata = playback_item_to_didl(&item);
@@ -1462,8 +1405,14 @@ impl ControlPoint {
return Err(err);
}
if self.runtime.uses_openhome_playlist(renderer_id) {
self.play_next_openhome(renderer_id)?;
let renderer = self.music_renderer_by_id(renderer_id).ok_or_else(|| {
anyhow!("Renderer {} not found", renderer_id.0)
})?;
// If renderer has a backend queue (OpenHome), use backend playback
if renderer.queue_snapshot()?.is_some() {
renderer.play_next_from_backend_queue()?;
self.runtime.set_playback_source(renderer_id, PlaybackSource::FromQueue);
return Ok(());
}
@@ -1494,6 +1443,10 @@ impl ControlPoint {
anyhow!("Renderer {} not found", renderer_id.0)
})?;
// Temporarily disable auto-advance to prevent race condition
// when renderer sends Stopped event during SetAVTransportURI
self.runtime.set_playback_source(renderer_id, PlaybackSource::None);
let playback = (|| -> anyhow::Result<()> {
let didl_metadata = playback_item_to_didl(&item);
renderer.play_uri(&item.uri, &didl_metadata)?;
@@ -1570,6 +1523,126 @@ impl ControlPoint {
Ok(())
}
/// Jumps to a specific index in the queue and starts playback.
///
/// For OpenHome renderers, this uses the playlist's native SeekId capability.
/// For internal queues, this updates the current index and starts playback.
pub fn play_queue_index(&self, renderer_id: &RendererId, index: usize) -> anyhow::Result<()> {
if !self.runtime.has_entry(renderer_id) {
let err = Self::runtime_entry_missing(renderer_id);
warn!(
renderer = renderer_id.0.as_str(),
"Cannot jump to index: renderer not registered in runtime"
);
return Err(err);
}
let renderer = self.music_renderer_by_id(renderer_id).ok_or_else(|| {
anyhow!("Renderer {} not found", renderer_id.0)
})?;
// For backend queues (OpenHome), use select_track_index which handles the track_id lookup
if renderer.queue_snapshot()?.is_some() {
info!(
renderer = renderer_id.0.as_str(),
index,
"Seeking to OpenHome queue index"
);
// Use the queue to select the track by index
self.runtime.with_music_queue_mut(renderer_id, |queue| {
if let crate::control_point::music_queue::MusicQueue::OpenHome(oh_queue) = queue {
oh_queue.select_track_index(index)
} else {
Err(anyhow!("Expected OpenHome queue"))
}
})?;
info!(
renderer = renderer_id.0.as_str(),
index,
"Successfully seeked to OpenHome queue index"
);
self.runtime
.set_playback_source(renderer_id, PlaybackSource::FromQueue);
return self.sync_openhome_playlist_for(renderer_id);
}
// For internal queue, set the index and play
self.runtime.with_music_queue_mut(renderer_id, |queue| {
queue.set_index(Some(index))
})?;
let Some((item, remaining)) = self.runtime.peek_current(renderer_id) else {
debug!(
renderer = renderer_id.0.as_str(),
index,
"play_queue_index: no item at index"
);
self.runtime
.set_playback_source(renderer_id, PlaybackSource::None);
return Ok(());
};
debug!(
renderer = renderer_id.0.as_str(),
index,
queue_len = remaining + 1,
uri = item.uri.as_str(),
"Playing item at index from queue"
);
let renderer = self.music_renderer_by_id(renderer_id).ok_or_else(|| {
warn!(
renderer = renderer_id.0.as_str(),
"Renderer disappeared before queue playback could start"
);
anyhow!("Renderer {} not found", renderer_id.0)
})?;
// Temporarily disable auto-advance to prevent race condition
// when renderer sends Stopped event during SetAVTransportURI
self.runtime.set_playback_source(renderer_id, PlaybackSource::None);
let playback = (|| -> anyhow::Result<()> {
let didl_metadata = playback_item_to_didl(&item);
renderer.play_uri(&item.uri, &didl_metadata)?;
Ok(())
})();
match playback {
Ok(()) => {
info!(
renderer = renderer_id.0.as_str(),
index,
uri = item.uri.as_str(),
"Queue playback started at index"
);
let metadata = playback_item_track_metadata(&item);
self.runtime.update_snapshot_with(renderer_id, |snapshot| {
snapshot.last_metadata = Some(metadata);
});
self.runtime
.set_playback_source(renderer_id, PlaybackSource::FromQueue);
Ok(())
}
Err(err) => {
error!(
renderer = renderer_id.0.as_str(),
index,
error = %err,
"Failed to start playback for queued item at index"
);
self.runtime
.set_playback_source(renderer_id, PlaybackSource::None);
Err(err)
}
}
}
fn start_queue_playback_if_idle(&self, renderer_id: &RendererId) -> anyhow::Result<()> {
let snapshot = self.runtime.snapshot_for(renderer_id);
let renderer_playing = matches!(snapshot.state, Some(PlaybackState::Playing));
@@ -1608,7 +1681,13 @@ impl ControlPoint {
return Ok(());
}
if self.runtime.uses_openhome_playlist(renderer_id) {
let renderer = self.music_renderer_by_id(renderer_id).ok_or_else(|| {
anyhow!("Renderer {} not found", renderer_id.0)
})?;
// For backend queues (OpenHome), play current (renderer tracks its own position)
// For local queues (AVTransport), dequeue next (we manage position locally)
if renderer.queue_snapshot()?.is_some() {
self.play_current_from_queue(renderer_id)
} else {
self.play_next_from_queue(renderer_id)
@@ -1708,22 +1787,26 @@ impl ControlPoint {
"Attaching new playlist: clearing renderer queue"
);
// Clear the renderer's queue for OpenHome renderers
// We also sync the local cache to reflect the empty state, which will trigger
// refresh_attached_queue_for() to use replace_entire_playlist() instead of gentle sync
if self.runtime.uses_openhome_playlist(renderer_id) {
let renderer = self.openhome_renderer(renderer_id)?;
renderer.openhome_playlist_clear()?;
// Sync local cache to reflect the empty renderer state
self.sync_openhome_playlist_for(renderer_id)?;
debug!(
renderer = renderer_id.0.as_str(),
"Cleared OpenHome renderer playlist and synced local cache"
);
} else {
// For non-OpenHome renderers, use the standard clear_queue
self.clear_queue(renderer_id)?;
}
// Prepare the renderer for the new playlist (backend-agnostic)
let renderer = self.music_renderer_by_id(renderer_id).ok_or_else(|| {
anyhow!("Renderer {} not found", renderer_id.0)
})?;
renderer.clear_for_playlist_attach()?;
// Sync backend state to local cache (backend-agnostic)
renderer.sync_queue_state()?;
// Clear the local queue (detach binding + clear runtime queue structure)
self.detach_playlist_binding(renderer_id, "attach_new_playlist");
self.runtime.with_music_queue_mut(renderer_id, |queue| {
queue.clear_queue()?;
Ok(())
})?;
debug!(
renderer = renderer_id.0.as_str(),
"Cleared renderer and local queue for new playlist"
);
let binding = PlaylistBinding {
server_id: server_id.clone(),
@@ -1750,7 +1833,14 @@ impl ControlPoint {
binding: Some(binding),
});
let mut auto_start_cb = |rid: &RendererId| self.start_queue_playback_if_idle(rid);
// For initial attach with auto_play, force playback start (don't check if idle)
let mut auto_start_cb = |rid: &RendererId| {
debug!(
renderer = rid.0.as_str(),
"Attach callback: forcing playback start (not checking if idle)"
);
self.play_current_from_queue(rid)
};
let callback: Option<&mut dyn FnMut(&RendererId) -> anyhow::Result<()>> = if auto_play {
Some(&mut auto_start_cb)
} else {
@@ -1887,176 +1977,6 @@ impl ControlPoint {
sync_openhome_playlist(&self.registry, &self.runtime, &self.event_bus, renderer_id)
}
fn enqueue_items_openhome(
&self,
renderer_id: &RendererId,
items: Vec<PlaybackItem>,
) -> anyhow::Result<()> {
if items.is_empty() {
return Ok(());
}
let renderer = self.openhome_renderer(renderer_id)?;
// Get the last track ID from the OpenHome native playlist
// This ensures we append to the end of the actual playlist, not just our local queue
let mut after_id = renderer
.openhome_playlist_ids()
.ok()
.and_then(|ids| ids.last().copied());
for item in items.iter() {
let metadata = playback_item_to_didl(item);
debug!(
uri = item.uri.as_str(),
protocol_info = item.protocol_info.as_str(),
metadata_len = metadata.len(),
"Inserting track to OpenHome playlist"
);
trace!(metadata = metadata.as_str(), "DIDL-Lite metadata");
after_id =
Some(renderer.openhome_playlist_add_track(&item.uri, &metadata, after_id, false)?);
}
self.sync_openhome_playlist_for(renderer_id)?;
Ok(())
}
fn play_current_openhome(&self, renderer_id: &RendererId) -> anyhow::Result<()> {
let renderer = self.openhome_renderer(renderer_id)?;
// MICRO-PATCH 5: Pour OpenHome, récupérer les données directement depuis la playlist native
// au lieu du miroir local (qui peut être vide ou obsolète)
match self.get_cached_openhome_playlist_snapshot(
renderer_id,
OPENHOME_SNAPSHOT_CACHE_TTL,
) {
Ok(snapshot) => {
if snapshot.tracks.is_empty() {
debug!(
renderer = renderer_id.0.as_str(),
"OpenHome playlist is empty"
);
self.runtime
.set_playback_source(renderer_id, PlaybackSource::None);
return Ok(());
}
// Trouver le track_id courant depuis le snapshot OpenHome
let target_track_id = if let Some(current_id) = snapshot.current_id {
Some(current_id)
} else if let Some(current_idx) = snapshot.current_index {
snapshot.tracks.get(current_idx).map(|track| track.id)
} else {
snapshot.tracks.first().map(|track| track.id)
};
if let Some(track_id) = target_track_id {
match renderer.openhome_playlist_play_id(track_id) {
Ok(()) => {
self.runtime
.set_playback_source(renderer_id, PlaybackSource::FromQueue);
self.sync_openhome_playlist_for(renderer_id)?;
info!(
renderer = renderer_id.0.as_str(),
track_id,
playlist_len = snapshot.tracks.len(),
"Started OpenHome playlist playback (current item)"
);
return Ok(());
}
Err(err) => {
warn!(
renderer = renderer_id.0.as_str(),
track_id,
error = %err,
"PlayId failed, falling back to Play()"
);
}
}
}
// Fallback: appeler Play() sans spécifier de track_id
renderer.play()?;
self.runtime
.set_playback_source(renderer_id, PlaybackSource::FromQueue);
self.sync_openhome_playlist_for(renderer_id)?;
info!(
renderer = renderer_id.0.as_str(),
playlist_len = snapshot.tracks.len(),
"Started OpenHome native playlist playback"
);
return Ok(());
}
Err(err) => {
warn!(
renderer = renderer_id.0.as_str(),
error = %err,
"Failed to fetch OpenHome playlist snapshot, playlist might be empty"
);
debug!(
renderer = renderer_id.0.as_str(),
"OpenHome playlist is empty or unavailable"
);
self.runtime
.set_playback_source(renderer_id, PlaybackSource::None);
return Ok(());
}
}
}
fn play_next_openhome(&self, renderer_id: &RendererId) -> anyhow::Result<()> {
let renderer = self.openhome_renderer(renderer_id)?;
// MICRO-PATCH 5: Récupérer les données directement depuis OpenHome au lieu du miroir local
let snapshot = self.get_cached_openhome_playlist_snapshot(
renderer_id,
OPENHOME_SNAPSHOT_CACHE_TTL,
)?;
if snapshot.tracks.is_empty() {
debug!(
renderer = renderer_id.0.as_str(),
"OpenHome playlist is empty, cannot play next"
);
self.runtime
.set_playback_source(renderer_id, PlaybackSource::None);
return Ok(());
}
// Déterminer le prochain track_id
let next_track_id = match snapshot.current_index {
Some(idx) => {
// Prendre la piste suivante si elle existe
snapshot
.tracks
.get(idx + 1)
.map(|track| track.id)
.or_else(|| snapshot.tracks.first().map(|track| track.id))
}
None => snapshot.tracks.first().map(|track| track.id),
};
let Some(track_id) = next_track_id else {
debug!(
renderer = renderer_id.0.as_str(),
"No OpenHome track available to advance to"
);
self.runtime
.set_playback_source(renderer_id, PlaybackSource::None);
return Ok(());
};
renderer.openhome_playlist_play_id(track_id)?;
self.runtime
.set_playback_source(renderer_id, PlaybackSource::FromQueue);
self.sync_openhome_playlist_for(renderer_id)?;
info!(
renderer = renderer_id.0.as_str(),
track_id, "Advanced OpenHome playlist to next track"
);
Ok(())
}
}
#[cfg(feature = "pmoserver")]
@@ -2600,6 +2520,16 @@ fn refresh_attached_queue_for(
}
};
debug!(
renderer = renderer_id.0.as_str(),
server = server_id.0.as_str(),
container = container_id.as_str(),
total_entries = entries.len(),
containers = entries.iter().filter(|e| e.is_container).count(),
items_count = entries.iter().filter(|e| !e.is_container).count(),
"Browse returned entries for playlist refresh"
);
// Step 4: Convert MediaEntry to PlaybackItem
let new_items: Vec<PlaybackItem> = entries
.iter()
@@ -2607,11 +2537,12 @@ fn refresh_attached_queue_for(
.collect();
if new_items.is_empty() {
debug!(
warn!(
renderer = renderer_id.0.as_str(),
server = server_id.0.as_str(),
container = container_id.as_str(),
"Refreshed playlist is empty, clearing queue"
total_entries = entries.len(),
"Refreshed playlist is empty, clearing queue - all entries were filtered out"
);
runtime.with_music_queue_mut(renderer_id, |queue| queue.clear_queue())?;
runtime.invalidate_openhome_cache(renderer_id);
@@ -2759,6 +2690,12 @@ fn refresh_attached_queue_for(
if auto_play {
if let Some(callback) = after_refresh.as_deref_mut() {
debug!(
renderer = renderer_id.0.as_str(),
server = server_id.0.as_str(),
container = container_id.as_str(),
"Auto-play enabled: calling callback to start playback"
);
if let Err(err) = callback(renderer_id) {
warn!(
renderer = renderer_id.0.as_str(),
@@ -2767,48 +2704,30 @@ fn refresh_attached_queue_for(
error = %err,
"Failed to auto-start playback after playlist refresh"
);
} else {
info!(
renderer = renderer_id.0.as_str(),
server = server_id.0.as_str(),
container = container_id.as_str(),
"Auto-play callback completed successfully"
);
}
} else {
debug!(
renderer = renderer_id.0.as_str(),
"Auto-play enabled but no callback provided"
);
}
} else {
debug!(
renderer = renderer_id.0.as_str(),
"Auto-play disabled, skipping playback start"
);
}
Ok(())
}
/// Helper to convert a MediaEntry to a PlaybackItem.
fn playback_item_from_entry(server: &MusicServer, entry: &MediaEntry) -> Option<PlaybackItem> {
// Ignore containers
if entry.is_container {
return None;
}
// Skip "live stream" entries (heuristic from example)
if entry.title.to_ascii_lowercase().contains("live stream") {
return None;
}
// Find an audio resource
let resource = entry.resources.iter().find(|res| res.is_audio())?;
let metadata = TrackMetadata {
title: Some(entry.title.clone()),
artist: entry.artist.clone(),
album: entry.album.clone(),
genre: entry.genre.clone(),
album_art_uri: entry.album_art_uri.clone(),
date: entry.date.clone(),
track_number: entry.track_number.clone(),
creator: entry.creator.clone(),
};
Some(PlaybackItem {
media_server_id: server.id().clone(),
didl_id: entry.id.clone(),
uri: resource.uri.clone(),
protocol_info: resource.protocol_info.clone(),
metadata: Some(metadata),
})
}
const OPENHOME_TRACK_PREFIX: &str = "openhome:";
fn playback_item_from_openhome_track(

View File

@@ -10,7 +10,7 @@ use crate::openhome_client::{
OPENHOME_PLAYLIST_HEAD_ID,
};
use crate::openhome_playlist::{OpenHomePlaylistSnapshot, OpenHomePlaylistTrack};
use crate::queue_backend::{PlaybackItem, QueueBackend, QueueSnapshot};
use crate::queue_backend::{EnqueueMode, PlaybackItem, QueueBackend, QueueSnapshot};
/// Local mirror of an OpenHome playlist for a single renderer.
#[derive(Clone, Debug)]
@@ -58,56 +58,15 @@ impl OpenHomeQueue {
track_ids.push(entry.id);
}
// Try multiple methods to determine the currently playing track, from most to least reliable:
// 1. Info.Id() - Direct ID query (fastest, but fails if track no longer in playlist)
// 2. Info.Track() - Returns URI, which we can search for (works even if track removed)
// 3. None - No current track can be determined
let current_id = self.playlist.id()?;
// {
// // Try Info.Id() first
// if let Ok(id) = client.id() {
// debug!(
// renderer = self.renderer_id.0.as_str(),
// track_id = id,
// "Detected current track via Info.Id()"
// );
// return Some(id);
// }
// Get the currently playing track ID from the renderer (may be None if no track is playing)
let current_id = self.playlist.id().ok();
// // If Id() fails, try Track() to get the URI and search for it
// if let Ok(track_info) = client.track() {
// debug!(
// renderer = self.renderer_id.0.as_str(),
// track_uri = track_info.uri.as_str(),
// "Info.Id() failed, searching for current track by URI from Info.Track()"
// );
// return entries
// .iter()
// .find(|entry| entry.uri == track_info.uri)
// .map(|entry| {
// debug!(
// renderer = self.renderer_id.0.as_str(),
// found_id = entry.id,
// found_uri = entry.uri.as_str(),
// "Found current track ID by matching URI"
// );
// entry.id
// });
// }
// debug!(
// renderer = self.renderer_id.0.as_str(),
// "Both Info.Id() and Info.Track() failed, cannot determine current track"
// );
// None
// });
// let current_index = current_id
// .and_then(|id| track_ids.iter().position(|entry_id| *entry_id == id));
// Find the index of the current track in the playlist
let current_index = current_id.and_then(|id| track_ids.iter().position(|entry_id| *entry_id == id));
self.items = items;
self.track_ids = track_ids;
self.current_index = Some(current_id as usize);
self.current_index = current_index;
Ok(())
}
@@ -163,6 +122,20 @@ impl OpenHomeQueue {
Ok(())
}
/// Selects and plays a track by its queue index (0-based).
pub fn select_track_index(&mut self, index: usize) -> Result<()> {
let track_id = self
.track_ids
.get(index)
.copied()
.ok_or_else(|| anyhow!("Index {} out of bounds (queue length: {})", index, self.track_ids.len()))?;
self.ensure_playlist_source_selected()?;
self.playlist.play_id(track_id)?;
self.current_index = Some(index);
Ok(())
}
pub fn clear(&mut self) -> Result<()> {
self.ensure_playlist_source_selected()?;
self.playlist.delete_all()?;
@@ -325,34 +298,10 @@ impl OpenHomeQueue {
playing_idx: usize,
playing_id: u32,
) -> Result<()> {
// Re-read the current playlist state to get fresh IDs
// This minimizes race conditions where IDs become invalid between our last refresh
// and now (due to UPnP events from the server)
let current_entries = self.playlist.read_all_tracks()?;
let current_ids: Vec<u32> = current_entries.iter().map(|e| e.id).collect();
debug!(
renderer = self.renderer_id.0.as_str(),
fresh_id_count = current_ids.len(),
cached_id_count = self.track_ids.len(),
"Re-read playlist before deletions to avoid stale ID errors"
);
// Find the playing track in the fresh list
let fresh_playing_idx = current_ids.iter().position(|&id| id == playing_id);
if fresh_playing_idx.is_none() {
debug!(
renderer = self.renderer_id.0.as_str(),
playing_id,
"Playing track not found in fresh playlist - renderer state may have changed, aborting modification"
);
// The playing track is gone - don't try to manipulate the playlist
return Ok(());
}
// Delete everything except the currently playing item (using fresh IDs)
for &track_id in current_ids.iter().rev() {
// Delete everything except the currently playing item
// Using delete_id_if_exists() to handle cases where another control point
// may have already modified the playlist
for &track_id in self.track_ids.iter().rev() {
if track_id != playing_id {
self.playlist.delete_id_if_exists(track_id)?;
}
@@ -395,53 +344,15 @@ impl OpenHomeQueue {
pivot_idx_new: usize,
pivot_id: u32,
) -> Result<()> {
debug!(
renderer = self.renderer_id.0.as_str(),
pivot_id,
pivot_idx_new,
new_playlist_len = new_items.len(),
"Starting replace_queue_with_pivot - will re-read from OpenHome"
);
// Find the pivot index in our current state
let pivot_idx = self.track_ids.iter().position(|&id| id == pivot_id)
.ok_or_else(|| anyhow!("Pivot track ID {} not found in playlist", pivot_id))?;
// Re-read the current playlist state from OpenHome (the ONLY source of truth)
// This is CRITICAL to avoid deleting IDs that no longer exist, which can
// put the renderer (upmpdcli) into a degraded state where Info.TransportState()
// starts returning HTTP 500 errors.
let current_entries = self.playlist.read_all_tracks()?;
// Convert entries to PlaybackItems - this is the REAL current state
let mut fresh_items = Vec::with_capacity(current_entries.len());
let mut fresh_ids = Vec::with_capacity(current_entries.len());
for entry in &current_entries {
fresh_items.push(self.playback_item_from_entry(entry));
fresh_ids.push(entry.id);
}
debug!(
renderer = self.renderer_id.0.as_str(),
fresh_count = fresh_items.len(),
"Re-read playlist from OpenHome (source of truth)"
);
// Find the pivot in the fresh list
let fresh_pivot_idx = fresh_ids.iter().position(|&id| id == pivot_id);
if fresh_pivot_idx.is_none() {
debug!(
renderer = self.renderer_id.0.as_str(),
pivot_id,
"Pivot track not found in fresh playlist - renderer state changed, aborting"
);
return Ok(());
}
let fresh_pivot_idx = fresh_pivot_idx.unwrap();
// Split fresh data at the pivot - use ONLY fresh data, ignore cache
let old_before: Vec<PlaybackItem> = fresh_items[..fresh_pivot_idx].to_vec();
let old_after: Vec<PlaybackItem> = fresh_items[fresh_pivot_idx + 1..].to_vec();
let old_ids_before: Vec<u32> = fresh_ids[..fresh_pivot_idx].to_vec();
let old_ids_after: Vec<u32> = fresh_ids[fresh_pivot_idx + 1..].to_vec();
// Split current data at the pivot
let old_before: Vec<PlaybackItem> = self.items[..pivot_idx].to_vec();
let old_after: Vec<PlaybackItem> = self.items[pivot_idx + 1..].to_vec();
let old_ids_before: Vec<u32> = self.track_ids[..pivot_idx].to_vec();
let old_ids_after: Vec<u32> = self.track_ids[pivot_idx + 1..].to_vec();
let new_before = &new_items[..pivot_idx_new];
let new_after = &new_items[pivot_idx_new + 1..];
@@ -841,43 +752,17 @@ impl QueueBackend for OpenHomeQueue {
// (e.g., manual edits from another control point) would keep the stale items.
self.refresh_from_openhome()?;
// Try to get the currently playing track ID from the renderer.
// Note: Some OpenHome renderers (like upmpdcli) don't reliably support Info.Id(),
// so we fall back to using our internal current_index pointer.
let currently_playing_id_from_renderer = self
.playlist.id().ok();
// Find the currently playing item in our local state.
// Priority: 1) Renderer-reported ID, 2) Our internal current_index
let playing_info = if let Some(id) = currently_playing_id_from_renderer {
// CASE: Renderer explicitly reported the playing track ID
// Get the currently playing track ID from the renderer
let playing_info = self.playlist.id().ok().and_then(|id| {
self.track_ids
.iter()
.position(|&tid| tid == id)
.map(|idx| (idx, id, self.items[idx].uri.clone()))
} else if let Some(idx) = self.current_index {
// CASE: Use our internal pointer (fallback for renderers without Info.Id() support)
if idx < self.track_ids.len() && idx < self.items.len() {
let id = self.track_ids[idx];
let uri = self.items[idx].uri.clone();
debug!(
renderer = self.renderer_id.0.as_str(),
current_index = idx,
track_id = id,
"Using internal current_index as fallback (renderer didn't report playing ID)"
);
Some((idx, id, uri))
} else {
None
}
} else {
None
};
});
debug!(
renderer = self.renderer_id.0.as_str(),
actual_items = self.items.len(),
currently_playing_id_from_renderer = ?currently_playing_id_from_renderer,
playing_info_detected = playing_info.is_some(),
"OpenHome playlist state refreshed before replace_queue"
);
@@ -952,4 +837,43 @@ impl QueueBackend for OpenHomeQueue {
self.track_ids[index] = new_id;
Ok(())
}
/// Override enqueue_items to add items directly to the OpenHome playlist.
fn enqueue_items(&mut self, items: Vec<PlaybackItem>, mode: EnqueueMode) -> Result<()> {
if items.is_empty() {
return Ok(());
}
match mode {
EnqueueMode::AppendToEnd => {
// Append to the end of the OpenHome playlist
let mut after_id = self.track_ids.last().copied();
for item in items {
after_id = Some(self.add_playback_item(item, after_id, false)?);
}
}
EnqueueMode::InsertAfterCurrent => {
// Insert after the current playing track
let after_id = if let Some(idx) = self.current_index {
self.track_ids.get(idx).copied()
} else {
None
};
let mut next_after_id = after_id;
for item in items {
next_after_id = Some(self.add_playback_item(item, next_after_id, false)?);
}
}
EnqueueMode::ReplaceAll => {
// Replace the entire playlist
self.replace_queue(items, None)?;
}
}
// Refresh local cache from OpenHome after modification
self.refresh_from_openhome()?;
Ok(())
}
}

View File

@@ -4,8 +4,11 @@ use anyhow::{Result, anyhow};
use pmodidl::{self, DIDLLite};
use pmoupnp::soap::SoapEnvelope;
use pmoupnp::soap::error_codes;
use tracing::{debug, warn};
use xmltree::{Element, XMLNode};
use crate::model::TrackMetadata;
use crate::queue_backend::PlaybackItem;
use crate::soap_client::{SoapCallResult, invoke_upnp_action_with_timeout};
/// Unique identifier for a media server registered by the control point.
@@ -42,15 +45,49 @@ impl MediaResource {
/// Returns true if this resource represents audio content.
pub fn is_audio(&self) -> bool {
let lower = self.protocol_info.to_ascii_lowercase();
// Standard case: audio/* MIME types
if lower.contains("audio/") {
return true;
}
// List of known audio format subtypes (the part after the /)
// These are recognized regardless of the MIME type prefix
const AUDIO_FORMATS: &[&str] = &[
"flac", "ogg", "opus", "vorbis",
"mp3", "mpeg", "mp4", "m4a", "aac",
"wav", "wave", "pcm",
"wma", "webm",
"ape", "alac", "aiff",
"dsd", "dsf", "dff",
];
// Check if any known audio format appears in the protocol_info
for format in AUDIO_FORMATS {
if lower.contains(format) {
return true;
}
}
// protocolInfo format: protocol:network:contentFormat:additionalInfo
lower
.split(':')
.nth(2)
.map(|mime| mime.starts_with("audio/"))
.unwrap_or(false)
// Extract the MIME type (3rd field) for more precise checking
if let Some(mime) = lower.split(':').nth(2) {
// Check if it's audio/* or contains a known audio format
if mime.starts_with("audio/") {
return true;
}
// Check the subtype (part after /) for known audio formats
if let Some(subtype) = mime.split('/').nth(1) {
for format in AUDIO_FORMATS {
if subtype.contains(format) {
return true;
}
}
}
}
false
}
}
@@ -145,6 +182,80 @@ impl MediaBrowser for MusicServer {
}
}
/// Helper to convert a MediaEntry to a PlaybackItem.
///
/// This function filters out containers and entries without audio resources,
/// returning None for items that cannot be played.
pub fn playback_item_from_entry(server: &MusicServer, entry: &MediaEntry) -> Option<PlaybackItem> {
// Ignore containers
if entry.is_container {
debug!(
server_id = server.id().0.as_str(),
entry_id = entry.id.as_str(),
title = entry.title.as_str(),
class = entry.class.as_str(),
"Skipping container entry"
);
return None;
}
// Skip "live stream" entries (heuristic from example)
if entry.title.to_ascii_lowercase().contains("live stream") {
debug!(
server_id = server.id().0.as_str(),
entry_id = entry.id.as_str(),
title = entry.title.as_str(),
"Skipping 'live stream' entry"
);
return None;
}
// Find an audio resource
let resource = entry.resources.iter().find(|res| res.is_audio());
if resource.is_none() {
warn!(
server_id = server.id().0.as_str(),
entry_id = entry.id.as_str(),
title = entry.title.as_str(),
class = entry.class.as_str(),
resource_count = entry.resources.len(),
resources = ?entry.resources.iter().map(|r| &r.protocol_info).collect::<Vec<_>>(),
"No audio resource found for entry"
);
return None;
}
let resource = resource.unwrap();
let metadata = TrackMetadata {
title: Some(entry.title.clone()),
artist: entry.artist.clone(),
album: entry.album.clone(),
genre: entry.genre.clone(),
album_art_uri: entry.album_art_uri.clone(),
date: entry.date.clone(),
track_number: entry.track_number.clone(),
creator: entry.creator.clone(),
};
debug!(
server_id = server.id().0.as_str(),
entry_id = entry.id.as_str(),
title = entry.title.as_str(),
uri = resource.uri.as_str(),
"Created playback item"
);
Some(PlaybackItem {
media_server_id: server.id().clone(),
didl_id: entry.id.clone(),
uri: resource.uri.clone(),
protocol_info: resource.protocol_info.clone(),
metadata: Some(metadata),
})
}
/// Single UPnP ContentDirectory backend implementation.
#[derive(Clone, Debug)]
pub struct UpnpMediaServer {

View File

@@ -13,16 +13,16 @@ use crate::control_point::RendererRuntimeStateMut;
use crate::control_point::music_queue::MusicQueue;
use crate::control_point::openhome_queue::didl_id_from_metadata;
use crate::media_server::ServerId;
use crate::model::{RendererId, RendererInfo, RendererProtocol};
use crate::model::{RendererId, RendererInfo, RendererProtocol, TrackMetadata};
use crate::openhome_client::parse_track_metadata_from_didl;
use crate::openhome_playlist::OpenHomePlaylistSnapshot;
use crate::queue_backend::PlaybackItem;
use crate::openhome_playlist::{OpenHomePlaylistSnapshot, OpenHomePlaylistTrack};
use crate::queue_backend::{PlaybackItem, QueueSnapshot};
use crate::{
ArylicTcpRenderer, ChromecastRenderer, DeviceRegistry, LinkPlayRenderer, OpenHomeRenderer,
PlaybackPosition, PlaybackState, TransportControl, UpnpRenderer, VolumeControl,
};
use anyhow::{Result, anyhow};
use tracing::warn;
use tracing::{debug, info, warn};
/// Backend-agnostic façade exposing transport, volume, and status contracts.
#[derive(Clone, Debug)]
@@ -269,6 +269,321 @@ impl MusicRenderer {
}
}
/// High-level method to prepare the renderer for attaching a new playlist.
///
/// This method handles backend-specific clearing logic:
/// - For OpenHome: clears the OpenHome playlist
/// - For AVTransport/Chromecast/etc.: stops the renderer (since they don't have a persistent queue)
///
/// This should be called by ControlPoint when attaching a new playlist, ensuring that:
/// - Any currently playing content is stopped
/// - The renderer is in a clean state ready to receive new content
pub fn clear_for_playlist_attach(&self) -> Result<()> {
match self {
MusicRenderer::OpenHome(_) => {
// For OpenHome: clear the playlist (DeleteAll also stops playback automatically)
// then explicitly stop to ensure clean state
self.openhome_playlist_clear()?;
self.stop().or_else(|err| -> Result<()> {
// If stop fails (e.g., already stopped), that's fine
warn!(
renderer = self.id().0.as_str(),
error = %err,
"Stop failed after clearing OpenHome playlist (continuing anyway)"
);
Ok(())
})
}
MusicRenderer::Upnp(_)
| MusicRenderer::Chromecast(_)
| MusicRenderer::LinkPlay(_)
| MusicRenderer::ArylicTcp(_)
| MusicRenderer::HybridUpnpArylic { .. } => {
// For AVTransport and other single-track renderers: stop playback
// This ensures we're not in the middle of playing when we start the new playlist
self.stop().or_else(|err| {
// If stop fails (e.g., already stopped), that's fine - we just want to ensure it's not playing
warn!(
renderer = self.id().0.as_str(),
error = %err,
"Stop failed when preparing for playlist attach (continuing anyway)"
);
Ok(())
})
}
}
}
/// Synchronize the local queue state with the backend's actual state.
///
/// - For OpenHome: fetches the playlist from the renderer and updates local cache
/// - For Internal queue/AVTransport: no-op (queue is already local)
pub fn sync_queue_state(&self) -> Result<()> {
match self {
MusicRenderer::OpenHome(_) => {
if let Some(provider) = OPENHOME_QUEUE_PROVIDER.get() {
// Fetch fresh playlist snapshot and update cache
provider.invalidate_openhome_cache(self.id())?;
}
Ok(())
}
MusicRenderer::Upnp(_)
| MusicRenderer::Chromecast(_)
| MusicRenderer::LinkPlay(_)
| MusicRenderer::ArylicTcp(_)
| MusicRenderer::HybridUpnpArylic { .. } => {
// No sync needed - queue is local only
Ok(())
}
}
}
/// Clear the queue on both the backend and in local state.
///
/// This ensures the backend renderer and local cache are consistent.
/// Should be called before queue mutations to ensure clean state.
pub fn clear_queue(&self) -> Result<()> {
match self {
MusicRenderer::OpenHome(_) => {
// For OpenHome: clear the playlist on the renderer itself
self.openhome_playlist_clear()
}
MusicRenderer::Upnp(_)
| MusicRenderer::Chromecast(_)
| MusicRenderer::LinkPlay(_)
| MusicRenderer::ArylicTcp(_)
| MusicRenderer::HybridUpnpArylic { .. } => {
// For other renderers: no persistent queue to clear
Ok(())
}
}
}
/// Add a track to the backend's queue, returning backend-specific track ID if applicable.
///
/// - For OpenHome: adds to OpenHome playlist and returns track ID
/// - For others: returns error (not supported for single-track renderers)
pub fn add_track_to_queue(
&self,
uri: &str,
metadata: &str,
after_id: Option<u32>,
play: bool,
) -> Result<Option<u32>> {
match self {
MusicRenderer::OpenHome(_) => {
let track_id = self.openhome_playlist_add_track(uri, metadata, after_id, play)?;
Ok(Some(track_id))
}
MusicRenderer::Upnp(_)
| MusicRenderer::Chromecast(_)
| MusicRenderer::LinkPlay(_)
| MusicRenderer::ArylicTcp(_)
| MusicRenderer::HybridUpnpArylic { .. } => Err(anyhow!(
"add_track_to_queue is not supported for {} backend",
self.unsupported_backend_name()
)),
}
}
/// Select and play a specific track from the backend's queue.
///
/// - For OpenHome: uses track ID to select from OpenHome playlist
/// - For others: returns error (use play_uri instead)
pub fn select_queue_track(&self, track_id: u32) -> Result<()> {
match self {
MusicRenderer::OpenHome(_) => self.openhome_playlist_play_id(track_id),
MusicRenderer::Upnp(_)
| MusicRenderer::Chromecast(_)
| MusicRenderer::LinkPlay(_)
| MusicRenderer::ArylicTcp(_)
| MusicRenderer::HybridUpnpArylic { .. } => Err(anyhow!(
"select_queue_track is not supported for {} backend",
self.unsupported_backend_name()
)),
}
}
/// Get the current queue state from the backend.
///
/// - For OpenHome: fetches current playlist snapshot
/// - For others: returns None (no persistent queue on backend)
pub fn queue_snapshot(&self) -> Result<Option<QueueSnapshot>> {
match self {
MusicRenderer::OpenHome(_) => {
let oh_snapshot = self.fetch_openhome_playlist_snapshot()?;
// Convert OpenHome tracks to PlaybackItems
let items: Vec<PlaybackItem> = oh_snapshot.tracks.iter().map(|track| {
Self::playback_item_from_openhome_track(self.id(), track)
}).collect();
let snapshot = QueueSnapshot {
items,
current_index: oh_snapshot.current_index,
};
Ok(Some(snapshot))
}
MusicRenderer::Upnp(_)
| MusicRenderer::Chromecast(_)
| MusicRenderer::LinkPlay(_)
| MusicRenderer::ArylicTcp(_)
| MusicRenderer::HybridUpnpArylic { .. } => {
// No backend queue for these renderers
Ok(None)
}
}
}
/// Play the current item from the backend queue.
///
/// - For OpenHome: Uses the native playlist to play the current track
/// - For others: Returns error (no backend queue)
pub fn play_current_from_backend_queue(&self) -> Result<()> {
match self {
MusicRenderer::OpenHome(_) => {
// Get the current OpenHome playlist snapshot
let snapshot = self.fetch_openhome_playlist_snapshot()?;
debug!(
renderer = self.id().0.as_str(),
tracks_count = snapshot.tracks.len(),
current_id = ?snapshot.current_id,
current_index = ?snapshot.current_index,
"play_current_from_backend_queue: OpenHome snapshot fetched"
);
if snapshot.tracks.is_empty() {
return Err(anyhow!("OpenHome playlist is empty"));
}
// Find the track_id to play (prefer current_id, then current_index, then first)
let target_track_id = if let Some(current_id) = snapshot.current_id {
debug!(
renderer = self.id().0.as_str(),
track_id = current_id,
"Using current_id for playback"
);
Some(current_id)
} else if let Some(current_idx) = snapshot.current_index {
let track_id = snapshot.tracks.get(current_idx).map(|track| track.id);
debug!(
renderer = self.id().0.as_str(),
current_idx,
track_id = ?track_id,
"Using current_index for playback"
);
track_id
} else {
let track_id = snapshot.tracks.first().map(|track| track.id);
debug!(
renderer = self.id().0.as_str(),
track_id = ?track_id,
"Using first track for playback"
);
track_id
};
if let Some(track_id) = target_track_id {
info!(
renderer = self.id().0.as_str(),
track_id,
"Calling openhome_playlist_play_id to start playback"
);
self.openhome_playlist_play_id(track_id)?;
info!(
renderer = self.id().0.as_str(),
track_id,
"Successfully called openhome_playlist_play_id"
);
Ok(())
} else {
Err(anyhow!("No track to play in OpenHome playlist"))
}
}
MusicRenderer::Upnp(_)
| MusicRenderer::Chromecast(_)
| MusicRenderer::LinkPlay(_)
| MusicRenderer::ArylicTcp(_)
| MusicRenderer::HybridUpnpArylic { .. } => Err(anyhow!(
"play_current_from_backend_queue is not supported for {} backend (no persistent queue)",
self.unsupported_backend_name()
)),
}
}
/// Play the next item from the backend queue.
///
/// - For OpenHome: Advances to the next track in the playlist
/// - For others: Returns error (no backend queue)
pub fn play_next_from_backend_queue(&self) -> Result<()> {
match self {
MusicRenderer::OpenHome(_) => {
// Get the current OpenHome playlist snapshot
let snapshot = self.fetch_openhome_playlist_snapshot()?;
if snapshot.tracks.is_empty() {
return Err(anyhow!("OpenHome playlist is empty, cannot play next"));
}
// Determine the next track_id
let next_track_id = match snapshot.current_index {
Some(idx) => {
// Take the next track if it exists, otherwise loop to first
snapshot
.tracks
.get(idx + 1)
.map(|track| track.id)
.or_else(|| snapshot.tracks.first().map(|track| track.id))
}
None => snapshot.tracks.first().map(|track| track.id),
};
if let Some(track_id) = next_track_id {
self.openhome_playlist_play_id(track_id)?;
Ok(())
} else {
Err(anyhow!("No track available to advance to in OpenHome playlist"))
}
}
MusicRenderer::Upnp(_)
| MusicRenderer::Chromecast(_)
| MusicRenderer::LinkPlay(_)
| MusicRenderer::ArylicTcp(_)
| MusicRenderer::HybridUpnpArylic { .. } => Err(anyhow!(
"play_next_from_backend_queue is not supported for {} backend (no persistent queue)",
self.unsupported_backend_name()
)),
}
}
/// Convert an OpenHome playlist track to a PlaybackItem.
fn playback_item_from_openhome_track(
renderer_id: &RendererId,
track: &OpenHomePlaylistTrack,
) -> PlaybackItem {
let metadata = TrackMetadata {
title: track.title.clone(),
artist: track.artist.clone(),
album: track.album.clone(),
genre: None,
album_art_uri: track.album_art_uri.clone(),
date: None,
track_number: None,
creator: None,
};
PlaybackItem {
media_server_id: ServerId(format!("openhome:{}", renderer_id.0)),
didl_id: format!("openhome:{}", track.id),
uri: track.uri.clone(),
// OpenHome tracks don't provide protocolInfo, use generic default
protocol_info: "http-get:*:audio/*:*".to_string(),
metadata: Some(metadata),
}
}
pub fn openhome_playlist_add_track(
&self,
uri: &str,

View File

@@ -280,6 +280,14 @@ pub struct PlayContentRequest {
pub object_id: String,
}
/// Requête pour sauter à un index spécifique dans la queue
#[cfg(feature = "pmoserver")]
#[derive(Debug, Clone, Deserialize, ToSchema)]
pub struct SeekQueueRequest {
/// Index de l'item dans la queue (0-based)
pub index: usize,
}
/// Réponse générique de succès
#[cfg(feature = "pmoserver")]
#[derive(Debug, Clone, Serialize, ToSchema)]

View File

@@ -146,10 +146,10 @@ impl OhPlaylistClient {
let envelope = ensure_success("Id", &call_result)?;
let response = find_child_with_suffix(&envelope.body.content, "IdResponse")
.ok_or_else(|| anyhow!("Missing IdResponse element in SOAP body"))?;
let id_text = extract_child_text(response, "Id")?;
let id_text = extract_child_text(response, "Value")?;
let id = id_text
.parse::<u32>()
.map_err(|_| anyhow!("Invalid Info.Id value: {}", id_text))?;
.map_err(|_| anyhow!("Invalid Playlist.Id value: {}", id_text))?;
Ok(id)
}

View File

@@ -118,63 +118,8 @@ impl OpenHomeRenderer {
let playlist = self.playlist_client_for("snapshot_openhome_playlist")?;
let entries = playlist.read_all_tracks()?;
// Essayer d'obtenir current_id depuis Info.Id()
let mut current_id = self
.playlist
.as_ref()
.and_then(|client| {
match client.id() {
Ok(id) => {
debug!(
renderer = self.info.id.0.as_str(),
current_id = id,
"OpenHome Info service returned current_id"
);
Some(id)
}
Err(err) => {
debug!(
renderer = self.info.id.0.as_str(),
error = %err,
"OpenHome Info.Id() failed, will try Info.Track()"
);
None
}
}
});
// Fallback: Si Info.Id() échoue, essayer Info.Track() et matcher l'URI
if current_id.is_none() {
if let Some(client) = self.info_client.as_ref() {
match client.track() {
Ok(track_info) => {
debug!(
renderer = self.info.id.0.as_str(),
track_uri = track_info.uri.as_str(),
"OpenHome Info.Track() returned, searching by URI"
);
// Trouver l'entry qui matche cet URI
current_id = entries.iter()
.find(|entry| entry.uri == track_info.uri)
.map(|entry| {
debug!(
renderer = self.info.id.0.as_str(),
found_id = entry.id,
"Found current_id by matching URI"
);
entry.id
});
}
Err(err) => {
debug!(
renderer = self.info.id.0.as_str(),
error = %err,
"OpenHome Info.Track() also failed"
);
}
}
}
}
// Get current track ID from the playlist service
let current_id = playlist.id().ok();
let current_index =
current_id.and_then(|id| entries.iter().position(|entry| entry.id == id));

View File

@@ -8,7 +8,9 @@ use crate::control_point::{
ControlPoint, OpenHomeAccessError, OPENHOME_SNAPSHOT_CACHE_TTL,
};
#[cfg(feature = "pmoserver")]
use crate::media_server::{MediaBrowser, MediaEntry, MusicServer, ServerId};
use crate::media_server::{
playback_item_from_entry, MediaBrowser, MediaEntry, MusicServer, ServerId,
};
#[cfg(feature = "pmoserver")]
use crate::model::{RendererCapabilities, RendererId, RendererProtocol, TrackMetadata};
#[cfg(feature = "pmoserver")]
@@ -16,7 +18,8 @@ use crate::openapi::{
AttachPlaylistRequest, AttachedPlaylistInfo, BrowseResponse, ContainerEntry, ErrorResponse,
FullRendererSnapshot, MediaServerSummary, OpenHomePlaylistAddRequest, OpenHomePlaylistSnapshot,
PlayContentRequest, QueueItem, QueueSnapshot, RendererCapabilitiesSummary,
RendererProtocolSummary, RendererState, RendererSummary, SuccessResponse, VolumeSetRequest,
RendererProtocolSummary, RendererState, RendererSummary, SeekQueueRequest, SuccessResponse,
VolumeSetRequest,
};
#[cfg(feature = "pmoserver")]
use crate::queue_backend::PlaybackItem;
@@ -621,6 +624,89 @@ async fn next_renderer(
}))
}
/// POST /control/renderers/{renderer_id}/queue/seek - Saute à un index spécifique dans la queue
#[cfg(feature = "pmoserver")]
#[utoipa::path(
post,
path = "/renderers/{renderer_id}/queue/seek",
tag = "control",
request_body = SeekQueueRequest,
responses(
(status = 200, description = "Lecture démarrée à l'index spécifié", body = SuccessResponse),
(status = 404, description = "Renderer non trouvé", body = ErrorResponse),
(status = 400, description = "Index invalide", body = ErrorResponse),
(status = 500, description = "Erreur interne", body = ErrorResponse),
)
)]
async fn seek_queue_index(
State(state): State<ControlPointState>,
Path(renderer_id): Path<String>,
Json(payload): Json<SeekQueueRequest>,
) -> Result<Json<SuccessResponse>, (StatusCode, Json<ErrorResponse>)> {
let rid = RendererId(renderer_id.clone());
state
.control_point
.music_renderer_by_id(&rid)
.ok_or_else(|| {
(
StatusCode::NOT_FOUND,
Json(ErrorResponse {
error: format!("Renderer {} not found", renderer_id),
}),
)
})?;
let control_point = Arc::clone(&state.control_point);
let rid_for_task = rid.clone();
let index = payload.index;
let seek_task = tokio::task::spawn_blocking(move || {
control_point.play_queue_index(&rid_for_task, index)
});
time::timeout(TRANSPORT_COMMAND_TIMEOUT, seek_task)
.await
.map_err(|_| {
warn!(
"Seek queue command for renderer {} exceeded {:?}",
renderer_id, TRANSPORT_COMMAND_TIMEOUT
);
(
StatusCode::GATEWAY_TIMEOUT,
Json(ErrorResponse {
error: format!(
"Seek command timed out after {}s",
TRANSPORT_COMMAND_TIMEOUT.as_secs()
),
}),
)
})?
.map_err(|e| {
warn!("Task join error during queue seek: {}", e);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Internal task error: {}", e),
}),
)
})?
.map_err(|e| {
warn!(
"Failed to seek to index {} for renderer {}: {}",
index, renderer_id, e
);
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(ErrorResponse {
error: format!("Failed to seek to index {}: {}", index, e),
}),
)
})?;
Ok(Json(SuccessResponse {
message: format!("Playing item at index {}", index),
}))
}
/// POST /control/renderers/{renderer_id}/volume/set - Définit le volume
#[cfg(feature = "pmoserver")]
#[utoipa::path(
@@ -1157,7 +1243,7 @@ async fn clear_openhome_playlist(
let rid_for_task = rid.clone();
let clear_task =
tokio::task::spawn_blocking(move || control_point.clear_openhome_playlist(&rid_for_task));
tokio::task::spawn_blocking(move || control_point.clear_renderer_queue(&rid_for_task));
time::timeout(QUEUE_COMMAND_TIMEOUT, clear_task)
.await
@@ -1229,13 +1315,14 @@ async fn add_openhome_playlist_item(
let rid_for_task = rid.clone();
let add_task = tokio::task::spawn_blocking(move || {
control_point.add_openhome_track(
control_point.add_track_to_renderer(
&rid_for_task,
&req.uri,
&req.metadata,
req.after_id,
req.play,
)
.map(|_| ()) // Ignore the track_id result for backward compatibility
});
time::timeout(QUEUE_COMMAND_TIMEOUT, add_task)
@@ -1316,7 +1403,7 @@ async fn play_openhome_track(
let rid_for_task = rid.clone();
let play_task = tokio::task::spawn_blocking(move || {
control_point.play_openhome_track_id(&rid_for_task, parsed_id)
control_point.select_renderer_track(&rid_for_task, parsed_id)
});
time::timeout(QUEUE_COMMAND_TIMEOUT, play_task)
@@ -1840,51 +1927,33 @@ fn fetch_playback_items(
// Browse the object to get entries
let entries = music_server.browse_children(object_id, 0, BROWSE_PAGE_SIZE)?;
debug!(
server_id = server_id.0.as_str(),
object_id = object_id,
total_entries = entries.len(),
containers = entries.iter().filter(|e| e.is_container).count(),
items_count = entries.iter().filter(|e| !e.is_container).count(),
"Browse returned entries"
);
// Convert to PlaybackItem
let items: Vec<PlaybackItem> = entries
.iter()
.filter_map(|entry| playback_item_from_entry(&music_server, entry))
.collect();
if items.is_empty() && !entries.is_empty() {
warn!(
server_id = server_id.0.as_str(),
object_id = object_id,
total_entries = entries.len(),
"No playable items found - all entries were filtered out"
);
}
Ok(items)
}
/// Helper to convert a MediaEntry to a PlaybackItem.
#[cfg(feature = "pmoserver")]
fn playback_item_from_entry(server: &MusicServer, entry: &MediaEntry) -> Option<PlaybackItem> {
// Ignore containers
if entry.is_container {
return None;
}
// Skip "live stream" entries
if entry.title.to_ascii_lowercase().contains("live stream") {
return None;
}
// Find an audio resource
let resource = entry.resources.iter().find(|res| res.is_audio())?;
let metadata = TrackMetadata {
title: Some(entry.title.clone()),
artist: entry.artist.clone(),
album: entry.album.clone(),
genre: entry.genre.clone(),
album_art_uri: entry.album_art_uri.clone(),
date: entry.date.clone(),
track_number: entry.track_number.clone(),
creator: entry.creator.clone(),
};
Some(PlaybackItem {
media_server_id: server.id().clone(),
didl_id: entry.id.clone(),
uri: resource.uri.clone(),
protocol_info: resource.protocol_info.clone(),
metadata: Some(metadata),
})
}
#[cfg(feature = "pmoserver")]
fn protocol_summary(protocol: &RendererProtocol) -> RendererProtocolSummary {
match protocol {
@@ -1939,6 +2008,8 @@ pub fn create_api_router(state: ControlPointState, control_point: Arc<ControlPoi
.route("/renderers/{renderer_id}/stop", post(stop_renderer))
.route("/renderers/{renderer_id}/resume", post(resume_renderer))
.route("/renderers/{renderer_id}/next", post(next_renderer))
// Queue control
.route("/renderers/{renderer_id}/queue/seek", post(seek_queue_index))
// Volume control
.route(
"/renderers/{renderer_id}/volume/set",

View File

@@ -19,6 +19,8 @@ use crate::control_point::ControlPoint;
#[cfg(feature = "pmoserver")]
use crate::model::{MediaServerEvent, RendererEvent};
#[cfg(feature = "pmoserver")]
use crate::registry::DeviceRegistryRead;
#[cfg(feature = "pmoserver")]
use async_stream::stream;
#[cfg(feature = "pmoserver")]
use axum::{
@@ -168,6 +170,32 @@ pub async fn renderer_events_sse(
});
let stream = stream! {
// INITIAL SNAPSHOT: Send Online events for all currently discovered renderers
// This ensures clients see devices that were discovered before they connected
let initial_renderers = {
let registry = control_point.registry();
let reg = registry.read().unwrap();
reg.list_renderers()
};
for info in initial_renderers {
if info.online {
let timestamp = chrono::Utc::now();
let payload = RendererEventPayload::Online {
renderer_id: info.id.0.clone(),
friendly_name: info.friendly_name.clone(),
model_name: info.model_name.clone(),
manufacturer: info.manufacturer.clone(),
timestamp,
};
if let Ok(json) = serde_json::to_string(&payload) {
yield Ok::<_, axum::Error>(Event::default().event("renderer").data(json));
}
}
}
// Then stream future events
while let Some(event) = rx_tokio.recv().await {
let timestamp = chrono::Utc::now();
@@ -285,6 +313,32 @@ pub async fn media_server_events_sse(
});
let stream = stream! {
// INITIAL SNAPSHOT: Send Online events for all currently discovered media servers
// This ensures clients see servers that were discovered before they connected
let initial_servers = {
let registry = control_point.registry();
let reg = registry.read().unwrap();
reg.list_servers()
};
for info in initial_servers {
if info.online {
let timestamp = chrono::Utc::now();
let payload = MediaServerEventPayload::Online {
server_id: info.id.0.clone(),
friendly_name: info.friendly_name.clone(),
model_name: info.model_name.clone(),
manufacturer: info.manufacturer.clone(),
timestamp,
};
if let Ok(json) = serde_json::to_string(&payload) {
yield Ok::<_, axum::Error>(Event::default().event("media_server").data(json));
}
}
}
// Then stream future events
while let Some(event) = rx_tokio.recv().await {
let timestamp = chrono::Utc::now();
@@ -370,6 +424,53 @@ pub async fn all_events_sse(State(control_point): State<Arc<ControlPoint>>) -> i
});
let stream = stream! {
// INITIAL SNAPSHOT: Send Online events for all currently discovered devices
// This ensures clients see devices that were discovered before they connected
let (initial_renderers, initial_servers) = {
let registry = control_point.registry();
let reg = registry.read().unwrap();
(reg.list_renderers(), reg.list_servers())
};
// Send renderer Online events
for info in initial_renderers {
if info.online {
let timestamp = chrono::Utc::now();
let renderer_payload = RendererEventPayload::Online {
renderer_id: info.id.0.clone(),
friendly_name: info.friendly_name.clone(),
model_name: info.model_name.clone(),
manufacturer: info.manufacturer.clone(),
timestamp,
};
let payload = UnifiedEventPayload::Renderer(renderer_payload);
if let Ok(json) = serde_json::to_string(&payload) {
yield Ok::<_, axum::Error>(Event::default().event("control").data(json));
}
}
}
// Send server Online events
for info in initial_servers {
if info.online {
let timestamp = chrono::Utc::now();
let server_payload = MediaServerEventPayload::Online {
server_id: info.id.0.clone(),
friendly_name: info.friendly_name.clone(),
model_name: info.model_name.clone(),
manufacturer: info.manufacturer.clone(),
timestamp,
};
let payload = UnifiedEventPayload::MediaServer(server_payload);
if let Ok(json) = serde_json::to_string(&payload) {
yield Ok::<_, axum::Error>(Event::default().event("control").data(json));
}
}
}
// Then stream future events
loop {
tokio::select! {
Some(event) = renderer_rx_tokio.recv() => {

View File

@@ -41,16 +41,13 @@ use pmoupnp::devices::Device;
/// }
/// ```
pub static MEDIA_RENDERER: Lazy<Arc<Device>> = Lazy::new(|| {
let mut device = Device::new(
let mut device = Device::new_from_config(
"PMO_MediaRenderer".to_string(),
"MediaRenderer".to_string(),
"PMOMusic Audio Renderer".to_string(),
"Audio Renderer".to_string(),
);
device.set_manufacturer("PMOMusic".to_string());
device.set_model_name("PMOMusic Audio Renderer".to_string());
device.set_model_description("UPnP AV MediaRenderer for audio streaming".to_string());
device.set_udn_prefix("pmomusic".to_string());
// Ajouter les trois services obligatoires
device

View File

@@ -37,16 +37,13 @@ use pmoupnp::devices::Device;
/// }
/// ```
pub static MEDIA_SERVER: Lazy<Arc<Device>> = Lazy::new(|| {
let mut device = Device::new(
let mut device = Device::new_from_config(
"PMO_MediaServer".to_string(),
"MediaServer".to_string(),
"PMOMusic Media Server".to_string(),
"Media Server".to_string(),
);
device.set_manufacturer("PMOMusic".to_string());
device.set_model_name("PMOMusic Media Server".to_string());
device.set_model_description("UPnP AV MediaServer for audio streaming".to_string());
device.set_udn_prefix("pmomusic".to_string());
// Ajouter les deux services obligatoires
device

View File

@@ -127,13 +127,16 @@ impl SourcesExt for Server {
tracing::info!("Initializing Qobuz source...");
// Obtenir l'URL de base du serveur
let base_url = self.base_url();
// Créer le client depuis la config
let client = QobuzClient::from_config()
.await
.map_err(|e| SourceInitError::QobuzError(format!("Failed to create client: {}", e)))?;
// Créer la source depuis le registry
let source = QobuzSource::from_registry(client)
let source = QobuzSource::from_registry(client, base_url)
.map_err(|e| SourceInitError::QobuzError(format!("Failed to create source: {}", e)))?;
// Enregistrer la source
@@ -154,13 +157,16 @@ impl SourcesExt for Server {
tracing::info!("Initializing Qobuz source with explicit credentials...");
// Obtenir l'URL de base du serveur
let base_url = self.base_url();
// Créer le client avec credentials
let client = QobuzClient::new(username, password)
.await
.map_err(|e| SourceInitError::QobuzError(format!("Failed to authenticate: {}", e)))?;
// Créer la source depuis le registry
let source = QobuzSource::from_registry(client)
let source = QobuzSource::from_registry(client, base_url)
.map_err(|e| SourceInitError::QobuzError(format!("Failed to create source: {}", e)))?;
// Enregistrer la source

View File

@@ -26,6 +26,9 @@ pub struct QobuzCredentials {
pub username: Option<String>,
/// Mot de passe Qobuz (optionnel, lu depuis la config si absent)
pub password: Option<String>,
/// URL de base du serveur (optionnelle, "http://localhost:8080" par défaut)
#[serde(default)]
pub base_url: Option<String>,
}
/// Paramètres pour Radio Paradise
@@ -75,6 +78,11 @@ async fn register_qobuz(Json(creds): Json<QobuzCredentials>) -> impl IntoRespons
use pmoqobuz::{QobuzClient, QobuzSource};
use pmosource::api::register_source;
// Utiliser l'URL de base depuis les params ou une valeur par défaut
let base_url = creds
.base_url
.unwrap_or_else(|| "http://localhost:8080".to_string());
// Créer le client selon les credentials fournis
let client_result = if let (Some(username), Some(password)) = (creds.username, creds.password) {
QobuzClient::new(&username, &password).await
@@ -96,7 +104,7 @@ async fn register_qobuz(Json(creds): Json<QobuzCredentials>) -> impl IntoRespons
};
// Créer et enregistrer la source depuis le registry
let source = match QobuzSource::from_registry(client) {
let source = match QobuzSource::from_registry(client, base_url) {
Ok(s) => Arc::new(s),
Err(e) => {
return (

View File

@@ -135,6 +135,14 @@ impl RadioParadiseSource {
let token = mgr.register_callback(move |event| {
let pid = pid_clone.clone();
if event.playlist_id == pid {
// Ignorer PkUpdated - pas de notification UPnP (évite le reload côté control point)
if matches!(event.kind, pmoplaylist::PlaylistEventKind::PkUpdated { .. }) {
tracing::debug!(
"PK swap in playlist {} - no UPnP notification (prevents playback restart)",
event.playlist_id
);
return;
}
// On ne réagit qu'aux mises à jour structurelles (ajout/suppression)
if !matches!(event.kind, pmoplaylist::PlaylistEventKind::Updated) {
return;

View File

@@ -527,7 +527,8 @@ impl WriteHandle {
manager
.rebuild_track_index(&self.playlist.id, &snapshot)
.await;
manager.notify_playlist_changed(&self.playlist.id);
// PK swap uniquement - pas de notification UPnP pour éviter le reload
manager.notify_playlist_pk_updated(&self.playlist.id, old_pk, new_pk);
}
Ok(())

View File

@@ -48,6 +48,9 @@ pub struct PlaylistEvent {
pub enum PlaylistEventKind {
/// La playlist a été modifiée (ajout/suppression/changement de config).
Updated,
/// Cache PK commuté (lazy→real) - pas de changement structurel.
/// N'émet PAS de ContainersUpdated UPnP pour éviter le reload.
PkUpdated { old_pk: String, new_pk: String },
/// Un morceau référencé par la playlist a été servi par le cache audio.
TrackPlayed { cache_pk: String, qualifier: String },
}
@@ -224,6 +227,28 @@ impl PlaylistManager {
.await
}
/// Récupère l'âge d'une playlist depuis sa création
pub async fn get_playlist_age(&self, id: &str) -> Result<Option<Duration>> {
use std::time::{SystemTime, UNIX_EPOCH};
let persistence = match &self.inner.persistence {
Some(p) => p,
None => return Ok(None),
};
let created_at_nanos = persistence.get_playlist_created_at(id).await?;
if let Some(nanos) = created_at_nanos {
let created_at = UNIX_EPOCH + Duration::from_nanos(nanos as u64);
let age = SystemTime::now()
.duration_since(created_at)
.unwrap_or(Duration::ZERO);
Ok(Some(age))
} else {
Ok(None)
}
}
/// Enregistre un callback d'évènement playlist (update, track joué).
///
/// Retourne un jeton (u64) pour désenregistrer plus tard.
@@ -248,6 +273,18 @@ impl PlaylistManager {
self.notify_playlist_event(id, PlaylistEventKind::Updated);
}
/// Notifie que des PK ont été swappés (lazy→real).
/// N'émet PAS de notification UPnP ContainersUpdated pour éviter le reload.
pub(crate) fn notify_playlist_pk_updated(&self, id: &str, old_pk: &str, new_pk: &str) {
self.notify_playlist_event(
id,
PlaylistEventKind::PkUpdated {
old_pk: old_pk.to_string(),
new_pk: new_pk.to_string(),
},
);
}
/// Notifie les callbacks qu'un morceau a été joué pour une playlist donnée.
pub(crate) fn notify_playlist_track_played(
&self,

View File

@@ -259,6 +259,28 @@ impl PersistenceManager {
Ok(ids)
}
/// Récupère le timestamp de création d'une playlist
pub async fn get_playlist_created_at(&self, id: &str) -> Result<Option<i64>> {
let conn = self.conn.lock().unwrap();
let mut stmt = conn
.prepare("SELECT created_at FROM playlists WHERE id = ?1")
.map_err(|e| {
crate::Error::PersistenceError(format!("Failed to prepare statement: {}", e))
})?;
let result = stmt.query_row(params![id], |row| row.get(0));
match result {
Ok(created_at) => Ok(Some(created_at)),
Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
Err(e) => Err(crate::Error::PersistenceError(format!(
"Failed to get created_at: {}",
e
))),
}
}
/// Supprime tous les tracks contenant un cache_pk donné
pub async fn remove_by_cache_pk(&self, cache_pk: &str) -> Result<()> {
let conn = self.conn.lock().unwrap();

View File

@@ -55,6 +55,10 @@ pub async fn playlist_events_sse(Query(params): Query<EventsQuery>) -> impl Into
let (kind, cache_pk, qualifier) = match &envelope.event.kind {
PlaylistEventKind::Updated => ("updated", None, None),
PlaylistEventKind::PkUpdated { old_pk: _, new_pk: _ } => {
// PK swap silencieux - envoyer quand même l'événement SSE pour debug/monitoring
("pk_updated", None, None)
}
PlaylistEventKind::TrackPlayed { cache_pk, qualifier } => {
("track_played", Some(cache_pk.as_str()), Some(qualifier.as_str()))
}

View File

@@ -149,6 +149,12 @@ struct FeaturedPlaylistsResponse {
playlists: PaginatedResponse<PlaylistResponse>,
}
/// Réponse artistes featured
#[derive(Debug, Deserialize)]
struct FeaturedArtistsResponse {
artists: PaginatedResponse<ArtistResponse>,
}
/// Réponse search
#[derive(Debug, Deserialize)]
struct SearchResponse {
@@ -397,6 +403,36 @@ impl QobuzApi {
.collect())
}
/// Récupère les artistes featured
pub async fn get_featured_artists(
&self,
genre_id: Option<&str>,
limit: Option<u32>,
offset: Option<u32>,
) -> Result<Vec<Artist>> {
debug!("Fetching featured artists");
let limit_str = limit.unwrap_or(100).to_string();
let offset_str = offset.unwrap_or(0).to_string();
let mut params = vec![
("type", "featured-artists"),
("limit", &limit_str),
("offset", &offset_str),
];
if let Some(gid) = genre_id {
params.push(("genre_ids", gid));
}
let response: FeaturedArtistsResponse = self.get("/artist/getFeatured", &params).await?;
Ok(response
.artists
.items
.into_iter()
.map(Self::parse_artist)
.collect())
}
/// Recherche dans le catalogue
pub async fn search(&self, query: &str, type_: Option<&str>) -> Result<SearchResult> {
debug!("Searching for '{}' (type: {:?})", query, type_);
@@ -443,6 +479,13 @@ impl QobuzApi {
// Fonctions de parsing publiques (utilisées aussi par le module user)
pub(crate) fn parse_album(response: AlbumResponse) -> Album {
// Log pour débugger les valeurs audio
if let Some(rate) = response.maximum_sampling_rate {
debug!("Album {} - maximum_sampling_rate: {} Hz", response.id, rate);
} else {
debug!("Album {} - maximum_sampling_rate: None", response.id);
}
Album {
id: response.id,
title: response.title,

View File

@@ -14,6 +14,8 @@ pub struct QobuzCache {
albums: Arc<MokaCache<String, Album>>,
/// Cache des tracks (TTL: 1 heure)
tracks: Arc<MokaCache<String, Track>>,
/// Cache des tracks d'un album complet (TTL: 1 heure)
album_tracks: Arc<MokaCache<String, Vec<Track>>>,
/// Cache des artistes (TTL: 1 heure)
artists: Arc<MokaCache<String, Artist>>,
/// Cache des playlists (TTL: 30 minutes)
@@ -45,6 +47,12 @@ impl QobuzCache {
.time_to_live(Duration::from_secs(3600)) // 1 heure
.build(),
),
album_tracks: Arc::new(
MokaCache::builder()
.max_capacity(max_capacity)
.time_to_live(Duration::from_secs(3600)) // 1 heure
.build(),
),
artists: Arc::new(
MokaCache::builder()
.max_capacity(max_capacity / 2)
@@ -106,6 +114,23 @@ impl QobuzCache {
self.tracks.invalidate(id).await;
}
// ============ Album Tracks (liste complète des tracks d'un album) ============
/// Récupère la liste complète des tracks d'un album depuis le cache
pub async fn get_album_tracks(&self, album_id: &str) -> Option<Vec<Track>> {
self.album_tracks.get(album_id).await
}
/// Ajoute la liste complète des tracks d'un album au cache
pub async fn put_album_tracks(&self, album_id: String, tracks: Vec<Track>) {
self.album_tracks.insert(album_id, tracks).await;
}
/// Invalide la liste des tracks d'un album du cache
pub async fn invalidate_album_tracks(&self, album_id: &str) {
self.album_tracks.invalidate(album_id).await;
}
// ============ Artists ============
/// Récupère un artiste depuis le cache
@@ -180,6 +205,7 @@ impl QobuzCache {
pub async fn clear_all(&self) {
self.albums.invalidate_all();
self.tracks.invalidate_all();
self.album_tracks.invalidate_all();
self.artists.invalidate_all();
self.playlists.invalidate_all();
self.searches.invalidate_all();
@@ -190,6 +216,7 @@ impl QobuzCache {
pub async fn stats(&self) -> CacheStats {
self.albums.run_pending_tasks().await;
self.tracks.run_pending_tasks().await;
self.album_tracks.run_pending_tasks().await;
self.artists.run_pending_tasks().await;
self.playlists.run_pending_tasks().await;
self.searches.run_pending_tasks().await;
@@ -198,6 +225,7 @@ impl QobuzCache {
CacheStats {
albums_count: self.albums.entry_count(),
tracks_count: self.tracks.entry_count(),
album_tracks_count: self.album_tracks.entry_count(),
artists_count: self.artists.entry_count(),
playlists_count: self.playlists.entry_count(),
searches_count: self.searches.entry_count(),
@@ -219,6 +247,8 @@ pub struct CacheStats {
pub albums_count: u64,
/// Nombre de tracks en cache
pub tracks_count: u64,
/// Nombre de listes complètes de tracks d'albums en cache
pub album_tracks_count: u64,
/// Nombre d'artistes en cache
pub artists_count: u64,
/// Nombre de playlists en cache
@@ -234,6 +264,7 @@ impl CacheStats {
pub fn total_count(&self) -> u64 {
self.albums_count
+ self.tracks_count
+ self.album_tracks_count
+ self.artists_count
+ self.playlists_count
+ self.searches_count

View File

@@ -584,14 +584,24 @@ impl QobuzClient {
/// Récupère les tracks d'un album
pub async fn get_album_tracks(&self, album_id: &str) -> Result<Vec<Track>> {
// Vérifier le cache d'abord
if let Some(tracks) = self.cache.get_album_tracks(album_id).await {
debug!("Album tracks for {} found in cache", album_id);
return Ok(tracks);
}
// Sinon, récupérer depuis l'API
let tracks = self
.call_with_auth_repair("get_album_tracks", || self.api.get_album_tracks(album_id))
.await?;
// Mettre les tracks en cache
// Mettre les tracks en cache (individuellement ET la liste complète)
for track in &tracks {
self.cache.put_track(track.id.clone(), track.clone()).await;
}
self.cache
.put_album_tracks(album_id.to_string(), tracks.clone())
.await;
Ok(tracks)
}
@@ -742,6 +752,19 @@ impl QobuzClient {
.await
}
/// Récupère les artistes featured
pub async fn get_featured_artists(
&self,
genre_id: Option<&str>,
limit: Option<u32>,
offset: Option<u32>,
) -> Result<Vec<Artist>> {
self.call_with_auth_repair("get_featured_artists", || {
self.api.get_featured_artists(genre_id, limit, offset)
})
.await
}
// ============ Recherche ============
/// Recherche dans le catalogue Qobuz

View File

@@ -27,10 +27,10 @@ impl ToDIDL for Album {
///
/// ```rust,ignore
/// let album = client.get_album("12345").await?;
/// let container = album.to_didl_container("0$qobuz$albums")?;
/// let container = album.to_didl_container("qobuz:favorites")?;
/// ```
fn to_didl_container(&self, parent_id: &str) -> Result<Container> {
let id = format!("0$qobuz$album${}", self.id);
let id = format!("qobuz:album:{}", self.id);
Ok(Container {
id,
@@ -71,10 +71,10 @@ impl ToDIDL for Track {
///
/// ```rust,ignore
/// let track = client.get_track("98765").await?;
/// let item = track.to_didl_item("0$qobuz$album$12345")?;
/// let item = track.to_didl_item("qobuz:album:12345")?;
/// ```
fn to_didl_item(&self, parent_id: &str) -> Result<Item> {
let id = format!("0$qobuz$track${}", self.id);
let id = format!("qobuz:track:{}", self.id);
// Déterminer l'artiste à afficher
let artist_name = self
@@ -128,7 +128,7 @@ impl ToDIDL for Track {
impl ToDIDL for Playlist {
/// Convertit une playlist en Container DIDL
fn to_didl_container(&self, parent_id: &str) -> Result<Container> {
let id = format!("0$qobuz$playlist${}", self.id);
let id = format!("qobuz:playlist:{}", self.id);
Ok(Container {
id,
@@ -211,7 +211,7 @@ mod tests {
};
let container = album.to_didl_container("parent").unwrap();
assert_eq!(container.id, "0$qobuz$album$123");
assert_eq!(container.id, "qobuz:album:123");
assert_eq!(container.parent_id, "parent");
assert!(container.title.contains("Test Album"));
}
@@ -234,7 +234,7 @@ mod tests {
};
let item = track.to_didl_item("parent").unwrap();
assert_eq!(item.id, "0$qobuz$track$789");
assert_eq!(item.id, "qobuz:track:789");
assert_eq!(item.parent_id, "parent");
assert_eq!(item.title, "Test Track");
}

View File

@@ -124,7 +124,7 @@
//! let cover_cache = Arc::new(CoverCache::new("./cache/covers", 500)?);
//! let audio_cache = Arc::new(AudioCache::new("./cache/audio", 100)?);
//!
//! let source = QobuzSource::new(client, cover_cache, audio_cache);
//! let source = QobuzSource::new(client, cover_cache, audio_cache, "http://localhost:8080");
//! # Ok(())
//! # }
//! ```
@@ -145,7 +145,7 @@
//! let cover_cache = Arc::new(CoverCache::new("./cache/covers", 500)?);
//! let audio_cache = Arc::new(AudioCache::new("./cache/audio", 100)?);
//!
//! let source = QobuzSource::new(client, cover_cache, audio_cache);
//! let source = QobuzSource::new(client, cover_cache, audio_cache, "http://localhost:8080");
//!
//! // Add a track with caching
//! let tracks = source.client().get_favorite_tracks().await?;

View File

@@ -269,7 +269,13 @@ impl Album {
/// Retourne un titre formaté avec les informations audio si disponibles
pub fn formatted_title(&self) -> String {
if let (Some(rate), Some(depth)) = (self.maximum_sampling_rate, self.maximum_bit_depth) {
format!("{} ({:.0}/{} bit)", self.title, rate / 1000.0, depth)
// Convertir Hz en kHz, en gérant les valeurs qui pourraient déjà être en kHz
let rate_khz = if rate > 1000.0 {
rate / 1000.0
} else {
rate
};
format!("{} ({:.1} kHz / {} bits)", self.title, rate_khz, depth)
} else {
self.title.clone()
}
@@ -316,3 +322,109 @@ impl SearchResult {
self.total_count() == 0
}
}
/// Types d'albums featured disponibles dans le catalogue Qobuz
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum FeaturedAlbumType {
/// Nouveautés
NewReleases,
/// Nouveautés complètes
NewReleasesFull,
/// Discographie idéale
IdealDiscography,
/// Qobuzissime
Qobuzissims,
/// Choix de l'éditeur
EditorPicks,
/// Prix de la presse
PressAwards,
}
impl FeaturedAlbumType {
/// Retourne l'identifiant API pour ce type
pub fn api_id(&self) -> &'static str {
match self {
Self::NewReleases => "new-releases",
Self::NewReleasesFull => "new-releases-full",
Self::IdealDiscography => "ideal-discography",
Self::Qobuzissims => "qobuzissims",
Self::EditorPicks => "editor-picks",
Self::PressAwards => "press-awards",
}
}
}
/// Tags de playlists disponibles dans le catalogue Qobuz
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PlaylistTag {
/// Hi-Res
HiRes,
/// Nouvelles
New,
/// Thématiques
Themes,
/// Choix d'artistes
ArtistsChoices,
/// Labels
Labels,
/// Humeurs
Moods,
/// Artistes
Artists,
/// Événements
Events,
/// Auditoriums
Auditoriums,
/// Populaires
Popular,
}
impl PlaylistTag {
/// Retourne l'identifiant API pour ce tag
pub fn api_id(&self) -> &'static str {
match self {
Self::HiRes => "hi-res",
Self::New => "new",
Self::Themes => "focus",
Self::ArtistsChoices => "danslecasque",
Self::Labels => "label",
Self::Moods => "mood",
Self::Artists => "artist",
Self::Events => "events",
Self::Auditoriums => "auditoriums",
Self::Popular => "popular",
}
}
/// Retourne le titre localisé pour ce tag
pub fn display_name(&self) -> &'static str {
match self {
Self::HiRes => "Playlists (Hi-Res)",
Self::New => "Playlists (New)",
Self::Themes => "Playlists (Themes)",
Self::ArtistsChoices => "Playlists (Artist's Choices)",
Self::Labels => "Playlists (Labels)",
Self::Moods => "Playlists (Moods)",
Self::Artists => "Playlists (Artists)",
Self::Events => "Playlists (Events)",
Self::Auditoriums => "Playlists (Auditoriums)",
Self::Popular => "Playlists (Popular)",
}
}
/// Retourne la liste de tous les tags
pub fn all() -> &'static [PlaylistTag] {
&[
Self::HiRes,
Self::New,
Self::Themes,
Self::ArtistsChoices,
Self::Labels,
Self::Moods,
Self::Artists,
Self::Events,
Self::Auditoriums,
Self::Popular,
]
}
}

File diff suppressed because it is too large Load Diff

View File

@@ -217,6 +217,29 @@ impl SourceCacheManager {
}
}
/// Obtenir l'URL absolue d'un fichier audio en cache à partir de son pk
///
/// # Arguments
///
/// * `pk` - Clé primaire du fichier audio dans le cache
///
/// # Returns
///
/// L'URL absolue complète du fichier audio (ex: http://localhost:8080/audio/flac/QOBUZ:123)
pub fn audio_url(&self, _pk: &str) -> Result<String> {
#[cfg(feature = "server")]
{
pmoupnp::cache_registry::build_audio_url(_pk, None)
.map_err(|e| MusicSourceError::CacheError(e.to_string()))
}
#[cfg(not(feature = "server"))]
{
Err(MusicSourceError::CacheError(
"Server feature not enabled - cannot build audio URL".to_string(),
))
}
}
/// Cacher une piste audio depuis une URL
///
/// Utilise la collection de cette source pour organiser les pistes.

View File

@@ -36,6 +36,7 @@ reqwest = "0.12.23"
utoipa = { version = "5.3", features = ["axum_extras"] }
socket2 = "0.5"
get_if_addrs = "0.5"
serde_yaml = "0.9"
[features]
default = ["server"]

128
pmoupnp/src/config_ext.rs Normal file
View File

@@ -0,0 +1,128 @@
//! Extension pour intégrer la configuration UPnP dans pmoconfig
//!
//! Ce module fournit le trait `UpnpConfigExt` qui permet d'ajouter facilement
//! des méthodes de configuration UPnP à pmoconfig::Config.
//!
//! Il suit le même pattern que pmocache/src/config_ext.rs pour la cohérence.
use anyhow::Result;
use pmoconfig::Config;
use serde_yaml::Value;
// Constantes par défaut pour les noms UPnP
const DEFAULT_MANUFACTURER: &str = "PMOMusic";
const DEFAULT_UDN_PREFIX: &str = "pmomusic";
const DEFAULT_MODEL_NAME_PREFIX: &str = "PMOMusic";
const DEFAULT_FRIENDLY_NAME_PREFIX: &str = "PMOMusic";
/// Trait d'extension pour ajouter la configuration UPnP à pmoconfig
///
/// Ce trait étend `pmoconfig::Config` avec des méthodes pour configurer
/// les noms et identifiants des devices UPnP.
///
/// # Exemple
///
/// ```rust,ignore
/// use pmoconfig::get_config;
/// use pmoupnp::UpnpConfigExt;
///
/// let config = get_config();
/// let manufacturer = config.get_upnp_manufacturer()?;
/// let udn_prefix = config.get_upnp_udn_prefix()?;
/// ```
pub trait UpnpConfigExt {
/// Récupère le fabricant pour les devices UPnP
///
/// # Returns
///
/// Le nom du fabricant à afficher dans les descripteurs UPnP (défaut: "PMOMusic")
fn get_upnp_manufacturer(&self) -> Result<String>;
/// Définit le fabricant pour les devices UPnP
fn set_upnp_manufacturer(&self, manufacturer: String) -> Result<()>;
/// Récupère le préfixe UDN pour les devices UPnP
///
/// # Returns
///
/// Le préfixe utilisé pour générer les UDN (défaut: "pmomusic")
fn get_upnp_udn_prefix(&self) -> Result<String>;
/// Définit le préfixe UDN pour les devices UPnP
fn set_upnp_udn_prefix(&self, prefix: String) -> Result<()>;
/// Récupère le préfixe pour les noms de modèle des devices UPnP
///
/// # Returns
///
/// Le préfixe utilisé pour construire les model names (défaut: "PMOMusic")
fn get_upnp_model_name_prefix(&self) -> Result<String>;
/// Définit le préfixe pour les noms de modèle des devices UPnP
fn set_upnp_model_name_prefix(&self, prefix: String) -> Result<()>;
/// Récupère le préfixe pour les noms conviviaux des devices UPnP
///
/// # Returns
///
/// Le préfixe utilisé pour construire les friendly names (défaut: "PMOMusic")
fn get_upnp_friendly_name_prefix(&self) -> Result<String>;
/// Définit le préfixe pour les noms conviviaux des devices UPnP
fn set_upnp_friendly_name_prefix(&self, prefix: String) -> Result<()>;
}
impl UpnpConfigExt for Config {
fn get_upnp_manufacturer(&self) -> Result<String> {
match self.get_value(&["host", "upnp", "manufacturer"]) {
Ok(Value::String(s)) if !s.is_empty() => Ok(s),
_ => Ok(DEFAULT_MANUFACTURER.to_string()),
}
}
fn set_upnp_manufacturer(&self, manufacturer: String) -> Result<()> {
self.set_value(
&["host", "upnp", "manufacturer"],
Value::String(manufacturer),
)
}
fn get_upnp_udn_prefix(&self) -> Result<String> {
match self.get_value(&["host", "upnp", "udn_prefix"]) {
Ok(Value::String(s)) if !s.is_empty() => Ok(s),
_ => Ok(DEFAULT_UDN_PREFIX.to_string()),
}
}
fn set_upnp_udn_prefix(&self, prefix: String) -> Result<()> {
self.set_value(&["host", "upnp", "udn_prefix"], Value::String(prefix))
}
fn get_upnp_model_name_prefix(&self) -> Result<String> {
match self.get_value(&["host", "upnp", "model_name_prefix"]) {
Ok(Value::String(s)) if !s.is_empty() => Ok(s),
_ => Ok(DEFAULT_MODEL_NAME_PREFIX.to_string()),
}
}
fn set_upnp_model_name_prefix(&self, prefix: String) -> Result<()> {
self.set_value(
&["host", "upnp", "model_name_prefix"],
Value::String(prefix),
)
}
fn get_upnp_friendly_name_prefix(&self) -> Result<String> {
match self.get_value(&["host", "upnp", "friendly_name_prefix"]) {
Ok(Value::String(s)) if !s.is_empty() => Ok(s),
_ => Ok(DEFAULT_FRIENDLY_NAME_PREFIX.to_string()),
}
}
fn set_upnp_friendly_name_prefix(&self, prefix: String) -> Result<()> {
self.set_value(
&["host", "upnp", "friendly_name_prefix"],
Value::String(prefix),
)
}
}

View File

@@ -122,6 +122,84 @@ impl Device {
}
}
/// Crée un nouveau modèle de device en utilisant la configuration.
///
/// Cette factory method charge les préfixes depuis pmoconfig et construit
/// automatiquement les noms finaux en combinant les préfixes avec les suffixes.
///
/// # Arguments
///
/// * `name` - Nom unique du device
/// * `device_type` - Type UPnP du device (ex: "MediaServer", "MediaRenderer")
/// * `friendly_name_suffix` - Suffixe pour le nom convivial (sera combiné avec le préfixe)
///
/// # Examples
///
/// ```ignore
/// use pmoupnp::devices::Device;
///
/// let device = Device::new_from_config(
/// "PMO_MediaServer".to_string(),
/// "MediaServer".to_string(),
/// "Media Server".to_string(),
/// );
/// // Avec config par défaut :
/// // - manufacturer = "PMOMusic"
/// // - udn_prefix = "pmomusic"
/// // - model_name = "PMOMusic Media Server"
/// // - friendly_name = "PMOMusic Media Server"
/// ```
pub fn new_from_config(
name: String,
device_type: String,
friendly_name_suffix: String,
) -> Self {
use crate::config_ext::UpnpConfigExt;
let config = pmoconfig::get_config();
// Charger les valeurs depuis la config (avec fallback aux defaults)
let manufacturer = config
.get_upnp_manufacturer()
.unwrap_or_else(|_| "PMOMusic".to_string());
let udn_prefix = config
.get_upnp_udn_prefix()
.unwrap_or_else(|_| "pmomusic".to_string());
let model_name_prefix = config
.get_upnp_model_name_prefix()
.unwrap_or_else(|_| "PMOMusic".to_string());
let friendly_name_prefix = config
.get_upnp_friendly_name_prefix()
.unwrap_or_else(|_| "PMOMusic".to_string());
// Construire les noms finaux
let model_name = format!("{} {}", model_name_prefix, device_type);
let friendly_name = format!("{} {}", friendly_name_prefix, friendly_name_suffix);
Self {
object: UpnpObjectType {
name: name.clone(),
object_type: "Device".to_string(),
},
device_type,
version: 1,
friendly_name,
manufacturer,
manufacturer_url: None,
model_description: None,
model_name,
model_number: None,
model_url: None,
serial_number: None,
udn_prefix,
upc: None,
icon_url: None,
presentation_url: None,
services: RwLock::new(HashMap::new()),
devices: RwLock::new(HashMap::new()),
}
}
/// Retourne le type de device UPnP.
///
/// Format: `urn:schemas-upnp-org:device:{type}:{version}`

View File

@@ -3,6 +3,7 @@ mod object_trait;
pub mod actions;
pub mod cache_registry;
pub mod config_ext;
pub mod devices;
pub mod services;
pub mod soap;
@@ -20,6 +21,7 @@ use std::{collections::HashMap, sync::Arc};
pub use pmoaudiocache::get_audio_cache;
pub use pmocovers::get_cover_cache;
pub use crate::config_ext::UpnpConfigExt;
pub use crate::object_trait::*;
pub use crate::upnp_server::UpnpServerExt;

View File

@@ -261,7 +261,12 @@ impl UpnpServerExt for Server {
if self.ssdp_enabled() {
let ssdp_opt = SSDP_SERVER.read().unwrap();
if let Some(ref ssdp) = *ssdp_opt {
let ssdp_device = di.to_ssdp_device("PMOMusic", "1.0");
use crate::config_ext::UpnpConfigExt;
let config = pmoconfig::get_config();
let manufacturer = config
.get_upnp_manufacturer()
.unwrap_or_else(|_| "PMOMusic".to_string());
let ssdp_device = di.to_ssdp_device(&manufacturer, "1.0");
ssdp.add_device(ssdp_device);
info!("✅ SSDP announcement for {}", di.udn());
}

46
rust-analyzer.toml Normal file
View File

@@ -0,0 +1,46 @@
# Configuration rust-analyzer pour gros workspace (23 crates)
# Optimisé pour éviter les crashs et limiter l'utilisation mémoire
# Limiter l'analyse en arrière-plan
[checkOnSave]
enable = true
# N'analyser que la cible par défaut (pas tous les targets)
allTargets = false
# Utiliser clippy au lieu de cargo check (optionnel, commentez si trop lent)
# command = "clippy"
# Désactiver les build scripts pour réduire la charge
[cargo]
buildScripts.enable = false
# Ne charger que les crates nécessaires
loadOutDirsFromCheck = false
# Désactiver les proc-macros si elles causent des problèmes
[procMacro]
enable = true
# Si les crashs persistent, passez à false ci-dessus
# Limiter la complétion
[completion]
limit = 50
# Désactiver certains diagnostics coûteux
[diagnostics]
disabled = [
"unresolved-proc-macro",
"macro-error",
]
# Optimisations de performance
[inlayHints]
# Réduire les hints pour améliorer les perfs
maxLength = 25
# Limiter la profondeur d'analyse des types
[typing]
autoClosingAngleBrackets.enable = false
# Pour les très gros workspaces, décommenter pour analyser moins de crates
# [linkedProjects]
# Spécifier uniquement les crates que vous éditez activement
# Par exemple : ["PMOMusic/Cargo.toml", "pmocontrol/Cargo.toml"]