"""
aiskills/workflow_skill.py - 工作流管理技能
============================================
为 AI Agent 提供工作流的创建、执行、查询、更新、删除等完整 CRUD 能力。

通过 action 参数区分操作类型：
  - list: 列出工作流
  - create: 创建工作流
  - get: 获取工作流详情
  - run: 执行工作流
  - update: 更新工作流
  - delete: 删除工作流
  - enable / disable: 启用/禁用工作流
  - runs: 查看运行记录

支持的步骤类型：
  - browser_open: 打开网页（参数: url, profile）
  - delay: 等待（参数: seconds）
  - shell: 执行命令（参数: command, timeout）
  - browser_extract: 提取页面内容（参数: selector, attribute）
  - browser_screenshot: 截图（参数: save_path）
  - api_call: 调用 HTTP API（参数: url, method, headers, body）
  - condition: 条件判断（参数: expression, on_true, on_false）
  - send_email: 发送邮件（占位实现）

模板变量：
  - {{steps.N.result}} / {{steps.N.error}}: 引用步骤结果
  - {{variables.xxx}}: 引用工作流变量
  - {{agent_path}} / {{now}}: 内置变量
"""
from __future__ import annotations

import json
import time
from typing import Any, Dict, List, Optional

from aiskills.base import Skill, SkillParameter, SkillResult
from core.logger import get_logger

logger = get_logger("myagent.skills.workflow")


