diff --git a/Cargo.lock b/Cargo.lock index f7bd9ded..fe7c03f8 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4,7 +4,7 @@ version = 4 [[package]] name = "PMOMusic" -version = "0.3.22" +version = "0.3.23" dependencies = [ "axum 0.8.7", "console-subscriber", @@ -1700,6 +1700,24 @@ dependencies = [ "simd-adler32", ] +[[package]] +name = "fdk-aac" +version = "0.8.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7bb67e142688083cb9afb63f2424203fc98c4e7afb494bf912b60b55513b177e" +dependencies = [ + "fdk-aac-sys", +] + +[[package]] +name = "fdk-aac-sys" +version = "0.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "24516d2611506d5cb1833555adc75f6baf9fe2706b9c13e6fc33a6b22c51ca83" +dependencies = [ + "cc", +] + [[package]] name = "find-msvc-tools" version = "0.1.5" @@ -4043,6 +4061,7 @@ version = "0.1.0" dependencies = [ "bytes", "claxon", + "fdk-aac", "lewton", "libc", "libflac-sys", diff --git a/PMOMusic/Cargo.toml b/PMOMusic/Cargo.toml index 1f499509..07676aaf 100644 --- a/PMOMusic/Cargo.toml +++ b/PMOMusic/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "PMOMusic" -version = "0.3.22" +version = "0.3.23" edition = "2024" [dependencies] diff --git a/pmoaudiocache/src/streaming.rs b/pmoaudiocache/src/streaming.rs index 53d98779..c93ddc3e 100644 --- a/pmoaudiocache/src/streaming.rs +++ b/pmoaudiocache/src/streaming.rs @@ -125,6 +125,7 @@ fn codec_to_string(codec: AudioCodec) -> String { AudioCodec::OggOpus => "ogg_opus", AudioCodec::Wav => "wav", AudioCodec::Aiff => "aiff", + AudioCodec::Aac => "aac", } .to_string() } diff --git a/pmoflac/Cargo.toml b/pmoflac/Cargo.toml index 9919892e..c23460ea 100755 --- a/pmoflac/Cargo.toml +++ b/pmoflac/Cargo.toml @@ -19,6 +19,7 @@ libflac-sys = { version = "0.3.3", default-features = false, features = ["build- lofty = "0.22" minimp3 = "0.5" opus = "0.3" +fdk-aac = "0.8" pmometadata = { path = "../pmometadata" } thiserror = { workspace = true } tokio = { workspace = true, features = ["rt", "rt-multi-thread", "macros", "sync", "io-util", "fs"] } diff --git a/pmoflac/src/aac.rs b/pmoflac/src/aac.rs new file mode 100644 index 00000000..67a71184 --- /dev/null +++ b/pmoflac/src/aac.rs @@ -0,0 +1,177 @@ +//! AAC Decoder Module (ADTS streaming) +//! +//! Provides asynchronous streaming AAC/ADTS decoding via libfdk-aac (statically linked). +//! Decodes AAC audio streams into PCM data (16-bit little-endian interleaved). +//! +//! Designed for live radio streams (e.g. Radio France icecast AAC 192kbps). +//! No seek required — pure linear streaming. +//! +//! ## Architecture +//! +//! ```text +//! AAC Input → [Ingest Task] → [Decode Task (blocking)] → [Writer Task] → PCM Output (AsyncRead) +//! ``` + +use fdk_aac::dec::{Decoder, DecoderError, Transport}; +use tokio::{ + io::AsyncRead, + sync::{mpsc, oneshot}, +}; + +use crate::{ + common::ChannelReader, + decoder_common::{ + spawn_ingest_task, spawn_writer_task, DecodedStream, DecoderError as PmoDecoderError, + CHANNEL_CAPACITY, DUPLEX_BUFFER_SIZE, + }, + pcm::StreamInfo, + stream::ManagedAsyncReader, +}; + +/// Errors that can occur while decoding AAC data. +pub type AacError = PmoDecoderError; + +/// Async decoded AAC stream. +pub type AacDecodedStream = DecodedStream; + +/// Decodes an AAC/ADTS stream into PCM audio data (16-bit little-endian interleaved). +/// +/// The input must be a raw ADTS stream (as produced by Radio France icecast). +/// No seek is required — the decoder processes frames linearly. +pub async fn decode_aac_stream(reader: R) -> Result +where + R: AsyncRead + Unpin + Send + 'static, +{ + let (ingest_tx, ingest_rx) = mpsc::channel(CHANNEL_CAPACITY); + spawn_ingest_task(reader, ingest_tx); + + let (pcm_tx, pcm_rx) = mpsc::channel(CHANNEL_CAPACITY); + let (pcm_reader, pcm_writer) = tokio::io::duplex(DUPLEX_BUFFER_SIZE); + let (info_tx, info_rx) = oneshot::channel::>(); + + let blocking_handle = tokio::task::spawn_blocking(move || -> Result<(), AacError> { + let mut channel_reader = ChannelReader::::new(ingest_rx); + let mut decoder = Decoder::new(Transport::Adts); + let mut info_tx = Some(info_tx); + + // Buffer de lecture — on lit des chunks et on les pousse au décodeur + let mut read_buf = vec![0u8; 8192]; + // Buffer de sortie PCM — fdk-aac écrit des frames entières + let mut pcm_out = vec![0i16; 8192]; + + use std::io::Read; + + loop { + // Lire des bytes depuis le stream ADTS + let n = match channel_reader.read(&mut read_buf) { + Ok(0) => break, // EOF + Ok(n) => n, + Err(e) => { + let err = AacError::Io { + kind: e.kind(), + message: e.to_string(), + }; + if let Some(tx) = info_tx.take() { + let _ = tx.send(Err(err.clone())); + } + return Err(err); + } + }; + + // Pousser les bytes au décodeur fdk-aac + let filled = match decoder.fill(&read_buf[..n]) { + Ok(filled) => filled, + Err(e) => { + let err = AacError::Decode(format!("fdk-aac fill error: {:?}", e)); + if let Some(tx) = info_tx.take() { + let _ = tx.send(Err(err.clone())); + } + return Err(err); + } + }; + + // Si fill n'a pas consommé tous les bytes, on les remet devant + // (fdk-aac peut ne pas consommer tout en une passe) + // Note: fdk-aac retourne le nombre de bytes non-consommés — on les ignore + // car le décodeur conserve son état interne entre appels. + let _ = filled; + + // Décoder les frames disponibles + loop { + // Adapter la taille du buffer PCM si nécessaire + let stream_info = decoder.stream_info(); + let frame_size = if stream_info.frameSize > 0 && stream_info.numChannels > 0 { + (stream_info.frameSize * stream_info.numChannels) as usize + } else { + 2048 // taille par défaut avant la première frame + }; + + if pcm_out.len() < frame_size { + pcm_out.resize(frame_size, 0i16); + } + + match decoder.decode_frame(&mut pcm_out) { + Ok(()) => { + let info_ref = decoder.stream_info(); + + // Envoyer les infos au premier décodage réussi + if let Some(tx) = info_tx.take() { + let info = StreamInfo { + sample_rate: info_ref.sampleRate as u32, + channels: info_ref.numChannels as u8, + bits_per_sample: 16, + total_samples: None, + max_block_size: 0, + min_block_size: 0, + }; + if tx.send(Ok(info)).is_err() { + return Ok(()); + } + } + + // Convertir i16 → bytes little-endian + let samples = &pcm_out[..frame_size]; + let mut bytes = Vec::with_capacity(samples.len() * 2); + for s in samples { + bytes.extend_from_slice(&s.to_le_bytes()); + } + + if pcm_tx.blocking_send(Ok(bytes)).is_err() { + return Ok(()); + } + } + Err(DecoderError::NOT_ENOUGH_BITS) => { + // Pas assez de données — besoin de plus d'input + break; + } + Err(DecoderError::TRANSPORT_SYNC_ERROR) => { + // Erreur de sync ADTS transitoire — continuer + break; + } + Err(e) => { + let err = AacError::Decode(format!("fdk-aac decode error: {:?}", e)); + if let Some(tx) = info_tx.take() { + let _ = tx.send(Err(err.clone())); + } + return Err(err); + } + } + } + } + + if let Some(tx) = info_tx.take() { + let err = AacError::Decode("AAC stream contained no decodable frames".into()); + let _ = tx.send(Err(err.clone())); + return Err(err); + } + + Ok(()) + }); + + let writer_handle = spawn_writer_task(pcm_rx, pcm_writer, blocking_handle, "aac-decode"); + + let info = info_rx.await.map_err(|_| AacError::ChannelClosed)??; + let reader = ManagedAsyncReader::new("aac-decode-writer", pcm_reader, writer_handle); + + Ok(DecodedStream::new(info, reader)) +} diff --git a/pmoflac/src/autodetect.rs b/pmoflac/src/autodetect.rs index dd1f551f..1c5f6a61 100755 --- a/pmoflac/src/autodetect.rs +++ b/pmoflac/src/autodetect.rs @@ -7,6 +7,7 @@ use std::{ use tokio::io::{AsyncRead, AsyncReadExt, ReadBuf}; use crate::{ + aac::{decode_aac_stream, AacDecodedStream, AacError}, decode_aiff_stream, decode_flac_stream, decode_mp3_stream, decode_ogg_opus_stream, decode_ogg_vorbis_stream, decode_wav_stream, pcm::StreamInfo, prefixed_reader::PrefixedReader, AiffDecodedStream, AiffError, FlacDecodedStream, FlacError, Mp3DecodedStream, Mp3Error, @@ -34,6 +35,8 @@ pub enum DecodeAudioError { Wav(WavError), #[error("AIFF decode error: {0}")] Aiff(AiffError), + #[error("AAC decode error: {0}")] + Aac(AacError), } pub async fn decode_audio_stream(reader: R) -> Result @@ -94,6 +97,12 @@ where .map_err(DecodeAudioError::Aiff)?; DecodedAudioStream::Aiff(stream) } + DetectedFormat::Aac => { + let stream = decode_aac_stream(prefixed) + .await + .map_err(DecodeAudioError::Aac)?; + DecodedAudioStream::Aac(stream) + } }; Ok(stream) @@ -106,6 +115,7 @@ pub enum DecodedAudioStream { OggOpus(OggOpusDecodedStream), Wav(WavDecodedStream), Aiff(AiffDecodedStream), + Aac(AacDecodedStream), } impl DecodedAudioStream { @@ -117,6 +127,7 @@ impl DecodedAudioStream { DecodedAudioStream::OggOpus(inner) => inner.info(), DecodedAudioStream::Wav(inner) => inner.info(), DecodedAudioStream::Aiff(inner) => inner.info(), + DecodedAudioStream::Aac(inner) => inner.info(), } } @@ -132,6 +143,7 @@ impl DecodedAudioStream { } DecodedAudioStream::Wav(inner) => inner.wait().await.map_err(DecodeAudioError::Wav), DecodedAudioStream::Aiff(inner) => inner.wait().await.map_err(DecodeAudioError::Aiff), + DecodedAudioStream::Aac(inner) => inner.wait().await.map_err(DecodeAudioError::Aac), } } @@ -161,6 +173,10 @@ impl DecodedAudioStream { let (info, reader) = inner.into_parts(); (info, DecodedReader::Aiff(reader)) } + DecodedAudioStream::Aac(inner) => { + let (info, reader) = inner.into_parts(); + (info, DecodedReader::Aac(reader)) + } } } } @@ -178,6 +194,7 @@ impl AsyncRead for DecodedAudioStream { DecodedAudioStream::OggOpus(inner) => Pin::new(inner).poll_read(cx, buf), DecodedAudioStream::Wav(inner) => Pin::new(inner).poll_read(cx, buf), DecodedAudioStream::Aiff(inner) => Pin::new(inner).poll_read(cx, buf), + DecodedAudioStream::Aac(inner) => Pin::new(inner).poll_read(cx, buf), } } } @@ -189,6 +206,7 @@ pub enum DecodedReader { OggOpus(crate::stream::ManagedAsyncReader), Wav(crate::stream::ManagedAsyncReader), Aiff(crate::stream::ManagedAsyncReader), + Aac(crate::stream::ManagedAsyncReader), } impl DecodedReader { @@ -200,6 +218,7 @@ impl DecodedReader { DecodedReader::OggOpus(inner) => inner.wait().await.map_err(DecodeAudioError::Opus), DecodedReader::Wav(inner) => inner.wait().await.map_err(DecodeAudioError::Wav), DecodedReader::Aiff(inner) => inner.wait().await.map_err(DecodeAudioError::Aiff), + DecodedReader::Aac(inner) => inner.wait().await.map_err(DecodeAudioError::Aac), } } } @@ -217,6 +236,7 @@ impl AsyncRead for DecodedReader { DecodedReader::OggOpus(inner) => Pin::new(inner).poll_read(cx, buf), DecodedReader::Wav(inner) => Pin::new(inner).poll_read(cx, buf), DecodedReader::Aiff(inner) => Pin::new(inner).poll_read(cx, buf), + DecodedReader::Aac(inner) => Pin::new(inner).poll_read(cx, buf), } } } @@ -242,12 +262,29 @@ fn detect_format(bytes: &[u8]) -> Option { if let Some(fmt) = detect_ogg(bytes) { return Some(fmt); } + // Tester AAC avant MP3 : ADTS (0xFF 0xF0/0xF8) est un sous-ensemble du syncword MP3 + // (0xFF 0xEx), donc la détection MP3 absorberait les flux AAC si elle passait en premier. + if is_adts_aac(bytes) { + return Some(DetectedFormat::Aac); + } if is_mp3(bytes) { return Some(DetectedFormat::Mp3); } None } +/// Détecte un flux AAC ADTS : syncword 0xFFF (12 bits) + layer = 00. +/// Structure du 2e octet : 1111 VLLL → V=version, L=layer(00 pour AAC), P=protection +/// MP3 a layer != 00 (01=III, 10=II, 11=I), AAC ADTS a layer == 00. +fn is_adts_aac(bytes: &[u8]) -> bool { + if bytes.len() < 2 { + return false; + } + // syncword = 0xFFF (12 bits) : byte0=0xFF, byte1[7:4]=0xF + // layer = byte1[2:1] == 0b00 (distingue AAC de MP3) + bytes[0] == 0xFF && (bytes[1] & 0xF6) == 0xF0 +} + fn detect_ogg(bytes: &[u8]) -> Option { if bytes.len() < 27 || &bytes[..4] != b"OggS" { return None; @@ -294,4 +331,5 @@ enum DetectedFormat { OggOpus, Wav, Aiff, + Aac, } diff --git a/pmoflac/src/lib.rs b/pmoflac/src/lib.rs index aeaf58e3..a74a741a 100755 --- a/pmoflac/src/lib.rs +++ b/pmoflac/src/lib.rs @@ -95,6 +95,7 @@ //! } //! ``` +pub mod aac; pub mod aiff; pub mod autodetect; mod common; @@ -114,6 +115,7 @@ pub mod transcode; mod util; pub mod wav; +pub use aac::{decode_aac_stream, AacDecodedStream, AacError}; pub use aiff::{decode_aiff_stream, AiffDecodedStream, AiffError}; pub use autodetect::{ decode_audio_stream, is_flac_magic_header, DecodeAudioError, DecodedAudioStream, DecodedReader, diff --git a/pmoflac/src/transcode.rs b/pmoflac/src/transcode.rs index 685da1fe..822f0be3 100755 --- a/pmoflac/src/transcode.rs +++ b/pmoflac/src/transcode.rs @@ -31,6 +31,7 @@ pub enum AudioCodec { OggOpus, Wav, Aiff, + Aac, } /// Options controlling how the transcoder operates. @@ -195,6 +196,9 @@ where DecodedAudioStream::Aiff(stream) => { transcode_from_decoded(AudioCodec::Aiff, stream, options.encoder_options).await } + DecodedAudioStream::Aac(stream) => { + transcode_from_decoded(AudioCodec::Aac, stream, options.encoder_options).await + } } } diff --git a/pmoflac/tests/aac_decode.rs b/pmoflac/tests/aac_decode.rs new file mode 100644 index 00000000..5f7c6ab7 --- /dev/null +++ b/pmoflac/tests/aac_decode.rs @@ -0,0 +1,43 @@ +use pmoflac::decode_aac_stream; +use tokio::io::AsyncReadExt; + +#[tokio::test] +async fn test_decode_adts_file() { + let file = tokio::fs::File::open("/tmp/test_adts.aac") + .await + .expect("test ADTS file not found — run: ffmpeg -i tests/SBRtestStereoHiBr.mp4 -vn -acodec copy -f adts /tmp/test_adts.aac"); + + let mut stream = decode_aac_stream(file) + .await + .expect("decode_aac_stream failed"); + + let info = stream.info().clone(); + println!("StreamInfo: {} Hz, {} ch, {} bps", info.sample_rate, info.channels, info.bits_per_sample); + + assert!(info.sample_rate > 0, "sample_rate should be > 0"); + assert!(info.channels == 1 || info.channels == 2, "channels should be 1 or 2"); + assert_eq!(info.bits_per_sample, 16); + + let mut pcm = Vec::new(); + stream.read_to_end(&mut pcm).await.expect("read_to_end failed"); + + println!("Decoded {} PCM bytes ({} samples)", pcm.len(), pcm.len() / 2); + assert!(pcm.len() > 0, "should have decoded some PCM data"); +} + +#[tokio::test] +async fn test_autodetect_adts() { + use pmoflac::decode_audio_stream; + + let file = tokio::fs::File::open("/tmp/test_adts.aac") + .await + .expect("test ADTS file not found"); + + let stream = decode_audio_stream(file) + .await + .expect("decode_audio_stream failed"); + + let info = stream.info().clone(); + println!("Autodetect: {} Hz, {} ch, {} bps", info.sample_rate, info.channels, info.bits_per_sample); + assert!(info.sample_rate > 0); +} diff --git a/version.txt b/version.txt index 0c4b4549..a1dad2aa 100644 --- a/version.txt +++ b/version.txt @@ -1 +1 @@ -0.3.22 +0.3.23