Compare commits

..
Author SHA1 Message Date
Josh Hawkins db921a14dd report progress during the ownership sweep 2026-08-24 11:36:09 -05:00
Josh Hawkins 3de88492ec unwrap hard-wrapped prose in the installation docs 2026-08-24 10:48:53 -05:00
Josh Hawkins e69795ce8f discard stdout for the unprivileged smoke nginx -t 2026-08-24 10:05:23 -05:00
Josh Hawkins 22b9e4cb45 re-own the nginx shm cache on service restart 2026-08-24 10:01:38 -05:00
Josh Hawkins 92d18a3e82 run smoke nginx -t and the write probe as the runtime user 2026-08-24 10:01:38 -05:00
Josh Hawkins 6caf872999 set HOME to /config for non-root services 2026-08-24 10:01:38 -05:00
Josh Hawkins 383f26a132 chown the s6 log pipe so non-root nginx can reopen /dev/stdout 2026-08-24 09:52:59 -05:00
Josh Hawkins 78e39e7a22 tolerate homekit config chown failures in the go2rtc run script 2026-08-23 18:26:13 -05:00
Josh Hawkins fa191b0faa only write the sweep sentinel when a media volume is mounted 2026-08-23 18:26:13 -05:00
Josh Hawkins 0744021809 Assert non-root services, JWT migration, and escape hatch in CI 2026-08-23 18:17:48 -05:00
Josh Hawkins 90d6712e47 Create /media/frigate after the ownership sweep 2026-08-23 18:17:48 -05:00
Josh Hawkins 94bd0b043f Document non-root operation and per-hardware device access 2026-08-23 18:03:02 -05:00
Josh Hawkins 8183da8020 Hand TensorRT model cache ownership to the runtime user 2026-08-23 18:01:10 -05:00
Josh Hawkins 59b4ff6fd6 Disable bandwidth stats gracefully when not running as root 2026-08-23 17:59:37 -05:00
Josh Hawkins 716326fdcc Run nginx as the frigate user with writable state in /tmp/nginx 2026-08-23 17:59:19 -05:00
Josh Hawkins cb3c3b0e8c Run go2rtc as its own restricted user 2026-08-23 17:59:19 -05:00
Josh Hawkins 74aa85c646 Run the frigate service as the frigate user 2026-08-23 17:59:19 -05:00
41 changed files with 527 additions and 1170 deletions
+74 -4
View File
@@ -59,10 +59,15 @@ jobs:
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }} GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
- name: Start container - name: Start container
run: | run: |
mkdir -p /tmp/frigate-config mkdir -p /tmp/frigate-config /tmp/frigate-media
printf 'mqtt:\n enabled: false\ncameras: {}\n' > /tmp/frigate-config/config.yml printf 'mqtt:\n enabled: false\ncameras: {}\n' > /tmp/frigate-config/config.yml
# simulate a root-era install: root-owned 0600 jwt secret pre-exists
docker run --rm -v /tmp/frigate-config:/config --entrypoint bash \
${{ steps.setup.outputs.image-name }}-amd64 \
-c "python3 -c 'import secrets; open(\"/config/.jwt_secret\",\"w\").write(secrets.token_hex(64))' && chmod 600 /config/.jwt_secret && chown 0:0 /config/.jwt_secret"
docker run -d --name frigate --shm-size 256m \ docker run -d --name frigate --shm-size 256m \
-v /tmp/frigate-config:/config \ -v /tmp/frigate-config:/config \
-v /tmp/frigate-media:/media/frigate \
-p 5000:5000 -p 8971:8971 \ -p 5000:5000 -p 8971:8971 \
${{ steps.setup.outputs.image-name }}-amd64 ${{ steps.setup.outputs.image-name }}-amd64
- name: Wait for API - name: Wait for API
@@ -91,16 +96,81 @@ jobs:
echo "response carries frame-ancestors, which breaks cross-origin iframe embedding" echo "response carries frame-ancestors, which breaks cross-origin iframe embedding"
exit 1 exit 1
fi fi
docker exec frigate /usr/local/nginx/sbin/nginx -t # -t must NOT run as root: ngx_create_paths would chown the live cache
# and temp dirs to the `user root` directive user, breaking the workers.
# stdout goes to /dev/null because -t reopens the config's
# error_log/access_log /dev/stdout by path, and the docker exec pipe
# is root-owned; -t reports on stderr, so nothing is lost
docker exec frigate /command/s6-setuidgid frigate bash -c '/usr/local/nginx/sbin/nginx -e stderr -t -c /tmp/nginx/conf/nginx.conf >/dev/null'
docker exec frigate stat -c %a /etc/letsencrypt/live/frigate/privkey.pem | grep -qx 600 docker exec frigate stat -c %a /etc/letsencrypt/live/frigate/privkey.pem | grep -qx 600
docker exec frigate stat -c %a /dev/shm/go2rtc.yaml | grep -qx 640 docker exec frigate stat -c %a /dev/shm/go2rtc.yaml | grep -qx 640
- name: Assert services run as non-root
run: |
ps_out=$(docker exec frigate ps -eo user=,comm=)
echo "$ps_out"
assert_nonroot() {
# the process must exist AND no instance of it may run as root
echo "$ps_out" | grep -qw "$1" || { echo "$1 is not running"; exit 1; }
if echo "$ps_out" | grep -w "$1" | grep -q '^root'; then
echo "$1 is running as root"; exit 1
fi
}
assert_nonroot python3
assert_nonroot go2rtc
assert_nonroot nginx
# root-era jwt secret must have been captured by the sweep and the
# auth stack must be functional: wrong creds => clean 401, not 500
docker exec frigate stat -c %u /config/.jwt_secret | grep -qx "$(docker exec frigate id -u frigate)"
code=$(curl -s -o /dev/null -w '%{http_code}' -X POST http://127.0.0.1:5000/api/login \
-H 'content-type: application/json' -d '{"user":"admin","password":"definitely-wrong"}')
[ "$code" = "401" ] || { echo "login endpoint returned $code"; exit 1; }
# nginx runtime state must belong to the runtime user (a root nginx -t
# in the step above would have chowned it to root)
owners=$(docker exec frigate stat -c %U /tmp/nginx /dev/shm/nginx_cache)
echo "$owners"
if echo "$owners" | grep -qvx frigate; then
echo "nginx runtime dirs are not owned by frigate"; exit 1
fi
# runtime user can write recordings storage
docker exec frigate /command/s6-setuidgid frigate touch /media/frigate/.write-probe
docker exec frigate rm /media/frigate/.write-probe
- name: Assert escape hatch restores root
run: |
mkdir -p /tmp/frigate-config-root
printf 'mqtt:\n enabled: false\ncameras: {}\n' > /tmp/frigate-config-root/config.yml
# pre-seed a sentinel: the assertion below is that the escape hatch
# DELETES it. Against a fresh dir the absence check passes vacuously
# and proves nothing about the rm -f in the prepare script.
echo "2:1000:1000" > /tmp/frigate-config-root/.permissions_version
docker run -d --name frigate-root --shm-size 256m \
-e FRIGATE_RUN_AS_ROOT=true \
-v /tmp/frigate-config-root:/config \
${{ steps.setup.outputs.image-name }}-amd64
up=0
for i in $(seq 1 60); do
docker exec frigate-root curl -fs http://127.0.0.1:5000/api/version && up=1 && break
sleep 5
done
if [ "$up" -ne 1 ]; then echo "escape hatch container never healthy"; docker logs frigate-root; exit 1; fi
ps_out=$(docker exec frigate-root ps -eo user=,comm=)
echo "$ps_out"
echo "$ps_out" | grep -w python3 | grep -q '^root'
echo "$ps_out" | grep -w go2rtc | grep -q '^root'
echo "$ps_out" | grep -w nginx | grep -q '^root'
# escape hatch must have DELETED the pre-seeded sentinel. written as
# an if because bash exempts a negated command from set -e
if docker exec frigate-root test -f /config/.permissions_version; then
echo "escape hatch did not delete the sweep sentinel"; exit 1
fi
docker rm -f frigate-root
- name: Assert PUID/PGID remapping - name: Assert PUID/PGID remapping
run: | run: |
mkdir -p /tmp/frigate-config-puid mkdir -p /tmp/frigate-config-puid /tmp/frigate-media-puid
printf 'mqtt:\n enabled: false\ncameras: {}\n' > /tmp/frigate-config-puid/config.yml printf 'mqtt:\n enabled: false\ncameras: {}\n' > /tmp/frigate-config-puid/config.yml
docker run -d --name frigate-puid --shm-size 256m \ docker run -d --name frigate-puid --shm-size 256m \
-e PUID=1500 -e PGID=1500 \ -e PUID=1500 -e PGID=1500 \
-v /tmp/frigate-config-puid:/config \ -v /tmp/frigate-config-puid:/config \
-v /tmp/frigate-media-puid:/media/frigate \
${{ steps.setup.outputs.image-name }}-amd64 ${{ steps.setup.outputs.image-name }}-amd64
up=0 up=0
for i in $(seq 1 60); do for i in $(seq 1 60); do
@@ -110,7 +180,7 @@ jobs:
if [ "$up" -ne 1 ]; then echo "PUID container never became healthy"; docker logs frigate-puid; exit 1; fi if [ "$up" -ne 1 ]; then echo "PUID container never became healthy"; docker logs frigate-puid; exit 1; fi
docker exec frigate-puid id -u frigate | grep -qx 1500 docker exec frigate-puid id -u frigate | grep -qx 1500
docker exec frigate-puid id -g frigate | grep -qx 1500 docker exec frigate-puid id -g frigate | grep -qx 1500
docker exec frigate-puid cat /config/.permissions_version | grep -qx "1:1500:1500" docker exec frigate-puid cat /config/.permissions_version | grep -qx "2:1500:1500"
# second boot must skip the sweep (sentinel hit). Poll rather than # second boot must skip the sweep (sentinel hit). Poll rather than
# sleep: the string can only come from the second boot (the first # sleep: the string can only come from the second boot (the first
# had no sentinel), so grepping the full log is unambiguous. # had no sentinel), so grepping the full log is unambiguous.
@@ -49,7 +49,7 @@ do
then then
echo "[INFO] Reloading nginx to refresh TLS certificate" echo "[INFO] Reloading nginx to refresh TLS certificate"
echo "$lefile: $leprint" echo "$lefile: $leprint"
/usr/local/nginx/sbin/nginx -s reload /usr/local/nginx/sbin/nginx -c /tmp/nginx/conf/nginx.conf -s reload
fi fi
sleep 60 sleep 60
@@ -4,6 +4,13 @@
set -o errexit -o nounset -o pipefail set -o errexit -o nounset -o pipefail
# $HOME is /root from the container env and survives s6-setuidgid, so cache
# and telemetry writes (huggingface, openvino) fail after the drop. Set it
# before opt_in_out so the opt-out marker lands where the service will look.
if [[ "${FRIGATE_RUN_AS_ROOT:-false}" != "true" ]]; then
export HOME=/config
fi
# opt out of openvino telemetry # opt out of openvino telemetry
if [ -e /usr/local/bin/opt_in_out ]; then if [ -e /usr/local/bin/opt_in_out ]; then
/usr/local/bin/opt_in_out --opt_out > /dev/null 2>&1 /usr/local/bin/opt_in_out --opt_out > /dev/null 2>&1
@@ -30,4 +37,8 @@ cd /opt/frigate || echo "[ERROR] Failed to change working directory to /opt/frig
# Replace the bash process with the Frigate process, redirecting stderr to stdout # Replace the bash process with the Frigate process, redirecting stderr to stdout
exec 2>&1 exec 2>&1
exec python3 -u -m frigate if [[ "$(id -u)" -ne 0 || "${FRIGATE_RUN_AS_ROOT:-false}" == "true" ]]; then
exec python3 -u -m frigate
else
exec s6-setuidgid frigate python3 -u -m frigate
fi
@@ -110,6 +110,16 @@ fi
readonly homekit_config_path="/config/go2rtc_homekit.yml" readonly homekit_config_path="/config/go2rtc_homekit.yml"
setup_homekit_config "${homekit_config_path}" setup_homekit_config "${homekit_config_path}"
if [[ "$(id -u)" -eq 0 && "${FRIGATE_RUN_AS_ROOT:-false}" != "true" ]]; then
chown go2rtc:go2rtc /dev/shm/go2rtc.yaml 2>/dev/null || true
# go2rtc rewrites this in place (os.WriteFile, no rename), so owning the
# file is enough; /config grants frigate-data traverse only. Tolerated so
# a chown-refusing mount (NFS root_squash) degrades pairing persistence
# instead of crash-looping the service
chown go2rtc:frigate-data "${homekit_config_path}" 2>/dev/null && chmod 664 "${homekit_config_path}" 2>/dev/null || \
echo "[WARN] Could not hand ${homekit_config_path} to the go2rtc user; HomeKit pairing changes may not persist"
fi
readonly config_path="/config" readonly config_path="/config"
if [[ -x "${config_path}/go2rtc" ]]; then if [[ -x "${config_path}/go2rtc" ]]; then
@@ -125,4 +135,8 @@ echo "[INFO] Starting go2rtc..."
# Use HomeKit config as the primary config so writebacks go there # Use HomeKit config as the primary config so writebacks go there
# The main config from Frigate will be loaded as a secondary config # The main config from Frigate will be loaded as a secondary config
exec 2>&1 exec 2>&1
exec "${binary_path}" -config="${homekit_config_path}" -config=/dev/shm/go2rtc.yaml if [[ "$(id -u)" -ne 0 || "${FRIGATE_RUN_AS_ROOT:-false}" == "true" ]]; then
exec "${binary_path}" -config="${homekit_config_path}" -config=/dev/shm/go2rtc.yaml
else
exec s6-setuidgid go2rtc "${binary_path}" -config="${homekit_config_path}" -config=/dev/shm/go2rtc.yaml
fi
@@ -2,4 +2,4 @@
set -e set -e
# Wait for PID file to exist. # Wait for PID file to exist.
while ! test -f /run/nginx.pid; do sleep 1; done while ! test -f /tmp/nginx/nginx.pid; do sleep 1; done
@@ -59,10 +59,14 @@ function set_worker_processes() {
cpus=4 cpus=4
fi fi
# we need to catch any errors because sed will fail if user has bind mounted a custom nginx file sed -i "s/worker_processes auto;/worker_processes ${cpus};/" /tmp/nginx/conf/nginx.conf
sed -i "s/worker_processes auto;/worker_processes ${cpus};/" /usr/local/nginx/conf/nginx.conf || true
} }
# copied whole so the conf tree's relative includes still resolve
mkdir -p /tmp/nginx/conf /tmp/nginx/client_body /tmp/nginx/proxy \
/tmp/nginx/fastcgi /tmp/nginx/uwsgi /tmp/nginx/scgi
cp -r /usr/local/nginx/conf/. /tmp/nginx/conf/
set_worker_processes set_worker_processes
# ensure the directory for ACME challenges exists # ensure the directory for ACME challenges exists
@@ -87,15 +91,40 @@ nginx_settings=$(python3 /usr/local/nginx/get_nginx_settings.py)
# build templates for optional FRIGATE_BASE_PATH environment variable # build templates for optional FRIGATE_BASE_PATH environment variable
echo "$nginx_settings" | \ echo "$nginx_settings" | \
tempio -template /usr/local/nginx/templates/base_path.gotmpl \ tempio -template /usr/local/nginx/templates/base_path.gotmpl \
-out /usr/local/nginx/conf/base_path.conf -out /tmp/nginx/conf/base_path.conf
# build templates for additional network settings # build templates for additional network settings
echo "$nginx_settings" | \ echo "$nginx_settings" | \
tempio -template /usr/local/nginx/templates/listen.gotmpl \ tempio -template /usr/local/nginx/templates/listen.gotmpl \
-out /usr/local/nginx/conf/listen.conf -out /tmp/nginx/conf/listen.conf
if [[ "$(id -u)" -eq 0 && "${FRIGATE_RUN_AS_ROOT:-false}" != "true" ]]; then
chown -R frigate:frigate /tmp/nginx
# heal the cache if a root `nginx -t` chowned it (ngx_create_paths chowns
# every cycle path to the `user` directive user when run as root)
if [ -d /dev/shm/nginx_cache ]; then
chown -R frigate:frigate /dev/shm/nginx_cache
fi
# error_log/access_log /dev/stdout make nginx REOPEN the s6 log pipe by
# path, and s6 created it root-owned 0600; without this the non-root
# master exits with "open() /dev/stdout failed (13: Permission denied)"
chown frigate /dev/stdout
# self-signed certs are root-generated; tolerant because mounted certs may be :ro
if [ -f "$letsencrypt_path/privkey.pem" ]; then
chown frigate:frigate "$letsencrypt_path/privkey.pem" "$letsencrypt_path/fullchain.pem" 2>/dev/null || true
fi
fi
# Replace the bash process with the NGINX process, redirecting stderr to stdout # Replace the bash process with the NGINX process, redirecting stderr to stdout
exec 2>&1 exec 2>&1
exec \ # -e stderr: the compile-time default error log under /usr/local/nginx/logs
# is not writable by the runtime user and would alert before config load
if [[ "$(id -u)" -ne 0 || "${FRIGATE_RUN_AS_ROOT:-false}" == "true" ]]; then
exec \
s6-notifyoncheck -t 30000 -n 1 \ s6-notifyoncheck -t 30000 -n 1 \
nginx nginx -e stderr -c /tmp/nginx/conf/nginx.conf
else
exec \
s6-notifyoncheck -t 30000 -n 1 \
s6-setuidgid frigate nginx -e stderr -c /tmp/nginx/conf/nginx.conf
fi
@@ -153,7 +153,24 @@ if [[ "$(id -u)" -eq 0 ]]; then
if [[ "${FRIGATE_RUN_AS_ROOT:-false}" == "true" ]]; then if [[ "${FRIGATE_RUN_AS_ROOT:-false}" == "true" ]]; then
rm -f /config/.permissions_version rm -f /config/.permissions_version
else else
/usr/local/bin/fix-ownership --sentinel /config/.permissions_version \ # Only record the sentinel when something is mounted under /media: a
# sweep blessed against the container-local dir created below would let
# a volume attached later skip the sweep forever
sentinel_args=(--sentinel /config/.permissions_version)
if ! awk '$2 == "/media" || $2 ~ /^\/media\//' /proc/mounts | grep -q .; then
sentinel_args=()
fi
/usr/local/bin/fix-ownership "${sentinel_args[@]}" \
"${PUID:-1000}" "${PGID:-1000}" /config /media/frigate "${PUID:-1000}" "${PGID:-1000}" /config /media/frigate
fi fi
fi fi
# Not in the image, and the runtime user cannot create it under root-owned
# /media. Must stay after the sweep, which reads an absent /media/frigate as an
# unmounted volume rather than a swept one
if [[ "$(id -u)" -eq 0 && ! -d /media/frigate ]]; then
mkdir -p /media/frigate
if [[ "${FRIGATE_RUN_AS_ROOT:-false}" != "true" ]]; then
chown "${PUID:-1000}:${PGID:-1000}" /media/frigate
fi
fi
+31 -5
View File
@@ -22,7 +22,7 @@ set -o errexit -o nounset -o pipefail
# Permissions-layout epoch. Bump to force a one-time re-sweep on upgrade # Permissions-layout epoch. Bump to force a one-time re-sweep on upgrade
# (e.g. when the privilege-drop release must capture files created as root # (e.g. when the privilege-drop release must capture files created as root
# since the previous sweep). # since the previous sweep).
schema=1 schema=2
dry_run=0 dry_run=0
sentinel="" sentinel=""
@@ -76,6 +76,8 @@ for path in "$@"; do
continue continue
fi fi
echo "[INFO] fix-ownership: scanning ${path} for ownership mismatches; this may take a while on large filesystems"
# find may fail mid-walk on a live volume (file deleted under it) or on a # find may fail mid-walk on a live volume (file deleted under it) or on a
# stale mount. Tolerate it rather than aborting under errexit, but never # stale mount. Tolerate it rather than aborting under errexit, but never
# read a failed scan as "nothing to do": that would record the sweep as # read a failed scan as "nothing to do": that would record the sweep as
@@ -98,17 +100,41 @@ for path in "$@"; do
echo "[WARN] fix-ownership: ${path} contains symlinked directories; ownership behind them is not managed and must be aligned by hand" echo "[WARN] fix-ownership: ${path} contains symlinked directories; ownership behind them is not managed and must be aligned by hand"
fi fi
echo "[WARN] fix-ownership: adjusting ownership of ${count} entries under ${path}; on large recordings volumes this can take a long time" echo "[WARN] fix-ownership: adjusting ownership of ${count} entries under ${path}"
if [[ "$dry_run" -eq 1 ]]; then if [[ "$dry_run" -eq 1 ]]; then
echo "[INFO] fix-ownership: dry run, not changing ${path}" echo "[INFO] fix-ownership: dry run, not changing ${path}"
continue continue
fi fi
find "$path" \( -not -uid "$target_uid" -o -not -gid "$target_gid" \) \ # -print feeds the progress counter while -exec {} + keeps the chown
-exec chown -h "${target_uid}:${target_gid}" {} + || { # batched; the scan above is what makes a real percentage possible
started=$SECONDS
if find "$path" \( -not -uid "$target_uid" -o -not -gid "$target_gid" \) \
-print -exec chown -h "${target_uid}:${target_gid}" {} + \
| awk -v total="$count" -v path="$path" '
BEGIN { next_pct = 5 }
{
pct = int(NR * 100 / total)
if (pct > 100) pct = 100
if (pct >= next_pct) {
printf "[INFO] fix-ownership: %s %d%% (%d/%d entries)\n", path, pct, NR, total
# mawk block-buffers to a pipe; without this the whole
# progress log arrives at once when the sweep ends
fflush()
while (next_pct <= pct) next_pct += 5
}
}'; then
elapsed=$((SECONDS - started))
if [[ "$elapsed" -ge 60 ]]; then
elapsed="$((elapsed / 60))m $((elapsed % 60))s"
else
elapsed="${elapsed}s"
fi
echo "[INFO] fix-ownership: finished ${path} in ${elapsed}"
else
swept_clean=0 swept_clean=0
echo "[WARN] fix-ownership: some entries under ${path} could not be updated (deleted mid-sweep or chown denied); will retry on next mismatch" echo "[WARN] fix-ownership: some entries under ${path} could not be updated (deleted mid-sweep or chown denied); will retry on next mismatch"
} fi
done done
# go2rtc (separate user) must be able to REACH its HomeKit state in /config. # go2rtc (separate user) must be able to REACH its HomeKit state in /config.
@@ -1,9 +1,15 @@
# Copied to /tmp/nginx/conf at startup and loaded with -c from there. Relative
# includes resolve against the -c file, but every other path directive resolves
# against the compile-time --prefix, so non-include paths must stay absolute.
daemon off; daemon off;
# Ignored by a non-root master; required by FRIGATE_RUN_AS_ROOT so workers
# stay root instead of the compiled-in default user
user root; user root;
worker_processes auto; worker_processes auto;
error_log /dev/stdout warn; error_log /dev/stdout warn;
pid /var/run/nginx.pid; pid /tmp/nginx/nginx.pid;
events { events {
worker_connections 1024; worker_connections 1024;
@@ -13,6 +19,12 @@ http {
map_hash_bucket_size 256; map_hash_bucket_size 256;
server_tokens off; server_tokens off;
client_body_temp_path /tmp/nginx/client_body;
proxy_temp_path /tmp/nginx/proxy;
fastcgi_temp_path /tmp/nginx/fastcgi;
uwsgi_temp_path /tmp/nginx/uwsgi;
scgi_temp_path /tmp/nginx/scgi;
include mime.types; include mime.types;
default_type application/octet-stream; default_type application/octet-stream;
+10
View File
@@ -36,6 +36,16 @@ if ! [[ "$puid" =~ ^[0-9]+$ && "$pgid" =~ ^[0-9]+$ ]]; then
fi fi
echo "[INFO] Using image ${IMAGE} (override with FRIGATE_IMAGE=...)" echo "[INFO] Using image ${IMAGE} (override with FRIGATE_IMAGE=...)"
if ! docker image inspect "${IMAGE}" >/dev/null 2>&1; then
echo "[INFO] ${IMAGE} is not present locally and has to be pulled first; this may take a while"
fi
if [[ -n "$dry_run_flag" ]]; then
echo "[INFO] Dry run: reporting what would change under ${config_dir} and ${media_dir}, changing nothing"
else
echo "[INFO] Aligning ${config_dir} and ${media_dir} to ${puid}:${pgid}; this may take a while on large filesystems"
fi
# shellcheck disable=SC2086 # shellcheck disable=SC2086
docker run --rm \ docker run --rm \
-v "${config_dir}:/config" \ -v "${config_dir}:/config" \
@@ -13,6 +13,16 @@ TRT_VER=${TRT_VER:-$(cat /etc/TENSORRT_VER)}
OUTPUT_FOLDER="${MODEL_CACHE_DIR}/${TRT_VER}" OUTPUT_FOLDER="${MODEL_CACHE_DIR}/${TRT_VER}"
YOLO_MODELS=${YOLO_MODELS:-""} YOLO_MODELS=${YOLO_MODELS:-""}
# This runs as root after prepare's sentinel-guarded sweep, so the dirs and
# engines it creates below are the runtime user's to fix up, on every exit path
function hand_off_ownership() {
if [[ "$(id -u)" -eq 0 && "${FRIGATE_RUN_AS_ROOT:-false}" != "true" ]]; then
/usr/local/bin/fix-ownership "${PUID:-1000}" "${PGID:-1000}" \
/config/model_cache "${MODEL_CACHE_DIR}"
fi
}
trap hand_off_ownership EXIT
# Create output folder # Create output folder
mkdir -p ${OUTPUT_FOLDER} mkdir -p ${OUTPUT_FOLDER}
+1 -2
View File
@@ -84,10 +84,9 @@ A camera is enabled by default but can be disabled by using `enabled: False`. Ca
Each role can only be assigned to one input per camera. The options for roles are as follows: Each role can only be assigned to one input per camera. The options for roles are as follows:
| Role | Description | | Role | Description |
| ------------ | ------------------------------------------------------------------------------------------------------------ | | -------- | ----------------------------------------------------------------------------------- |
| `detect` | Main feed for object detection. [docs](object_detectors.md) | | `detect` | Main feed for object detection. [docs](object_detectors.md) |
| `record` | Saves segments of the video feed based on configuration settings. [docs](record.md) | | `record` | Saves segments of the video feed based on configuration settings. [docs](record.md) |
| `record_sub` | Saves segments of a second, lower quality stream with its own retention. [docs](record.md#sub-stream-recording) |
| `audio` | Feed for audio based detection. [docs](audio_detectors.md) | | `audio` | Feed for audio based detection. [docs](audio_detectors.md) |
<ConfigTabs> <ConfigTabs>
+92
View File
@@ -0,0 +1,92 @@
---
id: non_root
title: Running as a non-root user
---
Frigate's services run as an unprivileged user inside the container. The main Frigate process and nginx run as `frigate`, and go2rtc runs as its own more restricted `go2rtc` user. Only the s6 init system and the certsync helper stay root.
By default the runtime user is uid/gid `1000:1000`. You can change it with `PUID`/`PGID`, or bypass Frigate's user handling entirely with Docker's own `user:`.
## Run modes
| Mode | How to enable | Ownership of `/config` and `/media/frigate` | `read_only: true` |
| -------------------------------- | -------------------------------------- | ---------------------------------------------------------- | ----------------- |
| Default | nothing, this is the default | Aligned to `1000:1000` on first boot | Not supported |
| `PUID`/`PGID` | `PUID=1001`, `PGID=1001` | Aligned to the values you set, on first boot | Not supported |
| Docker-native user | `user: "1001:1001"` | You own it, Frigate never changes ownership | Not supported |
| Root (escape hatch) | `FRIGATE_RUN_AS_ROOT=true` | Never touched | Not supported |
`PUID`/`PGID` remapping runs `usermod` at startup, which writes to `/etc/passwd`, so it can't work with a read-only root filesystem. That combination fails fast at startup with a message pointing here rather than failing obscurely later.
`FRIGATE_RUN_AS_ROOT` is matched against the exact lowercase string `true`. `True`, `TRUE`, and `1` are all ignored.
## Migrating an existing install
Volumes created by earlier versions of Frigate are owned by root. Ownership has to be aligned with the runtime user once.
This happens automatically on the first boot after upgrading, but on large recordings volumes it's much better to do it from the host beforehand. The boot sweep runs before any service starts, so a multi-terabyte `/media/frigate` can hold the container in startup long enough for Docker's healthcheck to mark it unhealthy, and orchestrators that react to health will restart it mid-sweep. If you'd rather not run the script, raise the healthcheck start period instead (`--start-period=1800s`, or `start_period: 1800s` under `healthcheck:` in compose).
Grab `fix-permissions.sh` from `docker/migration/` in the Frigate repo and dry run it first:
```bash
./fix-permissions.sh --dry-run /path/to/your/config /path/to/your/storage
```
That reports how many entries would change and touches nothing. When it looks right, run it without `--dry-run`:
```bash
./fix-permissions.sh /path/to/your/config /path/to/your/storage
```
Both the script and the boot sweep report progress as they go, so you can tell a slow sweep apart from a stuck one:
```
[INFO] fix-ownership: scanning /media/frigate for ownership mismatches; this may take a while on large filesystems
[WARN] fix-ownership: adjusting ownership of 4823941 entries under /media/frigate
[INFO] fix-ownership: /media/frigate 5% (241197/4823941 entries)
[INFO] fix-ownership: /media/frigate 10% (482394/4823941 entries)
[INFO] fix-ownership: finished /media/frigate in 12m 4s
```
The scan has no percentage behind it because the total isn't known until it finishes. Watch the boot sweep with `docker logs -f frigate`.
Pass `PUID` and `PGID` as the third and fourth arguments if you're not using the default `1000:1000`. The script wraps the same `fix-ownership` helper the container uses, so it's the same logic either way. Override the image it pulls with `FRIGATE_IMAGE=...` if you're not on `stable`.
Once the volumes are aligned, start Frigate normally. A sentinel at `/config/.permissions_version` records what was done, so later boots skip the sweep entirely unless you change `PUID`/`PGID`.
## Rolling back
Set `FRIGATE_RUN_AS_ROOT=true` and restart. Everything runs as root again, exactly as it did before.
The escape hatch never changes ownership, and it deletes the sweep sentinel on startup, so switching back to non-root later re-sweeps whatever root created in the meantime. Toggling in either direction is safe.
## Hardware device access
Supplementary groups can't open a device node that's `root:root` with mode `0600`. If your accelerator's node isn't group readable on the host, no amount of container configuration will fix it, so the fix belongs on the host.
Use `group_add` in compose (`--group-add` with `docker run`) to give the runtime user a host GID. `EXTRA_GROUPS` does the same thing and also covers the `go2rtc` user, which needs render and video access to run hardware accelerated restreams.
| Hardware | Device(s) | Non-root requirement |
| ------------------------- | ----------------------------------------------------------- | ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| Intel/AMD GPU (VAAPI/QSV) | `/dev/dri/renderD128` | `group_add: ["<host render GID>"]` from `getent group render`, or `EXTRA_GROUPS` |
| Intel/AMD NPU | `/dev/accel` | Host udev rule granting a group, then `group_add` that GID |
| Coral USB | `/dev/bus/usb` | Host udev rule granting plugdev, for example `SUBSYSTEMS=="usb", ATTRS{idVendor}=="1a6e", GROUP="plugdev"` and the same for `18d1` post-init, then `group_add` the plugdev GID |
| Coral PCIe | `/dev/apex_0` | Host udev rule `SUBSYSTEM=="apex", GROUP="apex", MODE="0660"`, then `group_add` that GID |
| Hailo | `/dev/hailo0` | Host udev rule granting a group, then `group_add` that GID |
| NVIDIA | nvidia runtime | Works non-root with the nvidia-container-toolkit defaults |
| AMD ROCm | `/dev/kfd`, `/dev/dri` | `group_add` the host `video` and `render` GIDs |
| Raspberry Pi | `/dev/video11` | `group_add` the host `video` GID |
| Rockchip | `/dev/dri`, `/dev/dma_heap`, `/dev/rga`, `/dev/mpp_service` | These are commonly `root:root` `0600`, so host udev rules are required. If you can't grant access to all four, use `FRIGATE_RUN_AS_ROOT=true` |
| Axera (AXCL) | `/dev/ax_*` per the AXCL driver docs | Unverified. Node ownership is driver dependent, check it on your hardware before assuming this works |
| Synaptics SL1680 | per the Synaptics docs | Unverified |
| MemryX | per the MemryX docs | Still requires `privileged: true`, which means root. Out of scope for non-root operation |
## Known limitations
`telemetry.stats.network_bandwidth` uses nethogs, which needs `CAP_NET_ADMIN` and `CAP_NET_RAW` and therefore root. The stat is disabled automatically when Frigate isn't running as root, with one warning in the log. Use `FRIGATE_RUN_AS_ROOT=true` if you need it.
go2rtc's ffmpeg processes no longer appear in Intel GPU stats. Frigate reads per-process GPU usage from `/proc/<pid>/fdinfo`, which the kernel won't let one user read for another user's processes, so anything go2rtc spawns is invisible to it. Overall GPU utilization is unaffected.
If you mount your own TLS certificate at `/etc/letsencrypt/live/frigate`, the private key has to be readable by the runtime user. Frigate won't change ownership of a certificate you supplied, since the mount may be read-only.
If you're debugging nginx, run the config check as the runtime user with stdout discarded: `docker exec frigate /command/s6-setuidgid frigate bash -c 'nginx -t -c /tmp/nginx/conf/nginx.conf >/dev/null'`. Running `nginx -t` as root hands nginx's runtime directories to root as a side effect, which breaks the running workers until the service restarts, and the config's `/dev/stdout` logs can't be reopened through a root-owned `docker exec` pipe (the results print on stderr either way).
+2 -2
View File
@@ -280,7 +280,7 @@ This configuration will retain recording segments that overlap with alerts and d
In addition to the main recording stream, Frigate can record a second, lower quality stream for each camera. This serves two purposes: In addition to the main recording stream, Frigate can record a second, lower quality stream for each camera. This serves two purposes:
- **Quality selection during playback**: A quality selector (`Auto`, `Original`, or `Low`) appears in History view for cameras with sub stream recording enabled. `Original` and `Low` play only that stream's recordings. Time ranges where the selected stream has no footage are skipped during playback, and the selector notes when the selected stream has no recordings at all in the viewed time range. With `Auto` (the default), playback prefers the original quality and automatically falls back to the low quality stream when the connection cannot keep up, or for time ranges where the original recordings have expired. The selector shows each stream's video codec and audio details beneath the options; footage recorded by older Frigate versions shows no details. - **Quality selection during playback**: A quality selector (`Auto`, `Original`, or `Low`) appears in History view for cameras with sub stream recording enabled. `Original` and `Low` play only that stream's recordings. Time ranges where the selected stream has no footage are skipped during playback, and the selector notes when the selected stream has no recordings at all in the viewed time range. With `Auto` (the default), playback prefers the original quality and automatically falls back to the low quality stream when the connection cannot keep up, or for time ranges where the original recordings have expired. The selector shows each stream's video codec and audio details beneath the options; footage recorded by older Frigate versions shows no details.
- **Extended retention**: Sub stream recordings have their own retention settings, fully independent of the main recordings. By giving the low quality recordings a longer retention period, you can keep weeks or months of low quality history using a fraction of the storage, and that history remains playable after the main recordings expire. Playback falls back to the low quality recordings automatically, and the timeline shows a muted treatment for time ranges where only low quality footage remains. Timeline previews are kept for as long as either stream still has recordings, so scrubbing works across the whole retained history. - **Extended retention**: Sub stream recordings have their own retention settings, fully independent of the main recordings. By giving the low quality recordings a longer retention period, you can keep weeks or months of low quality history using a fraction of the storage, and that history remains playable after the main recordings expire. Playback falls back to the low quality recordings automatically, and the timeline shows a muted treatment for time ranges where only low quality footage remains.
### Configuring sub stream recording ### Configuring sub stream recording
@@ -427,7 +427,7 @@ This table covers only features that read recordings from disk. Tracked object s
### Trade-offs ### Trade-offs
- Recording a second stream increases overall storage use. The increase is typically small relative to the main recordings, since the low quality stream is much smaller. Both streams are cached before being written to disk, so cache use goes up as well. See [the `/tmp/cache` area is separate](#the-tmpcache-area-is-separate) if you start seeing `No space left on device` errors after enabling it. - Recording a second stream increases overall storage use. The increase is typically small relative to the main recordings, since the low quality stream is much smaller.
- The go2rtc transcode approach continuously encodes the low quality stream, which uses CPU or GPU resources. This cost only applies to the transcode path; recording the camera's native sub stream does not re-encode. See the [go2rtc hardware acceleration documentation](https://github.com/AlexxIT/go2rtc?tab=readme-ov-file#source-ffmpeg) for accelerating the transcode. - The go2rtc transcode approach continuously encodes the low quality stream, which uses CPU or GPU resources. This cost only applies to the transcode path; recording the camera's native sub stream does not re-encode. See the [go2rtc hardware acceleration documentation](https://github.com/AlexxIT/go2rtc?tab=readme-ov-file#source-ffmpeg) for accelerating the transcode.
- Many camera sub streams do not include audio. If the source stream has no audio, the low quality recordings will not have audio. - Many camera sub streams do not include audio. If the source stream has no audio, the low quality recordings will not have audio.
- **Matching video codecs and audio settings between the two streams gives the smoothest playback.** When playback combines both qualities on one timeline (the default `Auto` behavior: for example original quality during events with low quality in between, or low quality history after the original recordings expire) and the streams use different video codecs or audio settings, for example H.265 on the main stream and H.264 on the sub stream, or 16 kHz audio on one and 8 kHz on the other, playback still works: Frigate inserts a decoder reset at each quality transition, which can cause a barely-perceptible pause there. Configuring both streams in the camera's firmware to use the same video codec, audio codec, and sample rate makes transitions fully seamless, and a mismatched audio sample rate can also be corrected with [sub stream output args](#sub-stream-output-args). If one stream has audio and the other does not, combined time ranges play **without audio**; selecting a single quality with the playback selector always keeps that stream's audio. - **Matching video codecs and audio settings between the two streams gives the smoothest playback.** When playback combines both qualities on one timeline (the default `Auto` behavior: for example original quality during events with low quality in between, or low quality history after the original recordings expire) and the streams use different video codecs or audio settings, for example H.265 on the main stream and H.264 on the sub stream, or 16 kHz audio on one and 8 kHz on the other, playback still works: Frigate inserts a decoder reset at each quality transition, which can cause a barely-perceptible pause there. Configuring both streams in the camera's firmware to use the same video codec, audio codec, and sample rate makes transitions fully seamless, and a mismatched audio sample rate can also be corrected with [sub stream output args](#sub-stream-output-args). If one stream has audio and the other does not, combined time ranges play **without audio**; selecting a single quality with the playback selector always keeps that stream's audio.
+5 -8
View File
@@ -548,9 +548,7 @@ services:
### Recommended security options ### Recommended security options
Frigate does not need elevated container privileges for most setups. The Frigate does not need elevated container privileges for most setups. The following hardens the container; add the `devices`/`group_add` entries your hardware requires (see the hardware acceleration docs):
following hardens the container; add the `devices`/`group_add` entries your
hardware requires (see the hardware acceleration docs):
```yaml ```yaml
services: services:
@@ -564,15 +562,14 @@ services:
:::note :::note
`telemetry.stats.network_bandwidth` uses nethogs, which requires root with `telemetry.stats.network_bandwidth` uses nethogs, which requires root with NET_ADMIN/NET_RAW capabilities. If you enable that stat, omit `cap_drop: [ALL]` or add `cap_add: [NET_ADMIN, NET_RAW]`.
NET_ADMIN/NET_RAW capabilities. If you enable that stat, omit `cap_drop: [ALL]`
or add `cap_add: [NET_ADMIN, NET_RAW]`.
Platforms that genuinely require `privileged: true` (MemryX, some QNAP setups) Platforms that genuinely require `privileged: true` (MemryX, some QNAP setups) are called out in their own sections and are unaffected by this guidance.
are called out in their own sections and are unaffected by this guidance.
::: :::
Frigate's services run as an unprivileged user inside the container. See [Running as a non-root user](../configuration/non_root.md) for the run modes, the one time volume ownership migration, and what each accelerator needs on the host.
**Docker CLI** **Docker CLI**
If you can't use Docker Compose, you can run the container with something similar to this: If you can't use Docker Compose, you can run the container with something similar to this:
+1 -1
View File
@@ -304,7 +304,7 @@ Topic with current state of notifications. Published values are `ON` and `OFF`.
### `frigate/<camera_name>/status/<role>` ### `frigate/<camera_name>/status/<role>`
Publishes the current health status of each role that is enabled (`audio`, `detect`, `record`, `record_sub`). `record_sub` is only published for cameras with [sub stream recording](/configuration/record#sub-stream-recording) enabled, and is tracked separately from `record` so a healthy main stream can't hide a stalled sub stream. Possible values are: Publishes the current health status of each role that is enabled (`audio`, `detect`, `record`). Possible values are:
- `online`: Stream is running and being processed - `online`: Stream is running and being processed
- `offline`: Stream is offline and is being restarted - `offline`: Stream is offline and is being restarted
+1
View File
@@ -122,6 +122,7 @@ const sidebars: SidebarsConfig = {
"configuration/ffmpeg_presets", "configuration/ffmpeg_presets",
"configuration/pwa", "configuration/pwa",
"configuration/tls", "configuration/tls",
"configuration/non_root",
], ],
}, },
{ {
+3 -5
View File
@@ -1476,12 +1476,10 @@ paths:
- Classification - Classification
summary: Get custom classification attributes summary: Get custom classification attributes
description: |- description: |-
**Access:** Any authenticated user. **Access:** Admin role required.
Returns custom classification attributes for a given object type. Returns custom classification attributes for a given object type.
Only includes models with classification_type set to 'attribute'. Only includes models with classification_type set to 'attribute'.
Callers without access to every camera only receive values that have been
recorded on the cameras they can access.
By default returns a flat sorted list of all attribute labels. By default returns a flat sorted list of all attribute labels.
If group_by_model is true, returns attributes grouped by model name. If group_by_model is true, returns attributes grouped by model name.
operationId: get_custom_attributes_classification_attributes_get operationId: get_custom_attributes_classification_attributes_get
@@ -1512,8 +1510,8 @@ paths:
schema: schema:
$ref: '#/components/schemas/HTTPValidationError' $ref: '#/components/schemas/HTTPValidationError'
security: security:
- frigateUserAuth: [] - frigateAdminAuth: []
x-required-role: any x-required-role: admin
/classification/{name}/train: /classification/{name}/train:
get: get:
tags: tags:
-1
View File
@@ -86,7 +86,6 @@ def require_admin_by_default():
"/categorized_object_names", "/categorized_object_names",
"/plus/models", "/plus/models",
"/recognized_license_plates", "/recognized_license_plates",
"/classification/attributes",
"/timeline", "/timeline",
"/timeline/hourly", "/timeline/hourly",
"/recordings/storage", "/recordings/storage",
+3 -96
View File
@@ -11,14 +11,10 @@ from typing import Any
import cv2 import cv2
from fastapi import APIRouter, Depends, Request, UploadFile from fastapi import APIRouter, Depends, Request, UploadFile
from fastapi.responses import JSONResponse from fastapi.responses import JSONResponse
from peewee import DoesNotExist, fn from peewee import DoesNotExist
from playhouse.shortcuts import model_to_dict from playhouse.shortcuts import model_to_dict
from frigate.api.auth import ( from frigate.api.auth import require_role
allow_any_authenticated,
get_allowed_cameras_for_filter,
require_role,
)
from frigate.api.defs.request.classification_body import ( from frigate.api.defs.request.classification_body import (
AudioTranscriptionBody, AudioTranscriptionBody,
DeleteFaceImagesBody, DeleteFaceImagesBody,
@@ -743,81 +739,18 @@ def get_classification_dataset(name: str):
) )
def get_observed_attributes(
model_attributes: dict[str, list[str]],
object_labels: set[str],
allowed_cameras: list[str],
) -> dict[str, set[str]]:
"""Get the attribute values recorded on the given cameras.
Args:
model_attributes: Labels each attribute model can emit, keyed by model name
object_labels: Object types those models run on
allowed_cameras: Cameras the caller has access to
Returns:
Values seen for each model, keyed by model name
"""
if not model_attributes or not object_labels or not allowed_cameras:
return {}
model_names = list(model_attributes.keys())
query = (
Event.select(
*[
fn.json_extract(Event.data, f'$."{model_name}"')
for model_name in model_names
]
)
.where(
(Event.camera << allowed_cameras) & (Event.label << sorted(object_labels))
)
.distinct()
.tuples()
)
targets = {
model_name: set(attributes)
for model_name, attributes in model_attributes.items()
}
observed: dict[str, set[str]] = {model_name: set() for model_name in model_names}
for row in query.iterator():
found = False
for model_name, value in zip(model_names, row):
if isinstance(value, str) and value not in observed[model_name]:
observed[model_name].add(value)
found = True
if found and all(
observed[model_name] >= targets[model_name] for model_name in model_names
):
break
return observed
@router.get( @router.get(
"/classification/attributes", "/classification/attributes",
dependencies=[Depends(allow_any_authenticated())],
summary="Get custom classification attributes", summary="Get custom classification attributes",
description="""Returns custom classification attributes for a given object type. description="""Returns custom classification attributes for a given object type.
Only includes models with classification_type set to 'attribute'. Only includes models with classification_type set to 'attribute'.
Callers without access to every camera only receive values that have been
recorded on the cameras they can access.
By default returns a flat sorted list of all attribute labels. By default returns a flat sorted list of all attribute labels.
If group_by_model is true, returns attributes grouped by model name.""", If group_by_model is true, returns attributes grouped by model name.""",
) )
def get_custom_attributes( def get_custom_attributes(
request: Request, request: Request, object_type: str = None, group_by_model: bool = False
object_type: str = None,
group_by_model: bool = False,
allowed_cameras: list[str] = Depends(get_allowed_cameras_for_filter),
): ):
models_with_attributes = {} models_with_attributes = {}
objects_by_model = {}
for ( for (
model_key, model_key,
@@ -848,32 +781,6 @@ def get_custom_attributes(
if attributes: if attributes:
model_name = model_config.name or model_key model_name = model_config.name or model_key
models_with_attributes[model_name] = sorted(attributes) models_with_attributes[model_name] = sorted(attributes)
objects_by_model[model_name] = model_objects
# the dataset holds every label a model can emit, including ones never
# applied to an event, so callers without full camera access are limited to
# the values actually recorded on the cameras they can see
all_cameras = set(request.app.frigate_config.cameras.keys())
if models_with_attributes and not all_cameras.issubset(allowed_cameras):
observed = get_observed_attributes(
models_with_attributes,
set().union(*objects_by_model.values()),
allowed_cameras,
)
models_with_attributes = {
model_name: [
attribute
for attribute in attributes
if attribute in observed.get(model_name, set())
]
for model_name, attributes in models_with_attributes.items()
}
models_with_attributes = {
model_name: attributes
for model_name, attributes in models_with_attributes.items()
if attributes
}
if group_by_model: if group_by_model:
return JSONResponse(content=models_with_attributes) return JSONResponse(content=models_with_attributes)
-16
View File
@@ -6,8 +6,6 @@ import logging
from collections.abc import Callable, Iterable from collections.abc import Callable, Iterable
from typing import Any, cast from typing import Any, cast
from peewee import IntegrityError
from frigate.camera import PTZMetrics from frigate.camera import PTZMetrics
from frigate.camera.activity_manager import AudioActivityManager, CameraActivityManager from frigate.camera.activity_manager import AudioActivityManager, CameraActivityManager
from frigate.comms.base_communicator import Communicator from frigate.comms.base_communicator import Communicator
@@ -256,21 +254,7 @@ class Dispatcher:
restart_frigate() restart_frigate()
def handle_insert_many_recordings() -> None: def handle_insert_many_recordings() -> None:
try:
Recordings.insert_many(payload).execute() Recordings.insert_many(payload).execute()
except IntegrityError:
logger.warning(
"Batch recording insert failed, inserting rows individually"
)
for recording in payload:
try:
Recordings.insert(recording).execute()
except IntegrityError:
logger.warning(
"Skipping recording that is already stored: %s",
recording.get(Recordings.path.name),
)
def handle_request_region_grid() -> Any: def handle_request_region_grid() -> Any:
camera = payload camera = payload
+1 -4
View File
@@ -18,10 +18,7 @@ class RecordingsDataTypeEnum(str, Enum):
class RecordingsDataPublisher(Publisher[Any]): class RecordingsDataPublisher(Publisher[Any]):
"""Publishes latest recording data. """Publishes latest recording data."""
Payloads are (camera, stream_type, timestamp, cache_path) on every topic.
"""
topic_base = "recordings/" topic_base = "recordings/"
+2 -7
View File
@@ -220,21 +220,16 @@ class CameraConfig(FrigateBaseModel):
# add roles to the input if there is only one # add roles to the input if there is only one
if len(config["ffmpeg"]["inputs"]) == 1: if len(config["ffmpeg"]["inputs"]) == 1:
existing_roles = config["ffmpeg"]["inputs"][0].get("roles", []) has_audio = "audio" in config["ffmpeg"]["inputs"][0].get("roles", [])
config["ffmpeg"]["inputs"][0]["roles"] = [ config["ffmpeg"]["inputs"][0]["roles"] = [
"record", "record",
"detect", "detect",
] ]
if "audio" in existing_roles: if has_audio:
config["ffmpeg"]["inputs"][0]["roles"].append("audio") config["ffmpeg"]["inputs"][0]["roles"].append("audio")
# kept so role validation can report the real problem rather than
# claiming the role was never assigned
if "record_sub" in existing_roles:
config["ffmpeg"]["inputs"][0]["roles"].append("record_sub")
super().__init__(**config) super().__init__(**config)
@property @property
+1 -8
View File
@@ -2,7 +2,7 @@ from enum import Enum
from pydantic import Field from pydantic import Field
from frigate.const import MAX_PRE_CAPTURE, STREAM_TYPE_SUB from frigate.const import MAX_PRE_CAPTURE
from frigate.review.types import SeverityEnum from frigate.review.types import SeverityEnum
from ..base import FrigateBaseModel from ..base import FrigateBaseModel
@@ -191,13 +191,6 @@ class RecordConfig(FrigateBaseModel):
description="Indicates whether recording was enabled in the original static configuration.", description="Indicates whether recording was enabled in the original static configuration.",
) )
def stream_enabled(self, stream_type: str) -> bool:
"""Whether the given record stream type should currently be recording."""
if stream_type == STREAM_TYPE_SUB:
return self.enabled and self.sub.enabled
return self.enabled
@property @property
def effective_alert_days(self) -> float: def effective_alert_days(self) -> float:
"""Alert retention extended to the sub stream window when sub is enabled. """Alert retention extended to the sub stream window when sub is enabled.
-6
View File
@@ -269,12 +269,6 @@ def verify_config_roles(camera_config: CameraConfig) -> None:
f"Camera {camera_config.name} has sub stream recording enabled, but record_sub is not assigned to an input." f"Camera {camera_config.name} has sub stream recording enabled, but record_sub is not assigned to an input."
) )
for ffmpeg_input in camera_config.ffmpeg.inputs:
if "record" in ffmpeg_input.roles and "record_sub" in ffmpeg_input.roles:
raise ValueError(
f"Camera {camera_config.name} has record and record_sub assigned to the same input, which would record the same stream twice."
)
if camera_config.audio.enabled and "audio" not in assigned_roles: if camera_config.audio.enabled and "audio" not in assigned_roles:
raise ValueError( raise ValueError(
f"Camera {camera_config.name} has audio events enabled, but audio is not assigned to an input." f"Camera {camera_config.name} has audio events enabled, but audio is not assigned to an input."
-3
View File
@@ -28,9 +28,6 @@ REDACTED_CREDENTIAL_SENTINEL = "__FRIGATE_SAVED_CREDENTIAL__"
STREAM_TYPE_MAIN = "main" STREAM_TYPE_MAIN = "main"
STREAM_TYPE_SUB = "sub" STREAM_TYPE_SUB = "sub"
SUB_CACHE_TAG = "@sub" SUB_CACHE_TAG = "@sub"
RECORD_STREAM_TYPES = (STREAM_TYPE_MAIN, STREAM_TYPE_SUB)
ROLE_TO_STREAM_TYPE = {"record": STREAM_TYPE_MAIN, "record_sub": STREAM_TYPE_SUB}
STREAM_TYPE_TO_ROLE = {v: k for k, v in ROLE_TO_STREAM_TYPE.items()}
# Attribute & Object constants # Attribute & Object constants
@@ -23,7 +23,6 @@ from frigate.const import (
ATTRIBUTE_LABEL_DISPLAY_MAP, ATTRIBUTE_LABEL_DISPLAY_MAP,
CACHE_DIR, CACHE_DIR,
CLIPS_DIR, CLIPS_DIR,
STREAM_TYPE_MAIN,
UPDATE_REVIEW_DESCRIPTION, UPDATE_REVIEW_DESCRIPTION,
) )
from frigate.data_processing.types import PostProcessDataEnum from frigate.data_processing.types import PostProcessDataEnum
@@ -442,7 +441,6 @@ class ReviewDescriptionProcessor(PostProcessorApi):
) )
.where((ts >= Recordings.start_time) & (ts <= Recordings.end_time)) .where((ts >= Recordings.start_time) & (ts <= Recordings.end_time))
.where(Recordings.camera == camera) .where(Recordings.camera == camera)
.where(Recordings.stream_type == STREAM_TYPE_MAIN)
.order_by(Recordings.start_time.desc()) .order_by(Recordings.start_time.desc())
.limit(1) .limit(1)
.get() .get()
+1 -3
View File
@@ -714,9 +714,7 @@ class EmbeddingMaintainer(threading.Thread):
topic = str(raw_topic) topic = str(raw_topic)
if topic.endswith(RecordingsDataTypeEnum.saved.value): if topic.endswith(RecordingsDataTypeEnum.saved.value):
camera, _stream_type, recordings_available_through_timestamp, _ = ( camera, recordings_available_through_timestamp, _ = payload
payload
)
self.recordings_available_through[camera] = ( self.recordings_available_through[camera] = (
recordings_available_through_timestamp recordings_available_through_timestamp
+7 -34
View File
@@ -149,12 +149,8 @@ class RecordingCleanup(threading.Thread):
detections_retain_mode: RetainModeEnum, detections_retain_mode: RetainModeEnum,
config: CameraConfig, config: CameraConfig,
reviews: list[Any], reviews: list[Any],
) -> tuple[set[Path], list[tuple[float, float]]]: ) -> set[Path]:
"""Delete recordings for one stream of an existing camera based on retention config. """Delete recordings for existing camera based on retention config."""
Returns the directories to check for emptiness and the segments that
were kept, which the caller feeds to expire_camera_previews.
"""
# Get the timestamp for cutoff of retained days # Get the timestamp for cutoff of retained days
# Get recordings to check for expiration # Get recordings to check for expiration
@@ -261,23 +257,9 @@ class RecordingCleanup(threading.Thread):
Recordings.id << deleted_recordings_list[i : i + max_deletes] Recordings.id << deleted_recordings_list[i : i + max_deletes]
).execute() ).execute()
return maybe_empty_dirs, kept_recordings # previews follow main retention, so only the main pass expires them
if stream_type != STREAM_TYPE_MAIN:
def expire_camera_previews( return maybe_empty_dirs
self,
config: CameraConfig,
continuous_expire_date: float,
motion_expire_date: float,
kept_recordings: list[tuple[float, float]],
) -> set[Path]:
"""Delete previews that no longer have recordings on any stream.
Previews aren't recorded per stream, so the cutoffs must be the oldest
of the per stream values and kept_recordings must cover every stream,
sorted by start time. Otherwise a short main retention expires previews
the sub recordings still need.
"""
maybe_empty_dirs: set[Path] = set()
previews = ( previews = (
Previews.select( Previews.select(
@@ -456,7 +438,7 @@ class RecordingCleanup(threading.Thread):
.namedtuples() .namedtuples()
) )
main_dirs, main_kept = self.expire_existing_camera_recordings( maybe_empty_dirs |= self.expire_existing_camera_recordings(
STREAM_TYPE_MAIN, STREAM_TYPE_MAIN,
continuous_expire_date, continuous_expire_date,
motion_expire_date, motion_expire_date,
@@ -470,11 +452,10 @@ class RecordingCleanup(threading.Thread):
config.record.detections.retain.days, config.record.detections.retain.days,
), ),
) )
maybe_empty_dirs |= main_dirs
# runs even when sub recording is disabled so old rows still # runs even when sub recording is disabled so old rows still
# expire # expire
sub_dirs, sub_kept = self.expire_existing_camera_recordings( maybe_empty_dirs |= self.expire_existing_camera_recordings(
STREAM_TYPE_SUB, STREAM_TYPE_SUB,
sub_continuous_expire_date, sub_continuous_expire_date,
sub_motion_expire_date, sub_motion_expire_date,
@@ -488,14 +469,6 @@ class RecordingCleanup(threading.Thread):
config.record.sub.detections.days, config.record.sub.detections.days,
), ),
) )
maybe_empty_dirs |= sub_dirs
maybe_empty_dirs |= self.expire_camera_previews(
config,
min(continuous_expire_date, sub_continuous_expire_date),
min(motion_expire_date, sub_motion_expire_date),
sorted(main_kept + sub_kept),
)
logger.debug(f"End camera: {camera}.") logger.debug(f"End camera: {camera}.")
logger.debug("End all cameras.") logger.debug("End all cameras.")
+31 -86
View File
@@ -79,54 +79,6 @@ def parse_cache_segment_name(basename: str) -> tuple[str, str, str] | None:
return (prefix, STREAM_TYPE_MAIN, date) return (prefix, STREAM_TYPE_MAIN, date)
def format_segment_details(cache_path: str, segment_info: dict[str, Any]) -> str:
"""Comma separated facts about a segment, for discard warnings."""
details: list[str] = []
duration = segment_info.get("duration", -1)
if duration != -1:
details.append(f"duration: {duration:.2f}s")
try:
details.append(f"size: {os.path.getsize(cache_path) / 1024:.1f} KB")
except OSError:
pass
details.append(f"video: {segment_info.get('video_codec') or 'none'}")
if segment_info.get("has_audio"):
audio = segment_info.get("audio_codec") or "unknown"
rate = segment_info.get("audio_rate")
details.append(f"audio: {audio} {rate}Hz" if rate else f"audio: {audio}")
else:
details.append("audio: none")
return ", ".join(details)
def segment_path_time(cache_path: str) -> datetime.datetime | None:
"""Timestamp a segment's recording path is built from, or None if unparsable.
Recording paths carry one second of resolution, and so does ffmpeg's cache
segment template, which makes a cache file name unique per camera stream
and second. Resolved start times are not: a stream cutting segments faster
than once a second resolves consecutive segments into the same second, and
building the path from those collides on the unique path index.
"""
parsed = parse_cache_segment_name(Path(cache_path).stem)
if parsed is None:
return None
try:
return datetime.datetime.strptime(parsed[2], CACHE_SEGMENT_FORMAT).astimezone(
datetime.UTC
)
except ValueError:
return None
class SegmentInfo: class SegmentInfo:
def __init__( def __init__(
self, self,
@@ -289,53 +241,49 @@ class RecordingMaintainer(threading.Thread):
and not d.startswith("preview_") and not d.startswith("preview_")
] ]
# publish newest cached segment per camera stream (including in use files) # publish newest cached segment per camera (including in use files)
newest_cache_segments: dict[tuple[str, str], dict[str, Any]] = {} newest_cache_segments: dict[str, dict[str, Any]] = {}
for cache in cache_files: for cache in cache_files:
cache_path = os.path.join(CACHE_DIR, cache) cache_path = os.path.join(CACHE_DIR, cache)
basename = os.path.splitext(cache)[0] basename = os.path.splitext(cache)[0]
parsed = parse_cache_segment_name(basename) parsed = parse_cache_segment_name(basename)
if parsed is None: if parsed is None:
if not self.unexpected_cache_files_logged: if not self.unexpected_cache_files_logged:
logger.warning(f"Skipping unexpected files in cache, e.g. {cache}") logger.warning("Skipping unexpected files in cache")
self.unexpected_cache_files_logged = True self.unexpected_cache_files_logged = True
continue continue
camera, stream_type, date = parsed camera, stream_type, date = parsed
# this topic feeds main-stream health/sync consumers only
if stream_type == STREAM_TYPE_SUB:
continue
start_time = datetime.datetime.strptime( start_time = datetime.datetime.strptime(
date, CACHE_SEGMENT_FORMAT date, CACHE_SEGMENT_FORMAT
).astimezone(datetime.UTC) ).astimezone(datetime.UTC)
key = (camera, stream_type)
if ( if (
key not in newest_cache_segments camera not in newest_cache_segments
or start_time > newest_cache_segments[key]["start_time"] or start_time > newest_cache_segments[camera]["start_time"]
): ):
newest_cache_segments[key] = { newest_cache_segments[camera] = {
"start_time": start_time, "start_time": start_time,
"cache_path": cache_path, "cache_path": cache_path,
} }
for (camera, stream_type), newest in newest_cache_segments.items(): for camera, newest in newest_cache_segments.items():
self.recordings_publisher.publish( self.recordings_publisher.publish(
( (
camera, camera,
stream_type,
newest["start_time"].timestamp(), newest["start_time"].timestamp(),
newest["cache_path"], newest["cache_path"],
), ),
RecordingsDataTypeEnum.latest.value, RecordingsDataTypeEnum.latest.value,
) )
# publish None for streams with no cache files (but only if we know the camera exists) # publish None for cameras with no cache files (but only if we know the camera exists)
for camera_name, camera_config in self.config.cameras.items(): for camera_name in self.config.cameras:
stream_types = [STREAM_TYPE_MAIN] if camera_name not in newest_cache_segments:
if camera_config.record.sub.enabled:
stream_types.append(STREAM_TYPE_SUB)
for stream_type in stream_types:
if (camera_name, stream_type) not in newest_cache_segments:
self.recordings_publisher.publish( self.recordings_publisher.publish(
(camera_name, stream_type, None, None), (camera_name, None, None),
RecordingsDataTypeEnum.latest.value, RecordingsDataTypeEnum.latest.value,
) )
@@ -366,7 +314,7 @@ class RecordingMaintainer(threading.Thread):
parsed = parse_cache_segment_name(basename) parsed = parse_cache_segment_name(basename)
if parsed is None: if parsed is None:
if not self.unexpected_cache_files_logged: if not self.unexpected_cache_files_logged:
logger.warning(f"Skipping unexpected files in cache, e.g. {cache}") logger.warning("Skipping unexpected files in cache")
self.unexpected_cache_files_logged = True self.unexpected_cache_files_logged = True
continue continue
camera, stream_type, date = parsed camera, stream_type, date = parsed
@@ -499,7 +447,6 @@ class RecordingMaintainer(threading.Thread):
self.recordings_publisher.publish( self.recordings_publisher.publish(
( (
camera, camera,
stream_type,
recordings[0]["start_time"].timestamp() recordings[0]["start_time"].timestamp()
if camera_cfg and camera_cfg.record.enabled if camera_cfg and camera_cfg.record.enabled
else None, else None,
@@ -593,11 +540,11 @@ class RecordingMaintainer(threading.Thread):
if not segment_info.get("has_valid_video", False): if not segment_info.get("has_valid_video", False):
logger.warning( logger.warning(
f"Invalid or missing video stream in segment {cache_path} " f"Invalid or missing video stream in segment {cache_path}. Discarding."
f"({format_segment_details(cache_path, segment_info)}). Discarding."
) )
if stream_type == STREAM_TYPE_MAIN:
self.recordings_publisher.publish( self.recordings_publisher.publish(
(camera, stream_type, start_time.timestamp(), cache_path), (camera, start_time.timestamp(), cache_path),
RecordingsDataTypeEnum.invalid.value, RecordingsDataTypeEnum.invalid.value,
) )
self.drop_segment(cache_path) self.drop_segment(cache_path)
@@ -636,20 +583,19 @@ class RecordingMaintainer(threading.Thread):
if duration == -1: if duration == -1:
logger.warning(f"Failed to probe corrupt segment {cache_path}") logger.warning(f"Failed to probe corrupt segment {cache_path}")
logger.warning( logger.warning(f"Discarding a corrupt recording segment: {cache_path}")
f"Discarding a corrupt recording segment: {cache_path} " if stream_type == STREAM_TYPE_MAIN:
f"({format_segment_details(cache_path, segment_info)})"
)
self.recordings_publisher.publish( self.recordings_publisher.publish(
(camera, stream_type, start_time.timestamp(), cache_path), (camera, start_time.timestamp(), cache_path),
RecordingsDataTypeEnum.invalid.value, RecordingsDataTypeEnum.invalid.value,
) )
self.drop_segment(cache_path) self.drop_segment(cache_path)
return None return None
# this segment has a valid duration and has video data, so publish an update # this segment has a valid duration and has video data, so publish an update
if stream_type == STREAM_TYPE_MAIN:
self.recordings_publisher.publish( self.recordings_publisher.publish(
(camera, stream_type, start_time.timestamp(), cache_path), (camera, start_time.timestamp(), cache_path),
RecordingsDataTypeEnum.valid.value, RecordingsDataTypeEnum.valid.value,
) )
@@ -917,20 +863,18 @@ class RecordingMaintainer(threading.Thread):
video_codec: str | None = None, video_codec: str | None = None,
keyframes: list[int] | None = None, keyframes: list[int] | None = None,
) -> dict[str, Any] | None: ) -> dict[str, Any] | None:
path_time = segment_path_time(cache_path) or start_time # directory will be in utc due to start_time being in utc
# directory will be in utc due to path_time being in utc
# sub segments get a tagged directory to avoid filename collisions # sub segments get a tagged directory to avoid filename collisions
directory = os.path.join( directory = os.path.join(
RECORD_DIR, RECORD_DIR,
path_time.strftime("%Y-%m-%d/%H"), start_time.strftime("%Y-%m-%d/%H"),
camera if stream_type == STREAM_TYPE_MAIN else f"{camera}{SUB_CACHE_TAG}", camera if stream_type == STREAM_TYPE_MAIN else f"{camera}{SUB_CACHE_TAG}",
) )
os.makedirs(directory, exist_ok=True) os.makedirs(directory, exist_ok=True)
# file will be in utc due to path_time being in utc # file will be in utc due to start_time being in utc
file_name = f"{path_time.strftime('%M.%S.mp4')}" file_name = f"{start_time.strftime('%M.%S.mp4')}"
file_path = os.path.join(directory, file_name) file_path = os.path.join(directory, file_name)
try: try:
@@ -1002,9 +946,10 @@ class RecordingMaintainer(threading.Thread):
Recordings.video_codec.name: video_codec, Recordings.video_codec.name: video_codec,
Recordings.keyframes.name: keyframes, Recordings.keyframes.name: keyframes,
} }
except Exception: except Exception as e:
logger.exception(f"Unable to store recording segment {cache_path}") logger.error(f"Unable to store recording segment {cache_path}")
Path(cache_path).unlink(missing_ok=True) Path(cache_path).unlink(missing_ok=True)
logger.error(e)
# clear end_time cache # clear end_time cache
self.end_time_cache.pop(cache_path, None) self.end_time_cache.pop(cache_path, None)
@@ -1,189 +0,0 @@
"""Tests for GET /classification/attributes."""
import os
import shutil
import unittest
from frigate.api.auth import get_allowed_cameras_for_filter
from frigate.const import CLIPS_DIR
from frigate.models import Event, Recordings, ReviewSegment
from frigate.test.http_api.base_http_test import AuthTestClient, BaseTestHttp
# "limited_user" only reaches front_door, so it never sees the values that were
# recorded on back_door.
_CONFIG = {
"mqtt": {"host": "mqtt"},
"auth": {"roles": {"limited_user": ["front_door"]}},
"classification": {
"custom": {
"delivery_service": {
"enabled": True,
"object_config": {
"objects": ["car"],
"classification_type": "attribute",
},
}
}
},
"cameras": {
"front_door": {
"ffmpeg": {
"inputs": [{"path": "rtsp://10.0.0.1:554/video", "roles": ["detect"]}]
},
"detect": {"height": 1080, "width": 1920, "fps": 5},
},
"back_door": {
"ffmpeg": {
"inputs": [{"path": "rtsp://10.0.0.2:554/video", "roles": ["detect"]}]
},
"detect": {"height": 1080, "width": 1920, "fps": 5},
},
},
}
class TestClassificationAttributesAccess(BaseTestHttp):
"""The attribute list is read from the training dataset on disk, which holds
every label a model can emit regardless of which camera recorded it. Callers
without full camera access are cut back to the values on their own cameras,
so these tests pin that scoping.
"""
def setUp(self):
super().setUp([Event, ReviewSegment, Recordings])
self.minimal_config = _CONFIG
self.app = super().create_app()
self.model_dir = os.path.join(CLIPS_DIR, "delivery_service")
for category in ("DHL", "Amazon", "Hermes", "none"):
os.makedirs(
os.path.join(self.model_dir, "dataset", category), exist_ok=True
)
def tearDown(self):
shutil.rmtree(self.model_dir, ignore_errors=True)
self.app.dependency_overrides.clear()
super().tearDown()
def _insert_event(self, event_id: str, camera: str, attribute: str | None):
data = {"type": "object", "score": 0.9}
if attribute is not None:
data["delivery_service"] = attribute
Event.insert(
id=event_id,
label="car",
camera=camera,
start_time=100,
end_time=200,
top_score=0.9,
score=0.9,
false_positive=False,
zones=[],
thumbnail="",
has_clip=True,
has_snapshot=True,
region=[],
box=[],
area=0,
retain_indefinitely=False,
ratio=1.0,
plus_id=None,
model_hash="",
detector_type="cpu",
model_type="ssd",
data=data,
).execute()
def _get(self, role: str, **params):
# the base class resolves every camera by default, so drop the override
# to exercise the real role to allowed-cameras resolution
self.app.dependency_overrides.pop(get_allowed_cameras_for_filter, None)
with AuthTestClient(self.app) as client:
return client.get(
"/classification/attributes",
params=params,
headers={"remote-user": "test", "remote-role": role},
)
def _insert_split_events(self):
self._insert_event("front", "front_door", "DHL")
self._insert_event("back", "back_door", "Amazon")
def test_admin_gets_every_trained_label(self):
self._insert_split_events()
assert self._get("admin").json() == ["Amazon", "DHL", "Hermes"]
def test_viewer_gets_every_trained_label(self):
self._insert_split_events()
assert self._get("viewer").json() == ["Amazon", "DHL", "Hermes"]
def test_restricted_role_only_gets_its_own_cameras(self):
self._insert_split_events()
assert self._get("limited_user").json() == ["DHL"]
def test_restricted_role_grouped_by_model(self):
self._insert_split_events()
assert self._get("limited_user", group_by_model="true").json() == {
"delivery_service": ["DHL"]
}
def test_restricted_role_with_no_recorded_values(self):
self._insert_event("back", "back_door", "Amazon")
assert self._get("limited_user").json() == []
assert self._get("limited_user", group_by_model="true").json() == {}
def test_restricted_role_ignores_events_without_the_attribute(self):
self._insert_event("front", "front_door", None)
assert self._get("limited_user").json() == []
def test_restricted_role_with_a_dotted_model_name(self):
# model names are unrestricted config keys, and an unquoted "." in the
# json path would be read as a nested lookup and match nothing
self.app.frigate_config.classification.custom["delivery.service"] = (
self.app.frigate_config.classification.custom.pop("delivery_service")
)
self.app.frigate_config.classification.custom[
"delivery.service"
].name = "delivery.service"
os.rename(self.model_dir, os.path.join(CLIPS_DIR, "delivery.service"))
self.model_dir = os.path.join(CLIPS_DIR, "delivery.service")
data = {"type": "object", "score": 0.9, "delivery.service": "DHL"}
Event.insert(
id="front",
label="car",
camera="front_door",
start_time=100,
end_time=200,
top_score=0.9,
score=0.9,
false_positive=False,
zones=[],
thumbnail="",
has_clip=True,
has_snapshot=True,
region=[],
box=[],
area=0,
retain_indefinitely=False,
ratio=1.0,
plus_id=None,
model_hash="",
detector_type="cpu",
model_type="ssd",
data=data,
).execute()
assert self._get("limited_user").json() == ["DHL"]
def test_object_type_filters_out_unrelated_models(self):
self._insert_split_events()
assert self._get("limited_user", object_type="person").json() == []
assert self._get("limited_user", object_type="car").json() == ["DHL"]
if __name__ == "__main__":
unittest.main()
+26
View File
@@ -0,0 +1,26 @@
"""Tests for bandwidth stats privilege handling."""
import unittest
from unittest.mock import MagicMock, patch
from frigate.util import services
class TestBandwidthStatsPrivileges(unittest.TestCase):
def setUp(self):
services._bandwidth_warning_logged = False
@patch("frigate.util.services.sp.run")
@patch("frigate.util.services.os.geteuid", return_value=1000)
def test_returns_empty_and_warns_once_without_root(self, _, sp_run):
config = MagicMock()
with self.assertLogs("frigate.util.services", level="WARNING") as logs:
assert services.get_bandwidth_stats(config) == {}
assert services.get_bandwidth_stats(config) == {}
sp_run.assert_not_called()
warnings = [m for m in logs.output if "require root" in m]
assert len(warnings) == 1
if __name__ == "__main__":
unittest.main()
-220
View File
@@ -1,220 +0,0 @@
"""Tests for per stream recording health tracking in the camera watchdog."""
import unittest
from datetime import UTC, datetime, timedelta
from unittest.mock import MagicMock, patch
from frigate.config import FrigateConfig
from frigate.const import STREAM_TYPE_MAIN, STREAM_TYPE_SUB
from frigate.video.ffmpeg import CameraWatchdog
class TestCameraWatchdogStreamHealth(unittest.TestCase):
def _build_watchdog(
self, sub_enabled: bool = True, output_args: dict | None = None
) -> CameraWatchdog:
config = FrigateConfig(
**{
"mqtt": {"host": "mqtt"},
"cameras": {
"front_door": {
"ffmpeg": {
"output_args": output_args or {},
"inputs": [
{
"path": "rtsp://10.0.0.1:554/video",
"roles": ["record"],
},
{
"path": "rtsp://10.0.0.1:554/video2",
"roles": ["detect", "record_sub"],
},
],
},
"record": {
"enabled": True,
"sub": {"enabled": sub_enabled},
},
}
},
}
)
camera_config = config.cameras["front_door"]
with (
patch("frigate.video.ffmpeg.LogPipe"),
patch("frigate.video.ffmpeg.InterProcessRequestor"),
patch("frigate.video.ffmpeg.RecordingsDataSubscriber"),
patch("frigate.video.ffmpeg.CameraConfigUpdateSubscriber"),
):
watchdog = CameraWatchdog(
camera_config,
1,
MagicMock(),
MagicMock(),
MagicMock(),
MagicMock(),
MagicMock(),
MagicMock(),
MagicMock(),
MagicMock(),
)
watchdog.requestor = MagicMock()
return watchdog
def test_stale_sub_does_not_mark_main_stale(self):
watchdog = self._build_watchdog()
now = datetime.now().astimezone(UTC)
stale = (now - timedelta(hours=1)).timestamp()
watchdog.latest_cache_segment_time[STREAM_TYPE_MAIN] = now.timestamp()
watchdog.latest_valid_segment_time[STREAM_TYPE_MAIN] = now.timestamp()
watchdog.latest_cache_segment_time[STREAM_TYPE_SUB] = stale
watchdog.latest_valid_segment_time[STREAM_TYPE_SUB] = stale
assert watchdog._stream_staleness(STREAM_TYPE_MAIN, now) is None
assert watchdog._stream_staleness(STREAM_TYPE_SUB, now) is not None
def test_stale_main_does_not_mark_sub_stale(self):
watchdog = self._build_watchdog()
now = datetime.now().astimezone(UTC)
stale = (now - timedelta(hours=1)).timestamp()
watchdog.latest_cache_segment_time[STREAM_TYPE_MAIN] = stale
watchdog.latest_valid_segment_time[STREAM_TYPE_MAIN] = stale
watchdog.latest_cache_segment_time[STREAM_TYPE_SUB] = now.timestamp()
watchdog.latest_valid_segment_time[STREAM_TYPE_SUB] = now.timestamp()
assert watchdog._stream_staleness(STREAM_TYPE_MAIN, now) is not None
assert watchdog._stream_staleness(STREAM_TYPE_SUB, now) is None
def test_grace_period_suppresses_staleness(self):
watchdog = self._build_watchdog()
now = datetime.now().astimezone(UTC)
watchdog.record_enable_time = now - timedelta(seconds=10)
watchdog.latest_cache_segment_time[STREAM_TYPE_SUB] = (
now - timedelta(hours=1)
).timestamp()
assert watchdog._stream_staleness(STREAM_TYPE_SUB, now) is None
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.requestor.send_data.assert_any_call(
"front_door/status/record", "online"
)
watchdog.requestor.send_data.assert_any_call(
"front_door/status/record_sub", "offline"
)
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)
assert watchdog.requestor.send_data.call_count == 2
def test_recorded_streams_follows_config(self):
watchdog = self._build_watchdog()
assert watchdog._recorded_streams(["record"]) == [STREAM_TYPE_MAIN]
assert watchdog._recorded_streams(["detect", "record_sub"]) == [STREAM_TYPE_SUB]
assert watchdog._recorded_streams(["detect"]) == []
disabled = self._build_watchdog(sub_enabled=False)
assert disabled._recorded_streams(["detect", "record_sub"]) == []
def test_restart_grace_suppresses_repeat_staleness(self):
watchdog = self._build_watchdog()
now = datetime.now().astimezone(UTC)
stale = (now - timedelta(hours=1)).timestamp()
watchdog.latest_cache_segment_time[STREAM_TYPE_MAIN] = stale
watchdog.latest_valid_segment_time[STREAM_TYPE_MAIN] = stale
assert watchdog._stream_staleness(STREAM_TYPE_MAIN, now) is not None
watchdog._grant_restart_grace([STREAM_TYPE_MAIN], now)
assert watchdog._stream_staleness(STREAM_TYPE_MAIN, now) is None
assert (
watchdog._stream_staleness(STREAM_TYPE_MAIN, now + timedelta(seconds=89))
is None
)
assert (
watchdog._stream_staleness(STREAM_TYPE_MAIN, now + timedelta(seconds=91))
is not None
)
def test_restart_grace_is_per_stream(self):
watchdog = self._build_watchdog()
now = datetime.now().astimezone(UTC)
stale = (now - timedelta(hours=1)).timestamp()
for stream_type in (STREAM_TYPE_MAIN, STREAM_TYPE_SUB):
watchdog.latest_cache_segment_time[stream_type] = stale
watchdog.latest_valid_segment_time[stream_type] = stale
watchdog._grant_restart_grace([STREAM_TYPE_MAIN], now)
assert watchdog._stream_staleness(STREAM_TYPE_MAIN, now) is None
assert watchdog._stream_staleness(STREAM_TYPE_SUB, now) is not None
def test_detect_reset_grants_the_shared_sub_stream_grace(self):
watchdog = self._build_watchdog()
watchdog.detect_process_records_sub = True
watchdog.ffmpeg_detect_process = MagicMock()
watchdog.capture_thread = MagicMock()
watchdog.capture_thread.is_alive.return_value = False
watchdog.start_ffmpeg_detect = MagicMock()
now = datetime.now().astimezone(UTC)
stale = (now - timedelta(hours=1)).timestamp()
watchdog.latest_cache_segment_time[STREAM_TYPE_SUB] = stale
watchdog.latest_valid_segment_time[STREAM_TYPE_SUB] = stale
assert watchdog._stream_staleness(STREAM_TYPE_SUB, now) is not None
watchdog.reset_capture_thread(terminate=False)
# the sub check runs later in the same tick against a stale can_restart,
# so without this grace it would kill the just-restarted process again
assert (
watchdog._stream_staleness(STREAM_TYPE_SUB, datetime.now().astimezone(UTC))
is None
)
def test_detect_reset_leaves_sub_alone_when_not_shared(self):
watchdog = self._build_watchdog()
watchdog.detect_process_records_sub = False
watchdog.ffmpeg_detect_process = MagicMock()
watchdog.capture_thread = MagicMock()
watchdog.capture_thread.is_alive.return_value = False
watchdog.start_ffmpeg_detect = MagicMock()
now = datetime.now().astimezone(UTC)
stale = (now - timedelta(hours=1)).timestamp()
watchdog.latest_cache_segment_time[STREAM_TYPE_SUB] = stale
watchdog.latest_valid_segment_time[STREAM_TYPE_SUB] = stale
watchdog.reset_capture_thread(terminate=False)
assert (
watchdog._stream_staleness(STREAM_TYPE_SUB, datetime.now().astimezone(UTC))
is not None
)
def test_stale_threshold_follows_each_stream_segment_time(self):
watchdog = self._build_watchdog(
output_args={
"record": "-f segment -segment_time 10 -c copy",
"record_sub": "-f segment -segment_time 60 -c copy",
}
)
assert watchdog.record_stale_threshold[STREAM_TYPE_MAIN] == 120
assert watchdog.record_stale_threshold[STREAM_TYPE_SUB] == 150
-30
View File
@@ -1229,36 +1229,6 @@ class TestConfig(unittest.TestCase):
lambda: FrigateConfig(**config).cameras, lambda: FrigateConfig(**config).cameras,
) )
def test_fails_on_record_and_record_sub_on_same_input(self):
config = self._sub_record_config()
config["cameras"]["back"]["ffmpeg"]["inputs"] = [
{
"path": "rtsp://10.0.0.1:554/video",
"roles": ["detect", "record", "record_sub"],
},
{"path": "rtsp://10.0.0.1:554/video2", "roles": ["audio"]},
]
self.assertRaisesRegex(
ValueError,
"record and record_sub assigned to the same input",
lambda: FrigateConfig(**config).cameras,
)
def test_fails_on_record_sub_with_a_single_input(self):
# the single input case has record forced onto it, so record_sub can
# only ever duplicate that same stream
config = self._sub_record_config()
config["cameras"]["back"]["ffmpeg"]["inputs"] = [
{"path": "rtsp://10.0.0.1:554/video", "roles": ["detect", "record_sub"]},
]
self.assertRaisesRegex(
ValueError,
"record and record_sub assigned to the same input",
lambda: FrigateConfig(**config).cameras,
)
def test_record_sub_segment_time_not_checked_when_disabled(self): def test_record_sub_segment_time_not_checked_when_disabled(self):
config = self._sub_record_config( config = self._sub_record_config(
{ {
@@ -1,67 +0,0 @@
"""Tests for the recordings batch insert handler."""
import unittest
from unittest.mock import MagicMock, patch
from playhouse.sqlite_ext import SqliteExtDatabase
from frigate.comms.dispatcher import Dispatcher
from frigate.const import INSERT_MANY_RECORDINGS
from frigate.models import Recordings
def _recording(id: str, path: str) -> dict:
return {
Recordings.id.name: id,
Recordings.camera.name: "front_door",
Recordings.stream_type.name: "main",
Recordings.path.name: path,
Recordings.start_time.name: 1000.0,
Recordings.end_time.name: 1010.0,
Recordings.duration.name: 10.0,
Recordings.motion.name: 0,
Recordings.objects.name: 0,
Recordings.dBFS.name: 0,
Recordings.segment_size.name: 1.0,
}
class TestInsertManyRecordings(unittest.TestCase):
"""A duplicate path must not cost the rest of the batch."""
def setUp(self):
self.db = SqliteExtDatabase(":memory:")
self.db.bind([Recordings])
self.db.create_tables([Recordings])
with (
patch("frigate.comms.dispatcher.CameraActivityManager"),
patch("frigate.comms.dispatcher.AudioActivityManager"),
):
self.dispatcher = Dispatcher(MagicMock(), MagicMock(), MagicMock(), {}, [])
def tearDown(self):
self.db.close()
def test_batch_with_duplicate_keeps_the_other_rows(self):
Recordings.insert(_recording("existing", "/rec/00.10.mp4")).execute()
self.dispatcher._receive(
INSERT_MANY_RECORDINGS,
[
_recording("a", "/rec/00.20.mp4"),
_recording("b", "/rec/00.10.mp4"),
_recording("c", "/rec/00.30.mp4"),
],
)
paths = {r.path for r in Recordings.select()}
self.assertEqual(paths, {"/rec/00.10.mp4", "/rec/00.20.mp4", "/rec/00.30.mp4"})
def test_clean_batch_inserts_every_row(self):
self.dispatcher._receive(
INSERT_MANY_RECORDINGS,
[_recording("a", "/rec/00.20.mp4"), _recording("b", "/rec/00.30.mp4")],
)
self.assertEqual(Recordings.select().count(), 2)
-69
View File
@@ -73,21 +73,6 @@ class TestRecordingCleanupSubRetention(unittest.TestCase):
stream_type=stream_type, stream_type=stream_type,
) )
def _insert_preview(
self, id: str, age_days: float, camera: str = "front_door"
) -> None:
end_time = (
datetime.datetime.now() - datetime.timedelta(days=age_days)
).timestamp()
Previews.create(
id=id,
camera=camera,
path=f"/media/frigate/previews/{id}.mp4",
start_time=end_time - 10,
end_time=end_time,
duration=10,
)
def test_sub_recordings_expire_independently(self): def test_sub_recordings_expire_independently(self):
# main retention 7 days, sub retention 30 days; rows 10 days old # main retention 7 days, sub retention 30 days; rows 10 days old
# -> main row deleted, sub row kept # -> main row deleted, sub row kept
@@ -106,60 +91,6 @@ class TestRecordingCleanupSubRetention(unittest.TestCase):
assert Recordings.get_or_none(Recordings.id == "m1") is None assert Recordings.get_or_none(Recordings.id == "m1") is None
assert Recordings.get_or_none(Recordings.id == "s1") is not None assert Recordings.get_or_none(Recordings.id == "s1") is not None
def test_previews_survive_while_sub_recordings_remain(self):
# main retention 7 days, sub retention 30 days; only the sub row
# survives at 10 days, and the preview covering it must survive too
cleanup = self._build_cleanup(
{
"enabled": True,
"continuous": {"days": 7},
"sub": {"enabled": True, "continuous": {"days": 30}},
}
)
self._insert_recording("m1", "main", 10)
self._insert_recording("s1", "sub", 10)
self._insert_preview("p1", 10)
cleanup.expire_recordings()
assert Recordings.get_or_none(Recordings.id == "m1") is None
assert Previews.get_or_none(Previews.id == "p1") is not None
def test_previews_expire_once_every_stream_has(self):
# both streams expired at 40 days -> the preview goes with them
cleanup = self._build_cleanup(
{
"enabled": True,
"continuous": {"days": 7},
"sub": {"enabled": True, "continuous": {"days": 30}},
}
)
self._insert_recording("m1", "main", 40)
self._insert_recording("s1", "sub", 40)
self._insert_preview("p1", 40)
cleanup.expire_recordings()
assert Recordings.get_or_none(Recordings.id == "s1") is None
assert Previews.get_or_none(Previews.id == "p1") is None
def test_preview_retention_unchanged_when_sub_disabled(self):
cleanup = self._build_cleanup(
{
"enabled": True,
"continuous": {"days": 7},
"sub": {"enabled": False},
}
)
self._insert_recording("m1", "main", 10)
self._insert_preview("p_old", 10)
self._insert_preview("p_new", 1)
cleanup.expire_recordings()
assert Previews.get_or_none(Previews.id == "p_old") is None
assert Previews.get_or_none(Previews.id == "p_new") is not None
def test_sub_recordings_expire_after_sub_retention(self): def test_sub_recordings_expire_after_sub_retention(self):
# sub retention 30 days; sub row 40 days old -> deleted # sub retention 30 days; sub row 40 days old -> deleted
cleanup = self._build_cleanup( cleanup = self._build_cleanup(
@@ -15,7 +15,6 @@ from frigate.record.maintainer import (
RecordingMaintainer, RecordingMaintainer,
SegmentInfo, SegmentInfo,
parse_cache_segment_name, parse_cache_segment_name,
segment_path_time,
) )
@@ -350,98 +349,6 @@ class TestSegmentAudioPresence(unittest.IsolatedAsyncioTestCase):
self.assertEqual(result[Recordings.video_codec.name], video_codec) self.assertEqual(result[Recordings.video_codec.name], video_codec)
class TestSegmentPathTime(unittest.IsolatedAsyncioTestCase):
"""The recording path must stay unique when segments are shorter than a second."""
def _build_maintainer(self) -> RecordingMaintainer:
camera_config = MagicMock()
camera_config.record.enabled = True
camera_config.record.continuous.days = 1
camera_config.record.motion.days = 0
config = MagicMock()
config.cameras = {"test_cam": camera_config}
maintainer = RecordingMaintainer.__new__(RecordingMaintainer)
maintainer.config = config
maintainer.end_time_cache = {}
maintainer.object_recordings_info = defaultdict(list)
maintainer.audio_recordings_info = defaultdict(list)
maintainer.recordings_publisher = MagicMock()
maintainer.last_segment_end = {("test_cam", "main"): 0.0}
return maintainer
def test_parses_main_and_sub_names(self):
expected = datetime.datetime(2026, 6, 10, 14, 30, 22, tzinfo=datetime.UTC)
self.assertEqual(
segment_path_time("/tmp/cache/test_cam@20260610143022+0000.mp4"), expected
)
self.assertEqual(
segment_path_time("/tmp/cache/test_cam@sub@20260610143022+0000.mp4"),
expected,
)
def test_returns_none_for_unparsable_names(self):
self.assertIsNone(segment_path_time("/tmp/cache/garbage.mp4"))
self.assertIsNone(segment_path_time("/tmp/cache/test_cam@notadate.mp4"))
async def test_sub_second_segments_get_distinct_paths(self):
# two cache files a second apart whose resolved starts both land in
# second 22; deriving the path from the resolved start collides
segments = [
("test_cam@20260610143022+0000.mp4", 100_000),
("test_cam@20260610143023+0000.mp4", 980_000),
]
paths = []
with tempfile.TemporaryDirectory() as tmpdir:
for name, microsecond in segments:
maintainer = self._build_maintainer()
maintainer.config.ffmpeg.ffmpeg_path = "ffmpeg"
start_time = datetime.datetime(
2026, 6, 10, 14, 30, 22, microsecond, tzinfo=datetime.UTC
)
cache_path = os.path.join(tmpdir, name)
with open(cache_path, "wb") as f:
f.write(b"\x00" * 16)
proc = MagicMock()
proc.returncode = 0
proc.wait = AsyncMock(return_value=0)
with (
patch(
"frigate.record.maintainer.RECORD_DIR",
os.path.join(tmpdir, "recordings"),
),
patch(
"frigate.record.maintainer.asyncio.create_subprocess_exec",
AsyncMock(return_value=proc),
),
):
result = await maintainer.move_segment(
"test_cam",
"main",
start_time,
start_time + datetime.timedelta(seconds=0.96),
0.96,
cache_path,
SegmentInfo(0, 0, 0, 0),
)
self.assertIsNotNone(result)
paths.append(result[Recordings.path.name])
# the row keeps the resolved start even though the path doesn't
self.assertEqual(
result[Recordings.start_time.name], start_time.timestamp()
)
self.assertEqual(len(set(paths)), 2, paths)
self.assertTrue(paths[0].endswith("30.22.mp4"), paths[0])
self.assertTrue(paths[1].endswith("30.23.mp4"), paths[1])
class TestSegmentStartChaining(unittest.IsolatedAsyncioTestCase): class TestSegmentStartChaining(unittest.IsolatedAsyncioTestCase):
"""Contiguous segments must chain start times across filename truncation. """Contiguous segments must chain start times across filename truncation.
+4 -16
View File
@@ -20,12 +20,7 @@ from typing import TYPE_CHECKING, Any
import numpy as np import numpy as np
from ruamel.yaml import YAML from ruamel.yaml import YAML
from frigate.const import ( from frigate.const import REGEX_HTTP_CAMERA_USER_PASS, REGEX_RTSP_CAMERA_USER_PASS
REGEX_HTTP_CAMERA_USER_PASS,
REGEX_RTSP_CAMERA_USER_PASS,
STREAM_TYPE_MAIN,
STREAM_TYPE_SUB,
)
if TYPE_CHECKING: if TYPE_CHECKING:
from frigate.config import CameraConfig from frigate.config import CameraConfig
@@ -142,16 +137,9 @@ def get_ffmpeg_arg_list(arg: Any) -> list:
DEFAULT_RECORD_SEGMENT_TIME = 10 DEFAULT_RECORD_SEGMENT_TIME = 10
def get_record_segment_time( def get_record_segment_time(config: "CameraConfig") -> int:
config: "CameraConfig", stream_type: str = STREAM_TYPE_MAIN """Extract -segment_time from the camera's record output args."""
) -> int: record_args = get_ffmpeg_arg_list(config.ffmpeg.output_args.record)
"""Extract -segment_time from the camera's record output args for a stream."""
output_args = (
config.ffmpeg.output_args.effective_record_sub
if stream_type == STREAM_TYPE_SUB
else config.ffmpeg.output_args.record
)
record_args = get_ffmpeg_arg_list(output_args)
if record_args and record_args[0].startswith("preset"): if record_args and record_args[0].startswith("preset"):
return DEFAULT_RECORD_SEGMENT_TIME return DEFAULT_RECORD_SEGMENT_TIME
+15
View File
@@ -187,8 +187,23 @@ def get_physical_interfaces(interfaces) -> list:
return physical_interfaces return physical_interfaces
_bandwidth_warning_logged = False
def get_bandwidth_stats(config) -> dict[str, dict]: def get_bandwidth_stats(config) -> dict[str, dict]:
"""Get bandwidth usages for each ffmpeg process id""" """Get bandwidth usages for each ffmpeg process id"""
global _bandwidth_warning_logged
if os.geteuid() != 0:
if not _bandwidth_warning_logged:
logger.warning(
"Network bandwidth stats require root (nethogs needs CAP_NET_ADMIN/CAP_NET_RAW) "
"and are disabled; set FRIGATE_RUN_AS_ROOT=true or disable "
"telemetry.stats.network_bandwidth to silence this warning"
)
_bandwidth_warning_logged = True
return {}
usages = {} usages = {}
top_command = ["nethogs", "-t", "-v0", "-c5", "-d1"] + get_physical_interfaces( top_command = ["nethogs", "-t", "-v0", "-c5", "-d1"] + get_physical_interfaces(
config.telemetry.network_interfaces config.telemetry.network_interfaces
+93 -163
View File
@@ -5,7 +5,7 @@ import queue
import subprocess as sp import subprocess as sp
import threading import threading
import time import time
from collections import defaultdict, deque from collections import deque
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
@@ -22,14 +22,7 @@ from frigate.config.camera.updater import (
CameraConfigUpdateEnum, CameraConfigUpdateEnum,
CameraConfigUpdateSubscriber, CameraConfigUpdateSubscriber,
) )
from frigate.const import ( from frigate.const import PROCESS_PRIORITY_HIGH
PROCESS_PRIORITY_HIGH,
RECORD_STREAM_TYPES,
ROLE_TO_STREAM_TYPE,
STREAM_TYPE_MAIN,
STREAM_TYPE_SUB,
STREAM_TYPE_TO_ROLE,
)
from frigate.log import LogPipe from frigate.log import LogPipe
from frigate.util.builtin import EventsPerSecond, get_record_segment_time from frigate.util.builtin import EventsPerSecond, get_record_segment_time
from frigate.util.ffmpeg import start_or_restart_ffmpeg, stop_ffmpeg from frigate.util.ffmpeg import start_or_restart_ffmpeg, stop_ffmpeg
@@ -41,8 +34,6 @@ from frigate.util.process import FrigateProcess
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
RECORD_GRACE_SECONDS = 90
def capture_frames( def capture_frames(
ffmpeg_process: sp.Popen[Any], ffmpeg_process: sp.Popen[Any],
@@ -159,26 +150,16 @@ class CameraWatchdog(threading.Thread):
self.was_record_sub_enabled = self.config.record.sub.enabled self.was_record_sub_enabled = self.config.record.sub.enabled
self.segment_subscriber = RecordingsDataSubscriber(RecordingsDataTypeEnum.all) self.segment_subscriber = RecordingsDataSubscriber(RecordingsDataTypeEnum.all)
self.latest_valid_segment_time: dict[str, float] = defaultdict(float) self.latest_valid_segment_time: float = 0
self.latest_invalid_segment_time: dict[str, float] = defaultdict(float) self.latest_invalid_segment_time: float = 0
self.latest_cache_segment_time: dict[str, float] = defaultdict(float) self.latest_cache_segment_time: float = 0
self.record_enable_time: datetime | None = None self.record_enable_time: datetime | None = None
self.stream_grace_until: dict[str, datetime] = {}
# `valid` segments are published with the segment's start time, so the # `valid` segments are published with the segment's start time, so the
# gap between consecutive publishes can reach 2 * segment_time. Pad the # gap between consecutive publishes can reach 2 * segment_time. Pad the
# staleness threshold so it's never tighter than that worst case. # staleness threshold so it's never tighter than that worst case.
self.record_stale_threshold: dict[str, int] = { segment_time = get_record_segment_time(self.config)
stream_type: max( self.record_stale_threshold = max(120, 2 * segment_time + 30)
120, 2 * get_record_segment_time(self.config, stream_type) + 30
)
for stream_type in RECORD_STREAM_TYPES
}
# the sub stream usually shares its input, and therefore its ffmpeg
# process, with detect, so it isn't in ffmpeg_other_processes and needs
# its own staleness check
self.detect_process_records_sub = False
# Stall tracking (based on last processed frame) # Stall tracking (based on last processed frame)
self._stall_timestamps: deque[float] = deque() self._stall_timestamps: deque[float] = deque()
@@ -186,7 +167,7 @@ class CameraWatchdog(threading.Thread):
# Status caching to reduce message volume # Status caching to reduce message volume
self._last_detect_status: str | None = None self._last_detect_status: str | None = None
self._last_record_status: dict[str, str] = {} self._last_record_status: str | None = None
self._last_status_update_time: float = 0.0 self._last_status_update_time: float = 0.0
def _send_detect_status(self, status: str, now: float) -> None: def _send_detect_status(self, status: str, now: float) -> None:
@@ -199,78 +180,16 @@ class CameraWatchdog(threading.Thread):
self._last_detect_status = status self._last_detect_status = status
self._last_status_update_time = now self._last_status_update_time = now
def _send_record_status(self, stream_type: str, status: str, now: float) -> None: def _send_record_status(self, status: str, now: float) -> None:
"""Send a record stream's status only if changed or retry_interval has elapsed.""" """Send record status only if changed or retry_interval has elapsed."""
if ( if (
status != self._last_record_status.get(stream_type) status != self._last_record_status
or (now - self._last_status_update_time) >= self.sleeptime or (now - self._last_status_update_time) >= self.sleeptime
): ):
self.requestor.send_data( self.requestor.send_data(f"{self.config.name}/status/record", status)
f"{self.config.name}/status/{STREAM_TYPE_TO_ROLE[stream_type]}", status self._last_record_status = status
)
self._last_record_status[stream_type] = status
self._last_status_update_time = now self._last_status_update_time = now
def _reset_segment_times(self) -> None:
self.latest_valid_segment_time.clear()
self.latest_invalid_segment_time.clear()
self.latest_cache_segment_time.clear()
self.stream_grace_until.clear()
def _grant_restart_grace(self, stream_types: list[str], now_utc: datetime) -> None:
for stream_type in stream_types:
self.stream_grace_until[stream_type] = now_utc + timedelta(
seconds=RECORD_GRACE_SECONDS
)
def _stream_staleness(self, stream_type: str, now_utc: datetime) -> str | None:
"""Return why the stream's segments are stale, or None if they're healthy."""
# ffmpeg needs time to create a first segment after recording is
# enabled and after a restart, per stream
in_grace_period = (
self.record_enable_time is not None
and (now_utc - self.record_enable_time)
< timedelta(seconds=RECORD_GRACE_SECONDS)
) or now_utc < self.stream_grace_until.get(stream_type, now_utc)
if in_grace_period:
return None
latest_cache = self.latest_cache_segment_time[stream_type]
latest_valid = self.latest_valid_segment_time[stream_type]
latest_invalid = self.latest_invalid_segment_time[stream_type]
def as_dt(timestamp: float) -> datetime:
if timestamp > 0:
return datetime.fromtimestamp(timestamp, tz=UTC)
return now_utc - timedelta(seconds=1)
stale_window = timedelta(seconds=self.record_stale_threshold[stream_type])
if now_utc > (as_dt(latest_cache) + stale_window):
return "No new recording segments were created"
if now_utc > (as_dt(latest_valid) + stale_window):
return "No new valid recording segments were created"
if (
latest_invalid > 0
and now_utc > (as_dt(latest_invalid) + stale_window)
and latest_valid <= latest_invalid
):
return "No valid segments created since last invalid segment"
return None
def _recorded_streams(self, roles: list[Any]) -> list[str]:
"""Record stream types the given roles cover that are currently recording."""
return [
stream_type
for role, stream_type in ROLE_TO_STREAM_TYPE.items()
if role in roles and self.config.record.stream_enabled(stream_type)
]
def _check_config_updates(self) -> dict[str, list[str]]: def _check_config_updates(self) -> dict[str, list[str]]:
"""Check for config updates and return the update dict.""" """Check for config updates and return the update dict."""
return self.config_subscriber.check_for_updates() return self.config_subscriber.check_for_updates()
@@ -326,11 +245,6 @@ class CameraWatchdog(threading.Thread):
self.logger.info("Restarting ffmpeg...") self.logger.info("Restarting ffmpeg...")
self.start_ffmpeg_detect() self.start_ffmpeg_detect()
# this process produces the sub stream's segments too, so it gets the
# same startup grace however the reset was triggered
if self.detect_process_records_sub:
self._grant_restart_grace([STREAM_TYPE_SUB], datetime.now().astimezone(UTC))
def run(self) -> None: def run(self) -> None:
if self._update_enabled_state(): if self._update_enabled_state():
self.start_all_ffmpeg() self.start_all_ffmpeg()
@@ -353,7 +267,9 @@ class CameraWatchdog(threading.Thread):
) )
self.stop_all_ffmpeg() self.stop_all_ffmpeg()
self.start_all_ffmpeg() self.start_all_ffmpeg()
self._reset_segment_times() self.latest_valid_segment_time = 0
self.latest_invalid_segment_time = 0
self.latest_cache_segment_time = 0
self.record_enable_time = datetime.now().astimezone(UTC) self.record_enable_time = datetime.now().astimezone(UTC)
last_restart_time = datetime.now().timestamp() last_restart_time = datetime.now().timestamp()
continue continue
@@ -365,7 +281,9 @@ class CameraWatchdog(threading.Thread):
self.start_all_ffmpeg() self.start_all_ffmpeg()
# reset all timestamps and record the enable time for grace period # reset all timestamps and record the enable time for grace period
self._reset_segment_times() self.latest_valid_segment_time = 0
self.latest_invalid_segment_time = 0
self.latest_cache_segment_time = 0
self.record_enable_time = datetime.now().astimezone(UTC) self.record_enable_time = datetime.now().astimezone(UTC)
else: else:
self.logger.debug(f"Disabling camera {self.config.name}") self.logger.debug(f"Disabling camera {self.config.name}")
@@ -375,10 +293,7 @@ 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._send_detect_status("disabled", now)
self._send_record_status(STREAM_TYPE_MAIN, "disabled", now) self._send_record_status("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.was_enabled = enabled self.was_enabled = enabled
continue continue
@@ -390,7 +305,9 @@ class CameraWatchdog(threading.Thread):
) )
self.stop_all_ffmpeg() self.stop_all_ffmpeg()
self.start_all_ffmpeg() self.start_all_ffmpeg()
self._reset_segment_times() self.latest_valid_segment_time = 0
self.latest_invalid_segment_time = 0
self.latest_cache_segment_time = 0
self.record_enable_time = datetime.now().astimezone(UTC) self.record_enable_time = datetime.now().astimezone(UTC)
last_restart_time = datetime.now().timestamp() last_restart_time = datetime.now().timestamp()
self.was_record_enabled_in_config = record_enabled_in_config self.was_record_enabled_in_config = record_enabled_in_config
@@ -406,7 +323,9 @@ class CameraWatchdog(threading.Thread):
) )
self.stop_all_ffmpeg() self.stop_all_ffmpeg()
self.start_all_ffmpeg() self.start_all_ffmpeg()
self._reset_segment_times() self.latest_valid_segment_time = 0
self.latest_invalid_segment_time = 0
self.latest_cache_segment_time = 0
self.record_enable_time = datetime.now().astimezone(UTC) self.record_enable_time = datetime.now().astimezone(UTC)
last_restart_time = datetime.now().timestamp() last_restart_time = datetime.now().timestamp()
self.was_record_sub_enabled = record_sub_enabled self.was_record_sub_enabled = record_sub_enabled
@@ -424,25 +343,26 @@ class CameraWatchdog(threading.Thread):
raw_topic, payload = update raw_topic, payload = update
if raw_topic and payload: if raw_topic and payload:
topic = str(raw_topic) topic = str(raw_topic)
camera, stream_type, segment_time, _ = payload camera, segment_time, _ = payload
if camera != self.config.name: if camera != self.config.name:
continue continue
if topic.endswith(RecordingsDataTypeEnum.invalid.value): if topic.endswith(RecordingsDataTypeEnum.invalid.value):
self.logger.warning( self.logger.warning(
f"Invalid recording segment detected for {camera} ({stream_type}) at {segment_time}" f"Invalid recording segment detected for {camera} at {segment_time}"
) )
self.latest_invalid_segment_time[stream_type] = segment_time self.latest_invalid_segment_time = segment_time
elif topic.endswith(RecordingsDataTypeEnum.valid.value): elif topic.endswith(RecordingsDataTypeEnum.valid.value):
self.logger.debug( self.logger.debug(
f"Latest valid recording segment time on {camera} ({stream_type}): {segment_time}" f"Latest valid recording segment time on {camera}: {segment_time}"
) )
self.latest_valid_segment_time[stream_type] = segment_time self.latest_valid_segment_time = segment_time
elif topic.endswith(RecordingsDataTypeEnum.latest.value): elif topic.endswith(RecordingsDataTypeEnum.latest.value):
self.latest_cache_segment_time[stream_type] = ( if segment_time is not None:
segment_time if segment_time is not None else 0 self.latest_cache_segment_time = segment_time
) else:
self.latest_cache_segment_time = 0
now = datetime.now().timestamp() now = datetime.now().timestamp()
@@ -489,26 +409,63 @@ class CameraWatchdog(threading.Thread):
for p in self.ffmpeg_other_processes: for p in self.ffmpeg_other_processes:
poll = p["process"].poll() poll = p["process"].poll()
recorded_streams = self._recorded_streams(p["roles"]) if self.config.record.enabled and "record" in p["roles"]:
if recorded_streams:
now_utc = datetime.now().astimezone(UTC) now_utc = datetime.now().astimezone(UTC)
# ensure segments are still being created and that they have # Check if we're within the grace period after enabling recording
# valid video data. each stream is tracked separately so a # Grace period: 90 seconds allows time for ffmpeg to start and create first segment
# healthy one can't mask a stalled one. in_grace_period = self.record_enable_time is not None and (
stale_stream = None now_utc - self.record_enable_time
stale_reason = None ) < timedelta(seconds=90)
for stream_type in recorded_streams:
stale_reason = self._stream_staleness(stream_type, now_utc)
if stale_reason is not None: latest_cache_dt = (
stale_stream = stream_type datetime.fromtimestamp(self.latest_cache_segment_time, tz=UTC)
break if self.latest_cache_segment_time > 0
else now_utc - timedelta(seconds=1)
)
latest_valid_dt = (
datetime.fromtimestamp(self.latest_valid_segment_time, tz=UTC)
if self.latest_valid_segment_time > 0
else now_utc - timedelta(seconds=1)
)
latest_invalid_dt = (
datetime.fromtimestamp(self.latest_invalid_segment_time, tz=UTC)
if self.latest_invalid_segment_time > 0
else now_utc - timedelta(seconds=1)
)
# ensure segments are still being created and that they have valid video data
# Skip checks during grace period to allow segments to start being created
stale_window = timedelta(seconds=self.record_stale_threshold)
cache_stale = not in_grace_period and now_utc > (
latest_cache_dt + stale_window
)
valid_stale = not in_grace_period and now_utc > (
latest_valid_dt + stale_window
)
invalid_stale_condition = (
self.latest_invalid_segment_time > 0
and not in_grace_period
and now_utc > (latest_invalid_dt + stale_window)
and self.latest_valid_segment_time
<= self.latest_invalid_segment_time
)
invalid_stale = invalid_stale_condition
if cache_stale or valid_stale or invalid_stale:
if cache_stale:
reason = "No new recording segments were created"
elif valid_stale:
reason = "No new valid recording segments were created"
else: # invalid_stale
reason = (
"No valid segments created since last invalid segment"
)
if stale_stream is not None and can_restart:
self.logger.error( self.logger.error(
f"{stale_reason} for {self.config.name} ({stale_stream}) in the last {self.record_stale_threshold[stale_stream]}s. Restarting the ffmpeg record process..." f"{reason} for {self.config.name} in the last {self.record_stale_threshold}s. Restarting the ffmpeg record process..."
) )
p["process"] = start_or_restart_ffmpeg( p["process"] = start_or_restart_ffmpeg(
p["cmd"], p["cmd"],
@@ -522,18 +479,10 @@ class CameraWatchdog(threading.Thread):
f"{self.config.name}/status/{role.value}", "offline" f"{self.config.name}/status/{role.value}", "offline"
) )
self._grant_restart_grace(recorded_streams, now_utc)
last_restart_time = now
continue continue
elif stale_stream is None: else:
for stream_type in recorded_streams: self._send_record_status("online", now)
self._send_record_status(stream_type, "online", now) p["latest_segment_time"] = self.latest_cache_segment_time
p["latest_segment_time"] = max(
self.latest_cache_segment_time[stream_type]
for stream_type in recorded_streams
)
if poll is None: if poll is None:
continue continue
@@ -548,25 +497,6 @@ 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 (
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 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._send_record_status(STREAM_TYPE_SUB, "offline", now)
self.reset_capture_thread()
last_restart_time = now
# Prune expired reconnect timestamps # Prune expired reconnect timestamps
now = datetime.now().timestamp() now = datetime.now().timestamp()
while ( while (
@@ -609,9 +539,9 @@ class CameraWatchdog(threading.Thread):
self.segment_subscriber.stop() self.segment_subscriber.stop()
def start_ffmpeg_detect(self): def start_ffmpeg_detect(self):
detect_cmd = [c for c in self.config.ffmpeg_cmds if "detect" in c["roles"]][0] ffmpeg_cmd = [
ffmpeg_cmd = detect_cmd["cmd"] c["cmd"] for c in self.config.ffmpeg_cmds if "detect" in c["roles"]
self.detect_process_records_sub = "record_sub" in detect_cmd["roles"] ][0]
self.ffmpeg_detect_process = start_or_restart_ffmpeg( self.ffmpeg_detect_process = start_or_restart_ffmpeg(
ffmpeg_cmd, self.logger, self.logpipe, self.frame_size ffmpeg_cmd, self.logger, self.logpipe, self.frame_size
) )