import asyncio
import json
import os
import re
import shlex
import logging
import traceback
import secrets
import time
import signal
import shutil
import glob
from datetime import datetime, timedelta
from typing import Optional, Callable

try:
    from aiohttp import web as aiohttp_web
except ImportError:
    aiohttp_web = None
    logging.warning("aiohttp 未安装，Web Chat API 不可用。安装: pip install aiohttp")

try:
    import bcrypt
except ImportError:
    bcrypt = None
    logging.warning("bcrypt 未安装，MyData 认证不可用。安装: pip install bcrypt")

# === 程序启动时清除环境变量 ===
os.environ.pop("ANTHROPIC_AUTH_TOKEN", None)

# === 配置信息 ===
# -- 服务器基础路径 --
WEB_ROOT       = "/var/www/html/"                         # Web 根目录
DATA_DIR       = f"{WEB_ROOT}/data"                       # 中间文件/上传文件目录
SITE_DOMAIN    = "xybcloud.online"                        # 站点域名

# -- Bot 运行时文件 --
BOT_DIR        = os.path.dirname(os.path.abspath(__file__))  # Bot 数据目录（自动获取）
LOG_FILE       = f"{BOT_DIR}/bot.log"                     # 日志文件
SESSION_FILE   = f"{BOT_DIR}/sessions.json"               # 会话存储文件
HISTORY_FILE   = f"{BOT_DIR}/chat_history.json"           # 聊天历史记录文件
PROFILE_DB     = f"{BOT_DIR}/profile.db"                  # MyData 用户凭证文件

# -- Node.js / Claude CLI 路径 --
REMOTE_USER    = "ubuntu"                                  # 远程服务器用户名
REMOTE_HOME    = f"/home/{REMOTE_USER}"                    # 远程用户主目录
NODE_VERSION   = "v22.22.1"
NVM_BASE       = f"{REMOTE_HOME}/.nvm/versions/node"
NODE_BIN       = f"{NVM_BASE}/{NODE_VERSION}/bin/node"    # node 二进制路径
CLAUDE_BIN     = "/home/ubuntu/.local/bin/claude"  # claude CLI 路径
NODE_BIN_DIR   = f"{NVM_BASE}/{NODE_VERSION}/bin"         # node bin 目录（加入 PATH）

# -- 其他路径 --
BASHRC_PATH    = f"{REMOTE_HOME}/.bashrc"                  # 模型配置来源
COMMANDS_DIR   = f"{REMOTE_HOME}/.claude/commands"         # Skill 命令目录
CLAUDE_LOG_DIR = f"{REMOTE_HOME}/.claude/projects/-home-{REMOTE_USER}"  # Claude CLI 对话日志目录
CLAUDE_DEBUG_DIR = f"{REMOTE_HOME}/.claude/debug"          # Claude CLI debug 日志目录
CLAUDE_BACKUP_DIR = f"{REMOTE_HOME}/.claude/backups"       # Claude CLI 备份目录

# -- 运行参数 --
HTTP_PORT = 8765              # Web Chat API 端口
MESSAGE_IDLE_TIMEOUT = 120    # 消息间隔超时时间（秒）
MAX_RETRIES = 3               # 最大重试次数

# === 日志配置 ===
logging.basicConfig(
    level=logging.INFO,
    format='%(asctime)s - %(name)s - %(levelname)s - %(message)s',
    handlers=[
        logging.FileHandler(LOG_FILE),
        logging.StreamHandler()
    ]
)
logger = logging.getLogger(__name__)



def _normalize_text(text: str) -> str:
    """去除 emoji 和特殊字符，只保留中英文数字用于比较"""
    # 移除 emoji
    emoji_pattern = re.compile(
        "["
        "\U0001F600-\U0001F64F"  # emoticons
        "\U0001F300-\U0001F5FF"  # symbols & pictographs
        "\U0001F680-\U0001F6FF"  # transport & map symbols
        "\U0001F1E0-\U0001F1FF"  # flags
        "\U00002702-\U000027B0"
        "\U000024C2-\U0001F251"
        "\U0001f926-\U0001f937"
        "\U00010000-\U0010ffff"
        "\u2640-\u2642"
        "\u2600-\u2B55"
        "\u200d"
        "\u23cf"
        "\u23e9"
        "\u231a"
        "\ufe0f"
        "\u3030"
        "]+",
        flags=re.UNICODE
    )
    text = emoji_pattern.sub('', text)
    # 只保留中英文和数字
    text = re.sub(r'[^\u4e00-\u9fa5a-zA-Z0-9]', '', text)
    return text


def _deduplicate_progress(progress_texts: list, final_reply: str) -> list:
    """去除重复的进度消息，保留有意义的进展信息

    逻辑：如果某条进度消息的规范化文本是最终回复或后续进度的子串，则跳过
    """
    if not progress_texts:
        return []

    final_norm = _normalize_text(final_reply)
    result = []

    for i, text in enumerate(progress_texts):
        text_norm = _normalize_text(text)
        if not text_norm:
            continue
        # 如果是最终回复的子串，跳过（它只是最终回复的截断版本）
        if text_norm in final_norm:
            continue
        # 如果是后续某条进度消息的子串，跳过
        is_subset = False
        for j in range(i + 1, len(progress_texts)):
            later_norm = _normalize_text(progress_texts[j])
            if text_norm in later_norm:
                is_subset = True
                break
        if not is_subset:
            result.append(text)

    return result



def _get_model_tag(model_id: int) -> str:
    """根据 model_id 返回 [vendor+模型] 标签，未指定时使用默认模型"""
    def _tag(g):
        vendor = g.get('MODEL_VENDOR', '').replace('*', '')
        model = g.get('ANTHROPIC_MODEL', '')
        if vendor and model:
            return f"🔮 {vendor}+{model} "
        return f"🔮 {vendor or model} " if (vendor or model) else ""

    if model_id is not None and 0 <= model_id < len(MODEL_GROUPS):
        return _tag(MODEL_GROUPS[model_id])
    if MODEL_GROUPS and 0 <= DEFAULT_MODEL_ID < len(MODEL_GROUPS):
        return _tag(MODEL_GROUPS[DEFAULT_MODEL_ID])
    return ""


# === 模型组解析 ===
MODEL_GROUPS = []       # 所有模型组列表
DEFAULT_MODEL_ID = 0    # 默认（当前激活）模型组索引
_GLOBAL_MODEL_EXPORTS = {}  # 模型区域中不属于任何模型组的全局环境变量


def _parse_all_model_groups():
    """解析 ~/.bashrc 中从模型配置区域开始到文件结尾的所有环境变量

    模型区域通过以下方式检测（优先级从高到低）：
    1. 包含"模型"的注释行
    2. 首次出现 MODEL_VENDOR / ANTHROPIC_ / CLAUDE_CODE 相关 export 行

    区域内所有 export 变量都会被捕获，按 MODEL_VENDOR 分组，
    不属于任何模型组的 export 变量存入 _GLOBAL_MODEL_EXPORTS。
    """
    global MODEL_GROUPS, DEFAULT_MODEL_ID, _GLOBAL_MODEL_EXPORTS
    bashrc_path = BASHRC_PATH
    groups = []
    current_group = {}
    global_exports = {}

    GROUP_KEY = 'MODEL_VENDOR'
    GROUP_FIELDS = ['MODEL_VENDOR', 'ANTHROPIC_API_KEY', 'ANTHROPIC_BASE_URL', 'ANTHROPIC_MODEL']

    try:
        with open(bashrc_path, 'r', encoding='utf-8') as f:
            lines = f.readlines()

        # === 第一遍：找到模型配置区域的起始行 ===
        section_start = len(lines)  # 默认：未找到
        for i, line in enumerate(lines):
            stripped = line.strip()
            # 优先：注释行中包含"模型"
            if stripped.startswith('#') and '模型' in stripped:
                section_start = i
                break
        if section_start == len(lines):
            # 回退：首次出现 MODEL_VENDOR / ANTHROPIC_ / CLAUDE_CODE 的 export 行
            for i, line in enumerate(lines):
                stripped = line.strip()
                m = re.match(r'(?:#?\s*)?export\s+(\w+)', stripped)
                if m:
                    key = m.group(1)
                    if key.startswith(('MODEL_VENDOR', 'ANTHROPIC_', 'CLAUDE_CODE')):
                        section_start = i
                        break

        logger.info(f"模型配置区域起始行: {section_start + 1}")

        # === 第二遍：从 section_start 到文件末尾，解析所有 export ===
        for line in lines[section_start:]:
            line = line.strip()
            if not line:
                continue

            # 尝试匹配 export KEY=VALUE（注释或未注释均可）
            m = re.match(r'(#?\s*)?export\s+(\w+)\s*=\s*["\']?([^"\'#\s]+)["\']?', line)
            if not m:
                # 非 export 行且非注释行 → 如果已有组则保存
                if not line.startswith('#'):
                    if current_group.get(GROUP_KEY):
                        groups.append(current_group)
                        current_group = {}
                continue

            key = m.group(2).strip()
            value = m.group(3).strip()

            if key in GROUP_FIELDS:
                # MODEL_VENDOR 开始新组
                if key == GROUP_KEY and current_group.get(GROUP_KEY):
                    groups.append(current_group)
                    current_group = {}
                if key == GROUP_KEY:
                    # * 字符标记默认模型，保留在名称中
                    if '*' in value:
                        current_group['_active'] = True
                    else:
                        current_group['_active'] = False
                current_group[key] = value
            elif current_group.get(GROUP_KEY):
                # 当前模型组内的其他变量
                current_group.setdefault('_extras', {})[key] = value
            else:
                # 不属于任何模型组的全局变量
                global_exports[key] = value

        # 最后一个组
        if current_group.get(GROUP_KEY):
            groups.append(current_group)

        _GLOBAL_MODEL_EXPORTS = global_exports
        MODEL_GROUPS = groups

        # 找到默认（激活的）组
        DEFAULT_MODEL_ID = 0
        for i, g in enumerate(MODEL_GROUPS):
            if g.get('_active', False):
                DEFAULT_MODEL_ID = i
                break

        logger.info(f"已加载 {len(MODEL_GROUPS)} 个模型组，默认: {MODEL_GROUPS[DEFAULT_MODEL_ID]['MODEL_VENDOR'].replace('*', '') if MODEL_GROUPS else 'none'}")
        if _GLOBAL_MODEL_EXPORTS:
            logger.info(f"全局模型变量: {list(_GLOBAL_MODEL_EXPORTS.keys())}")
        for g in MODEL_GROUPS:
            extras = list(g.get('_extras', {}).keys())
            if extras:
                logger.info(f"模型组 {g.get('MODEL_VENDOR', '?').replace('*', '')} 额外变量: {extras}")

    except Exception as e:
        logger.error(f"解析模型组失败: {e}")


