Compare commits

...
2 Commits
Author SHA1 Message Date
c4f3fbff69 Dynamically resolve Intel NPU (#23761)
* Add support for newer Intel NPU busy time counter

* Resolve Intel NPU device dynamically

---------

Co-authored-by: Filious Louis <1417132+fjlouis@users.noreply.github.com>
2026-07-29 16:35:46 -06:00
DoFabienandGitHub 0f636a17f3 Improve recording timeline and VOD query performance (#23862)
* Improve recording timeline and VOD query performance

* Add recording query boundary tests
2026-07-29 16:03:39 -06:00
5 changed files with 288 additions and 11 deletions
+4 -4
View File
@@ -574,11 +574,11 @@ async def vod_ts(
Recordings.start_time,
)
.where(
Recordings.start_time.between(start_ts, end_ts)
| Recordings.end_time.between(start_ts, end_ts)
| ((start_ts > Recordings.start_time) & (end_ts < Recordings.end_time))
Recordings.camera == camera_name,
Recordings.start_time >= start_ts - MAX_SEGMENT_DURATION,
Recordings.start_time <= end_ts,
Recordings.end_time >= start_ts,
)
.where(Recordings.camera == camera_name)
.order_by(Recordings.start_time.asc())
.iterator()
)
+2 -1
View File
@@ -25,7 +25,7 @@ from frigate.api.defs.query.recordings_query_parameters import (
)
from frigate.api.defs.response.generic_response import GenericResponse
from frigate.api.defs.tags import Tags
from frigate.const import RECORD_DIR
from frigate.const import MAX_SEGMENT_DURATION, RECORD_DIR
from frigate.models import Event, Recordings
from frigate.util.time import get_dst_transitions
@@ -243,6 +243,7 @@ async def recordings(
)
.where(
Recordings.camera == camera_name,
Recordings.start_time >= after - MAX_SEGMENT_DURATION,
Recordings.end_time >= after,
Recordings.start_time <= before,
)
+174
View File
@@ -1,15 +1,73 @@
"""Unit tests for recordings/media API endpoints."""
from dataclasses import dataclass
from datetime import UTC, datetime
from unittest.mock import patch
import pytz
from fastapi import Request
from frigate.api.auth import get_allowed_cameras_for_filter, get_current_user
from frigate.const import MAX_SEGMENT_DURATION
from frigate.models import Recordings
from frigate.test.http_api.base_http_test import AuthTestClient, BaseTestHttp
@dataclass(frozen=True)
class RangeCase:
"""Expected behavior for one segment relative to the requested range.
Offsets are seconds from REQUEST_START; the request ends at +100 seconds.
"""
name: str
start_offset: float
end_offset: float
included_in_recordings: bool
vod_clip_from_ms: int | None = None
vod_duration_ms: int | None = None
REQUEST_START = 1000
REQUEST_END = 1100
RANGE_CASES = (
RangeCase("before", -MAX_SEGMENT_DURATION + 1, -1, False),
RangeCase("meets_start", -10, 0, True),
RangeCase(
"overlaps_start",
-MAX_SEGMENT_DURATION + 0.5,
0.25,
True,
vod_clip_from_ms=599500,
vod_duration_ms=250,
),
RangeCase("starts_at_start", 0, 10, True, vod_duration_ms=10000),
RangeCase("inside", 20, 80, True, vod_duration_ms=60000),
RangeCase("ends_at_end", 90, 100, True, vod_duration_ms=10000),
RangeCase("matches_range", 0, 100, True, vod_duration_ms=100000),
RangeCase("starts_with_range", 0, 110, True, vod_duration_ms=100000),
RangeCase(
"covers_range",
-20,
120,
True,
vod_clip_from_ms=20000,
vod_duration_ms=100000,
),
RangeCase(
"ends_with_range",
-10,
100,
True,
vod_clip_from_ms=10000,
vod_duration_ms=100000,
),
RangeCase("overlaps_end", 95, 105, True, vod_duration_ms=5000),
RangeCase("starts_at_end", 100, 110, True),
RangeCase("after", 101, 110, False),
)
class TestHttpMedia(BaseTestHttp):
"""Test media API endpoints, particularly recordings with DST handling."""
@@ -44,6 +102,26 @@ class TestHttpMedia(BaseTestHttp):
self.app.dependency_overrides.clear()
super().tearDown()
def _assert_vod_response(
self,
response,
expected_clips: list[tuple[str, int | None, int]],
) -> None:
"""Assert VOD clip metadata and its derived duration fields."""
assert response.status_code == 200
vod = response.json()
assert [
(
clip["path"],
clip.get("clipFrom"),
clip["keyFrameDurations"][0],
)
for clip in vod["sequences"][0]["clips"]
] == expected_clips
expected_durations = [clip[2] for clip in expected_clips]
assert vod["durations"] == expected_durations
assert vod["segment_duration"] == max(expected_durations)
def test_recordings_summary_across_dst_spring_forward(self):
"""
Test recordings summary across spring DST transition (spring forward).
@@ -404,6 +482,102 @@ class TestHttpMedia(BaseTestHttp):
assert "2024-03-10" in summary
assert summary["2024-03-10"] is True
def test_recordings_handles_all_range_relations(self):
"""Recordings return every interval relation that touches the range."""
with AuthTestClient(self.app) as client:
for case in RANGE_CASES:
with self.subTest(case=case.name):
Recordings.delete().execute()
super().insert_mock_recording(
case.name,
REQUEST_START + case.start_offset,
REQUEST_START + case.end_offset,
)
response = client.get(
"/front_door/recordings",
params={"after": REQUEST_START, "before": REQUEST_END},
)
assert response.status_code == 200
expected_ids = [case.name] if case.included_in_recordings else []
assert [
recording["id"] for recording in response.json()
] == expected_ids
def test_vod_handles_all_range_relations(self):
"""VOD clips every interval relation with positive playback duration."""
with (
AuthTestClient(self.app) as client,
patch(
"frigate.api.media.get_keyframe_before",
side_effect=lambda _path, offset: offset,
),
):
for case in RANGE_CASES:
with self.subTest(case=case.name):
Recordings.delete().execute()
super().insert_mock_recording(
case.name,
REQUEST_START + case.start_offset,
REQUEST_START + case.end_offset,
)
response = client.get(
f"/vod/front_door/start/{REQUEST_START}/end/{REQUEST_END}"
)
if case.vod_duration_ms is None:
assert response.status_code == 404
continue
self._assert_vod_response(
response,
[
(
case.name,
case.vod_clip_from_ms,
case.vod_duration_ms,
)
],
)
def test_vod_handles_segment_ending_at_start_with_keyframe_fallbacks(self):
"""VOD keeps a boundary segment when keyframe lookup extends it."""
def keyframe_before(path: str, offset: int) -> int | None:
return offset - 1000 if path == "previous_keyframe" else None
with (
AuthTestClient(self.app) as client,
patch(
"frigate.api.media.get_keyframe_before",
side_effect=keyframe_before,
),
):
super().insert_mock_recording(
"previous_keyframe",
REQUEST_START - 10,
REQUEST_START,
)
super().insert_mock_recording(
"missing_keyframe",
REQUEST_START - 5,
REQUEST_START,
)
response = client.get(
f"/vod/front_door/start/{REQUEST_START}/end/{REQUEST_END}"
)
self._assert_vod_response(
response,
[
("previous_keyframe", 9000, 1000),
("missing_keyframe", None, 5000),
],
)
def test_recordings_unavailable_reports_gap_between_recordings(self):
"""A gap between two recordings is reported as an unavailable segment."""
with AuthTestClient(self.app) as client:
+88 -1
View File
@@ -1,7 +1,12 @@
import unittest
from io import StringIO
from unittest.mock import MagicMock, patch
from frigate.util.services import get_amd_gpu_stats, get_intel_gpu_stats
from frigate.util.services import (
get_amd_gpu_stats,
get_intel_gpu_stats,
get_openvino_npu_stats,
)
class TestGpuStats(unittest.TestCase):
@@ -17,6 +22,88 @@ class TestGpuStats(unittest.TestCase):
amd_stats = get_amd_gpu_stats()
assert amd_stats == {"gpu": "4.17%", "mem": "60.37%"}
@patch("frigate.util.services.time.sleep")
@patch("frigate.util.services.time.time", side_effect=[0.0, 1.0])
@patch(
"frigate.util.services.os.readlink",
return_value="/sys/bus/pci/drivers/intel_vpu",
)
@patch(
"frigate.util.services.glob.glob",
return_value=["/sys/class/accel/accel0"],
)
@patch(
"builtins.open",
side_effect=[StringIO("1000"), StringIO("1250")],
)
def test_openvino_npu_stats_discovers_accel0(
self, open_file, glob, readlink, time, sleep
):
assert get_openvino_npu_stats() == {"npu": "25.0", "mem": "-%"}
open_file.assert_any_call(
"/sys/class/accel/accel0/device/power/runtime_active_time"
)
@patch("frigate.util.services.time.sleep")
@patch("frigate.util.services.time.time", side_effect=[0.0, 1.0])
@patch(
"frigate.util.services.os.readlink",
side_effect=[
"/sys/bus/pci/drivers/other",
"/sys/bus/pci/drivers/intel_vpu",
],
)
@patch(
"frigate.util.services.glob.glob",
return_value=[
"/sys/class/accel/accel0",
"/sys/class/accel/accel1",
],
)
@patch(
"builtins.open",
side_effect=[StringIO("1000"), StringIO("1250")],
)
def test_openvino_npu_stats_skips_non_intel_accelerator(
self, open_file, glob, readlink, time, sleep
):
assert get_openvino_npu_stats() == {"npu": "25.0", "mem": "-%"}
open_file.assert_any_call(
"/sys/class/accel/accel1/device/power/runtime_active_time"
)
@patch(
"frigate.util.services.os.readlink",
return_value="/sys/bus/pci/drivers/other",
)
@patch(
"frigate.util.services.glob.glob",
return_value=["/sys/class/accel/accel0"],
)
@patch("builtins.open")
def test_openvino_npu_stats_no_intel_accelerator(self, open_file, glob, readlink):
assert get_openvino_npu_stats() is None
open_file.assert_not_called()
@patch(
"frigate.util.services.os.readlink",
return_value="/sys/bus/pci/drivers/intel_vpu",
)
@patch(
"frigate.util.services.glob.glob",
return_value=["/sys/class/accel/accel0"],
)
@patch("builtins.open", side_effect=FileNotFoundError)
def test_openvino_npu_stats_runtime_counter_unavailable(
self, open_file, glob, readlink
):
assert get_openvino_npu_stats() is None
open_file.assert_called_once_with(
"/sys/class/accel/accel0/device/power/runtime_active_time"
)
@patch("frigate.stats.intel_gpu_info.intel_gpu_name_resolver.get_names")
@patch("frigate.util.services.time.sleep")
@patch("frigate.util.services.time.monotonic")
+20 -5
View File
@@ -1,6 +1,7 @@
"""Utilities for services."""
import asyncio
import glob
import json
import logging
import os
@@ -670,19 +671,33 @@ def get_intel_gpu_stats(
def get_openvino_npu_stats() -> dict[str, str] | None:
"""Get NPU stats using openvino."""
NPU_RUNTIME_PATH = "/sys/devices/pci0000:00/0000:00:0b.0/power/runtime_active_time"
for accel_path in sorted(glob.glob("/sys/class/accel/accel*")):
try:
driver = os.path.basename(os.readlink(f"{accel_path}/device/driver"))
except OSError:
continue
if driver != "intel_vpu":
continue
try:
runtime_path = f"{accel_path}/device/power/runtime_active_time"
with open(runtime_path) as f:
initial_runtime = float(f.read().strip())
break
except (FileNotFoundError, PermissionError, ValueError):
continue
else:
return None
try:
with open(NPU_RUNTIME_PATH) as f:
initial_runtime = float(f.read().strip())
initial_time = time.time()
# Sleep for 1 second to get an accurate reading
time.sleep(1.0)
# Read runtime value again
with open(NPU_RUNTIME_PATH) as f:
with open(runtime_path) as f:
current_runtime = float(f.read().strip())
current_time = time.time()