Browse Source

refactor(frontend): 管理端对话复用同次引用来源

feature/test
wei-py 3 weeks ago
parent
commit
d36f63a262
  1. 28
      frontend/src/api/chat.ts
  2. 15
      frontend/src/components/MessageSources.vue
  3. 137
      frontend/src/sdk-test/SdkTestPanel.vue
  4. 18
      frontend/src/types/sse.ts
  5. 3
      frontend/src/utils/chatAdapter.ts
  6. 269
      frontend/src/utils/sse.ts
  7. 105
      frontend/src/views/ChatPanel.vue
  8. 34
      frontend/src/views/PipelineFlow.vue
  9. 164
      frontend/tests/chat-protocol.test.mjs

28
frontend/src/api/chat.ts

@ -1,6 +1,5 @@
import request from './request'
import { getToken } from '@/utils/token' import { getToken } from '@/utils/token'
import type { ApiResponse } from '@/types/api'
import type { ChatResult } from '@/types/sse'
const API_BASE = '' const API_BASE = ''
@ -28,13 +27,11 @@ export interface ChatOptions {
* @param message 用户消息 * @param message 用户消息
* @param chatId 会话 ID * @param chatId 会话 ID
* @param options 可选参数 * @param options 可选参数
* @param includeEnableRag 是否写入 enableRag 参数(sources 接口无意义,需排除)
*/ */
function buildChatQueryParams( function buildChatQueryParams(
message: string, message: string,
chatId: string, chatId: string,
options?: ChatOptions, options?: ChatOptions,
includeEnableRag = true,
): URLSearchParams { ): URLSearchParams {
const params = new URLSearchParams() const params = new URLSearchParams()
params.set('message', message) params.set('message', message)
@ -42,7 +39,7 @@ function buildChatQueryParams(
if (options) { if (options) {
if (options.roleId) params.set('roleId', options.roleId) if (options.roleId) params.set('roleId', options.roleId)
if (options.accountId) params.set('accountId', options.accountId) if (options.accountId) params.set('accountId', options.accountId)
if (includeEnableRag && options.enableRag === true) params.set('enableRag', 'true')
if (options.enableRag === true) params.set('enableRag', 'true')
if (options.rewriteStrategy) params.set('rewriteStrategy', options.rewriteStrategy) if (options.rewriteStrategy) params.set('rewriteStrategy', options.rewriteStrategy)
if (options.categoryId !== undefined && options.categoryId !== '') params.set('categoryId', String(options.categoryId)) if (options.categoryId !== undefined && options.categoryId !== '') params.set('categoryId', String(options.categoryId))
if (options.categoryIds && options.categoryIds.length > 0) params.set('categoryIds', options.categoryIds.join(',')) if (options.categoryIds && options.categoryIds.length > 0) params.set('categoryIds', options.categoryIds.join(','))
@ -51,10 +48,17 @@ function buildChatQueryParams(
return params return params
} }
/** 同步对话:GET /ai/chat */
export function chatSync(message: string, chatId: string, options?: ChatOptions): Promise<string> {
const url = `${API_BASE}/ai/chat?${buildChatQueryParams(message, chatId, options).toString()}`
return fetch(url, { headers: authHeaders() }).then(res => res.text())
/** 同次回答和引用来源:GET /ai/chat/result(直接 JSON,无 data 包裹)。 */
export function chatSync(message: string, chatId: string, options?: ChatOptions, signal?: AbortSignal): Promise<ChatResult> {
const url = `${API_BASE}/ai/chat/result?${buildChatQueryParams(message, chatId, options).toString()}`
return fetchChatResult(url, authHeaders(), signal)
}
/** SDK 测试页也复用此响应处理,保留其独立域名与 SDK Token。 */
export async function fetchChatResult(url: string, headers: Record<string, string>, signal?: AbortSignal): Promise<ChatResult> {
const res = await fetch(url, { headers, signal })
if (!res.ok) throw new Error((await res.text()) || 'HTTP ' + res.status)
return res.json()
} }
/** 获取 SSE 流式对话 URL:GET /ai/chat/stream */ /** 获取 SSE 流式对话 URL:GET /ai/chat/stream */
@ -62,12 +66,6 @@ export function chatSSEUrl(message: string, chatId: string, options?: ChatOption
return `${API_BASE}/ai/chat/stream?${buildChatQueryParams(message, chatId, options).toString()}` return `${API_BASE}/ai/chat/stream?${buildChatQueryParams(message, chatId, options).toString()}`
} }
/** 获取 RAG 引用来源:GET /ai/chat/sources(不传 enableRag,其余参数照传) */
export function ragSources(message: string, chatId: string, options?: ChatOptions): Promise<ApiResponse> {
const path = `/ai/chat/sources?${buildChatQueryParams(message, chatId, options, false).toString()}`
return request.get(path).then(r => r.data)
}
/** /**
* 获取 AI 推荐问题列表(管理后台 ChatPanel 用)。 * 获取 AI 推荐问题列表(管理后台 ChatPanel 用)。
* 调用 /conversation/{id}/suggestions,由管理后台 JwtAuthFilter 守卫。 * 调用 /conversation/{id}/suggestions,由管理后台 JwtAuthFilter 守卫。

15
frontend/src/components/MessageSources.vue

@ -5,9 +5,9 @@
<div v-for="(group, key) in grouped" :key="key" class="source-group"> <div v-for="(group, key) in grouped" :key="key" class="source-group">
<div class="source-doc-title">📄 {{ group.title || '文档 ' + key }}</div> <div class="source-doc-title">📄 {{ group.title || '文档 ' + key }}</div>
<div class="source-chunks"> <div class="source-chunks">
<div v-for="chunk in group.chunks" :key="chunk.chunkIndex" class="source-chunk">
<span class="chunk-badge">#{{ chunk.chunkIndex }}</span>
<span v-html="renderMarkdown(chunk.content?.substring(0, 300) + (chunk.content?.length > 300 ? '...' : ''))"></span>
<div v-for="(chunk, index) in group.chunks" :key="index" class="source-chunk">
<span v-if="chunk.chunkIndex !== null" class="chunk-badge">#{{ chunk.chunkIndex }}</span>
<span v-html="renderMarkdown(chunk.snippet || '')"></span>
</div> </div>
</div> </div>
</div> </div>
@ -19,17 +19,18 @@
<script setup lang="ts"> <script setup lang="ts">
import { computed } from 'vue' import { computed } from 'vue'
import { renderMarkdown } from '@/utils/markdown' import { renderMarkdown } from '@/utils/markdown'
import type { SourceReference } from '@/types/sse'
const props = defineProps<{ sources: any[] }>()
const props = defineProps<{ sources: SourceReference[] }>()
/** 按文档 ID 归并 chunk */ /** 按文档 ID 归并 chunk */
const grouped = computed(() => { const grouped = computed(() => {
const groups: Record<string, { title?: string; chunks: any[] }> = {}
const groups: Record<string, { title: string | null; chunks: SourceReference[] }> = {}
for (const s of (props.sources || [])) { for (const s of (props.sources || [])) {
const docId = s.documentId || s.metadata?.documentId || 'unknown'
const docId = s.documentId || s.sourceName || s.title || 'unknown'
if (!groups[docId]) { if (!groups[docId]) {
groups[docId] = { groups[docId] = {
title: s.metadata?.title || s.sourceName,
title: s.title || s.sourceName,
chunks: [], chunks: [],
} }
} }

137
frontend/src/sdk-test/SdkTestPanel.vue

@ -282,13 +282,13 @@
</template> </template>
<script setup lang="ts"> <script setup lang="ts">
import { ref, reactive, computed, onMounted, onBeforeUnmount, nextTick } from 'vue'
import { ref, reactive, computed, watch, onMounted, onBeforeUnmount, nextTick } from 'vue'
import { ChatList, ChatSender, ChatActionbar } from '@tdesign-vue-next/chat' import { ChatList, ChatSender, ChatActionbar } from '@tdesign-vue-next/chat'
import '@tdesign-vue-next/chat/es/style/index.css' import '@tdesign-vue-next/chat/es/style/index.css'
import { MessagePlugin } from 'tdesign-vue-next' import { MessagePlugin } from 'tdesign-vue-next'
import { readSSEStream } from '@/utils/sse'
import { readSSEStreamWithEvents } from '@/utils/sse'
import { fetchChatResult } from '@/api/chat'
import { renderMarkdown } from '@/utils/markdown' import { renderMarkdown } from '@/utils/markdown'
import { toast } from '@/utils/toast'
import { toChatData, type ChatMessage } from '@/utils/chatAdapter' import { toChatData, type ChatMessage } from '@/utils/chatAdapter'
// ==================== 常量 ==================== // ==================== 常量 ====================
@ -309,11 +309,11 @@ const modeOptions = [
{ label: '同步调用', value: 'sync' }, { label: '同步调用', value: 'sync' },
] ]
const strategyOptions = [ const strategyOptions = [
{ label: '不重写', value: 'NONE' },
{ label: '查询重写', value: 'REWRITE' },
{ label: '翻译扩展', value: 'TRANSLATION' },
{ label: '查询压缩', value: 'COMPRESSION' },
{ label: '多路扩展', value: 'MULTI_QUERY' },
{ label: '不重写(最快)', value: 'NONE' },
{ label: '查询重写(额外耗时)', value: 'REWRITE' },
{ label: '翻译扩展(额外耗时)', value: 'TRANSLATION' },
{ label: '查询压缩(额外耗时)', value: 'COMPRESSION' },
{ label: '多路扩展(额外耗时)', value: 'MULTI_QUERY' },
] ]
// SDK 可选参数的默认值(与 SDK parseConfig 对齐),用于判断是否需要在生成代码中输出 // SDK 可选参数的默认值(与 SDK parseConfig 对齐),用于判断是否需要在生成代码中输出
@ -609,7 +609,7 @@ function domAction(act: 'open' | 'close' | 'toggle'): void {
const activeTab = ref<'chat' | 'sdk'>('chat') const activeTab = ref<'chat' | 'sdk'>('chat')
const demoMode = ref<'sse' | 'sync'>('sse') const demoMode = ref<'sse' | 'sync'>('sse')
const demoRag = ref(true) const demoRag = ref(true)
const demoStrategy = ref('MULTI_QUERY')
const demoStrategy = ref('NONE')
const demoInput = ref('') const demoInput = ref('')
const demoChatId = ref('') const demoChatId = ref('')
const isSending = ref(false) const isSending = ref(false)
@ -626,7 +626,7 @@ function welcomeMessage(): ChatMessage {
return { return {
id: genId(), id: genId(),
role: 'assistant', role: 'assistant',
content: '您好,这里是新接口(GET /ai/chat、/ai/chat/stream、/ai/chat/sources)对话演示区。请先在左侧获取 SDK Token,然后输入问题开始对话。',
content: '您好,这里是对话演示区(GET /ai/chat/result、/ai/chat/stream),回答会同时返回实际引用来源。请先在左侧获取 SDK Token,然后输入问题开始对话。',
streaming: false, streaming: false,
time: fmtTime(), time: fmtTime(),
} }
@ -636,7 +636,7 @@ const demoMessages = ref<ChatMessage[]>([welcomeMessage()])
const demoData = computed(() => toChatData(demoMessages.value)) const demoData = computed(() => toChatData(demoMessages.value))
function buildDemoUrl( function buildDemoUrl(
path: 'chat' | 'chat/stream' | 'chat/sources',
path: 'chat/result' | 'chat/stream',
message: string, message: string,
chatId: string, chatId: string,
opts: { enableRag: boolean; rewriteStrategy?: string }, opts: { enableRag: boolean; rewriteStrategy?: string },
@ -648,7 +648,7 @@ function buildDemoUrl(
if (rid) p.set('roleId', rid) if (rid) p.set('roleId', rid)
const uid = config.userId.trim() const uid = config.userId.trim()
if (uid) p.set('accountId', uid) if (uid) p.set('accountId', uid)
if (path !== 'chat/sources' && opts.enableRag) p.set('enableRag', 'true')
if (opts.enableRag) p.set('enableRag', 'true')
if (opts.enableRag && opts.rewriteStrategy) p.set('rewriteStrategy', opts.rewriteStrategy) if (opts.enableRag && opts.rewriteStrategy) p.set('rewriteStrategy', opts.rewriteStrategy)
return demoDomain.value + '/ai/' + path + '?' + p.toString() return demoDomain.value + '/ai/' + path + '?' + p.toString()
} }
@ -679,48 +679,45 @@ async function sendDemo(val?: string): Promise<void> {
const requestMode = demoMode.value const requestMode = demoMode.value
const opts = { enableRag: demoRag.value, rewriteStrategy: demoRag.value ? demoStrategy.value : undefined } const opts = { enableRag: demoRag.value, rewriteStrategy: demoRag.value ? demoStrategy.value : undefined }
// 提前构建 URL,固定本轮角色、账号、域名和检索配置。 // 提前构建 URL,固定本轮角色、账号、域名和检索配置。
const url = buildDemoUrl(requestMode === 'sync' ? 'chat' : 'chat/stream', text, cid, opts)
const sourcesUrl = opts.enableRag ? buildDemoUrl('chat/sources', text, cid, opts) : ''
const url = buildDemoUrl(requestMode === 'sync' ? 'chat/result' : 'chat/stream', text, cid, opts)
const headers = sdkAuthHeaders() const headers = sdkAuthHeaders()
const controller = new AbortController()
abortController?.abort()
abortController = controller
const isCurrent = () => abortController === controller
&& !controller.signal.aborted && demoChatId.value === cid
&& demoMessages.value.some(msg => msg.id === assistantMsg.id)
await scrollDemoBottom() await scrollDemoBottom()
try { try {
if (abortController) abortController.abort()
abortController = new AbortController()
if (!isCurrent()) return
if (requestMode === 'sync') { if (requestMode === 'sync') {
const res = await fetch(url, { headers, signal: abortController.signal })
if (!res.ok) throw new Error('HTTP ' + res.status)
assistantMsg.content = await res.text()
const result = await fetchChatResult(url, headers, controller.signal)
if (!isCurrent()) return
assistantMsg.content = result.text
assistantMsg.sources = result.sources
} else { } else {
// 复用 readSSEStream(内部手写解析 OpenAI delta.content),逐 chunk 追加
await readSSEStream(
url,
(chunk: string) => {
await readSSEStreamWithEvents(url, {
onMessage: (chunk) => {
if (!isCurrent()) return
assistantMsg.content += chunk assistantMsg.content += chunk
demoMessages.value = [...demoMessages.value] demoMessages.value = [...demoMessages.value]
void scrollDemoBottom()
}, },
undefined,
headers,
abortController.signal,
)
}
// 引用来源后台补齐,正文完成后立即结束发送状态。
if (opts.enableRag) {
void fetch(sourcesUrl, { headers }).then(async sres => {
const json = await sres.json()
if (!sres.ok || !json?.success) {
throw new Error(json?.message || (sres.ok ? '引用来源获取失败' : 'HTTP ' + sres.status))
}
if (demoChatId.value !== cid || !demoMessages.value.some(msg => msg.id === assistantMsg.id)) return
assistantMsg.sources = json.data || []
onSources: (sources) => {
if (!isCurrent()) return
assistantMsg.sources = sources
demoMessages.value = [...demoMessages.value] demoMessages.value = [...demoMessages.value]
}).catch((e: any) => {
toast('引用来源获取失败: ' + (e.message || e), 'warning')
})
},
onError: (data) => {
if (!isCurrent()) return
assistantMsg.content += '\n\n' + (data.message || '工具调用出错')
demoMessages.value = [...demoMessages.value]
},
}, headers, controller.signal)
} }
} catch (e: any) { } catch (e: any) {
if (!isCurrent()) return
if (e.name === 'AbortError') { if (e.name === 'AbortError') {
assistantMsg.content = assistantMsg.content || '已取消' assistantMsg.content = assistantMsg.content || '已取消'
} else { } else {
@ -728,11 +725,13 @@ async function sendDemo(val?: string): Promise<void> {
assistantMsg.error = true assistantMsg.error = true
} }
} finally { } finally {
if (isCurrent()) {
assistantMsg.streaming = false assistantMsg.streaming = false
isSending.value = false isSending.value = false
demoMessages.value = [...demoMessages.value] demoMessages.value = [...demoMessages.value]
await scrollDemoBottom() await scrollDemoBottom()
} }
}
} }
function abortDemo(): void { function abortDemo(): void {
@ -752,6 +751,13 @@ function clearDemo(): void {
demoMessages.value = [welcomeMessage()] demoMessages.value = [welcomeMessage()]
} }
// Changing the SDK identity starts a new conversation and invalidates pending callbacks.
watch(
[() => config.integrateId, () => config.userId, demoDomain, sdkToken],
clearDemo,
{ flush: 'sync' },
)
async function onDemoAction(action: string, index: number): Promise<void> { async function onDemoAction(action: string, index: number): Promise<void> {
const msg = demoMessages.value[index] const msg = demoMessages.value[index]
if (!msg) return if (!msg) return
@ -922,41 +928,30 @@ const testResults = ref<TestCaseResult[]>([
const url = buildDemoUrl('chat/stream', '你好', 'verify_sse', { enableRag: false }) const url = buildDemoUrl('chat/stream', '你好', 'verify_sse', { enableRag: false })
apiCount.value++ apiCount.value++
const t0 = performance.now() const t0 = performance.now()
let res: Response
try {
res = await fetch(url, { headers: sdkAuthHeaders(), signal: AbortSignal.timeout(20000) })
} catch {
log('⚠ 请求失败', 'warn')
mark('skip', '后端不可用')
return
}
if (!res.ok) {
log('⚠ HTTP ' + res.status, 'warn')
mark('skip', '接口异常')
return
}
const reader = res.body!.getReader()
const decoder = new TextDecoder()
let total = '' let total = ''
let chunks = 0 let chunks = 0
while (true) {
const { done, value } = await reader.read()
if (done) break
total += decoder.decode(value, { stream: true })
chunks++
}
let metadataChunks = 0
let doneCalls = 0
await readSSEStreamWithEvents(url, {
onMessage: (text) => { total += text; chunks++ },
onSources: (sources) => {
metadataChunks++
assert(sources.length === 0, '普通对话应返回空 sources')
},
onDone: () => { doneCalls++ },
}, sdkAuthHeaders(), AbortSignal.timeout(20000))
const elapsed = Math.round(performance.now() - t0) const elapsed = Math.round(performance.now() - t0)
apiDurations.value.push(elapsed) apiDurations.value.push(elapsed)
assert(total.includes('choices') || total.includes('[DONE]'), '非 OpenAI Chat Completions 格式')
log('✓ OpenAI 格式(含 choices 或 [DONE])', 'pass')
log('✓ SSE ' + chunks + ' chunks, ' + total.length + ' chars', 'pass')
assert(metadataChunks === 1 && doneCalls === 1, '应收到一次引用元数据和一次完成事件')
log('OpenAI 正文和来源元数据分别解析,完成事件只触发一次', 'pass')
log('SSE ' + chunks + ' chunks, ' + total.length + ' chars', 'pass')
mark('pass', '通过 (' + elapsed + 'ms)') mark('pass', '通过 (' + elapsed + 'ms)')
}, },
}, },
{ {
id: 'T8', name: 'API RAG 引用来源', desc: 'GET /ai/chat/sources', phase: 'p1', status: 'idle', summary: '', logs: [],
id: 'T8', name: 'API 回答引用来源', desc: 'GET /ai/chat/result:同次回答携带 sources', phase: 'p1', status: 'idle', summary: '', logs: [],
fn: async (log, mark) => { fn: async (log, mark) => {
const url = buildDemoUrl('chat/sources', '请假', 'verify_sources', { enableRag: true, rewriteStrategy: 'REWRITE' })
const url = buildDemoUrl('chat/result', '请假', 'verify_sources', { enableRag: true, rewriteStrategy: 'NONE' })
apiCount.value++ apiCount.value++
let res: Response let res: Response
try { try {
@ -972,9 +967,9 @@ const testResults = ref<TestCaseResult[]>([
return return
} }
const json = await res.json() const json = await res.json()
log('返回 success=' + json.success + ' data=' + (json.data ? json.data.length : 0) + ' 条', 'info')
assert(json.success !== undefined, '应返回 success 字段')
log('✓ RAG 引用来源接口可用', 'pass')
log('返回回答 ' + (json.text?.length || 0) + ' 字,引用 ' + (json.sources?.length || 0) + ' 条', 'info')
assert(typeof json.text === 'string' && Array.isArray(json.sources), '应直接返回 text 和 sources 字段')
log('回答与引用来源由同一请求返回', 'pass')
mark('pass', '通过') mark('pass', '通过')
}, },
}, },

18
frontend/src/types/sse.ts

@ -1,6 +1,24 @@
/** 同次回答实际命中的知识库分块;ID 保持字符串精度。 */
export interface SourceReference {
documentId: string | null
title: string | null
sourceName: string | null
chunkIndex: number | null
score: number | null
snippet: string | null
}
export interface ChatResult {
text: string
mcpEvents: unknown[]
suggestions: string[]
sources: SourceReference[]
}
/** SSE 事件类型定义 */ /** SSE 事件类型定义 */
export interface SSECallbacks { export interface SSECallbacks {
onMessage?: (chunk: string) => void onMessage?: (chunk: string) => void
onSources?: (sources: SourceReference[]) => void
onToolCallStart?: (data: any) => void onToolCallStart?: (data: any) => void
onToolCallResult?: (data: any) => void onToolCallResult?: (data: any) => void
onError?: (data: any) => void onError?: (data: any) => void

3
frontend/src/utils/chatAdapter.ts

@ -7,6 +7,7 @@
*/ */
import type { TdChatItemMeta, AIMessageContent, UserMessageContent } from '@tdesign-vue-next/chat' import type { TdChatItemMeta, AIMessageContent, UserMessageContent } from '@tdesign-vue-next/chat'
import type { SourceReference } from '@/types/sse'
/** 附件信息(图片或文件) */ /** 附件信息(图片或文件) */
export interface Attachment { export interface Attachment {
@ -24,7 +25,7 @@ export interface ChatMessage {
content: string content: string
streaming: boolean streaming: boolean
time: string time: string
sources?: any[]
sources?: SourceReference[]
toolCalls?: any[] toolCalls?: any[]
error?: boolean error?: boolean
feedback?: string | null feedback?: string | null

269
frontend/src/utils/sse.ts

@ -1,206 +1,139 @@
/**
* SSE 流式读取工具 —— 从 utils.js 原封不动搬移
*
* 统一处理 Flux<String> / ServerSentEvent / SseEmitter 三种 SSE 接口。
* 此文件使用原生 fetch + ReadableStream API,零框架依赖。
*/
import type { SSECallbacks } from '@/types/sse'
/** Shared SSE reader for plain text, OpenAI chunks and tool events. */
import type { SSECallbacks, SourceReference } from '@/types/sse'
import { getToken } from '@/utils/token' import { getToken } from '@/utils/token'
/** 构建带认证的请求头 */
/** Preserve explicit SDK authorization instead of replacing it with the admin token. */
function authHeaders(extra?: Record<string, string>): Record<string, string> { function authHeaders(extra?: Record<string, string>): Record<string, string> {
const token = getToken() const token = getToken()
const base: Record<string, string> = extra ? { ...extra } : {}
// 仅当调用方未显式提供 Authorization 时才注入管理后台 token,避免覆盖调用方传入的鉴权头(如测试面板的 SDK Token)
if (token && !base['Authorization']) base['Authorization'] = `Bearer ${token}`
return base
const headers = { ...extra }
if (token && !headers['Authorization']) headers['Authorization'] = `Bearer ${token}`
return headers
} }
/** 尝试从 OpenAI Chat Completions chunk 中提取 delta.content;返回 undefined 表示非 OpenAI 格式(回退纯文本) */
function extractOpenAIDelta(text: string): string | undefined {
let obj: any
try { obj = JSON.parse(text) } catch { return undefined }
// OpenAI 错误 chunk({"error":{"message":...}}):提取错误信息作为内容,避免原始 JSON 泄漏到对话框
if (obj && obj.error) {
const msg = obj.error.message || obj.error.type
return typeof msg === 'string' && msg ? msg : '服务异常'
function isSourceReference(value: unknown): value is SourceReference {
if (!value || typeof value !== 'object') return false
const nullableString = (v: unknown) => v === null || typeof v === 'string'
const nullableNumber = (v: unknown) => v === null || typeof v === 'number'
return 'documentId' in value && nullableString(value.documentId)
&& 'title' in value && nullableString(value.title)
&& 'sourceName' in value && nullableString(value.sourceName)
&& 'chunkIndex' in value && nullableNumber(value.chunkIndex)
&& 'score' in value && nullableNumber(value.score)
&& 'snippet' in value && nullableString(value.snippet)
}
/** Returns false only for legacy plain-text payloads. Metadata never enters message text. */
function dispatchOpenAI(text: string, handlers: SSECallbacks): boolean {
let value: unknown
try { value = JSON.parse(text) } catch { return false }
if (!value || typeof value !== 'object') return false
if ('error' in value && value.error && typeof value.error === 'object') {
const error = value.error
const message = 'message' in error ? error.message : 'type' in error ? error.type : undefined
handlers.onMessage?.(typeof message === 'string' && message ? message : '服务异常')
return true
} }
if (obj && Array.isArray(obj.choices)) {
const delta = obj.choices[0]?.delta
if (delta && typeof delta.content === 'string') return delta.content
return '' // OpenAI 形状但无 content(role/finish_reason 空 chunk)→ 跳过
if (!('choices' in value) || !Array.isArray(value.choices)) return false
if ('sources' in value && Array.isArray(value.sources)) {
if (!value.sources.every(isSourceReference)) throw new Error('无效的引用来源数据')
handlers.onSources?.(value.sources)
} }
return undefined // 是 JSON 但非 OpenAI 形状 → 回退纯文本
const content = value.choices[0]?.delta?.content
if (typeof content === 'string' && content) handlers.onMessage?.(content)
return true
} }
/**
* 通用 SSE 流式读取 —— 统一处理 Flux<String> / ServerSentEvent / SseEmitter 三种 SSE 接口
*
* @param url 请求地址
* @param onChunk 每收到一段文本的回调
* @param onDone 流结束的回调
* @param headers 额外请求头
* @param signal AbortSignal 用于取消请求(组件卸载时必须传入以释放网络资源)
*/
export async function readSSEStream(
/** Text-only callers share the same framing, completion and cleanup semantics. */
export function readSSEStream(
url: string, url: string,
onChunk: (text: string) => void, onChunk: (text: string) => void,
onDone?: () => void, onDone?: () => void,
headers?: Record<string, string>, headers?: Record<string, string>,
signal?: AbortSignal
signal?: AbortSignal,
): Promise<void> { ): Promise<void> {
const res = await fetch(url, { headers: authHeaders(headers), signal, credentials: 'include' })
if (!res.ok) throw new Error('HTTP ' + res.status)
const reader = res.body!.getReader()
const decoder = new TextDecoder()
let buffer = ''
// SSE 规范:同一事件内的多行 data: 字段用 \n 拼接,事件间用空行分隔
let eventDataLines: string[] = []
let currentEvent = 'message'
const flushEvent = () => {
if (eventDataLines.length === 0) return
const text = eventDataLines.join('\n')
eventDataLines = []
const ev = currentEvent
currentEvent = 'message'
if (text === '[DONE]') return
// 跳过 status / faq 等系统事件,不显示在对话框中
if (ev === 'status') return
// OpenAI Chat Completions 格式:提取 delta.content;无 content 的空 chunk 跳过
const extracted = extractOpenAIDelta(text)
if (extracted !== undefined) {
if (extracted) onChunk(extracted)
return
}
// 空事件视为 LLM 流式输出的换行符(Spring 将 "\n" 编码为单条空 data: 事件)
onChunk(text || '\n')
}
while (true) {
const { done, value } = await reader.read()
if (done) break
buffer += decoder.decode(value, { stream: true })
const lines = buffer.split('\n')
buffer = lines.pop() || ''
for (let line of lines) {
// 兼容 \r\n 行结束符
if (line.endsWith('\r')) line = line.slice(0, -1)
if (line === '') {
// 空行 = SSE 事件边界
flushEvent()
} else if (line.startsWith('event:')) {
// 记录事件类型,用于 flushEvent 时过滤系统事件
currentEvent = line.slice(6).trim()
} else if (line.startsWith('data:')) {
// 累积 data 字段;仅剥离 data: 后的一个可选空格,保留 markdown 列表缩进
let data = line.slice(5)
if (data.startsWith(' ')) data = data.slice(1)
eventDataLines.push(data)
} else if (!line.startsWith(':')) {
// Flux<String> 模式(非标准 SSE),先把已累积的 SSE 事件 flush 再处理
flushEvent()
if (line.trim()) onChunk(line)
}
}
}
// 流结束,flush 末尾未以空行收尾的事件
flushEvent()
if (onDone) onDone()
return readSSEStreamWithEvents(url, { onMessage: onChunk, onDone }, headers, signal)
} }
/**
* 增强版 SSE 流式读取(支持事件类型分发)
* 解析 SSE 标准的 event: 字段,将不同类型事件分发到对应回调。
*
* @param url 请求地址
* @param handlers 回调对象:
* - onMessage(chunk): 普通文本内容(event: message 或无 event 的 data)
* - onToolCallStart(data): 工具调用开始(event: tool_call_start)
* - onToolCallResult(data): 工具调用结果(event: tool_call_result)
* - onError(data): 错误事件(event: error)
* - onDone(): 流结束
* @param headers 额外请求头
* @param signal AbortSignal 用于取消请求(组件卸载时必须传入以释放网络资源)
*/
/** Read complete SSE events; [DONE] terminates immediately without waiting for network EOF. */
export async function readSSEStreamWithEvents( export async function readSSEStreamWithEvents(
url: string, url: string,
handlers: SSECallbacks, handlers: SSECallbacks,
headers?: Record<string, string>, headers?: Record<string, string>,
signal?: AbortSignal
signal?: AbortSignal,
): Promise<void> { ): Promise<void> {
const { onMessage, onToolCallStart, onToolCallResult, onError, onDone } = handlers
const res = await fetch(url, { headers: authHeaders(headers), signal, credentials: 'include' }) const res = await fetch(url, { headers: authHeaders(headers), signal, credentials: 'include' })
if (!res.ok) throw new Error('HTTP ' + res.status)
const reader = res.body!.getReader()
if (!res.ok) throw new Error((await res.text()) || 'HTTP ' + res.status)
if (!res.body) throw new Error('响应流为空')
const reader = res.body.getReader()
const decoder = new TextDecoder() const decoder = new TextDecoder()
let buffer = '' let buffer = ''
let currentEvent = 'message'
// SSE 规范:同一事件内的多行 data: 字段用 \n 拼接
let eventDataLines: string[] = []
let event = 'message'
let data: string[] = []
let completed = false
let eof = false
const flushEvent = () => { const flushEvent = () => {
if (eventDataLines.length === 0) return
const ev = currentEvent
currentEvent = 'message'
const raw = eventDataLines.join('\n')
eventDataLines = []
if (raw === '[DONE]') return
// 空事件视为换行符(仅 message 事件);JSON 事件(工具调用等)空内容不应出现
const text = raw || '\n'
switch (ev) {
case 'tool_call_start':
if (onToolCallStart) onToolCallStart(JSON.parse(text))
break
case 'tool_call_result':
if (onToolCallResult) onToolCallResult(JSON.parse(text))
break
case 'error':
if (onError) onError(JSON.parse(text))
break
case 'status':
// 系统状态事件(generating / faq_hit 等),不显示在对话框中
break
case 'message':
default: {
// OpenAI Chat Completions 格式:提取 delta.content;无 content 的空 chunk 跳过
const extracted = extractOpenAIDelta(text)
if (extracted !== undefined) {
if (extracted && onMessage) onMessage(extracted)
} else if (onMessage) {
onMessage(text)
}
break
const currentEvent = event
event = 'message'
if (!data.length) return
const raw = data.join('\n')
data = []
if (raw === '[DONE]') {
completed = true
return
} }
switch (currentEvent) {
case 'status': return
case 'tool_call_start': handlers.onToolCallStart?.(JSON.parse(raw)); return
case 'tool_call_result': handlers.onToolCallResult?.(JSON.parse(raw)); return
case 'error': handlers.onError?.(JSON.parse(raw)); return
default:
if (!dispatchOpenAI(raw, handlers)) handlers.onMessage?.(raw || '\n')
} }
} }
while (true) {
const { done, value } = await reader.read()
if (done) break
buffer += decoder.decode(value, { stream: true })
const lines = buffer.split('\n')
buffer = lines.pop() || ''
for (let line of lines) {
if (line.endsWith('\r')) line = line.slice(0, -1)
if (line === '') {
// 空行 = SSE 事件边界
const consumeLine = (raw: string) => {
const line = raw.endsWith('\r') ? raw.slice(0, -1) : raw
if (!line) {
flushEvent() flushEvent()
} else if (line.startsWith('event:')) { } else if (line.startsWith('event:')) {
currentEvent = line.slice(6).trim()
} else if (line.startsWith('data:')) {
let data = line.slice(5)
if (data.startsWith(' ')) data = data.slice(1)
eventDataLines.push(data)
} else if (!line.startsWith(':')) {
// Flux<String> 模式(非标准 SSE)
event = line.slice(6).trim()
} else if (line === 'data' || line.startsWith('data:')) {
const value = line === 'data' ? '' : line.slice(5)
data.push(value.startsWith(' ') ? value.slice(1) : value)
} else if (!line.startsWith(':') && !line.startsWith('id:') && !line.startsWith('retry:')) {
// Legacy Flux<String> bodies may contain unframed lines.
flushEvent() flushEvent()
if (line.trim() && onMessage) onMessage(line)
if (!completed && line.trim()) handlers.onMessage?.(line)
} }
} }
try {
while (!completed) {
signal?.throwIfAborted()
const result = await reader.read()
signal?.throwIfAborted()
eof = result.done
buffer += eof ? decoder.decode() : decoder.decode(result.value, { stream: true })
let start = 0
let end: number
while (!completed && (end = buffer.indexOf('\n', start)) !== -1) {
consumeLine(buffer.slice(start, end))
start = end + 1
}
buffer = buffer.slice(start)
if (eof) {
if (!completed && buffer) consumeLine(buffer)
if (!completed) flushEvent()
break
}
}
handlers.onDone?.()
} finally {
// Cancellation may itself reject or stall; it must not delay DONE or mask the original error.
if (!eof) {
try { void reader.cancel().catch(() => {}) } catch { /* Preserve the stream outcome. */ }
}
try { reader.releaseLock() } catch { /* Preserve the stream outcome. */ }
} }
flushEvent()
if (onDone) onDone()
} }

105
frontend/src/views/ChatPanel.vue

@ -235,14 +235,14 @@ import {
ChatActionbar, ChatActionbar,
} from '@tdesign-vue-next/chat' } from '@tdesign-vue-next/chat'
import '@tdesign-vue-next/chat/es/style/index.css' import '@tdesign-vue-next/chat/es/style/index.css'
import { chatSync, chatSSEUrl, ragSources, fetchSuggestions, type ChatOptions } from '@/api/chat'
import { chatSync, chatSSEUrl, fetchSuggestions, type ChatOptions } from '@/api/chat'
import { getRoleList } from '@/api/role' import { getRoleList } from '@/api/role'
import { getActiveModelConfig } from '@/api/model-config' import { getActiveModelConfig } from '@/api/model-config'
import { truncateConversation } from '@/api/conversation' import { truncateConversation } from '@/api/conversation'
import { submitFeedback as submitFeedbackApi } from '@/api/feedback' import { submitFeedback as submitFeedbackApi } from '@/api/feedback'
import { uploadAttachment } from '@/api/upload' import { uploadAttachment } from '@/api/upload'
import { toast } from '@/utils/toast' import { toast } from '@/utils/toast'
import { readSSEStream, readSSEStreamWithEvents } from '@/utils/sse'
import { readSSEStreamWithEvents } from '@/utils/sse'
import { renderMarkdown } from '@/utils/markdown' import { renderMarkdown } from '@/utils/markdown'
import { useCategoryStore } from '@/stores/category' import { useCategoryStore } from '@/stores/category'
import { toChatData, isLastAssistant } from '@/utils/chatAdapter' import { toChatData, isLastAssistant } from '@/utils/chatAdapter'
@ -266,11 +266,11 @@ const QUICK_QUESTIONS = [
] ]
const ragStrategyOptions = [ const ragStrategyOptions = [
{ label: '不重写', value: 'NONE' },
{ label: '查询重写', value: 'REWRITE' },
{ label: '翻译扩展', value: 'TRANSLATION' },
{ label: '查询压缩', value: 'COMPRESSION' },
{ label: '多路扩展', value: 'MULTI_QUERY' },
{ label: '不重写(最快)', value: 'NONE' },
{ label: '查询重写(额外耗时)', value: 'REWRITE' },
{ label: '翻译扩展(额外耗时)', value: 'TRANSLATION' },
{ label: '查询压缩(额外耗时)', value: 'COMPRESSION' },
{ label: '多路扩展(额外耗时)', value: 'MULTI_QUERY' },
] ]
const modeOptions = [ const modeOptions = [
@ -284,7 +284,7 @@ const mode = ref('sse') // 默认 SSE 流式
const selectedRole = ref('general') const selectedRole = ref('general')
const roles = ref([FALLBACK_ROLE]) const roles = ref([FALLBACK_ROLE])
const isRagMode = ref(false) const isRagMode = ref(false)
const ragStrategy = ref('MULTI_QUERY')
const ragStrategy = ref('NONE')
const activeModel = ref<any>(null) const activeModel = ref<any>(null)
const modelLoadError = ref('') const modelLoadError = ref('')
const userInput = ref('') const userInput = ref('')
@ -376,6 +376,7 @@ function providerLabel(provider: string): string {
// ==================== 角色管理 ==================== // ==================== 角色管理 ====================
function selectRole(roleKey: string): void { function selectRole(roleKey: string): void {
abortChat()
selectedRole.value = roleKey selectedRole.value = roleKey
newChatId() newChatId()
currentSuggestions.value = [] // 切换角色时清空推荐问题 currentSuggestions.value = [] // 切换角色时清空推荐问题
@ -494,7 +495,7 @@ async function send(): Promise<void> {
const cid = chatId.value || ('web_' + Date.now()) const cid = chatId.value || ('web_' + Date.now())
chatId.value = cid chatId.value = cid
const requestMode = mode.value const requestMode = mode.value
// 固定本轮参数,后台引用请求不读取后续切换的角色或 RAG 配置。
// 固定本轮参数,后续角色或配置切换不影响已发出的回答。
const imageUrls = attachments.filter(a => a.type === 'image').map(a => a.url) const imageUrls = attachments.filter(a => a.type === 'image').map(a => a.url)
const chatOptions: ChatOptions = { const chatOptions: ChatOptions = {
roleId: currentRoleId(), roleId: currentRoleId(),
@ -516,64 +517,56 @@ async function send(): Promise<void> {
streaming: true, time: formatTime(), sources: [], toolCalls: [], streaming: true, time: formatTime(), sources: [], toolCalls: [],
} }
messages.value.push(assistantMsg) messages.value.push(assistantMsg)
const controller = new AbortController()
sseAbortController?.abort()
sseAbortController = controller
const isCurrent = () => sseAbortController === controller
&& !controller.signal.aborted && chatId.value === cid
&& messages.value.some(msg => msg.id === assistantMsg.id)
await scrollToBottom() await scrollToBottom()
try { try {
// 取消上一个 SSE 请求并创建新的 AbortController
if (sseAbortController) { sseAbortController.abort() }
sseAbortController = new AbortController()
const signal = sseAbortController.signal
if (!isCurrent()) return
if (requestMode === 'sync') { if (requestMode === 'sync') {
// 同步调用
assistantMsg.content = await chatSync(text, cid, chatOptions)
const result = await chatSync(text, cid, chatOptions, controller.signal)
if (!isCurrent()) return
assistantMsg.content = result.text
assistantMsg.sources = result.sources
} else { } else {
// SSE 流式
const url = chatSSEUrl(text, cid, chatOptions)
if (chatOptions.enableRag) {
await readSSEStreamWithEvents(url, {
onMessage: async (chunk: string) => {
await readSSEStreamWithEvents(chatSSEUrl(text, cid, chatOptions), {
onMessage: (chunk) => {
if (!isCurrent()) return
assistantMsg.content += chunk assistantMsg.content += chunk
messages.value = [...messages.value] messages.value = [...messages.value]
await scrollToBottom()
void scrollToBottom()
}, },
onToolCallStart: (data: any) => {
onSources: (sources) => {
if (!isCurrent()) return
assistantMsg.sources = sources
messages.value = [...messages.value]
},
onToolCallStart: (data) => {
if (!isCurrent()) return
assistantMsg.toolCalls!.push({ tool: data.tool, input: data.input, status: 'running', result: null }) assistantMsg.toolCalls!.push({ tool: data.tool, input: data.input, status: 'running', result: null })
scrollToBottom()
messages.value = [...messages.value]
void scrollToBottom()
}, },
onToolCallResult: (data: any) => {
onToolCallResult: (data) => {
if (!isCurrent()) return
const tc = assistantMsg.toolCalls!.find(t => t.tool === data.tool && t.status === 'running') const tc = assistantMsg.toolCalls!.find(t => t.tool === data.tool && t.status === 'running')
if (tc) { tc.status = 'done'; tc.result = data.result; tc.latencyMs = data.latencyMs } if (tc) { tc.status = 'done'; tc.result = data.result; tc.latencyMs = data.latencyMs }
messages.value = [...messages.value] messages.value = [...messages.value]
scrollToBottom()
void scrollToBottom()
}, },
onError: (data: any) => {
assistantMsg.content += '\n\n⚠️ ' + (data.message || '工具调用出错')
onError: (data) => {
if (!isCurrent()) return
assistantMsg.content += '\n\n' + (data.message || '工具调用出错')
messages.value = [...messages.value] messages.value = [...messages.value]
}, },
onDone: () => {},
}, undefined, signal)
} else {
await readSSEStream(url, async (chunk: string) => {
assistantMsg.content += chunk
messages.value = [...messages.value]
await scrollToBottom()
}, () => {}, undefined, signal)
}
}
// 引用来源后台补齐,正文完成后立即结束发送状态。
if (chatOptions.enableRag) {
void ragSources(text, cid, chatOptions).then(sj => {
if (!sj?.success) throw new Error(sj?.message || '引用来源获取失败')
if (chatId.value !== cid || !messages.value.some(msg => msg.id === assistantMsg.id)) return
assistantMsg.sources = sj.data || []
messages.value = [...messages.value]
}).catch((e: any) => {
toast('引用来源获取失败: ' + (e.message || e), 'warning')
})
}, undefined, controller.signal)
} }
} catch (e: any) { } catch (e: any) {
if (!isCurrent()) return
// AbortError 不是真正的错误,不显示错误信息 // AbortError 不是真正的错误,不显示错误信息
if (e.name === 'AbortError') { if (e.name === 'AbortError') {
assistantMsg.content = assistantMsg.content || '已取消' assistantMsg.content = assistantMsg.content || '已取消'
@ -583,19 +576,19 @@ async function send(): Promise<void> {
toast('对话失败:' + e.message, 'error') toast('对话失败:' + e.message, 'error')
} }
} finally { } finally {
if (isCurrent()) {
assistantMsg.streaming = false assistantMsg.streaming = false
isSending.value = false isSending.value = false
messages.value = [...messages.value] messages.value = [...messages.value]
// 拉取推荐问题(suggest-message-list);AbortError / 异常时不拉取
if (!sseAbortController?.signal.aborted && assistantMsg.content && !assistantMsg.error) {
fetchSuggestions(chatId.value).then(items => {
if (items.length) currentSuggestions.value = items
}).catch(() => {})
// 推荐问题仍独立获取,但旧会话/旧回答的结果不得覆盖新对话。
if (assistantMsg.content && !assistantMsg.error) {
void fetchSuggestions(cid).then(items => {
if (isCurrent() && items.length) currentSuggestions.value = items
})
} }
await scrollToBottom() await scrollToBottom()
} }
}
} }
// ==================== 停止生成 ==================== // ==================== 停止生成 ====================

34
frontend/src/views/PipelineFlow.vue

@ -5,7 +5,7 @@
<div> <div>
<span style="font-size:16px;font-weight:600;">🔀 AI 执行链</span> <span style="font-size:16px;font-weight:600;">🔀 AI 执行链</span>
<p style="font-size:12px;color:var(--td-text-color-placeholder);margin:4px 0 0;"> <p style="font-size:12px;color:var(--td-text-color-placeholder);margin:4px 0 0;">
下图展示从用户请求到 AI 回复的完整处理流程,包含意图路由、RAG 检索、熔断保护和 Advisor 链。
下图展示从用户请求到 AI 回复的完整处理流程,包含 FAQ 与本地路由、RAG 检索、熔断保护和 Advisor 链。
菱形节点 = 决策分支 · 虚线框 = 独立子系统 · 虚线箭头 = 降级/异步路径。 菱形节点 = 决策分支 · 虚线框 = 独立子系统 · 虚线箭头 = 降级/异步路径。
</p> </p>
</div> </div>
@ -76,14 +76,15 @@ mermaid.initialize({
// Mermaid 流程图 DSL 定义 // Mermaid 流程图 DSL 定义
// 节点类型: [矩形]=处理步骤, {菱形}=决策分支, subgraph=子系统 // 节点类型: [矩形]=处理步骤, {菱形}=决策分支, subgraph=子系统
// %%graph-meta: { updated: "2026-09-14", basedOn: "ChatPipeline v4, RagPipeline v2, AssistantApp v2", mermaidVersion: "flowchart-v2" }
// %%graph-meta: { updated: "2026-09-14", basedOn: "ChatPipeline v5 direct retrieval, RagPipeline v3, AssistantApp sources metadata", mermaidVersion: "flowchart-v2" }
const GRAPH_DEFINITION = ` const GRAPH_DEFINITION = `
flowchart TD flowchart TD
A["<b>用户请求</b><br/>message + roleId + accountId + chatId"] A["<b>用户请求</b><br/>message + roleId + accountId + chatId"]
A --> B{"<b>Controller</b><br/>鉴权 / 角色解析 / KB 隔离判断<br/>构建 ChatContext"} A --> B{"<b>Controller</b><br/>鉴权 / 角色解析 / KB 隔离判断<br/>构建 ChatContext"}
B --> C["<b>ChatPipeline.buildRequest</b><br/>编排决策入口"]
B --> CB{"<b>AssistantApp 熔断检查</b><br/>SimpleCircuitBreaker<br/>阈值: 连续 3 次失败 / 恢复: 5 分钟"}
CB -- "正常" --> C["<b>ChatPipeline.buildRequest</b><br/>编排决策入口"]
C --> D{"enableRag ?"} C --> D{"enableRag ?"}
@ -91,36 +92,36 @@ flowchart TD
D -- "true" --> H["<b>FAQ 优先匹配</b><br/>FaqMatchEngine 完整三级匹配<br/>精确 → 关键词 → 向量语义<br/>沿用角色分类隔离"] D -- "true" --> H["<b>FAQ 优先匹配</b><br/>FaqMatchEngine 完整三级匹配<br/>精确 → 关键词 → 向量语义<br/>沿用角色分类隔离"]
H -- "命中标准答案,跳过意图分类" --> T
H -- "未命中 / 异常" --> F["<b>IntentRouter</b><br/>寒暄词快速路径: 本地精确匹配(零 LLM)<br/>未命中则 LLM 意图分类<br/>FAQ / RAG / CHITCHAT"]
H -- "命中标准答案" --> T
H -- "未命中 / 异常" --> G{"<b>本地 isChitchat</b><br/>精确寒暄词匹配<br/>无意图分类 LLM"}
F --> G{"意图分类结果"}
G -- "寒暄命中" --> I["<b>模式: 纯对话</b><br/>跳过知识库检索<br/>不注入资料块"]
G -- "CHITCHAT<br/>confidence ≧ 0.6" --> I["<b>模式: 纯对话</b><br/>跳过知识库检索<br/>不注入资料块"]
G -- "FAQ / RAG / 降级<br/>其余情况" --> J["<b>RagPipeline.retrieve</b><br/>RAG 检索流水线入口"]
G -- "其余问题" --> J["<b>RagPipeline.retrieve</b><br/>默认零 LLM 预处理"]
subgraph RAG["📚 RAG 检索流水线(当前: 纯向量检索)"] subgraph RAG["📚 RAG 检索流水线(当前: 纯向量检索)"]
J --> K{"前置 FAQ 匹配<br/>completedCleanly ?"} J --> K{"前置 FAQ 匹配<br/>completedCleanly ?"}
K -- "true: 跳过重复 FAQ" --> L["<b>2. 查询重写</b><br/>REWRITE / TRANSLATION<br/>COMPRESSION / MULTI_QUERY"]
K -- "true: 跳过重复 FAQ" --> L{"<b>2. rewriteStrategy</b><br/>默认 NONE"}
K -- "false: 异常后重试" --> KR["<b>1. FAQ 匹配重试</b><br/>FaqMatchEngine 完整三级匹配"] K -- "false: 异常后重试" --> KR["<b>1. FAQ 匹配重试</b><br/>FaqMatchEngine 完整三级匹配"]
KR -- "命中标准答案" --> T KR -- "命中标准答案" --> T
KR -- "未命中 / 异常" --> L KR -- "未命中 / 异常" --> L
L --> M["<b>3. 向量检索</b><br/>PGVector similaritySearch<br/>topK=4 + 分类过滤"]
M --> S["<b>4. 构建资料块</b><br/>拼接检索文档<br/>注入 system prompt 末尾"]
L -- "NONE 原文直检索" --> M["<b>3. 向量检索</b><br/>PGVector similaritySearch<br/>topK=4 + 分类过滤"]
L -- "显式选择" --> LR["<b>查询重写 LLM</b><br/>REWRITE / TRANSLATION / MULTI_QUERY<br/>COMPRESSION: 最近 10 条历史补全指代"]
LR --> M
M --> S["<b>4. 构建资料块</b><br/>命中/未命中日志仅此层记录一次<br/>注入 system prompt + 保留命中文档"]
end end
I --> T I --> T
S --> T S --> T
E --> T E --> T
T["<b>组装 ChatRequest</b><br/>finalMessage + finalSystemPrompt<br/>+ faqAnswer (可选)"]
T["<b>组装 ChatRequest</b><br/>finalMessage + finalSystemPrompt<br/>+ faqAnswer / 当次命中文档"]
T --> CB{"<b>🔌 AI 熔断检查</b><br/>SimpleCircuitBreaker<br/>阈值: 连续 3 次失败 / 恢复: 5 分钟"}
T -- "FAQ 标准答案" --> FAQ["<b>FAQ 直接回复</b><br/>内容安全检查 + 写入会话记忆<br/>不调用答案 LLM"]
CB -- "熔断中" --> FALLBACK["<b>返回降级提示</b><br/>「AI 服务暂时不可用<br/>请稍后重试」"] CB -- "熔断中" --> FALLBACK["<b>返回降级提示</b><br/>「AI 服务暂时不可用<br/>请稍后重试」"]
CB -- "正常" --> U["<b>AssistantApp</b><br/>chat / chatStream<br/>构建 ChatClient + MCP 工具"]
T -- "需要生成" --> U["<b>AssistantApp</b><br/>chat / chatStream<br/>构建 ChatClient + MCP 工具"]
subgraph ADVISOR["🛡️ Advisor 链(环绕 LLM 调用)"] subgraph ADVISOR["🛡️ Advisor 链(环绕 LLM 调用)"]
U --> V["<b>ContentSafetyAdvisor</b><br/>🔽 before: DFA 敏感词检测<br/>用户输入 BLOCK/MASK"] U --> V["<b>ContentSafetyAdvisor</b><br/>🔽 before: DFA 敏感词检测<br/>用户输入 BLOCK/MASK"]
@ -131,8 +132,9 @@ flowchart TD
Z --> AA["<b>ContentSafetyAdvisor</b><br/>🔼 after: AI 输出检测<br/>BLOCK/MASK 违规内容"] Z --> AA["<b>ContentSafetyAdvisor</b><br/>🔼 after: AI 输出检测<br/>BLOCK/MASK 违规内容"]
end end
AA --> AB["<b>返回 AI 回复</b><br/>SSE 流式输出<br/>+ MCP 工具调用事件"]
AA --> AB["<b>返回 AI 回复 + 当次引用</b><br/>同步 JSON 或 SSE 正文 → sources metadata → stop / DONE<br/>引用不进入正文,不二次检索"]
FAQ --> AB
FALLBACK --> AB FALLBACK --> AB
AB -. "异步按需触发" .-> SG AB -. "异步按需触发" .-> SG

164
frontend/tests/chat-protocol.test.mjs

@ -0,0 +1,164 @@
import assert from 'node:assert/strict'
import { afterEach, beforeEach, test } from 'node:test'
import { fileURLToPath } from 'node:url'
import { build } from 'esbuild'
// Exercise the production TypeScript with Vite's existing esbuild dependency, no test framework.
const root = fileURLToPath(new URL('../', import.meta.url))
async function loadModule(entry) {
const result = await build({
entryPoints: [root + entry], bundle: true, write: false, format: 'esm', platform: 'node',
alias: { '@': root + 'src' },
})
return import('data:text/javascript;base64,' + Buffer.from(result.outputFiles[0].text).toString('base64'))
}
const { readSSEStream, readSSEStreamWithEvents } = await loadModule('src/utils/sse.ts')
const { chatSync, fetchChatResult } = await loadModule('src/api/chat.ts')
const originalFetch = globalThis.fetch
const originalStorage = Object.getOwnPropertyDescriptor(globalThis, 'localStorage')
beforeEach(() => {
Object.defineProperty(globalThis, 'localStorage', {
configurable: true, value: { getItem: () => 'admin-token' },
})
})
afterEach(() => {
globalThis.fetch = originalFetch
if (originalStorage) Object.defineProperty(globalThis, 'localStorage', originalStorage)
else delete globalThis.localStorage
})
const sources = [{
documentId: '9223372036854775806', title: '报销制度', sourceName: 'policy.pdf',
chunkIndex: 2, score: 0.125, snippet: '申请应在三十天内提交。',
}]
const metadata = JSON.stringify({ object: 'chat.completion.chunk', choices: [], sources })
const answer = JSON.stringify({ choices: [{ delta: { content: '答复。' } }] })
function responseFor(text, { open = false, fragment = false, cancel } = {}) {
const bytes = new TextEncoder().encode(text)
const body = new ReadableStream({
start(controller) {
if (fragment) for (const byte of bytes) controller.enqueue(Uint8Array.of(byte))
else controller.enqueue(bytes)
if (!open) controller.close()
},
cancel,
})
return new Response(body)
}
test('fragmented CRLF metadata preserves ID/snippet and DONE completes without EOF or cancellation settlement', { timeout: 1000 }, async () => {
let cancelled = 0
const response = responseFor(
`data: ${answer}\r\n\r\ndata: ${metadata}\r\n\r\ndata: [DONE]\r\n\r\ndata: ignored\r\n\r\n`,
{ open: true, fragment: true, cancel() { cancelled++; return new Promise(() => {}) } },
)
const requests = []
globalThis.fetch = async (...args) => { requests.push(args); return response }
const events = []
await readSSEStreamWithEvents('/ai/chat/stream', {
onMessage: text => events.push(['text', text]),
onSources: value => events.push(['sources', value]),
onDone: () => events.push(['done']),
}, { Authorization: 'Bearer sdk-token' })
assert.deepEqual(events, [['text', '答复。'], ['sources', sources], ['done']])
assert.equal(requests.length, 1)
assert.equal(requests[0][1].headers.Authorization, 'Bearer sdk-token')
assert.equal(cancelled, 1)
assert.equal(response.body.locked, false)
})
test('text-only facade uses identical DONE handling and never leaks empty choices or sources', { timeout: 1000 }, async () => {
globalThis.fetch = async () => responseFor(
`data: ${metadata}\n\ndata: ${answer}\n\ndata: [DONE]\n\ndata: [DONE]\n\ndata: trailing\n\n`,
{ open: true },
)
const text = []
let completed = 0
await readSSEStream('/ai/chat/stream', value => text.push(value), () => completed++)
assert.deepEqual(text, ['答复。'])
assert.equal(completed, 1)
})
test('EOF residuals, multiline data, tool events, status and legacy text retain their semantics', async () => {
globalThis.fetch = async () => responseFor(
': heartbeat\r\nid: event-1\r\nretry: 1000\r\nevent: status\r\ndata: generating\r\n\r\n'
+ 'event: tool_call_start\r\ndata: {"tool":"lookup"}\r\n\r\n'
+ 'event: tool_call_result\r\ndata: {"tool":"lookup","result":"ok"}\r\n\r\n'
+ 'event: error\r\ndata: {"message":"tool unavailable"}\r\n\r\n'
+ 'data: first\r\ndata: indented\r\n\r\ndata:\r\n\r\nlegacy\r\ndata: 尾部',
{ fragment: true },
)
const events = []
await readSSEStreamWithEvents('/stream', {
onMessage: text => events.push(text),
onToolCallStart: value => events.push(value),
onToolCallResult: value => events.push(value),
onError: value => events.push(value),
onDone: () => events.push('done'),
})
assert.deepEqual(events, [
{ tool: 'lookup' }, { tool: 'lookup', result: 'ok' }, { message: 'tool unavailable' },
'first\n indented', '\n', 'legacy', '尾部', 'done',
])
})
test('EOF flushes an unterminated plain-text line and an unterminated DONE event', async () => {
for (const body of ['legacy tail', 'data: [DONE]']) {
globalThis.fetch = async () => responseFor(body)
const text = []
let done = 0
await readSSEStream('/stream', value => text.push(value), () => done++)
assert.deepEqual(text, body.startsWith('data:') ? [] : ['legacy tail'])
assert.equal(done, 1)
}
})
test('cleanup failure cannot replace callback errors or call onDone on an error', async () => {
const original = new Error('consumer failed')
const response = responseFor(`data: ${answer}\n\n`, {
open: true, cancel() { throw new Error('cleanup failed') },
})
globalThis.fetch = async () => response
let done = 0
await assert.rejects(readSSEStream('/stream', () => { throw original }, () => done++), error => error === original)
assert.equal(done, 0)
assert.equal(response.body.locked, false)
})
test('empty sources for ordinary answers are delivered separately from content', async () => {
globalThis.fetch = async () => responseFor('data: {"choices":[],"sources":[]}\n\ndata: [DONE]\n\n')
const values = []
await readSSEStreamWithEvents('/stream', {
onMessage: () => assert.fail('metadata is not text'),
onSources: value => values.push(value),
})
assert.deepEqual(values, [[]])
})
test('synchronous chat makes one result request and retains direct answer, sources, explicit strategy and signal', async () => {
const result = { text: '答复。', mcpEvents: [], suggestions: [], sources }
const calls = []
globalThis.fetch = async (...args) => { calls.push(args); return Response.json(result) }
const controller = new AbortController()
assert.deepEqual(await chatSync('费用?', 'conversation-1', {
enableRag: true, roleId: '9223372036854775806', rewriteStrategy: 'MULTI_QUERY', categoryIds: ['123'],
}, controller.signal), result)
assert.equal(calls.length, 1)
const url = new URL(calls[0][0], 'https://test.invalid')
assert.equal(url.pathname, '/ai/chat/result')
assert.equal(url.searchParams.get('rewriteStrategy'), 'MULTI_QUERY')
assert.equal(url.searchParams.get('roleId'), '9223372036854775806')
assert.equal(calls[0][1].signal, controller.signal)
assert.equal(calls[0][1].headers.Authorization, 'Bearer admin-token')
})
test('synchronous SDK response uses the same direct contract and surfaces server errors', async () => {
const result = { text: '普通回答', mcpEvents: [], suggestions: [], sources: [] }
globalThis.fetch = async () => Response.json(result)
assert.deepEqual(await fetchChatResult('https://sdk.invalid/ai/chat/result', {}), result)
globalThis.fetch = async () => new Response('角色无访问权限', { status: 403 })
await assert.rejects(fetchChatResult('/ai/chat/result', {}), /角色无访问权限/)
await assert.rejects(readSSEStream('/ai/chat/stream', () => {}), /角色无访问权限/)
})
Loading…
Cancel
Save