use std::marker::PhantomData;
use std::sync::{Arc, Mutex};
use serde::Serialize;
use super::Encoder;
use crate::Result;
pub use super::Config;
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 consume(&self) -> moq_net::track::Subscriber {
self.inner.lock().unwrap().track.inner.subscribe(None)
}
pub fn demand(&self) -> moq_net::track::Demand {
self.inner.lock().unwrap().track.inner.demand()
}
}
impl<T: Serialize> Producer<T> {
pub fn new(track: moq_net::track::Producer, config: Config) -> Self {
Self {
inner: Arc::new(Mutex::new(Inner {
track: Track {
inner: track,
group: None,
},
encoder: Encoder::new(config),
})),
_marker: PhantomData,
}
}
pub fn append(&mut self, value: &T) -> Result<()> {
self.inner.lock().unwrap().append(value)
}
pub fn finish(&mut self) -> Result<()> {
self.inner.lock().unwrap().finish()
}
}
struct Inner<T> {
track: Track,
encoder: Encoder<T>,
}
impl<T: Serialize> Inner<T> {
fn append(&mut self, value: &T) -> Result<()> {
let Inner { track, encoder } = self;
let record = match encoder.encode(value) {
Ok(record) => record,
Err(err) => {
track.abort(moq_net::Error::Cancel);
return Err(err);
}
};
let opened = track.open();
let published = opened.is_ok();
let result = match opened {
Ok(()) => track.write(record.payload()),
Err(err) => Err(err),
};
if let Err(err) = result {
drop(record);
encoder.reset();
if published {
track.abort(err.clone());
}
return Err(err.into());
}
record.commit();
Ok(())
}
fn finish(&mut self) -> Result<()> {
Ok(self.track.finish()?)
}
}
struct Track {
inner: moq_net::track::Producer,
group: Option<moq_net::group::Producer>,
}
impl Track {
fn open(&mut self) -> std::result::Result<(), moq_net::Error> {
if self.group.is_none() {
self.group = Some(self.inner.append_group()?);
}
Ok(())
}
fn write(&mut self, payload: &bytes::Bytes) -> std::result::Result<(), moq_net::Error> {
let group = self.group.as_mut().expect("a group is open");
group.write_frame(moq_net::Timestamp::now(), payload.clone())
}
fn abort(&mut self, err: moq_net::Error) {
if let Some(group) = self.group.take() {
let _ = group.abort(err.clone());
}
let _ = self.inner.clone().abort(err);
}
fn finish(&mut self) -> std::result::Result<(), moq_net::Error> {
let group = match self.group.take() {
Some(group) => group.finish(),
None => Ok(()),
};
let track = self.inner.finish();
group?;
track?;
Ok(())
}
}