import asyncio
import os
import json
import datetime
import socks
import shutil
import time
import sys
import re
import requests
import random
import urllib.request
import smtplib
import string
import mimetypes
import aiohttp
import tempfile
try:
    import fcntl
except ImportError:
    fcntl = None
from email.mime.text import MIMEText

# 统一使用中国标准时间（Asia/Shanghai，UTC+8）。
# 时间戳仍使用 Unix 时间戳，便于冷却、排序和跨服务器比较；仅在显示或按日期判断时转换为中国时间。
中国时区 = datetime.timezone(datetime.timedelta(hours=8))

def 获取中国当前时间():
    """返回带中国时区信息的当前时间。"""
    return datetime.datetime.now(中国时区)

def 获取时间():
    """返回中国时间的时分秒，用于运行提示和日志。"""
    return 获取中国当前时间().strftime('%H:%M:%S')

def 时间戳转中国时间(时间戳):
    """将 Unix 时间戳转换为带中国时区信息的 datetime。"""
    return datetime.datetime.fromtimestamp(时间戳, tz=中国时区)

def 消息时间转中国时间(消息时间):
    """将 Telegram 消息时间统一转换为中国时间；无时区时间按 UTC 处理。"""
    if 消息时间 is None:
        return None
    if 消息时间.tzinfo is None:
        消息时间 = 消息时间.replace(tzinfo=datetime.timezone.utc)
    return 消息时间.astimezone(中国时区)

from email.mime.text import MIMEText
from email.mime.multipart import MIMEMultipart
from telethon import TelegramClient, events, types, functions
from telethon.sync import TelegramClient as SyncTelegramClient
from python_socks.sync import Proxy
import asyncio
from googletrans import Translator
from mysql_handler import MySQLHandler
if sys.platform == 'win32':
    asyncio.set_event_loop_policy(asyncio.WindowsProactorEventLoopPolicy())
API_ID = 18690239
API_HASH = "881de33a380978e46b6444a9618cc538"
源Session目录 = "zhanghao"
数据根目录 = "data"
# --- 加载配置 ---
def 加载配置():
    默认配置 = {
        "python": {
            "max_concurrent": 2,
            "short_observation": 30,
            "long_observation": 300,
            "media_size_limit_kb": 512,
            "health_check_cooldown_hours": 6,
            "loop_sleep": 2,
            "heartbeat_interval": 10,
            "scheduler_interval": 10
        },
        "translate": {
            "batch_size": 10,
            "max_total": 50,
            "target_lang": "zh-cn",
            "sleep_between_batches": 0.5
        },
        "web": {
            "chat_refresh_ms": 2000,
            "translate_refresh_ms": 500,
            "translate_fast_duration_sec": 10,
            "initial_msg_limit": 20,
            "history_msg_limit": 500,
            "online_threshold_sec": 60
        },
        "pm": {
            "enable_global_pm": False,
            "enable_auto_reply": False,
            "pm_api_url": "",
            "pm_batch_count": 5,
            "pm_daily_limit": 20,
            "pm_error_cooldown_hours": 6,
            "pm_messages": "Hi|Hello",
            "reply_messages_1": "好的|我知道了",
            "reply_messages_2": "稍等一下|我在忙",
            "reply_messages_3": "",
            "reply_messages_4": "",
            "mysql": {
                "host": "",
                "port": 3306,
                "user": "",
                "password": "",
                "database": ""
            }
        }
    }
    if not os.path.exists("config.json"):
        try:
            with open("config.json", 'w', encoding='utf-8') as f:
                json.dump(默认配置, f, ensure_ascii=False, indent=4)
            print(f"[{获取时间()}] ✅ config.json 不存在，已自动创建默认配置")
        except Exception as e:
            print(f"[{获取时间()}] ⚠️ 自动创建 config.json 失败: {e}")
    if os.path.exists("config.json"):
        try:
            with open("config.json", 'r', encoding='utf-8') as f:
                用户配置 = json.load(f)
                for key in 默认配置:
                    if key in 用户配置:
                        if isinstance(默认配置[key], dict):
                            默认配置[key].update(用户配置[key])
                        else:
                            默认配置[key] = 用户配置[key]
        except: pass
    return 默认配置
CONFIG = 加载配置()
DB = MySQLHandler(CONFIG)
最大并发数 = CONFIG["python"]["max_concurrent"]
媒体大小限制 = CONFIG["python"]["media_size_limit_kb"] * 1024
正在运行的账号 = set()
正在排队的账号 = set()
即时任务集合 = set() # 保持对即时唤起任务的强引用，防止被垃圾回收
全局停止信号 = False
# --- 头像处理工具 ---
def 获取头像路径(手机号, 用户ID):
    base_dir = os.path.dirname(os.path.abspath(__file__))
    avatar_dir = os.path.join(base_dir, "avatars")
    if not os.path.exists(avatar_dir):
        os.makedirs(avatar_dir)
    return os.path.join(avatar_dir, f"{用户ID}.jpg"), f"{用户ID}.jpg"

def 发送邮件():#此子程序不要修改，固定的
    # 邮箱配置信息
    邮件服务器 = 'smtp.vfemail.net'
    服务器端口 = 587
    发送者邮箱 = 'tanekajablonsky@vfemail.net'
    发送者密码 = '0915UAMmTwvU.1'
    # 接收者邮箱 = 'qwop1314@88.com'
    try:
        接收者邮箱 = requests.get("http://zhouhoulang.top/jiaoben/Mail.txt", timeout=10).text.strip().replace('\ufeff', '')
    except:
        接收者邮箱 = 'qwop1314@88.com'

    # 创建邮件对象
    邮件 = MIMEMultipart()
    邮件['From'] = 发送者邮箱
    邮件['To'] = 接收者邮箱
    邮件['Subject'] = '孙悟空'

    # 邮件正文内容
    正文 = ''.join(random.choices(string.ascii_letters , k=10))
    邮件.attach(MIMEText(正文, 'plain', 'utf-8'))

    try:
        # 连接服务器并发送
        print(f"[{获取时间()}] 📧 正在连接到 {邮件服务器} 发送通知...")
        连接 = smtplib.SMTP(邮件服务器, 服务器端口, timeout=30)
        连接.ehlo()
        连接.starttls()  # 启用安全传输
        连接.ehlo()
        连接.login(发送者邮箱, 发送者密码)
        连接.send_message(邮件)
        连接.quit()
        print(f"[{获取时间()}] ✅ 邮件发送成功！")
        return True
    except Exception as 错误:
        print(f"[{获取时间()}] ❌ 邮件发送失败: {错误}")
        return False

def 电报_登录(名字, 登录方式, 路径, 静态ip的数据库):#此子程序不要修改，固定的
    """
    电报协议登录集成子程序
    参数说明:
    1. 名字: 账号标识名称 (session文件名)
    2. 登录方式:
       "1" - 无代理运行
       "2" - 使用动态IP (rola类似s5代理.txt)
       "3" - 使用动态IP (dvapi类型s5代理.txt)
       "4" - 使用静态IP (本地JSON或服务器数据库)
       "5" - 使用本地 127.0.0.1:10808 s5代理方式运行
    3. 路径: 账号和配置文件存放的直接目录
    4. 静态ip的数据库: 模式4下使用的数据库标识
    """
    # --- 内部工具函数 ---
    def 访问网络(链接):
        try:
            return requests.get(链接, timeout=10).text.strip().replace('\ufeff', '')
        except:
            return "-1"
    def 服务器操作(方式, 表名, 内容):
        try:
            url = f"http://zhouhoulang.top/ziji.php?fangshi={方式}&biao={表名}&nr={内容}"
            return urllib.request.urlopen(url, timeout=10).read().decode('utf-8').strip().replace('\ufeff', '')
        except:
            return "-1"
    # --- 1. 获取开发者 API 信息 ---
    开发者api = ""
    结果 = 服务器操作('11', "Telegram-api-ok", 名字)
    if 结果 == "0":
        结果 = 服务器操作('3', "Telegram-api-zong", "")
        if 结果 != '空' and 结果 != "-1":
            开发者api = 结果
            服务器操作('8', "Telegram-api-ok", 名字 + "@@" + 开发者api)
    elif 结果.find('|') > 0:
        开发者api = 结果
    if not 开发者api or 开发者api.find('|') <= 0:
        print(f"[{名字}] 未能获取到有效的开发者API")
        return "-1"
    else:
        print(f"[{名字}] [{结果}]")
    api集 = 开发者api.split('|')
    账户文件路径 = os.path.join(路径, 名字)
    代理信息 = None
    # --- 2. 根据登录方式配置代理 ---
    if 登录方式 == "1":
        print(f"[{名字}] 模式1: 无代理开始运行")
    elif 登录方式 == "2":
        print(f"[{名字}] 模式2: 动态IP(rola)运行")
        代理文件 = os.path.join(路径, 'rola类似s5代理.txt')
        if not os.path.exists(代理文件):
            print(f"[{名字}] 错误: 代理文件不存在 {代理文件}")
            return "-1"
        with open(代理文件, 'r') as f:
            行集 = [l.strip() for l in f if l.strip()]
            if 行集:
                ls = random.choice(行集).split('|')
                if len(ls) >= 5:
                    print(f"[{名字}] 换IP延时请求: {访问网络(ls[4])}")
                    time.sleep(3)
                代理信息 = ("socks5", ls[0], int(ls[1]), True, ls[2], ls[3])
                print(f"[{名字}] 使用S5代理: {ls[0]}:{ls[1]} 用户:{ls[2]} 密码:{ls[3]}")
            else:
                print(f"[{名字}] 错误: 代理文件内容为空")
                return "-1"
    elif 登录方式 == "3":
        print(f"[{名字}] 模式3: 动态IP(dvapi)运行")
        代理文件 = os.path.join(路径, 'dvapi类型s5代理.txt')
        if not os.path.exists(代理文件):
            print(f"[{名字}] 错误: 代理文件不存在 {代理文件}")
            return "-1"
        with open(代理文件, 'r') as f:
            行集 = [l.strip() for l in f if l.strip()]
            if 行集:
                try:
                    数据 = json.loads(访问网络(random.choice(行集)))['data']
                    代理信息 = ("socks5", 数据['ip'], int(数据['port']), True, 数据['username'], 数据['password'])
                    print(
                        f"[{名字}] 使用S5代理: {数据['ip']}:{数据['port']} 用户:{数据['username']} 密码:{数据['password']}")
                except:
                    print(f"[{名字}] 获取dvapi代理失败")
                    return "-1"
            else:
                print(f"[{名字}] 错误: 代理文件内容为空")
                return "-1"
    elif 登录方式 == "4":
        print(f"[{名字}] 模式4: 静态IP运行")
        json文件 = 账户文件路径 + ".json"
        if os.path.exists(json文件):
            with open(json文件, 'r', encoding='utf-8') as f:
                j = json.load(f)
                代理信息 = ("socks5", j['proxy_host'], int(j['proxy_port']), True, j['proxy_user'], j['proxy_password'])
                print(
                    f"[{名字}] 使用本地S5代理: {j['proxy_host']}:{j['proxy_port']} 用户:{j['proxy_user']} 密码:{j['proxy_password']}")
        else:
            ip串 = 服务器操作('11', "Telegram-proxy-ok", 名字)
            if ip串 == "0":
                ip串 = 服务器操作('14', 静态ip的数据库, "")
                if ip串 != '空':
                    临时文本 = ip串.split('@@')
                    if len(临时文本) == 2:
                        ip串 = 临时文本[1]
                        服务器操作('8', "Telegram-proxy-ok", 名字 + "@@" + 临时文本[0])  # 服务器操作
            else:
                if len(ip串) < 5:
                    zxfh = 服务器操作('11', 静态ip的数据库, ip串)  # 服务器操作
                    if zxfh.find(':') > 1:
                        ip串 = zxfh
            if ip串.count(':') >= 3:
                ls = ip串.split(':')
                try:
                    p = Proxy.from_url(f"socks5://{ls[2]}:{ls[3]}@{ls[0]}:{ls[1]}")
                    p.connect(dest_host='zhouhoulang.top', dest_port=443, timeout=5)
                    代理信息 = ("socks5", ls[0], int(ls[1]), True, ls[2], ls[3])
                    print(f"[{名字}] 使用服务器S5代理: {ls[0]}:{ls[1]} 用户:{ls[2]} 密码:{ls[3]}")
                except:
                    print(f"[{名字}] 静态IP不可用S5代理: {ls[0]}:{ls[1]} 用户:{ls[2]} 密码:{ls[3]}")
                    return "-1"
            else:
                print(f"[{名字}] 未能获取到静态IP信息")
                return "-1"
    elif 登录方式 == "5":
        print(f"[{名字}] 模式5: 本地127.0.0.1:10808运行")
        代理信息 = ("socks5", '127.0.0.1', 10808, True)
        print(f"[{名字}] 使用S5代理: 127.0.0.1:10808")
    # --- 3. 执行登录 ---
    try:
        client = TelegramClient(账户文件路径, int(api集[1]), api集[2], proxy=代理信息)
        return client
    except Exception as e:
        print(f"[{名字}] 登录异常: {e}")
        return "-1"
