Skip to main content

LlHlsOrigin

Struct LlHlsOrigin 

Source
pub struct LlHlsOrigin { /* private fields */ }
Available on crate feature std only.
Expand description

The LL-HLS origin ServedEgress: renders playlists and resolves blocking-reload/part-availability requests for one stream, backed by a shared Trunk. See this module’s own doc for exactly what comes straight from the Trunk and what needs the small synced Window.

Implementations§

Source§

impl LlHlsOrigin

Source

pub fn new( trunk: Arc<Trunk>, target_duration_secs: f64, part_target_ms: u32, window_segments: NonZeroUsize, ) -> Self

Build a fresh origin over trunk, subscribing its one SegmentCursor immediately (so the window starts empty but never misses a segment published from this point on).

window_segments bounds how many closed segments this origin advertises in a rendered Media Playlist — independent of media_plane::trunk::TrunkConfig::segment_capacity (the Trunk’s own retention bound): a caller may legitimately want a shorter advertised window than the Trunk retains for other consumers (e.g. a DVR SegmentEgress reading the same Trunk).

Examples found in repository?
examples/client_stepping.rs (line 43)
38fn canned_playlist() -> String {
39    let trunk = Trunk::new(TrunkConfig::new(nz(16), nz(4), nz(8), nz(4), nz(16)));
40    let writer = trunk
41        .segment_writer()
42        .expect("first (and only) segment writer");
43    let origin = LlHlsOrigin::new(std::sync::Arc::clone(&trunk), 1.0, 500, nz(4));
44    origin.set_init(vec![0xAA; 32]);
45
46    writer.publish_part(PartEntry::new(
47        vec![0x01; 16],
48        1,
49        0,
50        Duration::from_millis(500),
51        true,
52    ));
53    writer.publish_segment(SegmentEntry::new(
54        vec![0x02; 32],
55        1,
56        Duration::from_secs(1),
57        Timestamp::from_nanos(0),
58        SegmentMeta {
59            discontinuous: false,
60        },
61    ));
62    writer.publish_part(PartEntry::new(
63        vec![0x03; 16],
64        2,
65        0,
66        Duration::from_millis(500),
67        true,
68    ));
69
70    match origin.resolve(
71        LlHlsRequest::Playlist {
72            track_id: DEFAULT_TRACK_ID,
73            query: BlockingQuery::default(),
74        },
75        Timestamp::from_nanos(0),
76        AwaitPolicy::new(Timestamp::from_nanos(0)),
77    ) {
78        EgressResponse::Ready {
79            body: LlHlsBody::Playlist(m),
80            ..
81        } => m,
82        other => panic!("expected Ready(Playlist), got {other:?}"),
83    }
84}
More examples
Hide additional examples
examples/origin_playlist.rs (lines 57-62)
52fn main() {
53    let trunk = Trunk::new(TrunkConfig::new(nz(16), nz(4), nz(8), nz(4), nz(16)));
54    let writer = trunk
55        .segment_writer()
56        .expect("first (and only) segment writer");
57    let origin = LlHlsOrigin::new(
58        std::sync::Arc::clone(&trunk),
59        TARGET_DURATION_SECS,
60        PART_TARGET_MS,
61        nz(WINDOW_SEGMENTS),
62    );
63    origin.set_init(vec![0xAA; 32]);
64
65    // Segment 1 closes with two parts.
66    writer.publish_part(PartEntry::new(
67        vec![0x01; 16],
68        1,
69        0,
70        Duration::from_millis(500),
71        true,
72    ));
73    writer.publish_part(PartEntry::new(
74        vec![0x02; 16],
75        1,
76        1,
77        Duration::from_millis(500),
78        false,
79    ));
80    writer.publish_segment(SegmentEntry::new(
81        vec![0x03; 32],
82        1,
83        Duration::from_secs(1),
84        Timestamp::from_nanos(0),
85        SegmentMeta {
86            discontinuous: false,
87        },
88    ));
89
90    // Segment 2 is still open, with only its first part landed so far.
91    writer.publish_part(PartEntry::new(
92        vec![0x04; 16],
93        2,
94        0,
95        Duration::from_millis(500),
96        true,
97    ));
98
99    println!("--- master.m3u8 ---");
100    println!("{}", master_playlist_m3u8("media.m3u8"));
101
102    println!("--- media.m3u8 ---");
103    match resolve(
104        &origin,
105        LlHlsRequest::Playlist {
106            track_id: DEFAULT_TRACK_ID,
107            query: BlockingQuery::default(),
108        },
109    ) {
110        EgressResponse::Ready {
111            body: LlHlsBody::Playlist(m),
112            ..
113        } => println!("{m}"),
114        other => panic!("expected Ready(Playlist), got {other:?}"),
115    }
116
117    // A plain (non-blocking) request is Ready immediately.
118    let outcome = resolve(
119        &origin,
120        LlHlsRequest::Playlist {
121            track_id: DEFAULT_TRACK_ID,
122            query: BlockingQuery::default(),
123        },
124    );
125    assert!(matches!(
126        outcome,
127        EgressResponse::Ready {
128            body: LlHlsBody::Playlist(_),
129            ..
130        }
131    ));
132    println!("resolve(Playlist, no query)     -> Ready");
133
134    // A blocking-reload request for a segment that hasn't closed yet: with
135    // `await_policy`'s deadline already at `now`, this immediately reports
136    // the awaited condition has run out of patience rather than serving a
137    // fabricated Ready.
138    let outcome = resolve(
139        &origin,
140        LlHlsRequest::Playlist {
141            track_id: DEFAULT_TRACK_ID,
142            query: BlockingQuery {
143                hls_msn: Some(5),
144                hls_part: None,
145            },
146        },
147    );
148    assert_eq!(outcome, EgressResponse::NotFound);
149    println!("resolve(Playlist, _HLS_msn=5)   -> NotFound (Await's patience already expired)");
150
151    // A `_HLS_msn` unreasonably far beyond the live edge is rejected outright
152    // (RFC 8216bis §6.2.5.2 abuse prevention) rather than ever Await-ing.
153    let outcome = resolve(
154        &origin,
155        LlHlsRequest::Playlist {
156            track_id: DEFAULT_TRACK_ID,
157            query: BlockingQuery {
158                hls_msn: Some(999),
159                hls_part: None,
160            },
161        },
162    );
163    assert!(matches!(outcome, EgressResponse::BadRequest { .. }));
164    println!("resolve(Playlist, _HLS_msn=999) -> BadRequest (abuse bound)");
165
166    // `Resource`: the init segment and the closed segment are Ready...
167    match resolve(
168        &origin,
169        LlHlsRequest::Resource {
170            name: "init-1.mp4".to_string(),
171        },
172    ) {
173        EgressResponse::Ready { .. } => println!("resolve(Resource, init-1.mp4)     -> Ready"),
174        other => panic!("expected Ready, got {other:?}"),
175    }
176    match resolve(
177        &origin,
178        LlHlsRequest::Resource {
179            name: "seg-1-1.m4s".to_string(),
180        },
181    ) {
182        EgressResponse::Ready { .. } => println!("resolve(Resource, seg-1-1.m4s)    -> Ready"),
183        other => panic!("expected Ready, got {other:?}"),
184    }
185    // ...a live part of the still-open segment is Ready too...
186    match resolve(
187        &origin,
188        LlHlsRequest::Resource {
189            name: "part-1-2.0.m4s".to_string(),
190        },
191    ) {
192        EgressResponse::Ready { .. } => println!("resolve(Resource, part-1-2.0.m4s) -> Ready"),
193        other => panic!("expected Ready, got {other:?}"),
194    }
195    // ...a preload-hinted part not yet produced reports NotFound once this
196    // call's `await_policy` has already expired (a real HTTP adapter would
197    // instead give it a real deadline and block on `Trunk::listen()`)...
198    match resolve(
199        &origin,
200        LlHlsRequest::Resource {
201            name: "part-1-2.1.m4s".to_string(),
202        },
203    ) {
204        EgressResponse::NotFound => {
205            println!(
206                "resolve(Resource, part-1-2.1.m4s) -> NotFound (Await's patience already expired)"
207            )
208        }
209        other => panic!("expected NotFound, got {other:?}"),
210    }
211    // ...and an unrecognised filename is a plain 404.
212    match resolve(
213        &origin,
214        LlHlsRequest::Resource {
215            name: "nope.txt".to_string(),
216        },
217    ) {
218        EgressResponse::NotFound => println!("resolve(Resource, nope.txt)       -> NotFound"),
219        other => panic!("expected NotFound, got {other:?}"),
220    }
221}
Source

