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

"""
daily_pipeline_runner.py

야간 수집·상세·초안·비공개 대기 통합 실행기.

운영 목표:
- 18:05 작업 스케줄러에서 이 파일 하나만 실행
- 세션 점검/자동재연동부터 전체 후보 탐색, 버퍼 보충, 상세수집, 초안 생성까지 순서대로 실행
- 현재 매물 확인이 끝난 뒤 종료 매물을 비교하여 비공개 대기 Queue를 생성
- detect_new_articles.py가 현재 매물 확인/last_seen_at 갱신까지 담당하므로 refresh_realtor_articles_seen.py는 기본 제외
- STEP94: detect_new_articles.py --realtor-id 직접 지원에 맞춰 daily_pipeline_runner도 특정 중개사 실행을 정확히 전달
- 다음 날 09:55 종료선 도달 시 다음 단계 실행을 중지하고 pending 상태를 유지
- 미완료 작업은 DB 상태를 유지하고 다음날 이어서 처리
- STEP98-A: password_error/naver_block/disconnected 계정은 낮 파이프라인 버퍼 보충에서 제외
- STEP351: 비정상 종료 후 남은 daily_pipeline_runner stale DB Lock 자동 복구
- STEP360: 운영 대상 중개사 수의 고정 제한 제거
- STEP361: 신규/기존 구분 없이 중개사별 최신 매물 1건만 상세수집·초안 생성

기본 실행:
python jobs/daily_pipeline_runner.py

테스트:
python jobs/daily_pipeline_runner.py --dry-run
python jobs/daily_pipeline_runner.py --realtor-id 1
python jobs/daily_pipeline_runner.py --end-hour 9 --end-minute 55
"""

import os
import sys
import json
import time
import socket
import argparse
import subprocess
from datetime import datetime, timedelta

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

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

from db import get_conn


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


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


def worker_server_name():
    try:
        return socket.gethostname()
    except Exception:
        return "unknown-worker"


def create_pipeline_run(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
                (
                    'daily_pipeline_runner',
                    'running',
                    NOW(),
                    'daily pipeline runner started',
                    %s,
                    NOW()
                )
            """, (
                safe_json_dumps(meta or {}),
            ))

            run_id = cur.lastrowid

        conn.commit()

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

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

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


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

    conn = None

    try:
        conn = get_conn()

        with conn.cursor() as cur:
            cur.execute("""
                UPDATE blog_pipeline_runs
                SET
                    status = %s,
                    finished_at = NOW(),
                    total_count = %s,
                    success_count = %s,
                    fail_count = %s,
                    skipped_count = %s,
                    message = %s,
                    meta_json = %s
                WHERE id = %s
            """, (
                str(status or "")[: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()

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

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

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


def add_pipeline_log(run_id=None, level="info", step_name="", message="", 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,
                    message,
                    context_json,
                    created_at
                )
                VALUES
                (
                    %s,
                    %s,
                    %s,
                    %s,
                    %s,
                    NOW()
                )
            """, (
                run_id,
                str(level or "info")[:20],
                str(step_name or "")[:100],
                str(message or "")[:2000],
                safe_json_dumps(context or {}),
            ))

        conn.commit()

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

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


def worker_token(prefix):
    return f"{prefix}:{worker_server_name()}:{os.getpid()}:{datetime.now().strftime('%Y%m%d%H%M%S')}"


def parse_lock_owner_token(owner_token):
    """
    owner_token 형식:
        <lock_name>:<hostname>:<pid>:<timestamp>

    과거 데이터나 예상하지 못한 형식은 안전하게 해석 실패로 처리한다.
    """
    parts = str(owner_token or "").split(":")
    if len(parts) < 4:
        return None

    try:
        pid = int(parts[-2])
    except Exception:
        return None

    return {
        "lock_name": ":".join(parts[:-3]),
        "hostname": parts[-3],
        "pid": pid,
        "timestamp": parts[-1],
    }


