use moq_mux::catalog::hang::Extra;
use moq_mux::import;
use crate::{Error, Id, NonZeroSlab};
enum Media {
Track(Box<import::Track<Extra>>),
Container(import::Container<Extra>),
}
#[derive(Default)]
pub struct Publish {
broadcasts: NonZeroSlab<(moq_net::broadcast::Producer, moq_mux::catalog::Producer<Extra>)>,
media: NonZeroSlab<Media>,
tracks: NonZeroSlab<moq_net::track::Producer>,
groups: NonZeroSlab<moq_net::group::Producer>,
json_snapshot: NonZeroSlab<moq_json::snapshot::Producer<serde_json::Value>>,
json_stream: NonZeroSlab<moq_json::stream::Producer<serde_json::Value>>,
}
impl Publish {
pub fn create(&mut self, mut broadcast: moq_net::broadcast::Producer) -> Result<Id, Error> {
let catalog =
moq_mux::catalog::Producer::with_catalog(&mut broadcast, moq_mux::catalog::hang::Catalog::default())?;
let id = self.broadcasts.insert((broadcast, catalog))?;
Ok(id)
}
pub fn set_announce(&mut self, broadcast: Id, announce: bool) -> Result<(), Error> {
let (broadcast, _) = self.broadcasts.get_mut(broadcast).ok_or(Error::BroadcastNotFound)?;
let route = broadcast.consume().route();
broadcast.set_route(route.with_announce(announce))?;
Ok(())
}
pub fn pair_mut(
&mut self,
id: Id,
) -> Result<
(
&mut moq_net::broadcast::Producer,
&mut moq_mux::catalog::Producer<Extra>,
),
Error,
> {
let (broadcast, catalog) = self.broadcasts.get_mut(id).ok_or(Error::BroadcastNotFound)?;
Ok((broadcast, catalog))
}
pub fn finish(&mut self, broadcast: Id) -> Result<(), Error> {
let (mut broadcast, mut catalog) = self.broadcasts.remove(broadcast).ok_or(Error::BroadcastNotFound)?;
broadcast.finish();
catalog.finish()?;
Ok(())
}
pub fn media(&mut self, broadcast: Id, format: &str, init: &[u8]) -> Result<Id, Error> {
let (broadcast, catalog) = self.broadcasts.get(broadcast).ok_or(Error::BroadcastNotFound)?;
let media = match import::Container::new(broadcast.clone(), catalog.reserve(), format, init) {
Ok(container) => Media::Container(container),
Err(moq_mux::Error::UnknownFormat(_)) => {
let mut broadcast = broadcast.clone();
let name = broadcast.unique_name(&format!(".{format}"));
let request = broadcast.reserve_track(name)?;
match import::Track::new(request, catalog.reserve(), import::Init::new(format, init.to_vec())) {
Ok(track) => Media::Track(Box::new(track)),
Err(moq_mux::Error::UnknownFormat(_)) => return Err(Error::UnknownFormat(format.to_string())),
Err(err) => return Err(err.into()),
}
}
Err(err) => return Err(err.into()),
};
let id = self.media.insert(media)?;
Ok(id)
}
pub fn media_frame(&mut self, media: Id, data: &[u8], timestamp: hang::container::Timestamp) -> Result<(), Error> {
let media = self.media.get_mut(media).ok_or(Error::MediaNotFound)?;
match media {
Media::Track(track) => track.decode(data, Some(timestamp))?,
Media::Container(container) => container.decode(data)?,
}
Ok(())
}
pub fn media_finish(&mut self, media: Id) -> Result<(), Error> {
let mut media = self.media.remove(media).ok_or(Error::MediaNotFound)?;
match &mut media {
Media::Track(track) => track.finish()?,
Media::Container(container) => container.finish()?,
}
Ok(())
}
pub fn video_config(&mut self, broadcast: Id, name: &str, config: hang::catalog::VideoConfig) -> Result<(), Error> {
let (_, catalog) = self.broadcasts.get_mut(broadcast).ok_or(Error::BroadcastNotFound)?;
catalog.lock().video.insert(name, config).map_err(Error::Hang)?;
Ok(())
}
pub fn audio_config(&mut self, broadcast: Id, name: &str, config: hang::catalog::AudioConfig) -> Result<(), Error> {
let (_, catalog) = self.broadcasts.get_mut(broadcast).ok_or(Error::BroadcastNotFound)?;
catalog.lock().audio.insert(name, config).map_err(Error::Hang)?;
Ok(())
}
pub fn video_remove(&mut self, broadcast: Id, name: &str) -> Result<(), Error> {
let (_, catalog) = self.broadcasts.get_mut(broadcast).ok_or(Error::BroadcastNotFound)?;
catalog.lock().video.remove(name);
Ok(())
}
pub fn audio_remove(&mut self, broadcast: Id, name: &str) -> Result<(), Error> {
let (_, catalog) = self.broadcasts.get_mut(broadcast).ok_or(Error::BroadcastNotFound)?;
catalog.lock().audio.remove(name);
Ok(())
}
pub fn catalog_section_set(&mut self, broadcast: Id, name: &str, value: serde_json::Value) -> Result<(), Error> {
let (_, catalog) = self.broadcasts.get_mut(broadcast).ok_or(Error::BroadcastNotFound)?;
catalog.lock().set_section(name.to_string(), value)?;
Ok(())
}
pub fn catalog_section_remove(&mut self, broadcast: Id, name: &str) -> Result<(), Error> {
let (_, catalog) = self.broadcasts.get_mut(broadcast).ok_or(Error::BroadcastNotFound)?;
catalog.lock().remove_section(name);
Ok(())
}
pub fn track(&mut self, broadcast: Id, name: &str, info: Option<moq_net::track::Info>) -> Result<Id, Error> {
let (broadcast, _) = self.broadcasts.get_mut(broadcast).ok_or(Error::BroadcastNotFound)?;
let track = broadcast.create_track(name, info)?;
self.tracks.insert(track)
}
pub fn track_group(&mut self, track: Id) -> Result<Id, Error> {
let track = self.tracks.get_mut(track).ok_or(Error::TrackNotFound)?;
let group = track.append_group()?;
self.groups.insert(group)
}
pub fn track_group_at(&mut self, track: Id, sequence: u64) -> Result<Id, Error> {
let track = self.tracks.get_mut(track).ok_or(Error::TrackNotFound)?;
let group = track.create_group(moq_net::group::Info { sequence })?;
self.groups.insert(group)
}
pub fn track_frame(&mut self, track: Id, timestamp: moq_net::Timestamp, payload: &[u8]) -> Result<(), Error> {
let track = self.tracks.get_mut(track).ok_or(Error::TrackNotFound)?;
track.write_frame(timestamp, bytes::Bytes::copy_from_slice(payload))?;
Ok(())
}
pub fn track_datagram(&mut self, track: Id, timestamp_us: u64, payload: &[u8]) -> Result<u64, Error> {
let track = self.tracks.get_mut(track).ok_or(Error::TrackNotFound)?;
let timestamp = moq_net::Timestamp::from_micros(timestamp_us)?;
Ok(track.append_datagram(timestamp, bytes::Bytes::copy_from_slice(payload))?)
}
pub fn track_finish(&mut self, track: Id) -> Result<(), Error> {
let mut track = self.tracks.remove(track).ok_or(Error::TrackNotFound)?;
if track.final_sequence().is_none() {
track.finish()?;
}
Ok(())
}
pub fn track_finish_at(&mut self, track: Id, final_sequence: u64) -> Result<(), Error> {
let track = self.tracks.get_mut(track).ok_or(Error::TrackNotFound)?;
track.finish_at(final_sequence)?;
Ok(())
}
pub fn track_abort(&mut self, track: Id, error_code: u16) -> Result<(), Error> {
let track = self.tracks.remove(track).ok_or(Error::TrackNotFound)?;
track.abort(moq_net::Error::App(error_code))?;
Ok(())
}
pub fn json_snapshot(
&mut self,
broadcast: Id,
name: &str,
config: moq_json::snapshot::ProducerConfig,
) -> Result<Id, Error> {
let (broadcast, _) = self.broadcasts.get_mut(broadcast).ok_or(Error::BroadcastNotFound)?;
let track = broadcast.create_track(name, None)?;
let producer = moq_json::snapshot::Producer::new(track, config);
self.json_snapshot.insert(producer)
}
pub fn json_snapshot_update(&mut self, json: Id, value: serde_json::Value) -> Result<(), Error> {
let producer = self.json_snapshot.get_mut(json).ok_or(Error::TrackNotFound)?;
producer.update(&value)?;
Ok(())
}
pub fn json_snapshot_finish(&mut self, json: Id) -> Result<(), Error> {
let mut producer = self.json_snapshot.remove(json).ok_or(Error::TrackNotFound)?;
producer.finish()?;
Ok(())
}
pub fn json_stream(
&mut self,
broadcast: Id,
name: &str,
config: moq_json::stream::ProducerConfig,
) -> Result<Id, Error> {
let (broadcast, _) = self.broadcasts.get_mut(broadcast).ok_or(Error::BroadcastNotFound)?;
let track = broadcast.create_track(name, None)?;
let producer = moq_json::stream::Producer::new(track, config);
self.json_stream.insert(producer)
}
pub fn json_stream_append(&mut self, stream: Id, value: serde_json::Value) -> Result<(), Error> {
let producer = self.json_stream.get_mut(stream).ok_or(Error::TrackNotFound)?;
producer.append(&value)?;
Ok(())
}
pub fn json_stream_finish(&mut self, stream: Id) -> Result<(), Error> {
let mut producer = self.json_stream.remove(stream).ok_or(Error::TrackNotFound)?;
producer.finish()?;
Ok(())
}
pub fn group_frame(&mut self, group: Id, timestamp: moq_net::Timestamp, payload: &[u8]) -> Result<(), Error> {
let group = self.groups.get_mut(group).ok_or(Error::GroupNotFound)?;
group.write_frame(timestamp, bytes::Bytes::copy_from_slice(payload))?;
Ok(())
}
pub fn group_finish(&mut self, group: Id) -> Result<(), Error> {
let mut group = self.groups.remove(group).ok_or(Error::GroupNotFound)?;
group.finish()?;
Ok(())
}
pub fn group_abort(&mut self, group: Id, error_code: u16) -> Result<(), Error> {
let group = self.groups.remove(group).ok_or(Error::GroupNotFound)?;
group.abort(moq_net::Error::App(error_code))?;
Ok(())
}
}