# -*- coding: utf-8 -*-
"""
STEP105 Realtor Pipeline Core

목표:
- 기존 운영 배치(daily_pipeline_runner 등)는 건드리지 않는다.
- 중개사 1명 기준 원스톱 엔진을 별도 core로 재구성한다.
- 가능한 것은 함수 직접 호출로 처리한다.
- 아직 함수화되지 않은 상세수집/발행은 adapter 형태로 분리해 이후 직접 함수 호출로 교체 가능하게 한다.

현재 STEP105-01 범위:
1) 중개사/정책 로드
2) 테스트 계정 정책 분기
3) 수집 후보 등록: jobs.schedule_daily_articles.process_realtor 직접 호출
4) 초안 생성: jobs.generate_blog_drafts.process_article 직접 호출
5) 큐 상태 점검
6) 발행은 publish_worker_step104.py를 임시 adapter로 호출
   - 다음 STEP105-02에서 publish_one_queue 함수로 분리 예정
"""

import os
import sys
import json
import time
import subprocess
from pathlib import Path
from datetime import datetime

BASE_DIR = Path(__file__).resolve().parents[1]
if str(BASE_DIR) not in sys.path:
    sys.path.append(str(BASE_DIR))

from db import get_conn


PYTHON_EXE = sys.executable


def clean_text(value):
    return str(value or "").strip()


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


def table_columns(conn, table_name):
    try:
        with conn.cursor() as cur:
            cur.execute(f"SHOW COLUMNS FROM {table_name}")
            rows = cur.fetchall()
        return {row["Field"] for row in rows}
    except Exception:
        return set()


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


def load_realtor(conn, realtor_id):
    with conn.cursor() as cur:
        cur.execute("SELECT * FROM blog_realtors WHERE id=%s LIMIT 1", (int(realtor_id),))
        row = cur.fetchone()
    if not row:
        raise RuntimeError(f"realtor not found: {realtor_id}")
    return row


def load_publish_setting(conn, realtor_id):
    if not table_exists(conn, "blog_publish_settings"):
        return {}

    with conn.cursor() as cur:
        cur.execute("SELECT * FROM blog_publish_settings WHERE realtor_id=%s LIMIT 1", (int(realtor_id),))
        return cur.fetchone() or {}


def get_policy(conn, realtor_id):
    realtor = load_realtor(conn, realtor_id)
    setting = load_publish_setting(conn, realtor_id)

    publish_mode = clean_text(realtor.get("publish_mode") or "server_auto")
    try:
        is_test_account = int(realtor.get("is_test_account") or 0)
    except Exception:
        is_test_account = 0

    try:
        auto_publish_enabled = int(setting.get("auto_publish_enabled") or 0)
    except Exception:
        auto_publish_enabled = 0

    return {
        "realtor": realtor,
        "setting": setting,
        "publish_mode": publish_mode,
        "is_test_account": is_test_account,
        "auto_publish_enabled": auto_publish_enabled,
    }


def set_publish_mode(conn, realtor_id, mode):
    if not mode:
        return

    columns = table_columns(conn, "blog_realtors")
    if "publish_mode" not in columns:
        log("[WARN] publish_mode column not found")
        return

    with conn.cursor() as cur:
        cur.execute("""
            UPDATE blog_realtors
            SET publish_mode=%s,
                updated_at=NOW()
            WHERE id=%s
        """, (str(mode), int(realtor_id)))
    conn.commit()


def ensure_publish_setting_row(conn, realtor_id):
    if not table_exists(conn, "blog_publish_settings"):
        return

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

    if row:
        return

    columns = table_columns(conn, "blog_publish_settings")
    data = {
        "realtor_id": int(realtor_id),
        "auto_publish_enabled": 0,
        "daily_post_limit": 1,
        "created_at": datetime.now(),
        "updated_at": datetime.now(),
    }
    filtered = {k: v for k, v in data.items() if k in columns}

    if not filtered:
        return

    keys = list(filtered.keys())
    with conn.cursor() as cur:
        cur.execute(
            f"INSERT INTO blog_publish_settings ({', '.join(keys)}) VALUES ({', '.join(['%s'] * len(keys))})",
            [filtered[k] for k in keys],
        )
    conn.commit()


def set_auto_publish_enabled(conn, realtor_id, enabled):
    if not table_exists(conn, "blog_publish_settings"):
        return

    ensure_publish_setting_row(conn, realtor_id)

    columns = table_columns(conn, "blog_publish_settings")
    if "auto_publish_enabled" not in columns:
        return

    sets = ["auto_publish_enabled=%s"]
    params = [1 if enabled else 0]

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

    params.append(int(realtor_id))

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


