mirror of
https://github.com/rustdesk/rustdesk.git
synced 2026-09-30 01:05:47 +00:00
* fix(audio): add streaming resampler * fix(audio): preserve stream resampling state * fix(audio): keep playback callback nonblocking * fix(audio): decouple capture conversion from dasp * fix(audio): support stateful samplerate backend * refactor(audio): isolate stream callback state * refactor(audio): group capture output options * fix(audio): clear stale playback state after startup failure Reset non-Linux playback state when stream startup fails to prevent new-format audio from using the previous stream or resampler. Add regression tests for failed format changes and successful playback. Signed-off-by: fufesou <linlong1266@gmail.com> * fix(audio): honor capture resampler selection and reuse buffers Use the selected resampling backend for fixed-frame capture. Convert samples directly into the input queue and reuse the PCM frame buffer. Add tests for anti-aliasing, thread transfer, and partial-frame draining. Signed-off-by: fufesou <linlong1266@gmail.com> * refact: reduce diffs Signed-off-by: fufesou <linlong1266@gmail.com> * test(audio): check resampler output count and passband energy Signed-off-by: fufesou <linlong1266@gmail.com> * fix(audio): reset incompatible Linux playback state on startup failure Preserve compatible output streams when replacement startup fails. Clear state when no compatible stream exists and cover both paths in tests. Signed-off-by: fufesou <linlong1266@gmail.com> * perf(audio): reuse PCM buffers in the capture pipeline - Reuse capture framing, resampling, and channel conversion buffers - Deliver borrowed packets and write Sinc output into reusable storage - Add allocation and output-equivalence regression tests Signed-off-by: fufesou <linlong1266@gmail.com> * fix(audio): smooth buffer discard discontinuities Signal receiver PCM discards and fade from the current playback output when the callback reaches the new timeline. Signed-off-by: fufesou <linlong1266@gmail.com> * fix(audio): add missing Cargo.toml Signed-off-by: fufesou <linlong1266@gmail.com> * perf(audio): move capture encoding off the CPAL callback Move Opus encoding and service delivery to a dedicated worker. Use a preallocated bounded PCM queue with explicit loss reporting. Add tests for callback allocations and queue saturation. Signed-off-by: fufesou <linlong1266@gmail.com> * fix(audio): smooth capture gaps and report losses during backlog Signed-off-by: fufesou <linlong1266@gmail.com> * feat(audio): report capture queue high-water mark Track peak queued PCM packets and log the approximate queued audio duration alongside capture loss statistics. Signed-off-by: fufesou <linlong1266@gmail.com> * refact(audio): reduce diffs Signed-off-by: fufesou <linlong1266@gmail.com> * fix(audio): avoid blocking capture on encoder queue contention Use preallocated queues with try_lock in the capture callback. Count and drop the current packet on contention, preserving drop-oldest behavior on overflow. Add regressions for paused workers, buffer reuse, and sequence wrap. Signed-off-by: fufesou <linlong1266@gmail.com> * fix: add the missing files Signed-off-by: fufesou <linlong1266@gmail.com> * fix(audio): isolate zero-gate state per encoder Signed-off-by: fufesou <linlong1266@gmail.com> * refact: reduce diffs Signed-off-by: fufesou <linlong1266@gmail.com> * refact(audio): simple refactor Signed-off-by: fufesou <linlong1266@gmail.com> * fix(audio): avoid waiting on playback callback locks Use one PCM try_lock attempt and preserve queued samples during contention. Replace readiness locking with per-stream atomic status and report callback errors from the receiving thread. Cover callback progress, retained audio, recovery, and poisoned-buffer handling. * fix(audio): restart capture after processing errors Stop further processing until the service recreates the stream. Document the guard as defensive recovery for an unconfirmed failure. Group capture and resampler submodules under their parent directories. Signed-off-by: fufesou <linlong1266@gmail.com> * audio: report capture queue contention drops separately - Add contention_dropped to loss reports while preserving total drop counts - Document packet rejection on contention even when buffers are available - Extend existing contention and saturation test assertions Signed-off-by: fufesou <linlong1266@gmail.com> * refact unit tests Signed-off-by: fufesou <linlong1266@gmail.com> --------- Signed-off-by: fufesou <linlong1266@gmail.com>
243 lines
8.8 KiB
Rust
243 lines
8.8 KiB
Rust
use hbb_common::thiserror;
|
|
|
|
#[cfg(test)]
|
|
pub(crate) mod allocation_tests;
|
|
|
|
#[cfg(all(feature = "use_samplerate", not(feature = "use_dasp")))]
|
|
mod sinc;
|
|
|
|
const INTERPOLATION_MARGIN_FRAMES: usize = 2;
|
|
const PENDING_PACKET_CAPACITY: usize = 2;
|
|
|
|
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
|
|
pub(crate) struct AudioResamplerConfig {
|
|
pub input_rate: u32,
|
|
pub output_rate: u32,
|
|
pub channels: u16,
|
|
}
|
|
|
|
#[derive(Debug, thiserror::Error, PartialEq, Eq)]
|
|
pub(crate) enum AudioResamplerError {
|
|
#[error(
|
|
"invalid audio resampler configuration: input_rate={}, output_rate={}, channels={}",
|
|
.0.input_rate, .0.output_rate, .0.channels
|
|
)]
|
|
InvalidConfig(AudioResamplerConfig),
|
|
#[error("invalid resampler output frame size: {output_frames}")]
|
|
InvalidOutputFrameSize { output_frames: usize },
|
|
#[error("audio resampler input length {samples} is not divisible by channel count {channels}")]
|
|
IncompleteFrame { samples: usize, channels: usize },
|
|
#[error("audio resampler output capacity overflow")]
|
|
CapacityOverflow,
|
|
#[cfg(all(feature = "use_samplerate", not(feature = "use_dasp")))]
|
|
#[error("audio resampler backend failed: {0}")]
|
|
Backend(String),
|
|
}
|
|
|
|
pub(crate) struct FixedFrameAudioResampler {
|
|
resampler: AudioResampler,
|
|
output_samples: usize,
|
|
pending_samples: Vec<f32>,
|
|
}
|
|
|
|
#[cfg(all(feature = "use_samplerate", not(feature = "use_dasp")))]
|
|
// SAFETY: libsamplerate's src_new state owns heap data and has no thread affinity.
|
|
// This wrapper never exposes or shares that state; processing requires &mut self.
|
|
unsafe impl Send for FixedFrameAudioResampler {}
|
|
|
|
impl FixedFrameAudioResampler {
|
|
pub(crate) fn new(
|
|
config: AudioResamplerConfig,
|
|
output_frames: usize,
|
|
) -> Result<Self, AudioResamplerError> {
|
|
if output_frames == 0 {
|
|
return Err(AudioResamplerError::InvalidOutputFrameSize { output_frames });
|
|
}
|
|
let channels = validate_config(config)?;
|
|
let output_samples = output_frames
|
|
.checked_mul(channels)
|
|
.ok_or(AudioResamplerError::CapacityOverflow)?;
|
|
let input_frames = output_frames
|
|
.checked_mul(config.input_rate as usize)
|
|
.ok_or(AudioResamplerError::CapacityOverflow)?
|
|
.div_ceil(config.output_rate as usize);
|
|
let capacity = output_samples
|
|
.checked_mul(PENDING_PACKET_CAPACITY)
|
|
.and_then(|samples| samples.checked_add(channels * INTERPOLATION_MARGIN_FRAMES))
|
|
.ok_or(AudioResamplerError::CapacityOverflow)?;
|
|
let mut resampler = AudioResampler::new(config)?;
|
|
resampler.reserve_input(input_frames)?;
|
|
Ok(Self {
|
|
resampler,
|
|
output_samples,
|
|
pending_samples: Vec::with_capacity(capacity),
|
|
})
|
|
}
|
|
|
|
pub(crate) fn process_with(
|
|
&mut self,
|
|
input: &[f32],
|
|
mut on_packet: impl FnMut(&[f32]),
|
|
) -> Result<(), AudioResamplerError> {
|
|
self.resampler
|
|
.process_into(input, &mut self.pending_samples)?;
|
|
let complete_samples =
|
|
self.pending_samples.len() / self.output_samples * self.output_samples;
|
|
for packet in self.pending_samples[..complete_samples].chunks_exact(self.output_samples) {
|
|
on_packet(packet);
|
|
}
|
|
self.pending_samples.drain(..complete_samples);
|
|
Ok(())
|
|
}
|
|
|
|
#[cfg(test)]
|
|
pub(crate) fn process(&mut self, input: &[f32]) -> Result<Vec<Vec<f32>>, AudioResamplerError> {
|
|
let mut packets = Vec::new();
|
|
self.process_with(input, |packet| packets.push(packet.to_owned()))?;
|
|
Ok(packets)
|
|
}
|
|
}
|
|
|
|
pub(crate) struct AudioResampler {
|
|
#[cfg(not(all(feature = "use_samplerate", not(feature = "use_dasp"))))]
|
|
backend: StreamingLinearAudioResampler,
|
|
#[cfg(all(feature = "use_samplerate", not(feature = "use_dasp")))]
|
|
backend: sinc::SincAudioResampler,
|
|
}
|
|
|
|
impl AudioResampler {
|
|
pub(crate) fn new(config: AudioResamplerConfig) -> Result<Self, AudioResamplerError> {
|
|
Ok(Self {
|
|
#[cfg(all(feature = "use_samplerate", not(feature = "use_dasp")))]
|
|
backend: sinc::SincAudioResampler::new(config)?,
|
|
#[cfg(not(all(feature = "use_samplerate", not(feature = "use_dasp"))))]
|
|
backend: StreamingLinearAudioResampler::new(config)?,
|
|
})
|
|
}
|
|
|
|
pub(crate) fn process(&mut self, input: &[f32]) -> Result<Vec<f32>, AudioResamplerError> {
|
|
let mut output = Vec::new();
|
|
self.process_into(input, &mut output)?;
|
|
Ok(output)
|
|
}
|
|
|
|
// Append samples so capture can retain an incomplete output packet in the same buffer.
|
|
fn process_into(
|
|
&mut self,
|
|
input: &[f32],
|
|
output: &mut Vec<f32>,
|
|
) -> Result<(), AudioResamplerError> {
|
|
self.backend.process_into(input, output)
|
|
}
|
|
|
|
fn reserve_input(&mut self, _frames: usize) -> Result<(), AudioResamplerError> {
|
|
#[cfg(not(all(feature = "use_samplerate", not(feature = "use_dasp"))))]
|
|
{
|
|
let capacity = _frames
|
|
.checked_add(INTERPOLATION_MARGIN_FRAMES)
|
|
.and_then(|frames| frames.checked_mul(self.backend.channels))
|
|
.ok_or(AudioResamplerError::CapacityOverflow)?;
|
|
self.backend.buffered_samples.reserve(capacity);
|
|
}
|
|
Ok(())
|
|
}
|
|
}
|
|
|
|
#[cfg(not(all(feature = "use_samplerate", not(feature = "use_dasp"))))]
|
|
struct StreamingLinearAudioResampler {
|
|
config: AudioResamplerConfig,
|
|
channels: usize,
|
|
buffered_samples: Vec<f32>,
|
|
next_position: u64,
|
|
}
|
|
|
|
#[cfg(not(all(feature = "use_samplerate", not(feature = "use_dasp"))))]
|
|
impl StreamingLinearAudioResampler {
|
|
fn new(config: AudioResamplerConfig) -> Result<Self, AudioResamplerError> {
|
|
Ok(Self {
|
|
config,
|
|
channels: validate_config(config)?,
|
|
buffered_samples: Vec::new(),
|
|
next_position: 0,
|
|
})
|
|
}
|
|
|
|
fn process_into(
|
|
&mut self,
|
|
input: &[f32],
|
|
output: &mut Vec<f32>,
|
|
) -> Result<(), AudioResamplerError> {
|
|
validate_input(input, self.channels)?;
|
|
let capacity = self.output_capacity(input.len())?;
|
|
output.reserve(capacity);
|
|
self.buffered_samples.extend_from_slice(input);
|
|
while self.write_next_frame(output) {
|
|
self.next_position += self.config.input_rate as u64;
|
|
}
|
|
self.discard_consumed_frames();
|
|
Ok(())
|
|
}
|
|
|
|
fn output_capacity(&self, input_samples: usize) -> Result<usize, AudioResamplerError> {
|
|
let input_frames = input_samples / self.channels;
|
|
let scaled_frames = input_frames
|
|
.checked_mul(self.config.output_rate as usize)
|
|
.ok_or(AudioResamplerError::CapacityOverflow)?
|
|
/ self.config.input_rate as usize;
|
|
scaled_frames
|
|
.checked_add(INTERPOLATION_MARGIN_FRAMES)
|
|
.and_then(|frames| frames.checked_mul(self.channels))
|
|
.ok_or(AudioResamplerError::CapacityOverflow)
|
|
}
|
|
|
|
fn write_next_frame(&self, output: &mut Vec<f32>) -> bool {
|
|
let output_rate = self.config.output_rate as u64;
|
|
let frame_count = self.buffered_samples.len() / self.channels;
|
|
let frame = (self.next_position / output_rate) as usize;
|
|
let fraction = self.next_position % output_rate;
|
|
if frame >= frame_count || (fraction != 0 && frame + 1 >= frame_count) {
|
|
return false;
|
|
}
|
|
let weight = fraction as f32 / output_rate as f32;
|
|
for channel in 0..self.channels {
|
|
let current = self.buffered_samples[frame * self.channels + channel];
|
|
let next_frame = frame + usize::from(fraction != 0);
|
|
let next = self.buffered_samples[next_frame * self.channels + channel];
|
|
output.push(current + (next - current) * weight);
|
|
}
|
|
true
|
|
}
|
|
|
|
fn discard_consumed_frames(&mut self) {
|
|
let output_rate = self.config.output_rate as u64;
|
|
let available_frames = self.buffered_samples.len() / self.channels;
|
|
let consumed_frames = ((self.next_position / output_rate) as usize).min(available_frames);
|
|
self.buffered_samples
|
|
.drain(0..consumed_frames * self.channels);
|
|
self.next_position -= consumed_frames as u64 * output_rate;
|
|
}
|
|
}
|
|
|
|
fn validate_config(config: AudioResamplerConfig) -> Result<usize, AudioResamplerError> {
|
|
if config.input_rate == 0 || config.output_rate == 0 || config.channels == 0 {
|
|
return Err(AudioResamplerError::InvalidConfig(config));
|
|
}
|
|
Ok(config.channels as usize)
|
|
}
|
|
|
|
fn validate_input(input: &[f32], channels: usize) -> Result<(), AudioResamplerError> {
|
|
if input.len() % channels != 0 {
|
|
return Err(AudioResamplerError::IncompleteFrame {
|
|
samples: input.len(),
|
|
channels,
|
|
});
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
#[cfg(all(test, not(all(feature = "use_samplerate", not(feature = "use_dasp")))))]
|
|
mod tests;
|
|
|
|
#[cfg(all(test, feature = "use_samplerate", not(feature = "use_dasp")))]
|
|
mod samplerate_tests;
|