Refactor FlacCacheSink for parallel write tasks to prevent file truncation

Problem: When TrackBoundary arrived, the pump was awaited before continuing,
causing file truncation when pcm_tx was dropped while data was still buffering.

Solution: Allow multiple pump tasks to run in parallel:
- Create dedicated channel (track_tx/track_rx) for each track's pump
- Main loop reads from rx and dispatches segments to current pump via track_tx
- When TrackBoundary arrives: drop track_tx (signals pump to finish) and immediately start new pump
- Old pump continues writing in background until all data is flushed

This prevents truncation in progressive cache scenario (radio streaming).

Changes in flac_cache_sink.rs:
- Replace pump_track_segments_owned() with pump_track_segments_from_channel()
- Remove rx ownership passing - each pump gets its own channel
- Dispatcher loop reads rx and forwards to active pump
- No await on pump completion - let it finish in background
This commit is contained in:
Claude
2025-11-07 13:54:23 +00:00
parent 7e81a8e777
commit 4ac81fecac

View File

@@ -154,17 +154,21 @@ impl NodeLogic for FlacCacheSinkLogic {
collection_ref,
);
// Spawner pump_future avec ownership de rx
// Cela permet d'attendre cache_future séparément et de pusher à la playlist immédiatement
let pump_handle = tokio::spawn(pump_track_segments_owned(
// Créer un channel dédié pour dispatcher les chunks vers ce pump
let (track_tx, track_rx) = mpsc::channel::<Arc<AudioSegment>>(16);
// Lancer le pump en arrière-plan avec son channel dédié
// Cela permet à plusieurs pumps de tourner simultanément (écriture parallèle)
let pump_handle = tokio::spawn(pump_track_segments_from_channel(
first_segment,
rx, // move ownership!
track_rx,
pcm_tx,
bits_per_sample,
sample_rate,
stop_token.clone(),
));
// Envoyer le first_segment déjà vers le track_tx est inutile car on l'a passé directement
// Attendre SEULEMENT le prebuffer (cache retourne après 512KB)
let start = std::time::Instant::now();
tracing::debug!("FlacCacheSink: Waiting for cache prebuffer to complete");
@@ -230,37 +234,71 @@ impl NodeLogic for FlacCacheSinkLogic {
tracing::info!("FlacCacheSink: Successfully pushed to playlist in {:?}", push_start.elapsed());
}
// MAINTENANT attendre que pump finisse (il continue en arrière-plan)
tracing::debug!("FlacCacheSink: Waiting for pump to complete");
let pump_result = pump_handle.await.map_err(|e| {
AudioError::ProcessingError(format!("Pump task panicked: {}", e))
})?;
// NE PAS attendre le pump - le laisser finir en arrière-plan
// Cela permet d'avoir plusieurs pumps en parallèle et évite la troncature
tracing::debug!("FlacCacheSink: Pump running in background, dispatching segments");
let (_chunks, _samples, _duration_sec, stop_reason, rx_returned) = pump_result?;
rx = rx_returned; // récupérer rx pour la prochaine track
tracing::debug!("FlacCacheSink: Pump completed");
// Si le fichier était déjà en cache (ChannelClosed), drainer les segments restants
// jusqu'au prochain TrackBoundary ou EndOfStream
// IMPORTANT: Faire ceci APRÈS l'ajout à la playlist pour ne pas bloquer la lecture
let stop_reason = if matches!(stop_reason, StopReason::ChannelClosed) {
tracing::debug!("File was already in cache, draining remaining segments");
drain_until_track_boundary(&mut rx, &stop_token).await?
} else {
stop_reason
// Dispatcher les segments depuis rx vers track_tx jusqu'au prochain TrackBoundary
loop {
let segment = tokio::select! {
result = rx.recv() => {
match result {
Some(seg) => seg,
None => {
// EOF sur rx - fin du stream, fermer le pump
drop(track_tx);
drop(pump_handle);
return Ok(());
}
}
}
_ = stop_token.cancelled() => {
drop(track_tx);
drop(pump_handle);
return Ok(());
}
};
// Vérifier le stop_reason pour savoir si on continue
match stop_reason {
StopReason::TrackBoundary(_metadata) => {
// Continuer avec la prochaine track
match &segment.segment {
_AudioSegment::Chunk(_) => {
// Dispatcher vers le pump actuel
if track_tx.send(segment).await.is_err() {
// Le pump est mort (channel fermé) - drainer jusqu'au TrackBoundary
tracing::warn!("FlacCacheSink: pump died, draining until TrackBoundary");
loop {
let seg = rx.recv().await;
match seg {
Some(s) if matches!(s.segment, _AudioSegment::Sync(ref m) if matches!(**m, SyncMarker::TrackBoundary { .. })) => {
track_number += 1;
continue;
break;
}
StopReason::EndOfStream | StopReason::ChannelClosed => {
// Fin de l'encodage
None => return Ok(()),
_ => continue,
}
}
break;
}
}
_AudioSegment::Sync(marker) => match &**marker {
SyncMarker::TrackBoundary { .. } => {
// Nouveau morceau - fermer le pump actuel et passer au suivant
drop(track_tx); // Ferme le channel, le pump va se terminer proprement
tracing::debug!("FlacCacheSink: TrackBoundary detected, pump will finish in background");
track_number += 1;
break;
}
SyncMarker::EndOfStream => {
tracing::debug!("FlacCacheSink: EndOfStream received");
drop(track_tx);
drop(pump_handle);
return Ok(());
}
_ => {
// Transmettre les autres syncmarkers au pump
let _ = track_tx.send(segment).await;
}
},
}
}
}
}
@@ -519,18 +557,17 @@ async fn pump_track_segments(
}
}
/// Pompe les segments pour une seule track (s'arrête au TrackBoundary).
/// Pompe les segments pour une seule track depuis un channel dédié.
///
/// Version qui prend ownership de rx pour permettre un await séparé du cache.
/// Retourne rx à la fin pour permettre le traitement des tracks suivantes.
async fn pump_track_segments_owned(
/// Cette version permet d'avoir plusieurs pumps en parallèle (pour cache progressif),
/// car chaque pump a son propre channel et ne bloque pas le traitement des tracks suivantes.
async fn pump_track_segments_from_channel(
first_segment: Arc<AudioSegment>,
mut rx: mpsc::Receiver<Arc<AudioSegment>>,
mut track_rx: mpsc::Receiver<Arc<AudioSegment>>,
pcm_tx: mpsc::Sender<Vec<u8>>,
bits_per_sample: u8,
expected_rate: u32,
stop_token: CancellationToken,
) -> Result<(u64, u64, f64, StopReason, mpsc::Receiver<Arc<AudioSegment>>), AudioError> {
) -> Result<(u64, u64, f64), AudioError> {
let mut chunks = 0u64;
let mut samples = 0u64;
let mut duration_sec = 0.0f64;
@@ -541,7 +578,8 @@ async fn pump_track_segments_owned(
if !pcm_bytes.is_empty() {
if pcm_tx.send(pcm_bytes).await.is_err() {
drop(pcm_tx);
return Ok((chunks, samples, duration_sec, StopReason::ChannelClosed, rx));
tracing::debug!("pump_track_segments_from_channel: pcm_tx closed on first segment");
return Ok((chunks, samples, duration_sec));
}
chunks += 1;
samples += chunk.len() as u64;
@@ -549,21 +587,15 @@ async fn pump_track_segments_owned(
}
}
// Boucle sur les segments suivants
// Boucle sur les segments depuis le channel dédié
loop {
let segment = tokio::select! {
result = rx.recv() => {
match result {
let segment = match track_rx.recv().await {
Some(seg) => seg,
None => {
// Channel fermé - la track est terminée (TrackBoundary a été reçu en amont)
drop(pcm_tx);
return Ok((chunks, samples, duration_sec, StopReason::ChannelClosed, rx));
}
}
}
_ = stop_token.cancelled() => {
drop(pcm_tx);
return Ok((chunks, samples, duration_sec, StopReason::ChannelClosed, rx));
tracing::debug!("pump_track_segments_from_channel: channel closed, track finished");
return Ok((chunks, samples, duration_sec));
}
};
@@ -583,31 +615,20 @@ async fn pump_track_segments_owned(
}
if pcm_tx.send(pcm_bytes).await.is_err() {
// Le cache a fermé le channel (erreur ou déjà en cache)
drop(pcm_tx);
return Ok((chunks, samples, duration_sec, StopReason::ChannelClosed, rx));
tracing::debug!("pump_track_segments_from_channel: pcm_tx closed");
return Ok((chunks, samples, duration_sec));
}
chunks += 1;
samples += chunk.len() as u64;
duration_sec += chunk.len() as f64 / expected_rate as f64;
}
_AudioSegment::Sync(marker) => match &**marker {
SyncMarker::TrackBoundary { metadata, .. } => {
drop(pcm_tx);
return Ok((
chunks,
samples,
duration_sec,
StopReason::TrackBoundary(metadata.clone()),
rx,
));
_AudioSegment::Sync(_marker) => {
// Ignorer les syncmarkers - le TrackBoundary est géré en amont
// Le channel sera fermé quand le TrackBoundary est détecté
}
SyncMarker::EndOfStream => {
drop(pcm_tx);
return Ok((chunks, samples, duration_sec, StopReason::EndOfStream, rx));
}
_ => {} // Ignorer les autres syncmarkers
},
}
}
}