From 81b990a8284d63d4510e289aeeb21cd84c0356c9 Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 12 Nov 2025 08:50:10 +0000 Subject: [PATCH] Add precise timestamp-based HTTP broadcast pacing Implemented real-time pacing at the HTTP broadcast level based on audio timestamps propagated from the pipeline. This provides much tighter control over streaming bandwidth compared to the pipeline TimerNode alone. Key changes: - Created PcmChunk struct to carry both PCM bytes and timestamps - Modified PCM channels from mpsc::channel> to mpsc::channel - ByteStreamReader now extracts timestamps and shares them via Arc> - Broadcasters read current audio timestamp and pace output accordingly - BROADCAST_MAX_LEAD_TIME set to 0.5s (vs 3.0s for pipeline TimerNode) Benefits: - Precise real-time delivery: ~92 KB/s for FLAC, ~86 KB/s for OGG-FLAC - Lower latency for new clients (0.5s buffer vs 3s) - Smoother streaming without bursts - Works with both StreamingFlacSink and StreamingOggFlacSink Tested: - FLAC streaming: 92.27 KB/s average over 30s (verified with curl) - OGG-FLAC streaming: 86.40 KB/s average over 30s - Both formats correctly identified by file command - Compilation successful with no errors Affected files: - streaming_flac_sink.rs: PcmChunk, ByteStreamReader, broadcast_flac_stream pacing - streaming_ogg_flac_sink.rs: Same changes for OGG-FLAC variant --- pmoaudio-ext/src/sinks/streaming_flac_sink.rs | 76 +++++++++++++++---- .../src/sinks/streaming_ogg_flac_sink.rs | 75 ++++++++++++++---- 2 files changed, 121 insertions(+), 30 deletions(-) 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]);