rusty_opus/parallel.rs
1//! Frame/chunk-parallel Opus encoding (R1) — the structural win that beats a
2//! single-threaded libopus on wall-clock.
3//!
4//! Opus carries real inter-frame state (SILK LTP/NSQ/NLSF/entropy, CELT
5//! pre-emphasis/overlap/prefilter/energy, the HP filter and input resampler), so
6//! a frame range cannot be encoded byte-identically from a cold encoder the way
7//! AAC/Vorbis frames can. Instead each worker **primes** its encoder by
8//! re-encoding `warmup` frames *before* its chunk (output discarded), which
9//! converges the state to the true continuous state — a stable encoder forgets
10//! its initial conditions over a few frames. The primed boundary is
11//! perceptually-neutral (PEAQ ΔODG ≤ 0.03 vs serial), not byte-identical, so this
12//! is an opt-in fast path, gated perceptually.
13//!
14//! Deterministic: fixed chunk boundaries → identical output across runs. Uses
15//! only `std::thread` (no rayon).
16
17use crate::{Application, OpusEncoder};
18
19/// Worker count for `requested` (0 = all cores). wasm without the `atomics`
20/// feature has no threads (`std::thread::spawn` panics there), so it is always
21/// 1 and every entry point below runs serially on the calling thread.
22fn worker_count(requested: usize) -> usize {
23 if cfg!(all(target_family = "wasm", not(target_feature = "atomics"))) {
24 return 1;
25 }
26 if requested == 0 {
27 std::thread::available_parallelism().map_or(1, std::num::NonZero::get)
28 } else {
29 requested
30 }
31}
32
33/// Configuration for a parallel encode; mirrors the knobs on [`OpusEncoder`].
34#[derive(Clone, Copy)]
35pub struct ParallelConfig {
36 /// Input sampling rate in Hz (8000, 12000, 16000, 24000 or 48000).
37 pub sample_rate: i32,
38 /// Channel count (1 or 2).
39 pub channels: usize,
40 /// Encoder application mode.
41 pub application: Application,
42 /// Target bitrate in bits per second.
43 pub bitrate_bps: i32,
44 /// Encoder complexity, 0-10.
45 pub complexity: i32,
46 /// Constant bitrate when `true`.
47 pub use_cbr: bool,
48 /// Frames of look-back each worker re-encodes to prime its state (discarded).
49 /// Must exceed the deepest inter-frame memory (SILK LTP lag + NSQ delay +
50 /// CELT overlap). 8 (~160 ms @20 ms frames) is a safe default; sweep down
51 /// under the PEAQ gate. `0` = no priming (equivalent to naive chunking).
52 pub warmup: usize,
53 /// Worker count; `0` selects `available_parallelism`.
54 pub threads: usize,
55}
56
57impl ParallelConfig {
58 /// Defaults: 64 kb/s VBR, complexity 9, 8 warm-up frames, one worker per core.
59 pub fn new(sample_rate: i32, channels: usize, application: Application) -> Self {
60 Self {
61 sample_rate,
62 channels,
63 application,
64 bitrate_bps: 64_000,
65 complexity: 9,
66 use_cbr: false,
67 warmup: 8,
68 threads: 0,
69 }
70 }
71}
72
73/// Encode `pcm` (interleaved f32, `channels`-interleaved) in `frame_size`
74/// samples-per-channel frames, across `cfg.threads` workers, returning one Opus
75/// packet per frame in order. Falls back to a single serial encoder when the
76/// input is too small to split usefully.
77///
78/// The serial equivalent is `encode_serial`; this returns the same *count* of
79/// packets and (with adequate `warmup`) a perceptually-identical bitstream.
80///
81/// # Panics
82///
83/// Panics if `cfg` is not a valid encoder configuration (see
84/// [`crate::OpusEncoder::new`]) or if encoding a frame fails, which can only
85/// happen for an invalid configuration. Validate the configuration with
86/// `OpusEncoder::new` first when it comes from untrusted input.
87/// Also panics if a worker thread panics.
88pub fn encode_parallel(cfg: &ParallelConfig, pcm: &[f32], frame_size: usize) -> Vec<Vec<u8>> {
89 let step = frame_size * cfg.channels;
90 if step == 0 {
91 return Vec::new();
92 }
93 let total_frames = pcm.len() / step;
94 if total_frames == 0 {
95 return Vec::new();
96 }
97
98 let threads = worker_count(cfg.threads);
99
100 // Each chunk must be ≫ warmup to keep the redundant-compute overhead small;
101 // require chunk ≥ 4·warmup (and ≥ 1). Cap the worker count accordingly.
102 let min_chunk = (cfg.warmup * 4).max(1);
103 let n_workers = threads.max(1).min((total_frames / min_chunk).max(1));
104 if n_workers <= 1 {
105 return encode_serial(cfg, pcm, frame_size);
106 }
107
108 // Contiguous, balanced frame ranges [start, end).
109 let base = total_frames / n_workers;
110 let rem = total_frames % n_workers;
111 let mut ranges = Vec::with_capacity(n_workers);
112 let mut start = 0usize;
113 for w in 0..n_workers {
114 let len = base + usize::from(w < rem);
115 ranges.push((start, start + len));
116 start += len;
117 }
118
119 let mut chunks: Vec<Vec<Vec<u8>>> = Vec::new();
120 std::thread::scope(|scope| {
121 let handles: Vec<_> = ranges
122 .iter()
123 .map(|&(cstart, cend)| {
124 scope.spawn(move || encode_chunk(cfg, pcm, frame_size, cstart, cend))
125 })
126 .collect();
127 for h in handles {
128 chunks.push(h.join().expect("opus parallel worker panicked"));
129 }
130 });
131
132 // Concatenate in range order.
133 let mut out = Vec::with_capacity(total_frames);
134 for c in chunks {
135 out.extend(c);
136 }
137 out
138}
139
140/// Encode frames `[cstart, cend)` with a fresh encoder primed by re-encoding the
141/// `warmup` frames before `cstart` (their packets discarded).
142fn encode_chunk(
143 cfg: &ParallelConfig,
144 pcm: &[f32],
145 frame_size: usize,
146 cstart: usize,
147 cend: usize,
148) -> Vec<Vec<u8>> {
149 let step = frame_size * cfg.channels;
150 let mut enc = new_encoder(cfg);
151 let warm_start = cstart.saturating_sub(cfg.warmup);
152 let mut buf = vec![0u8; 4000];
153 let mut packets = Vec::with_capacity(cend - cstart);
154 for f in warm_start..cend {
155 let frame = &pcm[f * step..(f + 1) * step];
156 let n = enc
157 .encode(frame, frame_size, &mut buf)
158 .expect("opus encode");
159 if f >= cstart {
160 packets.push(buf[..n].to_vec());
161 }
162 }
163 packets
164}
165
166/// Single-threaded reference: encode every frame with one continuous encoder.
167/// The correctness/quality anchor for [`encode_parallel`].
168///
169/// # Panics
170///
171/// Panics if `cfg` is not a valid encoder configuration (see
172/// [`crate::OpusEncoder::new`]) or if encoding a frame fails, which can only
173/// happen for an invalid configuration. Validate the configuration with
174/// `OpusEncoder::new` first when it comes from untrusted input.
175pub fn encode_serial(cfg: &ParallelConfig, pcm: &[f32], frame_size: usize) -> Vec<Vec<u8>> {
176 let step = frame_size * cfg.channels;
177 if step == 0 {
178 return Vec::new();
179 }
180 let total_frames = pcm.len() / step;
181 let mut enc = new_encoder(cfg);
182 let mut buf = vec![0u8; 4000];
183 let mut packets = Vec::with_capacity(total_frames);
184 for f in 0..total_frames {
185 let frame = &pcm[f * step..(f + 1) * step];
186 let n = enc
187 .encode(frame, frame_size, &mut buf)
188 .expect("opus encode");
189 packets.push(buf[..n].to_vec());
190 }
191 packets
192}
193
194/// **R1a — per-stream parallelism (byte-identical).** Encode several *independent*
195/// PCM streams concurrently, one serial encoder per worker. Each stream's output
196/// is exactly its serial encode (no chunk seams), so this is **bit-identical** to
197/// encoding them one-by-one — the right tool for batch/many-stream workloads
198/// (and for short streams too small to split internally with [`encode_parallel`]).
199///
200/// `streams[i]` is `(config, pcm, frame_size)`; returns `out[i]` = that stream's
201/// packets. Order preserved. Uses a bounded pool (`threads`, or all cores) so a
202/// thousand tiny streams don't spawn a thousand threads.
203///
204/// # Panics
205///
206/// Panics if `cfg` is not a valid encoder configuration (see
207/// [`crate::OpusEncoder::new`]) or if encoding a frame fails, which can only
208/// happen for an invalid configuration. Validate the configuration with
209/// `OpusEncoder::new` first when it comes from untrusted input.
210/// Also panics if a worker thread panics.
211pub fn encode_streams(
212 streams: &[(ParallelConfig, &[f32], usize)],
213 threads: usize,
214) -> Vec<Vec<Vec<u8>>> {
215 let n = streams.len();
216 let mut out: Vec<Vec<Vec<u8>>> = (0..n).map(|_| Vec::new()).collect();
217 if n == 0 {
218 return out;
219 }
220 let workers = worker_count(threads).max(1).min(n);
221 if workers == 1 {
222 // No pool for one worker (and no threads at all on plain wasm).
223 for ((cfg, pcm, frame_size), dst) in streams.iter().zip(out.iter_mut()) {
224 *dst = encode_serial(cfg, pcm, *frame_size);
225 }
226 return out;
227 }
228
229 let next = std::sync::atomic::AtomicUsize::new(0);
230 let out_slots: Vec<std::sync::Mutex<Option<Vec<Vec<u8>>>>> =
231 (0..n).map(|_| std::sync::Mutex::new(None)).collect();
232 std::thread::scope(|scope| {
233 for _ in 0..workers {
234 let next = &next;
235 let out_slots = &out_slots;
236 scope.spawn(move || {
237 loop {
238 let idx = next.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
239 if idx >= n {
240 break;
241 }
242 let (cfg, pcm, frame_size) = &streams[idx];
243 let pkts = encode_serial(cfg, pcm, *frame_size);
244 *out_slots[idx].lock().unwrap() = Some(pkts);
245 }
246 });
247 }
248 });
249 for (slot, dst) in out_slots.into_iter().zip(out.iter_mut()) {
250 *dst = slot.into_inner().unwrap().unwrap_or_default();
251 }
252 out
253}
254
255fn new_encoder(cfg: &ParallelConfig) -> OpusEncoder {
256 let mut enc = OpusEncoder::new(cfg.sample_rate, cfg.channels, cfg.application)
257 .expect("opus encoder init");
258 enc.bitrate_bps = cfg.bitrate_bps;
259 enc.complexity = cfg.complexity.clamp(0, 10);
260 enc.use_cbr = cfg.use_cbr;
261 enc
262}