def check_session_db(conn, realtor_id):
    """
    DB 세션 정보는 참고만 한다.
    실제 로그인 판정은 publish_worker/publisher가 브라우저에서 한다.
    """
    if not table_exists(conn, "blog_naver_sessions"):
        return {"exists": False, "ok": None, "message": "table missing"}

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

    if not row:
        return {"exists": False, "ok": False, "message": "session row missing"}

    status = clean_text(row.get("status") or row.get("naver_session_status") or row.get("login_status") or "")
    session_file = clean_text(row.get("session_file") or "")

    ok_statuses = {"login_ok", "linked", "ok", "active", "unknown"}

    return {
        "exists": True,
        "ok": status in ok_statuses if status else False,
        "status": status,
        "session_file": session_file,
        "message": f"status={status}, session_file={session_file}",
    }


def count_pending_work(conn, realtor_id):
    if not table_exists(conn, "blog_article_work_queue"):
        return 0
    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),))
        return int((cur.fetchone() or {}).get("cnt") or 0)



def count_latest_success_snapshot_articles(conn, realtor_id):
    """
    STEP107-23:
    V2 운영 기준 현재 매물 수는 blog_realtor_articles cache가 아니라
    최신 SUCCESS Snapshot item 수를 우선 사용한다.
    """
    if not table_exists(conn, "blog_realtor_article_snapshot_runs"):
        return None
    if not table_exists(conn, "blog_realtor_article_snapshot_items"):
        return None

    try:
        with conn.cursor() as cur:
            cur.execute("""
                SELECT id
                FROM blog_realtor_article_snapshot_runs
                WHERE realtor_id=%s
                  AND status='success'
                ORDER BY snapshot_date DESC, id DESC
                LIMIT 1
            """, (int(realtor_id),))
            run = cur.fetchone()

            if not run:
                return None

            cur.execute("""
                SELECT COUNT(DISTINCT article_no) AS cnt
                FROM blog_realtor_article_snapshot_items
                WHERE run_id=%s
                  AND COALESCE(article_no,'') <> ''
            """, (int(run.get("id")),))
            row = cur.fetchone() or {}

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

    except Exception:
        return None


def fetch_latest_success_snapshot_summary(conn, realtor_id):
    """
    운영 로그용 최신 Snapshot 요약.
    """
    if not table_exists(conn, "blog_realtor_article_snapshot_runs"):
        return None

    try:
        with conn.cursor() as cur:
            cur.execute("""
                SELECT *
                FROM blog_realtor_article_snapshot_runs
                WHERE realtor_id=%s
                  AND status='success'
                ORDER BY snapshot_date DESC, id DESC
                LIMIT 1
            """, (int(realtor_id),))
            run = cur.fetchone()

        if not run:
            return None

        return {
            "run_id": run.get("id"),
            "snapshot_date": str(run.get("snapshot_date")),
            "expected_count": run.get("expected_count"),
            "collected_count": run.get("collected_count"),
            "visible_total_count": run.get("visible_total_count"),
            "visible_sale_count": run.get("visible_sale_count"),
            "visible_lease_count": run.get("visible_lease_count"),
            "visible_rent_count": run.get("visible_rent_count"),
        }

    except Exception:
        return None


def count_collected_articles(conn, realtor_id):
    """
    STEP107-23:
    V2 운영 기준 현재 매물 수는 최신 SUCCESS Snapshot 기준이다.
    Snapshot이 아직 없을 때만 blog_realtor_articles active/public count로 fallback한다.
    """
    snapshot_count = count_latest_success_snapshot_articles(conn, realtor_id)
    if snapshot_count is not None:
        return snapshot_count

    if not table_exists(conn, "blog_realtor_articles"):
        return 0

    with conn.cursor() as cur:
        cur.execute("""
            SELECT COUNT(DISTINCT article_no) AS cnt
            FROM blog_realtor_articles
            WHERE realtor_id=%s
              AND COALESCE(article_no,'') <> ''
              AND COALESCE(article_status,'active') = 'active'
              AND COALESCE(visibility_status,'public') = 'public'
        """, (int(realtor_id),))
        return int((cur.fetchone() or {}).get("cnt") or 0)


