# -*- coding: utf-8 -*-
# STEP107-02 V2 Latest Full Collector
# V1 schedule_daily_articles_v2.py는 수정하지 않고 V2 전용 파일로 사용한다.

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

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_v2 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):
    """
    V2:
    - 전체 자동 운영에서는 active + auto_publish_enabled=1만 대상.
    - 단, --realtor-id 직접 지정 테스트/수동 실행에서는 auto_publish_enabled=0이어도 포함한다.
    """
    sql = """
        SELECT
            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
        WHERE r.status = 'active'
    """

    params = []

    if realtor_id:
        sql += " AND r.id = %s "
        params.append(int(realtor_id))
    else:
        sql += " AND COALESCE(s.auto_publish_enabled, 0) = 1 "

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

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

    if realtor_id:
        print(
            f"[V2 TARGET REALTOR DIRECT] realtor_id={realtor_id}, "
            f"rows={len(rows)}, auto_publish_filter=bypassed"
        )

    return rows


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"

    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

        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()





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 update_realtor_article_summary_v2(conn, realtor_id, article):
    """
    STEP107-17:
    전체수집 시점에 확보 가능한 매물 요약정보를 blog_realtor_articles에 추가 저장한다.
    DB에 존재하는 컬럼만 안전하게 업데이트한다.
    """
    if not table_exists(conn, "blog_realtor_articles"):
        return False

    cols = table_columns(conn, "blog_realtor_articles")
    article_no = str(article.get("article_no") or "").strip()
    if not article_no:
        return False

    mapping = {
        "trade_type": article.get("trade_type"),
        "trade_type_name": article.get("trade_type_name"),
        "real_estate_type": article.get("real_estate_type"),
        "real_estate_type_name": article.get("real_estate_type_name"),
        "complex_name": article.get("complex_name"),
        "building_name": article.get("building_name"),
        "price_text": article.get("price_text"),
        "area_text": article.get("area_text"),
        "floor_text": article.get("floor_text"),
        "direction": article.get("direction"),
        "confirm_date": article.get("confirm_date"),
        "short_description": article.get("short_description"),
        "snapshot_date": datetime.now().strftime("%Y-%m-%d"),
        "last_snapshot_date": datetime.now().strftime("%Y-%m-%d"),
        "last_current_checked_at": datetime.now(),
        "summary_collected_at": datetime.now(),
        "updated_at": datetime.now(),
    }

    data = {k: v for k, v in mapping.items() if k in cols and v not in (None, "")}
    if not data:
        return False

    sets = ", ".join([f"{k}=%s" for k in data.keys()])
    params = list(data.values()) + [int(realtor_id), article_no]

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

    return True



def ensure_v2_snapshot_tables(conn):
    """STEP107-19: 운영용 Snapshot 테이블 생성."""
    with conn.cursor() as cur:
        cur.execute("""
            CREATE TABLE IF NOT EXISTS blog_realtor_article_snapshot_runs (
                id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT PRIMARY KEY,
                realtor_id INT NOT NULL,
                snapshot_date DATE NOT NULL,
                expected_count INT DEFAULT NULL,
                collected_count INT DEFAULT 0,
                visible_sale_count INT DEFAULT NULL,
                visible_lease_count INT DEFAULT NULL,
                visible_rent_count INT DEFAULT NULL,
                visible_total_count INT DEFAULT NULL,
                status VARCHAR(20) NOT NULL DEFAULT 'running',
                error_message TEXT NULL,
                started_at DATETIME NOT NULL,
                finished_at DATETIME NULL,
                created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
                updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
                KEY idx_realtor_date_status (realtor_id, snapshot_date, status),
                KEY idx_started_at (started_at)
            ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4
        """)
        cur.execute("""
            CREATE TABLE IF NOT EXISTS blog_realtor_article_snapshot_items (
                id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT PRIMARY KEY,
                run_id BIGINT UNSIGNED NOT NULL,
                realtor_id INT NOT NULL,
                snapshot_date DATE NOT NULL,
                article_no VARCHAR(32) NOT NULL,
                sort_order INT DEFAULT NULL,
                trade_type VARCHAR(50) DEFAULT NULL,
                trade_type_name VARCHAR(50) DEFAULT NULL,
                real_estate_type VARCHAR(50) DEFAULT NULL,
                real_estate_type_name VARCHAR(100) DEFAULT NULL,
                complex_name VARCHAR(255) DEFAULT NULL,
                building_name VARCHAR(255) DEFAULT NULL,
                price_text VARCHAR(255) DEFAULT NULL,
                area_text VARCHAR(255) DEFAULT NULL,
                floor_text VARCHAR(100) DEFAULT NULL,
                direction VARCHAR(100) DEFAULT NULL,
                confirm_date VARCHAR(50) DEFAULT NULL,
                short_description TEXT NULL,
                source_url TEXT NULL,
                created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
                UNIQUE KEY uq_run_article (run_id, article_no),
                KEY idx_realtor_date_article (realtor_id, snapshot_date, article_no),
                KEY idx_realtor_date_sort (realtor_id, snapshot_date, sort_order)
            ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4
        """)
    conn.commit()


