# -*- coding: utf-8 -*-

"""
check_removed_articles.py

역할:
- 유료/상품/자동발행 조건을 충족하는 중개사만 대상으로 한다.
- 이미 네이버 블로그에 발행 완료된 article_no를 기준으로 한다.
- 현재 수집 파이프라인이 갱신한 blog_realtor_articles.last_seen_at을 기준으로
  네이버 부동산 전체매물에 없는 것으로 보이는 매물을 suspected_removed / removed_confirmed로 변경한다.
- removed_confirmed가 되면 blog_post_private_queue에 비공개 대기 큐를 생성한다.

중요:
- 이 worker는 블로그 비공개 실행을 하지 않는다.
- 실제 비공개 처리는 이후 publish_worker.py에서 같은 네이버 세션을 재사용해 처리한다.

사용 예:
python workers/check_removed_articles.py --dry-run
python workers/check_removed_articles.py --limit 100
python workers/check_removed_articles.py --realtor-id 2 --dry-run
python workers/check_removed_articles.py --confirm-count 2 --seen-hours 30
"""

import os
import sys
import json
import socket
import argparse
import traceback
from datetime import datetime

CURRENT_DIR = os.path.dirname(os.path.abspath(__file__))
ROOT_DIR = os.path.dirname(CURRENT_DIR)

if ROOT_DIR not in sys.path:
    sys.path.append(ROOT_DIR)

from db import get_conn


def log(message):
    print(f"[{datetime.now().strftime('%Y-%m-%d %H:%M:%S')}] {message}")


def safe_json_dumps(data):
    try:
        return json.dumps(data or {}, ensure_ascii=False, default=str)
    except Exception:
        return "{}"


def worker_token(prefix):
    try:
        host = socket.gethostname()
    except Exception:
        host = "unknown-host"
    return f"{prefix}:{host}:{os.getpid()}:{datetime.now().strftime('%Y%m%d%H%M%S')}"


def acquire_db_lock(lock_name, ttl_minutes=180):
    owner_token = worker_token(lock_name)
    conn = None

    try:
        conn = get_conn()

        with conn.cursor() as cur:
            cur.execute("""
                DELETE FROM blog_app_locks
                WHERE expires_at < NOW()
            """)

            cur.execute("""
                INSERT INTO blog_app_locks
                (
                    lock_name,
                    locked_at,
                    expires_at,
                    owner_token
                )
                VALUES
                (
                    %s,
                    NOW(),
                    DATE_ADD(NOW(), INTERVAL %s MINUTE),
                    %s
                )
            """, (
                str(lock_name)[:100],
                int(ttl_minutes),
                str(owner_token)[:100],
            ))

        conn.commit()
        log(f"[DB LOCK ACQUIRED] {lock_name} owner={owner_token}")
        return True, owner_token

    except Exception as e:
        try:
            if conn:
                conn.rollback()
        except Exception:
            pass

        log(f"[DB LOCKED OR ERROR] {lock_name} / {e}")
        return False, owner_token

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


def release_db_lock(lock_name, owner_token):
    conn = None

    try:
        conn = get_conn()

        with conn.cursor() as cur:
            cur.execute("""
                DELETE FROM blog_app_locks
                WHERE lock_name = %s
                  AND owner_token = %s
            """, (
                str(lock_name)[:100],
                str(owner_token)[:100],
            ))
            affected = cur.rowcount

        conn.commit()
        log(f"[DB LOCK RELEASED] {lock_name} affected={affected}")

    except Exception as e:
        log(f"[DB LOCK RELEASE ERROR] {lock_name} / {e}")

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


def create_pipeline_run(run_type="check_removed_articles", meta=None):
    conn = None

    try:
        conn = get_conn()

        with conn.cursor() as cur:
            cur.execute("""
                INSERT INTO blog_pipeline_runs
                (
                    run_type,
                    status,
                    started_at,
                    message,
                    meta_json,
                    created_at
                )
                VALUES
                (
                    %s,
                    'running',
                    NOW(),
                    %s,
                    %s,
                    NOW()
                )
            """, (
                run_type,
                "check_removed_articles started",
                safe_json_dumps(meta or {})
            ))
            run_id = cur.lastrowid

        conn.commit()
        log(f"[PIPELINE RUN START] run_id={run_id}, run_type={run_type}")
        return run_id

    except Exception as e:
        log(f"[PIPELINE RUN START ERROR] {e}")
        return None

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


