diff --git a/frigate/api/app.py b/frigate/api/app.py index 62676bbdef..9cd1c8fe3b 100644 --- a/frigate/api/app.py +++ b/frigate/api/app.py @@ -393,6 +393,11 @@ def config(request: Request): model_dict["non_logo_attributes"] = model.non_logo_attributes model_dict["labelmap"] = model.merged_labelmap + # report the configured reference rather than the resolved cache path, + # so saving the config back doesn't lose the Frigate+ model + if model.plus_id: + model_dict["path"] = f"plus://{model.plus_id}" + if not config["plus"]["enabled"]: continue diff --git a/frigate/api/config_util.py b/frigate/api/config_util.py index 0cb954af10..c69c242aed 100644 --- a/frigate/api/config_util.py +++ b/frigate/api/config_util.py @@ -51,6 +51,7 @@ def swap_runtime_config(app: FastAPI, config: FrigateConfig) -> None: if app.stats_emitter is not None: app.stats_emitter.config = config + app.stats_emitter.hardware_stats.set_config(config) if app.dispatcher is not None: app.dispatcher.config = config diff --git a/frigate/api/media.py b/frigate/api/media.py index f94886b48e..2305f06eb6 100644 --- a/frigate/api/media.py +++ b/frigate/api/media.py @@ -63,7 +63,6 @@ from frigate.util.recording_coverage import ( null_audio_glitches, plan_clip, resolve_coverage, - stream_has_audio, ) logger = logging.getLogger(__name__) @@ -681,15 +680,10 @@ async def _vod_response( end_ts, force_discontinuity, ) - intervals = resolve_coverage(camera_name, start_ts, end_ts) - # rows contradicting their stream's audio composition are # truncated-shutdown glitches - main_audio = stream_has_audio(intervals, main=True) - sub_audio = stream_has_audio(intervals, main=False) - spans = build_spans( - null_audio_glitches(intervals, main_audio, sub_audio), + null_audio_glitches(resolve_coverage(camera_name, start_ts, end_ts)), stream_preference, ) diff --git a/frigate/detectors/detector_config.py b/frigate/detectors/detector_config.py index 1193ba0740..2362c073cc 100644 --- a/frigate/detectors/detector_config.py +++ b/frigate/detectors/detector_config.py @@ -123,6 +123,7 @@ class ModelConfig(BaseModel): _all_attributes: list[str] = PrivateAttr() _all_attribute_logos: list[str] = PrivateAttr() _model_hash: str = PrivateAttr() + _plus_id: str | None = PrivateAttr(default=None) @property def merged_labelmap(self) -> dict[int, str]: @@ -148,6 +149,11 @@ class ModelConfig(BaseModel): def model_hash(self) -> str: return self._model_hash + @property + def plus_id(self) -> str | None: + """The Frigate+ model id, once a plus:// path has been resolved.""" + return self._plus_id + def __init__(self, **config): super().__init__(**config) @@ -178,6 +184,7 @@ class ModelConfig(BaseModel): os.makedirs(MODEL_CACHE_DIR, exist_ok=True) model_id = self.path[7:] + self._plus_id = model_id self.path = os.path.join(MODEL_CACHE_DIR, model_id) model_info_path = f"{self.path}.json" diff --git a/frigate/record/export.py b/frigate/record/export.py index 7c87c5fd83..7dfc2d739f 100644 --- a/frigate/record/export.py +++ b/frigate/record/export.py @@ -42,6 +42,8 @@ from frigate.util.ownership import chown_to_runtime from frigate.util.recording_coverage import ( build_spans, known_video_codecs, + null_audio_glitches, + realized_timeline, resolve_coverage, stream_media_summary, ) @@ -84,10 +86,20 @@ class StreamRun: @dataclass class _ChapterWindow: - """A merged-timeline slice, shaped like the recording rows chapters read.""" + """A merged-timeline slice, shaped like the recording rows chapters read. + + lead_in is the output time the slice's vod clip plays before start_time, + from snapping its first frame back to a keyframe. + """ start_time: float end_time: float + lead_in: float = 0.0 + + +def _lead_in(recording: Any) -> float: + """Output seconds a chapter source plays before its first wall second.""" + return recording.lead_in if isinstance(recording, _ChapterWindow) else 0.0 # Matches the setpts factor used in timelapse exports (e.g. setpts=0.04*PTS). @@ -376,15 +388,18 @@ class RecordingExporter(threading.Thread): def _resolve_coverage(self) -> tuple[list[list[Any]], set[str], bool]: """Resolve the export range into the spans the VOD manifest will serve. - Delegates to the same coverage resolution the manifest builder - uses, so what we plan around and what nginx-vod emits agree by - construction. Returns the spans (each [row, start, end, is_main]), - the known video codecs, and whether audio survives the range. + Delegates to the same coverage resolution and glitch nulling the + manifest builder uses, so what we plan around and what nginx-vod + emits agree by construction. Returns the spans (each [row, start, + end, is_main]), the known video codecs, and whether audio survives + the range. Memoized: several stages of the export ask the same question, and the recordings backing a finished range do not change under us. """ if self._coverage is None: - intervals = resolve_coverage(self.camera, self.start_time, self.end_time) + intervals = null_audio_glitches( + resolve_coverage(self.camera, self.start_time, self.end_time) + ) self._coverage = ( build_spans(intervals, self.pinned_stream), known_video_codecs(intervals), @@ -433,17 +448,57 @@ class RecordingExporter(threading.Thread): # hand-off to stage around return True - spans, codecs, keep_audio = self._resolve_coverage() - runs = self._stream_runs(spans) + _spans, codecs, keep_audio = self._resolve_coverage() + runs = self._planned_stream_runs() # a range one stream covers end to end has nothing to hand off, # so it stays on the existing path however long it is if len(runs) < 2: return True - runs = [piece for run in runs for piece in self._split_long_run(run)] return self._stage_stream_runs(runs, codecs, keep_audio) + def _planned_stream_runs(self) -> list[StreamRun]: + """The runs a mixed range is staged as, one pinned vod playlist each.""" + runs = self._stream_runs(self._merged_spans()) + + if len(runs) < 2: + return runs + + return [piece for run in runs for piece in self._split_long_run(run)] + + def _staged_chapter_windows(self) -> list[_ChapterWindow]: + """Chapter windows for the staged files as they were rendered. + + Each staged run comes from its own pinned vod playlist, whose first + clip snaps back to the preceding keyframe, so a staged file runs up + to a GOP longer than its slice of the merged timeline. Planning each + run the way its playlist does carries that lead-in into the chapter + offsets instead of letting it accumulate at every hand-off. + """ + windows: list[_ChapterWindow] = [] + + for run in self._planned_stream_runs(): + intervals = null_audio_glitches( + resolve_coverage(self.camera, run.start_time, run.end_time) + ) + + for clip in realized_timeline(intervals, run.stream_type): + # a skipped clip is absent from the playlist and the file + if clip["duration"] <= 0: + continue + + span = clip["end_time"] - clip["start_time"] + windows.append( + _ChapterWindow( + clip["start_time"], + clip["end_time"], + max(0.0, clip["duration"] / 1000 - span), + ) + ) + + return windows + def _stream_runs(self, spans: list[list[Any]]) -> list[StreamRun]: """Collapse the merged spans into contiguous runs of one stream type. @@ -840,6 +895,8 @@ class RecordingExporter(threading.Thread): clipped_end = min(float(rec.end_time), float(self.end_time)) if clipped_end <= clipped_start: continue + # a staged window's keyframe lead-in plays before it + output_offset += _lead_in(rec) windows.append((clipped_start, clipped_end, output_offset)) output_offset += clipped_end - clipped_start @@ -987,9 +1044,13 @@ class RecordingExporter(threading.Thread): if duration_ms <= 0: continue - title = datetime.datetime.fromtimestamp(clipped_start, tz=tz).isoformat( - timespec="seconds" - ) + # a staged window's keyframe lead-in opens its chapter, with + # frames captured that long before the window + lead_in = _lead_in(rec) + duration_ms += int(round(lead_in * 1000)) + title = datetime.datetime.fromtimestamp( + clipped_start - lead_in, tz=tz + ).isoformat(timespec="seconds") chapter_blocks.append( "[CHAPTER]\n" "TIMEBASE=1/1000\n" @@ -1128,12 +1189,12 @@ class RecordingExporter(threading.Thread): if self.staged_runs: # each run was already rendered to a temp file with a common # track timescale, so the concat demuxer has nothing left to - # reconcile and every chapter offset lines up with the merged - # timeline the staged files reproduce - recordings = [ - _ChapterWindow(span_start, span_end) - for _row, span_start, span_end, _is_main in self._merged_spans() - ] + # reconcile + recordings = ( + self._staged_chapter_windows() + if self.chapters not in (None, ChaptersEnum.none) + else [] + ) playlist_lines: list[str] = [f"file '{path}'" for path in self.staged_runs] ffmpeg_input = ( "-y -protocol_whitelist pipe,file -f concat -safe 0 -i /dev/stdin" @@ -1149,11 +1210,18 @@ class RecordingExporter(threading.Thread): # its own rows are the ones the chapters describe recordings = self._get_recordings_for_range(pin) else: - # never mix streams in one playlist; use main when available - # and fall back to sub for expired-main history - recordings = self._get_recordings_for_range(STREAM_TYPE_MAIN) + # an unstaged auto range resolves to at most one stream run, and + # its rows are the ones the chapters describe. Main rows the + # manifest drops (glitches, slivers at the edges of a sub range) + # must not stand in for it. + runs = self._stream_runs(self._merged_spans()) + recordings = self._get_recordings_for_range( + runs[0].stream_type if runs else STREAM_TYPE_MAIN + ) - if not recordings: + # never mix streams in one playlist; fall back to sub for + # expired-main history + if not recordings and not runs: recordings = self._get_recordings_for_range(STREAM_TYPE_SUB) playlist_lines = [] 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/stats/hardware.py b/frigate/stats/hardware.py index b8a877761c..dbe18cc548 100644 --- a/frigate/stats/hardware.py +++ b/frigate/stats/hardware.py @@ -256,6 +256,17 @@ class HardwareStats: ) self.update_config() + def set_config(self, config: FrigateConfig) -> None: + """Follow a runtime config swap and recalculate the monitored hardware. + + The camera update subscriber has to follow too, or later camera updates + would land on the discarded config. + """ + self.config = config + self._config_subscriber.config = config + self._config_subscriber.camera_configs = config.cameras + self.update_config() + def update_config(self) -> None: """Recalculate all hardware that needs to be monitored from the config.""" names = self._scan_ffmpeg() | self._scan_detectors() | self._scan_enrichments() diff --git a/frigate/stats/util.py b/frigate/stats/util.py index ef220e0125..620f8d691d 100644 --- a/frigate/stats/util.py +++ b/frigate/stats/util.py @@ -111,15 +111,18 @@ def get_detector_stats( ) -> dict[str, dict[str, Any]]: """Get stats for all detectors, including temperatures based on detector type.""" detector_stats: dict[str, dict[str, Any]] = {} - detector_type_indices: dict[str, int] = {} + # detector type -> device -> index into that type's temperatures + device_indices: dict[str, dict[str, int]] = {} for name, detector in stats_tracking["detectors"].items(): pid = detector.detect_process.pid if detector.detect_process else None detector_type = detector.detector_config.type - # Keep track of the index for each detector type to match temperatures correctly - current_index = detector_type_indices.get(detector_type, 0) - detector_type_indices[detector_type] = current_index + 1 + # temperatures are per physical unit, so a repeated device + # ("hailo:PCIe#2", see runner_names) shares its unit's reading + device = name.partition("#")[0] + type_devices = device_indices.setdefault(detector_type, {}) + current_index = type_devices.setdefault(device, len(type_devices)) detector_stat = { "inference_speed": round(detector.avg_inference_speed.value * 1000, 2), # type: ignore[attr-defined] diff --git a/frigate/test/http_api/test_http_app.py b/frigate/test/http_api/test_http_app.py index ef5b99ad08..3462a4f108 100644 --- a/frigate/test/http_api/test_http_app.py +++ b/frigate/test/http_api/test_http_app.py @@ -1,8 +1,10 @@ +import json +import os from unittest.mock import Mock, patch import frigate.genai from frigate.config import GenAIProviderEnum -from frigate.const import REDACTED_CREDENTIAL_SENTINEL +from frigate.const import MODEL_CACHE_DIR, REDACTED_CREDENTIAL_SENTINEL from frigate.genai import GenAIClient from frigate.models import Event, Recordings, ReviewSegment from frigate.stats.emitter import StatsEmitter @@ -90,6 +92,44 @@ class TestHttpApp(BaseTestHttp): mqtt = response.json()["mqtt"] assert mqtt["password"] == REDACTED_CREDENTIAL_SENTINEL + def test_config_response_keeps_plus_model_reference(self): + model_id = "test_plus_reference" + model_path = os.path.join(MODEL_CACHE_DIR, model_id) + os.makedirs(MODEL_CACHE_DIR, exist_ok=True) + + with open(model_path, "w") as f: + f.write("model") + + with open(f"{model_path}.json", "w") as f: + json.dump( + { + "id": model_id, + "type": "ssd", + "supportedDetectors": ["cpu"], + "width": 320, + "height": 320, + "inputShape": "nhwc", + "pixelFormat": "rgb", + "labelMap": {"0": "person"}, + }, + f, + ) + + self.addCleanup(os.remove, model_path) + self.addCleanup(os.remove, f"{model_path}.json") + self.minimal_config["models"] = [ + {"path": f"plus://{model_id}", "devices": ["cpu"]} + ] + app = super().create_app() + + with AuthTestClient(app) as client: + response = client.get("/config") + assert response.status_code == 200 + assert response.json()["models"][0]["path"] == f"plus://{model_id}" + + # detection still loads the resolved cache file + assert app.frigate_config.models[0].path == model_path + #################################################################################################################### ################################### POST /genai/probe Endpoint ################################################## #################################################################################################################### diff --git a/frigate/test/test_config_util.py b/frigate/test/test_config_util.py index 5888f26188..71d1eccdfd 100644 --- a/frigate/test/test_config_util.py +++ b/frigate/test/test_config_util.py @@ -26,6 +26,7 @@ class TestSwapRuntimeConfig(unittest.TestCase): app.genai_manager.update_config.assert_called_once_with(config) app.profile_manager.update_config.assert_called_once_with(config) self.assertIs(app.stats_emitter.config, config) + app.stats_emitter.hardware_stats.set_config.assert_called_once_with(config) self.assertIs(app.dispatcher.config, config) for comm in app.dispatcher.comms: self.assertIs(comm.config, config) diff --git a/frigate/test/test_detector_stats.py b/frigate/test/test_detector_stats.py new file mode 100644 index 0000000000..db7f265e0d --- /dev/null +++ b/frigate/test/test_detector_stats.py @@ -0,0 +1,39 @@ +"""Tests for per-detector stats.""" + +import unittest +from unittest.mock import MagicMock, patch + +from frigate.stats.util import get_detector_stats + + +def _detector(detector_type: str) -> MagicMock: + detector = MagicMock() + detector.detector_config.type = detector_type + detector.avg_inference_speed.value = 0.01 + detector.detection_start.value = 0.0 + detector.detect_process.pid = 1 + return detector + + +class TestDetectorTemperatures(unittest.TestCase): + def test_repeated_device_shares_its_unit_temperature(self): + stats_tracking = { + "detectors": { + "hailo:PCIe": _detector("hailo8l"), + "hailo:PCIe#2": _detector("hailo8l"), + "hailo:PCIe:1": _detector("hailo8l"), + } + } + + with patch( + "frigate.stats.util.get_hardware_temperatures", return_value=[50.0, 60.0] + ): + stats = get_detector_stats(stats_tracking) + + self.assertEqual(stats["hailo:PCIe"]["temperature"], 50.0) + self.assertEqual(stats["hailo:PCIe#2"]["temperature"], 50.0) + self.assertEqual(stats["hailo:PCIe:1"]["temperature"], 60.0) + + +if __name__ == "__main__": + unittest.main() diff --git a/frigate/test/test_export.py b/frigate/test/test_export.py index 4aee38fb6c..4d27e272df 100644 --- a/frigate/test/test_export.py +++ b/frigate/test/test_export.py @@ -543,6 +543,56 @@ class TestPinnedStream(unittest.TestCase): self.assertFalse(any("/vod/front/main/" in token for token in cmd)) +class TestExportTimelineAlignment(unittest.TestCase): + def test_unstaged_auto_reads_the_stream_it_serves(self) -> None: + # a sub-only range whose main rows are glitches the manifest drops + exporter = _make_exporter([_span("/s1.mp4", 1_000, 1_040, False)], {"h264"}) + streams: list[str] = [] + + def rows(stream: str) -> list: + streams.append(stream) + return [_FakeRow(f"/{stream}.mp4")] + + exporter._get_recordings_for_range = rows # type: ignore[method-assign] + exporter.get_record_export_command("/exports/out.mp4") + + self.assertEqual(streams, ["sub"]) + + def test_staged_chapters_carry_keyframe_lead_in(self) -> None: + exporter = _make_exporter( + [ + _span("/m1.mp4", 1_000, 1_020, True), + _span("/s1.mp4", 1_020, 1_040, False), + ], + {"h264"}, + ) + exporter.config.ui.timezone = None + + # each run's vod clip snaps 1.5s back to a keyframe + def timeline(_intervals: list, stream: str) -> list[dict]: + start, end = (1_000, 1_020) if stream == "main" else (1_020, 1_040) + return [ + { + "start_time": start, + "end_time": end, + "duration": (end - start + 1.5) * 1000, + } + ] + + with ( + patch("frigate.record.export.resolve_coverage", return_value=[]), + patch("frigate.record.export.realized_timeline", side_effect=timeline), + ): + windows = exporter._staged_chapter_windows() + + path = exporter._build_recording_segment_chapter_metadata_file(windows) + self.addCleanup(os.remove, path) + content = Path(path).read_text() + + self.assertIn("START=0\nEND=21500", content) + self.assertIn("START=21500\nEND=43000", content) + + class TestStagedFileCleanup(unittest.TestCase): """A staged path must be tracked before ffmpeg can write to it.""" diff --git a/frigate/test/test_hardware_stats.py b/frigate/test/test_hardware_stats.py index 2bb1e9cc9b..e75218d5d3 100644 --- a/frigate/test/test_hardware_stats.py +++ b/frigate/test/test_hardware_stats.py @@ -288,6 +288,17 @@ class TestUpdateConfig(HardwareStatsTestCase): self.assertEqual(set(stats._monitored), {"rockchip"}) + def test_follows_a_runtime_config_swap(self): + stats = self.make_stats(self.make_config()) + self.assertEqual(set(stats._monitored), set()) + + swapped = self.make_config("preset-rk-h264") + stats.set_config(swapped) + + self.assertEqual(set(stats._monitored), {"rockchip"}) + self.assertIs(self.subscriber.return_value.config, swapped) + self.assertIs(self.subscriber.return_value.camera_configs, swapped.cameras) + class TestUpdateStats(HardwareStatsTestCase): def run_stats(self, stats: HardwareStats) -> dict: 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) diff --git a/frigate/util/downloader.py b/frigate/util/downloader.py index 9a8e46ad9f..1755f421ff 100644 --- a/frigate/util/downloader.py +++ b/frigate/util/downloader.py @@ -14,6 +14,10 @@ from frigate.util.file import FileLock logger = logging.getLogger(__name__) +# (connect, read) seconds. The read timeout bounds each socket read rather than +# the whole download, so large models still finish. +DOWNLOAD_TIMEOUT = (15, 60) + # target path -> first line of the last download error for it; every existing # download function swallows its exceptions, so this is how the downloader # thread learns why a file is still missing @@ -124,7 +128,9 @@ class ModelDownloader: logger.info(f"Downloading model file from: {url}") try: - with requests.get(url, stream=True, allow_redirects=True) as r: + with requests.get( + url, stream=True, allow_redirects=True, timeout=DOWNLOAD_TIMEOUT + ) as r: r.raise_for_status() with open(temporary_filename, "wb") as f: for chunk in r.iter_content(chunk_size=8192): diff --git a/frigate/util/recording_coverage.py b/frigate/util/recording_coverage.py index d8bbd55ad6..96a3ec5292 100644 --- a/frigate/util/recording_coverage.py +++ b/frigate/util/recording_coverage.py @@ -213,17 +213,19 @@ def stream_has_audio(intervals: list[CoverageInterval], main: bool) -> bool: ) -def null_audio_glitches( - intervals: list[CoverageInterval], main_audio: bool, sub_audio: bool -) -> list[CoverageInterval]: +def null_audio_glitches(intervals: list[CoverageInterval]) -> list[CoverageInterval]: """Treat video-only glitch rows on audio-bearing streams as no recording. nginx-vod requires every clip in a sequence to carry the same track count, so a truncated video-only segment (a backend restart can flush a sub-second file before any audio packet landed) poisons every manifest that includes it. Nulling the row turns the glitch into a - hole the span builder skips like any recording gap. + hole the span builder skips like any recording gap. Every consumer of + a window's coverage (the vod manifest, its realized timelines, and + exports) goes through here, so they all agree on which rows exist. """ + main_audio = stream_has_audio(intervals, main=True) + sub_audio = stream_has_audio(intervals, main=False) result: list[CoverageInterval] = [] for interval in intervals: main = interval.main @@ -449,9 +451,7 @@ def realized_timelines( assembles each variant's realized spans. Keyframe snapping reads the per-row index stored at record time, so no file is touched. """ - main_audio = stream_has_audio(intervals, main=True) - sub_audio = stream_has_audio(intervals, main=False) - nulled = null_audio_glitches(intervals, main_audio, sub_audio) + nulled = null_audio_glitches(intervals) return { "auto": realized_timeline(nulled, None), diff --git a/web/e2e/specs/settings/detection-models.spec.ts b/web/e2e/specs/settings/detection-models.spec.ts index f764b70d7c..0dab62d6bf 100644 --- a/web/e2e/specs/settings/detection-models.spec.ts +++ b/web/e2e/specs/settings/detection-models.spec.ts @@ -217,6 +217,23 @@ test.describe("Detection models settings @high", () => { ).toBeDisabled(); }); + test("shareable hardware another model uses can still be picked", async ({ + frigateApp, + }) => { + await installRoutes(frigateApp.page, [ + { scene: "all", devices: ["openvino:GPU.0"] }, + { scene: "outdoor", devices: ["openvino:GPU.1"] }, + ]); + await openPage(frigateApp); + + await expect( + frigateApp.page.locator("#models-0-openvino\\:GPU\\.1").first(), + ).toBeEnabled(); + await expect(frigateApp.page.locator("#pageRoot")).not.toContainText( + "used by outdoor", + ); + }); + test("adding a model appends a card with an unused scene", async ({ frigateApp, }) => { @@ -248,15 +265,13 @@ test.describe("Detection models settings @high", () => { test("a saved Frigate+ model opens on the Frigate+ tab", async ({ frigateApp, }) => { - // the backend resolves plus:// to a cache path before serving the config - // back, so the plus metadata is the only signal the model is a Plus one await installRoutes( frigateApp.page, [ { scene: "all", devices: ["openvino:GPU.0"], - path: "/config/model_cache/abc123", + path: "plus://abc123", plus: PLUS_MODEL, }, ], @@ -302,6 +317,65 @@ test.describe("Detection models settings @high", () => { expect(saves.at(-1)?.config_data?.models?.[0].path).toBe("plus://abc123"); }); + test("saving a Frigate+ model keeps its reference without the Frigate+ fields", async ({ + frigateApp, + }) => { + // the backend fills these in from the Frigate+ model info when it loads + const saves = await installRoutes( + frigateApp.page, + [ + { + scene: "all", + devices: ["openvino:GPU.0"], + path: "plus://abc123", + plus: PLUS_MODEL, + width: 320, + height: 320, + input_tensor: "nchw", + input_dtype: "float", + model_type: "yolo-generic", + }, + ], + true, + ); + await openPage(frigateApp); + + await frigateApp.page.locator("#models-0-openvino\\:GPU\\.1").click(); + await frigateApp.page.getByRole("button", { name: /^Save$/ }).click(); + await expect.poll(() => saves.length).toBeGreaterThan(0); + + const model = saves.at(-1)?.config_data?.models?.[0]; + expect(model?.path).toBe("plus://abc123"); + expect(model?.devices).toEqual(["openvino:GPU.0", "openvino:GPU.1"]); + expect(model).not.toHaveProperty("width"); + expect(model).not.toHaveProperty("input_tensor"); + expect(model).not.toHaveProperty("model_type"); + // a leftover dtype from a custom model must not override the int default + expect(model).not.toHaveProperty("input_dtype"); + }); + + test("a Frigate+ model only shows its path without a Frigate+ API key", async ({ + frigateApp, + }) => { + await installRoutes(frigateApp.page, [ + { + scene: "all", + devices: ["openvino:GPU.0"], + path: "plus://abc123", + width: 320, + height: 320, + }, + ]); + await openPage(frigateApp); + + const root = frigateApp.page.locator("#pageRoot"); + await expect(root).toContainText("Custom object detector model path"); + await expect(root).not.toContainText("Object detection model input width"); + await expect(root).not.toContainText( + "Label map for custom object detector", + ); + }); + test("a Frigate+ Hailo model is listed by the device it was built for", async ({ frigateApp, }) => { diff --git a/web/public/locales/en/views/settings.json b/web/public/locales/en/views/settings.json index af823a68b1..2406442ec4 100644 --- a/web/public/locales/en/views/settings.json +++ b/web/public/locales/en/views/settings.json @@ -2063,7 +2063,7 @@ "unrecognized": "This model is configured for hardware that was not found on this system: {{devices}}", "description": "The hardware this model runs its detection on.", "detectorCountDescription": "How many detection processes to run on this hardware. More detectors keep up with more cameras, at the cost of extra device memory.", - "unitsDescription": "Each unit runs its own detection process. A unit already used by another model can not be selected." + "unitsDescription": "Each unit runs its own detection process. A unit that can't be shared, such as a Coral, can only be used by one model." }, "tabs": { "plus": "Frigate+", diff --git a/web/src/components/config-form/sections/section-special-cases.ts b/web/src/components/config-form/sections/section-special-cases.ts index 08f7d8c736..84edbb1d59 100644 --- a/web/src/components/config-form/sections/section-special-cases.ts +++ b/web/src/components/config-form/sections/section-special-cases.ts @@ -7,6 +7,7 @@ */ import { RJSFSchema } from "@rjsf/utils"; +import { omit } from "lodash"; import { applySchemaDefaults } from "@/lib/config-schema"; import { isJsonObject } from "@/lib/utils"; import { HiddenFieldContext, JsonObject, JsonValue } from "@/types/configForm"; @@ -352,6 +353,18 @@ export function synthesizeMissingFilters( return { ...(data as JsonObject), filters: newFilters }; } +// a Frigate+ model's config comes from its model info, so these are either +// redundant or stale leftovers from a custom model. A missing input_dtype +// in the model info means the int default. +const PLUS_SUPPLIED_MODEL_FIELDS = [ + "width", + "height", + "input_tensor", + "input_pixel_format", + "input_dtype", + "model_type", +]; + /** * Sanitize overrides payloads for section-specific quirks. */ @@ -360,6 +373,21 @@ export function sanitizeOverridesForSection( level: string, overrides: unknown, ): unknown { + // the models list is saved whole + if (sectionPath === "models" && Array.isArray(overrides)) { + return overrides.map((model) => { + if ( + !isJsonObject(model) || + typeof model.path !== "string" || + !model.path.startsWith("plus://") + ) { + return model; + } + + return omit(model, PLUS_SUPPLIED_MODEL_FIELDS); + }); + } + if (!overrides || !isJsonObject(overrides)) { return overrides; } diff --git a/web/src/components/config-form/theme/fields/HardwarePicker.tsx b/web/src/components/config-form/theme/fields/HardwarePicker.tsx index 98382a8bf7..fb2bcb72ef 100644 --- a/web/src/components/config-form/theme/fields/HardwarePicker.tsx +++ b/web/src/components/config-form/theme/fields/HardwarePicker.tsx @@ -21,7 +21,8 @@ type HardwarePickerProps = { // scopes the unit checkbox ids, since several models can list the same unit idPrefix: string; devices: string[]; - // device strings already taken by another model, mapped to that model's scene + // device strings already taken by another model, mapped to that model's + // scene. Only binding for hardware that can't be shared. claimedElsewhere: Record; cameraCount: number; disabled?: boolean; @@ -84,8 +85,11 @@ export function HardwarePicker({ return; } - // start with the first unit no other model has taken - const free = entry.units.find((unit) => !claimedElsewhere[unit.device]); + // start with the first unit no other model has taken, falling back to + // a taken one when the hardware can be shared + const free = + entry.units.find((unit) => !claimedElsewhere[unit.device]) ?? + (entry.unlimited ? entry.units[0] : undefined); if (!free) { onChange([]); @@ -184,7 +188,9 @@ export function HardwarePicker({ {t("detectionModels.hardware.unitsDescription")}

