# -*- coding: utf-8 -*-
"""
STEP109-01: V2 Workday Pipeline Orchestrator

목표:
- V1 daily_pipeline_runner.py의 안정적인 운영 구조를 V2에 맞게 계승한다.
- 오전 08:00~23:00 사이 중개사별로 독립 실행한다.
- 중개사 순서는 매 실행/매일 랜덤화한다.
- 테스트 기간에는 realtor_id 1,2를 사용한다.
- 검증 완료된 v2_publish_worker.py는 직접 호출하되, 그 전에 세션/스냅샷/수집/초안을 준비한다.
- 23:00 이후에는 비공개 처리 모드로 분리 실행한다.

중개사별 workday 흐름:
  1) session_check
  2) current snapshot / 전체매물 확인
  3) detail collect
  4) draft generation
  5) v2_publish_worker publish

중요 원칙:
- workers/v2_publish_worker.py는 건드리지 않는다.
- jobs/generate_blog_drafts_v2.py는 건드리지 않는다.
- services/naver_response_article_fetcher_v2.py는 건드리지 않는다.
- 이 파일은 전체 오케스트레이션만 담당한다.
"""

import argparse
import json
import os
import random
import socket
import subprocess
import sys
import time
from datetime import datetime
from pathlib import Path

ROOT = Path(__file__).resolve().parents[1]
if str(ROOT) not in sys.path:
    sys.path.insert(0, str(ROOT))

try:
    from db import get_conn
except Exception as exc:
    get_conn = None
    DB_IMPORT_ERROR = exc
else:
    DB_IMPORT_ERROR = None


DEFAULT_START_HOUR = 8.0
DEFAULT_UNTIL_HOUR = 23.0


def now_local():
    return datetime.now()


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


def safe_json(data):
    return json.dumps(data, ensure_ascii=False, default=str)


def hour_float(dt=None):
    dt = dt or now_local()
    return dt.hour + dt.minute / 60.0 + dt.second / 3600.0


def is_workday_window(start_hour, until_hour):
    h = hour_float()
    return float(start_hour) <= h < float(until_hour)


def is_private_window(until_hour):
    return hour_float() >= float(until_hour)


def parse_realtor_ids(value):
    value = str(value or "").strip()
    if not value:
        return []
    return [int(x.strip()) for x in value.replace(";", ",").split(",") if x.strip()]


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


def script_exists(rel_path):
    return (ROOT / rel_path).exists()


def get_table_columns(conn, table_name):
    try:
        with conn.cursor() as cur:
            cur.execute(f"SHOW COLUMNS FROM {table_name}")
            rows = cur.fetchall()
        return {str(r.get("Field") if isinstance(r, dict) else r[0]) for r in rows or [] if (r.get("Field") if isinstance(r, dict) else r[0])}
    except Exception:
        return set()


def acquire_db_lock(lock_name="v2_workday_pipeline", ttl_minutes=900):
    if get_conn is None:
        return False, f"db_import_failed:{DB_IMPORT_ERROR}"

    owner_token = f"{lock_name}:{worker_server_name()}:{os.getpid()}:{now_local().strftime('%Y%m%d%H%M%S')}"
    conn = None
    try:
        conn = get_conn()
        with conn.cursor() as cur:
            cur.execute("DELETE FROM blog_app_locks WHERE expires_at < NOW()")
            cur.execute("""
                INSERT INTO blog_app_locks
                (lock_name, locked_at, expires_at, owner_token)
                VALUES
                (%s, NOW(), DATE_ADD(NOW(), INTERVAL %s MINUTE), %s)
            """, (str(lock_name)[:100], int(ttl_minutes), str(owner_token)[:100]))
        conn.commit()
        log(f"[V2 DB LOCK ACQUIRED] {lock_name} owner={owner_token}")
        return True, owner_token
    except Exception as exc:
        try:
            if conn:
                conn.rollback()
        except Exception:
            pass
        log(f"[V2 DB LOCKED OR ERROR] {lock_name} / {exc}")
        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"[V2 DB LOCK RELEASED] {lock_name} affected={affected}")
    except Exception as exc:
        log(f"[V2 DB LOCK RELEASE ERROR] {exc}")
    finally:
        try:
            if conn:
                conn.close()
        except Exception:
            pass


