use std::sync::{Arc, Mutex};
use bytes::Bytes;
use crate::Result;
pub use super::Config;
#[derive(Clone)]
pub struct Producer {
inner: Arc<Mutex<Inner>>,
}
impl Producer {
pub fn new(track: moq_net::track::Producer, config: Config) -> Self {
Self {
inner: Arc::new(Mutex::new(Inner {
track,
group: None,
flate: config.compression.is_deflate().then(moq_flate::Encoder::new),
})),
}
}
pub fn consume(&self) -> moq_net::track::Subscriber {
self.inner.lock().unwrap().track.subscribe(None)
}
pub fn is_used(&self) -> bool {
self.inner.lock().unwrap().track.is_used()
}
pub fn append(&mut self, payload: impl Into<Bytes>) -> Result<()> {
self.inner.lock().unwrap().append(payload.into())
}
pub fn finish(&mut self) -> Result<()> {
self.inner.lock().unwrap().finish()
}
}
struct Inner {
track: moq_net::track::Producer,
group: Option<moq_net::group::Producer>,
flate: Option<moq_flate::Encoder>,
}
impl Inner {
fn append(&mut self, payload: Bytes) -> Result<()> {
if self.flate.is_some() && payload.len() as u64 > moq_flate::DEFAULT_MAX_FRAME_SIZE {
self.abort(moq_net::Error::FrameTooLarge);
return Err(moq_flate::Error::TooLarge(moq_flate::DEFAULT_MAX_FRAME_SIZE).into());
}
if self.group.is_none() {
self.group = Some(self.track.append_group()?);
}
let payload = match self.flate.as_mut() {
Some(flate) => flate.frame(&payload),
None => payload,
};
let group = self.group.as_mut().expect("a group is open");
let Err(err) = group.write_frame(moq_net::Timestamp::now(), payload) else {
return Ok(());
};
self.abort(err.clone());
Err(err.into())
}
fn abort(&mut self, err: moq_net::Error) {
if let Some(group) = self.group.take() {
let _ = group.abort(err.clone());
}
let _ = self.track.clone().abort(err);
}
fn finish(&mut self) -> Result<()> {
let group = match self.group.take() {
Some(group) => group.finish(),
None => Ok(()),
};
let track = self.track.finish();
group?;
track?;
Ok(())
}
}