nandi/jolt-nativepublic Fork 0
cfd3e3677bed92e80cfe447469a249cfdc0b1401
Commits
Clone
git clone https://git.rickub.com/nandi/jolt-native.git
git clone ssh://git@rickub.com/nandi/jolt-native.git

Host key fingerprint (ed25519): SHA256:iycHnxEyq0Q7uyVpB7JlznP0G7JrTPXLYRcAU5CSLhc — verify it before your first connect.

stream.rs · 1176 lines · 44.8 KBRust Blame HistoryRaw
Lift freeq's AV media plane out of sleek 90f8b89 nandi 19d ago1extern crate asio_sys as sys;
2extern crate num_traits;
3
4use crate::host::com;
5use crate::I24;
6
7use self::num_traits::{FromPrimitive, PrimInt};
8use super::Device;
9use crate::{
10 BackendSpecificError, BufferSize, BuildStreamError, Data, InputCallbackInfo,
11 OutputCallbackInfo, PauseStreamError, PlayStreamError, SampleFormat, StreamConfig, StreamError,
12};
13use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
14use std::sync::{Arc, Mutex};
15use std::time::Duration;
16
17pub struct Stream {
18 playing: Arc<AtomicBool>,
19 // Ensure the `Driver` does not terminate until the last stream is dropped.
20 driver: Arc<sys::Driver>,
21 #[allow(dead_code)]
22 asio_streams: Arc<Mutex<sys::AsioStreams>>,
23 callback_id: sys::BufferCallbackId,
24 driver_event_callback_id: sys::DriverEventCallbackId,
25}
26
27// Compile-time assertion that Stream is Send and Sync
28crate::assert_stream_send!(Stream);
29crate::assert_stream_sync!(Stream);
30
31impl Stream {
32 pub fn play(&self) -> Result<(), PlayStreamError> {
33 self.playing.store(true, Ordering::Release);
34 Ok(())
35 }
36
37 pub fn pause(&self) -> Result<(), PauseStreamError> {
38 self.playing.store(false, Ordering::Release);
39 Ok(())
40 }
41
42 pub fn buffer_size(&self) -> Option<crate::FrameCount> {
43 let streams = self.asio_streams.lock().ok()?;
44 streams
45 .output
46 .as_ref()
47 .or(streams.input.as_ref())
48 .map(|s| s.buffer_size as crate::FrameCount)
49 }
50}
51
52impl Device {
53 pub fn build_input_stream_raw<D, E>(
54 &self,
55 config: StreamConfig,
56 sample_format: SampleFormat,
57 mut data_callback: D,
58 error_callback: E,
59 _timeout: Option<Duration>,
60 ) -> Result<Stream, BuildStreamError>
61 where
62 D: FnMut(&Data, &InputCallbackInfo) + Send + 'static,
63 E: FnMut(StreamError) + Send + 'static,
64 {
65 com::com_initialized();
66 let description = self
67 .description()
68 .map_err(|_| BuildStreamError::DeviceNotAvailable)?;
69 let driver = super::GLOBAL_ASIO
70 .get()
71 .ok_or(BuildStreamError::DeviceNotAvailable)?
72 .load_driver(description.name())
73 .map_err(load_driver_err)?;
74
75 let stream_type = driver.input_data_type().map_err(build_stream_err)?;
76
77 // Ensure that the desired sample type is supported.
78 let expected_sample_format = super::device::convert_data_type(&stream_type)
79 .ok_or(BuildStreamError::StreamConfigNotSupported)?;
80 if sample_format != expected_sample_format {
81 return Err(BuildStreamError::StreamConfigNotSupported);
82 }
83
84 let num_channels = config.channels;
85 let buffer_size = self.get_or_create_input_stream(&driver, config, sample_format)?;
86 let cpal_num_samples = buffer_size * num_channels as usize;
87
88 // Create the buffer depending on the size of the data type.
89 let len_bytes = cpal_num_samples * sample_format.sample_size();
90 let mut interleaved = vec![0u8; len_bytes];
91
92 // Query hardware input latency (order matters: needs buffers created above).
93 // Wrapped in Arc<AtomicUsize> so the message callback can update it on
94 // kAsioLatenciesChanged without touching the buffer callback.
95 let hardware_input_latency = Arc::new(AtomicUsize::new(
96 driver
97 .latencies()
98 .map(|latencies| latencies.input.max(0) as usize)
99 .unwrap_or(0),
100 ));
101
102 let driver_event_callback_id = self.add_event_callback(
103 &driver,
104 error_callback,
105 Arc::clone(&hardware_input_latency),
106 true,
107 );
108
109 let stream_playing = Arc::new(AtomicBool::new(false));
110 let playing = Arc::clone(&stream_playing);
111 let asio_streams = self.asio_streams.clone();
112 let mut current_buffer_size = buffer_size as i32;
113 let mut last_buffer_index: i32 = -1;
114
115 // Set the input callback.
116 // This is most performance critical part of the ASIO bindings.
117 let callback_id = driver.add_callback(move |callback_info| unsafe {
118 // If not playing return early.
119 if !playing.load(Ordering::Acquire) {
120 return;
121 }
122
123 // Guard against non-conformant drivers (e.g. Focusrite USB ASIO, ReaRoute) that
124 // fire the buffer callback multiple times per buffer cycle with the same buffer
125 // index.
126 if callback_info.buffer_index == last_buffer_index {
127 return;
128 }
129 last_buffer_index = callback_info.buffer_index;
130
131 // There is 0% chance of lock contention the host only locks when recreating streams.
132 let stream_lock = asio_streams.lock().unwrap();
133 let asio_stream = match stream_lock.input {
134 Some(ref asio_stream) => asio_stream,
135 None => return,
136 };
137
138 // Resize the buffer only when the driver issues a buffer size change request.
139 // In normal operation this branch is never taken.
140 if asio_stream.buffer_size != current_buffer_size {
141 current_buffer_size = asio_stream.buffer_size;
142 interleaved.resize(
143 current_buffer_size as usize
144 * num_channels as usize
145 * sample_format.sample_size(),
146 0,
147 );
148 }
149
150 let hardware_input_latency = hardware_input_latency.load(Ordering::Relaxed);
151
152 /// 1. Write from the ASIO buffer to the interleaved CPAL buffer.
153 /// 2. Deliver the CPAL buffer to the user callback.
154 #[allow(clippy::too_many_arguments)]
155 unsafe fn process_input_callback<A, D, F>(
156 data_callback: &mut D,
157 interleaved: &mut [u8],
158 asio_stream: &sys::AsioStream,
159 asio_info: &sys::CallbackInfo,
160 sample_rate: crate::SampleRate,
161 format: SampleFormat,
162 from_endianness: F,
163 hardware_latency_frames: usize,
164 ) where
165 A: Copy,
166 D: FnMut(&Data, &InputCallbackInfo) + Send + 'static,
167 F: Fn(A) -> A,
168 {
169 // 1. Write the ASIO channels to the CPAL buffer.
170 let interleaved: &mut [A] = cast_slice_mut(interleaved);
171 let n_frames = asio_stream.buffer_size as usize;
172 let n_channels = interleaved.len() / n_frames;
173 let buffer_index = asio_info.buffer_index as usize;
174 for ch_ix in 0..n_channels {
175 let asio_channel =
176 asio_channel_slice::<A>(asio_stream, buffer_index, ch_ix, None);
177 for (frame, s_asio) in interleaved.chunks_mut(n_channels).zip(asio_channel) {
178 frame[ch_ix] = from_endianness(*s_asio);
179 }
180 }
181
182 // 2. Deliver the interleaved buffer to the callback.
183 apply_input_callback_to_data::<A, _>(
184 data_callback,
185 interleaved,
186 asio_info,
187 sample_rate,
188 format,
189 hardware_latency_frames,
190 );
191 }
192
193 match (&stream_type, sample_format) {
194 (&sys::AsioSampleType::ASIOSTInt16LSB, SampleFormat::I16) => {
195 process_input_callback::<i16, _, _>(
196 &mut data_callback,
197 &mut interleaved,
198 asio_stream,
199 callback_info,
200 config.sample_rate,
201 SampleFormat::I16,
202 from_le,
203 hardware_input_latency,
204 );
205 }
206 (&sys::AsioSampleType::ASIOSTInt16MSB, SampleFormat::I16) => {
207 process_input_callback::<i16, _, _>(
208 &mut data_callback,
209 &mut interleaved,
210 asio_stream,
211 callback_info,
212 config.sample_rate,
213 SampleFormat::I16,
214 from_be,
215 hardware_input_latency,
216 );
217 }
218
219 (&sys::AsioSampleType::ASIOSTFloat32LSB, SampleFormat::F32) => {
220 process_input_callback::<u32, _, _>(
221 &mut data_callback,
222 &mut interleaved,
223 asio_stream,
224 callback_info,
225 config.sample_rate,
226 SampleFormat::F32,
227 from_le,
228 hardware_input_latency,
229 );
230 }
231 (&sys::AsioSampleType::ASIOSTFloat32MSB, SampleFormat::F32) => {
232 process_input_callback::<u32, _, _>(
233 &mut data_callback,
234 &mut interleaved,
235 asio_stream,
236 callback_info,
237 config.sample_rate,
238 SampleFormat::F32,
239 from_be,
240 hardware_input_latency,
241 );
242 }
243
244 (&sys::AsioSampleType::ASIOSTInt32LSB, SampleFormat::I32) => {
245 process_input_callback::<i32, _, _>(
246 &mut data_callback,
247 &mut interleaved,
248 asio_stream,
249 callback_info,
250 config.sample_rate,
251 SampleFormat::I32,
252 from_le,
253 hardware_input_latency,
254 );
255 }
256 (&sys::AsioSampleType::ASIOSTInt32MSB, SampleFormat::I32) => {
257 process_input_callback::<i32, _, _>(
258 &mut data_callback,
259 &mut interleaved,
260 asio_stream,
261 callback_info,
262 config.sample_rate,
263 SampleFormat::I32,
264 from_be,
265 hardware_input_latency,
266 );
267 }
268
269 (&sys::AsioSampleType::ASIOSTFloat64LSB, SampleFormat::F64) => {
270 process_input_callback::<u64, _, _>(
271 &mut data_callback,
272 &mut interleaved,
273 asio_stream,
274 callback_info,
275 config.sample_rate,
276 SampleFormat::F64,
277 from_le,
278 hardware_input_latency,
279 );
280 }
281 (&sys::AsioSampleType::ASIOSTFloat64MSB, SampleFormat::F64) => {
282 process_input_callback::<u64, _, _>(
283 &mut data_callback,
284 &mut interleaved,
285 asio_stream,
286 callback_info,
287 config.sample_rate,
288 SampleFormat::F64,
289 from_be,
290 hardware_input_latency,
291 );
292 }
293
294 (&sys::AsioSampleType::ASIOSTInt24LSB, SampleFormat::I24) => {
295 process_input_callback_i24(
296 &mut data_callback,
297 &mut interleaved,
298 asio_stream,
299 callback_info,
300 config.sample_rate,
301 true,
302 hardware_input_latency,
303 );
304 }
305 (&sys::AsioSampleType::ASIOSTInt24MSB, SampleFormat::I24) => {
306 process_input_callback_i24(
307 &mut data_callback,
308 &mut interleaved,
309 asio_stream,
310 callback_info,
311 config.sample_rate,
312 false,
313 hardware_input_latency,
314 );
315 }
316
317 unsupported_format_pair => unreachable!(
318 "`build_input_stream_raw` should have returned with unsupported \
319 format {:?}",
320 unsupported_format_pair
321 ),
322 }
323 });
324
325 let driver = Arc::new(driver);
326 let asio_streams = self.asio_streams.clone();
327
328 driver.start().map_err(build_stream_err)?;
329
330 Ok(Stream {
331 playing: stream_playing,
332 driver,
333 asio_streams,
334 callback_id,
335 driver_event_callback_id,
336 })
337 }
338
339 pub fn build_output_stream_raw<D, E>(
340 &self,
341 config: StreamConfig,
342 sample_format: SampleFormat,
343 mut data_callback: D,
344 error_callback: E,
345 _timeout: Option<Duration>,
346 ) -> Result<Stream, BuildStreamError>
347 where
348 D: FnMut(&mut Data, &OutputCallbackInfo) + Send + 'static,
349 E: FnMut(StreamError) + Send + 'static,
350 {
351 com::com_initialized();
352 let description = self
353 .description()
354 .map_err(|_| BuildStreamError::DeviceNotAvailable)?;
355 let driver = super::GLOBAL_ASIO
356 .get()
357 .ok_or(BuildStreamError::DeviceNotAvailable)?
358 .load_driver(description.name())
359 .map_err(load_driver_err)?;
360
361 let stream_type = driver.output_data_type().map_err(build_stream_err)?;
362
363 // Ensure that the desired sample type is supported.
364 let expected_sample_format = super::device::convert_data_type(&stream_type)
365 .ok_or(BuildStreamError::StreamConfigNotSupported)?;
366 if sample_format != expected_sample_format {
367 return Err(BuildStreamError::StreamConfigNotSupported);
368 }
369
370 let num_channels = config.channels;
371 let buffer_size = self.get_or_create_output_stream(&driver, config, sample_format)?;
372 let cpal_num_samples = buffer_size * num_channels as usize;
373
374 // Create the buffer depending on data type.
375 let len_bytes = cpal_num_samples * sample_format.sample_size();
376 let mut interleaved = vec![0u8; len_bytes];
377 let current_callback_flag = self.current_callback_flag.clone();
378
379 // Query hardware output latency (order matters: needs buffers created above).
380 // Wrapped in Arc<AtomicUsize> so the message callback can update it on
381 // kAsioLatenciesChanged without touching the buffer callback.
382 let hardware_output_latency = Arc::new(AtomicUsize::new(
383 driver
384 .latencies()
385 .map(|latencies| latencies.output.max(0) as usize)
386 .unwrap_or(0),
387 ));
388
389 let driver_event_callback_id = self.add_event_callback(
390 &driver,
391 error_callback,
392 Arc::clone(&hardware_output_latency),
393 false,
394 );
395
396 let stream_playing = Arc::new(AtomicBool::new(false));
397 let playing = Arc::clone(&stream_playing);
398 let asio_streams = self.asio_streams.clone();
399 let mut current_buffer_size = buffer_size as i32;
400 let mut last_buffer_index: i32 = -1;
401
402 let callback_id = driver.add_callback(move |callback_info| unsafe {
403 // If not playing, return early.
404 if !playing.load(Ordering::Acquire) {
405 return;
406 }
407
408 // Guard against non-conformant drivers (e.g. Focusrite USB ASIO, ReaRoute) that
409 // fire the buffer callback multiple times per buffer cycle with the same buffer
410 // index.
411 if callback_info.buffer_index == last_buffer_index {
412 return;
413 }
414 last_buffer_index = callback_info.buffer_index;
415
416 // There is 0% chance of lock contention the host only locks when recreating streams.
417 let mut stream_lock = asio_streams.lock().unwrap();
418 let asio_stream = match stream_lock.output {
419 Some(ref mut asio_stream) => asio_stream,
420 None => return,
421 };
422
423 // Resize the buffer only when the driver issues a buffer size change request.
424 // In normal operation this branch is never taken.
425 if asio_stream.buffer_size != current_buffer_size {
426 current_buffer_size = asio_stream.buffer_size;
427 interleaved.resize(
428 current_buffer_size as usize
429 * num_channels as usize
430 * sample_format.sample_size(),
431 0,
432 );
433 }
434
435 let hardware_output_latency = hardware_output_latency.load(Ordering::Relaxed);
436
437 // Silence the ASIO buffer that is about to be used.
438 //
439 // Check if any other callbacks have already silenced the buffer associated with
440 // the current callback. The flag is updated once per buffer switch.
441 let silence =
442 current_callback_flag.load(Ordering::Acquire) != callback_info.callback_flag;
443
444 if silence {
445 current_callback_flag.store(callback_info.callback_flag, Ordering::Release);
446 }
447
448 /// 1. Render the given callback to the given buffer of interleaved samples.
449 /// 2. If required, silence the ASIO buffer.
450 /// 3. Finally, write the interleaved data to the non-interleaved ASIO buffer,
451 /// performing endianness conversions as necessary.
452 #[allow(clippy::too_many_arguments)]
453 unsafe fn process_output_callback<A, D, F>(
454 data_callback: &mut D,
455 interleaved: &mut [u8],
456 silence_asio_buffer: bool,
457 asio_stream: &mut sys::AsioStream,
458 asio_info: &sys::CallbackInfo,
459 sample_rate: crate::SampleRate,
460 format: SampleFormat,
461 mix_samples: F,
462 hardware_latency_frames: usize,
463 ) where
464 A: Copy,
465 D: FnMut(&mut Data, &OutputCallbackInfo) + Send + 'static,
466 F: Fn(A, A) -> A,
467 {
468 let interleaved: &mut [A] = cast_slice_mut(interleaved);
469 apply_output_callback_to_data::<A, _>(
470 data_callback,
471 interleaved,
472 asio_info,
473 sample_rate,
474 format,
475 hardware_latency_frames,
476 );
477 let n_channels = interleaved.len() / asio_stream.buffer_size as usize;
478 let buffer_index = asio_info.buffer_index as usize;
479
480 // Write interleaved samples to ASIO channels, one channel at a time.
481 for ch_ix in 0..n_channels {
482 let asio_channel =
483 asio_channel_slice_mut::<A>(asio_stream, buffer_index, ch_ix, None);
484 if silence_asio_buffer {
485 asio_channel.align_to_mut::<u8>().1.fill(0);
486 }
487 for (frame, s_asio) in interleaved.chunks(n_channels).zip(asio_channel) {
488 *s_asio = mix_samples(*s_asio, frame[ch_ix]);
489 }
490 }
491 }
492
493 match (sample_format, &stream_type) {
494 (SampleFormat::I16, &sys::AsioSampleType::ASIOSTInt16LSB) => {
495 process_output_callback::<i16, _, _>(
496 &mut data_callback,
497 &mut interleaved,
498 silence,
499 asio_stream,
500 callback_info,
501 config.sample_rate,
502 SampleFormat::I16,
503 |old_sample, new_sample| {
504 from_le(old_sample).saturating_add(new_sample).to_le()
505 },
506 hardware_output_latency,
507 );
508 }
509 (SampleFormat::I16, &sys::AsioSampleType::ASIOSTInt16MSB) => {
510 process_output_callback::<i16, _, _>(
511 &mut data_callback,
512 &mut interleaved,
513 silence,
514 asio_stream,
515 callback_info,
516 config.sample_rate,
517 SampleFormat::I16,
518 |old_sample, new_sample| {
519 from_be(old_sample).saturating_add(new_sample).to_be()
520 },
521 hardware_output_latency,
522 );
523 }
524 (SampleFormat::F32, &sys::AsioSampleType::ASIOSTFloat32LSB) => {
525 process_output_callback::<u32, _, _>(
526 &mut data_callback,
527 &mut interleaved,
528 silence,
529 asio_stream,
530 callback_info,
531 config.sample_rate,
532 SampleFormat::F32,
533 |old_sample, new_sample| {
534 (f32::from_bits(from_le(old_sample)) + f32::from_bits(new_sample))
535 .to_bits()
536 .to_le()
537 },
538 hardware_output_latency,
539 );
540 }
541
542 (SampleFormat::F32, &sys::AsioSampleType::ASIOSTFloat32MSB) => {
543 process_output_callback::<u32, _, _>(
544 &mut data_callback,
545 &mut interleaved,
546 silence,
547 asio_stream,
548 callback_info,
549 config.sample_rate,
550 SampleFormat::F32,
551 |old_sample, new_sample| {
552 (f32::from_bits(from_be(old_sample)) + f32::from_bits(new_sample))
553 .to_bits()
554 .to_be()
555 },
556 hardware_output_latency,
557 );
558 }
559
560 (SampleFormat::I32, &sys::AsioSampleType::ASIOSTInt32LSB) => {
561 process_output_callback::<i32, _, _>(
562 &mut data_callback,
563 &mut interleaved,
564 silence,
565 asio_stream,
566 callback_info,
567 config.sample_rate,
568 SampleFormat::I32,
569 |old_sample, new_sample| {
570 from_le(old_sample).saturating_add(new_sample).to_le()
571 },
572 hardware_output_latency,
573 );
574 }
575 (SampleFormat::I32, &sys::AsioSampleType::ASIOSTInt32MSB) => {
576 process_output_callback::<i32, _, _>(
577 &mut data_callback,
578 &mut interleaved,
579 silence,
580 asio_stream,
581 callback_info,
582 config.sample_rate,
583 SampleFormat::I32,
584 |old_sample, new_sample| {
585 from_be(old_sample).saturating_add(new_sample).to_be()
586 },
587 hardware_output_latency,
588 );
589 }
590
591 (SampleFormat::F64, &sys::AsioSampleType::ASIOSTFloat64LSB) => {
592 process_output_callback::<u64, _, _>(
593 &mut data_callback,
594 &mut interleaved,
595 silence,
596 asio_stream,
597 callback_info,
598 config.sample_rate,
599 SampleFormat::F64,
600 |old_sample, new_sample| {
601 (f64::from_bits(from_le(old_sample)) + f64::from_bits(new_sample))
602 .to_bits()
603 .to_le()
604 },
605 hardware_output_latency,
606 );
607 }
608
609 (SampleFormat::F64, &sys::AsioSampleType::ASIOSTFloat64MSB) => {
610 process_output_callback::<u64, _, _>(
611 &mut data_callback,
612 &mut interleaved,
613 silence,
614 asio_stream,
615 callback_info,
616 config.sample_rate,
617 SampleFormat::F64,
618 |old_sample, new_sample| {
619 (f64::from_bits(from_be(old_sample)) + f64::from_bits(new_sample))
620 .to_bits()
621 .to_be()
622 },
623 hardware_output_latency,
624 );
625 }
626
627 (SampleFormat::I24, &sys::AsioSampleType::ASIOSTInt24LSB) => {
628 process_output_callback_i24::<_>(
629 &mut data_callback,
630 &mut interleaved,
631 silence,
632 true,
633 asio_stream,
634 callback_info,
635 config.sample_rate,
636 hardware_output_latency,
637 );
638 }
639
640 (SampleFormat::I24, &sys::AsioSampleType::ASIOSTInt24MSB) => {
641 process_output_callback_i24::<_>(
642 &mut data_callback,
643 &mut interleaved,
644 silence,
645 false,
646 asio_stream,
647 callback_info,
648 config.sample_rate,
649 hardware_output_latency,
650 );
651 }
652
653 unsupported_format_pair => unreachable!(
654 "`build_output_stream_raw` should have returned with unsupported \
655 format {:?}",
656 unsupported_format_pair
657 ),
658 }
659 });
660
661 let driver = Arc::new(driver);
662 let asio_streams = self.asio_streams.clone();
663
664 driver.start().map_err(build_stream_err)?;
665
666 Ok(Stream {
667 playing: stream_playing,
668 driver,
669 asio_streams,
670 callback_id,
671 driver_event_callback_id,
672 })
673 }
674
675 /// Create a new CPAL Input Stream.
676 ///
677 /// If there is no existing ASIO Input Stream it will be created.
678 ///
679 /// On success, the buffer size of the stream is returned.
680 fn get_or_create_input_stream(
681 &self,
682 driver: &sys::Driver,
683 config: StreamConfig,
684 sample_format: SampleFormat,
685 ) -> Result<usize, BuildStreamError> {
686 let num_asio_channels = self
687 .default_input_config()
688 .map_err(|_| BuildStreamError::StreamConfigNotSupported)?
689 .channels;
690 check_config(driver, config, sample_format, num_asio_channels)?;
691 let num_channels = config.channels as usize;
692 let mut streams = self.asio_streams.lock().unwrap();
693
694 let buffer_size = match config.buffer_size {
695 BufferSize::Fixed(v) => Some(v as i32),
696 BufferSize::Default => None,
697 };
698
699 // Either create a stream if thers none or had back the
700 // size of the current one.
701 match streams.input {
702 Some(ref input) => Ok(input.buffer_size as usize),
703 None => {
704 let output = streams.output.take();
705 driver
706 .prepare_input_stream(output, num_channels, buffer_size)
707 .map(|new_streams| {
708 let bs = match new_streams.input {
709 Some(ref inp) => inp.buffer_size as usize,
710 None => unreachable!(),
711 };
712 *streams = new_streams;
713 bs
714 })
715 .map_err(|_| BuildStreamError::DeviceNotAvailable)
716 }
717 }
718 }
719
720 /// Create a new CPAL Output Stream.
721 ///
722 /// If there is no existing ASIO Output Stream it will be created.
723 fn get_or_create_output_stream(
724 &self,
725 driver: &sys::Driver,
726 config: StreamConfig,
727 sample_format: SampleFormat,
728 ) -> Result<usize, BuildStreamError> {
729 let num_asio_channels = self
730 .default_output_config()
731 .map_err(|_| BuildStreamError::StreamConfigNotSupported)?
732 .channels;
733 check_config(driver, config, sample_format, num_asio_channels)?;
734 let num_channels = config.channels as usize;
735 let mut streams = self.asio_streams.lock().unwrap();
736
737 let buffer_size = match config.buffer_size {
738 BufferSize::Fixed(v) => Some(v as i32),
739 BufferSize::Default => None,
740 };
741
742 // Either create a stream if thers none or had back the
743 // size of the current one.
744 match streams.output {
745 Some(ref output) => Ok(output.buffer_size as usize),
746 None => {
747 let input = streams.input.take();
748 driver
749 .prepare_output_stream(input, num_channels, buffer_size)
750 .map(|new_streams| {
751 let bs = match new_streams.output {
752 Some(ref out) => out.buffer_size as usize,
753 None => unreachable!(),
754 };
755 *streams = new_streams;
756 bs
757 })
758 .map_err(|_| BuildStreamError::DeviceNotAvailable)
759 }
760 }
761 }
762
763 fn add_event_callback<E>(
764 &self,
765 driver: &sys::Driver,
766 error_callback: E,
767 hardware_latency: Arc<AtomicUsize>,
768 is_input: bool,
769 ) -> sys::DriverEventCallbackId
770 where
771 E: FnMut(StreamError) + Send + 'static,
772 {
773 let error_callback_shared = Arc::new(Mutex::new(error_callback));
774 let configured_sample_rate = driver.sample_rate().ok().filter(|&r| r > 0.0);
775 let driver_for_latency = driver.clone();
776 let asio_streams_for_event = self.asio_streams.clone();
777
778 driver.add_event_callback(move |event| {
779 match event {
780 sys::AsioDriverEvent::Message {
781 selector: msg,
782 value,
783 } => match msg {
784 sys::AsioMessageSelectors::kAsioSelectorSupported => {
785 // Signal which selectors this stream opts into.
786 matches!(
787 sys::AsioMessageSelectors::from_i64(value as i64),
788 Some(sys::AsioMessageSelectors::kAsioBufferSizeChange)
789 )
790 }
791 sys::AsioMessageSelectors::kAsioResetRequest => {
792 if let Ok(mut cb) = error_callback_shared.lock() {
793 cb(StreamError::StreamInvalidated);
794 }
795 false
796 }
797 sys::AsioMessageSelectors::kAsioResyncRequest => {
798 if let Ok(mut cb) = error_callback_shared.lock() {
799 cb(StreamError::BufferUnderrun);
800 }
801 false
802 }
803 sys::AsioMessageSelectors::kAsioLatenciesChanged => {
804 if let Ok(latencies) = driver_for_latency.latencies() {
805 let latency = if is_input {
806 latencies.input
807 } else {
808 latencies.output
809 };
810 hardware_latency.store(latency.max(0) as usize, Ordering::Relaxed);
811 }
812 false
813 }
814 sys::AsioMessageSelectors::kAsioBufferSizeChange => {
815 if value > 0 {
816 if let Ok(mut streams) = asio_streams_for_event.lock() {
817 let stream = if is_input {
818 streams.input.as_mut()
819 } else {
820 streams.output.as_mut()
821 };
822 if let Some(s) = stream {
823 s.buffer_size = value;
824 }
825 }
826 }
827 true
828 }
829 _ => false,
830 },
831 sys::AsioDriverEvent::SampleRateChanged(new_rate) => {
832 if let Some(rate) = configured_sample_rate {
833 if (new_rate - rate).abs() >= 1.0 {
834 if let Ok(mut cb) = error_callback_shared.lock() {
835 cb(StreamError::StreamInvalidated);
836 }
837 }
838 }
839 false
840 }
841 }
842 })
843 }
844}
845
846impl Drop for Stream {
847 fn drop(&mut self) {
848 self.driver.remove_callback(self.callback_id);
849 self.driver
850 .remove_event_callback(self.driver_event_callback_id);
851 }
852}
853
854// Convert the given duration in frames at the given sample rate to a `std::time::Duration`.
855#[inline]
856fn frames_to_duration(frames: usize, rate: crate::SampleRate) -> std::time::Duration {
857 let secsf = frames as f64 / rate as f64;
858 let secs = secsf as u64;
859 let nanos = ((secsf - secs as f64) * 1_000_000_000.0) as u32;
860 std::time::Duration::new(secs, nanos)
861}
862
863/// Check whether or not the desired config is supported by the stream.
864///
865/// Checks sample rate, data type, number of channels, and buffer size.
866fn check_config(
867 driver: &sys::Driver,
868 config: StreamConfig,
869 sample_format: SampleFormat,
870 num_asio_channels: u16,
871) -> Result<(), BuildStreamError> {
872 let StreamConfig {
873 channels,
874 sample_rate,
875 buffer_size,
876 } = config;
877
878 // Validate buffer size if `Fixed` is specified. This is necessary because ASIO's
879 // `create_buffers` only validates the upper bound (returns `InvalidBufferSize` if > max) but
880 // does NOT validate the lower bound. Passing a buffer size below min would be accepted but
881 // behavior is unspecified.
882 if let BufferSize::Fixed(requested_size) = buffer_size {
883 let range = driver.buffersize_range().map_err(build_stream_err)?;
884 let requested_size_i32 = requested_size as i32;
885 if !(range.min..=range.max).contains(&requested_size_i32) {
886 return Err(BuildStreamError::StreamConfigNotSupported);
887 }
888 }
889
890 // Try and set the sample rate to what the user selected.
891 let sample_rate = sample_rate.into();
892 if sample_rate != driver.sample_rate().map_err(build_stream_err)? {
893 if driver
894 .can_sample_rate(sample_rate)
895 .map_err(build_stream_err)?
896 {
897 driver
898 .set_sample_rate(sample_rate)
899 .map_err(build_stream_err)?;
900 } else {
901 return Err(BuildStreamError::StreamConfigNotSupported);
902 }
903 }
904 // unsigned formats are not supported by asio
905 match sample_format {
906 SampleFormat::I16 | SampleFormat::I24 | SampleFormat::I32 | SampleFormat::F32 => (),
907 _ => return Err(BuildStreamError::StreamConfigNotSupported),
908 }
909 if channels > num_asio_channels {
910 return Err(BuildStreamError::StreamConfigNotSupported);
911 }
912 Ok(())
913}
914
915/// Cast a byte slice into a mutable slice of desired type.
916///
917/// Safety: it's up to the caller to ensure that the input slice has valid bit representations.
918unsafe fn cast_slice_mut<T>(v: &mut [u8]) -> &mut [T] {
919 debug_assert!(v.len() % std::mem::size_of::<T>() == 0);
920 std::slice::from_raw_parts_mut(v.as_mut_ptr() as *mut T, v.len() / std::mem::size_of::<T>())
921}
922
923/// Helper function to convert from little endianness.
924fn from_le<T: PrimInt>(t: T) -> T {
925 T::from_le(t)
926}
927
928/// Helper function to convert from little endianness.
929fn from_be<T: PrimInt>(t: T) -> T {
930 T::from_be(t)
931}
932
933/// Shorthand for retrieving the asio buffer slice associated with a channel.
934///
935/// The channel length is automatically inferred from the buffer size or some
936/// value can be passed to enforce a certain length (for odd sized sample formats)
937unsafe fn asio_channel_slice<T>(
938 asio_stream: &sys::AsioStream,
939 buffer_index: usize,
940 channel_index: usize,
941 requested_channel_length: Option<usize>,
942) -> &[T] {
943 let channel_length = requested_channel_length.unwrap_or(asio_stream.buffer_size as usize);
944 let buff_ptr: *const T =
945 asio_stream.buffer_infos[channel_index].buffers[buffer_index] as *const _;
946 std::slice::from_raw_parts(buff_ptr, channel_length)
947}
948
949/// Shorthand for retrieving the asio buffer slice associated with a channel.
950///
951/// The channel length is automatically inferred from the buffer size or some
952/// value can be passed to enforce a certain length (for odd sized sample formats)
953unsafe fn asio_channel_slice_mut<T>(
954 asio_stream: &mut sys::AsioStream,
955 buffer_index: usize,
956 channel_index: usize,
957 requested_channel_length: Option<usize>,
958) -> &mut [T] {
959 let channel_length = requested_channel_length.unwrap_or(asio_stream.buffer_size as usize);
960 let buff_ptr: *mut T = asio_stream.buffer_infos[channel_index].buffers[buffer_index] as *mut _;
961 std::slice::from_raw_parts_mut(buff_ptr, channel_length)
962}
963
964fn load_driver_err(e: sys::LoadDriverError) -> BuildStreamError {
965 match e {
966 sys::LoadDriverError::LoadDriverFailed | sys::LoadDriverError::DriverAlreadyExists => {
967 BuildStreamError::DeviceNotAvailable
968 }
969 sys::LoadDriverError::InitializationFailed(asio_err) => build_stream_err(asio_err),
970 }
971}
972
973fn build_stream_err(e: sys::AsioError) -> BuildStreamError {
974 match e {
975 sys::AsioError::NoDrivers | sys::AsioError::HardwareMalfunction => {
976 BuildStreamError::DeviceNotAvailable
977 }
978 sys::AsioError::InvalidInput | sys::AsioError::BadMode => BuildStreamError::InvalidArgument,
979 err => {
980 let description = format!("{}", err);
981 BackendSpecificError { description }.into()
982 }
983 }
984}
985
986/// Convert i24 bytes to i32
987fn i24_bytes_to_i32(i24_bytes: &[u8; 3], little_endian: bool) -> i32 {
988 let sample = if little_endian {
989 i32::from_le_bytes([i24_bytes[0], i24_bytes[1], i24_bytes[2], 0u8])
990 } else {
991 i32::from_le_bytes([i24_bytes[2], i24_bytes[1], i24_bytes[0], 0u8])
992 };
993 if sample & 0x800000 != 0 {
994 sample | -0x1000000
995 } else {
996 sample
997 }
998}
999
1000#[allow(clippy::too_many_arguments)]
1001unsafe fn process_output_callback_i24<D>(
1002 data_callback: &mut D,
1003 interleaved: &mut [u8],
1004 silence_asio_buffer: bool,
1005 little_endian: bool,
1006 asio_stream: &mut sys::AsioStream,
1007 asio_info: &sys::CallbackInfo,
1008 sample_rate: crate::SampleRate,
1009 hardware_latency_frames: usize,
1010) where
1011 D: FnMut(&mut Data, &OutputCallbackInfo) + Send + 'static,
1012{
1013 let format = SampleFormat::I24;
1014 let interleaved: &mut [I24] = cast_slice_mut(interleaved);
1015 apply_output_callback_to_data::<I24, _>(
1016 data_callback,
1017 interleaved,
1018 asio_info,
1019 sample_rate,
1020 format,
1021 hardware_latency_frames,
1022 );
1023
1024 // Size of samples in the ASIO buffer (has to be 3 in this case)
1025 let asio_sample_size_bytes = 3;
1026 let n_channels = interleaved.len() / asio_stream.buffer_size as usize;
1027 let buffer_index = asio_info.buffer_index as usize;
1028
1029 // Write interleaved samples to ASIO channels, one channel at a time.
1030 for ch_ix in 0..n_channels {
1031 // Take channel as u8 array ([u8; 3] packets to represent i24)
1032 let asio_channel = asio_channel_slice_mut(
1033 asio_stream,
1034 buffer_index,
1035 ch_ix,
1036 Some(asio_stream.buffer_size as usize * asio_sample_size_bytes),
1037 );
1038
1039 if silence_asio_buffer {
1040 asio_channel.align_to_mut::<u8>().1.fill(0);
1041 }
1042
1043 // Fill in every channel from the interleaved vector
1044 for (channel_sample, sample_in_buffer) in asio_channel
1045 .chunks_mut(asio_sample_size_bytes)
1046 .zip(interleaved.iter().skip(ch_ix).step_by(n_channels))
1047 {
1048 // Add samples from buffer if no silence was applied, otherwise just overwrite
1049 let result = if silence_asio_buffer {
1050 sample_in_buffer.inner()
1051 } else {
1052 let sample = i24_bytes_to_i32(
1053 &[channel_sample[0], channel_sample[1], channel_sample[2]],
1054 little_endian,
1055 );
1056 (sample_in_buffer.inner() + sample).clamp(-8388608, 8388607)
1057 };
1058 let bytes = result.to_le_bytes();
1059 if little_endian {
1060 channel_sample[0] = bytes[0];
1061 channel_sample[1] = bytes[1];
1062 channel_sample[2] = bytes[2];
1063 } else {
1064 channel_sample[2] = bytes[0];
1065 channel_sample[1] = bytes[1];
1066 channel_sample[0] = bytes[2];
1067 }
1068 }
1069 }
1070}
1071
1072unsafe fn process_input_callback_i24<D>(
1073 data_callback: &mut D,
1074 interleaved: &mut [u8],
1075 asio_stream: &sys::AsioStream,
1076 asio_info: &sys::CallbackInfo,
1077 sample_rate: crate::SampleRate,
1078 little_endian: bool,
1079 hardware_latency_frames: usize,
1080) where
1081 D: FnMut(&Data, &InputCallbackInfo) + Send + 'static,
1082{
1083 let format = SampleFormat::I24;
1084
1085 // 1. Write the ASIO channels to the CPAL buffer.
1086 let interleaved: &mut [I24] = cast_slice_mut(interleaved);
1087 let n_frames = asio_stream.buffer_size as usize;
1088 let n_channels = interleaved.len() / n_frames;
1089 let buffer_index = asio_info.buffer_index as usize;
1090 let asio_sample_size_bytes = 3;
1091
1092 for ch_ix in 0..n_channels {
1093 let asio_channel = asio_channel_slice::<u8>(
1094 asio_stream,
1095 buffer_index,
1096 ch_ix,
1097 Some(n_frames * asio_sample_size_bytes),
1098 );
1099 for (channel_sample, sample_in_buffer) in asio_channel
1100 .chunks(asio_sample_size_bytes)
1101 .zip(interleaved.iter_mut().skip(ch_ix).step_by(n_channels))
1102 {
1103 let sample = i24_bytes_to_i32(
1104 &[channel_sample[0], channel_sample[1], channel_sample[2]],
1105 little_endian,
1106 );
1107 *sample_in_buffer = I24::new(sample).unwrap();
1108 }
1109 }
1110
1111 // 2. Deliver the interleaved buffer to the callback.
1112 apply_input_callback_to_data::<I24, _>(
1113 data_callback,
1114 interleaved,
1115 asio_info,
1116 sample_rate,
1117 format,
1118 hardware_latency_frames,
1119 );
1120}
1121
1122/// Apply the output callback to the interleaved buffer.
1123unsafe fn apply_output_callback_to_data<A, D>(
1124 data_callback: &mut D,
1125 interleaved: &mut [A],
1126 asio_info: &sys::CallbackInfo,
1127 sample_rate: crate::SampleRate,
1128 sample_format: SampleFormat,
1129 hardware_latency_frames: usize,
1130) where
1131 A: Copy,
1132 D: FnMut(&mut Data, &OutputCallbackInfo) + Send + 'static,
1133{
1134 let mut data = Data::from_parts(
1135 interleaved.as_mut_ptr() as *mut (),
1136 interleaved.len(),
1137 sample_format,
1138 );
1139 let callback = crate::StreamInstant::from_nanos_i128(asio_info.system_time as i128)
1140 .expect("`system_time` out of range of `StreamInstant` representation");
1141 let delay = frames_to_duration(hardware_latency_frames, sample_rate);
1142 let playback = callback
1143 .add(delay)
1144 .expect("`playback` occurs beyond representation supported by `StreamInstant`");
1145 let timestamp = crate::OutputStreamTimestamp { callback, playback };
1146 let info = OutputCallbackInfo { timestamp };
1147 data_callback(&mut data, &info);
1148}
1149
1150/// Apply the input callback to the interleaved buffer.
1151unsafe fn apply_input_callback_to_data<A, D>(
1152 data_callback: &mut D,
1153 interleaved: &mut [A],
1154 asio_info: &sys::CallbackInfo,
1155 sample_rate: crate::SampleRate,
1156 format: SampleFormat,
1157 hardware_latency_frames: usize,
1158) where
1159 A: Copy,
1160 D: FnMut(&Data, &InputCallbackInfo) + Send + 'static,
1161{
1162 let data = Data::from_parts(
1163 interleaved.as_mut_ptr() as *mut (),
1164 interleaved.len(),
1165 format,
1166 );
1167 let callback = crate::StreamInstant::from_nanos_i128(asio_info.system_time as i128)
1168 .expect("`system_time` out of range of `StreamInstant` representation");
1169 let delay = frames_to_duration(hardware_latency_frames, sample_rate);
1170 let capture = callback
1171 .sub(delay)
1172 .expect("`capture` occurs before origin of alsa `StreamInstant`");
1173 let timestamp = crate::InputStreamTimestamp { callback, capture };
1174 let info = InputCallbackInfo { timestamp };
1175 data_callback(&data, &info);
1176}