mirror of
https://github.com/blakeblackshear/frigate.git
synced 2026-10-03 13:26:48 +03:00
Compare commits
4
Commits
overflow-menu
...
dev
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
9b2839f4fb | ||
|
|
c06bf97b5e | ||
|
|
f3723698cd | ||
|
|
38b87feece |
@@ -571,6 +571,8 @@ notifications:
|
|||||||
enabled: False
|
enabled: False
|
||||||
# Optional: Email for push service to reach out to
|
# Optional: Email for push service to reach out to
|
||||||
# NOTE: This is required to use notifications
|
# NOTE: This is required to use notifications
|
||||||
|
# NOTE: Email can be specified with an environment variable or docker secrets that must begin with 'FRIGATE_'.
|
||||||
|
# e.g. email: '{FRIGATE_NOTIFICATION_EMAIL}'
|
||||||
email: "admin@example.com"
|
email: "admin@example.com"
|
||||||
# Optional: Cooldown time for notifications in seconds (default: shown below)
|
# Optional: Cooldown time for notifications in seconds (default: shown below)
|
||||||
cooldown: 0
|
cooldown: 0
|
||||||
|
|||||||
+14
-3
@@ -311,9 +311,12 @@ def config(request: Request):
|
|||||||
mode="json", warnings="none", exclude_none=True
|
mode="json", warnings="none", exclude_none=True
|
||||||
)
|
)
|
||||||
|
|
||||||
# remove environment_vars for non-admin users
|
is_admin = request.headers.get("remote-role") == "admin"
|
||||||
if request.headers.get("remote-role") != "admin":
|
|
||||||
|
# hide environment_vars and the notification email from non-admin users
|
||||||
|
if not is_admin:
|
||||||
config.pop("environment_vars", None)
|
config.pop("environment_vars", None)
|
||||||
|
redact_credential(config["notifications"], "email")
|
||||||
|
|
||||||
# redact mqtt credentials
|
# redact mqtt credentials
|
||||||
redact_credential(config["mqtt"], "password")
|
redact_credential(config["mqtt"], "password")
|
||||||
@@ -370,7 +373,15 @@ def config(request: Request):
|
|||||||
camera_name
|
camera_name
|
||||||
)
|
)
|
||||||
if base_sections:
|
if base_sections:
|
||||||
camera_dict["base_config"] = base_sections
|
# copy so redaction below can't alter the profile manager's cache
|
||||||
|
camera_dict["base_config"] = copy.deepcopy(base_sections)
|
||||||
|
|
||||||
|
# cameras inherit the global notification email
|
||||||
|
if not is_admin:
|
||||||
|
redact_credential(camera_dict["notifications"], "email")
|
||||||
|
redact_credential(
|
||||||
|
camera_dict.get("base_config", {}).get("notifications", {}), "email"
|
||||||
|
)
|
||||||
|
|
||||||
# remove go2rtc stream passwords
|
# remove go2rtc stream passwords
|
||||||
go2rtc: dict[str, Any] = config_obj.go2rtc.model_dump(
|
go2rtc: dict[str, Any] = config_obj.go2rtc.model_dump(
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
from pydantic import Field
|
from pydantic import Field
|
||||||
|
|
||||||
from ..base import FrigateBaseModel
|
from ..base import FrigateBaseModel
|
||||||
|
from ..env import EnvString
|
||||||
|
|
||||||
__all__ = ["NotificationConfig"]
|
__all__ = ["NotificationConfig"]
|
||||||
|
|
||||||
@@ -11,7 +12,7 @@ class NotificationConfig(FrigateBaseModel):
|
|||||||
title="Enable notifications",
|
title="Enable notifications",
|
||||||
description="Enable or disable notifications for all cameras; can be overridden per-camera.",
|
description="Enable or disable notifications for all cameras; can be overridden per-camera.",
|
||||||
)
|
)
|
||||||
email: str | None = Field(
|
email: EnvString | None = Field(
|
||||||
default=None,
|
default=None,
|
||||||
title="Notification email",
|
title="Notification email",
|
||||||
description="Email address used for push notifications or required by certain notification providers.",
|
description="Email address used for push notifications or required by certain notification providers.",
|
||||||
|
|||||||
@@ -186,15 +186,28 @@ class ModelConfig(BaseModel):
|
|||||||
|
|
||||||
# download the model if it doesn't exist
|
# download the model if it doesn't exist
|
||||||
if not os.path.isfile(self.path):
|
if not os.path.isfile(self.path):
|
||||||
|
try:
|
||||||
download_url = plus_api.get_model_download_url(model_id)
|
download_url = plus_api.get_model_download_url(model_id)
|
||||||
r = requests.get(download_url)
|
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:
|
with open(self.path, "wb") as f:
|
||||||
f.write(r.content)
|
f.write(r.content)
|
||||||
|
|
||||||
# download the model info if it doesn't exist
|
# download the model info if it doesn't exist
|
||||||
if not os.path.isfile(model_info_path):
|
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:
|
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)
|
model_info = load_plus_model_info(model_id)
|
||||||
|
|
||||||
|
|||||||
+15
-4
@@ -9,7 +9,9 @@ from typing import Any
|
|||||||
import cv2
|
import cv2
|
||||||
import requests
|
import requests
|
||||||
from numpy import ndarray
|
from numpy import ndarray
|
||||||
|
from requests.adapters import HTTPAdapter
|
||||||
from requests.models import Response
|
from requests.models import Response
|
||||||
|
from urllib3.util.retry import Retry
|
||||||
|
|
||||||
from frigate.const import MODEL_CACHE_DIR, PLUS_API_HOST, PLUS_ENV_VAR
|
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._is_active: bool = self.key is not None
|
||||||
self._token_data: dict = {}
|
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:
|
def _refresh_token_if_needed(self) -> None:
|
||||||
if (
|
if (
|
||||||
self._token_data.get("expires") is None
|
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"
|
"Plus API key not set. See https://docs.frigate.video/integrations/plus#set-your-api-key"
|
||||||
)
|
)
|
||||||
parts = self.key.split(":")
|
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:
|
if not r.ok:
|
||||||
raise Exception(f"Unable to refresh API token: {r.text}")
|
raise Exception(f"Unable to refresh API token: {r.text}")
|
||||||
self._token_data = r.json()
|
self._token_data = r.json()
|
||||||
@@ -121,19 +132,19 @@ class PlusApi:
|
|||||||
return {"authorization": f"Bearer {self._token_data.get('accessToken')}"}
|
return {"authorization": f"Bearer {self._token_data.get('accessToken')}"}
|
||||||
|
|
||||||
def _get(self, path: str) -> Response:
|
def _get(self, path: str) -> Response:
|
||||||
return requests.get(
|
return self._session.get(
|
||||||
f"{self.host}/v1/{path}", headers=self._get_authorization_header()
|
f"{self.host}/v1/{path}", headers=self._get_authorization_header()
|
||||||
)
|
)
|
||||||
|
|
||||||
def _post(self, path: str, data: dict) -> Response:
|
def _post(self, path: str, data: dict) -> Response:
|
||||||
return requests.post(
|
return self._session.post(
|
||||||
f"{self.host}/v1/{path}",
|
f"{self.host}/v1/{path}",
|
||||||
headers=self._get_authorization_header(),
|
headers=self._get_authorization_header(),
|
||||||
json=data,
|
json=data,
|
||||||
)
|
)
|
||||||
|
|
||||||
def _put(self, path: str, data: dict) -> Response:
|
def _put(self, path: str, data: dict) -> Response:
|
||||||
return requests.put(
|
return self._session.put(
|
||||||
f"{self.host}/v1/{path}",
|
f"{self.host}/v1/{path}",
|
||||||
headers=self._get_authorization_header(),
|
headers=self._get_authorization_header(),
|
||||||
json=data,
|
json=data,
|
||||||
|
|||||||
+17
-1
@@ -236,6 +236,15 @@ def skipped_percent(skipped_fps: float, camera_fps: float, enabled: bool) -> flo
|
|||||||
return round(skipped_fps / camera_fps * 100, 1)
|
return round(skipped_fps / camera_fps * 100, 1)
|
||||||
|
|
||||||
|
|
||||||
|
def get_go2rtc_pid(cpu_usages: dict[str, dict[str, Any]]) -> int | None:
|
||||||
|
"""Find the pid of the running go2rtc process in the cpu usages."""
|
||||||
|
for pid, usage in cpu_usages.items():
|
||||||
|
if usage.get("cmdline", "").split(" ")[0].endswith("/go2rtc"):
|
||||||
|
return int(pid)
|
||||||
|
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
def stats_snapshot(
|
def stats_snapshot(
|
||||||
config: FrigateConfig,
|
config: FrigateConfig,
|
||||||
stats_tracking: StatsTrackingTypes,
|
stats_tracking: StatsTrackingTypes,
|
||||||
@@ -356,6 +365,14 @@ def stats_snapshot(
|
|||||||
|
|
||||||
stats["service"]["storage"]["/dev/shm"] = calculate_shm_requirements(config)
|
stats["service"]["storage"]["/dev/shm"] = calculate_shm_requirements(config)
|
||||||
|
|
||||||
|
cpu_usages = stats.get("cpu_usages", {})
|
||||||
|
|
||||||
|
# go2rtc is supervised by s6, so its pid changes when s6 restarts it
|
||||||
|
go2rtc_pid = get_go2rtc_pid(cpu_usages)
|
||||||
|
|
||||||
|
if go2rtc_pid is not None:
|
||||||
|
stats_tracking["processes"]["go2rtc"] = go2rtc_pid
|
||||||
|
|
||||||
stats["processes"] = {}
|
stats["processes"] = {}
|
||||||
for name, pid in stats_tracking["processes"].items():
|
for name, pid in stats_tracking["processes"].items():
|
||||||
stats["processes"][name] = {
|
stats["processes"][name] = {
|
||||||
@@ -364,7 +381,6 @@ def stats_snapshot(
|
|||||||
|
|
||||||
# Embed cpu/mem stats into detectors, cameras, and processes
|
# Embed cpu/mem stats into detectors, cameras, and processes
|
||||||
# so history consumers don't need the full cpu_usages dict
|
# so history consumers don't need the full cpu_usages dict
|
||||||
cpu_usages = stats.get("cpu_usages", {})
|
|
||||||
|
|
||||||
for det_stats in stats["detectors"].values():
|
for det_stats in stats["detectors"].values():
|
||||||
pid_str = str(det_stats.get("pid", ""))
|
pid_str = str(det_stats.get("pid", ""))
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ from unittest.mock import Mock, patch
|
|||||||
|
|
||||||
import frigate.genai
|
import frigate.genai
|
||||||
from frigate.config import GenAIProviderEnum
|
from frigate.config import GenAIProviderEnum
|
||||||
|
from frigate.config.env import FRIGATE_ENV_VARS
|
||||||
from frigate.const import MODEL_CACHE_DIR, REDACTED_CREDENTIAL_SENTINEL
|
from frigate.const import MODEL_CACHE_DIR, REDACTED_CREDENTIAL_SENTINEL
|
||||||
from frigate.genai import GenAIClient
|
from frigate.genai import GenAIClient
|
||||||
from frigate.models import Event, Recordings, ReviewSegment
|
from frigate.models import Event, Recordings, ReviewSegment
|
||||||
@@ -111,6 +112,30 @@ class TestHttpApp(BaseTestHttp):
|
|||||||
mqtt = response.json()["mqtt"]
|
mqtt = response.json()["mqtt"]
|
||||||
assert mqtt["password"] == REDACTED_CREDENTIAL_SENTINEL
|
assert mqtt["password"] == REDACTED_CREDENTIAL_SENTINEL
|
||||||
|
|
||||||
|
def test_config_response_hides_notification_email_from_viewers(self):
|
||||||
|
self.minimal_config["notifications"] = {"email": "{FRIGATE_TEST_EMAIL}"}
|
||||||
|
|
||||||
|
with patch.dict(FRIGATE_ENV_VARS, {"FRIGATE_TEST_EMAIL": "me@example.com"}):
|
||||||
|
app = super().create_app()
|
||||||
|
|
||||||
|
assert app.frigate_config.notifications.email == "me@example.com"
|
||||||
|
|
||||||
|
with AuthTestClient(app) as client:
|
||||||
|
response = client.get(
|
||||||
|
"/config",
|
||||||
|
headers={"remote-user": "viewer", "remote-role": "viewer"},
|
||||||
|
)
|
||||||
|
assert response.status_code == 200
|
||||||
|
config = response.json()
|
||||||
|
assert config["notifications"]["email"] == REDACTED_CREDENTIAL_SENTINEL
|
||||||
|
assert (
|
||||||
|
config["cameras"]["front_door"]["notifications"]["email"]
|
||||||
|
== REDACTED_CREDENTIAL_SENTINEL
|
||||||
|
)
|
||||||
|
|
||||||
|
response = client.get("/config")
|
||||||
|
assert response.json()["notifications"]["email"] == "me@example.com"
|
||||||
|
|
||||||
def test_config_response_keeps_plus_model_reference(self):
|
def test_config_response_keeps_plus_model_reference(self):
|
||||||
model_id = "test_plus_reference"
|
model_id = "test_plus_reference"
|
||||||
model_path = os.path.join(MODEL_CACHE_DIR, model_id)
|
model_path = os.path.join(MODEL_CACHE_DIR, model_id)
|
||||||
|
|||||||
@@ -59,7 +59,6 @@ def build_watchdog(
|
|||||||
MagicMock(),
|
MagicMock(),
|
||||||
)
|
)
|
||||||
|
|
||||||
watchdog.requestor = MagicMock()
|
|
||||||
return watchdog
|
return watchdog
|
||||||
|
|
||||||
|
|
||||||
@@ -108,8 +107,8 @@ class TestCameraWatchdogStreamHealth(unittest.TestCase):
|
|||||||
def test_status_goes_to_the_matching_role_topic(self):
|
def test_status_goes_to_the_matching_role_topic(self):
|
||||||
watchdog = self._build_watchdog()
|
watchdog = self._build_watchdog()
|
||||||
|
|
||||||
watchdog._send_record_status(STREAM_TYPE_MAIN, "online", 100.0)
|
watchdog.record_status[STREAM_TYPE_MAIN].send("online", 100.0)
|
||||||
watchdog._send_record_status(STREAM_TYPE_SUB, "offline", 100.0)
|
watchdog.record_status[STREAM_TYPE_SUB].send("offline", 100.0)
|
||||||
|
|
||||||
watchdog.requestor.send_data.assert_any_call(
|
watchdog.requestor.send_data.assert_any_call(
|
||||||
"front_door/status/record", "online"
|
"front_door/status/record", "online"
|
||||||
@@ -121,9 +120,9 @@ class TestCameraWatchdogStreamHealth(unittest.TestCase):
|
|||||||
def test_status_is_cached_per_stream(self):
|
def test_status_is_cached_per_stream(self):
|
||||||
watchdog = self._build_watchdog()
|
watchdog = self._build_watchdog()
|
||||||
|
|
||||||
watchdog._send_record_status(STREAM_TYPE_MAIN, "online", 100.0)
|
watchdog.record_status[STREAM_TYPE_MAIN].send("online", 100.0)
|
||||||
watchdog._send_record_status(STREAM_TYPE_SUB, "online", 100.0)
|
watchdog.record_status[STREAM_TYPE_SUB].send("online", 100.0)
|
||||||
watchdog._send_record_status(STREAM_TYPE_MAIN, "online", 100.0)
|
watchdog.record_status[STREAM_TYPE_MAIN].send("online", 100.0)
|
||||||
|
|
||||||
assert watchdog.requestor.send_data.call_count == 2
|
assert watchdog.requestor.send_data.call_count == 2
|
||||||
|
|
||||||
|
|||||||
@@ -6,6 +6,7 @@ from copy import deepcopy
|
|||||||
from unittest.mock import patch
|
from unittest.mock import patch
|
||||||
|
|
||||||
import numpy as np
|
import numpy as np
|
||||||
|
import requests
|
||||||
from pydantic import ValidationError
|
from pydantic import ValidationError
|
||||||
from ruamel.yaml.constructor import DuplicateKeyError
|
from ruamel.yaml.constructor import DuplicateKeyError
|
||||||
|
|
||||||
@@ -1595,6 +1596,31 @@ class TestConfig(unittest.TestCase):
|
|||||||
frigate_config = FrigateConfig(**config)
|
frigate_config = FrigateConfig(**config)
|
||||||
assert frigate_config.primary_model.merged_labelmap[0] == "amazon"
|
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):
|
def test_fails_on_invalid_role(self):
|
||||||
config = {
|
config = {
|
||||||
"mqtt": {"host": "mqtt"},
|
"mqtt": {"host": "mqtt"},
|
||||||
|
|||||||
@@ -0,0 +1,70 @@
|
|||||||
|
"""Tests for resolving the go2rtc pid from cpu usages."""
|
||||||
|
|
||||||
|
import unittest
|
||||||
|
from types import SimpleNamespace
|
||||||
|
from unittest.mock import Mock, patch
|
||||||
|
|
||||||
|
from frigate.stats.util import get_go2rtc_pid, stats_snapshot
|
||||||
|
|
||||||
|
|
||||||
|
class TestGo2rtcPid(unittest.TestCase):
|
||||||
|
def test_finds_go2rtc_by_binary_path(self):
|
||||||
|
cpu_usages = {
|
||||||
|
"frigate.full_system": {"cpu": "1.0", "mem": "2.0"},
|
||||||
|
"100": {"cmdline": "ffmpeg -i rtsp://127.0.0.1:8554/go2rtc_cam"},
|
||||||
|
"200": {
|
||||||
|
"cmdline": "/usr/local/go2rtc/bin/go2rtc -config=/dev/shm/go2rtc.yaml"
|
||||||
|
},
|
||||||
|
"300": {"cmdline": "frigate.recording"},
|
||||||
|
}
|
||||||
|
|
||||||
|
self.assertEqual(get_go2rtc_pid(cpu_usages), 200)
|
||||||
|
|
||||||
|
def test_finds_custom_go2rtc_binary(self):
|
||||||
|
self.assertEqual(get_go2rtc_pid({"42": {"cmdline": "/config/go2rtc"}}), 42)
|
||||||
|
|
||||||
|
def test_returns_none_when_go2rtc_is_not_running(self):
|
||||||
|
self.assertIsNone(get_go2rtc_pid({"100": {"cmdline": "ffmpeg -i x"}}))
|
||||||
|
self.assertIsNone(get_go2rtc_pid({}))
|
||||||
|
|
||||||
|
|
||||||
|
class TestGo2rtcPidInSnapshot(unittest.TestCase):
|
||||||
|
def snapshot(self, tracking: dict, go2rtc_pid: int) -> dict:
|
||||||
|
def update_stats(stats: dict) -> None:
|
||||||
|
stats["cpu_usages"] = {
|
||||||
|
str(go2rtc_pid): {
|
||||||
|
"cmdline": "/usr/local/go2rtc/bin/go2rtc -config=x",
|
||||||
|
"cpu": str(go2rtc_pid / 100),
|
||||||
|
"mem": str(go2rtc_pid / 10),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
config = SimpleNamespace(
|
||||||
|
cameras={},
|
||||||
|
telemetry=SimpleNamespace(stats=SimpleNamespace(network_bandwidth=False)),
|
||||||
|
)
|
||||||
|
hardware_stats = Mock()
|
||||||
|
hardware_stats.update_stats.side_effect = update_stats
|
||||||
|
|
||||||
|
with (
|
||||||
|
patch("frigate.stats.util.get_detector_stats", return_value={}),
|
||||||
|
patch("frigate.stats.util.embeddings_stats", return_value={}),
|
||||||
|
patch("frigate.stats.util.calculate_shm_requirements", return_value={}),
|
||||||
|
):
|
||||||
|
return stats_snapshot(config, tracking, hardware_stats)
|
||||||
|
|
||||||
|
def test_snapshot_follows_go2rtc_restart(self):
|
||||||
|
tracking = {
|
||||||
|
"camera_metrics": {},
|
||||||
|
"detectors": {},
|
||||||
|
"started": 0,
|
||||||
|
"latest_frigate_version": "",
|
||||||
|
"processes": {"go2rtc": 200, "recording": 50},
|
||||||
|
"storage_maintainer": None,
|
||||||
|
}
|
||||||
|
|
||||||
|
first = self.snapshot(tracking, 200)["processes"]["go2rtc"]
|
||||||
|
self.assertEqual(first, {"pid": 200, "cpu": "2.0", "mem": "20.0"})
|
||||||
|
|
||||||
|
restarted = self.snapshot(tracking, 300)["processes"]["go2rtc"]
|
||||||
|
self.assertEqual(restarted, {"pid": 300, "cpu": "3.0", "mem": "30.0"})
|
||||||
+47
-42
@@ -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,27 @@ 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:
|
||||||
|
"""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):
|
class CameraWatchdog(threading.Thread):
|
||||||
def __init__(
|
def __init__(
|
||||||
self,
|
self,
|
||||||
@@ -186,33 +208,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 +227,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 +393,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 +470,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 +482,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 +492,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 +502,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 +545,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 +561,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
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user