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 3c00abe36 fix(ai): stop the run speed report from racing the terminal 
run update (#5196)
3c00abe36 is described below

commit 3c00abe36a7c2674fb2fbcba5b89e36e8aa65682
Author: Zhao Jianing <[email protected]>
AuthorDate: Thu Oct 1 18:27:24 2026 +0800

    fix(ai): stop the run speed report from racing the terminal run update 
(#5196)
    
    reportSpeed read the full rmq_ai_run row, set one field and wrote the
    whole row back. The client sends the report exactly as its stream closes,
    which also happens mid-run when the connection drops while the run keeps
    executing, so the read can land just before finalizeRun's terminal update
    and the write just after it: the stale RUNNING state overwrites the
    terminal one, resurrects the run as active, and every later turn of that
    conversation is refused with 409 until the orphan sweep reaps the row —
    besides permanently losing the duration, tokens and finished_at the
    worker had written. Write only the id and the speed column instead, the
    same partial-update shape rememberRuntimeSession already uses.
    
    Signed-off-by: zjncs <[email protected]>
---
 .../studio/ops/ai/conversation/AiRunService.java   | 17 +++++++++++++-
 .../ops/ai/conversation/AiRunServiceTest.java      | 27 ++++++++++++++++++++++
 .../ops/ai/conversation/AiRunTestSupport.java      |  1 +
 3 files changed, 44 insertions(+), 1 deletion(-)

diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/conversation/AiRunService.java
 
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/conversation/AiRunService.java
index 7993ae50f..3b5a83589 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/conversation/AiRunService.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/conversation/AiRunService.java
@@ -393,13 +393,28 @@ public class AiRunService {
      * the same deltas but not the same clock — a client estimate is what the 
user watched, so the
      * replay should show that exact number. Idempotent: a duplicate report 
simply overwrites.
      *
+     * <p>The update is partial on purpose. The report arrives exactly when a 
run is finishing — the
+     * client sends it as its stream closes, which also happens mid-run when 
the connection drops
+     * while the run keeps executing — so a full-row write read before
+     * {@link AiRunExecutor#finalizeRun} updates the row races it and writes 
the stale
+     * {@code RUNNING} state back over the terminal one. That resurrects the 
run as active: every
+     * later turn of the conversation is refused with 409 until the orphan 
sweep reaps the row, and
+     * the non-null columns the stale read still carries — {@code status}, 
{@code end_seq} and
+     * {@code gmt_modified} — are written back over the terminal ones. The 
nullable terminal facts
+     * (duration, tokens, finished_at) survive: MyBatis-Plus defaults to 
{@code FieldStrategy.NOT_NULL},
+     * so a null field never reaches the SET clause. Only the speed column is 
touched, the same
+     * partial-update shape {@code rememberRuntimeSession} uses.
+     *
      * @throws BusinessException 404 when the run does not exist or belongs to 
somebody else
      */
     public RmqAiRun reportSpeed(Long runId, double tokensPerSecond) {
         String owner = AiConversationService.currentOwner();
         RmqAiRun run = requireOwnedRun(runId, owner);
+        RmqAiRun update = new RmqAiRun();
+        update.setId(run.getId());
+        update.setTokensPerSecond(tokensPerSecond);
+        runRepository.update(update);
         run.setTokensPerSecond(tokensPerSecond);
-        runRepository.update(run);
         return run;
     }
 
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/ops/ai/conversation/AiRunServiceTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/conversation/AiRunServiceTest.java
index 51a8a1a76..0629713c2 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/ops/ai/conversation/AiRunServiceTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/conversation/AiRunServiceTest.java
@@ -350,6 +350,33 @@ class AiRunServiceTest {
         assertThat(inserted).isEmpty();
     }
 
+    @Test
+    void reportSpeedShouldTouchOnlyTheSpeedColumnTest() {
+        // The client reports the speed as its stream closes, which also 
happens mid-run when the
+        // connection drops while the run keeps executing. A full-row write 
read before the worker
+        // finalises would race finalizeRun's terminal update and write the 
stale RUNNING state back
+        // over it, resurrecting the run as active and locking the 
conversation with 409s until the
+        // orphan sweep reaps the row. The update must therefore carry only 
the id and the speed.
+        RmqAiRun running = AiRunTestSupport.run(RUN_ID, CONVERSATION_ID, 1, 
RunStatus.RUNNING);
+        when(runRepository.findById(RUN_ID)).thenReturn(Optional.of(running));
+
+        service.reportSpeed(RUN_ID, 42.5);
+
+        assertThat(runUpdates).hasSize(1);
+        RmqAiRun update = runUpdates.get(0);
+        assertThat(update.getId()).isEqualTo(RUN_ID);
+        assertThat(update.getTokensPerSecond()).isEqualTo(42.5);
+        // Everything the finalize path owns must stay null so updateById 
cannot touch those columns.
+        assertThat(update.getStatus()).isNull();
+        assertThat(update.getStopReason()).isNull();
+        assertThat(update.getFinishedAt()).isNull();
+        assertThat(update.getDurationMs()).isNull();
+        assertThat(update.getEndSeq()).isNull();
+        assertThat(update.getInputTokens()).isNull();
+        assertThat(update.getOutputTokens()).isNull();
+        assertThat(update.getGmtModified()).isNull();
+    }
+
     @Test
     void aStaleStopShouldFailClosedAndNotTouchTheNewerRunTest() {
         RmqAiRun stale = AiRunTestSupport.run(RUN_ID, CONVERSATION_ID, 1, 
RunStatus.QUEUED);
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 c01f9495e..23860b958 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
@@ -350,6 +350,7 @@ final class AiRunTestSupport {
         copy.setStopReason(run.getStopReason());
         copy.setErrorCode(run.getErrorCode());
         copy.setErrorMessage(run.getErrorMessage());
+        copy.setTokensPerSecond(run.getTokensPerSecond());
         return copy;
     }
 

Reply via email to