def create_pipeline_run(run_type="v2_workday_pipeline", meta=None):
    if get_conn is None:
        return None

    conn = None
    try:
        conn = get_conn()
        cols = get_table_columns(conn, "blog_pipeline_runs")
        if not cols:
            return None

        values = {
            "run_type": run_type,
            "status": "running",
            "started_at": "__NOW__",
            "message": f"{run_type} started",
            "meta_json": safe_json(meta or {}),
            "created_at": "__NOW__",
        }

        insert_cols = []
        exprs = []
        params = []
        for col, val in values.items():
            if col not in cols:
                continue
            insert_cols.append(col)
            if val == "__NOW__":
                exprs.append("NOW()")
            else:
                exprs.append("%s")
                params.append(val)

        if not insert_cols:
            return None

        sql = f"INSERT INTO blog_pipeline_runs ({', '.join(insert_cols)}) VALUES ({', '.join(exprs)})"
        with conn.cursor() as cur:
            cur.execute(sql, params)
            run_id = cur.lastrowid
        conn.commit()
        log(f"[V2 PIPELINE RUN START] run_id={run_id}, run_type={run_type}")
        return run_id
    except Exception as exc:
        log(f"[V2 PIPELINE RUN START ERROR] {exc}")
        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 or get_conn is None:
        return

    conn = None
    try:
        conn = get_conn()
        cols = get_table_columns(conn, "blog_pipeline_runs")
        if not cols:
            return

        values = {
            "status": status,
            "finished_at": "__NOW__",
            "total_count": int(total or 0),
            "success_count": int(success or 0),
            "fail_count": int(failed or 0),
            "skipped_count": int(skipped or 0),
            "message": str(message or "")[:2000],
            "meta_json": safe_json(meta or {}),
        }

        sets = []
        params = []
        for col, val in values.items():
            if col not in cols:
                continue
            if val == "__NOW__":
                sets.append(f"{col}=NOW()")
            else:
                sets.append(f"{col}=%s")
                params.append(val)

        if not sets or "id" not in cols:
            return

        params.append(int(run_id))
        with conn.cursor() as cur:
            cur.execute(f"UPDATE blog_pipeline_runs SET {', '.join(sets)} WHERE id=%s", params)
        conn.commit()
        log(f"[V2 PIPELINE RUN FINISH] run_id={run_id}, status={status}")
    except Exception as exc:
        log(f"[V2 PIPELINE RUN FINISH ERROR] {exc}")
    finally:
        try:
            if conn:
                conn.close()
        except Exception:
            pass


def add_pipeline_log(run_id=None, level="info", step_name="", message="", context=None):
    if get_conn is None:
        return

    conn = None
    try:
        conn = get_conn()
        cols = get_table_columns(conn, "blog_pipeline_logs")
        if not cols:
            return

        values = {
            "run_id": run_id,
            "level": str(level or "info")[:20],
            "step_name": str(step_name or "")[:100],
            "message": str(message or "")[:2000],
            "context_json": safe_json(context or {}),
            "created_at": "__NOW__",
        }

        insert_cols = []
        exprs = []
        params = []
        for col, val in values.items():
            if col not in cols:
                continue
            insert_cols.append(col)
            if val == "__NOW__":
                exprs.append("NOW()")
            else:
                exprs.append("%s")
                params.append(val)

        if not insert_cols:
            return

        with conn.cursor() as cur:
            cur.execute(f"INSERT INTO blog_pipeline_logs ({', '.join(insert_cols)}) VALUES ({', '.join(exprs)})", params)
        conn.commit()
    except Exception as exc:
        log(f"[V2 PIPELINE LOG ERROR] {step_name} / {exc}")
    finally:
        try:
            if conn:
                conn.close()
        except Exception:
            pass


def fetch_realtors_from_db(realtor_ids=None, include_auto_off=False, ignore_paid_filter=False, limit=0):
    if get_conn is None:
        raise RuntimeError(f"db import failed: {DB_IMPORT_ERROR}")

    realtor_ids = realtor_ids or []
    conn = get_conn()
    try:
        with conn.cursor() as cur:
            cur.execute("SHOW TABLES LIKE 'blog_realtors'")
            if not cur.fetchone():
                raise RuntimeError("blog_realtors table not found")

            cols = get_table_columns(conn, "blog_realtors")
            office_col = "office_name" if "office_name" in cols else "name" if "name" in cols else None
            auto_col = "auto_publish_enabled" if "auto_publish_enabled" in cols else None
            active_col = "is_active" if "is_active" in cols else None
            status_col = "status" if "status" in cols else None
            paid_col = "is_paid" if "is_paid" in cols else "paid" if "paid" in cols else None
            test_col = "is_test_account" if "is_test_account" in cols else None

            select_cols = ["id AS realtor_id"]
            select_cols.append(f"{office_col} AS office_name" if office_col else "'' AS office_name")
            select_cols.append(f"{auto_col} AS auto_publish_enabled" if auto_col else "NULL AS auto_publish_enabled")
            select_cols.append(f"{test_col} AS is_test_account" if test_col else "NULL AS is_test_account")

            where = ["1=1"]
            params = []

            if realtor_ids:
                where.append("id IN (" + ",".join(["%s"] * len(realtor_ids)) + ")")
                params.extend([int(x) for x in realtor_ids])

            if active_col:
                where.append(f"COALESCE({active_col}, 1)=1")
            if status_col:
                where.append(f"COALESCE({status_col}, 'active')='active'")
            if paid_col and not ignore_paid_filter:
                where.append(f"COALESCE({paid_col}, 1)=1")
            if auto_col and not include_auto_off:
                where.append(f"COALESCE({auto_col}, 0)=1")

            sql = f"""
                SELECT {", ".join(select_cols)}
                FROM blog_realtors
                WHERE {" AND ".join(where)}
                ORDER BY id ASC
            """
            if limit and int(limit) > 0:
                sql += " LIMIT %s"
                params.append(int(limit))

            cur.execute(sql, params)
            rows = cur.fetchall()

        return [
            {
                "realtor_id": int(r.get("realtor_id")),
                "office_name": str(r.get("office_name") or ""),
                "auto_publish_enabled": r.get("auto_publish_enabled"),
                "is_test_account": r.get("is_test_account"),
            }
            for r in rows
        ]
    finally:
        conn.close()


