"""
Helpers to write screenshot events into the per-user manage file processing table
and keep the schema consistent with the backend Lambda.
"""

import json
import logging
import os
from typing import Any, Dict

from mysql.mysql_client import get_mysql_connection


def _table_name_for_user(user_id: str) -> str:
    """Return per-user table name: file_process_ + userId with dashes removed."""
    base = (user_id or '').replace('-', '')
    return f"file_process_{base}"


def ensure_user_table_exists(user_id: str) -> None:
    """Ensure the per-user manage table exists with the expected schema."""
    if not user_id:
        raise ValueError("user_id is required")
    table = _table_name_for_user(user_id)
    ddl = f"""
    CREATE TABLE IF NOT EXISTS {table} (
      id INT AUTO_INCREMENT PRIMARY KEY,
      flow_type VARCHAR(255),
      filename VARCHAR(255),
      description VARCHAR(255),
      type VARCHAR(50),
      step_number INT,
      processed TINYINT(1) DEFAULT 0,
      result LONGTEXT,
      matching_result LONGTEXT,
      prompt LONGTEXT,
      model VARCHAR(255),
      s3Path TEXT,
      error TEXT,
      fields LONGTEXT,
      options LONGTEXT,
      created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
      updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
      user_id VARCHAR(255)
    );
    """
    conn = get_mysql_connection()
    try:
        cur = conn.cursor()
        cur.execute(ddl)
        conn.commit()
    finally:
        try:
            cur.close()
        except Exception:
            pass
        try:
            conn.close()
        except Exception:
            pass


def add_or_update_file_record(
    *,
    user_id: str,
    filename: str,
    description: str = '',
    ftype: str = 'image',
    step_number: int = 1,
    s3_path: str | None = None,
    fields: Dict[str, Any] | None = None,
    options: Dict[str, Any] | None = None,
    model: str = 'gpt-4o',
    prompt: str = '',
    flow_type: str = 'analyze',
) -> Dict[str, Any]:
    """Upsert a file record for this user/step/flow (and filename when analyze).

    Returns a dict with {'updated': True} or {'inserted': True}.
    """
    if not user_id:
        raise ValueError("user_id is required")
    if not filename:
        raise ValueError("filename is required")
    table = _table_name_for_user(user_id)

    ensure_user_table_exists(user_id)

    fields = fields or {}
    options = options or {}
    processed = 2 if options.get('isExcelDocument') else 0

    conn = get_mysql_connection()
    try:
        cur = conn.cursor()
        # Check existing
        if flow_type == 'analyze':
            sel = f"SELECT id FROM {table} WHERE step_number = %s AND flow_type = %s AND filename = %s"
            cur.execute(sel, (step_number, flow_type, filename))
        else:
            sel = f"SELECT id FROM {table} WHERE step_number = %s AND flow_type = %s"
            cur.execute(sel, (step_number, flow_type))
        row = cur.fetchone()

        if row:
            # Update existing
            upd = (
                f"UPDATE {table} SET filename=%s, description=%s, type=%s, s3Path=%s, "
                f"fields=%s, options=%s, model=%s, prompt=%s, processed=%s, matching_result = NULL, result = NULL "
                f"WHERE step_number=%s AND flow_type=%s"
                + (" AND filename=%s" if flow_type == 'analyze' else '')
            )
            params = [
                filename,
                description,
                ftype,
                s3_path or '',
                json.dumps(fields),
                json.dumps(options),
                model,
                prompt,
                processed,
                step_number,
                flow_type,
            ]
            if flow_type == 'analyze':
                params.append(filename)
            cur.execute(upd, tuple(params))
            conn.commit()
            return {"updated": True}

        ins = (
            f"INSERT INTO {table} "
            f"(flow_type, filename, description, type, step_number, processed, prompt, s3Path, model, options, user_id, fields) "
            f"VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)"
        )
        cur.execute(
            ins,
            (
                flow_type,
                filename,
                description,
                ftype,
                step_number,
                processed,
                prompt,
                s3_path or '',
                model,
                json.dumps(options),
                user_id,
                json.dumps(fields),
            ),
        )
        conn.commit()
        return {"inserted": True}
    finally:
        try:
            cur.close()
        except Exception:
            pass
        try:
            conn.close()
        except Exception:
            pass