_parse_all_model_groups()





# === 会话锁管理 ===
class SessionLockManager:
    """管理每个会话的任务锁，实现串行处理"""

    _locks: dict = {}  # 会话锁
    _busy: dict = {}   # 会话忙碌状态
    _queues: dict = {} # 会话等待队列计数

    @classmethod
    def get_lock(cls, session_id: str) -> asyncio.Lock:
        """获取指定会话的锁"""
        if session_id not in cls._locks:
            cls._locks[session_id] = asyncio.Lock()
            cls._busy[session_id] = False
            cls._queues[session_id] = 0
        return cls._locks[session_id]

    @classmethod
    def is_busy(cls, session_id: str) -> bool:
        """检查会话是否正在处理任务"""
        return cls._busy.get(session_id, False)

    @classmethod
    def set_busy(cls, session_id: str, busy: bool):
        """设置会话忙碌状态"""
        cls._busy[session_id] = busy

    @classmethod
    def get_queue_count(cls, session_id: str) -> int:
        """获取等待队列中的任务数"""
        return cls._queues.get(session_id, 0)

    @classmethod
    def increment_queue(cls, session_id: str):
        """增加等待队列计数"""
        if session_id not in cls._queues:
            cls._queues[session_id] = 0
        cls._queues[session_id] += 1

    @classmethod
    def decrement_queue(cls, session_id: str):
        """减少等待队列计数"""
        if session_id in cls._queues and cls._queues[session_id] > 0:
            cls._queues[session_id] -= 1


# === 会话管理 ===
class SessionManager:
    """管理 Cutebot 会话 ID"""

    @staticmethod
    def _load_sessions() -> dict:
        """从文件加载会话数据"""
        try:
            if os.path.exists(SESSION_FILE):
                with open(SESSION_FILE, 'r') as f:
                    return json.load(f)
        except Exception as e:
            logger.warning(f"加载会话文件失败: {e}")
        return {}

    @staticmethod
    def _save_sessions(sessions: dict):
        """保存会话数据到文件"""
        try:
            os.makedirs(os.path.dirname(SESSION_FILE), exist_ok=True)
            with open(SESSION_FILE, 'w') as f:
                json.dump(sessions, f, indent=2)
        except Exception as e:
            logger.error(f"保存会话文件失败: {e}")

    @staticmethod
    def get_session_id(identifier: str) -> Optional[str]:
        """获取已保存的会话 ID"""
        sessions = SessionManager._load_sessions()
        return sessions.get(identifier)

    @staticmethod
    def save_session_id(identifier: str, session_id: str):
        """保存 Cutebot 返回的 session ID"""
        sessions = SessionManager._load_sessions()
        sessions[identifier] = session_id
        SessionManager._save_sessions(sessions)
        logger.info(f"保存会话 ID {identifier}: {session_id}")

    @staticmethod
    def reset_session(identifier: str):
        """重置会话 ID"""
        sessions = SessionManager._load_sessions()
        if identifier in sessions:
            del sessions[identifier]
            SessionManager._save_sessions(sessions)
            logger.info(f"重置会话: {identifier}")

    @staticmethod
    def should_reset_session(content: str) -> bool:
        """检查是否应该重置会话"""
        content_lower = content.lower()
        return "新对话" in content or "重新开始" in content or "新会话" in content or bool(re.search(r'/new\b', content_lower))

    @staticmethod
    def extract_session_id(content: str) -> tuple:
        """从消息中提取 session_id

        Returns:
            (session_id, cleaned_content): 提取的 session_id 和清理后的消息内容
        """
        # 匹配 session_id = xxx 或 session_id=xxx 格式
        pattern = r'session_id\s*=\s*([a-f0-9]{8}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{12})'
        match = re.search(pattern, content, re.IGNORECASE)
        if match:
            session_id = match.group(1)
            # 移除 session_id 部分，保留其他内容
            cleaned_content = re.sub(pattern, '', content, flags=re.IGNORECASE).strip()
            return session_id, cleaned_content
        return None, content



