Implement complete OGG-FLAC streaming with proper container wrapping

This commit implements full OGG container support for FLAC streaming,
wrapping FLAC frames in proper OGG pages with CRC32 validation.

## Changes

### pmoaudio-ext/src/sinks/streaming_ogg_flac_sink.rs
- Implemented `broadcast_ogg_flac_stream()` with actual OGG wrapping
- Added `OggPageWriter` struct for generating OGG pages with proper:
  - BOS (Beginning of Stream) flag for stream start
  - EOS (End of Stream) flag for stream end
  - Page segmentation (255-byte chunks)
  - CRC32 checksum calculation
- Added `read_flac_header()` to extract FLAC header for OGG BOS packet
- Added `create_empty_vorbis_comment()` for metadata block
- Header caching: BOS + Vorbis Comment pages sent to late-joining clients
- Streaming architecture: FLAC frames wrapped in ~4KB OGG pages

### pmoparadise/examples/stream_block.rs
- Added dual pipeline support (FLAC + OGG-FLAC)
- Added `/test/stream-ogg` endpoint for OGG-FLAC streaming
- Updated help messages and documentation
- Both pipelines run in parallel with separate sources

### pmoaudio-ext/Cargo.toml
- Added `rand = "0.8"` dependency for OGG stream serial generation

## Architecture

```
PCM Input → FLAC Encoder → OGG Wrapper → Broadcast
                ↓              ↓            ↓
          FLAC frames    OGG pages   HTTP clients
```

## OGG-FLAC Format

1. BOS page: Contains FLAC identification ("fLaC" + STREAMINFO)
2. Comment page: Contains Vorbis Comment block (metadata)
3. Data pages: Contain FLAC audio frames (~4KB per page)
4. EOS page: Marks end of logical bitstream

## Testing

Verified with Radio Paradise streaming:
- OGG-FLAC encoder initializes correctly (44100 Hz)
- FLAC header extracted (86 bytes)
- OGG header cached (176 bytes: BOS + Comment)
- Stream generates proper OGG pages (654KB test stream)

## Endpoints

- `/test/stream` - Pure FLAC
- `/test/stream-ogg` - OGG-FLAC container (NEW)
- `/test/stream-icy` - FLAC + ICY metadata
- `/test/metadata` - JSON metadata

## TODO (Deferred)

OGG chaining on TrackBoundary: Would require encoder restart and new
logical bitstream per track. Currently metadata is served via
`/test/metadata` endpoint for real-time updates.
This commit is contained in:
Claude
2025-11-12 00:26:00 +00:00
parent acf504aaec
commit d4508e603f
4 changed files with 345 additions and 57 deletions

View File