def add_or_append_screenshot_record(
    *,
    user_id: str,
    step_number: int,
    bucket: str,
    key: str,
    description: str = '',
    flow_type: str = 'analyze',
    di_step_details: Dict[str, Any] | None = None,
    extra_fields: Dict[str, Any] | None = None,
) -> Dict[str, Any]:
    """Append a screenshot key into a single per-step record with s3Path = {Bucket, keys: []}.

    We use a constant filename 'screenshots' so multiple screenshots for the same
    step aggregate into one row.
    """
    if not user_id:
        raise ValueError("user_id is required")
    table = _table_name_for_user(user_id)

    ensure_user_table_exists(user_id)

    # Map DI step details to manage table columns
    # - description <- step_details.description (when provided)
    # - fields      <- step_details.fields (raw structure)
    di_description = None
    di_fields = None
    if di_step_details and isinstance(di_step_details, dict):
        di_description = di_step_details.get('description')
        di_fields = di_step_details.get('fields')

    # Prefer DI mapping strictly for fields/description
    fields_payload = di_fields if di_fields is not None else (extra_fields or {})
    description = di_description if di_description is not None else description
    # Per requirement: do not add any extra fields; use DI step_details.fields only (or preserve existing)

    # Resolve effective flow_type from DI if present, else fallback to provided/default
    effective_flow = (
        (di_step_details.get('flowType') if isinstance(di_step_details, dict) else None)
        or (di_step_details.get('flow_type') if isinstance(di_step_details, dict) else None)
        or flow_type
        or 'analyze'
    )

    # Determine type from DI, defaulting to 'image'
    effective_type = (
        (di_step_details.get('type') if isinstance(di_step_details, dict) else None)
        or 'image'
    )

    conn = get_mysql_connection()
    try:
        cur = conn.cursor()
        # Look for existing row for this step/flow/constant filename
        filename = 'screenshots'
        sel = f"SELECT id, s3Path, fields FROM {table} WHERE step_number = %s AND flow_type = %s AND filename = %s"
        cur.execute(sel, (step_number, effective_flow, filename))
        row = cur.fetchone()

        s3_obj = {"Bucket": bucket, "keys": []}
        if row:
            # Merge existing s3Path->keys
            try:
                existing_s3 = row[1]
                if existing_s3:
                    parsed = json.loads(existing_s3) if isinstance(existing_s3, str) else existing_s3
                    if isinstance(parsed, dict):
                        s3_obj = parsed
                        if 'keys' not in s3_obj or not isinstance(s3_obj['keys'], list):
                            s3_obj['keys'] = []
                        # Keep latest bucket authoritative
                        s3_obj['Bucket'] = bucket or s3_obj.get('Bucket')
            except Exception:
                pass
            if key and key not in s3_obj['keys']:
                s3_obj['keys'].append(key)

            # Overwrite fields with DI mapping if provided (keep previous if not provided)
            if di_fields is None:
                try:
                    # Preserve existing fields when DI didn't provide fields this time
                    existing_fields = row[2]
                    if existing_fields:
                        fields_payload = existing_fields
                except Exception:
                    pass

            upd = (
                f"UPDATE {table} SET description=%s, type=%s, s3Path=%s, fields=%s, updated_at = CURRENT_TIMESTAMP "
                f"WHERE step_number=%s AND flow_type=%s AND filename=%s"
            )
            cur.execute(
                upd,
                (
                    description,
                    effective_type,
                    json.dumps(s3_obj),
                    json.dumps(fields_payload),
                    step_number,
                    effective_flow,
                    filename,
                ),
            )
            conn.commit()
            return {"updated": True, "appended": True}

        # No row yet – insert fresh aggregate row
        if key:
            s3_obj['keys'].append(key)
        ins = (
            f"INSERT INTO {table} "
            f"(flow_type, filename, description, type, step_number, processed, prompt, s3Path, model, options, user_id, fields) "
            f"VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)"
        )
        cur.execute(
            ins,
            (
                effective_flow,
                'screenshots',
                description,
                effective_type,
                step_number,
                0,
                '',
                json.dumps(s3_obj),
                'gpt-4o',
                json.dumps({"isExcelDocument": False}),
                user_id,
                json.dumps(fields_payload),
            ),
        )
        conn.commit()
        return {"inserted": True}
    finally:
        try:
            cur.close()
        except Exception:
            pass
        try:
            conn.close()
        except Exception:
            pass


def add_screenshot_record(
    *,
    user_id: str,
    step_number: int,
    bucket: str,
    key: str,
    description: str = '',
    di_step_details: Dict[str, Any] | None = None,
    flow_type_fallback: str = 'analyze',
) -> Dict[str, Any]:
    """Insert a NEW row for each screenshot (no aggregation).

    - flow_type: from DI step_details.flowType/flow_type else fallback
    - type: from DI step_details.type else 'image'
    - description: from DI step_details.description if provided
    - fields: DI step_details.fields strictly (no extras)
    - s3Path: { Bucket, keys: [key] }
    """
    if not user_id:
        raise ValueError("user_id is required")
    table = _table_name_for_user(user_id)
    ensure_user_table_exists(user_id)

    # Map from DI
    di_description = None
    di_fields = None
    effective_flow = flow_type_fallback or 'analyze'
    effective_type = 'image'
    if isinstance(di_step_details, dict):
        di_description = di_step_details.get('description')
        di_fields = di_step_details.get('fields')
        effective_flow = (
            di_step_details.get('flowType')
            or di_step_details.get('flow_type')
            or effective_flow
        )
        effective_type = di_step_details.get('type') or effective_type

    description = di_description if di_description is not None else description
    fields_payload = di_fields if di_fields is not None else {}

    s3_obj = {"Bucket": bucket, "keys": [key] if key else []}

    conn = get_mysql_connection()
    try:
        cur = conn.cursor()
        ins = (
            f"INSERT INTO {table} "
            f"(flow_type, filename, description, type, step_number, processed, prompt, s3Path, model, options, user_id, fields) "
            f"VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)"
        )
        filename = key.rsplit('/', 1)[-1] if key else 'screenshot.png'
        cur.execute(
            ins,
            (
                effective_flow,
                filename,
                description,
                effective_type,
                step_number,
                0,
                '',
                json.dumps(s3_obj),
                'gpt-4o',
                json.dumps({"isExcelDocument": False}),
                user_id,
                json.dumps(fields_payload),
            ),
        )
        conn.commit()
        return {"inserted": True}
    finally:
        try:
            cur.close()
        except Exception:
            pass
        try:
            conn.close()
        except Exception:
            pass


