Retour sur pmoaudio

This commit is contained in:
2025-11-01 21:10:57 +01:00
parent ec0a0e675e
commit 6bce26c4fc
26 changed files with 4481 additions and 1048 deletions

View File

@@ -0,0 +1,472 @@
//! Nodes de conversion de type pour AudioChunk
//!
//! Ces nodes permettent de convertir les chunks audio d'un type vers un autre
//! (I16, I24, I32, F32, F64). Toutes les conversions utilisent les fonctions
//! DSP optimisées SIMD du module `crate::conversions`.
//!
//! Le designer de pipeline doit insérer manuellement ces nodes pour gérer
//! les incompatibilités de type entre producers et consumers.
use crate::{
nodes::{AudioError, MultiSubscriberNode, TypedAudioNode},
type_constraints::{SampleType, TypeRequirement},
AudioSegment,
};
use std::sync::Arc;
use tokio::sync::mpsc;
/// Node de conversion vers I16
///
/// Convertit n'importe quel type de chunk audio vers I16 (16-bit signed integer).
/// Utilise les conversions DSP SIMD optimisées.
pub struct ToI16Node {
rx: mpsc::Receiver<Arc<AudioSegment>>,
subscribers: MultiSubscriberNode,
}
impl ToI16Node {
/// Crée un nouveau node de conversion vers I16
pub fn new() -> (Self, mpsc::Sender<Arc<AudioSegment>>) {
Self::with_channel_size(16)
}
/// Crée un nouveau node avec une taille de buffer spécifique
pub fn with_channel_size(channel_size: usize) -> (Self, mpsc::Sender<Arc<AudioSegment>>) {
let (tx, rx) = mpsc::channel(channel_size);
let node = Self {
rx,
subscribers: MultiSubscriberNode::new(),
};
(node, tx)
}
/// Ajoute un abonné qui recevra les segments audio convertis
pub fn add_subscriber(&mut self, tx: mpsc::Sender<Arc<AudioSegment>>) {
self.subscribers.add_subscriber(tx);
}
/// Lance le traitement de conversion
pub async fn run(mut self) -> Result<(), AudioError> {
while let Some(segment) = self.rx.recv().await {
// Si c'est un syncmarker, passer directement
if !segment.is_audio_chunk() {
self.subscribers.push(segment).await?;
continue;
}
// Convertir le chunk audio vers I16
let converted_segment = if let Some(chunk) = segment.as_chunk() {
let converted_chunk = chunk.to_i16();
Arc::new(AudioSegment {
order: segment.order,
timestamp_sec: segment.timestamp_sec,
segment: crate::_AudioSegment::Chunk(Arc::new(converted_chunk)),
})
} else {
segment
};
self.subscribers.push(converted_segment).await?;
}
Ok(())
}
}
impl TypedAudioNode for ToI16Node {
fn input_type(&self) -> Option<TypeRequirement> {
// Accepte n'importe quel type
Some(TypeRequirement::any())
}
fn output_type(&self) -> Option<TypeRequirement> {
// Produit uniquement I16
Some(TypeRequirement::specific(SampleType::I16))
}
}
impl Default for ToI16Node {
fn default() -> Self {
Self::new().0
}
}
/// Node de conversion vers I24
///
/// Convertit n'importe quel type de chunk audio vers I24 (24-bit signed integer).
/// Utilise les conversions DSP SIMD optimisées.
pub struct ToI24Node {
rx: mpsc::Receiver<Arc<AudioSegment>>,
subscribers: MultiSubscriberNode,
}
impl ToI24Node {
/// Crée un nouveau node de conversion vers I24
pub fn new() -> (Self, mpsc::Sender<Arc<AudioSegment>>) {
Self::with_channel_size(16)
}
/// Crée un nouveau node avec une taille de buffer spécifique
pub fn with_channel_size(channel_size: usize) -> (Self, mpsc::Sender<Arc<AudioSegment>>) {
let (tx, rx) = mpsc::channel(channel_size);
let node = Self {
rx,
subscribers: MultiSubscriberNode::new(),
};
(node, tx)
}
/// Ajoute un abonné qui recevra les segments audio convertis
pub fn add_subscriber(&mut self, tx: mpsc::Sender<Arc<AudioSegment>>) {
self.subscribers.add_subscriber(tx);
}
/// Lance le traitement de conversion
pub async fn run(mut self) -> Result<(), AudioError> {
while let Some(segment) = self.rx.recv().await {
if !segment.is_audio_chunk() {
self.subscribers.push(segment).await?;
continue;
}
let converted_segment = if let Some(chunk) = segment.as_chunk() {
let converted_chunk = chunk.to_i24();
Arc::new(AudioSegment {
order: segment.order,
timestamp_sec: segment.timestamp_sec,
segment: crate::_AudioSegment::Chunk(Arc::new(converted_chunk)),
})
} else {
segment
};
self.subscribers.push(converted_segment).await?;
}
Ok(())
}
}
impl TypedAudioNode for ToI24Node {
fn input_type(&self) -> Option<TypeRequirement> {
Some(TypeRequirement::any())
}
fn output_type(&self) -> Option<TypeRequirement> {
Some(TypeRequirement::specific(SampleType::I24))
}
}
impl Default for ToI24Node {
fn default() -> Self {
Self::new().0
}
}
/// Node de conversion vers I32
///
/// Convertit n'importe quel type de chunk audio vers I32 (32-bit signed integer).
/// Utilise les conversions DSP SIMD optimisées.
pub struct ToI32Node {
rx: mpsc::Receiver<Arc<AudioSegment>>,
subscribers: MultiSubscriberNode,
}
impl ToI32Node {
/// Crée un nouveau node de conversion vers I32
pub fn new() -> (Self, mpsc::Sender<Arc<AudioSegment>>) {
Self::with_channel_size(16)
}
/// Crée un nouveau node avec une taille de buffer spécifique
pub fn with_channel_size(channel_size: usize) -> (Self, mpsc::Sender<Arc<AudioSegment>>) {
let (tx, rx) = mpsc::channel(channel_size);
let node = Self {
rx,
subscribers: MultiSubscriberNode::new(),
};
(node, tx)
}
/// Ajoute un abonné qui recevra les segments audio convertis
pub fn add_subscriber(&mut self, tx: mpsc::Sender<Arc<AudioSegment>>) {
self.subscribers.add_subscriber(tx);
}
/// Lance le traitement de conversion
pub async fn run(mut self) -> Result<(), AudioError> {
while let Some(segment) = self.rx.recv().await {
if !segment.is_audio_chunk() {
self.subscribers.push(segment).await?;
continue;
}
let converted_segment = if let Some(chunk) = segment.as_chunk() {
let converted_chunk = chunk.to_i32();
Arc::new(AudioSegment {
order: segment.order,
timestamp_sec: segment.timestamp_sec,
segment: crate::_AudioSegment::Chunk(Arc::new(converted_chunk)),
})
} else {
segment
};
self.subscribers.push(converted_segment).await?;
}
Ok(())
}
}
impl TypedAudioNode for ToI32Node {
fn input_type(&self) -> Option<TypeRequirement> {
Some(TypeRequirement::any())
}
fn output_type(&self) -> Option<TypeRequirement> {
Some(TypeRequirement::specific(SampleType::I32))
}
}
impl Default for ToI32Node {
fn default() -> Self {
Self::new().0
}
}
/// Node de conversion vers F32
///
/// Convertit n'importe quel type de chunk audio vers F32 (32-bit floating point).
/// Utilise les conversions DSP SIMD optimisées.
pub struct ToF32Node {
rx: mpsc::Receiver<Arc<AudioSegment>>,
subscribers: MultiSubscriberNode,
}
impl ToF32Node {
/// Crée un nouveau node de conversion vers F32
pub fn new() -> (Self, mpsc::Sender<Arc<AudioSegment>>) {
Self::with_channel_size(16)
}
/// Crée un nouveau node avec une taille de buffer spécifique
pub fn with_channel_size(channel_size: usize) -> (Self, mpsc::Sender<Arc<AudioSegment>>) {
let (tx, rx) = mpsc::channel(channel_size);
let node = Self {
rx,
subscribers: MultiSubscriberNode::new(),
};
(node, tx)
}
/// Ajoute un abonné qui recevra les segments audio convertis
pub fn add_subscriber(&mut self, tx: mpsc::Sender<Arc<AudioSegment>>) {
self.subscribers.add_subscriber(tx);
}
/// Lance le traitement de conversion
pub async fn run(mut self) -> Result<(), AudioError> {
while let Some(segment) = self.rx.recv().await {
if !segment.is_audio_chunk() {
self.subscribers.push(segment).await?;
continue;
}
let converted_segment = if let Some(chunk) = segment.as_chunk() {
let converted_chunk = chunk.to_f32();
Arc::new(AudioSegment {
order: segment.order,
timestamp_sec: segment.timestamp_sec,
segment: crate::_AudioSegment::Chunk(Arc::new(converted_chunk)),
})
} else {
segment
};
self.subscribers.push(converted_segment).await?;
}
Ok(())
}
}
impl TypedAudioNode for ToF32Node {
fn input_type(&self) -> Option<TypeRequirement> {
Some(TypeRequirement::any())
}
fn output_type(&self) -> Option<TypeRequirement> {
Some(TypeRequirement::specific(SampleType::F32))
}
}
impl Default for ToF32Node {
fn default() -> Self {
Self::new().0
}
}
/// Node de conversion vers F64
///
/// Convertit n'importe quel type de chunk audio vers F64 (64-bit floating point).
/// Utilise les conversions DSP SIMD optimisées.
pub struct ToF64Node {
rx: mpsc::Receiver<Arc<AudioSegment>>,
subscribers: MultiSubscriberNode,
}
impl ToF64Node {
/// Crée un nouveau node de conversion vers F64
pub fn new() -> (Self, mpsc::Sender<Arc<AudioSegment>>) {
Self::with_channel_size(16)
}
/// Crée un nouveau node avec une taille de buffer spécifique
pub fn with_channel_size(channel_size: usize) -> (Self, mpsc::Sender<Arc<AudioSegment>>) {
let (tx, rx) = mpsc::channel(channel_size);
let node = Self {
rx,
subscribers: MultiSubscriberNode::new(),
};
(node, tx)
}
/// Ajoute un abonné qui recevra les segments audio convertis
pub fn add_subscriber(&mut self, tx: mpsc::Sender<Arc<AudioSegment>>) {
self.subscribers.add_subscriber(tx);
}
/// Lance le traitement de conversion
pub async fn run(mut self) -> Result<(), AudioError> {
while let Some(segment) = self.rx.recv().await {
if !segment.is_audio_chunk() {
self.subscribers.push(segment).await?;
continue;
}
let converted_segment = if let Some(chunk) = segment.as_chunk() {
let converted_chunk = chunk.to_f64();
Arc::new(AudioSegment {
order: segment.order,
timestamp_sec: segment.timestamp_sec,
segment: crate::_AudioSegment::Chunk(Arc::new(converted_chunk)),
})
} else {
segment
};
self.subscribers.push(converted_segment).await?;
}
Ok(())
}
}
impl TypedAudioNode for ToF64Node {
fn input_type(&self) -> Option<TypeRequirement> {
Some(TypeRequirement::any())
}
fn output_type(&self) -> Option<TypeRequirement> {
Some(TypeRequirement::specific(SampleType::F64))
}
}
impl Default for ToF64Node {
fn default() -> Self {
Self::new().0
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{AudioChunk, AudioChunkData};
#[tokio::test]
async fn test_to_f32_node_type_requirements() {
let (node, _tx) = ToF32Node::new();
// Vérifier les types d'entrée/sortie
assert_eq!(
node.input_type().unwrap().get_accepted_types().len(),
5,
"Should accept all 5 types"
);
assert_eq!(
node.output_type()
.unwrap()
.get_accepted_types()
.first()
.copied(),
Some(SampleType::F32),
"Should output F32 only"
);
}
#[tokio::test]
async fn test_to_i16_node_converts_from_i32() {
let (mut node, tx) = ToI16Node::new();
let (out_tx, mut out_rx) = mpsc::channel(16);
node.add_subscriber(out_tx);
// Lancer le node dans une tâche
let handle = tokio::spawn(async move { node.run().await });
// Créer et envoyer un chunk I32
let stereo = vec![[1_000_000i32 << 16, -500_000i32 << 16]; 100];
let chunk_data = AudioChunkData::new(stereo.clone(), 48_000, 0.0);
let chunk = AudioChunk::I32(chunk_data);
let segment = Arc::new(AudioSegment {
order: 0,
timestamp_sec: 0.0,
segment: crate::_AudioSegment::Chunk(Arc::new(chunk)),
});
tx.send(segment).await.unwrap();
drop(tx);
// Recevoir le chunk converti
let result = out_rx.recv().await.unwrap();
assert!(result.is_audio_chunk());
if let Some(converted) = result.as_chunk() {
assert_eq!(converted.type_name(), "i16");
assert_eq!(converted.len(), 100);
// Vérifier la conversion (downsampling de I32 vers I16)
if let AudioChunk::I16(data) = &**converted {
for (orig, converted_frame) in stereo.iter().zip(data.frames().iter()) {
let expected_l = (orig[0] >> 16) as i16;
let expected_r = (orig[1] >> 16) as i16;
assert_eq!(converted_frame[0], expected_l);
assert_eq!(converted_frame[1], expected_r);
}
}
}
handle.await.unwrap().unwrap();
}
#[tokio::test]
async fn test_syncmarkers_passthrough() {
let (mut node, tx) = ToI16Node::new();
let (out_tx, mut out_rx) = mpsc::channel(16);
node.add_subscriber(out_tx);
tokio::spawn(async move {
node.run().await.unwrap();
});
// Envoyer un syncmarker
let top_zero = AudioSegment::new_top_zero_sync();
tx.send(top_zero.clone()).await.unwrap();
drop(tx);
// Recevoir le syncmarker
let result = out_rx.recv().await.unwrap();
assert!(!result.is_audio_chunk());
assert!(result.as_sync_marker().is_some());
}
}

View File

@@ -0,0 +1,426 @@
use crate::{
nodes::{AudioError, MultiSubscriberNode, TypedAudioNode, DEFAULT_CHUNK_DURATION_MS},
type_constraints::TypeRequirement,
AudioChunk, AudioChunkData, AudioSegment, I24,
};
use pmoflac::{decode_audio_stream, AudioFileMetadata, StreamInfo};
use pmometadata::{MemoryTrackMetadata, TrackMetadata};
use std::{path::PathBuf, sync::Arc, time::Duration};
use tokio::{fs::File, io::AsyncReadExt, sync::mpsc};
/// FileSource - Lit un fichier audio et publie des `AudioSegment`
///
/// Cette source utilise `pmoflac` pour décoder le fichier (FLAC/MP3/OGG/WAV/AIFF)
/// puis transforme les échantillons PCM en `AudioSegment` stéréo avec le type approprié
/// (I16, I24, ou I32) selon la profondeur de bit du fichier source.
///
/// Le node émet trois types de syncmarkers :
/// - `TopZeroSync` au début du flux
/// - `TrackBoundary` avec les métadonnées du fichier
/// - `EndOfStream` à la fin du flux
pub struct FileSource {
path: PathBuf,
chunk_frames: usize,
subscribers: MultiSubscriberNode,
}
impl FileSource {
/// Crée une nouvelle source de fichier avec calcul automatique de la taille des chunks.
///
/// La taille des chunks sera calculée automatiquement pour obtenir environ 50ms
/// de latence par chunk, en fonction du sample rate du fichier.
///
/// * `path` - chemin du fichier audio à lire
pub fn new<P: Into<PathBuf>>(path: P) -> Self {
Self::with_chunk_size(path, 0) // 0 = auto-calculer
}
/// Crée une nouvelle source de fichier avec une taille de chunk spécifique.
///
/// * `path` - chemin du fichier audio à lire
/// * `chunk_frames` - nombre d'échantillons par canal par chunk (0 = auto)
pub fn with_chunk_size<P: Into<PathBuf>>(path: P, chunk_frames: usize) -> Self {
Self {
path: path.into(),
chunk_frames,
subscribers: MultiSubscriberNode::new(),
}
}
/// Ajoute un abonné qui recevra les segments audio.
pub fn add_subscriber(&mut self, tx: mpsc::Sender<Arc<AudioSegment>>) {
self.subscribers.add_subscriber(tx);
}
/// Lance la lecture du fichier et diffuse les segments audio.
pub async fn run(self) -> Result<(), AudioError> {
// Ouvrir le fichier
let file = File::open(&self.path).await.map_err(|e| {
AudioError::ProcessingError(format!("Failed to open {:?}: {}", self.path, e))
})?;
// Décoder le flux audio
let mut stream = decode_audio_stream(file)
.await
.map_err(|e| AudioError::ProcessingError(format!("Decode error: {}", e)))?;
let stream_info = stream.info().clone();
validate_stream(&stream_info)?;
// Calculer la taille des chunks si non spécifiée (0 = auto)
let chunk_frames = if self.chunk_frames == 0 {
// Calculer pour obtenir DEFAULT_CHUNK_DURATION_MS millisecondes
let frames =
(stream_info.sample_rate as f64 * DEFAULT_CHUNK_DURATION_MS / 1000.0) as usize;
// Arrondir à la puissance de 2 la plus proche pour optimiser les buffers
frames.next_power_of_two().max(256)
} else {
self.chunk_frames.max(1)
};
// Émettre TopZeroSync
let top_zero = AudioSegment::new_top_zero_sync();
self.subscribers.push(top_zero).await?;
// Extraire et émettre les métadonnées du fichier
match AudioFileMetadata::from_file(&self.path) {
Ok(file_metadata) => {
let mut metadata = MemoryTrackMetadata::new();
// Convertir AudioFileMetadata vers MemoryTrackMetadata
if let Some(title) = file_metadata.title {
let _ = metadata.set_title(Some(title)).await;
}
if let Some(artist) = file_metadata.artist {
let _ = metadata.set_artist(Some(artist)).await;
}
if let Some(album) = file_metadata.album {
let _ = metadata.set_album(Some(album)).await;
}
if let Some(year) = file_metadata.year {
let _ = metadata.set_year(Some(year)).await;
}
if let Some(duration_secs) = file_metadata.duration_secs {
let _ = metadata
.set_duration(Some(Duration::from_secs(duration_secs)))
.await;
}
// Émettre TrackBoundary
let track_boundary = AudioSegment::new_track_boundary(0, 0.0, Arc::new(metadata));
self.subscribers.push(track_boundary).await?;
}
Err(e) => {
eprintln!(
"Warning: Failed to extract metadata from {:?}: {}",
self.path, e
);
// Continuer sans métadonnées
}
}
// Préparer la lecture des chunks audio
let frame_bytes = stream_info.bytes_per_sample() * stream_info.channels as usize;
let chunk_byte_len = chunk_frames * frame_bytes;
let mut pending = Vec::new();
let mut read_buf = vec![0u8; frame_bytes * 512.max(chunk_frames)];
let mut chunk_index = 0u64;
let mut total_frames = 0u64;
// Lire et émettre les chunks audio
loop {
// Remplir le buffer
if pending.len() < chunk_byte_len {
let read = stream.read(&mut read_buf).await.map_err(|e| {
AudioError::ProcessingError(format!("I/O error while decoding: {}", e))
})?;
if read == 0 {
break;
}
pending.extend_from_slice(&read_buf[..read]);
}
if pending.is_empty() {
break;
}
// Extraire un chunk
let frames_in_pending = pending.len() / frame_bytes;
let frames_to_emit = frames_in_pending.min(chunk_frames);
let take_bytes = frames_to_emit * frame_bytes;
let chunk_bytes = pending.drain(..take_bytes).collect::<Vec<u8>>();
// Calculer le timestamp
let timestamp_sec = total_frames as f64 / stream_info.sample_rate as f64;
// Créer le segment audio
let segment = bytes_to_segment(
&chunk_bytes,
&stream_info,
frames_to_emit,
chunk_index,
timestamp_sec,
)?;
self.subscribers.push(segment).await?;
chunk_index += 1;
total_frames += frames_to_emit as u64;
}
// Traiter le reste éventuel (moins qu'un chunk complet)
if !pending.is_empty() {
let frames = pending.len() / frame_bytes;
if frames > 0 {
let timestamp_sec = total_frames as f64 / stream_info.sample_rate as f64;
let segment =
bytes_to_segment(&pending, &stream_info, frames, chunk_index, timestamp_sec)?;
self.subscribers.push(segment).await?;
total_frames += frames as u64;
chunk_index += 1;
}
}
// Émettre EndOfStream
let final_timestamp = total_frames as f64 / stream_info.sample_rate as f64;
let eos = AudioSegment::new_end_of_stream(chunk_index, final_timestamp);
self.subscribers.push(eos).await?;
// Attendre la fin du décodage
stream
.wait()
.await
.map_err(|e| AudioError::ProcessingError(format!("Decode task failed: {}", e)))?;
Ok(())
}
}
fn validate_stream(info: &StreamInfo) -> Result<(), AudioError> {
if !(1..=2).contains(&info.channels) {
return Err(AudioError::ProcessingError(format!(
"Unsupported channel count: {}",
info.channels
)));
}
match info.bits_per_sample {
8 | 16 | 24 | 32 => Ok(()),
other => Err(AudioError::ProcessingError(format!(
"Unsupported bit depth: {}",
other
))),
}
}
/// Convertit des bytes PCM en AudioSegment avec le type approprié
fn bytes_to_segment(
chunk_bytes: &[u8],
info: &StreamInfo,
frames: usize,
order: u64,
timestamp_sec: f64,
) -> Result<Arc<AudioSegment>, AudioError> {
let bytes_per_sample = info.bytes_per_sample();
let channels = info.channels as usize;
let frame_bytes = bytes_per_sample * channels;
// Créer le chunk du bon type selon la profondeur de bit
let chunk = match info.bits_per_sample {
16 => {
// Type I16
let mut stereo = Vec::with_capacity(frames);
for frame_idx in 0..frames {
let base = frame_idx * frame_bytes;
let l = i16::from_le_bytes(
chunk_bytes[base..base + bytes_per_sample]
.try_into()
.unwrap(),
);
let r = if channels == 1 {
l
} else {
i16::from_le_bytes(
chunk_bytes[base + bytes_per_sample..base + 2 * bytes_per_sample]
.try_into()
.unwrap(),
)
};
stereo.push([l, r]);
}
let chunk_data = AudioChunkData::new(stereo, info.sample_rate, 0.0);
AudioChunk::I16(chunk_data)
}
24 => {
// Type I24
let mut stereo = Vec::with_capacity(frames);
for frame_idx in 0..frames {
let base = frame_idx * frame_bytes;
let l_i32 = {
let mut buf = [0u8; 4];
buf[..3].copy_from_slice(&chunk_bytes[base..base + 3]);
// Sign extend
if chunk_bytes[base + 2] & 0x80 != 0 {
buf[3] = 0xFF;
}
i32::from_le_bytes(buf)
};
let l = I24::new(l_i32).ok_or_else(|| {
AudioError::ProcessingError(format!("Invalid I24 value: {}", l_i32))
})?;
let r = if channels == 1 {
l
} else {
let r_i32 = {
let mut buf = [0u8; 4];
buf[..3].copy_from_slice(
&chunk_bytes[base + bytes_per_sample..base + bytes_per_sample + 3],
);
// Sign extend
if chunk_bytes[base + bytes_per_sample + 2] & 0x80 != 0 {
buf[3] = 0xFF;
}
i32::from_le_bytes(buf)
};
I24::new(r_i32).ok_or_else(|| {
AudioError::ProcessingError(format!("Invalid I24 value: {}", r_i32))
})?
};
stereo.push([l, r]);
}
let chunk_data = AudioChunkData::new(stereo, info.sample_rate, 0.0);
AudioChunk::I24(chunk_data)
}
32 => {
// Type I32
let mut stereo = Vec::with_capacity(frames);
for frame_idx in 0..frames {
let base = frame_idx * frame_bytes;
let l = i32::from_le_bytes(
chunk_bytes[base..base + bytes_per_sample]
.try_into()
.unwrap(),
);
let r = if channels == 1 {
l
} else {
i32::from_le_bytes(
chunk_bytes[base + bytes_per_sample..base + 2 * bytes_per_sample]
.try_into()
.unwrap(),
)
};
stereo.push([l, r]);
}
let chunk_data = AudioChunkData::new(stereo, info.sample_rate, 0.0);
AudioChunk::I32(chunk_data)
}
other => {
return Err(AudioError::ProcessingError(format!(
"Unsupported bit depth: {}",
other
)))
}
};
// Créer le segment audio
Ok(Arc::new(AudioSegment {
order,
timestamp_sec,
segment: crate::_AudioSegment::Chunk(Arc::new(chunk)),
}))
}
impl TypedAudioNode for FileSource {
fn input_type(&self) -> Option<TypeRequirement> {
// FileSource est une source, elle ne consomme pas d'audio
None
}
fn output_type(&self) -> Option<TypeRequirement> {
// FileSource peut produire n'importe quel type entier (I16, I24, I32)
// selon la profondeur de bit du fichier source
Some(TypeRequirement::any_integer())
}
}
#[cfg(test)]
mod tests {
use super::*;
use pmoflac::{encode_flac_stream, EncoderOptions, PcmFormat};
use std::io::Cursor;
use tokio::io::AsyncWriteExt;
use tokio::sync::mpsc;
#[tokio::test]
async fn test_file_source_decodes_flac() {
let temp_dir = tempfile::tempdir().unwrap();
let flac_path = temp_dir.path().join("test.flac");
let sample_rate = 48_000;
let frames = 256;
let mut pcm = Vec::with_capacity(frames * 4);
for i in 0..frames {
let sample = ((i % 32) as f32 / 31.0 * 2.0 - 1.0) * 0.5; // simple ramp
let sample_i16 = (sample * 32767.0) as i16;
pcm.extend_from_slice(&sample_i16.to_le_bytes());
pcm.extend_from_slice(&sample_i16.to_le_bytes());
}
let format = PcmFormat {
sample_rate,
channels: 2,
bits_per_sample: 16,
};
let mut flac_stream =
encode_flac_stream(Cursor::new(pcm.clone()), format, EncoderOptions::default())
.await
.unwrap();
let mut file = File::create(&flac_path).await.expect("create flac file");
tokio::io::copy(&mut flac_stream, &mut file)
.await
.expect("write flac");
file.flush().await.expect("flush file");
flac_stream.wait().await.unwrap();
let mut source = FileSource::with_chunk_size(&flac_path, 64);
let (tx, mut rx) = mpsc::channel(16);
source.add_subscriber(tx);
tokio::spawn(async move {
source.run().await.unwrap();
});
let mut received_frames = 0usize;
let mut received_syncmarkers = 0usize;
let mut seen_top_zero = false;
let mut seen_eos = false;
while let Some(segment) = rx.recv().await {
if segment.is_audio_chunk() {
if let Some(chunk) = segment.as_chunk() {
received_frames += chunk.len();
assert_eq!(chunk.sample_rate(), sample_rate);
}
} else {
received_syncmarkers += 1;
if let Some(marker) = segment.as_sync_marker() {
match **marker {
crate::SyncMarker::TopZeroSync => seen_top_zero = true,
crate::SyncMarker::EndOfStream => seen_eos = true,
crate::SyncMarker::TrackBoundary { .. } => {}
_ => {}
}
}
}
}
// Vérifier que tous les frames ont été reçus
assert_eq!(received_frames, frames);
// Vérifier qu'on a bien reçu des syncmarkers
assert!(received_syncmarkers >= 2); // Au moins TopZeroSync et EndOfStream
assert!(seen_top_zero, "Should have received TopZeroSync");
assert!(seen_eos, "Should have received EndOfStream");
}
}

View File

@@ -0,0 +1,766 @@
use crate::{
nodes::{AudioError, TypedAudioNode, DEFAULT_CHANNEL_SIZE},
type_constraints::TypeRequirement,
AudioChunk, AudioSegment, SyncMarker,
};
use pmoflac::{encode_flac_stream, EncoderOptions, PcmFormat};
use std::{
collections::VecDeque,
path::{Path, PathBuf},
pin::Pin,
sync::Arc,
task::{Context, Poll},
};
use tokio::{
fs::File,
io::{self, AsyncRead, AsyncWriteExt, ReadBuf},
sync::mpsc,
};
/// Sink qui encode les `AudioSegment` reçus au format FLAC.
///
/// Ce sink :
/// - Filtre les chunks audio et ignore les autres syncmarkers (sauf TrackBoundary et EndOfStream)
/// - Crée un nouveau fichier FLAC pour chaque TrackBoundary rencontré
/// - Adapte automatiquement l'encodage FLAC selon la profondeur de bit du chunk (8/16/24/32-bit)
/// - Termine l'encodage proprement quand il reçoit EndOfStream
pub struct FlacFileSink {
rx: mpsc::Receiver<Arc<AudioSegment>>,
base_path: PathBuf,
encoder_options: EncoderOptions,
pcm_buffer_capacity: usize,
}
impl FlacFileSink {
/// Crée un sink FLAC avec les options par défaut (compression 5, buffer de 16 segments).
///
/// # Arguments
///
/// * `base_path` - Chemin de base pour les fichiers FLAC. Si des TrackBoundary sont reçus,
/// des fichiers seront créés avec des suffixes (_01, _02, etc.)
pub fn new<P: Into<PathBuf>>(base_path: P) -> (Self, mpsc::Sender<Arc<AudioSegment>>) {
Self::with_channel_size(base_path, DEFAULT_CHANNEL_SIZE)
}
/// Crée un sink FLAC avec une taille de buffer MPSC personnalisée.
///
/// # Arguments
///
/// * `base_path` - Chemin de base pour les fichiers FLAC
/// * `channel_size` - Taille du buffer MPSC (nombre de segments en attente avant backpressure)
pub fn with_channel_size<P: Into<PathBuf>>(
base_path: P,
channel_size: usize,
) -> (Self, mpsc::Sender<Arc<AudioSegment>>) {
Self::with_config(base_path, channel_size, EncoderOptions::default())
}
/// Crée un sink FLAC avec une configuration complète.
///
/// # Arguments
///
/// * `base_path` - Chemin de base pour les fichiers FLAC
/// * `channel_size` - Taille du buffer MPSC
/// * `encoder_options` - Options d'encodage FLAC (compression, etc.)
pub fn with_config<P: Into<PathBuf>>(
base_path: P,
channel_size: usize,
encoder_options: EncoderOptions,
) -> (Self, mpsc::Sender<Arc<AudioSegment>>) {
let (tx, rx) = mpsc::channel(channel_size);
let sink = Self {
rx,
base_path: base_path.into(),
encoder_options,
pcm_buffer_capacity: 8,
};
(sink, tx)
}
/// Lance l'encodage vers le(s) fichier(s) cible(s).
///
/// Cette méthode crée un nouveau fichier FLAC pour chaque TrackBoundary rencontré.
/// Les fichiers sont nommés selon la convention :
/// - Track 0 : base_path.flac
/// - Track 1 : base_path_01.flac
/// - Track 2 : base_path_02.flac, etc.
pub async fn run(self) -> Result<FlacFileSinkStats, AudioError> {
let FlacFileSink {
mut rx,
base_path,
encoder_options,
pcm_buffer_capacity,
} = self;
let mut all_tracks = Vec::new();
let mut track_number = 0;
loop {
// Attendre le premier chunk audio pour cette track, en capturant les métadonnées du TrackBoundary
let (first_segment, track_metadata) = match wait_for_first_audio_chunk_with_metadata(&mut rx).await {
Ok(result) => result,
Err(_) => {
// Plus d'audio disponible
if all_tracks.is_empty() {
return Err(AudioError::ProcessingError("No audio data received".into()));
}
break;
}
};
// Extraire les informations du premier chunk
let first_chunk = first_segment.as_chunk().unwrap();
let sample_rate = first_chunk.sample_rate();
let bits_per_sample = get_chunk_bit_depth(first_chunk);
let format = PcmFormat {
sample_rate,
channels: 2,
bits_per_sample,
};
if let Err(err) = format.validate() {
return Err(AudioError::ProcessingError(format!(
"Invalid PCM format: {}",
err
)));
}
// Générer le chemin du fichier pour cette track
let track_path = generate_track_path(&base_path, track_number);
// Créer le pipeline d'encodage pour cette track
let (pcm_tx, pcm_rx) = mpsc::channel::<Vec<u8>>(pcm_buffer_capacity);
// Préparer les options d'encodage avec les métadonnées du TrackBoundary
let mut options_with_metadata = encoder_options.clone();
options_with_metadata.metadata = track_metadata;
// Créer l'encoder et le fichier
let reader = ByteStreamReader::new(pcm_rx);
let mut flac_stream = encode_flac_stream(reader, format, options_with_metadata)
.await
.map_err(|e| {
AudioError::ProcessingError(format!("FLAC encode init failed: {}", e))
})?;
let mut output = File::create(&track_path).await.map_err(|e| {
AudioError::ProcessingError(format!("Failed to create {:?}: {}", track_path, e))
})?;
// Exécuter pump et copy en parallèle avec tokio::select! en boucle
let pump_future =
pump_track_segments(first_segment, &mut rx, pcm_tx, bits_per_sample, sample_rate);
let copy_future = async {
let copy_result = tokio::io::copy(&mut flac_stream, &mut output).await;
let flush_result = output.flush().await;
let wait_result = flac_stream.wait().await;
copy_result.map_err(|e| {
AudioError::ProcessingError(format!("FLAC write failed: {}", e))
})?;
flush_result
.map_err(|e| AudioError::ProcessingError(format!("Failed to flush: {}", e)))?;
wait_result
.map_err(|e| AudioError::ProcessingError(format!("Encoder failed: {}", e)))?;
Ok::<_, AudioError>(())
};
// Attendre les deux tâches en parallèle
let (copy_result, pump_result) = tokio::join!(copy_future, pump_future);
copy_result?;
let (chunks, samples, duration_sec, stop_reason) = pump_result?;
// Ajouter les stats de cette track
all_tracks.push(TrackStats {
path: track_path,
track_number,
chunks_received: chunks,
total_samples: samples,
total_duration_sec: duration_sec,
});
// 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
break;
}
}
}
Ok(FlacFileSinkStats { tracks: all_tracks })
}
}
/// Génère le chemin de fichier pour une track donnée.
/// - track 0 → base_path.flac
/// - track 1 → base_path_01.flac
/// - track 2 → base_path_02.flac, etc.
fn generate_track_path(base_path: &Path, track_number: usize) -> PathBuf {
if track_number == 0 {
base_path.to_path_buf()
} else {
let stem = base_path
.file_stem()
.and_then(|s| s.to_str())
.unwrap_or("output");
let extension = base_path
.extension()
.and_then(|s| s.to_str())
.unwrap_or("flac");
let parent = base_path.parent().unwrap_or(Path::new("."));
parent.join(format!("{}_{:02}.{}", stem, track_number, extension))
}
}
/// Signal retourné par pump_segments indiquant pourquoi l'encodage s'est arrêté.
enum StopReason {
TrackBoundary(Arc<dyn pmometadata::TrackMetadata + Send + Sync>),
EndOfStream,
ChannelClosed,
}
/// 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.
async fn wait_for_first_audio_chunk_with_metadata(
rx: &mut mpsc::Receiver<Arc<AudioSegment>>,
) -> Result<(Arc<AudioSegment>, Option<Arc<dyn pmometadata::TrackMetadata + Send + Sync>>), AudioError> {
let mut track_metadata: Option<Arc<dyn pmometadata::TrackMetadata + Send + Sync>> = None;
loop {
let segment = rx
.recv()
.await
.ok_or_else(|| AudioError::ProcessingError("No audio data received".into()))?;
match &segment.segment {
crate::_AudioSegment::Chunk(chunk) => {
if chunk.len() == 0 {
return Err(AudioError::ProcessingError("Received empty chunk".into()));
}
return Ok((segment, track_metadata));
}
crate::_AudioSegment::Sync(marker) => {
match **marker {
SyncMarker::TrackBoundary { ref metadata, .. } => {
// Capturer les métadonnées du TrackBoundary
track_metadata = Some(metadata.clone());
continue;
}
SyncMarker::EndOfStream => {
return Err(AudioError::ProcessingError(
"EndOfStream received before any audio".into(),
));
}
_ => {
// Ignorer TopZeroSync, Heartbeat, etc.
continue;
}
}
}
}
}
}
/// Pompe les segments pour une seule track (s'arrête au TrackBoundary).
async fn pump_track_segments(
first_segment: Arc<AudioSegment>,
rx: &mut mpsc::Receiver<Arc<AudioSegment>>,
pcm_tx: mpsc::Sender<Vec<u8>>,
bits_per_sample: u8,
expected_rate: u32,
) -> Result<(u64, u64, f64, StopReason), AudioError> {
let mut chunks = 0u64;
let mut samples = 0u64;
let mut duration_sec = 0.0f64;
// Traiter le premier segment
if let Some(chunk) = first_segment.as_chunk() {
let pcm_bytes = chunk_to_pcm_bytes(chunk, bits_per_sample)?;
if !pcm_bytes.is_empty() {
pcm_tx
.send(pcm_bytes)
.await
.map_err(|_| AudioError::SendError)?;
chunks += 1;
samples += chunk.len() as u64;
duration_sec += chunk.len() as f64 / expected_rate as f64;
}
}
// Boucle sur les segments suivants
loop {
let segment = match rx.recv().await {
Some(seg) => seg,
None => {
drop(pcm_tx); // Fermer le channel PCM
return Ok((chunks, samples, duration_sec, StopReason::ChannelClosed));
}
};
match &segment.segment {
crate::_AudioSegment::Chunk(chunk) => {
// Vérifier la cohérence du sample rate
if chunk.sample_rate() != expected_rate {
return Err(AudioError::ProcessingError(format!(
"FlacFileSink: inconsistent sample rate ({} vs {})",
chunk.sample_rate(),
expected_rate
)));
}
let pcm_bytes = chunk_to_pcm_bytes(chunk, bits_per_sample)?;
if pcm_bytes.is_empty() {
continue;
}
pcm_tx
.send(pcm_bytes)
.await
.map_err(|_| AudioError::SendError)?;
chunks += 1;
samples += chunk.len() as u64;
duration_sec += chunk.len() as f64 / expected_rate as f64;
}
crate::_AudioSegment::Sync(marker) => {
match &**marker {
SyncMarker::TrackBoundary { metadata, .. } => {
drop(pcm_tx); // Fermer le channel PCM
return Ok((
chunks,
samples,
duration_sec,
StopReason::TrackBoundary(metadata.clone()),
));
}
SyncMarker::EndOfStream => {
drop(pcm_tx); // Fermer le channel PCM
return Ok((chunks, samples, duration_sec, StopReason::EndOfStream));
}
_ => {} // Ignorer les autres syncmarkers
}
}
}
}
}
/// Détermine la profondeur de bit d'un chunk audio
fn get_chunk_bit_depth(chunk: &AudioChunk) -> u8 {
match chunk {
AudioChunk::I16(_) => 16,
AudioChunk::I24(_) => 24,
AudioChunk::I32(_) => 32,
AudioChunk::F32(_) => 32, // Les flottants seront convertis en 32-bit
AudioChunk::F64(_) => 32, // Les flottants seront convertis en 32-bit
}
}
/// Convertit un chunk audio en bytes PCM avec la profondeur de bit spécifiée
fn chunk_to_pcm_bytes(chunk: &AudioChunk, bits_per_sample: u8) -> Result<Vec<u8>, AudioError> {
// Vérifier que le chunk est de type entier
match chunk {
AudioChunk::F32(_) | AudioChunk::F64(_) => {
return Err(AudioError::ProcessingError(
"FlacFileSink only supports integer audio chunks (I16, I24, I32)".into(),
));
}
_ => {}
}
let len = chunk.len();
let bytes_per_frame = (bits_per_sample / 8) as usize * 2; // 2 channels
let mut bytes = Vec::with_capacity(len * bytes_per_frame);
// Convertir selon le type du chunk
match (chunk, bits_per_sample) {
// I16 source
(AudioChunk::I16(data), 16) => {
for frame in data.frames() {
bytes.extend_from_slice(&frame[0].to_le_bytes());
bytes.extend_from_slice(&frame[1].to_le_bytes());
}
}
(AudioChunk::I16(data), 24) => {
for frame in data.frames() {
let left = (frame[0] as i32) << 8;
let right = (frame[1] as i32) << 8;
bytes.extend_from_slice(&left.to_le_bytes()[..3]);
bytes.extend_from_slice(&right.to_le_bytes()[..3]);
}
}
(AudioChunk::I16(data), 32) => {
for frame in data.frames() {
let left = (frame[0] as i32) << 16;
let right = (frame[1] as i32) << 16;
bytes.extend_from_slice(&left.to_le_bytes());
bytes.extend_from_slice(&right.to_le_bytes());
}
}
// I24 source
(AudioChunk::I24(data), 16) => {
for frame in data.frames() {
let left = (frame[0].as_i32() >> 8) as i16;
let right = (frame[1].as_i32() >> 8) as i16;
bytes.extend_from_slice(&left.to_le_bytes());
bytes.extend_from_slice(&right.to_le_bytes());
}
}
(AudioChunk::I24(data), 24) => {
for frame in data.frames() {
bytes.extend_from_slice(&frame[0].as_i32().to_le_bytes()[..3]);
bytes.extend_from_slice(&frame[1].as_i32().to_le_bytes()[..3]);
}
}
(AudioChunk::I24(data), 32) => {
for frame in data.frames() {
let left = frame[0].as_i32() << 8;
let right = frame[1].as_i32() << 8;
bytes.extend_from_slice(&left.to_le_bytes());
bytes.extend_from_slice(&right.to_le_bytes());
}
}
// I32 source
(AudioChunk::I32(data), 16) => {
for frame in data.frames() {
let left = (frame[0] >> 16) as i16;
let right = (frame[1] >> 16) as i16;
bytes.extend_from_slice(&left.to_le_bytes());
bytes.extend_from_slice(&right.to_le_bytes());
}
}
(AudioChunk::I32(data), 24) => {
for frame in data.frames() {
let left = frame[0] >> 8;
let right = frame[1] >> 8;
bytes.extend_from_slice(&left.to_le_bytes()[..3]);
bytes.extend_from_slice(&right.to_le_bytes()[..3]);
}
}
(AudioChunk::I32(data), 32) => {
for frame in data.frames() {
bytes.extend_from_slice(&frame[0].to_le_bytes());
bytes.extend_from_slice(&frame[1].to_le_bytes());
}
}
_ => {
return Err(AudioError::ProcessingError(format!(
"Unsupported bits_per_sample: {}",
bits_per_sample
)));
}
}
Ok(bytes)
}
struct ByteStreamReader {
rx: mpsc::Receiver<Vec<u8>>,
buffer: VecDeque<u8>,
finished: bool,
}
impl ByteStreamReader {
fn new(rx: mpsc::Receiver<Vec<u8>>) -> Self {
Self {
rx,
buffer: VecDeque::new(),
finished: false,
}
}
}
impl AsyncRead for ByteStreamReader {
fn poll_read(
mut self: Pin<&mut Self>,
cx: &mut Context<'_>,
buf: &mut ReadBuf<'_>,
) -> Poll<io::Result<()>> {
loop {
if !self.buffer.is_empty() {
let to_copy = self.buffer.len().min(buf.remaining());
if to_copy == 0 {
return Poll::Ready(Ok(()));
}
// VecDeque::make_contiguous pour copier efficacement
let slice = self.buffer.make_contiguous();
buf.put_slice(&slice[..to_copy]);
self.buffer.drain(..to_copy);
return Poll::Ready(Ok(()));
}
if self.finished {
return Poll::Ready(Ok(()));
}
match Pin::new(&mut self.rx).poll_recv(cx) {
Poll::Ready(Some(bytes)) => {
if bytes.is_empty() {
continue;
}
self.buffer.extend(bytes);
}
Poll::Ready(None) => {
self.finished = true;
return Poll::Ready(Ok(()));
}
Poll::Pending => return Poll::Pending,
}
}
}
}
/// Statistiques pour une track individuelle.
#[derive(Debug, Clone)]
pub struct TrackStats {
pub path: PathBuf,
pub track_number: usize,
pub chunks_received: u64,
pub total_samples: u64,
pub total_duration_sec: f64,
}
/// Statistiques produites par le `FlacFileSink`.
#[derive(Debug, Clone)]
pub struct FlacFileSinkStats {
pub tracks: Vec<TrackStats>,
}
impl TypedAudioNode for FlacFileSink {
fn input_type(&self) -> Option<TypeRequirement> {
// FlacFileSink accepte n'importe quel type entier (I16, I24, I32)
// mais rejette les chunks flottants
Some(TypeRequirement::any_integer())
}
fn output_type(&self) -> Option<TypeRequirement> {
// FlacFileSink est un sink, il ne produit pas d'audio
None
}
}
#[cfg(test)]
mod tests {
use super::*;
use pmoflac::{decode_flac_stream, AudioFileMetadata};
use pmometadata::{MemoryTrackMetadata, TrackMetadata};
use tokio::io::AsyncReadExt;
#[tokio::test]
async fn test_flac_file_sink_writes_metadata() {
let temp_dir = tempfile::tempdir().unwrap();
let output_path = temp_dir.path().join("output_with_metadata.flac");
let sample_rate = 44_100;
let frames = 256;
// Créer le sink
let (sink, tx) = FlacFileSink::with_channel_size(&output_path, 16);
let sink_handle = tokio::spawn(async move { sink.run().await.unwrap() });
// Envoyer des segments avec métadonnées
tokio::spawn(async move {
// TopZeroSync
tx.send(crate::AudioSegment::new_top_zero_sync())
.await
.unwrap();
// TrackBoundary avec métadonnées
let mut metadata = MemoryTrackMetadata::new();
metadata.set_title(Some("Test Track Title".to_string())).await.unwrap();
metadata.set_artist(Some("Test Artist".to_string())).await.unwrap();
metadata.set_album(Some("Test Album".to_string())).await.unwrap();
metadata.set_year(Some(2024)).await.unwrap();
let track_boundary =
crate::AudioSegment::new_track_boundary(0, 0.0, std::sync::Arc::new(metadata));
tx.send(track_boundary).await.unwrap();
// Générer et envoyer des chunks audio
let chunk_frames = 64;
let mut order = 0u64;
let mut total_frames = 0u64;
for chunk_start in (0..frames).step_by(chunk_frames) {
let chunk_len = (frames - chunk_start).min(chunk_frames);
let mut stereo = Vec::with_capacity(chunk_len);
for i in 0..chunk_len {
let frame_idx = chunk_start + i;
let sample = ((frame_idx % 32) as f32 / 31.0 * 2.0 - 1.0) * 0.5;
let sample_i16 = (sample * 32767.0) as i16;
stereo.push([sample_i16, sample_i16]);
}
let timestamp = total_frames as f64 / sample_rate as f64;
let chunk_data = crate::AudioChunkData::new(stereo, sample_rate, 0.0);
let chunk = crate::AudioChunk::I16(chunk_data);
let segment = crate::AudioSegment {
order,
timestamp_sec: timestamp,
segment: crate::_AudioSegment::Chunk(std::sync::Arc::new(chunk)),
};
tx.send(std::sync::Arc::new(segment)).await.unwrap();
total_frames += chunk_len as u64;
order += 1;
}
// EndOfStream
let final_timestamp = total_frames as f64 / sample_rate as f64;
tx.send(crate::AudioSegment::new_end_of_stream(
order,
final_timestamp,
))
.await
.unwrap();
drop(tx);
});
sink_handle.await.unwrap();
// Vérifier que le fichier a été créé et contient les métadonnées
assert!(output_path.exists(), "Output file should exist");
// Lire les métadonnées du fichier FLAC généré
let file_metadata = AudioFileMetadata::from_file(&output_path).unwrap();
// Vérifier que les métadonnées ont été correctement écrites
assert_eq!(file_metadata.title, Some("Test Track Title".to_string()));
assert_eq!(file_metadata.artist, Some("Test Artist".to_string()));
assert_eq!(file_metadata.album, Some("Test Album".to_string()));
assert_eq!(file_metadata.year, Some(2024));
}
#[tokio::test]
async fn test_flac_file_sink_writes_audio() {
use pmoflac::{encode_flac_stream, EncoderOptions, PcmFormat};
use std::io::Cursor;
let temp_dir = tempfile::tempdir().unwrap();
let input_path = temp_dir.path().join("input.flac");
let output_path = temp_dir.path().join("output.flac");
// Créer un petit fichier FLAC de test (comme dans file_source test)
let sample_rate = 44_100;
let frames = 512;
let mut pcm = Vec::with_capacity(frames * 4);
for i in 0..frames {
let sample = ((i % 32) as f32 / 31.0 * 2.0 - 1.0) * 0.5;
let sample_i16 = (sample * 32767.0) as i16;
pcm.extend_from_slice(&sample_i16.to_le_bytes());
pcm.extend_from_slice(&sample_i16.to_le_bytes());
}
let format = PcmFormat {
sample_rate,
channels: 2,
bits_per_sample: 16,
};
let mut flac_stream =
encode_flac_stream(Cursor::new(pcm.clone()), format, EncoderOptions::default())
.await
.unwrap();
let mut input_file = File::create(&input_path).await.unwrap();
tokio::io::copy(&mut flac_stream, &mut input_file)
.await
.unwrap();
input_file.flush().await.unwrap();
flac_stream.wait().await.unwrap();
// Maintenant utiliser FlacFileSink pour réécrire le fichier
let (sink, tx) = FlacFileSink::with_channel_size(&output_path, 16);
let sink_handle = tokio::spawn(async move { sink.run().await.unwrap() });
// Lire le fichier input et envoyer les segments au sink
tokio::spawn(async move {
let source_file = File::open(&input_path).await.unwrap();
let mut decode_stream = pmoflac::decode_audio_stream(source_file).await.unwrap();
let info = decode_stream.info().clone();
// TopZeroSync
tx.send(crate::AudioSegment::new_top_zero_sync())
.await
.unwrap();
// Lire et envoyer les chunks
let mut buffer = vec![0u8; info.bytes_per_sample() * info.channels as usize * 256];
let mut total_frames = 0u64;
let mut order = 0u64;
loop {
let read = decode_stream.read(&mut buffer).await.unwrap();
if read == 0 {
break;
}
let chunk_frames = read / (info.bytes_per_sample() * info.channels as usize);
let timestamp = total_frames as f64 / info.sample_rate as f64;
// Créer un segment I16
let mut stereo = Vec::with_capacity(chunk_frames);
for i in 0..chunk_frames {
let offset = i * info.bytes_per_sample() * info.channels as usize;
let l = i16::from_le_bytes([buffer[offset], buffer[offset + 1]]);
let r = i16::from_le_bytes([buffer[offset + 2], buffer[offset + 3]]);
stereo.push([l, r]);
}
let chunk_data = crate::AudioChunkData::new(stereo, info.sample_rate, 0.0);
let chunk = crate::AudioChunk::I16(chunk_data);
let segment = crate::AudioSegment {
order,
timestamp_sec: timestamp,
segment: crate::_AudioSegment::Chunk(std::sync::Arc::new(chunk)),
};
tx.send(std::sync::Arc::new(segment)).await.unwrap();
total_frames += chunk_frames as u64;
order += 1;
}
// EndOfStream
let final_timestamp = total_frames as f64 / info.sample_rate as f64;
tx.send(crate::AudioSegment::new_end_of_stream(
order,
final_timestamp,
))
.await
.unwrap();
drop(tx);
decode_stream.wait().await.unwrap();
});
let stats = sink_handle.await.unwrap();
assert_eq!(stats.tracks.len(), 1);
assert!(stats.tracks[0].chunks_received > 0);
assert_eq!(stats.tracks[0].total_samples, frames as u64);
// Vérifier que le fichier de sortie est valide
let file = File::open(&output_path).await.unwrap();
let mut stream = decode_flac_stream(file).await.unwrap();
let info = stream.info().clone();
assert_eq!(info.channels, 2);
assert_eq!(info.sample_rate, sample_rate);
assert_eq!(info.bits_per_sample, 16);
let mut decoded = Vec::new();
stream.read_to_end(&mut decoded).await.unwrap();
stream.wait().await.unwrap();
assert!(decoded.len() > 0);
}
}

