feat(player): add direct USB audio and DSD transport

This commit is contained in:
zarzet committed 2026-09-27 15:26:10 +07:00
1 parent 4e464b5b6b
commit 608675a255
31 files changed
+2875 -12

No files matched your search

@@ -561,6 +561,8 @@ fn supported(path: &str) -> bool {
| "wav"
| "aiff"
| "aif"
| "dsf"
| "dff"
| "cue"
)
}
+4
View File
@@ -18,5 +18,9 @@ serde_json.workspace = true
zeroize.workspace = true
rustix.workspace = true
[target.'cfg(target_os = "android")'.dependencies]
libusb1-sys = { version = "=0.7.0", features = ["vendored"] }
libc = "0.2"
[lints]
workspace = true
+1
View File
@@ -13,6 +13,7 @@ mod metadata;
mod progress;
mod repository;
mod tags;
mod usb_audio;
uniffi::setup_scaffolding!();
@@ -0,0 +1,74 @@
# Direct USB audio
Android USB host transport for local FLAC/WAV PCM and uncompressed DSF/DSDIFF.
The existing Android 14 preferred-mixer path remains available. Direct USB is
opt-in under Settings > Library > Playback > USB bit-perfect audio.
## Ownership and playback
- Android asks for USB permission; the engine, not the Activity, owns the
connection. Permission denial and absent hardware leave PCM on normal output.
- One Rust worker owns libusb, its claimed interfaces, four isochronous output
transfers, and a separate explicit-feedback transfer. The bounded queue holds
approximately 200 ms of audio (minimum 256 KiB). There is no resampling or gain.
- UAC1 fixed rates are verified from descriptors; variable UAC1/UAC2 rates are
read back from the endpoint/clock. Implicit feedback, ambiguous clocks, and
multiple configurations fail closed. Endpoint capacities bound packet sizes.
- Fractional rates and 10.14/16.16 feedback determine packet lengths. Completed
data frames advance position; silence inserted during an underrun does not.
- Every native transfer is cancelled and acknowledged before its memory is
freed. The worker joins before Android closes the original USB descriptor.
Interfaces are released and the kernel driver reattached automatically.
- Pause/seek discard queued audio and reopen the decoder position at the last
completed frame. Disconnect pauses playback; it never redirects active USB
playback to the speaker. Pending permission requests can be cancelled by Next.
- ReplayGain, AutoMix, software volume and playback-rate processing are bypassed
while USB mode is selected. Volume is controlled at the DAC.
## DSD
`DsdFile.kt` normalizes DSF channel blocks/LSB ordering and DFF byte-interleaved
MSB data without loading the file into memory. Native U32 uses the device's
byte ordering. DoP uses 24-bit words, left-aligned when carried in 32-bit slots.
The transport assigns alternating 05/FA markers across all frames, including
inserted DSD silence. Seeking retains channel/frame alignment.
Native DSD initially recognizes these exact USB identities and alternate
settings; no product-name match or assumption that PCM support implies DSD:
| VID:PID | Alternate | Firmware | Encoding |
| --- | --- | --- | --- |
| 16d0:071a (Amanero Combo384) | 2 | 0199 | U32 LE |
| 16d0:071a (Amanero Combo384) | 2 | 019b, 0203 | U32 BE |
| 2772:0230 (Pro-Ject Pre Box S2 Digital) | 2 | any | U32 BE |
| 20b1:3089 (Mola-Mola) | 2 | any | U32 BE |
DoP is a separate opt-in for a DAC explicitly known to support it; USB audio
descriptors do not announce DoP support. Keep it disabled for ordinary USB
headsets. Unsupported DSD is stopped rather than interpreted as PCM audio.
WavPack DSD, DST compression, DSD-to-PCM conversion and vendor DAP outputs are not
implemented. SAF scanning reads basic DSD format/duration and uses the filename
for metadata; embedded DSD artwork/tags are not yet imported.
## Verification
- Rust descriptor/clock tests: UAC1 and UAC2 fixtures, malformed data, exact
native-DSD identity matching, precision rejection, fractional rates/feedback.
- JVM tests: DSF padding/bit order, DFF channel order, native LE/BE, DoP layout,
seeking, truncated containers and compressed-DST rejection.
- Android instrumentation: absent USB fallback, real native FFI error cleanup,
and integer WAV decode/seek after the original SAF descriptor is closed.
- Flutter tests: opt-in settings, transport options, cancellation, stale replies,
DSD failure without PCM fallback, existing playback/queue/DSP behavior.
Physical PCM/DoP/native-DSD output, DAC lock indication, hardware volume and long
playback stability still require real USB hardware. Host tests, emulator tests
and successful ARM32/ARM64 builds do not establish DAC compatibility.
## Reference and dependencies
USB identity/format facts were checked against
[Linux USB audio quirks](https://github.com/torvalds/linux/blob/master/sound/usb/quirks.c).
The Android-only `libusb1-sys` dependency builds its unmodified bundled libusb;
license texts and rebuild information are in `assets/licenses/usb.txt` and in
the app's license registry. No additional audio SDK or decoder library is added.
@@ -0,0 +1,380 @@
use super::UsbOutputFormat;
use std::collections::HashMap;
#[derive(Clone, Debug, Default)]
pub struct Endpoint {
pub address: u8,
pub attributes: u8,
pub max_packet: usize,
pub interval: u8,
pub sync_address: u8,
pub rate_control: bool,
}
#[derive(Clone, Debug, Default)]
pub struct Alternate {
pub control: u8,
pub interface: u8,
pub alternate: u8,
pub uac2: bool,
pub terminal: u8,
pub clock: u8,
pub channels: u8,
pub bits: u8,
pub subslot: u8,
pub pcm: bool,
pub rates: Vec<(u32, u32)>,
pub endpoints: Vec<Endpoint>,
}
#[derive(Default)]
pub struct Device {
pub vendor: u16,
pub product: u16,
pub revision: u16,
pub alternates: Vec<Alternate>,
}
fn u16le(b: &[u8]) -> u16 {
u16::from_le_bytes([b[0], b[1]])
}
fn u24le(b: &[u8]) -> u32 {
u32::from_le_bytes([b[0], b[1], b[2], 0])
}
/// Only UAC1/2 Type I, mono/stereo output. Ambiguous clock selectors and
/// implicit feedback are rejected instead of guessing a configuration.
pub fn parse(bytes: &[u8]) -> Result<Device, String> {
let mut device = Device::default();
let mut current: Option<Alternate> = None;
let mut control = 0;
let mut is_control = false;
let mut clocks = HashMap::new();
let mut offset = 0;
let mut configurations = 0;
while offset < bytes.len() {
let length = *bytes.get(offset).ok_or("truncated descriptor")? as usize;
if length < 2 || offset + length > bytes.len() {
return Err("malformed USB descriptor".into());
}
let d = &bytes[offset..offset + length];
offset += length;
match d[1] {
2 => {
configurations += 1;
if configurations > 1 {
return Err("Multiple USB configurations are not supported".into());
}
}
1 if length >= 18 => {
device.vendor = u16le(&d[8..]);
device.product = u16le(&d[10..]);
device.revision = u16le(&d[12..]);
}
4 if length >= 9 => {
if let Some(a) = current.take() {
device.alternates.push(a);
}
is_control = d[5] == 1 && d[6] == 1 && d[7] == 0x20;
if d[5] == 1 && d[6] == 1 {
control = d[2];
}
if d[5] == 1 && d[6] == 2 && d[3] > 0 && (d[7] == 0 || d[7] == 0x20) {
current = Some(Alternate {
control,
interface: d[2],
alternate: d[3],
uac2: d[7] == 0x20,
..Default::default()
});
}
}
0x24 if is_control && length >= 9 && d[2] == 2 => {
clocks.insert((control, d[3]), d[7]);
}
0x24 => {
if let Some(a) = current.as_mut() {
if length >= 7 && d[2] == 1 {
a.terminal = d[3];
if a.uac2 && length >= 16 && d[5] == 1 {
let formats = u32::from_le_bytes(d[6..10].try_into().unwrap());
a.pcm = formats & 1 != 0;
a.channels = d[10];
} else if !a.uac2 {
a.pcm = u16le(&d[5..]) == 1;
}
}
if length >= 6 && d[2] == 2 && d[3] == 1 {
if a.uac2 {
a.subslot = d[4];
a.bits = d[5];
} else if length >= 8 {
a.channels = d[4];
a.subslot = d[5];
a.bits = d[6];
let count = d[7] as usize;
if count == 0 && length >= 14 {
a.rates.push((u24le(&d[8..]), u24le(&d[11..])));
} else if length >= 8 + count * 3 {
for f in d[8..8 + count * 3].as_chunks::<3>().0 {
let rate = u24le(f);
a.rates.push((rate, rate));
}
}
}
}
}
}
5 if length >= 7 => {
if let Some(a) = current.as_mut() {
let packet = u16le(&d[4..]);
a.endpoints.push(Endpoint {
address: d[2],
attributes: d[3],
max_packet: (packet & 0x7ff) as usize * (1 + ((packet >> 11) & 3) as usize),
interval: d[6],
sync_address: d.get(8).copied().unwrap_or(0),
rate_control: false,
});
}
}
0x25 if length >= 4 && d[2] == 1 => {
if let Some(a) = current.as_mut()
&& !a.uac2
&& let Some(endpoint) = a.endpoints.last_mut()
{
endpoint.rate_control = d[3] & 1 != 0;
}
}
_ => {}
}
}
if let Some(a) = current {
device.alternates.push(a);
}
for a in &mut device.alternates {
a.clock = clocks.get(&(a.control, a.terminal)).copied().unwrap_or(0);
}
Ok(device)
}
/// A deliberately small exact-match registry based on documented USB transport
/// identities in Linux sound/usb/quirks.c, not a product-name heuristic.
fn native_encoding(device: &Device, a: &Alternate) -> Option<&'static str> {
if !a.uac2 || a.subslot != 4 {
return None;
}
match (device.vendor, device.product, a.alternate, device.revision) {
(0x16d0, 0x071a, 2, 0x0199) => Some("dsd_le"),
(0x16d0, 0x071a, 2, 0x019b | 0x0203) => Some("dsd_be"),
(0x2772, 0x0230, 2, _) | (0x20b1, 0x3089, 2, _) => Some("dsd_be"),
_ => None,
}
}
pub fn formats(
device: &Device,
rate: u32,
channels: u8,
bits: u8,
dsd: bool,
dop: bool,
) -> Vec<(Alternate, UsbOutputFormat)> {
let mut result = Vec::new();
for a in &device.alternates {
if a.channels != channels || !(1..=2).contains(&channels) || !(2..=4).contains(&a.subslot) {
continue;
}
let native = native_encoding(device, a);
let (wire_rate, encoding, precision) = if dsd {
if let Some(encoding) = native {
(rate / 32, encoding, 1)
} else if dop && a.pcm && a.bits >= 24 && a.subslot >= 3 {
(rate / 16, "dop", 24)
} else {
continue;
}
} else {
if !a.pcm || native.is_some() || !matches!(bits, 16 | 24 | 32) || a.bits < bits {
continue;
}
(rate, "pcm", a.bits)
};
if wire_rate == 0 || a.bits > a.subslot * 8 || (a.uac2 && a.clock == 0) {
continue;
}
if !a.uac2
&& !a
.rates
.iter()
.any(|(min, max)| wire_rate >= *min && wire_rate <= *max)
{
continue;
}
result.push((
a.clone(),
UsbOutputFormat {
sample_rate: wire_rate,
channels,
bits: precision,
subslot: a.subslot,
encoding: encoding.into(),
},
));
}
result.sort_by_key(|(_, f)| {
(
if f.encoding.starts_with("dsd_") { 0 } else { 1 },
f.subslot,
)
});
result
}
pub struct PacketClock {
nominal: f64,
feedback: Option<f64>,
remainder: f64,
}
impl PacketClock {
pub fn new(rate: u32, interval_us: u32) -> Self {
Self {
nominal: rate as f64 * interval_us as f64 / 1_000_000.0,
feedback: None,
remainder: 0.0,
}
}
pub fn feedback(&mut self, bytes: &[u8], high_speed: bool, interval_us: u32) -> bool {
let frames = match bytes.len() {
3 if !high_speed => u24le(bytes) as f64 / 16384.0,
4 => u32::from_le_bytes(bytes.try_into().unwrap()) as f64 / 65536.0,
_ => return false,
} * interval_us as f64
/ if high_speed { 125.0 } else { 1000.0 };
if !frames.is_finite() || (frames - self.nominal).abs() > self.nominal * 0.02 {
return false;
}
self.feedback = Some(frames);
true
}
pub fn next(&mut self) -> usize {
self.remainder += self.feedback.unwrap_or(self.nominal);
let frames = self.remainder.floor() as usize;
self.remainder -= frames as f64;
frames
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn parses_uac1_fixed_rate_headset_without_clock_controls() {
let raw = [
9, 4, 0, 0, 0, 1, 1, 0, 0, 9, 4, 1, 1, 1, 1, 2, 0, 0, 7, 0x24, 1, 1, 1, 1, 0, 11, 0x24,
2, 1, 2, 2, 16, 1, 0x80, 0xbb, 0, 9, 5, 1, 0x0d, 192, 0, 1, 0, 0, 7, 0x25, 1, 0, 0, 0,
0,
];
let device = parse(&raw).unwrap();
let matches = formats(&device, 48000, 2, 16, false, false);
assert_eq!(matches.len(), 1);
let a = &matches[0].0;
assert_eq!(a.interface, 1);
assert!(!a.endpoints[0].rate_control);
assert_eq!(a.endpoints[0].address, 1);
assert_eq!(a.endpoints[0].max_packet, 192);
assert_eq!(a.endpoints[0].interval, 1);
assert_eq!(a.endpoints[0].attributes, 0x0d);
assert_eq!(a.endpoints[0].sync_address, 0);
assert!(formats(&device, 44100, 2, 16, false, false).is_empty());
}
#[test]
fn resolves_uac2_terminal_clock_and_explicit_feedback() {
let raw = [
9, 4, 2, 0, 0, 1, 1, 0x20, 0, 17, 0x24, 2, 7, 1, 1, 0, 9, 2, 3, 0, 0, 0, 0, 0, 0, 0, 9,
4, 3, 1, 2, 1, 2, 0x20, 0, 16, 0x24, 1, 7, 0, 1, 1, 0, 0, 0, 2, 3, 0, 0, 0, 0, 6, 0x24,
2, 1, 4, 24, 7, 5, 1, 5, 0, 4, 1, 7, 5, 0x81, 0x11, 4, 0, 4,
];
let device = parse(&raw).unwrap();
let matches = formats(&device, 96000, 2, 24, false, false);
assert_eq!(matches.len(), 1);
assert_eq!(matches[0].0.control, 2);
assert_eq!(matches[0].0.clock, 9);
assert_eq!(matches[0].0.endpoints[1].address, 0x81);
assert_eq!(matches[0].1.subslot, 4);
assert!(formats(&device, 96000, 2, 32, false, false).is_empty());
}
#[test]
fn fractional_rates_do_not_drift() {
for (rate, interval, packets) in
[(44100, 125, 8000), (88200, 1000, 1000), (176400, 125, 8000)]
{
let mut clock = PacketClock::new(rate, interval);
assert_eq!(
(0..packets).map(|_| clock.next()).sum::<usize>(),
rate as usize
);
}
}
#[test]
fn invalid_feedback_cannot_overrun_endpoint() {
let mut clock = PacketClock::new(48000, 125);
assert!(!clock.feedback(&[255; 4], true, 125));
assert!(!clock.feedback(&[0; 4], true, 125));
assert!(clock.feedback(&(6u32 << 16).to_le_bytes(), true, 125));
assert_eq!(clock.next(), 6);
}
#[test]
fn untrusted_descriptors_are_bounded() {
for data in [vec![0, 4], vec![9, 4, 0], vec![1], vec![255; 18]] {
assert!(parse(&data).is_err());
}
assert!(parse(&[]).unwrap().alternates.is_empty());
}
#[test]
fn unknown_dsd_device_requires_explicit_dop() {
let a = Alternate {
uac2: true,
clock: 1,
channels: 2,
bits: 24,
subslot: 3,
pcm: true,
..Default::default()
};
let device = Device {
alternates: vec![a],
..Default::default()
};
assert!(formats(&device, 2822400, 2, 1, true, false).is_empty());
let f = formats(&device, 2822400, 2, 1, true, true);
assert_eq!(f[0].1.sample_rate, 176400);
assert_eq!(f[0].1.encoding, "dop");
assert!(formats(&device, 96000, 2, 32, false, false).is_empty());
}
#[test]
fn native_requires_exact_firmware_and_altsetting() {
let a = Alternate {
uac2: true,
alternate: 2,
clock: 1,
channels: 2,
subslot: 4,
bits: 32,
pcm: true,
..Default::default()
};
let mut d = Device {
vendor: 0x16d0,
product: 0x071a,
revision: 0x0199,
alternates: vec![a],
};
assert_eq!(
formats(&d, 2822400, 2, 1, true, false)[0].1.encoding,
"dsd_le"
);
d.revision = 0x100;
assert!(formats(&d, 2822400, 2, 1, true, false).is_empty());
}
}
@@ -0,0 +1,77 @@
use super::UsbOutputFormat;
/// Completes underruns without putting PCM zero words into a DSD stream.
/// DoP markers describe wire frames, so they must continue through silence too.
pub fn finish_transfer(
bytes: &mut [u8],
real_bytes: usize,
format: &UsbOutputFormat,
marker: &mut u8,
) {
bytes[real_bytes..].fill(if format.encoding == "pcm" { 0 } else { 0x69 });
if format.encoding == "dop" {
let subslot = format.subslot as usize;
let frame = subslot * format.channels as usize;
for chunk in bytes.chunks_exact_mut(frame) {
for sample in chunk.chunks_exact_mut(subslot) {
sample[subslot - 1] = *marker;
if subslot == 4 {
sample[0] = 0;
}
}
*marker ^= 0xff;
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn format(encoding: &str, subslot: u8) -> UsbOutputFormat {
UsbOutputFormat {
sample_rate: 176400,
channels: 2,
bits: 24,
subslot,
encoding: encoding.into(),
}
}
#[test]
fn dop_markers_alternate_across_transfers_and_underruns() {
let mut marker = 5;
let mut first = [1, 2, 0, 3, 4, 0];
finish_transfer(&mut first, 6, &format("dop", 3), &mut marker);
assert_eq!(first, [1, 2, 5, 3, 4, 5]);
let mut second = [0; 12];
finish_transfer(&mut second, 0, &format("dop", 3), &mut marker);
assert_eq!(
second,
[
0x69, 0x69, 0xfa, 0x69, 0x69, 0xfa, 0x69, 0x69, 5, 0x69, 0x69, 5
]
);
assert_eq!(marker, 0xfa);
}
#[test]
fn dop32_padding_and_channel_markers_match() {
let mut data = [0, 1, 2, 0, 0, 3, 4, 0, 0, 0, 0, 0, 0, 0, 0, 0];
finish_transfer(&mut data, 8, &format("dop", 4), &mut 5);
assert_eq!(
data,
[
0, 1, 2, 5, 0, 3, 4, 5, 0, 0x69, 0x69, 0xfa, 0, 0x69, 0x69, 0xfa
]
);
}
#[test]
fn native_dsd_and_pcm_keep_original_data_and_use_different_silence() {
for (encoding, silence) in [("pcm", 0), ("dsd_be", 0x69), ("dsd_le", 0x69)] {
let mut bytes = [
1, 2, 3, 4, 5, 6, 7, 8, 255, 255, 255, 255, 255, 255, 255, 255,
];
finish_transfer(&mut bytes, 8, &format(encoding, 4), &mut 5);
assert_eq!(bytes[..8], [1, 2, 3, 4, 5, 6, 7, 8]);
assert_eq!(bytes[8..], [silence; 8]);
}
}
}
@@ -0,0 +1,121 @@
//! Android-owned USB permission and a bounded, mixer-free USB audio transport.
//! No resampling or gain processing occurs at this boundary.
#[cfg(any(test, target_os = "android"))]
mod descriptors;
#[cfg(any(test, target_os = "android"))]
mod framing;
#[cfg(target_os = "android")]
#[allow(unsafe_code)] // libusb ownership is confined to this worker module.
mod transport;
use std::sync::Arc;
#[derive(Debug, thiserror::Error, uniffi::Error)]
#[uniffi(flat_error)]
pub enum UsbAudioError {
#[error("{message}")]
Failed { message: String },
}
impl From<String> for UsbAudioError {
fn from(message: String) -> Self {
Self::Failed { message }
}
}
#[derive(uniffi::Record, Clone)]
pub struct UsbOutputFormat {
pub sample_rate: u32,
pub channels: u8,
pub bits: u8,
pub subslot: u8,
/// pcm, dop, dsd_be or dsd_le. Native DSD is selected by exact device ID.
pub encoding: String,
}
#[derive(uniffi::Object)]
pub struct UsbDirectOutput {
#[cfg(target_os = "android")]
output: transport::Output,
}
#[uniffi::export]
impl UsbDirectOutput {
/// The caller retains its UsbDeviceConnection until shutdown() completes.
#[uniffi::constructor]
pub fn open(
fd: i32,
descriptors: Vec<u8>,
sample_rate: u32,
channels: u8,
bits: u8,
dsd: bool,
allow_dop: bool,
) -> Result<Arc<Self>, UsbAudioError> {
#[cfg(target_os = "android")]
{
let output = transport::Output::open(
fd,
descriptors,
sample_rate,
channels,
bits,
dsd,
allow_dop,
)?;
Ok(Arc::new(Self { output }))
}
#[cfg(not(target_os = "android"))]
{
let _ = (fd, descriptors, sample_rate, channels, bits, dsd, allow_dop);
Err("Direct USB audio requires Android".to_string().into())
}
}
pub fn format(&self) -> UsbOutputFormat {
#[cfg(target_os = "android")]
return self.output.format.clone();
#[cfg(not(target_os = "android"))]
unreachable!("Android-only constructor")
}
/// Nonblocking; returns zero when the bounded queue is full.
pub fn write(&self, data: Vec<u8>) -> Result<u32, UsbAudioError> {
#[cfg(target_os = "android")]
return self.output.write(data).map_err(Into::into);
#[cfg(not(target_os = "android"))]
{
let _ = data;
unreachable!("Android-only constructor")
}
}
pub fn start(&self) -> Result<(), UsbAudioError> {
#[cfg(target_os = "android")]
return self.output.start().map_err(Into::into);
#[cfg(not(target_os = "android"))]
unreachable!("Android-only constructor")
}
/// Cancels in-flight transfers, clears buffered audio and resets position.
pub fn flush(&self) -> Result<(), UsbAudioError> {
#[cfg(target_os = "android")]
return self.output.flush().map_err(Into::into);
#[cfg(not(target_os = "android"))]
unreachable!("Android-only constructor")
}
/// Counts completed USB frames, excluding inserted silence.
pub fn frames(&self) -> Result<u64, UsbAudioError> {
#[cfg(target_os = "android")]
return self.output.frames().map_err(Into::into);
#[cfg(not(target_os = "android"))]
unreachable!("Android-only constructor")
}
pub fn shutdown(&self) {
#[cfg(target_os = "android")]
self.output.close();
}
}
@@ -0,0 +1,584 @@
//! All raw libusb pointers belong to one thread. Transfer buffers remain at
//! stable addresses until callbacks complete, including cancellation on drop.
use super::{
UsbOutputFormat,
descriptors::{self, Alternate, Endpoint, PacketClock},
};
use libusb1_sys as usb;
use std::{
collections::VecDeque,
ptr,
sync::{Arc, Condvar, Mutex, mpsc},
thread::{self, JoinHandle},
time::{Duration, Instant},
};
#[derive(Default)]
struct State {
bytes: VecDeque<u8>,
playing: bool,
closed: bool,
frames: u64,
flush: u64,
flushed: u64,
error: Option<String>,
}
type Shared = Arc<(Mutex<State>, Condvar)>;
pub struct Output {
pub format: UsbOutputFormat,
shared: Shared,
thread: Mutex<Option<JoinHandle<()>>>,
}
impl Output {
pub fn open(
fd: i32,
raw: Vec<u8>,
rate: u32,
channels: u8,
bits: u8,
dsd: bool,
dop: bool,
) -> Result<Self, String> {
let device = descriptors::parse(&raw)?;
let candidates = descriptors::formats(&device, rate, channels, bits, dsd, dop);
if fd < 0 || candidates.len() > 32 {
return Err("Invalid USB device/configuration".into());
}
if candidates.is_empty() {
return Err("No compatible USB format".into());
}
let shared = Arc::new((Mutex::new(State::default()), Condvar::new()));
let state = shared.clone();
let (tx, rx) = mpsc::sync_channel(1);
let handle = thread::Builder::new()
.name("SpotiFLAC-USB-iso".into())
.spawn(move || {
// Best effort; failure changes scheduling only, never audio data.
unsafe {
libc::setpriority(libc::PRIO_PROCESS, 0, -16);
}
// SAF/device descriptors are opened by Android and stay owned by
// Kotlin until this worker has joined. No device discovery/root.
let result = Session::open(fd, candidates);
match result {
Ok(mut session) => {
let _ = tx.send(Ok(session.format.clone()));
if let Err(error) = session.run(&state) {
let mut s = state.0.lock().unwrap();
s.error = Some(error);
s.playing = false;
}
}
Err(error) => {
let _ = tx.send(Err(error));
}
}
state.1.notify_all();
})
.map_err(|e| e.to_string())?;
match rx.recv().map_err(|e| e.to_string())? {
Ok(format) => Ok(Self {
format,
shared,
thread: Mutex::new(Some(handle)),
}),
Err(e) => {
let _ = handle.join();
Err(e)
}
}
}
pub fn write(&self, data: Vec<u8>) -> Result<u32, String> {
let frame = self.format.subslot as usize * self.format.channels as usize;
if data.len() > 256 * 1024 || !data.len().is_multiple_of(frame) {
return Err("Invalid USB frame buffer".into());
}
let mut s = self.shared.0.lock().unwrap();
if let Some(e) = &s.error {
return Err(e.clone());
}
if s.closed {
return Err("USB output closed".into());
}
let capacity = (self.format.sample_rate as usize * frame / 5).max(256 * 1024);
let n = data.len().min(capacity.saturating_sub(s.bytes.len())) / frame * frame;
s.bytes.extend(&data[..n]);
self.shared.1.notify_all();
Ok(n as u32)
}
pub fn start(&self) -> Result<(), String> {
let mut s = self.shared.0.lock().unwrap();
if let Some(e) = &s.error {
return Err(e.clone());
}
if s.closed {
return Err("USB output closed".into());
}
s.playing = true;
self.shared.1.notify_all();
Ok(())
}
pub fn flush(&self) -> Result<(), String> {
let mut s = self.shared.0.lock().unwrap();
s.playing = false;
s.flush += 1;
self.shared.1.notify_all();
let (s, timeout) = self
.shared
.1
.wait_timeout_while(s, Duration::from_secs(2), |s| {
s.flushed != s.flush && s.error.is_none() && !s.closed
})
.unwrap();
if let Some(e) = &s.error {
return Err(e.clone());
}
if timeout.timed_out() {
return Err("USB cancellation timed out".into());
}
Ok(())
}
pub fn frames(&self) -> Result<u64, String> {
let s = self.shared.0.lock().unwrap();
if let Some(e) = &s.error {
return Err(e.clone());
}
Ok(s.frames)
}
pub fn close(&self) {
{
let mut s = self.shared.0.lock().unwrap();
s.closed = true;
self.shared.1.notify_all();
}
if let Some(t) = self.thread.lock().unwrap().take() {
let _ = t.join();
}
}
}
impl Drop for Output {
fn drop(&mut self) {
self.close();
}
}
fn check(code: i32) -> Result<(), String> {
if code < 0 {
Err(format!("USB error {code}"))
} else {
Ok(())
}
}
struct Slot {
transfer: *mut usb::libusb_transfer,
bytes: Vec<u8>,
pending: bool,
music_frames: u64,
submitted: bool,
}
extern "system" fn completed(transfer: *mut usb::libusb_transfer) {
// SAFETY: user_data points to its stable Box<Slot>; libusb callbacks are
// invoked only while this worker is handling events, before drop/free.
unsafe {
(*((*transfer).user_data as *mut Slot)).pending = false;
}
}
impl Slot {
fn new(packets: usize, max_bytes: usize) -> Result<Box<Self>, String> {
// SAFETY: libusb allocates space for exactly this many descriptors.
let transfer = unsafe { usb::libusb_alloc_transfer(packets as i32) };
if transfer.is_null() {
return Err("USB transfer allocation failed".into());
}
Ok(Box::new(Self {
transfer,
bytes: vec![0; packets * max_bytes],
pending: false,
music_frames: 0,
submitted: false,
}))
}
fn submit(
&mut self,
handle: *mut usb::libusb_device_handle,
address: u8,
lengths: &[usize],
) -> Result<(), String> {
// SAFETY: Box and Vec allocations never move/resize while pending;
// lengths are checked against the endpoint capacity before submission.
unsafe {
usb::libusb_fill_iso_transfer(
self.transfer,
handle,
address,
self.bytes.as_mut_ptr(),
lengths.iter().sum::<usize>() as i32,
lengths.len() as i32,
completed,
self as *mut Self as *mut _,
500,
);
for (i, n) in lengths.iter().enumerate() {
(*(*self.transfer).iso_packet_desc.as_mut_ptr().add(i)).length = *n as u32;
}
check(usb::libusb_submit_transfer(self.transfer))?;
}
self.pending = true;
self.submitted = true;
Ok(())
}
fn success(&self) -> bool {
unsafe {
(*self.transfer).status == 0
&& (0..(*self.transfer).num_iso_packets as usize).all(|i| {
let packet = &*(*self.transfer).iso_packet_desc.as_ptr().add(i);
packet.status == 0
&& ((*self.transfer).endpoint & 0x80 != 0
|| packet.actual_length == packet.length)
})
}
}
fn actual(&self) -> usize {
unsafe { (*(*self.transfer).iso_packet_desc.as_ptr()).actual_length as usize }
}
}
impl Drop for Slot {
fn drop(&mut self) {
assert!(
!self.pending,
"USB transfer freed before cancellation completed"
);
unsafe {
usb::libusb_free_transfer(self.transfer);
}
}
}
struct Session {
context: *mut usb::libusb_context,
handle: *mut usb::libusb_device_handle,
claimed: Vec<i32>,
alternate: Alternate,
endpoint: Endpoint,
feedback: Option<Endpoint>,
high_speed: bool,
interval: u32,
#[allow(clippy::vec_box)] // libusb user_data must survive Vec moves/reallocation.
slots: Vec<Box<Slot>>,
feedback_slot: Option<Box<Slot>>,
format: UsbOutputFormat,
}
impl Session {
fn open(fd: i32, candidates: Vec<(Alternate, UsbOutputFormat)>) -> Result<Self, String> {
let mut s = Self {
context: ptr::null_mut(),
handle: ptr::null_mut(),
claimed: vec![],
alternate: Alternate::default(),
endpoint: Endpoint::default(),
feedback: None,
high_speed: false,
interval: 0,
slots: vec![],
feedback_slot: None,
format: candidates[0].1.clone(),
};
unsafe {
// NO_DEVICE_DISCOVERY is required on Android; the USB permission
// grant is represented by the file descriptor, not /dev scanning.
check(usb::libusb_set_option(ptr::null_mut(), 2))?;
check(usb::libusb_init(&mut s.context))?;
check(usb::libusb_wrap_sys_device(
s.context,
fd as _,
&mut s.handle,
))?;
s.high_speed = usb::libusb_get_device_speed(usb::libusb_get_device(s.handle)) >= 3;
check(usb::libusb_set_auto_detach_kernel_driver(s.handle, 1))?;
}
let mut failure = "No supported USB endpoint".to_string();
for (a, f) in candidates {
if let Err(e) = s.configure(a, f) {
failure = e;
s.release();
continue;
}
let packets = (4000 / s.interval).clamp(1, 32) as usize;
for _ in 0..4 {
s.slots.push(Slot::new(packets, s.endpoint.max_packet)?);
}
if let Some(ep) = &s.feedback {
s.feedback_slot = Some(Slot::new(1, ep.max_packet)?);
}
return Ok(s);
}
Err(failure)
}
fn control(
&self,
input: bool,
request: u8,
value: u16,
index: u16,
data: &mut [u8],
endpoint: bool,
) -> Result<(), String> {
let kind = (if input { 0x80 } else { 0 }) | 0x20 | if endpoint { 2 } else { 1 };
let n = unsafe {
usb::libusb_control_transfer(
self.handle,
kind,
request,
value,
index,
data.as_mut_ptr(),
data.len() as u16,
500,
)
};
check(n)?;
if n as usize != data.len() {
return Err("Short USB control response".into());
}
Ok(())
}
fn configure(&mut self, a: Alternate, f: UsbOutputFormat) -> Result<(), String> {
let ep = a
.endpoints
.iter()
.find(|e| e.address & 0x80 == 0 && e.attributes & 3 == 1 && e.attributes & 0x30 == 0)
.ok_or("No isochronous audio output")?
.clone();
if !(1..=4).contains(&ep.interval) || ep.max_packet == 0 || ep.max_packet > 3072 {
return Err("Unsupported USB interval/packet size".into());
}
let interval = (if self.high_speed { 125 } else { 1000 }) * (1 << (ep.interval - 1));
let maximum_frames = (f.sample_rate as u64 * interval as u64).div_ceil(1_000_000) as usize;
if maximum_frames * f.channels as usize * f.subslot as usize > ep.max_packet {
return Err("USB endpoint bandwidth too small".into());
}
let feedback = a
.endpoints
.iter()
.find(|e| {
e.address & 0x80 != 0
&& e.attributes & 3 == 1
&& (e.address == ep.sync_address || e.attributes & 0x30 == 0x10)
})
.cloned();
if (ep.attributes >> 2) & 3 == 1 && feedback.is_none() {
return Err("Implicit feedback is not supported".into());
}
if let Some(e) = &feedback
&& !(3..=64).contains(&e.max_packet)
{
return Err("Invalid feedback endpoint".into());
}
for interface in [a.control, a.interface] {
if self.claimed.contains(&(interface as i32)) {
continue;
}
unsafe {
check(usb::libusb_claim_interface(self.handle, interface as i32))?;
}
self.claimed.push(interface as i32);
}
self.alternate = a.clone();
unsafe {
check(usb::libusb_set_interface_alt_setting(
self.handle,
a.interface as i32,
0,
))?;
}
if a.uac2 {
let index = ((a.clock as u16) << 8) | a.control as u16;
let mut rate = f.sample_rate.to_le_bytes();
// A read-only clock already at the exact rate needs no SET_CUR.
let mut current = [0; 4];
self.control(true, 1, 0x100, index, &mut current, false)?;
if current != rate {
self.control(false, 1, 0x100, index, &mut rate, false)?;
}
self.control(true, 1, 0x100, index, &mut current, false)?;
if current != rate {
return Err("USB clock rate could not be verified".into());
}
}
unsafe {
check(usb::libusb_set_interface_alt_setting(
self.handle,
a.interface as i32,
a.alternate as i32,
))?;
}
if !a.uac2 && ep.rate_control {
let mut rate = f.sample_rate.to_le_bytes()[..3].to_vec();
let mut current = [0; 3];
self.control(false, 1, 0x100, ep.address as u16, &mut rate, true)?;
self.control(true, 0x81, 0x100, ep.address as u16, &mut current, true)?;
if current != rate.as_slice() {
return Err("USB clock rate could not be verified".into());
}
} else if !a.uac2 && a.rates.as_slice() != [(f.sample_rate, f.sample_rate)] {
return Err("Variable USB clock has no rate control".into());
}
self.alternate = a;
self.endpoint = ep;
self.feedback = feedback;
self.interval = interval;
self.format = f;
Ok(())
}
fn events(&self) -> Result<(), String> {
let timeout = libc::timeval {
tv_sec: 0,
tv_usec: 1000,
};
check(unsafe { usb::libusb_handle_events_timeout(self.context, &timeout) })
}
fn drain(&mut self) {
for slot in self.slots.iter_mut().chain(self.feedback_slot.iter_mut()) {
if slot.pending {
unsafe {
usb::libusb_cancel_transfer(slot.transfer);
}
}
}
// Cancellation is asynchronous. Keep every buffer and the context alive
// until libusb acknowledges it, including USB detach and error paths.
while self
.slots
.iter()
.chain(self.feedback_slot.iter())
.any(|s| s.pending)
{
let _ = self.events();
}
}
fn run(&mut self, shared: &Shared) -> Result<(), String> {
let frame = self.format.subslot as usize * self.format.channels as usize;
let packets = (4000 / self.interval).clamp(1, 32) as usize;
let mut clock = PacketClock::new(self.format.sample_rate, self.interval);
let mut marker = 0x05;
let mut last_feedback = Instant::now();
let mut started = false;
loop {
self.events()?;
let mut s = shared.0.lock().unwrap();
if s.closed {
break;
}
if s.flush != s.flushed {
drop(s);
self.drain();
s = shared.0.lock().unwrap();
s.bytes.clear();
s.frames = 0;
s.flushed = s.flush;
for slot in &mut self.slots {
slot.music_frames = 0;
slot.submitted = false;
}
clock = PacketClock::new(self.format.sample_rate, self.interval);
marker = 0x05;
started = false;
shared.1.notify_all();
}
if !s.playing {
drop(shared.1.wait_timeout(s, Duration::from_millis(20)).unwrap());
continue;
}
if !started {
last_feedback = Instant::now();
started = true;
}
if let (Some(ep), Some(slot)) = (&self.feedback, &mut self.feedback_slot) {
if !slot.pending {
if slot.actual() > 0
&& slot.success()
&& clock.feedback(
&slot.bytes[..slot.actual().min(slot.bytes.len())],
self.high_speed,
self.interval,
)
{
last_feedback = Instant::now();
}
slot.submit(self.handle, ep.address, &[ep.max_packet])?;
}
if last_feedback.elapsed() > Duration::from_secs(2) {
return Err("USB feedback stopped or invalid".into());
}
}
for slot in &mut self.slots {
if slot.pending {
continue;
}
if slot.submitted {
if !slot.success() {
return Err("USB audio transfer failed".into());
}
s.frames += slot.music_frames;
slot.music_frames = 0;
}
let lengths: Vec<_> = (0..packets).map(|_| clock.next() * frame).collect();
if lengths.iter().any(|n| *n > self.endpoint.max_packet) {
return Err("USB feedback exceeded endpoint capacity".into());
}
let total: usize = lengths.iter().sum();
let available = total.min(s.bytes.len()) / frame * frame;
for b in &mut slot.bytes[..available] {
*b = s.bytes.pop_front().unwrap();
}
super::framing::finish_transfer(
&mut slot.bytes[..total],
available,
&self.format,
&mut marker,
);
slot.music_frames = (available / frame) as u64;
slot.submit(self.handle, self.endpoint.address, &lengths)?;
}
}
self.drain();
Ok(())
}
fn release(&mut self) {
if self.handle.is_null() {
return;
}
unsafe {
if self.alternate.alternate > 0 {
usb::libusb_set_interface_alt_setting(
self.handle,
self.alternate.interface as i32,
0,
);
}
for interface in self.claimed.drain(..).rev() {
usb::libusb_release_interface(self.handle, interface);
}
}
self.alternate = Alternate::default();
}
}
impl Drop for Session {
fn drop(&mut self) {
self.drain();
self.slots.clear();
self.feedback_slot = None;
self.release();
unsafe {
if !self.handle.is_null() {
usb::libusb_close(self.handle);
}
if !self.context.is_null() {
usb::libusb_exit(self.context);
}
}
}
}