{selected.units.map((unit) => { - const claimedBy = claimedElsewhere[unit.device]; + const claimedBy = selected.unlimited + ? undefined + : claimedElsewhere[unit.device]; return (
` to a local - // cache path before serving the config back + // Frigate+ metadata the backend attaches to a saved model plus?: { id: string } | null; // the detector this model runs on, used to filter incompatible Plus models detector?: string; diff --git a/web/src/components/config-form/theme/fields/ModelsField.tsx b/web/src/components/config-form/theme/fields/ModelsField.tsx index f03f3f494b..f7afbfa37c 100644 --- a/web/src/components/config-form/theme/fields/ModelsField.tsx +++ b/web/src/components/config-form/theme/fields/ModelsField.tsx @@ -60,6 +60,15 @@ const CUSTOM_MODEL_FIELDS = [ "model_type", ]; +/** + * The fields a model can edit. A Frigate+ model's size, format, type, and + * labels come from its model info, so only its path stays editable. + */ +const editableFields = (model: DetectionModel): string[] => + typeof model.path === "string" && model.path.startsWith("plus://") + ? ["path"] + : CUSTOM_MODEL_FIELDS; + /** The detector a model runs on, which is the prefix of its device strings. */ const detectorForModel = (model: DetectionModel): string | undefined => model.devices?.[0]?.split(":")[0]; @@ -169,20 +178,23 @@ export function ModelsField(props: FieldProps) { [savedModels], ); - // a model serves the cameras naming its scene, plus every camera that names - // no scene at all when it is the "all" model + // a model serves the cameras naming its scene, and like the backend, the + // "all" model also serves every camera whose scene has no model of its own const cameraCountForScene = useCallback( (scene: string | undefined): number => { if (!cameras) { return 0; } + const modelScenes = new Set(models.map((model) => model.scene ?? "all")); + return Object.values(cameras).filter((camera) => { - const cameraScene = camera?.detect?.scene; - return cameraScene ? cameraScene === scene : scene === "all"; + const cameraScene = camera?.detect?.scene ?? "all"; + const servedBy = modelScenes.has(cameraScene) ? cameraScene : "all"; + return servedBy === (scene ?? "all"); }).length; }, - [cameras], + [cameras, models], ); const claimedByOtherModels = useCallback( @@ -314,7 +326,9 @@ export function ModelsField(props: FieldProps) { ); return ( - + // keyed by scene, which is unique per model, so deleting a card + // doesn't hand its state (such as the model source tab) to the next + @@ -428,7 +442,7 @@ export function ModelsField(props: FieldProps) { detector={detectorForModel(model)} disabled={disabled || readonly} onPathChange={(path) => updateModel(index, { path })} - customFields={CUSTOM_MODEL_FIELDS.map((fieldName) => + customFields={editableFields(model).map((fieldName) => renderField(index, fieldName), )} /> diff --git a/web/src/components/player/HlsVideoPlayer.tsx b/web/src/components/player/HlsVideoPlayer.tsx index e0d606cffb..7b0a66004e 100644 --- a/web/src/components/player/HlsVideoPlayer.tsx +++ b/web/src/components/player/HlsVideoPlayer.tsx @@ -26,6 +26,7 @@ import { useIsAdmin } from "@/hooks/use-is-admin"; // Android native hls does not seek correctly const USE_NATIVE_HLS = false; const HLS_MIME_TYPE = "application/vnd.apple.mpegurl" as const; +const DEFAULT_MAX_BUFFER_LENGTH_S = 10; const unsupportedErrorCodes: number[] = [ MediaError.MEDIA_ERR_SRC_NOT_SUPPORTED, MediaError.MEDIA_ERR_DECODE, @@ -148,6 +149,42 @@ export default function HlsVideoPlayer({ // a ref rather than an effect-scoped counter so the element error // handler can hold its toast while a recovery is still possible const mediaRecoveryBudgetRef = useRef(0); + // hls.js and the element can both report one failure, so each source + // toasts at most once + const failureReportedRef = useRef(false); + + const reportPlaybackFailure = useCallback( + (code: number, message: string) => { + if (failureReportedRef.current) { + return; + } + + failureReportedRef.current = true; + toast.error(t("toast.error.playRecordingsFailed", { code, message }), { + position: "top-center", + }); + }, + [t], + ); + const reportPlaybackFailureRef = useRef(reportPlaybackFailure); + + useEffect(() => { + reportPlaybackFailureRef.current = reportPlaybackFailure; + }, [reportPlaybackFailure]); + + // a quality switch changes the buffer length a commit before the new + // source arrives, so it updates the live instance instead of rebuilding + // it on the outgoing playlist + const bufferLengthRef = useRef(bufferLength); + + useEffect(() => { + bufferLengthRef.current = bufferLength; + + if (hlsRef.current) { + hlsRef.current.config.maxBufferLength = + bufferLength ?? DEFAULT_MAX_BUFFER_LENGTH_S; + } + }, [bufferLength]); const applyVideoDimensions = useCallback( (width: number, height: number) => { @@ -225,6 +262,7 @@ export default function HlsVideoPlayer({ // the element already holds a decoded frame, and keeping it visible // bridges the gap while the new source loads const currentPlaybackRate = videoRef.current.playbackRate; + failureReportedRef.current = false; if (!useHlsCompat) { nativeRetryRef.current = 0; @@ -236,7 +274,7 @@ export default function HlsVideoPlayer({ // Base HLS configuration const hlsConfig: Partial = { - maxBufferLength: bufferLength ?? 10, + maxBufferLength: bufferLengthRef.current ?? DEFAULT_MAX_BUFFER_LENGTH_S, maxBufferSize: 20 * 1000 * 1000, startPosition: currentSource.startPosition, }; @@ -269,10 +307,18 @@ export default function HlsVideoPlayer({ data.details === Hls.ErrorDetails.BUFFER_INCOMPATIBLE_CODECS_ERROR || data.details === Hls.ErrorDetails.BUFFER_ADD_CODEC_ERROR; - if (isCodecError && qualitySignalsRef.current.onFatalCodecError?.()) { + if (isCodecError) { + // with no stream to fall back to, the element may never raise + // an error of its own, so report the failure here + if (!qualitySignalsRef.current.onFatalCodecError?.()) { + reportPlaybackFailureRef.current( + MediaError.MEDIA_ERR_SRC_NOT_SUPPORTED, + data.details, + ); + } return; } - if (!isCodecError && mediaRecoveryBudgetRef.current > 0) { + if (mediaRecoveryBudgetRef.current > 0) { mediaRecoveryBudgetRef.current -= 1; hls.recoverMediaError(); } @@ -308,7 +354,7 @@ export default function HlsVideoPlayer({ hlsRef.current.destroy(); } }; - }, [videoRef, hlsRef, useHlsCompat, currentSource, bufferLength]); + }, [videoRef, hlsRef, useHlsCompat, currentSource]); // state handling @@ -714,15 +760,7 @@ export default function HlsVideoPlayer({ } } - toast.error( - t("toast.error.playRecordingsFailed", { - code: mediaError.code, - message: mediaError.message, - }), - { - position: "top-center", - }, - ); + reportPlaybackFailure(mediaError.code, mediaError.message); }} />
diff --git a/web/src/components/player/dynamic/AutoQualityGovernor.ts b/web/src/components/player/dynamic/AutoQualityGovernor.ts index d7805eefe2..bd0871a8d5 100644 --- a/web/src/components/player/dynamic/AutoQualityGovernor.ts +++ b/web/src/components/player/dynamic/AutoQualityGovernor.ts @@ -340,8 +340,11 @@ export class AutoQualityGovernor { private triggerDownswitch(reason: DownswitchReason): boolean { const handled = this.requestDownswitch(reason); if (handled) { - // the low stream starts with a clean record + // the low stream starts with a clean record, and the probe lets a + // recovered connection (or a wrong downswitch) return to full + // quality mid-chunk rather than at the next boundary this.resetStallHistory(); + this.armUpswitchProbe(); } return handled; } diff --git a/web/src/components/player/dynamic/DynamicVideoPlayer.tsx b/web/src/components/player/dynamic/DynamicVideoPlayer.tsx index bdcbe9a350..2aaeff5f62 100644 --- a/web/src/components/player/dynamic/DynamicVideoPlayer.tsx +++ b/web/src/components/player/dynamic/DynamicVideoPlayer.tsx @@ -484,9 +484,6 @@ export default function DynamicVideoPlayer({ } setAutoLowQuality(true); setAutoLowReason(reason === "codec" ? "codec" : "bandwidth"); - // so a recovered connection (or a wrong downswitch) returns to - // full quality mid-chunk rather than at the next boundary - governor.armUpswitchProbe(); return true; }; tryUpswitchRef.current = () => { @@ -495,7 +492,7 @@ export default function DynamicVideoPlayer({ setAutoLowReason(undefined); } }; - }, [resolvedQuality, subAvailable, governor]); + }, [resolvedQuality, subAvailable]); // persisted across sessions so a device on a known-slow connection // starts low instead of paying the first stall to find out @@ -616,7 +613,19 @@ export default function DynamicVideoPlayer({ governor.sourceLoadStarted(); }, [source, isScrubbing, governor]); + const lastChunkRef = useRef(timeRange); useEffect(() => { + // the seed effect picks the starting quality, so only a later chunk + // boundary reconsiders it + const lastChunk = lastChunkRef.current; + lastChunkRef.current = timeRange; + if ( + lastChunk.after === timeRange.after && + lastChunk.before === timeRange.before + ) { + return; + } + // a chunk boundary is where full quality may be retried, and a // natural point to persist what the governor has learned setAutoLowQuality((prev) => prev && !governor.shouldRetryMain()); diff --git a/web/src/hooks/use-webrtc-availability.ts b/web/src/hooks/use-webrtc-availability.ts index 194cbac589..4e641e115d 100644 --- a/web/src/hooks/use-webrtc-availability.ts +++ b/web/src/hooks/use-webrtc-availability.ts @@ -79,7 +79,6 @@ export function evaluateStreamWebRTCAvailability(args: { /** Console detail for a reason decided by the global availability check. */ function describeGlobalReason( reason: WebRTCUnavailableReason, - testStream: string | undefined, probeDetail: string | undefined, ): string { switch (reason) { @@ -88,7 +87,7 @@ function describeGlobalReason( case "not-configured": return "No candidates or ice_servers are set under go2rtc.webrtc."; case "unreachable": - return `The connectivity probe against stream '${testStream}' failed: ${probeDetail ?? "unknown cause"}.`; + return `The connectivity probe failed against ${probeDetail ?? "an unknown stream"}.`; default: return ""; } @@ -98,9 +97,13 @@ function describeGlobalReason( let lastProbeSignature: string | null = null; /** - * Once-per-session: browser support, go2rtc config, and a live handshake probe. + * Browser support, go2rtc config, and a live handshake probe shared by the + * page. A failed probe is retried on a later mount or when the page becomes + * visible again once its cached result has expired. */ -export function useWebRTCGloballyAvailable(): GlobalAvailability { +export function useWebRTCGloballyAvailable( + preferredStream?: string, +): GlobalAvailability { const { data: config } = useSWR("config"); const [probe, setProbe] = useState<{ state: "pending" | "pass" | "fail"; @@ -118,11 +121,19 @@ export function useWebRTCGloballyAvailable(): GlobalAvailability { ); }, [config]); - // Representative restreamed stream to probe against. - const testStream = useMemo(() => { - const streams = config?.go2rtc?.streams ?? {}; - return Object.keys(streams)[0]; - }, [config]); + // the stream being viewed is the one most likely to be online + const testStreams = useMemo(() => { + const streams = Object.keys(config?.go2rtc?.streams ?? {}); + + if (!preferredStream || !streams.includes(preferredStream)) { + return streams; + } + + return [preferredStream, ...streams.filter((s) => s !== preferredStream)]; + }, [config, preferredStream]); + const testStreamsKey = testStreams.join("\n"); + + const [retryToken, setRetryToken] = useState(0); const iceServers = useMemo( () => webRTCIceServers(config?.go2rtc?.webrtc?.ice_servers), @@ -136,9 +147,9 @@ export function useWebRTCGloballyAvailable(): GlobalAvailability { JSON.stringify({ candidates: config?.go2rtc?.webrtc?.candidates ?? [], iceServers, - testStream, + streams: Object.keys(config?.go2rtc?.streams ?? {}), }), - [config, iceServers, testStream], + [config, iceServers], ); useEffect(() => { @@ -153,11 +164,29 @@ export function useWebRTCGloballyAvailable(): GlobalAvailability { }, [probeSignature]); useEffect(() => { - if (!browserOk || !configured || !testStream) { + if (probe.state !== "fail") { + return; + } + + const onVisibilityChange = () => { + if (document.visibilityState === "visible") { + setRetryToken((token) => token + 1); + } + }; + + document.addEventListener("visibilitychange", onVisibilityChange); + return () => + document.removeEventListener("visibilitychange", onVisibilityChange); + }, [probe.state]); + + useEffect(() => { + const streams = testStreamsKey ? testStreamsKey.split("\n") : []; + + if (!browserOk || !configured || streams.length === 0) { return; } let cancelled = false; - probeWebRTCAvailability(testStream, iceServers).then((result) => { + probeWebRTCAvailability(streams, iceServers).then((result) => { if (!cancelled) { setProbe({ state: result.ok ? "pass" : "fail", @@ -168,7 +197,7 @@ export function useWebRTCGloballyAvailable(): GlobalAvailability { return () => { cancelled = true; }; - }, [browserOk, configured, testStream, iceServers]); + }, [browserOk, configured, testStreamsKey, iceServers, retryToken]); const availability = useMemo(() => { if (!browserOk) { @@ -196,9 +225,9 @@ export function useWebRTCGloballyAvailable(): GlobalAvailability { logWebRTCUnavailable( undefined, reason, - describeGlobalReason(reason, testStream, probe.detail), + describeGlobalReason(reason, probe.detail), ); - }, [config, availability, testStream, probe.detail]); + }, [config, availability, probe.detail]); return availability; } @@ -208,7 +237,8 @@ export function useWebRTCAvailableForStream( metadata: LiveStreamMetadata | null | undefined, streamName?: string, ): StreamAvailability { - const { globallyAvailable, globalReason } = useWebRTCGloballyAvailable(); + const { globallyAvailable, globalReason } = + useWebRTCGloballyAvailable(streamName); const availability = useMemo( () => diff --git a/web/src/utils/webrtcProbe.ts b/web/src/utils/webrtcProbe.ts index 637415c4f6..1160e4117d 100644 --- a/web/src/utils/webrtcProbe.ts +++ b/web/src/utils/webrtcProbe.ts @@ -1,18 +1,28 @@ import { baseUrl } from "@/api/baseUrl"; /** - * Performs a single real WebRTC handshake against go2rtc to verify that a - * media connection can actually be established (validates candidates, port - * 8555 reachability, and STUN/TURN end-to-end). Result is cached per page - * session via a module-level promise. + * Performs a real WebRTC handshake against go2rtc to verify that a media + * connection can actually be established (validates candidates, port 8555 + * reachability, and STUN/TURN end-to-end). A success is cached for the page + * session, and a failure for PROBE_FAILURE_TTL_MS and only for probes that + * start from the same stream. */ export type WebRTCProbeResult = { ok: boolean; detail?: string; + // go2rtc could not open the stream's source, which says nothing about + // whether WebRTC itself can connect + streamError?: boolean; }; -let probePromise: Promise | null = null; +const PROBE_FAILURE_TTL_MS = 30_000; + +let cachedProbe: { + promise: Promise; + firstStream: string; + failedAt?: number; +} | null = null; function describeError(err: unknown): string { return err instanceof Error ? err.message : String(err); @@ -118,23 +128,83 @@ function runProbe( pc.setRemoteDescription({ type: "answer", sdp: msg.value }).catch( (err) => fail(`remote answer rejected: ${describeError(err)}`), ); + } else if (msg.type === "error") { + cleanup({ ok: false, streamError: true, detail: msg.value }); } }); }); } +async function runProbes( + testStreams: string[], + iceServers: RTCIceServer[], + timeoutMs: number, +): Promise { + let result: WebRTCProbeResult = { ok: false, detail: "no stream to probe" }; + + for (const stream of testStreams) { + result = await runProbe(stream, iceServers, timeoutMs); + + if (!result.ok) { + result.detail = `stream '${stream}': ${result.detail}`; + } + + if (!result.streamError) { + break; + } + } + + return result; +} + +/** Probes the streams in order, moving on only when go2rtc refuses one. */ export function probeWebRTCAvailability( - testStream: string, + testStreams: string[], iceServers: RTCIceServer[], timeoutMs: number = 5000, ): Promise { - if (!probePromise) { - probePromise = runProbe(testStream, iceServers, timeoutMs); + const firstStream = testStreams[0]; + + if (cachedProbe) { + const { promise, failedAt } = cachedProbe; + + if ( + cachedProbe.firstStream === firstStream && + (failedAt === undefined || Date.now() - failedAt < PROBE_FAILURE_TTL_MS) + ) { + return promise; + } + + // a pass from any stream proves the connection, but another stream's + // failure says nothing about this one + if (failedAt === undefined) { + return promise.then((result) => + result.ok + ? result + : probeWebRTCAvailability(testStreams, iceServers, timeoutMs), + ); + } } - return probePromise; + + const entry: NonNullable = { + firstStream, + promise: runProbes(testStreams, iceServers, timeoutMs) + .catch((err): WebRTCProbeResult => ({ + ok: false, + detail: describeError(err), + })) + .then((result) => { + if (!result.ok) { + entry.failedAt = Date.now(); + } + return result; + }), + }; + cachedProbe = entry; + return entry.promise; } /** Clears the cached probe result (e.g. when go2rtc config changes). */ export function resetWebRTCProbe(): void { - probePromise = null; + cachedProbe = null; } diff --git a/web/src/views/live/LiveCameraView.tsx b/web/src/views/live/LiveCameraView.tsx index a853402cd9..2bb6489599 100644 --- a/web/src/views/live/LiveCameraView.tsx +++ b/web/src/views/live/LiveCameraView.tsx @@ -27,10 +27,7 @@ import { } from "@/components/ui/popover"; import { useResizeObserver } from "@/hooks/resize-observer"; import useKeyboardListener from "@/hooks/use-keyboard-listener"; -import { - useWebRTCAvailableForStream, - useWebRTCGloballyAvailable, -} from "@/hooks/use-webrtc-availability"; +import { useWebRTCAvailableForStream } from "@/hooks/use-webrtc-availability"; import { CameraConfig, FrigateConfig } from "@/types/frigateConfig"; import { LivePlayerError, @@ -219,15 +216,16 @@ export default function LiveCameraView({ ); const isWebRTCAvailable = webRTCAvailability.available; - // Two-way talk is the sendonly backchannel: global, not per-stream. - const { globallyAvailable: webRTCGloballyAvailable } = - useWebRTCGloballyAvailable(); - // "checking" means the probe has not answered yet, and it re-enters that on // every mount, so treating it as unavailable downgrades the saved choice. const webRTCVerdictPending = webRTCAvailability.reason === "checking"; const webRTCUsable = isWebRTCAvailable || webRTCVerdictPending; + // Two-way talk is a sendonly backchannel, so playback audio that WebRTC + // can't carry doesn't rule it out, but the stream's video must connect + const talkAvailable = + isWebRTCAvailable || webRTCAvailability.reason === "audio-codec"; + // Resolves the saved preference without overwriting it. Transient error // fallbacks layer on top in preferredLiveMode. const resolvedUserMode = useMemo(() => { @@ -398,6 +396,14 @@ export default function LiveCameraView({ const [audio, setAudio] = useSessionPersistence("liveAudio", false); const [mic, setMic] = useState(false); + + // the mic only connects through the WebRTC player, so it can't stay on + // for a stream that has lost it + useEffect(() => { + if (!talkAvailable) { + setMic(false); + } + }, [talkAvailable]); const [webRTC, setWebRTC] = useState(false); const [pip, setPip] = useState(false); const [lowBandwidth, setLowBandwidth] = useState(false); @@ -424,7 +430,7 @@ export default function LiveCameraView({ }); const preferredLiveMode = useMemo(() => { - if (mic && isWebRTCAvailable) { + if (mic && talkAvailable) { return "webrtc"; } @@ -456,6 +462,7 @@ export default function LiveCameraView({ lowBandwidth, forceLowBandwidth, mic, + talkAvailable, webRTC, isRestreamed, resolvedUserMode, @@ -489,7 +496,7 @@ export default function LiveCameraView({ } break; case "t": - if (supports2WayTalk) { + if (supports2WayTalk && talkAvailable) { setMic(!mic); return true; } @@ -737,11 +744,15 @@ export default function LiveCameraView({ Icon={mic ? FaMicrophone : FaMicrophoneSlash} isActive={mic} title={ - !webRTCGloballyAvailable - ? t("twoWayTalk.requiresWebRTC", { ns: "views/live" }) - : mic - ? t("twoWayTalk.disable", { ns: "views/live" }) - : t("twoWayTalk.enable", { ns: "views/live" }) + webRTCVerdictPending + ? t("stream.technology.unavailable.checking", { + ns: "views/live", + }) + : !talkAvailable + ? t("twoWayTalk.requiresWebRTC", { ns: "views/live" }) + : mic + ? t("twoWayTalk.disable", { ns: "views/live" }) + : t("twoWayTalk.enable", { ns: "views/live" }) } onClick={() => { setMic(!mic); @@ -749,7 +760,7 @@ export default function LiveCameraView({ setAudio(true); } }} - disabled={!cameraEnabled || debug || !webRTCGloballyAvailable} + disabled={!cameraEnabled || debug || !talkAvailable} /> )} {supportsAudioOutput && preferredLiveMode != "jsmpeg" && (