1mod engine;
30mod playlist;
31
32use engine::{SegmentEngine, BANDWIDTH_HINT};
33
34#[cfg(feature = "mpegts")]
35#[cfg_attr(docsrs, doc(cfg(feature = "mpegts")))]
36mod mpegts;
37
38#[cfg(feature = "mpegts")]
39#[cfg_attr(docsrs, doc(cfg(feature = "mpegts")))]
40pub use mpegts::MpegTsMuxer;
41
42#[cfg(feature = "fmp4")]
43#[cfg_attr(docsrs, doc(cfg(feature = "fmp4")))]
44mod fmp4;
45
46#[cfg(feature = "fmp4")]
47#[cfg_attr(docsrs, doc(cfg(feature = "fmp4")))]
48pub use fmp4::Fmp4Muxer;
49
50pub use playlist::{render_master, HlsPlaylist, Part, Segment};
51
52#[cfg(feature = "dash")]
53#[cfg_attr(docsrs, doc(cfg(feature = "dash")))]
54mod dash;
55
56#[cfg(feature = "dash")]
57#[cfg_attr(docsrs, doc(cfg(feature = "dash")))]
58pub use dash::{DashManifest, DashPackager};
59
60use crate::traits::StorageBackend;
61use crate::{CodecId, FrameFlags, MediaFrame, Result};
62use async_trait::async_trait;
63use bytes::{BufMut, Bytes, BytesMut};
64
65pub trait Muxer: Send {
71 fn extension(&self) -> &'static str;
73
74 fn start_segment(&mut self) -> Result<()>;
76
77 fn write(&mut self, frame: &MediaFrame) -> Result<()>;
79
80 fn finish_segment(&mut self) -> Result<Bytes>;
82
83 fn take_partial(&mut self) -> Result<Option<Bytes>> {
92 Ok(None)
93 }
94
95 fn init_segment(&mut self, _codec: CodecId, _config_record: &[u8]) -> Result<Option<Bytes>> {
102 Ok(None)
103 }
104
105 fn build_init_from(&mut self, configs: &[(CodecId, Bytes)]) -> Result<Option<Bytes>> {
115 match configs.first() {
116 Some((codec, data)) => self.init_segment(*codec, data),
117 None => Ok(None),
118 }
119 }
120
121 fn codec_string(&self) -> Option<String> {
124 None
125 }
126}
127
128pub struct PassthroughMuxer {
133 ext: &'static str,
134 buf: BytesMut,
135}
136
137impl PassthroughMuxer {
138 pub fn new(ext: &'static str) -> Self {
140 Self {
141 ext,
142 buf: BytesMut::new(),
143 }
144 }
145}
146
147impl Muxer for PassthroughMuxer {
148 fn extension(&self) -> &'static str {
149 self.ext
150 }
151 fn start_segment(&mut self) -> Result<()> {
152 self.buf.clear();
153 Ok(())
154 }
155 fn write(&mut self, frame: &MediaFrame) -> Result<()> {
156 self.buf.put_slice(&frame.data);
157 Ok(())
158 }
159 fn finish_segment(&mut self) -> Result<Bytes> {
160 Ok(std::mem::take(&mut self.buf).freeze())
161 }
162 fn take_partial(&mut self) -> Result<Option<Bytes>> {
163 if self.buf.is_empty() {
164 return Ok(None);
165 }
166 Ok(Some(std::mem::take(&mut self.buf).freeze()))
167 }
168}
169
170#[async_trait]
172pub trait Packager: Send {
173 async fn push(&mut self, frame: &MediaFrame) -> Result<()>;
175 async fn finish(&mut self) -> Result<()>;
177}
178
179pub struct HlsSegmenter<M: Muxer, S: StorageBackend> {
185 engine: SegmentEngine<M, S>,
186 playlist: HlsPlaylist,
187 master_written: bool,
190 codec_string: Option<String>,
195 part_target_ms: i64,
198 part_idx: u64,
200 part_anchor_pts: Option<i64>,
202 last_pts: i64,
204 seg_part_bytes: Vec<Bytes>,
207}
208
209impl<M: Muxer, S: StorageBackend> HlsSegmenter<M, S> {
210 pub fn new(
213 muxer: M,
214 storage: S,
215 prefix: impl Into<String>,
216 target_duration: u64,
217 window: usize,
218 ) -> Self {
219 Self {
220 engine: SegmentEngine::new(muxer, storage, prefix, target_duration),
221 playlist: HlsPlaylist::new(target_duration, window),
222 master_written: false,
223 codec_string: None,
224 part_target_ms: 0,
225 part_idx: 0,
226 part_anchor_pts: None,
227 last_pts: 0,
228 seg_part_bytes: Vec::new(),
229 }
230 }
231
232 pub fn codec_string(&self) -> Option<String> {
234 self.codec_string
235 .clone()
236 .or_else(|| self.engine.muxer.codec_string())
237 }
238
239 async fn ensure_init_segment(&mut self, frame: &MediaFrame) -> Result<()> {
245 if self.codec_string.is_none()
246 && frame.is_video()
247 && frame.flags.contains(FrameFlags::CONFIG)
248 {
249 self.codec_string = crate::codec::dispatch::parse_config(frame.codec, &frame.data)
250 .and_then(|p| crate::codec::dispatch::hls_codec_string(frame.codec, &p));
251 }
252 if let Some(uri) = self.engine.ensure_init(frame).await? {
253 self.playlist.set_map(uri);
254 }
255 Ok(())
256 }
257
258 pub fn low_latency(mut self, part_target: f64) -> Self {
264 self.playlist = self.playlist.low_latency(part_target);
265 self.part_target_ms = (part_target.max(0.05) * 1000.0) as i64;
266 self
267 }
268
269 pub fn playlist_key(&self) -> String {
271 format!("{}/index.m3u8", self.engine.prefix)
272 }
273
274 pub fn master_key(&self) -> String {
276 format!("{}/master.m3u8", self.engine.prefix)
277 }
278
279 async fn ensure_master_playlist(&mut self) -> Result<()> {
284 if self.master_written {
285 return Ok(());
286 }
287 let Some(codecs) = self.codec_string() else {
288 return Ok(());
289 };
290 let body = playlist::render_master("index.m3u8", Some(&codecs), BANDWIDTH_HINT);
291 self.engine
292 .storage
293 .put(
294 &self.engine.key("master.m3u8"),
295 Bytes::from(body.into_bytes()),
296 )
297 .await?;
298 self.master_written = true;
299 Ok(())
300 }
301
302 async fn flush_part(&mut self, end_pts: i64) -> Result<()> {
306 let Some(bytes) = self.engine.muxer.take_partial()? else {
307 return Ok(());
308 };
309 let anchor = self.part_anchor_pts.unwrap_or(end_pts);
310 let duration = (end_pts - anchor).max(0) as f64 / 1000.0;
311 let uri = self.engine.part_uri(self.engine.seq, self.part_idx);
312 let key = self.engine.key(&uri);
313 self.engine.storage.put(&key, bytes.clone()).await?;
314 self.seg_part_bytes.push(bytes);
315 self.playlist.add_pending_part(Part {
316 uri,
317 duration,
318 independent: self.part_idx == 0,
320 });
321 self.playlist
322 .set_preload_hint(self.engine.part_uri(self.engine.seq, self.part_idx + 1));
323 self.part_idx += 1;
324 self.part_anchor_pts = Some(end_pts);
325 self.write_playlist().await
326 }
327
328 async fn cut(&mut self, duration: f64) -> Result<()> {
329 let uri = self.engine.segment_uri(self.engine.seq);
330
331 if self.part_target_ms > 0 {
332 self.flush_part(self.last_pts).await?;
335 let mut full = BytesMut::new();
336 for b in &self.seg_part_bytes {
337 full.put_slice(b);
338 }
339 let full = if full.is_empty() {
343 self.engine.muxer.finish_segment()?
344 } else {
345 full.freeze()
346 };
347 let key = self.engine.key(&uri);
348 self.engine.storage.put(&key, full).await?;
349 self.playlist.commit_segment(Segment {
350 seq: self.engine.seq,
351 duration,
352 uri,
353 discontinuity: false,
354 parts: Vec::new(), });
356 self.seg_part_bytes.clear();
357 self.part_idx = 0;
358 self.part_anchor_pts = None;
359 } else {
360 let bytes = self.engine.muxer.finish_segment()?;
361 let key = self.engine.key(&uri);
362 self.engine.storage.put(&key, bytes).await?;
363 self.playlist.push(Segment {
364 seq: self.engine.seq,
365 duration,
366 uri,
367 discontinuity: false,
368 parts: Vec::new(),
369 });
370 }
371 self.write_playlist().await?;
372 self.engine.seq += 1;
373 Ok(())
374 }
375
376 async fn write_playlist(&mut self) -> Result<()> {
377 let body = self.playlist.render();
378 self.engine
379 .storage
380 .put(
381 &self.engine.key("index.m3u8"),
382 Bytes::from(body.into_bytes()),
383 )
384 .await
385 }
386}
387
388#[async_trait]
389impl<M: Muxer, S: StorageBackend> Packager for HlsSegmenter<M, S> {
390 async fn push(&mut self, frame: &MediaFrame) -> Result<()> {
391 self.ensure_init_segment(frame).await?;
392 self.ensure_master_playlist().await?;
397 let decision = self.engine.observe(frame);
398 if decision.skip {
399 return Ok(());
400 }
401 if let Some(duration) = decision.cut_previous {
402 self.cut(duration).await?;
403 }
404 if decision.open_new {
405 self.engine.muxer.start_segment()?;
406 self.part_anchor_pts = Some(frame.pts);
408 self.part_idx = 0;
409 self.seg_part_bytes.clear();
410 }
411 self.engine.muxer.write(frame)?;
412 self.last_pts = frame.pts;
413
414 if self.part_target_ms > 0 {
417 if let Some(anchor) = self.part_anchor_pts {
418 if frame.pts - anchor >= self.part_target_ms {
419 self.flush_part(frame.pts).await?;
420 }
421 }
422 }
423 Ok(())
424 }
425
426 async fn finish(&mut self) -> Result<()> {
427 if let Some(duration) = self.engine.flush() {
428 self.cut(duration).await?;
429 }
430 self.playlist.finish();
431 self.write_playlist().await
432 }
433}
434
435#[cfg(test)]
436mod tests {
437 use super::*;
438 use crate::testing::{video_frame, InMemoryStorage};
439
440 #[tokio::test]
441 async fn segments_on_keyframe_after_target_and_writes_playlist() {
442 let store = InMemoryStorage::new();
443 let mut seg =
444 HlsSegmenter::new(PassthroughMuxer::new("ts"), store.clone(), "live/cam", 2, 5);
445
446 for i in 0..5 {
449 let pts = i * 1000;
450 seg.push(&video_frame(pts, true)).await.unwrap();
451 seg.push(&video_frame(pts + 500, false)).await.unwrap();
453 }
454 seg.finish().await.unwrap();
455
456 let playlist = store.get("live/cam/index.m3u8").await.unwrap();
457 let text = String::from_utf8(playlist.to_vec()).unwrap();
458 assert!(text.contains("#EXTM3U"));
459 assert!(text.contains("#EXT-X-ENDLIST"));
460 assert!(store.get("live/cam/seg0.ts").await.is_ok());
462 assert!(!store.get("live/cam/seg0.ts").await.unwrap().is_empty());
464 }
465
466 #[tokio::test]
467 async fn ll_hls_emits_part_files_and_ext_x_part() {
468 let store = InMemoryStorage::new();
469 let mut seg =
471 HlsSegmenter::new(PassthroughMuxer::new("m4s"), store.clone(), "live/ll", 2, 5)
472 .low_latency(0.5);
473
474 seg.push(&video_frame(0, true)).await.unwrap();
476 for i in 1..8 {
477 seg.push(&video_frame(i * 250, false)).await.unwrap();
478 }
479 seg.push(&video_frame(2000, true)).await.unwrap(); seg.push(&video_frame(2250, false)).await.unwrap();
481 seg.finish().await.unwrap();
482
483 let text =
484 String::from_utf8(store.get("live/ll/index.m3u8").await.unwrap().to_vec()).unwrap();
485 assert!(text.contains("#EXT-X-PART-INF:PART-TARGET=0.500"));
487 assert!(
488 text.contains("#EXT-X-PART:DURATION="),
489 "renders EXT-X-PART lines"
490 );
491 assert!(
492 text.contains("INDEPENDENT=YES"),
493 "first part is independent"
494 );
495 assert!(
497 store.get("live/ll/seg0.0.m4s").await.is_ok(),
498 "part 0 written"
499 );
500 assert!(
501 store.get("live/ll/seg0.1.m4s").await.is_ok(),
502 "part 1 written"
503 );
504 assert!(
505 store.get("live/ll/seg0.m4s").await.is_ok(),
506 "full segment 0 written"
507 );
508
509 let mut concat = Vec::new();
511 let mut p = 0;
512 while let Ok(part) = store.get(&format!("live/ll/seg0.{p}.m4s")).await {
513 concat.extend_from_slice(&part);
514 p += 1;
515 }
516 assert!(p >= 2, "segment 0 had multiple parts");
517 let full = store.get("live/ll/seg0.m4s").await.unwrap();
518 assert_eq!(&concat[..], &full[..], "full segment == concat of parts");
519 }
520
521 #[tokio::test]
522 async fn ll_hls_falls_back_to_full_segment_when_muxer_has_no_parts() {
523 let store = InMemoryStorage::new();
526 let mut seg = HlsSegmenter::new(
527 InitMuxer {
528 buf: BytesMut::new(),
529 },
530 store.clone(),
531 "live/nopart",
532 2,
533 5,
534 )
535 .low_latency(0.5);
536
537 seg.push(&video_frame(0, true)).await.unwrap();
538 for i in 1..6 {
539 seg.push(&video_frame(i * 1000, i % 2 == 0)).await.unwrap();
540 }
541 seg.finish().await.unwrap();
542
543 let seg0 = store.get("live/nopart/seg0.m4s").await.unwrap();
544 assert!(!seg0.is_empty(), "LL-HLS segment must not be empty");
545 }
546
547 #[tokio::test]
548 async fn writes_master_playlist_with_hevc_codecs() {
549 use crate::FrameFlags;
550 let store = InMemoryStorage::new();
551 let mut seg = HlsSegmenter::new(
552 InitMuxer {
553 buf: BytesMut::new(),
554 },
555 store.clone(),
556 "live/hevc",
557 2,
558 5,
559 );
560
561 let mut cfg = video_frame(0, true);
562 cfg.codec = CodecId::H265;
563 cfg.flags |= FrameFlags::CONFIG;
564 seg.push(&cfg).await.unwrap();
565 for i in 1..4 {
566 seg.push(&video_frame(i * 1000, true)).await.unwrap();
567 }
568 seg.finish().await.unwrap();
569
570 let master =
571 String::from_utf8(store.get("live/hevc/master.m3u8").await.unwrap().to_vec()).unwrap();
572 assert!(master.contains("#EXT-X-STREAM-INF:"));
573 assert!(
574 master.contains("CODECS=\"hvc1.1.6.L120.B0\""),
575 "master advertises HEVC codec string: {master}"
576 );
577 assert!(master.contains("index.m3u8"));
578 }
579
580 struct InitMuxer {
582 buf: BytesMut,
583 }
584 impl Muxer for InitMuxer {
585 fn extension(&self) -> &'static str {
586 "m4s"
587 }
588 fn start_segment(&mut self) -> Result<()> {
589 self.buf.clear();
590 Ok(())
591 }
592 fn write(&mut self, frame: &MediaFrame) -> Result<()> {
593 self.buf.put_slice(&frame.data);
594 Ok(())
595 }
596 fn finish_segment(&mut self) -> Result<Bytes> {
597 Ok(std::mem::take(&mut self.buf).freeze())
598 }
599 fn init_segment(&mut self, _codec: CodecId, config_record: &[u8]) -> Result<Option<Bytes>> {
600 Ok(Some(Bytes::copy_from_slice(config_record)))
602 }
603 fn codec_string(&self) -> Option<String> {
604 Some("hvc1.1.6.L120.B0".into())
605 }
606 }
607
608 #[tokio::test]
609 async fn fmp4_muxer_writes_init_segment_and_ext_x_map() {
610 use crate::FrameFlags;
611 let store = InMemoryStorage::new();
612 let mut seg = HlsSegmenter::new(
613 InitMuxer {
614 buf: BytesMut::new(),
615 },
616 store.clone(),
617 "live/hevc",
618 2,
619 5,
620 );
621 assert_eq!(seg.codec_string().as_deref(), Some("hvc1.1.6.L120.B0"));
622
623 let mut cfg = video_frame(0, true);
625 cfg.codec = CodecId::H265;
626 cfg.flags |= FrameFlags::CONFIG;
627 seg.push(&cfg).await.unwrap();
628 for i in 1..6 {
629 seg.push(&video_frame(i * 1000, true)).await.unwrap();
630 }
631 seg.finish().await.unwrap();
632
633 assert!(store.get("live/hevc/init.m4s").await.is_ok());
634 let pl =
635 String::from_utf8(store.get("live/hevc/index.m3u8").await.unwrap().to_vec()).unwrap();
636 assert!(pl.contains("#EXT-X-MAP:URI=\"init.m4s\""));
637 }
638}