ci / smoke (push) Canceled after 0s
Big docked truck bbox rests 1-2px past detection_polygon bottom edge, so 100%-containment made truck_in_area flicker False -> batch 15 split (6+276) on 2026-09-29 while the truck never left. v4-best.engine detected it at conf 0.92-0.97 in every replayed frame (no model miss, no retraining). - truck gate: frac_inside >= 0.5 instead of contains(bbox) - batch.truck_gone_tolerance_seconds (new, 30) authoritative; drop the hardcoded batch_mgr._truck_gone_tolerance = 30.0 override; batch.timeout_seconds documented as sack-idle pause only - [TRUCK] truck_in_area transition debug log
2061 lines
92 KiB
Python
2061 lines
92 KiB
Python
import os
|
|
os.environ["OPENCV_FFMPEG_CAPTURE_OPTIONS"] = "rtsp_transport;tcp|buffer_size;20480000|max_delay;500000|reorder_queue_size;500"
|
|
import argparse
|
|
import cv2
|
|
import numpy as np
|
|
import json
|
|
import threading
|
|
import time
|
|
from http.server import BaseHTTPRequestHandler, HTTPServer
|
|
from socketserver import ThreadingMixIn
|
|
from ultralytics import YOLO
|
|
from shapely.geometry import Point, Polygon, LineString, box
|
|
from collections import defaultdict, deque
|
|
import torch
|
|
import sqlite3
|
|
from datetime import datetime, timedelta
|
|
import logging
|
|
|
|
# Import repo rafan modules
|
|
from src.detection import SackDetector, TruckDetector, BoxDetector
|
|
from src.tracking import ByteTrackTracker
|
|
from src.stabilizer import BboxStabilizer
|
|
from src.truck_roi import TruckROITracker
|
|
from src.counting import LineCrossCounter, MultiClassLineCounter, drop_sacks_overlapping_boxes
|
|
from src.batch import BatchLifecycleManager, BatchRecord
|
|
from src.logger import CSVLogger
|
|
from src.dashboard import DashboardOverlay
|
|
from src.config_loader import (
|
|
Config,
|
|
check_legacy_batch_mode,
|
|
load_config,
|
|
read_zone_polygons,
|
|
resolve_active_mode,
|
|
)
|
|
|
|
# --- SQLite Database & State Configuration ---
|
|
if os.name == 'nt':
|
|
OUTPUT_DIR = 'd:/Belajar/menghitung karung'
|
|
DB_PATH = f'{OUTPUT_DIR}/jetson_counter.db'
|
|
STATE_FILE = f'{OUTPUT_DIR}/current_batch.json'
|
|
BATCH_MODE_FILE = f'{OUTPUT_DIR}/batch_mode.json'
|
|
LIVE_STREAM_FRAME_PATH = f'{OUTPUT_DIR}/live_frame.jpg'
|
|
else:
|
|
OUTPUT_DIR = os.getenv('OUTPUT_DIR', '/opt/jetson-counter')
|
|
DB_PATH = os.getenv('DB_PATH', f'{OUTPUT_DIR}/jetson_counter.db')
|
|
STATE_FILE = os.getenv('STATE_FILE', f'{OUTPUT_DIR}/current_batch.json')
|
|
BATCH_MODE_FILE = os.getenv('BATCH_MODE_FILE', f'{OUTPUT_DIR}/batch_mode.json')
|
|
LIVE_STREAM_FRAME_PATH = os.getenv('LIVE_STREAM_FRAME_PATH', '/dev/shm/jetson-counter/live_frame.jpg')
|
|
|
|
CAMERA_NAME = os.getenv('CAMERA_NAME', 'CC1')
|
|
OBJECT_LABEL = os.getenv('OBJECT_LABEL', 'karung-pakan')
|
|
DAILY_CUTOFF_TIME = os.getenv('DAILY_CUTOFF_TIME', '20:00')
|
|
BATCH_MERGE_THRESHOLD_SECONDS = int(os.getenv('BATCH_MERGE_THRESHOLD_SECONDS', '300'))
|
|
|
|
active_batch_info = None
|
|
batch_logger = None # CSVLogger for batch_summary.csv / sack_events.csv (set in run_prediction)
|
|
|
|
# --- CLI runtime flags (set from parse_args in __main__; defaults = production) ---
|
|
NO_DASHBOARD = False # --no-dashboard: skip live frame + live_status.json writes
|
|
NO_DB = False # --no-db: skip SQLite/state-file persistence
|
|
|
|
|
|
def parse_args():
|
|
"""CLI controls — zero flags reproduces systemd production behaviour."""
|
|
p = argparse.ArgumentParser(
|
|
description="Karung Counter — production pipeline (systemd service + dev CLI)"
|
|
)
|
|
p.add_argument("--source", type=str, default=None,
|
|
help="Video file path (overrides RTSP_URL)")
|
|
p.add_argument("--env", type=str, default=".env",
|
|
help="Path to .env file")
|
|
p.add_argument("--model", type=str, default=None,
|
|
help="Override MODEL_PATH (.pt/.engine)")
|
|
p.add_argument("--output-dir", type=str, default=None,
|
|
help="Override output dir for DB/state/live-frame (default: .env values)")
|
|
p.add_argument("--output-json", type=str, default=None,
|
|
help="Override batch report JSON path")
|
|
p.add_argument("--sack-conf", type=float, default=None,
|
|
help="Override sack detection confidence (default 0.35)")
|
|
p.add_argument("--truck-conf", type=float, default=None,
|
|
help="Override truck detection confidence (default 0.35)")
|
|
p.add_argument("--box-conf", type=float, default=None,
|
|
help="Override box detection confidence (default 0.35)")
|
|
p.add_argument("--box-model", type=str, default=None,
|
|
help="Override yolo11n sack+box model path (.pt/.engine)")
|
|
p.add_argument("--model-mode", type=str, default=None,
|
|
help="Model pipeline mode (default: config.yaml models.active_mode). "
|
|
"A=combined only; B=v4 truck + yolo11n sack+box; "
|
|
"C=A + yolo11n box-only; D=v4 truck + best sack + yolo11n box. "
|
|
"New modes from config.yaml need no code change.")
|
|
p.add_argument("--config", type=str, default=None,
|
|
help="Path to config.yaml (default: config.yaml next to predict.py)")
|
|
p.add_argument("--batch-timeout", type=float, default=None,
|
|
help="Override sack-idle + truck-gone timeouts (seconds)")
|
|
p.add_argument("--max-frames", type=int, default=None,
|
|
help="Stop after N frames (dev)")
|
|
p.add_argument("--no-dashboard", action="store_true",
|
|
help="Disable live frame publish + live_status.json")
|
|
p.add_argument("--no-db", action="store_true",
|
|
help="Disable SQLite/state-file persistence")
|
|
return p.parse_args()
|
|
|
|
|
|
def get_counting_date(dt=None):
|
|
if dt is None:
|
|
dt = datetime.now()
|
|
try:
|
|
cutoff = datetime.strptime(DAILY_CUTOFF_TIME, "%H:%M").time()
|
|
except Exception:
|
|
cutoff = datetime.strptime("20:00", "%H:%M").time()
|
|
|
|
if cutoff.hour == 0 and cutoff.minute == 0:
|
|
return dt.date().isoformat()
|
|
|
|
if dt.time() < cutoff:
|
|
return (dt.date() - timedelta(days=1)).isoformat()
|
|
return dt.date().isoformat()
|
|
|
|
def init_db():
|
|
try:
|
|
os.makedirs(os.path.dirname(DB_PATH), exist_ok=True)
|
|
os.makedirs(os.path.dirname(STATE_FILE), exist_ok=True)
|
|
os.makedirs(os.path.dirname(LIVE_STREAM_FRAME_PATH), exist_ok=True)
|
|
|
|
conn = sqlite3.connect(DB_PATH)
|
|
cur = conn.cursor()
|
|
cur.execute("""
|
|
CREATE TABLE IF NOT EXISTS batches (
|
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
counting_date TEXT NOT NULL,
|
|
batch_number INTEGER NOT NULL,
|
|
camera_name TEXT NOT NULL,
|
|
object_label TEXT NOT NULL,
|
|
count INTEGER NOT NULL,
|
|
start_time TEXT NOT NULL,
|
|
end_time TEXT NOT NULL,
|
|
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
|
|
UNIQUE(counting_date, batch_number, camera_name, object_label)
|
|
)
|
|
""")
|
|
cur.execute("""
|
|
CREATE TABLE IF NOT EXISTS daily_summaries (
|
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
counting_date TEXT NOT NULL,
|
|
camera_name TEXT NOT NULL,
|
|
object_label TEXT NOT NULL,
|
|
total_count INTEGER NOT NULL DEFAULT 0,
|
|
total_batches INTEGER NOT NULL DEFAULT 0,
|
|
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
|
|
UNIQUE(counting_date, camera_name, object_label)
|
|
)
|
|
""")
|
|
# Additive migration for box counting + model mode (idempotent).
|
|
cur.execute("PRAGMA table_info(batches)")
|
|
_cols = {r[1] for r in cur.fetchall()}
|
|
for _col, _typ in (
|
|
("box_loading", "INTEGER NOT NULL DEFAULT 0"),
|
|
("box_unloading", "INTEGER NOT NULL DEFAULT 0"),
|
|
("model_mode", "TEXT NOT NULL DEFAULT 'A'"),
|
|
):
|
|
if _col not in _cols:
|
|
cur.execute(f"ALTER TABLE batches ADD COLUMN {_col} {_typ}")
|
|
conn.commit()
|
|
conn.close()
|
|
print(f"[DB Info] Inisialisasi SQLite database berhasil: {DB_PATH}")
|
|
except Exception as e:
|
|
print(f"[DB Error] Gagal inisialisasi SQLite database: {e}")
|
|
|
|
def get_next_batch_number(counting_date):
|
|
try:
|
|
conn = sqlite3.connect(DB_PATH)
|
|
cur = conn.cursor()
|
|
cur.execute("""
|
|
SELECT COALESCE(MAX(batch_number), 0)
|
|
FROM batches
|
|
WHERE counting_date = ? AND camera_name = ? AND object_label = ?
|
|
""", (counting_date, CAMERA_NAME, OBJECT_LABEL))
|
|
row = cur.fetchone()
|
|
conn.close()
|
|
return row[0] + 1
|
|
except Exception as e:
|
|
print(f"[DB Error] Gagal mendapatkan batch_number: {e}")
|
|
return 1
|
|
|
|
def get_last_batch_info(counting_date):
|
|
try:
|
|
conn = sqlite3.connect(DB_PATH)
|
|
cur = conn.cursor()
|
|
cur.execute("""
|
|
SELECT batch_number, count, start_time, end_time
|
|
FROM batches
|
|
WHERE counting_date = ? AND camera_name = ? AND object_label = ?
|
|
ORDER BY batch_number DESC LIMIT 1
|
|
""", (counting_date, CAMERA_NAME, OBJECT_LABEL))
|
|
row = cur.fetchone()
|
|
conn.close()
|
|
if row:
|
|
return {
|
|
"batch_number": row[0],
|
|
"count": row[1],
|
|
"start_time": row[2],
|
|
"end_time": row[3]
|
|
}
|
|
except Exception as e:
|
|
print(f"[DB Error] Gagal mendapatkan batch terakhir: {e}")
|
|
return None
|
|
|
|
def save_active_batch_state():
|
|
global active_batch_info
|
|
if NO_DB:
|
|
return
|
|
if active_batch_info is None:
|
|
try:
|
|
if os.path.exists(STATE_FILE):
|
|
os.remove(STATE_FILE)
|
|
except Exception:
|
|
pass
|
|
return
|
|
try:
|
|
with open(STATE_FILE, 'w', encoding='utf-8') as f:
|
|
json.dump(active_batch_info, f, indent=2, ensure_ascii=False)
|
|
except Exception as e:
|
|
print(f"[DB Error] Gagal menulis {STATE_FILE}: {e}")
|
|
|
|
def log_crossing_event(ev, counter, batch_id, timestamp):
|
|
"""Print one crossing event (direction-aware wording) + append CSV event row.
|
|
|
|
Shared by the manual and auto event loops. CSV failure only warns —
|
|
the counting pipeline must never crash because of logging.
|
|
"""
|
|
label = "KARUNG" if ev.get("class_name", "sack") == "sack" else "BOX"
|
|
is_unloading = ev.get("direction") == "unloading"
|
|
if is_unloading:
|
|
total = counter.unloading_count if label == "KARUNG" else counter.box_unloading_count
|
|
print(f"[{label}] {label.capitalize()} #{ev['track_id']} keluar.")
|
|
print(f"[TOTAL] Total {label.lower()} keluar: {total}.")
|
|
else:
|
|
total = counter.loading_count if label == "KARUNG" else counter.box_loading_count
|
|
print(f"[{label}] {label.capitalize()} #{ev['track_id']} masuk.")
|
|
print(f"[TOTAL] Total {label.lower()} saat ini: {total}.")
|
|
if batch_logger:
|
|
try:
|
|
batch_logger.log_event(
|
|
int(batch_id or 0), ev["track_id"],
|
|
ev.get("direction", "loading"), timestamp,
|
|
ev.get("class_name", "sack"))
|
|
except Exception as e:
|
|
print(f"[WARN] CSV log_event gagal: {e}")
|
|
|
|
def log_manual_batch_stop():
|
|
"""CSV row for a manual batch stopped via dashboard button.
|
|
|
|
predict.py does not write manual batches to SQLite (counter_dashboard.py
|
|
owns that insert), so this is the only place the CSV row is produced.
|
|
"""
|
|
info = active_batch_info
|
|
if info is None or batch_logger is None:
|
|
return
|
|
final_count = int(info.get("count", 0) or 0)
|
|
box_final = int(info.get("box_count", 0) or 0)
|
|
if final_count == 0 and box_final == 0:
|
|
return # same discard rule as finalize_batch / dashboard stop
|
|
try:
|
|
batch_logger.log_batch(BatchRecord(
|
|
batch_id=int(info.get("batch_number", 0) or 0),
|
|
start_time=datetime.fromisoformat(
|
|
info.get("start_time") or datetime.now().isoformat()).timestamp(),
|
|
end_time=time.time(),
|
|
loading_count=final_count,
|
|
unloading_count=int(info.get("unloading", 0) or 0),
|
|
box_loading_count=box_final,
|
|
box_unloading_count=int(info.get("box_unloading", 0) or 0),
|
|
))
|
|
print(f"[CSV Info] Manual batch #{info.get('batch_number')} disimpan ke batch_summary.csv.")
|
|
except Exception as e:
|
|
print(f"[WARN] CSV log_batch manual gagal: {e}")
|
|
|
|
def finalize_batch(final_count, start_time_iso, end_time_iso,
|
|
box_final_count=0, box_unloading_count=0, model_mode="?",
|
|
unloading_count=0):
|
|
global active_batch_info
|
|
if active_batch_info is None:
|
|
return
|
|
net_sack = int(final_count) - int(unloading_count or 0)
|
|
net_box = int(box_final_count) - int(box_unloading_count or 0)
|
|
if NO_DB:
|
|
print(f"[DB Info] --no-db: batch #{active_batch_info.get('batch_number', 0)} "
|
|
f"({final_count} karung, {box_final_count} box) tidak disimpan.")
|
|
active_batch_info = None
|
|
save_active_batch_state()
|
|
return
|
|
if int(final_count) == 0 and int(box_final_count) == 0:
|
|
print(f"[BATCH] Batch #{active_batch_info.get('batch_number', 0)} bernilai 0 diabaikan (tidak disimpan ke database).")
|
|
active_batch_info = None
|
|
save_active_batch_state()
|
|
return
|
|
counting_date = active_batch_info["counting_date"]
|
|
batch_num = active_batch_info["batch_number"]
|
|
try:
|
|
conn = sqlite3.connect(DB_PATH)
|
|
cur = conn.cursor()
|
|
|
|
# 1. Insert completed batch
|
|
cur.execute("""
|
|
INSERT OR REPLACE INTO batches
|
|
(counting_date, batch_number, camera_name, object_label, count, start_time, end_time)
|
|
VALUES (?, ?, ?, ?, ?, ?, ?)
|
|
""", (counting_date, batch_num, CAMERA_NAME, OBJECT_LABEL, final_count, start_time_iso, end_time_iso))
|
|
|
|
# 2. Update daily summaries
|
|
cur.execute("""
|
|
SELECT SUM(count), COUNT(id)
|
|
FROM batches
|
|
WHERE counting_date = ? AND camera_name = ? AND object_label = ?
|
|
""", (counting_date, CAMERA_NAME, OBJECT_LABEL))
|
|
sum_row = cur.fetchone()
|
|
tot_count = sum_row[0] if sum_row[0] is not None else 0
|
|
tot_batches = sum_row[1] if sum_row[1] is not None else 0
|
|
|
|
cur.execute("""
|
|
INSERT OR REPLACE INTO daily_summaries
|
|
(counting_date, camera_name, object_label, total_count, total_batches, updated_at)
|
|
VALUES (?, ?, ?, ?, ?, CURRENT_TIMESTAMP)
|
|
""", (counting_date, CAMERA_NAME, OBJECT_LABEL, tot_count, tot_batches))
|
|
conn.commit()
|
|
|
|
# Box/DO columns may not exist on pre-migration DBs — tolerant update.
|
|
try:
|
|
cur2 = conn.cursor()
|
|
plate = active_batch_info.get("plate", "") or ""
|
|
do_numbers = json.dumps(active_batch_info.get("do_numbers", []) or [], ensure_ascii=False)
|
|
expected_sack = int(active_batch_info.get("expected_sack", 0) or 0)
|
|
expected_box = int(active_batch_info.get("expected_box", 0) or 0)
|
|
cur2.execute("""
|
|
UPDATE batches SET box_loading = ?, box_unloading = ?, model_mode = ?,
|
|
plate = ?, do_numbers = ?, expected_sack = ?, expected_box = ?,
|
|
net_sack = ?, net_box = ?
|
|
WHERE counting_date = ? AND batch_number = ? AND camera_name = ? AND object_label = ?
|
|
""", (box_final_count, box_unloading_count, model_mode,
|
|
plate, do_numbers, expected_sack, expected_box,
|
|
net_sack, net_box,
|
|
counting_date, batch_num, CAMERA_NAME, OBJECT_LABEL))
|
|
conn.commit()
|
|
except Exception:
|
|
pass
|
|
# Attach DO rows when batch_id known (best-effort).
|
|
try:
|
|
cur.execute(
|
|
"SELECT id FROM batches WHERE counting_date=? AND batch_number=? "
|
|
"AND camera_name=? AND object_label=?",
|
|
(counting_date, batch_num, CAMERA_NAME, OBJECT_LABEL))
|
|
row = cur.fetchone()
|
|
if row:
|
|
for do_id in (active_batch_info.get("do_ids") or []):
|
|
cur.execute(
|
|
"UPDATE delivery_orders SET batch_id=?, status='attached', updated_at=CURRENT_TIMESTAMP "
|
|
"WHERE id=?",
|
|
(row[0], do_id))
|
|
conn.commit()
|
|
except Exception:
|
|
pass
|
|
conn.close()
|
|
print(f"[DB Info] Sesi batch #{batch_num} disimpan ke database SQLite: "
|
|
f"{final_count} karung, {box_final_count} box (mode {model_mode}, "
|
|
f"net karung={net_sack}, net box={net_box}).")
|
|
except Exception as e:
|
|
print(f"[DB Error] Gagal menyimpan batch ke database: {e}")
|
|
# CSV row for every kept batch (discard/NO_DB early-returns above), even if
|
|
# the SQLite write failed — CSV is the fallback record.
|
|
if batch_logger:
|
|
try:
|
|
batch_logger.log_batch(BatchRecord(
|
|
batch_id=batch_num,
|
|
start_time=datetime.fromisoformat(start_time_iso).timestamp(),
|
|
end_time=datetime.fromisoformat(end_time_iso).timestamp(),
|
|
loading_count=int(final_count),
|
|
unloading_count=int(unloading_count or 0),
|
|
box_loading_count=int(box_final_count or 0),
|
|
box_unloading_count=int(box_unloading_count or 0),
|
|
))
|
|
print(f"[CSV Info] Batch #{batch_num} disimpan ke batch_summary.csv.")
|
|
except Exception as e:
|
|
print(f"[WARN] CSV log_batch gagal: {e}")
|
|
active_batch_info = None
|
|
save_active_batch_state()
|
|
|
|
def write_live_frame(frame):
|
|
if NO_DASHBOARD:
|
|
return
|
|
try:
|
|
os.makedirs(os.path.dirname(LIVE_STREAM_FRAME_PATH), exist_ok=True)
|
|
tmp_path = LIVE_STREAM_FRAME_PATH.replace(".jpg", ".tmp.jpg")
|
|
cv2.imwrite(tmp_path, frame, [cv2.IMWRITE_JPEG_QUALITY, 80])
|
|
os.replace(tmp_path, LIVE_STREAM_FRAME_PATH)
|
|
except PermissionError:
|
|
# Transient file lock on Windows when Flask dashboard reads it, safe to ignore
|
|
pass
|
|
except Exception as e:
|
|
print(f"[ERROR] Gagal menulis live frame: {e}")
|
|
|
|
def draw_annotations_on_frame(frame, bbox_list):
|
|
try:
|
|
h_f, w_f = frame.shape[:2]
|
|
|
|
# 1. Draw left and right limits (vertical lines)
|
|
cv2.line(frame, (left_limit, 0), (left_limit, h_f), (255, 0, 255), 2)
|
|
cv2.line(frame, (right_limit, 0), (right_limit, h_f), (255, 0, 255), 2)
|
|
|
|
# 2. Draw ZONA_PALET (Cyan)
|
|
if len(ZONA_PALET) >= 3:
|
|
cv2.polylines(frame, [ZONA_PALET], True, (255, 255, 0), 2)
|
|
|
|
# 3. Draw ZONA_TRUCK (Yellow/Red based on system_state)
|
|
if poly_truck is not None and not poly_truck.is_empty:
|
|
pts = np.array(poly_truck.exterior.coords, dtype=np.int32)
|
|
color_truck = (0, 204, 255) if system_state == STATE_WAITING_FOR_TRUCK else (0, 255, 0)
|
|
cv2.polylines(frame, [pts], True, color_truck, 2)
|
|
cv2.putText(frame, "ZONA TRUK BATCH", (pts[0][0], max(20, pts[0][1] - 8)), cv2.FONT_HERSHEY_SIMPLEX, 0.6, color_truck, 2)
|
|
elif len(ZONA_TRUCK) >= 3:
|
|
color_truck = (0, 204, 255) if system_state == STATE_WAITING_FOR_TRUCK else (0, 255, 0)
|
|
cv2.polylines(frame, [ZONA_TRUCK], True, color_truck, 2)
|
|
cv2.putText(frame, "ZONA TRUK BATCH", (ZONA_TRUCK[0][0], max(20, ZONA_TRUCK[0][1] - 8)), cv2.FONT_HERSHEY_SIMPLEX, 0.6, color_truck, 2)
|
|
|
|
# 4. Draw active bboxes and points from bbox_list
|
|
for item in bbox_list:
|
|
try:
|
|
if item.get("is_counted", False):
|
|
continue
|
|
color_str = item.get("color", "rgba(0, 255, 0, 1.0)")
|
|
if color_str.startswith("rgba"):
|
|
parts = color_str.replace("rgba(", "").replace(")", "").split(",")
|
|
r, g, b = int(parts[0]), int(parts[1]), int(parts[2])
|
|
bgr_color = (b, g, r)
|
|
else:
|
|
bgr_color = (0, 255, 0)
|
|
|
|
norm_box = item["bbox"]
|
|
x1 = int(norm_box[0] * w_f)
|
|
y1 = int(norm_box[1] * h_f)
|
|
x2 = int(norm_box[2] * w_f)
|
|
y2 = int(norm_box[3] * h_f)
|
|
|
|
cv2.rectangle(frame, (x1, y1), (x2, y2), bgr_color, 2)
|
|
|
|
cx, cy = int(item["centroid"][0] * w_f), int(item["centroid"][1] * h_f)
|
|
# Gambar Point di Tengah BBox
|
|
cv2.circle(frame, (cx, cy), 5, (0, 255, 255), -1)
|
|
|
|
label = f"{item['status']} #{item['id']}" if item["id"] != 9999 else item["status"]
|
|
cv2.putText(frame, label, (x1, max(15, y1 - 8)), cv2.FONT_HERSHEY_SIMPLEX, 0.5, bgr_color, 2)
|
|
|
|
# Jika ada entry_point (titik acuan awal & radius 50px)
|
|
if item.get("entry_point") is not None:
|
|
ex, ey = item["entry_point"]
|
|
is_counted = item.get("is_counted", False)
|
|
viz_color = (0, 255, 0) if is_counted else (0, 140, 255)
|
|
|
|
# 1. Gambar Titik Acuan Awal
|
|
cv2.circle(frame, (ex, ey), 4, viz_color, -1)
|
|
# 2. Gambar Lingkaran Radius 50px
|
|
cv2.circle(frame, (ex, ey), 50, viz_color, 2, lineType=cv2.LINE_AA)
|
|
# 3. Garis hubung ke centroid aktif
|
|
cv2.line(frame, (ex, ey), (cx, cy), viz_color, 1)
|
|
# 4. Label Jarak
|
|
dist_val = np.sqrt((cx - ex)**2 + (cy - ey)**2)
|
|
dist_label = "COUNTED (+1)" if is_counted else f"{dist_val:.0f}/50px"
|
|
cv2.putText(frame, dist_label, (ex - 20, max(15, ey - 10)), cv2.FONT_HERSHEY_SIMPLEX, 0.45, viz_color, 2)
|
|
except Exception:
|
|
pass
|
|
|
|
# 5. Draw active Duplicate Radius Circles on Port 8000 live stream (terkini 3.0 detik)
|
|
if 'counted_sack_positions' in globals() and counted_sack_positions:
|
|
rad_vis = DUPLICATE_CIRCLE_RADIUS if ('DUPLICATE_CIRCLE_RADIUS' in globals() and DUPLICATE_CIRCLE_RADIUS > 0) else 30
|
|
now_t = time.time()
|
|
active_circles = [p for p in counted_sack_positions if len(p) < 3 or (now_t - p[2]) <= 3.0]
|
|
for pos_item in active_circles:
|
|
px, py = pos_item[0], pos_item[1]
|
|
tid = pos_item[3] if len(pos_item) > 3 else 0
|
|
cv2.circle(frame, (int(px), int(py)), int(rad_vis), (0, 255, 255), 2, lineType=cv2.LINE_AA)
|
|
cv2.circle(frame, (int(px), int(py)), 4, (0, 255, 0), -1)
|
|
cv2.putText(frame, f"DEDUP #{tid}", (int(px) - 25, max(15, int(py) - int(rad_vis) - 5)), cv2.FONT_HERSHEY_SIMPLEX, 0.45, (0, 255, 255), 1)
|
|
|
|
cv2.putText(frame, f"RADIUS DEDUP: {rad_vis}px", (w_f - 270, 40), cv2.FONT_HERSHEY_SIMPLEX, 0.65, (0, 255, 255), 2)
|
|
|
|
# 6. Draw HUD stats on top left
|
|
total_in = metrics.get('total_masuk', 0)
|
|
total_out = metrics.get('total_keluar', 0)
|
|
net_cnt = total_in - total_out
|
|
cv2.putText(frame, f"STATUS: {system_state}", (20, 40), cv2.FONT_HERSHEY_SIMPLEX, 0.7, (255, 255, 255), 2)
|
|
cv2.putText(frame, f"IN: {total_in} OUT: {total_out} NET: {net_cnt}", (20, 75), cv2.FONT_HERSHEY_SIMPLEX, 0.7, (0, 255, 0), 2)
|
|
cv2.putText(frame, f"FPS: {current_fps:.2f}", (20, 110), cv2.FONT_HERSHEY_SIMPLEX, 0.7, (255, 204, 0), 2)
|
|
except Exception:
|
|
pass
|
|
|
|
# --- RUNNING LOCALLY (Colab patches removed) ---
|
|
# -----------------------------------------------
|
|
|
|
# =====================================================================
|
|
# 0. PARAMETER KALIBRASI & STATE MACHINE
|
|
# =====================================================================
|
|
MIN_VALID_AREA_REF = 15000
|
|
MIN_VALID_AREA = 15000
|
|
JARAK_ABSORBSI_GHOST = 50
|
|
|
|
# --- Logika masuk/keluar berbasis overlap + delay ---
|
|
ENTRY_OVERLAP_THRESHOLD = 0.20 # 20%
|
|
EXIT_OVERLAP_THRESHOLD = 0.05 # 5%
|
|
CONFIRM_DELAY_SEC = 0.5 # delay masuk
|
|
EXIT_CONFIRM_DELAY_SEC = 6.0 # delay keluar
|
|
COUNTED_DISPLAY_TIMEOUT_SEC = 0.5 # durasi tampil kotak hijau setelah terhitung
|
|
|
|
# --- Parameter Anti-Double Count (Spasial) ---
|
|
JARAK_TOLERANSI_DUPLIKAT_REF = 80
|
|
JARAK_TOLERANSI_DUPLIKAT = 80
|
|
TOLERANSI_FRAME_HILANG = 1200 # 1200 frame untuk Re-ID lost
|
|
MAX_REID_TRANSIT_DISTANCE_REF = 400
|
|
MAX_REID_TRANSIT_DISTANCE = 400 # Max pixel distance for Re-ID
|
|
|
|
# --- Parameter Bbox Muncul Tiba-tiba & Transit ---
|
|
MIN_DISPLACEMENT_START_IN_TRUCK = 15
|
|
MIN_LINEARITY_START_IN_TRUCK = 0.70
|
|
MIN_DISPLACEMENT_COUNTING_ZONE = 12
|
|
MIN_DY_COUNTING_ZONE = -3
|
|
|
|
# --- Parameter Lingkaran Duplikat Statis ---
|
|
DUPLICATE_CIRCLE_RADIUS_REF = 35
|
|
DUPLICATE_CIRCLE_RADIUS = 35
|
|
CROSS_CLASS_IOU = 0.4 # drop sack dets overlapping a box det (cross-model)
|
|
SHOW_ALL_BBOXES = False
|
|
counted_sack_positions = []
|
|
CIRCLE_STAY_TIMEOUT_SEC = 10.0
|
|
INFERENCE_STRIDE = 2
|
|
CAMERA_NOISE_DEADBAND = 50 # pixels deadband for static checks (diperbesar ke 50 sesuai permintaan user)
|
|
|
|
# --- Duo-Model Batching State Machine ---
|
|
STATE_WAITING_FOR_TRUCK = "WAITING_FOR_TRUCK"
|
|
STATE_COUNTING_SACKS = "COUNTING_SACKS"
|
|
STATE_TRUCK_FULL = "TRUCK_FULL"
|
|
STATE_TRUCK_LEAVING = "TRUCK_LEAVING"
|
|
|
|
system_state = STATE_WAITING_FOR_TRUCK
|
|
truck_static_frames = 0
|
|
truck_initial_bbox = None # [x1, y1, x2, y2]
|
|
truck_entry_point = None # (ex_t, ey_t)
|
|
entry_points = {} # track_id -> (ex, ey)
|
|
static_sack_visuals = {} # track_id -> HSV histogram representation
|
|
static_frames = defaultdict(int)
|
|
moving_frames = defaultdict(int)
|
|
sack_class_id = 0 # Default class ID for sack
|
|
counted_sacks = {}
|
|
lost_counted_sacks = {}
|
|
all_counted_sacks_map = {}
|
|
last_seen_near_person_frame = {}
|
|
blocked_due_to_duplicate = {}
|
|
|
|
# --- Model paths & config ---
|
|
# All weights live in models/ (see models/modelREADME.md for mode/filter matrix
|
|
# and config.yaml models.* for the canonical paths + per-mode presets).
|
|
_BASE_DIR = os.path.dirname(os.path.abspath(__file__))
|
|
_MODELS_DIR = os.path.join(_BASE_DIR, "models")
|
|
_DEFAULT_CONFIG_PATH = os.path.join(_BASE_DIR, "config.yaml")
|
|
|
|
# Active Config object (set in __main__/run_prediction before the pipeline starts).
|
|
_CFG: Config | None = None
|
|
|
|
|
|
def _legacy_combined_fallback() -> str:
|
|
"""Pre-YAML combined-model auto-pick (kept for --model/MODEL_PATH handling)."""
|
|
_env_model = os.getenv("MODEL_PATH")
|
|
if _env_model and os.path.exists(_env_model):
|
|
return _env_model
|
|
_candidates = [
|
|
os.path.join(_MODELS_DIR, "v4-best.engine"),
|
|
os.path.join(_MODELS_DIR, "v4-best.pt"),
|
|
os.path.join(_MODELS_DIR, "v4-best (1).pt"),
|
|
os.path.join(_MODELS_DIR, "model_karung_truk.engine"),
|
|
os.path.join(_MODELS_DIR, "model_karung_truk.pt"),
|
|
]
|
|
return next((p for p in _candidates if os.path.exists(p)),
|
|
os.path.join(_MODELS_DIR, "v4-best.pt"))
|
|
|
|
|
|
def apply_cli_overrides_to_config(cfg: Config, args) -> Config:
|
|
"""Fold legacy CLI/env model overrides into the Config object (in place).
|
|
|
|
--model/--box-model/MODEL_PATH keep their pre-YAML meaning:
|
|
a single v4 file used wherever a v4 engine is needed.
|
|
--sack-conf/--truck-conf/--box-conf override detection_params conf.
|
|
"""
|
|
v4_override = args.model or os.getenv("MODEL_PATH")
|
|
if v4_override:
|
|
if not os.path.exists(v4_override):
|
|
print(f"[WARN] --model/MODEL_PATH {v4_override} tidak ditemukan, "
|
|
f"dipakai langsung (YOLO bisa resolve).")
|
|
for key in ("combined", "truck_only"):
|
|
if key in cfg.models.paths:
|
|
cfg.models.paths[key] = v4_override
|
|
if getattr(args, "box_model", None):
|
|
cfg.models.paths["box"] = args.box_model
|
|
for cls_name, val in (("sack", args.sack_conf), ("truck", args.truck_conf),
|
|
("box", args.box_conf)):
|
|
if val is not None and cls_name in cfg.models.detection_params:
|
|
cfg.models.detection_params[cls_name].conf = float(val)
|
|
if getattr(args, "batch_timeout", None) is not None:
|
|
cfg.batch.timeout_seconds = float(args.batch_timeout)
|
|
return cfg
|
|
|
|
|
|
# --- Model pipeline modes -------------------------------------------------
|
|
# Modes are DATA in config.yaml models.modes (engines + class filters only).
|
|
# A=combined only; B=v4 truck + yolo11n sack+box; C=A + yolo11n box-only (default);
|
|
# D=v4 truck + best sack-only + yolo11n box-only. New modes need no code change.
|
|
# All engines verified to coexist (~16-24 MB peak of 7.6 GB).
|
|
|
|
|
|
def resolve_model_mode(explicit: str | None) -> str:
|
|
"""Precedence: --model-mode > MODEL_MODE env (deprecated) > config.yaml.
|
|
|
|
Dashboard switches persist to config.yaml models.active_mode and take effect
|
|
on next service restart (models are loaded once at startup).
|
|
batch_mode.json is legacy and no longer consulted.
|
|
"""
|
|
cfg = _CFG or load_config(_DEFAULT_CONFIG_PATH)
|
|
return resolve_active_mode(explicit, cfg)
|
|
|
|
|
|
def _bbox_area_ok(d, min_area: float) -> bool:
|
|
"""Permissive per-class guardrail from config.yaml detection_params."""
|
|
x1, y1, x2, y2 = d.bbox
|
|
return (x2 - x1) * (y2 - y1) >= min_area
|
|
|
|
|
|
def build_model_pipeline(mode, cfg, device, base_dir=None):
|
|
"""Instantiate YOLO handles + detectors/trackers for a mode preset.
|
|
|
|
Roles derive structurally from the preset's per-engine `classes`
|
|
(config.yaml models.modes) — no per-mode-letter branching, so new
|
|
modes (E, F, ...) work without code changes:
|
|
- tracker (primary): first engine covering `sack`
|
|
- truck_detector: first engine covering `truck`
|
|
(separate_truck_model = truck engine is not the primary)
|
|
- box_tracker: dedicated iff a non-primary engine covers `box`
|
|
(shared with the primary tracker otherwise, e.g. mode B)
|
|
|
|
Returns dict with keys: mode, truck_detector, tracker, box_tracker
|
|
(None when sack tracker covers boxes or no box engine), separate_truck_model,
|
|
class_filters, min_areas.
|
|
"""
|
|
base_dir = base_dir or _BASE_DIR
|
|
preset = cfg.models.modes[mode]
|
|
filters = {k: tuple(v) for k, v in preset.class_filters.items()}
|
|
dp_truck = cfg.detection_params_for("truck")
|
|
dp_sack = cfg.detection_params_for("sack")
|
|
dp_box = cfg.detection_params_for("box")
|
|
print(f"[INFO] Model mode: {mode} ({preset.description})")
|
|
dummy = np.zeros((720, 1280, 3), dtype=np.uint8)
|
|
|
|
handles: dict[str, object] = {}
|
|
|
|
def _load(path_key: str):
|
|
if path_key not in handles:
|
|
path = cfg.engine_path(path_key, base_dir)
|
|
print(f"[INFO] Memuat {path_key}: {path}")
|
|
m = YOLO(path)
|
|
_ = m(dummy, imgsz=640, device=device, verbose=False) # warm-up CUDA/TRT ctx
|
|
handles[path_key] = m
|
|
return handles[path_key]
|
|
|
|
primary_key = next((e.path for e in preset.engines if "sack" in e.classes), None)
|
|
truck_key = next((e.path for e in preset.engines if "truck" in e.classes), None)
|
|
box_keys = [e.path for e in preset.engines if "box" in e.classes]
|
|
|
|
if primary_key is None:
|
|
raise ValueError(
|
|
f"Mode {mode!r} has no engine covering class 'sack' — "
|
|
f"at least one engines[].classes must include sack."
|
|
)
|
|
|
|
tracker = ByteTrackTracker(_load(primary_key), dp_sack.conf, dp_sack.iou)
|
|
truck_detector = (
|
|
TruckDetector(_load(truck_key), dp_truck.conf,
|
|
class_filter=filters.get("truck") or ("truck",),
|
|
iou=dp_truck.iou)
|
|
if truck_key is not None else None
|
|
)
|
|
separate_truck_model = truck_key is not None and truck_key != primary_key
|
|
|
|
box_tracker = None
|
|
if box_keys and filters.get("box"):
|
|
dedicated = next((k for k in box_keys if k != primary_key), None)
|
|
if dedicated is not None:
|
|
box_tracker = ByteTrackTracker(_load(dedicated), dp_box.conf, dp_box.iou)
|
|
# else: boxes share the primary tracker (single ID space, e.g. mode B)
|
|
|
|
return {
|
|
"mode": mode,
|
|
"truck_detector": truck_detector,
|
|
"tracker": tracker,
|
|
"box_tracker": box_tracker,
|
|
"separate_truck_model": separate_truck_model,
|
|
"class_filters": filters,
|
|
"min_areas": {
|
|
"truck": dp_truck.min_bbox_area,
|
|
"sack": dp_sack.min_bbox_area,
|
|
"box": dp_box.min_bbox_area,
|
|
},
|
|
"confs": {
|
|
"truck": dp_truck.conf,
|
|
"sack": dp_sack.conf,
|
|
"box": dp_box.conf,
|
|
},
|
|
}
|
|
|
|
# =====================================================================
|
|
|
|
# =====================================================================
|
|
# 1. KONFIGURASI KOORDINAT ZONA
|
|
# =====================================================================
|
|
width = 1280
|
|
height = 720
|
|
scale_x = 1.0
|
|
scale_y = 1.0
|
|
|
|
def get_bottom_quarter(pts):
|
|
if len(pts) < 4:
|
|
return np.array([], dtype=np.int32)
|
|
pts_list = pts.tolist() if isinstance(pts, np.ndarray) else list(pts)
|
|
sorted_by_y = sorted(pts_list, key=lambda p: p[1])
|
|
tops = sorted_by_y[:2]
|
|
bottoms = sorted_by_y[2:]
|
|
|
|
tops_sorted = sorted(tops, key=lambda p: p[0])
|
|
tl = np.array(tops_sorted[0], dtype=np.float32)
|
|
tr = np.array(tops_sorted[1], dtype=np.float32)
|
|
|
|
bottoms_sorted = sorted(bottoms, key=lambda p: p[0])
|
|
bl = np.array(bottoms_sorted[0], dtype=np.float32)
|
|
br = np.array(bottoms_sorted[1], dtype=np.float32)
|
|
|
|
p_left = tl * 0.5 + bl * 0.5
|
|
p_right = tr * 0.5 + br * 0.5
|
|
|
|
return np.array([
|
|
[int(p_left[0]), int(p_left[1])],
|
|
[int(p_right[0]), int(p_right[1])],
|
|
[int(tr[0]), int(tr[1])],
|
|
[int(tl[0]), int(tl[1])]
|
|
], dtype=np.int32)
|
|
|
|
ZONES_JSON_PATH = "zones.json"
|
|
DEFAULT_PALET = []
|
|
DEFAULT_TRUCK = []
|
|
|
|
|
|
|
|
def load_zones():
|
|
"""Load zone geometry from zones.json (polygons only — knobs live in config.yaml).
|
|
|
|
Delegates to src.config_loader.read_zone_polygons: legacy knob keys in
|
|
zones.json are ignored with a warning (config.yaml counting.* wins).
|
|
ZONA_COUNTING_REF is derived from the truck polygon (bottom quarter),
|
|
not from the file's "counting" key — preserved pre-YAML behaviour.
|
|
"""
|
|
global ZONA_PALET_REF, ZONA_TRUCK_REF, GARIS_COUNTING_REF
|
|
global ZONA_COUNTING_REF, left_limit_ref, right_limit_ref, EXTERNAL_STREAM_URL_REF
|
|
if os.path.exists(ZONES_JSON_PATH):
|
|
try:
|
|
zones = read_zone_polygons(ZONES_JSON_PATH)
|
|
ZONA_PALET_REF = np.array(zones.get('palet', []), dtype=np.int32)
|
|
ZONA_TRUCK_REF = np.array(zones.get('truck', []), dtype=np.int32)
|
|
ZONA_COUNTING_REF = get_bottom_quarter(ZONA_TRUCK_REF)
|
|
left_limit_ref = float(zones.get('left_limit', 0.05))
|
|
right_limit_ref = float(zones.get('right_limit', 0.95))
|
|
GARIS_COUNTING_REF = ZONA_TRUCK_REF.copy()
|
|
EXTERNAL_STREAM_URL_REF = zones.get('external_stream_url', '') or 'http://192.168.192.96:8888/cam/'
|
|
print("[INFO] Berhasil memuat koordinat zona dari zones.json (via src.config_loader)")
|
|
return
|
|
except Exception as e:
|
|
print(f"[WARNING] Gagal memuat zones.json ({e}), menggunakan default.")
|
|
|
|
ZONA_PALET_REF = np.array(DEFAULT_PALET, dtype=np.int32)
|
|
ZONA_TRUCK_REF = np.array(DEFAULT_TRUCK, dtype=np.int32)
|
|
ZONA_COUNTING_REF = get_bottom_quarter(ZONA_TRUCK_REF)
|
|
left_limit_ref = 0.05
|
|
right_limit_ref = 0.95
|
|
GARIS_COUNTING_REF = ZONA_TRUCK_REF.copy()
|
|
|
|
load_zones()
|
|
|
|
DUPLICATE_CIRCLE_RADIUS = DUPLICATE_CIRCLE_RADIUS_REF
|
|
MIN_VALID_AREA = MIN_VALID_AREA_REF
|
|
JARAK_TOLERANSI_DUPLIKAT = JARAK_TOLERANSI_DUPLIKAT_REF
|
|
MAX_REID_TRANSIT_DISTANCE = MAX_REID_TRANSIT_DISTANCE_REF
|
|
|
|
ZONA_PALET = ZONA_PALET_REF.copy()
|
|
ZONA_TRUCK = ZONA_TRUCK_REF.copy()
|
|
ZONA_COUNTING = ZONA_COUNTING_REF.copy() if len(ZONA_COUNTING_REF) > 0 else np.array([], dtype=np.int32)
|
|
GARIS_COUNTING = GARIS_COUNTING_REF.copy()
|
|
|
|
poly_palet = Polygon(ZONA_PALET) if len(ZONA_PALET) >= 3 else None
|
|
poly_truck = Polygon(ZONA_TRUCK) if len(ZONA_TRUCK) >= 3 else None
|
|
poly_counting = Polygon(ZONA_COUNTING) if len(ZONA_COUNTING) >= 3 else None
|
|
|
|
left_limit = int(left_limit_ref * 1280)
|
|
right_limit = int(right_limit_ref * 1280)
|
|
|
|
from shapely.geometry import LineString
|
|
if len(ZONA_COUNTING) >= 4:
|
|
line_counting = LineString([ZONA_COUNTING[3], ZONA_COUNTING[2]])
|
|
else:
|
|
line_counting = poly_truck.boundary if poly_truck is not None else None
|
|
|
|
DEBOUNCE_FRAMES = 8 # tetap dipakai untuk label visual zona (PALET/AREA BEBAS), TIDAK untuk keputusan counting
|
|
|
|
# =====================================================================
|
|
# 2. STATE TRACKING
|
|
# =====================================================================
|
|
track_zone_history = defaultdict(lambda: deque(maxlen=DEBOUNCE_FRAMES))
|
|
track_confirmed_state = {} # dipakai untuk LABEL VISUAL saja (PALET/AREA BEBAS), bukan untuk counting
|
|
|
|
is_locked = defaultdict(bool)
|
|
already_counted = defaultdict(bool)
|
|
has_crossed_line = defaultdict(bool) # --- TAMBAHAN BARU ---
|
|
exit_crossed_line = defaultdict(bool) # --- TAMBAHAN BARU: LOGIKA KELUAR ---
|
|
track_areas = defaultdict(float) # --- TAMBAHAN BARU: LUAS BBOX ---
|
|
track_started_in_truck = defaultdict(bool)
|
|
outside_truck_frames = defaultdict(int)
|
|
counted_at_frame = {}
|
|
|
|
# --- FIX: pending timer terpisah untuk proses MASUK dan KELUAR, berbasis overlap, bukan jarak ---
|
|
pending_enter_since = defaultdict(lambda: None)
|
|
pending_exit_since = defaultdict(lambda: None)
|
|
|
|
track_positions = defaultdict(lambda: deque(maxlen=20))
|
|
lost_tracks = {}
|
|
prev_active_track_ids = set()
|
|
|
|
# --- STATE LINGKARAN DUPLIKAT STATIS ---
|
|
track_initial_truck_pos = {}
|
|
track_truck_entry_frame = {}
|
|
has_exited_circle = defaultdict(bool)
|
|
delay_completed = defaultdict(bool)
|
|
blocked_without_counting = defaultdict(bool)
|
|
track_is_valid_bag = defaultdict(bool)
|
|
|
|
metrics = {
|
|
"total_masuk": 0,
|
|
"total_keluar": 0,
|
|
"box_masuk": 0,
|
|
"box_keluar": 0
|
|
}
|
|
|
|
MAX_REID_DISTANCE = 120
|
|
MAX_REID_FRAMES = 200
|
|
|
|
# Warna
|
|
COKLAT = (19, 69, 139) # PENDING - baru masuk, menunggu konfirmasi 0.5s
|
|
HIJAU_TERVERIFIKASI = (100, 255, 100) # CONFIRMED - masuk sah
|
|
ORANYE_PENDING_KELUAR = (0, 165, 255) # PENDING - sedang menunggu konfirmasi keluar
|
|
BIRU_PALET = (255, 100, 100)
|
|
MERAH_BEBAS = (100, 100, 255)
|
|
ABU_FRAGMENT = (150, 150, 150)
|
|
|
|
|
|
# =====================================================================
|
|
# MULTI-THREADED REAL-TIME WEB DASHBOARD & STREAMING (ZERO DEPENDENCY)
|
|
# =====================================================================
|
|
import queue
|
|
|
|
EXTERNAL_STREAM_URL_REF = "http://192.168.192.96:8888/cam/"
|
|
current_fps = 0.0
|
|
save_queue = queue.Queue(maxsize=100)
|
|
|
|
DASHBOARD_HTML = ""
|
|
|
|
class RTSPStreamReader:
|
|
def __init__(self, source_path):
|
|
self.source_path = source_path
|
|
self.cap = None
|
|
self.frame = None
|
|
self.ret = False
|
|
self.new_frame_event = threading.Event()
|
|
self.running = True
|
|
self.lock = threading.Lock()
|
|
self._last_new_frame_time = time.time()
|
|
self._reconnect_timeout = 10.0 # detik tanpa frame baru sebelum reconnect
|
|
self._stale_threshold = 5.0 # detik tanpa frame baru = dianggap stale
|
|
|
|
self._connect()
|
|
|
|
self.thread = threading.Thread(target=self._update, daemon=True)
|
|
self.thread.start()
|
|
|
|
def _connect(self):
|
|
"""Establish connection to RTSP stream (GStreamer H.265 → H.264 → CPU fallback)."""
|
|
# 1. Pipeline GStreamer H.265 (Jetson NVDEC)
|
|
pipeline_h265 = (
|
|
f"rtspsrc location=\"{self.source_path}\" protocols=tcp latency=0 ! "
|
|
"rtph265depay ! h265parse ! nvv4l2decoder ! "
|
|
"nvvidconv ! video/x-raw, format=BGRx ! "
|
|
"videoconvert ! video/x-raw, format=BGR ! appsink drop=1"
|
|
)
|
|
|
|
# 2. Pipeline GStreamer H.264 (Jetson NVDEC)
|
|
pipeline_h264 = (
|
|
f"rtspsrc location=\"{self.source_path}\" protocols=tcp latency=0 ! "
|
|
"rtph264depay ! h264parse ! nvv4l2decoder ! "
|
|
"nvvidconv ! video/x-raw, format=BGRx ! "
|
|
"videoconvert ! video/x-raw, format=BGR ! appsink drop=1"
|
|
)
|
|
|
|
# Mencoba membuka dengan GStreamer H.265
|
|
print("[INFO] Mencoba GStreamer H.265 NVDEC di Jetson...")
|
|
self.cap = cv2.VideoCapture(pipeline_h265, cv2.CAP_GSTREAMER)
|
|
|
|
# Jika gagal, coba H.264
|
|
if self.cap is None or not self.cap.isOpened():
|
|
print("[INFO] GStreamer H.265 gagal, mencoba GStreamer H.264 NVDEC...")
|
|
self.cap = cv2.VideoCapture(pipeline_h264, cv2.CAP_GSTREAMER)
|
|
|
|
# Fallback ke default CPU OpenCV jika GStreamer tidak terpasang/gagal
|
|
if self.cap is None or not self.cap.isOpened():
|
|
print("[WARNING] GStreamer NVDEC gagal dibuka, menggunakan backend default OpenCV (CPU)...")
|
|
self.cap = cv2.VideoCapture(self.source_path)
|
|
|
|
def _reconnect(self):
|
|
"""Release current capture and reconnect to RTSP stream."""
|
|
print(f"[WARNING] Tidak ada frame baru selama {self._reconnect_timeout:.0f}s. Mencoba reconnect RTSP...")
|
|
try:
|
|
if self.cap is not None:
|
|
self.cap.release()
|
|
except Exception:
|
|
pass
|
|
time.sleep(1.0) # Jeda sebelum reconnect
|
|
self._connect()
|
|
self._last_new_frame_time = time.time()
|
|
if self.cap is not None and self.cap.isOpened():
|
|
print("[INFO] Reconnect RTSP berhasil.")
|
|
else:
|
|
print("[ERROR] Reconnect RTSP gagal. Akan dicoba lagi...")
|
|
|
|
def _update(self):
|
|
consecutive_failures = 0
|
|
while self.running:
|
|
# Auto-reconnect jika tidak ada frame baru terlalu lama
|
|
if time.time() - self._last_new_frame_time > self._reconnect_timeout:
|
|
self._reconnect()
|
|
consecutive_failures = 0
|
|
|
|
if self.cap is None or not self.cap.isOpened():
|
|
time.sleep(0.1)
|
|
continue
|
|
ret, frame = self.cap.read()
|
|
if not ret:
|
|
consecutive_failures += 1
|
|
# Jika gagal berturut-turut, tunggu lebih lama
|
|
if consecutive_failures > 100:
|
|
time.sleep(0.1)
|
|
else:
|
|
time.sleep(0.01)
|
|
continue
|
|
consecutive_failures = 0
|
|
with self.lock:
|
|
self.ret = ret
|
|
self.frame = frame
|
|
self._last_new_frame_time = time.time()
|
|
self.new_frame_event.set()
|
|
time.sleep(0.001)
|
|
|
|
@property
|
|
def is_stale(self):
|
|
"""True jika tidak ada frame baru selama > _stale_threshold detik."""
|
|
return time.time() - self._last_new_frame_time > self._stale_threshold
|
|
|
|
def read(self):
|
|
if self.new_frame_event.wait(timeout=1.0):
|
|
self.new_frame_event.clear()
|
|
with self.lock:
|
|
if self.frame is None:
|
|
return False, None
|
|
return self.ret, self.frame.copy()
|
|
else:
|
|
with self.lock:
|
|
if self.frame is None:
|
|
return False, None
|
|
return self.ret, self.frame.copy()
|
|
|
|
def isOpened(self):
|
|
return self.cap is not None and self.cap.isOpened()
|
|
|
|
def get(self, propId):
|
|
if self.cap is None:
|
|
return 0.0
|
|
return self.cap.get(propId)
|
|
|
|
def release(self):
|
|
self.running = False
|
|
if self.cap is not None and self.cap.isOpened():
|
|
self.cap.release()
|
|
|
|
# =====================================================================
|
|
# 2.5 UTILITY AKURASI (HSV HISTOGRAM & PERSPECTIVE PROFILE)
|
|
# =====================================================================
|
|
def get_visual_features(crop):
|
|
"""Mengekstrak fitur visual berupa histogram HSV (warna) dan grayscale image (struktur/tekstur) dari crop karung."""
|
|
if crop is None or crop.size == 0:
|
|
return None, None
|
|
try:
|
|
resized = cv2.resize(crop, (64, 64))
|
|
hsv = cv2.cvtColor(resized, cv2.COLOR_BGR2HSV)
|
|
# Ekstrak histogram H-S untuk ketahanan terhadap pencahayaan
|
|
hist = cv2.calcHist([hsv], [0, 1], None, [16, 16], [0, 180, 0, 256])
|
|
cv2.normalize(hist, hist, 0, 1, cv2.NORM_MINMAX)
|
|
|
|
# Fitur tekstur/struktur menggunakan grayscale thumbnail
|
|
gray = cv2.cvtColor(resized, cv2.COLOR_BGR2GRAY)
|
|
return hist, gray
|
|
except Exception as e:
|
|
print(f"[ERROR get_visual_features] {e}")
|
|
return None, None
|
|
|
|
def compare_visual_similarity(feat1, feat2):
|
|
"""Membandingkan kemiripan visual karung (gabungan korelasi warna HSV 60% dan struktur grayscale NCC 40%)."""
|
|
if feat1 is None or feat2 is None:
|
|
return 0.0
|
|
hist1, gray1 = feat1
|
|
hist2, gray2 = feat2
|
|
if hist1 is None or hist2 is None or gray1 is None or gray2 is None:
|
|
return 0.0
|
|
try:
|
|
# Kemiripan warna HSV
|
|
color_sim = cv2.compareHist(hist1, hist2, cv2.HISTCMP_CORREL)
|
|
color_sim = max(0.0, color_sim) if not np.isnan(color_sim) else 0.0
|
|
|
|
# Kemiripan tekstur/struktur menggunakan Template Matching Normalized Cross-Correlation (NCC)
|
|
res = cv2.matchTemplate(gray1, gray2, cv2.TM_CCOEFF_NORMED)
|
|
struct_sim = max(0.0, res[0][0]) if not np.isnan(res[0][0]) else 0.0
|
|
|
|
# Rata-rata tertimbang
|
|
return 0.6 * color_sim + 0.4 * struct_sim
|
|
except Exception:
|
|
return 0.0
|
|
|
|
def get_min_valid_area(cy):
|
|
"""Menghitung batas luas area minimum secara dinamis berdasarkan perspektif Y."""
|
|
global scale_x, scale_y
|
|
top_y = 200 * scale_y
|
|
top_area = 8000 * scale_x * scale_y
|
|
bot_y = 1080 * scale_y
|
|
bot_area = 25000 * scale_x * scale_y
|
|
|
|
if cy <= top_y:
|
|
return top_area
|
|
if cy >= bot_y:
|
|
return bot_area
|
|
|
|
ratio = (cy - top_y) / (bot_y - top_y)
|
|
return top_area + ratio * (bot_area - top_area)
|
|
# =====================================================================
|
|
|
|
|
|
def get_zone_name(point):
|
|
pt = Point(point)
|
|
if poly_palet is not None and not poly_palet.is_empty and poly_palet.contains(pt):
|
|
return "PALET"
|
|
elif poly_truck is not None and not poly_truck.is_empty and poly_truck.contains(pt):
|
|
return "TRUCK"
|
|
else:
|
|
return "BEBAS"
|
|
|
|
|
|
def update_zone_label(track_id, current_zone):
|
|
"""Update label visual zona (dengan debounce ringan), TIDAK memengaruhi logika counting."""
|
|
track_zone_history[track_id].append(current_zone)
|
|
history = list(track_zone_history[track_id])
|
|
if len(history) < DEBOUNCE_FRAMES:
|
|
track_confirmed_state[track_id] = current_zone
|
|
return
|
|
most_frequent_zone = max(set(history), key=history.count)
|
|
if history.count(most_frequent_zone) >= (DEBOUNCE_FRAMES - 2):
|
|
track_confirmed_state[track_id] = most_frequent_zone
|
|
|
|
|
|
# =====================================================================
|
|
# 3. RE-ID: PEMULIHAN ID SETELAH OKLUSI
|
|
# =====================================================================
|
|
def check_reid_recovery(new_id, current_centroid, overlap_ratio_now, frame_idx):
|
|
global lost_tracks, counted_at_frame, blocked_without_counting, track_is_valid_bag, static_frames, moving_frames, counted_sacks, lost_counted_sacks, blocked_due_to_duplicate
|
|
if not lost_tracks:
|
|
return False
|
|
|
|
closest_old_id = None
|
|
min_dist = float('inf')
|
|
|
|
pt = Point(current_centroid)
|
|
in_truck_zone = poly_truck is not None and not poly_truck.is_empty and poly_truck.contains(pt)
|
|
|
|
for old_id, info in lost_tracks.items():
|
|
frame_diff = frame_idx - info['frame_idx']
|
|
if frame_diff > MAX_REID_FRAMES:
|
|
continue
|
|
# JIKA track lama sudah terhitung, track baru tidak boleh berada di area palet untuk memulihkannya
|
|
if info['already_counted'] and poly_palet is not None and not poly_palet.is_empty and poly_palet.contains(pt):
|
|
continue
|
|
lc = info['last_centroid']
|
|
dist = np.sqrt((current_centroid[0] - lc[0]) ** 2 + (current_centroid[1] - lc[1]) ** 2)
|
|
|
|
is_consistent_direction = (current_centroid[1] < lc[1] + (50 * scale_y))
|
|
max_dist = MAX_REID_TRANSIT_DISTANCE * (1.0 + 0.01 * frame_diff)
|
|
if (dist < max_dist) and (is_consistent_direction or in_truck_zone):
|
|
if dist < min_dist:
|
|
min_dist = dist
|
|
closest_old_id = old_id
|
|
|
|
if closest_old_id is None:
|
|
return False
|
|
|
|
info = lost_tracks[closest_old_id]
|
|
|
|
# JIKA sudah terhitung (already_counted), langsung pulihkan ID tersebut agar tidak terhitung lagi
|
|
if info['already_counted']:
|
|
already_counted[new_id] = True
|
|
is_locked[new_id] = True
|
|
pending_enter_since[new_id] = None
|
|
pending_exit_since[new_id] = info['pending_exit_since']
|
|
track_zone_history[new_id] = info['zone_history'].copy()
|
|
track_positions[new_id] = info['positions'].copy()
|
|
has_crossed_line[new_id] = info.get('has_crossed_line', True)
|
|
exit_crossed_line[new_id] = info.get('exit_crossed_line', False)
|
|
track_areas[new_id] = info.get('box_area', 0.0)
|
|
track_started_in_truck[new_id] = info.get('started_in_truck', False)
|
|
counted_at_frame[new_id] = info.get('counted_at_frame')
|
|
|
|
# Pulihkan state lingkaran
|
|
track_initial_truck_pos[new_id] = info.get('initial_truck_pos')
|
|
track_truck_entry_frame[new_id] = info.get('truck_entry_frame')
|
|
has_exited_circle[new_id] = info.get('has_exited_circle', False)
|
|
delay_completed[new_id] = info.get('delay_completed', False)
|
|
blocked_without_counting[new_id] = info.get('blocked_without_counting', False)
|
|
track_is_valid_bag[new_id] = info.get('is_valid_bag', False)
|
|
blocked_due_to_duplicate[new_id] = info.get('blocked_due_to_duplicate', False)
|
|
|
|
# Pulihkan state stabilitas
|
|
static_frames[new_id] = info.get('static_frames', 0)
|
|
moving_frames[new_id] = info.get('moving_frames', 0)
|
|
|
|
# Pulihkan posisi terhitung aktif
|
|
if closest_old_id in lost_counted_sacks:
|
|
val = lost_counted_sacks.pop(closest_old_id)
|
|
counted_sacks[new_id] = (val[0], val[1])
|
|
|
|
del lost_tracks[closest_old_id]
|
|
return True
|
|
|
|
# Logika lama untuk yang belum terhitung (pending masuk dll)
|
|
was_counted_or_pending = info['already_counted'] or (info['pending_enter_since'] is not None)
|
|
if was_counted_or_pending and overlap_ratio_now < EXIT_OVERLAP_THRESHOLD:
|
|
if info['already_counted']:
|
|
del lost_tracks[closest_old_id]
|
|
return False
|
|
|
|
already_counted[new_id] = info['already_counted']
|
|
is_locked[new_id] = info['is_locked']
|
|
pending_enter_since[new_id] = info['pending_enter_since']
|
|
pending_exit_since[new_id] = info['pending_exit_since']
|
|
track_zone_history[new_id] = info['zone_history'].copy()
|
|
track_positions[new_id] = info['positions'].copy()
|
|
has_crossed_line[new_id] = info.get('has_crossed_line', False)
|
|
exit_crossed_line[new_id] = info.get('exit_crossed_line', False)
|
|
track_areas[new_id] = info.get('box_area', 0.0)
|
|
track_started_in_truck[new_id] = info.get('started_in_truck', False)
|
|
counted_at_frame[new_id] = info.get('counted_at_frame')
|
|
|
|
# Pulihkan state lingkaran
|
|
track_initial_truck_pos[new_id] = info.get('initial_truck_pos')
|
|
track_truck_entry_frame[new_id] = info.get('truck_entry_frame')
|
|
has_exited_circle[new_id] = info.get('has_exited_circle', False)
|
|
delay_completed[new_id] = info.get('delay_completed', False)
|
|
blocked_without_counting[new_id] = info.get('blocked_without_counting', False)
|
|
track_is_valid_bag[new_id] = info.get('is_valid_bag', False)
|
|
blocked_due_to_duplicate[new_id] = info.get('blocked_due_to_duplicate', False)
|
|
|
|
# Pulihkan state stabilitas
|
|
static_frames[new_id] = info.get('static_frames', 0)
|
|
moving_frames[new_id] = info.get('moving_frames', 0)
|
|
|
|
del lost_tracks[closest_old_id]
|
|
return True
|
|
|
|
|
|
# =====================================================================
|
|
# 3.5 FUNGSI ANTI-DOUBLE COUNT (SPASIAL DUPLIKASI) # --- TAMBAHAN BARU ---
|
|
# =====================================================================
|
|
def cek_duplikat_karung_locked(new_id, cx, cy, box_area, frame_idx, frame=None, bbox=None):
|
|
"""
|
|
Mengecek apakah bbox baru muncul di titik yang sangat dekat dengan
|
|
karung yang SUDAH DIHITUNG (locked), baik yang sedang aktif maupun yang baru hilang.
|
|
Menggunakan visual similarity (HSV Histogram & NCC Grayscale) untuk membedakan penumpukan karung.
|
|
"""
|
|
pt = Point(cx, cy)
|
|
# Karung di palet tidak boleh dideteksi duplikat
|
|
if poly_palet is not None and not poly_palet.is_empty and poly_palet.contains(pt):
|
|
return False
|
|
|
|
# Hanya lakukan duplicate checking jika centroid baru berada di area truk/counting
|
|
is_in_truck = poly_truck is not None and not poly_truck.is_empty and poly_truck.contains(pt)
|
|
if not is_in_truck:
|
|
return False
|
|
|
|
# Ekstrak fitur visual untuk deteksi baru
|
|
new_feat = None
|
|
if frame is not None and bbox is not None:
|
|
x1, y1, x2, y2 = bbox
|
|
crop = frame[max(0, int(y1)):min(frame.shape[0], int(y2)), max(0, int(x1)):min(frame.shape[1], int(x2))]
|
|
new_feat = get_visual_features(crop)
|
|
|
|
# 1. Cek dari track yang SEDANG AKTIF dan SUDAH COUNTED
|
|
for active_id in prev_active_track_ids:
|
|
if active_id != new_id and already_counted.get(active_id, False):
|
|
if active_id in track_positions and len(track_positions[active_id]) > 0:
|
|
last_cx, last_cy = track_positions[active_id][-1]
|
|
dist = np.sqrt((cx - last_cx)**2 + (cy - last_cy)**2)
|
|
|
|
# Cek perbandingan luas area box
|
|
old_area = track_areas.get(active_id, 0)
|
|
if old_area > 0 and box_area > 0:
|
|
area_ratio = min(box_area, old_area) / max(box_area, old_area)
|
|
else:
|
|
area_ratio = 1.0
|
|
|
|
is_similar_size = (dist < 120) or (area_ratio >= 0.40)
|
|
if is_similar_size and dist < JARAK_TOLERANSI_DUPLIKAT:
|
|
# Lakukan verifikasi visual jika fitur tersedia
|
|
old_feat = static_sack_visuals.get(active_id)
|
|
if new_feat is not None and old_feat is not None:
|
|
sim = compare_visual_similarity(new_feat, old_feat)
|
|
if sim > 0.85:
|
|
return True
|
|
else:
|
|
# Fallback jika tidak ada data visual, anggap duplikat secara spasial
|
|
return True
|
|
|
|
# 2. Cek dari track yang SUDAH HILANG (lost_tracks)
|
|
for lost_id, info in lost_tracks.items():
|
|
if info.get('already_counted', False):
|
|
frame_diff = frame_idx - info['frame_idx']
|
|
if frame_diff <= TOLERANSI_FRAME_HILANG:
|
|
last_cx, last_cy = info['last_centroid']
|
|
dist = np.sqrt((cx - last_cx)**2 + (cy - last_cy)**2)
|
|
|
|
# Cek perbandingan luas area box
|
|
old_area = info.get('box_area', 0)
|
|
if old_area > 0 and box_area > 0:
|
|
area_ratio = min(box_area, old_area) / max(box_area, old_area)
|
|
else:
|
|
area_ratio = 1.0
|
|
|
|
is_similar_size = (dist < 120) or (area_ratio >= 0.40)
|
|
if is_similar_size and dist < JARAK_TOLERANSI_DUPLIKAT:
|
|
# Lakukan verifikasi visual jika fitur tersedia
|
|
old_feat = static_sack_visuals.get(lost_id)
|
|
if new_feat is not None and old_feat is not None:
|
|
sim = compare_visual_similarity(new_feat, old_feat)
|
|
if sim > 0.85:
|
|
return True
|
|
else:
|
|
# Fallback
|
|
return True
|
|
|
|
return False
|
|
|
|
|
|
# =====================================================================
|
|
# 4. LOGIKA MASUK / KELUAR
|
|
# =====================================================================
|
|
def update_counting(track_id, overlap_ratio_counting, in_counting_zone, overlap_ratio_truck, frame_idx, required_frames, current_zone, required_exit_frames=None):
|
|
global has_crossed_line, exit_crossed_line, already_counted, pending_enter_since, pending_exit_since, is_locked, metrics, counted_at_frame
|
|
global has_exited_circle, delay_completed, track_started_in_truck
|
|
|
|
if required_exit_frames is None:
|
|
required_exit_frames = required_frames
|
|
|
|
if not already_counted[track_id]:
|
|
# Logika Masuk Baru Berdasarkan Zona:
|
|
# - Zona COUNTING: Centroid di area counting, overlap counting >= 70%
|
|
# - Zona TRUCK: Centroid di area truck, overlap truck >= 70%
|
|
if current_zone == "TRUCK":
|
|
# Jika mulai di dalam truk, kita ijinkan delay berjalan meskipun belum cross line
|
|
# agar saat keluar lingkaran bisa langsung dihitung jika delay sudah selesai.
|
|
is_qualifying_entry = (has_crossed_line[track_id] or track_started_in_truck[track_id]) and (overlap_ratio_truck >= ENTRY_OVERLAP_THRESHOLD)
|
|
else:
|
|
is_qualifying_entry = False
|
|
|
|
if is_qualifying_entry:
|
|
if pending_enter_since[track_id] is None:
|
|
pending_enter_since[track_id] = frame_idx
|
|
else:
|
|
elapsed = frame_idx - pending_enter_since[track_id]
|
|
if elapsed >= required_frames:
|
|
# JIKA masih di dalam lingkaran, jangan dulu counting, "simpan dulu"
|
|
if not has_exited_circle[track_id]:
|
|
delay_completed[track_id] = True
|
|
else:
|
|
metrics['total_masuk'] += 1
|
|
already_counted[track_id] = True
|
|
is_locked[track_id] = True
|
|
pending_enter_since[track_id] = None
|
|
counted_at_frame[track_id] = frame_idx
|
|
else:
|
|
# Jika tidak memenuhi kualifikasi masuk, reset pending timer
|
|
pending_enter_since[track_id] = None
|
|
else:
|
|
pass
|
|
|
|
|
|
# =====================================================================
|
|
# 5. PROSES PREDIKSI & VISUALISASI VIDEO
|
|
# =====================================================================
|
|
def _filter_sacks_in_roi(detections, roi):
|
|
"""Keep only sacks whose centroid X falls within the truck ROI."""
|
|
if roi is None:
|
|
return []
|
|
return [
|
|
d for d in detections
|
|
if roi.contains_x((d.bbox[0] + d.bbox[2]) / 2.0)
|
|
]
|
|
|
|
def run_prediction(model_path, source_path,
|
|
output_json_path="hasil_perhitungan.json", max_frames=None,
|
|
inference_stride=None, sack_conf=None, truck_conf=None,
|
|
box_conf=None, box_model_path=None, model_mode=None,
|
|
output_dir=None, batch_timeout=None,
|
|
config_path=None, cfg=None):
|
|
global prev_active_track_ids, lost_tracks, metrics, track_positions, counted_at_frame
|
|
global track_confirmed_state, already_counted, is_locked, has_crossed_line, exit_crossed_line, track_areas
|
|
global pending_enter_since, pending_exit_since, track_started_in_truck, outside_truck_frames
|
|
global track_initial_truck_pos, track_truck_entry_frame, has_exited_circle, delay_completed, blocked_without_counting, track_is_valid_bag
|
|
global current_fps, INFERENCE_STRIDE, all_counted_sacks_map, last_seen_near_person_frame, blocked_due_to_duplicate
|
|
global width, height, CONFIRM_DELAY_SEC, EXIT_CONFIRM_DELAY_SEC
|
|
global active_batch_info, system_state
|
|
global DUPLICATE_CIRCLE_RADIUS, MIN_VALID_AREA, JARAK_TOLERANSI_DUPLIKAT, MAX_REID_TRANSIT_DISTANCE
|
|
global ENTRY_OVERLAP_THRESHOLD, EXIT_OVERLAP_THRESHOLD, TOLERANSI_FRAME_HILANG
|
|
global MAX_REID_FRAMES, DEBOUNCE_FRAMES
|
|
global CROSS_CLASS_IOU
|
|
global DB_PATH, STATE_FILE, BATCH_MODE_FILE, LIVE_STREAM_FRAME_PATH
|
|
global batch_logger
|
|
global CAMERA_NAME, OBJECT_LABEL, DAILY_CUTOFF_TIME, BATCH_MERGE_THRESHOLD_SECONDS
|
|
global CIRCLE_STAY_TIMEOUT_SEC, _CFG
|
|
|
|
# --- Unified config (config.yaml canonical, .env for secrets/deployment) ---
|
|
global _CFG
|
|
_CFG = cfg or load_config(config_path or _DEFAULT_CONFIG_PATH)
|
|
cfg = _CFG
|
|
|
|
# Legacy path overrides keep their pre-YAML meaning (folded into cfg).
|
|
if model_path and model_path != _legacy_combined_fallback():
|
|
for _k in ("combined", "truck_only"):
|
|
if _k in cfg.models.paths:
|
|
cfg.models.paths[_k] = model_path
|
|
if box_model_path:
|
|
cfg.models.paths["box"] = box_model_path
|
|
for _cls, _v in (("sack", sack_conf), ("truck", truck_conf), ("box", box_conf)):
|
|
if _v is not None and _cls in cfg.models.detection_params:
|
|
cfg.models.detection_params[_cls].conf = float(_v)
|
|
|
|
# --- Path / identity / timing globals from config (output_dir wins) ---
|
|
if output_dir:
|
|
# Cross-platform: plain join, keep config layout when output_dir is None
|
|
DB_PATH = os.path.join(output_dir, "jetson_counter.db")
|
|
STATE_FILE = os.path.join(output_dir, "current_batch.json")
|
|
BATCH_MODE_FILE = os.path.join(output_dir, "batch_mode.json")
|
|
LIVE_STREAM_FRAME_PATH = os.path.join(output_dir, "live_frame.jpg")
|
|
print(f"[INFO] Output dir override: {output_dir}")
|
|
else:
|
|
DB_PATH = os.path.join(cfg.output.dir, cfg.output.db_name)
|
|
STATE_FILE = os.path.join(cfg.output.dir, cfg.output.state_file)
|
|
BATCH_MODE_FILE = os.path.join(cfg.output.dir, cfg.output.batch_mode_file)
|
|
LIVE_STREAM_FRAME_PATH = cfg.output.live_frame_path
|
|
CAMERA_NAME = cfg.camera.name
|
|
OBJECT_LABEL = cfg.camera.object_label
|
|
DAILY_CUTOFF_TIME = cfg.batch.daily_cutoff_time
|
|
BATCH_MERGE_THRESHOLD_SECONDS = cfg.batch.merge_threshold_seconds
|
|
|
|
# --- CSV logger (batch_summary.csv + sack_events.csv) next to the DB ---
|
|
# --no-db means "persist nothing"; construction/call failures only warn.
|
|
if NO_DB:
|
|
batch_logger = None
|
|
else:
|
|
_csv_dir = os.path.dirname(DB_PATH) or "."
|
|
try:
|
|
batch_logger = CSVLogger(_csv_dir)
|
|
print(f"[INFO] CSV logger aktif: {os.path.join(_csv_dir, 'batch_summary.csv')}, "
|
|
f"{os.path.join(_csv_dir, 'sack_events.csv')}")
|
|
except Exception as e:
|
|
batch_logger = None
|
|
print(f"[WARN] CSV logger gagal diinisialisasi: {e}")
|
|
|
|
# --- Counting-knob globals from config (values match legacy zones.json) ---
|
|
CONFIRM_DELAY_SEC = cfg.counting.confirm_delay_sec
|
|
EXIT_CONFIRM_DELAY_SEC = cfg.counting.exit_confirm_delay_sec
|
|
DUPLICATE_CIRCLE_RADIUS = cfg.counting.duplicate_circle_radius
|
|
MIN_VALID_AREA = cfg.counting.min_valid_area
|
|
JARAK_TOLERANSI_DUPLIKAT = cfg.counting.jarak_toleransi_duplikat
|
|
MAX_REID_TRANSIT_DISTANCE = cfg.counting.max_reid_transit_distance
|
|
CIRCLE_STAY_TIMEOUT_SEC = cfg.counting.circle_stay_timeout_sec
|
|
ENTRY_OVERLAP_THRESHOLD = cfg.counting.entry_overlap_threshold
|
|
EXIT_OVERLAP_THRESHOLD = cfg.counting.exit_overlap_threshold
|
|
TOLERANSI_FRAME_HILANG = cfg.counting.tolerance_missing_frames
|
|
MAX_REID_FRAMES = cfg.counting.max_reid_frames
|
|
DEBOUNCE_FRAMES = cfg.counting.debounce_frames
|
|
CROSS_CLASS_IOU = cfg.counting.cross_class_iou
|
|
|
|
# Legacy batch_mode.json no longer drives the mode — nudge once if stale.
|
|
_legacy_warn = check_legacy_batch_mode(cfg, BATCH_MODE_FILE)
|
|
if _legacy_warn:
|
|
print(f"[WARN] {_legacy_warn}")
|
|
|
|
mode = resolve_model_mode(model_mode)
|
|
|
|
# 1. Silencing YOLO logs
|
|
from ultralytics.utils import LOGGER
|
|
import logging
|
|
LOGGER.setLevel(logging.WARNING)
|
|
|
|
INFERENCE_STRIDE = inference_stride if inference_stride is not None else cfg.stream.inference_stride
|
|
|
|
saver = None
|
|
saver_thread = None
|
|
save_queue = None
|
|
|
|
# Reset lists and dicts
|
|
for k in metrics:
|
|
metrics[k] = 0
|
|
|
|
# Device
|
|
device = 'cuda' if torch.cuda.is_available() else 'cpu'
|
|
print(f"[INFO] Device inferensi diset ke: {device}")
|
|
|
|
# Reader
|
|
is_stream = any(str(source_path).startswith(p) for p in ["rtsp://", "rtmp://", "http://", "https://"])
|
|
if is_stream:
|
|
print("[INFO] Membuka RTSP stream menggunakan Threaded GStreamer NVDEC Reader...")
|
|
cap = RTSPStreamReader(source_path)
|
|
else:
|
|
print("[INFO] Membuka file video lokal...")
|
|
cap = cv2.VideoCapture(source_path)
|
|
|
|
if not cap.isOpened():
|
|
print(f"Error: Gagal membuka video source (RTSP stream/file) di {source_path}")
|
|
return
|
|
|
|
width = 1280
|
|
height = 720
|
|
fps = cap.get(cv2.CAP_PROP_FPS)
|
|
if fps <= 0 or np.isnan(fps):
|
|
fps = 25.0
|
|
|
|
print(f"[INFO] Resolusi Asli: {int(cap.get(cv2.CAP_PROP_FRAME_WIDTH))}x{int(cap.get(cv2.CAP_PROP_FRAME_HEIGHT))} @ {fps:.1f} FPS (Diresize ke 1280x720 untuk koordinat tetap)")
|
|
|
|
# Initialize components per mode preset (engines verified to coexist)
|
|
pipe = build_model_pipeline(mode, cfg, device, _BASE_DIR)
|
|
print("[INFO] Warm-up model selesai.")
|
|
|
|
truck_detector = pipe["truck_detector"]
|
|
tracker = pipe["tracker"]
|
|
box_tracker = pipe["box_tracker"]
|
|
separate_truck_model = pipe["separate_truck_model"]
|
|
min_areas = pipe["min_areas"]
|
|
confs = pipe["confs"]
|
|
stabilizer = BboxStabilizer(
|
|
ema_alpha=0.35,
|
|
max_hold_frames=10,
|
|
max_height_ratio=1.5,
|
|
min_height_ratio=0.70,
|
|
)
|
|
# Separate stabilizer for the box tracker's own ID space (modes C/D).
|
|
# Mode B shares the yolo11n tracker (single ID space) so this stays unused.
|
|
box_stabilizer = BboxStabilizer(
|
|
ema_alpha=0.35,
|
|
max_hold_frames=10,
|
|
max_height_ratio=1.5,
|
|
min_height_ratio=0.70,
|
|
)
|
|
|
|
# ================================================================
|
|
# HARDCODED COORDINATES FOR LOCAL (1280x720)
|
|
# Truck detector hanya untuk batch lifecycle (deteksi truk datang/pergi)
|
|
# Area di bawah ini FIXED, tidak tergantung deteksi truk.
|
|
# ================================================================
|
|
|
|
# Calculate scale factors from native 1920x1080 to 1280x720
|
|
scale_x = 1280.0 / 1920.0
|
|
scale_y = 720.0 / 1080.0
|
|
|
|
# 1. Detection Area (4-point Polygon)
|
|
detection_poly_pts = [
|
|
[int(574 * scale_x), int(50 * scale_y)],
|
|
[int(586 * scale_x), int(1077 * scale_y)],
|
|
[int(1418 * scale_x), int(1076 * scale_y)],
|
|
[int(1397 * scale_x), int(50 * scale_y)],
|
|
]
|
|
detection_polygon = Polygon(detection_poly_pts)
|
|
|
|
# 2. Count Line coordinates
|
|
static_line_y = int(330 * scale_y)
|
|
static_line_x_start = int(577 * scale_x)
|
|
static_line_x_end = int(1401 * scale_x)
|
|
|
|
# 3. Truck Area (4-point Polygon for presence check)
|
|
truck_poly_pts = [
|
|
[int(600 * scale_x), int(385 * scale_y)],
|
|
[int(609 * scale_x), int(1076 * scale_y)],
|
|
[int(1404 * scale_x), int(1078 * scale_y)],
|
|
[int(1381 * scale_x), int(343 * scale_y)],
|
|
]
|
|
truck_polygon = Polygon(truck_poly_pts)
|
|
|
|
from src.truck_roi import TruckROI
|
|
static_roi = TruckROI(
|
|
x1=int(600 * scale_x),
|
|
y1=int(343 * scale_y),
|
|
x2=int(1404 * scale_x),
|
|
y2=int(1078 * scale_y),
|
|
line_y=static_line_y,
|
|
confidence=1.0
|
|
)
|
|
|
|
# NOTE: counting-knob globals (DUPLICATE_CIRCLE_RADIUS, MIN_VALID_AREA,
|
|
# JARAK_TOLERANSI_DUPLIKAT, MAX_REID_TRANSIT_DISTANCE, ...) were set from
|
|
# config.yaml at the top of run_prediction — intentionally NOT reset to
|
|
# zones.json REFs here (config.yaml is canonical now).
|
|
|
|
counter = MultiClassLineCounter(
|
|
line_y=static_line_y,
|
|
line_x_start=static_line_x_start,
|
|
line_x_end=static_line_x_end,
|
|
margin=20,
|
|
dedup_radius=float(DUPLICATE_CIRCLE_RADIUS),
|
|
)
|
|
_idle_timeout = batch_timeout if batch_timeout is not None else cfg.batch.timeout_seconds
|
|
_truck_gone = (batch_timeout if batch_timeout is not None
|
|
else cfg.batch.truck_gone_tolerance_seconds)
|
|
batch_mgr = BatchLifecycleManager(
|
|
stabilize_seconds=0.0, # Start batch instantly when triggered by crossing
|
|
stabilize_threshold_px=9999.0, # Disable displacement threshold check
|
|
sack_idle_timeout=_idle_timeout, # tanpa karung masuk & tanpa karung di truk -> WAITING
|
|
min_batch_duration=5.0, # Short min duration
|
|
truck_gone_tolerance=_truck_gone, # WAITING: tanpa karung & tanpa sinyal truk -> END BATCH (config.yaml batch.truck_gone_tolerance_seconds)
|
|
)
|
|
dashboard = DashboardOverlay()
|
|
|
|
global active_batch_info
|
|
if NO_DB:
|
|
active_batch_info = None
|
|
elif os.path.exists(STATE_FILE):
|
|
try:
|
|
with open(STATE_FILE, "r") as sf:
|
|
active_batch_info = json.load(sf)
|
|
except Exception:
|
|
active_batch_info = None
|
|
else:
|
|
active_batch_info = None
|
|
|
|
frame_idx = 0
|
|
last_time = time.time()
|
|
current_fps = 0.0
|
|
prev_manual_active = (active_batch_info is not None and active_batch_info.get("batch_number") is not None)
|
|
prev_auto_active = False
|
|
_prev_truck_in_area = False
|
|
|
|
TRUCK_DET_INTERVAL_IDLE = 5 # Check truck every 5 frames when IDLE
|
|
TRUCK_DET_INTERVAL_STABILIZING = 1 # Check truck every frame when STABILIZING
|
|
|
|
try:
|
|
while cap.isOpened():
|
|
ret, frame = cap.read()
|
|
if not ret:
|
|
if is_stream:
|
|
time.sleep(0.01)
|
|
continue
|
|
else:
|
|
break
|
|
|
|
if frame is not None:
|
|
frame = cv2.resize(frame, (1280, 720))
|
|
|
|
timestamp = time.time()
|
|
frame_idx += 1
|
|
|
|
if max_frames is not None and frame_idx > max_frames:
|
|
break
|
|
|
|
# Stream freeze detection: bekukan batch state agar tidak ditutup prematur
|
|
if is_stream and hasattr(cap, 'is_stale') and cap.is_stale:
|
|
if frame_idx % 75 == 0: # Log setiap ~3 detik (25 FPS * 3)
|
|
print(f"[WARNING] Stream RTSP freeze terdeteksi. Batch timer dibekukan.")
|
|
# Tetap tulis frame terakhir ke dashboard agar tidak blank
|
|
if frame is not None:
|
|
write_live_frame(frame)
|
|
continue
|
|
|
|
# Track previous state for transition detection
|
|
prev_auto_active = batch_mgr.is_active
|
|
prev_state = batch_mgr.state
|
|
|
|
# ================================================================
|
|
# STATE-DRIVEN MODEL SWITCHING
|
|
# ================================================================
|
|
|
|
tracked_sacks = []
|
|
events = []
|
|
|
|
# Run tracker and stabilizer for sacks on every frame
|
|
if INFERENCE_STRIDE <= 1 or frame_idx % INFERENCE_STRIDE == 0 or 'last_raw_tracked_all' not in locals():
|
|
raw_tracked_all = tracker.update(frame, [])
|
|
last_raw_tracked_all = raw_tracked_all
|
|
else:
|
|
raw_tracked_all = last_raw_tracked_all
|
|
|
|
# Filter sack (+box, modes B-D) detections. Per-class conf floor
|
|
# comes from config.yaml detection_params (authoritative).
|
|
# Counter splits by class_name downstream; MultiClassLineCounter
|
|
# ignores anything that is not sack/box.
|
|
# min_bbox_area guardrails come from config.yaml detection_params.
|
|
raw_tracked_sacks = [
|
|
d for d in raw_tracked_all
|
|
if d.class_name in ("sack", "box") and d.confidence >= confs.get(d.class_name, 0.50)
|
|
and _bbox_area_ok(d, min_areas.get(d.class_name, 0))
|
|
]
|
|
|
|
# Truck candidates: from shared tracker (modes A/C) and/or the
|
|
# separate v4 truck model (modes B/D, every 5th frame, cached).
|
|
truck_candidates = [
|
|
d for d in raw_tracked_all
|
|
if d.class_name == "truck" and _bbox_area_ok(d, min_areas.get("truck", 0))
|
|
]
|
|
if separate_truck_model:
|
|
if frame_idx % 5 == 0 or 'last_detected_trucks' not in locals():
|
|
last_detected_trucks = truck_detector.detect(frame)
|
|
truck_candidates += last_detected_trucks
|
|
|
|
# Filter truck detections: Wajib 100% berada di dalam detection_polygon & ambil maksimal 1 bbox terbaik
|
|
valid_trucks = []
|
|
for d in truck_candidates:
|
|
if d.class_name == "truck" and d.confidence >= confs.get("truck", 0.45):
|
|
x1, y1, x2, y2 = d.bbox
|
|
# Overlap >= 50%, NOT 100% containment: a docked truck's bbox
|
|
# rests on the frame bottom (y2 exceeds the polygon's bottom
|
|
# edge by 1-2 px), which made truck_in_area flicker False while
|
|
# the truck never left (batch 15 split, 2026-09-29 14:07).
|
|
# ponytail: area-fraction gate; centroid-inside is the
|
|
# stricter upgrade path if background trucks leak into the band.
|
|
truck_bbox_poly = box(x1, y1, x2, y2)
|
|
area = (x2 - x1) * (y2 - y1)
|
|
if (area > 0 and
|
|
detection_polygon.intersection(truck_bbox_poly).area / area >= 0.5):
|
|
valid_trucks.append((area, d.confidence, d))
|
|
|
|
# Hanya ambil 1 truk (bbox dengan luas area / confidence terbesar) yang valid di area deteksi
|
|
if valid_trucks:
|
|
valid_trucks.sort(key=lambda x: (x[0], x[1]), reverse=True)
|
|
raw_tracked_trucks = [valid_trucks[0][2]]
|
|
else:
|
|
raw_tracked_trucks = []
|
|
|
|
stable = stabilizer.update(raw_tracked_sacks)
|
|
|
|
# Box stream: modes C/D run a dedicated box tracker (own ID space +
|
|
# own stabilizer); mode B shares the yolo11n tracker (single ID space).
|
|
stable_boxes = []
|
|
if box_tracker is not None:
|
|
if INFERENCE_STRIDE <= 1 or frame_idx % INFERENCE_STRIDE == 0 or 'last_raw_tracked_boxes' not in locals():
|
|
raw_tracked_boxes = box_tracker.update(frame, [])
|
|
last_raw_tracked_boxes = raw_tracked_boxes
|
|
else:
|
|
raw_tracked_boxes = last_raw_tracked_boxes
|
|
raw_boxes = [d for d in raw_tracked_boxes
|
|
if d.class_name == "box"
|
|
and _bbox_area_ok(d, min_areas.get("box", 0))]
|
|
stable_boxes = box_stabilizer.update(raw_boxes)
|
|
stable_boxes = [
|
|
d for d in stable_boxes
|
|
if detection_polygon.contains(Point((d.bbox[0] + d.bbox[2]) / 2.0, (d.bbox[1] + d.bbox[3]) / 2.0))
|
|
]
|
|
else:
|
|
# Mode B: boxes already stabilized alongside sacks; mode A: none.
|
|
stable_boxes = [d for d in stable if d.class_name == "box"]
|
|
stable = [d for d in stable if d.class_name != "box"]
|
|
|
|
# Filter using Detection Area (4-point Polygon)
|
|
stable = [
|
|
d for d in stable
|
|
if detection_polygon.contains(Point((d.bbox[0] + d.bbox[2]) / 2.0, (d.bbox[1] + d.bbox[3]) / 2.0))
|
|
]
|
|
|
|
# Combined list for ROI filter / counting / viz (counter splits by class).
|
|
# Cross-model suppression: v4 has no box class (white boxes labeled
|
|
# "sack"); box model always wins. cross_class_iou <= 0 disables.
|
|
if stable_boxes and CROSS_CLASS_IOU > 0:
|
|
stable = drop_sacks_overlapping_boxes(stable, stable_boxes, CROSS_CLASS_IOU)
|
|
stable_all = stable + stable_boxes
|
|
|
|
# Count sacks currently visible in the bottom 85% of truck area (for batch start/end condition)
|
|
min_ty, max_ty = truck_polygon.bounds[1], truck_polygon.bounds[3]
|
|
truck_height = max_ty - min_ty
|
|
truck_cutoff_y = min_ty + 0.15 * truck_height
|
|
|
|
sacks_in_truck_area = 0
|
|
for d in stable:
|
|
cx = (d.bbox[0] + d.bbox[2]) / 2.0
|
|
cy = (d.bbox[1] + d.bbox[3]) / 2.0
|
|
if truck_polygon.contains(Point(cx, cy)) and cy >= truck_cutoff_y:
|
|
sacks_in_truck_area += 1
|
|
|
|
# Cek apakah ada truk valid di area deteksi (1 truk, overlap >= 50% polygon)
|
|
truck_in_area = len(raw_tracked_trucks) > 0
|
|
if truck_in_area != _prev_truck_in_area:
|
|
info = ""
|
|
if raw_tracked_trucks:
|
|
_tb = raw_tracked_trucks[0]
|
|
info = " bbox={} conf={:.2f}".format(
|
|
tuple(int(v) for v in _tb.bbox), _tb.confidence)
|
|
print("[TRUCK] truck_in_area {} -> {} (state={}, valid={}{})".format(
|
|
int(_prev_truck_in_area), int(truck_in_area),
|
|
batch_mgr.state, len(raw_tracked_trucks), info))
|
|
_prev_truck_in_area = truck_in_area
|
|
|
|
# Cek Mode Batch: auto | do_manual | manual (default: auto)
|
|
current_batch_mode = "auto"
|
|
if os.path.exists(BATCH_MODE_FILE):
|
|
try:
|
|
with open(BATCH_MODE_FILE, "r") as bmf:
|
|
bm_data = json.load(bmf)
|
|
raw_mode = bm_data.get("mode", "auto")
|
|
current_batch_mode = raw_mode if raw_mode in (
|
|
"auto", "do_manual", "manual") else "auto"
|
|
except Exception:
|
|
pass
|
|
|
|
roi = static_roi
|
|
|
|
# ================================================================
|
|
# MODE OPERATOR (manual / do_manual: controlled by dashboard API)
|
|
# ================================================================
|
|
if current_batch_mode != "auto":
|
|
manual_batch_exists = False
|
|
if os.path.exists(STATE_FILE):
|
|
try:
|
|
with open(STATE_FILE, "r") as sf:
|
|
s_data = json.load(sf)
|
|
if s_data and s_data.get("batch_number"):
|
|
manual_batch_exists = True
|
|
active_batch_info = s_data
|
|
except Exception:
|
|
pass
|
|
else:
|
|
# Stop transition below needs the last batch id/counts/times
|
|
# (the dashboard already deleted STATE_FILE) — keep them one
|
|
# frame; the stop branch clears them.
|
|
if not prev_manual_active:
|
|
active_batch_info = None
|
|
|
|
is_batch_active = manual_batch_exists
|
|
|
|
# Transition handling for Manual Mode
|
|
if is_batch_active and not prev_manual_active:
|
|
counter.reset()
|
|
stabilizer.reset()
|
|
box_stabilizer.reset()
|
|
print(f"[BATCH] Manual Batch #{active_batch_info.get('batch_number')} dimulai via Tombol.")
|
|
system_state = STATE_COUNTING_SACKS
|
|
elif not is_batch_active and prev_manual_active:
|
|
print(f"[BATCH] Manual Batch dihentikan via Tombol.")
|
|
log_manual_batch_stop()
|
|
active_batch_info = None
|
|
system_state = STATE_WAITING_FOR_TRUCK
|
|
counter.reset()
|
|
stabilizer.reset()
|
|
box_stabilizer.reset()
|
|
|
|
prev_manual_active = is_batch_active
|
|
|
|
if is_batch_active:
|
|
tracked_sacks = _filter_sacks_in_roi(stable_all, static_roi)
|
|
events = counter.update(tracked_sacks)
|
|
for ev in events:
|
|
log_crossing_event(ev, counter,
|
|
active_batch_info["batch_number"], timestamp)
|
|
|
|
if 'cx' in ev and 'cy' in ev:
|
|
counted_sack_positions.append((ev['cx'], ev['cy'], time.time(), ev['track_id']))
|
|
|
|
if active_batch_info is not None:
|
|
active_batch_info["count"] = counter.loading_count
|
|
active_batch_info["box_count"] = counter.box_loading_count
|
|
active_batch_info["box_unloading"] = counter.box_unloading_count
|
|
active_batch_info["unloading"] = counter.unloading_count
|
|
active_batch_info["last_detection_time"] = datetime.now().isoformat()
|
|
save_active_batch_state()
|
|
else:
|
|
tracked_sacks = []
|
|
events = []
|
|
|
|
batch_num_disp = active_batch_info.get("batch_number", "--") if active_batch_info else "--"
|
|
viz_state = STATE_COUNTING_SACKS if is_batch_active else "IDLE"
|
|
viz_progress = 1.0 if is_batch_active else 0.0
|
|
|
|
# ================================================================
|
|
# MODE OTOMATIS (100% Original Algorithm: Truk + Sack Crossing)
|
|
# ================================================================
|
|
else:
|
|
# Run line crossing counter on every frame
|
|
tracked_sacks = _filter_sacks_in_roi(stable_all, static_roi)
|
|
events = counter.update(tracked_sacks)
|
|
has_crossing = len(events) > 0
|
|
|
|
# --- Sack-driven Batch Lifecycle Transitions ---
|
|
if batch_mgr.state in ("IDLE", "TRUCK_STABILIZING"):
|
|
batch_mgr.update_truck(has_crossing, (0.0, 0.0), timestamp)
|
|
|
|
if batch_mgr.state in ("COUNTING_SACKS", "WAITING_FOR_ACTIVITY"):
|
|
batch_mgr.update_sacks(
|
|
has_crossing_event=has_crossing,
|
|
sacks_in_area_count=sacks_in_truck_area,
|
|
timestamp=timestamp,
|
|
loading_count=counter.loading_count,
|
|
unloading_count=counter.unloading_count,
|
|
)
|
|
|
|
if batch_mgr.state == "WAITING_FOR_ACTIVITY":
|
|
activity_detected = has_crossing or (sacks_in_truck_area > 0) or truck_in_area
|
|
batch_mgr.update_truck(activity_detected, None, timestamp)
|
|
|
|
for ev in events:
|
|
log_crossing_event(ev, counter,
|
|
active_batch_info["batch_number"] if active_batch_info else 0,
|
|
timestamp)
|
|
|
|
if 'cx' in ev and 'cy' in ev:
|
|
counted_sack_positions.append((ev['cx'], ev['cy'], time.time(), ev['track_id']))
|
|
|
|
if active_batch_info is not None:
|
|
active_batch_info["count"] = counter.loading_count
|
|
active_batch_info["box_count"] = counter.box_loading_count
|
|
active_batch_info["box_unloading"] = counter.box_unloading_count
|
|
active_batch_info["unloading"] = counter.unloading_count
|
|
active_batch_info["last_detection_time"] = datetime.now().isoformat()
|
|
save_active_batch_state()
|
|
|
|
# Batch just started (STABILIZING -> COUNTING)
|
|
if batch_mgr.is_active and not prev_auto_active:
|
|
counting_date = get_counting_date()
|
|
batch_num = get_next_batch_number(counting_date)
|
|
now_iso = datetime.now().isoformat()
|
|
active_batch_info = {
|
|
"counting_date": counting_date,
|
|
"batch_number": batch_num,
|
|
"count": 0,
|
|
"box_count": 0,
|
|
"model_mode": mode,
|
|
"start_time": now_iso,
|
|
"last_detection_time": now_iso
|
|
}
|
|
save_active_batch_state()
|
|
print(f"[BATCH] Sesi batch #{batch_num} dimulai secara otomatis.")
|
|
system_state = STATE_COUNTING_SACKS
|
|
|
|
# Batch just ended (WAITING -> IDLE, truck left)
|
|
elif not batch_mgr.is_active and prev_auto_active:
|
|
final_count = counter.loading_count
|
|
start_iso = active_batch_info["start_time"] if active_batch_info else datetime.now().isoformat()
|
|
end_iso = datetime.now().isoformat()
|
|
batch_num = active_batch_info["batch_number"] if active_batch_info else 0
|
|
finalize_batch(final_count, start_iso, end_iso,
|
|
box_final_count=counter.box_loading_count,
|
|
box_unloading_count=counter.box_unloading_count,
|
|
model_mode=mode,
|
|
unloading_count=counter.unloading_count)
|
|
print(f"[BATCH] Truk pergi. Sesi batch #{batch_num} selesai secara otomatis. "
|
|
f"Total karung: {final_count}, box: {counter.box_loading_count}.")
|
|
system_state = STATE_WAITING_FOR_TRUCK
|
|
|
|
counter.reset()
|
|
stabilizer.reset()
|
|
box_stabilizer.reset()
|
|
|
|
if batch_mgr.state == "TRUCK_STABILIZING" and prev_state != "TRUCK_STABILIZING":
|
|
system_state = "TRUCK_STABILIZING"
|
|
elif batch_mgr.state == "WAITING_FOR_ACTIVITY" and prev_state != "WAITING_FOR_ACTIVITY":
|
|
system_state = "WAITING_FOR_ACTIVITY"
|
|
elif batch_mgr.state == "COUNTING_SACKS" and prev_state == "WAITING_FOR_ACTIVITY":
|
|
system_state = STATE_COUNTING_SACKS
|
|
|
|
prev_auto_active = batch_mgr.is_active
|
|
is_batch_active = batch_mgr.is_active
|
|
batch_num_disp = batch_mgr.current_batch_id
|
|
viz_state = batch_mgr.state
|
|
viz_progress = batch_mgr.stabilize_progress
|
|
|
|
# Sync counts to metrics so APIs get correct results
|
|
metrics['total_masuk'] = counter.loading_count
|
|
metrics['total_keluar'] = counter.unloading_count
|
|
metrics['box_masuk'] = counter.box_loading_count
|
|
metrics['box_keluar'] = counter.box_unloading_count
|
|
|
|
# 4. Draw Dashboard visualization overlay
|
|
viz = dashboard.draw(
|
|
frame=frame,
|
|
detections=stable if (is_batch_active and SHOW_ALL_BBOXES and 'stable' in dir()) else tracked_sacks,
|
|
roi=roi,
|
|
loading_count=counter.loading_count,
|
|
unloading_count=counter.unloading_count,
|
|
batch_id=batch_num_disp,
|
|
history=batch_mgr.history if current_batch_mode == "auto" else [],
|
|
system_state=viz_state,
|
|
batch_duration=batch_mgr.batch_duration if current_batch_mode == "auto" else 0,
|
|
idle_timer=batch_mgr.time_since_last_sack_activity if current_batch_mode == "auto" else 0,
|
|
stabilize_progress=viz_progress,
|
|
waiting_duration=batch_mgr.waiting_duration if current_batch_mode == "auto" else 0,
|
|
)
|
|
|
|
# Draw Truck Bounding Boxes on viz frame
|
|
if 'raw_tracked_trucks' in locals() and raw_tracked_trucks:
|
|
for trk in raw_tracked_trucks:
|
|
tx1, ty1, tx2, ty2 = [int(v) for v in trk.bbox]
|
|
t_label = f"TRUCK {trk.confidence:.0%}"
|
|
if trk.track_id is not None:
|
|
t_label = f"TRUCK #{trk.track_id} {trk.confidence:.0%}"
|
|
cv2.rectangle(viz, (tx1, ty1), (tx2, ty2), (0, 165, 255), 2) # Orange color
|
|
# Label tag background
|
|
lbl_size, _ = cv2.getTextSize(t_label, cv2.FONT_HERSHEY_SIMPLEX, 0.5, 1)
|
|
cv2.rectangle(viz, (tx1, max(0, ty1 - 20)), (tx1 + lbl_size[0] + 6, ty1), (0, 165, 255), -1)
|
|
cv2.putText(viz, t_label, (tx1 + 3, max(14, ty1 - 5)), cv2.FONT_HERSHEY_SIMPLEX, 0.5, (0, 0, 0), 1, cv2.LINE_AA)
|
|
|
|
# Draw active Duplicate Radius Circles on viz frame (terkini 3.0 detik)
|
|
if 'counted_sack_positions' in globals() and counted_sack_positions:
|
|
rad_vis = DUPLICATE_CIRCLE_RADIUS if ('DUPLICATE_CIRCLE_RADIUS' in globals() and DUPLICATE_CIRCLE_RADIUS > 0) else 60
|
|
now_t = time.time()
|
|
# Clean up expired entries in-place to avoid memory accumulation
|
|
counted_sack_positions[:] = [p for p in counted_sack_positions if len(p) >= 3 and (now_t - p[2]) <= 3.0]
|
|
for pos_item in counted_sack_positions:
|
|
px, py = pos_item[0], pos_item[1]
|
|
tid = pos_item[3] if len(pos_item) > 3 else 0
|
|
cv2.circle(viz, (int(px), int(py)), int(rad_vis), (0, 255, 255), 2, lineType=cv2.LINE_AA)
|
|
cv2.circle(viz, (int(px), int(py)), 4, (0, 255, 0), -1)
|
|
cv2.putText(viz, f"DEDUP #{tid}", (int(px) - 25, max(15, int(py) - int(rad_vis) - 5)), cv2.FONT_HERSHEY_SIMPLEX, 0.45, (0, 255, 255), 1)
|
|
|
|
cv2.putText(viz, f"RADIUS DEDUP: {rad_vis}px", (viz.shape[1] - 270, 40), cv2.FONT_HERSHEY_SIMPLEX, 0.65, (0, 255, 255), 2)
|
|
|
|
# Write live frame to RAM disk for Flask port 5000
|
|
if frame_idx % 2 == 0:
|
|
write_live_frame(viz)
|
|
|
|
# Compute FPS every 25 frames
|
|
if frame_idx % 25 == 0:
|
|
elapsed = time.time() - last_time
|
|
current_fps = 25.0 / elapsed if elapsed > 0 else 0
|
|
last_time = time.time()
|
|
if NO_DASHBOARD:
|
|
continue
|
|
try:
|
|
status_file = os.getenv('LIVE_STATUS_FILE', '/dev/shm/jetson-counter/live_status.json' if os.name != 'nt' else 'd:/Belajar/menghitung karung/live_status.json')
|
|
os.makedirs(os.path.dirname(status_file), exist_ok=True)
|
|
with open(status_file, 'w') as f:
|
|
json.dump({"fps": round(current_fps, 1)}, f)
|
|
except Exception:
|
|
pass
|
|
|
|
finally:
|
|
# Membersihkan dan menutup semua resource
|
|
if saver_thread is not None:
|
|
save_queue.put(None)
|
|
saver_thread.join(timeout=2.0)
|
|
if saver is not None:
|
|
saver.release()
|
|
cap.release()
|
|
cv2.destroyAllWindows()
|
|
|
|
net_count = metrics['total_masuk'] - metrics['total_keluar']
|
|
box_net = metrics['box_masuk'] - metrics['box_keluar']
|
|
final_results = {
|
|
"total_masuk_truck": metrics['total_masuk'],
|
|
"total_keluar_truck": metrics['total_keluar'],
|
|
"net_karung_di_truck": net_count,
|
|
"box_masuk_truck": metrics['box_masuk'],
|
|
"box_keluar_truck": metrics['box_keluar'],
|
|
"net_box_di_truck": box_net,
|
|
}
|
|
|
|
with open(output_json_path, 'w') as f:
|
|
json.dump(final_results, f, indent=4)
|
|
print("\n" + "=" * 50)
|
|
print("PROSES SELESAI!")
|
|
print(final_results)
|
|
|
|
|
|
if __name__ == "__main__":
|
|
args = parse_args()
|
|
|
|
# Load .env (systemd already injects env via EnvironmentFile; load_dotenv
|
|
# never overrides existing vars, so this is a safe no-op in production)
|
|
try:
|
|
from dotenv import load_dotenv
|
|
load_dotenv(args.env)
|
|
except ImportError:
|
|
print("[WARN] python-dotenv tidak tersedia, memakai environment apa adanya.")
|
|
|
|
NO_DASHBOARD = args.no_dashboard
|
|
NO_DB = args.no_db
|
|
|
|
# Unified config first (config.yaml canonical; .env supplies secrets/RTSP).
|
|
_CFG = load_config(args.config or _DEFAULT_CONFIG_PATH)
|
|
_CFG = apply_cli_overrides_to_config(_CFG, args)
|
|
|
|
if args.source:
|
|
SOURCE_INPUT = args.source
|
|
else:
|
|
_env_url = _CFG.stream.rtsp_url
|
|
SOURCE_INPUT = _env_url if _env_url else ("anomali.mp4" if (os.path.exists("anomali.mp4") and os.name == 'nt') else "rtsp://192.168.192.96:8554/cam")
|
|
OUTPUT_JSON = args.output_json or "hasil_perhitungan.json"
|
|
|
|
print(f"[INFO] config={args.config or _DEFAULT_CONFIG_PATH} source={SOURCE_INPUT} "
|
|
f"model_mode={resolve_model_mode(args.model_mode)} "
|
|
f"no_dashboard={NO_DASHBOARD} no_db={NO_DB}")
|
|
|
|
try:
|
|
run_prediction(
|
|
model_path=args.model,
|
|
source_path=SOURCE_INPUT,
|
|
output_json_path=OUTPUT_JSON,
|
|
max_frames=args.max_frames,
|
|
sack_conf=args.sack_conf,
|
|
truck_conf=args.truck_conf,
|
|
box_conf=args.box_conf,
|
|
box_model_path=args.box_model,
|
|
model_mode=args.model_mode,
|
|
output_dir=args.output_dir,
|
|
batch_timeout=args.batch_timeout,
|
|
config_path=args.config or _DEFAULT_CONFIG_PATH,
|
|
cfg=_CFG,
|
|
)
|
|
except KeyboardInterrupt:
|
|
print("\n" + "=" * 50)
|
|
print("[INFO] Program dihentikan secara manual (Ctrl+C).")
|
|
print("Membersihkan resource dan menyimpan hasil perhitungan terakhir...")
|
|
|
|
# Simpan hasil perhitungan parsial sebelum keluar
|
|
final_results = {
|
|
"total_masuk_truck": metrics['total_masuk'],
|
|
"total_keluar_truck": metrics['total_keluar'],
|
|
"net_karung_di_truck": metrics['total_masuk'] - metrics['total_keluar'],
|
|
"box_masuk_truck": metrics['box_masuk'],
|
|
"box_keluar_truck": metrics['box_keluar'],
|
|
"net_box_di_truck": metrics['box_masuk'] - metrics['box_keluar'],
|
|
}
|
|
with open(OUTPUT_JSON, 'w') as f:
|
|
json.dump(final_results, f, indent=4)
|
|
|
|
print("Hasil akhir yang disimpan:")
|
|
print(final_results)
|
|
print("=" * 50) |