""" Per-session video recorder (ByteTrack-counter style). Encode pending MP4 into SHM (tmpfs), then background-copy to VIDEO_OUTPUT_DIR when the session ends. Triggers: webhook OFF→IN/OUT start, IN/OUT→OFF stop. Uses ffmpeg subprocess when available (reliable on RK3588); falls back to OpenCV VideoWriter with multiple fourcc values. """ from __future__ import annotations import os import shutil import subprocess import threading import time from datetime import datetime, timedelta from pathlib import Path import cv2 _CODEC_CANDIDATES = ("avc1", "mp4v", "H264", "X264") class _FrameWriter: def write(self, frame) -> None: raise NotImplementedError def close(self) -> None: raise NotImplementedError def is_open(self) -> bool: raise NotImplementedError class _OpenCVFrameWriter(_FrameWriter): def __init__(self, path: str, w: int, h: int, fps: float, crf: int, preset: str): self._writer = None self._codec = None if crf > 0 or preset: opts = [] if crf > 0: opts.append(f"crf;{crf}") if preset: opts.append(f"preset;{preset}") os.environ["OPENCV_FFMPEG_WRITE_OPTIONS"] = "|".join(opts) for codec in _CODEC_CANDIDATES: writer = cv2.VideoWriter( path, cv2.VideoWriter_fourcc(*codec), fps, (int(w), int(h)) ) if writer.isOpened(): self._writer = writer self._codec = codec print(f"[RECORD] OpenCV writer codec={codec} fps={fps} size={w}x{h}") return if self._writer is not None: self._writer.release() def is_open(self) -> bool: return self._writer is not None and self._writer.isOpened() def write(self, frame) -> None: if self.is_open(): self._writer.write(frame) def close(self) -> None: if self._writer is not None: try: self._writer.release() except Exception: pass self._writer = None class _FfmpegPipeWriter(_FrameWriter): """Pipe raw BGR frames to ffmpeg (needs even width/height for libx264).""" def __init__(self, path: str, w: int, h: int, fps: float, crf: int, preset: str): self._proc = None self._w = int(w) - (int(w) % 2) self._h = int(h) - (int(h) % 2) self._use_crop = self._w != int(w) or self._h != int(h) if self._w < 2 or self._h < 2: return preset_v = preset or "medium" crf_v = str(crf) if crf > 0 else "23" cmd = [ "ffmpeg", "-y", "-hide_banner", "-loglevel", "error", "-f", "rawvideo", "-vcodec", "rawvideo", "-s", f"{self._w}x{self._h}", "-pix_fmt", "bgr24", "-r", str(fps), "-i", "-", "-an", "-c:v", "libx264", "-preset", preset_v, "-crf", crf_v, "-pix_fmt", "yuv420p", str(path), ] try: self._proc = subprocess.Popen( cmd, stdin=subprocess.PIPE, stdout=subprocess.DEVNULL, stderr=subprocess.PIPE, ) print( f"[RECORD] ffmpeg pipe libx264 fps={fps} size={self._w}x{self._h} → {path}" ) except FileNotFoundError: print("[RECORD] ffmpeg not found in PATH") except Exception as exc: print(f"[RECORD] ffmpeg pipe failed to start: {exc}") def is_open(self) -> bool: return self._proc is not None and self._proc.poll() is None def write(self, frame) -> None: if not self.is_open(): return if self._use_crop: frame = frame[: self._h, : self._w] try: self._proc.stdin.write(frame.tobytes()) except (BrokenPipeError, OSError) as exc: print(f"[RECORD] ffmpeg pipe write error: {exc}") self.close() def close(self) -> None: if self._proc is None: return try: if self._proc.stdin: self._proc.stdin.close() except Exception: pass try: self._proc.wait(timeout=30) except subprocess.TimeoutExpired: self._proc.kill() if self._proc.returncode not in (0, None): err = "" try: err = (self._proc.stderr.read() if self._proc.stderr else b"").decode( "utf-8", errors="replace" ).strip() except Exception: pass if err: print(f"[RECORD] ffmpeg exit {self._proc.returncode}: {err[:500]}") self._proc = None def _open_frame_writer(path: str, w: int, h: int, fps: float, crf: int, preset: str): if shutil.which("ffmpeg"): writer = _FfmpegPipeWriter(path, w, h, fps, crf, preset) if writer.is_open(): return writer writer.close() writer = _OpenCVFrameWriter(path, w, h, fps, crf, preset) if writer.is_open(): return writer writer.close() print(f"[RECORD] Failed to open any writer for {path}") return None def _copy_and_clean(src: str, dst: str) -> None: try: src_p = Path(src) if not src_p.is_file() or src_p.stat().st_size == 0: print(f"[RECORD] Skip move — missing or empty: {src}") src_p.unlink(missing_ok=True) return Path(dst).parent.mkdir(parents=True, exist_ok=True) shutil.copy2(src, dst) src_p.unlink(missing_ok=True) print(f"[RECORD] Moved: {src} -> {dst}") except Exception as exc: print(f"[RECORD] Move failed {src} -> {dst}: {exc}") class VideoSessionWriter: """IDLE → RECORDING → STOPPING → save/discard → IDLE.""" IDLE = "IDLE" RECORDING = "RECORDING" STOPPING = "STOPPING" def __init__( self, shm_dir: str, video_output_dir: str, w: int, h: int, fps: float, end_delay_sec: float = 1.0, retention_days: int = 3, crf: int = -1, preset: str = "", discard_empty: bool = True, ): self.shm_dir = Path(shm_dir) self.video_output_dir = Path(video_output_dir) self.w = int(w) self.h = int(h) self.fps = float(fps) if fps and fps > 1 else 15.0 self.end_delay_sec = float(end_delay_sec) self.retention_days = int(retention_days) self.crf = int(crf) self.preset = preset or "" self.discard_empty = bool(discard_empty) self.state = self.IDLE self.state_start = 0.0 self.writer: _FrameWriter | None = None self.current_path = None self.direction = "SESSION" self.counting_date = datetime.now().strftime("%Y-%m-%d") self.had_crossing = False self._frames_written = 0 self._move_threads: list[threading.Thread] = [] self.shm_dir.mkdir(parents=True, exist_ok=True) try: self.video_output_dir.mkdir(parents=True, exist_ok=True) except OSError as exc: print(f"[RECORD] WARNING: cannot create VIDEO_OUTPUT_DIR {video_output_dir}: {exc}") for stale in self.shm_dir.glob("pending_*.mp4"): try: stale.unlink() except OSError: pass self.cleanup_retention() def note_session_start(self, direction: str, counting_date: str | None = None) -> None: direction = (direction or "SESSION").strip().upper() if direction not in ("IN", "OUT"): direction = "SESSION" if counting_date: self.counting_date = counting_date if self.state != self.IDLE: print(f"[RECORD] Session start ignored (state={self.state})") return if not self._start_recording(direction): print("[RECORD] Session start failed — no encoder opened") return self.state = self.RECORDING self.state_start = time.monotonic() self.had_crossing = False self._frames_written = 0 print(f"[RECORD] Session active ({direction})") def note_crossing(self) -> None: if self.state in (self.RECORDING, self.STOPPING): self.had_crossing = True def note_session_end(self) -> None: if self.state != self.RECORDING: return if self.discard_empty and not self.had_crossing: print("[RECORD] Session end — discarding (no crossings)") self._stop_and_delete() self.state = self.IDLE return self.state = self.STOPPING self.state_start = time.monotonic() print(f"[RECORD] Session stopping (+{self.end_delay_sec}s tail)") def feed_frame(self, raw) -> None: now = time.monotonic() if self.state == self.STOPPING: if self.writer is not None and self.writer.is_open(): self.writer.write(raw) self._frames_written += 1 if now - self.state_start >= self.end_delay_sec: self._stop_and_move() self.state = self.IDLE return if self.state == self.RECORDING and self.writer is not None and self.writer.is_open(): self.writer.write(raw) self._frames_written += 1 def cleanup_retention(self) -> None: if self.retention_days <= 0: return cutoff = (datetime.now() - timedelta(days=self.retention_days)).strftime("%Y-%m-%d") try: for entry in self.video_output_dir.iterdir(): if entry.is_dir() and entry.name < cutoff: shutil.rmtree(entry, ignore_errors=True) print(f"[RECORD] Retention removed: {entry}") except OSError as exc: print(f"[RECORD] Retention scan failed: {exc}") def shutdown(self) -> None: prev = self.state self._close_writer() self.state = self.IDLE path = self.current_path self.current_path = None if not path: return keep = prev in (self.RECORDING, self.STOPPING) and ( not self.discard_empty or self.had_crossing ) if keep: dst = self._final_path() self._spawn_move(str(path), str(dst)) else: try: Path(path).unlink(missing_ok=True) print(f"[RECORD] Discarding on shutdown: {path}") except OSError: pass for t in self._move_threads: t.join(timeout=30.0) self._move_threads.clear() def _start_recording(self, direction: str) -> bool: self._close_writer() self.direction = direction ts = datetime.now().strftime("%Y%m%d_%H%M%S") path = self.shm_dir / f"pending_{direction}_{ts}.mp4" self.current_path = path self.writer = _open_frame_writer( str(path), self.w, self.h, self.fps, self.crf, self.preset ) if self.writer is not None and self.writer.is_open(): return True self._close_writer() self.current_path = None return False def _close_writer(self) -> None: if self.writer is not None: try: self.writer.close() except Exception: pass self.writer = None def _stop_and_delete(self) -> None: path = self.current_path self._close_writer() self.current_path = None if path: try: Path(path).unlink(missing_ok=True) print(f"[RECORD] Discarded: {path} ({self._frames_written} frames)") except OSError: pass def _stop_and_move(self) -> None: path = self.current_path self._close_writer() self.current_path = None if not path: return size = Path(path).stat().st_size if Path(path).is_file() else 0 print(f"[RECORD] Stopped: {path} ({self._frames_written} frames, {size} bytes)") if size == 0: Path(path).unlink(missing_ok=True) print(f"[RECORD] Skip move — empty file removed: {path}") return dst = self._final_path() self._spawn_move(str(path), str(dst)) def _final_path(self) -> Path: day = self.counting_date or datetime.now().strftime("%Y-%m-%d") date_part = day.replace("-", "") ts = datetime.now().strftime("%H%M%S") name = f"{self.direction}_{date_part}_{ts}.mp4" return self.video_output_dir / day / name def _spawn_move(self, src: str, dst: str) -> None: t = threading.Thread( target=_copy_and_clean, args=(src, dst), name="record-move", daemon=True, ) t.start() self._move_threads.append(t)