# -*- coding: utf-8 -*-
"""MaemulHOME v1.0 sequential shorts render worker (2026-08-26 r1).

Only v1 tables are used. Existing v0.1.9.1/server publishing and legacy
shorts tables are not imported or modified.
"""

from __future__ import annotations

import argparse
import json
import logging
import os
import sys
import time
import urllib.error
import urllib.request
from pathlib import Path
from typing import Any, Dict, List, Optional, Sequence, Tuple


FILE_PATH = Path(__file__).resolve()
for candidate in (FILE_PATH.parent, FILE_PATH.parent.parent):
    value = str(candidate)
    if value not in sys.path:
        sys.path.insert(0, value)

try:
    from db import get_conn as _get_conn  # type: ignore
except ImportError:
    _get_conn = None


PROJECT_TABLE = "blog_v1_shorts_projects"
SETTINGS_TABLE = "blog_v1_shorts_settings"
JOBS_TABLE = "blog_v1_shorts_jobs"
EVENTS_TABLE = "blog_v1_shorts_events"
ADVISORY_LOCK = "maemulhome_v1_shorts_render_worker"
WORKER_BUILD = "20260902-r2-v1-fixed-layout"
RENDER_API_URL = os.getenv("V1_SHORTS_RENDER_API_URL", "http://127.0.0.1:9000/v1/render-video")
API_KEY = os.getenv("V1_SHORTS_API_KEY", "honghee-shorts-secret-2026")
API_TIMEOUT = int(os.getenv("V1_SHORTS_RENDER_TIMEOUT", "900"))
STALE_MINUTES = int(os.getenv("V1_SHORTS_RENDER_STALE_MINUTES", "60"))
LOGGER = logging.getLogger("v1_shorts_render_worker")


def open_connection() -> Any:
    if _get_conn is None:
        raise RuntimeError("db.py를 찾을 수 없습니다. 이 파일을 jobs 폴더에 두고 실행해 주세요.")
    return _get_conn()


def json_text(value: Any) -> str:
    return json.dumps(value, ensure_ascii=False, separators=(",", ":"), default=str)


def json_object(value: Any) -> Dict[str, Any]:
    if isinstance(value, dict):
        return value
    try:
        parsed = json.loads(str(value or ""))
        return parsed if isinstance(parsed, dict) else {}
    except (TypeError, ValueError, json.JSONDecodeError):
        return {}


def table_exists(conn: Any, table_name: str) -> bool:
    with conn.cursor() as cur:
        cur.execute(
            """SELECT COUNT(*) AS cnt FROM information_schema.TABLES
                 WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=%s""",
            (table_name,),
        )
        row = cur.fetchone() or {}
    return int(row.get("cnt") or 0) > 0


def acquire_lock(conn: Any) -> bool:
    with conn.cursor() as cur:
        cur.execute("SELECT GET_LOCK(%s, 0) AS acquired", (ADVISORY_LOCK,))
        row = cur.fetchone() or {}
    return int(row.get("acquired") or 0) == 1


def release_lock(conn: Any) -> None:
    try:
        with conn.cursor() as cur:
            cur.execute("SELECT RELEASE_LOCK(%s)", (ADVISORY_LOCK,))
    except Exception:
        LOGGER.exception("렌더링 advisory lock 해제 실패")


def claim_project(conn: Any, project_id: int = 0) -> Optional[Dict[str, Any]]:
    conn.begin()
    try:
        where = """
            (
                (p.workflow_status='script_ready' AND p.script_status='done'
                 AND p.render_status IN ('pending','failed'))
                OR
                (p.workflow_status='rendering' AND p.render_status='rendering'
                 AND p.updated_at < DATE_SUB(NOW(), INTERVAL %s MINUTE))
            )
        """
        params: List[Any] = [STALE_MINUTES]
        if project_id > 0:
            where += " AND p.id=%s"
            params.append(project_id)

        with conn.cursor() as cur:
            cur.execute(
                f"""SELECT p.* FROM {PROJECT_TABLE} p
                     WHERE {where}
                     ORDER BY p.created_at ASC, p.id ASC
                     LIMIT 1 FOR UPDATE""",
                tuple(params),
            )
            project = cur.fetchone()
            if not project:
                conn.rollback()
                return None
            cur.execute(
                f"""UPDATE {PROJECT_TABLE}
                       SET project_status='processing', workflow_status='rendering',
                           render_status='rendering', render_started_at=NOW(),
                           error_message=NULL, updated_at=NOW()
                     WHERE id=%s""",
                (project["id"],),
            )
            if table_exists(conn, JOBS_TABLE):
                cur.execute(
                    f"""UPDATE {JOBS_TABLE}
                           SET queue_status='processing', started_at=COALESCE(started_at,NOW()),
                               heartbeat_at=NOW(), updated_at=NOW()
                         WHERE project_id=%s AND job_type='generate_video'""",
                    (project["id"],),
                )
        conn.commit()
        project["workflow_status"] = "rendering"
        project["render_status"] = "rendering"
        return project
    except Exception:
        conn.rollback()
        raise


