# -*- coding: utf-8 -*-
"""
STEP107-32: V2 Private Queue Router

목적:
- blog_post_private_queue를 서버/클라이언트 비공개 처리 공용 Queue로 사용한다.
- Snapshot 비교로 생성된 비공개 Queue를 안전하게 claim/heartbeat/complete/fail 한다.
- V1 흐름은 건드리지 않는다.
"""

import argparse
import json
import socket
import uuid

from db import get_conn


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


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 ensure_v2_private_queue_columns(conn):
    if not table_exists(conn, "blog_post_private_queue"):
        raise RuntimeError("blog_post_private_queue table not found")

    cols = table_columns(conn, "blog_post_private_queue")
    alters = []

    if "execution_target" not in cols:
        alters.append("ADD COLUMN execution_target VARCHAR(20) NOT NULL DEFAULT 'server'")
    if "claimed_by" not in cols:
        alters.append("ADD COLUMN claimed_by VARCHAR(100) DEFAULT NULL")
    if "claim_token" not in cols:
        alters.append("ADD COLUMN claim_token VARCHAR(80) DEFAULT NULL")
    if "claimed_at" not in cols:
        alters.append("ADD COLUMN claimed_at DATETIME DEFAULT NULL")
    if "heartbeat_at" not in cols:
        alters.append("ADD COLUMN heartbeat_at DATETIME DEFAULT NULL")
    if "last_claim_error" not in cols:
        alters.append("ADD COLUMN last_claim_error VARCHAR(1000) DEFAULT NULL")
    if "private_result_url" not in cols:
        alters.append("ADD COLUMN private_result_url VARCHAR(1000) DEFAULT NULL")
    if "private_completed_at" not in cols:
        alters.append("ADD COLUMN private_completed_at DATETIME DEFAULT NULL")

    if alters:
        with conn.cursor() as cur:
            for alter in alters:
                log("[V2 PRIVATE QUEUE ALTER]", alter)
                cur.execute(f"ALTER TABLE blog_post_private_queue {alter}")
        conn.commit()

    return True


def status_column(conn):
    cols = table_columns(conn, "blog_post_private_queue")
    if "queue_status" in cols:
        return "queue_status"
    if "status" in cols:
        return "status"
    raise RuntimeError("blog_post_private_queue has no queue_status/status column")


def worker_name(prefix="v2-private-worker"):
    try:
        host = socket.gethostname()
    except Exception:
        host = "unknown-host"
    return f"{prefix}:{host}"


def claim_next_private_queue(execution_target="server", realtor_id=None, worker_id=None, stale_minutes=30):
    execution_target = clean_text(execution_target) or "server"
    worker_id = clean_text(worker_id) or worker_name(f"{execution_target}-private")
    token = str(uuid.uuid4())

    conn = get_conn()
    try:
        ensure_v2_private_queue_columns(conn)
        st_col = status_column(conn)

        conn.begin()

        where = [
            f"COALESCE({st_col},'') IN ('pending','ready','waiting','failed')",
            "COALESCE(execution_target,'server') IN (%s, 'auto')",
            """(
                claimed_at IS NULL
                OR heartbeat_at IS NULL
                OR heartbeat_at < DATE_SUB(NOW(), INTERVAL %s MINUTE)
            )""",
        ]
        params = [execution_target, int(stale_minutes or 30)]

        if realtor_id:
            where.append("realtor_id=%s")
            params.append(int(realtor_id))

        with conn.cursor() as cur:
            cur.execute(f"""
                SELECT *
                FROM blog_post_private_queue
                WHERE {" AND ".join(where)}
                ORDER BY
                    COALESCE(created_at, id) ASC,
                    id ASC
                LIMIT 1
                FOR UPDATE
            """, params)
            row = cur.fetchone()

            if not row:
                conn.commit()
                result = {"ok": True, "claimed": False, "message": "no pending private queue"}
                log("[V2 PRIVATE QUEUE CLAIM]", json.dumps(result, ensure_ascii=False))
                return result

            queue_id = int(row.get("id"))

            cur.execute(f"""
                UPDATE blog_post_private_queue
                SET {st_col}='processing',
                    claimed_by=%s,
                    claim_token=%s,
                    claimed_at=NOW(),
                    heartbeat_at=NOW(),
                    updated_at=NOW(),
                    last_claim_error=NULL
                WHERE id=%s
            """, (worker_id, token, queue_id))

            cur.execute("""
                SELECT *
                FROM blog_post_private_queue
                WHERE id=%s
                LIMIT 1
            """, (queue_id,))
            claimed_row = cur.fetchone() or row

        conn.commit()

        result = {
            "ok": True,
            "claimed": True,
            "queue_id": queue_id,
            "article_no": claimed_row.get("article_no"),
            "realtor_id": claimed_row.get("realtor_id"),
            "execution_target": execution_target,
            "claim_token": token,
            "claimed_by": worker_id,
            "row": claimed_row,
        }
        log("[V2 PRIVATE QUEUE CLAIM]", json.dumps(result, ensure_ascii=False, default=str))
        return result

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


