1use 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 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 #[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 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 #[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 #[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#[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 #[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 #[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 #[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}