Skip to main content

citadel_sdk/media/
sender.rs

1use 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
11/// Outbound half of a media endpoint. `send_frame` is synchronous and never
12/// blocks; control messages (`announce`, `end_of_stream`) are async.
13pub 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    /// Tells the peer which tracks this side will send.
61    pub async fn announce(&mut self, tracks: &[MediaTrackDescriptor]) -> Result<(), NetworkError> {
62        self.send_control(&ControlMessage::AnnounceTracks(tracks.to_vec()))
63    }
64
65    /// Acknowledges a peer's announcement.
66    pub async fn accept(&mut self, tracks: &[MediaTrackDescriptor]) -> Result<(), NetworkError> {
67        self.send_control(&ControlMessage::AcceptTracks(tracks.to_vec()))
68    }
69
70    /// Marks `track` finished on the peer side, telling it how many frames
71    /// (`0..frames_sent`) were sent so it can drain in-flight media that races
72    /// this control message. Fails fast if frames are still parked in the send
73    /// queue after a drain attempt (a stale count would lie to the peer).
74    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    /// Queues a frame and drains the queue to the transport. Returns how many
87    /// frames the queue's drop policy evicted to make room (0 normally).
88    ///
89    /// Sequence numbers are assigned at packetization, so evicted frames never
90    /// consume one. If the transport rejects a datagram the frame stays at the
91    /// head of the queue (retried on the next call) and the error is returned.
92    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                // Placeholder: the packetizer assigns the real sequence on drain.
105                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}