Files
frigate/frigate/camera/maintainer.py
T

226 lines
8.5 KiB
Python
Raw Normal View History

2025-06-11 11:25:30 -06:00
"""Create and maintain camera processes / management."""
import logging
2025-06-12 12:12:34 -06:00
import multiprocessing as mp
2025-06-11 11:25:30 -06:00
import threading
from multiprocessing import Queue
from multiprocessing.managers import DictProxy, SyncManager
2025-06-11 11:25:30 -06:00
from multiprocessing.synchronize import Event as MpEvent
from frigate.camera import CameraMetrics, PTZMetrics
from frigate.config import FrigateConfig
from frigate.config.camera import CameraConfig
from frigate.config.camera.updater import (
CameraConfigUpdateEnum,
CameraConfigUpdateSubscriber,
)
from frigate.models import Regions
from frigate.util.builtin import empty_and_close_queue
from frigate.util.image import SharedMemoryFrameManager, UntrackedSharedMemory
from frigate.util.object import get_camera_regions_grid
from frigate.util.services import calculate_shm_requirements
2025-06-12 12:12:34 -06:00
from frigate.video import CameraCapture, CameraTracker
2025-06-11 11:25:30 -06:00
logger = logging.getLogger(__name__)
class CameraMaintainer(threading.Thread):
def __init__(
self,
config: FrigateConfig,
detection_queue: Queue,
detected_frames_queue: Queue,
2025-06-12 12:12:34 -06:00
camera_metrics: DictProxy,
2025-06-11 11:25:30 -06:00
ptz_metrics: dict[str, PTZMetrics],
stop_event: MpEvent,
metrics_manager: SyncManager,
2025-06-11 11:25:30 -06:00
):
super().__init__(name="camera_processor")
self.config = config
self.detection_queue = detection_queue
self.detected_frames_queue = detected_frames_queue
self.stop_event = stop_event
self.camera_metrics = camera_metrics
self.ptz_metrics = ptz_metrics
self.frame_manager = SharedMemoryFrameManager()
self.region_grids: dict[str, list[list[dict[str, int]]]] = {}
self.update_subscriber = CameraConfigUpdateSubscriber(
self.config,
{},
[
CameraConfigUpdateEnum.add,
CameraConfigUpdateEnum.remove,
],
)
self.shm_count = self.__calculate_shm_frame_count()
2025-06-12 12:12:34 -06:00
self.camera_processes: dict[str, mp.Process] = {}
self.capture_processes: dict[str, mp.Process] = {}
self.metrics_manager = metrics_manager
2025-06-11 11:25:30 -06:00
def __init_historical_regions(self) -> None:
# delete region grids for removed or renamed cameras
cameras = list(self.config.cameras.keys())
Regions.delete().where(~(Regions.camera << cameras)).execute()
# create or update region grids for each camera
for camera in self.config.cameras.values():
assert camera.name is not None
self.region_grids[camera.name] = get_camera_regions_grid(
camera.name,
camera.detect,
max(self.config.model.width, self.config.model.height),
)
def __calculate_shm_frame_count(self) -> int:
shm_stats = calculate_shm_requirements(self.config)
2025-06-11 11:25:30 -06:00
if not shm_stats:
# /dev/shm not available
2025-06-11 11:25:30 -06:00
return 0
logger.debug(
f"Calculated total camera size {shm_stats['available']} / "
f"{shm_stats['camera_frame_size']} :: {shm_stats['shm_frame_count']} "
f"frames for each camera in SHM"
2025-06-11 11:25:30 -06:00
)
if shm_stats["shm_frame_count"] < 20:
2025-06-11 11:25:30 -06:00
logger.warning(
f"The current SHM size of {shm_stats['total']}MB is too small, "
f"recommend increasing it to at least {shm_stats['min_shm']}MB."
2025-06-11 11:25:30 -06:00
)
return shm_stats["shm_frame_count"]
2025-06-11 11:25:30 -06:00
def __start_camera_processor(
self, name: str, config: CameraConfig, runtime: bool = False
) -> None:
if not config.enabled_in_config:
logger.info(f"Camera processor not started for disabled camera {name}")
return
if runtime:
self.camera_metrics[name] = CameraMetrics(self.metrics_manager)
2025-06-11 11:25:30 -06:00
self.ptz_metrics[name] = PTZMetrics(autotracker_enabled=False)
self.region_grids[name] = get_camera_regions_grid(
name,
config.detect,
max(self.config.model.width, self.config.model.height),
)
try:
largest_frame = max(
[
det.model.height * det.model.width * 3
if det.model is not None
else 320
for det in self.config.detectors.values()
]
)
UntrackedSharedMemory(name=f"out-{name}", create=True, size=20 * 6 * 4)
UntrackedSharedMemory(
name=name,
create=True,
size=largest_frame,
)
except FileExistsError:
pass
2025-06-12 12:12:34 -06:00
camera_process = CameraTracker(
config,
self.config.model,
self.config.model.merged_labelmap,
self.detection_queue,
self.detected_frames_queue,
self.camera_metrics[name],
self.ptz_metrics[name],
self.region_grids[name],
2025-06-24 11:41:11 -06:00
self.stop_event,
2025-11-17 08:12:05 -06:00
self.config.logger,
2025-06-11 11:25:30 -06:00
)
2025-06-12 12:12:34 -06:00
self.camera_processes[config.name] = camera_process
2025-06-11 11:25:30 -06:00
camera_process.start()
2025-06-12 12:12:34 -06:00
self.camera_metrics[config.name].process_pid.value = camera_process.pid
2025-06-11 11:25:30 -06:00
logger.info(f"Camera processor started for {config.name}: {camera_process.pid}")
def __start_camera_capture(
self, name: str, config: CameraConfig, runtime: bool = False
) -> None:
if not config.enabled_in_config:
logger.info(f"Capture process not started for disabled camera {name}")
return
# pre-create shms
2025-06-12 12:12:34 -06:00
count = 10 if runtime else self.shm_count
for i in range(count):
2025-06-11 11:25:30 -06:00
frame_size = config.frame_shape_yuv[0] * config.frame_shape_yuv[1]
self.frame_manager.create(f"{config.name}_frame{i}", frame_size)
2025-06-24 11:41:11 -06:00
capture_process = CameraCapture(
2025-11-17 08:12:05 -06:00
config,
count,
self.camera_metrics[name],
self.stop_event,
self.config.logger,
2025-06-24 11:41:11 -06:00
)
2025-06-11 11:25:30 -06:00
capture_process.daemon = True
2025-06-12 12:12:34 -06:00
self.capture_processes[name] = capture_process
2025-06-11 11:25:30 -06:00
capture_process.start()
2025-06-12 12:12:34 -06:00
self.camera_metrics[name].capture_process_pid.value = capture_process.pid
2025-06-11 11:25:30 -06:00
logger.info(f"Capture process started for {name}: {capture_process.pid}")
def __stop_camera_capture_process(self, camera: str) -> None:
2025-06-12 12:12:34 -06:00
capture_process = self.capture_processes[camera]
2025-06-11 11:25:30 -06:00
if capture_process is not None:
logger.info(f"Waiting for capture process for {camera} to stop")
capture_process.terminate()
capture_process.join()
def __stop_camera_process(self, camera: str) -> None:
2025-06-12 12:12:34 -06:00
camera_process = self.camera_processes[camera]
2025-06-11 11:25:30 -06:00
if camera_process is not None:
logger.info(f"Waiting for process for {camera} to stop")
camera_process.terminate()
camera_process.join()
logger.info(f"Closing frame queue for {camera}")
2025-06-12 12:12:34 -06:00
empty_and_close_queue(self.camera_metrics[camera].frame_queue)
2025-06-11 11:25:30 -06:00
def run(self):
self.__init_historical_regions()
# start camera processes
for camera, config in self.config.cameras.items():
self.__start_camera_processor(camera, config)
self.__start_camera_capture(camera, config)
while not self.stop_event.wait(1):
updates = self.update_subscriber.check_for_updates()
for update_type, updated_cameras in updates.items():
if update_type == CameraConfigUpdateEnum.add.name:
for camera in updated_cameras:
self.__start_camera_processor(
camera,
self.update_subscriber.camera_configs[camera],
runtime=True,
)
self.__start_camera_capture(
2025-06-12 12:12:34 -06:00
camera,
self.update_subscriber.camera_configs[camera],
runtime=True,
2025-06-11 11:25:30 -06:00
)
elif update_type == CameraConfigUpdateEnum.remove.name:
self.__stop_camera_capture_process(camera)
self.__stop_camera_process(camera)
# ensure the capture processes are done
2025-06-12 12:12:34 -06:00
for camera in self.camera_processes.keys():
2025-06-11 11:25:30 -06:00
self.__stop_camera_capture_process(camera)
# ensure the camera processors are done
2025-06-12 12:12:34 -06:00
for camera in self.capture_processes.keys():
2025-06-11 11:25:30 -06:00
self.__stop_camera_process(camera)
self.update_subscriber.stop()
self.frame_manager.cleanup()