1pub 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 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 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 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 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 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 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}