# -*- coding: utf-8 -*-
"""
STEP120 private_queue_worker.py

목적:
- schedule_daily_articles.py가 생성한 blog_post_private_queue pending 항목을
  publish_worker의 기존 비공개 처리 함수로 별도 처리한다.
- 23시 야간 수집/초안 이후에도, 신규 발행이 없어도 비공개 Queue를 처리할 수 있게 한다.
- 매물HOME v1.0 완료 전까지 server 발행 중개사의 비공개 Queue만 처리한다.
- client(매물HOME) 비공개 Queue는 삭제/변경하지 않고 pending 상태로 보존한다.

사용:
    python workers/private_queue_worker.py --limit-realtors 20

주의:
- 새 글 발행은 하지 않는다.
- 기존 publish_worker.py의 make_post_private/process_private_queues_for_realtor 로직을 재사용한다.
"""

import os
import sys
import argparse
import traceback
from datetime import datetime

from playwright.sync_api import sync_playwright

CURRENT_DIR = os.path.dirname(os.path.abspath(__file__))
ROOT_DIR = os.path.dirname(CURRENT_DIR)

if ROOT_DIR not in sys.path:
    sys.path.append(ROOT_DIR)

from db import get_conn

# publish_worker의 검증된 세션/비공개 처리 함수를 재사용한다.
from workers import publish_worker as pub


def get_target_realtor_ids(realtor_id=None, limit_realtors=50):
    conn = get_conn()

    try:
        sql = """
            SELECT pq.realtor_id
            FROM blog_post_private_queue pq
            WHERE pq.queue_status = 'pending'
              AND (pq.scheduled_at IS NULL OR pq.scheduled_at <= NOW())
              AND (pq.next_retry_at IS NULL OR pq.next_retry_at <= NOW())
              AND COALESCE(pq.blog_id, '') <> ''
              AND COALESCE(pq.blog_post_no, '') <> ''
              /* Queue 생성 당시의 발행 주체를 기준으로 처리한다.
                 현재 설정값을 사용하면 server로 발행된 과거 글이 이후 client 전환 시
                 영구 대기하는 문제가 생긴다. client Queue는 매물HOME이 처리한다. */
              AND COALESCE(NULLIF(LOWER(pq.execution_target), ''), 'server') = 'server'
        """

        params = []

        if realtor_id:
            sql += " AND pq.realtor_id = %s "
            params.append(int(realtor_id))

        sql += " GROUP BY pq.realtor_id ORDER BY MIN(pq.id) ASC "

        if limit_realtors and int(limit_realtors) > 0:
            sql += " LIMIT %s "
            params.append(int(limit_realtors))

        with conn.cursor() as cur:
            cur.execute(sql, params)
            rows = cur.fetchall() or []

        return [int(r.get("realtor_id")) for r in rows if r.get("realtor_id")]

    finally:
        conn.close()


def count_pending_private_queues(realtor_id=None):
    conn = get_conn()

    try:
        sql = """
            SELECT COUNT(*) AS cnt
            FROM blog_post_private_queue pq
            WHERE pq.queue_status = 'pending'
              AND (pq.scheduled_at IS NULL OR pq.scheduled_at <= NOW())
              AND (pq.next_retry_at IS NULL OR pq.next_retry_at <= NOW())
        """
        params = []

        if realtor_id:
            sql += " AND pq.realtor_id = %s "
            params.append(int(realtor_id))

        with conn.cursor() as cur:
            cur.execute(sql, params)
            row = cur.fetchone()

        return int((row or {}).get("cnt") or 0)

    finally:
        conn.close()


def record_private_completed_at(realtor_id, worker_started_at):
    """이번 Worker 실행에서 완료된 비공개 Queue의 완료시각을 기록한다."""
    conn = get_conn()

    try:
        with conn.cursor() as cur:
            cur.execute(
                """
                UPDATE blog_post_private_queue
                SET private_completed_at = COALESCE(finished_at, NOW())
                WHERE realtor_id = %s
                  AND queue_status = 'completed'
                  AND private_completed_at IS NULL
                  AND started_at >= %s
                """,
                (int(realtor_id), worker_started_at),
            )
            updated = int(cur.rowcount or 0)

        conn.commit()
        return updated

    finally:
        conn.close()


