Compare commits

...
2 Commits
Author SHA1 Message Date
Josh Hawkins e3955b18c1 fix tests 2026-10-02 07:22:55 -05:00
Josh Hawkins 996ffe27ed refactor camera status caching 2026-10-02 07:15:43 -05:00
2 changed files with 53 additions and 49 deletions
+5 -6
View File
@@ -59,7 +59,6 @@ def build_watchdog(
MagicMock(),
)
watchdog.requestor = MagicMock()
return watchdog
@@ -108,8 +107,8 @@ class TestCameraWatchdogStreamHealth(unittest.TestCase):
def test_status_goes_to_the_matching_role_topic(self):
watchdog = self._build_watchdog()
watchdog._send_record_status(STREAM_TYPE_MAIN, "online", 100.0)
watchdog._send_record_status(STREAM_TYPE_SUB, "offline", 100.0)
watchdog.record_status[STREAM_TYPE_MAIN].send("online", 100.0)
watchdog.record_status[STREAM_TYPE_SUB].send("offline", 100.0)
watchdog.requestor.send_data.assert_any_call(
"front_door/status/record", "online"
@@ -121,9 +120,9 @@ class TestCameraWatchdogStreamHealth(unittest.TestCase):
def test_status_is_cached_per_stream(self):
watchdog = self._build_watchdog()
watchdog._send_record_status(STREAM_TYPE_MAIN, "online", 100.0)
watchdog._send_record_status(STREAM_TYPE_SUB, "online", 100.0)
watchdog._send_record_status(STREAM_TYPE_MAIN, "online", 100.0)
watchdog.record_status[STREAM_TYPE_MAIN].send("online", 100.0)
watchdog.record_status[STREAM_TYPE_SUB].send("online", 100.0)
watchdog.record_status[STREAM_TYPE_MAIN].send("online", 100.0)
assert watchdog.requestor.send_data.call_count == 2
+48 -43
View File
@@ -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,27 @@ 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:
"""Publish a changed status or resend it after the configured interval."""
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 +208,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 +227,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 +393,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 +470,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 +482,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 +492,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 +502,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 +545,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 +561,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