def finish_pipeline_run(run_id, status, total=0, success=0, failed=0, skipped=0, message="", meta=None):
    if not run_id:
        return

    conn = None

    try:
        conn = get_conn()

        with conn.cursor() as cur:
            cur.execute("""
                UPDATE blog_pipeline_runs
                SET
                    status = %s,
                    finished_at = NOW(),
                    total_count = %s,
                    success_count = %s,
                    fail_count = %s,
                    skipped_count = %s,
                    message = %s,
                    meta_json = %s
                WHERE id = %s
            """, (
                status,
                int(total or 0),
                int(success or 0),
                int(failed or 0),
                int(skipped or 0),
                str(message or "")[:2000],
                safe_json_dumps(meta or {}),
                run_id
            ))

        conn.commit()
        log(f"[PIPELINE RUN FINISH] run_id={run_id}, status={status}")

    except Exception as e:
        log(f"[PIPELINE RUN FINISH ERROR] {e}")

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


def add_pipeline_log(run_id=None, level="info", step_name="", message="", article_no=None, realtor_id=None, context=None):
    conn = None

    try:
        conn = get_conn()

        with conn.cursor() as cur:
            cur.execute("""
                INSERT INTO blog_pipeline_logs
                (
                    run_id,
                    level,
                    step_name,
                    article_no,
                    realtor_id,
                    message,
                    context_json,
                    created_at
                )
                VALUES
                (
                    %s,
                    %s,
                    %s,
                    %s,
                    %s,
                    %s,
                    %s,
                    NOW()
                )
            """, (
                run_id,
                str(level or "info")[:20],
                str(step_name or "")[:100],
                str(article_no or "")[:30] if article_no else None,
                realtor_id,
                str(message or "")[:2000],
                safe_json_dumps(context or {})
            ))

        conn.commit()

    except Exception as e:
        log(f"[PIPELINE LOG ERROR] {step_name} / {e}")

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


def get_eligible_realtors(realtor_id=None, limit=100):
    """
    삭제매물 검증 대상 중개사.

    조건:
    - blog_realtors.status='active'
    - blog_publish_settings.auto_publish_enabled=1
    - 실제 발행 Queue(queue_status='published') 존재

    비공개 처리는 신규 발행 자격 심사가 아니라 이미 공개된 글의 사후관리다.
    따라서 현재 결제/서비스 기간으로 대상을 다시 제한하면, 결제 상태가 바뀐
    중개사의 기존 공개 글이 영구히 누락될 수 있으므로 membership 조건을 두지 않는다.
    """
    conn = get_conn()

    try:
        params = []

        sql = """
            SELECT
                r.id AS realtor_id,
                r.office_name,
                r.naver_realtor_id,
                COUNT(DISTINCT q.id) AS published_count
            FROM blog_realtors r
            INNER JOIN blog_publish_queue q
                ON q.realtor_id = r.id
               AND q.queue_status = 'published'
            LEFT JOIN blog_article_drafts d
                ON d.id = q.draft_id
            WHERE r.status = 'active'
              /* 구버전 서버발행은 blog_post_no를 초안 테이블에만 저장한 경우가 있다. */
              AND COALESCE(q.blog_post_no, d.blog_post_no, '') <> ''
        """

        if realtor_id:
            sql += " AND r.id = %s "
            params.append(int(realtor_id))

        sql += """
            GROUP BY r.id, r.office_name, r.naver_realtor_id
            ORDER BY r.id ASC
            LIMIT %s
        """
        params.append(int(limit))

        with conn.cursor() as cur:
            cur.execute(sql, params)
            rows = cur.fetchall()

        return rows or []

    finally:
        conn.close()