def start_v2_snapshot_run(conn, realtor_id, visible_counts=None, expected_count=None):
    ensure_v2_snapshot_tables(conn)
    now = datetime.now()
    snapshot_date = now.strftime("%Y-%m-%d")
    visible_counts = visible_counts or {}
    with conn.cursor() as cur:
        cur.execute("""
            INSERT INTO blog_realtor_article_snapshot_runs
                (realtor_id, snapshot_date, expected_count, visible_sale_count,
                 visible_lease_count, visible_rent_count, visible_total_count,
                 status, started_at, created_at, updated_at)
            VALUES (%s, %s, %s, %s, %s, %s, %s, 'running', %s, %s, %s)
        """, (
            int(realtor_id), snapshot_date, expected_count,
            visible_counts.get("sale"), visible_counts.get("lease"), visible_counts.get("rent"), visible_counts.get("total"),
            now, now, now,
        ))
        run_id = cur.lastrowid
    conn.commit()
    print("[V2 SNAPSHOT RUN START]", {"run_id": run_id, "realtor_id": realtor_id, "snapshot_date": snapshot_date, "expected_count": expected_count})
    return run_id, snapshot_date


def save_v2_snapshot_items(conn, run_id, realtor_id, snapshot_date, candidates):
    ensure_v2_snapshot_tables(conn)
    saved = 0
    with conn.cursor() as cur:
        for idx, article in enumerate(candidates or [], start=1):
            article_no = str(article.get("article_no") or "").strip()
            if not article_no:
                continue
            cur.execute("""
                INSERT INTO blog_realtor_article_snapshot_items
                    (run_id, realtor_id, snapshot_date, article_no, sort_order,
                     trade_type, trade_type_name, real_estate_type, real_estate_type_name,
                     complex_name, building_name, price_text, area_text, floor_text,
                     direction, confirm_date, short_description, source_url, created_at)
                VALUES
                    (%s, %s, %s, %s, %s,
                     %s, %s, %s, %s,
                     %s, %s, %s, %s, %s,
                     %s, %s, %s, %s, %s)
                ON DUPLICATE KEY UPDATE
                    sort_order=VALUES(sort_order),
                    trade_type=VALUES(trade_type),
                    trade_type_name=VALUES(trade_type_name),
                    real_estate_type=VALUES(real_estate_type),
                    real_estate_type_name=VALUES(real_estate_type_name),
                    complex_name=VALUES(complex_name),
                    building_name=VALUES(building_name),
                    price_text=VALUES(price_text),
                    area_text=VALUES(area_text),
                    floor_text=VALUES(floor_text),
                    direction=VALUES(direction),
                    confirm_date=VALUES(confirm_date),
                    short_description=VALUES(short_description),
                    source_url=VALUES(source_url)
            """, (
                int(run_id), int(realtor_id), snapshot_date, article_no, idx,
                article.get("trade_type"), article.get("trade_type_name"),
                article.get("real_estate_type"), article.get("real_estate_type_name"),
                article.get("complex_name"), article.get("building_name"),
                article.get("price_text"), article.get("area_text"),
                article.get("floor_text"), article.get("direction"),
                article.get("confirm_date"), article.get("short_description"),
                article.get("source_url"), datetime.now(),
            ))
            saved += 1
    conn.commit()
    print("[V2 SNAPSHOT ITEMS SAVED]", {"run_id": run_id, "saved": saved})
    return saved