def build_targets(args):
    realtor_ids = parse_realtor_ids(args.realtor_ids)
    if args.test_realtors:
        realtor_ids = [1, 2]

    targets = fetch_realtors_from_db(
        realtor_ids=realtor_ids,
        include_auto_off=args.include_auto_off or args.test_realtors,
        ignore_paid_filter=args.ignore_paid_filter,
        limit=args.limit,
    )

    seed = ""
    if args.shuffle_realtors:
        # 운영에서는 매일/매실행 달라지고, 테스트에서는 --shuffle-seed로 재현 가능.
        seed = args.shuffle_seed or f"{now_local().strftime('%Y%m%d')}_{socket.gethostname()}_{os.getpid()}_{now_local().strftime('%H%M%S')}"
        rnd = random.Random(str(seed))
        rnd.shuffle(targets)

    return targets, seed


def detect_internal_failure(stdout="", stderr=""):
    text = (str(stdout or "") + "\n" + str(stderr or ""))

    failure_markers = [
        '"queue_status": "failed"',
        "'queue_status': 'failed'",
        '"worker_status": "failed"',
        "'worker_status': 'failed'",
        "[V2 PUBLISH FAILURE RECORD]",
        "V2 pre-execute gate blocked final publish",
        "latest_success_snapshot_missing",
        "article_not_in_latest_success_snapshot",
        "[V2 PUBLISH ERROR CAUGHT]",
        "[HPE SESSION FAILED]",
        "Traceback (most recent call last)",
        "\"current_naver_count\": 0",
        "\"current_run_id\": null",
        "[STEP FAILED]",
        "[STEP ERROR]",
        "[STEP TIMEOUT]",
    ]

    hits = [m for m in failure_markers if m in text]

    published_success = (
        ('"queue_status": "published"' in text or "'queue_status': 'published'" in text)
        and ("https://blog.naver.com/" in text)
        and ("[V2 PUBLISH RESULT RECORD]" in text or "[V2 DB QUEUE UPDATED]" in text)
    )

    if published_success:
        # 발행 성공 후 post verify reference warning은 실패로 보지 않는다.
        hits = [
            h for h in hits
            if h not in [
                '"worker_status": "failed"',
                "'worker_status': 'failed'",
                "RuntimeError:",
            ]
        ]

    return {
        "failed": bool(hits),
        "markers": hits,
        "published_success_detected": published_success,
    }


def run_subprocess(cmd, dry_run=False, timeout_seconds=None):
    cmd_text = " ".join(f'"{x}"' if " " in str(x) else str(x) for x in cmd)
    if dry_run:
        return {
            "ok": True,
            "dry_run": True,
            "returncode": 0,
            "cmd": cmd_text,
            "stdout": "",
            "stderr": "",
            "elapsed_seconds": 0,
            "internal_failure": {"failed": False, "markers": []},
        }

    env = os.environ.copy()
    env["PYTHONIOENCODING"] = "utf-8"
    env["PYTHONUTF8"] = "1"

    started = time.time()
    try:
        proc = subprocess.run(
            cmd,
            cwd=str(ROOT),
            text=True,
            capture_output=True,
            encoding="utf-8",
            errors="replace",
            env=env,
            timeout=timeout_seconds,
        )
        internal_failure = detect_internal_failure(proc.stdout, proc.stderr)
        ok = proc.returncode == 0 and not internal_failure.get("failed")
        return {
            "ok": ok,
            "dry_run": False,
            "returncode": proc.returncode,
            "elapsed_seconds": round(time.time() - started, 2),
            "cmd": cmd_text,
            "stdout": proc.stdout,
            "stderr": proc.stderr,
            "internal_failure": internal_failure,
        }
    except subprocess.TimeoutExpired as exc:
        return {
            "ok": False,
            "dry_run": False,
            "returncode": 124,
            "elapsed_seconds": round(time.time() - started, 2),
            "cmd": cmd_text,
            "stdout": getattr(exc, "stdout", "") or "",
            "stderr": getattr(exc, "stderr", "") or "",
            "internal_failure": {"failed": True, "markers": ["subprocess_timeout"]},
            "error": f"TimeoutExpired:{timeout_seconds}",
        }
    except Exception as exc:
        return {
            "ok": False,
            "dry_run": False,
            "returncode": 1,
            "elapsed_seconds": round(time.time() - started, 2),
            "cmd": cmd_text,
            "stdout": "",
            "stderr": "",
            "internal_failure": {"failed": True, "markers": [f"{type(exc).__name__}:{exc}"]},
            "error": f"{type(exc).__name__}: {exc}",
        }


