diff --git a/pmoaudio-ext/src/sinks/direct_ogg_flac_sink.rs b/pmoaudio-ext/src/sinks/direct_ogg_flac_sink.rs index 4f94e248..343567f4 100644 --- a/pmoaudio-ext/src/sinks/direct_ogg_flac_sink.rs +++ b/pmoaudio-ext/src/sinks/direct_ogg_flac_sink.rs @@ -5,15 +5,17 @@ //! //! # Cycle de vie //! -//! - **Play** : le navigateur appelle `GET /stream`. `connect()` crée un nouveau -//! canal PCM + pipe duplex + encodeur FLAC + wrapper OGG, installe le sender -//! dans le sink, et notifie le sink via `client_notify`. Le flux reste ouvert : -//! les morceaux s'enchaînent en gapless. -//! - **Stop** : le navigateur ferme la connexion. Le pipe se rompt, l'encodeur -//! s'arrête. Le sink voit `pcm_tx.send()` échouer, passe le sender à `None`, -//! et **bloque** sur `client_notify` jusqu'au prochain Play. -//! - **Play suivant** : `connect()` → nouveau pipe → `client_notify.notify_one()` -//! → le sink se débloque et reprend la consommation des segments. +//! - **Play** : le navigateur appelle `GET /stream`. `connect()` crée un channel +//! Bytes (bytes_tx/rx), une task de forwarding qui copie les Bytes dans un +//! DuplexStream, et lance le premier encodeur OGG-FLAC. Le `bytes_tx` est +//! stocké dans le sink pour le chaining TrackBoundary. +//! - **Stop** : le navigateur ferme la connexion. Le DuplexStream se rompt, +//! la task de forwarding se termine, le bytes_tx devient invalide. Le sink +//! voit `pcm_tx.send()` échouer et bloque sur `client_notify`. +//! - **TrackBoundary** : le sink ferme le `pcm_tx` courant (EOF → encodeur écrit +//! EOS OGG), attend la fin de l'encodeur, puis relance un nouvel encodeur +//! dans le même `bytes_tx` (OGG chaining : nouvelle BOS OGG dans le même flux HTTP). +//! - **Play suivant** : `connect()` → nouveau DuplexStream + channel → nouveau pipe. //! //! # Architecture //! @@ -21,11 +23,10 @@ //! AudioSegment I24 @ 96 kHz //! ↓ NodeLogic::process() [bloque si pas de client] //! chunk_to_pcm_bytes() → PCM 24-bit LE -//! ↓ Arc>>> -//! ByteStreamReader (AsyncRead) -//! ↓ encode_flac_stream() -//! ↓ broadcast_ogg_flac_stream() → wrapping OGG pages -//! ↓ tokio::io::duplex pipe (256 KB) +//! ↓ SharedPcmTx +//! ByteStreamReader → encode_flac_stream() → OGG pages +//! ↓ SharedBytesTx (mpsc::Sender) ← persistant entre encodeurs +//! [forwarding task] → tokio::io::DuplexStream //! ↓ DirectOggFlacStream (AsyncRead) → Body HTTP //! ``` @@ -43,6 +44,7 @@ use pmoaudio::{ use pmoflac::{EncoderOptions, PcmFormat}; use tokio::io::{AsyncRead, ReadBuf}; use tokio::sync::{mpsc, watch, Mutex}; +use tokio::task::JoinHandle; use tokio_util::sync::CancellationToken; use tracing::{debug, warn}; @@ -57,10 +59,17 @@ pub const DIRECT_OGG_FLAC_BITS_PER_SAMPLE: u8 = 24; /// Capacité du pipe duplex (~256 KB). const PIPE_CAPACITY: usize = 256 * 1024; +/// Capacité du channel Bytes intermédiaire. +const BYTES_CHANNEL_CAPACITY: usize = 64; // ─── Shared state ───────────────────────────────────────────────────────────── type SharedPcmTx = Arc>>>; +/// Canal Bytes persistant entre les encodeurs successifs (OGG chaining). +/// Le sink y envoie les pages OGG ; une task de forwarding les copie dans le DuplexStream. +type SharedBytesTx = Arc>>>; +/// Handle de la task encodeur courante. +type SharedEncoderTask = Arc>>>; // ─── Handle public ──────────────────────────────────────────────────────────── @@ -72,48 +81,76 @@ pub struct DirectOggFlacHandle { client_notify_internal: Arc, first_byte_tx: Arc>, encoder_options: EncoderOptions, - /// Position de lecture courante (mise à jour par ByteStreamReader). current_timestamp: Arc>, + /// Canal Bytes persistant partagé avec la logic du sink pour le OGG chaining. + bytes_tx: SharedBytesTx, + /// Task encodeur courante partagée avec la logic du sink. + encoder_task: SharedEncoderTask, } impl DirectOggFlacHandle { - /// Crée un nouveau pipe OGG-FLAC et retourne le flux côté lecture. - /// Débloque le sink s'il attendait un client. + /// Crée un nouveau flux OGG-FLAC pour le client HTTP. + /// Remplace toute connexion précédente. pub async fn connect(&self) -> DirectOggFlacStream { let connect_count_before = *self.client_connect_tx.borrow(); debug!("DirectOggFlacHandle::connect() called, connect_count={}", connect_count_before); - let (pcm_tx, pcm_rx) = mpsc::channel::(8); - // Réinitialiser le timestamp à 0 pour la nouvelle connexion + // Annuler l'encodeur précédent s'il tourne encore + if let Some(old_task) = self.encoder_task.lock().await.take() { + old_task.abort(); + } + + // Créer le channel Bytes persistant (OGG chaining) + let (bytes_tx, bytes_rx) = mpsc::channel::(BYTES_CHANNEL_CAPACITY); + *self.bytes_tx.lock().await = Some(bytes_tx.clone()); + + // Créer le DuplexStream vers le client HTTP + let (mut pipe_writer, pipe_reader) = tokio::io::duplex(PIPE_CAPACITY); + + // Task de forwarding : Bytes → DuplexStream + // Se termine quand bytes_rx est fermé (bytes_tx droppé) ou pipe cassé + tokio::spawn(async move { + use tokio::io::AsyncWriteExt; + let mut rx = bytes_rx; + while let Some(bytes) = rx.recv().await { + if pipe_writer.write_all(&bytes).await.is_err() { + debug!("DirectOggFlacStream forwarder: pipe broken, client disconnected"); + break; + } + } + debug!("DirectOggFlacStream forwarder: done"); + }); + + // Réinitialiser les signaux + let _ = self.first_byte_tx.send(false); *self.current_timestamp.write().await = 0.0; + + // Créer le premier pcm_tx + encodeur + let (pcm_tx, pcm_rx) = mpsc::channel::(8); let current_dur = Arc::new(tokio::sync::RwLock::new(0.0f64)); - // Partager current_timestamp avec ByteStreamReader : il sera mis à jour - // avec le timestamp absolu du segment audio (position dans le fichier source). let pcm_reader = ByteStreamReader::new(pcm_rx, self.current_timestamp.clone(), current_dur); - let (pipe_writer, pipe_reader) = tokio::io::duplex(PIPE_CAPACITY); - - let _ = self.first_byte_tx.send(false); - debug!("DirectOggFlacHandle::connect() first_byte reset to false"); - *self.pcm_tx.lock().await = Some(pcm_tx); - debug!("DirectOggFlacHandle::connect() pcm_tx installed"); let new_count = connect_count_before.wrapping_add(1); let _ = self.client_connect_tx.send(new_count); - debug!("DirectOggFlacHandle::connect() client_connect_count -> {}", new_count); self.client_notify_internal.notify_one(); + debug!("DirectOggFlacHandle::connect() client_connect_count -> {}", new_count); let options = self.encoder_options.clone(); let current_timestamp = self.current_timestamp.clone(); - tokio::spawn(async move { - debug!("DirectOggFlacHandle: encoder+ogg task started"); - if let Err(e) = run_ogg_encoder(pcm_reader, pipe_writer, options, current_timestamp).await { - debug!("DirectOggFlacStream encoder stopped: {}", e); + let shared_bytes_tx = self.bytes_tx.clone(); + + let handle = tokio::spawn(async move { + debug!("DirectOggFlacHandle: initial encoder task started"); + if let Err(e) = run_ogg_encoder(pcm_reader, bytes_tx, shared_bytes_tx, options, current_timestamp).await { + debug!("DirectOggFlacHandle: initial encoder stopped: {}", e); } - debug!("DirectOggFlacHandle: encoder+ogg task ended"); + debug!("DirectOggFlacHandle: initial encoder task ended"); }); + *self.encoder_task.lock().await = Some(handle); + debug!("DirectOggFlacHandle::connect() returning DirectOggFlacStream"); DirectOggFlacStream { inner: pipe_reader, @@ -125,7 +162,6 @@ impl DirectOggFlacHandle { self.first_byte_tx.subscribe() } - /// Retourne la position de lecture courante en secondes. pub async fn current_position_sec(&self) -> f64 { *self.current_timestamp.read().await } @@ -172,6 +208,13 @@ impl AsyncRead for DirectOggFlacStream { struct DirectOggFlacSinkLogic { pcm_tx: SharedPcmTx, client_notify: Arc, + encoder_options: EncoderOptions, + encoder_task: SharedEncoderTask, + bytes_tx: SharedBytesTx, + current_timestamp: Arc>, + /// Vrai dès qu'au moins un chunk audio a été encodé dans le stream courant. + /// Empêche le OGG chaining sur le TrackBoundary initial (avant tout audio). + has_encoded_frames: bool, } #[async_trait] @@ -234,9 +277,21 @@ impl NodeLogic for DirectOggFlacSinkLogic { seg.timestamp_sec, ); *self.pcm_tx.lock().await = None; + self.has_encoded_frames = false; + } else { + self.has_encoded_frames = true; } } _AudioSegment::Sync(marker) => match marker.as_ref() { + SyncMarker::TrackBoundary { .. } => { + if self.has_encoded_frames { + debug!("DirectOggFlacSink: TrackBoundary — OGG chaining"); + self.has_encoded_frames = false; + self.do_track_boundary().await; + } else { + debug!("DirectOggFlacSink: TrackBoundary ignored (no frames encoded yet)"); + } + } SyncMarker::EndOfStream => { debug!("DirectOggFlacSink: EndOfStream"); } @@ -256,11 +311,69 @@ impl NodeLogic for DirectOggFlacSinkLogic { } } +impl DirectOggFlacSinkLogic { + /// OGG chaining : ferme l'encodeur courant (EOF → EOS OGG), attend sa fin, + /// puis relance un nouvel encodeur dans le même channel Bytes (nouvelle BOS OGG). + async fn do_track_boundary(&mut self) { + // 1. Fermer le pcm_tx courant → EOF dans ByteStreamReader → encodeur écrit EOS OGG + { + let mut guard = self.pcm_tx.lock().await; + *guard = None; + } + + // 2. Attendre la fin de la task encodeur courante + let old_task = self.encoder_task.lock().await.take(); + if let Some(handle) = old_task { + let _ = handle.await; + debug!("DirectOggFlacSink: previous encoder task joined"); + } + + // 3. Vérifier que le bytes_tx est encore valide (client pas déconnecté) + let bytes_tx = { + let guard = self.bytes_tx.lock().await; + guard.clone() + }; + let Some(bytes_tx) = bytes_tx else { + debug!("DirectOggFlacSink: bytes_tx gone (client disconnected), skip OGG chaining"); + return; + }; + + // 4. Nouveau pcm_tx + ByteStreamReader + let (pcm_tx, pcm_rx) = mpsc::channel::(8); + let current_dur = Arc::new(tokio::sync::RwLock::new(0.0f64)); + let pcm_reader = ByteStreamReader::new( + pcm_rx, + self.current_timestamp.clone(), + current_dur, + ); + *self.pcm_tx.lock().await = Some(pcm_tx); + + // 5. Relancer l'encodeur dans le même bytes_tx (OGG chaining : nouvelle BOS OGG) + let options = self.encoder_options.clone(); + let current_timestamp = self.current_timestamp.clone(); + let shared_bytes_tx = self.bytes_tx.clone(); + + let handle = tokio::spawn(async move { + debug!("DirectOggFlacSink: chained encoder task started"); + if let Err(e) = run_ogg_encoder(pcm_reader, bytes_tx, shared_bytes_tx, options, current_timestamp).await { + debug!("DirectOggFlacSink: chained encoder stopped: {}", e); + } + debug!("DirectOggFlacSink: chained encoder task ended"); + }); + + *self.encoder_task.lock().await = Some(handle); + debug!("DirectOggFlacSink: OGG chaining complete, new encoder started"); + } +} + // ─── Encodeur FLAC + wrapper OGG ───────────────────────────────────────────── +/// Encode PCM → FLAC → OGG et envoie les pages OGG dans `bytes_tx`. +/// Quand le channel devient invalide (client déconnecté), nettoie `shared_bytes_tx`. async fn run_ogg_encoder( pcm_reader: ByteStreamReader, - mut pipe_writer: tokio::io::DuplexStream, + bytes_tx: mpsc::Sender, + shared_bytes_tx: SharedBytesTx, options: EncoderOptions, _current_timestamp: Arc>, ) -> Result<(), AudioError> { @@ -274,56 +387,56 @@ async fn run_ogg_encoder( .await .map_err(|e| AudioError::ProcessingError(format!("FLAC encoder init: {}", e)))?; - // Lire le header FLAC et construire les pages OGG d'en-tête let flac_header = read_flac_header(&mut flac_stream).await?; - let sample_rate = extract_sample_rate_from_streaminfo(&flac_header)?; + let _sample_rate = extract_sample_rate_from_streaminfo(&flac_header)?; let stream_serial: u32 = rand::random(); let mut ogg = OggPageWriter::new(stream_serial); - // Page BOS (identification OGG-FLAC) let ogg_flac_id = create_ogg_flac_identification(&flac_header)?; let bos_page = Bytes::from(ogg.create_page(&ogg_flac_id, true, false, false)); - - // Page Vorbis Comment let vorbis_comment = create_empty_vorbis_comment(); let comment_page = Bytes::from(ogg.create_page(&vorbis_comment, false, false, false)); - pipe_writer.write_all(&bos_page).await - .map_err(|e| AudioError::IoError(format!("OGG BOS write: {}", e)))?; - pipe_writer.write_all(&comment_page).await - .map_err(|e| AudioError::IoError(format!("OGG comment write: {}", e)))?; + macro_rules! send_or_cleanup { + ($page:expr) => { + if bytes_tx.send($page).await.is_err() { + debug!("DirectOggFlacSink: bytes_tx broken, client disconnected"); + *shared_bytes_tx.lock().await = None; + return Ok(()); + } + }; + } + + send_or_cleanup!(bos_page); + send_or_cleanup!(comment_page); + + use tokio::io::AsyncReadExt; + use crate::sinks::flac_frame_utils::{validate_frame_header_crc, parse_flac_block_size}; - // Lire les frames FLAC et les encapsuler dans des pages OGG - let sample_rate_f64 = sample_rate as f64; let mut encoded_samples = 0u64; let mut read_buffer = vec![0u8; 16384]; let mut accumulator: Vec = Vec::with_capacity(32768); - use tokio::io::{AsyncReadExt, AsyncWriteExt}; loop { match flac_stream.read(&mut read_buffer).await { Ok(0) => { // EOF : page EOS finale let eos_page = Bytes::from(ogg.create_page(&accumulator, false, true, false)); - let _ = pipe_writer.write_all(&eos_page).await; + let _ = bytes_tx.send(eos_page).await; break; } Ok(n) => { accumulator.extend_from_slice(&read_buffer[..n]); loop { - if accumulator.len() < 4 { - break; - } + if accumulator.len() < 4 { break; } - // Trouver les positions de sync FLAC let mut sync_data: Vec<(usize, u32)> = Vec::new(); for i in 0..accumulator.len() - 1 { let b1 = accumulator[i]; let b2 = accumulator[i + 1]; if b1 == 0xFF && b2 >= 0xF8 && b2 <= 0xFE { - use crate::sinks::flac_frame_utils::{validate_frame_header_crc, parse_flac_block_size}; if validate_frame_header_crc(&accumulator, i) { if let Some(samples) = parse_flac_block_size(&accumulator, i) { sync_data.push((i, samples)); @@ -332,9 +445,7 @@ async fn run_ogg_encoder( } } - if sync_data.len() < 2 { - break; - } + if sync_data.len() < 2 { break; } let first_start = sync_data[0].0; let first_samples = sync_data[0].1; @@ -346,19 +457,17 @@ async fn run_ogg_encoder( } let frame: Vec = accumulator.drain(0..second_start).collect(); - encoded_samples = encoded_samples.saturating_add(first_samples as u64); ogg.add_samples(first_samples as u64); let ogg_page = Bytes::from(ogg.create_page(&frame, false, false, false)); - if pipe_writer.write_all(&ogg_page).await.is_err() { - // Client déconnecté — le pipe HTTP s'est rompu + if bytes_tx.send(ogg_page).await.is_err() { warn!( - samples = encoded_samples, - "DirectOggFlacSink: OGG pipe broken after {} samples ({:.3}s), client disconnected", + "DirectOggFlacSink: bytes_tx broken after {} samples ({:.3}s), client disconnected", encoded_samples, encoded_samples as f64 / DIRECT_OGG_FLAC_SAMPLE_RATE as f64, ); + *shared_bytes_tx.lock().await = None; return Ok(()); } } @@ -375,7 +484,7 @@ async fn run_ogg_encoder( Ok(()) } -// ─── OGG helpers (copiés de streaming_ogg_flac_sink) ───────────────────────── +// ─── OGG helpers ───────────────────────────────────────────────────────────── struct OggPageWriter { stream_serial: u32, @@ -521,10 +630,17 @@ impl DirectOggFlacSink { let (first_byte_tx, _) = watch::channel(false); let first_byte_tx = Arc::new(first_byte_tx); let current_timestamp = Arc::new(tokio::sync::RwLock::new(0.0f64)); + let bytes_tx: SharedBytesTx = Arc::new(Mutex::new(None)); + let encoder_task: SharedEncoderTask = Arc::new(Mutex::new(None)); let logic = DirectOggFlacSinkLogic { pcm_tx: pcm_tx.clone(), client_notify: client_notify_internal.clone(), + encoder_options: encoder_options.clone(), + encoder_task: encoder_task.clone(), + bytes_tx: bytes_tx.clone(), + current_timestamp: current_timestamp.clone(), + has_encoded_frames: false, }; let sink = Self { @@ -538,6 +654,8 @@ impl DirectOggFlacSink { first_byte_tx, encoder_options, current_timestamp, + bytes_tx, + encoder_task, }; (sink, handle) diff --git a/pmowebrenderer/src/pipeline.rs b/pmowebrenderer/src/pipeline.rs index 4f8548d0..faa178a0 100644 --- a/pmowebrenderer/src/pipeline.rs +++ b/pmowebrenderer/src/pipeline.rs @@ -10,8 +10,8 @@ use std::sync::Arc; use pmoaudio::{AudioSegment, PositionTrackerNode, ResamplingNode, ToI24Node}; use pmometadata::{MemoryTrackMetadata, TrackMetadata}; use pmoaudio_ext::sinks::{ - OggFlacStreamHandle, StreamingOggFlacSink, - DIRECT_OGG_FLAC_BITS_PER_SAMPLE, DIRECT_OGG_FLAC_SAMPLE_RATE, + DirectOggFlacHandle, DirectOggFlacSink, + DIRECT_OGG_FLAC_SAMPLE_RATE, }; use pmoaudio_ext::UriSource; use pmoflac::EncoderOptions; @@ -59,10 +59,10 @@ impl PipelineHandle { /// Pipeline audio complet pour une instance WebRenderer. /// /// Créé au `POST /register`. Le flux OGG-FLAC est accessible via `flac_handle` -/// (multi-clients, chaque `subscribe()` retourne un nouveau flux depuis le point courant). +/// (mono-client, chaque `connect()` crée un nouveau flux avec backpressure TCP). pub struct InstancePipeline { - /// Handle vers le sink OGG-FLAC — clonable, chaque subscribe() donne un flux live. - pub flac_handle: OggFlacStreamHandle, + /// Handle vers le sink OGG-FLAC — clonable, connect() crée un nouveau flux. + pub flac_handle: DirectOggFlacHandle, pub pipeline_handle: PipelineHandle, } @@ -82,10 +82,7 @@ impl InstancePipeline { // La backpressure remonte naturellement depuis le pipe duplex jusqu'à la source. use pmoaudio::pipeline::AudioPipelineNode; - let (sink, flac_handle) = StreamingOggFlacSink::new( - EncoderOptions::default(), - DIRECT_OGG_FLAC_BITS_PER_SAMPLE, - ); + let (sink, flac_handle) = DirectOggFlacSink::new(EncoderOptions::default()); // Nœud de suivi de position : lit le timestamp des chunks sortant vers le sink let (mut position_tracker, position_handle) = PositionTrackerNode::new(); diff --git a/pmowebrenderer/src/registry.rs b/pmowebrenderer/src/registry.rs index b49e2c06..41cb9fb7 100644 --- a/pmowebrenderer/src/registry.rs +++ b/pmowebrenderer/src/registry.rs @@ -29,8 +29,8 @@ pub struct WebRendererInstance { pub udn: String, pub device_instance: Arc, pub state: SharedState, - /// Handle vers le sink OGG-FLAC — clonable, chaque subscribe() donne un flux live. - pub flac_handle: pmoaudio_ext::sinks::OggFlacStreamHandle, + /// Handle vers le sink OGG-FLAC — clonable, connect() crée un nouveau flux avec backpressure. + pub flac_handle: pmoaudio_ext::sinks::DirectOggFlacHandle, pub pipeline: PipelineHandle, pub created_at: SystemTime, } @@ -116,11 +116,11 @@ impl RendererRegistry { Ok((stream_url, udn)) } - /// Retourne le OggFlacStreamHandle pour l'endpoint /stream (clonable). + /// Retourne le DirectOggFlacHandle pour l'endpoint /stream (clonable). pub fn get_flac_handle( &self, instance_id: &str, - ) -> Option { + ) -> Option { self.instances .read() .get(instance_id) diff --git a/pmowebrenderer/src/stream.rs b/pmowebrenderer/src/stream.rs index fa28a19a..b147c986 100644 --- a/pmowebrenderer/src/stream.rs +++ b/pmowebrenderer/src/stream.rs @@ -44,7 +44,7 @@ pub async fn stream_handler( } }; - let stream = handle.subscribe(); + let stream = handle.connect().await; info!(instance_id = %instance_id, "FLAC stream started");