"""
성림 라이브 제어기 - 웹 서버 (FastAPI + WebSocket)
브라우저에서 카메라 스트리밍 및 PTZ 제어를 위한 백엔드 서버입니다.
"""

import asyncio
import base64
import json
import threading
import time
import os
from contextlib import asynccontextmanager

import cv2
import requests
from requests.auth import HTTPBasicAuth
from requests.exceptions import RequestException

from fastapi import FastAPI, WebSocket, WebSocketDisconnect
from fastapi.responses import JSONResponse, FileResponse
import uvicorn

# ─────────────────── 카메라 설정 ───────────────────
CAMERA_CONFIGS = {
    1: {
        "name": "예배당 좌측 (1번)",
        "rtsp_main": "rtsp://mev.o-r.kr:20001/stream1",
        "ctrl_ip": "mev.o-r.kr",
        "ctrl_port": 20004,
    },
    2: {
        "name": "예배당 중앙 (2번)",
        "rtsp_main": "rtsp://mev.o-r.kr:20002/stream1",
        "ctrl_ip": "mev.o-r.kr",
        "ctrl_port": 20005,
    },
    3: {
        "name": "예배당 우측 (3번)",
        "rtsp_main": "rtsp://mev.o-r.kr:20003/stream1",
        "ctrl_ip": "mev.o-r.kr",
        "ctrl_port": 20006,
    },
}

INTERNAL_CONFIGS = {
    1: {
        "rtsp_main": "rtsp://192.168.0.91:554/stream1",
        "ctrl_ip": "192.168.0.91",
        "ctrl_port": 80
    },
    2: {
        "rtsp_main": "rtsp://192.168.0.92:554/stream1",
        "ctrl_ip": "192.168.0.92",
        "ctrl_port": 80
    },
    3: {
        "rtsp_main": "rtsp://192.168.0.90:554/stream1",
        "ctrl_ip": "192.168.0.90",
        "ctrl_port": 80
    },
}

# ─────────────────── 스트림 관리 ───────────────────