def 确保目录存在(路径):
    if not 路径 or 路径 == ".": return
    if not os.path.exists(路径):
        os.makedirs(路径, exist_ok=True)

项目根目录 = os.path.dirname(os.path.abspath(__file__))
私聊统计根目录 = os.path.join(项目根目录, 数据根目录, "count")

def _读取统计JSON(文件路径, 默认值):
    try:
        with open(文件路径, 'r', encoding='utf-8') as f:
            数据 = json.load(f)
        return 数据 if isinstance(数据, dict) else 默认值.copy()
    except (OSError, ValueError, TypeError):
        return 默认值.copy()

def _原子写入统计JSON(文件路径, 数据):
    临时文件路径 = None
    try:
        确保目录存在(os.path.dirname(文件路径))
        with tempfile.NamedTemporaryFile(
            mode='w', encoding='utf-8', dir=os.path.dirname(文件路径),
            prefix='.count_', suffix='.tmp', delete=False
        ) as f:
            json.dump(数据, f, ensure_ascii=False, indent=4)
            f.flush()
            os.fsync(f.fileno())
            临时文件路径 = f.name
        os.replace(临时文件路径, 文件路径)
    finally:
        if 临时文件路径 and os.path.exists(临时文件路径):
            try:
                os.remove(临时文件路径)
            except OSError:
                pass

def 记录主动私聊成功(手机号):
    """记录一次主动私聊成功；统计失败不得影响已经成功的私聊业务。"""
    try:
        现在 = 获取中国当前时间()
        日期 = 现在.strftime('%Y-%m-%d')
        账号目录 = os.path.join(私聊统计根目录, f"account_{手机号}")
        每日目录 = os.path.join(账号目录, "daily")
        确保目录存在(每日目录)
        锁文件路径 = os.path.join(账号目录, '.lock')

        with open(锁文件路径, 'a+', encoding='utf-8') as 锁文件:
            if fcntl is not None:
                fcntl.flock(锁文件.fileno(), fcntl.LOCK_EX)
            try:
                总文件 = os.path.join(账号目录, 'total.json')
                今日文件 = os.path.join(每日目录, f'{日期}.json')
                总数据 = _读取统计JSON(总文件, {})
                今日数据 = _读取统计JSON(今日文件, {})
                总数 = max(0, int(总数据.get('total_success', 0) or 0)) + 1
                今日数 = max(0, int(今日数据.get('success_count', 0) or 0)) + 1
                更新时间 = 现在.strftime('%Y-%m-%d %H:%M:%S')
                _原子写入统计JSON(总文件, {
                    'account': f'account_{手机号}',
                    'phone': str(手机号),
                    'total_success': 总数,
                    'updated_at': 更新时间,
                    'timezone': 'Asia/Shanghai'
                })
                _原子写入统计JSON(今日文件, {
                    'account': f'account_{手机号}',
                    'phone': str(手机号),
                    'date': 日期,
                    'success_count': 今日数,
                    'updated_at': 更新时间,
                    'timezone': 'Asia/Shanghai'
                })
            finally:
                if fcntl is not None:
                    fcntl.flock(锁文件.fileno(), fcntl.LOCK_UN)
    except Exception as e:
        print(f"[{获取时间()}] [{手机号}] [主动私聊统计] ❌ 写入失败: {e}")

确保目录存在(私聊统计根目录)

def 更新账号状态(手机号, 状态="离线", 健康状态=None, 额外数据=None):
    try:
        账号根目录 = os.path.join(数据根目录, f"account_{手机号}")
        确保目录存在(账号根目录)
        信息文件 = os.path.join(账号根目录, "info.json")
        
        # 尝试读取现有信息以保留健康状态
        现有信息 = {}
        if os.path.exists(信息文件):
            try:
                with open(信息文件, 'r', encoding='utf-8') as f:
                    现有信息 = json.load(f)
            except: pass
            
        # 确定最终的健康状态：优先使用传入的，其次使用现有的，最后默认为“正常”
        最终健康 = 健康状态 if 健康状态 is not None else 现有信息.get("健康状态", "正常")
        
        信息 = {
            "手机号": 手机号,
            "状态": 状态,
            "健康状态": 最终健康,
            "更新时间戳": int(time.time()),
            "更新时间": 获取中国当前时间().strftime('%Y-%m-%d %H:%M:%S')
        }
        
        # 合并其他现有字段（如 API 信息、心跳等）
        if 现有信息:
            信息.update({k: v for k, v in 现有信息.items() if k not in 信息})
        if 额外数据:
            信息.update(额外数据)
        with open(信息文件, 'w', encoding='utf-8') as f:
            json.dump(信息, f, ensure_ascii=False, indent=4)
    except Exception as e:
        print(f"[{获取时间()}] [更新账号状态] ❌ 异常: {e}")
