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 5a390ee56 fix(ai): keep a run whose worker died with an Error out of
COMPLETED (#5118)
5a390ee56 is described below
commit 5a390ee56fd6cb35b6faf057941d3f6b07c4f2a8
Author: 烤化の初雪 <[email protected]>
AuthorDate: Thu Oct 1 17:30:16 2026 +0800
fix(ai): keep a run whose worker died with an Error out of COMPLETED (#5118)
execute() catches LlmGatewayException and RuntimeException and then decides
the
terminal state in finally(). An Error (an OOM, or a stack overflow while
decoding a deeply nested tool payload) matches neither catch, so the outcome
was empty, decide() fell through to its last branch and the run row was
written
COMPLETED with no answer — the transcript then replays it as a finished run.
Catch the Error as well: record it in the same outcome field as an
unexpected
RuntimeException, write the FAILED / PROVIDER_ERROR / ai.run.internal_error
terminal row, and let the Error propagate to the worker's uncaught handler.
Test:
AiRunExecutorTest#anErrorFromTheWorkerShouldNotBeRecordedAsACompletedRunTest
(before the change the row reads COMPLETED; the sibling RuntimeException
case was
already covered).
Co-authored-by: unbridled-41
<[email protected]>
---
.../studio/ops/ai/conversation/AiRunExecutor.java | 10 +++++++++-
.../studio/ops/ai/conversation/AiRunExecutorTest.java | 16 ++++++++++++++++
.../studio/ops/ai/conversation/AiRunTestSupport.java | 8 ++++++++
3 files changed, 33 insertions(+), 1 deletion(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/conversation/AiRunExecutor.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/conversation/AiRunExecutor.java
index 756d51895..7d6b46839 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/conversation/AiRunExecutor.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/conversation/AiRunExecutor.java
@@ -283,6 +283,14 @@ public class AiRunExecutor {
} catch (RuntimeException exception) {
log.error("agent run {} failed unexpectedly", run.getId(),
exception);
outcome.unexpected = exception;
+ } catch (Error error) {
+ // A RuntimeException is a failure of the run; an Error (an OOM or
a stack overflow while
+ // decoding a tool payload, say) is a failure of the worker. It
still has to reach the
+ // terminal write below — without this branch it skipped both
catches, and decide() read an
+ // empty outcome and marked a crashed run COMPLETED.
+ log.error("agent run {} died with an error", run.getId(), error);
+ outcome.unexpected = error;
+ throw error;
} finally {
try {
finalizeRun(context, decide(context, outcome));
@@ -724,7 +732,7 @@ public class AiRunExecutor {
private boolean successTerminalProjected;
private boolean stopRacedSuccess;
private LlmGatewayException gatewayFailure;
- private RuntimeException unexpected;
+ private Throwable unexpected;
/**
* Drops everything the attempt that asked for the missing session
reported, so the retry's own
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/ai/conversation/AiRunExecutorTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/conversation/AiRunExecutorTest.java
index 31e1f87c0..ff47659a8 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/ops/ai/conversation/AiRunExecutorTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/conversation/AiRunExecutorTest.java
@@ -360,6 +360,22 @@ class AiRunExecutorTest {
assertThat(emitters.get(0).completed()).isTrue();
}
+ @Test
+ void anErrorFromTheWorkerShouldNotBeRecordedAsACompletedRunTest() {
+ provider.emit(new AgentEvent.TextDelta("half an answer"));
+ provider.error = new StackOverflowError("nested tool output");
+
+ startAndRun();
+
+ // The run must not read as a finished answer just because the failure
was an Error: the
+ // executor has to write the same terminal state it writes for a
RuntimeException.
+ assertThat(runRow().getStatus()).isEqualTo(RunStatus.FAILED.name());
+
assertThat(runRow().getStopReason()).isEqualTo(StopReason.PROVIDER_ERROR.name());
+
assertThat(runRow().getErrorCode()).isEqualTo(AiRunExecutor.ERROR_CODE_INTERNAL);
+ assertThat(runRow().getFinishedAt()).isNotNull();
+ assertThat(emitters.get(0).completed()).isTrue();
+ }
+
@Test
void finalisationShouldPersistInTheFixedOrderTest() {
// A provider failure, so the terminal event is the one the EXECUTOR
writes: the projector only
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/ai/conversation/AiRunTestSupport.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/conversation/AiRunTestSupport.java
index 5d5f2e36c..c01f9495e 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/ops/ai/conversation/AiRunTestSupport.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/conversation/AiRunTestSupport.java
@@ -275,6 +275,11 @@ final class AiRunTestSupport {
Consumer<AgentStreamOptions> beforeStream = options -> {
};
RuntimeException failure;
+ /**
+ * Thrown instead of any RuntimeException, so a test can reach the
paths that only an {@link Error}
+ * takes — the executor has to write the terminal row for those too.
+ */
+ Error error;
String lastPrompt;
AgentStreamOptions lastOptions;
int calls;
@@ -309,6 +314,9 @@ final class AiRunTestSupport {
lastOptions = options;
beforeStream.accept(options);
scripted.forEach(sink);
+ if (error != null) {
+ throw error;
+ }
if (failure != null) {
throw failure;
}