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