def get_windows_process_info(pid):
    """
    Windows PID의 프로세스명/명령행을 조회한다.

    반환값:
    - found=True  : 해당 PID가 존재함
    - found=False : 해당 PID가 존재하지 않음
    - ok=False    : 조회 자체가 실패함. 이 경우 유효 Lock으로 보수적으로 처리한다.
    """
    try:
        pid = int(pid)
    except Exception:
        return {"ok": True, "found": False, "name": "", "command_line": ""}

    if pid <= 0:
        return {"ok": True, "found": False, "name": "", "command_line": ""}

    ps_script = (
        f"$p = Get-CimInstance Win32_Process -Filter \"ProcessId={pid}\"; "
        "if ($null -eq $p) { exit 3 }; "
        "$p | Select-Object Name,CommandLine | ConvertTo-Json -Compress"
    )

    try:
        result = subprocess.run(
            ["powershell", "-NoProfile", "-Command", ps_script],
            capture_output=True,
            text=True,
            encoding="utf-8",
            errors="ignore",
            timeout=10,
            check=False,
        )
    except Exception as e:
        return {
            "ok": False,
            "found": False,
            "name": "",
            "command_line": "",
            "error": str(e),
        }

    if result.returncode == 3:
        return {"ok": True, "found": False, "name": "", "command_line": ""}

    if result.returncode != 0:
        return {
            "ok": False,
            "found": False,
            "name": "",
            "command_line": "",
            "error": (result.stderr or result.stdout or f"returncode={result.returncode}").strip()[:500],
        }

    raw = (result.stdout or "").strip()
    if not raw:
        return {"ok": True, "found": False, "name": "", "command_line": ""}

    try:
        data = json.loads(raw)
    except Exception as e:
        return {
            "ok": False,
            "found": False,
            "name": "",
            "command_line": "",
            "error": f"process json parse failed: {e}",
        }

    return {
        "ok": True,
        "found": True,
        "name": str((data or {}).get("Name") or ""),
        "command_line": str((data or {}).get("CommandLine") or ""),
    }


def lock_owner_is_active(owner_token, expected_script="daily_pipeline_runner.py"):
    """
    현재 서버의 기존 Lock 소유 프로세스가 실제 daily_pipeline_runner인지 확인한다.

    다른 서버의 Lock 또는 프로세스 조회 실패는 오삭제 방지를 위해 active=True로 본다.
    """
    parsed = parse_lock_owner_token(owner_token)
    if not parsed:
        return True, "owner_token 형식 확인 불가"

    owner_host = str(parsed.get("hostname") or "").lower()
    local_host = str(worker_server_name() or "").lower()

    if owner_host and local_host and owner_host != local_host:
        return True, f"다른 서버 Lock owner_host={owner_host}"

    process_info = get_windows_process_info(parsed.get("pid"))
    if not process_info.get("ok"):
        return True, f"프로세스 조회 실패: {process_info.get('error') or '-'}"

    if not process_info.get("found"):
        return False, f"PID {parsed.get('pid')} 프로세스 없음"

    command_line = str(process_info.get("command_line") or "").lower().replace("/", "\\")
    expected = str(expected_script or "").lower().replace("/", "\\")

    if expected and expected not in command_line:
        return False, (
            f"PID {parsed.get('pid')} 재사용/다른 프로세스 "
            f"name={process_info.get('name') or '-'}"
        )

    return True, f"PID {parsed.get('pid')} daily_pipeline_runner 실행 중"


def acquire_db_lock(lock_name="daily_pipeline_runner", ttl_minutes=900):
    owner_token = worker_token(lock_name)
    conn = None
    stale_owner = ""
    stale_reason = ""

    try:
        conn = get_conn()

        with conn.cursor() as cur:
            # 동일 Lock 한 건만 트랜잭션 잠금하여 동시 실행 경쟁을 막는다.
            cur.execute("""
                SELECT
                    lock_name,
                    locked_at,
                    expires_at,
                    owner_token,
                    CASE WHEN expires_at < NOW() THEN 1 ELSE 0 END AS is_expired
                FROM blog_app_locks
                WHERE lock_name = %s
                FOR UPDATE
            """, (str(lock_name)[:100],))

            existing = cur.fetchone()

            if existing:
                if isinstance(existing, dict):
                    existing_owner = str(existing.get("owner_token") or "")
                    is_expired = int(existing.get("is_expired") or 0) == 1
                else:
                    existing_owner = str(existing[3] or "")
                    is_expired = int(existing[4] or 0) == 1

                if is_expired:
                    stale_owner = existing_owner
                    stale_reason = "DB expires_at 만료"
                else:
                    active, reason = lock_owner_is_active(existing_owner)
                    if active:
                        conn.rollback()
                        log(
                            f"[DB LOCK ACTIVE] {lock_name} owner={existing_owner} / {reason}"
                        )
                        return False, owner_token

                    stale_owner = existing_owner
                    stale_reason = reason

                cur.execute("""
                    DELETE FROM blog_app_locks
                    WHERE lock_name = %s
                      AND owner_token = %s
                """, (
                    str(lock_name)[:100],
                    str(existing_owner)[:100],
                ))

                if cur.rowcount != 1:
                    raise RuntimeError(
                        f"stale lock delete failed: lock_name={lock_name}, affected={cur.rowcount}"
                    )

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

        if stale_owner:
            log(
                f"[STALE DB LOCK REMOVED] {lock_name} "
                f"owner={stale_owner} / reason={stale_reason}"
            )

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

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

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

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


