From 628f1b53fb27ffa64b5025af8b9a909c08636a16 Mon Sep 17 00:00:00 2001
From: kong <135376906+3316891527@users.noreply.github.com>
Date: Wed, 9 Sep 2026 17:00:40 +0000
Subject: [PATCH 1/2] fix(chat): move long-history work off main thread
---
.../operit/core/chat/AIMessageManager.kt | 77 +++++-----
.../operit/services/ChatServiceCore.kt | 2 +-
.../core/MessageCoordinationDelegate.kt | 145 ++++++++++++------
.../features/chat/viewmodel/ChatViewModel.kt | 59 +++----
4 files changed, 164 insertions(+), 119 deletions(-)
diff --git a/app/src/main/java/com/ai/assistance/operit/core/chat/AIMessageManager.kt b/app/src/main/java/com/ai/assistance/operit/core/chat/AIMessageManager.kt
index 1f961f22a2..388b2a49cf 100644
--- a/app/src/main/java/com/ai/assistance/operit/core/chat/AIMessageManager.kt
+++ b/app/src/main/java/com/ai/assistance/operit/core/chat/AIMessageManager.kt
@@ -376,28 +376,28 @@ object AIMessageManager {
lastActiveChatKey = chatKey
activeEnhancedAiServiceByChatId[chatKey] = enhancedAiService
- val buildMemoryStartTime = messageTimingNow()
- val memory = getMemoryFromMessages(
- messages = chatHistory,
- splitByRole = splitHistoryByRole,
- targetRoleName = currentRoleName,
- groupOrchestrationMode = groupOrchestrationMode
- )
- logMessageTiming(
- stage = "sendMessage.buildMemory",
- startTimeMs = buildMemoryStartTime,
- details = "chatKey=$chatKey, source=${chatHistory.size}, result=${memory.size}, splitByRole=$splitHistoryByRole, groupOrchestration=$groupOrchestrationMode"
- )
- if (splitHistoryByRole && !currentRoleName.isNullOrBlank()) {
- val assistantCount = memory.count { it.kind == PromptTurnKind.ASSISTANT }
- val userCount = memory.count { it.kind == PromptTurnKind.USER }
- AppLogger.d(
- TAG,
- "按角色拆解历史: role=$currentRoleName, assistant=$assistantCount, user=$userCount, total=${memory.size}"
+ return withContext(Dispatchers.IO) {
+ val buildMemoryStartTime = messageTimingNow()
+ val memory = getMemoryFromMessages(
+ messages = chatHistory,
+ splitByRole = splitHistoryByRole,
+ targetRoleName = currentRoleName,
+ groupOrchestrationMode = groupOrchestrationMode
)
- }
+ logMessageTiming(
+ stage = "sendMessage.buildMemory",
+ startTimeMs = buildMemoryStartTime,
+ details = "chatKey=$chatKey, source=${chatHistory.size}, result=${memory.size}, splitByRole=$splitHistoryByRole, groupOrchestration=$groupOrchestrationMode"
+ )
+ if (splitHistoryByRole && !currentRoleName.isNullOrBlank()) {
+ val assistantCount = memory.count { it.kind == PromptTurnKind.ASSISTANT }
+ val userCount = memory.count { it.kind == PromptTurnKind.USER }
+ AppLogger.d(
+ TAG,
+ "按角色拆解历史: role=$currentRoleName, assistant=$assistantCount, user=$userCount, total=${memory.size}"
+ )
+ }
- return withContext(Dispatchers.IO) {
val limitHistoryStartTime = messageTimingNow()
val maxImageHistoryUserTurns = apiPreferences.maxImageHistoryUserTurnsFlow.first()
val maxMediaHistoryUserTurns = apiPreferences.maxMediaHistoryUserTurnsFlow.first()
@@ -702,7 +702,7 @@ object AIMessageManager {
return null
}
- val memoryTagRegex = Regex(".*?", RegexOption.DOT_MATCHES_ALL)
+ val memoryTagRegex = ChatMarkupRegex.memoryTag
val conversationReviewEntries = mutableListOf>()
fun normalizeForReview(text: String): String {
return text
@@ -1393,11 +1393,7 @@ object AIMessageManager {
targetRoleName: String,
removeStatusTags: (String) -> String
): PromptTurn? {
- // 清理思考内容
- val cleanedContent = ChatUtils.removeThinkingContent(message.content).trim()
- val contentWithoutStatus = removeStatusTags(cleanedContent)
-
- // 非角色隔离模式:直接返回 assistant 消息
+ // 非角色隔离模式和当前角色消息无需清洗;只有跨角色桥接才需要移除思考内容。
if (!isRoleScopedMode) {
return PromptTurn(
kind = PromptTurnKind.ASSISTANT,
@@ -1407,25 +1403,24 @@ object AIMessageManager {
// 角色隔离模式:判断是当前角色还是其他角色
val messageRoleName = message.roleName.trim()
- return if (messageRoleName == targetRoleName) {
- // 当前角色的消息:作为 assistant 返回
- PromptTurn(
+ if (messageRoleName == targetRoleName) {
+ return PromptTurn(
kind = PromptTurnKind.ASSISTANT,
content = message.content
)
- } else {
- // 其他角色的消息:转换为 user 消息,添加角色标签
- val roleLabel = if (messageRoleName.isNotBlank()) messageRoleName else "unknown"
- val bridgedContent = removeStatusTags(cleanedContent)
- if (bridgedContent.isBlank()) {
- null
- } else {
- PromptTurn(
- kind = PromptTurnKind.USER,
- content = "[From role: $roleLabel]\n$bridgedContent"
- )
- }
}
+
+ val cleanedContent = ChatUtils.removeThinkingContent(message.content).trim()
+ val bridgedContent = removeStatusTags(cleanedContent)
+ // 其他角色的消息:转换为 user 消息,添加角色标签
+ val roleLabel = if (messageRoleName.isNotBlank()) messageRoleName else "unknown"
+ if (bridgedContent.isBlank()) {
+ return null
+ }
+ return PromptTurn(
+ kind = PromptTurnKind.USER,
+ content = "[From role: $roleLabel]\n$bridgedContent"
+ )
}
private fun processUserMessage(
diff --git a/app/src/main/java/com/ai/assistance/operit/services/ChatServiceCore.kt b/app/src/main/java/com/ai/assistance/operit/services/ChatServiceCore.kt
index 4d318611e3..19043188b2 100644
--- a/app/src/main/java/com/ai/assistance/operit/services/ChatServiceCore.kt
+++ b/app/src/main/java/com/ai/assistance/operit/services/ChatServiceCore.kt
@@ -239,7 +239,7 @@ class ChatServiceCore(
messageProcessingDelegate.cancelMessageForDestructiveMutation(chatId)
}
chatHistoryDelegate.setAfterDestructiveHistoryMutation { chatId ->
- messageCoordinationDelegate.refreshStableContextWindow(chatId = chatId)
+ messageCoordinationDelegate.scheduleStableContextWindowRefresh(chatId = chatId)
}
initialized = true
diff --git a/app/src/main/java/com/ai/assistance/operit/services/core/MessageCoordinationDelegate.kt b/app/src/main/java/com/ai/assistance/operit/services/core/MessageCoordinationDelegate.kt
index 0db56264f7..cac61a4435 100644
--- a/app/src/main/java/com/ai/assistance/operit/services/core/MessageCoordinationDelegate.kt
+++ b/app/src/main/java/com/ai/assistance/operit/services/core/MessageCoordinationDelegate.kt
@@ -46,7 +46,6 @@ import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.flow.first
import kotlinx.coroutines.launch
-import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.withContext
import kotlinx.coroutines.CancellationException
import kotlinx.coroutines.withTimeoutOrNull
@@ -121,6 +120,7 @@ class MessageCoordinationDelegate(
private val pendingAutoContinuationByChatId =
ConcurrentHashMap()
+ private val stableWindowRefreshJobsByChatId = ConcurrentHashMap()
init {
ensureNonFatalErrorCollectorStarted()
@@ -270,9 +270,9 @@ class MessageCoordinationDelegate(
chatModelConfigIdOverride: String? = null,
chatModelIndexOverride: Int? = null,
memorySpaceIdOverride: String? = null
- ): Long? {
- val targetChatId = chatId ?: chatHistoryDelegate.currentChatId.value ?: return null
- val service = resolveWindowEstimateService(targetChatId) ?: return null
+ ): Long? = withContext(Dispatchers.IO) {
+ val targetChatId = chatId ?: chatHistoryDelegate.currentChatId.value ?: return@withContext null
+ val service = resolveWindowEstimateService(targetChatId) ?: return@withContext null
val effectiveRoleCardId = resolveRoleCardId(targetChatId, roleCardId)
val effectivePromptFunctionType = promptFunctionType ?: currentPromptFunctionType
val effectiveChatModelConfigIdOverride =
@@ -317,7 +317,43 @@ class MessageCoordinationDelegate(
"input=$inputTokens, output=$outputTokens, service=${service.javaClass.simpleName}, " +
"promptType=$effectivePromptFunctionType"
)
- return newWindowSize
+ newWindowSize
+ }
+
+ fun scheduleStableContextWindowRefresh(
+ chatId: String? = null,
+ roleCardId: String? = null,
+ promptFunctionType: PromptFunctionType? = null,
+ groupOrchestrationMode: Boolean = false,
+ groupParticipantNamesText: String? = null,
+ chatModelConfigIdOverride: String? = null,
+ chatModelIndexOverride: Int? = null,
+ memorySpaceIdOverride: String? = null
+ ): Job? {
+ val targetChatId = chatId ?: chatHistoryDelegate.currentChatId.value ?: return null
+ val refreshJob = coroutineScope.launch(Dispatchers.IO) {
+ try {
+ refreshStableContextWindow(
+ chatId = targetChatId,
+ roleCardId = roleCardId,
+ promptFunctionType = promptFunctionType,
+ groupOrchestrationMode = groupOrchestrationMode,
+ groupParticipantNamesText = groupParticipantNamesText,
+ chatModelConfigIdOverride = chatModelConfigIdOverride,
+ chatModelIndexOverride = chatModelIndexOverride,
+ memorySpaceIdOverride = memorySpaceIdOverride
+ )
+ } catch (e: CancellationException) {
+ throw e
+ } catch (e: Exception) {
+ AppLogger.w(TAG, "异步刷新上下文窗口失败: chatId=$targetChatId", e)
+ }
+ }
+ stableWindowRefreshJobsByChatId.put(targetChatId, refreshJob)?.cancel()
+ refreshJob.invokeOnCompletion {
+ stableWindowRefreshJobsByChatId.remove(targetChatId, refreshJob)
+ }
+ return refreshJob
}
/**
@@ -376,18 +412,20 @@ class MessageCoordinationDelegate(
)
}
} else {
- // 已有对话,直接发送消息
- sendMessageInternal(
- promptFunctionType,
- roleCardIdOverride = roleCardIdOverride,
- preferActiveRoleCard = preferActiveRoleCard,
- chatIdOverride = chatIdOverride,
- messageTextOverride = messageTextOverride,
- proxySenderNameOverride = proxySenderNameOverride,
- chatModelConfigIdOverride = chatModelConfigIdOverride,
- chatModelIndexOverride = chatModelIndexOverride,
- turnOptions = turnOptions
- )
+ // 已有对话,异步进入发送流程;内部的配置和历史读取会切到 IO
+ coroutineScope.launch {
+ sendMessageInternal(
+ promptFunctionType,
+ roleCardIdOverride = roleCardIdOverride,
+ preferActiveRoleCard = preferActiveRoleCard,
+ chatIdOverride = chatIdOverride,
+ messageTextOverride = messageTextOverride,
+ proxySenderNameOverride = proxySenderNameOverride,
+ chatModelConfigIdOverride = chatModelConfigIdOverride,
+ chatModelIndexOverride = chatModelIndexOverride,
+ turnOptions = turnOptions
+ )
+ }
}
}
@@ -515,7 +553,7 @@ class MessageCoordinationDelegate(
/**
* 内部发送消息的逻辑
*/
- private fun sendMessageInternal(
+ private suspend fun sendMessageInternal(
promptFunctionType: PromptFunctionType,
isContinuation: Boolean = false,
skipSummaryCheck: Boolean = false,
@@ -615,7 +653,7 @@ class MessageCoordinationDelegate(
val currentAttachments =
if (shouldReadComposerState) attachmentDelegate.attachments.value else emptyList()
// 手动发送必须与选择器一致;后台和定向消息仍由窗口绑定决定角色卡。
- val roleCardId = runBlocking {
+ val roleCardId = withContext(Dispatchers.IO) {
resolveRoleCardId(
chatId = chatId,
roleCardId = roleCardIdOverride,
@@ -623,32 +661,34 @@ class MessageCoordinationDelegate(
)
}
val resolvedOverrides = try {
- if (promptFunctionType == PromptFunctionType.CHAT) {
- val (resolvedChatModelConfigIdOverride, resolvedChatModelIndexOverride) =
- when {
- !chatModelConfigIdOverride.isNullOrBlank() -> {
- Pair(chatModelConfigIdOverride, (chatModelIndexOverride ?: 0).coerceAtLeast(0))
- }
- isAutoContinuation -> {
- Pair(currentChatModelConfigIdOverride, currentChatModelIndexOverride)
+ withContext(Dispatchers.IO) {
+ if (promptFunctionType == PromptFunctionType.CHAT) {
+ val (resolvedChatModelConfigIdOverride, resolvedChatModelIndexOverride) =
+ when {
+ !chatModelConfigIdOverride.isNullOrBlank() -> {
+ Pair(chatModelConfigIdOverride, (chatModelIndexOverride ?: 0).coerceAtLeast(0))
+ }
+ isAutoContinuation -> {
+ Pair(currentChatModelConfigIdOverride, currentChatModelIndexOverride)
+ }
+ else -> {
+ resolveRoleCardChatModelOverrides(roleCardId)
+ }
}
- else -> {
- resolveRoleCardChatModelOverrides(roleCardId)
+ val resolvedMemorySpaceIdOverride =
+ when {
+ !memorySpaceIdOverride.isNullOrBlank() -> memorySpaceIdOverride
+ isAutoContinuation -> currentMemorySpaceIdOverride
+ else -> roleCardId?.let { resolveRoleCardMemoryProfileOverride(it) }
}
- }
- val resolvedMemorySpaceIdOverride =
- when {
- !memorySpaceIdOverride.isNullOrBlank() -> memorySpaceIdOverride
- isAutoContinuation -> currentMemorySpaceIdOverride
- else -> roleCardId?.let { resolveRoleCardMemoryProfileOverride(it) }
- }
- Triple(
- resolvedChatModelConfigIdOverride,
- resolvedChatModelIndexOverride,
- resolvedMemorySpaceIdOverride
- )
- } else {
- Triple(null, null, null)
+ Triple(
+ resolvedChatModelConfigIdOverride,
+ resolvedChatModelIndexOverride,
+ resolvedMemorySpaceIdOverride
+ )
+ } else {
+ Triple(null, null, null)
+ }
}
} catch (e: Exception) {
AppLogger.e(TAG, "解析角色卡对话模型绑定失败", e)
@@ -661,7 +701,7 @@ class MessageCoordinationDelegate(
val resolvedChatModelIndexOverride = resolvedOverrides.second
val resolvedMemorySpaceIdOverride = resolvedOverrides.third
val chatContextSettings =
- runBlocking {
+ withContext(Dispatchers.IO) {
resolveChatContextSettingsForRequest(resolvedChatModelConfigIdOverride)
}
@@ -681,7 +721,10 @@ class MessageCoordinationDelegate(
// 如果不是续写,检查是否需要总结
if (turnOptions.persistTurn && !isBackgroundSend && !isContinuation && !skipSummaryCheck) {
- val currentMessages = runBlocking { chatHistoryDelegate.getCurrentRuntimeChatHistorySnapshot() }
+ val currentMessages =
+ withContext(Dispatchers.IO) {
+ chatHistoryDelegate.getCurrentRuntimeChatHistorySnapshot()
+ }
val currentTokens = tokenStatsDelegate.currentWindowSizeFlow.value
val isShouldGenerateSummary = AIMessageManager.shouldGenerateSummary(
@@ -763,7 +806,7 @@ class MessageCoordinationDelegate(
}
}
- private fun shouldRunGroupOrchestration(
+ private suspend fun shouldRunGroupOrchestration(
promptFunctionType: PromptFunctionType,
isContinuation: Boolean,
isAutoContinuation: Boolean,
@@ -779,7 +822,7 @@ class MessageCoordinationDelegate(
if (!proxySenderNameOverride.isNullOrBlank()) return false
if (!messageTextOverride.isNullOrBlank()) return false
if (!chatIdOverride.isNullOrBlank()) return false
- val activePrompt = runBlocking { activePromptManager.getActivePrompt() }
+ val activePrompt = withContext(Dispatchers.IO) { activePromptManager.getActivePrompt() }
if (activePrompt !is ActivePrompt.CharacterGroup) return false
return true
}
@@ -1416,8 +1459,8 @@ class MessageCoordinationDelegate(
}
}
- private fun resolveRoleCardChatModelOverrides(roleCardId: String): Pair {
- val roleCard = runBlocking { characterCardManager.getCharacterCardFlow(roleCardId).first() }
+ private suspend fun resolveRoleCardChatModelOverrides(roleCardId: String): Pair {
+ val roleCard = characterCardManager.getCharacterCardFlow(roleCardId).first()
val bindingMode = CharacterCardChatModelBindingMode.normalize(roleCard.chatModelBindingMode)
return if (
bindingMode == CharacterCardChatModelBindingMode.FIXED_CONFIG &&
@@ -1429,8 +1472,8 @@ class MessageCoordinationDelegate(
}
}
- private fun resolveRoleCardMemoryProfileOverride(roleCardId: String): String? {
- val roleCard = runBlocking { characterCardManager.getCharacterCardFlow(roleCardId).first() }
+ private suspend fun resolveRoleCardMemoryProfileOverride(roleCardId: String): String? {
+ val roleCard = characterCardManager.getCharacterCardFlow(roleCardId).first()
val bindingMode =
CharacterCardMemoryProfileBindingMode.normalize(roleCard.memoryProfileBindingMode)
return if (
diff --git a/app/src/main/java/com/ai/assistance/operit/ui/features/chat/viewmodel/ChatViewModel.kt b/app/src/main/java/com/ai/assistance/operit/ui/features/chat/viewmodel/ChatViewModel.kt
index c4506c04fa..238e883b04 100644
--- a/app/src/main/java/com/ai/assistance/operit/ui/features/chat/viewmodel/ChatViewModel.kt
+++ b/app/src/main/java/com/ai/assistance/operit/ui/features/chat/viewmodel/ChatViewModel.kt
@@ -873,12 +873,14 @@ class ChatViewModel(private val context: Context) : ViewModel() {
val beforeTimestamp = if (message.sender == "ai") message.timestamp else null
val afterTimestamp = if (message.sender == "user") message.timestamp else null
val messagesToSummarize =
- chatHistoryDelegate
- .loadMessagesForSummaryInsertion(
- chatId = currentChatId,
- beforeTimestampExclusive = afterTimestamp,
- upToTimestampInclusive = beforeTimestamp,
- ).filter { it.sender == "user" || it.sender == "ai" }
+ withContext(Dispatchers.IO) {
+ chatHistoryDelegate
+ .loadMessagesForSummaryInsertion(
+ chatId = currentChatId,
+ beforeTimestampExclusive = afterTimestamp,
+ upToTimestampInclusive = beforeTimestamp,
+ ).filter { it.sender == "user" || it.sender == "ai" }
+ }
if (messagesToSummarize.isEmpty()) {
uiStateDelegate.showToast(context.getString(R.string.chat_no_messages_to_summarize))
@@ -899,25 +901,30 @@ class ChatViewModel(private val context: Context) : ViewModel() {
// 检查是否是群聊
val currentChat = chatHistoryDelegate.chatHistories.value.firstOrNull { it.id == currentChatId }
val isGroupChat = currentChat?.characterGroupId != null
- val summaryConfig = messageCoordinationDelegate.readSummaryConfig()
-
- val summaryMessage = AIMessageManager.summarizeMemory(
- enhancedAiService!!,
- messagesToSummarize,
- autoContinue = false,
- isGroupChat = isGroupChat,
- summaryConfig = summaryConfig
- )
+ val summaryConfig = withContext(Dispatchers.IO) {
+ messageCoordinationDelegate.readSummaryConfig()
+ }
- if (summaryMessage != null) {
- // 插入总结消息
- chatHistoryDelegate.addSummaryMessage(
- summaryMessage = summaryMessage,
- beforeTimestamp = beforeTimestamp,
- afterTimestamp = afterTimestamp,
+ val summaryMessage = withContext(Dispatchers.IO) {
+ AIMessageManager.summarizeMemory(
+ enhancedAiService!!,
+ messagesToSummarize,
+ autoContinue = false,
+ isGroupChat = isGroupChat,
+ summaryConfig = summaryConfig
)
+ }
- messageCoordinationDelegate.refreshStableContextWindow(chatId = currentChatId)
+ if (summaryMessage != null) {
+ // 插入总结消息
+ withContext(Dispatchers.IO) {
+ chatHistoryDelegate.addSummaryMessage(
+ summaryMessage = summaryMessage,
+ beforeTimestamp = beforeTimestamp,
+ afterTimestamp = afterTimestamp,
+ )
+ }
+ messageCoordinationDelegate.scheduleStableContextWindowRefresh(chatId = currentChatId)
uiStateDelegate.showToast(context.getString(R.string.chat_summary_inserted))
} else {
@@ -1112,8 +1119,8 @@ class ChatViewModel(private val context: Context) : ViewModel() {
// 直接在数据库中更新该条消息
chatHistoryDelegate.addMessageToChat(editedMessage)
-
- messageCoordinationDelegate.refreshStableContextWindow(chatId = currentChatId.value)
+ // 变更历史后异步刷新窗口,不阻塞编辑操作
+ messageCoordinationDelegate.scheduleStableContextWindowRefresh(chatId = currentChatId.value)
// 显示成功提示
uiStateDelegate.showToast(context.getString(R.string.chat_message_updated))
@@ -1153,7 +1160,7 @@ class ChatViewModel(private val context: Context) : ViewModel() {
return@launch
}
chatHistoryDelegate.selectMessageVariant(targetMessage.timestamp, targetVariantIndex)
- messageCoordinationDelegate.refreshStableContextWindow(chatId = currentChatId.value)
+ messageCoordinationDelegate.scheduleStableContextWindowRefresh(chatId = currentChatId.value)
} catch (e: Exception) {
AppLogger.e(TAG, "切换回答版本失败", e)
uiStateDelegate.showErrorMessage(
@@ -1186,7 +1193,7 @@ class ChatViewModel(private val context: Context) : ViewModel() {
timestamp = targetMessage.timestamp,
variantIndex = targetMessage.selectedVariantIndex,
)
- messageCoordinationDelegate.refreshStableContextWindow(chatId = currentChatId.value)
+ messageCoordinationDelegate.scheduleStableContextWindowRefresh(chatId = currentChatId.value)
} catch (e: Exception) {
AppLogger.e(TAG, "删除当前消息候选失败", e)
uiStateDelegate.showErrorMessage(
From d0fdd163cbd6cf7b5fa5f25c7c1a425f99e6ecf2 Mon Sep 17 00:00:00 2001
From: 3316891527 <135376906+3316891527@users.noreply.github.com>
Date: Sat, 26 Sep 2026 03:11:48 +0000
Subject: [PATCH 2/2] fix(chat): serialize window refresh and handle async send
errors
---
.../core/MessageCoordinationDelegate.kt | 203 +++++++++++-------
1 file changed, 121 insertions(+), 82 deletions(-)
diff --git a/app/src/main/java/com/ai/assistance/operit/services/core/MessageCoordinationDelegate.kt b/app/src/main/java/com/ai/assistance/operit/services/core/MessageCoordinationDelegate.kt
index cac61a4435..62aa84bad0 100644
--- a/app/src/main/java/com/ai/assistance/operit/services/core/MessageCoordinationDelegate.kt
+++ b/app/src/main/java/com/ai/assistance/operit/services/core/MessageCoordinationDelegate.kt
@@ -37,7 +37,9 @@ import com.ai.assistance.operit.util.ChatMarkupRegex
import com.ai.assistance.operit.util.ChatUtils
import com.ai.assistance.operit.util.LocaleUtils
import com.ai.assistance.operit.data.repository.MemoryAutoSaveCandidateRepository
+import kotlinx.coroutines.CancellationException
import kotlinx.coroutines.CoroutineScope
+import kotlinx.coroutines.CoroutineStart
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.Job
import kotlinx.coroutines.delay
@@ -47,7 +49,9 @@ import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.flow.first
import kotlinx.coroutines.launch
import kotlinx.coroutines.withContext
-import kotlinx.coroutines.CancellationException
+import kotlinx.coroutines.ensureActive
+import kotlinx.coroutines.sync.Mutex
+import kotlinx.coroutines.sync.withLock
import kotlinx.coroutines.withTimeoutOrNull
import org.json.JSONArray
import org.json.JSONObject
@@ -121,6 +125,7 @@ class MessageCoordinationDelegate(
private val pendingAutoContinuationByChatId =
ConcurrentHashMap()
private val stableWindowRefreshJobsByChatId = ConcurrentHashMap()
+ private val stableWindowRefreshMutexesByChatId = ConcurrentHashMap()
init {
ensureNonFatalErrorCollectorStarted()
@@ -269,55 +274,55 @@ class MessageCoordinationDelegate(
groupParticipantNamesText: String? = null,
chatModelConfigIdOverride: String? = null,
chatModelIndexOverride: Int? = null,
- memorySpaceIdOverride: String? = null
+ memorySpaceIdOverride: String? = null,
): Long? = withContext(Dispatchers.IO) {
val targetChatId = chatId ?: chatHistoryDelegate.currentChatId.value ?: return@withContext null
- val service = resolveWindowEstimateService(targetChatId) ?: return@withContext null
- val effectiveRoleCardId = resolveRoleCardId(targetChatId, roleCardId)
- val effectivePromptFunctionType = promptFunctionType ?: currentPromptFunctionType
- val effectiveChatModelConfigIdOverride =
- chatModelConfigIdOverride ?: currentChatModelConfigIdOverride
- val effectiveChatModelIndexOverride =
- chatModelIndexOverride ?: currentChatModelIndexOverride
- val effectiveMemorySpaceIdOverride =
- memorySpaceIdOverride
- ?: currentMemorySpaceIdOverride
- ?: effectiveRoleCardId?.let { resolveRoleCardMemoryProfileOverride(it) }
-
- val newWindowSize =
- recalculateStableWindowSize(
- service = service,
- chatId = targetChatId,
- roleCardId = effectiveRoleCardId,
- promptFunctionType = effectivePromptFunctionType,
- groupOrchestrationMode = groupOrchestrationMode,
- groupParticipantNamesText = groupParticipantNamesText,
- chatModelConfigIdOverride = effectiveChatModelConfigIdOverride,
- chatModelIndexOverride = effectiveChatModelIndexOverride,
- memorySpaceIdOverride = effectiveMemorySpaceIdOverride
+ // 同一聊天的估算、持久化和界面写回整体串行,旧结果不会覆盖新结果。
+ stableWindowRefreshMutexesByChatId.getOrPut(targetChatId) { Mutex() }.withLock {
+ val service = resolveWindowEstimateService(targetChatId) ?: return@withLock null
+ val effectiveRoleCardId = resolveRoleCardId(targetChatId, roleCardId)
+ val effectivePromptFunctionType = promptFunctionType ?: currentPromptFunctionType
+ val effectiveChatModelConfigIdOverride =
+ chatModelConfigIdOverride ?: currentChatModelConfigIdOverride
+ val effectiveChatModelIndexOverride =
+ chatModelIndexOverride ?: currentChatModelIndexOverride
+ val effectiveMemorySpaceIdOverride =
+ memorySpaceIdOverride
+ ?: currentMemorySpaceIdOverride
+ ?: effectiveRoleCardId?.let { resolveRoleCardMemoryProfileOverride(it) }
+
+ val newWindowSize =
+ recalculateStableWindowSize(
+ service = service,
+ chatId = targetChatId,
+ roleCardId = effectiveRoleCardId,
+ promptFunctionType = effectivePromptFunctionType,
+ groupOrchestrationMode = groupOrchestrationMode,
+ groupParticipantNamesText = groupParticipantNamesText,
+ chatModelConfigIdOverride = effectiveChatModelConfigIdOverride,
+ chatModelIndexOverride = effectiveChatModelIndexOverride,
+ memorySpaceIdOverride = effectiveMemorySpaceIdOverride
+ )
+ coroutineContext.ensureActive()
+ val (inputTokens, outputTokens) = tokenStatsDelegate.getCumulativeTokenCounts(targetChatId)
+ chatHistoryDelegate.saveCurrentChat(
+ inputTokens = inputTokens,
+ outputTokens = outputTokens,
+ actualContextWindowSize = newWindowSize,
+ chatIdOverride = targetChatId
)
- val (inputTokens, outputTokens) = tokenStatsDelegate.getCumulativeTokenCounts(targetChatId)
- chatHistoryDelegate.saveCurrentChat(
- inputTokens = inputTokens,
- outputTokens = outputTokens,
- actualContextWindowSize = newWindowSize,
- chatIdOverride = targetChatId
- )
- withContext(Dispatchers.Main) {
- tokenStatsDelegate.setTokenCounts(
- targetChatId,
- inputTokens,
- outputTokens,
- newWindowSize
+ withContext(Dispatchers.Main) {
+ coroutineContext.ensureActive()
+ tokenStatsDelegate.setTokenCounts(targetChatId, inputTokens, outputTokens, newWindowSize)
+ }
+ AppLogger.d(
+ TAG,
+ "上下文窗口已刷新: chatId=$targetChatId, window=$newWindowSize, " +
+ "input=$inputTokens, output=$outputTokens, service=${service.javaClass.simpleName}, " +
+ "promptType=$effectivePromptFunctionType"
)
+ newWindowSize
}
- AppLogger.d(
- TAG,
- "上下文窗口已刷新: chatId=$targetChatId, window=$newWindowSize, " +
- "input=$inputTokens, output=$outputTokens, service=${service.javaClass.simpleName}, " +
- "promptType=$effectivePromptFunctionType"
- )
- newWindowSize
}
fun scheduleStableContextWindowRefresh(
@@ -331,29 +336,57 @@ class MessageCoordinationDelegate(
memorySpaceIdOverride: String? = null
): Job? {
val targetChatId = chatId ?: chatHistoryDelegate.currentChatId.value ?: return null
- val refreshJob = coroutineScope.launch(Dispatchers.IO) {
- try {
- refreshStableContextWindow(
- chatId = targetChatId,
- roleCardId = roleCardId,
- promptFunctionType = promptFunctionType,
- groupOrchestrationMode = groupOrchestrationMode,
- groupParticipantNamesText = groupParticipantNamesText,
- chatModelConfigIdOverride = chatModelConfigIdOverride,
- chatModelIndexOverride = chatModelIndexOverride,
- memorySpaceIdOverride = memorySpaceIdOverride
- )
- } catch (e: CancellationException) {
- throw e
- } catch (e: Exception) {
- AppLogger.w(TAG, "异步刷新上下文窗口失败: chatId=$targetChatId", e)
+ // 先登记惰性任务再启动,避免并发调度时旧任务反过来取消新任务。
+ return synchronized(stableWindowRefreshJobsByChatId) {
+ val refreshJob = coroutineScope.launch(Dispatchers.IO, start = CoroutineStart.LAZY) {
+ try {
+ refreshStableContextWindow(
+ chatId = targetChatId,
+ roleCardId = roleCardId,
+ promptFunctionType = promptFunctionType,
+ groupOrchestrationMode = groupOrchestrationMode,
+ groupParticipantNamesText = groupParticipantNamesText,
+ chatModelConfigIdOverride = chatModelConfigIdOverride,
+ chatModelIndexOverride = chatModelIndexOverride,
+ memorySpaceIdOverride = memorySpaceIdOverride
+ )
+ } catch (e: CancellationException) {
+ throw e
+ } catch (e: Exception) {
+ AppLogger.w(TAG, "异步刷新上下文窗口失败: chatId=$targetChatId", e)
+ }
}
+ stableWindowRefreshJobsByChatId.put(targetChatId, refreshJob)?.cancel()
+ refreshJob.invokeOnCompletion {
+ stableWindowRefreshJobsByChatId.remove(targetChatId, refreshJob)
+ }
+ refreshJob.start()
+ refreshJob
+ }
+ }
+
+ private fun launchMessageSend(block: suspend CoroutineScope.() -> Unit): Job =
+ coroutineScope.launch(Dispatchers.Main.immediate) {
+ runMessageSend { block() }
+ }
+
+ private suspend fun runMessageSend(block: suspend () -> Unit) {
+ try {
+ block()
+ } catch (e: CancellationException) {
+ throw e
+ } catch (e: Exception) {
+ reportSendFailure(e)
}
- stableWindowRefreshJobsByChatId.put(targetChatId, refreshJob)?.cancel()
- refreshJob.invokeOnCompletion {
- stableWindowRefreshJobsByChatId.remove(targetChatId, refreshJob)
+ }
+
+ private suspend fun reportSendFailure(error: Exception) {
+ AppLogger.e(TAG, "发送消息时出错", error)
+ withContext(Dispatchers.Main) {
+ uiStateDelegate.showErrorMessage(
+ context.getString(R.string.message_send_failed, error.message.orEmpty())
+ )
}
- return refreshJob
}
/**
@@ -376,7 +409,7 @@ class MessageCoordinationDelegate(
AppLogger.d(TAG, "当前没有活跃对话,自动创建新对话")
// 使用 coroutineScope 启动协程
- coroutineScope.launch {
+ launchMessageSend {
// 使用现有的createNewChat方法创建新对话
chatHistoryDelegate.createNewChat()
@@ -390,7 +423,7 @@ class MessageCoordinationDelegate(
if (chatHistoryDelegate.currentChatId.value == null) {
AppLogger.e(TAG, "创建新对话超时,无法发送消息")
uiStateDelegate.showErrorMessage(context.getString(R.string.chat_cannot_create_new))
- return@launch
+ return@launchMessageSend
}
AppLogger.d(
@@ -413,7 +446,7 @@ class MessageCoordinationDelegate(
}
} else {
// 已有对话,异步进入发送流程;内部的配置和历史读取会切到 IO
- coroutineScope.launch {
+ launchMessageSend {
sendMessageInternal(
promptFunctionType,
roleCardIdOverride = roleCardIdOverride,
@@ -607,7 +640,7 @@ class MessageCoordinationDelegate(
chatIdOverride = chatIdOverride
)
) {
- coroutineScope.launch {
+ launchMessageSend {
val handled = runCatching {
orchestrateGroupConversation(
chatId = chatId,
@@ -615,9 +648,11 @@ class MessageCoordinationDelegate(
turnOptions = turnOptions
)
}.getOrElse { throwable ->
+ if (throwable is CancellationException) throw throwable
AppLogger.e(TAG, "群组编排失败,回退普通发送", throwable)
false
}
+ coroutineContext.ensureActive()
if (!handled) {
sendMessageInternal(
promptFunctionType = promptFunctionType,
@@ -690,6 +725,8 @@ class MessageCoordinationDelegate(
Triple(null, null, null)
}
}
+ } catch (e: CancellationException) {
+ throw e
} catch (e: Exception) {
AppLogger.e(TAG, "解析角色卡对话模型绑定失败", e)
uiStateDelegate.showErrorMessage(
@@ -1411,7 +1448,7 @@ class MessageCoordinationDelegate(
} catch (e: CancellationException) {
throw e
} catch (e: Exception) {
- AppLogger.e(TAG, "派发自动续聊时出错: ${e.message}", e)
+ reportSendFailure(e)
if (isSamePendingAutoContinuation(chatId, request)) {
removePendingAutoContinuation(chatId)
messageProcessingDelegate.setSuppressIdleCompletedStateForChat(chatId, false)
@@ -2028,18 +2065,20 @@ class MessageCoordinationDelegate(
} else {
messageProcessingDelegate.setSuppressIdleCompletedStateForChat(currentChatId, false)
AppLogger.d(TAG, "总结成功,自动继续对话...")
- sendMessageInternal(
- promptFunctionType = continuationPromptType,
- isContinuation = true,
- isAutoContinuation = true,
- roleCardIdOverride = roleCardIdOverride,
- chatIdOverride = currentChatId,
- chatModelConfigIdOverride = effectiveChatModelConfigIdOverride,
- chatModelIndexOverride = effectiveChatModelIndexOverride,
- memorySpaceIdOverride = effectiveMemorySpaceIdOverride,
- isGroupOrchestrationTurn = isGroupOrchestrationTurn,
- groupParticipantNamesText = groupParticipantNamesText
- )
+ runMessageSend {
+ sendMessageInternal(
+ promptFunctionType = continuationPromptType,
+ isContinuation = true,
+ isAutoContinuation = true,
+ roleCardIdOverride = roleCardIdOverride,
+ chatIdOverride = currentChatId,
+ chatModelConfigIdOverride = effectiveChatModelConfigIdOverride,
+ chatModelIndexOverride = effectiveChatModelIndexOverride,
+ memorySpaceIdOverride = effectiveMemorySpaceIdOverride,
+ isGroupOrchestrationTurn = isGroupOrchestrationTurn,
+ groupParticipantNamesText = groupParticipantNamesText
+ )
+ }
if (!messageProcessingDelegate.isChatLoading(currentChatId)) {
AppLogger.w(TAG, "自动续聊未能启动,恢复Idle: chatId=$currentChatId")
restoreIdleIfCurrentlySummarizing(currentChatId)