# -*- coding: utf-8 -*-
"""
STEP107-20 Snapshot 기반 삭제매물 Sync

기준:
- 삭제매물 판단은 blog_realtor_articles 현재 cache가 아니라
  blog_realtor_article_snapshot_runs/items 의 SUCCESS Snapshot끼리 비교한다.
- 오늘 SUCCESS Snapshot과 이전 SUCCESS Snapshot을 비교해,
  이전에는 있었고 오늘은 없는 article_no만 삭제매물 후보로 본다.
- 발행된 글이 있는 article_no만 blog_post_private_queue에 넣는다.
"""

import argparse
import json
from datetime import datetime, date

from core.db import get_conn


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() or []
        return {str(row.get("Field") or "").strip() for row in rows if row.get("Field")}
    except Exception:
        return set()


def get_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() or {}


def get_latest_success_snapshot_run(conn, realtor_id, before_run_id=None, snapshot_date=None):
    where = ["realtor_id=%s", "status='success'"]
    params = [int(realtor_id)]

    if before_run_id is not None:
        where.append("id < %s")
        params.append(int(before_run_id))

    if snapshot_date is not None:
        where.append("snapshot_date <= %s")
        params.append(snapshot_date)

    sql = f"""
        SELECT *
        FROM blog_realtor_article_snapshot_runs
        WHERE {' AND '.join(where)}
        ORDER BY snapshot_date DESC, id DESC
        LIMIT 1
    """

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


def get_previous_success_snapshot_run(conn, realtor_id, current_run):
    if not current_run:
        return None

    return get_latest_success_snapshot_run(
        conn=conn,
        realtor_id=realtor_id,
        before_run_id=current_run.get("id"),
        snapshot_date=current_run.get("snapshot_date"),
    )


def fetch_snapshot_article_nos(conn, run_id):
    with conn.cursor() as cur:
        cur.execute("""
            SELECT article_no
            FROM blog_realtor_article_snapshot_items
            WHERE run_id=%s
              AND COALESCE(article_no,'') <> ''
        """, (int(run_id),))
        rows = cur.fetchall() or []

    return {clean_text(row.get("article_no")) for row in rows if clean_text(row.get("article_no"))}


def fetch_published_by_article_nos(conn, realtor_id, article_nos):
    article_nos = [clean_text(x) for x in article_nos if clean_text(x)]
    if not article_nos:
        return {}

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

    with conn.cursor() as cur:
        cur.execute(f"""
            SELECT *
            FROM blog_publish_queue
            WHERE realtor_id=%s
              AND article_no IN ({placeholders})
              AND COALESCE(queue_status,'') IN ('published','success','done')
              AND COALESCE(publish_result_url,'') <> ''
            ORDER BY id DESC
        """, [int(realtor_id)] + article_nos)
        rows = cur.fetchall() or []

    result = {}
    for row in rows:
        article_no = clean_text(row.get("article_no"))
        if article_no and article_no not in result:
            result[article_no] = row
    return result


def private_queue_status_expr(conn):
    cols = table_columns(conn, "blog_post_private_queue")
    if "queue_status" in cols:
        return "queue_status"
    if "status" in cols:
        return "status"
    return None


def private_queue_exists(conn, realtor_id, article_no):
    if not table_exists(conn, "blog_post_private_queue"):
        return False

    with conn.cursor() as cur:
        cur.execute("""
            SELECT id
            FROM blog_post_private_queue
            WHERE realtor_id=%s
              AND article_no=%s
            LIMIT 1
        """, (int(realtor_id), clean_text(article_no)))
        return cur.fetchone() is not None


def insert_private_queue(conn, realtor_id, article_no, published_row, current_run, previous_run):
    """
    운영 DB 컬럼 차이를 고려해 존재 컬럼만 INSERT한다.
    """
    if not table_exists(conn, "blog_post_private_queue"):
        raise RuntimeError("blog_post_private_queue 테이블이 없습니다.")

    cols = table_columns(conn, "blog_post_private_queue")
    now = datetime.now()

    data = {
        "realtor_id": int(realtor_id),
        "article_no": clean_text(article_no),
        "publish_queue_id": published_row.get("id"),
        "draft_id": published_row.get("draft_id"),
        "blog_url": published_row.get("publish_result_url"),
        "publish_result_url": published_row.get("publish_result_url"),
        "queue_status": "pending",
        "status": "pending",
        "reason": "snapshot_missing",
        "error_message": None,
        "snapshot_run_id": current_run.get("id") if current_run else None,
        "previous_snapshot_run_id": previous_run.get("id") if previous_run else None,
        "created_at": now,
        "updated_at": now,
    }

    insert_data = {}
    for k, v in data.items():
        if k in cols and v is not None:
            insert_data[k] = v

    if not insert_data:
        raise RuntimeError("blog_post_private_queue에 INSERT 가능한 컬럼이 없습니다.")

    fields = list(insert_data.keys())
    placeholders = ",".join(["%s"] * len(fields))

    with conn.cursor() as cur:
        cur.execute(
            f"""
            INSERT INTO blog_post_private_queue
                ({",".join(fields)})
            VALUES
                ({placeholders})
            """,
            [insert_data[k] for k in fields],
        )

    return True


