From 9595afcc99aa365bdb9830bffd88b5d176bce041 Mon Sep 17 00:00:00 2001 From: Michael SCHAL Date: Wed, 10 Jun 2026 14:34:31 +0200 Subject: [PATCH] fix(streaming): stop streams dying mid-playback - reaper no longer kills exited sessions immediately: a finished VOD transcode was reaped (segments deleted) while still being watched - reuse completed VOD/recording sessions instead of re-transcoding - live input: add -rw_timeout 30s and -reconnect_at_eof so a stalled upstream triggers reconnect instead of wedging ffmpeg forever - reaper watchdog restarts live sessions whose playlist stopped updating - waitForPlaylist fails fast on clean ffmpeg exit without output - direct .ts proxy: close http.Client on stream end/error (socket leak) Co-Authored-By: Claude Fable 5 --- bin/api/streaming_handler.dart | 30 ++++++++++- bin/services/ffmpeg_session_manager.dart | 64 +++++++++++++++++++++--- 2 files changed, 85 insertions(+), 9 deletions(-) diff --git a/bin/api/streaming_handler.dart b/bin/api/streaming_handler.dart index 6d8303c..d454679 100644 --- a/bin/api/streaming_handler.dart +++ b/bin/api/streaming_handler.dart @@ -1,3 +1,4 @@ +import 'dart:async'; import 'dart:io'; import 'package:shelf/shelf.dart'; import 'package:shelf_router/shelf_router.dart'; @@ -215,7 +216,11 @@ Handler createLiveStreamHandler( if (useNvidiaGpu && quality != 'source') ...['-hwaccel', 'cuda'], '-headers', 'User-Agent: VLC/3.0.18 LibVLC/3.0.18\r\n', '-reconnect', '1', '-reconnect_streamed', '1', + '-reconnect_at_eof', '1', '-reconnect_delay_max', '10', + // Abort reads stuck for 30s so a stalled upstream triggers the + // reconnect logic instead of wedging the process forever. + '-rw_timeout', '30000000', '-i', targetUrl, ..._liveVideoArgs(quality, useNvidiaGpu), ..._audioArgs(withFilters: false), @@ -276,11 +281,32 @@ Handler createLiveStreamHandler( proxyRequest.headers['User-Agent'] = 'VLC/3.0.18 LibVLC/3.0.18'; proxyRequest.headers['Accept'] = '*/*'; - final response = await client.send(proxyRequest); + final http.StreamedResponse response; + try { + response = await client.send(proxyRequest); + } catch (e) { + client.close(); + return Response(502, body: 'Upstream connection failed'); + } + + // Close the client when the upstream stream ends or errors, otherwise + // each proxied stream leaks a socket. + final body = response.stream.transform>( + StreamTransformer.fromHandlers( + handleDone: (sink) { + sink.close(); + client.close(); + }, + handleError: (error, stackTrace, sink) { + sink.close(); + client.close(); + }, + ), + ); return Response( response.statusCode, - body: response.stream, + body: body, headers: { 'Content-Type': 'video/mp2t', 'Connection': 'keep-alive', diff --git a/bin/services/ffmpeg_session_manager.dart b/bin/services/ffmpeg_session_manager.dart index 7adb8f0..2c5af8e 100644 --- a/bin/services/ffmpeg_session_manager.dart +++ b/bin/services/ffmpeg_session_manager.dart @@ -8,6 +8,7 @@ class FfmpegSession { final Process process; final Directory dir; final bool isLive; + final DateTime startedAt = DateTime.now(); DateTime lastAccess = DateTime.now(); bool exited = false; int? exitCode; @@ -40,6 +41,9 @@ class FfmpegSessionManager { static const liveIdleTimeout = Duration(minutes: 4); static const vodIdleTimeout = Duration(minutes: 15); + /// Live playlist not rewritten for this long => FFmpeg is wedged. + static const liveStallTimeout = Duration(seconds: 45); + FfmpegSessionManager(this.baseDir); /// Wipe orphan session dirs and start the reaper. Call once at startup. @@ -90,6 +94,12 @@ class FfmpegSessionManager { return existing; } if (existing != null) { + // A VOD/recording transcode that finished cleanly is still fully + // playable from its segments — reuse it instead of re-transcoding. + if (!isLive && existing.exitCode == 0 && _playlistComplete(existing)) { + existing.touch(); + return existing; + } // Process died: clean up before restarting killSession(id); } @@ -128,6 +138,16 @@ class FfmpegSessionManager { return session; } + bool _playlistComplete(FfmpegSession session) { + final playlistFile = File('${session.dir.path}/playlist.m3u8'); + if (!playlistFile.existsSync()) return false; + try { + return playlistFile.readAsStringSync().contains('#EXT-X-ENDLIST'); + } catch (_) { + return false; + } + } + /// Waits until the session's playlist references at least one segment. /// Fails fast when the process dies before producing output, returning /// the recent stderr for diagnostics. @@ -139,17 +159,17 @@ class FfmpegSessionManager { final deadline = DateTime.now().add(timeout); while (DateTime.now().isBefore(deadline)) { - if (session.exited && session.exitCode != 0) { + if (playlistFile.existsSync() && + playlistFile.readAsStringSync().contains('.ts')) { + return (ready: true, error: null); + } + if (session.exited) { return ( ready: false, error: 'FFmpeg exited (${session.exitCode}): ' '${session.recentStderr.join().trim()}' ); } - if (playlistFile.existsSync() && - playlistFile.readAsStringSync().contains('.ts')) { - return (ready: true, error: null); - } await Future.delayed(const Duration(milliseconds: 500)); } return (ready: false, error: 'Timeout waiting for transcoder'); @@ -180,13 +200,43 @@ class FfmpegSessionManager { final now = DateTime.now(); for (final session in _sessions.values.toList()) { final timeout = session.isLive ? liveIdleTimeout : vodIdleTimeout; - if (session.exited || now.difference(session.lastAccess) > timeout) { + final idle = now.difference(session.lastAccess); + + // Exited sessions keep their segments until the idle timeout: a + // finished VOD transcode is usually still being watched, and killing + // it here would delete the segments mid-playback. + if (idle > timeout) { print( '[FFmpegManager] Reaping idle session ${session.id} ' - '(idle ${now.difference(session.lastAccess).inSeconds}s)', + '(idle ${idle.inSeconds}s)', ); killSession(session.id); + continue; + } + + // Watchdog: a live FFmpeg that is running but has not updated its + // playlist recently is wedged on a stalled upstream. Kill it so the + // player's next playlist poll restarts the transcoder. + if (session.isLive && !session.exited && _isStalled(session, now)) { + print('[FFmpegManager] Restarting stalled live session ${session.id}'); + killSession(session.id); } } } + + /// A live session is stalled when its playlist exists but has not been + /// rewritten for [liveStallTimeout] (FFmpeg rewrites it on every segment). + bool _isStalled(FfmpegSession session, DateTime now) { + final playlistFile = File('${session.dir.path}/playlist.m3u8'); + if (!playlistFile.existsSync()) { + // Never produced output: give it until liveStallTimeout after start. + return now.difference(session.startedAt) > liveStallTimeout; + } + try { + return now.difference(playlistFile.lastModifiedSync()) > + liveStallTimeout; + } catch (_) { + return false; + } + } }