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 10a7b4937 fix(ai): preserve whitespace-only streaming chunks (#4109)
10a7b4937 is described below
commit 10a7b49378877ae5ab13f4601f005cb147d10e16
Author: qiyu <[email protected]>
AuthorDate: Tue Sep 8 15:48:19 2026 +0800
fix(ai): preserve whitespace-only streaming chunks (#4109)
---
.../studio/ops/ai/OpenAiCompatibleLlmClient.java | 3 +-
.../studio/ops/ai/OpenAiCompatibleLlmGateway.java | 3 +-
.../ops/ai/OpenAiCompatibleLlmClientTest.java | 51 ++++++++++++++++++++++
.../ops/ai/OpenAiCompatibleLlmGatewayTest.java | 36 +++++++++++++++
4 files changed, 91 insertions(+), 2 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/OpenAiCompatibleLlmClient.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/OpenAiCompatibleLlmClient.java
index 654727609..63966eb36 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/OpenAiCompatibleLlmClient.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/OpenAiCompatibleLlmClient.java
@@ -270,7 +270,8 @@ public class OpenAiCompatibleLlmClient {
return true;
}
String token = parseDelta(data);
- if (StringUtils.hasText(token)) {
+ // Whitespace-only deltas carry formatting and must reach the consumer.
+ if (StringUtils.hasLength(token)) {
tokenConsumer.accept(token);
}
return false;
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/OpenAiCompatibleLlmGateway.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/OpenAiCompatibleLlmGateway.java
index cc1d55da7..07b29dc4a 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/OpenAiCompatibleLlmGateway.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/OpenAiCompatibleLlmGateway.java
@@ -208,7 +208,8 @@ public class OpenAiCompatibleLlmGateway implements
LlmGateway {
}
private void emitEnhanceChunk(LlmSseSession session, String chunk) {
- if (!StringUtils.hasText(chunk)) {
+ // Preserve whitespace in the streamed prompt preview.
+ if (!StringUtils.hasLength(chunk)) {
return;
}
try {
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/ai/OpenAiCompatibleLlmClientTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/OpenAiCompatibleLlmClientTest.java
index 619284bc8..86583f40d 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/ops/ai/OpenAiCompatibleLlmClientTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/OpenAiCompatibleLlmClientTest.java
@@ -23,6 +23,8 @@ import com.sun.net.httpserver.HttpServer;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
import java.io.IOException;
import java.net.InetSocketAddress;
@@ -155,6 +157,55 @@ class OpenAiCompatibleLlmClientTest {
assertThat(requestBody.get().path("stream").asBoolean()).isTrue();
}
+ @Test
+ void streamShouldIgnoreNonContentDeltasAndStopAtDoneTest() {
+ server.createContext("/v1/chat/completions", exchange ->
respond(exchange, 200, """
+ data: {"choices":[{"delta":{"role":"assistant"}}]}
+
+ data: {"choices":[{"delta":{"content":""}}]}
+
+ data: {"choices":[{"delta":{"content":null}}]}
+
+ data: {"choices":[{"delta":{"content":"hello"}}]}
+
+ data: {"choices":[{"delta":{},"finish_reason":"stop"}]}
+
+ data: {"choices":[],"usage":{"completion_tokens":1}}
+
+ data: [DONE]
+
+ data: {"choices":[{"delta":{"content":"ignored"}}]}
+
+ """, "text/event-stream"));
+
+ List<String> tokens = new ArrayList<>();
+ client.stream(config("openai", "sk-test"), "hello", null, tokens::add);
+
+ assertThat(tokens).containsExactly("hello");
+ }
+
+ @ParameterizedTest
+ @ValueSource(strings = {" ", "\n", "\r\n", "\t", " ", "\n\n"})
+ void streamShouldPreserveWhitespaceOnlyDeltasTest(String whitespace)
throws Exception {
+ String whitespaceJson = objectMapper.writeValueAsString(whitespace);
+ server.createContext("/v1/chat/completions", exchange ->
respond(exchange, 200, """
+ data: {"choices":[{"delta":{"content":"hello"}}]}
+
+ data: {"choices":[{"delta":{"content":%s}}]}
+
+ data: {"choices":[{"delta":{"content":"world"}}]}
+
+ data: [DONE]
+
+ """.formatted(whitespaceJson), "text/event-stream"));
+
+ List<String> tokens = new ArrayList<>();
+ client.stream(config("openai", "sk-test"), "hello", null, tokens::add);
+
+ assertThat(tokens).containsExactly("hello", whitespace, "world");
+ assertThat(String.join("", tokens)).isEqualTo("hello" + whitespace +
"world");
+ }
+
@Test
void streamShouldExposeErrorEnvelopeFromSuccessfulResponse() {
server.createContext("/v1/chat/completions", exchange ->
respond(exchange, 200, """
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/ai/OpenAiCompatibleLlmGatewayTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/OpenAiCompatibleLlmGatewayTest.java
index 0745de523..4e1bca231 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/ops/ai/OpenAiCompatibleLlmGatewayTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/OpenAiCompatibleLlmGatewayTest.java
@@ -266,6 +266,42 @@ class OpenAiCompatibleLlmGatewayTest {
}
}
+ @Test
+ void httpEnhanceShouldPreserveWhitespaceOnlyChunksTest() throws Exception {
+ ExecutorService executor = singleChatExecutor();
+ List<RecordingSseEmitter> emitters = new CopyOnWriteArrayList<>();
+ OpenAiCompatibleLlmGateway testedGateway = gateway(executor, emitters);
+ LlmConfigVO config = config("openai", "sk-test");
+ when(configService.getConfig()).thenReturn(config);
+ when(llmClient.supports(config)).thenReturn(true);
+ doAnswer(invocation -> {
+ @SuppressWarnings("unchecked")
+ Consumer<String> consumer = invocation.getArgument(3,
Consumer.class);
+ List.of("", "hello", " ", "world", "\n", "\r\n", "\t", " ",
"next").forEach(consumer);
+ return null;
+ }).doAnswer(invocation -> {
+ @SuppressWarnings("unchecked")
+ Consumer<String> consumer = invocation.getArgument(3,
Consumer.class);
+ consumer.accept("answer");
+ return null;
+ }).when(llmClient).stream(any(), any(), any(), any());
+ try {
+ testedGateway.chat(ChatDTO.builder().message("raw
prompt").enhance(true).build());
+ RecordingSseEmitter emitter = emitters.get(0);
+ assertThat(emitter.completedLatch.await(5,
TimeUnit.SECONDS)).isTrue();
+
+ assertThat(emitter.eventCount("event:enhance")).isEqualTo(8);
+ assertThat(emitter.eventText())
+ .contains("{\"delta\":\" \"}", "{\"delta\":\"\\n\"}",
"{\"delta\":\"\\r\\n\"}",
+ "{\"delta\":\"\\t\"}", "{\"delta\":\" \"}",
"{\"text\":\"answer\"}")
+ .doesNotContain("{\"delta\":\"\"}", "event:error");
+ assertThat(emitter.eventCount("event:done")).isEqualTo(1);
+ verify(llmClient).stream(eq(config), eq("hello world\n\r\n\t
next"), any(), any());
+ } finally {
+ testedGateway.destroy();
+ }
+ }
+
private OpenAiCompatibleLlmGateway gateway(ExecutorService executor,
List<RecordingSseEmitter>
emitters) {
return new OpenAiCompatibleLlmGateway(