From c7f67e7c7c07a5c995eb4d2c2869d30fc92375d6 Mon Sep 17 00:00:00 2001 From: Eric Coissac Date: Thu, 11 Jun 2026 09:36:21 +0200 Subject: [PATCH] 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. --- .DS_Store | Bin 18436 -> 18436 bytes .../pmoqobuz_ameliorations_qbz.md | 119 +++++++++++++ Cargo.lock | 1 + pmoqobuz/Cargo.toml | 3 +- pmoqobuz/src/api_rest.rs | 30 ++++ pmoqobuz/src/client.rs | 45 ++++- pmoqobuz/src/cmaf/mod.rs | 161 +++++++++++++++--- pmoqobuz/src/lazy_provider.rs | 12 ++ pmoqobuz/src/pmoserver_impl.rs | 1 + 9 files changed, 348 insertions(+), 24 deletions(-) create mode 100644 Blackboard/Architecture/pmoqobuz_ameliorations_qbz.md diff --git a/.DS_Store b/.DS_Store index 58e2f0f507227499b386c330e07dc0c68d077235..06359845aaf749f8201b6b54df7765981768c1a9 100644 GIT binary patch delta 161 zcmZpfz}PZ@aRa}CA{T=mLjglBLq0JHTL;Bw=j7()_fB45EHBvy)QHvK&3_b3 h_*j4zrB8k!Dze#M--UnjJA>%S&-FZ!MD4e7002UKEMEWs delta 95 zcmZpfz}PZ@aRa}CASZ(!LjglBLp~6fG88jpPG*y6;xIHZ(@`)qG@U#_!d@6d2}AN` iW Router { // Tracks .route("/tracks/:id", axum::routing::get(get_track)) .route("/tracks/:id/stream", axum::routing::get(get_stream_url)) + .route("/tracks/:id/flac", axum::routing::get(get_flac_stream)) // Artists .route("/artists/:id/albums", axum::routing::get(get_artist_albums)) .route( @@ -154,6 +155,35 @@ async fn get_stream_url( Ok(Json(serde_json::json!({ "url": url }))) } +/// `GET /tracks/:id/flac` — flux FLAC déchiffré en streaming progressif via CMAF. +/// +/// Retourne les bytes FLAC directement dans le corps de la réponse HTTP, au fur +/// et à mesure que les segments CMAF sont téléchargés et déchiffrés. Le client +/// (typiquement `pmocache::add_from_url`) peut commencer à consommer le contenu +/// avant la fin du téléchargement, préservant le progressive caching. +/// +/// `Content-Length` est une estimation calculée depuis la table des segments du +/// segment init — la valeur réelle peut différer de quelques octets. +#[cfg(feature = "pmoserver")] +async fn get_flac_stream( + State(state): State, + Path(id): Path, +) -> Result { + use axum::http::header; + use tokio_util::io::ReaderStream; + + let (reader, estimated_size) = state.client.open_cmaf_stream(&id).await?; + let stream = ReaderStream::new(reader); + let body = axum::body::Body::from_stream(stream); + + axum::response::Response::builder() + .status(axum::http::StatusCode::OK) + .header(header::CONTENT_TYPE, "audio/flac") + .header(header::CONTENT_LENGTH, estimated_size) + .body(body) + .map_err(|e| AppError(QobuzError::Other(e.to_string()))) +} + #[cfg(feature = "pmoserver")] async fn get_artist_albums( State(state): State, diff --git a/pmoqobuz/src/client.rs b/pmoqobuz/src/client.rs index fb97c991..3666516e 100644 --- a/pmoqobuz/src/client.rs +++ b/pmoqobuz/src/client.rs @@ -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. /// /// 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> { let format_id = self.api.format().id() as u32; diff --git a/pmoqobuz/src/cmaf/mod.rs b/pmoqobuz/src/cmaf/mod.rs index 1fff9e0f..fd6963b1 100644 --- a/pmoqobuz/src/cmaf/mod.rs +++ b/pmoqobuz/src/cmaf/mod.rs @@ -25,6 +25,7 @@ pub use parser::{ }; use std::sync::Arc; +use tokio::io::AsyncRead; use tokio::sync::Semaphore; use tracing::{debug, info, warn}; @@ -159,35 +160,45 @@ async fn fetch_all_segments( 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. +fn decrypt_one_segment( + seg_data: &[u8], + content_key: &[u8; 16], + output: &mut Vec, + seg_idx: usize, +) -> Result<()> { + let crypto = parse_segment_crypto(seg_data) + .map_err(|e| QobuzError::Other(format!("CMAF seg {} parse: {}", seg_idx, e)))?; + + let mut data_pos = crypto.data_offset; + for entry in &crypto.entries { + let frame_end = data_pos + entry.size as usize; + if frame_end > seg_data.len() { + return Err(QobuzError::Other(format!("CMAF seg {} débordement frame", seg_idx))); + } + let start = output.len(); + output.extend_from_slice(&seg_data[data_pos..frame_end]); + if entry.flags != 0 { + decrypt_frame(content_key, &entry.iv, &mut output[start..]); + } + data_pos = frame_end; + } + if data_pos < crypto.mdat_end && crypto.mdat_end <= seg_data.len() { + output.extend_from_slice(&seg_data[data_pos..crypto.mdat_end]); + } + Ok(()) +} + +/// Déchiffre une séquence de segments CMAF chiffrés et écrit les frames FLAC dans `output`. pub fn decrypt_segments_into( segments: &[Vec], content_key: &[u8; 16], output: &mut Vec, ) -> Result<()> { for (seg_idx, seg_data) in segments.iter().enumerate() { - let log_idx = seg_idx + 1; - let crypto = parse_segment_crypto(seg_data) - .map_err(|e| QobuzError::Other(format!("CMAF seg {} parse: {}", log_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", log_idx))); - } - let output_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[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]); - } + decrypt_one_segment(seg_data, content_key, output, seg_idx + 1)?; } 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, + bit_depth: Option, +) -> Result<(impl AsyncRead + Send + Unpin + 'static, u64)> { + use tokio::io::AsyncWriteExt; + + let setup = setup_streaming( + url_template, &key_str, &infos, n_segments, + format_id, sampling_rate, bit_depth, + ).await?; + + let estimated_size = (setup.flac_header.len() + + setup.segment_table.iter().map(|s| s.byte_len as usize).sum::()) as u64; + + // Pipe 256 KB : assez grand pour absorber un segment FLAC Hi-Res typique + // sans bloquer le producteur, assez petit pour ne pas sur-allouer. + let (mut writer, reader) = tokio::io::duplex(256 * 1024); + + tokio::spawn(async move { + if let Err(e) = writer.write_all(&setup.flac_header).await { + warn!("[CMAF-STREAM] écriture header FLAC: {}", e); + return; + } + + let http = match build_cdn_client() { + Ok(c) => c, + Err(e) => { warn!("[CMAF-STREAM] client HTTP: {}", e); return; } + }; + + // Lancer tous les fetches avec concurrence bornée par le sémaphore. + // Les handles sont stockés dans l'ordre — on les consomme en ordre. + let sem = Arc::new(Semaphore::new(CMAF_PREFETCH_CONCURRENCY)); + let mut handles = Vec::with_capacity(setup.n_segments as usize); + + for seg_idx in 1u8..=setup.n_segments { + let sem = sem.clone(); + let http = http.clone(); + let url = setup.url_template.replace("$SEGMENT$", &seg_idx.to_string()); + + handles.push(tokio::spawn(async move { + let permit = sem.acquire_owned().await + .map_err(|e| format!("sémaphore: {}", e))?; + let result = fetch_bytes_with_retry(&http, &url, &format!("CMAF-STREAM seg {}", seg_idx)) + .await + .map_err(|e| format!("seg {} fetch: {}", seg_idx, e)); + // Cooldown CDN Akamai : retenir le slot 500ms avant de libérer, + // identique à fetch_all_segments, pour éviter le rate-limiting. + tokio::time::sleep(std::time::Duration::from_millis(500)).await; + drop(permit); + result + })); + } + + // Consommer dans l'ordre : segment 1 d'abord, puis 2, etc. + // Les segments prêts en avance attendent dans leur JoinHandle. + let mut buf = Vec::new(); + for (i, handle) in handles.into_iter().enumerate() { + let seg_idx = i + 1; + let seg_data = match handle.await { + Ok(Ok(data)) => data, + Ok(Err(e)) => { warn!("[CMAF-STREAM] seg {}: {}", seg_idx, e); return; } + Err(e) => { warn!("[CMAF-STREAM] seg {} panique: {}", seg_idx, e); return; } + }; + + buf.clear(); + if let Err(e) = decrypt_one_segment(&seg_data, &setup.content_key, &mut buf, seg_idx) { + warn!("[CMAF-STREAM] seg {}: déchiffrement: {}", seg_idx, e); + return; + } + if let Err(e) = writer.write_all(&buf).await { + // Le lecteur a fermé le pipe (ex: playback stoppé) — arrêt silencieux. + debug!("[CMAF-STREAM] seg {}: pipe fermé ({})", seg_idx, e); + return; + } + debug!("[CMAF-STREAM] seg {}/{} → {} B", seg_idx, setup.n_segments, buf.len()); + } + + info!("[CMAF-STREAM] complet : {} segments, {:.2} MB estimés", + setup.n_segments, + estimated_size as f64 / (1024.0 * 1024.0)); + // writer dropped ici → EOF propre sur le reader + }); + + Ok((reader, estimated_size)) +} + /// Télécharge un track CMAF complet et retourne les bytes FLAC déchiffrés. pub async fn download_full( url_template: String, diff --git a/pmoqobuz/src/lazy_provider.rs b/pmoqobuz/src/lazy_provider.rs index cb83a36e..8615a9f7 100644 --- a/pmoqobuz/src/lazy_provider.rs +++ b/pmoqobuz/src/lazy_provider.rs @@ -68,6 +68,18 @@ impl LazyProvider for QobuzLazyProvider { async fn get_url(&self, lazy_pk: &str) -> Result { let track_id = self.track_id_from_lazy(lazy_pk)?; + + // L'endpoint CMAF local n'existe que si la feature `pmoserver` est active + // (route /qobuz/tracks/:id/flac enregistrée). Sans cette feature, fallback + // sur l'ancienne API Qobuz directe. + #[cfg(feature = "pmoserver")] + { + let base = std::env::var("PMO_SERVER_URL") + .unwrap_or_else(|_| "http://localhost:8080".to_string()); + return Ok(format!("{}/qobuz/tracks/{}/flac", base.trim_end_matches('/'), track_id)); + } + + #[cfg(not(feature = "pmoserver"))] self.client .get_stream_url(track_id) .await diff --git a/pmoqobuz/src/pmoserver_impl.rs b/pmoqobuz/src/pmoserver_impl.rs index 17961002..224f713d 100644 --- a/pmoqobuz/src/pmoserver_impl.rs +++ b/pmoqobuz/src/pmoserver_impl.rs @@ -29,6 +29,7 @@ use crate::api_rest::{create_router, QobuzState}; use crate::client::QobuzClient; +use crate::config_ext::QobuzConfigExt; use crate::pmoserver_ext::QobuzServerExt; use anyhow::Result; use pmoconfig::Config;