use std::collections::HashMap;
use hang::catalog::{Video, VideoConfig};
fn same_picture(a: &VideoConfig, b: &VideoConfig) -> bool {
let mut a = a.clone();
let mut b = b.clone();
a.stalled = None;
b.stalled = None;
a == b
}
use crate::catalog::{self, Names, Published};
use crate::feed::Feed;
use crate::{Config, Error, active, rung};
pub(crate) struct Pipeline {
source: moq_net::broadcast::Consumer,
config: Config,
active: active::Producer,
name: String,
rendition: VideoConfig,
feed: Feed,
rungs: Vec<Published>,
names: Names,
serving: HashMap<String, tokio::sync::watch::Sender<bool>>,
}
impl Pipeline {
pub(crate) async fn new(
source: moq_net::broadcast::Consumer,
config: Config,
active: active::Producer,
name: String,
rendition: VideoConfig,
) -> Result<Self, Error> {
let feed = Feed::new(source.track(&name)?, rendition.clone(), config.feed_decoder());
let mut ladder = Self {
source,
config,
active,
name,
rendition,
feed,
rungs: Vec::new(),
names: Names::default(),
serving: HashMap::new(),
};
let (name, rendition) = (ladder.name.clone(), ladder.rendition.clone());
let rungs = ladder.resolve(&name, &rendition, &[]).await?;
tracing::info!(source = %ladder.name, rungs = rungs.len(), "transcoding");
ladder.active.declare(rungs.iter().map(|published| &published.rung));
ladder.rungs = rungs;
Ok(ladder)
}
async fn resolve(
&mut self,
source_name: &str,
source: &VideoConfig,
carry: &[Published],
) -> Result<Vec<Published>, Error> {
let resolved = catalog::resolve_rungs(&self.config.ladder, source_name, source)?;
let mut published = Vec::new();
for mut rung in resolved {
let reused = carry.iter().find(|other| other.rung.same_shape(&rung)).cloned();
let entry = match reused {
Some(reused) => {
rung.name = reused.rung.name;
let mut entry = reused.entry;
entry.optimize_for_latency = source.optimize_for_latency;
entry
}
None => {
rung.name = self.names.mint(rung.height);
catalog::rung_entry(&rung, source, &self.config.encoder).await?
}
};
published.push(Published { rung, entry });
}
catalog::inherit_stalled(&mut published, source);
Ok(published)
}
pub(crate) fn rungs(&self) -> &[Published] {
&self.rungs
}
pub(crate) fn rung(&mut self, name: &str) -> Result<Option<rung::Rung>, Error> {
let Some(published) = self.rungs.iter().find(|published| published.rung.name == name) else {
return Ok(None);
};
let (retired, retire) = rung::Retire::channel();
self.serving.insert(published.rung.name.clone(), retired);
Ok(Some(rung::Rung {
source: self.source.track(&self.name)?,
feed: self.feed.clone(),
broadcast: self.source.clone(),
config: self.rendition.clone(),
encoder: self.config.encoder.clone(),
decoder: self.config.decoder.clone(),
resize: self.config.resize,
active: self.active.clone(),
info: published.rung.clone(),
retire,
}))
}
pub(crate) async fn follow(&mut self, video: &Video) -> Result<(), Error> {
let (name, rendition) = match catalog::follow_source(video, &self.name) {
Ok(chosen) => chosen,
Err(err) => {
tracing::debug!(%err, "no transcodable rendition in the catalog update");
return Ok(());
}
};
if name == self.name && same_picture(&rendition, &self.rendition) {
catalog::inherit_stalled(&mut self.rungs, &rendition);
self.rendition = rendition;
return Ok(());
}
let rebuilt = name != self.name || !catalog::same_stream(&self.rendition, &rendition);
let carry = match rebuilt {
true => Vec::new(),
false => self.rungs.clone(),
};
let rungs = match self.resolve(&name, &rendition, &carry).await {
Ok(rungs) => rungs,
Err(err) => {
tracing::warn!(%err, source = %name, "could not resolve a ladder for the new source");
return Ok(());
}
};
if rebuilt {
for retired in self.serving.values() {
let _ = retired.send(true);
}
self.serving.clear();
self.feed = Feed::new(self.source.track(&name)?, rendition.clone(), self.config.feed_decoder());
}
self.name = name;
self.rendition = rendition;
for published in &self.rungs {
if rungs.iter().any(|other| other.rung == published.rung) {
continue;
}
if let Some(retired) = self.serving.remove(&published.rung.name) {
let _ = retired.send(true);
}
}
tracing::info!(source = %self.name, rungs = rungs.len(), "source changed; ladder resolved again");
self.active.declare(rungs.iter().map(|published| &published.rung));
self.rungs = rungs;
Ok(())
}
}