def get_published_articles_for_realtor(realtor_id, seen_hours=30):
    """
    해당 중개사의 발행 완료 글 목록과 현재 수집 상태를 조회한다.

    current_seen_recent:
    - blog_realtor_articles.last_seen_at이 seen_hours 이내면 현재 전체매물 수집에서 발견된 것으로 간주한다.
    """
    conn = get_conn()

    try:
        with conn.cursor() as cur:
            cur.execute("""
                SELECT
                    d.id AS draft_id,
                    q.realtor_id,
                    COALESCE(q.article_no, d.article_no) AS article_no,
                    COALESCE(d.blog_url, q.publish_result_url) AS blog_url,
                    COALESCE(q.blog_id, d.blog_id) AS blog_id,
                    COALESCE(q.blog_post_no, d.blog_post_no) AS blog_post_no,
                    COALESCE(q.finished_at, d.published_at) AS published_at,
                    COALESCE(NULLIF(LOWER(q.execution_target), ''), 'server') AS execution_target,
                    ra.id AS realtor_article_id,
                    ra.article_name,
                    ra.building_name,
                    ra.trade_type,
                    ra.price_text,
                    ra.area_info,
                    ra.first_posted_at,
                    ra.article_status,
                    ra.visibility_status,
                    ra.removal_status,
                    ra.removal_reason,
                    ra.reposted_candidate_article_no,
                    COALESCE(ra.removal_confirm_count, 0) AS removal_confirm_count,
                    ra.last_seen_at,
                    ra.removed_detected_at,
                    CASE
                        WHEN ra.last_seen_at IS NOT NULL
                         AND ra.last_seen_at >= DATE_SUB(NOW(), INTERVAL %s HOUR)
                        THEN 1
                        ELSE 0
                    END AS current_seen_recent
                FROM blog_publish_queue q
                LEFT JOIN blog_article_drafts d
                    ON d.id = q.draft_id
                LEFT JOIN blog_realtor_articles ra
                    ON CONVERT(ra.article_no USING utf8mb4) COLLATE utf8mb4_unicode_ci
                       = CONVERT(COALESCE(q.article_no, d.article_no) USING utf8mb4) COLLATE utf8mb4_unicode_ci
                   AND ra.realtor_id = q.realtor_id
                WHERE q.realtor_id = %s
                  AND q.queue_status = 'published'
                  AND COALESCE(q.blog_post_no, d.blog_post_no, '') <> ''
                  AND q.id = (
                      SELECT MAX(q2.id)
                      FROM blog_publish_queue q2
                      WHERE q2.realtor_id = q.realtor_id
                        AND CONVERT(COALESCE(q2.article_no, '') USING utf8mb4) COLLATE utf8mb4_unicode_ci
                            = CONVERT(COALESCE(q.article_no, d.article_no, '') USING utf8mb4) COLLATE utf8mb4_unicode_ci
                        AND q2.queue_status = 'published'
                  )
                ORDER BY q.id ASC
            """, (
                int(seen_hours),
                int(realtor_id),
            ))

            return cur.fetchall() or []

    finally:
        conn.close()


def mark_article_active(row, dry_run=False):
    article_no = str(row.get("article_no") or "")
    realtor_id = int(row.get("realtor_id") or 0)

    if dry_run:
        log(f"[DRY ACTIVE] realtor_id={realtor_id}, article_no={article_no}")
        return

    conn = get_conn()

    try:
        with conn.cursor() as cur:
            cur.execute("""
                UPDATE blog_realtor_articles
                SET
                    article_status = 'active',
                    visibility_status = 'public',
                    removal_status = NULL,
                    removal_reason = NULL,
                    reposted_candidate_article_no = NULL,
                    removal_confirm_count = 0,
                    removed_detected_at = NULL,
                    removal_confirmed_at = NULL,
                    updated_at = NOW()
                WHERE realtor_id = %s
                  AND article_no = %s
            """, (
                realtor_id,
                article_no
            ))

        conn.commit()

    finally:
        conn.close()


def get_article_age_days_for_private(row):
    """
    블로그 비공개 큐 등록 최소 경과일 판단.
    first_posted_at이 있으면 네이버 최초게재일 기준,
    없으면 블로그 published_at 기준으로 보수적으로 판단한다.
    """
    base_date = row.get("first_posted_at") or row.get("published_at")

    if not base_date:
        return None

    try:
        return (datetime.now() - base_date).days
    except Exception:
        return None


