# -*- coding: utf-8 -*-
"""
STEP114: V2 Workday Pipeline uses V2 Publish Engine

변경점:
- 기존 STEP106 workday_pipeline.py는 발행 시 workers/publish_one_queue_step105.py를 직접 subprocess 호출했다.
- STEP114부터는 core.publish_engine.publish_queue(queue_id)를 통해 V2 Publish Worker를 호출한다.
- publish_engine 내부에서 검증/프로필/V2 Worker 호출/결과 확인을 담당한다.

기존 운영 publish_worker.py는 건드리지 않는다.
"""

import sys
import json
from pathlib import Path
from datetime import datetime

BASE_DIR = Path(__file__).resolve().parents[1]
if str(BASE_DIR) not in sys.path:
    sys.path.append(str(BASE_DIR))

from db import get_conn
from core.realtor_pipeline_core import (
    get_policy,
    set_publish_mode,
    set_auto_publish_enabled,
    check_session_db,
    print_counts,
    collect_candidate_for_realtor,
    fetch_details_adapter,
    generate_draft_direct,
    fetch_latest_publish_queue,
    mark_v2_fetched_articles_collected,
)
from core.publish_engine import publish_queue
from core.deleted_article_sync import sync_deleted_articles_for_realtor



def fetch_v2_current_naver_count(conn, realtor_id):
    """
    STEP107-09:
    V2 전체수집 후 blog_realtor_articles 기준 현재 active article_no 수를 가져온다.
    deleted_sync 안전장치 expected_current_count로 사용한다.
    """
    try:
        with conn.cursor() as cur:
            cur.execute("""
                SELECT COUNT(DISTINCT article_no) AS cnt
                FROM blog_realtor_articles
                WHERE realtor_id = %s
                  AND COALESCE(article_no, '') <> ''
                  AND COALESCE(article_status, 'active') = 'active'
                  AND COALESCE(visibility_status, 'public') = 'public'
            """, (int(realtor_id),))
            row = cur.fetchone()
        return int((row or {}).get("cnt") or 0)
    except Exception:
        return None


def log(*args):
    print(*args, flush=True)


def fetch_next_pending_queue(conn, realtor_id):
    with conn.cursor() as cur:
        cur.execute("""
            SELECT *
            FROM blog_publish_queue
            WHERE realtor_id = %s
              AND COALESCE(queue_status, '') IN ('pending', 'waiting', 'ready', 'failed')
            ORDER BY id DESC
            LIMIT 1
        """, (int(realtor_id),))
        return cur.fetchone()



def log_v2_workday_contract():
    """
    STEP107-22:
    V2 서버 Workday 운영 계약.
    실제 네이버 블로그 발행/비공개는 클라이언트가 담당한다.
    서버 Workday는 발행대기 Queue 생성까지가 원칙이다.
    """
    log("[V2 WORKDAY CONTRACT]", "08:00~23:00 server pipeline")
    log("[V2 WORKDAY CONTRACT]", "session_check -> latest_full_collect -> snapshot -> cache -> latest_one_detail -> draft -> publish_queue")
    log("[V2 WORKDAY CONTRACT]", "client handles actual publish/private actions")