def count_ready_drafts(conn, realtor_id):
    if not table_exists(conn, "blog_article_drafts"):
        return 0
    with conn.cursor() as cur:
        cur.execute("""
            SELECT COUNT(*) AS cnt
            FROM blog_article_drafts
            WHERE realtor_id=%s
              AND COALESCE(status,'ready') IN ('ready','approved','pending')
        """, (int(realtor_id),))
        return int((cur.fetchone() or {}).get("cnt") or 0)


def count_pending_publish(conn, realtor_id):
    if not table_exists(conn, "blog_publish_queue"):
        return 0
    with conn.cursor() as cur:
        cur.execute("""
            SELECT COUNT(*) AS cnt
            FROM blog_publish_queue
            WHERE realtor_id=%s
              AND COALESCE(queue_status,'') IN ('pending','waiting','ready','failed')
        """, (int(realtor_id),))
        return int((cur.fetchone() or {}).get("cnt") or 0)


def count_published_posts(conn, realtor_id):
    if not table_exists(conn, "blog_publish_queue"):
        return 0
    with conn.cursor() as cur:
        cur.execute("""
            SELECT COUNT(DISTINCT article_no) AS cnt
            FROM blog_publish_queue
            WHERE realtor_id=%s
              AND COALESCE(article_no,'') <> ''
              AND COALESCE(queue_status,'')='published'
        """, (int(realtor_id),))
        return int((cur.fetchone() or {}).get("cnt") or 0)


def count_private_queue(conn, realtor_id):
    """
    STEP107-12-FIX:
    blog_post_private_queue 테이블에 status 컬럼이 없는 운영 DB가 있다.
    없는 컬럼을 COALESCE 안에 넣어도 MySQL은 Unknown column 에러를 낸다.
    따라서 실제 존재 컬럼 기준으로 status expression을 만든다.
    """
    if not table_exists(conn, "blog_post_private_queue"):
        return 0

    cols = table_columns(conn, "blog_post_private_queue")
    if "queue_status" in cols:
        status_expr = "COALESCE(queue_status,'')"
    elif "status" in cols:
        status_expr = "COALESCE(status,'')"
    else:
        return 0

    with conn.cursor() as cur:
        cur.execute(f"""
            SELECT COUNT(*) AS cnt
            FROM blog_post_private_queue
            WHERE realtor_id=%s
              AND {status_expr} IN ('pending','waiting','ready','processing','failed')
        """, (int(realtor_id),))
        return int((cur.fetchone() or {}).get("cnt") or 0)


def count_private_done(conn, realtor_id):
    """
    STEP107-12-FIX:
    private_queue 완료 카운트도 실제 존재 컬럼 기준으로 조회한다.
    """
    if not table_exists(conn, "blog_post_private_queue"):
        return 0

    cols = table_columns(conn, "blog_post_private_queue")
    if "queue_status" in cols:
        status_expr = "COALESCE(queue_status,'')"
    elif "status" in cols:
        status_expr = "COALESCE(status,'')"
    else:
        return 0

    with conn.cursor() as cur:
        cur.execute(f"""
            SELECT COUNT(*) AS cnt
            FROM blog_post_private_queue
            WHERE realtor_id=%s
              AND {status_expr} IN ('private_done','done','success')
        """, (int(realtor_id),))
        return int((cur.fetchone() or {}).get("cnt") or 0)


def count_current_cache_raw(conn, realtor_id):
    if not table_exists(conn, "blog_realtor_articles"):
        return 0
    with conn.cursor() as cur:
        cur.execute("""
            SELECT COUNT(DISTINCT article_no) AS cnt
            FROM blog_realtor_articles
            WHERE realtor_id=%s
              AND COALESCE(article_no,'') <> ''
        """, (int(realtor_id),))
        return int((cur.fetchone() or {}).get("cnt") or 0)


def build_v2_ops_counts(conn, realtor_id):
    snapshot_summary = fetch_latest_success_snapshot_summary(conn, realtor_id)
    return {
        "pending_work": count_pending_work(conn, realtor_id),
        "current_naver_articles": count_collected_articles(conn, realtor_id),
        "current_cache_raw": count_current_cache_raw(conn, realtor_id),
        "latest_snapshot": snapshot_summary,
        "ready_drafts": count_ready_drafts(conn, realtor_id),
        "pending_publish": count_pending_publish(conn, realtor_id),
        "published_posts": count_published_posts(conn, realtor_id),
        "private_queue": count_private_queue(conn, realtor_id),
        "private_done": count_private_done(conn, realtor_id),
    }