def is_private_queue_allowed(row, min_days_before_private=30):
    """
    신규 글이 오탐지로 바로 비공개되는 것을 막기 위한 최종 안전장치.
    """
    age_days = get_article_age_days_for_private(row)

    if age_days is None:
        return False, "first_posted_at/published_at 없음"

    if age_days < int(min_days_before_private):
        return False, f"최소 {min_days_before_private}일 미경과 age_days={age_days}"

    return True, f"age_days={age_days}"



def upsert_private_queue(row, dry_run=False, min_days_before_private=30):
    article_no = str(row.get("article_no") or "")
    realtor_id = int(row.get("realtor_id") or 0)
    draft_id = row.get("draft_id")
    blog_url = str(row.get("blog_url") or "")
    blog_id = str(row.get("blog_id") or "")
    blog_post_no = str(row.get("blog_post_no") or "")
    execution_target = (
        "client"
        if str(row.get("execution_target") or "").strip().lower() == "client"
        else "server"
    )

    if not article_no or not realtor_id or not blog_post_no:
        return False, "private_queue 생성 불가: article_no/realtor_id/blog_post_no 누락"

    allowed, allowed_message = is_private_queue_allowed(
        row,
        min_days_before_private=min_days_before_private,
    )

    if not allowed:
        log(
            f"[PRIVATE QUEUE HOLD] realtor_id={realtor_id}, article_no={article_no}, "
            f"reason={allowed_message}"
        )
        return False, "private_queue 보류: " + allowed_message

    removal_reason = str(row.get("removal_reason") or "").strip()
    reposted_candidate_article_no = str(row.get("reposted_candidate_article_no") or "").strip()

    if not removal_reason:
        removal_reason, reposted_candidate_article_no = estimate_removal_reason(row)

    queue_reason = removal_reason or "removed_confirmed"

    if reposted_candidate_article_no:
        queue_reason = f"{queue_reason}:reposted_candidate={reposted_candidate_article_no}"

    if dry_run:
        log(
            f"[DRY PRIVATE QUEUE] realtor_id={realtor_id}, article_no={article_no}, "
            f"blog_post_no={blog_post_no}, reason={queue_reason}"
        )
        return True, "dry_run"

    conn = get_conn()

    try:
        with conn.cursor() as cur:
            cur.execute("""
                INSERT INTO blog_post_private_queue
                (
                    realtor_id,
                    article_no,
                    draft_id,
                    blog_url,
                    blog_id,
                    blog_post_no,
                    execution_target,
                    reason,
                    queue_status,
                    detected_at,
                    scheduled_at,
                    created_at,
                    updated_at
                )
                VALUES
                (
                    %s,
                    %s,
                    %s,
                    %s,
                    %s,
                    %s,
                    %s,
                    %s,
                    'pending',
                    NOW(),
                    NULL,
                    NOW(),
                    NOW()
                )
                ON DUPLICATE KEY UPDATE
                    draft_id = VALUES(draft_id),
                    blog_url = VALUES(blog_url),
                    blog_id = VALUES(blog_id),
                    blog_post_no = VALUES(blog_post_no),
                    execution_target = VALUES(execution_target),
                    reason = VALUES(reason),
                    queue_status = IF(queue_status IN ('completed', 'private_done', 'processing'), queue_status, 'pending'),
                    updated_at = NOW()
            """, (
                realtor_id,
                article_no,
                draft_id,
                blog_url,
                blog_id,
                blog_post_no,
                execution_target,
                str(queue_reason or "removed_confirmed")[:100]
            ))

        conn.commit()
        return True, "queued"

    finally:
        conn.close()



