refactor camera status caching

This commit is contained in:
Josh Hawkins
2026-10-02 07:15:43 -05:00
parent 9b02a077e9
commit 996ffe27ed
+47 -43
View File
@@ -6,6 +6,7 @@ import subprocess as sp
import threading import threading
import time import time
from collections import defaultdict, deque from collections import defaultdict, deque
from dataclasses import dataclass
from datetime import UTC, datetime, timedelta from datetime import UTC, datetime, timedelta
from multiprocessing import Queue, Value from multiprocessing import Queue, Value
from multiprocessing.synchronize import Event as MpEvent 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 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): class CameraWatchdog(threading.Thread):
def __init__( def __init__(
self, self,
@@ -186,33 +207,16 @@ class CameraWatchdog(threading.Thread):
self._stall_active: bool = False self._stall_active: bool = False
# Status caching to reduce message volume # Status caching to reduce message volume
self._last_detect_status: str | None = None self.detect_status = self._role_status("detect")
self._last_record_status: dict[str, str] = {} self.record_status = {
self._last_detect_status_update_time: float = 0.0 stream_type: self._role_status(role)
self._last_record_status_update_time: dict[str, float] = defaultdict(float) for stream_type, role in STREAM_TYPE_TO_ROLE.items()
}
def _send_detect_status(self, status: str, now: float) -> None: def _role_status(self, role: str) -> RoleStatus:
"""Send detect status only if changed or retry_interval has elapsed.""" return RoleStatus(
if ( self.requestor, f"{self.config.name}/status/{role}", self.sleeptime
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 _send_roles_offline(self, roles: list[CameraRoleEnum], now: float) -> None: def _send_roles_offline(self, roles: list[CameraRoleEnum], now: float) -> None:
"""Send offline status for each role of a restarted ffmpeg process.""" """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) stream_type = ROLE_TO_STREAM_TYPE.get(role.value)
if stream_type is not None: if stream_type is not None:
self._send_record_status(stream_type, "offline", now) self.record_status[stream_type].send("offline", now)
else: else:
self.requestor.send_data( self.requestor.send_data(
f"{self.config.name}/status/{role.value}", "offline" f"{self.config.name}/status/{role.value}", "offline"
@@ -388,11 +392,11 @@ class CameraWatchdog(threading.Thread):
# update camera status # update camera status
now = datetime.now().timestamp() now = datetime.now().timestamp()
self._send_detect_status("disabled", now) self.detect_status.send("disabled", now)
self._send_record_status(STREAM_TYPE_MAIN, "disabled", now) self.record_status[STREAM_TYPE_MAIN].send("disabled", now)
# cameras without a sub stream never get a record_sub topic # cameras without a sub stream never get a record_sub topic
if self.config.record.sub.enabled: 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 self.was_enabled = enabled
continue continue
@@ -465,7 +469,7 @@ class CameraWatchdog(threading.Thread):
can_restart = time_since_last_restart >= self.sleeptime can_restart = time_since_last_restart >= self.sleeptime
if not self.capture_thread.is_alive(): 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.camera_fps.value = 0
self.logger.error( self.logger.error(
f"Ffmpeg process crashed unexpectedly for {self.config.name}." f"Ffmpeg process crashed unexpectedly for {self.config.name}."
@@ -477,7 +481,7 @@ class CameraWatchdog(threading.Thread):
self.fps_overflow_count += 1 self.fps_overflow_count += 1
if self.fps_overflow_count == 3: if self.fps_overflow_count == 3:
self._send_detect_status("offline", now) self.detect_status.send("offline", now)
self.fps_overflow_count = 0 self.fps_overflow_count = 0
self.camera_fps.value = 0 self.camera_fps.value = 0
self.logger.info( self.logger.info(
@@ -487,7 +491,7 @@ class CameraWatchdog(threading.Thread):
self.reset_capture_thread(drain_output=False) self.reset_capture_thread(drain_output=False)
last_restart_time = now last_restart_time = now
elif now - self.capture_thread.current_frame.value > 20: 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.camera_fps.value = 0
self.logger.info( self.logger.info(
f"No frames received from {self.config.name} in 20 seconds. Exiting ffmpeg..." 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 last_restart_time = now
else: else:
# process is running normally # process is running normally
self._send_detect_status("online", now) self.detect_status.send("online", now)
self.fps_overflow_count = 0 self.fps_overflow_count = 0
for p in self.ffmpeg_other_processes: for p in self.ffmpeg_other_processes:
@@ -540,7 +544,7 @@ class CameraWatchdog(threading.Thread):
elif stale_stream is None: elif stale_stream is None:
if poll is None: if poll is None:
for stream_type in recorded_streams: 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( p["latest_segment_time"] = max(
self.latest_cache_segment_time[stream_type] 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"] p["cmd"], self.logger, p["logpipe"], ffmpeg_process=p["process"]
) )
if ( if self.detect_process_records_sub and self.config.record.stream_enabled(
self.detect_process_records_sub STREAM_TYPE_SUB
and self.config.record.stream_enabled(STREAM_TYPE_SUB)
and self.capture_thread is not None
and self.capture_thread.is_alive()
): ):
now_utc = datetime.now().astimezone(UTC) now_utc = datetime.now().astimezone(UTC)
stale_reason = self._stream_staleness(STREAM_TYPE_SUB, now_utc) stale_reason = self._stream_staleness(STREAM_TYPE_SUB, now_utc)
if stale_reason is None: if self.detect_status.last_status == "offline":
self._send_record_status(STREAM_TYPE_SUB, "online", now) # 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: elif can_restart:
self.logger.error( 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..." 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() self.reset_capture_thread()
last_restart_time = now last_restart_time = now