def print_tail(label, result, tail_lines=300):
    """
    하위 프로세스 stdout/stderr 출력 보조.
    - --tail-lines 300: 마지막 300줄 출력
    - --tail-lines 0 또는 음수: 전체 출력
    """
    log(label + " " + safe_json({
        "ok": result.get("ok"),
        "returncode": result.get("returncode"),
        "elapsed_seconds": result.get("elapsed_seconds"),
        "dry_run": result.get("dry_run"),
        "cmd": result.get("cmd"),
        "internal_failure": result.get("internal_failure"),
        "error": result.get("error"),
    }))

    try:
        tail_n = int(tail_lines)
    except Exception:
        tail_n = 300

    for stream_name in ["stdout", "stderr"]:
        stream_text = str(result.get(stream_name) or "")
        if not stream_text:
            continue

        lines = stream_text.splitlines()
        total_lines = len(lines)

        if tail_n <= 0 or tail_n >= total_lines:
            selected_lines = lines
            log(f"{label} {stream_name.upper()} FULL lines={total_lines}")
        else:
            selected_lines = lines[-tail_n:]
            omitted = max(0, total_lines - tail_n)
            log(f"{label} {stream_name.upper()} TAIL last={tail_n} total={total_lines} omitted={omitted}")

        for line in selected_lines:
            print(line, flush=True)



def maybe_add_realtor_arg(cmd, realtor_id):
    if realtor_id:
        cmd.extend(["--realtor-id", str(int(realtor_id))])
    return cmd


