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>
|
<modelVersion>4.0.0</modelVersion>
|
||||||
<modules>
|
<modules>
|
||||||
<module>muse-module-ai-api</module>
|
<module>muse-module-ai-api</module>
|
||||||
<module>muse-module-ai-contract-server</module>
|
|
||||||
<module>muse-module-ai-server</module>
|
<module>muse-module-ai-server</module>
|
||||||
</modules>
|
</modules>
|
||||||
<packaging>pom</packaging>
|
<packaging>pom</packaging>
|
||||||
|
|||||||
@ -315,19 +315,23 @@ export const aiHandlers = [
|
|||||||
// 模拟多批次的流式文本块推送
|
// 模拟多批次的流式文本块推送
|
||||||
const chunks = ['围绕', '“', task.prompt.slice(0, 12), '”', '展开', '新的', '情节', '转折。'];
|
const chunks = ['围绕', '“', task.prompt.slice(0, 12), '”', '展开', '新的', '情节', '转折。'];
|
||||||
|
|
||||||
for (const chunk of chunks) {
|
for (let index = 0; index < chunks.length; index += 1) {
|
||||||
// 模拟 Token 渲染的轻微延时效果
|
// 模拟 Token 渲染的轻微延时效果
|
||||||
await new Promise((resolve) => setTimeout(resolve, 150));
|
await new Promise((resolve) => setTimeout(resolve, 150));
|
||||||
|
|
||||||
// 按 SSE 格式要求,以 "data: " 开头并以 "\n\n" 结尾进行数据组装与入队
|
// 对齐 OpenAPI SSEChunkEvent:事件名走 SSE event 字段、data 为独立 JSON,sequenceNo 从 1 递增
|
||||||
controller.enqueue(
|
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(
|
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();
|
controller.close();
|
||||||
},
|
},
|
||||||
|
|||||||
@ -29,8 +29,8 @@ describe('AIPanel task stream contract', () => {
|
|||||||
const encoder = new TextEncoder();
|
const encoder = new TextEncoder();
|
||||||
const stream = new ReadableStream({
|
const stream = new ReadableStream({
|
||||||
start(controller) {
|
start(controller) {
|
||||||
controller.enqueue(encoder.encode('data: {"type":"chunk","content":"星海","sequenceNo":1}\n\n'));
|
controller.enqueue(encoder.encode('event: chunk\ndata: {"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: done\ndata: {"taskId":123,"suggestionId":456,"summary":"完成"}\n\n'));
|
||||||
controller.close();
|
controller.close();
|
||||||
},
|
},
|
||||||
});
|
});
|
||||||
|
|||||||
@ -1,5 +1,5 @@
|
|||||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
|
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
|
||||||
import { connectEventStream } from './sse';
|
import { connectAIStream, connectEventStream } from './sse';
|
||||||
|
|
||||||
type FetchCallRecorder = {
|
type FetchCallRecorder = {
|
||||||
mock: {
|
mock: {
|
||||||
@ -340,3 +340,108 @@ describe('connectEventStream fetch SSE', () => {
|
|||||||
expect(vi.getTimerCount()).toBe(0);
|
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;
|
onChunk?: (data: { content: string; sequenceNo: number }) => void;
|
||||||
/** 接收到 AI 影子层的质检结果推送 */
|
/** 接收到 AI 影子层的质检结果推送 */
|
||||||
onQualityCheck?: (data: { dimension: string; score: number; passed: boolean }) => void;
|
onQualityCheck?: (data: { dimension: string; score: number; passed: boolean }) => void;
|
||||||
/** AI 生成任务全部结束 */
|
/** AI 生成任务全部结束(载荷对齐 OpenAPI SSEDoneEvent.data) */
|
||||||
onDone?: (data: { generationId: string; candidateId: string }) => void;
|
onDone?: (data: { taskId: number; suggestionId: number; summary?: string }) => void;
|
||||||
/** 连接或解析发生错误 */
|
/** 连接或解析发生错误 */
|
||||||
onError?: (data: { code: string; message: 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 的限制。 */
|
/** 将浏览器内相对路径补齐为绝对 URL,兼容 Node/Vitest fetch 对相对 URL 的限制。 */
|
||||||
function normalizeStreamUrl(url: string): string {
|
function normalizeStreamUrl(url: string): string {
|
||||||
try {
|
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 生成专用流式通道连接 (独立连接,单次生命周期,支持中止控制器)
|
* 建立 AI 生成专用流式通道连接 (独立连接,单次生命周期,支持中止控制器)
|
||||||
* @param url 请求的目标流地址
|
* @param url 请求的目标流地址
|
||||||
@ -272,27 +259,34 @@ export function connectAIStream(
|
|||||||
throw new Error(`HTTP 异常: ${response.status}`);
|
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();
|
const reader = response.body?.getReader();
|
||||||
if (!reader) {
|
if (!reader) {
|
||||||
// 部分测试环境或 polyfill 不暴露 reader,但仍能一次性读取 text。
|
// 部分测试环境或 polyfill 不暴露 reader,但仍能一次性读取 text。
|
||||||
dispatchSSEText(await response.text(), handlers);
|
parser.consume(await response.text());
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
const decoder = new TextDecoder();
|
const decoder = new TextDecoder();
|
||||||
let buffer = '';
|
|
||||||
|
|
||||||
while (true) {
|
while (true) {
|
||||||
const { done, value } = await reader.read();
|
const { done, value } = await reader.read();
|
||||||
if (done) break;
|
if (done) break;
|
||||||
|
|
||||||
// 解码获取的二进制包并存入缓冲区
|
// 解码获取的二进制包并交给 parser 累积解析
|
||||||
buffer = dispatchSSEText(decoder.decode(value, { stream: true }), handlers, buffer);
|
parser.consume(decoder.decode(value, { stream: true }));
|
||||||
}
|
}
|
||||||
|
|
||||||
const tail = decoder.decode();
|
const tail = decoder.decode();
|
||||||
if (tail || buffer) {
|
if (tail) {
|
||||||
dispatchSSEText(tail, handlers, buffer);
|
parser.consume(tail);
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
.catch((err) => {
|
.catch((err) => {
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user