def release_db_lock(lock_name, owner_token):
    conn = None

    try:
        conn = get_conn()

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

            affected = cur.rowcount

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

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

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


def get_end_datetime(end_hour=9, end_minute=55):
    now = datetime.now()
    end_dt = now.replace(
        hour=int(end_hour),
        minute=int(end_minute),
        second=0,
        microsecond=0,
    )

    # 23시 이후 시작해 다음 날 오전에 끝나는 야간 파이프라인만 날짜를 넘긴다.
    # 오전 종료 시각을 지난 뒤에는 다시 하루를 더하지 않아 11:55 종료선이 유지된다.
    if now.hour >= 12 and int(end_hour) < 12:
        end_dt += timedelta(days=1)

    return end_dt


def is_past_end_time(end_hour=9, end_minute=55):
    return datetime.now() >= get_end_datetime(end_hour, end_minute)


def seconds_until_end(end_hour=9, end_minute=55):
    return max(0, int((get_end_datetime(end_hour, end_minute) - datetime.now()).total_seconds()))


def command_exists(script_rel_path):
    return os.path.exists(os.path.join(BASE_DIR, script_rel_path))


def run_command(step_name, command, run_id=None, dry_run=False, timeout_seconds=None):
    log("-" * 80)
    log(f"[STEP START] {step_name}")
    log("[COMMAND] " + " ".join(command))

    add_pipeline_log(
        run_id=run_id,
        level="info",
        step_name=step_name,
        message=f"step start: {step_name}",
        context={
            "command": command,
            "dry_run": dry_run,
            "timeout_seconds": timeout_seconds,
        }
    )

    if dry_run:
        log(f"[DRY RUN SKIP] {step_name}")
        add_pipeline_log(
            run_id=run_id,
            level="info",
            step_name=step_name,
            message=f"dry-run skipped: {step_name}",
            context={"command": command}
        )
        return True, 0

    started = time.time()

    try:
        proc = subprocess.run(
            command,
            cwd=BASE_DIR,
            text=True,
            encoding="utf-8",
            errors="replace",
            timeout=timeout_seconds,
        )

        elapsed = max(0, int(time.time() - started))

        ok = proc.returncode == 0

        if ok:
            log(f"[STEP DONE] {step_name} elapsed={elapsed}s")
        else:
            log(f"[STEP FAILED] {step_name} code={proc.returncode}, elapsed={elapsed}s")

        add_pipeline_log(
            run_id=run_id,
            level="info" if ok else "error",
            step_name=step_name,
            message=f"step {'done' if ok else 'failed'}: {step_name}, code={proc.returncode}, elapsed={elapsed}s",
            context={
                "command": command,
                "returncode": proc.returncode,
                "elapsed": elapsed,
            }
        )

        return ok, proc.returncode

    except subprocess.TimeoutExpired:
        elapsed = max(0, int(time.time() - started))
        log(f"[STEP TIMEOUT] {step_name} timeout={timeout_seconds}s elapsed={elapsed}s")

        add_pipeline_log(
            run_id=run_id,
            level="error",
            step_name=step_name,
            message=f"step timeout: {step_name}",
            context={
                "command": command,
                "timeout_seconds": timeout_seconds,
                "elapsed": elapsed,
            }
        )

        return False, 124

    except Exception as e:
        elapsed = max(0, int(time.time() - started))
        log(f"[STEP ERROR] {step_name} {e}")

        add_pipeline_log(
            run_id=run_id,
            level="error",
            step_name=step_name,
            message=f"step exception: {e}",
            context={
                "command": command,
                "elapsed": elapsed,
                "error": str(e),
            }
        )

        return False, 1



# ---------------------------------------------------------
# STEP91-A/B 운영형 버퍼 보충 내부 단계
# - detect_new_articles.py가 blog_realtor_articles에 전체 후보를 넣는다.
# - 여기서는 blog_realtor_articles 기준으로 버퍼 부족분만 blog_article_work_queue에 넣는다.
# - 저장/상세수집/초안생성은 전체 매물이 아니라 부족분만 처리한다.
# ---------------------------------------------------------

