Files
karung-counting-feedmill-se…/rpo_iki/predict_rpo_iki.py
T

626 lines
24 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 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).")