use std::net::SocketAddr;
use std::ops::ControlFlow;
use std::sync::Arc;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::net::{TcpListener, TcpStream, UdpSocket};
use tokio_util::sync::CancellationToken;
use tracing::{debug, info, warn};
use super::message::{InterleavedFrame, RtspRequest};
use super::sdp::Sdp;
use crate::auth::{Credentials, EgressGate};
use crate::bus::PlaybackRegistry;
use crate::inbound::IngestContext;
use crate::protocol::rtp::{
AccessUnit, DepacketizeError, H264Depacketizer, H265Depacketizer, RtpHeader, RtpPacketizer,
};
use crate::{CodecId, MediaFrame, Result, StreamKey};
const MAX_CONSECUTIVE_DEPACK_FAILURES: u32 = 64;
fn publish_au(
session: &crate::inbound::PublishSession,
codec: CodecId,
au: AccessUnit,
sdp_config: &Option<bytes::Bytes>,
) -> Result<()> {
let pts = (au.timestamp / 90) as i64; if au.keyframe {
let cfg =
crate::codec::dispatch::parameter_sets(codec, &au.data).or_else(|| sdp_config.clone());
if let Some(params) = cfg {
let mut frame = MediaFrame::new_video(pts, pts, params, codec, true);
frame.flags |= crate::FrameFlags::CONFIG;
session.publish_frame(frame)?;
}
}
let mf = MediaFrame::new_video(pts, pts, au.data, codec, au.keyframe);
session.publish_frame(mf)?;
Ok(())
}
fn publish_config(
session: &crate::inbound::PublishSession,
codec: CodecId,
annexb: bytes::Bytes,
) -> Result<()> {
let mut cfg = MediaFrame::new_video(0, 0, annexb, codec, true);
cfg.flags |= crate::FrameFlags::CONFIG;
session.publish_frame(cfg)?;
Ok(())
}
enum VideoDepack {
H264(H264Depacketizer),
H265(H265Depacketizer),
}
impl VideoDepack {
fn new(codec: CodecId) -> Self {
match codec {
CodecId::H265 => VideoDepack::H265(H265Depacketizer::new()),
_ => VideoDepack::H264(H264Depacketizer::new()),
}
}
fn codec(&self) -> CodecId {
match self {
VideoDepack::H264(_) => CodecId::H264,
VideoDepack::H265(_) => CodecId::H265,
}
}
fn push(
&mut self,
payload: &[u8],
marker: bool,
timestamp: u32,
sequence: u16,
) -> std::result::Result<Option<AccessUnit>, DepacketizeError> {
match self {
VideoDepack::H264(d) => d.push(payload, marker, timestamp, sequence),
VideoDepack::H265(d) => d.push(payload, marker, timestamp, sequence),
}
}
}
fn video_codec_from_sdp(body: &[u8]) -> CodecId {
let sdp = Sdp::parse(&String::from_utf8_lossy(body));
let enc = sdp
.media
.iter()
.find(|m| m.media == "video")
.and_then(|m| m.encoding.as_deref())
.map(|e| e.to_ascii_uppercase());
match enc.as_deref() {
Some("H265") | Some("HEVC") => CodecId::H265,
_ => CodecId::H264,
}
}
struct UdpMedia {
rtp: Arc<UdpSocket>,
client_rtp: SocketAddr,
server_rtp_port: u16,
}
fn parse_client_ports(transport: &str) -> Option<(u16, u16)> {
let v = transport
.split(';')
.find_map(|p| p.trim().strip_prefix("client_port="))?;
let (a, b) = v.split_once('-')?;
Some((a.trim().parse().ok()?, b.trim().parse().ok()?))
}
pub struct RtspServer {
playback: Arc<dyn PlaybackRegistry>,
ingest: IngestContext,
bind: SocketAddr,
gate: Option<EgressGate>,
}
impl RtspServer {
pub fn new(
playback: Arc<dyn PlaybackRegistry>,
ingest: IngestContext,
bind: SocketAddr,
) -> Self {
Self {
playback,
ingest,
bind,
gate: None,
}
}
pub fn with_gate(mut self, gate: EgressGate) -> Self {
self.gate = Some(gate);
self
}
pub async fn run(self, shutdown: CancellationToken) -> Result<()> {
let listener = TcpListener::bind(self.bind).await?;
info!(bind = %self.bind, "rtsp server listening");
loop {
tokio::select! {
_ = shutdown.cancelled() => break,
accepted = listener.accept() => {
let (sock, peer) = match accepted {
Ok(v) => v,
Err(e) => { warn!(error = %e, "rtsp accept failed"); continue; }
};
let playback = Arc::clone(&self.playback);
let ingest = self.ingest.clone();
let gate = self.gate.clone();
let shutdown = shutdown.clone();
tokio::spawn(async move {
if let Err(e) = serve_connection(sock, playback, ingest, gate, shutdown).await {
debug!(%peer, error = %e, "rtsp connection ended");
}
});
}
}
}
Ok(())
}
}
fn uri_token(uri: &str) -> Option<String> {
crate::auth::token_from_query(uri)
}
const PT_H264: u8 = 96;
const RTP_CHANNEL: u8 = 0;
async fn serve_connection(
mut sock: TcpStream,
playback: Arc<dyn PlaybackRegistry>,
ingest: IngestContext,
gate: Option<EgressGate>,
shutdown: CancellationToken,
) -> Result<()> {
let mut buf = Vec::with_capacity(2048);
let session = new_session_id();
let mut announce_key: Option<StreamKey> = None;
let mut announce_codec = CodecId::H264;
let mut announce_config: Option<bytes::Bytes> = None;
let peer = sock.peer_addr().ok();
let peer_ip = peer.map(|a| a.ip());
let mut udp: Option<UdpMedia> = None;
loop {
let Some((req, body)) = read_request(&mut sock, &mut buf).await? else {
return Ok(()); };
let key = stream_key_from_uri(&req.uri);
match req.method.as_str() {
"OPTIONS" => {
sock.write_all(options_response(req.cseq).as_bytes())
.await?
}
"ANNOUNCE" => match &key {
Some(k) => {
announce_key = Some(k.clone());
let sdp = Sdp::parse(&String::from_utf8_lossy(&body));
announce_codec = video_codec_from_sdp(&body);
announce_config = sdp.video_config_annexb();
sock.write_all(simple_ok(req.cseq, &session).as_bytes())
.await?;
}
None => sock.write_all(not_found(req.cseq).as_bytes()).await?,
},
"RECORD" => match announce_key.take().or_else(|| key.clone()) {
Some(k) => {
sock.write_all(simple_ok(req.cseq, &session).as_bytes())
.await?;
let mut creds = Credentials::default();
creds.params.push(("proto".into(), "rtsp".into()));
creds.token = uri_token(&req.uri);
creds.addr = peer;
let sess = ingest.open_publish_checked(k.clone(), &creds).await?;
let config = announce_config.take();
if let Some(cfg) = &config {
publish_config(&sess, announce_codec, cfg.clone())?;
}
match udp.take() {
Some(u) => {
info!(stream = %k, codec = ?announce_codec, "rtsp ingest (RECORD/UDP) started");
return record_udp(sock, sess, u, announce_codec, config, shutdown)
.await;
}
None => {
info!(stream = %k, codec = ?announce_codec, "rtsp ingest (RECORD/TCP) started");
return record(sock, sess, buf, announce_codec, config, shutdown).await;
}
}
}
None => sock.write_all(not_found(req.cseq).as_bytes()).await?,
},
"DESCRIBE" => match key.as_ref().and_then(|k| playback.get_stream(k).ok()) {
Some(handle) => {
let codec = stream_video_codec(&handle);
let sdp = build_sdp(codec);
sock.write_all(describe_response(req.cseq, &req.uri, &sdp).as_bytes())
.await?;
}
None => sock.write_all(not_found(req.cseq).as_bytes()).await?,
},
"SETUP" => {
let transport = req
.headers
.iter()
.find(|(n, _)| n.eq_ignore_ascii_case("transport"))
.map(|(_, v)| v.as_str())
.unwrap_or("");
let is_tcp = transport.is_empty()
|| transport.contains("TCP")
|| transport.contains("interleaved");
if is_tcp {
udp = None;
sock.write_all(setup_response(req.cseq, &session, transport).as_bytes())
.await?;
} else if let Some((cl_rtp, cl_rtcp)) = parse_client_ports(transport) {
match (peer_ip, UdpSocket::bind("0.0.0.0:0").await.ok()) {
(Some(ip), Some(rtp_sock)) => {
let server_rtp_port =
rtp_sock.local_addr().map(|a| a.port()).unwrap_or(0);
udp = Some(UdpMedia {
rtp: Arc::new(rtp_sock),
client_rtp: SocketAddr::new(ip, cl_rtp),
server_rtp_port,
});
let record = transport.to_ascii_lowercase().contains("mode=record");
sock.write_all(
setup_response_udp(
req.cseq,
&session,
cl_rtp,
cl_rtcp,
server_rtp_port,
server_rtp_port.wrapping_add(1),
record,
)
.as_bytes(),
)
.await?;
}
_ => {
sock.write_all(unsupported_transport(req.cseq).as_bytes())
.await?;
}
}
} else {
sock.write_all(unsupported_transport(req.cseq).as_bytes())
.await?;
}
}
"PLAY" => {
let live = key.filter(|k| playback.get_stream(k).is_ok());
match live {
Some(key) => {
let allowed = match gate.as_ref() {
Some(g) => g(key.clone(), uri_token(&req.uri), peer).await,
None => true,
};
if !allowed {
sock.write_all(unauthorized(req.cseq).as_bytes()).await?;
continue;
}
sock.write_all(play_response(req.cseq, &session).as_bytes())
.await?;
return match udp.take() {
Some(u) => play_udp(sock, &playback, key, u, shutdown).await,
None => play(sock, &playback, key, shutdown).await,
};
}
None => sock.write_all(not_found(req.cseq).as_bytes()).await?,
}
}
"TEARDOWN" => {
sock.write_all(simple_ok(req.cseq, &session).as_bytes())
.await?;
return Ok(());
}
other => {
debug!(method = other, "rtsp: unsupported method");
sock.write_all(not_implemented(req.cseq).as_bytes()).await?;
}
}
}
}
async fn record(
mut sock: TcpStream,
session: crate::inbound::PublishSession,
mut carry: Vec<u8>,
codec: CodecId,
config: Option<bytes::Bytes>,
shutdown: CancellationToken,
) -> Result<()> {
let mut depack = VideoDepack::new(codec);
let mut fails = 0u32;
drain_record(&mut carry, &mut depack, &session, &mut fails, &config)?;
let mut read = [0u8; 16 * 1024];
loop {
tokio::select! {
_ = shutdown.cancelled() => break,
n = sock.read(&mut read) => {
let n = n?;
if n == 0 { break; }
carry.extend_from_slice(&read[..n]);
drain_record(&mut carry, &mut depack, &session, &mut fails, &config)?;
}
}
}
session.finish().await
}
async fn record_udp(
mut sock: TcpStream,
session: crate::inbound::PublishSession,
udp: UdpMedia,
codec: CodecId,
config: Option<bytes::Bytes>,
shutdown: CancellationToken,
) -> Result<()> {
let mut depack = VideoDepack::new(codec);
let mut dgram = vec![0u8; 64 * 1024];
let mut ctl = [0u8; 4096];
let mut fails = 0u32;
loop {
tokio::select! {
_ = shutdown.cancelled() => break,
n = sock.read(&mut ctl) => {
if n.unwrap_or(0) == 0 { break; }
}
r = udp.rtp.recv_from(&mut dgram) => {
let Ok((n, _from)) = r else { continue };
let Some(header) = RtpHeader::parse(&dgram[..n]) else { continue };
let payload = &dgram[header.payload_offset..n];
match depack.push(payload, header.marker, header.timestamp, header.sequence) {
Ok(Some(au)) => {
fails = 0;
publish_au(&session, codec, au, &config)?;
}
Ok(None) => fails = 0,
Err(_) => {
fails += 1;
if fails >= MAX_CONSECUTIVE_DEPACK_FAILURES {
warn!(?codec, "rtsp udp ingest: repeated depacketize \
failures — aborting (codec mismatch?)");
break;
}
}
}
}
}
}
session.finish().await
}
async fn play_udp(
mut sock: TcpStream,
playback: &Arc<dyn PlaybackRegistry>,
key: StreamKey,
udp: UdpMedia,
shutdown: CancellationToken,
) -> Result<()> {
let handle = playback.get_stream(&key)?;
let codec = stream_video_codec(&handle);
let mut sub = handle.subscribe_resilient();
let mut packetizer = packetizer_for(codec);
let mut pkts: Vec<Vec<u8>> = Vec::new();
let mut ctl = [0u8; 4096];
loop {
tokio::select! {
_ = shutdown.cancelled() => break,
n = sock.read(&mut ctl) => {
if n.unwrap_or(0) == 0 { break; } }
frame = sub.recv() => {
let Some(frame) = frame else { break };
if !frame.is_video() || frame.codec != codec {
continue;
}
let timestamp = (frame.pts.max(0) as u64).wrapping_mul(90) as u32;
packetizer.packetize_into(&frame.data, timestamp, &mut pkts);
for pkt in pkts.iter() {
if udp.rtp.send_to(pkt, udp.client_rtp).await.is_err() {
return Ok(());
}
}
}
}
}
let _ = udp.server_rtp_port; Ok(())
}
fn drain_record(
buf: &mut Vec<u8>,
depack: &mut VideoDepack,
session: &crate::inbound::PublishSession,
fails: &mut u32,
config: &Option<bytes::Bytes>,
) -> Result<()> {
if buf.first().is_some_and(|&b| b != b'$') {
return Err(crate::StreamError::protocol(
"rtsp ingest: non-interleaved data on control connection",
));
}
let codec = depack.codec();
let mut consumed = 0;
while let Some((frame, len)) = InterleavedFrame::parse(&buf[consumed..]) {
consumed += len;
if frame.channel != 0 {
continue;
}
let Some(header) = RtpHeader::parse(frame.payload) else {
continue;
};
let payload = &frame.payload[header.payload_offset..];
match depack.push(payload, header.marker, header.timestamp, header.sequence) {
Ok(Some(au)) => {
*fails = 0;
publish_au(session, codec, au, config)?;
}
Ok(None) => *fails = 0,
Err(_) => {
*fails += 1;
if *fails >= MAX_CONSECUTIVE_DEPACK_FAILURES {
warn!(
?codec,
"rtsp ingest: repeated depacketize failures — \
aborting (codec mismatch or corrupt RTP?)"
);
return Err(crate::StreamError::protocol(
"rtsp ingest: depacketize failure threshold",
));
}
}
}
}
buf.drain(..consumed);
Ok(())
}
async fn play(
mut sock: TcpStream,
playback: &Arc<dyn PlaybackRegistry>,
key: StreamKey,
shutdown: CancellationToken,
) -> Result<()> {
let handle = playback.get_stream(&key)?;
let codec = stream_video_codec(&handle);
let packetizer = packetizer_for(codec);
let mut sink = RtspSink {
sock: &mut sock,
packetizer,
codec,
pkts: Vec::new(),
wbuf: Vec::with_capacity(1500),
};
handle.drive_to(&shutdown, &mut sink).await
}
struct RtspSink<'a> {
sock: &'a mut TcpStream,
packetizer: RtpPacketizer,
codec: CodecId,
pkts: Vec<Vec<u8>>,
wbuf: Vec<u8>,
}
#[async_trait::async_trait]
impl crate::bus::FrameSink for RtspSink<'_> {
async fn send(&mut self, frame: Arc<MediaFrame>) -> Result<ControlFlow<()>> {
match send_frame(
self.sock,
&mut self.packetizer,
self.codec,
&mut self.pkts,
&mut self.wbuf,
&frame,
)
.await
{
Ok(()) => Ok(ControlFlow::Continue(())),
Err(_) => Ok(ControlFlow::Break(())),
}
}
}
async fn send_frame(
sock: &mut TcpStream,
packetizer: &mut RtpPacketizer,
codec: CodecId,
pkts: &mut Vec<Vec<u8>>,
wbuf: &mut Vec<u8>,
frame: &MediaFrame,
) -> std::io::Result<()> {
if !frame.is_video() || frame.codec != codec {
return Ok(()); }
let timestamp = (frame.pts.max(0) as u64).wrapping_mul(90) as u32; packetizer.packetize_into(&frame.data, timestamp, pkts);
wbuf.clear();
for pkt in pkts.iter() {
frame_interleaved(RTP_CHANNEL, pkt, wbuf);
}
sock.write_all(wbuf).await
}
fn frame_interleaved(channel: u8, rtp: &[u8], out: &mut Vec<u8>) {
out.push(b'$');
out.push(channel);
out.extend_from_slice(&(rtp.len() as u16).to_be_bytes());
out.extend_from_slice(rtp);
}
async fn read_request(
sock: &mut TcpStream,
buf: &mut Vec<u8>,
) -> Result<Option<(RtspRequest, Vec<u8>)>> {
let mut tmp = [0u8; 1024];
loop {
if let Some(end) = find_double_crlf(buf) {
let head = String::from_utf8_lossy(&buf[..end]).into_owned();
let Some(req) = RtspRequest::parse(&head) else {
buf.drain(..end);
return Ok(None);
};
let want = content_length(&req);
while buf.len() < end + want {
let n = sock.read(&mut tmp).await?;
if n == 0 {
return Ok(None);
}
buf.extend_from_slice(&tmp[..n]);
}
let body = buf[end..end + want].to_vec();
buf.drain(..end + want);
return Ok(Some((req, body)));
}
let n = sock.read(&mut tmp).await?;
if n == 0 {
return Ok(None);
}
buf.extend_from_slice(&tmp[..n]);
if buf.len() > 64 * 1024 {
return Err(crate::StreamError::protocol("rtsp request too large"));
}
}
}
fn content_length(req: &RtspRequest) -> usize {
req.headers
.iter()
.find(|(n, _)| n.eq_ignore_ascii_case("content-length"))
.and_then(|(_, v)| v.parse().ok())
.unwrap_or(0)
}
fn find_double_crlf(buf: &[u8]) -> Option<usize> {
buf.windows(4).position(|w| w == b"\r\n\r\n").map(|p| p + 4)
}
fn stream_key_from_uri(uri: &str) -> Option<StreamKey> {
let rest = uri.strip_prefix("rtsp://")?;
let path = rest.split_once('/').map(|(_, p)| p)?;
let path = path.split(['?', ';']).next().unwrap_or(path);
let mut segs = path.split('/').filter(|s| !s.is_empty());
let app = segs.next()?;
let stream = segs.next()?;
Some(StreamKey::new(app, stream))
}
fn options_response(cseq: u32) -> String {
format!(
"RTSP/1.0 200 OK\r\nCSeq: {cseq}\r\n\
Public: OPTIONS, DESCRIBE, SETUP, PLAY, TEARDOWN\r\n\r\n"
)
}
fn build_sdp(codec: CodecId) -> String {
let enc = if codec == CodecId::H265 {
"H265"
} else {
"H264"
};
format!(
"v=0\r\n\
o=- 0 0 IN IP4 0.0.0.0\r\n\
s=arcly-stream\r\n\
t=0 0\r\n\
m=video 0 RTP/AVP 96\r\n\
a=rtpmap:96 {enc}/90000\r\n\
a=control:streamid=0\r\n"
)
}
fn stream_video_codec(handle: &crate::bus::StreamHandle) -> CodecId {
handle
.cached_configs()
.0
.map(|f| f.codec)
.unwrap_or(CodecId::H264)
}
fn packetizer_for(codec: CodecId) -> RtpPacketizer {
if codec == CodecId::H265 {
RtpPacketizer::new_h265(PT_H264, 0x5254_5350, 1400)
} else {
RtpPacketizer::new(PT_H264, 0x5254_5350, 1400)
}
}
fn describe_response(cseq: u32, uri: &str, sdp: &str) -> String {
format!(
"RTSP/1.0 200 OK\r\nCSeq: {cseq}\r\n\
Content-Base: {uri}\r\nContent-Type: application/sdp\r\n\
Content-Length: {}\r\n\r\n{sdp}",
sdp.len()
)
}
fn setup_response(cseq: u32, session: &str, client_transport: &str) -> String {
let interleaved = client_transport
.split(';')
.find_map(|p| p.trim().strip_prefix("interleaved="))
.unwrap_or("0-1");
let mode =
if client_transport.contains("mode=record") || client_transport.contains("mode=RECORD") {
";mode=record"
} else {
""
};
format!(
"RTSP/1.0 200 OK\r\nCSeq: {cseq}\r\n\
Transport: RTP/AVP/TCP;unicast;interleaved={interleaved}{mode}\r\n\
Session: {session}\r\n\r\n"
)
}
#[allow(clippy::too_many_arguments)]
fn setup_response_udp(
cseq: u32,
session: &str,
client_rtp: u16,
client_rtcp: u16,
server_rtp: u16,
server_rtcp: u16,
record: bool,
) -> String {
let mode = if record { ";mode=record" } else { "" };
format!(
"RTSP/1.0 200 OK\r\nCSeq: {cseq}\r\n\
Transport: RTP/AVP;unicast;client_port={client_rtp}-{client_rtcp};\
server_port={server_rtp}-{server_rtcp}{mode}\r\n\
Session: {session}\r\n\r\n"
)
}
fn unsupported_transport(cseq: u32) -> String {
format!("RTSP/1.0 461 Unsupported Transport\r\nCSeq: {cseq}\r\n\r\n")
}
fn play_response(cseq: u32, session: &str) -> String {
format!("RTSP/1.0 200 OK\r\nCSeq: {cseq}\r\nSession: {session}\r\nRange: npt=0.000-\r\n\r\n")
}
fn simple_ok(cseq: u32, session: &str) -> String {
format!("RTSP/1.0 200 OK\r\nCSeq: {cseq}\r\nSession: {session}\r\n\r\n")
}
fn not_implemented(cseq: u32) -> String {
format!("RTSP/1.0 501 Not Implemented\r\nCSeq: {cseq}\r\n\r\n")
}
fn not_found(cseq: u32) -> String {
format!("RTSP/1.0 404 Not Found\r\nCSeq: {cseq}\r\n\r\n")
}
fn unauthorized(cseq: u32) -> String {
format!("RTSP/1.0 401 Unauthorized\r\nCSeq: {cseq}\r\n\r\n")
}
fn new_session_id() -> String {
use std::collections::hash_map::RandomState;
use std::hash::{BuildHasher, Hasher};
use std::time::{SystemTime, UNIX_EPOCH};
let mut h = RandomState::new().build_hasher();
h.write_u128(
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_nanos())
.unwrap_or(0),
);
format!("{:016X}", h.finish())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::protocol::rtsp::message::{InterleavedFrame, RtspResponse};
#[test]
fn interleaved_frame_round_trips() {
let rtp = [0x80u8, 0xE0, 0x00, 0x01, 0xAA, 0xBB];
let mut framed = Vec::new();
frame_interleaved(0, &rtp, &mut framed);
assert_eq!(framed[0], b'$');
let (f, used) = InterleavedFrame::parse(&framed).expect("parse");
assert_eq!(used, framed.len());
assert_eq!(f.channel, 0);
assert_eq!(f.payload, &rtp);
}
#[test]
fn stream_key_parsed_from_uri_variants() {
let k = stream_key_from_uri("rtsp://host:554/live/cam").unwrap();
assert_eq!((k.app.as_str(), k.stream_id.as_str()), ("live", "cam"));
let k = stream_key_from_uri("rtsp://h/live/cam/streamid=0?x=1").unwrap();
assert_eq!((k.app.as_str(), k.stream_id.as_str()), ("live", "cam"));
assert!(stream_key_from_uri("rtsp://host/onlyapp").is_none());
}
#[test]
fn uri_token_extracted() {
assert_eq!(
uri_token("rtsp://h/live/cam?token=abc").as_deref(),
Some("abc")
);
assert_eq!(uri_token("rtsp://h/live/cam").as_deref(), None);
}
#[test]
fn announce_request_parses_with_body_len() {
let req = RtspRequest::parse(
"ANNOUNCE rtsp://h/live/cam RTSP/1.0\r\nCSeq: 2\r\nContent-Length: 5\r\n",
)
.unwrap();
assert_eq!(req.method, "ANNOUNCE");
assert_eq!(content_length(&req), 5);
assert_eq!(
stream_key_from_uri(&req.uri).map(|k| (k.app.to_string(), k.stream_id.to_string())),
Some(("live".into(), "cam".into()))
);
}
#[test]
fn responses_are_well_formed() {
assert!(options_response(2).contains("Public: OPTIONS"));
let d = describe_response(3, "rtsp://h/live/cam", &build_sdp(CodecId::H264));
assert!(build_sdp(CodecId::H265).contains("H265/90000"));
let parsed = RtspResponse::parse(
d.split("\r\n\r\n").next().unwrap(),
d.split("\r\n\r\n").nth(1).unwrap_or("").to_string(),
)
.expect("response parses");
assert_eq!(parsed.header("Content-Type"), Some("application/sdp"));
let setup = setup_response(4, "DEADBEEF", "RTP/AVP/TCP;unicast;interleaved=0-1");
assert!(setup.contains("interleaved=0-1"));
assert!(setup.contains("Session: DEADBEEF"));
let rec = setup_response(5, "S", "RTP/AVP/TCP;unicast;interleaved=0-1;mode=record");
assert!(rec.contains("mode=record"));
assert!(unsupported_transport(6).starts_with("RTSP/1.0 461"));
assert!(not_found(7).starts_with("RTSP/1.0 404"));
}
}