Compare commits

..
Author SHA1 Message Date
Josh Hawkins f6a0e23782 retry and report Frigate+ connection failures at startup
A Frigate+ model that wasn't cached yet needed api.frigate.video at startup, and when it couldn't be reached (a network that comes up late, a DNS blip) the requests ConnectionError wasn't a validation error, so Frigate crashed with a traceback before it could start. PlusApi requests now go through a session that retries connection failures for about 30 seconds, and a connection failure that outlasts that is raised as a ValueError so it shows up as a clear config validation error instead.
2026-10-02 08:33:57 -05:00
5 changed files with 106 additions and 60 deletions
+16 -3
View File
@@ -186,15 +186,28 @@ class ModelConfig(BaseModel):
# download the model if it doesn't exist
if not os.path.isfile(self.path):
download_url = plus_api.get_model_download_url(model_id)
r = requests.get(download_url)
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
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(plus_api.get_model_info(model_id), f)
json.dump(model_info, f)
model_info = load_plus_model_info(model_id)
+15 -4
View File
@@ -9,7 +9,9 @@ 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
@@ -101,6 +103,13 @@ 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
@@ -111,7 +120,9 @@ class PlusApi:
"Plus API key not set. See https://docs.frigate.video/integrations/plus#set-your-api-key"
)
parts = self.key.split(":")
r = requests.get(f"{self.host}/v1/auth/token", auth=(parts[0], parts[1]))
r = self._session.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()
@@ -121,19 +132,19 @@ class PlusApi:
return {"authorization": f"Bearer {self._token_data.get('accessToken')}"}
def _get(self, path: str) -> Response:
return requests.get(
return self._session.get(
f"{self.host}/v1/{path}", headers=self._get_authorization_header()
)
def _post(self, path: str, data: dict) -> Response:
return requests.post(
return self._session.post(
f"{self.host}/v1/{path}",
headers=self._get_authorization_header(),
json=data,
)
def _put(self, path: str, data: dict) -> Response:
return requests.put(
return self._session.put(
f"{self.host}/v1/{path}",
headers=self._get_authorization_header(),
json=data,
+6 -5
View File
@@ -59,6 +59,7 @@ def build_watchdog(
MagicMock(),
)
watchdog.requestor = MagicMock()
return watchdog
@@ -107,8 +108,8 @@ class TestCameraWatchdogStreamHealth(unittest.TestCase):
def test_status_goes_to_the_matching_role_topic(self):
watchdog = self._build_watchdog()
watchdog.record_status[STREAM_TYPE_MAIN].send("online", 100.0)
watchdog.record_status[STREAM_TYPE_SUB].send("offline", 100.0)
watchdog._send_record_status(STREAM_TYPE_MAIN, "online", 100.0)
watchdog._send_record_status(STREAM_TYPE_SUB, "offline", 100.0)
watchdog.requestor.send_data.assert_any_call(
"front_door/status/record", "online"
@@ -120,9 +121,9 @@ class TestCameraWatchdogStreamHealth(unittest.TestCase):
def test_status_is_cached_per_stream(self):
watchdog = self._build_watchdog()
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)
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)
assert watchdog.requestor.send_data.call_count == 2
+26
View File
@@ -6,6 +6,7 @@ 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
@@ -1595,6 +1596,31 @@ 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"},
+43 -48
View File
@@ -6,7 +6,6 @@ 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
@@ -109,27 +108,6 @@ 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,
@@ -208,16 +186,33 @@ class CameraWatchdog(threading.Thread):
self._stall_active: bool = False
# Status caching to reduce message volume
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()
}
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)
def _role_status(self, role: str) -> RoleStatus:
return RoleStatus(
self.requestor, f"{self.config.name}/status/{role}", self.sleeptime
)
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 _send_roles_offline(self, roles: list[CameraRoleEnum], now: float) -> None:
"""Send offline status for each role of a restarted ffmpeg process."""
@@ -227,7 +222,7 @@ class CameraWatchdog(threading.Thread):
stream_type = ROLE_TO_STREAM_TYPE.get(role.value)
if stream_type is not None:
self.record_status[stream_type].send("offline", now)
self._send_record_status(stream_type, "offline", now)
else:
self.requestor.send_data(
f"{self.config.name}/status/{role.value}", "offline"
@@ -393,11 +388,11 @@ class CameraWatchdog(threading.Thread):
# update camera status
now = datetime.now().timestamp()
self.detect_status.send("disabled", now)
self.record_status[STREAM_TYPE_MAIN].send("disabled", now)
self._send_detect_status("disabled", now)
self._send_record_status(STREAM_TYPE_MAIN, "disabled", now)
# cameras without a sub stream never get a record_sub topic
if self.config.record.sub.enabled:
self.record_status[STREAM_TYPE_SUB].send("disabled", now)
self._send_record_status(STREAM_TYPE_SUB, "disabled", now)
self.was_enabled = enabled
continue
@@ -470,7 +465,7 @@ class CameraWatchdog(threading.Thread):
can_restart = time_since_last_restart >= self.sleeptime
if not self.capture_thread.is_alive():
self.detect_status.send("offline", now)
self._send_detect_status("offline", now)
self.camera_fps.value = 0
self.logger.error(
f"Ffmpeg process crashed unexpectedly for {self.config.name}."
@@ -482,7 +477,7 @@ class CameraWatchdog(threading.Thread):
self.fps_overflow_count += 1
if self.fps_overflow_count == 3:
self.detect_status.send("offline", now)
self._send_detect_status("offline", now)
self.fps_overflow_count = 0
self.camera_fps.value = 0
self.logger.info(
@@ -492,7 +487,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.detect_status.send("offline", now)
self._send_detect_status("offline", now)
self.camera_fps.value = 0
self.logger.info(
f"No frames received from {self.config.name} in 20 seconds. Exiting ffmpeg..."
@@ -502,7 +497,7 @@ class CameraWatchdog(threading.Thread):
last_restart_time = now
else:
# process is running normally
self.detect_status.send("online", now)
self._send_detect_status("online", now)
self.fps_overflow_count = 0
for p in self.ffmpeg_other_processes:
@@ -545,7 +540,7 @@ class CameraWatchdog(threading.Thread):
elif stale_stream is None:
if poll is None:
for stream_type in recorded_streams:
self.record_status[stream_type].send("online", now)
self._send_record_status(stream_type, "online", now)
p["latest_segment_time"] = max(
self.latest_cache_segment_time[stream_type]
@@ -561,22 +556,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
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()
):
now_utc = datetime.now().astimezone(UTC)
stale_reason = self._stream_staleness(STREAM_TYPE_SUB, now_utc)
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)
if stale_reason is None:
self._send_record_status(STREAM_TYPE_SUB, "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.record_status[STREAM_TYPE_SUB].send("offline", now)
self._send_record_status(STREAM_TYPE_SUB, "offline", now)
self.reset_capture_thread()
last_restart_time = now