use anyhow::Result; use bytes::Bytes; use futures::stream::Stream; use std::io::{self, Read}; use std::pin::Pin; use std::sync::mpsc::{sync_channel, Receiver, RecvError, SyncSender}; const CHANNEL_BUFFER_SIZE: usize = 16; pub const CHUNK_SIZE_FRAMES: usize = 4096; pub struct ChannelReader { receiver: Receiver>, current_chunk: Option, position: usize, } impl ChannelReader { pub fn new( stream: Pin> + Send>>, ) -> Self { let (tx, rx) = sync_channel(CHANNEL_BUFFER_SIZE); tokio::spawn(Self::stream_feeder(stream, tx)); Self { receiver: rx, current_chunk: None, position: 0, } } async fn stream_feeder( mut stream: Pin> + Send>>, tx: SyncSender>, ) { use futures::StreamExt; while let Some(result) = stream.next().await { let to_send = result.map_err(|e| e.to_string()); if tx.send(to_send).is_err() { break; } } } } impl Read for ChannelReader { fn read(&mut self, buf: &mut [u8]) -> io::Result { loop { if let Some(chunk) = &self.current_chunk { if self.position < chunk.len() { let available = chunk.len() - self.position; let to_copy = available.min(buf.len()); buf[..to_copy].copy_from_slice(&chunk[self.position..self.position + to_copy]); self.position += to_copy; return Ok(to_copy); } } match self.receiver.recv() { Ok(Ok(bytes)) => { self.current_chunk = Some(bytes); self.position = 0; } Ok(Err(e)) => return Err(io::Error::new(io::ErrorKind::Other, e)), Err(RecvError) => return Ok(0), } } } } #[derive(Debug, Clone)] pub struct PCMChunk { pub samples: Vec, pub position_ms: u64, pub sample_rate: u32, pub channels: u32, } pub struct StreamingPCMDecoder { reader: claxon::FlacReader>, sample_rate: u32, channels: u32, bits_per_sample: u32, total_samples_decoded: u64, done: bool, } impl StreamingPCMDecoder { /// Create a new decoder from an HTTP stream with default chunk size pub fn new(http_stream: crate::stream::BlockStream) -> anyhow::Result { Self::with_chunk_size(http_stream, CHUNK_SIZE_FRAMES) } pub fn with_chunk_size( http_stream: crate::stream::BlockStream, _chunk_size: usize, ) -> anyhow::Result { let channel_reader = ChannelReader::new(http_stream.into_inner()); let buffered = std::io::BufReader::new(channel_reader); let reader = claxon::FlacReader::new(buffered) .map_err(|e| anyhow::anyhow!("FLAC reader error: {}", e))?; let info = reader.streaminfo(); Ok(Self { reader, sample_rate: info.sample_rate, channels: info.channels, bits_per_sample: info.bits_per_sample, total_samples_decoded: 0, done: false, }) } /// Get the sample rate (e.g., 44100 Hz) pub fn sample_rate(&self) -> u32 { self.sample_rate } /// Get the number of channels (e.g., 2 for stereo) pub fn channels(&self) -> u32 { self.channels } /// Get bits per sample (e.g., 16) pub fn bits_per_sample(&self) -> u32 { self.bits_per_sample } pub fn decode_chunk(&mut self) -> anyhow::Result> { if self.done { return Ok(None); } // Crée le FrameReader à la volée (emprunt de self.reader) let mut frames = self.reader.blocks(); // API claxon 0.6.x : il FAUT fournir un Vec par valeur let buf: Vec = Vec::new(); let frame = match frames.read_next_or_eof(buf) { Ok(None) => { self.done = true; return Ok(None); } Ok(Some(f)) => f, Err(e) => return Err(anyhow::anyhow!("FLAC decode error: {}", e)), }; let samples: Vec = frame.into_buffer(); if samples.is_empty() { self.done = true; return Ok(None); } let position_ms = { let frames = self.total_samples_decoded / self.channels as u64; (frames * 1000) / self.sample_rate as u64 }; self.total_samples_decoded += samples.len() as u64; Ok(Some(PCMChunk { samples, position_ms, sample_rate: self.sample_rate, channels: self.channels, })) } } pub fn ms_to_frames(ms: u64, sample_rate: u32) -> usize { ((ms as u128 * sample_rate as u128) / 1000) as usize } pub fn frames_to_ms(frames: usize, sample_rate: u32) -> u64 { ((frames as u128 * 1000) / sample_rate as u128) as u64 }