#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""
=============================================================================
THẬP NGŨ NIÊN (THE FIFTEEN SPRINGS) - MULTI-WORKER PARALLEL ORCHESTRATOR
=============================================================================
Hệ thống điều phối sản xuất video song song đa phiên (Multi-Worker Engine):
- Quản lý đồng thời 5-6 phiên Muse.ai cô lập hoàn toàn (muse_w1 ... muse_w6).
- Mỗi Worker gắn với 1 tài khoản Muse riêng biệt và thư mục Download riêng.
- Triệt tiêu 100% Race Condition và xung đột file tải về.
- Bảo toàn tuyệt đối quy chuẩn Head-Tail Chaining nội bộ trong từng Scene.
- Phân phối tải thông minh (Work-Stealing Dynamic Scene Queue).
=============================================================================
"""

import os
import sys
import re
import time
import json
import argparse
import subprocess
import threading
from queue import Queue, Empty
from pathlib import Path
from typing import Dict, List, Optional, Tuple
from concurrent.futures import ThreadPoolExecutor, as_completed

# Console UTF-8 Windows
if sys.platform == "win32":
    try:
        sys.stdout.reconfigure(encoding="utf-8")
        sys.stderr.reconfigure(encoding="utf-8")
    except Exception:
        pass

BASE_DIR = Path(__file__).resolve().parent.parent
PROMPTS_DIR = BASE_DIR / "02_AI_Prompts"
ASSETS_DIR = BASE_DIR / "04_Assets"
VIDEOS_DIR = ASSETS_DIR / "videos"
KEYFRAMES_DIR = ASSETS_DIR / "keyframes"
TEMP_DOWNLOADS_DIR = ASSETS_DIR / "temp_downloads"
MUSE_PROMPTS_FILE = PROMPTS_DIR / "muse_ai_video_prompts.json"

# Nạp các module liên kết
from production_orchestrator import (
    load_muse_prompts,
    get_all_shots,
    get_shots_for_scene,
    find_rendered_video,
    find_tail_frame,
    concat_scene_shots,
    batch_render_scene
)

try:
    from muse_api_client import MuseApiClient
    HAS_MUSE_API = True
except Exception:
    try:
        from .muse_api_client import MuseApiClient
        HAS_MUSE_API = True
    except Exception:
        HAS_MUSE_API = False
        MuseApiClient = None

print_lock = threading.Lock()

def safe_print(*args, **kwargs):
    """Thread-safe logging để log giữa các worker không bị xé dòng."""
    with print_lock:
        print(*args, **kwargs)
        sys.stdout.flush()


class MuseWorker:
    def __init__(self, worker_id: int, engine: str = "auto"):
        self.worker_id = worker_id
        self.engine = engine
        # Worker 1 tái sử dụng session 'muse' đã đăng nhập sẵn; Worker 2..N dùng 'muse_w2'..'muse_wN'
        self.session_name = "muse" if worker_id == 1 else f"muse_w{worker_id}"
        self.download_dir = TEMP_DOWNLOADS_DIR / f"w{worker_id}"
        self.download_dir.mkdir(parents=True, exist_ok=True)
        self.current_scene: Optional[str] = None
        self.status = "IDLE"
        self.completed_scenes: List[str] = []
        self.failed_scenes: List[str] = []

    def open_headed_setup(self):
        """Mở cửa sổ trình duyệt có giao diện để người dùng đăng nhập tài khoản Muse."""
        safe_print(f"\n[Worker {self.worker_id}] 🌐 Đang mở trình duyệt đăng nhập cho session '{self.session_name}'...")
        safe_print(f"👉 Vui lòng đăng nhập tài khoản Muse tương ứng vào cửa sổ vừa xuất hiện.")
        cmd = ["agent-browser", "--session", self.session_name, "--headed", "open", "https://muse.ai"]
        subprocess.run(cmd)

    def verify_auth_status(self) -> Tuple[bool, str]:
        """Kiểm tra tức thì trạng thái worker qua daemon metadata của agent-browser (0ms)."""
        dot_agent_dir = Path.home() / ".agent-browser"
        target_file = dot_agent_dir / f"{self.session_name}.target"
        pid_file = dot_agent_dir / f"{self.session_name}.pid"

        if not target_file.exists() or not pid_file.exists():
            return False, "Chưa khởi tạo session (Chạy --setup)"

        try:
            with open(target_file, "r", encoding="utf-8") as f:
                target_data = json.load(f)
            url = target_data.get("url", "")
            if not url or "about:blank" in url:
                return False, "Đang mở tab trống (Cần chạy --setup để đăng nhập)"
            elif "muse.ai" in url:
                return True, f"SẴN SÀNG (Muse.ai)"
            else:
                return False, f"Đang ở URL khác: {url[:30]}"
        except Exception as e:
            return False, f"Lỗi đọc target: {e}"

    def process_scene(self, scene_id: str, force: bool = False) -> bool:
        """Thực thi toàn bộ chuỗi Head-Tail của Scene bằng session riêng biệt."""
        self.current_scene = scene_id
        self.status = f"RUNNING ({scene_id})"
        safe_print(f"\n🚀 [Worker {self.worker_id} - {self.session_name}] TIẾP NHẬN CẢNH: {scene_id} (Engine: {self.engine})")

        try:
            success = batch_render_scene(
                scene_id=scene_id,
                force=force,
                auto_concat=True,
                session_name=self.session_name,
                download_dir=self.download_dir,
                engine=self.engine
            )
            if success:
                self.completed_scenes.append(scene_id)
                self.status = "IDLE"
                safe_print(f"🎉 [Worker {self.worker_id}] HOÀN TẤT XUẤT SẮC CẢNH: {scene_id}!")
                
                # Tự động cập nhật manifest của tập phim tương ứng
                try:
                    m_ep = re.search(r"^(ep\d+)", scene_id, re.IGNORECASE)
                    if m_ep:
                        ep_id = m_ep.group(1).lower()
                        from episode_manager import get_episode_assets_manifest, EPISODES_DIR
                        prompts_file = EPISODES_DIR / ep_id / "prompts" / "muse_prompts.json"
                        manifest_file = EPISODES_DIR / ep_id / "manifest.json"
                        if prompts_file.exists():
                            pdata = json.loads(prompts_file.read_text(encoding="utf-8"))
                            new_m = get_episode_assets_manifest(ep_id, pdata.get("motion_prompts", {}))
                            manifest_file.write_text(json.dumps(new_m, ensure_ascii=False, indent=2), encoding="utf-8")
                except Exception:
                    pass

                return True
            else:
                self.failed_scenes.append(scene_id)
                self.status = f"FAILED ({scene_id})"
                safe_print(f"❌ [Worker {self.worker_id}] THẤT BẠI TẠI CẢNH: {scene_id}")
                return False
        except Exception as e:
            self.failed_scenes.append(scene_id)
            self.status = f"ERROR: {str(e)[:20]}"
            safe_print(f"💥 [Worker {self.worker_id}] Lỗi ngoại lệ tại {scene_id}: {e}")
            return False
        finally:
            self.current_scene = None


class MultiWorkerOrchestrator:
    def __init__(self, num_workers: int = 6, engine: str = "auto"):
        self.num_workers = num_workers
        self.engine = engine
        self.workers = [MuseWorker(i + 1, engine=engine) for i in range(num_workers)]
        self.scene_queue = Queue()

    def setup_single_worker(self, worker_id: int):
        if 1 <= worker_id <= self.num_workers:
            self.workers[worker_id - 1].open_headed_setup()
        else:
            print(f"[!] Worker ID không hợp lệ: {worker_id} (Phải từ 1 đến {self.num_workers})")

    def verify_all_workers(self):
        print("\n" + "=" * 75)
        print("🔍 KIỂM TRA TRẠNG THÁI HẠ TẦNG SẢN XUẤT (MUSE2API & WORKER POOL)")
        print("=" * 75)

        # 1. Kiểm tra Muse2API Gateway (Ưu tiên số 1)
        if HAS_MUSE_API and MuseApiClient is not None:
            try:
                client = MuseApiClient()
                if client.is_available():
                    health = client.check_health()
                    acc = health.get("accounts", {})
                    print(f"⚡ MUSE2API GATEWAY: ĐANG TRỰC TUYẾN ({client.base_url})")
                    print(f"   ✓ Account Pool: {acc.get('available', 0)}/{acc.get('total', 0)} tài khoản khả dụng ngầm (Headless Multi-Account).")
                    print(f"   ✓ Trạng thái: Sẵn sàng tự động xoay vòng không cần mở trình duyệt thủ công.")
                    print("-" * 75)
            except Exception as e:
                print(f"⚠️ MUSE2API GATEWAY: Ngoại lệ kiểm tra: {e}")
                print("-" * 75)

        print(f"{'WORKER':<12} | {'SESSION NAME':<16} | {'THƯ MỤC DOWNLOAD':<22} | TRẠNG THÁI")
        print("-" * 75)

        results = {}
        with ThreadPoolExecutor(max_workers=self.num_workers) as executor:
            future_to_worker = {executor.submit(w.verify_auth_status): w for w in self.workers}
            for fut in as_completed(future_to_worker):
                w = future_to_worker[fut]
                try:
                    is_ok, msg = fut.result()
                    results[w.worker_id] = (is_ok, msg)
                except Exception as e:
                    results[w.worker_id] = (False, f"Lỗi: {e}")

        for w in self.workers:
            is_ok, msg = results.get(w.worker_id, (False, "Chưa rõ"))
            status_str = f"✓ {msg}" if is_ok else f"⚠️ {msg}"
            print(f"Worker {w.worker_id:<5} | {w.session_name:<16} | temp_downloads/w{w.worker_id:<8} | {status_str}")
        print("=" * 75 + "\n")

    def get_pending_scenes(self) -> List[str]:
        """Lấy danh sách các scene còn thiếu video Master hoặc chưa hoàn tất shot."""
        shots = get_all_shots()
        scenes = set()
        for shot_id in shots.keys():
            m = re.search(r"^(ep\d+_scene\d+)", shot_id, re.IGNORECASE)
            if m:
                scenes.add(m.group(1))

        # Sắp xếp thứ tự Scene
        sorted_scenes = sorted(list(scenes))
        pending = []
        for s in sorted_scenes:
            master_candidates = list(VIDEOS_DIR.glob(f"{s}*master*.mp4"))
            scene_shots = get_shots_for_scene(s)
            has_missing_shots = any(find_rendered_video(sid) is None for sid, _ in scene_shots)
            
            if not master_candidates or has_missing_shots:
                pending.append(s)
        return pending

    def dispatch_work_stealing_queue(self, scenes: Optional[List[str]] = None, force: bool = False):
        """
        Mô hình Dynamic Work-Stealing:
        Đẩy tất cả scene cần làm vào hàng đợi tập trung.
        Tất cả worker đồng thời rút việc để xử lý song song.
        """
        target_scenes = scenes if scenes else self.get_pending_scenes()
        if not target_scenes:
            safe_print("\n[✓] Tuyệt vời! Toàn bộ các Scene đều đã có Master và đã hoàn thành 100%!")
            return

        safe_print("\n" + "=" * 75)
        safe_print(f"🎬 KÍCH HOẠT HỆ THỐNG RENDER SONG SONG {self.num_workers} WORKERS")
        safe_print(f"   Tổng số Scene cần xử lý: {len(target_scenes)}")
        safe_print("=" * 75)

        for sc in target_scenes:
            self.scene_queue.put(sc)

        def worker_loop(worker: MuseWorker):
            while not self.scene_queue.empty():
                try:
                    scene_id = self.scene_queue.get_nowait()
                except Empty:
                    break

                worker.process_scene(scene_id, force=force)
                self.scene_queue.task_done()

        # Khởi chạy Pool
        with ThreadPoolExecutor(max_workers=self.num_workers) as executor:
            futures = [executor.submit(worker_loop, w) for w in self.workers]
            for f in as_completed(futures):
                try:
                    f.result()
                except Exception as e:
                    safe_print(f"[!] Worker thread encountered error: {e}")

        safe_print("\n" + "=" * 75)
        safe_print("🏁 TẤT CẢ CÁC WORKER ĐÃ HOÀN TẤT NHIỆM VỤ TRONG HÀNG ĐỢI!")
        safe_print("=" * 75)
        self.print_summary()

    def dispatch_by_episode(self, force: bool = False):
        """
        Phân phối tĩnh theo Tập: Mỗi Worker phụ trách trọn vẹn 1 Tập (EP01 -> EP06)
        """
        safe_print("\n" + "=" * 75)
        safe_print(f"🎬 KÍCH HOẠT PHÂN PHỐI 6 TẬP ĐỒNG THỜI (EP01 -> EP06)")
        safe_print("=" * 75)

        ep_assignments: Dict[int, List[str]] = {i + 1: [] for i in range(self.num_workers)}
        pending = self.get_pending_scenes()

        for sc in pending:
            m = re.search(r"^ep(\d+)_scene", sc, re.IGNORECASE)
            if m:
                ep_num = int(m.group(1))
                worker_idx = ((ep_num - 1) % self.num_workers) + 1
                ep_assignments[worker_idx].append(sc)

        def run_assigned_worker(w: MuseWorker, scenes_to_run: List[str]):
            for sc in scenes_to_run:
                w.process_scene(sc, force=force)

        with ThreadPoolExecutor(max_workers=self.num_workers) as executor:
            futures = []
            for w in self.workers:
                scenes_for_w = ep_assignments[w.worker_id]
                if scenes_for_w:
                    safe_print(f"   • Worker {w.worker_id} ({w.session_name}) nhận {len(scenes_for_w)} cảnh.")
                    futures.append(executor.submit(run_assigned_worker, w, scenes_for_w))
                else:
                    safe_print(f"   • Worker {w.worker_id} ({w.session_name}) không có cảnh nào còn thiếu.")

            for f in as_completed(futures):
                f.result()

        safe_print("\n" + "=" * 75)
        safe_print("🏁 HOÀN TẤT QUÁ TRÌNH RENDER 6 TẬP SONG SONG!")
        safe_print("=" * 75)
        self.print_summary()

    def dispatch_single_episode(self, episode_id: str, force: bool = False):
        """
        Dồn lực nhiều worker song song (đa tài khoản Muse) cho một Tập cụ thể.
        Ví dụ: Phân chia 15 cảnh của Tập 1 cho 3-6 worker chạy đồng thời.
        """
        ep_prefix = episode_id.lower()
        all_pending = self.get_pending_scenes()
        ep_scenes = [sc for sc in all_pending if sc.lower().startswith(ep_prefix)]

        if not ep_scenes:
            safe_print(f"\n[✓] Toàn bộ các cảnh của {episode_id.upper()} đều đã hoàn tất!")
            return

        safe_print("\n" + "=" * 75)
        safe_print(f"🎬 KÍCH HOẠT {self.num_workers} WORKERS SONG SONG CHO TẬP: {episode_id.upper()}")
        safe_print(f"   Số cảnh cần sản xuất: {len(ep_scenes)}")
        safe_print("=" * 75)

        self.dispatch_work_stealing_queue(scenes=ep_scenes, force=force)

    def print_summary(self):
        safe_print("\n📊 BÁO CÁO KẾT QUẢ SẢN XUẤT CỦA CÁC WORKER:")
        safe_print("-" * 65)
        for w in self.workers:
            comp_count = len(w.completed_scenes)
            fail_count = len(w.failed_scenes)
            safe_print(f"Worker {w.worker_id} ({w.session_name:<10}): {comp_count} cảnh hoàn thành | {fail_count} lỗi")
        safe_print("-" * 65 + "\n")


if __name__ == "__main__":
    parser = argparse.ArgumentParser(description="Multi-Worker Parallel Orchestrator for Muse.ai (Thập Ngũ Niên)")
    parser.add_argument("--workers", type=int, default=6, help="Số lượng worker chạy song song (mặc định: 6)")
    parser.add_argument("--setup", type=int, help="Mở trình duyệt có giao diện đăng nhập cho Worker ID (1..N)")
    parser.add_argument("--verify", action="store_true", help="Kiểm tra trạng thái đăng nhập của toàn bộ worker")
    parser.add_argument("--dispatch-episodes", action="store_true", help="Chạy song song 6 Tập (mỗi worker 1 tập)")
    parser.add_argument("--dispatch-queue", action="store_true", help="Chạy song song hàng đợi Scene linh hoạt (Work Stealing)")
    parser.add_argument("--episode", type=str, help="Chạy song song đa worker cho một tập cụ thể (vd: ep01)")
    parser.add_argument("--scenes", nargs="+", help="Danh sách scene cụ thể cần chạy song song")
    parser.add_argument("--force", action="store_true", help="Bắt buộc render lại kể cả đã có video")
    parser.add_argument("--engine", default="auto", choices=["auto", "muse", "muse_api", "gradio_ltx"], help="Động cơ render (mặc định: auto - tự động phân luồng theo góc máy)")

    args = parser.parse_args()
    orchestrator = MultiWorkerOrchestrator(num_workers=args.workers, engine=args.engine)

    if args.setup:
        orchestrator.setup_single_worker(args.setup)
    elif args.verify:
        orchestrator.verify_all_workers()
    elif args.dispatch_episodes:
        orchestrator.dispatch_by_episode(force=args.force)
    elif args.episode:
        orchestrator.dispatch_single_episode(args.episode, force=args.force)
    elif args.dispatch_queue:
        orchestrator.dispatch_work_stealing_queue(scenes=args.scenes, force=args.force)
    elif args.scenes:
        orchestrator.dispatch_work_stealing_queue(scenes=args.scenes, force=args.force)
    else:
        # Mặc định in trạng thái và hướng dẫn
        orchestrator.verify_all_workers()
        pending = orchestrator.get_pending_scenes()
        print(f"Số Scene đang chờ sản xuất: {len(pending)}")
        if pending:
            print("Các scene cần render:", ", ".join(pending[:10]), "..." if len(pending) > 10 else "")
            print("\n💡 Các lệnh vận hành nhanh:")
            print("  python 05_Production_Pipeline/multi_worker_orchestrator.py --setup 1          # Đăng nhập Worker 1")
            print("  python 05_Production_Pipeline/multi_worker_orchestrator.py --verify           # Kiểm tra tất cả Worker")
            print("  python 05_Production_Pipeline/multi_worker_orchestrator.py --episode ep01    # Dồn đa worker render Tập 1")
            print("  python 05_Production_Pipeline/multi_worker_orchestrator.py --dispatch-queue   # Bắt đầu chạy song song")
