push-xyrsxtokwrlu #81

Merged
eric merged 8 commits from push-xyrsxtokwrlu into main 2026-03-25 13:09:23 +01:00
25 changed files with 913 additions and 217 deletions

BIN
.DS_Store vendored

Binary file not shown.

1
.gitignore vendored
View File

@@ -31,6 +31,7 @@ xxx
C/src/soxr-0.1.3/Release/tests
**/Release/
**/Debug/
qobuz_debug
OLD-GO-CODE/
xxx
xx

2
Cargo.lock generated
View File

@@ -4,7 +4,7 @@ version = 4
[[package]]
name = "PMOMusic"
version = "0.3.25"
version = "0.3.27"
dependencies = [
"axum 0.8.7",
"console-subscriber",

View File

@@ -32,8 +32,9 @@ RUN apt-get update && apt-get install -y \
cmake \
&& rm -rf /var/lib/apt/lists/*
# Copy Cargo workspace files
# Copy Cargo workspace files and registry configuration
COPY Cargo.toml Cargo.lock ./
COPY .cargo/ ./.cargo/
# Copy all crates
COPY PMOMusic/ ./PMOMusic/

View File

@@ -1,6 +1,6 @@
[package]
name = "PMOMusic"
version = "0.3.26"
version = "0.3.27"
edition = "2024"
[dependencies]

View File

@@ -449,8 +449,7 @@ async function handleTransferQueue(event: Event, targetRendererId: string) {
@media (max-width: 768px) and (orientation: portrait) {
.drawer-backdrop {
right: 0; /* Mobile portrait: backdrop prend tout l'écran */
background: rgba(0, 0, 0, 0.4); /* Plus sombre sur mobile */
display: none; /* Mobile portrait: le drawer couvre 100vw, pas besoin de backdrop */
}
}

View File

