Create a notice when frames are backed up in the tracked object process (#24619)

* Create a notice when frames are backed up in the tracked object processor

* Launch notice in thread to avoid more delay
This commit is contained in:
Nicolas Mowen
2026-10-10 11:59:15 -05:00
committed by GitHub
parent 01f3369659
commit e8e6326d37
4 changed files with 174 additions and 0 deletions
+9
View File
@@ -73,6 +73,15 @@ _KINDS = (
"detect_high_cpu", NoticeSeverity.warning, "camera", link="/system#cameras"
),
NoticeKind("shm_too_low", NoticeSeverity.warning, "system", link="/system#storage"),
# every camera process raises it when the shared queue backs up, so repeats
# wait for the next flush instead of writing once per camera
NoticeKind(
"object_processing_behind",
NoticeSeverity.warning,
"system",
link="/system#cameras",
batch_repeats=True,
),
# one row per user per burst; the login log lines carry the address
NoticeKind(
"failed_login",
+103
View File
@@ -0,0 +1,103 @@
"""Tests for the notice raised when object processing drops camera frames."""
import threading
import unittest
from unittest.mock import patch
from frigate.const import REPLAY_CAMERA_PREFIX
from frigate.notices.types import NOTICE_KINDS
from frigate.video.detect import (
DROPPED_FRAMES_NOTICE_COUNT,
DROPPED_FRAMES_NOTICE_INTERVAL_S,
DROPPED_FRAMES_WINDOW_S,
DroppedFrameTracker,
)
class TestDroppedFrameTracker(unittest.TestCase):
def setUp(self):
raise_patch = patch("frigate.video.detect.raise_notice")
self.raise_notice = raise_patch.start()
self.addCleanup(raise_patch.stop)
def _drop(self, tracker: DroppedFrameTracker, count: int, start: float) -> None:
for i in range(count):
tracker.dropped(start + i * 0.2)
# the notice is sent on a background thread
if tracker._sender is not None:
tracker._sender.join(timeout=5)
def test_kind_is_registered(self):
self.assertIn("object_processing_behind", NOTICE_KINDS)
def test_a_few_drops_raise_nothing(self):
tracker = DroppedFrameTracker("front_door")
self._drop(tracker, DROPPED_FRAMES_NOTICE_COUNT - 1, 0.0)
self.raise_notice.assert_not_called()
def test_enough_drops_raise_the_notice(self):
tracker = DroppedFrameTracker("front_door")
self._drop(tracker, DROPPED_FRAMES_NOTICE_COUNT, 0.0)
self.raise_notice.assert_called_once_with("object_processing_behind")
def test_drops_outside_the_window_do_not_add_up(self):
tracker = DroppedFrameTracker("front_door")
for i in range(DROPPED_FRAMES_NOTICE_COUNT):
tracker.dropped(i * (DROPPED_FRAMES_WINDOW_S + 1.0))
self.raise_notice.assert_not_called()
def test_repeats_wait_for_the_interval(self):
tracker = DroppedFrameTracker("front_door")
self._drop(tracker, DROPPED_FRAMES_NOTICE_COUNT, 0.0)
self._drop(tracker, DROPPED_FRAMES_NOTICE_COUNT * 2, 5.0)
self.assertEqual(self.raise_notice.call_count, 1)
self._drop(
tracker, DROPPED_FRAMES_NOTICE_COUNT, DROPPED_FRAMES_NOTICE_INTERVAL_S
)
self.assertEqual(self.raise_notice.call_count, 2)
def test_a_pending_send_does_not_block_or_start_another(self):
sending = threading.Event()
release = threading.Event()
def wait_for_reply(*_) -> None:
sending.set()
release.wait(5)
self.raise_notice.side_effect = wait_for_reply
tracker = DroppedFrameTracker("front_door")
# the first notice gets no reply while the second burst arrives
for i in range(DROPPED_FRAMES_NOTICE_COUNT):
tracker.dropped(i * 0.2)
self.assertTrue(sending.wait(5))
for i in range(DROPPED_FRAMES_NOTICE_COUNT):
tracker.dropped(DROPPED_FRAMES_NOTICE_INTERVAL_S + i * 0.2)
self.assertEqual(self.raise_notice.call_count, 1)
release.set()
assert tracker._sender is not None
tracker._sender.join(timeout=5)
def test_replay_camera_never_raises(self):
tracker = DroppedFrameTracker(f"{REPLAY_CAMERA_PREFIX}front_door")
self._drop(tracker, DROPPED_FRAMES_NOTICE_COUNT * 2, 0.0)
self.raise_notice.assert_not_called()
if __name__ == "__main__":
unittest.main()
+61
View File
@@ -2,7 +2,9 @@
import logging
import queue
import threading
import time
from collections import deque
from datetime import UTC, datetime
from multiprocessing import Queue
from multiprocessing.synchronize import Event as MpEvent
@@ -20,10 +22,12 @@ from frigate.config.camera.updater import (
)
from frigate.const import (
PROCESS_PRIORITY_HIGH,
REPLAY_CAMERA_PREFIX,
REQUEST_REGION_GRID,
)
from frigate.motion import MotionDetector
from frigate.motion.improved_motion import ImprovedMotionDetector
from frigate.notices import raise_notice
from frigate.object_detection.base import RemoteObjectDetector
from frigate.ptz.autotrack import ptz_moving_at_frame_time
from frigate.track import ObjectTracker
@@ -52,6 +56,61 @@ from frigate.util.time import get_tomorrow_at_time
logger = logging.getLogger(__name__)
# this many frames dropped within the window because the shared detected frames
# queue was full means tracked object processing is falling behind
DROPPED_FRAMES_NOTICE_COUNT = 5
DROPPED_FRAMES_WINDOW_S = 30
DROPPED_FRAMES_NOTICE_INTERVAL_S = 60
class DroppedFrameTracker:
"""Raises a notice when a camera drops frames because object processing is behind."""
def __init__(self, camera: str) -> None:
# a debug replay can feed frames faster than real time
self._enabled = not camera.startswith(REPLAY_CAMERA_PREFIX)
self._drops: deque[float] = deque()
self._last_notice: float | None = None
# raising a notice waits on the main process, which answers every
# process's requests one at a time, so it runs off the frame loop
self._sender: threading.Thread | None = None
def dropped(self, now: float) -> None:
"""Record a dropped frame and raise the notice if enough were dropped.
Args:
now: Monotonic time of the drop in seconds
"""
if not self._enabled:
return
self._drops.append(now)
while now - self._drops[0] > DROPPED_FRAMES_WINDOW_S:
self._drops.popleft()
if (
len(self._drops) < DROPPED_FRAMES_NOTICE_COUNT
or (
self._last_notice is not None
and now - self._last_notice < DROPPED_FRAMES_NOTICE_INTERVAL_S
)
or (self._sender is not None and self._sender.is_alive())
):
return
self._drops.clear()
self._last_notice = now
self._sender = threading.Thread(
target=raise_notice,
args=("object_processing_behind",),
name="dropped_frames_notice",
daemon=True,
)
self._sender.start()
class CameraTracker(FrigateProcess):
def __init__(
@@ -207,6 +266,7 @@ def process_frames(
fps_tracker = EventsPerSecond()
fps_tracker.start()
dropped_frames = DroppedFrameTracker(camera_config.name)
startup_scan = True
stationary_frame_counter = 0
@@ -542,6 +602,7 @@ def process_frames(
# add to the queue if not full
if detected_objects_queue.full():
frame_manager.close(frame_name)
dropped_frames.dropped(time.monotonic())
continue
else:
fps_tracker.update()
+1
View File
@@ -53,6 +53,7 @@
"ffmpeg_high_cpu": "FFmpeg CPU usage is high ({{cpu}}% average)",
"detect_high_cpu": "Detection CPU usage is high ({{cpu}}% average)",
"shm_too_low": "/dev/shm allocation ({{total}} MB) should be increased to at least {{min}} MB",
"object_processing_behind": "Tracked object processing could not keep up, so camera frames were dropped. Check CPU utilization, and whether CPU limits or pinning on the container or VM are leaving Frigate too little CPU",
"failed_login_one": "Failed login attempt for {{user}}",
"failed_login_other": "Failed login attempts for {{user}}",
"update_available": "Frigate {{version}} is available"