Files
frigate/frigate/record/cleanup.py
T
Josh Hawkins 194e9bef69 Recording fixes (#24072)
* pin genai review frames to the main stream

* retain previews as long as either stream has recordings

* watch sub stream recording health separately from main

* reject record_sub on the same input as record and document the role

* derive recording paths from the cache segment timestamp

Recording paths carry one second of resolution, but since sub stream recording start times are resolved to fractional wall clock, anchored to the cache file mtime and chained to the previous segment's end. A stream cutting segments faster than once a second resolves consecutive segments into the same second, so two rows collide on the unique path index and the batch insert fails. The cache segment name is unique per camera stream and second by construction because ffmpeg names segments with strftime, so the recording path is now built from that timestamp while the row keeps the resolved start time. This also restores the path semantics from before sub stream recording, when start times came straight from the cache filename.

Nothing derives times from recording paths: playback offsets, stream switching, and export all use the row's start time, which is unchanged, and the recordings sync matches files by exact path string.

* keep the rest of a recording batch when one row conflicts

* only publish record_sub status when a sub stream is configured

* don't shadow camera_cfg when publishing empty cache streams

* back off restarts when a recording stream goes stale

* give the shared sub stream grace on any capture thread reset

* include segment details in recording discard warnings
2026-08-27 20:30:35 -05:00

525 lines
20 KiB
Python

"""Cleanup recordings that are expired based on retention config."""
import datetime
import itertools
import logging
import os
import threading
from multiprocessing.synchronize import Event as MpEvent
from pathlib import Path
from typing import Any
from playhouse.sqlite_ext import SqliteExtDatabase
from frigate.config import CameraConfig, FrigateConfig, RetainModeEnum
from frigate.const import (
CACHE_DIR,
CLIPS_DIR,
MAX_WAL_SIZE,
RECORD_DIR,
STREAM_TYPE_MAIN,
STREAM_TYPE_SUB,
)
from frigate.models import Previews, Recordings, ReviewSegment, UserReviewStatus
from frigate.util.builtin import clear_and_unlink
from frigate.util.media import remove_empty_directories
logger = logging.getLogger(__name__)
def _filter_reviews_for_pass(
reviews: list[Any],
now: datetime.datetime,
alerts_days: float,
detections_days: float,
) -> list[Any]:
"""Limit reviews to those still within this pass's per-severity retention window.
Review rows survive to the longer of the main and sub retention windows,
so a pass that honored all of them would let extended sub retention keep
main recordings alive too. Filtering preserves sort order for the overlap
loop in expire_existing_camera_recordings.
"""
alert_cutoff = (now - datetime.timedelta(days=alerts_days)).timestamp()
detection_cutoff = (now - datetime.timedelta(days=detections_days)).timestamp()
return [
r
for r in reviews
if r.end_time is None
or (r.end_time >= (alert_cutoff if r.severity == "alert" else detection_cutoff))
]
class RecordingCleanup(threading.Thread):
"""Cleanup existing recordings based on retention config."""
def __init__(self, config: FrigateConfig, stop_event: MpEvent) -> None:
super().__init__(name="recording_cleanup")
self.config = config
self.stop_event = stop_event
def clean_tmp_previews(self) -> None:
"""delete any previews in the cache that are more than 1 hour old."""
for p in Path(CACHE_DIR).rglob("preview_*.mp4"):
logger.debug(f"Checking preview {p}.")
if p.stat().st_mtime < (datetime.datetime.now().timestamp() - 60 * 60):
logger.debug("Deleting preview.")
clear_and_unlink(p)
def clean_tmp_clips(self) -> None:
"""delete any clips in the cache that are more than 1 hour old."""
for p in Path(os.path.join(CLIPS_DIR, "cache")).rglob("clip_*.mp4"):
logger.debug(f"Checking tmp clip {p}.")
if p.stat().st_mtime < (datetime.datetime.now().timestamp() - 60 * 60):
logger.debug("Deleting tmp clip.")
clear_and_unlink(p)
def truncate_wal(self) -> None:
"""check if the WAL needs to be manually truncated."""
# by default the WAL should be check-pointed automatically
# however, high levels of activity can prevent an opportunity
# for the checkpoint to be finished which means the WAL will grow
# without bound
# with auto checkpoint most users should never hit this
if (
os.stat(f"{self.config.database.path}-wal").st_size / (1024 * 1024)
) > MAX_WAL_SIZE:
db = SqliteExtDatabase(self.config.database.path)
db.execute_sql("PRAGMA wal_checkpoint(TRUNCATE);")
db.close()
def expire_review_segments(
self, config: CameraConfig, now: datetime.datetime
) -> set[Path]:
"""Delete review segments that are expired"""
# review rows survive to the longer of the main and sub windows so
# they stay visible while either stream still has recordings
alert_days = config.record.effective_alert_days
detection_days = config.record.effective_detection_days
alert_expire_date = (now - datetime.timedelta(days=alert_days)).timestamp()
detection_expire_date = (
now - datetime.timedelta(days=detection_days)
).timestamp()
expired_reviews = (
ReviewSegment.select(ReviewSegment.id, ReviewSegment.thumb_path)
.where(ReviewSegment.camera == config.name)
.where(
(
(ReviewSegment.severity == "alert")
& (ReviewSegment.end_time < alert_expire_date)
)
| (
(ReviewSegment.severity == "detection")
& (ReviewSegment.end_time < detection_expire_date)
)
)
.namedtuples()
)
maybe_empty_dirs = set()
thumbs_to_delete = list(map(lambda x: x[1], expired_reviews))
for thumb_path in thumbs_to_delete:
thumb_path = Path(thumb_path)
thumb_path.unlink(missing_ok=True)
maybe_empty_dirs.add(thumb_path.parent)
max_deletes = 100000
deleted_reviews_list = list(map(lambda x: x[0], expired_reviews))
for i in range(0, len(deleted_reviews_list), max_deletes):
ReviewSegment.delete().where(
ReviewSegment.id << deleted_reviews_list[i : i + max_deletes]
).execute()
UserReviewStatus.delete().where(
UserReviewStatus.review_segment
<< deleted_reviews_list[i : i + max_deletes]
).execute()
return maybe_empty_dirs
def expire_existing_camera_recordings(
self,
stream_type: str,
continuous_expire_date: float,
motion_expire_date: float,
alerts_retain_mode: RetainModeEnum,
detections_retain_mode: RetainModeEnum,
config: CameraConfig,
reviews: list[Any],
) -> tuple[set[Path], list[tuple[float, float]]]:
"""Delete recordings for one stream of an existing camera based on retention config.
Returns the directories to check for emptiness and the segments that
were kept, which the caller feeds to expire_camera_previews.
"""
# Get the timestamp for cutoff of retained days
# Get recordings to check for expiration
recordings = (
Recordings.select(
Recordings.id,
Recordings.start_time,
Recordings.end_time,
Recordings.path,
Recordings.objects,
Recordings.motion,
Recordings.dBFS,
)
.where(
(Recordings.camera == config.name)
& (Recordings.stream_type == stream_type)
& (
Recordings.start_time
< max(continuous_expire_date, motion_expire_date)
)
& (
(
(Recordings.end_time < continuous_expire_date)
& (Recordings.motion == 0)
& (Recordings.dBFS == 0)
)
| (Recordings.end_time < motion_expire_date)
)
)
.order_by(Recordings.start_time)
.namedtuples()
.iterator()
)
maybe_empty_dirs = set()
# loop over recordings and see if they overlap with any non-expired reviews
# TODO: expire segments based on segment stats according to config
review_start = 0
deleted_recordings = set()
kept_recordings: list[tuple[float, float]] = []
for recording in recordings:
keep = False
mode = None
# Now look for a reason to keep this recording segment
for idx in range(review_start, len(reviews)):
review = reviews[idx]
severity = review.severity
pre_capture = config.record.get_review_pre_capture(severity)
post_capture = config.record.get_review_post_capture(severity)
# if the review starts in the future, stop checking reviews
# and let this recording segment expire
if review.start_time - pre_capture > recording.end_time:
keep = False
break
# if the review is in progress or ends after the recording starts, keep it
# and stop looking at reviews
if (
review.end_time is None
or review.end_time + post_capture >= recording.start_time
):
keep = True
mode = (
alerts_retain_mode
if review.severity == "alert"
else detections_retain_mode
)
break
# if the review ends before this recording segment starts, skip
# this review and check the next review for an overlap.
# since the review and recordings are sorted, we can skip review
# that end before the previous recording segment started on future segments
if review.end_time + post_capture < recording.start_time:
review_start = idx
# Delete recordings outside of the retention window or based on the retention mode
if (
not keep
or (
mode == RetainModeEnum.motion
and recording.motion == 0
and recording.objects == 0
and recording.dBFS == 0
)
or (mode == RetainModeEnum.active_objects and recording.objects == 0)
):
recording_path = Path(recording.path)
recording_path.unlink(missing_ok=True)
deleted_recordings.add(recording.id)
maybe_empty_dirs.add(recording_path.parent)
else:
kept_recordings.append((recording.start_time, recording.end_time))
# expire recordings
logger.debug(f"Expiring {len(deleted_recordings)} recordings")
# delete up to 100,000 at a time
max_deletes = 100000
deleted_recordings_list = list(deleted_recordings)
for i in range(0, len(deleted_recordings_list), max_deletes):
Recordings.delete().where(
Recordings.id << deleted_recordings_list[i : i + max_deletes]
).execute()
return maybe_empty_dirs, kept_recordings
def expire_camera_previews(
self,
config: CameraConfig,
continuous_expire_date: float,
motion_expire_date: float,
kept_recordings: list[tuple[float, float]],
) -> set[Path]:
"""Delete previews that no longer have recordings on any stream.
Previews aren't recorded per stream, so the cutoffs must be the oldest
of the per stream values and kept_recordings must cover every stream,
sorted by start time. Otherwise a short main retention expires previews
the sub recordings still need.
"""
maybe_empty_dirs: set[Path] = set()
previews = (
Previews.select(
Previews.id,
Previews.start_time,
Previews.end_time,
Previews.path,
)
.where(
(Previews.camera == config.name)
& (Previews.end_time < continuous_expire_date)
& (Previews.end_time < motion_expire_date)
)
.order_by(Previews.start_time)
.namedtuples()
.iterator()
)
# expire previews
recording_start = 0
deleted_previews = set()
for preview in previews:
keep = False
# look for a reason to keep this preview
for idx in range(recording_start, len(kept_recordings)):
start_time, end_time = kept_recordings[idx]
# if the recording starts in the future, stop checking recordings
# and let this preview expire
if start_time > preview.end_time:
keep = False
break
# if the recording ends after the preview starts, keep it
# and stop looking at recordings
if end_time >= preview.start_time:
keep = True
break
# if the recording ends before this preview starts, skip
# this recording and check the next recording for an overlap.
# since the kept recordings and previews are sorted, we can skip recordings
# that end before the current preview started
if end_time < preview.start_time:
recording_start = idx
# Delete previews without any relevant recordings
if not keep:
preview_path = Path(preview.path)
preview_path.unlink(missing_ok=True)
deleted_previews.add(preview.id)
maybe_empty_dirs.add(preview_path.parent)
# expire previews
logger.debug(f"Expiring {len(deleted_previews)} previews")
# delete up to 100,000 at a time
max_deletes = 100000
deleted_previews_list = list(deleted_previews)
for i in range(0, len(deleted_previews_list), max_deletes):
Previews.delete().where(
Previews.id << deleted_previews_list[i : i + max_deletes]
).execute()
return maybe_empty_dirs
def expire_recordings(self) -> set[Path]:
"""Delete recordings based on retention config."""
logger.debug("Start expire recordings.")
logger.debug("Start deleted cameras.")
# Handle deleted cameras
expire_days = max(
self.config.record.continuous.days, self.config.record.motion.days
)
expire_before = (
datetime.datetime.now() - datetime.timedelta(days=expire_days)
).timestamp()
# enumerate the distinct cameras with one index seek each
db_cameras: list[str] = []
last_camera: str | None = None
while True:
query = Recordings.select(Recordings.camera)
if last_camera is not None:
query = query.where(Recordings.camera > last_camera)
next_camera = query.order_by(Recordings.camera.asc()).limit(1).scalar()
if next_camera is None:
break
db_cameras.append(next_camera)
last_camera = next_camera
maybe_empty_dirs = set()
deleted_recordings = set()
for camera in db_cameras:
if camera in self.config.cameras:
continue
no_camera_recordings = (
Recordings.select(
Recordings.id,
Recordings.path,
)
.where(
Recordings.camera == camera,
Recordings.end_time < expire_before,
)
.namedtuples()
.iterator()
)
for recording in no_camera_recordings:
recording_path = Path(recording.path)
recording_path.unlink(missing_ok=True)
deleted_recordings.add(recording.id)
maybe_empty_dirs.add(recording_path.parent)
logger.debug(f"Expiring {len(deleted_recordings)} recordings")
# delete up to 100,000 at a time
max_deletes = 100000
deleted_recordings_list = list(deleted_recordings)
for i in range(0, len(deleted_recordings_list), max_deletes):
Recordings.delete().where(
Recordings.id << deleted_recordings_list[i : i + max_deletes]
).execute()
logger.debug("End deleted cameras.")
logger.debug("Start all cameras.")
for camera, config in self.config.cameras.items():
logger.debug(f"Start camera: {camera}.")
now = datetime.datetime.now()
maybe_empty_dirs |= self.expire_review_segments(config, now)
continuous_expire_date = (
now - datetime.timedelta(days=config.record.continuous.days)
).timestamp()
motion_expire_date = (
now
- datetime.timedelta(
days=max(
config.record.motion.days, config.record.continuous.days
) # can't keep motion for less than continuous
)
).timestamp()
# computed here so the reviews window below covers both passes
sub_continuous_expire_date = (
now - datetime.timedelta(days=config.record.sub.continuous.days)
).timestamp()
sub_motion_expire_date = (
now
- datetime.timedelta(
days=max(
config.record.sub.motion.days,
config.record.sub.continuous.days,
) # can't keep motion for less than continuous
)
).timestamp()
# Get all the reviews to check against
reviews = (
ReviewSegment.select(
ReviewSegment.start_time,
ReviewSegment.end_time,
ReviewSegment.severity,
)
.where(
ReviewSegment.camera == camera,
# candidate recordings reach the later of the two passes'
# continuous cutoffs, so reviews must cover that whole
# range or segments overlapping recent alerts get deleted
ReviewSegment.start_time
< max(continuous_expire_date, sub_continuous_expire_date),
)
.order_by(ReviewSegment.start_time)
.namedtuples()
)
main_dirs, main_kept = self.expire_existing_camera_recordings(
STREAM_TYPE_MAIN,
continuous_expire_date,
motion_expire_date,
config.record.alerts.retain.mode,
config.record.detections.retain.mode,
config,
_filter_reviews_for_pass(
reviews,
now,
config.record.alerts.retain.days,
config.record.detections.retain.days,
),
)
maybe_empty_dirs |= main_dirs
# runs even when sub recording is disabled so old rows still
# expire
sub_dirs, sub_kept = self.expire_existing_camera_recordings(
STREAM_TYPE_SUB,
sub_continuous_expire_date,
sub_motion_expire_date,
config.record.sub.alerts.mode,
config.record.sub.detections.mode,
config,
_filter_reviews_for_pass(
reviews,
now,
config.record.sub.alerts.days,
config.record.sub.detections.days,
),
)
maybe_empty_dirs |= sub_dirs
maybe_empty_dirs |= self.expire_camera_previews(
config,
min(continuous_expire_date, sub_continuous_expire_date),
min(motion_expire_date, sub_motion_expire_date),
sorted(main_kept + sub_kept),
)
logger.debug(f"End camera: {camera}.")
logger.debug("End all cameras.")
logger.debug("End expire recordings.")
return maybe_empty_dirs
def run(self) -> None:
if self.config.safe_mode:
logger.info("Safe mode enabled, skipping recording cleanup")
return
# Expire tmp clips every minute, recordings and clean directories every hour.
for counter in itertools.cycle(range(self.config.record.expire_interval)):
if self.stop_event.wait(60):
logger.info("Exiting recording cleanup...")
break
self.clean_tmp_previews()
if counter == 0:
self.clean_tmp_clips()
maybe_empty_dirs = self.expire_recordings()
remove_empty_directories(Path(RECORD_DIR), maybe_empty_dirs)
self.truncate_wal()