diff --git a/pmoaudio-ext/src/sinks/flac_cache_sink.rs b/pmoaudio-ext/src/sinks/flac_cache_sink.rs index 0bbb38da..81dde7a9 100755 --- a/pmoaudio-ext/src/sinks/flac_cache_sink.rs +++ b/pmoaudio-ext/src/sinks/flac_cache_sink.rs @@ -299,6 +299,7 @@ impl NodeLogic for FlacCacheSinkLogic { // Phase 3: Continuer à dispatcher jusqu'au TrackBoundary tracing::debug!("FlacCacheSink: Continuing dispatch until TrackBoundary (pump runs in background)"); + let mut pump_closed = false; loop { let segment = tokio::select! { result = rx.recv() => { @@ -306,75 +307,79 @@ impl NodeLogic for FlacCacheSinkLogic { Some(seg) => seg, None => { // EOF sur rx - drop(track_tx); - drop(pump_handle); + if !pump_closed { + drop(track_tx); + drop(pump_handle); + } return Ok(()); } } } _ = stop_token.cancelled() => { - drop(track_tx); - drop(pump_handle); + if !pump_closed { + drop(track_tx); + drop(pump_handle); + } return Ok(()); } }; match &segment.segment { _AudioSegment::Chunk(_) => { - // Continuer à dispatcher vers le pump - if track_tx.send(segment).await.is_err() { - // Le pump a fermé son channel - cela peut arriver si le fichier - // était déjà en cache (add_from_reader retourne immédiatement) - tracing::debug!("FlacCacheSink: pump closed track_tx, checking pump status"); - drop(track_tx); + // Continuer à dispatcher vers le pump (sauf si déjà fermé) + if !pump_closed { + if track_tx.send(segment).await.is_err() { + // Le pump a fermé son channel - cela peut arriver si le fichier + // était déjà en cache (add_from_reader retourne immédiatement) + tracing::debug!("FlacCacheSink: pump closed track_tx, checking pump status"); + drop(track_tx); - // Attendre que le pump se termine et vérifier le résultat - match pump_handle.await { - Ok(Ok(_)) => { - // Le pump s'est terminé proprement (fichier était en cache) - tracing::debug!("FlacCacheSink: pump completed successfully, draining remaining segments"); - // Drainer les segments restants jusqu'au TrackBoundary - match drain_until_track_boundary(&mut rx, &stop_token).await? { - StopReason::TrackBoundary(_) => { - track_number += 1; - break; // Continue avec la prochaine track - } - StopReason::EndOfStream | StopReason::ChannelClosed => { - return Ok(()); - } + // Attendre que le pump se termine et vérifier le résultat + match pump_handle.await { + Ok(Ok(_)) => { + // Le pump s'est terminé proprement (fichier était en cache) + tracing::debug!("FlacCacheSink: pump completed successfully, ignoring remaining chunks until TrackBoundary"); + pump_closed = true; + } + Ok(Err(e)) => { + // Le pump a rencontré une erreur + tracing::error!("FlacCacheSink: pump died with error: {}", e); + return Err(e); + } + Err(e) => { + // Le pump task a paniqué + tracing::error!("FlacCacheSink: pump task panicked: {}", e); + return Err(AudioError::ProcessingError("Pump task panicked".to_string())); } - } - Ok(Err(e)) => { - // Le pump a rencontré une erreur - tracing::error!("FlacCacheSink: pump died with error: {}", e); - return Err(e); - } - Err(e) => { - // Le pump task a paniqué - tracing::error!("FlacCacheSink: pump task panicked: {}", e); - return Err(AudioError::ProcessingError("Pump task panicked".to_string())); } } } + // Si pump_closed, ignorer silencieusement le chunk } _AudioSegment::Sync(marker) => match &**marker { SyncMarker::TrackBoundary { .. } => { - // Nouveau morceau - fermer le pump et passer au suivant + // Nouveau morceau - fermer le pump si pas déjà fermé tracing::debug!("FlacCacheSink: TrackBoundary received, closing pump"); - drop(track_tx); // Ferme le channel, le pump se termine proprement - drop(pump_handle); + if !pump_closed { + drop(track_tx); // Ferme le channel, le pump se termine proprement + drop(pump_handle); + } track_number += 1; break; // Sort de la Phase 3, retour à la loop externe pour next track } SyncMarker::EndOfStream => { tracing::debug!("FlacCacheSink: EndOfStream received"); - drop(track_tx); - drop(pump_handle); + if !pump_closed { + drop(track_tx); + drop(pump_handle); + } return Ok(()); } _ => { - // Transmettre les autres syncmarkers au pump - let _ = track_tx.send(segment).await; + // Transmettre les autres syncmarkers au pump (sauf si fermé) + if !pump_closed { + let _ = track_tx.send(segment).await; + } } }, }