# -*- coding: utf-8 -*-
"""
STEP103-02 V2 SchedulerStateManager

역할:
- blog_v2_scheduler_state 테이블을 통해 V2 Scheduler 상태를 저장/조회한다.
- 08:00~23:00 연속 실행 중 현재 cursor를 저장한다.
- 23:00 종료 또는 서버 재시작 후 다음 실행 때 이어서 처리할 수 있게 한다.

주의:
- V1 테이블은 수정하지 않는다.
"""

import json
from datetime import date


class SchedulerStateManager:
    table_name = "blog_v2_scheduler_state"

    def __init__(self, conn, scheduler_key="server_default", scheduler_type="server"):
        self.conn = conn
        self.scheduler_key = scheduler_key
        self.scheduler_type = scheduler_type

    def ensure_state(self):
        sql = f"""
            INSERT INTO {self.table_name}
            (
                scheduler_key,
                scheduler_type,
                run_date,
                status,
                created_at,
                updated_at
            )
            VALUES
            (
                %s,
                %s,
                CURDATE(),
                'idle',
                NOW(),
                NOW()
            )
            ON DUPLICATE KEY UPDATE
                scheduler_type = VALUES(scheduler_type),
                updated_at = NOW()
        """
        with self.conn.cursor() as cur:
            cur.execute(sql, (self.scheduler_key, self.scheduler_type))
        self.conn.commit()

    def get_state(self):
        self.ensure_state()

        sql = f"""
            SELECT *
            FROM {self.table_name}
            WHERE scheduler_key = %s
            LIMIT 1
        """
        with self.conn.cursor() as cur:
            cur.execute(sql, (self.scheduler_key,))
            row = cur.fetchone()

        return dict(row) if row else None

    def start_run(self, total_realtors=0, meta=None):
        self.ensure_state()

        meta_text = json.dumps(meta or {}, ensure_ascii=False, default=str)

        sql = f"""
            UPDATE {self.table_name}
            SET
                run_date = CURDATE(),
                status = 'running',
                total_realtors = %s,
                processed_realtors = 0,
                success_count = 0,
                failed_count = 0,
                skipped_count = 0,
                last_message = 'scheduler started',
                last_error_message = NULL,
                started_at = NOW(),
                stopped_at = NULL,
                last_heartbeat_at = NOW(),
                meta_text = %s,
                updated_at = NOW()
            WHERE scheduler_key = %s
        """

        with self.conn.cursor() as cur:
            cur.execute(sql, (int(total_realtors or 0), meta_text, self.scheduler_key))
        self.conn.commit()

    def heartbeat(self, message=None):
        sql = f"""
            UPDATE {self.table_name}
            SET
                last_heartbeat_at = NOW(),
                last_message = COALESCE(%s, last_message),
                updated_at = NOW()
            WHERE scheduler_key = %s
        """
        with self.conn.cursor() as cur:
            cur.execute(sql, (message, self.scheduler_key))
        self.conn.commit()

    def update_cursor(
        self,
        realtor_id=None,
        stage=None,
        article_no=None,
        draft_id=None,
        publish_queue_id=None,
        pipeline_id=None,
        message=None,
    ):
        sql = f"""
            UPDATE {self.table_name}
            SET
                cursor_realtor_id = %s,
                cursor_stage = %s,
                cursor_article_no = %s,
                cursor_draft_id = %s,
                cursor_publish_queue_id = %s,
                last_pipeline_id = %s,
                last_message = %s,
                last_heartbeat_at = NOW(),
                updated_at = NOW()
            WHERE scheduler_key = %s
        """

        with self.conn.cursor() as cur:
            cur.execute(
                sql,
                (
                    self._to_int_or_none(realtor_id),
                    stage,
                    str(article_no) if article_no else None,
                    self._to_int_or_none(draft_id),
                    self._to_int_or_none(publish_queue_id),
                    pipeline_id,
                    message,
                    self.scheduler_key,
                ),
            )
        self.conn.commit()

    def mark_realtor_done(self, success=True, skipped=False, message=None):
        if skipped:
            count_sql = "skipped_count = skipped_count + 1"
        elif success:
            count_sql = "success_count = success_count + 1"
        else:
            count_sql = "failed_count = failed_count + 1"

        sql = f"""
            UPDATE {self.table_name}
            SET
                processed_realtors = processed_realtors + 1,
                {count_sql},
                last_message = COALESCE(%s, last_message),
                last_heartbeat_at = NOW(),
                updated_at = NOW()
            WHERE scheduler_key = %s
        """

        with self.conn.cursor() as cur:
            cur.execute(sql, (message, self.scheduler_key))
        self.conn.commit()

    def stop_run(self, status="stopped", message=None, error_message=None):
        sql = f"""
            UPDATE {self.table_name}
            SET
                status = %s,
                last_message = COALESCE(%s, last_message),
                last_error_message = %s,
                stopped_at = NOW(),
                last_heartbeat_at = NOW(),
                updated_at = NOW()
            WHERE scheduler_key = %s
        """

        with self.conn.cursor() as cur:
            cur.execute(sql, (status, message, error_message, self.scheduler_key))
        self.conn.commit()

    def reset_cursor(self):
        sql = f"""
            UPDATE {self.table_name}
            SET
                cursor_realtor_id = NULL,
                cursor_stage = NULL,
                cursor_article_no = NULL,
                cursor_draft_id = NULL,
                cursor_publish_queue_id = NULL,
                last_pipeline_id = NULL,
                updated_at = NOW()
            WHERE scheduler_key = %s
        """
        with self.conn.cursor() as cur:
            cur.execute(sql, (self.scheduler_key,))
        self.conn.commit()

    def _to_int_or_none(self, value):
        try:
            if value is None or value == "":
                return None
            return int(value)
        except Exception:
            return None
