arcly_stream/protocol/rtsp/
egress.rs1use 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
28pub struct RtspServer {
30 playback: Arc<dyn PlaybackRegistry>,
31 bind: SocketAddr,
32}
33
34impl RtspServer {
35 pub fn new(playback: Arc<dyn PlaybackRegistry>, bind: SocketAddr) -> Self {
37 Self { playback, bind }
38 }
39
40 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
66const PT_H264: u8 = 96;
68const 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 let session = new_session_id();
79 loop {
80 let Some(req) = read_request(&mut sock, &mut buf).await? else {
81 return Ok(()); };
83 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 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
129async 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); 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
151struct RtspSink<'a> {
153 sock: &'a mut TcpStream,
154 packetizer: RtpPacketizer,
155 pkts: Vec<Vec<u8>>,
157 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 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(()); }
190 let timestamp = (frame.pts.max(0) as u64).wrapping_mul(90) as u32; packetizer.packetize_into(&frame.data, timestamp, pkts);
194 wbuf.clear();
197 for pkt in pkts.iter() {
198 frame_interleaved(RTP_CHANNEL, pkt, wbuf);
199 }
200 sock.write_all(wbuf).await
201}
202
203fn 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
212async 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
237fn 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
248fn 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 "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
302fn 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 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}