mirror of
https://github.com/blakeblackshear/frigate.git
synced 2026-09-29 19:36:57 +03:00
Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
c6d8f067be | ||
|
|
7767f67c67 | ||
|
|
cdda89a11e |
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
]
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
@@ -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"]
|
||||
|
||||
@@ -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.**
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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"]] = (
|
||||
|
||||
@@ -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):
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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()
|
||||
Reference in New Issue
Block a user