# === 聊天历史管理 ===
class ChatHistoryManager:
    """管理聊天历史记录

    使用内存缓存避免每次操作都读写磁盘：
    - 首次访问从磁盘加载到内存，后续操作直接读写缓存
    - add_message / delete 等关键操作立即写盘（原子写入）
    - update_message 进度更新仅改内存，最终回复时才写盘
    - flush() 可手动将缓存写盘
    """

    _cache: dict = None   # 内存缓存，None 表示尚未加载
    _dirty: bool = False  # 缓存已修改但未写盘

    @classmethod
    def _load_history(cls) -> dict:
        """获取历史记录（优先使用内存缓存）"""
        if cls._cache is not None:
            return cls._cache
        try:
            if os.path.exists(HISTORY_FILE):
                with open(HISTORY_FILE, 'r', encoding='utf-8') as f:
                    cls._cache = json.load(f)
                    return cls._cache
        except Exception as e:
            logger.warning(f"加载历史记录失败: {e}")
        cls._cache = {}
        return cls._cache

    @classmethod
    def _flush_to_disk(cls):
        """原子写入：先写临时文件，再 replace，防止崩溃时数据损坏"""
        if cls._cache is None:
            return
        try:
            os.makedirs(os.path.dirname(HISTORY_FILE), exist_ok=True)
            tmp = HISTORY_FILE + '.tmp'
            with open(tmp, 'w', encoding='utf-8') as f:
                json.dump(cls._cache, f, ensure_ascii=False, indent=2)
            os.replace(tmp, HISTORY_FILE)
            cls._dirty = False
        except Exception as e:
            logger.error(f"保存历史记录失败: {e}")

    @classmethod
    def flush(cls):
        """将未写盘的缓存刷入磁盘（仅在有脏数据时）"""
        if cls._dirty and cls._cache is not None:
            cls._flush_to_disk()

    @classmethod
    def _ensure_entry(cls, history: dict, session_key: str) -> dict:
        """确保 session_key 对应的条目存在且格式正确"""
        if session_key not in history:
            history[session_key] = {"title": "", "messages": []}
        entry = history[session_key]
        # 兼容旧格式：如果存的是 list，自动迁移
        if isinstance(entry, list):
            history[session_key] = {"title": "", "messages": entry}
        return history[session_key]

    @classmethod
    def add_message(cls, session_key: str, role: str, content: str, source: str = "web"):
        """添加一条消息到历史记录（立即写盘）"""
        history = cls._load_history()
        entry = cls._ensure_entry(history, session_key)

        message = {
            "id": secrets.token_hex(8),
            "role": role,
            "content": content,
            "source": source,
            "timestamp": time.time(),
            "time_str": time.strftime("%Y-%m-%d %H:%M:%S")
        }
        entry["messages"].append(message)

        # 限制每个会话最多保存500条消息
        if len(entry["messages"]) > 500:
            entry["messages"] = entry["messages"][-500:]

        cls._flush_to_disk()
        return message["id"]

    @classmethod
    def update_message(cls, session_key: str, message_id: str, content: str, update_timestamp: bool = False, usage: dict = None):
        """更新已有消息的内容

        进度更新（update_timestamp=False）仅修改内存缓存，不写盘；
        最终回复（update_timestamp=True）写盘。
        """
        history = cls._load_history()
        if session_key not in history:
            return False
        entry = cls._ensure_entry(history, session_key)

        for msg in entry["messages"]:
            if msg.get("id") == message_id:
                msg["content"] = content
                msg["updated_at"] = time.time()
                if update_timestamp:
                    msg["timestamp"] = time.time()
                    msg["time_str"] = time.strftime("%Y-%m-%d %H:%M:%S")
                if usage:
                    msg["usage"] = usage
                if update_timestamp:
                    cls._flush_to_disk()
                else:
                    cls._dirty = True
                return True
        return False

    @classmethod
    def get_messages(cls, session_key: str, limit: int = 100, offset: int = 0) -> list:
        """获取历史消息"""
        history = cls._load_history()
        entry = history.get(session_key, {})
        messages = entry.get("messages", []) if isinstance(entry, dict) else entry
        messages = sorted(messages, key=lambda x: x.get("timestamp", 0), reverse=True)
        return messages[offset:offset + limit]

    @classmethod
    def get_all_messages(cls, limit: int = 200, offset: int = 0) -> list:
        """获取所有来源的消息（合并显示）"""
        history = cls._load_history()
        all_messages = []
        for session_key, entry in history.items():
            msgs = entry.get("messages", []) if isinstance(entry, dict) else entry
            for msg in msgs:
                msg_copy = msg.copy()
                msg_copy["session_key"] = session_key
                all_messages.append(msg_copy)
        all_messages.sort(key=lambda x: x.get("timestamp", 0), reverse=True)
        return all_messages[offset:offset + limit]

    @classmethod
    def delete_message(cls, session_key: str, message_id: str) -> bool:
        """删除单条消息（立即写盘）"""
        history = cls._load_history()
        if session_key not in history:
            return False
        entry = cls._ensure_entry(history, session_key)
        original_len = len(entry["messages"])
        entry["messages"] = [m for m in entry["messages"] if m.get("id") != message_id]
        if len(entry["messages"]) < original_len:
            cls._flush_to_disk()
            return True
        return False

    @classmethod
    def delete_messages(cls, session_key: str, message_ids: list) -> int:
        """批量删除消息（立即写盘）"""
        history = cls._load_history()
        if session_key not in history:
            return 0
        entry = cls._ensure_entry(history, session_key)
        original_len = len(entry["messages"])
        ids_set = set(message_ids)
        entry["messages"] = [m for m in entry["messages"] if m.get("id") not in ids_set]
        deleted = original_len - len(entry["messages"])
        if deleted > 0:
            cls._flush_to_disk()
        return deleted

    @classmethod
    def delete_all(cls, session_key: str = None) -> int:
        """删除所有消息或指定会话的消息（立即写盘）"""
        history = cls._load_history()
        if session_key:
            if session_key in history:
                entry = history[session_key]
                count = len(entry.get("messages", [])) if isinstance(entry, dict) else len(entry)
                del history[session_key]
                cls._flush_to_disk()
                return count
            return 0
        else:
            count = sum(len(e.get("messages", []) if isinstance(e, dict) else e) for e in history.values())
            cls._cache = {}
            cls._flush_to_disk()
            return count

    @classmethod
    def delete_by_ids_global(cls, message_ids: list) -> int:
        """全局按ID删除消息（立即写盘）"""
        history = cls._load_history()
        ids_set = set(message_ids)
        total_deleted = 0
        for session_key in history:
            entry = cls._ensure_entry(history, session_key)
            original_len = len(entry["messages"])
            entry["messages"] = [m for m in entry["messages"] if m.get("id") not in ids_set]
            total_deleted += original_len - len(entry["messages"])
        if total_deleted > 0:
            cls._flush_to_disk()
        return total_deleted

    @classmethod
    def rename_session(cls, old_key: str, new_key: str):
        """将一个会话键重命名为新键（立即写盘），保留 title"""
        history = cls._load_history()
        if old_key not in history:
            return
        old_entry = history.pop(old_key)
        old_msgs = old_entry.get("messages", []) if isinstance(old_entry, dict) else old_entry
        old_title = old_entry.get("title", "") if isinstance(old_entry, dict) else ""

        if new_key in history:
            new_entry = cls._ensure_entry(history, new_key)
            new_entry["messages"].extend(old_msgs)
        else:
            history[new_key] = {"title": old_title, "messages": old_msgs}
        cls._flush_to_disk()
        logger.info(f"会话键重命名: {old_key} -> {new_key}")

    @classmethod
    def set_title(cls, session_key: str, title: str):
        """设置对话标题（立即写盘）"""
        history = cls._load_history()
        entry = cls._ensure_entry(history, session_key)
        entry["title"] = title
        cls._flush_to_disk()

    @classmethod
    def get_title(cls, session_key: str) -> str:
        """获取对话标题"""
        history = cls._load_history()
        entry = history.get(session_key, {})
        return entry.get("title", "") if isinstance(entry, dict) else ""