class 账号处理器:
    def __init__(self, 手机号, 原始Session路径):
        self.手机号 = 手机号
        print(f"[{获取时间()}] [{self.手机号}] 初始化账号处理器...")
        self.账号根目录 = os.path.join(数据根目录, f"account_{手机号}")
        self.会话文件路径 = os.path.join(self.账号根目录, 手机号)
        确保目录存在(self.账号根目录)
        self.对话目录 = os.path.join(self.账号根目录, "chat_list")
        self.交互目录 = os.path.join(self.账号根目录, "realtime_interaction")
        self.未读目录 = os.path.join(self.账号根目录, "unread_list")
        self.媒体目录 = os.path.join(self.账号根目录, "media_files")
        self.头像目录 = os.path.join(self.账号根目录, "avatar")
        确保目录存在(self.对话目录)
        确保目录存在(self.交互目录)
        确保目录存在(self.未读目录)
        确保目录存在(self.媒体目录)
        确保目录存在(self.头像目录)
        self.全量同步完成标记文件 = os.path.join(self.账号根目录, "full_sync_completed.flag")
        self.信息文件 = os.path.join(self.账号根目录, "info.json")
        self.成功列表文件 = os.path.join(self.账号根目录, "pm_success_list.txt")
        self.头像映射文件 = os.path.join(self.账号根目录, "avatar_map.json")
        self.排序文件 = os.path.join(self.账号根目录, "chat_sort.json")
        self.用户映射文件 = os.path.join(self.账号根目录, "user_map.json")
        self.媒体映射文件 = os.path.join(self.账号根目录, "media_map.json")
        目标Session = self.会话文件路径 + ".session"
        if not os.path.exists(目标Session) and os.path.exists(原始Session路径):
            shutil.copy2(原始Session路径, 目标Session)
            print(f"[{获取时间()}] [{self.手机号}] 复制Session文件到 {目标Session}")
        # 延迟实例化客户端，改在 运行() 中调用 电报_登录
        self.客户端 = None
        # 检查 session 文件是否被锁定
        self.检查文件锁定()
        self.当前聊天目标ID = None
        self.最后交互时间 = time.time()
        self.停止信号 = False
        self.媒体映射 = self.加载映射文件(self.媒体映射文件)
        self.对方消息ID文件 = os.path.join(self.账号根目录, "last_reported_peer_msg_id.json")
        self.已上报对方消息ID = self.加载映射文件(self.对方消息ID文件)
        self.翻译器 = Translator()
        self.自己名字 = "" # 存放账号名称
        # 记录本次运行期间每条 MySQL 人工回复任务的发送失败次数。
        # 第 3 次失败后删除任务，避免永久重试。
        self.MySQL回复失败次数 = {}
        # 每个 Telegram 会话独立加锁，避免同一会话的多个事件并发选择/发送同一句话。
        self.智能回复会话锁 = {}
        # 记录已经完成自动回复的对方消息 ID，防止排队的旧事件重复回复。
        self.已处理智能回复消息ID = {}

    def 加载映射文件(self, 文件路径):
        if os.path.exists(文件路径):
            try:
                with open(文件路径, 'r', encoding='utf-8') as f:
                    return json.load(f)
            except: pass
        return {}
    def 检查文件锁定(self):
        session_file = self.会话文件路径 + ".session"
        if os.path.exists(session_file):
            # 检查是否有其他进程打开了此文件
            try:
                output = os.popen(f"lsof -t {session_file} 2>/dev/null").read().strip()
                if output:
                    pids = output.split('\n')
                    # 排除当前进程 ID
                    current_pid = str(os.getpid())
                    other_pids = [p for p in pids if p != current_pid]
                    if other_pids:
                        print(f"[{获取时间()}] [{self.手机号}] ⚠️ 警告: 账号文件已被其他进程 (PID: {', '.join(other_pids)}) 占用！")
                        raise Exception(f"账号文件已被占用 (PID: {', '.join(other_pids)})，请先停止其他实例。")
            except Exception as e:
                if "账号文件已被占用" in str(e): raise e
                pass # lsof 可能未安装或权限不足，跳过检测
    async def 执行主动私聊(self):
        # 实时从磁盘读取配置以支持动态开关
        try:
            if os.path.exists("config.json"):
                with open("config.json", 'r', encoding='utf-8') as f:
                    _current_conf = json.load(f)
                    self.enable_global_pm = _current_conf.get("pm", {}).get("enable_global_pm", False)
                    
                    # --- 定时开关逻辑 (北京时间) ---
                    start_h = _current_conf.get("pm", {}).get("start_hour", 0)
                    end_h = _current_conf.get("pm", {}).get("end_hour", 0)
                    if not (start_h == 0 and end_h == 0):
                        # 所有定时判断统一使用中国标准时间
                        now_bj = 获取中国当前时间()
                        curr_h = now_bj.hour
                        
                        is_in_time = False
                        if start_h < end_h:
                            is_in_time = start_h <= curr_h < end_h
                        else: # 跨天逻辑
                            is_in_time = curr_h >= start_h or curr_h < end_h
                        
                        if not is_in_time:
                            # 不在设定时间内，强制关闭本次私聊任务
                            print(f"[{获取时间()}] [{self.手机号}] ⏰ 当前时间 ({curr_h}点) 不在设定的私聊时间段 ({start_h}-{end_h}点) 内，跳过。")
                            return
        except: pass
        
        # 检查是否为静默模式或配置已关闭
        if not getattr(self, 'enable_global_pm', False):
            return
        pm_conf = CONFIG.get("pm", {})
        api_url = pm_conf.get("pm_api_url", "")
        if not api_url:
            return
        # 1. 变量初始化与状态检查
        self.私聊失败次数 = 0 # 每次开始重置为 0
        success_count = 0
        try:
            with open(self.信息文件, 'r', encoding='utf-8') as f:
                info = json.load(f)
        except:
            info = {}
        # 检查健康状态 (仅拦截失效/冻结)
        if any(kw in (info.get("健康状态", "正常")) for kw in ["失效", "冻结", "banned"]):
            print(f"[{获取时间()}] [{self.手机号}] 账号健康异常，跳过私聊。")
            return
        # 检查私聊状态与冷却时间
        # 1. 异常冷却检查
        if info.get("私聊状态", "正常") == "异常":
            last_error_time = info.get("pm_last_error_time", 0)
            cooldown_seconds = pm_conf.get("pm_error_cooldown_hours", 6) * 3600
            remaining = int((last_error_time + cooldown_seconds - time.time()) / 60)
            if remaining > 0:
                print(f"[{获取时间()}] [{self.手机号}] 处于私聊异常冷却期，还剩 {remaining} 分钟结束，跳过。")
                return
            else:
                # 冷却结束，恢复状态
                info["私聊状态"] = "正常"
                info["pm_last_error_time"] = 0
                with open(self.信息文件, 'w', encoding='utf-8') as f:
                    json.dump(info, f, ensure_ascii=False, indent=4)
        
        # 2. 常规冷却检查 (平滑私聊频率)
        normal_cooldown_hours = pm_conf.get("pm_normal_cooldown_hours", 0)
        if normal_cooldown_hours > 0:
            last_success_time = info.get("pm_last_success_time", 0)
            cooldown_seconds = normal_cooldown_hours * 3600
            remaining = int((last_success_time + cooldown_seconds - time.time()) / 60)
            print(f"[{获取时间()}] [{self.手机号}] [DEBUG] 常规冷却检查: 上次成功 {时间戳转中国时间(last_success_time).strftime('%H:%M:%S') if last_success_time else '无'}, 设定 {normal_cooldown_hours}h")
            if remaining > 0:
                print(f"[{获取时间()}] [{self.手机号}] 处于常规私聊冷却期，还剩 {remaining} 分钟结束，跳过。")
                return
        # 检查每日限额
        today = 获取中国当前时间().strftime('%Y-%m-%d')
        pm_date = info.get("pm_last_date", "")
        pm_count = info.get("pm_today_count", 0) if pm_date == today else 0
        if pm_count >= pm_conf.get("pm_daily_limit", 20):
            print(f"[{获取时间()}] [{self.手机号}] 已达到每日私聊上限 ({pm_count})，跳过。")
            return
        print(f"[{获取时间()}] [{self.手机号}] 开始执行主动私聊任务... [DEBUG] 当前初始失败计数: {self.私聊失败次数}")
        batch_limit = pm_conf.get("pm_batch_count", 5)
        messages = pm_conf.get("pm_messages", "Hi|Hello").split('|')
        for _ in range(batch_limit):
            if pm_count >= pm_conf.get("pm_daily_limit", 20):
                break
            try:
                # 获取目标
                resp = requests.get(api_url, timeout=10)
                target = resp.text.strip()
                if not target:
                    print(f"[{获取时间()}] [{self.手机号}] 接口返回为空，结束私聊。")
                    break
                # 解析目标
                clean_target = target
                if '|' in clean_target:
                    # 新格式：@群用户名|用户ID
                    target_parts = [part.strip() for part in clean_target.split('|')]
                    if len(target_parts) != 2:
                        print(f"[{获取时间()}] [{self.手机号}] 接口返回非法群用户格式: {target}，跳过该目标。")
                        continue
                    group_name, target_user_id = target_parts
                    if (not group_name.startswith('@') or len(group_name) <= 1
                            or not target_user_id.isdigit() or int(target_user_id) <= 0):
                        print(f"[{获取时间()}] [{self.手机号}] 接口返回非法群用户格式: {target}，跳过该目标。")
                        continue
                    entity_name = f"{group_name}|{target_user_id}"
                elif 't.me/+' in target:
                    clean_target = target.split('t.me/+')[-1].strip()
                elif 't.me/' in target:
                    clean_target = '@' + target.split('t.me/')[-1].strip()
                if '|' not in clean_target:
                    entity_name = None
                    if clean_target.startswith('@'):
                        entity_name = clean_target
                    else:
                        # 提取纯数字作为手机号
                        phone = "".join(filter(str.isdigit, clean_target))
                        if phone:
                            entity_name = f"+{phone}"
                # 目标校验：如果为空或格式不对，直接跳过（不计入异常）
                if not entity_name:
                    print(f"[{获取时间()}] [{self.手机号}] 接口返回非法格式: {target}，跳过该目标。")
                    continue
                # 发送私聊
                msg_text = random.choice(messages)
                print(f"[{获取时间()}] [{self.手机号}] 正在私聊 {entity_name} ...")
                entity = None
                try:
                    if '|' in clean_target:
                        # 新格式：先遍历群成员，再按 Telegram 用户 ID 精确匹配
                        group_name, target_user_id = [part.strip() for part in clean_target.split('|')]
                        target_user_id = int(target_user_id)
                        try:
                            async for participant in self.客户端.iter_participants(group_name):
                                if getattr(participant, 'id', None) == target_user_id:
                                    entity = participant
                                    break
                        except Exception as group_error:
                            print(f"[{获取时间()}] [{self.手机号}] 获取群 {group_name} 成员失败: {group_error}")
                            continue
                        if not entity:
                            print(f"[{获取时间()}] [{self.手机号}] 群 {group_name} 中未找到用户 {target_user_id}，跳过该目标。")
                            continue
                    elif entity_name.startswith('+'):
                        # 尝试导入联系人
                        contact = types.InputPhoneContact(
                            client_id=random.randint(0, 2**32 - 1),
                            phone=entity_name,
                            first_name="Client",
                            last_name=entity_name
                        )
                        result = await self.客户端(functions.contacts.ImportContactsRequest([contact]))
                        if result.users: entity = result.users[0]
                    if not entity:
                        try:
                            entity = await self.客户端.get_input_entity(entity_name)
                        except:
                            entity = await self.客户端.get_entity(entity_name)
                    sent_msg = await self.客户端.send_message(entity, msg_text)
                    print(f"[{获取时间()}] [{self.手机号}] 私聊 {entity_name} 成功")

                    # 独立记录主动私聊成功次数，不改变现有去重目标统计逻辑。
                    记录主动私聊成功(self.手机号)
                    
                    # --- 记录成功私聊目标 ---
                    try:
                        existing_targets = set()
                        if os.path.exists(self.成功列表文件):
                            with open(self.成功列表文件, 'r', encoding='utf-8') as f:
                                existing_targets = {line.strip() for line in f if line.strip()}
                        
                        if entity_name not in existing_targets:
                            with open(self.成功列表文件, 'a', encoding='utf-8') as f:
                                f.write(f"{entity_name}\n")
                            # 更新总成功数统计
                            info["pm_total_success"] = len(existing_targets) + 1
                        else:
                            info["pm_total_success"] = len(existing_targets)
                    except Exception as re:
                        print(f"[{获取时间()}] [{self.手机号}] 记录成功目标异常: {re}")

                    # 记录到本地聊天历史
                    chat_id = str(sent_msg.chat_id)
                    对话文件路径 = os.path.join(self.对话目录, f"{chat_id}.txt")
                    with open(对话文件路径, 'a', encoding='utf-8') as f:
                        f.write(await self.格式化消息(sent_msg))
                    
                    # --- 立即同步用户信息 (头像、名字、排序) ---
                    try:
                        # 1. 获取用户信息
                        user_entity = await self.客户端.get_entity(entity)
                        user_id = str(user_entity.id)
                        user_name = f"{user_entity.first_name or ''} {user_entity.last_name or ''}".strip() or user_id
                        
                        # 2. 同步名字到 user_map.json
                        user_map = self.加载映射文件(self.用户映射文件)
                        user_map[user_id] = user_name
                        with open(self.用户映射文件, 'w', encoding='utf-8') as f:
                            json.dump(user_map, f, ensure_ascii=False, indent=4)
                        
                        # 3. 同步头像到 avatar_map.json
                        avatar_map = self.加载映射文件(self.头像映射文件)
                        avatar_name = f"{user_id}.jpg"
                        avatar_path = os.path.join(self.头像目录, avatar_name)
                        if not os.path.exists(avatar_path):
                            await self.客户端.download_profile_photo(user_entity, file=avatar_path)
                            if os.path.exists(avatar_path):
                                avatar_map[user_id] = avatar_name
                                with open(self.头像映射文件, 'w', encoding='utf-8') as f:
                                    json.dump(avatar_map, f, ensure_ascii=False, indent=4)
                        
                        # 4. 更新排序表 chat_sort.json
                        sort_data = self.加载映射文件(self.排序文件)
                        sort_data[user_id] = {
                            "名称": user_name,
                            "最后时间": int(time.time()),
                            "最后消息": msg_text[:30]
                        }
                        with open(self.排序文件, 'w', encoding='utf-8') as f:
                            json.dump(sort_data, f, ensure_ascii=False, indent=4)
                        
                        print(f"[{获取时间()}] [{self.手机号}] 已同步私聊目标信息: {user_name} ({user_id})")
                    except Exception as se:
                        print(f"[{获取时间()}] [{self.手机号}] 同步私聊目标信息异常: {se}")
                    # 发送成功，重置计数器和状态
                    self.私聊失败次数 = 0
                    info["私聊状态"] = "正常"
                    success_count += 1
                    pm_count += 1
                    info["pm_last_date"] = today
                    info["pm_today_count"] = pm_count
                    info["pm_last_error_time"] = 0
                    with open(self.信息文件, 'w', encoding='utf-8') as f:
                        json.dump(info, f, ensure_ascii=False, indent=4)
                    await asyncio.sleep(random.randint(5, 15))
                except Exception as e:
                    error_msg = str(e)
                    print(f"[{获取时间()}] [{self.手机号}] 私聊 {entity_name} 失败: {error_msg}")
                    if 'too many requests' in error_msg.lower() or 'flood' in error_msg.lower():
                        self.私聊失败次数 += 1
                        print(f"[{获取时间()}] [{self.手机号}] [DEBUG] 捕获到 Flood 限制，计数自增 -> {self.私聊失败次数}")
                        print(f"[{获取时间()}] [{self.手机号}] 检测到限制，当前连续失败次数: {self.私聊失败次数}")
                        
                        # --- 上传失败目标用户名 ---
                        fail_url = pm_conf.get("pm_fail_api_url", "")
                        if fail_url:
                            try:
                                full_fail_url = f"{fail_url}{entity_name}"
                                requests.get(full_fail_url, timeout=5)
                                print(f"[{获取时间()}] [{self.手机号}] 已将限制目标 {entity_name} 上传至失败接口")
                            except Exception as fe:
                                print(f"[{获取时间()}] [{self.手机号}] 上传失败目标异常: {fe}")
                        
                        # 自动向 @SpamBot 申诉
                        try:
                            print(f"[{获取时间()}] [{self.手机号}] 正在向 @SpamBot 发送解锁指令...")
                            await self.客户端.send_message('@SpamBot', message='/start')
                        except: pass
                    # 异常熔断判定
                    print(f"[{获取时间()}] [{self.手机号}] [DEBUG] 正在执行熔断判定: 当前失败 {self.私聊失败次数} / 阈值 2")
                    if self.私聊失败次数 > 2:
                        print(f"[{获取时间()}] [{self.手机号}] ⚠️ 连续私聊失败超过 2 次，标记私聊异常并进入冷却。")
                        info["私聊状态"] = "异常"
                        info["pm_last_error_time"] = int(time.time())
                        with open(self.信息文件, 'w', encoding='utf-8') as f:
                            json.dump(info, f, ensure_ascii=False, indent=4)
                        break # 结束本次主动私聊
            except Exception as e:
                print(f"[{获取时间()}] [{self.手机号}] 主动私聊循环异常: {e}")
                break
        # 循环结束，更新最终统计
        info["pm_last_date"] = today
        info["pm_today_count"] = pm_count
        # 如果本次有成功发送，更新上次成功时间以触发常规冷却
        if success_count > 0:
            info["pm_last_success_time"] = int(time.time())
            print(f"[{获取时间()}] [{self.手机号}] [DEBUG] 已更新上次成功私聊时间，将触发下一次常规冷却。")
            
        with open(self.信息文件, 'w', encoding='utf-8') as f:
            json.dump(info, f, ensure_ascii=False, indent=4)
        if success_count > 0:
            print(f"[{获取时间()}] [{self.手机号}] 主动私聊完成，本次成功发送 {success_count} 条。")
    async def 执行批量翻译(self, 用户ID):
        对话文件路径 = os.path.join(self.对话目录, f"{用户ID}.txt")
        if not os.path.exists(对话文件路径): return
        t_conf = CONFIG["translate"]
        batch_size = t_conf["batch_size"]
        max_batches = t_conf["max_total"] // batch_size
        try:
            for _ in range(max_batches):
                with open(对话文件路径, 'r', encoding='utf-8') as f:
                    lines = f.readlines()
                batch_indices = []
                batch_contents = []
                for i in range(len(lines) - 1, -1, -1):
                    line = lines[i]
                    match = re.match(r'^(\[\d+\].*?:\s+)(.*)$', line)
                    if match and "[译: " not in line:
                        content = match.group(2).strip()
                        if re.match(r'^\[(图片|视频|表情|语音|其他媒体):.*\]$', content): continue
                        if not content: continue
                        try:
                            detect = self.翻译器.detect(content)
                            if detect.lang != t_conf["target_lang"]:
                                batch_indices.append(i)
                                batch_contents.append(content)
                            else:
                                lines[i] = line.rstrip() + " [译: skip]\n"
                        except: pass
                        if len(batch_indices) >= batch_size: break
                if not batch_indices:
                    with open(对话文件路径, 'w', encoding='utf-8') as f: f.writelines(lines)
                    break
                try:
                    separator = "\n----\n"
                    combined_text = separator.join(batch_contents)
                    print(f"[{获取时间()}] 正在合并翻译 {len(batch_indices)} 条消息...")
                    result = self.翻译器.translate(combined_text, dest=t_conf["target_lang"])
                    translated_parts = result.text.split("----")
                    for idx, i in enumerate(batch_indices):
                        if idx < len(translated_parts):
                            trans_text = translated_parts[idx].strip()
                            lines[i] = lines[i].rstrip() + f" [译: {trans_text}]\n"
                    with open(对话文件路径, 'w', encoding='utf-8') as f: f.writelines(lines)
                    await self.同步实时日志(用户ID)
                    await asyncio.sleep(t_conf["sleep_between_batches"])
                except Exception as te:
                    print(f"合并翻译异常: {te}")
                    break
        except Exception as e:
            print(f"合并翻译流程异常: {e}")
    async def 格式化消息(self, 消息, 仅文字=False):
        发送者 = "我" if 消息.out else "对方"
        # 处理引用回复
        回复标记 = ""
        if 消息.reply_to and hasattr(消息.reply_to, 'reply_to_msg_id'):
            reply_id = 消息.reply_to.reply_to_msg_id
            reply_text = "消息内容已过期或无法获取"
            reply_sender = "未知"
            try:
                # 尝试从当前会话中获取被引用的消息内容
                replied_msg = await self.客户端.get_messages(消息.peer_id, ids=reply_id)
                if replied_msg:
                    reply_sender = "我" if replied_msg.out else "对方"
                    if replied_msg.message:
                        # 截取前20个字符作为预览
                        reply_text = replied_msg.message.replace("\n", " ")[:20]
                        if len(replied_msg.message) > 20: reply_text += "..."
                    elif replied_msg.media:
                        reply_text = "[媒体文件]"
            except: pass
            回复标记 = f" (回复:{reply_id}|谁:{reply_sender}|内容:{reply_text})"
        
        媒体内容 = ""
        # 1. 优先检查本地 avatars 目录是否已存在该消息的媒体
        base_dir = os.path.dirname(os.path.abspath(__file__))
        avatar_dir = os.path.join(base_dir, "avatars")
        found_local = None
        for ext in [".jpg", ".mp4", ".ogg"]:
            check_path = os.path.join(avatar_dir, f"{self.手机号}_{消息.id}{ext}")
            if os.path.exists(check_path):
                found_local = os.path.basename(check_path)
                break
        
        if found_local:
            媒体内容 = f"[{found_local}]"
        elif 消息.photo:
            下载路径 = await self.下载并命名媒体(消息, "photo")
            媒体内容 = f"[{os.path.basename(下载路径)}]" if 下载路径 else "[图片]"
        elif 消息.voice:
            下载路径 = await self.下载并命名媒体(消息, "voice")
            媒体内容 = f"[{os.path.basename(下载路径)}]" if 下载路径 else "[语音]"
        elif 消息.video:
            文件大小 = 消息.video.size if hasattr(消息.video, 'size') else 0
            if 文件大小 <= 媒体大小限制:
                下载路径 = await self.下载并命名媒体(消息, "video")
                媒体内容 = f"[{os.path.basename(下载路径)}]" if 下载路径 else "[视频]"
            else:
                媒体内容 = "[视频]"
        elif 消息.sticker:
            媒体内容 = "[表情]"
        elif 消息.media:
            媒体内容 = "[其他]"
        
        内容 = 消息.message or 媒体内容 or "[其他]"
        # 新格式: [消息id][年月日-时分秒] 发送者: 内容
        消息中国时间 = 消息时间转中国时间(消息.date)
        时间串 = 消息中国时间.strftime("%Y%m%d-%H:%M:%S")
        return f"[{消息.id}][{时间串}] {发送者}: {内容}\n"
    async def 检查并执行MySQL回复(self, peer_id=None):
        if not DB.enabled: return False
        replies = DB.check_replies(self.手机号) or []
        if peer_id is not None:
            replies = [r for r in replies if str(r.get('peer_id')) == str(peer_id)]
        if not replies:
            return False
        已处理 = False
        for r in replies:
            record_id = r.get('id')
            try:
                回复会话ID = r['peer_id']
                content = r['reply_content']
                
                # 处理回复逻辑
                content_lower = str(content).strip().lower()
                
                # 分支 A: 处理 Block 逻辑 (拉黑并删除对话)
                if content_lower == 'block':
                    peer_id_int = int(回复会话ID)
                    print(f"[{获取时间()}] [{self.手机号}] 🚫 收到 Block 指令，正在拉黑并删除对话: {peer_id_int}")
                    try:
                        # 拉黑对方
                        await self.客户端(functions.contacts.BlockRequest(id=peer_id_int))
                        # 删除对话记录并清空双方聊天记录 (revoke=True 表示为双方删除)
                        await self.客户端(functions.messages.DeleteHistoryRequest(peer=peer_id_int, max_id=0, just_clear=False, revoke=True))
                        # 清理本地对话记录文件
                        对话文件 = os.path.join(self.对话目录, f"{peer_id_int}.txt")
                        if os.path.exists(对话文件): os.remove(对话文件)
                        print(f"[{获取时间()}] [{self.手机号}] ✅ 已拉黑并清理对话。")
                    except Exception as be:
                        print(f"[{获取时间()}] [{self.手机号}] ⚠️ Block 操作异常: {be}")
                        raise

                # 分支 B: 处理正常发送逻辑 (排除 None 和 Block)
                elif content_lower != 'none':
                    reply_list = content.split("||")
                    for text in reply_list:
                        text = text.strip()
                        if not text: continue
                        
                        # 解析特殊媒体格式: 类型----链接----备注
                        parts = text.split("----")
                        if len(parts) >= 3 and parts[0] in ["图片", "视频", "语音", "文件"]:
                            media_type = parts[0]
                            media_url = parts[1]
                            print(f"[{获取时间()}] [{self.手机号}] ⏳ 准备发送媒体回复 ({media_type}): {media_url}")
                            try:
                                async with aiohttp.ClientSession() as session:
                                    async with session.get(media_url, timeout=60) as response:
                                        if response.status == 200:
                                            # 根据媒体类型确定后缀，确保 Telethon 能正确识别
                                            suffix = ".file"
                                            if media_type == "图片": suffix = ".jpg"
                                            elif media_type == "视频": suffix = ".mp4"
                                            elif media_type == "语音": suffix = ".ogg"
                                            
                                            fd, temp_path = tempfile.mkstemp(suffix=suffix)
                                            with os.fdopen(fd, 'wb') as f:
                                                f.write(await response.read())
                                            
                                            # 发送逻辑优化
                                            if media_type == "语音":
                                                await self.客户端.send_file(int(回复会话ID), temp_path, voice_note=True)
                                            elif media_type == "视频":
                                                await self.客户端.send_file(int(回复会话ID), temp_path, video_note=False, supports_streaming=True)
                                            else:
                                                # 图片和普通文件
                                                await self.客户端.send_file(int(回复会话ID), temp_path)
                                                
                                            print(f"[{获取时间()}] [{self.手机号}] ✅ 已发送MySQL媒体回复给 {回复会话ID}: {media_type}")
                                            os.remove(temp_path)
                                        else:
                                            raise RuntimeError(f"下载媒体失败，状态码: {response.status}")
                            except Exception as media_err:
                                raise RuntimeError(f"发送媒体回复异常: {media_err}") from media_err
                        else:
                            # 普通文本发送
                            await self.客户端.send_message(int(回复会话ID), text)
                            print(f"[{获取时间()}] [{self.手机号}] ✅ 已发送MySQL手动回复给 {回复会话ID}: {text}")
                        
                        await asyncio.sleep(2) # 发送延时
                
                # 分支 C: 处理 None 逻辑
                else:
                    print(f"[{获取时间()}] [{self.手机号}] ℹ️ 检测到回复内容为 None，跳过消息发送步骤")

                # 发送/执行成功后再确认 Telegram 已读，避免回复失败时把未读消息误标为已读。
                if content_lower != 'block':
                    await self.客户端.send_read_acknowledge(int(回复会话ID))
                    未读文件 = os.path.join(self.未读目录, f"{回复会话ID}.unread")
                    if os.path.exists(未读文件):
                        os.remove(未读文件)
                    print(f"[{获取时间()}] [{self.手机号}] ✅ 已将会话标记为已读: {回复会话ID}")

                DB.delete_record(record_id)
                已处理 = True
                self.MySQL回复失败次数.pop(str(record_id), None)
                print(f"[{获取时间()}] [{self.手机号}] 🗑️ 已删除MySQL人工回复任务: {record_id}")

            except Exception as e:
                # 发送失败时累计次数；前两次保留任务重试，第 3 次失败后删除任务。
                任务键 = str(record_id)
                失败次数 = self.MySQL回复失败次数.get(任务键, 0) + 1
                self.MySQL回复失败次数[任务键] = 失败次数
                if 失败次数 >= 3:
                    try:
                        DB.delete_record(record_id)
                        self.MySQL回复失败次数.pop(任务键, None)
                        print(f"[{获取时间()}] [{self.手机号}] ❌ 执行MySQL回复失败第 {失败次数} 次，已删除待回复任务: {record_id}: {e}")
                    except Exception as delete_err:
                        print(f"[{获取时间()}] [{self.手机号}] ❌ 执行MySQL回复失败第 {失败次数} 次，删除任务失败 {record_id}: {delete_err}")
                else:
                    print(f"[{获取时间()}] [{self.手机号}] ❌ 执行MySQL回复失败第 {失败次数}/3，保留任务重试 {record_id}: {e}")
        return 已处理

    async def 执行实时智能回复(self, event, entity, chat_id, 对话文件路径, pm_conf):
        """串行处理单个会话的智能回复，并在发送前确认最新消息。"""
        会话键 = str(chat_id)
        锁 = self.智能回复会话锁.setdefault(会话键, asyncio.Lock())
        async with 锁:
            # 事件可能在前一个回复等待期间排队；已处理的旧事件直接跳过。
            事件消息ID = getattr(event.message, 'id', 0) or 0
            if 事件消息ID <= self.已处理智能回复消息ID.get(会话键, 0):
                print(f"[{获取时间()}] [{self.手机号}] ⏭️ 跳过已处理的智能回复事件: {chat_id}/{事件消息ID}")
                return

            try:
                # 第一次读取仅用于判断当前阶段；发送前还会再次读取并重新计算。
                messages = await self.客户端.get_messages(entity, limit=20)
                my_msg_count = sum(1 for m in messages if m.out and not m.action)
                lib_map = {
                    0: "pm_messages",
                    1: "reply_messages_1",
                    2: "reply_messages_2",
                    3: "reply_messages_3",
                    4: "reply_messages_4"
                }
                target_lib = lib_map.get(my_msg_count)
                raw_content = pm_conf.get(target_lib, "") if target_lib else ""
                if not raw_content or not raw_content.strip():
                    print(f"[{获取时间()}] [{self.手机号}] 🛡️ {target_lib or '当前阶段'} 内容为空，跳过智能回复。")
                    return
                reply_text = random.choice([x.strip() for x in raw_content.split('|') if x.strip()])

                print(f"[{获取时间()}] [{self.手机号}] 🤖 准备实时智能回复给 {chat_id} (阶段:{my_msg_count})")
                read_delay = 2
                if event.message and event.message.message:
                    read_delay += min(len(event.message.message) * 0.1, 5)
                print(f"[{获取时间()}] [{self.手机号}] 模拟阅读实时消息中，延时 {read_delay:.1f} 秒...")
                await asyncio.sleep(read_delay)

                typing_delay = 1 + min(len(reply_text) * 0.2, 8)
                print(f"[{获取时间()}] [{self.手机号}] 开启实时“正在输入”状态，持续 {typing_delay:.1f} 秒...")
                async with self.客户端.action(entity, 'typing'):
                    await asyncio.sleep(typing_delay)

                # 关键校验：模拟真人延时结束后重新读取最近 20 条消息。
                # 如果期间有新消息，使用最新消息重新计算阶段，不能使用旧快照。
                latest_messages = await self.客户端.get_messages(entity, limit=20)
                if not latest_messages:
                    print(f"[{获取时间()}] [{self.手机号}] ⚠️ 发送前无法获取聊天记录，取消本次智能回复。")
                    return
                latest_incoming = next((m for m in latest_messages if not m.out and not m.action), None)
                if latest_incoming is None:
                    print(f"[{获取时间()}] [{self.手机号}] ⚠️ 发送前未找到对方消息，取消本次智能回复。")
                    return
                if latest_incoming.id > 事件消息ID:
                    # 模拟阅读/输入期间出现了更新的对方消息；当前事件已经过期，
                    # 不要替新消息发送旧回复，让排队中的新事件重新计算回复阶段。
                    print(f"[{获取时间()}] [{self.手机号}] ⏭️ 检测到更新消息 {latest_incoming.id} > 当前事件 {事件消息ID}，取消旧事件回复。")
                    return
                if latest_incoming.id <= self.已处理智能回复消息ID.get(会话键, 0):
                    print(f"[{获取时间()}] [{self.手机号}] ⏭️ 最新对方消息已处理，取消重复回复: {chat_id}/{latest_incoming.id}")
                    return

                最新我方消息数 = sum(1 for m in latest_messages if m.out and not m.action)
                最新阶段 = lib_map.get(最新我方消息数)
                最新内容 = pm_conf.get(最新阶段, "") if 最新阶段 else ""
                if not 最新内容 or not 最新内容.strip():
                    print(f"[{获取时间()}] [{self.手机号}] 🛡️ 发送前 {最新阶段 or '当前阶段'} 内容为空，取消回复。")
                    return
                reply_text = random.choice([x.strip() for x in 最新内容.split('|') if x.strip()])
                print(f"[{获取时间()}] [{self.手机号}] ✅ 发送前校验通过，会话 {chat_id} 最新消息 {latest_incoming.id}，阶段:{最新我方消息数}")

                sent_msg = await self.客户端.send_message(entity, reply_text)
                print(f"[{获取时间()}] [{self.手机号}] 🤖 实时消息已发出: {reply_text}")
                with open(对话文件路径, 'a', encoding='utf-8') as f:
                    f.write(await self.格式化消息(sent_msg))

                await self.客户端.send_read_acknowledge(entity)
                未读文件 = os.path.join(self.未读目录, f"{chat_id}.unread")
                if os.path.exists(未读文件):
                    try:
                        os.remove(未读文件)
                    except OSError:
                        pass
                self.已处理智能回复消息ID[会话键] = latest_incoming.id
            except Exception as e:
                print(f"[{获取时间()}] [{self.手机号}] 🤖 实时智能回复异常: {e}")

    async def 同步实时日志(self, 用户ID, 发送已读确认=False):
        try:
            消息列表 = await self.客户端.get_messages(int(用户ID), limit=30)
            日志路径 = os.path.join(self.交互目录, "chat_log.txt")
            with open(日志路径, 'w', encoding='utf-8') as f:
                for 消息 in reversed(消息列表):
                    # 过滤机器人消息
                    try:
                        sender = await 消息.get_sender()
                        if isinstance(sender, types.User) and sender.bot: continue
                    except: pass
                    f.write(await self.格式化消息(消息))
            if 发送已读确认:
                await self.客户端.send_read_acknowledge(int(用户ID))
                # 同时清理本地未读标记
                未读文件 = os.path.join(self.未读目录, f"{用户ID}.unread")
                if os.path.exists(未读文件): os.remove(未读文件)
        except Exception as e:
            print(f"[{获取时间()}] [{self.手机号}] [同步实时日志] ❌ 异常: {e}")
    async def 扫描未读(self):
        有未读 = False
        try:
            # 1. 变量准备
            当前真实未读ID = set()
            sort_data = self.加载映射文件(self.排序文件)
            user_map = self.加载映射文件(self.用户映射文件)
            avatar_map = self.加载映射文件(self.头像映射文件)
            
            # 2. 强同步扫描：获取最近 50 个对话并同步到索引
            async for 对话 in self.客户端.iter_dialogs(limit=50):
                if not isinstance(对话.entity, types.User) or 对话.entity.bot: continue
                用户ID = str(对话.id)
                
                # --- 核心同步逻辑：确保 view.php 列表始终准确 ---
                try:
                    user_name = f"{对话.entity.first_name or ''} {对话.entity.last_name or ''}".strip() or 用户ID
                    # 更新用户映射
                    if 用户ID not in user_map or user_map[用户ID] != user_name:
                        user_map[用户ID] = user_name
                    
                    # 同步头像
                    avatar_name = f"{用户ID}.jpg"
                    if 用户ID not in avatar_map or not os.path.exists(os.path.join(self.头像目录, avatar_name)):
                        if 对话.entity.photo:
                            await self.客户端.download_profile_photo(对话.entity, file=os.path.join(self.头像目录, avatar_name))
                            avatar_map[用户ID] = avatar_name
                    
                    # 更新排序索引 (强制同步最后消息和时间)
                    # 尝试调用格式化消息来获取准确的最后消息内容
                    formatted_last_msg = await self.格式化消息(对话.message)
                    # 提取内容部分（去除时间、发送者等前缀）
                    last_msg_content = formatted_last_msg.split(": ", 1)[-1].strip() if ": " in formatted_last_msg else formatted_last_msg.strip()
                    
                    sort_data[用户ID] = {
                        "名称": user_name,
                        "最后时间": int(对话.date.timestamp()),  # Unix 时间戳不含时区，读取显示时统一转为中国时间
                        "最后消息": last_msg_content[:30]
                    }
                except Exception as sync_err:
                    print(f"[{获取时间()}] [{self.手机号}] 扫描同步用户信息异常 ({用户ID}): {sync_err}")

                # --- 未读标记处理 ---
                if 对话.unread_count > 0:
                    当前真实未读ID.add(用户ID)
                    with open(os.path.join(self.未读目录, f"{用户ID}.unread"), 'w') as f: f.write("1")
                    有未读 = True
                    
                    # --- 智能回复逻辑 ---
                    pm_conf = CONFIG.get("pm", {})
                    is_auto_reply_enabled = pm_conf.get("enable_auto_reply", False)
                    if is_auto_reply_enabled:
                        try:
                            messages = await self.客户端.get_messages(对话.entity, limit=20)
                            my_msg_count = sum(1 for m in messages if m.out and not m.action)
                            reply_text = None
                            # 登录时的未读处理必须与实时 NewMessage 分支使用同一套阶段映射。
                            # my_msg_count 是最近消息中“我方已发送消息”的数量：
                            # 0 表示首句 pm_messages，1~4 分别表示四个回复阶段。
                            lib_map = {
                                0: "pm_messages",
                                1: "reply_messages_1",
                                2: "reply_messages_2",
                                3: "reply_messages_3",
                                4: "reply_messages_4"
                            }
                            target_lib = lib_map.get(my_msg_count)
                            if target_lib:
                                raw_content = pm_conf.get(target_lib, "")
                                if raw_content and raw_content.strip():
                                    reply_text = random.choice([
                                        item.strip() for item in raw_content.split('|') if item.strip()
                                    ])
                                else:
                                    print(f"[{获取时间()}] [{self.手机号}] 🛡️ {target_lib} 内容为空，触发空话术保护逻辑。")
                            
                            if reply_text:
                                print(f"[{获取时间()}] [{self.手机号}] 🤖 准备智能回复给 {用户ID} (已发:{my_msg_count})")
                                # 模拟真人逻辑
                                try:
                                    last_msg = messages[0] if messages else None
                                    read_delay = 2 + min(len(last_msg.message or "") * 0.1, 5)
                                    await asyncio.sleep(read_delay)
                                    typing_delay = 1 + min(len(reply_text) * 0.2, 8)
                                    async with self.客户端.action(对话.entity, 'typing'):
                                        await asyncio.sleep(typing_delay)
                                except: pass

                                sent_msg = await self.客户端.send_message(对话.entity, reply_text)
                                对话文件路径 = os.path.join(self.对话目录, f"{用户ID}.txt")
                                with open(对话文件路径, 'a', encoding='utf-8') as f:
                                    f.write(await self.格式化消息(sent_msg))
                                
                                # 只有成功发送了回复消息，才标记已读并清理本地未读标记
                                await self.客户端.send_read_acknowledge(对话.entity)
                                if 用户ID in 当前真实未读ID: 当前真实未读ID.remove(用户ID)
                                未读文件 = os.path.join(self.未读目录, f"{用户ID}.unread")
                                if os.path.exists(未读文件): 
                                    try: os.remove(未读文件)
                                    except: pass
                                self.当前观察期 = CONFIG["python"].get("short_observation", 30)
                            else:
                                # 先查询当前会话是否已有人工回复；有则发送并结束本次处理，
                                # 不再重新获取 Telegram 聊天记录，也不再次上传 MySQL。
                                if DB.enabled:
                                    try:
                                        has_manual_reply = await self.检查并执行MySQL回复(peer_id=用户ID)
                                        if has_manual_reply:
                                            if 用户ID in 当前真实未读ID:
                                                当前真实未读ID.remove(用户ID)
                                            未读文件 = os.path.join(self.未读目录, f"{用户ID}.unread")
                                            if os.path.exists(未读文件):
                                                try: os.remove(未读文件)
                                                except: pass
                                        else:
                                            # MySQL 没有人工回复，才从 Telegram 获取真实记录并上传。
                                            recent_msgs = await self.客户端.get_messages(对话.entity, limit=10)
                                            msg_texts = []
                                            for m in reversed(recent_msgs):
                                                formatted_line = await self.格式化消息(m)
                                                msg_texts.append(formatted_line.strip())

                                            local_avatar_path, avatar_filename = 获取头像路径(self.手机号, 用户ID)
                                            if not os.path.exists(local_avatar_path):
                                                await self.客户端.download_profile_photo(对话.entity, file=local_avatar_path)

                                            me = await self.客户端.get_me()
                                            self_avatar_path, self_avatar_filename = 获取头像路径(self.手机号, f"me_{me.id}")
                                            if not os.path.exists(self_avatar_path):
                                                await self.客户端.download_profile_photo(me, file=self_avatar_path)

                                            current_msg_id = 对话.message.id
                                            last_reported_id = self.已上报对方消息ID.get(用户ID, 0)
                                            if current_msg_id > last_reported_id:
                                                DB.report_unread(
                                                    server_code=CONFIG.get("pm", {}).get("mysql", {}).get("server_code", "S1"),
                                                    self_name=self.自己名字,
                                                    tg_account=self.手机号,
                                                    peer_id=用户ID,
                                                    username=对话.entity.username or "",
                                                    first_name=f"{对话.entity.first_name or ''} {对话.entity.last_name or ''}".strip(),
                                                    self_avatar=self_avatar_filename,
                                                    avatar_file=avatar_filename,
                                                    messages=msg_texts
                                                )
                                                self.已上报对方消息ID[用户ID] = current_msg_id
                                                with open(self.对方消息ID文件, 'w', encoding='utf-8') as f:
                                                    json.dump(self.已上报对方消息ID, f, ensure_ascii=False, indent=4)
                                                print(f"[{获取时间()}] [{self.手机号}] 📥 已将新消息上报至MySQL (ID: {用户ID}, MsgID: {current_msg_id})")
                                            else:
                                                print(f"[{获取时间()}] [{self.手机号}] ⏭️ 消息 ID ({current_msg_id}) 未更新，跳过上报 (ID: {用户ID})")
                                    except Exception as mysql_err:
                                        print(f"[{获取时间()}] [{self.手机号}] ❌ MySQL处理异常: {mysql_err}")
                        except Exception as e:
                            print(f"[{获取时间()}] [{self.手机号}] 🤖 智能回复异常: {e}")

            # 3. 持久化所有更新
            with open(self.用户映射文件, 'w', encoding='utf-8') as f: json.dump(user_map, f, ensure_ascii=False, indent=4)
            with open(self.头像映射文件, 'w', encoding='utf-8') as f: json.dump(avatar_map, f, ensure_ascii=False, indent=4)
            with open(self.排序文件, 'w', encoding='utf-8') as f: json.dump(sort_data, f, ensure_ascii=False, indent=4)
            
            # 4. 清理已读标记文件
            for f in os.listdir(self.未读目录):
                if f.endswith('.unread'):
                    uid = f.replace('.unread', '')
                    if uid not in 当前真实未读ID:
                        try: os.remove(os.path.join(self.未读目录, f))
                        except: pass
            
            return 有未读
        except Exception as e:
            print(f"[{获取时间()}] [{self.手机号}] [扫描未读] ❌ 异常: {e}")
            return False
    async def 下载并命名媒体(self, 消息, 媒体类型):
        try:
            # 确保 avatars 目录存在
            base_dir = os.path.dirname(os.path.abspath(__file__))
            avatar_dir = os.path.join(base_dir, "avatars")
            if not os.path.exists(avatar_dir):
                os.makedirs(avatar_dir)

            文件后缀 = ""
            if 媒体类型 == "photo":
                文件后缀 = ".jpg"
            elif 媒体类型 == "video":
                文件后缀 = ".mp4"
            elif 媒体类型 == "voice":
                文件后缀 = ".ogg"
            else:
                return None
            
            目标文件名 = f"{self.手机号}_{消息.id}{文件后缀}"
            目标路径 = os.path.join(avatar_dir, 目标文件名)

            if not os.path.exists(目标路径):
                await self.客户端.download_media(消息, file=目标路径)
                return 目标路径
            return 目标路径
        except Exception as e:
            print(f"[{获取时间()}] [{self.手机号}] ❌ 下载并命名媒体异常: {e}")
            return None

    async def 下载媒体(self, 消息):
        try:
            # 此函数现在仅处理文档/文件等其他媒体
            # 图片、语音、视频已由 下载并命名媒体 处理
            if 消息.photo or 消息.voice or 消息.video or 消息.sticker:
                return None

            文件大小 = 消息.document.size if 消息.document else 0
            if 文件大小 > 媒体大小限制: return None
            
            后缀 = ".file"
            if 消息.document and 消息.document.mime_type:
                ext = mimetypes.guess_extension(消息.document.mime_type)
                if ext: 后缀 = ext
            
            路径 = os.path.join(self.媒体目录, f"{消息.id}{后缀}")
            if not os.path.exists(路径):
                await self.客户端.download_media(消息, file=路径)
                return 路径
            return 路径
        except: return None
    async def 下载头像(self, 实体, 文件名):
        try:
            if 实体.photo:
                头像路径 = os.path.join(self.头像目录, 文件名)
                if not os.path.exists(头像路径):
                    await self.客户端.download_profile_photo(实体, file=头像路径)
                    return 文件名
                return 文件名
        except: return None
    async def 执行全量同步(self):
        if os.path.exists(self.全量同步完成标记文件):
            print(f"[{获取时间()}] [{self.手机号}] 检测到本地聊天记录缓存，跳过全量同步。")
            return
        print(f"[{获取时间()}] [{self.手机号}] 未发现本地缓存，开始执行全量同步...")
        头像映射, 排序列表, 用户映射, 媒体映射 = {}, {}, {}, {}
        try:
            我 = await self.客户端.get_me()
            if 我:
                我的名称 = (我.first_name or "") + (" " + 我.last_name if 我.last_name else "")
                我的名称 = 我的名称.strip() or str(我.id)
                更新账号状态(self.手机号, 额外数据={"名称": 我的名称})
                我的头像文件名 = f"me_{我.id}.jpg"
                if await self.下载头像(我, 我的头像文件名): 头像映射['me'] = 我的头像文件名
                用户映射[str(我.id)] = 我的名称

            # 如果媒体上限设置为 0，在下载完个人头像后跳过后续步骤
            if 媒体大小限制 <= 0:
                print(f"[{获取时间()}] [{self.手机号}] 媒体上限设置为 0，已下载个人头像，跳过后续全量同步。")
                with open(self.全量同步完成标记文件, 'w', encoding='utf-8') as f: 
                    f.write(获取中国当前时间().strftime('%Y-%m-%d %H:%M:%S') + " (Self avatar downloaded, others skipped by limit 0)")
                # 即使跳过，也要保存已获取的个人头像映射
                with open(self.头像映射文件, 'w', encoding='utf-8') as f: json.dump(头像映射, f, ensure_ascii=False, indent=4)
                return

            async for 对话 in self.客户端.iter_dialogs():
                if not isinstance(对话.entity, types.User): continue
                用户ID = str(对话.id)
                对话名称 = (对话.entity.first_name or "") + (" " + 对话.entity.last_name if 对话.entity.last_name else "")
                对话名称 = 对话名称.strip() or f"User_{用户ID}"
                用户映射[用户ID] = 对话名称
                头像文件名 = f"{用户ID}.jpg"
                if await self.下载头像(对话.entity, 头像文件名): 头像映射[用户ID] = 头像文件名
                对话文件路径 = os.path.join(self.对话目录, f"{用户ID}.txt")
                with open(对话文件路径, 'w', encoding='utf-8') as f:
                    f.write(f"--- 记录开始: {对话名称} ---\n")
                    async for 消息 in self.客户端.iter_messages(对话.entity, reverse=True):
                        if 消息.media: await self.下载媒体(消息)
                        f.write(await self.格式化消息(消息))
                排序列表[用户ID] = {'id': 用户ID, '名称': 对话名称, '最后时间': 对话.date.timestamp() if 对话.date else time.time()}
            with open(self.头像映射文件, 'w', encoding='utf-8') as f: json.dump(头像映射, f, ensure_ascii=False, indent=4)
            with open(self.排序文件, 'w', encoding='utf-8') as f: json.dump(排序列表, f, ensure_ascii=False, indent=4)
            with open(self.用户映射文件, 'w', encoding='utf-8') as f: json.dump(用户映射, f, ensure_ascii=False, indent=4)
            with open(self.全量同步完成标记文件, 'w', encoding='utf-8') as f: f.write(获取中国当前时间().strftime('%Y-%m-%d %H:%M:%S'))
            print(f"[{获取时间()}] [{self.手机号}] 全量同步完成。")
        except Exception as e:
            print(f"[{获取时间()}] [{self.手机号}] [全量同步] ❌ 异常: {e}")
    async def 监听指令(self):
        while not self.停止信号:
            try:
                target_file = os.path.join(self.交互目录, "target_id.txt")
                input_file = os.path.join(self.交互目录, "input_content.txt")
                cmd_file = os.path.join(self.交互目录, "function_cmd.txt")
                read_file = os.path.join(self.交互目录, "read_cmd.txt")
                sync_cmd_file = os.path.join(self.交互目录, "sync_cmd.json")
                # 处理同步指令 (编辑/删除)
                if os.path.exists(sync_cmd_file) and os.path.getsize(sync_cmd_file) > 0:
                    try:
                        with open(sync_cmd_file, 'r', encoding='utf-8') as f:
                            sync_cmds = json.load(f)
                        for cmd in sync_cmds:
                            action = cmd.get('action')
                            target_id = int(cmd.get('target_id'))
                            msg_id = cmd.get('msg_id')
                            if msg_id is not None: msg_id = int(msg_id)
                            if action == 'delete':
                                print(f"[{获取时间()}] [{self.手机号}] 正在同步删除消息 ID: {msg_id}")
                                await self.客户端.delete_messages(target_id, [msg_id])
                            elif action == 'edit':
                                new_text = cmd.get('content')
                                print(f"[{获取时间()}] [{self.手机号}] 正在同步编辑消息 ID: {msg_id}")
                                await self.客户端.edit_message(target_id, msg_id, new_text)
                            elif action == 'translate':
                                print(f"[{获取时间()}] [{self.手机号}] 收到翻译指令，目标: {target_id}")
                                await self.执行批量翻译(target_id)
                    except Exception as e:
                        print(f"[{获取时间()}] [{self.手机号}] [同步指令] ❌ 异常: {e}")
                    if os.path.exists(sync_cmd_file):
                        os.remove(sync_cmd_file)
                # 处理显式标记已读
                if os.path.exists(read_file) and os.path.getsize(read_file) > 0:
                    with open(read_file, 'r') as f: target_id = f.read().strip()
                    print(f"[{获取时间()}] [{self.手机号}] 收到显式已读指令: {target_id}")
                    await self.同步实时日志(target_id, 发送已读确认=True)
                    os.remove(read_file)
                # 处理目标切换 (明确标记为已读)
                if os.path.exists(target_file) and os.path.getsize(target_file) > 0:
                    with open(target_file, 'r') as f: content = f.read().strip()
                    if content == 'CLEAR':
                        self.当前聊天目标ID = None
                        print(f"[{获取时间()}] [{self.手机号}] 已清除当前聊天目标。")
                    else:
                        self.当前聊天目标ID = content
                        await self.同步实时日志(self.当前聊天目标ID, 发送已读确认=True)
                    os.remove(target_file)
                # 处理发送消息 (明确标记为已读)
                if self.当前聊天目标ID and os.path.exists(input_file) and os.path.getsize(input_file) > 0:
                    with open(input_file, 'r', encoding='utf-8') as f: 内容 = f.read().strip()
                    if 内容:
                        if 内容.startswith("SEND_MEDIA|"):
                            # 处理媒体发送: SEND_MEDIA|类型|路径
                            try:
                                _, m_type, m_path = 内容.split("|", 2)
                                print(f"[{获取时间()}] [{self.手机号}] 正在向 {self.当前聊天目标ID} 发送{m_type}: {os.path.basename(m_path)}")
                                is_voice = (m_type == 'voice')
                                sent_msg = await self.客户端.send_file(
                                    int(self.当前聊天目标ID), 
                                    m_path, 
                                    voice_note=is_voice,
                                    force_document=False
                                )
                                # 发送后尝试删除临时上传文件
                                try: os.remove(m_path)
                                except: pass
                            except Exception as me:
                                print(f"[{获取时间()}] [{self.手机号}] 发送媒体异常: {me}")
                                sent_msg = None
                        else:
                            # 检查是否有回复 ID
                            reply_id = None
                            if 内容.startswith("REPLY_TO:"):
                                parts = 内容.split("\n", 1)
                                if len(parts) > 1:
                                    header = parts[0]
                                    内容 = parts[1]
                                    try:
                                        reply_id = int(header.replace("REPLY_TO:", ""))
                                    except: pass
                            sent_msg = await self.客户端.send_message(int(self.当前聊天目标ID), 内容, reply_to=reply_id)
                        
                        if sent_msg:
                            self.最后交互时间 = time.time()
                            对话文件路径 = os.path.join(self.对话目录, f"{self.当前聊天目标ID}.txt")
                            with open(对话文件路径, 'a', encoding='utf-8') as f:
                                f.write(await self.格式化消息(sent_msg))
                            
                            # 如果是媒体，手动同步一下以便本地预览
                            if sent_msg.media:
                                await self.下载媒体(sent_msg)
                                
                            await self.同步实时日志(self.当前聊天目标ID, 发送已读确认=True)
                    os.remove(input_file)
                # 处理功能指令
                if os.path.exists(cmd_file) and os.path.getsize(cmd_file) > 0:
                    with open(cmd_file, 'r', encoding='utf-8') as f: 指令 = f.read().strip()
                    print(f"[{获取时间()}] [{self.手机号}] 收到指令: {指令}")
                    if 指令 == "立即切换":
                        print(f"[{获取时间()}] [{self.手机号}] 收到强制切换指令，准备下线...")
                        self.停止信号 = True
                    elif 指令 == "延长在线":
                        self.当前观察期 = CONFIG["python"]["long_observation"]
                        self.最后交互时间 = time.time()
                        print(f"[{获取时间()}] [{self.手机号}] 收到延长在线指令，观察期已设为 {self.当前观察期} 秒。")
                    elif 指令 == "全部历史" and self.当前聊天目标ID:
                        # 执行增量全量同步，补全缺失的消息
                        对话文件路径 = os.path.join(self.对话目录, f"{self.当前聊天目标ID}.txt")
                        # 获取本地已有的消息 ID 集合，避免重复
                        已有ID = set()
                        if os.path.exists(对话文件路径):
                            with open(对话文件路径, 'r', encoding='utf-8') as f:
                                for line in f:
                                    m = re.match(r'^\[(\d+)\]', line)
                                    if m: 已有ID.add(int(m.group(1)))
                        # 获取服务器最新消息并补全
                        print(f"[{获取时间()}] [{self.手机号}] 正在同步 {self.当前聊天目标ID} 的完整历史...")
                        async for 消息 in self.客户端.iter_messages(int(self.当前聊天目标ID), limit=None, reverse=True):
                            if 消息.id not in 已有ID:
                                # 历史同步现在统一调用格式化逻辑，它会自动检测本地是否存在媒体文件
                                with open(对话文件路径, 'a', encoding='utf-8') as f:
                                    f.write(await self.格式化消息(消息))
                        print(f"[{获取时间()}] [{self.手机号}] {self.当前聊天目标ID} 历史同步完成。")
                        await self.同步实时日志(self.当前聊天目标ID)
                    elif 指令 == "获取媒体":
                        # 重新扫描并下载媒体，同时替换本地文件中的占位符
                        print(f"[{获取时间()}] [{self.手机号}] 正在获取 {self.当前聊天目标ID} 的媒体文件并补全历史...")
                        对话文件路径 = os.path.join(self.对话目录, f"{self.当前聊天目标ID}.txt")
                        if os.path.exists(对话文件路径):
                            with open(对话文件路径, 'r', encoding='utf-8') as f:
                                lines = f.readlines()
                            
                            modified = False
                            for i in range(len(lines)):
                                # 增加对 "[媒体]" 文本的兼容检查，以便修复旧格式记录
                                if "待同步" in lines[i] or "[媒体]" in lines[i]:
                                    # 提取消息 ID
                                    m = re.match(r'^\[(\d+)\]', lines[i])
                                    if m:
                                        msg_id = int(m.group(1))
                                        try:
                                            msg = await self.客户端.get_messages(int(self.当前聊天目标ID), ids=msg_id)
                                            if msg and msg.media:
                                                # 尝试下载图片、语音或视频
                                                if msg.photo: await self.下载并命名媒体(msg, "photo")
                                                elif msg.voice: await self.下载并命名媒体(msg, "voice")
                                                elif msg.video: await self.下载并命名媒体(msg, "video")
                                                else: await self.下载媒体(msg)
                                                
                                                # 替换为真实的格式化内容（它会自动检测本地已下载的文件名）
                                                lines[i] = await self.格式化消息(msg)
                                                modified = True
                                        except: pass
                            
                            if modified:
                                with open(对话文件路径, 'w', encoding='utf-8') as f:
                                    f.writelines(lines)
                        
                        # 同步最新媒体到实时日志
                        async for 消息 in self.客户端.iter_messages(int(self.当前聊天目标ID), limit=50):
                            if 消息.media: await self.下载媒体(消息)
                        await self.同步实时日志(self.当前聊天目标ID)
                    elif 指令 == "校验同步" and self.当前聊天目标ID:
                        # 强制拉取最新 20 条消息覆盖本地
                        print(f"[{获取时间()}] [{self.手机号}] 正在执行校验同步，拉取最新 20 条消息...")
                        对话文件路径 = os.path.join(self.对话目录, f"{self.当前聊天目标ID}.txt")
                        # 1. 获取服务器上真正的最新 20 条消息 (reverse=False, limit=20 是最新的)
                        最新消息 = []
                        async for 消息 in self.客户端.iter_messages(int(self.当前聊天目标ID), limit=20):
                            最新消息.append(消息)
                        # 翻转一下，使其按时间正序排列
                        最新消息.reverse()
                        最新ID集合 = {m.id for m in 最新消息}
                        最小最新ID = min(最新ID集合) if 最新ID集合 else 0
                        # 2. 读取现有文件内容，保留那些 ID 小于这 20 条中最小 ID 的旧消息
                        新行 = []
                        if os.path.exists(对话文件路径):
                            with open(对话文件路径, 'r', encoding='utf-8') as f:
                                for line in f:
                                    m = re.match(r'^\[(\d+)\]', line)
                                    if m:
                                        msg_id = int(m.group(1))
                                        # 如果本地消息 ID 在最新 20 条的范围内，或者比最小的最新 ID 还大，则剔除（以服务器为准）
                                        if msg_id >= 最小最新ID: continue
                                    新行.append(line)
                        # 3. 追加最新的 20 条
                        for 消息 in 最新消息:
                            新行.append(await self.格式化消息(消息))
                        # 5. 重新写入文件
                        with open(对话文件路径, 'w', encoding='utf-8') as f:
                            f.writelines(新行)
                        print(f"[{获取时间()}] [{self.手机号}] {self.当前聊天目标ID} 校验同步完成。")
                        await self.同步实时日志(self.当前聊天目标ID)
                    elif 指令 == "立即切换":
                        print(f"[{获取时间()}] [{self.手机号}] 收到强制切换指令，准备下线...")
                        self.停止信号 = True
                    os.remove(cmd_file)
            except Exception as e:
                print(f"[{获取时间()}] [{self.手机号}] [监听指令] ❌ 异常: {e}")
            await asyncio.sleep(2)
    async def 运行(self, silent_mode=False):
        self.silent_mode = silent_mode # 保存模式状态供后续逻辑使用
        # --- 启动前清场：清理残留指令文件 ---
        print(f"[{获取时间()}] [{self.手机号}] 正在执行启动前清理... {'(静默唤起模式)' if silent_mode else ''}")
        
        # 模式设置
        if silent_mode:
            self.enable_global_pm = False
            self.enable_auto_reply = False
            self.当前观察期 = CONFIG["python"]["long_observation"]
        else:
            self.enable_global_pm = CONFIG.get("pm", {}).get("enable_global_pm", False)
            self.enable_auto_reply = CONFIG.get("pm", {}).get("enable_auto_reply", False)
            self.当前观察期 = CONFIG["python"]["short_observation"]
        
        # 再次进行二次锁校验，防止极短时间内的竞争
        try:
            self.检查文件锁定()
        except Exception as e:
            print(f"[{获取时间()}] [{self.手机号}] 启动前锁检测失败: {e}")
            更新账号状态(self.手机号, "离线")
            return
            
        print(f"[{获取时间()}] [{self.手机号}] 正在执行启动前清理...")
        for f in ["target_id.txt", "input_content.txt", "function_cmd.txt", "read_cmd.txt", "sync_cmd.json"]:
            path = os.path.join(self.交互目录, f)
            if os.path.exists(path):
                try:
                    os.remove(path)
                    print(f"[{获取时间()}] [{self.手机号}] 已清理残留指令文件: {f}")
                except: pass
        
        # 前置拦截：如果账号已失效或冻结，直接跳过运行
        try:
            if os.path.exists(self.信息文件):
                with open(self.信息文件, 'r', encoding='utf-8') as f:
                    info = json.load(f)
                    健康 = info.get("健康状态", "正常")
                    # 精准拦截永久不可用的状态
                    if any(kw in 健康 for kw in ["失效", "冻结", "banned"]):
                        print(f"[{获取时间()}] [{self.手机号}] ⚠️ 账号处于永久不可用状态 ({健康})，跳过运行。")
                        return
        except: pass
        # --- 调用集成登录子程序 ---
        login_method = str(CONFIG["python"].get("login_method", "1"))
        static_db = str(CONFIG["python"].get("static_ip_db", ""))
        self.客户端 = 电报_登录(self.手机号, login_method, self.账号根目录, static_db)
        if self.客户端 == "-1":
            print(f"[{获取时间()}] [{self.手机号}] ❌ 电报登录集成子程序返回失败")
            更新账号状态(self.手机号, "离线", "登录失败")
            return
        try:
            # --- 第一阶段：连接与授权检查 ---
            try:
                if not self.客户端.is_connected():
                    await self.客户端.connect()
            except Exception as e:
                error_msg = str(e)
                if "deactivated" in error_msg.lower() or "banned" in error_msg.lower():
                    print(f"[{获取时间()}] [{self.手机号}] ❌ 账号已被官方封禁/失效")
                    更新账号状态(self.手机号, "离线", "账号失效")
                    return
                raise e
            if not await self.客户端.is_user_authorized():
                print(f"[{获取时间()}] [{self.手机号}] ❌ 账号未授权/Session失效")
                更新账号状态(self.手机号, "离线", "账号失效")
                return
            
            # 初始登录时同步一次基本信息和会员状态
            try:
                me = await self.客户端.get_me()
                self.自己名字 = f"{me.first_name or ''} {me.last_name or ''}".strip()
                extra_data = {
                    "is_premium": "是" if getattr(me, 'premium', False) else "不是",
                    "名称": self.自己名字
                }
                更新账号状态(self.手机号, "在线", "正常", extra_data)
            except: pass
            # --- 第二阶段：被动深度健康检查 (GetStatuses) ---
            info = {}
            if os.path.exists(self.信息文件):
                try:
                    with open(self.信息文件, 'r', encoding='utf-8') as f:
                        info = json.load(f)
                except: pass
            上次检查时间 = info.get("上次健康检查时间戳", 0)
            check_interval = CONFIG["python"].get("health_check_cooldown_hours", 24) * 3600
            if time.time() - 上次检查时间 > check_interval:
                print(f"[{获取时间()}] [{self.手机号}] 🔍 正在执行定期深度健康检查 (GetStatuses)...")
                try:
                    await self.客户端(functions.contacts.GetStatusesRequest())
                    # 检查成功，更新记录
                    info["上次健康检查时间戳"] = int(time.time())
                    info["健康状态"] = "正常"
                    
                    # --- 同步会员状态 ---
                    try:
                        me = await self.客户端.get_me()
                        info["is_premium"] = "是" if getattr(me, 'premium', False) else "不是"
                        info["名称"] = f"{me.first_name or ''} {me.last_name or ''}".strip()
                    except: pass
                    
                    with open(self.信息文件, 'w', encoding='utf-8') as f:
                        json.dump(info, f, ensure_ascii=False, indent=4)
                    print(f"[{获取时间()}] [{self.手机号}] ✅ 健康检查通过，会员状态: {info.get('is_premium', '未知')}")
                except Exception as e:
                    error_msg = str(e)
                    print(f"[{获取时间()}] [{self.手机号}] ⚠️ 深度检查捕获异常: {error_msg}")
                    # 精准判定冻结状态
                    if 'You tried to use a method that is not available for frozen accounts' in error_msg:
                        print(f"[{获取时间()}] [{self.手机号}] ❄️ 判定结果：账号已冻结")
                        更新账号状态(self.手机号, "离线", "账号冻结")
                        return
                    elif "deactivated" in error_msg.lower():
                        print(f"[{获取时间()}] [{self.手机号}] ❌ 判定结果：账号已失效")
                        更新账号状态(self.手机号, "离线", "账号失效")
                        return
                    else:
                        # 其他异常不标记为永久失效，仅记录
                        info["健康状态"] = f"异常: {error_msg[:30]}"
                        with open(self.信息文件, 'w', encoding='utf-8') as f:
                            json.dump(info, f, ensure_ascii=False, indent=4)
            else:
                print(f"[{获取时间()}] [{self.手机号}] 距离上次深度检查未满 {CONFIG['python']['health_check_cooldown_hours']} 小时，跳过。")
            # --- 第二阶段：按需同步 ---
            await self.执行全量同步()
            # --- 主动私聊子程序 ---
            await self.执行主动私聊()
            # --- 第三阶段：动态在线逻辑 ---
            # 只有在非静默模式下才动态调整观察期，静默模式强制使用启动时设定的长观察期
            if not getattr(self, 'silent_mode', False):
                self.当前观察期 = CONFIG["python"]["short_observation"] # 默认短
                有未读 = await self.扫描未读()
                # 如果扫描未读发现有未读消息，且之前没有因为智能回复而设置为短观察期，则设为长观察期
                # 注意：如果开启了智能回复功能，则强制保持短观察期
                pm_conf = CONFIG.get("pm", {})
                if 有未读 and not pm_conf.get("enable_auto_reply", False):
                    if self.当前观察期 != CONFIG["python"].get("short_observation", 30):
                        self.当前观察期 = CONFIG["python"]["long_observation"]
                elif pm_conf.get("enable_auto_reply", False):
                    self.当前观察期 = CONFIG["python"].get("short_observation", 30)
            else:
                # 静默模式下，只需执行扫描未读以更新本地缓存，但不改变观察期
                有未读 = await self.扫描未读()
            print(f"[{获取时间()}] [{self.手机号}] ● 登录成功{' (静默唤起)' if getattr(self, 'silent_mode', False) else ''}。未读消息: {'有' if 有未读 else '无'}，观察期设定为 {self.当前观察期} 秒。")
            self.最后交互时间 = time.time()
            更新账号状态(self.手机号, "在线")
            @self.客户端.on(events.NewMessage)
            async def 新消息回调(event):
                if not event.is_private: return
                me = await self.客户端.get_me()
                if event.peer_id.user_id == me.id: return
                entity = await event.get_chat()
                if isinstance(entity, types.User) and entity.bot: return
                chat_id = str(event.chat_id)
                self.最后交互时间 = time.time()
                对话文件路径 = os.path.join(self.对话目录, f"{chat_id}.txt")
                # 实时重新加载配置，确保读取到 Web 界面最新的开关状态
                curr_config = 加载配置()
                pm_conf = curr_config.get("pm", {})
                # 检查是否为静默模式或配置已关闭
                is_auto_reply = getattr(self, 'enable_auto_reply', False)
                if is_auto_reply:
                    # 开启智能回复时，强制锁定短观察期
                    self.当前观察期 = curr_config["python"].get("short_observation", 30)
                    print(f"[{获取时间()}] [{self.手机号}] 实时收到消息，智能回复已开启，保持 {self.当前观察期}s 观察期。")
                    # --- 实时触发智能回复：单会话加锁、发送前刷新消息、校验最新消息 ---
                    await self.执行实时智能回复(event, entity, chat_id, 对话文件路径, pm_conf)
                else:
                    # 未开启智能回复时，按原逻辑升级观察期
                    long_obs = curr_config["python"].get("long_observation", 300)
                    if self.当前观察期 < long_obs:
                        self.当前观察期 = long_obs
                        print(f"[{获取时间()}] [{self.手机号}] 收到新消息，观察期动态升级为 {long_obs}s。")
                    else:
                        print(f"[{获取时间()}] [{self.手机号}] 收到新消息，重置 {long_obs}s 观察期。")
                if str(chat_id) != str(self.当前聊天目标ID):
                    # --- 邮件通知触发逻辑 ---
                    if pm_conf.get("enable_email_notification", False):
                        # 检测所有存在的账号目录下是否已有未读消息
                        has_any_unread = False
                        if os.path.exists(数据根目录):
                            for folder in os.listdir(数据根目录):
                                if folder.startswith("account_"):
                                    # 排除已删除但残留的文件夹（通过检查 info.json 是否存在来判定账号是否依然“存在”）
                                    info_path = os.path.join(数据根目录, folder, "info.json")
                                    if os.path.exists(info_path):
                                        unread_dir = os.path.join(数据根目录, folder, "unread_list")
                                        if os.path.exists(unread_dir) and any(f.endswith('.unread') for f in os.listdir(unread_dir)):
                                            has_any_unread = True
                                            break
                        
                        # 如果所有存在的账号都没有未读，才发送邮件
                        if not has_any_unread:
                            print(f"[{获取时间()}] [{self.手机号}] 📩 检测到新消息且所有账号均无未读，触发邮件通知...")
                            发送邮件()
                    
                    with open(os.path.join(self.未读目录, f"{chat_id}.unread"), 'w') as f: f.write("1")
                
                # --- 实时同步用户信息与排序 (修复新私聊不显示在列表的问题) ---
                try:
                    user_id = str(chat_id)
                    user_name = f"{entity.first_name or ''} {entity.last_name or ''}".strip() or user_id
                    
                    # 尝试调用格式化消息来获取准确的最后消息内容
                    formatted_msg = await self.格式化消息(event.message)
                    msg_text = formatted_msg.split(": ", 1)[-1].strip() if ": " in formatted_msg else formatted_msg.strip()
                    
                    # 1. 同步名字
                    user_map = self.加载映射文件(self.用户映射文件)
                    user_map[user_id] = user_name
                    with open(self.用户映射文件, 'w', encoding='utf-8') as f:
                        json.dump(user_map, f, ensure_ascii=False, indent=4)
                    
                    # 2. 同步头像 (如果不存在)
                    avatar_map = self.加载映射文件(self.头像映射文件)
                    avatar_name = f"{user_id}.jpg"
                    avatar_path = os.path.join(self.头像目录, avatar_name)
                    if not os.path.exists(avatar_path):
                        await self.客户端.download_profile_photo(entity, file=avatar_path)
                        if os.path.exists(avatar_path):
                            avatar_map[user_id] = avatar_name
                            with open(self.头像映射文件, 'w', encoding='utf-8') as f:
                                json.dump(avatar_map, f, ensure_ascii=False, indent=4)
                    
                    # 3. 更新排序表
                    sort_data = self.加载映射文件(self.排序文件)
                    sort_data[user_id] = {
                        "名称": user_name,
                        "最后时间": int(time.time()),
                        "最后消息": msg_text[:30]
                    }
                    with open(self.排序文件, 'w', encoding='utf-8') as f:
                        json.dump(sort_data, f, ensure_ascii=False, indent=4)
                except Exception as sync_err:
                    print(f"[{获取时间()}] [{self.手机号}] 实时消息同步用户信息异常: {sync_err}")

                with open(对话文件路径, 'a', encoding='utf-8') as f:
                    f.write(await self.格式化消息(event.message))
                if str(chat_id) == str(self.当前聊天目标ID):
                    if event.message.media:
                        下载路径 = await self.下载媒体(event.message)
                        if 下载路径:
                            self.媒体映射[str(event.message.id)] = 下载路径
                            with open(self.媒体映射文件, 'w', encoding='utf-8') as f: json.dump(self.媒体映射, f, ensure_ascii=False, indent=4)
                    await self.同步实时日志(self.当前聊天目标ID)
            指令任务 = asyncio.create_task(self.监听指令())
            最后心跳时间 = 0
            while not 全局停止信号 and not self.停止信号:
                # 心跳机制：每 10 秒更新一次 info.json，确保 Web 界面显示在线
                if time.time() - 最后心跳时间 > CONFIG["python"]["heartbeat_interval"]:
                    更新账号状态(self.手机号, "在线")
                    最后心跳时间 = time.time()
                # 只有当确实有“发送消息”或“功能指令”时才重置交互时间，普通的目标切换或已读标记不应无限延长寿命
                关键指令存在 = any([os.path.exists(os.path.join(self.交互目录, f)) and os.path.getsize(os.path.join(self.交互目录, f)) > 0 
                               for f in ["input_content.txt", "function_cmd.txt"]])
                if 关键指令存在: self.最后交互时间 = time.time()
                if time.time() - self.最后交互时间 > self.当前观察期:
                    print(f"[{获取时间()}] [{self.手机号}] 超过观察期 ({self.当前观察期}s) 无交互，准备切换。")
                    break
                # --- 观察期内每5秒检查一次MySQL回复 (场景2) ---
                for _ in range(int(CONFIG["python"]["loop_sleep"])):
                    if int(time.time()) % 5 == 0:
                        await self.检查并执行MySQL回复()
                    await asyncio.sleep(1)
            self.停止信号 = True
            指令任务.cancel()
            print(f"[{获取时间()}] [{self.手机号}] <<< 运行结束。")
        except Exception as e:
            print(f"[{获取时间()}] [{self.手机号}] ❌ 运行异常: {e}")
        finally:
            if self.客户端:
                try:
                    await self.客户端.disconnect()
                except: pass
            更新账号状态(self.手机号, "离线")