def finish_v2_snapshot_run(conn, run_id, status, collected_count=0, error_message=None):
    now = datetime.now()
    with conn.cursor() as cur:
        cur.execute("""
            UPDATE blog_realtor_article_snapshot_runs
            SET status=%s, collected_count=%s, error_message=%s, finished_at=%s, updated_at=%s
            WHERE id=%s
        """, (status, int(collected_count or 0), error_message, now, now, int(run_id)))
    conn.commit()
    print("[V2 SNAPSHOT RUN FINISH]", {"run_id": run_id, "status": status, "collected_count": collected_count})

def cache_all_current_articles_v2(conn, realtor, candidates):
    """
    V2 핵심:
    - queue_target과 무관하게 현재 네이버 매물 후보 전체를 blog_realtor_articles에 cache한다.
    - 삭제매물 판단은 이 current cache가 충분히 확보된 뒤에만 사용한다.
    """
    cached = 0
    failed = 0

    for article in candidates or []:
        article_no = str(article.get("article_no") or "").strip()
        if not article_no:
            continue

        try:
            upsert_realtor_article_cache(
                conn=conn,
                realtor=realtor,
                article=article,
            )
            try:
                update_realtor_article_summary_v2(
                    conn=conn,
                    realtor_id=int(realtor.get("id")),
                    article=article,
                )
            except Exception as e:
                print(f"[V2 SUMMARY UPDATE WARN] {article_no} / {str(e)}")
            cached += 1
        except Exception as e:
            failed += 1
            try:
                conn.rollback()
            except Exception:
                pass
            print(f"[V2 CURRENT CACHE ERROR] {article_no} / {str(e)}")

    try:
        conn.commit()
    except Exception:
        pass

    print(f"[V2 CURRENT CACHE DONE] cached={cached}, failed={failed}")
    return {
        "cached": cached,
        "failed": failed,
    }


def warn_if_current_count_suspicious_v2(candidates, realtor=None):
    """
    현재 수집 건수가 너무 적으면 삭제매물 Sync에서 사용하면 위험하다.
    이 함수는 우선 로그 경고만 남긴다.
    """
    count = len(candidates or [])
    office_name = (realtor or {}).get("office_name", "")
    naver_realtor_id = (realtor or {}).get("naver_realtor_id", "")

    if count <= 25:
        print(
            "[V2 CURRENT COUNT WARN] "
            f"current_naver_count={count}, office={office_name}, realtorId={naver_realtor_id}. "
            "전체 매물 수집이 덜 되었을 가능성이 있으므로 deleted_sync 판단에 주의 필요."
        )

    return count

