pub struct LlHlsOrigin { /* private fields */ }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
impl LlHlsOrigin
Sourcepub fn new(
trunk: Arc<Trunk>,
target_duration_secs: f64,
part_target_ms: u32,
window_segments: NonZeroUsize,
) -> Self
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?
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
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}Sourcepub fn set_init(&self, bytes: impl Into<Bytes>)
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?
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
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}Sourcepub fn init_bytes(&self) -> Option<Bytes>
pub fn init_bytes(&self) -> Option<Bytes>
The fMP4 init segment bytes, if set.
Trait Implementations§
Source§impl ServedEgress for LlHlsOrigin
impl ServedEgress for LlHlsOrigin
Source§type Request = LlHlsRequest
type Request = LlHlsRequest
Source§type Body = LlHlsBody
type Body = LlHlsBody
Source§fn resolve(
&self,
request: LlHlsRequest,
now: Timestamp,
await_policy: AwaitPolicy,
) -> EgressResponse<LlHlsBody>
fn resolve( &self, request: LlHlsRequest, now: Timestamp, await_policy: AwaitPolicy, ) -> EgressResponse<LlHlsBody>
request against this implementation’s current state at
time now, bounded by await_policy — see
EgressResponse::Await’s contract.