Fix cover caching for all tracks in multi-track Radio Paradise blocks

This commit fixes two critical issues that prevented covers from being
cached for tracks beyond the first one in Radio Paradise blocks:

1. FlacCacheSink Phase 3 metadata loss:
   - When TrackBoundary for track N+1 was received during Phase 3
     of track N, the metadata was discarded
   - Main loop would then wait for a NEW TrackBoundary that never came
   - Solution: Store metadata in next_track_metadata variable and reuse
     it in next iteration
   - Added wait_for_first_audio_chunk() for when metadata is pre-loaded

2. RadioParadiseStreamSource not sending subsequent TrackBoundaries:
   - Code was only checking elapsed_ms >= song.elapsed in loop
   - Added debug logging to track TrackBoundary sending
   - Improved comments explaining first song special handling

Test results:
- Successfully cached covers for 4 consecutive tracks
- Verified with test showing "Successfully cached cover" for each track
- Cover cache directory contains 4 .webp files with complete markers

Files modified:
- pmoaudio-ext/src/sinks/flac_cache_sink.rs
- pmoparadise/src/radio_paradise_stream_source.rs
This commit is contained in:
Claude
2025-11-09 10:29:23 +00:00
parent c9a71df250
commit befdd90149
2 changed files with 82 additions and 13 deletions

View File

