forked from dsutanto/zenai-kpc-python
Add Control mechanism
This commit is contained in:
1 parent
6afffb2184
commit
8f7ecb4321
5 files changed
+456
-2
No files matched your search
+145
-1
@@ -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
|
||||
@@ -121,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
|
||||
@@ -702,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
|
||||
@@ -1095,6 +1200,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()
|
||||
@@ -1144,8 +1266,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)
|
||||
@@ -1367,6 +1498,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,
|
||||
):
|
||||
@@ -1438,6 +1577,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()
|
||||
|
||||
|
||||
Reference in new issue
Block a user