From 76c051145f36445aa1d6b76ab364d5719924d932 Mon Sep 17 00:00:00 2001 From: Doloro1978 Date: Fri, 10 Jul 2026 13:33:46 +0100 Subject: [PATCH] fix: rtmp can timeout now --- crates/server/src/rtmp.rs | 19 +++++++++++++------ 1 file changed, 13 insertions(+), 6 deletions(-) diff --git a/crates/server/src/rtmp.rs b/crates/server/src/rtmp.rs index 7edc73b..c15dc48 100644 --- a/crates/server/src/rtmp.rs +++ b/crates/server/src/rtmp.rs @@ -1,4 +1,4 @@ -use std::{error::Error, sync::Arc}; +use std::{error::Error, sync::Arc, time::Duration}; use async_broadcast::broadcast; use chrono::Utc; @@ -12,7 +12,7 @@ use sea_orm::{DatabaseConnection, IntoActiveModel, sqlx::types::chrono::Local}; use tokio::{ io::{AsyncReadExt, AsyncWriteExt}, net::{TcpListener, TcpStream}, - time::Instant, + time::{Instant, timeout}, }; use tracing::{debug, info, warn}; @@ -116,6 +116,7 @@ impl Rtmp { video_tx.set_overflow(true); audio_tx.set_overflow(true); let mut parser: Option> = None; + let mut stream_id: Option = None; let mut aac_parser = AACParser::new(); let mut audio_proc = AudioProcesser::new(); let db = db.clone(); @@ -143,22 +144,27 @@ impl Rtmp { loop { let mut buf = [0u8; 4096]; - let n = match socket.read(&mut buf).await { - Ok(0) => { + const IDLE_TIMEOUT: Duration = Duration::from_secs(15); + let n = match timeout(IDLE_TIMEOUT, socket.read(&mut buf)).await { + Ok(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) => { + Ok(Ok(n)) => n, + Ok(Err(e)) => { warn!("RTMP read error: {e}"); if let Some(id) = current_stream_key_id { cleanup(id).await; } return; } + Err(_) => { + warn!("RTMP timed out (15 sec no packets) id: {:?}", stream_id); + return; + } }; let events = match session.handle_input(&buf[..n]) { Ok(e) => e, @@ -247,6 +253,7 @@ impl Rtmp { ) .await .unwrap(); + stream_id = Some(key.id); write_outbound(&mut socket, reply).await; } ServerSessionEvent::PublishStreamFinished {