def print_counts(conn, realtor_id, label):
    """
    STEP107-12:
    V2 운영 로그 카운트를 한 화면에서 보이도록 통일한다.
    """
    counts = build_v2_ops_counts(conn, realtor_id)

    log("-" * 80)
    log(f"[V2 OPS COUNTS] {label}")
    log("pending_work:", counts["pending_work"])
    log("current_naver_articles:", counts["current_naver_articles"], "(latest SUCCESS snapshot 기준)")
    log("current_cache_raw:", counts["current_cache_raw"], "(참고용 legacy cache)")
    if counts.get("latest_snapshot"):
        log("latest_snapshot:", json.dumps(counts["latest_snapshot"], ensure_ascii=False, default=str))
    log("ready_drafts:", counts["ready_drafts"])
    log("pending_publish:", counts["pending_publish"])
    log("published_posts:", counts["published_posts"])
    log("private_queue:", counts["private_queue"])
    log("private_done:", counts["private_done"])
    log("-" * 80)


def fetch_latest_publish_queue(conn, realtor_id):
    if not table_exists(conn, "blog_publish_queue"):
        return None
    with conn.cursor() as cur:
        cur.execute("""
            SELECT *
            FROM blog_publish_queue
            WHERE realtor_id=%s
            ORDER BY id DESC
            LIMIT 1
        """, (int(realtor_id),))
        return cur.fetchone()


def shift_today_published_for_test(conn, realtor_id):
    if not table_exists(conn, "blog_publish_queue"):
        return 0
    with conn.cursor() as cur:
        cur.execute("""
            UPDATE blog_publish_queue
            SET started_at = CASE WHEN started_at IS NOT NULL THEN DATE_SUB(started_at, INTERVAL 1 DAY) ELSE started_at END,
                finished_at = CASE WHEN finished_at IS NOT NULL THEN DATE_SUB(finished_at, INTERVAL 1 DAY) ELSE finished_at END,
                updated_at = NOW()
            WHERE realtor_id=%s
              AND queue_status='published'
              AND finished_at >= CURDATE()
              AND finished_at < DATE_ADD(CURDATE(), INTERVAL 1 DAY)
        """, (int(realtor_id),))
        affected = cur.rowcount
    conn.commit()
    return int(affected or 0)


# ---------------------------------------------------------------------
# STEP105 direct function adapters
# ---------------------------------------------------------------------

def collect_candidate_for_realtor(conn, realtor, draft_buffer=1, replenish_count=1, wait_ms=2500, max_pages=30):
    """
    STEP107-10:
    V2 Workday 수집은 V1 jobs.schedule_daily_articles.py를 사용하지 않는다.

    반드시 V2 전용 수집기:
      jobs.schedule_daily_articles_v2.process_realtor
      services.naver_realtor_response_finder_v2.NaverRealtorResponseFinder

    를 사용한다.

    이유:
    - V1 수집기는 내부 스크롤을 바닥까지 밀지 못해 22건에서 멈출 수 있다.
    - V2 수집기는 실제 네이버 표시 총계(예: 매매47/전세1/월세5=53)와 일치하도록 전체 수집한다.
    """
    from jobs.schedule_daily_articles_v2 import NaverRealtorResponseFinder, process_realtor

    finder = NaverRealtorResponseFinder(headless=True)

    try:
        print("[V2 CORE COLLECT START] use schedule_daily_articles_v2", flush=True)
        finder.start()
        process_realtor(
            conn=conn,
            finder=finder,
            realtor=realtor,
            wait_ms=wait_ms,
            draft_buffer_default=draft_buffer,
            replenish_count=replenish_count,
            max_pages=max_pages,
        )
        print("[V2 CORE COLLECT DONE]", flush=True)
    finally:
        try:
            finder.close()
        except Exception:
            pass


def fetch_details_adapter(limit=3, realtor_id=None):
    """
    STEP107-28:
    V2 전용 상세수집기 호출.
    - V1 naver_response_article_fetcher.py는 수정하지 않는다.
    - services/naver_response_article_fetcher_v2.py를 호출한다.
    - realtor_id가 있으면 해당 중개사의 최신 SUCCESS Snapshot 순서 기준으로 상세수집한다.
    """
    cmd = [
        PYTHON_EXE,
        "services/naver_response_article_fetcher_v2.py",
        "--limit",
        str(int(limit or 1)),
        "--no-db-lock",
    ]

    if realtor_id is not None:
        cmd.extend(["--realtor-id", str(int(realtor_id))])

    return run_subprocess_adapter("fetch_details_adapter", cmd, timeout=1800)