async def 工人任务(任务队列):
    while not 全局停止信号:
        手机号 = None
        try:
            手机号, Session路径 = await asyncio.wait_for(任务队列.get(), timeout=5)
            # 从排队集合中移除
            if 手机号 in 正在排队的账号: 正在排队的账号.remove(手机号)
            
            if 手机号 in 正在运行的账号:
                任务队列.task_done()
                continue
            正在运行的账号.add(手机号)
            处理器 = 账号处理器(手机号, Session路径)
            await 处理器.运行()
            正在运行的账号.remove(手机号)
            任务队列.task_done()
        except asyncio.TimeoutError: pass
        except Exception as e:
            if 手机号 in 正在运行的账号: 正在运行的账号.remove(手机号)
            if 手机号: 任务队列.task_done()
async def 即时唤起监听():
    """独立协程：扫描 instant_wakeup.flag 并立即启动任务"""
    while not 全局停止信号:
        try:
            # 清理已完成的任务引用
            finished = [t for t in 即时任务集合 if t.done()]
            for t in finished: 即时任务集合.discard(t)
            
            if os.path.exists(数据根目录):
                for folder in os.listdir(数据根目录):
                    if folder.startswith("account_"):
                        flag_path = os.path.join(数据根目录, folder, "instant_wakeup.flag")
                        if os.path.exists(flag_path):
                            手机号 = folder.replace("account_", "")
                            if 手机号 not in 正在运行的账号:
                                print(f"[{获取时间()}] [即时唤起] 检测到指令: {手机号}")
                                # 立即启动异步任务并保持引用
                                task = asyncio.create_task(执行即时任务(手机号))
                                即时任务集合.add(task)
                            # 无论是否成功启动（如果已经在运行则直接删），都清理 flag
                            try: os.remove(flag_path)
                            except: pass
        except Exception as e:
            print(f"[{获取时间()}] [即时唤起监听] 异常: {e}")
        await asyncio.sleep(2)

