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

"""
refresh_realtor_articles_seen.py

역할:
- 유료/상품/자동발행 조건을 충족하는 중개사만 대상으로 한다.
- 중개사 네이버 부동산 전체매물 페이지에서 현재 articleNo 목록을 수집한다.
- 발견된 articleNo는 blog_realtor_articles.last_seen_at = NOW() 로 갱신한다.
- 이후 check_removed_articles.py가 last_seen_at 기준으로 삭제 의심/확정을 판단한다.

사용 예:
python workers/refresh_realtor_articles_seen.py --limit 10 --dry-run
python workers/refresh_realtor_articles_seen.py --limit 10
python workers/refresh_realtor_articles_seen.py --realtor-id 2 --headless 0 --dry-run
"""

import os
import sys
import json
import socket
import argparse
import traceback
import urllib.request
import urllib.parse
import re
from datetime import datetime

CURRENT_DIR = os.path.dirname(os.path.abspath(__file__))
ROOT_DIR = os.path.dirname(CURRENT_DIR)

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

from db import get_conn
from services.naver_realtor_response_finder import NaverRealtorResponseFinder


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


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


def worker_token(prefix):
    try:
        host = socket.gethostname()
    except Exception:
        host = "unknown-host"
    return f"{prefix}:{host}:{os.getpid()}:{datetime.now().strftime('%Y%m%d%H%M%S')}"


def acquire_db_lock(lock_name, ttl_minutes=180):
    owner_token = worker_token(lock_name)
    conn = None

    try:
        conn = get_conn()

        with conn.cursor() as cur:
            cur.execute("""
                DELETE FROM blog_app_locks
                WHERE expires_at < NOW()
            """)

            cur.execute("""
                INSERT INTO blog_app_locks
                (
                    lock_name,
                    locked_at,
                    expires_at,
                    owner_token
                )
                VALUES
                (
                    %s,
                    NOW(),
                    DATE_ADD(NOW(), INTERVAL %s MINUTE),
                    %s
                )
            """, (
                str(lock_name)[:100],
                int(ttl_minutes),
                str(owner_token)[:100],
            ))

        conn.commit()
        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] {lock_name} / {e}")

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


