From 9f34b4109bd30ce7e612c88012d877f92c2b79cf Mon Sep 17 00:00:00 2001 From: mofeng-git Date: Thu, 10 Sep 2026 22:43:36 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20=E4=BF=AE=E6=AD=A3=20UAC=20=E9=93=BE?= =?UTF-8?q?=E6=8E=A5=E9=A1=BA=E5=BA=8F=E5=8F=8A=E5=8E=9F=E7=94=9F=20ALSA?= =?UTF-8?q?=20=E5=90=AF=E5=8A=A8=E6=81=A2=E5=A4=8D?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/audio/uac/playback.rs | 408 ++++++++++++++++++++++++++------------ src/otg/manager.rs | 132 ++++++++---- 2 files changed, 371 insertions(+), 169 deletions(-) diff --git a/src/audio/uac/playback.rs b/src/audio/uac/playback.rs index 98fb9805..ad3abfd0 100644 --- a/src/audio/uac/playback.rs +++ b/src/audio/uac/playback.rs @@ -9,9 +9,11 @@ use tracing::{info, warn}; use crate::error::{AppError, Result}; const RETRY_BACKOFF: Duration = Duration::from_secs(1); -const PERIOD_FRAMES: Frames = 960; -const BUFFER_FRAMES: Frames = 4_800; -const START_THRESHOLD_PERIODS: Frames = 4; +const PERIOD_FRAMES: Frames = 1_024; +// Request the same compatibility buffer as the known-working ALSA player. +// The gadget driver may negotiate a smaller buffer; always use its result. +const BUFFER_FRAMES: Frames = 32_768; +const IDLE_REOPEN_TIMEOUT: Duration = Duration::from_secs(5); const SINK_STALL_TIMEOUT: Duration = Duration::from_millis(200); #[derive(Debug, Clone, Copy, PartialEq, Eq)] @@ -58,9 +60,17 @@ struct PlaybackInner { } enum SessionSink { - Closed { retry_at: Option }, - Probing { pcm: PCM, stalled: bool }, - Active { pcm: PCM, last_progress: Instant }, + Closed { + retry_at: Option, + }, + Probing { + pcm: PlaybackPcm, + stalled: bool, + }, + Active { + pcm: PlaybackPcm, + last_progress: Instant, + }, } impl SessionSink { @@ -77,12 +87,14 @@ impl SessionSink { struct SessionRuntime { sink: SessionSink, + last_frame: Option, } impl SessionRuntime { fn new() -> Self { Self { sink: SessionSink::Closed { retry_at: None }, + last_frame: None, } } @@ -92,12 +104,22 @@ impl SessionRuntime { fn close(&mut self) { self.sink = SessionSink::Closed { retry_at: None }; + self.last_frame = None; } /// Advance playback only when a WebSocket frame arrives. All ALSA handles /// are non-blocking, so a slow or absent USB host drops the current frame /// instead of occupying a worker thread or accumulating stale speech. fn write(&mut self, config: &UacPlaybackConfig, samples: &[i16]) -> bool { + // Reopen on resume without an idle timer thread. stop()/session drop + // still close the PCM synchronously, even when no frames arrive. + if self + .last_frame + .is_some_and(|last| last.elapsed() >= IDLE_REOPEN_TIMEOUT) + { + self.close(); + } + self.last_frame = Some(Instant::now()); let sink = std::mem::replace(&mut self.sink, SessionSink::Closed { retry_at: None }); let (next_sink, accepted) = drive_sink(sink, config, samples); self.sink = next_sink; @@ -222,8 +244,8 @@ fn drive_sink( return (SessionSink::Closed { retry_at }, false); } - match open_pcm(config).and_then(|pcm| { - prime_pcm_with_silence(&pcm, config.channels as usize)?; + match open_pcm(config).and_then(|mut pcm| { + pcm.prime_with_silence(config.channels as usize)?; Ok(pcm) }) { Ok(pcm) => drive_probe(pcm, false, config, samples), @@ -246,64 +268,71 @@ fn drive_sink( } fn drive_probe( - pcm: PCM, + mut pcm: PlaybackPcm, stalled: bool, config: &UacPlaybackConfig, samples: &[i16], ) -> (SessionSink, bool) { - match sink_is_consuming(&pcm) { + match pcm.consumption_progress() { Ok(false) => (SessionSink::Probing { pcm, stalled }, false), Ok(true) => { - if let Err(error) = reset_pcm_buffer(&pcm) { - warn!("Failed to activate UAC playback; retrying later: {error}"); - return retry_later(); - } + // Keep the stream that has just started consuming. Dropping and + // preparing it here creates another startup/underrun window. info!("UAC target started consuming microphone audio"); drive_active(pcm, Instant::now(), config, samples) } - Err(error) => { - warn!("Failed to probe UAC playback; retrying later: {error}"); - retry_later() - } + Err(error) => recover_sink(pcm, config, error), } } fn drive_active( - pcm: PCM, + mut pcm: PlaybackPcm, last_progress: Instant, config: &UacPlaybackConfig, samples: &[i16], ) -> (SessionSink, bool) { - match write_pcm_nonblocking(&pcm, samples, config.channels as usize) { - Ok(WriteOutcome::Progress) => ( - SessionSink::Active { - pcm, - last_progress: Instant::now(), - }, - true, - ), - Ok(WriteOutcome::Recovered) => ( - SessionSink::Active { - pcm, - last_progress: Instant::now(), - }, - false, - ), - Ok(WriteOutcome::Blocked) if last_progress.elapsed() < SINK_STALL_TIMEOUT => { - (SessionSink::Active { pcm, last_progress }, false) + let last_progress = match pcm.consumption_progress() { + Ok(true) => Instant::now(), + Ok(false) => last_progress, + Err(error) => return recover_sink(pcm, config, error), + }; + if last_progress.elapsed() >= SINK_STALL_TIMEOUT { + // Discard queued speech before probing an unavailable host again. + if let Err(error) = pcm.reset_and_prime(config.channels as usize) { + warn!("Failed to reset stalled UAC playback: {error}"); + return retry_later(); } - Ok(WriteOutcome::Blocked) => { - if let Err(error) = reset_pcm_buffer(&pcm) - .and_then(|_| prime_pcm_with_silence(&pcm, config.channels as usize)) - { - warn!("Failed to reset stalled UAC playback: {error}"); - return retry_later(); + info!("UAC target stopped consuming audio; waiting for playback activity"); + return (SessionSink::Probing { pcm, stalled: true }, false); + } + + match pcm.write_samples(samples, config.channels as usize) { + Ok(accepted) => (SessionSink::Active { pcm, last_progress }, accepted), + Err(error) => recover_sink(pcm, config, error), + } +} + +fn recover_sink( + mut pcm: PlaybackPcm, + config: &UacPlaybackConfig, + error: alsa::Error, +) -> (SessionSink, bool) { + match error.errno() { + libc::EAGAIN | libc::EINTR => (SessionSink::Probing { pcm, stalled: true }, false), + libc::EPIPE | libc::ESTRPIPE => { + // prepare restarts after XRUN/suspend without snd_pcm_recover's + // potentially unbounded resume loop. Start again with silence, + // and require fresh consumption before reporting Active. + match pcm.reset_and_prime(config.channels as usize) { + Ok(()) => (SessionSink::Probing { pcm, stalled: true }, false), + Err(error) => { + warn!("Failed to recover UAC playback: {error}"); + retry_later() + } } - info!("UAC target stopped consuming audio; waiting for playback activity"); - (SessionSink::Probing { pcm, stalled: true }, false) } - Err(error) => { - warn!("UAC playback write failed; retrying later: {error}"); + _ => { + warn!("UAC playback failed; reopening later: {error}"); retry_later() } } @@ -318,7 +347,7 @@ fn retry_later() -> (SessionSink, bool) { ) } -fn open_pcm(config: &UacPlaybackConfig) -> Result { +fn open_pcm(config: &UacPlaybackConfig) -> Result { let pcm = PCM::new(&config.device_name, Direction::Playback, true).map_err(|error| { AppError::AudioError(format!( "Failed to open UAC device {}: {error}", @@ -348,10 +377,9 @@ fn open_pcm(config: &UacPlaybackConfig) -> Result { let params = pcm.sw_params_current().map_err(|error| { AppError::AudioError(format!("Failed to read UAC SwParams: {error}")) })?; - let start_threshold = - (period_frames as Frames * START_THRESHOLD_PERIODS).min(buffer_frames as Frames); params - .set_start_threshold(start_threshold) + .set_start_threshold(buffer_frames as Frames) + .and_then(|_| params.set_stop_threshold(buffer_frames as Frames)) .and_then(|_| params.set_avail_min(period_frames as Frames)) .and_then(|_| pcm.sw_params(¶ms)) .map_err(|error| { @@ -365,101 +393,219 @@ fn open_pcm(config: &UacPlaybackConfig) -> Result { "UAC playback opened on {} (buffer={} frames, period={} frames)", config.device_name, buffer_frames, period_frames ); - Ok(pcm) + Ok(PlaybackPcm { + pcm, + buffer_frames: buffer_frames as Frames, + period_frames: period_frames as Frames, + submitted_frames: 0, + consumed_frames: 0, + }) } -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -enum WriteOutcome { - Progress, - Blocked, - Recovered, +struct PlaybackPcm { + pcm: PCM, + buffer_frames: Frames, + period_frames: Frames, + submitted_frames: u64, + consumed_frames: u64, } -fn write_pcm_nonblocking(pcm: &PCM, samples: &[i16], channels: usize) -> Result { +impl PlaybackPcm { + fn consumption_progress(&mut self) -> std::result::Result { + // avail synchronizes the hardware pointer. Successful writes alone + // only show that the ring buffer has room, not that USB is consuming. + let available = self.pcm.avail()?; + match self.pcm.state() { + State::XRun => return Err(alsa::Error::new("UAC PCM state", libc::EPIPE)), + State::Suspended => return Err(alsa::Error::new("UAC PCM state", libc::ESTRPIPE)), + State::Disconnected => return Err(alsa::Error::new("UAC PCM state", libc::ENODEV)), + State::Running => {} + _ => return Ok(false), + } + let consumed = consumed_frames(self.submitted_frames, self.buffer_frames, available); + let progressed = consumed > self.consumed_frames; + self.consumed_frames = consumed; + Ok(progressed) + } + + fn write_samples( + &mut self, + samples: &[i16], + channels: usize, + ) -> std::result::Result { + let io = self.pcm.io_i16()?; + let written = write_frames(samples, channels, self.period_frames as usize, |chunk| { + let written = io.writei(chunk)?; + self.submitted_frames += written as u64; + Ok(written) + })?; + Ok(written == samples.len() / channels) + } + + fn prime_with_silence(&mut self, channels: usize) -> Result<()> { + // Use the negotiated capacity, not BUFFER_FRAMES. A near request is + // often clamped by u_audio's DMA buffer limit. + let silence = vec![0i16; self.buffer_frames as usize * channels]; + let complete = self + .write_samples(&silence, channels) + .map_err(|error| AppError::AudioError(format!("Failed to prime UAC PCM: {error}")))?; + if !complete { + return Err(AppError::AudioError( + "UAC PCM priming was interrupted".into(), + )); + } + // Most hardware starts automatically at the threshold. Some PCM + // plugins remain Prepared despite accepting the complete prefill. + // Start explicitly only after priming, and never restart a running PCM. + if self.pcm.state() == State::Prepared { + self.pcm.start().map_err(|error| { + AppError::AudioError(format!("Failed to start primed UAC PCM: {error}")) + })?; + } + Ok(()) + } + + fn reset_and_prime(&mut self, channels: usize) -> Result<()> { + self.pcm + .drop() + .and_then(|_| self.pcm.prepare()) + .map_err(|error| AppError::AudioError(format!("Failed to reset UAC PCM: {error}")))?; + self.submitted_frames = 0; + self.consumed_frames = 0; + self.prime_with_silence(channels) + } +} + +fn consumed_frames(submitted: u64, buffer: Frames, available: Frames) -> u64 { + let queued = (buffer - available.clamp(0, buffer)) as u64; + submitted.saturating_sub(queued) +} + +/// Bound each write to one negotiated period and advance by actual frames, +/// including short writes. Never wait for space or retain stale audio. +fn write_frames( + samples: &[i16], + channels: usize, + period_frames: usize, + mut write: impl FnMut(&[i16]) -> std::result::Result, +) -> std::result::Result { let total_frames = samples.len() / channels; - match pcm.avail() { - Ok(available) if available < total_frames as Frames => return Ok(WriteOutcome::Blocked), - Ok(_) => {} - Err(error) => { - recover_pcm(pcm, error)?; - return Ok(WriteOutcome::Recovered); - } - } - - let io = pcm - .io_i16() - .map_err(|error| AppError::AudioError(format!("UAC PCM I/O failed: {error}")))?; - match io.writei(samples) { - Ok(0) => Ok(WriteOutcome::Blocked), - Ok(_) => Ok(WriteOutcome::Progress), - Err(error) if error.errno() == libc::EAGAIN => Ok(WriteOutcome::Blocked), - Err(error) => { - recover_pcm(pcm, error)?; - Ok(WriteOutcome::Recovered) - } - } -} - -/// Once a full playback buffer gains at least one period of free space, the -/// USB host has enabled the UAC streaming interface and is consuming samples. -fn sink_is_consuming(pcm: &PCM) -> Result { - if pcm.state() == State::XRun { - return Ok(true); - } - - match pcm.avail() { - Ok(available) => Ok(available >= PERIOD_FRAMES), - Err(error) if error.errno() == libc::EPIPE => Ok(true), - Err(error) => Err(AppError::AudioError(format!( - "Failed to query UAC playback availability: {error}" - ))), - } -} - -fn recover_pcm(pcm: &PCM, error: alsa::Error) -> Result<()> { - let errno = error.errno(); - pcm.try_recover(error, true).map_err(|recover_error| { - AppError::AudioError(format!("Failed to recover UAC playback: {recover_error}")) - })?; - if matches!(errno, libc::EPIPE | libc::ESTRPIPE) { - warn!("Recovered UAC playback after ALSA error {errno}"); - } - Ok(()) -} - -fn reset_pcm_buffer(pcm: &PCM) -> Result<()> { - pcm.drop() - .and_then(|_| pcm.prepare()) - .map_err(|error| AppError::AudioError(format!("Failed to reset UAC PCM: {error}"))) -} - -/// Prime the non-blocking ALSA buffer with silence. Subsequent WebSocket -/// frames inspect buffer progress to detect when the USB host starts reading. -fn prime_pcm_with_silence(pcm: &PCM, channels: usize) -> Result<()> { - let silence = vec![0i16; BUFFER_FRAMES as usize * channels]; - let io = pcm - .io_i16() - .map_err(|error| AppError::AudioError(format!("UAC PCM I/O failed: {error}")))?; - let mut frame_offset = 0usize; - while frame_offset < BUFFER_FRAMES as usize { - match io.writei(&silence[frame_offset * channels..]) { + let mut offset = 0; + while offset < total_frames { + let end = (offset + period_frames).min(total_frames); + match write(&samples[offset * channels..end * channels]) { Ok(0) => break, - Ok(written) => frame_offset += written, - Err(error) if error.errno() == libc::EAGAIN => break, - Err(error) => { - return Err(AppError::AudioError(format!( - "Failed to prime UAC PCM with silence: {error}" - ))); - } + Ok(written) => offset += written, + Err(error) if matches!(error.errno(), libc::EAGAIN | libc::EINTR) => break, + Err(error) => return Err(error), } } - Ok(()) + Ok(offset) } #[cfg(test)] mod tests { use super::*; + #[test] + fn writes_short_frames_without_skipping_stereo_samples() { + let samples: Vec = (0..24).collect(); + let mut received = Vec::new(); + let written = write_frames(&samples, 2, 4, |chunk| { + assert!(chunk.len() <= 8); + // Simulate a device accepting only one frame per write. + received.extend_from_slice(&chunk[..2]); + Ok(1) + }) + .unwrap(); + assert_eq!(written, 12); + assert_eq!(received, samples); + } + + #[test] + fn full_device_stops_writing_without_waiting_or_claiming_whole_packet() { + let mut calls = 0; + let written = write_frames(&[0; 24], 2, 4, |_| { + calls += 1; + if calls == 1 { + Ok(2) + } else { + Err(alsa::Error::new("test write", libc::EAGAIN)) + } + }) + .unwrap(); + assert_eq!(written, 2); + assert_eq!(calls, 2); + } + + #[test] + fn writing_into_free_space_is_not_host_consumption() { + // A partially filled or full buffer can exist without any USB I/O. + assert_eq!(consumed_frames(1024, 4096, 3072), 0); + assert_eq!(consumed_frames(4096, 4096, 0), 0); + // Consuming a period, followed by filling it again, preserves progress. + assert_eq!(consumed_frames(4096, 4096, 1024), 1024); + assert_eq!(consumed_frames(5120, 4096, 0), 1024); + } + + fn null_config() -> UacPlaybackConfig { + UacPlaybackConfig { + device_name: "null".into(), + sample_rate: 48_000, + channels: 2, + } + } + + #[test] + fn native_pcm_uses_negotiated_start_threshold_and_recovers_to_probing() { + // ALSA's null plugin exercises real libasound configuration and I/O + // without requiring a USB controller. It cannot verify DWC3 behavior. + let config = null_config(); + let mut pcm = open_pcm(&config).unwrap(); + assert_eq!( + pcm.pcm + .sw_params_current() + .unwrap() + .get_start_threshold() + .unwrap(), + pcm.buffer_frames + ); + pcm.prime_with_silence(2).unwrap(); + assert!(pcm.consumption_progress().unwrap()); + let (sink, accepted) = + recover_sink(pcm, &config, alsa::Error::new("test xrun", libc::EPIPE)); + assert!(!accepted); + assert!(matches!(sink, SessionSink::Probing { stalled: true, .. })); + } + + #[test] + fn stop_closes_an_open_native_pcm_before_returning() { + let playback = UacPlayback::start(null_config()).unwrap(); + let session = playback.acquire_session().unwrap(); + session.try_write(&[0; 2048]).unwrap(); + assert!(!matches!( + session.runtime.lock().unwrap().sink, + SessionSink::Closed { .. } + )); + playback.stop(); + assert!(matches!( + session.runtime.lock().unwrap().sink, + SessionSink::Closed { retry_at: None } + )); + } + + #[test] + fn idle_resume_reopens_instead_of_reusing_previous_sink() { + let config = null_config(); + let mut runtime = SessionRuntime::new(); + runtime.sink = SessionSink::Closed { + retry_at: Some(Instant::now() + Duration::from_secs(60)), + }; + runtime.last_frame = Some(Instant::now() - IDLE_REOPEN_TIMEOUT); + runtime.write(&config, &[0; 2048]); + assert!(!matches!(runtime.sink, SessionSink::Closed { .. })); + } + #[test] fn permits_only_one_microphone_session() { let playback = UacPlayback::start(UacPlaybackConfig::default()).unwrap(); diff --git a/src/otg/manager.rs b/src/otg/manager.rs index 147c6c7a..8bf4899a 100644 --- a/src/otg/manager.rs +++ b/src/otg/manager.rs @@ -3,9 +3,9 @@ use std::path::PathBuf; use tracing::{debug, error, info, warn}; use super::configfs::{ - configfs_path, create_dir, create_symlink, find_udc, is_configfs_available, remove_dir, - remove_file, write_file, write_file_if_exists, DEFAULT_GADGET_NAME, DEFAULT_USB_BCD_DEVICE, - DEFAULT_USB_PRODUCT_ID, DEFAULT_USB_VENDOR_ID, USB_BCD_USB, + configfs_path, create_dir, find_udc, is_configfs_available, remove_dir, write_file, + write_file_if_exists, DEFAULT_GADGET_NAME, DEFAULT_USB_BCD_DEVICE, DEFAULT_USB_PRODUCT_ID, + DEFAULT_USB_VENDOR_ID, USB_BCD_USB, }; use super::function::GadgetFunction; use super::hid::HidFunction; @@ -221,9 +221,7 @@ impl OtgGadgetManager { } pub fn bind(&mut self, udc: &str) -> Result<()> { - if let Err(e) = self.recreate_config_links() { - warn!("Failed to recreate gadget config links before bind: {}", e); - } + self.recreate_config_links()?; debug!("Binding gadget to UDC: {}", udc); write_file(&self.gadget_path.join("UDC"), &udc)?; @@ -385,39 +383,22 @@ impl OtgGadgetManager { return Ok(()); } - let entries = std::fs::read_dir(&functions_path).map_err(|e| { - AppError::Internal(format!( - "Failed to read functions directory {}: {}", - functions_path.display(), - e - )) - })?; - - for entry in entries.flatten() { - let name = entry.file_name(); - let name = match name.to_str() { - Some(n) => n, - None => continue, - }; - if !name.contains(".usb") { - continue; + // ConfigFS binds functions in link insertion order. Preserve the + // setup order (UAC before HID), including on rebind: directory + // iteration order is unspecified and can change endpoint allocation. + for func in &self.functions { + let dest = self.config_path.join(func.name()); + if dest.symlink_metadata().is_ok() { + fs::remove_file(&dest).map_err(|error| { + AppError::Internal(format!( + "Failed to remove config link {}: {error}", + dest.display() + )) + })?; } - - let src = functions_path.join(name); - let dest = self.config_path.join(name); - - if dest.exists() { - if let Err(e) = remove_file(&dest) { - warn!( - "Failed to remove existing config link {}: {}", - dest.display(), - e - ); - continue; - } - } - - create_symlink(&src, &dest)?; + } + for func in &self.functions { + func.link(&self.config_path, &self.gadget_path)?; } Ok(()) @@ -470,6 +451,81 @@ pub async fn wait_for_hid_devices(device_paths: &[PathBuf], timeout_ms: u64) -> #[cfg(test)] mod tests { use super::*; + use std::path::Path; + use std::sync::{Arc, Mutex}; + + struct RecordedFunction { + name: &'static str, + links: Arc>>, + } + + impl GadgetFunction for RecordedFunction { + fn name(&self) -> &str { + self.name + } + fn create(&self, _: &Path) -> Result<()> { + Ok(()) + } + fn link(&self, config: &Path, gadget: &Path) -> Result<()> { + super::super::configfs::create_symlink( + &gadget.join("functions").join(self.name), + &config.join(self.name), + )?; + self.links.lock().unwrap().push(self.name.into()); + Ok(()) + } + fn unlink(&self, _: &Path) -> Result<()> { + Ok(()) + } + fn cleanup(&self, _: &Path) -> Result<()> { + Ok(()) + } + } + + #[test] + fn rebind_links_uac_before_hid_in_registration_order() { + let temp = tempfile::tempdir().unwrap(); + let mut manager = OtgGadgetManager::new(); + manager.gadget_path = temp.path().to_path_buf(); + manager.config_path = temp.path().join("configs/c.1"); + fs::create_dir_all(&manager.config_path).unwrap(); + let links = Arc::new(Mutex::new(Vec::new())); + for name in ["uac1.usb0", "hid.usb0", "mass_storage.usb0"] { + fs::create_dir_all(temp.path().join("functions").join(name)).unwrap(); + // Include a dangling pre-existing link, as well as testing rebind. + std::os::unix::fs::symlink("/nonexistent-uac-test", manager.config_path.join(name)) + .unwrap(); + manager + .add_function(Box::new(RecordedFunction { + name, + links: Arc::clone(&links), + })) + .unwrap(); + } + for _ in 0..2 { + links.lock().unwrap().clear(); + manager.recreate_config_links().unwrap(); + assert_eq!( + *links.lock().unwrap(), + ["uac1.usb0", "hid.usb0", "mass_storage.usb0"] + ); + } + } + + #[test] + fn link_failure_prevents_udc_binding() { + let temp = tempfile::tempdir().unwrap(); + let mut manager = OtgGadgetManager::new(); + manager.gadget_path = temp.path().to_path_buf(); + manager.config_path = temp.path().join("configs/c.1"); + fs::create_dir_all(temp.path().join("functions/hid.usb0")).unwrap(); + // A directory occupying a config link cannot be removed as a file. + fs::create_dir_all(manager.config_path.join("hid.usb0")).unwrap(); + fs::write(temp.path().join("UDC"), "").unwrap(); + manager.add_keyboard(false).unwrap(); + assert!(manager.bind("test-udc").is_err()); + assert_eq!(fs::read_to_string(temp.path().join("UDC")).unwrap(), ""); + } #[test] fn test_manager_creation() {