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