CoolFace
Apppublic

Pamudu13/gemma-3-chat

sourceHugging Faceupdated 25d agoView on Hugging Face
0likes
feishu_bot_manager.py602 linesDownload Raw Back to py
1# feishu_bot_manager.py2import asyncio3import json4import random5import threading6from typing import Optional, List, Dict7import weakref8import aiohttp9import io10import base6411import logging12import re13from pydantic import BaseModel, Field14from openai import AsyncOpenAI15 16from py.get_setting import convert_to_opus_simple, get_port, load_settings17from py.behavior_engine import BehaviorItem, global_behavior_engine, BehaviorSettings18 19# 飞书机器人配置模型20class FeishuBotConfig(BaseModel):21    FeishuAgent: str          # LLM模型名22    memoryLimit: int          # 记忆条数限制23    appid: str                # 飞书APP_ID24    secret: str               # 飞书APP_SECRET25    separators: List[str]     # 消息分段符26    reasoningVisible: bool    # 是否显示推理过程27    quickRestart: bool        # 快速重启指令开关28    enableTTS: bool           # 是否启用TTS29    wakeWord: str             # 唤醒词30    behaviorSettings: Optional[BehaviorSettings] = None31    behaviorTargetChatIds: List[str] = Field(default_factory=list)32 33class FeishuBotManager:34    def __init__(self):35        self.bot_thread: Optional[threading.Thread] = None36        self.bot_client: Optional[FeishuClient] = None37        self.is_running = False38        self.config = None39        self.loop = None40        self._shutdown_event = threading.Event()41        self._startup_complete = threading.Event()42        self._ready_complete = threading.Event()43        self._startup_error = None44        self.ws = None  45        self._stop_requested = False  46        47    def start_bot(self, config):48        if self.is_running:49            raise Exception("飞书机器人已在运行")50            51        self.config = config52        self._shutdown_event.clear()53        self._startup_complete.clear()54        self._ready_complete.clear()55        self._startup_error = None56        self._stop_requested = False57        58        self.bot_thread = threading.Thread(59            target=self._run_bot_thread,60            args=(config,),61            daemon=True,62            name="FeishuBotThread"63        )64        self.bot_thread.start()65        66        if not self._startup_complete.wait(timeout=30):67            self.stop_bot()68            raise Exception("飞书机器人连接超时")69            70        if self._startup_error:71            self.stop_bot()72            raise Exception(f"飞书机器人启动失败: {self._startup_error}")73        74        if not self._ready_complete.wait(timeout=30):75            self.stop_bot()76            raise Exception("飞书机器人就绪超时,请检查网络连接和配置")77            78        if not self.is_running:79            self.stop_bot()80            raise Exception("飞书机器人未能正常运行")81    82    def _run_bot_thread(self, config):83        self.loop = None84        try:85            self.loop = asyncio.new_event_loop()86            asyncio.set_event_loop(self.loop)87            88            self.bot_client = FeishuClient()89            self.bot_client.FeishuAgent = config.FeishuAgent90            self.bot_client.memoryLimit = config.memoryLimit91            self.bot_client.separators = config.separators if config.separators else []92            self.bot_client.reasoningVisible = config.reasoningVisible93            self.bot_client.quickRestart = config.quickRestart94            self.bot_client.appid = config.appid95            self.bot_client.secret = config.secret96            self.bot_client.enableTTS = config.enableTTS97            self.bot_client.wakeWord = config.wakeWord98            self.bot_client._manager_ref = weakref.ref(self)99            self.bot_client._ready_callback = self._on_bot_ready100 101            try:102                settings = asyncio.run(load_settings())103                behavior_data = settings.get("behaviorSettings", {})104                target_ids = config.behaviorTargetChatIds or settings.get("feishuBotConfig", {}).get("behaviorTargetChatIds", [])105                106                if behavior_data:107                    logging.info(f"飞书同步行为配置中... 目标: {len(target_ids)}")108                    global_behavior_engine.update_config(behavior_data, {"feishu": target_ids})109            except Exception as e:110                logging.error(f"飞书行为配置同步失败: {e}")111 112            import lark_oapi as lark113            lark_client = lark.Client.builder().app_id(config.appid).app_secret(config.secret).build()114            self.bot_client.lark_client = lark_client115            116            event_dispatcher = lark.EventDispatcherHandler.builder("", "")\117                .register_p2_im_message_receive_v1(self.bot_client.sync_handle_message)\118                .build()119                120            self.ws = lark.ws.Client(121                config.appid, 122                config.secret,123                event_handler=event_dispatcher,124                auto_reconnect=False125            )126            127            self.loop.run_until_complete(self._async_run_websocket())128            129        except Exception as e:130            if not self._stop_requested:131                self._startup_error = str(e)132            if not self._startup_complete.is_set(): self._startup_complete.set()133            if not self._ready_complete.is_set(): self._ready_complete.set()134        finally:135            self._cleanup()  136 137    async def _async_run_websocket(self):138        try:139            await self.ws._connect()140            self._startup_complete.set()141            self._ready_complete.set()142            self.is_running = True143            144            ping_task = asyncio.create_task(self.ws._ping_loop())145            receive_task = asyncio.create_task(self._message_receive_loop())146            147            if global_behavior_engine.is_running:148                global_behavior_engine.stop()149                await asyncio.sleep(0.5)150 151            behavior_task = asyncio.create_task(global_behavior_engine.start())152            await asyncio.gather(ping_task, receive_task, behavior_task, return_exceptions=True)153        except Exception as e:154            if not self._stop_requested: self._startup_error = str(e)155            raise156 157    async def _message_receive_loop(self):158        while not self._stop_requested and not self._shutdown_event.is_set():159            if self.ws._conn is None: break160            try:161                msg = await asyncio.wait_for(self.ws._conn.recv(), timeout=1.0)162                asyncio.create_task(self.ws._handle_message(msg))163            except asyncio.TimeoutError: continue164            except: break165    166    def _on_bot_ready(self):167        self.is_running = True168        self._ready_complete.set()169 170    def _cleanup(self):171        self.is_running = False172        if global_behavior_engine.is_running: global_behavior_engine.stop()173        if self.ws and self.loop and not self.loop.is_closed():174            try:175                if asyncio.iscoroutinefunction(self.ws._disconnect):176                    self.loop.run_until_complete(self.ws._disconnect())177            except: pass178        if self.loop and not self.loop.is_closed():179            try:180                for task in asyncio.all_tasks(self.loop): task.cancel()181                self.loop.close()182            except: pass183        self._shutdown_event.set()184 185    def stop_bot(self):186        if not self.is_running and not self.bot_thread: return187        self._stop_requested = True188        self._shutdown_event.set()189        self.is_running = False190        191        # 取消所有正在进行的对话任务192        if self.bot_client:193            for task in self.bot_client.active_tasks.values():194                task.cancel()195 196        if self.loop and not self.loop.is_closed():197            try:198                if self.ws and hasattr(self.ws, '_disconnect'):199                    asyncio.run_coroutine_threadsafe(self.ws._disconnect(), self.loop)200            except: pass201        202        if self.bot_thread and self.bot_thread.is_alive():203            self.bot_thread.join(timeout=5)204        self._stop_requested = False205 206    def get_status(self):207        return {208            "is_running": self.is_running,209            "thread_alive": self.bot_thread.is_alive() if self.bot_thread else False,210            "config": self.config.model_dump() if self.config else None,211            "startup_error": self._startup_error212        }213 214    def update_behavior_config(self, config: FeishuBotConfig):215        self.config = config216        if self.bot_client:217            self.bot_client.FeishuAgent = config.FeishuAgent 218            self.bot_client.enableTTS = config.enableTTS219            self.bot_client.wakeWord = config.wakeWord220        global_behavior_engine.update_config(config.behaviorSettings, {"feishu": config.behaviorTargetChatIds})221 222class FeishuClient:223    def __init__(self):224        self.FeishuAgent = "super-model"225        self.memoryLimit = 10226        self.memoryList = {}227        self.asyncToolsID = {}228        self.fileLinks = {}229        self.separators = []230        self.reasoningVisible = False231        self.quickRestart = True232        self._is_ready = False233        self.appid = None234        self.secret = None235        self.lark_client = None236        self.port = get_port()237        self._shutdown_requested = False238        self._manager_ref = None239        self._ready_callback = None240        self.enableTTS = False241        self.wakeWord = None242        # 核心:追踪当前任务,实现打断功能243        self.active_tasks: Dict[str, asyncio.Task] = {}244        global_behavior_engine.register_handler("feishu", self.execute_behavior_event)245        246    def sync_handle_message(self, data) -> None:247        if self._shutdown_requested: return248        if self._manager_ref and self._manager_ref()._stop_requested: return249        250        try:251            loop = asyncio.get_event_loop()252            if loop.is_closed(): return253            # 在协程中处理消息以便支持 cancel254            asyncio.run_coroutine_threadsafe(self.handle_message(data), loop)255        except Exception as e:256            logging.error(f"同步处理异常: {e}")257 258    async def handle_message(self, data) -> None:259        """接收并调度消息处理任务(包含打断逻辑)"""260        if self._shutdown_requested: return261        if not self._is_ready:262            self._is_ready = True263            if self._ready_callback: self._ready_callback()264        265        chat_id = data.event.message.chat_id266        msg_type = data.event.message.message_type267        global_behavior_engine.report_activity("feishu", chat_id)268 269        # 1. 快捷指令与任务打断检查270        if self.quickRestart and msg_type == "text":271            try:272                text = json.loads(data.event.message.content).get("text", "").strip()273                # 停止逻辑274                if text in ["/停止", "/stop"]:275                    if chat_id in self.active_tasks:276                        self.active_tasks[chat_id].cancel()277                        await self._send_text(data.event.message, "已停止当前输出。")278                    return279                # 重启逻辑280                if text in ["/重启", "/restart"]:281                    if chat_id in self.active_tasks:282                        self.active_tasks[chat_id].cancel()283                    self.memoryList[chat_id] = []284                    await self._send_text(data.event.message, "对话记录已重置。")285                    return286            except: pass287 288        # 2. 如果当前有任务正在运行,直接打断289        if chat_id in self.active_tasks:290            logging.info(f"检测到新消息,打断会话 {chat_id} 的旧任务")291            self.active_tasks[chat_id].cancel()292 293        # 3. 创建处理任务并记录294        current_task = asyncio.create_task(self._do_handle_message(data))295        self.active_tasks[chat_id] = current_task296        297        try:298            await current_task299        except asyncio.CancelledError:300            logging.info(f"会话 {chat_id} 的旧任务已安全退出")301        finally:302            if self.active_tasks.get(chat_id) == current_task:303                self.active_tasks.pop(chat_id, None)304 305    async def _do_handle_message(self, data) -> None:306        """实际的消息处理逻辑(原 handle_message 的全部内容)"""307        msg = data.event.message308        msg_type = msg.message_type309        chat_id = msg.chat_id310        311        client = AsyncOpenAI(api_key="super-secret-key", base_url=f"http://127.0.0.1:{self.port}/v1")312        settings = await load_settings()313 314        if chat_id not in self.memoryList: self.memoryList[chat_id] = []315            316        user_content = []317        user_text = ""318        has_image = False319        320        # --- 解析逻辑保持不变 ---321        if msg_type == "text":322            text = json.loads(msg.content).get("text", "")323            if "/id" in text.lower():324                await self._send_text(msg, f"🤖 **会话信息**\n\nChatID:\n`{chat_id}`")325                return326            user_text = text327            if self.wakeWord and self.wakeWord not in user_text: return328        elif msg_type == "image":329            image_key = json.loads(msg.content).get("image_key", "")330            if image_key:331                from lark_oapi.api.im.v1 import GetMessageResourceRequest as ResReq332                res_resp = self.lark_client.im.v1.message_resource.get(ResReq.builder().message_id(msg.message_id).file_key(image_key).type("image").build())333                if res_resp.success():334                    img_bin = res_resp.file.read()335                    base64_data = base64.b64encode(img_bin).decode("utf-8")336                    has_image = True337                    user_content.append({"type": "image_url", "image_url": {"url": f"data:image/jpeg;base64,{base64_data}"}})338        elif msg_type == "post":339            content_json = json.loads(msg.content)340            user_text = self._extract_text_from_post(content_json)341            for image_key in self._extract_images_from_post(content_json):342                from lark_oapi.api.im.v1 import GetMessageResourceRequest as ResReq343                res_resp = self.lark_client.im.v1.message_resource.get(ResReq.builder().message_id(msg.message_id).file_key(image_key).type("image").build())344                if res_resp.success():345                    img_bin = res_resp.file.read()346                    base64_data = base64.b64encode(img_bin).decode("utf-8")347                    has_image = True348                    user_content.append({"type": "image_url", "image_url": {"url": f"data:image/jpeg;base64,{base64_data}"}})349        elif msg_type == "audio":350            file_key = json.loads(msg.content).get("file_key", "")351            if file_key:352                from lark_oapi.api.im.v1 import GetMessageResourceRequest as ResReq353                res_resp = self.lark_client.im.v1.message_resource.get(ResReq.builder().message_id(msg.message_id).file_key(file_key).type("file").build())354                if res_resp.success():355                    user_text = await self._transcribe_audio(res_resp.file.read(), file_key)356                    if not user_text or (self.wakeWord and self.wakeWord not in user_text): return357        else: return358 359        if has_image:360            if user_text: user_content.append({"type": "text", "text": user_text})361            self.memoryList[chat_id].append({"role": "user", "content": user_content})362        else:363            if user_text: self.memoryList[chat_id].append({"role": "user", "content": user_text})364            else: return365 366        # AI 请求逻辑367        state = {"text_buffer": "", "image_buffer": "", "image_cache": [], "audio_buffer": []}368        try:369            asyncToolsID = self.asyncToolsID.setdefault(chat_id, [])370            fileLinks = self.fileLinks.setdefault(chat_id, [])371            372            stream = await client.chat.completions.create(373                model=self.FeishuAgent,374                messages=self.memoryList[chat_id],375                stream=True,376                extra_body={"asyncToolsID": asyncToolsID, "fileLinks": fileLinks, "is_app_bot": True, "platform": "feishu"}377            )378            379            full_response = []380            async for chunk in stream:381                if chunk.choices:382                    delta = chunk.choices[0].delta383                    if hasattr(delta, "audio") and delta.audio and "data" in delta.audio:384                        state["audio_buffer"].append(delta.audio["data"])385                    if hasattr(delta, "async_tool_id") and delta.async_tool_id:386                        tid = delta.async_tool_id387                        if tid not in self.asyncToolsID[chat_id]: self.asyncToolsID[chat_id].append(tid)388                        else: self.asyncToolsID[chat_id].remove(tid)389                    if hasattr(delta, "tool_link") and delta.tool_link:390                        if settings["tools"]["toolMemorandum"]["enabled"]: self.fileLinks[chat_id].append(delta.tool_link)391 392                    content = delta.content or ""393                    if self.reasoningVisible and hasattr(delta, "reasoning_content") and delta.reasoning_content:394                        content = delta.reasoning_content395                    396                    full_response.append(delta.content or "")397                    state["text_buffer"] += content398                    state["image_buffer"] += content399                    400                    # 分段发送401                    if state["text_buffer"]:402                        buffer = state["text_buffer"]403                        split_pos = -1404                        for sep in self.separators:405                            pos = buffer.find(sep)406                            if pos != -1:407                                split_pos = pos + len(sep)408                                break409                        if split_pos != -1:410                            chunk_to_send = buffer[:split_pos]411                            state["text_buffer"] = buffer[split_pos:]412                            clean = self._clean_text(chunk_to_send)413                            if clean: await self._send_text(msg, clean)414            415            # 处理收尾416            self._extract_images(state)417            if state["text_buffer"]:418                clean = self._clean_text(state["text_buffer"])419                if clean: await self._send_text(msg, clean)420            for img_url in state["image_cache"]: await self._send_image(msg, img_url)421            422            # Omni音频处理423            has_omni = False424            if state["audio_buffer"]:425                final_audio, is_opus = await asyncio.to_thread(convert_to_opus_simple, base64.b64decode("".join(state["audio_buffer"])))426                await self._send_omni_response(msg, final_audio, is_opus)427                has_omni = True428 429            full_content = "".join(full_response)430            if self.enableTTS and not has_omni: await self._send_voice(msg, full_content)431            self.memoryList[chat_id].append({"role": "assistant", "content": full_content})432            433            if self.memoryLimit > 0:434                while len(self.memoryList[chat_id]) > self.memoryLimit * 2:435                    self.memoryList[chat_id].pop(0)436 437        except Exception as e:438            if not isinstance(e, asyncio.CancelledError):439                logging.error(f"对话异常: {e}")440                await self._send_text(msg, f"对话中断: {str(e)}")441 442    # --- 后续所有辅助工具方法均保持原样(Omni, TTS, ASR, Upload等) ---443 444    async def _send_omni_response(self, original_msg, audio_data: bytes, is_opus: bool):445        try:446            file_type = "opus" if is_opus else "wav"447            file_name = f"reply.{file_type}"448            msg_type = "audio" if is_opus else "file"449            450            from lark_oapi.api.im.v1 import CreateFileRequest, CreateFileRequestBody451            upload_resp = self.lark_client.im.v1.file.create(CreateFileRequest.builder().request_body(CreateFileRequestBody.builder().file_type(file_type).file_name(file_name).file(io.BytesIO(audio_data)).build()).build())452            if not upload_resp.success(): return453            454            content_str = json.dumps({"file_key": upload_resp.data.file_key})455            from lark_oapi.api.im.v1 import CreateMessageRequest, CreateMessageRequestBody, ReplyMessageRequest, ReplyMessageRequestBody456            if original_msg.chat_type == "p2p":457                req = CreateMessageRequest.builder().receive_id_type("chat_id").request_body(CreateMessageRequestBody.builder().receive_id(original_msg.chat_id).msg_type(msg_type).content(content_str).build()).build()458                self.lark_client.im.v1.message.create(req)459            else:460                req = ReplyMessageRequest.builder().message_id(original_msg.message_id).request_body(ReplyMessageRequestBody.builder().msg_type(msg_type).content(content_str).build()).build()461                self.lark_client.im.v1.message.reply(req)462        except: pass463 464    async def _transcribe_audio(self, audio_data: bytes, file_key: str) -> str:465        try:466            form_data = aiohttp.FormData()467            form_data.add_field('audio', io.BytesIO(audio_data), filename=f"{file_key}.ogg", content_type='audio/ogg')468            form_data.add_field('format', 'auto')469            async with aiohttp.ClientSession() as session:470                async with session.post(f"http://127.0.0.1:{self.port}/asr", data=form_data, timeout=60) as resp:471                    if resp.status == 200:472                        res = await resp.json()473                        return res.get("text", "").strip() if res.get("success") else None474        except: return None475 476    def clean_markdown(self, buffer):477        buffer = re.sub(r'#{1,6}\s', '', buffer, flags=re.MULTILINE)478        buffer = re.sub(r'[*_~`]+', '', buffer)479        buffer = re.sub(r'^\s*[-*]\s', '', buffer, flags=re.MULTILINE)480        buffer = re.sub(r'[\u2600-\u27BF\U0001F300-\U0001F9FF]', '', buffer)481        buffer = re.sub(r'!\[.*?\]\(.*?\)', '', buffer)482        buffer = re.sub(r'\[(.*?)\]\(.*?\)', r'\1', buffer)483        return buffer.strip()484 485    async def _send_voice(self, original_msg, text):486        try:487            settings = await load_settings()488            clean_t = self.clean_markdown(text)489            payload = {"text": clean_t, "voice": "default", "ttsSettings": settings.get("ttsSettings", {}), "index": 0, "mobile_optimized": True, "format": "opus"}490            async with aiohttp.ClientSession() as session:491                async with session.post(f"http://127.0.0.1:{self.port}/tts", json=payload, timeout=90) as resp:492                    if resp.status != 200: return493                    opus_data = await resp.read()494                    from lark_oapi.api.im.v1 import CreateFileRequest, CreateFileRequestBody495                    upload_resp = self.lark_client.im.v1.file.create(CreateFileRequest.builder().request_body(CreateFileRequestBody.builder().file_type("opus").file_name("v.opus").file(io.BytesIO(opus_data)).build()).build())496                    if not upload_resp.success(): return497                    key = upload_resp.data.file_key498                    from lark_oapi.api.im.v1 import CreateMessageRequest, CreateMessageRequestBody, ReplyMessageRequest, ReplyMessageRequestBody499                    if original_msg.chat_type == "p2p":500                        self.lark_client.im.v1.message.create(CreateMessageRequest.builder().receive_id_type("chat_id").request_body(CreateMessageRequestBody.builder().receive_id(original_msg.chat_id).msg_type("audio").content(json.dumps({"file_key": key})).build()).build())501                    else:502                        self.lark_client.im.v1.message.reply(ReplyMessageRequest.builder().message_id(original_msg.message_id).request_body(ReplyMessageRequestBody.builder().msg_type("audio").content(json.dumps({"file_key": key})).build()).build())503        except: pass504 505    def _extract_text_from_post(self, post_content):506        res = []507        try:508            if isinstance(post_content, dict):509                if post_content.get("title"): res.append(post_content["title"])510                if "content" in post_content:511                    for para in post_content["content"]:512                        for el in para:513                            if el.get("tag") in ["text", "a"]: res.append(el.get("text", ""))514                            elif el.get("tag") == "at": res.append(f"@{el.get('user_name', '')}")515        except: pass516        return "\n".join(res)517 518    def _extract_images_from_post(self, post_content):519        keys = []520        try:521            if isinstance(post_content, dict) and "content" in post_content:522                for para in post_content["content"]:523                    for el in para:524                        if el.get("tag") in ["img", "media"] and el.get("image_key"): keys.append(el["image_key"])525        except: pass526        return keys527 528    def _extract_images(self, state):529        pattern = r'!\[.*?\]\((https?://[^\s\)]+)'530        for match in re.finditer(pattern, state["image_buffer"]): state["image_cache"].append(match.group(1))531    532    def _clean_text(self, text: str) -> str:533        text = re.sub(r"!\[.*?\]\(.*?\)", "", text)534        text = re.sub(r'<.*?>', '', text)535        return text.strip()536    537    async def _send_text(self, original_msg, text):538        try:539            if not text: return540            content = json.dumps({"zh_cn": {"content": [[{"tag": "md", "text": text}]]}})541            from lark_oapi.api.im.v1 import CreateMessageRequest, CreateMessageRequestBody, ReplyMessageRequest, ReplyMessageRequestBody542            if original_msg.chat_type == "p2p":543                req = CreateMessageRequest.builder().receive_id_type("chat_id").request_body(CreateMessageRequestBody.builder().receive_id(original_msg.chat_id).msg_type("post").content(content).build()).build()544                self.lark_client.im.v1.message.create(req)545            else:546                req = ReplyMessageRequest.builder().message_id(original_msg.message_id).request_body(ReplyMessageRequestBody.builder().msg_type("post").content(content).build()).build()547                self.lark_client.im.v1.message.reply(req)548        except: pass549                550    async def _send_image(self, original_msg, image_url):551        try:552            async with aiohttp.ClientSession() as session:553                async with session.get(image_url) as r:554                    if r.status != 200: return555                    data = await r.read()556            from lark_oapi.api.im.v1 import CreateImageRequest, CreateImageRequestBody557            up_resp = self.lark_client.im.v1.image.create(CreateImageRequest.builder().request_body(CreateImageRequestBody.builder().image_type("message").image(io.BytesIO(data)).build()).build())558            if not up_resp.success(): return559            key = up_resp.data.image_key560            from lark_oapi.api.im.v1 import CreateMessageRequest, CreateMessageRequestBody, ReplyMessageRequest, ReplyMessageRequestBody561            if original_msg.chat_type == "p2p":562                self.lark_client.im.v1.message.create(CreateMessageRequest.builder().receive_id_type("chat_id").request_body(CreateMessageRequestBody.builder().receive_id(original_msg.chat_id).msg_type("image").content(json.dumps({"image_key": key})).build()).build())563            else:564                self.lark_client.im.v1.message.reply(ReplyMessageRequest.builder().message_id(original_msg.message_id).request_body(ReplyMessageRequestBody.builder().msg_type("image").content(json.dumps({"image_key": key})).build()).build())565        except: pass566 567    async def execute_behavior_event(self, chat_id: str, behavior_item: BehaviorItem):568        prompt = await self._resolve_behavior_prompt(behavior_item)569        if not prompt: return570        class Mock:571            def __init__(self, cid): self.chat_id = cid; self.message_id = None; self.chat_type = "p2p"572        mock = Mock(chat_id)573        if chat_id not in self.memoryList: self.memoryList[chat_id] = []574        messages = self.memoryList[chat_id].copy()575        messages.append({"role": "user", "content": f"[system]: {prompt}"})576        try:577            client = AsyncOpenAI(api_key="super-secret-key", base_url=f"http://127.0.0.1:{self.port}/v1")578            resp = await client.chat.completions.create(model=self.FeishuAgent, messages=messages, stream=False, extra_body={"is_app_bot": True, "platform": "feishu", "behavior_trigger": True})579            reply = resp.choices[0].message.content580            if reply:581                await self._send_text(mock, reply)582                self.memoryList[chat_id].append({"role": "user", "content": f"[system]: {prompt}"})583                self.memoryList[chat_id].append({"role": "assistant", "content": reply})584                if self.enableTTS: await self._send_voice(mock, reply)585        except: pass586 587    async def _resolve_behavior_prompt(self, behavior: BehaviorItem) -> str:588        action = behavior.action589        if action.type == "prompt": return action.prompt590        elif action.type == "random" and action.random and action.random.events:591            events = action.random.events592            if action.random.type == "random": return random.choice(events)593            idx = action.random.orderIndex594            selected = events[idx % len(events)]595            action.random.orderIndex = (idx + 1) % len(events)596            return selected597        return None            598 599    def __del__(self):600        try:601            for t in self.active_tasks.values(): t.cancel()602        except: pass