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.
This commit is contained in:
@@ -245,6 +245,15 @@ pub trait QobuzConfigExt {
|
|||||||
|
|
||||||
/// Persiste la version du bundle après une extraction réussie.
|
/// Persiste la version du bundle après une extraction réussie.
|
||||||
fn set_qobuz_bundle_version(&self, version: &str) -> Result<()>;
|
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 {
|
impl QobuzConfigExt for Config {
|
||||||
@@ -504,4 +513,13 @@ impl QobuzConfigExt for Config {
|
|||||||
Value::String(version.to_string()),
|
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,
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -115,6 +115,9 @@ struct QobuzSourceInner {
|
|||||||
/// Base URL for streaming server (e.g., "http://192.168.0.138:8080")
|
/// Base URL for streaming server (e.g., "http://192.168.0.138:8080")
|
||||||
base_url: String,
|
base_url: String,
|
||||||
|
|
||||||
|
/// Nombre de workers concurrents pour register_tracks_lazy (configurable)
|
||||||
|
register_concurrency: usize,
|
||||||
|
|
||||||
/// Update tracking
|
/// Update tracking
|
||||||
update_counter: tokio::sync::RwLock<u32>,
|
update_counter: tokio::sync::RwLock<u32>,
|
||||||
last_change: tokio::sync::RwLock<SystemTime>,
|
last_change: tokio::sync::RwLock<SystemTime>,
|
||||||
@@ -142,15 +145,18 @@ impl QobuzSource {
|
|||||||
/// Returns an error if the caches are not initialized in the registry
|
/// Returns an error if the caches are not initialized in the registry
|
||||||
#[cfg(feature = "server")]
|
#[cfg(feature = "server")]
|
||||||
pub fn from_registry(client: QobuzClient, base_url: impl Into<String>) -> Result<Self> {
|
pub fn from_registry(client: QobuzClient, base_url: impl Into<String>) -> Result<Self> {
|
||||||
|
use crate::config_ext::QobuzConfigExt;
|
||||||
let cache_manager = SourceCacheManager::from_registry("qobuz".to_string())?;
|
let cache_manager = SourceCacheManager::from_registry("qobuz".to_string())?;
|
||||||
let client = Arc::new(client);
|
let client = Arc::new(client);
|
||||||
cache_manager.register_lazy_provider(Arc::new(QobuzLazyProvider::new(client.clone())));
|
cache_manager.register_lazy_provider(Arc::new(QobuzLazyProvider::new(client.clone())));
|
||||||
|
let register_concurrency = pmoconfig::get_config().get_qobuz_register_concurrency();
|
||||||
|
|
||||||
Ok(Self {
|
Ok(Self {
|
||||||
inner: Arc::new(QobuzSourceInner {
|
inner: Arc::new(QobuzSourceInner {
|
||||||
client,
|
client,
|
||||||
cache_manager,
|
cache_manager,
|
||||||
base_url: base_url.into(),
|
base_url: base_url.into(),
|
||||||
|
register_concurrency,
|
||||||
update_counter: tokio::sync::RwLock::new(0),
|
update_counter: tokio::sync::RwLock::new(0),
|
||||||
last_change: tokio::sync::RwLock::new(SystemTime::now()),
|
last_change: tokio::sync::RwLock::new(SystemTime::now()),
|
||||||
}),
|
}),
|
||||||
@@ -180,6 +186,7 @@ impl QobuzSource {
|
|||||||
client,
|
client,
|
||||||
cache_manager,
|
cache_manager,
|
||||||
base_url: base_url.into(),
|
base_url: base_url.into(),
|
||||||
|
register_concurrency: 4,
|
||||||
update_counter: tokio::sync::RwLock::new(0),
|
update_counter: tokio::sync::RwLock::new(0),
|
||||||
last_change: tokio::sync::RwLock::new(SystemTime::now()),
|
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.
|
/// 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.
|
/// Retourne la liste des lazy PKs enregistrés avec succès.
|
||||||
async fn register_tracks_lazy(&self, tracks: &[crate::models::Track]) -> Vec<String> {
|
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.
|
// Configurable via accounts.qobuz.register_concurrency (défaut 4).
|
||||||
// Les covers déjà cachées sont retournées immédiatement (pas d'HTTP), donc même
|
// SQLite sérialise les écritures — au-delà de ~4 workers on accumule
|
||||||
// 600 tracks ne génèrent que ~N_albums_uniques téléchargements réels.
|
// des threads en attente du mutex DB sans gain de débit.
|
||||||
let sem = std::sync::Arc::new(tokio::sync::Semaphore::new(16));
|
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
|
// 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).
|
// d'origine après complétion parallèle (JoinSet retourne dans l'ordre de fin).
|
||||||
|
|||||||
Reference in New Issue
Block a user