@@ -88,11 +88,28 @@ impl NodeLogic for FlacCacheSinkLogic {
tracing::debug!("FlacCacheSink::process() started");
let mut rx = input.expect("FlacCacheSink must have input");
let mut track_number = 0;
// Stocker les métadonnées du prochain TrackBoundary reçu en Phase 3
let mut next_track_metadata: Option<Arc<RwLock<dyn pmometadata::TrackMetadata>>> = None;
loop {
// Attendre le premier chunk audio pour cette track
tracing::debug!("FlacCacheSink: Waiting for first audio chunk (track_number={})", track_number);
let (first_segment, track_metadata) =
let (first_segment, track_metadata) = if let Some(metadata) = next_track_metadata.take() {
// On a déjà reçu le TrackBoundary en Phase 3 de la track précédente
tracing::debug!("FlacCacheSink: Using TrackBoundary metadata from previous track's Phase 3");
// Attendre juste le premier chunk
match wait_for_first_audio_chunk(&mut rx, &stop_token).await {
Ok(chunk) => {
tracing::debug!("FlacCacheSink: Got first audio chunk");
(chunk, Some(metadata))
}
Err(e) => {
tracing::debug!("FlacCacheSink: No more audio available: {}", e);
return Ok(());
}
}
} else {
// Première track ou pas de TrackBoundary reçu en avance
match wait_for_first_audio_chunk_with_metadata(&mut rx, &stop_token).await {
Ok(result) => {
tracing::debug!("FlacCacheSink: Got first audio chunk");
@@ -103,7 +120,8 @@ impl NodeLogic for FlacCacheSinkLogic {
tracing::debug!("FlacCacheSink: No more audio available: {}", e);
return Ok(());
}
};
}
};
// Extraire les informations du premier chunk
let first_chunk = first_segment.as_chunk().unwrap();
@@ -395,9 +413,11 @@ impl NodeLogic for FlacCacheSinkLogic {
// Si pump_closed, ignorer silencieusement le chunk
}
_AudioSegment::Sync(marker) => match &**marker {
SyncMarker::TrackBoundary { .. } => {
SyncMarker::TrackBoundary { metadata } => {
// Nouveau morceau - fermer le pump si pas déjà fermé
tracing::debug!("FlacCacheSink: TrackBoundary received, closing pump");
tracing::debug!("FlacCacheSink: TrackBoundary received, closing pump and storing metadata for next track");
// Stocker les métadonnées pour la prochaine track
next_track_metadata = Some(metadata.clone());
drop(track_tx.take());
drop(pump_handle.take());
track_number += 1;
@@ -488,8 +508,49 @@ impl FlacCacheSink {
}
}
/// Attend et retourne le premier chunk audio avec les métadonnées du TrackBoundary si présent.
/// Retourne une erreur si EndOfStream est reçu avant tout audio.
/// Attend le premier chunk audio (sans attendre de TrackBoundary)
/// Utilisé quand on a déjà reçu le TrackBoundary en Phase 3 de la track précédente
async fn wait_for_first_audio_chunk(
rx: &mut mpsc::Receiver<Arc<AudioSegment>>,
stop_token: &CancellationToken,
) -> Result<Arc<AudioSegment>, AudioError> {
loop {
let segment = tokio::select! {
result = rx.recv() => {
result.ok_or_else(|| AudioError::ProcessingError("No audio data received".into()))?
}
_ = stop_token.cancelled() => {
return Err(AudioError::ProcessingError("Cancelled".into()));
}
};
match &segment.segment {
_AudioSegment::Chunk(chunk) => {
if chunk.len() == 0 {
return Err(AudioError::ProcessingError("Received empty chunk".into()));
}
return Ok(segment);
}
_AudioSegment::Sync(marker) => match &**marker {
SyncMarker::TrackBoundary { .. } => {
// On ne devrait pas recevoir de TrackBoundary ici car on l'a déjà
tracing::warn!("FlacCacheSink: Unexpected TrackBoundary while waiting for first chunk");
continue;
}
SyncMarker::EndOfStream => {
return Err(AudioError::ProcessingError(
"EndOfStream received before any audio".into(),
));
}
_ => {
// Ignorer TopZeroSync, Heartbeat, etc.
continue;
}
},
}
}
}
async fn wait_for_first_audio_chunk_with_metadata(
rx: &mut mpsc::Receiver<Arc<AudioSegment>>,
stop_token: &CancellationToken,

View File

@@ -137,19 +137,22 @@ impl RadioParadiseStreamSourceLogic {
self.send_to_children(output, top_zero).await?;
tracing::debug!("TopZeroSync sent");
// Envoyer TrackBoundary pour la première song AVANT le premier chunk
// Cela garantit que FlacCacheSink reçoit les métadonnées dès le début
let mut next_song: Option<(usize, &Song)> = if let Some((_idx, song)) = songs.get(0).copied() {
tracing::debug!("Sending TrackBoundary for first song before audio chunks");
// Envoyer TrackBoundary pour la première song AVANT le premier chunk audio
// Même si son elapsed > 0, cela garantit que FlacCacheSink a des métadonnées
// dès le début (sinon il attendrait indéfiniment un TrackBoundary)
let mut next_song: Option<(usize, &Song)> = if let Some((idx, song)) = songs.get(0).copied() {
tracing::debug!("Sending TrackBoundary for first song (idx={}, elapsed={}ms) at timestamp 0",
idx, song.elapsed);
let metadata = song_to_metadata(song, block).await;
let track_boundary = AudioSegment::new_track_boundary(
*order,
0.0, // timestamp = 0 au début du bloc
0.0, // timestamp = 0 au début du stream
metadata,
);
self.send_to_children(output, track_boundary).await?;
song_index = 1;
songs.get(1).copied() // Passer à la song suivante
// Le prochain TrackBoundary sera pour la deuxième song quand elapsed_ms >= song.elapsed
songs.get(1).copied()
} else {
None
};
@@ -197,11 +200,15 @@ impl RadioParadiseStreamSourceLogic {
let chunk_len = (pcm_data.len() / (bytes_per_sample * 2)) as u64; // 2 = stereo
// Vérifier si on doit insérer un TrackBoundary avant ce chunk
if let Some((_idx, song)) = next_song {
if let Some((idx, song)) = next_song {
let elapsed_ms = (total_samples * 1000) / sample_rate as u64;
if elapsed_ms >= song.elapsed {
// Envoyer TrackBoundary AVANT le chunk (avec le même order)
tracing::debug!(
"Sending TrackBoundary for song {} at elapsed_ms={} (song.elapsed={}, timestamp_sec={:.2})",
idx, elapsed_ms, song.elapsed, (total_samples as f64 / sample_rate as f64)
);
let metadata = song_to_metadata(song, block).await;
let timestamp_sec = total_samples as f64 / sample_rate as f64;
let track_boundary = AudioSegment::new_track_boundary(
@@ -214,6 +221,7 @@ impl RadioParadiseStreamSourceLogic {
// Passer à la song suivante
song_index += 1;
next_song = songs.get(song_index).copied();
tracing::debug!("Moved to next song, song_index={}, next_song present={}", song_index, next_song.is_some());
}
}