def mark_v2_fetched_articles_collected(conn, realtor_id, limit=20):
    """
    STEP107-26:
    상세수집이 성공했는데 blog_realtor_articles.collect_status가 pending으로 남는 경우를 보정한다.

    기준:
    - blog_article_images에 이미지가 있거나
    - 상세 JSON/상세 수집 시각 컬럼이 채워져 있으면
    collect_status='collected'로 올린다.

    이 함수는 fetch_details_adapter 실행 직후 호출하기 위한 안전 보정 함수다.
    """
    if not table_exists(conn, "blog_realtor_articles"):
        return {"updated": 0, "checked": 0}

    cols = table_columns(conn, "blog_realtor_articles")
    if "collect_status" not in cols:
        return {"updated": 0, "checked": 0, "reason": "collect_status column missing"}

    # 후보: 최신 Snapshot 순서상 pending인 매물
    snapshot_join = ""
    params_prefix = []
    order_sql = "ORDER BY CAST(a.article_no AS UNSIGNED) DESC"

    if table_exists(conn, "blog_realtor_article_snapshot_runs") and table_exists(conn, "blog_realtor_article_snapshot_items"):
        snapshot_join = """
            INNER JOIN (
                SELECT si.article_no, si.sort_order
                FROM blog_realtor_article_snapshot_items si
                INNER JOIN blog_realtor_article_snapshot_runs sr
                    ON sr.id = si.run_id
                WHERE sr.realtor_id = %s
                  AND sr.status = 'success'
                  AND sr.id = (
                      SELECT id
                      FROM blog_realtor_article_snapshot_runs
                      WHERE realtor_id = %s
                        AND status = 'success'
                      ORDER BY snapshot_date DESC, id DESC
                      LIMIT 1
                  )
            ) snap
              ON CAST(snap.article_no AS UNSIGNED) = CAST(a.article_no AS UNSIGNED)
        """
        params_prefix = [int(realtor_id), int(realtor_id)]
        order_sql = "ORDER BY COALESCE(snap.sort_order, 999999) ASC"

    with conn.cursor() as cur:
        cur.execute(f"""
            SELECT a.article_no
            FROM blog_realtor_articles a
            {snapshot_join}
            WHERE a.realtor_id=%s
              AND COALESCE(a.article_no,'') <> ''
              AND COALESCE(a.collect_status,'') NOT IN ('collected','detail_collected','done','success')
            {order_sql}
            LIMIT %s
        """, params_prefix + [int(realtor_id), int(limit or 20)])
        candidates = cur.fetchall() or []

    updated = 0
    checked = 0

    for row in candidates:
        article_no = clean_text(row.get("article_no"))
        if not article_no:
            continue

        checked += 1
        evidence = False

        # 1) 이미지 테이블 근거
        if table_exists(conn, "blog_article_images"):
            try:
                with conn.cursor() as cur:
                    cur.execute("""
                        SELECT COUNT(*) AS cnt
                        FROM blog_article_images
                        WHERE article_no=%s
                    """, (article_no,))
                    if int((cur.fetchone() or {}).get("cnt") or 0) > 0:
                        evidence = True
            except Exception:
                pass

        # 2) realtor_articles 자체 상세 컬럼 근거
        if not evidence:
            detail_checks = []
            for col in ["article_detail_json", "detail_json", "detail_collected_at", "last_detail_collected_at"]:
                if col in cols:
                    detail_checks.append(f"{col} IS NOT NULL AND COALESCE({col},'') <> ''")

            if detail_checks:
                try:
                    with conn.cursor() as cur:
                        cur.execute(f"""
                            SELECT article_no
                            FROM blog_realtor_articles
                            WHERE realtor_id=%s
                              AND article_no=%s
                              AND ({" OR ".join(detail_checks)})
                            LIMIT 1
                        """, (int(realtor_id), article_no))
                        if cur.fetchone():
                            evidence = True
                except Exception:
                    pass

        if not evidence:
            continue

        sets = ["collect_status='collected'"]
        if "detail_collected" in cols:
            sets.append("detail_collected=1")
        if "detail_collected_at" in cols:
            sets.append("detail_collected_at=NOW()")
        if "last_detail_collected_at" in cols:
            sets.append("last_detail_collected_at=NOW()")
        if "updated_at" in cols:
            sets.append("updated_at=NOW()")

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

        updated += 1
        log("[V2 FETCH STATUS SYNC]", article_no, "collect_status=collected")

    conn.commit()
    result = {"updated": updated, "checked": checked}
    log("[V2 FETCH STATUS SYNC RESULT]", result)
    return result


