use std::sync::Arc;
use bytes::Bytes;
use hang::catalog::VideoConfig;
use moq_mux::container::Container as _;
use tokio::sync::Semaphore;
use crate::Error;
use crate::catalog::Resolved;
use crate::feed::{Feed, Item};
const MAX_CONCURRENT_FETCHES: usize = 4;
#[derive(Clone)]
pub(crate) struct Rung {
pub info: Resolved,
pub source: moq_net::track::Consumer,
pub feed: Feed,
pub broadcast: moq_net::broadcast::Consumer,
pub config: VideoConfig,
pub encoder: moq_video::encode::Kind,
pub decoder: moq_video::decode::Kind,
}
impl Rung {
fn pipeline(&self) -> Result<Pipeline, Error> {
Pipeline::new(self)
}
fn container(&self) -> Result<moq_mux::catalog::hang::Container, Error> {
Ok(moq_mux::catalog::hang::Container::try_from(&self.config.container)?)
}
fn encode(&self) -> Result<moq_video::encode::Encoder, Error> {
let mut config =
moq_video::encode::Config::new(self.info.size.width, self.info.size.height, self.info.framerate);
config.bitrate = Some(self.info.bitrate);
config.kind = self.encoder.clone();
config.gop = self.info.framerate.saturating_mul(8).max(1);
Ok(moq_video::encode::Encoder::new(&config)?)
}
}
pub(crate) async fn serve(rung: Rung, request: moq_net::track::Request) -> Result<(), Error> {
let dynamic = request.dynamic();
let info = hang::container::track_info();
let mut producer = request.accept(info);
let result = tokio::select! {
res = live(&rung, &mut producer) => res,
res = fetches(&rung, &dynamic) => res,
};
if result.is_err() {
let _ = producer.abort(moq_net::Error::Cancel);
}
result
}
async fn live(rung: &Rung, producer: &mut moq_net::track::Producer) -> Result<(), Error> {
let demand = producer.demand();
loop {
tokio::select! {
used = demand.used() => if used.is_err() {
return Ok(());
},
err = rung.broadcast.closed() => {
producer.clone().abort(err)?;
return Ok(());
}
}
let mut listener = rung.feed.listen();
let mut encoder = rung.encode()?;
let mut current: Option<moq_net::group::Producer> = None;
let mut first = true;
'session: loop {
let item = tokio::select! {
item = listener.recv() => item,
_ = demand.unused() => {
if let Some(output) = current.take() {
output.abort(moq_net::Error::Cancel)?;
}
break 'session;
}
};
match item {
Some(Item::Group(sequence)) => {
if let Some(output) = current.take() {
output.abort(moq_net::Error::Cancel)?;
}
first = true;
let info = moq_net::group::Info { sequence };
current = match producer.create_group(info) {
Ok(output) => Some(output),
Err(moq_net::Error::Duplicate) => None,
Err(err) => return Err(err.into()),
};
}
Some(Item::Frame(frame)) => {
let Some(output) = &mut current else { continue };
let keyframe = first;
first = false;
let encoded = if frame.size == rung.info.size {
encoder.encode(&frame, keyframe)?
} else {
let scaled = frame.resize(rung.info.size)?;
encoder.encode(&scaled, keyframe)?
};
let timestamp = frame.timestamp;
write(output, encoded.into_iter().map(|packet| (timestamp, packet)).collect())?;
}
Some(Item::End) => {
if let Some(mut output) = current.take() {
output.finish()?;
}
}
Some(Item::Lagged) => {
if let Some(output) = current.take() {
output.abort(moq_net::Error::Cancel)?;
}
}
Some(Item::Finished) => {
if let Some(output) = current.take() {
output.abort(moq_net::Error::Cancel)?;
}
producer.finish()?;
return Ok(());
}
None => {
if let Some(output) = current.take() {
let _ = output.abort(moq_net::Error::Cancel);
}
producer.clone().abort(moq_net::Error::Cancel)?;
return Ok(());
}
}
}
}
}
async fn fetches(rung: &Rung, dynamic: &moq_net::track::Dynamic) -> Result<(), Error> {
let limit = Arc::new(Semaphore::new(MAX_CONCURRENT_FETCHES));
let mut tasks = tokio::task::JoinSet::new();
loop {
while tasks.try_join_next().is_some() {}
let Ok(request) = dynamic.requested_group().await else {
return Ok(());
};
let Ok(permit) = limit.clone().acquire_owned().await else {
return Ok(());
};
let rung = rung.clone();
tasks.spawn(async move {
let _permit = permit;
let sequence = request.sequence();
if let Err(err) = fetch(rung, request).await {
tracing::warn!(%err, sequence, "transcode fetch failed");
}
});
}
}
async fn fetch(rung: Rung, request: moq_net::track::GroupRequest) -> Result<(), Error> {
let options = moq_net::group::Fetch::default().with_priority(request.priority());
let mut source = match rung.source.fetch_group(request.sequence(), options).await {
Ok(source) => source,
Err(err) => {
request.reject(err.clone());
return Err(err.into());
}
};
let (pipeline, container) = match rung.pipeline().and_then(|p| rung.container().map(|c| (p, c))) {
Ok(built) => built,
Err(err) => {
request.reject(moq_net::Error::Cancel);
return Err(err);
}
};
let output = match request.accept(None) {
Ok(output) => output,
Err(err) => return Err(err.into()),
};
transcode_group(pipeline, &container, &mut source, output).await?;
Ok(())
}
async fn transcode_group(
pipeline: Pipeline,
container: &moq_mux::catalog::hang::Container,
source: &mut moq_net::group::Consumer,
mut output: moq_net::group::Producer,
) -> Result<(), Error> {
match transcode_group_inner(pipeline, container, source, &mut output).await {
Ok(()) => {
output.finish()?;
Ok(())
}
Err(err) => {
let _ = output.abort(moq_net::Error::Cancel);
Err(err)
}
}
}
async fn transcode_group_inner(
mut pipeline: Pipeline,
container: &moq_mux::catalog::hang::Container,
source: &mut moq_net::group::Consumer,
output: &mut moq_net::group::Producer,
) -> Result<(), Error> {
let mut first = true;
let mut last_timestamp: Option<moq_net::Timestamp> = None;
while let Some(frames) = container.read(source).await? {
for frame in frames {
let timestamp = frame.timestamp;
last_timestamp = Some(last_timestamp.map_or(timestamp, |last| last.max(timestamp)));
let keyframe = frame.keyframe || first;
first = false;
write(output, pipeline.process(&frame.payload, timestamp, keyframe)?)?;
}
}
if let Some(last_timestamp) = last_timestamp {
write(output, pipeline.finish(last_timestamp)?)?;
}
Ok(())
}
fn write(output: &mut moq_net::group::Producer, packets: Vec<(moq_net::Timestamp, Bytes)>) -> Result<(), Error> {
for (timestamp, payload) in packets {
let frame = hang::container::Frame { timestamp, payload };
frame.write_to(output)?;
}
Ok(())
}
struct Pipeline {
decoder: moq_video::decode::Decoder,
encoder: moq_video::encode::Encoder,
size: moq_video::Size,
}
impl Pipeline {
fn new(rung: &Rung) -> Result<Self, Error> {
let mut decode = moq_video::decode::Config::new();
decode.kind = rung.decoder.clone();
decode.resize = Some(rung.info.size);
let decoder = moq_video::decode::Decoder::new(&rung.config, &decode)?;
Ok(Self {
decoder,
encoder: rung.encode()?,
size: rung.info.size,
})
}
fn process(
&mut self,
payload: &Bytes,
timestamp: moq_net::Timestamp,
keyframe: bool,
) -> Result<Vec<(moq_net::Timestamp, Bytes)>, Error> {
let mut packets = Vec::new();
for raw in self.decoder.decode(payload, timestamp, keyframe)? {
let raw_timestamp = raw.timestamp;
let encoded = if raw.size == self.size {
self.encoder.encode(&raw, keyframe)?
} else {
self.encoder.encode(&raw.resize(self.size)?, keyframe)?
};
for packet in encoded {
packets.push((raw_timestamp, packet));
}
}
Ok(packets)
}
fn finish(self, timestamp: moq_net::Timestamp) -> Result<Vec<(moq_net::Timestamp, Bytes)>, Error> {
Ok(self
.encoder
.finish()?
.into_iter()
.map(|packet| (timestamp, packet))
.collect())
}
}