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["non_logo_attributes"] = model.non_logo_attributes
|
||||||
model_dict["labelmap"] = model.merged_labelmap
|
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"]:
|
if not config["plus"]["enabled"]:
|
||||||
continue
|
continue
|
||||||
|
|
||||||
|
|||||||
@@ -51,6 +51,7 @@ def swap_runtime_config(app: FastAPI, config: FrigateConfig) -> None:
|
|||||||
|
|
||||||
if app.stats_emitter is not None:
|
if app.stats_emitter is not None:
|
||||||
app.stats_emitter.config = config
|
app.stats_emitter.config = config
|
||||||
|
app.stats_emitter.hardware_stats.set_config(config)
|
||||||
|
|
||||||
if app.dispatcher is not None:
|
if app.dispatcher is not None:
|
||||||
app.dispatcher.config = config
|
app.dispatcher.config = config
|
||||||
|
|||||||
@@ -63,7 +63,6 @@ from frigate.util.recording_coverage import (
|
|||||||
null_audio_glitches,
|
null_audio_glitches,
|
||||||
plan_clip,
|
plan_clip,
|
||||||
resolve_coverage,
|
resolve_coverage,
|
||||||
stream_has_audio,
|
|
||||||
)
|
)
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
@@ -681,15 +680,10 @@ async def _vod_response(
|
|||||||
end_ts,
|
end_ts,
|
||||||
force_discontinuity,
|
force_discontinuity,
|
||||||
)
|
)
|
||||||
intervals = resolve_coverage(camera_name, start_ts, end_ts)
|
|
||||||
|
|
||||||
# rows contradicting their stream's audio composition are
|
# rows contradicting their stream's audio composition are
|
||||||
# truncated-shutdown glitches
|
# truncated-shutdown glitches
|
||||||
main_audio = stream_has_audio(intervals, main=True)
|
|
||||||
sub_audio = stream_has_audio(intervals, main=False)
|
|
||||||
|
|
||||||
spans = build_spans(
|
spans = build_spans(
|
||||||
null_audio_glitches(intervals, main_audio, sub_audio),
|
null_audio_glitches(resolve_coverage(camera_name, start_ts, end_ts)),
|
||||||
stream_preference,
|
stream_preference,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|||||||
@@ -123,6 +123,7 @@ class ModelConfig(BaseModel):
|
|||||||
_all_attributes: list[str] = PrivateAttr()
|
_all_attributes: list[str] = PrivateAttr()
|
||||||
_all_attribute_logos: list[str] = PrivateAttr()
|
_all_attribute_logos: list[str] = PrivateAttr()
|
||||||
_model_hash: str = PrivateAttr()
|
_model_hash: str = PrivateAttr()
|
||||||
|
_plus_id: str | None = PrivateAttr(default=None)
|
||||||
|
|
||||||
@property
|
@property
|
||||||
def merged_labelmap(self) -> dict[int, str]:
|
def merged_labelmap(self) -> dict[int, str]:
|
||||||
@@ -148,6 +149,11 @@ class ModelConfig(BaseModel):
|
|||||||
def model_hash(self) -> str:
|
def model_hash(self) -> str:
|
||||||
return self._model_hash
|
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):
|
def __init__(self, **config):
|
||||||
super().__init__(**config)
|
super().__init__(**config)
|
||||||
|
|
||||||
@@ -178,6 +184,7 @@ class ModelConfig(BaseModel):
|
|||||||
os.makedirs(MODEL_CACHE_DIR, exist_ok=True)
|
os.makedirs(MODEL_CACHE_DIR, exist_ok=True)
|
||||||
|
|
||||||
model_id = self.path[7:]
|
model_id = self.path[7:]
|
||||||
|
self._plus_id = model_id
|
||||||
self.path = os.path.join(MODEL_CACHE_DIR, model_id)
|
self.path = os.path.join(MODEL_CACHE_DIR, model_id)
|
||||||
model_info_path = f"{self.path}.json"
|
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 (
|
from frigate.util.recording_coverage import (
|
||||||
build_spans,
|
build_spans,
|
||||||
known_video_codecs,
|
known_video_codecs,
|
||||||
|
null_audio_glitches,
|
||||||
|
realized_timeline,
|
||||||
resolve_coverage,
|
resolve_coverage,
|
||||||
stream_media_summary,
|
stream_media_summary,
|
||||||
)
|
)
|
||||||
@@ -84,10 +86,20 @@ class StreamRun:
|
|||||||
|
|
||||||
@dataclass
|
@dataclass
|
||||||
class _ChapterWindow:
|
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
|
start_time: float
|
||||||
end_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).
|
# 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]:
|
def _resolve_coverage(self) -> tuple[list[list[Any]], set[str], bool]:
|
||||||
"""Resolve the export range into the spans the VOD manifest will serve.
|
"""Resolve the export range into the spans the VOD manifest will serve.
|
||||||
|
|
||||||
Delegates to the same coverage resolution the manifest builder
|
Delegates to the same coverage resolution and glitch nulling the
|
||||||
uses, so what we plan around and what nginx-vod emits agree by
|
manifest builder uses, so what we plan around and what nginx-vod
|
||||||
construction. Returns the spans (each [row, start, end, is_main]),
|
emits agree by construction. Returns the spans (each [row, start,
|
||||||
the known video codecs, and whether audio survives the range.
|
end, is_main]), the known video codecs, and whether audio survives
|
||||||
|
the range.
|
||||||
Memoized: several stages of the export ask the same question, and
|
Memoized: several stages of the export ask the same question, and
|
||||||
the recordings backing a finished range do not change under us.
|
the recordings backing a finished range do not change under us.
|
||||||
"""
|
"""
|
||||||
if self._coverage is None:
|
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 = (
|
self._coverage = (
|
||||||
build_spans(intervals, self.pinned_stream),
|
build_spans(intervals, self.pinned_stream),
|
||||||
known_video_codecs(intervals),
|
known_video_codecs(intervals),
|
||||||
@@ -433,17 +448,57 @@ class RecordingExporter(threading.Thread):
|
|||||||
# hand-off to stage around
|
# hand-off to stage around
|
||||||
return True
|
return True
|
||||||
|
|
||||||
spans, codecs, keep_audio = self._resolve_coverage()
|
_spans, codecs, keep_audio = self._resolve_coverage()
|
||||||
runs = self._stream_runs(spans)
|
runs = self._planned_stream_runs()
|
||||||
|
|
||||||
# a range one stream covers end to end has nothing to hand off,
|
# a range one stream covers end to end has nothing to hand off,
|
||||||
# so it stays on the existing path however long it is
|
# so it stays on the existing path however long it is
|
||||||
if len(runs) < 2:
|
if len(runs) < 2:
|
||||||
return True
|
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)
|
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]:
|
def _stream_runs(self, spans: list[list[Any]]) -> list[StreamRun]:
|
||||||
"""Collapse the merged spans into contiguous runs of one stream type.
|
"""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))
|
clipped_end = min(float(rec.end_time), float(self.end_time))
|
||||||
if clipped_end <= clipped_start:
|
if clipped_end <= clipped_start:
|
||||||
continue
|
continue
|
||||||
|
# a staged window's keyframe lead-in plays before it
|
||||||
|
output_offset += _lead_in(rec)
|
||||||
windows.append((clipped_start, clipped_end, output_offset))
|
windows.append((clipped_start, clipped_end, output_offset))
|
||||||
output_offset += clipped_end - clipped_start
|
output_offset += clipped_end - clipped_start
|
||||||
|
|
||||||
@@ -987,9 +1044,13 @@ class RecordingExporter(threading.Thread):
|
|||||||
if duration_ms <= 0:
|
if duration_ms <= 0:
|
||||||
continue
|
continue
|
||||||
|
|
||||||
title = datetime.datetime.fromtimestamp(clipped_start, tz=tz).isoformat(
|
# a staged window's keyframe lead-in opens its chapter, with
|
||||||
timespec="seconds"
|
# 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_blocks.append(
|
||||||
"[CHAPTER]\n"
|
"[CHAPTER]\n"
|
||||||
"TIMEBASE=1/1000\n"
|
"TIMEBASE=1/1000\n"
|
||||||
@@ -1128,12 +1189,12 @@ class RecordingExporter(threading.Thread):
|
|||||||
if self.staged_runs:
|
if self.staged_runs:
|
||||||
# each run was already rendered to a temp file with a common
|
# each run was already rendered to a temp file with a common
|
||||||
# track timescale, so the concat demuxer has nothing left to
|
# track timescale, so the concat demuxer has nothing left to
|
||||||
# reconcile and every chapter offset lines up with the merged
|
# reconcile
|
||||||
# timeline the staged files reproduce
|
recordings = (
|
||||||
recordings = [
|
self._staged_chapter_windows()
|
||||||
_ChapterWindow(span_start, span_end)
|
if self.chapters not in (None, ChaptersEnum.none)
|
||||||
for _row, span_start, span_end, _is_main in self._merged_spans()
|
else []
|
||||||
]
|
)
|
||||||
playlist_lines: list[str] = [f"file '{path}'" for path in self.staged_runs]
|
playlist_lines: list[str] = [f"file '{path}'" for path in self.staged_runs]
|
||||||
ffmpeg_input = (
|
ffmpeg_input = (
|
||||||
"-y -protocol_whitelist pipe,file -f concat -safe 0 -i /dev/stdin"
|
"-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
|
# its own rows are the ones the chapters describe
|
||||||
recordings = self._get_recordings_for_range(pin)
|
recordings = self._get_recordings_for_range(pin)
|
||||||
else:
|
else:
|
||||||
# never mix streams in one playlist; use main when available
|
# an unstaged auto range resolves to at most one stream run, and
|
||||||
# and fall back to sub for expired-main history
|
# its rows are the ones the chapters describe. Main rows the
|
||||||
recordings = self._get_recordings_for_range(STREAM_TYPE_MAIN)
|
# 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)
|
recordings = self._get_recordings_for_range(STREAM_TYPE_SUB)
|
||||||
|
|
||||||
playlist_lines = []
|
playlist_lines = []
|
||||||
|
|||||||
@@ -490,9 +490,17 @@ class RecordingMaintainer(threading.Thread):
|
|||||||
)
|
)
|
||||||
reviews = reviews_by_camera[camera]
|
reviews = reviews_by_camera[camera]
|
||||||
|
|
||||||
tasks.extend(
|
# probes run concurrently, but each segment's start chains off the
|
||||||
[self.validate_and_move_segment(camera, reviews, r) for r in recordings]
|
# 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
|
# publish most recently available recording time and None if disabled
|
||||||
if stream_type == STREAM_TYPE_MAIN:
|
if stream_type == STREAM_TYPE_MAIN:
|
||||||
@@ -550,12 +558,33 @@ class RecordingMaintainer(threading.Thread):
|
|||||||
while info and info[0][0] < expire_before:
|
while info and info[0][0] < expire_before:
|
||||||
info.pop(0)
|
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:
|
def drop_segment(self, cache_path: str) -> None:
|
||||||
Path(cache_path).unlink(missing_ok=True)
|
Path(cache_path).unlink(missing_ok=True)
|
||||||
self.end_time_cache.pop(cache_path, None)
|
self.end_time_cache.pop(cache_path, None)
|
||||||
|
|
||||||
async def validate_and_move_segment(
|
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:
|
) -> dict[str, Any] | None:
|
||||||
cache_path: str = recording["cache_path"]
|
cache_path: str = recording["cache_path"]
|
||||||
start_time: datetime.datetime = recording["start_time"]
|
start_time: datetime.datetime = recording["start_time"]
|
||||||
@@ -617,6 +646,9 @@ class RecordingMaintainer(threading.Thread):
|
|||||||
async with self.probe_semaphore:
|
async with self.probe_semaphore:
|
||||||
keyframes = await get_keyframe_offsets(cache_path)
|
keyframes = await get_keyframe_offsets(cache_path)
|
||||||
|
|
||||||
|
if previous_start is not None:
|
||||||
|
await previous_start.wait()
|
||||||
|
|
||||||
start_time = self._resolve_segment_start(
|
start_time = self._resolve_segment_start(
|
||||||
camera, stream_type, start_time, duration, cache_path
|
camera, stream_type, start_time, duration, cache_path
|
||||||
)
|
)
|
||||||
@@ -654,6 +686,11 @@ class RecordingMaintainer(threading.Thread):
|
|||||||
RecordingsDataTypeEnum.valid.value,
|
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
|
record_config = self.config.cameras[camera].record
|
||||||
|
|
||||||
# sub's alerts/detections carry the retain mode directly, unlike
|
# sub's alerts/detections carry the retain mode directly, unlike
|
||||||
|
|||||||
@@ -256,6 +256,17 @@ class HardwareStats:
|
|||||||
)
|
)
|
||||||
self.update_config()
|
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:
|
def update_config(self) -> None:
|
||||||
"""Recalculate all hardware that needs to be monitored from the config."""
|
"""Recalculate all hardware that needs to be monitored from the config."""
|
||||||
names = self._scan_ffmpeg() | self._scan_detectors() | self._scan_enrichments()
|
names = self._scan_ffmpeg() | self._scan_detectors() | self._scan_enrichments()
|
||||||
|
|||||||
@@ -111,15 +111,18 @@ def get_detector_stats(
|
|||||||
) -> dict[str, dict[str, Any]]:
|
) -> dict[str, dict[str, Any]]:
|
||||||
"""Get stats for all detectors, including temperatures based on detector type."""
|
"""Get stats for all detectors, including temperatures based on detector type."""
|
||||||
detector_stats: dict[str, dict[str, Any]] = {}
|
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():
|
for name, detector in stats_tracking["detectors"].items():
|
||||||
pid = detector.detect_process.pid if detector.detect_process else None
|
pid = detector.detect_process.pid if detector.detect_process else None
|
||||||
detector_type = detector.detector_config.type
|
detector_type = detector.detector_config.type
|
||||||
|
|
||||||
# Keep track of the index for each detector type to match temperatures correctly
|
# temperatures are per physical unit, so a repeated device
|
||||||
current_index = detector_type_indices.get(detector_type, 0)
|
# ("hailo:PCIe#2", see runner_names) shares its unit's reading
|
||||||
detector_type_indices[detector_type] = current_index + 1
|
device = name.partition("#")[0]
|
||||||
|
type_devices = device_indices.setdefault(detector_type, {})
|
||||||
|
current_index = type_devices.setdefault(device, len(type_devices))
|
||||||
|
|
||||||
detector_stat = {
|
detector_stat = {
|
||||||
"inference_speed": round(detector.avg_inference_speed.value * 1000, 2), # type: ignore[attr-defined]
|
"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
|
from unittest.mock import Mock, patch
|
||||||
|
|
||||||
import frigate.genai
|
import frigate.genai
|
||||||
from frigate.config import GenAIProviderEnum
|
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.genai import GenAIClient
|
||||||
from frigate.models import Event, Recordings, ReviewSegment
|
from frigate.models import Event, Recordings, ReviewSegment
|
||||||
from frigate.stats.emitter import StatsEmitter
|
from frigate.stats.emitter import StatsEmitter
|
||||||
@@ -90,6 +92,44 @@ class TestHttpApp(BaseTestHttp):
|
|||||||
mqtt = response.json()["mqtt"]
|
mqtt = response.json()["mqtt"]
|
||||||
assert mqtt["password"] == REDACTED_CREDENTIAL_SENTINEL
|
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 ##################################################
|
################################### POST /genai/probe Endpoint ##################################################
|
||||||
####################################################################################################################
|
####################################################################################################################
|
||||||
|
|||||||
@@ -26,6 +26,7 @@ class TestSwapRuntimeConfig(unittest.TestCase):
|
|||||||
app.genai_manager.update_config.assert_called_once_with(config)
|
app.genai_manager.update_config.assert_called_once_with(config)
|
||||||
app.profile_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)
|
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)
|
self.assertIs(app.dispatcher.config, config)
|
||||||
for comm in app.dispatcher.comms:
|
for comm in app.dispatcher.comms:
|
||||||
self.assertIs(comm.config, config)
|
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))
|
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):
|
class TestStagedFileCleanup(unittest.TestCase):
|
||||||
"""A staged path must be tracked before ffmpeg can write to it."""
|
"""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"})
|
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):
|
class TestUpdateStats(HardwareStatsTestCase):
|
||||||
def run_stats(self, stats: HardwareStats) -> dict:
|
def run_stats(self, stats: HardwareStats) -> dict:
|
||||||
|
|||||||
@@ -1,7 +1,7 @@
|
|||||||
import datetime
|
import datetime
|
||||||
import sys
|
import sys
|
||||||
import unittest
|
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
|
# Mock complex imports before importing maintainer, saving originals so we can
|
||||||
# restore them after import and avoid polluting sys.modules for other tests.
|
# 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=[]
|
"frigate.record.maintainer.psutil.process_iter", return_value=[]
|
||||||
):
|
):
|
||||||
with patch("frigate.record.maintainer.logger.warning") as warn:
|
with patch("frigate.record.maintainer.logger.warning") as warn:
|
||||||
# Mock validate_and_move_segment to avoid further logic
|
# Mock validate_and_move_segment to avoid further logic.
|
||||||
maintainer.validate_and_move_segment = MagicMock()
|
# 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:
|
try:
|
||||||
await maintainer.move_files()
|
await maintainer.move_files()
|
||||||
|
|||||||
@@ -1,5 +1,6 @@
|
|||||||
"""Tests for sub stream cache segment handling in the recording maintainer."""
|
"""Tests for sub stream cache segment handling in the recording maintainer."""
|
||||||
|
|
||||||
|
import asyncio
|
||||||
import datetime
|
import datetime
|
||||||
import os
|
import os
|
||||||
import tempfile
|
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[2].timestamp(), self.T0 + 10.4, places=3)
|
||||||
self.assertAlmostEqual(calls[1].args[3].timestamp(), self.T0 + 20.8, 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):
|
async def test_genuine_gap_is_not_snapped(self):
|
||||||
maintainer = _build_chaining_maintainer(self.T0)
|
maintainer = _build_chaining_maintainer(self.T0)
|
||||||
|
|
||||||
|
|||||||
@@ -14,6 +14,10 @@ from frigate.util.file import FileLock
|
|||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
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
|
# target path -> first line of the last download error for it; every existing
|
||||||
# download function swallows its exceptions, so this is how the downloader
|
# download function swallows its exceptions, so this is how the downloader
|
||||||
# thread learns why a file is still missing
|
# thread learns why a file is still missing
|
||||||
@@ -124,7 +128,9 @@ class ModelDownloader:
|
|||||||
logger.info(f"Downloading model file from: {url}")
|
logger.info(f"Downloading model file from: {url}")
|
||||||
|
|
||||||
try:
|
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()
|
r.raise_for_status()
|
||||||
with open(temporary_filename, "wb") as f:
|
with open(temporary_filename, "wb") as f:
|
||||||
for chunk in r.iter_content(chunk_size=8192):
|
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(
|
def null_audio_glitches(intervals: list[CoverageInterval]) -> list[CoverageInterval]:
|
||||||
intervals: list[CoverageInterval], main_audio: bool, sub_audio: bool
|
|
||||||
) -> list[CoverageInterval]:
|
|
||||||
"""Treat video-only glitch rows on audio-bearing streams as no recording.
|
"""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
|
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
|
count, so a truncated video-only segment (a backend restart can flush
|
||||||
a sub-second file before any audio packet landed) poisons every
|
a sub-second file before any audio packet landed) poisons every
|
||||||
manifest that includes it. Nulling the row turns the glitch into a
|
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] = []
|
result: list[CoverageInterval] = []
|
||||||
for interval in intervals:
|
for interval in intervals:
|
||||||
main = interval.main
|
main = interval.main
|
||||||
@@ -449,9 +451,7 @@ def realized_timelines(
|
|||||||
assembles each variant's realized spans. Keyframe snapping reads the
|
assembles each variant's realized spans. Keyframe snapping reads the
|
||||||
per-row index stored at record time, so no file is touched.
|
per-row index stored at record time, so no file is touched.
|
||||||
"""
|
"""
|
||||||
main_audio = stream_has_audio(intervals, main=True)
|
nulled = null_audio_glitches(intervals)
|
||||||
sub_audio = stream_has_audio(intervals, main=False)
|
|
||||||
nulled = null_audio_glitches(intervals, main_audio, sub_audio)
|
|
||||||
|
|
||||||
return {
|
return {
|
||||||
"auto": realized_timeline(nulled, None),
|
"auto": realized_timeline(nulled, None),
|
||||||
|
|||||||
@@ -217,6 +217,23 @@ test.describe("Detection models settings @high", () => {
|
|||||||
).toBeDisabled();
|
).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 ({
|
test("adding a model appends a card with an unused scene", async ({
|
||||||
frigateApp,
|
frigateApp,
|
||||||
}) => {
|
}) => {
|
||||||
@@ -248,15 +265,13 @@ test.describe("Detection models settings @high", () => {
|
|||||||
test("a saved Frigate+ model opens on the Frigate+ tab", async ({
|
test("a saved Frigate+ model opens on the Frigate+ tab", async ({
|
||||||
frigateApp,
|
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(
|
await installRoutes(
|
||||||
frigateApp.page,
|
frigateApp.page,
|
||||||
[
|
[
|
||||||
{
|
{
|
||||||
scene: "all",
|
scene: "all",
|
||||||
devices: ["openvino:GPU.0"],
|
devices: ["openvino:GPU.0"],
|
||||||
path: "/config/model_cache/abc123",
|
path: "plus://abc123",
|
||||||
plus: PLUS_MODEL,
|
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");
|
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 ({
|
test("a Frigate+ Hailo model is listed by the device it was built for", async ({
|
||||||
frigateApp,
|
frigateApp,
|
||||||
}) => {
|
}) => {
|
||||||
|
|||||||
@@ -2053,7 +2053,7 @@
|
|||||||
"unrecognized": "This model is configured for hardware that was not found on this system: {{devices}}",
|
"unrecognized": "This model is configured for hardware that was not found on this system: {{devices}}",
|
||||||
"description": "The hardware this model runs its detection on.",
|
"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.",
|
"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": {
|
"tabs": {
|
||||||
"plus": "Frigate+",
|
"plus": "Frigate+",
|
||||||
|
|||||||
@@ -7,6 +7,7 @@
|
|||||||
*/
|
*/
|
||||||
|
|
||||||
import { RJSFSchema } from "@rjsf/utils";
|
import { RJSFSchema } from "@rjsf/utils";
|
||||||
|
import { omit } from "lodash";
|
||||||
import { applySchemaDefaults } from "@/lib/config-schema";
|
import { applySchemaDefaults } from "@/lib/config-schema";
|
||||||
import { isJsonObject } from "@/lib/utils";
|
import { isJsonObject } from "@/lib/utils";
|
||||||
import { HiddenFieldContext, JsonObject, JsonValue } from "@/types/configForm";
|
import { HiddenFieldContext, JsonObject, JsonValue } from "@/types/configForm";
|
||||||
@@ -352,6 +353,17 @@ export function synthesizeMissingFilters(
|
|||||||
return { ...(data as JsonObject), filters: newFilters };
|
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.
|
* Sanitize overrides payloads for section-specific quirks.
|
||||||
*/
|
*/
|
||||||
@@ -360,6 +372,21 @@ export function sanitizeOverridesForSection(
|
|||||||
level: string,
|
level: string,
|
||||||
overrides: unknown,
|
overrides: unknown,
|
||||||
): 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)) {
|
if (!overrides || !isJsonObject(overrides)) {
|
||||||
return overrides;
|
return overrides;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -21,7 +21,8 @@ type HardwarePickerProps = {
|
|||||||
// scopes the unit checkbox ids, since several models can list the same unit
|
// scopes the unit checkbox ids, since several models can list the same unit
|
||||||
idPrefix: string;
|
idPrefix: string;
|
||||||
devices: 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>;
|
claimedElsewhere: Record<string, string>;
|
||||||
cameraCount: number;
|
cameraCount: number;
|
||||||
disabled?: boolean;
|
disabled?: boolean;
|
||||||
@@ -84,8 +85,11 @@ export function HardwarePicker({
|
|||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
// start with the first unit no other model has taken
|
// start with the first unit no other model has taken, falling back to
|
||||||
const free = entry.units.find((unit) => !claimedElsewhere[unit.device]);
|
// 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) {
|
if (!free) {
|
||||||
onChange([]);
|
onChange([]);
|
||||||
@@ -184,7 +188,9 @@ export function HardwarePicker({
|
|||||||
{t("detectionModels.hardware.unitsDescription")}
|
{t("detectionModels.hardware.unitsDescription")}
|
||||||
</p>
|
</p>
|
||||||
{selected.units.map((unit) => {
|
{selected.units.map((unit) => {
|
||||||
const claimedBy = claimedElsewhere[unit.device];
|
const claimedBy = selected.unlimited
|
||||||
|
? undefined
|
||||||
|
: claimedElsewhere[unit.device];
|
||||||
|
|
||||||
return (
|
return (
|
||||||
<div
|
<div
|
||||||
|
|||||||
@@ -38,9 +38,7 @@ function plusModelId(path: unknown): string | undefined {
|
|||||||
|
|
||||||
type ModelSourcePickerProps = {
|
type ModelSourcePickerProps = {
|
||||||
path: unknown;
|
path: unknown;
|
||||||
// Frigate+ metadata the backend attaches to a saved model, and the only
|
// Frigate+ metadata the backend attaches to a saved model
|
||||||
// reliable signal that one is active: it resolves `plus://<id>` to a local
|
|
||||||
// cache path before serving the config back
|
|
||||||
plus?: { id: string } | null;
|
plus?: { id: string } | null;
|
||||||
// the detector this model runs on, used to filter incompatible Plus models
|
// the detector this model runs on, used to filter incompatible Plus models
|
||||||
detector?: string;
|
detector?: string;
|
||||||
|
|||||||
@@ -169,20 +169,23 @@ export function ModelsField(props: FieldProps) {
|
|||||||
[savedModels],
|
[savedModels],
|
||||||
);
|
);
|
||||||
|
|
||||||
// a model serves the cameras naming its scene, plus every camera that names
|
// a model serves the cameras naming its scene, and like the backend, the
|
||||||
// no scene at all when it is the "all" model
|
// "all" model also serves every camera whose scene has no model of its own
|
||||||
const cameraCountForScene = useCallback(
|
const cameraCountForScene = useCallback(
|
||||||
(scene: string | undefined): number => {
|
(scene: string | undefined): number => {
|
||||||
if (!cameras) {
|
if (!cameras) {
|
||||||
return 0;
|
return 0;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const modelScenes = new Set(models.map((model) => model.scene ?? "all"));
|
||||||
|
|
||||||
return Object.values(cameras).filter((camera) => {
|
return Object.values(cameras).filter((camera) => {
|
||||||
const cameraScene = camera?.detect?.scene;
|
const cameraScene = camera?.detect?.scene ?? "all";
|
||||||
return cameraScene ? cameraScene === scene : scene === "all";
|
const servedBy = modelScenes.has(cameraScene) ? cameraScene : "all";
|
||||||
|
return servedBy === (scene ?? "all");
|
||||||
}).length;
|
}).length;
|
||||||
},
|
},
|
||||||
[cameras],
|
[cameras, models],
|
||||||
);
|
);
|
||||||
|
|
||||||
const claimedByOtherModels = useCallback(
|
const claimedByOtherModels = useCallback(
|
||||||
@@ -314,7 +317,9 @@ export function ModelsField(props: FieldProps) {
|
|||||||
);
|
);
|
||||||
|
|
||||||
return (
|
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
|
<Collapsible
|
||||||
open={open}
|
open={open}
|
||||||
onOpenChange={(nextOpen) =>
|
onOpenChange={(nextOpen) =>
|
||||||
|
|||||||
@@ -26,6 +26,7 @@ import { useIsAdmin } from "@/hooks/use-is-admin";
|
|||||||
// Android native hls does not seek correctly
|
// Android native hls does not seek correctly
|
||||||
const USE_NATIVE_HLS = false;
|
const USE_NATIVE_HLS = false;
|
||||||
const HLS_MIME_TYPE = "application/vnd.apple.mpegurl" as const;
|
const HLS_MIME_TYPE = "application/vnd.apple.mpegurl" as const;
|
||||||
|
const DEFAULT_MAX_BUFFER_LENGTH_S = 10;
|
||||||
const unsupportedErrorCodes: number[] = [
|
const unsupportedErrorCodes: number[] = [
|
||||||
MediaError.MEDIA_ERR_SRC_NOT_SUPPORTED,
|
MediaError.MEDIA_ERR_SRC_NOT_SUPPORTED,
|
||||||
MediaError.MEDIA_ERR_DECODE,
|
MediaError.MEDIA_ERR_DECODE,
|
||||||
@@ -148,6 +149,42 @@ export default function HlsVideoPlayer({
|
|||||||
// a ref rather than an effect-scoped counter so the element error
|
// a ref rather than an effect-scoped counter so the element error
|
||||||
// handler can hold its toast while a recovery is still possible
|
// handler can hold its toast while a recovery is still possible
|
||||||
const mediaRecoveryBudgetRef = useRef(0);
|
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(
|
const applyVideoDimensions = useCallback(
|
||||||
(width: number, height: number) => {
|
(width: number, height: number) => {
|
||||||
@@ -225,6 +262,7 @@ export default function HlsVideoPlayer({
|
|||||||
// the element already holds a decoded frame, and keeping it visible
|
// the element already holds a decoded frame, and keeping it visible
|
||||||
// bridges the gap while the new source loads
|
// bridges the gap while the new source loads
|
||||||
const currentPlaybackRate = videoRef.current.playbackRate;
|
const currentPlaybackRate = videoRef.current.playbackRate;
|
||||||
|
failureReportedRef.current = false;
|
||||||
|
|
||||||
if (!useHlsCompat) {
|
if (!useHlsCompat) {
|
||||||
nativeRetryRef.current = 0;
|
nativeRetryRef.current = 0;
|
||||||
@@ -236,7 +274,7 @@ export default function HlsVideoPlayer({
|
|||||||
|
|
||||||
// Base HLS configuration
|
// Base HLS configuration
|
||||||
const hlsConfig: Partial<HlsConfig> = {
|
const hlsConfig: Partial<HlsConfig> = {
|
||||||
maxBufferLength: bufferLength ?? 10,
|
maxBufferLength: bufferLengthRef.current ?? DEFAULT_MAX_BUFFER_LENGTH_S,
|
||||||
maxBufferSize: 20 * 1000 * 1000,
|
maxBufferSize: 20 * 1000 * 1000,
|
||||||
startPosition: currentSource.startPosition,
|
startPosition: currentSource.startPosition,
|
||||||
};
|
};
|
||||||
@@ -269,10 +307,18 @@ export default function HlsVideoPlayer({
|
|||||||
data.details ===
|
data.details ===
|
||||||
Hls.ErrorDetails.BUFFER_INCOMPATIBLE_CODECS_ERROR ||
|
Hls.ErrorDetails.BUFFER_INCOMPATIBLE_CODECS_ERROR ||
|
||||||
data.details === Hls.ErrorDetails.BUFFER_ADD_CODEC_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;
|
return;
|
||||||
}
|
}
|
||||||
if (!isCodecError && mediaRecoveryBudgetRef.current > 0) {
|
if (mediaRecoveryBudgetRef.current > 0) {
|
||||||
mediaRecoveryBudgetRef.current -= 1;
|
mediaRecoveryBudgetRef.current -= 1;
|
||||||
hls.recoverMediaError();
|
hls.recoverMediaError();
|
||||||
}
|
}
|
||||||
@@ -308,7 +354,7 @@ export default function HlsVideoPlayer({
|
|||||||
hlsRef.current.destroy();
|
hlsRef.current.destroy();
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
}, [videoRef, hlsRef, useHlsCompat, currentSource, bufferLength]);
|
}, [videoRef, hlsRef, useHlsCompat, currentSource]);
|
||||||
|
|
||||||
// state handling
|
// state handling
|
||||||
|
|
||||||
@@ -714,15 +760,7 @@ export default function HlsVideoPlayer({
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
toast.error(
|
reportPlaybackFailure(mediaError.code, mediaError.message);
|
||||||
t("toast.error.playRecordingsFailed", {
|
|
||||||
code: mediaError.code,
|
|
||||||
message: mediaError.message,
|
|
||||||
}),
|
|
||||||
{
|
|
||||||
position: "top-center",
|
|
||||||
},
|
|
||||||
);
|
|
||||||
}}
|
}}
|
||||||
/>
|
/>
|
||||||
</div>
|
</div>
|
||||||
|
|||||||
@@ -340,8 +340,11 @@ export class AutoQualityGovernor {
|
|||||||
private triggerDownswitch(reason: DownswitchReason): boolean {
|
private triggerDownswitch(reason: DownswitchReason): boolean {
|
||||||
const handled = this.requestDownswitch(reason);
|
const handled = this.requestDownswitch(reason);
|
||||||
if (handled) {
|
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.resetStallHistory();
|
||||||
|
this.armUpswitchProbe();
|
||||||
}
|
}
|
||||||
return handled;
|
return handled;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -484,9 +484,6 @@ export default function DynamicVideoPlayer({
|
|||||||
}
|
}
|
||||||
setAutoLowQuality(true);
|
setAutoLowQuality(true);
|
||||||
setAutoLowReason(reason === "codec" ? "codec" : "bandwidth");
|
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;
|
return true;
|
||||||
};
|
};
|
||||||
tryUpswitchRef.current = () => {
|
tryUpswitchRef.current = () => {
|
||||||
@@ -495,7 +492,7 @@ export default function DynamicVideoPlayer({
|
|||||||
setAutoLowReason(undefined);
|
setAutoLowReason(undefined);
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
}, [resolvedQuality, subAvailable, governor]);
|
}, [resolvedQuality, subAvailable]);
|
||||||
|
|
||||||
// persisted across sessions so a device on a known-slow connection
|
// persisted across sessions so a device on a known-slow connection
|
||||||
// starts low instead of paying the first stall to find out
|
// starts low instead of paying the first stall to find out
|
||||||
@@ -616,7 +613,19 @@ export default function DynamicVideoPlayer({
|
|||||||
governor.sourceLoadStarted();
|
governor.sourceLoadStarted();
|
||||||
}, [source, isScrubbing, governor]);
|
}, [source, isScrubbing, governor]);
|
||||||
|
|
||||||
|
const lastChunkRef = useRef(timeRange);
|
||||||
useEffect(() => {
|
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
|
// a chunk boundary is where full quality may be retried, and a
|
||||||
// natural point to persist what the governor has learned
|
// natural point to persist what the governor has learned
|
||||||
setAutoLowQuality((prev) => prev && !governor.shouldRetryMain());
|
setAutoLowQuality((prev) => prev && !governor.shouldRetryMain());
|
||||||
|
|||||||
@@ -79,7 +79,6 @@ export function evaluateStreamWebRTCAvailability(args: {
|
|||||||
/** Console detail for a reason decided by the global availability check. */
|
/** Console detail for a reason decided by the global availability check. */
|
||||||
function describeGlobalReason(
|
function describeGlobalReason(
|
||||||
reason: WebRTCUnavailableReason,
|
reason: WebRTCUnavailableReason,
|
||||||
testStream: string | undefined,
|
|
||||||
probeDetail: string | undefined,
|
probeDetail: string | undefined,
|
||||||
): string {
|
): string {
|
||||||
switch (reason) {
|
switch (reason) {
|
||||||
@@ -88,7 +87,7 @@ function describeGlobalReason(
|
|||||||
case "not-configured":
|
case "not-configured":
|
||||||
return "No candidates or ice_servers are set under go2rtc.webrtc.";
|
return "No candidates or ice_servers are set under go2rtc.webrtc.";
|
||||||
case "unreachable":
|
case "unreachable":
|
||||||
return `The connectivity probe against stream '${testStream}' failed: ${probeDetail ?? "unknown cause"}.`;
|
return `The connectivity probe failed against ${probeDetail ?? "an unknown stream"}.`;
|
||||||
default:
|
default:
|
||||||
return "";
|
return "";
|
||||||
}
|
}
|
||||||
@@ -98,9 +97,13 @@ function describeGlobalReason(
|
|||||||
let lastProbeSignature: string | null = null;
|
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 { data: config } = useSWR<FrigateConfig>("config");
|
||||||
const [probe, setProbe] = useState<{
|
const [probe, setProbe] = useState<{
|
||||||
state: "pending" | "pass" | "fail";
|
state: "pending" | "pass" | "fail";
|
||||||
@@ -118,11 +121,19 @@ export function useWebRTCGloballyAvailable(): GlobalAvailability {
|
|||||||
);
|
);
|
||||||
}, [config]);
|
}, [config]);
|
||||||
|
|
||||||
// Representative restreamed stream to probe against.
|
// the stream being viewed is the one most likely to be online
|
||||||
const testStream = useMemo(() => {
|
const testStreams = useMemo(() => {
|
||||||
const streams = config?.go2rtc?.streams ?? {};
|
const streams = Object.keys(config?.go2rtc?.streams ?? {});
|
||||||
return Object.keys(streams)[0];
|
|
||||||
}, [config]);
|
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(
|
const iceServers = useMemo(
|
||||||
() => webRTCIceServers(config?.go2rtc?.webrtc?.ice_servers),
|
() => webRTCIceServers(config?.go2rtc?.webrtc?.ice_servers),
|
||||||
@@ -136,9 +147,9 @@ export function useWebRTCGloballyAvailable(): GlobalAvailability {
|
|||||||
JSON.stringify({
|
JSON.stringify({
|
||||||
candidates: config?.go2rtc?.webrtc?.candidates ?? [],
|
candidates: config?.go2rtc?.webrtc?.candidates ?? [],
|
||||||
iceServers,
|
iceServers,
|
||||||
testStream,
|
streams: Object.keys(config?.go2rtc?.streams ?? {}),
|
||||||
}),
|
}),
|
||||||
[config, iceServers, testStream],
|
[config, iceServers],
|
||||||
);
|
);
|
||||||
|
|
||||||
useEffect(() => {
|
useEffect(() => {
|
||||||
@@ -153,11 +164,29 @@ export function useWebRTCGloballyAvailable(): GlobalAvailability {
|
|||||||
}, [probeSignature]);
|
}, [probeSignature]);
|
||||||
|
|
||||||
useEffect(() => {
|
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;
|
return;
|
||||||
}
|
}
|
||||||
let cancelled = false;
|
let cancelled = false;
|
||||||
probeWebRTCAvailability(testStream, iceServers).then((result) => {
|
probeWebRTCAvailability(streams, iceServers).then((result) => {
|
||||||
if (!cancelled) {
|
if (!cancelled) {
|
||||||
setProbe({
|
setProbe({
|
||||||
state: result.ok ? "pass" : "fail",
|
state: result.ok ? "pass" : "fail",
|
||||||
@@ -168,7 +197,7 @@ export function useWebRTCGloballyAvailable(): GlobalAvailability {
|
|||||||
return () => {
|
return () => {
|
||||||
cancelled = true;
|
cancelled = true;
|
||||||
};
|
};
|
||||||
}, [browserOk, configured, testStream, iceServers]);
|
}, [browserOk, configured, testStreamsKey, iceServers, retryToken]);
|
||||||
|
|
||||||
const availability = useMemo<GlobalAvailability>(() => {
|
const availability = useMemo<GlobalAvailability>(() => {
|
||||||
if (!browserOk) {
|
if (!browserOk) {
|
||||||
@@ -196,9 +225,9 @@ export function useWebRTCGloballyAvailable(): GlobalAvailability {
|
|||||||
logWebRTCUnavailable(
|
logWebRTCUnavailable(
|
||||||
undefined,
|
undefined,
|
||||||
reason,
|
reason,
|
||||||
describeGlobalReason(reason, testStream, probe.detail),
|
describeGlobalReason(reason, probe.detail),
|
||||||
);
|
);
|
||||||
}, [config, availability, testStream, probe.detail]);
|
}, [config, availability, probe.detail]);
|
||||||
|
|
||||||
return availability;
|
return availability;
|
||||||
}
|
}
|
||||||
@@ -208,7 +237,8 @@ export function useWebRTCAvailableForStream(
|
|||||||
metadata: LiveStreamMetadata | null | undefined,
|
metadata: LiveStreamMetadata | null | undefined,
|
||||||
streamName?: string,
|
streamName?: string,
|
||||||
): StreamAvailability {
|
): StreamAvailability {
|
||||||
const { globallyAvailable, globalReason } = useWebRTCGloballyAvailable();
|
const { globallyAvailable, globalReason } =
|
||||||
|
useWebRTCGloballyAvailable(streamName);
|
||||||
|
|
||||||
const availability = useMemo(
|
const availability = useMemo(
|
||||||
() =>
|
() =>
|
||||||
|
|||||||
@@ -1,18 +1,29 @@
|
|||||||
import { baseUrl } from "@/api/baseUrl";
|
import { baseUrl } from "@/api/baseUrl";
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Performs a single real WebRTC handshake against go2rtc to verify that a
|
* Performs a real WebRTC handshake against go2rtc to verify that a media
|
||||||
* media connection can actually be established (validates candidates, port
|
* connection can actually be established (validates candidates, port 8555
|
||||||
* 8555 reachability, and STUN/TURN end-to-end). Result is cached per page
|
* reachability, and STUN/TURN end-to-end). A success is cached for the page
|
||||||
* session via a module-level promise.
|
* session, a failure only for PROBE_FAILURE_TTL_MS.
|
||||||
*/
|
*/
|
||||||
|
|
||||||
export type WebRTCProbeResult = {
|
export type WebRTCProbeResult = {
|
||||||
ok: boolean;
|
ok: boolean;
|
||||||
detail?: string;
|
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 {
|
function describeError(err: unknown): string {
|
||||||
return err instanceof Error ? err.message : String(err);
|
return err instanceof Error ? err.message : String(err);
|
||||||
@@ -118,23 +129,67 @@ function runProbe(
|
|||||||
pc.setRemoteDescription({ type: "answer", sdp: msg.value }).catch(
|
pc.setRemoteDescription({ type: "answer", sdp: msg.value }).catch(
|
||||||
(err) => fail(`remote answer rejected: ${describeError(err)}`),
|
(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(
|
export function probeWebRTCAvailability(
|
||||||
testStream: string,
|
testStreams: string[],
|
||||||
iceServers: RTCIceServer[],
|
iceServers: RTCIceServer[],
|
||||||
timeoutMs: number = 5000,
|
timeoutMs: number = 5000,
|
||||||
): Promise<WebRTCProbeResult> {
|
): Promise<WebRTCProbeResult> {
|
||||||
if (!probePromise) {
|
if (
|
||||||
probePromise = runProbe(testStream, iceServers, timeoutMs);
|
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). */
|
/** Clears the cached probe result (e.g. when go2rtc config changes). */
|
||||||
export function resetWebRTCProbe(): void {
|
export function resetWebRTCProbe(): void {
|
||||||
probePromise = null;
|
cachedProbe = null;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -27,10 +27,7 @@ import {
|
|||||||
} from "@/components/ui/popover";
|
} from "@/components/ui/popover";
|
||||||
import { useResizeObserver } from "@/hooks/resize-observer";
|
import { useResizeObserver } from "@/hooks/resize-observer";
|
||||||
import useKeyboardListener from "@/hooks/use-keyboard-listener";
|
import useKeyboardListener from "@/hooks/use-keyboard-listener";
|
||||||
import {
|
import { useWebRTCAvailableForStream } from "@/hooks/use-webrtc-availability";
|
||||||
useWebRTCAvailableForStream,
|
|
||||||
useWebRTCGloballyAvailable,
|
|
||||||
} from "@/hooks/use-webrtc-availability";
|
|
||||||
import { CameraConfig, FrigateConfig } from "@/types/frigateConfig";
|
import { CameraConfig, FrigateConfig } from "@/types/frigateConfig";
|
||||||
import {
|
import {
|
||||||
LivePlayerError,
|
LivePlayerError,
|
||||||
@@ -219,15 +216,16 @@ export default function LiveCameraView({
|
|||||||
);
|
);
|
||||||
const isWebRTCAvailable = webRTCAvailability.available;
|
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
|
// "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.
|
// every mount, so treating it as unavailable downgrades the saved choice.
|
||||||
const webRTCVerdictPending = webRTCAvailability.reason === "checking";
|
const webRTCVerdictPending = webRTCAvailability.reason === "checking";
|
||||||
const webRTCUsable = isWebRTCAvailable || webRTCVerdictPending;
|
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
|
// Resolves the saved preference without overwriting it. Transient error
|
||||||
// fallbacks layer on top in preferredLiveMode.
|
// fallbacks layer on top in preferredLiveMode.
|
||||||
const resolvedUserMode = useMemo<LivePlayerMode>(() => {
|
const resolvedUserMode = useMemo<LivePlayerMode>(() => {
|
||||||
@@ -398,6 +396,14 @@ export default function LiveCameraView({
|
|||||||
|
|
||||||
const [audio, setAudio] = useSessionPersistence("liveAudio", false);
|
const [audio, setAudio] = useSessionPersistence("liveAudio", false);
|
||||||
const [mic, setMic] = useState(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 [webRTC, setWebRTC] = useState(false);
|
||||||
const [pip, setPip] = useState(false);
|
const [pip, setPip] = useState(false);
|
||||||
const [lowBandwidth, setLowBandwidth] = useState(false);
|
const [lowBandwidth, setLowBandwidth] = useState(false);
|
||||||
@@ -424,7 +430,7 @@ export default function LiveCameraView({
|
|||||||
});
|
});
|
||||||
|
|
||||||
const preferredLiveMode = useMemo(() => {
|
const preferredLiveMode = useMemo(() => {
|
||||||
if (mic && isWebRTCAvailable) {
|
if (mic && talkAvailable) {
|
||||||
return "webrtc";
|
return "webrtc";
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -456,6 +462,7 @@ export default function LiveCameraView({
|
|||||||
lowBandwidth,
|
lowBandwidth,
|
||||||
forceLowBandwidth,
|
forceLowBandwidth,
|
||||||
mic,
|
mic,
|
||||||
|
talkAvailable,
|
||||||
webRTC,
|
webRTC,
|
||||||
isRestreamed,
|
isRestreamed,
|
||||||
resolvedUserMode,
|
resolvedUserMode,
|
||||||
@@ -489,7 +496,7 @@ export default function LiveCameraView({
|
|||||||
}
|
}
|
||||||
break;
|
break;
|
||||||
case "t":
|
case "t":
|
||||||
if (supports2WayTalk) {
|
if (supports2WayTalk && talkAvailable) {
|
||||||
setMic(!mic);
|
setMic(!mic);
|
||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
@@ -737,11 +744,15 @@ export default function LiveCameraView({
|
|||||||
Icon={mic ? FaMicrophone : FaMicrophoneSlash}
|
Icon={mic ? FaMicrophone : FaMicrophoneSlash}
|
||||||
isActive={mic}
|
isActive={mic}
|
||||||
title={
|
title={
|
||||||
!webRTCGloballyAvailable
|
webRTCVerdictPending
|
||||||
? t("twoWayTalk.requiresWebRTC", { ns: "views/live" })
|
? t("stream.technology.unavailable.checking", {
|
||||||
: mic
|
ns: "views/live",
|
||||||
? t("twoWayTalk.disable", { ns: "views/live" })
|
})
|
||||||
: t("twoWayTalk.enable", { ns: "views/live" })
|
: !talkAvailable
|
||||||
|
? t("twoWayTalk.requiresWebRTC", { ns: "views/live" })
|
||||||
|
: mic
|
||||||
|
? t("twoWayTalk.disable", { ns: "views/live" })
|
||||||
|
: t("twoWayTalk.enable", { ns: "views/live" })
|
||||||
}
|
}
|
||||||
onClick={() => {
|
onClick={() => {
|
||||||
setMic(!mic);
|
setMic(!mic);
|
||||||
@@ -749,7 +760,7 @@ export default function LiveCameraView({
|
|||||||
setAudio(true);
|
setAudio(true);
|
||||||
}
|
}
|
||||||
}}
|
}}
|
||||||
disabled={!cameraEnabled || debug || !webRTCGloballyAvailable}
|
disabled={!cameraEnabled || debug || !talkAvailable}
|
||||||
/>
|
/>
|
||||||
)}
|
)}
|
||||||
{supportsAudioOutput && preferredLiveMode != "jsmpeg" && (
|
{supportsAudioOutput && preferredLiveMode != "jsmpeg" && (
|
||||||
|
|||||||
Reference in New Issue
Block a user