Compare 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
5 changed files with 60 additions and 106 deletions
+3 -16
View File
@@ -186,28 +186,15 @@ class ModelConfig(BaseModel):
# download the model if it doesn't exist
if not os.path.isfile(self.path):
try:
download_url = plus_api.get_model_download_url(model_id)
r = requests.get(download_url)
except requests.exceptions.ConnectionError as e:
raise ValueError(
f"Unable to connect to Frigate+ to download model {model_id}"
) from e
download_url = plus_api.get_model_download_url(model_id)
r = requests.get(download_url)
with open(self.path, "wb") as f:
f.write(r.content)
# download the model info if it doesn't exist
if not os.path.isfile(model_info_path):
try:
model_info = plus_api.get_model_info(model_id)
except requests.exceptions.ConnectionError as e:
raise ValueError(
f"Unable to connect to Frigate+ to download model info for {model_id}"
) from e
with open(model_info_path, "w") as f:
json.dump(model_info, f)
json.dump(plus_api.get_model_info(model_id), f)
model_info = load_plus_model_info(model_id)
+4 -15
View File
@@ -9,9 +9,7 @@ from typing import Any
import cv2
import requests
from numpy import ndarray
from requests.adapters import HTTPAdapter
from requests.models import Response
from urllib3.util.retry import Retry
from frigate.const import MODEL_CACHE_DIR, PLUS_API_HOST, PLUS_ENV_VAR
@@ -103,13 +101,6 @@ class PlusApi:
self._is_active: bool = self.key is not None
self._token_data: dict = {}
# Retry connection failures so a network that comes up late at startup
# doesn't fail the Frigate+ model download
self._session = requests.Session()
self._session.mount(
self.host, HTTPAdapter(max_retries=Retry(connect=5, backoff_factor=1))
)
def _refresh_token_if_needed(self) -> None:
if (
self._token_data.get("expires") is None
@@ -120,9 +111,7 @@ class PlusApi:
"Plus API key not set. See https://docs.frigate.video/integrations/plus#set-your-api-key"
)
parts = self.key.split(":")
r = self._session.get(
f"{self.host}/v1/auth/token", auth=(parts[0], parts[1])
)
r = requests.get(f"{self.host}/v1/auth/token", auth=(parts[0], parts[1]))
if not r.ok:
raise Exception(f"Unable to refresh API token: {r.text}")
self._token_data = r.json()
@@ -132,19 +121,19 @@ class PlusApi:
return {"authorization": f"Bearer {self._token_data.get('accessToken')}"}
def _get(self, path: str) -> Response:
return self._session.get(
return requests.get(
f"{self.host}/v1/{path}", headers=self._get_authorization_header()
)
def _post(self, path: str, data: dict) -> Response:
return self._session.post(
return requests.post(
f"{self.host}/v1/{path}",
headers=self._get_authorization_header(),
json=data,
)
def _put(self, path: str, data: dict) -> Response:
return self._session.put(
return requests.put(
f"{self.host}/v1/{path}",
headers=self._get_authorization_header(),
json=data,
+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
-26
View File
@@ -6,7 +6,6 @@ from copy import deepcopy
from unittest.mock import patch
import numpy as np
import requests
from pydantic import ValidationError
from ruamel.yaml.constructor import DuplicateKeyError
@@ -1596,31 +1595,6 @@ class TestConfig(unittest.TestCase):
frigate_config = FrigateConfig(**config)
assert frigate_config.primary_model.merged_labelmap[0] == "amazon"
@patch(
"frigate.plus.PlusApi.get_model_download_url",
side_effect=requests.exceptions.ConnectionError,
)
def test_plus_unreachable_is_validation_error(self, _):
config = {
"mqtt": {"host": "mqtt"},
"models": [{"path": "plus://unreachable", "devices": ["cpu"]}],
"cameras": {
"back": {
"ffmpeg": {
"inputs": [
{
"path": "rtsp://10.0.0.1:554/video",
"roles": ["detect"],
},
]
},
}
},
}
with self.assertRaisesRegex(ValidationError, "Unable to connect to Frigate+"):
FrigateConfig(**config)
def test_fails_on_invalid_role(self):
config = {
"mqtt": {"host": "mqtt"},
+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