forked from dsutanto/zenai-kpc-python
update recording feature saat status IN/OUT
This commit is contained in:
1 parent
3c30ded269
commit
c80984ba80
6 files changed
+535
-13
No files matched your search
@@ -26,7 +26,7 @@ sudo apt update
|
||||
sudo apt install -y python3 python3-venv python3-pip ffmpeg libgl1
|
||||
```
|
||||
|
||||
`ffmpeg` is required for low-latency RTSP capture via OpenCV. `libgl1` is often needed for `opencv-python` on headless systems.
|
||||
`ffmpeg` is required for low-latency RTSP capture via OpenCV and for session MP4 recording (`RECORD_VIDEO=true`). `libgl1` is often needed for `opencv-python` on headless systems.
|
||||
|
||||
### RKNN model
|
||||
|
||||
|
||||
+15
-4
@@ -211,9 +211,20 @@ MAX_RECONNECT_ATTEMPTS=0
|
||||
TRACKED_PRUNE_SEC=300
|
||||
|
||||
# --- Video recording ---
|
||||
# Save annotated frames to segmented MP4 files (true/false)
|
||||
RECORD_VIDEO=false
|
||||
# Duration in seconds of each video segment file
|
||||
# Session mode (default): ByteTrack-style — encode pending MP4 in SHM_DIR, start on
|
||||
# webhook OFF→IN/OUT (or control resume), stop on IN/OUT→OFF, then background-copy to
|
||||
# VIDEO_OUTPUT_DIR/<counting_date>/. Point VIDEO_OUTPUT_DIR at NFS for remote store.
|
||||
# segment mode: legacy continuous dumps under OUTPUT_DIR every VIDEO_SEGMENT_SEC.
|
||||
RECORD_VIDEO=true
|
||||
RECORD_MODE=session
|
||||
SHM_DIR=/dev/shm/zenai-kpc-counter
|
||||
VIDEO_OUTPUT_DIR=/opt/bytetrack-counter
|
||||
RECORD_END_DELAY=1
|
||||
RECORD_DISCARD_EMPTY=true
|
||||
RECORD_RETENTION_DAYS=3
|
||||
RECORD_CRF=18
|
||||
RECORD_PRESET=slow
|
||||
# Legacy segment length (RECORD_MODE=segment only)
|
||||
VIDEO_SEGMENT_SEC=3600
|
||||
# Output video FPS (fallback if source FPS is unknown or ≤ 1)
|
||||
OUTPUT_FPS=15
|
||||
@@ -222,7 +233,7 @@ OUTPUT_FPS=15
|
||||
# Periodically write the latest annotated frame as JPEG for an external web server
|
||||
LIVE_STREAM_ENABLED=false
|
||||
# Path to the shared-memory snapshot file (served by nginx / lighttpd)
|
||||
LIVE_STREAM_FRAME_PATH=/dev/shm/byetrack-counter/live_frame.jpg
|
||||
LIVE_STREAM_FRAME_PATH=/dev/shm/zenai-kpc-counter/live_frame.jpg
|
||||
# JPEG quality (1–100)
|
||||
LIVE_STREAM_QUALITY=75
|
||||
# Write the snapshot every N frames (lower = more frequent updates)
|
||||
|
||||
+75
-6
@@ -9,6 +9,7 @@ import cv2
|
||||
import csv
|
||||
import json
|
||||
import os
|
||||
import shutil
|
||||
import signal
|
||||
import socket
|
||||
import threading
|
||||
@@ -23,6 +24,7 @@ load_dotenv()
|
||||
|
||||
from rknnlite.api import RKNNLite
|
||||
from counter_store import CounterStore
|
||||
from video_session_writer import VideoSessionWriter
|
||||
|
||||
# --- config (override via env / .env) ---
|
||||
OUTPUT_DIR = os.getenv("OUTPUT_DIR", "/opt/jetson-counter")
|
||||
@@ -96,12 +98,23 @@ RECONNECT_DELAY_SEC = int(os.getenv("RECONNECT_DELAY_SEC", "3"))
|
||||
MAX_RECONNECT_ATTEMPTS = int(os.getenv("MAX_RECONNECT_ATTEMPTS", "0"))
|
||||
TRACKED_PRUNE_SEC = int(os.getenv("TRACKED_PRUNE_SEC", "300"))
|
||||
RECORD_VIDEO = os.getenv("RECORD_VIDEO", "false").lower() == "true"
|
||||
VIDEO_SEGMENT_SEC = int(os.getenv("VIDEO_SEGMENT_SEC", "3600"))
|
||||
VIDEO_SEGMENT_SEC = int(os.getenv("VIDEO_SEGMENT_SEC", "3600")) # legacy segment mode only
|
||||
OUTPUT_FPS = int(os.getenv("OUTPUT_FPS", "15"))
|
||||
# ByteTrack-style session recording: encode in SHM, move to VIDEO_OUTPUT_DIR on session end.
|
||||
SHM_DIR = os.getenv("SHM_DIR", "/dev/shm/zenai-kpc-counter")
|
||||
VIDEO_OUTPUT_DIR = os.getenv("VIDEO_OUTPUT_DIR", OUTPUT_DIR)
|
||||
RECORD_END_DELAY = float(os.getenv("RECORD_END_DELAY", "1"))
|
||||
RECORD_DISCARD_EMPTY = os.getenv("RECORD_DISCARD_EMPTY", "true").lower() == "true"
|
||||
RECORD_RETENTION_DAYS = int(os.getenv("RECORD_RETENTION_DAYS", "3"))
|
||||
_crf_raw = os.getenv("RECORD_CRF", "").strip()
|
||||
RECORD_CRF = int(_crf_raw) if _crf_raw else -1
|
||||
RECORD_PRESET = os.getenv("RECORD_PRESET", "").strip()
|
||||
# segment = old hourly dump to OUTPUT_DIR; session = SHM + webhook/control lifecycle
|
||||
RECORD_MODE = os.getenv("RECORD_MODE", "session").strip().lower()
|
||||
|
||||
LIVE_STREAM_ENABLED = os.getenv("LIVE_STREAM_ENABLED", "false").lower() == "true"
|
||||
LIVE_STREAM_FRAME_PATH = os.getenv(
|
||||
"LIVE_STREAM_FRAME_PATH", "/dev/shm/jetson-counter/live_frame.jpg"
|
||||
"LIVE_STREAM_FRAME_PATH", f"{SHM_DIR}/live_frame.jpg"
|
||||
)
|
||||
LIVE_STREAM_QUALITY = int(os.getenv("LIVE_STREAM_QUALITY", "75"))
|
||||
LIVE_STREAM_EVERY_N = int(os.getenv("LIVE_STREAM_EVERY_N", "2"))
|
||||
@@ -145,7 +158,7 @@ CONTROL_SOCKET_PORT = int(os.getenv("CONTROL_SOCKET_PORT", "5090"))
|
||||
STATUS_WEBHOOK_ENABLED = os.getenv("STATUS_WEBHOOK_ENABLED", "false").lower() == "true"
|
||||
STATUS_WEBHOOK_FILE = os.getenv(
|
||||
"STATUS_WEBHOOK_FILE",
|
||||
"/opt/zenai-kpc-python/status_webhook_state",
|
||||
"/opt/zenai-kpc-python/status_webhook_state.json",
|
||||
)
|
||||
_active_raw = os.getenv("STATUS_WEBHOOK_ACTIVE_VALUES", "IN,OUT")
|
||||
STATUS_WEBHOOK_ACTIVE_VALUES = frozenset(
|
||||
@@ -1297,7 +1310,36 @@ def run():
|
||||
print(f"State: {STATE_FILE}")
|
||||
|
||||
if RECORD_VIDEO:
|
||||
video_writer = VideoSegmentWriter(OUTPUT_DIR, w, h, fps, VIDEO_SEGMENT_SEC)
|
||||
if RECORD_MODE == "segment":
|
||||
video_writer = VideoSegmentWriter(OUTPUT_DIR, w, h, fps, VIDEO_SEGMENT_SEC)
|
||||
print(
|
||||
f"Recording (segment): {OUTPUT_DIR} every {VIDEO_SEGMENT_SEC}s"
|
||||
)
|
||||
else:
|
||||
video_writer = VideoSessionWriter(
|
||||
shm_dir=SHM_DIR,
|
||||
video_output_dir=VIDEO_OUTPUT_DIR,
|
||||
w=w,
|
||||
h=h,
|
||||
fps=fps if fps and fps > 1 else OUTPUT_FPS,
|
||||
end_delay_sec=RECORD_END_DELAY,
|
||||
retention_days=RECORD_RETENTION_DAYS,
|
||||
crf=RECORD_CRF,
|
||||
preset=RECORD_PRESET,
|
||||
discard_empty=RECORD_DISCARD_EMPTY,
|
||||
)
|
||||
print(
|
||||
f"Recording (session): SHM={SHM_DIR} → out={VIDEO_OUTPUT_DIR} | "
|
||||
f"end_delay={RECORD_END_DELAY}s discard_empty={RECORD_DISCARD_EMPTY} | "
|
||||
f"ffmpeg={'yes' if shutil.which('ffmpeg') else 'NO'}"
|
||||
)
|
||||
# If webhook already active at startup, begin a session immediately.
|
||||
if STATUS_WEBHOOK_ENABLED and status_nyala_active:
|
||||
video_writer.note_session_start(
|
||||
status_webhook_value, store.get_counting_date()
|
||||
)
|
||||
elif not STATUS_WEBHOOK_ENABLED and counting_active:
|
||||
video_writer.note_session_start("SESSION", store.get_counting_date())
|
||||
|
||||
reconnect_count = 0
|
||||
|
||||
@@ -1338,14 +1380,37 @@ def run():
|
||||
f"[{now_str()}] Counting "
|
||||
f"{'RESUMED' if counting_active else 'PAUSED'} via control file"
|
||||
)
|
||||
if (
|
||||
isinstance(video_writer, VideoSessionWriter)
|
||||
and not STATUS_WEBHOOK_ENABLED
|
||||
):
|
||||
if counting_active:
|
||||
video_writer.note_session_start(
|
||||
"SESSION", store.get_counting_date()
|
||||
)
|
||||
else:
|
||||
video_writer.note_session_end()
|
||||
if STATUS_WEBHOOK_ENABLED:
|
||||
new_status = read_status_webhook()
|
||||
if new_status != status_webhook_value:
|
||||
prev_status = status_webhook_value
|
||||
status_webhook_value = new_status
|
||||
status_nyala_active = status_webhook_allows_counting(
|
||||
status_webhook_value
|
||||
)
|
||||
print(f"[{now_str()}] Status webhook {status_webhook_value}")
|
||||
if isinstance(video_writer, VideoSessionWriter):
|
||||
was_active = status_webhook_allows_counting(prev_status)
|
||||
if status_nyala_active and not was_active:
|
||||
video_writer.note_session_start(
|
||||
status_webhook_value, store.get_counting_date()
|
||||
)
|
||||
elif was_active and not status_nyala_active:
|
||||
video_writer.note_session_end()
|
||||
|
||||
# Raw frames only (before overlay). Session writer encodes here every frame.
|
||||
if isinstance(video_writer, VideoSessionWriter):
|
||||
video_writer.feed_frame(frame)
|
||||
|
||||
counting_allowed = counting_active and status_nyala_active
|
||||
skip_inference = False
|
||||
@@ -1578,6 +1643,8 @@ def run():
|
||||
object_crossed_frame = True
|
||||
cross_events_frame.append((tid, direction))
|
||||
crossing_times.append(mono)
|
||||
if isinstance(video_writer, VideoSessionWriter):
|
||||
video_writer.note_crossing()
|
||||
object_cross_flash1[tid] = CROSS_FLASH_FRAMES
|
||||
object_cross_flash2[tid] = CROSS_FLASH_FRAMES
|
||||
popups.append(
|
||||
@@ -1667,7 +1734,7 @@ def run():
|
||||
count_in_pulse = max(0, count_in_pulse - 1)
|
||||
count_out_pulse = max(0, count_out_pulse - 1)
|
||||
|
||||
if video_writer is not None:
|
||||
if isinstance(video_writer, VideoSegmentWriter):
|
||||
video_writer.write(frame)
|
||||
|
||||
if LIVE_STREAM_ENABLED and frame_idx % LIVE_STREAM_EVERY_N == 0:
|
||||
@@ -1723,7 +1790,9 @@ def run():
|
||||
prune_stale_tracks(object_tracked, mono)
|
||||
|
||||
cap.release()
|
||||
if video_writer is not None:
|
||||
if isinstance(video_writer, VideoSessionWriter):
|
||||
video_writer.shutdown()
|
||||
elif video_writer is not None:
|
||||
video_writer.release()
|
||||
if cross_logger:
|
||||
cross_logger.close()
|
||||
|
||||
+20
-2
@@ -87,14 +87,32 @@ MOTION_MIN_AREA_FRAC=0.002
|
||||
MOTION_HEARTBEAT_FRAMES=15
|
||||
MOTION_THRESHOLD=5.0
|
||||
|
||||
# --- Video recording ---
|
||||
# --- Video recording (ByteTrack-style session) ---
|
||||
# When true, record raw frames for each counting session
|
||||
RECORD_VIDEO=false
|
||||
# session = encode in SHM, start on webhook IN/OUT (or control ON), stop on OFF
|
||||
# segment = legacy continuous hourly files under OUTPUT_DIR
|
||||
RECORD_MODE=session
|
||||
# Pending MP4 while session is active (tmpfs / RAM)
|
||||
SHM_DIR=/dev/shm/zenai-kpc-counter
|
||||
# Finished MP4 destination (use NFS mount here if desired). Defaults to OUTPUT_DIR.
|
||||
VIDEO_OUTPUT_DIR=/opt/zenai-kpc-counter
|
||||
# Extra seconds to keep recording after session end (IN/OUT → OFF)
|
||||
RECORD_END_DELAY=1
|
||||
# Delete pending clip if no object was counted during the session
|
||||
RECORD_DISCARD_EMPTY=true
|
||||
# Delete day folders under VIDEO_OUTPUT_DIR older than this many days (0 = off)
|
||||
RECORD_RETENTION_DAYS=3
|
||||
# H.264 quality via OpenCV FFmpeg writer (empty = default). Lower CRF = better.
|
||||
RECORD_CRF=18
|
||||
RECORD_PRESET=slow
|
||||
# Legacy segment length (only when RECORD_MODE=segment)
|
||||
VIDEO_SEGMENT_SEC=3600
|
||||
OUTPUT_FPS=15
|
||||
|
||||
# --- Runtime control (start/stop counting on the fly) ---
|
||||
CONTROL_ENABLED=false
|
||||
CONTROL_FILE=/opt/bytetrack-counter/control.json
|
||||
CONTROL_FILE=/opt/zenai-kpc-counter/control.json
|
||||
CONTROL_DEFAULT_COUNTING=true
|
||||
CONTROL_POLL_SEC=1.0
|
||||
CONTROL_SOCKET_ENABLED=false
|
||||
|
||||
@@ -1,3 +1,7 @@
|
||||
# System packages (apt — not installable via pip):
|
||||
# sudo apt install -y python3 python3-venv python3-pip ffmpeg libgl1
|
||||
# ffmpeg: RTSP capture via OpenCV + session MP4 recording (VideoWriter)
|
||||
|
||||
numpy<2
|
||||
rknn-toolkit-lite2
|
||||
opencv-python
|
||||
|
||||
@@ -0,0 +1,420 @@
|
||||
"""
|
||||
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)
|
||||
Reference in new issue
Block a user