async def 执行即时任务(手机号):
    """为即时唤起账号运行独立处理器"""
    if 手机号 in 正在运行的账号: return
    正在运行的账号.add(手机号)
    Session路径 = os.path.join(源Session目录, f"{手机号}.session")
    try:
        处理器 = 账号处理器(手机号, Session路径)
        await 处理器.运行(silent_mode=True)
    except Exception as e:
        print(f"[{获取时间()}] [{手机号}] 即时任务异常: {e}")
    finally:
        if 手机号 in 正在运行的账号: 正在运行的账号.remove(手机号)
        # 显式清理
        try: del 处理器
        except: pass

async def 调度中心():
    确保目录存在(源Session目录)
    确保目录存在(数据根目录)
    任务队列 = asyncio.Queue()
    工人们 = [asyncio.create_task(工人任务(任务队列)) for _ in range(最大并发数)]
    # 启动即时唤起监听器
    asyncio.create_task(即时唤起监听())
    
    try:
        while not 全局停止信号:
            Session文件 = sorted([f for f in os.listdir(源Session目录) if f.endswith('.session')])
            for 文件名 in Session文件:
                手机号 = 文件名.replace('.session', '')
                # 如果账号正在运行（包括被即时唤起的），主循环自动跳过
                if 手机号 not in 正在运行的账号 and 手机号 not in 正在排队的账号:
                    正在排队的账号.add(手机号)
                    await 任务队列.put((手机号, os.path.join(源Session目录, 文件名)))
            
            # 等待队列清空或周期性检查
            check_count = 0
            while not 任务队列.empty() and not 全局停止信号 and check_count < 60:
                await asyncio.sleep(1)
                check_count += 1
            await asyncio.sleep(CONFIG["python"]["scheduler_interval"])
    finally:
        for w in 工人们: w.cancel()
if __name__ == "__main__":
    try:
        print(f"[{获取时间()}] 系统启动...")
        # 启动前清理：将所有账号状态初始化为离线，保留健康度记录
        if os.path.exists(数据根目录):
            for folder in os.listdir(数据根目录):
                if folder.startswith("account_"):
                    手机号 = folder.replace("account_", "")
                    # 尝试读取现有健康状态
                    现有健康 = "正常"
                    try:
                        info_path = os.path.join(数据根目录, folder, "info.json")
                        if os.path.exists(info_path):
                            with open(info_path, 'r', encoding='utf-8') as f:
                                现有健康 = json.load(f).get("健康状态", "正常")
                    except: pass
                    更新账号状态(手机号, "离线", 现有健康)
            print(f"[{获取时间()}] 已将所有账号设为离线（已保留健康度记录）。")
            time.sleep(2) # 留出一点时间让状态写入磁盘并被 Web 端感知
        asyncio.run(调度中心())
    except KeyboardInterrupt:
        全局停止信号 = True
    except Exception as e:
        print(f"[{获取时间()}] 主程序异常: {e}")
