# -*- coding: utf-8 -*-
"""
STEP105: publish_one_queue_step105.py

기존 운영 workers/publish_worker.py는 건드리지 않는다.
STEP104 테스트본 workers/publish_worker_step104.py를 이용해 특정 queue_id 1건만 발행한다.

현재 버전은 "안전 wrapper" 방식:
1) queue_id의 realtor_id 확인
2) 같은 중개사의 다른 pending/waiting/ready/failed 큐를 잠시 hold 처리
3) 테스트 계정이면 오늘 발행 이력 시간을 하루 전으로 이동하여 daily limit 회피
4) publish_worker_step104 실행 동안만 auto_publish_enabled 임시 ON
5) 실행 후 auto_publish_enabled 원복
6) 다른 큐 상태 원복

주의:
- 운영 publish_worker.py는 수정하지 않는다.
- true publish_one(queue_id) 함수 분리는 다음 단계에서 진행한다.
"""

import os
import sys
import json
import time
import argparse
import subprocess
from pathlib import Path
from datetime import datetime

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

from db import get_conn

PYTHON_EXE = sys.executable


def clean_text(value):
    return str(value or "").strip()


def table_exists(conn, table_name):
    try:
        with conn.cursor() as cur:
            cur.execute("SHOW TABLES LIKE %s", (table_name,))
            return cur.fetchone() is not None
    except Exception:
        return False


def table_columns(conn, table_name):
    try:
        with conn.cursor() as cur:
            cur.execute(f"SHOW COLUMNS FROM {table_name}")
            rows = cur.fetchall()
        return {row["Field"] for row in rows}
    except Exception:
        return set()


def fetch_queue(conn, queue_id):
    with conn.cursor() as cur:
        cur.execute("""
            SELECT *
            FROM blog_publish_queue
            WHERE id = %s
            LIMIT 1
        """, (int(queue_id),))
        row = cur.fetchone()

    if not row:
        raise RuntimeError(f"publish queue not found: {queue_id}")

    return row


def fetch_realtor(conn, realtor_id):
    with conn.cursor() as cur:
        cur.execute("""
            SELECT *
            FROM blog_realtors
            WHERE id = %s
            LIMIT 1
        """, (int(realtor_id),))
        row = cur.fetchone()

    if not row:
        raise RuntimeError(f"realtor not found: {realtor_id}")

    return row


def fetch_publish_setting(conn, realtor_id):
    if not table_exists(conn, "blog_publish_settings"):
        return {}

    with conn.cursor() as cur:
        cur.execute("""
            SELECT *
            FROM blog_publish_settings
            WHERE realtor_id = %s
            LIMIT 1
        """, (int(realtor_id),))
        return cur.fetchone() or {}


def ensure_publish_setting_row(conn, realtor_id):
    if not table_exists(conn, "blog_publish_settings"):
        return

    columns = table_columns(conn, "blog_publish_settings")

    with conn.cursor() as cur:
        cur.execute("""
            SELECT id
            FROM blog_publish_settings
            WHERE realtor_id = %s
            LIMIT 1
        """, (int(realtor_id),))
        row = cur.fetchone()

    if row:
        return

    data = {
        "realtor_id": int(realtor_id),
        "auto_publish_enabled": 0,
        "daily_post_limit": 1,
        "created_at": datetime.now(),
        "updated_at": datetime.now(),
    }

    filtered = {k: v for k, v in data.items() if k in columns}

    if not filtered:
        return

    keys = list(filtered.keys())

    with conn.cursor() as cur:
        cur.execute(
            f"""
            INSERT INTO blog_publish_settings
            ({", ".join(keys)})
            VALUES ({", ".join(["%s"] * len(keys))})
            """,
            [filtered[k] for k in keys],
        )

    conn.commit()