async def ask_claude(question: str, session_id: Optional[str], on_progress: Optional[Callable] = None, task_key: Optional[str] = None, model_id: Optional[int] = None) -> tuple:
    """调用 Cutebot CLI，解析流式 JSON 输出，返回回复文本、session ID 和 token 使用统计

    使用消息间隔超时而非总时间超时，允许长任务运行

    Args:
        question: 用户问题（可能包含 session_id=xxx 格式）
        session_id: 会话 ID（如果为 None 则创建新会话）
        on_progress: 可选回调，收到中间 assistant 文本时调用 await on_progress(text)
        task_key: 任务标识符，用于支持任务取消（仅 Web 端使用）
        model_id: 模型组索引（None 则使用默认模型）

    Returns:
        (result_text, returned_session_id, usage_info)
        usage_info 格式: {
            'input_tokens': int,
            'output_tokens': int,
            'cache_creation_input_tokens': int,
            'cache_read_input_tokens': int
        }
    """
    # 从问题中提取用户指定的 session_id（如果有）
    user_specified_session_id, cleaned_question = SessionManager.extract_session_id(question)
    if user_specified_session_id:
        logger.info(f"用户指定 session_id: {user_specified_session_id}")
        # 优先使用用户指定的 session_id
        session_id = user_specified_session_id
        question = cleaned_question
        # 发送恢复提示消息
        if on_progress:
            try:
                await on_progress(f"🔄 正在恢复会话：session_id = {user_specified_session_id}")
            except Exception as e:
                logger.error(f"发送恢复提示失败: {e}")

    logger.info(f"正在调用 Cutebot，问题: {question[:100]}")
    task_start_time = asyncio.get_event_loop().time()

    for retry in range(MAX_RETRIES):
        try:
            # 构建 Cutebot 命令
            # v2.1.63+ 版本需要设置 API key 环境变量，同时清除嵌套检测变量

            # 构建 claude 命令字符串（使用绝对路径 node 调用，避免 shebang /usr/bin/env node 找不到 node 的问题）
            claude_cmd = f"{CLAUDE_BIN} --allow-dangerously-skip-permissions --permission-mode bypassPermissions"
            claude_cmd += " -p --output-format stream-json --verbose"

            # 如果有 session_id，使用 --resume 恢复之前的对话
            if session_id:
                claude_cmd += f" --resume {session_id}"
                logger.info(f"[DEBUG] 恢复会话，传入 session_id: {session_id}")
            else:
                logger.info("[DEBUG] 创建新会话，未传入 session_id")

            # 使用 -- 分隔符将问题作为位置参数传递，
            # 避免以 - 开头的内容被 CLI 误解析为命令行选项
            claude_cmd += f" -- {shlex.quote(question)}"

            # 根据选择的模型组构建环境变量
            if model_id is not None and 0 <= model_id < len(MODEL_GROUPS):
                group = MODEL_GROUPS[model_id]
            elif MODEL_GROUPS:
                group = MODEL_GROUPS[DEFAULT_MODEL_ID]
            else:
                group = {}

            env_exports = []
            # 1. 先加载全局模型变量（不属于任何模型组的）
            for key, value in _GLOBAL_MODEL_EXPORTS.items():
                env_exports.append(f'export {key}="{value}"')
            # 2. 加载模型组的核心字段（会覆盖同名全局变量）
            for key in ['ANTHROPIC_API_KEY', 'ANTHROPIC_BASE_URL', 'ANTHROPIC_MODEL']:
                if key in group:
                    env_exports.append(f'export {key}="{group[key]}"')
            # 3. 加载模型组的额外变量
            for key, value in group.get('_extras', {}).items():
                env_exports.append(f'export {key}="{value}"')

            env_str = " && ".join(env_exports) if env_exports else ""
            unset_str = "unset CLAUDE_CODE CLAUDECODE CLAUDE_CODE_ENTRY_POINT CLAUDE_CODE_ENTRYPOINT CLAUDE_CODE_DISABLE_NONESSENTIAL_TRAFFIC"

            node_path = f"export PATH={NODE_BIN_DIR}:$PATH"
            if env_str:
                bash_cmd = f"cd {REMOTE_HOME} && {node_path} && {env_str} && {unset_str} && {claude_cmd}"
            else:
                bash_cmd = f'cd {REMOTE_HOME} && {node_path} && eval "$(grep \'^export ANTHROPIC\' {BASHRC_PATH} | grep -v \'^#\')" && {unset_str} && {claude_cmd}'

            logger.info(f"使用模型: {group.get('MODEL_VENDOR', 'default').replace('*', '')}")
            all_extras = {**_GLOBAL_MODEL_EXPORTS, **group.get('_extras', {})}
            logger.info(f"环境变量: API_KEY={'***' if group.get('ANTHROPIC_API_KEY') else 'None'}, BASE_URL={group.get('ANTHROPIC_BASE_URL', 'None')}, MODEL={group.get('ANTHROPIC_MODEL', 'None')}, 额外变量={list(all_extras.keys())}")
            logger.debug(f"执行命令: {bash_cmd[:200]}...")

            cmd_args = [
                "bash", "-c", bash_cmd
            ]

            # 创建子进程（使用进程组以便能杀掉所有子进程）
            proc = await asyncio.create_subprocess_exec(
                *cmd_args,
                stdout=asyncio.subprocess.PIPE,
                stderr=asyncio.subprocess.PIPE,
                limit=100 * 1024 * 1024,
                preexec_fn=os.setsid,  # 创建新的进程组
            )

            # 注册任务到 _active_tasks（仅 Web 端使用）
            if task_key:
                _active_tasks[task_key] = {'proc': proc, 'cancelled': False}

            result_text = ""
            last_sent_text = ""
            last_message_time = asyncio.get_event_loop().time()
            returned_session_id = session_id
            session_id_sent = [False]  # 使用列表以便在内部函数中修改
            usage_info = {
                'input_tokens': 0,
                'output_tokens': 0,
                'cache_creation_input_tokens': 0,
                'cache_read_input_tokens': 0,
                'api_calls': 0
            }

            try:
                async for line in proc.stdout:
                    # 检查任务是否被取消
                    if task_key and task_key in _active_tasks and _active_tasks[task_key].get('cancelled'):
                        logger.info(f"任务 {task_key} 被用户取消，终止 Cutebot 进程")
                        try:
                            # 杀掉整个进程组（包括所有子进程）
                            if proc.returncode is None:
                                try:
                                    os.killpg(os.getpgid(proc.pid), signal.SIGTERM)
                                    await asyncio.sleep(0.5)
                                    # 如果还没退出，强制杀掉
                                    if proc.returncode is None:
                                        os.killpg(os.getpgid(proc.pid), signal.SIGKILL)
                                except ProcessLookupError:
                                    # 进程已经不存在了
                                    pass
                        except Exception as e:
                            logger.warning(f"终止进程失败: {e}")
                        # 清理任务注册
                        if task_key in _active_tasks:
                            del _active_tasks[task_key]
                        return "⏹️ 任务已被取消", returned_session_id, usage_info

                    try:
                        # 更新最后收到消息的时间
                        current_time = asyncio.get_event_loop().time()
                        time_since_last_msg = current_time - last_message_time
                        
                        # 如果距离上次消息超过 idle timeout，检查进程状态
                        if time_since_last_msg > MESSAGE_IDLE_TIMEOUT:
                            logger.warning(f"Cutebot 已 {time_since_last_msg:.0f}秒未输出消息，检查进程")
                        
                        last_message_time = current_time
                        
                        decoded = line.decode("utf-8").strip()
                        if not decoded:
                            continue
                        logger.debug(f"Cutebot stdout: {decoded[:200]}")
                        
                        data = json.loads(decoded)
                    except json.JSONDecodeError as e:
                        logger.warning(f"JSON 解析错误: {e}")
                        continue
                    except UnicodeDecodeError as e:
                        logger.warning(f"编码错误: {e}")
                        continue

                    # 从 system 类型的消息中提取 session_id（只在新会话时更新）
                    if data.get("type") == "system" and data.get("session_id"):
                        cli_session_id = data["session_id"]
                        logger.info(f"[DEBUG] Claude CLI 返回 session_id (system): {cli_session_id}, 传入的 session_id: {session_id}")
                        if not session_id:  # 只在新会话时更新
                            returned_session_id = cli_session_id
                            logger.info(f"从 system 消息获取 session_id: {returned_session_id}")
                        else:
                            logger.warning(f"[DEBUG] session_id 不匹配！传入: {session_id}, CLI 返回: {cli_session_id}")
                        # 首次获取到 session_id 时，通过回调通知用户
                        if not session_id_sent[0] and on_progress:
                            try:
                                await on_progress(f"📋 已创建任务：session_id = {returned_session_id}")
                                session_id_sent[0] = True
                            except Exception as e:
                                logger.error(f"发送 session_id 通知失败: {e}")

                    # 提取 assistant 文本及 tool_use，实时回调
                    if data.get("type") == "assistant":
                        usage_info['api_calls'] += 1
                        msg = data.get("message", {})
                        for block in msg.get("content", []):
                            if block.get("type") == "text" and block.get("text"):
                                text = block["text"]
                                # 只发送新增的、有意义的文本
                                if text != last_sent_text and on_progress:
                                    try:
                                        await on_progress(f"🔄 {text[:500]}")
                                        last_sent_text = text
                                    except Exception as e:
                                        logger.error(f"发送中间消息失败: {e}")
                                result_text = text
                            elif block.get("type") == "tool_use" and on_progress:
                                tool_name = block.get("name", "")
                                tool_input = block.get("input", {})
                                # 提取命令前几十个字符作为进度提示
                                cmd = tool_input.get("command", "") or tool_input.get("code", "")
                                cmd_preview = cmd.strip()[:80].replace("\n", " ") if cmd else tool_name
                                try:
                                    await on_progress(f"⚙️ 正在执行：{cmd_preview}...")
                                except Exception as e:
                                    logger.error(f"发送 tool_use 进度失败: {e}")

                    # 最终 result 行包含完整回复
                    if data.get("type") == "result" and data.get("result"):
                        result_text = data["result"]
                        # result 中也有 session_id（只在新会话时更新）
                        if data.get("session_id"):
                            cli_session_id = data["session_id"]
                            logger.info(f"[DEBUG] Claude CLI 返回 session_id (result): {cli_session_id}, 传入的 session_id: {session_id}")
                            if not session_id:
                                returned_session_id = cli_session_id
                            else:
                                logger.warning(f"[DEBUG] result 中 session_id 不匹配！传入: {session_id}, CLI 返回: {cli_session_id}")
                        # 提取 token 使用统计
                        if data.get("usage"):
                            usage = data["usage"]
                            usage_info['input_tokens'] = usage.get('input_tokens', 0)
                            usage_info['output_tokens'] = usage.get('output_tokens', 0)
                            usage_info['cache_creation_input_tokens'] = usage.get('cache_creation_input_tokens', 0)
                            usage_info['cache_read_input_tokens'] = usage.get('cache_read_input_tokens', 0)
                            logger.info(f"Token 使用: 调用={usage_info['api_calls']}次, 输入={usage_info['input_tokens']}, 输出={usage_info['output_tokens']}, "
                                      f"缓存创建={usage_info['cache_creation_input_tokens']}, 缓存读取={usage_info['cache_read_input_tokens']}")

                # 异步等待进程退出，设置总的等待超时
                try:
                    await asyncio.wait_for(proc.wait(), timeout=10)
                except asyncio.TimeoutError:
                    logger.warning("等待 Cutebot 进程退出超时")
                    try:
                        proc.kill()
                    except:
                        pass

            except Exception as e:
                logger.error(f"Cutebot 输出处理异常: {e}\n{traceback.format_exc()}")
                try:
                    proc.kill()
                except:
                    pass

            # 读取错误输出
            try:
                stderr_data = await asyncio.wait_for(proc.stderr.read(), timeout=5)
                if stderr_data:
                    stderr_str = stderr_data.decode('utf-8', errors='ignore')[:500]
                    logger.warning(f"Cutebot stderr: {stderr_str}")
            except asyncio.TimeoutError:
                logger.warning("读取 stderr 超时")

            if proc.returncode != 0:
                logger.warning(f"Cutebot 进程非正常退出，码: {proc.returncode}")

            # 清理任务注册
            if task_key and task_key in _active_tasks:
                del _active_tasks[task_key]

            # 计算总用时并附加到结果文本
            elapsed = asyncio.get_event_loop().time() - task_start_time
            if elapsed >= 60:
                elapsed_str = f"\n\n⏱️ 总用时: {int(elapsed // 60)}分{int(elapsed % 60)}秒"
            else:
                elapsed_str = f"\n\n⏱️ 总用时: {int(elapsed)}秒"

            # 添加 session 完成信息（前面加空行）
            mtag = _get_model_tag(model_id)
            session_info = f"\n\n✅ session_id = {returned_session_id} 已完成\n{mtag.strip()}" if returned_session_id else ""
            final_text = (result_text or "⚠️ Cutebot 暂时无法回复，请稍后再试") + session_info + elapsed_str

            return final_text, returned_session_id, usage_info

        except FileNotFoundError:
            # 清理任务注册
            if task_key and task_key in _active_tasks:
                del _active_tasks[task_key]
            logger.error(f"Cutebot CLI 未找到，请检查 {CLAUDE_BIN} 是否存在")
            return "⚠️ Cutebot 服务不可用", session_id, {}
        except Exception as e:
            logger.error(f"Cutebot 调用异常 (尝试 {retry + 1}/{MAX_RETRIES}): {e}\n{traceback.format_exc()}")
            if retry < MAX_RETRIES - 1:
                await asyncio.sleep(2)
                continue
            # 清理任务注册
            if task_key and task_key in _active_tasks:
                del _active_tasks[task_key]
            return f"⚠️ 调用失败: {str(e)[:100]}", session_id, {}

    # 清理任务注册
    if task_key and task_key in _active_tasks:
        del _active_tasks[task_key]
    return "⚠️ Cutebot 多次重试失败，请检查服务状态", session_id, {}


# === Web Chat HTTP API ===
_web_tokens = {}  # {token: {'session_key': str, 'created': float}}
_active_tasks = {}  # {task_key: {'proc': Process, 'cancelled': bool}}  任务取消管理


