use std::collections::BTreeMap;
use moq_mux::catalog::hang::Extra;
use moq_mux::import;
use tokio::sync::oneshot;
use crate::ffi::OnStatus;
use crate::{Error, Id, NonZeroSlab, State, moq_demand};
struct TaskEntry {
close: Option<oneshot::Sender<()>>,
callback: OnStatus,
}
struct TrackRequest {
broadcast: Id,
request: moq_net::track::Request,
}
enum Request {
Track(TrackRequest),
Group(moq_net::group::Request),
}
enum Dynamic {
Broadcast(moq_net::broadcast::Dynamic, Id),
Track(moq_net::track::Dynamic),
}
struct Broadcast {
producer: moq_net::broadcast::Producer,
catalog: moq_mux::catalog::Producer<Extra>,
video: BTreeMap<String, hang::catalog::VideoConfig>,
audio: BTreeMap<String, hang::catalog::AudioConfig>,
}
#[derive(Default)]
pub struct Publish {
broadcasts: NonZeroSlab<Broadcast>,
media: NonZeroSlab<Box<import::Track>>,
containers: NonZeroSlab<import::Container<Extra>>,
tracks: NonZeroSlab<moq_net::track::Producer>,
groups: NonZeroSlab<moq_net::group::Producer>,
json_snapshot: NonZeroSlab<moq_mux::json::Snapshot<serde_json::Value, Extra>>,
json_stream: NonZeroSlab<moq_mux::json::Stream<serde_json::Value, Extra>>,
binary_snapshot: NonZeroSlab<moq_mux::binary::Snapshot<Extra>>,
binary_stream: NonZeroSlab<moq_mux::binary::Stream<Extra>>,
demand: NonZeroSlab<Option<TaskEntry>>,
dynamic: NonZeroSlab<Option<TaskEntry>>,
track_request: NonZeroSlab<TrackRequest>,
group_request: NonZeroSlab<moq_net::group::Request>,
}
impl Publish {
pub fn create(&mut self, mut broadcast: moq_net::broadcast::Producer) -> Result<Id, Error> {
let config = moq_mux::catalog::Config::default()
.with_catalog(moq_mux::catalog::hang::Catalog::<moq_mux::catalog::hang::Extra>::default());
let catalog = moq_mux::catalog::Producer::new(&mut broadcast, config)?;
let id = self.broadcasts.insert(Broadcast {
producer: broadcast,
catalog,
video: BTreeMap::new(),
audio: BTreeMap::new(),
})?;
Ok(id)
}
pub fn announce(&mut self, broadcast: Id, route: moq_net::origin::Route) -> Result<(), Error> {
let broadcast = self.broadcasts.get_mut(broadcast).ok_or(Error::BroadcastNotFound)?;
broadcast.producer.announce(route)?;
Ok(())
}
pub fn unannounce(&mut self, broadcast: Id) -> Result<(), Error> {
let broadcast = self.broadcasts.get_mut(broadcast).ok_or(Error::BroadcastNotFound)?;
broadcast.producer.unannounce();
Ok(())
}
pub(crate) fn producer(&mut self, id: Id) -> Result<&mut moq_net::broadcast::Producer, Error> {
Ok(&mut self.broadcasts.get_mut(id).ok_or(Error::BroadcastNotFound)?.producer)
}
fn catalog(&mut self, id: Id) -> Result<&mut moq_mux::catalog::Producer<Extra>, Error> {
Ok(&mut self.broadcasts.get_mut(id).ok_or(Error::BroadcastNotFound)?.catalog)
}
#[cfg(test)]
pub fn catalog_snapshot(&mut self, id: Id) -> Result<moq_mux::catalog::hang::Catalog<Extra>, Error> {
Ok(self.catalog(id)?.snapshot())
}
pub fn pair_mut(
&mut self,
id: Id,
) -> Result<
(
&mut moq_net::broadcast::Producer,
&mut moq_mux::catalog::Producer<Extra>,
),
Error,
> {
let broadcast = self.broadcasts.get_mut(id).ok_or(Error::BroadcastNotFound)?;
Ok((&mut broadcast.producer, &mut broadcast.catalog))
}
pub fn finish(&mut self, broadcast: Id) -> Result<(), Error> {
let Broadcast {
producer,
mut catalog,
video,
audio,
..
} = self.broadcasts.remove(broadcast).ok_or(Error::BroadcastNotFound)?;
{
let mut guard = catalog.modify()?;
for name in video.keys() {
guard.video.renditions.remove(name);
}
for name in audio.keys() {
guard.audio.renditions.remove(name);
}
guard.commit()?;
}
producer.finish();
catalog.finish()?;
Ok(())
}
pub fn audio(&mut self, broadcast: Id, init: import::AudioInit) -> Result<Id, Error> {
let Broadcast {
producer: broadcast,
catalog,
..
} = self.broadcasts.get(broadcast).ok_or(Error::BroadcastNotFound)?;
let broadcast = broadcast.clone();
let name = broadcast.unique_name(&format!(".{}", init.format));
let request = broadcast.reserve_track(name)?;
let track = import::Track::audio(request, catalog.reserve(), init)?;
let id = self.media.insert(Box::new(track))?;
Ok(id)
}
pub fn video(&mut self, broadcast: Id, init: import::VideoInit) -> Result<Id, Error> {
let Broadcast {
producer: broadcast,
catalog,
..
} = self.broadcasts.get(broadcast).ok_or(Error::BroadcastNotFound)?;
let broadcast = broadcast.clone();
let name = broadcast.unique_name(&format!(".{}", init.format));
let request = broadcast.reserve_track(name)?;
let track = import::Track::video(request, catalog.reserve(), init)?;
let id = self.media.insert(Box::new(track))?;
Ok(id)
}
pub fn container(&mut self, broadcast: Id, init: import::ContainerInit) -> Result<Id, Error> {
let Broadcast {
producer: broadcast,
catalog,
..
} = self.broadcasts.get(broadcast).ok_or(Error::BroadcastNotFound)?;
let container = import::Container::new(broadcast.clone(), catalog.reserve(), &init)?;
let id = self.containers.insert(container)?;
Ok(id)
}
pub fn media_frame(&mut self, media: Id, data: &[u8], timestamp: hang::container::Timestamp) -> Result<(), Error> {
let track = self.media.get_mut(media).ok_or(Error::MediaNotFound)?;
track.decode(data, Some(timestamp))?;
Ok(())
}
pub fn media_flush(&mut self, media: Id, timestamp: hang::container::Timestamp) -> Result<(), Error> {
let track = self.media.get_mut(media).ok_or(Error::MediaNotFound)?;
track.flush(timestamp, std::time::Instant::now())?;
Ok(())
}
pub fn media_cut(&mut self, media: Id) -> Result<(), Error> {
let track = self.media.get_mut(media).ok_or(Error::MediaNotFound)?;
track.cut(None)?;
Ok(())
}
pub fn media_seek(&mut self, media: Id, sequence: u64) -> Result<(), Error> {
let track = self.media.get_mut(media).ok_or(Error::MediaNotFound)?;
track.seek(sequence)?;
Ok(())
}
pub fn media_finish(&mut self, media: Id) -> Result<(), Error> {
let mut track = self.media.remove(media).ok_or(Error::MediaNotFound)?;
track.finish()?;
Ok(())
}
pub fn container_write(&mut self, container: Id, data: &[u8]) -> Result<(), Error> {
let container = self.containers.get_mut(container).ok_or(Error::MediaNotFound)?;
container.decode(data)?;
Ok(())
}
pub fn container_cut(&mut self, container: Id) -> Result<(), Error> {
let container = self.containers.get_mut(container).ok_or(Error::MediaNotFound)?;
container.cut();
Ok(())
}
pub fn container_seek(&mut self, container: Id, sequence: u64) -> Result<(), Error> {
let container = self.containers.get_mut(container).ok_or(Error::MediaNotFound)?;
container.seek(sequence)?;
Ok(())
}
pub fn container_finish(&mut self, container: Id) -> Result<(), Error> {
let mut container = self.containers.remove(container).ok_or(Error::MediaNotFound)?;
container.finish()?;
Ok(())
}
pub fn video_config(&mut self, broadcast: Id, name: &str, config: hang::catalog::VideoConfig) -> Result<(), Error> {
let broadcast = self.broadcasts.get_mut(broadcast).ok_or(Error::BroadcastNotFound)?;
if !broadcast.video.contains_key(name) && broadcast.catalog.is_claimed::<hang::catalog::VideoConfig>(name) {
return Err(Error::Hang(hang::Error::Duplicate(name.to_string())));
}
broadcast.video.insert(name.to_string(), config.clone());
let mut catalog = broadcast.catalog.modify()?;
catalog.video.renditions.insert(name.to_string(), config);
catalog.commit()?;
Ok(())
}
pub fn audio_config(&mut self, broadcast: Id, name: &str, config: hang::catalog::AudioConfig) -> Result<(), Error> {
let broadcast = self.broadcasts.get_mut(broadcast).ok_or(Error::BroadcastNotFound)?;
if !broadcast.audio.contains_key(name) && broadcast.catalog.is_claimed::<hang::catalog::AudioConfig>(name) {
return Err(Error::Hang(hang::Error::Duplicate(name.to_string())));
}
broadcast.audio.insert(name.to_string(), config.clone());
let mut catalog = broadcast.catalog.modify()?;
catalog.audio.renditions.insert(name.to_string(), config);
catalog.commit()?;
Ok(())
}
pub fn video_remove(&mut self, broadcast: Id, name: &str) -> Result<(), Error> {
let broadcast = self.broadcasts.get_mut(broadcast).ok_or(Error::BroadcastNotFound)?;
if broadcast.video.remove(name).is_some() {
let mut catalog = broadcast.catalog.modify()?;
catalog.video.renditions.remove(name);
catalog.commit()?;
}
Ok(())
}
pub fn audio_remove(&mut self, broadcast: Id, name: &str) -> Result<(), Error> {
let broadcast = self.broadcasts.get_mut(broadcast).ok_or(Error::BroadcastNotFound)?;
if broadcast.audio.remove(name).is_some() {
let mut catalog = broadcast.catalog.modify()?;
catalog.audio.renditions.remove(name);
catalog.commit()?;
}
Ok(())
}
pub fn video_properties(&mut self, broadcast: Id, properties: hang::catalog::VideoProperties) -> Result<(), Error> {
let catalog = self.catalog(broadcast)?;
let mut catalog = catalog.modify()?;
catalog.video.set_properties(properties)?;
catalog.commit()?;
Ok(())
}
pub fn catalog_section_set(&mut self, broadcast: Id, name: &str, value: serde_json::Value) -> Result<(), Error> {
let catalog = self.catalog(broadcast)?;
let mut guard = catalog.modify()?;
guard.set_section(name.to_string(), value)?;
guard.commit()?;
Ok(())
}
pub fn catalog_section_remove(&mut self, broadcast: Id, name: &str) -> Result<(), Error> {
let catalog = self.catalog(broadcast)?;
let mut guard = catalog.modify()?;
guard.remove_section(name);
guard.commit()?;
Ok(())
}
pub fn track_demand(&self, track: Id) -> Result<moq_net::track::Demand, Error> {
Ok(self.tracks.get(track).ok_or(Error::TrackNotFound)?.demand())
}
pub fn media_demand(&self, media: Id) -> Result<moq_net::track::Demand, Error> {
Ok(self.media.get(media).ok_or(Error::MediaNotFound)?.demand())
}
pub fn demand(&mut self, demand: moq_net::track::Demand, on_demand: OnStatus) -> Result<Id, Error> {
let channel = oneshot::channel();
let id = self.demand.insert(Some(TaskEntry {
close: Some(channel.0),
callback: on_demand,
}))?;
tokio::spawn(async move {
let res = Self::run_demand(on_demand, demand, channel.1).await;
let entry = State::lock().publish.demand.remove(id).flatten();
if let Some(entry) = entry {
entry.callback.call(res);
}
});
Ok(id)
}
pub(crate) async fn run_demand(
callback: OnStatus,
demand: moq_net::track::Demand,
mut close: oneshot::Receiver<()>,
) -> Result<(), Error> {
let mut used = tokio::select! {
res = demand.used() => match res {
Ok(()) => true,
Err(moq_net::Error::Dropped) => return Ok(()),
Err(err) => return Err(err.into()),
},
res = demand.unused() => match res {
Ok(()) => false,
Err(moq_net::Error::Dropped) => return Ok(()),
Err(err) => return Err(err.into()),
},
};
loop {
let state = if used {
moq_demand::MOQ_DEMAND_USED
} else {
moq_demand::MOQ_DEMAND_UNUSED
};
callback.call(state as i32);
let flipped = async {
if used {
demand.unused().await
} else {
demand.used().await
}
};
tokio::select! {
biased;
_ = &mut close => return Ok(()),
res = flipped => match res {
Ok(()) => used = !used,
Err(moq_net::Error::Dropped) => return Ok(()),
Err(err) => return Err(err.into()),
},
}
}
}
pub fn demand_close(&mut self, watcher: Id) -> Result<(), Error> {
self.demand
.get_mut(watcher)
.and_then(|entry| entry.as_mut())
.ok_or(Error::NotFound)?
.close
.take()
.ok_or(Error::NotFound)?;
Ok(())
}
pub fn dynamic(&mut self, broadcast: Id, on_request: OnStatus) -> Result<Id, Error> {
let dynamic = self.producer(broadcast)?.dynamic();
self.spawn_dynamic(Dynamic::Broadcast(dynamic, broadcast), on_request)
}
pub fn track_dynamic(&mut self, track: Id, on_group: OnStatus) -> Result<Id, Error> {
let dynamic = self.tracks.get(track).ok_or(Error::TrackNotFound)?.dynamic();
self.spawn_dynamic(Dynamic::Track(dynamic), on_group)
}
pub fn track_request_dynamic(&mut self, request: Id, on_group: OnStatus) -> Result<Id, Error> {
let dynamic = self
.track_request
.get(request)
.ok_or(Error::NotFound)?
.request
.dynamic();
self.spawn_dynamic(Dynamic::Track(dynamic), on_group)
}
fn spawn_dynamic(&mut self, dynamic: Dynamic, callback: OnStatus) -> Result<Id, Error> {
let channel = oneshot::channel();
let id = self.dynamic.insert(Some(TaskEntry {
close: Some(channel.0),
callback,
}))?;
tokio::spawn(async move {
let res = Self::run_dynamic(callback, dynamic, channel.1).await;
let entry = State::lock().publish.dynamic.remove(id).flatten();
if let Some(entry) = entry {
entry.callback.call(res);
}
});
Ok(id)
}
async fn run_dynamic(
callback: OnStatus,
mut dynamic: Dynamic,
mut close: oneshot::Receiver<()>,
) -> Result<(), Error> {
loop {
let res = tokio::select! {
biased;
_ = &mut close => return Ok(()),
res = kio::wait(|waiter| match &mut dynamic {
Dynamic::Broadcast(dynamic, broadcast) => dynamic
.poll_requested_track(waiter)
.map_ok(|request| Request::Track(TrackRequest { broadcast: *broadcast, request })),
Dynamic::Track(dynamic) => dynamic.poll_requested_group(waiter).map_ok(Request::Group),
}) => res,
};
let request = match res {
Ok(request) => request,
Err(moq_net::Error::Closed | moq_net::Error::Dropped) => return Ok(()),
Err(err) => return Err(err.into()),
};
let mut state = State::lock();
let id = match request {
Request::Track(request) => state.publish.track_request.insert(request)?,
Request::Group(request) => state.publish.group_request.insert(request)?,
};
drop(state);
callback.call(id);
}
}
pub fn dynamic_close(&mut self, dynamic: Id) -> Result<(), Error> {
self.dynamic
.get_mut(dynamic)
.and_then(|entry| entry.as_mut())
.ok_or(Error::NotFound)?
.close
.take()
.ok_or(Error::NotFound)?;
Ok(())
}
pub fn track_request_name(&self, request: Id, dst: &mut crate::moq_string) -> Result<(), Error> {
let name = self.track_request.get(request).ok_or(Error::NotFound)?.request.name();
*dst = crate::moq_string {
data: name.as_ptr().cast::<std::ffi::c_char>(),
len: name.len(),
};
Ok(())
}
pub fn track_request_accept(&mut self, request: Id, info: moq_net::track::Info) -> Result<Id, Error> {
let request = self.track_request.remove(request).ok_or(Error::NotFound)?;
self.tracks.insert(request.request.accept(info))
}
pub fn track_request_audio(&mut self, request: Id, init: import::AudioInit) -> Result<Id, Error> {
let TrackRequest { broadcast, request } = self.track_request.remove(request).ok_or(Error::NotFound)?;
let catalog = self.catalog(broadcast)?;
let track = import::Track::audio(request, catalog.reserve(), init)?;
self.media.insert(Box::new(track))
}
pub fn track_request_video(&mut self, request: Id, init: import::VideoInit) -> Result<Id, Error> {
let TrackRequest { broadcast, request } = self.track_request.remove(request).ok_or(Error::NotFound)?;
let catalog = self.catalog(broadcast)?;
let track = import::Track::video(request, catalog.reserve(), init)?;
self.media.insert(Box::new(track))
}
pub fn track_request_abort(&mut self, request: Id, error_code: u16) -> Result<(), Error> {
let request = self.track_request.remove(request).ok_or(Error::NotFound)?;
request.request.reject(moq_net::Error::App(error_code));
Ok(())
}
pub fn track_request_free(&mut self, request: Id) -> Result<(), Error> {
self.track_request.remove(request).ok_or(Error::NotFound)?;
Ok(())
}
pub fn group_request_info(&self, request: Id) -> Result<(u64, u8, u64), Error> {
let request = self.group_request.get(request).ok_or(Error::NotFound)?;
Ok((request.sequence(), request.priority(), request.frame_start()))
}
pub fn group_request_accept(&mut self, request: Id) -> Result<Id, Error> {
let request = self.group_request.remove(request).ok_or(Error::NotFound)?;
let frame_start = request.frame_start();
let mut group = request.accept(None)?;
if let Err(err) = group.start_at(frame_start) {
let _ = group.abort(err.clone());
return Err(err.into());
}
self.groups.insert(group)
}
pub fn group_request_abort(&mut self, request: Id, error_code: u16) -> Result<(), Error> {
let request = self.group_request.remove(request).ok_or(Error::NotFound)?;
request.reject(moq_net::Error::App(error_code));
Ok(())
}
pub fn group_request_free(&mut self, request: Id) -> Result<(), Error> {
self.group_request.remove(request).ok_or(Error::NotFound)?;
Ok(())
}
pub fn track(&mut self, broadcast: Id, name: &str, info: Option<moq_net::track::Info>) -> Result<Id, Error> {
let broadcast = self.producer(broadcast)?;
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 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(())
}
fn data_track<T>(
&mut self,
broadcast: Id,
name: &str,
publish: impl FnOnce(&moq_mux::catalog::Producer<Extra>, moq_net::track::Producer) -> moq_mux::Result<T>,
) -> Result<T, Error> {
let broadcast = self.broadcasts.get_mut(broadcast).ok_or(Error::BroadcastNotFound)?;
let track = broadcast.producer.create_track(name, None)?;
Ok(publish(&broadcast.catalog, track)?)
}
pub fn json_snapshot(&mut self, broadcast: Id, name: &str, config: moq_mux::json::Config) -> Result<Id, Error> {
let producer = self.data_track(broadcast, name, |catalog, track| catalog.json_snapshot(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 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_mux::json::Config) -> Result<Id, Error> {
let producer = self.data_track(broadcast, name, |catalog, track| catalog.json_stream(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 producer = self.json_stream.remove(stream).ok_or(Error::TrackNotFound)?;
producer.finish()?;
Ok(())
}
pub fn binary_snapshot(&mut self, broadcast: Id, name: &str, config: moq_mux::binary::Config) -> Result<Id, Error> {
let producer = self.data_track(broadcast, name, |catalog, track| catalog.binary_snapshot(track, config))?;
self.binary_snapshot.insert(producer)
}
pub fn binary_snapshot_update(&mut self, binary: Id, payload: &[u8]) -> Result<(), Error> {
let producer = self.binary_snapshot.get_mut(binary).ok_or(Error::TrackNotFound)?;
producer.update(bytes::Bytes::copy_from_slice(payload))?;
Ok(())
}
pub fn binary_snapshot_finish(&mut self, binary: Id) -> Result<(), Error> {
let producer = self.binary_snapshot.remove(binary).ok_or(Error::TrackNotFound)?;
producer.finish()?;
Ok(())
}
pub fn binary_stream(&mut self, broadcast: Id, name: &str, config: moq_mux::binary::Config) -> Result<Id, Error> {
let producer = self.data_track(broadcast, name, |catalog, track| catalog.binary_stream(track, config))?;
self.binary_stream.insert(producer)
}
pub fn binary_stream_append(&mut self, stream: Id, payload: &[u8]) -> Result<(), Error> {
let producer = self.binary_stream.get_mut(stream).ok_or(Error::TrackNotFound)?;
producer.append(bytes::Bytes::copy_from_slice(payload))?;
Ok(())
}
pub fn binary_stream_finish(&mut self, stream: Id) -> Result<(), Error> {
let producer = self.binary_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 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(())
}
}