View File

@@ -0,0 +1,827 @@
use crate::{
nodes::{AudioError, MultiSubscriberNode, TypedAudioNode, DEFAULT_CHUNK_DURATION_MS},
type_constraints::TypeRequirement,
AudioChunk, AudioChunkData, AudioSegment, I24,
};
use futures_util::StreamExt;
use pmoflac::{decode_audio_stream, StreamInfo};
use pmometadata::{MemoryTrackMetadata, TrackMetadata};
use std::sync::Arc;
use tokio::sync::mpsc;
use tokio_util::io::StreamReader;
/// HttpSource - Récupère un fichier audio via HTTP et publie des `AudioSegment`
///
/// Cette source télécharge un fichier audio depuis une URL HTTP/HTTPS,
/// utilise `pmoflac` pour le décoder (FLAC/MP3/OGG/WAV/AIFF) puis transforme
/// les échantillons PCM en `AudioSegment` stéréo avec le type approprié.
///
/// Le node émet trois types de syncmarkers :
/// - `TopZeroSync` au début du flux
/// - `TrackBoundary` avec les métadonnées extraites des headers HTTP
/// - `EndOfStream` à la fin du flux
///
/// # Métadonnées HTTP
///
/// Les métadonnées suivantes sont extraites des headers HTTP lorsqu'elles sont disponibles:
/// - `icy-name`: nom du stream (Icecast/Shoutcast) → utilisé comme titre
/// - `icy-url`: URL du stream source
/// - `content-type`: type MIME du contenu (ex: audio/flac, audio/mpeg)
///
/// Si aucun header `icy-name` n'est présent, le nom du fichier est extrait de l'URL
/// et utilisé comme titre.
///
/// # Exemples
///
/// ## Lecture d'un fichier FLAC distant
///
/// ```no_run
/// use pmoaudio::HttpSource;
/// use tokio::sync::mpsc;
///
/// #[tokio::main]
/// async fn main() {
/// let mut source = HttpSource::new("http://example.com/audio.flac");
/// let (tx, mut rx) = mpsc::channel(16);
/// source.add_subscriber(tx);
///
/// // Lancer la lecture dans une tâche séparée
/// tokio::spawn(async move {
/// source.run().await.unwrap();
/// });
///
/// // Recevoir et traiter les segments audio
/// while let Some(segment) = rx.recv().await {
/// if segment.is_audio_chunk() {
/// println!("Chunk reçu à {}s", segment.timestamp_sec);
/// }
/// }
/// }
/// ```
///
/// ## Stream Icecast/Shoutcast
///
/// ```no_run
/// use pmoaudio::HttpSource;
/// use tokio::sync::mpsc;
///
/// #[tokio::main]
/// async fn main() {
/// // Les métadonnées icy-name seront extraites automatiquement
/// let mut source = HttpSource::new("http://stream.example.com:8000/stream");
/// let (tx, rx) = mpsc::channel(32);
/// source.add_subscriber(tx);
///
/// tokio::spawn(async move {
/// source.run().await.unwrap();
/// });
/// }
/// ```
///
/// # Gestion des erreurs
///
/// La méthode `run()` peut retourner les erreurs suivantes:
/// - `AudioError::ProcessingError`: échec de connexion HTTP, status code non-200,
/// erreur de décodage audio, ou format non supporté
///
/// # Performance
///
/// - Le téléchargement et le décodage sont effectués en streaming
/// - Pas de buffering complet du fichier en mémoire
/// - La taille des chunks audio est calculée automatiquement pour ~50ms de latence
/// - Compatible avec les streams infinis (radios web, etc.)
pub struct HttpSource {
url: String,
chunk_frames: usize,
subscribers: MultiSubscriberNode,
}
impl HttpSource {
/// Crée une nouvelle source HTTP avec calcul automatique de la taille des chunks.
///
/// La taille des chunks sera calculée automatiquement pour obtenir environ 50ms
/// de latence par chunk, en fonction du sample rate du fichier distant.
///
/// # Arguments
///
/// * `url` - URL HTTP ou HTTPS du fichier audio à télécharger
///
/// # Exemples
///
/// ```no_run
/// use pmoaudio::HttpSource;
///
/// let source = HttpSource::new("http://example.com/music.flac");
/// ```
pub fn new<S: Into<String>>(url: S) -> Self {
Self::with_chunk_size(url, 0)
}
/// Crée une nouvelle source HTTP avec une taille de chunk spécifique.
///
/// # Arguments
///
/// * `url` - URL HTTP ou HTTPS du fichier audio à télécharger
/// * `chunk_frames` - nombre d'échantillons par canal par chunk (0 = auto-calcul)
///
/// # Exemples
///
/// ```no_run
/// use pmoaudio::HttpSource;
///
/// // Utiliser des chunks de 2048 frames
/// let source = HttpSource::with_chunk_size("http://example.com/music.mp3", 2048);
/// ```
pub fn with_chunk_size<S: Into<String>>(url: S, chunk_frames: usize) -> Self {
Self {
url: url.into(),
chunk_frames,
subscribers: MultiSubscriberNode::new(),
}
}
/// Ajoute un abonné qui recevra les segments audio.
///
/// Chaque abonné recevra une copie (via `Arc`) de tous les segments audio
/// produits par cette source, y compris les syncmarkers.
///
/// # Arguments
///
/// * `tx` - Channel sender pour recevoir les `AudioSegment`
///
/// # Exemples
///
/// ```no_run
/// use pmoaudio::HttpSource;
/// use tokio::sync::mpsc;
///
/// let mut source = HttpSource::new("http://example.com/audio.flac");
/// let (tx, rx) = mpsc::channel(16);
/// source.add_subscriber(tx);
/// ```
pub fn add_subscriber(&mut self, tx: mpsc::Sender<Arc<AudioSegment>>) {
self.subscribers.add_subscriber(tx);
}
/// Lance le téléchargement et la lecture du flux audio.
///
/// Cette méthode consomme `self` et exécute le pipeline complet:
/// 1. Effectue la requête HTTP GET vers l'URL spécifiée
/// 2. Vérifie le status HTTP (doit être 2xx)
/// 3. Extrait les métadonnées des headers HTTP
/// 4. Décode le flux audio en streaming
/// 5. Émet les syncmarkers et chunks audio vers les abonnés
///
/// La méthode se termine quand le flux est complètement lu ou en cas d'erreur.
///
/// # Erreurs
///
/// Retourne `AudioError::ProcessingError` si:
/// - La requête HTTP échoue (réseau, DNS, etc.)
/// - Le serveur retourne un status code non-2xx
/// - Le format audio n'est pas supporté
/// - Le décodage échoue
/// - Le fichier a un nombre de canaux non supporté (doit être 1 ou 2)
/// - La profondeur de bit n'est pas supportée (doit être 8, 16, 24 ou 32 bits)
///
/// # Exemples
///
/// ```no_run
/// use pmoaudio::HttpSource;
/// use tokio::sync::mpsc;
///
/// #[tokio::main]
/// async fn main() -> Result<(), Box<dyn std::error::Error>> {
/// let mut source = HttpSource::new("http://example.com/audio.flac");
/// let (tx, mut rx) = mpsc::channel(16);
/// source.add_subscriber(tx);
///
/// let handle = tokio::spawn(async move {
/// source.run().await
/// });
///
/// // Traiter les segments
/// while let Some(segment) = rx.recv().await {
/// println!("Segment reçu: order={}", segment.order);
/// }
///
/// handle.await??;
/// Ok(())
/// }
/// ```
pub async fn run(self) -> Result<(), AudioError> {
// Effectuer la requête HTTP
let response = reqwest::get(&self.url)
.await
.map_err(|e| {
AudioError::ProcessingError(format!("HTTP request failed for {}: {}", self.url, e))
})?;
// Vérifier le status
if !response.status().is_success() {
return Err(AudioError::ProcessingError(format!(
"HTTP request returned status {}: {}",
response.status(),
self.url
)));
}
// Extraire les métadonnées depuis les headers HTTP
let metadata = extract_metadata_from_headers(&response, &self.url).await;
// Convertir le stream de bytes en AsyncRead
let bytes_stream = response.bytes_stream();
let stream_reader = StreamReader::new(bytes_stream.map(|result| {
result.map_err(|e| std::io::Error::new(std::io::ErrorKind::Other, e))
}));
// Décoder le flux audio
let mut stream = decode_audio_stream(stream_reader)
.await
.map_err(|e| AudioError::ProcessingError(format!("Decode error: {}", e)))?;
let stream_info = stream.info().clone();
validate_stream(&stream_info)?;
// Calculer la taille des chunks si non spécifiée (0 = auto)
let chunk_frames = if self.chunk_frames == 0 {
let frames =
(stream_info.sample_rate as f64 * DEFAULT_CHUNK_DURATION_MS / 1000.0) as usize;
frames.next_power_of_two().max(256)
} else {
self.chunk_frames.max(1)
};
// Émettre TopZeroSync
let top_zero = AudioSegment::new_top_zero_sync();
self.subscribers.push(top_zero).await?;
// Émettre TrackBoundary avec les métadonnées HTTP
let track_boundary = AudioSegment::new_track_boundary(0, 0.0, Arc::new(metadata));
self.subscribers.push(track_boundary).await?;
// Préparer la lecture des chunks audio
let frame_bytes = stream_info.bytes_per_sample() * stream_info.channels as usize;
let chunk_byte_len = chunk_frames * frame_bytes;
let mut pending = Vec::new();
let mut read_buf = vec![0u8; frame_bytes * 512.max(chunk_frames)];
let mut chunk_index = 0u64;
let mut total_frames = 0u64;
// Lire et émettre les chunks audio
loop {
// Remplir le buffer
if pending.len() < chunk_byte_len {
use tokio::io::AsyncReadExt;
let read = stream.read(&mut read_buf).await.map_err(|e| {
AudioError::ProcessingError(format!("I/O error while decoding: {}", e))
})?;
if read == 0 {
break;
}
pending.extend_from_slice(&read_buf[..read]);
}
if pending.is_empty() {
break;
}
// Extraire un chunk
let frames_in_pending = pending.len() / frame_bytes;
let frames_to_emit = frames_in_pending.min(chunk_frames);
let take_bytes = frames_to_emit * frame_bytes;
let chunk_bytes = pending.drain(..take_bytes).collect::<Vec<u8>>();
// Calculer le timestamp
let timestamp_sec = total_frames as f64 / stream_info.sample_rate as f64;
// Créer le segment audio
let segment = bytes_to_segment(
&chunk_bytes,
&stream_info,
frames_to_emit,
chunk_index,
timestamp_sec,
)?;
self.subscribers.push(segment).await?;
chunk_index += 1;
total_frames += frames_to_emit as u64;
}
// Traiter le reste éventuel
if !pending.is_empty() {
let frames = pending.len() / frame_bytes;
if frames > 0 {
let timestamp_sec = total_frames as f64 / stream_info.sample_rate as f64;
let segment =
bytes_to_segment(&pending, &stream_info, frames, chunk_index, timestamp_sec)?;
self.subscribers.push(segment).await?;
total_frames += frames as u64;
chunk_index += 1;
}
}
// Émettre EndOfStream
let final_timestamp = total_frames as f64 / stream_info.sample_rate as f64;
let eos = AudioSegment::new_end_of_stream(chunk_index, final_timestamp);
self.subscribers.push(eos).await?;
// Attendre la fin du décodage
stream
.wait()
.await
.map_err(|e| AudioError::ProcessingError(format!("Decode task failed: {}", e)))?;
Ok(())
}
}
/// Extrait les métadonnées disponibles depuis les headers HTTP
async fn extract_metadata_from_headers(
response: &reqwest::Response,
url: &str,
) -> MemoryTrackMetadata {
let mut metadata = MemoryTrackMetadata::new();
let headers = response.headers();
// Icecast/Shoutcast stream name
if let Some(name) = headers.get("icy-name").and_then(|v| v.to_str().ok()) {
let _ = metadata.set_title(Some(name.to_string())).await;
}
// Icecast/Shoutcast stream URL (peut être utilisé comme source)
if let Some(stream_url) = headers.get("icy-url").and_then(|v| v.to_str().ok()) {
// On pourrait stocker ça dans un champ custom si nécessaire
eprintln!("Stream URL: {}", stream_url);
}
// Content-Type pour déterminer le format
if let Some(content_type) = headers.get("content-type").and_then(|v| v.to_str().ok()) {
eprintln!("Content-Type: {}", content_type);
// On pourrait utiliser ça pour valider le format attendu
}
// Si aucune métadonnée spécifique n'est trouvée, utiliser l'URL comme titre
if metadata.get_title().await.ok().flatten().is_none() {
// Extraire le nom du fichier depuis l'URL
if let Some(filename) = url.rsplit('/').next() {
if !filename.is_empty() {
let _ = metadata.set_title(Some(filename.to_string())).await;
}
}
}
metadata
}
fn validate_stream(info: &StreamInfo) -> Result<(), AudioError> {
if !(1..=2).contains(&info.channels) {
return Err(AudioError::ProcessingError(format!(
"Unsupported channel count: {}",
info.channels
)));
}
match info.bits_per_sample {
8 | 16 | 24 | 32 => Ok(()),
other => Err(AudioError::ProcessingError(format!(
"Unsupported bit depth: {}",
other
))),
}
}
/// Convertit des bytes PCM en AudioSegment avec le type approprié
fn bytes_to_segment(
chunk_bytes: &[u8],
info: &StreamInfo,
frames: usize,
order: u64,
timestamp_sec: f64,
) -> Result<Arc<AudioSegment>, AudioError> {
let bytes_per_sample = info.bytes_per_sample();
let channels = info.channels as usize;
let frame_bytes = bytes_per_sample * channels;
// Créer le chunk du bon type selon la profondeur de bit
let chunk = match info.bits_per_sample {
16 => {
// Type I16
let mut stereo = Vec::with_capacity(frames);
for frame_idx in 0..frames {
let base = frame_idx * frame_bytes;
let l = i16::from_le_bytes(
chunk_bytes[base..base + bytes_per_sample]
.try_into()
.unwrap(),
);
let r = if channels == 1 {
l
} else {
i16::from_le_bytes(
chunk_bytes[base + bytes_per_sample..base + 2 * bytes_per_sample]
.try_into()
.unwrap(),
)
};
stereo.push([l, r]);
}
let chunk_data = AudioChunkData::new(stereo, info.sample_rate, 0.0);
AudioChunk::I16(chunk_data)
}
24 => {
// Type I24
let mut stereo = Vec::with_capacity(frames);
for frame_idx in 0..frames {
let base = frame_idx * frame_bytes;
let l_i32 = {
let mut buf = [0u8; 4];
buf[..3].copy_from_slice(&chunk_bytes[base..base + 3]);
// Sign extend
if chunk_bytes[base + 2] & 0x80 != 0 {
buf[3] = 0xFF;
}
i32::from_le_bytes(buf)
};
let l = I24::new(l_i32).ok_or_else(|| {
AudioError::ProcessingError(format!("Invalid I24 value: {}", l_i32))
})?;
let r = if channels == 1 {
l
} else {
let r_i32 = {
let mut buf = [0u8; 4];
buf[..3].copy_from_slice(
&chunk_bytes[base + bytes_per_sample..base + bytes_per_sample + 3],
);
// Sign extend
if chunk_bytes[base + bytes_per_sample + 2] & 0x80 != 0 {
buf[3] = 0xFF;
}
i32::from_le_bytes(buf)
};
I24::new(r_i32).ok_or_else(|| {
AudioError::ProcessingError(format!("Invalid I24 value: {}", r_i32))
})?
};
stereo.push([l, r]);
}
let chunk_data = AudioChunkData::new(stereo, info.sample_rate, 0.0);
AudioChunk::I24(chunk_data)
}
32 => {
// Type I32
let mut stereo = Vec::with_capacity(frames);
for frame_idx in 0..frames {
let base = frame_idx * frame_bytes;
let l = i32::from_le_bytes(
chunk_bytes[base..base + bytes_per_sample]
.try_into()
.unwrap(),
);
let r = if channels == 1 {
l
} else {
i32::from_le_bytes(
chunk_bytes[base + bytes_per_sample..base + 2 * bytes_per_sample]
.try_into()
.unwrap(),
)
};
stereo.push([l, r]);
}
let chunk_data = AudioChunkData::new(stereo, info.sample_rate, 0.0);
AudioChunk::I32(chunk_data)
}
other => {
return Err(AudioError::ProcessingError(format!(
"Unsupported bit depth: {}",
other
)))
}
};
// Créer le segment audio
Ok(Arc::new(AudioSegment {
order,
timestamp_sec,
segment: crate::_AudioSegment::Chunk(Arc::new(chunk)),
}))
}
impl TypedAudioNode for HttpSource {
fn input_type(&self) -> Option<TypeRequirement> {
// HttpSource est une source, elle ne consomme pas d'audio
None
}
fn output_type(&self) -> Option<TypeRequirement> {
// HttpSource peut produire n'importe quel type entier (I16, I24, I32)
// selon la profondeur de bit du fichier source
Some(TypeRequirement::any_integer())
}
}
#[cfg(test)]
mod tests {
use super::*;
use pmoflac::{encode_flac_stream, EncoderOptions, PcmFormat};
use std::io::Cursor;
use tokio::sync::mpsc;
use wiremock::{
matchers::{method, path},
Mock, MockServer, ResponseTemplate,
};
/// Test de création basique de HttpSource
#[test]
fn test_http_source_creation() {
let source = HttpSource::new("http://example.com/audio.flac");
assert_eq!(source.url, "http://example.com/audio.flac");
assert_eq!(source.chunk_frames, 0);
}
/// Test de création avec taille de chunk personnalisée
#[test]
fn test_http_source_with_chunk_size() {
let source = HttpSource::with_chunk_size("http://example.com/audio.mp3", 1024);
assert_eq!(source.url, "http://example.com/audio.mp3");
assert_eq!(source.chunk_frames, 1024);
}
/// Test de téléchargement et décodage d'un fichier FLAC via HTTP
#[tokio::test]
async fn test_http_source_downloads_and_decodes_flac() {
// Créer un serveur HTTP mock
let mock_server = MockServer::start().await;
// Générer un petit fichier FLAC de test
let sample_rate = 48_000;
let frames = 256;
let mut pcm = Vec::with_capacity(frames * 4);
for i in 0..frames {
let sample = ((i % 32) as f32 / 31.0 * 2.0 - 1.0) * 0.5;
let sample_i16 = (sample * 32767.0) as i16;
pcm.extend_from_slice(&sample_i16.to_le_bytes());
pcm.extend_from_slice(&sample_i16.to_le_bytes());
}
let format = PcmFormat {
sample_rate,
channels: 2,
bits_per_sample: 16,
};
let mut flac_stream =
encode_flac_stream(Cursor::new(pcm.clone()), format, EncoderOptions::default())
.await
.unwrap();
// Lire le FLAC encodé dans un buffer
let mut flac_data = Vec::new();
tokio::io::copy(&mut flac_stream, &mut flac_data)
.await
.unwrap();
flac_stream.wait().await.unwrap();
// Configurer le mock pour servir le fichier FLAC
Mock::given(method("GET"))
.and(path("/test.flac"))
.respond_with(
ResponseTemplate::new(200)
.set_body_bytes(flac_data)
.insert_header("content-type", "audio/flac"),
)
.mount(&mock_server)
.await;
// Créer la source HTTP pointant vers le mock
let url = format!("{}/test.flac", mock_server.uri());
let mut source = HttpSource::with_chunk_size(&url, 64);
let (tx, mut rx) = mpsc::channel(16);
source.add_subscriber(tx);
// Lancer le téléchargement et le décodage
tokio::spawn(async move {
source.run().await.unwrap();
});
// Vérifier les segments reçus
let mut received_frames = 0usize;
let mut seen_top_zero = false;
let mut seen_track_boundary = false;
let mut seen_eos = false;
while let Some(segment) = rx.recv().await {
if segment.is_audio_chunk() {
if let Some(chunk) = segment.as_chunk() {
received_frames += chunk.len();
assert_eq!(chunk.sample_rate(), sample_rate);
}
} else if let Some(marker) = segment.as_sync_marker() {
match **marker {
crate::SyncMarker::TopZeroSync => seen_top_zero = true,
crate::SyncMarker::TrackBoundary { .. } => seen_track_boundary = true,
crate::SyncMarker::EndOfStream => seen_eos = true,
_ => {}
}
}
}
// Vérifications
assert_eq!(received_frames, frames, "Tous les frames doivent être reçus");
assert!(seen_top_zero, "TopZeroSync doit être émis");
assert!(seen_track_boundary, "TrackBoundary doit être émis");
assert!(seen_eos, "EndOfStream doit être émis");
}
/// Test de l'extraction des métadonnées depuis les headers HTTP
#[tokio::test]
async fn test_http_source_extracts_icy_metadata() {
let mock_server = MockServer::start().await;
// Créer un fichier FLAC minimal
let sample_rate = 48_000;
let frames = 128;
let mut pcm = Vec::with_capacity(frames * 4);
for i in 0..frames {
let sample_i16 = ((i % 100) as i16) * 100;
pcm.extend_from_slice(&sample_i16.to_le_bytes());
pcm.extend_from_slice(&sample_i16.to_le_bytes());
}
let format = PcmFormat {
sample_rate,
channels: 2,
bits_per_sample: 16,
};
let mut flac_stream = encode_flac_stream(Cursor::new(pcm), format, EncoderOptions::default())
.await
.unwrap();
let mut flac_data = Vec::new();
tokio::io::copy(&mut flac_stream, &mut flac_data)
.await
.unwrap();
flac_stream.wait().await.unwrap();
// Configurer le mock avec headers Icecast
Mock::given(method("GET"))
.and(path("/stream"))
.respond_with(
ResponseTemplate::new(200)
.set_body_bytes(flac_data)
.insert_header("content-type", "audio/flac")
.insert_header("icy-name", "Test Radio Stream")
.insert_header("icy-url", "http://example.com/radio"),
)
.mount(&mock_server)
.await;
let url = format!("{}/stream", mock_server.uri());
let mut source = HttpSource::new(&url);
let (tx, mut rx) = mpsc::channel(16);
source.add_subscriber(tx);
tokio::spawn(async move {
source.run().await.unwrap();
});
// Chercher le TrackBoundary pour vérifier les métadonnées
let mut found_metadata = false;
while let Some(segment) = rx.recv().await {
if let Some(marker) = segment.as_sync_marker() {
if let crate::SyncMarker::TrackBoundary { metadata, .. } = &**marker {
// Vérifier que le titre extrait est "Test Radio Stream"
if let Some(title) = metadata.get_title().await.ok().flatten() {
assert_eq!(title, "Test Radio Stream");
found_metadata = true;
}
}
}
}
assert!(found_metadata, "Les métadonnées ICY doivent être extraites");
}
/// Test du comportement en cas d'erreur HTTP 404
#[tokio::test]
async fn test_http_source_handles_404_error() {
let mock_server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/notfound.flac"))
.respond_with(ResponseTemplate::new(404))
.mount(&mock_server)
.await;
let url = format!("{}/notfound.flac", mock_server.uri());
let mut source = HttpSource::new(&url);
let (tx, _rx) = mpsc::channel(16);
source.add_subscriber(tx);
let result = source.run().await;
assert!(result.is_err(), "Doit retourner une erreur pour HTTP 404");
if let Err(AudioError::ProcessingError(msg)) = result {
assert!(msg.contains("404"), "Le message d'erreur doit mentionner le code 404");
} else {
panic!("Le type d'erreur doit être ProcessingError");
}
}
/// Test du comportement avec un format audio invalide
#[tokio::test]
async fn test_http_source_handles_invalid_audio_format() {
let mock_server = MockServer::start().await;
// Envoyer des données invalides (pas un fichier audio)
Mock::given(method("GET"))
.and(path("/invalid.flac"))
.respond_with(
ResponseTemplate::new(200)
.set_body_bytes(b"This is not a valid audio file")
.insert_header("content-type", "audio/flac"),
)
.mount(&mock_server)
.await;
let url = format!("{}/invalid.flac", mock_server.uri());
let mut source = HttpSource::new(&url);
let (tx, _rx) = mpsc::channel(16);
source.add_subscriber(tx);
let result = source.run().await;
assert!(
result.is_err(),
"Doit retourner une erreur pour un format invalide"
);
}
/// Test de l'extraction du nom de fichier depuis l'URL quand pas de header icy-name
#[tokio::test]
async fn test_http_source_uses_filename_as_title() {
let mock_server = MockServer::start().await;
let sample_rate = 48_000;
let frames = 128;
let mut pcm = Vec::with_capacity(frames * 4);
for i in 0..frames {
let sample_i16 = (i % 100) as i16;
pcm.extend_from_slice(&sample_i16.to_le_bytes());
pcm.extend_from_slice(&sample_i16.to_le_bytes());
}
let format = PcmFormat {
sample_rate,
channels: 2,
bits_per_sample: 16,
};
let mut flac_stream = encode_flac_stream(Cursor::new(pcm), format, EncoderOptions::default())
.await
.unwrap();
let mut flac_data = Vec::new();
tokio::io::copy(&mut flac_stream, &mut flac_data)
.await
.unwrap();
flac_stream.wait().await.unwrap();
// Sans header icy-name
Mock::given(method("GET"))
.and(path("/my-song.flac"))
.respond_with(
ResponseTemplate::new(200)
.set_body_bytes(flac_data)
.insert_header("content-type", "audio/flac"),
)
.mount(&mock_server)
.await;
let url = format!("{}/my-song.flac", mock_server.uri());
let mut source = HttpSource::new(&url);
let (tx, mut rx) = mpsc::channel(16);
source.add_subscriber(tx);
tokio::spawn(async move {
source.run().await.unwrap();
});
let mut found_title = false;
while let Some(segment) = rx.recv().await {
if let Some(marker) = segment.as_sync_marker() {
if let crate::SyncMarker::TrackBoundary { metadata, .. } = &**marker {
if let Some(title) = metadata.get_title().await.ok().flatten() {
assert_eq!(title, "my-song.flac");
found_title = true;
}
}
}
}
assert!(found_title, "Le nom du fichier doit être utilisé comme titre");
}
}

