From 1355e3460ecca9b654ef985a4f5ff45685fd4e8a Mon Sep 17 00:00:00 2001 From: InfinityPacer <160988576+InfinityPacer@users.noreply.github.com> Date: Wed, 12 Aug 2026 14:35:44 +0800 Subject: [PATCH] feat(agent): render protected stream deliveries (#669) --- src/components/agent/AgentAssistantPanel.vue | 214 ++++- .../__tests__/AgentAssistantPanel.spec.ts | 742 ++++++++++++++++++ 2 files changed, 935 insertions(+), 21 deletions(-) diff --git a/src/components/agent/AgentAssistantPanel.vue b/src/components/agent/AgentAssistantPanel.vue index b15bb538..7f056c03 100644 --- a/src/components/agent/AgentAssistantPanel.vue +++ b/src/components/agent/AgentAssistantPanel.vue @@ -111,6 +111,17 @@ interface AgentStreamEvent { session_id?: string } +const AGENT_STREAM_EVENT_TYPES = new Set([ + 'start', + 'delta', + 'tool', + 'attachment', + 'choice', + 'message_update', + 'done', + 'error', +]) + interface AgentPendingAttachment { id: string file: File @@ -155,6 +166,18 @@ interface AgentStreamReadResult { receivedTerminalEvent: boolean } +interface AgentProtectedDelivery { + /** 一次性投递标识,只用于当前可见范围去重。 */ + id: string + /** 仅驻留组件内存并按纯文本渲染的受保护内容。 */ + content: string +} + +interface ParsedSseBlock { + eventName: string + data: string +} + interface AgentSlashCommand { command: string description: string @@ -211,6 +234,7 @@ const STREAM_STATE_PERSIST_DELAY = 1000 const inputText = ref('') const messages = ref([]) +const protectedDeliveries = ref([]) const historySessions = ref([]) const sessionId = ref('') const sending = ref(false) @@ -250,6 +274,8 @@ let userAbortRequested = false let streamRecoveryAbortRequested = false let streamRecoveryTimer: number | null = null let activeStreamStartedAt = 0 +let protectedDeliveryGeneration = 0 +const protectedDeliveryIds = new Set() // 汇总实时请求与后台恢复状态,保证恢复期间仍展示处理中并锁定会话操作。 const isBusy = computed(() => sending.value || Boolean(pendingStreamRecovery.value)) @@ -298,7 +324,7 @@ const recordingTimeText = computed(() => { const drawerWidth = computed(() => (display.mdAndDown.value || fullscreen.value ? '100vw' : '30rem')) // 仅桌面宽屏展示全屏开关,窄屏已默认占满视口。 const canToggleFullscreen = computed(() => !display.mdAndDown.value) -const hasMessages = computed(() => messages.value.length > 0) +const hasConversationContent = computed(() => messages.value.length > 0 || protectedDeliveries.value.length > 0) const hasHistorySessions = computed(() => historySessions.value.length > 0) const currentUserName = computed(() => userStore.getUserName || t('common.user')) const isOpen = computed({ @@ -320,6 +346,18 @@ function createSessionId() { return `web-${Date.now()}-${Math.random().toString(16).slice(2)}` } +// 受保护内容只属于当前可见范围;递增代次可阻止旧流在清理后重新写入。 +function invalidateProtectedDeliveries() { + protectedDeliveryGeneration += 1 + protectedDeliveries.value = [] + protectedDeliveryIds.clear() +} + +// 识别前端需要隔离展示和持久化的确认控制文本,授权判断仍由后端完成。 +function isReservedConfirmationControl(value: string) { + return /^(确认|取消) ([A-Za-z0-9]{4}-[A-Za-z0-9]{4})$/.test(value) +} + // 将未知字段安全转换为可展示文本。 function stringifyChoiceField(value: unknown) { if (typeof value === 'string') return value.trim() @@ -1506,20 +1544,71 @@ function queueStreamEvent(event: AgentStreamEvent, assistantMessage: AgentChatMe schedulePendingStreamDeltaFlush() } -// 解析一个 SSE 数据块。 -function parseSseBlock(block: string) { - const data = block - .split('\n') - .filter(line => line.startsWith('data:')) - .map(line => line.slice(5).trimStart()) - .join('\n') +// 拆分 SSE event name 与 data,确保受保护 frame 在普通事件解析前完成分流。 +function splitSseBlock(block: string) { + let eventName = '' + const dataLines: string[] = [] - if (!data) return null - return JSON.parse(data) as AgentStreamEvent + for (const rawLine of block.split(/\r?\n/)) { + if (!rawLine || rawLine.startsWith(':')) continue + + const separatorIndex = rawLine.indexOf(':') + const field = separatorIndex >= 0 ? rawLine.slice(0, separatorIndex) : rawLine + let value = separatorIndex >= 0 ? rawLine.slice(separatorIndex + 1) : '' + if (value.startsWith(' ')) value = value.slice(1) + + if (field === 'event') eventName = value + if (field === 'data') dataLines.push(value) + } + + if (!dataLines.length) return null + + const data = dataLines.join('\n') + if (!data.trim()) return null + return { eventName, data } satisfies ParsedSseBlock +} + +function parseProtectedTransportFrame(data: string): AgentProtectedDelivery | null { + try { + const frame = JSON.parse(data) as Record + const schemaVersion = typeof frame.schema_version === 'string' ? frame.schema_version : '' + const versionMatch = /^(\d+)\.(\d+)$/.exec(schemaVersion) + const deliveryId = frame.delivery_id + + if ( + !versionMatch || + versionMatch[1] !== '1' || + typeof deliveryId !== 'string' || + !deliveryId.trim() || + frame.content_type !== 'text/plain' || + typeof frame.content !== 'string' + ) { + return null + } + + return { id: deliveryId, content: frame.content } + } catch { + return null + } +} + +function consumeProtectedTransportFrame(data: string, streamGeneration: number) { + if (!isOpen.value || streamGeneration !== protectedDeliveryGeneration) return + + const delivery = parseProtectedTransportFrame(data) + if (!delivery || protectedDeliveryIds.has(delivery.id)) return + + protectedDeliveryIds.add(delivery.id) + protectedDeliveries.value.push(delivery) + nextTick(() => scheduleMessageScrollerUpdate({ toBottom: messageScrollerShouldFollow })) } // 读取并应用智能助手 SSE 响应流。 -async function readAgentStream(response: Response, assistantMessage: AgentChatMessage): Promise { +async function readAgentStream( + response: Response, + assistantMessage: AgentChatMessage, + streamGeneration: number, +): Promise { if (!response.body) { throw new Error(t('agentAssistant.noStream')) } @@ -1530,8 +1619,36 @@ async function readAgentStream(response: Response, assistantMessage: AgentChatMe let receivedTerminalEvent = false // 应用事件并记录服务端是否明确结束本轮流,区分正常完成与无异常的后台断流。 - const consumeEvent = (event: AgentStreamEvent | null) => { - if (!event) return + const consumeBlock = (block: string) => { + const parsedBlock = splitSseBlock(block) + if (!parsedBlock) return + + if (parsedBlock.eventName === 'interaction-protected') { + consumeProtectedTransportFrame(parsedBlock.data, streamGeneration) + return + } + if ( + parsedBlock.eventName && + parsedBlock.eventName !== 'message' && + !AGENT_STREAM_EVENT_TYPES.has(parsedBlock.eventName as AgentStreamEvent['type']) + ) { + return + } + + const parsedEvent = JSON.parse(parsedBlock.data) as unknown + if (!parsedEvent || typeof parsedEvent !== 'object' || Array.isArray(parsedEvent)) { + return + } + + const eventRecord = parsedEvent as Record + if ( + typeof eventRecord.type !== 'string' || + !AGENT_STREAM_EVENT_TYPES.has(eventRecord.type as AgentStreamEvent['type']) + ) + return + + const event = eventRecord as unknown as AgentStreamEvent + if (parsedBlock.eventName && parsedBlock.eventName !== 'message' && parsedBlock.eventName !== event.type) return queueStreamEvent(event, assistantMessage) if (event.type === 'done' || event.type === 'error') receivedTerminalEvent = true @@ -1547,13 +1664,13 @@ async function readAgentStream(response: Response, assistantMessage: AgentChatMe buffer = blocks.pop() || '' for (const block of blocks) { - consumeEvent(parseSseBlock(block)) + consumeBlock(block) } } buffer += decoder.decode() if (buffer.trim()) { - consumeEvent(parseSseBlock(buffer)) + consumeBlock(buffer) } } finally { flushPendingStreamDelta() @@ -1753,6 +1870,7 @@ async function streamAgentMessage( userAbortRequested = false streamRecoveryAbortRequested = false const streamStartedAt = Date.now() + const streamProtectedDeliveryGeneration = protectedDeliveryGeneration activeStreamStartedAt = streamStartedAt let shouldFollowBottomAfterStream = true let shouldSaveClientSnapshot = true @@ -1762,6 +1880,7 @@ async function streamAgentMessage( method: 'POST', headers: buildAgentRequestHeaders({ 'Content-Type': 'application/json', + 'X-MoviePilot-Agent-Interaction': '1', }), body: JSON.stringify({ text: content, @@ -1783,10 +1902,11 @@ async function streamAgentMessage( throw new Error(await resolveAgentResponseErrorMessage(response)) } - const streamResult = await readAgentStream(response, assistantMessage) + const streamResult = await readAgentStream(response, assistantMessage, streamProtectedDeliveryGeneration) shouldFollowBottomAfterStream = isMessageScrollerNearBottom() if (!streamResult.receivedTerminalEvent) { shouldSaveClientSnapshot = false + invalidateProtectedDeliveries() beginStreamRecovery(sessionId.value, streamStartedAt) refreshMessageList() if (document.visibilityState === 'visible') scheduleStreamRecovery(0) @@ -1817,6 +1937,7 @@ async function streamAgentMessage( if (isRecoverableStreamDisconnect(error)) { shouldSaveClientSnapshot = false + invalidateProtectedDeliveries() beginStreamRecovery(sessionId.value, streamStartedAt) assistantMessage.status = 'streaming' refreshMessageList() @@ -1849,7 +1970,8 @@ async function streamAgentMessage( // 发送输入框中的文本和附件。 async function sendMessage() { - const text = inputText.value.trim() + const rawText = inputText.value + const text = rawText.trim() const attachments = [...pendingAttachments.value] if ((!text && !attachments.length) || isBusy.value) return @@ -1861,7 +1983,10 @@ async function sendMessage() { try { const prepared = await prepareAgentAttachments(attachments) - await streamAgentMessage(text, prepared.images, prepared.files, prepared.audioRefs, prepared.userAttachments) + const reservedControl = attachments.length === 0 && isReservedConfirmationControl(rawText) + await streamAgentMessage(text, prepared.images, prepared.files, prepared.audioRefs, prepared.userAttachments, { + echoUser: !reservedControl, + }) } catch (error: any) { // 附件准备失败同样落到对话消息里,底部提示条只保留给没有消息承载的本地错误。 addMessage('assistant', error?.message || t('agentAssistant.uploadFailed'), 'error') @@ -2109,6 +2234,7 @@ function stopGeneration() { // 开始新的空白会话。 function startNewSession() { + invalidateProtectedDeliveries() stopGeneration() sessionId.value = createSessionId() messages.value = [] @@ -2127,6 +2253,7 @@ async function loadHistorySession(targetSessionId: string) { if (!historySession) return try { + invalidateProtectedDeliveries() stopGeneration() if (!historySession.messages.length) { historySession = await loadServerHistorySession(targetSessionId) @@ -2180,6 +2307,7 @@ function formatHistoryTime(timestamp: number) { // 关闭智能助手面板。 function closeDrawer() { + invalidateProtectedDeliveries() isOpen.value = false } @@ -2262,7 +2390,12 @@ watch(drawerWidth, () => { }) watch(isOpen, open => { - if (open) scrollToBottom() + if (open) { + scrollToBottom() + return + } + + invalidateProtectedDeliveries() }) watch(isBusy, value => emit('thinking-change', value), { immediate: true }) @@ -2285,6 +2418,7 @@ onScopeDispose(clearMessageScrollFrame) onScopeDispose(clearStreamPersistTimer) onScopeDispose(clearPendingStreamDelta) onScopeDispose(clearStreamRecoveryTimer) +onScopeDispose(invalidateProtectedDeliveries) onScopeDispose(() => { if (typeof window === 'undefined') return @@ -2437,11 +2571,11 @@ onScopeDispose(() => {
-
+
@@ -2600,6 +2734,18 @@ onScopeDispose(() => {
+ +
+
+
+
@@ -3155,6 +3301,32 @@ onScopeDispose(() => { align-items: flex-start; } +.agent-assistant-protected-deliveries { + display: grid; + gap: 0.5rem; + inline-size: min(100%, 34rem); + margin-block-end: 1rem; +} + +.agent-assistant-protected-delivery { + display: grid; + align-items: start; + border: 1px solid rgba(var(--v-theme-warning), 0.42); + border-radius: var(--app-control-radius); + background: rgba(var(--v-theme-warning), 0.08); + color: rgba(var(--v-theme-on-surface), 0.9); + column-gap: 0.5rem; + font-family: ui-monospace, SFMono-Regular, Menlo, Monaco, Consolas, monospace; + font-size: 0.82rem; + grid-template-columns: auto minmax(0, 1fr); + line-height: 1.55; + min-inline-size: 0; + overflow-wrap: anywhere; + padding-block: 0.65rem; + padding-inline: 0.75rem; + white-space: pre-wrap; +} + .agent-assistant-message__meta { display: inline-flex; align-items: center; diff --git a/src/components/agent/__tests__/AgentAssistantPanel.spec.ts b/src/components/agent/__tests__/AgentAssistantPanel.spec.ts index 2400f7db..da769339 100644 --- a/src/components/agent/__tests__/AgentAssistantPanel.spec.ts +++ b/src/components/agent/__tests__/AgentAssistantPanel.spec.ts @@ -1,4 +1,5 @@ import { flushPromises, shallowMount } from '@vue/test-utils' +import { defineComponent, h } from 'vue' import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' import AgentAssistantPanel from '@/components/agent/AgentAssistantPanel.vue' @@ -29,6 +30,11 @@ interface MockServerSession { messages: Array> } +interface SyntheticSseFrame { + eventName?: string + data: unknown +} + const agentMarkdownContentStub = { props: ['content', 'variant'], template: @@ -43,9 +49,98 @@ function createAgentResponse(data: unknown) { } as unknown as Response } +function legacySseFrame(data: Record): SyntheticSseFrame { + return { data } +} + +function protectedSseFrame(data: Record): SyntheticSseFrame { + return { eventName: 'interaction-protected', data } +} + +function createAgentStreamResponse(frames: SyntheticSseFrame[]) { + const body = frames + .map(frame => `${frame.eventName ? `event: ${frame.eventName}\n` : ''}data: ${JSON.stringify(frame.data)}\n\n`) + .join('') + + return new Response(body, { + status: 200, + headers: { 'Content-Type': 'text/event-stream' }, + }) +} + +function createRawAgentStreamResponse(body: string) { + return new Response(body, { + status: 200, + headers: { 'Content-Type': 'text/event-stream' }, + }) +} + +function createControllableAgentStream() { + const encoder = new TextEncoder() + let streamController: ReadableStreamDefaultController + const response = new Response( + new ReadableStream({ + start(controller) { + streamController = controller + }, + }), + { status: 200, headers: { 'Content-Type': 'text/event-stream' } }, + ) + + return { + response, + emit(frame: SyntheticSseFrame) { + const block = `${frame.eventName ? `event: ${frame.eventName}\n` : ''}data: ${JSON.stringify(frame.data)}\n\n` + streamController.enqueue(encoder.encode(block)) + }, + close() { + streamController.close() + }, + fail(error: Error) { + streamController.error(error) + }, + } +} + +const slotContainerStub = { template: '
' } +const menuStub = defineComponent({ + setup(_props, { slots }) { + return () => h('div', [slots.activator?.({ props: {} }), slots.default?.()]) + }, +}) +const virtualScrollStub = defineComponent({ + props: { items: { type: Array, default: () => [] } }, + setup(props, { slots }) { + return () => + h( + 'div', + props.items.map(item => slots.default?.({ item, itemRef: () => undefined })), + ) + }, +}) + +function mountPanel() { + return shallowMount(AgentAssistantPanel, { + props: { modelValue: true }, + global: { + stubs: { + AgentMarkdownContent: agentMarkdownContentStub, + IconBtn: { template: '' }, + PerfectScrollbar: { template: '
' }, + VCard: slotContainerStub, + VInfiniteScroll: slotContainerStub, + VIcon: true, + VMenu: menuStub, + VVirtualScroll: virtualScrollStub, + }, + }, + }) +} + describe('AgentAssistantPanel stream recovery', () => { beforeEach(() => { vi.useFakeTimers() + localStorage.clear() }) afterEach(() => { @@ -349,6 +444,653 @@ describe('AgentAssistantPanel stream recovery', () => { wrapper.unmount() }) + it('renders protected delivery only as transient literal text and advertises the stream capability', async () => { + const protectedMarker = 'MP-PROTECTED-MARKER **not bold** ' + const fetchMock = vi.fn(async (input: RequestInfo | URL, init?: RequestInit) => { + const url = String(input) + if (url.endsWith('/message/agent/stream') && init?.method === 'POST') { + return createAgentStreamResponse([ + legacySseFrame({ type: 'start', session_id: 'web-agent:protected' }), + legacySseFrame({ type: 'delta', content: '普通回复' }), + protectedSseFrame({ + schema_version: '1.0', + delivery_id: 'delivery-1', + content_type: 'text/plain', + content: protectedMarker, + }), + legacySseFrame({ type: 'tool', message: '(查询了 1 次数据)' }), + legacySseFrame({ type: 'done' }), + ]) + } + + return createAgentResponse([]) + }) + vi.stubGlobal('fetch', fetchMock) + + const wrapper = mountPanel() + await wrapper.find('textarea').setValue('读取受保护结果') + await wrapper.find('textarea').trigger('keydown', { key: 'Enter' }) + await flushPromises() + + const protectedNode = wrapper.get('[data-protected-delivery-id="delivery-1"]') + expect(protectedNode.text()).toBe(protectedMarker) + expect(protectedNode.find('script').exists()).toBe(false) + const assistantBubble = wrapper.get('.agent-assistant-message--assistant .agent-assistant-message__bubble') + expect(assistantBubble.text()).not.toContain(protectedMarker) + expect(assistantBubble.text()).toContain('普通回复') + + const streamCall = fetchMock.mock.calls.find(([input]) => String(input).endsWith('/message/agent/stream')) + const streamHeaders = new Headers(streamCall?.[1]?.headers) + expect(streamHeaders.get('X-MoviePilot-Agent-Interaction')).toBe('1') + + const displayCalls = fetchMock.mock.calls.filter(([input]) => String(input).includes('/display')) + expect(displayCalls.length).toBeGreaterThan(0) + displayCalls.forEach(([, init]) => { + expect(new Headers(init?.headers).has('X-MoviePilot-Agent-Interaction')).toBe(false) + expect(String(init?.body)).not.toContain(protectedMarker) + }) + expect(localStorage.getItem('moviepilot-agent-assistant-state')).not.toContain(protectedMarker) + expect(localStorage.getItem('moviepilot-agent-assistant-history')).not.toContain(protectedMarker) + expect(JSON.stringify(wrapper.emitted('assistant-preview') || [])).not.toContain(protectedMarker) + + wrapper.unmount() + }) + + it('ignores empty data heartbeats without interrupting ordinary stream events', async () => { + vi.stubGlobal( + 'fetch', + vi.fn(async (input: RequestInfo | URL, init?: RequestInit) => { + if (String(input).endsWith('/message/agent/stream') && init?.method === 'POST') { + return createRawAgentStreamResponse( + [ + 'data:', + '', + 'data: ', + '', + 'data', + '', + 'data: {"type":"delta","content":"心跳后的普通回复"}', + '', + 'data: {"type":"done"}', + '', + ].join('\r\n'), + ) + } + + return createAgentResponse([]) + }), + ) + + const wrapper = mountPanel() + await wrapper.find('textarea').setValue('测试空数据心跳') + await wrapper.find('textarea').trigger('keydown', { key: 'Enter' }) + await flushPromises() + + expect(wrapper.get('.agent-assistant-message--assistant').text()).toContain('心跳后的普通回复') + expect(wrapper.find('.agent-assistant-message--error').exists()).toBe(false) + + wrapper.unmount() + }) + + it('ignores unknown and mismatched named protocols before the ordinary event queue', async () => { + const rejectedMarker = 'MP-REJECTED-NAMED-EVENT' + const fetchMock = vi.fn(async (input: RequestInfo | URL, init?: RequestInit) => { + if (String(input).endsWith('/message/agent/stream') && init?.method === 'POST') { + return createRawAgentStreamResponse( + [ + 'event: ping', + 'data: keepalive', + '', + 'event: interaction', + `data: {"type":"interaction","content":"${rejectedMarker}"}`, + '', + 'event: delta', + `data: {"type":"tool","message":"${rejectedMarker}"}`, + '', + 'event: message', + `data: {"type":"future_event","content":"${rejectedMarker}"}`, + '', + 'event: delta', + 'data: {"type":"delta","content":"未知协议后的普通回复"}', + '', + 'event: done', + 'data: {"type":"done"}', + '', + ].join('\n'), + ) + } + + return createAgentResponse([]) + }) + vi.stubGlobal('fetch', fetchMock) + + const wrapper = mountPanel() + await wrapper.find('textarea').setValue('测试未知具名协议') + await wrapper.find('textarea').trigger('keydown', { key: 'Enter' }) + await flushPromises() + + expect(wrapper.get('.agent-assistant-message--assistant').text()).toContain('未知协议后的普通回复') + expect(wrapper.text()).not.toContain(rejectedMarker) + expect(wrapper.find('.agent-assistant-message--error').exists()).toBe(false) + expect(localStorage.getItem('moviepilot-agent-assistant-state')).not.toContain(rejectedMarker) + expect(localStorage.getItem('moviepilot-agent-assistant-history')).not.toContain(rejectedMarker) + expect(JSON.stringify(wrapper.emitted('assistant-preview') || [])).not.toContain(rejectedMarker) + fetchMock.mock.calls + .filter(([input]) => String(input).includes('/display')) + .forEach(([, init]) => expect(String(init?.body)).not.toContain(rejectedMarker)) + + wrapper.unmount() + }) + + it('ignores null and non-object JSON frames without interrupting ordinary stream events', async () => { + vi.stubGlobal( + 'fetch', + vi.fn(async (input: RequestInfo | URL, init?: RequestInit) => { + if (String(input).endsWith('/message/agent/stream') && init?.method === 'POST') { + return createRawAgentStreamResponse( + [ + 'data: null', + '', + 'data: 42', + '', + 'data: "heartbeat"', + '', + 'data: []', + '', + 'data: {"type":"delta","content":"空值帧后的普通回复"}', + '', + 'data: {"type":"done"}', + '', + ].join('\n'), + ) + } + + return createAgentResponse([]) + }), + ) + + const wrapper = mountPanel() + await wrapper.find('textarea').setValue('测试 JSON 空值帧') + await wrapper.find('textarea').trigger('keydown', { key: 'Enter' }) + await flushPromises() + + expect(wrapper.get('.agent-assistant-message--assistant').text()).toContain('空值帧后的普通回复') + expect(wrapper.find('.agent-assistant-message--error').exists()).toBe(false) + + wrapper.unmount() + }) + + it('keeps consuming ordinary JSON frames with explicit SSE event names', async () => { + vi.stubGlobal( + 'fetch', + vi.fn(async (input: RequestInfo | URL, init?: RequestInit) => { + if (String(input).endsWith('/message/agent/stream') && init?.method === 'POST') { + return createAgentStreamResponse([ + { eventName: 'message', data: { type: 'start', session_id: 'web-agent:named-events' } }, + { eventName: 'delta', data: { type: 'delta', content: '具名普通回复' } }, + { eventName: 'done', data: { type: 'done' } }, + ]) + } + + return createAgentResponse([]) + }), + ) + + const wrapper = mountPanel() + await wrapper.find('textarea').setValue('测试普通具名事件') + await wrapper.find('textarea').trigger('keydown', { key: 'Enter' }) + await flushPromises() + + expect(wrapper.get('.agent-assistant-message--assistant').text()).toContain('具名普通回复') + const persistedState = JSON.parse(localStorage.getItem('moviepilot-agent-assistant-state') || '{}') + expect(persistedState.streamRecovery).toBeNull() + + wrapper.unmount() + }) + + it('ignores malformed and duplicate protected delivery events without creating a normal message', async () => { + const acceptedMarker = 'MP-ACCEPTED-PROTECTED-MARKER' + const rejectedMarker = 'MP-REJECTED-PROTECTED-MARKER' + vi.stubGlobal( + 'fetch', + vi.fn(async (input: RequestInfo | URL, init?: RequestInit) => { + if (String(input).endsWith('/message/agent/stream') && init?.method === 'POST') { + return createAgentStreamResponse([ + legacySseFrame({ type: 'start', session_id: 'web-agent:protected-validation' }), + protectedSseFrame({ + schema_version: '2.0', + delivery_id: 'wrong-major', + content_type: 'text/plain', + content: rejectedMarker, + }), + protectedSseFrame({ + schema_version: '1.0', + delivery_id: '', + content_type: 'text/plain', + content: rejectedMarker, + }), + protectedSseFrame({ + schema_version: '1.0', + delivery_id: 'wrong-content-type', + content_type: 'text/html', + content: rejectedMarker, + }), + protectedSseFrame({ + schema_version: '1.0', + delivery_id: 'bad-content', + content_type: 'text/plain', + content: 42, + }), + { eventName: 'interaction', data: { content: rejectedMarker } }, + legacySseFrame({ type: 'protected_delivery', delivery_id: 'legacy-shape', content: rejectedMarker }), + protectedSseFrame({ + schema_version: '1.7', + delivery_id: 'delivery-valid', + content_type: 'text/plain', + content: acceptedMarker, + future_field: true, + }), + protectedSseFrame({ + schema_version: '1.0', + delivery_id: 'delivery-valid', + content_type: 'text/plain', + content: rejectedMarker, + }), + legacySseFrame({ type: 'done' }), + ]) + } + + return createAgentResponse([]) + }), + ) + + const wrapper = mountPanel() + await wrapper.find('textarea').setValue('读取一次性结果') + await wrapper.find('textarea').trigger('keydown', { key: 'Enter' }) + await flushPromises() + + const protectedNodes = wrapper.findAll('.agent-assistant-protected-delivery') + expect(protectedNodes).toHaveLength(1) + expect(protectedNodes[0].text()).toBe(acceptedMarker) + expect(wrapper.text()).not.toContain(rejectedMarker) + expect(wrapper.find('.agent-assistant-message--assistant').exists()).toBe(false) + + wrapper.unmount() + }) + + it('clears protected delivery on close and rejects a late frame from the old stream after reopen', async () => { + const protectedMarker = 'MP-LATE-PROTECTED-MARKER' + const stream = createControllableAgentStream() + vi.stubGlobal( + 'fetch', + vi.fn(async (input: RequestInfo | URL, init?: RequestInit) => { + if (String(input).endsWith('/message/agent/stream') && init?.method === 'POST') return stream.response + return createAgentResponse([]) + }), + ) + + const wrapper = mountPanel() + await wrapper.find('textarea').setValue('读取晚到结果') + await wrapper.find('textarea').trigger('keydown', { key: 'Enter' }) + await flushPromises() + + stream.emit(legacySseFrame({ type: 'start', session_id: 'web-agent:late-protected' })) + stream.emit( + protectedSseFrame({ + schema_version: '1.0', + delivery_id: 'delivery-before-close', + content_type: 'text/plain', + content: protectedMarker, + }), + ) + await flushPromises() + expect(wrapper.text()).toContain(protectedMarker) + + await wrapper.setProps({ modelValue: false }) + expect(wrapper.text()).not.toContain(protectedMarker) + + stream.emit( + protectedSseFrame({ + schema_version: '1.0', + delivery_id: 'delivery-after-close', + content_type: 'text/plain', + content: protectedMarker, + }), + ) + await wrapper.setProps({ modelValue: true }) + await flushPromises() + expect(wrapper.text()).not.toContain(protectedMarker) + + stream.close() + await flushPromises() + expect(JSON.stringify(localStorage)).not.toContain(protectedMarker) + wrapper.unmount() + }) + + it('clears protected delivery when the stream ends without a terminal event and enters recovery', async () => { + const protectedMarker = 'MP-EOF-PROTECTED-MARKER' + const stream = createControllableAgentStream() + const fetchMock = vi.fn(async (input: RequestInfo | URL, init?: RequestInit) => { + if (String(input).endsWith('/message/agent/stream') && init?.method === 'POST') return stream.response + return createAgentResponse([]) + }) + vi.stubGlobal('fetch', fetchMock) + + const wrapper = mountPanel() + await wrapper.find('textarea').setValue('读取后模拟断流') + await wrapper.find('textarea').trigger('keydown', { key: 'Enter' }) + await flushPromises() + + stream.emit( + protectedSseFrame({ + schema_version: '1.0', + delivery_id: 'delivery-before-eof', + content_type: 'text/plain', + content: protectedMarker, + }), + ) + await flushPromises() + expect(wrapper.text()).toContain(protectedMarker) + + stream.close() + await flushPromises() + + const persistedState = JSON.parse(localStorage.getItem('moviepilot-agent-assistant-state') || '{}') + expect(wrapper.text()).not.toContain(protectedMarker) + expect(persistedState.streamRecovery).toMatchObject({ attempts: 0 }) + expect(JSON.stringify(persistedState)).not.toContain(protectedMarker) + expect(localStorage.getItem('moviepilot-agent-assistant-history')).not.toContain(protectedMarker) + expect(JSON.stringify(wrapper.emitted('assistant-preview') || [])).not.toContain(protectedMarker) + fetchMock.mock.calls + .filter(([input]) => String(input).includes('/display')) + .forEach(([, init]) => expect(String(init?.body)).not.toContain(protectedMarker)) + + wrapper.unmount() + }) + + it('clears protected delivery on a recoverable stream error while preserving ordinary recovery state', async () => { + const protectedMarker = 'MP-NETWORK-PROTECTED-MARKER' + const stream = createControllableAgentStream() + const fetchMock = vi.fn(async (input: RequestInfo | URL, init?: RequestInit) => { + if (String(input).endsWith('/message/agent/stream') && init?.method === 'POST') return stream.response + return createAgentResponse([]) + }) + vi.stubGlobal('fetch', fetchMock) + + const wrapper = mountPanel() + await wrapper.find('textarea').setValue('读取后模拟网络断开') + await wrapper.find('textarea').trigger('keydown', { key: 'Enter' }) + await flushPromises() + + stream.emit( + protectedSseFrame({ + schema_version: '1.0', + delivery_id: 'delivery-before-network-error', + content_type: 'text/plain', + content: protectedMarker, + }), + ) + await flushPromises() + expect(wrapper.text()).toContain(protectedMarker) + + stream.fail(new TypeError('Failed to fetch')) + await flushPromises() + + const persistedState = JSON.parse(localStorage.getItem('moviepilot-agent-assistant-state') || '{}') + expect(wrapper.text()).not.toContain(protectedMarker) + expect(persistedState.streamRecovery).toMatchObject({ attempts: 0 }) + expect(persistedState.messages.at(-1)).toMatchObject({ role: 'assistant', status: 'streaming' }) + expect(JSON.stringify(persistedState)).not.toContain(protectedMarker) + expect(localStorage.getItem('moviepilot-agent-assistant-history')).not.toContain(protectedMarker) + expect(JSON.stringify(wrapper.emitted('assistant-preview') || [])).not.toContain(protectedMarker) + fetchMock.mock.calls + .filter(([input]) => String(input).includes('/display')) + .forEach(([, init]) => expect(String(init?.body)).not.toContain(protectedMarker)) + + wrapper.unmount() + }) + + it('clears protected delivery for a new session and never restores it after remount', async () => { + const protectedMarker = 'MP-SESSION-PROTECTED-MARKER' + const fetchMock = vi.fn(async (input: RequestInfo | URL, init?: RequestInit) => { + if (String(input).endsWith('/message/agent/stream') && init?.method === 'POST') { + return createAgentStreamResponse([ + protectedSseFrame({ + schema_version: '1.0', + delivery_id: 'delivery-session', + content_type: 'text/plain', + content: protectedMarker, + }), + legacySseFrame({ type: 'done' }), + ]) + } + return createAgentResponse([]) + }) + vi.stubGlobal('fetch', fetchMock) + + const wrapper = mountPanel() + await wrapper.find('textarea').setValue('读取会话结果') + await wrapper.find('textarea').trigger('keydown', { key: 'Enter' }) + await flushPromises() + expect(wrapper.text()).toContain(protectedMarker) + + await wrapper.get('[title="agentAssistant.newChat"]').trigger('click') + expect(wrapper.text()).not.toContain(protectedMarker) + + wrapper.unmount() + const remounted = mountPanel() + await flushPromises() + expect(remounted.text()).not.toContain(protectedMarker) + expect(localStorage.getItem('moviepilot-agent-assistant-state')).not.toContain(protectedMarker) + expect(localStorage.getItem('moviepilot-agent-assistant-history')).not.toContain(protectedMarker) + fetchMock.mock.calls + .filter(([input]) => String(input).includes('/display')) + .forEach(([, init]) => expect(String(init?.body)).not.toContain(protectedMarker)) + remounted.unmount() + }) + + it('clears protected delivery before loading another history session', async () => { + const protectedMarker = 'MP-HISTORY-PROTECTED-MARKER' + const historySessionId = 'web-agent:history-target' + const historyMessage = { + id: 'history-user', + role: 'user', + content: '历史会话内容', + createdAt: Date.now() - 1000, + status: 'done', + attachments: [], + choices: [], + tools: [], + } + const serverSession = { + session_id: historySessionId, + client_session_id: 'history-target', + title: '历史会话', + updated_at: new Date().toISOString(), + is_processing: false, + messages: [historyMessage], + } + const fetchMock = vi.fn(async (input: RequestInfo | URL, init?: RequestInit) => { + const url = String(input) + if (url.endsWith('/message/agent/stream') && init?.method === 'POST') { + return createAgentStreamResponse([ + protectedSseFrame({ + schema_version: '1.0', + delivery_id: 'delivery-history', + content_type: 'text/plain', + content: protectedMarker, + }), + legacySseFrame({ type: 'done' }), + ]) + } + if (url.includes(`/sessions/${encodeURIComponent(historySessionId)}`)) return createAgentResponse(serverSession) + if (url.includes('/message/agent/sessions?')) return createAgentResponse([serverSession]) + return createAgentResponse([]) + }) + vi.stubGlobal('fetch', fetchMock) + + const wrapper = mountPanel() + await flushPromises() + await wrapper.find('textarea').setValue('读取后切换历史') + await wrapper.find('textarea').trigger('keydown', { key: 'Enter' }) + await flushPromises() + expect(wrapper.text()).toContain(protectedMarker) + + await wrapper.get('.agent-assistant-history-item').trigger('click') + await flushPromises() + expect(wrapper.text()).not.toContain(protectedMarker) + expect(wrapper.text()).toContain('历史会话内容') + wrapper.unmount() + }) + + it('keeps exact confirmation controls out of Web history while preserving the old-backend response', async () => { + const controlText = '确认 aB12-cD34' + const backendFeedback = '确认无效或已过期' + const fetchMock = vi.fn(async (input: RequestInfo | URL, init?: RequestInit) => { + if (String(input).endsWith('/message/agent/stream') && init?.method === 'POST') { + return createAgentStreamResponse([ + legacySseFrame({ type: 'start', session_id: 'web-agent:control' }), + legacySseFrame({ type: 'delta', content: backendFeedback }), + legacySseFrame({ type: 'done' }), + ]) + } + + return createAgentResponse([]) + }) + vi.stubGlobal('fetch', fetchMock) + + const wrapper = mountPanel() + await wrapper.find('textarea').setValue(controlText) + await wrapper.find('textarea').trigger('keydown', { key: 'Enter' }) + await flushPromises() + + const streamCall = fetchMock.mock.calls.find(([input]) => String(input).endsWith('/message/agent/stream')) + const streamBody = JSON.parse(String(streamCall?.[1]?.body || '{}')) + expect(streamBody).toMatchObject({ text: controlText, display_text: controlText, echo_user: false }) + expect(wrapper.text()).toContain(backendFeedback) + expect(wrapper.find('.agent-assistant-message--user').exists()).toBe(false) + expect(localStorage.getItem('moviepilot-agent-assistant-state')).not.toContain(controlText) + expect(localStorage.getItem('moviepilot-agent-assistant-history')).not.toContain(controlText) + fetchMock.mock.calls + .filter(([input]) => String(input).includes('/display')) + .forEach(([, init]) => expect(String(init?.body)).not.toContain(controlText)) + + wrapper.unmount() + }) + + it.each([ + ['leading space', ' 确认 AB12-CD34'], + ['trailing space', '确认 AB12-CD34 '], + ['repeated separator', '确认 AB12-CD34'], + ['trailing newline', '确认 AB12-CD34\n'], + ['non-breaking space', '确认\u00a0AB12-CD34'], + ['full-width space', '确认 AB12-CD34'], + ['extra prose', '确认 AB12-CD34 请执行'], + ['malformed code', '确认 AB12-CD3'], + ['different verb', '同意 AB12-CD34'], + ])('keeps the %s control-like text on the ordinary echo path', async (_caseName, controlLikeText) => { + const fetchMock = vi.fn(async (input: RequestInfo | URL, init?: RequestInit) => { + if (String(input).endsWith('/message/agent/stream') && init?.method === 'POST') { + return createAgentStreamResponse([legacySseFrame({ type: 'done' })]) + } + return createAgentResponse([]) + }) + vi.stubGlobal('fetch', fetchMock) + + const wrapper = mountPanel() + await wrapper.find('textarea').setValue(controlLikeText) + await wrapper.find('textarea').trigger('keydown', { key: 'Enter' }) + await flushPromises() + + const streamCall = fetchMock.mock.calls.find(([input]) => String(input).endsWith('/message/agent/stream')) + const streamBody = JSON.parse(String(streamCall?.[1]?.body || '{}')) + expect(streamBody.echo_user).toBe(true) + expect(wrapper.find('.agent-assistant-message--user').exists()).toBe(true) + wrapper.unmount() + }) + + it('keeps an exact confirmation string with an attachment on the ordinary echo path', async () => { + const controlText = '确认 AB12-CD34' + const fetchMock = vi.fn(async (input: RequestInfo | URL, init?: RequestInit) => { + const url = String(input) + if (url.endsWith('/message/agent/upload') && init?.method === 'POST') { + return createAgentResponse({ + ref: 'attachment-ref', + url: '/api/v1/message/agent/attachments/attachment-ref', + name: 'proof.txt', + mime_type: 'text/plain', + size: 5, + kind: 'file', + }) + } + if (url.endsWith('/message/agent/stream') && init?.method === 'POST') { + return createAgentStreamResponse([legacySseFrame({ type: 'done' })]) + } + return createAgentResponse([]) + }) + vi.stubGlobal('fetch', fetchMock) + + const wrapper = mountPanel() + const fileInput = wrapper.get('input[type="file"]') + Object.defineProperty(fileInput.element, 'files', { + configurable: true, + value: [new File(['proof'], 'proof.txt', { type: 'text/plain' })], + }) + await fileInput.trigger('change') + await wrapper.find('textarea').setValue(controlText) + await wrapper.find('textarea').trigger('keydown', { key: 'Enter' }) + await flushPromises() + + const streamCall = fetchMock.mock.calls.find(([input]) => String(input).endsWith('/message/agent/stream')) + const streamBody = JSON.parse(String(streamCall?.[1]?.body || '{}')) + expect(streamBody.echo_user).toBe(true) + expect(wrapper.find('.agent-assistant-message--user').text()).toContain(controlText) + + const uploadCall = fetchMock.mock.calls.find(([input]) => String(input).endsWith('/message/agent/upload')) + const uploadHeaders = new Headers(uploadCall?.[1]?.headers) + expect(uploadHeaders.has('X-MoviePilot-Agent-Interaction')).toBe(false) + wrapper.unmount() + }) + + it('clears protected delivery through Escape and the close button', async () => { + const protectedMarker = 'MP-CLOSE-PROTECTED-MARKER' + const fetchMock = vi.fn(async (input: RequestInfo | URL, init?: RequestInit) => { + if (String(input).endsWith('/message/agent/stream') && init?.method === 'POST') { + return createAgentStreamResponse([ + protectedSseFrame({ + schema_version: '1.0', + delivery_id: `delivery-${fetchMock.mock.calls.length}`, + content_type: 'text/plain', + content: protectedMarker, + }), + legacySseFrame({ type: 'done' }), + ]) + } + return createAgentResponse([]) + }) + vi.stubGlobal('fetch', fetchMock) + + const wrapper = mountPanel() + await wrapper.find('textarea').setValue('读取后按 Escape') + await wrapper.find('textarea').trigger('keydown', { key: 'Enter' }) + await flushPromises() + expect(wrapper.text()).toContain(protectedMarker) + + window.dispatchEvent(new KeyboardEvent('keydown', { key: 'Escape' })) + await flushPromises() + expect(wrapper.text()).not.toContain(protectedMarker) + + await wrapper.setProps({ modelValue: true }) + await wrapper.find('textarea').setValue('读取后点击关闭') + await wrapper.find('textarea').trigger('keydown', { key: 'Enter' }) + await flushPromises() + expect(wrapper.text()).toContain(protectedMarker) + + await wrapper.get('[title="common.close"]').trigger('click') + expect(wrapper.text()).not.toContain(protectedMarker) + wrapper.unmount() + }) + it('renders interleaved assistant text and tool events in their SSE order', async () => { const serverSessionId = 'web-agent:ordered-segments' const streamEvents = [