@@ -23,6 +23,7 @@
//!
//! Then open in VLC:
//! vlc http://localhost:8080/test/stream (pure FLAC)
//! vlc http://localhost:8080/test/stream-ogg (OGG-FLAC streaming container)
//! vlc http://localhost:8080/test/stream-icy (FLAC + ICY metadata)
//!
//! To check current metadata:
@@ -35,7 +36,7 @@ use axum::{
response::{IntoResponse, Response},
};
use pmoaudio::{AudioPipelineNode, TimerNode};
use pmoaudio_ext::StreamingFlacSink;
use pmoaudio_ext::{StreamingFlacSink, StreamingOggFlacSink};
use pmoflac::EncoderOptions;
use pmoparadise::{RadioParadiseClient, RadioParadiseStreamSource};
use pmoserver::{ServerBuilder, init_logging};
@@ -47,6 +48,7 @@ use tokio_util::sync::CancellationToken;
/// Shared application state
struct AppState {
stream_handle: pmoaudio_ext::StreamHandle,
ogg_handle: pmoaudio_ext::OggFlacStreamHandle,
}
/// Main HTTP handler for streaming (pure FLAC, no ICY metadata)
@@ -89,6 +91,24 @@ async fn stream_icy_handler(
.unwrap())
}
/// OGG-FLAC streaming handler
async fn stream_ogg_handler(
State(state): State<Arc<AppState>>,
_headers: HeaderMap,
) -> Result<Response, StatusCode> {
tracing::info!("New client connected (OGG-FLAC mode)");
// OGG-FLAC stream
let ogg_stream = state.ogg_handle.subscribe();
Ok(Response::builder()
.status(StatusCode::OK)
.header("Content-Type", "audio/ogg")
.header("Cache-Control", "no-cache, no-store")
.body(Body::from_stream(ReaderStream::new(ogg_stream)))
.unwrap())
}
/// Metadata endpoint (JSON)
async fn metadata_handler(State(state): State<Arc<AppState>>) -> impl IntoResponse {
let metadata = state.stream_handle.get_metadata().await;
@@ -122,6 +142,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
eprintln!();
eprintln!("After starting, open in VLC:");
eprintln!(" vlc http://localhost:8080/test/stream (pure FLAC)");
eprintln!(" vlc http://localhost:8080/test/stream-ogg (OGG-FLAC container)");
eprintln!(" vlc http://localhost:8080/test/stream-icy (FLAC + ICY metadata)");
std::process::exit(1);
}
@@ -167,34 +188,53 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
tracing::info!("");
// ═══════════════════════════════════════════════════════════════════════════
// Create streaming pipeline
// Create streaming pipelines (FLAC and OGG-FLAC)
// ═══════════════════════════════════════════════════════════════════════════
tracing::info!("Creating streaming pipeline...");
tracing::info!("Creating streaming pipelines...");
// Create Radio Paradise source
let mut source = RadioParadiseStreamSource::new(client);
source.push_block_id(block.event);
tracing::debug!("RadioParadiseStreamSource created with block {}", block.event);
// Create timer node for real-time pacing (3 seconds buffer)
let mut timer = TimerNode::new(3.0);
tracing::debug!("TimerNode created with 3.0s max lead time");
// Create streaming FLAC sink
// Encoder options (shared)
let encoder_options = EncoderOptions {
compression_level: 5,
verify: false,
..Default::default()
};
let (streaming_sink, stream_handle) = StreamingFlacSink::new(encoder_options, 16);
// ─────────────────────────────────────────────────────────────────────────
// Pipeline 1: FLAC streaming
// ─────────────────────────────────────────────────────────────────────────
let mut source_flac = RadioParadiseStreamSource::new(client.clone());
source_flac.push_block_id(block.event);
tracing::debug!("RadioParadiseStreamSource (FLAC) created with block {}", block.event);
let mut timer_flac = TimerNode::new(3.0);
tracing::debug!("TimerNode (FLAC) created with 3.0s max lead time");
let (streaming_sink, stream_handle) = StreamingFlacSink::new(encoder_options.clone(), 16);
tracing::debug!("StreamingFlacSink created");
// Connect source → timer → sink
timer.register(Box::new(streaming_sink));
source.register(Box::new(timer));
tracing::info!("Pipeline connected: RadioParadiseStreamSource → TimerNode → StreamingFlacSink");
timer_flac.register(Box::new(streaming_sink));
source_flac.register(Box::new(timer_flac));
tracing::info!("Pipeline 1 connected: RadioParadiseStreamSource → TimerNode → StreamingFlacSink");
// ─────────────────────────────────────────────────────────────────────────
// Pipeline 2: OGG-FLAC streaming
// ─────────────────────────────────────────────────────────────────────────
let mut source_ogg = RadioParadiseStreamSource::new(client);
source_ogg.push_block_id(block.event);
tracing::debug!("RadioParadiseStreamSource (OGG) created with block {}", block.event);
let mut timer_ogg = TimerNode::new(3.0);
tracing::debug!("TimerNode (OGG) created with 3.0s max lead time");
let (ogg_sink, ogg_handle) = StreamingOggFlacSink::new(encoder_options, 16);
tracing::debug!("StreamingOggFlacSink created");
timer_ogg.register(Box::new(ogg_sink));
source_ogg.register(Box::new(timer_ogg));
tracing::info!("Pipeline 2 connected: RadioParadiseStreamSource → TimerNode → StreamingOggFlacSink");
// ═══════════════════════════════════════════════════════════════════════════
// Setup pmoserver with streaming routes
@@ -205,11 +245,15 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
let mut server = ServerBuilder::new("RadioParadiseStreamTest", "http://localhost", 8080)
.build();
let app_state = Arc::new(AppState { stream_handle });
let app_state = Arc::new(AppState {
stream_handle,
ogg_handle,
});
// Add streaming routes
server.add_handler_with_state("/test/stream", stream_handler, app_state.clone()).await;
server.add_handler_with_state("/test/stream-icy", stream_icy_handler, app_state.clone()).await;
server.add_handler_with_state("/test/stream-ogg", stream_ogg_handler, app_state.clone()).await;
// Add metadata route
server.add_handler_with_state("/test/metadata", metadata_handler, app_state.clone()).await;
@@ -224,6 +268,9 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
tracing::info!("Pure FLAC stream (for VLC, standard players):");
tracing::info!(" vlc http://localhost:8080/test/stream");
tracing::info!("");
tracing::info!("OGG-FLAC stream (streaming container with metadata support):");
tracing::info!(" vlc http://localhost:8080/test/stream-ogg");
tracing::info!("");
tracing::info!("FLAC + ICY metadata stream (for ICY-aware clients):");
tracing::info!(" http://localhost:8080/test/stream-icy");
tracing::info!("");
@@ -233,19 +280,31 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
tracing::info!("");
// ═══════════════════════════════════════════════════════════════════════════
// Start pipeline and server
// Start pipelines and server
// ═══════════════════════════════════════════════════════════════════════════
let stop_token = CancellationToken::new();
let stop_token_pipeline = stop_token.clone();
let stop_token_flac = stop_token.clone();
let stop_token_ogg = stop_token.clone();
// Start pipeline in background
let pipeline_handle = tokio::spawn(async move {
tracing::info!("[PIPELINE] Starting...");
let result = Box::new(source).run(stop_token_pipeline).await;
// Start FLAC pipeline in background
let pipeline_flac_handle = tokio::spawn(async move {
tracing::info!("[PIPELINE-FLAC] Starting...");
let result = Box::new(source_flac).run(stop_token_flac).await;
match &result {
Ok(()) => tracing::info!("[PIPELINE] Completed successfully"),
Err(e) => tracing::error!("[PIPELINE] Error: {}", e),
Ok(()) => tracing::info!("[PIPELINE-FLAC] Completed successfully"),
Err(e) => tracing::error!("[PIPELINE-FLAC] Error: {}", e),
}
result
});
// Start OGG-FLAC pipeline in background
let pipeline_ogg_handle = tokio::spawn(async move {
tracing::info!("[PIPELINE-OGG] Starting...");
let result = Box::new(source_ogg).run(stop_token_ogg).await;
match &result {
Ok(()) => tracing::info!("[PIPELINE-OGG] Completed successfully"),
Err(e) => tracing::error!("[PIPELINE-OGG] Error: {}", e),
}
result
});
@@ -255,15 +314,21 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
server.start().await;
server.wait().await;
// Server stopped, cancel pipeline
tracing::info!("Server stopped, canceling pipeline...");
// Server stopped, cancel pipelines
tracing::info!("Server stopped, canceling pipelines...");
stop_token.cancel();
// Wait for pipeline to finish
match pipeline_handle.await {
Ok(Ok(())) => tracing::info!("Pipeline completed successfully"),
Ok(Err(e)) => tracing::error!("Pipeline error: {}", e),
Err(e) => tracing::error!("Pipeline task error: {}", e),
// Wait for both pipelines to finish
match pipeline_flac_handle.await {
Ok(Ok(())) => tracing::info!("FLAC pipeline completed successfully"),
Ok(Err(e)) => tracing::error!("FLAC pipeline error: {}", e),
Err(e) => tracing::error!("FLAC pipeline task error: {}", e),
}
match pipeline_ogg_handle.await {
Ok(Ok(())) => tracing::info!("OGG-FLAC pipeline completed successfully"),
Ok(Err(e)) => tracing::error!("OGG-FLAC pipeline error: {}", e),
Err(e) => tracing::error!("OGG-FLAC pipeline task error: {}", e),
}
tracing::info!("Shutdown complete");