Compare commits

...
5 Commits
Author SHA1 Message Date
Martin WeineltandGitHub 48d09efd41 Merge b0588a02f9 into 48aaafba3c 2026-07-16 22:42:21 +02:00
Josh HawkinsandGitHub 48aaafba3c re-apply runtime overrides to the config without re-broadcasting them (#23739)
/api/config/set and camera deletion re-parse yaml into a fresh FrigateConfig and swap it in, then re-layered the persisted runtime toggle overrides so a camera the user turned off wouldn't come back on. That re-layer ran apply_runtime_state, which replays each override through the command handlers, so every save re-published a ZMQ config update, a retained MQTT state message, and a runtime-state disk write for every camera with a stored toggle. All of it was redundant: the worker processes were never swapped and still hold the live toggle values, so only the in-process config object the API and dispatcher read was out of date. The extra traffic churned the retained MQTT topics, amplified disk writes, and co-drained enabled updates with other topics on the config socket.

Add Dispatcher.reapply_runtime_state_to_config, which corrects only the swapped-in config object, mirroring the field mutations and gates of the _on_*_command handlers with no ZMQ, MQTT, or disk writes. swap_runtime_config now calls it instead of apply_runtime_state; apply_runtime_state is unchanged and still used at startup, where the workers genuinely must be told.
2026-07-16 13:30:41 -05:00
Josh HawkinsandGitHub f1028d0c36 Fix persisted runtime camera toggles (#23734)
* preserve runtime camera toggles across config saves

Runtime toggles (camera on/off, detect, recordings, snapshots, audio) mutate the in-memory config and persist an override to .runtime_state.json. /api/config/set re-parses yaml into a fresh FrigateConfig and swaps it in, re-applying the yaml and profile layers but dropping the runtime layer, so a camera turned off from the dashboard came back on when an unrelated camera was saved. The workers were never notified, so it only appeared to come back: the UI streamed go2rtc while ffmpeg stayed stopped.

Extract the startup replay into Dispatcher.apply_runtime_state() and call it from config_set after the swap, re-layering the overrides and republishing them so workers and the UI reconverge.

Remove the broad clear_runtime_state() from ProfileManager.update_config, which is only ever reached from config_set: with a profile active, every save wiped every camera's overrides from disk. The broad wipe stays in activate_profile, where a real profile switch does invalidate the steady state. Saves still clear the keys they rewrote via clear_runtime_state_for_yaml_keys, so yaml wins where the two disagree.

* sync runtime config on camera delete and prune its overrides

Deleting a camera re-parsed yaml into a fresh FrigateConfig but only rebound app.frigate_config and genai_manager, never dispatcher.config (nor profile_manager, stats_emitter, or the runtime overrides). The API and the dispatcher then drifted onto different config objects until the next config save re-synced them, so the API reported surviving cameras with their yaml enabled state while the dispatcher still acted on their real runtime state.

Extract the config swap that config_set already does into a shared swap_runtime_config helper and call it from both sites, so every collaborator is rebound and the surviving cameras' runtime toggles are re-layered. Also drop the deleted camera's persisted overrides via a new clear_camera so a camera later added under the same name does not inherit them.
2026-07-16 09:05:21 -05:00
Nicolas MowenandGitHub 70d629bf93 Update OpenVINO model generation (#23733)
CI / ARM Build (push) Waiting to run
CI / Jetson Jetpack 6 (push) Waiting to run
CI / AMD64 Extra Build (push) Blocked by required conditions
CI / ARM Extra Build (push) Blocked by required conditions
CI / Synaptics Build (push) Blocked by required conditions
CI / Assemble and push default build (push) Blocked by required conditions
CI / AMD64 Build (push) Waiting to run
2026-07-16 08:00:52 -06:00
Martin Weinelt b0588a02f9 Replace blocking I/O in async functions
Replaces aiofiles with anyio, because anyio.Path is much more complete
and comparable to the Pathlib API.
2026-07-06 22:24:19 +02:00
22 changed files with 799 additions and 113 deletions
+2 -2
View File
@@ -81,10 +81,10 @@ RUN --mount=type=bind,source=docker/main/install_tempio.sh,target=/deps/install_
FROM base_host AS ov-converter
ARG DEBIAN_FRONTEND
# Install OpenVino Runtime and Dev library
# Install OpenVINO for model conversion
COPY docker/main/requirements-ov.txt /requirements-ov.txt
RUN apt-get -qq update \
&& apt-get -qq install -y wget python3 python3-dev python3-distutils gcc pkg-config libhdf5-dev \
&& apt-get -qq install -y wget python3 python3-distutils \
&& wget -q https://bootstrap.pypa.io/get-pip.py -O get-pip.py \
&& sed -i 's/args.append("setuptools")/args.append("setuptools==77.0.3")/' get-pip.py \
&& python3 get-pip.py "pip" \
+39 -8
View File
@@ -1,11 +1,42 @@
import openvino as ov
from openvino.tools import mo
"""Convert the default SSDLite MobileNet v2 model to OpenVINO IR.
ov_model = mo.convert_model(
Replaces the legacy openvino-dev Model Optimizer conversion. The TensorFlow
frontend converts the Object Detection API frozen graph natively; the four TF
outputs are then repacked into the single [1, 1, 100, 7] DetectionOutput-style
tensor that Frigate's OpenVINO detector expects, and the input is flipped to
BGR to match the legacy reverse_input_channels behavior.
"""
import numpy as np
import openvino as ov
from openvino import opset8 as ops
from openvino.preprocess import PrePostProcessor
model = ov.convert_model(
"/models/ssdlite_mobilenet_v2_coco_2018_05_09/frozen_inference_graph.pb",
compress_to_fp16=True,
transformations_config="/usr/local/lib/python3.11/dist-packages/openvino/tools/mo/front/tf/ssd_v2_support.json",
tensorflow_object_detection_api_pipeline_config="/models/ssdlite_mobilenet_v2_coco_2018_05_09/pipeline.config",
reverse_input_channels=True,
input=[("image_tensor:0", [1, 300, 300, 3])],
)
ov.save_model(ov_model, "/models/ssdlite_mobilenet_v2.xml")
# rows of (image_id, class_id, score, xmin, ymin, xmax, ymax)
boxes = model.output("detection_boxes:0").get_node().input_value(0)
classes = model.output("detection_classes:0").get_node().input_value(0)
scores = model.output("detection_scores:0").get_node().input_value(0)
# (ymin,xmin,ymax,xmax) -> (xmin,ymin,xmax,ymax)
boxes = ops.gather(boxes, [1, 0, 3, 2], 2)
classes = ops.unsqueeze(classes, 2)
scores = ops.unsqueeze(scores, 2)
image_id = ops.multiply(scores, np.float32(0.0))
detections = ops.concat([image_id, classes, scores, boxes], 2)
detections = ops.unsqueeze(detections, 1)
detections.output(0).get_tensor().set_names({"detection_out"})
model = ov.Model([detections], model.get_parameters(), "ssdlite_mobilenet_v2")
ppp = PrePostProcessor(model)
ppp.input().tensor().set_layout(ov.Layout("NHWC"))
ppp.input().preprocess().reverse_channels()
model = ppp.build()
ov.save_model(model, "/models/ssdlite_mobilenet_v2.xml", compress_to_fp16=True)
+1 -2
View File
@@ -1,3 +1,2 @@
numpy
tensorflow
openvino-dev>=2024.0.0
openvino >= 2026.2.0
+1 -1
View File
@@ -1,4 +1,4 @@
aiofiles == 24.1.*
anyio == 4.14.*
click == 8.1.*
# FastAPI
aiohttp == 3.12.*
+5 -16
View File
@@ -14,8 +14,8 @@ from io import StringIO
from pathlib import Path as FilePath
from typing import Any
import aiofiles
import ruamel.yaml
from anyio import open_file as aopen
from fastapi import APIRouter, Body, Path, Request, Response
from fastapi.encoders import jsonable_encoder
from fastapi.params import Depends
@@ -31,6 +31,7 @@ from frigate.api.auth import (
get_allowed_cameras_for_filter,
require_role,
)
from frigate.api.config_util import swap_runtime_config
from frigate.api.defs.query.app_query_parameters import AppTimelineHourlyQueryParameters
from frigate.api.defs.request.app_body import (
AppConfigSetBody,
@@ -915,19 +916,7 @@ def config_set(request: Request, body: AppConfigSetBody):
if body.requires_restart == 0 or body.update_topic:
old_config: FrigateConfig = request.app.frigate_config
request.app.frigate_config = config
request.app.genai_manager.update_config(config)
if request.app.profile_manager is not None:
request.app.profile_manager.update_config(config)
if request.app.stats_emitter is not None:
request.app.stats_emitter.config = config
if request.app.dispatcher is not None:
request.app.dispatcher.config = config
for comm in request.app.dispatcher.comms:
comm.config = config
swap_runtime_config(request.app, config)
if body.update_topic:
if body.update_topic.startswith("config/cameras/"):
@@ -1056,7 +1045,7 @@ async def logs(
"""Asynchronously stream log lines."""
buffer = ""
try:
async with aiofiles.open(file_path) as file:
async with await aopen(file_path) as file:
await file.seek(0, 2)
while True:
line = await file.readline()
@@ -1094,7 +1083,7 @@ async def logs(
# For full logs initially
try:
async with aiofiles.open(service_location) as file:
async with await aopen(service_location) as file:
contents = await file.read()
total_lines, log_lines = process_logs(contents, service, start, end)
+21 -12
View File
@@ -10,6 +10,7 @@ from urllib.parse import quote_plus
import httpx
import requests
from anyio import open_file as aopen
from fastapi import APIRouter, Depends, Query, Request, Response
from fastapi.responses import JSONResponse
from filelock import FileLock, Timeout
@@ -25,6 +26,7 @@ from frigate.api.auth import (
require_go2rtc_stream_access,
require_role,
)
from frigate.api.config_util import swap_runtime_config
from frigate.api.defs.request.app_body import CameraSetBody
from frigate.api.defs.tags import Tags
from frigate.config import FrigateConfig
@@ -1187,15 +1189,17 @@ async def delete_camera(
try:
with lock:
with open(config_file) as f:
old_raw_config = f.read()
async with await aopen(config_file) as f:
old_raw_config = await f.read()
try:
yaml = YAML()
yaml.indent(mapping=2, sequence=4, offset=2)
with open(config_file) as f:
data = yaml.load(f)
async with await aopen(config_file) as f:
text = await f.read()
data = yaml.load(text)
# Remove camera from config
if "cameras" in data and camera_name in data["cameras"]:
@@ -1220,17 +1224,17 @@ async def delete_camera(
for role_name in empty_roles:
del auth["roles"][role_name]
with open(config_file, "w") as f:
async with await aopen(config_file, "w") as f:
yaml.dump(data, f)
with open(config_file) as f:
new_raw_config = f.read()
async with await aopen(config_file) as f:
new_raw_config = await f.read()
try:
config = FrigateConfig.parse(new_raw_config)
except Exception:
with open(config_file, "w") as f:
f.write(old_raw_config)
async with await aopen(config_file, "w") as f:
await f.write(old_raw_config)
logger.exception(
"Config error after removing camera %s",
camera_name,
@@ -1254,9 +1258,14 @@ async def delete_camera(
status_code=500,
)
# Update runtime config
request.app.frigate_config = config
request.app.genai_manager.update_config(config)
# rebind every collaborator to the new config and re-layer runtime
# toggles for the surviving cameras, same as /api/config/set
swap_runtime_config(request.app, config)
# drop the deleted camera's persisted overrides so a camera later
# added under the same name doesn't inherit them
if request.app.dispatcher is not None:
request.app.dispatcher.clear_runtime_state_for_camera(camera_name)
# Publish removal to stop ffmpeg processes and clean up runtime state
request.app.config_publisher.publish_update(
+35
View File
@@ -0,0 +1,35 @@
"""Shared helpers for applying a freshly parsed config to the running app."""
from fastapi import FastAPI
from frigate.config import FrigateConfig
def swap_runtime_config(app: FastAPI, config: FrigateConfig) -> None:
"""Point every long-lived collaborator at a newly parsed config object.
Both /api/config/set and camera deletion re-parse yaml into a fresh
FrigateConfig and must rebind the same set of references, or the API and
the dispatcher drift onto different objects (the API reports one camera
state while the dispatcher acts on another). Runtime toggle overrides are
re-layered last: the swap rebuilt every camera from yaml, so without this a
camera the user turned off would silently come back on.
"""
app.frigate_config = config
app.genai_manager.update_config(config)
if app.profile_manager is not None:
app.profile_manager.update_config(config)
if app.stats_emitter is not None:
app.stats_emitter.config = config
if app.dispatcher is not None:
app.dispatcher.config = config
for comm in app.dispatcher.comms:
comm.config = config
# workers still hold the live toggle values, so correct only the
# config object here rather than re-broadcasting every override
app.dispatcher.reapply_runtime_state_to_config()
+3 -2
View File
@@ -13,6 +13,7 @@ from pathlib import Path
from urllib.parse import unquote
import numpy as np
from anyio import Path as AsyncPath
from fastapi import APIRouter, Request
from fastapi.params import Depends
from fastapi.responses import JSONResponse
@@ -1455,10 +1456,10 @@ async def set_attributes(
dataset_dir = os.path.join(CLIPS_DIR, sanitize_filename(model_key), "dataset")
available_labels = set()
if os.path.exists(dataset_dir):
if await AsyncPath(dataset_dir).exists():
for category_name in os.listdir(dataset_dir):
category_dir = os.path.join(dataset_dir, category_name)
if os.path.isdir(category_dir):
if await AsyncPath(category_dir).is_dir():
available_labels.add(category_name)
if not available_labels:
+13 -11
View File
@@ -15,6 +15,8 @@ from urllib.parse import unquote
import cv2
import numpy as np
import pytz
from anyio import Path as AsyncPath
from anyio import open_file as aopen
from fastapi import APIRouter, Depends, Path, Query, Request, Response
from fastapi.responses import FileResponse, JSONResponse, StreamingResponse
from pathvalidate import sanitize_filename
@@ -497,18 +499,18 @@ async def recording_clip(
file_name = sanitize_filename(f"playlist_{camera_name}_{start_ts}-{end_ts}.txt")
file_path = os.path.join(CACHE_DIR, file_name)
with open(file_path, "w") as file:
async with await aopen(file_path, "w") as file:
clip: Recordings
for clip in recordings:
file.write(f"file '{clip.path}'\n")
await file.write(f"file '{clip.path}'\n")
# if this is the starting clip, add an inpoint
if clip.start_time < start_ts:
file.write(f"inpoint {int(start_ts - clip.start_time)}\n")
await file.write(f"inpoint {int(start_ts - clip.start_time)}\n")
# if this is the ending clip, add an outpoint
if clip.end_time > end_ts:
file.write(f"outpoint {int(end_ts - clip.start_time)}\n")
await file.write(f"outpoint {int(end_ts - clip.start_time)}\n")
if len(file_name) > 1000:
return JSONResponse(
@@ -1149,8 +1151,8 @@ async def event_snapshot_clean(request: Request, event_id: str, download: bool =
)
if image_path.endswith(".webp"):
with open(image_path, "rb") as image_file:
webp_bytes = image_file.read()
async with await aopen(image_path, "rb") as image_file:
webp_bytes = await image_file.read()
else:
image = load_event_snapshot_image(event, clean_only=True)[0]
if image is None:
@@ -1366,7 +1368,7 @@ async def preview_gif(
# need to generate from existing images
preview_dir = os.path.join(CACHE_DIR, "preview_frames")
if not os.path.isdir(preview_dir):
if not await AsyncPath(preview_dir).is_dir():
return JSONResponse(
content={"success": False, "message": "Preview not found"},
status_code=404,
@@ -1555,7 +1557,7 @@ async def preview_mp4(
# need to generate from existing images
preview_dir = os.path.join(CACHE_DIR, "preview_frames")
if not os.path.isdir(preview_dir):
if not await AsyncPath(preview_dir).is_dir():
return JSONResponse(
content={"success": False, "message": "Preview not found"},
status_code=404,
@@ -1633,7 +1635,7 @@ async def preview_mp4(
"Content-Description": "File Transfer",
"Cache-Control": f"private, max-age={_resolve_cache_age(max_cache_age)}",
"Content-Type": "video/mp4",
"Content-Length": str(os.path.getsize(path)),
"Content-Length": str((await AsyncPath(path).stat()).st_size),
# nginx: https://nginx.org/en/docs/http/ngx_http_proxy_module.html#proxy_ignore_headers
"X-Accel-Redirect": f"/cache/{file_name}",
}
@@ -1707,10 +1709,10 @@ async def preview_thumbnail(request: Request, file_name: str):
preview_dir = os.path.join(CACHE_DIR, "preview_frames")
try:
with open(
async with await aopen(
os.path.join(preview_dir, safe_file_name_current), "rb"
) as image_file:
jpg_bytes = image_file.read()
jpg_bytes = await image_file.read()
except FileNotFoundError:
return JSONResponse(
content=({"success": False, "message": "Image file not found"}),
+2 -2
View File
@@ -4,9 +4,9 @@ import datetime as dt
import logging
from datetime import datetime, timedelta
from functools import reduce
from pathlib import Path
from urllib.parse import unquote
from anyio import Path as AsyncPath
from fastapi import APIRouter, Depends, Request
from fastapi import Path as PathParam
from fastapi.responses import JSONResponse
@@ -443,7 +443,7 @@ async def delete_recordings(
recording_ids.append(recording["id"])
try:
Path(recording["path"]).unlink(missing_ok=True)
await AsyncPath(recording["path"]).unlink(missing_ok=True)
deleted_count += 1
except Exception as e:
logger.error(f"Failed to delete recording file {recording['path']}: {e}")
+83 -7
View File
@@ -404,38 +404,64 @@ class Dispatcher:
for comm in self.comms:
comm.stop()
def restore_runtime_state(self) -> None:
def apply_runtime_state(self) -> dict[str, dict[str, bool]]:
"""Replay persisted runtime overrides through the camera settings handlers.
Called once after Frigate startup completes so processing threads can
receive the resulting ``config_updater`` broadcasts. Unknown cameras
and topics are skipped; handler exceptions are logged and replay
continues for remaining entries.
Routing through the handlers (rather than mutating config directly) is
deliberate: they publish the ``config_updater`` broadcast and the
retained MQTT state as a side effect, so worker processes and the UI
converge on the replayed value. Unknown cameras and topics are skipped;
handler exceptions are logged and replay continues for the rest.
Returns:
The entries handed to a handler without raising, keyed by camera
then topic. A handler can still refuse the value internally (an ON
payload for a camera that is not enabled_in_config, for example),
so this is not proof the override took effect.
"""
state = self._runtime_state.load()
applied: dict[str, dict[str, bool]] = {}
for camera_name, features in state.items():
if camera_name not in self.config.cameras:
continue
for topic, value in features.items():
handler = self._camera_settings_handlers.get(topic)
if handler is None:
continue
payload = "ON" if value else "OFF"
try:
handler(camera_name, payload)
except Exception:
logger.exception(
"Failed to restore runtime state %s.%s=%s",
"Failed to apply runtime state %s.%s=%s",
camera_name,
topic,
payload,
)
continue
applied.setdefault(camera_name, {})[topic] = value
return applied
def restore_runtime_state(self) -> None:
"""Replay persisted runtime overrides once Frigate startup completes.
Called after every ``config_updater`` subscriber is up so the resulting
broadcasts are not dropped by ZMQ PUB/SUB.
"""
for camera_name, features in self.apply_runtime_state().items():
for topic, value in features.items():
logger.info(
"Restored runtime state: %s.%s=%s",
camera_name,
topic,
payload,
"ON" if value else "OFF",
)
def clear_runtime_state_for_yaml_keys(self, dotted_keys: Iterable[str]) -> None:
@@ -458,6 +484,56 @@ class Dispatcher:
"""
self._runtime_state.clear_all()
def clear_runtime_state_for_camera(self, camera: str) -> None:
"""Drop all persisted runtime overrides for a deleted camera.
Called by camera deletion so a camera later added under the same name
does not inherit the removed camera's stale toggles.
"""
self._runtime_state.clear_camera(camera)
def reapply_runtime_state_to_config(self) -> None:
"""Re-apply persisted runtime overrides to the swapped-in config object.
After config/set (or a camera delete) parses fresh yaml and swaps the
config, the worker processes still hold the live toggle values and the
overrides are already on disk, so only the in-process config object is
out of date. Unlike apply_runtime_state (used at startup, where workers
must be told), this makes no ZMQ, MQTT, or disk writes, it just corrects
the config the API and dispatcher read.
The field mutations and gates mirror the _on_*_command handlers; keep
the two in sync if a tracked toggle is added or its gate changes.
"""
state = self._runtime_state.load()
for camera_name, features in state.items():
camera = self.config.cameras.get(camera_name)
if camera is None:
continue
for topic, value in features.items():
if topic == "enabled":
if value and not camera.enabled_in_config:
continue
camera.enabled = value
elif topic == "detect":
camera.detect.enabled = value
# detection requires motion, mirror the handler coupling
if value and not camera.motion.enabled:
camera.motion.enabled = True
elif topic == "snapshots":
camera.snapshots.enabled = value
elif topic == "recordings":
if value and not camera.record.enabled_in_config:
continue
camera.record.enabled = value
elif topic == "audio":
if value and not camera.audio.enabled_in_config:
continue
camera.audio.enabled = value
def _on_detect_command(self, camera_name: str, payload: str) -> None:
"""Callback for detect topic."""
detect_settings = self.config.cameras[camera_name].detect
+19
View File
@@ -96,6 +96,25 @@ class RuntimeStatePersistence:
except OSError:
logger.exception("Failed to clear runtime state")
def clear_camera(self, camera: str) -> None:
"""Drop every stored override for a single camera.
Called when a camera is deleted so a camera later added under the same
name does not inherit the removed camera's stale toggles.
"""
try:
with FileLock(self._lock_path, timeout=self._lock_timeout):
data = self._read_locked()
cameras = data.get("cameras")
if not isinstance(cameras, dict) or camera not in cameras:
return
del cameras[camera]
self._write_locked(data)
except Timeout:
logger.error("Timed out clearing runtime state for camera")
except OSError:
logger.exception("Failed to clear runtime state for camera")
def clear_for_yaml_keys(self, dotted_keys: Iterable[str]) -> None:
"""Remove stored entries whose YAML key was just rewritten.
+5 -4
View File
@@ -141,6 +141,11 @@ class ProfileManager:
Preserves active profile state: re-snapshots base configs from the new
(freshly parsed) config, then re-applies profile overrides if a profile
was active.
Deliberately does not clear the dispatcher's runtime overrides. This is
the config-save path, not a profile switch: the save only invalidates
the toggles it rewrote in yaml, which /api/config/set already clears by
key. The broad wipe belongs to activate_profile alone.
"""
current_active = self.config.active_profile
self.config = new_config
@@ -164,10 +169,6 @@ class ProfileManager:
self.config.active_profile = None
self._persist_active_profile(None)
# drop all runtime overrides so they don't replay stale values on restart
if self.dispatcher is not None:
self.dispatcher.clear_runtime_state()
def activate_profile(
self,
profile_name: str | None,
+13 -10
View File
@@ -15,6 +15,7 @@ from typing import Any
import numpy as np
import psutil
from anyio import Path as AsyncPath
from frigate.comms.detections_updater import DetectionSubscriber, DetectionTypeEnum
from frigate.comms.inter_process import InterProcessRequestor
@@ -105,11 +106,11 @@ class RecordingMaintainer(threading.Thread):
async def move_files(self) -> None:
cache_files = [
d
for d in os.listdir(CACHE_DIR)
if os.path.isfile(os.path.join(CACHE_DIR, d))
and d.endswith(".mp4")
and not d.startswith("preview_")
path.name
async for path in AsyncPath(CACHE_DIR).iterdir()
if await path.is_file()
and path.suffix == ".mp4"
and not path.name.startswith("preview_")
]
# publish newest cached segment per camera (including in use files)
@@ -229,7 +230,7 @@ class RecordingMaintainer(threading.Thread):
to_remove = grouped_recordings[camera][:-keep_count]
for rec in to_remove:
cache_path = rec["cache_path"]
Path(cache_path).unlink(missing_ok=True)
await AsyncPath(cache_path).unlink(missing_ok=True)
self.end_time_cache.pop(cache_path, None)
grouped_recordings[camera] = grouped_recordings[camera][-keep_count:]
@@ -244,7 +245,7 @@ class RecordingMaintainer(threading.Thread):
to_remove = grouped_recordings[camera][:-keep_count]
for rec in to_remove:
cache_path = rec["cache_path"]
Path(cache_path).unlink(missing_ok=True)
await AsyncPath(cache_path).unlink(missing_ok=True)
self.end_time_cache.pop(cache_path, None)
grouped_recordings[camera] = grouped_recordings[camera][-keep_count:]
@@ -634,7 +635,7 @@ class RecordingMaintainer(threading.Thread):
file_path = os.path.join(directory, file_name)
try:
if not os.path.exists(file_path):
if not await AsyncPath(file_path).exists():
start_frame = datetime.datetime.now().timestamp()
# add faststart to kept segments to improve metadata reading
@@ -670,7 +671,9 @@ class RecordingMaintainer(threading.Thread):
# get the segment size of the cache file
# file without faststart is same size
segment_size = round(
float(os.path.getsize(cache_path)) / pow(2, 20), 2
float((await AsyncPath(cache_path).stat()).st_size)
/ pow(2, 20),
2,
)
except OSError:
segment_size = 0
@@ -698,7 +701,7 @@ class RecordingMaintainer(threading.Thread):
}
except Exception as e:
logger.error(f"Unable to store recording segment {cache_path}")
Path(cache_path).unlink(missing_ok=True)
await AsyncPath(cache_path).unlink(missing_ok=True)
logger.error(e)
# clear end_time cache
+132
View File
@@ -0,0 +1,132 @@
"""Tests for the camera delete endpoint's runtime config handling."""
import os
import tempfile
import unittest
from unittest.mock import MagicMock, Mock, patch
import ruamel.yaml
from frigate.config import FrigateConfig
from frigate.config.camera.updater import CameraConfigUpdatePublisher
from frigate.models import Event, Recordings, ReviewSegment
from frigate.test.http_api.base_http_test import AuthTestClient, BaseTestHttp
class TestDeleteCameraRuntimeConfig(BaseTestHttp):
"""Deleting a camera must keep the API and dispatcher on the same config."""
def setUp(self):
super().setUp(models=[Event, Recordings, ReviewSegment])
self.minimal_config = {
"mqtt": {"host": "mqtt"},
"cameras": {
"front_door": {
"ffmpeg": {
"inputs": [
{"path": "rtsp://10.0.0.1:554/video", "roles": ["detect"]}
]
},
"detect": {"height": 1080, "width": 1920, "fps": 5},
},
"back_yard": {
"ffmpeg": {
"inputs": [
{"path": "rtsp://10.0.0.2:554/video", "roles": ["detect"]}
]
},
"detect": {"height": 720, "width": 1280, "fps": 10},
},
},
}
def _write_config_file(self):
yaml = ruamel.yaml.YAML()
f = tempfile.NamedTemporaryFile(mode="w", suffix=".yml", delete=False)
yaml.dump(self.minimal_config, f)
f.close()
return f.name
def _create_app_with_dispatcher(self, dispatcher):
from fastapi import Request
from frigate.api.auth import get_allowed_cameras_for_filter, get_current_user
from frigate.api.fastapi_app import create_fastapi_app
mock_publisher = Mock(spec=CameraConfigUpdatePublisher)
mock_publisher.publisher = MagicMock()
app = create_fastapi_app(
FrigateConfig(**self.minimal_config),
self.db,
None,
None,
None,
None,
None,
None,
mock_publisher,
None,
dispatcher=dispatcher,
enforce_default_admin=False,
)
async def mock_get_current_user(request: Request):
return {
"username": request.headers.get("remote-user"),
"role": request.headers.get("remote-role"),
}
async def mock_get_allowed_cameras_for_filter(request: Request):
return list(self.minimal_config.get("cameras", {}).keys())
app.dependency_overrides[get_current_user] = mock_get_current_user
app.dependency_overrides[get_allowed_cameras_for_filter] = (
mock_get_allowed_cameras_for_filter
)
return app, mock_publisher
@patch("frigate.api.camera.requests.delete")
@patch("frigate.api.camera.cleanup_camera_files")
@patch("frigate.api.camera.cleanup_camera_db")
@patch("frigate.api.camera.find_config_file")
def test_delete_syncs_dispatcher_and_prunes_runtime_state(
self, mock_find_config, mock_cleanup_db, mock_cleanup_files, mock_go2rtc_delete
):
"""Deleting a camera swaps every config reference and prunes its state."""
config_path = self._write_config_file()
mock_find_config.return_value = config_path
mock_cleanup_db.return_value = ({}, [])
dispatcher = MagicMock()
dispatcher.comms = []
try:
app, _ = self._create_app_with_dispatcher(dispatcher)
with AuthTestClient(app) as client:
resp = client.delete("/cameras/front_door")
self.assertEqual(resp.status_code, 200)
self.assertTrue(resp.json()["success"])
# the dispatcher must be moved onto the same new object the API
# now serves, and that object must no longer contain the camera
self.assertIs(dispatcher.config, app.frigate_config)
self.assertNotIn("front_door", dispatcher.config.cameras)
self.assertIn("back_yard", dispatcher.config.cameras)
# surviving cameras' overrides are re-layered onto the new object
dispatcher.reapply_runtime_state_to_config.assert_called_once_with()
# the deleted camera's persisted overrides are pruned
dispatcher.clear_runtime_state_for_camera.assert_called_once_with(
"front_door"
)
finally:
os.unlink(config_path)
if __name__ == "__main__":
unittest.main()
@@ -91,6 +91,123 @@ class TestConfigSetWildcardPropagation(BaseTestHttp):
return app, mock_publisher
def _create_app_with_dispatcher(self, dispatcher):
"""Create app with a mocked config publisher and a real-ish dispatcher."""
from fastapi import Request
from frigate.api.auth import get_allowed_cameras_for_filter, get_current_user
from frigate.api.fastapi_app import create_fastapi_app
mock_publisher = Mock(spec=CameraConfigUpdatePublisher)
mock_publisher.publisher = MagicMock()
app = create_fastapi_app(
FrigateConfig(**self.minimal_config),
self.db,
None,
None,
None,
None,
None,
None,
mock_publisher,
None,
dispatcher=dispatcher,
enforce_default_admin=False,
)
async def mock_get_current_user(request: Request):
username = request.headers.get("remote-user")
role = request.headers.get("remote-role")
return {"username": username, "role": role}
async def mock_get_allowed_cameras_for_filter(request: Request):
return list(self.minimal_config.get("cameras", {}).keys())
app.dependency_overrides[get_current_user] = mock_get_current_user
app.dependency_overrides[get_allowed_cameras_for_filter] = (
mock_get_allowed_cameras_for_filter
)
return app, mock_publisher
@patch("frigate.api.app.find_config_file")
def test_runtime_disabled_camera_survives_unrelated_save(self, mock_find_config):
"""A camera turned off at runtime stays off when another camera is saved."""
config_path = self._write_config_file()
mock_find_config.return_value = config_path
dispatcher = MagicMock()
dispatcher.comms = []
# front_door was turned off via the UI: the override is on disk, and
# yaml still says enabled: true. Stand in for the real replay, which
# reads dispatcher.config - the object the endpoint just swapped in.
def fake_reapply():
dispatcher.config.cameras["front_door"].enabled = False
dispatcher.reapply_runtime_state_to_config.side_effect = fake_reapply
try:
app, _ = self._create_app_with_dispatcher(dispatcher)
with AuthTestClient(app) as client:
resp = client.put(
"/config/set",
json={
"config_data": {
"cameras": {"back_yard": {"detect": {"fps": 7}}}
},
"requires_restart": 0,
},
)
self.assertEqual(resp.status_code, 200)
self.assertTrue(resp.json()["success"])
# the swap must be repaired: the new config object the API and
# dispatcher now share has to still show front_door as off
dispatcher.reapply_runtime_state_to_config.assert_called_once_with()
self.assertFalse(app.frigate_config.cameras["front_door"].enabled)
self.assertIs(dispatcher.config, app.frigate_config)
# yaml-wins ordering: the surgical clear for rewritten keys
# must run before the replay, or a save that rewrote a toggle
# would have its old override resurrected
call_names = [name for name, _, _ in dispatcher.mock_calls]
self.assertLess(
call_names.index("clear_runtime_state_for_yaml_keys"),
call_names.index("reapply_runtime_state_to_config"),
)
finally:
os.unlink(config_path)
@patch("frigate.api.app.find_config_file")
def test_no_reapply_when_config_is_not_swapped(self, mock_find_config):
"""A restart-required save with no update topic never swaps, so no replay."""
config_path = self._write_config_file()
mock_find_config.return_value = config_path
dispatcher = MagicMock()
dispatcher.comms = []
try:
app, _ = self._create_app_with_dispatcher(dispatcher)
with AuthTestClient(app) as client:
resp = client.put(
"/config/set",
json={
"config_data": {"mqtt": {"host": "other"}},
"requires_restart": 1,
},
)
self.assertEqual(resp.status_code, 200)
dispatcher.reapply_runtime_state_to_config.assert_not_called()
finally:
os.unlink(config_path)
def _write_config_file(self):
"""Write the minimal config to a temp YAML file and return the path."""
yaml = ruamel.yaml.YAML()
+55
View File
@@ -0,0 +1,55 @@
"""Tests for the shared runtime config swap helper."""
import unittest
from unittest.mock import MagicMock
from frigate.api.config_util import swap_runtime_config
class TestSwapRuntimeConfig(unittest.TestCase):
"""swap_runtime_config rebinds every collaborator to the new config."""
def _make_app(self) -> MagicMock:
app = MagicMock()
app.dispatcher.comms = [MagicMock(), MagicMock()]
return app
def test_rebinds_all_references(self) -> None:
app = self._make_app()
config = MagicMock(name="new_config")
swap_runtime_config(app, config)
self.assertIs(app.frigate_config, config)
app.genai_manager.update_config.assert_called_once_with(config)
app.profile_manager.update_config.assert_called_once_with(config)
self.assertIs(app.stats_emitter.config, config)
self.assertIs(app.dispatcher.config, config)
for comm in app.dispatcher.comms:
self.assertIs(comm.config, config)
def test_reapplies_runtime_state_after_swap(self) -> None:
app = self._make_app()
config = MagicMock(name="new_config")
swap_runtime_config(app, config)
# the swap rebuilds cameras from yaml, so overrides must be re-layered
app.dispatcher.reapply_runtime_state_to_config.assert_called_once_with()
def test_tolerates_missing_optional_collaborators(self) -> None:
app = MagicMock()
app.profile_manager = None
app.stats_emitter = None
app.dispatcher = None
config = MagicMock(name="new_config")
# must not raise when the optional collaborators are absent
swap_runtime_config(app, config)
self.assertIs(app.frigate_config, config)
app.genai_manager.update_config.assert_called_once_with(config)
if __name__ == "__main__":
unittest.main()
@@ -126,6 +126,40 @@ class TestRestoreRuntimeState(unittest.TestCase):
self.dispatcher.restore_runtime_state()
self.handler_mocks["detect"].assert_called_once_with("front_door", "ON")
def test_apply_runtime_state_replays_through_handlers(self) -> None:
"""The extracted method replays every stored entry."""
with patch.object(
self.dispatcher._runtime_state,
"load",
return_value={"front_door": {"enabled": False, "detect": True}},
):
self.dispatcher.apply_runtime_state()
self.handler_mocks["enabled"].assert_called_once_with("front_door", "OFF")
self.handler_mocks["detect"].assert_called_once_with("front_door", "ON")
def test_apply_runtime_state_returns_applied_entries(self) -> None:
"""Callers get back what was replayed, for logging and assertions."""
with patch.object(
self.dispatcher._runtime_state,
"load",
return_value={"front_door": {"enabled": False}, "nope": {"enabled": True}},
):
applied = self.dispatcher.apply_runtime_state()
self.assertEqual(applied, {"front_door": {"enabled": False}})
def test_restore_runtime_state_still_replays(self) -> None:
"""The startup entry point keeps working after the extraction."""
with patch.object(
self.dispatcher._runtime_state,
"load",
return_value={"back_yard": {"snapshots": False}},
):
self.dispatcher.restore_runtime_state()
self.handler_mocks["snapshots"].assert_called_once_with("back_yard", "OFF")
class TestHandlersPersistViaSet(unittest.TestCase):
"""Verify each in-scope handler writes to the runtime state on success."""
@@ -212,6 +246,122 @@ class TestClearPassthrough(unittest.TestCase):
dispatcher.clear_runtime_state()
dispatcher._runtime_state.clear_all.assert_called_once_with()
def test_clear_runtime_state_for_camera_passthrough(self) -> None:
dispatcher = _build_dispatcher({})
dispatcher._runtime_state = MagicMock(spec=RuntimeStatePersistence)
dispatcher.clear_runtime_state_for_camera("front_door")
dispatcher._runtime_state.clear_camera.assert_called_once_with("front_door")
class TestReapplyRuntimeStateToConfig(unittest.TestCase):
"""The silent re-apply corrects the config object with no side effects."""
def _dispatcher_with(
self, cameras: dict[str, MagicMock], state: dict
) -> Dispatcher:
dispatcher = _build_dispatcher(cameras)
dispatcher._runtime_state = MagicMock(spec=RuntimeStatePersistence)
dispatcher._runtime_state.load.return_value = state
dispatcher.publish = MagicMock()
return dispatcher
def test_mutates_every_tracked_field(self) -> None:
cameras = {"front_door": _make_camera_mock()}
dispatcher = self._dispatcher_with(
cameras,
{
"front_door": {
"enabled": False,
"detect": False,
"snapshots": False,
"recordings": False,
"audio": False,
}
},
)
dispatcher.reapply_runtime_state_to_config()
cam = cameras["front_door"]
self.assertFalse(cam.enabled)
self.assertFalse(cam.detect.enabled)
self.assertFalse(cam.snapshots.enabled)
self.assertFalse(cam.record.enabled)
self.assertFalse(cam.audio.enabled)
def test_makes_no_zmq_mqtt_or_disk_writes(self) -> None:
dispatcher = self._dispatcher_with(
{"front_door": _make_camera_mock()},
{"front_door": {"enabled": False}},
)
dispatcher.reapply_runtime_state_to_config()
dispatcher.config_updater.publish_update.assert_not_called()
dispatcher._runtime_state.set.assert_not_called()
dispatcher.publish.assert_not_called()
def test_respects_enabled_in_config_gate(self) -> None:
# an ON override for a camera disabled in yaml must not enable it
cameras = {
"front_door": _make_camera_mock(enabled=False, enabled_in_config=False)
}
dispatcher = self._dispatcher_with(cameras, {"front_door": {"enabled": True}})
dispatcher.reapply_runtime_state_to_config()
self.assertFalse(cameras["front_door"].enabled)
def test_respects_recordings_and_audio_gates(self) -> None:
# ON overrides for recordings/audio not enabled in yaml must be ignored
cameras = {
"front_door": _make_camera_mock(
record_enabled=False,
record_enabled_in_config=False,
audio_enabled=False,
audio_enabled_in_config=False,
)
}
dispatcher = self._dispatcher_with(
cameras, {"front_door": {"recordings": True, "audio": True}}
)
dispatcher.reapply_runtime_state_to_config()
self.assertFalse(cameras["front_door"].record.enabled)
self.assertFalse(cameras["front_door"].audio.enabled)
def test_applies_on_override_when_gate_passes(self) -> None:
# a camera off in yaml but enabled_in_config keeps its runtime-on state
cameras = {
"front_door": _make_camera_mock(enabled=False, enabled_in_config=True)
}
dispatcher = self._dispatcher_with(cameras, {"front_door": {"enabled": True}})
dispatcher.reapply_runtime_state_to_config()
self.assertTrue(cameras["front_door"].enabled)
def test_detect_on_couples_motion(self) -> None:
cam = _make_camera_mock(detect_enabled=False)
cam.motion.enabled = False
dispatcher = self._dispatcher_with(
{"front_door": cam}, {"front_door": {"detect": True}}
)
dispatcher.reapply_runtime_state_to_config()
self.assertTrue(cam.detect.enabled)
self.assertTrue(cam.motion.enabled)
def test_skips_camera_not_in_config(self) -> None:
dispatcher = self._dispatcher_with(
{"front_door": _make_camera_mock()}, {"ghost": {"enabled": False}}
)
# a stale entry for a deleted camera must be ignored, not raise
dispatcher.reapply_runtime_state_to_config()
if __name__ == "__main__":
unittest.main()
+44 -31
View File
@@ -1,7 +1,7 @@
import datetime
import sys
import unittest
from unittest.mock import MagicMock, patch
from unittest.mock import AsyncMock, MagicMock, patch
# Mock complex imports before importing maintainer, saving originals so we can
# restore them after import and avoid polluting sys.modules for other tests.
@@ -42,38 +42,51 @@ class TestMaintainer(unittest.IsolatedAsyncioTestCase):
# One bad file, one good file
files = ["bad_filename.mp4", "camera@20210101000000+0000.mp4"]
with patch("os.listdir", return_value=files):
with patch("os.path.isfile", return_value=True):
with patch(
"frigate.record.maintainer.psutil.process_iter", return_value=[]
):
with patch("frigate.record.maintainer.logger.warning") as warn:
# Mock validate_and_move_segment to avoid further logic
maintainer.validate_and_move_segment = MagicMock()
mock_paths = []
for filename in files:
path = MagicMock()
path.name = filename
path.suffix = ".mp4"
path.is_file = AsyncMock(return_value=True)
mock_paths.append(path)
try:
await maintainer.move_files()
except ValueError as e:
if "not enough values to unpack" in str(e):
self.fail("move_files() crashed on bad filename!")
raise e
except Exception:
# Ignore other errors (like DB connection) as we only care about the unpack crash
pass
async def mock_iterdir():
for path in mock_paths:
yield path
# The bad filename is encountered in multiple loops, but should only warn once.
matching = [
c
for c in warn.call_args_list
if c.args
and isinstance(c.args[0], str)
and "Skipping unexpected files in cache" in c.args[0]
]
self.assertEqual(
1,
len(matching),
f"Expected a single warning for unexpected files, got {len(matching)}",
)
with patch("frigate.record.maintainer.AsyncPath") as mock_async_path:
mock_async_path.return_value.iterdir = mock_iterdir
with patch(
"frigate.record.maintainer.psutil.process_iter", return_value=[]
):
with patch("frigate.record.maintainer.logger.warning") as warn:
# Mock validate_and_move_segment to avoid further logic
maintainer.validate_and_move_segment = MagicMock()
try:
await maintainer.move_files()
except ValueError as e:
if "not enough values to unpack" in str(e):
self.fail("move_files() crashed on bad filename!")
raise e
except Exception:
# Ignore other errors (like DB connection) as we only care about the unpack crash
pass
# The bad filename is encountered in multiple loops, but should only warn once.
matching = [
c
for c in warn.call_args_list
if c.args
and isinstance(c.args[0], str)
and "Skipping unexpected files in cache" in c.args[0]
]
self.assertEqual(
1,
len(matching),
f"Expected a single warning for unexpected files, got {len(matching)}",
)
async def test_drops_quiet_segment_when_only_motion_retention(self):
# Regression: when motion retention is enabled but a segment has no
+23 -3
View File
@@ -786,8 +786,15 @@ class TestProfileManager(unittest.TestCase):
dispatcher.clear_runtime_state.assert_not_called()
@patch.object(ProfileManager, "_persist_active_profile")
def test_update_config_clears_when_active_profile_reapplies(self, mock_persist):
"""After /api/config/set, an active-profile re-application drops state."""
def test_update_config_preserves_runtime_state_with_active_profile(
self, mock_persist
):
"""A config/set save must not wipe overrides it never rewrote.
The save path clears matching entries itself via
clear_runtime_state_for_yaml_keys; a broad wipe here would drop
overrides for unrelated cameras.
"""
dispatcher = MagicMock()
manager = ProfileManager(self.config, self.mock_updater, dispatcher)
manager.activate_profile("armed")
@@ -795,7 +802,20 @@ class TestProfileManager(unittest.TestCase):
new_config = FrigateConfig(**self.config_data)
manager.update_config(new_config)
dispatcher.clear_runtime_state.assert_called_once_with()
dispatcher.clear_runtime_state.assert_not_called()
@patch.object(ProfileManager, "_persist_active_profile")
def test_update_config_still_reapplies_active_profile(self, mock_persist):
"""Dropping the wipe must not disturb profile re-application."""
dispatcher = MagicMock()
manager = ProfileManager(self.config, self.mock_updater, dispatcher)
manager.activate_profile("armed")
new_config = FrigateConfig(**self.config_data)
manager.update_config(new_config)
self.assertEqual(manager.config, new_config)
self.assertEqual(new_config.active_profile, "armed")
@patch.object(ProfileManager, "_persist_active_profile")
def test_update_config_does_not_clear_when_no_active_profile(self, mock_persist):
+19
View File
@@ -131,6 +131,25 @@ class TestRuntimeStatePersistence(unittest.TestCase):
self.store.clear_all()
self.assertEqual(self.store.load(), {})
def test_clear_camera_removes_only_that_camera(self) -> None:
self.store.set("front_door", "enabled", False)
self.store.set("front_door", "detect", False)
self.store.set("back_yard", "audio", False)
self.store.clear_camera("front_door")
self.assertEqual(self.store.load(), {"back_yard": {"audio": False}})
def test_clear_camera_is_noop_for_unknown_camera(self) -> None:
self.store.set("front_door", "enabled", False)
self.store.clear_camera("side_gate")
self.assertEqual(self.store.load(), {"front_door": {"enabled": False}})
def test_clear_camera_is_safe_when_file_missing(self) -> None:
# No prior set() calls, so the file does not exist
self.store.clear_camera("front_door")
self.assertEqual(self.store.load(), {})
if __name__ == "__main__":
unittest.main()
+17 -2
View File
@@ -2,5 +2,20 @@
target-version = "py311"
[tool.ruff.lint]
ignore = ["E501","E711","E712","UP031","UP032","UP042","G004"]
extend-select = ["I", "UP", "G", "ASYNC210", "B904"]
ignore = [
"ASYNC109", # Async function definition with a timeout parameter
"E501", # line-too-long
"E711", # none-comparison
"E712", # true-false-comparison
"UP031", # printf-string-formatting
"UP032", # f-string
"UP042", # replace-str-enum
"G004", # logging-f-string
]
extend-select = [
"ASYNC", # https://docs.astral.sh/ruff/rules/#flake8-async-async
"B904", # https://docs.astral.sh/ruff/rules/raise-without-from-inside-except/
"G", # https://docs.astral.sh/ruff/rules/#flake8-logging-format-g
"I", # https://docs.astral.sh/ruff/rules/#isort-i
"UP", # https://docs.astral.sh/ruff/rules/#pyupgrade-up
]