diff --git a/pmoaudio-ext/src/sinks/flac_cache_sink.rs b/pmoaudio-ext/src/sinks/flac_cache_sink.rs index f1924a9c..ec3877d2 100755 --- a/pmoaudio-ext/src/sinks/flac_cache_sink.rs +++ b/pmoaudio-ext/src/sinks/flac_cache_sink.rs @@ -127,50 +127,30 @@ impl NodeLogic for FlacCacheSinkLogic { // Créer l'encoder let reader = ByteStreamReader::new(pcm_rx); - let mut flac_stream = encode_flac_stream(reader, format, options_with_metadata) + let flac_stream = encode_flac_stream(reader, format, options_with_metadata) .await .map_err(|e| { AudioError::ProcessingError(format!("FLAC encode init failed: {}", e)) })?; - // Créer un buffer pour collecter le FLAC encodé - let mut flac_buffer = Vec::new(); - - // Exécuter pump et copy en parallèle - let pump_future = pump_track_segments( + // Lancer pump_track_segments en parallèle + let pump_handle = tokio::spawn(pump_track_segments( first_segment, &mut rx, pcm_tx, bits_per_sample, sample_rate, - &stop_token, - ); - let copy_future = async { - tokio::io::copy(&mut flac_stream, &mut flac_buffer) - .await - .map_err(|e| { - AudioError::ProcessingError(format!("FLAC write failed: {}", e)) - })?; - flac_stream - .wait() - .await - .map_err(|e| AudioError::ProcessingError(format!("Encoder failed: {}", e)))?; - Ok::<_, AudioError>(()) - }; + stop_token.clone(), + )); - // Attendre les deux tâches en parallèle - let (copy_result, pump_result) = tokio::join!(copy_future, pump_future); - copy_result?; - let (_chunks, _samples, _duration_sec, stop_reason) = pump_result?; - - // Ingérer le FLAC dans le cache - let flac_reader = Cursor::new(flac_buffer.clone()); + // Ingérer le FLAC progressivement dans le cache + // add_from_reader retourne dès que le prebuffer (512 KB) est atteint let collection_ref = self.collection.as_deref(); let pk = self.cache .add_from_reader( None, - flac_reader, - Some(flac_buffer.len() as u64), + flac_stream, + None, // Taille inconnue car streaming collection_ref, ) .await @@ -178,6 +158,13 @@ impl NodeLogic for FlacCacheSinkLogic { AudioError::ProcessingError(format!("Failed to add to cache: {}", e)) })?; + tracing::debug!("Track added to cache with pk {}, prebuffer complete", pk); + + // Attendre la fin du pump + let pump_result = pump_handle.await + .map_err(|e| AudioError::ProcessingError(format!("Pump task failed: {}", e)))?; + let (_chunks, _samples, _duration_sec, stop_reason) = pump_result?; + // Copier les métadonnées du TrackBoundary dans le cache if let Some(src_metadata) = track_metadata { let dest_metadata = self.cache.track_metadata(&pk); diff --git a/pmoparadise/examples/play_and_cache.rs b/pmoparadise/examples/play_and_cache.rs index 57354fa6..e971c843 100644 --- a/pmoparadise/examples/play_and_cache.rs +++ b/pmoparadise/examples/play_and_cache.rs @@ -233,9 +233,9 @@ async fn main() -> Result<(), Box> { }); let playback_handle = tokio::spawn(async move { - // Attendre un peu que le premier track soit disponible - tokio::time::sleep(tokio::time::Duration::from_secs(2)).await; - tracing::info!("[PLAYBACK] Pipeline starting..."); + // Pas de sleep - le cache progressif permet de démarrer immédiatement + // dès que le prebuffer (512 KB) est atteint + tracing::info!("[PLAYBACK] Pipeline starting (will wait for prebuffer)..."); let result = Box::new(playlist_source).run(stop_token_playback).await; match &result { Ok(()) => tracing::info!("[PLAYBACK] Pipeline completed successfully"),