# -*- coding: utf-8 -*-
"""
STEP107-29: V2 Publish Queue Router

목적:
- blog_publish_queue를 서버 발행 / 클라이언트 발행 공용 Queue로 사용한다.
- 기존 V1 publish_worker는 건드리지 않는다.
- V2 worker 또는 클라이언트 셋업 프로그램이 안전하게 queue를 claim/release/complete 할 수 있게 한다.

추가 컬럼은 없으면 자동 생성한다.
"""

import argparse
import json
import socket
import uuid
from datetime import datetime

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_publish_queue_columns(conn):
    """
    V2 Queue 라우팅/점유 컬럼.
    기존 컬럼은 건드리지 않고, 없는 컬럼만 추가한다.
    """
    if not table_exists(conn, "blog_publish_queue"):
        raise RuntimeError("blog_publish_queue table not found")

    cols = table_columns(conn, "blog_publish_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 alters:
        with conn.cursor() as cur:
            for alter in alters:
                log("[V2 PUBLISH QUEUE ALTER]", alter)
                cur.execute(f"ALTER TABLE blog_publish_queue {alter}")
        conn.commit()

    return True


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


def set_queue_execution_target(queue_id=None, realtor_id=None, article_no=None, execution_target="server"):
    """
    execution_target:
    - server: 서버 발행 worker가 처리
    - client: 중개사 PC 클라이언트가 처리
    - auto: 나중에 정책에 따라 자동 선택
    """
    execution_target = clean_text(execution_target) or "server"
    if execution_target not in ["server", "client", "auto"]:
        raise ValueError("execution_target must be server/client/auto")

    conn = get_conn()
    try:
        ensure_v2_publish_queue_columns(conn)

        where = []
        params = []

        if queue_id:
            where.append("id=%s")
            params.append(int(queue_id))
        if realtor_id:
            where.append("realtor_id=%s")
            params.append(int(realtor_id))
        if article_no:
            where.append("article_no=%s")
            params.append(clean_text(article_no))

        if not where:
            raise ValueError("queue_id or realtor_id/article_no is required")

        with conn.cursor() as cur:
            cur.execute(f"""
                UPDATE blog_publish_queue
                SET execution_target=%s,
                    updated_at=NOW()
                WHERE {" AND ".join(where)}
            """, [execution_target] + params)
            affected = cur.rowcount

        conn.commit()
        result = {"ok": True, "affected": int(affected or 0), "execution_target": execution_target}
        log("[V2 PUBLISH QUEUE TARGET SET]", json.dumps(result, ensure_ascii=False))
        return result
    finally:
        conn.close()



def latest_snapshot_exists_sql(alias="pq"):
    """
    최신 SUCCESS Snapshot에 해당 article_no가 존재하는지 확인하는 EXISTS SQL.
    Collation 문제 방지를 위해 CAST(... AS UNSIGNED) 비교.
    """
    return f"""
        EXISTS (
            SELECT 1
            FROM blog_realtor_article_snapshot_items si
            INNER JOIN blog_realtor_article_snapshot_runs sr
                ON sr.id = si.run_id
            WHERE sr.realtor_id = {alias}.realtor_id
              AND sr.status = 'success'
              AND sr.id = (
                  SELECT id
                  FROM blog_realtor_article_snapshot_runs
                  WHERE realtor_id = {alias}.realtor_id
                    AND status = 'success'
                  ORDER BY snapshot_date DESC, id DESC
                  LIMIT 1
              )
              AND CAST(si.article_no AS UNSIGNED) = CAST({alias}.article_no AS UNSIGNED)
        )
    """


def has_snapshot_tables(conn):
    return (
        table_exists(conn, "blog_realtor_article_snapshot_runs")
        and table_exists(conn, "blog_realtor_article_snapshot_items")
    )


def mark_missing_snapshot_publish_queues_failed(realtor_id=None, execution_target=None, dry_run=False, limit=100):
    """
    STEP107-31:
    pending/ready/waiting 상태지만 최신 SUCCESS Snapshot에 없는 매물은 발행 대상에서 제외한다.
    이미 네이버에 없는 매물을 발행하지 않기 위한 안전장치다.

    queue_status는 기존 호환을 위해 failed로 둔다.
    worker_status='missing_current_snapshot'
    """
    conn = get_conn()
    try:
        ensure_v2_publish_queue_columns(conn)

        if not has_snapshot_tables(conn):
            result = {"ok": False, "error": "snapshot tables missing"}
            log("[V2 PUBLISH QUEUE MISSING SNAPSHOT CLEANUP]", json.dumps(result, ensure_ascii=False))
            return result

        where = [
            "COALESCE(pq.queue_status,'') IN ('pending','ready','waiting')",
            f"NOT {latest_snapshot_exists_sql('pq')}",
        ]
        params = []

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

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

        with conn.cursor() as cur:
            cur.execute(f"""
                SELECT pq.id, pq.realtor_id, pq.article_no, pq.queue_status, pq.execution_target
                FROM blog_publish_queue pq
                WHERE {" AND ".join(where)}
                ORDER BY pq.id ASC
                LIMIT %s
            """, params + [int(limit or 100)])
            rows = cur.fetchall() or []

        if dry_run:
            result = {
                "ok": True,
                "dry_run": True,
                "count": len(rows),
                "rows": rows,
            }
            log("[V2 PUBLISH QUEUE MISSING SNAPSHOT CLEANUP]", json.dumps(result, ensure_ascii=False, default=str))
            return result

        ids = [int(row.get("id")) for row in rows if row.get("id")]
        affected = 0

        if ids:
            placeholders = ",".join(["%s"] * len(ids))
            with conn.cursor() as cur:
                cur.execute(f"""
                    UPDATE blog_publish_queue
                    SET queue_status='failed',
                        worker_status='missing_current_snapshot',
                        error_message='최신 SUCCESS Snapshot에 없는 매물이므로 발행 제외',
                        last_claim_error='missing_current_snapshot',
                        updated_at=NOW()
                    WHERE id IN ({placeholders})
                """, ids)
                affected = cur.rowcount
            conn.commit()

        result = {
            "ok": True,
            "dry_run": False,
            "checked": len(rows),
            "affected": int(affected or 0),
            "ids": ids,
        }
        log("[V2 PUBLISH QUEUE MISSING SNAPSHOT CLEANUP]", json.dumps(result, ensure_ascii=False, default=str))
        return result

    finally:
        conn.close()


def claim_next_publish_queue(execution_target="server", realtor_id=None, worker_id=None, stale_minutes=30, require_current_snapshot=True):
    """
    다음 발행 Queue 1건을 원자적으로 점유한다.
    """
    execution_target = clean_text(execution_target) or "server"
    worker_id = clean_text(worker_id) or worker_name(f"{execution_target}-publisher")
    token = str(uuid.uuid4())

    conn = get_conn()
    try:
        ensure_v2_publish_queue_columns(conn)
        conn.begin()

        where = [
            "COALESCE(queue_status,'') 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)]

        # STEP107-31:
        # 최신 SUCCESS Snapshot에 없는 매물은 claim하지 않는다.
        # 과거 테스트/삭제매물 pending queue가 발행되는 것을 차단한다.
        if require_current_snapshot and has_snapshot_tables(conn):
            where.append(latest_snapshot_exists_sql("blog_publish_queue"))

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

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

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

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

            cur.execute("""
                UPDATE blog_publish_queue
                SET queue_status='processing',
                    worker_status='claimed',
                    claimed_by=%s,
                    claim_token=%s,
                    claimed_at=NOW(),
                    heartbeat_at=NOW(),
                    started_at=COALESCE(started_at, NOW()),
                    updated_at=NOW(),
                    last_claim_error=NULL
                WHERE id=%s
            """, (worker_id, token, queue_id))

        conn.commit()

        # UPDATE 이후 상태를 다시 읽어 claimed_at/heartbeat_at이 null로 보이는 혼선을 없앤다.
        with conn.cursor() as cur:
            cur.execute("""
                SELECT *
                FROM blog_publish_queue
                WHERE id=%s
                LIMIT 1
            """, (queue_id,))
            claimed_row = cur.fetchone() or row

        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 PUBLISH 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_publish_queue(queue_id, claim_token):
    conn = get_conn()
    try:
        ensure_v2_publish_queue_columns(conn)

        with conn.cursor() as cur:
            cur.execute("""
                UPDATE blog_publish_queue
                SET heartbeat_at=NOW(),
                    updated_at=NOW()
                WHERE id=%s
                  AND claim_token=%s
                  AND COALESCE(queue_status,'')='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 PUBLISH QUEUE HEARTBEAT]", json.dumps(result, ensure_ascii=False))
        return result
    finally:
        conn.close()


def complete_publish_queue(queue_id, claim_token, publish_result_url="", blog_id="", blog_post_no=""):
    conn = get_conn()
    try:
        ensure_v2_publish_queue_columns(conn)

        with conn.cursor() as cur:
            cur.execute("""
                UPDATE blog_publish_queue
                SET queue_status='published',
                    worker_status='success',
                    publish_result_url=%s,
                    blog_id=%s,
                    blog_post_no=%s,
                    finished_at=NOW(),
                    heartbeat_at=NOW(),
                    updated_at=NOW(),
                    error_message=NULL
                WHERE id=%s
                  AND claim_token=%s
                  AND COALESCE(queue_status,'')='processing'
            """, (
                clean_text(publish_result_url),
                clean_text(blog_id),
                clean_text(blog_post_no),
                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 PUBLISH QUEUE COMPLETE]", json.dumps(result, ensure_ascii=False))
        return result
    finally:
        conn.close()


def fail_publish_queue(queue_id, claim_token, error_message="", retry=True):
    conn = get_conn()
    try:
        ensure_v2_publish_queue_columns(conn)
        next_status = "failed" if retry else "error"

        with conn.cursor() as cur:
            cur.execute("""
                UPDATE blog_publish_queue
                SET queue_status=%s,
                    worker_status='failed',
                    finished_at=NOW(),
                    heartbeat_at=NOW(),
                    retry_count=COALESCE(retry_count,0)+1,
                    error_message=%s,
                    last_claim_error=%s,
                    updated_at=NOW()
                WHERE id=%s
                  AND claim_token=%s
            """, (
                next_status,
                clean_text(error_message)[:2000],
                clean_text(error_message)[:1000],
                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), "queue_status": next_status}
        log("[V2 PUBLISH QUEUE FAIL]", json.dumps(result, ensure_ascii=False))
        return result
    finally:
        conn.close()


def release_stale_claims(stale_minutes=30, execution_target=None):
    conn = get_conn()
    try:
        ensure_v2_publish_queue_columns(conn)

        where = [
            "COALESCE(queue_status,'')='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))

        with conn.cursor() as cur:
            cur.execute(f"""
                UPDATE blog_publish_queue
                SET queue_status='failed',
                    worker_status='stale_released',
                    retry_count=COALESCE(retry_count,0)+1,
                    last_claim_error='stale claim released',
                    claimed_by=NULL,
                    claim_token=NULL,
                    claimed_at=NULL,
                    heartbeat_at=NULL,
                    updated_at=NOW()
                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 PUBLISH QUEUE STALE RELEASE]", json.dumps(result, ensure_ascii=False))
        return result
    finally:
        conn.close()



def fetch_publish_queue_status(queue_id):
    conn = get_conn()
    try:
        ensure_v2_publish_queue_columns(conn)

        with conn.cursor() as cur:
            cur.execute("""
                SELECT *
                FROM blog_publish_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 PUBLISH QUEUE STATUS]", json.dumps(result, ensure_ascii=False, default=str))
        return result
    finally:
        conn.close()


def release_publish_queue_claim(queue_id, claim_token, reset_status="pending"):
    """
    테스트/클라이언트 중단/서버 중단 시 점유만 해제한다.
    기본은 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_publish_queue_columns(conn)

        with conn.cursor() as cur:
            cur.execute("""
                UPDATE blog_publish_queue
                SET queue_status=%s,
                    worker_status='claim_released',
                    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(queue_status,'')='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 PUBLISH QUEUE RELEASE CLAIM]", json.dumps(result, ensure_ascii=False))
        return result
    finally:
        conn.close()



def record_publish_success(queue_id=None, claim_token="", article_no=None, realtor_id=None, publish_result_url="", blog_id="", blog_post_no="", worker_id=""):
    """
    STEP107-34:
    서버/클라이언트 발행 성공 결과를 blog_publish_queue에 확정 기록한다.
    이후 deleted_article_sync가 이 Queue를 published_count로 인식한다.
    """
    publish_result_url = clean_text(publish_result_url)
    if not publish_result_url:
        raise ValueError("publish_result_url is required")

    conn = get_conn()
    try:
        ensure_v2_publish_queue_columns(conn)

        target = None

        if queue_id:
            with conn.cursor() as cur:
                cur.execute("SELECT * FROM blog_publish_queue WHERE id=%s LIMIT 1", (int(queue_id),))
                target = cur.fetchone()
        else:
            where = []
            params = []

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

            if not where:
                raise ValueError("queue_id 또는 realtor_id/article_no가 필요합니다.")

            with conn.cursor() as cur:
                cur.execute(f"""
                    SELECT *
                    FROM blog_publish_queue
                    WHERE {" AND ".join(where)}
                      AND COALESCE(queue_status,'') IN ('processing','pending','ready','waiting','failed')
                    ORDER BY
                        CASE
                            WHEN COALESCE(queue_status,'')='processing' THEN 1
                            WHEN COALESCE(queue_status,'') IN ('pending','ready','waiting') THEN 2
                            WHEN COALESCE(queue_status,'')='failed' THEN 3
                            ELSE 9
                        END,
                        id DESC
                    LIMIT 1
                """, params)
                target = cur.fetchone()

        if not target:
            result = {
                "ok": False,
                "reason": "queue_not_found",
                "queue_id": queue_id,
                "article_no": article_no,
                "realtor_id": realtor_id,
            }
            log("[V2 PUBLISH RESULT RECORD]", json.dumps(result, ensure_ascii=False))
            return result

        if claim_token:
            existing_token = clean_text(target.get("claim_token"))
            if existing_token and existing_token != clean_text(claim_token):
                result = {"ok": False, "reason": "claim_token_mismatch", "queue_id": target.get("id")}
                log("[V2 PUBLISH RESULT RECORD]", json.dumps(result, ensure_ascii=False))
                return result

        cols = table_columns(conn, "blog_publish_queue")

        sets = [
            "queue_status='published'",
            "worker_status='success'",
            "publish_result_url=%s",
            "finished_at=NOW()",
            "heartbeat_at=NOW()",
            "updated_at=NOW()",
            "error_message=NULL",
            "last_claim_error=NULL",
        ]
        params = [publish_result_url]

        if "blog_id" in cols:
            sets.append("blog_id=%s")
            params.append(clean_text(blog_id))
        if "blog_post_no" in cols:
            sets.append("blog_post_no=%s")
            params.append(clean_text(blog_post_no))
        if "worker_server" in cols:
            sets.append("worker_server=%s")
            params.append(clean_text(worker_id) or worker_name("publish-result"))

        with conn.cursor() as cur:
            cur.execute(f"""
                UPDATE blog_publish_queue
                SET {", ".join(sets)}
                WHERE id=%s
            """, params + [int(target.get("id"))])
            affected = cur.rowcount

        conn.commit()

        result = {
            "ok": bool(affected),
            "queue_id": int(target.get("id")),
            "article_no": target.get("article_no"),
            "realtor_id": target.get("realtor_id"),
            "publish_result_url": publish_result_url,
            "affected": int(affected or 0),
        }
        log("[V2 PUBLISH RESULT RECORD]", json.dumps(result, ensure_ascii=False, default=str))
        return result

    finally:
        conn.close()


def record_publish_failure(queue_id=None, claim_token="", article_no=None, realtor_id=None, error_message="", retry=True):
    """
    STEP107-34:
    서버/클라이언트 발행 실패 결과 기록.
    """
    conn = get_conn()
    try:
        ensure_v2_publish_queue_columns(conn)

        target = None

        if queue_id:
            with conn.cursor() as cur:
                cur.execute("SELECT * FROM blog_publish_queue WHERE id=%s LIMIT 1", (int(queue_id),))
                target = cur.fetchone()
        else:
            where = []
            params = []

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

            if not where:
                raise ValueError("queue_id 또는 realtor_id/article_no가 필요합니다.")

            with conn.cursor() as cur:
                cur.execute(f"""
                    SELECT *
                    FROM blog_publish_queue
                    WHERE {" AND ".join(where)}
                      AND COALESCE(queue_status,'') IN ('processing','pending','ready','waiting')
                    ORDER BY id DESC
                    LIMIT 1
                """, params)
                target = cur.fetchone()

        if not target:
            result = {
                "ok": False,
                "reason": "queue_not_found",
                "queue_id": queue_id,
                "article_no": article_no,
                "realtor_id": realtor_id,
            }
            log("[V2 PUBLISH FAILURE RECORD]", json.dumps(result, ensure_ascii=False))
            return result

        if claim_token:
            existing_token = clean_text(target.get("claim_token"))
            if existing_token and existing_token != clean_text(claim_token):
                result = {"ok": False, "reason": "claim_token_mismatch", "queue_id": target.get("id")}
                log("[V2 PUBLISH FAILURE RECORD]", json.dumps(result, ensure_ascii=False))
                return result

        next_status = "failed" if retry else "error"

        with conn.cursor() as cur:
            cur.execute("""
                UPDATE blog_publish_queue
                SET queue_status=%s,
                    worker_status='failed',
                    retry_count=COALESCE(retry_count,0)+1,
                    error_message=%s,
                    last_claim_error=%s,
                    finished_at=NOW(),
                    heartbeat_at=NOW(),
                    updated_at=NOW()
                WHERE id=%s
            """, (
                next_status,
                clean_text(error_message)[:2000],
                clean_text(error_message)[:1000],
                int(target.get("id")),
            ))
            affected = cur.rowcount

        conn.commit()

        result = {
            "ok": bool(affected),
            "queue_id": int(target.get("id")),
            "article_no": target.get("article_no"),
            "realtor_id": target.get("realtor_id"),
            "affected": int(affected or 0),
            "queue_status": next_status,
        }
        log("[V2 PUBLISH FAILURE RECORD]", 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-stale", action="store_true")
    parser.add_argument("--release-claim", action="store_true")
    parser.add_argument("--status", action="store_true")
    parser.add_argument("--record-success", action="store_true")
    parser.add_argument("--record-failure", action="store_true")
    parser.add_argument("--cleanup-missing-current", action="store_true")
    parser.add_argument("--dry-run", action="store_true")
    parser.add_argument("--no-current-snapshot-check", 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("--article-no", default="")
    parser.add_argument("--claim-token", default="")
    parser.add_argument("--worker-id", default="")
    parser.add_argument("--url", default="")
    parser.add_argument("--blog-id", default="")
    parser.add_argument("--blog-post-no", 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"])
    parser.add_argument("--limit", type=int, default=100)

    args = parser.parse_args()

    if args.ensure:
        conn = get_conn()
        try:
            ensure_v2_publish_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_claims(args.stale_minutes, args.execution_target), ensure_ascii=False, default=str))
        return

    if args.record_success:
        print("[RESULT]", json.dumps(record_publish_success(
            queue_id=args.queue_id,
            claim_token=args.claim_token,
            article_no=args.article_no,
            realtor_id=args.realtor_id,
            publish_result_url=args.url,
            blog_id=args.blog_id,
            blog_post_no=args.blog_post_no,
            worker_id=args.worker_id,
        ), ensure_ascii=False, default=str))
        return

    if args.record_failure:
        print("[RESULT]", json.dumps(record_publish_failure(
            queue_id=args.queue_id,
            claim_token=args.claim_token,
            article_no=args.article_no,
            realtor_id=args.realtor_id,
            error_message=args.error_message,
        ), ensure_ascii=False, default=str))
        return

    if args.cleanup_missing_current:
        print("[RESULT]", json.dumps(mark_missing_snapshot_publish_queues_failed(args.realtor_id, args.execution_target, args.dry_run, args.limit), ensure_ascii=False, default=str))
        return

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

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

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

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

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

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

    parser.print_help()


if __name__ == "__main__":
    main()
