diff --git a/pytests/webui/test_jargon_routes.py b/pytests/webui/test_jargon_routes.py index c1962d44be..5c8ebb80b7 100644 --- a/pytests/webui/test_jargon_routes.py +++ b/pytests/webui/test_jargon_routes.py @@ -622,6 +622,7 @@ def test_get_chat_list_includes_chat_session_without_jargon(client: TestClient, "session_id": sample_chat_session.session_id, "chat_name": sample_chat_session.group_name, "platform": sample_chat_session.platform, + "account_id": sample_chat_session.account_id, "is_group": True, } ] diff --git a/pytests/webui/test_memory_routes.py b/pytests/webui/test_memory_routes.py index c5c03c9dba..709ccaa602 100644 --- a/pytests/webui/test_memory_routes.py +++ b/pytests/webui/test_memory_routes.py @@ -638,6 +638,7 @@ def test_webui_memory_timeline_returns_chat_scoped_events(client: TestClient, mo platform="qq", group_id="100", user_id=None, + account_id=None, group_name="测试群", user_cardname=None, user_nickname=None, @@ -708,6 +709,7 @@ def test_webui_memory_timeline_filters_types_and_limit(client: TestClient, monke platform="qq", group_id="100", user_id=None, + account_id=None, group_name="测试群", user_cardname=None, user_nickname=None, @@ -761,6 +763,7 @@ def test_webui_memory_timeline_deleted_paragraph_prefers_delete_operation(client platform="qq", group_id="100", user_id=None, + account_id=None, group_name="测试群", user_cardname=None, user_nickname=None, @@ -791,6 +794,7 @@ def test_webui_memory_timeline_uses_latest_message_snapshot(client: TestClient, platform="qq", group_id=None, user_id="user-1", + account_id=None, group_name=None, user_cardname=None, user_nickname=None, @@ -845,6 +849,7 @@ def test_webui_memory_timeline_handles_json_bytes_zero_timestamp_and_batches_ite platform="qq", group_id="100", user_id=None, + account_id=None, group_name="测试群", user_cardname=None, user_nickname=None, diff --git a/src/chat/heart_flow/heartflow_manager.py b/src/chat/heart_flow/heartflow_manager.py index a31a1215de..6c6d4dedc3 100644 --- a/src/chat/heart_flow/heartflow_manager.py +++ b/src/chat/heart_flow/heartflow_manager.py @@ -91,6 +91,14 @@ async def _evict_chat(self, session_id: str, *, reason: str) -> None: except Exception as exc: logger.warning(f"淘汰心流聊天 {session_id} 失败: {exc}", exc_info=True) + async def release_chat(self, session_id: str, *, reason: str = "released") -> None: + """停止并移除指定会话的心流实例。 + + 供聊天流删除等外部路径复用,与 LRU 淘汰共用同一条 stop 逻辑, + 避免绕过 stop 直接弹出字典导致后台任务泄漏。 + """ + await self._evict_chat(session_id, reason=reason) + def adjust_talk_frequency(self, session_id: str, frequency: float) -> None: """调整指定聊天流的说话频率。""" chat = self.heartflow_chat_list.get(session_id) diff --git a/src/chat/message_receive/uni_message_sender.py b/src/chat/message_receive/uni_message_sender.py index 93bd2c3c5b..29f8b15136 100644 --- a/src/chat/message_receive/uni_message_sender.py +++ b/src/chat/message_receive/uni_message_sender.py @@ -8,7 +8,6 @@ from src.chat.message_receive.message import SessionMessage from src.chat.utils.utils import calculate_typing_time, truncate_message from src.common.data_models.message_component_data_model import ReplyComponent -from src.common.database.database import get_db_session from src.common.logger import get_logger from src.common.message_server.api import get_global_api from src.common.utils.utils_message import MessageUtils @@ -354,9 +353,9 @@ async def send_message( # message.processed_plain_text = modified_message.plain_text if storage_message: - with get_db_session() as db_session: - MessageUtils.fill_reply_frequency_if_available(message) - db_session.add(message.to_db_instance()) + # 必须走 MessageUtils 的串行写入路径:直接同步写库会绕过 + # 进程级写锁并阻塞事件循环,DB 争用时消息已发出但落库失败。 + await MessageUtils.store_sent_message_to_db_async(message) try: from src.services.memory_flow_service import memory_automation_service diff --git a/src/common/utils/utils_message.py b/src/common/utils/utils_message.py index 23af394768..d3a5ea6e22 100644 --- a/src/common/utils/utils_message.py +++ b/src/common/utils/utils_message.py @@ -41,7 +41,7 @@ logger = get_logger("message_utils") -# 串行化 store_message_to_db / update_message_id 的 SQLite 写入: +# 串行化 store_message_to_db / store_sent_message_to_db / update_message_id 的 SQLite 写入: # 底层 SQLite WAL 仅允许单写,busy_timeout 1s。bot 进程不只一个 event loop # (bot.py 主 loop、WebUI 在另一个线程的独立 loop、临时 asyncio.run 调用等), # 因此 lock 必须是进程级的 threading.Lock 而不是 asyncio.Lock;后者只能互斥 @@ -231,6 +231,30 @@ async def store_message_to_db_async(message: "SessionMessage") -> None: """ await asyncio.to_thread(MessageUtils.store_message_to_db, message) + @staticmethod + def store_sent_message_to_db(message: "SessionMessage") -> None: + """存储 bot 已发送的消息到数据库。 + + 与 `store_message_to_db` 的差别是不做图片组件落盘(发送侧组件 + 不携带待持久化的二进制数据);写入同样必须持有 + `_DB_WRITE_THREAD_LOCK`,参见锁注释。 + """ + from src.common.database.database import get_db_session + + with _DB_WRITE_THREAD_LOCK: + with get_db_session() as session: + MessageUtils.fill_reply_frequency_if_available(message) + session.add(message.to_db_instance()) + + @staticmethod + async def store_sent_message_to_db_async(message: "SessionMessage") -> None: + """异步存储 bot 已发送的消息。 + + 把同步 SQLAlchemy session 移出事件循环;锁逻辑在 + `store_sent_message_to_db` 本体里持有,本方法仅做 `to_thread` 透传。 + """ + await asyncio.to_thread(MessageUtils.store_sent_message_to_db, message) + @staticmethod def fill_reply_frequency_if_available(message: "SessionMessage") -> None: """在消息入库前补充当前会话的生效回复频率。""" diff --git a/src/config/config.py b/src/config/config.py index 2c0bd10352..7a917d4b25 100644 --- a/src/config/config.py +++ b/src/config/config.py @@ -9,6 +9,7 @@ from src.common.i18n import t from src.common.logger import get_logger +from src.common.runtime_loop import run_on_main_loop from src.common.version import read_project_version from .config_base import AttributeData, ConfigBase, Field @@ -427,6 +428,22 @@ async def reload_config(self, changed_scopes: Sequence[str] | None = None) -> bo logger.debug("配置热重载未命中有效范围,已跳过") return True + # WebUI 运行在独立线程的第二事件循环上,也会触发热重载;asyncio.Lock + # 只能互斥同一个 loop 内的协程,跨 loop 争用时要么抛 RuntimeError 要么 + # 完全失去互斥。统一投递到主循环执行,保证重载与回调串行化,且回调 + # (如对主循环 asyncio.Event 的 set)总是在其所属的事件循环上运行。 + return await run_on_main_loop(self._reload_config_on_main_loop(normalized_scopes)) + + async def _reload_config_on_main_loop(self, normalized_scopes: tuple[str, ...]) -> bool: + """在主循环上串行执行配置热重载。 + + Args: + normalized_scopes: 已规范化的配置变更范围。 + + Returns: + bool: 是否重载成功。 + """ + async with self._reload_lock: try: global_config_new = self.global_config diff --git a/src/webui/routers/chat/routes.py b/src/webui/routers/chat/routes.py index 1249810cb7..ac17b44dd4 100644 --- a/src/webui/routers/chat/routes.py +++ b/src/webui/routers/chat/routes.py @@ -29,6 +29,7 @@ ToolRecord, ) from src.common.logger import get_logger +from src.common.runtime_loop import run_on_main_loop from src.common.utils.utils_config import ( BehaviorConfigUtils, ChatConfigUtils, @@ -1097,14 +1098,23 @@ def _delete_or_unlink_jargons(session: Any, session_id: str) -> Dict[str, int]: } -def _release_deleted_chat_runtime(session_id: str) -> None: - """移除运行期缓存,避免定时保存把已删除聊天流重新写回数据库。""" +async def _release_deleted_chat_runtime(session_id: str) -> None: + """在主循环上停止心流实例并移除运行期缓存。 - core_chat_manager.sessions.pop(session_id, None) - heartflow_manager.heartflow_chat_list.pop(session_id, None) + 一是避免定时保存把已删除聊天流重新写回数据库;二是这些对象归主循环 + 持有,直接在 WebUI 线程弹出会留下仍在运行的后台任务(幽灵 runtime), + 并与主循环的字典写入竞态,因此统一投递到主循环,并复用 + heartflow_manager 的淘汰路径完成 stop。 + """ + + async def _stop_and_release() -> None: + core_chat_manager.sessions.pop(session_id, None) + await heartflow_manager.release_chat(session_id, reason="webui_delete") + + await run_on_main_loop(_stop_and_release()) -def _delete_chat_session_scope(session_id: str) -> Dict[str, Any]: +async def _delete_chat_session_scope(session_id: str) -> Dict[str, Any]: """删除聊天流及所有直接归属该 session_id 的数据库记录。""" with get_db_session() as session: @@ -1132,7 +1142,7 @@ def _delete_chat_session_scope(session_id: str) -> Dict[str, Any]: total_deleted += deleted_count items.append({"key": key, "label": label, "count": deleted_count}) - _release_deleted_chat_runtime(session_id) + await _release_deleted_chat_runtime(session_id) logger.warning( "已删除聊天流及关联数据: " f"session_id={session_id} total_deleted={total_deleted} items={items}" @@ -1313,7 +1323,7 @@ async def delete_chat_session(session_id: str) -> Dict[str, object]: if not normalized_session_id: raise HTTPException(status_code=400, detail="缺少聊天流 session_id") - return _delete_chat_session_scope(normalized_session_id) + return await _delete_chat_session_scope(normalized_session_id) @router.put("/sessions/{session_id}/talk-frequency")