def mark_realtor_private_session_hold(realtor_id, message, limit=20):
    queues = pub.get_pending_private_queues(
        realtor_id=realtor_id,
        limit=limit,
    )

    for q in queues:
        try:
            pub.update_private_queue_status(
                q.get("id"),
                "hold",
                worker_status="session_hold",
                error_message=f"[PRIVATE SESSION_HOLD] {message}",
                last_error_type="session_hold",
            )
        except Exception as e:
            print("[PRIVATE HOLD UPDATE ERROR]", realtor_id, q.get("id"), str(e))


def process_private_for_realtor(pw, realtor_id, per_realtor_limit=3, headless=False):
    session_file = pub.get_session_file(realtor_id)
    worker_started_at = datetime.now()

    print("=" * 80)
    print("[PRIVATE REALTOR START]", f"realtor_id={realtor_id}")
    print("[PRIVATE SESSION FILE]", session_file)

    if not os.path.exists(session_file):
        message = f"세션 파일 없음: {session_file}"
        print("[PRIVATE SESSION MISSING]", message)
        try:
            pub.update_naver_session_status(
                realtor_id,
                "missing",
                message,
                session_file=session_file,
            )
        except Exception:
            pass

        mark_realtor_private_session_hold(
            realtor_id,
            message,
            limit=per_realtor_limit,
        )
        return {"processed": 0, "success": 0, "failed": 1, "session_hold": 1}

    browser = None

    try:
        browser = pw.chromium.launch(
            headless=bool(headless),
            args=[
                "--disable-blink-features=AutomationControlled",
                "--no-sandbox",
            ],
        )

        context = browser.new_context(
            storage_state=session_file,
            locale="ko-KR",
            viewport={
                "width": 1400,
                "height": 1000,
            },
        )

        page = context.new_page()

        print("[PRIVATE OPEN BLOG HOME] https://blog.naver.com")

        page.goto(
            "https://blog.naver.com",
            wait_until="domcontentloaded",
            timeout=60000,
        )

        page.wait_for_timeout(1200)

        print("[PRIVATE CURRENT URL]", page.url)

        if pub.is_login_page_url(page.url):
            recovered, recover_message = pub.recover_login_page_once(
                page=page,
                context=context,
                realtor_id=realtor_id,
                session_file=session_file,
                reason=f"비공개 처리 전 블로그 홈 로그인 페이지 이동: {page.url}",
                target_url="https://blog.naver.com",
            )

            if not recovered:
                print("[PRIVATE SESSION RECOVER FAIL]", recover_message)
                mark_realtor_private_session_hold(
                    realtor_id,
                    recover_message,
                    limit=per_realtor_limit,
                )
                return {"processed": 0, "success": 0, "failed": 1, "session_hold": 1}

            print("[PRIVATE SESSION RECOVERED] blog home")

        try:
            context.storage_state(path=session_file)
            pub.update_naver_session_status(
                realtor_id,
                "linked",
                "비공개 처리 전 blog.naver.com 접속 성공 / 세션 재저장",
                session_file=session_file,
            )
            print("[PRIVATE SESSION REFRESHED]")
        except Exception as e:
            print("[PRIVATE SESSION REFRESH ERROR]", str(e))

        try:
            pub.close_guides(page)
        except Exception:
            pass

        stats = pub.process_private_queues_for_realtor(
            page=page,
            context=context,
            realtor_id=realtor_id,
            session_file=session_file,
            limit=per_realtor_limit,
        )

        completed_at_updated = record_private_completed_at(
            realtor_id=realtor_id,
            worker_started_at=worker_started_at,
        )
        stats["private_completed_at_updated"] = completed_at_updated
        print("[PRIVATE COMPLETED AT UPDATED]", completed_at_updated)

        print("[PRIVATE REALTOR RESULT]", stats)
        return stats

    except Exception as e:
        traceback.print_exc()
        print("[PRIVATE REALTOR ERROR]", realtor_id, str(e))
        return {"processed": 0, "success": 0, "failed": 1, "error": str(e)}

    finally:
        try:
            if browser:
                browser.close()
        except Exception:
            pass