def heartbeat_private_queue(queue_id, claim_token):
    conn = get_conn()
    try:
        ensure_v2_private_queue_columns(conn)
        st_col = status_column(conn)

        with conn.cursor() as cur:
            cur.execute(f"""
                UPDATE blog_post_private_queue
                SET heartbeat_at=NOW(),
                    updated_at=NOW()
                WHERE id=%s
                  AND claim_token=%s
                  AND COALESCE({st_col},'')='processing'
            """, (int(queue_id), clean_text(claim_token)))
            affected = cur.rowcount

        conn.commit()
        result = {"ok": bool(affected), "queue_id": int(queue_id), "affected": int(affected or 0)}
        log("[V2 PRIVATE QUEUE HEARTBEAT]", json.dumps(result, ensure_ascii=False))
        return result
    finally:
        conn.close()


def complete_private_queue(queue_id, claim_token, private_result_url=""):
    conn = get_conn()
    try:
        ensure_v2_private_queue_columns(conn)
        st_col = status_column(conn)

        sets = [
            f"{st_col}='done'",
            "private_result_url=%s",
            "private_completed_at=NOW()",
            "heartbeat_at=NOW()",
            "updated_at=NOW()",
            "last_claim_error=NULL",
        ]

        with conn.cursor() as cur:
            cur.execute(f"""
                UPDATE blog_post_private_queue
                SET {", ".join(sets)}
                WHERE id=%s
                  AND claim_token=%s
                  AND COALESCE({st_col},'')='processing'
            """, (
                clean_text(private_result_url),
                int(queue_id),
                clean_text(claim_token),
            ))
            affected = cur.rowcount

        conn.commit()
        result = {"ok": bool(affected), "queue_id": int(queue_id), "affected": int(affected or 0)}
        log("[V2 PRIVATE QUEUE COMPLETE]", json.dumps(result, ensure_ascii=False))
        return result
    finally:
        conn.close()


def fail_private_queue(queue_id, claim_token, error_message="", retry=True):
    conn = get_conn()
    try:
        ensure_v2_private_queue_columns(conn)
        st_col = status_column(conn)
        next_status = "failed" if retry else "error"
        cols = table_columns(conn, "blog_post_private_queue")

        sets = [
            f"{st_col}=%s",
            "heartbeat_at=NOW()",
            "last_claim_error=%s",
            "updated_at=NOW()",
        ]

        params = [
            next_status,
            clean_text(error_message)[:1000],
        ]

        if "error_message" in cols:
            sets.append("error_message=%s")
            params.append(clean_text(error_message)[:2000])

        if "retry_count" in cols:
            sets.append("retry_count=COALESCE(retry_count,0)+1")

        with conn.cursor() as cur:
            cur.execute(f"""
                UPDATE blog_post_private_queue
                SET {", ".join(sets)}
                WHERE id=%s
                  AND claim_token=%s
            """, params + [int(queue_id), clean_text(claim_token)])
            affected = cur.rowcount

        conn.commit()
        result = {"ok": bool(affected), "queue_id": int(queue_id), "affected": int(affected or 0), "status": next_status}
        log("[V2 PRIVATE QUEUE FAIL]", json.dumps(result, ensure_ascii=False))
        return result
    finally:
        conn.close()


def release_private_queue_claim(queue_id, claim_token, reset_status="pending"):
    reset_status = clean_text(reset_status) or "pending"
    if reset_status not in ["pending", "failed", "ready", "waiting"]:
        raise ValueError("reset_status must be pending/failed/ready/waiting")

    conn = get_conn()
    try:
        ensure_v2_private_queue_columns(conn)
        st_col = status_column(conn)

        with conn.cursor() as cur:
            cur.execute(f"""
                UPDATE blog_post_private_queue
                SET {st_col}=%s,
                    claimed_by=NULL,
                    claim_token=NULL,
                    claimed_at=NULL,
                    heartbeat_at=NULL,
                    last_claim_error=NULL,
                    updated_at=NOW()
                WHERE id=%s
                  AND claim_token=%s
                  AND COALESCE({st_col},'')='processing'
            """, (
                reset_status,
                int(queue_id),
                clean_text(claim_token),
            ))
            affected = cur.rowcount

        conn.commit()
        result = {"ok": bool(affected), "queue_id": int(queue_id), "affected": int(affected or 0), "reset_status": reset_status}
        log("[V2 PRIVATE QUEUE RELEASE CLAIM]", json.dumps(result, ensure_ascii=False))
        return result
    finally:
        conn.close()