def set_auto_publish(conn, realtor_id, enabled):
    if not table_exists(conn, "blog_publish_settings"):
        return

    columns = table_columns(conn, "blog_publish_settings")
    if "auto_publish_enabled" not in columns:
        return

    ensure_publish_setting_row(conn, realtor_id)

    sets = ["auto_publish_enabled = %s"]
    params = [1 if enabled else 0]

    if "updated_at" in columns:
        sets.append("updated_at = NOW()")

    params.append(int(realtor_id))

    with conn.cursor() as cur:
        cur.execute(
            f"""
            UPDATE blog_publish_settings
            SET {", ".join(sets)}
            WHERE realtor_id = %s
            """,
            params,
        )

    conn.commit()


def shift_today_published_for_test(conn, realtor_id):
    with conn.cursor() as cur:
        cur.execute("""
            UPDATE blog_publish_queue
            SET started_at = CASE
                    WHEN started_at IS NOT NULL THEN DATE_SUB(started_at, INTERVAL 1 DAY)
                    ELSE started_at
                END,
                finished_at = CASE
                    WHEN finished_at IS NOT NULL THEN DATE_SUB(finished_at, INTERVAL 1 DAY)
                    ELSE finished_at
                END,
                updated_at = NOW()
            WHERE realtor_id = %s
              AND queue_status = 'published'
              AND finished_at >= CURDATE()
              AND finished_at < DATE_ADD(CURDATE(), INTERVAL 1 DAY)
        """, (int(realtor_id),))
        affected = cur.rowcount

    conn.commit()
    return int(affected or 0)


def hold_other_pending_queues(conn, realtor_id, queue_id):
    """
    같은 realtor의 다른 발행대기 큐를 임시 hold 처리.
    원복을 위해 기존 상태를 반환한다.
    """
    with conn.cursor() as cur:
        cur.execute("""
            SELECT id, queue_status
            FROM blog_publish_queue
            WHERE realtor_id = %s
              AND id <> %s
              AND COALESCE(queue_status, '') IN ('pending', 'waiting', 'ready', 'failed')
        """, (int(realtor_id), int(queue_id)))
        rows = cur.fetchall() or []

    if not rows:
        return []

    ids = [int(r["id"]) for r in rows]

    placeholders = ",".join(["%s"] * len(ids))

    with conn.cursor() as cur:
        cur.execute(
            f"""
            UPDATE blog_publish_queue
            SET queue_status = 'hold',
                updated_at = NOW()
            WHERE id IN ({placeholders})
            """,
            ids,
        )

    conn.commit()

    print("[OTHER QUEUES TEMP HOLD]", json.dumps(rows, ensure_ascii=False, default=str))
    return rows


def restore_held_queues(conn, rows):
    if not rows:
        return

    for row in rows:
        qid = int(row["id"])
        status = clean_text(row.get("queue_status") or "pending") or "pending"

        with conn.cursor() as cur:
            cur.execute("""
                UPDATE blog_publish_queue
                SET queue_status = %s,
                    updated_at = NOW()
                WHERE id = %s
                  AND queue_status = 'hold'
            """, (status, qid))

    conn.commit()
    print("[OTHER QUEUES RESTORED]", len(rows))


def make_target_queue_pending(conn, queue_id):
    with conn.cursor() as cur:
        cur.execute("""
            UPDATE blog_publish_queue
            SET queue_status = 'pending',
                started_at = NULL,
                finished_at = NULL,
                error_message = NULL,
                updated_at = NOW()
            WHERE id = %s
              AND COALESCE(queue_status, '') IN ('pending', 'waiting', 'ready', 'failed', 'hold')
        """, (int(queue_id),))

    conn.commit()


