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;
             }

Reply via email to