Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions pytests/webui/test_jargon_routes.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
}
]
Expand Down
5 changes: 5 additions & 0 deletions pytests/webui/test_memory_routes.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down
8 changes: 8 additions & 0 deletions src/chat/heart_flow/heartflow_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
7 changes: 3 additions & 4 deletions src/chat/message_receive/uni_message_sender.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
26 changes: 25 additions & 1 deletion src/common/utils/utils_message.py
Original file line number Diff line number Diff line change
Expand Up @@ -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;后者只能互斥
Expand Down Expand Up @@ -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:
"""在消息入库前补充当前会话的生效回复频率。"""
Expand Down
17 changes: 17 additions & 0 deletions src/config/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
24 changes: 17 additions & 7 deletions src/webui/routers/chat/routes.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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}"
Expand Down Expand Up @@ -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")
Expand Down