def fetch_settings(conn: Any, realtor_id: int) -> Dict[str, Any]:
    if not table_exists(conn, SETTINGS_TABLE):
        return {}
    with conn.cursor() as cur:
        cur.execute(f"SELECT * FROM {SETTINGS_TABLE} WHERE realtor_id=%s LIMIT 1", (realtor_id,))
        return cur.fetchone() or {}


def extract_source(project: Dict[str, Any]) -> Dict[str, Any]:
    request_data = json_object(project.get("script_request_json"))
    payload = request_data.get("payload") if isinstance(request_data.get("payload"), dict) else request_data
    source = request_data.get("source") if isinstance(request_data.get("source"), dict) else {}
    images = source.get("image_urls") or payload.get("image_urls") or []
    if not isinstance(images, list):
        images = []
    images = [str(value).strip() for value in images if str(value).strip().startswith(("http://", "https://"))]
    return {
        "image_urls": list(dict.fromkeys(images)),
        "source_url": str(project.get("blog_post_url") or source.get("blog_url") or payload.get("source_url") or ""),
    }


def enabled_platforms(settings: Dict[str, Any]) -> List[str]:
    mapping = (
        ("youtube_enabled", "youtube"),
        ("instagram_enabled", "instagram"),
        ("tiktok_enabled", "tiktok"),
        ("naver_shorts_enabled", "naver_shorts"),
    )
    result = [platform for column, platform in mapping if int(settings.get(column) or 0) == 1]
    return result or ["youtube"]


def build_payload(project: Dict[str, Any], settings: Dict[str, Any]) -> Dict[str, Any]:
    source = extract_source(project)
    if not str(project.get("shorts_script") or "").strip():
        raise RuntimeError("숏츠 대본이 없습니다.")
    return {
        "project_id": int(project.get("id") or 0),
        "render_key": f"v1_{int(project.get('id') or 0)}",
        "realtor_id": int(project.get("realtor_id") or 0),
        "article_no": str(project.get("article_no") or ""),
        "title": str(
            project.get("source_post_title")
            or project.get("post_title")
            or project.get("shorts_title")
            or "부동산 매물 안내"
        ),
        "script": str(project.get("shorts_script") or ""),
        "video_prompt": str(project.get("video_prompt") or ""),
        "hashtags": str(project.get("hashtags") or ""),
        "source_type": str(project.get("source_type") or "blog_published"),
        "source_blog_url": source["source_url"],
        "post_url": source["source_url"],
        "image_urls": source["image_urls"],
        "platforms": enabled_platforms(settings),
        "duration_seconds": int(settings.get("video_duration_seconds") or 45),
        "video_width": int(settings.get("video_width") or 1080),
        "video_height": int(settings.get("video_height") or 1920),
        "voice_code": str(settings.get("voice_code") or ""),
        "bgm_mode": str(settings.get("bgm_mode") or "random"),
        "caption_style": str(settings.get("caption_style") or "default"),
        "transition_style": str(settings.get("transition_style") or "random"),
        "save_on_python_server": True,
        "storage_policy": {
            "owner": "python_flask",
            "php_should_download_video": False,
            "php_should_store_video_file": False,
            "php_should_store_returned_path_only": True,
        },
    }


