4 Commits

Author SHA1 Message Date
bcfffa6df8 feat: implement batch track loading via track/getList endpoint
Replaces sequential track requests with a chunked batch enrichment workflow using MD5-signed POST requests (50-ID windows). Introduces cache-first fetching, order preservation via HashMap lookups, and graceful fallbacks for playlist and favorite loading.
2026-06-11 11:07:09 +02:00
e670e8e0d8 feat: enhance Qobuz bundle extraction with caching and retry logic
Refactor Qobuz bundle fetching to improve resilience and performance. Add HTTP timeout and retry configuration, and implement a retry loop that returns version metadata. Introduce bundle version caching and persistence to skip redundant extractions when the bundle remains unchanged. Replace console logging with structured tracing macros and update the architectural tracking document to mark these improvements as complete.
2026-06-11 09:51:15 +02:00
c7f67e7c7c feat(pmoqobuz): add progressive CMAF FLAC streaming support
Introduce a `/tracks/:id/flac` endpoint and `open_cmaf_stream` client method to enable incremental decryption and caching of FLAC audio. Refactor segment decryption into a reusable helper and stream decrypted frames via a bounded tokio channel. Add `tokio-util` as an optional dependency, update the `pmoserver` feature, and implement conditional routing for local proxy versus direct API access. Update the architecture roadmap to document these improvements.
2026-06-11 09:43:44 +02:00
a97bfe6148 feat: implement Qobuz CMAF streaming and AES decryption
Introduces a new CMAF streaming module for Qobuz that handles progressive segment fetching and full-track decryption. The implementation adds HKDF session key derivation, AES-128-CBC key unwrapping, and AES-128-CTR frame decryption, alongside an ISO BMFF parser for extracting FLAC headers and segment metadata. It also integrates automatic session renewal with double-checked locking, a generic async retry mechanism for transient HTTP errors, and centralized error handling. Finally, it exposes `CmafStreamInfo` and async streaming methods in the public API, updating cryptographic dependencies accordingly.
2026-06-11 08:52:28 +02:00
24 changed files with 2093 additions and 117 deletions

BIN
.DS_Store vendored

Binary file not shown.

View File

@@ -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 |

46
Cargo.lock generated
View File

@@ -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",

View File

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

View File

@@ -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)

View File

