fix: rtmp being unreliable and keep session alive when there are read errors

This commit is contained in:
2026-07-08 14:08:07 +01:00
parent f0a77de6f4
commit 888eace0bf
4 changed files with 78 additions and 24 deletions
+2
View File
@@ -43,4 +43,6 @@ EXPOSE 1935
EXPOSE 3000 EXPOSE 3000
EXPOSE 6969/udp EXPOSE 6969/udp
ENV RUST_LOG="warn"
CMD ["/rtmp-to-whip"] CMD ["/rtmp-to-whip"]
+1 -1
View File
@@ -396,7 +396,7 @@ async fn stream_handler(
} }
} else if answer.0 == request_id_clone { } else if answer.0 == request_id_clone {
return Response::builder() return Response::builder()
.status(StatusCode::NOT_FOUND) .status(StatusCode::UNSUPPORTED_MEDIA_TYPE)
.body(String::new()) .body(String::new())
.unwrap(); .unwrap();
} }
+2 -1
View File
@@ -14,6 +14,7 @@ use sea_orm::{
sqlx::types::chrono::{self, Local}, sqlx::types::chrono::{self, Local},
}; };
use tokio::{ use tokio::{
fs,
io::{AsyncReadExt, AsyncWriteExt}, io::{AsyncReadExt, AsyncWriteExt},
net::TcpListener, net::TcpListener,
sync::Mutex, sync::Mutex,
@@ -55,7 +56,6 @@ pub struct StreamSession {
pub stream_key_label: String, pub stream_key_label: String,
pub frame_channel: async_broadcast::Sender<Arc<VideoFrame>>, pub frame_channel: async_broadcast::Sender<Arc<VideoFrame>>,
pub audio_channel: async_broadcast::Sender<Arc<OpusAudioFrame>>, pub audio_channel: async_broadcast::Sender<Arc<OpusAudioFrame>>,
//
pub codec: Option<StreamCodec>, pub codec: Option<StreamCodec>,
} }
@@ -144,6 +144,7 @@ async fn main() -> Result<(), Box<dyn Error>> {
res = workers.join_next() => { res = workers.join_next() => {
if let Some(Err(e)) = res { if let Some(Err(e)) = res {
tracing::error!("worker panicked: {:?}", e); tracing::error!("worker panicked: {:?}", e);
// fs::File::
} else { } else {
tracing::error!("a worker exited unexpectedly"); tracing::error!("a worker exited unexpectedly");
} }
+73 -22
View File
@@ -18,10 +18,7 @@ use crate::{
StreamCodec, StreamSession, StreamCodec, StreamSession,
audio::{AACParser, AudioProcesser, OpusAudioFrame}, audio::{AACParser, AudioProcesser, OpusAudioFrame},
codec::{ codec::{
CodecParser, VideoFrame, CodecParser, VideoFrame, av1::Av1CodecParser, h264::H264CodecParser, h265::H265CodecParser,
av1::Av1CodecParser,
h264::H264CodecParser,
h265::H265CodecParser,
}, },
}; };
@@ -70,21 +67,21 @@ impl Rtmp {
let mut server = Handshake::new(PeerType::Server); let mut server = Handshake::new(PeerType::Server);
let mut c0_c1 = [0u8; 1537]; let mut c0_c1 = [0u8; 1537];
socket.read_exact(&mut c0_c1).await.unwrap(); socket.read_exact(&mut c0_c1).await?;
let s0_s1_s2 = match server.process_bytes(&c0_c1) { let s0_s1_s2 = match server.process_bytes(&c0_c1) {
Ok(HandshakeProcessResult::InProgress { response_bytes }) => response_bytes, Ok(HandshakeProcessResult::InProgress { response_bytes }) => response_bytes,
_ => panic!("handshake failed"), x => return Err(format!("unexpected handshake state: {:?}", x).into()),
}; };
socket.write_all(&s0_s1_s2).await.unwrap(); socket.write_all(&s0_s1_s2).await?;
let mut c2 = [0u8; 1536]; let mut c2 = [0u8; 1536];
socket.read_exact(&mut c2).await.unwrap(); socket.read_exact(&mut c2).await?;
match server.process_bytes(&c2) { match server.process_bytes(&c2) {
Ok(HandshakeProcessResult::Completed { .. }) => {} Ok(HandshakeProcessResult::Completed { .. }) => {}
Ok(HandshakeProcessResult::InProgress { response_bytes }) => { Ok(HandshakeProcessResult::InProgress { response_bytes }) => {
socket.write_all(&response_bytes).await.unwrap() socket.write_all(&response_bytes).await?
} }
x => panic!("Unexpected process_bytes response: {:?}", x), x => return Err(format!("unexpected handshake state: {:?}", x).into()),
} }
let (rtmp_session, init_bytes) = ServerSession::new(ServerSessionConfig::new()).unwrap(); let (rtmp_session, init_bytes) = ServerSession::new(ServerSessionConfig::new()).unwrap();
@@ -102,7 +99,13 @@ impl Rtmp {
loop { loop {
let (socket, peer_addr) = listener.accept().await.unwrap(); let (socket, peer_addr) = listener.accept().await.unwrap();
info!(%peer_addr, "RTMP connection accepted"); info!(%peer_addr, "RTMP connection accepted");
let (mut session, mut socket) = Rtmp::handshake(socket).await.unwrap(); let (mut session, mut socket) = match Rtmp::handshake(socket).await {
Ok(v) => v,
Err(e) => {
warn!(%peer_addr, "RTMP handshake failed: {e}");
continue;
}
};
info!(%peer_addr, "RTMP handshake complete"); info!(%peer_addr, "RTMP handshake complete");
let (mut video_tx, mut video_rx) = broadcast::<Arc<VideoFrame>>(32); let (mut video_tx, mut video_rx) = broadcast::<Arc<VideoFrame>>(32);
let (mut audio_tx, mut audio_rx) = broadcast::<Arc<OpusAudioFrame>>(32); let (mut audio_tx, mut audio_rx) = broadcast::<Arc<OpusAudioFrame>>(32);
@@ -119,14 +122,52 @@ impl Rtmp {
tokio::spawn(async move { tokio::spawn(async move {
let mut current_stream_key_id: Option<i32> = None; let mut current_stream_key_id: Option<i32> = None;
let mut codec_stamped = false; let mut codec_stamped = false;
let cleanup = |id: i32| {
let stream_sessions = stream_sessions.clone();
let db = db.clone();
async move {
stream_sessions.remove(&id);
if let Ok(Some(s)) =
stream_session::Model::get_active_by_stream_key_id(&db, id).await
{
s.into_active_model()
.finish_stream_session(&db, Local::now().into())
.await
.ok();
}
}
};
loop { loop {
let mut buf = [0u8; 4096]; let mut buf = [0u8; 4096];
let n = socket.read(&mut buf).await.unwrap(); let n = match socket.read(&mut buf).await {
if n == 0 { Ok(0) => {
debug!("RTMP connection closed by peer"); debug!("RTMP connection closed by peer");
return; if let Some(id) = current_stream_key_id {
} cleanup(id).await;
let events = session.handle_input(&buf[..n]).unwrap(); }
return;
}
Ok(n) => n,
Err(e) => {
warn!("RTMP read error: {e}");
if let Some(id) = current_stream_key_id {
cleanup(id).await;
}
return;
}
};
let events = match session.handle_input(&buf[..n]) {
Ok(e) => e,
Err(e) => {
warn!("RTMP session error: {e}");
if let Some(id) = current_stream_key_id {
cleanup(id).await;
}
return;
}
};
// Blankly using it, so it doesnt drop // Blankly using it, so it doesnt drop
video_rx.is_closed(); video_rx.is_closed();
audio_rx.is_closed(); audio_rx.is_closed();
@@ -234,7 +275,9 @@ impl Rtmp {
Ok(codec) => { Ok(codec) => {
if !codec_stamped { if !codec_stamped {
if let Some(id) = current_stream_key_id { if let Some(id) = current_stream_key_id {
if let Some(mut session) = stream_sessions.get_mut(&id) { if let Some(mut session) =
stream_sessions.get_mut(&id)
{
session.codec = Some(codec.clone()); session.codec = Some(codec.clone());
} }
} }
@@ -245,7 +288,8 @@ impl Rtmp {
let p = parser.get_or_insert_with(|| { let p = parser.get_or_insert_with(|| {
Box::new(H264CodecParser::new()) Box::new(H264CodecParser::new())
}); });
if let Some(frame) = p.parse(&data, timestamp.value) { if let Some(frame) = p.parse(&data, timestamp.value)
{
video_tx.broadcast(Arc::new(frame)).await.ok(); video_tx.broadcast(Arc::new(frame)).await.ok();
} }
} }
@@ -253,7 +297,8 @@ impl Rtmp {
let p = parser.get_or_insert_with(|| { let p = parser.get_or_insert_with(|| {
Box::new(H265CodecParser::new()) Box::new(H265CodecParser::new())
}); });
if let Some(frame) = p.parse(&data, timestamp.value) { if let Some(frame) = p.parse(&data, timestamp.value)
{
video_tx.broadcast(Arc::new(frame)).await.ok(); video_tx.broadcast(Arc::new(frame)).await.ok();
} }
} }
@@ -271,10 +316,16 @@ impl Rtmp {
bytes = frame.data.len(), bytes = frame.data.len(),
"AV1 frame → broadcast" "AV1 frame → broadcast"
); );
video_tx.broadcast(Arc::new(frame)).await.ok(); video_tx
.broadcast(Arc::new(frame))
.await
.ok();
} }
None => { None => {
debug!(pkt_type, "AV1 packet produced no frame (seq header or unknown type)"); debug!(
pkt_type,
"AV1 packet produced no frame (seq header or unknown type)"
);
} }
} }
} }