def build_step_commands(args, realtor):
    """
    중개사 1명에 대한 V2 workday 준비+발행 단계.
    실제 파일명은 현재 프로젝트에서 사용 중인 V2 파일을 우선하고,
    없으면 V1 호환 파일로 fallback한다.
    """
    py = sys.executable
    realtor_id = int(realtor["realtor_id"])

    steps = []

    # 1) Session check
    if script_exists("workers/naver_session_check_worker.py"):
        cmd = [
            py,
            "workers/naver_session_check_worker.py",
            "--limit",
            "1",
            "--auto-relogin",
            "--realtor-id",
            str(realtor_id),
        ]
        # STEP110-01
        if args.include_auto_off or args.test_realtors:
            cmd.append("--include-disabled")
        steps.append({
            "name": "session_check",
            "command": cmd,
            "required": False,
            "timeout_seconds": args.session_timeout_seconds,
        })

    # 2-A) Current article collect / schedule
    # STEP109-03:
    # deleted_article_sync만으로는 신규/전환 중개사의 현재 매물 snapshot이 비어 있을 수 있다.
    # 먼저 schedule_daily_articles_v2로 네이버 전체매물 페이지를 실제로 열어 current article 목록을 확보한다.
    if script_exists("jobs/schedule_daily_articles_v2.py"):
        cmd = [
            py,
            "jobs/schedule_daily_articles_v2.py",
            "--realtor-id",
            str(realtor_id),
            "--max-pages",
            str(args.schedule_max_pages),
            "--draft-buffer",
            str(args.schedule_snapshot_buffer if args.force_schedule_snapshot else args.schedule_draft_buffer),
            "--replenish-count",
            str(args.schedule_replenish_count),
            "--no-db-lock",
        ]
        if args.schedule_headless:
            cmd.extend(["--headless", "1"])
        steps.append({
            "name": "schedule_current_articles_v2",
            "command": cmd,
            "required": True,
            "timeout_seconds": args.schedule_timeout_seconds,
        })

    # 2-B) Snapshot / current articles
    # V2 snapshot이 현재 core.deleted_article_sync --use-snapshot 기반으로도 snapshot을 생성/비교하므로 우선 사용한다.
    # 별도 snapshot 전용 모듈이 있으면 그것을 우선한다.
    if script_exists("core/realtor_article_snapshot.py"):
        cmd = [py, "-m", "core.realtor_article_snapshot", "--realtor-id", str(realtor_id)]
        if args.snapshot_dry_run:
            cmd.append("--dry-run")
        steps.append({
            "name": "snapshot_current_articles",
            "command": cmd,
            "required": True,
            "timeout_seconds": args.snapshot_timeout_seconds,
        })
    elif script_exists("core/deleted_article_sync.py"):
        cmd = [
            py,
            "-m",
            "core.deleted_article_sync",
            "--realtor-id",
            str(realtor_id),
            "--use-snapshot",
        ]

        # STEP109-02:
        # Workday 파이프라인의 snapshot 단계는 비공개 비교가 아니라
        # v2_publish_worker snapshot guard 통과용 SUCCESS baseline 확보가 목적이다.
        # 특히 신규 V2 전환 중개사(realtor_id=2 세종파라곤 등)는 이전 SUCCESS snapshot이 없을 수 있으므로
        # migration baseline 전용 모드로 현재 전체매물 snapshot을 먼저 확정한다.
        if args.snapshot_baseline_only:
            cmd.append("--migration-baseline-only")

        if args.snapshot_dry_run:
            cmd.append("--dry-run")

        steps.append({
            "name": "snapshot_current_articles",
            "command": cmd,
            "required": True,
            "timeout_seconds": args.snapshot_timeout_seconds,
        })

    # 3) Detail collect
    if script_exists("services/naver_response_article_fetcher_v2.py"):
        cmd = [
            py,
            "services/naver_response_article_fetcher_v2.py",
            "--realtor-id",
            str(realtor_id),
            "--limit",
            str(args.fetch_limit),
        ]
        steps.append({
            "name": "detail_collect_v2",
            "command": cmd,
            "required": True,
            "timeout_seconds": args.fetch_timeout_seconds,
        })
    elif script_exists("services/naver_response_article_fetcher.py"):
        cmd = [
            py,
            "services/naver_response_article_fetcher.py",
            "--limit",
            str(args.fetch_limit),
        ]
        # V1 fetcher가 realtor-id를 지원하는 환경이면 전달한다.
        if args.pass_realtor_id_to_legacy_fetcher:
            cmd.extend(["--realtor-id", str(realtor_id)])
        steps.append({
            "name": "detail_collect_legacy",
            "command": cmd,
            "required": True,
            "timeout_seconds": args.fetch_timeout_seconds,
        })

    # 4) Draft generation
    if script_exists("jobs/generate_blog_drafts_v2.py"):
        cmd = [
            py,
            "-m",
            "jobs.generate_blog_drafts_v2",
            "--realtor-id",
            str(realtor_id),
            "--limit",
            str(args.draft_limit),
        ]
        if args.force_draft:
            cmd.append("--force")
        steps.append({
            "name": "draft_generate_v2",
            "command": cmd,
            "required": True,
            "timeout_seconds": args.draft_timeout_seconds,
        })
    elif script_exists("jobs/generate_blog_drafts.py"):
        cmd = [
            py,
            "jobs/generate_blog_drafts.py",
            "--realtor-id",
            str(realtor_id),
            "--limit",
            str(args.draft_limit),
        ]
        steps.append({
            "name": "draft_generate_legacy",
            "command": cmd,
            "required": True,
            "timeout_seconds": args.draft_timeout_seconds,
        })

    # 5) Publish
    cmd = [
        py,
        "-m",
        "workers.v2_publish_worker",
        "--realtor-id",
        str(realtor_id),
        "--limit",
        str(args.publish_limit),
    ]

    if args.publish_dry_run:
        cmd.append("--check-final-publish")
    else:
        cmd.append("--execute")

    if args.hpe_mode:
        cmd.extend(["--hpe-mode", str(args.hpe_mode)])

    if args.headless:
        cmd.append("--headless")

    steps.append({
        "name": "publish_v2_worker",
        "command": cmd,
        "required": True,
        "timeout_seconds": args.publish_timeout_seconds,
    })

    if args.only_step:
        names = {x.strip() for x in str(args.only_step).split(",") if x.strip()}
        steps = [s for s in steps if s["name"] in names]

    if args.skip_session:
        steps = [s for s in steps if s["name"] != "session_check"]
    if args.skip_schedule:
        steps = [s for s in steps if s["name"] != "schedule_current_articles_v2"]
    if args.skip_snapshot:
        steps = [s for s in steps if s["name"] != "snapshot_current_articles"]
    if args.skip_fetch:
        steps = [s for s in steps if not s["name"].startswith("detail_collect")]
    if args.skip_draft:
        steps = [s for s in steps if not s["name"].startswith("draft_generate")]
    if args.skip_publish:
        steps = [s for s in steps if s["name"] != "publish_v2_worker"]

    return steps


