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 from src.batch import BatchLifecycleManager, BatchRecord 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 # --- 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 finalize_batch(final_count, start_time_iso, end_time_iso, box_final_count=0, box_unloading_count=0, model_mode="?"): global active_batch_info if active_batch_info is None: return 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 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() # Box columns may not exist on pre-migration DBs — tolerant update. try: cur2 = conn.cursor() cur2.execute(""" UPDATE batches SET box_loading = ?, box_unloading = ?, model_mode = ? WHERE counting_date = ? AND batch_number = ? AND camera_name = ? AND object_label = ? """, (box_final_count, box_unloading_count, model_mode, counting_date, batch_num, CAMERA_NAME, OBJECT_LABEL)) 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}).") 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): 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 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, }, } # ===================================================================== # ===================================================================== # 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 DB_PATH, STATE_FILE, BATCH_MODE_FILE, LIVE_STREAM_FRAME_PATH 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 # --- 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 # 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"] 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 (confidence >= 0.50). # 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 >= 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 >= 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). 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' vs 'manual' (default: 'manual') current_batch_mode = "manual" if os.path.exists(BATCH_MODE_FILE): try: with open(BATCH_MODE_FILE, "r") as bmf: bm_data = json.load(bmf) current_batch_mode = bm_data.get("mode", "manual") except Exception: pass roi = static_roi # ================================================================ # MODE MANUAL (Controlled by Operator Button / Dashboard API) # ================================================================ if current_batch_mode == "manual": 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: 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.") 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: _label = "KARUNG" if ev.get("class_name", "sack") == "sack" else "BOX" _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 '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["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: _label = "KARUNG" if ev.get("class_name", "sack") == "sack" else "BOX" _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 '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["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) 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)