# -*- coding: utf-8 -*-
"""
STEP104-04: process_one_realtor_pipeline.py

중개사 1명 기준 원스톱 테스트 파이프라인.

흐름:
  1) publish_mode / test 계정 확인
  2) 세션 DB 상태 확인
  3) 수집 후보 등록(schedule_daily_articles)
  4) 상세수집(naver_response_article_fetcher)
  5) 초안 생성(generate_blog_drafts)
  6) publish_mode 분기
     - test/server_auto: publish_worker_step104.py 실행
     - client_assist: 발행대기 큐까지만 확인하고 종료
"""

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_realtor(conn, realtor_id):
    with conn.cursor() as cur:
        cur.execute("""
            SELECT *
            FROM blog_realtors
            WHERE id = %s
            LIMIT 1
        """, (int(realtor_id),))
        return cur.fetchone()


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")
    if "realtor_id" not in columns:
        return

    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}
    keys = list(filtered.keys())

    if not keys:
        return

    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 set_realtor_publish_mode(conn, realtor_id, mode):
    if not mode:
        return

    columns = table_columns(conn, "blog_realtors")
    if "publish_mode" not in columns:
        print("[WARN] blog_realtors.publish_mode column not found")
        return

    with conn.cursor() as cur:
        cur.execute("""
            UPDATE blog_realtors
            SET publish_mode = %s,
                updated_at = NOW()
            WHERE id = %s
        """, (str(mode), int(realtor_id)))

    conn.commit()


def check_session_db(conn, realtor_id):
    if not table_exists(conn, "blog_naver_sessions"):
        return {
            "exists": False,
            "ok": None,
            "message": "blog_naver_sessions table not exists",
        }

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

    if not row:
        return {
            "exists": False,
            "ok": False,
            "message": "session row not found",
        }

    status = clean_text(
        row.get("status")
        or row.get("naver_session_status")
        or row.get("login_status")
        or ""
    )

    session_file = clean_text(row.get("session_file") or "")

    ok_statuses = {
        "login_ok",
        "linked",
        "ok",
        "active",
        "unknown",
    }

    return {
        "exists": True,
        "ok": status in ok_statuses,
        "status": status,
        "session_file": session_file,
        "message": f"status={status}, session_file={session_file}",
    }


def count_pending_work(conn, realtor_id):
    if not table_exists(conn, "blog_article_work_queue"):
        return 0

    with conn.cursor() as cur:
        cur.execute("""
            SELECT COUNT(*) AS cnt
            FROM blog_article_work_queue
            WHERE realtor_id = %s
              AND work_type = 'new_article'
              AND work_status IN ('pending', 'processing')
        """, (int(realtor_id),))
        row = cur.fetchone() or {}

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


def count_collected_articles(conn, realtor_id):
    if not table_exists(conn, "blog_realtor_articles"):
        return 0

    with conn.cursor() as cur:
        cur.execute("""
            SELECT COUNT(*) AS cnt
            FROM blog_realtor_articles
            WHERE realtor_id = %s
              AND COALESCE(collect_status, '') = 'collected'
              AND COALESCE(detail_collected, 0) = 1
        """, (int(realtor_id),))
        row = cur.fetchone() or {}

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


def count_ready_drafts(conn, realtor_id):
    if not table_exists(conn, "blog_article_drafts"):
        return 0

    with conn.cursor() as cur:
        cur.execute("""
            SELECT COUNT(*) AS cnt
            FROM blog_article_drafts
            WHERE realtor_id = %s
              AND COALESCE(status, 'ready') IN ('ready', 'approved', 'pending')
        """, (int(realtor_id),))
        row = cur.fetchone() or {}

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


def count_pending_publish(conn, realtor_id):
    if not table_exists(conn, "blog_publish_queue"):
        return 0

    with conn.cursor() as cur:
        cur.execute("""
            SELECT COUNT(*) AS cnt
            FROM blog_publish_queue
            WHERE realtor_id = %s
              AND COALESCE(queue_status, '') IN ('pending', 'waiting', 'ready', 'failed')
        """, (int(realtor_id),))
        row = cur.fetchone() or {}

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


def fetch_latest_publish_queue(conn, realtor_id):
    if not table_exists(conn, "blog_publish_queue"):
        return None

    with conn.cursor() as cur:
        cur.execute("""
            SELECT *
            FROM blog_publish_queue
            WHERE realtor_id = %s
            ORDER BY id DESC
            LIMIT 1
        """, (int(realtor_id),))
        return cur.fetchone()


def shift_today_published_for_test(conn, realtor_id):
    if not table_exists(conn, "blog_publish_queue"):
        return 0

    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 run_cmd(label, cmd, timeout=None):
    print("=" * 80)
    print(f"[RUN START] {label}")
    print("[CMD]", " ".join([str(x) for x in cmd]))
    print("=" * 80)

    started = time.time()

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

    elapsed = time.time() - started

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

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

    return proc.stdout or ""