def _read_mydata_profile():
    """读取 MyData profile.db 文件，返回 (username, password_hash)"""
    try:
        if not os.path.exists(PROFILE_DB):
            logger.warning(f"profile.db 不存在: {PROFILE_DB}")
            return None, None

        with open(PROFILE_DB, 'r', encoding='utf-8') as f:
            lines = f.read().strip().split('\n')
            if len(lines) < 2:
                logger.warning("profile.db 格式错误：行数不足")
                return None, None

            username = lines[0].strip()
            password_hash = lines[1].strip()

            if not username or not password_hash:
                logger.warning("profile.db 格式错误：用户名或密码哈希为空")
                return None, None

            return username, password_hash
    except Exception as e:
        logger.error(f"读取 profile.db 失败: {e}")
        return None, None


def _verify_mydata_password(password: str, password_hash: str) -> bool:
    """验证密码是否匹配 bcrypt 哈希"""
    if bcrypt is None:
        logger.error("bcrypt 未安装，无法验证密码")
        return False

    try:
        return bcrypt.checkpw(password.encode('utf-8'), password_hash.encode('utf-8'))
    except Exception as e:
        logger.error(f"密码验证失败: {e}")
        return False


def _cors_headers():
    """CORS 响应头"""
    return {
        'Access-Control-Allow-Origin': '*',
        'Access-Control-Allow-Methods': 'GET, POST, OPTIONS',
        'Access-Control-Allow-Headers': 'Content-Type, Authorization',
    }


def _verify_web_token(request):
    """验证 Web auth token，返回 token 信息或 None"""
    token = request.headers.get('Authorization', '').replace('Bearer ', '')
    info = _web_tokens.get(token)
    if info and time.time() - info['created'] < 86400:
        return info
    if info:
        del _web_tokens[token]
    return None


async def handle_options(request):
    """CORS 预检请求"""
    return aiohttp_web.Response(headers=_cors_headers())


async def handle_login(request):
    """Web 登录（使用 MyData profile.db 认证）

    如果 profile.db 不存在，则视为首次设置：用输入的密码创建 profile.db 并登录。
    """
    try:
        data = await request.json()
        password = data.get('password', '').strip()

        if not password:
            return aiohttp_web.json_response(
                {'ok': False, 'error': '密码不能为空'},
                status=400,
                headers=_cors_headers()
            )

        # 读取 profile.db
        username, password_hash = _read_mydata_profile()

        # profile.db 不存在：首次设置密码
        if not username or not password_hash:
            if bcrypt is None:
                logger.error("bcrypt 未安装，无法设置密码")
                return aiohttp_web.json_response(
                    {'ok': False, 'error': '服务器缺少 bcrypt 依赖'},
                    status=500,
                    headers=_cors_headers()
                )
            try:
                new_username = "user"
                new_hash = bcrypt.hashpw(password.encode('utf-8'), bcrypt.gensalt()).decode('utf-8')
                os.makedirs(os.path.dirname(PROFILE_DB), exist_ok=True)
                with open(PROFILE_DB, 'w', encoding='utf-8') as f:
                    f.write(f"{new_username}\n{new_hash}\n")
                logger.info(f"首次设置密码，已创建 profile.db，用户名: {new_username}")
                username = new_username
            except Exception as e:
                logger.error(f"创建 profile.db 失败: {e}")
                return aiohttp_web.json_response(
                    {'ok': False, 'error': '密码设置失败'},
                    status=500,
                    headers=_cors_headers()
                )
        else:
            # profile.db 存在：验证密码
            if not _verify_mydata_password(password, password_hash):
                logger.warning(f"用户 {username} 登录失败：密码错误")
                return aiohttp_web.json_response(
                    {'ok': False, 'error': '密码错误'},
                    status=401,
                    headers=_cors_headers()
                )

        # 登录成功，生成 token
        token = secrets.token_hex(32)
        _web_tokens[token] = {
            'session_key': f"web:{username}",
            'created': time.time(),
            'username': username,
        }
        logger.info(f"用户 {username} 登录成功")
        return aiohttp_web.json_response(
            {'ok': True, 'token': token, 'username': username},
            headers=_cors_headers()
        )
    except Exception as e:
        logger.error(f"Web 登录处理失败: {e}\n{traceback.format_exc()}")
        return aiohttp_web.json_response(
            {'ok': False, 'error': '服务器错误'},
            status=500,
            headers=_cors_headers()
        )


async def handle_send(request):
    """Web 消息发送（SSE 流式响应），复用 ask_claude + SessionLockManager"""
    auth = _verify_web_token(request)
    if not auth:
        return aiohttp_web.json_response({'error': '未登录或已过期'}, status=401, headers=_cors_headers())

    try:
        data = await request.json()
    except Exception:
        return aiohttp_web.json_response({'error': '无效请求'}, status=400, headers=_cors_headers())

    message = data.get('message', '').strip()
    model_id = data.get('model_id')  # 用户选择的模型组索引
    client_session_id = data.get('session_id')  # 前端传来的 session_id
    is_new = data.get('new', False)  # 前端标记是否新对话
    if not message:
        return aiohttp_web.json_response({'error': '消息为空'}, status=400, headers=_cors_headers())

    session_key = auth['session_key']
    logger.info(f"[Web] 收到消息: {message[:100]}, session_id={client_session_id}, new={is_new}")

    # 检查会话重置：前端 new=true 或消息包含 /new
    if is_new or SessionManager.should_reset_session(message):
        SessionManager.reset_session(session_key)
        message = re.sub(r'/new\b', '', message.replace("新对话", "").replace("重新开始", "").replace("新会话", "")).strip()
        client_session_id = None  # 强制新建
        if not message:
            try:
                headers = _cors_headers()
                headers['Content-Type'] = 'text/event-stream'
                headers['Cache-Control'] = 'no-cache'
                response = aiohttp_web.StreamResponse(headers=headers)
                await response.prepare(request)
                await response.write(f"data: {json.dumps({'type': 'result', 'text': '已开始新对话'}, ensure_ascii=False)}\n\n".encode('utf-8'))
                await response.write_eof()
                return response
            except (ConnectionResetError, ConnectionError, BrokenPipeError):
                logger.info("[Web] 发送新对话确认时客户端已断开")
                return aiohttp_web.Response(status=499)
            except Exception as e:
                logger.warning(f"[Web] 发送新对话确认失败: {e}")
                return aiohttp_web.json_response({'ok': True}, headers=_cors_headers())

    # 获取 session_id：优先使用前端传来的，否则用 SessionManager
    if client_session_id:
        session_id = client_session_id
        SessionManager.save_session_id(session_key, client_session_id)
    else:
        session_id = SessionManager.get_session_id(session_key)

    # SSE 流式响应
    headers = _cors_headers()
    headers['Content-Type'] = 'text/event-stream'
    headers['Cache-Control'] = 'no-cache'
    headers['X-Accel-Buffering'] = 'no'
    response = aiohttp_web.StreamResponse(headers=headers)
    try:
        await response.prepare(request)
    except (ConnectionResetError, ConnectionError) as e:
        logger.warning(f"[Web] SSE 连接准备失败（客户端已断开）: {e}")
        return response

    # 立即回复确认消息
    client_disconnected = False
    try:
        payload = json.dumps({'type': 'progress', 'text': '已经收到您的消息，正在处理中..'}, ensure_ascii=False)
        await response.write(f"data: {payload}\n\n".encode('utf-8'))
    except (ConnectionResetError, ConnectionError, BrokenPipeError):
        client_disconnected = True
        logger.info("[Web] 发送确认消息时客户端已断开")

    # 串行处理
    lock = SessionLockManager.get_lock(session_key)
    is_busy = SessionLockManager.is_busy(session_key)

    # 如果当前会话正在忙碌，先发送排队提示
    if is_busy:
        queue_count = SessionLockManager.get_queue_count(session_key)
        try:
            payload = json.dumps({'type': 'queue', 'text': f'⏳ 正在忙碌中，您的消息已排队（前面还有 {queue_count} 个任务），请稍候...'}, ensure_ascii=False)
            await response.write(f"data: {payload}\n\n".encode('utf-8'))
        except (ConnectionResetError, ConnectionError, BrokenPipeError):
            client_disconnected = True
            logger.info("[Web] 发送排队提示时客户端已断开")
        except Exception as e:
            logger.error(f"[Web] 发送排队提示失败: {e}")

        SessionLockManager.increment_queue(session_key)
        logger.info(f"[Web] 会话 {session_key} 正在忙碌，消息已排队")

    async with lock:
        # 如果之前在排队，减少计数并通知开始处理
        if is_busy:
            SessionLockManager.decrement_queue(session_key)
            if not client_disconnected:
                try:
                    payload = json.dumps({'type': 'start', 'text': '▶️ 开始处理您的消息...'}, ensure_ascii=False)
                    await response.write(f"data: {payload}\n\n".encode('utf-8'))
                except (ConnectionResetError, ConnectionError, BrokenPipeError):
                    client_disconnected = True
                    logger.info("[Web] 发送开始处理提示时客户端已断开")
                except Exception as e:
                    logger.error(f"[Web] 发送开始处理提示失败: {e}")

        SessionLockManager.set_busy(session_key, True)
        task_key = f"task:{session_key}"  # 任务标识符，用于取消任务
        try:
            # 确定存储键：已有 session_id 直接用，新对话用临时键
            storage_key = session_id if session_id else f"pending:{secrets.token_hex(8)}"

            # 立即保存用户消息（对话开始时就能看到）
            ChatHistoryManager.add_message(storage_key, "user", message, "web")
            # 预先创建助手消息占位符，用于实时更新
            assistant_msg_id = ChatHistoryManager.add_message(storage_key, "assistant", "...", "web")
            progress_texts = []  # 收集所有进度消息

            async def on_progress(text):
                nonlocal client_disconnected, storage_key
                if client_disconnected:
                    return
                try:
                    # 检测 session_id 通知：仅新对话时发送 session 事件给前端
                    if text.startswith("📋 已创建任务：session_id = "):
                        sid = text.split("= ", 1)[1].strip()
                        if sid and storage_key.startswith("pending:"):
                            sid_payload = json.dumps({'type': 'session', 'session_id': sid}, ensure_ascii=False)
                            await response.write(f"data: {sid_payload}\n\n".encode('utf-8'))
                            ChatHistoryManager.rename_session(storage_key, sid)
                            storage_key = sid
                            SessionManager.save_session_id(session_key, sid)
                            default_title = re.sub(r'^/new\b\s*', '', message).strip()[:50]
                            if default_title:
                                default_title = _transform_skill_title(default_title)
                            ChatHistoryManager.set_title(storage_key, default_title or "新对话")
                    # session_id 通知也作为普通 progress 显示
                    payload = json.dumps({'type': 'progress', 'text': text}, ensure_ascii=False)
                    await response.write(f"data: {payload}\n\n".encode('utf-8'))
                    # 实时更新助手消息内容
                    clean_text = text.lstrip("🔄 ").strip() if text.startswith("🔄") else text
                    # 收集去重的进度消息
                    if not progress_texts or _normalize_text(clean_text) != _normalize_text(progress_texts[-1]):
                        progress_texts.append(clean_text)
                    ChatHistoryManager.update_message(storage_key, assistant_msg_id, clean_text)
                except (ConnectionResetError, ConnectionError, BrokenPipeError):
                    client_disconnected = True
                    logger.info("[Web] 客户端已断开连接")
                except Exception as e:
                    logger.error(f"[Web] SSE 发送进度失败: {e}")

            reply, returned_session_id, usage_info = await ask_claude(message, session_id, on_progress=on_progress, task_key=task_key, model_id=model_id)

            # 新对话：如果 on_progress 中还没处理过，将临时键重命名为真实 session_id
            if returned_session_id and not session_id and storage_key != returned_session_id:
                ChatHistoryManager.rename_session(storage_key, returned_session_id)
                storage_key = returned_session_id
                SessionManager.save_session_id(session_key, returned_session_id)
                default_title = re.sub(r'^/new\b\s*', '', message).strip()[:50]
                if default_title:
                    default_title = _transform_skill_title(default_title)
                ChatHistoryManager.set_title(storage_key, default_title or "新对话")
            # 去重进度消息，合并为完整的回复内容
            unique_progress = _deduplicate_progress(progress_texts, reply)
            if unique_progress:
                combined_reply = "\n\n".join(f"🔄 {p}" for p in unique_progress) + "\n\n---\n\n" + reply
            else:
                combined_reply = reply
            # 更新助手消息为最终回复（同时更新时间戳和 token 使用统计）
            ChatHistoryManager.update_message(storage_key, assistant_msg_id, combined_reply, update_timestamp=True, usage=usage_info)

            # 记录 token 使用统计
            if usage_info:
                logger.info(f"[Web] Token 统计 - 输入: {usage_info.get('input_tokens', 0)}, "
                          f"输出: {usage_info.get('output_tokens', 0)}, "
                          f"缓存创建: {usage_info.get('cache_creation_input_tokens', 0)}, "
                          f"缓存读取: {usage_info.get('cache_read_input_tokens', 0)}")

            if not client_disconnected:
                try:
                    # 在最终结果中包含 session_id、token 使用信息
                    final_session_id = session_id if session_id else returned_session_id
                    payload = json.dumps({
                        'type': 'result',
                        'text': reply,
                        'session_id': final_session_id,
                        'usage': usage_info
                    }, ensure_ascii=False)
                    await response.write(f"data: {payload}\n\n".encode('utf-8'))
                except (ConnectionResetError, ConnectionError, BrokenPipeError):
                    client_disconnected = True
                    logger.info("[Web] 发送结果时客户端已断开")
        except Exception as e:
            logger.error(f"[Web] 消息处理失败: {e}\n{traceback.format_exc()}")
            if not client_disconnected:
                try:
                    payload = json.dumps({'type': 'error', 'text': '处理失败，请重试'}, ensure_ascii=False)
                    await response.write(f"data: {payload}\n\n".encode('utf-8'))
                except:
                    pass
        finally:
            SessionLockManager.set_busy(session_key, False)

    try:
        await response.write_eof()
    except Exception:
        pass
    logger.info("[Web] 回复完成")
    return response


