From dba9f668f6d673467779d77a0857ee5baf3bce77 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 7 Nov 2025 14:36:19 +0000 Subject: [PATCH] Fix critical deadlock in FlacCacheSink parallel write architecture MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Problem: The dispatcher was placed AFTER the prebuffer await, causing a deadlock: - Pump waits for data on track_rx - Cache waits for pump to produce PCM data - Code awaits cache completion before reaching dispatcher - Dispatcher never runs → pump never receives data → deadlock Solution: Use tokio::select! to dispatch segments in parallel with prebuffer wait Architecture now has 3 phases: 1. Phase 1: Dispatch chunks + await prebuffer (in parallel via select!) 2. Phase 2: Copy metadata + push to playlist (after prebuffer complete) 3. Phase 3: Continue dispatching until TrackBoundary This fixes the "sans musique" blocking issue where the system would freeze waiting for prebuffer that could never complete. --- pmoaudio-ext/src/sinks/flac_cache_sink.rs | 118 +++++++++++++++------- 1 file changed, 84 insertions(+), 34 deletions(-) diff --git a/pmoaudio-ext/src/sinks/flac_cache_sink.rs b/pmoaudio-ext/src/sinks/flac_cache_sink.rs index 11f67519..d0a0d359 100755 --- a/pmoaudio-ext/src/sinks/flac_cache_sink.rs +++ b/pmoaudio-ext/src/sinks/flac_cache_sink.rs @@ -167,18 +167,81 @@ impl NodeLogic for FlacCacheSinkLogic { sample_rate, )); - // 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) + // Dispatcher les segments vers track_tx en parallèle de l'attente du prebuffer + // Utiliser tokio::select! pour éviter le deadlock let start = std::time::Instant::now(); - tracing::debug!("FlacCacheSink: Waiting for cache prebuffer to complete"); - let pk = cache_future.await.map_err(|e| { - AudioError::ProcessingError(format!("Failed to add to cache: {}", e)) - })?; + tracing::debug!("FlacCacheSink: Starting dispatcher loop with prebuffer wait"); - let prebuffer_time = start.elapsed(); - tracing::info!("FlacCacheSink: Prebuffer complete with pk {} in {:?}, pushing to playlist NOW", pk, prebuffer_time); + // Pin la future pour pouvoir l'utiliser dans select! + tokio::pin!(cache_future); + // Phase 1: Dispatcher jusqu'à ce que le prebuffer soit terminé + let pk = loop { + tokio::select! { + // Attendre le prebuffer + result = &mut cache_future => { + match result { + Ok(pk) => { + let prebuffer_time = start.elapsed(); + tracing::info!("FlacCacheSink: Prebuffer complete with pk {} in {:?}", pk, prebuffer_time); + break pk; // Sort de la loop pour faire les métadonnées et le push + } + Err(e) => { + return Err(AudioError::ProcessingError(format!("Failed to add to cache: {}", e))); + } + } + } + + // Dispatcher les segments depuis rx vers track_tx + result = rx.recv() => { + match result { + Some(segment) => { + 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())); + } + } + _AudioSegment::Sync(marker) => match &**marker { + SyncMarker::TrackBoundary { .. } => { + // TrackBoundary avant fin du prebuffer - track trop courte + tracing::error!("FlacCacheSink: TrackBoundary received before prebuffer complete - track too short"); + 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(()); + } + _ => { + // Transmettre les autres syncmarkers au pump + let _ = track_tx.send(segment).await; + } + }, + } + } + None => { + // EOF sur rx + drop(track_tx); + drop(pump_handle); + return Ok(()); + } + } + } + + _ = stop_token.cancelled() => { + drop(track_tx); + drop(pump_handle); + return Ok(()); + } + } + }; + + // Phase 2: Prebuffer terminé! Copier les métadonnées et pusher à la playlist // Copier les métadonnées du TrackBoundary dans le cache // IMPORTANT: Faire ceci AVANT d'ajouter à la playlist pour que les métadonnées soient disponibles if let Some(src_metadata) = track_metadata.clone() { @@ -234,18 +297,15 @@ impl NodeLogic for FlacCacheSinkLogic { tracing::info!("FlacCacheSink: Successfully pushed to playlist in {:?}", push_start.elapsed()); } - // 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"); - - // Dispatcher les segments depuis rx vers track_tx jusqu'au prochain TrackBoundary + // Phase 3: Continuer à dispatcher jusqu'au TrackBoundary + tracing::debug!("FlacCacheSink: Continuing dispatch until TrackBoundary (pump runs in background)"); loop { let segment = tokio::select! { result = rx.recv() => { match result { Some(seg) => seg, None => { - // EOF sur rx - fin du stream, fermer le pump + // EOF sur rx drop(track_tx); drop(pump_handle); return Ok(()); @@ -261,31 +321,21 @@ impl NodeLogic for FlacCacheSinkLogic { match &segment.segment { _AudioSegment::Chunk(_) => { - // Dispatcher vers le pump actuel + // Continuer à dispatcher vers le pump 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; + // Le pump est mort - erreur + tracing::error!("FlacCacheSink: pump died during post-prebuffer phase"); + return Err(AudioError::ProcessingError("Pump task died".to_string())); } } _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"); + // Nouveau morceau - fermer le pump et passer au suivant + tracing::debug!("FlacCacheSink: TrackBoundary received, closing pump"); + drop(track_tx); // Ferme le channel, le pump se termine proprement + drop(pump_handle); track_number += 1; - break; + break; // Sort de la Phase 3, retour à la loop externe pour next track } SyncMarker::EndOfStream => { tracing::debug!("FlacCacheSink: EndOfStream received");