def estimate_removal_reason(row):
    """
    articleNo가 현재 목록에서 사라진 이유를 추정한다.

    - expired_estimated: 네이버 자동만료 가능성. first_posted_at 또는 published_at 기준 30일 이상
    - reposted_estimated: 동일/유사 매물이 새 articleNo로 재등록된 것으로 추정
    - missing_from_naver: 현재 네이버 전체매물 목록에서 단순 미확인
    - unknown: 판단 불가

    이 값은 확정 사유가 아니라 운영 참고용이다.
    """
    article_no = str(row.get("article_no") or "")
    realtor_id = int(row.get("realtor_id") or 0)

    first_posted_at = row.get("first_posted_at")
    published_at = row.get("published_at")

    base_date = first_posted_at or published_at

    try:
        if base_date:
            # PyMySQL은 datetime 객체로 주는 경우가 많다.
            age_days = (datetime.now() - base_date).days

            if age_days >= 30:
                return "expired_estimated", ""
    except Exception:
        pass

    candidate_article_no = find_reposted_candidate_article_no(row)

    if candidate_article_no:
        return "reposted_estimated", candidate_article_no

    if article_no and realtor_id:
        return "missing_from_naver", ""

    return "unknown", ""


def find_reposted_candidate_article_no(row):
    """
    동일/유사 매물이 새 articleNo로 재등록됐는지 보수적으로 추정한다.

    현재는 안전하게 다음 기준만 사용한다.
    - 같은 realtor_id
    - 최근 seen_hours 내 last_seen_at 갱신
    - article_no 다름
    - building_name 또는 article_name이 같음
    - 가능하면 trade_type / price_text / area_info 중 일부 일치

    이 함수는 확정 판단이 아니라 reason 기록용이다.
    """
    realtor_id = int(row.get("realtor_id") or 0)
    article_no = str(row.get("article_no") or "")

    building_name = str(row.get("building_name") or "").strip()
    article_name = str(row.get("article_name") or "").strip()
    trade_type = str(row.get("trade_type") or "").strip()
    price_text = str(row.get("price_text") or "").strip()
    area_info = str(row.get("area_info") or "").strip()

    if not realtor_id or not article_no:
        return ""

    if not building_name and not article_name:
        return ""

    conn = get_conn()

    try:
        params = [
            realtor_id,
            article_no,
        ]

        sql = """
            SELECT
                article_no,
                article_name,
                building_name,
                trade_type,
                price_text,
                area_info,
                last_seen_at
            FROM blog_realtor_articles
            WHERE realtor_id = %s
              AND article_no <> %s
              AND last_seen_at IS NOT NULL
              AND last_seen_at >= DATE_SUB(NOW(), INTERVAL 30 HOUR)
              AND article_status = 'active'
        """

        conditions = []

        if building_name:
            conditions.append("building_name = %s")
            params.append(building_name)

        if article_name:
            conditions.append("article_name = %s")
            params.append(article_name)

        if not conditions:
            return ""

        sql += " AND (" + " OR ".join(conditions) + ") "

        sql += """
            ORDER BY last_seen_at DESC
            LIMIT 20
        """

        with conn.cursor() as cur:
            cur.execute(sql, params)
            candidates = cur.fetchall() or []

        for cand in candidates:
            score = 0

            if building_name and str(cand.get("building_name") or "").strip() == building_name:
                score += 3

            if article_name and str(cand.get("article_name") or "").strip() == article_name:
                score += 2

            if trade_type and str(cand.get("trade_type") or "").strip() == trade_type:
                score += 1

            if price_text and str(cand.get("price_text") or "").strip() == price_text:
                score += 1

            if area_info and str(cand.get("area_info") or "").strip() == area_info:
                score += 1

            # 너무 약한 유사도는 오판 위험이 있으므로 기록하지 않는다.
            if score >= 4:
                return str(cand.get("article_no") or "")

        return ""

    finally:
        conn.close()


