use std::sync::atomic::{AtomicBool, AtomicU32, Ordering};
use std::sync::Arc;
use std::thread::{self, JoinHandle};
use std::time::Duration;
use cpal::traits::{DeviceTrait, HostTrait, StreamTrait};
use cpal::{SampleFormat, StreamConfig as CpalStreamConfig};
use crossbeam_queue::ArrayQueue;
use crate::resample::StreamingResampler;
use super::OscilloscopeConfig;
use crate::backend::{DacBackend, FifoBackend, WriteOutcome};
use crate::buffer_estimate::{BufferEstimator, QueueDepthSource, RuntimeAuthorityEstimator};
use crate::device::{DacCapabilities, DacType};
use crate::error::{Error, Result};
use crate::point::LaserPoint;
const MUTE_RAMP_MS: f32 = 3.0;
struct RuntimeState {
queue: ArrayQueue<(f32, f32)>,
muted: AtomicBool,
connected: AtomicBool,
sample_rate: u32,
last_l_bits: AtomicU32,
last_r_bits: AtomicU32,
}
impl RuntimeState {
fn new(capacity: usize, sample_rate: u32) -> Self {
Self {
queue: ArrayQueue::new(capacity),
muted: AtomicBool::new(true),
connected: AtomicBool::new(false),
sample_rate,
last_l_bits: AtomicU32::new(0.0f32.to_bits()),
last_r_bits: AtomicU32::new(0.0f32.to_bits()),
}
}
fn remaining_capacity(&self) -> usize {
self.queue.capacity().saturating_sub(self.queue.len())
}
fn has_capacity_for(&self, count: usize) -> bool {
count == 0 || self.remaining_capacity() >= count
}
fn queued_points(&self) -> u64 {
self.queue.len() as u64
}
fn clear_queue(&self) {
while self.queue.pop().is_some() {}
}
fn last_output(&self) -> (f32, f32) {
(
f32::from_bits(self.last_l_bits.load(Ordering::Relaxed)),
f32::from_bits(self.last_r_bits.load(Ordering::Relaxed)),
)
}
fn set_last_output(&self, l: f32, r: f32) {
self.last_l_bits.store(l.to_bits(), Ordering::Relaxed);
self.last_r_bits.store(r.to_bits(), Ordering::Relaxed);
}
fn next_output(&self) -> (f32, f32) {
let muted = self.muted.load(Ordering::Relaxed);
let queued = self.queue.pop();
let (last_l, last_r) = self.last_output();
let (l, r) = if muted {
let k = (1000.0 / (MUTE_RAMP_MS * self.sample_rate as f32)).clamp(0.0, 1.0);
(last_l + (0.0 - last_l) * k, last_r + (0.0 - last_r) * k)
} else {
match queued {
Some((l, r)) => (l, r),
None => (last_l, last_r),
}
};
self.set_last_output(l, r);
(l, r)
}
}
impl QueueDepthSource for RuntimeState {
fn queued_points(&self) -> u64 {
RuntimeState::queued_points(self)
}
fn sample_rate(&self) -> u32 {
self.sample_rate
}
}
struct AudioThread {
handle: JoinHandle<()>,
stop_flag: Arc<AtomicBool>,
}
pub struct OscilloscopeBackend {
device_name: String,
sample_rate: u32,
config: OscilloscopeConfig,
caps: DacCapabilities,
runtime: Option<Arc<RuntimeState>>,
audio_thread: Option<AudioThread>,
sample_buffer: Vec<(f32, f32)>,
estimator: RuntimeAuthorityEstimator,
resampler: StreamingResampler<(f32, f32)>,
engine: Arc<dyn AudioEngine>,
}
impl OscilloscopeBackend {
pub fn new(device_name: String, sample_rate: u32) -> Self {
Self::with_engine(device_name, sample_rate, Arc::new(CpalAudioEngine))
}
fn with_engine(device_name: String, sample_rate: u32, engine: Arc<dyn AudioEngine>) -> Self {
Self {
device_name,
sample_rate,
config: OscilloscopeConfig::default(),
caps: super::capabilities(sample_rate),
runtime: None,
audio_thread: None,
sample_buffer: Vec::new(),
estimator: RuntimeAuthorityEstimator::new(),
resampler: StreamingResampler::new(1, 1),
engine,
}
}
#[cfg(test)]
fn with_engine_for_test(
device_name: String,
sample_rate: u32,
engine: Arc<dyn AudioEngine>,
) -> Self {
Self::with_engine(device_name, sample_rate, engine)
}
pub fn set_config(&mut self, config: OscilloscopeConfig) {
self.config = config;
}
fn point_to_samples(p: &LaserPoint, config: &OscilloscopeConfig) -> (f32, f32) {
let mut l = p.x * config.gain + config.dc_offset;
let mut r = p.y * config.gain + config.dc_offset;
if config.clip {
l = l.clamp(-1.0, 1.0);
r = r.clamp(-1.0, 1.0);
}
(l, r)
}
fn buffer_capacity(&self) -> usize {
super::buffer_capacity(self.sample_rate)
}
fn start_audio_thread(&self, runtime: &Arc<RuntimeState>) -> Result<AudioThread> {
let device_name = self.device_name.clone();
let sample_rate = self.sample_rate;
let runtime = Arc::clone(runtime);
let engine = Arc::clone(&self.engine);
let stop_flag = Arc::new(AtomicBool::new(false));
let stop_flag_clone = Arc::clone(&stop_flag);
let handle = thread::Builder::new()
.name(format!("oscilloscope-{}", device_name))
.spawn(move || {
if let Err(e) = run_audio_thread(
&engine,
&device_name,
sample_rate,
&runtime,
&stop_flag_clone,
) {
log::error!("Oscilloscope audio thread error: {}", e);
}
runtime.connected.store(false, Ordering::Release);
})
.map_err(|e| {
Error::backend(std::io::Error::other(format!(
"Failed to spawn oscilloscope thread: {}",
e
)))
})?;
Ok(AudioThread { handle, stop_flag })
}
}
trait RunningAudioStream {}
trait AudioEngine: Send + Sync {
fn open_stream(
&self,
device_name: &str,
sample_rate: u32,
runtime: Arc<RuntimeState>,
) -> Result<Box<dyn RunningAudioStream>>;
}
struct CpalRunningStream {
_stream: cpal::Stream,
}
impl RunningAudioStream for CpalRunningStream {}
struct CpalAudioEngine;
impl AudioEngine for CpalAudioEngine {
fn open_stream(
&self,
device_name: &str,
sample_rate: u32,
runtime: Arc<RuntimeState>,
) -> Result<Box<dyn RunningAudioStream>> {
let host = cpal::default_host();
let device = host
.output_devices()
.map_err(|e| {
Error::backend(std::io::Error::other(format!(
"Failed to enumerate devices: {}",
e
)))
})?
.find(|d| d.name().map(|n| n == device_name).unwrap_or(false))
.ok_or_else(|| {
Error::disconnected(format!("Audio device '{}' not found", device_name))
})?;
let supported_config = device
.supported_output_configs()
.map_err(|e| {
Error::backend(std::io::Error::other(format!(
"Failed to get audio configs: {}",
e
)))
})?
.find(|c| {
c.channels() == 2
&& c.min_sample_rate().0 <= sample_rate
&& c.max_sample_rate().0 >= sample_rate
})
.ok_or_else(|| {
Error::invalid_config(format!(
"Audio device doesn't support stereo output at {} Hz",
sample_rate
))
})?;
let sample_format = supported_config.sample_format();
let config = CpalStreamConfig {
channels: 2,
sample_rate: cpal::SampleRate(sample_rate),
buffer_size: cpal::BufferSize::Default,
};
let runtime_err = Arc::clone(&runtime);
let err_callback = move |err: cpal::StreamError| {
log::error!("Oscilloscope stream error: {}", err);
if matches!(err, cpal::StreamError::DeviceNotAvailable) {
runtime_err.connected.store(false, Ordering::Release);
}
};
let stream = match sample_format {
SampleFormat::F32 => {
let rt = Arc::clone(&runtime);
device
.build_output_stream(
&config,
move |data: &mut [f32], _| fill_f32_output(data, &rt),
err_callback,
None,
)
.map_err(|e| {
Error::backend(std::io::Error::other(format!(
"Failed to build audio stream: {}",
e
)))
})?
}
SampleFormat::I16 => {
let rt = Arc::clone(&runtime);
device
.build_output_stream(
&config,
move |data: &mut [i16], _| fill_i16_output(data, &rt),
err_callback,
None,
)
.map_err(|e| {
Error::backend(std::io::Error::other(format!(
"Failed to build audio stream: {}",
e
)))
})?
}
format => {
return Err(Error::invalid_config(format!(
"Unsupported audio sample format: {:?}",
format
)));
}
};
stream.play().map_err(|e| {
Error::backend(std::io::Error::other(format!(
"Failed to start audio stream: {}",
e
)))
})?;
Ok(Box::new(CpalRunningStream { _stream: stream }))
}
}
fn fill_f32_output(data: &mut [f32], runtime: &RuntimeState) {
for chunk in data.chunks_mut(2) {
let (l, r) = runtime.next_output();
chunk[0] = l;
chunk[1] = r;
}
}
fn fill_i16_output(data: &mut [i16], runtime: &RuntimeState) {
for chunk in data.chunks_mut(2) {
let (l, r) = runtime.next_output();
chunk[0] = (l * i16::MAX as f32) as i16;
chunk[1] = (r * i16::MAX as f32) as i16;
}
}
fn run_audio_thread(
engine: &Arc<dyn AudioEngine>,
device_name: &str,
sample_rate: u32,
runtime: &Arc<RuntimeState>,
stop_flag: &AtomicBool,
) -> Result<()> {
let stream = engine.open_stream(device_name, sample_rate, Arc::clone(runtime))?;
runtime.connected.store(true, Ordering::Release);
log::info!(
"Oscilloscope started for '{}' at {} Hz",
device_name,
sample_rate
);
while !stop_flag.load(Ordering::Relaxed) {
thread::sleep(Duration::from_millis(10));
}
drop(stream);
log::info!("Oscilloscope stopped for '{}'", device_name);
Ok(())
}
impl DacBackend for OscilloscopeBackend {
fn dac_type(&self) -> DacType {
DacType::Oscilloscope
}
fn caps(&self) -> &DacCapabilities {
&self.caps
}
fn connect(&mut self) -> Result<()> {
if self.is_connected() {
return Ok(());
}
let runtime = Arc::new(RuntimeState::new(self.buffer_capacity(), self.sample_rate));
self.resampler.reset();
let audio_thread = self.start_audio_thread(&runtime)?;
let start = std::time::Instant::now();
while !runtime.connected.load(Ordering::Acquire) {
if start.elapsed() > Duration::from_secs(5) {
audio_thread.stop_flag.store(true, Ordering::Release);
let _ = audio_thread.handle.join();
return Err(Error::backend(std::io::Error::other(
"Timeout waiting for oscilloscope connection",
)));
}
thread::sleep(Duration::from_millis(10));
}
self.estimator
.set_source(Arc::clone(&runtime) as Arc<dyn QueueDepthSource>);
self.runtime = Some(runtime);
self.audio_thread = Some(audio_thread);
log::info!(
"Connected to oscilloscope '{}' at {} Hz",
self.device_name,
self.sample_rate
);
Ok(())
}
fn disconnect(&mut self) -> Result<()> {
if let Some(audio_thread) = self.audio_thread.take() {
audio_thread.stop_flag.store(true, Ordering::Release);
let _ = audio_thread.handle.join();
}
if let Some(runtime) = self.runtime.take() {
runtime.clear_queue();
}
self.estimator.clear_source();
self.resampler.reset();
Ok(())
}
fn is_connected(&self) -> bool {
self.runtime
.as_ref()
.is_some_and(|rt| rt.connected.load(Ordering::Relaxed))
}
fn stop(&mut self) -> Result<()> {
if let Some(runtime) = &self.runtime {
runtime.muted.store(true, Ordering::Release);
}
self.resampler.reset();
Ok(())
}
fn set_shutter(&mut self, open: bool) -> Result<()> {
if let Some(runtime) = &self.runtime {
runtime.muted.store(!open, Ordering::Release);
}
Ok(())
}
}
impl FifoBackend for OscilloscopeBackend {
fn try_write_points(&mut self, pps: u32, points: &[LaserPoint]) -> Result<WriteOutcome> {
let runtime = self
.runtime
.as_ref()
.ok_or_else(|| Error::disconnected("Not connected"))?;
if !runtime.connected.load(Ordering::Acquire) {
return Err(Error::disconnected("Not connected"));
}
if points.is_empty() {
return Ok(WriteOutcome::Written);
}
self.resampler.set_rates(pps.max(1), self.sample_rate);
let output_len = self.resampler.pending_output_count(points.len());
if !runtime.has_capacity_for(output_len) {
return Ok(WriteOutcome::WouldBlock);
}
self.sample_buffer.clear();
self.sample_buffer.extend(
points
.iter()
.map(|p| Self::point_to_samples(p, &self.config)),
);
let samples = std::mem::take(&mut self.sample_buffer);
self.resampler.process(&samples, |s| {
let _ = runtime.queue.push(s);
});
self.sample_buffer = samples;
Ok(WriteOutcome::Written)
}
fn estimator(&self) -> &dyn BufferEstimator {
&self.estimator
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::collections::VecDeque;
use std::sync::atomic::AtomicUsize;
use std::sync::mpsc::{self, RecvTimeoutError};
use std::sync::Mutex;
use std::time::Instant;
const RATE: u32 = 48_000;
#[test]
fn point_to_samples_maps_x_left_y_right_without_inversion() {
let cfg = OscilloscopeConfig::default();
let p = LaserPoint::new(0.25, -0.5, 65535, 0, 0, 65535);
let (l, r) = OscilloscopeBackend::point_to_samples(&p, &cfg);
assert!((l - 0.25).abs() < 1e-6);
assert!((r + 0.5).abs() < 1e-6);
}
#[test]
fn point_to_samples_applies_gain_and_dc_offset() {
let cfg = OscilloscopeConfig {
gain: 0.5,
dc_offset: 0.1,
clip: false,
};
let p = LaserPoint::new(0.4, -0.4, 0, 0, 0, 0);
let (l, r) = OscilloscopeBackend::point_to_samples(&p, &cfg);
assert!((l - 0.3).abs() < 1e-6); assert!((r - (-0.1)).abs() < 1e-6); }
#[test]
fn point_to_samples_clips_when_enabled() {
let cfg = OscilloscopeConfig {
gain: 4.0,
dc_offset: 0.0,
clip: true,
};
let p = LaserPoint::new(0.5, -0.5, 0, 0, 0, 0);
let (l, r) = OscilloscopeBackend::point_to_samples(&p, &cfg);
assert_eq!(l, 1.0); assert_eq!(r, -1.0); }
#[test]
fn point_to_samples_passes_through_when_clip_disabled() {
let cfg = OscilloscopeConfig {
gain: 4.0,
dc_offset: 0.0,
clip: false,
};
let p = LaserPoint::new(0.5, -0.5, 0, 0, 0, 0);
let (l, r) = OscilloscopeBackend::point_to_samples(&p, &cfg);
assert_eq!(l, 2.0);
assert_eq!(r, -2.0);
}
#[test]
fn point_to_samples_blanked_still_outputs_position() {
let cfg = OscilloscopeConfig::default();
let p = LaserPoint::blanked(0.3, -0.2);
let (l, r) = OscilloscopeBackend::point_to_samples(&p, &cfg);
assert!((l - 0.3).abs() < 1e-6);
assert!((r + 0.2).abs() < 1e-6);
}
#[test]
fn fill_f32_output_unmuted_emits_samples_and_drains() {
let rt = RuntimeState::new(8, RATE);
rt.muted.store(false, Ordering::Release);
rt.queue.push((0.5, -0.5)).unwrap();
rt.queue.push((0.25, 0.75)).unwrap();
let mut data = [9.0f32; 4];
fill_f32_output(&mut data, &rt);
assert_eq!(data, [0.5, -0.5, 0.25, 0.75]);
assert_eq!(rt.queued_points(), 0);
}
#[test]
fn fill_f32_output_muted_emits_zero_but_still_drains() {
let rt = RuntimeState::new(8, RATE); rt.queue.push((0.5, -0.5)).unwrap();
let mut data = [9.0f32; 2];
fill_f32_output(&mut data, &rt);
assert_eq!(data, [0.0, 0.0]);
assert_eq!(rt.queued_points(), 0);
}
#[test]
fn fill_f32_output_underrun_emits_zero() {
let rt = RuntimeState::new(8, RATE);
rt.muted.store(false, Ordering::Release);
let mut data = [9.0f32; 4];
fill_f32_output(&mut data, &rt);
assert!(data.iter().all(|&v| v == 0.0));
}
#[test]
fn fill_i16_output_scales_unmuted() {
let rt = RuntimeState::new(8, RATE);
rt.muted.store(false, Ordering::Release);
rt.queue.push((1.0, -1.0)).unwrap();
let mut data = [123i16; 2];
fill_i16_output(&mut data, &rt);
assert_eq!(data[0], i16::MAX);
assert_eq!(data[1], -i16::MAX);
assert_eq!(rt.queued_points(), 0);
}
#[test]
fn fill_i16_output_muted_and_underrun_emit_zero() {
let rt = RuntimeState::new(8, RATE); rt.queue.push((1.0, -1.0)).unwrap();
let mut data = [123i16; 4];
fill_i16_output(&mut data, &rt);
assert_eq!(data, [0, 0, 0, 0]);
assert_eq!(rt.queued_points(), 0);
}
#[test]
fn runtime_state_capacity_accounting() {
let rt = RuntimeState::new(4, RATE);
assert!(rt.has_capacity_for(0));
assert!(rt.has_capacity_for(4));
assert!(!rt.has_capacity_for(5));
rt.queue.push((0.0, 0.0)).unwrap();
assert_eq!(rt.queued_points(), 1);
assert!(rt.has_capacity_for(3));
assert!(!rt.has_capacity_for(4));
rt.clear_queue();
assert_eq!(rt.queued_points(), 0);
assert!(rt.has_capacity_for(4));
}
#[test]
fn buffer_capacity_scales_with_rate_and_has_floor() {
assert_eq!(super::super::buffer_capacity(48_000), 4_800);
assert_eq!(super::super::buffer_capacity(10_000), 4_096);
}
struct FakeRunningStream {
stop_tx: Option<mpsc::Sender<()>>,
handle: Option<JoinHandle<()>>,
}
impl Drop for FakeRunningStream {
fn drop(&mut self) {
if let Some(tx) = self.stop_tx.take() {
let _ = tx.send(());
}
if let Some(handle) = self.handle.take() {
let _ = handle.join();
}
}
}
impl RunningAudioStream for FakeRunningStream {}
struct FakeAudioEngine {
fail_open: AtomicBool,
paused: Arc<AtomicBool>,
captured_frames: Arc<Mutex<VecDeque<(f32, f32)>>>,
frame_count: Arc<AtomicUsize>,
}
impl FakeAudioEngine {
fn new() -> Self {
Self {
fail_open: AtomicBool::new(false),
paused: Arc::new(AtomicBool::new(false)),
captured_frames: Arc::new(Mutex::new(VecDeque::new())),
frame_count: Arc::new(AtomicUsize::new(0)),
}
}
fn snapshot(&self) -> Vec<(f32, f32)> {
self.captured_frames
.lock()
.unwrap()
.iter()
.cloned()
.collect()
}
fn frame_count(&self) -> usize {
self.frame_count.load(Ordering::Acquire)
}
}
impl AudioEngine for FakeAudioEngine {
fn open_stream(
&self,
_device_name: &str,
_sample_rate: u32,
runtime: Arc<RuntimeState>,
) -> Result<Box<dyn RunningAudioStream>> {
if self.fail_open.load(Ordering::Acquire) {
return Err(Error::backend(std::io::Error::other("mock open failure")));
}
let (stop_tx, stop_rx) = mpsc::channel();
let captured_frames = Arc::clone(&self.captured_frames);
let frame_count = Arc::clone(&self.frame_count);
let paused = Arc::clone(&self.paused);
let handle = thread::spawn(move || loop {
match stop_rx.recv_timeout(Duration::from_millis(2)) {
Ok(_) => break,
Err(RecvTimeoutError::Disconnected) => break,
Err(RecvTimeoutError::Timeout) => {
if paused.load(Ordering::Acquire) {
continue;
}
let mut frame = [0.0f32; 2];
fill_f32_output(&mut frame, &runtime);
if let Ok(mut captured) = captured_frames.lock() {
captured.push_back((frame[0], frame[1]));
if captured.len() > 8192 {
captured.pop_front();
}
}
frame_count.fetch_add(1, Ordering::Release);
}
}
});
Ok(Box::new(FakeRunningStream {
stop_tx: Some(stop_tx),
handle: Some(handle),
}))
}
}
fn wait_for_frame_count(engine: &FakeAudioEngine, target: usize) {
let deadline = Instant::now() + Duration::from_millis(500);
while Instant::now() < deadline {
if engine.frame_count() >= target {
return;
}
thread::sleep(Duration::from_millis(1));
}
panic!(
"timed out waiting for frame count {}, got {}",
target,
engine.frame_count()
);
}
fn backend_with(engine: Arc<FakeAudioEngine>) -> OscilloscopeBackend {
let dyn_engine: Arc<dyn AudioEngine> = engine;
OscilloscopeBackend::with_engine_for_test("scope".to_string(), RATE, dyn_engine)
}
#[test]
fn connect_disconnect_lifecycle() {
let fake = Arc::new(FakeAudioEngine::new());
let mut backend = backend_with(fake.clone());
assert!(!backend.is_connected());
backend.connect().unwrap();
assert!(backend.is_connected());
backend.connect().unwrap();
assert!(backend.is_connected());
backend.disconnect().unwrap();
assert!(!backend.is_connected());
}
#[test]
fn write_before_connect_errors() {
let fake = Arc::new(FakeAudioEngine::new());
let mut backend = backend_with(fake);
let err = backend
.try_write_points(RATE, &[LaserPoint::blanked(0.0, 0.0)])
.unwrap_err();
assert!(err.to_string().contains("Not connected"));
}
#[test]
fn run_audio_thread_propagates_open_failure() {
let fake = Arc::new(FakeAudioEngine::new());
fake.fail_open.store(true, Ordering::Release);
let engine: Arc<dyn AudioEngine> = fake;
let runtime = Arc::new(RuntimeState::new(16, RATE));
let stop = AtomicBool::new(false);
let err = run_audio_thread(&engine, "scope", RATE, &runtime, &stop).unwrap_err();
assert!(err.to_string().contains("mock open failure"));
assert!(!runtime.connected.load(Ordering::Acquire));
}
#[test]
fn end_to_end_write_drains_to_output_when_unmuted() {
let fake = Arc::new(FakeAudioEngine::new());
let mut backend = backend_with(fake.clone());
backend.connect().unwrap();
backend.set_shutter(true).unwrap();
let points = vec![LaserPoint::new(0.5, -0.5, 0, 0, 0, 0); 16];
assert_eq!(
backend.try_write_points(RATE, &points).unwrap(),
WriteOutcome::Written
);
wait_for_frame_count(fake.as_ref(), 16);
let frames = fake.snapshot();
assert!(frames
.iter()
.any(|&(l, r)| (l - 0.5).abs() < 1e-3 && (r + 0.5).abs() < 1e-3));
backend.disconnect().unwrap();
}
#[test]
fn end_to_end_muted_output_is_silent_but_drains() {
let fake = Arc::new(FakeAudioEngine::new());
let mut backend = backend_with(fake.clone());
backend.connect().unwrap(); let rt = backend.runtime.as_ref().unwrap().clone();
let points = vec![LaserPoint::new(0.5, -0.5, 0, 0, 0, 0); 16];
assert_eq!(
backend.try_write_points(RATE, &points).unwrap(),
WriteOutcome::Written
);
let deadline = Instant::now() + Duration::from_millis(500);
while rt.queued_points() != 0 && Instant::now() < deadline {
thread::sleep(Duration::from_millis(1));
}
assert_eq!(rt.queued_points(), 0);
let frames = fake.snapshot();
assert!(frames.iter().all(|&(l, r)| l == 0.0 && r == 0.0));
backend.disconnect().unwrap();
}
#[test]
fn wouldblock_when_full_then_recovers_after_drain() {
let fake = Arc::new(FakeAudioEngine::new());
fake.paused.store(true, Ordering::Release); let mut backend = backend_with(fake.clone());
backend.connect().unwrap();
backend.set_shutter(true).unwrap();
let cap = super::super::buffer_capacity(RATE);
let fill = vec![LaserPoint::blanked(0.0, 0.0); cap];
assert_eq!(
backend.try_write_points(RATE, &fill).unwrap(),
WriteOutcome::Written
);
assert_eq!(
backend
.try_write_points(RATE, &[LaserPoint::blanked(0.0, 0.0)])
.unwrap(),
WriteOutcome::WouldBlock
);
fake.paused.store(false, Ordering::Release);
thread::sleep(Duration::from_millis(40));
assert_eq!(
backend
.try_write_points(RATE, &[LaserPoint::blanked(0.0, 0.0)])
.unwrap(),
WriteOutcome::Written
);
backend.disconnect().unwrap();
}
#[test]
fn stop_and_set_shutter_toggle_mute() {
let fake = Arc::new(FakeAudioEngine::new());
let mut backend = backend_with(fake.clone());
backend.connect().unwrap();
let rt = backend.runtime.as_ref().unwrap().clone();
assert!(rt.muted.load(Ordering::Acquire));
backend.set_shutter(true).unwrap();
assert!(!rt.muted.load(Ordering::Acquire));
backend.stop().unwrap();
assert!(rt.muted.load(Ordering::Acquire));
backend.set_shutter(true).unwrap();
assert!(!rt.muted.load(Ordering::Acquire));
backend.set_shutter(false).unwrap();
assert!(rt.muted.load(Ordering::Acquire));
backend.disconnect().unwrap();
}
#[test]
fn passthrough_when_pps_matches_sample_rate() {
let fake = Arc::new(FakeAudioEngine::new());
fake.paused.store(true, Ordering::Release);
let mut backend = backend_with(fake.clone());
backend.connect().unwrap();
let rt = backend.runtime.as_ref().unwrap().clone();
let points = vec![
LaserPoint::new(-1.0, -1.0, 0, 0, 0, 0),
LaserPoint::new(0.0, 0.0, 0, 0, 0, 0),
LaserPoint::new(1.0, 1.0, 0, 0, 0, 0),
];
assert_eq!(
backend.try_write_points(RATE, &points).unwrap(),
WriteOutcome::Written
);
assert_eq!(rt.queued_points(), 3);
backend.disconnect().unwrap();
}
#[test]
fn resamples_up_to_sample_rate() {
let fake = Arc::new(FakeAudioEngine::new());
fake.paused.store(true, Ordering::Release);
let mut backend = backend_with(fake.clone());
backend.connect().unwrap();
let rt = backend.runtime.as_ref().unwrap().clone();
let points = vec![
LaserPoint::new(-1.0, -1.0, 0, 0, 0, 0),
LaserPoint::new(1.0, 1.0, 0, 0, 0, 0),
];
assert_eq!(
backend.try_write_points(24_000, &points).unwrap(),
WriteOutcome::Written
);
assert_eq!(rt.queued_points(), 3);
let first = rt.queue.pop().unwrap();
assert!((first.0 - (-1.0)).abs() < 0.01);
let _ = rt.queue.pop();
let last = rt.queue.pop().unwrap();
assert!((last.0 - 1.0).abs() < 0.01);
backend.disconnect().unwrap();
}
#[test]
fn resamples_down_to_sample_rate() {
let fake = Arc::new(FakeAudioEngine::new());
fake.paused.store(true, Ordering::Release);
let mut backend = backend_with(fake.clone());
backend.connect().unwrap();
let rt = backend.runtime.as_ref().unwrap().clone();
let points = vec![
LaserPoint::new(-1.0, -1.0, 0, 0, 0, 0),
LaserPoint::new(-0.5, -0.5, 0, 0, 0, 0),
LaserPoint::new(0.5, 0.5, 0, 0, 0, 0),
LaserPoint::new(1.0, 1.0, 0, 0, 0, 0),
];
assert_eq!(
backend.try_write_points(96_000, &points).unwrap(),
WriteOutcome::Written
);
assert_eq!(rt.queued_points(), 2);
backend.disconnect().unwrap();
}
#[test]
fn wouldblock_on_resampled_overflow() {
let fake = Arc::new(FakeAudioEngine::new());
fake.paused.store(true, Ordering::Release);
let mut backend = backend_with(fake.clone());
backend.connect().unwrap();
let rt = backend.runtime.as_ref().unwrap().clone();
let cap = super::super::buffer_capacity(RATE);
let fill = vec![LaserPoint::blanked(0.0, 0.0); cap - 1];
assert_eq!(
backend.try_write_points(RATE, &fill).unwrap(),
WriteOutcome::Written
);
let points = vec![
LaserPoint::new(0.0, 0.0, 0, 0, 0, 0),
LaserPoint::new(1.0, 1.0, 0, 0, 0, 0),
];
assert_eq!(
backend.try_write_points(24_000, &points).unwrap(),
WriteOutcome::WouldBlock
);
assert_eq!(rt.queued_points(), cap as u64 - 1);
backend.disconnect().unwrap();
}
#[test]
fn estimator_reports_queue_depth_and_clears_on_disconnect() {
let fake = Arc::new(FakeAudioEngine::new());
fake.paused.store(true, Ordering::Release);
let mut backend = backend_with(fake.clone());
let now = Instant::now();
assert_eq!(backend.estimator().estimated_fullness(now, RATE), 0);
backend.connect().unwrap();
let points = vec![LaserPoint::blanked(0.0, 0.0); 10];
backend.try_write_points(RATE, &points).unwrap();
assert_eq!(backend.estimator().estimated_fullness(now, RATE), 10);
backend.disconnect().unwrap();
assert_eq!(backend.estimator().estimated_fullness(now, RATE), 0);
}
#[test]
fn disconnect_clears_queue() {
let fake = Arc::new(FakeAudioEngine::new());
fake.paused.store(true, Ordering::Release);
let mut backend = backend_with(fake.clone());
backend.connect().unwrap();
let rt = backend.runtime.as_ref().unwrap().clone();
backend
.try_write_points(RATE, &vec![LaserPoint::blanked(0.0, 0.0); 20])
.unwrap();
assert_eq!(rt.queued_points(), 20);
backend.disconnect().unwrap();
assert_eq!(rt.queued_points(), 0);
}
}