def run_realtor_workday(args, realtor, run_id=None):
    realtor_id = int(realtor["realtor_id"])
    office_name = realtor.get("office_name") or ""

    log("-" * 80)
    log("[V2 REALTOR WORKDAY START] " + safe_json({
        "realtor_id": realtor_id,
        "office_name": office_name,
        "auto_publish_enabled": realtor.get("auto_publish_enabled"),
    }))

    add_pipeline_log(
        run_id=run_id,
        level="info",
        step_name="v2_realtor_workday_start",
        message=f"realtor workday start: {realtor_id}",
        context={"realtor": realtor},
    )

    steps = build_step_commands(args, realtor)

    results = []
    success = 0
    failed = 0
    skipped = 0

    for step in steps:
        if args.stop_when_until_hour and hour_float() >= float(args.until_hour):
            skipped += 1
            result = {
                "ok": True,
                "skipped": True,
                "reason": "until_hour_reached",
                "step": step["name"],
            }
            results.append(result)
            log("[V2 REALTOR STEP SKIP UNTIL HOUR] " + safe_json(result))
            break

        name = step["name"]
        command = step["command"]
        timeout_seconds = step.get("timeout_seconds")

        log("[V2 REALTOR STEP START] " + safe_json({
            "realtor_id": realtor_id,
            "step": name,
            "command": " ".join(command),
            "timeout_seconds": timeout_seconds,
        }))

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

        result = run_subprocess(
            command,
            dry_run=args.dry_run,
            timeout_seconds=timeout_seconds,
        )
        if name == "session_check":
            session_file = ROOT / "storage" / "naver_sessions" / f"naver_session_{realtor_id}.json"
            if session_file.exists():
                log(f"[SESSION JSON OK] {session_file}")
            else:
                log(f"[SESSION JSON MISSING] {session_file}")
        print_tail("[V2 REALTOR STEP RESULT]", result, tail_lines=args.tail_lines)

        item = {
            "step": name,
            "ok": result.get("ok"),
            "returncode": result.get("returncode"),
            "elapsed_seconds": result.get("elapsed_seconds"),
            "cmd": result.get("cmd"),
            "internal_failure": result.get("internal_failure"),
        }
        results.append(item)

        add_pipeline_log(
            run_id=run_id,
            level="info" if result.get("ok") else "error",
            step_name=name,
            message=f"step {'done' if result.get('ok') else 'failed'}: {name}",
            context={"realtor_id": realtor_id, "result": item},
        )

        if result.get("ok"):
            success += 1
        else:
            failed += 1
            if step.get("required") or args.stop_realtor_on_step_error:
                log("[V2 REALTOR STOP ON STEP ERROR] " + safe_json({
                    "realtor_id": realtor_id,
                    "step": name,
                    "required": step.get("required"),
                }))
                break

    ok = failed == 0
    summary = {
        "realtor_id": realtor_id,
        "office_name": office_name,
        "ok": ok,
        "steps_total": len(steps),
        "success": success,
        "failed": failed,
        "skipped": skipped,
        "results": results,
    }

    log("[V2 REALTOR WORKDAY DONE] " + safe_json(summary))
    return summary


def run_workday(args):
    if args.enforce_time_window and not is_workday_window(args.start_hour, args.until_hour):
        result = {
            "ok": True,
            "skipped": True,
            "reason": "outside_workday_window",
            "now": now_local().strftime("%Y-%m-%d %H:%M:%S"),
            "start_hour": args.start_hour,
            "until_hour": args.until_hour,
        }
        log("[V2 WORKDAY RUNNER SKIPPED] " + safe_json(result))
        return result

    targets, seed = build_targets(args)

    run_id = create_pipeline_run(
        run_type="v2_workday_pipeline",
        meta={
            "mode": "workday",
            "targets_count": len(targets),
            "shuffle_seed": seed,
            "test_realtors": args.test_realtors,
            "include_auto_off": args.include_auto_off or args.test_realtors,
            "start_hour": args.start_hour,
            "until_hour": args.until_hour,
            "dry_run": args.dry_run,
            "publish_dry_run": args.publish_dry_run,
            "worker_server": worker_server_name(),
        },
    )

    log("=" * 80)
    log("[V2 WORKDAY RUNNER START]")
    log("TIME: " + now_local().strftime("%Y-%m-%d %H:%M:%S"))
    log("targets_count: " + str(len(targets)))
    log("shuffle: " + str(args.shuffle_realtors) + " seed: " + str(seed))
    log("test_realtors: " + str(args.test_realtors))
    log("include_auto_off: " + str(args.include_auto_off or args.test_realtors))
    log("window: " + f"{args.start_hour}~{args.until_hour}")
    log("=" * 80)

    for idx, t in enumerate(targets, start=1):
        log("[V2 WORKDAY TARGET] " + safe_json({"index": idx, "total": len(targets), **t}))

    results = []
    success = 0
    failed = 0
    skipped = 0

    for idx, realtor in enumerate(targets, start=1):
        if args.stop_when_until_hour and hour_float() >= float(args.until_hour):
            skipped += 1
            log("[V2 WORKDAY RUNNER STOP] " + safe_json({
                "reason": "until_hour_reached",
                "next_realtor_id": realtor.get("realtor_id"),
                "now": now_local().strftime("%Y-%m-%d %H:%M:%S"),
            }))
            break

        result = run_realtor_workday(args, realtor, run_id=run_id)
        results.append(result)

        if result.get("ok"):
            success += 1
        else:
            failed += 1
            if args.stop_on_error:
                break

        if args.sleep_between_realtors_max > 0:
            lo = max(0, float(args.sleep_between_realtors_min))
            hi = max(lo, float(args.sleep_between_realtors_max))
            sleep_sec = random.uniform(lo, hi)
            log(f"[V2 BETWEEN REALTOR SLEEP] {sleep_sec:.2f}s")
            if not args.dry_run:
                time.sleep(sleep_sec)

    ok = failed == 0
    summary = {
        "ok": ok,
        "targets_count": len(targets),
        "processed": len(results),
        "success": success,
        "failed": failed,
        "skipped": skipped,
        "shuffle_seed": seed,
        "results": results,
    }

    finish_pipeline_run(
        run_id=run_id,
        status="success" if ok else "failed",
        total=len(targets),
        success=success,
        failed=failed,
        skipped=skipped,
        message="v2 workday pipeline completed" if ok else "v2 workday pipeline completed with failures",
        meta=summary,
    )

    log("=" * 80)
    log("[V2 WORKDAY RUNNER DONE] " + safe_json(summary))
    log("=" * 80)

    return summary