def fetch_target_articles_v2(conn, realtor_id, limit=1):
    """
    STEP107-25:
    V2 전용 Draft 대상 선택.

    기준:
    - 최신 SUCCESS Snapshot 순서
    - blog_realtor_articles 상세수집 완료
    - blog_drafts에 초안 없음
    - blog_publish_queue에서 published/success/done 아님
    - Collation 문제 방지를 위해 article_no 비교는 CAST(... AS UNSIGNED) 사용
    """
    if not table_exists(conn, "blog_realtor_articles"):
        return []

    article_cols = table_columns(conn, "blog_realtor_articles")

    # 상세수집 완료 판단:
    # generate_blog_drafts.is_valid_collected_article_for_draft()가 collect_status를 최종 검증하므로
    # collect_status 컬럼이 있으면 반드시 collected 계열만 대상이 되게 한다.
    if "collect_status" in article_cols:
        detail_sql = "COALESCE(a.collect_status,'') IN ('collected','detail_collected','done','success')"
    else:
        detail_conditions = []
        if "detail_collected" in article_cols:
            detail_conditions.append("COALESCE(a.detail_collected,0)=1")
        if "article_detail_json" in article_cols:
            detail_conditions.append("COALESCE(a.article_detail_json,'') <> ''")
        if "detail_json" in article_cols:
            detail_conditions.append("COALESCE(a.detail_json,'') <> ''")
        if "detail_collected_at" in article_cols:
            detail_conditions.append("a.detail_collected_at IS NOT NULL")
        if "last_detail_collected_at" in article_cols:
            detail_conditions.append("a.last_detail_collected_at IS NOT NULL")

        if not detail_conditions:
            log("[V2 DRAFT TARGETS]", "no detail status columns found")
            return []

        detail_sql = "(" + " OR ".join(detail_conditions) + ")"

    # 최신 SUCCESS Snapshot subquery
    snapshot_join = ""
    snapshot_where = ""
    params_prefix = []
    order_sql = "ORDER BY CAST(a.article_no AS UNSIGNED) DESC"

    if table_exists(conn, "blog_realtor_article_snapshot_runs") and table_exists(conn, "blog_realtor_article_snapshot_items"):
        snapshot_join = """
            INNER JOIN (
                SELECT si.article_no, si.sort_order
                FROM blog_realtor_article_snapshot_items si
                INNER JOIN blog_realtor_article_snapshot_runs sr
                    ON sr.id = si.run_id
                WHERE sr.realtor_id = %s
                  AND sr.status = 'success'
                  AND sr.id = (
                      SELECT id
                      FROM blog_realtor_article_snapshot_runs
                      WHERE realtor_id = %s
                        AND status = 'success'
                      ORDER BY snapshot_date DESC, id DESC
                      LIMIT 1
                  )
            ) snap
              ON CAST(snap.article_no AS UNSIGNED) = CAST(a.article_no AS UNSIGNED)
        """
        params_prefix = [int(realtor_id), int(realtor_id)]
        order_sql = "ORDER BY COALESCE(snap.sort_order, 999999) ASC, CAST(a.article_no AS UNSIGNED) DESC"

    draft_join = ""
    draft_where = ""
    if table_exists(conn, "blog_drafts"):
        draft_cols = table_columns(conn, "blog_drafts")
        if "article_no" in draft_cols:
            draft_join = """
                LEFT JOIN blog_drafts d
                  ON d.realtor_id = a.realtor_id
                 AND CAST(d.article_no AS UNSIGNED) = CAST(a.article_no AS UNSIGNED)
            """
            draft_where = "AND d.id IS NULL"

    publish_join = ""
    publish_where = ""
    if table_exists(conn, "blog_publish_queue"):
        publish_join = """
            LEFT JOIN blog_publish_queue pq_done
              ON pq_done.realtor_id = a.realtor_id
             AND CAST(pq_done.article_no AS UNSIGNED) = CAST(a.article_no AS UNSIGNED)
             AND COALESCE(pq_done.queue_status,'') IN ('published','success','done')
        """
        publish_where = "AND pq_done.id IS NULL"

    sql = f"""
        SELECT a.*
        FROM blog_realtor_articles a
        {snapshot_join}
        {draft_join}
        {publish_join}
        WHERE a.realtor_id = %s
          AND COALESCE(a.article_no,'') <> ''
          AND {detail_sql}
          {draft_where}
          {publish_where}
        {order_sql}
        LIMIT %s
    """

    params = params_prefix + [int(realtor_id), int(limit or 1)]

    try:
        with conn.cursor() as cur:
            cur.execute(sql, params)
            rows = cur.fetchall() or []
    except Exception as e:
        log("[V2 DRAFT TARGETS SQL ERROR]", str(e))
        log("[V2 DRAFT TARGETS SQL]", sql)
        raise

    log("[V2 DRAFT TARGETS]", len(rows))
    for row in rows:
        log("[V2 DRAFT TARGET]", clean_text(row.get("article_no")))

    return rows


