Skip to repository content

tenant.openagents/omega

No repository description is available.

OpenAgents Git authority 2026-07-28T05:14:46.779Z Public web read
NIP-34 coordinate30617:7649603503856e5148d571eac2766b288a8ff1e9e35d380337a1d2b0015b4f92:omega
MaintainersHidden in public view
References2 branches · 1 tag
Read-only clonegit clone https://openagents.com/git/tenant.openagents/omega.git
Browse files

playback.rs

1039 lines · 40.6 KB · rust
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(&timestamped.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
Served at tenant.openagents/omega Member data and write actions are omitted.