rusty-opus 1.0.0

Pure-Rust Opus audio codec (RFC 6716): encoder and decoder, SILK/CELT/Hybrid, SIMD-accelerated, zero dependencies
Documentation
//! Frame/chunk-parallel Opus encoding (R1) — the structural win that beats a
//! single-threaded libopus on wall-clock.
//!
//! Opus carries real inter-frame state (SILK LTP/NSQ/NLSF/entropy, CELT
//! pre-emphasis/overlap/prefilter/energy, the HP filter and input resampler), so
//! a frame range cannot be encoded byte-identically from a cold encoder the way
//! AAC/Vorbis frames can. Instead each worker **primes** its encoder by
//! re-encoding `warmup` frames *before* its chunk (output discarded), which
//! converges the state to the true continuous state — a stable encoder forgets
//! its initial conditions over a few frames. The primed boundary is
//! perceptually-neutral (PEAQ ΔODG ≤ 0.03 vs serial), not byte-identical, so this
//! is an opt-in fast path, gated perceptually.
//!
//! Deterministic: fixed chunk boundaries → identical output across runs. Uses
//! only `std::thread` (no rayon).

use crate::{Application, OpusEncoder};

/// Worker count for `requested` (0 = all cores). wasm without the `atomics`
/// feature has no threads (`std::thread::spawn` panics there), so it is always
/// 1 and every entry point below runs serially on the calling thread.
fn worker_count(requested: usize) -> usize {
    if cfg!(all(target_family = "wasm", not(target_feature = "atomics"))) {
        return 1;
    }
    if requested == 0 {
        std::thread::available_parallelism().map_or(1, std::num::NonZero::get)
    } else {
        requested
    }
}

/// Configuration for a parallel encode; mirrors the knobs on [`OpusEncoder`].
#[derive(Clone, Copy)]
pub struct ParallelConfig {
    /// Input sampling rate in Hz (8000, 12000, 16000, 24000 or 48000).
    pub sample_rate: i32,
    /// Channel count (1 or 2).
    pub channels: usize,
    /// Encoder application mode.
    pub application: Application,
    /// Target bitrate in bits per second.
    pub bitrate_bps: i32,
    /// Encoder complexity, 0-10.
    pub complexity: i32,
    /// Constant bitrate when `true`.
    pub use_cbr: bool,
    /// Frames of look-back each worker re-encodes to prime its state (discarded).
    /// Must exceed the deepest inter-frame memory (SILK LTP lag + NSQ delay +
    /// CELT overlap). 8 (~160 ms @20 ms frames) is a safe default; sweep down
    /// under the PEAQ gate. `0` = no priming (equivalent to naive chunking).
    pub warmup: usize,
    /// Worker count; `0` selects `available_parallelism`.
    pub threads: usize,
}

impl ParallelConfig {
    /// Defaults: 64 kb/s VBR, complexity 9, 8 warm-up frames, one worker per core.
    pub fn new(sample_rate: i32, channels: usize, application: Application) -> Self {
        Self {
            sample_rate,
            channels,
            application,
            bitrate_bps: 64_000,
            complexity: 9,
            use_cbr: false,
            warmup: 8,
            threads: 0,
        }
    }
}

/// Encode `pcm` (interleaved f32, `channels`-interleaved) in `frame_size`
/// samples-per-channel frames, across `cfg.threads` workers, returning one Opus
/// packet per frame in order. Falls back to a single serial encoder when the
/// input is too small to split usefully.
///
/// The serial equivalent is `encode_serial`; this returns the same *count* of
/// packets and (with adequate `warmup`) a perceptually-identical bitstream.
///
/// # Panics
///
/// Panics if `cfg` is not a valid encoder configuration (see
/// [`crate::OpusEncoder::new`]) or if encoding a frame fails, which can only
/// happen for an invalid configuration. Validate the configuration with
/// `OpusEncoder::new` first when it comes from untrusted input.
/// Also panics if a worker thread panics.
pub fn encode_parallel(cfg: &ParallelConfig, pcm: &[f32], frame_size: usize) -> Vec<Vec<u8>> {
    let step = frame_size * cfg.channels;
    if step == 0 {
        return Vec::new();
    }
    let total_frames = pcm.len() / step;
    if total_frames == 0 {
        return Vec::new();
    }

    let threads = worker_count(cfg.threads);

    // Each chunk must be ≫ warmup to keep the redundant-compute overhead small;
    // require chunk ≥ 4·warmup (and ≥ 1). Cap the worker count accordingly.
    let min_chunk = (cfg.warmup * 4).max(1);
    let n_workers = threads.max(1).min((total_frames / min_chunk).max(1));
    if n_workers <= 1 {
        return encode_serial(cfg, pcm, frame_size);
    }

    // Contiguous, balanced frame ranges [start, end).
    let base = total_frames / n_workers;
    let rem = total_frames % n_workers;
    let mut ranges = Vec::with_capacity(n_workers);
    let mut start = 0usize;
    for w in 0..n_workers {
        let len = base + usize::from(w < rem);
        ranges.push((start, start + len));
        start += len;
    }

    let mut chunks: Vec<Vec<Vec<u8>>> = Vec::new();
    std::thread::scope(|scope| {
        let handles: Vec<_> = ranges
            .iter()
            .map(|&(cstart, cend)| {
                scope.spawn(move || encode_chunk(cfg, pcm, frame_size, cstart, cend))
            })
            .collect();
        for h in handles {
            chunks.push(h.join().expect("opus parallel worker panicked"));
        }
    });

    // Concatenate in range order.
    let mut out = Vec::with_capacity(total_frames);
    for c in chunks {
        out.extend(c);
    }
    out
}

