返回 CodeWhale
audio.rs
根目录 / crates / tui / src / tui / pet_watch / audio.rs
1 //! A bounded presentation sink for the shared core's stereo PCM. No score,
2 //! event interpretation or simulation clock lives in the audio player.
3 use std::io::{self, Write};
4 use std::process::{Child, Command, Stdio};
5 use std::sync::atomic::{AtomicBool, Ordering};
6 use std::sync::{Arc, Mutex, mpsc};
7 use std::time::{Duration, Instant};
8
9 pub const SAMPLE_RATE: usize = 48_000;
10 pub const MAX_FRAMES: usize = SAMPLE_RATE / 2;
11
12 pub(super) struct Packet {
13 pub(super) bytes: Vec<u8>,
14 created: Instant,
15 }
16
17 #[derive(Default)]
18 struct Control {
19 child: Mutex<Option<Child>>,
20 cancelled: AtomicBool,
21 failed: AtomicBool,
22 }
23
24 #[derive(Clone)]
25 pub struct Target {
26 tx: mpsc::SyncSender<Packet>,
27 control: Arc<Control>,
28 requested: Instant,
29 }
30
31 impl Target {
32 pub fn same_stream(&self, other: &Self) -> bool {
33 Arc::ptr_eq(&self.control, &other.control)
34 }
35
36 pub fn active(&self) -> bool {
37 !self.control.cancelled.load(Ordering::Acquire)
38 && !self.control.failed.load(Ordering::Acquire)
39 }
40
41 pub fn fail(&self) {
42 self.control.failed.store(true, Ordering::Release);
43 }
44
45 pub fn current(&self) -> bool {
46 self.requested.elapsed() <= Duration::from_millis(500)
47 }
48
49 /// `Ok` means the samples were valid and the stream is live, not that they
50 /// were played: a full queue or a stale request drops them. Live audio is
51 /// a bounded real-time sink — late samples never queue behind the world
52 /// clock — so an overrun is a dropout, not a stream failure.
53 pub fn send(&self, channels: [Vec<f32>; 2]) -> Result<(), ()> {
54 let length = channels[0].len();
55 if length > MAX_FRAMES
56 || channels[1].len() != length
57 || channels
58 .iter()
59 .flatten()
60 .any(|s| !s.is_finite() || s.abs() > 1.0)
61 {
62 self.fail();
63 return Err(());
64 }
65 if !self.active() {
66 return Err(());
67 }
68 if !self.current() {
69 return Ok(());
70 }
71 let mut bytes = Vec::with_capacity(length * 8);
72 for (left, right) in channels[0].iter().zip(&channels[1]) {
73 bytes.extend_from_slice(&left.to_le_bytes());
74 bytes.extend_from_slice(&right.to_le_bytes());
75 }
76 match self.tx.try_send(Packet {
77 bytes,
78 created: self.requested,
79 }) {
80 Ok(()) | Err(mpsc::TrySendError::Full(_)) => Ok(()),
81 Err(mpsc::TrySendError::Disconnected(_)) => {
82 self.fail();
83 Err(())
84 }
85 }
86 }
87 }
88
89 pub struct Output {
90 target: Target,
91 }
92
93 impl Output {
94 #[cfg(test)]
95 pub(super) fn capture() -> (Self, mpsc::Receiver<Packet>) {
96 let (tx, rx) = mpsc::sync_channel(4);
97 (
98 Self {
99 target: Target {
100 tx,
101 control: Arc::new(Control::default()),
102 requested: Instant::now(),
103 },
104 },
105 rx,
106 )
107 }
108
109 #[cfg(not(test))]
110 pub fn start() -> io::Result<Self> {
111 Self::spawn(player_command()?)
112 }
113
114 #[cfg(test)]
115 pub fn start() -> io::Result<Self> {
116 // Product/library tests never open the user's audio output.
117 Err(io::Error::new(
118 io::ErrorKind::Unsupported,
119 "audio disabled in tests",
120 ))
121 }
122
123 fn spawn(mut command: Command) -> io::Result<Self> {
124 let (tx, rx) = mpsc::sync_channel(4);
125 let control = Arc::new(Control::default());
126 let thread_control = Arc::clone(&control);
127 std::thread::Builder::new()
128 .name("pet-audio".into())
129 .spawn(move || {
130 let result = play(&mut command, rx, &thread_control);
131 if result.is_err() && !thread_control.cancelled.load(Ordering::Acquire) {
132 thread_control.failed.store(true, Ordering::Release);
133 }
134 let child = thread_control
135 .child
136 .lock()
137 .ok()
138 .and_then(|mut slot| slot.take());
139 if let Some(mut child) = child {
140 let _ = child.kill();
141 let _ = child.wait();
142 }
143 })?;
144 Ok(Self {
145 target: Target {
146 tx,
147 control,
148 requested: Instant::now(),
149 },
150 })
151 }
152
153 pub fn target(&self) -> Target {
154 Target {
155 requested: Instant::now(),
156 ..self.target.clone()
157 }
158 }
159 pub fn failed(&self) -> bool {
160 self.target.control.failed.load(Ordering::Acquire)
161 }
162 }
163
164 impl Drop for Output {
165 fn drop(&mut self) {
166 self.target.control.cancelled.store(true, Ordering::Release);
167 // Closing or muting Watch interrupts even a blocked pipe write. Reaping
168 // happens in the audio thread, never in the terminal event loop.
169 if let Ok(mut slot) = self.target.control.child.lock()
170 && let Some(child) = slot.as_mut()
171 {
172 let _ = child.kill();
173 }
174 }
175 }
176
177 /// The pet's PCM player: `ffplay` from a fixed install prefix, never from the
178 /// ambient `PATH` (see [`crate::notify::audio::trusted_system_executable`]). A player
179 /// that is not installed there fails the start and Watch reports audio as
180 /// unavailable.
181 fn player_command() -> io::Result<Command> {
182 let mut command = Command::new(crate::notify::audio::trusted_system_executable("ffplay")?);
183 command.args([
184 "-nodisp",
185 "-autoexit",
186 "-loglevel",
187 "error",
188 "-probesize",
189 "32",
190 "-analyzeduration",
191 "0",
192 "-f",
193 "f32le",
194 "-sample_rate",
195 "48000",
196 "-ch_layout",
197 "stereo",
198 "-i",
199 "pipe:0",
200 ]);
201 #[cfg(windows)]
202 {
203 use std::os::windows::process::CommandExt;
204 command.creation_flags(0x08000000); // CREATE_NO_WINDOW
205 }
206 Ok(command)
207 }
208
209 fn play(command: &mut Command, rx: mpsc::Receiver<Packet>, control: &Control) -> io::Result<()> {
210 if control.cancelled.load(Ordering::Acquire) {
211 return Ok(());
212 }
213 let mut child = command
214 .stdin(Stdio::piped())
215 .stdout(Stdio::null())
216 .stderr(Stdio::null())
217 .spawn()?;
218 let Some(mut input) = child.stdin.take() else {
219 let _ = child.kill();
220 let _ = child.wait();
221 return Err(io::Error::other("audio pipe unavailable"));
222 };
223 match control.child.lock() {
224 Ok(mut slot) => *slot = Some(child),
225 Err(_) => {
226 let _ = child.kill();
227 let _ = child.wait();
228 return Err(io::Error::other("audio lock failed"));
229 }
230 }
231 while !control.cancelled.load(Ordering::Acquire) && !control.failed.load(Ordering::Acquire) {
232 if control
233 .child
234 .lock()
235 .map_err(|_| io::Error::other("audio lock failed"))?
236 .as_mut()
237 .is_none_or(|c| c.try_wait().map_or(true, |status| status.is_some()))
238 {
239 return Err(io::Error::other("audio player exited"));
240 }
241 match rx.recv_timeout(Duration::from_millis(50)) {
242 Ok(packet) => {
243 if packet.created.elapsed() > Duration::from_millis(500) {
244 continue;
245 }
246 input.write_all(&packet.bytes)?;
247 }
248 Err(mpsc::RecvTimeoutError::Timeout) => {}
249 Err(mpsc::RecvTimeoutError::Disconnected) => break,
250 }
251 }
252 Ok(())
253 }
254
255 #[cfg(test)]
256 mod tests {
257 use super::*;
258
259 fn wait_until(mut ready: impl FnMut() -> bool) {
260 let until = Instant::now() + Duration::from_secs(5);
261 while !ready() {
262 assert!(Instant::now() < until, "audio process did not settle");
263 std::thread::sleep(Duration::from_millis(5));
264 }
265 }
266
267 fn receiver() -> (Target, mpsc::Receiver<Packet>) {
268 let (tx, rx) = mpsc::sync_channel(4);
269 (
270 Target {
271 tx,
272 control: Arc::new(Control::default()),
273 requested: Instant::now(),
274 },
275 rx,
276 )
277 }
278
279 #[test]
280 fn pcm_boundary_preserves_samples_rejects_invalid_input_and_discards_backlog() {
281 let (target, rx) = receiver();
282 target.send([vec![0.125, -0.5], vec![0.25, 0.75]]).unwrap();
283 let bytes: Vec<_> = [0.125_f32, 0.25, -0.5, 0.75]
284 .into_iter()
285 .flat_map(f32::to_le_bytes)
286 .collect();
287 assert_eq!(rx.recv().unwrap().bytes, bytes);
288 for channels in [
289 [vec![0.0], vec![]],
290 [vec![f32::NAN], vec![0.0]],
291 [vec![1.01], vec![0.0]],
292 [vec![0.0; MAX_FRAMES + 1], vec![0.0; MAX_FRAMES + 1]],
293 ] {
294 let (target, rx) = receiver();
295 assert!(target.send(channels).is_err());
296 assert!(!target.active());
297 assert!(rx.try_recv().is_err());
298 }
299 let (mut target, rx) = receiver();
300 target.requested = Instant::now() - Duration::from_secs(1);
301 assert!(target.send([vec![0.0], vec![0.0]]).is_ok());
302 assert!(target.active());
303 assert!(rx.try_recv().is_err());
304 let (target, rx) = receiver();
305 for _ in 0..4 {
306 target.send([vec![0.0], vec![0.0]]).unwrap();
307 }
308 assert!(target.send([vec![0.0], vec![0.0]]).is_ok());
309 assert!(target.active());
310 assert_eq!(rx.try_iter().count(), 4);
311 }
312
313 /// Only a test subprocess with this explicit environment enters the sink.
314 /// Ordinary library/nextest runs return without launching or playing audio.
315 #[test]
316 fn pet_audio_fake_player() {
317 let Some(path) = std::env::var_os("CODEWHALE_TEST_PET_AUDIO_OUTPUT") else {
318 return;
319 };
320 let mut file = std::fs::File::create(path).unwrap();
321 if std::env::var_os("CODEWHALE_TEST_PET_AUDIO_HOLD").is_some() {
322 std::thread::sleep(Duration::from_secs(30));
323 } else {
324 std::io::copy(&mut std::io::stdin().lock(), &mut file).unwrap();
325 }
326 }
327
328 fn fake_player(path: &std::path::Path, hold: bool) -> Command {
329 let module = module_path!().split_once("::").unwrap().1;
330 let mut command = Command::new(std::env::current_exe().unwrap());
331 command
332 .args([
333 "--exact",
334 &format!("{module}::pet_audio_fake_player"),
335 "--nocapture",
336 ])
337 .env("CODEWHALE_TEST_PET_AUDIO_OUTPUT", path);
338 if hold {
339 command.env("CODEWHALE_TEST_PET_AUDIO_HOLD", "1");
340 }
341 command
342 }
343
344 #[test]
345 fn muting_reaps_the_player_even_with_a_retained_target_and_blocked_pipe() {
346 for hold in [false, true] {
347 let dir = tempfile::tempdir().unwrap();
348 let path = dir.path().join("received.f32");
349 let output = Output::spawn(fake_player(&path, hold)).unwrap();
350 wait_until(|| path.exists());
351 let target = output.target();
352 let control = Arc::clone(&target.control);
353 target
354 .send([vec![0.125; MAX_FRAMES], vec![-0.25; MAX_FRAMES]])
355 .unwrap();
356 if !hold {
357 wait_until(|| std::fs::metadata(&path).unwrap().len() == (MAX_FRAMES * 8) as u64);
358 }
359 drop(output);
360 assert!(!target.active());
361 assert!(target.send([vec![0.0], vec![0.0]]).is_err());
362 // Only this test and the deliberately retained target remain. The
363 // player thread must have returned after killing and reaping it.
364 wait_until(|| Arc::strong_count(&control) == 2);
365 assert!(control.child.lock().unwrap().is_none());
366 }
367 }
368
369 /// U06-06: the player is an absolute path from a trusted prefix or it is
370 /// refused — never a bare name the ambient `PATH` resolves, where an empty
371 /// or repository-local entry would run a workspace-planted `ffplay`.
372 #[test]
373 fn pet_player_is_never_resolved_through_path() {
374 match player_command() {
375 Ok(command) => assert!(
376 std::path::Path::new(command.get_program()).is_absolute(),
377 "player resolved through PATH: {:?}",
378 command.get_program()
379 ),
380 Err(error) => assert_eq!(error.kind(), io::ErrorKind::NotFound),
381 }
382 }
383
384 #[test]
385 fn missing_player_fails_without_opening_a_device_or_blocking_the_caller() {
386 assert!(Output::start().is_err());
387 let dir = tempfile::tempdir().unwrap();
388 let output = Output::spawn(Command::new(dir.path().join("missing-player"))).unwrap();
389 wait_until(|| output.failed());
390 assert!(!output.target().active());
391 }
392 }
393
393 lines RUST