def call_render_api(payload: Dict[str, Any]) -> Tuple[Dict[str, Any], str]:
    request = urllib.request.Request(
        RENDER_API_URL,
        data=json_text(payload).encode("utf-8"),
        headers={"Content-Type": "application/json; charset=utf-8", "X-API-KEY": API_KEY},
        method="POST",
    )
    try:
        with urllib.request.urlopen(request, timeout=API_TIMEOUT) as response:
            raw = response.read().decode("utf-8", errors="replace")
            status = int(getattr(response, "status", 200))
    except urllib.error.HTTPError as exc:
        raw = exc.read().decode("utf-8", errors="replace")
        raise RuntimeError(f"render-video API HTTP {exc.code}: {raw[:1500]}") from exc
    except urllib.error.URLError as exc:
        raise RuntimeError(f"render-video API 연결 실패: {exc}") from exc
    if status >= 400:
        raise RuntimeError(f"render-video API HTTP {status}: {raw[:1500]}")
    try:
        data = json.loads(raw)
    except json.JSONDecodeError as exc:
        raise RuntimeError(f"render-video API 응답이 JSON이 아닙니다: {raw[:1500]}") from exc
    if not isinstance(data, dict) or data.get("ok") is False:
        message = data.get("error") or data.get("message") if isinstance(data, dict) else raw
        raise RuntimeError(str(message or "영상 렌더링 실패"))
    return data, raw


def pick_value(data: Dict[str, Any], paths: Sequence[Sequence[str]]) -> str:
    for path in paths:
        value: Any = data
        for key in path:
            if not isinstance(value, dict) or key not in value:
                value = None
                break
            value = value[key]
        if isinstance(value, str) and value.strip():
            return value.strip()
    return ""


def save_success(conn: Any, project: Dict[str, Any], payload: Dict[str, Any], data: Dict[str, Any], raw: str) -> None:
    project_id = int(project["id"])
    video_path = pick_value(data, (("video_path",), ("path",), ("result", "video_path"), ("data", "video_path"), ("output", "video_path")))
    video_url = pick_value(data, (("video_url",), ("url",), ("result", "video_url"), ("data", "video_url"), ("output", "video_url")))
    thumb_path = pick_value(data, (("thumbnail_path",), ("thumb_path",), ("result", "thumbnail_path"), ("data", "thumbnail_path")))
    thumb_url = pick_value(data, (("thumbnail_url",), ("result", "thumbnail_url"), ("data", "thumbnail_url")))
    if not video_path and not video_url:
        raise RuntimeError(f"render-video 응답에 video_path 또는 video_url이 없습니다: {raw[:1500]}")
    with conn.cursor() as cur:
        cur.execute(
            f"""UPDATE {PROJECT_TABLE}
                   SET project_status='processing', workflow_status='video_ready',
                       render_status='done', video_path=%s, video_url=%s,
                       thumbnail_path=%s, thumbnail_url=%s,
                       render_request_json=%s, render_response_json=%s,
                       render_completed_at=NOW(), error_message=NULL, updated_at=NOW()
                 WHERE id=%s""",
            (video_path or None, video_url or None, thumb_path or None, thumb_url or None,
             json_text(payload), raw, project_id),
        )
        if table_exists(conn, JOBS_TABLE):
            cur.execute(
                f"""UPDATE {JOBS_TABLE}
                       SET queue_status='done', finished_at=NOW(), heartbeat_at=NOW(),
                           last_error_type=NULL, last_error_message=NULL, updated_at=NOW()
                     WHERE project_id=%s AND job_type='generate_video'""",
                (project_id,),
            )
        if table_exists(conn, EVENTS_TABLE):
            cur.execute(
                f"""INSERT INTO {EVENTS_TABLE}
                    (project_id,realtor_id,event_type,event_status,message,context_json,created_by,created_at)
                    VALUES (%s,%s,'render_completed','success','v1 숏츠 영상 렌더링 완료',%s,'v1_render_worker',NOW())""",
                (project_id, int(project.get("realtor_id") or 0), json_text({"video_path": video_path, "video_url": video_url})),
            )
    conn.commit()


def save_failure(conn: Any, project: Dict[str, Any], error: Exception) -> None:
    project_id = int(project["id"])
    message = str(error).strip()[:60000]
    try:
        with conn.cursor() as cur:
            cur.execute(
                f"""UPDATE {PROJECT_TABLE}
                       SET project_status='failed', workflow_status='failed', render_status='failed',
                           error_message=%s, updated_at=NOW() WHERE id=%s""",
                (message, project_id),
            )
            if table_exists(conn, JOBS_TABLE):
                cur.execute(
                    f"""UPDATE {JOBS_TABLE}
                           SET queue_status='failed', finished_at=NOW(), retry_count=retry_count+1,
                               last_error_type='render_error', last_error_message=%s, updated_at=NOW()
                         WHERE project_id=%s AND job_type='generate_video'""",
                    (message, project_id),
                )
        conn.commit()
    except Exception:
        conn.rollback()
        LOGGER.exception("렌더링 실패 상태 저장 실패: project_id=%s", project_id)


