#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""웹진 전용 숏츠 대본 Worker. 매물HOME v1.0 운영 Queue는 변경하지 않는다."""

import argparse
import html
import json
import logging
import os
import re
import sys
import urllib.error
import urllib.request
import time
from pathlib import Path

BASE_DIR=Path(__file__).resolve().parents[1]
if str(BASE_DIR) not in sys.path: sys.path.insert(0,str(BASE_DIR))
from db import get_conn

API_URL=os.getenv('V1_SHORTS_SCRIPT_API_URL','http://127.0.0.1:9000/v1/generate-script')
API_KEY=os.getenv('V1_SHORTS_API_KEY','honghee-shorts-secret-2026')
TIMEOUT=int(os.getenv('V1_SHORTS_SCRIPT_TIMEOUT','240'))
LOCK_NAME='multi_webzine_shorts_script_worker'
LOGGER=logging.getLogger('multi_shorts_worker')


def clean(value,limit=14000):
    text=str(value or '')
    text=re.sub(r'<script\b[^>]*>.*?</script>',' ',text,flags=re.I|re.S)
    text=re.sub(r'<style\b[^>]*>.*?</style>',' ',text,flags=re.I|re.S)
    text=re.sub(r'<br\s*/?>','\n',text,flags=re.I)
    text=re.sub(r'</(?:p|div|li|tr|h[1-6])>','\n',text,flags=re.I)
    text=html.unescape(re.sub(r'<[^>]+>',' ',text))
    text=re.sub(r'[ \t]+',' ',text); text=re.sub(r' *\n *','\n',text)
    return re.sub(r'\n{3,}','\n\n',text).strip()[:limit]


def claim(conn,job_id=0):
    conn.begin()
    try:
        sql="SELECT * FROM multi_shorts_jobs WHERE queue_status IN ('pending','retry_wait') AND (next_retry_at IS NULL OR next_retry_at<=NOW())"
        params=[]
        if job_id: sql+=' AND id=%s'; params.append(job_id)
        sql+=' ORDER BY scheduled_at,id LIMIT 1 FOR UPDATE'
        with conn.cursor() as cur:
            cur.execute(sql,params); row=cur.fetchone()
            if not row: conn.rollback(); return None
            cur.execute("UPDATE multi_shorts_jobs SET queue_status='scripting',started_at=COALESCE(started_at,NOW()),error_message=NULL,updated_at=NOW() WHERE id=%s",(row['id'],))
        conn.commit(); return row
    except Exception: conn.rollback(); raise


def source_bundle(conn,job):
    with conn.cursor() as cur:
        cur.execute("""SELECT ma.id,ma.source_article_no,ma.title,ma.summary,ma.body_html,
                              ma.primary_image_url,s.slug,r.office_name,r.representative_name,
                              r.office_phone,r.mobile_phone,r.region
                         FROM multi_articles ma
                         INNER JOIN multi_sites s ON s.id=ma.site_id
                         INNER JOIN blog_realtors r ON r.id=ma.realtor_id
                        WHERE ma.id=%s AND ma.site_id=%s AND ma.status='published' LIMIT 1""",(job['article_id'],job['site_id']))
        row=cur.fetchone() or {}
        if not row: raise RuntimeError('공개 상태의 웹진 원본 기사를 찾지 못했습니다.')
        cur.execute("SELECT image_url FROM multi_article_images WHERE article_id=%s ORDER BY sort_order,id LIMIT 30",(row['id'],))
        image_rows=cur.fetchall() or []
    images=[]
    for value in [row.get('primary_image_url'),*[item.get('image_url') for item in image_rows]]:
        value=str(value or '').strip()
        if value.startswith(('http://','https://')) and value not in images: images.append(value)
    content=clean(row.get('body_html') or row.get('summary') or row.get('title'))
    return {**row,'content':content,'image_urls':images,'article_url':f"http://multi.hongheemarketing.com/{row['slug']}/article/{row['id']}"}


