mirror of
https://github.com/blakeblackshear/frigate.git
synced 2026-10-02 21:06:52 +03:00
Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e3955b18c1 | ||
|
|
996ffe27ed |
@@ -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
@@ -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,
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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
@@ -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
|
||||
|
||||
|
||||
Reference in New Issue
Block a user