async def handle_new_session(request):
    """Web 新对话"""
    try:
        auth = _verify_web_token(request)
        if not auth:
            return aiohttp_web.json_response({'error': '未登录'}, status=401, headers=_cors_headers())
        SessionManager.reset_session(auth['session_key'])
        logger.info(f"[Web] 重置会话: {auth['session_key']}")
        return aiohttp_web.json_response({'ok': True}, headers=_cors_headers())
    except Exception as e:
        logger.error(f"[Web] 重置会话失败: {e}")
        return aiohttp_web.json_response({'ok': False, 'error': '服务器错误'}, status=500, headers=_cors_headers())



async def handle_history(request):
    """获取历史消息"""
    auth = _verify_web_token(request)
    if not auth:
        return aiohttp_web.json_response({'error': '未登录'}, status=401, headers=_cors_headers())

    try:
        all_sources = request.query.get('all', 'false').lower() == 'true'
        limit = min(int(request.query.get('limit', 100)), 500)
        offset = int(request.query.get('offset', 0))
        filter_sid = request.query.get('session_id', '')

        if filter_sid:
            # session_id 就是存储键，直接查找
            messages = ChatHistoryManager.get_messages(filter_sid, limit, offset)
        elif all_sources:
            messages = ChatHistoryManager.get_all_messages(limit, offset)
        else:
            # 默认：返回当前活跃会话的消息
            session_key = auth['session_key']
            current_sid = SessionManager.get_session_id(session_key)
            if current_sid:
                messages = ChatHistoryManager.get_messages(current_sid, limit, offset)
            else:
                messages = []

        return aiohttp_web.json_response({
            'ok': True,
            'messages': messages,
        }, headers=_cors_headers())
    except Exception as e:
        logger.error(f"获取历史消息失败: {e}")
        return aiohttp_web.json_response({'error': '获取失败'}, status=500, headers=_cors_headers())


def _transform_skill_title(raw_title: str) -> str:
    """转换 skill 命令为友好的标题格式

    例如：/news-xxx → 年月日-XXX新闻
    """
    raw = raw_title.strip()
    print(f"[DEBUG] _transform_skill_title 输入: {repr(raw)}", flush=True)
    # 匹配 /xxx-yyy 格式，skill 名称至少2个字符
    if raw.startswith('/') and '-' in raw:
        parts = raw[1:].split('-', 1)
        if len(parts) == 2:
            skill_name = parts[0]
            content = parts[1]
            date_str = time.strftime('%Y%m%d')
            skill_map = {
                'news': '新闻',
                'weather': '天气',
                'search': '搜索',
            }
            suffix = skill_map.get(skill_name, skill_name)
            result = f"{date_str}-{content}{suffix}"
            print(f"[DEBUG] _transform_skill_title 输出: {repr(result)}", flush=True)
            return result
    print(f"[DEBUG] _transform_skill_title 未转换，返回原值", flush=True)
    return raw_title


async def handle_conversations(request):
    """获取对话列表（history 格式: {sid: {title, messages}}）"""
    auth = _verify_web_token(request)
    if not auth:
        return aiohttp_web.json_response({'error': '未登录'}, status=401, headers=_cors_headers())

    try:
        history = ChatHistoryManager._load_history()
        conversations = []

        for sid, entry in history.items():
            # 兼容旧格式
            if isinstance(entry, list):
                entry = {"title": "", "messages": entry}
            msgs = entry.get("messages", [])
            if not msgs:
                continue
            # 跳过临时键（正在处理中的新对话）
            if sid.startswith('pending:'):
                continue

            title = entry.get("title", "")
            custom_title = bool(title)
            if not title:
                sorted_m = sorted(msgs, key=lambda m: m.get('timestamp', 0))
                for m in sorted_m:
                    if m.get('role') == 'user':
                        raw = m.get('content', '').strip()
                        raw = re.sub(r'^/new\b\s*', '', raw).strip()
                        if raw:
                            title = _transform_skill_title(raw)[:50]
                            break
            if not title:
                title = '新对话'

            last_msg = max(msgs, key=lambda m: m.get('timestamp', 0))
            conversations.append({
                'session_id': sid,
                'title': title,
                'custom_title': custom_title,
                'last_time': last_msg.get('timestamp', 0),
                'last_time_str': last_msg.get('time_str', ''),
                'msg_count': len(msgs),
                'source': 'web',
            })
        conversations.sort(key=lambda c: c['last_time'], reverse=True)
        return aiohttp_web.json_response({'ok': True, 'conversations': conversations}, headers=_cors_headers())
    except Exception as e:
        logger.error(f"获取对话列表失败: {e}\n{traceback.format_exc()}")
        return aiohttp_web.json_response({'error': '获取失败'}, status=500, headers=_cors_headers())