def fetch_enabled_realtors_for_buffer(realtor_id=None):
    conn = None

    try:
        conn = get_conn()

        sql = """
            SELECT DISTINCT
                r.*,
                COALESCE(s.auto_publish_enabled, 0) AS auto_publish_enabled,
                COALESCE(s.daily_post_limit, 1) AS daily_post_limit,
                COALESCE(a.login_status, '') AS login_status
            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 "

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

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


def get_realtor_buffer_target(realtor, default_target=45):
    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):
    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):
    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 count_collected_ready_articles(conn, realtor_id):
    """
    상세수집은 끝났지만 아직 초안/발행큐가 없는 매물 수.
    이 매물은 곧 초안 재고로 전환될 수 있으므로 버퍼 계산에 포함한다.
    """
    with conn.cursor() as cur:
        cur.execute("""
            SELECT COUNT(DISTINCT a.article_no) AS cnt
            FROM blog_realtor_articles a
            LEFT JOIN blog_article_drafts d
                ON CONVERT(d.article_no USING utf8mb4) COLLATE utf8mb4_unicode_ci
                   =
                   CONVERT(a.article_no USING utf8mb4) COLLATE utf8mb4_unicode_ci
            LEFT JOIN blog_publish_queue q
                ON CONVERT(q.article_no USING utf8mb4) COLLATE utf8mb4_unicode_ci
                   =
                   CONVERT(a.article_no USING utf8mb4) COLLATE utf8mb4_unicode_ci
            WHERE a.realtor_id = %s
              AND COALESCE(a.collect_status, '') = 'collected'
              AND COALESCE(a.detail_collected, 0) = 1
              AND d.id IS NULL
              AND q.id IS NULL
        """, (int(realtor_id),))
        row = cur.fetchone()

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


def work_queue_exists(conn, realtor_id, article_no):
    with conn.cursor() as cur:
        cur.execute("""
            SELECT id
            FROM blog_article_work_queue
            WHERE realtor_id = %s
              AND article_no = %s
              AND work_type = 'new_article'
              AND work_status IN ('pending', 'processing', 'done')
            LIMIT 1
        """, (int(realtor_id), 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 id
            FROM blog_publish_queue
            WHERE article_no = %s
              AND queue_status IN (
                    'waiting',
                    'pending',
                    'ready',
                    'processing',
                    'published',
                    'done',
                    'hold'
              )
            LIMIT 1
        """, (str(article_no),))
        row = cur.fetchone()

    return bool(row)


def load_realtor_article_candidates_for_buffer(conn, realtor_id):
    """
    detect_new_articles가 확보한 blog_realtor_articles 후보 중
    상세수집 큐에 넣을 수 있는 매물을 최신순으로 가져온다.
    """
    with conn.cursor() as cur:
        cur.execute("""
            SELECT
                id,
                realtor_id,
                article_no,
                collect_status,
                detail_collected,
                article_name,
                price_text,
                first_posted_at,
                last_seen_at,
                created_at,
                updated_at
            FROM blog_realtor_articles
            WHERE realtor_id = %s
              AND COALESCE(article_status, 'active') = 'active'
              AND COALESCE(visibility_status, 'public') = 'public'
              AND (
                    COALESCE(detail_collected, 0) = 0
                    OR COALESCE(collect_status, '') IN ('pending', 'failed', '')
                    OR COALESCE(article_name, '') IN (
                        '',
                        '네이버 부동산 후보 매물',
                        '네이버 부동산 현재 매물'
                    )
              )
            ORDER BY
                COALESCE(first_posted_at, last_seen_at, created_at, updated_at) DESC,
                id DESC
        """, (int(realtor_id),))
        return cur.fetchall()


