Skip to main content

arcly_stream/protocol/rtsp/
egress.rs

1//! RTSP **egress server** — serve a live stream to RTSP players (VLC, ffmpeg).
2//!
3//! The serving counterpart to the client-pull [ingest handler](super::RtspHandler):
4//! [`RtspServer`] accepts TCP connections, runs the
5//! `OPTIONS → DESCRIBE → SETUP → PLAY → TEARDOWN` state machine, and streams the
6//! requested stream's H.264 access units as RTP over the TCP-interleaved
7//! transport (RFC 2326 §10.12), reusing the shared
8//! [`RtpPacketizer`](crate::protocol::rtp::RtpPacketizer).
9//!
10//! Interleaved (RTP-over-TCP) transport only — the universally-supported path
11//! that needs no separate UDP ports. The stream is selected from the request
12//! URI's `/app/stream` path.
13
14use std::net::SocketAddr;
15use std::ops::ControlFlow;
16use std::sync::Arc;
17
18use tokio::io::{AsyncReadExt, AsyncWriteExt};
19use tokio::net::{TcpListener, TcpStream};
20use tokio_util::sync::CancellationToken;
21use tracing::{debug, info, warn};
22
23use super::message::RtspRequest;
24use crate::bus::PlaybackRegistry;
25use crate::protocol::rtp::RtpPacketizer;
26use crate::{CodecId, MediaFrame, Result, StreamKey};
27
28/// An RTSP server that serves live streams to pulling players.
29pub struct RtspServer {
30    playback: Arc<dyn PlaybackRegistry>,
31    bind: SocketAddr,
32}
33
34impl RtspServer {
35    /// Serve streams from `playback`, listening on `bind`.
36    pub fn new(playback: Arc<dyn PlaybackRegistry>, bind: SocketAddr) -> Self {
37        Self { playback, bind }
38    }
39
40    /// Accept and serve RTSP connections until `shutdown` fires.
41    pub async fn run(self, shutdown: CancellationToken) -> Result<()> {
42        let listener = TcpListener::bind(self.bind).await?;
43        info!(bind = %self.bind, "rtsp egress server listening");
44        loop {
45            tokio::select! {
46                _ = shutdown.cancelled() => break,
47                accepted = listener.accept() => {
48                    let (sock, peer) = match accepted {
49                        Ok(v) => v,
50                        Err(e) => { warn!(error = %e, "rtsp accept failed"); continue; }
51                    };
52                    let playback = Arc::clone(&self.playback);
53                    let shutdown = shutdown.clone();
54                    tokio::spawn(async move {
55                        if let Err(e) = serve_connection(sock, playback, shutdown).await {
56                            debug!(%peer, error = %e, "rtsp egress connection ended");
57                        }
58                    });
59                }
60            }
61        }
62        Ok(())
63    }
64}
65
66/// The RTP payload type advertised for H.264 (must match the SDP `rtpmap`).
67const PT_H264: u8 = 96;
68/// The interleaved channel for RTP (RTCP would be channel 1).
69const RTP_CHANNEL: u8 = 0;
70
71async fn serve_connection(
72    mut sock: TcpStream,
73    playback: Arc<dyn PlaybackRegistry>,
74    shutdown: CancellationToken,
75) -> Result<()> {
76    let mut buf = Vec::with_capacity(2048);
77    // A distinct session id per connection (a real server tracks state per id).
78    let session = new_session_id();
79    loop {
80        let Some(req) = read_request(&mut sock, &mut buf).await? else {
81            return Ok(()); // peer closed
82        };
83        // The stream a request addresses, if its URI names a live one.
84        let live = stream_key_from_uri(&req.uri).filter(|k| playback.get_stream(k).is_ok());
85        match req.method.as_str() {
86            "OPTIONS" => {
87                sock.write_all(options_response(req.cseq).as_bytes())
88                    .await?
89            }
90            "DESCRIBE" => match live {
91                Some(_) => {
92                    let sdp = build_sdp();
93                    sock.write_all(describe_response(req.cseq, &req.uri, &sdp).as_bytes())
94                        .await?;
95                }
96                None => {
97                    sock.write_all(not_found(req.cseq).as_bytes()).await?;
98                }
99            },
100            "SETUP" => {
101                sock.write_all(setup_response(req.cseq, &session).as_bytes())
102                    .await?;
103            }
104            "PLAY" => match live {
105                Some(key) => {
106                    sock.write_all(play_response(req.cseq, &session).as_bytes())
107                        .await?;
108                    return play(sock, &playback, key, shutdown).await;
109                }
110                None => {
111                    // Unknown stream: 404, but keep the connection so a player
112                    // can DESCRIBE/SETUP another stream.
113                    sock.write_all(not_found(req.cseq).as_bytes()).await?;
114                }
115            },
116            "TEARDOWN" => {
117                sock.write_all(simple_ok(req.cseq, &session).as_bytes())
118                    .await?;
119                return Ok(());
120            }
121            other => {
122                debug!(method = other, "rtsp: unsupported method");
123                sock.write_all(not_implemented(req.cseq).as_bytes()).await?;
124            }
125        }
126    }
127}
128
129/// Stream `key`'s H.264 access units as interleaved RTP until the stream ends,
130/// the client disconnects, or `shutdown` fires.
131async fn play(
132    mut sock: TcpStream,
133    playback: &Arc<dyn PlaybackRegistry>,
134    key: StreamKey,
135    shutdown: CancellationToken,
136) -> Result<()> {
137    let handle = playback.get_stream(&key)?;
138    let packetizer = RtpPacketizer::new(PT_H264, 0x5254_5350, 1400); // "RTSP"
139
140    // Replay the instant-start buffer then forward live frames; a write error
141    // means the player disconnected, which ends the session gracefully.
142    let mut sink = RtspSink {
143        sock: &mut sock,
144        packetizer,
145        pkts: Vec::new(),
146        wbuf: Vec::with_capacity(1500),
147    };
148    handle.drive_to(&shutdown, &mut sink).await
149}
150
151/// Streams a stream's frames to one RTSP player as interleaved RTP-over-TCP.
152struct RtspSink<'a> {
153    sock: &'a mut TcpStream,
154    packetizer: RtpPacketizer,
155    /// Reused across frames: the per-frame RTP packet buffers.
156    pkts: Vec<Vec<u8>>,
157    /// Reused across frames: the interleaved-framing write buffer.
158    wbuf: Vec<u8>,
159}
160
161#[async_trait::async_trait]
162impl crate::bus::FrameSink for RtspSink<'_> {
163    async fn send(&mut self, frame: Arc<MediaFrame>) -> Result<ControlFlow<()>> {
164        // A write error means the player disconnected — stop the drive cleanly.
165        match send_frame(
166            self.sock,
167            &mut self.packetizer,
168            &mut self.pkts,
169            &mut self.wbuf,
170            &frame,
171        )
172        .await
173        {
174            Ok(()) => Ok(ControlFlow::Continue(())),
175            Err(_) => Ok(ControlFlow::Break(())),
176        }
177    }
178}
179
180async fn send_frame(
181    sock: &mut TcpStream,
182    packetizer: &mut RtpPacketizer,
183    pkts: &mut Vec<Vec<u8>>,
184    wbuf: &mut Vec<u8>,
185    frame: &MediaFrame,
186) -> std::io::Result<()> {
187    if !frame.is_video() || frame.codec != CodecId::H264 {
188        return Ok(()); // egress is H.264 video over RTSP for now
189    }
190    // Clamp to 0 before scaling: a negative PTS cast straight to u64 would become
191    // an enormous value and make the RTP timestamp jump wildly.
192    let timestamp = (frame.pts.max(0) as u64).wrapping_mul(90) as u32; // ms → 90 kHz
193    packetizer.packetize_into(&frame.data, timestamp, pkts);
194    // Coalesce every interleaved RTP packet of this access unit into one write
195    // (and one reused buffer) rather than a syscall + allocation per packet.
196    wbuf.clear();
197    for pkt in pkts.iter() {
198        frame_interleaved(RTP_CHANNEL, pkt, wbuf);
199    }
200    sock.write_all(wbuf).await
201}
202
203/// Append an RTP packet framed for the RTSP TCP-interleaved transport to `out`:
204/// `$` + 1-byte channel + 2-byte big-endian length + the RTP packet.
205fn frame_interleaved(channel: u8, rtp: &[u8], out: &mut Vec<u8>) {
206    out.push(b'$');
207    out.push(channel);
208    out.extend_from_slice(&(rtp.len() as u16).to_be_bytes());
209    out.extend_from_slice(rtp);
210}
211
212/// Read one RTSP request (headers terminated by a blank line) from `sock`,
213/// returning `None` on a clean EOF.
214async fn read_request(sock: &mut TcpStream, buf: &mut Vec<u8>) -> Result<Option<RtspRequest>> {
215    let mut tmp = [0u8; 1024];
216    loop {
217        if let Some(end) = find_double_crlf(buf) {
218            let head = String::from_utf8_lossy(&buf[..end]).into_owned();
219            buf.drain(..end);
220            return Ok(RtspRequest::parse(&head).map(Some).unwrap_or(None));
221        }
222        let n = sock.read(&mut tmp).await?;
223        if n == 0 {
224            return Ok(None);
225        }
226        buf.extend_from_slice(&tmp[..n]);
227        if buf.len() > 64 * 1024 {
228            return Err(crate::StreamError::protocol("rtsp request too large"));
229        }
230    }
231}
232
233fn find_double_crlf(buf: &[u8]) -> Option<usize> {
234    buf.windows(4).position(|w| w == b"\r\n\r\n").map(|p| p + 4)
235}
236
237/// Map an RTSP URI (`rtsp://host[:port]/app/stream[/trackID][?q]`) to a stream.
238fn stream_key_from_uri(uri: &str) -> Option<StreamKey> {
239    let rest = uri.strip_prefix("rtsp://")?;
240    let path = rest.split_once('/').map(|(_, p)| p)?;
241    let path = path.split(['?', ';']).next().unwrap_or(path);
242    let mut segs = path.split('/').filter(|s| !s.is_empty());
243    let app = segs.next()?;
244    let stream = segs.next()?;
245    Some(StreamKey::new(app, stream))
246}
247
248// ── Response builders ────────────────────────────────────────────────────────
249
250fn options_response(cseq: u32) -> String {
251    format!(
252        "RTSP/1.0 200 OK\r\nCSeq: {cseq}\r\n\
253         Public: OPTIONS, DESCRIBE, SETUP, PLAY, TEARDOWN\r\n\r\n"
254    )
255}
256
257fn build_sdp() -> String {
258    // Minimal H.264 elementary description; the player learns SPS/PPS in-band.
259    "v=0\r\n\
260     o=- 0 0 IN IP4 0.0.0.0\r\n\
261     s=arcly-stream\r\n\
262     t=0 0\r\n\
263     m=video 0 RTP/AVP 96\r\n\
264     a=rtpmap:96 H264/90000\r\n\
265     a=control:streamid=0\r\n"
266        .to_string()
267}
268
269fn describe_response(cseq: u32, uri: &str, sdp: &str) -> String {
270    format!(
271        "RTSP/1.0 200 OK\r\nCSeq: {cseq}\r\n\
272         Content-Base: {uri}\r\nContent-Type: application/sdp\r\n\
273         Content-Length: {}\r\n\r\n{sdp}",
274        sdp.len()
275    )
276}
277
278fn setup_response(cseq: u32, session: &str) -> String {
279    format!(
280        "RTSP/1.0 200 OK\r\nCSeq: {cseq}\r\n\
281         Transport: RTP/AVP/TCP;unicast;interleaved=0-1\r\n\
282         Session: {session}\r\n\r\n"
283    )
284}
285
286fn play_response(cseq: u32, session: &str) -> String {
287    format!("RTSP/1.0 200 OK\r\nCSeq: {cseq}\r\nSession: {session}\r\nRange: npt=0.000-\r\n\r\n")
288}
289
290fn simple_ok(cseq: u32, session: &str) -> String {
291    format!("RTSP/1.0 200 OK\r\nCSeq: {cseq}\r\nSession: {session}\r\n\r\n")
292}
293
294fn not_implemented(cseq: u32) -> String {
295    format!("RTSP/1.0 501 Not Implemented\r\nCSeq: {cseq}\r\n\r\n")
296}
297
298fn not_found(cseq: u32) -> String {
299    format!("RTSP/1.0 404 Not Found\r\nCSeq: {cseq}\r\n\r\n")
300}
301
302/// A fresh, hard-to-guess RTSP session id (hex), unique per connection. Seeded
303/// from the OS-backed `RandomState`, so no RNG dependency is pulled in.
304fn new_session_id() -> String {
305    use std::collections::hash_map::RandomState;
306    use std::hash::{BuildHasher, Hasher};
307    use std::time::{SystemTime, UNIX_EPOCH};
308    let mut h = RandomState::new().build_hasher();
309    h.write_u128(
310        SystemTime::now()
311            .duration_since(UNIX_EPOCH)
312            .map(|d| d.as_nanos())
313            .unwrap_or(0),
314    );
315    format!("{:016X}", h.finish())
316}
317
318#[cfg(test)]
319mod tests {
320    use super::*;
321    use crate::protocol::rtsp::message::{InterleavedFrame, RtspResponse};
322
323    #[test]
324    fn interleaved_frame_round_trips() {
325        let rtp = [0x80u8, 0xE0, 0x00, 0x01, 0xAA, 0xBB];
326        let mut framed = Vec::new();
327        frame_interleaved(0, &rtp, &mut framed);
328        assert_eq!(framed[0], b'$');
329        let (f, used) = InterleavedFrame::parse(&framed).expect("parse");
330        assert_eq!(used, framed.len());
331        assert_eq!(f.channel, 0);
332        assert_eq!(f.payload, &rtp);
333    }
334
335    #[test]
336    fn stream_key_parsed_from_uri_variants() {
337        let k = stream_key_from_uri("rtsp://host:554/live/cam").unwrap();
338        assert_eq!((k.app.as_str(), k.stream_id.as_str()), ("live", "cam"));
339        // trailing track/control + query are stripped.
340        let k = stream_key_from_uri("rtsp://h/live/cam/streamid=0?x=1").unwrap();
341        assert_eq!((k.app.as_str(), k.stream_id.as_str()), ("live", "cam"));
342        assert!(stream_key_from_uri("rtsp://host/onlyapp").is_none());
343    }
344
345    #[test]
346    fn responses_are_well_formed() {
347        assert!(options_response(2).contains("Public: OPTIONS"));
348        let d = describe_response(3, "rtsp://h/live/cam", &build_sdp());
349        let parsed = RtspResponse::parse(
350            d.split("\r\n\r\n").next().unwrap(),
351            d.split("\r\n\r\n").nth(1).unwrap_or("").to_string(),
352        )
353        .expect("response parses");
354        assert_eq!(parsed.header("Content-Type"), Some("application/sdp"));
355        let setup = setup_response(4, "DEADBEEF");
356        assert!(setup.contains("interleaved=0-1"));
357        assert!(setup.contains("Session: DEADBEEF"));
358        assert!(not_found(7).starts_with("RTSP/1.0 404"));
359    }
360}