Files
karung-counting-feedmill-se…/predict.py
T

2044 lines
91 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
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=_idle_timeout, # tanpa truk & tanpa karung -> END BATCH
)
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
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
# Bounding box truk 100% harus berada di dalam detection_polygon
truck_bbox_poly = box(x1, y1, x2, y2)
if detection_polygon.contains(truck_bbox_poly):
area = (x2 - x1) * (y2 - y1)
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 dalam area deteksi (tepat 1 truk yang 100% di dalam area)
truck_in_area = len(raw_tracked_trucks) > 0
# 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._truck_gone_tolerance = 30.0
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)