def enqueue_article_work_from_cache(conn, realtor_id, article_no, priority=100):
    with conn.cursor() as cur:
        cur.execute("""
            INSERT IGNORE 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),
        ))

        return cur.rowcount


def run_enqueue_buffer_from_articles(args, run_id=None, dry_run=False):
    """
    STEP91-A 내부 실행 단계.

    detect_new_articles.py가 전체 후보를 넓게 확보한 뒤,
    여기서는 버퍼 부족분만 상세수집 work_queue로 보낸다.
    """
    conn = None
    total_realtors = 0
    total_queued = 0
    total_checked = 0

    try:
        realtors = fetch_enabled_realtors_for_buffer(
            realtor_id=args.realtor_id,
        )

        total_realtors = len(realtors)

        log(f"[BUFFER ENQUEUE REALTORS] {total_realtors}")

        conn = get_conn()

        for realtor in realtors:
            realtor_id = int(realtor["id"])
            office_name = realtor.get("office_name", "")
            target = get_realtor_buffer_target(
                realtor,
                default_target=args.draft_buffer,
            )

            active_drafts = count_active_draft_buffer(conn, realtor_id)
            pending_work = count_pending_work_queue(conn, realtor_id)
            collected_ready = count_collected_ready_articles(conn, realtor_id)

            current_buffer = active_drafts + pending_work + collected_ready
            need = max(0, int(target) - int(current_buffer))

            # STEP361:
            # 신규/기존 중개사를 구분하지 않고 이번 실행에서는 최신 매물
            # 1건만 Queue에 넣는다. 전체 중개사 수에는 상한을 두지 않는다.
            queue_target = min(need, 1)
            fill_mode = "latest_one"

            log("-" * 80)
            log(
                "[BUFFER STATUS] "
                f"realtor_id={realtor_id}, "
                f"office={office_name}, "
                f"active_drafts={active_drafts}, "
                f"pending_work={pending_work}, "
                f"collected_ready={collected_ready}, "
                f"current={current_buffer}, "
                f"target={target}, "
                f"need={need}, "
                f"queue_target={queue_target}, "
                f"fill_mode={fill_mode}"
            )

            if queue_target <= 0:
                log("[BUFFER SKIP] enough")
                continue

            candidates = load_realtor_article_candidates_for_buffer(
                conn,
                realtor_id=realtor_id,
            )

            log(f"[BUFFER CANDIDATES] realtor_id={realtor_id}, rows={len(candidates)}")

            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
                total_checked += 1

                if work_queue_exists(conn, realtor_id, article_no):
                    log(f"[BUFFER SKIP WORK EXISTS] {article_no}")
                    continue

                if draft_exists(conn, article_no):
                    log(f"[BUFFER SKIP DRAFT EXISTS] {article_no}")
                    continue

                if publish_exists(conn, article_no):
                    log(f"[BUFFER SKIP PUBLISH EXISTS] {article_no}")
                    continue

                if dry_run:
                    queued += 1
                    total_queued += 1
                    log(f"[DRY BUFFER QUEUE] {article_no}")
                    continue

                affected = enqueue_article_work_from_cache(
                    conn,
                    realtor_id=realtor_id,
                    article_no=article_no,
                    priority=100,
                )

                conn.commit()

                if affected:
                    queued += 1
                    total_queued += 1
                    log(f"[BUFFER QUEUED] {article_no}")
                else:
                    log(f"[BUFFER QUEUE EXISTS/IGNORED] {article_no}")

            log(
                "[BUFFER SUMMARY] "
                f"realtor_id={realtor_id}, "
                f"checked={checked}, "
                f"queued={queued}, "
                f"queue_target={queue_target}, "
                f"need={need}"
            )

        add_pipeline_log(
            run_id=run_id,
            level="info",
            step_name="enqueue_buffer_from_articles",
            message=(
                f"buffer enqueue done: realtors={total_realtors}, "
                f"checked={total_checked}, queued={total_queued}"
            ),
            context={
                "total_realtors": total_realtors,
                "total_checked": total_checked,
                "total_queued": total_queued,
            }
        )

        return True, 0

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

        log(f"[BUFFER ENQUEUE ERROR] {e}")

        add_pipeline_log(
            run_id=run_id,
            level="error",
            step_name="enqueue_buffer_from_articles",
            message=str(e),
            context={"error": str(e)}
        )

        return False, 1

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


def count_detail_targets(realtor_id=None):
    """현재 상세수집이 필요한 작업 전체 수를 반환한다."""
    conn = None

    try:
        conn = get_conn()
        sql = """
            SELECT COUNT(DISTINCT a.id) AS cnt
            FROM blog_realtor_articles a
            WHERE (
                    COALESCE(a.detail_collected, 0) = 0
                    OR COALESCE(a.collect_status, '') = 'pending'
                    OR COALESCE(a.article_name, '') IN (
                        '',
                        '네이버 부동산 후보 매물',
                        '네이버 부동산 현재 매물'
                    )
                  )
              AND EXISTS (
                    SELECT 1
                    FROM blog_article_work_queue w
                    WHERE w.realtor_id = a.realtor_id
                      AND BINARY w.article_no = BINARY a.article_no
                      AND w.work_type = 'new_article'
                      AND w.work_status IN ('pending', 'processing')
                  )
        """
        params = []

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

        with conn.cursor() as cur:
            cur.execute(sql, params)
            row = cur.fetchone()

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

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


def count_enabled_realtors(realtor_id=None):
    """자동수집·자동발행이 켜진 운영 중개사 전체 수를 반환한다."""
    conn = None

    try:
        conn = get_conn()
        sql = """
            SELECT COUNT(DISTINCT r.id) AS cnt
            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))

        with conn.cursor() as cur:
            cur.execute(sql, params)
            row = cur.fetchone()

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

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


