Skip to main content

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}