use super::image_from_gst_buffer;
use crate::stream::error::StreamCaptureError;
use circular_buffer::FixedCircularBuffer;
use gstreamer::prelude::*;
use kornia_image::{Image, ImageSize};
use std::sync::{Arc, Mutex};
struct FrameBuffer {
buffer: gstreamer::Buffer,
width: i32,
height: i32,
}
pub enum StreamerState {
Null,
Ready,
Paused,
Playing,
}
impl From<gstreamer::State> for StreamerState {
fn from(value: gstreamer::State) -> Self {
match value {
gstreamer::State::VoidPending => StreamerState::Null,
gstreamer::State::Null => StreamerState::Null,
gstreamer::State::Ready => StreamerState::Ready,
gstreamer::State::Paused => StreamerState::Paused,
gstreamer::State::Playing => StreamerState::Playing,
}
}
}
pub struct StreamCapture {
pub(crate) pipeline: gstreamer::Pipeline,
circular_buffer: Arc<Mutex<FixedCircularBuffer<FrameBuffer, 5>>>,
fps: Arc<Mutex<gstreamer::Fraction>>,
}
impl StreamCapture {
pub fn new(pipeline_desc: &str) -> Result<Self, StreamCaptureError> {
if !gstreamer::INITIALIZED.load(std::sync::atomic::Ordering::Relaxed) {
gstreamer::init()?;
}
let pipeline = gstreamer::parse::launch(pipeline_desc)?
.dynamic_cast::<gstreamer::Pipeline>()
.map_err(StreamCaptureError::DowncastPipelineError)?;
let appsink = pipeline
.by_name("sink")
.ok_or_else(|| StreamCaptureError::GetElementByNameError)?
.dynamic_cast::<gstreamer_app::AppSink>()
.map_err(StreamCaptureError::DowncastPipelineError)?;
let circular_buffer = Arc::new(Mutex::new(FixedCircularBuffer::new()));
let fps = Arc::new(Mutex::new(gstreamer::Fraction::new(1, 1)));
appsink.set_callbacks(
gstreamer_app::AppSinkCallbacks::builder()
.new_sample({
let circular_buffer = circular_buffer.clone();
let fps = fps.clone();
move |sink| {
Self::extract_frame_buffer(sink)
.map_err(|_| gstreamer::FlowError::Eos)
.and_then(|(frame_buffer, fps_fraction)| {
circular_buffer
.lock()
.map_err(|_| gstreamer::FlowError::Error)?
.push_back(frame_buffer);
*fps.lock().map_err(|_| gstreamer::FlowError::Error)? =
fps_fraction;
Ok(gstreamer::FlowSuccess::Ok)
})
}
})
.build(),
);
Ok(Self {
pipeline,
circular_buffer,
fps,
})
}
pub fn get_fps(&self) -> Option<f64> {
self.fps
.lock()
.ok()
.map(|fps| fps.numer() as f64 / fps.denom() as f64)
}
pub fn get_state(&self) -> StreamerState {
self.pipeline.current_state().into()
}
pub fn start(&self) -> Result<(), StreamCaptureError> {
self.circular_buffer
.lock()
.map_err(|_| StreamCaptureError::MutexPoisonError)?
.clear();
self.pipeline.set_state(gstreamer::State::Playing)?;
Ok(())
}
pub fn grab_rgb8(&mut self) -> Result<Option<Image<u8, 3>>, StreamCaptureError> {
let mut circular_buffer = self
.circular_buffer
.lock()
.map_err(|_| StreamCaptureError::MutexPoisonError)?;
let Some(frame_buffer) = circular_buffer.pop_front() else {
return Ok(None);
};
let width = frame_buffer.width;
let height = frame_buffer.height;
let buffer = frame_buffer.buffer;
let mapped_buffer = buffer
.into_mapped_buffer_readable()
.map_err(|_| StreamCaptureError::GetBufferError)?;
let image = image_from_gst_buffer(
ImageSize {
width: width as usize,
height: height as usize,
},
mapped_buffer,
)?;
Ok(Some(image))
}
pub fn close(&self) -> Result<(), StreamCaptureError> {
let res = self.pipeline.send_event(gstreamer::event::Eos::new());
if !res {
return Err(StreamCaptureError::SendEosError);
}
self.pipeline.set_state(gstreamer::State::Null)?;
self.circular_buffer
.lock()
.map_err(|_| StreamCaptureError::MutexPoisonError)?
.clear();
Ok(())
}
fn extract_frame_buffer(
appsink: &gstreamer_app::AppSink,
) -> Result<(FrameBuffer, gstreamer::Fraction), StreamCaptureError> {
let sample = appsink.pull_sample()?;
let caps = sample.caps().ok_or_else(|| {
StreamCaptureError::GetCapsError("Failed to get the caps".to_string())
})?;
let structure = caps.structure(0).ok_or_else(|| {
StreamCaptureError::GetCapsError("Failed to get the structure".to_string())
})?;
let height = structure
.get::<i32>("height")
.map_err(|e| StreamCaptureError::GetCapsError(e.to_string()))?;
let width = structure
.get::<i32>("width")
.map_err(|e| StreamCaptureError::GetCapsError(e.to_string()))?;
let fps = structure
.get::<gstreamer::Fraction>("framerate")
.map_err(|e| StreamCaptureError::GetCapsError(e.to_string()))?;
let buffer = sample
.buffer_owned()
.ok_or_else(|| StreamCaptureError::GetBufferError)?;
let frame_buffer = FrameBuffer {
buffer,
width,
height,
};
Ok((frame_buffer, fps))
}
}
impl Drop for StreamCapture {
fn drop(&mut self) {
if let Err(e) = self.close() {
log::warn!("Failed to close stream safely on drop: {:?}", e);
}
}
}