def payload(job,source):
    lines=['부동산 숏츠 제작용 웹진 원본 정보',f"매물번호: {source['source_article_no']}",f"제목: {source['title']}",f"웹진 URL: {source['article_url']}",f"중개사무소: {source.get('office_name') or ''}",f"대표자: {source.get('representative_name') or ''}",f"지역: {source.get('region') or ''}",'','웹진 본문:',source['content'],'','작성 조건:','- 45초 세로형 부동산 숏츠','- 첫 3초에 핵심 매력과 가격 제시','- 허위·과장·확정 수익 표현 금지','- 짧고 읽기 쉬운 한국어 자막','- 마지막은 중개사무소 상담 안내']
    content='\n'.join(lines)
    # v1/기존 대본 API가 사용하는 키를 모두 보낸다. 기존에는 title/content만 보내
    # 일부 API에서 입력이 빈 값으로 해석되어 HTTP 500 또는 빈 대본이 발생했다.
    return {'project_id':int(job['id']),'realtor_id':int(job['realtor_id']),'source_type':'multi_webzine','title':source['title'],'post_title':source['title'],'source_url':source['article_url'],'post_url':source['article_url'],'article_no':str(source['source_article_no']),'content':content,'post_summary':content,'region':source.get('region') or '','purpose':'매물 안내 및 중개사 상담 문의 유도','script_type':'basic','generation_mode':'render_video','platforms':['youtube'],'duration_seconds':55,'video_width':1080,'video_height':1920,'image_urls':source['image_urls']}


def fallback_script(data):
    title=clean(data.get('post_title') or data.get('title'),120) or '부동산 매물'
    body=clean(data.get('post_summary') or data.get('content'),4000)
    candidates=[]
    for sentence in re.split(r'(?<=[.!?。？！])\s+|\n+',body):
        sentence=clean(sentence,100)
        if sentence and sentence not in candidates and not sentence.startswith(('매물번호:','웹진 URL:','작성 조건:','-')):
            candidates.append(sentence)
    office='등록 공인중개사사무소'
    match=re.search(r'중개사무소:\s*([^\n]+)',body)
    if match:office=clean(match.group(1),60) or office
    parts=[f'{title}, 핵심 조건을 확인해 보세요.']+candidates[:5]+[f'자세한 조건과 현장 확인은 {office}로 문의해 주세요.']
    return clean(' '.join(parts),900)

def call_api(data):
    urls=[API_URL]
    if '/v1/generate-script' in API_URL:
        urls.append(API_URL.replace('/v1/generate-script','/generate-script'))
    errors=[]
    for url in dict.fromkeys(urls):
        for attempt in range(2):
            request=urllib.request.Request(url,data=json.dumps(data,ensure_ascii=False).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=TIMEOUT) as response:raw=response.read().decode('utf-8','replace')
                result=json.loads(raw)
                if not isinstance(result,dict) or result.get('ok') is False:raise RuntimeError(str(result.get('error') or result.get('message') or '숏츠 대본 생성 실패'))
                if str(pick(result,'shorts_script','script','text')).strip():return result,raw
                raise RuntimeError('대본 API가 빈 대본을 반환했습니다.')
            except urllib.error.HTTPError as exc:errors.append(f'{url} HTTP {exc.code}: {exc.read().decode("utf-8","replace")[:500]}')
            except Exception as exc:errors.append(f'{url}: {exc}')
            if attempt==0:time.sleep(2)
    # Ollama/API 일시 장애가 수동·자동 작업 전체를 failed로 만들지 않도록
    # 기사 원문 기반 안전 대본으로 계속 렌더한다.
    script=fallback_script(data)
    result={'ok':True,'script':script,'shorts_title':data.get('title') or '부동산 매물 숏츠','hashtags':'#부동산 #매물 #Shorts','fallback_used':True,'fallback_reason':' | '.join(errors)[-3000:]}
    return result,json.dumps(result,ensure_ascii=False)