# 호환용 이름 유지
def fetch_v2_draft_targets_fallback(conn, realtor_id, limit=1):
    return fetch_target_articles_v2(conn, realtor_id, limit=limit)


def generate_draft_direct(conn, realtor_id, limit=1, draft_buffer=1, replenish_count=1):
    """
    generate_blog_drafts.py의 fetch_target_articles/process_article을 직접 호출한다.
    """
    from jobs.generate_blog_drafts import fetch_target_articles, is_valid_collected_article_for_draft, process_article

    # STEP107-25:
    # V2 Workday에서는 Snapshot 중심 V2 대상 선택을 우선 사용한다.
    articles = fetch_target_articles_v2(
        conn=conn,
        realtor_id=int(realtor_id),
        limit=int(limit or 1),
    )

    log("[DRAFT DIRECT V2 TARGETS]", len(articles))

    if not articles:
        log("[DRAFT DIRECT LEGACY FALLBACK]", "V2 target query returned 0, trying legacy fetch_target_articles")
        articles = fetch_target_articles(
            conn=conn,
            realtor_id=int(realtor_id),
            article_no=None,
            limit=int(limit or 1),
            force=False,
            draft_buffer=int(draft_buffer or 1),
            replenish_count=int(replenish_count or 1),
            limit_realtors=1,
        )
        log("[DRAFT DIRECT LEGACY TARGETS]", len(articles))

    success = 0
    failed = 0
    skipped = 0

    for row in articles:
        valid, reason = is_valid_collected_article_for_draft(row)
        article_no = clean_text(row.get("article_no"))

        if not valid:
            skipped += 1
            log("[DRAFT DIRECT SKIP]", article_no, reason)
            continue

        try:
            process_article(conn, row)
            success += 1
        except Exception as e:
            failed += 1
            try:
                conn.rollback()
            except Exception:
                pass
            log("[DRAFT DIRECT ERROR]", article_no, str(e))

    return {"success": success, "failed": failed, "skipped": skipped, "total": len(articles)}


def run_subprocess_adapter(label, cmd, timeout=None):
    """
    완전 함수화 전 임시 adapter.
    STEP105 원칙상 최소화하며, fetcher/publisher를 순차적으로 함수화하며 제거한다.
    """
    log("=" * 80)
    log(f"[ADAPTER RUN START] {label}")
    log("[CMD]", " ".join([str(x) for x in cmd]))
    log("=" * 80)

    started = time.time()

    proc = subprocess.run(
        [str(x) for x in cmd],
        cwd=str(BASE_DIR),
        text=True,
        encoding="utf-8",
        errors="replace",
        stdout=subprocess.PIPE,
        stderr=subprocess.STDOUT,
        timeout=timeout,
    )

    elapsed = time.time() - started

    print(proc.stdout or "", flush=True)
    log("=" * 80)
    log(f"[ADAPTER RUN END] {label} returncode={proc.returncode} elapsed={elapsed:.1f}s")
    log("=" * 80)

    if proc.returncode != 0:
        raise RuntimeError(f"{label} failed with returncode={proc.returncode}")

    return proc.stdout or ""


def publish_server_step104_adapter(realtor_id, limit=1):
    """
    발행은 아직 publish_one_queue 함수화 전이므로 STEP104 worker adapter 사용.
    """
    cmd = [
        PYTHON_EXE,
        "workers/publish_worker_step104.py",
        "--realtor-id",
        str(int(realtor_id)),
        "--limit",
        str(int(limit or 1)),
        "--no-db-lock",
    ]
    return run_subprocess_adapter("publish_server_step104_adapter", cmd, timeout=3600)


