Pamudu13/gemma-3-chat
0
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