import { envConfig } from '@/config/env' export interface AiStreamEvent { type: 'start' | 'delta' | 'done' | 'error' | string content?: string code?: string message?: string seq?: number metadata?: Record timestamp?: number } export interface StreamAiSceneOptions { sceneCode: string inputs?: Record onStart?: (event: AiStreamEvent) => void onDelta?: (delta: string, output: string, event: AiStreamEvent) => void onDone?: (event: AiStreamEvent, output: string) => void onError?: (message: string, event?: AiStreamEvent) => void } interface RuntimeLog { status?: string outputText?: string errorMessage?: string } const createRequestId = () => `web-${Date.now()}-${Math.random().toString(36).slice(2, 10)}` const sleep = (ms: number) => new Promise(resolve => setTimeout(resolve, ms)) const parseSseFrame = (frame: string): AiStreamEvent | null => { const parsed = { type: 'message', data: '' } frame.split(/\r?\n/).forEach((line) => { if (line.startsWith('event:')) parsed.type = line.slice(6).trim() if (line.startsWith('data:')) parsed.data += line.slice(5).trim() }) if (!parsed.data) return null try { return JSON.parse(parsed.data) } catch { return { type: parsed.type, content: parsed.data } } } const authHeaders = () => { const token = localStorage.getItem('access_token') return token ? { Authorization: `Bearer ${token}` } : {} } const queryRuntimeResult = async (requestId: string): Promise => { const response = await fetch(`${envConfig.apiBaseUrl}/ai/runtime/result?requestId=${encodeURIComponent(requestId)}`, { headers: authHeaders() }) const payload = await response.json().catch(() => null) if (response.ok && payload?.code === 200 && payload.data) { return payload.data } throw new Error(payload?.message || 'AI 结果还在生成中') } const recoverRuntimeOutput = async ( requestId: string, callbacks: Pick ): Promise => { const maxAttempts = 150 for (let index = 0; index < maxAttempts; index++) { try { const log = await queryRuntimeResult(requestId) const output = String(log.outputText || '').trim() if (output) { callbacks.onDelta?.(output, output, { type: 'delta', metadata: { recovered: true, requestId } }) callbacks.onDone?.({ type: 'done', metadata: { recovered: true, requestId, source: 'call-log' } }, output) return output } if (log.status === 'failed') { throw new Error(log.errorMessage || 'AI 生成失败') } } catch (error) { if (index === maxAttempts - 1) { throw error } } await sleep(2000) } throw new Error('AI 生成结果暂时没有返回') } export const streamAiScene = async ({ sceneCode, inputs = {}, onStart, onDelta, onDone, onError }: StreamAiSceneOptions): Promise<{ output: string }> => { const requestId = String(inputs.requestId || createRequestId()) const requestInputs = { ...inputs, requestId } const controller = new AbortController() let buffer = '' let output = '' let closed = false let recovered = false let recoveryTimer: ReturnType | undefined let recoveryPromise: Promise | undefined const clearRecoveryTimer = () => { if (recoveryTimer) { clearTimeout(recoveryTimer) recoveryTimer = undefined } } const recoverOnce = () => { if (!recoveryPromise) { recoveryPromise = recoverRuntimeOutput(requestId, { onDelta, onDone }) } return recoveryPromise } const completeFromRecoveredOutput = async () => { if (closed) return try { const recoveredOutput = await recoverOnce() if (closed) return output = recoveredOutput recovered = true closed = true clearRecoveryTimer() controller.abort() } catch { if (!closed) recoveryPromise = undefined } } recoveryTimer = setTimeout(() => { completeFromRecoveredOutput() }, 8000) const finishRecovered = (event?: AiStreamEvent, message?: string) => { if (!output.trim()) return false recovered = true closed = true clearRecoveryTimer() onDone?.({ type: 'done', metadata: { ...(event?.metadata || {}), recovered: true, warningCode: event?.code, warningMessage: message || event?.message } }, output) return true } const recoverOrThrow = async (message?: string, event?: AiStreamEvent) => { if (finishRecovered(event, message)) return try { output = await recoverOnce() recovered = true closed = true clearRecoveryTimer() } catch (error: any) { const finalMessage = message || error?.message || 'AI 生成结果暂时没有返回' onError?.(finalMessage, event) throw new Error(finalMessage) } } const consumeText = (text: string) => { buffer += text const frames = buffer.split(/\r?\n\r?\n/) buffer = frames.pop() || '' frames.forEach((frame) => { const event = parseSseFrame(frame) if (!event) return if (event.type === 'start') { onStart?.(event) } else if (event.type === 'delta') { const delta = event.content || '' output += delta onDelta?.(delta, output, event) } else if (event.type === 'done') { closed = true clearRecoveryTimer() onDone?.(event, output) } else if (event.type === 'error') { throw Object.assign(new Error(event.message || event.code || 'AI 流式请求失败'), { event }) } }) } let response: Response try { response = await fetch(`${envConfig.apiBaseUrl}/ai/runtime/stream`, { method: 'POST', headers: { 'Content-Type': 'application/json', ...authHeaders() }, body: JSON.stringify({ sceneCode, requestId, inputs: requestInputs }), signal: controller.signal }) } catch (error: any) { if (!closed) await recoverOrThrow(error?.message) return { output } } if (!response.ok || !response.body) { await recoverOrThrow(`AI 流式请求失败(${response.status})`) return { output } } const reader = response.body.getReader() const decoder = new TextDecoder('utf-8') try { while (true) { const { value, done } = await reader.read() if (done) break consumeText(decoder.decode(value, { stream: true })) if (closed || recovered) break } if (!closed && !recovered) { consumeText(decoder.decode()) if (buffer.trim()) consumeText('\n\n') } } catch (error: any) { if (!closed) { await recoverOrThrow(error?.message, error?.event) } } finally { clearRecoveryTimer() } if (!output.trim() && !closed) { await recoverOrThrow('AI 生成结果暂时没有返回') } return { output } } export default { streamAiScene }