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.
This commit is contained in:
119
Blackboard/Architecture/pmoqobuz_ameliorations_qbz.md
Normal file
119
Blackboard/Architecture/pmoqobuz_ameliorations_qbz.md
Normal file
@@ -0,0 +1,119 @@
|
|||||||
|
# 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 — **À FAIRE** (priorité haute)
|
||||||
|
|
||||||
|
**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` — **À FAIRE** (priorité haute)
|
||||||
|
|
||||||
|
**Problème** : lors du chargement d'une playlist, `pmoqobuz` fait N requêtes `track/get` individuelles
|
||||||
|
(une par track). Sur une playlist de 200 tracks, c'est 200 requêtes séquentielles.
|
||||||
|
|
||||||
|
**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)
|
||||||
|
|
||||||
|
**Impact pour pmoqobuz** : réduire le temps de chargement d'une playlist de plusieurs minutes à
|
||||||
|
quelques secondes. À brancher dans `catalog.rs` et dans le chargement des favoris.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## 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) | À faire |
|
||||||
|
| 3 | Batch `track/getList` | Faible | Élevé (performances) | À faire |
|
||||||
|
| 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 |
|
||||||
1
Cargo.lock
generated
1
Cargo.lock
generated
@@ -4098,6 +4098,7 @@ dependencies = [
|
|||||||
"thiserror 2.0.17",
|
"thiserror 2.0.17",
|
||||||
"tokio",
|
"tokio",
|
||||||
"tokio-test",
|
"tokio-test",
|
||||||
|
"tokio-util",
|
||||||
"tracing",
|
"tracing",
|
||||||
"tracing-subscriber",
|
"tracing-subscriber",
|
||||||
"utoipa",
|
"utoipa",
|
||||||
|
|||||||
@@ -61,6 +61,7 @@ pmodidl = { path = "../pmodidl" }
|
|||||||
# Intégration avec pmoserver pour l'API HTTP
|
# Intégration avec pmoserver pour l'API HTTP
|
||||||
pmoserver = { path = "../pmoserver", optional = true }
|
pmoserver = { path = "../pmoserver", optional = true }
|
||||||
axum = { version = "0.8", optional = true }
|
axum = { version = "0.8", optional = true }
|
||||||
|
tokio-util = { workspace = true, optional = true }
|
||||||
rusqlite = { version = "0.37", features = ["bundled"], optional = true }
|
rusqlite = { version = "0.37", features = ["bundled"], optional = true }
|
||||||
|
|
||||||
# Documentation OpenAPI
|
# Documentation OpenAPI
|
||||||
@@ -76,7 +77,7 @@ pmocache = { path = "../pmocache" }
|
|||||||
[features]
|
[features]
|
||||||
default = ["async-trait-support"]
|
default = ["async-trait-support"]
|
||||||
# Feature pour activer les extensions pmoserver
|
# 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)
|
# Feature pour activer le support serveur (cache registry)
|
||||||
server = ["pmosource/server"]
|
server = ["pmosource/server"]
|
||||||
# Feature cache (deprecated - toujours actif maintenant)
|
# Feature cache (deprecated - toujours actif maintenant)
|
||||||
|
|||||||
@@ -74,6 +74,7 @@ pub fn create_router(state: QobuzState) -> Router {
|
|||||||
// Tracks
|
// Tracks
|
||||||
.route("/tracks/:id", axum::routing::get(get_track))
|
.route("/tracks/:id", axum::routing::get(get_track))
|
||||||
.route("/tracks/:id/stream", axum::routing::get(get_stream_url))
|
.route("/tracks/:id/stream", axum::routing::get(get_stream_url))
|
||||||
|
.route("/tracks/:id/flac", axum::routing::get(get_flac_stream))
|
||||||
// Artists
|
// Artists
|
||||||
.route("/artists/:id/albums", axum::routing::get(get_artist_albums))
|
.route("/artists/:id/albums", axum::routing::get(get_artist_albums))
|
||||||
.route(
|
.route(
|
||||||
@@ -154,6 +155,35 @@ async fn get_stream_url(
|
|||||||
Ok(Json(serde_json::json!({ "url": 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")]
|
#[cfg(feature = "pmoserver")]
|
||||||
async fn get_artist_albums(
|
async fn get_artist_albums(
|
||||||
State(state): State<QobuzState>,
|
State(state): State<QobuzState>,
|
||||||
|
|||||||
@@ -1037,10 +1037,53 @@ impl QobuzClient {
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// 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.
|
/// 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,
|
/// Bloquant jusqu'à la fin du téléchargement. Pour les gros fichiers Hi-Res,
|
||||||
/// préférer `get_cmaf_stream_info` + streaming segment par segment.
|
/// 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>> {
|
pub async fn download_cmaf_full(&self, track_id: &str) -> Result<Vec<u8>> {
|
||||||
let format_id = self.api.format().id() as u32;
|
let format_id = self.api.format().id() as u32;
|
||||||
|
|
||||||
|
|||||||
@@ -25,6 +25,7 @@ pub use parser::{
|
|||||||
};
|
};
|
||||||
|
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
|
use tokio::io::AsyncRead;
|
||||||
use tokio::sync::Semaphore;
|
use tokio::sync::Semaphore;
|
||||||
use tracing::{debug, info, warn};
|
use tracing::{debug, info, warn};
|
||||||
|
|
||||||
@@ -159,35 +160,45 @@ async fn fetch_all_segments(
|
|||||||
Ok(segments.into_iter().map(|(_, data)| data).collect())
|
Ok(segments.into_iter().map(|(_, data)| data).collect())
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Déchiffre une séquence de segments CMAF chiffrés et écrit les frames FLAC dans `output`.
|
/// Déchiffre un segment CMAF et ajoute les frames FLAC dans `output`.
|
||||||
///
|
///
|
||||||
/// Optimisation hot-path : extend + decrypt in-place plutôt que copie triple.
|
/// Optimisation : extend + decrypt in-place, zéro allocation par frame.
|
||||||
pub fn decrypt_segments_into(
|
fn decrypt_one_segment(
|
||||||
segments: &[Vec<u8>],
|
seg_data: &[u8],
|
||||||
content_key: &[u8; 16],
|
content_key: &[u8; 16],
|
||||||
output: &mut Vec<u8>,
|
output: &mut Vec<u8>,
|
||||||
|
seg_idx: usize,
|
||||||
) -> Result<()> {
|
) -> Result<()> {
|
||||||
for (seg_idx, seg_data) in segments.iter().enumerate() {
|
|
||||||
let log_idx = seg_idx + 1;
|
|
||||||
let crypto = parse_segment_crypto(seg_data)
|
let crypto = parse_segment_crypto(seg_data)
|
||||||
.map_err(|e| QobuzError::Other(format!("CMAF seg {} parse: {}", log_idx, e)))?;
|
.map_err(|e| QobuzError::Other(format!("CMAF seg {} parse: {}", seg_idx, e)))?;
|
||||||
|
|
||||||
let mut data_pos = crypto.data_offset;
|
let mut data_pos = crypto.data_offset;
|
||||||
for entry in &crypto.entries {
|
for entry in &crypto.entries {
|
||||||
let frame_end = data_pos + entry.size as usize;
|
let frame_end = data_pos + entry.size as usize;
|
||||||
if frame_end > seg_data.len() {
|
if frame_end > seg_data.len() {
|
||||||
return Err(QobuzError::Other(format!("CMAF seg {} débordement frame", log_idx)));
|
return Err(QobuzError::Other(format!("CMAF seg {} débordement frame", seg_idx)));
|
||||||
}
|
}
|
||||||
let output_start = output.len();
|
let start = output.len();
|
||||||
output.extend_from_slice(&seg_data[data_pos..frame_end]);
|
output.extend_from_slice(&seg_data[data_pos..frame_end]);
|
||||||
if entry.flags != 0 {
|
if entry.flags != 0 {
|
||||||
decrypt_frame(content_key, &entry.iv, &mut output[output_start..]);
|
decrypt_frame(content_key, &entry.iv, &mut output[start..]);
|
||||||
}
|
}
|
||||||
data_pos = frame_end;
|
data_pos = frame_end;
|
||||||
}
|
}
|
||||||
if data_pos < crypto.mdat_end && crypto.mdat_end <= seg_data.len() {
|
if data_pos < crypto.mdat_end && crypto.mdat_end <= seg_data.len() {
|
||||||
output.extend_from_slice(&seg_data[data_pos..crypto.mdat_end]);
|
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(())
|
Ok(())
|
||||||
}
|
}
|
||||||
@@ -246,6 +257,112 @@ pub async fn setup_streaming(
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// 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.
|
/// Télécharge un track CMAF complet et retourne les bytes FLAC déchiffrés.
|
||||||
pub async fn download_full(
|
pub async fn download_full(
|
||||||
url_template: String,
|
url_template: String,
|
||||||
|
|||||||
@@ -68,6 +68,18 @@ impl LazyProvider for QobuzLazyProvider {
|
|||||||
|
|
||||||
async fn get_url(&self, lazy_pk: &str) -> Result<String> {
|
async fn get_url(&self, lazy_pk: &str) -> Result<String> {
|
||||||
let track_id = self.track_id_from_lazy(lazy_pk)?;
|
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
|
self.client
|
||||||
.get_stream_url(track_id)
|
.get_stream_url(track_id)
|
||||||
.await
|
.await
|
||||||
|
|||||||
@@ -29,6 +29,7 @@
|
|||||||
|
|
||||||
use crate::api_rest::{create_router, QobuzState};
|
use crate::api_rest::{create_router, QobuzState};
|
||||||
use crate::client::QobuzClient;
|
use crate::client::QobuzClient;
|
||||||
|
use crate::config_ext::QobuzConfigExt;
|
||||||
use crate::pmoserver_ext::QobuzServerExt;
|
use crate::pmoserver_ext::QobuzServerExt;
|
||||||
use anyhow::Result;
|
use anyhow::Result;
|
||||||
use pmoconfig::Config;
|
use pmoconfig::Config;
|
||||||
|
|||||||
Reference in New Issue
Block a user