""" Batch Video Cropper — Production 24/7 Rekam video RTSP per sesi batch truk. Ringan, tanpa GUI, auto-reconnect. Output: ~/reTraining/data/archive/{YYYY-MM-DD}/batch_{N}_{HH-MM-SS}.mp4 Menjalankan: cd ~/reTraining/algoritma-batch python3 batch_video_cropper.py """ import os # KRITIS: Konfigurasi RTSP transport — HARUS sebelum import cv2 # Tanpa ini, OpenCV pakai UDP (default) yang sering drop koneksi os.environ["OPENCV_FFMPEG_CAPTURE_OPTIONS"] = ( "rtsp_transport;tcp|buffer_size;20480000|max_delay;500000|reorder_queue_size;500" ) import json import shutil import signal import sys import urllib.parse import urllib.request import cv2 import numpy as np import time import threading import platform from datetime import datetime, timedelta from shapely.geometry import Point, Polygon from ultralytics import YOLO from src.tracking import ByteTrackTracker from src.stabilizer import BboxStabilizer from src.truck_roi import TruckROI from src.counting import LineCrossCounter from src.batch import BatchLifecycleManager, BatchState # ===================================================================== # KONFIGURASI # ===================================================================== IS_WINDOWS = platform.system() == "Windows" BASE_DIR = os.path.dirname(os.path.abspath(__file__)) if IS_WINDOWS: MODEL_PATH = os.path.join(BASE_DIR, "v3-best.pt") ARCHIVE_BASE = os.path.join(BASE_DIR, "archive_output") RTSP_URL = "video truk.mp4" # Testing lokal video else: MODEL_PATH = os.path.join(BASE_DIR, "v3-best.pt") ARCHIVE_BASE = os.path.expanduser("~/reTraining/data/archive") RTSP_URL = "rtsp://192.168.192.96:8554/cam" # Production RTSP stream (.105) DAILY_CUTOFF_TIME = "00:00" # State Machine SACK_IDLE_TIMEOUT = 5.0 MIN_BATCH_DURATION = 2.0 TRUCK_GONE_TOLERANCE = 5.0 # Pengambilan video dari MediaMTX (REQ-170) # Jetson merekam 24/7 apa adanya; skrip ini hanya menentukan potongannya. # Sebelumnya frame di-encode ulang ke mpeg4 di sini: 4,7x lebih besar dari # sumbernya, kualitas turun, dan fps-nya salah. Sekarang potongan diunduh # sebagai salinan, jadi codec, fps dan waktunya persis seperti kamera. PLAYBACK_URL = os.getenv("PLAYBACK_URL", "http://192.168.192.96:9996/get") PLAYBACK_PATH = os.getenv("PLAYBACK_PATH", "cam") FETCH_PAD_BEFORE = 3.0 # detik diambil sebelum truk terdeteksi FETCH_PAD_AFTER = 3.0 # dan sesudahnya, supaya tidak terpotong FETCH_RETRIES = 3 FETCH_RETRY_DELAY = 20.0 # Video Recording # Diambil dari stream yang diterima, bukan ditebak. Angka 10.0 yang dulu # di-hardcode membuat SETIAP file di arsip punya timebase salah: kamera # mengirim 25 fps, file mengaku 10 fps, jadi rekaman 15,4 menit tersimpan # sebagai 38,3 menit dan diputar 2,49x lebih lambat dari kenyataan. # Dipakai hanya kalau fps stream tidak terbaca. FALLBACK_FPS = 25.0 MIN_FPS, MAX_FPS = 1.0, 60.0 VIDEO_CODEC = "mp4v" # Reconnect RECONNECT_DELAY = 5 # Detik menunggu sebelum reconnect RTSP MAX_EMPTY_FRAMES = 300 # Maks frame kosong sebelum reconnect (~30 detik) # Matikan tampilan visualisasi agar program sangat ringan 24/7 SHOW_DISPLAY = False # ===================================================================== # THREADED RTSP READER (selalu ambil frame terbaru, anti-lag) # ===================================================================== class RTSPReader: def __init__(self, url): self.url = url self.fps = FALLBACK_FPS self.cap = None self.frame = None self.ret = False self.running = True self.lock = threading.Lock() self.event = threading.Event() self._connect() self.thread = threading.Thread(target=self._loop, daemon=True) self.thread.start() def _connect(self): if self.cap and self.cap.isOpened(): self.cap.release() self.cap = cv2.VideoCapture(self.url) if self.cap.isOpened(): self.fps = self._read_fps() log(f"RTSP terhubung ({self.fps:.1f} fps)") else: log("RTSP gagal terhubung") def _read_fps(self): """Fps yang diumumkan stream, dibatasi ke rentang masuk akal.""" try: reported = float(self.cap.get(cv2.CAP_PROP_FPS) or 0.0) except Exception: reported = 0.0 if MIN_FPS <= reported <= MAX_FPS: return reported log(f"Fps stream tidak masuk akal ({reported}), pakai {FALLBACK_FPS}") return FALLBACK_FPS def _loop(self): empty = 0 while self.running: if not self.cap or not self.cap.isOpened(): log(f"RTSP terputus, reconnect dalam {RECONNECT_DELAY}s...") time.sleep(RECONNECT_DELAY) self._connect() empty = 0 continue ret, frame = self.cap.read() if not ret: empty += 1 if empty > MAX_EMPTY_FRAMES: log(f"RTSP {empty} frame kosong, reconnect...") self._connect() empty = 0 time.sleep(0.01) continue empty = 0 with self.lock: self.ret, self.frame = ret, frame self.event.set() time.sleep(0.001) def read(self): if self.event.wait(timeout=2.0): self.event.clear() with self.lock: return self.ret, self.frame.copy() if self.frame is not None else (False, None) return False, None def release(self): self.running = False if self.cap: self.cap.release() # ===================================================================== # UTILITAS # ===================================================================== def log(msg): ts = datetime.now().strftime("%Y-%m-%d %H:%M:%S") print(f"[{ts}] {msg}", flush=True) def get_counting_date(): 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().isoformat() return (dt.date() + timedelta(days=1)).isoformat() def get_next_batch_number_from_files(counting_date): folder = os.path.join(ARCHIVE_BASE, counting_date) if not os.path.exists(folder): return 1 max_num = 0 try: import re pattern = re.compile(r'^batch_?(\d+)\.mp4$', re.IGNORECASE) for filename in os.listdir(folder): match = pattern.match(filename) if match: num = int(match.group(1)) max_num = max(max_num, num) except Exception as e: log(f"Error scanning folder for batch files: {e}") return max_num + 1 # ===================================================================== # VIDEO RECORDER # ===================================================================== class SessionFetcher: """Tandai kapan satu sesi truk mulai dan selesai, lalu unduh potongannya. Mengunduh dilakukan di thread terpisah supaya loop deteksi tidak berhenti menunggu jaringan — satu sesi 40 menit bisa ratusan MB. Kalau gagal, dicoba lagi; buffer di Jetson menyimpan 24 jam, jadi ada banyak waktu untuk pulih. """ def __init__(self): self.batch_num = None self.counting_date = None self.started_at = None self.path = None def start(self, batch_num, counting_date, w=1280, h=720, fps=None): self.batch_num = batch_num self.counting_date = counting_date self.started_at = datetime.now() folder = os.path.join(ARCHIVE_BASE, counting_date) os.makedirs(folder, exist_ok=True) self.path = os.path.join(folder, f"batch{batch_num:03d}.mp4") log(f"REC MARK START -> {self.path} @ {self.started_at:%H:%M:%S}") def write(self, frame): """Tidak ada yang ditulis per frame lagi — Jetson yang merekam.""" def stop(self, discard=False): if self.started_at is None: return started, path = self.started_at, self.path ended = datetime.now() self.started_at = self.path = None if discard: log(f"REC DISCARD -> {path} tidak diunduh (batch tidak valid)") return threading.Thread(target=self._fetch, args=(path, started, ended), daemon=True).start() def _fetch(self, path, started, ended): begin = started - timedelta(seconds=FETCH_PAD_BEFORE) duration = (ended - started).total_seconds() + FETCH_PAD_BEFORE + FETCH_PAD_AFTER # '+' pada offset zona waktu wajib di-encode; kalau tidak, ia terbaca # sebagai spasi dan MediaMTX menolak dengan "invalid start". start_param = urllib.parse.quote(begin.astimezone().isoformat(timespec="seconds"), safe="") url = (f"{PLAYBACK_URL}?path={PLAYBACK_PATH}&start={start_param}" f"&duration={duration:.0f}&format=mp4") for attempt in range(1, FETCH_RETRIES + 1): try: tmp = f"{path}.part" with urllib.request.urlopen(url, timeout=600) as response: if response.status != 200: raise IOError(f"HTTP {response.status}") with open(tmp, "wb") as handle: shutil.copyfileobj(response, handle) size = os.path.getsize(tmp) if size < 1024: raise IOError(f"hasil terlalu kecil ({size} byte)") os.replace(tmp, path) _write_sidecar(path, begin, duration) log(f"REC FETCHED -> {path} ({size/1e6:.0f} MB, {duration:.0f} detik)") return except Exception as exc: log(f"REC FETCH gagal ({attempt}/{FETCH_RETRIES}) {path}: {exc}") try: os.remove(f"{path}.part") except OSError: pass if attempt < FETCH_RETRIES: time.sleep(FETCH_RETRY_DELAY) log(f"REC FETCH MENYERAH -> {path}. Rekaman masih ada di buffer Jetson " f"selama 24 jam sejak {begin:%Y-%m-%d %H:%M:%S}") @property def is_recording(self): return self.started_at is not None def _write_sidecar(video_path, begin, duration): """Waktu sebenarnya, di sebelah videonya. Aplikasi tidak perlu lagi membaca jam dari overlay dengan OCR untuk file baru: waktunya datang dari server rekaman, tepat sampai detik. """ payload = { "started_at": begin.strftime("%Y-%m-%d %H:%M:%S"), "duration_seconds": round(duration, 1), "source": "mediamtx-playback", "written_at": datetime.now().strftime("%Y-%m-%d %H:%M:%S"), } sidecar = os.path.splitext(video_path)[0] + ".json" try: with open(sidecar, "w", encoding="utf-8") as handle: json.dump(payload, handle) except OSError as exc: log(f"Gagal menulis sidecar {sidecar}: {exc}") # ===================================================================== # MAIN LOOP # ===================================================================== shutdown_flag = False def handle_signal(sig, _): global shutdown_flag log(f"Signal {sig} diterima, menutup program...") shutdown_flag = True signal.signal(signal.SIGINT, handle_signal) signal.signal(signal.SIGTERM, handle_signal) def run(): global shutdown_flag log("=" * 50) log("BATCH VIDEO CROPPER — Production 24/7") log(f"Model : {MODEL_PATH}") log(f"RTSP : {RTSP_URL}") log(f"Archive : {ARCHIVE_BASE}") log(f"Toleransi batch: {TRUCK_GONE_TOLERANCE}s (truk+karung)") log("=" * 50) os.makedirs(ARCHIVE_BASE, exist_ok=True) # Load model log("Memuat model YOLO...") model = YOLO(MODEL_PATH) # Detect device device = "cpu" try: import torch if torch.cuda.is_available(): device = "cuda" except ImportError: pass # Warm-up dummy = np.zeros((720, 1280, 3), dtype=np.uint8) _ = model(dummy, imgsz=640, device=device, verbose=False) log(f"Model siap. Device: {device}") # Components tracker = ByteTrackTracker(model, conf=0.25) stabilizer = BboxStabilizer(ema_alpha=0.35, max_hold_frames=10, max_height_ratio=1.5, min_height_ratio=0.70) # Koordinat zona (1920x1080 → 1280x720) sx, sy = 1280.0 / 1920.0, 720.0 / 1080.0 detection_polygon = Polygon([ [int(574*sx), int(50*sy)], [int(586*sx), int(1077*sy)], [int(1418*sx), int(1076*sy)], [int(1397*sx), int(50*sy)], ]) truck_polygon = Polygon([ [int(600*sx), int(385*sy)], [int(609*sx), int(1076*sy)], [int(1404*sx), int(1078*sy)], [int(1381*sx), int(343*sy)], ]) line_y = int(330 * sy) line_x1 = int(577 * sx) line_x2 = int(1401 * sx) static_roi = TruckROI( x1=int(600*sx), y1=int(343*sy), x2=int(1404*sx), y2=int(1078*sy), line_y=line_y, confidence=1.0, ) counter = LineCrossCounter(line_y=line_y, line_x_start=line_x1, line_x_end=line_x2, margin=20, dedup_radius=60.0) batch_mgr = BatchLifecycleManager( stabilize_seconds=0.0, stabilize_threshold_px=9999.0, sack_idle_timeout=SACK_IDLE_TIMEOUT, min_batch_duration=MIN_BATCH_DURATION, truck_gone_tolerance=3.0, ) recorder = SessionFetcher() # RTSP Stream log(f"Membuka RTSP: {RTSP_URL}") cap = RTSPReader(RTSP_URL) batch_counter = 0 frame_idx = 0 last_status_time = time.time() truck_seen_in_current_batch = False last_frame_time = time.time() NO_FRAME_BATCH_TIMEOUT = 30.0 # Akhiri batch jika tidak ada frame 30 detik log("Loop utama dimulai...") try: while not shutdown_flag: ret, frame = cap.read() if not ret or frame is None: # Saat tidak ada frame DAN batch aktif, cek timeout if batch_mgr.is_active: no_frame_duration = time.time() - last_frame_time if no_frame_duration >= NO_FRAME_BATCH_TIMEOUT: log(f"RTSP drop {no_frame_duration:.0f}s. Force-end BATCH #{batch_counter}. Karung: {counter.loading_count}") # Force-end: langsung reset state machine (bypass update_truck) batch_mgr._state = BatchState.IDLE batch_mgr._current_batch_id = None batch_mgr._truck_is_stable = False recorder.stop(discard=True) counter.reset() stabilizer.reset() last_frame_time = time.time() # Reset timer agar tidak spam time.sleep(0.01) continue last_frame_time = time.time() frame = cv2.resize(frame, (1280, 720)) timestamp = time.time() frame_idx += 1 prev_active = batch_mgr.is_active prev_state = batch_mgr.state # --- Deteksi --- raw_all = tracker.update(frame, []) # Hanya proses objek yang pusatnya berada di dalam area deteksi (poligon ungu) raw_all_filtered = [ d for d in raw_all if detection_polygon.contains(Point((d.bbox[0] + d.bbox[2]) / 2.0, (d.bbox[1] + d.bbox[3]) / 2.0)) ] sacks = [d for d in raw_all_filtered if d.class_name == "sack"] trucks = [d for d in raw_all_filtered if d.class_name == "truck"] if batch_mgr.is_active and len(trucks) > 0: truck_seen_in_current_batch = True # Stabilizer stable = stabilizer.update(sacks) # Karung di 70% area truk ty_min, ty_max = truck_polygon.bounds[1], truck_polygon.bounds[3] cutoff_y = ty_min + 0.30 * (ty_max - ty_min) sacks_in_area = sum( 1 for d in stable if truck_polygon.contains(Point((d.bbox[0]+d.bbox[2])/2, (d.bbox[1]+d.bbox[3])/2)) and (d.bbox[1]+d.bbox[3])/2 >= cutoff_y ) # Line crossing in_roi = [d for d in stable if static_roi.contains_x((d.bbox[0]+d.bbox[2])/2)] events = counter.update(in_roi) has_crossing = len(events) > 0 # --- State Machine --- 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.update_sacks( has_crossing_event=has_crossing, sacks_in_area_count=sacks_in_area, timestamp=timestamp, loading_count=counter.loading_count, unloading_count=counter.unloading_count, ) if batch_mgr.state == "WAITING_FOR_ACTIVITY": # Sinkronkan _truck_last_seen agar countdown toleransi # mulai dari saat WAITING dimulai, bukan dari TRUCK_STABILIZING if batch_mgr._truck_last_seen < batch_mgr._waiting_since: batch_mgr._truck_last_seen = batch_mgr._waiting_since anything = (sacks_in_area > 0) or (len(trucks) > 0) batch_mgr._truck_gone_tolerance = TRUCK_GONE_TOLERANCE batch_mgr.update_truck(anything, None, timestamp) # --- Transisi Batch --- if batch_mgr.is_active and not prev_active: cd = get_counting_date() batch_counter = get_next_batch_number_from_files(cd) truck_seen_in_current_batch = False log(f"BATCH #{batch_counter} DIMULAI (tanggal: {cd})") recorder.start(batch_counter, cd, fps=cap.fps) elif not batch_mgr.is_active and prev_active: final_count = counter.loading_count # Tentukan apakah batch valid (truk harus terdeteksi minimal sekali DAN hitungan karung > 0) is_valid = (final_count > 0) and truck_seen_in_current_batch if is_valid: log(f"BATCH #{batch_counter} SELESAI. Karung: {final_count}") recorder.stop(discard=False) else: log(f"BATCH #{batch_counter} DIABAIKAN (Karung={final_count}, Truk Terdeteksi={truck_seen_in_current_batch})") recorder.stop(discard=True) # Kembalikan nomor counter batch karena batch ini dianulir batch_counter = max(0, batch_counter - 1) counter.reset() stabilizer.reset() if batch_mgr.state != prev_state: log(f"STATE: {prev_state} -> {batch_mgr.state}") # Tulis frame ke video if batch_mgr.is_active: recorder.write(frame) # Log karung crossing for ev in events: log(f"KARUNG #{ev['track_id']} crossing. Total: {counter.loading_count}") # Status log setiap 5 menit if timestamp - last_status_time >= 300: last_status_time = timestamp rec = "REC" if recorder.is_recording else "---" log(f"STATUS: frames={frame_idx} batches={batch_counter} " f"state={batch_mgr.state} {rec}") # --- Visualisasi Live Predict (Lokal Windows saja) --- if SHOW_DISPLAY: display = frame.copy() # Gambar detection polygon (magenta) det_pts = np.array([ [int(574*sx), int(50*sy)], [int(586*sx), int(1077*sy)], [int(1418*sx), int(1076*sy)], [int(1397*sx), int(50*sy)] ], dtype=np.int32) cv2.polylines(display, [det_pts], True, (255, 0, 255), 2) # Gambar truck polygon (orange) trk_pts = np.array([ [int(600*sx), int(385*sy)], [int(609*sx), int(1076*sy)], [int(1404*sx), int(1078*sy)], [int(1381*sx), int(343*sy)] ], dtype=np.int32) cv2.polylines(display, [trk_pts], True, (0, 165, 255), 2) # Gambar count line (cyan) cv2.line(display, (line_x1, line_y), (line_x2, line_y), (255, 255, 0), 2) cv2.putText(display, "COUNTING LINE", (line_x1 + 10, line_y - 8), cv2.FONT_HERSHEY_SIMPLEX, 0.5, (255, 255, 0), 1) # Gambar bbox TRUCK (hijau) for d in trucks: x1, y1, x2, y2 = [int(v) for v in d.bbox] cv2.rectangle(display, (x1, y1), (x2, y2), (0, 200, 0), 2) lbl = f"truck #{d.track_id} ({d.confidence:.2f})" cv2.putText(display, lbl, (x1, y1 - 5), cv2.FONT_HERSHEY_SIMPLEX, 0.5, (0, 200, 0), 1) # Gambar bbox SACK (cyan untuk stable) for d in stable: x1, y1, x2, y2 = [int(v) for v in d.bbox] cv2.rectangle(display, (x1, y1), (x2, y2), (255, 255, 0), 2) lbl = f"sack #{d.track_id}" cv2.putText(display, lbl, (x1, y2 + 15), cv2.FONT_HERSHEY_SIMPLEX, 0.4, (255, 255, 0), 1) # Overlay HUD overlay = display.copy() cv2.rectangle(overlay, (5, 5), (380, 150), (0, 0, 0), -1) cv2.addWeighted(overlay, 0.65, display, 0.35, 0, display) state_str = batch_mgr.state rec_str = "● RECORDING" if recorder.is_recording else "○ IDLE" color_state = (0, 250, 0) if batch_mgr.is_active else (0, 165, 255) cv2.putText(display, f"State: {state_str}", (15, 30), cv2.FONT_HERSHEY_SIMPLEX, 0.6, color_state, 2) cv2.putText(display, f"Total Counted: {counter.loading_count}", (15, 55), cv2.FONT_HERSHEY_SIMPLEX, 0.5, (0, 255, 255), 1) cv2.putText(display, f"Sacks in Area: {sacks_in_area}", (15, 80), cv2.FONT_HERSHEY_SIMPLEX, 0.5, (0, 255, 255), 1) cv2.putText(display, f"Current Batch: #{batch_counter} ({rec_str})", (15, 105), cv2.FONT_HERSHEY_SIMPLEX, 0.5, (255, 255, 255), 1) cv2.putText(display, f"Sacks: {len(stable)} | Trucks: {len(trucks)}", (15, 130), cv2.FONT_HERSHEY_SIMPLEX, 0.5, (200, 200, 200), 1) # Resize agar muat layar lokal resized = cv2.resize(display, (960, 540)) cv2.imshow("Batch Video Cropper - Local Predict", resized) if cv2.waitKey(1) & 0xFF == ord('q'): log("Dihentikan secara manual melalui tombol 'q'") break except Exception as e: log(f"ERROR: {e}") import traceback traceback.print_exc() finally: if recorder.is_recording: log("Menyimpan rekaman batch terakhir...") recorder.stop() cap.release() if SHOW_DISPLAY: cv2.destroyAllWindows() log(f"SELESAI. Total batch: {batch_counter}, frames: {frame_idx}") if __name__ == "__main__": run()