# -*- coding: utf-8 -*-
"""
STEP103-22 V2 Scheduler Engine with ResumeExecutor

역할:
- 08:00~23:00 사이 V2 Pipeline을 연속 실행한다.
- Scheduler 시작 시 ResumeInspector로 마지막 실행 상태를 확인한다.
- ResumeExecutor가 PipelineContext에 resume_from 등을 자동 주입한다.
- Cursor를 blog_v2_scheduler_state에 저장하여 다음 실행 때 이어서 처리한다.
- 기본은 dry-run publish이며, execute_publish=True일 때만 실제 발행한다.

주의:
- 45건 버퍼 정책은 사용하지 않는다.
- V2 TargetSelector 기준으로 하루 처리 대상 1건을 결정한다.
"""

from datetime import datetime, time as dt_time

from services.v2.pipeline_context import PipelineContext
from services.v2.runner import build_default_manager, make_pipeline_id
from services.v2.scheduler_state_manager import SchedulerStateManager
from services.v2.resume_inspector import ResumeInspector
from services.v2.resume_executor import ResumeExecutor


class V2SchedulerEngine:
    def __init__(
        self,
        conn,
        scheduler_key="server_default",
        scheduler_type="server",
        window_start="08:00",
        window_end="23:00",
        execute_publish=False,
    ):
        self.conn = conn
        self.scheduler_key = scheduler_key
        self.scheduler_type = scheduler_type
        self.window_start = self._parse_hhmm(window_start)
        self.window_end = self._parse_hhmm(window_end)
        self.execute_publish = bool(execute_publish)
        self.state = SchedulerStateManager(
            conn,
            scheduler_key=scheduler_key,
            scheduler_type=scheduler_type,
        )
        self.resume_executor = ResumeExecutor()

    def run(self, max_realtors=None, realtor_id=None):
        resume_inspection = self.inspect_resume()
        resume_strategy_meta = None

        realtors = self.fetch_target_realtors(realtor_id=realtor_id)
        total = len(realtors)

        self.state.start_run(
            total_realtors=total,
            meta={
                "execute_publish": self.execute_publish,
                "window_start": self.window_start.strftime("%H:%M"),
                "window_end": self.window_end.strftime("%H:%M"),
                "max_realtors": max_realtors,
                "realtor_id": realtor_id,
                "resume_inspection": self.resume_inspection_to_dict(resume_inspection),
            },
        )

        if resume_inspection and resume_inspection.should_resume:
            self.state.heartbeat(
                message=(
                    f"resume detected: reason={resume_inspection.reason}, "
                    f"realtor_id={resume_inspection.realtor_id}, stage={resume_inspection.stage}"
                )
            )
            print(
                "[V2 RESUME DETECTED]",
                f"reason={resume_inspection.reason}",
                f"realtor_id={resume_inspection.realtor_id}",
                f"stage={resume_inspection.stage}",
                f"article_no={resume_inspection.article_no}",
            )
        else:
            print(
                "[V2 RESUME NOT REQUIRED]",
                f"reason={(resume_inspection.reason if resume_inspection else 'none')}"
            )

        processed = 0
        success = 0
        failed = 0
        skipped = 0

        try:
            start_index = self.resolve_start_index(realtors, resume_inspection)

            for idx in range(start_index, total):
                if max_realtors is not None and processed >= int(max_realtors):
                    self.state.stop_run(
                        status="stopped",
                        message=f"max_realtors reached: {max_realtors}",
                    )
                    break

                if not self.is_within_window():
                    self.state.stop_run(
                        status="paused",
                        message="window ended, resume next run",
                    )
                    break

                realtor = realtors[idx]
                rid = int(realtor.get("id") or 0)
                if not rid:
                    skipped += 1
                    continue

                pipeline_id = make_pipeline_id(rid)

                self.state.update_cursor(
                    realtor_id=rid,
                    stage="start",
                    pipeline_id=pipeline_id,
                    message=f"start realtor_id={rid}",
                )

                context = PipelineContext(
                    pipeline_id=pipeline_id,
                    realtor_id=rid,
                    version="v2",
                )
                context.put("scheduler_key", self.scheduler_key)
                context.put("execute_publish", self.execute_publish)
                context.put("draft_limit", 1)
                context.put("replenish_count", 1)

                # ResumeExecutor 연결
                strategy_result = None
                if resume_inspection and resume_inspection.should_resume and resume_inspection.realtor_id == rid:
                    strategy_result = self.resume_executor.prepare_context(
                        context,
                        resume_inspection,
                    )
                    resume_strategy_meta = self.resume_executor.to_meta(strategy_result)

                    print(
                        "[V2 RESUME STRATEGY]",
                        f"action={strategy_result.action}",
                        f"resume_from={strategy_result.resume_from}",
                        f"message={strategy_result.message}",
                    )

                try:
                    manager = build_default_manager()
                    manager.run(context)

                    processed += 1
                    success += 1

                    self.state.update_cursor(
                        realtor_id=rid,
                        stage=context.current_stage,
                        article_no=self._selected_article_no(context),
                        draft_id=self._selected_draft_id(context),
                        publish_queue_id=self._selected_publish_queue_id(context),
                        pipeline_id=pipeline_id,
                        message=f"done realtor_id={rid} stage={context.current_stage}",
                    )
                    self.state.mark_realtor_done(
                        success=True,
                        message=f"success realtor_id={rid}",
                    )

                except Exception as e:
                    processed += 1
                    failed += 1

                    self.state.update_cursor(
                        realtor_id=rid,
                        stage="failed",
                        pipeline_id=pipeline_id,
                        message=f"failed realtor_id={rid}: {e}",
                    )
                    self.state.mark_realtor_done(
                        success=False,
                        message=f"failed realtor_id={rid}: {e}",
                    )

                self.state.heartbeat(
                    message=f"processed={processed}/{total} success={success} failed={failed}"
                )

            else:
                self.state.stop_run(
                    status="success",
                    message=f"scheduler completed processed={processed}",
                )

            return {
                "total": total,
                "processed": processed,
                "success": success,
                "failed": failed,
                "skipped": skipped,
                "resume_inspection": self.resume_inspection_to_dict(resume_inspection),
                "resume_strategy": resume_strategy_meta,
            }

        except Exception as e:
            self.state.stop_run(
                status="failed",
                message="scheduler failed",
                error_message=str(e),
            )
            raise

    def inspect_resume(self):
        try:
            inspector = ResumeInspector(
                self.conn,
                scheduler_key=self.scheduler_key,
            )
            return inspector.inspect()
        except Exception as e:
            print("[V2 RESUME INSPECT ERROR]", str(e))
            return None

    def resume_inspection_to_dict(self, result):
        if not result:
            return None

        return {
            "should_resume": result.should_resume,
            "reason": result.reason,
            "pipeline_id": result.pipeline_id,
            "realtor_id": result.realtor_id,
            "stage": result.stage,
            "article_no": result.article_no,
            "draft_id": result.draft_id,
            "publish_queue_id": result.publish_queue_id,
        }

    def fetch_target_realtors(self, realtor_id=None):
        sql = """
            SELECT
                r.id,
                r.office_name,
                r.status
            FROM blog_realtors r
            WHERE COALESCE(r.status, 'active') = 'active'
        """
        params = []

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

        sql += " ORDER BY r.id ASC "

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

        return [dict(row) for row in rows]

    def resolve_start_index(self, realtors, resume_inspection=None):
        state = self.state.get_state() or {}
        cursor_realtor_id = state.get("cursor_realtor_id")

        if len(realtors) == 1:
            return 0

        if resume_inspection and resume_inspection.should_resume and resume_inspection.realtor_id:
            for idx, realtor in enumerate(realtors):
                rid = int(realtor.get("id") or 0)
                if rid == int(resume_inspection.realtor_id):
                    return idx

        if not cursor_realtor_id:
            return 0

        for idx, realtor in enumerate(realtors):
            rid = int(realtor.get("id") or 0)
            if rid > int(cursor_realtor_id):
                return idx

        return 0

    def is_within_window(self):
        now = datetime.now().time()
        return self.window_start <= now <= self.window_end

    def _parse_hhmm(self, value):
        hh, mm = str(value).split(":")[:2]
        return dt_time(int(hh), int(mm))

    def _selected_article_no(self, context):
        data = context.get("publish_data", {}) or {}
        target = data.get("target", {}) or {}
        return (
            target.get("article_no")
            or context.get("article_no")
            or context.get("resume_article_no")
            or None
        )

    def _selected_draft_id(self, context):
        data = context.get("publish_data", {}) or {}
        target = data.get("target", {}) or {}
        return (
            target.get("draft_id")
            or context.get("draft_id")
            or context.get("resume_draft_id")
            or None
        )

    def _selected_publish_queue_id(self, context):
        data = context.get("publish_data", {}) or {}
        target = data.get("target", {}) or {}
        return (
            target.get("publish_queue_id")
            or context.get("publish_queue_id")
            or context.get("resume_publish_queue_id")
            or None
        )
