diff --git a/pmoaudio-ext/src/sinks/streaming_flac_sink.rs b/pmoaudio-ext/src/sinks/streaming_flac_sink.rs index d8e4af22..9f8e1df7 100644 --- a/pmoaudio-ext/src/sinks/streaming_flac_sink.rs +++ b/pmoaudio-ext/src/sinks/streaming_flac_sink.rs @@ -86,6 +86,20 @@ const DEFAULT_ICY_METAINT: usize = 16000; /// while keeping metadata synchronized (larger buffers cause metadata drift). const BROADCAST_CAPACITY: usize = 128; +/// Maximum lead time for HTTP broadcast pacing (in seconds). +/// The broadcaster will sleep if it's ahead of real-time by more than this amount. +/// This is much smaller than the pipeline TimerNode's 3.0s to provide tighter control. +const BROADCAST_MAX_LEAD_TIME: f64 = 0.5; + +/// PCM chunk with audio data and timestamp for precise pacing. +#[derive(Debug)] +struct PcmChunk { + /// Raw PCM audio bytes + bytes: Vec, + /// Timestamp in seconds (from AudioSegment) + timestamp_sec: f64, +} + /// Snapshot of track metadata at a point in time. /// /// This structure is shared between the sink and clients to provide @@ -511,8 +525,8 @@ struct EncoderState { struct StreamingFlacSinkLogic { encoder_options: EncoderOptions, bits_per_sample: u8, - pcm_tx: mpsc::Sender>, - pcm_rx: Option>>, + pcm_tx: mpsc::Sender, + pcm_rx: Option>, metadata: Arc>, flac_broadcast: broadcast::Sender, flac_header: Arc>>, @@ -534,8 +548,11 @@ impl StreamingFlacSinkLogic { AudioError::ProcessingError("PCM receiver already consumed".into()) })?; + // Create shared timestamp for pacing + let current_timestamp = Arc::new(RwLock::new(0.0f64)); + // Create ByteStreamReader for the encoder - let pcm_reader = ByteStreamReader::new(pcm_rx); + let pcm_reader = ByteStreamReader::new(pcm_rx, current_timestamp.clone()); // Create PCM format let pcm_format = PcmFormat { @@ -551,11 +568,11 @@ impl StreamingFlacSinkLogic { info!("FLAC encoder initialized successfully"); - // Spawn broadcaster task + // Spawn broadcaster task with timestamp for pacing let flac_broadcast = self.flac_broadcast.clone(); let flac_header = self.flac_header.clone(); let broadcaster_task = tokio::spawn(async move { - if let Err(e) = broadcast_flac_stream(flac_stream, flac_broadcast, flac_header).await { + if let Err(e) = broadcast_flac_stream(flac_stream, flac_broadcast, flac_header, current_timestamp).await { error!("Broadcaster task error: {}", e); } }); @@ -659,8 +676,12 @@ impl NodeLogic for StreamingFlacSinkLogic { seg.timestamp_sec ); - // Send to FLAC encoder - if let Err(e) = self.pcm_tx.send(pcm_bytes).await { + // Send to FLAC encoder with timestamp + let pcm_chunk = PcmChunk { + bytes: pcm_bytes, + timestamp_sec: seg.timestamp_sec, + }; + if let Err(e) = self.pcm_tx.send(pcm_chunk).await { warn!("Failed to send PCM data to encoder: {}", e); break; } @@ -707,16 +728,19 @@ impl NodeLogic for StreamingFlacSinkLogic { } /// Broadcaster task: reads FLAC bytes from encoder and broadcasts to all clients. +/// Implements precise real-time pacing based on audio timestamps. async fn broadcast_flac_stream( mut flac_stream: FlacEncodedStream, broadcast_tx: broadcast::Sender, header_cache: Arc>>, + current_timestamp: Arc>, ) -> Result<(), AudioError> { - info!("Broadcaster task started"); + info!("Broadcaster task started with precise timestamp-based pacing"); let mut buffer = vec![0u8; 8192]; // 8KB buffer for reading let mut total_bytes = 0u64; let mut header_captured = false; + let start_time = std::time::Instant::now(); loop { match flac_stream.read(&mut buffer).await { @@ -731,6 +755,20 @@ async fn broadcast_flac_stream( info!("Read {} bytes from FLAC encoder (total: {})", n, total_bytes); } + // Precise pacing based on audio timestamp + let audio_timestamp = *current_timestamp.read().await; + let elapsed = start_time.elapsed().as_secs_f64(); + let lead_time = audio_timestamp - elapsed; + + if lead_time > BROADCAST_MAX_LEAD_TIME { + let sleep_duration = lead_time - BROADCAST_MAX_LEAD_TIME; + debug!( + "Broadcaster pacing: sleeping {:.3}s (audio_ts={:.3}s, elapsed={:.3}s, lead={:.3}s)", + sleep_duration, audio_timestamp, elapsed, lead_time + ); + tokio::time::sleep(tokio::time::Duration::from_secs_f64(sleep_duration)).await; + } + // Broadcast to all clients let bytes = Bytes::copy_from_slice(&buffer[..n]); @@ -800,7 +838,7 @@ impl StreamingFlacSink { } // Create PCM channel (bounded for backpressure) - let (pcm_tx, pcm_rx) = mpsc::channel::>(16); + let (pcm_tx, pcm_rx) = mpsc::channel::(16); // Shared metadata let metadata = Arc::new(RwLock::new(MetadataSnapshot::default())); @@ -965,19 +1003,23 @@ fn chunk_to_pcm_bytes(chunk: &AudioChunk, bits_per_sample: u8) -> Result Ok(bytes) } -/// AsyncRead adapter for mpsc::Receiver>. +/// AsyncRead adapter for mpsc::Receiver. +/// Extracts bytes from PcmChunk and provides them to the FLAC encoder. struct ByteStreamReader { - rx: mpsc::Receiver>, + rx: mpsc::Receiver, buffer: VecDeque, finished: bool, + /// Shared timestamp for broadcaster pacing + current_timestamp: Arc>, } impl ByteStreamReader { - fn new(rx: mpsc::Receiver>) -> Self { + fn new(rx: mpsc::Receiver, current_timestamp: Arc>) -> Self { Self { rx, buffer: VecDeque::new(), finished: false, + current_timestamp, } } } @@ -1006,11 +1048,15 @@ impl AsyncRead for ByteStreamReader { } match Pin::new(&mut self.rx).poll_recv(cx) { - Poll::Ready(Some(bytes)) => { - if bytes.is_empty() { + Poll::Ready(Some(chunk)) => { + if chunk.bytes.is_empty() { continue; } - self.buffer.extend(bytes); + // Update shared timestamp for broadcaster pacing + if let Ok(mut ts) = self.current_timestamp.try_write() { + *ts = chunk.timestamp_sec; + } + self.buffer.extend(chunk.bytes); } Poll::Ready(None) => { self.finished = true; diff --git a/pmoaudio-ext/src/sinks/streaming_ogg_flac_sink.rs b/pmoaudio-ext/src/sinks/streaming_ogg_flac_sink.rs index fa03f727..20053c40 100644 --- a/pmoaudio-ext/src/sinks/streaming_ogg_flac_sink.rs +++ b/pmoaudio-ext/src/sinks/streaming_ogg_flac_sink.rs @@ -71,6 +71,19 @@ use tracing::{debug, error, info, trace, warn}; /// Same as StreamingFlacSink for consistency. const BROADCAST_CAPACITY: usize = 128; +/// Maximum lead time for HTTP broadcast pacing (in seconds). +/// The broadcaster will sleep if it's ahead of real-time by more than this amount. +const BROADCAST_MAX_LEAD_TIME: f64 = 0.5; + +/// PCM chunk with audio data and timestamp for precise pacing. +#[derive(Debug)] +struct PcmChunk { + /// Raw PCM audio bytes + bytes: Vec, + /// Timestamp in seconds (from AudioSegment) + timestamp_sec: f64, +} + /// Snapshot of track metadata (reuse from streaming_flac_sink) pub use super::streaming_flac_sink::MetadataSnapshot; @@ -227,8 +240,8 @@ struct EncoderState { struct StreamingOggFlacSinkLogic { encoder_options: EncoderOptions, bits_per_sample: u8, - pcm_tx: mpsc::Sender>, - pcm_rx: Option>>, + pcm_tx: mpsc::Sender, + pcm_rx: Option>, metadata: Arc>, ogg_broadcast: broadcast::Sender, ogg_header: Arc>>, @@ -250,8 +263,11 @@ impl StreamingOggFlacSinkLogic { AudioError::ProcessingError("PCM receiver already consumed".into()) })?; + // Create shared timestamp for pacing + let current_timestamp = Arc::new(RwLock::new(0.0f64)); + // Create ByteStreamReader for the encoder - let pcm_reader = ByteStreamReader::new(pcm_rx); + let pcm_reader = ByteStreamReader::new(pcm_rx, current_timestamp.clone()); // Create PCM format let pcm_format = PcmFormat { @@ -267,11 +283,11 @@ impl StreamingOggFlacSinkLogic { info!("OGG-FLAC encoder initialized successfully"); - // Spawn OGG wrapper + broadcaster task + // Spawn OGG wrapper + broadcaster task with timestamp for pacing let ogg_broadcast = self.ogg_broadcast.clone(); let ogg_header = self.ogg_header.clone(); let broadcaster_task = tokio::spawn(async move { - if let Err(e) = broadcast_ogg_flac_stream(flac_stream, ogg_broadcast, ogg_header).await { + if let Err(e) = broadcast_ogg_flac_stream(flac_stream, ogg_broadcast, ogg_header, current_timestamp).await { error!("OGG broadcaster task error: {}", e); } }); @@ -374,8 +390,12 @@ impl NodeLogic for StreamingOggFlacSinkLogic { seg.timestamp_sec ); - // Send to FLAC encoder - if let Err(e) = self.pcm_tx.send(pcm_bytes).await { + // Send to FLAC encoder with timestamp + let pcm_chunk = PcmChunk { + bytes: pcm_bytes, + timestamp_sec: seg.timestamp_sec, + }; + if let Err(e) = self.pcm_tx.send(pcm_chunk).await { warn!("Failed to send PCM data to encoder: {}", e); break; } @@ -450,7 +470,7 @@ impl StreamingOggFlacSink { } // Create PCM channel (bounded for backpressure) - let (pcm_tx, pcm_rx) = mpsc::channel::>(16); + let (pcm_tx, pcm_rx) = mpsc::channel::(16); // Shared metadata let metadata = Arc::new(RwLock::new(MetadataSnapshot::default())); @@ -522,19 +542,23 @@ impl TypedAudioNode for StreamingOggFlacSink { } } -/// AsyncRead adapter for mpsc::Receiver>. +/// AsyncRead adapter for mpsc::Receiver. +/// Extracts bytes from PcmChunk and provides them to the FLAC encoder. struct ByteStreamReader { - rx: mpsc::Receiver>, + rx: mpsc::Receiver, buffer: VecDeque, finished: bool, + /// Shared timestamp for broadcaster pacing + current_timestamp: Arc>, } impl ByteStreamReader { - fn new(rx: mpsc::Receiver>) -> Self { + fn new(rx: mpsc::Receiver, current_timestamp: Arc>) -> Self { Self { rx, buffer: VecDeque::new(), finished: false, + current_timestamp, } } } @@ -563,11 +587,15 @@ impl AsyncRead for ByteStreamReader { } match Pin::new(&mut self.rx).poll_recv(cx) { - Poll::Ready(Some(bytes)) => { - if bytes.is_empty() { + Poll::Ready(Some(chunk)) => { + if chunk.bytes.is_empty() { continue; } - self.buffer.extend(bytes); + // Update shared timestamp for broadcaster pacing + if let Ok(mut ts) = self.current_timestamp.try_write() { + *ts = chunk.timestamp_sec; + } + self.buffer.extend(chunk.bytes); } Poll::Ready(None) => { self.finished = true; @@ -673,18 +701,21 @@ fn chunk_to_pcm_bytes(chunk: &AudioChunk, bits_per_sample: u8) -> Result } /// OGG wrapper + broadcaster task: reads FLAC bytes from encoder, wraps in OGG pages, and broadcasts. +/// Implements precise real-time pacing based on audio timestamps. async fn broadcast_ogg_flac_stream( mut flac_stream: FlacEncodedStream, broadcast_tx: broadcast::Sender, header_cache: Arc>>, + current_timestamp: Arc>, ) -> Result<(), AudioError> { - info!("OGG-FLAC broadcaster task started"); + info!("OGG-FLAC broadcaster task started with precise timestamp-based pacing"); let stream_serial = rand::random::(); let mut ogg_writer = OggPageWriter::new(stream_serial); let mut total_ogg_bytes = 0u64; let mut header_captured = false; + let start_time = std::time::Instant::now(); // Step 1: Read FLAC header (fLaC + metadata blocks) let flac_header = read_flac_header(&mut flac_stream).await?; @@ -746,6 +777,20 @@ async fn broadcast_ogg_flac_stream( break; } Ok(n) => { + // Precise pacing based on audio timestamp + let audio_timestamp = *current_timestamp.read().await; + let elapsed = start_time.elapsed().as_secs_f64(); + let lead_time = audio_timestamp - elapsed; + + if lead_time > BROADCAST_MAX_LEAD_TIME { + let sleep_duration = lead_time - BROADCAST_MAX_LEAD_TIME; + debug!( + "OGG broadcaster pacing: sleeping {:.3}s (audio_ts={:.3}s, elapsed={:.3}s, lead={:.3}s)", + sleep_duration, audio_timestamp, elapsed, lead_time + ); + tokio::time::sleep(tokio::time::Duration::from_secs_f64(sleep_duration)).await; + } + // Accumulate FLAC data flac_data.extend_from_slice(&read_buffer[..n]);