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#[derive(Debug, Clone, PartialEq, Eq)]
16pub enum MediaEvent {
17 Tracks(Vec<MediaTrackDescriptor>),
19 Frame(MediaFrame),
21 Gap {
24 track: TrackId,
25 missing_from: u32,
26 missing_to: u32,
27 },
28 EndOfStream(TrackId),
30 Closed,
32}
33
34enum Input {
35 Datagram(Option<SecBuffer>),
36 Deadline,
37}
38
39pub struct MediaReceiver {
44 kind: MediaTransportKind,
45 lanes: SelectAll<BoxedSource>,
47 reassembler: Reassembler,
48 jitter: JitterBuffer,
49 ready: VecDeque<MediaEvent>,
51 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 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 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 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 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 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}