def build_steps(args):
    py = sys.executable

    common_realtor_args = []

    if args.realtor_id:
        common_realtor_args = ["--realtor-id", str(args.realtor_id)]

    # 하위 스크립트는 숫자 인자를 요구한다. 고정 상한 대신 실행 시점의
    # 실제 전체 대상 수를 전달해 신규 중개사가 번호와 무관하게 포함되게 한다.
    enabled_realtor_count = count_enabled_realtors(args.realtor_id)
    if args.limit_realtors:
        # 운영 스케줄에서는 값을 전달하지 않는다. 관리자 화면의
        # "2개 중개사 테스트"처럼 명시적인 수동 테스트에서만 제한한다.
        enabled_realtor_count = min(enabled_realtor_count, int(args.limit_realtors))
    limit_realtors = str(enabled_realtor_count)

    # STEP361 운영 정책:
    # 전체 대상 중개사를 빠짐없이 순회하되 중개사별 후보/상세/초안은
    # 최신 1건만 처리한다. 신규 중개사의 initial_fill 예외도 두지 않는다.
    per_realtor_target = "1"

    detect_args = [
        py,
        "jobs/detect_new_articles.py",
        "--limit-realtors",
        limit_realtors,
        "--only-enabled",
    ]

    # STEP94:
    # detect_new_articles.py가 --realtor-id를 지원하므로
    # 테스트/장애복구 시 특정 중개사만 정확히 전체 후보 탐색한다.
    if args.realtor_id:
        detect_args += ["--realtor-id", str(args.realtor_id)]

    steps = [
        {
            "name": "session_check",
            "script": "workers/naver_session_check_worker.py",
            "command": [
                py,
                "workers/naver_session_check_worker.py",
                "--limit",
                limit_realtors,
                "--auto-relogin",
            ] + common_realtor_args,
            "required": False,
        },
        {
            "name": "schedule_daily_articles_buffer",
            "script": "jobs/schedule_daily_articles.py",
            "command": [
                py,
                "jobs/schedule_daily_articles.py",
                "--limit-realtors",
                limit_realtors,
                "--headless",
                "1",
                "--draft-buffer",
                per_realtor_target,
                "--replenish-count",
                per_realtor_target,
                "--max-pages",
                str(args.fast_max_pages),
                "--disable-private-sync",
                "--lock-ttl-minutes",
                "300",
            ] + common_realtor_args,
            # 기본수집이 일부 시간 초과되어도 이미 확보된 Queue의
            # 상세수집·초안생성은 반드시 계속한다.
            "required": False,
            "continue_on_fail": True,
            "may_run_long": True,
            "timeout_cap_seconds": 14400,
        },
        {
            "name": "check_removed_articles",
            "script": "workers/check_removed_articles.py",
            "command": [
                py,
                "workers/check_removed_articles.py",
                "--limit",
                limit_realtors,
                "--confirm-count",
                "2",
                "--min-days-before-private",
                "30",
            ] + common_realtor_args,
            "required": False,
            "may_run_long": True,
            # 전체 중개사의 종료 매물 확인·비공개 대기 Queue 생성에는
            # 약 3시간이 필요하므로 최대 4시간까지 허용한다.
            # 실제 완료 시에는 기다리지 않고 즉시 다음 단계로 진행한다.
            "timeout_cap_seconds": 14400,
        },
        {
            "name": "naver_response_article_fetcher",
            "script": "services/naver_response_article_fetcher.py",
            "command": [
                py,
                "services/naver_response_article_fetcher.py",
                "--limit",
                "1",
            ],
            "dynamic_detail_limit": True,
            # C4-21:
            # 상세수집은 중요하지만, 이 단계가 timeout/실패해도
            # 이미 detail_collected=1 로 저장된 매물이 있으면 초안/Queue 생성은 계속 진행해야 한다.
            # 클라이언트 발행 구조에서는 "초안 생성 중단"이 더 큰 장애가 된다.
            "required": False,
            "continue_on_fail": True,
            "may_run_long": True,
            "timeout_cap_seconds": 5400,
        },
        {
            "name": "generate_blog_drafts",
            "script": "jobs/generate_blog_drafts.py",
            "command": [
                py,
                "jobs/generate_blog_drafts.py",
                "--limit",
                str(enabled_realtor_count),
                "--draft-buffer",
                per_realtor_target,
                "--replenish-count",
                per_realtor_target,
                "--limit-realtors",
                str(enabled_realtor_count),
            ] + common_realtor_args,
            "required": True,
            "may_run_long": True,
        },
    ]

    if args.use_refresh_seen:
        steps.insert(3, {
            "name": "refresh_realtor_articles_seen_optional",
            "script": "workers/refresh_realtor_articles_seen.py",
            "command": [
                py,
                "workers/refresh_realtor_articles_seen.py",
                "--limit",
                limit_realtors,
                "--headless",
                "1",
            ] + common_realtor_args,
            "required": False,
        })

    if args.use_legacy_scheduler:
        steps.insert(3, {
            "name": "schedule_daily_articles_legacy",
            "script": "jobs/schedule_daily_articles.py",
            "command": [
                py,
                "jobs/schedule_daily_articles.py",
                "--limit-realtors",
                limit_realtors,
                "--headless",
                "1",
                "--draft-buffer",
                per_realtor_target,
                "--replenish-count",
                per_realtor_target,
            ] + common_realtor_args,
            "required": False,
        })

    return steps

