Fix streaming and cache progressive in play_and_cache example
Cette correction implémente le cache progressif et le streaming pour permettre un démarrage quasi immédiat de la lecture pendant le téléchargement. ## Changements dans FlacCacheSink (pmoaudio-ext) Avant : - Accumulait tout le FLAC en mémoire dans un buffer - Attendait la fin complète de l'encodage avant d'ajouter au cache - Ajoutait à la playlist seulement après ingestion complète Après : - Passe le flux FLAC directement à add_from_reader - add_from_reader retourne dès que le prebuffer (512 KB) est atteint - Le PK est ajouté à la playlist immédiatement après le prebuffer - L'encodage et l'écriture continuent en arrière-plan ## Changements dans play_and_cache.rs - Suppression du sleep de 2 secondes avant le démarrage de la lecture - Ajout de commentaire expliquant le mécanisme de prebuffer - La lecture démarre dès que le prebuffer est atteint (~1-2 secondes) ## Résultat La musique démarre maintenant presque immédiatement après le début du téléchargement (temps du prebuffer) au lieu d'attendre la fin du téléchargement complet du premier morceau.
This commit is contained in:
@@ -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);
|
||||
|
||||
@@ -233,9 +233,9 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
|
||||
});
|
||||
|
||||
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"),
|
||||
|
||||
Reference in New Issue
Block a user