diff --git a/docs/docs/integrations/mqtt.md b/docs/docs/integrations/mqtt.md index 002feafe3d..6363021188 100644 --- a/docs/docs/integrations/mqtt.md +++ b/docs/docs/integrations/mqtt.md @@ -304,7 +304,7 @@ Topic with current state of notifications. Published values are `ON` and `OFF`. ### `frigate//status/` -Publishes the current health status of each role that is enabled (`audio`, `detect`, `record`). Possible values are: +Publishes the current health status of each role that is enabled (`audio`, `detect`, `record`, `record_sub`). `record_sub` is only published for cameras with [sub stream recording](/configuration/record#sub-stream-recording) enabled, and is tracked separately from `record` so a healthy main stream can't hide a stalled sub stream. Possible values are: - `online`: Stream is running and being processed - `offline`: Stream is offline and is being restarted diff --git a/frigate/comms/recordings_updater.py b/frigate/comms/recordings_updater.py index 249c2f6076..88645dd648 100644 --- a/frigate/comms/recordings_updater.py +++ b/frigate/comms/recordings_updater.py @@ -18,7 +18,10 @@ class RecordingsDataTypeEnum(str, Enum): class RecordingsDataPublisher(Publisher[Any]): - """Publishes latest recording data.""" + """Publishes latest recording data. + + Payloads are (camera, stream_type, timestamp, cache_path) on every topic. + """ topic_base = "recordings/" diff --git a/frigate/config/camera/record.py b/frigate/config/camera/record.py index 341299113d..d9453257d2 100644 --- a/frigate/config/camera/record.py +++ b/frigate/config/camera/record.py @@ -2,7 +2,7 @@ from enum import Enum from pydantic import Field -from frigate.const import MAX_PRE_CAPTURE +from frigate.const import MAX_PRE_CAPTURE, STREAM_TYPE_SUB from frigate.review.types import SeverityEnum from ..base import FrigateBaseModel @@ -191,6 +191,13 @@ class RecordConfig(FrigateBaseModel): description="Indicates whether recording was enabled in the original static configuration.", ) + def stream_enabled(self, stream_type: str) -> bool: + """Whether the given record stream type should currently be recording.""" + if stream_type == STREAM_TYPE_SUB: + return self.enabled and self.sub.enabled + + return self.enabled + @property def effective_alert_days(self) -> float: """Alert retention extended to the sub stream window when sub is enabled. diff --git a/frigate/const.py b/frigate/const.py index c0cf560b99..899324adcb 100644 --- a/frigate/const.py +++ b/frigate/const.py @@ -28,6 +28,9 @@ REDACTED_CREDENTIAL_SENTINEL = "__FRIGATE_SAVED_CREDENTIAL__" STREAM_TYPE_MAIN = "main" STREAM_TYPE_SUB = "sub" SUB_CACHE_TAG = "@sub" +RECORD_STREAM_TYPES = (STREAM_TYPE_MAIN, STREAM_TYPE_SUB) +ROLE_TO_STREAM_TYPE = {"record": STREAM_TYPE_MAIN, "record_sub": STREAM_TYPE_SUB} +STREAM_TYPE_TO_ROLE = {v: k for k, v in ROLE_TO_STREAM_TYPE.items()} # Attribute & Object constants diff --git a/frigate/embeddings/maintainer.py b/frigate/embeddings/maintainer.py index 03309a7ad4..17b09383d9 100644 --- a/frigate/embeddings/maintainer.py +++ b/frigate/embeddings/maintainer.py @@ -714,7 +714,9 @@ class EmbeddingMaintainer(threading.Thread): topic = str(raw_topic) if topic.endswith(RecordingsDataTypeEnum.saved.value): - camera, recordings_available_through_timestamp, _ = payload + camera, _stream_type, recordings_available_through_timestamp, _ = ( + payload + ) self.recordings_available_through[camera] = ( recordings_available_through_timestamp diff --git a/frigate/record/maintainer.py b/frigate/record/maintainer.py index c347566b07..5d57ed9c30 100644 --- a/frigate/record/maintainer.py +++ b/frigate/record/maintainer.py @@ -241,8 +241,8 @@ class RecordingMaintainer(threading.Thread): and not d.startswith("preview_") ] - # publish newest cached segment per camera (including in use files) - newest_cache_segments: dict[str, dict[str, Any]] = {} + # publish newest cached segment per camera stream (including in use files) + newest_cache_segments: dict[tuple[str, str], dict[str, Any]] = {} for cache in cache_files: cache_path = os.path.join(CACHE_DIR, cache) basename = os.path.splitext(cache)[0] @@ -254,38 +254,42 @@ class RecordingMaintainer(threading.Thread): continue camera, stream_type, date = parsed - # this topic feeds main-stream health/sync consumers only - if stream_type == STREAM_TYPE_SUB: - continue - start_time = datetime.datetime.strptime( date, CACHE_SEGMENT_FORMAT ).astimezone(datetime.UTC) + key = (camera, stream_type) if ( - camera not in newest_cache_segments - or start_time > newest_cache_segments[camera]["start_time"] + key not in newest_cache_segments + or start_time > newest_cache_segments[key]["start_time"] ): - newest_cache_segments[camera] = { + newest_cache_segments[key] = { "start_time": start_time, "cache_path": cache_path, } - for camera, newest in newest_cache_segments.items(): + for (camera, stream_type), newest in newest_cache_segments.items(): self.recordings_publisher.publish( ( camera, + stream_type, newest["start_time"].timestamp(), newest["cache_path"], ), RecordingsDataTypeEnum.latest.value, ) - # publish None for cameras with no cache files (but only if we know the camera exists) - for camera_name in self.config.cameras: - if camera_name not in newest_cache_segments: - self.recordings_publisher.publish( - (camera_name, None, None), - RecordingsDataTypeEnum.latest.value, - ) + # publish None for streams with no cache files (but only if we know the camera exists) + for camera_name, camera_cfg in self.config.cameras.items(): + stream_types = [STREAM_TYPE_MAIN] + + if camera_cfg.record.sub.enabled: + stream_types.append(STREAM_TYPE_SUB) + + for stream_type in stream_types: + if (camera_name, stream_type) not in newest_cache_segments: + self.recordings_publisher.publish( + (camera_name, stream_type, None, None), + RecordingsDataTypeEnum.latest.value, + ) files_in_use = [] for process in psutil.process_iter(): @@ -447,6 +451,7 @@ class RecordingMaintainer(threading.Thread): self.recordings_publisher.publish( ( camera, + stream_type, recordings[0]["start_time"].timestamp() if camera_cfg and camera_cfg.record.enabled else None, @@ -542,11 +547,10 @@ class RecordingMaintainer(threading.Thread): logger.warning( f"Invalid or missing video stream in segment {cache_path}. Discarding." ) - if stream_type == STREAM_TYPE_MAIN: - self.recordings_publisher.publish( - (camera, start_time.timestamp(), cache_path), - RecordingsDataTypeEnum.invalid.value, - ) + self.recordings_publisher.publish( + (camera, stream_type, start_time.timestamp(), cache_path), + RecordingsDataTypeEnum.invalid.value, + ) self.drop_segment(cache_path) return None @@ -584,20 +588,18 @@ class RecordingMaintainer(threading.Thread): logger.warning(f"Failed to probe corrupt segment {cache_path}") logger.warning(f"Discarding a corrupt recording segment: {cache_path}") - if stream_type == STREAM_TYPE_MAIN: - self.recordings_publisher.publish( - (camera, start_time.timestamp(), cache_path), - RecordingsDataTypeEnum.invalid.value, - ) + self.recordings_publisher.publish( + (camera, stream_type, start_time.timestamp(), cache_path), + RecordingsDataTypeEnum.invalid.value, + ) self.drop_segment(cache_path) return None # this segment has a valid duration and has video data, so publish an update - if stream_type == STREAM_TYPE_MAIN: - self.recordings_publisher.publish( - (camera, start_time.timestamp(), cache_path), - RecordingsDataTypeEnum.valid.value, - ) + self.recordings_publisher.publish( + (camera, stream_type, start_time.timestamp(), cache_path), + RecordingsDataTypeEnum.valid.value, + ) record_config = self.config.cameras[camera].record diff --git a/frigate/test/test_camera_watchdog.py b/frigate/test/test_camera_watchdog.py new file mode 100644 index 0000000000..3ae9ca8332 --- /dev/null +++ b/frigate/test/test_camera_watchdog.py @@ -0,0 +1,142 @@ +"""Tests for per stream recording health tracking in the camera watchdog.""" + +import unittest +from datetime import UTC, datetime, timedelta +from unittest.mock import MagicMock, patch + +from frigate.config import FrigateConfig +from frigate.const import STREAM_TYPE_MAIN, STREAM_TYPE_SUB +from frigate.video.ffmpeg import CameraWatchdog + + +class TestCameraWatchdogStreamHealth(unittest.TestCase): + def _build_watchdog( + self, sub_enabled: bool = True, output_args: dict | None = None + ) -> CameraWatchdog: + config = FrigateConfig( + **{ + "mqtt": {"host": "mqtt"}, + "cameras": { + "front_door": { + "ffmpeg": { + "output_args": output_args or {}, + "inputs": [ + { + "path": "rtsp://10.0.0.1:554/video", + "roles": ["record"], + }, + { + "path": "rtsp://10.0.0.1:554/video2", + "roles": ["detect", "record_sub"], + }, + ], + }, + "record": { + "enabled": True, + "sub": {"enabled": sub_enabled}, + }, + } + }, + } + ) + camera_config = config.cameras["front_door"] + + with ( + patch("frigate.video.ffmpeg.LogPipe"), + patch("frigate.video.ffmpeg.InterProcessRequestor"), + patch("frigate.video.ffmpeg.RecordingsDataSubscriber"), + patch("frigate.video.ffmpeg.CameraConfigUpdateSubscriber"), + ): + watchdog = CameraWatchdog( + camera_config, + 1, + MagicMock(), + MagicMock(), + MagicMock(), + MagicMock(), + MagicMock(), + MagicMock(), + MagicMock(), + MagicMock(), + ) + + watchdog.requestor = MagicMock() + return watchdog + + def test_stale_sub_does_not_mark_main_stale(self): + watchdog = self._build_watchdog() + now = datetime.now().astimezone(UTC) + stale = (now - timedelta(hours=1)).timestamp() + + watchdog.latest_cache_segment_time[STREAM_TYPE_MAIN] = now.timestamp() + watchdog.latest_valid_segment_time[STREAM_TYPE_MAIN] = now.timestamp() + watchdog.latest_cache_segment_time[STREAM_TYPE_SUB] = stale + watchdog.latest_valid_segment_time[STREAM_TYPE_SUB] = stale + + assert watchdog._stream_staleness(STREAM_TYPE_MAIN, now) is None + assert watchdog._stream_staleness(STREAM_TYPE_SUB, now) is not None + + def test_stale_main_does_not_mark_sub_stale(self): + watchdog = self._build_watchdog() + now = datetime.now().astimezone(UTC) + stale = (now - timedelta(hours=1)).timestamp() + + watchdog.latest_cache_segment_time[STREAM_TYPE_MAIN] = stale + watchdog.latest_valid_segment_time[STREAM_TYPE_MAIN] = stale + watchdog.latest_cache_segment_time[STREAM_TYPE_SUB] = now.timestamp() + watchdog.latest_valid_segment_time[STREAM_TYPE_SUB] = now.timestamp() + + assert watchdog._stream_staleness(STREAM_TYPE_MAIN, now) is not None + assert watchdog._stream_staleness(STREAM_TYPE_SUB, now) is None + + def test_grace_period_suppresses_staleness(self): + watchdog = self._build_watchdog() + now = datetime.now().astimezone(UTC) + watchdog.record_enable_time = now - timedelta(seconds=10) + watchdog.latest_cache_segment_time[STREAM_TYPE_SUB] = ( + now - timedelta(hours=1) + ).timestamp() + + assert watchdog._stream_staleness(STREAM_TYPE_SUB, now) is None + + def test_status_goes_to_the_matching_role_topic(self): + watchdog = self._build_watchdog() + + watchdog._send_record_status(STREAM_TYPE_MAIN, "online", 100.0) + watchdog._send_record_status(STREAM_TYPE_SUB, "offline", 100.0) + + watchdog.requestor.send_data.assert_any_call( + "front_door/status/record", "online" + ) + watchdog.requestor.send_data.assert_any_call( + "front_door/status/record_sub", "offline" + ) + + def test_status_is_cached_per_stream(self): + watchdog = self._build_watchdog() + + watchdog._send_record_status(STREAM_TYPE_MAIN, "online", 100.0) + watchdog._send_record_status(STREAM_TYPE_SUB, "online", 100.0) + watchdog._send_record_status(STREAM_TYPE_MAIN, "online", 100.0) + + assert watchdog.requestor.send_data.call_count == 2 + + def test_recorded_streams_follows_config(self): + watchdog = self._build_watchdog() + assert watchdog._recorded_streams(["record"]) == [STREAM_TYPE_MAIN] + assert watchdog._recorded_streams(["detect", "record_sub"]) == [STREAM_TYPE_SUB] + assert watchdog._recorded_streams(["detect"]) == [] + + disabled = self._build_watchdog(sub_enabled=False) + assert disabled._recorded_streams(["detect", "record_sub"]) == [] + + def test_stale_threshold_follows_each_stream_segment_time(self): + watchdog = self._build_watchdog( + output_args={ + "record": "-f segment -segment_time 10 -c copy", + "record_sub": "-f segment -segment_time 60 -c copy", + } + ) + + assert watchdog.record_stale_threshold[STREAM_TYPE_MAIN] == 120 + assert watchdog.record_stale_threshold[STREAM_TYPE_SUB] == 150 diff --git a/frigate/util/builtin.py b/frigate/util/builtin.py index 6bf5c6f380..fbc53e2cd9 100644 --- a/frigate/util/builtin.py +++ b/frigate/util/builtin.py @@ -20,7 +20,12 @@ from typing import TYPE_CHECKING, Any import numpy as np from ruamel.yaml import YAML -from frigate.const import REGEX_HTTP_CAMERA_USER_PASS, REGEX_RTSP_CAMERA_USER_PASS +from frigate.const import ( + REGEX_HTTP_CAMERA_USER_PASS, + REGEX_RTSP_CAMERA_USER_PASS, + STREAM_TYPE_MAIN, + STREAM_TYPE_SUB, +) if TYPE_CHECKING: from frigate.config import CameraConfig @@ -137,9 +142,16 @@ def get_ffmpeg_arg_list(arg: Any) -> list: DEFAULT_RECORD_SEGMENT_TIME = 10 -def get_record_segment_time(config: "CameraConfig") -> int: - """Extract -segment_time from the camera's record output args.""" - record_args = get_ffmpeg_arg_list(config.ffmpeg.output_args.record) +def get_record_segment_time( + config: "CameraConfig", stream_type: str = STREAM_TYPE_MAIN +) -> int: + """Extract -segment_time from the camera's record output args for a stream.""" + output_args = ( + config.ffmpeg.output_args.effective_record_sub + if stream_type == STREAM_TYPE_SUB + else config.ffmpeg.output_args.record + ) + record_args = get_ffmpeg_arg_list(output_args) if record_args and record_args[0].startswith("preset"): return DEFAULT_RECORD_SEGMENT_TIME diff --git a/frigate/video/ffmpeg.py b/frigate/video/ffmpeg.py index d4645efc78..1877ec2bbf 100644 --- a/frigate/video/ffmpeg.py +++ b/frigate/video/ffmpeg.py @@ -5,7 +5,7 @@ import queue import subprocess as sp import threading import time -from collections import deque +from collections import defaultdict, deque from datetime import UTC, datetime, timedelta from multiprocessing import Queue, Value from multiprocessing.synchronize import Event as MpEvent @@ -22,7 +22,13 @@ from frigate.config.camera.updater import ( CameraConfigUpdateEnum, CameraConfigUpdateSubscriber, ) -from frigate.const import PROCESS_PRIORITY_HIGH +from frigate.const import ( + PROCESS_PRIORITY_HIGH, + RECORD_STREAM_TYPES, + ROLE_TO_STREAM_TYPE, + STREAM_TYPE_SUB, + STREAM_TYPE_TO_ROLE, +) from frigate.log import LogPipe from frigate.util.builtin import EventsPerSecond, get_record_segment_time from frigate.util.ffmpeg import start_or_restart_ffmpeg, stop_ffmpeg @@ -150,16 +156,25 @@ class CameraWatchdog(threading.Thread): self.was_record_sub_enabled = self.config.record.sub.enabled self.segment_subscriber = RecordingsDataSubscriber(RecordingsDataTypeEnum.all) - self.latest_valid_segment_time: float = 0 - self.latest_invalid_segment_time: float = 0 - self.latest_cache_segment_time: float = 0 + self.latest_valid_segment_time: dict[str, float] = defaultdict(float) + self.latest_invalid_segment_time: dict[str, float] = defaultdict(float) + self.latest_cache_segment_time: dict[str, float] = defaultdict(float) self.record_enable_time: datetime | None = None # `valid` segments are published with the segment's start time, so the # gap between consecutive publishes can reach 2 * segment_time. Pad the # staleness threshold so it's never tighter than that worst case. - segment_time = get_record_segment_time(self.config) - self.record_stale_threshold = max(120, 2 * segment_time + 30) + self.record_stale_threshold: dict[str, int] = { + stream_type: max( + 120, 2 * get_record_segment_time(self.config, stream_type) + 30 + ) + for stream_type in RECORD_STREAM_TYPES + } + + # the sub stream usually shares its input, and therefore its ffmpeg + # process, with detect, so it isn't in ffmpeg_other_processes and needs + # its own staleness check + self.detect_process_records_sub = False # Stall tracking (based on last processed frame) self._stall_timestamps: deque[float] = deque() @@ -167,7 +182,7 @@ class CameraWatchdog(threading.Thread): # Status caching to reduce message volume self._last_detect_status: str | None = None - self._last_record_status: str | None = None + self._last_record_status: dict[str, str] = {} self._last_status_update_time: float = 0.0 def _send_detect_status(self, status: str, now: float) -> None: @@ -180,16 +195,69 @@ class CameraWatchdog(threading.Thread): self._last_detect_status = status self._last_status_update_time = now - def _send_record_status(self, status: str, now: float) -> None: - """Send record status only if changed or retry_interval has elapsed.""" + def _send_record_status(self, stream_type: str, status: str, now: float) -> None: + """Send a record stream's status only if changed or retry_interval has elapsed.""" if ( - status != self._last_record_status + status != self._last_record_status.get(stream_type) or (now - self._last_status_update_time) >= self.sleeptime ): - self.requestor.send_data(f"{self.config.name}/status/record", status) - self._last_record_status = status + self.requestor.send_data( + f"{self.config.name}/status/{STREAM_TYPE_TO_ROLE[stream_type]}", status + ) + self._last_record_status[stream_type] = status self._last_status_update_time = now + def _reset_segment_times(self) -> None: + self.latest_valid_segment_time.clear() + self.latest_invalid_segment_time.clear() + self.latest_cache_segment_time.clear() + + def _stream_staleness(self, stream_type: str, now_utc: datetime) -> str | None: + """Return why the stream's segments are stale, or None if they're healthy.""" + # Grace period: 90 seconds allows time for ffmpeg to start and create + # the first segment + in_grace_period = self.record_enable_time is not None and ( + now_utc - self.record_enable_time + ) < timedelta(seconds=90) + + if in_grace_period: + return None + + latest_cache = self.latest_cache_segment_time[stream_type] + latest_valid = self.latest_valid_segment_time[stream_type] + latest_invalid = self.latest_invalid_segment_time[stream_type] + + def as_dt(timestamp: float) -> datetime: + if timestamp > 0: + return datetime.fromtimestamp(timestamp, tz=UTC) + + return now_utc - timedelta(seconds=1) + + stale_window = timedelta(seconds=self.record_stale_threshold[stream_type]) + + if now_utc > (as_dt(latest_cache) + stale_window): + return "No new recording segments were created" + + if now_utc > (as_dt(latest_valid) + stale_window): + return "No new valid recording segments were created" + + if ( + latest_invalid > 0 + and now_utc > (as_dt(latest_invalid) + stale_window) + and latest_valid <= latest_invalid + ): + return "No valid segments created since last invalid segment" + + return None + + def _recorded_streams(self, roles: list[Any]) -> list[str]: + """Record stream types the given roles cover that are currently recording.""" + return [ + stream_type + for role, stream_type in ROLE_TO_STREAM_TYPE.items() + if role in roles and self.config.record.stream_enabled(stream_type) + ] + def _check_config_updates(self) -> dict[str, list[str]]: """Check for config updates and return the update dict.""" return self.config_subscriber.check_for_updates() @@ -267,9 +335,7 @@ class CameraWatchdog(threading.Thread): ) self.stop_all_ffmpeg() self.start_all_ffmpeg() - self.latest_valid_segment_time = 0 - self.latest_invalid_segment_time = 0 - self.latest_cache_segment_time = 0 + self._reset_segment_times() self.record_enable_time = datetime.now().astimezone(UTC) last_restart_time = datetime.now().timestamp() continue @@ -281,9 +347,7 @@ class CameraWatchdog(threading.Thread): self.start_all_ffmpeg() # reset all timestamps and record the enable time for grace period - self.latest_valid_segment_time = 0 - self.latest_invalid_segment_time = 0 - self.latest_cache_segment_time = 0 + self._reset_segment_times() self.record_enable_time = datetime.now().astimezone(UTC) else: self.logger.debug(f"Disabling camera {self.config.name}") @@ -293,7 +357,8 @@ class CameraWatchdog(threading.Thread): # update camera status now = datetime.now().timestamp() self._send_detect_status("disabled", now) - self._send_record_status("disabled", now) + for stream_type in RECORD_STREAM_TYPES: + self._send_record_status(stream_type, "disabled", now) self.was_enabled = enabled continue @@ -305,9 +370,7 @@ class CameraWatchdog(threading.Thread): ) self.stop_all_ffmpeg() self.start_all_ffmpeg() - self.latest_valid_segment_time = 0 - self.latest_invalid_segment_time = 0 - self.latest_cache_segment_time = 0 + self._reset_segment_times() self.record_enable_time = datetime.now().astimezone(UTC) last_restart_time = datetime.now().timestamp() self.was_record_enabled_in_config = record_enabled_in_config @@ -323,9 +386,7 @@ class CameraWatchdog(threading.Thread): ) self.stop_all_ffmpeg() self.start_all_ffmpeg() - self.latest_valid_segment_time = 0 - self.latest_invalid_segment_time = 0 - self.latest_cache_segment_time = 0 + self._reset_segment_times() self.record_enable_time = datetime.now().astimezone(UTC) last_restart_time = datetime.now().timestamp() self.was_record_sub_enabled = record_sub_enabled @@ -343,26 +404,25 @@ class CameraWatchdog(threading.Thread): raw_topic, payload = update if raw_topic and payload: topic = str(raw_topic) - camera, segment_time, _ = payload + camera, stream_type, segment_time, _ = payload if camera != self.config.name: continue if topic.endswith(RecordingsDataTypeEnum.invalid.value): self.logger.warning( - f"Invalid recording segment detected for {camera} at {segment_time}" + f"Invalid recording segment detected for {camera} ({stream_type}) at {segment_time}" ) - self.latest_invalid_segment_time = segment_time + self.latest_invalid_segment_time[stream_type] = segment_time elif topic.endswith(RecordingsDataTypeEnum.valid.value): self.logger.debug( - f"Latest valid recording segment time on {camera}: {segment_time}" + f"Latest valid recording segment time on {camera} ({stream_type}): {segment_time}" ) - self.latest_valid_segment_time = segment_time + self.latest_valid_segment_time[stream_type] = segment_time elif topic.endswith(RecordingsDataTypeEnum.latest.value): - if segment_time is not None: - self.latest_cache_segment_time = segment_time - else: - self.latest_cache_segment_time = 0 + self.latest_cache_segment_time[stream_type] = ( + segment_time if segment_time is not None else 0 + ) now = datetime.now().timestamp() @@ -409,63 +469,26 @@ class CameraWatchdog(threading.Thread): for p in self.ffmpeg_other_processes: poll = p["process"].poll() - if self.config.record.enabled and "record" in p["roles"]: + recorded_streams = self._recorded_streams(p["roles"]) + + if recorded_streams: now_utc = datetime.now().astimezone(UTC) - # Check if we're within the grace period after enabling recording - # Grace period: 90 seconds allows time for ffmpeg to start and create first segment - in_grace_period = self.record_enable_time is not None and ( - now_utc - self.record_enable_time - ) < timedelta(seconds=90) + # ensure segments are still being created and that they have + # valid video data. each stream is tracked separately so a + # healthy one can't mask a stalled one. + stale_stream = None + stale_reason = None + for stream_type in recorded_streams: + stale_reason = self._stream_staleness(stream_type, now_utc) - latest_cache_dt = ( - datetime.fromtimestamp(self.latest_cache_segment_time, tz=UTC) - if self.latest_cache_segment_time > 0 - else now_utc - timedelta(seconds=1) - ) - - latest_valid_dt = ( - datetime.fromtimestamp(self.latest_valid_segment_time, tz=UTC) - if self.latest_valid_segment_time > 0 - else now_utc - timedelta(seconds=1) - ) - - latest_invalid_dt = ( - datetime.fromtimestamp(self.latest_invalid_segment_time, tz=UTC) - if self.latest_invalid_segment_time > 0 - else now_utc - timedelta(seconds=1) - ) - - # ensure segments are still being created and that they have valid video data - # Skip checks during grace period to allow segments to start being created - stale_window = timedelta(seconds=self.record_stale_threshold) - cache_stale = not in_grace_period and now_utc > ( - latest_cache_dt + stale_window - ) - valid_stale = not in_grace_period and now_utc > ( - latest_valid_dt + stale_window - ) - invalid_stale_condition = ( - self.latest_invalid_segment_time > 0 - and not in_grace_period - and now_utc > (latest_invalid_dt + stale_window) - and self.latest_valid_segment_time - <= self.latest_invalid_segment_time - ) - invalid_stale = invalid_stale_condition - - if cache_stale or valid_stale or invalid_stale: - if cache_stale: - reason = "No new recording segments were created" - elif valid_stale: - reason = "No new valid recording segments were created" - else: # invalid_stale - reason = ( - "No valid segments created since last invalid segment" - ) + if stale_reason is not None: + stale_stream = stream_type + break + if stale_stream is not None: self.logger.error( - f"{reason} for {self.config.name} in the last {self.record_stale_threshold}s. Restarting the ffmpeg record process..." + f"{stale_reason} for {self.config.name} ({stale_stream}) in the last {self.record_stale_threshold[stale_stream]}s. Restarting the ffmpeg record process..." ) p["process"] = start_or_restart_ffmpeg( p["cmd"], @@ -481,8 +504,13 @@ class CameraWatchdog(threading.Thread): continue else: - self._send_record_status("online", now) - p["latest_segment_time"] = self.latest_cache_segment_time + for stream_type in recorded_streams: + self._send_record_status(stream_type, "online", now) + + p["latest_segment_time"] = max( + self.latest_cache_segment_time[stream_type] + for stream_type in recorded_streams + ) if poll is None: continue @@ -497,6 +525,25 @@ class CameraWatchdog(threading.Thread): p["cmd"], self.logger, p["logpipe"], ffmpeg_process=p["process"] ) + if ( + self.detect_process_records_sub + and self.config.record.stream_enabled(STREAM_TYPE_SUB) + and self.capture_thread is not None + and self.capture_thread.is_alive() + ): + now_utc = datetime.now().astimezone(UTC) + stale_reason = self._stream_staleness(STREAM_TYPE_SUB, now_utc) + + if stale_reason is None: + self._send_record_status(STREAM_TYPE_SUB, "online", now) + elif can_restart: + self.logger.error( + f"{stale_reason} for {self.config.name} (sub, shared with detect) in the last {self.record_stale_threshold[STREAM_TYPE_SUB]}s. Restarting ffmpeg..." + ) + self._send_record_status(STREAM_TYPE_SUB, "offline", now) + self.reset_capture_thread() + last_restart_time = now + # Prune expired reconnect timestamps now = datetime.now().timestamp() while ( @@ -539,9 +586,9 @@ class CameraWatchdog(threading.Thread): self.segment_subscriber.stop() def start_ffmpeg_detect(self): - ffmpeg_cmd = [ - c["cmd"] for c in self.config.ffmpeg_cmds if "detect" in c["roles"] - ][0] + detect_cmd = [c for c in self.config.ffmpeg_cmds if "detect" in c["roles"]][0] + ffmpeg_cmd = detect_cmd["cmd"] + self.detect_process_records_sub = "record_sub" in detect_cmd["roles"] self.ffmpeg_detect_process = start_or_restart_ffmpeg( ffmpeg_cmd, self.logger, self.logpipe, self.frame_size )