Files
frigate/frigate/data_processing/post/review_descriptions.py
T

615 lines
22 KiB
Python
Raw Normal View History

"""Post processor for review items to get descriptions."""
2025-08-10 05:57:54 -06:00
import copy
import datetime
import logging
import math
2025-08-10 05:57:54 -06:00
import os
import shutil
import threading
from pathlib import Path
2025-08-12 16:27:35 -06:00
from typing import Any
2025-08-10 05:57:54 -06:00
import cv2
from peewee import DoesNotExist
2025-11-26 06:23:51 -07:00
from titlecase import titlecase
2025-08-10 05:57:54 -06:00
2025-08-12 16:27:35 -06:00
from frigate.comms.embeddings_updater import EmbeddingsRequestEnum
2025-08-10 05:57:54 -06:00
from frigate.comms.inter_process import InterProcessRequestor
from frigate.config import FrigateConfig
2025-11-10 11:03:56 -06:00
from frigate.config.camera import CameraConfig
from frigate.config.camera.review import GenAIReviewConfig, ImageSourceEnum
2026-04-07 08:16:19 -05:00
from frigate.const import (
ATTRIBUTE_LABEL_DISPLAY_MAP,
CACHE_DIR,
CLIPS_DIR,
UPDATE_REVIEW_DESCRIPTION,
)
from frigate.data_processing.types import PostProcessDataEnum
2025-08-10 05:57:54 -06:00
from frigate.genai import GenAIClient
2026-04-03 17:13:52 -06:00
from frigate.genai.manager import GenAIClientManager
from frigate.models import Recordings, ReviewSegment
2025-08-10 05:57:54 -06:00
from frigate.util.builtin import EventsPerSecond, InferenceSpeed
from frigate.util.image import get_image_from_recording
from ..post.api import PostProcessorApi
2025-08-10 05:57:54 -06:00
from ..types import DataProcessorMetrics
logger = logging.getLogger(__name__)
2025-10-28 07:28:36 -06:00
RECORDING_BUFFER_EXTENSION_PERCENT = 0.10
2025-11-10 11:03:56 -06:00
MIN_RECORDING_DURATION = 10
2026-04-25 16:38:18 -06:00
MAX_IMAGE_TOKENS = 24000
2026-04-26 16:09:35 -06:00
MAX_FRAMES_PER_SECOND = 1
class ReviewDescriptionProcessor(PostProcessorApi):
2025-08-10 05:57:54 -06:00
def __init__(
self,
config: FrigateConfig,
requestor: InterProcessRequestor,
metrics: DataProcessorMetrics,
2026-04-03 17:13:52 -06:00
genai_manager: GenAIClientManager,
2025-08-10 05:57:54 -06:00
):
super().__init__(config, metrics, None)
2025-08-10 05:57:54 -06:00
self.requestor = requestor
self.metrics = metrics
2026-04-03 17:13:52 -06:00
self.genai_manager = genai_manager
2025-08-10 05:57:54 -06:00
self.review_desc_speed = InferenceSpeed(self.metrics.review_desc_speed)
2026-03-26 12:54:12 -06:00
self.review_desc_dps = EventsPerSecond()
self.review_desc_dps.start()
def calculate_frame_count(
2025-10-30 08:52:55 -06:00
self,
camera: str,
2026-04-25 16:38:18 -06:00
duration: float,
2025-10-30 08:52:55 -06:00
image_source: ImageSourceEnum = ImageSourceEnum.preview,
height: int = 480,
) -> int:
2026-04-25 16:38:18 -06:00
"""Calculate optimal number of frames based on event duration, context size,
image source, and resolution.
2025-10-30 08:52:55 -06:00
2026-04-25 16:38:18 -06:00
Per-image token cost is asked of the GenAI provider so providers that know
their model's true cost (e.g. llama.cpp can probe the loaded mmproj) can
diverge from the default ~1-token-per-1250-pixels heuristic. The frame
budget is bounded by:
- remaining context window after prompt + response reservations
- a fixed MAX_IMAGE_TOKENS ceiling
- MAX_FRAMES_PER_SECOND x duration, to avoid drowning short events in
near-duplicate frames where the model latches onto the redundant middle
and skips the start/end action
2025-10-30 08:52:55 -06:00
"""
2026-04-03 17:13:52 -06:00
client = self.genai_manager.description_client
if client is None:
return 3
context_size = client.get_context_size()
2025-10-30 08:52:55 -06:00
camera_config = self.config.cameras[camera]
detect_width = camera_config.detect.width
detect_height = camera_config.detect.height
2026-03-26 12:54:12 -06:00
if not detect_width or not detect_height:
aspect_ratio = 16 / 9
else:
aspect_ratio = detect_width / detect_height
2025-10-02 09:17:25 -06:00
if image_source == ImageSourceEnum.recordings:
2025-10-30 08:52:55 -06:00
if aspect_ratio >= 1:
# Landscape or square: constrain height
width = int(height * aspect_ratio)
else:
2025-10-30 08:52:55 -06:00
# Portrait: constrain width
width = height
height = int(width / aspect_ratio)
2025-10-02 09:17:25 -06:00
else:
2025-10-30 08:52:55 -06:00
if aspect_ratio >= 1:
# Landscape or square: constrain height
target_height = 180
width = int(target_height * aspect_ratio)
height = target_height
else:
2025-10-30 08:52:55 -06:00
# Portrait: constrain width
target_width = 180
width = target_width
height = int(target_width / aspect_ratio)
2026-04-25 16:38:18 -06:00
tokens_per_image = client.estimate_image_tokens(width, height)
2025-12-26 07:45:03 -07:00
prompt_tokens = 3800
2025-11-08 13:13:40 -07:00
response_tokens = 300
2026-04-25 16:38:18 -06:00
context_budget = context_size - prompt_tokens - response_tokens
image_token_budget = min(context_budget, MAX_IMAGE_TOKENS)
max_frames_by_tokens = int(image_token_budget / tokens_per_image)
max_frames_by_duration = int(duration * MAX_FRAMES_PER_SECOND)
max_frames = min(max_frames_by_tokens, max_frames_by_duration)
return max(max_frames, 3)
2025-10-02 09:17:25 -06:00
2026-03-26 12:54:12 -06:00
def process_data(
self, data: dict[str, Any], data_type: PostProcessDataEnum
) -> None:
self.metrics.review_desc_dps.value = self.review_desc_dps.eps()
2025-08-10 05:57:54 -06:00
if data_type != PostProcessDataEnum.review:
return
2026-04-03 17:13:52 -06:00
if self.genai_manager.description_client is None:
return
2025-08-10 07:38:04 -06:00
camera = data["after"]["camera"]
camera_config = self.config.cameras[camera]
2025-08-10 07:38:04 -06:00
if not camera_config.review.genai.enabled:
2025-08-10 07:38:04 -06:00
return
2025-08-10 05:57:54 -06:00
id = data["after"]["id"]
if data["type"] == "new" or data["type"] == "update":
return
else:
final_data = data["after"]
2025-08-10 05:57:54 -06:00
if (
final_data["severity"] == "alert"
and not camera_config.review.genai.alerts
):
return
elif (
final_data["severity"] == "detection"
and not camera_config.review.genai.detections
):
return
2025-08-10 05:57:54 -06:00
image_source = camera_config.review.genai.image_source
2025-08-10 05:57:54 -06:00
if image_source == ImageSourceEnum.recordings:
2025-10-28 07:28:36 -06:00
duration = final_data["end_time"] - final_data["start_time"]
2025-11-17 08:12:05 -06:00
buffer_extension = min(5, duration * RECORDING_BUFFER_EXTENSION_PERCENT)
2025-11-10 11:03:56 -06:00
# Ensure minimum total duration for short review items
# This provides better context for brief events
total_duration = duration + (2 * buffer_extension)
if total_duration < MIN_RECORDING_DURATION:
2025-11-17 08:12:05 -06:00
# Expand buffer to reach minimum duration, still respecting max of 5s per side
2025-11-10 11:03:56 -06:00
additional_buffer_per_side = (MIN_RECORDING_DURATION - duration) / 2
2025-11-17 08:12:05 -06:00
buffer_extension = min(5, additional_buffer_per_side)
2025-10-28 07:28:36 -06:00
final_data["start_time"] -= buffer_extension
final_data["end_time"] += buffer_extension
thumbs = self.get_recording_frames(
camera,
final_data["start_time"],
final_data["end_time"],
height=480, # Use 480p for good balance between quality and token usage
2025-08-10 05:57:54 -06:00
)
if not thumbs:
# Fallback to preview frames if no recordings available
logger.warning(
f"No recording frames found for {camera}, falling back to preview frames"
)
thumbs = self.get_preview_frames_as_bytes(
camera,
final_data["start_time"],
final_data["end_time"],
final_data["thumb_path"],
id,
camera_config.review.genai.debug_save_thumbnails,
)
elif camera_config.review.genai.debug_save_thumbnails:
# Save debug thumbnails for recordings
Path(os.path.join(CLIPS_DIR, "genai-requests", id)).mkdir(
2025-08-10 05:57:54 -06:00
parents=True, exist_ok=True
)
for idx, frame_bytes in enumerate(thumbs):
with open(
os.path.join(CLIPS_DIR, f"genai-requests/{id}/{idx}.jpg"),
"wb",
) as f:
f.write(frame_bytes)
else:
# Use preview frames
thumbs = self.get_preview_frames_as_bytes(
camera,
final_data["start_time"],
final_data["end_time"],
final_data["thumb_path"],
id,
camera_config.review.genai.debug_save_thumbnails,
)
2025-08-10 05:57:54 -06:00
# kickoff analysis
2026-03-26 12:54:12 -06:00
self.review_desc_dps.update()
2025-08-10 05:57:54 -06:00
threading.Thread(
target=run_analysis,
args=(
self.requestor,
2026-04-03 17:13:52 -06:00
self.genai_manager.description_client,
2025-08-10 05:57:54 -06:00
self.review_desc_speed,
2025-11-10 11:03:56 -06:00
camera_config,
2025-08-10 05:57:54 -06:00
final_data,
thumbs,
2025-08-12 16:27:35 -06:00
camera_config.review.genai,
2025-08-15 07:25:49 -06:00
list(self.config.model.merged_labelmap.values()),
self.config.model.all_attributes,
2025-08-10 05:57:54 -06:00
),
).start()
2026-03-26 12:54:12 -06:00
def handle_request(self, topic: str, request_data: dict[str, Any]) -> str | None:
2025-08-12 16:27:35 -06:00
if topic == EmbeddingsRequestEnum.summarize_review.value:
start_ts = request_data["start_ts"]
end_ts = request_data["end_ts"]
2025-09-25 20:05:22 -06:00
logger.debug(
f"Found GenAI Review Summary request for {start_ts} to {end_ts}"
)
2025-12-04 12:19:07 -06:00
# Query all review segments with camera and time information
segments: list[dict[str, Any]] = [
{
"camera": r["camera"].replace("_", " ").title(),
"start_time": r["start_time"],
"end_time": r["end_time"],
"metadata": r["data"]["metadata"],
}
2025-08-12 16:27:35 -06:00
for r in (
2025-12-04 12:19:07 -06:00
ReviewSegment.select(
ReviewSegment.camera,
ReviewSegment.start_time,
ReviewSegment.end_time,
ReviewSegment.data,
)
2025-08-12 16:27:35 -06:00
.where(
(ReviewSegment.data["metadata"].is_null(False))
& (ReviewSegment.start_time < end_ts)
& (ReviewSegment.end_time > start_ts)
)
.order_by(ReviewSegment.start_time.asc())
.dicts()
.iterator()
)
]
2025-12-04 12:19:07 -06:00
if len(segments) == 0:
2025-08-12 16:27:35 -06:00
logger.debug("No review items with metadata found during time period")
2025-12-04 12:19:07 -06:00
return "No activity was found during this time period."
2025-08-12 16:27:35 -06:00
2025-12-04 12:19:07 -06:00
# Identify primary items (important items that need review)
primary_segments = [
seg
for seg in segments
if seg["metadata"].get("potential_threat_level", 0) > 0
or seg["metadata"].get("other_concerns")
]
2025-08-12 16:27:35 -06:00
2025-12-04 12:19:07 -06:00
if not primary_segments:
2025-08-12 16:27:35 -06:00
return "No concerns were found during this time period."
2025-12-11 08:23:34 -06:00
# Build hierarchical structure: each primary event with its contextual items
events_with_context = []
2025-12-04 12:19:07 -06:00
for primary_seg in primary_segments:
2025-12-11 08:23:34 -06:00
# Start building the primary event structure
2025-12-04 12:19:07 -06:00
primary_item = copy.deepcopy(primary_seg["metadata"])
2025-12-11 08:23:34 -06:00
primary_item["camera"] = primary_seg["camera"]
primary_item["start_time"] = primary_seg["start_time"]
primary_item["end_time"] = primary_seg["end_time"]
2025-12-04 12:19:07 -06:00
# Find overlapping contextual items from other cameras
primary_start = primary_seg["start_time"]
primary_end = primary_seg["end_time"]
primary_camera = primary_seg["camera"]
2025-12-11 08:23:34 -06:00
contextual_items = []
seen_contextual_cameras = set()
2025-12-04 12:19:07 -06:00
for seg in segments:
seg_camera = seg["camera"]
if seg_camera == primary_camera:
continue
if seg in primary_segments:
continue
seg_start = seg["start_time"]
seg_end = seg["end_time"]
if seg_start < primary_end and primary_start < seg_end:
2025-12-11 08:23:34 -06:00
# Avoid duplicates if same camera has multiple overlapping segments
if seg_camera not in seen_contextual_cameras:
contextual_item = copy.deepcopy(seg["metadata"])
contextual_item["camera"] = seg_camera
contextual_item["start_time"] = seg_start
contextual_item["end_time"] = seg_end
contextual_items.append(contextual_item)
seen_contextual_cameras.add(seg_camera)
2025-12-04 12:19:07 -06:00
2025-12-11 08:23:34 -06:00
# Add context array to primary item
primary_item["context"] = contextual_items
events_with_context.append(primary_item)
2025-12-04 12:19:07 -06:00
2025-12-11 08:23:34 -06:00
total_context_items = sum(
len(event.get("context", [])) for event in events_with_context
)
2025-12-04 12:19:07 -06:00
logger.debug(
2025-12-11 08:23:34 -06:00
f"Summary includes {len(events_with_context)} primary events with "
f"{total_context_items} total contextual items"
2025-12-04 12:19:07 -06:00
)
2025-09-25 20:05:22 -06:00
if self.config.review.genai.debug_save_thumbnails:
Path(
os.path.join(CLIPS_DIR, "genai-requests", f"{start_ts}-{end_ts}")
).mkdir(parents=True, exist_ok=True)
2026-04-03 17:13:52 -06:00
client = self.genai_manager.description_client
if client is None:
return None
return client.generate_review_summary(
2025-09-25 20:05:22 -06:00
start_ts,
end_ts,
2025-12-11 08:23:34 -06:00
events_with_context,
2025-12-20 17:30:34 -07:00
self.config.review.genai.preferred_language,
2025-09-25 20:05:22 -06:00
self.config.review.genai.debug_save_thumbnails,
2025-08-12 16:27:35 -06:00
)
else:
return None
2025-08-10 05:57:54 -06:00
def get_cache_frames(
2025-08-15 07:25:49 -06:00
self,
camera: str,
start_time: float,
end_time: float,
) -> list[str]:
preview_dir = os.path.join(CACHE_DIR, "preview_frames")
2026-03-23 10:22:52 -06:00
file_start = f"preview_{camera}-"
start_file = f"{file_start}{start_time}.webp"
end_file = f"{file_start}{end_time}.webp"
2026-04-30 11:53:34 -06:00
camera_files = [
entry.name
for entry in os.scandir(preview_dir)
if entry.name.startswith(file_start)
]
camera_files.sort()
2026-03-26 12:54:12 -06:00
all_frames: list[str] = []
2026-04-30 11:53:34 -06:00
for file in camera_files:
if file < start_file:
2025-08-15 07:25:49 -06:00
if len(all_frames):
all_frames[0] = os.path.join(preview_dir, file)
else:
all_frames.append(os.path.join(preview_dir, file))
continue
if file > end_file:
2025-08-15 07:25:49 -06:00
all_frames.append(os.path.join(preview_dir, file))
break
all_frames.append(os.path.join(preview_dir, file))
frame_count = len(all_frames)
2026-04-25 16:38:18 -06:00
desired_frame_count = self.calculate_frame_count(
camera, duration=end_time - start_time
)
2025-10-02 09:17:25 -06:00
2025-08-15 07:25:49 -06:00
if frame_count <= desired_frame_count:
return all_frames
selected_frames = []
2025-08-15 07:25:49 -06:00
step_size = (frame_count - 1) / (desired_frame_count - 1)
2025-08-15 07:25:49 -06:00
for i in range(desired_frame_count):
index = round(i * step_size)
selected_frames.append(all_frames[index])
return selected_frames
def get_recording_frames(
self,
camera: str,
start_time: float,
end_time: float,
height: int = 480,
) -> list[bytes]:
"""Get frames from recordings at specified timestamps."""
duration = end_time - start_time
2025-10-30 08:52:55 -06:00
desired_frame_count = self.calculate_frame_count(
2026-04-25 16:38:18 -06:00
camera, duration, ImageSourceEnum.recordings, height
2025-10-30 08:52:55 -06:00
)
# Calculate evenly spaced timestamps throughout the duration
if desired_frame_count == 1:
timestamps = [start_time + duration / 2]
else:
step = duration / (desired_frame_count - 1)
timestamps = [start_time + (i * step) for i in range(desired_frame_count)]
def extract_frame_from_recording(ts: float) -> bytes | None:
"""Extract a single frame from recording at given timestamp."""
try:
recording = (
Recordings.select(
Recordings.path,
Recordings.start_time,
)
.where((ts >= Recordings.start_time) & (ts <= Recordings.end_time))
.where(Recordings.camera == camera)
.order_by(Recordings.start_time.desc())
.limit(1)
.get()
)
time_in_segment = ts - recording.start_time
return get_image_from_recording(
self.config.ffmpeg,
recording.path,
time_in_segment,
"mjpeg",
height=height,
)
except DoesNotExist:
return None
frames = []
for timestamp in timestamps:
try:
# Try to extract frame at exact timestamp
image_data = extract_frame_from_recording(timestamp)
if not image_data:
# Try with rounded timestamp as fallback
rounded_timestamp = math.ceil(timestamp)
image_data = extract_frame_from_recording(rounded_timestamp)
if image_data:
frames.append(image_data)
else:
logger.warning(
f"No recording found for {camera} at timestamp {timestamp}"
)
except Exception as e:
logger.error(
f"Error extracting frame from recording for {camera} at {timestamp}: {e}"
)
continue
return frames
def get_preview_frames_as_bytes(
self,
camera: str,
start_time: float,
end_time: float,
thumb_path_fallback: str,
review_id: str,
save_debug: bool,
) -> list[bytes]:
"""Get preview frames and convert them to JPEG bytes.
Args:
camera: Camera name
start_time: Start timestamp
end_time: End timestamp
thumb_path_fallback: Fallback thumbnail path if no preview frames found
review_id: Review item ID for debug saving
save_debug: Whether to save debug thumbnails
Returns:
List of JPEG image bytes
"""
frame_paths = self.get_cache_frames(camera, start_time, end_time)
if not frame_paths:
frame_paths = [thumb_path_fallback]
thumbs = []
for idx, thumb_path in enumerate(frame_paths):
thumb_data = cv2.imread(thumb_path)
2026-03-10 13:26:45 -06:00
if thumb_data is None:
2026-03-26 12:54:12 -06:00
logger.warning( # type: ignore[unreachable]
2026-03-10 13:26:45 -06:00
"Could not read preview frame at %s, skipping", thumb_path
)
continue
ret, jpg = cv2.imencode(
".jpg", thumb_data, [int(cv2.IMWRITE_JPEG_QUALITY), 100]
)
if ret:
thumbs.append(jpg.tobytes())
if save_debug:
Path(os.path.join(CLIPS_DIR, "genai-requests", review_id)).mkdir(
parents=True, exist_ok=True
)
shutil.copy(
thumb_path,
os.path.join(CLIPS_DIR, f"genai-requests/{review_id}/{idx}.webp"),
)
return thumbs
2025-08-10 05:57:54 -06:00
def run_analysis(
requestor: InterProcessRequestor,
genai_client: GenAIClient,
review_inference_speed: InferenceSpeed,
2025-11-10 11:03:56 -06:00
camera_config: CameraConfig,
2026-03-26 12:54:12 -06:00
final_data: dict[str, Any],
2025-08-10 05:57:54 -06:00
thumbs: list[bytes],
2025-08-12 16:27:35 -06:00
genai_config: GenAIReviewConfig,
2025-08-15 07:25:49 -06:00
labelmap_objects: list[str],
attribute_labels: list[str],
2025-08-10 05:57:54 -06:00
) -> None:
start = datetime.datetime.now().timestamp()
2025-11-10 11:03:56 -06:00
# Format zone names using zone config friendly names if available
formatted_zones = []
for zone_name in final_data["data"]["zones"]:
if zone_name in camera_config.zones:
formatted_zones.append(
camera_config.zones[zone_name].get_formatted_name(zone_name)
)
2025-08-15 07:25:49 -06:00
analytics_data = {
"id": final_data["id"],
2025-11-10 11:03:56 -06:00
"camera": camera_config.get_formatted_name(),
"zones": formatted_zones,
2025-08-15 07:25:49 -06:00
"start": datetime.datetime.fromtimestamp(final_data["start_time"]).strftime(
"%A, %I:%M %p"
),
2025-10-02 09:17:25 -06:00
"duration": round(final_data["end_time"] - final_data["start_time"]),
2025-08-15 07:25:49 -06:00
}
unified_objects = []
2025-08-15 07:25:49 -06:00
objects_list = final_data["data"]["objects"]
sub_labels_list = final_data["data"]["sub_labels"]
for i, verified_label in enumerate(final_data["data"]["verified_objects"]):
object_type = verified_label.replace("-verified", "").replace("_", " ")
2025-11-26 06:23:51 -07:00
name = titlecase(sub_labels_list[i].replace("_", " "))
2026-03-19 10:39:24 -06:00
unified_objects.append(f"{name}{object_type}")
for label in objects_list:
2025-08-15 07:25:49 -06:00
if "-verified" in label:
continue
elif label in labelmap_objects:
2026-04-07 08:16:19 -05:00
object_type = label.replace("_", " ")
if label in attribute_labels:
2026-04-07 08:16:19 -05:00
display_name = ATTRIBUTE_LABEL_DISPLAY_MAP.get(label, object_type)
unified_objects.append(f"{display_name} (delivery/service)")
else:
unified_objects.append(object_type)
analytics_data["unified_objects"] = unified_objects
2025-08-15 07:25:49 -06:00
2025-08-10 05:57:54 -06:00
metadata = genai_client.generate_review_description(
2025-08-15 07:25:49 -06:00
analytics_data,
2025-08-10 05:57:54 -06:00
thumbs,
2025-08-12 16:27:35 -06:00
genai_config.additional_concerns,
genai_config.preferred_language,
genai_config.debug_save_thumbnails,
2025-09-30 17:07:16 -06:00
genai_config.activity_context_prompt,
2025-08-10 05:57:54 -06:00
)
review_inference_speed.update(datetime.datetime.now().timestamp() - start)
if not metadata:
return None
prev_data = copy.deepcopy(final_data)
final_data["data"]["metadata"] = metadata.model_dump()
requestor.send_data(
UPDATE_REVIEW_DESCRIPTION,
{
"type": "genai",
"before": {k: v for k, v in prev_data.items()},
"after": {k: v for k, v in final_data.items()},
},
)