@@ -1,5 +1,5 @@
<script setup lang="ts">
import { ref, computed, watch, reactive } from "vue";
import { ref, computed, watch, reactive, onMounted, onBeforeUnmount } from "vue";
import { useRouter } from "vue-router";
import {
X,
@@ -20,7 +20,6 @@ import type {
MediaServerSummary,
ContainerEntry,
} from "@/services/pmocontrol/types";
import type { BrowseState } from "@/composables/useMediaServers";
const props = defineProps<{
modelValue: boolean; // v-model pour contrôler l'ouverture
@@ -35,6 +34,10 @@ const {
allServers,
fetchServers,
browseContainer,
getBrowseCached,
loadMoreBrowse,
hasMore,
loadingMore,
currentPath,
setPath,
clearPath,
@@ -47,9 +50,51 @@ const router = useRouter();
// État de navigation
const currentServer = ref<MediaServerSummary | null>(null);
const browseData = ref<BrowseState | null>(null);
const currentContainerId = ref<string | null>(null);
const browseData = computed(() =>
currentServer.value && currentContainerId.value
? (getBrowseCached(currentServer.value.id, currentContainerId.value) ?? null)
: null,
);
const isLoading = ref(false);
// Infinite scroll
const sentinelRef = ref<HTMLElement | null>(null);
const drawerContentRef = ref<HTMLElement | null>(null);
let observer: IntersectionObserver | null = null;
const canLoadMore = computed(() =>
currentServer.value && currentContainerId.value
? hasMore(currentServer.value.id, currentContainerId.value)
: false,
);
function setupObserver() {
if (observer) observer.disconnect();
observer = new IntersectionObserver(
(entries) => {
if (entries[0]?.isIntersecting && canLoadMore.value && !loadingMore.value) {
if (currentServer.value && currentContainerId.value) {
loadMoreBrowse(currentServer.value.id, currentContainerId.value);
}
}
},
{ root: drawerContentRef.value, threshold: 0.1 },
);
if (sentinelRef.value) observer.observe(sentinelRef.value);
}
onMounted(() => setupObserver());
onBeforeUnmount(() => observer?.disconnect());
watch(sentinelRef, (el) => {
if (el) setupObserver();
});
watch(drawerContentRef, (el) => {
if (el) setupObserver();
});
// État du menu dropdown (pour chaque item, on stocke si son menu est ouvert)
const openMenuId = ref<string | null>(null);
@@ -85,7 +130,7 @@ watch(
} else {
// Reset navigation quand on ferme
currentServer.value = null;
browseData.value = null;
currentContainerId.value = null;
clearPath();
closeMenu();
imageStates.clear();
@@ -143,11 +188,12 @@ async function handleServerClick(server: MediaServerSummary) {
try {
// Browse racine (containerId = "0")
browseData.value = await browseContainer(server.id, "0");
await browseContainer(server.id, "0");
currentContainerId.value = "0";
setPath([{ id: "0", title: server.friendly_name }]);
} catch (error) {
console.error("[ServerDrawer] Erreur browse racine:", error);
browseData.value = null;
currentContainerId.value = null;
} finally {
isLoading.value = false;
}
@@ -155,7 +201,7 @@ async function handleServerClick(server: MediaServerSummary) {
function goBack() {
currentServer.value = null;
browseData.value = null;
currentContainerId.value = null;
clearPath();
}
@@ -164,10 +210,8 @@ async function handleContainerClick(item: ContainerEntry) {
isLoading.value = true;
try {
browseData.value = await browseContainer(
currentServer.value.id,
item.id,
);
await browseContainer(currentServer.value.id, item.id);
currentContainerId.value = item.id;
setPath([...currentPath.value, { id: item.id, title: item.title }]);
} catch (error) {
console.error("[ServerDrawer] Erreur browse container:", error);
@@ -185,10 +229,8 @@ async function handleBreadcrumbClick(index: number) {
isLoading.value = true;
try {
browseData.value = await browseContainer(
currentServer.value.id,
targetCrumb.id,
);
await browseContainer(currentServer.value.id, targetCrumb.id);
currentContainerId.value = targetCrumb.id;
setPath(currentPath.value.slice(0, index + 1));
} catch (error) {
console.error("[ServerDrawer] Erreur breadcrumb navigation:", error);
@@ -386,7 +428,7 @@ function handleSettingsClick() {
</nav>
<!-- Contenu -->
<div class="drawer-content">
<div ref="drawerContentRef" class="drawer-content">
<!-- Liste des serveurs -->
<div v-if="!isNavigating">
<!-- Servers online -->
@@ -627,8 +669,16 @@ function handleSettingsClick() {
</li>
</ul>
<!-- Sentinel infinite scroll -->
<div v-if="canLoadMore" ref="sentinelRef" class="scroll-sentinel" />
<!-- Spinner load more -->
<div v-if="loadingMore" class="load-more-spinner">
<div class="spinner"></div>
</div>
<!-- Vide -->
<div v-else class="empty-servers">
<div v-else-if="!browseData" class="empty-servers">
<Folder :size="48" />
<p>Dossier vide</p>
</div>
@@ -666,8 +716,7 @@ function handleSettingsClick() {
@media (max-width: 768px) and (orientation: portrait) {
.drawer-backdrop {
left: 0; /* Mobile portrait: backdrop commence à gauche car drawer prend 100vw */
background: rgba(0, 0, 0, 0.4); /* Plus sombre sur mobile */
display: none; /* Mobile portrait: le drawer couvre 100vw, pas besoin de backdrop */
}
}
@@ -1257,6 +1306,17 @@ function handleSettingsClick() {
transform: translateY(0);
}
.scroll-sentinel {
height: 1px;
}
.load-more-spinner {
display: flex;
justify-content: center;
padding: var(--spacing-md);
color: var(--color-text-secondary);
}
/* Animations */
.backdrop-enter-active {
transition: opacity 0.3s ease-out;

View File

@@ -20,6 +20,9 @@ fn map_db_err(err: rusqlite::Error) -> MetadataError {
pub struct AudioCacheTrackMetadata {
cache: Arc<crate::Cache>,
pk: String,
/// Real pk si le lazy pk a été téléchargé, None sinon.
/// Résolu une seule fois à la construction pour éviter N appels DB par champ.
real_pk: Option<String>,
}
impl AudioCacheTrackMetadata {
@@ -28,13 +31,26 @@ impl AudioCacheTrackMetadata {
/// Le type implémente ensuite toutes les méthodes du trait `pmometadata::TrackMetadata`
/// en stockant les données dans la base SQLite de `pmocache`.
pub fn new(cache: Arc<crate::Cache>, pk: impl Into<String>) -> Self {
Self {
cache,
pk: pk.into(),
}
let pk_str: String = pk.into();
// Résolution lazy→real une seule fois à la construction (une seule query DB).
let real_pk = if pmocache::is_lazy_pk(&pk_str) {
cache.db.get_pk_by_lazy_pk(&pk_str).ok().flatten()
} else {
None
};
Self { cache, pk: pk_str, real_pk }
}
fn read_raw(&self, key: &str) -> Result<Option<Value>, MetadataError> {
// Si le lazy pk a été téléchargé, lire d'abord sous le real_pk
// (métadonnées écrites par FlacCacheSink), puis fallback sous le lazy_pk
// pour les clés semées avant téléchargement (cover_pk, etc.).
if let Some(real_pk) = &self.real_pk {
if let Ok(Some(v)) = self.cache.db.get_a_metadata(real_pk, key) {
return Ok(Some(v));
}
return self.cache.db.get_a_metadata(&self.pk, key).map_err(map_db_err);
}
self.cache
.db
.get_a_metadata(&self.pk, key)

View File

@@ -404,7 +404,20 @@ impl<C: CacheConfig + 'static> Cache<C> {
) -> Result<Self> {
let directory = PathBuf::from(dir);
std::fs::create_dir_all(&directory)?;
let db = DB::init(&directory.join("cache.db"))?;
let (db, was_reset) = DB::init(&directory.join("cache.db"))?;
// Si le schéma a changé, effacer tous les fichiers du cache (cohérence DB/fichiers)
if was_reset {
tracing::warn!("Cache schema changed, clearing all cache files in {:?}", directory);
if let Ok(entries) = std::fs::read_dir(&directory) {
for entry in entries.flatten() {
let path = entry.path();
if path.is_file() && path.file_name().map_or(false, |n| n != "cache.db") {
std::fs::remove_file(&path).ok();
}
}
}
}
// Créer un channel pour les events (capacité de 100 events en buffer)
let (served_tx, _) = broadcast::channel(100);

View File

@@ -12,6 +12,13 @@ use tracing::{trace, warn};
use std::path::Path;
use std::str::FromStr;
use std::sync::{Mutex, MutexGuard};
/// Version du schéma de la base de données du cache.
///
/// Incrémenter cette constante à chaque modification incompatible du schéma.
/// Cela provoquera la suppression de la DB **et de tous les fichiers du cache**
/// au prochain démarrage.
pub const SCHEMA_VERSION: u32 = 1;
use std::time::Instant;
#[cfg(feature = "openapi")]
@@ -118,7 +125,33 @@ impl DB {
///
/// let db = DB::init(Path::new("cache.db")).unwrap();
/// ```
pub fn init(path: &Path) -> Result<Self, rusqlite::Error> {
/// Initialise la DB. Retourne `(db, was_reset)` où `was_reset` indique si la DB
/// a été supprimée et recréée suite à un changement de version de schéma.
/// Dans ce cas, l'appelant doit aussi effacer les fichiers du cache.
pub fn init(path: &Path) -> Result<(Self, bool), rusqlite::Error> {
let was_reset = if path.exists() {
if let Ok(conn) = Connection::open(path) {
let version: u32 = conn
.query_row("PRAGMA user_version", [], |r| r.get(0))
.unwrap_or(0);
if version != SCHEMA_VERSION {
drop(conn);
warn!(
"Cache DB schema version mismatch (found {}, expected {}), recreating",
version, SCHEMA_VERSION
);
std::fs::remove_file(path).ok();
true
} else {
false
}
} else {
false
}
} else {
false
};
let conn = Connection::open(path)?;
conn.execute("PRAGMA foreign_keys = ON", [])?;
@@ -184,9 +217,10 @@ impl DB {
[],
)?;
Ok(Self {
conn: Mutex::new(conn),
})
// Inscrire la version du schéma
conn.execute_batch(&format!("PRAGMA user_version = {}", SCHEMA_VERSION))?;
Ok((Self { conn: Mutex::new(conn) }, was_reset))
}
/// Ajoute ou met à jour une entrée dans la base de données
@@ -1103,7 +1137,23 @@ impl DB {
return Err(Error::QueryReturnedNoRows);
}
tx.commit()
// Compter les métadonnées encore sous l'ancien lazy_pk avant le commit
let meta_under_lazy: i64 = tx
.query_row(
"SELECT COUNT(*) FROM metadata WHERE pk = ?1",
[lazy_pk],
|r| r.get(0),
)
.unwrap_or(0);
tx.commit()?;
tracing::debug!(
"update_lazy_to_downloaded: {} → {} ({} metadata rows still under lazy_pk)",
lazy_pk, real_pk, meta_under_lazy
);
Ok(())
}
/// Recherche une entry par son origin_url

View File

@@ -1538,41 +1538,62 @@ fn refresh_attached_queue_for(
return Ok(());
}
// Step 3: Browse container
// Step 3: Notify UI that the renderer is loading (Transitioning state)
event_bus.broadcast(RendererEvent::StateChanged {
id: renderer_id.clone(),
state: PlaybackState::Transitioning,
});
// Step 4: Browse container (renamed from Step 3 for clarity)
const MAX_BROWSE_ATTEMPTS: usize = 3;
const BROWSE_RETRY_DELAY_MS: u64 = 200;
let mut attempt = 1;
let entries = loop {
match music_server.browse_children(&container_id, 0, 64) {
Ok(e) => break e,
Err(err) => {
if attempt >= MAX_BROWSE_ATTEMPTS {
warn!(
renderer = renderer_id.0.as_str(),
server = server_id.0.as_str(),
container = container_id.as_str(),
attempts = attempt,
error = %err,
"Failed to browse playlist container for refresh"
);
return Err(err);
}
const BROWSE_PAGE_SIZE: u32 = 64;
debug!(
renderer = renderer_id.0.as_str(),
server = server_id.0.as_str(),
container = container_id.as_str(),
attempt,
error = %err,
"Browse attempt failed, retrying"
);
thread::sleep(Duration::from_millis(
BROWSE_RETRY_DELAY_MS * attempt as u64,
));
attempt += 1;
// Paginated browse — une playlist peut dépasser BROWSE_PAGE_SIZE items
let entries = {
let mut all_entries = Vec::new();
let mut offset = 0u32;
loop {
let mut attempt = 1;
let page = loop {
match music_server.browse_children(&container_id, offset, BROWSE_PAGE_SIZE) {
Ok(e) => break e,
Err(err) => {
if attempt >= MAX_BROWSE_ATTEMPTS {
warn!(
renderer = renderer_id.0.as_str(),
server = server_id.0.as_str(),
container = container_id.as_str(),
attempts = attempt,
error = %err,
"Failed to browse playlist container for refresh"
);
return Err(err);
}
debug!(
renderer = renderer_id.0.as_str(),
server = server_id.0.as_str(),
container = container_id.as_str(),
attempt,
error = %err,
"Browse attempt failed, retrying"
);
thread::sleep(Duration::from_millis(
BROWSE_RETRY_DELAY_MS * attempt as u64,
));
attempt += 1;
}
}
};
let fetched = page.len() as u32;
all_entries.extend(page);
if fetched < BROWSE_PAGE_SIZE {
break;
}
offset += fetched;
}
all_entries
};
debug!(

View File

@@ -2234,8 +2234,19 @@ fn fetch_playback_items(
let object_metadata = server.browse_object(object_id)?;
let entries = if object_metadata.is_container {
// For containers, browse children to get all items
server.browse_children(object_id, 0, BROWSE_DEFAULT_LIMIT)?
// For containers, browse all children with pagination
let mut all_entries = Vec::new();
let mut offset = 0u32;
loop {
let page = server.browse_children(object_id, offset, BROWSE_DEFAULT_LIMIT)?;
let fetched = page.len() as u32;
all_entries.extend(page);
if fetched < BROWSE_DEFAULT_LIMIT {
break;
}
offset += fetched;
}
all_entries
} else {
// For items, use the object itself
vec![object_metadata]

View File

@@ -185,11 +185,12 @@ impl OhPlaylistClient {
return Ok(Vec::new());
}
// OpenHome spec: IdList is space-delimited (not comma-delimited)
let id_list_csv = id_list
.iter()
.map(|id| id.to_string())
.collect::<Vec<_>>()
.join(",");
.join(" ");
let args = [("IdList", id_list_csv.as_str())];
let call_result =

View File

@@ -70,6 +70,8 @@ impl ReadHandle {
let role = self.playlist.role().await;
let cover_pk = self.playlist.cover_pk().await;
let artist = self.playlist.artist().await;
let source = self.playlist.source().await;
let source_version = self.playlist.source_version().await;
let core = self.playlist.core.read().await;
let _ = persistence
.save_playlist(
@@ -78,6 +80,8 @@ impl ReadHandle {
&role,
cover_pk.as_deref(),
artist.as_deref(),
source.as_deref(),
source_version.as_deref(),
&core.config,
&core.tracks,
)
@@ -170,6 +174,21 @@ impl ReadHandle {
&self.playlist.id
}
/// Nombre total de records (sans validation cache)
pub async fn len(&self) -> usize {
self.playlist.core.read().await.len()
}
/// Retourne la source externe.
pub async fn source(&self) -> Option<String> {
self.playlist.source().await
}
/// Retourne la version de la source.
pub async fn source_version(&self) -> Option<String> {
self.playlist.source_version().await
}
/// Génère un Container DIDL-Lite
pub async fn to_container(&self) -> Result<Container> {
if !self.playlist.is_alive() {
@@ -199,6 +218,86 @@ impl ReadHandle {
})
}
/// Génère des Items DIDL-Lite avec pagination offset-based (sans déplacer le curseur)
///
/// Retourne `(items, total)` où `total` est le nombre total de records dans la playlist.
/// Contrairement à `to_items()`, cette méthode ignore le curseur et utilise `offset`.
pub async fn to_items_paged(&self, offset: usize, limit: usize) -> Result<(Vec<Item>, usize)> {
if !self.playlist.is_alive() {
return Err(crate::Error::PlaylistDeleted(self.playlist.id.clone()));
}
let core = self.playlist.core.read().await;
let total = core.len();
let cache = crate::manager::audio_cache()?;
let mut items = Vec::new();
let mut idx = 0usize;
for i in offset..core.len() {
if items.len() >= limit {
break;
}
let record = match core.get(i) {
Some(r) => r,
None => continue,
};
if !cache.is_valid_pk(&record.cache_pk).await {
continue;
}
use pmoaudiocache::metadata_ext::{AudioTrackMetadataExt, TrackMetadataDidlExt};
let track_meta = cache.track_metadata(&record.cache_pk);
let meta = track_meta.read().await;
let url = cache.route_for(&record.cache_pk, None);
let resource = meta.to_didl_resource(url).await;
let title = meta
.get_title()
.await
.ok()
.flatten()
.unwrap_or_else(|| "Unknown".to_string());
let artist = meta.get_artist().await.ok().flatten();
let album = meta.get_album().await.ok().flatten();
let genre = meta.get_genre().await.ok().flatten();
let year = meta.get_year().await.ok().flatten();
let track_number = meta.get_track_number().await.ok().flatten();
let cover_pk = meta.get_cover_pk().await.ok().flatten();
let cover_url = if let Some(pk) = cover_pk.as_ref() {
Some(pmocache::covers_route_for(pk, None))
} else {
meta.get_cover_url().await.ok().flatten()
};
let item = Item {
id: format!("{}:{}", self.playlist.id, offset + idx),
parent_id: self.playlist.id.clone(),
restricted: Some("1".to_string()),
title: title.clone(),
creator: artist.clone(),
class: "object.item.audioItem.musicTrack".to_string(),
artist,
album,
genre,
album_art: cover_url,
album_art_pk: cover_pk,
date: year.map(|y| y.to_string()),
original_track_number: track_number.map(|n| n.to_string()),
resources: vec![resource],
descriptions: vec![],
};
items.push(item);
idx += 1;
}
Ok((items, total))
}
/// Génère des Items DIDL-Lite depuis la position actuelle
pub async fn to_items(&self, limit: usize) -> Result<Vec<Item>> {
if !self.playlist.is_alive() {

View File

@@ -553,6 +553,50 @@ impl WriteHandle {
// Helpers internes
/// Met à jour la source externe de la playlist.
pub async fn set_source(&self, source: Option<String>) -> Result<()> {
if !self.playlist.is_alive() {
return Err(crate::Error::PlaylistDeleted(self.playlist.id.clone()));
}
self.playlist.set_source(source).await;
if self.playlist.persistent {
self.save_to_db().await?;
}
crate::manager::PlaylistManager().notify_playlist_changed(&self.playlist.id);
Ok(())
}
/// Met à jour la version de la source (ex: updated_at).
pub async fn set_source_version(&self, version: Option<String>) -> Result<()> {
if !self.playlist.is_alive() {
return Err(crate::Error::PlaylistDeleted(self.playlist.id.clone()));
}
self.playlist.set_source_version(version).await;
if self.playlist.persistent {
self.save_to_db().await?;
}
crate::manager::PlaylistManager().notify_playlist_changed(&self.playlist.id);
Ok(())
}
/// Retourne la source externe.
pub async fn source(&self) -> Option<String> {
self.playlist.source().await
}
/// Retourne la version de la source.
pub async fn source_version(&self) -> Option<String> {
self.playlist.source_version().await
}
async fn save_to_db(&self) -> Result<()> {
let manager = crate::manager::PlaylistManager();
let persistence = manager
@@ -567,6 +611,8 @@ impl WriteHandle {
let cover_pk = self.playlist.cover_pk().await;
let artist = self.playlist.artist().await;
let source = self.playlist.source().await;
let source_version = self.playlist.source_version().await;
persistence
.save_playlist(
@@ -575,6 +621,8 @@ impl WriteHandle {
&role,
cover_pk.as_deref(),
artist.as_deref(),
source.as_deref(),
source_version.as_deref(),
config,
tracks,
)

View File

@@ -215,6 +215,8 @@ impl PlaylistManager {
&role,
cover_pk.as_deref(),
artist.as_deref(),
None,
None,
&core.config,
&core.tracks,
)
@@ -522,7 +524,7 @@ impl PlaylistManager {
// Pas en mémoire, essayer de charger depuis la DB
if let Some(persistence) = &self.inner.persistence {
if let Some((title, role, config, cover_pk, artist, tracks)) =
if let Some((title, role, config, cover_pk, artist, source, source_version, tracks)) =
persistence.load_playlist(&id).await?
{
// Reconstruire la playlist
@@ -537,10 +539,16 @@ impl PlaylistManager {
cover_pk,
));
// Restaurer l'artiste si présent
// Restaurer les métadonnées optionnelles
if let Some(artist_name) = artist {
playlist.set_artist(Some(artist_name)).await;
}
if source.is_some() {
playlist.set_source(source).await;
}
if source_version.is_some() {
playlist.set_source_version(source_version).await;
}
// Restaurer les tracks
{
@@ -583,7 +591,7 @@ impl PlaylistManager {
// Pas en m<>moire, essayer de ressusciter depuis la DB
if let Some(persistence) = &self.inner.persistence {
if let Some((title, role, config, cover_pk, artist, tracks)) =
if let Some((title, role, config, cover_pk, artist, source, source_version, tracks)) =
persistence.load_playlist(id).await?
{
// Reconstruire la playlist
@@ -598,10 +606,16 @@ impl PlaylistManager {
cover_pk,
));
// Restaurer l'artiste si présent
// Restaurer les métadonnées optionnelles
if let Some(artist_name) = artist {
playlist.set_artist(Some(artist_name)).await;
}
if source.is_some() {
playlist.set_source(source).await;
}
if source_version.is_some() {
playlist.set_source_version(source_version).await;
}
// Restaurer les tracks
{
@@ -995,6 +1009,8 @@ impl PlaylistManager {
let role = playlist.role().await;
let cover_pk = playlist.cover_pk().await;
let artist = playlist.artist().await;
let source = playlist.source().await;
let source_version = playlist.source_version().await;
let core = playlist.core.read().await;
let _ = persistence
.save_playlist(
@@ -1003,6 +1019,8 @@ impl PlaylistManager {
&role,
cover_pk.as_deref(),
artist.as_deref(),
source.as_deref(),
source_version.as_deref(),
&core.config,
&core.tracks,
)

View File

@@ -11,6 +11,13 @@ use std::str::FromStr;
use std::sync::{Arc, Mutex};
use std::time::{Duration, SystemTime, UNIX_EPOCH};
/// Version du schéma de la base de données des playlists.
///
/// Incrémenter cette constante à chaque modification incompatible du schéma
/// (ajout/suppression de colonne non nullable, changement de type, etc.).
/// Cela provoquera la suppression et la recréation automatique de la DB au démarrage.
const SCHEMA_VERSION: u32 = 1;
/// Gestionnaire de persistance (une base pour toutes les playlists)
pub struct PersistenceManager {
conn: Arc<Mutex<Connection>>,
@@ -26,6 +33,26 @@ impl PersistenceManager {
})?;
}
// Vérifier la version du schéma — supprimer la DB si incompatible
if db_path.exists() {
if let Ok(conn) = Connection::open(db_path) {
let version: u32 = conn
.query_row("PRAGMA user_version", [], |r| r.get(0))
.unwrap_or(0);
if version != SCHEMA_VERSION {
drop(conn);
tracing::warn!(
"Playlist DB schema version mismatch (found {}, expected {}), recreating",
version,
SCHEMA_VERSION
);
std::fs::remove_file(db_path).map_err(|e| {
crate::Error::PersistenceError(format!("Failed to remove old DB: {}", e))
})?;
}
}
}
let conn = Connection::open(db_path).map_err(|e| {
crate::Error::PersistenceError(format!("Failed to open database: {}", e))
})?;
@@ -38,6 +65,8 @@ impl PersistenceManager {
role TEXT NOT NULL,
cover_pk TEXT,
artist TEXT,
source TEXT,
source_version TEXT,
max_size INTEGER,
default_ttl_secs INTEGER,
created_at INTEGER NOT NULL,
@@ -75,6 +104,12 @@ impl PersistenceManager {
)
.map_err(|e| crate::Error::PersistenceError(format!("Failed to create index: {}", e)))?;
// Inscrire la version du schéma
conn.execute_batch(&format!("PRAGMA user_version = {}", SCHEMA_VERSION))
.map_err(|e| {
crate::Error::PersistenceError(format!("Failed to set schema version: {}", e))
})?;
Ok(Self {
conn: Arc::new(Mutex::new(conn)),
})
@@ -88,6 +123,8 @@ impl PersistenceManager {
role: &PlaylistRole,
cover_pk: Option<&str>,
artist: Option<&str>,
source: Option<&str>,
source_version: Option<&str>,
config: &PlaylistConfig,
tracks: &VecDeque<Arc<Record>>,
) -> Result<()> {
@@ -100,16 +137,18 @@ impl PersistenceManager {
// Upsert playlist metadata
conn.execute(
"INSERT OR REPLACE INTO playlists (id, title, role, cover_pk, artist, max_size, default_ttl_secs, created_at, last_modified)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7,
COALESCE((SELECT created_at FROM playlists WHERE id = ?1), ?8),
?8)",
"INSERT OR REPLACE INTO playlists (id, title, role, cover_pk, artist, source, source_version, max_size, default_ttl_secs, created_at, last_modified)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9,
COALESCE((SELECT created_at FROM playlists WHERE id = ?1), ?10),
?10)",
params![
id,
title,
role.as_str(),
cover_pk,
artist,
source,
source_version,
config.max_size.map(|s| s as i64),
config.default_ttl.map(|d| d.as_secs() as i64),
now_nanos,
@@ -154,6 +193,8 @@ impl PersistenceManager {
PlaylistConfig,
Option<String>,
Option<String>,
Option<String>,
Option<String>,
VecDeque<Arc<Record>>,
)>,
> {
@@ -161,7 +202,7 @@ impl PersistenceManager {
// Charger les métadonnées
let mut stmt = conn.prepare(
"SELECT title, role, cover_pk, artist, max_size, default_ttl_secs FROM playlists WHERE id = ?1",
"SELECT title, role, cover_pk, artist, source, source_version, max_size, default_ttl_secs FROM playlists WHERE id = ?1",
)
.map_err(|e| {
crate::Error::PersistenceError(format!("Failed to prepare statement: {}", e))
@@ -172,8 +213,10 @@ impl PersistenceManager {
let role_raw: String = row.get(1)?;
let cover_pk: Option<String> = row.get(2)?;
let artist: Option<String> = row.get(3)?;
let max_size: Option<i64> = row.get(4)?;
let default_ttl_secs: Option<i64> = row.get(5)?;
let source: Option<String> = row.get(4)?;
let source_version: Option<String> = row.get(5)?;
let max_size: Option<i64> = row.get(6)?;
let default_ttl_secs: Option<i64> = row.get(7)?;
Ok((
title,
@@ -185,10 +228,12 @@ impl PersistenceManager {
},
cover_pk,
artist,
source,
source_version,
))
});
let (title, role, config, cover_pk, artist) = match result {
let (title, role, config, cover_pk, artist, source, source_version) = match result {
Ok(data) => data,
Err(rusqlite::Error::QueryReturnedNoRows) => return Ok(None),
Err(e) => {
@@ -231,7 +276,7 @@ impl PersistenceManager {
tracks.push_back(Arc::new(record));
}
Ok(Some((title, role, config, cover_pk, artist, tracks)))
Ok(Some((title, role, config, cover_pk, artist, source, source_version, tracks)))
}
/// Supprime une playlist

View File

@@ -36,6 +36,10 @@ pub struct Playlist {
role: RwLock<PlaylistRole>,
cover_pk: RwLock<Option<String>>,
artist: RwLock<Option<String>>,
/// Source externe (ex: "qobuz")
source: RwLock<Option<String>>,
/// Version de la source (ex: timestamp updated_at de Qobuz)
source_version: RwLock<Option<String>>,
state: Arc<AtomicU8>,
pub core: Arc<RwLock<PlaylistCore>>,
pub persistent: bool,
@@ -59,6 +63,8 @@ impl Playlist {
role: RwLock::new(role),
cover_pk: RwLock::new(cover_pk),
artist: RwLock::new(None),
source: RwLock::new(None),
source_version: RwLock::new(None),
state: Arc::new(AtomicU8::new(PlaylistState::Active as u8)),
core: Arc::new(RwLock::new(PlaylistCore::new(config))),
persistent,
@@ -127,6 +133,28 @@ impl Playlist {
self.touch().await;
}
/// Retourne la source externe associée à la playlist.
pub async fn source(&self) -> Option<String> {
self.source.read().await.clone()
}
/// Modifie la source de la playlist.
pub async fn set_source(&self, value: Option<String>) {
*self.source.write().await = value;
self.touch().await;
}
/// Retourne la version de la source (ex: updated_at timestamp).
pub async fn source_version(&self) -> Option<String> {
self.source_version.read().await.clone()
}
/// Modifie la version de la source.
pub async fn set_source_version(&self, value: Option<String>) {
*self.source_version.write().await = value;
self.touch().await;
}
/// Timestamp du dernier changement
pub async fn last_change(&self) -> SystemTime {
*self.last_change.read().await

View File

@@ -32,6 +32,8 @@ pub(crate) struct AlbumResponse {
#[serde(default)]
release_date_original: Option<String>,
#[serde(default)]
released_at: Option<i64>,
#[serde(default)]
image: Option<ImageResponse>,
#[serde(default = "default_streamable")]
streamable: bool,
@@ -120,6 +122,8 @@ pub(crate) struct PlaylistResponse {
#[serde(default)]
owner: Option<OwnerResponse>,
#[serde(default)]
updated_at: Option<i64>,
#[serde(default)]
tracks: Option<PaginatedResponse<TrackResponse>>,
}
@@ -317,6 +321,16 @@ impl QobuzApi {
.collect())
}
/// Récupère les métadonnées d'une playlist sans les tracks (appel léger)
///
/// Utile pour vérifier `updated_at` avant de décider si le cache est valide.
pub async fn get_playlist_metadata(&self, playlist_id: &str) -> Result<Playlist> {
debug!("Fetching playlist metadata {}", playlist_id);
let params = [("playlist_id", playlist_id)];
let response: PlaylistResponse = self.get("/playlist/get", &params).await?;
Ok(Self::parse_playlist(response))
}
/// Récupère les détails d'une playlist
pub async fn get_playlist(&self, playlist_id: &str) -> Result<Playlist> {
debug!("Fetching playlist {}", playlist_id);
@@ -328,18 +342,36 @@ impl QobuzApi {
/// Récupère les tracks d'une playlist
pub async fn get_playlist_tracks(&self, playlist_id: &str) -> Result<Vec<Track>> {
debug!("Fetching tracks for playlist {}", playlist_id);
let params = [("playlist_id", playlist_id), ("extra", "tracks")];
let response: PlaylistResponse = self.get("/playlist/get", &params).await?;
const PAGE_SIZE: u32 = 50;
let mut all_tracks = Vec::new();
let mut offset = 0u32;
if let Some(tracks) = response.tracks {
Ok(tracks
.items
.into_iter()
.map(|t| Self::parse_track(t, None))
.collect())
} else {
Ok(Vec::new())
loop {
let offset_str = offset.to_string();
let limit_str = PAGE_SIZE.to_string();
let params = [
("playlist_id", playlist_id),
("extra", "tracks"),
("offset", offset_str.as_str()),
("limit", limit_str.as_str()),
];
let response: PlaylistResponse = self.get("/playlist/get", &params).await?;
if let Some(tracks) = response.tracks {
let total = tracks.total.unwrap_or(0);
let count = tracks.items.len() as u32;
all_tracks.extend(tracks.items.into_iter().map(|t| Self::parse_track(t, None)));
offset += count;
if count == 0 || offset >= total {
break;
}
} else {
break;
}
}
debug!("Fetched {} tracks total for playlist {}", all_tracks.len(), playlist_id);
Ok(all_tracks)
}
/// Récupère la liste des genres
@@ -493,6 +525,7 @@ impl QobuzApi {
tracks_count: response.tracks_count,
duration: response.duration,
release_date: response.release_date_original,
released_at: response.released_at,
image: response.image.and_then(|i| i.large),
image_cached: None,
streamable: response.streamable,
@@ -547,6 +580,7 @@ impl QobuzApi {
image: response.images300.and_then(|imgs| imgs.first().cloned()),
image_cached: None,
is_public: response.is_public,
updated_at: response.updated_at,
owner: response.owner.map(|o| PlaylistOwner {
id: o.id.parse().unwrap_or(0),
name: o.name,

View File

@@ -17,6 +17,49 @@ use std::sync::RwLock;
use std::time::Duration;
use tracing::debug;
/// Si la variable d'environnement `QOBUZ_DEBUG_DIR` est définie, sauvegarde
/// le JSON brut de chaque réponse API dans ce dossier.
fn debug_save_response(endpoint: &str, params: &[(&str, &str)], text: &str) {
let Ok(dir) = std::env::var("QOBUZ_DEBUG_DIR") else {
return;
};
let dir = std::path::Path::new(&dir);
if let Err(e) = std::fs::create_dir_all(dir) {
debug!("QOBUZ_DEBUG_DIR: cannot create dir: {}", e);
return;
}
// Construire un nom de fichier lisible : endpoint + params clés
let endpoint_slug = endpoint.trim_start_matches('/').replace('/', "_");
let param_slug: String = params
.iter()
.filter(|(k, _)| !["app_id", "user_auth_token", "request_ts", "request_sig"].contains(k))
.map(|(k, v)| format!("{}-{}", k, v.chars().take(20).collect::<String>()))
.collect::<Vec<_>>()
.join("_");
let ts = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis())
.unwrap_or(0);
let filename = format!("{}_{}_{}.json", endpoint_slug, param_slug, ts);
let path = dir.join(&filename);
// Pretty-print si possible
let content = serde_json::from_str::<Value>(text)
.ok()
.and_then(|v| serde_json::to_string_pretty(&v).ok())
.unwrap_or_else(|| text.to_string());
if let Err(e) = std::fs::write(&path, content) {
debug!("QOBUZ_DEBUG_DIR: cannot write {}: {}", filename, e);
} else {
debug!("QOBUZ_DEBUG_DIR: saved {}", filename);
}
}
pub use spoofer::Spoofer;
/// URL de base de l'API Qobuz
@@ -270,7 +313,7 @@ impl QobuzApi {
// Envoyer la requête
let response = request.send().await?;
self.handle_response(response, endpoint).await
self.handle_response(response, endpoint, params).await
}
/// Traite la réponse HTTP
@@ -278,6 +321,7 @@ impl QobuzApi {
&self,
response: Response,
endpoint: &str,
params: &[(&str, &str)],
) -> Result<T> {
let status = response.status();
let status_code = status.as_u16();
@@ -294,6 +338,7 @@ impl QobuzApi {
}
let text = response.text().await?;
debug_save_response(endpoint, params, &text);
// Vérifier si la réponse contient une erreur Qobuz
if let Ok(json) = serde_json::from_str::<Value>(&text) {

View File

@@ -54,6 +54,9 @@ pub struct Album {
/// Date de sortie (format ISO 8601)
#[serde(default)]
pub release_date: Option<String>,
/// Timestamp Unix de sortie (released_at)
#[serde(default)]
pub released_at: Option<i64>,
/// URL de l'image de couverture
#[serde(default)]
pub image: Option<String>,
@@ -141,6 +144,9 @@ pub struct Playlist {
/// Indique si c'est une playlist publique
#[serde(default)]
pub is_public: bool,
/// Timestamp de dernière modification (Unix epoch secondes)
#[serde(default)]
pub updated_at: Option<i64>,
/// Propriétaire de la playlist
#[serde(default)]
pub owner: Option<PlaylistOwner>,

View File

@@ -14,10 +14,9 @@ use pmosource::SourceCacheManager;
use pmosource::{async_trait, BrowseResult, MusicSource, MusicSourceError, Result};
use serde_json::json;
use std::sync::Arc;
use std::time::{Duration, SystemTime};
use std::time::SystemTime;
/// TTL pour les playlists d'albums (7 jours)
const ALBUM_PLAYLIST_TTL: Duration = Duration::from_secs(7 * 24 * 3600);
/// Trait pour les types dont on peut cacher la cover image.
trait CoverCacheable {
@@ -563,54 +562,41 @@ impl QobuzSource {
}
/// Vérifie si une playlist d'album existe et est valide (non expirée ET non vide)
async fn is_album_playlist_valid(&self, playlist_id: &str) -> Result<bool> {
/// Retourne la source_version d'une pmoplaylist si elle existe en mémoire ou en DB.
async fn cached_source_version(&self, playlist_id: &str) -> Option<String> {
let playlist_manager = pmoplaylist::PlaylistManager();
if !playlist_manager.exists(playlist_id).await {
return Ok(false);
}
// Vérifier l'âge
match playlist_manager.get_playlist_age(playlist_id).await {
Ok(Some(age)) if age < ALBUM_PLAYLIST_TTL => {
// Playlist non expirée, vérifier qu'elle contient des tracks
match playlist_manager.get_read_handle(playlist_id).await {
Ok(reader) => {
let count = reader.remaining().await.unwrap_or(0);
Ok(count > 0) // Valide seulement si non vide
}
Err(_) => Ok(false),
}
}
_ => Ok(false),
return None;
}
playlist_manager
.get_read_handle(playlist_id)
.await
.ok()?
.source_version()
.await
}
/// Adapte les Items d'une playlist pour correspondre au schéma UPnP Qobuz
async fn adapt_playlist_items_to_qobuz(
/// Adapte des Items DIDL issus d'une pmoplaylist pour le schéma UPnP Qobuz.
///
/// - Résout l'ID de la track (`qobuz:track:<id>`) depuis les métadonnées cache
/// - Rend absolues les URLs relatives (resource + cover)
/// - Affecte le `parent_id` fourni
async fn adapt_items_to_qobuz(
&self,
items: Vec<Item>,
album_id: &str,
parent_id: &str,
) -> Result<Vec<Item>> {
use tracing::warn;
let parent_id = format!("qobuz:album:{}", album_id);
let mut adapted = Vec::with_capacity(items.len());
for mut item in items {
// Extraire cache_pk depuis l'URL du resource
let cache_pk = if let Some(resource) = item.resources.first() {
resource
.url
.strip_prefix("/audio/flac/")
.map(|s| s.to_string())
} else {
None
};
let cache_pk = item
.resources
.first()
.and_then(|r| r.url.strip_prefix("/audio/flac/").map(|s| s.to_string()));
if let Some(pk) = cache_pk {
// Récupérer track_id depuis metadata
if let Ok(Some(track_id_value)) = self
.inner
.cache_manager
@@ -625,9 +611,6 @@ impl QobuzSource {
warn!("No qobuz_track_id metadata for {}", pk);
}
// Convertir URL relative en URL absolue
// From: /audio/flac/QOBUZ:123
// To: http://192.168.0.138:8080/audio/flac/QOBUZ:123
if let Some(resource) = item.resources.first_mut() {
if resource.url.starts_with('/') {
resource.url = format!("{}{}", self.inner.base_url, resource.url);
@@ -635,21 +618,59 @@ impl QobuzSource {
}
}
// Convertir l'URL de la cover en URL absolue si elle est relative
if let Some(art) = item.album_art.as_mut() {
if art.starts_with('/') {
*art = format!("{}{}", self.inner.base_url, art);
}
}
item.parent_id = parent_id.clone();
item.parent_id = parent_id.to_string();
adapted.push(item);
}
Ok(adapted)
}
/// Récupère ou crée une playlist lazy pour un album
/// Ouvre un WriteHandle sur une pmoplaylist persistante, en la créant si absente
/// ou en la vidant (flush) si elle existe déjà.
async fn get_or_flush_playlist_writer(
&self,
playlist_id: &str,
role: pmoplaylist::PlaylistRole,
) -> Result<pmoplaylist::WriteHandle> {
let playlist_manager = pmoplaylist::PlaylistManager();
if playlist_manager.exists(playlist_id).await {
let writer = playlist_manager
.get_persistent_write_handle(playlist_id.to_string())
.await
.map_err(|e| MusicSourceError::PlaylistError(e.to_string()))?;
writer
.flush()
.await
.map_err(|e| MusicSourceError::PlaylistError(e.to_string()))?;
Ok(writer)
} else {
playlist_manager
.create_persistent_playlist_with_role(playlist_id.to_string(), role)
.await
.map_err(|e| MusicSourceError::PlaylistError(e.to_string()))
}
}
/// Met en cache la cover si l'URL est présente.
async fn cache_cover_opt(&self, image_url: Option<&str>) -> Option<String> {
self.inner
.cache_manager
.cache_cover(image_url?)
.await
.ok()
}
/// Récupère ou crée une playlist lazy pour un album, avec invalidation par version.
///
/// `source_version` = `"{released_at}_{tracks_count}"` — change uniquement en cas de
/// réédition avec pistes supplémentaires. Cache valide à vie si la version correspond.
async fn get_or_create_album_playlist_items(
&self,
album_id: &str,
@@ -659,30 +680,11 @@ impl QobuzSource {
let playlist_id = format!("qobuz-album-{}", album_id);
let playlist_manager = pmoplaylist::PlaylistManager();
let parent_id = format!("qobuz:album:{}", album_id);
// Vérifier validité (existe ET non expirée)
let is_valid = self.is_album_playlist_valid(&playlist_id).await?;
let cached_version = self.cached_source_version(&playlist_id).await;
if is_valid {
debug!("Album playlist {} found and valid", playlist_id);
let reader = playlist_manager
.get_read_handle(&playlist_id)
.await
.map_err(|e| MusicSourceError::PlaylistError(e.to_string()))?;
let items = reader
.to_items(limit)
.await
.map_err(|e| MusicSourceError::PlaylistError(e.to_string()))?;
return self.adapt_playlist_items_to_qobuz(items, album_id).await;
}
// Playlist invalide/inexistante : (re)créer
info!("Album playlist {} creating/refreshing", playlist_id);
// 1. Métadonnées album
// Charger les métadonnées album (contient déjà les tracks)
let album = self
.inner
.client
@@ -690,47 +692,49 @@ impl QobuzSource {
.await
.map_err(|e| MusicSourceError::BrowseError(e.to_string()))?;
// 2. Cache cover
let cover_pk = if let Some(ref image_url) = album.image {
self.inner.cache_manager.cache_cover(image_url).await.ok()
} else {
None
};
let album_version = Some(format!(
"{}_{}",
album.released_at.unwrap_or(0),
album.tracks_count.unwrap_or(0)
));
// 3. Créer ou récupérer playlist
let writer = if playlist_manager.exists(&playlist_id).await {
let writer = playlist_manager
.get_persistent_write_handle(playlist_id.clone())
if cached_version.is_some() && cached_version == album_version {
debug!("Album playlist {} cache valid (version {:?})", playlist_id, album_version);
let reader = playlist_manager
.get_read_handle(&playlist_id)
.await
.map_err(|e| MusicSourceError::PlaylistError(e.to_string()))?;
writer
.flush()
let items = reader
.to_items(limit)
.await
.map_err(|e| MusicSourceError::PlaylistError(e.to_string()))?;
return self.adapt_items_to_qobuz(items, &parent_id).await;
}
writer
} else {
playlist_manager
.create_persistent_playlist_with_role(
playlist_id.clone(),
pmoplaylist::PlaylistRole::Album,
)
.await
.map_err(|e| MusicSourceError::PlaylistError(e.to_string()))?
};
info!("Album playlist {} creating/refreshing (version {:?})", playlist_id, album_version);
let cover_pk = self.cache_cover_opt(album.image.as_deref()).await;
let writer = self
.get_or_flush_playlist_writer(&playlist_id, pmoplaylist::PlaylistRole::Album)
.await?;
// 4. Métadonnées playlist
writer
.set_title(album.title.clone())
.await
.map_err(|e| MusicSourceError::PlaylistError(e.to_string()))?;
writer
.set_artist(Some(album.artist.name.clone()))
.await
.map_err(|e| MusicSourceError::PlaylistError(e.to_string()))?;
writer
.set_source(Some("qobuz".to_string()))
.await
.map_err(|e| MusicSourceError::PlaylistError(e.to_string()))?;
writer
.set_source_version(album_version)
.await
.map_err(|e| MusicSourceError::PlaylistError(e.to_string()))?;
if let Some(pk) = cover_pk {
writer
.set_cover_pk(Some(pk))
@@ -741,22 +745,122 @@ impl QobuzSource {
// IMPORTANT: Libérer le write lock avant d'appeler add_album_to_playlist
drop(writer);
// 5. Ajouter tracks (réutilise add_album_to_playlist existant)
self.add_album_to_playlist(&playlist_id, album_id).await?;
// 6. Récupérer items
let reader = playlist_manager
.get_read_handle(&playlist_id)
.await
.map_err(|e| MusicSourceError::PlaylistError(e.to_string()))?;
let items = reader
.to_items(limit)
.await
.map_err(|e| MusicSourceError::PlaylistError(e.to_string()))?;
// 7. Adapter IDs
self.adapt_playlist_items_to_qobuz(items, album_id).await
self.adapt_items_to_qobuz(items, &parent_id).await
}
/// Récupère ou crée une pmoplaylist pour une playlist Qobuz, avec invalidation par updated_at.
///
/// - Si la playlist est en cache et que `updated_at` n'a pas changé → retourne le cache.
/// - Sinon → recharge toutes les tracks depuis Qobuz et reconstruit la pmoplaylist.
async fn get_or_create_qobuz_playlist_items(
&self,
qobuz_playlist_id: &str,
) -> Result<Vec<Item>> {
use tracing::{debug, info};
let pmo_playlist_id = format!("qobuz-playlist-{}", qobuz_playlist_id);
let playlist_manager = pmoplaylist::PlaylistManager();
let cached_version = self.cached_source_version(&pmo_playlist_id).await;
// Obtenir les métadonnées depuis Qobuz pour comparer updated_at
let qobuz_meta: crate::models::Playlist = self
.inner
.client
.get_playlist(qobuz_playlist_id)
.await
.map_err(|e: crate::error::QobuzError| MusicSourceError::BrowseError(e.to_string()))?;
let qobuz_version = qobuz_meta.updated_at.map(|t: i64| t.to_string());
// Cache valide si version correspond
let cache_valid = cached_version.is_some() && cached_version == qobuz_version;
let parent_id = format!("qobuz:playlist:{}", qobuz_playlist_id);
if cache_valid {
debug!("Qobuz playlist {} cache valid (version {:?})", pmo_playlist_id, qobuz_version);
let reader = playlist_manager
.get_read_handle(&pmo_playlist_id)
.await
.map_err(|e| MusicSourceError::PlaylistError(e.to_string()))?;
let (items, total) = reader
.to_items_paged(0, usize::MAX)
.await
.map_err(|e| MusicSourceError::PlaylistError(e.to_string()))?;
// Si des PKs ne sont pas enregistrés dans le cache audio (e.g. playlist créée
// avec l'ancien code), on force un refresh pour les ré-enregistrer proprement.
if items.len() == total {
return self.adapt_items_to_qobuz(items, &parent_id).await;
}
info!(
"Qobuz playlist {} has {}/{} valid items, forcing refresh to repair missing entries",
pmo_playlist_id, items.len(), total
);
}
info!("Qobuz playlist {} creating/refreshing (version {:?})", pmo_playlist_id, qobuz_version);
let tracks = self
.inner
.client
.get_playlist_tracks(qobuz_playlist_id)
.await
.map_err(|e| MusicSourceError::BrowseError(e.to_string()))?;
let cover_pk = self.cache_cover_opt(qobuz_meta.image.as_deref()).await;
let writer = self
.get_or_flush_playlist_writer(&pmo_playlist_id, pmoplaylist::PlaylistRole::Source)
.await?;
writer
.set_title(qobuz_meta.name.clone())
.await
.map_err(|e| MusicSourceError::PlaylistError(e.to_string()))?;
writer
.set_source(Some("qobuz".to_string()))
.await
.map_err(|e| MusicSourceError::PlaylistError(e.to_string()))?;
writer
.set_source_version(qobuz_version)
.await
.map_err(|e| MusicSourceError::PlaylistError(e.to_string()))?;
if let Some(pk) = cover_pk {
writer
.set_cover_pk(Some(pk))
.await
.map_err(|e| MusicSourceError::PlaylistError(e.to_string()))?;
}
let lazy_pks = self.register_tracks_lazy(&tracks).await;
writer
.push_lazy_batch(lazy_pks)
.await
.map_err(|e| MusicSourceError::PlaylistError(e.to_string()))?;
drop(writer);
let reader = playlist_manager
.get_read_handle(&pmo_playlist_id)
.await
.map_err(|e| MusicSourceError::PlaylistError(e.to_string()))?;
let (items, _total) = reader
.to_items_paged(0, usize::MAX)
.await
.map_err(|e| MusicSourceError::PlaylistError(e.to_string()))?;
self.adapt_items_to_qobuz(items, &parent_id).await
}
/// Increment update counter (called on catalog changes)
@@ -937,6 +1041,104 @@ impl QobuzSource {
.await
}
/// Enregistre une liste de tracks comme lazy entries dans l'audio cache (en parallèle).
///
/// Contrairement à `add_track_lazy`, cette méthode ne fait PAS d'appel à `get_stream_url`
/// (coûteux pour de grandes playlists). L'URL audio est résolue à la demande via
/// `QobuzLazyProvider` lors de la première lecture.
///
/// Pour chaque track : cache la cover, enregistre la lazy entry, stocke les métadonnées.
/// Retourne la liste des lazy PKs enregistrés avec succès.
async fn register_tracks_lazy(&self, tracks: &[crate::models::Track]) -> Vec<String> {
// Limite la concurrence pour ne pas saturer l'API Qobuz ni la connexion réseau.
// Les covers déjà cachées sont retournées immédiatement (pas d'HTTP), donc même
// 600 tracks ne génèrent que ~N_albums_uniques téléchargements réels.
let sem = std::sync::Arc::new(tokio::sync::Semaphore::new(16));
// On attache l'index original à chaque future pour pouvoir retrier dans l'ordre
// d'origine après complétion parallèle (JoinSet retourne dans l'ordre de fin).
let futs: Vec<_> = tracks.iter().enumerate().map(|(idx, track)| {
let source = self.clone();
let sem = sem.clone();
let track_id = track.id.clone();
let track_title = track.title.clone();
let cover_image_url = track.album.as_ref().and_then(|a| a.image.clone());
let metadata = pmoaudiocache::AudioMetadata {
title: Some(track.title.clone()),
artist: track.performer.as_ref().map(|p| p.name.clone()),
album: track.album.as_ref().map(|a| a.title.clone()),
duration_secs: Some(track.duration as u64),
year: track.album.as_ref().and_then(|a| {
a.release_date.as_ref().and_then(|d| d.split('-').next()?.parse().ok())
}),
track_number: Some(track.track_number),
track_total: track.album.as_ref().and_then(|a| a.tracks_count),
disc_number: Some(track.media_number),
disc_total: None,
genre: track.album.as_ref().and_then(|a| {
if !a.genres.is_empty() { Some(a.genres.join(", ")) } else { None }
}),
sample_rate: track.sample_rate,
channels: track.channels,
bitrate: None,
conversion: None,
};
async move {
let _permit = sem.acquire().await.ok()?;
let lazy_pk = format!("QOBUZ:{}", track_id);
// 1. Cache cover eagerly
let cover_pk = if let Some(ref url) = cover_image_url {
source.inner.cache_manager.cache_cover(url).await.ok()
} else {
None
};
// 2. Register lazy entry + set cover_pk + seed metadata.
// Si le performer est absent (réponse API incomplète), on passe None pour
// que le provider appelle get_track et récupère les métadonnées complètes.
if metadata.artist.is_none() {
tracing::warn!(
"register_tracks_lazy: no performer for track {} (id={}), will call provider",
track_title, track_id
);
}
if cover_pk.is_none() && cover_image_url.is_some() {
tracing::warn!(
"register_tracks_lazy: cover download failed for track {} (id={}), will call provider for cover",
track_title, track_id
);
}
let meta_hint = if metadata.artist.is_some() { Some(metadata) } else { None };
match source.inner.cache_manager
.cache_audio_lazy_with_provider(&lazy_pk, meta_hint, cover_pk)
.await
{
Ok(pk) => {
let _ = source.inner.cache_manager.set_audio_metadata(
&pk, "qobuz_track_id", json!(track_id),
);
Some((idx, pk))
}
Err(e) => {
tracing::warn!("Failed to register lazy track {}: {}", track_title, e);
None
}
}
}
}).collect();
// Retrier par index original pour conserver l'ordre de la playlist Qobuz
let mut results: Vec<(usize, String)> = tokio::task::JoinSet::from_iter(futs)
.join_all()
.await
.into_iter()
.flatten()
.collect();
results.sort_unstable_by_key(|(i, _)| *i);
results.into_iter().map(|(_, pk)| pk).collect()
}
/// Cache les covers d'une liste d'artistes en parallèle.
async fn cache_artist_covers(&self, artists: Vec<crate::models::Artist>) -> Vec<crate::models::Artist> {
let futs: Vec<_> = artists.into_iter().map(|mut artist| {
@@ -1623,22 +1825,9 @@ impl MusicSource for QobuzSource {
}
ObjectIdType::Playlist(playlist_id) => {
let tracks = self
.inner
.client
.get_playlist_tracks(&playlist_id)
.await
.map_err(|e| MusicSourceError::BrowseError(e.to_string()))?;
let tracks = self.cache_covers(tracks).await;
let items: Vec<Item> = tracks
.into_iter()
.filter_map(|track| {
track
.to_didl_item(&format!("qobuz:playlist:{}", playlist_id))
.ok()
})
.collect();
let items = self
.get_or_create_qobuz_playlist_items(&playlist_id)
.await?;
Ok(BrowseResult::Items(items))
}

View File

@@ -26,9 +26,6 @@ use tokio::sync::RwLock;
#[cfg(feature = "cache")]
use pmocovers::Cache as CoverCache;
#[cfg(feature = "cache")]
use pmocache::cache_trait::FileCache;
// ============================================================================
// CachedMetadata
// ============================================================================
@@ -583,11 +580,7 @@ impl MetadataCache {
}
Err(discover_err) => {
#[cfg(feature = "logging")]
tracing::error!(
"Rediscovery failed for '{}': {}",
slug,
discover_err
);
tracing::error!("Rediscovery failed for '{}': {}", slug, discover_err);
return Err(e);
}
}
@@ -697,12 +690,7 @@ impl MetadataCache {
match self.client.rediscover_station(slug).await {
Ok((id, url)) => {
#[cfg(feature = "logging")]
tracing::info!(
"Re-discovered station '{}': id={}, url={}",
slug,
id,
url
);
tracing::info!("Re-discovered station '{}': id={}, url={}", slug, id, url);
self.persist_station_mapping();
// Invalider le cache mémoire pour forcer un refresh des métadonnées
let mut cache = self.cache.write().await;
@@ -801,11 +789,18 @@ mod tests {
assert!(base.is_expired());
// end_time dans le futur → non expiré
let not_expired = CachedMetadata { end_time: Some(now + 3600), ..base.clone() };
let not_expired = CachedMetadata {
end_time: Some(now + 3600),
..base.clone()
};
assert!(!not_expired.is_expired());
// Pas de end_time, fetched_at récent → non expiré (fallback TTL)
let no_end_fresh = CachedMetadata { end_time: None, fetched_at: now, ..base.clone() };
let no_end_fresh = CachedMetadata {
end_time: None,
fetched_at: now,
..base.clone()
};
assert!(!no_end_fresh.is_expired());
// Pas de end_time, fetched_at ancien → expiré

View File

@@ -328,6 +328,7 @@ impl SourceCacheManager {
None
};
let metadata_is_some = metadata.is_some();
let mut final_metadata = metadata;
if final_metadata.is_none() {
if let Some(data) = provider_data.as_ref() {
@@ -352,10 +353,25 @@ impl SourceCacheManager {
if let Some(meta) = final_metadata.as_ref() {
self.seed_audio_metadata(lazy_pk, meta);
} else {
#[cfg(feature = "server")]
tracing::warn!(
"cache_audio_lazy_with_provider: no metadata to seed for {} \
(meta_hint={}, provider_data={})",
lazy_pk,
metadata_is_some,
provider_data.as_ref().map(|d| d.metadata.is_some()).unwrap_or(false)
);
}
if let Some(cover_pk) = final_cover_pk {
let _ = self.set_audio_metadata(lazy_pk, "cover_pk", json!(cover_pk));
} else {
#[cfg(feature = "server")]
tracing::debug!(
"cache_audio_lazy_with_provider: no cover_pk seeded for {}",
lazy_pk
);
}
Ok(lazy_pk.to_string())

View File

@@ -1 +1 @@
0.3.26
0.3.27