diff --git a/frigate/video/ffmpeg.py b/frigate/video/ffmpeg.py index 969f28746d..2313564c5d 100644 --- a/frigate/video/ffmpeg.py +++ b/frigate/video/ffmpeg.py @@ -6,6 +6,7 @@ import subprocess as sp import threading import time from collections import defaultdict, deque +from dataclasses import dataclass from datetime import UTC, datetime, timedelta from multiprocessing import Queue, Value from multiprocessing.synchronize import Event as MpEvent @@ -108,6 +109,26 @@ def capture_frames( frame_index = 0 if frame_index == shm_frame_count - 1 else frame_index + 1 +@dataclass +class RoleStatus: + """Publishes a role's status when it changes or the resend interval elapses.""" + + requestor: InterProcessRequestor + topic: str + resend_interval: float + last_status: str | None = None + last_update_time: float = 0.0 + + def send(self, status: str, now: float) -> None: + if ( + status != self.last_status + or (now - self.last_update_time) >= self.resend_interval + ): + self.requestor.send_data(self.topic, status) + self.last_status = status + self.last_update_time = now + + class CameraWatchdog(threading.Thread): def __init__( self, @@ -186,33 +207,16 @@ class CameraWatchdog(threading.Thread): self._stall_active: bool = False # Status caching to reduce message volume - self._last_detect_status: str | None = None - self._last_record_status: dict[str, str] = {} - self._last_detect_status_update_time: float = 0.0 - self._last_record_status_update_time: dict[str, float] = defaultdict(float) + self.detect_status = self._role_status("detect") + self.record_status = { + stream_type: self._role_status(role) + for stream_type, role in STREAM_TYPE_TO_ROLE.items() + } - def _send_detect_status(self, status: str, now: float) -> None: - """Send detect status only if changed or retry_interval has elapsed.""" - if ( - status != self._last_detect_status - or (now - self._last_detect_status_update_time) >= self.sleeptime - ): - self.requestor.send_data(f"{self.config.name}/status/detect", status) - self._last_detect_status = status - self._last_detect_status_update_time = now - - def _send_record_status(self, stream_type: str, status: str, now: float) -> None: - """Send a record stream's status only if changed or retry_interval has elapsed.""" - if ( - status != self._last_record_status.get(stream_type) - or (now - self._last_record_status_update_time[stream_type]) - >= self.sleeptime - ): - self.requestor.send_data( - f"{self.config.name}/status/{STREAM_TYPE_TO_ROLE[stream_type]}", status - ) - self._last_record_status[stream_type] = status - self._last_record_status_update_time[stream_type] = now + def _role_status(self, role: str) -> RoleStatus: + return RoleStatus( + self.requestor, f"{self.config.name}/status/{role}", self.sleeptime + ) def _send_roles_offline(self, roles: list[CameraRoleEnum], now: float) -> None: """Send offline status for each role of a restarted ffmpeg process.""" @@ -222,7 +226,7 @@ class CameraWatchdog(threading.Thread): stream_type = ROLE_TO_STREAM_TYPE.get(role.value) if stream_type is not None: - self._send_record_status(stream_type, "offline", now) + self.record_status[stream_type].send("offline", now) else: self.requestor.send_data( f"{self.config.name}/status/{role.value}", "offline" @@ -388,11 +392,11 @@ class CameraWatchdog(threading.Thread): # update camera status now = datetime.now().timestamp() - self._send_detect_status("disabled", now) - self._send_record_status(STREAM_TYPE_MAIN, "disabled", now) + self.detect_status.send("disabled", now) + self.record_status[STREAM_TYPE_MAIN].send("disabled", now) # cameras without a sub stream never get a record_sub topic if self.config.record.sub.enabled: - self._send_record_status(STREAM_TYPE_SUB, "disabled", now) + self.record_status[STREAM_TYPE_SUB].send("disabled", now) self.was_enabled = enabled continue @@ -465,7 +469,7 @@ class CameraWatchdog(threading.Thread): can_restart = time_since_last_restart >= self.sleeptime if not self.capture_thread.is_alive(): - self._send_detect_status("offline", now) + self.detect_status.send("offline", now) self.camera_fps.value = 0 self.logger.error( f"Ffmpeg process crashed unexpectedly for {self.config.name}." @@ -477,7 +481,7 @@ class CameraWatchdog(threading.Thread): self.fps_overflow_count += 1 if self.fps_overflow_count == 3: - self._send_detect_status("offline", now) + self.detect_status.send("offline", now) self.fps_overflow_count = 0 self.camera_fps.value = 0 self.logger.info( @@ -487,7 +491,7 @@ class CameraWatchdog(threading.Thread): self.reset_capture_thread(drain_output=False) last_restart_time = now elif now - self.capture_thread.current_frame.value > 20: - self._send_detect_status("offline", now) + self.detect_status.send("offline", now) self.camera_fps.value = 0 self.logger.info( f"No frames received from {self.config.name} in 20 seconds. Exiting ffmpeg..." @@ -497,7 +501,7 @@ class CameraWatchdog(threading.Thread): last_restart_time = now else: # process is running normally - self._send_detect_status("online", now) + self.detect_status.send("online", now) self.fps_overflow_count = 0 for p in self.ffmpeg_other_processes: @@ -540,7 +544,7 @@ class CameraWatchdog(threading.Thread): elif stale_stream is None: if poll is None: for stream_type in recorded_streams: - self._send_record_status(stream_type, "online", now) + self.record_status[stream_type].send("online", now) p["latest_segment_time"] = max( self.latest_cache_segment_time[stream_type] @@ -556,22 +560,22 @@ class CameraWatchdog(threading.Thread): p["cmd"], self.logger, p["logpipe"], ffmpeg_process=p["process"] ) - if ( - self.detect_process_records_sub - and self.config.record.stream_enabled(STREAM_TYPE_SUB) - and self.capture_thread is not None - and self.capture_thread.is_alive() + if self.detect_process_records_sub and self.config.record.stream_enabled( + STREAM_TYPE_SUB ): now_utc = datetime.now().astimezone(UTC) stale_reason = self._stream_staleness(STREAM_TYPE_SUB, now_utc) - if stale_reason is None: - self._send_record_status(STREAM_TYPE_SUB, "online", now) + if self.detect_status.last_status == "offline": + # the sub stream is down whenever the detect process is + self.record_status[STREAM_TYPE_SUB].send("offline", now) + elif stale_reason is None: + self.record_status[STREAM_TYPE_SUB].send("online", now) elif can_restart: self.logger.error( f"{stale_reason} for {self.config.name} (sub, shared with detect) in the last {self.record_stale_threshold[STREAM_TYPE_SUB]}s. Restarting ffmpeg..." ) - self._send_record_status(STREAM_TYPE_SUB, "offline", now) + self.record_status[STREAM_TYPE_SUB].send("offline", now) self.reset_capture_thread() last_restart_time = now