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