/// Encode frames `[cstart, cend)` with a fresh encoder primed by re-encoding the
/// `warmup` frames before `cstart` (their packets discarded).
fn encode_chunk(
    cfg: &ParallelConfig,
    pcm: &[f32],
    frame_size: usize,
    cstart: usize,
    cend: usize,
) -> Vec<Vec<u8>> {
    let step = frame_size * cfg.channels;
    let mut enc = new_encoder(cfg);
    let warm_start = cstart.saturating_sub(cfg.warmup);
    let mut buf = vec![0u8; 4000];
    let mut packets = Vec::with_capacity(cend - cstart);
    for f in warm_start..cend {
        let frame = &pcm[f * step..(f + 1) * step];
        let n = enc
            .encode(frame, frame_size, &mut buf)
            .expect("opus encode");
        if f >= cstart {
            packets.push(buf[..n].to_vec());
        }
    }
    packets
}

/// Single-threaded reference: encode every frame with one continuous encoder.
/// The correctness/quality anchor for [`encode_parallel`].
///
/// # Panics
///
/// Panics if `cfg` is not a valid encoder configuration (see
/// [`crate::OpusEncoder::new`]) or if encoding a frame fails, which can only
/// happen for an invalid configuration. Validate the configuration with
/// `OpusEncoder::new` first when it comes from untrusted input.
pub fn encode_serial(cfg: &ParallelConfig, pcm: &[f32], frame_size: usize) -> Vec<Vec<u8>> {
    let step = frame_size * cfg.channels;
    if step == 0 {
        return Vec::new();
    }
    let total_frames = pcm.len() / step;
    let mut enc = new_encoder(cfg);
    let mut buf = vec![0u8; 4000];
    let mut packets = Vec::with_capacity(total_frames);
    for f in 0..total_frames {
        let frame = &pcm[f * step..(f + 1) * step];
        let n = enc
            .encode(frame, frame_size, &mut buf)
            .expect("opus encode");
        packets.push(buf[..n].to_vec());
    }
    packets
}

/// **R1a — per-stream parallelism (byte-identical).** Encode several *independent*
/// PCM streams concurrently, one serial encoder per worker. Each stream's output
/// is exactly its serial encode (no chunk seams), so this is **bit-identical** to
/// encoding them one-by-one — the right tool for batch/many-stream workloads
/// (and for short streams too small to split internally with [`encode_parallel`]).
///
/// `streams[i]` is `(config, pcm, frame_size)`; returns `out[i]` = that stream's
/// packets. Order preserved. Uses a bounded pool (`threads`, or all cores) so a
/// thousand tiny streams don't spawn a thousand threads.
///
/// # Panics
///
/// Panics if `cfg` is not a valid encoder configuration (see
/// [`crate::OpusEncoder::new`]) or if encoding a frame fails, which can only
/// happen for an invalid configuration. Validate the configuration with
/// `OpusEncoder::new` first when it comes from untrusted input.
/// Also panics if a worker thread panics.
pub fn encode_streams(
    streams: &[(ParallelConfig, &[f32], usize)],
    threads: usize,
) -> Vec<Vec<Vec<u8>>> {
    let n = streams.len();
    let mut out: Vec<Vec<Vec<u8>>> = (0..n).map(|_| Vec::new()).collect();
    if n == 0 {
        return out;
    }
    let workers = worker_count(threads).max(1).min(n);
    if workers == 1 {
        // No pool for one worker (and no threads at all on plain wasm).
        for ((cfg, pcm, frame_size), dst) in streams.iter().zip(out.iter_mut()) {
            *dst = encode_serial(cfg, pcm, *frame_size);
        }
        return out;
    }

    let next = std::sync::atomic::AtomicUsize::new(0);
    let out_slots: Vec<std::sync::Mutex<Option<Vec<Vec<u8>>>>> =
        (0..n).map(|_| std::sync::Mutex::new(None)).collect();
    std::thread::scope(|scope| {
        for _ in 0..workers {
            let next = &next;
            let out_slots = &out_slots;
            scope.spawn(move || {
                loop {
                    let idx = next.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
                    if idx >= n {
                        break;
                    }
                    let (cfg, pcm, frame_size) = &streams[idx];
                    let pkts = encode_serial(cfg, pcm, *frame_size);
                    *out_slots[idx].lock().unwrap() = Some(pkts);
                }
            });
        }
    });
    for (slot, dst) in out_slots.into_iter().zip(out.iter_mut()) {
        *dst = slot.into_inner().unwrap().unwrap_or_default();
    }
    out
}

fn new_encoder(cfg: &ParallelConfig) -> OpusEncoder {
    let mut enc = OpusEncoder::new(cfg.sample_rate, cfg.channels, cfg.application)
        .expect("opus encoder init");
    enc.bitrate_bps = cfg.bitrate_bps;
    enc.complexity = cfg.complexity.clamp(0, 10);
    enc.use_cbr = cfg.use_cbr;
    enc
}