224
pmoaudio/src/nodes/mod.rs Normal file
View File

@@ -0,0 +1,224 @@
//! Nodes du pipeline audio
//!
//! Ce module contient tous les types de nodes disponibles pour construire
//! un pipeline audio, ainsi que les traits et structures de support.
use std::sync::Arc;
use tokio::sync::mpsc;
use crate::type_constraints::{TypeMismatch, TypeRequirement};
use crate::AudioSegment;
/// Taille par défaut du buffer de channel MPSC pour les nodes
/// Cette valeur détermine combien de segments audio peuvent être mis en attente
/// avant que le producteur soit bloqué (backpressure).
pub const DEFAULT_CHANNEL_SIZE: usize = 16;
/// Durée par défaut des chunks audio en millisecondes
/// Cette valeur détermine la latence de traitement et le compromis efficacité/réactivité.
/// 50ms offre un bon équilibre pour la plupart des applications de lecture audio.
pub const DEFAULT_CHUNK_DURATION_MS: f64 = 50.0;
// Modules actifs
pub mod converter_nodes;
pub mod file_source;
pub mod flac_file_sink;
pub mod http_source;
// Modules temporairement désactivés
/*
pub mod buffer_node;
pub mod chromecast_sink;
pub mod decoder_node;
pub mod disk_sink;
pub mod dsp_node;
pub mod mpd_sink;
pub mod sink_node;
pub mod source_node;
pub mod timer_node;
pub mod volume_node;
*/
/// Trait de base pour tous les nodes audio
///
/// Tous les nodes du pipeline implémentent ce trait pour permettre
/// une interface uniforme de traitement des chunks audio.
#[async_trait::async_trait]
pub trait AudioNode: Send + Sync {
/// Push un chunk vers ce node
///
/// # Erreurs
///
/// Retourne `AudioError::SendError` si l'envoi échoue
async fn push(&mut self, chunk: Arc<AudioSegment>) -> Result<(), AudioError>;
/// Ferme le node proprement
async fn close(&mut self);
}
/// Trait pour les nodes qui déclarent leurs types acceptés/produits
///
/// Ce trait permet de vérifier la compatibilité des types entre nodes
/// avant de les connecter dans un pipeline.
///
/// # Exemples
///
/// ```no_run
/// use pmoaudio::{FileSource, FlacFileSink, TypedAudioNode};
/// use pmoaudio::type_constraints::check_compatibility;
///
/// // Vérifier la compatibilité avant de connecter
/// let source = FileSource::new("input.flac");
/// let (sink, tx) = FlacFileSink::new("output.flac");
///
/// let source_output = source.output_type();
/// let sink_input = sink.input_type();
///
/// match check_compatibility(&source_output, &sink_input) {
/// Ok(()) => println!("Types compatibles!"),
/// Err(e) => eprintln!("Types incompatibles: {}", e),
/// }
/// ```
pub trait TypedAudioNode {
/// Retourne les types que ce node peut accepter en entrée
///
/// Pour les sources (qui ne consomment rien), retourne `None`.
fn input_type(&self) -> Option<TypeRequirement>;
/// Retourne les types que ce node peut produire en sortie
///
/// Pour les sinks (qui ne produisent rien), retourne `None`.
fn output_type(&self) -> Option<TypeRequirement>;
/// Vérifie si ce node peut accepter les chunks d'un producer donné
///
/// # Erreurs
///
/// Retourne `AudioError::TypeMismatch` si les types sont incompatibles
fn can_accept_from(&self, producer: &dyn TypedAudioNode) -> Result<(), AudioError> {
match (producer.output_type(), self.input_type()) {
(Some(prod), Some(cons)) => crate::type_constraints::check_compatibility(&prod, &cons)
.map_err(|e| AudioError::TypeMismatch(e)),
(None, Some(_)) => Err(AudioError::TypeMismatch(TypeMismatch {
producer: TypeRequirement::any(), // Placeholder
consumer: self.input_type().unwrap(),
incompatible_type: None,
})),
_ => Ok(()), // Si pas de contrainte, toujours compatible
}
}
}
/// Node avec un seul abonné (pas de clone inutile)
///
/// Optimisé pour les cas où un node n'a qu'un seul destinataire.
/// Le Arc du chunk est simplement transféré sans clonage supplémentaire.
///
/// # Exemples
///
/// ```
/// use pmoaudio::SingleSubscriberNode;
/// use tokio::sync::mpsc;
///
/// let (tx, rx) = mpsc::channel(10);
/// let node = SingleSubscriberNode::new(tx);
/// ```
pub struct SingleSubscriberNode {
tx: mpsc::Sender<Arc<AudioSegment>>,
}
impl SingleSubscriberNode {
pub fn new(tx: mpsc::Sender<Arc<AudioSegment>>) -> Self {
Self { tx }
}
pub async fn push(&self, chunk: Arc<AudioSegment>) -> Result<(), AudioError> {
self.tx.send(chunk).await.map_err(|_| AudioError::SendError)
}
}
/// Node avec plusieurs abonnés (partage le même Arc)
///
/// Permet de broadcaster un chunk à plusieurs destinations.
/// Tous les abonnés reçoivent le même `Arc<AudioSegment>`, donc pas de copie
/// des données audio - seul le compteur de référence Arc est incrémenté.
///
/// # Exemples
///
/// ```
/// use pmoaudio::MultiSubscriberNode;
/// use tokio::sync::mpsc;
///
/// let mut node = MultiSubscriberNode::new();
/// let (tx1, rx1) = mpsc::channel(10);
/// let (tx2, rx2) = mpsc::channel(10);
///
/// node.add_subscriber(tx1);
/// node.add_subscriber(tx2);
/// // Les deux abonnés recevront les mêmes chunks
/// ```
pub struct MultiSubscriberNode {
subscribers: Vec<mpsc::Sender<Arc<AudioSegment>>>,
}
impl MultiSubscriberNode {
pub fn new() -> Self {
Self {
subscribers: Vec::new(),
}
}
pub fn add_subscriber(&mut self, tx: mpsc::Sender<Arc<AudioSegment>>) {
self.subscribers.push(tx);
}
pub async fn push(&self, chunk: Arc<AudioSegment>) -> Result<(), AudioError> {
for tx in &self.subscribers {
// On partage le même Arc avec tous les abonnés
tx.send(chunk.clone())
.await
.map_err(|_| AudioError::SendError)?;
}
Ok(())
}
pub async fn try_push(&self, chunk: Arc<AudioSegment>) -> Result<(), AudioError> {
for tx in &self.subscribers {
// try_send non-bloquant, ignore si saturé
let _ = tx.try_send(chunk.clone());
}
Ok(())
}
}
impl Default for MultiSubscriberNode {
fn default() -> Self {
Self::new()
}
}
/// Erreurs possibles dans le pipeline audio
#[derive(Debug, Clone)]
pub enum AudioError {
/// Échec d'envoi d'un chunk à travers un channel
SendError,
/// Échec de réception d'un chunk depuis un channel
ReceiveError,
/// Erreur de traitement avec message descriptif
ProcessingError(String),
/// Incompatibilité de types entre nodes
TypeMismatch(TypeMismatch),
}
impl std::fmt::Display for AudioError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
AudioError::SendError => write!(f, "Failed to send audio chunk"),
AudioError::ReceiveError => write!(f, "Failed to receive audio chunk"),
AudioError::ProcessingError(msg) => write!(f, "Processing error: {}", msg),
AudioError::TypeMismatch(tm) => write!(f, "{}", tm),
}
}
}
impl std::error::Error for AudioError {}