Compare commits

10 Commits
13 changed files with 1590 additions and 221 deletions

No files matched your search

+29 -12
View File
@@ -4,10 +4,11 @@ Production deployment for **RK3588** (or compatible RKNN NPU) edge devices runni
| Component | Script | systemd unit |
|-----------|--------|--------------|
| Status webhook (IN/OUT/OFF gate) | `status_webhook.py` | `zenai-kpc-status-webhook.service` |
| RTSP counter (RKNN + ByteTrack) | `counter_live_rknn.py` | `zenai-kpc-counter.service` |
| Web dashboard (Flask) | `counter_dashboard.py` | `zenai-kpc-dashboard.service` |
Both processes share a single `.env` file and read/write the same SQLite database and state JSON.
All three processes share a single `.env` file. The counter and dashboard read/write the same SQLite database and state JSON. The webhook writes `status_webhook_state.json`, which the counter polls when `STATUS_WEBHOOK_ENABLED=true`.
---
@@ -23,10 +24,10 @@ Both processes share a single `.env` file and read/write the same SQLite databas
```bash
sudo apt update
sudo apt install -y python3 python3-venv python3-pip ffmpeg libgl1
sudo apt install -y python3 python3-venv python3-pip ffmpeg libgl1 openssl
```
`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
@@ -49,6 +50,7 @@ Default paths used by the service files and `env.example`:
├── counter_live_rknn.py
├── counter_dashboard.py
├── counter_store.py
├── status_webhook.py
├── templates/
├── venv/ # Python virtual environment (created during install)
├── .env # runtime config (not in git)
@@ -123,21 +125,23 @@ sudo mkdir -p /opt/zenai-kpc-counter /opt/models /dev/shm/zenai-kpc-counter
```bash
cd /opt/zenai-kpc-python
sudo cp zenai-kpc-counter.service zenai-kpc-dashboard.service /etc/systemd/system/
sudo cp zenai-kpc-status-webhook.service zenai-kpc-counter.service zenai-kpc-dashboard.service /etc/systemd/system/
sudo systemctl daemon-reload
sudo systemctl enable zenai-kpc-counter zenai-kpc-dashboard
sudo systemctl enable zenai-kpc-status-webhook zenai-kpc-counter zenai-kpc-dashboard
sudo systemctl start zenai-kpc-status-webhook
sudo systemctl start zenai-kpc-counter
sudo systemctl start zenai-kpc-dashboard
```
The dashboard unit starts **after** the counter unit (`After=zenai-kpc-counter.service`).
The dashboard unit starts **after** the counter unit (`After=zenai-kpc-counter.service`). Enable the webhook so it comes up on boot with the other two.
### Verify
```bash
systemctl status zenai-kpc-status-webhook
systemctl status zenai-kpc-counter
systemctl status zenai-kpc-dashboard
journalctl -u zenai-kpc-counter -f
journalctl -u zenai-kpc-status-webhook -f
```
Open the dashboard in a browser:
@@ -148,6 +152,14 @@ http://<device-ip>:5000
(Port is set by `DASHBOARD_PORT` in `.env`, default `5000`.)
Status webhook endpoints (fixed ports):
```text
http://<device-ip>:8002/?lokasi=IN|OUT&status=ON|OFF
https://<device-ip>:8443/?lokasi=IN|OUT&status=ON|OFF
http://<device-ip>:8002/status
```
---
## 5. Operations
@@ -155,6 +167,7 @@ http://<device-ip>:5000
### Restart after config change
```bash
sudo systemctl restart zenai-kpc-status-webhook
sudo systemctl restart zenai-kpc-counter
sudo systemctl restart zenai-kpc-dashboard
```
@@ -162,6 +175,7 @@ sudo systemctl restart zenai-kpc-dashboard
### View logs
```bash
journalctl -u zenai-kpc-status-webhook -n 100 --no-pager
journalctl -u zenai-kpc-counter -n 100 --no-pager
journalctl -u zenai-kpc-dashboard -n 100 --no-pager
```
@@ -169,7 +183,7 @@ journalctl -u zenai-kpc-dashboard -n 100 --no-pager
### Stop services
```bash
sudo systemctl stop zenai-kpc-dashboard zenai-kpc-counter
sudo systemctl stop zenai-kpc-dashboard zenai-kpc-counter zenai-kpc-status-webhook
```
The counter handles `SIGTERM` gracefully — it finishes the current frame, persists state to SQLite, then exits.
@@ -178,10 +192,10 @@ The counter handles `SIGTERM` gracefully — it finishes the current frame, pers
```bash
cd /opt/zenai-kpc-python
sudo systemctl stop zenai-kpc-dashboard zenai-kpc-counter
sudo systemctl stop zenai-kpc-dashboard zenai-kpc-counter zenai-kpc-status-webhook
# rsync or git pull new code
sudo ./venv/bin/pip install -r requirements.txt # if dependencies changed
sudo systemctl start zenai-kpc-counter zenai-kpc-dashboard
sudo systemctl start zenai-kpc-status-webhook zenai-kpc-counter zenai-kpc-dashboard
```
---
@@ -191,8 +205,10 @@ sudo systemctl start zenai-kpc-counter zenai-kpc-dashboard
| Symptom | Things to check |
|---------|-----------------|
| Counter won't start | `journalctl -u zenai-kpc-counter`; verify `MODEL_PATH` exists; RKNN drivers installed |
| Status webhook won't start | `journalctl -u zenai-kpc-status-webhook`; ports 8002/8443 free; `openssl` installed if certs are missing |
| No RTSP frames | Ping camera; test with `ffplay <SOURCE>`; check `OPENCV_FFMPEG_CAPTURE_OPTIONS` |
| Dashboard shows 0 count | `STATE_FILE` in `.env` must match between counter and dashboard; check file exists |
| Counting stays OFF | Check `STATUS_WEBHOOK_ENABLED` and `GET http://<device-ip>:8002/status`; `STATUS_WEBHOOK_FILE` must match the webhook state file |
| Live video blank | `LIVE_STREAM_ENABLED=true`; path matches `LIVE_STREAM_FRAME_PATH` in both processes; if using nginx, set `proxy_read_timeout` (see §7) |
| Wrong counts | Tune `LINE_Y1_FRAC`/`LINE_Y2_FRAC`, `CONF`, ByteTrack thresholds; enable `DEBUG_TRACKING=true` temporarily |
| Service keeps restarting | `journalctl -u zenai-kpc-counter -e`; often missing model, bad RTSP URL, or venv not created |
@@ -202,8 +218,9 @@ sudo systemctl start zenai-kpc-counter zenai-kpc-dashboard
```bash
cd /opt/zenai-kpc-python
source venv/bin/activate
python counter_live_rknn.py # terminal 1
python counter_dashboard.py # terminal 2
python status_webhook.py # terminal 1
python counter_live_rknn.py # terminal 2
python counter_dashboard.py # terminal 3
```
---
+32 -4
View File
@@ -186,6 +186,19 @@ CONTROL_SOCKET_HOST=127.0.0.1
# Control socket TCP port.
CONTROL_SOCKET_PORT=5090
# --- Status webhook gate (count only when status is NYALA) ---
# When true, counting runs only while the webhook state file reports NYALA
# (or STATUS_WEBHOOK_ACTIVE_VALUE). Any other/missing status pauses counting.
STATUS_WEBHOOK_ENABLED=false
# JSON file written by the external status webhook.
# Example active: {"status": "IN / OUT"}
# Example inactive: {"status": "OFF"}
STATUS_WEBHOOK_FILE=/opt/zenai-kpc-python/status_webhook_state
# Status value that allows counting (case-insensitive match).
STATUS_WEBHOOK_ACTIVE_VALUE=IN, OUT
# Expected inactive status from the webhook (shown on overlay when paused).
STATUS_WEBHOOK_INACTIVE_VALUE=OFF
# Sliding window in seconds for computing the crossing rate (objects/minute)
RATE_WINDOW_SEC=60
# Number of frames to discard at startup to let the stream buffer stabilise
@@ -194,13 +207,28 @@ WARMUP_FRAMES=30
RECONNECT_DELAY_SEC=3
# Maximum reconnection attempts (0 = infinite)
MAX_RECONNECT_ATTEMPTS=0
# Max stream-error log lines per outage (Cannot open / Stream dropped).
# After this, further errors are silent until the stream recovers; then the
# budget resets for the next outage. 0 = unlimited.
STREAM_ERROR_LOG_MAX=3
# Seconds after which a tracked but unseen object is pruned from the active set
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
@@ -209,7 +237,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)
+310 -62
View File
@@ -43,6 +43,8 @@ CONTROL_FILE = os.getenv("CONTROL_FILE", f"{_DEFAULT_DIR}/control.json")
CONTROL_DEFAULT_COUNTING = os.getenv("CONTROL_DEFAULT_COUNTING", "true").lower() == "true"
SITE_NAME = os.getenv("SITE_NAME", "LIVE")
CAMERA_NAME = os.getenv("CAMERA_NAME", "CC1")
OBJECT_LABEL = os.getenv("OBJECT_LABEL", "object")
DASHBOARD_PORT = int(os.getenv("DASHBOARD_PORT", "5000"))
DASHBOARD_HOST = os.getenv("DASHBOARD_HOST", "0.0.0.0")
@@ -266,46 +268,319 @@ def snapshots_page():
return render_template("snapshots.html", site_name=SITE_NAME, show_detect=SAVE_DETECT_SNAPSHOT)
@app.route("/api/current-counter")
def api_current_counter():
try:
with open(CURRENT_COUNTER_PATH, "r") as f:
def _empty_current():
return {
"counting_date": get_counting_date(),
"count": 0,
"count_in": 0,
"count_out": 0,
"start_time": None,
"last_detection_time": None,
}
def _load_current_state():
"""Load live counter state from current_counter.json."""
with open(CURRENT_COUNTER_PATH, "r", encoding="utf-8") as f:
data = json.load(f)
return {
"counting_date": data.get("counting_date") or get_counting_date(),
"count": int(data.get("count", 0) or 0),
"count_in": int(data.get("count_in", 0) or 0),
"count_out": int(data.get("count_out", 0) or 0),
"start_time": data.get("start_time"),
"last_detection_time": data.get("last_detection_time"),
}
def _zero_day_row(date_str, camera_name=None, object_label=None):
"""Stub row for a counting day with no activity."""
return {
"date": date_str,
"camera_name": camera_name or CAMERA_NAME,
"object_label": object_label or OBJECT_LABEL,
"total_count": 0,
"total_in": 0,
"total_out": 0,
"diff": 0,
"start_time": None,
"end_time": None,
"updated_at": None,
}
def _fill_date_gaps(rows, date_from=None, date_to=None, descending=True):
"""Ensure every date in [date_from, date_to] is present; missing days are 0."""
by_date = {r["date"]: r for r in rows}
if date_to is None:
date_to = get_counting_date()
if date_from is None:
if not by_date:
return rows
date_from = min(by_date.keys())
sample = next(iter(by_date.values()), None)
cam = sample["camera_name"] if sample else CAMERA_NAME
lab = sample["object_label"] if sample else OBJECT_LABEL
filled = []
d = datetime.strptime(date_from, "%Y-%m-%d").date()
end = datetime.strptime(date_to, "%Y-%m-%d").date()
if d > end:
return rows
while d <= end:
key = d.isoformat()
if key in by_date:
filled.append(by_date[key])
else:
filled.append(_zero_day_row(key, cam, lab))
d += timedelta(days=1)
if descending:
filled.reverse()
return filled
def _query_history(date_from=None, date_to=None, days=None, limit=None, offset=0):
"""Query daily_counters with optional date range and pagination.
Missing dates in the requested range are filled with zero totals so the API
always returns a continuous series.
"""
conn = get_db()
cur = conn.cursor()
clauses = []
params = []
if days is not None and date_from is None and date_to is None:
date_from = (datetime.now() - timedelta(days=days)).date().isoformat()
if date_from:
clauses.append("counting_date >= ?")
params.append(date_from)
if date_to:
clauses.append("counting_date <= ?")
params.append(date_to)
where = f"WHERE {' AND '.join(clauses)}" if clauses else ""
sql = f"""
SELECT counting_date, camera_name, object_label,
total_count, total_in, total_out, start_time, end_time, updated_at
FROM daily_counters
{where}
ORDER BY counting_date ASC
"""
cur.execute(sql, params)
rows = [
{
"date": row["counting_date"],
"camera_name": row["camera_name"],
"object_label": row["object_label"],
"total_count": row["total_count"],
"total_in": row["total_in"],
"total_out": row["total_out"],
"diff": (row["total_in"] or 0) - (row["total_out"] or 0),
"start_time": row["start_time"],
"end_time": row["end_time"],
"updated_at": row["updated_at"],
}
for row in cur.fetchall()
]
conn.close()
fill_to = date_to or get_counting_date()
rows = _fill_date_gaps(rows, date_from=date_from, date_to=fill_to, descending=True)
total = len(rows)
if limit is not None:
rows = rows[offset: offset + limit]
return rows, total
@app.route("/api/current")
@app.route("/api/current-counter")
def api_current():
"""Current counting-day totals from live state file."""
try:
state = _load_current_state()
return jsonify(
{
"success": True,
"counting_date": data.get("counting_date"),
"count": data.get("count", 0),
"count_in": data.get("count_in", 0),
"count_out": data.get("count_out", 0),
"start_time": data.get("start_time"),
"last_detection_time": data.get("last_detection_time"),
"active": True,
"site_name": SITE_NAME,
"counting": _read_counting_flag() if CONTROL_ENABLED else True,
**state,
}
)
except FileNotFoundError:
return jsonify(
{
"success": False,
"error": "No active counter",
"count": 0,
"count_in": 0,
"count_out": 0,
"counting_date": None,
"success": True,
"site_name": SITE_NAME,
"counting": _read_counting_flag() if CONTROL_ENABLED else True,
"active": False,
"error": "No active counter — state file missing",
**_empty_current(),
}
), 200
except Exception as e:
return jsonify(
{
"success": False,
"error": str(e),
"count": 0,
"count_in": 0,
"count_out": 0,
"counting_date": None,
"error": f"Failed to read current counter from {CURRENT_COUNTER_PATH}: {e}",
"site_name": SITE_NAME,
**_empty_current(),
}
), 500
@app.route("/api/history")
def api_history():
"""
Historical daily counters.
Query params:
days – last N calendar days (default 30; ignored if date_from/date_to set)
date_from – inclusive YYYY-MM-DD
date_to – inclusive YYYY-MM-DD
limit – page size (default: all matching)
offset – page offset (default 0)
"""
try:
days = request.args.get("days", type=int)
date_from = request.args.get("date_from")
date_to = request.args.get("date_to")
limit = request.args.get("limit", type=int)
offset = request.args.get("offset", 0, type=int)
if date_from:
try:
datetime.strptime(date_from, "%Y-%m-%d")
except ValueError:
return jsonify({
"success": False,
"error": {
"code": "VALIDATION_ERROR",
"message": f"Invalid date_from '{date_from}'. Use YYYY-MM-DD.",
},
}), 400
if date_to:
try:
datetime.strptime(date_to, "%Y-%m-%d")
except ValueError:
return jsonify({
"success": False,
"error": {
"code": "VALIDATION_ERROR",
"message": f"Invalid date_to '{date_to}'. Use YYYY-MM-DD.",
},
}), 400
if days is None and date_from is None and date_to is None:
days = 30
if offset < 0:
offset = 0
if limit is not None and limit < 1:
return jsonify({
"success": False,
"error": {
"code": "VALIDATION_ERROR",
"message": "limit must be a positive integer",
},
}), 400
rows, total = _query_history(
date_from=date_from,
date_to=date_to,
days=days,
limit=limit,
offset=offset,
)
payload = {
"success": True,
"site_name": SITE_NAME,
"filters": {
"days": days,
"date_from": date_from,
"date_to": date_to,
},
"count": len(rows),
"data": rows,
}
if limit is not None:
payload["pagination"] = {
"offset": offset,
"limit": limit,
"total": total,
"has_next": offset + limit < total,
"has_prev": offset > 0,
}
else:
payload["total"] = total
return jsonify(payload)
except sqlite3.OperationalError as e:
return jsonify({
"success": False,
"error": {
"code": "DATABASE_UNAVAILABLE",
"message": f"Database unavailable at {DB_PATH}: {e}",
},
"data": [],
"count": 0,
"total": 0,
}), 200
except Exception as e:
return jsonify({
"success": False,
"error": {
"code": "INTERNAL_ERROR",
"message": str(e),
},
}), 500
@app.route("/api/history/<counting_date>")
def api_history_day(counting_date):
"""Single counting-day record by YYYY-MM-DD."""
try:
datetime.strptime(counting_date, "%Y-%m-%d")
except ValueError:
return jsonify({
"success": False,
"error": {
"code": "VALIDATION_ERROR",
"message": f"Invalid counting_date '{counting_date}'. Use YYYY-MM-DD.",
},
}), 400
try:
rows, _ = _query_history(date_from=counting_date, date_to=counting_date)
data = rows[0] if rows else _zero_day_row(counting_date)
return jsonify({
"success": True,
"site_name": SITE_NAME,
"data": data,
})
except sqlite3.OperationalError as e:
return jsonify({
"success": False,
"error": {
"code": "DATABASE_UNAVAILABLE",
"message": f"Database unavailable at {DB_PATH}: {e}",
},
}), 503
except Exception as e:
return jsonify({
"success": False,
"error": {
"code": "INTERNAL_ERROR",
"message": str(e),
},
}), 500
@app.route("/api/summary")
def api_summary():
try:
@@ -395,27 +670,19 @@ def api_daily_data():
try:
days = request.args.get("days", 30, type=int)
date_from = (datetime.now() - timedelta(days=days)).date().isoformat()
conn = get_db()
cur = conn.cursor()
cur.execute(
"""
SELECT counting_date, total_count, total_in, total_out
FROM daily_counters
WHERE counting_date >= ?
ORDER BY counting_date ASC
""",
(date_from,),
)
date_to = get_counting_date()
rows, _ = _query_history(date_from=date_from, date_to=date_to)
# ASC for charts
rows = list(reversed(rows))
daily_data = [
{
"date": row["counting_date"],
"date": row["date"],
"total_count": row["total_count"],
"total_in": row["total_in"],
"total_out": row["total_out"],
}
for row in cur.fetchall()
for row in rows
]
conn.close()
return jsonify(daily_data)
except sqlite3.OperationalError:
return jsonify([]), 200
@@ -426,27 +693,18 @@ def api_daily_data():
@app.route("/api/available-dates")
def api_available_dates():
try:
conn = get_db()
cur = conn.cursor()
cur.execute(
"""
SELECT counting_date, total_count, total_in, total_out, start_time, end_time
FROM daily_counters
ORDER BY counting_date DESC
"""
)
rows, _ = _query_history()
dates = [
{
"date": row["counting_date"],
"date": row["date"],
"total_count": row["total_count"],
"total_in": row["total_in"],
"total_out": row["total_out"],
"start_time": row["start_time"],
"end_time": row["end_time"],
}
for row in cur.fetchall()
for row in rows
]
conn.close()
return jsonify(dates)
except sqlite3.OperationalError:
return jsonify([]), 200
@@ -496,19 +754,9 @@ def export_daily_xlsx():
try:
days = request.args.get("days", 30, type=int)
date_from = (datetime.now() - timedelta(days=days)).date().isoformat()
conn = get_db()
cur = conn.cursor()
cur.execute(
"""
SELECT counting_date, total_count, total_in, total_out, start_time, end_time
FROM daily_counters
WHERE counting_date >= ?
ORDER BY counting_date ASC
""",
(date_from,),
)
rows = cur.fetchall()
conn.close()
date_to = get_counting_date()
rows, _ = _query_history(date_from=date_from, date_to=date_to)
rows = list(reversed(rows)) # ASC for export
except sqlite3.OperationalError as e:
return jsonify({"success": False, "error": f"Database unavailable: {e}"}), 503
except Exception as e:
@@ -520,11 +768,11 @@ def export_daily_xlsx():
_style_header(ws, [("A", "Date"), ("B", "Total"), ("C", "In"), ("D", "Out"), ("E", "Diff (In-Out)"), ("F", "First Count"), ("G", "Last Count")])
for r_idx, row in enumerate(rows, 2):
ws.cell(row=r_idx, column=1, value=row["counting_date"])
ws.cell(row=r_idx, column=1, value=row["date"])
ws.cell(row=r_idx, column=2, value=row["total_count"])
ws.cell(row=r_idx, column=3, value=row["total_in"])
ws.cell(row=r_idx, column=4, value=row["total_out"])
ws.cell(row=r_idx, column=5, value=(row["total_in"] or 0) - (row["total_out"] or 0))
ws.cell(row=r_idx, column=5, value=row["diff"])
ws.cell(row=r_idx, column=6, value=row["start_time"])
ws.cell(row=r_idx, column=7, value=row["end_time"])
+244 -47
View File
@@ -4,13 +4,16 @@ Runs on RK3588 hardware with RKNN model (320×320 input).
Uses ByteTrack (Kalman filter + two-stage IoU association) for tracking.
"""
import contextlib
import numpy as np
import cv2
import csv
import json
import os
import shutil
import signal
import socket
import sys
import threading
import time
from collections import deque
@@ -23,6 +26,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")
@@ -94,14 +98,29 @@ RATE_WINDOW_SEC = int(os.getenv("RATE_WINDOW_SEC", "60"))
WARMUP_FRAMES = int(os.getenv("WARMUP_FRAMES", "30"))
RECONNECT_DELAY_SEC = int(os.getenv("RECONNECT_DELAY_SEC", "3"))
MAX_RECONNECT_ATTEMPTS = int(os.getenv("MAX_RECONNECT_ATTEMPTS", "0"))
# Max stream-error prints per outage (Cannot open / Stream dropped). After this,
# further errors are silent until the stream recovers, then the budget resets.
# 0 = unlimited (old behavior).
STREAM_ERROR_LOG_MAX = int(os.getenv("STREAM_ERROR_LOG_MAX", "3"))
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"))
@@ -139,6 +158,20 @@ CONTROL_SOCKET_ENABLED = os.getenv("CONTROL_SOCKET_ENABLED", "false").lower() ==
CONTROL_SOCKET_HOST = os.getenv("CONTROL_SOCKET_HOST", "127.0.0.1")
CONTROL_SOCKET_PORT = int(os.getenv("CONTROL_SOCKET_PORT", "5090"))
# --- Status webhook gate (direction-locked when IN or OUT) ---
# When enabled, counting runs only while status is IN or OUT, and only that
# direction is counted (IN → in crossings, OUT → out crossings). OFF pauses all.
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.json",
)
_active_raw = os.getenv("STATUS_WEBHOOK_ACTIVE_VALUES", "IN,OUT")
STATUS_WEBHOOK_ACTIVE_VALUES = frozenset(
part.strip().upper() for part in _active_raw.split(",") if part.strip()
)
STATUS_WEBHOOK_INACTIVE_VALUE = os.getenv("STATUS_WEBHOOK_INACTIVE_VALUE", "OFF").strip().upper()
CROSS_FLASH_FRAMES = 12
POPUP_LIFETIME = 20
LINE_PULSE_FRAMES = 12
@@ -750,6 +783,21 @@ def read_counting_flag(default=True):
return default
def read_status_webhook():
"""Return normalized status from webhook file, or inactive value on error."""
try:
with open(STATUS_WEBHOOK_FILE, "r", encoding="utf-8") as f:
data = json.load(f)
status = str(data.get("status", "")).strip().upper()
return status or STATUS_WEBHOOK_INACTIVE_VALUE
except Exception:
return STATUS_WEBHOOK_INACTIVE_VALUE
def status_webhook_allows_counting(status):
return status in STATUS_WEBHOOK_ACTIVE_VALUES
def write_control_file(counting):
"""Create/update the control file atomically (used to seed defaults)."""
try:
@@ -825,9 +873,36 @@ def start_control_socket():
return srv
@contextlib.contextmanager
def _quiet_opencv_stderr():
"""Silence OpenCV/FFmpeg stderr (DESCRIBE failed, cap.cpp WARN, etc.)."""
saved_fd = None
devnull_fd = None
try:
stderr_fd = sys.stderr.fileno()
saved_fd = os.dup(stderr_fd)
devnull_fd = os.open(os.devnull, os.O_WRONLY)
os.dup2(devnull_fd, stderr_fd)
except (AttributeError, OSError, ValueError):
yield
return
try:
yield
finally:
os.dup2(saved_fd, stderr_fd)
os.close(saved_fd)
os.close(devnull_fd)
def open_capture(source):
if source.lower().startswith(("rtsp://", "http://")):
os.environ["OPENCV_FFMPEG_CAPTURE_OPTIONS"] = RTSP_FFMPEG_OPTIONS
# When STREAM_ERROR_LOG_MAX > 0, native OpenCV/FFmpeg stderr is suppressed;
# stream_error_log emits at most that many user-facing messages per outage.
if STREAM_ERROR_LOG_MAX > 0:
with _quiet_opencv_stderr():
cap = cv2.VideoCapture(source, cv2.CAP_FFMPEG)
else:
cap = cv2.VideoCapture(source, cv2.CAP_FFMPEG)
cap.set(cv2.CAP_PROP_BUFFERSIZE, 1)
return cap
@@ -1126,6 +1201,46 @@ def draw_popups(img, popups, frame_idx):
return alive
class StreamErrorLog:
"""Rate-limit stream error prints to STREAM_ERROR_LOG_MAX per outage.
After the cap, messages are suppressed until note_up() when the stream
is healthy again; the next drop starts a fresh budget.
"""
def __init__(self, max_logs=STREAM_ERROR_LOG_MAX):
self.max_logs = int(max_logs)
self._emitted = 0
self._suppressed = 0
def emit(self, msg: str) -> None:
if self.max_logs <= 0:
print(msg)
return
if self._emitted < self.max_logs:
self._emitted += 1
print(msg)
if self._emitted >= self.max_logs:
print(
f"[{now_str()}] Stream errors capped at {self.max_logs} "
f"this outage — further messages suppressed until stream recovers"
)
else:
self._suppressed += 1
def note_up(self) -> None:
if self._suppressed > 0:
print(
f"[{now_str()}] Stream recovered "
f"(suppressed {self._suppressed} error log(s) during outage)"
)
self._emitted = 0
self._suppressed = 0
stream_error_log = StreamErrorLog()
def connect_stream(source, warmup=WARMUP_FRAMES):
attempts = 0
while not shutdown_requested:
@@ -1136,7 +1251,9 @@ def connect_stream(source, warmup=WARMUP_FRAMES):
raise RuntimeError(
f"Cannot open source after {attempts} attempts: {source}"
)
print(f"Cannot open source, retry in {RECONNECT_DELAY_SEC}s...")
stream_error_log.emit(
f"Cannot open source, retry in {RECONNECT_DELAY_SEC}s..."
)
time.sleep(RECONNECT_DELAY_SEC)
continue
if warmup > 0 and source.lower().startswith(("rtsp://", "http://")):
@@ -1146,6 +1263,7 @@ def connect_stream(source, warmup=WARMUP_FRAMES):
fps = cap.get(cv2.CAP_PROP_FPS)
if not fps or fps <= 1:
fps = OUTPUT_FPS
stream_error_log.note_up()
return cap, w, h, fps
return None, 0, 0, OUTPUT_FPS
@@ -1215,11 +1333,14 @@ def run():
frames_since_infer = 0
video_writer = None
crossing_times = deque()
counter_in = 0
counter_out = 0
# Overlay shows daily store totals (resume-safe, resets on counting-day change).
counter_in = store.current_count_in
counter_out = store.current_count_out
last_snapshot_cleanup = 0.0
counting_active = True
status_nyala_active = True
status_webhook_value = STATUS_WEBHOOK_INACTIVE_VALUE
last_control_poll = 0.0
control_socket = None
if CONTROL_ENABLED:
@@ -1235,6 +1356,13 @@ def run():
control_socket = start_control_socket()
except Exception as exc:
print(f"[{now_str()}] Failed to start control socket: {exc}")
if STATUS_WEBHOOK_ENABLED:
status_webhook_value = read_status_webhook()
status_nyala_active = status_webhook_allows_counting(status_webhook_value)
print(
f"Status webhook enabled | file={STATUS_WEBHOOK_FILE} | "
f"status={status_webhook_value}"
)
cap, w, h, fps = connect_stream(SOURCE)
if cap is None:
@@ -1258,7 +1386,36 @@ def run():
print(f"State: {STATE_FILE}")
if RECORD_VIDEO:
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
@@ -1268,8 +1425,9 @@ def run():
if not IS_LIVE:
break
reconnect_count += 1
print(
f"Stream dropped (attempt {reconnect_count}), reconnecting in {RECONNECT_DELAY_SEC}s..."
stream_error_log.emit(
f"Stream dropped (attempt {reconnect_count}), "
f"reconnecting in {RECONNECT_DELAY_SEC}s..."
)
cap.release()
time.sleep(RECONNECT_DELAY_SEC)
@@ -1287,15 +1445,53 @@ def run():
cross_events_frame = []
detect_events_frame = []
if CONTROL_ENABLED and (now - last_control_poll) >= CONTROL_POLL_SEC:
if (CONTROL_ENABLED or STATUS_WEBHOOK_ENABLED) and (
now - last_control_poll
) >= CONTROL_POLL_SEC:
last_control_poll = now
if CONTROL_ENABLED:
new_flag = read_counting_flag(CONTROL_DEFAULT_COUNTING)
if new_flag != counting_active:
counting_active = new_flag
print(f"[{now_str()}] Counting {'RESUMED' if counting_active else 'PAUSED'} via control file")
print(
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
if not counting_active:
if not counting_allowed:
skip_inference = True
elif MOTION_DETECTION_ENABLED:
gray = cv2.cvtColor(frame, cv2.COLOR_BGR2GRAY)
@@ -1357,29 +1553,23 @@ def run():
if object_det_to_track
else ""
)
crossing_str1 = ""
crossing_str2 = ""
crossing_str = ""
if object_cross_state:
seen_in_ids = sorted(
line1_ids = sorted(
tid for tid, st in object_cross_state.items() if st.get("seen_in")
)
seen_out_ids = sorted(
line2_ids = sorted(
tid for tid, st in object_cross_state.items() if st.get("seen_out")
)
counted_ids = sorted(
tid for tid, st in object_cross_state.items() if st.get("counted")
)
if seen_in_ids:
crossing_str1 = f" seen_in: {seen_in_ids}"
if seen_out_ids:
crossing_str2 = f" seen_out: {seen_out_ids}"
if counted_ids:
crossing_str2 += f" counted: {counted_ids}"
if line1_ids:
crossing_str += f" line1_crossed: {line1_ids}"
if line2_ids:
crossing_str += f" line2_crossed: {line2_ids}"
print(
f"[DEBUG F{frame_idx}] dets={len(object_boxes_xyxy)} "
f"tracks={len(object_track_map)} "
f"line_y1={line_y1} line_y2={line_y2}{scores_str}"
f"{tracks_str}{crossing_str1}{crossing_str2}"
f"{tracks_str}{crossing_str}"
)
# --- crossing detection only on DETECTED objects this frame ---
@@ -1455,6 +1645,12 @@ def run():
# Prefer decisive motion if both count triggers fire in one jump.
count_in = in_down and not st["counted"]
count_out = out_up and not st["counted"]
if STATUS_WEBHOOK_ENABLED:
# Webhook IN → only IN counts; OUT → only OUT counts.
if status_webhook_value == "IN":
count_out = False
elif status_webhook_value == "OUT":
count_in = False
if count_in and count_out:
if cy >= prev_cy:
count_out = False
@@ -1490,36 +1686,28 @@ def run():
object_tracked[tid] = (cx, cy, mono)
continue
seq = (
"out→in"
if direction == "in" and st["seen_out"]
else (
"in→out"
if direction == "out" and st["seen_in"]
else ("in-only" if direction == "in" else "out-only")
)
)
if os.getenv("DEBUG_TRACKING", "").lower() == "true":
line_n = 1 if direction == "in" else 2
line_y = line_y1 if direction == "in" else line_y2
print(
f"[DEBUG F{frame_idx}] CROSS DETECTED: tid={tid} "
f"prev_cy={prev_cy:.1f} -> cy={cy:.1f} dir={direction} "
f"seq={seq}"
+ (
" (inherit)"
if inherited_this_frame
else ""
)
f"prev_cy={prev_cy:.1f} -> cy={cy:.1f} "
f"line={line_n} y={line_y} dir={direction}"
)
recent.append((frame_idx, cx))
st["counted"] = True
_, _, counted = store.record_object_crossing(tid, direction)
counter_in = store.current_count_in
counter_out = store.current_count_out
if not counted:
# Dedup in the store rejected this event; skip overlay/CSV.
object_tracked[tid] = (cx, cy, mono)
continue
if direction == "in":
counter_in += 1
count_in_pulse = COUNT_PULSE_FRAMES
else:
counter_out += 1
count_out_pulse = COUNT_PULSE_FRAMES
store.record_object_crossing(tid, direction)
if cross_logger:
cross_logger.write_row(
[
@@ -1532,6 +1720,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(
@@ -1568,7 +1758,9 @@ def run():
if object_crossed_frame:
line_pulse = LINE_PULSE_FRAMES
display_total = store.display_total()
# Keep overlay aligned with daily store (also resets after cutoff).
counter_in = store.current_count_in
counter_out = store.current_count_out
while crossing_times and mono - crossing_times[0] > RATE_WINDOW_SEC:
crossing_times.popleft()
rate = (len(crossing_times) / RATE_WINDOW_SEC * 60) if crossing_times else 0.0
@@ -1597,7 +1789,10 @@ def run():
)
popups = draw_popups(frame, popups, frame_idx)
if CONTROL_ENABLED and not counting_active:
if (CONTROL_ENABLED or STATUS_WEBHOOK_ENABLED) and not counting_allowed:
if STATUS_WEBHOOK_ENABLED and not status_nyala_active:
badge = status_webhook_value
else:
badge = "COUNTING PAUSED"
(bw, bh), _ = cv2.getTextSize(badge, cv2.FONT_HERSHEY_SIMPLEX, 0.6, 2)
bx = w // 2 - bw // 2
@@ -1616,7 +1811,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:
@@ -1672,7 +1867,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()
+87 -4
View File
@@ -41,6 +41,12 @@ class CounterStore:
self.db = sqlite3.connect(db_path, check_same_thread=False)
self._init_db()
self.current_state = self._load_state()
with self.state_lock:
self._fill_missing_days()
counting_date = self.get_counting_date()
if self.current_state is None:
self._start_new_day(counting_date)
self._persist_day()
def _init_db(self):
cur = self.db.cursor()
@@ -88,11 +94,23 @@ class CounterStore:
state.setdefault('count_in', 0)
state.setdefault('count_out', 0)
state.setdefault('count', state['count_in'] + state['count_out'])
state.setdefault('counted_event_ids', [])
# Track IDs are process-local and restart from 1 after every process
# start. Persisted dedup keys like "5_in" would silently block new
# counts that reuse those IDs, so clear them on resume while keeping
# the day's totals.
stale_ids = state.get('counted_event_ids') or []
if stale_ids:
self.log(
f"Cleared {len(stale_ids)} persisted track dedup keys "
f"(track IDs reset on restart)"
)
state['counted_event_ids'] = []
self.log(
f"Resumed {current_date} with total={state['count']} "
f"(in={state['count_in']} out={state['count_out']})"
)
with open(self.state_file, 'w', encoding='utf-8') as f:
json.dump(state, f, indent=2, ensure_ascii=False)
return state
except Exception as exc:
self.log(f'Failed to load state file: {exc}')
@@ -125,8 +143,65 @@ class CounterStore:
self.save_state()
self.log(f'Started counting day {counting_date} ({self.object_label})')
def _ensure_day_row(self, counting_date, start_time=None, end_time=None):
"""Insert a zero row for counting_date if it does not already exist."""
now = datetime.now().isoformat()
cur = self.db.cursor()
cur.execute(
"""
INSERT OR IGNORE INTO daily_counters
(counting_date, camera_name, object_label,
total_count, total_in, total_out, start_time, end_time)
VALUES (?, ?, ?, 0, 0, 0, ?, ?)
""",
(
counting_date,
self.camera_name,
self.object_label,
start_time or now,
end_time or now,
),
)
self.db.commit()
def _fill_missing_days(self):
"""Backfill any missing counting dates from first DB row through today as 0."""
today = self.get_counting_date()
cur = self.db.cursor()
cur.execute(
"""
SELECT counting_date FROM daily_counters
WHERE camera_name = ? AND object_label = ?
ORDER BY counting_date ASC
""",
(self.camera_name, self.object_label),
)
existing = {row[0] for row in cur.fetchall()}
if not existing:
self._ensure_day_row(today)
return
start = datetime.strptime(min(existing), '%Y-%m-%d').date()
end = datetime.strptime(today, '%Y-%m-%d').date()
filled = 0
d = start
while d <= end:
key = d.isoformat()
if key not in existing:
self._ensure_day_row(key)
filled += 1
d += timedelta(days=1)
if filled:
self.log(
f'Backfilled {filled} zero-activity day(s) '
f'{start.isoformat()}..{end.isoformat()}'
)
def record_object_crossing(self, track_id, direction):
"""Record an object crossing a counting line. direction: 'in' | 'out'."""
"""Record an object crossing a counting line. direction: 'in' | 'out'.
Returns (total_count, day_started, counted).
"""
with self.state_lock:
counting_date = self.get_counting_date()
day_started = False
@@ -134,6 +209,7 @@ class CounterStore:
self._start_new_day(counting_date)
day_started = True
counted = False
event_key = f"{track_id}_{direction}"
if event_key not in self.current_state['counted_event_ids']:
self.current_state['count'] += 1
@@ -142,6 +218,7 @@ class CounterStore:
else:
self.current_state['count_out'] += 1
self.current_state['counted_event_ids'].append(event_key)
counted = True
self.log(
f'Counted {direction} (track {track_id}) | {counting_date} '
f'total: {self.current_state["count"]} '
@@ -151,7 +228,7 @@ class CounterStore:
self.current_state['last_detection_time'] = datetime.now().isoformat()
self.save_state()
return self.current_state['count'], day_started
return self.current_state['count'], day_started, counted
def _persist_day(self):
state = self.current_state
@@ -182,13 +259,19 @@ class CounterStore:
while not self.shutdown_event.is_set():
time.sleep(60)
with self.state_lock:
self._fill_missing_days()
counting_date = self.get_counting_date()
if self.current_state is None:
self._start_new_day(counting_date)
self._persist_day()
continue
if self.current_state['counting_date'] != self.get_counting_date():
if self.current_state['counting_date'] != counting_date:
self.log('Daily cutoff reached - finalizing day totals')
self._persist_day()
self.current_state = None
self.save_state()
self._start_new_day(counting_date)
self._persist_day()
def start_cutoff_watcher(self):
t = threading.Thread(target=self.cutoff_watcher_loop, daemon=True)
+91 -92
View File
@@ -1,152 +1,151 @@
# =============================================================================
# Edge RK3588 production counter + dashboard
# Shared config for: counter_live_rknn_bytetrack.py + counter_dashboard.py
# Copy to .env on device: cp config.env.example .env && nano .env
# ZenAI KPC edge counter + dashboard
# Shared config for: counter_live_rknn.py + counter_dashboard.py
#
# On device:
# cp env.example .env && nano .env
#
# Install path (systemd): /opt/zenai-kpc-python
# Data path: /opt/zenai-kpc-counter
# See DEPLOY.md for full setup instructions.
# =============================================================================
# --- Core paths ---
# Root output directory (logs, DB, video, CSV)
OUTPUT_DIR=/opt/zenai-kpc-bt-counter
# SQLite database path for daily counter records & crossing logs
DB_PATH=/tmp/counter.db
# JSON file persisting the current active counting day state
STATE_FILE=/tmp/current_counter.json
OUTPUT_DIR=/opt/zenai-kpc-counter
DB_PATH=/opt/zenai-kpc-counter/counter.db
STATE_FILE=/opt/zenai-kpc-counter/current_counter.json
# --- Input source ---
# RTSP / HTTP live stream, or a local video file path
#SOURCE=rtsp://user:pass@192.168.0.100:554/stream1
SOURCE=rtsp://10.38.30.64:8554/my_stream
# FFmpeg capture options passed to cv2.VideoCapture (RTSP low-latency flags)
SOURCE=rtsp://192.168.192.250:8554/my_stream_1
OPENCV_FFMPEG_CAPTURE_OPTIONS=rtsp_transport;tcp|fflags;nobuffer|flags;low_delay
# --- RKNN model ---
# Path to exported .rknn model (YOLO format, e.g. yolo11n.rknn)
MODEL_PATH=/opt/models/zenai_kac_sukawarna_20260702.rknn
# Input image size for the model (square, e.g. 320 → 320×320)
MODEL_PATH=/opt/models/zenai_kac_sukawarna_20260716.rknn
IMGSZ=320
# Use FP16 inference on NPU (true/false); currently unused in ByteTrack variant
HALF=false
# NPU core mask: 1=core0, 2=core1, 3=core0+core1, 7=all three
CORE_MASK=1
# Compute device index (reserved; not used at runtime)
CORE_MASK=7
DEVICE=0
# --- YOLO decoder ---
# Number of object classes the model outputs
NUM_CLASSES=4
# Apply sigmoid to raw class scores (true/false); set true if model head uses BCE logits
NUM_CLASSES=1
SCORE_SIGMOID=false
# --- Detection ---
# Confidence threshold – detections below this are discarded before NMS
CONF=0.5
CONF=0.6
NMS_IOU=0.9
# --- ByteTrack tracking ---
# Detections with score >= this get priority matching in the first association stage
TRACK_HIGH_THRESH=0.5
# Detections with score between this and TRACK_HIGH_THRESH are matched in the second stage
TRACK_LOW_THRESH=0.3
# IoU threshold for the first-stage association (0–1). Higher = stricter overlap required
TRACK_MATCH_THRESH=0.7
# Frames a track survives without a match before being permanently removed
TRACK_HIGH_THRESH=0.6
TRACK_LOW_THRESH=0.4
TRACK_MATCH_THRESH=0.8
TRACK_BUFFER=60
# Minimum consecutive (or total) hits needed before a track is considered confirmed
TRACK_MIN_HITS=3
# --- ID-switch counting guards ---
DEDUP_FRAMES=15
DEDUP_PX=60
INHERIT_SEC=1.0
INHERIT_PX=60
INHERIT_CROSS_MAX_DY=40
# --- Display ---
# Site name shown on the dashboard header (top-right)
SITE_NAME=ZenAi
# --- Object class names ---
# Camera / location identifier shown in HUD and stored in DB
CAMERA_NAME=ZenAi
# Label used for batch grouping in the database
OBJECT_LABEL=karung
# Class name for the counted object (must match model class order)
CLASS_OBJECT=karung
# Model class ID for the object being counted (default 0)
OBJECT_CLASS_ID=0
# --- Line crossing ---
# Two horizontal counting lines with sequence logic:
# Line 1 (IN): top→down counts IN (OUT→IN sequence, or IN-only)
# Line 2 (OUT): bottom→up counts OUT (IN→OUT sequence, or OUT-only)
# Passing the opposite line first only arms; the destination line still counts
# even if the opposite line was never crossed.
# Fixed y-coordinate for line 1/IN (overrides LINE_Y1_FRAC if set)
LINE_Y1=
# Fraction of frame height for line 1 (default 0.33)
LINE_Y1_FRAC=0.70
# Fixed y-coordinate for line 2 (overrides LINE_Y2_FRAC if set)
LINE_Y2=
# Fraction of frame height for line 2 (default 0.66)
LINE_Y2_FRAC=0.30
# --- Counting day management ---
# Daily cutoff time (HH:MM) – a new counting day starts after this time and the
# previous day's counter_in / counter_out totals are finalized in the database.
# CUTOFF_TIME is an alias used by the dashboard; DAILY_CUTOFF_TIME takes priority in counter_live_rknn.py.
DAILY_CUTOFF_TIME=20:00
CUTOFF_TIME=20:00
DAILY_CUTOFF_TIME=00:00
CUTOFF_TIME=00:00
# --- CSV export ---
# Write per-crossing events to a CSV file (true/false)
EXPORT_CSV=false
# Path where the crossing CSV is written
CROSS_CSV=/tmp/crossings.csv
CROSS_CSV=/opt/zenai-kpc-counter/crossings.csv
# --- Crossing snapshots ---
SAVE_CROSS_SNAPSHOT=true
SAVE_DETECT_SNAPSHOT=false
CROSS_SNAPSHOT_DIR=/opt/zenai-kpc-counter/snapshots
CROSS_SNAPSHOT_QUALITY=85
CROSS_SNAPSHOT_MAX_FILES=1000
CROSS_SNAPSHOT_MAX_AGE_DAYS=3
CROSS_SNAPSHOT_CLEANUP_SEC=3600
# --- Rate / performance ---
# Enable motion detection pre-filter: skip inference on frames with no movement
# (true/false, default: false). When enabled, frames below MOTION_THRESHOLD are
# skipped, saving NPU/CPU load.
MOTION_DETECTION_ENABLED=true
# Mean absolute pixel difference threshold (0–255) to consider a frame as having
# motion. Lower = more sensitive. Default 5.0.
MOTION_PIXEL_DELTA=25
MOTION_MIN_AREA_FRAC=0.002
MOTION_HEARTBEAT_FRAMES=15
MOTION_THRESHOLD=5.0
# Sliding window in seconds for computing the crossing rate (objects/minute)
RATE_WINDOW_SEC=60
# Number of frames to discard at startup to let the stream buffer stabilise
WARMUP_FRAMES=30
# Delay in seconds between stream reconnection attempts
RECONNECT_DELAY_SEC=3
# Maximum reconnection attempts (0 = infinite)
MAX_RECONNECT_ATTEMPTS=0
# Seconds after which a tracked but unseen object is pruned from the active set
TRACKED_PRUNE_SEC=300
# --- Video recording ---
# Save annotated frames to segmented MP4 files (true/false)
# --- Video recording (ByteTrack-style session) ---
# When true, record raw frames for each counting session
RECORD_VIDEO=false
# Duration in seconds of each video segment file
# 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 video FPS (fallback if source FPS is unknown or ≤ 1)
OUTPUT_FPS=15
# --- Runtime control (start/stop counting on the fly) ---
CONTROL_ENABLED=false
CONTROL_FILE=/opt/zenai-kpc-counter/control.json
CONTROL_DEFAULT_COUNTING=true
CONTROL_POLL_SEC=1.0
CONTROL_SOCKET_ENABLED=false
CONTROL_SOCKET_HOST=127.0.0.1
CONTROL_SOCKET_PORT=5090
# --- Status webhook gate (count only when status is IN or OUT) ---
STATUS_WEBHOOK_ENABLED=true
STATUS_WEBHOOK_FILE=/opt/zenai-kpc-python/status_webhook_state.json
STATUS_WEBHOOK_ACTIVE_VALUES=IN,OUT
STATUS_WEBHOOK_INACTIVE_VALUE=OFF
# Sliding window in seconds for computing the crossing rate (objects/minute)
RATE_WINDOW_SEC=60
WARMUP_FRAMES=30
RECONNECT_DELAY_SEC=3
MAX_RECONNECT_ATTEMPTS=0
# Max stream-error logs per outage; resets when stream recovers. 0 = unlimited.
STREAM_ERROR_LOG_MAX=3
TRACKED_PRUNE_SEC=300
# --- Live stream snapshot ---
# Periodically write the latest annotated frame as JPEG for an external web server
LIVE_STREAM_ENABLED=true
# Path to the shared-memory snapshot file (served by nginx / lighttpd)
LIVE_STREAM_FRAME_PATH=/dev/shm/byetrack-counter/live_frame.jpg
# JPEG quality (1–100)
LIVE_STREAM_FRAME_PATH=/dev/shm/zenai-kpc-counter/live_frame.jpg
LIVE_STREAM_QUALITY=75
# Write the snapshot every N frames (lower = more frequent updates)
LIVE_STREAM_EVERY_N=2
# --- Dashboard (counter_dashboard.py) ---
# Flask secret key for session/cookie signing — change in production!
SECRET_KEY=change-me-in-production
# Bind address for the Flask web server
DASHBOARD_HOST=0.0.0.0
# Listen port for the dashboard web UI
DASHBOARD_PORT=5000
# Enable Flask debug mode (true/false) — auto-reloads on code changes; disable in production
DASHBOARD_PORT=7000
FLASK_DEBUG=false
# Fallback name for the active counting-day JSON state file used by the dashboard
CURRENT_COUNTER_PATH=/tmp/bytetrack_current_counter.json
SAVE_CROSS_SNAPSHOT=true
CROSS_SNAPSHOT_DIR=/opt/zenai-kpc-snaps/snapshots
CROSS_SNAPSHOT_CLEANUP_SEC=3600
CROSS_SNAPSHOT_MAX_AGE_DAYS=3
DEDUP_PX=-1
# Debug: set DEBUG_TRACKING=true to log per-frame tracking details to stdout
DEBUG_TRACKING=true
+4
View File
@@ -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
+293
View File
@@ -0,0 +1,293 @@
#!/usr/bin/env python3
"""
Standalone status webhook (stdlib only, HTTP + HTTPS).
Receives:
GET /status -> {"status": "IN"|"OUT"|"OFF"}
GET /?lokasi=IN|OUT&status=ON|OFF -> success | fail
State stored in JSON as IN | OUT | OFF.
Success (200) only when a valid transition applies; otherwise 400 and JSON unchanged:
- lokasi=IN, status=ON + current OFF -> IN
- lokasi=OUT, status=ON + current OFF -> OUT
- lokasi=IN, status=OFF + current IN -> OFF
- lokasi=OUT, status=OFF + current OUT -> OFF
Fixed ports: HTTP 8002, HTTPS 8443 (TLS cannot share a port with plain HTTP).
"""
from __future__ import annotations
import json
import os
import signal
import ssl
import subprocess
import threading
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from pathlib import Path
from urllib.parse import parse_qs, urlparse
# Fixed ports — always the same (do not override via env).
HTTP_PORT = 8002
HTTPS_PORT = 8443
BASE_DIR = Path(os.path.dirname(os.path.abspath(__file__)))
STATE_FILE = Path(
os.getenv("STATUS_WEBHOOK_STATE", str(BASE_DIR / "status_webhook_state.json"))
)
CERT_FILE = Path(os.getenv("STATUS_WEBHOOK_CERT", str(BASE_DIR / "status_webhook_cert.pem")))
KEY_FILE = Path(os.getenv("STATUS_WEBHOOK_KEY", str(BASE_DIR / "status_webhook_key.pem")))
ALLOWED_LOKASI = frozenset({"IN", "OUT"})
ALLOWED_STATUS = frozenset({"ON", "OFF"})
ALLOWED_STORED = frozenset({"IN", "OUT", "OFF"})
_lock = threading.Lock()
def normalize_lokasi(value: str | None) -> str | None:
if value is None:
return None
value = value.strip().upper()
return value if value in ALLOWED_LOKASI else None
def normalize_request_status(value: str | None) -> str | None:
if value is None:
return None
value = value.strip().upper()
return value if value in ALLOWED_STATUS else None
def normalize_stored(value: str | None) -> str | None:
if value is None:
return None
value = str(value).strip().upper()
return value if value in ALLOWED_STORED else None
def load_status() -> str | None:
try:
data = json.loads(STATE_FILE.read_text(encoding="utf-8"))
except FileNotFoundError:
return None
except (OSError, json.JSONDecodeError) as exc:
print(f"status_webhook: failed to read state from {STATE_FILE}: {exc}")
return None
return normalize_stored(data.get("status"))
def is_run_status(status: str | None = None) -> bool:
"""True when counter should run (IN or OUT)."""
if status is None:
status = load_status()
return normalize_stored(status) in ("IN", "OUT")
def save_status(status: str) -> None:
STATE_FILE.parent.mkdir(parents=True, exist_ok=True)
tmp = STATE_FILE.with_suffix(".tmp")
payload = {"status": status}
tmp.write_text(json.dumps(payload) + "\n", encoding="utf-8")
tmp.replace(STATE_FILE)
def apply_transition(
lokasi: str, status: str, current: str | None
) -> tuple[str | None, bool]:
"""Return (new_status, ok). ok True only for a valid transition that stores."""
# Treat missing/empty state as OFF so first ON can start a session.
effective = current if current is not None else "OFF"
if status == "ON":
if effective == "OFF":
return lokasi, True
return None, False
# status == "OFF"
if effective == lokasi:
return "OFF", True
return None, False
def ensure_tls_certs() -> None:
if CERT_FILE.is_file() and KEY_FILE.is_file():
return
CERT_FILE.parent.mkdir(parents=True, exist_ok=True)
print(
f"status_webhook: generating self-signed TLS cert at {CERT_FILE} "
f"(set STATUS_WEBHOOK_CERT / STATUS_WEBHOOK_KEY to use your own)"
)
subprocess.run(
[
"openssl",
"req",
"-x509",
"-newkey",
"rsa:2048",
"-keyout",
str(KEY_FILE),
"-out",
str(CERT_FILE),
"-days",
"365",
"-nodes",
"-subj",
"/CN=status-webhook",
],
check=True,
capture_output=True,
text=True,
)
class StatusHandler(BaseHTTPRequestHandler):
def log_message(self, fmt: str, *args) -> None:
print(f"status_webhook: {self.address_string()} - {fmt % args}")
def do_GET(self) -> None:
parsed = urlparse(self.path)
path = parsed.path.rstrip("/") or "/"
if path == "/status":
current = load_status() or "OFF"
self._send_json(200, {"status": current})
return
if path not in ("/", ""):
self._send_text(404, "Not Found")
return
params = parse_qs(parsed.query)
lokasi_values = params.get("lokasi", [])
status_values = params.get("status", [])
if (
not lokasi_values
or lokasi_values[0] == ""
or not status_values
or status_values[0] == ""
):
self._send_text(400, "fail")
return
lokasi = normalize_lokasi(lokasi_values[0])
status = normalize_request_status(status_values[0])
if lokasi is None or status is None:
self._send_text(400, "fail")
return
with _lock:
current = load_status()
new_status, ok = apply_transition(lokasi, status, current)
if not ok or new_status is None:
self._send_text(400, "fail")
return
save_status(new_status)
self._send_text(200, "success")
def _send_text(self, code: int, body: str) -> None:
raw = body.encode("utf-8")
self.send_response(code)
self.send_header("Content-Type", "text/plain; charset=utf-8")
self.send_header("Content-Length", str(len(raw)))
self.end_headers()
self.wfile.write(raw)
def _send_json(self, code: int, payload: dict) -> None:
raw = json.dumps(payload).encode("utf-8")
self.send_response(code)
self.send_header("Content-Type", "application/json; charset=utf-8")
self.send_header("Content-Length", str(len(raw)))
self.end_headers()
self.wfile.write(raw)
def make_http_server(port: int = HTTP_PORT) -> ThreadingHTTPServer:
return ThreadingHTTPServer(("0.0.0.0", port), StatusHandler)
def make_https_server(port: int = HTTPS_PORT) -> ThreadingHTTPServer:
ensure_tls_certs()
server = ThreadingHTTPServer(("0.0.0.0", port), StatusHandler)
ctx = ssl.SSLContext(ssl.PROTOCOL_TLS_SERVER)
ctx.load_cert_chain(certfile=str(CERT_FILE), keyfile=str(KEY_FILE))
server.socket = ctx.wrap_socket(server.socket, server_side=True)
return server
def _try_make_server(kind: str, factory) -> ThreadingHTTPServer | None:
try:
return factory()
except OSError as exc:
print(f"status_webhook: {kind} port already in use ({exc}); skipping bind")
return None
def start_background() -> tuple[ThreadingHTTPServer | None, ThreadingHTTPServer | None]:
"""Start HTTP + HTTPS servers on daemon threads. Safe if ports are taken."""
http_server = _try_make_server("HTTP", make_http_server)
https_server = _try_make_server("HTTPS", make_https_server)
if http_server is not None:
threading.Thread(
target=http_server.serve_forever,
name="status-webhook-http",
daemon=True,
).start()
print(
f"status_webhook HTTP http://0.0.0.0:{HTTP_PORT}/"
f"?lokasi=IN|OUT&status=ON|OFF"
)
if https_server is not None:
threading.Thread(
target=https_server.serve_forever,
name="status-webhook-https",
daemon=True,
).start()
print(
f"status_webhook HTTPS https://0.0.0.0:{HTTPS_PORT}/"
f"?lokasi=IN|OUT&status=ON|OFF"
)
print(f"status_webhook state file: {STATE_FILE}")
return http_server, https_server
def stop_servers(
http_server: ThreadingHTTPServer | None,
https_server: ThreadingHTTPServer | None,
) -> None:
for server in (http_server, https_server):
if server is None:
continue
try:
server.shutdown()
server.server_close()
except Exception as exc:
print(f"status_webhook: shutdown error: {exc}")
def main() -> None:
http_server, https_server = start_background()
if http_server is None and https_server is None:
raise SystemExit("status_webhook: neither HTTP nor HTTPS could bind")
stop = threading.Event()
def _handle_stop(_signum, _frame) -> None:
stop.set()
signal.signal(signal.SIGTERM, _handle_stop)
signal.signal(signal.SIGINT, _handle_stop)
try:
# serve_forever already running on threads; park main until stop.
while not stop.is_set():
stop.wait(3600)
finally:
print("status_webhook stopped")
stop_servers(http_server, https_server)
if __name__ == "__main__":
main()
+19
View File
@@ -0,0 +1,19 @@
-----BEGIN CERTIFICATE-----
MIIDEzCCAfugAwIBAgIUFz+BTq4D+gxdrVY45TcZBoN8PnYwDQYJKoZIhvcNAQEL
BQAwGTEXMBUGA1UEAwwOc3RhdHVzLXdlYmhvb2swHhcNMjYwODA2MDQ0NjE4WhcN
MjcwODA2MDQ0NjE4WjAZMRcwFQYDVQQDDA5zdGF0dXMtd2ViaG9vazCCASIwDQYJ
KoZIhvcNAQEBBQADggEPADCCAQoCggEBALyKLPWaDeepc9n2edmHniHtsRkUGjl3
5SmkWcSbMfYUO3MeRif1HysnFd6VpZP9xP3GuxpMl6/Zpabcs5t1hKeQjIa9rl2z
OYr9ly4Lcwyl1UYwqy7yjNOtPrsNUNoWQkSuX7pHSdxDgh7Oa6KbKqYCsljjuT6z
ZBSptFINMJpP0n1/BEVYclk10nRqQ67YLngBCYBoSHtgvSSYlRDOUYmXja++Vlu4
tsJyguAb+YuZ+79Addsc+7uuO622IvUEcO55pzf4dBdenkxDEoPW7fppikh301Ck
k+wyCPv3jbDR/QOiH5K7YtdovFQuw+BAVRH/ddarBZwSaywhVCf7TXMCAwEAAaNT
MFEwHQYDVR0OBBYEFPZtrVQMH8TCxTNzRA9LtiBdFL//MB8GA1UdIwQYMBaAFPZt
rVQMH8TCxTNzRA9LtiBdFL//MA8GA1UdEwEB/wQFMAMBAf8wDQYJKoZIhvcNAQEL
BQADggEBABqAXwBSS3JcpBRGmIYwulOC2LUPi/qr58j9VvNFTLlRwUQkmXSXbcW/
n/cPjJxg0RG9xihCxUCClbXuC76ZKeTyfmp3kB0mWRJ7lRonf674/35/EG5+xXcM
2Vs6zkVrrcbebuGHzAv8z0GfcFbyN5iecHzpOD64CDRkxNb6hVf2u8Q7oC9TB9r6
PHrErjrT8pYs15LbZlB7925PtWZB5z8x695qYpO1c0y6IBsbVsXm9lVxww/BldK9
nJSIetKGSp1xoHgGsaiW1oNMkdpML0f9sk3CVzO4e3MG+dWJ9YCzTE+78YniN1Wr
U4dXZz9jT9eMxiav/KxrA3JVThmU8Qo=
-----END CERTIFICATE-----
+28
View File
@@ -0,0 +1,28 @@
-----BEGIN PRIVATE KEY-----
MIIEvQIBADANBgkqhkiG9w0BAQEFAASCBKcwggSjAgEAAoIBAQC8iiz1mg3nqXPZ
9nnZh54h7bEZFBo5d+UppFnEmzH2FDtzHkYn9R8rJxXelaWT/cT9xrsaTJev2aWm
3LObdYSnkIyGva5dszmK/ZcuC3MMpdVGMKsu8ozTrT67DVDaFkJErl+6R0ncQ4Ie
zmuimyqmArJY47k+s2QUqbRSDTCaT9J9fwRFWHJZNdJ0akOu2C54AQmAaEh7YL0k
mJUQzlGJl42vvlZbuLbCcoLgG/mLmfu/QHXbHPu7rjuttiL1BHDueac3+HQXXp5M
QxKD1u36aYpId9NQpJPsMgj7942w0f0Doh+Su2LXaLxULsPgQFUR/3XWqwWcEmss
IVQn+01zAgMBAAECggEAC96vIe9G/NTARHKuDTHqlLxAMBIB7KhNtydvt18F8DYp
3/+B7zYRdkgJqm/FcuHBKzD9ypQT4LBVK4ItlJX7egkxr7H1blTARK3efLmfzqYK
HVcnD9eZYiJAFsqp0nEgTu6jfDjMv59Ia+QXBq+6KaV10P7VRMtKe7qLbbcC3lQY
bzpwQr691t0NUy7xXbqjyx/akuEOAGZPYDkovdBu3eQrCzjGSosOKmwxZHUIUyDI
CYlIiKViR5qVvftioPgEHE/lM9e/2JWxQZB5vZlpjsWut5RNJihoZeWxKBNWxEqj
75zPiUa1KHAke4AiDKqxL37/1Dd4HWF/nz0PMPv6gQKBgQDtd4RNriSgruq2olh6
EAOlg70t1Qv8GOxrmOQTRm1nXa4N6Mb8eOpB7BH8twyRUwV5ftoIlG4GX2k/0rHT
3o/+C8BVUG8+t65K2fku8o2mVowjshWSgEKyEG/mZ4jGUlJYwi4rZM9opz/hhY8u
+0+yOOBmOuPzRXAL7eKQSjEWUwKBgQDLQR3c3IHWCGxaeNTTIj4KxKLOeGHT7Qt5
t9ktmKtTwMHbJWksQ0QNAWoXLr30bxyQgJjtd3TTi93nl/ALT956zp7rLFj0017O
1tGwEKgDgKuMO5wolnWUCMw+w071jdBpLC/hb6zJuLQHITrY+IAYOgZooWe/w9vD
lk6/BivIYQKBgQDZrQQXXPlwXccD8V9fTMy67T7+A1xAE+ysWPNBA/8HkKUbVPUK
vCAom6iFWppnoI3VKEXfNYiByPYmrhGaYFroCoec7OV8vU1EifjUYz0bbBx8ICOM
Lox0w4J/1wpWmWGowR8nYfqKOT3ikdaFv5L3kRGKRJNuDYm/NanIkGncxwKBgFAC
VAUK8DkWm8CJbA2onw+SFBx+mtPXrfq9+knOnTKc4DKp6Vq5J+KOufpiNfgwfOgN
FyXzLhPQLQvrbVymlgd1qm0cye+l/N4jBevuwpSOY/kRxgjcIXCiffP+4egbaPzd
ngN5+GR3xrY/yHB8ccAXp0osrzB3oty9IEZl4XpBAoGAMEXHcPfb7u7maXm4w6bR
BGGlSmxoUqueS252NSfuBIpeBpDdqJS8+euh2s7mAZOrBUlw+XTzC6HbyPoThLl/
dPyagyL48dwjZMn65s6T8Ia5uzDzSfjREKI2yFkAQqArF7JbBY/AwHkcko8o/zXc
F1TYtACipQQyFRK6FYmjK+M=
-----END PRIVATE KEY-----
+1
View File
@@ -0,0 +1 @@
{"status": "OFF"}
+420
View File
@@ -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)
+32
View File
@@ -0,0 +1,32 @@
[Unit]
Description=ZenAI KPC Status Webhook (HTTP + HTTPS)
Documentation=file:///opt/zenai-kpc-python/DEPLOY.md
After=network-online.target
Wants=network-online.target
[Service]
Type=simple
User=root
Group=root
WorkingDirectory=/opt/zenai-kpc-python
EnvironmentFile=/opt/zenai-kpc-python/.env
Environment=PYTHONNOUSERSITE=1
Environment=PATH=/opt/zenai-kpc-python/venv/bin:/usr/local/bin:/usr/bin:/bin
ExecStart=/opt/zenai-kpc-python/venv/bin/python status_webhook.py
TimeoutStopSec=15
KillSignal=SIGTERM
Restart=always
RestartSec=5
StartLimitInterval=60s
StartLimitBurst=3
NoNewPrivileges=true
ProtectHome=true
PrivateTmp=false
[Install]
WantedBy=multi-user.target