forked from dsutanto/zenai-kpc-python
421 lines
13 KiB
Python
421 lines
13 KiB
Python
"""
|
|
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)
|