Skip to main content

servo_media_gstreamer/
lib.rs

1/* This Source Code Form is subject to the terms of the Mozilla Public
2 * License, v. 2.0. If a copy of the MPL was not distributed with this
3 * file, You can obtain one at https://mozilla.org/MPL/2.0/. */
4
5pub mod audio_decoder;
6pub mod audio_sink;
7pub mod audio_stream_reader;
8mod datachannel;
9mod device_monitor;
10pub mod media_capture;
11pub mod media_stream;
12mod media_stream_source;
13pub mod player;
14mod registry_scanner;
15mod render;
16mod source;
17pub mod webrtc;
18
19use std::collections::HashMap;
20use std::path::PathBuf;
21use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
22use std::sync::mpsc::{self, Sender};
23use std::sync::{Arc, LazyLock, Mutex, OnceLock, Weak};
24use std::thread;
25use std::vec::Vec;
26
27use device_monitor::GStreamerDeviceMonitor;
28use gstreamer::prelude::*;
29use log::warn;
30use media_stream::GStreamerMediaStream;
31use mime::Mime;
32use registry_scanner::GSTREAMER_REGISTRY_SCANNER;
33use servo_base::generic_channel::GenericCallback;
34use servo_media::{Backend, BackendDeInit, BackendInit, MediaInstanceError, SupportsMediaType};
35use servo_media_audio::context::{AudioContext, AudioContextOptions};
36use servo_media_audio::decoder::AudioDecoder;
37use servo_media_audio::sink::AudioSinkError;
38use servo_media_audio::{AudioBackend, AudioStreamReader};
39use servo_media_player::audio::AudioRenderer;
40use servo_media_player::context::PlayerGLContext;
41use servo_media_player::video::VideoFrameRenderer;
42use servo_media_player::{Player, PlayerEvent, StreamType};
43use servo_media_streams::capture::MediaTrackConstraintSet;
44use servo_media_streams::device_monitor::MediaDeviceMonitor;
45use servo_media_streams::registry::MediaStreamId;
46use servo_media_streams::{MediaOutput, MediaSocket, MediaStreamType};
47use servo_media_traits::{BackendMsg, ClientContextId, MediaInstance};
48use servo_media_webrtc::{WebRtcBackend, WebRtcController, WebRtcSignaller};
49
50static BACKEND_BASE_TIME: LazyLock<gstreamer::ClockTime> =
51    LazyLock::new(|| gstreamer::SystemClock::obtain().time());
52
53static BACKEND_THREAD: OnceLock<bool> = OnceLock::new();
54
55pub type WeakMediaInstance = Weak<Mutex<dyn MediaInstance>>;
56pub type WeakMediaInstanceHashMap = HashMap<ClientContextId, Vec<(usize, WeakMediaInstance)>>;
57
58pub struct GStreamerBackend {
59    capture_mocking: AtomicBool,
60    instances: Arc<Mutex<WeakMediaInstanceHashMap>>,
61    next_instance_id: AtomicUsize,
62    /// Channel to communicate media instances with its owner Backend.
63    backend_chan: Arc<Mutex<Sender<BackendMsg>>>,
64}
65
66#[derive(Debug)]
67#[allow(dead_code)]
68pub struct ErrorLoadingPlugins<'a>(Vec<&'a str>);
69
70impl GStreamerBackend {
71    pub fn init_with_plugins<'a, T: AsRef<str>>(
72        plugin_dir: PathBuf,
73        plugins: &'a [T],
74    ) -> Result<Box<dyn Backend>, ErrorLoadingPlugins<'a>> {
75        gstreamer::init().unwrap();
76
77        // GStreamer between 1.19.1 and 1.22.7 will not send messages like "end of stream"
78        // to GstPlayer unless there is a GLib main loop running somewhere. We should remove
79        // this workaround when we raise of required version of GStreamer.
80        // See https://github.com/servo/media/pull/393.
81        let needs_background_glib_main_loop = {
82            let (major, minor, micro, _) = gstreamer::version();
83            (major, minor, micro) >= (1, 19, 1) && (major, minor, micro) <= (1, 22, 7)
84        };
85
86        if needs_background_glib_main_loop {
87            BACKEND_THREAD.get_or_init(|| {
88                thread::spawn(|| glib::MainLoop::new(None, false).run());
89                true
90            });
91        }
92
93        #[cfg(target_os = "linux")]
94        {
95            // Remove the va plugins that have some issues on certain platforms.
96            // When the driver is fixed we will remove this.
97            let registry = gstreamer::Registry::get();
98            if let Some(plugin) = registry.find_plugin("va") {
99                registry.remove_plugin(&plugin);
100            }
101            if let Some(plugin) = registry.find_plugin("vaapi") {
102                registry.remove_plugin(&plugin);
103            }
104        }
105
106        let mut errors = vec![];
107        for plugin in plugins {
108            let mut path = plugin_dir.clone();
109            path.push(plugin.as_ref());
110            let registry = gstreamer::Registry::get();
111            if gstreamer::Plugin::load_file(&path)
112                .is_ok_and(|plugin| registry.add_plugin(&plugin).is_ok())
113            {
114                continue;
115            }
116            errors.push(plugin.as_ref());
117        }
118
119        if !errors.is_empty() {
120            return Err(ErrorLoadingPlugins(errors));
121        }
122
123        type MediaInstancesVec = Vec<(usize, Weak<Mutex<dyn MediaInstance>>)>;
124        let instances: HashMap<ClientContextId, MediaInstancesVec> = Default::default();
125        let instances = Arc::new(Mutex::new(instances));
126
127        let instances_ = instances.clone();
128        let (backend_chan, recvr) = mpsc::channel();
129        thread::Builder::new()
130            .name("GStreamerBackend ShutdownThread".to_owned())
131            .spawn(move || {
132                match recvr.recv().unwrap() {
133                    BackendMsg::Shutdown {
134                        context,
135                        id,
136                        tx_ack,
137                    } => {
138                        let mut instances_ = instances_.lock().unwrap();
139                        if let Some(vec) = instances_.get_mut(&context) {
140                            vec.retain(|m| m.0 != id);
141                            if vec.is_empty() {
142                                instances_.remove(&context);
143                            }
144                        }
145                        // tell caller we are done removing this instance
146                        let _ = tx_ack.send(());
147                    },
148                };
149            })
150            .unwrap();
151
152        Ok(Box::new(GStreamerBackend {
153            capture_mocking: AtomicBool::new(false),
154            instances,
155            next_instance_id: AtomicUsize::new(0),
156            backend_chan: Arc::new(Mutex::new(backend_chan)),
157        }))
158    }
159
160    fn media_instance_action(
161        &self,
162        id: &ClientContextId,
163        cb: &dyn Fn(&dyn MediaInstance) -> Result<(), MediaInstanceError>,
164    ) {
165        let mut instances = self.instances.lock().unwrap();
166        match instances.get_mut(id) {
167            Some(vec) => vec.retain(|(_, weak)| match weak.upgrade() {
168                Some(instance) => {
169                    if cb(&*(instance.lock().unwrap())).is_err() {
170                        warn!("Error executing media instance action");
171                    }
172                    true
173                },
174                _ => false,
175            }),
176            None => {
177                warn!("Trying to exec media action on an unknown client context");
178            },
179        }
180    }
181}
182
183impl Backend for GStreamerBackend {
184    #[expect(clippy::redundant_clone, reason = "False positive")]
185    fn create_player(
186        &self,
187        context_id: &ClientContextId,
188        stream_type: StreamType,
189        sender: GenericCallback<PlayerEvent>,
190        renderer: Option<Arc<Mutex<dyn VideoFrameRenderer>>>,
191        audio_renderer: Option<Arc<Mutex<dyn AudioRenderer>>>,
192        gl_context: Box<dyn PlayerGLContext>,
193    ) -> Arc<Mutex<dyn Player>> {
194        let id = self.next_instance_id.fetch_add(1, Ordering::Relaxed);
195        let player = Arc::new(Mutex::new(player::GStreamerPlayer::new(
196            id,
197            context_id,
198            self.backend_chan.clone(),
199            stream_type,
200            sender,
201            renderer,
202            audio_renderer,
203            gl_context,
204        )));
205        let mut instances = self.instances.lock().unwrap();
206        let entry = instances.entry(*context_id).or_default();
207        entry.push((id, Arc::downgrade(&player).clone()));
208        player
209    }
210
211    #[expect(clippy::redundant_clone, reason = "False positive")]
212    fn create_audio_context(
213        &self,
214        client_context_id: &ClientContextId,
215        options: AudioContextOptions,
216    ) -> Result<Arc<Mutex<AudioContext>>, AudioSinkError> {
217        let id = self.next_instance_id.fetch_add(1, Ordering::Relaxed);
218        let audio_context =
219            AudioContext::new::<Self>(id, client_context_id, self.backend_chan.clone(), options)?;
220
221        let audio_context = Arc::new(Mutex::new(audio_context));
222
223        let mut instances = self.instances.lock().unwrap();
224        let entry = instances.entry(*client_context_id).or_default();
225        entry.push((id, Arc::downgrade(&audio_context).clone()));
226
227        Ok(audio_context)
228    }
229
230    fn create_webrtc(&self, signaller: Box<dyn WebRtcSignaller>) -> WebRtcController {
231        WebRtcController::new::<Self>(signaller)
232    }
233
234    fn create_audiostream(&self) -> MediaStreamId {
235        GStreamerMediaStream::create_audio()
236    }
237
238    fn create_videostream(&self) -> MediaStreamId {
239        GStreamerMediaStream::create_video()
240    }
241
242    fn create_stream_output(&self) -> Box<dyn MediaOutput> {
243        Box::new(media_stream::MediaSink::default())
244    }
245
246    fn create_stream_and_socket(
247        &self,
248        ty: MediaStreamType,
249    ) -> (Box<dyn MediaSocket>, MediaStreamId) {
250        let (id, socket) = GStreamerMediaStream::create_proxy(ty);
251        (Box::new(socket), id)
252    }
253
254    fn create_audioinput_stream(&self, set: MediaTrackConstraintSet) -> Option<MediaStreamId> {
255        if self.capture_mocking.load(Ordering::Acquire) {
256            // XXXManishearth we should caps filter this
257            return Some(self.create_audiostream());
258        }
259        media_capture::create_audioinput_stream(set)
260    }
261
262    fn create_videoinput_stream(&self, set: MediaTrackConstraintSet) -> Option<MediaStreamId> {
263        if self.capture_mocking.load(Ordering::Acquire) {
264            // XXXManishearth we should caps filter this
265            return Some(self.create_videostream());
266        }
267        media_capture::create_videoinput_stream(set)
268    }
269
270    fn can_play_type(&self, media_type: &str) -> SupportsMediaType {
271        if let Ok(mime) = media_type.parse::<Mime>() {
272            let mime_type = mime.type_().as_str().to_owned() + "/" + mime.subtype().as_str();
273            let codecs = match mime.get_param("codecs") {
274                Some(codecs) => codecs
275                    .as_str()
276                    .split(',')
277                    .map(|codec| codec.trim())
278                    .collect(),
279                None => vec![],
280            };
281
282            if GSTREAMER_REGISTRY_SCANNER.is_container_type_supported(&mime_type) {
283                if codecs.is_empty() {
284                    return SupportsMediaType::Maybe;
285                } else if GSTREAMER_REGISTRY_SCANNER.are_all_codecs_supported(&codecs) {
286                    return SupportsMediaType::Probably;
287                } else {
288                    return SupportsMediaType::No;
289                }
290            }
291        }
292        SupportsMediaType::No
293    }
294
295    fn set_capture_mocking(&self, mock: bool) {
296        self.capture_mocking.store(mock, Ordering::Release)
297    }
298
299    fn mute(&self, id: &ClientContextId, val: bool) {
300        self.media_instance_action(
301            id,
302            &(move |instance: &dyn MediaInstance| instance.mute(val)),
303        );
304    }
305
306    fn suspend(&self, id: &ClientContextId) {
307        self.media_instance_action(id, &|instance: &dyn MediaInstance| instance.suspend());
308    }
309
310    fn resume(&self, id: &ClientContextId) {
311        self.media_instance_action(id, &|instance: &dyn MediaInstance| instance.resume());
312    }
313
314    fn get_device_monitor(&self) -> Box<dyn MediaDeviceMonitor> {
315        Box::new(GStreamerDeviceMonitor::new())
316    }
317}
318
319impl AudioBackend for GStreamerBackend {
320    type Sink = audio_sink::GStreamerAudioSink;
321    fn make_decoder() -> Box<dyn AudioDecoder> {
322        Box::new(audio_decoder::GStreamerAudioDecoder::new())
323    }
324    fn make_sink() -> Result<Self::Sink, AudioSinkError> {
325        audio_sink::GStreamerAudioSink::new()
326    }
327
328    fn make_streamreader(
329        id: MediaStreamId,
330        sample_rate: f32,
331    ) -> Result<Box<dyn AudioStreamReader + Send>, AudioSinkError> {
332        let reader = audio_stream_reader::GStreamerAudioStreamReader::new(id, sample_rate)
333            .map_err(AudioSinkError::Backend)?;
334        Ok(Box::new(reader))
335    }
336}
337
338impl WebRtcBackend for GStreamerBackend {
339    type Controller = webrtc::GStreamerWebRtcController;
340
341    fn construct_webrtc_controller(
342        signaller: Box<dyn WebRtcSignaller>,
343        thread: WebRtcController,
344    ) -> Self::Controller {
345        webrtc::construct(signaller, thread).expect("WebRTC creation failed")
346    }
347}
348
349impl BackendInit for GStreamerBackend {
350    fn init() -> Box<dyn Backend> {
351        Self::init_with_plugins::<&str>(PathBuf::new(), &[]).unwrap()
352    }
353}
354
355impl BackendDeInit for GStreamerBackend {
356    fn deinit(&self) {
357        let to_shutdown: Vec<(ClientContextId, usize)> = {
358            let map = self.instances.lock().unwrap();
359            map.iter()
360                .flat_map(|(ctx, v)| v.iter().map(move |(id, _)| (*ctx, *id)))
361                .collect()
362        };
363
364        for (ctx, id) in to_shutdown {
365            let (tx_ack, rx_ack) = mpsc::channel();
366            let _ = self
367                .backend_chan
368                .lock()
369                .unwrap()
370                .send(BackendMsg::Shutdown {
371                    context: ctx,
372                    id,
373                    tx_ack,
374                });
375            let _ = rx_ack.recv();
376        }
377    }
378}