class StreamManager:
    """
    RTSP 캡처를 백그라운드 스레드에서 유지하고,
    연결된 모든 WebSocket에 JPEG 프레임을 브로드캐스트합니다.
    (단일 스레드 구조로 스레드 누수 원천 차단)
    """

    def __init__(self):
        self.current_camera_id = 3
        self.network_mode = "external"
        self._status = "초기화 중..."
        
        # 💡 카메라 전환 신호용
        self._target_camera_id = 3
        self._target_network_mode = "external"
        self._switch_event = threading.Event()
        
        # 스레드 공유용 버퍼 및 락
        self._latest_raw_frame = None
        self._frame_lock = threading.Lock()
        self._frame_seq = 0
        
        self._subscribers: list[asyncio.Queue] = []
        self._loop = None  # 메인 이벤트 루프 참조 (브로드캐스트용)
        
        self._running = True
        self._capture_thread = threading.Thread(target=self._capture_loop, daemon=True)
        self._broadcast_thread = threading.Thread(target=self._broadcast_loop, daemon=True)

    def start_threads(self):
        """앱 시작 시 딱 한 번만 스레드 구동"""
        if not self._capture_thread.is_alive():
            self._capture_thread.start()
        if not self._broadcast_thread.is_alive():
            self._broadcast_thread.start()

    def set_camera(self, camera_id: int):
        """카메라 전환 요청 (스레드에 신호 전달)"""
        self._target_camera_id = camera_id
        self._switch_event.set()
        
    def set_network_mode(self, mode: str):
        """네트워크 모드 전환 요청"""
        self._target_network_mode = mode
        self._switch_event.set()

    def _clear_all_queues(self):
        """구독 중인 웹소켓 큐들에 남아 있는 전송 대기 프레임들을 전부 비워 섞임 현상 방지"""
        if self._loop is None:
            return
        
        def flush_queues():
            for q in list(self._subscribers):
                try:
                    while not q.empty():
                        q.get_nowait()
                except Exception:
                    pass
                    
        self._loop.call_soon_threadsafe(flush_queues)

    def _get_urls(self):
        if self.network_mode == "internal":
            cfg = INTERNAL_CONFIGS[self.current_camera_id]
        else:
            cfg = CAMERA_CONFIGS[self.current_camera_id]
        return cfg["rtsp_main"]

    def _capture_loop(self):
        is_connected = False
        cap = None
        last_retrieve_time = 0

        if "OPENCV_FFMPEG_CAPTURE_OPTIONS" in os.environ:
            del os.environ["OPENCV_FFMPEG_CAPTURE_OPTIONS"]

        while self._running:
            # 💡 카메라 전환 신호 감지 시
            if self._switch_event.is_set():
                self._switch_event.clear()
                self.current_camera_id = self._target_camera_id
                self.network_mode = self._target_network_mode
                is_connected = False
                
                if cap:
                    try:
                        cap.release()
                    except Exception:
                        pass
                    cap = None
                    
                with self._frame_lock:
                    self._latest_raw_frame = None
                    self._frame_seq = 0
                self._clear_all_queues()
                self._set_status("카메라 변경 중...")

            if not is_connected:
                rtsp_url = self._get_urls()
                self._set_status("카메라 연결 시도 중...")
                
                cap = cv2.VideoCapture(rtsp_url, cv2.CAP_FFMPEG)
                cap.set(cv2.CAP_PROP_BUFFERSIZE, 1) # OpenCV 내장 버퍼 크기 최소화
                
                if cap.isOpened():
                    ret, frame = cap.read()
                    if ret and frame is not None:
                        is_connected = True
                        self._set_status("RTSP 연결 완료")
                        with self._frame_lock:
                            self._latest_raw_frame = frame

                if not is_connected:
                    if cap:
                        cap.release()
                        cap = None
                    self._set_status("카메라 연결 실패 (대기 중...)")
                    # 연결 재시도 전 대기 (이 도중에도 전환 신호가 오면 즉시 깸)
                    for _ in range(20):
                        if self._switch_event.is_set():
                            break
                        time.sleep(0.1)
                    continue

            # 메인 캡처 루프
            try:
                # grab()만 수행하면 압축 해제(디코딩) 없이 네트워크 패킷만 빨리 소모함
                grabbed = cap.grab()
                if grabbed:
                    now = time.time()
                    # 0.1초(10fps) 주기로만 실제 디코딩(retrieve)을 수행하여 라즈베리파이 CPU 부하를 방지하고 지연시간 제거
                    if now - last_retrieve_time >= 0.1:
                        ret_rt, frame = cap.retrieve()
                        if ret_rt and frame is not None:
                            with self._frame_lock:
                                self._latest_raw_frame = frame
                        last_retrieve_time = now
                else:
                    is_connected = False
                    if cap:
                        cap.release()
                        cap = None
                    self._set_status("연결 유실 (재접속 대기)")
                    time.sleep(1)
                
                # time.sleep을 없애서 패킷 소모 속도를 극대화 (grab은 블로킹 함수라 CPU 100% 안침)
                
            except Exception:
                is_connected = False
                if cap:
                    try:
                        cap.release()
                    except Exception:
                        pass
                    cap = None
                self._set_status("캡처 오류 (재접속 대기)")
                time.sleep(1)

        if cap:
            try:
                cap.release()
            except Exception:
                pass

    def _broadcast_loop(self):
        """별도 스레드에서 프레임을 처리하여 전송 대기자에게 전달"""
        while self._running:
            frame = None
            with self._frame_lock:
                if self._latest_raw_frame is not None:
                    frame = self._latest_raw_frame.copy()
            
            if frame is not None:
                self._push_frame(frame)
            
            # 10~12 FPS 수준 유지 (부하 및 지연 최적화)
            time.sleep(0.09)

    def _set_status(self, msg: str):
        self._status = msg
        self._broadcast_status(msg)

    def _push_frame(self, frame):
        """프레임을 고화질(1280x720, Quality 85)로 인코딩 후 브로드캐스트"""
        h, w = frame.shape[:2]
        target_w = 1280
        if w > target_w:
            scale = target_w / w
            frame = cv2.resize(frame, (target_w, int(h * scale)), interpolation=cv2.INTER_AREA)

        ret, buf = cv2.imencode(".jpg", frame, [cv2.IMWRITE_JPEG_QUALITY, 85])
        if not ret:
            return

        jpeg_bytes = buf.tobytes()
        b64 = base64.b64encode(jpeg_bytes).decode("ascii")
        
        self._frame_seq += 1
        msg = json.dumps({
            "type": "frame", 
            "data": b64,
            "seq": self._frame_seq
        })
        self._broadcast(msg)

    def _broadcast(self, msg: str):
        if self._loop is None:
            return
        for q in list(self._subscribers):
            try:
                self._loop.call_soon_threadsafe(q.put_nowait, msg)
            except Exception:
                pass

    def _broadcast_status(self, status: str):
        msg = json.dumps({"type": "status", "status": status, "camera_id": self.current_camera_id})
        self._broadcast(msg)

    def subscribe(self) -> asyncio.Queue:
        q: asyncio.Queue = asyncio.Queue(maxsize=5)
        self._subscribers.append(q)
        return q

    def unsubscribe(self, q: asyncio.Queue):
        try:
            self._subscribers.remove(q)
        except ValueError:
            pass

    @property
    def status(self):
        return self._status


stream_manager = StreamManager()

# ─────────────────── FastAPI 앱 ───────────────────

