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 1e74d140 fix(ai): bound CLI completion output and probe engine 
availability (#1721)
1e74d140 is described below

commit 1e74d1401e9c8a37058daf4a13af7373ba044535
Author: youngkermit8-coder <[email protected]>
AuthorDate: Thu Aug 13 19:48:43 2026 +0800

    fix(ai): bound CLI completion output and probe engine availability (#1721)
    
    Consolidates #1721 and #1736: cap buffered CLI completion output at 5MiB
    and terminate the subprocess on overflow, and test the selected CLI
    engine availability through the provider registry.
---
 .../rocketmq/studio/ops/ai/CliAgentProvider.java   | 63 ++++++++++++++++++++--
 .../rocketmq/studio/ops/ai/LlmConfigService.java   | 20 +++++++
 .../studio/ops/ai/CliAgentProviderTest.java        | 31 +++++++++++
 .../studio/ops/ai/LlmConfigServiceTest.java        | 59 ++++++++++++++++++--
 4 files changed, 164 insertions(+), 9 deletions(-)

diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/CliAgentProvider.java 
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/CliAgentProvider.java
index 47b37d0c..22411ffe 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/CliAgentProvider.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/CliAgentProvider.java
@@ -19,12 +19,17 @@ package org.apache.rocketmq.studio.ops.ai;
 import lombok.extern.slf4j.Slf4j;
 import org.springframework.util.StringUtils;
 
+import java.io.ByteArrayOutputStream;
 import java.io.IOException;
+import java.io.InputStream;
 import java.nio.charset.StandardCharsets;
 import java.util.List;
 import java.util.Map;
 import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CompletionException;
+import java.util.concurrent.ExecutionException;
 import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
 
 /**
  * Base class for CLI-based agent providers: spawns the vendor CLI in a
@@ -35,6 +40,7 @@ import java.util.concurrent.TimeUnit;
 public abstract class CliAgentProvider implements AgentProvider {
 
     private static final long TIMEOUT_SECONDS = 180;
+    private static final int MAX_OUTPUT_BYTES = 5 * 1024 * 1024;
 
     protected abstract List<String> buildCommand(LlmConfigVO config, String 
prompt, String modelOverride);
 
@@ -84,10 +90,10 @@ public abstract class CliAgentProvider implements 
AgentProvider {
                     "Check that the CLI binary is installed and executable.", 
exception);
         }
         CompletableFuture<String> outputFuture = 
CompletableFuture.supplyAsync(() -> {
-            try (java.io.InputStream in = process.getInputStream()) {
-                return new String(in.readAllBytes(), StandardCharsets.UTF_8);
-            } catch (IOException ignored) {
-                return "";
+            try {
+                return readOutput(process);
+            } catch (IOException exception) {
+                throw new CompletionException(exception);
             }
         });
         boolean finished;
@@ -102,7 +108,20 @@ public abstract class CliAgentProvider implements 
AgentProvider {
         String output;
         try {
             output = outputFuture.get(10, TimeUnit.SECONDS);
-        } catch (Exception exception) {
+        } catch (ExecutionException exception) {
+            if (exception.getCause() instanceof OutputLimitException 
outputLimitException) {
+                throw new LlmGatewayException(502, 
"llm.provider.output_too_large",
+                        binaryName() + " CLI output exceeded the maximum of "
+                                + outputLimitException.limitBytes() + " bytes",
+                        "Retry with a shorter prompt or reduce the provider 
response size.", exception);
+            }
+            output = "";
+        } catch (InterruptedException exception) {
+            process.destroyForcibly();
+            Thread.currentThread().interrupt();
+            throw new LlmGatewayException(502, "llm.provider.interrupted",
+                    binaryName() + " CLI output collection was interrupted", 
"Retry the request.", exception);
+        } catch (TimeoutException exception) {
             output = "";
         }
         if (!finished) {
@@ -126,6 +145,27 @@ public abstract class CliAgentProvider implements 
AgentProvider {
         return result;
     }
 
+    private String readOutput(Process process) throws IOException {
+        int limitBytes = outputLimitBytes();
+        try (InputStream input = process.getInputStream();
+             ByteArrayOutputStream output = new 
ByteArrayOutputStream(Math.min(limitBytes, 8192))) {
+            byte[] buffer = new byte[8192];
+            int read;
+            while ((read = input.read(buffer)) != -1) {
+                if (read > limitBytes - output.size()) {
+                    process.destroyForcibly();
+                    throw new OutputLimitException(limitBytes);
+                }
+                output.write(buffer, 0, read);
+            }
+            return output.toString(StandardCharsets.UTF_8);
+        }
+    }
+
+    int outputLimitBytes() {
+        return MAX_OUTPUT_BYTES;
+    }
+
     private String abbreviate(String value) {
         if (value == null) {
             return "";
@@ -133,4 +173,17 @@ public abstract class CliAgentProvider implements 
AgentProvider {
         String trimmed = value.trim();
         return trimmed.length() <= 500 ? trimmed : trimmed.substring(0, 500) + 
"...";
     }
+
+    private static final class OutputLimitException extends IOException {
+        private final int limitBytes;
+
+        private OutputLimitException(int limitBytes) {
+            super("CLI output exceeds " + limitBytes + " bytes");
+            this.limitBytes = limitBytes;
+        }
+
+        private int limitBytes() {
+            return limitBytes;
+        }
+    }
 }
diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/LlmConfigService.java 
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/LlmConfigService.java
index 7cf2b8fd..e2b9ad66 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/LlmConfigService.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/LlmConfigService.java
@@ -75,6 +75,7 @@ public class LlmConfigService {
 
     private final SettingsService settingsService;
     private final OpenAiCompatibleLlmClient llmClient;
+    private final AgentProviderRegistry agentProviders;
     private final LlmProperties llmProperties;
     private LlmConfigVO overrides;
 
@@ -128,6 +129,9 @@ public class LlmConfigService {
         if (validation.getStatus() != 0) {
             return validation;
         }
+        if 
(!LlmConfigVO.ENGINE_HTTP.equalsIgnoreCase(normalized.normalizeEngine())) {
+            return testCliEngine(normalized.normalizeEngine());
+        }
         if (llmClient.supports(normalized)) {
             try {
                 llmClient.listModels(normalized);
@@ -138,6 +142,22 @@ public class LlmConfigService {
         return LlmOperationResultVO.success("Connection successful");
     }
 
+    private LlmOperationResultVO testCliEngine(String engine) {
+        try {
+            AgentProvider provider = agentProviders.forEngine(engine);
+            if (provider.available()) {
+                return LlmOperationResultVO.success("CLI is available");
+            }
+            return LlmOperationResultVO.failure(
+                    "llm.provider.cli_missing",
+                    "Agent CLI is not available for engine: " + engine,
+                    "Install the selected CLI in the server runtime or select 
the HTTP engine.");
+        } catch (LlmGatewayException exception) {
+            return LlmOperationResultVO.failure(
+                    exception.getCode(), exception.getMessage(), 
exception.getHint());
+        }
+    }
+
     private LlmOperationResultVO validate(LlmConfigVO normalized) {
         String provider = normalized.getProvider();
         if (!isValidApiBase(normalized.getApiBase())) {
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/ops/ai/CliAgentProviderTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/CliAgentProviderTest.java
index 9a1fcdcd..54a2226f 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/ops/ai/CliAgentProviderTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/CliAgentProviderTest.java
@@ -16,14 +16,21 @@ import java.util.List;
 import java.util.Map;
 
 import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
 
 class CliAgentProviderTest {
 
     private static final class FakeCli extends CliAgentProvider {
         private final String script;
+        private final int outputLimitBytes;
 
         FakeCli(String script) {
+            this(script, Integer.MAX_VALUE);
+        }
+
+        FakeCli(String script, int outputLimitBytes) {
             this.script = script;
+            this.outputLimitBytes = outputLimitBytes;
         }
 
         @Override
@@ -45,6 +52,11 @@ class CliAgentProviderTest {
         protected String binaryName() {
             return "sh";
         }
+
+        @Override
+        int outputLimitBytes() {
+            return outputLimitBytes;
+        }
     }
 
     @Test
@@ -63,4 +75,23 @@ class CliAgentProviderTest {
         // well past the 64 KiB pipe buffer that used to deadlock the 
sequential reads.
         assertThat(result.length()).isGreaterThan(500_000);
     }
+
+    @Test
+    void completeRejectsOutputBeyondConfiguredLimit() {
+        FakeCli cli = new FakeCli("yes 0123456789abcdef | head -c 4096", 1024);
+
+        assertThatThrownBy(() -> cli.complete(null, "prompt", null))
+                .isInstanceOfSatisfying(LlmGatewayException.class, exception 
-> {
+                    assertThat(exception.getStatusCode()).isEqualTo(502);
+                    
assertThat(exception.getCode()).isEqualTo("llm.provider.output_too_large");
+                    assertThat(exception.getMessage()).contains("1024 bytes");
+                });
+    }
+
+    @Test
+    void completeAllowsOutputAtConfiguredLimit() {
+        FakeCli cli = new FakeCli("yes x | head -c 1024", 1024);
+
+        assertThat(cli.complete(null, "prompt", null)).isNotEmpty();
+    }
 }
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/ops/ai/LlmConfigServiceTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/LlmConfigServiceTest.java
index 584870fe..f942b2f4 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/ops/ai/LlmConfigServiceTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/LlmConfigServiceTest.java
@@ -37,12 +37,14 @@ import static org.mockito.Mockito.doThrow;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoInteractions;
 import static org.mockito.Mockito.when;
 
 class LlmConfigServiceTest {
 
     private SettingsService settingsService;
     private OpenAiCompatibleLlmClient llmClient;
+    private AgentProviderRegistry agentProviders;
     private LlmConfigService llmConfigService;
 
     @BeforeEach
@@ -61,7 +63,9 @@ class LlmConfigServiceTest {
                 .baseUrl("https://api.openai.com/v1";)
                 .build());
         llmClient = mock(OpenAiCompatibleLlmClient.class);
-        llmConfigService = new LlmConfigService(settingsService, llmClient, 
new LlmProperties());
+        agentProviders = mock(AgentProviderRegistry.class);
+        llmConfigService = new LlmConfigService(
+                settingsService, llmClient, agentProviders, new 
LlmProperties());
     }
 
     @Test
@@ -80,7 +84,7 @@ class LlmConfigServiceTest {
     void envTokenShouldOverrideApiKeyAtRuntimeButNeverBePersisted() {
         LlmProperties properties = new LlmProperties();
         properties.setToken("env-token");
-        LlmConfigService service = new LlmConfigService(settingsService, 
llmClient, properties);
+        LlmConfigService service = new LlmConfigService(settingsService, 
llmClient, agentProviders, properties);
 
         LlmConfigVO config = service.getConfig();
         assertThat(config.getApiKey()).isEqualTo("env-token");
@@ -207,7 +211,7 @@ class LlmConfigServiceTest {
         assertThat(persisted.getAwsRegion()).isEqualTo("eu-west-1");
 
         when(settingsService.getGeneralSettings()).thenReturn(persisted);
-        LlmConfigVO reloaded = new LlmConfigService(settingsService, 
llmClient, new LlmProperties()).getConfig();
+        LlmConfigVO reloaded = new LlmConfigService(settingsService, 
llmClient, agentProviders, new LlmProperties()).getConfig();
         assertThat(reloaded.getDeploymentName()).isEqualTo("production-gpt");
         assertThat(reloaded.getApiVersion()).isEqualTo("2024-06-01");
         assertThat(reloaded.getAwsRegion()).isEqualTo("eu-west-1");
@@ -321,6 +325,7 @@ class LlmConfigServiceTest {
     void testConfigShouldAllowOllamaWithoutApiKey() {
         LlmOperationResultVO result = 
llmConfigService.testConfig(LlmConfigVO.builder()
                 .provider("ollama")
+                .engine("http")
                 .apiBase("http://localhost:11434/v1";)
                 .model("llama3")
                 .build());
@@ -337,6 +342,7 @@ class LlmConfigServiceTest {
 
         LlmOperationResultVO result = 
llmConfigService.testConfig(LlmConfigVO.builder()
                 .provider("openai")
+                .engine("http")
                 .apiBase("https://api.openai.com/v1";)
                 .model("gpt-4o")
                 .maxTokens(2048)
@@ -354,13 +360,14 @@ class LlmConfigServiceTest {
     void testConfigShouldPreferEnvironmentTokenOverStoredApiKey() {
         LlmProperties properties = new LlmProperties();
         properties.setToken("env-token");
-        LlmConfigService service = new LlmConfigService(settingsService, 
llmClient, properties);
+        LlmConfigService service = new LlmConfigService(settingsService, 
llmClient, agentProviders, properties);
         
when(llmClient.supports(org.mockito.ArgumentMatchers.any())).thenReturn(true);
         
when(llmClient.listModels(org.mockito.ArgumentMatchers.any())).thenReturn(List.of(
                 new LlmModelItemVO("gpt-4o", "GPT-4o")));
 
         LlmOperationResultVO result = service.testConfig(LlmConfigVO.builder()
                 .provider("openai")
+                .engine(LlmConfigVO.ENGINE_HTTP)
                 .apiBase("https://api.openai.com/v1";)
                 .model("gpt-4o")
                 .maxTokens(2048)
@@ -385,6 +392,7 @@ class LlmConfigServiceTest {
 
         LlmOperationResultVO result = 
llmConfigService.testConfig(LlmConfigVO.builder()
                 .provider("openai")
+                .engine("http")
                 .apiKey("sk-bad")
                 .apiBase("https://api.openai.com/v1";)
                 .model("gpt-4o")
@@ -398,6 +406,49 @@ class LlmConfigServiceTest {
         assertThat(result.getHint()).contains("credentials");
     }
 
+    @Test
+    void testConfigShouldNotProbeHttpModelsForCliEngine() {
+        AgentProvider provider = mock(AgentProvider.class);
+        when(agentProviders.forEngine("claude-code")).thenReturn(provider);
+        when(provider.available()).thenReturn(true);
+
+        LlmOperationResultVO result = 
llmConfigService.testConfig(LlmConfigVO.builder()
+                .provider("openai")
+                .engine("claude-code")
+                .apiBase("https://api.openai.com/v1";)
+                .model("claude-sonnet-4")
+                .maxTokens(2048)
+                .temperature(1.0)
+                .build());
+
+        assertThat(result.getStatus()).isZero();
+        assertThat(result.getMsg()).isEqualTo("CLI is available");
+        verify(agentProviders).forEngine("claude-code");
+        verifyNoInteractions(llmClient);
+    }
+
+    @Test
+    void testConfigShouldReportMissingCliEngine() {
+        AgentProvider provider = mock(AgentProvider.class);
+        when(agentProviders.forEngine("qoder")).thenReturn(provider);
+        when(provider.available()).thenReturn(false);
+
+        LlmOperationResultVO result = 
llmConfigService.testConfig(LlmConfigVO.builder()
+                .provider("openai")
+                .engine("qoder")
+                .apiBase("https://api.openai.com/v1";)
+                .model("qoder-model")
+                .maxTokens(2048)
+                .temperature(1.0)
+                .build());
+
+        assertThat(result.getStatus()).isEqualTo(1);
+        assertThat(result.getCode()).isEqualTo("llm.provider.cli_missing");
+        assertThat(result.getErrMsg()).contains("qoder");
+        assertThat(result.getHint()).contains("Install").contains("HTTP 
engine");
+        verifyNoInteractions(llmClient);
+    }
+
     @Test
     void testConfigShouldRejectInvalidApiBase() {
         LlmOperationResultVO result = 
llmConfigService.testConfig(LlmConfigVO.builder()

Reply via email to