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,
compression: config.compression.is_deflate(),
})),
}
}
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 update(&mut self, payload: impl Into<Bytes>) -> Result<()> {
self.inner.lock().unwrap().update(payload.into())
}
pub fn finish(&mut self) -> Result<()> {
self.inner.lock().unwrap().finish()
}
}
struct Inner {
track: moq_net::track::Producer,
compression: bool,
}
impl Inner {
fn update(&mut self, payload: Bytes) -> Result<()> {
let payload = match self.compression {
true => {
if payload.len() as u64 > moq_flate::DEFAULT_MAX_FRAME_SIZE {
return Err(moq_flate::Error::TooLarge(moq_flate::DEFAULT_MAX_FRAME_SIZE).into());
}
moq_flate::Encoder::new().frame(&payload)
}
false => payload,
};
if payload.len() as u64 > moq_net::group::MAX_CACHE_BYTES {
return Err(moq_net::Error::FrameTooLarge.into());
}
let mut group = self.track.append_group()?;
if let Err(err) = group.write_frame(moq_net::Timestamp::now(), payload) {
let _ = group.finish();
return Err(err.into());
}
group.finish()?;
Ok(())
}
fn finish(&mut self) -> Result<()> {
self.track.finish()?;
Ok(())
}
}