- AI Agent
- 人工智能
- 大模型
- AI 应用
- 工具调用
- 本地部署
- MCP Clients
- Agent 记忆
【免费下载链接】Operit
The most powerful AI agent and AI chat software on Android/Operit是一款Android上能力最为强大、发展最久的AI Agent
导读
本文围绕 Operit 项目中 issue-741「中断回合统计」修复展开,剖析用户中断 AI 输出后「部分内容、Token、耗时、完成时间丢失」这一缺陷的根因,以及官方采用的两层修正方案:其一,由MessageProcessingDelegate的每会话运行态(ChatRuntime)持有当前流式 AI 消息,取消收尾直接持久化同一对象,不再依赖数据库实体推断运行态身份;其二,协调 Web 端删除会话与丢弃式取消的执行顺序,避免部分消息写入已删除会话。读完本文,你将掌握 Operit 聊天运行时中断收尾的完整调用链、回合编号互斥机制与对应的单元测试验证方式。
一、问题背景:取消收尾为何丢失中断统计
在 Operit 的聊天消息处理层中,用户中断模型输出是一个高频交互动作。修复前(旧实现)的取消收尾流程是:
- 消息处理层在取消模型输出之前,已经从当前服务读取了本轮 Token 与耗时快照;
- 取消任务结束后,收尾逻辑重新从 Room 数据库加载聊天记录;
- 通过只存在于内存中的
contentStream字段查找「当前流式 AI 消息」; - 对该消息回写部分回复内容、Token、耗时与完成时间。
问题在于第 3 步:contentStream是一个瞬态字段。查看 ChatMessage.kt 可以看到,该字段被@Transient注解标记:
@Transient var contentStream: Stream<String>? = null // 修改为Stream<String>类型,与EnhancedAIService.sendMessage返回类型匹配@Transient意味着该字段不属于数据库实体、不参与 Room 持久化。从 Room 重载后的所有消息contentStream恒为null,因此「按contentStream非空定位流式消息」的查找永远无法命中——于是中断时已生成的部分回复、输入/输出 Token、等待耗时、输出耗时以及completedAt全部不会写回,用户界面上表现为「中断即丢失」。
这正是 index.md 中描述的原始状况:取消前快照已读取,收尾却因身份识别失败而无法落盘。
二、修正总览:运行态消息所有权
官方修正的核心思想是把「当前可持久化的流式 AI 消息」的所有权从数据库模型迁移到每个会话的运行态对象,彻底切断收尾逻辑对 Room 重载结果的依赖。方案共五步(见 01-bind-cancellation-to-runtime-message.md):
- 在
ChatRuntime中保存当前流式 AI 消息及 Waifu 分段列表; - 取消任务前读取该对象与统计快照,任务停止后将两者一起交给收尾逻辑;
- 收尾逻辑只从数据库读取匹配的用户消息(用于回写回合统计),不再用数据库对象识别流式 AI 消息;
- 通过回合编号与取消互斥保证旧取消请求不会清理新回合;
- 正常完成与取消完成后均清理运行态消息引用。
从源码结构看,MessageProcessingDelegate为每个会话维护了一个ChatRuntime实例(ConcurrentHashMap<String, ChatRuntime>),其中关键字段定义在 MessageProcessingDelegate.kt:
private data class ActiveStreamingTurn( val message: ChatMessage, val segmentedMessages: MutableList<ChatMessage>? = null, ) private data class ChatRuntime( var sendJob: Job? = null, var responseStream: SharedStream<String>? = null, // 取消收尾必须持有运行态对象;Room 不保存 contentStream,重载后无法识别当前流消息。 var activeStreamingTurn: ActiveStreamingTurn? = null, var streamCollectionJob: Job? = null, var stateCollectionJob: Job? = null, var currentTurnOptions: ChatTurnOptions = ChatTurnOptions(), var requestSentAt: Long = 0L, var requestStartElapsed: Long = 0L, var firstResponseElapsed: Long? = null, val turnSequence: AtomicLong = AtomicLong(0L), @Volatile var activeTurnId: Long = 0L, val cancellationMutex: Mutex = Mutex(), @Volatile var cancellationInProgress: Boolean = false, val isLoading: MutableStateFlow<Boolean> = MutableStateFlow(false) )其中ActiveStreamingTurn同时持有单条流式 AI 消息(message)与Waifu 分段消息列表(segmentedMessages),这是 Waifu(分段长回复)模式下多个已形成分段能被统一统计的前提。注释也直接点明了设计动机:Room 不保存contentStream,重载后无法识别当前流消息。
三、取消收尾完整调用链
3.1 取消入口:普通取消 vs 丢弃式取消
MessageProcessingDelegate提供两个取消入口(见 MessageProcessingDelegate.kt):
fun cancelMessage(chatId: String) { val expectedTurnId = runtimeFor(chatId).activeTurnId coroutineScope.launch(Dispatchers.IO) { cancelMessageInternal( chatId = chatId, keepPartialResponse = true, expectedTurnId = expectedTurnId, ) } } suspend fun cancelMessageForDestructiveMutation(chatId: String) { cancelMessageInternal(chatId, keepPartialResponse = false) }cancelMessage(chatId):用户点击「停止」触发的普通取消,keepPartialResponse = true,会持久化部分回复;cancelMessageForDestructiveMutation(chatId):删除会话等破坏性操作前的丢弃式取消,keepPartialResponse = false,不持久化部分回复。
3.2 取消互斥与回合编号守卫
cancelMessageInternal的骨架(MessageProcessingDelegate.kt)体现了两个并发安全设计:
val chatRuntime = runtimeFor(chatId) chatRuntime.cancellationMutex.withLock { val turnId = expectedTurnId ?: chatRuntime.activeTurnId if (!chatRuntime.isLoading.value || chatRuntime.activeTurnId != turnId) { return@withLock } chatRuntime.cancellationInProgress = true ... }- 回合编号守卫:
turnSequence.incrementAndGet()在每次sendUserMessage开始时产生新turnId并写入activeTurnId。取消时把发起时刻的activeTurnId作为expectedTurnId传入,进入临界区后若发现当前activeTurnId已变化(即用户已发起新回合),则直接返回,保证旧取消请求永远不会清理新回合; - 取消互斥:整个取消过程持有
cancellationMutex,同一会话同一时刻只有一个取消操作在执行,避免重复取消导致运行态被多次清理。
3.3 快照捕获:取消前的 Token 与耗时
在任务真正停止之前,代码通过readCurrentTurnCancellationSnapshot(chatId)捕获统计快照(MessageProcessingDelegate.kt)。快照类型为:
internal data class TurnCancellationSnapshot( val inputTokens: Long, val outputTokens: Long, val cachedInputTokens: Long, val sentAt: Long, val outputDurationMs: Long, val waitDurationMs: Long, )各字段来源:
inputTokens/outputTokens/cachedInputTokens:来自service.captureCurrentTurnTokenSnapshot(),即模型服务当前回合的实时计数(含缓存命中 token);sentAt:runtime.requestSentAt,本轮请求发送时间戳;waitDurationMs:(firstResponseElapsed - requestStartElapsed),即「等待首包耗时」;outputDurationMs:(messageTimingNow() - firstResponseElapsed),即「首包到取消时刻的输出耗时」。
快照读取失败时仅记录AppLogger.w并返回null,不会阻断取消流程。
3.4 任务停止与收尾持久化
取消在cancellationMutex内依次完成:清除工具调用计数、调用AIMessageManager.cancelOperation(chatId)、逐个cancel()并join()发送/状态收集/流收集协程;随后若activeTurn != null(即keepPartialResponse = true且存在活跃流),调用detachStreamingAiMessage完成收尾持久化(MessageProcessingDelegate.kt):
- 解析最终内容:
resolveFinalContent从contentStream的 replayCache 或事件载体拼接已生成的全部文本,写回streamingMessage.content; - 构造完成消息:
completeInterruptedMessage将快照各统计字段复制到消息并写入completedAt,同时把contentStream置空(MessageProcessingDelegate.kt); - 回写用户消息统计:从
getRuntimeChatHistory(chatId)中按sender == "user"且sentAt匹配找到本轮用户消息,用withTurnMetrics写入相同的 Token 与耗时——这就是「中断前已写入的用户消息获得相同的回合统计」; - Waifu 分段处理:若
activeTurn.segmentedMessages非空,则不再落单条消息,而是把每一条已持久化分段copy上相同的统计字段与completedAt后逐一写回,保证 Waifu 模式下已形成的分段消息获得相同统计; - 持久化:若
currentTurnOptions.persistTurn为真,调用saveCurrentChat()落盘。
3.5 运行态清理
finally块中,只有当chatRuntime.activeTurnId == turnId(仍是当前回合)时才执行清理:将sendJob、stateCollectionJob、streamCollectionJob、responseStream、activeStreamingTurn全部置空,重置时间戳字段与isLoading,刷新全局加载状态。这样无论是正常完成还是取消完成,运行态消息引用都会被回收,不会残留影响下一回合。ChatRuntime的cancellationInProgress标志在finally中先复位,供上层判断取消是否仍在进行。
四、协调破坏性删除:Web 端删除会话的执行顺序
第二个独立缺陷(见 02-coordinate-destructive-delete.md):旧实现中,Web 接口发现目标会话仍在输出时,会异步请求取消并立即删除数据库会话。由于中断收尾需要写入部分消息,这两个操作可能交错,造成消息写入一个已删除的会话,产生悬挂数据。
修正后的执行顺序为:Web 删除接口先等待总结与消息处理的丢弃式取消完全结束,再执行会话删除。丢弃式取消(keepPartialResponse = false)不持久化部分回复,因此中断统计的普通保存路径完全不会参与破坏性删除。
代码层面的落地有两处:
ChatServiceCore初始化时注册了破坏性变更前钩子(ChatServiceCore.kt):
chatHistoryDelegate.setBeforeDestructiveHistoryMutation { chatId -> messageCoordinationDelegate.cancelSummaryForDestructiveMutation(chatId) messageProcessingDelegate.cancelMessageForDestructiveMutation(chatId) }- Web HTTP 桥的
handleDeleteChat在删除前先挂起等待取消完成(WebChatHttpBridge.kt):
if (core.activeStreamingChatIds.value.contains(chatId)) { runBlocking { core.cancelMessageForDestructiveMutation(chatId) } } val deleted = runBlocking { chatHistoryManager.deleteChatHistory(chatId) }由于cancelMessageForDestructiveMutation是suspend函数且内部使用job.join()等任务停止其协程,runBlocking会阻塞到丢弃式取消完整结束后才继续执行deleteChatHistory,从执行顺序上根除了「消息写入已删除会话」的竞态。
五、验收与测试验证
修复的验收标准(01-bind-cancellation-to-runtime-message.md)共四条:
- 中断消息写入部分内容、
completedAt、Token 和耗时; - 中断前已写入的用户消息获得相同的回合统计;
- Waifu 模式下已持久化的分段消息获得相同统计;
- 流式消息身份不依赖 Room 中不存在的字段(即不再依赖
contentStream)。
仓库为此新增了 MessageProcessingDelegateTest.kt,其中completeInterruptedMessage_appliesTurnSnapshotAndCompletesPartialContent用例直接构造一个带contentStream = emptyStream()的 AI 消息与TurnCancellationSnapshot,断言completeInterruptedMessage返回的消息满足:
content被替换为部分响应文本"partial response";contentStream被置为null;inputTokens、outputTokens、cachedInputTokens与快照一致(120 / 34 / 56);sentAt、outputDurationMs、waitDurationMs、completedAt与快照一致(30 / 4000 / 500 / 5000)。
该测试从纯函数层面锁定了「快照应用 + 部分内容完成 + 流引用清理」的核心契约。按仓库本地工作约束(index.md 验证记录注明「本次未执行构建或测试命令」),读者可在自己的构建环境中通过./gradlew :app:testDebugUnitTest --tests "com.ai.assistance.operit.services.core.MessageProcessingDelegateTest"运行该用例。
六、设计要点小结
| 关注点 | 旧实现 | 修正后 |
|---|---|---|
| 流式 AI 消息身份 | 取消收尾从 Room 重载后按contentStream非空查找 | ChatRuntime.activeStreamingTurn运行态持有 |
| 统计快照来源 | 取消前读取,收尾时丢失 | TurnCancellationSnapshot与消息引用一起传递 |
| 用户消息统计回写 | 无(因 AI 消息未命中) | 按sentAt匹配用户消息并withTurnMetrics回写 |
| Waifu 分段 | 无 | segmentedMessages逐条复制统计与completedAt |
| 并发安全 | 无 | 回合编号activeTurnId+cancellationMutex互斥 |
| Web 删除会话 | 异步取消与删除交错 | runBlocking等待丢弃式取消结束后再删除 |
| 运行态清理 | 无 | finally中按回合校验后清理全部引用 |
整体上,这一修复把「谁拥有当前流式消息」从持久化模型的身份推断改为运行态对象的显式持有,配合回合编号互斥与破坏性操作顺序协调,同时保障了普通取消的统计保留、Waifu 分段统计一致性与删除会话的数据完整性。若要进一步深入,可继续阅读 preserve_interrupted_ai_output_20260814(网络失败消息定稿的同类收尾问题)以及 chat_runtime_foreground_service_plan.md(ChatRuntime运行态的服务化演进),两者与本文共用同一套运行态消息生命周期模型。
- AI Agent
- 人工智能
- 大模型
- AI 应用
- 工具调用
- 本地部署
- MCP Clients
- Agent 记忆
【免费下载链接】Operit
The most powerful AI agent and AI chat software on Android/Operit是一款Android上能力最为强大、发展最久的AI Agent
相关推荐
Operit 中断回合统计修复:将取消快照与运行态流式消息绑定,保留部分回复与 Token/耗时数据
Operit 中断回合统计修复:将取消快照与运行态流式消息绑定,保留部分回复与 Token/耗时数据 中断回合统计是 AI 聊天应用中一个极易被忽略、却直接影响
AI Agent人工智能大模型AI 应用工具调用本地部署MCP ClientsAgent 记忆GUI 自动化Operit 会话破坏性删除与中断收尾的协调:丢弃式取消如何阻止消息写入已删会话
Operit 会话破坏性删除与中断收尾的协调:丢弃式取消如何阻止消息写入已删会话 导读 本文围绕 Operit 中"中断回合统计修复"(issue 741)的第
AI Agent人工智能大模型AI 应用工具调用本地部署MCP ClientsAgent 记忆GUI 自动化Operit 超大聊天消息读取修复:基于 Room 事务分块读取规避 Android CursorWindow 溢出
Operit 超大聊天消息读取修复:基于 Room 事务分块读取规避 Android CursorWindow 溢出 本指南聚焦 Operit(Android
AI Agent人工智能大模型AI 应用工具调用本地部署MCP ClientsAgent 记忆GUI 自动化
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考