Files
asus 5c7c122105 feat: add counting bench, triage, and dataset modules
This commit includes major additions and updates to the frontend and backend architectures, introducing new dataset management, live counting features, batch processing, and triage logic. Includes new UI pages, components, and API routes.
2026-08-14 16:28:52 +07:00

1405 lines
60 KiB
Python

import os
os.environ["OPENCV_FFMPEG_CAPTURE_OPTIONS"] = "rtsp_transport;tcp|buffer_size;20480000|max_delay;500000|reorder_queue_size;500"
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
from src.tracking import ByteTrackTracker
from src.stabilizer import BboxStabilizer
from src.truck_roi import TruckROITracker
from src.counting import LineCrossCounter
from src.batch import BatchLifecycleManager, BatchRecord
from src.dashboard import DashboardOverlay
# --- 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'
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')
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')
# A shift runs 06:00 to 06:00, so that is where the counting day turns over and
# batch_number restarts at 1. The old default of 20:00 labelled the whole day
# shift as the *previous* date: a batch at 08:27 on the 13th was filed under the
# 12th. It disagreed with the archive's cycles in 3 of 8 boundary cases tested;
# at 06:00 the two agree exactly. Still overridable per deployment.
DAILY_CUTOFF_TIME = os.getenv('DAILY_CUTOFF_TIME', '06:00')
BATCH_MERGE_THRESHOLD_SECONDS = int(os.getenv('BATCH_MERGE_THRESHOLD_SECONDS', '300'))
active_batch_info = None
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)
)
""")
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 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 finalize_batch(final_count, start_time_iso, end_time_iso):
global active_batch_info
if active_batch_info is None:
return
if 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()
conn.close()
print(f"[DB Info] Sesi batch #{batch_num} disimpan ke database SQLite: {final_count} karung.")
except Exception as e:
print(f"[DB Error] Gagal menyimpan batch ke database: {e}")
active_batch_info = None
save_active_batch_state()
def write_live_frame(frame):
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 35
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
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 ---
# Model gabungan karung + truk (menggantikan truck-detector dan best)
if os.name == 'nt':
COMBINED_MODEL_PATH = r"D:\Belajar\Menghitung karung\v4-best (1).pt"
else:
COMBINED_MODEL_PATH = "model_karung_trukf.engine" if os.path.exists("model_karung_trukf.engine") else "v4-best (1).pt"
# =====================================================================
# =====================================================================
# 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():
global ZONA_PALET_REF, ZONA_TRUCK_REF, GARIS_COUNTING_REF, DUPLICATE_CIRCLE_RADIUS_REF
global MIN_VALID_AREA_REF, JARAK_TOLERANSI_DUPLIKAT_REF, MAX_REID_TRANSIT_DISTANCE_REF
global CIRCLE_STAY_TIMEOUT_SEC, INFERENCE_STRIDE, CONFIRM_DELAY_SEC, EXIT_CONFIRM_DELAY_SEC
global ZONA_COUNTING_REF, left_limit_ref, right_limit_ref, EXTERNAL_STREAM_URL_REF
if os.path.exists(ZONES_JSON_PATH):
try:
with open(ZONES_JSON_PATH, 'r') as f:
data = json.load(f)
ZONA_PALET_REF = np.array(data.get('palet', []), dtype=np.int32)
ZONA_TRUCK_REF = np.array(data.get('truck', []), dtype=np.int32)
ZONA_COUNTING_REF = get_bottom_quarter(ZONA_TRUCK_REF)
left_limit_ref = float(data.get('left_limit', 0.05))
right_limit_ref = float(data.get('right_limit', 0.95))
GARIS_COUNTING_REF = ZONA_TRUCK_REF.copy()
DUPLICATE_CIRCLE_RADIUS_REF = data.get('duplicate_circle_radius', 60)
MIN_VALID_AREA_REF = data.get('min_valid_area', 15000)
JARAK_TOLERANSI_DUPLIKAT_REF = data.get('jarak_toleransi_duplikat', 20)
MAX_REID_TRANSIT_DISTANCE_REF = data.get('max_reid_transit_distance', 400)
CIRCLE_STAY_TIMEOUT_SEC = data.get('circle_stay_timeout_sec', 10.0)
INFERENCE_STRIDE = data.get('inference_stride', 2)
CONFIRM_DELAY_SEC = data.get('confirm_delay_sec', 0.5)
EXIT_CONFIRM_DELAY_SEC = data.get('exit_confirm_delay_sec', 6.0)
EXTERNAL_STREAM_URL_REF = data.get('external_stream_url', 'http://192.168.192.96:8888/cam/')
if not EXTERNAL_STREAM_URL_REF:
EXTERNAL_STREAM_URL_REF = 'http://192.168.192.96:8888/cam/'
print("[INFO] Berhasil memuat koordinat zona dan parameter kalibrasi dari zones.json")
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
}
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
# 1. Pipeline GStreamer H.265 (Jetson NVDEC)
pipeline_h265 = (
f"rtspsrc location=\"{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=\"{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(source_path)
self.frame = None
self.ret = False
self.new_frame_event = threading.Event()
self.running = True
self.lock = threading.Lock()
self.thread = threading.Thread(target=self._update, daemon=True)
self.thread.start()
def _update(self):
while self.running:
if not self.cap.isOpened():
time.sleep(0.1)
continue
ret, frame = self.cap.read()
if not ret:
time.sleep(0.01)
continue
with self.lock:
self.ret = ret
self.frame = frame
self.new_frame_event.set()
time.sleep(0.001)
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.isOpened()
def get(self, propId):
return self.cap.get(propId)
def release(self):
self.running = False
if 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=2):
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
# 1. Silencing YOLO logs
from ultralytics.utils import LOGGER
import logging
LOGGER.setLevel(logging.WARNING)
INFERENCE_STRIDE = 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 from repo rafan using shared model
print(f"[INFO] Memuat model YOLO gabungan dari: {model_path}")
shared_model = YOLO(model_path)
# Warm-up model to initialize CUDA/TensorRT execution context and prevent segfaults on tracking
print("[INFO] Melakukan warm-up model YOLO...")
dummy_frame = np.zeros((720, 1280, 3), dtype=np.uint8)
_ = shared_model(dummy_frame, imgsz=640, device=device, verbose=False)
print("[INFO] Warm-up model selesai.")
truck_detector = TruckDetector(shared_model, 0.45)
tracker = ByteTrackTracker(shared_model, 0.45)
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
)
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
counter = LineCrossCounter(
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),
)
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=10.0, # Waiting state timeout
min_batch_duration=5.0, # Short min duration
truck_gone_tolerance=15.0, # Time to wait (in seconds) after all sacks disappear before closing batch
)
dashboard = DashboardOverlay()
global active_batch_info
active_batch_info = None
save_active_batch_state()
frame_idx = 0
last_time = time.time()
current_fps = 0.0
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
# Track previous state for transition detection
prev_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 detections (confidence >= 0.50)
raw_tracked_sacks = [d for d in raw_tracked_all if d.class_name == "sack" and d.confidence >= 0.50]
# Filter truck detections for batch lifecycle guard (confidence >= 0.45)
raw_tracked_trucks = [d for d in raw_tracked_all if d.class_name == "truck" and d.confidence >= 0.45]
stable = stabilizer.update(raw_tracked_sacks)
# 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))
]
# Count sacks currently visible in the bottom 70% 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.30 * 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 truk masih terdeteksi di area truk (untuk batch lifecycle guard)
truck_in_area = any(
truck_polygon.contains(Point((d.bbox[0] + d.bbox[2]) / 2.0, (d.bbox[1] + d.bbox[3]) / 2.0))
for d in raw_tracked_trucks
)
# Run line crossing counter on every frame
tracked_sacks = _filter_sacks_in_roi(stable, 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"):
# Trigger batch start ONLY when a sack actually crosses the line
batch_mgr.update_truck(has_crossing, (0.0, 0.0), timestamp)
if batch_mgr.state in ("COUNTING_SACKS", "WAITING_FOR_ACTIVITY"):
# Toleransi dinamis berdasarkan jumlah karung terhitung
# (toleransi ini hanya aktif setelah truk TIDAK terdeteksi di kamera)
current_count = counter.loading_count
if current_count < 20:
batch_mgr._truck_gone_tolerance = 60.0 # 60 detik jika < 20 karung
elif current_count >= 40:
batch_mgr._truck_gone_tolerance = 30.0 # 30 detik jika >= 40 karung
else:
batch_mgr._truck_gone_tolerance = 45.0 # 45 detik jika di antara 20 - 39
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 waiting, check if sacks are completely gone to finalize batch
if batch_mgr.state == "WAITING_FOR_ACTIVITY":
# Batch tetap terbuka selama truk ATAU karung masih terdeteksi di area
# Ini mencegah batch ditutup prematur saat karung di blind spot kamera
activity_detected = sacks_in_truck_area > 0 or truck_in_area
batch_mgr.update_truck(activity_detected, None, timestamp)
# Process crossing events
for ev in events:
print(f"[KARUNG] Karung #{ev['track_id']} masuk.")
print(f"[TOTAL] Total karung saat ini: {counter.loading_count}.")
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["last_detection_time"] = datetime.now().isoformat()
save_active_batch_state()
# ================================================================
# BATCH TRANSITION HANDLING
# ================================================================
# ROI and counting line are always visible
roi = static_roi
# Batch just started (STABILIZING -> COUNTING)
if batch_mgr.is_active and not prev_active:
counting_date = get_counting_date()
last_batch = get_last_batch_info(counting_date)
# Disable batch resume/merge logic. Every session is a brand new batch.
should_resume = False
if should_resume:
batch_num = last_batch["batch_number"]
prev_count = last_batch["count"]
start_iso = last_batch["start_time"]
# Set the counter's starting count
counter._loading_count = prev_count
counter._unloading_count = 0 # Assuming loading session
# Resume the batch in the lifecycle manager
try:
start_dt = datetime.fromisoformat(start_iso)
start_ts = start_dt.timestamp()
except Exception:
start_ts = timestamp
batch_mgr.resume_batch(
batch_id=batch_num,
start_time=start_ts,
loading_count=prev_count,
unloading_count=0
)
active_batch_info = {
"counting_date": counting_date,
"batch_number": batch_num,
"count": prev_count,
"start_time": start_iso,
"last_detection_time": datetime.now().isoformat()
}
save_active_batch_state()
print(f"[BATCH] Melanjutkan sesi batch #{batch_num} (selisih waktu: {gap_seconds:.1f}s < {BATCH_MERGE_THRESHOLD_SECONDS}s). Mulai dari {prev_count} karung.")
else:
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,
"start_time": now_iso,
"last_detection_time": now_iso
}
save_active_batch_state()
print(f"[BATCH] Sesi batch #{batch_num} dimulai.")
system_state = STATE_COUNTING_SACKS
# Batch just ended (WAITING -> IDLE, truck left)
elif not batch_mgr.is_active and prev_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)
print(f"[BATCH] Truk pergi. Sesi batch #{batch_num} selesai. Total karung: {final_count}.")
system_state = STATE_WAITING_FOR_TRUCK
# Reset counter and trackers for next batch
counter.reset()
stabilizer.reset()
# Update system_state for display
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 # Resumed from waiting
# Sync counts to metrics so APIs get correct results
metrics['total_masuk'] = counter.loading_count
metrics['total_keluar'] = counter.unloading_count
# 4. Draw Dashboard visualization overlay
viz = dashboard.draw(
frame=frame,
detections=stable if (batch_mgr.is_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_mgr.current_batch_id,
history=batch_mgr.history,
system_state=batch_mgr.state,
batch_duration=batch_mgr.batch_duration,
idle_timer=batch_mgr.time_since_last_sack_activity,
stabilize_progress=batch_mgr.stabilize_progress,
waiting_duration=batch_mgr.waiting_duration,
)
# 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()
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']
final_results = {
"total_masuk_truck": metrics['total_masuk'],
"total_keluar_truck": metrics['total_keluar'],
"net_karung_di_truck": net_count
}
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__":
MODEL_FILE = COMBINED_MODEL_PATH
# Default to local sample video 0727.mp4 on Windows
SOURCE_INPUT = "0727s7.mp4" if (os.path.exists("0727s7.mp4") and os.name == 'nt') else "rtsp://192.168.192.96:8554/cam"
OUTPUT_JSON = "hasil_perhitungan.json"
try:
run_prediction(
model_path=MODEL_FILE,
source_path=SOURCE_INPUT,
output_json_path=OUTPUT_JSON,
max_frames=None
)
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']
}
with open(OUTPUT_JSON, 'w') as f:
json.dump(final_results, f, indent=4)
print("Hasil akhir yang disimpan:")
print(final_results)
print("=" * 50)