From 386b3d31aaf4c4d48ed70c289ad2bc361930687a Mon Sep 17 00:00:00 2001 From: Josh Hawkins <32435876+hawkeye217@users.noreply.github.com> Date: Mon, 28 Sep 2026 09:04:11 -0500 Subject: [PATCH] resolve segment start times in segment order A camera stream's cached segments are probed concurrently, and each one chained its start off `last_segment_end` as soon as its own probe finished. When segments backed up in the cache and a later probe finished first, it chained off the wrong segment and the earlier one then moved `last_segment_end` backwards, so rows lost their exact adjacency. Probes still run concurrently, but each segment now waits for the one before it to settle its start before resolving its own. --- frigate/record/maintainer.py | 45 ++++++++++++++++-- frigate/test/test_maintainer.py | 9 ++-- frigate/test/test_record_sub_maintainer.py | 54 ++++++++++++++++++++++ 3 files changed, 101 insertions(+), 7 deletions(-) diff --git a/frigate/record/maintainer.py b/frigate/record/maintainer.py index 7bb90ae823..e4151069f9 100644 --- a/frigate/record/maintainer.py +++ b/frigate/record/maintainer.py @@ -490,9 +490,17 @@ class RecordingMaintainer(threading.Thread): ) reviews = reviews_by_camera[camera] - tasks.extend( - [self.validate_and_move_segment(camera, reviews, r) for r in recordings] - ) + # probes run concurrently, but each segment's start chains off the + # previous segment's end, so starts resolve in segment order + previous: asyncio.Event | None = None + for recording in recordings: + resolved = asyncio.Event() + tasks.append( + self._validate_in_order( + camera, reviews, recording, previous, resolved + ) + ) + previous = resolved # publish most recently available recording time and None if disabled if stream_type == STREAM_TYPE_MAIN: @@ -550,12 +558,33 @@ class RecordingMaintainer(threading.Thread): while info and info[0][0] < expire_before: info.pop(0) + async def _validate_in_order( + self, + camera: str, + reviews: Any, + recording: dict[str, Any], + previous_start: asyncio.Event | None, + start_resolved: asyncio.Event, + ) -> dict[str, Any] | None: + """Validate a segment, always releasing the next one in its stream.""" + try: + return await self.validate_and_move_segment( + camera, reviews, recording, previous_start, start_resolved + ) + finally: + start_resolved.set() + def drop_segment(self, cache_path: str) -> None: Path(cache_path).unlink(missing_ok=True) self.end_time_cache.pop(cache_path, None) async def validate_and_move_segment( - self, camera: str, reviews: Any, recording: dict[str, Any] + self, + camera: str, + reviews: Any, + recording: dict[str, Any], + previous_start: asyncio.Event | None = None, + start_resolved: asyncio.Event | None = None, ) -> dict[str, Any] | None: cache_path: str = recording["cache_path"] start_time: datetime.datetime = recording["start_time"] @@ -617,6 +646,9 @@ class RecordingMaintainer(threading.Thread): async with self.probe_semaphore: keyframes = await get_keyframe_offsets(cache_path) + if previous_start is not None: + await previous_start.wait() + start_time = self._resolve_segment_start( camera, stream_type, start_time, duration, cache_path ) @@ -654,6 +686,11 @@ class RecordingMaintainer(threading.Thread): RecordingsDataTypeEnum.valid.value, ) + # the start is settled, so the next segment of the stream can chain + # off it while this one waits on retention and the move + if start_resolved is not None: + start_resolved.set() + record_config = self.config.cameras[camera].record # sub's alerts/detections carry the retain mode directly, unlike diff --git a/frigate/test/test_maintainer.py b/frigate/test/test_maintainer.py index ab7f7608df..b1f78845b1 100644 --- a/frigate/test/test_maintainer.py +++ b/frigate/test/test_maintainer.py @@ -1,7 +1,7 @@ import datetime import sys import unittest -from unittest.mock import MagicMock, patch +from unittest.mock import AsyncMock, MagicMock, patch # Mock complex imports before importing maintainer, saving originals so we can # restore them after import and avoid polluting sys.modules for other tests. @@ -48,8 +48,11 @@ class TestMaintainer(unittest.IsolatedAsyncioTestCase): "frigate.record.maintainer.psutil.process_iter", return_value=[] ): with patch("frigate.record.maintainer.logger.warning") as warn: - # Mock validate_and_move_segment to avoid further logic - maintainer.validate_and_move_segment = MagicMock() + # Mock validate_and_move_segment to avoid further logic. + # The requestor is real when another test imported the + # maintainer first, and it would block on a reply. + maintainer.validate_and_move_segment = AsyncMock() + maintainer.requestor = MagicMock() try: await maintainer.move_files() diff --git a/frigate/test/test_record_sub_maintainer.py b/frigate/test/test_record_sub_maintainer.py index a831e6ecc0..42a814aef0 100644 --- a/frigate/test/test_record_sub_maintainer.py +++ b/frigate/test/test_record_sub_maintainer.py @@ -1,5 +1,6 @@ """Tests for sub stream cache segment handling in the recording maintainer.""" +import asyncio import datetime import os import tempfile @@ -532,6 +533,59 @@ class TestSegmentStartChaining(unittest.IsolatedAsyncioTestCase): self.assertAlmostEqual(calls[1].args[2].timestamp(), self.T0 + 10.4, places=3) self.assertAlmostEqual(calls[1].args[3].timestamp(), self.T0 + 20.8, places=3) + async def test_out_of_order_probes_chain_in_segment_order(self): + maintainer = _build_chaining_maintainer(self.T0) + + async def probe(_ffmpeg, cache_path, get_duration=False): + # the earlier segment's probe finishes last + if "chain0" in cache_path: + await asyncio.sleep(0.05) + + return {"has_valid_video": True, "duration": 10.4} + + recordings = [ + { + "start_time": datetime.datetime.fromtimestamp( + self.T0 + offset, tz=datetime.UTC + ), + "cache_path": f"/tmp/cache/test_cam@chain{offset}.mp4", + "stream_type": "main", + } + for offset in (0, 10) + ] + first_resolved = asyncio.Event() + + with ( + patch("frigate.record.maintainer.get_video_properties", probe), + patch( + "frigate.record.maintainer.get_keyframe_offsets", + AsyncMock(return_value=[0]), + ), + patch( + "frigate.record.maintainer.os.path.getmtime", + MagicMock(side_effect=OSError("missing")), + ), + ): + await asyncio.gather( + maintainer._validate_in_order( + "test_cam", [], recordings[0], None, first_resolved + ), + maintainer._validate_in_order( + "test_cam", [], recordings[1], first_resolved, asyncio.Event() + ), + ) + + starts = sorted( + call.args[2].timestamp() for call in maintainer.move_segment.await_args_list + ) + self.assertEqual(starts[0], self.T0) + self.assertAlmostEqual(starts[1], self.T0 + 10.4, places=3) + self.assertAlmostEqual( + maintainer.last_segment_end[("test_cam", "main")], + self.T0 + 20.8, + places=3, + ) + async def test_genuine_gap_is_not_snapped(self): maintainer = _build_chaining_maintainer(self.T0)