class WorkflowSkill(Skill):
    """工作流管理 — 创建、执行、管理可重复运行的自动化工作流。

    类似 GitHub Actions，可以定义多步骤自动化任务（打开浏览器、执行命令、
    调用 API、提取数据等），支持模板变量、条件判断、运行记录。

    示例用法：
      创建工作流: action="create", name="每日签到", steps=[...]
      执行工作流: action="run", workflow_id="wf_xxx"
      列出工作流: action="list"
    """

    name = "workflow"
    description = (
        "工作流管理 — 创建、执行、管理可重复运行的自动化工作流。"
        "支持多步骤自动化（浏览器操作、Shell 命令、API 调用、数据提取等），"
        "类似 GitHub Actions。创建工作流后可反复执行，适合定时任务、批量操作、"
        "自动化测试等场景。"
    )
    category = "workflow"
    parameters = [
        SkillParameter(
            name="action",
            type="string",
            description=(
                "操作类型: list(列出), create(创建), get(详情), "
                "run(执行), update(更新), delete(删除), "
                "enable(启用), disable(禁用), runs(运行记录)"
            ),
            required=True,
        ),
        SkillParameter(
            name="agent_path",
            type="string",
            description="Agent 路径（用于工作流分组），默认 'default'",
            required=False,
            default="default",
        ),
        SkillParameter(
            name="workflow_id",
            type="string",
            description="工作流 ID（get/run/update/delete/enable/disable 时必填）",
            required=False,
            default="",
        ),
        SkillParameter(
            name="name",
            type="string",
            description="工作流名称（create/update 时使用）",
            required=False,
            default="",
        ),
        SkillParameter(
            name="description",
            type="string",
            description="工作流描述（create/update 时使用）",
            required=False,
            default="",
        ),
        SkillParameter(
            name="steps",
            type="string",
            description=(
                "步骤列表的 JSON 字符串（create/update 时使用）。"
                "每个步骤: {id, type, name, params}。"
                "type 可选: browser_open, delay, shell, browser_extract, "
                "browser_screenshot, api_call, condition, send_email"
            ),
            required=False,
            default="",
        ),
        SkillParameter(
            name="variables",
            type="string",
            description="工作流变量的 JSON 字符串（create/update 时使用），如 {\"url\": \"https://x.com\"}",
            required=False,
            default="",
        ),
        SkillParameter(
            name="run_variables",
            type="string",
            description="运行时覆盖变量的 JSON 字符串（run 时使用），如 {\"post_content\": \"Hello!\"}",
            required=False,
            default="",
        ),
        SkillParameter(
            name="enabled",
            type="boolean",
            description="启用/禁用（enable/disable 时使用）",
            required=False,
            default="true",
        ),
    ]

    # 步骤类型说明（供 LLM 参考）
    STEP_TYPES_HELP = {
        "browser_open": "打开网页。params: {url, profile(可选)}",
        "delay": "等待。params: {seconds(0.1-300)}",
        "shell": "执行 Shell 命令。params: {command, timeout(可选,1-600)}",
        "browser_extract": "提取页面内容。params: {selector, attribute(可选,默认text)}",
        "browser_screenshot": "截图。params: {save_path}",
        "api_call": "调用 HTTP API。params: {url, method(可选), headers(可选), body(可选)}",
        "condition": "条件判断。params: {expression, on_true(可选), on_false(可选)}",
        "send_email": "发送邮件（占位）。params: {to, subject, body}",
    }

    def _get_engine(self):
        from core.workflow_engine import WorkflowEngine
        return WorkflowEngine()

    def _parse_json(self, value: str, field_name: str) -> Any:
        """安全解析 JSON 字符串。"""
        if not value or not value.strip():
            return None
        try:
            return json.loads(value)
        except json.JSONDecodeError as e:
            raise ValueError(f"{field_name} JSON 格式错误: {e}") from e

    async def execute(self, action: str = "", **kwargs) -> SkillResult:
        action = (action or "").strip().lower()
        if not action:
            return SkillResult(success=False, error="缺少 action 参数")

        agent_path = kwargs.get("agent_path", "default") or "default"
        workflow_id = kwargs.get("workflow_id", "") or ""

        handler = {
            "list": self._action_list,
            "create": self._action_create,
            "get": self._action_get,
            "run": self._action_run,
            "update": self._action_update,
            "delete": self._action_delete,
            "enable": lambda ap, wid, kw: self._action_toggle(ap, wid, True),
            "disable": lambda ap, wid, kw: self._action_toggle(ap, wid, False),
            "runs": self._action_runs,
        }.get(action)

        if not handler:
            return SkillResult(
                success=False,
                error=f"未知 action: {action}。支持: list, create, get, run, update, delete, enable, disable, runs",
            )

        try:
            return await handler(agent_path, workflow_id, kwargs)
        except FileNotFoundError as e:
            return SkillResult(success=False, error=str(e))
        except ValueError as e:
            return SkillResult(success=False, error=str(e))
        except Exception as e:
            logger.error(f"工作流操作失败 ({action}): {e}", exc_info=True)
            return SkillResult(success=False, error=f"工作流操作失败: {e}")

    # ── Action handlers ──────────────────────────────────────────────

    async def _action_list(self, agent_path: str, _wid: str, _kw: Dict) -> SkillResult:
        engine = self._get_engine()
        workflows = engine.list_workflows(agent_path)
        if not workflows:
            return SkillResult(
                success=True,
                message=f"Agent [{agent_path}] 暂无工作流",
                output="[]",
            )

        # 构建摘要
        lines = [f"Agent [{agent_path}] 共 {len(workflows)} 个工作流:\n"]
        for wf in workflows:
            wfid = wf.get("id", "")
            wname = wf.get("name", "未命名")
            enabled = "✅" if wf.get("enabled", True) else "❌"
            steps_count = len(wf.get("steps", []))
            run_count = wf.get("run_count", 0)
            last_status = wf.get("last_status", "-")
            status_icon = {"completed": "✅", "failed": "❌", "running": "🔄", "stopped": "⏹"}.get(last_status, "➖")
            lines.append(
                f"  {enabled} [{wfid}] {wname} ({steps_count}步, 运行{run_count}次, 最近:{status_icon} {last_status})"
            )

        return SkillResult(
            success=True,
            message="\n".join(lines),
            output=json.dumps(workflows, ensure_ascii=False),
        )

    async def _action_create(self, agent_path: str, _wid: str, kw: Dict) -> SkillResult:
        name = kw.get("name", "").strip()
        if not name:
            return SkillResult(success=False, error="创建工作流需要 name 参数")

        description = kw.get("description", "").strip()

        # 解析步骤
        steps_raw = kw.get("steps", "")
        steps = self._parse_json(steps_raw, "steps") if steps_raw else []
        if steps is None:
            steps = []

        # 验证步骤格式
        validated_steps = self._validate_steps(steps)

        # 解析变量
        variables = self._parse_json(kw.get("variables", ""), "variables") or {}

        engine = self._get_engine()
        workflow = engine.create_workflow(agent_path, {
            "name": name,
            "description": description,
            "steps": validated_steps,
            "variables": variables,
            "enabled": True,
        })

        step_summary = "\n".join(
            f"    步骤 {s['id']}: [{s['type']}] {s.get('name', '')}"
            for s in validated_steps
        )
        output = json.dumps(workflow, ensure_ascii=False, indent=2)

        return SkillResult(
            success=True,
            message=(
                f"工作流已创建: [{workflow['id']}] {name}\n"
                f"  Agent: {agent_path}\n"
                f"  步骤数: {len(validated_steps)}\n"
                f"  变量: {json.dumps(variables, ensure_ascii=False) if variables else '无'}\n"
                f"步骤列表:\n{step_summary or '    (无步骤)'}\n"
                f"\n使用 action='run', workflow_id='{workflow['id']}' 执行此工作流。"
            ),
            output=output,
        )

    async def _action_get(self, agent_path: str, workflow_id: str, _kw: Dict) -> SkillResult:
        if not workflow_id:
            return SkillResult(success=False, error="获取工作流需要 workflow_id 参数")

        engine = self._get_engine()
        wf = engine.get_workflow(agent_path, workflow_id)
        if not wf:
            return SkillResult(success=False, error=f"工作流不存在: {workflow_id}")

        steps = wf.get("steps", [])
        step_lines = []
        for s in steps:
            params_str = json.dumps(s.get("params", {}), ensure_ascii=False)
            step_lines.append(
                f"  步骤 {s['id']}: [{s['type']}] {s.get('name', '')}\n    参数: {params_str}"
            )

        output = json.dumps(wf, ensure_ascii=False, indent=2)
        return SkillResult(
            success=True,
            message=(
                f"工作流详情: [{wf['id']}] {wf.get('name', '')}\n"
                f"  描述: {wf.get('description', '无')}\n"
                f"  启用: {'是' if wf.get('enabled') else '否'}\n"
                f"  运行次数: {wf.get('run_count', 0)}\n"
                f"  最近状态: {wf.get('last_status', '-')}\n"
                f"  变量: {json.dumps(wf.get('variables', {}), ensure_ascii=False)}\n"
                f"步骤列表:\n{chr(10).join(step_lines) or '    (无步骤)'}"
            ),
            output=output,
        )

    async def _action_run(self, agent_path: str, workflow_id: str, kw: Dict) -> SkillResult:
        if not workflow_id:
            return SkillResult(success=False, error="执行工作流需要 workflow_id 参数")

        # 解析运行时变量
        run_vars = self._parse_json(kw.get("run_variables", ""), "run_variables") or None

        engine = self._get_engine()

        # 先获取工作流信息（用于日志）
        wf = engine.get_workflow(agent_path, workflow_id)
        if not wf:
            return SkillResult(success=False, error=f"工作流不存在: {workflow_id}")

        wf_name = wf.get("name", "未命名")
        step_count = len(wf.get("steps", []))

        logger.info(f"开始执行工作流: [{workflow_id}] {wf_name} ({step_count}步)")
        run_record = await engine.run_workflow(agent_path, workflow_id, run_vars)

        # 构建结果摘要
        status = run_record.get("status", "unknown")
        steps_log = run_record.get("steps_log", [])
        started = run_record.get("started_at", "")
        finished = run_record.get("finished_at", "")

        status_icon = {"completed": "✅", "failed": "❌", "running": "🔄", "stopped": "⏹"}.get(status, "➖")

        step_lines = []
        for sl in steps_log:
            s_icon = {"success": "✅", "failed": "❌", "skipped": "⏭"}.get(sl.get("status", ""), "➖")
            s_name = sl.get("name", "") or sl.get("type", "")
            s_result = str(sl.get("result", ""))[:100] if sl.get("result") else ""
            s_error = str(sl.get("error", ""))[:100] if sl.get("error") else ""
            detail = s_result or s_error
            step_lines.append(f"  {s_icon} 步骤 {sl['step_id']}: {s_name} — {detail}")

        output = json.dumps(run_record, ensure_ascii=False, indent=2)

        return SkillResult(
            success=(status == "completed"),
            message=(
                f"工作流执行完毕: {status_icon} [{workflow_id}] {wf_name}\n"
                f"  状态: {status}\n"
                f"  时间: {started} → {finished}\n"
                f"  步骤结果:\n{chr(10).join(step_lines) or '    (无步骤)'}"
            ),
            output=output,
        )

    async def _action_update(self, agent_path: str, workflow_id: str, kw: Dict) -> SkillResult:
        if not workflow_id:
            return SkillResult(success=False, error="更新工作流需要 workflow_id 参数")

        engine = self._get_engine()
        wf = engine.get_workflow(agent_path, workflow_id)
        if not wf:
            return SkillResult(success=False, error=f"工作流不存在: {workflow_id}")

        update_data = {}
        for field in ("name", "description"):
            val = kw.get(field, "").strip()
            if val:
                update_data[field] = val

        # 解析步骤
        steps_raw = kw.get("steps", "")
        if steps_raw and steps_raw.strip():
            steps = self._parse_json(steps_raw, "steps")
            if steps is not None:
                update_data["steps"] = self._validate_steps(steps)

        # 解析变量
        variables_raw = kw.get("variables", "")
        if variables_raw and variables_raw.strip():
            variables = self._parse_json(variables_raw, "variables")
            if variables is not None:
                update_data["variables"] = variables

        if not update_data:
            return SkillResult(success=False, error="没有提供要更新的字段")

        updated = engine.update_workflow(agent_path, workflow_id, update_data)
        return SkillResult(
            success=True,
            message=f"工作流已更新: [{workflow_id}] {updated.get('name', '')}",
            output=json.dumps(updated, ensure_ascii=False, indent=2),
        )

    async def _action_delete(self, agent_path: str, workflow_id: str, _kw: Dict) -> SkillResult:
        if not workflow_id:
            return SkillResult(success=False, error="删除工作流需要 workflow_id 参数")

        engine = self._get_engine()
        ok = engine.delete_workflow(agent_path, workflow_id)
        if ok:
            return SkillResult(success=True, message=f"工作流已删除: {workflow_id}")
        return SkillResult(success=False, error=f"工作流不存在: {workflow_id}")

    async def _action_toggle(self, agent_path: str, workflow_id: str, enabled: bool) -> SkillResult:
        if not workflow_id:
            return SkillResult(success=False, error="需要 workflow_id 参数")

        engine = self._get_engine()
        updated = engine.toggle_workflow(agent_path, workflow_id, enabled)
        state = "启用" if enabled else "禁用"
        return SkillResult(
            success=True,
            message=f"工作流已{state}: [{workflow_id}] {updated.get('name', '')}",
            output=json.dumps(updated, ensure_ascii=False),
        )

    async def _action_runs(self, agent_path: str, workflow_id: str, _kw: Dict) -> SkillResult:
        engine = self._get_engine()
        runs = engine.list_runs(agent_path, workflow_id, limit=20)

        if not runs:
            return SkillResult(
                success=True,
                message=f"暂无运行记录 (agent={agent_path}, workflow={workflow_id or '全部'})",
                output="[]",
            )

        lines = [f"运行记录 ({len(runs)} 条):\n"]
        for r in runs:
            rid = r.get("run_id", "")
            wid = r.get("workflow_id", "")
            status = r.get("status", "?")
            started = r.get("started_at", "")
            finished = r.get("finished_at", "")
            steps_count = len(r.get("steps_log", []))
            status_icon = {"completed": "✅", "failed": "❌", "running": "🔄", "stopped": "⏹"}.get(status, "➖")
            lines.append(
                f"  {status_icon} [{rid}] workflow={wid} ({steps_count}步) {started}"
            )

        return SkillResult(
            success=True,
            message="\n".join(lines),
            output=json.dumps(runs, ensure_ascii=False),
        )

    # ── 辅助方法 ─────────────────────────────────────────────────────

    def _validate_steps(self, steps: List[Dict]) -> List[Dict]:
        """验证并规范化步骤列表。"""
        valid_types = set(self.STEP_TYPES_HELP.keys())
        validated = []

        for i, step in enumerate(steps):
            if not isinstance(step, dict):
                raise ValueError(f"步骤 #{i} 必须是对象")

            step_type = step.get("type", "")
            if step_type not in valid_types:
                raise ValueError(
                    f"步骤 #{i} 类型无效: '{step_type}'。"
                    f"可选类型: {', '.join(sorted(valid_types))}"
                )

            validated.append({
                "id": step.get("id", i),
                "type": step_type,
                "name": step.get("name", ""),
                "params": step.get("params", {}),
            })

        return validated
