use alloc::collections::{BTreeMap, BTreeSet, VecDeque};
use alloc::string::{String, ToString};
use alloc::vec::Vec;
use broadcast_common::Unpackage;
use transmux::hls::{ByteRange, MapTag, PreloadHintType};
use transmux::{Fmp4Demux, MediaPlaylist, MediaSegment, OpenSegment, TrackSpec, TsDemux};
use super::action::{Action, BlockingReload, ResourceId};
use super::error::{Error, Result};
use super::output::Output;
use super::url;
const TS_SYNC_BYTE: u8 = 0x47;
#[derive(Debug)]
pub struct LlHlsClient {
playlist_url: String,
pending_actions: VecDeque<Action>,
pending_outputs: VecDeque<Output>,
init_uri: Option<String>,
init_bytes: Option<Vec<u8>>,
init_emitted: bool,
pending_demux: VecDeque<(ResourceId, Vec<u8>)>,
requested: BTreeSet<ResourceId>,
delivered_parts: BTreeSet<(u64, u64)>,
delivered_segments: BTreeSet<u64>,
discontinuous_msns: BTreeSet<u64>,
discontinuity_emitted: BTreeSet<u64>,
byte_range_cursor: BTreeMap<String, u64>,
outstanding_fetches: u64,
saw_endlist: bool,
end_emitted: bool,
last_full_playlist: Option<MediaPlaylist>,
}
impl LlHlsClient {
pub fn new(playlist_url: impl Into<String>) -> Self {
let playlist_url = playlist_url.into();
let mut pending_actions = VecDeque::new();
pending_actions.push_back(Action::FetchPlaylist {
url: playlist_url.clone(),
blocking: None,
skip: false,
});
Self {
playlist_url,
pending_actions,
pending_outputs: VecDeque::new(),
init_uri: None,
init_bytes: None,
init_emitted: false,
pending_demux: VecDeque::new(),
requested: BTreeSet::new(),
delivered_parts: BTreeSet::new(),
delivered_segments: BTreeSet::new(),
discontinuous_msns: BTreeSet::new(),
discontinuity_emitted: BTreeSet::new(),
byte_range_cursor: BTreeMap::new(),
outstanding_fetches: 0,
saw_endlist: false,
end_emitted: false,
last_full_playlist: None,
}
}
pub fn playlist_url(&self) -> &str {
&self.playlist_url
}
pub fn poll(&mut self) -> Option<Action> {
self.pending_actions.pop_front()
}
pub fn next_output(&mut self) -> Option<Output> {
self.pending_outputs.pop_front()
}
pub fn on_playlist(&mut self, bytes: &[u8]) -> Result<()> {
let text = core::str::from_utf8(bytes)?;
let playlist = MediaPlaylist::parse(text)?;
let playlist = self.merge_delta(playlist);
for (i, seg) in playlist.segments.iter().enumerate() {
let msn = playlist.media_sequence + i as u64;
self.process_closed_segment(msn, seg)?;
}
let next_msn = playlist.media_sequence + playlist.segments.len() as u64;
if let Some(open) = &playlist.open_segment {
self.process_open_segment(next_msn, open)?;
}
let map = playlist
.open_segment
.as_ref()
.and_then(|o| o.map.as_ref())
.or_else(|| playlist.segments.last().and_then(|s| s.map.as_ref()));
if let Some(map) = map {
self.ensure_init_requested(map)?;
}
if let Some(ll) = &playlist.low_latency {
if let Some(hint_uri) = &ll.preload_hint_part {
match ll.preload_hint_type {
PreloadHintType::Part => {
let part_idx = playlist
.open_segment
.as_ref()
.map(|o| o.parts.len() as u64)
.unwrap_or(0);
let id = ResourceId::Part {
msn: next_msn,
part: part_idx,
};
let url = url::resolve(&self.playlist_url, hint_uri);
let byte_range = self.resolve_hint_byte_range(&url, ll);
self.request_resource(id, url, byte_range);
}
PreloadHintType::Map => {
let map = MapTag {
uri: hint_uri.clone(),
byte_range: ll.preload_hint_byte_range_length.map(|length| ByteRange {
length,
offset: ll.preload_hint_byte_range_start,
}),
};
self.ensure_init_requested(&map)?;
}
_ => {
}
}
}
}
if playlist.endlist {
self.saw_endlist = true;
} else {
let blocking = playlist
.low_latency
.as_ref()
.filter(|ll| ll.can_block_reload)
.map(|_| {
let part = playlist
.open_segment
.as_ref()
.map(|o| o.parts.len() as u64)
.unwrap_or(0);
BlockingReload {
msn: next_msn,
part: Some(part),
}
});
let can_skip = playlist
.low_latency
.as_ref()
.and_then(|ll| ll.can_skip_until)
.is_some();
let skip = can_skip && self.last_full_playlist.is_some();
self.pending_actions.push_back(Action::FetchPlaylist {
url: self.playlist_url.clone(),
blocking,
skip,
});
if blocking.is_none() {
let wait_ms = (u64::from(playlist.target_duration.max(1)) * 1000) / 2;
self.pending_actions.push_back(Action::WaitMs(wait_ms));
}
}
if playlist.skip.is_none() {
self.last_full_playlist = Some(playlist);
}
self.maybe_emit_end_of_stream();
Ok(())
}
pub fn on_resource(&mut self, id: ResourceId, bytes: &[u8]) -> Result<()> {
let was_requested = match id {
ResourceId::Init => self.init_uri.is_some(),
ResourceId::Part { .. } | ResourceId::Segment { .. } => self.requested.contains(&id),
};
if !was_requested {
return Err(Error::UnrequestedResource { id });
}
self.outstanding_fetches = self.outstanding_fetches.saturating_sub(1);
match id {
ResourceId::Init => {
self.init_bytes = Some(bytes.to_vec());
if !self.init_emitted {
self.pending_outputs.push_back(Output::Init(bytes.to_vec()));
self.init_emitted = true;
}
let buffered: Vec<_> = self.pending_demux.drain(..).collect();
for (bid, bbytes) in buffered {
self.finish_media_resource(bid, &bbytes)?;
}
}
ResourceId::Part { .. } | ResourceId::Segment { .. } => {
if self.is_ts_segment(bytes) {
self.finish_ts_resource(id, bytes)?;
} else if self.init_bytes.is_none() {
self.pending_demux.push_back((id, bytes.to_vec()));
} else {
self.finish_media_resource(id, bytes)?;
}
}
}
self.maybe_emit_end_of_stream();
Ok(())
}
fn finish_media_resource(&mut self, id: ResourceId, bytes: &[u8]) -> Result<()> {
match id {
ResourceId::Part { msn, part } => {
self.emit_discontinuity_if_needed(msn);
self.demux_and_emit(id, bytes)?;
self.delivered_parts.insert((msn, part));
}
ResourceId::Segment { msn } => {
self.emit_discontinuity_if_needed(msn);
self.demux_and_emit(id, bytes)?;
self.delivered_segments.insert(msn);
}
ResourceId::Init => {}
}
Ok(())
}
fn finish_ts_resource(&mut self, id: ResourceId, bytes: &[u8]) -> Result<()> {
match id {
ResourceId::Part { msn, part } => {
self.emit_discontinuity_if_needed(msn);
self.demux_and_emit_ts(id, bytes)?;
self.delivered_parts.insert((msn, part));
}
ResourceId::Segment { msn } => {
self.emit_discontinuity_if_needed(msn);
self.demux_and_emit_ts(id, bytes)?;
self.delivered_segments.insert(msn);
}
ResourceId::Init => {}
}
Ok(())
}
fn is_ts_segment(&self, bytes: &[u8]) -> bool {
self.init_uri.is_none() && bytes.first() == Some(&TS_SYNC_BYTE)
}
pub fn on_error(&mut self, id: Option<ResourceId>) {
if let Some(id) = id {
self.outstanding_fetches = self.outstanding_fetches.saturating_sub(1);
match id {
ResourceId::Init => self.init_uri = None,
other => {
self.requested.remove(&other);
}
}
}
self.maybe_emit_end_of_stream();
}
fn merge_delta(&self, playlist: MediaPlaylist) -> MediaPlaylist {
let Some(skip) = &playlist.skip else {
return playlist;
};
if skip.skipped_segments == 0 {
return playlist;
}
let Some(prev) = &self.last_full_playlist else {
return playlist;
};
if playlist.media_sequence < prev.media_sequence {
return playlist;
}
let prefix_start = (playlist.media_sequence - prev.media_sequence) as usize;
let prefix_end = prefix_start + skip.skipped_segments as usize;
let Some(prefix) = prev.segments.get(prefix_start..prefix_end) else {
return playlist;
};
let mut merged = playlist;
let mut segments = prefix.to_vec();
segments.extend(merged.segments);
merged.segments = segments;
merged
}
fn process_closed_segment(&mut self, msn: u64, seg: &MediaSegment) -> Result<()> {
if seg.discontinuous {
self.discontinuous_msns.insert(msn);
}
if self.delivered_segments.contains(&msn) {
return Ok(());
}
if seg.parts.is_empty() {
let already_have_parts = self
.delivered_parts
.range((msn, 0)..(msn + 1, 0))
.next()
.is_some();
if already_have_parts {
self.delivered_segments.insert(msn);
return Ok(());
}
let id = ResourceId::Segment { msn };
if !self.requested.contains(&id) {
let url = url::resolve(&self.playlist_url, &seg.uri);
let byte_range = self.resolve_byte_range(&url, &seg.byte_range);
self.request_resource(id, url, byte_range);
}
return Ok(());
}
let mut fully_accounted = true;
for (i, part) in seg.parts.iter().enumerate() {
let i = i as u64;
if part.gap || self.delivered_parts.contains(&(msn, i)) {
continue;
}
fully_accounted = false;
let id = ResourceId::Part { msn, part: i };
if !self.requested.contains(&id) {
let url = url::resolve(&self.playlist_url, &part.uri);
let byte_range = self.resolve_byte_range(&url, &part.byte_range);
self.request_resource(id, url, byte_range);
}
}
if fully_accounted {
self.delivered_segments.insert(msn);
}
Ok(())
}
fn process_open_segment(&mut self, msn: u64, open: &OpenSegment) -> Result<()> {
for (i, part) in open.parts.iter().enumerate() {
let i = i as u64;
if part.gap || self.delivered_parts.contains(&(msn, i)) {
continue;
}
let id = ResourceId::Part { msn, part: i };
if !self.requested.contains(&id) {
let url = url::resolve(&self.playlist_url, &part.uri);
let byte_range = self.resolve_byte_range(&url, &part.byte_range);
self.request_resource(id, url, byte_range);
}
}
Ok(())
}
fn ensure_init_requested(&mut self, map: &MapTag) -> Result<()> {
let url = url::resolve(&self.playlist_url, &map.uri);
if self.init_uri.as_deref() == Some(url.as_str()) {
return Ok(());
}
self.init_uri = Some(url.clone());
self.init_bytes = None;
self.init_emitted = false;
let byte_range = self.resolve_byte_range(&url, &map.byte_range);
self.pending_actions.push_back(Action::FetchResource {
id: ResourceId::Init,
url,
byte_range,
});
self.outstanding_fetches += 1;
Ok(())
}
fn request_resource(&mut self, id: ResourceId, url: String, byte_range: Option<(u64, u64)>) {
self.requested.insert(id);
self.outstanding_fetches += 1;
self.pending_actions.push_back(Action::FetchResource {
id,
url,
byte_range,
});
}
fn resolve_byte_range(&mut self, url: &str, br: &Option<ByteRange>) -> Option<(u64, u64)> {
let br = br.as_ref()?;
let offset = br
.offset
.unwrap_or_else(|| *self.byte_range_cursor.get(url).unwrap_or(&0));
self.byte_range_cursor
.insert(url.to_string(), offset + br.length);
Some((offset, br.length))
}
fn resolve_hint_byte_range(
&mut self,
url: &str,
ll: &transmux::hls::LowLatencyConfig,
) -> Option<(u64, u64)> {
let length = ll.preload_hint_byte_range_length?;
let br = ByteRange {
length,
offset: ll.preload_hint_byte_range_start,
};
self.resolve_byte_range(url, &Some(br))
}
fn demux_and_emit(&mut self, id: ResourceId, bytes: &[u8]) -> Result<()> {
let init = self
.init_bytes
.as_ref()
.ok_or(Error::InitNotYetAvailable { id })?;
let mut combined = Vec::with_capacity(init.len() + bytes.len());
combined.extend_from_slice(init);
combined.extend_from_slice(bytes);
let mut demux = Fmp4Demux::new();
let media = demux
.unpackage(combined.as_slice())
.map_err(|source| Error::Demux { id, source })?;
for track in media.tracks {
if !track.samples.is_empty() {
self.pending_outputs.push_back(Output::Samples {
track_id: track.spec.track_id,
samples: track.samples,
});
}
}
Ok(())
}
fn demux_and_emit_ts(&mut self, id: ResourceId, bytes: &[u8]) -> Result<()> {
let mut demux = TsDemux::new();
let media = demux
.demux(bytes)
.map_err(|source| Error::Demux { id, source })?;
if !self.init_emitted {
let specs: Vec<TrackSpec> = media.tracks.iter().map(|t| t.spec.clone()).collect();
let init_bytes = transmux::build_init_segment(&specs, media.movie_timescale)
.map_err(|source| Error::Demux { id, source })?;
self.pending_outputs.push_back(Output::Init(init_bytes));
self.init_emitted = true;
}
for track in media.tracks {
if !track.samples.is_empty() {
self.pending_outputs.push_back(Output::Samples {
track_id: track.spec.track_id,
samples: track.samples,
});
}
}
Ok(())
}
fn emit_discontinuity_if_needed(&mut self, msn: u64) {
if self.discontinuous_msns.contains(&msn) && !self.discontinuity_emitted.contains(&msn) {
self.pending_outputs.push_back(Output::Discontinuity);
self.discontinuity_emitted.insert(msn);
}
}
fn maybe_emit_end_of_stream(&mut self) {
if self.saw_endlist && !self.end_emitted && self.outstanding_fetches == 0 {
self.pending_outputs.push_back(Output::EndOfStream);
self.end_emitted = true;
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn on_resource_rejects_a_never_requested_id() {
let mut client = LlHlsClient::new("http://example.com/playlist.m3u8");
let id = ResourceId::Segment { msn: 0 };
let err = client
.on_resource(id, b"some bytes")
.expect_err("an id the client never requested must be rejected");
assert!(
matches!(err, Error::UnrequestedResource { id: got } if got == id),
"wrong error variant: {err:?}"
);
let err = client
.on_resource(ResourceId::Init, b"init bytes")
.expect_err("an unrequested Init must be rejected");
assert!(
matches!(
err,
Error::UnrequestedResource {
id: ResourceId::Init
}
),
"wrong error variant: {err:?}"
);
}
#[test]
fn on_resource_accepts_a_previously_requested_id() {
let mut client = LlHlsClient::new("http://example.com/playlist.m3u8");
let id = ResourceId::Segment { msn: 0 };
client.request_resource(id, "http://example.com/seg0.m4s".to_string(), None);
let result = client.on_resource(id, b"some bytes");
assert!(
result.is_ok(),
"a requested id must be accepted: {result:?}"
);
assert!(
client.pending_demux.iter().any(|(bid, _)| *bid == id),
"expected the resource to be buffered pending the init segment"
);
}
#[test]
fn is_ts_segment_true_when_no_map_seen_and_sync_byte_present() {
let client = LlHlsClient::new("http://example.com/playlist.m3u8");
assert!(client.is_ts_segment(&[TS_SYNC_BYTE, 0x40, 0x11, 0x00]));
}
#[test]
fn is_ts_segment_false_for_an_isobmff_resource_with_no_map_seen() {
let client = LlHlsClient::new("http://example.com/playlist.m3u8");
let ftyp_box = b"\x00\x00\x00\x18ftypiso5\x00\x00\x02\x00iso5iso6mp41";
assert!(!client.is_ts_segment(ftyp_box));
}
#[test]
fn is_ts_segment_false_once_a_map_has_been_requested() {
let mut client = LlHlsClient::new("http://example.com/playlist.m3u8");
client
.ensure_init_requested(&MapTag {
uri: "init.mp4".to_string(),
byte_range: None,
})
.expect("ensure_init_requested succeeds");
assert!(!client.is_ts_segment(&[TS_SYNC_BYTE, 0x40, 0x11, 0x00]));
}
}