import { reactive, ref } from 'vue' /** * AG-UI 客户端(自实现 SSE,容错解析) * * 请求:POST 标准 AG-UI RunAgentInput(threadId/runId/messages/tools/context/state) * 响应:text/event-stream,逐条按 `data:{...}` 解析,根据 event.type 分发渲染。 * * 不使用 @ag-ui/client 的内置 verifyEvents 严格校验,避免后端事件顺序与其预期 * 不完全一致时整条流被中断、界面无输出。字段做 camelCase/snake_case 双写兼容。 */ export function useAgui() { // 展示用消息:{ id, role: 'user'|'assistant'|'tool'|'error', content } const messages = reactive([]) const running = ref(false) const error = ref('') // 发给后端的会话历史(仅 user / assistant 文本) const history = [] // 当前会话 ID(sessionId / threadId);开启新对话时会重新生成 const sessionId = ref(crypto.randomUUID()) const pick = (obj, ...keys) => { for (const k of keys) if (obj && obj[k] != null) return obj[k] return undefined } function findOrCreate(id, role) { let m = messages.find((x) => x.id === id) if (!m) { m = { id: id || crypto.randomUUID(), role, content: '' } messages.push(m) } return m } // 根据一条 AG-UI 事件更新界面 function dispatch(ev) { const type = (ev.type || ev.event || '').toUpperCase() switch (type) { case 'TEXT_MESSAGE_START': { const id = pick(ev, 'messageId', 'message_id') || crypto.randomUUID() findOrCreate(id, 'assistant') break } case 'TEXT_MESSAGE_CONTENT': { const id = pick(ev, 'messageId', 'message_id') const delta = pick(ev, 'delta', 'content') ?? '' findOrCreate(id, 'assistant').content += delta break } case 'TEXT_MESSAGE_END': { const id = pick(ev, 'messageId', 'message_id') const m = messages.find((x) => x.id === id) if (m) history.push({ id: m.id, role: 'assistant', content: m.content }) break } case 'TOOL_CALL_START': { const id = pick(ev, 'toolCallId', 'tool_call_id') || crypto.randomUUID() const name = pick(ev, 'toolCallName', 'tool_call_name') || '未知工具' findOrCreate(id, 'tool').content = `🔧 调用工具:${name}` break } case 'TOOL_CALL_ARGS': { const id = pick(ev, 'toolCallId', 'tool_call_id') const delta = pick(ev, 'delta', 'args') ?? '' const m = messages.find((x) => x.id === id) if (m) m.content += delta break } case 'TOOL_CALL_RESULT': { const id = pick(ev, 'toolCallId', 'tool_call_id') const content = pick(ev, 'content', 'result') ?? '' const m = messages.find((x) => x.id === id) if (m) m.content += `\n↳ 结果:${content}` break } case 'RUN_ERROR': { error.value = pick(ev, 'message', 'error') || '智能体运行出错' break } // RUN_STARTED / RUN_FINISHED / STATE_* / TOOL_CALL_END 等无需特殊处理 default: break } } // 解析 SSE 文本块(可能包含多条事件) function parseChunk(buffer) { // 以空行分隔事件;返回剩余未完成的尾部 const parts = buffer.split(/\r?\n\r?\n/) const rest = parts.pop() // 最后一段可能不完整,留到下次 for (const block of parts) { const dataLines = [] for (const raw of block.split(/\r?\n/)) { const line = raw.replace(/\r$/, '') if (!line || line.startsWith(':')) continue // 空行或注释/心跳 if (line.startsWith('data:')) { dataLines.push(line.slice(5).replace(/^ /, '')) // 去掉 data: 及一个可选空格 } // event:/id:/retry: 行忽略,type 从 data 的 JSON 里取 } if (!dataLines.length) continue const payload = dataLines.join('\n').trim() if (!payload || payload === '[DONE]') continue try { dispatch(JSON.parse(payload)) } catch (e) { // 单条解析失败不影响后续 console.warn('[AG-UI] 无法解析事件:', payload) } } return rest } async function send(url, text, userId) { const content = (text || '').trim() if (!content) return if (!url) { error.value = '请先选择或输入智能体 URL' return } if (running.value) return error.value = '' const userMsg = { id: crypto.randomUUID(), role: 'user', content } messages.push({ ...userMsg }) history.push(userMsg) const body = { threadId: sessionId.value, runId: crypto.randomUUID(), state: {}, messages: [userMsg], tools: [], context: [], forwardedProps: {} } running.value = true try { const headers = { 'Content-Type': 'application/json', Accept: 'text/event-stream' } if (userId) headers['x-user-id'] = userId if (sessionId.value) headers['x-session-id'] = sessionId.value const res = await fetch(url, { method: 'POST', headers, body: JSON.stringify(body) }) if (!res.ok) { error.value = `请求失败:HTTP ${res.status} ${res.statusText}` return } if (!res.body) { error.value = '响应无数据流(body 为空)' return } const reader = res.body.getReader() const decoder = new TextDecoder('utf-8') let buffer = '' // 逐块读取并解析 // eslint-disable-next-line no-constant-condition while (true) { const { value, done } = await reader.read() if (done) break buffer += decoder.decode(value, { stream: true }) buffer = parseChunk(buffer) } // 冲刷尾部(补一个空行确保最后一条事件被处理) parseChunk(buffer + '\n\n') } catch (e) { error.value = e?.message || String(e) } finally { running.value = false } } function reset() { messages.splice(0, messages.length) history.splice(0, history.length) error.value = '' sessionId.value = crypto.randomUUID() } // 开启新对话:清空界面与历史,并生成一个新的 sessionId function newChat() { reset() return sessionId.value } return { messages, running, error, sessionId, send, reset, newChat } }