1use std::time::{Duration, Instant};
15
16use bytes::Bytes;
17use hang::catalog::{AudioCodec, VideoCodec};
18use moq_mux::catalog::hang::Catalog;
19use str0m::format::Codec;
20use str0m::media::{Frequency, MediaTime, Mid, Pt};
21use tokio::sync::mpsc;
22
23use crate::{Error, Result, codec};
24
25pub struct WriteRequest {
30 pub mid: Mid,
32 pub pt: Pt,
34 pub time: MediaTime,
36 pub payload: Bytes,
38}
39
40#[derive(Default)]
48pub(crate) struct EgressClock {
49 anchor: Option<(Duration, Instant)>,
50}
51
52impl EgressClock {
53 pub(crate) fn wallclock(&mut self, time: MediaTime, now: Instant) -> Instant {
55 let presentation = Duration::from(time);
56 let Some((anchor_presentation, anchor_wallclock)) = self.anchor else {
57 self.anchor = Some((presentation, now));
58 return now;
59 };
60
61 if presentation >= anchor_presentation {
62 let delta = presentation - anchor_presentation;
63 let Some(mapped) = anchor_wallclock.checked_add(delta) else {
64 self.anchor = Some((presentation, now));
65 return now;
66 };
67 if mapped > now {
68 self.anchor = Some((presentation, now));
72 now
73 } else {
74 mapped
75 }
76 } else {
77 anchor_wallclock
78 .checked_sub(anchor_presentation - presentation)
79 .unwrap_or(now)
80 }
81 }
82}
83
84pub struct EgressSource {
86 broadcast: moq_net::BroadcastConsumer,
87 catalog: Catalog,
91 writes_tx: mpsc::Sender<WriteRequest>,
92 writes_rx: Option<mpsc::Receiver<WriteRequest>>,
93}
94
95impl EgressSource {
96 pub async fn new(broadcast: moq_net::BroadcastConsumer) -> Result<Self> {
102 let catalog_track = broadcast.subscribe_track(&moq_net::Track::new(hang::Catalog::DEFAULT_NAME))?;
103 let mut consumer = moq_mux::catalog::hang::Consumer::new(catalog_track);
104 let catalog = consumer
105 .next()
106 .await
107 .map_err(|err| Error::Other(anyhow::anyhow!("catalog subscribe: {err}")))?
108 .ok_or_else(|| Error::Other(anyhow::anyhow!("catalog closed before first snapshot")))?;
109
110 let (tx, rx) = mpsc::channel(64);
111 Ok(Self {
112 broadcast,
113 catalog,
114 writes_tx: tx,
115 writes_rx: Some(rx),
116 })
117 }
118
119 pub fn take_writes(&mut self) -> mpsc::Receiver<WriteRequest> {
122 self.writes_rx.take().expect("EgressSource writes_rx already taken")
123 }
124
125 pub fn on_track(&mut self, mid: Mid, codec: Codec, pt: Pt, clock_rate: Frequency) -> Result<()> {
131 let tx = self.writes_tx.clone();
134 let broadcast = self.broadcast.clone();
135 let catalog = self.catalog.clone();
136 tokio::spawn(async move {
137 let track = match pick_track(&broadcast, &catalog, codec).await {
138 Ok(Some(t)) => t,
139 Ok(None) => {
140 tracing::warn!(?codec, "no matching catalog rendition; egress track ignored");
141 return;
142 }
143 Err(err) => {
144 tracing::warn!(?codec, %err, "egress track subscribe failed");
145 return;
146 }
147 };
148 pump(mid, pt, clock_rate, track, tx).await;
149 });
150 Ok(())
151 }
152
153 pub fn catalog_codecs(&self) -> Vec<Codec> {
157 let mut out = Vec::new();
158 if self
159 .catalog
160 .audio
161 .renditions
162 .values()
163 .any(|r| matches!(r.codec, AudioCodec::Opus))
164 {
165 out.push(Codec::Opus);
166 }
167 for rendition in self.catalog.video.renditions.values() {
168 if let Some(c) = video_codec(&rendition.codec)
169 && !out.contains(&c)
170 {
171 out.push(c);
172 }
173 }
174 out
175 }
176}
177
178fn video_codec(codec: &VideoCodec) -> Option<Codec> {
180 match codec {
181 VideoCodec::H264(_) => Some(Codec::H264),
182 VideoCodec::H265(_) => Some(Codec::H265),
183 VideoCodec::VP8 => Some(Codec::Vp8),
184 VideoCodec::VP9(_) => Some(Codec::Vp9),
185 VideoCodec::AV1(_) => Some(Codec::Av1),
186 _ => None,
187 }
188}
189
190async fn pick_track(
193 broadcast: &moq_net::BroadcastConsumer,
194 catalog: &Catalog,
195 codec: Codec,
196) -> Result<Option<codec::Track>> {
197 match codec {
198 Codec::Opus => {
199 let Some((name, _config)) = catalog
200 .audio
201 .renditions
202 .iter()
203 .find(|(_, c)| matches!(c.codec, AudioCodec::Opus))
204 else {
205 return Ok(None);
206 };
207 Ok(Some(codec::Track::opus(broadcast, name).await?))
208 }
209 Codec::H264 | Codec::H265 | Codec::Vp8 | Codec::Vp9 | Codec::Av1 => {
210 let Some((name, config)) = catalog
211 .video
212 .renditions
213 .iter()
214 .find(|(_, c)| video_codec(&c.codec) == Some(codec))
215 else {
216 return Ok(None);
217 };
218 Ok(Some(codec::Track::video(broadcast, name, config).await?))
219 }
220 other => Err(Error::UnsupportedCodec(format!("{other:?}"))),
221 }
222}
223
224async fn pump(mid: Mid, pt: Pt, clock_rate: Frequency, mut track: codec::Track, tx: mpsc::Sender<WriteRequest>) {
227 loop {
228 let frame = match track.next().await {
229 Ok(Some(f)) => f,
230 Ok(None) => {
231 tracing::debug!(?mid, "egress track ended");
232 return;
233 }
234 Err(err) => {
235 tracing::warn!(?mid, %err, "egress track error");
236 return;
237 }
238 };
239 let ticks = us_to_ticks(frame.timestamp_us, clock_rate);
240 let time = MediaTime::new(ticks, clock_rate);
241 let req = WriteRequest {
242 mid,
243 pt,
244 time,
245 payload: frame.payload,
246 };
247 if tx.send(req).await.is_err() {
248 return;
250 }
251 }
252}
253
254fn us_to_ticks(timestamp_us: u64, clock_rate: Frequency) -> u64 {
257 let rate = clock_rate.get() as u128;
258 ((timestamp_us as u128 * rate) / 1_000_000) as u64
259}
260
261pub fn dispatch(rtc: &mut str0m::Rtc, request: WriteRequest, wallclock: Instant) {
268 let Some(writer) = rtc.writer(request.mid) else {
269 tracing::debug!(?request.mid, "egress write before media available");
270 return;
271 };
272 let WriteRequest {
273 pt,
274 time,
275 payload,
276 mid: _,
277 } = request;
278 if let Err(err) = writer.write(pt, wallclock, time, payload.to_vec()) {
279 tracing::warn!(%err, "egress write rejected by str0m");
280 }
281}
282
283#[cfg(test)]
284mod tests {
285 use super::*;
286
287 #[test]
288 fn egress_clock_ignores_cross_track_dequeue_jitter() {
289 let mut clock = EgressClock::default();
290 let t0 = Instant::now();
291
292 assert_eq!(clock.wallclock(MediaTime::from_millis(1_000), t0), t0);
293 assert_eq!(
294 clock.wallclock(MediaTime::from_millis(1_100), t0 + Duration::from_millis(100)),
295 t0 + Duration::from_millis(100)
296 );
297
298 let audio = clock.wallclock(MediaTime::from_millis(1_200), t0 + Duration::from_millis(250));
301 let video = clock.wallclock(MediaTime::from_millis(1_200), t0 + Duration::from_millis(300));
302 assert_eq!(audio, t0 + Duration::from_millis(200));
303 assert_eq!(video, audio);
304 }
305
306 #[test]
307 fn egress_clock_moves_epoch_earlier_for_catch_up_bursts() {
308 let mut clock = EgressClock::default();
309 let t0 = Instant::now();
310
311 assert_eq!(clock.wallclock(MediaTime::from_millis(1_000), t0), t0);
312 assert_eq!(clock.wallclock(MediaTime::from_millis(1_100), t0), t0);
315
316 assert_eq!(
319 clock.wallclock(MediaTime::from_millis(1_100), t0 + Duration::from_millis(50)),
320 t0
321 );
322 }
323}