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