def run_private_queue_worker(
    realtor_id=None,
    limit_realtors=50,
    per_realtor_limit=3,
    headless=False,
    use_db_lock=True,
    lock_ttl_minutes=720,
):
    run_id = pub.create_pipeline_run(
        run_type="private_queue_worker",
        meta={
            "realtor_id": realtor_id,
            "limit_realtors": limit_realtors,
            "per_realtor_limit": per_realtor_limit,
            "headless": headless,
            "worker_server": pub.get_worker_server_name(),
        },
    )

    lock_owner = None
    final_status = "success"
    final_message = "private_queue_worker completed"

    processed = 0
    success = 0
    failed = 0
    skipped = 0

    try:
        if use_db_lock:
            db_locked, lock_owner = pub.acquire_db_lock(
                "private_queue_worker",
                ttl_minutes=lock_ttl_minutes,
            )

            if not db_locked:
                final_status = "skipped"
                final_message = "private_queue_worker already running or lock exists"
                print("[PRIVATE DB LOCKED]", final_message)
                return

        before_count = count_pending_private_queues(realtor_id=realtor_id)
        print("[PRIVATE PENDING BEFORE]", before_count)

        target_ids = get_target_realtor_ids(
            realtor_id=realtor_id,
            limit_realtors=limit_realtors,
        )

        print("[PRIVATE TARGET REALTORS]", target_ids)

        if not target_ids:
            final_message = "pending private queue 없음"
            print("[PRIVATE NONE]", final_message)
            return

        with sync_playwright() as pw:
            for rid in target_ids:
                stats = process_private_for_realtor(
                    pw=pw,
                    realtor_id=rid,
                    per_realtor_limit=per_realtor_limit,
                    headless=headless,
                )

                processed += int(stats.get("processed") or 0)
                success += int(stats.get("success") or 0)
                failed += int(stats.get("failed") or 0)

        after_count = count_pending_private_queues(realtor_id=realtor_id)
        print("[PRIVATE PENDING AFTER]", after_count)

        if failed > 0:
            final_status = "failed"
            final_message = f"completed with failures: success={success}, failed={failed}"

    except Exception as e:
        traceback.print_exc()
        final_status = "failed"
        final_message = str(e)[:2000]
        failed += 1

    finally:
        try:
            if lock_owner:
                pub.release_db_lock("private_queue_worker", lock_owner)
        except Exception:
            pass

        try:
            pub.finish_pipeline_run(
                run_id=run_id,
                status=final_status,
                total=processed,
                success=success,
                failed=failed,
                skipped=skipped,
                message=final_message,
                meta={
                    "realtor_id": realtor_id,
                    "limit_realtors": limit_realtors,
                    "per_realtor_limit": per_realtor_limit,
                    "processed": processed,
                    "success": success,
                    "failed": failed,
                    "skipped": skipped,
                    "worker_server": pub.get_worker_server_name(),
                },
            )
        except Exception:
            pass

        print("=" * 80)
        print("[PRIVATE QUEUE WORKER DONE]")
        print("[PROCESSED]", processed)
        print("[SUCCESS]", success)
        print("[FAILED]", failed)
        print("=" * 80)


def main():
    parser = argparse.ArgumentParser()
    parser.add_argument("--realtor-id", type=int, default=None)
    parser.add_argument("--limit-realtors", type=int, default=50)
    parser.add_argument("--per-realtor-limit", type=int, default=3)
    parser.add_argument("--headless", type=int, default=0)
    parser.add_argument("--no-db-lock", action="store_true")
    parser.add_argument("--lock-ttl-minutes", type=int, default=720)

    args = parser.parse_args()

    run_private_queue_worker(
        realtor_id=args.realtor_id,
        limit_realtors=args.limit_realtors,
        per_realtor_limit=args.per_realtor_limit,
        headless=bool(args.headless),
        use_db_lock=not args.no_db_lock,
        lock_ttl_minutes=args.lock_ttl_minutes,
    )


if __name__ == "__main__":
    main()
