use crate::analysis::DETECT_SIZE;
use crate::{Application, Bandwidth, OpusEncoder, RateControl, Result, Signal};
pub const DEFAULT_WARMUP_MS: u32 = DETECT_SIZE as u32 * 20;
const MIN_CHUNK_WARMUPS: usize = 4;
use crate::encoder::MAX_PACKET_BYTES as MAX_PACKET;
#[derive(Debug, Clone, Copy)]
#[non_exhaustive]
pub struct ParallelConfig {
pub sample_rate: i32,
pub channels: usize,
pub application: Application,
pub bitrate_bps: i32,
pub complexity: i32,
pub rate_control: RateControl,
pub use_inband_fec: bool,
pub use_dtx: bool,
pub packet_loss_perc: i32,
pub force_bandwidth: Option<Bandwidth>,
pub signal_type: Option<Signal>,
pub max_bandwidth: Bandwidth,
pub lsb_depth: i32,
pub warmup_ms: u32,
pub threads: usize,
}
impl ParallelConfig {
pub fn new(sample_rate: i32, channels: usize, application: Application) -> Self {
ParallelConfig {
sample_rate,
channels,
application,
bitrate_bps: 64_000,
complexity: 9,
rate_control: RateControl::ConstrainedVbr,
use_inband_fec: false,
use_dtx: false,
packet_loss_perc: 0,
force_bandwidth: None,
signal_type: None,
max_bandwidth: Bandwidth::Fullband,
lsb_depth: 24,
warmup_ms: DEFAULT_WARMUP_MS,
threads: 0,
}
}
pub fn warmup_frames(&self, frame_size: usize) -> usize {
if frame_size == 0 || self.warmup_ms == 0 || self.sample_rate <= 0 {
return 0;
}
let frame_us = (frame_size as u64 * 1_000_000) / self.sample_rate as u64;
if frame_us == 0 {
return 0;
}
((self.warmup_ms as u64 * 1000).div_ceil(frame_us)) as usize
}
pub fn plan(&self, total_frames: usize, frame_size: usize) -> ParallelPlan {
let warmup = self.warmup_frames(frame_size);
let requested = if self.threads == 0 {
std::thread::available_parallelism()
.map(|n| n.get())
.unwrap_or(1)
} else {
self.threads
};
let min_chunk = (warmup * MIN_CHUNK_WARMUPS).max(1);
let workers = requested.max(1).min((total_frames / min_chunk).max(1));
let mut ranges = Vec::with_capacity(workers);
if total_frames > 0 {
let (base, rem) = (total_frames / workers, total_frames % workers);
let mut start = 0usize;
for w in 0..workers {
let len = base + usize::from(w < rem);
ranges.push((start, start + len));
start += len;
}
}
let redundant_frames = ranges.iter().map(|&(start, _)| start.min(warmup)).sum();
ParallelPlan {
workers,
warmup_frames: warmup,
redundant_frames,
ranges,
}
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub struct ParallelPlan {
pub workers: usize,
pub warmup_frames: usize,
pub redundant_frames: usize,
pub ranges: Vec<(usize, usize)>,
}
impl ParallelPlan {
pub fn overhead(&self) -> f64 {
let useful: usize = self.ranges.iter().map(|&(s, e)| e - s).sum();
if useful == 0 {
0.0
} else {
self.redundant_frames as f64 / useful as f64
}
}
}
pub fn encode_parallel(
cfg: &ParallelConfig,
pcm: &[f32],
frame_size: usize,
) -> Result<Vec<Vec<u8>>> {
let step = frame_size * cfg.channels;
if step == 0 {
return Ok(Vec::new());
}
let total_frames = pcm.len() / step;
if total_frames == 0 {
return Ok(Vec::new());
}
let plan = cfg.plan(total_frames, frame_size);
if plan.workers <= 1 {
return encode_serial(cfg, pcm, frame_size);
}
let warmup = plan.warmup_frames;
let mut chunks: Vec<Result<Vec<Vec<u8>>>> = Vec::with_capacity(plan.workers);
std::thread::scope(|scope| {
let handles: Vec<_> = plan
.ranges
.iter()
.map(|&(cstart, cend)| {
scope.spawn(move || encode_chunk(cfg, pcm, frame_size, warmup, cstart, cend))
})
.collect();
for h in handles {
match h.join() {
Ok(r) => chunks.push(r),
Err(payload) => std::panic::resume_unwind(payload),
}
}
});
let mut out = Vec::with_capacity(total_frames);
for c in chunks {
out.extend(c?);
}
Ok(out)
}
fn encode_chunk(
cfg: &ParallelConfig,
pcm: &[f32],
frame_size: usize,
warmup: usize,
cstart: usize,
cend: usize,
) -> Result<Vec<Vec<u8>>> {
let step = frame_size * cfg.channels;
let mut enc = new_encoder(cfg)?;
let warm_start = cstart.saturating_sub(warmup);
let mut buf = vec![0u8; MAX_PACKET];
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)?;
if f >= cstart {
packets.push(buf[..n].to_vec());
}
}
Ok(packets)
}
fn encode_serial(cfg: &ParallelConfig, pcm: &[f32], frame_size: usize) -> Result<Vec<Vec<u8>>> {
let step = frame_size * cfg.channels;
if step == 0 {
return Ok(Vec::new());
}
let total_frames = pcm.len() / step;
let mut enc = new_encoder(cfg)?;
let mut buf = vec![0u8; MAX_PACKET];
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)?;
packets.push(buf[..n].to_vec());
}
Ok(packets)
}
fn new_encoder(cfg: &ParallelConfig) -> Result<OpusEncoder> {
let mut enc = OpusEncoder::new(cfg.sample_rate, cfg.channels, cfg.application)?;
enc.bitrate_bps = cfg.bitrate_bps;
enc.complexity = cfg.complexity;
enc.rate_control = cfg.rate_control;
enc.use_inband_fec = cfg.use_inband_fec;
enc.use_dtx = cfg.use_dtx;
enc.packet_loss_perc = cfg.packet_loss_perc;
enc.force_bandwidth = cfg.force_bandwidth;
enc.signal_type = cfg.signal_type;
enc.max_bandwidth = cfg.max_bandwidth;
enc.lsb_depth = cfg.lsb_depth;
Ok(enc)
}