def mark_article_missing(row, confirm_count=2, dry_run=False, allow_removed_confirm=True, min_days_before_private=30):
    """
    현재 전체매물에서 발견되지 않은 발행 글을 단계별로 상태 변경한다.
    """
    article_no = str(row.get("article_no") or "")
    realtor_id = int(row.get("realtor_id") or 0)
    current_count = int(row.get("removal_confirm_count") or 0)
    previous_status = str(row.get("removal_status") or row.get("article_status") or "")

    next_count = current_count + 1

    private_allowed, private_allowed_message = is_private_queue_allowed(
        row,
        min_days_before_private=min_days_before_private,
    )

    if next_count >= int(confirm_count) and allow_removed_confirm and private_allowed:
        next_status = "removed_confirmed"
        visibility_status = "private_queued"
    else:
        # 안전장치:
        # - 현재매물 수집 자체가 실패한 경우(recent_seen_count=0)
        # - 최근 수집 매물이 너무 많은 경우(page=2 이후 누락 가능성)
        # - 최초게재/발행 후 최소 유예기간 미경과
        # 위 경우 삭제 확정 및 private_queue 생성을 막고 suspected_removed까지만 허용한다.
        next_status = "suspected_removed"
        visibility_status = str(row.get("visibility_status") or "public")

    removal_reason, reposted_candidate_article_no = estimate_removal_reason(row)

    if dry_run:
        log(
            f"[DRY MISSING] realtor_id={realtor_id}, article_no={article_no}, "
            f"prev={previous_status}, count={current_count}->{next_count}, next={next_status}, "
            f"allow_removed_confirm={allow_removed_confirm}, "
            f"private_allowed={private_allowed}, private_reason={private_allowed_message}, "
            f"reason={removal_reason}, reposted_candidate={reposted_candidate_article_no}"
        )

        if next_status == "removed_confirmed":
            upsert_private_queue(row, dry_run=True, min_days_before_private=min_days_before_private)

        return next_status

    conn = get_conn()

    try:
        with conn.cursor() as cur:
            cur.execute("""
                UPDATE blog_realtor_articles
                SET
                    article_status = %s,
                    removal_status = %s,
                    removal_reason = %s,
                    reposted_candidate_article_no = %s,
                    removal_confirm_count = %s,
                    removed_detected_at = IF(removed_detected_at IS NULL, NOW(), removed_detected_at),
                    removal_confirmed_at = CASE
                        WHEN %s = 'removed_confirmed' THEN NOW()
                        ELSE removal_confirmed_at
                    END,
                    visibility_status = %s,
                    updated_at = NOW()
                WHERE realtor_id = %s
                  AND article_no = %s
            """, (
                next_status,
                next_status,
                removal_reason,
                reposted_candidate_article_no,
                next_count,
                next_status,
                visibility_status,
                realtor_id,
                article_no
            ))

        conn.commit()

    finally:
        conn.close()

    if next_status == "removed_confirmed":
        ok, message = upsert_private_queue(row, dry_run=False, min_days_before_private=min_days_before_private)

        if not ok:
            log(f"[PRIVATE QUEUE ERROR] realtor_id={realtor_id}, article_no={article_no}, {message}")

    return next_status



def get_recent_seen_count_for_realtor(realtor_id, seen_hours=30):
    """
    최근 seen_hours 이내에 last_seen_at이 갱신된 현재매물 수.
    현재 수집기가 page=1 기준 20건을 안정적으로 잡기 때문에,
    20건 이상이면 목록이 더 있을 가능성이 있어 삭제확정은 보수적으로 처리한다.
    """
    conn = get_conn()

    try:
        with conn.cursor() as cur:
            cur.execute("""
                SELECT COUNT(*) AS cnt
                FROM blog_realtor_articles
                WHERE realtor_id = %s
                  AND last_seen_at IS NOT NULL
                  AND last_seen_at >= DATE_SUB(NOW(), INTERVAL %s HOUR)
            """, (
                int(realtor_id),
                int(seen_hours),
            ))

            row = cur.fetchone() or {}

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

    finally:
        conn.close()