def process_realtor(
    conn,
    finder,
    realtor,
    wait_ms=2500,
    draft_buffer_default=1,
    replenish_count=1,
    max_pages=30,
):
    """
    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
    buffer_need = max(0, draft_buffer_target - current_buffer)

    # STEP101-10:
    # 최초 등록/초기 재고가 0인 중개사는 draft_buffer_target(기본 30)까지 한 번에 채운다.
    # 이미 재고가 있는 기존 중개사는 replenish_count(운영 기본 1)만 보충한다.
    # STEP101-14:
    # 신규/초기 재고가 0인 중개사도 최초 1건만 확보한다.
    # 기존 30건/45건 버퍼 방식은 수집·초안 생성 시간이 길고,
    # 실제 운영 정책도 '최신 미발행 매물 1건씩 확보'로 변경되었다.
    if current_buffer <= 0:
        queue_target = 1 if buffer_need > 0 else 0
        fill_mode = "initial_fill_one"
    else:
        queue_target = min(buffer_need, replenish_count)
        fill_mode = "replenish"

    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}"
    )

    # STEP109-05:
    # 여기서 return하면 전체매물 find_all_articles / snapshot / current cache가 모두 실행되지 않는다.
    # V2 Workday에서는 초안 버퍼가 충분해도 삭제매물/발행 guard용 최신 snapshot은 반드시 필요하다.
    # 따라서 queue_target <= 0이어도 전체매물 수집과 snapshot 저장은 계속 진행하고,
    # 아래 큐 등록 단계에서만 신규 큐잉을 건너뛴다.
    buffer_enough_skip_queue = queue_target <= 0
    if buffer_enough_skip_queue:
        print("[BUFFER ENOUGH] queue registration will be skipped after snapshot/current cache")

    # 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)),
        scroll_steps=1,
    )

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

    print(f"[CANDIDATE COUNT] {len(candidates)}")
    print(f"[ALL ARTICLE COUNT] {result.get('count') or len(candidates)}")
    print(f"[V2 CURRENT NAVER COUNT] {len(candidates)}")

    warn_if_current_count_suspicious_v2(candidates, realtor=realtor)

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

    # STEP107-19: 성공/실패가 남는 운영용 Snapshot 저장.
    snapshot_run_id = None
    try:
        visible_counts = result.get("visible_trade_counts") or {}
        expected_count = visible_counts.get("total") or len(candidates)
        snapshot_run_id, snapshot_date = start_v2_snapshot_run(
            conn=conn,
            realtor_id=int(realtor.get("id")),
            visible_counts=visible_counts,
            expected_count=expected_count,
        )
        save_v2_snapshot_items(
            conn=conn,
            run_id=snapshot_run_id,
            realtor_id=int(realtor.get("id")),
            snapshot_date=snapshot_date,
            candidates=candidates,
        )
    except Exception as e:
        if snapshot_run_id:
            try:
                finish_v2_snapshot_run(conn, snapshot_run_id, "failed", 0, str(e))
            except Exception:
                pass
        print("[V2 SNAPSHOT ERROR]", str(e))

    # STEP107-02:
    # 현재 네이버 매물 전체 목록은 queue_target과 무관하게 먼저 cache한다.
    # 이후 queue 등록은 최신순 미처리 매물 1건만 진행해도 된다.
    cache_result = cache_all_current_articles_v2(conn, realtor, candidates)
    print("[V2 CURRENT CACHE RESULT]", json.dumps(cache_result, ensure_ascii=False, default=str))

    if snapshot_run_id:
        try:
            failed_count = int(cache_result.get("failed") or 0)
            finish_v2_snapshot_run(
                conn=conn,
                run_id=snapshot_run_id,
                status="success" if failed_count == 0 else "partial",
                collected_count=len(candidates),
                error_message=None if failed_count == 0 else json.dumps(cache_result, ensure_ascii=False, default=str),
            )
        except Exception as e:
            print("[V2 SNAPSHOT FINISH WARN]", str(e))

    if buffer_enough_skip_queue:
        print(
            f"[QUEUE SKIP] draft buffer enough after snapshot/current cache. "
            f"checked_candidates={len(candidates)}, queue_target={queue_target}, buffer_need={buffer_need}"
        )
        return

    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

        # V2: 전체 current cache는 이미 완료했다. 여기서는 큐 등록만 한다.
        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=30)
    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_v2"
    lock_owner = None
    run_id = None
    success = 0
    failed = 0
    skipped = 0
    total = 0
    final_status = "success"
    final_message = "schedule_daily_articles_v2 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_v2 already running")
            return

    conn = get_conn()

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

    try:
        run_id = create_pipeline_run(
            "schedule_daily_articles_v2",
            message="schedule_daily_articles_v2 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,
            },
        )

        print("=" * 80)
        print("[START] schedule_daily_articles_v2")
        print("[STEP107-07] V2 bottom-based internal scroll / visible count check / latest sort")
        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,
                )
                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_v2")
        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()