diff --git a/.DS_Store b/.DS_Store index 58e2f0f5..06359845 100644 Binary files a/.DS_Store and b/.DS_Store differ diff --git a/Blackboard/Architecture/pmoqobuz_ameliorations_qbz.md b/Blackboard/Architecture/pmoqobuz_ameliorations_qbz.md new file mode 100644 index 00000000..47302dc2 --- /dev/null +++ b/Blackboard/Architecture/pmoqobuz_ameliorations_qbz.md @@ -0,0 +1,127 @@ +# Améliorations pmoqobuz inspirées de qbz + +Analyse comparative avec le projet [qbz](../../../qbz) (`crates/qbz-qobuz/`), client Qobuz Rust plus avancé. +Les items sont classés par priorité et état d'avancement. + +--- + +## 1. Streaming CMAF — **FAIT** + +**Problème** : l'endpoint legacy `/track/getFileUrl` est en cours de dépréciation. Qobuz bascule vers +CMAF (Common Media Application Format) : segments AES-CTR chiffrés sur CDN Akamai. + +**Implémentation réalisée** : +- `pmoqobuz::cmaf` — pipeline complet : dérivation HKDF, dérobage AES-CBC, déchiffrement AES-CTR par frame +- `pmoqobuz::retry` — retry exponentiel avec classification Transient/Terminal +- `QobuzClient::open_cmaf_stream()` — `AsyncRead` progressif via `tokio::io::duplex`, 3 segments en vol +- `GET /qobuz/tracks/:id/flac` — endpoint REST local qui expose le flux ; `LazyProvider::get_url()` + retourne cette URL locale (via `PMO_SERVER_URL`) → le progressive caching de pmocache est préservé sans + modification + +**Référence qbz** : `crates/qbz-qobuz/src/cmaf.rs` + +--- + +## 2. Extraction du bundle Qobuz — **Fait** + +**Problème** : `pmoqobuz` utilise un `app_id` et un `configvalue` statiques, hardcodés ou configurés +manuellement. Qobuz peut les invalider à tout moment en changeant son bundle JS. + +**Ce que fait qbz** (`bundle.rs`) : +- Télécharge la page `https://play.qobuz.com/login`, extrait l'URL du bundle JS +- Parse le bundle (~7 MB) avec des regex pour en extraire `app_id`, les secrets, et la `private_key` OAuth +- Met en cache les tokens sur disque avec un hash de version du bundle (`bundle_version`) +- Revalide automatiquement si la version change (rotation silencieuse de Qobuz) +- Timeout de 45 s sur le fetch, 2 retry sur extraction + +**Impact pour pmoqobuz** : +- Le `Spoofer` actuel fait un fetch similaire mais sans cache disque ni détection de version +- Ajouter `CachedBundle` (version + tokens + timestamp) dans `pmoconfig` ou dans le répertoire de données +- Relire le cache au démarrage, re-extraire seulement si la version du bundle a changé + +**Avantage** : ne jamais tomber en panne quand Qobuz rotate ses secrets sans préavis. + +--- + +## 3. Chargement batch de tracks — `track/getList` — **FAIT** + +**Problème** : les tracks retournées par `/playlist/get?extra=tracks` et +`/favorite/getUserFavorites` ont des métadonnées incomplètes (parfois sans `performer`, +jamais de `sample_rate`/`bit_depth`/`channels`). Cela déclenchait des appels individuels +`get_track` lazy à la lecture. + +**Implémentation réalisée** : +- `signing::sign_track_get_list(ids_csv, timestamp, secret)` — signature MD5 pour `track/getList` +- `QobuzApi::post_json_with_query` — POST avec auth headers + query sig + JSON body +- `QobuzApi::get_tracks_batch(&[&str])` — fenêtres de 50 IDs, appels en série +- `QobuzClient::get_tracks_batch` — wrapper avec cache (skip les IDs déjà en cache) +- `get_playlist_tracks` : phase 1 pagination existante, phase 2 enrichissement si secret disponible +- `get_favorite_tracks` : même enrichissement en phase 2 +- Fallback gracieux si le secret est absent ou si `track/getList` échoue + +**Ce que fait qbz** (`get_tracks_batch`, l.1323) : +``` +POST /track/getList +{ "tracks_id": [id1, id2, ..., id50] } +→ { "tracks": { "total": N, "items": [...Track] } } +``` +- Fenêtre de 50 IDs max par appel (limite API Qobuz) +- Les fenêtres supérieures à 50 sont découpées et appelées en série (respecte les quotas) + +--- + +## 4. Pagination concurrente des playlists — **À FAIRE** (priorité moyenne) + +**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. + +**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. + +--- + +## 5. Endpoint `release_watch` — **À FAIRE** (priorité basse) + +**Problème** : pmoqobuz ne supporte pas les nouvelles sorties d'artistes suivis. + +**Ce que fait qbz** (`get_release_watch`, l.809) : +``` +GET /favorite/getNewReleases?type=album&limit=50&offset=0 +→ { has_more: bool, items: [...Album] } +``` +- Types disponibles : `album`, `live`, `ep_single` +- Pas de champ `total` dans la réponse — pagination via `has_more` + +**Impact pour pmoqobuz** : permettre un "quoi de neuf" dans l'interface — albums des artistes favoris +sortis récemment. Utile pour le catalogue de la webapp. + +--- + +## 6. Chargement `playlist/get?extra=track_ids` — **À FAIRE** (priorité basse) + +**Ce que fait qbz** (`get_playlist_track_ids`, l.1296) : +- Variante légère de `playlist/get` qui retourne uniquement les IDs (pas les objets Track complets) +- Utile pour vérifier si une playlist a changé sans tout recharger +- Combiné avec `get_tracks_batch` pour un chargement optimal en deux passes : + 1. `playlist/get?extra=track_ids` → liste d'IDs + 2. `track/getList` par fenêtres de 50 → objets Track complets + +--- + +## Résumé de priorités + +| # | Amélioration | Effort | Impact | État | +|---|---|---|---|---| +| 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 | +| 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 450d1a02..5f159b5d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4,7 +4,7 @@ version = 4 [[package]] name = "PMOMusic" -version = "0.3.49" +version = "0.3.51" dependencies = [ "axum 0.8.7", "console-subscriber", @@ -715,6 +715,15 @@ dependencies = [ "generic-array", ] +[[package]] +name = "block-padding" +version = "0.3.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a8894febbff9f758034a5b8e12d87918f56dfc64a8e1fe757d65e29041538d93" +dependencies = [ + "generic-array", +] + [[package]] name = "block2" version = "0.6.2" @@ -808,6 +817,15 @@ dependencies = [ "rustversion", ] +[[package]] +name = "cbc" +version = "0.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "26b52a9543ae338f279b96b0b9fed9c8093744685043739079ce85cd58f289a6" +dependencies = [ + "cipher", +] + [[package]] name = "cc" version = "1.2.46" @@ -1317,6 +1335,7 @@ checksum = "9ed9a281f7bc9b7576e61468ba615a66a5c8cfdff42420a70aa82701a3b1e292" dependencies = [ "block-buffer", "crypto-common", + "subtle", ] [[package]] @@ -2025,6 +2044,24 @@ version = "0.4.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7f24254aa9a54b5c858eaee2f5bccdb46aaf0e486a595ed5fd8f86ba55232a70" +[[package]] +name = "hkdf" +version = "0.12.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7b5f8eb2ad728638ea2c7d47a21db23b7b58a72ed6a38256b8a1849f15fbbdf7" +dependencies = [ + "hmac", +] + +[[package]] +name = "hmac" +version = "0.12.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6c49c37c09c17a53d937dfbb742eb3a961d65a994e6bcdcf37e7399d0cc8ab5e" +dependencies = [ + "digest", +] + [[package]] name = "home" version = "0.5.12" @@ -2408,6 +2445,7 @@ version = "0.1.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "879f10e63c20629ecabbb64a8010319738c66a5cd0c29b02d63d272b03751d01" dependencies = [ + "block-padding", "generic-array", ] @@ -4026,12 +4064,16 @@ dependencies = [ name = "pmoqobuz" version = "0.1.0" dependencies = [ + "aes", "anyhow", "async-trait", "axum 0.8.7", "base64 0.22.1", + "cbc", "chrono", + "ctr", "hex", + "hkdf", "indexmap 2.12.0", "md-5", "mockito", @@ -4051,10 +4093,12 @@ dependencies = [ "serde_json", "serde_yaml", "sha1", + "sha2", "tempfile", "thiserror 2.0.17", "tokio", "tokio-test", + "tokio-util", "tracing", "tracing-subscriber", "utoipa", diff --git a/PMOMusic/Cargo.toml b/PMOMusic/Cargo.toml index b37f9d90..2f8ba0ef 100644 --- a/PMOMusic/Cargo.toml +++ b/PMOMusic/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "PMOMusic" -version = "0.3.50" +version = "0.3.51" edition = "2024" [dependencies] diff --git a/pmoqobuz/Cargo.toml b/pmoqobuz/Cargo.toml index cfd82d24..71011b26 100644 --- a/pmoqobuz/Cargo.toml +++ b/pmoqobuz/Cargo.toml @@ -29,9 +29,17 @@ sha1 = "0.10" hex = "0.4" md-5 = "0.10" +# Crypto CMAF (AES-CTR déchiffrement frames, AES-CBC dérobage clé, HKDF dérivation) +aes = "0.8" +cbc = "0.1" +ctr = "0.9" +hkdf = "0.12" +sha2 = "0.10" + # Cache en mémoire avec TTL moka = { version = "0.12", features = ["future"] } + # Logging tracing = { workspace = true } @@ -53,6 +61,7 @@ pmodidl = { path = "../pmodidl" } # Intégration avec pmoserver pour l'API HTTP pmoserver = { path = "../pmoserver", optional = true } axum = { version = "0.8", optional = true } +tokio-util = { workspace = true, optional = true } rusqlite = { version = "0.37", features = ["bundled"], optional = true } # Documentation OpenAPI @@ -68,7 +77,7 @@ pmocache = { path = "../pmocache" } [features] default = ["async-trait-support"] # Feature pour activer les extensions pmoserver -pmoserver = ["dep:pmoserver", "dep:axum", "dep:utoipa"] +pmoserver = ["dep:pmoserver", "dep:axum", "dep:utoipa", "dep:tokio-util"] # Feature pour activer le support serveur (cache registry) server = ["pmosource/server"] # Feature cache (deprecated - toujours actif maintenant) diff --git a/pmoqobuz/src/api/catalog.rs b/pmoqobuz/src/api/catalog.rs index b0a3b10d..683e67bf 100644 --- a/pmoqobuz/src/api/catalog.rs +++ b/pmoqobuz/src/api/catalog.rs @@ -4,6 +4,7 @@ use super::QobuzApi; use crate::error::{QobuzError, Result}; use crate::models::*; use serde::Deserialize; +use std::collections::HashMap; use tracing::debug; /// Réponse paginée de l'API @@ -182,6 +183,17 @@ struct FileUrlResponse { format_id: u8, } +/// Réponse de l'endpoint track/getList +#[derive(Debug, Deserialize)] +struct TrackListResponse { + tracks: TrackListItems, +} + +#[derive(Debug, Deserialize)] +struct TrackListItems { + items: Vec, +} + fn default_streamable() -> bool { true } @@ -213,6 +225,66 @@ impl QobuzApi { } } + /// Récupère les détails de plusieurs tracks en une ou plusieurs requêtes batch. + /// + /// Utilise l'endpoint `track/getList` (max 50 IDs par appel). Les fenêtres + /// supérieures à 50 sont découpées et appellées en série. L'ordre de sortie + /// correspond à l'ordre des `track_ids` en entrée. + /// + /// Requiert que le secret s4 soit configuré. + pub async fn get_tracks_batch(&self, track_ids: &[&str]) -> Result> { + const MAX_PER_CALL: usize = 50; + if track_ids.is_empty() { + return Ok(Vec::new()); + } + if track_ids.len() <= MAX_PER_CALL { + return self.get_tracks_batch_chunk(track_ids).await; + } + debug!("get_tracks_batch: {} IDs en fenêtres de {}", track_ids.len(), MAX_PER_CALL); + let mut all = Vec::with_capacity(track_ids.len()); + for chunk in track_ids.chunks(MAX_PER_CALL) { + let mut tracks = self.get_tracks_batch_chunk(chunk).await?; + all.append(&mut tracks); + } + Ok(all) + } + + async fn get_tracks_batch_chunk(&self, track_ids: &[&str]) -> Result> { + use super::signing; + + let secret = self.secret().ok_or_else(|| { + QobuzError::Configuration("Secret s4 requis pour track/getList".to_string()) + })?; + + let ids_csv = track_ids.join(","); + let timestamp = signing::get_timestamp(); + let signature = signing::sign_track_get_list(&ids_csv, ×tamp, &secret); + + let query_params = [ + ("request_ts", timestamp.as_str()), + ("request_sig", signature.as_str()), + ]; + + // Les IDs sont envoyés comme tableau d'entiers dans le body JSON + let ids_as_numbers: Vec = track_ids + .iter() + .filter_map(|id| id.parse().ok()) + .collect(); + let body = serde_json::json!({ "tracks_id": ids_as_numbers }); + + debug!("get_tracks_batch_chunk POST {} IDs", track_ids.len()); + let response: TrackListResponse = self + .post_json_with_query("/track/getList", &query_params, body) + .await?; + + Ok(response + .tracks + .items + .into_iter() + .map(|t| Self::parse_track(t, None)) + .collect()) + } + /// Récupère les détails d'une track pub async fn get_track(&self, track_id: &str) -> Result { debug!("Fetching track {}", track_id); @@ -339,13 +411,20 @@ impl QobuzApi { Ok(Self::parse_playlist(response)) } - /// Récupère les tracks d'une playlist + /// 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). pub async fn get_playlist_tracks(&self, playlist_id: &str) -> Result> { debug!("Fetching tracks for playlist {}", playlist_id); const PAGE_SIZE: u32 = 50; - let mut all_tracks = Vec::new(); + let mut ordered_ids: Vec = Vec::new(); + let mut fallback_tracks: Vec = Vec::new(); let mut offset = 0u32; + // Phase 1 : pagination pour collecter les IDs et les tracks de base loop { let offset_str = offset.to_string(); let limit_str = PAGE_SIZE.to_string(); @@ -360,7 +439,10 @@ impl QobuzApi { 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))); + for t in tracks.items { + ordered_ids.push(t.id.clone()); + fallback_tracks.push(Self::parse_track(t, None)); + } offset += count; if count == 0 || offset >= total { break; @@ -370,8 +452,23 @@ impl QobuzApi { } } - debug!("Fetched {} tracks total for playlist {}", all_tracks.len(), playlist_id); - Ok(all_tracks) + debug!("Fetched {} track IDs for playlist {}", ordered_ids.len(), playlist_id); + + if ordered_ids.is_empty() { + return Ok(Vec::new()); + } + + // Phase 2 : enrichissement via track/getList pour métadonnées complètes + let id_refs: Vec<&str> = ordered_ids.iter().map(|s| s.as_str()).collect(); + let full_tracks = self.get_tracks_batch(&id_refs).await?; + let mut track_map: HashMap = + full_tracks.into_iter().map(|t| (t.id.clone(), t)).collect(); + let enriched: Vec = ordered_ids + .iter() + .filter_map(|id| track_map.remove(id.as_str())) + .collect(); + debug!("Fetched {} tracks for playlist {} via track/getList", enriched.len(), playlist_id); + Ok(enriched) } /// Récupère la liste des genres diff --git a/pmoqobuz/src/api/cmaf.rs b/pmoqobuz/src/api/cmaf.rs new file mode 100644 index 00000000..aaa9cf65 --- /dev/null +++ b/pmoqobuz/src/api/cmaf.rs @@ -0,0 +1,261 @@ +//! Endpoints API CMAF de Qobuz : session/start et file/url. +//! +//! Ces endpoints implémentent le nouveau pipeline de streaming Qobuz +//! qui remplace progressivement `/track/getFileUrl`. + +use serde::Deserialize; +use std::time::{SystemTime, UNIX_EPOCH}; +use tokio::sync::RwLock; +use tracing::info; + +use crate::cmaf::crypto::{compute_request_sig, CMAF_SEED}; +use crate::error::{QobuzError, Result}; + +/// État d'une session CMAF active. +pub struct CmafSession { + pub session_id: String, + pub infos: String, + pub expires_at: u64, +} + +/// Réponse de l'endpoint /session/start. +#[derive(Debug, Deserialize)] +struct SessionStartResponse { + session_id: String, + #[serde(default)] + infos: Option, + expires_at: u64, +} + +/// Réponse de l'endpoint /file/url (CMAF). +#[derive(Debug, Deserialize)] +pub struct CmafFileUrlResponse { + /// Modèle d'URL avec le placeholder `$SEGMENT$`. + #[serde(default)] + pub url_template: Option, + /// Clé de contenu enveloppée, format `"qbz-1.wrapped_b64url.iv_b64url"`. + #[serde(default)] + pub key: Option, + /// Nombre de segments audio (hors segment init). + pub n_segments: u8, + #[serde(default)] + pub format_id: Option, + #[serde(default)] + pub mime_type: Option, + #[serde(default)] + pub sampling_rate: Option, + /// Profondeur de bits (champ v1 de l'API). + #[serde(default)] + pub bits_depth: Option, + /// Profondeur de bits (champ v2 de l'API). + #[serde(default)] + pub bit_depth: Option, +} + +impl CmafFileUrlResponse { + /// Retourne la profondeur de bits en préférant `bits_depth` puis `bit_depth`. + pub fn resolved_bit_depth(&self) -> Option { + self.bits_depth.or(self.bit_depth) + } +} + +const BASE_URL: &str = "https://www.qobuz.com/api.json/0.2"; + +fn current_timestamp() -> u64 { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap_or_default() + .as_secs() +} + +/// Signe une requête /session/start. +fn sign_session_start(timestamp: u64) -> String { + let mut args = std::collections::BTreeMap::new(); + args.insert("profile", "qbz-1".to_string()); + compute_request_sig("sessionstart", &args, ×tamp.to_string(), CMAF_SEED) +} + +/// Signe une requête /file/url. +fn sign_file_url(track_id: &str, format_id: u32, timestamp: u64) -> String { + let mut args = std::collections::BTreeMap::new(); + args.insert("format_id", format_id.to_string()); + args.insert("intent", "stream".to_string()); + args.insert("track_id", track_id.to_string()); + compute_request_sig("fileurl", &args, ×tamp.to_string(), CMAF_SEED) +} + +/// Gestionnaire de session CMAF avec renouvellement automatique. +/// +/// Utilise le pattern double-checked lock : fast path sous read guard, +/// slow path sous write guard exclusif pour éviter les renouvellements +/// concurrents qui produiraient des `infos` incohérents avec les `key` +/// retournées par `/file/url`. +pub struct CmafSessionManager { + session: RwLock>, +} + +impl CmafSessionManager { + pub fn new() -> Self { + Self { session: RwLock::new(None) } + } + + /// Retourne `(session_id, infos)` d'une session valide, en en démarrant + /// une nouvelle si la session courante est absente ou expire dans < 60s. + pub async fn ensure_session( + &self, + http: &reqwest::Client, + app_id: &str, + auth_token: &str, + ) -> Result<(String, String)> { + let now = current_timestamp(); + + // Fast path : session existante avec > 60s restants. + { + let guard = self.session.read().await; + if let Some(ref cs) = *guard { + if cs.expires_at > now + 60 { + return Ok((cs.session_id.clone(), cs.infos.clone())); + } + } + } + + // Slow path : prend le verrou d'écriture et vérifie à nouveau. + let mut guard = self.session.write().await; + if let Some(ref cs) = *guard { + if cs.expires_at > now + 60 { + return Ok((cs.session_id.clone(), cs.infos.clone())); + } + } + + info!("[CMAF] Démarrage d'une nouvelle session"); + let timestamp = current_timestamp(); + let sig = sign_session_start(timestamp); + + let url = format!("{}/session/start", BASE_URL); + let response = http + .post(&url) + .header("X-App-Id", app_id) + .header("X-User-Auth-Token", auth_token) + .form(&[ + ("profile", "qbz-1"), + ("request_ts", ×tamp.to_string()), + ("request_sig", &sig), + ]) + .send() + .await + .map_err(QobuzError::Http)?; + + let status = response.status(); + if !status.is_success() { + let body = response.text().await.unwrap_or_default(); + return Err(QobuzError::ApiError { + code: status.as_u16(), + message: format!("session/start échoué: {}", body), + }); + } + + let resp: SessionStartResponse = response.json().await.map_err(|e| { + QobuzError::Other(format!("parse session/start: {}", e)) + })?; + + let infos = resp.infos.unwrap_or_default(); + info!( + "[CMAF] Session démarrée: id={}..., expires_at={}", + &resp.session_id[..resp.session_id.len().min(8)], + resp.expires_at, + ); + + let session_id = resp.session_id.clone(); + let infos_clone = infos.clone(); + + *guard = Some(CmafSession { + session_id: resp.session_id, + infos, + expires_at: resp.expires_at, + }); + + Ok((session_id, infos_clone)) + } + + /// Invalide la session courante (utile en cas d'erreur de déchiffrement). + pub async fn invalidate(&self) { + *self.session.write().await = None; + } +} + +/// Récupère l'URL de fichier CMAF pour un track. +/// +/// Inclut le retry sur les erreurs transitoires (5xx, 429, erreurs réseau). +/// Un 404 retourne immédiatement une erreur `NotFound`. +pub async fn get_file_url( + http: &reqwest::Client, + app_id: &str, + auth_token: &str, + session_id: &str, + track_id: &str, + format_id: u32, +) -> Result { + use crate::retry::{classify_reqwest, classify_status, retry_transient, DEFAULT_MAX_ATTEMPTS}; + + let url = format!("{}/file/url", BASE_URL); + + let result = retry_transient( + DEFAULT_MAX_ATTEMPTS, + "CMAF file/url", + |e: &QobuzError| matches!(e, QobuzError::RateLimitExceeded | QobuzError::ApiError { code: 500..=599, .. }), + |_attempt| { + let url = url.clone(); + async move { + let timestamp = current_timestamp(); + let sig = sign_file_url(track_id, format_id, timestamp); + + let response = http + .get(&url) + .header("X-App-Id", app_id) + .header("X-User-Auth-Token", auth_token) + .header("X-Session-Id", session_id) + .query(&[ + ("track_id", track_id), + ("format_id", &format_id.to_string()), + ("intent", "stream"), + ("request_ts", ×tamp.to_string()), + ("request_sig", &sig), + ]) + .send() + .await + .map_err(QobuzError::Http)?; + + let status = response.status(); + tracing::info!("[CMAF] file/url track_id={} format_id={} status={}", track_id, format_id, status); + + if !status.is_success() { + let code = status.as_u16(); + return Err(match code { + 404 => QobuzError::NotFound(format!("track {} non disponible", track_id)), + 429 => QobuzError::RateLimitExceeded, + _ => QobuzError::ApiError { + code, + message: format!("file/url status {}", code), + }, + }); + } + + let file_url: CmafFileUrlResponse = response.json().await.map_err(|e| { + QobuzError::Other(format!("parse file/url: {}", e)) + })?; + + tracing::info!( + "[CMAF] file/url: n_segments={}, mime={:?}, sampling_rate={:?}", + file_url.n_segments, + file_url.mime_type, + file_url.sampling_rate, + ); + + Ok(file_url) + } + }, + ) + .await?; + + Ok(result) +} diff --git a/pmoqobuz/src/api/mod.rs b/pmoqobuz/src/api/mod.rs index 88acdd1b..e570cae4 100644 --- a/pmoqobuz/src/api/mod.rs +++ b/pmoqobuz/src/api/mod.rs @@ -4,10 +4,12 @@ pub mod auth; pub mod catalog; +pub mod cmaf; pub mod signing; pub mod spoofer; pub mod user; +use crate::api::cmaf::CmafSessionManager; use crate::error::{QobuzError, Result}; use crate::models::AudioFormat; use reqwest::{Client, Response}; @@ -95,6 +97,8 @@ pub struct QobuzApi { user_id: RwLock>, /// Format audio par défaut format_id: AudioFormat, + /// Gestionnaire de session CMAF (renouvellement automatique thread-safe) + pub(crate) cmaf_session: CmafSessionManager, } impl QobuzApi { @@ -114,6 +118,7 @@ impl QobuzApi { user_auth_token: RwLock::new(None), user_id: RwLock::new(None), format_id: AudioFormat::default(), + cmaf_session: CmafSessionManager::new(), }) } @@ -269,6 +274,37 @@ impl QobuzApi { self.request("POST", endpoint, params).await } + /// Effectue un POST avec signature en query params et données en JSON body. + /// + /// Utilisé par les endpoints qui attendent une structure JSON complexe + /// (ex: `track/getList` avec `{"tracks_id": [...]}`). + pub(crate) async fn post_json_with_query( + &self, + endpoint: &str, + query_params: &[(&str, &str)], + body: serde_json::Value, + ) -> Result { + let url = format!("{}{}", API_BASE_URL, endpoint); + debug!("POST JSON {} with {} query params", url, query_params.len()); + + let app_id = self.app_id.read().unwrap().clone(); + let mut builder = self + .client + .post(&url) + .header("X-App-Id", &app_id) + .header("Accept-Language", "en,en-US;q=0.8,ko;q=0.6,zh;q=0.4,zh-CN;q=0.2") + .header("Access-Control-Request-Headers", "x-user-auth-token,x-app-id") + .query(query_params) + .json(&body); + + if let Some(token) = self.auth_token() { + builder = builder.header("X-User-Auth-Token", token); + } + + let response = builder.send().await?; + self.handle_response(response, endpoint, query_params).await + } + /// Effectue une requête à l'API (générique) async fn request( &self, @@ -316,6 +352,57 @@ impl QobuzApi { self.handle_response(response, endpoint, params).await } + /// Retourne le client HTTP interne (utilisé par les endpoints CMAF). + pub(crate) fn http_client(&self) -> &Client { + &self.client + } + + /// Assure qu'une session CMAF valide existe et retourne `(session_id, infos)`. + /// + /// Crée ou renouvelle automatiquement la session si elle est absente ou expire + /// dans moins de 60 secondes. Les appels concurrents sont sérialisés pour + /// éviter des sessions incohérentes. + pub async fn ensure_cmaf_session(&self) -> Result<(String, String)> { + let app_id = self.app_id.read().unwrap().clone(); + let auth_token = self + .user_auth_token + .read() + .unwrap() + .clone() + .ok_or_else(|| QobuzError::Unauthorized("Token d'authentification manquant pour CMAF".into()))?; + + self.cmaf_session + .ensure_session(&self.client, &app_id, &auth_token) + .await + } + + /// Récupère l'URL de fichier CMAF pour un track avec le format donné. + pub async fn get_cmaf_file_url( + &self, + track_id: &str, + format_id: u32, + ) -> Result { + let app_id = self.app_id.read().unwrap().clone(); + let auth_token = self + .user_auth_token + .read() + .unwrap() + .clone() + .ok_or_else(|| QobuzError::Unauthorized("Token d'authentification manquant pour CMAF".into()))?; + + let (session_id, _infos) = self.ensure_cmaf_session().await?; + + crate::api::cmaf::get_file_url( + &self.client, + &app_id, + &auth_token, + &session_id, + track_id, + format_id, + ) + .await + } + /// Traite la réponse HTTP async fn handle_response( &self, diff --git a/pmoqobuz/src/api/signing.rs b/pmoqobuz/src/api/signing.rs index e4f5889c..f2909fb7 100644 --- a/pmoqobuz/src/api/signing.rs +++ b/pmoqobuz/src/api/signing.rs @@ -88,6 +88,20 @@ pub fn sign_track_get_file_url( /// # Returns /// /// Signature MD5 hexadécimale +/// Signe une requête track/getList +/// +/// Chaîne signée : `"trackgetList" + "tracks_id" + ids_csv + timestamp + secret` +/// où `ids_csv` est la liste des IDs séparés par des virgules. +pub fn sign_track_get_list(ids_csv: &str, timestamp: &str, secret: &[u8]) -> String { + let mut hasher = Md5::new(); + hasher.update(b"trackgetList"); + hasher.update(b"tracks_id"); + hasher.update(ids_csv.as_bytes()); + hasher.update(timestamp.as_bytes()); + hasher.update(secret); + format!("{:x}", hasher.finalize()) +} + pub fn sign_userlib_get_albums(timestamp: &str, secret: &[u8]) -> String { let mut hasher = Md5::new(); diff --git a/pmoqobuz/src/api/spoofer.rs b/pmoqobuz/src/api/spoofer.rs index 5ef3ffa1..81cc038a 100644 --- a/pmoqobuz/src/api/spoofer.rs +++ b/pmoqobuz/src/api/spoofer.rs @@ -3,36 +3,40 @@ use base64::{engine::general_purpose::STANDARD, Engine}; use indexmap::IndexMap; use regex::Regex; use reqwest::Client; +use std::time::Duration; +use tracing::{debug, info, warn}; + +/// Timeout par requête HTTP — le bundle fait ~7 MB et le CDN Qobuz peut être lent. +const BUNDLE_FETCH_TIMEOUT: Duration = Duration::from_secs(45); + +/// Tentatives supplémentaires après un échec d'extraction. +const BUNDLE_EXTRACTION_RETRIES: usize = 2; pub struct Spoofer { bundle: String, + /// Version du bundle extrait, ex. `"8.1.0-b019"`. + bundle_version: String, seed_timezone_regex: Regex, info_extras_regex_template: String, app_id_regex: Regex, } impl Spoofer { - /// Crée un nouveau Spoofer et télécharge le bundle.js - pub async fn new() -> Result { - // Expressions régulières (équivalent Python) - let seed_timezone_regex = Regex::new( - r#"[a-z]\.initialSeed\("(?P[\w=]+)",window\.utimezone\.(?P[a-z]+)\)"#, - )?; + /// Version du bundle Qobuz actuellement chargé. + pub fn bundle_version(&self) -> &str { + &self.bundle_version + } - let info_extras_regex_template = - r#"name:"\w+/(?P{timezones})",info:"(?P[\w=]+)",extras:"(?P[\w=]+)""# - .to_string(); - - let app_id_regex = Regex::new( - r#"production:\{api:\{appId:"(?P\d{9})",appSecret:"(?P\w{32})"\},braze:.\(.\(\{\},.\),\{\},\{apiKey:"([-0-9a-fA-F]{36})"\}\),extra:.\}"#, - )?; - - // Créer un client HTTP + /// Récupère uniquement la version du bundle courant sans télécharger les 7 MB. + /// + /// Utile pour savoir si le bundle a changé avant de déclencher une extraction + /// complète. Ne télécharge que la page de login (~5 KB). + pub async fn fetch_current_bundle_version() -> Result { let client = Client::builder() .user_agent("Mozilla/5.0 (compatible; PMOMusic/1.0)") + .timeout(BUNDLE_FETCH_TIMEOUT) .build()?; - println!("Récupération de la page de login..."); let login_page = client .get("https://play.qobuz.com/login") .send() @@ -40,31 +44,80 @@ impl Spoofer { .text() .await?; - // Extraire l'URL du bundle - let bundle_url_regex = Regex::new( - r#""#, + let re = Regex::new( + r#""#, )?; - let bundle_url = bundle_url_regex - .captures(&login_page) + re.captures(&login_page) .and_then(|cap| cap.get(1)) - .ok_or_else(|| anyhow::anyhow!("Impossible de trouver l'URL du bundle"))? - .as_str(); - - println!("Téléchargement du bundle depuis: {}", bundle_url); - let bundle_full_url = format!("https://play.qobuz.com{}", bundle_url); - let bundle = client.get(&bundle_full_url).send().await?.text().await?; - - println!("Bundle téléchargé ({} bytes)", bundle.len()); - - Ok(Self { - bundle, - seed_timezone_regex, - info_extras_regex_template, - app_id_regex, - }) + .map(|m| m.as_str().to_string()) + .ok_or_else(|| anyhow::anyhow!("Version bundle introuvable dans la page de login")) } - /// Extrait l'App ID depuis le bundle + /// Télécharge et parse le bundle Qobuz. Retry jusqu'à `BUNDLE_EXTRACTION_RETRIES` + /// fois en cas d'échec réseau ou d'extraction. + pub async fn new() -> Result { + let seed_timezone_regex = Regex::new( + r#"[a-z]\.initialSeed\("(?P[\w=]+)",window\.utimezone\.(?P[a-z]+)\)"#, + )?; + let info_extras_regex_template = + r#"name:"\w+/(?P{timezones})",info:"(?P[\w=]+)",extras:"(?P[\w=]+)""# + .to_string(); + let app_id_regex = Regex::new( + r#"production:\{api:\{appId:"(?P\d{9})",appSecret:"(?P\w{32})"\},braze:.\(.\(\{\},.\),\{\},\{apiKey:"([-0-9a-fA-F]{36})"\}\),extra:.\}"#, + )?; + + let client = Client::builder() + .user_agent("Mozilla/5.0 (compatible; PMOMusic/1.0)") + .timeout(BUNDLE_FETCH_TIMEOUT) + .build()?; + + let mut last_err: Option = None; + for attempt in 1..=(BUNDLE_EXTRACTION_RETRIES + 1) { + match Self::fetch_bundle(&client).await { + Ok((bundle, bundle_version)) => { + info!("[Spoofer] Bundle {} téléchargé ({} bytes)", bundle_version, bundle.len()); + return Ok(Self { + bundle, + bundle_version, + seed_timezone_regex, + info_extras_regex_template, + app_id_regex, + }); + } + Err(e) => { + warn!("[Spoofer] Tentative {}/{} échouée : {}", attempt, BUNDLE_EXTRACTION_RETRIES + 1, e); + last_err = Some(e); + } + } + } + Err(last_err.unwrap_or_else(|| anyhow::anyhow!("Échec téléchargement bundle"))) + } + + async fn fetch_bundle(client: &Client) -> Result<(String, String)> { + let login_page = client + .get("https://play.qobuz.com/login") + .send() + .await? + .text() + .await?; + + let bundle_url_regex = Regex::new( + r#""#, + )?; + let caps = bundle_url_regex + .captures(&login_page) + .ok_or_else(|| anyhow::anyhow!("URL bundle introuvable dans la page de login"))?; + let bundle_path = caps.get(1).unwrap().as_str(); + let bundle_version = caps.get(2).unwrap().as_str().to_string(); + + let bundle_url = format!("https://play.qobuz.com{}", bundle_path); + debug!("[Spoofer] Téléchargement bundle depuis {}", bundle_url); + + let bundle = client.get(&bundle_url).send().await?.text().await?; + Ok((bundle, bundle_version)) + } + + /// Extrait l'App ID depuis le bundle. pub fn get_app_id(&self) -> Result { let captures = self .app_id_regex @@ -78,10 +131,7 @@ impl Spoofer { .to_string()) } - /// Extrait l'appSecret depuis le bundle (secret MD5 à 32 caractères) - /// - /// Ce secret est utilisé directement par Qobuz (nouvelle méthode) - /// au lieu d'être XORé avec l'app_id + /// Extrait l'appSecret depuis le bundle (secret MD5 à 32 caractères). pub fn get_app_secret(&self) -> Result { let captures = self .app_id_regex @@ -95,37 +145,25 @@ impl Spoofer { .to_string()) } - /// Extrait les secrets depuis le bundle + /// Extrait les secrets timezone depuis le bundle. pub fn get_secrets(&self) -> Result> { - // Étape 1: Extraire tous les seed/timezone pairs let mut secrets: IndexMap> = IndexMap::new(); for captures in self.seed_timezone_regex.captures_iter(&self.bundle) { - let seed = captures - .name("seed") - .ok_or_else(|| anyhow::anyhow!("Groupe seed non trouvé"))? - .as_str(); - let timezone = captures - .name("timezone") - .ok_or_else(|| anyhow::anyhow!("Groupe timezone non trouvé"))? - .as_str(); - + let seed = captures.name("seed").unwrap().as_str(); + let timezone = captures.name("timezone").unwrap().as_str(); secrets .entry(timezone.to_string()) - .or_insert_with(Vec::new) + .or_default() .push(seed.to_string()); } - println!("Timezones trouvées: {:?}", secrets.keys()); + debug!("[Spoofer] Timezones trouvées : {:?}", secrets.keys().collect::>()); - // Étape 2: Réordonner - on met la deuxième timezone en premier - // (comme le fait le code Python avec move_to_end) if secrets.len() >= 2 { let keys: Vec = secrets.keys().cloned().collect(); let second_key = keys[1].clone(); let second_value = secrets.get(&second_key).unwrap().clone(); - - // Retirer et réinsérer pour le mettre en premier secrets.shift_remove(&second_key); let mut new_secrets = IndexMap::new(); new_secrets.insert(second_key, second_value); @@ -135,7 +173,6 @@ impl Spoofer { secrets = new_secrets; } - // Étape 3: Construire la regex pour info/extras let timezones_capitalized: Vec = secrets .keys() .map(|tz| { @@ -150,24 +187,12 @@ impl Spoofer { let info_extras_regex_str = self .info_extras_regex_template .replace("{timezones}", &timezones_capitalized.join("|")); - let info_extras_regex = Regex::new(&info_extras_regex_str)?; - // Étape 4: Extraire info et extras pour chaque timezone for captures in info_extras_regex.captures_iter(&self.bundle) { - let timezone_cap = captures - .name("timezone") - .ok_or_else(|| anyhow::anyhow!("Groupe timezone non trouvé"))? - .as_str(); - let info = captures - .name("info") - .ok_or_else(|| anyhow::anyhow!("Groupe info non trouvé"))? - .as_str(); - let extras = captures - .name("extras") - .ok_or_else(|| anyhow::anyhow!("Groupe extras non trouvé"))? - .as_str(); - + let timezone_cap = captures.name("timezone").unwrap().as_str(); + let info = captures.name("info").unwrap().as_str(); + let extras = captures.name("extras").unwrap().as_str(); let timezone_lower = timezone_cap.to_lowercase(); if let Some(vec) = secrets.get_mut(&timezone_lower) { vec.push(info.to_string()); @@ -175,31 +200,19 @@ impl Spoofer { } } - // Étape 5: Décoder les secrets en base64 let mut decoded_secrets = IndexMap::new(); for (timezone, parts) in secrets { let concatenated = parts.join(""); - - // Retirer les 44 derniers caractères (comme Python [:-44]) if concatenated.len() > 44 { let trimmed = &concatenated[..concatenated.len() - 44]; - - // Décoder en base64 match STANDARD.decode(trimmed) { Ok(decoded_bytes) => match String::from_utf8(decoded_bytes) { Ok(decoded_str) => { decoded_secrets.insert(timezone, decoded_str); } - Err(e) => { - eprintln!("Erreur UTF-8 pour timezone {}: {}", timezone, e); - } + Err(e) => warn!("[Spoofer] UTF-8 invalide pour timezone {}: {}", timezone, e), }, - Err(e) => { - eprintln!( - "Erreur de décodage base64 pour timezone {}: {}", - timezone, e - ); - } + Err(e) => warn!("[Spoofer] Base64 invalide pour timezone {}: {}", timezone, e), } } } diff --git a/pmoqobuz/src/api/user.rs b/pmoqobuz/src/api/user.rs index c8b23ced..2d636def 100644 --- a/pmoqobuz/src/api/user.rs +++ b/pmoqobuz/src/api/user.rs @@ -5,6 +5,7 @@ use super::QobuzApi; use crate::error::{QobuzError, Result}; use crate::models::*; use serde::Deserialize; +use std::collections::HashMap; use tracing::debug; /// Réponse paginée @@ -86,7 +87,11 @@ impl QobuzApi { } } - /// Récupère les tracks favorites de l'utilisateur + /// Récupère les tracks favorites de l'utilisateur. + /// + /// Si le secret s4 est disponible, les données de base retournées par + /// `/favorite/getUserFavorites` sont enrichies via `track/getList` pour + /// obtenir les métadonnées complètes (performer, sample_rate, bit_depth). pub async fn get_favorite_tracks(&self) -> Result> { let user_id = self.ensure_authenticated()?; debug!("Fetching favorite tracks for user {}", user_id); @@ -99,16 +104,32 @@ impl QobuzApi { let response: FavoritesResponse = self.get("/favorite/getUserFavorites", ¶ms).await?; - if let Some(tracks) = response.tracks { - Ok(tracks + let base_tracks: Vec = match response.tracks { + Some(tracks) => tracks .items .into_iter() .map(|t| QobuzApi::parse_track(t, None)) .filter(|t| t.streamable) - .collect()) - } else { - Ok(Vec::new()) + .collect(), + None => return Ok(Vec::new()), + }; + + if base_tracks.is_empty() { + return Ok(Vec::new()); } + + // Enrichissement via track/getList pour métadonnées complètes + let ordered_ids: Vec = base_tracks.iter().map(|t| t.id.clone()).collect(); + let id_refs: Vec<&str> = ordered_ids.iter().map(|s| s.as_str()).collect(); + let full_tracks = self.get_tracks_batch(&id_refs).await?; + let mut track_map: HashMap = + full_tracks.into_iter().map(|t| (t.id.clone(), t)).collect(); + let enriched: Vec = ordered_ids + .iter() + .filter_map(|id| track_map.remove(id.as_str())) + .collect(); + debug!("Fetched {} favorite tracks via track/getList", enriched.len()); + Ok(enriched) } /// Récupère les playlists de l'utilisateur diff --git a/pmoqobuz/src/api_rest.rs b/pmoqobuz/src/api_rest.rs index e1bcb500..02088ddd 100644 --- a/pmoqobuz/src/api_rest.rs +++ b/pmoqobuz/src/api_rest.rs @@ -74,6 +74,7 @@ pub fn create_router(state: QobuzState) -> Router { // Tracks .route("/tracks/:id", axum::routing::get(get_track)) .route("/tracks/:id/stream", axum::routing::get(get_stream_url)) + .route("/tracks/:id/flac", axum::routing::get(get_flac_stream)) // Artists .route("/artists/:id/albums", axum::routing::get(get_artist_albums)) .route( @@ -154,6 +155,35 @@ async fn get_stream_url( Ok(Json(serde_json::json!({ "url": url }))) } +/// `GET /tracks/:id/flac` — flux FLAC déchiffré en streaming progressif via CMAF. +/// +/// Retourne les bytes FLAC directement dans le corps de la réponse HTTP, au fur +/// et à mesure que les segments CMAF sont téléchargés et déchiffrés. Le client +/// (typiquement `pmocache::add_from_url`) peut commencer à consommer le contenu +/// avant la fin du téléchargement, préservant le progressive caching. +/// +/// `Content-Length` est une estimation calculée depuis la table des segments du +/// segment init — la valeur réelle peut différer de quelques octets. +#[cfg(feature = "pmoserver")] +async fn get_flac_stream( + State(state): State, + Path(id): Path, +) -> Result { + use axum::http::header; + use tokio_util::io::ReaderStream; + + let (reader, estimated_size) = state.client.open_cmaf_stream(&id).await?; + let stream = ReaderStream::new(reader); + let body = axum::body::Body::from_stream(stream); + + axum::response::Response::builder() + .status(axum::http::StatusCode::OK) + .header(header::CONTENT_TYPE, "audio/flac") + .header(header::CONTENT_LENGTH, estimated_size) + .body(body) + .map_err(|e| AppError(QobuzError::Other(e.to_string()))) +} + #[cfg(feature = "pmoserver")] async fn get_artist_albums( State(state): State, diff --git a/pmoqobuz/src/client.rs b/pmoqobuz/src/client.rs index 4ad8b8e4..8b976415 100644 --- a/pmoqobuz/src/client.rs +++ b/pmoqobuz/src/client.rs @@ -236,14 +236,32 @@ impl QobuzClient { /// - Quand aucun appid/secret n'est configuré /// - Quand les credentials configurés sont invalides/expirés async fn try_spoofer_fallback(config: &Config) -> Result { + // Vérification bon marché : si la version du bundle n'a pas changé, + // re-télécharger les 7 MB ne donnera pas de meilleurs secrets. + // On court-circuite l'extraction et on passe directement au DEFAULT_APP_ID. + if let Ok(Some(cached_version)) = config.get_qobuz_bundle_version() { + match crate::api::Spoofer::fetch_current_bundle_version().await { + Ok(current) if current == cached_version => { + info!( + "[Spoofer] Bundle inchangé ({}) — secret invalide pour une autre raison, skip extraction", + current + ); + return QobuzApi::new(DEFAULT_APP_ID); + } + Ok(new_version) => { + info!("[Spoofer] Bundle rotaté : {} → {}, re-extraction", cached_version, new_version); + } + Err(e) => { + debug!("[Spoofer] Impossible de vérifier la version du bundle : {}", e); + } + } + } + if let Some((app_id, secret)) = Self::fetch_spoofer_credentials(config).await? { - // Use raw secret from Spoofer (no XOR) return QobuzApi::with_raw_secret(app_id, &secret); } - info!( - "✗ No valid secret found from Spoofer, falling back to DEFAULT_APP_ID without secret" - ); + info!("✗ No valid secret found from Spoofer, falling back to DEFAULT_APP_ID without secret"); QobuzApi::new(DEFAULT_APP_ID) } @@ -288,18 +306,15 @@ impl QobuzClient { timezone ); - // Save both appid and the working secret - if let Err(e) = config.set_qobuz_appid(&app_id) - { + // Save appid, secret, and bundle version + if let Err(e) = config.set_qobuz_appid(&app_id) { debug!("Could not save appid: {}", e); } - if let Err(e) = - config.set_qobuz_spoofer_secret(secret) - { - debug!( - "Could not save spoofer secret: {}", - e - ); + if let Err(e) = config.set_qobuz_spoofer_secret(secret) { + debug!("Could not save spoofer secret: {}", e); + } + if let Err(e) = config.set_qobuz_bundle_version(spoofer.bundle_version()) { + debug!("Could not save bundle version: {}", e); } return Ok(Some(( @@ -625,6 +640,44 @@ impl QobuzClient { Ok(track) } + /// Récupère les détails d'un batch de tracks via track/getList. + /// + /// Les tracks retournées sont mises en cache individuellement. + /// Voir `QobuzApi::get_tracks_batch` pour le détail du comportement. + pub async fn get_tracks_batch(&self, track_ids: &[&str]) -> Result> { + // Séparer les IDs déjà en cache de ceux à récupérer + let mut cached: std::collections::HashMap = + std::collections::HashMap::new(); + let mut missing_ids: Vec<&str> = Vec::new(); + + for &id in track_ids { + if let Some(track) = self.cache.get_track(id).await { + cached.insert(id.to_string(), track); + } else { + missing_ids.push(id); + } + } + + if !missing_ids.is_empty() { + let fetched = self + .call_with_auth_repair("get_tracks_batch", || { + self.api.get_tracks_batch(&missing_ids) + }) + .await?; + + for track in fetched { + self.cache.put_track(track.id.clone(), track.clone()).await; + cached.insert(track.id.clone(), track); + } + } + + // Restituer dans l'ordre d'entrée + Ok(track_ids + .iter() + .filter_map(|id| cached.remove(*id)) + .collect()) + } + /// Récupère l'URL de streaming d'une track pub async fn get_stream_url(&self, track_id: &str) -> Result { // Vérifier le cache d'abord @@ -973,6 +1026,152 @@ impl QobuzClient { }) .await } + + // ============ CMAF (streaming moderne) ============ + + /// Prépare le streaming CMAF pour un track : dérive les clés et fetche + /// le segment d'initialisation. Retourne une `CmafStreamInfo` prête à l'emploi. + /// + /// Pour télécharger la totalité du FLAC déchiffré d'un coup, utilisez + /// plutôt `download_cmaf_full`. Pour streamer segment par segment, utilisez + /// les données de `CmafStreamInfo` avec `cmaf::fetch_all_segments`. + /// + /// # Erreurs + /// + /// * `QobuzError::NotFound` — track non disponible sur Qobuz + /// * `QobuzError::Unauthorized` — session expirée (réparée automatiquement) + /// * `QobuzError::Other` — erreur CMAF (clé invalide, segment init corrompu, etc.) + pub async fn get_cmaf_stream_info( + &self, + track_id: &str, + ) -> Result { + let format_id = self.api.format().id() as u32; + + let file_url = self + .call_with_auth_repair("get_cmaf_file_url", || { + self.api.get_cmaf_file_url(track_id, format_id) + }) + .await?; + + let bit_depth = file_url.resolved_bit_depth(); + let url_template = file_url + .url_template + .ok_or_else(|| QobuzError::Other("CMAF file/url: url_template absent".into()))?; + let key_str = file_url + .key + .ok_or_else(|| QobuzError::Other("CMAF file/url: key absent".into()))?; + + let (_session_id, infos) = self + .call_with_auth_repair("ensure_cmaf_session", || { + self.api.ensure_cmaf_session() + }) + .await?; + + let setup = crate::cmaf::setup_streaming( + url_template, + &key_str, + &infos, + file_url.n_segments, + file_url.format_id.unwrap_or(format_id), + file_url.sampling_rate, + bit_depth, + ) + .await?; + + Ok(crate::models::CmafStreamInfo { + url_template: setup.url_template, + n_segments: setup.n_segments, + content_key: setup.content_key, + flac_header: setup.flac_header, + segment_table: setup.segment_table, + format_id: setup.format_id, + sampling_rate: setup.sampling_rate, + bit_depth: setup.bit_depth, + }) + } + + /// Ouvre un flux FLAC déchiffré en mode progressif pour un track CMAF. + /// + /// Retourne un `(AsyncRead, taille_estimée)`. Le flux commence à produire + /// du FLAC dès que le premier segment est déchiffré — le cache peut commencer + /// à servir avant la fin du téléchargement (progressive caching préservé). + /// + /// C'est la méthode à utiliser pour alimenter `pmocache::add_from_reader`. + pub async fn open_cmaf_stream( + &self, + track_id: &str, + ) -> Result<(impl tokio::io::AsyncRead + Send + Unpin + 'static, u64)> { + let format_id = self.api.format().id() as u32; + + let file_url = self + .call_with_auth_repair("get_cmaf_file_url", || { + self.api.get_cmaf_file_url(track_id, format_id) + }) + .await?; + + let bit_depth = file_url.resolved_bit_depth(); + let url_template = file_url + .url_template + .ok_or_else(|| QobuzError::Other("CMAF: url_template absent".into()))?; + let key_str = file_url + .key + .ok_or_else(|| QobuzError::Other("CMAF: key absent".into()))?; + + let (_session_id, infos) = self + .call_with_auth_repair("ensure_cmaf_session", || self.api.ensure_cmaf_session()) + .await?; + + crate::cmaf::open_flac_stream( + url_template, + key_str, + infos, + file_url.n_segments, + file_url.format_id.unwrap_or(format_id), + file_url.sampling_rate, + bit_depth, + ) + .await + } + + /// Télécharge un track complet via CMAF et retourne les bytes FLAC déchiffrés. + /// + /// Bloquant jusqu'à la fin du téléchargement. Pour les gros fichiers Hi-Res, + /// préférer `open_cmaf_stream` qui préserve le progressive caching. + pub async fn download_cmaf_full(&self, track_id: &str) -> Result> { + let format_id = self.api.format().id() as u32; + + let file_url = self + .call_with_auth_repair("get_cmaf_file_url", || { + self.api.get_cmaf_file_url(track_id, format_id) + }) + .await?; + + let bit_depth = file_url.resolved_bit_depth(); + let url_template = file_url + .url_template + .ok_or_else(|| QobuzError::Other("CMAF file/url: url_template absent".into()))?; + let key_str = file_url + .key + .ok_or_else(|| QobuzError::Other("CMAF file/url: key absent".into()))?; + + let (_session_id, infos) = self + .call_with_auth_repair("ensure_cmaf_session", || { + self.api.ensure_cmaf_session() + }) + .await?; + + crate::cmaf::download_full( + url_template, + &key_str, + &infos, + file_url.n_segments, + file_url.format_id.unwrap_or(format_id), + file_url.sampling_rate, + bit_depth, + None, + ) + .await + } } #[cfg(test)] diff --git a/pmoqobuz/src/cmaf/crypto.rs b/pmoqobuz/src/cmaf/crypto.rs new file mode 100644 index 00000000..dc27ffc9 --- /dev/null +++ b/pmoqobuz/src/cmaf/crypto.rs @@ -0,0 +1,154 @@ +use aes::cipher::{BlockDecryptMut, KeyIvInit, StreamCipher}; +use base64::{Engine, engine::general_purpose::URL_SAFE_NO_PAD}; +use hkdf::Hkdf; +use md5::{Digest, Md5}; +use sha2::Sha256; + +use super::error::CmafError; + +type Aes128CbcDec = cbc::Decryptor; +type Aes128Ctr = ctr::Ctr128BE; + +/// Seed publique extraite du bundle web Qobuz. Valeur IKM pour HKDF. +pub const CMAF_SEED: &str = "abb21364945c0583309667d13ca3d93a"; + +fn hex_decode(hex: &str) -> Vec { + (0..hex.len()) + .step_by(2) + .map(|i| u8::from_str_radix(&hex[i..i + 2], 16).unwrap_or(0)) + .collect() +} + +/// Dérive la clé de session 16 octets depuis le champ `infos` de session/start. +/// +/// Format infos : `"salt_b64url.info_b64url"` +/// `seed` est le CMAF_SEED hex-encodé utilisé comme IKM HKDF. +pub fn derive_session_key(seed: &str, infos: &str) -> Result<[u8; 16], CmafError> { + let parts: Vec<&str> = infos.split('.').collect(); + if parts.len() < 2 { + return Err(CmafError::InvalidInfos( + "session infos doit avoir au moins 2 parties séparées par des points".into(), + )); + } + + let salt = URL_SAFE_NO_PAD.decode(parts[0])?; + let info = URL_SAFE_NO_PAD.decode(parts[1])?; + let ikm = hex_decode(seed); + + let hk = Hkdf::::new(Some(&salt), &ikm); + let mut okm = [0u8; 16]; + hk.expand(&info, &mut okm).map_err(|_| CmafError::HkdfExpand)?; + + Ok(okm) +} + +/// Déroule la clé de contenu par track avec la clé de session. +/// +/// Format key_str : `"qbz-1.wrapped_key_b64url.iv_b64url"` +pub fn unwrap_content_key(session_key: &[u8; 16], key_str: &str) -> Result<[u8; 16], CmafError> { + let parts: Vec<&str> = key_str.split('.').collect(); + if parts.len() < 3 { + return Err(CmafError::InvalidKey( + "key string doit avoir au moins 3 parties séparées par des points".into(), + )); + } + + let wrapped = URL_SAFE_NO_PAD.decode(parts[1])?; + let iv = URL_SAFE_NO_PAD.decode(parts[2])?; + + if iv.len() != 16 { + return Err(CmafError::InvalidKey(format!( + "IV de dérobage doit faire 16 octets, reçu {}", + iv.len() + ))); + } + + let mut buf = wrapped.clone(); + let decrypted = + Aes128CbcDec::new(session_key.into(), iv.as_slice().into()) + .decrypt_padded_mut::(&mut buf) + .map_err(|e| CmafError::AesDecrypt(format!("AES-CBC unwrap échoué: {e}")))?; + + if decrypted.len() != 16 { + return Err(CmafError::InvalidKey(format!( + "clé déroullée doit faire 16 octets, obtenu {}", + decrypted.len() + ))); + } + + let mut key = [0u8; 16]; + key.copy_from_slice(decrypted); + Ok(key) +} + +/// Déchiffre une frame FLAC en place avec AES-128-CTR. +/// +/// `iv_8` = IV 8 octets du segment UUID box, complété à zéro jusqu'à 16 octets. +pub fn decrypt_frame(content_key: &[u8; 16], iv_8: &[u8; 8], data: &mut [u8]) { + let mut nonce = [0u8; 16]; + nonce[..8].copy_from_slice(iv_8); + Aes128Ctr::new(content_key.into(), &nonce.into()).apply_keystream(data); +} + +/// Calcule la signature MD5 pour les appels API CMAF de Qobuz. +/// +/// Concatène method + paires clé-valeur triées + timestamp + seed, +/// puis retourne le digest MD5 hexadécimal minuscule. +pub fn compute_request_sig( + method: &str, + args: &std::collections::BTreeMap<&str, String>, + timestamp: &str, + seed: &str, +) -> String { + let mut hasher = Md5::new(); + hasher.update(method.as_bytes()); + for (k, v) in args { + hasher.update(k.as_bytes()); + hasher.update(v.as_bytes()); + } + hasher.update(timestamp.as_bytes()); + hasher.update(seed.as_bytes()); + + format!("{:x}", hasher.finalize()) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_compute_request_sig_deterministe() { + let mut args = std::collections::BTreeMap::new(); + args.insert("profile", "qbz-1".to_string()); + let sig1 = compute_request_sig("sessionstart", &args, "1775500000", CMAF_SEED); + let sig2 = compute_request_sig("sessionstart", &args, "1775500000", CMAF_SEED); + assert_eq!(sig1.len(), 32); + assert_eq!(sig1, sig2); + } + + #[test] + fn test_decrypt_frame_aller_retour() { + let key = [1u8, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15, 16]; + let iv = [1u8, 2, 3, 4, 5, 6, 7, 8]; + let original = b"Hello FLAC frame data here!".to_vec(); + let mut data = original.clone(); + decrypt_frame(&key, &iv, &mut data); + assert_ne!(data, original); + // AES-CTR est son propre inverse + decrypt_frame(&key, &iv, &mut data); + assert_eq!(data, original); + } + + #[test] + fn test_derive_session_key_infos_invalide() { + let result = derive_session_key(CMAF_SEED, "pas_de_point"); + assert!(result.is_err()); + } + + #[test] + fn test_unwrap_content_key_format_invalide() { + let key = [0u8; 16]; + let result = unwrap_content_key(&key, "seulement.deux"); + assert!(result.is_err()); + } +} diff --git a/pmoqobuz/src/cmaf/error.rs b/pmoqobuz/src/cmaf/error.rs new file mode 100644 index 00000000..f8e13f70 --- /dev/null +++ b/pmoqobuz/src/cmaf/error.rs @@ -0,0 +1,17 @@ +use thiserror::Error; + +#[derive(Debug, Error)] +pub enum CmafError { + #[error("Invalid infos format: {0}")] + InvalidInfos(String), + #[error("Invalid key format: {0}")] + InvalidKey(String), + #[error("Base64 decode error: {0}")] + Base64(#[from] base64::DecodeError), + #[error("HKDF expand error")] + HkdfExpand, + #[error("AES decrypt error: {0}")] + AesDecrypt(String), + #[error("CMAF parse error: {0}")] + ParseError(String), +} diff --git a/pmoqobuz/src/cmaf/mod.rs b/pmoqobuz/src/cmaf/mod.rs new file mode 100644 index 00000000..fd6963b1 --- /dev/null +++ b/pmoqobuz/src/cmaf/mod.rs @@ -0,0 +1,395 @@ +//! Pipeline CMAF (Common Media Application Format) pour Qobuz. +//! +//! Qobuz utilise CMAF avec chiffrement AES-CTR par frame sur CDN Akamai. +//! C'est le pipeline de l'app Android v9.7+ qui remplace l'endpoint legacy +//! `/track/getFileUrl`. +//! +//! # Pipeline +//! +//! 1. `/session/start` → `{ session_id, infos, expires_at }` +//! 2. `/file/url` → `{ url_template, key (enveloppé), n_segments, ... }` +//! 3. Session key = `HKDF(CMAF_SEED, infos)` +//! 4. Content key = AES-CBC-unwrap(session_key, key) +//! 5. Segment init (s=0) → header FLAC + table des segments +//! 6. Pour chaque s=1..n_segments : fetch → parse crypto boxes → déchiffrement AES-CTR + +pub mod crypto; +pub mod error; +pub mod parser; + +pub use crypto::{compute_request_sig, decrypt_frame, derive_session_key, unwrap_content_key, CMAF_SEED}; +pub use error::CmafError; +pub use parser::{ + parse_init_segment, parse_segment_crypto, FrameEntry, InitInfo, SegmentCrypto, + SegmentTableEntry, +}; + +use std::sync::Arc; +use tokio::io::AsyncRead; +use tokio::sync::Semaphore; +use tracing::{debug, info, warn}; + +use crate::error::{QobuzError, Result}; +use crate::retry::{classify_reqwest, classify_status, retry_transient, FetchError, DEFAULT_MAX_ATTEMPTS}; + +/// Concurrence max pour le fetch de segments CMAF. +/// 3 segments en vol est le compromis optimal — le CDN Akamai rate-limite +/// au-delà de ~5 requêtes parallèles par IP sur des fenêtres de 1s. +pub const CMAF_PREFETCH_CONCURRENCY: usize = 3; + +/// Callback de progression pour les fonctions de téléchargement. +pub type CmafProgressCallback = Arc; + +/// Un tick de progression. `segments_completed` est cumulatif (1..=n). +#[derive(Debug, Clone, Copy)] +pub struct CmafProgressUpdate { + pub segments_completed: u32, + pub n_segments: u32, + pub bytes_this_segment: u64, +} + +/// Info réunies depuis le segment init, suffisantes pour démarrer le streaming. +pub struct CmafStreamingInfo { + pub url_template: String, + pub n_segments: u8, + pub content_key: [u8; 16], + pub flac_header: Vec, + pub segment_table: Vec, + pub format_id: u32, + pub sampling_rate: Option, + pub bit_depth: Option, +} + +/// Construit un client reqwest dédié aux fetches CDN Akamai. +fn build_cdn_client() -> Result { + reqwest::Client::builder() + .connect_timeout(std::time::Duration::from_secs(10)) + .build() + .map_err(|e| QobuzError::Http(e)) +} + +/// Fetch une URL CDN en bytes avec retry sur les erreurs transitoires. +/// Un 404/403 échoue immédiatement (terminal). 5xx et 429 → retry avec backoff. +async fn fetch_bytes_with_retry( + http: &reqwest::Client, + url: &str, + log_tag: &str, +) -> std::result::Result, FetchError> { + retry_transient( + DEFAULT_MAX_ATTEMPTS, + log_tag, + FetchError::is_transient, + |_attempt| async move { + let response = http + .get(url) + .header("User-Agent", "Mozilla/5.0") + .send() + .await + .map_err(|e| classify_reqwest(&e, "fetch"))?; + let status = response.status(); + if !status.is_success() { + return Err(classify_status(status, "fetch")); + } + response + .bytes() + .await + .map(|b| b.to_vec()) + .map_err(|e| classify_reqwest(&e, "lecture")) + }, + ) + .await +} + +/// Fetch les segments 1..=n_segments avec contrôle de concurrence. +/// Déclenche le callback de progression une fois par segment complété. +/// Les résultats sont retournés triés par index de segment. +async fn fetch_all_segments( + http: &reqwest::Client, + url_template: &str, + n_segments: u8, + log_tag: &str, + on_progress: Option, +) -> Result>> { + let semaphore = Arc::new(Semaphore::new(CMAF_PREFETCH_CONCURRENCY)); + let completed_count = Arc::new(std::sync::atomic::AtomicU32::new(0)); + let mut handles = Vec::with_capacity(n_segments as usize); + + for seg_idx in 1u8..=n_segments { + let sem = semaphore.clone(); + let http = http.clone(); + let seg_url = url_template.replace("$SEGMENT$", &seg_idx.to_string()); + let log_tag = log_tag.to_string(); + let progress = on_progress.clone(); + let counter = completed_count.clone(); + + handles.push(tokio::spawn(async move { + let permit = sem.acquire_owned().await + .map_err(|e| format!("semaphore: {}", e))?; + + let seg_data = fetch_bytes_with_retry(&http, &seg_url, &format!("{} seg {}", log_tag, seg_idx)) + .await + .map_err(|e| format!("[{}] seg {} fetch: {}", log_tag, seg_idx, e))?; + + let bytes_this_segment = seg_data.len() as u64; + if let Some(cb) = progress { + let done = counter.fetch_add(1, std::sync::atomic::Ordering::Relaxed) + 1; + cb(CmafProgressUpdate { + segments_completed: done, + n_segments: n_segments as u32, + bytes_this_segment, + }); + } + + // Pause avant de libérer le slot pour respecter les limites CDN + tokio::time::sleep(std::time::Duration::from_millis(500)).await; + drop(permit); + + Ok::<(u8, Vec), String>((seg_idx, seg_data)) + })); + } + + let mut segments: Vec<(u8, Vec)> = Vec::with_capacity(handles.len()); + for handle in handles { + let (idx, data) = handle + .await + .map_err(|e| QobuzError::Other(format!("[{}] panic task: {}", log_tag, e)))? + .map_err(|e| QobuzError::Other(format!("[{}] téléchargement échoué: {}", log_tag, e)))?; + segments.push((idx, data)); + } + segments.sort_by_key(|(idx, _)| *idx); + Ok(segments.into_iter().map(|(_, data)| data).collect()) +} + +/// Déchiffre un segment CMAF et ajoute les frames FLAC dans `output`. +/// +/// Optimisation : extend + decrypt in-place, zéro allocation par frame. +fn decrypt_one_segment( + seg_data: &[u8], + content_key: &[u8; 16], + output: &mut Vec, + seg_idx: usize, +) -> Result<()> { + let crypto = parse_segment_crypto(seg_data) + .map_err(|e| QobuzError::Other(format!("CMAF seg {} parse: {}", seg_idx, e)))?; + + let mut data_pos = crypto.data_offset; + for entry in &crypto.entries { + let frame_end = data_pos + entry.size as usize; + if frame_end > seg_data.len() { + return Err(QobuzError::Other(format!("CMAF seg {} débordement frame", seg_idx))); + } + let start = output.len(); + output.extend_from_slice(&seg_data[data_pos..frame_end]); + if entry.flags != 0 { + decrypt_frame(content_key, &entry.iv, &mut output[start..]); + } + data_pos = frame_end; + } + if data_pos < crypto.mdat_end && crypto.mdat_end <= seg_data.len() { + output.extend_from_slice(&seg_data[data_pos..crypto.mdat_end]); + } + Ok(()) +} + +/// Déchiffre une séquence de segments CMAF chiffrés et écrit les frames FLAC dans `output`. +pub fn decrypt_segments_into( + segments: &[Vec], + content_key: &[u8; 16], + output: &mut Vec, +) -> Result<()> { + for (seg_idx, seg_data) in segments.iter().enumerate() { + decrypt_one_segment(seg_data, content_key, output, seg_idx + 1)?; + } + Ok(()) +} + +/// Prépare le streaming CMAF : dérive les clés, fetche le segment init. +/// Ne télécharge PAS les segments audio — l'appelant les streame en arrière-plan. +pub async fn setup_streaming( + url_template: String, + key_str: &str, + infos: &str, + n_segments: u8, + format_id: u32, + sampling_rate: Option, + bit_depth: Option, +) -> Result { + let session_key = derive_session_key(CMAF_SEED, infos) + .map_err(|e| QobuzError::Other(format!("dérivation clé session: {}", e)))?; + let content_key = unwrap_content_key(&session_key, key_str) + .map_err(|e| QobuzError::Other(format!("dérobage clé contenu: {}", e)))?; + + let http = build_cdn_client()?; + let init_url = url_template.replace("$SEGMENT$", "0"); + + info!("[CMAF] Fetch segment init: {}", &init_url[..init_url.len().min(60)]); + + let init_data = fetch_bytes_with_retry(&http, &init_url, "CMAF init") + .await + .map_err(|e| QobuzError::Other(format!("fetch segment init: {}", e)))?; + + let init_info = parse_init_segment(&init_data) + .map_err(|e| QobuzError::Other(format!("parse segment init: {}", e)))?; + + info!( + "[CMAF] Init: header FLAC {}B, {} segments dans la table, n_segments API={}", + init_info.flac_header.len(), + init_info.segment_table.len(), + n_segments, + ); + if init_info.segment_table.len() != n_segments as usize { + warn!( + "[CMAF] ÉCART: table={} entrées mais API dit n_segments={}", + init_info.segment_table.len(), + n_segments, + ); + } + + Ok(CmafStreamingInfo { + url_template, + n_segments, + content_key, + flac_header: init_info.flac_header, + segment_table: init_info.segment_table, + format_id, + sampling_rate, + bit_depth, + }) +} + +/// Ouvre un flux FLAC déchiffré en mode progressif. +/// +/// Retourne un `AsyncRead` qui produit les bytes FLAC au fur et à mesure que les +/// segments CMAF sont téléchargés et déchiffrés. Le lecteur peut commencer à +/// consommer le flux (et le cache à le servir) avant que tous les segments +/// soient téléchargés — le progressive caching est préservé. +/// +/// Le deuxième élément est la taille totale estimée en bytes, calculée depuis +/// le header FLAC et la table des segments du segment init. +/// +/// # Erreurs +/// +/// Retourne une erreur si la dérivation des clés ou le fetch du segment init +/// échoue. Les erreurs de segments ultérieurs provoquent la fermeture du pipe +/// (le lecteur verra un EOF prématuré). +pub async fn open_flac_stream( + url_template: String, + key_str: String, + infos: String, + n_segments: u8, + format_id: u32, + sampling_rate: Option, + bit_depth: Option, +) -> Result<(impl AsyncRead + Send + Unpin + 'static, u64)> { + use tokio::io::AsyncWriteExt; + + let setup = setup_streaming( + url_template, &key_str, &infos, n_segments, + format_id, sampling_rate, bit_depth, + ).await?; + + let estimated_size = (setup.flac_header.len() + + setup.segment_table.iter().map(|s| s.byte_len as usize).sum::()) as u64; + + // Pipe 256 KB : assez grand pour absorber un segment FLAC Hi-Res typique + // sans bloquer le producteur, assez petit pour ne pas sur-allouer. + let (mut writer, reader) = tokio::io::duplex(256 * 1024); + + tokio::spawn(async move { + if let Err(e) = writer.write_all(&setup.flac_header).await { + warn!("[CMAF-STREAM] écriture header FLAC: {}", e); + return; + } + + let http = match build_cdn_client() { + Ok(c) => c, + Err(e) => { warn!("[CMAF-STREAM] client HTTP: {}", e); return; } + }; + + // Lancer tous les fetches avec concurrence bornée par le sémaphore. + // Les handles sont stockés dans l'ordre — on les consomme en ordre. + let sem = Arc::new(Semaphore::new(CMAF_PREFETCH_CONCURRENCY)); + let mut handles = Vec::with_capacity(setup.n_segments as usize); + + for seg_idx in 1u8..=setup.n_segments { + let sem = sem.clone(); + let http = http.clone(); + let url = setup.url_template.replace("$SEGMENT$", &seg_idx.to_string()); + + handles.push(tokio::spawn(async move { + let permit = sem.acquire_owned().await + .map_err(|e| format!("sémaphore: {}", e))?; + let result = fetch_bytes_with_retry(&http, &url, &format!("CMAF-STREAM seg {}", seg_idx)) + .await + .map_err(|e| format!("seg {} fetch: {}", seg_idx, e)); + // Cooldown CDN Akamai : retenir le slot 500ms avant de libérer, + // identique à fetch_all_segments, pour éviter le rate-limiting. + tokio::time::sleep(std::time::Duration::from_millis(500)).await; + drop(permit); + result + })); + } + + // Consommer dans l'ordre : segment 1 d'abord, puis 2, etc. + // Les segments prêts en avance attendent dans leur JoinHandle. + let mut buf = Vec::new(); + for (i, handle) in handles.into_iter().enumerate() { + let seg_idx = i + 1; + let seg_data = match handle.await { + Ok(Ok(data)) => data, + Ok(Err(e)) => { warn!("[CMAF-STREAM] seg {}: {}", seg_idx, e); return; } + Err(e) => { warn!("[CMAF-STREAM] seg {} panique: {}", seg_idx, e); return; } + }; + + buf.clear(); + if let Err(e) = decrypt_one_segment(&seg_data, &setup.content_key, &mut buf, seg_idx) { + warn!("[CMAF-STREAM] seg {}: déchiffrement: {}", seg_idx, e); + return; + } + if let Err(e) = writer.write_all(&buf).await { + // Le lecteur a fermé le pipe (ex: playback stoppé) — arrêt silencieux. + debug!("[CMAF-STREAM] seg {}: pipe fermé ({})", seg_idx, e); + return; + } + debug!("[CMAF-STREAM] seg {}/{} → {} B", seg_idx, setup.n_segments, buf.len()); + } + + info!("[CMAF-STREAM] complet : {} segments, {:.2} MB estimés", + setup.n_segments, + estimated_size as f64 / (1024.0 * 1024.0)); + // writer dropped ici → EOF propre sur le reader + }); + + Ok((reader, estimated_size)) +} + +/// Télécharge un track CMAF complet et retourne les bytes FLAC déchiffrés. +pub async fn download_full( + url_template: String, + key_str: &str, + infos: &str, + n_segments: u8, + format_id: u32, + sampling_rate: Option, + bit_depth: Option, + on_progress: Option, +) -> Result> { + let setup = setup_streaming(url_template, key_str, infos, n_segments, format_id, sampling_rate, bit_depth).await?; + let http = build_cdn_client()?; + + let total_size: usize = setup.flac_header.len() + + setup.segment_table.iter().map(|s| s.byte_len as usize).sum::(); + + let segments = fetch_all_segments(&http, &setup.url_template, setup.n_segments, "CMAF-FULL", on_progress).await?; + + let mut output = Vec::with_capacity(total_size); + output.extend_from_slice(&setup.flac_header); + decrypt_segments_into(&segments, &setup.content_key, &mut output)?; + + debug!( + "[CMAF-FULL] Complet: {:.2} MB FLAC, attendu {:.2} MB", + output.len() as f64 / (1024.0 * 1024.0), + total_size as f64 / (1024.0 * 1024.0), + ); + Ok(output) +} diff --git a/pmoqobuz/src/cmaf/parser.rs b/pmoqobuz/src/cmaf/parser.rs new file mode 100644 index 00000000..5a234ab2 --- /dev/null +++ b/pmoqobuz/src/cmaf/parser.rs @@ -0,0 +1,253 @@ +use super::error::CmafError; + +const QBZ_INIT_UUID: [u8; 16] = [ + 0xc7, 0xc7, 0x5d, 0xf0, 0xfd, 0xd9, 0x51, 0xe9, + 0x8f, 0xc2, 0x29, 0x71, 0xe4, 0xac, 0xf8, 0xd2, +]; +const QBZ_SEGMENT_UUID: [u8; 16] = [ + 0x3b, 0x42, 0x12, 0x92, 0x56, 0xf3, 0x5f, 0x75, + 0x92, 0x36, 0x63, 0xb6, 0x9a, 0x1f, 0x52, 0xb2, +]; +const FLAC_MAGIC: &[u8; 4] = b"fLaC"; + +/// Taille en octets et nombre d'échantillons d'un segment. +#[derive(Debug, Clone)] +pub struct SegmentTableEntry { + /// Taille des données FLAC déchiffrées de ce segment. + pub byte_len: u32, + /// Nombre d'échantillons audio dans ce segment. + pub sample_count: u32, +} + +/// Header FLAC et table des segments extraits du segment d'initialisation. +pub struct InitInfo { + pub flac_header: Vec, + /// Tailles par segment (indices 0..n_segments-1 correspondent aux segments 1..n_segments). + pub segment_table: Vec, +} + +/// Une entrée de frame dans le segment UUID box. +pub struct FrameEntry { + pub size: u32, + pub flags: u16, + pub iv: [u8; 8], +} + +/// Informations crypto parsées depuis le UUID box d'un segment audio. +pub struct SegmentCrypto { + /// Offset vers le début des données audio (payload mdat). + pub data_offset: usize, + /// Fin du contenu de la mdat box. + pub mdat_end: usize, + pub entries: Vec, +} + +/// Parcourt les boxes ISO BMFF et trouve le premier UUID box correspondant à `target_uuid`. +/// Retourne `(payload_start, box_end)` où payload_start est après les 16 octets UUID. +fn find_uuid_box(data: &[u8], target_uuid: &[u8; 16]) -> Option<(usize, usize)> { + let mut pos = 0; + while pos + 8 <= data.len() { + let size = read_box_size(data, pos); + if size < 8 || pos + size > data.len() { + break; + } + if &data[pos + 4..pos + 8] == b"uuid" && pos + 24 <= data.len() { + if &data[pos + 8..pos + 24] == target_uuid.as_ref() { + return Some((pos + 24, pos + size)); + } + } + pos += size; + } + None +} + +/// Parse le segment d'initialisation (segment 0) pour extraire le header FLAC et la table des segments. +pub fn parse_init_segment(data: &[u8]) -> Result { + let (payload_start, box_end) = find_uuid_box(data, &QBZ_INIT_UUID) + .ok_or_else(|| CmafError::ParseError("segment init: QBZ_INIT_UUID box non trouvé".into()))?; + + let payload = &data[payload_start..box_end]; + parse_init_uuid_payload(payload) +} + +/// Parse un segment audio pour extraire les informations crypto par frame. +pub fn parse_segment_crypto(data: &[u8]) -> Result { + let mut uuid_box_start: Option = None; + let mut mdat_end = data.len(); + + let mut pos = 0; + while pos + 8 <= data.len() { + let size = read_box_size(data, pos); + if size < 8 || pos + size > data.len() { + break; + } + let box_type = &data[pos + 4..pos + 8]; + if box_type == b"uuid" && pos + 24 <= data.len() { + if &data[pos + 8..pos + 24] == QBZ_SEGMENT_UUID.as_ref() { + uuid_box_start = Some(pos); + } + } else if box_type == b"mdat" { + mdat_end = pos + size; + } + pos += size; + } + + let box_start = uuid_box_start + .ok_or_else(|| CmafError::ParseError("segment audio: QBZ_SEGMENT_UUID box non trouvé".into()))?; + + parse_segment_uuid_payload(data, box_start, mdat_end) +} + +fn parse_init_uuid_payload(payload: &[u8]) -> Result { + // Layout payload: + // [4B padding/version] + // [4B track_id] + // [4B file_id] + // [4B sample_rate] + // [1B bits_per_sample] + // [1B channels + 2B padding] + // [6B total_samples_count] + // [2B raw_data_len] + // [raw_data_len bytes: contient le header FLAC] + // [1B key_id_len] + // [key_id_len bytes: key_id] + // [2B segment_count] + // Par segment: [4B byte_len][4B sample_count] + + if payload.len() < 28 { + return Err(CmafError::ParseError("payload init UUID trop court".into())); + } + + let mut a = 4; // version/padding + a += 4; // track_id + a += 4; // file_id + a += 4; // sample_rate + a += 1; // bits_per_sample + a += 3; // channels + padding + a += 6; // total_samples_count + + if a + 2 > payload.len() { + return Err(CmafError::ParseError("payload init UUID tronqué au raw_len".into())); + } + let raw_len = u16::from_be_bytes([payload[a], payload[a + 1]]) as usize; + a += 2; + + let raw_data = &payload[a..a + raw_len.min(payload.len() - a)]; + a += raw_len; + + let flac_pos = raw_data + .windows(4) + .position(|w| w == FLAC_MAGIC) + .ok_or_else(|| CmafError::ParseError("payload init UUID: magic fLaC non trouvé".into()))?; + + // fLaC (4) + STREAMINFO block header (4) + STREAMINFO data (34) = 42 octets + let header_len = 4 + 4 + 34; + if flac_pos + header_len > raw_data.len() { + return Err(CmafError::ParseError("payload init UUID: STREAMINFO tronqué".into())); + } + + let mut flac_header = raw_data[flac_pos..flac_pos + header_len].to_vec(); + // Marquer le dernier bloc de métadonnées + flac_header[4] |= 0x80; + + if a + 1 > payload.len() { + return Ok(InitInfo { flac_header, segment_table: Vec::new() }); + } + let key_id_len = payload[a] as usize; + a += 1 + key_id_len; + + let mut segment_table = Vec::new(); + if a + 2 <= payload.len() { + let seg_count = u16::from_be_bytes([payload[a], payload[a + 1]]) as usize; + a += 2; + + for _ in 0..seg_count { + if a + 8 > payload.len() { + break; + } + let byte_len = u32::from_be_bytes([payload[a], payload[a + 1], payload[a + 2], payload[a + 3]]); + a += 4; + let sample_count = u32::from_be_bytes([payload[a], payload[a + 1], payload[a + 2], payload[a + 3]]); + a += 4; + segment_table.push(SegmentTableEntry { byte_len, sample_count }); + } + } + + tracing::debug!( + "Init UUID: {} segments dans la table, header FLAC {} octets", + segment_table.len(), + flac_header.len() + ); + + Ok(InitInfo { flac_header, segment_table }) +} + +fn parse_segment_uuid_payload( + data: &[u8], + uuid_box_start: usize, + mdat_end: usize, +) -> Result { + // Layout après box header (8) + UUID (16) = offset 24 depuis uuid_box_start: + // [4B version/padding] + // [4B data_offset_raw] — offset depuis uuid_box_start vers les données audio + // [1B iv_size] + // [3B frame_count (24-bit BE)] + // Par frame: [4B size][2B skip][2B flags][iv_size bytes IV] + + let base = uuid_box_start + 24; + if base + 12 > data.len() { + return Err(CmafError::ParseError( + "payload segment UUID trop court pour le header".into(), + )); + } + + let mut a = base + 4; // skip 4-byte version/padding + + let data_offset_raw = u32::from_be_bytes([data[a], data[a + 1], data[a + 2], data[a + 3]]); + let data_offset = uuid_box_start + data_offset_raw as usize; + a += 4; + + let iv_size = data[a] as usize; + a += 1; + + let frame_count = + ((data[a] as usize) << 16) | ((data[a + 1] as usize) << 8) | (data[a + 2] as usize); + a += 3; + + let entry_size = 4 + 2 + 2 + iv_size; + if a + frame_count * entry_size > data.len() { + return Err(CmafError::ParseError(format!( + "segment UUID: données insuffisantes pour {frame_count} entrées de {entry_size} octets" + ))); + } + + let mut entries = Vec::with_capacity(frame_count); + for _ in 0..frame_count { + let size = u32::from_be_bytes([data[a], data[a + 1], data[a + 2], data[a + 3]]); + a += 4; + a += 2; // 2 octets inconnus + let flags = u16::from_be_bytes([data[a], data[a + 1]]); + a += 2; + + let mut iv = [0u8; 8]; + let copy_len = iv_size.min(8); + iv[..copy_len].copy_from_slice(&data[a..a + copy_len]); + a += iv_size; + + entries.push(FrameEntry { size, flags, iv }); + } + + Ok(SegmentCrypto { data_offset, mdat_end, entries }) +} + +fn read_box_size(data: &[u8], pos: usize) -> usize { + if pos + 8 > data.len() { + return 0; + } + let s = u32::from_be_bytes([data[pos], data[pos + 1], data[pos + 2], data[pos + 3]]); + match s { + 0 => data.len() - pos, + 1..=7 => 0, + s => s as usize, + } +} diff --git a/pmoqobuz/src/config_ext.rs b/pmoqobuz/src/config_ext.rs index 085a63a8..f1535883 100644 --- a/pmoqobuz/src/config_ext.rs +++ b/pmoqobuz/src/config_ext.rs @@ -238,6 +238,13 @@ pub trait QobuzConfigExt { /// Active ou désactive le rate limiting fn set_qobuz_rate_limiting_enabled(&self, enabled: bool) -> Result<()>; + + /// Version du bundle Qobuz extrait en dernier (ex : `"8.1.0-b019"`). + /// Permet de détecter une rotation de bundle sans télécharger les 7 MB. + fn get_qobuz_bundle_version(&self) -> Result>; + + /// Persiste la version du bundle après une extraction réussie. + fn set_qobuz_bundle_version(&self, version: &str) -> Result<()>; } impl QobuzConfigExt for Config { @@ -483,4 +490,18 @@ impl QobuzConfigExt for Config { Value::Bool(enabled), ) } + + fn get_qobuz_bundle_version(&self) -> Result> { + match self.get_value(&["accounts", "qobuz", "bundle_version"]) { + Ok(Value::String(s)) if !s.is_empty() => Ok(Some(s)), + _ => Ok(None), + } + } + + fn set_qobuz_bundle_version(&self, version: &str) -> Result<()> { + self.set_value( + &["accounts", "qobuz", "bundle_version"], + Value::String(version.to_string()), + ) + } } diff --git a/pmoqobuz/src/lazy_provider.rs b/pmoqobuz/src/lazy_provider.rs index cb83a36e..8615a9f7 100644 --- a/pmoqobuz/src/lazy_provider.rs +++ b/pmoqobuz/src/lazy_provider.rs @@ -68,6 +68,18 @@ impl LazyProvider for QobuzLazyProvider { async fn get_url(&self, lazy_pk: &str) -> Result { let track_id = self.track_id_from_lazy(lazy_pk)?; + + // L'endpoint CMAF local n'existe que si la feature `pmoserver` est active + // (route /qobuz/tracks/:id/flac enregistrée). Sans cette feature, fallback + // sur l'ancienne API Qobuz directe. + #[cfg(feature = "pmoserver")] + { + let base = std::env::var("PMO_SERVER_URL") + .unwrap_or_else(|_| "http://localhost:8080".to_string()); + return Ok(format!("{}/qobuz/tracks/{}/flac", base.trim_end_matches('/'), track_id)); + } + + #[cfg(not(feature = "pmoserver"))] self.client .get_stream_url(track_id) .await diff --git a/pmoqobuz/src/lib.rs b/pmoqobuz/src/lib.rs index ce3b7f5f..e5d82d26 100644 --- a/pmoqobuz/src/lib.rs +++ b/pmoqobuz/src/lib.rs @@ -208,6 +208,7 @@ pub mod api; pub mod cache; +pub mod cmaf; pub mod client; pub mod config_ext; pub mod didl; @@ -216,6 +217,7 @@ pub mod disk_cache; pub mod error; mod lazy_provider; pub mod models; +pub mod retry; pub mod source; // Extension pmoserver (feature-gated) diff --git a/pmoqobuz/src/models.rs b/pmoqobuz/src/models.rs index 50e81400..71932c57 100644 --- a/pmoqobuz/src/models.rs +++ b/pmoqobuz/src/models.rs @@ -208,6 +208,31 @@ pub struct StreamInfo { pub expires_at: DateTime, } +/// Informations pour le streaming CMAF d'un track. +/// +/// Retourné par `QobuzClient::get_cmaf_stream_info`. +/// L'appelant peut ensuite appeler `pmoqobuz::cmaf::download_full` pour +/// obtenir les bytes FLAC déchiffrés, ou streamer les segments manuellement. +#[derive(Debug)] +pub struct CmafStreamInfo { + /// Modèle d'URL des segments, avec le placeholder `$SEGMENT$`. + pub url_template: String, + /// Nombre de segments audio (hors segment init s=0). + pub n_segments: u8, + /// Clé AES-128 déchiffrée pour décoder les frames FLAC. + pub content_key: [u8; 16], + /// Header FLAC extrait du segment init (à placer en tête du flux décodé). + pub flac_header: Vec, + /// Table des segments avec taille et compte d'échantillons. + pub segment_table: Vec, + /// Format ID Qobuz. + pub format_id: u32, + /// Fréquence d'échantillonnage (Hz). + pub sampling_rate: Option, + /// Profondeur de bits. + pub bit_depth: Option, +} + /// Format audio demandé pour le streaming #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] #[repr(u8)] diff --git a/pmoqobuz/src/pmoserver_impl.rs b/pmoqobuz/src/pmoserver_impl.rs index 17961002..224f713d 100644 --- a/pmoqobuz/src/pmoserver_impl.rs +++ b/pmoqobuz/src/pmoserver_impl.rs @@ -29,6 +29,7 @@ use crate::api_rest::{create_router, QobuzState}; use crate::client::QobuzClient; +use crate::config_ext::QobuzConfigExt; use crate::pmoserver_ext::QobuzServerExt; use anyhow::Result; use pmoconfig::Config; diff --git a/pmoqobuz/src/retry.rs b/pmoqobuz/src/retry.rs new file mode 100644 index 00000000..660a798d --- /dev/null +++ b/pmoqobuz/src/retry.rs @@ -0,0 +1,194 @@ +//! Helper de retry avec backoff exponentiel pour les fetches réseau transitoires. +//! +//! Un blip réseau transitoire (5xx, timeout, connexion reset, 429) sur le +//! `file/url` du prochain track ou un segment CMAF ne doit pas être fatal. +//! Ce module retente les échecs *transitoires* avec backoff exponentiel +//! et laisse les échecs *terminaux* (404 "disparu définitivement", erreurs auth) +//! se propager immédiatement. + +use std::future::Future; +use std::time::Duration; + +/// Nombre de tentatives : 1 initiale + 2 retrys. +pub const DEFAULT_MAX_ATTEMPTS: u32 = 3; + +/// Erreur de fetch taguée selon son caractère retryable. +#[derive(Debug)] +pub enum FetchError { + /// Vaut la peine de retenter : erreur réseau/timeout/connect/body, 5xx, ou 429. + Transient(String), + /// Ne vaut pas la peine de retenter : 4xx (sauf 429), ou échec définitif. + Terminal(String), +} + +impl FetchError { + pub fn is_transient(&self) -> bool { + matches!(self, FetchError::Transient(_)) + } +} + +impl std::fmt::Display for FetchError { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + FetchError::Transient(s) | FetchError::Terminal(s) => write!(f, "{}", s), + } + } +} + +/// Vrai pour les erreurs reqwest qui valent la peine d'être retentées. +pub fn reqwest_is_transient(e: &reqwest::Error) -> bool { + e.is_timeout() || e.is_connect() || e.is_request() || e.is_body() +} + +/// Classifie une erreur reqwest en `FetchError`. +/// Toutes les erreurs transport reqwest sont traitées comme transitoires. +pub fn classify_reqwest(e: &reqwest::Error, context: &str) -> FetchError { + FetchError::Transient(format!("{}: {}", context, e)) +} + +/// Classifie un status HTTP non-succès en `FetchError`. +/// 5xx et 429 → transitoire ; tout le reste (404, 403, ...) → terminal. +pub fn classify_status(status: reqwest::StatusCode, context: &str) -> FetchError { + let msg = format!("{}: HTTP {}", context, status); + if status.is_server_error() || status == reqwest::StatusCode::TOO_MANY_REQUESTS { + FetchError::Transient(msg) + } else { + FetchError::Terminal(msg) + } +} + +/// Backoff exponentiel avec jitter pour la N-ième tentative (base 1) : +/// ~250 ms, ~500 ms, ~1 s, plafonné à 2 s, plus jusqu'à +25% de jitter. +/// Le jitter est dérivé de l'horloge pour éviter une dépendance à `rand`. +fn backoff_delay(attempt: u32) -> Duration { + let exp = attempt.saturating_sub(1).min(3); + let base_ms = 250u64.saturating_mul(1u64 << exp).min(2000); + let jitter_span = base_ms / 4; + let jitter = if jitter_span > 0 { + let nanos = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map(|d| d.subsec_nanos() as u64) + .unwrap_or(0); + nanos % (jitter_span + 1) + } else { + 0 + }; + Duration::from_millis(base_ms + jitter) +} + +/// Exécute `op` (qui prend le numéro de tentative base 1) et retente +/// tant qu'il retourne une erreur transitoire, avec backoff entre les tentatives. +/// Les erreurs terminales et la dernière tentative retournent immédiatement. +pub async fn retry_transient( + max_attempts: u32, + log_tag: &str, + is_transient: impl Fn(&E) -> bool, + mut op: F, +) -> std::result::Result +where + F: FnMut(u32) -> Fut, + Fut: Future>, + E: std::fmt::Display, +{ + let mut attempt = 1; + loop { + match op(attempt).await { + Ok(value) => return Ok(value), + Err(err) => { + if attempt >= max_attempts || !is_transient(&err) { + return Err(err); + } + let delay = backoff_delay(attempt); + tracing::warn!( + "[{}] erreur transitoire tentative {}/{}: {} — retry dans {}ms", + log_tag, + attempt, + max_attempts, + err, + delay.as_millis() + ); + tokio::time::sleep(delay).await; + attempt += 1; + } + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use std::sync::atomic::{AtomicU32, Ordering}; + use std::sync::Arc; + + #[tokio::test] + async fn reussit_au_premier_essai() { + let calls = Arc::new(AtomicU32::new(0)); + let c = calls.clone(); + let r: std::result::Result = + retry_transient(3, "test", FetchError::is_transient, |_| { + let c = c.clone(); + async move { + c.fetch_add(1, Ordering::Relaxed); + Ok(42) + } + }) + .await; + assert_eq!(r.unwrap(), 42); + assert_eq!(calls.load(Ordering::Relaxed), 1); + } + + #[tokio::test] + async fn retente_transitoire_puis_reussit() { + let calls = Arc::new(AtomicU32::new(0)); + let c = calls.clone(); + let r: std::result::Result = + retry_transient(3, "test", FetchError::is_transient, |attempt| { + let c = c.clone(); + async move { + c.fetch_add(1, Ordering::Relaxed); + if attempt < 3 { + Err(FetchError::Transient("503".into())) + } else { + Ok(7) + } + } + }) + .await; + assert_eq!(r.unwrap(), 7); + assert_eq!(calls.load(Ordering::Relaxed), 3); + } + + #[tokio::test] + async fn terminal_ne_retente_pas() { + let calls = Arc::new(AtomicU32::new(0)); + let c = calls.clone(); + let r: std::result::Result = + retry_transient(3, "test", FetchError::is_transient, |_| { + let c = c.clone(); + async move { + c.fetch_add(1, Ordering::Relaxed); + Err(FetchError::Terminal("404".into())) + } + }) + .await; + assert!(r.is_err()); + assert_eq!(calls.load(Ordering::Relaxed), 1); + } + + #[tokio::test] + async fn abandonne_apres_max_tentatives() { + let calls = Arc::new(AtomicU32::new(0)); + let c = calls.clone(); + let r: std::result::Result = + retry_transient(3, "test", FetchError::is_transient, |_| { + let c = c.clone(); + async move { + c.fetch_add(1, Ordering::Relaxed); + Err(FetchError::Transient("timeout".into())) + } + }) + .await; + assert!(r.is_err()); + assert_eq!(calls.load(Ordering::Relaxed), 3); + } +} diff --git a/version.txt b/version.txt index cd3dcaaf..6b9aa4e6 100644 --- a/version.txt +++ b/version.txt @@ -1 +1 @@ -0.3.50 +0.3.51