def create_pipeline_run(run_type="refresh_realtor_articles_seen", meta=None):
    conn = None

    try:
        conn = get_conn()

        with conn.cursor() as cur:
            cur.execute("""
                INSERT INTO blog_pipeline_runs
                (
                    run_type,
                    status,
                    started_at,
                    message,
                    meta_json,
                    created_at
                )
                VALUES
                (
                    %s,
                    'running',
                    NOW(),
                    %s,
                    %s,
                    NOW()
                )
            """, (
                run_type,
                "refresh_realtor_articles_seen started",
                safe_json_dumps(meta or {})
            ))
            run_id = cur.lastrowid

        conn.commit()
        log(f"[PIPELINE RUN START] run_id={run_id}, run_type={run_type}")
        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
            """, (
                status,
                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 {}),
                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="", article_no=None, realtor_id=None, context=None):
    conn = None

    try:
        conn = get_conn()

        with conn.cursor() as cur:
            cur.execute("""
                INSERT INTO blog_pipeline_logs
                (
                    run_id,
                    level,
                    step_name,
                    article_no,
                    realtor_id,
                    message,
                    context_json,
                    created_at
                )
                VALUES
                (
                    %s,
                    %s,
                    %s,
                    %s,
                    %s,
                    %s,
                    %s,
                    NOW()
                )
            """, (
                run_id,
                str(level or "info")[:20],
                str(step_name or "")[:100],
                str(article_no or "")[:30] if article_no else None,
                realtor_id,
                str(message or "")[:2000],
                safe_json_dumps(context or {})
            ))

        conn.commit()

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

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


def get_eligible_realtors(realtor_id=None, limit=100):
    conn = get_conn()

    try:
        params = []

        sql = """
            SELECT
                r.id AS realtor_id,
                r.office_name,
                r.naver_realtor_id
            FROM blog_realtors r
            INNER JOIN blog_publish_settings ps
                ON ps.realtor_id = r.id
               AND ps.auto_publish_enabled = 1
            INNER JOIN blog_paid_memberships pm
                ON pm.realtor_id = r.id
               AND pm.membership_status = 'active'
               AND pm.payment_status IN ('paid', 'partial')
               AND pm.service_type IN ('blog_auto', 'package')
               AND pm.start_date IS NOT NULL
               AND pm.start_date <= CURDATE()
               AND COALESCE(pm.end_date, pm.auto_end_date) IS NOT NULL
               AND COALESCE(pm.end_date, pm.auto_end_date) >= CURDATE()
            WHERE r.status = 'active'
              AND COALESCE(r.naver_realtor_id, '') <> ''
        """

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

        sql += """
            GROUP BY r.id, r.office_name, r.naver_realtor_id
            ORDER BY r.id ASC
            LIMIT %s
        """
        params.append(int(limit))

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

        return rows or []

    finally:
        conn.close()


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

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


def upsert_seen_article(realtor, article_no, dry_run=False):
    article_no = str(article_no or "").strip()

    if not article_no:
        return False

    realtor_id = int(realtor.get("realtor_id") or 0)

    if dry_run:
        log(f"[DRY SEEN] realtor_id={realtor_id}, article_no={article_no}")
        return True

    conn = get_conn()

    try:
        columns = table_columns(conn, "blog_realtor_articles")

        data = {}

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

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

        if "article_name" in columns:
            data["article_name"] = "네이버 부동산 현재 매물"

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

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

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

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

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

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

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

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

        if "collect_status" in columns:
            data["collect_status"] = "pending"

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

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

        keys = list(data.keys())

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

        update_sets = []

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

            if key == "last_seen_at":
                update_sets.append("last_seen_at = NOW()")
            elif key in ["removed_detected_at", "removal_confirmed_at"]:
                update_sets.append(f"{key} = NULL")
            elif key == "removal_status":
                update_sets.append("removal_status = NULL")
            else:
                update_sets.append(f"{key} = VALUES({key})")

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

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

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

        conn.commit()
        return True

    finally:
        conn.close()



USER_AGENT = (
    "Mozilla/5.0 (Windows NT 10.0; Win64; x64) "
    "AppleWebKit/537.36 (KHTML, like Gecko) "
    "Chrome/136.0.0.0 Safari/537.36"
)


def uniq_keep_order(items):
    seen = set()
    result = []

    for item in items:
        item = str(item or "").strip()

        if not item or item in seen:
            continue

        seen.add(item)
        result.append(item)

    return result


def extract_article_nos_from_text(text):
    text = str(text or "")

    patterns = [
        r'"articleNo"\s*:\s*"([0-9]{8,12})"',
        r'"articleNo"\s*:\s*([0-9]{8,12})',
        r'"atclNo"\s*:\s*"([0-9]{8,12})"',
        r'"atclNo"\s*:\s*([0-9]{8,12})',
        r'"articleNumber"\s*:\s*"([0-9]{8,12})"',
        r'"articleNumber"\s*:\s*([0-9]{8,12})',
        r"articleNo=([0-9]{8,12})",
    ]

    found = []

    for pattern in patterns:
        found.extend(re.findall(pattern, text, flags=re.I))

    return uniq_keep_order(found)


def fast_fetch_realtor_article_nos(naver_realtor_id, max_pages=10, timeout=10):
    """
    Playwright 없이 /api/articles 를 직접 호출한다.
    실패하면 refresh worker가 Playwright fallback을 사용한다.
    """
    naver_realtor_id = str(naver_realtor_id or "").strip()

    if not naver_realtor_id:
        return {
            "ok": False,
            "article_nos": [],
            "error": "naver_realtor_id 없음",
            "debug": [],
        }

    all_nos = []
    debug = []

    for page_no in range(1, int(max_pages or 1) + 1):
        query = urllib.parse.urlencode({
            "realtorId": naver_realtor_id,
            "page": page_no,
            "order": "rank",
            "tradeType": "",
            "realEstateType": "",
            "isFixed": "false",
        })

        url = f"https://new.land.naver.com/api/articles?{query}"

        req = urllib.request.Request(
            url,
            headers={
                "User-Agent": USER_AGENT,
                "Accept": "application/json, text/plain, */*",
                "Referer": f"https://new.land.naver.com/offices?realtorId={urllib.parse.quote(naver_realtor_id)}",
                "Origin": "https://new.land.naver.com",
            }
        )

        try:
            with urllib.request.urlopen(req, timeout=int(timeout or 10)) as res:
                status = getattr(res, "status", 0)
                raw = res.read()

            text = raw.decode("utf-8", errors="ignore")
            article_nos = extract_article_nos_from_text(text)

            debug.append({
                "url": url,
                "status": status,
                "count": len(article_nos),
                "article_nos": article_nos[:50],
            })

            log(f"[FAST API RESPONSE] page={page_no}, status={status}, count={len(article_nos)}")

            if status != 200:
                return {
                    "ok": False,
                    "article_nos": uniq_keep_order(all_nos),
                    "error": f"status={status}",
                    "debug": debug,
                }

            if not article_nos:
                break

            all_nos.extend(article_nos)

            if len(article_nos) < 20:
                break

        except Exception as e:
            debug.append({
                "url": url,
                "error": str(e),
            })

            return {
                "ok": False,
                "article_nos": uniq_keep_order(all_nos),
                "error": str(e),
                "debug": debug,
            }

    article_nos = uniq_keep_order(all_nos)

    return {
        "ok": bool(article_nos),
        "article_nos": article_nos,
        "count": len(article_nos),
        "error": "" if article_nos else "articleNo 없음",
        "debug": debug,
    }


class LazyNaverRealtorFinder:
    def __init__(self, headless=True):
        self.headless = headless
        self.inner = None

    def find_all_articles(self, *args, **kwargs):
        if self.inner is None:
            self.inner = NaverRealtorResponseFinder(headless=self.headless)
            self.inner.start()
        return self.inner.find_all_articles(*args, **kwargs)

    def close(self):
        if self.inner:
            self.inner.close()
            self.inner = None


def process_realtor(finder, realtor, wait_ms=2500, scroll_steps=1, dry_run=False, run_id=None, fast_api=False, fast_max_pages=10):
    realtor_id = int(realtor.get("realtor_id") or 0)
    office_name = str(realtor.get("office_name") or "")
    naver_realtor_id = str(realtor.get("naver_realtor_id") or "").strip()

    log("-" * 80)
    log(f"[REALTOR START] realtor_id={realtor_id}, office={office_name}, naver_realtor_id={naver_realtor_id}")

    result = None
    article_nos = []
    fetch_mode = "playwright"

    if fast_api:
        fast_result = fast_fetch_realtor_article_nos(
            naver_realtor_id=naver_realtor_id,
            max_pages=fast_max_pages,
        )

        if fast_result.get("ok"):
            result = fast_result
            article_nos = fast_result.get("article_nos") or []
            fetch_mode = "fast_api"
            log(f"[FAST API OK] realtor_id={realtor_id}, count={len(article_nos)}")
        else:
            log(f"[FAST API FALLBACK] realtor_id={realtor_id}, error={fast_result.get('error')}")

    if not article_nos:
        result = finder.find_all_articles(
            naver_realtor_id=naver_realtor_id,
            wait_ms=wait_ms,
            scroll_steps=scroll_steps,
            max_pages=fast_max_pages,
        )

        article_nos = result.get("article_nos") or []
        fetch_mode = "playwright"

    log(f"[CURRENT ARTICLE COUNT] realtor_id={realtor_id}, mode={fetch_mode}, count={len(article_nos)}")

    saved = 0

    for article_no in article_nos:
        if upsert_seen_article(realtor, article_no, dry_run=dry_run):
            saved += 1

        add_pipeline_log(
            run_id=run_id,
            level="info",
            step_name="article_seen",
            message="current article seen",
            article_no=article_no,
            realtor_id=realtor_id,
            context={
                "dry_run": dry_run,
                "naver_realtor_id": naver_realtor_id,
                "fetch_mode": fetch_mode,
            }
        )

    log(f"[REALTOR DONE] realtor_id={realtor_id}, seen={len(article_nos)}, saved={saved}")

    return {
        "seen": len(article_nos),
        "saved": saved,
    }


def run(limit=100, realtor_id=None, headless=True, wait_ms=2500, scroll_steps=1, dry_run=False, no_db_lock=False, lock_ttl_minutes=180, fast_api=True, fast_max_pages=10):
    lock_name = "refresh_realtor_articles_seen"
    lock_owner = None
    run_id = None
    total_realtors = 0
    total_seen = 0
    total_saved = 0
    failed = 0

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

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

        run_id = create_pipeline_run(
            run_type="refresh_realtor_articles_seen",
            meta={
                "limit": limit,
                "realtor_id": realtor_id,
                "headless": headless,
                "wait_ms": wait_ms,
                "scroll_steps": scroll_steps,
                "dry_run": dry_run,
                "fast_api": fast_api,
                "fast_max_pages": fast_max_pages,
            }
        )

        realtors = get_eligible_realtors(
            realtor_id=realtor_id,
            limit=limit
        )

        log("=" * 80)
        log(f"[TARGET REALTORS] {len(realtors)}")
        log(f"[HEADLESS] {headless}")
        log(f"[DRY RUN] {'ON' if dry_run else 'OFF'}")
        log(f"[FAST API] {'ON' if fast_api else 'OFF'}")
        log(f"[FAST MAX PAGES] {fast_max_pages}")
        log("=" * 80)

        finder = LazyNaverRealtorFinder(headless=headless)

        try:
            for realtor in realtors:
                total_realtors += 1

                try:
                    stats = process_realtor(
                        finder=finder,
                        realtor=realtor,
                        wait_ms=wait_ms,
                        scroll_steps=scroll_steps,
                        dry_run=dry_run,
                        run_id=run_id,
                        fast_api=fast_api,
                        fast_max_pages=fast_max_pages,
                    )

                    total_seen += int(stats.get("seen") or 0)
                    total_saved += int(stats.get("saved") or 0)

                except Exception as e:
                    failed += 1
                    traceback.print_exc()
                    add_pipeline_log(
                        run_id=run_id,
                        level="error",
                        step_name="realtor_exception",
                        message=str(e)[:2000],
                        realtor_id=realtor.get("realtor_id"),
                        context={"traceback": traceback.format_exc()}
                    )

        finally:
            finder.close()

        final_message = (
            f"realtors={total_realtors}, seen={total_seen}, "
            f"saved={total_saved}, failed={failed}"
        )

        finish_pipeline_run(
            run_id,
            status="success" if failed == 0 else "failed",
            total=total_realtors,
            success=total_realtors - failed,
            failed=failed,
            skipped=0,
            message=final_message,
            meta={
                "total_realtors": total_realtors,
                "total_seen": total_seen,
                "total_saved": total_saved,
                "failed": failed,
                "dry_run": dry_run,
                "fast_api": fast_api,
                "fast_max_pages": fast_max_pages,
            }
        )

        log("=" * 80)
        log("[REFRESH SEEN DONE]")
        log(final_message)
        log("=" * 80)

    finally:
        if lock_owner:
            release_db_lock(lock_name, lock_owner)


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

    parser.add_argument("--limit", type=int, default=100)
    parser.add_argument("--realtor-id", type=int, default=None)
    parser.add_argument("--headless", type=int, default=1)
    parser.add_argument("--wait-ms", type=int, default=2500)
    parser.add_argument("--scroll-steps", type=int, default=1)
    parser.add_argument("--dry-run", action="store_true")
    # STEP92:
    # 운영 기본값은 fast_api ON.
    # /api/articles page=1~N 직접 조회로 전체 현재매물 확인을 먼저 시도한다.
    # 실패 시에만 Playwright fallback을 사용한다.
    parser.add_argument("--use-fast-api", action="store_true", default=True)
    parser.add_argument("--no-fast-api", action="store_true", help="fast api를 끄고 Playwright만 사용")
    parser.add_argument("--fast-max-pages", type=int, default=10)
    parser.add_argument("--no-db-lock", action="store_true")
    parser.add_argument("--lock-ttl-minutes", type=int, default=180)

    args = parser.parse_args()

    run(
        limit=args.limit,
        realtor_id=args.realtor_id,
        headless=bool(args.headless),
        wait_ms=args.wait_ms,
        scroll_steps=args.scroll_steps,
        dry_run=args.dry_run,
        no_db_lock=args.no_db_lock,
        lock_ttl_minutes=args.lock_ttl_minutes,
        fast_api=(False if args.no_fast_api else bool(args.use_fast_api)),
        fast_max_pages=args.fast_max_pages,
    )


if __name__ == "__main__":
    main()