async def handle_conv_rename(request):
    """重命名对话"""
    auth = _verify_web_token(request)
    if not auth:
        return aiohttp_web.json_response({'error': '未登录'}, status=401, headers=_cors_headers())
    try:
        data = await request.json()
        sid = data.get('session_id', '')
        title = data.get('title', '').strip()[:50]
        if not sid or not title:
            return aiohttp_web.json_response({'error': '参数不完整'}, status=400, headers=_cors_headers())
        ChatHistoryManager.set_title(sid, title)
        return aiohttp_web.json_response({'ok': True}, headers=_cors_headers())
    except Exception as e:
        logger.error(f"重命名对话失败: {e}")
        return aiohttp_web.json_response({'error': '操作失败'}, status=500, headers=_cors_headers())


async def handle_conv_delete(request):
    """删除整个对话（按 session_id 键直接删除）"""
    auth = _verify_web_token(request)
    if not auth:
        return aiohttp_web.json_response({'error': '未登录'}, status=401, headers=_cors_headers())
    try:
        data = await request.json()
        target_sid = data.get('session_id', '')
        if not target_sid:
            return aiohttp_web.json_response({'error': '参数不完整'}, status=400, headers=_cors_headers())

        # session_id 就是存储键，直接删除（title 随 entry 一起删除）
        count = ChatHistoryManager.delete_all(target_sid)

        # 同步删除 Claude CLI 对话日志
        _delete_session_logs(target_sid)

        # 如果删除的是当前活跃会话，重置 SessionManager
        session_key = auth['session_key']
        if SessionManager.get_session_id(session_key) == target_sid:
            SessionManager.reset_session(session_key)

        return aiohttp_web.json_response({'ok': True, 'deleted': count}, headers=_cors_headers())
    except Exception as e:
        logger.error(f"删除对话失败: {e}")
        return aiohttp_web.json_response({'error': '删除失败'}, status=500, headers=_cors_headers())


async def handle_delete_messages(request):
    """删除消息"""
    auth = _verify_web_token(request)
    if not auth:
        return aiohttp_web.json_response({'error': '未登录'}, status=401, headers=_cors_headers())
    
    try:
        data = await request.json()
        action = data.get('action', 'selected')
        message_ids = data.get('ids', [])
        
        if action == 'all':
            count = ChatHistoryManager.delete_all()
            logger.info(f"[Web] 删除全部消息: {count}条")
            return aiohttp_web.json_response({'ok': True, 'deleted': count}, headers=_cors_headers())
        elif action == 'selected' and message_ids:
            count = ChatHistoryManager.delete_by_ids_global(message_ids)
            logger.info(f"[Web] 删除选中消息: {count}条")
            return aiohttp_web.json_response({'ok': True, 'deleted': count}, headers=_cors_headers())
        else:
            return aiohttp_web.json_response({'error': '无效操作'}, status=400, headers=_cors_headers())
    except Exception as e:
        logger.error(f"删除消息失败: {e}")
        return aiohttp_web.json_response({'error': '删除失败'}, status=500, headers=_cors_headers())


async def handle_models(request):
    """返回可用模型列表（不含 API Key）"""
    auth = _verify_web_token(request)
    if not auth:
        return aiohttp_web.json_response({'error': '未登录'}, status=401, headers=_cors_headers())

    models = []
    for i, g in enumerate(MODEL_GROUPS):
        vendor = g.get('MODEL_VENDOR', 'unknown')
        model = g.get('ANTHROPIC_MODEL', '')
        label = f"{vendor} + {model}" if model else vendor
        models.append({
            'id': i,
            'label': label,
            'vendor': vendor,
            'model': model,
        })

    return aiohttp_web.json_response({
        'ok': True,
        'models': models,
        'default': DEFAULT_MODEL_ID,
    }, headers=_cors_headers())


async def handle_get_username(request):
    """获取 MyData 用户名（用于前端自动填充）"""
    try:
        username, _ = _read_mydata_profile()
        if username:
            return aiohttp_web.json_response(
                {'ok': True, 'username': username},
                headers=_cors_headers()
            )
        else:
            return aiohttp_web.json_response(
                {'ok': False, 'error': '无法读取用户信息'},
                status=500,
                headers=_cors_headers()
            )
    except Exception as e:
        logger.error(f"获取用户名失败: {e}")
        return aiohttp_web.json_response(
            {'ok': False, 'error': '服务器错误'},
            status=500,
            headers=_cors_headers()
        )


async def handle_logout(request):
    """Web 退出"""
    try:
        token = request.headers.get('Authorization', '').replace('Bearer ', '')
        if token in _web_tokens:
            del _web_tokens[token]
            logger.info("[Web] 用户退出")
    except Exception as e:
        logger.error(f"[Web] 退出处理异常: {e}")
    return aiohttp_web.json_response({'ok': True}, headers=_cors_headers())


async def handle_cancel(request):
    """取消正在进行的 Cutebot 任务"""
    auth = _verify_web_token(request)
    if not auth:
        return aiohttp_web.json_response({'error': '未登录'}, status=401, headers=_cors_headers())

    session_key = auth['session_key']
    task_key = f"task:{session_key}"

    try:
        if task_key in _active_tasks:
            task_info = _active_tasks[task_key]
            # 标记任务为已取消
            task_info['cancelled'] = True
            # 强制终止进程组（包括所有子进程）
            proc = task_info.get('proc')
            if proc and proc.returncode is None:
                try:
                    # 杀掉整个进程组
                    os.killpg(os.getpgid(proc.pid), signal.SIGTERM)
                    logger.info(f"[Web] 已发送 SIGTERM 到进程组 {proc.pid}")
                    # 等待一小段时间
                    await asyncio.sleep(0.3)
                    # 如果还没退出，强制杀掉
                    if proc.returncode is None:
                        os.killpg(os.getpgid(proc.pid), signal.SIGKILL)
                        logger.info(f"[Web] 已发送 SIGKILL 到进程组 {proc.pid}")
                except ProcessLookupError:
                    logger.info(f"[Web] 进程 {proc.pid} 已经不存在")
                except Exception as e:
                    logger.warning(f"[Web] 终止进程失败: {e}")
            logger.info(f"[Web] 任务 {task_key} 已标记为取消")
            return aiohttp_web.json_response({'ok': True, 'message': '任务已取消'}, headers=_cors_headers())
        else:
            logger.info(f"[Web] 没有找到活动任务: {task_key}")
            return aiohttp_web.json_response({'ok': True, 'message': '没有进行中的任务'}, headers=_cors_headers())
    except Exception as e:
        logger.error(f"[Web] 取消任务失败: {e}")
        return aiohttp_web.json_response({'error': '取消失败'}, status=500, headers=_cors_headers())


async def handle_events(request):
    """SSE 长连接，向已登录的 Web 客户端推送新消息（来自 QQ 或其他来源）"""
    auth = _verify_web_token(request)
    if not auth:
        return aiohttp_web.json_response({'error': '未登录'}, status=401, headers=_cors_headers())

    # since 参数：客户端已知的最新消息时间戳，只推送更新的消息
    try:
        since = float(request.query.get('since', time.time()))
    except (ValueError, TypeError):
        since = time.time()

    headers = _cors_headers()
    headers['Content-Type'] = 'text/event-stream; charset=utf-8'
    headers['Cache-Control'] = 'no-cache'
    headers['X-Accel-Buffering'] = 'no'
    response = aiohttp_web.StreamResponse(headers=headers)
    try:
        await response.prepare(request)
    except Exception as e:
        logger.warning(f"[Events] 准备 SSE 响应失败: {e}")
        return response

    # 发送心跳注释保持连接活跃
    async def send_heartbeat():
        try:
            await response.write(b": heartbeat\n\n")
        except Exception:
            return False
        return True

    last_since = since
    try:
        while True:
            # 每 2 秒检查一次新消息
            await asyncio.sleep(2)
            try:
                all_msgs = ChatHistoryManager.get_all_messages(200)
                new_msgs = [m for m in all_msgs if m.get('timestamp', 0) > last_since]
                if new_msgs:
                    # 按时间正序排列
                    new_msgs.sort(key=lambda x: x.get('timestamp', 0))
                    for msg in new_msgs:
                        payload = json.dumps(msg, ensure_ascii=False)
                        await response.write(f"data: {payload}\n\n".encode('utf-8'))
                    last_since = new_msgs[-1].get('timestamp', last_since)
                else:
                    # 每 2 次循环（4s）发一次心跳
                    ok = await send_heartbeat()
                    if not ok:
                        break
            except (ConnectionResetError, ConnectionError, BrokenPipeError):
                logger.info("[Events] 客户端断开连接")
                break
            except Exception as e:
                logger.warning(f"[Events] 推送消息异常: {e}")
                break
    except asyncio.CancelledError:
        pass
    return response


def _delete_session_logs(session_id: str):
    """删除指定 session_id 对应的 Claude CLI 日志文件和子目录"""
    if not session_id or not os.path.isdir(CLAUDE_LOG_DIR):
        return
    jsonl_path = os.path.join(CLAUDE_LOG_DIR, f"{session_id}.jsonl")
    subdir_path = os.path.join(CLAUDE_LOG_DIR, session_id)
    removed = []
    if os.path.isfile(jsonl_path):
        try:
            os.remove(jsonl_path)
            removed.append(f"{session_id}.jsonl")
        except Exception as e:
            logger.warning(f"[清理] 删除日志文件失败 {jsonl_path}: {e}")
    if os.path.isdir(subdir_path):
        try:
            shutil.rmtree(subdir_path)
            removed.append(f"{session_id}/")
        except Exception as e:
            logger.warning(f"[清理] 删除日志目录失败 {subdir_path}: {e}")
    if removed:
        logger.info(f"[清理] 已删除 session {session_id} 的 CLI 日志: {', '.join(removed)}")


