mirror of
https://github.com/blakeblackshear/frigate.git
synced 2026-09-28 19:06:52 +03:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
7b26ecc1fd | ||
|
|
a30f9b41fe | ||
|
|
dce06d5852 | ||
|
|
0928324840 | ||
|
|
386b3d31aa | ||
|
|
3f00c75bbd | ||
|
|
dd5ba8bd8e | ||
|
|
f2bfcebbc4 | ||
|
|
90426cfba8 | ||
|
|
3adf2cf48e | ||
|
|
225e014979 | ||
|
|
85a2731510 | ||
|
|
71d7741e12 | ||
|
|
7164f4b4a2 | ||
|
|
3f7e2c4686 |
@@ -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
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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,
|
||||
)
|
||||
|
||||
|
||||
@@ -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"
|
||||
|
||||
|
||||
+90
-22
@@ -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 = []
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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]
|
||||
|
||||
@@ -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 ##################################################
|
||||
####################################################################################################################
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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()
|
||||
@@ -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."""
|
||||
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
@@ -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):
|
||||
|
||||
@@ -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),
|
||||
|
||||
@@ -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,40 @@ 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",
|
||||
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");
|
||||
});
|
||||
|
||||
test("a Frigate+ Hailo model is listed by the device it was built for", async ({
|
||||
frigateApp,
|
||||
}) => {
|
||||
|
||||
@@ -2053,7 +2053,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+",
|
||||
|
||||
@@ -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,17 @@ export function synthesizeMissingFilters(
|
||||
return { ...(data as JsonObject), filters: newFilters };
|
||||
}
|
||||
|
||||
// the backend fills these from the Frigate+ model info when it loads a
|
||||
// plus:// model, so saving them would only pin values Frigate+ owns
|
||||
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 +372,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;
|
||||
}
|
||||
|
||||
@@ -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<string, string>;
|
||||
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")}
|
||||
</p>
|
||||
{selected.units.map((unit) => {
|
||||
const claimedBy = claimedElsewhere[unit.device];
|
||||
const claimedBy = selected.unlimited
|
||||
? undefined
|
||||
: claimedElsewhere[unit.device];
|
||||
|
||||
return (
|
||||
<div
|
||||
|
||||
@@ -38,9 +38,7 @@ function plusModelId(path: unknown): string | undefined {
|
||||
|
||||
type ModelSourcePickerProps = {
|
||||
path: unknown;
|
||||
// Frigate+ metadata the backend attaches to a saved model, and the only
|
||||
// reliable signal that one is active: it resolves `plus://<id>` 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;
|
||||
|
||||
@@ -169,20 +169,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 +317,9 @@ export function ModelsField(props: FieldProps) {
|
||||
);
|
||||
|
||||
return (
|
||||
<Card key={`${baseId}-${index}`} className="w-full">
|
||||
// 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
|
||||
<Card key={`${baseId}-${model.scene ?? index}`} className="w-full">
|
||||
<Collapsible
|
||||
open={open}
|
||||
onOpenChange={(nextOpen) =>
|
||||
|
||||
@@ -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<HlsConfig> = {
|
||||
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);
|
||||
}}
|
||||
/>
|
||||
</div>
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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());
|
||||
|
||||
@@ -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<FrigateConfig>("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<GlobalAvailability>(() => {
|
||||
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(
|
||||
() =>
|
||||
|
||||
@@ -1,18 +1,29 @@
|
||||
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, a failure only for PROBE_FAILURE_TTL_MS.
|
||||
*/
|
||||
|
||||
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<WebRTCProbeResult> | null = null;
|
||||
const PROBE_FAILURE_TTL_MS = 30_000;
|
||||
|
||||
// streams tried per probe when go2rtc refuses the earlier ones
|
||||
const MAX_PROBE_STREAMS = 3;
|
||||
|
||||
let cachedProbe: {
|
||||
promise: Promise<WebRTCProbeResult>;
|
||||
failedAt?: number;
|
||||
} | null = null;
|
||||
|
||||
function describeError(err: unknown): string {
|
||||
return err instanceof Error ? err.message : String(err);
|
||||
@@ -118,23 +129,67 @@ 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<WebRTCProbeResult> {
|
||||
let result: WebRTCProbeResult = { ok: false, detail: "no stream to probe" };
|
||||
|
||||
for (const stream of testStreams.slice(0, MAX_PROBE_STREAMS)) {
|
||||
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<WebRTCProbeResult> {
|
||||
if (!probePromise) {
|
||||
probePromise = runProbe(testStream, iceServers, timeoutMs);
|
||||
if (
|
||||
cachedProbe &&
|
||||
(cachedProbe.failedAt === undefined ||
|
||||
Date.now() - cachedProbe.failedAt < PROBE_FAILURE_TTL_MS)
|
||||
) {
|
||||
return cachedProbe.promise;
|
||||
}
|
||||
return probePromise;
|
||||
|
||||
const entry: NonNullable<typeof cachedProbe> = {
|
||||
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;
|
||||
}
|
||||
|
||||
@@ -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<LivePlayerMode>(() => {
|
||||
@@ -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,7 +744,11 @@ export default function LiveCameraView({
|
||||
Icon={mic ? FaMicrophone : FaMicrophoneSlash}
|
||||
isActive={mic}
|
||||
title={
|
||||
!webRTCGloballyAvailable
|
||||
webRTCVerdictPending
|
||||
? t("stream.technology.unavailable.checking", {
|
||||
ns: "views/live",
|
||||
})
|
||||
: !talkAvailable
|
||||
? t("twoWayTalk.requiresWebRTC", { ns: "views/live" })
|
||||
: mic
|
||||
? t("twoWayTalk.disable", { ns: "views/live" })
|
||||
@@ -749,7 +760,7 @@ export default function LiveCameraView({
|
||||
setAudio(true);
|
||||
}
|
||||
}}
|
||||
disabled={!cameraEnabled || debug || !webRTCGloballyAvailable}
|
||||
disabled={!cameraEnabled || debug || !talkAvailable}
|
||||
/>
|
||||
)}
|
||||
{supportsAudioOutput && preferredLiveMode != "jsmpeg" && (
|
||||
|
||||
Reference in New Issue
Block a user