def release_stale_private_claims(stale_minutes=30, execution_target=None):
    conn = get_conn()
    try:
        ensure_v2_private_queue_columns(conn)
        st_col = status_column(conn)

        where = [
            f"COALESCE({st_col},'')='processing'",
            "heartbeat_at < DATE_SUB(NOW(), INTERVAL %s MINUTE)",
        ]
        params = [int(stale_minutes or 30)]

        if execution_target:
            where.append("COALESCE(execution_target,'server')=%s")
            params.append(clean_text(execution_target))

        cols = table_columns(conn, "blog_post_private_queue")
        sets = [
            f"{st_col}='failed'",
            "claimed_by=NULL",
            "claim_token=NULL",
            "claimed_at=NULL",
            "heartbeat_at=NULL",
            "last_claim_error='stale claim released'",
            "updated_at=NOW()",
        ]

        if "retry_count" in cols:
            sets.append("retry_count=COALESCE(retry_count,0)+1")

        with conn.cursor() as cur:
            cur.execute(f"""
                UPDATE blog_post_private_queue
                SET {", ".join(sets)}
                WHERE {" AND ".join(where)}
            """, params)
            affected = cur.rowcount

        conn.commit()
        result = {"ok": True, "released": int(affected or 0), "stale_minutes": int(stale_minutes or 30)}
        log("[V2 PRIVATE QUEUE STALE RELEASE]", json.dumps(result, ensure_ascii=False))
        return result
    finally:
        conn.close()


def fetch_private_queue_status(queue_id):
    conn = get_conn()
    try:
        ensure_v2_private_queue_columns(conn)

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

        result = {"ok": row is not None, "queue_id": int(queue_id), "row": row}
        log("[V2 PRIVATE QUEUE STATUS]", json.dumps(result, ensure_ascii=False, default=str))
        return result
    finally:
        conn.close()


def main():
    parser = argparse.ArgumentParser()
    parser.add_argument("--ensure", action="store_true")
    parser.add_argument("--claim", action="store_true")
    parser.add_argument("--heartbeat", action="store_true")
    parser.add_argument("--complete", action="store_true")
    parser.add_argument("--fail", action="store_true")
    parser.add_argument("--release-claim", action="store_true")
    parser.add_argument("--release-stale", action="store_true")
    parser.add_argument("--status", action="store_true")

    parser.add_argument("--execution-target", default="server", choices=["server", "client", "auto"])
    parser.add_argument("--realtor-id", type=int, default=None)
    parser.add_argument("--queue-id", type=int, default=None)
    parser.add_argument("--claim-token", default="")
    parser.add_argument("--worker-id", default="")
    parser.add_argument("--url", default="")
    parser.add_argument("--error-message", default="")
    parser.add_argument("--stale-minutes", type=int, default=30)
    parser.add_argument("--reset-status", default="pending", choices=["pending", "failed", "ready", "waiting"])

    args = parser.parse_args()

    if args.ensure:
        conn = get_conn()
        try:
            ensure_v2_private_queue_columns(conn)
            print("[RESULT]", json.dumps({"ok": True, "ensured": True}, ensure_ascii=False))
        finally:
            conn.close()
        return

    if args.release_stale:
        print("[RESULT]", json.dumps(release_stale_private_claims(args.stale_minutes, args.execution_target), ensure_ascii=False, default=str))
        return

    if args.status:
        print("[RESULT]", json.dumps(fetch_private_queue_status(args.queue_id), ensure_ascii=False, default=str))
        return

    if args.claim:
        print("[RESULT]", json.dumps(claim_next_private_queue(args.execution_target, args.realtor_id, args.worker_id, args.stale_minutes), ensure_ascii=False, default=str))
        return

    if args.heartbeat:
        print("[RESULT]", json.dumps(heartbeat_private_queue(args.queue_id, args.claim_token), ensure_ascii=False, default=str))
        return

    if args.complete:
        print("[RESULT]", json.dumps(complete_private_queue(args.queue_id, args.claim_token, args.url), ensure_ascii=False, default=str))
        return

    if args.fail:
        print("[RESULT]", json.dumps(fail_private_queue(args.queue_id, args.claim_token, args.error_message), ensure_ascii=False, default=str))
        return

    if args.release_claim:
        print("[RESULT]", json.dumps(release_private_queue_claim(args.queue_id, args.claim_token, args.reset_status), ensure_ascii=False, default=str))
        return

    parser.print_help()


if __name__ == "__main__":
    main()
