back off restarts when a recording stream goes stale

This commit is contained in:
Josh Hawkins
2026-08-23 17:26:21 -05:00
parent ad8eab74e6
commit ccfe5cff46
2 changed files with 58 additions and 7 deletions
+35
View File
@@ -130,6 +130,41 @@ class TestCameraWatchdogStreamHealth(unittest.TestCase):
disabled = self._build_watchdog(sub_enabled=False) disabled = self._build_watchdog(sub_enabled=False)
assert disabled._recorded_streams(["detect", "record_sub"]) == [] assert disabled._recorded_streams(["detect", "record_sub"]) == []
def test_restart_grace_suppresses_repeat_staleness(self):
watchdog = self._build_watchdog()
now = datetime.now().astimezone(UTC)
stale = (now - timedelta(hours=1)).timestamp()
watchdog.latest_cache_segment_time[STREAM_TYPE_MAIN] = stale
watchdog.latest_valid_segment_time[STREAM_TYPE_MAIN] = stale
assert watchdog._stream_staleness(STREAM_TYPE_MAIN, now) is not None
watchdog._grant_restart_grace([STREAM_TYPE_MAIN], now)
assert watchdog._stream_staleness(STREAM_TYPE_MAIN, now) is None
assert (
watchdog._stream_staleness(STREAM_TYPE_MAIN, now + timedelta(seconds=89))
is None
)
assert (
watchdog._stream_staleness(STREAM_TYPE_MAIN, now + timedelta(seconds=91))
is not None
)
def test_restart_grace_is_per_stream(self):
watchdog = self._build_watchdog()
now = datetime.now().astimezone(UTC)
stale = (now - timedelta(hours=1)).timestamp()
for stream_type in (STREAM_TYPE_MAIN, STREAM_TYPE_SUB):
watchdog.latest_cache_segment_time[stream_type] = stale
watchdog.latest_valid_segment_time[stream_type] = stale
watchdog._grant_restart_grace([STREAM_TYPE_MAIN], now)
assert watchdog._stream_staleness(STREAM_TYPE_MAIN, now) is None
assert watchdog._stream_staleness(STREAM_TYPE_SUB, now) is not None
def test_stale_threshold_follows_each_stream_segment_time(self): def test_stale_threshold_follows_each_stream_segment_time(self):
watchdog = self._build_watchdog( watchdog = self._build_watchdog(
output_args={ output_args={
+23 -7
View File
@@ -41,6 +41,8 @@ from frigate.util.process import FrigateProcess
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
RECORD_GRACE_SECONDS = 90
def capture_frames( def capture_frames(
ffmpeg_process: sp.Popen[Any], ffmpeg_process: sp.Popen[Any],
@@ -161,6 +163,7 @@ class CameraWatchdog(threading.Thread):
self.latest_invalid_segment_time: dict[str, float] = defaultdict(float) self.latest_invalid_segment_time: dict[str, float] = defaultdict(float)
self.latest_cache_segment_time: dict[str, float] = defaultdict(float) self.latest_cache_segment_time: dict[str, float] = defaultdict(float)
self.record_enable_time: datetime | None = None self.record_enable_time: datetime | None = None
self.stream_grace_until: dict[str, datetime] = {}
# `valid` segments are published with the segment's start time, so the # `valid` segments are published with the segment's start time, so the
# gap between consecutive publishes can reach 2 * segment_time. Pad the # gap between consecutive publishes can reach 2 * segment_time. Pad the
@@ -212,14 +215,23 @@ class CameraWatchdog(threading.Thread):
self.latest_valid_segment_time.clear() self.latest_valid_segment_time.clear()
self.latest_invalid_segment_time.clear() self.latest_invalid_segment_time.clear()
self.latest_cache_segment_time.clear() self.latest_cache_segment_time.clear()
self.stream_grace_until.clear()
def _grant_restart_grace(self, stream_types: list[str], now_utc: datetime) -> None:
for stream_type in stream_types:
self.stream_grace_until[stream_type] = now_utc + timedelta(
seconds=RECORD_GRACE_SECONDS
)
def _stream_staleness(self, stream_type: str, now_utc: datetime) -> str | None: def _stream_staleness(self, stream_type: str, now_utc: datetime) -> str | None:
"""Return why the stream's segments are stale, or None if they're healthy.""" """Return why the stream's segments are stale, or None if they're healthy."""
# Grace period: 90 seconds allows time for ffmpeg to start and create # ffmpeg needs time to create a first segment after recording is
# the first segment # enabled and after a restart, per stream
in_grace_period = self.record_enable_time is not None and ( in_grace_period = (
now_utc - self.record_enable_time self.record_enable_time is not None
) < timedelta(seconds=90) and (now_utc - self.record_enable_time)
< timedelta(seconds=RECORD_GRACE_SECONDS)
) or now_utc < self.stream_grace_until.get(stream_type, now_utc)
if in_grace_period: if in_grace_period:
return None return None
@@ -489,7 +501,7 @@ class CameraWatchdog(threading.Thread):
stale_stream = stream_type stale_stream = stream_type
break break
if stale_stream is not None: if stale_stream is not None and can_restart:
self.logger.error( self.logger.error(
f"{stale_reason} for {self.config.name} ({stale_stream}) in the last {self.record_stale_threshold[stale_stream]}s. Restarting the ffmpeg record process..." f"{stale_reason} for {self.config.name} ({stale_stream}) in the last {self.record_stale_threshold[stale_stream]}s. Restarting the ffmpeg record process..."
) )
@@ -505,8 +517,11 @@ class CameraWatchdog(threading.Thread):
f"{self.config.name}/status/{role.value}", "offline" f"{self.config.name}/status/{role.value}", "offline"
) )
self._grant_restart_grace(recorded_streams, now_utc)
last_restart_time = now
continue continue
else: elif stale_stream is None:
for stream_type in recorded_streams: for stream_type in recorded_streams:
self._send_record_status(stream_type, "online", now) self._send_record_status(stream_type, "online", now)
@@ -545,6 +560,7 @@ class CameraWatchdog(threading.Thread):
) )
self._send_record_status(STREAM_TYPE_SUB, "offline", now) self._send_record_status(STREAM_TYPE_SUB, "offline", now)
self.reset_capture_thread() self.reset_capture_thread()
self._grant_restart_grace([STREAM_TYPE_SUB], now_utc)
last_restart_time = now last_restart_time = now
# Prune expired reconnect timestamps # Prune expired reconnect timestamps