def sync_deleted_articles_by_snapshot(realtor_id, dry_run=False, current_run_id=None, limit=None):
    conn = get_conn()

    try:
        if not table_exists(conn, "blog_realtor_article_snapshot_runs") or not table_exists(conn, "blog_realtor_article_snapshot_items"):
            return {
                "ok": False,
                "error": "snapshot 테이블이 없습니다.",
                "realtor_id": realtor_id,
                "dry_run": dry_run,
            }

        realtor = get_realtor(conn, realtor_id)
        office_name = realtor.get("office_name") or realtor.get("name") or ""

        if current_run_id:
            with conn.cursor() as cur:
                cur.execute("""
                    SELECT *
                    FROM blog_realtor_article_snapshot_runs
                    WHERE id=%s AND realtor_id=%s AND status='success'
                    LIMIT 1
                """, (int(current_run_id), int(realtor_id)))
                current_run = cur.fetchone()
        else:
            current_run = get_latest_success_snapshot_run(conn, realtor_id)

        previous_run = get_previous_success_snapshot_run(conn, realtor_id, current_run)

        if not current_run:
            return {
                "ok": False,
                "error": "현재 SUCCESS snapshot이 없습니다.",
                "realtor_id": realtor_id,
                "office_name": office_name,
                "dry_run": dry_run,
            }

        if not previous_run:
            return {
                "ok": True,
                "realtor_id": realtor_id,
                "office_name": office_name,
                "current_run_id": current_run.get("id"),
                "previous_run_id": None,
                "message": "이전 SUCCESS snapshot이 없어 삭제비교를 건너뜁니다.",
                "deleted_count": 0,
                "private_queue_created": 0,
                "private_queue_skipped": 0,
                "dry_run": dry_run,
            }

        current_set = fetch_snapshot_article_nos(conn, current_run.get("id"))
        previous_set = fetch_snapshot_article_nos(conn, previous_run.get("id"))

        missing_article_nos = sorted(list(previous_set - current_set), key=lambda x: int(x) if str(x).isdigit() else 0, reverse=True)

        if limit is not None:
            missing_article_nos = missing_article_nos[:int(limit)]

        published_map = fetch_published_by_article_nos(conn, realtor_id, missing_article_nos)

        created = []
        skipped = []
        errors = []

        for article_no in missing_article_nos:
            published_row = published_map.get(article_no)

            if not published_row:
                skipped.append({
                    "article_no": article_no,
                    "reason": "not_published_or_no_blog_url",
                })
                continue

            if private_queue_exists(conn, realtor_id, article_no):
                skipped.append({
                    "article_no": article_no,
                    "reason": "already_private_queued",
                })
                continue

            if dry_run:
                created.append({
                    "article_no": article_no,
                    "blog_url": published_row.get("publish_result_url"),
                    "dry_run": True,
                })
                continue

            try:
                insert_private_queue(
                    conn=conn,
                    realtor_id=realtor_id,
                    article_no=article_no,
                    published_row=published_row,
                    current_run=current_run,
                    previous_run=previous_run,
                )
                created.append({
                    "article_no": article_no,
                    "blog_url": published_row.get("publish_result_url"),
                })
            except Exception as e:
                errors.append({
                    "article_no": article_no,
                    "error": str(e),
                })

        if not dry_run:
            conn.commit()

        result = {
            "ok": len(errors) == 0,
            "realtor_id": int(realtor_id),
            "office_name": office_name,
            "current_run_id": current_run.get("id"),
            "current_snapshot_date": str(current_run.get("snapshot_date")),
            "current_count": len(current_set),
            "previous_run_id": previous_run.get("id"),
            "previous_snapshot_date": str(previous_run.get("snapshot_date")),
            "previous_count": len(previous_set),
            "deleted_count": len(missing_article_nos),
            "private_queue_created": len(created),
            "private_queue_skipped": len(skipped),
            "errors": errors,
            "created": created,
            "skipped": skipped,
            "dry_run": dry_run,
        }

        print("[V2 SNAPSHOT DELETED SYNC]", json.dumps(result, ensure_ascii=False, default=str))
        return result

    finally:
        try:
            conn.close()
        except Exception:
            pass


def main():
    parser = argparse.ArgumentParser()
    parser.add_argument("--realtor-id", type=int, required=True)
    parser.add_argument("--current-run-id", type=int, default=None)
    parser.add_argument("--dry-run", action="store_true")
    parser.add_argument("--limit", type=int, default=None)
    args = parser.parse_args()

    result = sync_deleted_articles_by_snapshot(
        realtor_id=args.realtor_id,
        current_run_id=args.current_run_id,
        dry_run=args.dry_run,
        limit=args.limit,
    )
    print("[RESULT]", json.dumps(result, ensure_ascii=False, default=str))


if __name__ == "__main__":
    main()
