diff --git a/Dockerfile b/Dockerfile index 6786fe5..c942b1d 100644 --- a/Dockerfile +++ b/Dockerfile @@ -43,4 +43,6 @@ EXPOSE 1935 EXPOSE 3000 EXPOSE 6969/udp +ENV RUST_LOG="warn" + CMD ["/rtmp-to-whip"] diff --git a/crates/server/src/http.rs b/crates/server/src/http.rs index 30d7d5e..e7cc5d2 100644 --- a/crates/server/src/http.rs +++ b/crates/server/src/http.rs @@ -396,7 +396,7 @@ async fn stream_handler( } } else if answer.0 == request_id_clone { return Response::builder() - .status(StatusCode::NOT_FOUND) + .status(StatusCode::UNSUPPORTED_MEDIA_TYPE) .body(String::new()) .unwrap(); } diff --git a/crates/server/src/main.rs b/crates/server/src/main.rs index 6d3f5f3..bd2b960 100644 --- a/crates/server/src/main.rs +++ b/crates/server/src/main.rs @@ -14,6 +14,7 @@ use sea_orm::{ sqlx::types::chrono::{self, Local}, }; use tokio::{ + fs, io::{AsyncReadExt, AsyncWriteExt}, net::TcpListener, sync::Mutex, @@ -55,7 +56,6 @@ pub struct StreamSession { pub stream_key_label: String, pub frame_channel: async_broadcast::Sender>, pub audio_channel: async_broadcast::Sender>, - // pub codec: Option, } @@ -144,6 +144,7 @@ async fn main() -> Result<(), Box> { res = workers.join_next() => { if let Some(Err(e)) = res { tracing::error!("worker panicked: {:?}", e); + // fs::File:: } else { tracing::error!("a worker exited unexpectedly"); } diff --git a/crates/server/src/rtmp.rs b/crates/server/src/rtmp.rs index 892c434..eadf16e 100644 --- a/crates/server/src/rtmp.rs +++ b/crates/server/src/rtmp.rs @@ -18,10 +18,7 @@ use crate::{ StreamCodec, StreamSession, audio::{AACParser, AudioProcesser, OpusAudioFrame}, codec::{ - CodecParser, VideoFrame, - av1::Av1CodecParser, - h264::H264CodecParser, - h265::H265CodecParser, + CodecParser, VideoFrame, av1::Av1CodecParser, h264::H264CodecParser, h265::H265CodecParser, }, }; @@ -70,21 +67,21 @@ impl Rtmp { let mut server = Handshake::new(PeerType::Server); 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) { 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]; - socket.read_exact(&mut c2).await.unwrap(); + socket.read_exact(&mut c2).await?; match server.process_bytes(&c2) { Ok(HandshakeProcessResult::Completed { .. }) => {} 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(); @@ -102,7 +99,13 @@ impl Rtmp { loop { let (socket, peer_addr) = listener.accept().await.unwrap(); 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"); let (mut video_tx, mut video_rx) = broadcast::>(32); let (mut audio_tx, mut audio_rx) = broadcast::>(32); @@ -119,14 +122,52 @@ impl Rtmp { tokio::spawn(async move { let mut current_stream_key_id: Option = None; 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 { let mut buf = [0u8; 4096]; - let n = socket.read(&mut buf).await.unwrap(); - if n == 0 { - debug!("RTMP connection closed by peer"); - return; - } - let events = session.handle_input(&buf[..n]).unwrap(); + let n = match socket.read(&mut buf).await { + Ok(0) => { + debug!("RTMP connection closed by peer"); + if let Some(id) = current_stream_key_id { + cleanup(id).await; + } + 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 video_rx.is_closed(); audio_rx.is_closed(); @@ -234,7 +275,9 @@ impl Rtmp { Ok(codec) => { if !codec_stamped { 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()); } } @@ -245,7 +288,8 @@ impl Rtmp { let p = parser.get_or_insert_with(|| { 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(); } } @@ -253,7 +297,8 @@ impl Rtmp { let p = parser.get_or_insert_with(|| { 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(); } } @@ -271,10 +316,16 @@ impl Rtmp { bytes = frame.data.len(), "AV1 frame → broadcast" ); - video_tx.broadcast(Arc::new(frame)).await.ok(); + video_tx + .broadcast(Arc::new(frame)) + .await + .ok(); } 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)" + ); } } }