Skip to main content

gstreamer_app/
app_src.rs

1// Take a look at the license at the top of the repository in the LICENSE file.
2
3use std::{mem, panic, ptr};
4
5#[cfg(not(panic = "abort"))]
6use std::sync::atomic::{AtomicBool, Ordering};
7
8use glib::{
9    ffi::{gboolean, gpointer},
10    prelude::*,
11    translate::*,
12};
13
14use crate::{AppSrc, ffi};
15
16#[cfg(feature = "futures")]
17pub use crate::app_src_futures::AppSrcSink;
18
19#[allow(clippy::type_complexity)]
20pub struct AppSrcCallbacks {
21    need_data: Option<Box<dyn FnMut(&AppSrc, u32) + Send + 'static>>,
22    enough_data: Option<Box<dyn Fn(&AppSrc) + Send + Sync + 'static>>,
23    seek_data: Option<Box<dyn Fn(&AppSrc, u64) -> bool + Send + Sync + 'static>>,
24    #[cfg(not(panic = "abort"))]
25    panicked: AtomicBool,
26    callbacks: ffi::GstAppSrcCallbacks,
27}
28
29unsafe impl Send for AppSrcCallbacks {}
30unsafe impl Sync for AppSrcCallbacks {}
31
32impl AppSrcCallbacks {
33    pub fn builder() -> AppSrcCallbacksBuilder {
34        skip_assert_initialized!();
35
36        AppSrcCallbacksBuilder {
37            need_data: None,
38            enough_data: None,
39            seek_data: None,
40        }
41    }
42}
43
44#[allow(clippy::type_complexity)]
45#[must_use = "The builder must be built to be used"]
46pub struct AppSrcCallbacksBuilder {
47    need_data: Option<Box<dyn FnMut(&AppSrc, u32) + Send + 'static>>,
48    enough_data: Option<Box<dyn Fn(&AppSrc) + Send + Sync + 'static>>,
49    seek_data: Option<Box<dyn Fn(&AppSrc, u64) -> bool + Send + Sync + 'static>>,
50}
51
52impl AppSrcCallbacksBuilder {
53    pub fn need_data<F: FnMut(&AppSrc, u32) + Send + 'static>(self, need_data: F) -> Self {
54        Self {
55            need_data: Some(Box::new(need_data)),
56            ..self
57        }
58    }
59
60    pub fn need_data_if<F: FnMut(&AppSrc, u32) + Send + 'static>(
61        self,
62        need_data: F,
63        predicate: bool,
64    ) -> Self {
65        if predicate {
66            self.need_data(need_data)
67        } else {
68            self
69        }
70    }
71
72    pub fn need_data_if_some<F: FnMut(&AppSrc, u32) + Send + 'static>(
73        self,
74        need_data: Option<F>,
75    ) -> Self {
76        if let Some(need_data) = need_data {
77            self.need_data(need_data)
78        } else {
79            self
80        }
81    }
82
83    pub fn enough_data<F: Fn(&AppSrc) + Send + Sync + 'static>(self, enough_data: F) -> Self {
84        Self {
85            enough_data: Some(Box::new(enough_data)),
86            ..self
87        }
88    }
89
90    pub fn enough_data_if<F: Fn(&AppSrc) + Send + Sync + 'static>(
91        self,
92        enough_data: F,
93        predicate: bool,
94    ) -> Self {
95        if predicate {
96            self.enough_data(enough_data)
97        } else {
98            self
99        }
100    }
101
102    pub fn enough_data_if_some<F: Fn(&AppSrc) + Send + Sync + 'static>(
103        self,
104        enough_data: Option<F>,
105    ) -> Self {
106        if let Some(enough_data) = enough_data {
107            self.enough_data(enough_data)
108        } else {
109            self
110        }
111    }
112
113    pub fn seek_data<F: Fn(&AppSrc, u64) -> bool + Send + Sync + 'static>(
114        self,
115        seek_data: F,
116    ) -> Self {
117        Self {
118            seek_data: Some(Box::new(seek_data)),
119            ..self
120        }
121    }
122
123    pub fn seek_data_if<F: Fn(&AppSrc, u64) -> bool + Send + Sync + 'static>(
124        self,
125        seek_data: F,
126        predicate: bool,
127    ) -> Self {
128        if predicate {
129            self.seek_data(seek_data)
130        } else {
131            self
132        }
133    }
134
135    pub fn seek_data_if_some<F: Fn(&AppSrc, u64) -> bool + Send + Sync + 'static>(
136        self,
137        seek_data: Option<F>,
138    ) -> Self {
139        if let Some(seek_data) = seek_data {
140            self.seek_data(seek_data)
141        } else {
142            self
143        }
144    }
145
146    #[must_use = "Building the callbacks without using them has no effect"]
147    pub fn build(self) -> AppSrcCallbacks {
148        let have_need_data = self.need_data.is_some();
149        let have_enough_data = self.enough_data.is_some();
150        let have_seek_data = self.seek_data.is_some();
151
152        AppSrcCallbacks {
153            need_data: self.need_data,
154            enough_data: self.enough_data,
155            seek_data: self.seek_data,
156            #[cfg(not(panic = "abort"))]
157            panicked: AtomicBool::new(false),
158            callbacks: ffi::GstAppSrcCallbacks {
159                need_data: if have_need_data {
160                    Some(trampoline_need_data)
161                } else {
162                    None
163                },
164                enough_data: if have_enough_data {
165                    Some(trampoline_enough_data)
166                } else {
167                    None
168                },
169                seek_data: if have_seek_data {
170                    Some(trampoline_seek_data)
171                } else {
172                    None
173                },
174                _gst_reserved: [
175                    ptr::null_mut(),
176                    ptr::null_mut(),
177                    ptr::null_mut(),
178                    ptr::null_mut(),
179                ],
180            },
181        }
182    }
183}
184
185unsafe extern "C" fn trampoline_need_data(
186    appsrc: *mut ffi::GstAppSrc,
187    length: u32,
188    callbacks: gpointer,
189) {
190    unsafe {
191        let callbacks = callbacks as *mut AppSrcCallbacks;
192        let element: Borrowed<AppSrc> = from_glib_borrow(appsrc);
193
194        #[cfg(not(panic = "abort"))]
195        if (*callbacks).panicked.load(Ordering::Relaxed) {
196            let element: Borrowed<AppSrc> = from_glib_borrow(appsrc);
197            gst::subclass::post_panic_error_message(
198                element.upcast_ref(),
199                element.upcast_ref(),
200                None,
201            );
202            return;
203        }
204
205        if let Some(ref mut need_data) = (*callbacks).need_data {
206            let result =
207                panic::catch_unwind(panic::AssertUnwindSafe(|| need_data(&element, length)));
208            match result {
209                Ok(result) => result,
210                Err(err) => {
211                    #[cfg(panic = "abort")]
212                    {
213                        unreachable!("{err:?}");
214                    }
215                    #[cfg(not(panic = "abort"))]
216                    {
217                        (*callbacks).panicked.store(true, Ordering::Relaxed);
218                        gst::subclass::post_panic_error_message(
219                            element.upcast_ref(),
220                            element.upcast_ref(),
221                            Some(err),
222                        );
223                    }
224                }
225            }
226        }
227    }
228}
229
230unsafe extern "C" fn trampoline_enough_data(appsrc: *mut ffi::GstAppSrc, callbacks: gpointer) {
231    unsafe {
232        let callbacks = callbacks as *const AppSrcCallbacks;
233        let element: Borrowed<AppSrc> = from_glib_borrow(appsrc);
234
235        #[cfg(not(panic = "abort"))]
236        if (*callbacks).panicked.load(Ordering::Relaxed) {
237            let element: Borrowed<AppSrc> = from_glib_borrow(appsrc);
238            gst::subclass::post_panic_error_message(
239                element.upcast_ref(),
240                element.upcast_ref(),
241                None,
242            );
243            return;
244        }
245
246        if let Some(ref enough_data) = (*callbacks).enough_data {
247            let result = panic::catch_unwind(panic::AssertUnwindSafe(|| enough_data(&element)));
248            match result {
249                Ok(result) => result,
250                Err(err) => {
251                    #[cfg(panic = "abort")]
252                    {
253                        unreachable!("{err:?}");
254                    }
255                    #[cfg(not(panic = "abort"))]
256                    {
257                        (*callbacks).panicked.store(true, Ordering::Relaxed);
258                        gst::subclass::post_panic_error_message(
259                            element.upcast_ref(),
260                            element.upcast_ref(),
261                            Some(err),
262                        );
263                    }
264                }
265            }
266        }
267    }
268}
269
270unsafe extern "C" fn trampoline_seek_data(
271    appsrc: *mut ffi::GstAppSrc,
272    offset: u64,
273    callbacks: gpointer,
274) -> gboolean {
275    unsafe {
276        let callbacks = callbacks as *const AppSrcCallbacks;
277        let element: Borrowed<AppSrc> = from_glib_borrow(appsrc);
278
279        #[cfg(not(panic = "abort"))]
280        if (*callbacks).panicked.load(Ordering::Relaxed) {
281            let element: Borrowed<AppSrc> = from_glib_borrow(appsrc);
282            gst::subclass::post_panic_error_message(
283                element.upcast_ref(),
284                element.upcast_ref(),
285                None,
286            );
287            return false.into_glib();
288        }
289
290        let ret = if let Some(ref seek_data) = (*callbacks).seek_data {
291            let result =
292                panic::catch_unwind(panic::AssertUnwindSafe(|| seek_data(&element, offset)));
293            match result {
294                Ok(result) => result,
295                Err(err) => {
296                    #[cfg(panic = "abort")]
297                    {
298                        unreachable!("{err:?}");
299                    }
300                    #[cfg(not(panic = "abort"))]
301                    {
302                        (*callbacks).panicked.store(true, Ordering::Relaxed);
303                        gst::subclass::post_panic_error_message(
304                            element.upcast_ref(),
305                            element.upcast_ref(),
306                            Some(err),
307                        );
308
309                        false
310                    }
311                }
312            }
313        } else {
314            false
315        };
316
317        ret.into_glib()
318    }
319}
320
321unsafe extern "C" fn destroy_callbacks(ptr: gpointer) {
322    unsafe {
323        let _ = Box::<AppSrcCallbacks>::from_raw(ptr as *mut _);
324    }
325}
326
327impl AppSrc {
328    // rustdoc-stripper-ignore-next
329    /// Creates a new builder-pattern struct instance to construct [`AppSrc`] objects.
330    ///
331    /// This method returns an instance of [`AppSrcBuilder`](crate::builders::AppSrcBuilder) which can be used to create [`AppSrc`] objects.
332    pub fn builder<'a>() -> AppSrcBuilder<'a> {
333        assert_initialized_main_thread!();
334        AppSrcBuilder {
335            builder: gst::Object::builder(),
336            callbacks: None,
337            automatic_eos: None,
338        }
339    }
340
341    /// Set callbacks which will be executed when data is needed, enough data has
342    /// been collected or when a seek should be performed.
343    /// This is an alternative to using the signals, it has lower overhead and is thus
344    /// less expensive, but also less flexible.
345    ///
346    /// If callbacks are installed, no signals will be emitted for performance
347    /// reasons.
348    ///
349    /// Before 1.16.3 it was not possible to change the callbacks in a thread-safe
350    /// way.
351    ///
352    /// Since 1.28.3 it is allowed to set the `callbacks` to [`None`] to unset them.
353    ///
354    /// Note that `gst_app_src_set_callbacks()` and
355    /// `gst_app_src_set_simple_callbacks()` are mutually exclusive and setting one
356    /// will unset the other.
357    /// ## `callbacks`
358    /// the callbacks
359    /// ## `notify`
360    /// a destroy notify function
361    #[doc(alias = "gst_app_src_set_callbacks")]
362    pub fn set_callbacks(&self, callbacks: AppSrcCallbacks) {
363        unsafe {
364            let src = self.to_glib_none().0;
365            #[allow(clippy::manual_dangling_ptr)]
366            #[cfg(not(feature = "v1_18"))]
367            {
368                static SET_ONCE_QUARK: std::sync::OnceLock<glib::Quark> =
369                    std::sync::OnceLock::new();
370
371                let set_once_quark = SET_ONCE_QUARK
372                    .get_or_init(|| glib::Quark::from_str("gstreamer-rs-app-src-callbacks"));
373
374                // This is not thread-safe before 1.16.3, see
375                // https://gitlab.freedesktop.org/gstreamer/gst-plugins-base/merge_requests/570
376                if gst::version() < (1, 16, 3, 0) {
377                    if !glib::gobject_ffi::g_object_get_qdata(
378                        src as *mut _,
379                        set_once_quark.into_glib(),
380                    )
381                    .is_null()
382                    {
383                        panic!("AppSrc callbacks can only be set once");
384                    }
385
386                    glib::gobject_ffi::g_object_set_qdata(
387                        src as *mut _,
388                        set_once_quark.into_glib(),
389                        1 as *mut _,
390                    );
391                }
392            }
393
394            ffi::gst_app_src_set_callbacks(
395                src,
396                mut_override(&callbacks.callbacks),
397                Box::into_raw(Box::new(callbacks)) as *mut _,
398                Some(destroy_callbacks),
399            );
400        }
401    }
402
403    /// Configure the `min` and `max` latency in `src`. If `min` is set to -1, the
404    /// default latency calculations for pseudo-live sources will be used.
405    /// ## `min`
406    /// the min latency
407    /// ## `max`
408    /// the max latency
409    #[doc(alias = "gst_app_src_set_latency")]
410    pub fn set_latency(
411        &self,
412        min: impl Into<Option<gst::ClockTime>>,
413        max: impl Into<Option<gst::ClockTime>>,
414    ) {
415        unsafe {
416            ffi::gst_app_src_set_latency(
417                self.to_glib_none().0,
418                min.into().into_glib(),
419                max.into().into_glib(),
420            );
421        }
422    }
423
424    /// Retrieve the min and max latencies in `min` and `max` respectively.
425    ///
426    /// # Returns
427    ///
428    ///
429    /// ## `min`
430    /// the min latency
431    ///
432    /// ## `max`
433    /// the max latency
434    #[doc(alias = "get_latency")]
435    #[doc(alias = "gst_app_src_get_latency")]
436    pub fn latency(&self) -> (Option<gst::ClockTime>, Option<gst::ClockTime>) {
437        unsafe {
438            let mut min = mem::MaybeUninit::uninit();
439            let mut max = mem::MaybeUninit::uninit();
440            ffi::gst_app_src_get_latency(self.to_glib_none().0, min.as_mut_ptr(), max.as_mut_ptr());
441            (from_glib(min.assume_init()), from_glib(max.assume_init()))
442        }
443    }
444
445    #[doc(alias = "do-timestamp")]
446    #[doc(alias = "gst_base_src_set_do_timestamp")]
447    pub fn set_do_timestamp(&self, timestamp: bool) {
448        unsafe {
449            gst_base::ffi::gst_base_src_set_do_timestamp(
450                self.as_ptr() as *mut gst_base::ffi::GstBaseSrc,
451                timestamp.into_glib(),
452            );
453        }
454    }
455
456    #[doc(alias = "do-timestamp")]
457    #[doc(alias = "gst_base_src_get_do_timestamp")]
458    pub fn do_timestamp(&self) -> bool {
459        unsafe {
460            from_glib(gst_base::ffi::gst_base_src_get_do_timestamp(
461                self.as_ptr() as *mut gst_base::ffi::GstBaseSrc
462            ))
463        }
464    }
465
466    #[doc(alias = "do-timestamp")]
467    pub fn connect_do_timestamp_notify<F: Fn(&Self) + Send + Sync + 'static>(
468        &self,
469        f: F,
470    ) -> glib::SignalHandlerId {
471        unsafe extern "C" fn notify_do_timestamp_trampoline<
472            F: Fn(&AppSrc) + Send + Sync + 'static,
473        >(
474            this: *mut ffi::GstAppSrc,
475            _param_spec: glib::ffi::gpointer,
476            f: glib::ffi::gpointer,
477        ) {
478            unsafe {
479                let f: &F = &*(f as *const F);
480                f(&AppSrc::from_glib_borrow(this))
481            }
482        }
483        unsafe {
484            let f: Box<F> = Box::new(f);
485            glib::signal::connect_raw(
486                self.as_ptr() as *mut _,
487                b"notify::do-timestamp\0".as_ptr() as *const _,
488                Some(mem::transmute::<*const (), unsafe extern "C" fn()>(
489                    notify_do_timestamp_trampoline::<F> as *const (),
490                )),
491                Box::into_raw(f),
492            )
493        }
494    }
495
496    #[doc(alias = "set-automatic-eos")]
497    #[doc(alias = "gst_base_src_set_automatic_eos")]
498    pub fn set_automatic_eos(&self, automatic_eos: bool) {
499        unsafe {
500            gst_base::ffi::gst_base_src_set_automatic_eos(
501                self.as_ptr() as *mut gst_base::ffi::GstBaseSrc,
502                automatic_eos.into_glib(),
503            );
504        }
505    }
506
507    #[cfg(feature = "futures")]
508    pub fn sink(&self) -> AppSrcSink {
509        AppSrcSink::new(self)
510    }
511}
512
513// rustdoc-stripper-ignore-next
514/// A [builder-pattern] type to construct [`AppSrc`] objects.
515///
516/// [builder-pattern]: https://doc.rust-lang.org/1.0.0/style/ownership/builders.html
517#[must_use = "The builder must be built to be used"]
518pub struct AppSrcBuilder<'a> {
519    builder: gst::gobject::GObjectBuilder<'a, AppSrc>,
520    callbacks: Option<AppSrcCallbacks>,
521    automatic_eos: Option<bool>,
522}
523
524impl<'a> AppSrcBuilder<'a> {
525    // rustdoc-stripper-ignore-next
526    /// Build the [`AppSrc`].
527    ///
528    /// # Panics
529    ///
530    /// This panics if the [`AppSrc`] doesn't have all the given properties or
531    /// property values of the wrong type are provided.
532    #[must_use = "Building the object from the builder is usually expensive and is not expected to have side effects"]
533    pub fn build(self) -> AppSrc {
534        let appsrc = self.builder.build().unwrap();
535
536        if let Some(callbacks) = self.callbacks {
537            appsrc.set_callbacks(callbacks);
538        }
539
540        if let Some(automatic_eos) = self.automatic_eos {
541            appsrc.set_automatic_eos(automatic_eos);
542        }
543
544        appsrc
545    }
546
547    pub fn automatic_eos(self, automatic_eos: bool) -> Self {
548        Self {
549            automatic_eos: Some(automatic_eos),
550            ..self
551        }
552    }
553
554    pub fn block(self, block: bool) -> Self {
555        Self {
556            builder: self.builder.property("block", block),
557            ..self
558        }
559    }
560
561    pub fn callbacks(self, callbacks: AppSrcCallbacks) -> Self {
562        Self {
563            callbacks: Some(callbacks),
564            ..self
565        }
566    }
567
568    pub fn caps(self, caps: &'a gst::Caps) -> Self {
569        Self {
570            builder: self.builder.property("caps", caps),
571            ..self
572        }
573    }
574
575    pub fn do_timestamp(self, do_timestamp: bool) -> Self {
576        Self {
577            builder: self.builder.property("do-timestamp", do_timestamp),
578            ..self
579        }
580    }
581
582    pub fn duration(self, duration: u64) -> Self {
583        Self {
584            builder: self.builder.property("duration", duration),
585            ..self
586        }
587    }
588
589    pub fn format(self, format: gst::Format) -> Self {
590        Self {
591            builder: self.builder.property("format", format),
592            ..self
593        }
594    }
595
596    #[cfg(feature = "v1_18")]
597    #[cfg_attr(docsrs, doc(cfg(feature = "v1_18")))]
598    pub fn handle_segment_change(self, handle_segment_change: bool) -> Self {
599        Self {
600            builder: self
601                .builder
602                .property("handle-segment-change", handle_segment_change),
603            ..self
604        }
605    }
606
607    pub fn is_live(self, is_live: bool) -> Self {
608        Self {
609            builder: self.builder.property("is-live", is_live),
610            ..self
611        }
612    }
613
614    #[cfg(feature = "v1_20")]
615    #[cfg_attr(docsrs, doc(cfg(feature = "v1_20")))]
616    pub fn leaky_type(self, leaky_type: crate::AppLeakyType) -> Self {
617        Self {
618            builder: self.builder.property("leaky-type", leaky_type),
619            ..self
620        }
621    }
622
623    #[cfg(feature = "v1_20")]
624    #[cfg_attr(docsrs, doc(cfg(feature = "v1_20")))]
625    pub fn max_buffers(self, max_buffers: u64) -> Self {
626        Self {
627            builder: self.builder.property("max-buffers", max_buffers),
628            ..self
629        }
630    }
631
632    pub fn max_bytes(self, max_bytes: u64) -> Self {
633        Self {
634            builder: self.builder.property("max-bytes", max_bytes),
635            ..self
636        }
637    }
638
639    pub fn max_latency(self, max_latency: i64) -> Self {
640        Self {
641            builder: self.builder.property("max-latency", max_latency),
642            ..self
643        }
644    }
645
646    #[cfg(feature = "v1_20")]
647    #[cfg_attr(docsrs, doc(cfg(feature = "v1_20")))]
648    pub fn max_time(self, max_time: gst::ClockTime) -> Self {
649        Self {
650            builder: self.builder.property("max-time", max_time),
651            ..self
652        }
653    }
654
655    pub fn min_latency(self, min_latency: i64) -> Self {
656        Self {
657            builder: self.builder.property("min-latency", min_latency),
658            ..self
659        }
660    }
661
662    pub fn min_percent(self, min_percent: u32) -> Self {
663        Self {
664            builder: self.builder.property("min-percent", min_percent),
665            ..self
666        }
667    }
668
669    pub fn size(self, size: i64) -> Self {
670        Self {
671            builder: self.builder.property("size", size),
672            ..self
673        }
674    }
675
676    pub fn stream_type(self, stream_type: crate::AppStreamType) -> Self {
677        Self {
678            builder: self.builder.property("stream-type", stream_type),
679            ..self
680        }
681    }
682
683    #[cfg(feature = "v1_28")]
684    #[cfg_attr(docsrs, doc(cfg(feature = "v1_28")))]
685    pub fn silent(self, silent: bool) -> Self {
686        Self {
687            builder: self.builder.property("silent", silent),
688            ..self
689        }
690    }
691
692    // rustdoc-stripper-ignore-next
693    /// Sets property `name` to the given value `value`.
694    ///
695    /// Overrides any default or previously defined value for `name`.
696    #[inline]
697    pub fn property(self, name: &'a str, value: impl Into<glib::Value> + 'a) -> Self {
698        Self {
699            builder: self.builder.property(name, value),
700            ..self
701        }
702    }
703
704    // rustdoc-stripper-ignore-next
705    /// Sets property `name` to the given string value `value`.
706    #[inline]
707    pub fn property_from_str(self, name: &'a str, value: &'a str) -> Self {
708        Self {
709            builder: self.builder.property_from_str(name, value),
710            ..self
711        }
712    }
713
714    gst::impl_builder_gvalue_extra_setters!(property_and_name);
715}
716
717#[cfg(test)]
718mod tests {
719    use gst::prelude::*;
720
721    use super::*;
722
723    #[test]
724    fn builder_caps_lt() {
725        gst::init().unwrap();
726
727        let caps = &gst::Caps::new_any();
728        {
729            let stream_type = "random-access".to_owned();
730            let appsrc = AppSrc::builder()
731                .property_from_str("stream-type", &stream_type)
732                .caps(caps)
733                .build();
734            assert_eq!(
735                appsrc.property::<crate::AppStreamType>("stream-type"),
736                crate::AppStreamType::RandomAccess
737            );
738            assert!(appsrc.property::<gst::Caps>("caps").is_any());
739        }
740
741        let stream_type = &"random-access".to_owned();
742        {
743            let caps = &gst::Caps::new_any();
744            let appsrc = AppSrc::builder()
745                .property_from_str("stream-type", stream_type)
746                .caps(caps)
747                .build();
748            assert_eq!(
749                appsrc.property::<crate::AppStreamType>("stream-type"),
750                crate::AppStreamType::RandomAccess
751            );
752            assert!(appsrc.property::<gst::Caps>("caps").is_any());
753        }
754    }
755}