def print_counts(conn, realtor_id, label):
    print("-" * 80)
    print(f"[COUNTS] {label}")
    print("pending_work:", count_pending_work(conn, realtor_id))
    print("collected_articles:", count_collected_articles(conn, realtor_id))
    print("ready_drafts:", count_ready_drafts(conn, realtor_id))
    print("pending_publish:", count_pending_publish(conn, realtor_id))
    print("-" * 80)


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

    parser.add_argument("--realtor-id", type=int, required=True)
    parser.add_argument("--mode", choices=["auto", "server", "client", "draft-only"], default="auto")
    parser.add_argument("--set-publish-mode", default="")
    parser.add_argument("--require-test-mode", action="store_true")
    parser.add_argument("--allow-test-override", action="store_true")
    parser.add_argument("--shift-today-published", action="store_true")
    parser.add_argument("--skip-schedule", action="store_true")
    parser.add_argument("--skip-fetch", action="store_true")
    parser.add_argument("--skip-draft", action="store_true")
    parser.add_argument("--skip-publish", action="store_true")
    parser.add_argument("--draft-buffer", type=int, default=1)
    parser.add_argument("--replenish-count", type=int, default=1)
    parser.add_argument("--fetch-limit", type=int, default=3)
    parser.add_argument("--draft-limit", type=int, default=1)
    parser.add_argument("--publish-limit", type=int, default=1)

    args = parser.parse_args()

    realtor_id = int(args.realtor_id)
    conn = get_conn()

    schedule_auto_overridden = False

    try:
        if args.set_publish_mode:
            set_realtor_publish_mode(conn, realtor_id, args.set_publish_mode)

        realtor = fetch_realtor(conn, realtor_id)
        if not realtor:
            raise RuntimeError(f"realtor not found: {realtor_id}")

        setting = fetch_publish_setting(conn, realtor_id)

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

        print("=" * 80)
        print("[STEP104 ONE REALTOR PIPELINE START]")
        print("TIME:", datetime.now().strftime("%Y-%m-%d %H:%M:%S"))
        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:", auto_enabled)
        print("=" * 80)

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

        session = check_session_db(conn, realtor_id)
        print("[SESSION DB CHECK]", json.dumps(session, ensure_ascii=False, default=str))

        if session.get("exists") and session.get("ok") is False:
            print("[SESSION DB WARN]", f"session db not ready but publish_worker will verify login: {session.get('message')}")

        if args.shift_today_published:
            shifted = shift_today_published_for_test(conn, realtor_id)
            print(f"[TEST SHIFT TODAY PUBLISHED] affected={shifted}")

        print_counts(conn, realtor_id, "before")

        if not args.skip_schedule:
            if auto_enabled == 0 and args.allow_test_override:
                set_auto_publish(conn, realtor_id, True)
                schedule_auto_overridden = True
                print("[AUTO PUBLISH TEMP ON] schedule_daily_articles only")

            try:
                run_cmd(
                    "schedule_daily_articles",
                    [
                        PYTHON_EXE,
                        "services/naver_response_article_fetcher.py",
                        "--limit",
                        str(args.fetch_limit),
                        "--no-db-lock",
                    ],
                    timeout=900,
                )
            finally:
                if schedule_auto_overridden:
                    set_auto_publish(conn, realtor_id, False)
                    print("[AUTO PUBLISH RESTORED] OFF")

            conn.close()
            conn = get_conn()
            print_counts(conn, realtor_id, "after schedule")

        if not args.skip_fetch:
            run_cmd(
                "naver_response_article_fetcher",
                [
                    PYTHON_EXE,
                    "services/naver_response_article_fetcher.py",
                    "--realtor-id",
                    str(realtor_id),
                    "--limit",
                    str(args.fetch_limit),
                    "--no-db-lock",
                ],
                timeout=1800,
            )
            conn.close()
            conn = get_conn()
            print_counts(conn, realtor_id, "after detail fetch")

        if not args.skip_draft:
            run_cmd(
                "generate_blog_drafts",
                [
                    PYTHON_EXE,
                    "jobs/generate_blog_drafts.py",
                    "--realtor-id",
                    str(realtor_id),
                    "--limit",
                    str(args.draft_limit),
                    "--draft-buffer",
                    str(args.draft_buffer),
                    "--replenish-count",
                    str(args.replenish_count),
                    "--no-db-lock",
                ],
                timeout=2400,
            )
            conn.close()
            conn = get_conn()
            print_counts(conn, realtor_id, "after draft")

        effective_mode = args.mode

        if effective_mode == "auto":
            if publish_mode == "client_assist":
                effective_mode = "client"
            elif publish_mode == "manual":
                effective_mode = "draft-only"
            else:
                effective_mode = "server"

        print("[EFFECTIVE MODE]", effective_mode)

        if args.skip_publish:
            print("[SKIP PUBLISH] requested")
        elif effective_mode == "client":
            pending = count_pending_publish(conn, realtor_id)
            print("[CLIENT MODE READY]", f"pending_publish={pending}")
            print("[CLIENT MODE] 서버는 수집/초안/큐까지만 처리하고 발행은 클라이언트가 수행합니다.")
        elif effective_mode == "draft-only":
            print("[DRAFT ONLY MODE] publish skipped")
        elif effective_mode == "server":
            run_cmd(
                "publish_worker_step104",
                [
                    PYTHON_EXE,
                    "workers/publish_worker_step104.py",
                    "--realtor-id",
                    str(realtor_id),
                    "--limit",
                    str(args.publish_limit),
                    "--no-db-lock",
                ],
                timeout=3600,
            )
            conn.close()
            conn = get_conn()
            print_counts(conn, realtor_id, "after publish")

            latest = fetch_latest_publish_queue(conn, realtor_id)
            print("[LATEST QUEUE]", json.dumps(latest or {}, ensure_ascii=False, default=str))
        else:
            raise RuntimeError(f"unknown effective mode: {effective_mode}")

        print("=" * 80)
        print("[STEP104 ONE REALTOR PIPELINE DONE]")
        print("=" * 80)

    finally:
        try:
            if schedule_auto_overridden:
                set_auto_publish(conn, realtor_id, False)
                print("[AUTO PUBLISH RESTORED IN FINALLY] OFF")
        except Exception:
            pass

        try:
            conn.close()
        except Exception:
            pass


if __name__ == "__main__":
    main()
