| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198 |
- 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 }
- }
|