Integrate state changes with review items

This commit is contained in:
Nicolas Mowen
2026-09-29 09:47:36 -06:00
parent d092189b8e
commit cdda89a11e
15 changed files with 528 additions and 34 deletions
@@ -85,6 +85,14 @@ An optional config, `save_attempts`, can be set as a key under the model name. T
</TabItem> </TabItem>
</ConfigTabs> </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/review_summaries) 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 ## Training the model
Creating and training the model is done within the Frigate UI using the `Classification` page. The process consists of three steps: 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`. 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. 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 :::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` - The severity changes from `detection` to `alert`
- Additional objects are detected - Additional objects are detected
- An object is recognized via face, lpr, etc. - 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. 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"], "objects": ["person", "car"],
"sub_labels": [], "sub_labels": [],
"zones": [], "zones": [],
"audio": [] "audio": [],
"classification_state_changes": []
} }
}, },
"after": { "after": {
@@ -254,7 +256,16 @@ When the review activity has ended a final `end` message is published.
"objects": ["person", "car"], "objects": ["person", "car"],
"sub_labels": ["Bob"], "sub_labels": ["Bob"],
"zones": ["front_yard"], "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" video = "video"
audio = "audio" audio = "audio"
lpr = "lpr" lpr = "lpr"
classification_state = "classification_state"
class DetectionPublisher(Publisher): class DetectionPublisher(Publisher):
@@ -1,10 +1,10 @@
"""Frame annotations derived from object tracking data. """Frame annotations derived from object tracking data.
Builds short notes describing what changed during a review item, keyed to the 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 frames sampled from it. Everything here comes from data already recorded
in the database (each event's `path_data` trajectory and the timeline's (each event's `path_data` trajectory, the timeline's stationary/active
stationary/active changes), so the notes can be stated to the model as fact changes, and the review item's state classification changes), so the notes
rather than as something it must perceive. can be stated to the model as fact rather than as something it must perceive.
""" """
import logging import logging
@@ -95,6 +95,15 @@ def event_name(event: dict[str, Any]) -> str:
return f"{article} {label}" 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]: def path_legs(points: list[Point]) -> list[Leg]:
"""Split a trajectory into runs of travel in a consistent direction. """Split a trajectory into runs of travel in a consistent direction.
@@ -359,13 +368,15 @@ def get_state_changes(detection_ids: list[str]) -> list[dict[str, Any]]:
def build_frame_captions( def build_frame_captions(
detection_ids: list[str], detection_ids: list[str],
frame_times: list[float], frame_times: list[float],
classification_changes: Sequence[dict[str, Any]] = (),
) -> list[str]: ) -> list[str]:
"""A caption for each sampled frame, in frame order. """A caption for each sampled frame, in frame order.
Every frame gets its index and elapsed time so the model can tell them 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 apart; frames where something changed also carry the tracker and state
that moment. Returns an empty list when there is nothing to say, which classification notes for that moment. Returns an empty list when there is
callers treat as a reason to fall back to sending plain frames. nothing to say, which callers treat as a reason to fall back to sending
plain frames.
""" """
if not frame_times: if not frame_times:
return [] return []
@@ -376,8 +387,19 @@ def build_frame_captions(
logger.debug("No tracked events found for review item, skipping annotations") logger.debug("No tracked events found for review item, skipping annotations")
return [] return []
timeline = build_timeline(events, frame_times[-1], get_state_changes(detection_ids)) span_end = frame_times[-1]
buckets = annotations_by_frame(timeline, frame_times) timeline = [
(timestamp, f"[tracker] {note}")
for timestamp, note in build_timeline(
events, span_end, get_state_changes(detection_ids)
)
]
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: if not buckets:
return [] return []
@@ -388,7 +410,7 @@ def build_frame_captions(
for index, timestamp in enumerate(frame_times): for index, timestamp in enumerate(frame_times):
lines = [f"Frame {index + 1} of {total} (+{timestamp - origin:.1f}s):"] 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)) captions.append("\n".join(lines))
return captions return captions
@@ -40,7 +40,7 @@ from frigate.util.image import get_image_from_recording
from ..post.api import PostProcessorApi from ..post.api import PostProcessorApi
from ..types import DataProcessorMetrics 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__) logger = logging.getLogger(__name__)
@@ -254,6 +254,10 @@ class ReviewDescriptionProcessor(PostProcessorApi):
"start_time": r["start_time"], "start_time": r["start_time"],
"end_time": r["end_time"], "end_time": r["end_time"],
"metadata": r["data"]["metadata"], "metadata": r["data"]["metadata"],
"state_changes": [
describe_classification_change(change)
for change in sorted_classification_state_changes(r["data"])
],
} }
for r in ( for r in (
ReviewSegment.select( ReviewSegment.select(
@@ -298,6 +302,9 @@ class ReviewDescriptionProcessor(PostProcessorApi):
primary_item["start_time"] = primary_seg["start_time"] primary_item["start_time"] = primary_seg["start_time"]
primary_item["end_time"] = primary_seg["end_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 # Find overlapping contextual items from other cameras
primary_start = primary_seg["start_time"] primary_start = primary_seg["start_time"]
primary_end = primary_seg["end_time"] primary_end = primary_seg["end_time"]
@@ -324,6 +331,9 @@ class ReviewDescriptionProcessor(PostProcessorApi):
contextual_item["camera"] = seg_camera contextual_item["camera"] = seg_camera
contextual_item["start_time"] = seg_start contextual_item["start_time"] = seg_start
contextual_item["end_time"] = seg_end contextual_item["end_time"] = seg_end
if seg["state_changes"]:
contextual_item["state_changes"] = seg["state_changes"]
contextual_items.append(contextual_item) contextual_items.append(contextual_item)
seen_contextual_cameras.add(seg_camera) seen_contextual_cameras.add(seg_camera)
@@ -439,6 +449,7 @@ class ReviewDescriptionProcessor(PostProcessorApi):
captions = build_frame_captions( captions = build_frame_captions(
final_data["data"].get("detections") or [], final_data["data"].get("detections") or [],
[timestamp for _, timestamp in frames], [timestamp for _, timestamp in frames],
sorted_classification_state_changes(final_data["data"]),
) )
if not captions: if not captions:
@@ -688,6 +699,39 @@ def get_recording_buffer_extension(duration: float) -> float:
return buffer_extension 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( def run_analysis(
requestor: InterProcessRequestor, requestor: InterProcessRequestor,
genai_client: GenAIClient, genai_client: GenAIClient,
@@ -743,6 +787,13 @@ def run_analysis(
unified_objects.append(object_type) unified_objects.append(object_type)
analytics_data["unified_objects"] = unified_objects 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( metadata = genai_client.generate_review_description(
analytics_data, analytics_data,
@@ -134,15 +134,20 @@ class CustomStateClassificationProcessor(DeferredRealtimeProcessorApi):
# Don't save if state is stable (detected_state == current_state) AND score is 100% # Don't save if state is stable (detected_state == current_state) AND score is 100%
return False 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. 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: if camera not in self.state_history:
self.state_history[camera] = { self.state_history[camera] = {
"current_state": None, "current_state": None,
"pending_state": None, "pending_state": None,
"pending_since": 0.0,
"consecutive_count": 0, "consecutive_count": 0,
} }
@@ -157,12 +162,14 @@ class CustomStateClassificationProcessor(DeferredRealtimeProcessorApi):
verification["consecutive_count"] += 1 verification["consecutive_count"] += 1
if verification["consecutive_count"] >= 3: if verification["consecutive_count"] >= 3:
previous_state = verification["current_state"]
verification["current_state"] = detected_state verification["current_state"] = detected_state
verification["pending_state"] = None verification["pending_state"] = None
verification["consecutive_count"] = 0 verification["consecutive_count"] = 0
return detected_state return previous_state, verification["pending_since"]
else: else:
verification["pending_state"] = detected_state verification["pending_state"] = detected_state
verification["pending_since"] = timestamp
verification["consecutive_count"] = 1 verification["consecutive_count"] = 1
logger.debug( logger.debug(
f"New state '{detected_state}' detected for {camera}, need {3 - verification['consecutive_count']} more consecutive detections" f"New state '{detected_state}' detected for {camera}, need {3 - verification['consecutive_count']} more consecutive detections"
@@ -340,16 +347,19 @@ class CustomStateClassificationProcessor(DeferredRealtimeProcessorApi):
) )
return 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( self._emit_result(
{ {
"type": "classification", "type": "classification",
"processor": "state", "processor": "state",
"model_name": self.model_config.name, "model_name": self.model_config.name,
"camera": camera, "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 peewee import DoesNotExist
from frigate.comms.config_updater import ConfigSubscriber 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 ( from frigate.comms.embeddings_updater import (
EmbeddingsRequestEnum, EmbeddingsRequestEnum,
EmbeddingsResponder, EmbeddingsResponder,
@@ -168,6 +172,7 @@ class EmbeddingMaintainer(threading.Thread):
) )
self.review_subscriber = ReviewDataSubscriber("") self.review_subscriber = ReviewDataSubscriber("")
self.detection_subscriber = DetectionSubscriber(DetectionTypeEnum.video.value) self.detection_subscriber = DetectionSubscriber(DetectionTypeEnum.video.value)
self.detection_publisher = DetectionPublisher(DetectionTypeEnum.all.value)
self.embeddings_responder = EmbeddingsResponder() self.embeddings_responder = EmbeddingsResponder()
self.frame_manager = SharedMemoryFrameManager() self.frame_manager = SharedMemoryFrameManager()
@@ -356,6 +361,7 @@ class EmbeddingMaintainer(threading.Thread):
self.event_end_subscriber.stop() self.event_end_subscriber.stop()
self.recordings_subscriber.stop() self.recordings_subscriber.stop()
self.detection_subscriber.stop() self.detection_subscriber.stop()
self.detection_publisher.stop()
self.event_metadata_publisher.stop() self.event_metadata_publisher.stop()
self.event_metadata_subscriber.stop() self.event_metadata_subscriber.stop()
self.embeddings_responder.stop() self.embeddings_responder.stop()
@@ -851,6 +857,21 @@ class EmbeddingMaintainer(threading.Thread):
f"{result['camera']}/classification/{result['model_name']}", f"{result['camera']}/classification/{result['model_name']}",
result["state"], 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": elif result["processor"] == "object":
object_id = result["object_id"] object_id = result["object_id"]
camera = result["camera"] camera = result["camera"]
+28 -2
View File
@@ -104,6 +104,21 @@ def build_review_description_prompt(
else: else:
return "\n- (No objects detected)" 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) fields = get_review_field_guidelines(response_style)
frame_guidance = f"\n{FRAME_ANNOTATION_GUIDANCE}" if frame_captions else "" 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"]} - Camera: {review_data["camera"]}
- Total frames: {len(thumbnails)} (Frame 1 = earliest, Frame {len(thumbnails)} = latest){frame_guidance} - 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 - 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 ## Objects in Scene
@@ -196,6 +211,17 @@ def build_review_summary_prompt(
f" to " f" to "
f"{datetime.datetime.fromtimestamp(end_ts).strftime('%B %d, %Y at %I:%M %p')}" 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""" prompt = f"""
You are a security officer writing a concise security report. 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: 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" - "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.** **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 ( elif (
topic == DetectionTypeEnum.api.value topic == DetectionTypeEnum.api.value
or topic == DetectionTypeEnum.lpr.value or topic == DetectionTypeEnum.lpr.value
or topic == DetectionTypeEnum.classification_state.value
): ):
continue continue
+69 -12
View File
@@ -40,6 +40,10 @@ logger = logging.getLogger(__name__)
THUMB_HEIGHT = 180 THUMB_HEIGHT = 180
THUMB_WIDTH = 320 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: class PendingReviewSegment:
def __init__( def __init__(
@@ -61,6 +65,7 @@ class PendingReviewSegment:
self.sub_labels = sub_labels self.sub_labels = sub_labels
self.zones = zones self.zones = zones
self.audio = audio self.audio = audio
self.classification_state_changes: list[dict[str, Any]] = []
self.thumb_time: float | None = None self.thumb_time: float | None = None
self.last_alert_time: float | None = None self.last_alert_time: float | None = None
self.last_detection_time: float = frame_time self.last_detection_time: float = frame_time
@@ -162,6 +167,7 @@ class PendingReviewSegment:
"sub_labels": list(self.sub_labels.values()), "sub_labels": list(self.sub_labels.values()),
"zones": self.zones, "zones": self.zones,
"audio": list(self.audio), "audio": list(self.audio),
"classification_state_changes": self.classification_state_changes,
"thumb_time": self.thumb_time, "thumb_time": self.thumb_time,
"metadata": None, "metadata": None,
}, },
@@ -293,6 +299,9 @@ class ReviewSegmentMaintainer(threading.Thread):
# manual events # manual events
self.indefinite_events: dict[str, dict[str, Any]] = {} 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 # ensure dirs
Path(os.path.join(CLIPS_DIR, "review")).mkdir(exist_ok=True) 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 self.active_review_segments[segment.camera] = None
return end_time 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: def forcibly_end_segment(self, camera: str) -> Any:
"""Forcibly end the pending segment for a camera.""" """Forcibly end the pending segment for a camera."""
segment = self.active_review_segments.get(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.""" """Close out a deleted camera's segment so a reused name cannot inherit it."""
self.forcibly_end_segment(camera) self.forcibly_end_segment(camera)
self.indefinite_events.pop(camera, None) self.indefinite_events.pop(camera, None)
self.recent_classification_state_changes.pop(camera, None)
def update_existing_segment( def update_existing_segment(
self, self,
@@ -579,7 +626,7 @@ class ReviewSegmentMaintainer(threading.Thread):
audio=set(), audio=set(),
zones=list(new_zones), zones=list(new_zones),
) )
self.active_review_segments[segment.camera] = new_segment self._activate_segment(new_segment)
self._publish_segment_start(new_segment) self._publish_segment_start(new_segment)
new_segment.last_detection_time = last_detection_time new_segment.last_detection_time = last_detection_time
elif segment.severity == SeverityEnum.detection and frame_time > ( elif segment.severity == SeverityEnum.detection and frame_time > (
@@ -639,7 +686,7 @@ class ReviewSegmentMaintainer(threading.Thread):
audio=set(), audio=set(),
zones=zones, zones=zones,
) )
self.active_review_segments[camera] = new_segment self._activate_segment(new_segment)
try: try:
yuv_frame = self.frame_manager.get( yuv_frame = self.frame_manager.get(
@@ -713,6 +760,10 @@ class ReviewSegmentMaintainer(threading.Thread):
if camera not in self.indefinite_events: if camera not in self.indefinite_events:
self.indefinite_events[camera] = {} self.indefinite_events[camera] = {}
elif topic == DetectionTypeEnum.classification_state.value:
(camera, classification_change) = data
else:
continue
if camera not in self.config.cameras: if camera not in self.config.cameras:
continue continue
@@ -723,6 +774,10 @@ class ReviewSegmentMaintainer(threading.Thread):
): ):
continue continue
if topic == DetectionTypeEnum.classification_state:
self.handle_classification_state_change(camera, classification_change)
continue
current_segment = self.active_review_segments.get(camera) current_segment = self.active_review_segments.get(camera)
# Check if the current segment should be processed based on enabled settings # Check if the current segment should be processed based on enabled settings
@@ -864,14 +919,16 @@ class ReviewSegmentMaintainer(threading.Thread):
severity = SeverityEnum.detection severity = SeverityEnum.detection
if severity: if severity:
self.active_review_segments[camera] = PendingReviewSegment( self._activate_segment(
camera, PendingReviewSegment(
frame_time, camera,
severity, frame_time,
{}, severity,
{}, {},
[], {},
detections, [],
detections,
)
) )
elif topic == DetectionTypeEnum.api: elif topic == DetectionTypeEnum.api:
severity = self.get_manual_event_severity( severity = self.get_manual_event_severity(
@@ -888,7 +945,7 @@ class ReviewSegmentMaintainer(threading.Thread):
[], [],
set(), set(),
) )
self.active_review_segments[camera] = api_segment self._activate_segment(api_segment)
if manual_info["state"] == ManualEventState.start: if manual_info["state"] == ManualEventState.start:
self.indefinite_events[camera][manual_info["event_id"]] = ( self.indefinite_events[camera][manual_info["event_id"]] = (
@@ -915,7 +972,7 @@ class ReviewSegmentMaintainer(threading.Thread):
[], [],
set(), set(),
) )
self.active_review_segments[camera] = lpr_segment self._activate_segment(lpr_segment)
if manual_info["state"] == ManualEventState.start: if manual_info["state"] == ManualEventState.start:
self.indefinite_events[camera][manual_info["event_id"]] = ( self.indefinite_events[camera][manual_info["event_id"]] = (
+4
View File
@@ -179,12 +179,16 @@ class TestReviewMaintainerRemoval(unittest.TestCase):
maintainer = ReviewSegmentMaintainer.__new__(ReviewSegmentMaintainer) maintainer = ReviewSegmentMaintainer.__new__(ReviewSegmentMaintainer)
maintainer.active_review_segments = {"deleted_cam": MagicMock()} maintainer.active_review_segments = {"deleted_cam": MagicMock()}
maintainer.indefinite_events = {"deleted_cam": {"1234.5-abcdef": 1.0}} 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.forcibly_end_segment = MagicMock()
maintainer._handle_camera_removed("deleted_cam") maintainer._handle_camera_removed("deleted_cam")
maintainer.forcibly_end_segment.assert_called_once_with("deleted_cam") maintainer.forcibly_end_segment.assert_called_once_with("deleted_cam")
self.assertNotIn("deleted_cam", maintainer.indefinite_events) self.assertNotIn("deleted_cam", maintainer.indefinite_events)
self.assertNotIn("deleted_cam", maintainer.recent_classification_state_changes)
class TestAutotrackerMoveQueue(unittest.TestCase): class TestAutotrackerMoveQueue(unittest.TestCase):
+45
View File
@@ -7,6 +7,7 @@ from frigate.genai.prompts import (
REVIEW_DESCRIPTION_FIELD_GUIDELINES, REVIEW_DESCRIPTION_FIELD_GUIDELINES,
REVIEW_RESPONSE_STYLES, REVIEW_RESPONSE_STYLES,
build_review_description_prompt, build_review_description_prompt,
build_review_summary_prompt,
get_review_field_guidelines, 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__": if __name__ == "__main__":
unittest.main() unittest.main()
+54
View File
@@ -1,10 +1,13 @@
"""Tests for tracker-derived review frame annotations.""" """Tests for tracker-derived review frame annotations."""
import unittest import unittest
from unittest.mock import patch
from frigate.data_processing.post.review_annotations import ( from frigate.data_processing.post.review_annotations import (
annotations_by_frame, annotations_by_frame,
build_frame_captions,
build_timeline, build_timeline,
describe_classification_change,
describe_heading, describe_heading,
describe_position, describe_position,
event_name, event_name,
@@ -312,5 +315,56 @@ class TestFrameBucketing(unittest.TestCase):
self.assertEqual(annotations_by_frame([(1.0, "x")], []), {}) 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_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__": if __name__ == "__main__":
unittest.main() unittest.main()
@@ -0,0 +1,181 @@
"""Tests for attaching state classification changes to review items."""
import unittest
from unittest.mock import MagicMock
from frigate.config import FrigateConfig
from frigate.data_processing.post.review_descriptions import (
format_classification_state_changes,
)
from frigate.data_processing.real_time.custom_classification import (
CustomStateClassificationProcessor,
)
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))
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",
],
)
if __name__ == "__main__":
unittest.main()