use transmux::hls::{LowLatencyConfig, MediaPlaylist, MediaSegment, OpenSegment, PartSpec};
use super::store::MediaStore;
pub const DEFAULT_TRACK_ID: u32 = 1;
const PLACEHOLDER_BANDWIDTH_BPS: u64 = 5_000_000;
const ABUSE_MSN_FUTURE_BOUND: u64 = 4;
const LL_HLS_VERSION: u8 = 9;
const PART_HOLD_BACK_MULTIPLIER: f64 = 3.0;
pub fn media_playlist_m3u8(store: &MediaStore, track_id: u32) -> String {
let target_duration_secs = store.target_duration_secs();
let max_segment_duration = store.max_segment_duration();
store.with_segments_and_parts(|store_segments, live_parts| {
let media_sequence = store_segments
.front()
.map(|s| u64::from(s.segment_seq))
.or_else(|| live_parts.first().map(|p| u64::from(p.segment_seq)))
.unwrap_or(1);
let segments: Vec<MediaSegment> = store_segments
.iter()
.map(|s| MediaSegment {
uri: format!("seg-{track_id}-{}.m4s", s.segment_seq),
duration: s.duration,
discontinuous: false,
parts: Vec::new(),
..Default::default()
})
.collect();
let part_target = f64::from(store.part_target_ms()) / 1000.0;
let open_seq = live_parts.first().map(|p| p.segment_seq);
let open_segment = open_seq.map(|seq| {
OpenSegment::new(
live_parts
.iter()
.filter(|p| p.segment_seq == seq)
.map(|p| PartSpec {
uri: format!("part-{track_id}-{}.{}.m4s", p.segment_seq, p.part_index),
duration: p.duration,
independent: p.independent,
..Default::default()
})
.collect(),
)
});
let next_part_hint = open_seq.map(|seq| {
let next_idx = live_parts
.iter()
.filter(|p| p.segment_seq == seq)
.map(|p| p.part_index)
.max()
.map(|idx| idx + 1)
.unwrap_or(0);
format!("part-{track_id}-{seq}.{next_idx}.m4s")
});
let target_duration = target_duration_secs.max(max_segment_duration).round() as u32;
let playlist = MediaPlaylist {
version: LL_HLS_VERSION,
target_duration,
media_sequence,
discontinuity_sequence: 0,
segments,
open_segment,
endlist: false,
extra_tags: vec![format!("#EXT-X-MAP:URI=\"init-{track_id}.mp4\"")],
low_latency: Some(LowLatencyConfig {
part_target,
part_hold_back: part_target * PART_HOLD_BACK_MULTIPLIER,
preload_hint_part: next_part_hint,
..Default::default()
}),
iframes_only: false,
..Default::default()
};
playlist.to_m3u8()
})
}
pub fn master_playlist_m3u8(media_playlist_name: &str) -> String {
format!(
"#EXTM3U\n#EXT-X-STREAM-INF:BANDWIDTH={PLACEHOLDER_BANDWIDTH_BPS}\n{media_playlist_name}\n"
)
}
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
pub struct BlockingQuery {
pub hls_msn: Option<u64>,
pub hls_part: Option<u32>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum PlaylistOutcome {
Ready(String),
WouldBlock,
BadRequest,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum CachePolicy {
Immutable,
NoCache,
}
impl CachePolicy {
pub fn name(&self) -> &'static str {
match self {
CachePolicy::Immutable => "immutable",
CachePolicy::NoCache => "no-cache",
}
}
}
broadcast_common::impl_spec_display!(CachePolicy);
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum ResourceOutcome {
Ready {
bytes: Vec<u8>,
cache: CachePolicy,
},
WouldBlock,
NotFound,
}
impl MediaStore {
pub fn resolve_playlist(&self, track_id: u32, query: BlockingQuery) -> PlaylistOutcome {
if query.hls_part.is_some() && query.hls_msn.is_none() {
return PlaylistOutcome::BadRequest;
}
if let Some(msn) = query.hls_msn {
let (current_max_msn, _) = self.latest_progress();
if msn > u64::from(current_max_msn) + ABUSE_MSN_FUTURE_BOUND {
return PlaylistOutcome::BadRequest;
}
let satisfied = match query.hls_part {
Some(part) => {
let (in_progress_seg_seq, part_count) = self.latest_progress();
u64::from(in_progress_seg_seq) > msn
|| (u64::from(in_progress_seg_seq) == msn && part_count > part)
}
None => u64::from(self.last_closed_segment_seq()) >= msn,
};
if !satisfied {
return PlaylistOutcome::WouldBlock;
}
}
PlaylistOutcome::Ready(media_playlist_m3u8(self, track_id))
}
pub fn resolve_resource(&self, name: &str) -> ResourceOutcome {
if let Some((seq, idx)) = parse_part(name) {
return match self.part_bytes(seq, idx) {
Some(bytes) => ResourceOutcome::Ready {
bytes,
cache: CachePolicy::Immutable,
},
None => {
let (in_progress_seg_seq, _) = self.latest_progress();
if in_progress_seg_seq > seq || self.segment_bytes(seq).is_some() {
ResourceOutcome::NotFound
} else {
ResourceOutcome::WouldBlock
}
}
};
}
match resolve_file(self, name) {
Some(bytes) => ResourceOutcome::Ready {
bytes,
cache: CachePolicy::Immutable,
},
None => ResourceOutcome::NotFound,
}
}
}
fn parse_part(file: &str) -> Option<(u32, u32)> {
let rest = file.strip_prefix("part-")?.strip_suffix(".m4s")?;
let (track_seq, idx) = rest.rsplit_once('.')?;
let (track, seq) = track_seq.split_once('-')?;
track.parse::<u32>().ok()?;
Some((seq.parse().ok()?, idx.parse().ok()?))
}
fn resolve_file(store: &MediaStore, file: &str) -> Option<Vec<u8>> {
if let Some(rest) = file.strip_prefix("init-") {
let track = rest.strip_suffix(".mp4")?;
track.parse::<u32>().ok()?;
return store.init_bytes();
}
if let Some(rest) = file.strip_prefix("seg-") {
let rest = rest.strip_suffix(".m4s")?;
let (track, seq) = rest.split_once('-')?;
track.parse::<u32>().ok()?;
let seq: u32 = seq.parse().ok()?;
return store.segment_bytes(seq);
}
None
}
#[cfg(test)]
mod tests {
use super::*;
use transmux::ll_hls::{PartInfo, SegmentInfo};
fn part(seq: u32, idx: u32) -> PartInfo {
PartInfo {
bytes: vec![0x10 + idx as u8; 4],
duration: 0.5,
independent: idx == 0,
segment_seq: seq,
part_index: idx,
}
}
fn seg(seq: u32) -> SegmentInfo {
SegmentInfo {
bytes: vec![0x20 + seq as u8; 8],
duration: 4.0,
segment_seq: seq,
part_count: 2,
}
}
fn make_store() -> MediaStore {
let store = MediaStore::new(4.0, 500, 4);
store.set_init(vec![0xAA; 8]);
store.add_segment(seg(1));
store.add_part(part(2, 0));
store.add_part(part(2, 1));
store
}
#[test]
fn cache_policy_name_and_display_agree() {
for (policy, label) in [
(CachePolicy::Immutable, "immutable"),
(CachePolicy::NoCache, "no-cache"),
] {
assert_eq!(policy.name(), label);
assert_eq!(policy.to_string(), label);
}
}
#[test]
fn master_playlist_has_stream_inf() {
let m = master_playlist_m3u8("media.m3u8");
assert!(m.contains("#EXTM3U"));
assert!(m.contains("#EXT-X-STREAM-INF"));
assert!(m.contains("media.m3u8"));
}
#[test]
fn master_playlist_points_at_configured_playlist_name() {
let m = master_playlist_m3u8("index.m3u8");
assert!(m.contains("index.m3u8"));
assert!(!m.contains("media.m3u8"));
}
#[test]
fn resolve_playlist_no_query_is_ready_now() {
let store = make_store();
let outcome = store.resolve_playlist(DEFAULT_TRACK_ID, BlockingQuery::default());
match outcome {
PlaylistOutcome::Ready(body) => assert!(body.contains("#EXT-X-PART"), "body: {body}"),
other => panic!("expected Ready, got {other:?}"),
}
}
#[test]
fn resolve_playlist_already_satisfied_earlier_msn_is_ready() {
let store = make_store();
let outcome = store.resolve_playlist(
DEFAULT_TRACK_ID,
BlockingQuery {
hls_msn: Some(1),
hls_part: Some(0),
},
);
assert!(matches!(outcome, PlaylistOutcome::Ready(_)));
}
#[test]
fn resolve_playlist_already_satisfied_same_msn_lower_part_is_ready() {
let store = make_store();
let outcome = store.resolve_playlist(
DEFAULT_TRACK_ID,
BlockingQuery {
hls_msn: Some(2),
hls_part: Some(1),
},
);
assert!(matches!(outcome, PlaylistOutcome::Ready(_)));
}
#[test]
fn resolve_playlist_msn_only_waits_for_closed_segment_not_just_open_parts() {
let store = make_store();
let outcome = store.resolve_playlist(
DEFAULT_TRACK_ID,
BlockingQuery {
hls_msn: Some(2),
hls_part: None,
},
);
assert_eq!(outcome, PlaylistOutcome::WouldBlock);
store.add_segment(seg(2));
let outcome = store.resolve_playlist(
DEFAULT_TRACK_ID,
BlockingQuery {
hls_msn: Some(2),
hls_part: None,
},
);
match outcome {
PlaylistOutcome::Ready(body) => assert!(
body.contains("seg-1-2.m4s"),
"resolved playlist must show segment 2 as closed: {body}"
),
other => panic!("expected Ready after close, got {other:?}"),
}
}
#[test]
fn resolve_playlist_msn_within_bound_would_block_until_part_lands() {
let store = make_store(); let outcome = store.resolve_playlist(
DEFAULT_TRACK_ID,
BlockingQuery {
hls_msn: Some(2),
hls_part: Some(2),
},
);
assert_eq!(outcome, PlaylistOutcome::WouldBlock);
store.add_part(part(2, 2));
let outcome = store.resolve_playlist(
DEFAULT_TRACK_ID,
BlockingQuery {
hls_msn: Some(2),
hls_part: Some(2),
},
);
assert!(matches!(outcome, PlaylistOutcome::Ready(_)));
}
#[test]
fn resolve_playlist_far_future_msn_rejected() {
let store = make_store();
let outcome = store.resolve_playlist(
DEFAULT_TRACK_ID,
BlockingQuery {
hls_msn: Some(1002),
hls_part: None,
},
);
assert_eq!(outcome, PlaylistOutcome::BadRequest);
}
#[test]
fn resolve_playlist_part_without_msn_rejected() {
let store = make_store();
let outcome = store.resolve_playlist(
DEFAULT_TRACK_ID,
BlockingQuery {
hls_msn: None,
hls_part: Some(0),
},
);
assert_eq!(outcome, PlaylistOutcome::BadRequest);
}
#[test]
fn resolve_resource_init_present() {
let store = make_store();
let outcome = store.resolve_resource("init-1.mp4");
match outcome {
ResourceOutcome::Ready { bytes, cache } => {
assert_eq!(bytes, vec![0xAA; 8]);
assert_eq!(cache, CachePolicy::Immutable);
}
other => panic!("expected Ready, got {other:?}"),
}
}
#[test]
fn resolve_resource_segment_present_and_absent() {
let store = make_store();
match store.resolve_resource("seg-1-1.m4s") {
ResourceOutcome::Ready { bytes, .. } => assert_eq!(bytes, vec![0x21; 8]),
other => panic!("expected Ready, got {other:?}"),
}
assert_eq!(
store.resolve_resource("seg-1-99.m4s"),
ResourceOutcome::NotFound
);
}
#[test]
fn resolve_resource_part_present() {
let store = make_store();
match store.resolve_resource("part-1-2.0.m4s") {
ResourceOutcome::Ready { bytes, .. } => assert_eq!(bytes, vec![0x10; 4]),
other => panic!("expected Ready, got {other:?}"),
}
}
#[test]
fn resolve_resource_part_not_yet_produced_would_block() {
let store = make_store();
assert_eq!(
store.resolve_resource("part-1-2.2.m4s"),
ResourceOutcome::WouldBlock
);
store.add_part(part(2, 2));
match store.resolve_resource("part-1-2.2.m4s") {
ResourceOutcome::Ready { bytes, .. } => assert_eq!(bytes, vec![0x12; 4]),
other => panic!("expected Ready once produced, got {other:?}"),
}
}
#[test]
fn resolve_resource_part_not_found_once_segment_closes_without_it() {
let store = make_store();
assert_eq!(
store.resolve_resource("part-1-2.9.m4s"),
ResourceOutcome::WouldBlock,
"not yet decidable while segment 2 is still open"
);
store.add_segment(seg(2));
assert_eq!(
store.resolve_resource("part-1-2.9.m4s"),
ResourceOutcome::NotFound,
"must resolve NotFound once segment 2 has closed without producing it"
);
}
#[test]
fn resolve_resource_part_served_from_recent_after_close() {
let store = make_store();
store.add_segment(seg(2)); match store.resolve_resource("part-1-2.1.m4s") {
ResourceOutcome::Ready { bytes, .. } => assert_eq!(bytes, vec![0x11; 4]),
other => panic!("a just-closed segment's part must still resolve Ready, got {other:?}"),
}
}
#[test]
fn resolve_resource_part_of_old_segment_not_found() {
let store = make_store();
assert_eq!(
store.resolve_resource("part-1-1.0.m4s"),
ResourceOutcome::NotFound
);
}
#[test]
fn resolve_resource_unmatched_filename_not_found() {
let store = make_store();
assert_eq!(
store.resolve_resource("not-a-thing.txt"),
ResourceOutcome::NotFound
);
}
fn plain_seg(seq: u32, parts: u32) -> SegmentInfo {
SegmentInfo {
bytes: vec![seq as u8; 8],
duration: 4.0,
segment_seq: seq,
part_count: parts,
}
}
fn plain_part(seq: u32, idx: u32) -> PartInfo {
PartInfo {
bytes: vec![idx as u8; 4],
duration: 0.5,
independent: idx == 0,
segment_seq: seq,
part_index: idx,
}
}
#[test]
fn playlist_has_llhls_tags_and_parts() {
let s = MediaStore::new(4.0, 500, 4);
s.set_init(vec![0; 4]);
s.add_part(plain_part(1, 0));
s.add_part(plain_part(1, 1));
let m = media_playlist_m3u8(&s, 1);
assert!(m.contains("#EXT-X-PART-INF"), "PART-INF present");
assert!(
m.contains("#EXT-X-SERVER-CONTROL"),
"SERVER-CONTROL present"
);
assert!(m.contains("#EXT-X-PART"), "at least one PART");
assert!(
m.contains("part-1-1.0.m4s") || m.contains("part-1-1.1.m4s"),
"part URI"
);
}
#[test]
fn open_segment_has_parts_but_no_extinf() {
let s = MediaStore::new(4.0, 500, 4);
s.set_init(vec![0; 4]);
s.add_part(plain_part(1, 0));
s.add_part(plain_part(1, 1));
let m = media_playlist_m3u8(&s, 1);
assert!(m.contains("#EXT-X-PART"), "at least one PART line");
assert!(m.contains("part-1-1.0.m4s"), "part 0 URI present");
assert!(m.contains("part-1-1.1.m4s"), "part 1 URI present");
assert!(
!m.contains("seg-1-1.m4s"),
"no full-segment URI for the open segment: {m}"
);
assert!(
!m.contains("#EXTINF"),
"no EXTINF for the open segment: {m}"
);
}
#[test]
fn final_part_fetchable_after_its_segment_closes() {
let s = MediaStore::new(4.0, 500, 4);
s.set_init(vec![0; 4]);
s.add_part(plain_part(1, 0));
s.add_part(plain_part(1, 1)); s.add_segment(plain_seg(1, 2)); assert_eq!(
s.resolve_resource("part-1-1.1.m4s"),
ResourceOutcome::Ready {
bytes: vec![1; 4],
cache: CachePolicy::Immutable
},
"final part of a just-closed segment must still be individually fetchable"
);
assert_eq!(
s.resolve_resource("part-1-1.0.m4s"),
ResourceOutcome::Ready {
bytes: vec![0; 4],
cache: CachePolicy::Immutable
},
"earlier parts too"
);
assert_eq!(
s.resolve_resource("part-1-1.9.m4s"),
ResourceOutcome::NotFound
);
let m = media_playlist_m3u8(&s, 1);
assert!(
m.contains("seg-1-1.m4s"),
"closed segment rendered whole: {m}"
);
assert!(
!m.contains("part-1-1."),
"closed parts not rendered as open: {m}"
);
}
#[test]
fn live_parts_capped_when_segment_never_closes() {
let s = MediaStore::new(4.0, 500, 4);
let cap = super::super::store::compute_max_live_parts(4.0, 500);
assert_eq!(cap, 12, "sanity-check the expected cap for these params");
s.set_init(vec![0; 4]);
for i in 0..(cap as u32 * 5) {
s.add_part(plain_part(1, i));
}
assert_eq!(
s.live_part_count(),
cap,
"live_parts must stay capped even though the segment never closed"
);
let m = media_playlist_m3u8(&s, 1);
assert!(m.contains("#EXT-X-PART"), "still has PART lines: {m}");
let last_idx = cap as u32 * 5 - 1;
assert!(
m.contains(&format!("part-1-1.{last_idx}.m4s")),
"most recent part must survive the cap: {m}"
);
let first_idx = cap as u32 * 5 - cap as u32;
assert!(
!m.contains(&format!("part-1-1.{}.m4s", first_idx - 1)),
"an older part beyond the cap must have been dropped: {m}"
);
}
#[test]
fn target_duration_is_max_of_configured_and_actual_segment_duration() {
let s = MediaStore::new(4.0, 500, 4);
s.set_init(vec![0; 4]);
let mut long_seg = plain_seg(1, 2);
long_seg.duration = 7.5;
s.add_segment(long_seg);
let m = media_playlist_m3u8(&s, 1);
assert!(
m.contains("#EXT-X-TARGETDURATION:8"),
"TARGETDURATION must be round(7.5)=8, not the configured target (4): {m}"
);
}
#[test]
fn target_duration_falls_back_to_configured_when_segments_are_short() {
let s = MediaStore::new(4.0, 500, 4);
s.set_init(vec![0; 4]);
s.add_segment(plain_seg(1, 2)); let m = media_playlist_m3u8(&s, 1);
assert!(
m.contains("#EXT-X-TARGETDURATION:4"),
"unchanged behaviour when no segment exceeds the configured target: {m}"
);
}
}