﻿# -*- coding: utf-8 -*-
# STEP119_V15_PRIVATE_SYNC_PATCH
# STEP117_V15_FULL_SCROLL_COLLECT_PATCH

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


def configure_utf8_console():
    """
    Windows 콘솔 또는 로그 파일로 출력할 때 한글이 깨지지 않도록
    표준출력과 표준오류를 UTF-8로 통일한다.
    """
    for stream_name in ("stdout", "stderr"):
        stream = getattr(sys, stream_name, None)
        if stream is None or not hasattr(stream, "reconfigure"):
            continue

        try:
            stream.reconfigure(encoding="utf-8", errors="replace")
        except Exception:
            pass


configure_utf8_console()

BASE_DIR = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
sys.path.append(BASE_DIR)

from db import get_conn


# ---------------------------------------------------------
# STEP29 OPERATION LOCK / PIPELINE LOG HELPERS
# - 기존 성공 흐름은 유지하고, DB 작업락/로그만 보강한다.
# - blog_app_locks 구조: lock_name, locked_at, expires_at, owner_token
# - JSON 타입 미지원 환경을 고려해 meta_json/context_json에는 문자열 저장
# ---------------------------------------------------------

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=None):
    """
    DB 기반 작업락.
    - 만료된 락은 제거
    - 같은 lock_name이 살아 있으면 실행하지 않음
    - 실패해도 예외를 밖으로 던지지 않고 False 반환
    """
    owner_token = owner_token or _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()
        print(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

        print(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()
        print(f"[DB LOCK RELEASED] {lock_name} affected={affected}")

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

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


def create_pipeline_run(run_type, message="", 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()
                )
            """, (
                str(run_type)[:50],
                str(message or "")[:2000],
                _safe_json_dumps(meta or {}),
            ))

            run_id = cur.lastrowid

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

    except Exception as e:
        print(f"[PIPELINE RUN START ERROR] {run_type} / {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
            """, (
                str(status)[:30],
                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 {}),
                int(run_id),
            ))

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

    except Exception as e:
        print(f"[PIPELINE RUN FINISH ERROR] run_id={run_id} / {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:
        print(f"[PIPELINE LOG ERROR] {step_name} / {e}")

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


from services.naver_realtor_response_finder import (
    NaverRealtorResponseFinder,
)


def table_columns(conn, table_name):
    with conn.cursor() as cur:
        cur.execute(f"SHOW COLUMNS FROM {table_name}")
        rows = cur.fetchall()

    return [row["Field"] for row in rows]


def fetch_target_realtors(conn, limit_realtors=100, realtor_id=None):
    sql = """
        SELECT DISTINCT
            r.*,
            COALESCE(s.auto_publish_enabled, 0) AS auto_publish_enabled,
            COALESCE(s.daily_post_limit, 1) AS daily_post_limit
        FROM blog_realtors r
        LEFT JOIN blog_publish_settings s
            ON r.id = s.realtor_id
        LEFT JOIN blog_naver_accounts a
            ON r.id = a.realtor_id
        WHERE r.status = 'active'
          AND COALESCE(r.is_collect_enabled, 0) = 1
          AND COALESCE(TRIM(r.naver_realtor_id), '') <> ''
    """

    params = []

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

    sql += " ORDER BY r.id ASC LIMIT %s "
    params.append(int(limit_realtors))

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


def get_today_scheduled_count(conn, realtor_id):
    with conn.cursor() as cur:
        cur.execute("""
            SELECT COUNT(*) AS cnt
            FROM blog_article_work_queue
            WHERE realtor_id = %s
              AND work_type = 'new_article'
              AND work_status IN ('pending', 'processing', 'done')
              AND DATE(created_at) = CURDATE()
        """, (int(realtor_id),))
        row = cur.fetchone()

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


def queue_exists(conn, article_no):
    with conn.cursor() as cur:
        cur.execute("""
            SELECT id
            FROM blog_article_work_queue
            WHERE article_no = %s
              AND work_status IN ('pending', 'processing', 'done')
            LIMIT 1
        """, (str(article_no),))
        row = cur.fetchone()

    return bool(row)


def draft_exists(conn, article_no):
    with conn.cursor() as cur:
        cur.execute("""
            SELECT id
            FROM blog_article_drafts
            WHERE article_no = %s
            LIMIT 1
        """, (str(article_no),))
        row = cur.fetchone()

    return bool(row)


def publish_exists(conn, article_no):
    with conn.cursor() as cur:
        cur.execute("""
            SELECT q.id
            FROM blog_publish_queue q
            INNER JOIN blog_article_drafts d
                ON q.draft_id = d.id
            WHERE d.article_no = %s
              AND q.queue_status IN (
                    'waiting',
                    'processing',
                    'published',
                    'done',
                    'hold'
              )
            LIMIT 1
        """, (str(article_no),))
        row = cur.fetchone()

    return bool(row)



def get_draft_buffer_target(realtor, default_target=45):
    """
    중개사별 초안 버퍼 목표.
    blog_realtors.draft_buffer_target 값이 있으면 우선 사용한다.
    """
    try:
        value = int((realtor or {}).get("draft_buffer_target") or 0)
        if value > 0:
            return value
    except Exception:
        pass

    try:
        default_target = int(default_target or 45)
    except Exception:
        default_target = 45

    return max(1, default_target)


def count_active_draft_buffer(conn, realtor_id):
    """
    현재 초안 재고 수.
    published/done/private_done 완료 큐는 재고로 보지 않는다.
    draft는 있으나 queue가 없는 경우도 재고로 본다.
    """
    with conn.cursor() as cur:
        cur.execute("""
            SELECT COUNT(DISTINCT d.id) AS cnt
            FROM blog_article_drafts d
            LEFT JOIN blog_publish_queue q
                ON (
                    q.draft_id = d.id
                    OR (
                        q.article_no IS NOT NULL
                        AND CONVERT(q.article_no USING utf8mb4) COLLATE utf8mb4_unicode_ci
                            =
                            CONVERT(d.article_no USING utf8mb4) COLLATE utf8mb4_unicode_ci
                    )
                )
            WHERE d.realtor_id = %s
              AND COALESCE(d.status, 'ready') IN ('ready', 'approved', 'pending')
              AND (
                    q.id IS NULL
                    OR COALESCE(q.queue_status, 'waiting') IN (
                        'waiting',
                        'pending',
                        'ready',
                        'processing',
                        'failed',
                        'hold'
                    )
              )
        """, (int(realtor_id),))
        row = cur.fetchone()

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


def count_pending_work_queue(conn, realtor_id):
    """
    상세수집 대기/진행 중인 work_queue 수.
    아직 초안은 아니지만 초안 재고로 전환될 예정이므로 버퍼 계산에 포함한다.
    """
    with conn.cursor() as cur:
        cur.execute("""
            SELECT COUNT(*) AS cnt
            FROM blog_article_work_queue
            WHERE realtor_id = %s
              AND work_type = 'new_article'
              AND work_status IN ('pending', 'processing')
        """, (int(realtor_id),))
        row = cur.fetchone()

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


def is_processed_article(conn, article_no):
    if queue_exists(conn, article_no):
        return True

    if draft_exists(conn, article_no):
        return True

    if publish_exists(conn, article_no):
        return True

    return False


def upsert_realtor_article_cache(conn, realtor, article):
    columns = table_columns(conn, "blog_realtor_articles")

    realtor_id = int(realtor["id"])
    article_no = str(article.get("article_no") or "").strip()

    if not article_no:
        return False

    data = {}

    if "realtor_id" in columns:
        data["realtor_id"] = realtor_id

    if "search_run_id" in columns:
        data["search_run_id"] = None

    if "article_no" in columns:
        data["article_no"] = article_no

    if "article_name" in columns:
        data["article_name"] = article.get("article_name") or "네이버 부동산 후보 매물"

    if "trade_type" in columns:
        data["trade_type"] = article.get("trade_type", "")

    if "real_estate_type" in columns:
        data["real_estate_type"] = article.get("real_estate_type", "")

    if "price_text" in columns:
        data["price_text"] = article.get("price_text", "")

    if "building_name" in columns:
        data["building_name"] = article.get("building_name", "")

    if "floor_info" in columns:
        data["floor_info"] = article.get("floor_info", "")

    if "area_info" in columns:
        data["area_info"] = article.get("area_info", "")

    if "article_url" in columns:
        data["article_url"] = (
            "https://new.land.naver.com/offices"
            f"?articleNo={article_no}"
        )

    if "matched_office_name" in columns:
        data["matched_office_name"] = article.get("matched_office_name", "")

    if "matched_representative_name" in columns:
        data["matched_representative_name"] = article.get("matched_representative_name", "")

    if "matched_phone" in columns:
        data["matched_phone"] = article.get("matched_phone", "")

    if "match_score" in columns:
        data["match_score"] = int(article.get("match_score") or 100)

    if "article_status" in columns:
        data["article_status"] = "active"

    if "visibility_status" in columns:
        data["visibility_status"] = "public"

    # STEP117:
    # 전체 스크롤 목록에서 확인된 현재 매물은 last_seen_at을 갱신한다.
    # check_removed_articles.py가 삭제/비공개 판단을 할 때 현재 살아있는 매물 기준으로 사용한다.
    if "last_seen_at" in columns:
        data["last_seen_at"] = datetime.now()

    # 네이버 목록 응답에서 받은 요약정보도 가능한 컬럼에만 안전 저장한다.
    if "confirm_date" in columns:
        data["confirm_date"] = article.get("confirm_date") or ""

    if "article_feature_desc" in columns:
        data["article_feature_desc"] = article.get("short_description") or ""

    if "collect_status" in columns:
        # STEP89-B:
        # schedule 단계는 "상세수집 후보 등록"만 한다.
        # 여기서 detail_collected=1 상태가 남아 있으면 fetcher가 스킵하고
        # generate_blog_drafts가 "네이버 부동산 후보 매물"로 초안을 만들어버린다.
        data["collect_status"] = "pending"

    if "detail_collected" in columns:
        data["detail_collected"] = 0

    if "detail_collected_at" in columns:
        data["detail_collected_at"] = None

    if "detail_error" in columns:
        data["detail_error"] = None

    if "created_at" in columns:
        data["created_at"] = datetime.now()

    if "updated_at" in columns:
        data["updated_at"] = datetime.now()

    keys = list(data.keys())

    if not keys:
        return False

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

    update_sets = []

    for key in keys:
        if key in ["article_no", "created_at"]:
            continue

        # STEP117 안전장치:
        # 전체 목록 스냅샷 저장 시 기존 상세수집 완료/실패 상태를 덮어쓰면 안 된다.
        # 신규 insert에는 pending/detail_collected=0을 넣되, duplicate update에서는 상태 컬럼을 보존한다.
        if key in ["collect_status", "detail_collected", "detail_collected_at", "detail_error"]:
            continue

        update_sets.append(f"{key} = VALUES({key})")

    if "updated_at" in columns:
        update_sets.append("updated_at = NOW()")

    sql = f"""
        INSERT INTO blog_realtor_articles
        ({insert_cols})
        VALUES ({placeholders})
        ON DUPLICATE KEY UPDATE
            {", ".join(update_sets)}
    """

    with conn.cursor() as cur:
        cur.execute(sql, [data[k] for k in keys])

    return True


def enqueue_article(conn, realtor_id, article_no, priority=100):
    with conn.cursor() as cur:
        cur.execute("""
            INSERT INTO blog_article_work_queue
            (
                realtor_id,
                article_no,
                work_type,
                work_status,
                priority,
                created_at,
                updated_at
            )
            VALUES
            (
                %s,
                %s,
                'new_article',
                'pending',
                %s,
                NOW(),
                NOW()
            )
        """, (
            int(realtor_id),
            str(article_no),
            int(priority),
        ))

    conn.commit()



# ---------------------------------------------------------
# STEP119 PRIVATE SYNC HELPERS
# - STEP118에서 확보한 전체 현재 article_no 목록을 기준으로,
#   이미 발행된 블로그 글 중 현재 네이버 목록에서 사라진 매물을
#   blog_post_private_queue에 등록한다.
# - 실제 비공개 처리는 publish_worker.py의 기존 private queue 처리 로직이 담당한다.
# - 수집 커버리지가 낮으면 오탐 방지를 위해 아무 것도 만들지 않는다.
# ---------------------------------------------------------


def _column_set(conn, table_name):
    try:
        return set(table_columns(conn, table_name))
    except Exception:
        return set()


def _table_exists(conn, table_name):
    try:
        with conn.cursor() as cur:
            cur.execute("SHOW TABLES LIKE %s", (str(table_name),))
            return cur.fetchone() is not None
    except Exception:
        return False


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
              AND queue_status IN ('pending', 'processing', 'completed', 'hold')
            LIMIT 1
        """, (int(realtor_id), str(article_no)))
        row = cur.fetchone()

    return bool(row)


def _insert_private_queue(conn, row, reason):
    cols = _column_set(conn, "blog_post_private_queue")
    if not cols:
        print("[PRIVATE SYNC SKIP] blog_post_private_queue missing or no columns")
        return False

    data = {}

    mapping = {
        "realtor_id": row.get("realtor_id"),
        "article_no": row.get("article_no"),
        "draft_id": row.get("draft_id"),
        "blog_url": row.get("blog_url") or row.get("publish_result_url"),
        "blog_id": row.get("blog_id"),
        "blog_post_no": row.get("blog_post_no"),
        # private Queue에도 원 발행 주체를 보존한다.
        # 실제 컬럼이 있는 운영 DB에서만 동적으로 INSERT되므로 구버전 스키마와도 호환된다.
        "execution_target": (
            "client"
            if str(row.get("execution_target") or "").strip().lower() == "client"
            else "server"
        ),
        "reason": reason,
        "queue_status": "pending",
        "worker_status": "pending",
        "retry_count": 0,
        "error_message": None,
        "last_error_type": None,
    }

    for key, value in mapping.items():
        if key in cols:
            data[key] = value

    if "scheduled_at" in cols:
        data["scheduled_at"] = datetime.now()

    if "created_at" in cols:
        data["created_at"] = datetime.now()

    if "updated_at" in cols:
        data["updated_at"] = datetime.now()

    if not data:
        return False

    keys = list(data.keys())
    sql = f"""
        INSERT INTO blog_post_private_queue
        ({', '.join(keys)})
        VALUES ({', '.join(['%s'] * len(keys))})
    """

    with conn.cursor() as cur:
        cur.execute(sql, [data[k] for k in keys])

    return True


def _mark_article_private_pending(conn, realtor_id, article_no):
    if not _table_exists(conn, "blog_realtor_articles"):
        return

    cols = _column_set(conn, "blog_realtor_articles")
    sets = []
    params = []

    if "visibility_status" in cols:
        sets.append("visibility_status = %s")
        params.append("private_pending")

    if "article_status" in cols:
        sets.append("article_status = %s")
        params.append("removed")

    if "updated_at" in cols:
        sets.append("updated_at = NOW()")

    if not sets:
        return

    params.extend([int(realtor_id), str(article_no)])

    with conn.cursor() as cur:
        cur.execute(f"""
            UPDATE blog_realtor_articles
            SET {', '.join(sets)}
            WHERE realtor_id = %s
              AND article_no = %s
        """, params)


def sync_deleted_articles_to_private_queue(
    conn,
    realtor,
    current_article_nos,
    visible_counts=None,
    min_coverage_ratio=0.9,
    enabled=True,
):
    """
    현재 네이버 목록에 없는 발행 완료 매물을 비공개 Queue로 등록한다.

    안전 조건:
    - 현재 수집 article_no가 비어 있으면 skip
    - 네이버 UI의 visible_total이 있으면 수집 수가 visible_total의 min_coverage_ratio 미만이면 skip
    - blog_id/blog_post_no가 없는 발행건은 비공개 처리 진입이 불안정하므로 skip
    """
    realtor_id = int(realtor.get("id"))

    if not enabled:
        print("[PRIVATE SYNC SKIP] disabled")
        return {"created": 0, "skipped": 0, "reason": "disabled"}

    if not _table_exists(conn, "blog_post_private_queue"):
        print("[PRIVATE SYNC SKIP] blog_post_private_queue table missing")
        return {"created": 0, "skipped": 0, "reason": "table_missing"}

    current_set = {str(x).strip() for x in (current_article_nos or []) if str(x or "").strip()}
    current_count = len(current_set)

    visible_counts = visible_counts or {}
    try:
        visible_total = int(visible_counts.get("total") or 0)
    except Exception:
        visible_total = 0

    if current_count <= 0:
        print("[PRIVATE SYNC SKIP] current article list empty")
        return {"created": 0, "skipped": 0, "reason": "current_empty"}

    if visible_total > 0:
        try:
            ratio = current_count / max(1, visible_total)
        except Exception:
            ratio = 0

        if ratio < float(min_coverage_ratio or 0.9):
            print(
                "[PRIVATE SYNC SKIP] low coverage",
                f"current={current_count}",
                f"visible_total={visible_total}",
                f"ratio={ratio:.3f}",
                f"min={float(min_coverage_ratio or 0.9):.3f}",
            )
            return {"created": 0, "skipped": 0, "reason": "low_coverage"}

    print(
        "[PRIVATE SYNC START]",
        f"realtor_id={realtor_id}",
        f"current={current_count}",
        f"visible_total={visible_total}",
    )

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

    sql = f"""
        SELECT
            q.id AS publish_queue_id,
            q.realtor_id,
            q.draft_id,
            COALESCE(q.article_no, d.article_no) AS article_no,
            COALESCE(NULLIF(q.execution_target, ''), 'server') AS execution_target,
            q.publish_result_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(d.blog_url, q.publish_result_url) AS blog_url
        FROM blog_publish_queue q
        LEFT JOIN blog_article_drafts d
            ON d.id = q.draft_id
        WHERE q.realtor_id = %s
          AND q.queue_status = 'published'
          AND COALESCE(q.article_no, d.article_no, '') <> ''
          AND COALESCE(q.article_no, d.article_no) NOT IN ({placeholders})
          AND COALESCE(q.blog_id, d.blog_id, '') <> ''
          AND COALESCE(q.blog_post_no, d.blog_post_no, '') <> ''
    """

    params = [realtor_id] + sorted(current_set)

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

    created = 0
    skipped = 0

    for row in rows:
        article_no = str(row.get("article_no") or "").strip()
        if not article_no:
            skipped += 1
            continue

        if _private_queue_exists(conn, realtor_id, article_no):
            print(f"[PRIVATE SYNC SKIP EXISTS] {article_no}")
            skipped += 1
            continue

        reason = "naver_article_missing_from_current_scroll_snapshot"

        try:
            if _insert_private_queue(conn, row, reason=reason):
                _mark_article_private_pending(conn, realtor_id, article_no)
                created += 1
                print(
                    "[PRIVATE SYNC QUEUED]",
                    f"article_no={article_no}",
                    f"blog_id={row.get('blog_id')}",
                    f"blog_post_no={row.get('blog_post_no')}",
                )
            else:
                skipped += 1
        except Exception as e:
            conn.rollback()
            skipped += 1
            print(f"[PRIVATE SYNC ERROR] {article_no} / {str(e)}")
            continue

    try:
        conn.commit()
    except Exception:
        pass

    print(
        "[PRIVATE SYNC DONE]",
        f"created={created}",
        f"skipped={skipped}",
        f"candidates={len(rows)}",
    )

    return {"created": created, "skipped": skipped, "candidates": len(rows)}


def process_realtor(
    conn,
    finder,
    realtor,
    wait_ms=2500,
    draft_buffer_default=1,
    replenish_count=1,
    max_pages=5,
    enable_private_sync=True,
    private_sync_min_coverage=0.9,
):
    """
    STEP85 운영형 스케줄러.

    핵심 변경:
    - daily_post_limit은 발행 목표값으로만 참고한다.
    - blog_article_work_queue는 '상세수집 대기열'로 사용한다.
    - 초안 재고가 draft_buffer_target 미만이면 부족분만큼 후보를 큐에 넣는다.
    - 한 번에 너무 많이 큐잉하지 않도록 replenish_count로 상한을 둔다.
    """
    realtor_id = int(realtor["id"])
    office_name = realtor.get("office_name", "")
    daily_limit = int(realtor.get("daily_post_limit") or 1)
    naver_realtor_id = str(realtor.get("naver_realtor_id") or "").strip()

    draft_buffer_target = get_draft_buffer_target(
        realtor,
        default_target=draft_buffer_default,
    )

    try:
        replenish_count = int(replenish_count or 0)
    except Exception:
        replenish_count = 5

    if replenish_count <= 0:
        replenish_count = draft_buffer_target

    print("-" * 80)
    print(f"[REALTOR] {office_name}")
    print(f"[REALTOR ID] {realtor_id}")
    print(f"[NAVER REALTOR ID] {naver_realtor_id}")
    print(f"[DAILY LIMIT] {daily_limit}")
    print(f"[DRAFT BUFFER TARGET] {draft_buffer_target}")
    print(f"[REPLENISH COUNT] {replenish_count}")

    if not naver_realtor_id:
        print("[SKIP] naver_realtor_id 없음")
        return

    today_count = get_today_scheduled_count(conn, realtor_id)
    active_drafts = count_active_draft_buffer(conn, realtor_id)
    pending_work = count_pending_work_queue(conn, realtor_id)

    current_buffer = active_drafts + pending_work
    # STEP362: 기존 초안/발행대기는 상세수집을 막지 않는다.
    # 파이프라인을 실행할 때마다 최신순 미처리 매물 1건을 상세수집 Queue에 넣는다.
    buffer_need = 1
    queue_target = 1
    fill_mode = "latest_each_run"

    print(f"[TODAY WORK QUEUE COUNT] {today_count}")
    print(f"[BUFFER FILL MODE] {fill_mode}")
    print(
        f"[DRAFT BUFFER STATUS] "
        f"active_drafts={active_drafts}, "
        f"pending_work={pending_work}, "
        f"current={current_buffer}, "
        f"target={draft_buffer_target}, "
        f"need={buffer_need}, "
        f"queue_target={queue_target}"
    )

    buffer_enough_skip_queue = False

    # STEP90-A:
    # 기존 find_latest_article()은 최초 /api/articles 응답 1회에서 보통 20건만 잡고 종료했다.
    # 운영형 버퍼에서는 "저장/큐 등록은 부족분만" 하되,
    # 후보 탐색은 중개사의 전체 노출 매물 범위까지 넓게 본다.
    result = finder.find_all_articles(
        naver_realtor_id=naver_realtor_id,
        wait_ms=wait_ms,
        max_pages=max(1, int(max_pages or 5)),
        # STEP117: 네이버 목록은 page 방식이 아니라 내부 UI 스크롤 방식이다.
        # max_pages 값은 기존 옵션명을 유지하되 스크롤 라운드 힌트로 사용한다.
        # 비공개 동기화가 켜진 전체수집은 기존 최소 8회를 유지한다.
        # 야간 초안 보충처럼 비공개 동기화를 끈 실행은 최신 1건 확보가 목적이므로
        # 최소 2회만 스크롤해 상세수집/초안생성 시간을 확보한다.
        scroll_steps=(
            max(8, int(max_pages or 8))
            if enable_private_sync
            else max(2, int(max_pages or 2))
        ),
    )

    candidates = result.get("candidates") or []

    print(f"[CANDIDATE COUNT] {len(candidates)}")
    print(f"[ALL ARTICLE COUNT] {result.get('count') or len(candidates)}")
    print(f"[VISIBLE TRADE COUNTS] {result.get('visible_trade_counts') or {}}")

    if not candidates:
        print("[NO ARTICLE]")
        return

    # STEP117 핵심:
    # 상세수집/초안 큐는 최신 1건만 넣더라도, 현재 살아있는 전체 article_no는 모두 캐시에 저장한다.
    # 이 데이터가 삭제매물 비공개 처리의 기준이 된다.
    cached_all = 0
    for article in candidates:
        article_no = str(article.get("article_no") or "").strip()
        if not article_no:
            continue
        try:
            if upsert_realtor_article_cache(conn=conn, realtor=realtor, article=article):
                cached_all += 1
        except Exception as e:
            conn.rollback()
            print(f"[CURRENT CACHE ERROR] {article_no} / {str(e)}")
            continue

    try:
        conn.commit()
    except Exception:
        pass

    print(f"[CURRENT ARTICLE CACHE SAVED] realtor_id={realtor_id}, count={cached_all}")

    # STEP119:
    # 현재 목록 전체 확보 후, 이미 발행되었지만 현재 목록에서 사라진 매물을 비공개 Queue로 보낸다.
    try:
        current_article_nos = [str(a.get("article_no") or "").strip() for a in candidates if str(a.get("article_no") or "").strip()]
        private_sync_result = sync_deleted_articles_to_private_queue(
            conn=conn,
            realtor=realtor,
            current_article_nos=current_article_nos,
            visible_counts=result.get("visible_trade_counts") or {},
            min_coverage_ratio=private_sync_min_coverage,
            enabled=enable_private_sync,
        )
        print("[PRIVATE SYNC RESULT]", private_sync_result)
    except Exception as e:
        conn.rollback()
        print(f"[PRIVATE SYNC WARN] realtor_id={realtor_id} / {str(e)}")

    queued = 0
    checked = 0

    for article in candidates:
        if queued >= queue_target:
            break

        article_no = str(article.get("article_no") or "").strip()

        if not article_no:
            continue

        checked += 1

        if is_processed_article(conn, article_no):
            print(f"[SKIP PROCESSED] {article_no}")
            continue

        try:
            upsert_realtor_article_cache(
                conn=conn,
                realtor=realtor,
                article=article,
            )

            conn.commit()

        except Exception as e:
            conn.rollback()
            print(f"[CACHE ERROR] {article_no} / {str(e)}")
            continue

        enqueue_article(
            conn=conn,
            realtor_id=realtor_id,
            article_no=article_no,
            priority=100,
        )

        queued += 1

        print(f"[QUEUED BUFFER] {article_no}")

    print(
        f"[SUMMARY] checked={checked}, queued={queued}, "
        f"queue_target={queue_target}, buffer_need={buffer_need}"
    )

    if queued == 0:
        print("[NO NEW ARTICLE] 후보는 있으나 모두 처리된 매물입니다.")


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

    parser.add_argument("--limit-realtors", type=int, default=100)
    parser.add_argument("--realtor-id", type=int, default=None)
    parser.add_argument("--headless", type=int, default=1)
    parser.add_argument("--wait-ms", type=int, default=2500)
    parser.add_argument("--draft-buffer", type=int, default=1)
    parser.add_argument("--replenish-count", type=int, default=1)
    parser.add_argument("--max-pages", type=int, default=5)
    parser.add_argument("--disable-private-sync", action="store_true", help="현재 목록 기반 삭제매물 비공개 Queue 생성을 끕니다")
    parser.add_argument("--private-sync-min-coverage", type=float, default=0.9, help="비공개 Sync 실행 최소 수집 커버리지. 기본 0.9")
    parser.add_argument("--lock-ttl-minutes", type=int, default=180)
    parser.add_argument("--no-db-lock", action="store_true")

    args = parser.parse_args()

    lock_name = "schedule_daily_articles"
    lock_owner = None
    run_id = None
    success = 0
    failed = 0
    skipped = 0
    total = 0
    final_status = "success"
    final_message = "schedule_daily_articles completed"

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

        if not locked:
            print("[SKIP] schedule_daily_articles already running")
            return

    conn = get_conn()

    finder = NaverRealtorResponseFinder(
        headless=bool(args.headless),
    )

    try:
        run_id = create_pipeline_run(
            "schedule_daily_articles",
            message="schedule_daily_articles started",
            meta={
                "limit_realtors": args.limit_realtors,
                "realtor_id": args.realtor_id,
                "draft_buffer": args.draft_buffer,
                "replenish_count": args.replenish_count,
                "headless": args.headless,
                "wait_ms": args.wait_ms,
                "draft_buffer": args.draft_buffer,
                "replenish_count": args.replenish_count,
                "enable_private_sync": not args.disable_private_sync,
                "private_sync_min_coverage": args.private_sync_min_coverage,
            },
        )

        print("=" * 80)
        print("[START] schedule_daily_articles")
        print("[STEP101-14] latest unpublished article only / initial fill=1 / replenish=1")
        print("TIME:", datetime.now().strftime("%Y-%m-%d %H:%M:%S"))
        print("=" * 80)

        realtors = fetch_target_realtors(
            conn=conn,
            limit_realtors=args.limit_realtors,
            realtor_id=args.realtor_id,
        )

        total = len(realtors)
        print(f"[TARGET REALTORS] {total}")

        add_pipeline_log(
            run_id=run_id,
            level="info",
            step_name="target_realtors",
            message=f"target realtors: {total}",
            realtor_id=args.realtor_id,
            context={"target_count": total},
        )

        finder.start()

        for realtor in realtors:
            try:
                process_realtor(
                    conn=conn,
                    finder=finder,
                    realtor=realtor,
                    wait_ms=args.wait_ms,
                    draft_buffer_default=args.draft_buffer,
                    replenish_count=args.replenish_count,
                    max_pages=args.max_pages,
                    enable_private_sync=not args.disable_private_sync,
                    private_sync_min_coverage=args.private_sync_min_coverage,
                )
                success += 1

            except Exception as e:
                failed += 1
                conn.rollback()
                print(f"[ERROR] realtor_id={realtor.get('id')} / {str(e)}")
                add_pipeline_log(
                    run_id=run_id,
                    level="error",
                    step_name="process_realtor",
                    message=str(e),
                    realtor_id=realtor.get("id"),
                    context={"traceback": traceback.format_exc()},
                )

        print("=" * 80)
        print("[DONE] schedule_daily_articles")
        print("=" * 80)

        if failed > 0:
            final_status = "failed"
            final_message = f"completed with failures: success={success}, failed={failed}"

    except Exception as e:
        final_status = "failed"
        final_message = str(e)[:2000]
        print("[FATAL]", final_message)
        traceback.print_exc()
        add_pipeline_log(
            run_id=run_id,
            level="error",
            step_name="fatal",
            message=final_message,
            context={"traceback": traceback.format_exc()},
        )

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

        try:
            conn.close()
        except Exception:
            pass

        finish_pipeline_run(
            run_id,
            final_status,
            total=total,
            success=success,
            failed=failed,
            skipped=skipped,
            message=final_message,
            meta={
                "limit_realtors": args.limit_realtors,
                "realtor_id": args.realtor_id,
                "draft_buffer": args.draft_buffer,
                "replenish_count": args.replenish_count,
            },
        )

        if lock_owner:
            release_db_lock(lock_name, lock_owner)


if __name__ == "__main__":
    main()
