use std::marker::PhantomData;
use std::sync::{Arc, Mutex};
use serde::Serialize;
use serde_json::Value;
use super::{Encoded, Encoder, ProducerConfig};
use crate::Result;
pub struct Producer<T> {
inner: Arc<Mutex<Inner<T>>>,
_marker: PhantomData<fn(T)>,
}
impl<T> Clone for Producer<T> {
fn clone(&self) -> Self {
Self {
inner: self.inner.clone(),
_marker: PhantomData,
}
}
}
impl<T> Producer<T> {
pub fn new(track: moq_net::track::Producer, config: ProducerConfig) -> Self {
Self {
inner: Arc::new(Mutex::new(Inner {
track: Track {
inner: track,
group: None,
},
encoder: Encoder::new(config),
finished: false,
})),
_marker: PhantomData,
}
}
pub fn consume(&self) -> moq_net::track::Subscriber {
self.inner.lock().unwrap().track.inner.subscribe(None)
}
pub fn window(&self) -> Vec<Value> {
self.inner.lock().unwrap().encoder.window()
}
pub fn range(&self) -> std::ops::Range<u64> {
self.inner.lock().unwrap().encoder.range()
}
pub fn pop(&mut self, count: u64) -> Result<()> {
self.inner.lock().unwrap().pop(count)
}
pub fn finish(self) -> Result<()> {
self.inner.lock().unwrap().finish()
}
}
impl<T: Serialize> Producer<T> {
pub fn push(&mut self, value: &T) -> Result<()> {
self.inner.lock().unwrap().push(value)
}
}
struct Inner<T> {
track: Track,
encoder: Encoder<T>,
finished: bool,
}
impl<T> Inner<T> {
fn pop(&mut self, count: u64) -> Result<()> {
self.ensure_open()?;
let Inner { track, encoder, .. } = self;
let Some(frame) = encoder.pop(count)? else {
return Ok(());
};
track.write(&frame)?;
frame.commit();
Ok(())
}
fn finish(&mut self) -> Result<()> {
if self.finished {
return Ok(());
}
self.finished = true;
self.track.finish()
}
fn ensure_open(&self) -> Result<()> {
if self.finished {
return Err(moq_net::Error::Closed.into());
}
Ok(())
}
}
impl<T: Serialize> Inner<T> {
fn push(&mut self, value: &T) -> Result<()> {
self.ensure_open()?;
let Inner { track, encoder, .. } = self;
let frame = encoder.push(value)?;
track.write(&frame)?;
frame.commit();
Ok(())
}
}
struct Track {
inner: moq_net::track::Producer,
group: Option<moq_net::group::Producer>,
}
impl Track {
fn write(&mut self, encoded: &Encoded) -> Result<()> {
match encoded.keyframe {
true => self.write_header(encoded.payload.clone()),
false => self.write_op(encoded.payload.clone()),
}
}
fn write_header(&mut self, payload: bytes::Bytes) -> Result<()> {
if let Some(mut group) = self.group.take() {
group.finish()?;
}
let mut group = self.inner.append_group()?;
if let Err(err) = group.write_frame(moq_net::Timestamp::now(), payload) {
let _ = group.finish();
return Err(err.into());
}
self.group = Some(group);
Ok(())
}
fn write_op(&mut self, payload: bytes::Bytes) -> Result<()> {
self.group
.as_mut()
.expect("the encoder only emits an op after a header opened a group")
.write_frame(moq_net::Timestamp::now(), payload)?;
Ok(())
}
fn finish(&mut self) -> Result<()> {
if let Some(mut group) = self.group.take() {
group.finish()?;
}
self.inner.finish()?;
Ok(())
}
}