def build_private_cmd(args, realtor_id):
    cmd = [
        sys.executable,
        "-m",
        "core.deleted_article_sync",
        "--realtor-id",
        str(int(realtor_id)),
        "--use-snapshot",
    ]
    if args.private_dry_run or args.dry_run:
        cmd.append("--dry-run")
    return cmd


def run_private(args):
    if args.enforce_time_window and not is_private_window(args.until_hour):
        result = {
            "ok": True,
            "skipped": True,
            "reason": "before_private_window",
            "now": now_local().strftime("%Y-%m-%d %H:%M:%S"),
            "until_hour": args.until_hour,
        }
        log("[V2 PRIVATE RUNNER SKIPPED] " + safe_json(result))
        return result

    targets, seed = build_targets(args)

    run_id = create_pipeline_run(
        run_type="v2_private_pipeline",
        meta={
            "mode": "private",
            "targets_count": len(targets),
            "shuffle_seed": seed,
            "dry_run": args.dry_run or args.private_dry_run,
            "worker_server": worker_server_name(),
        },
    )

    log("=" * 80)
    log("[V2 PRIVATE RUNNER START]")
    log("TIME: " + now_local().strftime("%Y-%m-%d %H:%M:%S"))
    log("targets_count: " + str(len(targets)))
    log("shuffle: " + str(args.shuffle_realtors) + " seed: " + str(seed))
    log("=" * 80)

    results = []
    success = 0
    failed = 0

    for idx, realtor in enumerate(targets, start=1):
        realtor_id = int(realtor["realtor_id"])
        cmd = build_private_cmd(args, realtor_id)

        log("-" * 80)
        log("[V2 PRIVATE REALTOR START] " + safe_json({
            "index": idx,
            "total": len(targets),
            "realtor_id": realtor_id,
            "office_name": realtor.get("office_name"),
            "cmd": " ".join(cmd),
        }))

        result = run_subprocess(cmd, dry_run=False, timeout_seconds=args.private_timeout_seconds)
        print_tail("[V2 PRIVATE REALTOR RESULT]", result, tail_lines=args.tail_lines)

        item = {
            "realtor_id": realtor_id,
            "office_name": realtor.get("office_name"),
            "ok": result.get("ok"),
            "returncode": result.get("returncode"),
            "elapsed_seconds": result.get("elapsed_seconds"),
            "cmd": result.get("cmd"),
            "internal_failure": result.get("internal_failure"),
        }
        results.append(item)

        if result.get("ok"):
            success += 1
        else:
            failed += 1
            if args.stop_on_error:
                break

        if args.sleep_between_realtors_max > 0:
            lo = max(0, float(args.sleep_between_realtors_min))
            hi = max(lo, float(args.sleep_between_realtors_max))
            sleep_sec = random.uniform(lo, hi)
            log(f"[V2 BETWEEN PRIVATE REALTOR SLEEP] {sleep_sec:.2f}s")
            if not args.dry_run:
                time.sleep(sleep_sec)

    ok = failed == 0
    summary = {
        "ok": ok,
        "targets_count": len(targets),
        "processed": len(results),
        "success": success,
        "failed": failed,
        "shuffle_seed": seed,
        "results": results,
    }

    finish_pipeline_run(
        run_id=run_id,
        status="success" if ok else "failed",
        total=len(targets),
        success=success,
        failed=failed,
        skipped=0,
        message="v2 private pipeline completed" if ok else "v2 private pipeline completed with failures",
        meta=summary,
    )

    log("=" * 80)
    log("[V2 PRIVATE RUNNER DONE] " + safe_json(summary))
    log("=" * 80)
    return summary