def pick(data,*names):
    for container in (data,data.get('result') or {},data.get('data') or {}):
        if isinstance(container,dict):
            for name in names:
                value=container.get(name)
                if value not in (None,'',[],{}): return value
    return ''


def finish(conn,job,source,request_data,result,raw):
    script=pick(result,'shorts_script','script','text')
    if not str(script).strip(): raise RuntimeError('숏츠 대본 API 응답에 script가 없습니다.')
    hashtags=pick(result,'hashtags'); hashtags=' '.join(map(str,hashtags)) if isinstance(hashtags,list) else str(hashtags or '')
    with conn.cursor() as cur:
        cur.execute("""UPDATE multi_shorts_jobs SET queue_status='script_ready',shorts_title=%s,
                       shorts_script=%s,hashtags=%s,video_prompt=%s,source_payload_json=%s,
                       script_response_json=%s,script_generated_at=NOW(),error_message=NULL,updated_at=NOW()
                       WHERE id=%s""",(str(pick(result,'shorts_title','title') or source['title'])[:255],str(script),hashtags,str(pick(result,'video_prompt') or ''),json.dumps(request_data,ensure_ascii=False,default=str),raw,job['id']))
        cur.execute("INSERT INTO multi_shorts_job_events (shorts_job_id,event_type,from_status,to_status,message,created_at) VALUES (%s,'script_completed','scripting','script_ready','웹진 전용 숏츠 대본 생성 완료',NOW())",(job['id'],))
    conn.commit()


def fail(conn,job,error):
    conn.rollback()
    retry=int(job.get('retry_count') or 0)+1
    status='failed' if retry>=3 else 'retry_wait'
    with conn.cursor() as cur: cur.execute("UPDATE multi_shorts_jobs SET queue_status=%s,retry_count=%s,error_message=%s,next_retry_at=IF(%s='retry_wait',DATE_ADD(NOW(),INTERVAL 5 MINUTE),next_retry_at),finished_at=IF(%s='failed',NOW(),NULL),updated_at=NOW() WHERE id=%s",(status,retry,str(error)[:60000],status,status,job['id']))
    conn.commit()


def main():
    parser=argparse.ArgumentParser(); parser.add_argument('--job-id',type=int,default=0); parser.add_argument('--dry-run',action='store_true'); args=parser.parse_args()
    logging.basicConfig(level=logging.INFO,format='%(asctime)s [%(levelname)s] %(message)s')
    conn=get_conn(); job=None; locked=False
    try:
        with conn.cursor() as cur: cur.execute('SELECT GET_LOCK(%s,0) acquired',(LOCK_NAME,)); locked=int((cur.fetchone() or {}).get('acquired') or 0)==1
        if not locked: print('[MULTI SHORTS WORKER SKIP] another_worker_running'); return
        job=claim(conn,args.job_id)
        if not job: print('[MULTI SHORTS WORKER DONE] no_pending_job'); return
        source=source_bundle(conn,job); data=payload(job,source)
        if args.dry_run:
            with conn.cursor() as cur: cur.execute("UPDATE multi_shorts_jobs SET queue_status='pending',started_at=NULL,updated_at=NOW() WHERE id=%s",(job['id'],))
            conn.commit(); print(json.dumps({'ok':True,'dry_run':True,'job_id':job['id'],'source':{'title':source['title'],'image_count':len(source['image_urls']),'article_url':source['article_url']}},ensure_ascii=False,indent=2)); return
        result,raw=call_api(data); finish(conn,job,source,data,result,raw); print(f"[MULTI SHORTS WORKER DONE] job_id={job['id']} status=script_ready images={len(source['image_urls'])}")
    except Exception as exc:
        if job: fail(conn,job,exc)
        LOGGER.exception('웹진 숏츠 대본 생성 실패: %s',exc); raise SystemExit(1)
    finally:
        if locked:
            try:
                with conn.cursor() as cur: cur.execute('SELECT RELEASE_LOCK(%s)',(LOCK_NAME,))
            except Exception: pass
        conn.close()


if __name__=='__main__': main()
