This is an automated email from the ASF dual-hosted git repository.
lizhimins pushed a commit to branch rocketmq-studio
in repository https://gitbox.apache.org/repos/asf/rocketmq-dashboard.git
The following commit(s) were added to refs/heads/rocketmq-studio by this push:
new d9a1ac3 fix: preserve SSE framing for AI chat (#438)
d9a1ac3 is described below
commit d9a1ac3e96432d305afa7e590322d8cc98f37081
Author: btlqql <[email protected]>
AuthorDate: Fri Jul 17 17:15:31 2026 +0800
fix: preserve SSE framing for AI chat (#438)
---
web/src/api/ai.test.ts | 97 ++++++++++++++++++++++++++++++++++++++++++++++++++
web/src/api/ai.ts | 69 ++++++++++++++++++++++++++++-------
2 files changed, 153 insertions(+), 13 deletions(-)
diff --git a/web/src/api/ai.test.ts b/web/src/api/ai.test.ts
new file mode 100644
index 0000000..1b03825
--- /dev/null
+++ b/web/src/api/ai.test.ts
@@ -0,0 +1,97 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+import { afterEach, describe, expect, it, vi } from 'vitest';
+import { chatStream } from './ai';
+
+const encoder = new TextEncoder();
+
+function streamResponse(chunks: string[]): Response {
+ const body = new ReadableStream<Uint8Array>({
+ start(controller) {
+ chunks.forEach((chunk) => controller.enqueue(encoder.encode(chunk)));
+ controller.close();
+ },
+ });
+ return new Response(body, { status: 200 });
+}
+
+describe('AI chat SSE stream', () => {
+ afterEach(() => {
+ vi.unstubAllGlobals();
+ });
+
+ it('reassembles an event split across network chunks', async () => {
+ vi.stubGlobal(
+ 'fetch',
+ vi
+ .fn()
+ .mockResolvedValue(
+ streamResponse([
+ 'event: message\r\ndata: {"text":"hel',
+ 'lo"}\r\n\r\nevent: done\r\ndata: [DONE]\r\n\r\n',
+ ]),
+ ),
+ );
+ vi.stubGlobal('localStorage', { getItem: vi.fn().mockReturnValue('token')
});
+ const chunks: string[] = [];
+
+ await chatStream({ message: 'hello', mode: 'chat', model: 'stub' }, (text)
=>
+ chunks.push(text),
+ );
+
+ expect(chunks).toEqual(['hello']);
+ });
+
+ it('dispatches multiple events delivered in one network chunk', async () => {
+ vi.stubGlobal(
+ 'fetch',
+ vi
+ .fn()
+ .mockResolvedValue(
+ streamResponse([
+ 'data: {"content":"first"}\n\ndata: {"content":"second"}\n\ndata:
[DONE]\n\n',
+ ]),
+ ),
+ );
+ vi.stubGlobal('localStorage', { getItem: vi.fn().mockReturnValue(null) });
+ const chunks: string[] = [];
+
+ await chatStream({ message: 'hello', mode: 'chat', model: 'stub' }, (text)
=>
+ chunks.push(text),
+ );
+
+ expect(chunks).toEqual(['first', 'second']);
+ });
+
+ it('supports multiline and raw SSE data at end of stream', async () => {
+ vi.stubGlobal(
+ 'fetch',
+ vi
+ .fn()
+ .mockResolvedValue(streamResponse(['data: {"content":\ndata:
"hello"}\n\ndata: raw text'])),
+ );
+ vi.stubGlobal('localStorage', { getItem: vi.fn().mockReturnValue(null) });
+ const chunks: string[] = [];
+
+ await chatStream({ message: 'hello', mode: 'chat', model: 'stub' }, (text)
=>
+ chunks.push(text),
+ );
+
+ expect(chunks).toEqual(['hello', 'raw text']);
+ });
+});
diff --git a/web/src/api/ai.ts b/web/src/api/ai.ts
index f38bb02..547b3b4 100644
--- a/web/src/api/ai.ts
+++ b/web/src/api/ai.ts
@@ -31,6 +31,49 @@ export interface AiExecuteRequest {
tools?: string[];
}
+interface AiStreamPayload {
+ content?: unknown;
+ text?: unknown;
+}
+
+function getEventBoundary(buffer: string): { index: number; length: number } |
null {
+ const match = /\r\n\r\n|\n\n|\r\r/.exec(buffer);
+ return match ? { index: match.index, length: match[0].length } : null;
+}
+
+function getEventData(event: string): string | null {
+ const dataLines = event
+ .split(/\r\n|\r|\n/)
+ .filter((line) => line.startsWith('data:'))
+ .map((line) => {
+ const value = line.slice(5);
+ return value.startsWith(' ') ? value.slice(1) : value;
+ });
+
+ return dataLines.length ? dataLines.join('\n') : null;
+}
+
+function emitEvent(event: string, onChunk: (text: string) => void): boolean {
+ const payload = getEventData(event);
+ if (payload === null) return false;
+ if (payload === '[DONE]') return true;
+
+ try {
+ const parsed = JSON.parse(payload) as AiStreamPayload;
+ const text =
+ typeof parsed.content === 'string'
+ ? parsed.content
+ : typeof parsed.text === 'string'
+ ? parsed.text
+ : null;
+ if (text !== null) onChunk(text);
+ } catch {
+ onChunk(payload);
+ }
+
+ return false;
+}
+
// ─── AI ─────────────────────────────────────────────────────────
export async function chatStream(
data: AiExecuteRequest,
@@ -53,24 +96,24 @@ export async function chatStream(
const reader = response.body.getReader();
const decoder = new TextDecoder();
+ let buffer = '';
+
while (true) {
const { done, value } = await reader.read();
if (done) break;
- const chunk = decoder.decode(value, { stream: true });
- // SSE format: "data: ..."
- const lines = chunk.split('\n').filter((l) => l.startsWith('data: '));
- for (const line of lines) {
- const payload = line.slice(6);
- if (payload === '[DONE]') return;
- try {
- const parsed = JSON.parse(payload);
- if (parsed.content) onChunk(parsed.content);
- } catch {
- // raw text chunk
- onChunk(payload);
- }
+
+ buffer += decoder.decode(value, { stream: true });
+ let boundary = getEventBoundary(buffer);
+ while (boundary) {
+ const event = buffer.slice(0, boundary.index);
+ buffer = buffer.slice(boundary.index + boundary.length);
+ if (emitEvent(event, onChunk)) return;
+ boundary = getEventBoundary(buffer);
}
}
+
+ buffer += decoder.decode();
+ if (buffer && emitEvent(buffer, onChunk)) return;
}
export async function executeAiCommand(data: AiExecuteRequest) {