Compare commits

..
Author SHA1 Message Date
Nicolas Mowen c6d8f067be Improve report context management 2026-09-29 10:18:47 -06:00
Nicolas Mowen 7767f67c67 Cleanup and fixes 2026-09-29 10:13:04 -06:00
Nicolas Mowen cdda89a11e Integrate state changes with review items 2026-09-29 09:47:36 -06:00
15 changed files with 658 additions and 45 deletions
@@ -85,6 +85,14 @@ An optional config, `save_attempts`, can be set as a key under the model name. T
</TabItem>
</ConfigTabs>
## Review items
When a model's state changes while its camera has an active review item, the change is recorded on that review item. This includes changes in the few seconds before the item starts, such as a garage door opening just before the car is detected. State changes never create or extend review items on their own, and the first state reported after Frigate starts is not recorded as a change.
Recorded changes appear in the review item's data as `classification_state_changes` (see the [`frigate/reviews`](/integrations/mqtt#frigatereviews) MQTT topic) and are passed to [GenAI review summaries](/configuration/genai/genai_review) as facts, so a description can note that a gate was opened during the activity.
Change times are most accurate with `motion: true`. A model that only runs on an `interval` notices a change at its next run, so the change may be recorded late or attached to a later review item.
## Training the model
Creating and training the model is done within the Frigate UI using the `Classification` page. The process consists of three steps:
@@ -201,6 +201,8 @@ Review items are sent to the model as a sequence of still frames. Some models fo
The notes come from tracking data rather than from the images, so they describe activity the model may not have picked up on its own. In testing with a person carrying three waste bins to the curb one at a time, `gemma4` described a single trip on every attempt with `frames`, and consistently described multiple trips with `annotated_frames`. Models that already handle these sequences well, such as the `qwen3-vl` family, gain little and should stay on `frames`.
Changes reported by [state classification](/configuration/custom_classification/state_classification#review-items) models during the review item are listed in the prompt in both modes. `annotated_frames` also notes each change before the frame where it happened.
Annotated mode also caps the number of frames, since the notes already establish the order of events and extra near-duplicate frames tend to crowd out the middle of a clip. Longer review items are sampled more sparsely as a result, and typically use fewer tokens than `frames` mode for the same item.
:::note
+13 -2
View File
@@ -212,6 +212,7 @@ An `update` with the same ID will be published when:
- The severity changes from `detection` to `alert`
- Additional objects are detected
- An object is recognized via face, lpr, etc.
- A [state classification](/configuration/custom_classification/state_classification#review-items) model changes state
When the review activity has ended a final `end` message is published.
@@ -235,7 +236,8 @@ When the review activity has ended a final `end` message is published.
"objects": ["person", "car"],
"sub_labels": [],
"zones": [],
"audio": []
"audio": [],
"classification_state_changes": []
}
},
"after": {
@@ -254,7 +256,16 @@ When the review activity has ended a final `end` message is published.
"objects": ["person", "car"],
"sub_labels": ["Bob"],
"zones": ["front_yard"],
"audio": []
"audio": [],
"classification_state_changes": [
// verified changes of state classification models on this camera
{
"model": "front_gate",
"from": "closed",
"to": "open",
"timestamp": 1718987131.52
}
]
}
}
}
+1
View File
@@ -12,6 +12,7 @@ class DetectionTypeEnum(str, Enum):
video = "video"
audio = "audio"
lpr = "lpr"
classification_state = "classification_state"
class DetectionPublisher(Publisher):
@@ -1,10 +1,10 @@
"""Frame annotations derived from object tracking data.
Builds short notes describing what changed during a review item, keyed to the
frames sampled from it. Everything here comes from tracked object data already
in the database (each event's `path_data` trajectory and the timeline's
stationary/active changes), so the notes can be stated to the model as fact
rather than as something it must perceive.
frames sampled from it. Everything here comes from data already recorded
(each event's `path_data` trajectory, the timeline's stationary/active
changes, and the review item's state classification changes), so the notes
can be stated to the model as fact rather than as something it must perceive.
"""
import logging
@@ -95,6 +95,15 @@ def event_name(event: dict[str, Any]) -> str:
return f"{article} {label}"
def describe_classification_change(change: dict[str, Any]) -> str:
"""Phrase a state classification change, e.g. 'front gate changed from
closed to open'."""
model = str(change["model"]).replace("_", " ")
before = str(change["from"]).replace("_", " ")
after = str(change["to"]).replace("_", " ")
return f"{model} changed from {before} to {after}"
def path_legs(points: list[Point]) -> list[Leg]:
"""Split a trajectory into runs of travel in a consistent direction.
@@ -359,25 +368,38 @@ def get_state_changes(detection_ids: list[str]) -> list[dict[str, Any]]:
def build_frame_captions(
detection_ids: list[str],
frame_times: list[float],
classification_changes: Sequence[dict[str, Any]] = (),
) -> list[str]:
"""A caption for each sampled frame, in frame order.
Every frame gets its index and elapsed time so the model can tell them
apart; frames where something changed also carry the tracker notes for
that moment. Returns an empty list when there is nothing to say, which
callers treat as a reason to fall back to sending plain frames.
apart; frames where something changed also carry the tracker and state
classification notes for that moment. Returns an empty list when there is
nothing to say, which callers treat as a reason to fall back to sending
plain frames.
"""
if not frame_times:
return []
span_end = frame_times[-1]
timeline: list[tuple[float, str]] = []
events = get_tracked_events(detection_ids)
if not events:
logger.debug("No tracked events found for review item, skipping annotations")
return []
# audio and manual review items can have state changes but no tracked objects
if events:
timeline.extend(
(timestamp, f"[tracker] {note}")
for timestamp, note in build_timeline(
events, span_end, get_state_changes(detection_ids)
)
)
timeline = build_timeline(events, frame_times[-1], get_state_changes(detection_ids))
buckets = annotations_by_frame(timeline, frame_times)
timeline.extend(
(change["timestamp"], f"[state] {describe_classification_change(change)}")
for change in classification_changes
if change["timestamp"] <= span_end
)
buckets = annotations_by_frame(sorted(timeline, key=lambda m: m[0]), frame_times)
if not buckets:
return []
@@ -388,7 +410,7 @@ def build_frame_captions(
for index, timestamp in enumerate(frame_times):
lines = [f"Frame {index + 1} of {total} (+{timestamp - origin:.1f}s):"]
lines.extend(f"[tracker] {note}" for note in buckets.get(index, []))
lines.extend(buckets.get(index, []))
captions.append("\n".join(lines))
return captions
@@ -40,7 +40,7 @@ from frigate.util.image import get_image_from_recording
from ..post.api import PostProcessorApi
from ..types import DataProcessorMetrics
from .review_annotations import build_frame_captions
from .review_annotations import build_frame_captions, describe_classification_change
logger = logging.getLogger(__name__)
@@ -254,6 +254,10 @@ class ReviewDescriptionProcessor(PostProcessorApi):
"start_time": r["start_time"],
"end_time": r["end_time"],
"metadata": r["data"]["metadata"],
"state_changes": [
describe_classification_change(change)
for change in sorted_classification_state_changes(r["data"])
],
}
for r in (
ReviewSegment.select(
@@ -298,6 +302,9 @@ class ReviewDescriptionProcessor(PostProcessorApi):
primary_item["start_time"] = primary_seg["start_time"]
primary_item["end_time"] = primary_seg["end_time"]
if primary_seg["state_changes"]:
primary_item["state_changes"] = primary_seg["state_changes"]
# Find overlapping contextual items from other cameras
primary_start = primary_seg["start_time"]
primary_end = primary_seg["end_time"]
@@ -318,12 +325,23 @@ class ReviewDescriptionProcessor(PostProcessorApi):
seg_end = seg["end_time"]
if seg_start < primary_end and primary_start < seg_end:
# Avoid duplicates if same camera has multiple overlapping segments
if seg_camera not in seen_contextual_cameras:
# Avoid duplicates if same camera has multiple overlapping
# segments. One with state changes is kept as its own item
# so each change stays within its item's time range.
if (
seg_camera in seen_contextual_cameras
and not seg["state_changes"]
):
continue
contextual_item = copy.deepcopy(seg["metadata"])
contextual_item["camera"] = seg_camera
contextual_item["start_time"] = seg_start
contextual_item["end_time"] = seg_end
if seg["state_changes"]:
contextual_item["state_changes"] = seg["state_changes"]
contextual_items.append(contextual_item)
seen_contextual_cameras.add(seg_camera)
@@ -439,6 +457,7 @@ class ReviewDescriptionProcessor(PostProcessorApi):
captions = build_frame_captions(
final_data["data"].get("detections") or [],
[timestamp for _, timestamp in frames],
sorted_classification_state_changes(final_data["data"]),
)
if not captions:
@@ -688,6 +707,39 @@ def get_recording_buffer_extension(duration: float) -> float:
return buffer_extension
def sorted_classification_state_changes(
review_data: dict[str, Any],
) -> list[dict[str, Any]]:
"""A review item's state classification changes in time order."""
return sorted(
review_data.get("classification_state_changes") or [],
key=lambda change: change["timestamp"],
)
def format_classification_state_changes(
changes: list[dict[str, Any]], start_time: float, end_time: float
) -> list[str]:
"""Phrase state classification changes with their timing in the activity.
Changes are attached while the review item is active, which runs past its
end_time by the review cutoff, and a few seconds before its start.
"""
lines = []
for change in changes:
if change["timestamp"] < start_time:
when = "just before the activity started"
elif change["timestamp"] > end_time:
when = "after the activity ended"
else:
when = f"{round(change['timestamp'] - start_time)}s into the activity"
lines.append(f"{describe_classification_change(change)}, {when}")
return lines
def run_analysis(
requestor: InterProcessRequestor,
genai_client: GenAIClient,
@@ -743,6 +795,13 @@ def run_analysis(
unified_objects.append(object_type)
analytics_data["unified_objects"] = unified_objects
analytics_data["classification_state_changes"] = (
format_classification_state_changes(
sorted_classification_state_changes(final_data["data"]),
final_data["start_time"],
final_data["end_time"],
)
)
metadata = genai_client.generate_review_description(
analytics_data,
@@ -91,8 +91,23 @@ class CustomStateClassificationProcessor(DeferredRealtimeProcessorApi):
self.tensor_input_details = self.interpreter.get_input_details()
self.tensor_output_details = self.interpreter.get_output_details()
self.labelmap = load_labels(labelmap_path, prefill=0, indexed=False)
self._forget_unknown_states()
self.classifications_per_second.start()
def _forget_unknown_states(self) -> None:
"""Drop verified states that are not labels of the loaded model.
A retrained model can rename or remove labels. Keeping a state it can
no longer produce would report its first verified state as a change
from that obsolete label.
"""
labels = set(self.labelmap.values())
self.state_history = {
camera: history
for camera, history in self.state_history.items()
if history["current_state"] in labels
}
def __update_metrics(self, duration: float) -> None:
self.classifications_per_second.update()
if self.inference_speed:
@@ -134,15 +149,20 @@ class CustomStateClassificationProcessor(DeferredRealtimeProcessorApi):
# Don't save if state is stable (detected_state == current_state) AND score is 100%
return False
def verify_state_change(self, camera: str, detected_state: str) -> str | None:
def verify_state_change(
self, camera: str, detected_state: str, timestamp: float
) -> tuple[str | None, float] | None:
"""
Verify state change requires 3 consecutive identical states before publishing.
Returns state to publish or None if verification not complete.
Returns (previous state, time the new state was first seen) once verified,
or None if verification not complete. The previous state is None for the
first state verified on a camera.
"""
if camera not in self.state_history:
self.state_history[camera] = {
"current_state": None,
"pending_state": None,
"pending_since": 0.0,
"consecutive_count": 0,
}
@@ -157,12 +177,14 @@ class CustomStateClassificationProcessor(DeferredRealtimeProcessorApi):
verification["consecutive_count"] += 1
if verification["consecutive_count"] >= 3:
previous_state = verification["current_state"]
verification["current_state"] = detected_state
verification["pending_state"] = None
verification["consecutive_count"] = 0
return detected_state
return previous_state, verification["pending_since"]
else:
verification["pending_state"] = detected_state
verification["pending_since"] = timestamp
verification["consecutive_count"] = 1
logger.debug(
f"New state '{detected_state}' detected for {camera}, need {3 - verification['consecutive_count']} more consecutive detections"
@@ -340,16 +362,19 @@ class CustomStateClassificationProcessor(DeferredRealtimeProcessorApi):
)
return
verified_state = self.verify_state_change(camera, detected_state)
verified = self.verify_state_change(camera, detected_state, timestamp)
if verified_state is not None:
if verified is not None:
previous_state, changed_at = verified
self._emit_result(
{
"type": "classification",
"processor": "state",
"model_name": self.model_config.name,
"camera": camera,
"state": verified_state,
"state": detected_state,
"previous_state": previous_state,
"timestamp": changed_at,
}
)
+22 -1
View File
@@ -11,7 +11,11 @@ from typing import Any
from peewee import DoesNotExist
from frigate.comms.config_updater import ConfigSubscriber
from frigate.comms.detections_updater import DetectionSubscriber, DetectionTypeEnum
from frigate.comms.detections_updater import (
DetectionPublisher,
DetectionSubscriber,
DetectionTypeEnum,
)
from frigate.comms.embeddings_updater import (
EmbeddingsRequestEnum,
EmbeddingsResponder,
@@ -168,6 +172,7 @@ class EmbeddingMaintainer(threading.Thread):
)
self.review_subscriber = ReviewDataSubscriber("")
self.detection_subscriber = DetectionSubscriber(DetectionTypeEnum.video.value)
self.detection_publisher = DetectionPublisher(DetectionTypeEnum.all.value)
self.embeddings_responder = EmbeddingsResponder()
self.frame_manager = SharedMemoryFrameManager()
@@ -356,6 +361,7 @@ class EmbeddingMaintainer(threading.Thread):
self.event_end_subscriber.stop()
self.recordings_subscriber.stop()
self.detection_subscriber.stop()
self.detection_publisher.stop()
self.event_metadata_publisher.stop()
self.event_metadata_subscriber.stop()
self.embeddings_responder.stop()
@@ -851,6 +857,21 @@ class EmbeddingMaintainer(threading.Thread):
f"{result['camera']}/classification/{result['model_name']}",
result["state"],
)
# the first state verified after startup is not a change
if result["previous_state"] is not None:
self.detection_publisher.publish(
(
result["camera"],
{
"model": result["model_name"],
"from": result["previous_state"],
"to": result["state"],
"timestamp": result["timestamp"],
},
),
DetectionTypeEnum.classification_state.value,
)
elif result["processor"] == "object":
object_id = result["object_id"]
camera = result["camera"]
+28 -2
View File
@@ -104,6 +104,21 @@ def build_review_description_prompt(
else:
return "\n- (No objects detected)"
def get_state_changes_section() -> str:
# empty when nothing changed so the prompt is otherwise unaffected
changes = review_data.get("classification_state_changes")
if not changes:
return ""
return (
"\n\n## State Changes\n\n"
"The camera's state classifiers watch fixed areas of the scene and "
"reported these changes. They come from the classifiers rather than "
"from the images, and they are reliable. Describe each one where it "
"fits in the sequence of events.\n- " + "\n- ".join(changes)
)
fields = get_review_field_guidelines(response_style)
frame_guidance = f"\n{FRAME_ANNOTATION_GUIDANCE}" if frame_captions else ""
@@ -145,7 +160,7 @@ Respond with a JSON object matching the provided schema. Field-specific guidance
- Camera: {review_data["camera"]}
- Total frames: {len(thumbnails)} (Frame 1 = earliest, Frame {len(thumbnails)} = latest){frame_guidance}
- Activity started at {review_data["start"]} and lasted {review_data["duration"]} seconds
- Zones involved: {", ".join(review_data["zones"]) if review_data["zones"] else "None"}
- Zones involved: {", ".join(review_data["zones"]) if review_data["zones"] else "None"}{get_state_changes_section()}
## Objects in Scene
@@ -196,6 +211,17 @@ def build_review_summary_prompt(
f" to "
f"{datetime.datetime.fromtimestamp(end_ts).strftime('%B %d, %Y at %I:%M %p')}"
)
has_state_changes = any(
"state_changes" in item
for event in events
for item in [event, *event.get("context", [])]
)
state_changes_format = (
'\n- "state_changes" (only on some events): changes to monitored areas '
"reported by the camera's state classifiers, which are reliable"
if has_state_changes
else ""
)
prompt = f"""
You are a security officer writing a concise security report.
@@ -203,7 +229,7 @@ Time range: {time_range}
Input format: Each event is a JSON object with:
- "title", "scene", "confidence", "potential_threat_level" (0-2), "other_concerns", "camera", "time", "start_time", "end_time"
- "context": array of related events from other cameras that occurred during overlapping time periods
- "context": array of related events from other cameras that occurred during overlapping time periods{state_changes_format}
**Note: Use the "scene" field for event descriptions in the report. Ignore any "shortSummary" field if present.**
+1
View File
@@ -1127,6 +1127,7 @@ class RecordingMaintainer(threading.Thread):
elif (
topic == DetectionTypeEnum.api.value
or topic == DetectionTypeEnum.lpr.value
or topic == DetectionTypeEnum.classification_state.value
):
continue
+62 -5
View File
@@ -40,6 +40,10 @@ logger = logging.getLogger(__name__)
THUMB_HEIGHT = 180
THUMB_WIDTH = 320
# seconds before a review item starts that a state classification change is
# still attached to it, e.g. a garage door opening before the car is visible
CLASSIFICATION_STATE_PRE_ROLL = 5
class PendingReviewSegment:
def __init__(
@@ -61,6 +65,7 @@ class PendingReviewSegment:
self.sub_labels = sub_labels
self.zones = zones
self.audio = audio
self.classification_state_changes: list[dict[str, Any]] = []
self.thumb_time: float | None = None
self.last_alert_time: float | None = None
self.last_detection_time: float = frame_time
@@ -162,6 +167,7 @@ class PendingReviewSegment:
"sub_labels": list(self.sub_labels.values()),
"zones": self.zones,
"audio": list(self.audio),
"classification_state_changes": self.classification_state_changes,
"thumb_time": self.thumb_time,
"metadata": None,
},
@@ -293,6 +299,9 @@ class ReviewSegmentMaintainer(threading.Thread):
# manual events
self.indefinite_events: dict[str, dict[str, Any]] = {}
# state classification changes seen while a camera had no review item
self.recent_classification_state_changes: dict[str, list[dict[str, Any]]] = {}
# ensure dirs
Path(os.path.join(CLIPS_DIR, "review")).mkdir(exist_ok=True)
@@ -374,6 +383,43 @@ class ReviewSegmentMaintainer(threading.Thread):
self.active_review_segments[segment.camera] = None
return end_time
def _activate_segment(self, segment: PendingReviewSegment) -> None:
"""Make a segment the camera's active one, attaching any state
classification changes seen just before it started."""
self.active_review_segments[segment.camera] = segment
recent = self.recent_classification_state_changes.pop(segment.camera, [])
segment.classification_state_changes.extend(
c
for c in recent
if c["timestamp"] >= segment.start_time - CLASSIFICATION_STATE_PRE_ROLL
)
def handle_classification_state_change(
self, camera: str, change: dict[str, Any]
) -> None:
"""Attach a verified state classification change to the active segment.
State changes never start, extend, or upgrade a segment. A change seen
with no active segment is held briefly for a segment starting right
after it.
"""
segment = self.active_review_segments.get(camera)
if segment is None:
cutoff = change["timestamp"] - CLASSIFICATION_STATE_PRE_ROLL
self.recent_classification_state_changes[camera] = [
c
for c in self.recent_classification_state_changes.get(camera, [])
if c["timestamp"] >= cutoff
] + [change]
return
prev_data = segment.get_data(False)
segment.classification_state_changes.append(change)
self._publish_segment_update(
segment, self.config.cameras[camera], None, [], prev_data
)
def forcibly_end_segment(self, camera: str) -> Any:
"""Forcibly end the pending segment for a camera."""
segment = self.active_review_segments.get(camera)
@@ -422,6 +468,7 @@ class ReviewSegmentMaintainer(threading.Thread):
"""Close out a deleted camera's segment so a reused name cannot inherit it."""
self.forcibly_end_segment(camera)
self.indefinite_events.pop(camera, None)
self.recent_classification_state_changes.pop(camera, None)
def update_existing_segment(
self,
@@ -579,7 +626,7 @@ class ReviewSegmentMaintainer(threading.Thread):
audio=set(),
zones=list(new_zones),
)
self.active_review_segments[segment.camera] = new_segment
self._activate_segment(new_segment)
self._publish_segment_start(new_segment)
new_segment.last_detection_time = last_detection_time
elif segment.severity == SeverityEnum.detection and frame_time > (
@@ -639,7 +686,7 @@ class ReviewSegmentMaintainer(threading.Thread):
audio=set(),
zones=zones,
)
self.active_review_segments[camera] = new_segment
self._activate_segment(new_segment)
try:
yuv_frame = self.frame_manager.get(
@@ -713,6 +760,10 @@ class ReviewSegmentMaintainer(threading.Thread):
if camera not in self.indefinite_events:
self.indefinite_events[camera] = {}
elif topic == DetectionTypeEnum.classification_state.value:
(camera, classification_change) = data
else:
continue
if camera not in self.config.cameras:
continue
@@ -723,6 +774,10 @@ class ReviewSegmentMaintainer(threading.Thread):
):
continue
if topic == DetectionTypeEnum.classification_state:
self.handle_classification_state_change(camera, classification_change)
continue
current_segment = self.active_review_segments.get(camera)
# Check if the current segment should be processed based on enabled settings
@@ -864,7 +919,8 @@ class ReviewSegmentMaintainer(threading.Thread):
severity = SeverityEnum.detection
if severity:
self.active_review_segments[camera] = PendingReviewSegment(
self._activate_segment(
PendingReviewSegment(
camera,
frame_time,
severity,
@@ -873,6 +929,7 @@ class ReviewSegmentMaintainer(threading.Thread):
[],
detections,
)
)
elif topic == DetectionTypeEnum.api:
severity = self.get_manual_event_severity(
camera, manual_info["label"]
@@ -888,7 +945,7 @@ class ReviewSegmentMaintainer(threading.Thread):
[],
set(),
)
self.active_review_segments[camera] = api_segment
self._activate_segment(api_segment)
if manual_info["state"] == ManualEventState.start:
self.indefinite_events[camera][manual_info["event_id"]] = (
@@ -915,7 +972,7 @@ class ReviewSegmentMaintainer(threading.Thread):
[],
set(),
)
self.active_review_segments[camera] = lpr_segment
self._activate_segment(lpr_segment)
if manual_info["state"] == ManualEventState.start:
self.indefinite_events[camera][manual_info["event_id"]] = (
+4
View File
@@ -179,12 +179,16 @@ class TestReviewMaintainerRemoval(unittest.TestCase):
maintainer = ReviewSegmentMaintainer.__new__(ReviewSegmentMaintainer)
maintainer.active_review_segments = {"deleted_cam": MagicMock()}
maintainer.indefinite_events = {"deleted_cam": {"1234.5-abcdef": 1.0}}
maintainer.recent_classification_state_changes = {
"deleted_cam": [{"model": "gate", "from": "a", "to": "b", "timestamp": 1.0}]
}
maintainer.forcibly_end_segment = MagicMock()
maintainer._handle_camera_removed("deleted_cam")
maintainer.forcibly_end_segment.assert_called_once_with("deleted_cam")
self.assertNotIn("deleted_cam", maintainer.indefinite_events)
self.assertNotIn("deleted_cam", maintainer.recent_classification_state_changes)
class TestAutotrackerMoveQueue(unittest.TestCase):
+45
View File
@@ -7,6 +7,7 @@ from frigate.genai.prompts import (
REVIEW_DESCRIPTION_FIELD_GUIDELINES,
REVIEW_RESPONSE_STYLES,
build_review_description_prompt,
build_review_summary_prompt,
get_review_field_guidelines,
)
@@ -83,5 +84,49 @@ class TestReviewResponseStyle(unittest.TestCase):
)
class TestClassificationStateChanges(unittest.TestCase):
def _build_prompt(self, **extra) -> str:
review_data = {
"camera": "Front Door",
"start": "Monday, 09:30 AM",
"duration": 25,
"zones": [],
"unified_objects": ["person"],
**extra,
}
return build_review_description_prompt(
review_data, [b"fake-image"], [], None, "activity context"
)
def test_no_changes_leaves_prompt_unchanged(self):
self.assertEqual(
self._build_prompt(classification_state_changes=[]),
self._build_prompt(),
)
self.assertNotIn("## State Changes", self._build_prompt())
def test_changes_are_listed_before_objects(self):
change = "front gate changed from closed to open, 12s into the activity"
prompt = self._build_prompt(classification_state_changes=[change])
self.assertIn(f"\n- {change}\n\n## Objects in Scene", prompt)
self.assertLess(
prompt.index("## Sequence Details"), prompt.index("## State Changes")
)
def test_summary_describes_state_changes_only_when_present(self):
event = {"title": "Person at door", "camera": "Front Door", "context": []}
without = build_review_summary_prompt(0, 3600, [event], None)
self.assertNotIn('"state_changes"', without)
context_event = {
**event,
"context": [
{"camera": "Driveway", "state_changes": ["gate changed from a to b"]}
],
}
with_changes = build_review_summary_prompt(0, 3600, [context_event], None)
self.assertIn('- "state_changes"', with_changes)
if __name__ == "__main__":
unittest.main()
+73
View File
@@ -1,10 +1,13 @@
"""Tests for tracker-derived review frame annotations."""
import unittest
from unittest.mock import patch
from frigate.data_processing.post.review_annotations import (
annotations_by_frame,
build_frame_captions,
build_timeline,
describe_classification_change,
describe_heading,
describe_position,
event_name,
@@ -312,5 +315,75 @@ class TestFrameBucketing(unittest.TestCase):
self.assertEqual(annotations_by_frame([(1.0, "x")], []), {})
class TestClassificationChangeCaptions(unittest.TestCase):
def setUp(self):
person = track(
"1789481994.684479-lpyc2z",
"person",
0.0,
straight_path((0.9, 0.6), (0.4, 0.3), 10, 0.0),
)
patcher = patch(
"frigate.data_processing.post.review_annotations.get_tracked_events",
return_value=[person],
)
patcher.start()
self.addCleanup(patcher.stop)
patcher = patch(
"frigate.data_processing.post.review_annotations.get_state_changes",
return_value=[],
)
patcher.start()
self.addCleanup(patcher.stop)
def test_change_is_phrased_without_underscores(self):
self.assertEqual(
describe_classification_change(
{"model": "trash_day", "from": "no_bins", "to": "bins_at_curb"}
),
"trash day changed from no bins to bins at curb",
)
def test_change_is_noted_before_the_frame_it_precedes(self):
change = {"model": "front_gate", "from": "closed", "to": "open"}
captions = build_frame_captions(
["1789481994.684479-lpyc2z"],
[0.0, 10.0, 20.0],
[{**change, "timestamp": 12.0}],
)
self.assertEqual(
captions[2],
"Frame 3 of 3 (+20.0s):\n[state] front gate changed from closed to open",
)
self.assertTrue(captions[0].splitlines()[1].startswith("[tracker] "))
def test_change_is_noted_without_tracked_objects(self):
with patch(
"frigate.data_processing.post.review_annotations.get_tracked_events",
return_value=[],
):
captions = build_frame_captions(
[],
[0.0, 10.0],
[{"model": "gate", "from": "a", "to": "b", "timestamp": 5.0}],
)
self.assertEqual(
captions,
[
"Frame 1 of 2 (+0.0s):",
"Frame 2 of 2 (+10.0s):\n[state] gate changed from a to b",
],
)
def test_change_after_the_last_frame_is_dropped(self):
captions = build_frame_captions(
["1789481994.684479-lpyc2z"],
[0.0, 10.0, 20.0],
[{"model": "gate", "from": "a", "to": "b", "timestamp": 25.0}],
)
self.assertFalse(any("[state]" in caption for caption in captions))
if __name__ == "__main__":
unittest.main()
@@ -0,0 +1,258 @@
"""Tests for attaching state classification changes to review items."""
import unittest
from unittest.mock import MagicMock, patch
from frigate.comms.embeddings_updater import EmbeddingsRequestEnum
from frigate.config import FrigateConfig
from frigate.data_processing.post.review_descriptions import (
ReviewDescriptionProcessor,
format_classification_state_changes,
)
from frigate.data_processing.real_time.custom_classification import (
CustomStateClassificationProcessor,
)
from frigate.models import ReviewSegment
from frigate.review.maintainer import (
CLASSIFICATION_STATE_PRE_ROLL,
PendingReviewSegment,
ReviewSegmentMaintainer,
)
from frigate.review.types import SeverityEnum
CONFIG = """
mqtt:
enabled: False
cameras:
front_door:
ffmpeg:
inputs:
- path: rtsp://10.0.0.1:554/video
roles:
- detect
detect:
width: 1920
height: 1080
fps: 5
"""
def gate_change(timestamp: float, before: str = "closed", after: str = "open"):
return {"model": "front_gate", "from": before, "to": after, "timestamp": timestamp}
class TestVerifyStateChange(unittest.TestCase):
def setUp(self):
self.processor = CustomStateClassificationProcessor.__new__(
CustomStateClassificationProcessor
)
self.processor.state_history = {}
def verify(self, state: str, timestamp: float):
return self.processor.verify_state_change("front_door", state, timestamp)
def test_first_verified_state_has_no_previous_state(self):
self.assertIsNone(self.verify("closed", 1.0))
self.assertIsNone(self.verify("closed", 2.0))
self.assertEqual(self.verify("closed", 3.0), (None, 1.0))
def test_change_reports_previous_state_and_first_sighting(self):
for timestamp in (1.0, 2.0, 3.0):
self.verify("closed", timestamp)
self.assertIsNone(self.verify("open", 10.0))
self.assertIsNone(self.verify("open", 11.0))
self.assertEqual(self.verify("open", 12.0), ("closed", 10.0))
def test_interrupted_verification_restarts_the_first_sighting(self):
for timestamp in (1.0, 2.0, 3.0):
self.verify("closed", timestamp)
self.verify("open", 10.0)
self.verify("closed", 11.0)
self.verify("open", 20.0)
self.verify("open", 21.0)
self.assertEqual(self.verify("open", 22.0), ("closed", 20.0))
def test_reload_forgets_states_the_model_no_longer_has(self):
for timestamp in (1.0, 2.0, 3.0):
self.verify("closed", timestamp)
self.processor.state_history["back_door"] = {"current_state": "open"}
self.processor.labelmap = {0: "open", 1: "shut"}
self.processor._forget_unknown_states()
self.assertEqual(list(self.processor.state_history), ["back_door"])
self.assertIsNone(self.verify("shut", 10.0))
self.assertIsNone(self.verify("shut", 11.0))
self.assertEqual(self.verify("shut", 12.0), (None, 10.0))
class TestReviewSegmentAttachment(unittest.TestCase):
def setUp(self):
self.maintainer = ReviewSegmentMaintainer.__new__(ReviewSegmentMaintainer)
self.maintainer.config = FrigateConfig.parse_yaml(CONFIG)
self.maintainer.active_review_segments = {}
self.maintainer.recent_classification_state_changes = {}
self.maintainer._publish_segment_update = MagicMock()
def segment(self, start_time: float) -> PendingReviewSegment:
return PendingReviewSegment(
"front_door",
start_time,
SeverityEnum.alert,
{"1.0-abcdef": "person"},
{},
[],
set(),
)
def test_change_during_segment_is_attached_and_published(self):
segment = self.segment(100.0)
self.maintainer.active_review_segments["front_door"] = segment
self.maintainer.handle_classification_state_change(
"front_door", gate_change(110.0)
)
self.assertEqual(segment.classification_state_changes, [gate_change(110.0)])
self.maintainer._publish_segment_update.assert_called_once()
prev_data = self.maintainer._publish_segment_update.call_args.args[4]
self.assertEqual(prev_data["data"]["classification_state_changes"], [])
self.assertEqual(
segment.get_data(False)["data"]["classification_state_changes"],
[gate_change(110.0)],
)
def test_change_never_starts_a_segment(self):
self.maintainer.handle_classification_state_change(
"front_door", gate_change(110.0)
)
self.assertIsNone(self.maintainer.active_review_segments.get("front_door"))
self.maintainer._publish_segment_update.assert_not_called()
def test_change_just_before_a_segment_is_attached_when_it_starts(self):
self.maintainer.handle_classification_state_change(
"front_door", gate_change(98.0)
)
segment = self.segment(100.0)
self.maintainer._activate_segment(segment)
self.assertIs(self.maintainer.active_review_segments["front_door"], segment)
self.assertEqual(segment.classification_state_changes, [gate_change(98.0)])
self.assertNotIn(
"front_door", self.maintainer.recent_classification_state_changes
)
def test_change_long_before_a_segment_is_not_attached(self):
self.maintainer.handle_classification_state_change(
"front_door", gate_change(100.0 - CLASSIFICATION_STATE_PRE_ROLL - 1)
)
segment = self.segment(100.0)
self.maintainer._activate_segment(segment)
self.assertEqual(segment.classification_state_changes, [])
def test_held_changes_are_pruned(self):
self.maintainer.handle_classification_state_change(
"front_door", gate_change(10.0)
)
self.maintainer.handle_classification_state_change(
"front_door", gate_change(50.0, "open", "closed")
)
self.assertEqual(
self.maintainer.recent_classification_state_changes["front_door"],
[gate_change(50.0, "open", "closed")],
)
class TestChangeTiming(unittest.TestCase):
def test_changes_are_placed_relative_to_the_activity(self):
lines = format_classification_state_changes(
[
gate_change(98.0),
gate_change(112.4, "open", "closed"),
gate_change(140.0),
],
start_time=100.0,
end_time=130.0,
)
self.assertEqual(
lines,
[
"front gate changed from closed to open, "
"just before the activity started",
"front gate changed from open to closed, 12s into the activity",
"front gate changed from closed to open, after the activity ended",
],
)
class TestSummaryContext(unittest.TestCase):
def row(self, camera, start, end, threat, changes=()):
return {
"camera": camera,
"start_time": start,
"end_time": end,
"data": {
"metadata": {"title": camera, "potential_threat_level": threat},
"classification_state_changes": list(changes),
},
}
def summarize(self, rows):
processor = ReviewDescriptionProcessor.__new__(ReviewDescriptionProcessor)
processor.config = FrigateConfig.parse_yaml(CONFIG)
processor.genai_manager = MagicMock()
client = processor.genai_manager.description_client
with patch.object(ReviewSegment, "select") as select:
query = select.return_value.where.return_value.order_by.return_value
query.dicts.return_value.iterator.return_value = iter(rows)
processor.handle_request(
EmbeddingsRequestEnum.summarize_review.value,
{"start_ts": 0, "end_ts": 100},
)
return client.generate_review_summary.call_args.args[2]
def test_context_state_changes_stay_with_their_review(self):
events = self.summarize(
[
self.row("front_door", 10, 60, 1),
self.row("driveway", 15, 25, 0, [gate_change(20.0)]),
self.row("driveway", 30, 40, 0, [gate_change(35.0, "open", "closed")]),
]
)
self.assertEqual(
[
(item["start_time"], item["end_time"], item["state_changes"])
for item in events[0]["context"]
],
[
(15, 25, ["front gate changed from closed to open"]),
(30, 40, ["front gate changed from open to closed"]),
],
)
def test_later_context_review_without_changes_is_deduplicated(self):
events = self.summarize(
[
self.row("front_door", 10, 60, 1),
self.row("driveway", 15, 25, 0, [gate_change(20.0)]),
self.row("driveway", 30, 40, 0),
]
)
self.assertEqual(len(events[0]["context"]), 1)
self.assertEqual(events[0]["context"][0]["start_time"], 15)
if __name__ == "__main__":
unittest.main()