use std::sync::Arc;
use axum::Router;
use axum::extract::State;
use axum::http::{StatusCode, header};
use axum::response::{IntoResponse, Response};
use axum::routing::get;
use broadcast_common::Package;
use hls_runtime::server::DEFAULT_TRACK_ID;
use media_plane::egress::{AwaitPolicy, CachePolicy, EgressResponse, ServedEgress};
use transmux::{Addressing, LlDashPackager, Media, Track, TrackSegments};
use crate::http::{self, BLOCKING_RELOAD_TIMEOUT};
use crate::origin::resource::cors_preflight;
use crate::output::dash::{DASH_MANIFEST_CONTENT_TYPE, format_iso8601, select_representable_track};
use crate::output::{Output, OutputKind};
use crate::route::RouteHandle;
pub const LL_DASH_MANIFEST_NAME: &str = "manifest-ll.mpd";
const LATENCY_TARGET_PART_MULTIPLE: u64 = 3;
pub struct LlDashOutput;
impl Output for LlDashOutput {
fn kind(&self) -> OutputKind {
OutputKind::LlDash
}
fn manifest_routes(&self, route: Arc<RouteHandle>) -> Router {
Router::new()
.route(
&format!("/{LL_DASH_MANIFEST_NAME}"),
get(manifest).options(cors_preflight),
)
.with_state(route)
}
}
async fn manifest(State(route): State<Arc<RouteHandle>>) -> Response {
let serving = match http::resolve_route_program(&route) {
Ok(serving) => serving,
Err(resp) => return *resp,
};
let trunk = serving.trunk();
let origin = LlDashOrigin { route };
let resp = http::resolve_blocking(&trunk, &origin, (), BLOCKING_RELOAD_TIMEOUT, || ()).await;
http::into_response(resp, StatusCode::SERVICE_UNAVAILABLE, |body| {
([(header::CONTENT_TYPE, DASH_MANIFEST_CONTENT_TYPE)], body).into_response()
})
}
struct LlDashOrigin {
route: Arc<RouteHandle>,
}
impl ServedEgress for LlDashOrigin {
type Request = ();
type Body = String;
fn resolve(
&self,
_request: (),
_now: broadcast_common::Timestamp,
_await_policy: AwaitPolicy,
) -> EgressResponse<String> {
match render_ll_dash_mpd(&self.route) {
Some(body) => EgressResponse::Ready {
body,
cache: CachePolicy::NoCache,
},
None => EgressResponse::NotFound,
}
}
}
fn render_ll_dash_mpd(route: &RouteHandle) -> Option<String> {
let specs = route.track_specs(crate::route::SPTS_PROGRAM_ID);
let mut spec = select_representable_track(&specs)?;
spec.track_id = DEFAULT_TRACK_ID;
let timescale = spec.timescale.max(1);
let window = route.window_segments(crate::route::SPTS_PROGRAM_ID);
let start_number = window
.first()
.map(|s| u64::from(s.segment_seq))
.unwrap_or(1);
let target_duration_secs = route.target_duration_secs();
let nominal_duration_ticks =
((target_duration_secs * f64::from(timescale)).round() as u64).max(1);
let part_target_ms = route.part_target_ms().max(1);
let chunk_duration_secs = f64::from(part_target_ms) / 1000.0;
let latency_target_ms =
u32::try_from(u64::from(part_target_ms) * LATENCY_TARGET_PART_MULTIPLE).unwrap_or(u32::MAX);
let media = Media::new(vec![Track::new(spec, Vec::new())], timescale);
let mut packager = LlDashPackager::new(
target_duration_secs,
chunk_duration_secs,
latency_target_ms,
format_iso8601(route.created_at()),
)
.ok()?;
packager.base.addressing = Addressing::Number;
packager.base.start_number = start_number;
packager.base.init_template = "init-$RepresentationID$.mp4".to_string();
packager.base.media_template = "seg-$RepresentationID$-$Number$.m4s".to_string();
packager.base.minimum_update_period = Some(format!("PT{chunk_duration_secs}S"));
let time_shift_buffer_depth_secs = target_duration_secs * (window.len().max(1) as f64);
packager.base.time_shift_buffer_depth = Some(format!("PT{time_shift_buffer_depth_secs}S"));
packager.base.segments = vec![TrackSegments {
track_id: DEFAULT_TRACK_ID,
durations: vec![nominal_duration_ticks],
}];
packager.package(&media).ok()
}
#[cfg(test)]
mod tests {
use super::*;
use transmux::CodecConfig;
use transmux::TrackSpec;
use transmux::ll_hls::PartInfo;
fn video_spec(track_id: u32) -> TrackSpec {
TrackSpec::new(
track_id,
90_000,
CodecConfig::Vp8 {
width: 1280,
height: 720,
},
)
}
fn part(seq: u32, idx: u32) -> PartInfo {
PartInfo {
bytes: vec![0x40 + idx as u8; 4],
duration: 0.5,
independent: idx == 0,
segment_seq: seq,
part_index: idx,
}
}
fn seg(seq: u32, duration: f64) -> transmux::ll_hls::SegmentInfo {
transmux::ll_hls::SegmentInfo {
bytes: vec![seq as u8; 8],
duration,
segment_seq: seq,
part_count: 2,
}
}
#[test]
fn render_ll_dash_mpd_none_without_track_specs() {
let route = RouteHandle::new(4.0, 500, 4);
assert!(
render_ll_dash_mpd(&route).is_none(),
"no track specs recorded yet -> nothing to describe"
);
}
#[test]
fn render_ll_dash_mpd_valid_before_any_segment_closes() {
let route = RouteHandle::new(4.0, 500, 4);
route.publish_new_program(crate::route::SPTS_PROGRAM_ID);
route.set_track_specs(crate::route::SPTS_PROGRAM_ID, vec![video_spec(7)]);
let mpd = render_ll_dash_mpd(&route).expect("must render even with an empty window");
assert!(mpd.contains("<MPD"));
assert!(mpd.contains("type=\"dynamic\""));
}
#[test]
fn render_ll_dash_mpd_carries_required_ll_dash_elements() {
let route = RouteHandle::new(4.0, 500, 4);
route.publish_new_program(crate::route::SPTS_PROGRAM_ID);
route.set_track_specs(crate::route::SPTS_PROGRAM_ID, vec![video_spec(1)]);
route.add_part(crate::route::SPTS_PROGRAM_ID, part(1, 0));
let mpd = render_ll_dash_mpd(&route).unwrap();
assert!(
mpd.contains("availabilityTimeComplete=\"false\""),
"availabilityTimeComplete must be present and false: {mpd}"
);
assert!(
mpd.contains("availabilityTimeOffset=\"3.5\""),
"availabilityTimeOffset must reflect segment-chunk duration: {mpd}"
);
assert!(
mpd.contains("<ServiceDescription"),
"ServiceDescription must be present: {mpd}"
);
assert!(mpd.contains("<Latency target="), "{mpd}");
assert!(
mpd.contains("PT0.5S"),
"minimumUpdatePeriod must be tuned to the part/chunk target: {mpd}"
);
}
#[test]
fn render_ll_dash_mpd_addresses_whole_segments_not_parts() {
let route = RouteHandle::new(4.0, 500, 4);
route.publish_new_program(crate::route::SPTS_PROGRAM_ID);
route.set_track_specs(crate::route::SPTS_PROGRAM_ID, vec![video_spec(1)]);
route.add_part(crate::route::SPTS_PROGRAM_ID, part(5, 0));
let mpd = render_ll_dash_mpd(&route).unwrap();
assert!(
mpd.contains("seg-$RepresentationID$-$Number$.m4s"),
"media template must address whole segments (the true chunked-transfer \
design serves parts internally, never in the MPD itself): {mpd}"
);
assert!(
!mpd.contains("part-"),
"no part-addressed URI should ever appear in the MPD: {mpd}"
);
}
#[test]
fn render_ll_dash_mpd_start_number_tracks_window() {
let route = RouteHandle::new(4.0, 500, 2);
route.publish_new_program(crate::route::SPTS_PROGRAM_ID);
route.set_track_specs(crate::route::SPTS_PROGRAM_ID, vec![video_spec(1)]);
route.add_segment(crate::route::SPTS_PROGRAM_ID, seg(1, 4.0));
route.add_segment(crate::route::SPTS_PROGRAM_ID, seg(2, 4.0));
route.add_segment(crate::route::SPTS_PROGRAM_ID, seg(3, 4.0));
let mpd = render_ll_dash_mpd(&route).unwrap();
assert!(
mpd.contains("startNumber=\"2\""),
"startNumber must track the window's oldest retained segment_seq (2): {mpd}"
);
}
#[test]
fn render_ll_dash_mpd_carries_time_shift_buffer_depth() {
let route = RouteHandle::new(2.0, 500, 4);
route.publish_new_program(crate::route::SPTS_PROGRAM_ID);
route.set_track_specs(crate::route::SPTS_PROGRAM_ID, vec![video_spec(1)]);
route.add_segment(crate::route::SPTS_PROGRAM_ID, seg(1, 2.0));
let mpd = render_ll_dash_mpd(&route).unwrap();
assert!(mpd.contains("timeShiftBufferDepth=\"PT2S\""), "{mpd}");
}
#[test]
fn render_ll_dash_mpd_forces_representation_id_to_default_track() {
let route = RouteHandle::new(4.0, 500, 4);
route.publish_new_program(crate::route::SPTS_PROGRAM_ID);
route.set_track_specs(crate::route::SPTS_PROGRAM_ID, vec![video_spec(7)]);
route.add_part(crate::route::SPTS_PROGRAM_ID, part(1, 0));
let mpd = render_ll_dash_mpd(&route).unwrap();
assert!(
mpd.contains(&format!("id=\"{DEFAULT_TRACK_ID}\"")),
"Representation @id must be the DEFAULT_TRACK_ID, not the source's own \
track_id (7): {mpd}"
);
assert!(!mpd.contains("id=\"7\""), "source track_id leaked: {mpd}");
}
#[tokio::test]
async fn manifest_handler_503_before_track_specs_known() {
let route = Arc::new(RouteHandle::new(4.0, 500, 4));
route.publish_new_program(crate::route::SPTS_PROGRAM_ID);
let resp = manifest(State(route)).await;
assert_eq!(resp.status(), StatusCode::SERVICE_UNAVAILABLE);
}
#[tokio::test]
async fn manifest_handler_200_with_dash_content_type() {
let route = Arc::new(RouteHandle::new(4.0, 500, 4));
route.publish_new_program(crate::route::SPTS_PROGRAM_ID);
route.set_track_specs(crate::route::SPTS_PROGRAM_ID, vec![video_spec(1)]);
route.add_part(crate::route::SPTS_PROGRAM_ID, part(1, 0));
let resp = manifest(State(route)).await;
assert_eq!(resp.status(), StatusCode::OK);
assert_eq!(
resp.headers().get(header::CONTENT_TYPE).unwrap(),
DASH_MANIFEST_CONTENT_TYPE
);
}
#[tokio::test]
async fn manifest_not_yet_announced_is_503_not_404() {
let route = Arc::new(RouteHandle::new(4.0, 500, 4));
let resp = manifest(State(route)).await;
assert_eq!(
resp.status(),
StatusCode::SERVICE_UNAVAILABLE,
"a route with no program announced yet must be 503 (not ready), not 404 (gone)"
);
}
}