def run(args):
    lock_name = "daily_pipeline_runner"
    lock_owner = None
    run_id = None

    total = 0
    success = 0
    failed = 0
    noncritical_failed = 0
    skipped = 0

    final_status = "success"
    final_message = "daily pipeline completed"

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

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

        run_id = create_pipeline_run(
            meta={
                "end_hour": args.end_hour,
                "end_minute": args.end_minute,
                "processing_limit_policy": "all_current_targets",
                "draft_buffer": args.draft_buffer,
                "replenish_count": args.replenish_count,
                "fast_max_pages": args.fast_max_pages,
                "realtor_id": args.realtor_id,
                "dry_run": args.dry_run,
                "use_legacy_scheduler": args.use_legacy_scheduler,
                "use_refresh_seen": args.use_refresh_seen,
                "use_refresh_seen": args.use_refresh_seen,
                "worker_server": worker_server_name(),
                "noncritical_failed": noncritical_failed,
            }
        )

        log("=" * 80)
        log("[DAILY PIPELINE START]")
        log(f"BASE_DIR={BASE_DIR}")
        log(f"END_TIME={args.end_hour:02d}:{args.end_minute:02d}")
        log(f"SECONDS_UNTIL_END={seconds_until_end(args.end_hour, args.end_minute)}")
        log("=" * 80)

        steps = build_steps(args)

        for step in steps:
            total += 1

            step_name = step["name"]
            script = step["script"]
            command = step["command"]

            if is_past_end_time(args.end_hour, args.end_minute):
                skipped += 1
                final_status = "stopped"
                final_message = f"end time reached before step: {step_name}"
                log(f"[END TIME STOP] before step={step_name}")
                add_pipeline_log(
                    run_id=run_id,
                    level="warning",
                    step_name=step_name,
                    message=final_message,
                    context={"end_hour": args.end_hour, "end_minute": args.end_minute}
                )
                break

            if step.get("type") == "internal_enqueue_buffer":
                ok, code = run_enqueue_buffer_from_articles(
                    args=args,
                    run_id=run_id,
                    dry_run=args.dry_run,
                )

            else:
                if not command_exists(script):
                    skipped += 1
                    msg = f"script not found: {script}"
                    log("[STEP SKIP] " + msg)

                    add_pipeline_log(
                        run_id=run_id,
                        level="warning",
                        step_name=step_name,
                        message=msg,
                        context={"script": script}
                    )

                    if step.get("required"):
                        failed += 1
                        final_status = "failed"
                        final_message = msg
                        break

                    continue

                if step.get("dynamic_detail_limit"):
                    # schedule_daily_articles가 이번 실행에서 Queue를 만든 뒤의
                    # 최신 pending 전체 수를 사용해야 신규 중개사가 빠지지 않는다.
                    detail_target_count = count_detail_targets(args.realtor_id)
                    if args.fetch_limit:
                        detail_target_count = min(detail_target_count, int(args.fetch_limit))
                    command = list(command)
                    command[command.index("--limit") + 1] = str(detail_target_count)
                    log(f"[DETAIL ALL TARGETS] count={detail_target_count}")

                timeout_seconds = None

                if step.get("may_run_long"):
                    # 운영 종료선까지 남은 시간만큼 허용한다.
                    remain = seconds_until_end(args.end_hour, args.end_minute)
                    timeout_seconds = remain if remain > 60 else 60
                    timeout_cap_seconds = step.get("timeout_cap_seconds")
                    if timeout_cap_seconds:
                        timeout_seconds = min(timeout_seconds, int(timeout_cap_seconds))

                ok, code = run_command(
                    step_name=step_name,
                    command=command,
                    run_id=run_id,
                    dry_run=args.dry_run,
                    timeout_seconds=timeout_seconds,
                )

            if ok:
                success += 1
            else:
                if step.get("required"):
                    failed += 1
                    final_status = "failed"
                    final_message = f"required step failed: {step_name}, code={code}"
                    break

                # C4-21:
                # 비필수 단계는 실패/timeout 되어도 다음 단계로 진행한다.
                # 특히 naver_response_article_fetcher 실패 때문에 generate_blog_drafts가 막히면
                # 클라이언트 발행할 pending Queue가 생성되지 않는다.
                noncritical_failed += 1
                log(f"[STEP WARNING CONTINUE] {step_name} code={code}")
                add_pipeline_log(
                    run_id=run_id,
                    level="warning",
                    step_name=step_name,
                    message=f"non-critical step failed but pipeline continues: {step_name}, code={code}",
                    context={"step_name": step_name, "code": code}
                )
                continue

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

        if final_status == "success" and noncritical_failed > 0:
            final_message = f"daily pipeline completed with warnings: noncritical_failed={noncritical_failed}"

        log("=" * 80)
        log("[DAILY PIPELINE DONE]")
        log(f"status={final_status}")
        log(f"total={total}, success={success}, failed={failed}, noncritical_failed={noncritical_failed}, skipped={skipped}")
        log("=" * 80)

    except Exception as e:
        final_status = "failed"
        final_message = str(e)[:2000]
        failed += 1
        log(f"[DAILY PIPELINE ERROR] {e}")

        add_pipeline_log(
            run_id=run_id,
            level="error",
            step_name="daily_pipeline_exception",
            message=final_message,
            context={"error": str(e)}
        )

    finally:
        finish_pipeline_run(
            run_id=run_id,
            status=final_status,
            total=total,
            success=success,
            failed=failed,
            skipped=skipped,
            message=final_message,
            meta={
                "end_hour": args.end_hour,
                "end_minute": args.end_minute,
                "processing_limit_policy": "all_current_targets",
                "draft_buffer": args.draft_buffer,
                "replenish_count": args.replenish_count,
                "fast_max_pages": args.fast_max_pages,
                "realtor_id": args.realtor_id,
                "dry_run": args.dry_run,
                "use_legacy_scheduler": args.use_legacy_scheduler,
                "worker_server": worker_server_name(),
            }
        )

        if lock_owner:
            release_db_lock(lock_name, lock_owner)


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

    parser.add_argument("--end-hour", type=int, default=9)
    parser.add_argument("--end-minute", type=int, default=55)
    parser.add_argument("--limit-realtors", type=int, default=None, help="수동 테스트 전용. 운영 스케줄에서는 사용하지 않음")
    parser.add_argument("--fetch-limit", type=int, default=None, help="수동 테스트 전용. 운영 스케줄에서는 사용하지 않음")
    parser.add_argument("--draft-limit", type=int, default=None, help="이전 테스트 명령 호환용")
    parser.add_argument("--draft-buffer", type=int, default=1, help="호환 인자. 운영 정책은 중개사별 최신 1건")
    parser.add_argument("--replenish-count", type=int, default=1, help="호환 인자. 운영 정책은 중개사별 최신 1건")
    parser.add_argument("--fast-max-pages", type=int, default=1)
    parser.add_argument("--realtor-id", type=int, default=None)
    parser.add_argument("--dry-run", action="store_true")
    parser.add_argument("--lock-ttl-minutes", type=int, default=900)
    parser.add_argument("--no-db-lock", action="store_true")
    parser.add_argument("--use-legacy-scheduler", action="store_true", help="기존 schedule_daily_articles.py도 함께 실행")
    parser.add_argument("--use-refresh-seen", action="store_true", help="기존 refresh_realtor_articles_seen.py를 추가 실행")

    args = parser.parse_args()

    run(args)


if __name__ == "__main__":
    main()
