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<>();

Reply via email to