gstreamer_app/
app_sink_futures.rs1use 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 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}