//! 发送与队列管理 — sendMessage/approveToolCall + 队列(drain/cancel/clear) + stopChat //! //! B-260616-02 L2 发送韧性:三级降级 //! L0 normal:后端 idle → 直接发送 //! L1 queued: 后端 busy → 入队等待(≤30s 正常排队) //! L2 force: 排队>30s → 弹 confirm → 用户确认后调 ai_chat_force_send 复位残留并强制发出 //! //! 模块级私有: //! - _pendingApprovalIds:审批按钮防抖守卫 //! //! 耦合: //! - sendMessage 调 useAiEvents.startListener(useAiEvents 导出但不在解构集中暴露给组件, //! 故 startListener 同时经 useAiEvents 导出,sendMessage 直接 import 调用) //! - sendMessage 调 useAiStream.resetStreamWatchdog //! - drainQueue 由 handleEvent(AiCompleted) 经 ai-drain-queue 事件总线触发(非直接调用, //! 打破 events↔send 循环依赖);本模块模块级监听该事件 import { listen } from '@tauri-apps/api/event' import { invoke } from '@tauri-apps/api/core' import { ref } from 'vue' import { aiApi } from '@/api' import { state } from '@/stores/ai' import type { ContentPart, AiMessage } from '@/api/types' import i18n from '@/i18n' import { resetStreamWatchdog, clearStreamWatchdog } from './useAiStream' import { startListener } from './useAiEvents' import { nextMsgId, findToolCall, startApprovalTimer, clearApprovalTimer, clearAllApprovalTimers, resolveAiLang } from './aiShared' const t = ((i18n as any).global.t as (k: string, named?: Record) => string).bind((i18n as any).global) /// 待发送队列上限(超过抛错提示用户) const QUEUE_LIMIT = 10 /// 排队超时阈值(ms) — 超过后弹 confirm 让用户选择"强制发送"或继续等(B-260616-02) const QUEUE_TIMEOUT_MS = 30_000 /** * F-01 阶段6: 用户指定模型 override(主对话专用,模块级单例)。 * * null=自动模式(路由器选);非空字符串=用户从下拉指定的 model_id。 * 仅主对话(agentic)生效:doSend/sendMessage/regenerate/editMessage 透传给后端, * 标题/扫描/灵感等内部调用不读此字段(后端各自独立 IPC 不带 override)。 * * 后端兜底(agentic.rs):override 非空且在该 provider model_configs 池中才用, * 否则落回路由结果——绝不让 override 导致无模型。 * * 模块级(非组件级):AiChat.vue 顶部下拉读写此 ref,发送链路读此 ref 透传, * 单实例 AiChat 不需 props/inject 串联。切换对话时后端清 session.model_override, * 前端应同步清(modelOverride.value = null)保持 UI 与后端一致。 */ const modelOverride = ref(null) // 审批计时器(startApprovalTimer/clearApprovalTimer/clearAllApprovalTimers)已下沉到 aiShared.ts // 打破 useAiEvents ↔ useAiSend 循环依赖。本模块 re-export 供组件使用。 /** * F-260616-06: 审批按钮防抖守卫——同一 id 短期多次点击只发一次 IPC。 * IPC 往返 <100ms,300ms 窗口足够覆盖竞态。 */ const _pendingApprovalIds = new Set() /** * 入队时记录时间戳,供 UI 展示排队时长和触发超时提示。 * queue item 扩展: 原 { text, skill? } → { text, skill?, enqueuedAt, parts? } * 兼容 drainQueue/cancelQueued/clearQueue 等已有消费方(enqueuedAt/parts 仅读不写)。 */ // ── 内核发送:推送气泡+调 IPC,被 normal/force 两条路径共用 ── /** * F-260614-05 Phase 2b: parts 仅影响本地 push 的 user 消息(渲染多模态图)。 * * 后端 ai_chat_send IPC 当前不接 parts(Phase 2c 接入);本轮 doSend 仅把 parts 写入本地 push 的 * user 消息让用户气泡可见图,后端请求仍走 content 文本路径(无图)。 * Phase 2c 后端接 parts 后,本函数透传给 aiApi.sendMessage/forceSend 即可全链路生效。 */ async function doSend(text: string, skill?: string, force = false, parts?: ContentPart[]) { const userMsgId = `user-${nextMsgId()}` state.messages.push({ id: userMsgId, role: 'user', content: text.trim(), // 仅在非空时挂 parts(纯文本消息保持 undefined,渲染走原 content 路径零回归) parts: parts && parts.length > 0 ? parts : undefined, timestamp: Date.now(), }) const aiMsgId = `ai-${nextMsgId()}` state.messages.push({ id: aiMsgId, role: 'assistant', content: '', timestamp: Date.now(), }) state.streaming = true state.currentText = '' resetStreamWatchdog() // 启动流式看门狗,无数据超时兜底 await startListener() const lang = resolveAiLang() try { // force=true 走 force_send IPC(复位残留 generating),否则走正常 send // F-01 阶段6: 透传 modelOverride(主对话专用,后端兜底校验池内才用) const override = modelOverride.value if (force) { // F-05 Phase 2c: 透传 parts(图片)给后端 ai_chat_force_send,多模态全链生效 // F-260616-09 B 批4(决策 e):传 activeConversationId,操作指定 conv 的 per_conv。 await aiApi.forceSend(text.trim(), lang, skill, override, parts, state.activeConversationId) } else { await aiApi.sendMessage(text.trim(), lang, skill, override, parts, state.activeConversationId) } } catch (e) { // IPC 失败(spawn 前/provider 配置错等):回滚 streaming 防光标卡死 + 移除 user 消息与空气泡占位; // 重新抛出由 handleSend 回填输入框,用户可重试 state.streaming = false clearStreamWatchdog() // 清看门狗,防 130s 后 onStreamTimeout 误推错误气泡 state.messages = state.messages.filter(m => m.id !== aiMsgId && m.id !== userMsgId) throw e } } /** * 重新生成最后一条 AI 回复(UX-02)。 * * 与 doSend 的区别: * - 不 push 新 user 消息(触发它的 user 消息已在历史末尾) * - 仅 push 一个空 AI 气泡占位 + 置 streaming,后端弹出旧 AI 回复后重跑 loop * - 失败回滚只移除占位气泡(不移除 user 消息) * * 后端 ai_regenerate 占用 generating 防并发双发,此处不再额外入队/预检。 */ async function regenerate() { // 末尾须有 AI 回复可重生成(后端也会校验,前端先挡避免空跑) const msgs = state.messages const last = msgs[msgs.length - 1] if (!last || last.role !== 'assistant') { throw new Error(t('aiChat.regenerateUnavailable')) } // 移除末尾旧 AI 气泡(后端也弹 history;前端先弹避免占位前残留旧文),push 空气泡占位 state.messages.pop() const aiMsgId = `ai-${nextMsgId()}` state.messages.push({ id: aiMsgId, role: 'assistant', content: '', timestamp: Date.now(), }) state.streaming = true state.currentText = '' state.queue = [] // 重生成视同新发,清残留队列 resetStreamWatchdog() await startListener() const lang = resolveAiLang() const convId = state.activeConversationId if (!convId) { // 无活跃对话:回滚占位,报错 state.streaming = false clearStreamWatchdog() state.messages = state.messages.filter(m => m.id !== aiMsgId) throw new Error(t('aiChat.regenerateUnavailable')) } try { await aiApi.regenerate(convId, lang, modelOverride.value) } catch (e) { // IPC 失败:回滚 streaming + 移除占位气泡(旧回复后端已弹但前端先弹了,保持弹后态) state.streaming = false clearStreamWatchdog() state.messages = state.messages.filter(m => m.id !== aiMsgId) throw e } } /** * 编辑最后一条 user 消息并重新生成(UX-09)。 * * 与 regenerate 的区别: * - regenerate:末条须是 assistant,pop 旧 AI 回复后重跑 * - editMessage:末条须是 user(中间编辑被后端拒绝),就地改其 content,截断其后所有消息, * push 空气泡占位后重跑。被截断消息后端标 truncated(软删保留 DB),前端视图层同步移除。 * * 后端 ai_chat_edit 占用 generating 防并发双发,此处不再额外入队/预检。 * * @param newMessage 编辑后的新内容(替换末条 user 消息 content) */ async function editMessage(newMessage: string) { if (!newMessage.trim()) return // 末条须是 user(中间编辑语义复杂,拒绝);后端也会校验,前端先挡避免空跑 const msgs = state.messages const last = msgs[msgs.length - 1] if (!last || last.role !== 'user') { throw new Error(t('aiChat.editNotLastUser')) } // 截断末条 user 之后的消息(此时 last 即末条 user,其后应已有 assistant 回复或为空) // 找到末条 user 的下标,移除其后所有消息 let lastUserIdx = -1 for (let i = msgs.length - 1; i >= 0; i--) { if (msgs[i].role === 'user') { lastUserIdx = i; break } } if (lastUserIdx < 0) { throw new Error(t('aiChat.editNotLastUser')) } // 截断其后(保留末条 user 本身,改其 content) state.messages.splice(lastUserIdx + 1) msgs[lastUserIdx].content = newMessage.trim() // push 空气泡占位(新 AI 回复将流入) const aiMsgId = `ai-${nextMsgId()}` state.messages.push({ id: aiMsgId, role: 'assistant', content: '', timestamp: Date.now(), }) state.streaming = true state.currentText = '' state.queue = [] // 编辑重生成视同新发,清残留队列 resetStreamWatchdog() await startListener() const lang = resolveAiLang() const convId = state.activeConversationId if (!convId) { state.streaming = false clearStreamWatchdog() state.messages = state.messages.filter(m => m.id !== aiMsgId) throw new Error(t('aiChat.editNotLastUser')) } try { await aiApi.editMessage(convId, newMessage.trim(), lang, modelOverride.value) } catch (e) { // IPC 失败:回滚 streaming + 移除占位气泡(被截断的旧回复后端已标 truncated,前端已移除) state.streaming = false clearStreamWatchdog() state.messages = state.messages.filter(m => m.id !== aiMsgId) throw e } } /** * 取出队首并发送(AiCompleted 触发,此时 streaming 已 false)。 * * B-260617-02: 原 void sendMessage(...) fire-and-forget 吞没了 doSend IPC 失败 throw, * 既无 AiCompleted 触发下次 drain(队列永久卡死),也无任何用户反馈。 * 现显式 catch:复位 streaming + 推错误气泡(复用 AiError 的 isError 气泡模式) + * 回填失败消息到队首并停 drain(不静默丢用户输入,保留剩余队列供手动重试/编辑/取消)。 * 成功路径不受影响——doSend 成功后端会再 emit AiCompleted 续 drain 链路。 */ export function drainQueue() { if (state.queue.length === 0) return const next = state.queue.shift()! void (async () => { try { await sendMessage(next.text, next.skill, false, next.parts) } catch (e) { state.streaming = false clearStreamWatchdog() const errMsg = e instanceof Error ? e.message : String(e) state.queue.unshift(next) // 回填失败消息到队首,保留剩余队列(不静默丢用户输入) state.messages.push({ id: `err-${nextMsgId()}`, role: 'assistant', content: t('ai.queuedSendFailed', { error: errMsg }), isError: true, timestamp: Date.now(), } as AiMessage) console.error('[AI] drainQueue 续发失败:', e) } })() } /** * 发送消息 — 三级降级(B-260616-02 L2 发送韧性): * * L0 normal: 后端 idle → 直接 doSend() * L1 queued: 后端 busy → 入队等 AiCompleted 续发(≤30s 正常等待) * L2 force: 排队>30s → 返回特殊标记,由调用方(handleSend)弹 confirm; * 用户确认后重新进 sendMessage(forceMode=true)走 forceSend IPC * * forceMode=true 时跳过入队,直接走 ai_chat_force_send。 */ async function sendMessage(text: string, skill?: string, forceMode = false, parts?: ContentPart[]) { if (!text.trim() && !(parts && parts.length)) return // ── L2 强制模式:跳过所有预检,直接 force_send ── if (forceMode) { await doSend(text, skill, true, parts) return } // B-260615-22(方案 A):发送前 IPC 查后端真实 generating。 // F-260616-09 B 批4(决策 e):传 activeConversationId 精确查当前 conv(不读全局单值)。 let backendGenerating = false try { backendGenerating = await invoke('ai_is_generating', { conversationId: state.activeConversationId || null }) } catch { backendGenerating = false } // ── L1:后端 busy → 入队 ── if (backendGenerating || state.streaming) { if (state.queue.length >= QUEUE_LIMIT) { throw new Error(t('ai.queueFull', { limit: QUEUE_LIMIT })) } state.streaming = true // F-260614-05 Phase 2b: 入队项挂 parts(供 drainQueue 续发时本地 user 消息渲染图) state.queue.push({ text: text.trim(), skill: skill || undefined, enqueuedAt: Date.now(), parts: parts && parts.length > 0 ? parts : undefined, }) return } // ── L0:normal 直接发送 ── await doSend(text, skill, false, parts) } /** 排队是否已超时(任一条超即算,因队首阻塞导致后续全堵) */ export function isQueueTimedOut(): boolean { const first = state.queue[0] return first !== undefined && (Date.now() - first.enqueuedAt > QUEUE_TIMEOUT_MS) } /** * 超时时弹 confirm,用户确认后以 force_mode 重发队首消息;取消则保持排队。 * * B-260617-04: force_send 失败时,原实现队首已 shift 致消息静默丢失。 * 现失败把队首 unshift 回队列(保留 UI 队列可见,用户可改普通发送/编辑/取消), * 并复位 streaming/clearStreamWatchdog(对齐 doSend 失败回填模式)让 UI 脱卡死态。 * 不自动重试 force(避免循环),交用户决定下一步。 */ export async function tryForceSend(confirmFn: (msg: string) => Promise): Promise { const first = state.queue[0] if (!first) return false const confirmed = await confirmFn(t('ai.forceSendConfirm')) if (!confirmed) return false // 从队列移除队首,以 force_mode 发送 state.queue.shift() // 清看门狗:doSend 内 force_send 路径会重设 streaming=true 并重启 watchdog, // 此处无需前置 false(前置 false 会触发 AiChat watch 清流式块再重建,产生瞬态抖动) clearStreamWatchdog() try { await sendMessage(first.text, first.skill, true, first.parts) return true } catch (e) { // force_send 也失败:回填队首保消息不丢,复位 streaming 让 UI 可继续操作 state.queue.unshift({ text: first.text, skill: first.skill, enqueuedAt: first.enqueuedAt, parts: first.parts, }) state.streaming = false clearStreamWatchdog() console.error('[AI] tryForceSend force_send 失败,消息已回填队首:', e) return false } } /** 工具审批:保持 pending_approval 让审批卡片可见(按钮 loading 由 ToolCard 本地 ref 持有), * IPC 失败时改 completed+错误文案。后端回事件后由 useAiEvents 转 completed/rejected。 * B-260616-08:不再乐观置 running——原写法让 .ai-tool-approval 整块消失(切骨架屏), * 按钮无 loading 中间态、用户无重审入口;loading 现下沉到 ToolCard 局部 ref,语义正确。 */ async function approveToolCall(toolCallId: string, approved: boolean) { // F-260616-06: 防抖——同一 id 仍在处理中则跳过(竞态点击/网络慢重试) if (_pendingApprovalIds.has(toolCallId)) return _pendingApprovalIds.add(toolCallId) try { // 复用 findToolCall(反向扫描)避免重复实现查找逻辑;仅缓存引用用于 IPC 失败回滚 const tc = findToolCall(toolCallId) // 重启看门狗覆盖审批执行→续生成窗口(useAiEvents.ts:162 AiApprovalRequired 已 clear, // 审批态无心跳兜底;approve 后后端要跑工具+续生成,期间任何后端异常不回则按钮永久 loading。 // 复用通用 watchdog 而非加审批专用超时:复用同一收尾路径(onStreamTimeout 兜底复位 streaming), // 避免新增独立计时器与状态分支) resetStreamWatchdog() try { await aiApi.approve(toolCallId, approved) } catch (e) { // IPC 未送达:用户已点过按钮,清看门狗(避免 130s 后误触发假错误)再回滚到失败态 clearStreamWatchdog() state.queue = [] // B-32:IPC 失败即对话结束,清队列防生成中入队消息静默丢失 // 仅 IPC 真失败(审批已处理/网络断)走到这里:后端工具失败已改走 emit completed,不进此分支 // 不回滚 pending_approval(会卡死按钮),改设 completed + 错误提示,让用户知晓失败 console.error('[AI] 审批操作未送达后端:', e) if (tc) { tc.status = 'completed' tc.result = t('ai.approvalNotDelivered', { error: e instanceof Error ? e.message : String(e) }) } throw e } } finally { _pendingApprovalIds.delete(toolCallId) } } /** 批量审批:遍历 pendingApprovals 逐个调用 approveToolCall */ async function batchApprove(decision: 'approve' | 'reject') { const approved = decision === 'approve' // 快照当前待审批列表(遍历中 state.pendingApprovals 会因事件回调而缩短) const ids = [...state.pendingApprovals].map(p => p.id) for (const id of ids) { try { await approveToolCall(id, approved) } catch { // 单条失败不中断批量操作(已由 approveToolCall 内部回滚该条状态) // 继续处理剩余项 } } } /** 取消队列中指定位置的消息 */ function cancelQueued(index: number) { state.queue.splice(index, 1) } /** 清空整个待发送队列 */ function clearQueue() { state.queue = [] } /** * UX-260616-06: 编辑队列项的 text(决策只编 text)。 * * 设计: * - 只改 text,skill 编辑态不暴露(skill 来自发送时技能联想,编辑态改 skill 复杂化交互无收益)。 * - 边界:index 越界 / newText 空白 → 直接忽略(no-op),不抛错(UI 不会触发,防御)。 * - queue item 类型 { text, skill?, enqueuedAt } 不改。 */ function editQueued(index: number, newText: string) { if (index < 0 || index >= state.queue.length) return const trimmed = newText.trim() if (!trimmed) return state.queue[index].text = trimmed } /** * UX-260616-07: 队列项立即发送(决策 a:插队=打断当前+发本条,复用 UX-05 stop 链路)。 * * 流程(顺序关键): * 1. 先 splice 出该项(记下 {text, skill})——必须在 stopChat 之前, * 否则 stopChat→AiCompleted→drainQueue 续发时该项还在队列致重复发。 * 2. 调 stopChat() 打断当前生成(stopChat UX-05 后已不清队列,故 splice 步骤前置必要)。 * 3. 直接 sendMessage(splicedItem.text, skill) 立即发这条(不等 drainQueue 轮到)。 * 当前未完成回复已由 stop 截断保留(agentic.rs:185 save_conversation),无需额外处理。 */ async function sendQueuedNow(index: number) { if (index < 0 || index >= state.queue.length) return const spliced = state.queue.splice(index, 1)[0] await stopChat() // B-260617-01: forceMode=true 跳过 L1 预检(ai_is_generating)直走 force_send。 // stopChat() 仅 await stop IPC 发出(置 stop_flag=true),不等后端 loop 跑到 // agentic.rs:444/:819 检测点+guard.reset()(generating 才复位)。紧接 sendMessage // 命中 L1 的 backendGenerating(后端真值仍 true)→ spliced 被入队而非立即发, // 与"立即发送"语义不符。force_send(commands.rs:742-748)原子复位 generating=false // 再发,无竞态窗口;stopChat 的 stop_flag 让旧 loop 在下个检测点退出,不冲突。 await sendMessage(spliced.text, spliced.skill, true, spliced.parts) } /** 停止当前生成:本地先复位 streaming,再发停止信号。 * 不依赖后端 AiCompleted 收尾——审批态看门狗已 clear(useAiEvents),若 AiCompleted 竞态丢失则 streaming 永久 true 卡死,故本地兜底。 * * UX-260617-19: stop IPC 失败时的兜底。本地已先置 streaming=false,但若 stop 信号未送达后端, * 后端会继续生成并发回 delta/Completed 事件,与用户已感知的「已停止」语义冲突;更坏情况是 * 后端 hang 既不收尾也不发事件(stop_flag 未置位、无 watchdog 兜底则无恢复路径)。 * 故 stop 失败时不 clearStreamWatchdog,反而重启看门狗:后端若继续生成,活跃 delta 会不断重置它 * 不误触发;后端若真 hang,130s 后 onStreamTimeout 兜底复位 streaming+清队列+推错误气泡。 * 同时推一条 stopLoopFailed 错误消息让用户知晓停止未生效(可再点或等待收尾)。 */ async function stopChat() { state.streaming = false // 本地先复位,不依赖后端 AiCompleted(审批态看门狗已 clear,竞态丢失则卡死) clearStreamWatchdog() // 清残留看门狗 // UX-260616-05 决策 a(逐条续发):不再清队列,保留供 AiCompleted 链式触发 drainQueue 续发。 // 后端 agentic.rs:190 emit AiCompleted 注释明确「保证前端收事件时后端已可接下一条,队列续发不被拒」; // 当前未完成回复已入库(agentic.rs:185 save_conversation)保留为截断回复。 try { // F-260616-09 B 批4(决策 e):传 activeConversationId,停止仅作用于当前 conv。 await aiApi.stopChat(state.activeConversationId) } catch (e) { // stop 信号未送达后端:重启看门狗(后端继续生成时活跃事件会重置它,真 hang 则 130s 后兜底收尾), // 避免后端仍在跑却无任何恢复路径(原 clearStreamWatchdog 已把兜底关掉)。 resetStreamWatchdog() const errMsg = e instanceof Error ? e.message : String(e) state.messages.push({ id: `stop-err-${nextMsgId()}`, role: 'assistant', content: t('aiChat.stopLoopFailed', { msg: errMsg }), isError: true, timestamp: Date.now(), } as AiMessage) console.error('[AI] stopChat IPC 失败,已重启看门狗兜底:', e) } } export function useAiSend() { return { sendMessage, regenerate, editMessage, approveToolCall, batchApprove, drainQueue, cancelQueued, clearQueue, editQueued, // UX-260616-06 sendQueuedNow, // UX-260616-07 stopChat, isQueueTimedOut, tryForceSend, // F-01 阶段6: 模型 override(主对话专用),顶部下拉读写 modelOverride, // 审批计时器已下沉到 aiShared.ts,此处 re-export 供组件使用 startApprovalTimer, clearApprovalTimer, clearAllApprovalTimers, } } // ARC-06: 模块级监听 ai-drain-queue 事件(由 useAiEvents AiCompleted case 发射) // 打破 useAiEvents ↔ useAiSend 循环依赖:events 不再直接 import 调用 drainQueue let _drainUnlisten: (() => void) | null = null export async function initDrainQueueListener(): Promise { if (_drainUnlisten) return _drainUnlisten = await listen('ai-drain-queue', () => { drainQueue() }) }