Skip to main content

gstreamer_app/
app_sink_futures.rs

1// Take a look at the license at the top of the repository in the LICENSE file.
2
3use futures_core::Stream;
4use glib::object::ObjectExt as _;
5use std::{
6    pin::Pin,
7    sync::{Arc, Mutex},
8    task::Waker,
9    task::{Context, Poll},
10};
11
12use crate::{AppSink, AppSinkCallbacks};
13
14#[derive(Debug)]
15pub struct AppSinkStream {
16    app_sink: glib::WeakRef<AppSink>,
17    waker_reference: Arc<Mutex<Option<Waker>>>,
18}
19
20impl AppSinkStream {
21    pub(crate) fn new(app_sink: &AppSink) -> Self {
22        skip_assert_initialized!();
23
24        let waker_reference = Arc::new(Mutex::new(None as Option<Waker>));
25
26        app_sink.set_callbacks(
27            AppSinkCallbacks::builder()
28                .new_sample({
29                    let waker_reference = Arc::clone(&waker_reference);
30
31                    move |_| {
32                        if let Some(waker) = waker_reference.lock().unwrap().take() {
33                            waker.wake();
34                        }
35
36                        Ok(gst::FlowSuccess::Ok)
37                    }
38                })
39                .eos({
40                    let waker_reference = Arc::clone(&waker_reference);
41
42                    move |_| {
43                        if let Some(waker) = waker_reference.lock().unwrap().take() {
44                            waker.wake();
45                        }
46                    }
47                })
48                .build(),
49        );
50
51        Self {
52            app_sink: app_sink.downgrade(),
53            waker_reference,
54        }
55    }
56}
57
58impl Drop for AppSinkStream {
59    fn drop(&mut self) {
60        #[cfg(not(feature = "v1_18"))]
61        {
62            // This is not thread-safe before 1.16.3, see
63            // https://gitlab.freedesktop.org/gstreamer/gst-plugins-base/merge_requests/570
64            if gst::version() >= (1, 16, 3, 0)
65                && let Some(app_sink) = self.app_sink.upgrade()
66            {
67                app_sink.set_callbacks(AppSinkCallbacks::builder().build());
68            }
69        }
70    }
71}
72impl Stream for AppSinkStream {
73    type Item = gst::Sample;
74
75    fn poll_next(self: Pin<&mut Self>, context: &mut Context) -> Poll<Option<Self::Item>> {
76        let mut waker = self.waker_reference.lock().unwrap();
77
78        let Some(app_sink) = self.app_sink.upgrade() else {
79            return Poll::Ready(None);
80        };
81
82        app_sink
83            .try_pull_sample(gst::ClockTime::ZERO)
84            .map(|sample| Poll::Ready(Some(sample)))
85            .unwrap_or_else(|| {
86                if app_sink.is_eos() {
87                    return Poll::Ready(None);
88                }
89
90                waker.replace(context.waker().to_owned());
91
92                Poll::Pending
93            })
94    }
95}
96
97#[cfg(test)]
98mod tests {
99    use futures_util::StreamExt;
100    use gst::prelude::*;
101
102    use super::*;
103
104    #[test]
105    fn test_app_sink_stream() {
106        gst::init().unwrap();
107
108        let videotestsrc = gst::ElementFactory::make("videotestsrc")
109            .property("num-buffers", 5)
110            .build()
111            .unwrap();
112        let appsink = gst::ElementFactory::make("appsink").build().unwrap();
113
114        let pipeline = gst::Pipeline::new();
115        pipeline.add(&videotestsrc).unwrap();
116        pipeline.add(&appsink).unwrap();
117
118        videotestsrc.link(&appsink).unwrap();
119
120        let app_sink_stream = appsink.dynamic_cast::<AppSink>().unwrap().stream();
121        let samples_future = app_sink_stream.collect::<Vec<gst::Sample>>();
122
123        pipeline.set_state(gst::State::Playing).unwrap();
124        let samples = futures_executor::block_on(samples_future);
125        pipeline.set_state(gst::State::Null).unwrap();
126
127        assert_eq!(samples.len(), 5);
128    }
129}