Code/Resolver/Live/RtspTunnel.cs
using System.Globalization;
using System.IO;
using System.Net.Http;
using System.Text;
using System.Threading;
namespace Bimp.Resolver.Live;
/// <summary>
/// RTSP tunnelled over HTTP (the QuickTime / Apple scheme most IP cameras support) - the only way to reach an RTSP
/// camera from sandboxed code, which gets no TCP or UDP sockets.
/// <para>
/// One long-lived HTTP GET carries everything from the camera (RTSP responses and the RTP media, interleaved as
/// in RTSP-over-TCP); each RTSP request goes up as a separate HTTP POST whose body is the request in base64.
/// Both carry the same x-sessioncookie. rtsp://, rtspt:// and rtsps:// all map to this.
/// </para>
/// </summary>
public sealed class RtspTunnel
{
readonly Uri rtspUrl;
readonly LiveSegmenter sink;
readonly string user, password;
readonly string cookie = Guid.NewGuid().ToString( "N" )[..22];
string httpUrl;
int cseq;
/// <summary>
/// Requests go up one long POST, as in Apple's scheme (VRCDN only reads the first POST of a tunnel), or a POST
/// each (some cameras only take that). <see cref="LiveStream"/> tries the other on a reconnect.
/// </summary>
public bool LongPost { get; }
/// <summary> The session got as far as PLAY, so this way of tunnelling works with the server. </summary>
public bool Played { get; private set; }
bool answered;
UpstreamStream upstream;
string session;
int sessionTimeout = 60;
string contentBase;
(string realm, string nonce, string qop, bool digest)? auth;
int nonceCount;
readonly Dictionary<int, TaskCompletionSource<Response>> waiting = new();
readonly Dictionary<int, Track> channels = new();
readonly List<Track> tracks = new();
readonly System.Diagnostics.Stopwatch clock = System.Diagnostics.Stopwatch.StartNew();
sealed class Response
{
public int Status;
public string Reason;
public Dictionary<string, string> Headers = new( StringComparer.OrdinalIgnoreCase );
public string Body = "";
}
public RtspTunnel( Uri url, LiveSegmenter sink, bool longPost = true )
{
LongPost = longPost;
rtspUrl = url;
this.sink = sink;
if ( !string.IsNullOrEmpty( url.UserInfo ) )
{
var parts = url.UserInfo.Split( ':', 2 );
user = Uri.UnescapeDataString( parts[0] );
password = parts.Length > 1 ? Uri.UnescapeDataString( parts[1] ) : "";
}
}
/// <summary> The rtsp url as sent in requests: rtsp://, no credentials. </summary>
string RequestUrl => $"rtsp://{rtspUrl.Host}{(rtspUrl.IsDefaultPort || rtspUrl.Port <= 0 ? "" : $":{rtspUrl.Port}")}{rtspUrl.PathAndQuery}";
/// <summary> The last failure is worth reconnecting for (a lost answer or a dropped tunnel, not a refusal). </summary>
public bool Retryable { get; private set; }
Action<int> received;
public async Task Run( CancellationToken ct, Action<string> describe, Action<int> bytesReceived = null )
{
// The tunnel is HTTP on the camera's web port, or on the RTSP port itself (servers like VRCDN tell RTSP
// and HTTP apart on 554): the port in the link if there is one, else 554 then 80 (443 for rtsps).
// A link with the RTSP port (554) that refuses HTTP is retried on 80.
var secure = rtspUrl.Scheme.Equals( "rtsps", StringComparison.OrdinalIgnoreCase );
var explicitPort = rtspUrl.IsDefaultPort || rtspUrl.Port <= 0 ? -1 : rtspUrl.Port;
var ports = explicitPort > 0 ? (explicitPort == 554 ? new[] { 554, 80 } : new[] { explicitPort }) : secure ? new[] { 443 } : new[] { 554, 80 };
Stream stream = null;
Exception last = null;
foreach ( var port in ports )
{
httpUrl = $"{(secure ? "https" : "http")}://{rtspUrl.Host}:{port}{rtspUrl.PathAndQuery}";
// a port nothing listens on can hang until the TCP timeout - give up on it sooner
using var connect = CancellationTokenSource.CreateLinkedTokenSource( ct );
connect.CancelAfter( ports.Length > 1 ? 6000 : 20000 );
try
{
stream = await Http.RequestStreamAsync( httpUrl, headers: new()
{
["x-sessioncookie"] = cookie,
["Accept"] = "application/x-rtsp-tunnelled",
["Pragma"] = "no-cache",
["Cache-Control"] = "no-cache",
}, cancellationToken: connect.Token );
break;
}
catch ( InvalidOperationException ) { throw; } // not allowed (private network / raw IP)
catch ( Exception e ) when ( !ct.IsCancellationRequested )
{
last = e;
}
}
if ( stream is null )
throw new ResolveException( $"Couldn't open an RTSP-over-HTTP tunnel to {rtspUrl.Host} ({last?.Message}). The camera must support RTSP over HTTP tunnelling - put its HTTP port in the link." );
using ( stream )
{
received = bytesReceived;
// kept rather than left unobserved when a request fails first (the tunnel closing then faults the read)
var reader = Observed( ReadLoop( stream, ct ) );
var describeResponse = await Request( "DESCRIBE", RequestUrl, new() { ["Accept"] = "application/sdp" }, ct );
if ( describeResponse.Status != 200 ) throw new ResolveException( $"The camera refused DESCRIBE: {describeResponse.Status} {describeResponse.Reason}" );
// some servers (VRCDN) hand out the session at DESCRIBE already
if ( describeResponse.Headers.TryGetValue( "Session", out var describeSession ) ) session = describeSession.Split( ';' )[0].Trim();
contentBase = describeResponse.Headers.GetValueOrDefault( "Content-Base" ) ?? describeResponse.Headers.GetValueOrDefault( "Content-Location" ) ?? RequestUrl;
var sdp = Sdp.Parse( describeResponse.Body );
SetupTracks( sdp );
if ( tracks.Count == 0 ) throw new ResolveException( "The camera offers no stream that can be played (needs H.264, AV1 or MJPEG video, or AAC / G.711 audio)." );
var channel = 0;
foreach ( var track in tracks )
{
var headers = new Dictionary<string, string> { ["Transport"] = $"RTP/AVP/TCP;unicast;interleaved={channel}-{channel + 1}" };
var r = await Request( "SETUP", ControlUrl( sdp.Control, track.Media.Control ), headers, ct );
if ( r.Status != 200 ) throw new ResolveException( $"The camera refused SETUP: {r.Status} {r.Reason}" );
if ( r.Headers.TryGetValue( "Session", out var s ) )
{
var parts = s.Split( ';' );
session = parts[0].Trim();
foreach ( var p in parts.Skip( 1 ) )
if ( p.Trim().StartsWith( "timeout=" ) && int.TryParse( p.Trim()[8..], out var t ) ) sessionTimeout = t;
}
// the camera may pick other channels
var transport = r.Headers.GetValueOrDefault( "Transport" ) ?? "";
var m = System.Text.RegularExpressions.Regex.Match( transport, "interleaved=(\\d+)" );
var used = m.Success ? int.Parse( m.Groups[1].Value ) : channel;
channels[used] = track;
channel = used + 2;
}
describe( "RTSP " + string.Join( " + ", tracks.Select( t => t.Name ) ) );
var play = await Request( "PLAY", ControlUrl( sdp.Control, null ), new() { ["Range"] = "npt=0.000-" }, ct );
if ( play.Status != 200 ) throw new ResolveException( $"The camera refused PLAY: {play.Status} {play.Reason}" );
Played = true;
_ = KeepAlive( ct );
await reader;
if ( readError is not null and not OperationCanceledException ) throw readError;
}
Retryable = true;
throw new ResolveException( "The camera closed the stream." );
}
Exception readError;
async Task Observed( Task task )
{
try { await task; }
catch ( Exception e ) { readError = e; }
}
string ControlUrl( string sessionControl, string trackControl )
{
var control = trackControl ?? sessionControl;
if ( string.IsNullOrEmpty( control ) || control == "*" ) return contentBase;
if ( control.StartsWith( "rtsp://", StringComparison.OrdinalIgnoreCase ) ) return control;
return contentBase.TrimEnd( '/' ) + "/" + control.TrimStart( '/' );
}
void SetupTracks( Sdp sdp )
{
// AV1 first: the engine's AV1 decoder presents more evenly than its H.264 one (see VideoFormat)
var av1 = sdp.Media.FirstOrDefault( m => m.Type == "video" && m.Encoding.Equals( "AV1", StringComparison.OrdinalIgnoreCase ) );
var video = sdp.Media.FirstOrDefault( m => m.Type == "video" && m.Encoding.Equals( "H264", StringComparison.OrdinalIgnoreCase ) );
var jpeg = sdp.Media.FirstOrDefault( m => m.Type == "video" && m.Encoding.Equals( "JPEG", StringComparison.OrdinalIgnoreCase ) );
if ( av1 is not null ) tracks.Add( new Av1Track( av1, sink ) );
else if ( video is not null ) tracks.Add( new H264Track( video, sink ) );
else if ( jpeg is not null ) tracks.Add( new MjpegTrack( jpeg, sink ) );
else if ( sdp.Media.FirstOrDefault( m => m.Type == "video" ) is { } other )
{
if ( other.Encoding.Equals( "H265", StringComparison.OrdinalIgnoreCase ) || other.Encoding.Equals( "HEVC", StringComparison.OrdinalIgnoreCase ) )
throw new ResolveException( "This camera's video is H.265 (HEVC), which the engine can't decode - set it to H.264, AV1 or MJPEG." );
Log.Warning( $"[bimp] rtsp: video is {other.Encoding}, only H.264, AV1 and MJPEG can be played" );
}
var audio = sdp.Media.Where( m => m.Type == "audio" ).ToList();
var aac = audio.FirstOrDefault( m => m.Encoding.Equals( "mpeg4-generic", StringComparison.OrdinalIgnoreCase ) );
var g711 = audio.FirstOrDefault( m => m.Encoding is "PCMU" or "PCMA" || m.PayloadType is 0 or 8 );
if ( aac is not null ) tracks.Add( new AacTrack( aac, sink ) );
else if ( g711 is not null ) tracks.Add( new G711Track( g711, sink ) );
else if ( audio.Count > 0 ) Log.Warning( $"[bimp] rtsp: audio is {audio[0].Encoding}, which isn't supported - video only" );
}
async Task KeepAlive( CancellationToken ct )
{
var interval = Math.Clamp( sessionTimeout / 2, 10, 60 );
while ( !ct.IsCancellationRequested )
{
await Task.Delay( interval * 1000, ct );
try { await Request( "OPTIONS", RequestUrl, null, ct ); }
catch ( Exception e ) when ( e is not OperationCanceledException ) { Log.Trace( $"[bimp] rtsp keepalive: {e.Message}" ); }
}
}
/// <summary>
/// Send an RTSP request up a POST and wait for its response on the GET. Retries once with credentials on 401.
/// </summary>
async Task<Response> Request( string method, string url, Dictionary<string, string> headers, CancellationToken ct, bool retried = false )
{
var id = ++cseq;
var sb = new StringBuilder();
sb.Append( $"{method} {url} RTSP/1.0\r\nCSeq: {id}\r\n" );
if ( session is not null ) sb.Append( $"Session: {session}\r\n" );
if ( Authorization( method, url ) is { } a ) sb.Append( $"Authorization: {a}\r\n" );
foreach ( var (k, v) in headers ?? new() ) sb.Append( $"{k}: {v}\r\n" );
// padded to a multiple of 3 bytes, so no base64 '=' lands in the middle of the long POST's body
var userAgent = "User-Agent: bimp";
while ( (sb.Length + userAgent.Length + 4) % 3 != 0 ) userAgent += " ";
sb.Append( userAgent + "\r\n\r\n" );
var base64 = Encoding.ASCII.GetBytes( Convert.ToBase64String( Encoding.ASCII.GetBytes( sb.ToString() ) ) );
var tcs = new TaskCompletionSource<Response>( TaskCreationOptions.RunContinuationsAsynchronously );
waiting[id] = tcs;
// The server never answers the POST itself (answers come down the GET), so don't wait for it. A POST of
// its own is dropped once the answer came down the GET.
using var postCts = CancellationTokenSource.CreateLinkedTokenSource( ct );
var persistent = LongPost;
if ( persistent )
{
// its declared length runs out after an hour or so of keep-alives: start another
var padded = UpstreamStream.Padded( base64 );
if ( upstream is not null && !upstream.Fits( padded.Length ) )
{
upstream.Complete();
upstream = null;
}
if ( upstream is null )
{
upstream = new UpstreamStream();
var longPost = new StreamContent( upstream, 64 * 1024 );
longPost.Headers.ContentLength = UpstreamStream.DeclaredLength;
longPost.Headers.TryAddWithoutValidation( "Content-Type", "application/x-rtsp-tunnelled" );
_ = Post( longPost, ct );
}
upstream.Send( padded );
}
else
{
var content = new ByteArrayContent( base64 );
content.Headers.TryAddWithoutValidation( "Content-Type", "application/x-rtsp-tunnelled" );
_ = Post( content, postCts.Token );
}
Response response;
// the first answer shows whether this way of tunnelling works - don't wait the full 10 s for it
using ( var timeout = new CancellationTokenSource( !answered ? 4000 : 10000 ) )
using ( timeout.Token.Register( () => tcs.TrySetCanceled() ) )
using ( ct.Register( () => tcs.TrySetCanceled() ) )
{
try
{
response = await tcs.Task;
}
catch ( OperationCanceledException )
{
ct.ThrowIfCancellationRequested();
Retryable = true;
throw new ResolveException( $"The camera didn't answer {method} (is RTSP over HTTP enabled on it?)." );
}
finally
{
waiting.Remove( id );
postCts.Cancel();
}
}
answered = true;
if ( response.Status == 401 && !retried && user is not null && response.Headers.TryGetValue( "WWW-Authenticate", out var challenge ) )
{
ParseChallenge( challenge );
return await Request( method, url, headers, ct, true );
}
if ( response.Status == 401 ) throw new ResolveException( user is null ? "The camera needs a login - put it in the link (rtsp://user:password@camera/...)." : "The camera rejected the login." );
return response;
}
async Task Post( HttpContent content, CancellationToken ct )
{
try
{
using var response = await Http.RequestAsync( httpUrl, "POST", content, new() { ["x-sessioncookie"] = cookie, ["Pragma"] = "no-cache", ["Cache-Control"] = "no-cache" }, ct );
}
catch ( Exception e ) when ( e is OperationCanceledException or HttpRequestException )
{
// expected: dropped once the answer came down the GET
}
finally
{
content.Dispose();
}
}
void ParseChallenge( string header )
{
// may hold several challenges ("Digest ..., Basic ...") joined - prefer Digest
var digest = header.Contains( "Digest", StringComparison.OrdinalIgnoreCase );
string Field( string name )
{
var m = System.Text.RegularExpressions.Regex.Match( header, name + "=\"?([^\",]*)\"?", System.Text.RegularExpressions.RegexOptions.IgnoreCase );
return m.Success ? m.Groups[1].Value : null;
}
auth = (Field( "realm" ) ?? "", Field( "nonce" ) ?? "", Field( "qop" ), digest);
nonceCount = 0;
}
string Authorization( string method, string uri )
{
if ( auth is not { } a || user is null ) return null;
if ( !a.digest ) return "Basic " + Convert.ToBase64String( Encoding.UTF8.GetBytes( $"{user}:{password}" ) );
static string Md5( string s ) => Convert.ToHexString( System.Security.Cryptography.MD5.HashData( Encoding.UTF8.GetBytes( s ) ) ).ToLowerInvariant();
var ha1 = Md5( $"{user}:{a.realm}:{password}" );
var ha2 = Md5( $"{method}:{uri}" );
if ( a.qop is not null && a.qop.Split( ',' ).Any( q => q.Trim() == "auth" ) )
{
var nc = (++nonceCount).ToString( "x8" );
var cnonce = Guid.NewGuid().ToString( "N" )[..16];
var response = Md5( $"{ha1}:{a.nonce}:{nc}:{cnonce}:auth:{ha2}" );
return $"Digest username=\"{user}\", realm=\"{a.realm}\", nonce=\"{a.nonce}\", uri=\"{uri}\", response=\"{response}\", qop=auth, nc={nc}, cnonce=\"{cnonce}\"";
}
return $"Digest username=\"{user}\", realm=\"{a.realm}\", nonce=\"{a.nonce}\", uri=\"{uri}\", response=\"{Md5( $"{ha1}:{a.nonce}:{ha2}" )}\"";
}
/// <summary>
/// Everything from the camera: RTSP responses and '$'-framed interleaved RTP/RTCP packets.
/// </summary>
async Task ReadLoop( Stream stream, CancellationToken ct )
{
var buf = new byte[256 * 1024];
int start = 0, end = 0;
while ( !ct.IsCancellationRequested )
{
// make room
if ( end == buf.Length )
{
if ( start > 0 ) { Buffer.BlockCopy( buf, start, buf, 0, end - start ); end -= start; start = 0; }
else Array.Resize( ref buf, buf.Length * 2 );
}
var n = await stream.ReadAsync( buf.AsMemory( end ), ct );
if ( n <= 0 ) return;
received?.Invoke( n );
end += n;
while ( end - start > 0 )
{
if ( buf[start] == (byte)'$' )
{
if ( end - start < 4 ) break;
var ch = buf[start + 1];
var len = (buf[start + 2] << 8) | buf[start + 3];
if ( end - start < 4 + len ) break;
var packet = buf.AsSpan( start + 4, len );
if ( (ch & 1) == 0 )
{
if ( TrackFor( ch, packet ) is { } track ) track.OnRtp( packet, clock.Elapsed.TotalSeconds );
}
else if ( channels.TryGetValue( ch - 1, out var rtcpTrack ) ) rtcpTrack.OnRtcp( packet );
start += 4 + len;
continue;
}
// text: a response (or a request from the server) up to the blank line, then its body
var headerEnd = IndexOf( buf, start, end, "\r\n\r\n"u8 );
if ( headerEnd < 0 )
{
if ( end - start > 64 * 1024 ) start = end; // garbage
break;
}
var headerText = Encoding.ASCII.GetString( buf, start, headerEnd - start );
var lines = headerText.Split( "\r\n" );
var response = new Response();
foreach ( var line in lines.Skip( 1 ) )
{
var colon = line.IndexOf( ':' );
if ( colon > 0 ) response.Headers[line[..colon].Trim()] = line[(colon + 1)..].Trim();
}
var bodyLength = int.TryParse( response.Headers.GetValueOrDefault( "Content-Length" ), out var cl ) ? cl : 0;
if ( end - (headerEnd + 4) < bodyLength ) break;
response.Body = Encoding.UTF8.GetString( buf, headerEnd + 4, bodyLength );
start = headerEnd + 4 + bodyLength;
if ( !lines[0].StartsWith( "RTSP/" ) ) continue; // a request from the server - ignore
var status = lines[0].Split( ' ', 3 );
response.Status = status.Length > 1 && int.TryParse( status[1], out var st ) ? st : 0;
response.Reason = status.Length > 2 ? status[2] : "";
// Line the tracks up before any of their packets are handled (they follow right after)
if ( response.Headers.TryGetValue( "RTP-Info", out var rtpInfo ) ) ApplyRtpInfo( rtpInfo );
if ( int.TryParse( response.Headers.GetValueOrDefault( "CSeq" ), out var id ) && waiting.TryGetValue( id, out var tcs ) )
tcs.TrySetResult( response );
}
if ( start == end ) start = end = 0;
}
}
readonly HashSet<int> checkedChannels = new();
/// <summary>
/// The track an RTP channel carries. Checked once per channel against the packet's payload type: VRCDN's SETUP
/// answers name each track's channels swapped (video announced on 2-3 arrives on 0-1).
/// </summary>
Track TrackFor( int ch, ReadOnlySpan<byte> packet )
{
channels.TryGetValue( ch, out var track );
if ( packet.Length < 2 || !checkedChannels.Add( ch ) ) return track;
var pt = packet[1] & 0x7F;
if ( track is not null && track.Media.PayloadType == pt ) return track;
var byType = tracks.Where( t => t.Media.PayloadType == pt ).ToList();
if ( byType.Count != 1 ) return track;
// swap the channel pairs round: whatever was on this one moves to the byType track's old one
var other = channels.FirstOrDefault( kv => kv.Value == byType[0] && kv.Key != ch ).Key;
if ( track is not null && channels.ContainsKey( other ) ) channels[other] = track;
channels[ch] = byType[0];
Log.Info( $"[bimp] rtsp: channel {ch} carries {byType[0].Name} (payload type {pt}), not what SETUP said" );
return byType[0];
}
/// <summary> RTP-Info: url=...;seq=...;rtptime=..., one per track - the RTP time of the play start. </summary>
void ApplyRtpInfo( string header )
{
foreach ( var entry in header.Split( ',' ) )
{
string url = null;
long? rtptime = null;
foreach ( var part in entry.Split( ';' ) )
{
var kv = part.Trim().Split( '=', 2 );
if ( kv.Length != 2 ) continue;
if ( kv[0] == "url" ) url = kv[1];
else if ( kv[0] == "rtptime" && long.TryParse( kv[1], out var t ) ) rtptime = t;
}
if ( rtptime is null ) continue;
var track = tracks.FirstOrDefault( t => url is not null && t.Media.Control is not null && (url.EndsWith( t.Media.Control ) || url == t.Media.Control) )
?? (tracks.Count == 1 ? tracks[0] : null);
track?.SetOrigin( rtptime.Value );
}
}
static int IndexOf( byte[] buf, int start, int end, ReadOnlySpan<byte> pattern )
{
var i = buf.AsSpan( start, end - start ).IndexOf( pattern );
return i < 0 ? -1 : start + i;
}
/// <summary>
/// The body of the long POST: requests as they're queued. The declared length is just large, so the connection
/// stays open for the session.
/// <para>
/// HttpClient holds back body writes that fit its 4 KB connection buffer until the body ends - which a tunnel's
/// never does - and a custom HttpContent that flushes needs a type the sandbox doesn't allow. A single write of
/// over 8 KB goes straight to the socket, so each request goes up as one write, after enough CRLFs (which RTSP
/// servers skip between messages) to make it that big.
/// </para>
/// </summary>
sealed class UpstreamStream : Stream
{
public const int DeclaredLength = 100_000_000;
/// <summary> base64 of 9216 bytes of CRLF (a multiple of 3, so it decodes on its own). </summary>
static readonly byte[] Padding = Encoding.ASCII.GetBytes( Convert.ToBase64String( Encoding.ASCII.GetBytes( string.Concat( Enumerable.Repeat( "\r\n", 4608 ) ) ) ) );
public static byte[] Padded( byte[] request )
{
var result = new byte[Padding.Length + request.Length];
Padding.CopyTo( result, 0 );
request.CopyTo( result, Padding.Length );
return result;
}
readonly System.Threading.Channels.Channel<byte[]> queue = System.Threading.Channels.Channel.CreateUnbounded<byte[]>();
byte[] current;
int at, total, queued;
public bool Fits( int length ) => queued + length <= DeclaredLength;
public void Send( byte[] data )
{
queued += data.Length;
queue.Writer.TryWrite( data );
}
public void Complete() => queue.Writer.TryComplete();
public override async ValueTask<int> ReadAsync( Memory<byte> buffer, CancellationToken ct = default )
{
while ( current is null || at >= current.Length )
{
if ( total >= DeclaredLength || !await queue.Reader.WaitToReadAsync( ct ) ) return 0;
queue.Reader.TryRead( out current );
at = 0;
}
var n = Math.Min( Math.Min( buffer.Length, current.Length - at ), DeclaredLength - total );
current.AsMemory( at, n ).CopyTo( buffer );
at += n;
total += n;
return n;
}
public override Task<int> ReadAsync( byte[] buffer, int offset, int count, CancellationToken ct ) => ReadAsync( buffer.AsMemory( offset, count ), ct ).AsTask();
public override int Read( byte[] buffer, int offset, int count ) => throw new NotSupportedException();
public override bool CanRead => true;
public override bool CanSeek => false;
public override bool CanWrite => false;
public override long Length => throw new NotSupportedException();
public override long Position { get => total; set => throw new NotSupportedException(); }
public override void Flush() { }
public override long Seek( long offset, SeekOrigin origin ) => throw new NotSupportedException();
public override void SetLength( long value ) => throw new NotSupportedException();
public override void Write( byte[] buffer, int offset, int count ) => throw new NotSupportedException();
}
//
// SDP
//
sealed class Sdp
{
public string Control;
public readonly List<MediaDesc> Media = new();
public static Sdp Parse( string text )
{
var sdp = new Sdp();
MediaDesc current = null;
foreach ( var raw in text.Split( '\n' ) )
{
var line = raw.Trim();
if ( line.Length < 2 || line[1] != '=' ) continue;
var value = line[2..];
if ( line[0] == 'm' )
{
var p = value.Split( ' ' );
current = new MediaDesc { Type = p[0], PayloadType = p.Length > 3 && int.TryParse( p[3], out var pt ) ? pt : -1 };
// static payload types have no rtpmap
if ( current.PayloadType == 0 ) { current.Encoding = "PCMU"; current.ClockRate = 8000; }
if ( current.PayloadType == 8 ) { current.Encoding = "PCMA"; current.ClockRate = 8000; }
if ( current.PayloadType == 26 ) { current.Encoding = "JPEG"; current.ClockRate = 90000; }
sdp.Media.Add( current );
}
else if ( line[0] == 'a' )
{
if ( value.StartsWith( "control:" ) )
{
if ( current is null ) sdp.Control = value[8..];
else current.Control = value[8..];
}
else if ( value.StartsWith( "rtpmap:" ) && current is not null )
{
// rtpmap:96 H264/90000 or rtpmap:97 mpeg4-generic/48000/2
var p = value[7..].Split( ' ', 2 );
if ( p.Length == 2 && int.TryParse( p[0], out var pt ) && pt == current.PayloadType )
{
var enc = p[1].Split( '/' );
current.Encoding = enc[0];
if ( enc.Length > 1 && int.TryParse( enc[1], out var rate ) ) current.ClockRate = rate;
if ( enc.Length > 2 && int.TryParse( enc[2], out var ch ) ) current.Channels = ch;
}
}
else if ( value.StartsWith( "fmtp:" ) && current is not null )
{
var p = value[5..].Split( ' ', 2 );
if ( p.Length == 2 )
{
foreach ( var kv in p[1].Split( ';' ) )
{
var pair = kv.Trim().Split( '=', 2 );
if ( pair.Length == 2 ) current.Fmtp[pair[0].Trim().ToLowerInvariant()] = pair[1].Trim();
}
}
}
}
}
return sdp;
}
}
sealed class MediaDesc
{
public string Type;
public int PayloadType;
public string Encoding = "";
public int ClockRate = 90000;
public int Channels = 1;
public string Control;
public readonly Dictionary<string, string> Fmtp = new();
}
//
// RTP depacketizing
//
abstract class Track
{
public MediaDesc Media;
protected readonly LiveSegmenter Sink;
long? origin;
double originTime;
long lastTs = -1, wraps;
protected int LastSeq = -1;
protected Track( MediaDesc media, LiveSegmenter sink )
{
Media = media;
Sink = sink;
}
public abstract string Name { get; }
/// <summary> From RTP-Info: this RTP time is stream time 0. </summary>
public void SetOrigin( long rtptime )
{
origin = rtptime;
originTime = 0;
}
/// <summary> RTP timestamp to stream time in 90 kHz units. </summary>
protected long Time( uint ts, double arrival )
{
// unwrap 32 bits
long t = ts;
if ( lastTs >= 0 && t < lastTs - 0x80000000L ) wraps++;
lastTs = t;
t += wraps << 32;
if ( origin is null )
{
// no RTP-Info: line tracks up by when their first packet arrived
origin = t;
originTime = arrival;
}
return (long)(originTime * 90000) + (t - origin.Value) * 90000 / Media.ClockRate;
}
public void OnRtp( ReadOnlySpan<byte> p, double arrival )
{
if ( p.Length < 12 || (p[0] >> 6) != 2 ) return;
var padding = (p[0] & 0x20) != 0;
var extension = (p[0] & 0x10) != 0;
var csrc = p[0] & 0x0F;
var marker = (p[1] & 0x80) != 0;
var seq = (p[2] << 8) | p[3];
var ts = (uint)((p[4] << 24) | (p[5] << 16) | (p[6] << 8) | p[7]);
var at = 12 + csrc * 4;
if ( extension )
{
if ( at + 4 > p.Length ) return;
at += 4 + ((p[at + 2] << 8) | p[at + 3]) * 4;
}
var end = p.Length - (padding ? p[^1] : 0);
if ( at >= end ) return;
var lost = LastSeq >= 0 && seq != ((LastSeq + 1) & 0xFFFF);
LastSeq = seq;
Payload( p[at..end], ts, marker, lost, arrival );
}
protected abstract void Payload( ReadOnlySpan<byte> data, uint ts, bool marker, bool lost, double arrival );
/// <summary>
/// RTCP from the server: sender reports tie an RTP time to the sender's wall clock (NTP), which is how the
/// latency readout knows when a frame was sent.
/// </summary>
public void OnRtcp( ReadOnlySpan<byte> p )
{
var at = 0;
while ( at + 8 <= p.Length )
{
var length = (((p[at + 2] << 8) | p[at + 3]) + 1) * 4;
if ( p[at + 1] == 200 && at + 20 <= p.Length && origin is not null )
{
var ntpSeconds = System.Buffers.Binary.BinaryPrimitives.ReadUInt32BigEndian( p[(at + 8)..] );
var ntpFraction = System.Buffers.Binary.BinaryPrimitives.ReadUInt32BigEndian( p[(at + 12)..] );
var rtp = System.Buffers.Binary.BinaryPrimitives.ReadUInt32BigEndian( p[(at + 16)..] );
var unix = ntpSeconds - 2208988800.0 + ntpFraction / 4294967296.0;
OnSenderReport( StreamTimeOf( rtp ), unix );
}
at += length;
}
}
/// <summary> Like <see cref="Time"/>, without moving the unwrap state (a report's time can be older). </summary>
long StreamTimeOf( uint ts )
{
long t = ts + (wraps << 32);
if ( lastTs >= 0 && ts < lastTs - 0x80000000L ) t += 1L << 32;
else if ( lastTs >= 0 && ts > lastTs + 0x80000000L ) t -= 1L << 32;
return (long)(originTime * 90000) + (t - origin.Value) * 90000 / Media.ClockRate;
}
/// <summary> Video's sender reports give the latency readout; audio's line the audio up with it (lip sync). </summary>
protected virtual bool IsVideo => false;
void OnSenderReport( long streamTime, double unixSeconds ) => Sink.SetWallClock( streamTime, unixSeconds, IsVideo );
}
/// <summary> RFC 6184: single NAL units, STAP-A aggregates and FU-A fragments. </summary>
sealed class H264Track : Track
{
readonly List<byte[]> nals = new();
List<byte> fragment;
uint auTs;
bool haveAu;
public override string Name => Sink.VideoInfo ?? "H.264";
protected override bool IsVideo => true;
public H264Track( MediaDesc media, LiveSegmenter sink ) : base( media, sink )
{
sink.ExpectVideo();
// parameter sets from the SDP, in case the camera doesn't repeat them in band
if ( media.Fmtp.TryGetValue( "sprop-parameter-sets", out var sets ) )
{
byte[] sps = null, pps = null;
foreach ( var b64 in sets.Split( ',' ) )
{
try
{
var nal = Convert.FromBase64String( b64.Trim() );
if ( H264.NalType( nal ) == H264.NalSps ) sps = nal;
else if ( H264.NalType( nal ) == H264.NalPps ) pps = nal;
}
catch ( FormatException ) { }
}
sink.SetParameterSets( sps, pps );
}
}
protected override void Payload( ReadOnlySpan<byte> d, uint ts, bool marker, bool lost, double arrival )
{
// a new timestamp means the previous access unit is complete (if its marker got lost)
if ( haveAu && ts != auTs ) Emit( arrival );
auTs = ts;
haveAu = true;
if ( lost ) fragment = null;
var type = d[0] & 0x1F;
if ( type is >= 1 and <= 23 )
{
nals.Add( d.ToArray() );
}
else if ( type == 24 ) // STAP-A
{
var i = 1;
while ( i + 2 <= d.Length )
{
var size = (d[i] << 8) | d[i + 1];
i += 2;
if ( i + size > d.Length ) break;
nals.Add( d.Slice( i, size ).ToArray() );
i += size;
}
}
else if ( type == 28 && d.Length > 2 ) // FU-A
{
var header = d[1];
if ( (header & 0x80) != 0 )
{
fragment = new List<byte> { (byte)((d[0] & 0xE0) | (header & 0x1F)) };
}
if ( fragment is not null )
{
for ( int i = 2; i < d.Length; i++ ) fragment.Add( d[i] );
if ( (header & 0x40) != 0 )
{
nals.Add( fragment.ToArray() );
fragment = null;
}
}
}
if ( marker ) Emit( arrival );
}
void Emit( double arrival )
{
haveAu = false;
if ( nals.Count == 0 ) return;
var t = Time( auTs, arrival );
// RTP carries presentation times only; cameras don't use B-frames, so decode time = presentation time
Sink.AddVideo( t, t, nals.ToList() );
nals.Clear();
}
}
/// <summary>
/// AV1 (the AOM RTP payload format): an aggregation header (Z: the first OBU element continues the last packet's,
/// Y: the last continues in the next one, W: element count, the last without a length), then OBU elements.
/// OBUs travel without size fields; they get them back here, which is how MP4 stores them.
/// </summary>
sealed class Av1Track : Track
{
readonly List<byte> unit = new();
List<byte> fragment;
uint unitTs;
bool haveUnit;
public override string Name => Sink.VideoInfo ?? "AV1";
protected override bool IsVideo => true;
public Av1Track( MediaDesc media, LiveSegmenter sink ) : base( media, sink )
{
sink.ExpectVideo();
}
protected override void Payload( ReadOnlySpan<byte> d, uint ts, bool marker, bool lost, double arrival )
{
if ( haveUnit && ts != unitTs ) Emit( arrival );
unitTs = ts;
haveUnit = true;
if ( lost ) fragment = null;
if ( d.Length < 2 ) return;
var z = (d[0] & 0x80) != 0;
var y = (d[0] & 0x40) != 0;
var w = (d[0] >> 4) & 3;
var i = 1;
var element = 0;
while ( i < d.Length )
{
element++;
int size;
if ( w == 0 || element < w )
{
var (length, n) = Resolver.Media.Av1.Leb128( d[i..] );
if ( n == 0 ) break;
i += n;
size = (int)Math.Min( length, (ulong)(d.Length - i) );
}
else size = d.Length - i;
var part = d.Slice( i, size );
i += size;
var continues = i >= d.Length && y;
if ( element == 1 && z )
{
if ( fragment is null ) continue; // its start was lost
Append( fragment, part );
if ( !continues ) { Finish( fragment.ToArray() ); fragment = null; }
}
else if ( continues )
{
fragment = new List<byte>( size * 2 );
Append( fragment, part );
}
else Finish( part );
}
if ( marker ) Emit( arrival );
}
static void Append( List<byte> list, ReadOnlySpan<byte> data )
{
foreach ( var b in data ) list.Add( b );
}
/// <summary> One whole OBU: into the temporal unit with a size field (dropping delimiters and padding). </summary>
void Finish( ReadOnlySpan<byte> obu )
{
if ( obu.Length < 1 ) return;
var type = (obu[0] >> 3) & 0xF;
if ( type is Resolver.Media.Av1.ObuTemporalDelimiter or Resolver.Media.Av1.ObuTileList or Resolver.Media.Av1.ObuPadding ) return;
if ( (obu[0] & 0x02) != 0 ) { Append( unit, obu ); return; }
var headerLength = (obu[0] & 0x04) != 0 ? 2 : 1;
if ( obu.Length < headerLength ) return;
Resolver.Media.Av1.AppendSized( unit, obu[..headerLength], obu[headerLength..] );
}
void Emit( double arrival )
{
haveUnit = false;
if ( unit.Count == 0 ) return;
Sink.AddAv1( Time( unitTs, arrival ), unit.ToArray() );
unit.Clear();
}
}
/// <summary> RFC 2435 Motion JPEG: each frame reassembled into a JPEG and shown as it arrives. </summary>
sealed class MjpegTrack : Track
{
readonly RtpJpeg jpeg = new();
public override string Name => Sink.VideoInfo ?? "Motion JPEG";
public MjpegTrack( MediaDesc media, LiveSegmenter sink ) : base( media, sink )
{
sink.ExpectMjpeg();
}
protected override bool IsVideo => true;
protected override void Payload( ReadOnlySpan<byte> d, uint ts, bool marker, bool lost, double arrival )
{
var t = Time( ts, arrival );
if ( jpeg.Add( d, ts, marker, lost ) is { } frame ) Sink.AddJpeg( t, frame );
}
}
/// <summary> RFC 3640 AAC (mpeg4-generic, AAC-hbr): AU headers, then the access units. </summary>
sealed class AacTrack : Track
{
readonly int sizeLength, indexLength;
readonly AacConfig config;
public override string Name => $"AAC {config?.SampleRate}Hz";
public AacTrack( MediaDesc media, LiveSegmenter sink ) : base( media, sink )
{
sizeLength = int.TryParse( media.Fmtp.GetValueOrDefault( "sizelength" ), out var s ) ? s : 13;
indexLength = int.TryParse( media.Fmtp.GetValueOrDefault( "indexlength" ), out var i ) ? i : 3;
if ( media.Fmtp.TryGetValue( "config", out var hex ) )
{
try { config = AacConfig.FromAsc( Convert.FromHexString( hex ) ); } catch ( FormatException ) { }
}
if ( config is not null ) sink.SetAacConfig( config );
}
protected override void Payload( ReadOnlySpan<byte> d, uint ts, bool marker, bool lost, double arrival )
{
if ( config is null || d.Length < 2 ) return;
var headerBits = (d[0] << 8) | d[1];
var headerBytes = (headerBits + 7) / 8;
var per = sizeLength + indexLength;
var count = per > 0 ? headerBits / per : 0;
var dataAt = 2 + headerBytes;
var baseTime = Time( ts, arrival );
for ( int k = 0; k < count; k++ )
{
var size = ReadBits( d.Slice( 2, headerBytes ), k * per, sizeLength );
if ( dataAt + size > d.Length ) break;
var pts = baseTime + (long)k * 1024 * 90000 / config.SampleRate;
Sink.AddAudio( pts, d.Slice( dataAt, size ).ToArray() );
dataAt += size;
}
}
static int ReadBits( ReadOnlySpan<byte> d, int bit, int count )
{
var v = 0;
for ( int i = 0; i < count; i++, bit++ )
v = (v << 1) | ((d[bit >> 3] >> (7 - (bit & 7))) & 1);
return v;
}
}
/// <summary> G.711 mu-law / A-law: decoded here, played through a SoundStream beside the video. </summary>
sealed class G711Track : Track
{
readonly bool alaw;
public override string Name => alaw ? "G.711 A-law" : "G.711 mu-law";
public G711Track( MediaDesc media, LiveSegmenter sink ) : base( media, sink )
{
alaw = media.Encoding == "PCMA" || media.PayloadType == 8;
if ( media.ClockRate <= 0 || media.ClockRate == 90000 ) media.ClockRate = 8000;
}
protected override void Payload( ReadOnlySpan<byte> d, uint ts, bool marker, bool lost, double arrival )
{
Sink.AddPcm( Time( ts, arrival ), G711.Decode( d, alaw ), Media.ClockRate );
}
}
}