pub fn set_init(&self, bytes: impl Into<Bytes>)

Store the fMP4 init segment bytes — see this module’s doc for why an init segment is not something any Trunk ring holds.

Examples found in repository?
examples/client_stepping.rs (line 44)
38fn canned_playlist() -> String {
39    let trunk = Trunk::new(TrunkConfig::new(nz(16), nz(4), nz(8), nz(4), nz(16)));
40    let writer = trunk
41        .segment_writer()
42        .expect("first (and only) segment writer");
43    let origin = LlHlsOrigin::new(std::sync::Arc::clone(&trunk), 1.0, 500, nz(4));
44    origin.set_init(vec![0xAA; 32]);
45
46    writer.publish_part(PartEntry::new(
47        vec![0x01; 16],
48        1,
49        0,
50        Duration::from_millis(500),
51        true,
52    ));
53    writer.publish_segment(SegmentEntry::new(
54        vec![0x02; 32],
55        1,
56        Duration::from_secs(1),
57        Timestamp::from_nanos(0),
58        SegmentMeta {
59            discontinuous: false,
60        },
61    ));
62    writer.publish_part(PartEntry::new(
63        vec![0x03; 16],
64        2,
65        0,
66        Duration::from_millis(500),
67        true,
68    ));
69
70    match origin.resolve(
71        LlHlsRequest::Playlist {
72            track_id: DEFAULT_TRACK_ID,
73            query: BlockingQuery::default(),
74        },
75        Timestamp::from_nanos(0),
76        AwaitPolicy::new(Timestamp::from_nanos(0)),
77    ) {
78        EgressResponse::Ready {
79            body: LlHlsBody::Playlist(m),
80            ..
81        } => m,
82        other => panic!("expected Ready(Playlist), got {other:?}"),
83    }
84}
More examples
Hide additional examples
examples/origin_playlist.rs (line 63)
52fn main() {
53    let trunk = Trunk::new(TrunkConfig::new(nz(16), nz(4), nz(8), nz(4), nz(16)));
54    let writer = trunk
55        .segment_writer()
56        .expect("first (and only) segment writer");
57    let origin = LlHlsOrigin::new(
58        std::sync::Arc::clone(&trunk),
59        TARGET_DURATION_SECS,
60        PART_TARGET_MS,
61        nz(WINDOW_SEGMENTS),
62    );
63    origin.set_init(vec![0xAA; 32]);
64
65    // Segment 1 closes with two parts.
66    writer.publish_part(PartEntry::new(
67        vec![0x01; 16],
68        1,
69        0,
70        Duration::from_millis(500),
71        true,
72    ));
73    writer.publish_part(PartEntry::new(
74        vec![0x02; 16],
75        1,
76        1,
77        Duration::from_millis(500),
78        false,
79    ));
80    writer.publish_segment(SegmentEntry::new(
81        vec![0x03; 32],
82        1,
83        Duration::from_secs(1),
84        Timestamp::from_nanos(0),
85        SegmentMeta {
86            discontinuous: false,
87        },
88    ));
89
90    // Segment 2 is still open, with only its first part landed so far.
91    writer.publish_part(PartEntry::new(
92        vec![0x04; 16],
93        2,
94        0,
95        Duration::from_millis(500),
96        true,
97    ));
98
99    println!("--- master.m3u8 ---");
100    println!("{}", master_playlist_m3u8("media.m3u8"));
101
102    println!("--- media.m3u8 ---");
103    match resolve(
104        &origin,
105        LlHlsRequest::Playlist {
106            track_id: DEFAULT_TRACK_ID,
107            query: BlockingQuery::default(),
108        },
109    ) {
110        EgressResponse::Ready {
111            body: LlHlsBody::Playlist(m),
112            ..
113        } => println!("{m}"),
114        other => panic!("expected Ready(Playlist), got {other:?}"),
115    }
116
117    // A plain (non-blocking) request is Ready immediately.
118    let outcome = resolve(
119        &origin,
120        LlHlsRequest::Playlist {
121            track_id: DEFAULT_TRACK_ID,
122            query: BlockingQuery::default(),
123        },
124    );
125    assert!(matches!(
126        outcome,
127        EgressResponse::Ready {
128            body: LlHlsBody::Playlist(_),
129            ..
130        }
131    ));
132    println!("resolve(Playlist, no query)     -> Ready");
133
134    // A blocking-reload request for a segment that hasn't closed yet: with
135    // `await_policy`'s deadline already at `now`, this immediately reports
136    // the awaited condition has run out of patience rather than serving a
137    // fabricated Ready.
138    let outcome = resolve(
139        &origin,
140        LlHlsRequest::Playlist {
141            track_id: DEFAULT_TRACK_ID,
142            query: BlockingQuery {
143                hls_msn: Some(5),
144                hls_part: None,
145            },
146        },
147    );
148    assert_eq!(outcome, EgressResponse::NotFound);
149    println!("resolve(Playlist, _HLS_msn=5)   -> NotFound (Await's patience already expired)");
150
151    // A `_HLS_msn` unreasonably far beyond the live edge is rejected outright
152    // (RFC 8216bis §6.2.5.2 abuse prevention) rather than ever Await-ing.
153    let outcome = resolve(
154        &origin,
155        LlHlsRequest::Playlist {
156            track_id: DEFAULT_TRACK_ID,
157            query: BlockingQuery {
158                hls_msn: Some(999),
159                hls_part: None,
160            },
161        },
162    );
163    assert!(matches!(outcome, EgressResponse::BadRequest { .. }));
164    println!("resolve(Playlist, _HLS_msn=999) -> BadRequest (abuse bound)");
165
166    // `Resource`: the init segment and the closed segment are Ready...
167    match resolve(
168        &origin,
169        LlHlsRequest::Resource {
170            name: "init-1.mp4".to_string(),
171        },
172    ) {
173        EgressResponse::Ready { .. } => println!("resolve(Resource, init-1.mp4)     -> Ready"),
174        other => panic!("expected Ready, got {other:?}"),
175    }
176    match resolve(
177        &origin,
178        LlHlsRequest::Resource {
179            name: "seg-1-1.m4s".to_string(),
180        },
181    ) {
182        EgressResponse::Ready { .. } => println!("resolve(Resource, seg-1-1.m4s)    -> Ready"),
183        other => panic!("expected Ready, got {other:?}"),
184    }
185    // ...a live part of the still-open segment is Ready too...
186    match resolve(
187        &origin,
188        LlHlsRequest::Resource {
189            name: "part-1-2.0.m4s".to_string(),
190        },
191    ) {
192        EgressResponse::Ready { .. } => println!("resolve(Resource, part-1-2.0.m4s) -> Ready"),
193        other => panic!("expected Ready, got {other:?}"),
194    }
195    // ...a preload-hinted part not yet produced reports NotFound once this
196    // call's `await_policy` has already expired (a real HTTP adapter would
197    // instead give it a real deadline and block on `Trunk::listen()`)...
198    match resolve(
199        &origin,
200        LlHlsRequest::Resource {
201            name: "part-1-2.1.m4s".to_string(),
202        },
203    ) {
204        EgressResponse::NotFound => {
205            println!(
206                "resolve(Resource, part-1-2.1.m4s) -> NotFound (Await's patience already expired)"
207            )
208        }
209        other => panic!("expected NotFound, got {other:?}"),
210    }
211    // ...and an unrecognised filename is a plain 404.
212    match resolve(
213        &origin,
214        LlHlsRequest::Resource {
215            name: "nope.txt".to_string(),
216        },
217    ) {
218        EgressResponse::NotFound => println!("resolve(Resource, nope.txt)       -> NotFound"),
219        other => panic!("expected NotFound, got {other:?}"),
220    }
221}
Source

pub fn init_bytes(&self) -> Option<Bytes>

The fMP4 init segment bytes, if set.

Trait Implementations§

Source§

impl ServedEgress for LlHlsOrigin

Source§

type Request = LlHlsRequest

The protocol-specific request shape (e.g. LL-HLS’s Media Sequence Number + part index; a DASH segment/manifest request; a catch-up time range) — owned by the implementing crate, not this one.
Source§

type Body = LlHlsBody

The protocol-specific resolved body (a playlist string, manifest XML, segment bytes, …).
Source§

fn resolve( &self, request: LlHlsRequest, now: Timestamp, await_policy: AwaitPolicy, ) -> EgressResponse<LlHlsBody>

Resolve request against this implementation’s current state at time now, bounded by await_policy — see EgressResponse::Await’s contract.

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> PolicyExt for T
where T: ?Sized,

Source§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow only if self and other return Action::Follow. Read more
Source§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow if either self or other returns Action::Follow. Read more
Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V

Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more