diff --git a/pmoaudio-ext/src/sinks/flac_cache_sink.rs b/pmoaudio-ext/src/sinks/flac_cache_sink.rs index 1c16555a..11f67519 100755 --- a/pmoaudio-ext/src/sinks/flac_cache_sink.rs +++ b/pmoaudio-ext/src/sinks/flac_cache_sink.rs @@ -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::>(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,36 +234,70 @@ 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"); + // 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(()); + } + }; - // 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 - }; - - // Vérifier le stop_reason pour savoir si on continue - match stop_reason { - StopReason::TrackBoundary(_metadata) => { - // Continuer avec la prochaine track - track_number += 1; - continue; - } - StopReason::EndOfStream | StopReason::ChannelClosed => { - // Fin de l'encodage - return Ok(()); + 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; + break; + } + 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, - mut rx: mpsc::Receiver>, + mut track_rx: mpsc::Receiver>, pcm_tx: mpsc::Sender>, bits_per_sample: u8, expected_rate: u32, - stop_token: CancellationToken, -) -> Result<(u64, u64, f64, StopReason, mpsc::Receiver>), 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 { - Some(seg) => seg, - None => { - drop(pcm_tx); - return Ok((chunks, samples, duration_sec, StopReason::ChannelClosed, rx)); - } - } - } - _ = stop_token.cancelled() => { + 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)); + 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, - )); - } - SyncMarker::EndOfStream => { - drop(pcm_tx); - return Ok((chunks, samples, duration_sec, StopReason::EndOfStream, rx)); - } - _ => {} // Ignorer les autres syncmarkers - }, + _AudioSegment::Sync(_marker) => { + // Ignorer les syncmarkers - le TrackBoundary est géré en amont + // Le channel sera fermé quand le TrackBoundary est détecté + } } } }