watch sub stream recording health separately from main

This commit is contained in:
Josh Hawkins
2026-08-23 15:48:25 -05:00
parent 27291ef9c8
commit 864772440d
9 changed files with 350 additions and 132 deletions
+1 -1
View File
@@ -304,7 +304,7 @@ Topic with current state of notifications. Published values are `ON` and `OFF`.
### `frigate/<camera_name>/status/<role>` ### `frigate/<camera_name>/status/<role>`
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 - `online`: Stream is running and being processed
- `offline`: Stream is offline and is being restarted - `offline`: Stream is offline and is being restarted
+4 -1
View File
@@ -18,7 +18,10 @@ class RecordingsDataTypeEnum(str, Enum):
class RecordingsDataPublisher(Publisher[Any]): 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/" topic_base = "recordings/"
+8 -1
View File
@@ -2,7 +2,7 @@ from enum import Enum
from pydantic import Field 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 frigate.review.types import SeverityEnum
from ..base import FrigateBaseModel from ..base import FrigateBaseModel
@@ -191,6 +191,13 @@ class RecordConfig(FrigateBaseModel):
description="Indicates whether recording was enabled in the original static configuration.", 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 @property
def effective_alert_days(self) -> float: def effective_alert_days(self) -> float:
"""Alert retention extended to the sub stream window when sub is enabled. """Alert retention extended to the sub stream window when sub is enabled.
+3
View File
@@ -28,6 +28,9 @@ REDACTED_CREDENTIAL_SENTINEL = "__FRIGATE_SAVED_CREDENTIAL__"
STREAM_TYPE_MAIN = "main" STREAM_TYPE_MAIN = "main"
STREAM_TYPE_SUB = "sub" STREAM_TYPE_SUB = "sub"
SUB_CACHE_TAG = "@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 # Attribute & Object constants
+3 -1
View File
@@ -714,7 +714,9 @@ class EmbeddingMaintainer(threading.Thread):
topic = str(raw_topic) topic = str(raw_topic)
if topic.endswith(RecordingsDataTypeEnum.saved.value): 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] = ( self.recordings_available_through[camera] = (
recordings_available_through_timestamp recordings_available_through_timestamp
+34 -32
View File
@@ -241,8 +241,8 @@ class RecordingMaintainer(threading.Thread):
and not d.startswith("preview_") and not d.startswith("preview_")
] ]
# publish newest cached segment per camera (including in use files) # publish newest cached segment per camera stream (including in use files)
newest_cache_segments: dict[str, dict[str, Any]] = {} newest_cache_segments: dict[tuple[str, str], dict[str, Any]] = {}
for cache in cache_files: for cache in cache_files:
cache_path = os.path.join(CACHE_DIR, cache) cache_path = os.path.join(CACHE_DIR, cache)
basename = os.path.splitext(cache)[0] basename = os.path.splitext(cache)[0]
@@ -254,38 +254,42 @@ class RecordingMaintainer(threading.Thread):
continue continue
camera, stream_type, date = parsed 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( start_time = datetime.datetime.strptime(
date, CACHE_SEGMENT_FORMAT date, CACHE_SEGMENT_FORMAT
).astimezone(datetime.UTC) ).astimezone(datetime.UTC)
key = (camera, stream_type)
if ( if (
camera not in newest_cache_segments key not in newest_cache_segments
or start_time > newest_cache_segments[camera]["start_time"] or start_time > newest_cache_segments[key]["start_time"]
): ):
newest_cache_segments[camera] = { newest_cache_segments[key] = {
"start_time": start_time, "start_time": start_time,
"cache_path": cache_path, "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( self.recordings_publisher.publish(
( (
camera, camera,
stream_type,
newest["start_time"].timestamp(), newest["start_time"].timestamp(),
newest["cache_path"], newest["cache_path"],
), ),
RecordingsDataTypeEnum.latest.value, RecordingsDataTypeEnum.latest.value,
) )
# publish None for cameras with no cache files (but only if we know the camera exists) # publish None for streams with no cache files (but only if we know the camera exists)
for camera_name in self.config.cameras: for camera_name, camera_cfg in self.config.cameras.items():
if camera_name not in newest_cache_segments: stream_types = [STREAM_TYPE_MAIN]
self.recordings_publisher.publish(
(camera_name, None, None), if camera_cfg.record.sub.enabled:
RecordingsDataTypeEnum.latest.value, 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 = [] files_in_use = []
for process in psutil.process_iter(): for process in psutil.process_iter():
@@ -447,6 +451,7 @@ class RecordingMaintainer(threading.Thread):
self.recordings_publisher.publish( self.recordings_publisher.publish(
( (
camera, camera,
stream_type,
recordings[0]["start_time"].timestamp() recordings[0]["start_time"].timestamp()
if camera_cfg and camera_cfg.record.enabled if camera_cfg and camera_cfg.record.enabled
else None, else None,
@@ -542,11 +547,10 @@ class RecordingMaintainer(threading.Thread):
logger.warning( logger.warning(
f"Invalid or missing video stream in segment {cache_path}. Discarding." f"Invalid or missing video stream in segment {cache_path}. Discarding."
) )
if stream_type == STREAM_TYPE_MAIN: self.recordings_publisher.publish(
self.recordings_publisher.publish( (camera, stream_type, start_time.timestamp(), cache_path),
(camera, start_time.timestamp(), cache_path), RecordingsDataTypeEnum.invalid.value,
RecordingsDataTypeEnum.invalid.value, )
)
self.drop_segment(cache_path) self.drop_segment(cache_path)
return None return None
@@ -584,20 +588,18 @@ class RecordingMaintainer(threading.Thread):
logger.warning(f"Failed to probe corrupt segment {cache_path}") logger.warning(f"Failed to probe corrupt segment {cache_path}")
logger.warning(f"Discarding a corrupt recording segment: {cache_path}") logger.warning(f"Discarding a corrupt recording segment: {cache_path}")
if stream_type == STREAM_TYPE_MAIN: self.recordings_publisher.publish(
self.recordings_publisher.publish( (camera, stream_type, start_time.timestamp(), cache_path),
(camera, start_time.timestamp(), cache_path), RecordingsDataTypeEnum.invalid.value,
RecordingsDataTypeEnum.invalid.value, )
)
self.drop_segment(cache_path) self.drop_segment(cache_path)
return None return None
# this segment has a valid duration and has video data, so publish an update # this segment has a valid duration and has video data, so publish an update
if stream_type == STREAM_TYPE_MAIN: self.recordings_publisher.publish(
self.recordings_publisher.publish( (camera, stream_type, start_time.timestamp(), cache_path),
(camera, start_time.timestamp(), cache_path), RecordingsDataTypeEnum.valid.value,
RecordingsDataTypeEnum.valid.value, )
)
record_config = self.config.cameras[camera].record record_config = self.config.cameras[camera].record
+142
View File
@@ -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
+16 -4
View File
@@ -20,7 +20,12 @@ from typing import TYPE_CHECKING, Any
import numpy as np import numpy as np
from ruamel.yaml import YAML 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: if TYPE_CHECKING:
from frigate.config import CameraConfig from frigate.config import CameraConfig
@@ -137,9 +142,16 @@ def get_ffmpeg_arg_list(arg: Any) -> list:
DEFAULT_RECORD_SEGMENT_TIME = 10 DEFAULT_RECORD_SEGMENT_TIME = 10
def get_record_segment_time(config: "CameraConfig") -> int: def get_record_segment_time(
"""Extract -segment_time from the camera's record output args.""" config: "CameraConfig", stream_type: str = STREAM_TYPE_MAIN
record_args = get_ffmpeg_arg_list(config.ffmpeg.output_args.record) ) -> 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"): if record_args and record_args[0].startswith("preset"):
return DEFAULT_RECORD_SEGMENT_TIME return DEFAULT_RECORD_SEGMENT_TIME
+139 -92
View File
@@ -5,7 +5,7 @@ import queue
import subprocess as sp import subprocess as sp
import threading import threading
import time import time
from collections import deque from collections import defaultdict, deque
from datetime import UTC, datetime, timedelta from datetime import UTC, datetime, timedelta
from multiprocessing import Queue, Value from multiprocessing import Queue, Value
from multiprocessing.synchronize import Event as MpEvent from multiprocessing.synchronize import Event as MpEvent
@@ -22,7 +22,13 @@ from frigate.config.camera.updater import (
CameraConfigUpdateEnum, CameraConfigUpdateEnum,
CameraConfigUpdateSubscriber, 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.log import LogPipe
from frigate.util.builtin import EventsPerSecond, get_record_segment_time from frigate.util.builtin import EventsPerSecond, get_record_segment_time
from frigate.util.ffmpeg import start_or_restart_ffmpeg, stop_ffmpeg 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.was_record_sub_enabled = self.config.record.sub.enabled
self.segment_subscriber = RecordingsDataSubscriber(RecordingsDataTypeEnum.all) self.segment_subscriber = RecordingsDataSubscriber(RecordingsDataTypeEnum.all)
self.latest_valid_segment_time: float = 0 self.latest_valid_segment_time: dict[str, float] = defaultdict(float)
self.latest_invalid_segment_time: float = 0 self.latest_invalid_segment_time: dict[str, float] = defaultdict(float)
self.latest_cache_segment_time: float = 0 self.latest_cache_segment_time: dict[str, float] = defaultdict(float)
self.record_enable_time: datetime | None = None self.record_enable_time: datetime | None = None
# `valid` segments are published with the segment's start time, so the # `valid` segments are published with the segment's start time, so the
# gap between consecutive publishes can reach 2 * segment_time. Pad the # gap between consecutive publishes can reach 2 * segment_time. Pad the
# staleness threshold so it's never tighter than that worst case. # staleness threshold so it's never tighter than that worst case.
segment_time = get_record_segment_time(self.config) self.record_stale_threshold: dict[str, int] = {
self.record_stale_threshold = max(120, 2 * segment_time + 30) 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) # Stall tracking (based on last processed frame)
self._stall_timestamps: deque[float] = deque() self._stall_timestamps: deque[float] = deque()
@@ -167,7 +182,7 @@ class CameraWatchdog(threading.Thread):
# Status caching to reduce message volume # Status caching to reduce message volume
self._last_detect_status: str | None = None 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 self._last_status_update_time: float = 0.0
def _send_detect_status(self, status: str, now: float) -> None: 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_detect_status = status
self._last_status_update_time = now self._last_status_update_time = now
def _send_record_status(self, status: str, now: float) -> None: def _send_record_status(self, stream_type: str, status: str, now: float) -> None:
"""Send record status only if changed or retry_interval has elapsed.""" """Send a record stream's status only if changed or retry_interval has elapsed."""
if ( if (
status != self._last_record_status status != self._last_record_status.get(stream_type)
or (now - self._last_status_update_time) >= self.sleeptime or (now - self._last_status_update_time) >= self.sleeptime
): ):
self.requestor.send_data(f"{self.config.name}/status/record", status) self.requestor.send_data(
self._last_record_status = status 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 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]]: def _check_config_updates(self) -> dict[str, list[str]]:
"""Check for config updates and return the update dict.""" """Check for config updates and return the update dict."""
return self.config_subscriber.check_for_updates() return self.config_subscriber.check_for_updates()
@@ -267,9 +335,7 @@ class CameraWatchdog(threading.Thread):
) )
self.stop_all_ffmpeg() self.stop_all_ffmpeg()
self.start_all_ffmpeg() self.start_all_ffmpeg()
self.latest_valid_segment_time = 0 self._reset_segment_times()
self.latest_invalid_segment_time = 0
self.latest_cache_segment_time = 0
self.record_enable_time = datetime.now().astimezone(UTC) self.record_enable_time = datetime.now().astimezone(UTC)
last_restart_time = datetime.now().timestamp() last_restart_time = datetime.now().timestamp()
continue continue
@@ -281,9 +347,7 @@ class CameraWatchdog(threading.Thread):
self.start_all_ffmpeg() self.start_all_ffmpeg()
# reset all timestamps and record the enable time for grace period # reset all timestamps and record the enable time for grace period
self.latest_valid_segment_time = 0 self._reset_segment_times()
self.latest_invalid_segment_time = 0
self.latest_cache_segment_time = 0
self.record_enable_time = datetime.now().astimezone(UTC) self.record_enable_time = datetime.now().astimezone(UTC)
else: else:
self.logger.debug(f"Disabling camera {self.config.name}") self.logger.debug(f"Disabling camera {self.config.name}")
@@ -293,7 +357,8 @@ class CameraWatchdog(threading.Thread):
# update camera status # update camera status
now = datetime.now().timestamp() now = datetime.now().timestamp()
self._send_detect_status("disabled", now) 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 self.was_enabled = enabled
continue continue
@@ -305,9 +370,7 @@ class CameraWatchdog(threading.Thread):
) )
self.stop_all_ffmpeg() self.stop_all_ffmpeg()
self.start_all_ffmpeg() self.start_all_ffmpeg()
self.latest_valid_segment_time = 0 self._reset_segment_times()
self.latest_invalid_segment_time = 0
self.latest_cache_segment_time = 0
self.record_enable_time = datetime.now().astimezone(UTC) self.record_enable_time = datetime.now().astimezone(UTC)
last_restart_time = datetime.now().timestamp() last_restart_time = datetime.now().timestamp()
self.was_record_enabled_in_config = record_enabled_in_config self.was_record_enabled_in_config = record_enabled_in_config
@@ -323,9 +386,7 @@ class CameraWatchdog(threading.Thread):
) )
self.stop_all_ffmpeg() self.stop_all_ffmpeg()
self.start_all_ffmpeg() self.start_all_ffmpeg()
self.latest_valid_segment_time = 0 self._reset_segment_times()
self.latest_invalid_segment_time = 0
self.latest_cache_segment_time = 0
self.record_enable_time = datetime.now().astimezone(UTC) self.record_enable_time = datetime.now().astimezone(UTC)
last_restart_time = datetime.now().timestamp() last_restart_time = datetime.now().timestamp()
self.was_record_sub_enabled = record_sub_enabled self.was_record_sub_enabled = record_sub_enabled
@@ -343,26 +404,25 @@ class CameraWatchdog(threading.Thread):
raw_topic, payload = update raw_topic, payload = update
if raw_topic and payload: if raw_topic and payload:
topic = str(raw_topic) topic = str(raw_topic)
camera, segment_time, _ = payload camera, stream_type, segment_time, _ = payload
if camera != self.config.name: if camera != self.config.name:
continue continue
if topic.endswith(RecordingsDataTypeEnum.invalid.value): if topic.endswith(RecordingsDataTypeEnum.invalid.value):
self.logger.warning( 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): elif topic.endswith(RecordingsDataTypeEnum.valid.value):
self.logger.debug( 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): elif topic.endswith(RecordingsDataTypeEnum.latest.value):
if segment_time is not None: self.latest_cache_segment_time[stream_type] = (
self.latest_cache_segment_time = segment_time segment_time if segment_time is not None else 0
else: )
self.latest_cache_segment_time = 0
now = datetime.now().timestamp() now = datetime.now().timestamp()
@@ -409,63 +469,26 @@ class CameraWatchdog(threading.Thread):
for p in self.ffmpeg_other_processes: for p in self.ffmpeg_other_processes:
poll = p["process"].poll() 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) now_utc = datetime.now().astimezone(UTC)
# Check if we're within the grace period after enabling recording # ensure segments are still being created and that they have
# Grace period: 90 seconds allows time for ffmpeg to start and create first segment # valid video data. each stream is tracked separately so a
in_grace_period = self.record_enable_time is not None and ( # healthy one can't mask a stalled one.
now_utc - self.record_enable_time stale_stream = None
) < timedelta(seconds=90) stale_reason = None
for stream_type in recorded_streams:
stale_reason = self._stream_staleness(stream_type, now_utc)
latest_cache_dt = ( if stale_reason is not None:
datetime.fromtimestamp(self.latest_cache_segment_time, tz=UTC) stale_stream = stream_type
if self.latest_cache_segment_time > 0 break
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_stream is not None:
self.logger.error( 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["process"] = start_or_restart_ffmpeg(
p["cmd"], p["cmd"],
@@ -481,8 +504,13 @@ class CameraWatchdog(threading.Thread):
continue continue
else: else:
self._send_record_status("online", now) for stream_type in recorded_streams:
p["latest_segment_time"] = self.latest_cache_segment_time 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: if poll is None:
continue continue
@@ -497,6 +525,25 @@ class CameraWatchdog(threading.Thread):
p["cmd"], self.logger, p["logpipe"], ffmpeg_process=p["process"] 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 # Prune expired reconnect timestamps
now = datetime.now().timestamp() now = datetime.now().timestamp()
while ( while (
@@ -539,9 +586,9 @@ class CameraWatchdog(threading.Thread):
self.segment_subscriber.stop() self.segment_subscriber.stop()
def start_ffmpeg_detect(self): def start_ffmpeg_detect(self):
ffmpeg_cmd = [ detect_cmd = [c for c in self.config.ffmpeg_cmds if "detect" in c["roles"]][0]
c["cmd"] for c in self.config.ffmpeg_cmds if "detect" in c["roles"] ffmpeg_cmd = detect_cmd["cmd"]
][0] self.detect_process_records_sub = "record_sub" in detect_cmd["roles"]
self.ffmpeg_detect_process = start_or_restart_ffmpeg( self.ffmpeg_detect_process = start_or_restart_ffmpeg(
ffmpeg_cmd, self.logger, self.logpipe, self.frame_size ffmpeg_cmd, self.logger, self.logpipe, self.frame_size
) )