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.
This commit is contained in:
Josh Hawkins
2026-09-28 09:28:09 -05:00
parent 3f00c75bbd
commit 386b3d31aa
3 changed files with 101 additions and 7 deletions
+41 -4
View File
@@ -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
+6 -3
View File
@@ -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()
@@ -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)