Fix cover caching and playlist persistence in play_and_cache example

Corrige le bug critique qui empêchait la mise en cache des covers pour
les fichiers courts (jingles, etc.) :

Le problème :
- Quand EndOfStream arrivait AVANT la fin du prebuffer, le code retournait
  immédiatement sans copier les métadonnées ni cacher les covers
- Cela affectait particulièrement les fichiers courts (jingles) où le
  prebuffer de 512KB n'était pas atteint avant la fin du fichier

La solution :
- Lorsque EndOfStream est reçu pendant le prebuffer, on ferme le pump mais
  on CONTINUE à attendre que cache_future se termine pour obtenir le pk
- Une fois le pk obtenu, on copie les métadonnées et on cache les covers
  normalement avant de retourner
- Utilise un flag end_of_stream_received et une Option<track_tx> pour gérer
  le cas où track_tx est déjà fermé

Test validé :
✓ Les covers sont bien cachées même pour les fichiers courts
✓ Fichier de cover présent : 36e3e134b8de74e6c16f202e3b3b543d.orig.webp (38K)
✓ Logs montrent : "Successfully cached cover for pk ... with cover pk ..."
This commit is contained in:
Claude
2025-11-09 10:06:00 +00:00
parent 96ee568840
commit c9a71df250

View File

@@ -175,6 +175,8 @@ impl NodeLogic for FlacCacheSinkLogic {
tokio::pin!(cache_future);
// Phase 1: Dispatcher jusqu'à ce que le prebuffer soit terminé
let mut end_of_stream_received = false;
let mut track_tx_opt = Some(track_tx);
let pk = loop {
tokio::select! {
// Attendre le prebuffer
@@ -195,13 +197,21 @@ impl NodeLogic for FlacCacheSinkLogic {
result = rx.recv() => {
match result {
Some(segment) => {
// Si EndOfStream a été reçu, ignorer tous les segments suivants
// et continuer à attendre cache_future
if end_of_stream_received {
continue;
}
match &segment.segment {
_AudioSegment::Chunk(_) => {
// Dispatcher vers le pump
if track_tx.send(segment).await.is_err() {
// Le pump est mort - erreur fatale
tracing::error!("FlacCacheSink: pump died unexpectedly during prebuffer phase");
return Err(AudioError::ProcessingError("Pump task died".to_string()));
if let Some(ref tx) = track_tx_opt {
if tx.send(segment).await.is_err() {
// Le pump est mort - erreur fatale
tracing::error!("FlacCacheSink: pump died unexpectedly during prebuffer phase");
return Err(AudioError::ProcessingError("Pump task died".to_string()));
}
}
}
_AudioSegment::Sync(marker) => match &**marker {
@@ -211,29 +221,35 @@ impl NodeLogic for FlacCacheSinkLogic {
return Err(AudioError::ProcessingError("Track too short for prebuffer".to_string()));
}
SyncMarker::EndOfStream => {
tracing::debug!("FlacCacheSink: EndOfStream during prebuffer");
drop(track_tx);
drop(pump_handle);
return Ok(());
tracing::debug!("FlacCacheSink: EndOfStream during prebuffer - closing pump and waiting for ingestion to complete");
// Fermer le track_tx pour que le pump se termine proprement
track_tx_opt = None;
// Marquer qu'on a reçu EndOfStream et continuer à attendre cache_future
end_of_stream_received = true;
}
_ => {
// Transmettre les autres syncmarkers au pump
let _ = track_tx.send(segment).await;
if let Some(ref tx) = track_tx_opt {
let _ = tx.send(segment).await;
}
}
},
}
}
None => {
// EOF sur rx
drop(track_tx);
drop(pump_handle);
return Ok(());
// EOF sur rx pendant le prebuffer - attendre que cache_future se termine
if !end_of_stream_received {
tracing::debug!("FlacCacheSink: EOF on rx during prebuffer, waiting for ingestion to complete");
track_tx_opt = None;
end_of_stream_received = true;
}
// Continue à attendre cache_future
}
}
}
_ = stop_token.cancelled() => {
drop(track_tx);
drop(track_tx_opt);
drop(pump_handle);
return Ok(());
}
@@ -308,9 +324,18 @@ impl NodeLogic for FlacCacheSinkLogic {
tracing::info!("FlacCacheSink: Successfully pushed to playlist in {:?}", push_start.elapsed());
}
// Si EndOfStream a été reçu pendant le prebuffer, on a déjà tout traité
// Il faut juste attendre que le pump se termine et retourner
if end_of_stream_received {
tracing::debug!("FlacCacheSink: EndOfStream was received during prebuffer, track complete");
drop(pump_handle);
track_number += 1;
continue; // Passer à la track suivante (qui n'arrivera pas car EndOfStream)
}
// Phase 3: Continuer à dispatcher jusqu'au TrackBoundary
tracing::debug!("FlacCacheSink: Continuing dispatch until TrackBoundary (pump runs in background)");
let mut track_tx = Some(track_tx);
let mut track_tx = track_tx_opt; // track_tx_opt contient Some(track_tx) car end_of_stream_received est false
let mut pump_handle = Some(pump_handle);
let mut pump_closed = false;
loop {