From a30186485f62912564aa136ae226d66688775131 Mon Sep 17 00:00:00 2001
From: Eric Coissac
Date: Thu, 11 Jun 2026 14:47:40 +0200
Subject: [PATCH 1/5] feat: make qobuz register concurrency configurable
Replace the hardcoded semaphore capacity of 16 with a configurable `register_concurrency` setting (defaulting to 4). This mitigates SQLite write contention and optimizes concurrent API and network requests during parallel track caching.
---
pmoqobuz/src/config_ext.rs | 18 ++++++++++++++++++
pmoqobuz/src/source.rs | 15 +++++++++++----
2 files changed, 29 insertions(+), 4 deletions(-)
diff --git a/pmoqobuz/src/config_ext.rs b/pmoqobuz/src/config_ext.rs
index f1535883..fb96f767 100644
--- a/pmoqobuz/src/config_ext.rs
+++ b/pmoqobuz/src/config_ext.rs
@@ -245,6 +245,15 @@ pub trait QobuzConfigExt {
/// Persiste la version du bundle après une extraction réussie.
fn set_qobuz_bundle_version(&self, version: &str) -> Result<()>;
+
+ /// Nombre de workers concurrents pour l'enregistrement des tracks en cache.
+ ///
+ /// Contrôle le semaphore dans `register_tracks_lazy` : plus la valeur est
+ /// haute, plus les covers sont téléchargées en parallèle, mais plus la
+ /// contention sur le mutex SQLite est forte.
+ ///
+ /// Défaut : 4 (adapté à une machine sous contrainte mémoire / Docker).
+ fn get_qobuz_register_concurrency(&self) -> usize;
}
impl QobuzConfigExt for Config {
@@ -504,4 +513,13 @@ impl QobuzConfigExt for Config {
Value::String(version.to_string()),
)
}
+
+ fn get_qobuz_register_concurrency(&self) -> usize {
+ match self.get_value(&["accounts", "qobuz", "register_concurrency"]) {
+ Ok(Value::Number(n)) if n.as_u64().unwrap_or(0) >= 1 => {
+ n.as_u64().unwrap() as usize
+ }
+ _ => 4,
+ }
+ }
}
diff --git a/pmoqobuz/src/source.rs b/pmoqobuz/src/source.rs
index 905d70bd..c8c21a40 100644
--- a/pmoqobuz/src/source.rs
+++ b/pmoqobuz/src/source.rs
@@ -115,6 +115,9 @@ struct QobuzSourceInner {
/// Base URL for streaming server (e.g., "http://192.168.0.138:8080")
base_url: String,
+ /// Nombre de workers concurrents pour register_tracks_lazy (configurable)
+ register_concurrency: usize,
+
/// Update tracking
update_counter: tokio::sync::RwLock,
last_change: tokio::sync::RwLock,
@@ -142,15 +145,18 @@ impl QobuzSource {
/// Returns an error if the caches are not initialized in the registry
#[cfg(feature = "server")]
pub fn from_registry(client: QobuzClient, base_url: impl Into) -> Result {
+ use crate::config_ext::QobuzConfigExt;
let cache_manager = SourceCacheManager::from_registry("qobuz".to_string())?;
let client = Arc::new(client);
cache_manager.register_lazy_provider(Arc::new(QobuzLazyProvider::new(client.clone())));
+ let register_concurrency = pmoconfig::get_config().get_qobuz_register_concurrency();
Ok(Self {
inner: Arc::new(QobuzSourceInner {
client,
cache_manager,
base_url: base_url.into(),
+ register_concurrency,
update_counter: tokio::sync::RwLock::new(0),
last_change: tokio::sync::RwLock::new(SystemTime::now()),
}),
@@ -180,6 +186,7 @@ impl QobuzSource {
client,
cache_manager,
base_url: base_url.into(),
+ register_concurrency: 4,
update_counter: tokio::sync::RwLock::new(0),
last_change: tokio::sync::RwLock::new(SystemTime::now()),
}),
@@ -1062,10 +1069,10 @@ impl QobuzSource {
/// 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 {
- // 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));
+ // Configurable via accounts.qobuz.register_concurrency (défaut 4).
+ // SQLite sérialise les écritures — au-delà de ~4 workers on accumule
+ // des threads en attente du mutex DB sans gain de débit.
+ let sem = std::sync::Arc::new(tokio::sync::Semaphore::new(self.inner.register_concurrency));
// 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).
--
2.49.1
From 1d3a9379a15ea69209e4daaddabf276d8c8345fe Mon Sep 17 00:00:00 2001
From: Eric Coissac
Date: Thu, 11 Jun 2026 14:49:44 +0200
Subject: [PATCH 2/5] feat(qobuz): implement concurrent playlist pagination
Parallelize Qobuz playlist track fetching by increasing the page size to 500 and processing remaining pages concurrently via `futures::try_join_all` with a configurable semaphore (default 3). Results are offset-sorted to preserve original order. Adds a `page_concurrency` configuration option, updates the API client initialization, and introduces the `futures` dependency. This reduces large playlist latency from ~1.6s to ~0.7s.
---
.../pmoqobuz_ameliorations_qbz.md | 23 ++--
Cargo.lock | 1 +
pmoqobuz/Cargo.toml | 1 +
pmoqobuz/src/api/catalog.rs | 127 +++++++++++++-----
pmoqobuz/src/api/mod.rs | 8 ++
pmoqobuz/src/client.rs | 2 +
pmoqobuz/src/config_ext.rs | 18 +++
7 files changed, 132 insertions(+), 48 deletions(-)
diff --git a/Blackboard/Architecture/pmoqobuz_ameliorations_qbz.md b/Blackboard/Architecture/pmoqobuz_ameliorations_qbz.md
index 47302dc2..69b9902f 100644
--- a/Blackboard/Architecture/pmoqobuz_ameliorations_qbz.md
+++ b/Blackboard/Architecture/pmoqobuz_ameliorations_qbz.md
@@ -70,20 +70,17 @@ POST /track/getList
---
-## 4. Pagination concurrente des playlists — **À FAIRE** (priorité moyenne)
+## 4. Pagination concurrente des playlists — **FAIT**
-**Problème** : `pmoqobuz` charge les pages de tracks d'une playlist séquentiellement (offset=0, puis
-offset=500, etc.). Chaque requête attend la précédente.
+**Implémentation réalisée** dans `QobuzApi::get_playlist_tracks` :
+- Page size augmentée de 50 → **500** (réduit le nombre de pages de 10×)
+- Page 1 séquentielle pour obtenir `total`
+- Pages 2..N lancées en parallèle via `futures::try_join_all` + `Semaphore(3)`
+- Résultats triés par offset avant fusion — ordre playlist garanti
+- Suivi de phase 2 (`track/getList`) inchangé
-**Ce que fait qbz** (`get_playlist`, l.1397) :
-- Page 1 → récupère les métadonnées + `total` track count
-- Pages 2..N → lancées **concurremment** via `join_all` dès que `total` est connu
-- Résultats ré-ordonnés par offset avant fusion
-
-**Impact pour pmoqobuz** : une playlist de 2 000 tracks (4 pages de 500) passe de 4 requêtes
-séquentielles (~1,6 s) à 1 + 3 en parallèle (~0,7 s).
-
-**Note** : à implémenter avec un semaphore (comme le CMAF) pour ne pas surcharger l'API Qobuz.
+**Impact** : playlist de 2 000 tracks (4 pages de 500) → 1 séquentielle + 3 parallèles ≈ 0,7 s
+au lieu de 4 séquentielles ≈ 1,6 s. Playlists ≤ 500 tracks : 1 seule requête.
---
@@ -122,6 +119,6 @@ sortis récemment. Utile pour le catalogue de la webapp.
| 1 | Streaming CMAF | Élevé | Critique (pipeline futur) | **Fait** |
| 2 | Bundle extraction avec cache disque | Moyen | Élevé (résilience) | **Fait** |
| 3 | Batch `track/getList` | Faible | Élevé (performances) | **Fait** |
-| 4 | Pagination concurrente playlists | Faible | Moyen | À faire |
+| 4 | Pagination concurrente playlists | Faible | Moyen | **Fait** |
| 5 | Release watch endpoint | Faible | Faible (catalogue) | À faire |
| 6 | `extra=track_ids` + batch à deux passes | Faible | Faible (optimisation) | À faire |
diff --git a/Cargo.lock b/Cargo.lock
index 5f159b5d..b6106781 100644
--- a/Cargo.lock
+++ b/Cargo.lock
@@ -4072,6 +4072,7 @@ dependencies = [
"cbc",
"chrono",
"ctr",
+ "futures",
"hex",
"hkdf",
"indexmap 2.12.0",
diff --git a/pmoqobuz/Cargo.toml b/pmoqobuz/Cargo.toml
index 71011b26..fba6b29e 100644
--- a/pmoqobuz/Cargo.toml
+++ b/pmoqobuz/Cargo.toml
@@ -14,6 +14,7 @@ reqwest = { version = "0.12", features = ["json", "cookies"] }
# Gestion asynchrone
tokio = { workspace = true }
+futures = { workspace = true }
# Sérialisation/Désérialisation JSON
serde = { workspace = true }
diff --git a/pmoqobuz/src/api/catalog.rs b/pmoqobuz/src/api/catalog.rs
index 683e67bf..7aa14bec 100644
--- a/pmoqobuz/src/api/catalog.rs
+++ b/pmoqobuz/src/api/catalog.rs
@@ -413,46 +413,100 @@ impl QobuzApi {
/// Récupère les tracks d'une playlist.
///
- /// Phase 1 : pagination de `/playlist/get?extra=tracks` pour collecter les IDs
- /// et les données de base.
- /// Phase 2 (si secret disponible) : enrichissement via `track/getList` pour
- /// obtenir les métadonnées complètes (performer, sample_rate, bit_depth, channels).
+ /// Phase 1 — pagination concurrente :
+ /// - Page 1 séquentielle pour obtenir `total`
+ /// - Pages 2..N lancées en parallèle (semaphore 3) dès que `total` est connu
+ /// - Résultats triés par offset avant fusion
+ ///
+ /// Phase 2 — enrichissement via `track/getList` pour métadonnées complètes
+ /// (performer, sample_rate, bit_depth, channels).
pub async fn get_playlist_tracks(&self, playlist_id: &str) -> Result> {
+ use futures::future::try_join_all;
+ use std::sync::Arc;
+ use tokio::sync::Semaphore;
+
+ const PAGE_SIZE: u32 = 500;
+ const LIMIT_STR: &str = "500";
+ // Configurable via accounts.qobuz.page_concurrency (défaut 3).
+ let max_concurrent_pages = self.page_concurrency;
+
debug!("Fetching tracks for playlist {}", playlist_id);
- const PAGE_SIZE: u32 = 50;
- let mut ordered_ids: Vec = Vec::new();
- let mut fallback_tracks: Vec
-
-
-
Aucun résultat
-
- -
-
-
![]()
-
-
-
-
-
-
-
{{ item.title }}
-
{{ item.artist }}
-
-
-
-
-
-
-
- ([])
-const searchResults = ref(null)
-const searchQuery = ref('')
const CACHE_DURATION_MS = 2000
const BROWSE_WINDOW_SIZE = 200
@@ -215,14 +213,14 @@ export function useMediaServers() {
}
// Recherche dans un serveur — retourne l'ID du container virtuel de résultats
- async function searchServer(serverId: string, query: string, context?: string): Promise {
+ async function searchServer(serverId: string, query: string): Promise {
if (!query.trim()) return null
try {
loading.value = true
error.value = null
- const data = await api.searchServer(serverId, query, context)
+ const data = await api.searchServer(serverId, query)
// data.container_id est l'ID virtuel réel (ex: "qobuz:search:catalog:all:camille")
const key = browseCacheKey(serverId, data.container_id)
browseCache.value.set(key, {
@@ -240,11 +238,6 @@ export function useMediaServers() {
}
}
- function clearSearch() {
- searchResults.value = null
- searchQuery.value = ''
- }
-
// Getters
function getServerById(id: string) {
return serversCache.value.get(id)
@@ -299,13 +292,9 @@ export function useMediaServers() {
getServerById,
getBrowseCached,
hasMore,
- // Search
- searchResults,
- searchQuery,
- searchServer,
- clearSearch,
// Actions
fetchServers,
+ searchServer,
browseContainer,
loadMoreBrowse,
setPath,
diff --git a/pmocontrol/src/pmoserver_ext.rs b/pmocontrol/src/pmoserver_ext.rs
index 1bcfdce4..ecaac228 100644
--- a/pmocontrol/src/pmoserver_ext.rs
+++ b/pmocontrol/src/pmoserver_ext.rs
@@ -8,6 +8,8 @@ use crate::control_point::ControlPoint;
#[cfg(feature = "pmoserver")]
use crate::media_server::{MediaBrowser, playback_item_from_entry};
#[cfg(feature = "pmoserver")]
+use crate::MediaEntry;
+#[cfg(feature = "pmoserver")]
use crate::model::{RendererCapabilities, RendererProtocol};
#[cfg(feature = "pmoserver")]
use crate::openapi::{
@@ -2392,6 +2394,22 @@ struct SearchQuery {
q: String,
}
+#[cfg(feature = "pmoserver")]
+fn search_result_container_id(entries: &[ContainerEntry]) -> String {
+ entries
+ .iter()
+ .find(|entry| entry.is_container)
+ .and_then(|entry| {
+ let parts: Vec<&str> = entry.id.splitn(5, ':').collect();
+ if parts.len() == 5 && parts[0] == "qobuz" && parts[1] == "search" {
+ Some(format!("qobuz:search:{}:all:{}", parts[2], parts[4]))
+ } else {
+ None
+ }
+ })
+ .unwrap_or_else(|| "search".to_string())
+}
+
/// GET /control/servers/{server_id}/search?q= - Recherche dans un serveur
#[cfg(feature = "pmoserver")]
#[utoipa::path(
@@ -2503,8 +2521,10 @@ async fn search_server(
})
.collect();
+ let container_id = search_result_container_id(&container_entries);
+
Ok(Json(BrowseResponse {
- container_id: "search".to_string(),
+ container_id,
entries: container_entries,
total_count,
offset: 0,
diff --git a/pmoqobuz/src/source.rs b/pmoqobuz/src/source.rs
index f234119c..a37b52f4 100644
--- a/pmoqobuz/src/source.rs
+++ b/pmoqobuz/src/source.rs
@@ -1898,7 +1898,8 @@ impl MusicSource for QobuzSource {
}
ObjectIdType::SearchResult(scope, media_type, query) => {
- tracing::debug!(scope = ?scope, media_type = ?media_type, query, "browse → execute_search");
+ tracing::debug!(scope = ?scope, media_type = ?media_type, query, "browse → search");
+ let is_all_search = media_type == MediaSearchType::All;
let sq = SearchQuery {
text: query,
media_type,
@@ -1906,9 +1907,11 @@ impl MusicSource for QobuzSource {
limit: 200,
offset: 0,
};
- let result = self.execute_search(&sq).await;
- tracing::debug!(ok = result.is_ok(), "execute_search returned");
- result
+ if is_all_search {
+ self.search_grouped(&sq).await
+ } else {
+ self.execute_search(&sq).await
+ }
}
ObjectIdType::Track(_) => {
@@ -2410,10 +2413,11 @@ impl QobuzSource {
};
let (n_albums, n_artists, n_tracks, n_playlists) = counts;
+ let parent_id = format!("qobuz:search:{}:all:{}", scope_str, query.text);
let mk_container = |type_str: &str, title: &str, count: usize| Container {
id: format!("qobuz:search:{}:{}:{}", scope_str, type_str, query.text),
- parent_id: "qobuz:search".to_string(),
+ parent_id: parent_id.clone(),
restricted: Some("1".to_string()),
child_count: Some(count.to_string()),
searchable: None,
@@ -2466,7 +2470,6 @@ impl QobuzSource {
use tracing::debug;
debug!(text = %query.text, scope = ?query.scope, media_type = ?query.media_type, "execute_search entry");
let text = &query.text;
- let limit = query.limit;
let scope_str = match query.scope {
SearchScope::Catalog => "catalog",
SearchScope::UserLibrary => "favorites",
diff --git a/version.txt b/version.txt
index 6b9aa4e6..d57e08b5 100644
--- a/version.txt
+++ b/version.txt
@@ -1 +1 @@
-0.3.51
+0.3.52
--
2.49.1