@asynccontextmanager
async def lifespan(app: FastAPI):
    stream_manager._loop = asyncio.get_event_loop()
    stream_manager.start_threads()
    yield
    stream_manager._running = False


app = FastAPI(title="성림 라이브 제어기 웹", lifespan=lifespan)

WEB_DIR = os.path.join(os.path.dirname(os.path.abspath(__file__)), "web")


@app.get("/")
async def index():
    return FileResponse(os.path.join(WEB_DIR, "index.html"))


@app.get("/api/cameras")
async def get_cameras():
    """카메라 목록 및 현재 상태 반환"""
    result = {}
    configs = INTERNAL_CONFIGS if stream_manager.network_mode == "internal" else CAMERA_CONFIGS
    for cam_id, cfg in configs.items():
        result[cam_id] = {"name": cfg["name"] if "name" in cfg else f"카메라 {cam_id}", "id": cam_id}
    return JSONResponse({
        "cameras": result,
        "current_camera_id": stream_manager.current_camera_id,
        "network_mode": stream_manager.network_mode,
        "status": stream_manager.status,
    })


@app.post("/api/switch_camera/{camera_id}")
async def switch_camera(camera_id: int):
    """카메라 전환"""
    configs = INTERNAL_CONFIGS if stream_manager.network_mode == "internal" else CAMERA_CONFIGS
    if camera_id not in configs:
        return JSONResponse({"error": "Invalid camera ID"}, status_code=400)
    stream_manager.set_camera(camera_id)
    return JSONResponse({"ok": True, "camera_id": camera_id})


@app.post("/api/network_mode/{mode}")
async def set_network_mode(mode: str):
    """네트워크 모드 전환 (external / internal)"""
    if mode not in ("external", "internal"):
        return JSONResponse({"error": "Invalid mode"}, status_code=400)
    stream_manager.set_network_mode(mode)
    return JSONResponse({"ok": True, "mode": mode})


@app.post("/api/ptz/{camera_id}/{command}")
async def ptz_command(camera_id: int, command: str, speed1: int = 10, speed2: int = 10):
    """PTZ 제어 명령 프록시"""
    configs = INTERNAL_CONFIGS if stream_manager.network_mode == "internal" else CAMERA_CONFIGS
    if camera_id not in configs:
        return JSONResponse({"error": "Invalid camera ID"}, status_code=400)

    cfg = configs[camera_id]
    ip = cfg["ctrl_ip"]
    port = cfg["ctrl_port"]
    base_url = f"http://{ip}:{port}/cgi-bin/ptzctrl.cgi"

    # 명령 URL 구성
    if command == "stop":
        url = f"{base_url}?ptzcmd&ptzstop&0&0"
    elif command == "zoomstop":
        url = f"{base_url}?ptzcmd&zoomstop&5"
    elif command in ("zoomin", "zoomout"):
        url = f"{base_url}?ptzcmd&{command}&{speed1}"
    elif command == "preset":
        url = f"{base_url}?ptzcmd&poscall&{speed1}"
    elif command in ("up", "down", "left", "right"):
        url = f"{base_url}?ptzcmd&{command}&{speed1}&{speed2}"
    else:
        return JSONResponse({"error": "Unknown command"}, status_code=400)

    def do_request():
        try:
            resp = requests.get(url, auth=HTTPBasicAuth("admin", "admin"), timeout=3)
            return resp.status_code == 200
        except RequestException:
            return False

    loop = asyncio.get_event_loop()
    ok = await loop.run_in_executor(None, do_request)
    return JSONResponse({"ok": ok, "url": url})


@app.websocket("/ws/stream")
async def websocket_stream(websocket: WebSocket):
    """카메라 영상 프레임을 WebSocket으로 스트리밍"""
    await websocket.accept()
    q = stream_manager.subscribe()

    # 현재 상태 즉시 전송
    await websocket.send_text(json.dumps({
        "type": "status",
        "status": stream_manager.status,
        "camera_id": stream_manager.current_camera_id,
    }))

    try:
        while True:
            msg = await asyncio.wait_for(q.get(), timeout=10.0)
            await websocket.send_text(msg)
    except (WebSocketDisconnect, asyncio.TimeoutError):
        pass
    except Exception:
        pass
    finally:
        stream_manager.unsubscribe(q)
        try:
            await websocket.close()
        except Exception:
            pass


# ─────────────────── 진입점 ───────────────────

if __name__ == "__main__":
    print("=" * 55)
    print("  성림 라이브 제어기 - 웹 서버")
    print("=" * 55)
    print("  로컬 접속:  http://localhost:8000")
    print("  LAN  접속:  http://<이 PC의 IP>:8000")
    print("  외부 접속:  공유기 포트포워딩 후 http://<공인IP>:8000")
    print("  종료:       Ctrl+C")
    print("=" * 55)
    uvicorn.run(app, host="0.0.0.0", port=8000, log_level="warning")