def process_realtor(row, seen_hours=30, confirm_count=2, dry_run=False, run_id=None, min_days_before_private=30):
    realtor_id = int(row.get("realtor_id") or 0)
    office_name = str(row.get("office_name") or "")

    log("-" * 80)
    log(f"[REALTOR START] realtor_id={realtor_id}, office={office_name}")

    articles = get_published_articles_for_realtor(
        realtor_id=realtor_id,
        seen_hours=seen_hours
    )

    log(f"[PUBLISHED ARTICLES] realtor_id={realtor_id}, count={len(articles)}")

    recent_seen_count = get_recent_seen_count_for_realtor(
        realtor_id=realtor_id,
        seen_hours=seen_hours
    )

    if recent_seen_count <= 0:
        log(
            f"[REMOVAL CHECK SKIP] realtor_id={realtor_id}, "
            f"recent_seen_count={recent_seen_count}, fresh listing snapshot 없음 / 상태 변경 안 함"
        )

        # 최신 전체매물 snapshot이 없으면 모든 발행 글을 '사라짐'으로 오판한다.
        # 보류 로그만 남기고 removal_count/status/비공개 Queue를 절대 변경하지 않는다.
        add_pipeline_log(
            run_id=run_id,
            level="warning",
            step_name="skip_no_fresh_snapshot",
            message="fresh listing snapshot 없음 / 비공개 판정 전체 건너뜀",
            realtor_id=realtor_id,
            context={
                "recent_seen_count": recent_seen_count,
                "seen_hours": seen_hours,
                "published_articles": len(articles),
            },
        )
        return {
            "active": 0,
            "suspected_removed": 0,
            "removed_confirmed": 0,
            "skipped_no_article_row": 0,
            "skipped_no_fresh_snapshot": len(articles),
        }

    # 현재 수집기는 UI 스크롤로 전체 목록을 수집한다. 과거 page=1 수집기에서
    # 사용하던 20건 이상 확정 금지는 정상 대형 중개사의 비공개 Queue를
    # 영구 차단하므로 적용하지 않는다.
    allow_removed_confirm = True
    log(f"[REMOVAL CONFIRM ALLOWED] realtor_id={realtor_id}, recent_seen_count={recent_seen_count}")

    stats = {
        "active": 0,
        "suspected_removed": 0,
        "removed_confirmed": 0,
        "skipped_no_article_row": 0,
        "skipped_no_fresh_snapshot": 0,
    }

    for article in articles:
        article_no = str(article.get("article_no") or "")

        if not article.get("realtor_article_id"):
            stats["skipped_no_article_row"] += 1
            log(f"[SKIP NO REALTOR ARTICLE] realtor_id={realtor_id}, article_no={article_no}")
            add_pipeline_log(
                run_id=run_id,
                level="warning",
                step_name="skip_no_realtor_article",
                message="blog_realtor_articles row missing",
                article_no=article_no,
                realtor_id=realtor_id,
                context=article
            )
            continue

        if int(article.get("current_seen_recent") or 0) == 1:
            mark_article_active(article, dry_run=dry_run)
            stats["active"] += 1

            add_pipeline_log(
                run_id=run_id,
                level="info",
                step_name="article_active",
                message="article currently exists",
                article_no=article_no,
                realtor_id=realtor_id,
                context={
                    "last_seen_at": article.get("last_seen_at"),
                    "dry_run": dry_run,
                }
            )

        else:
            next_status = mark_article_missing(
                article,
                confirm_count=confirm_count,
                dry_run=dry_run,
                allow_removed_confirm=allow_removed_confirm,
                min_days_before_private=min_days_before_private
            )

            stats[next_status] = stats.get(next_status, 0) + 1

            add_pipeline_log(
                run_id=run_id,
                level="warning",
                step_name=next_status,
                message=f"article not seen recently: status={next_status}",
                article_no=article_no,
                realtor_id=realtor_id,
                context={
                    "last_seen_at": article.get("last_seen_at"),
                    "removal_confirm_count": article.get("removal_confirm_count"),
                    "seen_hours": seen_hours,
                    "confirm_count": confirm_count,
                    "recent_seen_count": recent_seen_count,
                    "allow_removed_confirm": allow_removed_confirm,
                    "min_days_before_private": min_days_before_private,
                    "removal_reason": estimate_removal_reason(article)[0],
                    "reposted_candidate_article_no": estimate_removal_reason(article)[1],
                    "dry_run": dry_run,
                }
            )

    log(f"[REALTOR DONE] realtor_id={realtor_id}, stats={json.dumps(stats, ensure_ascii=False)}")
    return stats


