Skip to repository content1039 lines · 40.6 KB · rust
tenant.openagents/omega
No repository description is available.
OpenAgents Git authority 2026-07-28T04:09:47.630Z Public web read
NIP-34 coordinate
30617:7649603503856e5148d571eac2766b288a8ff1e9e35d380337a1d2b0015b4f92:omegaMaintainersHidden in public view
References2 branches · 1 tag
Read-only clone
git clone https://openagents.com/git/tenant.openagents/omega.gitBrowse files
playback.rs
1use anyhow::{Context as _, Result};
2
3use audio::{AudioSettings, CHANNEL_COUNT, SAMPLE_RATE};
4use cpal::DeviceId;
5use cpal::traits::{DeviceTrait, StreamTrait as _};
6use futures::channel::mpsc::Sender;
7use futures::{Stream, StreamExt as _};
8use gpui::{
9 AsyncApp, BackgroundExecutor, Priority, ScreenCaptureFrame, ScreenCaptureSource,
10 ScreenCaptureStream, Task,
11};
12use libwebrtc::native::{apm, audio_mixer, audio_resampler};
13use livekit::track;
14
15use livekit::webrtc::{
16 audio_frame::AudioFrame,
17 audio_source::{AudioSourceOptions, RtcAudioSource, native::NativeAudioSource},
18 audio_stream::native::NativeAudioStream,
19 video_frame::{VideoBuffer, VideoFrame, VideoRotation},
20 video_source::{RtcVideoSource, VideoResolution, native::NativeVideoSource},
21 video_stream::native::NativeVideoStream,
22};
23use log::info;
24use parking_lot::Mutex;
25use serde::{Deserialize, Serialize};
26use settings::Settings;
27use std::cell::RefCell;
28use std::sync::Weak;
29use std::sync::atomic::{AtomicI32, AtomicU64, Ordering};
30use std::time::{Duration, Instant};
31use std::{borrow::Cow, collections::VecDeque, sync::Arc};
32use util::{ResultExt as _, maybe};
33
34struct TimestampedFrame {
35 frame: AudioFrame<'static>,
36 captured_at: Instant,
37}
38
39pub(crate) struct AudioStack {
40 executor: BackgroundExecutor,
41 apm: Arc<Mutex<apm::AudioProcessingModule>>,
42 mixer: Arc<Mutex<audio_mixer::AudioMixer>>,
43 _output_task: RefCell<Weak<Task<()>>>,
44 next_ssrc: AtomicI32,
45}
46
47impl AudioStack {
48 pub(crate) fn new(executor: BackgroundExecutor) -> Self {
49 // AGC2's `adaptive_digital` is what actually levels speech toward a target;
50 // the `gain_controller2.enabled` master switch alone leaves it off, which
51 // historically meant capture was effectively unleveled. Defaults match
52 // what Chrome/Meet ship with -- in particular `max_gain_db = 50` paired
53 // with `max_output_noise_level_dbfs = -50`, which lets the AGC reach
54 // very quiet talkers while the noise-level estimator backs off before
55 // boosting amplifies the noise floor.
56 let apm = Arc::new(Mutex::new(apm::AudioProcessingModule::new(
57 apm::AudioProcessingConfig {
58 echo_canceller_enabled: true,
59 gain_controller2: apm::GainController2Config {
60 enabled: true,
61 adaptive_digital: apm::AdaptiveDigitalConfig {
62 enabled: true,
63 ..Default::default()
64 },
65 ..Default::default()
66 },
67 high_pass_filter_enabled: true,
68 noise_suppression_enabled: true,
69 },
70 )));
71 let mixer = Arc::new(Mutex::new(audio_mixer::AudioMixer::new()));
72 Self {
73 executor,
74 apm,
75 mixer,
76 _output_task: RefCell::new(Weak::new()),
77 next_ssrc: AtomicI32::new(1),
78 }
79 }
80
81 pub(crate) fn play_remote_audio_track(
82 &self,
83 track: &livekit::track::RemoteAudioTrack,
84 output_audio_device: Option<DeviceId>,
85 ) -> AudioStream {
86 let output_task = self.start_output(output_audio_device);
87
88 let next_ssrc = self.next_ssrc.fetch_add(1, Ordering::Relaxed);
89 let source = AudioMixerSource {
90 ssrc: next_ssrc,
91 sample_rate: SAMPLE_RATE.get(),
92 num_channels: CHANNEL_COUNT.get() as u32,
93 buffer: Arc::default(),
94 };
95 self.mixer.lock().add_source(source.clone());
96
97 let mut stream = NativeAudioStream::new(
98 track.rtc_track(),
99 source.sample_rate as i32,
100 source.num_channels as i32,
101 );
102
103 let receive_task = self.executor.spawn_with_priority(Priority::RealtimeAudio, {
104 let source = source.clone();
105 async move {
106 while let Some(frame) = stream.next().await {
107 source.receive(frame);
108 }
109 }
110 });
111
112 let mixer = self.mixer.clone();
113 let on_drop = util::defer(move || {
114 mixer.lock().remove_source(source.ssrc);
115 drop(receive_task);
116 drop(output_task);
117 });
118
119 AudioStream::Output {
120 _drop: Box::new(on_drop),
121 }
122 }
123
124 fn start_output(&self, output_audio_device: Option<DeviceId>) -> Arc<Task<()>> {
125 if let Some(task) = self._output_task.borrow().upgrade() {
126 return task;
127 }
128 let task = Arc::new(self.executor.spawn({
129 let apm = self.apm.clone();
130 let mixer = self.mixer.clone();
131 let executor = self.executor.clone();
132 async move {
133 Self::play_output(
134 executor,
135 apm,
136 mixer,
137 SAMPLE_RATE.get(),
138 CHANNEL_COUNT.get().into(),
139 output_audio_device,
140 )
141 .await
142 .log_err();
143 }
144 }));
145 *self._output_task.borrow_mut() = Arc::downgrade(&task);
146 task
147 }
148
149 pub(crate) fn capture_local_microphone_track(
150 &self,
151 user_name: String,
152 is_staff: bool,
153 cx: &AsyncApp,
154 ) -> Result<(crate::LocalAudioTrack, AudioStream, Arc<AtomicU64>)> {
155 let source = NativeAudioSource::new(
156 // n.b. this struct's options are always ignored, noise cancellation is provided by apm.
157 AudioSourceOptions::default(),
158 SAMPLE_RATE.get(),
159 CHANNEL_COUNT.get().into(),
160 10,
161 );
162
163 let speaker = Speaker {
164 name: user_name,
165 is_staff,
166 };
167 log::info!("Microphone speaker: {speaker:?}");
168 let track_name = serde_urlencoded::to_string(speaker)
169 .context("Could not encode user information in track name")?;
170
171 let track = track::LocalAudioTrack::create_audio_track(
172 &track_name,
173 RtcAudioSource::Native(source.clone()),
174 );
175
176 let apm = self.apm.clone();
177
178 let input_lag_us = Arc::new(AtomicU64::new(0));
179 let (frame_tx, mut frame_rx) = futures::channel::mpsc::channel::<TimestampedFrame>(1);
180 let transmit_task = self.executor.spawn_with_priority(Priority::RealtimeAudio, {
181 let input_lag_us = input_lag_us.clone();
182 async move {
183 while let Some(timestamped) = frame_rx.next().await {
184 let lag = timestamped.captured_at.elapsed();
185 input_lag_us.store(lag.as_micros() as u64, Ordering::Relaxed);
186 source.capture_frame(×tamped.frame).await.log_err();
187 }
188 }
189 });
190 let capture_task = {
191 let input_audio_device =
192 AudioSettings::try_read_global(cx, |settings| settings.input_audio_device.clone())
193 .flatten();
194 let executor = self.executor.clone();
195 self.executor.spawn(async move {
196 Self::capture_input(
197 executor,
198 apm,
199 frame_tx,
200 SAMPLE_RATE.get(), // TODO(audio): was legacy removed for now
201 CHANNEL_COUNT.get().into(),
202 input_audio_device,
203 )
204 .await
205 })
206 };
207
208 let on_drop = util::defer(|| {
209 drop(transmit_task);
210 drop(capture_task);
211 });
212 Ok((
213 super::LocalAudioTrack(track),
214 AudioStream::Output {
215 _drop: Box::new(on_drop),
216 },
217 input_lag_us,
218 ))
219 }
220
221 async fn play_output(
222 executor: BackgroundExecutor,
223 apm: Arc<Mutex<apm::AudioProcessingModule>>,
224 mixer: Arc<Mutex<audio_mixer::AudioMixer>>,
225 sample_rate: u32,
226 _num_channels: u32,
227 output_audio_device: Option<DeviceId>,
228 ) -> Result<()> {
229 // Prevent App Nap from throttling audio playback on macOS.
230 // This guard is held for the entire duration of audio output.
231 #[cfg(target_os = "macos")]
232 let _prevent_app_nap = PreventAppNapGuard::new();
233
234 loop {
235 let mut device_change_listener = DeviceChangeListener::new(false)?;
236 let (output_device, output_config) =
237 crate::default_device(false, output_audio_device.as_ref())?;
238 info!("Output config: {output_config:?}");
239 let (end_on_drop_tx, end_on_drop_rx) = std::sync::mpsc::channel::<()>();
240 let mixer = mixer.clone();
241 let apm = apm.clone();
242 let mut resampler = audio_resampler::AudioResampler::default();
243 let mut buf = Vec::new();
244
245 executor
246 .spawn_with_priority(Priority::RealtimeAudio, async move {
247 let output_stream = output_device.build_output_stream(
248 &output_config.config(),
249 {
250 move |mut data, _info| {
251 while data.len() > 0 {
252 if data.len() <= buf.len() {
253 let rest = buf.split_off(data.len());
254 data.copy_from_slice(&buf);
255 buf = rest;
256 return;
257 }
258 if buf.len() > 0 {
259 let (prefix, suffix) = data.split_at_mut(buf.len());
260 prefix.copy_from_slice(&buf);
261 data = suffix;
262 }
263
264 let mut mixer = mixer.lock();
265 let mixed = mixer.mix(output_config.channels() as usize);
266 let sampled = resampler.remix_and_resample(
267 mixed,
268 sample_rate / 100,
269 // We need to assume output number of channels as otherwise we will
270 // crash in process_reverse_stream otherwise as livekit's audio resampler
271 // does not seem to support non-matching channel counts.
272 // NOTE: you can verify this by debug printing buf.len() after this stage.
273 // For 2->4 channel upmix, we should see buf.len=1920, buf we get only 960.
274 output_config.channels() as u32,
275 sample_rate,
276 output_config.channels() as u32,
277 output_config.sample_rate(),
278 );
279 buf = sampled.to_vec();
280 apm.lock()
281 .process_reverse_stream(
282 &mut buf,
283 output_config.sample_rate() as i32,
284 output_config.channels() as i32,
285 )
286 .ok();
287 }
288 }
289 },
290 |error| log::error!("error playing audio track: {:?}", error),
291 Some(Duration::from_millis(100)),
292 );
293
294 let Some(output_stream) = output_stream.log_err() else {
295 return;
296 };
297
298 output_stream.play().log_err();
299 // Block forever to keep the output stream alive
300 end_on_drop_rx.recv().ok();
301 })
302 .detach();
303
304 device_change_listener.next().await;
305 drop(end_on_drop_tx)
306 }
307 }
308
309 async fn capture_input(
310 executor: BackgroundExecutor,
311 apm: Arc<Mutex<apm::AudioProcessingModule>>,
312 frame_tx: Sender<TimestampedFrame>,
313 sample_rate: u32,
314 num_channels: u32,
315 input_audio_device: Option<DeviceId>,
316 ) -> Result<()> {
317 loop {
318 let mut device_change_listener = DeviceChangeListener::new(true)?;
319 let (device, config) = crate::default_device(true, input_audio_device.as_ref())?;
320 let (end_on_drop_tx, end_on_drop_rx) = std::sync::mpsc::channel::<()>();
321 let apm = apm.clone();
322 let mut frame_tx = frame_tx.clone();
323 let mut resampler = audio_resampler::AudioResampler::default();
324
325 executor
326 .spawn_with_priority(Priority::RealtimeAudio, async move {
327 maybe!({
328 if let Some(desc) = device.description().ok() {
329 log::info!("Using microphone: {}", desc.name())
330 } else {
331 log::info!("Using microphone: <unknown>");
332 }
333
334 let ten_ms_buffer_size =
335 (config.channels() as u32 * config.sample_rate() / 100) as usize;
336 let mut buf: Vec<i16> = Vec::with_capacity(ten_ms_buffer_size);
337
338 let stream = device
339 .build_input_stream_raw(
340 &config.config(),
341 config.sample_format(),
342 move |data, _: &_| {
343 let captured_at = Instant::now();
344 let data = crate::get_sample_data(config.sample_format(), data)
345 .log_err();
346 let Some(data) = data else {
347 return;
348 };
349 let mut data = data.as_slice();
350
351 while data.len() > 0 {
352 let remainder =
353 (buf.capacity() - buf.len()).min(data.len());
354 buf.extend_from_slice(&data[..remainder]);
355 data = &data[remainder..];
356
357 if buf.capacity() == buf.len() {
358 let mut sampled = resampler
359 .remix_and_resample(
360 buf.as_slice(),
361 config.sample_rate() / 100,
362 config.channels() as u32,
363 config.sample_rate(),
364 num_channels,
365 sample_rate,
366 )
367 .to_owned();
368
369 apm.lock()
370 .process_stream(
371 &mut sampled,
372 sample_rate as i32,
373 num_channels as i32,
374 )
375 .log_err();
376 buf.clear();
377
378 frame_tx
379 .try_send(TimestampedFrame {
380 frame: AudioFrame {
381 data: Cow::Owned(sampled),
382 sample_rate,
383 num_channels,
384 samples_per_channel: sample_rate / 100,
385 },
386 captured_at,
387 })
388 .ok();
389 }
390 }
391 },
392 |err| log::error!("error capturing audio track: {:?}", err),
393 Some(Duration::from_millis(100)),
394 )
395 .context("failed to build input stream")?;
396
397 stream.play()?;
398 // Keep the thread alive and holding onto the `stream`
399 end_on_drop_rx.recv().ok();
400 anyhow::Ok(Some(()))
401 })
402 .log_err();
403 })
404 .detach();
405
406 device_change_listener.next().await;
407 drop(end_on_drop_tx)
408 }
409 }
410}
411
412#[derive(Serialize, Deserialize, Debug)]
413pub struct Speaker {
414 pub name: String,
415 pub is_staff: bool,
416}
417
418use super::LocalVideoTrack;
419
420pub enum AudioStream {
421 Input { _task: Task<()> },
422 Output { _drop: Box<dyn std::any::Any> },
423}
424
425pub(crate) async fn capture_local_video_track(
426 capture_source: &dyn ScreenCaptureSource,
427 cx: &mut gpui::AsyncApp,
428) -> Result<(crate::LocalVideoTrack, Box<dyn ScreenCaptureStream>)> {
429 let metadata = capture_source.metadata()?;
430 let track_source = gpui_tokio::Tokio::spawn(cx, async move {
431 NativeVideoSource::new(
432 VideoResolution {
433 width: metadata.resolution.width.0 as u32,
434 height: metadata.resolution.height.0 as u32,
435 },
436 true,
437 )
438 })
439 .await?;
440
441 let capture_stream = capture_source
442 .stream(cx.foreground_executor(), {
443 let track_source = track_source.clone();
444 Box::new(move |frame| {
445 if let Some(buffer) = video_frame_buffer_to_webrtc(frame) {
446 track_source.capture_frame(&VideoFrame {
447 rotation: VideoRotation::VideoRotation0,
448 timestamp_us: 0,
449 buffer,
450 });
451 }
452 })
453 })
454 .await??;
455
456 Ok((
457 LocalVideoTrack(track::LocalVideoTrack::create_video_track(
458 "screen share",
459 RtcVideoSource::Native(track_source),
460 )),
461 capture_stream,
462 ))
463}
464
465#[derive(Clone)]
466struct AudioMixerSource {
467 ssrc: i32,
468 sample_rate: u32,
469 num_channels: u32,
470 buffer: Arc<Mutex<VecDeque<Vec<i16>>>>,
471}
472
473impl AudioMixerSource {
474 fn receive(&self, frame: AudioFrame) {
475 assert_eq!(
476 frame.data.len() as u32,
477 self.sample_rate * self.num_channels / 100
478 );
479
480 let mut buffer = self.buffer.lock();
481 buffer.push_back(frame.data.to_vec());
482 while buffer.len() > 10 {
483 buffer.pop_front();
484 }
485 }
486}
487
488impl libwebrtc::native::audio_mixer::AudioMixerSource for AudioMixerSource {
489 fn ssrc(&self) -> i32 {
490 self.ssrc
491 }
492
493 fn preferred_sample_rate(&self) -> u32 {
494 self.sample_rate
495 }
496
497 fn get_audio_frame_with_info<'a>(&self, target_sample_rate: u32) -> Option<AudioFrame<'_>> {
498 assert_eq!(self.sample_rate, target_sample_rate);
499 let buf = self.buffer.lock().pop_front()?;
500 Some(AudioFrame {
501 data: Cow::Owned(buf),
502 sample_rate: self.sample_rate,
503 num_channels: self.num_channels,
504 samples_per_channel: self.sample_rate / 100,
505 })
506 }
507}
508
509pub fn play_remote_video_track(
510 track: &crate::RemoteVideoTrack,
511 executor: &BackgroundExecutor,
512) -> impl Stream<Item = RemoteVideoFrame> + use<> {
513 #[cfg(target_os = "macos")]
514 {
515 _ = executor;
516 let mut pool = None;
517 let most_recent_frame_size = (0, 0);
518 NativeVideoStream::new(track.0.rtc_track()).filter_map(move |frame| {
519 if pool == None
520 || most_recent_frame_size != (frame.buffer.width(), frame.buffer.height())
521 {
522 pool = create_buffer_pool(frame.buffer.width(), frame.buffer.height()).log_err();
523 }
524 let pool = pool.clone();
525 async move {
526 if frame.buffer.width() < 10 && frame.buffer.height() < 10 {
527 // when the remote stops sharing, we get an 8x8 black image.
528 // In a lil bit, the unpublish will come through and close the view,
529 // but until then, don't flash black.
530 return None;
531 }
532
533 video_frame_buffer_from_webrtc(pool?, frame.buffer)
534 }
535 })
536 }
537 #[cfg(not(target_os = "macos"))]
538 {
539 let executor = executor.clone();
540 NativeVideoStream::new(track.0.rtc_track()).filter_map(move |frame| {
541 executor.spawn(async move { video_frame_buffer_from_webrtc(frame.buffer) })
542 })
543 }
544}
545
546#[cfg(target_os = "macos")]
547fn create_buffer_pool(
548 width: u32,
549 height: u32,
550) -> Result<core_video::pixel_buffer_pool::CVPixelBufferPool> {
551 use core_foundation::{base::TCFType, number::CFNumber, string::CFString};
552 use core_video::pixel_buffer;
553 use core_video::{
554 pixel_buffer::kCVPixelFormatType_420YpCbCr8BiPlanarFullRange,
555 pixel_buffer_io_surface::kCVPixelBufferIOSurfaceCoreAnimationCompatibilityKey,
556 pixel_buffer_pool::{self},
557 };
558
559 let width_key: CFString =
560 unsafe { CFString::wrap_under_get_rule(pixel_buffer::kCVPixelBufferWidthKey) };
561 let height_key: CFString =
562 unsafe { CFString::wrap_under_get_rule(pixel_buffer::kCVPixelBufferHeightKey) };
563 let animation_key: CFString = unsafe {
564 CFString::wrap_under_get_rule(kCVPixelBufferIOSurfaceCoreAnimationCompatibilityKey)
565 };
566 let format_key: CFString =
567 unsafe { CFString::wrap_under_get_rule(pixel_buffer::kCVPixelBufferPixelFormatTypeKey) };
568
569 let yes: CFNumber = 1.into();
570 let width: CFNumber = (width as i32).into();
571 let height: CFNumber = (height as i32).into();
572 let format: CFNumber = (kCVPixelFormatType_420YpCbCr8BiPlanarFullRange as i64).into();
573
574 let buffer_attributes = core_foundation::dictionary::CFDictionary::from_CFType_pairs(&[
575 (width_key, width.into_CFType()),
576 (height_key, height.into_CFType()),
577 (animation_key, yes.into_CFType()),
578 (format_key, format.into_CFType()),
579 ]);
580
581 pixel_buffer_pool::CVPixelBufferPool::new(None, Some(&buffer_attributes)).map_err(|cv_return| {
582 anyhow::anyhow!("failed to create pixel buffer pool: CVReturn({cv_return})",)
583 })
584}
585
586#[cfg(target_os = "macos")]
587pub type RemoteVideoFrame = core_video::pixel_buffer::CVPixelBuffer;
588
589#[cfg(target_os = "macos")]
590fn video_frame_buffer_from_webrtc(
591 pool: core_video::pixel_buffer_pool::CVPixelBufferPool,
592 buffer: Box<dyn VideoBuffer>,
593) -> Option<RemoteVideoFrame> {
594 use core_foundation::base::TCFType;
595 use core_video::{pixel_buffer::CVPixelBuffer, r#return::kCVReturnSuccess};
596 use livekit::webrtc::native::yuv_helper::i420_to_nv12;
597
598 if let Some(native) = buffer.as_native() {
599 let pixel_buffer = native.get_cv_pixel_buffer();
600 if pixel_buffer.is_null() {
601 return None;
602 }
603 return unsafe { Some(CVPixelBuffer::wrap_under_get_rule(pixel_buffer as _)) };
604 }
605
606 let i420_buffer = buffer.as_i420()?;
607 let pixel_buffer = pool.create_pixel_buffer().log_err()?;
608
609 let image_buffer = unsafe {
610 if pixel_buffer.lock_base_address(0) != kCVReturnSuccess {
611 return None;
612 }
613
614 let dst_y = pixel_buffer.get_base_address_of_plane(0);
615 let dst_y_stride = pixel_buffer.get_bytes_per_row_of_plane(0);
616 let dst_y_len = pixel_buffer.get_height_of_plane(0) * dst_y_stride;
617 let dst_uv = pixel_buffer.get_base_address_of_plane(1);
618 let dst_uv_stride = pixel_buffer.get_bytes_per_row_of_plane(1);
619 let dst_uv_len = pixel_buffer.get_height_of_plane(1) * dst_uv_stride;
620 let width = pixel_buffer.get_width();
621 let height = pixel_buffer.get_height();
622 let dst_y_buffer = std::slice::from_raw_parts_mut(dst_y as *mut u8, dst_y_len);
623 let dst_uv_buffer = std::slice::from_raw_parts_mut(dst_uv as *mut u8, dst_uv_len);
624
625 let (stride_y, stride_u, stride_v) = i420_buffer.strides();
626 let (src_y, src_u, src_v) = i420_buffer.data();
627 i420_to_nv12(
628 src_y,
629 stride_y,
630 src_u,
631 stride_u,
632 src_v,
633 stride_v,
634 dst_y_buffer,
635 dst_y_stride as u32,
636 dst_uv_buffer,
637 dst_uv_stride as u32,
638 width as i32,
639 height as i32,
640 );
641
642 if pixel_buffer.unlock_base_address(0) != kCVReturnSuccess {
643 return None;
644 }
645
646 pixel_buffer
647 };
648
649 Some(image_buffer)
650}
651
652#[cfg(not(target_os = "macos"))]
653pub type RemoteVideoFrame = Arc<gpui::RenderImage>;
654
655#[cfg(not(target_os = "macos"))]
656fn video_frame_buffer_from_webrtc(buffer: Box<dyn VideoBuffer>) -> Option<RemoteVideoFrame> {
657 use gpui::RenderImage;
658 use image::{Frame, RgbaImage};
659 use livekit::webrtc::prelude::VideoFormatType;
660 use smallvec::SmallVec;
661 use std::alloc::{Layout, alloc};
662
663 let width = buffer.width();
664 let height = buffer.height();
665 let stride = width * 4;
666 let byte_len = (stride * height) as usize;
667 let argb_image = unsafe {
668 // Motivation for this unsafe code is to avoid initializing the frame data, since to_argb
669 // will write all bytes anyway.
670 let start_ptr = alloc(Layout::array::<u8>(byte_len).log_err()?);
671 if start_ptr.is_null() {
672 return None;
673 }
674 let argb_frame_slice = std::slice::from_raw_parts_mut(start_ptr, byte_len);
675 buffer.to_argb(
676 VideoFormatType::ARGB,
677 argb_frame_slice,
678 stride,
679 width as i32,
680 height as i32,
681 );
682 Vec::from_raw_parts(start_ptr, byte_len, byte_len)
683 };
684
685 // TODO: Unclear why providing argb_image to RgbaImage works properly.
686 let image = RgbaImage::from_raw(width, height, argb_image)
687 .with_context(|| "Bug: not enough bytes allocated for image.")
688 .log_err()?;
689
690 Some(Arc::new(RenderImage::new(SmallVec::from_elem(
691 Frame::new(image),
692 1,
693 ))))
694}
695
696#[cfg(target_os = "macos")]
697fn video_frame_buffer_to_webrtc(frame: ScreenCaptureFrame) -> Option<impl AsRef<dyn VideoBuffer>> {
698 use livekit::webrtc;
699
700 let pixel_buffer = frame.0.as_concrete_TypeRef();
701 std::mem::forget(frame.0);
702 unsafe {
703 Some(webrtc::video_frame::native::NativeBuffer::from_cv_pixel_buffer(pixel_buffer as _))
704 }
705}
706
707#[cfg(not(target_os = "macos"))]
708fn video_frame_buffer_to_webrtc(frame: ScreenCaptureFrame) -> Option<impl AsRef<dyn VideoBuffer>> {
709 use libwebrtc::native::yuv_helper::{abgr_to_nv12, argb_to_nv12};
710 use livekit::webrtc::prelude::NV12Buffer;
711 match frame.0 {
712 scap::frame::Frame::BGRx(frame) => {
713 let mut buffer = NV12Buffer::new(frame.width as u32, frame.height as u32);
714 let (stride_y, stride_uv) = buffer.strides();
715 let (data_y, data_uv) = buffer.data_mut();
716 argb_to_nv12(
717 &frame.data,
718 frame.width as u32 * 4,
719 data_y,
720 stride_y,
721 data_uv,
722 stride_uv,
723 frame.width,
724 frame.height,
725 );
726 Some(buffer)
727 }
728 scap::frame::Frame::RGBx(frame) => {
729 let mut buffer = NV12Buffer::new(frame.width as u32, frame.height as u32);
730 let (stride_y, stride_uv) = buffer.strides();
731 let (data_y, data_uv) = buffer.data_mut();
732 abgr_to_nv12(
733 &frame.data,
734 frame.width as u32 * 4,
735 data_y,
736 stride_y,
737 data_uv,
738 stride_uv,
739 frame.width,
740 frame.height,
741 );
742 Some(buffer)
743 }
744 scap::frame::Frame::YUVFrame(yuvframe) => {
745 let mut buffer = NV12Buffer::with_strides(
746 yuvframe.width as u32,
747 yuvframe.height as u32,
748 yuvframe.luminance_stride as u32,
749 yuvframe.chrominance_stride as u32,
750 );
751 let (luminance, chrominance) = buffer.data_mut();
752 luminance.copy_from_slice(yuvframe.luminance_bytes.as_slice());
753 chrominance.copy_from_slice(yuvframe.chrominance_bytes.as_slice());
754 Some(buffer)
755 }
756 _ => {
757 log::error!(
758 "Expected BGRx or YUV frame from scap screen capture but got some other format."
759 );
760 None
761 }
762 }
763}
764
765trait DeviceChangeListenerApi: Stream<Item = ()> + Sized {
766 fn new(input: bool) -> Result<Self>;
767}
768
769#[cfg(target_os = "macos")]
770mod macos {
771 use cocoa::{
772 base::{id, nil},
773 foundation::{NSProcessInfo, NSString},
774 };
775 use coreaudio::sys::{
776 AudioObjectAddPropertyListener, AudioObjectID, AudioObjectPropertyAddress,
777 AudioObjectRemovePropertyListener, OSStatus, kAudioHardwarePropertyDefaultInputDevice,
778 kAudioHardwarePropertyDefaultOutputDevice, kAudioObjectPropertyElementMaster,
779 kAudioObjectPropertyScopeGlobal, kAudioObjectSystemObject,
780 };
781 use futures::{StreamExt, channel::mpsc::UnboundedReceiver};
782 use objc::{msg_send, sel, sel_impl};
783
784 /// A guard that prevents App Nap while held.
785 ///
786 /// On macOS, App Nap can throttle background apps to save power. This can cause
787 /// audio artifacts when the app is not in the foreground. This guard tells macOS
788 /// that we're doing latency-sensitive work and should not be throttled.
789 ///
790 /// See Apple's documentation on prioritizing work at the app level:
791 /// https://developer.apple.com/library/archive/documentation/Performance/Conceptual/power_efficiency_guidelines_osx/PrioritizeWorkAtTheAppLevel.html
792 pub struct PreventAppNapGuard {
793 activity: id,
794 }
795
796 // The activity token returned by NSProcessInfo is thread-safe
797 unsafe impl Send for PreventAppNapGuard {}
798
799 // From NSProcessInfo.h
800 const NS_ACTIVITY_IDLE_SYSTEM_SLEEP_DISABLED: u64 = 1 << 20;
801 const NS_ACTIVITY_USER_INITIATED: u64 = 0x00FFFFFF | NS_ACTIVITY_IDLE_SYSTEM_SLEEP_DISABLED;
802 const NS_ACTIVITY_USER_INITIATED_ALLOWING_IDLE_SYSTEM_SLEEP: u64 =
803 NS_ACTIVITY_USER_INITIATED & !NS_ACTIVITY_IDLE_SYSTEM_SLEEP_DISABLED;
804
805 impl PreventAppNapGuard {
806 pub fn new() -> Self {
807 unsafe {
808 let process_info = NSProcessInfo::processInfo(nil);
809 #[allow(clippy::disallowed_methods)]
810 let reason = NSString::alloc(nil).init_str("Audio playback in progress");
811 let activity: id = msg_send![process_info, beginActivityWithOptions:NS_ACTIVITY_USER_INITIATED_ALLOWING_IDLE_SYSTEM_SLEEP reason:reason];
812 let _: () = msg_send![reason, release];
813 let _: () = msg_send![activity, retain];
814 Self { activity }
815 }
816 }
817 }
818
819 impl Drop for PreventAppNapGuard {
820 fn drop(&mut self) {
821 unsafe {
822 let process_info = NSProcessInfo::processInfo(nil);
823 let _: () = msg_send![process_info, endActivity:self.activity];
824 let _: () = msg_send![self.activity, release];
825 }
826 }
827 }
828
829 /// Implementation from: https://github.com/zed-industries/cpal/blob/fd8bc2fd39f1f5fdee5a0690656caff9a26d9d50/src/host/coreaudio/macos/property_listener.rs#L15
830 pub struct CoreAudioDefaultDeviceChangeListener {
831 rx: UnboundedReceiver<()>,
832 callback: Box<PropertyListenerCallbackWrapper>,
833 input: bool,
834 device_id: AudioObjectID, // Store the device ID to properly remove listeners
835 }
836
837 trait _AssertSend: Send {}
838 impl _AssertSend for CoreAudioDefaultDeviceChangeListener {}
839
840 struct PropertyListenerCallbackWrapper(Box<dyn FnMut() + Send>);
841
842 unsafe extern "C" fn property_listener_handler_shim(
843 _: AudioObjectID,
844 _: u32,
845 _: *const AudioObjectPropertyAddress,
846 callback: *mut ::std::os::raw::c_void,
847 ) -> OSStatus {
848 let wrapper = callback as *mut PropertyListenerCallbackWrapper;
849 unsafe { (*wrapper).0() };
850 0
851 }
852
853 impl super::DeviceChangeListenerApi for CoreAudioDefaultDeviceChangeListener {
854 fn new(input: bool) -> anyhow::Result<Self> {
855 let (tx, rx) = futures::channel::mpsc::unbounded();
856
857 let callback = Box::new(PropertyListenerCallbackWrapper(Box::new(move || {
858 tx.unbounded_send(()).ok();
859 })));
860
861 // Get the current default device ID
862 let device_id = unsafe {
863 // Listen for default device changes
864 coreaudio::Error::from_os_status(AudioObjectAddPropertyListener(
865 kAudioObjectSystemObject,
866 &AudioObjectPropertyAddress {
867 mSelector: if input {
868 kAudioHardwarePropertyDefaultInputDevice
869 } else {
870 kAudioHardwarePropertyDefaultOutputDevice
871 },
872 mScope: kAudioObjectPropertyScopeGlobal,
873 mElement: kAudioObjectPropertyElementMaster,
874 },
875 Some(property_listener_handler_shim),
876 &*callback as *const _ as *mut _,
877 ))?;
878
879 // Also listen for changes to the device configuration
880 let device_id = if input {
881 let mut input_device: AudioObjectID = 0;
882 let mut prop_size = std::mem::size_of::<AudioObjectID>() as u32;
883 let result = coreaudio::sys::AudioObjectGetPropertyData(
884 kAudioObjectSystemObject,
885 &AudioObjectPropertyAddress {
886 mSelector: kAudioHardwarePropertyDefaultInputDevice,
887 mScope: kAudioObjectPropertyScopeGlobal,
888 mElement: kAudioObjectPropertyElementMaster,
889 },
890 0,
891 std::ptr::null(),
892 &mut prop_size as *mut _,
893 &mut input_device as *mut _ as *mut _,
894 );
895 if result != 0 {
896 log::warn!("Failed to get default input device ID");
897 0
898 } else {
899 input_device
900 }
901 } else {
902 let mut output_device: AudioObjectID = 0;
903 let mut prop_size = std::mem::size_of::<AudioObjectID>() as u32;
904 let result = coreaudio::sys::AudioObjectGetPropertyData(
905 kAudioObjectSystemObject,
906 &AudioObjectPropertyAddress {
907 mSelector: kAudioHardwarePropertyDefaultOutputDevice,
908 mScope: kAudioObjectPropertyScopeGlobal,
909 mElement: kAudioObjectPropertyElementMaster,
910 },
911 0,
912 std::ptr::null(),
913 &mut prop_size as *mut _,
914 &mut output_device as *mut _ as *mut _,
915 );
916 if result != 0 {
917 log::warn!("Failed to get default output device ID");
918 0
919 } else {
920 output_device
921 }
922 };
923
924 if device_id != 0 {
925 // Listen for format changes on the device
926 coreaudio::Error::from_os_status(AudioObjectAddPropertyListener(
927 device_id,
928 &AudioObjectPropertyAddress {
929 mSelector: coreaudio::sys::kAudioDevicePropertyStreamFormat,
930 mScope: if input {
931 coreaudio::sys::kAudioObjectPropertyScopeInput
932 } else {
933 coreaudio::sys::kAudioObjectPropertyScopeOutput
934 },
935 mElement: kAudioObjectPropertyElementMaster,
936 },
937 Some(property_listener_handler_shim),
938 &*callback as *const _ as *mut _,
939 ))?;
940 }
941
942 device_id
943 };
944
945 Ok(Self {
946 rx,
947 callback,
948 input,
949 device_id,
950 })
951 }
952 }
953
954 impl Drop for CoreAudioDefaultDeviceChangeListener {
955 fn drop(&mut self) {
956 unsafe {
957 // Remove the system-level property listener
958 AudioObjectRemovePropertyListener(
959 kAudioObjectSystemObject,
960 &AudioObjectPropertyAddress {
961 mSelector: if self.input {
962 kAudioHardwarePropertyDefaultInputDevice
963 } else {
964 kAudioHardwarePropertyDefaultOutputDevice
965 },
966 mScope: kAudioObjectPropertyScopeGlobal,
967 mElement: kAudioObjectPropertyElementMaster,
968 },
969 Some(property_listener_handler_shim),
970 &*self.callback as *const _ as *mut _,
971 );
972
973 // Remove the device-specific property listener if we have a valid device ID
974 if self.device_id != 0 {
975 AudioObjectRemovePropertyListener(
976 self.device_id,
977 &AudioObjectPropertyAddress {
978 mSelector: coreaudio::sys::kAudioDevicePropertyStreamFormat,
979 mScope: if self.input {
980 coreaudio::sys::kAudioObjectPropertyScopeInput
981 } else {
982 coreaudio::sys::kAudioObjectPropertyScopeOutput
983 },
984 mElement: kAudioObjectPropertyElementMaster,
985 },
986 Some(property_listener_handler_shim),
987 &*self.callback as *const _ as *mut _,
988 );
989 }
990 }
991 }
992 }
993
994 impl futures::Stream for CoreAudioDefaultDeviceChangeListener {
995 type Item = ();
996
997 fn poll_next(
998 mut self: std::pin::Pin<&mut Self>,
999 cx: &mut std::task::Context<'_>,
1000 ) -> std::task::Poll<Option<Self::Item>> {
1001 self.rx.poll_next_unpin(cx)
1002 }
1003 }
1004}
1005
1006#[cfg(target_os = "macos")]
1007type DeviceChangeListener = macos::CoreAudioDefaultDeviceChangeListener;
1008#[cfg(target_os = "macos")]
1009use macos::PreventAppNapGuard;
1010
1011#[cfg(not(target_os = "macos"))]
1012mod noop_change_listener {
1013 use std::task::Poll;
1014
1015 use super::DeviceChangeListenerApi;
1016
1017 pub struct NoopOutputDeviceChangelistener {}
1018
1019 impl DeviceChangeListenerApi for NoopOutputDeviceChangelistener {
1020 fn new(_input: bool) -> anyhow::Result<Self> {
1021 Ok(NoopOutputDeviceChangelistener {})
1022 }
1023 }
1024
1025 impl futures::Stream for NoopOutputDeviceChangelistener {
1026 type Item = ();
1027
1028 fn poll_next(
1029 self: std::pin::Pin<&mut Self>,
1030 _cx: &mut std::task::Context<'_>,
1031 ) -> Poll<Option<Self::Item>> {
1032 Poll::Pending
1033 }
1034 }
1035}
1036
1037#[cfg(not(target_os = "macos"))]
1038type DeviceChangeListener = noop_change_listener::NoopOutputDeviceChangelistener;
1039