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 5d9e140a [ISSUE #1564] Enforce OpenAI-compatible SSE body timeout
(#1567)
5d9e140a is described below
commit 5d9e140a8c4f47df7de6c030b504a69980106b4d
Author: youngkermit8-coder <[email protected]>
AuthorDate: Tue Aug 11 20:24:14 2026 +0800
[ISSUE #1564] Enforce OpenAI-compatible SSE body timeout (#1567)
Signed-off-by: youngkermit8-coder <[email protected]>
---
.../studio/ops/ai/OpenAiCompatibleLlmClient.java | 48 +++++++++++++++++++++-
.../ops/ai/OpenAiCompatibleLlmClientTest.java | 37 +++++++++++++++++
2 files changed, 83 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 f8940149..3417c7ef 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
@@ -25,6 +25,7 @@ import org.springframework.util.StringUtils;
import java.io.BufferedReader;
import java.io.IOException;
+import java.io.InputStream;
import java.io.InputStreamReader;
import java.net.URI;
import java.net.URISyntaxException;
@@ -40,6 +41,10 @@ import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.Set;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.FutureTask;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
import java.util.function.Consumer;
@Component
@@ -130,7 +135,7 @@ public class OpenAiCompatibleLlmClient {
throw upstreamException(response.statusCode(),
new String(response.body().readAllBytes(),
StandardCharsets.UTF_8));
}
- parseStream(response, tokenConsumer);
+ parseStreamWithTimeout(response, tokenConsumer);
} catch (HttpTimeoutException exception) {
throw new LlmGatewayException(504, "llm.provider.timeout",
"LLM provider stream timed out",
@@ -146,7 +151,46 @@ public class OpenAiCompatibleLlmClient {
}
}
- private void parseStream(HttpResponse<java.io.InputStream> response,
Consumer<String> tokenConsumer)
+ private void parseStreamWithTimeout(HttpResponse<InputStream> response,
Consumer<String> tokenConsumer)
+ throws IOException, InterruptedException {
+ FutureTask<Void> readerTask = new FutureTask<>(() -> {
+ parseStream(response, tokenConsumer);
+ return null;
+ });
+ Thread.ofVirtual().name("openai-sse-reader").start(readerTask);
+ try {
+ readerTask.get(requestTimeout.toNanos(), TimeUnit.NANOSECONDS);
+ } catch (TimeoutException exception) {
+ throw new HttpTimeoutException("LLM provider stream timed out");
+ } catch (ExecutionException exception) {
+ Throwable cause = exception.getCause();
+ if (cause instanceof IOException ioException) {
+ throw ioException;
+ }
+ if (cause instanceof RuntimeException runtimeException) {
+ throw runtimeException;
+ }
+ if (cause instanceof Error error) {
+ throw error;
+ }
+ throw new IOException("Failed to consume LLM provider stream",
cause);
+ } finally {
+ if (!readerTask.isDone()) {
+ closeQuietly(response.body());
+ readerTask.cancel(true);
+ }
+ }
+ }
+
+ private void closeQuietly(InputStream stream) {
+ try {
+ stream.close();
+ } catch (IOException ignored) {
+ // Preserve the timeout or interruption that caused stream
cancellation.
+ }
+ }
+
+ private void parseStream(HttpResponse<InputStream> response,
Consumer<String> tokenConsumer)
throws IOException {
try (BufferedReader reader = new BufferedReader(
new InputStreamReader(response.body(),
StandardCharsets.UTF_8))) {
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 6fc960b7..ba122a2c 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
@@ -152,6 +152,43 @@ class OpenAiCompatibleLlmClientTest {
assertThat(requestBody.get().path("stream").asBoolean()).isTrue();
}
+ @Test
+ void streamShouldEnforceTimeoutWhileReadingResponseBody() {
+ OpenAiCompatibleLlmClient timeoutClient = new
OpenAiCompatibleLlmClient(
+ objectMapper,
+
HttpClient.newBuilder().connectTimeout(Duration.ofMillis(300)).build(),
+ Duration.ofMillis(300));
+ server.createContext("/v1/chat/completions", exchange -> {
+ exchange.getRequestBody().readAllBytes();
+ exchange.getResponseHeaders().set("Content-Type",
"text/event-stream");
+ exchange.sendResponseHeaders(200, 0);
+ try {
+ exchange.getResponseBody().write("""
+ data: {"choices":[{"delta":{"content":"first"}}]}
+
+ """.getBytes(StandardCharsets.UTF_8));
+ exchange.getResponseBody().flush();
+ Thread.sleep(1_500);
+ } catch (InterruptedException exception) {
+ Thread.currentThread().interrupt();
+ } finally {
+ exchange.close();
+ }
+ });
+ List<String> tokens = new ArrayList<>();
+
+ assertThatThrownBy(() -> timeoutClient.stream(
+ config("openai", "sk-test"), "hello", null, tokens::add))
+ .isInstanceOf(LlmGatewayException.class)
+ .hasMessage("LLM provider stream timed out")
+ .satisfies(exception -> {
+ LlmGatewayException gatewayException =
(LlmGatewayException) exception;
+
assertThat(gatewayException.getStatusCode()).isEqualTo(504);
+
assertThat(gatewayException.getCode()).isEqualTo("llm.provider.timeout");
+ });
+ assertThat(tokens).containsExactly("first");
+ }
+
@Test
void ollamaShouldAllowMissingApiKeyAndOmitAuthorizationHeader() {
AtomicReference<String> authorization = new AtomicReference<>();