From 3dfe36f6abb279dd8c9b04a8e2ddc70ecc2cf29f Mon Sep 17 00:00:00 2001 From: Roderick van Domburg Date: Sun, 20 Sep 2026 12:48:37 +0200 Subject: [PATCH 1/8] fix(coreaudio): track iOS buffer depth in seconds for timestamps and stop() --- CHANGELOG.md | 1 + src/host/coreaudio/ios/mod.rs | 145 +++++++++--------- .../coreaudio/ios/session_event_manager.rs | 23 ++- src/host/mod.rs | 18 ++- 4 files changed, 105 insertions(+), 82 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 2870a9625..265435b90 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -58,6 +58,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - **CoreAudio**: Fix the device running at a different sample rate from the stream on hardware that reports a continuous rate range. - **CoreAudio**: Fix `supported_configs()` only reporting `F32`, even on hardware that also supports other sample formats. - **CoreAudio**: Fix sample rate changes timing out early when the device reports other rates first. +- **iOS**: Fix timestamps and `buffer_size()` being off when the stream sample rate differs from the hardware rate. - **JACK**: Channel enumeration is capped at the physical system port count again. - **JACK**: Streams no longer panic when the server delivers a larger period than the negotiated buffer size. - **PipeWire**: Fix an empty chunk being emitted when a cycle requests no frames. diff --git a/src/host/coreaudio/ios/mod.rs b/src/host/coreaudio/ios/mod.rs index 9a1eaf898..b7d1ed405 100644 --- a/src/host/coreaudio/ios/mod.rs +++ b/src/host/coreaudio/ios/mod.rs @@ -5,7 +5,7 @@ use std::{ ptr::NonNull, sync::{ Arc, Mutex, - atomic::{AtomicBool, AtomicUsize, Ordering}, + atomic::{AtomicBool, AtomicU64, Ordering}, }, time::Duration, }; @@ -30,7 +30,7 @@ use crate::{ SupportedStreamConfigRange, host::{ ErrorCallbackArc, equilibrium::fill_equilibrium, frames_to_duration, latch::Latch, - try_emit_error, + secs_to_nanos, try_emit_error, wait_for_drain, }, traits::{DeviceTrait, HostTrait, StreamTrait}, }; @@ -221,11 +221,11 @@ impl DeviceTrait for Device { let mut audio_unit = setup_stream_audio_unit(config, sample_format, true)?; // Buffer depth used to offset capture timestamps. AVAudioSession auto-reroutes to a new device, // which changes this depth, and is then refreshed by the session event manager. - let latency_frames = Arc::new(AtomicUsize::new(input_latency_frames())); + let latency_nanos = Arc::new(AtomicU64::new(input_latency_nanos())); let error_callback: ErrorCallbackArc = Arc::new(Mutex::new(error_callback)); let session_error_callback = error_callback.clone(); - let session_latency_frames = latency_frames.clone(); + let session_latency_nanos = latency_nanos.clone(); let draining = Arc::new(AtomicBool::new(false)); @@ -234,7 +234,7 @@ impl DeviceTrait for Device { &mut audio_unit, sample_format, config.sample_rate, - latency_frames, + latency_nanos, draining.clone(), data_callback, move |e| { @@ -249,15 +249,16 @@ impl DeviceTrait for Device { let session_manager = SessionEventManager::new( session_error_callback, Latch::new(), - Some((session_latency_frames, true)), + (session_latency_nanos, true), Arc::downgrade(&inner), ); let stream = Stream { inner, session_manager, - // Draining is meaningless for capture, so the drain window is zero. draining, - drain_window: Duration::ZERO, + // Capture never drains, so there is nothing to wait out. + drain_nanos: Arc::new(AtomicU64::new(0)), + sample_rate: config.sample_rate, }; stream.signal_ready(); Ok(stream) @@ -286,22 +287,21 @@ impl DeviceTrait for Device { // Buffer depth used to offset playback timestamps. AVAudioSession auto-reroutes to a new device, // which changes this depth, and is then refreshed by the session event manager. - let latency_frames = Arc::new(AtomicUsize::new(output_latency_frames())); + let latency_nanos = Arc::new(AtomicU64::new(output_latency_nanos())); let error_callback: ErrorCallbackArc = Arc::new(Mutex::new(error_callback)); let session_error_callback = error_callback.clone(); - let session_latency_frames = latency_frames.clone(); + let session_latency_nanos = latency_nanos.clone(); let draining = Arc::new(AtomicBool::new(false)); - let drain_window = - frames_to_duration(get_device_buffer_frames() as FrameCount, config.sample_rate); + let drain_nanos = latency_nanos.clone(); // Set up output callback setup_output_callback( &mut audio_unit, sample_format, config.sample_rate, - latency_frames, + latency_nanos, draining.clone(), data_callback, move |e| { @@ -316,14 +316,15 @@ impl DeviceTrait for Device { let session_manager = SessionEventManager::new( session_error_callback, Latch::new(), - Some((session_latency_frames, false)), + (session_latency_nanos, false), Arc::downgrade(&inner), ); let stream = Stream { inner, session_manager, draining, - drain_window, + drain_nanos, + sample_rate: config.sample_rate, }; stream.signal_ready(); Ok(stream) @@ -334,7 +335,9 @@ pub struct Stream { inner: Arc>, session_manager: SessionEventManager, draining: Arc, - drain_window: Duration, + // Shared with the session manager so stop() sees the current depth. Zero on capture: no drain. + drain_nanos: Arc, + sample_rate: SampleRate, } impl Stream { @@ -390,12 +393,10 @@ impl StreamTrait for Stream { fn stop(&self, timeout: Option) -> Result<(), Error> { self.draining.store(true, Ordering::Relaxed); - if timeout != Some(Duration::ZERO) { - let wait = timeout.map_or(self.drain_window, |t| self.drain_window.min(t)); - if !wait.is_zero() { - std::thread::sleep(wait); - } - } + wait_for_drain( + Duration::from_nanos(self.drain_nanos.load(Ordering::Relaxed)), + timeout, + ); self.halt() } @@ -406,7 +407,7 @@ impl StreamTrait for Stream { } fn buffer_size(&self) -> Result { - Ok(get_device_buffer_frames() as FrameCount) + Ok(device_buffer_frames(self.sample_rate) as FrameCount) } } @@ -465,40 +466,34 @@ fn set_audio_session_buffer_size( Ok(()) } -/// Get the actual buffer size from AVAudioSession. +/// IO buffer size in frames at `sample_rate`. /// -/// This queries the current IO buffer duration from AVAudioSession and converts -/// it to frames based on the current sample rate. -fn get_device_buffer_frames() -> usize { +/// AVAudioSession reports a duration, and its own rate is the hardware's; the audio unit resamples +/// between that and the stream's, so convert at the stream's rate to get frames per callback. +fn device_buffer_frames(sample_rate: SampleRate) -> usize { // SAFETY: AVAudioSession methods are safe to call on the singleton instance - unsafe { - let audio_session = AVAudioSession::sharedInstance(); - let buffer_duration = audio_session.IOBufferDuration(); - let sample_rate = audio_session.sampleRate(); - // Round: the duration comes from an integer frame count, so the product is a whole number - // that floating-point can render as N.9999, which truncation would drop to N-1. - (buffer_duration * sample_rate).round() as usize - } + let buffer_duration = unsafe { AVAudioSession::sharedInstance().IOBufferDuration() }; + // Round: the duration comes from an integer frame count, so the product is a whole number + // that floating-point can render as N.9999, which truncation would drop to N-1. + (buffer_duration * sample_rate as f64).round() as usize } -/// Total capture buffer depth in frames: the IO buffer plus the hardware input latency. -pub(super) fn input_latency_frames() -> usize { +/// Total capture buffer depth: the IO buffer plus the hardware input latency. +pub(super) fn input_latency_nanos() -> u64 { // SAFETY: AVAudioSession methods are safe to call on the singleton instance - let extra = unsafe { + unsafe { let audio_session = AVAudioSession::sharedInstance(); - (audio_session.inputLatency() * audio_session.sampleRate()).round() as usize - }; - get_device_buffer_frames() + extra + secs_to_nanos(audio_session.IOBufferDuration() + audio_session.inputLatency()) + } } -/// Total playback buffer depth in frames: the IO buffer plus the hardware output latency. -pub(super) fn output_latency_frames() -> usize { +/// Total playback buffer depth: the IO buffer plus the hardware output latency. +pub(super) fn output_latency_nanos() -> u64 { // SAFETY: AVAudioSession methods are safe to call on the singleton instance - let extra = unsafe { + unsafe { let audio_session = AVAudioSession::sharedInstance(); - (audio_session.outputLatency() * audio_session.sampleRate()).round() as usize - }; - get_device_buffer_frames() + extra + secs_to_nanos(audio_session.IOBufferDuration() + audio_session.outputLatency()) + } } // Typical iOS hardware buffer frame limits according to Apple Technical Q&A QA1631. @@ -625,12 +620,32 @@ unsafe fn extract_audio_buffer( (buffer, data) } +/// Buffer depth for one callback. +/// +/// Refreshed on route changes by the session event manager; falls back to this buffer's own frame +/// count, converted at the stream's rate, when the depth is unknown (zero). +#[inline] +fn resolve_delay( + latency_nanos: &AtomicU64, + len: usize, + channels: usize, + sample_rate: SampleRate, +) -> Duration { + match latency_nanos.load(Ordering::Relaxed) { + 0 => frames_to_duration( + len.checked_div(channels).unwrap_or(0) as FrameCount, + sample_rate, + ), + n => Duration::from_nanos(n), + } +} + /// Setup input callback with proper latency calculation. fn setup_input_callback( audio_unit: &mut AudioUnit, sample_format: SampleFormat, sample_rate: SampleRate, - latency_frames: Arc, + latency_nanos: Arc, draining: Arc, mut data_callback: D, mut error_callback: E, @@ -656,16 +671,12 @@ where Ok(cb) => cb, }; - // Refreshed on route changes by the session event manager; fall back to this buffer's - // own frame count if the depth is unknown (zero). - let latency_frames = match latency_frames.load(Ordering::Relaxed) { - 0 => { - let channels = buffer.mNumberChannels as usize; - data.len().checked_div(channels).unwrap_or(0) - } - n => n, - }; - let delay = frames_to_duration(latency_frames as FrameCount, sample_rate); + let delay = resolve_delay( + &latency_nanos, + data.len(), + buffer.mNumberChannels as usize, + sample_rate, + ); let capture = callback.checked_sub(delay).unwrap_or(StreamInstant::ZERO); let timestamp = StreamTimestamp { callback, @@ -691,7 +702,7 @@ fn setup_output_callback( audio_unit: &mut AudioUnit, sample_format: SampleFormat, sample_rate: SampleRate, - latency_frames: Arc, + latency_nanos: Arc, draining: Arc, mut data_callback: D, mut error_callback: E, @@ -730,16 +741,12 @@ where Ok(cb) => cb, }; - // Refreshed on route changes by the session event manager; fall back to this buffer's - // own frame count if the depth is unknown (zero). - let latency_frames = match latency_frames.load(Ordering::Relaxed) { - 0 => { - let channels = buffer.mNumberChannels as usize; - data.len().checked_div(channels).unwrap_or(0) - } - n => n, - }; - let delay = frames_to_duration(latency_frames as FrameCount, sample_rate); + let delay = resolve_delay( + &latency_nanos, + data.len(), + buffer.mNumberChannels as usize, + sample_rate, + ); let playback = callback + delay; let timestamp = StreamTimestamp { callback, diff --git a/src/host/coreaudio/ios/session_event_manager.rs b/src/host/coreaudio/ios/session_event_manager.rs index e6a696b80..1ca497966 100644 --- a/src/host/coreaudio/ios/session_event_manager.rs +++ b/src/host/coreaudio/ios/session_event_manager.rs @@ -5,7 +5,7 @@ use std::{ ptr::NonNull, sync::{ Arc, Mutex, Weak, - atomic::{AtomicUsize, Ordering}, + atomic::{AtomicU64, Ordering}, }, }; @@ -19,7 +19,7 @@ use objc2_avf_audio::{ }; use objc2_foundation::{NSNotification, NSNotificationCenter, NSNumber, NSString}; -use super::{StreamInner, input_latency_frames, output_latency_frames}; +use super::{StreamInner, input_latency_nanos, output_latency_nanos}; use crate::{ Error, ErrorKind, host::{ErrorCallbackArc, emit_error, latch::Latch}, @@ -43,7 +43,7 @@ fn user_info_number(notification: &NSNotification, key: Option<&NSString>) -> Op /// Shared buffer-depth value to refresh on route changes, paired with `is_input` to select the /// input or output latency. `true` means an input stream. -type LatencyRefresh = (Arc, bool); +type LatencyRefresh = (Arc, bool); fn route_change_error(notification: &NSNotification) -> Option { let key = unsafe { AVAudioSessionRouteChangeReasonKey }; @@ -86,7 +86,7 @@ impl SessionEventManager { pub(super) fn new( error_callback: ErrorCallbackArc, latch: Latch, - latency_refresh: Option, + latency_refresh: LatencyRefresh, stream: Weak>, ) -> Self { let nc = NSNotificationCenter::defaultCenter(); @@ -137,14 +137,13 @@ impl SessionEventManager { if w.is_released() { // The route may have changed the active device; recompute the buffer depth so // capture/playback timestamps track the new latency. - if let Some((frames, is_input)) = &latency_refresh { - let depth = if *is_input { - input_latency_frames() - } else { - output_latency_frames() - }; - frames.store(depth, Ordering::Relaxed); - } + let (nanos, is_input) = &latency_refresh; + let depth = if *is_input { + input_latency_nanos() + } else { + output_latency_nanos() + }; + nanos.store(depth, Ordering::Relaxed); let notif = unsafe { notif.as_ref() }; if let Some(err) = route_change_error(notif) { emit_error(&cb, err); diff --git a/src/host/mod.rs b/src/host/mod.rs index 653ee771b..f80c2f4b1 100644 --- a/src/host/mod.rs +++ b/src/host/mod.rs @@ -274,6 +274,19 @@ pub(crate) const fn frames_to_duration( std::time::Duration::new(secs, nanos as u32) } +/// Converts a duration in seconds, as reported by the platform's audio session, to nanoseconds. +/// +/// Returns 0 for a value that is not finite and positive. +#[cfg(all(target_vendor = "apple", not(target_os = "macos")))] +#[inline] +pub(crate) fn secs_to_nanos(secs: f64) -> u64 { + if secs.is_finite() && secs > 0.0 { + (secs * 1_000_000_000.0).round() as u64 + } else { + 0 + } +} + /// Waits out `window` of buffered audio, cut short by `timeout` when it is the smaller of the two. /// /// Implements [`StreamTrait::stop`]'s timeout contract for the backends that approximate a drain @@ -281,7 +294,10 @@ pub(crate) const fn frames_to_duration( /// `Some(Duration::ZERO)` returns immediately. /// /// [`StreamTrait::stop`]: crate::traits::StreamTrait::stop -#[cfg(windows)] +#[cfg(any( + all(windows, feature = "asio"), + all(target_vendor = "apple", not(target_os = "macos")) +))] pub(crate) fn wait_for_drain(window: std::time::Duration, timeout: Option) { let wait = timeout.map_or(window, |t| window.min(t)); if !wait.is_zero() { From 9dbf7953ad04cf51e6c9d155c63bdb27cad253c8 Mon Sep 17 00:00:00 2001 From: Roderick van Domburg Date: Sun, 20 Sep 2026 12:49:11 +0200 Subject: [PATCH 2/8] fix(coreaudio): drain macOS stop() by the current buffer depth --- src/host/coreaudio/macos/device.rs | 16 ++++++++++------ src/host/coreaudio/macos/mod.rs | 30 +++++++++++++++++++----------- src/host/mod.rs | 5 +---- 3 files changed, 30 insertions(+), 21 deletions(-) diff --git a/src/host/coreaudio/macos/device.rs b/src/host/coreaudio/macos/device.rs index 4402a1004..1fc8eb84f 100644 --- a/src/host/coreaudio/macos/device.rs +++ b/src/host/coreaudio/macos/device.rs @@ -924,7 +924,14 @@ impl Device { error_callback_disconnect, pending_xrun_overload, )?); - let stream = Stream::new(inner_arc, monitor, draining, Duration::ZERO); + // Capture never drains, so there are no frames to wait out. + let stream = Stream::new( + inner_arc, + monitor, + draining, + Arc::new(AtomicUsize::new(0)), + sample_rate, + ); stream.signal_ready(); Ok(stream) } @@ -1001,10 +1008,7 @@ impl Device { let callback_latency_frames = latency_frames.clone(); let draining = Arc::new(AtomicBool::new(false)); let draining_render = draining.clone(); - let drain_window = frames_to_duration( - (device_buffer_frames.unwrap_or(0) + extra_latency_frames) as FrameCount, - sample_rate, - ); + let drain_frames = latency_frames.clone(); type Args = render_callback::Args; audio_unit.set_render_callback(move |args: Args| unsafe { @@ -1082,7 +1086,7 @@ impl Device { pending_xrun_overload, )?) }; - let stream = Stream::new(inner_arc, monitor, draining, drain_window); + let stream = Stream::new(inner_arc, monitor, draining, drain_frames, sample_rate); stream.signal_ready(); Ok(stream) } diff --git a/src/host/coreaudio/macos/mod.rs b/src/host/coreaudio/macos/mod.rs index 7b53a40a5..ece9f57e8 100644 --- a/src/host/coreaudio/macos/mod.rs +++ b/src/host/coreaudio/macos/mod.rs @@ -19,8 +19,11 @@ use property_listener::AudioObjectPropertyListener; pub use self::enumerate::{Devices, default_input_device, default_output_device}; use super::{OSStatus, asbd_from_config, check_os_status, host_time_to_stream_instant}; use crate::{ - Error, ErrorKind, FrameCount, ResultExt, StreamInstant, - host::{coreaudio::macos::loopback::LoopbackDevice, emit_error, latch::Latch}, + Error, ErrorKind, FrameCount, ResultExt, SampleRate, StreamInstant, + host::{ + coreaudio::macos::loopback::LoopbackDevice, emit_error, frames_to_duration, latch::Latch, + wait_for_drain, + }, traits::{HostTrait, StreamTrait}, }; @@ -414,7 +417,9 @@ pub struct Stream { inner: Arc>, monitor: Box, draining: Arc, - drain_window: Duration, + // Shared with the reroute monitor so stop() sees the current depth. Zero on capture: no drain. + drain_frames: Arc, + sample_rate: SampleRate, } impl Stream { @@ -422,13 +427,15 @@ impl Stream { inner: Arc>, monitor: Box, draining: Arc, - drain_window: Duration, + drain_frames: Arc, + sample_rate: SampleRate, ) -> Self { Self { inner, monitor, draining, - drain_window, + drain_frames, + sample_rate, } } @@ -463,12 +470,13 @@ impl StreamTrait for Stream { fn stop(&self, timeout: Option) -> Result<(), Error> { self.draining.store(true, Ordering::Relaxed); - if timeout != Some(Duration::ZERO) { - let wait = timeout.map_or(self.drain_window, |t| self.drain_window.min(t)); - if !wait.is_zero() { - std::thread::sleep(wait); - } - } + wait_for_drain( + frames_to_duration( + self.drain_frames.load(Ordering::Relaxed) as FrameCount, + self.sample_rate, + ), + timeout, + ); self.inner .lock() diff --git a/src/host/mod.rs b/src/host/mod.rs index f80c2f4b1..8e05d77be 100644 --- a/src/host/mod.rs +++ b/src/host/mod.rs @@ -294,10 +294,7 @@ pub(crate) fn secs_to_nanos(secs: f64) -> u64 { /// `Some(Duration::ZERO)` returns immediately. /// /// [`StreamTrait::stop`]: crate::traits::StreamTrait::stop -#[cfg(any( - all(windows, feature = "asio"), - all(target_vendor = "apple", not(target_os = "macos")) -))] +#[cfg(any(all(windows, feature = "asio"), all(target_vendor = "apple")))] pub(crate) fn wait_for_drain(window: std::time::Duration, timeout: Option) { let wait = timeout.map_or(window, |t| window.min(t)); if !wait.is_zero() { From 510d2014383ac5d9eb4058a29539428d1ca1bf5c Mon Sep 17 00:00:00 2001 From: Roderick van Domburg Date: Sun, 20 Sep 2026 12:49:48 +0200 Subject: [PATCH 3/8] fix(coreaudio): refresh macOS timestamp latency when the device buffer size changes --- CHANGELOG.md | 1 + src/host/coreaudio/macos/device.rs | 59 +++- src/host/coreaudio/macos/mod.rs | 450 ++++++++++++++++++----------- 3 files changed, 335 insertions(+), 175 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 265435b90..623362d52 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -58,6 +58,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - **CoreAudio**: Fix the device running at a different sample rate from the stream on hardware that reports a continuous rate range. - **CoreAudio**: Fix `supported_configs()` only reporting `F32`, even on hardware that also supports other sample formats. - **CoreAudio**: Fix sample rate changes timing out early when the device reports other rates first. +- **CoreAudio**: Fix timestamps going stale when the device's buffer size changes during a stream. - **iOS**: Fix timestamps and `buffer_size()` being off when the stream sample rate differs from the hardware rate. - **JACK**: Channel enumeration is capped at the physical system port count again. - **JACK**: Streams no longer panic when the server delivers a larger period than the negotiated buffer size. diff --git a/src/host/coreaudio/macos/device.rs b/src/host/coreaudio/macos/device.rs index 1fc8eb84f..26e82c2ec 100644 --- a/src/host/coreaudio/macos/device.rs +++ b/src/host/coreaudio/macos/device.rs @@ -866,6 +866,11 @@ impl Device { let (bytes_per_channel, sample_rate, device_buffer_frames, extra_latency_frames) = setup_callback_vars(&audio_unit, config, sample_format, Scope::Input); + // A resized device buffer changes this depth, and is then refreshed by DisconnectManager. + let latency_frames = Arc::new(AtomicUsize::new( + device_buffer_frames.map_or(0, |frames| frames + extra_latency_frames), + )); + let callback_latency_frames = latency_frames.clone(); let draining = Arc::new(AtomicBool::new(false)); let draining_input = draining.clone(); @@ -892,9 +897,12 @@ impl Device { } Ok(cb) => cb, }; - let buffer_frames = len / channels as usize; - let latency_frames = - device_buffer_frames.unwrap_or(buffer_frames) + extra_latency_frames; + let latency_frames = resolve_latency_frames( + &callback_latency_frames, + len, + channels as usize, + extra_latency_frames, + ); let delay = frames_to_duration(latency_frames as FrameCount, sample_rate); let capture = callback.checked_sub(delay).unwrap_or(StreamInstant::ZERO); let timestamp = StreamTimestamp { @@ -922,6 +930,7 @@ impl Device { self.audio_device_id, weak_inner, error_callback_disconnect, + (latency_frames, Scope::Input), pending_xrun_overload, )?); // Capture never drains, so there are no frames to wait out. @@ -1040,12 +1049,12 @@ impl Device { } Ok(cb) => cb, }; - let latency_frames = match callback_latency_frames.load(Ordering::Relaxed) { - // Depth unknown (query failed): estimate the device buffer from this callback, but - // still add the safety offset and device latency, matching the input path. - 0 => len / channels as usize + extra_latency_frames, - n => n, - }; + let latency_frames = resolve_latency_frames( + &callback_latency_frames, + len, + channels as usize, + extra_latency_frames, + ); let delay = frames_to_duration(latency_frames as FrameCount, sample_rate); let playback = callback + delay; let timestamp = StreamTimestamp { @@ -1075,14 +1084,16 @@ impl Device { Box::new(DefaultOutputMonitor::new( weak_inner, error_callback, - Some((latency_frames, Scope::Output)), + (latency_frames.clone(), Scope::Output), pending_xrun_overload, )?) } else { + // An explicit device never reroutes, so a resized buffer is the only depth change. Box::new(DisconnectManager::new( self.audio_device_id, weak_inner, error_callback, + (latency_frames.clone(), Scope::Output), pending_xrun_overload, )?) }; @@ -1186,6 +1197,34 @@ pub(crate) fn get_device_extra_latency_frames(audio_unit: &AudioUnit, scope: Sco (device_latency + safety_offset) as usize } +/// Total buffer depth in frames: the IO buffer plus the device's own latency, or 0 if the +/// device buffer size cannot be queried. +pub(crate) fn device_latency_frames(audio_unit: &AudioUnit, scope: Scope) -> usize { + get_device_buffer_frame_size(audio_unit) + .ok() + .map_or(0, |buffer| { + buffer + get_device_extra_latency_frames(audio_unit, scope) + }) +} + +/// Buffer depth for one callback. +/// +/// Refreshed by the stream's monitor when the device buffer is resized; falls back to estimating +/// the device buffer from this callback, plus the fixed device latency and safety offset, when +/// the depth is unknown (zero). +#[inline] +fn resolve_latency_frames( + cached: &AtomicUsize, + len: usize, + channels: usize, + extra_latency_frames: usize, +) -> usize { + match cached.load(Ordering::Relaxed) { + 0 => len.checked_div(channels).unwrap_or(0) + extra_latency_frames, + n => n, + } +} + /// Setup common callback variables, querying both the I/O buffer size and extra hardware latency. /// /// Returns `(bytes_per_channel, sample_rate, device_buffer_frames, extra_latency_frames)` diff --git a/src/host/coreaudio/macos/mod.rs b/src/host/coreaudio/macos/mod.rs index ece9f57e8..d9fb3aee4 100644 --- a/src/host/coreaudio/macos/mod.rs +++ b/src/host/coreaudio/macos/mod.rs @@ -9,7 +9,8 @@ use std::{ use coreaudio::audio_unit::{AudioUnit, Scope}; use objc2_core_audio::{ - AudioDeviceID, AudioObjectID, AudioObjectPropertyAddress, kAudioDeviceProcessorOverload, + AudioDeviceID, AudioObjectID, AudioObjectPropertyAddress, AudioObjectPropertySelector, + kAudioDeviceProcessorOverload, kAudioDevicePropertyBufferFrameSize, kAudioDevicePropertyDeviceIsAlive, kAudioDevicePropertyNominalSampleRate, kAudioHardwarePropertyDefaultOutputDevice, kAudioObjectPropertyElementMain, kAudioObjectPropertyScopeGlobal, kAudioObjectSystemObject, @@ -68,6 +69,38 @@ impl HostTrait for Host { /// Type alias for the error callback to reduce complexity type ErrorCallback = dyn FnMut(Error) + Send; +/// The cached buffer depth to refresh when the device's buffer frame size changes, and the scope +/// it was measured in. +type LatencyRefresh = (Arc, Scope); + +/// What a device's property listeners report to its delivery thread. +enum MonitorEvent { + /// The device went away, or changed such that the stream is no longer valid. + Lost(Error), + /// The device's buffer frame size changed, so the cached buffer depth is stale. + LatencyChanged, +} + +/// Device-scoped property address, the shape every listener here uses. +fn device_address(selector: AudioObjectPropertySelector) -> AudioObjectPropertyAddress { + AudioObjectPropertyAddress { + mSelector: selector, + mScope: kAudioObjectPropertyScopeGlobal, + mElement: kAudioObjectPropertyElementMain, + } +} + +/// Recomputes the buffer depth for `refresh` from the stream's current device. +fn refresh_latency(stream: &Mutex, refresh: &LatencyRefresh) { + let (frames, scope) = refresh; + if let Ok(inner) = stream.lock() { + frames.store( + device::device_latency_frames(&inner.audio_unit, *scope), + Ordering::Relaxed, + ); + } +} + /// Spawns a dedicated thread that registers a single property listener, calling `on_change` on /// each firing. The listener is deregistered when the returned `Sender<()>` is dropped. fn spawn_property_listener_thread( @@ -104,6 +137,74 @@ where Ok(shutdown_tx) } +/// Spawns the delivery thread shared by both monitors. +/// +/// It waits for the owning `Stream` to reach the caller, then applies each event until the stream +/// is dropped or every sender is gone. Property listeners must not call back into CoreAudio, so +/// they only post events and this thread does the work. +fn spawn_delivery_thread( + name: &str, + latch: &mut Latch, + events: mpsc::Receiver, + stream_weak: Weak>, + mut on_event: F, +) -> Result<(), Error> +where + E: Send + 'static, + F: FnMut(E, &Arc>) + Send + 'static, +{ + let waiter = latch.waiter(); + let handle = std::thread::Builder::new() + .name(name.to_owned()) + .spawn(move || { + // If the Latch is dropped without being released (error path), exit cleanly. + if !waiter.wait() { + return; + } + while let Ok(event) = events.recv() { + let Some(stream) = stream_weak.upgrade() else { + break; + }; + on_event(event, &stream); + } + }) + .map_err(|e| { + Error::with_message( + ErrorKind::ResourceExhausted, + format!("failed to spawn {name} thread: {e}"), + ) + })?; + latch.add_thread(handle.thread().clone()); + Ok(()) +} + +/// Halts the stream and reports why, for the cases where it cannot keep running. +fn report_lost( + stream: &Mutex, + error_callback: &Arc>, + err: Error, +) { + if let Ok(mut inner) = stream.try_lock() { + let _ = inner.pause(); + } + emit_error(error_callback, err); +} + +/// Registers an overload listener for `device_id`. These fire on the RT thread, so the callback +/// only sets a flag. +fn spawn_overload_listener( + device_id: AudioDeviceID, + pending_xrun: Arc, +) -> Result, Error> { + spawn_property_listener_thread( + device_id, + device_address(kAudioDeviceProcessorOverload), + move || { + pending_xrun.store(true, Ordering::Relaxed); + }, + ) +} + /// A device monitor that can signal when the owning `Stream` handle has been returned to the /// caller, allowing the delivery thread to start processing events. pub(super) trait Monitor: Send + Sync { @@ -130,61 +231,71 @@ impl DisconnectManager { device_id: AudioDeviceID, stream_weak: Weak>, error_callback: Arc>, + latency_refresh: LatencyRefresh, pending_xrun: Arc, ) -> Result { let (shutdown_tx, shutdown_rx) = mpsc::channel(); - let (disconnect_tx, disconnect_rx) = mpsc::channel::(); + let (disconnect_tx, disconnect_rx) = mpsc::channel::(); let (ready_tx, ready_rx) = mpsc::channel(); // Spawn a dedicated thread to own all listeners. CoreAudio requires that // AudioObjectPropertyListeners are added and removed on the same thread. let disconnect_tx_alive = disconnect_tx.clone(); - let disconnect_tx_rate = disconnect_tx; + let disconnect_tx_rate = disconnect_tx.clone(); + let disconnect_tx_buffer = disconnect_tx; std::thread::spawn(move || { - let alive_address = AudioObjectPropertyAddress { - mSelector: kAudioDevicePropertyDeviceIsAlive, - mScope: kAudioObjectPropertyScopeGlobal, - mElement: kAudioObjectPropertyElementMain, - }; - let alive_listener = - AudioObjectPropertyListener::new(device_id, alive_address, move || { - let _ = disconnect_tx_alive.send(Error::with_message( + let alive_listener = AudioObjectPropertyListener::new( + device_id, + device_address(kAudioDevicePropertyDeviceIsAlive), + move || { + let _ = disconnect_tx_alive.send(MonitorEvent::Lost(Error::with_message( ErrorKind::DeviceNotAvailable, "Device disconnected", - )); - }); - - let rate_address = AudioObjectPropertyAddress { - mSelector: kAudioDevicePropertyNominalSampleRate, - mScope: kAudioObjectPropertyScopeGlobal, - mElement: kAudioObjectPropertyElementMain, - }; - let rate_listener = - AudioObjectPropertyListener::new(device_id, rate_address, move || { - let _ = disconnect_tx_rate.send(Error::with_message( + ))); + }, + ); + + let rate_listener = AudioObjectPropertyListener::new( + device_id, + device_address(kAudioDevicePropertyNominalSampleRate), + move || { + let _ = disconnect_tx_rate.send(MonitorEvent::Lost(Error::with_message( ErrorKind::StreamInvalidated, "Device sample rate changed", - )); - }); + ))); + }, + ); + + // Device-global on macOS: another process changing it resizes our IO buffer too. + let buffer_size_listener = AudioObjectPropertyListener::new( + device_id, + device_address(kAudioDevicePropertyBufferFrameSize), + move || { + let _ = disconnect_tx_buffer.send(MonitorEvent::LatencyChanged); + }, + ); // Overload notifications fire on the RT thread. - let overload_address = AudioObjectPropertyAddress { - mSelector: kAudioDeviceProcessorOverload, - mScope: kAudioObjectPropertyScopeGlobal, - mElement: kAudioObjectPropertyElementMain, - }; - let overload_listener = - AudioObjectPropertyListener::new(device_id, overload_address, move || { + let overload_listener = AudioObjectPropertyListener::new( + device_id, + device_address(kAudioDeviceProcessorOverload), + move || { pending_xrun.store(true, Ordering::Relaxed); - }); - - match (alive_listener, rate_listener, overload_listener) { - (Ok(_alive), Ok(_rate), Ok(_overload)) => { + }, + ); + + match ( + alive_listener, + rate_listener, + buffer_size_listener, + overload_listener, + ) { + (Ok(_alive), Ok(_rate), Ok(_buffer), Ok(_overload)) => { let _ = ready_tx.send(Ok(())); // Block until the stream is dropped; listeners are removed on drop. let _ = shutdown_rx.recv(); } - (Err(e), _, _) | (_, Err(e), _) | (_, _, Err(e)) => { + (Err(e), ..) | (_, Err(e), ..) | (_, _, Err(e), _) | (_, _, _, Err(e)) => { let _ = ready_tx.send(Err(e)); } } @@ -198,34 +309,17 @@ impl DisconnectManager { })??; let mut latch = Latch::new(); - let waiter = latch.waiter(); + spawn_delivery_thread( + "cpal-coreaudio-disconnect", + &mut latch, + disconnect_rx, + stream_weak, + move |event, stream| match event { + MonitorEvent::Lost(err) => report_lost(stream, &error_callback, err), + MonitorEvent::LatencyChanged => refresh_latency(stream, &latency_refresh), + }, + )?; - let handle = std::thread::Builder::new() - .name("cpal-coreaudio-disconnect".into()) - .spawn(move || { - // If the Latch is dropped without being released (error path), exit cleanly. - if !waiter.wait() { - return; - } - while let Ok(err) = disconnect_rx.recv() { - if let Some(stream_arc) = stream_weak.upgrade() { - if let Ok(mut stream_inner) = stream_arc.try_lock() { - let _ = stream_inner.pause(); - } - emit_error(&error_callback, err); - } else { - break; - } - } - }) - .map_err(|e| { - Error::with_message( - ErrorKind::ResourceExhausted, - format!("Failed to spawn disconnect thread: {e}"), - ) - })?; - - latch.add_thread(handle.thread().clone()); Ok(DisconnectManager { latch, _shutdown_tx: shutdown_tx, @@ -247,131 +341,157 @@ impl Monitor for DisconnectManager { struct DefaultOutputMonitor { latch: Latch, _shutdown_tx: mpsc::Sender<()>, + // Both are held here rather than in the delivery thread: a sender that thread owned, directly + // or through a listener it kept alive, would hold the event channel open so its loop could + // never end. Dropping these is what lets it exit. + _event_tx: Arc>, + buffer_size_listener: BufferSizeListener, +} + +/// Shutdown handle for the buffer-size listener, re-registered by the delivery thread on reroute. +type BufferSizeListener = Arc>>>; + +/// What the default-output listeners report to that monitor's delivery thread. +enum DefaultOutputEvent { + /// The system default output device changed. + DeviceChanged, + /// The current device's buffer frame size changed, so the cached buffer depth is stale. + LatencyChanged, +} + +/// Registers a buffer-size listener for `device_id`, reporting through `event_tx`. +fn spawn_buffer_size_listener( + device_id: AudioDeviceID, + event_tx: mpsc::Sender, +) -> Result, Error> { + spawn_property_listener_thread( + device_id, + device_address(kAudioDevicePropertyBufferFrameSize), + move || { + let _ = event_tx.send(DefaultOutputEvent::LatencyChanged); + }, + ) +} + +/// Replaces the buffer-size listener, dropping the previous one so its thread exits. +fn set_buffer_size_listener(listener: &BufferSizeListener, next: Option>) { + *listener.lock().unwrap_or_else(|e| e.into_inner()) = next; +} + +impl Drop for DefaultOutputMonitor { + fn drop(&mut self) { + // Release the listener before `_event_tx`, so every sender is gone and the delivery + // thread's loop ends. + set_buffer_size_listener(&self.buffer_size_listener, None); + } } impl DefaultOutputMonitor { fn new( stream_weak: Weak>, error_callback: Arc>, - latency_refresh: Option<(Arc, Scope)>, + latency_refresh: LatencyRefresh, pending_xrun: Arc, ) -> Result { - let (change_tx, change_rx) = mpsc::channel::<()>(); - let shutdown_tx = spawn_property_listener_thread( - kAudioObjectSystemObject as AudioObjectID, - AudioObjectPropertyAddress { - mSelector: kAudioHardwarePropertyDefaultOutputDevice, - mScope: kAudioObjectPropertyScopeGlobal, - mElement: kAudioObjectPropertyElementMain, - }, - move || { - let _ = change_tx.send(()); - }, - )?; + let (change_tx, change_rx) = mpsc::channel::(); + let event_tx = Arc::new(change_tx); + let shutdown_tx = { + let change_tx = mpsc::Sender::clone(&event_tx); + spawn_property_listener_thread( + kAudioObjectSystemObject as AudioObjectID, + device_address(kAudioHardwarePropertyDefaultOutputDevice), + move || { + let _ = change_tx.send(DefaultOutputEvent::DeviceChanged); + }, + )? + }; - // The overload listener targets a specific device, so it must be re-registered against + // These listeners target a specific device, so they must be re-registered against // whatever device is current whenever the default output reroutes. // Held only to shut down the previous listener thread on drop when reassigned below. + let buffer_size_listener: BufferSizeListener = Arc::new(Mutex::new(None)); let mut _overload_shutdown_tx = match default_output_device() { - Some(device) => Some(spawn_property_listener_thread( - device.audio_device_id, - AudioObjectPropertyAddress { - mSelector: kAudioDeviceProcessorOverload, - mScope: kAudioObjectPropertyScopeGlobal, - mElement: kAudioObjectPropertyElementMain, - }, - { - let pending_xrun = pending_xrun.clone(); - move || { - pending_xrun.store(true, Ordering::Relaxed); - } - }, - )?), + Some(device) => { + set_buffer_size_listener( + &buffer_size_listener, + Some(spawn_buffer_size_listener( + device.audio_device_id, + mpsc::Sender::clone(&event_tx), + )?), + ); + Some(spawn_overload_listener( + device.audio_device_id, + pending_xrun.clone(), + )?) + } None => None, }; - let mut latch = Latch::new(); - let waiter = latch.waiter(); + // Weak, so the thread can hand a sender to a new listener without owning one itself. + let event_tx_weak = Arc::downgrade(&event_tx); + let buffer_size_listener_thread = buffer_size_listener.clone(); - let handle = std::thread::Builder::new() - .name("cpal-coreaudio-default-output".into()) - .spawn(move || { - if !waiter.wait() { + let mut latch = Latch::new(); + spawn_delivery_thread( + "cpal-coreaudio-default-output", + &mut latch, + change_rx, + stream_weak, + move |event, stream| { + if matches!(event, DefaultOutputEvent::LatencyChanged) { + // Same device, resized buffer: refresh the depth only. The listeners still + // target the right device, and no route changed to report. + refresh_latency(stream, &latency_refresh); return; } - while let Ok(()) = change_rx.recv() { - let Some(stream) = stream_weak.upgrade() else { - break; - }; - match default_output_device() { - None => { - _overload_shutdown_tx = None; - if let Ok(mut inner) = stream.try_lock() { - let _ = inner.pause(); - } - emit_error( - &error_callback, - Error::with_message( - ErrorKind::DeviceNotAvailable, - "no default output device", - ), - ); - } - Some(device) => { - // DefaultOutput AudioUnit rerouted automatically: recompute and notify - // the buffer depth for the new device. - if let Some((frames, scope)) = &latency_refresh { - if let Ok(inner) = stream.lock() { - let depth = - device::get_device_buffer_frame_size(&inner.audio_unit) - .ok() - .map_or(0, |buffer| { - buffer - + device::get_device_extra_latency_frames( - &inner.audio_unit, - *scope, - ) - }); - frames.store(depth, Ordering::Relaxed); - } - } - _overload_shutdown_tx = spawn_property_listener_thread( - device.audio_device_id, - AudioObjectPropertyAddress { - mSelector: kAudioDeviceProcessorOverload, - mScope: kAudioObjectPropertyScopeGlobal, - mElement: kAudioObjectPropertyElementMain, - }, - { - let pending_xrun = pending_xrun.clone(); - move || { - pending_xrun.store(true, Ordering::Relaxed); - } - }, - ) - .ok(); - emit_error( - &error_callback, - Error::with_message( - ErrorKind::DeviceChanged, - "default output device changed", - ), - ); - } + match default_output_device() { + None => { + _overload_shutdown_tx = None; + set_buffer_size_listener(&buffer_size_listener_thread, None); + report_lost( + stream, + &error_callback, + Error::with_message( + ErrorKind::DeviceNotAvailable, + "no default output device", + ), + ); + } + Some(device) => { + // DefaultOutput AudioUnit rerouted automatically: recompute and notify + // the buffer depth for the new device. + refresh_latency(stream, &latency_refresh); + _overload_shutdown_tx = + spawn_overload_listener(device.audio_device_id, pending_xrun.clone()) + .ok(); + // Skipped once the monitor is dropped: there is nothing left to notify. + set_buffer_size_listener( + &buffer_size_listener_thread, + event_tx_weak.upgrade().and_then(|event_tx| { + spawn_buffer_size_listener( + device.audio_device_id, + mpsc::Sender::clone(&event_tx), + ) + .ok() + }), + ); + emit_error( + &error_callback, + Error::with_message( + ErrorKind::DeviceChanged, + "default output device changed", + ), + ); } } - }) - .map_err(|e| { - Error::with_message( - ErrorKind::ResourceExhausted, - format!("failed to spawn default-output monitor thread: {e}"), - ) - })?; - - latch.add_thread(handle.thread().clone()); + }, + )?; + Ok(DefaultOutputMonitor { latch, _shutdown_tx: shutdown_tx, + _event_tx: event_tx, + buffer_size_listener, }) } } From cf61a77a2ba7a034119ea8563cd0b5aca886912c Mon Sep 17 00:00:00 2001 From: Roderick van Domburg Date: Sun, 20 Sep 2026 13:54:18 +0200 Subject: [PATCH 4/8] fix(coreaudio): end the sample rate wait when the device disconnects --- CHANGELOG.md | 1 + src/host/coreaudio/macos/device.rs | 144 +++++++++++++++++++++-------- 2 files changed, 109 insertions(+), 36 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 623362d52..7aefd7018 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -58,6 +58,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - **CoreAudio**: Fix the device running at a different sample rate from the stream on hardware that reports a continuous rate range. - **CoreAudio**: Fix `supported_configs()` only reporting `F32`, even on hardware that also supports other sample formats. - **CoreAudio**: Fix sample rate changes timing out early when the device reports other rates first. +- **CoreAudio**: Fix sample rate changes waiting on a device that has been disconnected. - **CoreAudio**: Fix timestamps going stale when the device's buffer size changes during a stream. - **iOS**: Fix timestamps and `buffer_size()` being off when the stream sample rate differs from the hardware rate. - **JACK**: Channel enumeration is capped at the physical system port count again. diff --git a/src/host/coreaudio/macos/device.rs b/src/host/coreaudio/macos/device.rs index 26e82c2ec..363d104da 100644 --- a/src/host/coreaudio/macos/device.rs +++ b/src/host/coreaudio/macos/device.rs @@ -14,8 +14,8 @@ use coreaudio::audio_unit::{ AudioUnit, Element, SampleFormat as CoreAudioSampleFormat, Scope, StreamFormat, audio_format::LinearPcmFlags, macos_helpers::{ - RateListener, audio_unit_from_device_id_uninitialized, find_matching_physical_format, - get_device_name, get_supported_physical_stream_formats, set_device_physical_stream_format, + audio_unit_from_device_id_uninitialized, find_matching_physical_format, get_device_name, + get_supported_physical_stream_formats, set_device_physical_stream_format, }, render_callback::{self, data}, }; @@ -27,15 +27,15 @@ use objc2_core_audio::{ AudioObjectID, AudioObjectPropertyAddress, AudioObjectPropertyScope, AudioObjectSetPropertyData, kAudioAggregateDeviceClassID, kAudioDevicePropertyAvailableNominalSampleRates, kAudioDevicePropertyBufferFrameSize, - kAudioDevicePropertyBufferFrameSizeRange, kAudioDevicePropertyDeviceUID, - kAudioDevicePropertyLatency, kAudioDevicePropertyNominalSampleRate, - kAudioDevicePropertySafetyOffset, kAudioDevicePropertyStreamConfiguration, - kAudioDevicePropertyStreamFormat, kAudioDevicePropertyTransportType, - kAudioDeviceTransportTypeAVB, kAudioDeviceTransportTypeAggregate, - kAudioDeviceTransportTypeAirPlay, kAudioDeviceTransportTypeBluetooth, - kAudioDeviceTransportTypeBluetoothLE, kAudioDeviceTransportTypeBuiltIn, - kAudioDeviceTransportTypeDisplayPort, kAudioDeviceTransportTypeFireWire, - kAudioDeviceTransportTypeHDMI, kAudioDeviceTransportTypePCI, + kAudioDevicePropertyBufferFrameSizeRange, kAudioDevicePropertyDeviceIsAlive, + kAudioDevicePropertyDeviceUID, kAudioDevicePropertyLatency, + kAudioDevicePropertyNominalSampleRate, kAudioDevicePropertySafetyOffset, + kAudioDevicePropertyStreamConfiguration, kAudioDevicePropertyStreamFormat, + kAudioDevicePropertyTransportType, kAudioDeviceTransportTypeAVB, + kAudioDeviceTransportTypeAggregate, kAudioDeviceTransportTypeAirPlay, + kAudioDeviceTransportTypeBluetooth, kAudioDeviceTransportTypeBluetoothLE, + kAudioDeviceTransportTypeBuiltIn, kAudioDeviceTransportTypeDisplayPort, + kAudioDeviceTransportTypeFireWire, kAudioDeviceTransportTypeHDMI, kAudioDeviceTransportTypePCI, kAudioDeviceTransportTypeThunderbolt, kAudioDeviceTransportTypeUSB, kAudioDeviceTransportTypeVirtual, kAudioObjectPropertyClass, kAudioObjectPropertyElementMain, kAudioObjectPropertyScopeGlobal, kAudioObjectPropertyScopeInput, @@ -49,7 +49,7 @@ use objc2_core_foundation::{CFRetained, CFString}; pub use super::enumerate::{SupportedInputConfigs, SupportedOutputConfigs}; use super::{ DefaultOutputMonitor, DisconnectManager, Monitor, Stream, asbd_from_config, check_os_status, - host_time_to_stream_instant, + host_time_to_stream_instant, property_listener::AudioObjectPropertyListener, }; use crate::{ BufferSize, CallbackInfo, ChannelCount, Data, DeviceDescription, DeviceDescriptionBuilder, @@ -94,17 +94,12 @@ fn set_physical_format( set_device_physical_stream_format(device_id, asbd).map(|_| asbd) } -/// Set the device's nominal sample rate via `kAudioDevicePropertyNominalSampleRate`. +/// Read the device's current nominal sample rate. /// -/// Unlike [`set_physical_format`], this only changes the device clock rate. The AudioUnit bridges -/// any remaining format difference to the virtual stream format seen by the callback. -fn set_sample_rate( - audio_device_id: AudioObjectID, - target_sample_rate: SampleRate, - timeout: Option, -) -> Result<(), Error> { - // Get the current sample rate. - let mut property_address = AudioObjectPropertyAddress { +/// "Nominal" is CoreAudio's term for the rate the device is configured to run at, as opposed to +/// the actual rate measured from its hardware clock (`kAudioDevicePropertyActualSampleRate`). +fn nominal_sample_rate(audio_device_id: AudioObjectID) -> Result { + let property_address = AudioObjectPropertyAddress { mSelector: kAudioDevicePropertyNominalSampleRate, mScope: kAudioObjectPropertyScopeGlobal, mElement: kAudioObjectPropertyElementMain, @@ -122,6 +117,24 @@ fn set_sample_rate( ) }; coreaudio::Error::from_os_status(status)?; + Ok(sample_rate) +} + +/// Set the device's nominal sample rate via `kAudioDevicePropertyNominalSampleRate`. +/// +/// Unlike [`set_physical_format`], this only changes the device clock rate. The AudioUnit bridges +/// any remaining format difference to the virtual stream format seen by the callback. +fn set_sample_rate( + audio_device_id: AudioObjectID, + target_sample_rate: SampleRate, + timeout: Option, +) -> Result<(), Error> { + let sample_rate = nominal_sample_rate(audio_device_id)?; + let mut property_address = AudioObjectPropertyAddress { + mSelector: kAudioDevicePropertyNominalSampleRate, + mScope: kAudioObjectPropertyScopeGlobal, + mElement: kAudioObjectPropertyElementMain, + }; // If the requested sample rate is different to the device sample rate, update the device. if (sample_rate - target_sample_rate as f64).abs() >= 1.0 { @@ -167,10 +180,32 @@ fn set_sample_rate( )); } - // Register the listener before setting the property so we don't miss the notification. - let (sender, receiver) = channel::(); - let mut listener = RateListener::new(audio_device_id, Some(sender)); - listener.register()?; + // Hook up both listeners before setting the rate, so that neither the new rate nor a + // disconnect can be missed while we wait. + let (sender, receiver) = channel::(); + let alive_address = AudioObjectPropertyAddress { + mSelector: kAudioDevicePropertyDeviceIsAlive, + mScope: kAudioObjectPropertyScopeGlobal, + mElement: kAudioObjectPropertyElementMain, + }; + let alive_sender = sender.clone(); + let _alive_listener = + AudioObjectPropertyListener::new(audio_device_id, alive_address, move || { + let _ = alive_sender.send(RateEvent::Unavailable); + })?; + let rate_address = AudioObjectPropertyAddress { + mSelector: kAudioDevicePropertyNominalSampleRate, + mScope: kAudioObjectPropertyScopeGlobal, + mElement: kAudioObjectPropertyElementMain, + }; + let _rate_listener = + AudioObjectPropertyListener::new(audio_device_id, rate_address, move || { + let event = match nominal_sample_rate(audio_device_id) { + Ok(rate) => RateEvent::Changed(rate), + Err(_) => RateEvent::Unavailable, + }; + let _ = sender.send(event); + })?; // Set the nominal sample rate. property_address.mSelector = kAudioDevicePropertyNominalSampleRate; @@ -190,17 +225,26 @@ fn set_sample_rate( // Wait for the reported_rate to change. This should not take longer than a few ms. wait_for_rate(&receiver, target_sample_rate, timeout)?; - // listener dropped here; its Drop impl calls unregister() automatically. + // listeners are removed when they drop here } Ok(()) } +/// What the device reports while a sample rate change is pending. +enum RateEvent { + /// The nominal sample rate is now this value. + Changed(f64), + /// The device disconnected, or can no longer be queried. + Unavailable, +} + /// Block until the rate listener reports `target_sample_rate`, giving up after `timeout`. /// /// Notifications carrying some other rate can arrive first, so `timeout` bounds the whole wait -/// rather than each individual receive. A `timeout` of `None` waits indefinitely. +/// rather than each individual receive. A `timeout` of `None` waits indefinitely, ending early +/// only if the device disappears. fn wait_for_rate( - receiver: &Receiver, + receiver: &Receiver, target_sample_rate: SampleRate, timeout: Option, ) -> Result<(), Error> { @@ -222,11 +266,17 @@ fn wait_for_rate( }; match received { - Ok(reported_rate) => { + Ok(RateEvent::Changed(reported_rate)) => { if (reported_rate - target_sample_rate as f64).abs() < 1.0 { return Ok(()); } } + Ok(RateEvent::Unavailable) => { + return Err(Error::with_message( + ErrorKind::DeviceNotAvailable, + "Device disconnected while updating sample rate", + )); + } Err(RecvTimeoutError::Timeout) => { return Err(Error::with_message( ErrorKind::DeviceNotAvailable, @@ -1270,7 +1320,7 @@ mod tests { use std::sync::mpsc::channel; use std::time::{Duration, Instant}; - use super::wait_for_rate; + use super::{RateEvent, wait_for_rate}; /// A listener can report rates other than the target before it reports the new one, e.g. a /// device stepping through rates. The whole timeout must remain available across those. @@ -1278,9 +1328,9 @@ mod tests { fn wait_for_rate_honours_the_full_timeout_across_repeated_events() { const TIMEOUT: Duration = Duration::from_millis(50); - let (sender, receiver) = channel::(); + let (sender, receiver) = channel::(); let feeder = std::thread::spawn(move || { - while sender.send(44_100.0).is_ok() { + while sender.send(RateEvent::Changed(44_100.0)).is_ok() { std::thread::sleep(Duration::from_millis(1)); } }); @@ -1300,10 +1350,32 @@ mod tests { #[test] fn wait_for_rate_returns_when_the_target_rate_is_reported() { - let (sender, receiver) = channel::(); - sender.send(44_100.0).unwrap(); - sender.send(48_000.0).unwrap(); + let (sender, receiver) = channel::(); + sender.send(RateEvent::Changed(44_100.0)).unwrap(); + sender.send(RateEvent::Changed(48_000.0)).unwrap(); assert!(wait_for_rate(&receiver, 48_000, Some(Duration::from_secs(5))).is_ok()); } + + /// The wait must end when the device disappears, whether it is unbounded or has a long timeout. + #[test] + fn wait_for_rate_ends_early_on_disconnect() { + for timeout in [None, Some(Duration::from_secs(30))] { + let (sender, receiver) = channel::(); + let disconnect = std::thread::spawn(move || { + std::thread::sleep(Duration::from_millis(20)); + sender.send(RateEvent::Unavailable).unwrap(); + }); + + let start = Instant::now(); + assert!(wait_for_rate(&receiver, 48_000, timeout).is_err()); + let elapsed = start.elapsed(); + disconnect.join().unwrap(); + + assert!( + elapsed < Duration::from_secs(5), + "waited {elapsed:?} with timeout {timeout:?}" + ); + } + } } From 48d4c8dc66a13e4265759c8dfe36c403b5d480d9 Mon Sep 17 00:00:00 2001 From: Roderick van Domburg Date: Sun, 20 Sep 2026 13:54:37 +0200 Subject: [PATCH 5/8] fix(coreaudio): bound format and rate changes by one shared timeout --- CHANGELOG.md | 2 + src/host/coreaudio/macos/device.rs | 122 +++++++++++++++++++++++++---- 2 files changed, 108 insertions(+), 16 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 7aefd7018..ad3276100 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -32,6 +32,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - Renamed the `wasm-beep` and `audioworklet-beep` examples to `webaudio` and `audioworklet`. - **ALSA**: Update `alsa` dependency to 0.12. - **CoreAudio**: `DeviceDescription::interface_type()` now reports the device transport instead of only marking aggregate devices. +- **CoreAudio**: A `None` timeout now waits indefinitely for sample rate and format changes, instead of giving up after 1 or 2 seconds. - **Linux**: `realtime` can now promote threads without requiring `realtime-dbus`. - **PipeWire**: Set `node.rate` property so that `default.clock.allowed-rates` PipeWire config works. @@ -59,6 +60,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - **CoreAudio**: Fix `supported_configs()` only reporting `F32`, even on hardware that also supports other sample formats. - **CoreAudio**: Fix sample rate changes timing out early when the device reports other rates first. - **CoreAudio**: Fix sample rate changes waiting on a device that has been disconnected. +- **CoreAudio**: Fix stream creation taking longer than the timeout while the device changes format. - **CoreAudio**: Fix timestamps going stale when the device's buffer size changes during a stream. - **iOS**: Fix timestamps and `buffer_size()` being off when the stream sample rate differs from the hardware rate. - **JACK**: Channel enumeration is capped at the physical system port count again. diff --git a/src/host/coreaudio/macos/device.rs b/src/host/coreaudio/macos/device.rs index 363d104da..7e265fd52 100644 --- a/src/host/coreaudio/macos/device.rs +++ b/src/host/coreaudio/macos/device.rs @@ -15,7 +15,7 @@ use coreaudio::audio_unit::{ audio_format::LinearPcmFlags, macos_helpers::{ audio_unit_from_device_id_uninitialized, find_matching_physical_format, get_device_name, - get_supported_physical_stream_formats, set_device_physical_stream_format, + get_supported_physical_stream_formats, }, render_callback::{self, data}, }; @@ -39,7 +39,7 @@ use objc2_core_audio::{ kAudioDeviceTransportTypeThunderbolt, kAudioDeviceTransportTypeUSB, kAudioDeviceTransportTypeVirtual, kAudioObjectPropertyClass, kAudioObjectPropertyElementMain, kAudioObjectPropertyScopeGlobal, kAudioObjectPropertyScopeInput, - kAudioObjectPropertyScopeOutput, + kAudioObjectPropertyScopeOutput, kAudioStreamPropertyPhysicalFormat, }; use objc2_core_audio_types::{ AudioBuffer, AudioBufferList, AudioStreamBasicDescription, AudioValueRange, @@ -65,6 +65,81 @@ use crate::{ traits::DeviceTrait, }; +const PHYSICAL_FORMAT_ADDRESS: AudioObjectPropertyAddress = AudioObjectPropertyAddress { + mSelector: kAudioStreamPropertyPhysicalFormat, + mScope: kAudioObjectPropertyScopeGlobal, + mElement: kAudioObjectPropertyElementMain, +}; + +fn physical_format( + device_id: AudioDeviceID, +) -> Result { + let address = PHYSICAL_FORMAT_ADDRESS; + let mut asbd = mem::MaybeUninit::::zeroed(); + let mut data_size = size_of::() as u32; + let status = unsafe { + AudioObjectGetPropertyData( + device_id, + NonNull::from(&address), + 0, + null(), + NonNull::from(&mut data_size), + NonNull::from(&mut asbd).cast(), + ) + }; + coreaudio::Error::from_os_status(status)?; + Ok(unsafe { asbd.assume_init() }) +} + +fn asbds_are_equal( + left: &AudioStreamBasicDescription, + right: &AudioStreamBasicDescription, +) -> bool { + left.mSampleRate as u32 == right.mSampleRate as u32 + && left.mFormatID == right.mFormatID + && left.mFormatFlags == right.mFormatFlags + && left.mBytesPerPacket == right.mBytesPerPacket + && left.mFramesPerPacket == right.mFramesPerPacket + && left.mBytesPerFrame == right.mBytesPerFrame + && left.mChannelsPerFrame == right.mChannelsPerFrame + && left.mBitsPerChannel == right.mBitsPerChannel +} + +/// Set the device's physical stream format and wait until it reports the new one. +/// +/// Gives up at `deadline`, or waits indefinitely if it is `None`. A failed read, such as after +/// the device disconnects, ends the wait. +fn set_physical_stream_format( + device_id: AudioDeviceID, + new_asbd: AudioStreamBasicDescription, + deadline: Option, +) -> Result<(), coreaudio::Error> { + if asbds_are_equal(&physical_format(device_id)?, &new_asbd) { + return Ok(()); + } + + let address = PHYSICAL_FORMAT_ADDRESS; + let status = unsafe { + AudioObjectSetPropertyData( + device_id, + NonNull::from(&address), + 0, + null(), + size_of::() as u32, + NonNull::from(&new_asbd).cast(), + ) + }; + coreaudio::Error::from_os_status(status)?; + + while !asbds_are_equal(&physical_format(device_id)?, &new_asbd) { + if deadline.is_some_and(|deadline| Instant::now() >= deadline) { + return Err(coreaudio::Error::UnsupportedStreamFormat); + } + std::thread::sleep(Duration::from_millis(5)); + } + Ok(()) +} + /// Try to find a matching physical stream format on the device and apply it. /// /// Setting the physical format ensures the hardware runs at the requested bit depth and sample @@ -74,6 +149,7 @@ fn set_physical_format( sample_rate: SampleRate, channels: ChannelCount, sample_format: SampleFormat, + deadline: Option, ) -> Result { let core_format = match sample_format { SampleFormat::I8 => CoreAudioSampleFormat::I8, @@ -91,7 +167,7 @@ fn set_physical_format( }; let asbd = find_matching_physical_format(device_id, stream_format) .ok_or(coreaudio::Error::UnsupportedStreamFormat)?; - set_device_physical_stream_format(device_id, asbd).map(|_| asbd) + set_physical_stream_format(device_id, asbd, deadline).map(|_| asbd) } /// Read the device's current nominal sample rate. @@ -127,7 +203,7 @@ fn nominal_sample_rate(audio_device_id: AudioObjectID) -> Result, + deadline: Option, ) -> Result<(), Error> { let sample_rate = nominal_sample_rate(audio_device_id)?; let mut property_address = AudioObjectPropertyAddress { @@ -224,7 +300,7 @@ fn set_sample_rate( coreaudio::Error::from_os_status(status)?; // Wait for the reported_rate to change. This should not take longer than a few ms. - wait_for_rate(&receiver, target_sample_rate, timeout)?; + wait_for_rate(&receiver, target_sample_rate, deadline)?; // listeners are removed when they drop here } Ok(()) @@ -238,18 +314,16 @@ enum RateEvent { Unavailable, } -/// Block until the rate listener reports `target_sample_rate`, giving up after `timeout`. +/// Block until the rate listener reports `target_sample_rate`, giving up at `deadline`. /// -/// Notifications carrying some other rate can arrive first, so `timeout` bounds the whole wait -/// rather than each individual receive. A `timeout` of `None` waits indefinitely, ending early +/// Notifications carrying some other rate can arrive first, so `deadline` bounds the whole wait +/// rather than each individual receive. A `deadline` of `None` waits indefinitely, ending early /// only if the device disappears. fn wait_for_rate( receiver: &Receiver, target_sample_rate: SampleRate, - timeout: Option, + deadline: Option, ) -> Result<(), Error> { - let deadline = timeout.and_then(|timeout| Instant::now().checked_add(timeout)); - loop { let received = match deadline { Some(deadline) => { @@ -860,6 +934,9 @@ impl Device { { crate::validate_stream_config(&config)?; + // One budget for every wait below, not one per step. + let deadline = timeout.and_then(|timeout| Instant::now().checked_add(timeout)); + // Input is not automatically rerouted, so its buffer depth is constant and its timestamp monotonic. // Set the physical stream format (bit depth + sample rate) on the hardware device. @@ -871,10 +948,11 @@ impl Device { config.sample_rate, config.channels, sample_format, + deadline, ) .is_ok_and(|asbd| (asbd.mSampleRate - config.sample_rate as f64).abs() < 1.0) { - set_sample_rate(self.audio_device_id, config.sample_rate, timeout)?; + set_sample_rate(self.audio_device_id, config.sample_rate, deadline)?; } let mut loopback_aggregate: Option = None; @@ -1009,6 +1087,9 @@ impl Device { { crate::validate_stream_config(&config)?; + // One budget for every wait below, not one per step. + let deadline = timeout.and_then(|timeout| Instant::now().checked_add(timeout)); + // Keep `playback` monotonic: a default output device reroute can lower the device buffer depth, // pulling `playback` backward. let mut data_callback = crate::host::monotonic_output_callback(data_callback); @@ -1022,10 +1103,11 @@ impl Device { config.sample_rate, config.channels, sample_format, + deadline, ) .is_ok_and(|asbd| (asbd.mSampleRate - config.sample_rate as f64).abs() < 1.0) { - set_sample_rate(self.audio_device_id, config.sample_rate, timeout)?; + set_sample_rate(self.audio_device_id, config.sample_rate, deadline)?; } let mode = if self.is_default_output { @@ -1336,7 +1418,7 @@ mod tests { }); let start = Instant::now(); - assert!(wait_for_rate(&receiver, 48_000, Some(TIMEOUT)).is_err()); + assert!(wait_for_rate(&receiver, 48_000, Some(start + TIMEOUT)).is_err()); let elapsed = start.elapsed(); drop(receiver); @@ -1354,7 +1436,14 @@ mod tests { sender.send(RateEvent::Changed(44_100.0)).unwrap(); sender.send(RateEvent::Changed(48_000.0)).unwrap(); - assert!(wait_for_rate(&receiver, 48_000, Some(Duration::from_secs(5))).is_ok()); + assert!( + wait_for_rate( + &receiver, + 48_000, + Some(Instant::now() + Duration::from_secs(5)) + ) + .is_ok() + ); } /// The wait must end when the device disappears, whether it is unbounded or has a long timeout. @@ -1368,7 +1457,8 @@ mod tests { }); let start = Instant::now(); - assert!(wait_for_rate(&receiver, 48_000, timeout).is_err()); + let deadline = timeout.map(|timeout| start + timeout); + assert!(wait_for_rate(&receiver, 48_000, deadline).is_err()); let elapsed = start.elapsed(); disconnect.join().unwrap(); From 81c55d13996e049bcf20182b16e04fa2323975cf Mon Sep 17 00:00:00 2001 From: Roderick van Domburg Date: Sun, 20 Sep 2026 13:54:38 +0200 Subject: [PATCH 6/8] refactor(coreaudio): wait for the physical format change via listeners --- src/host/coreaudio/macos/device.rs | 246 ++++++++++++++++++----------- 1 file changed, 157 insertions(+), 89 deletions(-) diff --git a/src/host/coreaudio/macos/device.rs b/src/host/coreaudio/macos/device.rs index 7e265fd52..29cfb9e15 100644 --- a/src/host/coreaudio/macos/device.rs +++ b/src/host/coreaudio/macos/device.rs @@ -25,13 +25,13 @@ use objc2_audio_toolbox::{ use objc2_core_audio::{ AudioClassID, AudioDeviceID, AudioObjectGetPropertyData, AudioObjectGetPropertyDataSize, AudioObjectID, AudioObjectPropertyAddress, AudioObjectPropertyScope, - AudioObjectSetPropertyData, kAudioAggregateDeviceClassID, + AudioObjectSetPropertyData, AudioStreamID, kAudioAggregateDeviceClassID, kAudioDevicePropertyAvailableNominalSampleRates, kAudioDevicePropertyBufferFrameSize, kAudioDevicePropertyBufferFrameSizeRange, kAudioDevicePropertyDeviceIsAlive, kAudioDevicePropertyDeviceUID, kAudioDevicePropertyLatency, kAudioDevicePropertyNominalSampleRate, kAudioDevicePropertySafetyOffset, kAudioDevicePropertyStreamConfiguration, kAudioDevicePropertyStreamFormat, - kAudioDevicePropertyTransportType, kAudioDeviceTransportTypeAVB, + kAudioDevicePropertyStreams, kAudioDevicePropertyTransportType, kAudioDeviceTransportTypeAVB, kAudioDeviceTransportTypeAggregate, kAudioDeviceTransportTypeAirPlay, kAudioDeviceTransportTypeBluetooth, kAudioDeviceTransportTypeBluetoothLE, kAudioDeviceTransportTypeBuiltIn, kAudioDeviceTransportTypeDisplayPort, @@ -71,15 +71,71 @@ const PHYSICAL_FORMAT_ADDRESS: AudioObjectPropertyAddress = AudioObjectPropertyA mElement: kAudioObjectPropertyElementMain, }; -fn physical_format( +const NOMINAL_SAMPLE_RATE_ADDRESS: AudioObjectPropertyAddress = AudioObjectPropertyAddress { + mSelector: kAudioDevicePropertyNominalSampleRate, + mScope: kAudioObjectPropertyScopeGlobal, + mElement: kAudioObjectPropertyElementMain, +}; + +const DEVICE_IS_ALIVE_ADDRESS: AudioObjectPropertyAddress = AudioObjectPropertyAddress { + mSelector: kAudioDevicePropertyDeviceIsAlive, + mScope: kAudioObjectPropertyScopeGlobal, + mElement: kAudioObjectPropertyElementMain, +}; + +/// Resolve the first `AudioStreamID` a device exposes in `scope` (Input or Output). +/// +/// `kAudioStreamPropertyPhysicalFormat` belongs to the stream object, not the device: reads +/// against the device happen to be forwarded by the HAL, but property-changed notifications are +/// not, so listeners must be registered on the stream object directly. +fn first_stream_id( device_id: AudioDeviceID, + scope: AudioObjectPropertyScope, +) -> Result { + let address = AudioObjectPropertyAddress { + mSelector: kAudioDevicePropertyStreams, + mScope: scope, + mElement: kAudioObjectPropertyElementMain, + }; + let mut data_size = 0u32; + let status = unsafe { + AudioObjectGetPropertyDataSize( + device_id, + NonNull::from(&address), + 0, + null(), + NonNull::from(&mut data_size), + ) + }; + coreaudio::Error::from_os_status(status)?; + let n_streams = data_size as usize / size_of::(); + let mut stream_ids: Vec = vec![0; n_streams]; + let status = unsafe { + AudioObjectGetPropertyData( + device_id, + NonNull::from(&address), + 0, + null(), + NonNull::from(&mut data_size), + NonNull::new(stream_ids.as_mut_ptr()).unwrap().cast(), + ) + }; + coreaudio::Error::from_os_status(status)?; + stream_ids + .into_iter() + .next() + .ok_or(coreaudio::Error::UnsupportedStreamFormat) +} + +fn physical_format( + stream_id: AudioStreamID, ) -> Result { let address = PHYSICAL_FORMAT_ADDRESS; let mut asbd = mem::MaybeUninit::::zeroed(); let mut data_size = size_of::() as u32; let status = unsafe { AudioObjectGetPropertyData( - device_id, + stream_id, NonNull::from(&address), 0, null(), @@ -107,21 +163,27 @@ fn asbds_are_equal( /// Set the device's physical stream format and wait until it reports the new one. /// -/// Gives up at `deadline`, or waits indefinitely if it is `None`. A failed read, such as after -/// the device disconnects, ends the wait. +/// Gives up at `deadline`, or waits indefinitely if it is `None`. The device disconnecting ends +/// the wait. fn set_physical_stream_format( device_id: AudioDeviceID, + scope: AudioObjectPropertyScope, new_asbd: AudioStreamBasicDescription, deadline: Option, -) -> Result<(), coreaudio::Error> { - if asbds_are_equal(&physical_format(device_id)?, &new_asbd) { +) -> Result<(), Error> { + let stream_id = first_stream_id(device_id, scope)?; + if asbds_are_equal(&physical_format(stream_id)?, &new_asbd) { return Ok(()); } + // Listen before setting the format, so the change can't be missed. + let (receiver, _listeners) = + watch_property(device_id, stream_id, PHYSICAL_FORMAT_ADDRESS, physical_format)?; + let address = PHYSICAL_FORMAT_ADDRESS; let status = unsafe { AudioObjectSetPropertyData( - device_id, + stream_id, NonNull::from(&address), 0, null(), @@ -131,13 +193,9 @@ fn set_physical_stream_format( }; coreaudio::Error::from_os_status(status)?; - while !asbds_are_equal(&physical_format(device_id)?, &new_asbd) { - if deadline.is_some_and(|deadline| Instant::now() >= deadline) { - return Err(coreaudio::Error::UnsupportedStreamFormat); - } - std::thread::sleep(Duration::from_millis(5)); - } - Ok(()) + wait_for_property(&receiver, deadline, "physical format", |asbd| { + asbds_are_equal(asbd, &new_asbd) + }) } /// Try to find a matching physical stream format on the device and apply it. @@ -146,18 +204,19 @@ fn set_physical_stream_format( /// rate without unnecessary conversions. fn set_physical_format( device_id: AudioDeviceID, + scope: AudioObjectPropertyScope, sample_rate: SampleRate, channels: ChannelCount, sample_format: SampleFormat, deadline: Option, -) -> Result { +) -> Result { let core_format = match sample_format { SampleFormat::I8 => CoreAudioSampleFormat::I8, SampleFormat::I16 => CoreAudioSampleFormat::I16, SampleFormat::I24 => CoreAudioSampleFormat::I24, SampleFormat::I32 => CoreAudioSampleFormat::I32, SampleFormat::F32 => CoreAudioSampleFormat::F32, - _ => return Err(coreaudio::Error::UnsupportedStreamFormat), + _ => return Err(coreaudio::Error::UnsupportedStreamFormat.into()), }; let stream_format = StreamFormat { sample_rate: sample_rate as f64, @@ -167,7 +226,7 @@ fn set_physical_format( }; let asbd = find_matching_physical_format(device_id, stream_format) .ok_or(coreaudio::Error::UnsupportedStreamFormat)?; - set_physical_stream_format(device_id, asbd, deadline).map(|_| asbd) + set_physical_stream_format(device_id, scope, asbd, deadline).map(|_| asbd) } /// Read the device's current nominal sample rate. @@ -175,11 +234,7 @@ fn set_physical_format( /// "Nominal" is CoreAudio's term for the rate the device is configured to run at, as opposed to /// the actual rate measured from its hardware clock (`kAudioDevicePropertyActualSampleRate`). fn nominal_sample_rate(audio_device_id: AudioObjectID) -> Result { - let property_address = AudioObjectPropertyAddress { - mSelector: kAudioDevicePropertyNominalSampleRate, - mScope: kAudioObjectPropertyScopeGlobal, - mElement: kAudioObjectPropertyElementMain, - }; + let property_address = NOMINAL_SAMPLE_RATE_ADDRESS; let mut sample_rate: f64 = 0.0; let mut data_size = mem::size_of::() as u32; let status = unsafe { @@ -206,11 +261,7 @@ fn set_sample_rate( deadline: Option, ) -> Result<(), Error> { let sample_rate = nominal_sample_rate(audio_device_id)?; - let mut property_address = AudioObjectPropertyAddress { - mSelector: kAudioDevicePropertyNominalSampleRate, - mScope: kAudioObjectPropertyScopeGlobal, - mElement: kAudioObjectPropertyElementMain, - }; + let mut property_address = NOMINAL_SAMPLE_RATE_ADDRESS; // If the requested sample rate is different to the device sample rate, update the device. if (sample_rate - target_sample_rate as f64).abs() >= 1.0 { @@ -256,32 +307,13 @@ fn set_sample_rate( )); } - // Hook up both listeners before setting the rate, so that neither the new rate nor a - // disconnect can be missed while we wait. - let (sender, receiver) = channel::(); - let alive_address = AudioObjectPropertyAddress { - mSelector: kAudioDevicePropertyDeviceIsAlive, - mScope: kAudioObjectPropertyScopeGlobal, - mElement: kAudioObjectPropertyElementMain, - }; - let alive_sender = sender.clone(); - let _alive_listener = - AudioObjectPropertyListener::new(audio_device_id, alive_address, move || { - let _ = alive_sender.send(RateEvent::Unavailable); - })?; - let rate_address = AudioObjectPropertyAddress { - mSelector: kAudioDevicePropertyNominalSampleRate, - mScope: kAudioObjectPropertyScopeGlobal, - mElement: kAudioObjectPropertyElementMain, - }; - let _rate_listener = - AudioObjectPropertyListener::new(audio_device_id, rate_address, move || { - let event = match nominal_sample_rate(audio_device_id) { - Ok(rate) => RateEvent::Changed(rate), - Err(_) => RateEvent::Unavailable, - }; - let _ = sender.send(event); - })?; + // Listen before setting the rate, so that neither the new rate nor a disconnect is missed. + let (receiver, _listeners) = watch_property( + audio_device_id, + audio_device_id, + NOMINAL_SAMPLE_RATE_ADDRESS, + nominal_sample_rate, + )?; // Set the nominal sample rate. property_address.mSelector = kAudioDevicePropertyNominalSampleRate; @@ -306,33 +338,61 @@ fn set_sample_rate( Ok(()) } -/// What the device reports while a sample rate change is pending. -enum RateEvent { - /// The nominal sample rate is now this value. - Changed(f64), +/// What a device reports while one of its properties is changing. +enum PropertyEvent { + /// The property now has this value. + Changed(T), /// The device disconnected, or can no longer be queried. Unavailable, } -/// Block until the rate listener reports `target_sample_rate`, giving up at `deadline`. +/// Report changes to a device property, and the device disconnecting, as [`PropertyEvent`]s. /// -/// Notifications carrying some other rate can arrive first, so `deadline` bounds the whole wait -/// rather than each individual receive. A `deadline` of `None` waits indefinitely, ending early -/// only if the device disappears. -fn wait_for_rate( - receiver: &Receiver, - target_sample_rate: SampleRate, +/// The listeners stay registered for as long as the returned guards are alive. +/// `alive_id` is the device watched for disconnection. `target_id` is the object that actually +/// owns `address`: the device itself, or one of its streams. +fn watch_property( + alive_id: AudioDeviceID, + target_id: AudioObjectID, + address: AudioObjectPropertyAddress, + read: fn(AudioObjectID) -> Result, +) -> Result<(Receiver>, [AudioObjectPropertyListener; 2]), Error> { + let (sender, receiver) = channel(); + let alive_sender = sender.clone(); + let alive = AudioObjectPropertyListener::new(alive_id, DEVICE_IS_ALIVE_ADDRESS, move || { + let _ = alive_sender.send(PropertyEvent::Unavailable); + })?; + let changed = AudioObjectPropertyListener::new(target_id, address, move || { + let event = read(target_id).map_or(PropertyEvent::Unavailable, PropertyEvent::Changed); + let _ = sender.send(event); + })?; + Ok((receiver, [alive, changed])) +} + +/// Block until a reported value satisfies `is_target`, giving up at `deadline`. +/// +/// Other values can be reported first, so `deadline` bounds the whole wait rather than each +/// individual receive. A `deadline` of `None` waits indefinitely, ending early only if the device +/// disappears. +fn wait_for_property( + receiver: &Receiver>, deadline: Option, + what: &str, + mut is_target: impl FnMut(&T) -> bool, ) -> Result<(), Error> { + let timed_out = || { + Error::with_message( + ErrorKind::DeviceNotAvailable, + format!("Timed out waiting for the {what} to change"), + ) + }; + loop { let received = match deadline { Some(deadline) => { let remaining = deadline.saturating_duration_since(Instant::now()); if remaining.is_zero() { - return Err(Error::with_message( - ErrorKind::DeviceNotAvailable, - "Sample rate update timed out", - )); + return Err(timed_out()); } receiver.recv_timeout(remaining) } @@ -340,33 +400,39 @@ fn wait_for_rate( }; match received { - Ok(RateEvent::Changed(reported_rate)) => { - if (reported_rate - target_sample_rate as f64).abs() < 1.0 { + Ok(PropertyEvent::Changed(value)) => { + if is_target(&value) { return Ok(()); } } - Ok(RateEvent::Unavailable) => { - return Err(Error::with_message( - ErrorKind::DeviceNotAvailable, - "Device disconnected while updating sample rate", - )); - } - Err(RecvTimeoutError::Timeout) => { + Ok(PropertyEvent::Unavailable) => { return Err(Error::with_message( ErrorKind::DeviceNotAvailable, - "Sample rate update timed out", + format!("Device disconnected while updating the {what}"), )); } + Err(RecvTimeoutError::Timeout) => return Err(timed_out()), Err(RecvTimeoutError::Disconnected) => { return Err(Error::with_message( ErrorKind::StreamInvalidated, - "Sample rate listener disconnected unexpectedly", + format!("Listener for the {what} disconnected unexpectedly"), )); } } } } +/// Block until the device reports `target_sample_rate`, giving up at `deadline`. +fn wait_for_rate( + receiver: &Receiver>, + target_sample_rate: SampleRate, + deadline: Option, +) -> Result<(), Error> { + wait_for_property(receiver, deadline, "sample rate", |rate| { + (rate - target_sample_rate as f64).abs() < 1.0 + }) +} + #[derive(Clone, Copy)] enum AudioUnitMode { /// HAL Output AudioUnit with input enabled, pinned to a specific device. @@ -945,6 +1011,7 @@ impl Device { // if the closest match found doesn't actually run at the requested rate. if !set_physical_format( self.audio_device_id, + kAudioObjectPropertyScopeInput, config.sample_rate, config.channels, sample_format, @@ -1100,6 +1167,7 @@ impl Device { // closest match found doesn't actually run at the requested rate. if !set_physical_format( self.audio_device_id, + kAudioObjectPropertyScopeOutput, config.sample_rate, config.channels, sample_format, @@ -1402,7 +1470,7 @@ mod tests { use std::sync::mpsc::channel; use std::time::{Duration, Instant}; - use super::{RateEvent, wait_for_rate}; + use super::{PropertyEvent, wait_for_rate}; /// A listener can report rates other than the target before it reports the new one, e.g. a /// device stepping through rates. The whole timeout must remain available across those. @@ -1410,9 +1478,9 @@ mod tests { fn wait_for_rate_honours_the_full_timeout_across_repeated_events() { const TIMEOUT: Duration = Duration::from_millis(50); - let (sender, receiver) = channel::(); + let (sender, receiver) = channel::>(); let feeder = std::thread::spawn(move || { - while sender.send(RateEvent::Changed(44_100.0)).is_ok() { + while sender.send(PropertyEvent::Changed(44_100.0)).is_ok() { std::thread::sleep(Duration::from_millis(1)); } }); @@ -1432,9 +1500,9 @@ mod tests { #[test] fn wait_for_rate_returns_when_the_target_rate_is_reported() { - let (sender, receiver) = channel::(); - sender.send(RateEvent::Changed(44_100.0)).unwrap(); - sender.send(RateEvent::Changed(48_000.0)).unwrap(); + let (sender, receiver) = channel::>(); + sender.send(PropertyEvent::Changed(44_100.0)).unwrap(); + sender.send(PropertyEvent::Changed(48_000.0)).unwrap(); assert!( wait_for_rate( @@ -1450,10 +1518,10 @@ mod tests { #[test] fn wait_for_rate_ends_early_on_disconnect() { for timeout in [None, Some(Duration::from_secs(30))] { - let (sender, receiver) = channel::(); + let (sender, receiver) = channel::>(); let disconnect = std::thread::spawn(move || { std::thread::sleep(Duration::from_millis(20)); - sender.send(RateEvent::Unavailable).unwrap(); + sender.send(PropertyEvent::Unavailable).unwrap(); }); let start = Instant::now(); From 8a01d7f47736f50987b407bb4eefc44327b3ef86 Mon Sep 17 00:00:00 2001 From: Roderick van Domburg Date: Sun, 20 Sep 2026 17:21:07 +0200 Subject: [PATCH 7/8] refactor(coreaudio): extract get_property/get_property_array helpers --- src/host/coreaudio/macos/device.rs | 216 ++++----------------------- src/host/coreaudio/macos/mod.rs | 1 + src/host/coreaudio/macos/property.rs | 75 ++++++++++ 3 files changed, 109 insertions(+), 183 deletions(-) create mode 100644 src/host/coreaudio/macos/property.rs diff --git a/src/host/coreaudio/macos/device.rs b/src/host/coreaudio/macos/device.rs index 29cfb9e15..d664ac03a 100644 --- a/src/host/coreaudio/macos/device.rs +++ b/src/host/coreaudio/macos/device.rs @@ -49,7 +49,9 @@ use objc2_core_foundation::{CFRetained, CFString}; pub use super::enumerate::{SupportedInputConfigs, SupportedOutputConfigs}; use super::{ DefaultOutputMonitor, DisconnectManager, Monitor, Stream, asbd_from_config, check_os_status, - host_time_to_stream_instant, property_listener::AudioObjectPropertyListener, + host_time_to_stream_instant, + property::{get_property, get_property_array}, + property_listener::AudioObjectPropertyListener, }; use crate::{ BufferSize, CallbackInfo, ChannelCount, Data, DeviceDescription, DeviceDescriptionBuilder, @@ -97,30 +99,8 @@ fn first_stream_id( mScope: scope, mElement: kAudioObjectPropertyElementMain, }; - let mut data_size = 0u32; - let status = unsafe { - AudioObjectGetPropertyDataSize( - device_id, - NonNull::from(&address), - 0, - null(), - NonNull::from(&mut data_size), - ) - }; - coreaudio::Error::from_os_status(status)?; - let n_streams = data_size as usize / size_of::(); - let mut stream_ids: Vec = vec![0; n_streams]; - let status = unsafe { - AudioObjectGetPropertyData( - device_id, - NonNull::from(&address), - 0, - null(), - NonNull::from(&mut data_size), - NonNull::new(stream_ids.as_mut_ptr()).unwrap().cast(), - ) - }; - coreaudio::Error::from_os_status(status)?; + // SAFETY: kAudioDevicePropertyStreams is documented to return an array of AudioStreamID. + let stream_ids: Vec = unsafe { get_property_array(device_id, address) }?; stream_ids .into_iter() .next() @@ -130,21 +110,8 @@ fn first_stream_id( fn physical_format( stream_id: AudioStreamID, ) -> Result { - let address = PHYSICAL_FORMAT_ADDRESS; - let mut asbd = mem::MaybeUninit::::zeroed(); - let mut data_size = size_of::() as u32; - let status = unsafe { - AudioObjectGetPropertyData( - stream_id, - NonNull::from(&address), - 0, - null(), - NonNull::from(&mut data_size), - NonNull::from(&mut asbd).cast(), - ) - }; - coreaudio::Error::from_os_status(status)?; - Ok(unsafe { asbd.assume_init() }) + // SAFETY: kAudioStreamPropertyPhysicalFormat is documented to return an AudioStreamBasicDescription. + unsafe { get_property(stream_id, PHYSICAL_FORMAT_ADDRESS) } } fn asbds_are_equal( @@ -234,21 +201,8 @@ fn set_physical_format( /// "Nominal" is CoreAudio's term for the rate the device is configured to run at, as opposed to /// the actual rate measured from its hardware clock (`kAudioDevicePropertyActualSampleRate`). fn nominal_sample_rate(audio_device_id: AudioObjectID) -> Result { - let property_address = NOMINAL_SAMPLE_RATE_ADDRESS; - let mut sample_rate: f64 = 0.0; - let mut data_size = mem::size_of::() as u32; - let status = unsafe { - AudioObjectGetPropertyData( - audio_device_id, - NonNull::from(&property_address), - 0, - null(), - NonNull::from(&mut data_size), - NonNull::from(&mut sample_rate).cast(), - ) - }; - coreaudio::Error::from_os_status(status)?; - Ok(sample_rate) + // SAFETY: kAudioDevicePropertyNominalSampleRate is documented to return an f64. + unsafe { get_property(audio_device_id, NOMINAL_SAMPLE_RATE_ADDRESS) } } /// Set the device's nominal sample rate via `kAudioDevicePropertyNominalSampleRate`. @@ -267,33 +221,10 @@ fn set_sample_rate( if (sample_rate - target_sample_rate as f64).abs() >= 1.0 { // Get available sample rate ranges. property_address.mSelector = kAudioDevicePropertyAvailableNominalSampleRates; - let mut data_size = 0u32; - let status = unsafe { - AudioObjectGetPropertyDataSize( - audio_device_id, - NonNull::from(&property_address), - 0, - null(), - NonNull::from(&mut data_size), - ) - }; - coreaudio::Error::from_os_status(status)?; - let n_ranges = data_size as usize / mem::size_of::(); - let mut ranges: Vec = Vec::with_capacity(n_ranges); - let status = unsafe { - AudioObjectGetPropertyData( - audio_device_id, - NonNull::from(&property_address), - 0, - null(), - NonNull::from(&mut data_size), - NonNull::new(ranges.as_mut_ptr()).unwrap().cast(), - ) - }; - coreaudio::Error::from_os_status(status)?; - unsafe { - ranges.set_len(n_ranges); - } + // SAFETY: kAudioDevicePropertyAvailableNominalSampleRates is documented to return an + // array of AudioValueRange. + let ranges: Vec = + unsafe { get_property_array(audio_device_id, property_address) }?; // Now that we have the available ranges, pick the one matching the desired rate. let sample_rate = target_sample_rate; @@ -479,22 +410,8 @@ fn get_io_buffer_frame_size_range(device_id: AudioDeviceID) -> Result() as u32; - let status = unsafe { - AudioObjectGetPropertyData( - device_id, - NonNull::from(&property_address), - 0, - null(), - NonNull::from(&mut data_size), - NonNull::from(&mut range).cast(), - ) - }; - check_os_status(status)?; + // SAFETY: kAudioDevicePropertyBufferFrameSizeRange is documented to return an AudioValueRange. + let range: AudioValueRange = unsafe { get_property(device_id, property_address) }?; Ok(SupportedBufferSize::Range { min: range.mMinimum as u32, max: range.mMaximum as u32, @@ -601,24 +518,9 @@ impl Device { mElement: kAudioObjectPropertyElementMain, }; - let mut class_id: AudioClassID = 0; - let data_size = size_of::() as u32; - - // SAFETY: AudioObjectGetPropertyData is documented to write an AudioClassID - // for kAudioObjectPropertyClass. We check the status before using the value. - let status = unsafe { - AudioObjectGetPropertyData( - self.audio_device_id, - NonNull::from(&property_address), - 0, - null(), - NonNull::from(&data_size), - NonNull::from(&mut class_id).cast(), - ) - }; - - // If successful, check if it's an aggregate device - status == 0 && class_id == kAudioAggregateDeviceClassID + // SAFETY: kAudioObjectPropertyClass is documented to return an AudioClassID. + unsafe { get_property::(self.audio_device_id, property_address) } + .is_ok_and(|class_id| class_id == kAudioAggregateDeviceClassID) } /// `None` when the property is unavailable or names a transport with no @@ -630,24 +532,11 @@ impl Device { mElement: kAudioObjectPropertyElementMain, }; - let mut transport: u32 = 0; - let mut data_size = size_of::() as u32; - - // SAFETY: AudioObjectGetPropertyData writes a UInt32 for - // kAudioDevicePropertyTransportType. The status is checked before use. - let status = unsafe { - AudioObjectGetPropertyData( - self.audio_device_id, - NonNull::from(&property_address), - 0, - null(), - NonNull::from(&mut data_size), - NonNull::from(&mut transport).cast(), - ) - }; - if status != 0 { + // SAFETY: kAudioDevicePropertyTransportType is documented to return a UInt32. + let Ok(transport) = (unsafe { get_property::(self.audio_device_id, property_address) }) + else { return None; - } + }; #[allow(non_upper_case_globals)] match transport { @@ -706,22 +595,9 @@ impl Device { }; // CFString is returned under the create rule, so take ownership of the +1 reference. - let mut uid: *mut CFString = std::ptr::null_mut(); - let mut data_size = size_of::<*mut CFString>() as u32; - - // SAFETY: AudioObjectGetPropertyData is documented to write a CFString pointer - // for kAudioDevicePropertyDeviceUID. We check the status code before use. - let status = unsafe { - AudioObjectGetPropertyData( - self.audio_device_id, - NonNull::from(&property_address), - 0, - null(), - NonNull::from(&mut data_size), - NonNull::from(&mut uid).cast(), - ) - }; - check_os_status(status)?; + // SAFETY: kAudioDevicePropertyDeviceUID is documented to return a CFString pointer. + let uid: *mut CFString = + unsafe { get_property(self.audio_device_id, property_address) }?; // SAFETY: Status was successful, meaning the API call succeeded. // We now check if the returned uid is non-null before use. @@ -796,29 +672,10 @@ impl Device { // Get available sample rate ranges. property_address.mSelector = kAudioDevicePropertyAvailableNominalSampleRates; - let mut data_size = 0u32; - let status = AudioObjectGetPropertyDataSize( - self.audio_device_id, - NonNull::from(&property_address), - 0, - null(), - NonNull::from(&mut data_size), - ); - check_os_status(status)?; - - let n_ranges = data_size as usize / mem::size_of::(); - let mut ranges: Vec = Vec::with_capacity(n_ranges); - let status = AudioObjectGetPropertyData( - self.audio_device_id, - NonNull::from(&property_address), - 0, - null(), - NonNull::from(&mut data_size), - NonNull::new(ranges.as_mut_ptr()).unwrap().cast(), - ); - check_os_status(status)?; - - ranges.set_len(n_ranges); + // SAFETY: kAudioDevicePropertyAvailableNominalSampleRates is documented to return an + // array of AudioValueRange. + let ranges: Vec = + get_property_array(self.audio_device_id, property_address)?; #[allow(non_upper_case_globals)] match scope { @@ -903,17 +760,10 @@ impl Device { }; unsafe { - let mut asbd: AudioStreamBasicDescription = mem::zeroed(); - let mut data_size = mem::size_of::() as u32; - let status = AudioObjectGetPropertyData( - self.audio_device_id, - NonNull::from(&property_address), - 0, - null(), - NonNull::from(&mut data_size), - NonNull::from(&mut asbd).cast(), - ); - check_os_status(status)?; + // SAFETY: kAudioDevicePropertyStreamFormat is documented to return an + // AudioStreamBasicDescription. + let asbd: AudioStreamBasicDescription = + get_property(self.audio_device_id, property_address)?; let sample_format = { let audio_format = coreaudio::audio_unit::AudioFormat::from_format_and_flag( diff --git a/src/host/coreaudio/macos/mod.rs b/src/host/coreaudio/macos/mod.rs index d9fb3aee4..c3b06c018 100644 --- a/src/host/coreaudio/macos/mod.rs +++ b/src/host/coreaudio/macos/mod.rs @@ -31,6 +31,7 @@ use crate::{ mod device; pub mod enumerate; mod loopback; +mod property; mod property_listener; pub use device::Device; diff --git a/src/host/coreaudio/macos/property.rs b/src/host/coreaudio/macos/property.rs new file mode 100644 index 000000000..f96a9d32f --- /dev/null +++ b/src/host/coreaudio/macos/property.rs @@ -0,0 +1,75 @@ +//! Helper code for reading CoreAudio object properties. +use std::{ + mem, + ptr::{NonNull, null}, +}; + +use objc2_core_audio::{ + AudioObjectGetPropertyData, AudioObjectGetPropertyDataSize, AudioObjectID, + AudioObjectPropertyAddress, +}; + +/// Read a single fixed-size property value. +/// +/// # Safety +/// +/// `T` must match the binary layout CoreAudio writes for `address`. +pub unsafe fn get_property( + object_id: AudioObjectID, + address: AudioObjectPropertyAddress, +) -> Result { + let mut value = mem::MaybeUninit::::zeroed(); + let mut data_size = mem::size_of::() as u32; + let status = unsafe { + AudioObjectGetPropertyData( + object_id, + NonNull::from(&address), + 0, + null(), + NonNull::from(&mut data_size), + NonNull::from(&mut value).cast(), + ) + }; + coreaudio::Error::from_os_status(status)?; + Ok(unsafe { value.assume_init() }) +} + +/// Read a variable-length array property. +/// +/// # Safety +/// +/// `T` must match the binary layout CoreAudio writes for each element of `address`. +pub unsafe fn get_property_array( + object_id: AudioObjectID, + address: AudioObjectPropertyAddress, +) -> Result, coreaudio::Error> { + let mut data_size = 0u32; + let status = unsafe { + AudioObjectGetPropertyDataSize( + object_id, + NonNull::from(&address), + 0, + null(), + NonNull::from(&mut data_size), + ) + }; + coreaudio::Error::from_os_status(status)?; + + let n = data_size as usize / mem::size_of::(); + let mut values: Vec = Vec::with_capacity(n); + let status = unsafe { + AudioObjectGetPropertyData( + object_id, + NonNull::from(&address), + 0, + null(), + NonNull::from(&mut data_size), + NonNull::new(values.as_mut_ptr()).unwrap().cast(), + ) + }; + coreaudio::Error::from_os_status(status)?; + // SAFETY: the size query above reported room for exactly `n` elements, and the status check + // confirms CoreAudio filled the buffer it was given. + unsafe { values.set_len(n) }; + Ok(values) +} From fb969a33819d42397704bb2a394f65b4bcb72326 Mon Sep 17 00:00:00 2001 From: Roderick van Domburg Date: Sun, 20 Sep 2026 17:30:11 +0200 Subject: [PATCH 8/8] refactor(coreaudio): extract format/rate negotiation into format.rs --- src/host/coreaudio/macos/device.rs | 397 +--------------------------- src/host/coreaudio/macos/format.rs | 401 +++++++++++++++++++++++++++++ src/host/coreaudio/macos/mod.rs | 1 + src/host/mod.rs | 2 +- 4 files changed, 414 insertions(+), 387 deletions(-) create mode 100644 src/host/coreaudio/macos/format.rs diff --git a/src/host/coreaudio/macos/device.rs b/src/host/coreaudio/macos/device.rs index d664ac03a..7e489567a 100644 --- a/src/host/coreaudio/macos/device.rs +++ b/src/host/coreaudio/macos/device.rs @@ -1,20 +1,17 @@ use std::{ fmt, - mem::{self, size_of}, ptr::{NonNull, null}, sync::{ Arc, Mutex, atomic::{AtomicBool, AtomicUsize, Ordering}, - mpsc::{Receiver, RecvTimeoutError, channel}, }, time::{Duration, Instant}, }; use coreaudio::audio_unit::{ - AudioUnit, Element, SampleFormat as CoreAudioSampleFormat, Scope, StreamFormat, - audio_format::LinearPcmFlags, + AudioUnit, Element, SampleFormat as CoreAudioSampleFormat, Scope, macos_helpers::{ - audio_unit_from_device_id_uninitialized, find_matching_physical_format, get_device_name, + audio_unit_from_device_id_uninitialized, get_device_name, get_supported_physical_stream_formats, }, render_callback::{self, data}, @@ -24,14 +21,12 @@ use objc2_audio_toolbox::{ }; use objc2_core_audio::{ AudioClassID, AudioDeviceID, AudioObjectGetPropertyData, AudioObjectGetPropertyDataSize, - AudioObjectID, AudioObjectPropertyAddress, AudioObjectPropertyScope, - AudioObjectSetPropertyData, AudioStreamID, kAudioAggregateDeviceClassID, + AudioObjectPropertyAddress, AudioObjectPropertyScope, kAudioAggregateDeviceClassID, kAudioDevicePropertyAvailableNominalSampleRates, kAudioDevicePropertyBufferFrameSize, - kAudioDevicePropertyBufferFrameSizeRange, kAudioDevicePropertyDeviceIsAlive, - kAudioDevicePropertyDeviceUID, kAudioDevicePropertyLatency, - kAudioDevicePropertyNominalSampleRate, kAudioDevicePropertySafetyOffset, + kAudioDevicePropertyBufferFrameSizeRange, kAudioDevicePropertyDeviceUID, + kAudioDevicePropertyLatency, kAudioDevicePropertySafetyOffset, kAudioDevicePropertyStreamConfiguration, kAudioDevicePropertyStreamFormat, - kAudioDevicePropertyStreams, kAudioDevicePropertyTransportType, kAudioDeviceTransportTypeAVB, + kAudioDevicePropertyTransportType, kAudioDeviceTransportTypeAVB, kAudioDeviceTransportTypeAggregate, kAudioDeviceTransportTypeAirPlay, kAudioDeviceTransportTypeBluetooth, kAudioDeviceTransportTypeBluetoothLE, kAudioDeviceTransportTypeBuiltIn, kAudioDeviceTransportTypeDisplayPort, @@ -39,7 +34,7 @@ use objc2_core_audio::{ kAudioDeviceTransportTypeThunderbolt, kAudioDeviceTransportTypeUSB, kAudioDeviceTransportTypeVirtual, kAudioObjectPropertyClass, kAudioObjectPropertyElementMain, kAudioObjectPropertyScopeGlobal, kAudioObjectPropertyScopeInput, - kAudioObjectPropertyScopeOutput, kAudioStreamPropertyPhysicalFormat, + kAudioObjectPropertyScopeOutput, }; use objc2_core_audio_types::{ AudioBuffer, AudioBufferList, AudioStreamBasicDescription, AudioValueRange, @@ -49,9 +44,9 @@ use objc2_core_foundation::{CFRetained, CFString}; pub use super::enumerate::{SupportedInputConfigs, SupportedOutputConfigs}; use super::{ DefaultOutputMonitor, DisconnectManager, Monitor, Stream, asbd_from_config, check_os_status, + format::{set_physical_format, set_sample_rate}, host_time_to_stream_instant, property::{get_property, get_property_array}, - property_listener::AudioObjectPropertyListener, }; use crate::{ BufferSize, CallbackInfo, ChannelCount, Data, DeviceDescription, DeviceDescriptionBuilder, @@ -67,303 +62,6 @@ use crate::{ traits::DeviceTrait, }; -const PHYSICAL_FORMAT_ADDRESS: AudioObjectPropertyAddress = AudioObjectPropertyAddress { - mSelector: kAudioStreamPropertyPhysicalFormat, - mScope: kAudioObjectPropertyScopeGlobal, - mElement: kAudioObjectPropertyElementMain, -}; - -const NOMINAL_SAMPLE_RATE_ADDRESS: AudioObjectPropertyAddress = AudioObjectPropertyAddress { - mSelector: kAudioDevicePropertyNominalSampleRate, - mScope: kAudioObjectPropertyScopeGlobal, - mElement: kAudioObjectPropertyElementMain, -}; - -const DEVICE_IS_ALIVE_ADDRESS: AudioObjectPropertyAddress = AudioObjectPropertyAddress { - mSelector: kAudioDevicePropertyDeviceIsAlive, - mScope: kAudioObjectPropertyScopeGlobal, - mElement: kAudioObjectPropertyElementMain, -}; - -/// Resolve the first `AudioStreamID` a device exposes in `scope` (Input or Output). -/// -/// `kAudioStreamPropertyPhysicalFormat` belongs to the stream object, not the device: reads -/// against the device happen to be forwarded by the HAL, but property-changed notifications are -/// not, so listeners must be registered on the stream object directly. -fn first_stream_id( - device_id: AudioDeviceID, - scope: AudioObjectPropertyScope, -) -> Result { - let address = AudioObjectPropertyAddress { - mSelector: kAudioDevicePropertyStreams, - mScope: scope, - mElement: kAudioObjectPropertyElementMain, - }; - // SAFETY: kAudioDevicePropertyStreams is documented to return an array of AudioStreamID. - let stream_ids: Vec = unsafe { get_property_array(device_id, address) }?; - stream_ids - .into_iter() - .next() - .ok_or(coreaudio::Error::UnsupportedStreamFormat) -} - -fn physical_format( - stream_id: AudioStreamID, -) -> Result { - // SAFETY: kAudioStreamPropertyPhysicalFormat is documented to return an AudioStreamBasicDescription. - unsafe { get_property(stream_id, PHYSICAL_FORMAT_ADDRESS) } -} - -fn asbds_are_equal( - left: &AudioStreamBasicDescription, - right: &AudioStreamBasicDescription, -) -> bool { - left.mSampleRate as u32 == right.mSampleRate as u32 - && left.mFormatID == right.mFormatID - && left.mFormatFlags == right.mFormatFlags - && left.mBytesPerPacket == right.mBytesPerPacket - && left.mFramesPerPacket == right.mFramesPerPacket - && left.mBytesPerFrame == right.mBytesPerFrame - && left.mChannelsPerFrame == right.mChannelsPerFrame - && left.mBitsPerChannel == right.mBitsPerChannel -} - -/// Set the device's physical stream format and wait until it reports the new one. -/// -/// Gives up at `deadline`, or waits indefinitely if it is `None`. The device disconnecting ends -/// the wait. -fn set_physical_stream_format( - device_id: AudioDeviceID, - scope: AudioObjectPropertyScope, - new_asbd: AudioStreamBasicDescription, - deadline: Option, -) -> Result<(), Error> { - let stream_id = first_stream_id(device_id, scope)?; - if asbds_are_equal(&physical_format(stream_id)?, &new_asbd) { - return Ok(()); - } - - // Listen before setting the format, so the change can't be missed. - let (receiver, _listeners) = - watch_property(device_id, stream_id, PHYSICAL_FORMAT_ADDRESS, physical_format)?; - - let address = PHYSICAL_FORMAT_ADDRESS; - let status = unsafe { - AudioObjectSetPropertyData( - stream_id, - NonNull::from(&address), - 0, - null(), - size_of::() as u32, - NonNull::from(&new_asbd).cast(), - ) - }; - coreaudio::Error::from_os_status(status)?; - - wait_for_property(&receiver, deadline, "physical format", |asbd| { - asbds_are_equal(asbd, &new_asbd) - }) -} - -/// Try to find a matching physical stream format on the device and apply it. -/// -/// Setting the physical format ensures the hardware runs at the requested bit depth and sample -/// rate without unnecessary conversions. -fn set_physical_format( - device_id: AudioDeviceID, - scope: AudioObjectPropertyScope, - sample_rate: SampleRate, - channels: ChannelCount, - sample_format: SampleFormat, - deadline: Option, -) -> Result { - let core_format = match sample_format { - SampleFormat::I8 => CoreAudioSampleFormat::I8, - SampleFormat::I16 => CoreAudioSampleFormat::I16, - SampleFormat::I24 => CoreAudioSampleFormat::I24, - SampleFormat::I32 => CoreAudioSampleFormat::I32, - SampleFormat::F32 => CoreAudioSampleFormat::F32, - _ => return Err(coreaudio::Error::UnsupportedStreamFormat.into()), - }; - let stream_format = StreamFormat { - sample_rate: sample_rate as f64, - sample_format: core_format, - flags: LinearPcmFlags::empty(), - channels: channels as u32, - }; - let asbd = find_matching_physical_format(device_id, stream_format) - .ok_or(coreaudio::Error::UnsupportedStreamFormat)?; - set_physical_stream_format(device_id, scope, asbd, deadline).map(|_| asbd) -} - -/// Read the device's current nominal sample rate. -/// -/// "Nominal" is CoreAudio's term for the rate the device is configured to run at, as opposed to -/// the actual rate measured from its hardware clock (`kAudioDevicePropertyActualSampleRate`). -fn nominal_sample_rate(audio_device_id: AudioObjectID) -> Result { - // SAFETY: kAudioDevicePropertyNominalSampleRate is documented to return an f64. - unsafe { get_property(audio_device_id, NOMINAL_SAMPLE_RATE_ADDRESS) } -} - -/// Set the device's nominal sample rate via `kAudioDevicePropertyNominalSampleRate`. -/// -/// Unlike [`set_physical_format`], this only changes the device clock rate. The AudioUnit bridges -/// any remaining format difference to the virtual stream format seen by the callback. -fn set_sample_rate( - audio_device_id: AudioObjectID, - target_sample_rate: SampleRate, - deadline: Option, -) -> Result<(), Error> { - let sample_rate = nominal_sample_rate(audio_device_id)?; - let mut property_address = NOMINAL_SAMPLE_RATE_ADDRESS; - - // If the requested sample rate is different to the device sample rate, update the device. - if (sample_rate - target_sample_rate as f64).abs() >= 1.0 { - // Get available sample rate ranges. - property_address.mSelector = kAudioDevicePropertyAvailableNominalSampleRates; - // SAFETY: kAudioDevicePropertyAvailableNominalSampleRates is documented to return an - // array of AudioValueRange. - let ranges: Vec = - unsafe { get_property_array(audio_device_id, property_address) }?; - - // Now that we have the available ranges, pick the one matching the desired rate. - let sample_rate = target_sample_rate; - if !ranges - .iter() - .any(|r| sample_rate as f64 >= r.mMinimum && sample_rate as f64 <= r.mMaximum) - { - return Err(Error::with_message( - ErrorKind::UnsupportedConfig, - format!("Sample rate {sample_rate} Hz is not supported"), - )); - } - - // Listen before setting the rate, so that neither the new rate nor a disconnect is missed. - let (receiver, _listeners) = watch_property( - audio_device_id, - audio_device_id, - NOMINAL_SAMPLE_RATE_ADDRESS, - nominal_sample_rate, - )?; - - // Set the nominal sample rate. - property_address.mSelector = kAudioDevicePropertyNominalSampleRate; - let rate = sample_rate as f64; - let data_size = mem::size_of::() as u32; - let status = unsafe { - AudioObjectSetPropertyData( - audio_device_id, - NonNull::from(&property_address), - 0, - null(), - data_size, - NonNull::from(&rate).cast(), - ) - }; - coreaudio::Error::from_os_status(status)?; - - // Wait for the reported_rate to change. This should not take longer than a few ms. - wait_for_rate(&receiver, target_sample_rate, deadline)?; - // listeners are removed when they drop here - } - Ok(()) -} - -/// What a device reports while one of its properties is changing. -enum PropertyEvent { - /// The property now has this value. - Changed(T), - /// The device disconnected, or can no longer be queried. - Unavailable, -} - -/// Report changes to a device property, and the device disconnecting, as [`PropertyEvent`]s. -/// -/// The listeners stay registered for as long as the returned guards are alive. -/// `alive_id` is the device watched for disconnection. `target_id` is the object that actually -/// owns `address`: the device itself, or one of its streams. -fn watch_property( - alive_id: AudioDeviceID, - target_id: AudioObjectID, - address: AudioObjectPropertyAddress, - read: fn(AudioObjectID) -> Result, -) -> Result<(Receiver>, [AudioObjectPropertyListener; 2]), Error> { - let (sender, receiver) = channel(); - let alive_sender = sender.clone(); - let alive = AudioObjectPropertyListener::new(alive_id, DEVICE_IS_ALIVE_ADDRESS, move || { - let _ = alive_sender.send(PropertyEvent::Unavailable); - })?; - let changed = AudioObjectPropertyListener::new(target_id, address, move || { - let event = read(target_id).map_or(PropertyEvent::Unavailable, PropertyEvent::Changed); - let _ = sender.send(event); - })?; - Ok((receiver, [alive, changed])) -} - -/// Block until a reported value satisfies `is_target`, giving up at `deadline`. -/// -/// Other values can be reported first, so `deadline` bounds the whole wait rather than each -/// individual receive. A `deadline` of `None` waits indefinitely, ending early only if the device -/// disappears. -fn wait_for_property( - receiver: &Receiver>, - deadline: Option, - what: &str, - mut is_target: impl FnMut(&T) -> bool, -) -> Result<(), Error> { - let timed_out = || { - Error::with_message( - ErrorKind::DeviceNotAvailable, - format!("Timed out waiting for the {what} to change"), - ) - }; - - loop { - let received = match deadline { - Some(deadline) => { - let remaining = deadline.saturating_duration_since(Instant::now()); - if remaining.is_zero() { - return Err(timed_out()); - } - receiver.recv_timeout(remaining) - } - None => receiver.recv().map_err(|_| RecvTimeoutError::Disconnected), - }; - - match received { - Ok(PropertyEvent::Changed(value)) => { - if is_target(&value) { - return Ok(()); - } - } - Ok(PropertyEvent::Unavailable) => { - return Err(Error::with_message( - ErrorKind::DeviceNotAvailable, - format!("Device disconnected while updating the {what}"), - )); - } - Err(RecvTimeoutError::Timeout) => return Err(timed_out()), - Err(RecvTimeoutError::Disconnected) => { - return Err(Error::with_message( - ErrorKind::StreamInvalidated, - format!("Listener for the {what} disconnected unexpectedly"), - )); - } - } - } -} - -/// Block until the device reports `target_sample_rate`, giving up at `deadline`. -fn wait_for_rate( - receiver: &Receiver>, - target_sample_rate: SampleRate, - deadline: Option, -) -> Result<(), Error> { - wait_for_property(receiver, deadline, "sample rate", |rate| { - (rate - target_sample_rate as f64).abs() < 1.0 - }) -} - #[derive(Clone, Copy)] enum AudioUnitMode { /// HAL Output AudioUnit with input enabled, pinned to a specific device. @@ -533,7 +231,8 @@ impl Device { }; // SAFETY: kAudioDevicePropertyTransportType is documented to return a UInt32. - let Ok(transport) = (unsafe { get_property::(self.audio_device_id, property_address) }) + let Ok(transport) = + (unsafe { get_property::(self.audio_device_id, property_address) }) else { return None; }; @@ -596,8 +295,7 @@ impl Device { // CFString is returned under the create rule, so take ownership of the +1 reference. // SAFETY: kAudioDevicePropertyDeviceUID is documented to return a CFString pointer. - let uid: *mut CFString = - unsafe { get_property(self.audio_device_id, property_address) }?; + let uid: *mut CFString = unsafe { get_property(self.audio_device_id, property_address) }?; // SAFETY: Status was successful, meaning the API call succeeded. // We now check if the returned uid is non-null before use. @@ -1314,76 +1012,3 @@ pub(crate) fn get_device_buffer_frame_size( )?; Ok(frames as usize) } - -#[cfg(test)] -mod tests { - use std::sync::mpsc::channel; - use std::time::{Duration, Instant}; - - use super::{PropertyEvent, wait_for_rate}; - - /// A listener can report rates other than the target before it reports the new one, e.g. a - /// device stepping through rates. The whole timeout must remain available across those. - #[test] - fn wait_for_rate_honours_the_full_timeout_across_repeated_events() { - const TIMEOUT: Duration = Duration::from_millis(50); - - let (sender, receiver) = channel::>(); - let feeder = std::thread::spawn(move || { - while sender.send(PropertyEvent::Changed(44_100.0)).is_ok() { - std::thread::sleep(Duration::from_millis(1)); - } - }); - - let start = Instant::now(); - assert!(wait_for_rate(&receiver, 48_000, Some(start + TIMEOUT)).is_err()); - let elapsed = start.elapsed(); - - drop(receiver); - let _ = feeder.join(); - - assert!( - elapsed >= TIMEOUT - Duration::from_millis(10), - "gave up after {elapsed:?}, well before the {TIMEOUT:?} timeout" - ); - } - - #[test] - fn wait_for_rate_returns_when_the_target_rate_is_reported() { - let (sender, receiver) = channel::>(); - sender.send(PropertyEvent::Changed(44_100.0)).unwrap(); - sender.send(PropertyEvent::Changed(48_000.0)).unwrap(); - - assert!( - wait_for_rate( - &receiver, - 48_000, - Some(Instant::now() + Duration::from_secs(5)) - ) - .is_ok() - ); - } - - /// The wait must end when the device disappears, whether it is unbounded or has a long timeout. - #[test] - fn wait_for_rate_ends_early_on_disconnect() { - for timeout in [None, Some(Duration::from_secs(30))] { - let (sender, receiver) = channel::>(); - let disconnect = std::thread::spawn(move || { - std::thread::sleep(Duration::from_millis(20)); - sender.send(PropertyEvent::Unavailable).unwrap(); - }); - - let start = Instant::now(); - let deadline = timeout.map(|timeout| start + timeout); - assert!(wait_for_rate(&receiver, 48_000, deadline).is_err()); - let elapsed = start.elapsed(); - disconnect.join().unwrap(); - - assert!( - elapsed < Duration::from_secs(5), - "waited {elapsed:?} with timeout {timeout:?}" - ); - } - } -} diff --git a/src/host/coreaudio/macos/format.rs b/src/host/coreaudio/macos/format.rs new file mode 100644 index 000000000..9b9cb509a --- /dev/null +++ b/src/host/coreaudio/macos/format.rs @@ -0,0 +1,401 @@ +//! Negotiate the device's physical stream format and nominal sample rate. +use std::{ + mem::{self, size_of}, + ptr::{NonNull, null}, + sync::mpsc::{Receiver, RecvTimeoutError, channel}, + time::Instant, +}; + +use coreaudio::audio_unit::{ + SampleFormat as CoreAudioSampleFormat, StreamFormat, audio_format::LinearPcmFlags, + macos_helpers::find_matching_physical_format, +}; +use objc2_core_audio::{ + AudioDeviceID, AudioObjectID, AudioObjectPropertyAddress, AudioObjectPropertyScope, + AudioObjectSetPropertyData, AudioStreamID, kAudioDevicePropertyAvailableNominalSampleRates, + kAudioDevicePropertyDeviceIsAlive, kAudioDevicePropertyNominalSampleRate, + kAudioDevicePropertyStreams, kAudioObjectPropertyElementMain, kAudioObjectPropertyScopeGlobal, + kAudioStreamPropertyPhysicalFormat, +}; +use objc2_core_audio_types::{AudioStreamBasicDescription, AudioValueRange}; + +use super::{ + property::{get_property, get_property_array}, + property_listener::AudioObjectPropertyListener, +}; +use crate::{ChannelCount, Error, ErrorKind, SampleFormat, SampleRate}; + +const PHYSICAL_FORMAT_ADDRESS: AudioObjectPropertyAddress = AudioObjectPropertyAddress { + mSelector: kAudioStreamPropertyPhysicalFormat, + mScope: kAudioObjectPropertyScopeGlobal, + mElement: kAudioObjectPropertyElementMain, +}; + +const NOMINAL_SAMPLE_RATE_ADDRESS: AudioObjectPropertyAddress = AudioObjectPropertyAddress { + mSelector: kAudioDevicePropertyNominalSampleRate, + mScope: kAudioObjectPropertyScopeGlobal, + mElement: kAudioObjectPropertyElementMain, +}; + +const DEVICE_IS_ALIVE_ADDRESS: AudioObjectPropertyAddress = AudioObjectPropertyAddress { + mSelector: kAudioDevicePropertyDeviceIsAlive, + mScope: kAudioObjectPropertyScopeGlobal, + mElement: kAudioObjectPropertyElementMain, +}; + +/// Resolve the first `AudioStreamID` a device exposes in `scope` (Input or Output). +/// +/// `kAudioStreamPropertyPhysicalFormat` belongs to the stream object, not the device: reads +/// against the device happen to be forwarded by the HAL, but property-changed notifications are +/// not, so listeners must be registered on the stream object directly. +fn first_stream_id( + device_id: AudioDeviceID, + scope: AudioObjectPropertyScope, +) -> Result { + let address = AudioObjectPropertyAddress { + mSelector: kAudioDevicePropertyStreams, + mScope: scope, + mElement: kAudioObjectPropertyElementMain, + }; + // SAFETY: kAudioDevicePropertyStreams is documented to return an array of AudioStreamID. + let stream_ids: Vec = unsafe { get_property_array(device_id, address) }?; + stream_ids + .into_iter() + .next() + .ok_or(coreaudio::Error::UnsupportedStreamFormat) +} + +fn physical_format( + stream_id: AudioStreamID, +) -> Result { + // SAFETY: kAudioStreamPropertyPhysicalFormat is documented to return an AudioStreamBasicDescription. + unsafe { get_property(stream_id, PHYSICAL_FORMAT_ADDRESS) } +} + +fn asbds_are_equal( + left: &AudioStreamBasicDescription, + right: &AudioStreamBasicDescription, +) -> bool { + left.mSampleRate as u32 == right.mSampleRate as u32 + && left.mFormatID == right.mFormatID + && left.mFormatFlags == right.mFormatFlags + && left.mBytesPerPacket == right.mBytesPerPacket + && left.mFramesPerPacket == right.mFramesPerPacket + && left.mBytesPerFrame == right.mBytesPerFrame + && left.mChannelsPerFrame == right.mChannelsPerFrame + && left.mBitsPerChannel == right.mBitsPerChannel +} + +/// Set the device's physical stream format and wait until it reports the new one. +/// +/// Gives up at `deadline`, or waits indefinitely if it is `None`. If the device disconnects, the +/// wait ends early. +fn set_physical_stream_format( + device_id: AudioDeviceID, + scope: AudioObjectPropertyScope, + new_asbd: AudioStreamBasicDescription, + deadline: Option, +) -> Result<(), Error> { + let stream_id = first_stream_id(device_id, scope)?; + if asbds_are_equal(&physical_format(stream_id)?, &new_asbd) { + return Ok(()); + } + + // Listen before setting the format, so the change can't be missed. + let (receiver, _listeners) = watch_property( + device_id, + stream_id, + PHYSICAL_FORMAT_ADDRESS, + physical_format, + )?; + + let address = PHYSICAL_FORMAT_ADDRESS; + let status = unsafe { + AudioObjectSetPropertyData( + stream_id, + NonNull::from(&address), + 0, + null(), + size_of::() as u32, + NonNull::from(&new_asbd).cast(), + ) + }; + coreaudio::Error::from_os_status(status)?; + + wait_for_property(&receiver, deadline, "physical format", |asbd| { + asbds_are_equal(asbd, &new_asbd) + }) +} + +/// Try to find a matching physical stream format on the device and apply it. +/// +/// This makes the hardware run at the requested bit depth and sample rate directly, without +/// unnecessary conversions. +pub fn set_physical_format( + device_id: AudioDeviceID, + scope: AudioObjectPropertyScope, + sample_rate: SampleRate, + channels: ChannelCount, + sample_format: SampleFormat, + deadline: Option, +) -> Result { + let core_format = match sample_format { + SampleFormat::I8 => CoreAudioSampleFormat::I8, + SampleFormat::I16 => CoreAudioSampleFormat::I16, + SampleFormat::I24 => CoreAudioSampleFormat::I24, + SampleFormat::I32 => CoreAudioSampleFormat::I32, + SampleFormat::F32 => CoreAudioSampleFormat::F32, + _ => return Err(coreaudio::Error::UnsupportedStreamFormat.into()), + }; + let stream_format = StreamFormat { + sample_rate: sample_rate as f64, + sample_format: core_format, + flags: LinearPcmFlags::empty(), + channels: channels as u32, + }; + let asbd = find_matching_physical_format(device_id, stream_format) + .ok_or(coreaudio::Error::UnsupportedStreamFormat)?; + set_physical_stream_format(device_id, scope, asbd, deadline).map(|_| asbd) +} + +/// Read the device's current nominal sample rate. +/// +/// "Nominal" is CoreAudio's term for the rate the device is configured to run at, as opposed to +/// the actual rate measured from its hardware clock (`kAudioDevicePropertyActualSampleRate`). +fn nominal_sample_rate(audio_device_id: AudioObjectID) -> Result { + // SAFETY: kAudioDevicePropertyNominalSampleRate is documented to return an f64. + unsafe { get_property(audio_device_id, NOMINAL_SAMPLE_RATE_ADDRESS) } +} + +/// Set the device's nominal sample rate via `kAudioDevicePropertyNominalSampleRate`. +/// +/// Unlike [`set_physical_format`], this only changes the device clock rate. The AudioUnit bridges +/// any remaining format difference to the virtual stream format the callback sees. +pub fn set_sample_rate( + audio_device_id: AudioObjectID, + target_sample_rate: SampleRate, + deadline: Option, +) -> Result<(), Error> { + let sample_rate = nominal_sample_rate(audio_device_id)?; + let mut property_address = NOMINAL_SAMPLE_RATE_ADDRESS; + + // If the requested sample rate is different to the device sample rate, update the device. + if (sample_rate - target_sample_rate as f64).abs() >= 1.0 { + // Get available sample rate ranges. + property_address.mSelector = kAudioDevicePropertyAvailableNominalSampleRates; + // SAFETY: kAudioDevicePropertyAvailableNominalSampleRates is documented to return an + // array of AudioValueRange. + let ranges: Vec = + unsafe { get_property_array(audio_device_id, property_address) }?; + + // Now that we have the available ranges, pick the one matching the desired rate. + let sample_rate = target_sample_rate; + if !ranges + .iter() + .any(|r| sample_rate as f64 >= r.mMinimum && sample_rate as f64 <= r.mMaximum) + { + return Err(Error::with_message( + ErrorKind::UnsupportedConfig, + format!("Sample rate {sample_rate} Hz is not supported"), + )); + } + + // Listen before setting the rate, so that neither the new rate nor a disconnect is missed. + let (receiver, _listeners) = watch_property( + audio_device_id, + audio_device_id, + NOMINAL_SAMPLE_RATE_ADDRESS, + nominal_sample_rate, + )?; + + // Set the nominal sample rate. + property_address.mSelector = kAudioDevicePropertyNominalSampleRate; + let rate = sample_rate as f64; + let data_size = mem::size_of::() as u32; + let status = unsafe { + AudioObjectSetPropertyData( + audio_device_id, + NonNull::from(&property_address), + 0, + null(), + data_size, + NonNull::from(&rate).cast(), + ) + }; + coreaudio::Error::from_os_status(status)?; + + // Wait for the reported_rate to change. This should not take longer than a few ms. + wait_for_rate(&receiver, target_sample_rate, deadline)?; + // listeners are removed when they drop here + } + Ok(()) +} + +/// What a device reports while one of its properties is changing. +enum PropertyEvent { + /// The property now has this value. + Changed(T), + /// The device disconnected, or can no longer be queried. + Unavailable, +} + +/// Report a property change, or a device disconnect, as a [`PropertyEvent`]. +/// +/// The listeners stay registered for as long as the returned guards are alive. +/// `alive_id` is the device to watch for a disconnect. `target_id` is the object that actually +/// owns `address`: the device itself, or one of its streams. +fn watch_property( + alive_id: AudioDeviceID, + target_id: AudioObjectID, + address: AudioObjectPropertyAddress, + read: fn(AudioObjectID) -> Result, +) -> Result<(Receiver>, [AudioObjectPropertyListener; 2]), Error> { + let (sender, receiver) = channel(); + let alive_sender = sender.clone(); + let alive = AudioObjectPropertyListener::new(alive_id, DEVICE_IS_ALIVE_ADDRESS, move || { + let _ = alive_sender.send(PropertyEvent::Unavailable); + })?; + let changed = AudioObjectPropertyListener::new(target_id, address, move || { + let event = read(target_id).map_or(PropertyEvent::Unavailable, PropertyEvent::Changed); + let _ = sender.send(event); + })?; + Ok((receiver, [alive, changed])) +} + +/// Block until a reported value satisfies `is_target`, giving up at `deadline`. +/// +/// Other values can be reported first, so `deadline` bounds the whole wait rather than each +/// individual receive. A `deadline` of `None` waits indefinitely, ending early only if the device +/// disappears. +fn wait_for_property( + receiver: &Receiver>, + deadline: Option, + what: &str, + mut is_target: impl FnMut(&T) -> bool, +) -> Result<(), Error> { + let timed_out = || { + Error::with_message( + ErrorKind::DeviceNotAvailable, + format!("Timed out waiting for the {what} to change"), + ) + }; + + loop { + let received = match deadline { + Some(deadline) => { + let remaining = deadline.saturating_duration_since(Instant::now()); + if remaining.is_zero() { + return Err(timed_out()); + } + receiver.recv_timeout(remaining) + } + None => receiver.recv().map_err(|_| RecvTimeoutError::Disconnected), + }; + + match received { + Ok(PropertyEvent::Changed(value)) => { + if is_target(&value) { + return Ok(()); + } + } + Ok(PropertyEvent::Unavailable) => { + return Err(Error::with_message( + ErrorKind::DeviceNotAvailable, + format!("Device disconnected while updating the {what}"), + )); + } + Err(RecvTimeoutError::Timeout) => return Err(timed_out()), + Err(RecvTimeoutError::Disconnected) => { + return Err(Error::with_message( + ErrorKind::StreamInvalidated, + format!("Listener for the {what} disconnected unexpectedly"), + )); + } + } + } +} + +/// Block until the device reports `target_sample_rate`, giving up at `deadline`. +fn wait_for_rate( + receiver: &Receiver>, + target_sample_rate: SampleRate, + deadline: Option, +) -> Result<(), Error> { + wait_for_property(receiver, deadline, "sample rate", |rate| { + (rate - target_sample_rate as f64).abs() < 1.0 + }) +} + +#[cfg(test)] +mod tests { + use std::sync::mpsc::channel; + use std::time::{Duration, Instant}; + + use super::{PropertyEvent, wait_for_rate}; + + /// A listener can report rates other than the target before it reports the new one: for + /// example, the device might step through several rates first. The whole timeout must remain + /// available across those. + #[test] + fn wait_for_rate_honours_the_full_timeout_across_repeated_events() { + const TIMEOUT: Duration = Duration::from_millis(50); + + let (sender, receiver) = channel::>(); + let feeder = std::thread::spawn(move || { + while sender.send(PropertyEvent::Changed(44_100.0)).is_ok() { + std::thread::sleep(Duration::from_millis(1)); + } + }); + + let start = Instant::now(); + assert!(wait_for_rate(&receiver, 48_000, Some(start + TIMEOUT)).is_err()); + let elapsed = start.elapsed(); + + drop(receiver); + let _ = feeder.join(); + + assert!( + elapsed >= TIMEOUT - Duration::from_millis(10), + "gave up after {elapsed:?}, well before the {TIMEOUT:?} timeout" + ); + } + + #[test] + fn wait_for_rate_returns_when_the_target_rate_is_reported() { + let (sender, receiver) = channel::>(); + sender.send(PropertyEvent::Changed(44_100.0)).unwrap(); + sender.send(PropertyEvent::Changed(48_000.0)).unwrap(); + + assert!( + wait_for_rate( + &receiver, + 48_000, + Some(Instant::now() + Duration::from_secs(5)) + ) + .is_ok() + ); + } + + /// The wait must end when the device disappears, whether it is unbounded or has a long timeout. + #[test] + fn wait_for_rate_ends_early_on_disconnect() { + for timeout in [None, Some(Duration::from_secs(30))] { + let (sender, receiver) = channel::>(); + let disconnect = std::thread::spawn(move || { + std::thread::sleep(Duration::from_millis(20)); + sender.send(PropertyEvent::Unavailable).unwrap(); + }); + + let start = Instant::now(); + let deadline = timeout.map(|timeout| start + timeout); + assert!(wait_for_rate(&receiver, 48_000, deadline).is_err()); + let elapsed = start.elapsed(); + disconnect.join().unwrap(); + + assert!( + elapsed < Duration::from_secs(5), + "waited {elapsed:?} with timeout {timeout:?}" + ); + } + } +} diff --git a/src/host/coreaudio/macos/mod.rs b/src/host/coreaudio/macos/mod.rs index c3b06c018..057246086 100644 --- a/src/host/coreaudio/macos/mod.rs +++ b/src/host/coreaudio/macos/mod.rs @@ -30,6 +30,7 @@ use crate::{ mod device; pub mod enumerate; +mod format; mod loopback; mod property; mod property_listener; diff --git a/src/host/mod.rs b/src/host/mod.rs index 8e05d77be..3d4446d86 100644 --- a/src/host/mod.rs +++ b/src/host/mod.rs @@ -294,7 +294,7 @@ pub(crate) fn secs_to_nanos(secs: f64) -> u64 { /// `Some(Duration::ZERO)` returns immediately. /// /// [`StreamTrait::stop`]: crate::traits::StreamTrait::stop -#[cfg(any(all(windows, feature = "asio"), all(target_vendor = "apple")))] +#[cfg(any(all(windows, feature = "asio"), target_vendor = "apple"))] pub(crate) fn wait_for_drain(window: std::time::Duration, timeout: Option) { let wait = timeout.map_or(window, |t| window.min(t)); if !wait.is_zero() {