def _get_active_session_ids() -> set:
    """获取当前活跃的 session_id 集合（chat_history + sessions.json 的并集）"""
    active = set()
    # 从 chat_history.json 获取
    history = ChatHistoryManager._load_history()
    for key in history:
        # session_id 格式是 UUID，跳过 pending: 前缀的临时键
        if not key.startswith('pending:'):
            active.add(key)
    # 从 sessions.json 获取
    try:
        if os.path.exists(SESSION_FILE):
            with open(SESSION_FILE, 'r', encoding='utf-8') as f:
                sessions = json.load(f)
            for sid in sessions.values():
                if sid:
                    active.add(sid)
    except Exception as e:
        logger.warning(f"[清理] 读取 sessions.json 失败: {e}")
    return active


async def _cleanup_cli_logs():
    """定期清理 Claude CLI 日志（每天凌晨执行）"""
    while True:
        try:
            # 计算到下一个凌晨 00:00 的秒数
            now = datetime.now()
            tomorrow = (now + timedelta(days=1)).replace(hour=0, minute=0, second=0, microsecond=0)
            wait_seconds = (tomorrow - now).total_seconds()
            await asyncio.sleep(wait_seconds)

            logger.info("[清理] 开始每日 CLI 日志清理...")

            # 1. 清理 debug 日志（全部删除）
            debug_removed = 0
            if os.path.isdir(CLAUDE_DEBUG_DIR):
                for f in os.listdir(CLAUDE_DEBUG_DIR):
                    fp = os.path.join(CLAUDE_DEBUG_DIR, f)
                    try:
                        if os.path.isfile(fp):
                            os.remove(fp)
                            debug_removed += 1
                    except Exception:
                        pass
            if debug_removed:
                logger.info(f"[清理] 已删除 {debug_removed} 个 debug 日志文件")

            # 2. 清理孤立对话日志（不在活跃集合中的 .jsonl 和子目录）
            if os.path.isdir(CLAUDE_LOG_DIR):
                active = _get_active_session_ids()
                orphan_removed = 0
                # 保留最近 3 天的文件作为安全网（防止误删终端直接用 claude 的对话）
                cutoff = time.time() - 3 * 86400
                for entry in os.listdir(CLAUDE_LOG_DIR):
                    # 提取 session_id
                    sid = entry.replace('.jsonl', '') if entry.endswith('.jsonl') else entry
                    if sid in active:
                        continue
                    fp = os.path.join(CLAUDE_LOG_DIR, entry)
                    # 检查修改时间，保留最近 3 天的
                    try:
                        mtime = os.path.getmtime(fp)
                        if mtime > cutoff:
                            continue
                    except Exception:
                        continue
                    # 删除孤立文件/目录
                    try:
                        if os.path.isfile(fp):
                            os.remove(fp)
                            orphan_removed += 1
                        elif os.path.isdir(fp):
                            # 跳过非 UUID 目录（如 memory 等系统目录）
                            if len(sid) == 36 and sid.count('-') == 4:
                                shutil.rmtree(fp)
                                orphan_removed += 1
                    except Exception as e:
                        logger.warning(f"[清理] 删除孤立日志失败 {entry}: {e}")
                if orphan_removed:
                    logger.info(f"[清理] 已删除 {orphan_removed} 个孤立对话日志")

            # 3. 清理备份目录（保留最近 7 天的）
            backup_removed = 0
            if os.path.isdir(CLAUDE_BACKUP_DIR):
                backup_cutoff = time.time() - 7 * 86400
                for f in os.listdir(CLAUDE_BACKUP_DIR):
                    fp = os.path.join(CLAUDE_BACKUP_DIR, f)
                    try:
                        if os.path.getmtime(fp) < backup_cutoff:
                            if os.path.isfile(fp):
                                os.remove(fp)
                                backup_removed += 1
                            elif os.path.isdir(fp):
                                shutil.rmtree(fp)
                                backup_removed += 1
                    except Exception:
                        pass
            if backup_removed:
                logger.info(f"[清理] 已删除 {backup_removed} 个过期备份文件")

            logger.info("[清理] 每日 CLI 日志清理完成")

        except asyncio.CancelledError:
            break
        except Exception as e:
            logger.error(f"[清理] CLI 日志清理异常: {e}")


async def _cleanup_expired_tokens():
    """定期清理过期的 Web 登录 token"""
    while True:
        try:
            await asyncio.sleep(3600)
            now = time.time()
            expired = [t for t, info in _web_tokens.items() if now - info['created'] > 86400]
            for t in expired:
                del _web_tokens[t]
            if expired:
                logger.info(f"[Web] 清理 {len(expired)} 个过期 token")
        except asyncio.CancelledError:
            break
        except Exception as e:
            logger.error(f"[Web] token 清理异常: {e}")


async def _periodic_history_flush():
    """每 30 秒将内存中的脏历史记录刷入磁盘"""
    while True:
        try:
            await asyncio.sleep(30)
            ChatHistoryManager.flush()
        except asyncio.CancelledError:
            break
        except Exception as e:
            logger.error(f"定期刷盘异常: {e}")


async def handle_status(request):
    """查询当前会话是否正在执行任务"""
    auth = _verify_web_token(request)
    if not auth:
        return aiohttp_web.json_response({'error': '未登录'}, status=401, headers=_cors_headers())

    session_key = auth['session_key']
    busy = SessionLockManager.is_busy(session_key)
    queue_count = SessionLockManager.get_queue_count(session_key) if busy else 0
    return aiohttp_web.json_response({
        'ok': True,
        'busy': busy,
        'queue_count': queue_count
    }, headers=_cors_headers())



async def handle_skills(request):
    """返回可用的 skill 列表（扫描 commands 目录下的 .md 文件）"""
    auth = _verify_web_token(request)
    if not auth:
        return aiohttp_web.json_response({'error': '未登录'}, status=401, headers=_cors_headers())

    skills = []
    try:
        if os.path.isdir(COMMANDS_DIR):
            for fname in sorted(os.listdir(COMMANDS_DIR)):
                if not fname.endswith('.md'):
                    continue
                name = fname[:-3]  # 去掉 .md
                # 读取第一行作为描述
                desc = ''
                try:
                    with open(os.path.join(COMMANDS_DIR, fname), 'r', encoding='utf-8') as f:
                        first_line = f.readline().strip()
                    desc = first_line.lstrip('#').strip()[:80]
                except Exception:
                    pass
                skills.append({'name': name, 'desc': desc})
    except Exception as e:
        logger.warning(f"扫描 commands 目录失败: {e}")

    return aiohttp_web.json_response({
        'ok': True,
        'skills': skills,
    }, headers=_cors_headers())


async def start_http_server():
    """启动 Web Chat HTTP API"""
    if aiohttp_web is None:
        logger.error("aiohttp 未安装，Web Chat API 无法启动")
        return
    try:
        app = aiohttp_web.Application()
        app.router.add_route('OPTIONS', '/api/{path:.*}', handle_options)
        app.router.add_get('/api/username', handle_get_username)
        app.router.add_post('/api/login', handle_login)
        app.router.add_post('/api/send', handle_send)
        app.router.add_post('/api/new', handle_new_session)
        app.router.add_post('/api/logout', handle_logout)
        app.router.add_get('/api/history', handle_history)
        app.router.add_get('/api/conversations', handle_conversations)
        app.router.add_post('/api/conv/rename', handle_conv_rename)
        app.router.add_post('/api/conv/delete', handle_conv_delete)
        app.router.add_post('/api/delete', handle_delete_messages)
        app.router.add_post('/api/cancel', handle_cancel)
        app.router.add_get('/api/models', handle_models)
        app.router.add_get('/api/events', handle_events)
        app.router.add_get('/api/status', handle_status)
        app.router.add_get('/api/skills', handle_skills)
        runner = aiohttp_web.AppRunner(app)
        await runner.setup()
        site = aiohttp_web.TCPSite(runner, '0.0.0.0', HTTP_PORT)
        await site.start()
        logger.info(f"Web Chat API 已启动: http://0.0.0.0:{HTTP_PORT}")
        asyncio.create_task(_cleanup_expired_tokens())
        asyncio.create_task(_periodic_history_flush())
        asyncio.create_task(_cleanup_cli_logs())
    except OSError as e:
        logger.error(f"Web Chat API 启动失败（端口 {HTTP_PORT} 可能被占用）: {e}")
    except Exception as e:
        logger.error(f"Web Chat API 启动失败: {e}\n{traceback.format_exc()}")


if __name__ == "__main__":
    async def main():
        logger.info("Cutebot Web 服务启动中...")
        await start_http_server()
        # 保持运行
        try:
            while True:
                await asyncio.sleep(3600)
        except asyncio.CancelledError:
            pass

    try:
        asyncio.run(main())
    except KeyboardInterrupt:
        logger.info("Cutebot 被手动关闭")
    except Exception as e:
        logger.error(f"Cutebot 启动失败: {e}\n{traceback.format_exc()}")
    finally:
        ChatHistoryManager.flush()
        logger.info("已将聊天历史缓存写入磁盘")
