diff --git a/Cargo.lock b/Cargo.lock index 384098fe4f..54658621ee 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -9991,7 +9991,7 @@ dependencies = [ [[package]] name = "rockbox-playback" -version = "0.4.2" +version = "0.5.0" dependencies = [ "cpal", "reqwest", diff --git a/crates/rockbox-ffi/Cargo.toml b/crates/rockbox-ffi/Cargo.toml index 3d7cd65785..762c655a3d 100644 --- a/crates/rockbox-ffi/Cargo.toml +++ b/crates/rockbox-ffi/Cargo.toml @@ -28,6 +28,6 @@ player = ["dep:rockbox-playback"] rockbox-codecs = { version = "0.2.1", path = "../rockbox-codecs" } rockbox-dsp = { version = "0.3.0", path = "../rockbox-dsp" } rockbox-metadata = { version = "0.1.0", path = "../rockbox-metadata" } -rockbox-playback = { version = "0.4.2", path = "../rockbox-playback", optional = true } +rockbox-playback = { version = "0.5.0", path = "../rockbox-playback", optional = true } serde = { workspace = true } serde_json = { workspace = true } diff --git a/crates/rockbox-ffi/src/player.rs b/crates/rockbox-ffi/src/player.rs index 29980271cd..8b6b53c31a 100644 --- a/crates/rockbox-ffi/src/player.rs +++ b/crates/rockbox-ffi/src/player.rs @@ -9,8 +9,9 @@ use crate::meta::MetadataJson; use crate::util::{cstr, into_cstring}; use rockbox_playback::{ m3u, BassEnhancement, ChannelMode, Compressor, CrossfadeMode, CrossfadeSettings, Crossfeed, - CrossfeedMode, DspSettings, EqBand, EqPreset, InsertPosition, MixMode, PlaybackState, Player, - PlayerConfig, RepeatMode, ReplayGainMode, ResumeState, Status, Surround, ToneControls, + CrossfeedMode, DspSettings, EqBand, EqPreset, InsertPosition, MixMode, OutputConfig, + PlaybackState, Player, PlayerConfig, RepeatMode, ReplayGainMode, ResumeState, Status, Surround, + ToneControls, }; use serde::Serialize; use std::os::raw::c_char; @@ -108,6 +109,7 @@ pub extern "C" fn rb_player_new_with_config( mix_mode_v: i32, ) -> *mut Player { let cfg = PlayerConfig { + output: OutputConfig::default(), sample_rate: (sample_rate != 0).then_some(sample_rate), buffer_seconds, crossfade: crossfade( @@ -166,6 +168,87 @@ pub extern "C" fn rb_player_new_with_config_ex( Duration::from_millis(resume_save_interval_ms as u64) }; let cfg = PlayerConfig { + output: OutputConfig::default(), + sample_rate: (sample_rate != 0).then_some(sample_rate), + buffer_seconds, + crossfade: crossfade( + xfade_mode_v, + fo_delay_ms, + fo_dur_ms, + fi_delay_ms, + fi_dur_ms, + mix_mode_v, + ), + replaygain_mode: rg_mode(rg_mode_v), + replaygain_preamp_db: rg_preamp_db, + replaygain_prevent_clipping: rg_prevent_clipping, + dsp: DspSettings::default(), + shuffle: false, + repeat: RepeatMode::Off, + volume, + resume_file, + resume_save_interval: interval, + }; + match Player::with_config(cfg) { + Ok(p) => Box::into_raw(Box::new(p)), + Err(_) => std::ptr::null_mut(), + } +} + +/// Create a player whose **output backend** is chosen by an `output` spec +/// string, otherwise identical to [`rb_player_new_with_config_ex`]. The spec +/// mirrors [`OutputConfig`]'s syntax: +/// +/// ```text +/// cpal system audio device (same as the other ctors) +/// stdout (or -) raw S16LE stereo on stdout +/// fifo:/tmp/snapfifo raw S16LE into a named FIFO +/// unix:/tmp/rb.sock raw S16LE, Unix socket (listen) +/// unix-connect:/tmp/rb.sock raw S16LE, Unix socket (connect out) +/// tcp:0.0.0.0:9000 raw S16LE, TCP (listen) +/// tcp-connect:host:9000 raw S16LE, TCP (connect out) +/// ``` +/// +/// A null/empty `output` defaults to `cpal`. Non-`cpal` backends emit raw +/// S16LE stereo at `sample_rate` (0 → 44100 Hz), paced to real time. A +/// *listening* socket blocks until a client connects. Returns null on an +/// invalid spec or any open/connect failure. +#[no_mangle] +#[allow(clippy::too_many_arguments)] +pub extern "C" fn rb_player_new_with_output( + output: *const c_char, + sample_rate: u32, + buffer_seconds: f32, + volume: f32, + rg_mode_v: i32, + rg_preamp_db: f32, + rg_prevent_clipping: bool, + xfade_mode_v: i32, + fo_delay_ms: u32, + fo_dur_ms: u32, + fi_delay_ms: u32, + fi_dur_ms: u32, + mix_mode_v: i32, + resume_file: *const c_char, + resume_save_interval_ms: u32, +) -> *mut Player { + let output = match cstr(output).filter(|s| !s.is_empty()) { + Some(spec) => match spec.parse::() { + Ok(cfg) => cfg, + Err(_) => return std::ptr::null_mut(), + }, + None => OutputConfig::default(), + }; + let resume_file = cstr(resume_file) + .filter(|s| !s.is_empty()) + .map(PathBuf::from); + let interval = if resume_save_interval_ms == 0 { + Duration::from_secs(5) + } else { + Duration::from_millis(resume_save_interval_ms as u64) + }; + let cfg = PlayerConfig { + output, sample_rate: (sample_rate != 0).then_some(sample_rate), buffer_seconds, crossfade: crossfade( diff --git a/crates/rockbox-playback/Cargo.toml b/crates/rockbox-playback/Cargo.toml index e87a628695..16666a0577 100644 --- a/crates/rockbox-playback/Cargo.toml +++ b/crates/rockbox-playback/Cargo.toml @@ -1,7 +1,7 @@ [package] name = "rockbox-playback" -version = "0.4.2" -description = "Audio playback engine over rockbox-codecs + rockbox-dsp + cpal — native ReplayGain and Rockbox crossfade" +version = "0.5.0" +description = "Audio playback engine over rockbox-codecs + rockbox-dsp — native ReplayGain and Rockbox crossfade, with configurable output (cpal device, stdout, FIFO, Unix/TCP socket)" authors = { workspace = true } edition = { workspace = true } license = "GPL-2.0-or-later" @@ -11,16 +11,20 @@ keywords = ["audio", "player", "playback", "crossfade", "rockbox"] categories = ["multimedia::audio"] [features] -default = ["http"] +default = ["http", "cpal"] # HTTP(S) remote media support: stream/decode `http(s)://` URLs via ranged # requests, backed by `reqwest`. Disable for a leaner, local-file-only build. http = ["dep:reqwest", "dep:tempfile"] +# System-audio output via `cpal` (the default `OutputConfig::Cpal` backend). +# Disable to drop the `cpal` dependency and use only the byte-stream backends +# (stdout / FIFO / Unix / TCP), e.g. on a headless host with no audio device. +cpal = ["dep:cpal"] [dependencies] rockbox-metadata = { version = "0.1.0", path = "../rockbox-metadata" } rockbox-codecs = { version = "0.2.1", path = "../rockbox-codecs" } rockbox-dsp = { version = "0.3.0", path = "../rockbox-dsp" } -cpal = "0.15" +cpal = { version = "0.15", optional = true } reqwest = { version = "0.12", default-features = false, features = [ "blocking", "rustls-tls-native-roots", diff --git a/crates/rockbox-playback/README.md b/crates/rockbox-playback/README.md index f990afbd94..2804cb56e5 100644 --- a/crates/rockbox-playback/README.md +++ b/crates/rockbox-playback/README.md @@ -224,6 +224,56 @@ Position resolution mirrors `apps/playlist.c:playlist_update_resume_info` (`resume_index` + `resume_elapsed`); the saved elapsed is the true *playback* position, corrected for the decode-ahead buffer. +## Output backends + +By default the player opens the system audio device via +[`cpal`](https://crates.io/crates/cpal). Set `PlayerConfig.output` (an +[`OutputConfig`]) to send audio elsewhere instead — every non-`cpal` +backend emits the **same** raw interleaved **S16LE stereo** byte stream, +paced to real time so a consumer that doesn't clock the stream itself +still plays at the right speed. + +| Backend | Where audio goes | +| ---------------------- | --------------------------------------------------------- | +| `OutputConfig::Cpal` | System audio device (default; needs the `cpal` feature). | +| `Stdout` | Raw S16LE on stdout — pipe to any player. | +| `Fifo(path)` | Raw S16LE into a named FIFO (e.g. a Snapcast pipe). | +| `Unix { path, mode }` | Raw S16LE over a Unix-domain socket (listen or connect). | +| `Tcp { addr, mode }` | Raw S16LE over TCP (listen or connect). | + +`OutputConfig` also parses from a compact string (used by the `play` +example and the FFI layer): `cpal`, `stdout` (or `-`), +`fifo:/tmp/snapfifo`, `unix:/path` / `unix-connect:/path`, +`tcp:0.0.0.0:9000` / `tcp-connect:host:9000`. Bare `tcp:`/`unix:` **listen** +(a player connects in); the `-connect` forms **dial out** to a receiver +that is already up. A listening backend blocks in `with_config` until a +client connects. + +```rust +use rockbox_playback::{OutputConfig, PlayerConfig, SocketMode}; + +// Stream raw PCM over TCP; a player connects to us. +let player = PlayerConfig::builder() + .output(OutputConfig::Tcp { + addr: "0.0.0.0:9000".into(), + mode: SocketMode::Listen, + }) + .open()?; +# Ok::<(), rockbox_playback::Error>(()) +``` + +**stdout mode** turns fd 1 into the PCM stream, so pipe it straight to a +player — but the host program must keep stdout otherwise clean (send all +logs to **stderr**): + +```sh +my-app --output stdout song.flac | ffplay -f s16le -ar 44100 -ac 2 - +``` + +Disable the `cpal` default feature for a leaner headless build with only +the byte-stream backends (no audio-device dependency); the default +`OutputConfig` then becomes `Stdout`. + ## DSP chain Beyond ReplayGain, the player exposes Rockbox's entire DSP pipeline. Every diff --git a/crates/rockbox-playback/examples/play.rs b/crates/rockbox-playback/examples/play.rs index b6c0aa5d79..905d02e23d 100644 --- a/crates/rockbox-playback/examples/play.rs +++ b/crates/rockbox-playback/examples/play.rs @@ -18,18 +18,29 @@ //! # mix local + remote, with 2 s crossfade + track ReplayGain //! cargo run --release --example play -- --crossfade 2 --replaygain track \ //! a.flac https://example.com/b.mp3 +//! +//! # send raw S16LE PCM to stdout and pipe it to ffplay (note: all of this +//! # program's own logging goes to stderr, so stdout stays a clean stream) +//! cargo run --release --example play -- --output stdout a.flac \ +//! | ffplay -f s16le -ar 44100 -ac 2 - +//! +//! # stream over TCP; connect a player to it +//! cargo run --release --example play -- --output tcp:0.0.0.0:9000 a.flac +//! ffplay -f s16le -ar 44100 -ac 2 tcp://127.0.0.1:9000 //! ``` use std::time::Duration; use rockbox_playback::{ - is_url, CrossfadeMode, CrossfadeSettings, EqPreset, PlaybackState, Player, ReplayGainMode, + is_url, CrossfadeMode, CrossfadeSettings, EqPreset, OutputConfig, PlaybackState, PlayerConfig, + ReplayGainMode, }; fn main() -> Result<(), Box> { let mut crossfade_secs = 0u64; let mut replaygain = ReplayGainMode::Off; let mut volume = 1.0f32; + let mut output = OutputConfig::Cpal; // Tracks are kept as strings so `http(s)://` URLs pass through unchanged // alongside local file paths — the engine dispatches on the string. let mut tracks: Vec = Vec::new(); @@ -50,6 +61,13 @@ fn main() -> Result<(), Box> { "--volume" | "-v" => { volume = args.next().and_then(|s| s.parse().ok()).unwrap_or(1.0); } + "--output" | "-o" => { + let spec = args.next().unwrap_or_default(); + output = spec.parse().unwrap_or_else(|e| { + eprintln!("{e}"); + std::process::exit(2); + }); + } _ => tracks.push(arg), } } @@ -57,19 +75,22 @@ fn main() -> Result<(), Box> { if tracks.is_empty() { eprintln!( "usage: play [--volume 0..1] [--crossfade SECS] [--replaygain track|album] \ - " + [--output cpal|stdout|fifo:PATH|unix:PATH|tcp:ADDR] " ); std::process::exit(2); } let has_stream = tracks.iter().any(|t| is_url(t)); - let player = Player::new()?; - println!("output: {} Hz", player.sample_rate()); + // In stdout mode fd 1 carries the raw PCM stream, so ALL human output — + // now-playing, progress, diagnostics — goes to stderr. We do that + // unconditionally so the example is correct for every backend. + let player = PlayerConfig::builder().output(output).open()?; + eprintln!("output: {} Hz", player.sample_rate()); // NOTE: volume 0.0 pauses ring consumption (a click-free mute), so playback // looks frozen. Use a non-zero volume to actually hear/advance. player.set_volume(volume); - println!("volume: {volume}"); + eprintln!("volume: {volume}"); if crossfade_secs > 0 { player.set_crossfade(CrossfadeSettings { @@ -78,24 +99,24 @@ fn main() -> Result<(), Box> { fade_out_duration: Duration::from_secs(crossfade_secs), ..Default::default() }); - println!("crossfade: {crossfade_secs}s"); + eprintln!("crossfade: {crossfade_secs}s"); } if replaygain != ReplayGainMode::Off { player.set_replaygain(replaygain, 0.0, true); - println!("replaygain: {replaygain:?}"); + eprintln!("replaygain: {replaygain:?}"); } // DSP: apply the Bass Boost equalizer preset and a +3 dB bass/treble lift. player.set_eq_preset(EqPreset::BassBoost); player.set_bass(7); player.set_treble(4); - println!("eq: BassBoost preset, bass +7 dB, treble +4 dB"); + eprintln!("eq: BassBoost preset, bass +7 dB, treble +4 dB"); for t in &tracks { - println!("{}: {t}", if is_url(t) { "url" } else { "file" }); + eprintln!("{}: {t}", if is_url(t) { "url" } else { "file" }); } if has_stream { - println!("(live streams play until Ctrl-C)"); + eprintln!("(live streams play until Ctrl-C)"); } player.set_queue(tracks.clone()); @@ -121,9 +142,9 @@ fn main() -> Result<(), Box> { eprintln!("\nfailed to start playback (could not open the source)"); std::process::exit(1); } - print!("\rconnecting… "); + eprint!("\rconnecting… "); use std::io::Write; - std::io::stdout().flush().ok(); + std::io::stderr().flush().ok(); continue; } // Once started, a return to Stopped means the queue finished. @@ -189,12 +210,12 @@ fn main() -> Result<(), Box> { clock, ); if line != last_line { - print!("\r{line}"); + eprint!("\r{line}"); use std::io::Write; - std::io::stdout().flush().ok(); + std::io::stderr().flush().ok(); last_line = line; } } - println!("\ndone"); + eprintln!("\ndone"); Ok(()) } diff --git a/crates/rockbox-playback/src/lib.rs b/crates/rockbox-playback/src/lib.rs index 2f70ad63dc..c94f9f3e09 100644 --- a/crates/rockbox-playback/src/lib.rs +++ b/crates/rockbox-playback/src/lib.rs @@ -35,11 +35,13 @@ mod crossfade; pub mod m3u; +pub mod output; mod resume; pub mod source; pub use crossfade::{CrossfadeMode, CrossfadeSettings, MixMode}; pub use m3u::M3uEntry; +pub use output::{OutputConfig, ParseOutputError, SocketMode}; pub use resume::ResumeState; pub use rockbox_codecs::Decoder; pub use rockbox_metadata::Metadata; @@ -59,10 +61,12 @@ use std::path::PathBuf; use std::sync::atomic::{ AtomicBool, AtomicI32, AtomicU32, AtomicU64, AtomicU8, AtomicUsize, Ordering, }; +use std::io::Write; use std::sync::mpsc::{Receiver, Sender}; use std::sync::{Arc, Mutex}; use std::time::{Duration, Instant}; +#[cfg(feature = "cpal")] use cpal::traits::{DeviceTrait, HostTrait, StreamTrait}; /// How the queue repeats when a track (or the whole queue) finishes. @@ -553,9 +557,14 @@ pub struct Status { /// Configuration for [`Player::with_config`]. pub struct PlayerConfig { + /// Where audio is sent: the system device via `cpal` (default), or a + /// raw **S16LE** stereo stream to stdout / a FIFO / a Unix or TCP + /// socket. See [`OutputConfig`]. + pub output: OutputConfig, /// Output sample rate. `None` uses the output device's default; every /// track is resampled to this rate by the DSP so mixed-rate queues - /// work. + /// work. For the non-`cpal` byte-stream backends there is no device to + /// query, so `None` falls back to 44100 Hz. pub sample_rate: Option, /// Seconds of audio to decode ahead into the ring buffer. pub buffer_seconds: f32, @@ -586,6 +595,7 @@ pub struct PlayerConfig { impl Default for PlayerConfig { fn default() -> Self { PlayerConfig { + output: OutputConfig::default(), sample_rate: None, buffer_seconds: 4.0, crossfade: CrossfadeSettings::default(), @@ -631,6 +641,12 @@ pub struct PlayerConfigBuilder { } impl PlayerConfigBuilder { + /// Where audio is sent (system device, stdout, FIFO, Unix or TCP + /// socket). See [`OutputConfig`]; defaults to [`OutputConfig::Cpal`]. + pub fn output(mut self, output: OutputConfig) -> Self { + self.cfg.output = output; + self + } /// Output sample rate in Hz. Unset (the default) uses the output device's /// native rate; every track is resampled to this rate. pub fn sample_rate(mut self, hz: u32) -> Self { @@ -819,10 +835,15 @@ impl Shared { /// Errors constructing a [`Player`]. #[derive(Debug)] pub enum Error { - /// No output audio device available. + /// No output audio device available (`cpal` backend). NoOutputDevice, - /// cpal could not build or start the output stream. + /// cpal could not build/start the stream, or a byte-stream backend + /// (stdout / FIFO / Unix / TCP) failed to open or connect. Stream(String), + /// The requested [`OutputConfig`] backend was compiled out (e.g. the + /// `cpal` backend without the `cpal` feature, or a FIFO/Unix socket on + /// a non-Unix platform). + UnsupportedBackend(String), } impl std::fmt::Display for Error { @@ -830,19 +851,50 @@ impl std::fmt::Display for Error { match self { Error::NoOutputDevice => write!(f, "no output audio device"), Error::Stream(e) => write!(f, "audio stream error: {e}"), + Error::UnsupportedBackend(e) => write!(f, "unsupported output backend: {e}"), } } } impl std::error::Error for Error {} +/// The live output backend, kept on the [`Player`] handle so output stays +/// alive for the player's lifetime. `cpal` streams are non-`Send`, which is +/// why the handle (not the engine thread) owns them. The variants are held +/// purely for their `Drop` side effects (stop the stream / join the writer). +#[allow(dead_code)] +enum OutputHandle { + /// System audio device: a `cpal` stream whose callback drains the ring. + #[cfg(feature = "cpal")] + Cpal(cpal::Stream), + /// A byte-stream sink (stdout / FIFO / Unix / TCP): a writer thread + /// drains the ring, converts to S16LE and paces to real time. Dropping + /// this signals the thread to stop and joins it. + Stream(StreamSink), +} + +/// Owns the writer thread for a byte-stream backend and stops it on drop. +struct StreamSink { + stop: Arc, + thread: Option>, +} + +impl Drop for StreamSink { + fn drop(&mut self) { + self.stop.store(true, Ordering::Relaxed); + if let Some(t) = self.thread.take() { + let _ = t.join(); + } + } +} + /// The player handle. Cloneable-free but `Send` controls are issued -/// through it; the cpal stream lives here and keeps output alive for the +/// through it; the output backend lives here and keeps output alive for the /// player's lifetime. pub struct Player { tx: Sender, shared: Arc, - _stream: cpal::Stream, + _output: OutputHandle, engine: Option>, /// Resume file path (mirrors `PlayerConfig::resume_file`) so /// [`Player::resume`] can read it on the handle side. @@ -856,70 +908,44 @@ impl Player { } /// Create a player with explicit configuration. + /// + /// The output backend is chosen by [`PlayerConfig::output`]. Note that a + /// *listening* socket backend ([`SocketMode::Listen`]) blocks here until + /// a client connects, and a *connecting* one requires the receiver to be + /// up already. pub fn with_config(config: PlayerConfig) -> Result { - let host = cpal::default_host(); - let device = host.default_output_device().ok_or(Error::NoOutputDevice)?; - let default_cfg = device - .default_output_config() + #[cfg(feature = "cpal")] + if config.output == OutputConfig::Cpal { + let host = cpal::default_host(); + let device = host.default_output_device().ok_or(Error::NoOutputDevice)?; + let default_cfg = device + .default_output_config() + .map_err(|e| Error::Stream(e.to_string()))?; + let rate = config + .sample_rate + .unwrap_or_else(|| default_cfg.sample_rate().0); + let shared = make_shared(&config, rate); + let stream = build_stream(&device, rate, Arc::clone(&shared))?; + stream.play().map_err(|e| Error::Stream(e.to_string()))?; + return Ok(assemble(config, rate, shared, OutputHandle::Cpal(stream))); + } + #[cfg(not(feature = "cpal"))] + if config.output == OutputConfig::Cpal { + return Err(Error::UnsupportedBackend( + "cpal backend requires the `cpal` feature".into(), + )); + } + + // Byte-stream backends (stdout / FIFO / Unix / TCP): no device to + // query, so an unset sample rate falls back to CD quality. + let rate = config.sample_rate.unwrap_or(44100); + let writer = config + .output + .open_writer() .map_err(|e| Error::Stream(e.to_string()))?; - - let rate = config - .sample_rate - .unwrap_or_else(|| default_cfg.sample_rate().0); - - let shared = Arc::new(Shared { - state: AtomicU8::new(ST_STOPPED), - decode_pos_ms: AtomicU64::new(0), - duration_ms: AtomicU64::new(0), - index: AtomicUsize::new(usize::MAX), - queue_len: AtomicUsize::new(0), - target_amp: AtomicU32::new(0f32.to_bits()), - volume: AtomicU32::new(config.volume.clamp(0.0, 1.0).to_bits()), - balance: AtomicI32::new(0), - balance_gain_l: AtomicU32::new(1f32.to_bits()), - balance_gain_r: AtomicU32::new(1f32.to_bits()), - output_rate: AtomicU32::new(rate), - shuffle: AtomicBool::new(config.shuffle), - repeat: AtomicU8::new(config.repeat.to_u8()), - ring: Mutex::new(VecDeque::new()), - meta: Mutex::new(None), - dsp: Mutex::new(config.dsp.clone()), - queue: Mutex::new(Vec::new()), - }); - - let stream = build_stream(&device, rate, Arc::clone(&shared))?; - stream.play().map_err(|e| Error::Stream(e.to_string()))?; - - let (tx, rx) = std::sync::mpsc::channel(); - let engine_shared = Arc::clone(&shared); - let resume_file = config.resume_file.clone(); - let engine_cfg = EngineConfig { - output_rate: rate, - buffer_frames: (config.buffer_seconds.max(0.5) * rate as f32) as usize, - crossfade: config.crossfade, - replaygain: ReplayGainConfig { - mode: config.replaygain_mode, - preamp_db: config.replaygain_preamp_db, - prevent_clipping: config.replaygain_prevent_clipping, - }, - dsp: config.dsp, - shuffle: config.shuffle, - repeat: config.repeat, - resume_file: resume_file.clone(), - resume_save_interval: config.resume_save_interval, - }; - let engine = std::thread::Builder::new() - .name("rbplayback".into()) - .spawn(move || Engine::new(engine_shared, rx, engine_cfg).run()) - .expect("spawn engine thread"); - - Ok(Player { - tx, - shared, - _stream: stream, - engine: Some(engine), - resume_file, - }) + let shared = make_shared(&config, rate); + let sink = spawn_stream_writer(writer, rate, Arc::clone(&shared)); + Ok(assemble(config, rate, shared, OutputHandle::Stream(sink))) } /// The output sample rate everything is resampled to. @@ -1333,10 +1359,159 @@ impl Drop for Player { } } +/// Allocate the [`Shared`] state block from a config and the resolved output +/// rate. Backend-agnostic — used by every [`OutputConfig`] path. +fn make_shared(config: &PlayerConfig, rate: u32) -> Arc { + Arc::new(Shared { + state: AtomicU8::new(ST_STOPPED), + decode_pos_ms: AtomicU64::new(0), + duration_ms: AtomicU64::new(0), + index: AtomicUsize::new(usize::MAX), + queue_len: AtomicUsize::new(0), + target_amp: AtomicU32::new(0f32.to_bits()), + volume: AtomicU32::new(config.volume.clamp(0.0, 1.0).to_bits()), + balance: AtomicI32::new(0), + balance_gain_l: AtomicU32::new(1f32.to_bits()), + balance_gain_r: AtomicU32::new(1f32.to_bits()), + output_rate: AtomicU32::new(rate), + shuffle: AtomicBool::new(config.shuffle), + repeat: AtomicU8::new(config.repeat.to_u8()), + ring: Mutex::new(VecDeque::new()), + meta: Mutex::new(None), + dsp: Mutex::new(config.dsp.clone()), + queue: Mutex::new(Vec::new()), + }) +} + +/// Spawn the engine thread and wrap everything up into a [`Player`]. Shared +/// by every backend; `output` is the already-started output handle. +fn assemble(config: PlayerConfig, rate: u32, shared: Arc, output: OutputHandle) -> Player { + let (tx, rx) = std::sync::mpsc::channel(); + let engine_shared = Arc::clone(&shared); + let resume_file = config.resume_file.clone(); + let engine_cfg = EngineConfig { + output_rate: rate, + buffer_frames: (config.buffer_seconds.max(0.5) * rate as f32) as usize, + crossfade: config.crossfade, + replaygain: ReplayGainConfig { + mode: config.replaygain_mode, + preamp_db: config.replaygain_preamp_db, + prevent_clipping: config.replaygain_prevent_clipping, + }, + dsp: config.dsp, + shuffle: config.shuffle, + repeat: config.repeat, + resume_file: resume_file.clone(), + resume_save_interval: config.resume_save_interval, + }; + let engine = std::thread::Builder::new() + .name("rbplayback".into()) + .spawn(move || Engine::new(engine_shared, rx, engine_cfg).run()) + .expect("spawn engine thread"); + + Player { + tx, + shared, + _output: output, + engine: Some(engine), + resume_file, + } +} + +/// Drain the ring into a raw **S16LE** stereo byte stream (stdout / FIFO / +/// Unix / TCP), paced to real time with a monotonic clock so a consumer +/// that does *not* clock the stream itself (a FIFO, a socket, `ffplay -`) +/// still plays at the correct speed. +/// +/// Mirrors the `cpal` callback's fade/balance semantics: while paused +/// (`target_amp == 0`) it emits **silence without draining the ring**, so +/// the buffered audio is frozen and resume is click-free — and the byte +/// stream keeps flowing so a permanent reader (Snapcast, `ffplay`) never +/// sees a gap or EOF. +fn spawn_stream_writer( + mut writer: Box, + rate: u32, + shared: Arc, +) -> StreamSink { + let stop = Arc::new(AtomicBool::new(false)); + let stop_thread = Arc::clone(&stop); + // ~20 ms chunks: small enough for responsive transport, large enough to + // keep syscall overhead negligible. + let chunk_frames = (rate / 50).max(1) as usize; + // ~1/3 second to fade the full 0..1 range, matching pcmbuf_fade_tick. + let step = 3.0 / rate as f32; + let frame_dur = Duration::from_secs_f64(chunk_frames as f64 / rate as f64); + + let thread = std::thread::Builder::new() + .name("rbplayback-out".into()) + .spawn(move || { + // 2 channels * 2 bytes/sample. + let mut buf = vec![0u8; chunk_frames * 4]; + let mut cur_amp = 0.0f32; + let mut next = Instant::now() + frame_dur; + + while !stop_thread.load(Ordering::Relaxed) { + let target = f32::from_bits(shared.target_amp.load(Ordering::Relaxed)); + let gain_l = f32::from_bits(shared.balance_gain_l.load(Ordering::Relaxed)); + let gain_r = f32::from_bits(shared.balance_gain_r.load(Ordering::Relaxed)); + + if target == 0.0 && cur_amp == 0.0 { + // Paused/stopped: emit silence, freeze the ring. + buf.iter_mut().for_each(|b| *b = 0); + } else { + let mut ring = shared.ring.lock().unwrap(); + for frame in buf.chunks_mut(4) { + if cur_amp < target { + cur_amp = (cur_amp + step).min(target); + } else if cur_amp > target { + cur_amp = (cur_amp - step).max(target); + } + // Fade-out just completed: stop draining, freeze. + if cur_amp == 0.0 && target == 0.0 { + frame.fill(0); + continue; + } + let l = ring.pop_front().unwrap_or(0); + let r = ring.pop_front().unwrap_or(0); + let lv = ((l as f32) * cur_amp * gain_l).clamp(-32768.0, 32767.0) as i16; + let rv = ((r as f32) * cur_amp * gain_r).clamp(-32768.0, 32767.0) as i16; + frame[0..2].copy_from_slice(&lv.to_le_bytes()); + frame[2..4].copy_from_slice(&rv.to_le_bytes()); + } + } + + // A write/flush error means the consumer went away (pipe + // closed, socket reset): stop cleanly rather than spin. + if writer.write_all(&buf).is_err() || writer.flush().is_err() { + break; + } + + // Pace to real time. If we fell behind (scheduling hiccup), + // resync the deadline instead of trying to "catch up" in a + // burst. + let now = Instant::now(); + if next > now { + std::thread::sleep(next - now); + } + next += frame_dur; + if next < now { + next = now + frame_dur; + } + } + }) + .expect("spawn output writer thread"); + + StreamSink { + stop, + thread: Some(thread), + } +} + /// Build the cpal output stream: drains the ring buffer, converts i16 → /// f32, and applies a per-sample amplitude ramp toward `target_amp` /// (~⅓ s full-range, matching Rockbox's pause/stop fade) for click-free /// transitions. +#[cfg(feature = "cpal")] fn build_stream( device: &cpal::Device, rate: u32, diff --git a/crates/rockbox-playback/src/output.rs b/crates/rockbox-playback/src/output.rs new file mode 100644 index 0000000000..bdeb0d67ee --- /dev/null +++ b/crates/rockbox-playback/src/output.rs @@ -0,0 +1,288 @@ +//! Output-backend selection for [`Player`](crate::Player). +//! +//! The engine thread always produces the same thing — interleaved **stereo +//! `i16` at the output rate** — into a ring buffer. [`OutputConfig`] chooses +//! who drains that ring and where the samples go: +//! +//! | Backend | What it does | +//! | ---------------------- | -------------------------------------------------------- | +//! | [`OutputConfig::Cpal`] | System audio device via `cpal` (default; needs `cpal`). | +//! | `Stdout` | Raw **S16LE** stereo on stdout — pipe to a player. | +//! | `Fifo` | Raw S16LE into a named FIFO (e.g. a Snapcast pipe). | +//! | `Unix` | Raw S16LE over a Unix-domain socket (listen or connect). | +//! | `Tcp` | Raw S16LE over TCP (listen or connect). | +//! +//! Every non-`cpal` backend emits the exact same byte stream, real-time +//! paced, so it drops straight into a terminal player: +//! +//! ```sh +//! # stdout → ffplay +//! my-app --output stdout song.flac | ffplay -f s16le -ar 44100 -ac 2 - +//! ``` +//! +//! **Important:** in `Stdout` mode fd 1 carries the PCM stream, so the host +//! program must keep stdout clean — send *all* logs/diagnostics to stderr. +//! +//! Backends parse from a compact string via [`str::parse`], which the FFI +//! layer and the `play` example use: +//! +//! ```text +//! cpal +//! stdout (or "-") +//! fifo:/tmp/snapfifo +//! unix:/tmp/rb.sock (listen) unix-connect:/tmp/rb.sock +//! tcp:127.0.0.1:9000 (listen) tcp-connect:192.168.1.9:9000 +//! ``` + +use std::io::{self, Write}; +use std::path::PathBuf; + +/// Whether a socket backend binds and accepts clients, or dials out to a +/// receiver that is already listening. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum SocketMode { + /// Bind the address/path and stream to whoever connects. This is the + /// default for `tcp:` / `unix:` — a player connects *in* (e.g. + /// `ffplay tcp://127.0.0.1:9000`). + Listen, + /// Dial out to an address/path that is already listening (e.g. a + /// Snapcast TCP source). The receiver must be up at construction time. + Connect, +} + +/// Where [`Player`](crate::Player) sends its audio. See the [module +/// docs](self) for the string syntax and the byte format. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum OutputConfig { + /// System audio device via `cpal`. Requires the `cpal` feature (on by + /// default). This is the [`Default`]. + Cpal, + /// Raw interleaved **S16LE** stereo on stdout (fd 1). Pipe to any + /// player that reads raw PCM. The host program must not write anything + /// else to stdout. + Stdout, + /// Raw S16LE stereo into a named FIFO at this path. The FIFO must + /// already exist (`mkfifo`); it is opened read+write so a permanent + /// writer reference is held and readers never see EOF between tracks. + Fifo(PathBuf), + /// Raw S16LE stereo over a Unix-domain socket. + Unix { + /// Filesystem path of the socket. + path: PathBuf, + /// Listen (bind + accept) or connect (dial out). + mode: SocketMode, + }, + /// Raw S16LE stereo over TCP. + Tcp { + /// `host:port` (or `:port` to bind all interfaces when listening). + addr: String, + /// Listen (bind + accept) or connect (dial out). + mode: SocketMode, + }, +} + +impl OutputConfig { + /// Open the byte-stream writer for a non-`cpal` backend. The `Cpal` + /// variant has no writer (it is device-driven) and returns an error. + /// + /// For a *listening* socket ([`SocketMode::Listen`]) this **blocks** + /// until a client connects; for a *connecting* one the receiver must + /// already be up. + pub(crate) fn open_writer(&self) -> io::Result> { + match self { + OutputConfig::Cpal => Err(io::Error::new( + io::ErrorKind::Unsupported, + "cpal backend has no byte-stream writer", + )), + // fd 1. The host program must keep stdout otherwise clean — + // logs go to stderr — so the raw PCM is pipeable to a player. + OutputConfig::Stdout => Ok(Box::new(io::stdout())), + OutputConfig::Fifo(path) => open_fifo(path), + OutputConfig::Unix { path, mode } => open_unix(path, *mode), + OutputConfig::Tcp { addr, mode } => open_tcp(addr, *mode), + } + } +} + +/// Open a named FIFO for writing. Opened **read+write** so a permanent +/// writer reference is held — a reader (Snapcast) never sees EOF between +/// tracks. The FIFO must already exist (`mkfifo`). +#[cfg(unix)] +fn open_fifo(path: &std::path::Path) -> io::Result> { + let f = std::fs::OpenOptions::new().read(true).write(true).open(path)?; + Ok(Box::new(f)) +} + +#[cfg(not(unix))] +fn open_fifo(_path: &std::path::Path) -> io::Result> { + Err(io::Error::new( + io::ErrorKind::Unsupported, + "FIFO output is only supported on Unix", + )) +} + +#[cfg(unix)] +fn open_unix(path: &std::path::Path, mode: SocketMode) -> io::Result> { + use std::os::unix::net::{UnixListener, UnixStream}; + match mode { + SocketMode::Listen => { + // Remove a stale socket file so bind() doesn't fail with EADDRINUSE. + let _ = std::fs::remove_file(path); + let listener = UnixListener::bind(path)?; + let (stream, _) = listener.accept()?; + Ok(Box::new(stream)) + } + SocketMode::Connect => Ok(Box::new(UnixStream::connect(path)?)), + } +} + +#[cfg(not(unix))] +fn open_unix(_path: &std::path::Path, _mode: SocketMode) -> io::Result> { + Err(io::Error::new( + io::ErrorKind::Unsupported, + "Unix-socket output is only supported on Unix", + )) +} + +fn open_tcp(addr: &str, mode: SocketMode) -> io::Result> { + use std::net::{TcpListener, TcpStream}; + let stream = match mode { + SocketMode::Listen => { + let listener = TcpListener::bind(addr)?; + listener.accept()?.0 + } + SocketMode::Connect => TcpStream::connect(addr)?, + }; + // Latency over throughput: PCM chunks are small and time-sensitive. + stream.set_nodelay(true).ok(); + Ok(Box::new(stream)) +} + +impl Default for OutputConfig { + // Not derivable: the default is feature-dependent (with only one `cfg` + // arm compiled in, clippy sees a constant and thinks it could be + // `#[derive(Default)]`, but the chosen variant differs per feature). + #[allow(clippy::derivable_impls)] + fn default() -> Self { + // When `cpal` is compiled out there is no device to open, so the + // sensible zero-config default becomes the raw stdout stream. + #[cfg(feature = "cpal")] + { + OutputConfig::Cpal + } + #[cfg(not(feature = "cpal"))] + { + OutputConfig::Stdout + } + } +} + +/// Error returned when an [`OutputConfig`] string cannot be parsed. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ParseOutputError(pub String); + +impl std::fmt::Display for ParseOutputError { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!(f, "invalid audio output spec: {}", self.0) + } +} + +impl std::error::Error for ParseOutputError {} + +impl std::str::FromStr for OutputConfig { + type Err = ParseOutputError; + + fn from_str(s: &str) -> Result { + let s = s.trim(); + // `scheme:rest` where the scheme is everything up to the first ':'. + // stdout/cpal have no argument. + let (scheme, rest) = match s.split_once(':') { + Some((a, b)) => (a, Some(b)), + None => (s, None), + }; + let arg = |rest: Option<&str>| -> Result { + match rest { + Some(r) if !r.is_empty() => Ok(r.to_string()), + _ => Err(ParseOutputError(format!("`{scheme}` needs an argument"))), + } + }; + match scheme { + "cpal" => Ok(OutputConfig::Cpal), + "stdout" | "-" => Ok(OutputConfig::Stdout), + "fifo" => Ok(OutputConfig::Fifo(PathBuf::from(arg(rest)?))), + "unix" | "unix-listen" => Ok(OutputConfig::Unix { + path: PathBuf::from(arg(rest)?), + mode: SocketMode::Listen, + }), + "unix-connect" => Ok(OutputConfig::Unix { + path: PathBuf::from(arg(rest)?), + mode: SocketMode::Connect, + }), + "tcp" | "tcp-listen" => Ok(OutputConfig::Tcp { + addr: arg(rest)?, + mode: SocketMode::Listen, + }), + "tcp-connect" => Ok(OutputConfig::Tcp { + addr: arg(rest)?, + mode: SocketMode::Connect, + }), + other => Err(ParseOutputError(format!("unknown backend `{other}`"))), + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn parses_every_backend() { + assert_eq!("cpal".parse::().unwrap(), OutputConfig::Cpal); + assert_eq!("-".parse::().unwrap(), OutputConfig::Stdout); + assert_eq!( + "stdout".parse::().unwrap(), + OutputConfig::Stdout + ); + assert_eq!( + "fifo:/tmp/snapfifo".parse::().unwrap(), + OutputConfig::Fifo(PathBuf::from("/tmp/snapfifo")) + ); + assert_eq!( + "unix:/tmp/rb.sock".parse::().unwrap(), + OutputConfig::Unix { + path: PathBuf::from("/tmp/rb.sock"), + mode: SocketMode::Listen, + } + ); + assert_eq!( + "unix-connect:/tmp/rb.sock".parse::().unwrap(), + OutputConfig::Unix { + path: PathBuf::from("/tmp/rb.sock"), + mode: SocketMode::Connect, + } + ); + assert_eq!( + "tcp:127.0.0.1:9000".parse::().unwrap(), + OutputConfig::Tcp { + addr: "127.0.0.1:9000".to_string(), + mode: SocketMode::Listen, + } + ); + assert_eq!( + "tcp-connect:192.168.1.9:9000" + .parse::() + .unwrap(), + OutputConfig::Tcp { + addr: "192.168.1.9:9000".to_string(), + mode: SocketMode::Connect, + } + ); + } + + #[test] + fn rejects_garbage_and_missing_args() { + assert!("bogus".parse::().is_err()); + assert!("fifo".parse::().is_err()); + assert!("tcp:".parse::().is_err()); + } +}