1use std::time::Duration;
15use thiserror::Error;
16use tokio::io::{AsyncReadExt, AsyncWriteExt};
17use tokio::net::UnixStream;
18
19pub const CTRL_PATH: &str = "/tmp/audio.a2dp_ctrl";
20pub const DATA_PATH: &str = "/tmp/audio.a2dp_data";
21
22const CMD_CHECK_READY: u8 = 0x01;
24const CMD_START: u8 = 0x02;
25const CMD_STOP: u8 = 0x03;
26const ACK_SUCCESS: u8 = 0x00;
28const BTSERVICE_TIMEOUT_SECS: u64 = 3;
31
32#[derive(Debug, Error)]
33pub enum A2dpError {
34 #[error("connect {path}: {source}")]
35 Connect {
36 path: String,
37 source: std::io::Error,
38 },
39 #[error("ctrl command {cmd:#04x} ack {ack:#04x} (expected {expected:#04x})")]
40 BadAck { cmd: u8, ack: u8, expected: u8 },
41 #[error("io: {0}")]
42 Io(#[from] std::io::Error),
43}
44
45pub const TARGET_RATE: usize = 44_100;
48pub const TARGET_CHANNELS: usize = 2;
49
50pub struct AospA2dpSink {
51 ctrl: UnixStream,
52 data: Option<UnixStream>,
53 started: bool,
54}
55
56impl AospA2dpSink {
57 pub async fn open() -> Result<Self, A2dpError> {
60 let ctrl = UnixStream::connect(CTRL_PATH)
61 .await
62 .map_err(|e| A2dpError::Connect {
63 path: CTRL_PATH.into(),
64 source: e,
65 })?;
66 let mut s = AospA2dpSink {
67 ctrl,
68 data: None,
69 started: false,
70 };
71 s.ctrl_cmd(CMD_CHECK_READY, "CHECK_READY").await?;
72 Ok(s)
73 }
74
75 async fn ctrl_cmd(&mut self, cmd: u8, _name: &str) -> Result<u8, A2dpError> {
76 self.ctrl.write_all(&[cmd]).await?;
77 let mut ack = [0u8; 1];
78 match tokio::time::timeout(
81 Duration::from_secs(BTSERVICE_TIMEOUT_SECS),
82 self.ctrl.read_exact(&mut ack),
83 )
84 .await
85 {
86 Ok(Ok(_)) => {}
87 Ok(Err(e)) => return Err(A2dpError::Io(e)),
88 Err(_) => {
89 return Err(A2dpError::Io(std::io::Error::new(
90 std::io::ErrorKind::TimedOut,
91 "ctrl ack timeout",
92 )))
93 }
94 }
95 if (cmd == CMD_CHECK_READY || cmd == CMD_START) && ack[0] != ACK_SUCCESS {
96 return Err(A2dpError::BadAck {
97 cmd,
98 ack: ack[0],
99 expected: ACK_SUCCESS,
100 });
101 }
102 Ok(ack[0])
103 }
104
105 pub async fn start(&mut self) -> Result<(), A2dpError> {
107 self.ctrl_cmd(CMD_START, "START").await?;
108 self.data = Some(
109 UnixStream::connect(DATA_PATH)
110 .await
111 .map_err(|e| A2dpError::Connect {
112 path: DATA_PATH.into(),
113 source: e,
114 })?,
115 );
116 self.started = true;
117 Ok(())
118 }
119
120 pub async fn write_pcm(&mut self, pcm: &[u8]) -> Result<(), A2dpError> {
123 match self.data.as_mut() {
124 Some(d) => tokio::time::timeout(
125 Duration::from_secs(BTSERVICE_TIMEOUT_SECS),
126 d.write_all(pcm),
127 )
128 .await
129 .map_err(|_| {
130 A2dpError::Io(std::io::Error::new(
131 std::io::ErrorKind::TimedOut,
132 "a2dp data write timeout (btservice not draining)",
133 ))
134 })?
135 .map_err(A2dpError::Io),
136 None => Err(A2dpError::Io(std::io::Error::new(
137 std::io::ErrorKind::NotConnected,
138 "a2dp data socket not open (call start() first)",
139 ))),
140 }
141 }
142
143 pub fn unsent_bytes(&self) -> usize {
147 use std::os::unix::io::AsRawFd;
148 if let Some(ref data) = self.data {
149 let fd = data.as_raw_fd();
150 let mut value: std::ffi::c_int = 0;
151 unsafe {
155 libc::ioctl(fd, libc::TIOCOUTQ, &mut value);
156 }
157 value.max(0) as usize
158 } else {
159 0
160 }
161 }
162
163 pub async fn stop(&mut self) -> Result<(), A2dpError> {
165 if self.started {
166 log::warn!("a2dp: best-effort STOP (error ignored if stream already closing)");
168 let _ = self.ctrl_cmd(CMD_STOP, "STOP").await;
169 self.started = false;
170 }
171 Ok(())
172 }
173}
174
175impl Drop for AospA2dpSink {
176 fn drop(&mut self) {
177 use std::os::unix::net::UnixStream as StdStream;
179 if self.started {
180 if let Ok(mut c) = StdStream::connect(CTRL_PATH) {
181 use std::io::Write;
182 if let Err(e) = c.write_all(&[CMD_STOP]) {
185 log::warn!("a2dp drop STOP write failed: {e}");
186 }
187 }
188 }
189 }
190}