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 time import sqlite3 import threading from pathlib import Path from dataclasses import replace from datetime import datetime from ultralytics import YOLO # Import kustom dari repositori rpo iki (Count Engine buatan teman) from count import ( LineCounter, BoundarySettings, SACK_CLASS_ID, DEFAULT_MASK_ALPHA, annotate_tracks, draw_persisted_sacks, draw_count_hud, draw_truck_counter_box, draw_count_flashes, draw_boundary, draw_blind_truck_overlay, tick_flashes, track_points, tracking_point, CountFlash, load_counting_params ) # ===================================================================== # PATH DATABASES & CONFIGURATION FOR DASHBOARD (PORT 5000 & 8000) # ===================================================================== if os.name == 'nt': _DEFAULT_DIR = "d:/Belajar/menghitung karung" else: _DEFAULT_DIR = "/opt/jetson-counter" DB_PATH = os.getenv('DB_PATH', f"{_DEFAULT_DIR}/jetson_counter.db") STATE_FILE = os.getenv('STATE_FILE', f"{_DEFAULT_DIR}/current_batch.json") LIVE_STREAM_FRAME_PATH = os.getenv('LIVE_STREAM_FRAME_PATH', f"{_DEFAULT_DIR}/live_frame.jpg") SHM_LIVE_FRAME_PATH = "/dev/shm/jetson-counter/live_frame.jpg" CAMERA_NAME = "CC1" OBJECT_LABEL = "Karung Feedmill (RPO IKI Engine)" class RTSPBufferlessCapture: """Thread-safe RTSP Reader untuk Jetson / Windows.""" def __init__(self, source_path): self.source_path = source_path self.lock = threading.Lock() self.cap = cv2.VideoCapture(source_path, cv2.CAP_FFMPEG) self.frame = None self.ret = False self.running = True if self.cap.isOpened(): self.cap.set(cv2.CAP_PROP_BUFFERSIZE, 1) self.width = int(self.cap.get(cv2.CAP_PROP_FRAME_WIDTH)) self.height = int(self.cap.get(cv2.CAP_PROP_FRAME_HEIGHT)) self.fps = self.cap.get(cv2.CAP_PROP_FPS) else: self.width, self.height, self.fps = 1920, 1080, 25.0 if self.fps <= 0 or np.isnan(self.fps): self.fps = 25.0 self.thread = threading.Thread(target=self._update, daemon=True) self.thread.start() def _update(self): while self.running: if self.cap is None or not self.cap.isOpened(): time.sleep(0.05) continue ret, frame = self.cap.read() if ret and frame is not None: with self.lock: self.frame = frame self.ret = True else: time.sleep(0.005) def isOpened(self): return self.cap is not None and self.cap.isOpened() def get(self, propId): if propId == cv2.CAP_PROP_FRAME_WIDTH: return self.width elif propId == cv2.CAP_PROP_FRAME_HEIGHT: return self.height elif propId == cv2.CAP_PROP_FPS: return self.fps return 0 def read(self): with self.lock: if self.ret and self.frame is not None: return True, self.frame.copy() return False, None def release(self): self.running = False if self.cap is not None: self.cap.release() self.cap = None def init_db(): try: os.makedirs(os.path.dirname(DB_PATH), exist_ok=True) conn = sqlite3.connect(DB_PATH) cursor = conn.cursor() cursor.execute(""" CREATE TABLE IF NOT EXISTS batches ( id INTEGER PRIMARY KEY AUTOINCREMENT, batch_number TEXT UNIQUE, start_time TEXT, end_time TEXT, total_masuk INTEGER, total_keluar INTEGER, net_count INTEGER, status TEXT ) """) cursor.execute(""" CREATE TABLE IF NOT EXISTS daily_summaries ( id INTEGER PRIMARY KEY AUTOINCREMENT, date TEXT UNIQUE, camera_name TEXT, object_label TEXT, total_count INTEGER, total_batches INTEGER, updated_at TEXT ) """) 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}") from http.server import HTTPServer, BaseHTTPRequestHandler from socketserver import ThreadingMixIn streaming_frame = None streaming_lock = threading.Lock() live_stream_enabled = True class StreamingHandler(BaseHTTPRequestHandler): def log_message(self, format, *args): pass def do_GET(self): global streaming_frame if self.path == '/' or self.path == '/stream.mjpg' or self.path == '/stream': self.send_response(200) self.send_header('Age', '0') self.send_header('Cache-Control', 'no-cache, private') self.send_header('Pragma', 'no-cache') self.send_header('Content-Type', 'multipart/x-mixed-replace; boundary=frame') self.end_headers() try: while True: with streaming_lock: frame_to_stream = streaming_frame.copy() if streaming_frame is not None else None if frame_to_stream is None: time.sleep(0.05) continue h, w = frame_to_stream.shape[:2] if w > 960: frame_to_stream = cv2.resize(frame_to_stream, (960, int(h * 960 / w))) ret, jpeg = cv2.imencode('.jpg', frame_to_stream, [cv2.IMWRITE_JPEG_QUALITY, 75]) if not ret: time.sleep(0.05) continue frame_bytes = jpeg.tobytes() self.wfile.write(b'--frame\r\n') self.send_header('Content-Type', 'image/jpeg') self.send_header('Content-Length', len(frame_bytes)) self.end_headers() self.wfile.write(frame_bytes) self.wfile.write(b'\r\n') time.sleep(0.04) # ~25 FPS except Exception: pass else: self.send_error(404, "Path not found") class ThreadedHTTPServer(ThreadingMixIn, HTTPServer): allow_reuse_address = True daemon_threads = True def start_streaming_server(port=8000): try: server = ThreadedHTTPServer(('0.0.0.0', port), StreamingHandler) server_thread = threading.Thread(target=server.serve_forever, daemon=True) server_thread.start() print(f"[RPO IKI] Live View Server HTTP berjalan di http://0.0.0.0:{port}/") except Exception as e: print(f"[WARNING] Gagal membuka HTTP Streaming Server di port {port}: {e}") def write_live_frame(frame): """Simpan frame preview live ke memori & disk untuk Web Server Port 8000.""" global streaming_frame with streaming_lock: streaming_frame = frame try: if os.path.exists("/dev/shm"): os.makedirs("/dev/shm/jetson-counter", exist_ok=True) cv2.imwrite(SHM_LIVE_FRAME_PATH, frame, [cv2.IMWRITE_JPEG_QUALITY, 80]) os.makedirs(os.path.dirname(LIVE_STREAM_FRAME_PATH), exist_ok=True) cv2.imwrite(LIVE_STREAM_FRAME_PATH, frame, [cv2.IMWRITE_JPEG_QUALITY, 80]) except Exception: pass def save_active_batch_state(count=0, status="COUNTING_SACKS", truck_state="LOCKED"): try: os.makedirs(os.path.dirname(STATE_FILE), exist_ok=True) state_data = { "batch_number": "BATCH-RPO-01", "start_time": datetime.now().isoformat(), "count": count, "status": status, "truck_state": truck_state, "last_detection_time": datetime.now().isoformat(), "engine": "rpo_iki" } with open(STATE_FILE, 'w') as f: json.dump(state_data, f, indent=2) except Exception: pass def resolve_live_boundary(width=960, height=540): """Muat koordinat zona dari zones.json (Web Dashboard) atau configs/area_truk.json.""" # 1. Cek zones.json dari Web Dashboard zones_json = Path(_DEFAULT_DIR) / "zones.json" if not zones_json.exists(): zones_json = Path(__file__).parent.parent / "zones.json" if zones_json.exists(): try: with open(zones_json, 'r') as f: zdata = json.load(f) if 'truck' in zdata and len(zdata['truck']) >= 3: pts = np.array(zdata['truck'], dtype=np.float32) scale_x = width / 1920.0 scale_y = height / 1080.0 x_min = int(np.min(pts[:, 0]) * scale_x) x_max = int(np.max(pts[:, 0]) * scale_x) y_min = int(np.min(pts[:, 1]) * scale_y) y_max = int(np.max(pts[:, 1]) * scale_y) line_y = int(y_max - 5) print(f"[RPO IKI] Memuat Zona Dinamis dari Web (zones.json): x={x_min}..{x_max}, y={y_min}..{y_max}, line_y={line_y}") return BoundarySettings( line_pos=line_y, orientation="horizontal", direction="bottom_to_top", zone_x_min=x_min, zone_x_max=x_max, zone_y_min=y_min, zone_y_max=y_max ) except Exception as e: print(f"[WARNING] Gagal membaca zones.json web: {e}") # 2. Fallback ke configs/area_truk.json config_file = Path(__file__).parent / "configs" / "area_truk.json" if config_file.exists(): try: cdata = json.loads(config_file.read_text(encoding="utf-8")) vw = cdata.get("video_width", 1920) vh = cdata.get("video_height", 1080) scale_x = width / float(vw) scale_y = height / float(vh) x_min = int(cdata.get("zone_x_min", 858) * scale_x) x_max = int(cdata.get("zone_x_max", 1231) * scale_x) y_min = int(cdata.get("zone_y_min", 7) * scale_y) y_max = int(cdata.get("zone_y_max", 581) * scale_y) line_y = int(cdata.get("line_y", 580) * scale_y) print(f"[RPO IKI] Memuat & Rescale boundary dari configs/area_truk.json: x={x_min}..{x_max}, y={y_min}..{y_max}, line_y={line_y} (scale={scale_x:.2f})") return BoundarySettings( line_pos=line_y, orientation=cdata.get("orientation", "horizontal"), direction=cdata.get("direction", "bottom_to_top"), zone_x_min=x_min, zone_x_max=x_max, zone_y_min=y_min, zone_y_max=y_max ) except Exception as e: print(f"[WARNING] Gagal membaca configs/area_truk.json: {e}") return BoundarySettings( line_pos=290 if height <= 600 else 580, orientation="horizontal", direction="bottom_to_top", zone_x_min=429 if width <= 960 else 858, zone_x_max=615 if width <= 960 else 1231, zone_y_min=4 if height <= 600 else 7, zone_y_max=290 if height <= 600 else 581 ) def draw_friend_hud_overlay(frame, count, truck_state, current_fps): """Visualisasi HUD Glassmorphic persis buatan rpo iki (step02_count_live.py).""" scale = frame.shape[1] / 1280.0 card_w = int(430 * scale) card_h = int(210 * scale) cx1, cy1 = int(20 * scale), int(20 * scale) cx2, cy2 = cx1 + card_w, cy1 + card_h overlay = frame.copy() cv2.rectangle(overlay, (cx1, cy1), (cx2, cy2), (20, 24, 33), -1) cv2.addWeighted(overlay, 0.72, frame, 0.28, 0, frame) accent_w = int(6 * scale) accent_color = (0, 180, 255) # Orange default if truck_state == "LOCKED": accent_color = (16, 185, 129) # Emerald Green elif truck_state == "WAITING": accent_color = (59, 130, 246) # Blue cv2.rectangle(frame, (cx1, cy1), (cx1 + accent_w, cy2), accent_color, -1) cv2.rectangle(frame, (cx1, cy1), (cx2, cy2), (64, 74, 95), max(1, int(1 * scale))) tx = cx1 + int(18 * scale) cv2.putText( frame, f"AI COUNTER MONITOR | CC1 | LATCH: {truck_state}", (tx, cy1 + int(24 * scale)), cv2.FONT_HERSHEY_SIMPLEX, 0.44 * scale, (148, 163, 184), max(1, int(1 * scale)), ) cv2.putText( frame, "MUAT TRUK:", (tx, cy1 + int(64 * scale)), cv2.FONT_HERSHEY_SIMPLEX, 0.55 * scale, (226, 232, 240), max(1, int(1 * scale)), ) cv2.putText( frame, f"{count}", (tx + int(130 * scale), cy1 + int(72 * scale)), cv2.FONT_HERSHEY_SIMPLEX, 1.35 * scale, accent_color, max(1, int(3 * scale)), ) cv2.putText( frame, f"Status Truk: {truck_state}", (tx, cy1 + int(108 * scale)), cv2.FONT_HERSHEY_SIMPLEX, 0.56 * scale, (241, 245, 249), max(1, int(1 * scale)), ) cv2.putText( frame, f"FPS: {current_fps:.1f} | Engine: RPO IKI (Friends)", (tx, cy1 + int(192 * scale)), cv2.FONT_HERSHEY_SIMPLEX, 0.44 * scale, (100, 116, 139), max(1, int(1 * scale)), ) def run_rpo_iki_prediction(): print("=" * 60) print("MEMULAI LIVE PREDICTION ENGINE (RPO IKI - FRIENDS ALGORITHM)") print("============================================================") init_db() start_streaming_server(port=8000) # Path model karung model_path = Path(__file__).parent / "DATA" / "models" / "karung-dimuat-seg-200e.pt" if not model_path.exists(): model_path = Path(_DEFAULT_DIR) / "karung-dimuat-detection-di-feedmill-yolo26n-seg-200e.pt" if not model_path.exists(): model_path = Path("karung-dimuat-detection-di-feedmill-yolo26n-seg-200e.pt") print(f"[RPO IKI] Memuat Model YOLO Karung: {model_path}") model = YOLO(str(model_path)) # Path model truk truck_model_path = Path(_DEFAULT_DIR) / "truck-detector.pt" if not truck_model_path.exists(): truck_model_path = Path("/home/jetson/karung/truck-detector.pt") if not truck_model_path.exists(): truck_model_path = Path(__file__).parent.parent / "truck-detector.pt" if not truck_model_path.exists(): truck_model_path = Path("truck-detector.pt") truck_model = None truck_class_id = 0 if truck_model_path.exists(): print(f"[RPO IKI AI Latch] Memuat model detektor truk: {truck_model_path}") truck_model = YOLO(str(truck_model_path)) truck_class_id = 0 if "truck-detector" in str(truck_model_path) else 7 else: print(f"[RPO IKI] Model detektor truk tidak ditemukan, menggunakan mode ROI Statis Locked.") # Tentukan device inferensi import torch device = 'cuda' if torch.cuda.is_available() else 'cpu' print(f"[RPO IKI] Device inferensi diset ke: {device}") # Source RTSP source_path = os.getenv("RTSP_URL", "rtsp://admin:K0l0r4n123@10.38.250.21/cam/realmonitor?channel=1&subtype=0") is_stream = any(str(source_path).startswith(p) for p in ["rtsp://", "rtmp://", "http://", "https://"]) if is_stream: print(f"[RPO IKI] Membuka RTSP Stream via Threaded Bufferless Reader: {source_path}") cap = RTSPBufferlessCapture(source_path) else: print(f"[RPO IKI] Membuka File Video Lokal: {source_path}") cap = cv2.VideoCapture(source_path) if not cap.isOpened(): print(f"[ERROR] Gagal membuka stream / video: {source_path}") return raw_width = int(cap.get(cv2.CAP_PROP_FRAME_WIDTH)) or 1920 raw_height = int(cap.get(cv2.CAP_PROP_FRAME_HEIGHT)) or 1080 fps = cap.get(cv2.CAP_PROP_FPS) or 25.0 is_1080p = (raw_width == 1920 and raw_height == 1080) width = 960 if is_1080p else raw_width height = 540 if is_1080p else raw_height boundary = resolve_live_boundary(width, height) params = load_counting_params() counter = LineCounter( line_pos=boundary.line_pos, orientation=boundary.orientation, direction=boundary.direction, min_approach_depth=params.get("min_approach_depth", 15.0), same_sack_radius=params.get("same_sack_radius", 70.0), stack_sack_radius=params.get("stack_sack_radius", 58.0), outside_confirm_frames=params.get("outside_confirm_frames", 2), count_cooldown_dist=params.get("count_cooldown_dist", 85.0), count_cooldown_frames=params.get("count_cooldown_frames", 40), staging_cooldown_frames=params.get("staging_cooldown_frames", 15), min_staging_depth=params.get("min_staging_depth", 40.0), min_track_frames=params.get("min_track_frames", 8), ghost_track_frames=params.get("ghost_track_frames", 0), clip_warmup_frames=params.get("clip_warmup_frames", 25), zone_x_min=boundary.zone_x_min, zone_x_max=boundary.zone_x_max, zone_y_min=boundary.zone_y_min, zone_y_max=boundary.zone_y_max, ) # State machine truk truck_state = "LOCKED" # Default locked agar langsung menghitung karung consecutive_truck_detections = 0 locked_bbox = None flashes: list[CountFlash] = [] frame_idx = 0 last_time = time.time() current_fps = 0.0 print(f"[RPO IKI] Counter Siap! Resized Frame: {width}x{height}, Line: {boundary.orientation} @ {counter._line_pos}, Zone: x={counter.zone_x_min}..{counter.zone_x_max}, y={counter.zone_y_min}..{counter.zone_y_max}") config_file = Path(__file__).parent / "configs" / "area_truk.json" last_config_mtime = config_file.stat().st_mtime if config_file.exists() else 0.0 while cap.isOpened(): ret, frame = cap.read() if not ret or frame is None: time.sleep(0.01) continue frame_idx += 1 # Hot-reload konfigurasi jika configs/area_truk.json diperbarui di disk / web if frame_idx % 25 == 0 and config_file.exists(): try: mtime = config_file.stat().st_mtime if mtime > last_config_mtime: last_config_mtime = mtime boundary = resolve_live_boundary(width, height) counter._line_pos = boundary.line_pos counter.zone_x_min = boundary.zone_x_min counter.zone_x_max = boundary.zone_x_max counter.zone_y_min = boundary.zone_y_min counter.zone_y_max = boundary.zone_y_max print(f"[Live Config Reload] Boundary diperbarui secara dinamis: x={boundary.zone_x_min}..{boundary.zone_x_max}, line_y={boundary.line_pos}") except Exception: pass # Resize ke 960x540 jika 1080p agar presisi dengan area_truk.json rpo iki if is_1080p and frame.shape[1] == 1920 and frame.shape[0] == 1080: frame = cv2.resize(frame, (960, 540)) annotated = frame.copy() # 1. Dynamic Truk Latch State Machine (Jika truck_model aktif) if truck_model is not None and truck_state == "WAITING": if frame_idx % 10 == 0: results_t = truck_model(frame, classes=[truck_class_id], conf=0.40, verbose=False) boxes_t = results_t[0].boxes if len(boxes_t) > 0: sorted_boxes = sorted(boxes_t, key=lambda b: (b.xyxy[0][2] - b.xyxy[0][0]) * (b.xyxy[0][3] - b.xyxy[0][1]), reverse=True) t_box = sorted_boxes[0].xyxy[0].cpu().numpy() tx1, ty1, tx2, ty2 = map(int, t_box) consecutive_truck_detections += 1 if consecutive_truck_detections >= 3: locked_bbox = (tx1, ty1, tx2, ty2) truck_state = "LOCKED" consecutive_truck_detections = 0 line_y = int(ty2 - 5) boundary = replace( boundary, line_pos=line_y, zone_x_min=tx1, zone_x_max=tx2, zone_y_min=ty1, zone_y_max=ty2, ) counter._line_pos = boundary.line_pos counter.zone_x_min = boundary.zone_x_min counter.zone_x_max = boundary.zone_x_max counter.zone_y_min = boundary.zone_y_min counter.zone_y_max = boundary.zone_y_max print(f"[AI Latch] Truk Terdeteksi Stabil! Mengunci ROI Bak Truk: x={tx1}..{tx2}, y={ty1}..{ty2}, line_y={line_y}") # 2. Tracking Karung & Counting (Saat truck_state == "LOCKED") boxes = None masks = None if truck_state in ("LOCKED", "DEPARTING"): results = model.track( frame, persist=True, classes=[SACK_CLASS_ID], conf=params.get("conf", 0.15), tracker="bytetrack.yaml", device=device, verbose=False ) boxes = results[0].boxes masks = results[0].masks if boxes is not None and boxes.id is not None: for box_coord, track_id in zip(boxes.xyxy.cpu().numpy(), boxes.id.int().cpu().tolist()): cross, foot = track_points(box_coord) count_event = counter.update(track_id, cross, frame_idx, foot) if count_event is not None: flashes.append(CountFlash(number=counter.count, x=int(cross[0]), y=int(cross[1]), frames_left=int(fps * 0.35))) print(f"[RPO IKI COUNTER] Karung #{track_id} TERHITUNG! Total: {counter.count}") save_active_batch_state(count=counter.count, truck_state=truck_state) # 3. Render Visualisasi Asli rpo iki (Garis Hijau Batas, Kotak Truk Orange, HUD) draw_blind_truck_overlay(annotated, boundary, boundary.line_pos) draw_boundary( annotated, boundary.line_pos, boundary.orientation, boundary.zone_x_min, boundary.zone_x_max, boundary.zone_y_min, boundary.zone_y_max, ) # Gambar BBox Karung yang sedang mendekati garis if boxes is not None and boxes.id is not None: annotate_tracks( annotated, boxes, counter, flashes, fps, frame_idx, blind_truck=True, masks=masks, show_mask=True ) tick_flashes(flashes) draw_count_flashes(annotated, flashes, fps) draw_friend_hud_overlay(annotated, counter.count, truck_state, current_fps) # Hitung FPS if frame_idx % 25 == 0: now = time.time() elapsed = now - last_time if elapsed > 0: current_fps = 25.0 / elapsed last_time = now print(f"[INFO] Frame {frame_idx} - State: {truck_state} - Terhitung: {counter.count} karung ({current_fps:.2f} FPS)") save_active_batch_state(count=counter.count, truck_state=truck_state) # Tulis live frame preview untuk Web Server Port 8000 write_live_frame(annotated) if __name__ == "__main__": try: run_rpo_iki_prediction() except KeyboardInterrupt: print("\n[RPO IKI] Program dihentikan secara manual (Ctrl+C).")