@@ -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<TrackResponse>,
}
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<Vec<Track>> {
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<Vec<Track>> {
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, &timestamp, &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<u64> = 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<Track> {
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<Vec<Track>> {
debug!("Fetching tracks for playlist {}", playlist_id);
const PAGE_SIZE: u32 = 50;
let mut all_tracks = Vec::new();
let mut ordered_ids: Vec<String> = Vec::new();
let mut fallback_tracks: Vec<Track> = 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<String, Track> =
full_tracks.into_iter().map(|t| (t.id.clone(), t)).collect();
let enriched: Vec<Track> = 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

261
pmoqobuz/src/api/cmaf.rs Normal file
View File

@@ -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<String>,
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<String>,
/// Clé de contenu enveloppée, format `"qbz-1.wrapped_b64url.iv_b64url"`.
#[serde(default)]
pub key: Option<String>,
/// Nombre de segments audio (hors segment init).
pub n_segments: u8,
#[serde(default)]
pub format_id: Option<u32>,
#[serde(default)]
pub mime_type: Option<String>,
#[serde(default)]
pub sampling_rate: Option<u32>,
/// Profondeur de bits (champ v1 de l'API).
#[serde(default)]
pub bits_depth: Option<u32>,
/// Profondeur de bits (champ v2 de l'API).
#[serde(default)]
pub bit_depth: Option<u32>,
}
impl CmafFileUrlResponse {
/// Retourne la profondeur de bits en préférant `bits_depth` puis `bit_depth`.
pub fn resolved_bit_depth(&self) -> Option<u32> {
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, &timestamp.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, &timestamp.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<Option<CmafSession>>,
}
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", &timestamp.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<CmafFileUrlResponse> {
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", &timestamp.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)
}

View File

@@ -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<Option<String>>,
/// 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<T: DeserializeOwned>(
&self,
endpoint: &str,
query_params: &[(&str, &str)],
body: serde_json::Value,
) -> Result<T> {
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<T: DeserializeOwned>(
&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<crate::api::cmaf::CmafFileUrlResponse> {
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<T: DeserializeOwned>(
&self,

View File

@@ -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();

View File

@@ -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<Self> {
// Expressions régulières (équivalent Python)
let seed_timezone_regex = Regex::new(
r#"[a-z]\.initialSeed\("(?P<seed>[\w=]+)",window\.utimezone\.(?P<timezone>[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<timezone>{timezones})",info:"(?P<info>[\w=]+)",extras:"(?P<extras>[\w=]+)""#
.to_string();
let app_id_regex = Regex::new(
r#"production:\{api:\{appId:"(?P<app_id>\d{9})",appSecret:"(?P<secret>\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<String> {
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#"<script src="(/resources/\d+\.\d+\.\d+-[a-z]\d{3}/bundle\.js)"></script>"#,
let re = Regex::new(
r#"<script src="/resources/(\d+\.\d+\.\d+-[a-z]\d{3})/bundle\.js"></script>"#,
)?;
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<Self> {
let seed_timezone_regex = Regex::new(
r#"[a-z]\.initialSeed\("(?P<seed>[\w=]+)",window\.utimezone\.(?P<timezone>[a-z]+)\)"#,
)?;
let info_extras_regex_template =
r#"name:"\w+/(?P<timezone>{timezones})",info:"(?P<info>[\w=]+)",extras:"(?P<extras>[\w=]+)""#
.to_string();
let app_id_regex = Regex::new(
r#"production:\{api:\{appId:"(?P<app_id>\d{9})",appSecret:"(?P<secret>\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<anyhow::Error> = 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#"<script src="(/resources/(\d+\.\d+\.\d+-[a-z]\d{3})/bundle\.js)"></script>"#,
)?;
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<String> {
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<String> {
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<IndexMap<String, String>> {
// Étape 1: Extraire tous les seed/timezone pairs
let mut secrets: IndexMap<String, Vec<String>> = 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::<Vec<_>>());
// É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<String> = 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<String> = 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),
}
}
}

View File

@@ -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<Vec<Track>> {
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", &params).await?;
if let Some(tracks) = response.tracks {
Ok(tracks
let base_tracks: Vec<Track> = 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<String> = 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<String, Track> =
full_tracks.into_iter().map(|t| (t.id.clone(), t)).collect();
let enriched: Vec<Track> = 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

View File

@@ -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<QobuzState>,
Path(id): Path<String>,
) -> Result<Response, AppError> {
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<QobuzState>,

View File

@@ -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<QobuzApi> {
// 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<Vec<Track>> {
// Séparer les IDs déjà en cache de ceux à récupérer
let mut cached: std::collections::HashMap<String, Track> =
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<String> {
// 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<crate::models::CmafStreamInfo> {
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<Vec<u8>> {
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)]

154
pmoqobuz/src/cmaf/crypto.rs Normal file
View File

@@ -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<aes::Aes128>;
type Aes128Ctr = ctr::Ctr128BE<aes::Aes128>;
/// Seed publique extraite du bundle web Qobuz. Valeur IKM pour HKDF.
pub const CMAF_SEED: &str = "abb21364945c0583309667d13ca3d93a";
fn hex_decode(hex: &str) -> Vec<u8> {
(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::<Sha256>::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::<aes::cipher::block_padding::Pkcs7>(&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());
}
}

View File

@@ -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),
}

395
pmoqobuz/src/cmaf/mod.rs Normal file
View File

@@ -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<dyn Fn(CmafProgressUpdate) + Send + Sync>;
/// 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<u8>,
pub segment_table: Vec<SegmentTableEntry>,
pub format_id: u32,
pub sampling_rate: Option<u32>,
pub bit_depth: Option<u32>,
}
/// Construit un client reqwest dédié aux fetches CDN Akamai.
fn build_cdn_client() -> Result<reqwest::Client> {
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<Vec<u8>, 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<CmafProgressCallback>,
) -> Result<Vec<Vec<u8>>> {
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<u8>), String>((seg_idx, seg_data))
}));
}
let mut segments: Vec<(u8, Vec<u8>)> = 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<u8>,
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<u8>],
content_key: &[u8; 16],
output: &mut Vec<u8>,
) -> 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<u32>,
bit_depth: Option<u32>,
) -> Result<CmafStreamingInfo> {
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<u32>,
bit_depth: Option<u32>,
) -> 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::<usize>()) 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<u32>,
bit_depth: Option<u32>,
on_progress: Option<CmafProgressCallback>,
) -> Result<Vec<u8>> {
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::<usize>();
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)
}

253
pmoqobuz/src/cmaf/parser.rs Normal file
View File

@@ -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<u8>,
/// Tailles par segment (indices 0..n_segments-1 correspondent aux segments 1..n_segments).
pub segment_table: Vec<SegmentTableEntry>,
}
/// 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<FrameEntry>,
}
/// 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<InitInfo, CmafError> {
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<SegmentCrypto, CmafError> {
let mut uuid_box_start: Option<usize> = 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<InitInfo, CmafError> {
// 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<SegmentCrypto, CmafError> {
// 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,
}
}

View File

@@ -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<Option<String>>;
/// 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<Option<String>> {
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()),
)
}
}

View File

@@ -68,6 +68,18 @@ impl LazyProvider for QobuzLazyProvider {
async fn get_url(&self, lazy_pk: &str) -> Result<String> {
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

View File

@@ -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)

View File

@@ -208,6 +208,31 @@ pub struct StreamInfo {
pub expires_at: DateTime<Utc>,
}
/// 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<u8>,
/// Table des segments avec taille et compte d'échantillons.
pub segment_table: Vec<crate::cmaf::SegmentTableEntry>,
/// Format ID Qobuz.
pub format_id: u32,
/// Fréquence d'échantillonnage (Hz).
pub sampling_rate: Option<u32>,
/// Profondeur de bits.
pub bit_depth: Option<u32>,
}
/// Format audio demandé pour le streaming
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[repr(u8)]

View File

@@ -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;

194
pmoqobuz/src/retry.rs Normal file
View File

@@ -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<F, Fut, T, E>(
max_attempts: u32,
log_tag: &str,
is_transient: impl Fn(&E) -> bool,
mut op: F,
) -> std::result::Result<T, E>
where
F: FnMut(u32) -> Fut,
Fut: Future<Output = std::result::Result<T, E>>,
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<u32, FetchError> =
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<u32, FetchError> =
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<u32, FetchError> =
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<u32, FetchError> =
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);
}
}

View File

@@ -1 +1 @@
0.3.50
0.3.51