def project_state(conn: Any, project_id: int) -> Optional[Dict[str, Any]]:
    with conn.cursor() as cur:
        cur.execute(
            f"""SELECT id,realtor_id,article_no,project_status,workflow_status,
                       script_status,render_status,error_message,updated_at
                  FROM {PROJECT_TABLE} WHERE id=%s LIMIT 1""",
            (project_id,),
        )
        return cur.fetchone()


def process_once(project_id: int = 0, dry_run: bool = False) -> int:
    conn = open_connection()
    project: Optional[Dict[str, Any]] = None
    locked = False
    try:
        if not table_exists(conn, PROJECT_TABLE):
            raise RuntimeError(f"{PROJECT_TABLE} 테이블이 없습니다.")
        locked = acquire_lock(conn)
        if not locked:
            LOGGER.info("다른 v1 렌더링 Worker가 실행 중이므로 건너뜁니다.")
            return 0
        project = claim_project(conn, project_id)
        if not project:
            if project_id > 0:
                state = project_state(conn, project_id)
                if state:
                    LOGGER.warning("project_id=%s 상태가 렌더링 조건과 다릅니다: %s", project_id, json_text(state))
                else:
                    LOGGER.warning("현재 연결 DB에 project_id=%s가 없습니다.", project_id)
            else:
                LOGGER.info("처리할 script_ready 렌더링 작업이 없습니다.")
            return 0
        settings = fetch_settings(conn, int(project.get("realtor_id") or 0))
        payload = build_payload(project, settings)
        LOGGER.info("숏츠 렌더링 시작: project_id=%s, images=%s, size=%sx%s, duration=%s",
                    project["id"], len(payload["image_urls"]), payload["video_width"],
                    payload["video_height"], payload["duration_seconds"])
        if dry_run:
            with conn.cursor() as cur:
                cur.execute(
                    f"""UPDATE {PROJECT_TABLE}
                           SET project_status='processing', workflow_status='script_ready',
                               render_status='pending', render_started_at=NULL, updated_at=NOW()
                         WHERE id=%s""",
                    (project["id"],),
                )
                if table_exists(conn, JOBS_TABLE):
                    cur.execute(
                        f"""UPDATE {JOBS_TABLE} SET queue_status='pending', started_at=NULL,
                               heartbeat_at=NULL, updated_at=NOW()
                             WHERE project_id=%s AND job_type='generate_video'""",
                        (project["id"],),
                    )
            conn.commit()
            print(json.dumps({"ok": True, "dry_run": True, "payload": payload}, ensure_ascii=False, indent=2, default=str))
            return 0
        data, raw = call_render_api(payload)
        save_success(conn, project, payload, data, raw)
        LOGGER.info("숏츠 렌더링 완료: project_id=%s", project["id"])
        return 1
    except Exception as exc:
        conn.rollback()
        if project:
            save_failure(conn, project, exc)
        LOGGER.exception("숏츠 렌더링 실패: %s", exc)
        return 2
    finally:
        if locked:
            release_lock(conn)
        conn.close()


def parse_args() -> argparse.Namespace:
    parser = argparse.ArgumentParser(description="MaemulHOME v1.0 sequential shorts render worker")
    parser.add_argument("--project-id", type=int, default=0, help="특정 프로젝트만 렌더링")
    parser.add_argument("--dry-run", action="store_true", help="렌더링 요청 자료만 검증하고 상태 복원")
    parser.add_argument("--loop", action="store_true", help="순차 Worker 계속 실행")
    parser.add_argument("--poll-seconds", type=int, default=30, help="대기 작업 확인 간격")
    parser.add_argument("--log-level", default="INFO", choices=("DEBUG", "INFO", "WARNING", "ERROR"))
    return parser.parse_args()


def main() -> int:
    args = parse_args()
    logging.basicConfig(level=getattr(logging, args.log_level), format="%(asctime)s [%(levelname)s] %(message)s")
    LOGGER.info("v1 숏츠 렌더 Worker 빌드: %s / table=%s", WORKER_BUILD, PROJECT_TABLE)
    if not args.loop:
        return process_once(args.project_id, args.dry_run)
    LOGGER.info("v1 숏츠 순차 렌더 Worker 시작: poll=%s초", max(5, args.poll_seconds))
    while True:
        process_once(args.project_id, args.dry_run)
        if args.project_id > 0:
            return 0
        time.sleep(max(5, args.poll_seconds))


if __name__ == "__main__":
    raise SystemExit(main())
