mirror of
https://github.com/blakeblackshear/frigate.git
synced 2026-09-29 19:36:57 +03:00
Update base image to trixie
This commit is contained in:
@@ -14,9 +14,8 @@ concurrency:
|
||||
group: ${{ github.ref }}
|
||||
cancel-in-progress: true
|
||||
|
||||
env:
|
||||
PYTHON_VERSION: 3.11
|
||||
|
||||
# Note: jetson (jp6), rockchip, and synaptics builds are disabled until their
|
||||
# vendor runtimes support Python 3.13 (or a community member updates them).
|
||||
jobs:
|
||||
amd64_build:
|
||||
runs-on: ubuntu-22.04
|
||||
@@ -77,35 +76,6 @@ jobs:
|
||||
rpi.tags=${{ steps.setup.outputs.image-name }}-rpi
|
||||
*.cache-from=type=registry,ref=${{ steps.setup.outputs.cache-name }}-arm64
|
||||
*.cache-to=type=registry,ref=${{ steps.setup.outputs.cache-name }}-arm64,mode=max
|
||||
jetson_jp6_build:
|
||||
runs-on: ubuntu-22.04-arm
|
||||
name: Jetson Jetpack 6
|
||||
steps:
|
||||
- name: Check out code
|
||||
uses: actions/checkout@v6
|
||||
with:
|
||||
persist-credentials: false
|
||||
- name: Set up QEMU and Buildx
|
||||
id: setup
|
||||
uses: ./.github/actions/setup
|
||||
with:
|
||||
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
|
||||
- name: Build and push TensorRT (Jetson, Jetpack 6)
|
||||
env:
|
||||
ARCH: arm64
|
||||
BASE_IMAGE: nvcr.io/nvidia/tensorrt:23.12-py3-igpu
|
||||
SLIM_BASE: nvcr.io/nvidia/tensorrt:23.12-py3-igpu
|
||||
TRT_BASE: nvcr.io/nvidia/tensorrt:23.12-py3-igpu
|
||||
uses: docker/bake-action@v7
|
||||
with:
|
||||
source: .
|
||||
push: true
|
||||
targets: tensorrt
|
||||
files: docker/tensorrt/trt.hcl
|
||||
set: |
|
||||
tensorrt.tags=${{ steps.setup.outputs.image-name }}-tensorrt-jp6
|
||||
*.cache-from=type=registry,ref=${{ steps.setup.outputs.cache-name }}-jp6
|
||||
*.cache-to=type=registry,ref=${{ steps.setup.outputs.cache-name }}-jp6,mode=max
|
||||
amd64_extra_builds:
|
||||
runs-on: ubuntu-22.04
|
||||
name: AMD64 Extra Build
|
||||
@@ -147,56 +117,6 @@ jobs:
|
||||
rocm.tags=${{ steps.setup.outputs.image-name }}-rocm
|
||||
*.cache-to=type=registry,ref=${{ steps.setup.outputs.cache-name }}-rocm,mode=max
|
||||
*.cache-from=type=registry,ref=${{ steps.setup.outputs.cache-name }}-rocm
|
||||
arm64_extra_builds:
|
||||
runs-on: ubuntu-22.04-arm
|
||||
name: ARM Extra Build
|
||||
needs:
|
||||
- arm64_build
|
||||
steps:
|
||||
- name: Check out code
|
||||
uses: actions/checkout@v6
|
||||
with:
|
||||
persist-credentials: false
|
||||
- name: Set up QEMU and Buildx
|
||||
id: setup
|
||||
uses: ./.github/actions/setup
|
||||
with:
|
||||
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
|
||||
- name: Build and push Rockchip build
|
||||
uses: docker/bake-action@v7
|
||||
with:
|
||||
source: .
|
||||
push: true
|
||||
targets: rk
|
||||
files: docker/rockchip/rk.hcl
|
||||
set: |
|
||||
rk.tags=${{ steps.setup.outputs.image-name }}-rk
|
||||
*.cache-from=type=gha
|
||||
synaptics_build:
|
||||
runs-on: ubuntu-22.04-arm
|
||||
name: Synaptics Build
|
||||
needs:
|
||||
- arm64_build
|
||||
steps:
|
||||
- name: Check out code
|
||||
uses: actions/checkout@v6
|
||||
with:
|
||||
persist-credentials: false
|
||||
- name: Set up QEMU and Buildx
|
||||
id: setup
|
||||
uses: ./.github/actions/setup
|
||||
with:
|
||||
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
|
||||
- name: Build and push Synaptics build
|
||||
uses: docker/bake-action@v7
|
||||
with:
|
||||
source: .
|
||||
push: true
|
||||
targets: synaptics
|
||||
files: docker/synaptics/synaptics.hcl
|
||||
set: |
|
||||
synaptics.tags=${{ steps.setup.outputs.image-name }}-synaptics
|
||||
*.cache-from=type=gha
|
||||
# The majority of users running arm64 are rpi users, so the rpi
|
||||
# build should be the primary arm64 image
|
||||
assemble_default_build:
|
||||
|
||||
@@ -9,7 +9,7 @@ on:
|
||||
- ".github/ISSUE_TEMPLATE/**"
|
||||
|
||||
env:
|
||||
DEFAULT_PYTHON: 3.11
|
||||
DEFAULT_PYTHON: 3.13
|
||||
|
||||
jobs:
|
||||
web_lint:
|
||||
|
||||
@@ -39,14 +39,14 @@ jobs:
|
||||
STABLE_TAG=${BASE}:stable
|
||||
PULL_TAG=${BASE}:${BUILD_TAG}
|
||||
docker run --rm -v $HOME/.docker/config.json:/config.json quay.io/skopeo/stable:latest copy --authfile /config.json --multi-arch all docker://${PULL_TAG} docker://${VERSION_TAG}
|
||||
for variant in standard-arm64 tensorrt tensorrt-jp6 rk rocm synaptics; do
|
||||
for variant in standard-arm64 tensorrt rocm; do
|
||||
docker run --rm -v $HOME/.docker/config.json:/config.json quay.io/skopeo/stable:latest copy --authfile /config.json --multi-arch all docker://${PULL_TAG}-${variant} docker://${VERSION_TAG}-${variant}
|
||||
done
|
||||
|
||||
# stable tag
|
||||
if [[ "${BUILD_TYPE}" == "stable" ]]; then
|
||||
docker run --rm -v $HOME/.docker/config.json:/config.json quay.io/skopeo/stable:latest copy --authfile /config.json --multi-arch all docker://${PULL_TAG} docker://${STABLE_TAG}
|
||||
for variant in standard-arm64 tensorrt tensorrt-jp6 rk rocm synaptics; do
|
||||
for variant in standard-arm64 tensorrt rocm; do
|
||||
docker run --rm -v $HOME/.docker/config.json:/config.json quay.io/skopeo/stable:latest copy --authfile /config.json --multi-arch all docker://${PULL_TAG}-${variant} docker://${STABLE_TAG}-${variant}
|
||||
done
|
||||
fi
|
||||
|
||||
+14
-17
@@ -6,8 +6,8 @@ ARG DEBIAN_FRONTEND=noninteractive
|
||||
# Globally set pip break-system-packages option to avoid having to specify it every time
|
||||
ARG PIP_BREAK_SYSTEM_PACKAGES=1
|
||||
|
||||
ARG BASE_IMAGE=debian:12
|
||||
ARG SLIM_BASE=debian:12-slim
|
||||
ARG BASE_IMAGE=debian:13
|
||||
ARG SLIM_BASE=debian:13-slim
|
||||
|
||||
# A hook that allows us to inject commands right after the base images
|
||||
ARG BASE_HOOK=
|
||||
@@ -18,7 +18,7 @@ ARG BASE_HOOK
|
||||
|
||||
RUN sh -c "$BASE_HOOK"
|
||||
|
||||
FROM --platform=${BUILDPLATFORM} debian:12 AS base_host
|
||||
FROM --platform=${BUILDPLATFORM} debian:13 AS base_host
|
||||
ARG PIP_BREAK_SYSTEM_PACKAGES
|
||||
|
||||
FROM ${SLIM_BASE} AS slim-base
|
||||
@@ -52,14 +52,6 @@ RUN --mount=type=tmpfs,target=/tmp --mount=type=tmpfs,target=/var/cache/apt \
|
||||
--mount=type=cache,target=/root/.ccache \
|
||||
/deps/build_sqlite_vec.sh
|
||||
|
||||
# Build intel-media-driver from source against bookworm's system libva so it
|
||||
# works with Debian 12's glibc/libstdc++ (pre-built noble/trixie packages
|
||||
# require glibc 2.38 which is not available on bookworm).
|
||||
FROM base AS intel-media-driver
|
||||
ARG DEBIAN_FRONTEND
|
||||
RUN --mount=type=bind,source=docker/main/build_intel_media_driver.sh,target=/deps/build_intel_media_driver.sh \
|
||||
/deps/build_intel_media_driver.sh
|
||||
|
||||
FROM scratch AS go2rtc
|
||||
ARG TARGETARCH
|
||||
WORKDIR /rootfs/usr/local/go2rtc/bin
|
||||
@@ -80,11 +72,12 @@ RUN --mount=type=bind,source=docker/main/install_tempio.sh,target=/deps/install_
|
||||
# Download and Convert OpenVino model
|
||||
FROM base_host AS ov-converter
|
||||
ARG DEBIAN_FRONTEND
|
||||
ARG PIP_BREAK_SYSTEM_PACKAGES
|
||||
|
||||
# Install OpenVINO for model conversion
|
||||
COPY docker/main/requirements-ov.txt /requirements-ov.txt
|
||||
RUN apt-get -qq update \
|
||||
&& apt-get -qq install -y wget python3 python3-distutils \
|
||||
&& apt-get -qq install -y wget python3 \
|
||||
&& wget -q https://bootstrap.pypa.io/get-pip.py -O get-pip.py \
|
||||
&& sed -i 's/args.append("setuptools")/args.append("setuptools==77.0.3")/' get-pip.py \
|
||||
&& python3 get-pip.py "pip" \
|
||||
@@ -157,6 +150,7 @@ FROM base AS wheels
|
||||
ARG DEBIAN_FRONTEND
|
||||
ARG TARGETARCH
|
||||
ARG DEBUG=false
|
||||
ARG PIP_BREAK_SYSTEM_PACKAGES
|
||||
|
||||
# Use a separate container to build wheels to prevent build dependencies in final image
|
||||
RUN apt-get -qq update \
|
||||
@@ -164,8 +158,8 @@ RUN apt-get -qq update \
|
||||
apt-transport-https wget unzip \
|
||||
&& apt-get -qq update \
|
||||
&& apt-get -qq install -y \
|
||||
python3.11 \
|
||||
python3.11-dev \
|
||||
python3 \
|
||||
python3-dev \
|
||||
# opencv dependencies
|
||||
build-essential cmake git pkg-config libgtk-3-dev \
|
||||
libavcodec-dev libavformat-dev libswscale-dev libv4l-dev \
|
||||
@@ -179,8 +173,6 @@ RUN apt-get -qq update \
|
||||
gcc gfortran libopenblas-dev liblapack-dev && \
|
||||
rm -rf /var/lib/apt/lists/*
|
||||
|
||||
RUN update-alternatives --install /usr/bin/python3 python3 /usr/bin/python3.11 1
|
||||
|
||||
RUN wget -q https://bootstrap.pypa.io/get-pip.py -O get-pip.py \
|
||||
&& sed -i 's/args.append("setuptools")/args.append("setuptools==77.0.3")/' get-pip.py \
|
||||
&& python3 get-pip.py "pip"
|
||||
@@ -200,6 +192,10 @@ RUN pip3 wheel --wheel-dir=/wheels -r /requirements-wheels.txt && \
|
||||
pip3 wheel --wheel-dir=/wheels -r /requirements-dev.txt; \
|
||||
fi
|
||||
|
||||
# Build norfair separately with its numpy pin lifted
|
||||
RUN --mount=type=bind,source=docker/main/build_norfair.sh,target=/deps/build_norfair.sh \
|
||||
/deps/build_norfair.sh
|
||||
|
||||
# Install HailoRT & Wheels
|
||||
RUN --mount=type=bind,source=docker/main/install_hailort.sh,target=/deps/install_hailort.sh \
|
||||
/deps/install_hailort.sh
|
||||
@@ -208,7 +204,6 @@ RUN --mount=type=bind,source=docker/main/install_hailort.sh,target=/deps/install
|
||||
FROM scratch AS deps-rootfs
|
||||
COPY --from=nginx /usr/local/nginx/ /usr/local/nginx/
|
||||
COPY --from=sqlite-vec /usr/local/lib/ /usr/local/lib/
|
||||
COPY --from=intel-media-driver /rootfs/ /
|
||||
COPY --from=go2rtc /rootfs/ /
|
||||
COPY --from=libusb-build /usr/local/lib /usr/local/lib
|
||||
COPY --from=tempio /rootfs/ /
|
||||
@@ -222,6 +217,7 @@ COPY docker/main/rootfs/ /
|
||||
FROM slim-base AS deps
|
||||
ARG TARGETARCH
|
||||
ARG BASE_IMAGE
|
||||
ARG PIP_BREAK_SYSTEM_PACKAGES
|
||||
|
||||
ARG DEBIAN_FRONTEND
|
||||
# http://stackoverflow.com/questions/48162574/ddg#49462622
|
||||
@@ -306,6 +302,7 @@ HEALTHCHECK --start-period=300s --start-interval=5s --interval=15s --timeout=5s
|
||||
|
||||
# Frigate deps with Node.js and NPM for devcontainer
|
||||
FROM deps AS devcontainer
|
||||
ARG PIP_BREAK_SYSTEM_PACKAGES
|
||||
|
||||
# Do not start the actual Frigate service on devcontainer as it will be started by VS Code
|
||||
# But start a fake service for simulating the logs
|
||||
|
||||
@@ -1,48 +0,0 @@
|
||||
#!/bin/bash
|
||||
|
||||
set -euxo pipefail
|
||||
|
||||
# Intel media driver is x86_64-only. Create empty rootfs on other arches so
|
||||
# the downstream COPY --from has a valid source.
|
||||
if [ "$(uname -m)" != "x86_64" ]; then
|
||||
mkdir -p /rootfs
|
||||
exit 0
|
||||
fi
|
||||
|
||||
MEDIA_DRIVER_VERSION="intel-media-25.2.6"
|
||||
GMMLIB_VERSION="intel-gmmlib-22.7.2"
|
||||
|
||||
apt-get -qq update
|
||||
apt-get -qq install -y wget gnupg ca-certificates cmake g++ make pkg-config
|
||||
|
||||
# Use Intel's jammy repo for newer libva-dev (2.22) which provides the
|
||||
# VVC/VVC-decode headers required by media-driver 25.x
|
||||
wget -qO - https://repositories.intel.com/gpu/intel-graphics.key | gpg --yes --dearmor --output /usr/share/keyrings/intel-graphics.gpg
|
||||
echo "deb [arch=amd64 signed-by=/usr/share/keyrings/intel-graphics.gpg] https://repositories.intel.com/gpu/ubuntu jammy client" > /etc/apt/sources.list.d/intel-gpu-jammy.list
|
||||
apt-get -qq update
|
||||
apt-get -qq install -y libva-dev
|
||||
|
||||
# Build gmmlib (required by media-driver)
|
||||
wget -qO gmmlib.tar.gz "https://github.com/intel/gmmlib/archive/refs/tags/${GMMLIB_VERSION}.tar.gz"
|
||||
mkdir /tmp/gmmlib
|
||||
tar -xf gmmlib.tar.gz -C /tmp/gmmlib --strip-components 1
|
||||
cmake -S /tmp/gmmlib -B /tmp/gmmlib/build -DCMAKE_BUILD_TYPE=Release
|
||||
make -C /tmp/gmmlib/build -j"$(nproc)"
|
||||
make -C /tmp/gmmlib/build install
|
||||
|
||||
# Build intel-media-driver
|
||||
wget -qO media-driver.tar.gz "https://github.com/intel/media-driver/archive/refs/tags/${MEDIA_DRIVER_VERSION}.tar.gz"
|
||||
mkdir /tmp/media-driver
|
||||
tar -xf media-driver.tar.gz -C /tmp/media-driver --strip-components 1
|
||||
cmake -S /tmp/media-driver -B /tmp/media-driver/build \
|
||||
-DCMAKE_BUILD_TYPE=Release \
|
||||
-DENABLE_KERNELS=ON \
|
||||
-DENABLE_NONFREE_KERNELS=ON \
|
||||
-DCMAKE_INSTALL_PREFIX=/usr \
|
||||
-DCMAKE_INSTALL_LIBDIR=/usr/lib/x86_64-linux-gnu \
|
||||
-DCMAKE_C_FLAGS="-Wno-error" \
|
||||
-DCMAKE_CXX_FLAGS="-Wno-error"
|
||||
make -C /tmp/media-driver/build -j"$(nproc)"
|
||||
|
||||
# Install driver to rootfs for COPY --from
|
||||
make -C /tmp/media-driver/build install DESTDIR=/rootfs
|
||||
@@ -8,9 +8,8 @@ SECURE_TOKEN_MODULE_VERSION="1.5"
|
||||
SET_MISC_MODULE_VERSION="v0.33"
|
||||
NGX_DEVEL_KIT_VERSION="v0.3.3"
|
||||
|
||||
source /etc/os-release
|
||||
|
||||
if [[ "$VERSION_ID" == "12" ]]; then
|
||||
# enable deb-src entries in either deb822 or legacy sources format
|
||||
if [[ -f /etc/apt/sources.list.d/debian.sources ]]; then
|
||||
sed -i '/^Types:/s/deb/& deb-src/' /etc/apt/sources.list.d/debian.sources
|
||||
else
|
||||
cp /etc/apt/sources.list /etc/apt/sources.list.d/sources-src.list
|
||||
|
||||
Executable
+36
@@ -0,0 +1,36 @@
|
||||
#!/bin/bash
|
||||
|
||||
# norfair 2.3.0 declares numpy < 2 in its wheel metadata, but the code is
|
||||
# numpy 2 compatible. Build it without deps and lift the pin until upstream
|
||||
# publishes a numpy 2 compatible release.
|
||||
|
||||
set -euxo pipefail
|
||||
|
||||
norfair_version="2.3.0"
|
||||
|
||||
pip3 wheel --wheel-dir=/norfair-wheel --no-deps "norfair==${norfair_version}"
|
||||
|
||||
python3 - <<'EOF'
|
||||
import glob
|
||||
import os
|
||||
import re
|
||||
import zipfile
|
||||
|
||||
path = glob.glob("/norfair-wheel/norfair-*.whl")[0]
|
||||
patched = path + ".patched"
|
||||
|
||||
with zipfile.ZipFile(path) as zin, zipfile.ZipFile(patched, "w", zipfile.ZIP_DEFLATED) as zout:
|
||||
for item in zin.infolist():
|
||||
data = zin.read(item.filename)
|
||||
if item.filename.endswith(".dist-info/METADATA"):
|
||||
data = re.sub(
|
||||
rb"Requires-Dist: numpy.*",
|
||||
b"Requires-Dist: numpy (>=1.23.0)",
|
||||
data,
|
||||
)
|
||||
zout.writestr(item, data)
|
||||
|
||||
os.replace(patched, path)
|
||||
EOF
|
||||
|
||||
mv /norfair-wheel/*.whl /wheels/
|
||||
@@ -4,9 +4,8 @@ set -euxo pipefail
|
||||
|
||||
SQLITE_VEC_VERSION="0.1.9"
|
||||
|
||||
source /etc/os-release
|
||||
|
||||
if [[ "$VERSION_ID" == "12" ]]; then
|
||||
# enable deb-src entries in either deb822 or legacy sources format
|
||||
if [[ -f /etc/apt/sources.list.d/debian.sources ]]; then
|
||||
sed -i '/^Types:/s/deb/& deb-src/' /etc/apt/sources.list.d/debian.sources
|
||||
else
|
||||
cp /etc/apt/sources.list /etc/apt/sources.list.d/sources-src.list
|
||||
|
||||
+10
-23
@@ -12,36 +12,32 @@ apt-get -qq install --no-install-recommends -y \
|
||||
lbzip2 \
|
||||
procps vainfo \
|
||||
unzip locales tzdata libxml2 xz-utils \
|
||||
python3.11 \
|
||||
python3 \
|
||||
curl \
|
||||
lsof \
|
||||
jq \
|
||||
nethogs \
|
||||
libgl1 \
|
||||
libglib2.0-0 \
|
||||
libusb-1.0.0 \
|
||||
libusb-1.0-0 \
|
||||
python3-h2 \
|
||||
libgomp1 # memryx detector
|
||||
|
||||
update-alternatives --install /usr/bin/python3 python3 /usr/bin/python3.11 1
|
||||
|
||||
mkdir -p -m 600 /root/.gnupg
|
||||
|
||||
# install coral runtime
|
||||
wget -q -O /tmp/libedgetpu1-max.deb "https://github.com/feranick/libedgetpu/releases/download/16.0TF2.17.1-1/libedgetpu1-max_16.0tf2.17.1-1.bookworm_${TARGETARCH}.deb"
|
||||
wget -q -O /tmp/libedgetpu1-max.deb "https://github.com/feranick/libedgetpu/releases/download/16.0TF2.17.1-1/libedgetpu1-max_16.0tf2.17.1-1.trixie_${TARGETARCH}.deb"
|
||||
unset DEBIAN_FRONTEND
|
||||
yes | dpkg -i /tmp/libedgetpu1-max.deb && export DEBIAN_FRONTEND=noninteractive
|
||||
rm /tmp/libedgetpu1-max.deb
|
||||
|
||||
# install mesa-teflon-delegate from bookworm-backports
|
||||
# install mesa-teflon-delegate
|
||||
# Only available for arm64 at the moment
|
||||
if [[ "${TARGETARCH}" == "arm64" ]]; then
|
||||
if [[ "${BASE_IMAGE}" == *"nvcr.io/nvidia/tensorrt"* ]]; then
|
||||
echo "Info: Skipping apt-get commands because BASE_IMAGE includes 'nvcr.io/nvidia/tensorrt' for arm64."
|
||||
else
|
||||
echo "deb http://deb.debian.org/debian bookworm-backports main" | tee /etc/apt/sources.list.d/bookworm-backbacks.list
|
||||
apt-get -qq update
|
||||
apt-get -qq install --no-install-recommends --no-install-suggests -y mesa-teflon-delegate/bookworm-backports
|
||||
apt-get -qq install --no-install-recommends --no-install-suggests -y mesa-teflon-delegate
|
||||
fi
|
||||
fi
|
||||
|
||||
@@ -79,10 +75,11 @@ fi
|
||||
|
||||
# arch specific packages
|
||||
if [[ "${TARGETARCH}" == "amd64" ]]; then
|
||||
# Install non-free version of i965 driver
|
||||
# Install non-free intel media and i965 drivers
|
||||
sed -i -E "/^Components: main$/s/main/main contrib non-free non-free-firmware/" "/etc/apt/sources.list.d/debian.sources" \
|
||||
&& apt-get -qq update \
|
||||
&& apt-get install --no-install-recommends --no-install-suggests -y i965-va-driver-shaders \
|
||||
&& apt-get install --no-install-recommends --no-install-suggests -y \
|
||||
i965-va-driver-shaders intel-media-va-driver-non-free \
|
||||
&& sed -i -E "/^Components: main contrib non-free non-free-firmware$/s/main contrib non-free non-free-firmware/main/" "/etc/apt/sources.list.d/debian.sources" \
|
||||
&& apt-get update
|
||||
|
||||
@@ -92,28 +89,18 @@ if [[ "${TARGETARCH}" == "amd64" ]]; then
|
||||
libva-drm2 \
|
||||
mesa-va-drivers radeontop
|
||||
|
||||
# intel packages use zst compression so we need to update dpkg
|
||||
apt-get install -y dpkg
|
||||
|
||||
# use intel apt repo for libmfx1 (legacy QSV, pre-Gen12)
|
||||
wget -qO - https://repositories.intel.com/gpu/intel-graphics.key | gpg --yes --dearmor --output /usr/share/keyrings/intel-graphics.gpg
|
||||
echo "deb [arch=amd64 signed-by=/usr/share/keyrings/intel-graphics.gpg] https://repositories.intel.com/gpu/ubuntu jammy client" | tee /etc/apt/sources.list.d/intel-gpu-jammy.list
|
||||
apt-get -qq update
|
||||
|
||||
# intel-media-va-driver-non-free is built from source in the
|
||||
# intel-media-driver Dockerfile stage for Battlemage (Xe2) support
|
||||
apt-get -qq install --no-install-recommends --no-install-suggests -y \
|
||||
libmfx1
|
||||
rm -f /usr/share/keyrings/intel-graphics.gpg
|
||||
rm -f /etc/apt/sources.list.d/intel-gpu-jammy.list
|
||||
|
||||
# upgrade libva2, oneVPL runtime, and libvpl2 from trixie for Battlemage support
|
||||
echo "deb http://deb.debian.org/debian trixie main" > /etc/apt/sources.list.d/trixie.list
|
||||
apt-get -qq update
|
||||
apt-get -qq install -y -t trixie libva2 libva-drm2 libzstd1
|
||||
apt-get -qq install -y -t trixie libmfx-gen1.2 libvpl2
|
||||
rm -f /etc/apt/sources.list.d/trixie.list
|
||||
apt-get -qq update
|
||||
# oneVPL runtime for QSV support
|
||||
apt-get -qq install -y libva2 libva-drm2 libmfx-gen1.2 libvpl2
|
||||
apt-get -qq install -y ocl-icd-libopencl1
|
||||
|
||||
# install libtbb12 for NPU support
|
||||
|
||||
@@ -10,5 +10,5 @@ elif [[ "${TARGETARCH}" == "arm64" ]]; then
|
||||
arch="aarch64"
|
||||
fi
|
||||
|
||||
wget -qO- "https://github.com/frigate-nvr/hailort/releases/download/v${hailo_version}/hailort-debian12-${TARGETARCH}.tar.gz" | tar -C / -xzf -
|
||||
wget -P /wheels/ "https://github.com/frigate-nvr/hailort/releases/download/v${hailo_version}/hailort-${hailo_version}-cp311-cp311-linux_${arch}.whl"
|
||||
wget -qO- "https://github.com/frigate-nvr/hailort/releases/download/v${hailo_version}/hailort-debian13-${TARGETARCH}.tar.gz" | tar -C / -xzf -
|
||||
wget -P /wheels/ "https://github.com/frigate-nvr/hailort/releases/download/v${hailo_version}/hailort-${hailo_version}-cp313-cp313-linux_${arch}.whl"
|
||||
|
||||
@@ -13,10 +13,10 @@ pathvalidate == 3.3.*
|
||||
markupsafe == 3.0.*
|
||||
python-multipart == 0.0.26
|
||||
# Classification Model Training
|
||||
tensorflow == 2.19.* ; platform_machine == 'aarch64'
|
||||
tensorflow-cpu == 2.19.* ; platform_machine == 'x86_64'
|
||||
tensorflow == 2.20.* ; platform_machine == 'aarch64'
|
||||
tensorflow-cpu == 2.20.* ; platform_machine == 'x86_64'
|
||||
# General
|
||||
mypy == 1.6.1
|
||||
mypy == 1.20.*
|
||||
onvif-zeep-async == 4.0.*
|
||||
paho-mqtt == 2.1.*
|
||||
pandas == 2.2.*
|
||||
@@ -31,13 +31,15 @@ ruamel.yaml == 0.18.*
|
||||
tzlocal == 5.2
|
||||
requests == 2.32.*
|
||||
types-requests == 2.32.*
|
||||
norfair == 2.3.*
|
||||
# norfair is built separately in the Dockerfile with its numpy pin lifted;
|
||||
# filterpy is its only dependency not otherwise present
|
||||
filterpy == 1.4.*
|
||||
setproctitle == 1.3.*
|
||||
ws4py == 0.5.*
|
||||
unidecode == 1.3.*
|
||||
titlecase == 2.4.*
|
||||
# Image Manipulation
|
||||
numpy == 1.26.*
|
||||
numpy == 2.2.*
|
||||
opencv-python-headless == 4.11.0.*
|
||||
opencv-contrib-python == 4.11.0.*
|
||||
scipy == 1.16.*
|
||||
@@ -63,23 +65,20 @@ argcomplete==2.0.*
|
||||
contextlib2==0.6.*
|
||||
distlib==0.3.*
|
||||
filelock==3.8.*
|
||||
future==0.18.*
|
||||
importlib-metadata==5.1.*
|
||||
importlib-resources==5.1.*
|
||||
netaddr==0.8.*
|
||||
netifaces==0.10.*
|
||||
verboselogs==1.7.*
|
||||
virtualenv==20.17.*
|
||||
prometheus-client == 0.21.*
|
||||
# TFLite
|
||||
tflite_runtime @ https://github.com/frigate-nvr/TFlite-builds/releases/download/v2.17.1/tflite_runtime-2.17.1-cp311-cp311-linux_x86_64.whl; platform_machine == 'x86_64'
|
||||
tflite_runtime @ https://github.com/feranick/TFlite-builds/releases/download/v2.17.1/tflite_runtime-2.17.1-cp311-cp311-linux_aarch64.whl; platform_machine == 'aarch64'
|
||||
ai-edge-litert == 2.1.*
|
||||
# audio transcription
|
||||
sherpa-onnx==1.12.*
|
||||
faster-whisper==1.1.*
|
||||
librosa==0.11.*
|
||||
soundfile==0.13.*
|
||||
# DeGirum detector
|
||||
degirum == 0.16.*
|
||||
degirum == 0.20.*
|
||||
# Memory profiling
|
||||
memray == 1.15.*
|
||||
|
||||
@@ -31,15 +31,9 @@ RUN echo /opt/rocm/lib|tee /opt/rocm-dist/etc/ld.so.conf.d/rocm.conf
|
||||
#######################################################################
|
||||
FROM deps AS deps-prelim
|
||||
|
||||
COPY docker/rocm/debian-backports.sources /etc/apt/sources.list.d/debian-backports.sources
|
||||
# install_deps.sh upgraded libstdc++6 from trixie for Battlemage; the matching
|
||||
# -dev package must also come from trixie or apt refuses to satisfy it.
|
||||
RUN echo "deb http://deb.debian.org/debian trixie main" > /etc/apt/sources.list.d/trixie.list && \
|
||||
apt-get update && \
|
||||
RUN apt-get update && \
|
||||
apt-get install -y libnuma1 && \
|
||||
apt-get install -qq -y -t bookworm-backports mesa-va-drivers mesa-vulkan-drivers && \
|
||||
apt-get install -qq -y -t trixie libstdc++-14-dev && \
|
||||
rm -f /etc/apt/sources.list.d/trixie.list && \
|
||||
apt-get install -qq -y mesa-va-drivers mesa-vulkan-drivers libstdc++-14-dev && \
|
||||
rm -rf /var/lib/apt/lists/*
|
||||
|
||||
WORKDIR /opt/frigate
|
||||
|
||||
@@ -1,6 +0,0 @@
|
||||
Types: deb
|
||||
URIs: http://deb.debian.org/debian
|
||||
Suites: bookworm-backports
|
||||
Components: main
|
||||
Enabled: yes
|
||||
Signed-By: /usr/share/keyrings/debian-archive-keyring.gpg
|
||||
@@ -1 +1 @@
|
||||
onnxruntime-migraphx @ https://github.com/NickM-27/frigate-onnxruntime-rocm/releases/download/v7.2.3-1/onnxruntime_migraphx-1.24.4-cp311-cp311-linux_x86_64.whl
|
||||
onnxruntime-migraphx @ https://github.com/NickM-27/frigate-onnxruntime-rocm/releases/download/v7.2.3-1/onnxruntime_migraphx-1.24.4-cp313-cp313-linux_x86_64.whl
|
||||
@@ -18,14 +18,14 @@ apt-get -qq install --no-install-recommends -y \
|
||||
mkdir -p -m 600 /root/.gnupg
|
||||
|
||||
# enable non-free repo
|
||||
echo "deb http://deb.debian.org/debian bookworm main contrib non-free non-free-firmware" | tee -a /etc/apt/sources.list
|
||||
echo "deb http://deb.debian.org/debian trixie main contrib non-free non-free-firmware" | tee -a /etc/apt/sources.list
|
||||
apt update
|
||||
|
||||
# ffmpeg -> arm64
|
||||
if [[ "${TARGETARCH}" == "arm64" ]]; then
|
||||
# add raspberry pi repo
|
||||
gpg --no-default-keyring --keyring /usr/share/keyrings/raspbian.gpg --keyserver keyserver.ubuntu.com --recv-keys 82B129927FA3303E
|
||||
echo "deb [signed-by=/usr/share/keyrings/raspbian.gpg] https://archive.raspberrypi.org/debian/ bookworm main" | tee /etc/apt/sources.list.d/raspi.list
|
||||
echo "deb [signed-by=/usr/share/keyrings/raspbian.gpg] https://archive.raspberrypi.org/debian/ trixie main" | tee /etc/apt/sources.list.d/raspi.list
|
||||
apt-get -qq update
|
||||
apt-get -qq install --no-install-recommends --no-install-suggests -y ffmpeg
|
||||
mkdir -p /usr/lib/ffmpeg/rpi/bin
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
/usr/local/lib/python3.11/dist-packages/nvidia/cudnn/lib
|
||||
/usr/local/lib/python3.11/dist-packages/nvidia/cuda_runtime/lib
|
||||
/usr/local/lib/python3.11/dist-packages/nvidia/cublas/lib
|
||||
/usr/local/lib/python3.11/dist-packages/nvidia/cufft/lib
|
||||
/usr/local/lib/python3.11/dist-packages/nvidia/curand/lib/
|
||||
/usr/local/lib/python3.11/dist-packages/nvidia/cuda_nvrtc/lib/
|
||||
/usr/local/lib/python3.13/dist-packages/nvidia/cudnn/lib
|
||||
/usr/local/lib/python3.13/dist-packages/nvidia/cuda_runtime/lib
|
||||
/usr/local/lib/python3.13/dist-packages/nvidia/cublas/lib
|
||||
/usr/local/lib/python3.13/dist-packages/nvidia/cufft/lib
|
||||
/usr/local/lib/python3.13/dist-packages/nvidia/curand/lib/
|
||||
/usr/local/lib/python3.13/dist-packages/nvidia/cuda_nvrtc/lib/
|
||||
@@ -13,6 +13,5 @@ nvidia-cusolver-cu12==11.7.3.90; platform_machine == 'x86_64'
|
||||
nvidia-cusparse-cu12==12.5.8.93; platform_machine == 'x86_64'
|
||||
nvidia-nccl-cu12==2.26.2.post1; platform_machine == 'x86_64'
|
||||
nvidia-nvjitlink-cu12==12.8.93; platform_machine == 'x86_64'
|
||||
onnx==1.16.*; platform_machine == 'x86_64'
|
||||
onnx==1.18.*; platform_machine == 'x86_64'
|
||||
onnxruntime-gpu==1.24.*; platform_machine == 'x86_64'
|
||||
protobuf==3.20.3; platform_machine == 'x86_64'
|
||||
|
||||
@@ -19,9 +19,9 @@ If you already have Frigate installed through Docker or through a Home Assistant
|
||||
|
||||
## Setting up hardware
|
||||
|
||||
This section guides you through setting up a server with Debian Bookworm and Docker.
|
||||
This section guides you through setting up a server with Debian Trixie and Docker.
|
||||
|
||||
### Install Debian 12 (Bookworm)
|
||||
### Install Debian 13 (Trixie)
|
||||
|
||||
There are many guides on how to install Debian Server, so this will be an abbreviated guide. Connect a temporary monitor and keyboard to your device so you can install a minimal server without a desktop environment.
|
||||
|
||||
|
||||
@@ -26,7 +26,7 @@ VISUAL_WEIGHT = 0.65
|
||||
DESCRIPTION_WEIGHT = 0.35
|
||||
|
||||
|
||||
def chunk_content(content: str, chunk_size: int = 80) -> Generator[str, None, None]:
|
||||
def chunk_content(content: str, chunk_size: int = 80) -> Generator[str]:
|
||||
"""Yield content in word-aware chunks for streaming."""
|
||||
if not content:
|
||||
return
|
||||
|
||||
@@ -57,14 +57,14 @@ class PTZMetrics:
|
||||
reset: Event
|
||||
|
||||
def __init__(self, *, autotracker_enabled: bool):
|
||||
self.autotracker_enabled = mp.Value("i", autotracker_enabled) # type: ignore[assignment]
|
||||
self.autotracker_enabled = mp.Value("i", autotracker_enabled)
|
||||
|
||||
self.start_time = mp.Value("d", 0) # type: ignore[assignment]
|
||||
self.stop_time = mp.Value("d", 0) # type: ignore[assignment]
|
||||
self.frame_time = mp.Value("d", 0) # type: ignore[assignment]
|
||||
self.zoom_level = mp.Value("d", 0) # type: ignore[assignment]
|
||||
self.max_zoom = mp.Value("d", 0) # type: ignore[assignment]
|
||||
self.min_zoom = mp.Value("d", 0) # type: ignore[assignment]
|
||||
self.start_time = mp.Value("d", 0)
|
||||
self.stop_time = mp.Value("d", 0)
|
||||
self.frame_time = mp.Value("d", 0)
|
||||
self.zoom_level = mp.Value("d", 0)
|
||||
self.max_zoom = mp.Value("d", 0)
|
||||
self.min_zoom = mp.Value("d", 0)
|
||||
|
||||
self.tracking_active = mp.Event()
|
||||
self.motor_stopped = mp.Event()
|
||||
|
||||
@@ -38,7 +38,7 @@ class FaceRecognizer(ABC):
|
||||
def classify(self, face_image: np.ndarray) -> tuple[str, float] | None:
|
||||
pass
|
||||
|
||||
@redirect_output_to_logger(logger, logging.DEBUG) # type: ignore[misc]
|
||||
@redirect_output_to_logger(logger, logging.DEBUG) # type: ignore[untyped-decorator]
|
||||
def init_landmark_detector(self) -> None:
|
||||
landmark_model = os.path.join(MODEL_CACHE_DIR, "facedet/landmarkdet.yaml")
|
||||
|
||||
|
||||
@@ -10,7 +10,7 @@ import random
|
||||
import re
|
||||
import string
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
from typing import Any, cast
|
||||
|
||||
import cv2
|
||||
import numpy as np
|
||||
@@ -579,8 +579,8 @@ class LicensePlateProcessingMixin:
|
||||
boxes = []
|
||||
scores = []
|
||||
|
||||
for index in range(len(contours)): # type: ignore[arg-type]
|
||||
contour = contours[index] # type: ignore[index]
|
||||
for index in range(len(contours)):
|
||||
contour = contours[index]
|
||||
|
||||
# get minimum bounding box (rotated rectangle) around the contour and the smallest side length.
|
||||
points, sside = self._get_min_boxes(contour)
|
||||
@@ -1199,7 +1199,7 @@ class LicensePlateProcessingMixin:
|
||||
"""Look for license plates in image."""
|
||||
self.metrics.alpr_pps.value = self.plates_rec_second.eps()
|
||||
self.metrics.yolov9_lpr_pps.value = self.plates_det_second.eps()
|
||||
camera = obj_data if dedicated_lpr else obj_data["camera"]
|
||||
camera = cast(str, obj_data if dedicated_lpr else obj_data["camera"])
|
||||
current_time = int(datetime.datetime.now().timestamp())
|
||||
debug_frame_id = int(datetime.datetime.now().timestamp() * 1000)
|
||||
|
||||
|
||||
@@ -15,15 +15,11 @@ from frigate.config import FrigateConfig
|
||||
from frigate.const import MODEL_CACHE_DIR
|
||||
from frigate.log import suppress_stderr_during
|
||||
from frigate.util.image import calculate_region
|
||||
from frigate.util.tflite import Interpreter
|
||||
|
||||
from ..types import DataProcessorMetrics
|
||||
from .api import RealTimeProcessorApi
|
||||
|
||||
try:
|
||||
from tflite_runtime.interpreter import Interpreter
|
||||
except ModuleNotFoundError:
|
||||
from ai_edge_litert.interpreter import Interpreter
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
|
||||
@@ -18,15 +18,11 @@ from frigate.log import suppress_stderr_during
|
||||
from frigate.util.builtin import EventsPerSecond, InferenceSpeed, load_labels
|
||||
from frigate.util.image import calculate_region
|
||||
from frigate.util.object import box_overlaps
|
||||
from frigate.util.tflite import Interpreter
|
||||
|
||||
from ..types import DataProcessorMetrics
|
||||
from .api import DeferredRealtimeProcessorApi
|
||||
|
||||
try:
|
||||
from tflite_runtime.interpreter import Interpreter
|
||||
except ModuleNotFoundError:
|
||||
from ai_edge_litert.interpreter import Interpreter
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
MAX_OBJECT_CLASSIFICATIONS = 16
|
||||
|
||||
@@ -168,7 +168,7 @@ class FaceRealTimeProcessor(RealTimeProcessorApi):
|
||||
h: int = int(raw_bbox[3] / scale_factor)
|
||||
bbox = (x, y, x + w, y + h)
|
||||
|
||||
if face is None or area(bbox) > area(face): # type: ignore[unreachable]
|
||||
if face is None or area(bbox) > area(face):
|
||||
face = bbox
|
||||
|
||||
return face
|
||||
@@ -434,7 +434,7 @@ class FaceRealTimeProcessor(RealTimeProcessorApi):
|
||||
img = cv2.imread(current_file)
|
||||
|
||||
if img is None:
|
||||
return { # type: ignore[unreachable]
|
||||
return {
|
||||
"message": "Invalid image file.",
|
||||
"success": False,
|
||||
}
|
||||
|
||||
@@ -3,11 +3,7 @@ import os
|
||||
|
||||
import numpy as np
|
||||
|
||||
try:
|
||||
from tflite_runtime.interpreter import Interpreter, load_delegate
|
||||
except ModuleNotFoundError:
|
||||
from ai_edge_litert.interpreter import Interpreter, load_delegate
|
||||
|
||||
from frigate.util.tflite import Interpreter, load_delegate
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
@@ -6,15 +6,10 @@ from pydantic import ConfigDict, Field
|
||||
from frigate.detectors.detection_api import DetectionApi
|
||||
from frigate.detectors.detector_config import BaseDetectorConfig
|
||||
from frigate.log import suppress_stderr_during
|
||||
from frigate.util.tflite import Interpreter
|
||||
|
||||
from ..detector_utils import tflite_detect_raw, tflite_init
|
||||
|
||||
try:
|
||||
from tflite_runtime.interpreter import Interpreter
|
||||
except ModuleNotFoundError:
|
||||
from ai_edge_litert.interpreter import Interpreter
|
||||
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
DETECTOR_KEY = "cpu"
|
||||
|
||||
@@ -9,11 +9,7 @@ from pydantic import ConfigDict, Field
|
||||
|
||||
from frigate.detectors.detection_api import DetectionApi
|
||||
from frigate.detectors.detector_config import BaseDetectorConfig, ModelTypeEnum
|
||||
|
||||
try:
|
||||
from tflite_runtime.interpreter import Interpreter, load_delegate
|
||||
except ModuleNotFoundError:
|
||||
from ai_edge_litert.interpreter import Interpreter, load_delegate
|
||||
from frigate.util.tflite import Interpreter, load_delegate
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
@@ -10,15 +10,11 @@ from frigate.detectors.detection_runners import get_optimized_runner
|
||||
from frigate.embeddings.types import EnrichmentModelTypeEnum
|
||||
from frigate.log import suppress_stderr_during
|
||||
from frigate.util.downloader import ModelDownloader
|
||||
from frigate.util.tflite import Interpreter
|
||||
|
||||
from ...config import FaceRecognitionConfig
|
||||
from .base_embedding import BaseEmbedding
|
||||
|
||||
try:
|
||||
from tflite_runtime.interpreter import Interpreter
|
||||
except ModuleNotFoundError:
|
||||
from ai_edge_litert.interpreter import Interpreter
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
ARCFACE_INPUT_SIZE = 112
|
||||
|
||||
@@ -40,12 +40,7 @@ from frigate.log import LogPipe, suppress_stderr_during
|
||||
from frigate.util.builtin import get_ffmpeg_arg_list, load_labels
|
||||
from frigate.util.ffmpeg import start_or_restart_ffmpeg, stop_ffmpeg
|
||||
from frigate.util.process import FrigateProcess
|
||||
|
||||
try:
|
||||
from tflite_runtime.interpreter import Interpreter
|
||||
except ModuleNotFoundError:
|
||||
from ai_edge_litert.interpreter import Interpreter
|
||||
|
||||
from frigate.util.tflite import Interpreter
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
@@ -201,7 +201,7 @@ class EventCleanup(threading.Thread):
|
||||
self.config.record.alerts.retain.days,
|
||||
self.config.record.detections.retain.days,
|
||||
)
|
||||
file_extension = None # mp4 clips are no longer stored in /clips
|
||||
file_extension: str | None = None # mp4 clips are no longer stored in /clips
|
||||
update_params = {"has_clip": False}
|
||||
|
||||
# get expiration time for this label
|
||||
|
||||
@@ -94,7 +94,7 @@ class EventProcessor(threading.Thread):
|
||||
if update == None:
|
||||
continue
|
||||
|
||||
source_type, event_type, camera, _, event_data = update # type: ignore[misc]
|
||||
source_type, event_type, camera, _, event_data = update
|
||||
|
||||
logger.debug(
|
||||
f"Event received: {source_type} {event_type} {camera} {event_data['id']}"
|
||||
|
||||
@@ -37,13 +37,13 @@ __all__ = [
|
||||
"register_genai_provider",
|
||||
]
|
||||
|
||||
PROVIDERS = {}
|
||||
PROVIDERS: dict[GenAIProviderEnum, type["GenAIClient"]] = {}
|
||||
|
||||
|
||||
def register_genai_provider(key: GenAIProviderEnum) -> Callable:
|
||||
"""Register a GenAI provider."""
|
||||
|
||||
def decorator(cls: type) -> type:
|
||||
def decorator(cls: type["GenAIClient"]) -> type["GenAIClient"]:
|
||||
PROVIDERS[key] = cls
|
||||
return cls
|
||||
|
||||
@@ -414,7 +414,7 @@ class GenAIClient:
|
||||
tools: list[dict[str, Any]] | None = None,
|
||||
tool_choice: str | None = "auto",
|
||||
enable_thinking: bool | None = None,
|
||||
) -> AsyncGenerator[tuple[str, Any], None]:
|
||||
) -> AsyncGenerator[tuple[str, Any]]:
|
||||
"""Streaming counterpart to `chat_with_tools`.
|
||||
|
||||
Yields ``(kind, value)`` tuples where ``kind`` is one of:
|
||||
|
||||
@@ -432,7 +432,7 @@ class GeminiClient(GenAIClient):
|
||||
tools: list[dict[str, Any]] | None = None,
|
||||
tool_choice: str | None = "auto",
|
||||
enable_thinking: bool | None = None,
|
||||
) -> AsyncGenerator[tuple[str, Any], None]:
|
||||
) -> AsyncGenerator[tuple[str, Any]]:
|
||||
"""
|
||||
Stream chat with tools; yields content deltas then final message.
|
||||
|
||||
|
||||
@@ -853,7 +853,7 @@ class LlamaCppClient(GenAIClient):
|
||||
tools: list[dict[str, Any]] | None = None,
|
||||
tool_choice: str | None = "auto",
|
||||
enable_thinking: bool | None = None,
|
||||
) -> AsyncGenerator[tuple[str, Any], None]:
|
||||
) -> AsyncGenerator[tuple[str, Any]]:
|
||||
"""Stream chat with tools via OpenAI-compatible streaming API."""
|
||||
if self.provider is None:
|
||||
logger.warning(
|
||||
|
||||
@@ -430,7 +430,7 @@ class OllamaClient(GenAIClient):
|
||||
tools: list[dict[str, Any]] | None = None,
|
||||
tool_choice: str | None = "auto",
|
||||
enable_thinking: bool | None = None,
|
||||
) -> AsyncGenerator[tuple[str, Any], None]:
|
||||
) -> AsyncGenerator[tuple[str, Any]]:
|
||||
"""Stream chat with tools; yields content deltas then final message.
|
||||
|
||||
When tools are provided, Ollama streaming does not include tool_calls
|
||||
|
||||
@@ -311,7 +311,7 @@ class OpenAIClient(GenAIClient):
|
||||
tools: list[dict[str, Any]] | None = None,
|
||||
tool_choice: str | None = "auto",
|
||||
enable_thinking: bool | None = None,
|
||||
) -> AsyncGenerator[tuple[str, Any], None]:
|
||||
) -> AsyncGenerator[tuple[str, Any]]:
|
||||
"""
|
||||
Stream chat with tools; yields content deltas then final message.
|
||||
|
||||
|
||||
@@ -621,7 +621,7 @@ class MotionSearchRunner(threading.Thread):
|
||||
) -> list[MotionSearchResult]:
|
||||
"""Run detection while firing throttled progress as frames are scanned."""
|
||||
|
||||
def _gen() -> Generator[tuple[int, np.ndarray], None, None]:
|
||||
def _gen() -> Generator[tuple[int, np.ndarray]]:
|
||||
for i, frame in indexed_frames:
|
||||
if not self._should_stop():
|
||||
self._emit_progress(timestamp_fn(i))
|
||||
|
||||
@@ -186,7 +186,7 @@ def _run_vod_decode(
|
||||
skip_nonkey: bool,
|
||||
fps_rate: float | None,
|
||||
software_retry: bool,
|
||||
) -> Generator[np.ndarray, None, None]:
|
||||
) -> Generator[np.ndarray]:
|
||||
"""Run one VOD decode, yielding raw frames; retry in software if empty."""
|
||||
cmd = build_vod_decode_command(
|
||||
ffmpeg_path,
|
||||
@@ -258,7 +258,7 @@ def iter_vod_frames(
|
||||
*,
|
||||
skip_nonkey: bool,
|
||||
fps_rate: float | None,
|
||||
) -> Generator[np.ndarray, None, None]:
|
||||
) -> Generator[np.ndarray]:
|
||||
"""Decode a VOD HLS URL and yield raw frames in order.
|
||||
|
||||
Pair keyframe-mode output with probed keyframe PTS; pair fallback output with
|
||||
|
||||
+2
-2
@@ -202,7 +202,7 @@ class LogRedirect(io.StringIO):
|
||||
|
||||
|
||||
@contextmanager
|
||||
def __redirect_fd_to_queue(queue: Queue[str]) -> Generator[None, None, None]:
|
||||
def __redirect_fd_to_queue(queue: Queue[str]) -> Generator[None]:
|
||||
"""Redirect file descriptor 1 (stdout) to a pipe and capture output in a queue."""
|
||||
stdout_fd = os.dup(1)
|
||||
read_fd, write_fd = os.pipe()
|
||||
@@ -329,7 +329,7 @@ def suppress_os_output(func: Callable) -> Callable:
|
||||
|
||||
|
||||
@contextmanager
|
||||
def suppress_stderr_during(operation_name: str) -> Generator[None, None, None]:
|
||||
def suppress_stderr_during(operation_name: str) -> Generator[None]:
|
||||
"""
|
||||
Context manager to suppress stderr output during a specific operation.
|
||||
|
||||
|
||||
+1
-4
@@ -1,5 +1,5 @@
|
||||
[mypy]
|
||||
python_version = 3.11
|
||||
python_version = 3.13
|
||||
show_error_codes = true
|
||||
follow_imports = normal
|
||||
ignore_missing_imports = true
|
||||
@@ -45,9 +45,6 @@ ignore_errors = true
|
||||
[mypy-frigate.embeddings.*]
|
||||
ignore_errors = true
|
||||
|
||||
[mypy-frigate.http]
|
||||
ignore_errors = true
|
||||
|
||||
[mypy-frigate.ptz.*]
|
||||
ignore_errors = true
|
||||
|
||||
|
||||
@@ -353,7 +353,7 @@ class ObjectDetectProcess:
|
||||
logging.info("Detection process has exited...")
|
||||
|
||||
def start_or_restart(self) -> None:
|
||||
self.detection_start.value = 0.0 # type: ignore[attr-defined]
|
||||
self.detection_start.value = 0.0
|
||||
if (self.detect_process is not None) and self.detect_process.is_alive():
|
||||
self.stop()
|
||||
|
||||
|
||||
@@ -501,7 +501,7 @@ class RecordingMaintainer(threading.Thread):
|
||||
return None
|
||||
|
||||
def _compute_motion_heatmap(
|
||||
self, camera: str, motion_boxes: list[tuple[int, int, int, int]]
|
||||
self, camera: str, motion_boxes: list[tuple[int, ...]]
|
||||
) -> dict[str, int] | None:
|
||||
"""Compute a 16x16 motion intensity heatmap from motion boxes.
|
||||
|
||||
@@ -560,7 +560,7 @@ class RecordingMaintainer(threading.Thread):
|
||||
active_count = 0
|
||||
region_count = 0
|
||||
motion_count = 0
|
||||
all_motion_boxes: list[tuple[int, int, int, int]] = []
|
||||
all_motion_boxes: list[tuple[int, ...]] = []
|
||||
|
||||
for frame in self.object_recordings_info[camera]:
|
||||
# frame is after end time of segment
|
||||
|
||||
@@ -1,4 +0,0 @@
|
||||
from .multiprocessing import ServiceProcess
|
||||
from .service import Service, ServiceManager
|
||||
|
||||
__all__ = ["Service", "ServiceProcess", "ServiceManager"]
|
||||
@@ -1,162 +0,0 @@
|
||||
import asyncio
|
||||
import faulthandler
|
||||
import logging
|
||||
import multiprocessing as mp
|
||||
import signal
|
||||
import sys
|
||||
import threading
|
||||
from abc import ABC, abstractmethod
|
||||
from asyncio.exceptions import TimeoutError
|
||||
from logging.handlers import QueueHandler
|
||||
from types import FrameType
|
||||
|
||||
import frigate.log
|
||||
|
||||
from .multiprocessing_waiter import wait as mp_wait
|
||||
from .service import Service, ServiceManager
|
||||
|
||||
DEFAULT_STOP_TIMEOUT = 10 # seconds
|
||||
|
||||
|
||||
class BaseServiceProcess(Service, ABC):
|
||||
"""A Service the manages a multiprocessing.Process."""
|
||||
|
||||
_process: mp.Process | None
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
*,
|
||||
name: str | None = None,
|
||||
manager: ServiceManager | None = None,
|
||||
) -> None:
|
||||
super().__init__(name=name, manager=manager)
|
||||
|
||||
self._process = None
|
||||
|
||||
async def on_start(self) -> None:
|
||||
if self._process is not None:
|
||||
if self._process.is_alive():
|
||||
return # Already started.
|
||||
else:
|
||||
self._process.close()
|
||||
|
||||
# At this point, the process is either stopped or dead, so we can recreate it.
|
||||
self._process = mp.Process(target=self._run)
|
||||
self._process.name = self.name
|
||||
self._process.daemon = True
|
||||
self.before_start()
|
||||
self._process.start()
|
||||
self.after_start()
|
||||
|
||||
self.manager.logger.info(f"Started {self.name} (pid: {self._process.pid})")
|
||||
|
||||
async def on_stop(
|
||||
self,
|
||||
*,
|
||||
force: bool = False,
|
||||
timeout: float | None = None,
|
||||
) -> None:
|
||||
if timeout is None:
|
||||
timeout = DEFAULT_STOP_TIMEOUT
|
||||
|
||||
if self._process is None:
|
||||
return # Already stopped.
|
||||
|
||||
running = True
|
||||
|
||||
if not force:
|
||||
self._process.terminate()
|
||||
try:
|
||||
await asyncio.wait_for(mp_wait(self._process), timeout)
|
||||
running = False
|
||||
except TimeoutError:
|
||||
self.manager.logger.warning(
|
||||
f"{self.name} is still running after {timeout} seconds. Killing."
|
||||
)
|
||||
|
||||
if running:
|
||||
self._process.kill()
|
||||
await mp_wait(self._process)
|
||||
|
||||
self._process.close()
|
||||
self._process = None
|
||||
|
||||
self.manager.logger.info(f"{self.name} stopped")
|
||||
|
||||
@property
|
||||
def pid(self) -> int | None:
|
||||
return self._process.pid if self._process else None
|
||||
|
||||
def _run(self) -> None:
|
||||
self.before_run()
|
||||
self.run()
|
||||
self.after_run()
|
||||
|
||||
def before_start(self) -> None:
|
||||
pass
|
||||
|
||||
def after_start(self) -> None:
|
||||
pass
|
||||
|
||||
def before_run(self) -> None:
|
||||
pass
|
||||
|
||||
def after_run(self) -> None:
|
||||
pass
|
||||
|
||||
@abstractmethod
|
||||
def run(self) -> None:
|
||||
pass
|
||||
|
||||
def __getstate__(self) -> dict:
|
||||
return {
|
||||
k: v
|
||||
for k, v in self.__dict__.items()
|
||||
if not (k.startswith("_Service__") or k == "_process")
|
||||
}
|
||||
|
||||
|
||||
class ServiceProcess(BaseServiceProcess):
|
||||
logger: logging.Logger
|
||||
|
||||
@property
|
||||
def stop_event(self) -> threading.Event:
|
||||
# Lazily create the stop_event. This allows the signal handler to tell if anyone is
|
||||
# monitoring the stop event, and to raise a SystemExit if not.
|
||||
if "stop_event" not in self.__dict__:
|
||||
stop_event = threading.Event()
|
||||
self.__dict__["stop_event"] = stop_event
|
||||
else:
|
||||
stop_event = self.__dict__["stop_event"]
|
||||
assert isinstance(stop_event, threading.Event)
|
||||
|
||||
return stop_event
|
||||
|
||||
def before_start(self) -> None:
|
||||
if frigate.log.log_listener is None:
|
||||
raise RuntimeError("Logging has not yet been set up.")
|
||||
self.__log_queue = frigate.log.log_listener.queue
|
||||
|
||||
def before_run(self) -> None:
|
||||
super().before_run()
|
||||
|
||||
faulthandler.enable()
|
||||
|
||||
def receiveSignal(signalNumber: int, frame: FrameType | None) -> None:
|
||||
# Get the stop_event through the dict to bypass lazy initialization.
|
||||
stop_event = self.__dict__.get("stop_event")
|
||||
if stop_event is not None:
|
||||
# Someone is monitoring stop_event. We should set it.
|
||||
stop_event.set()
|
||||
else:
|
||||
# Nobody is monitoring stop_event. We should raise SystemExit.
|
||||
sys.exit()
|
||||
|
||||
signal.signal(signal.SIGTERM, receiveSignal)
|
||||
signal.signal(signal.SIGINT, receiveSignal)
|
||||
|
||||
self.logger = logging.getLogger(self.name)
|
||||
|
||||
logging.basicConfig(handlers=[], force=True)
|
||||
logging.getLogger().addHandler(QueueHandler(self.__log_queue))
|
||||
del self.__log_queue
|
||||
@@ -1,150 +0,0 @@
|
||||
import asyncio
|
||||
import functools
|
||||
import logging
|
||||
import multiprocessing as mp
|
||||
import queue
|
||||
import threading
|
||||
from multiprocessing.connection import Connection
|
||||
from multiprocessing.connection import wait as mp_wait
|
||||
from socket import socket
|
||||
from typing import Any
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class MultiprocessingWaiter(threading.Thread):
|
||||
"""A background thread that manages futures for the multiprocessing.connection.wait() method."""
|
||||
|
||||
def __init__(self) -> None:
|
||||
super().__init__(daemon=True)
|
||||
|
||||
# Queue of objects to wait for and futures to set results for.
|
||||
self._queue: queue.Queue[tuple[Any, asyncio.Future[None]]] = queue.Queue()
|
||||
|
||||
# This is required to get mp_wait() to wake up when new objects to wait for are received.
|
||||
receive, send = mp.Pipe(duplex=False)
|
||||
self._receive_connection = receive
|
||||
self._send_connection = send
|
||||
|
||||
def wait_for_sentinel(self, sentinel: Any) -> asyncio.Future[None]:
|
||||
"""Create an asyncio.Future tracking a sentinel for multiprocessing.connection.wait()
|
||||
|
||||
Warning: This method is NOT thread-safe.
|
||||
"""
|
||||
# This would be incredibly stupid, but you never know.
|
||||
assert sentinel != self._receive_connection
|
||||
|
||||
# Send the future to the background thread for processing.
|
||||
future = asyncio.get_running_loop().create_future()
|
||||
self._queue.put((sentinel, future))
|
||||
|
||||
# Notify the background thread.
|
||||
#
|
||||
# This is the non-thread-safe part, but since this method is not really meant to be called
|
||||
# by users, we can get away with not adding a lock at this point (to avoid adding 2 locks).
|
||||
self._send_connection.send_bytes(b".")
|
||||
|
||||
return future
|
||||
|
||||
def run(self) -> None:
|
||||
logger.debug("Started background thread")
|
||||
|
||||
wait_dict: dict[Any, set[asyncio.Future[None]]] = {
|
||||
self._receive_connection: set()
|
||||
}
|
||||
while True:
|
||||
for ready_obj in mp_wait(wait_dict.keys()):
|
||||
# Make sure we never remove the receive connection from the wait dict
|
||||
if ready_obj is self._receive_connection:
|
||||
continue
|
||||
|
||||
logger.debug(
|
||||
f"Sentinel {ready_obj!r} is ready. "
|
||||
f"Notifying {len(wait_dict[ready_obj])} future(s)."
|
||||
)
|
||||
|
||||
# Go over all the futures attached to this object and mark them as ready.
|
||||
for fut in wait_dict.pop(ready_obj):
|
||||
if fut.cancelled():
|
||||
logger.debug(
|
||||
f"A future for sentinel {ready_obj!r} is ready, "
|
||||
"but the future is cancelled. Skipping."
|
||||
)
|
||||
else:
|
||||
fut.get_loop().call_soon_threadsafe(
|
||||
# Note: We need to check fut.cancelled() again, since it might
|
||||
# have been set before the event loop's definition of "soon".
|
||||
functools.partial(
|
||||
lambda fut: fut.cancelled() or fut.set_result(None), fut
|
||||
)
|
||||
)
|
||||
|
||||
# Check for cancellations in the remaining futures.
|
||||
done_objects = []
|
||||
for obj, fut_set in wait_dict.items():
|
||||
if obj is self._receive_connection:
|
||||
continue
|
||||
|
||||
# Find any cancelled futures and remove them.
|
||||
cancelled = [fut for fut in fut_set if fut.cancelled()]
|
||||
fut_set.difference_update(cancelled)
|
||||
logger.debug(
|
||||
f"Removing {len(cancelled)} future(s) from sentinel: {obj!r}"
|
||||
)
|
||||
|
||||
# Mark objects with no remaining futures for removal.
|
||||
if len(fut_set) == 0:
|
||||
done_objects.append(obj)
|
||||
|
||||
# Remove any objects that are done after removing cancelled futures.
|
||||
for obj in done_objects:
|
||||
logger.debug(
|
||||
f"Sentinel {obj!r} no longer has any futures waiting for it."
|
||||
)
|
||||
del wait_dict[obj]
|
||||
|
||||
# Get new objects to wait for from the queue.
|
||||
while True:
|
||||
try:
|
||||
obj, fut = self._queue.get_nowait()
|
||||
self._receive_connection.recv_bytes(maxlength=1)
|
||||
self._queue.task_done()
|
||||
|
||||
logger.debug(f"Received new sentinel: {obj!r}")
|
||||
|
||||
wait_dict.setdefault(obj, set()).add(fut)
|
||||
except queue.Empty:
|
||||
break
|
||||
|
||||
|
||||
waiter_lock = threading.Lock()
|
||||
waiter_thread: MultiprocessingWaiter | None = None
|
||||
|
||||
|
||||
async def wait(object: mp.Process | Connection | socket) -> None:
|
||||
"""Wait for the supplied object to be ready.
|
||||
|
||||
Under the hood, this uses multiprocessing.connection.wait() and a background thread manage the
|
||||
returned futures.
|
||||
"""
|
||||
global waiter_thread, waiter_lock
|
||||
|
||||
sentinel: Connection | socket | int
|
||||
if isinstance(object, mp.Process):
|
||||
sentinel = object.sentinel
|
||||
elif isinstance(object, Connection) or isinstance(object, socket):
|
||||
sentinel = object
|
||||
else:
|
||||
raise ValueError(f"Cannot wait for object of type {type(object).__qualname__}")
|
||||
|
||||
with waiter_lock:
|
||||
if waiter_thread is None:
|
||||
# Start a new waiter thread.
|
||||
waiter_thread = MultiprocessingWaiter()
|
||||
waiter_thread.start()
|
||||
|
||||
# Create the future while still holding the lock,
|
||||
# since wait_for_sentinel() is not thread safe.
|
||||
fut = waiter_thread.wait_for_sentinel(sentinel)
|
||||
|
||||
await fut
|
||||
@@ -1,445 +0,0 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import atexit
|
||||
import logging
|
||||
import threading
|
||||
from abc import ABC, abstractmethod
|
||||
from collections.abc import Coroutine
|
||||
from contextvars import ContextVar
|
||||
from dataclasses import dataclass
|
||||
from functools import partial
|
||||
from typing import Self, cast
|
||||
|
||||
|
||||
class Service(ABC):
|
||||
"""An abstract service instance."""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
*,
|
||||
name: str | None = None,
|
||||
manager: ServiceManager | None = None,
|
||||
):
|
||||
if name:
|
||||
self.__dict__["name"] = name
|
||||
|
||||
self.__manager = manager or ServiceManager.current()
|
||||
self.__lock = asyncio.Lock(loop=self.__manager._event_loop) # type: ignore[call-arg]
|
||||
self.__manager._register(self)
|
||||
|
||||
@property
|
||||
def name(self) -> str:
|
||||
try:
|
||||
return cast(str, self.__dict__["name"])
|
||||
except KeyError:
|
||||
return type(self).__qualname__
|
||||
|
||||
@property
|
||||
def manager(self) -> ServiceManager:
|
||||
"""The service manager this service is registered with."""
|
||||
try:
|
||||
return self.__manager
|
||||
except AttributeError:
|
||||
raise RuntimeError("Cannot access associated service manager") from None
|
||||
|
||||
def start(
|
||||
self,
|
||||
*,
|
||||
wait: bool = False,
|
||||
wait_timeout: float | None = None,
|
||||
) -> Self:
|
||||
"""Start this service.
|
||||
|
||||
:param wait: If set, this function will block until the task is complete.
|
||||
:param wait_timeout: If set, this function will not return until the task is complete or the
|
||||
specified timeout has elapsed.
|
||||
"""
|
||||
|
||||
self.manager.run_task(
|
||||
self.on_start(),
|
||||
wait=wait,
|
||||
wait_timeout=wait_timeout,
|
||||
lock=self.__lock,
|
||||
)
|
||||
|
||||
return self
|
||||
|
||||
def stop(
|
||||
self,
|
||||
*,
|
||||
force: bool = False,
|
||||
timeout: float | None = None,
|
||||
wait: bool = False,
|
||||
wait_timeout: float | None = None,
|
||||
) -> Self:
|
||||
"""Stop this service.
|
||||
|
||||
:param force: If set, the service will be killed immediately.
|
||||
:param timeout: Maximum amount of time to wait before force-killing the service.
|
||||
|
||||
:param wait: If set, this function will block until the task is complete.
|
||||
:param wait_timeout: If set, this function will not return until the task is complete or the
|
||||
specified timeout has elapsed.
|
||||
"""
|
||||
|
||||
self.manager.run_task(
|
||||
self.on_stop(force=force, timeout=timeout),
|
||||
wait=wait,
|
||||
wait_timeout=wait_timeout,
|
||||
lock=self.__lock,
|
||||
)
|
||||
|
||||
return self
|
||||
|
||||
def restart(
|
||||
self,
|
||||
*,
|
||||
force: bool = False,
|
||||
stop_timeout: float | None = None,
|
||||
wait: bool = False,
|
||||
wait_timeout: float | None = None,
|
||||
) -> Self:
|
||||
"""Restart this service.
|
||||
|
||||
:param force: If set, the service will be killed immediately.
|
||||
:param timeout: Maximum amount of time to wait before force-killing the service.
|
||||
|
||||
:param wait: If set, this function will block until the task is complete.
|
||||
:param wait_timeout: If set, this function will not return until the task is complete or the
|
||||
specified timeout has elapsed.
|
||||
"""
|
||||
|
||||
self.manager.run_task(
|
||||
self.on_restart(force=force, stop_timeout=stop_timeout),
|
||||
wait=wait,
|
||||
wait_timeout=wait_timeout,
|
||||
lock=self.__lock,
|
||||
)
|
||||
|
||||
return self
|
||||
|
||||
@abstractmethod
|
||||
async def on_start(self) -> None:
|
||||
pass
|
||||
|
||||
@abstractmethod
|
||||
async def on_stop(
|
||||
self,
|
||||
*,
|
||||
force: bool = False,
|
||||
timeout: float | None = None,
|
||||
) -> None:
|
||||
pass
|
||||
|
||||
async def on_restart(
|
||||
self,
|
||||
*,
|
||||
force: bool = False,
|
||||
stop_timeout: float | None = None,
|
||||
) -> None:
|
||||
await self.on_stop(force=force, timeout=stop_timeout)
|
||||
await self.on_start()
|
||||
|
||||
|
||||
default_service_manager_lock = threading.Lock()
|
||||
default_service_manager: ServiceManager | None = None
|
||||
|
||||
current_service_manager: ContextVar[ServiceManager] = ContextVar(
|
||||
"current_service_manager"
|
||||
)
|
||||
|
||||
|
||||
@dataclass
|
||||
class Command:
|
||||
"""A coroutine to execute in the service manager thread.
|
||||
|
||||
Attributes:
|
||||
coro: The coroutine to execute.
|
||||
lock: An async lock to acquire before calling the coroutine.
|
||||
done: If specified, the service manager will set this event after the command completes.
|
||||
"""
|
||||
|
||||
coro: Coroutine
|
||||
lock: asyncio.Lock | None = None
|
||||
done: threading.Event | None = None
|
||||
|
||||
|
||||
class ServiceManager:
|
||||
"""A set of services, along with the global state required to manage them efficiently.
|
||||
|
||||
Typically users of the service infrastructure will not interact with a service manager directly,
|
||||
but rather through individual Service subclasses that will automatically manage a service
|
||||
manager instance.
|
||||
|
||||
Each service manager instance has a background thread in which service lifecycle tasks are
|
||||
executed in an async executor. This is done to avoid head-of-line blocking in the business logic
|
||||
that spins up individual services. This thread is automatically started when the service manager
|
||||
is created and stopped either manually, or on application exit.
|
||||
|
||||
All (public) service manager methods are thread-safe.
|
||||
"""
|
||||
|
||||
_name: str
|
||||
_logger: logging.Logger
|
||||
|
||||
# The set of services this service manager knows about.
|
||||
_services: dict[str, Service]
|
||||
_services_lock: threading.Lock
|
||||
|
||||
# Commands will be queued with associated event loop. Queueing `None` signals shutdown.
|
||||
_command_queue: asyncio.Queue[Command | None]
|
||||
_event_loop: asyncio.AbstractEventLoop
|
||||
|
||||
# The pending command counter is used to ensure all commands have been queued before shutdown.
|
||||
_pending_commands: AtomicCounter
|
||||
|
||||
# The set of pending tasks after they have been received by the background thread and spawned.
|
||||
_tasks: set
|
||||
|
||||
# Fired after the async runtime starts. Object initialization completes after this is set.
|
||||
_setup_event: threading.Event
|
||||
|
||||
# Will be acquired to ensure the shutdown sentinel is sent only once. Never released.
|
||||
_shutdown_lock: threading.Lock
|
||||
|
||||
def __init__(self, *, name: str | None = None):
|
||||
self._name = name if name is not None else (__package__ or __name__)
|
||||
self._logger = logging.getLogger(self.name)
|
||||
|
||||
self._services = dict()
|
||||
self._services_lock = threading.Lock()
|
||||
|
||||
self._pending_commands = AtomicCounter()
|
||||
self._tasks = set()
|
||||
|
||||
self._shutdown_lock = threading.Lock()
|
||||
|
||||
# --- Start the manager thread and wait for it to be ready. ---
|
||||
|
||||
self._setup_event = threading.Event()
|
||||
|
||||
async def start_manager() -> None:
|
||||
self._event_loop = asyncio.get_running_loop()
|
||||
self._command_queue = asyncio.Queue()
|
||||
|
||||
self._setup_event.set()
|
||||
await self._monitor_command_queue()
|
||||
|
||||
self._manager_thread = threading.Thread(
|
||||
name=self.name,
|
||||
target=lambda: asyncio.run(start_manager()),
|
||||
daemon=True,
|
||||
)
|
||||
|
||||
self._manager_thread.start()
|
||||
atexit.register(partial(self.shutdown, wait=True))
|
||||
|
||||
self._setup_event.wait()
|
||||
|
||||
@property
|
||||
def name(self) -> str:
|
||||
"""The name of this service manager. Primarily intended for logging purposes."""
|
||||
return self._name
|
||||
|
||||
@property
|
||||
def logger(self) -> logging.Logger:
|
||||
"""The logger used by this service manager."""
|
||||
return self._logger
|
||||
|
||||
@classmethod
|
||||
def current(cls) -> ServiceManager:
|
||||
"""The service manager set in the current context (async task or thread).
|
||||
|
||||
A global default service manager will be automatically created on first access."""
|
||||
|
||||
global default_service_manager
|
||||
|
||||
current = current_service_manager.get(None)
|
||||
if current is None:
|
||||
with default_service_manager_lock:
|
||||
if default_service_manager is None:
|
||||
default_service_manager = cls()
|
||||
|
||||
current = default_service_manager
|
||||
current_service_manager.set(current)
|
||||
return current
|
||||
|
||||
def make_current(self) -> None:
|
||||
"""Make this the current service manager."""
|
||||
|
||||
current_service_manager.set(self)
|
||||
|
||||
def run_task(
|
||||
self,
|
||||
coro: Coroutine,
|
||||
*,
|
||||
wait: bool = False,
|
||||
wait_timeout: float | None = None,
|
||||
lock: asyncio.Lock | None = None,
|
||||
) -> None:
|
||||
"""Run an async task in the service manager thread.
|
||||
|
||||
:param wait: If set, this function will block until the task is complete.
|
||||
:param wait_timeout: If set, this function will not return until the task is complete or the
|
||||
specified timeout has elapsed.
|
||||
"""
|
||||
|
||||
if not isinstance(coro, Coroutine):
|
||||
raise TypeError(f"Cannot schedule task for object of type {type(coro)}")
|
||||
|
||||
cmd = Command(coro=coro, lock=lock)
|
||||
if wait or wait_timeout is not None:
|
||||
cmd.done = threading.Event()
|
||||
|
||||
self._send_command(cmd)
|
||||
|
||||
if cmd.done is not None:
|
||||
cmd.done.wait(timeout=wait_timeout)
|
||||
|
||||
def shutdown(
|
||||
self, *, wait: bool = False, wait_timeout: float | None = None
|
||||
) -> None:
|
||||
"""Shutdown the service manager thread.
|
||||
|
||||
After the shutdown process completes, any subsequent calls to the service manager will
|
||||
produce an error.
|
||||
|
||||
:param wait: If set, this function will block until the shutdown process is complete.
|
||||
:param wait_timeout: If set, this function will not return until the shutdown process is
|
||||
complete or the specified timeout has elapsed.
|
||||
"""
|
||||
|
||||
if self._shutdown_lock.acquire(blocking=False):
|
||||
self._send_command(None)
|
||||
if wait:
|
||||
self._manager_thread.join(timeout=wait_timeout)
|
||||
|
||||
def _ensure_running(self) -> None:
|
||||
self._setup_event.wait()
|
||||
if not self._manager_thread.is_alive():
|
||||
raise RuntimeError(f"ServiceManager {self.name} is not running")
|
||||
|
||||
def _send_command(self, command: Command | None) -> None:
|
||||
self._ensure_running()
|
||||
|
||||
async def queue_command() -> None:
|
||||
await self._command_queue.put(command)
|
||||
self._pending_commands.sub()
|
||||
|
||||
self._pending_commands.add()
|
||||
asyncio.run_coroutine_threadsafe(queue_command(), self._event_loop)
|
||||
|
||||
def _register(self, service: Service) -> None:
|
||||
"""Register a service with the service manager. This is done by the service constructor."""
|
||||
|
||||
self._ensure_running()
|
||||
with self._services_lock:
|
||||
name_conflict: Service | None = next(
|
||||
(
|
||||
existing
|
||||
for name, existing in self._services.items()
|
||||
if name == service.name
|
||||
),
|
||||
None,
|
||||
)
|
||||
|
||||
if name_conflict is service:
|
||||
raise RuntimeError(f"Attempt to re-register service: {service.name}")
|
||||
elif name_conflict is not None:
|
||||
raise RuntimeError(f"Duplicate service name: {service.name}")
|
||||
|
||||
self.logger.debug(f"Registering service: {service.name}")
|
||||
self._services[service.name] = service
|
||||
|
||||
def _run_command(self, command: Command) -> None:
|
||||
"""Execute a command and add it to the tasks set."""
|
||||
|
||||
def task_done(task: asyncio.Task) -> None:
|
||||
exc = task.exception()
|
||||
if exc:
|
||||
self.logger.exception("Exception in service manager task", exc_info=exc)
|
||||
self._tasks.discard(task)
|
||||
if command.done is not None:
|
||||
command.done.set()
|
||||
|
||||
async def task_harness() -> None:
|
||||
if command.lock is not None:
|
||||
async with command.lock:
|
||||
await command.coro
|
||||
else:
|
||||
await command.coro
|
||||
|
||||
task = asyncio.create_task(task_harness())
|
||||
task.add_done_callback(task_done)
|
||||
self._tasks.add(task)
|
||||
|
||||
async def _monitor_command_queue(self) -> None:
|
||||
"""The main function of the background thread."""
|
||||
|
||||
self.logger.info("Started service manager")
|
||||
|
||||
# Main command processing loop.
|
||||
while (command := await self._command_queue.get()) is not None:
|
||||
self._run_command(command)
|
||||
|
||||
# Send a stop command to all services. We don't have a status command yet, so we can just
|
||||
# stop everything and be done with it.
|
||||
with self._services_lock:
|
||||
self.logger.debug(f"Stopping {len(self._services)} services")
|
||||
for service in self._services.values():
|
||||
service.stop()
|
||||
|
||||
# Wait for all commands to finish executing.
|
||||
await self._shutdown()
|
||||
|
||||
self.logger.info("Exiting service manager")
|
||||
|
||||
async def _shutdown(self) -> None:
|
||||
"""Ensure all commands have been queued & executed."""
|
||||
|
||||
while True:
|
||||
command = None
|
||||
try:
|
||||
# Try and get a command from the queue.
|
||||
command = self._command_queue.get_nowait()
|
||||
except asyncio.QueueEmpty:
|
||||
if self._pending_commands.value > 0:
|
||||
# If there are pending commands to queue, await them.
|
||||
command = await self._command_queue.get()
|
||||
elif self._tasks:
|
||||
# If there are still pending tasks, wait for them. These tasks might queue
|
||||
# commands though, so we have to loop again.
|
||||
await asyncio.wait(self._tasks)
|
||||
else:
|
||||
# Nothing is pending at this point, so we're done here.
|
||||
break
|
||||
|
||||
# If we got a command, run it.
|
||||
if command is not None:
|
||||
self._run_command(command)
|
||||
|
||||
|
||||
class AtomicCounter:
|
||||
"""A lock-protected atomic counter."""
|
||||
|
||||
# Modern CPUs have atomics, but python doesn't seem to include them in the standard library.
|
||||
# Besides, the performance penalty is negligible compared to, well, using python.
|
||||
# So this will do just fine.
|
||||
|
||||
def __init__(self, initial: int = 0):
|
||||
self._lock = threading.Lock()
|
||||
self._value = initial
|
||||
|
||||
def add(self, value: int = 1) -> None:
|
||||
with self._lock:
|
||||
self._value += value
|
||||
|
||||
def sub(self, value: int = 1) -> None:
|
||||
with self._lock:
|
||||
self._value -= value
|
||||
|
||||
@property
|
||||
def value(self) -> int:
|
||||
with self._lock:
|
||||
return self._value
|
||||
+6
-33
@@ -3,9 +3,7 @@
|
||||
import datetime
|
||||
import logging
|
||||
import subprocess as sp
|
||||
import threading
|
||||
from abc import ABC, abstractmethod
|
||||
from multiprocessing import resource_tracker as _mprt
|
||||
from multiprocessing import shared_memory as _mpshm
|
||||
from string import printable
|
||||
from typing import Any, AnyStr
|
||||
@@ -1015,9 +1013,12 @@ class FrameManager(ABC):
|
||||
|
||||
|
||||
class UntrackedSharedMemory(_mpshm.SharedMemory):
|
||||
# https://github.com/python/cpython/issues/82300#issuecomment-2169035092
|
||||
"""SharedMemory that is not registered with the resource tracker.
|
||||
|
||||
__lock = threading.Lock()
|
||||
Frigate manages the lifecycle of its shared memory segments across
|
||||
processes, so resource tracker cleanup (and its noisy leak warnings)
|
||||
is unwanted. https://github.com/python/cpython/issues/82300
|
||||
"""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
@@ -1027,35 +1028,7 @@ class UntrackedSharedMemory(_mpshm.SharedMemory):
|
||||
*,
|
||||
track: bool = False,
|
||||
) -> None:
|
||||
self._track = track
|
||||
|
||||
# if tracking, normal init will suffice
|
||||
if track:
|
||||
return super().__init__(name=name, create=create, size=size)
|
||||
|
||||
# lock so that other threads don't attempt to use the
|
||||
# register function during this time
|
||||
with self.__lock:
|
||||
# temporarily disable registration during initialization
|
||||
orig_register = _mprt.register
|
||||
_mprt.register = self.__tmp_register
|
||||
|
||||
# initialize; ensure original register function is
|
||||
# re-instated
|
||||
try:
|
||||
super().__init__(name=name, create=create, size=size)
|
||||
finally:
|
||||
_mprt.register = orig_register
|
||||
|
||||
@staticmethod
|
||||
def __tmp_register(*args, **kwargs) -> None:
|
||||
return
|
||||
|
||||
def unlink(self) -> None:
|
||||
if _mpshm._USE_POSIX and self._name:
|
||||
_mpshm._posixshmem.shm_unlink(self._name)
|
||||
if self._track:
|
||||
_mprt.unregister(self._name, "shared_memory")
|
||||
super().__init__(name=name, create=create, size=size, track=track)
|
||||
|
||||
|
||||
class SharedMemoryFrameManager(FrameManager):
|
||||
|
||||
@@ -0,0 +1,8 @@
|
||||
"""TFLite interpreter imports, preferring LiteRT with tflite_runtime fallback."""
|
||||
|
||||
try:
|
||||
from ai_edge_litert.interpreter import Interpreter, load_delegate
|
||||
except ModuleNotFoundError:
|
||||
from tflite_runtime.interpreter import Interpreter, load_delegate
|
||||
|
||||
__all__ = ["Interpreter", "load_delegate"]
|
||||
@@ -68,14 +68,14 @@ def get_dst_transitions(
|
||||
current = start_time
|
||||
|
||||
# Get initial offset
|
||||
dt = datetime.datetime.utcfromtimestamp(current).replace(tzinfo=pytz.UTC)
|
||||
dt = datetime.datetime.fromtimestamp(current, tz=pytz.UTC)
|
||||
local_dt = dt.astimezone(tz)
|
||||
prev_offset = local_dt.utcoffset().total_seconds()
|
||||
period_start = start_time
|
||||
|
||||
# Check each day for offset changes
|
||||
while current <= end_time:
|
||||
dt = datetime.datetime.utcfromtimestamp(current).replace(tzinfo=pytz.UTC)
|
||||
dt = datetime.datetime.fromtimestamp(current, tz=pytz.UTC)
|
||||
local_dt = dt.astimezone(tz)
|
||||
current_offset = local_dt.utcoffset().total_seconds()
|
||||
|
||||
|
||||
+1
-3
@@ -118,9 +118,7 @@ class FrigateWatchdog(threading.Thread):
|
||||
|
||||
# check the detection processes
|
||||
for detector in self.detectors.values():
|
||||
detection_start = detector.detection_start.value # type: ignore[attr-defined]
|
||||
# issue https://github.com/python/typeshed/issues/8799
|
||||
# from mypy 0.981 onwards
|
||||
detection_start = detector.detection_start.value
|
||||
if detection_start > 0.0 and now - detection_start > 10:
|
||||
logger.info(
|
||||
"Detection appears to be stuck. Restarting detection process..."
|
||||
|
||||
+2
-2
@@ -1,6 +1,6 @@
|
||||
[tool.ruff]
|
||||
target-version = "py311"
|
||||
target-version = "py313"
|
||||
|
||||
[tool.ruff.lint]
|
||||
ignore = ["E501","E711","E712","UP031","UP032","UP042","G004"]
|
||||
ignore = ["E501","E711","E712","UP031","UP032","UP042","UP046","UP047","G004"]
|
||||
extend-select = ["I", "UP", "G", "ASYNC210", "B904"]
|
||||
|
||||
Reference in New Issue
Block a user