def run_step104_worker(realtor_id):
    cmd = [
        PYTHON_EXE,
        "workers/publish_worker_step104.py",
        "--realtor-id",
        str(int(realtor_id)),
        "--limit",
        "1",
        "--no-db-lock",
    ]

    print("=" * 80)
    print("[RUN publish_worker_step104]")
    print("[CMD]", " ".join(cmd))
    print("=" * 80)

    started = time.time()

    proc = subprocess.run(
        cmd,
        cwd=str(BASE_DIR),
        text=True,
        encoding="utf-8",
        errors="replace",
        stdout=subprocess.PIPE,
        stderr=subprocess.STDOUT,
        timeout=3600,
    )

    elapsed = time.time() - started

    print(proc.stdout or "")
    print("=" * 80)
    print(f"[RUN DONE] returncode={proc.returncode} elapsed={elapsed:.1f}s")
    print("=" * 80)

    if proc.returncode != 0:
        raise RuntimeError(f"publish_worker_step104 failed: {proc.returncode}")

    return proc.stdout or ""


def main():
    parser = argparse.ArgumentParser()

    parser.add_argument("--queue-id", type=int, required=True)
    parser.add_argument("--require-test-account", action="store_true")
    parser.add_argument("--no-shift-today-published", action="store_true")
    parser.add_argument("--keep-auto-on", action="store_true", help="테스트 후 auto_publish_enabled를 원복하지 않음")

    args = parser.parse_args()

    conn = get_conn()

    held_rows = []
    original_auto = None
    auto_changed = False

    try:
        queue = fetch_queue(conn, args.queue_id)
        realtor_id = int(queue.get("realtor_id") or 0)

        if not realtor_id:
            raise RuntimeError(f"queue has no realtor_id: {args.queue_id}")

        realtor = fetch_realtor(conn, realtor_id)
        setting = fetch_publish_setting(conn, realtor_id)

        is_test_account = int(realtor.get("is_test_account") or 0) if "is_test_account" in realtor else 0
        publish_mode = clean_text(realtor.get("publish_mode") or "")
        original_auto = int(setting.get("auto_publish_enabled") or 0)

        print("=" * 80)
        print("[STEP105 PUBLISH ONE QUEUE START]")
        print("queue_id:", args.queue_id)
        print("realtor_id:", realtor_id)
        print("office_name:", realtor.get("office_name"))
        print("publish_mode:", publish_mode)
        print("is_test_account:", is_test_account)
        print("auto_publish_enabled:", original_auto)
        print("article_no:", queue.get("article_no"))
        print("queue_status:", queue.get("queue_status"))
        print("=" * 80)

        if args.require_test_account and not (is_test_account == 1 or publish_mode == "test"):
            raise RuntimeError("require-test-account enabled, but realtor is not test account")

        held_rows = hold_other_pending_queues(conn, realtor_id, args.queue_id)

        make_target_queue_pending(conn, args.queue_id)

        if not args.no_shift_today_published and (is_test_account == 1 or publish_mode == "test"):
            affected = shift_today_published_for_test(conn, realtor_id)
            print("[TEST SHIFT TODAY PUBLISHED]", f"affected={affected}")

        if original_auto == 0:
            set_auto_publish(conn, realtor_id, True)
            auto_changed = True
            print("[AUTO PUBLISH TEMP ON] publish only")

        # publish_worker_step104가 DB를 읽도록 커밋 후 실행
        try:
            conn.commit()
        except Exception:
            pass

        run_step104_worker(realtor_id)

        # 결과 확인
        conn.close()
        conn = get_conn()

        result = fetch_queue(conn, args.queue_id)

        print("[QUEUE RESULT]", json.dumps(result, ensure_ascii=False, default=str))

        print("=" * 80)
        print("[STEP105 PUBLISH ONE QUEUE DONE]")
        print("=" * 80)

    finally:
        try:
            if auto_changed and not args.keep_auto_on:
                set_auto_publish(conn, realtor_id, False)
                print("[AUTO PUBLISH RESTORED] OFF")
        except Exception as e:
            print("[AUTO RESTORE ERROR]", str(e))

        try:
            restore_held_queues(conn, held_rows)
        except Exception as e:
            print("[QUEUE RESTORE ERROR]", str(e))

        try:
            conn.close()
        except Exception:
            pass


if __name__ == "__main__":
    main()
