Merge branch 'dsutanto-main' to main

This commit is contained in:
Alberto-Audrix committed 2026-07-12 15:08:56 +07:00
commit d9c04d36c0
5 files changed
+478 -6

No files matched your search

+154 -3
View File
@@ -7,8 +7,11 @@ Uses ByteTrack (Kalman filter + two-stage IoU association) for tracking.
import numpy as np
import cv2
import csv
import json
import os
import signal
import socket
import threading
import time
from collections import deque
from datetime import datetime
@@ -39,6 +42,7 @@ LINE_Y2_FRAC = float(os.getenv("LINE_Y2_FRAC", "0.66"))
IMGSZ = int(os.getenv("IMGSZ", "320"))
HALF = os.getenv("HALF", "false").lower() == "true"
CONF = float(os.getenv("CONF", "0.3"))
NMS_IOU = float(os.getenv("NMS_IOU", "0.45"))
DEVICE = int(os.getenv("DEVICE", "0"))
# RKNN NPU core mask
@@ -120,6 +124,21 @@ MOTION_MIN_AREA_FRAC = float(os.getenv("MOTION_MIN_AREA_FRAC", "0.002"))
# slow/stationary object is never missed for long.
MOTION_HEARTBEAT_FRAMES = int(os.getenv("MOTION_HEARTBEAT_FRAMES", "15"))
# --- Runtime control (toggle counting on/off on the fly) ---
# When enabled, the process watches a small JSON control file and honors its
# "counting" flag. Set false to always count (ignore the control file).
CONTROL_ENABLED = os.getenv("CONTROL_ENABLED", "false").lower() == "true"
CONTROL_FILE = os.getenv("CONTROL_FILE", f"{OUTPUT_DIR}/control.json")
# Whether counting is active on startup when no control file exists yet.
CONTROL_DEFAULT_COUNTING = os.getenv("CONTROL_DEFAULT_COUNTING", "true").lower() == "true"
# Re-read the control file at most every N seconds.
CONTROL_POLL_SEC = float(os.getenv("CONTROL_POLL_SEC", "1.0"))
# Optional TCP control socket. When enabled, the counter listens for line-based
# commands so counting can be toggled over the network (in addition to the file).
CONTROL_SOCKET_ENABLED = os.getenv("CONTROL_SOCKET_ENABLED", "false").lower() == "true"
CONTROL_SOCKET_HOST = os.getenv("CONTROL_SOCKET_HOST", "127.0.0.1")
CONTROL_SOCKET_PORT = int(os.getenv("CONTROL_SOCKET_PORT", "5090"))
CROSS_FLASH_FRAMES = 12
POPUP_LIFETIME = 20
LINE_PULSE_FRAMES = 12
@@ -701,6 +720,93 @@ def now_str():
return datetime.now().strftime("%Y-%m-%d %H:%M:%S")
def read_counting_flag(default=True):
"""Read the 'counting' flag from the control file. Returns default on any error."""
try:
with open(CONTROL_FILE, "r", encoding="utf-8") as f:
data = json.load(f)
return bool(data.get("counting", default))
except FileNotFoundError:
return default
except Exception:
return default
def write_control_file(counting):
"""Create/update the control file atomically (used to seed defaults)."""
try:
Path(CONTROL_FILE).parent.mkdir(parents=True, exist_ok=True)
tmp = f"{CONTROL_FILE}.tmp"
with open(tmp, "w", encoding="utf-8") as f:
json.dump({"counting": bool(counting)}, f)
os.replace(tmp, CONTROL_FILE)
except Exception as exc:
print(f"[{now_str()}] Failed to write control file: {exc}")
def start_control_socket():
"""Start a TCP server for runtime control. Commands (newline-terminated):
START | RESUME | ON -> counting on
STOP | PAUSE | OFF -> counting off
TOGGLE -> flip
STATUS | GET -> report current state
It writes the shared control file, so the main loop's file-poll applies it.
Returns the server socket (call .close() to stop)."""
srv = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
srv.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
srv.bind((CONTROL_SOCKET_HOST, CONTROL_SOCKET_PORT))
srv.listen(5)
def handle(conn, addr):
with conn:
conn.settimeout(30)
try:
buf = b""
while not shutdown_requested:
try:
chunk = conn.recv(256)
except socket.timeout:
break
if not chunk:
break
buf += chunk
while b"\n" in buf:
line, buf = buf.split(b"\n", 1)
cmd = line.decode("utf-8", "ignore").strip().upper()
if not cmd:
continue
current = read_counting_flag(CONTROL_DEFAULT_COUNTING)
if cmd in ("START", "RESUME", "ON"):
write_control_file(True)
resp = "OK counting=on"
elif cmd in ("STOP", "PAUSE", "OFF"):
write_control_file(False)
resp = "OK counting=off"
elif cmd == "TOGGLE":
write_control_file(not current)
resp = f"OK counting={'off' if current else 'on'}"
elif cmd in ("STATUS", "GET"):
resp = f"OK counting={'on' if current else 'off'}"
else:
resp = "ERR unknown command"
conn.sendall((resp + "\n").encode("utf-8"))
except Exception:
pass
def loop():
print(f"Control socket listening on {CONTROL_SOCKET_HOST}:{CONTROL_SOCKET_PORT}")
while not shutdown_requested:
try:
conn, addr = srv.accept()
except OSError:
break
t = threading.Thread(target=handle, args=(conn, addr), daemon=True)
t.start()
threading.Thread(target=loop, daemon=True).start()
return srv
def open_capture(source):
if source.lower().startswith(("rtsp://", "http://")):
os.environ["OPENCV_FFMPEG_CAPTURE_OPTIONS"] = RTSP_FFMPEG_OPTIONS
@@ -1067,6 +1173,7 @@ def run():
core_mask=CORE_MASK,
imgsz=IMGSZ,
conf=CONF,
iou=NMS_IOU,
num_classes=NUM_CLASSES,
score_sigmoid=SCORE_SIGMOID,
)
@@ -1105,6 +1212,23 @@ def run():
counter_out = 0
last_snapshot_cleanup = 0.0
counting_active = True
last_control_poll = 0.0
control_socket = None
if CONTROL_ENABLED:
if not Path(CONTROL_FILE).exists():
write_control_file(CONTROL_DEFAULT_COUNTING)
counting_active = read_counting_flag(CONTROL_DEFAULT_COUNTING)
print(
f"Runtime control enabled | file={CONTROL_FILE} | "
f"counting={'ON' if counting_active else 'OFF'}"
)
if CONTROL_SOCKET_ENABLED:
try:
control_socket = start_control_socket()
except Exception as exc:
print(f"[{now_str()}] Failed to start control socket: {exc}")
cap, w, h, fps = connect_stream(SOURCE)
if cap is None:
store.shutdown()
@@ -1154,8 +1278,17 @@ def run():
cross_events_frame = []
detect_events_frame = []
if CONTROL_ENABLED and (now - last_control_poll) >= CONTROL_POLL_SEC:
last_control_poll = now
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")
skip_inference = False
if MOTION_DETECTION_ENABLED:
if not counting_active:
skip_inference = True
elif MOTION_DETECTION_ENABLED:
gray = cv2.cvtColor(frame, cv2.COLOR_BGR2GRAY)
if prev_gray is not None:
diff = cv2.absdiff(gray, prev_gray)
@@ -1377,6 +1510,14 @@ def run():
)
popups = draw_popups(frame, popups, frame_idx)
if CONTROL_ENABLED and not counting_active:
badge = "COUNTING PAUSED"
(bw, bh), _ = cv2.getTextSize(badge, cv2.FONT_HERSHEY_SIMPLEX, 0.6, 2)
bx = w // 2 - bw // 2
overlay_rect(frame, bx - 14, 48, bx + bw + 14, 48 + bh + 18, C_PANEL, alpha=0.75)
cv2.rectangle(frame, (bx - 14, 48), (bx + bw + 14, 48 + bh + 18), C_ACCENT, 2)
cv2.putText(frame, badge, (bx, 48 + bh + 6), cv2.FONT_HERSHEY_SIMPLEX, 0.6, C_ACCENT, 2, cv2.LINE_AA)
for flash_store in (
object_cross_flash1, object_cross_flash2,
):
@@ -1393,10 +1534,15 @@ def run():
if LIVE_STREAM_ENABLED and frame_idx % LIVE_STREAM_EVERY_N == 0:
try:
_, jpeg = cv2.imencode(
Path(LIVE_STREAM_FRAME_PATH).parent.mkdir(parents=True, exist_ok=True)
ok, jpeg = cv2.imencode(
".jpg", frame, [cv2.IMWRITE_JPEG_QUALITY, LIVE_STREAM_QUALITY]
)
_write_live_frame_atomic(LIVE_STREAM_FRAME_PATH, jpeg.tobytes())
if ok:
tmp_path = f"{LIVE_STREAM_FRAME_PATH}.tmp"
with open(tmp_path, "wb") as f:
f.write(jpeg.tobytes())
os.replace(tmp_path, LIVE_STREAM_FRAME_PATH)
except Exception:
pass
@@ -1443,6 +1589,11 @@ def run():
video_writer.release()
if cross_logger:
cross_logger.close()
if control_socket is not None:
try:
control_socket.close()
except Exception:
pass
model.release()
store.shutdown()