def run_one_realtor_workday(
    realtor_id,
    set_mode="test",
    require_test_mode=True,
    collect=True,
    fetch=True,
    draft=True,
    publish=True,
    draft_buffer=1,
    replenish_count=1,
    fetch_limit=3,
    draft_limit=1,
    require_test_account=True,
    publish_dry_run=False,
    deleted_sync_after_publish=False,
):
    conn = get_conn()
    auto_temp_on = False
    publish_result = None
    deleted_sync_result = None

    try:
        if set_mode:
            set_publish_mode(conn, realtor_id, set_mode)

        policy = get_policy(conn, realtor_id)
        realtor = policy["realtor"]
        publish_mode = policy["publish_mode"]
        is_test_account = int(policy["is_test_account"] or 0)
        auto_enabled = int(policy["auto_publish_enabled"] or 0)

        log("=" * 80)
        log("[V2 WORKDAY PIPELINE START]")
        log("TIME:", datetime.now().strftime("%Y-%m-%d %H:%M:%S"))
        log("realtor_id:", realtor_id)
        log("office_name:", realtor.get("office_name"))
        log("publish_mode:", publish_mode)
        log("is_test_account:", is_test_account)
        log("auto_publish_enabled:", auto_enabled)
        log("=" * 80)

        if require_test_mode and not (publish_mode == "test" or is_test_account == 1):
            raise RuntimeError("require_test_mode enabled, but realtor is not test")

        session = check_session_db(conn, realtor_id)
        log("[SESSION DB CHECK]", json.dumps(session, ensure_ascii=False, default=str))
        if session.get("exists") and session.get("ok") is False:
            log("[SESSION DB WARN] actual login will be verified by publisher:", session.get("message"))

        print_counts(conn, realtor_id, "before workday")

        if collect:
            if auto_enabled == 0:
                set_auto_publish_enabled(conn, realtor_id, True)
                auto_temp_on = True
                log("[AUTO PUBLISH TEMP ON] workday collect only")

            collect_candidate_for_realtor(
                conn=conn,
                realtor=realtor,
                draft_buffer=draft_buffer,
                replenish_count=replenish_count,
            )

            if auto_temp_on:
                set_auto_publish_enabled(conn, realtor_id, False)
                auto_temp_on = False
                log("[AUTO PUBLISH RESTORED] OFF")

            print_counts(conn, realtor_id, "after collect")

        if fetch:
            try:
                conn.commit()
            except Exception:
                pass

            fetch_details_adapter(limit=fetch_limit, realtor_id=realtor_id)

            conn.close()
            conn = get_conn()
            try:
                log("[V2 FETCH STATUS SYNC CALL]", "after fetch_details_adapter")
                mark_v2_fetched_articles_collected(
                    conn=conn,
                    realtor_id=int(realtor_id),
                    limit=max(int(fetch_limit or 3) * 3, 10),
                )
            except Exception as e:
                log("[V2 FETCH STATUS SYNC WARN]", str(e))

            print_counts(conn, realtor_id, "after fetch")

        if draft:
            result = generate_draft_direct(
                conn=conn,
                realtor_id=realtor_id,
                limit=draft_limit,
                draft_buffer=draft_buffer,
                replenish_count=replenish_count,
            )
            log("[V2 WORKDAY DRAFT RESULT]", json.dumps(result, ensure_ascii=False, default=str))
            print_counts(conn, realtor_id, "after draft")

        queue = fetch_next_pending_queue(conn, realtor_id)
        if not queue:
            log("[V2 WORKDAY QUEUE NONE] pending publish queue not found")
            latest = fetch_latest_publish_queue(conn, realtor_id)
            log("[LATEST QUEUE]", json.dumps(latest or {}, ensure_ascii=False, default=str))
            return {
                "ok": False,
                "reason": "pending_queue_not_found",
                "latest_queue": latest,
            }

        log("[V2 WORKDAY QUEUE SELECTED]", json.dumps(queue, ensure_ascii=False, default=str))

        if publish:
            # 연결을 닫고 publish_engine이 자체 DB context를 사용하도록 한다.
            try:
                conn.close()
            except Exception:
                pass

            publish_result = publish_queue(
                queue_id=int(queue["id"]),
                require_test_account=require_test_account,
                dry_run=publish_dry_run,
            )

            # STEP107-22: 08~23 Workday에서는 삭제매물 비공개 Queue를 만들지 않는다.
            # 삭제매물 Sync는 23:00 이후 Snapshot 기반 야간 작업에서만 실행한다.
            if deleted_sync_after_publish and (not publish_dry_run) and publish_result and publish_result.get("ok"):
                try:
                    # STEP107-09:
                    # 방금 V2 전체수집으로 확보된 현재 네이버 매물 수를 expected_current_count로 넘긴다.
                    # 수집이 일부만 된 상태라면 deleted_sync가 safety_blocked=True로 비공개 Queue 생성을 막는다.
                    expected_current_count = None
                    conn_for_count = None
                    try:
                        conn_for_count = get_conn()
                        expected_current_count = fetch_v2_current_naver_count(conn_for_count, int(realtor_id))
                    finally:
                        try:
                            if conn_for_count:
                                conn_for_count.close()
                        except Exception:
                            pass

                    deleted_sync_result = sync_deleted_articles_for_realtor(
                        realtor_id=int(realtor_id),
                        dry_run=False,
                        expected_current_count=expected_current_count,
                        min_current_ratio=0.95,
                    )
                    log("[V2 WORKDAY DELETED SYNC RESULT]", json.dumps(deleted_sync_result, ensure_ascii=False, default=str))
                except Exception as e:
                    deleted_sync_result = {"ok": False, "error": str(e)}
                    log("[V2 WORKDAY DELETED SYNC WARN]", str(e))
            else:
                if publish_result and publish_result.get("ok"):
                    log("[V2 WORKDAY DELETED SYNC SKIP]", "deleted sync is reserved for after 23:00 snapshot job")

            conn = get_conn()
            print_counts(conn, realtor_id, "after publish")

            with conn.cursor() as cur:
                cur.execute("SELECT * FROM blog_publish_queue WHERE id=%s LIMIT 1", (int(queue["id"]),))
                final_queue = cur.fetchone() or {}

            log("[V2 WORKDAY FINAL QUEUE]", json.dumps(final_queue, ensure_ascii=False, default=str))

        log("=" * 80)
        log("[V2 WORKDAY PIPELINE DONE]")
        log("=" * 80)

        return {
            "ok": True,
            "queue_id": int(queue["id"]),
            "article_no": queue.get("article_no"),
            "publish_result": publish_result,
            "deleted_sync_result": deleted_sync_result,
        }

    finally:
        try:
            if auto_temp_on:
                set_auto_publish_enabled(conn, realtor_id, False)
                log("[AUTO PUBLISH RESTORED IN FINALLY] OFF")
        except Exception:
            pass

        try:
            conn.close()
        except Exception:
            pass