def run(limit=100, realtor_id=None, seen_hours=30, confirm_count=2, dry_run=False, no_db_lock=False, lock_ttl_minutes=180, min_days_before_private=30):
    lock_name = "check_removed_articles"
    lock_owner = None
    run_id = None

    total_realtors = 0
    total_articles = 0
    total_removed_confirmed = 0
    total_suspected = 0
    total_active = 0
    total_failed = 0

    try:
        if not no_db_lock:
            locked, lock_owner = acquire_db_lock(lock_name, ttl_minutes=lock_ttl_minutes)

            if not locked:
                log("[SKIP] check_removed_articles already running")
                return

        run_id = create_pipeline_run(
            run_type="check_removed_articles",
            meta={
                "limit": limit,
                "realtor_id": realtor_id,
                "seen_hours": seen_hours,
                "confirm_count": confirm_count,
                "min_days_before_private": min_days_before_private,
                "dry_run": dry_run,
            }
        )

        realtors = get_eligible_realtors(
            realtor_id=realtor_id,
            limit=limit
        )

        log("=" * 80)
        log(f"[TARGET REALTORS] {len(realtors)}")
        log(f"[SEEN HOURS] {seen_hours}")
        log(f"[CONFIRM COUNT] {confirm_count}")
        log(f"[MIN DAYS BEFORE PRIVATE] {min_days_before_private}")
        log(f"[DRY RUN] {'ON' if dry_run else 'OFF'}")
        log("=" * 80)

        for row in realtors:
            total_realtors += 1

            try:
                stats = process_realtor(
                    row,
                    seen_hours=seen_hours,
                    confirm_count=confirm_count,
                    dry_run=dry_run,
                    run_id=run_id,
                    min_days_before_private=min_days_before_private
                )

                total_active += int(stats.get("active") or 0)
                total_suspected += int(stats.get("suspected_removed") or 0)
                total_removed_confirmed += int(stats.get("removed_confirmed") or 0)
                total_articles += sum(int(v or 0) for v in stats.values())

            except Exception as e:
                total_failed += 1
                traceback.print_exc()
                add_pipeline_log(
                    run_id=run_id,
                    level="error",
                    step_name="realtor_exception",
                    message=str(e)[:2000],
                    realtor_id=row.get("realtor_id"),
                    context={"traceback": traceback.format_exc()}
                )

        final_message = (
            f"realtors={total_realtors}, articles={total_articles}, "
            f"active={total_active}, suspected={total_suspected}, "
            f"removed_confirmed={total_removed_confirmed}, failed={total_failed}"
        )

        finish_pipeline_run(
            run_id,
            status="success" if total_failed == 0 else "failed",
            total=total_articles,
            success=total_active + total_suspected + total_removed_confirmed,
            failed=total_failed,
            skipped=0,
            message=final_message,
            meta={
                "total_realtors": total_realtors,
                "total_articles": total_articles,
                "total_active": total_active,
                "total_suspected": total_suspected,
                "total_removed_confirmed": total_removed_confirmed,
                "total_failed": total_failed,
                "min_days_before_private": min_days_before_private,
                "dry_run": dry_run,
            }
        )

        log("=" * 80)
        log("[CHECK REMOVED DONE]")
        log(final_message)
        log("=" * 80)

    finally:
        if lock_owner:
            release_db_lock(lock_name, lock_owner)


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

    parser.add_argument("--limit", type=int, default=100)
    parser.add_argument("--realtor-id", type=int, default=None)
    parser.add_argument("--seen-hours", type=int, default=30)
    parser.add_argument("--confirm-count", type=int, default=2)
    parser.add_argument("--min-days-before-private", type=int, default=30)
    parser.add_argument("--dry-run", action="store_true")
    parser.add_argument("--no-db-lock", action="store_true")
    parser.add_argument("--lock-ttl-minutes", type=int, default=180)

    args = parser.parse_args()

    run(
        limit=args.limit,
        realtor_id=args.realtor_id,
        seen_hours=args.seen_hours,
        confirm_count=args.confirm_count,
        dry_run=args.dry_run,
        no_db_lock=args.no_db_lock,
        lock_ttl_minutes=args.lock_ttl_minutes,
        min_days_before_private=args.min_days_before_private,
    )


if __name__ == "__main__":
    main()
