Skip to main content

citadel_sdk/media/
receiver.rs

1use super::eos::EosTracker;
2use super::transport::{BoxedSource, MediaTransportKind};
3use citadel_io::time::{sleep_until, Duration, Instant};
4use citadel_io::ErrorCode;
5use citadel_media::{
6    ControlMessage, JitterBuffer, MediaConfig, MediaFrame, MediaInstant, MediaStats,
7    MediaTrackDescriptor, PopResult, PushResult, ReassembleOutcome, Reassembler, TrackId,
8};
9use citadel_proto::prelude::{NetworkError, SecBuffer};
10use futures::stream::{select_all, SelectAll};
11use futures::StreamExt;
12use std::collections::VecDeque;
13
14/// What [`MediaReceiver::next_event`] yields.
15#[derive(Debug, Clone, PartialEq, Eq)]
16pub enum MediaEvent {
17    /// The peer announced (or accepted) these tracks.
18    Tracks(Vec<MediaTrackDescriptor>),
19    /// A frame in sequence order for its track.
20    Frame(MediaFrame),
21    /// Frames `missing_from..=missing_to` on `track` were skipped; the next
22    /// event is the frame that follows the gap.
23    Gap {
24        track: TrackId,
25        missing_from: u32,
26        missing_to: u32,
27    },
28    /// The peer finished `track`.
29    EndOfStream(TrackId),
30    /// A transport stream ended; no further events will follow.
31    Closed,
32}
33
34enum Input {
35    Datagram(Option<SecBuffer>),
36    Deadline,
37}
38
39/// Inbound half of a media endpoint.
40///
41/// Owns the transport receive halves. In unreliable mode that includes the UDP
42/// receive half, so dropping this receiver signals `DisconnectUDP` to the peer.
43pub struct MediaReceiver {
44    kind: MediaTransportKind,
45    /// Media and control lanes merged; ends only once every lane has ended.
46    lanes: SelectAll<BoxedSource>,
47    reassembler: Reassembler,
48    jitter: JitterBuffer,
49    /// Events decoded ahead of delivery (frame following a gap, control).
50    ready: VecDeque<MediaEvent>,
51    /// End-of-stream announcements waiting for their track to fully drain.
52    eos: EosTracker,
53    start: Instant,
54    closed: bool,
55    stats: MediaStats,
56}
57
58impl std::fmt::Debug for MediaReceiver {
59    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
60        f.debug_struct("MediaReceiver")
61            .field("kind", &self.kind)
62            .field("closed", &self.closed)
63            .field("stats", &self.stats)
64            .finish()
65    }
66}
67
68impl MediaReceiver {
69    pub(crate) fn new(
70        kind: MediaTransportKind,
71        media: BoxedSource,
72        control: Option<BoxedSource>,
73        config: MediaConfig,
74        start: Instant,
75    ) -> Result<Self, NetworkError> {
76        Ok(Self {
77            kind,
78            lanes: select_all(std::iter::once(media).chain(control)),
79            reassembler: Reassembler::new(config).map_err(super::error::from_media_error)?,
80            jitter: JitterBuffer::new(config).map_err(super::error::from_media_error)?,
81            ready: VecDeque::new(),
82            eos: EosTracker::new(config.jitter_depth_micros),
83            start,
84            closed: false,
85            stats: MediaStats::new(),
86        })
87    }
88
89    pub fn kind(&self) -> MediaTransportKind {
90        self.kind
91    }
92
93    pub fn stats(&self) -> MediaStats {
94        self.stats
95    }
96
97    /// The only place wall time is read for the inbound media path.
98    fn now(&self) -> MediaInstant {
99        MediaInstant::from_micros(self.start.elapsed().as_micros() as u64)
100    }
101
102    fn to_instant(&self, at: MediaInstant) -> Instant {
103        self.start + Duration::from_micros(at.as_micros())
104    }
105
106    /// Waits for the next event. After [`MediaEvent::Closed`] every call
107    /// returns `Closed` again.
108    pub async fn next_event(&mut self) -> MediaEvent {
109        loop {
110            if let Some(event) = self.ready.pop_front() {
111                return event;
112            }
113            if let Some(event) = self.pop_jitter() {
114                return event;
115            }
116            if self.resolve_eos(false) {
117                continue;
118            }
119            if self.closed {
120                // Flush frames still parked behind a gap before reporting closure.
121                match self.jitter.next_deadline() {
122                    Some(at) if self.jitter.buffered_len() > 0 => {
123                        sleep_until(self.to_instant(at)).await;
124                        continue;
125                    }
126                    _ => {}
127                }
128                // No more input can arrive: expire pending end-of-stream
129                // records now instead of waiting out their deadlines.
130                if self.resolve_eos(true) {
131                    continue;
132                }
133                return MediaEvent::Closed;
134            }
135            match self.wait_input().await {
136                Input::Datagram(Some(buf)) => self.ingest(buf.as_ref()),
137                Input::Datagram(None) => self.closed = true,
138                Input::Deadline => {
139                    let now = self.now();
140                    self.stats.frames_evicted_incomplete +=
141                        self.reassembler.evict_stale(now) as u64;
142                }
143            }
144        }
145    }
146
147    /// Bridges [`EosTracker::resolve`] to this receiver's state.
148    fn resolve_eos(&mut self, force: bool) -> bool {
149        let now = self.now();
150        let Self {
151            eos,
152            jitter,
153            ready,
154            stats,
155            ..
156        } = self;
157        eos.resolve(
158            now,
159            force,
160            |track| jitter.next_expected(track),
161            ready,
162            stats,
163        )
164    }
165
166    fn pop_jitter(&mut self) -> Option<MediaEvent> {
167        let now = self.now();
168        match self.jitter.pop_ready(now) {
169            PopResult::Frame(frame) => {
170                self.stats.frames_delivered += 1;
171                Some(MediaEvent::Frame(frame))
172            }
173            PopResult::Gap {
174                track,
175                missing_from,
176                missing_to,
177                next,
178            } => {
179                self.stats.gaps_skipped += 1;
180                self.stats.frames_missing +=
181                    u64::from(missing_to.wrapping_sub(missing_from).wrapping_add(1));
182                self.stats.frames_delivered += 1;
183                self.ready.push_back(MediaEvent::Frame(next));
184                Some(MediaEvent::Gap {
185                    track,
186                    missing_from,
187                    missing_to,
188                })
189            }
190            PopResult::NotReady => None,
191        }
192    }
193
194    async fn wait_input(&mut self) -> Input {
195        let deadline = self
196            .jitter
197            .next_deadline()
198            .into_iter()
199            .chain(self.eos.next_deadline())
200            .min()
201            .map(|at| self.to_instant(at));
202        let timer = async move {
203            match deadline {
204                Some(at) => sleep_until(at).await,
205                None => futures::future::pending().await,
206            }
207        };
208        citadel_io::tokio::select! {
209            item = self.lanes.next() => Input::Datagram(item),
210            _ = timer => Input::Deadline,
211        }
212    }
213
214    fn ingest(&mut self, datagram: &[u8]) {
215        let now = self.now();
216        self.stats.fragments_received += 1;
217        self.stats.bytes_received += datagram.len() as u64;
218        match self.reassembler.push(datagram, now) {
219            ReassembleOutcome::Complete(frame) => {
220                self.stats.frames_completed += 1;
221                match self.jitter.push(frame, now) {
222                    PushResult::Buffered => {}
223                    PushResult::Late => self.stats.frames_late += 1,
224                    PushResult::TooOld => self.stats.frames_too_old += 1,
225                    PushResult::Duplicate => self.stats.fragments_duplicate += 1,
226                }
227            }
228            ReassembleOutcome::Partial { .. } => {}
229            ReassembleOutcome::Duplicate => self.stats.fragments_duplicate += 1,
230            ReassembleOutcome::Control(body) => match Self::decode_control(&body) {
231                Ok(ControlMessage::AnnounceTracks(list) | ControlMessage::AcceptTracks(list)) => {
232                    self.ready.push_back(MediaEvent::Tracks(list))
233                }
234                Ok(ControlMessage::EndOfStream { track, frames_sent }) => {
235                    self.eos.record(track, frames_sent, now)
236                }
237                Err(err) => {
238                    self.stats.fragments_rejected += 1;
239                    log::warn!(target: "citadel", "{err}");
240                }
241            },
242            ReassembleOutcome::Rejected(err) => {
243                self.stats.fragments_rejected += 1;
244                log::debug!(target: "citadel", "media fragment rejected: {err}");
245            }
246        }
247    }
248
249    fn decode_control(body: &[u8]) -> Result<ControlMessage, NetworkError> {
250        ControlMessage::decode(body)
251            .map_err(|err| citadel_io::error!(ErrorCode::MediaControlDecode, err))
252    }
253}