Guard lookups when adding/deleting cameras at runtime (#23994)

* Guard object processor queue handlers against unknown cameras

* Skip embeddings post processing for removed cameras

* End review segments for removed cameras

* Drop queued autotracker moves for removed cameras

* Release tracked event thumbnails when skipping a removed camera

* Add locked accessors for camera states

* Read camera states through the processor accessors

* Guard output and recording paths against cameras not yet known

* Resolve camera state once in ONVIF, notification, and transcription paths
This commit is contained in:
Josh Hawkins
2026-09-12 07:30:04 -06:00
committed by Nicolas Mowen
parent d6a18e79aa
commit d78f6a7a98
16 changed files with 463 additions and 128 deletions
+60 -21
View File
@@ -68,6 +68,7 @@ class TrackedObjectProcessor(threading.Thread):
self.tracked_objects_queue = tracked_objects_queue
self.stop_event: MpEvent = stop_event
self.camera_states: dict[str, CameraState] = {}
self.camera_states_lock = threading.Lock()
self.frame_manager = SharedMemoryFrameManager()
self.last_motion_detected: dict[str, float] = {}
self.ptz_autotracker_thread = ptz_autotracker_thread
@@ -236,7 +237,9 @@ class TrackedObjectProcessor(threading.Thread):
camera_state.on("end", end)
camera_state.on("snapshot", snapshot)
camera_state.on("camera_activity", camera_activity)
self.camera_states[camera] = camera_state
with self.camera_states_lock:
self.camera_states[camera] = camera_state
def should_save_snapshot(self, camera: str, obj: TrackedObject) -> bool:
if obj.false_positive:
@@ -324,9 +327,22 @@ class TrackedObjectProcessor(threading.Thread):
# reset the last_motion so redundant `off` commands aren't sent
self.last_motion_detected[camera] = 0
def get_camera_state(self, camera: str) -> CameraState | None:
"""Returns the state for a camera, or None if it does not exist."""
with self.camera_states_lock:
return self.camera_states.get(camera)
def get_camera_states(self) -> list[CameraState]:
"""Returns a snapshot of camera states that is safe to iterate."""
with self.camera_states_lock:
return list(self.camera_states.values())
def get_best(self, camera: str, label: str) -> dict[str, Any]:
# TODO: need a lock here
camera_state = self.camera_states[camera]
camera_state = self.get_camera_state(camera)
if camera_state is None:
return {}
if label in camera_state.best_objects:
best_obj = camera_state.best_objects[label]
@@ -350,17 +366,21 @@ class TrackedObjectProcessor(threading.Thread):
(self.config.birdseye.height * 3 // 2, self.config.birdseye.width),
)
if camera not in self.camera_states:
camera_state = self.get_camera_state(camera)
if camera_state is None:
return None
return self.camera_states[camera].get_current_frame(draw_options)
return camera_state.get_current_frame(draw_options)
def get_current_frame_time(self, camera: str) -> float:
"""Returns the latest frame time for a given camera."""
if camera not in self.camera_states:
camera_state = self.get_camera_state(camera)
if camera_state is None:
return 0.0
return self.camera_states[camera].current_frame_time
return camera_state.current_frame_time
def set_sub_label(
self, event_id: str, sub_label: str | None, score: float | None
@@ -498,14 +518,18 @@ class TrackedObjectProcessor(threading.Thread):
# save the snapshot image
(frame, event_id, camera) = payload
camera_state = self.camera_states.get(camera)
if camera_state is None:
logger.debug("Discarding LPR snapshot for unknown camera %s", camera)
return
img = cv2.imdecode(
np.frombuffer(base64.b64decode(frame), dtype=np.uint8),
cv2.IMREAD_COLOR,
)
self.camera_states[camera].save_manual_event_image(
img, event_id, "license_plate", {}
)
camera_state.save_manual_event_image(img, event_id, "license_plate", {})
def create_manual_event(self, payload: tuple) -> None:
(
@@ -522,13 +546,17 @@ class TrackedObjectProcessor(threading.Thread):
pre_capture,
) = payload
camera_state = self.camera_states.get(camera_name)
if camera_state is None:
logger.debug("Discarding manual event for unknown camera %s", camera_name)
return
# save the snapshot image
self.camera_states[camera_name].save_manual_event_image(
None, event_id, label, draw
)
camera_state.save_manual_event_image(None, event_id, label, draw)
end_time = frame_time + duration if duration is not None else None
start_time = (
frame_time - self.config.cameras[camera_name].record.event_pre_capture
frame_time - camera_state.camera_config.record.event_pre_capture
if pre_capture is None
else frame_time - pre_capture
)
@@ -548,7 +576,7 @@ class TrackedObjectProcessor(threading.Thread):
"camera": camera_name,
"start_time": start_time,
"end_time": end_time,
"has_clip": self.config.cameras[camera_name].record.enabled
"has_clip": camera_state.camera_config.record.enabled
and include_recording,
"has_snapshot": True,
"snapshot_clean": True,
@@ -591,6 +619,12 @@ class TrackedObjectProcessor(threading.Thread):
plate,
) = payload
camera_state = self.camera_states.get(camera_name)
if camera_state is None:
logger.debug("Discarding LPR event for unknown camera %s", camera_name)
return
# send event to event maintainer
self.event_sender.publish(
(
@@ -605,9 +639,9 @@ class TrackedObjectProcessor(threading.Thread):
"score": score,
"camera": camera_name,
"start_time": frame_time
- self.config.cameras[camera_name].record.event_pre_capture,
- camera_state.camera_config.record.event_pre_capture,
"end_time": None,
"has_clip": self.config.cameras[camera_name].record.enabled
"has_clip": camera_state.camera_config.record.enabled
and include_recording,
"has_snapshot": True,
"snapshot_clean": True,
@@ -699,7 +733,10 @@ class TrackedObjectProcessor(threading.Thread):
continue
camera_state.shutdown()
self.camera_states.pop(camera)
with self.camera_states_lock:
self.camera_states.pop(camera)
self.camera_activity.pop(camera, None)
self.last_motion_detected.pop(camera, None)
@@ -715,8 +752,6 @@ class TrackedObjectProcessor(threading.Thread):
if camera_state is None:
continue
camera_state = self.camera_states[camera]
if camera_state.prev_enabled and not current_enabled:
logger.debug(f"Not processing objects for disabled camera {camera}")
self.force_end_all_events(camera, camera_state)
@@ -812,7 +847,11 @@ class TrackedObjectProcessor(threading.Thread):
break
event_id, camera, _ = update
self.camera_states[camera].finished(event_id)
camera_state = self.camera_states.get(camera)
# the camera may have been removed while its event was pending
if camera_state is not None:
camera_state.finished(event_id)
# shut down camera states
for state in self.camera_states.values():