def main():
    parser = argparse.ArgumentParser(description="STEP109-01 V2 Workday Pipeline Orchestrator")

    parser.add_argument("--mode", choices=["workday", "private", "both"], default="workday")

    parser.add_argument("--realtor-ids", default="")
    parser.add_argument("--test-realtors", action="store_true", help="테스트 기간 realtor_id 1,2 실행")
    parser.add_argument("--limit", type=int, default=0)

    parser.add_argument("--shuffle-realtors", action="store_true", default=True)
    parser.add_argument("--no-shuffle-realtors", dest="shuffle_realtors", action="store_false")
    parser.add_argument("--shuffle-seed", default="")

    parser.add_argument("--include-auto-off", action="store_true")
    parser.add_argument("--ignore-paid-filter", action="store_true")

    parser.add_argument("--start-hour", type=float, default=DEFAULT_START_HOUR)
    parser.add_argument("--until-hour", type=float, default=DEFAULT_UNTIL_HOUR)
    parser.add_argument("--enforce-time-window", action="store_true")
    parser.add_argument("--stop-when-until-hour", action="store_true", default=True)
    parser.add_argument("--no-stop-when-until-hour", dest="stop_when_until_hour", action="store_false")

    parser.add_argument("--dry-run", action="store_true")
    parser.add_argument("--publish-dry-run", action="store_true")
    parser.add_argument("--private-dry-run", action="store_true")
    # Workday 기본은 snapshot commit + baseline 확보.
    # 실제 비공개 비교는 --mode private 에서 수행한다.
    parser.add_argument("--snapshot-dry-run", action="store_true", default=False)
    parser.add_argument("--snapshot-commit", dest="snapshot_dry_run", action="store_false")
    parser.add_argument("--snapshot-baseline-only", action="store_true", default=True)
    parser.add_argument("--snapshot-compare-mode", dest="snapshot_baseline_only", action="store_false")

    parser.add_argument("--skip-session", action="store_true")
    parser.add_argument("--skip-schedule", action="store_true")
    parser.add_argument("--skip-snapshot", action="store_true")
    parser.add_argument("--skip-fetch", action="store_true")
    parser.add_argument("--skip-draft", action="store_true")
    parser.add_argument("--skip-publish", action="store_true")
    parser.add_argument("--only-step", default="", help="쉼표 구분: session_check,schedule_current_articles_v2,snapshot_current_articles,detail_collect_v2,draft_generate_v2,publish_v2_worker")

    parser.add_argument("--schedule-max-pages", type=int, default=30)
    parser.add_argument("--schedule-draft-buffer", type=int, default=45)

    # STEP109-04:
    # V2 baseline snapshot이 없는 중개사는 draft buffer가 충분해도
    # schedule_daily_articles_v2가 [SKIP] draft buffer enough로 끝나면
    # 네이버 전체매물 페이지를 열지 않아 SUCCESS snapshot이 생성되지 않는다.
    # Workday에서는 snapshot 확보가 필수이므로 기본적으로 큰 buffer target을 임시 전달한다.
    parser.add_argument("--force-schedule-snapshot", action="store_true", default=True)
    parser.add_argument("--no-force-schedule-snapshot", dest="force_schedule_snapshot", action="store_false")
    parser.add_argument("--schedule-snapshot-buffer", type=int, default=9999)

    parser.add_argument("--schedule-replenish-count", type=int, default=1)
    parser.add_argument("--schedule-timeout-seconds", type=int, default=1800)
    parser.add_argument("--schedule-headless", action="store_true", default=True)
    parser.add_argument("--no-schedule-headless", dest="schedule_headless", action="store_false")

    parser.add_argument("--fetch-limit", type=int, default=5)
    parser.add_argument("--draft-limit", type=int, default=5)
    parser.add_argument("--publish-limit", type=int, default=1)
    parser.add_argument("--force-draft", action="store_true")

    parser.add_argument("--hpe-mode", default="test")
    parser.add_argument("--headless", action="store_true")

    parser.add_argument("--session-timeout-seconds", type=int, default=240)
    parser.add_argument("--snapshot-timeout-seconds", type=int, default=600)
    parser.add_argument("--fetch-timeout-seconds", type=int, default=900)
    parser.add_argument("--draft-timeout-seconds", type=int, default=900)
    parser.add_argument("--publish-timeout-seconds", type=int, default=2400)
    parser.add_argument("--private-timeout-seconds", type=int, default=900)

    parser.add_argument("--sleep-between-realtors-min", type=float, default=7)
    parser.add_argument("--sleep-between-realtors-max", type=float, default=35)

    parser.add_argument("--stop-on-error", action="store_true")
    parser.add_argument("--stop-realtor-on-step-error", action="store_true", default=True)
    parser.add_argument("--no-stop-realtor-on-step-error", dest="stop_realtor_on_step_error", action="store_false")
    parser.add_argument("--tail-lines", type=int, default=300, help="하위 프로세스 stdout/stderr 출력 줄 수. 0이면 전체 출력")

    parser.add_argument("--lock-ttl-minutes", type=int, default=900)
    parser.add_argument("--no-db-lock", action="store_true")

    # 구형 fetcher 호환 옵션
    parser.add_argument("--pass-realtor-id-to-legacy-fetcher", action="store_true")

    args = parser.parse_args()

    lock_owner = None
    lock_name = "v2_workday_pipeline" if args.mode != "private" else "v2_private_pipeline"

    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:
                result = {"ok": True, "skipped": True, "reason": "lock_not_acquired", "lock_name": lock_name}
                print("[RESULT]", safe_json(result))
                return

        if args.mode == "workday":
            result = run_workday(args)
        elif args.mode == "private":
            result = run_private(args)
        else:
            result = {
                "workday": run_workday(args),
                "private": run_private(args),
            }

        print("[RESULT]", safe_json(result))

    finally:
        if lock_owner:
            release_db_lock(lock_name, lock_owner)


if __name__ == "__main__":
    main()