def run_one_realtor_pipeline(
    realtor_id,
    mode="auto",
    publish_mode_to_set="",
    require_test_mode=False,
    allow_test_override=True,
    shift_today_published=True,
    skip_collect=False,
    skip_fetch=False,
    skip_draft=False,
    skip_publish=False,
    draft_buffer=1,
    replenish_count=1,
    fetch_limit=3,
    draft_limit=1,
    publish_limit=1,
):
    conn = get_conn()
    auto_overridden = False

    try:
        if publish_mode_to_set:
            set_publish_mode(conn, realtor_id, publish_mode_to_set)

        policy = get_policy(conn, realtor_id)
        realtor = policy["realtor"]
        publish_mode = policy["publish_mode"]
        is_test_account = policy["is_test_account"]
        auto_enabled = policy["auto_publish_enabled"]

        log("=" * 80)
        log("[STEP105 ONE REALTOR PIPELINE START]")
        log("TIME:", datetime.now().strftime("%Y-%m-%d %H:%M:%S"))
        log("realtor_id:", realtor_id)
        log("office_name:", realtor.get("office_name"))
        log("publish_mode:", publish_mode)
        log("is_test_account:", is_test_account)
        log("auto_publish_enabled:", auto_enabled)
        log("=" * 80)

        if require_test_mode and not (publish_mode == "test" or is_test_account == 1):
            raise RuntimeError("require_test_mode enabled, but realtor is not test")

        session = check_session_db(conn, realtor_id)
        log("[SESSION DB CHECK]", json.dumps(session, ensure_ascii=False, default=str))
        if session.get("exists") and session.get("ok") is False:
            log("[SESSION DB WARN]", "publish adapter will verify actual login:", session.get("message"))

        if shift_today_published and (publish_mode == "test" or is_test_account == 1):
            affected = shift_today_published_for_test(conn, realtor_id)
            log("[TEST SHIFT TODAY PUBLISHED]", f"affected={affected}")

        print_counts(conn, realtor_id, "before")

        # collect
        if not skip_collect:
            if auto_enabled == 0 and allow_test_override:
                set_auto_publish_enabled(conn, realtor_id, True)
                auto_overridden = True
                log("[AUTO PUBLISH TEMP ON] collect_candidate only")

                # reload realtor with setting not necessary, schedule fetch_target uses setting only before direct process
                # direct process_realtor itself doesn't check setting.

            collect_candidate_for_realtor(
                conn=conn,
                realtor=realtor,
                draft_buffer=draft_buffer,
                replenish_count=replenish_count,
            )

            if auto_overridden:
                set_auto_publish_enabled(conn, realtor_id, False)
                auto_overridden = False
                log("[AUTO PUBLISH RESTORED] OFF")

            print_counts(conn, realtor_id, "after collect")

        # fetch details
        if not skip_fetch:
            # commit pending candidates before fetcher adapter reads them
            try:
                conn.commit()
            except Exception:
                pass

            fetch_details_adapter(limit=fetch_limit)

            conn.close()
            conn = get_conn()
            print_counts(conn, realtor_id, "after fetch")

        # draft
        if not skip_draft:
            result = generate_draft_direct(
                conn=conn,
                realtor_id=realtor_id,
                limit=draft_limit,
                draft_buffer=draft_buffer,
                replenish_count=replenish_count,
            )
            log("[DRAFT DIRECT RESULT]", json.dumps(result, ensure_ascii=False, default=str))

            print_counts(conn, realtor_id, "after draft")

        effective_mode = mode
        if effective_mode == "auto":
            if publish_mode == "client_assist":
                effective_mode = "client"
            elif publish_mode == "manual":
                effective_mode = "draft-only"
            else:
                effective_mode = "server"

        log("[EFFECTIVE MODE]", effective_mode)

        if skip_publish:
            log("[SKIP PUBLISH] requested")
        elif effective_mode == "client":
            log("[CLIENT READY]", f"pending_publish={count_pending_publish(conn, realtor_id)}")
        elif effective_mode == "draft-only":
            log("[DRAFT ONLY] publish skipped")
        elif effective_mode == "server":
            publish_server_step104_adapter(realtor_id, limit=publish_limit)
            conn.close()
            conn = get_conn()
            print_counts(conn, realtor_id, "after publish")
            latest = fetch_latest_publish_queue(conn, realtor_id)
            log("[LATEST QUEUE]", json.dumps(latest or {}, ensure_ascii=False, default=str))
        else:
            raise RuntimeError(f"unknown effective mode: {effective_mode}")

        log("=" * 80)
        log("[STEP105 ONE REALTOR PIPELINE DONE]")
        log("=" * 80)

    finally:
        try:
            if auto_overridden:
                set_auto_publish_enabled(conn, realtor_id, False)
                log("[AUTO PUBLISH RESTORED IN FINALLY] OFF")
        except Exception:
            pass

        try:
            conn.close()
        except Exception:
            pass
