citadel_sdk/media/
sender.rs1use super::error::MediaResultExt;
2use super::transport::{BoxedSink, MediaTransportKind};
3use bytes::{Bytes, BytesMut};
4use citadel_media::wire::encode_control;
5use citadel_media::{
6 ControlMessage, FrameFlags, FrameHeader, MediaConfig, MediaFrame, MediaStats,
7 MediaTrackDescriptor, Packetizer, SendQueue, TrackId, TrackKind,
8};
9use citadel_proto::prelude::NetworkError;
10
11pub struct MediaSender {
14 kind: MediaTransportKind,
15 media: BoxedSink,
16 control: BoxedSink,
17 packetizer: Packetizer,
18 queue: SendQueue,
19 scratch: BytesMut,
20 stats: MediaStats,
21}
22
23impl std::fmt::Debug for MediaSender {
24 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
25 f.debug_struct("MediaSender")
26 .field("kind", &self.kind)
27 .field("queued", &self.queue.len())
28 .field("stats", &self.stats)
29 .finish()
30 }
31}
32
33impl MediaSender {
34 pub(crate) fn new(
35 kind: MediaTransportKind,
36 media: BoxedSink,
37 control: BoxedSink,
38 config: MediaConfig,
39 send_queue_frames: usize,
40 ) -> Result<Self, NetworkError> {
41 Ok(Self {
42 kind,
43 media,
44 control,
45 packetizer: Packetizer::new(config).net()?,
46 queue: SendQueue::new(send_queue_frames).net()?,
47 scratch: BytesMut::with_capacity(config.max_fragment_payload * 2),
48 stats: MediaStats::new(),
49 })
50 }
51
52 pub fn kind(&self) -> MediaTransportKind {
53 self.kind
54 }
55
56 pub fn stats(&self) -> MediaStats {
57 self.stats
58 }
59
60 pub async fn announce(&mut self, tracks: &[MediaTrackDescriptor]) -> Result<(), NetworkError> {
62 self.send_control(&ControlMessage::AnnounceTracks(tracks.to_vec()))
63 }
64
65 pub async fn accept(&mut self, tracks: &[MediaTrackDescriptor]) -> Result<(), NetworkError> {
67 self.send_control(&ControlMessage::AcceptTracks(tracks.to_vec()))
68 }
69
70 pub async fn end_of_stream(&mut self, track: TrackId) -> Result<(), NetworkError> {
75 self.drain()?;
76 if !self.queue.is_empty() {
77 return Err(citadel_io::error!(
78 citadel_io::ErrorCode::MediaTransportClosed,
79 "cannot end stream: unsent frames remain in the send queue"
80 ));
81 }
82 let frames_sent = self.packetizer.next_sequence(track);
83 self.send_control(&ControlMessage::EndOfStream { track, frames_sent })
84 }
85
86 pub fn send_frame(
93 &mut self,
94 track: TrackId,
95 kind: TrackKind,
96 timestamp: u32,
97 flags: FrameFlags,
98 payload: Bytes,
99 ) -> Result<usize, NetworkError> {
100 let frame = MediaFrame {
101 header: FrameHeader {
102 track,
103 kind,
104 sequence: 0,
106 timestamp,
107 flags,
108 },
109 payload,
110 };
111 let dropped = usize::from(self.queue.push(frame).is_some());
112 self.stats.frames_dropped_on_send += dropped as u64;
113 self.drain()?;
114 Ok(dropped)
115 }
116
117 fn drain(&mut self) -> Result<(), NetworkError> {
118 loop {
119 let Some(frame) = self.queue.iter().next().cloned() else {
120 return Ok(());
121 };
122 self.send_one(&frame)?;
123 drop(self.queue.pop());
124 }
125 }
126
127 fn send_one(&mut self, frame: &MediaFrame) -> Result<(), NetworkError> {
128 let h = frame.header;
129 let fragments = self
130 .packetizer
131 .packetize(h.track, h.kind, h.timestamp, h.flags, frame.payload.clone())
132 .net()?;
133 for fragment in fragments {
134 let len = fragment.wire_len();
135 self.scratch.reserve(len);
136 fragment.write_into(&mut self.scratch);
137 self.media.send_datagram(self.scratch.split())?;
138 self.stats.fragments_sent += 1;
139 self.stats.bytes_sent += len as u64;
140 }
141 self.stats.frames_sent += 1;
142 Ok(())
143 }
144
145 fn send_control(&mut self, msg: &ControlMessage) -> Result<(), NetworkError> {
146 let body = msg.encode().net()?;
147 self.control
148 .send_datagram(BytesMut::from(encode_control(&body).as_slice()))
149 }
150}