opus_pure/parallel.rs
1//! Chunk-parallel Opus encoding: split the input into contiguous frame ranges
2//! and encode them on separate threads.
3//!
4//! Opus carries real inter-frame state — SILK's LTP, noise shaping, NLSF
5//! interpolation and bit reservoir, CELT's pre-emphasis, overlap, prefilter and
6//! energy prediction, the high-pass filter, the input resampler, and the content
7//! analysis that chooses between them — so a frame range cannot be encoded from
8//! a cold encoder and dropped into the middle of a stream. Each worker therefore
9//! **primes** its encoder by re-encoding the audio immediately before its chunk
10//! and discarding those packets, so that by the time it reaches its own first
11//! frame its state approximates the state a continuous encoder would have had.
12//!
13//! Priming is neither free nor exact, and what it costs and buys is the whole
14//! subject of this module.
15//!
16//! # What priming converges, and what it does not
17//!
18//! The signal path converges quickly. Its memory is a handful of frames — the
19//! LTP lag, the noise-shaping delay, one MDCT overlap — and a stable filter
20//! forgets its initial conditions. Around 160 ms of priming settles it.
21//!
22//! **The content analysis does not.** It keeps a hundred-entry ring of 20 ms
23//! observations and averages the music/speech probability over it, and the
24//! encoder applies hysteresis on top of that when it picks between SILK, hybrid
25//! and CELT. That memory is two seconds deep, and until it fills, a worker's
26//! mode decision is its own rather than the one a continuous encoder would have
27//! made. The consequence is not a seam at the boundary: it is the **entire
28//! chunk** coded in a different mode.
29//!
30//! Measured on 120 s of synthetic speech at 16 kb/s split four ways, where a
31//! continuous encoder settles on CELT and stays there:
32//!
33//! | `warmup_ms` | worst chunk, frames in a mode the serial encoder did not use | worst frame vs serial |
34//! | --- | --- | --- |
35//! | 160 | 36 of 1500 | −14.82 dB |
36//! | 500 | 16 of 1500 | −14.53 dB |
37//! | 1000 | 0 of 1500 | −4.10 dB |
38//! | 2000 (the default) | 0 of 1500 | −4.10 dB |
39//!
40//! Hence [`DEFAULT_WARMUP_MS`], which is the analysis's own history depth rather
41//! than a tuned number. A caller who pins [`ParallelConfig::signal_type`] takes
42//! the analysis out of the mode decision entirely and can prime far less; that
43//! is the cheapest way to buy a short warm-up.
44//!
45//! # What is left after priming
46//!
47//! Even fully primed, a chunked encode is not the serial encode. The rate
48//! controllers — CELT's VBR reservoir and drift, SILK's bit reservoir — are
49//! deliberately long-memory integrators, and a worker's differs from the
50//! continuous encoder's for some frames after its boundary. That shows as a
51//! bitrate dip of a few percent lasting tens of frames, and as a per-frame SNR
52//! difference at the boundary of a few dB against a serial encode of the same
53//! audio. Constant bitrate removes nearly all of it, there being no reservoir to
54//! be wrong about.
55//!
56//! This part is inherent rather than a defect awaiting a fix: at a chunk
57//! boundary one packet was produced by an encoder that did not produce the
58//! packet before it, and no amount of priming changes that. It is why this is an
59//! opt-in path and not what [`crate::OpusEncoder`] does by itself.
60//!
61//! `reference/parallel/` measures all of the above, and is where the numbers
62//! here come from.
63//!
64//! # Cost
65//!
66//! Every worker but the first re-encodes `warmup_ms` of audio it will not emit.
67//! With `w` workers that is `(w - 1) * warmup_ms` of redundant encoding, and
68//! [`ParallelConfig::plan`] reports it before any of it is done. The worker
69//! count is capped so redundancy stays at or below a quarter of the useful work,
70//! which with the default warm-up means one worker per 8 s of audio.
71//!
72//! Deterministic: fixed chunk boundaries mean identical output across runs. Uses
73//! only `std::thread`.
74
75use crate::analysis::DETECT_SIZE;
76use crate::{Application, Bandwidth, OpusEncoder, RateControl, Result, Signal};
77
78/// Default priming length, in milliseconds.
79///
80/// This is the depth of the content analysis's history: `DETECT_SIZE`
81/// observations of 20 ms each, from `src/analysis.rs`. Priming for less leaves a
82/// worker choosing its coding mode from a partly-filled analysis, which is what
83/// [the module documentation](self#what-priming-converges-and-what-it-does-not)
84/// measures.
85pub const DEFAULT_WARMUP_MS: u32 = DETECT_SIZE as u32 * 20;
86
87/// The redundancy ceiling that caps the worker count: a chunk is never shorter
88/// than this many warm-ups, which holds re-encoded audio at or below a quarter
89/// of the useful work.
90const MIN_CHUNK_WARMUPS: usize = 4;
91
92/// The largest packet [`OpusEncoder::encode`] can produce.
93///
94/// Re-exported here under its old local name because a worker allocates one of
95/// these once and reuses it for every frame; see [`MAX_PACKET_BYTES`] for why
96/// the size has to be exact rather than merely generous.
97use crate::encoder::MAX_PACKET_BYTES as MAX_PACKET;
98
99/// Configuration for a parallel encode.
100///
101/// Every setting [`OpusEncoder`] exposes appears here, because a worker's
102/// encoder is built from this struct and nothing else: a setting missing from
103/// this list is one the parallel path cannot reach at all. The defaults are
104/// [`OpusEncoder`]'s own, apart from `bitrate_bps` and `complexity`, which have
105/// none there.
106#[derive(Debug, Clone, Copy)]
107#[non_exhaustive]
108pub struct ParallelConfig {
109 /// Rate of the PCM being encoded, as [`OpusEncoder::new`] takes it.
110 pub sample_rate: i32,
111 /// Channels in that PCM, interleaved. 1 or 2.
112 pub channels: usize,
113 /// What to optimise for, as [`OpusEncoder::new`] takes it.
114 pub application: Application,
115 /// [`OpusEncoder::bitrate_bps`]. Defaults to 64000 here, where the encoder
116 /// itself has no default of its own.
117 pub bitrate_bps: i32,
118 /// [`OpusEncoder::complexity`]. Defaults to 9 here, where the encoder
119 /// itself has no default of its own.
120 pub complexity: i32,
121 /// [`OpusEncoder::rate_control`].
122 pub rate_control: RateControl,
123 /// [`OpusEncoder::use_inband_fec`].
124 pub use_inband_fec: bool,
125 /// Discontinuous transmission. Its trigger counts *consecutive* inactive
126 /// frames and each worker starts that count at zero, so a chunked encode
127 /// emits fewer DTX packets than a serial one over the same silence.
128 pub use_dtx: bool,
129 /// [`OpusEncoder::packet_loss_perc`].
130 pub packet_loss_perc: i32,
131 /// Pinning this, or [`Self::signal_type`], takes the content analysis out of
132 /// the coding-mode decision and lets `warmup_ms` be much shorter.
133 pub force_bandwidth: Option<Bandwidth>,
134 /// See [`Self::force_bandwidth`].
135 pub signal_type: Option<Signal>,
136 /// [`OpusEncoder::max_bandwidth`]. Unlike [`Self::force_bandwidth`] this
137 /// only caps the automatic choice, so it does not remove the analysis from
138 /// the decision and does not shorten the warm-up.
139 pub max_bandwidth: Bandwidth,
140 /// [`OpusEncoder::lsb_depth`].
141 pub lsb_depth: i32,
142 /// Audio each worker re-encodes before its own chunk to prime encoder state,
143 /// then discards. Milliseconds rather than frames because what it has to
144 /// cover is a time constant; see [`DEFAULT_WARMUP_MS`]. `0` disables priming,
145 /// which is naive chunking and audibly wrong.
146 pub warmup_ms: u32,
147 /// Worker count; `0` selects `available_parallelism`. Capped by the
148 /// redundancy ceiling — ask [`Self::plan`] what will actually be used.
149 pub threads: usize,
150}
151
152impl ParallelConfig {
153 /// A configuration with the defaults described on this type.
154 ///
155 /// The three arguments are the ones [`OpusEncoder::new`] fixes for an
156 /// encoder's life; the remaining fields are ordinary and can be set on the
157 /// returned value before it is handed to [`encode_parallel`].
158 pub fn new(sample_rate: i32, channels: usize, application: Application) -> Self {
159 ParallelConfig {
160 sample_rate,
161 channels,
162 application,
163 bitrate_bps: 64_000,
164 complexity: 9,
165 rate_control: RateControl::ConstrainedVbr,
166 use_inband_fec: false,
167 use_dtx: false,
168 packet_loss_perc: 0,
169 force_bandwidth: None,
170 signal_type: None,
171 max_bandwidth: Bandwidth::Fullband,
172 lsb_depth: 24,
173 warmup_ms: DEFAULT_WARMUP_MS,
174 threads: 0,
175 }
176 }
177
178 /// Warm-up expressed in frames of `frame_size` samples per channel, rounded
179 /// up so a frame duration that does not divide `warmup_ms` still covers it.
180 pub fn warmup_frames(&self, frame_size: usize) -> usize {
181 if frame_size == 0 || self.warmup_ms == 0 || self.sample_rate <= 0 {
182 return 0;
183 }
184 let frame_us = (frame_size as u64 * 1_000_000) / self.sample_rate as u64;
185 if frame_us == 0 {
186 return 0;
187 }
188 ((self.warmup_ms as u64 * 1000).div_ceil(frame_us)) as usize
189 }
190
191 /// How [`encode_parallel`] will divide `total_frames`, without encoding
192 /// anything.
193 ///
194 /// Worth asking before a large job: the worker count is capped by the
195 /// redundancy ceiling, so a short clip or a long warm-up can leave far fewer
196 /// threads in use than were requested, or fall back to serial encoding
197 /// altogether.
198 pub fn plan(&self, total_frames: usize, frame_size: usize) -> ParallelPlan {
199 let warmup = self.warmup_frames(frame_size);
200 let requested = if self.threads == 0 {
201 std::thread::available_parallelism()
202 .map(|n| n.get())
203 .unwrap_or(1)
204 } else {
205 self.threads
206 };
207 let min_chunk = (warmup * MIN_CHUNK_WARMUPS).max(1);
208 let workers = requested.max(1).min((total_frames / min_chunk).max(1));
209
210 // Contiguous, balanced frame ranges [start, end).
211 let mut ranges = Vec::with_capacity(workers);
212 if total_frames > 0 {
213 let (base, rem) = (total_frames / workers, total_frames % workers);
214 let mut start = 0usize;
215 for w in 0..workers {
216 let len = base + usize::from(w < rem);
217 ranges.push((start, start + len));
218 start += len;
219 }
220 }
221 // Worker 0 begins at frame 0 with nothing before it to prime from.
222 let redundant_frames = ranges.iter().map(|&(start, _)| start.min(warmup)).sum();
223
224 ParallelPlan {
225 workers,
226 warmup_frames: warmup,
227 redundant_frames,
228 ranges,
229 }
230 }
231}
232
233/// How a parallel encode will be divided up, from [`ParallelConfig::plan`].
234#[derive(Clone, Debug, PartialEq, Eq)]
235#[non_exhaustive]
236pub struct ParallelPlan {
237 /// Threads that will actually run. `1` means the encode is serial, whatever
238 /// was requested.
239 pub workers: usize,
240 /// Frames of priming each worker after the first re-encodes and discards.
241 pub warmup_frames: usize,
242 /// Total frames that will be encoded and thrown away.
243 pub redundant_frames: usize,
244 /// Half-open `[start, end)` frame range per worker, in output order.
245 pub ranges: Vec<(usize, usize)>,
246}
247
248impl ParallelPlan {
249 /// Redundant work as a fraction of the useful work; `0.0` when serial.
250 pub fn overhead(&self) -> f64 {
251 let useful: usize = self.ranges.iter().map(|&(s, e)| e - s).sum();
252 if useful == 0 {
253 0.0
254 } else {
255 self.redundant_frames as f64 / useful as f64
256 }
257 }
258}
259
260/// Encode `pcm` (interleaved f32, `channels`-interleaved) in `frame_size`
261/// samples-per-channel frames across several threads, returning one Opus packet
262/// per frame in order.
263///
264/// Output is deterministic and always the same length as a serial encode, but it
265/// is a *different* encode: see [the module documentation](self) for what
266/// differs and by how much. Falls back to a single encoder when the input is too
267/// short to split, which [`ParallelConfig::plan`] will say in advance.
268pub fn encode_parallel(
269 cfg: &ParallelConfig,
270 pcm: &[f32],
271 frame_size: usize,
272) -> Result<Vec<Vec<u8>>> {
273 let step = frame_size * cfg.channels;
274 if step == 0 {
275 return Ok(Vec::new());
276 }
277 let total_frames = pcm.len() / step;
278 if total_frames == 0 {
279 return Ok(Vec::new());
280 }
281
282 let plan = cfg.plan(total_frames, frame_size);
283 if plan.workers <= 1 {
284 return encode_serial(cfg, pcm, frame_size);
285 }
286 let warmup = plan.warmup_frames;
287
288 let mut chunks: Vec<Result<Vec<Vec<u8>>>> = Vec::with_capacity(plan.workers);
289 std::thread::scope(|scope| {
290 let handles: Vec<_> = plan
291 .ranges
292 .iter()
293 .map(|&(cstart, cend)| {
294 scope.spawn(move || encode_chunk(cfg, pcm, frame_size, warmup, cstart, cend))
295 })
296 .collect();
297 for h in handles {
298 // A worker returning `Err` is the caller's argument coming back and
299 // is propagated below. A worker that *panicked* is a bug in this
300 // crate, and `resume_unwind` carries the original payload and
301 // message up rather than replacing it with one of ours.
302 match h.join() {
303 Ok(r) => chunks.push(r),
304 Err(payload) => std::panic::resume_unwind(payload),
305 }
306 }
307 });
308
309 // Concatenate in range order. A worker's error is the whole encode's:
310 // the chunks are contiguous audio, so a gap in the middle is not a result
311 // anybody can use.
312 let mut out = Vec::with_capacity(total_frames);
313 for c in chunks {
314 out.extend(c?);
315 }
316 Ok(out)
317}
318
319/// Encode frames `[cstart, cend)` with a fresh encoder primed by re-encoding the
320/// `warmup` frames before `cstart`, whose packets are discarded.
321fn encode_chunk(
322 cfg: &ParallelConfig,
323 pcm: &[f32],
324 frame_size: usize,
325 warmup: usize,
326 cstart: usize,
327 cend: usize,
328) -> Result<Vec<Vec<u8>>> {
329 let step = frame_size * cfg.channels;
330 let mut enc = new_encoder(cfg)?;
331 let warm_start = cstart.saturating_sub(warmup);
332 let mut buf = vec![0u8; MAX_PACKET];
333 let mut packets = Vec::with_capacity(cend - cstart);
334 for f in warm_start..cend {
335 let frame = &pcm[f * step..(f + 1) * step];
336 let n = enc.encode(frame, frame_size, &mut buf)?;
337 if f >= cstart {
338 packets.push(buf[..n].to_vec());
339 }
340 }
341 Ok(packets)
342}
343
344/// Single-threaded reference: every frame through one continuous encoder. The
345/// correctness and quality anchor for [`encode_parallel`].
346fn encode_serial(cfg: &ParallelConfig, pcm: &[f32], frame_size: usize) -> Result<Vec<Vec<u8>>> {
347 let step = frame_size * cfg.channels;
348 if step == 0 {
349 return Ok(Vec::new());
350 }
351 let total_frames = pcm.len() / step;
352 let mut enc = new_encoder(cfg)?;
353 let mut buf = vec![0u8; MAX_PACKET];
354 let mut packets = Vec::with_capacity(total_frames);
355 for f in 0..total_frames {
356 let frame = &pcm[f * step..(f + 1) * step];
357 let n = enc.encode(frame, frame_size, &mut buf)?;
358 packets.push(buf[..n].to_vec());
359 }
360 Ok(packets)
361}
362
363fn new_encoder(cfg: &ParallelConfig) -> Result<OpusEncoder> {
364 let mut enc = OpusEncoder::new(cfg.sample_rate, cfg.channels, cfg.application)?;
365 enc.bitrate_bps = cfg.bitrate_bps;
366 enc.complexity = cfg.complexity;
367 enc.rate_control = cfg.rate_control;
368 enc.use_inband_fec = cfg.use_inband_fec;
369 enc.use_dtx = cfg.use_dtx;
370 enc.packet_loss_perc = cfg.packet_loss_perc;
371 enc.force_bandwidth = cfg.force_bandwidth;
372 enc.signal_type = cfg.signal_type;
373 enc.max_bandwidth = cfg.max_bandwidth;
374 enc.lsb_depth = cfg.lsb_depth;
375 Ok(enc)
376}