Compare commits
2 Commits
dev/2.0.0
...
chore/stud
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
69fcf6f925 | ||
|
|
adc7382acb |
@ -1,31 +0,0 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<project xmlns="http://maven.apache.org/POM/4.0.0"
|
||||
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
|
||||
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
|
||||
<parent>
|
||||
<groupId>cn.iocoder.cloud</groupId>
|
||||
<artifactId>muse-module-ai</artifactId>
|
||||
<version>${revision}</version>
|
||||
</parent>
|
||||
<modelVersion>4.0.0</modelVersion>
|
||||
<artifactId>muse-module-ai-contract-server</artifactId>
|
||||
<packaging>jar</packaging>
|
||||
<name>${project.artifactId}</name>
|
||||
<description>AI 模块 P1 Muse API 合同入口,避免聚合服务强依赖完整 AI Runtime。</description>
|
||||
|
||||
<dependencies>
|
||||
<dependency>
|
||||
<groupId>cn.iocoder.cloud</groupId>
|
||||
<artifactId>muse-module-ai-api</artifactId>
|
||||
<version>${revision}</version>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>cn.iocoder.cloud</groupId>
|
||||
<artifactId>muse-spring-boot-starter-security</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>cn.iocoder.cloud</groupId>
|
||||
<artifactId>muse-spring-boot-starter-mybatis</artifactId>
|
||||
</dependency>
|
||||
</dependencies>
|
||||
</project>
|
||||
@ -10,7 +10,6 @@
|
||||
<modelVersion>4.0.0</modelVersion>
|
||||
<modules>
|
||||
<module>muse-module-ai-api</module>
|
||||
<module>muse-module-ai-contract-server</module>
|
||||
<module>muse-module-ai-server</module>
|
||||
</modules>
|
||||
<packaging>pom</packaging>
|
||||
|
||||
@ -315,19 +315,23 @@ export const aiHandlers = [
|
||||
// 模拟多批次的流式文本块推送
|
||||
const chunks = ['围绕', '“', task.prompt.slice(0, 12), '”', '展开', '新的', '情节', '转折。'];
|
||||
|
||||
for (const chunk of chunks) {
|
||||
for (let index = 0; index < chunks.length; index += 1) {
|
||||
// 模拟 Token 渲染的轻微延时效果
|
||||
await new Promise((resolve) => setTimeout(resolve, 150));
|
||||
|
||||
// 按 SSE 格式要求,以 "data: " 开头并以 "\n\n" 结尾进行数据组装与入队
|
||||
// 对齐 OpenAPI SSEChunkEvent:事件名走 SSE event 字段、data 为独立 JSON,sequenceNo 从 1 递增
|
||||
controller.enqueue(
|
||||
encoder.encode(`data: ${JSON.stringify({ type: 'chunk', content: chunk, sequenceNo: chunks.indexOf(chunk) })}\n\n`)
|
||||
encoder.encode(
|
||||
`event: chunk\ndata: ${JSON.stringify({ content: chunks[index], sequenceNo: index + 1 })}\n\n`
|
||||
)
|
||||
);
|
||||
}
|
||||
|
||||
// 发送流结束标记完成整个推送
|
||||
// 发送结束事件:对齐 OpenAPI SSEDoneEvent.data(taskId/suggestionId 为 int64,summary 可选)
|
||||
controller.enqueue(
|
||||
encoder.encode(`data: ${JSON.stringify({ type: 'done', generationId: 'g-mock-123', candidateId: 'c-mock-456' })}\n\n`)
|
||||
encoder.encode(
|
||||
`event: done\ndata: ${JSON.stringify({ taskId: 90001, suggestionId: 90002, summary: 'AI 生成完成' })}\n\n`
|
||||
)
|
||||
);
|
||||
controller.close();
|
||||
},
|
||||
|
||||
@ -29,8 +29,8 @@ describe('AIPanel task stream contract', () => {
|
||||
const encoder = new TextEncoder();
|
||||
const stream = new ReadableStream({
|
||||
start(controller) {
|
||||
controller.enqueue(encoder.encode('data: {"type":"chunk","content":"星海","sequenceNo":1}\n\n'));
|
||||
controller.enqueue(encoder.encode('data: {"type":"done","generationId":"g-1","candidateId":"c-1"}\n\n'));
|
||||
controller.enqueue(encoder.encode('event: chunk\ndata: {"content":"星海","sequenceNo":1}\n\n'));
|
||||
controller.enqueue(encoder.encode('event: done\ndata: {"taskId":123,"suggestionId":456,"summary":"完成"}\n\n'));
|
||||
controller.close();
|
||||
},
|
||||
});
|
||||
|
||||
@ -1,5 +1,5 @@
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
|
||||
import { connectEventStream } from './sse';
|
||||
import { connectAIStream, connectEventStream } from './sse';
|
||||
|
||||
type FetchCallRecorder = {
|
||||
mock: {
|
||||
@ -340,3 +340,108 @@ describe('connectEventStream fetch SSE', () => {
|
||||
expect(vi.getTimerCount()).toBe(0);
|
||||
});
|
||||
});
|
||||
|
||||
describe('connectAIStream fetch SSE (OpenAPI 契约对齐)', () => {
|
||||
afterEach(() => {
|
||||
vi.restoreAllMocks();
|
||||
});
|
||||
|
||||
it('should_dispatchOnChunk_fromEventNamedChunk', async () => {
|
||||
vi.spyOn(globalThis, 'fetch').mockResolvedValue(
|
||||
createStreamResponse(['event: chunk\ndata: {"content":"星海","sequenceNo":1}\n\n'])
|
||||
);
|
||||
const onChunk = vi.fn();
|
||||
|
||||
connectAIStream('/app-api/muse/ai/tasks/1/stream', { onChunk });
|
||||
await flushMicrotasks();
|
||||
|
||||
// 事件类型来自 SSE event 字段,data 是独立 JSON,回调只拿到 data 载荷
|
||||
expect(onChunk).toHaveBeenCalledWith({ content: '星海', sequenceNo: 1 });
|
||||
});
|
||||
|
||||
it('should_dispatchOnDone_withTaskIdSuggestionIdSummary', async () => {
|
||||
vi.spyOn(globalThis, 'fetch').mockResolvedValue(
|
||||
createStreamResponse(['event: done\ndata: {"taskId":123,"suggestionId":456,"summary":"完成"}\n\n'])
|
||||
);
|
||||
const onDone = vi.fn();
|
||||
|
||||
connectAIStream('/app-api/muse/ai/tasks/1/stream', { onDone });
|
||||
await flushMicrotasks();
|
||||
|
||||
// done 载荷必须是契约 SSEDoneEvent.data 形状,而非旧的 generationId/candidateId
|
||||
expect(onDone).toHaveBeenCalledWith({ taskId: 123, suggestionId: 456, summary: '完成' });
|
||||
});
|
||||
|
||||
it('should_dispatchOnQualityCheck_fromEventNamedQualityCheck', async () => {
|
||||
vi.spyOn(globalThis, 'fetch').mockResolvedValue(
|
||||
createStreamResponse(['event: quality_check\ndata: {"dimension":"safety","score":0.9,"passed":true}\n\n'])
|
||||
);
|
||||
const onQualityCheck = vi.fn();
|
||||
|
||||
connectAIStream('/app-api/muse/ai/tasks/1/stream', { onQualityCheck });
|
||||
await flushMicrotasks();
|
||||
|
||||
expect(onQualityCheck).toHaveBeenCalledWith({ dimension: 'safety', score: 0.9, passed: true });
|
||||
});
|
||||
|
||||
it('should_dispatchOnError_fromEventNamedError', async () => {
|
||||
vi.spyOn(globalThis, 'fetch').mockResolvedValue(
|
||||
createStreamResponse(['event: error\ndata: {"code":"MUSE-AI-001-0001","message":"模型超时"}\n\n'])
|
||||
);
|
||||
const onError = vi.fn();
|
||||
|
||||
connectAIStream('/app-api/muse/ai/tasks/1/stream', { onError });
|
||||
await flushMicrotasks();
|
||||
|
||||
expect(onError).toHaveBeenCalledWith({ code: 'MUSE-AI-001-0001', message: '模型超时' });
|
||||
});
|
||||
|
||||
it('should_dispatchChunksThenDone_inOrder', async () => {
|
||||
vi.spyOn(globalThis, 'fetch').mockResolvedValue(
|
||||
createStreamResponse([
|
||||
'event: chunk\ndata: {"content":"星","sequenceNo":1}\n\n',
|
||||
'event: chunk\ndata: {"content":"海","sequenceNo":2}\n\n',
|
||||
'event: done\ndata: {"taskId":7,"suggestionId":8}\n\n',
|
||||
])
|
||||
);
|
||||
const onChunk = vi.fn();
|
||||
const onDone = vi.fn();
|
||||
|
||||
connectAIStream('/app-api/muse/ai/tasks/1/stream', { onChunk, onDone });
|
||||
await flushMicrotasks();
|
||||
|
||||
expect(onChunk).toHaveBeenNthCalledWith(1, { content: '星', sequenceNo: 1 });
|
||||
expect(onChunk).toHaveBeenNthCalledWith(2, { content: '海', sequenceNo: 2 });
|
||||
expect(onDone).toHaveBeenCalledWith({ taskId: 7, suggestionId: 8 });
|
||||
});
|
||||
|
||||
it('should_notDispatch_whenLegacyInnerTypeShapeWithoutEventName', async () => {
|
||||
// 回归护栏:旧的“内层 data.type、无 SSE event 名”形状不应再被识别为业务事件
|
||||
vi.spyOn(globalThis, 'fetch').mockResolvedValue(
|
||||
createStreamResponse(['data: {"type":"chunk","content":"x","sequenceNo":1}\n\n'])
|
||||
);
|
||||
const onChunk = vi.fn();
|
||||
const onError = vi.fn();
|
||||
|
||||
connectAIStream('/app-api/muse/ai/tasks/1/stream', { onChunk, onError });
|
||||
await flushMicrotasks();
|
||||
|
||||
expect(onChunk).not.toHaveBeenCalled();
|
||||
expect(onError).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('should_dispatchOnError_whenDataJsonInvalid', async () => {
|
||||
vi.spyOn(globalThis, 'fetch').mockResolvedValue(
|
||||
createStreamResponse(['event: chunk\ndata: {not-json}\n\n'])
|
||||
);
|
||||
const onError = vi.fn();
|
||||
|
||||
connectAIStream('/app-api/muse/ai/tasks/1/stream', { onError });
|
||||
await flushMicrotasks();
|
||||
|
||||
expect(onError).toHaveBeenCalledWith({
|
||||
code: 'SSE_PARSE_ERROR',
|
||||
message: expect.any(String),
|
||||
});
|
||||
});
|
||||
});
|
||||
|
||||
@ -6,56 +6,12 @@ export interface SSEEventHandler {
|
||||
onChunk?: (data: { content: string; sequenceNo: number }) => void;
|
||||
/** 接收到 AI 影子层的质检结果推送 */
|
||||
onQualityCheck?: (data: { dimension: string; score: number; passed: boolean }) => void;
|
||||
/** AI 生成任务全部结束 */
|
||||
onDone?: (data: { generationId: string; candidateId: string }) => void;
|
||||
/** AI 生成任务全部结束(载荷对齐 OpenAPI SSEDoneEvent.data) */
|
||||
onDone?: (data: { taskId: number; suggestionId: number; summary?: string }) => void;
|
||||
/** 连接或解析发生错误 */
|
||||
onError?: (data: { code: string; message: string }) => void;
|
||||
}
|
||||
|
||||
/** 解析单行 SSE data 事件,并派发到应用层处理器。 */
|
||||
function dispatchSSELine(line: string, handlers: SSEEventHandler): void {
|
||||
const trimmed = line.trim();
|
||||
if (!trimmed.startsWith('data: ')) return;
|
||||
|
||||
try {
|
||||
const rawData = trimmed.slice(6);
|
||||
const data = JSON.parse(rawData);
|
||||
|
||||
// 触发相应的应用层回调
|
||||
switch (data.type) {
|
||||
case 'chunk':
|
||||
handlers.onChunk?.(data);
|
||||
break;
|
||||
case 'quality_check':
|
||||
handlers.onQualityCheck?.(data);
|
||||
break;
|
||||
case 'done':
|
||||
handlers.onDone?.(data);
|
||||
break;
|
||||
case 'error':
|
||||
handlers.onError?.(data);
|
||||
break;
|
||||
}
|
||||
} catch (parseErr) {
|
||||
handlers.onError?.({
|
||||
code: 'SSE_PARSE_ERROR',
|
||||
message: parseErr instanceof Error ? parseErr.message : 'JSON 解析失败',
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
/** 按 SSE 换行协议解析文本块,返回尚未闭合的半行缓冲。 */
|
||||
function dispatchSSEText(text: string, handlers: SSEEventHandler, previousBuffer = ''): string {
|
||||
const lines = `${previousBuffer}${text}`.split('\n');
|
||||
const nextBuffer = lines.pop() || '';
|
||||
|
||||
for (const line of lines) {
|
||||
dispatchSSELine(line, handlers);
|
||||
}
|
||||
|
||||
return nextBuffer;
|
||||
}
|
||||
|
||||
/** 将浏览器内相对路径补齐为绝对 URL,兼容 Node/Vitest fetch 对相对 URL 的限制。 */
|
||||
function normalizeStreamUrl(url: string): string {
|
||||
try {
|
||||
@ -248,6 +204,37 @@ function createEventStreamParser(callbacks: {
|
||||
};
|
||||
}
|
||||
|
||||
/** 按 OpenAPI SSE 契约,将命名事件(event 字段)派发到 AI 流处理器。 */
|
||||
function dispatchAIStreamMessage(message: ParsedEventStreamMessage, handlers: SSEEventHandler): void {
|
||||
let data: unknown;
|
||||
try {
|
||||
// event 名与 data 解耦:data 是独立 JSON 对象,先解析再按事件名派发
|
||||
data = JSON.parse(message.data);
|
||||
} catch (parseErr) {
|
||||
handlers.onError?.({
|
||||
code: 'SSE_PARSE_ERROR',
|
||||
message: parseErr instanceof Error ? parseErr.message : 'JSON 解析失败',
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
// 事件类型来自 SSE event 字段,而非内层 data.type;未声明的事件名(如 notification)忽略
|
||||
switch (message.event) {
|
||||
case 'chunk':
|
||||
handlers.onChunk?.(data as { content: string; sequenceNo: number });
|
||||
break;
|
||||
case 'quality_check':
|
||||
handlers.onQualityCheck?.(data as { dimension: string; score: number; passed: boolean });
|
||||
break;
|
||||
case 'done':
|
||||
handlers.onDone?.(data as { taskId: number; suggestionId: number; summary?: string });
|
||||
break;
|
||||
case 'error':
|
||||
handlers.onError?.(data as { code: string; message: string });
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 建立 AI 生成专用流式通道连接 (独立连接,单次生命周期,支持中止控制器)
|
||||
* @param url 请求的目标流地址
|
||||
@ -272,27 +259,34 @@ export function connectAIStream(
|
||||
throw new Error(`HTTP 异常: ${response.status}`);
|
||||
}
|
||||
|
||||
// 复用与全局事件流一致的 SSE parser,按 event/id/data 规范解析;AI 流为单次生成,不做断线续传,故忽略 id。
|
||||
const parser = createEventStreamParser({
|
||||
onMessage: (message) => dispatchAIStreamMessage(message, handlers),
|
||||
onLastEventId: () => {
|
||||
/* AI 单次生成流无需断线续传游标,忽略 id 字段 */
|
||||
},
|
||||
});
|
||||
|
||||
const reader = response.body?.getReader();
|
||||
if (!reader) {
|
||||
// 部分测试环境或 polyfill 不暴露 reader,但仍能一次性读取 text。
|
||||
dispatchSSEText(await response.text(), handlers);
|
||||
parser.consume(await response.text());
|
||||
return;
|
||||
}
|
||||
|
||||
const decoder = new TextDecoder();
|
||||
let buffer = '';
|
||||
|
||||
while (true) {
|
||||
const { done, value } = await reader.read();
|
||||
if (done) break;
|
||||
|
||||
// 解码获取的二进制包并存入缓冲区
|
||||
buffer = dispatchSSEText(decoder.decode(value, { stream: true }), handlers, buffer);
|
||||
// 解码获取的二进制包并交给 parser 累积解析
|
||||
parser.consume(decoder.decode(value, { stream: true }));
|
||||
}
|
||||
|
||||
const tail = decoder.decode();
|
||||
if (tail || buffer) {
|
||||
dispatchSSEText(tail, handlers, buffer);
|
||||
if (tail) {
|
||||
parser.consume(tail);
|
||||
}
|
||||
})
|
||||
.catch((err) => {
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user