diff --git a/app/src/main/java/com/stand/standapp/ui/chat/ChatViewModel.kt b/app/src/main/java/com/stand/standapp/ui/chat/ChatViewModel.kt index d56a2bb..8c44f96 100644 --- a/app/src/main/java/com/stand/standapp/ui/chat/ChatViewModel.kt +++ b/app/src/main/java/com/stand/standapp/ui/chat/ChatViewModel.kt @@ -13,11 +13,16 @@ import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.asStateFlow import kotlinx.coroutines.flow.catch import kotlinx.coroutines.flow.onCompletion +import kotlinx.coroutines.flow.update import kotlinx.coroutines.launch +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock import org.json.JSONArray import org.json.JSONObject import java.util.UUID +private const val STREAM_FLUSH_INTERVAL_MS = 50L + class ChatViewModel(application: Application) : AndroidViewModel(application) { private val repository = AiChatRepository() @@ -30,14 +35,49 @@ class ChatViewModel(application: Application) : AndroidViewModel(application) { private val _messages = MutableStateFlow>(emptyList()) val messages: StateFlow> = _messages.asStateFlow() - // 👈 声明一个正在加载的状态,保证上一个问题没结束前无法再次发送提问! + private val _streamingContent = MutableStateFlow>(emptyMap()) + val streamingContent: StateFlow> = _streamingContent.asStateFlow() + private val _isLoading = MutableStateFlow(false) val isLoading: StateFlow = _isLoading.asStateFlow() - - // 👈 额外声明一个当前协程控制作业,用于物理取消/中止 SSE 对话! + private var activeChatJob: kotlinx.coroutines.Job? = null private val prefs = application.getSharedPreferences("ai_chat_cache", Context.MODE_PRIVATE) + private val chunkBuffer = java.util.concurrent.ConcurrentHashMap() + private val lastFlushAt = java.util.concurrent.ConcurrentHashMap() + private val flushScope = kotlinx.coroutines.CoroutineScope( + kotlinx.coroutines.SupervisorJob() + kotlinx.coroutines.Dispatchers.Main.immediate + ) + private val flushMutex = kotlinx.coroutines.sync.Mutex() + + private fun scheduleFlush(messageId: String) { + val now = android.os.SystemClock.uptimeMillis() + val last = lastFlushAt[messageId] ?: 0 + val delta = now - last + if (delta >= STREAM_FLUSH_INTERVAL_MS) { + flushStreamingMessage(messageId) + } else { + flushScope.launch { + kotlinx.coroutines.delay(STREAM_FLUSH_INTERVAL_MS - delta) + flushStreamingMessage(messageId) + } + } + } + + private fun flushStreamingMessage(messageId: String) { + flushScope.launch { + flushMutex.withLock { + val buffer = chunkBuffer.remove(messageId) ?: return@withLock + val pending = buffer.toString() + if (pending.isEmpty()) return@withLock + _streamingContent.update { current -> + current + (messageId to (current[messageId].orEmpty() + pending)) + } + lastFlushAt[messageId] = android.os.SystemClock.uptimeMillis() + } + } + } init { // 👈 将数据加载移到IO线程,避免阻塞主线程导致UI卡顿 @@ -132,18 +172,21 @@ class ChatViewModel(application: Application) : AndroidViewModel(application) { // 👈 用户手动中断当前正在生成的 AI 回答的方法! fun cancelActiveStreaming() { - activeChatJob?.cancel() // 物理取消协程作业,断开 SSE OkHttp 连接并关闭流! + activeChatJob?.cancel() activeChatJob = null - - // 将当前正在流式输出的 AI 消息状态重置为完成,防止气泡卡死在 streaming 样式 + _messages.value = _messages.value.map { msg -> if (msg.isStreaming) { - msg.copy(isStreaming = false) + val finalContent = _streamingContent.value[msg.id] ?: msg.content + msg.copy(content = finalContent, isStreaming = false) } else msg } - - _isLoading.value = false // 释放锁,允许用户立刻开始下一次提问! - saveCacheToLocal() // 强制存盘归档 + _streamingContent.value = emptyMap() + chunkBuffer.clear() + lastFlushAt.clear() + + _isLoading.value = false + saveCacheToLocal() } fun createNewSession() { @@ -232,20 +275,22 @@ class ChatViewModel(application: Application) : AndroidViewModel(application) { } private fun appendAiMessageChunk(messageId: String, chunk: String) { - _messages.value = _messages.value.map { msg -> - if (msg.id == messageId) { - msg.copy(content = msg.content + chunk) - } else msg - } + chunkBuffer.computeIfAbsent(messageId) { StringBuilder() }.append(chunk) + scheduleFlush(messageId) } private fun finalizeAiMessage(messageId: String) { + flushStreamingMessage(messageId) _messages.value = _messages.value.map { msg -> if (msg.id == messageId) { - msg.copy(isStreaming = false) + val finalContent = _streamingContent.value[messageId] ?: msg.content + _streamingContent.update { it - messageId } + msg.copy(content = finalContent, isStreaming = false) } else msg } - saveCacheToLocal() // 👈 AI 回复完毕后,状态变更为非 streaming 并归档存盘 + chunkBuffer.remove(messageId) + lastFlushAt.remove(messageId) + saveCacheToLocal() } private fun updateSessionTitleIfFirstMessage(sessionId: String, content: String) { diff --git a/app/src/main/java/com/stand/standapp/ui/chat/components/ChatScaffold.kt b/app/src/main/java/com/stand/standapp/ui/chat/components/ChatScaffold.kt index 708a645..13a9086 100644 --- a/app/src/main/java/com/stand/standapp/ui/chat/components/ChatScaffold.kt +++ b/app/src/main/java/com/stand/standapp/ui/chat/components/ChatScaffold.kt @@ -31,9 +31,19 @@ fun ChatScaffold( val sessions by viewModel.sessions.collectAsState() val currentSessionId by viewModel.currentSessionId.collectAsState() val messages by viewModel.messages.collectAsState() + val streamingContent by viewModel.streamingContent.collectAsState() var showClearDialog by remember { mutableStateOf(false) } - val currentMessages = messages.filter { it.sessionId == currentSessionId } + val currentMessages = remember(messages, streamingContent, currentSessionId) { + val streaming = streamingContent + messages.asSequence() + .filter { it.sessionId == currentSessionId } + .map { msg -> + val live = streaming[msg.id] + if (msg.isStreaming && live != null) msg.copy(content = live) else msg + } + .toList() + } // 清除确认对话框 if (showClearDialog) {