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 b47f2b451 fix(ai): preserve Topic identity in group progress output
(#5102)
b47f2b451 is described below
commit b47f2b4517e13781ed653cceebd86b095a3f37af
Author: 风起 <[email protected]>
AuthorDate: Thu Oct 1 17:29:57 2026 +0800
fix(ai): preserve Topic identity in group progress output (#5102)
Keep queue progress attributable when a group consumes multiple topics on
the same broker and queue ID. Preserve topicless cloud aggregates and refresh
the matching output schema and generated catalog digest.
Fixes #5101
---
rmqctl/internal/catalog/catalog_gen.go | 2 +-
.../ai/tool/contract/group/GroupDetailOutput.java | 2 +
.../main/resources/tool-catalog/tools/group.yaml | 3 ++
.../group/ConsumerGroupReadToolHandlersTest.java | 56 ++++++++++++++++++++++
.../tool/service/ToolOutputSchemaContractTest.java | 5 +-
5 files changed, 65 insertions(+), 3 deletions(-)
diff --git a/rmqctl/internal/catalog/catalog_gen.go
b/rmqctl/internal/catalog/catalog_gen.go
index b8eebdcc5..df9af561a 100644
--- a/rmqctl/internal/catalog/catalog_gen.go
+++ b/rmqctl/internal/catalog/catalog_gen.go
@@ -21,7 +21,7 @@ package catalog
var defaultDocument = Document{
Version: "2.0.0",
MinimumClientVersion: "2.0.0",
- Digest:
"378a711779f7ea73703c07cfc8fa757300c8fe0c77d967970564ef409834a85d",
+ Digest:
"098b84cd35cfe2dde3b71a1d4381863a54e187517564bbe004b6ea329e8f1814",
Tools: []Tool{
{
Name: "rmq.acl.list",
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/tool/contract/group/GroupDetailOutput.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/tool/contract/group/GroupDetailOutput.java
index c54f291da..ca5216136 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/tool/contract/group/GroupDetailOutput.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/tool/contract/group/GroupDetailOutput.java
@@ -162,6 +162,7 @@ public record GroupDetailOutput(
@JsonInclude(JsonInclude.Include.NON_NULL)
public record QueueProgress(
+ String topic,
String broker,
int queueId,
long brokerOffset,
@@ -170,6 +171,7 @@ public record GroupDetailOutput(
static QueueProgress from(QueueProgressVO source) {
return new QueueProgress(
+ source.getTopic(),
source.getBroker(),
source.getQueueId(),
source.getBrokerOffset(),
diff --git a/server/src/main/resources/tool-catalog/tools/group.yaml
b/server/src/main/resources/tool-catalog/tools/group.yaml
index adbc57bc8..3ac090089 100644
--- a/server/src/main/resources/tool-catalog/tools/group.yaml
+++ b/server/src/main/resources/tool-catalog/tools/group.yaml
@@ -267,6 +267,9 @@ tools:
- lag
additionalProperties: false
properties:
+ topic:
+ type: string
+ description: Topic owning this queue; omitted for
provider-level aggregate rows.
broker:
type: string
queueId:
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/ai/tool/handler/group/ConsumerGroupReadToolHandlersTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/tool/handler/group/ConsumerGroupReadToolHandlersTest.java
index 82f2e06ea..1f7eedabd 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/ops/ai/tool/handler/group/ConsumerGroupReadToolHandlersTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/tool/handler/group/ConsumerGroupReadToolHandlersTest.java
@@ -16,6 +16,8 @@
*/
package org.apache.rocketmq.studio.ops.ai.tool.handler.group;
+import com.fasterxml.jackson.databind.JsonNode;
+import org.apache.rocketmq.studio.common.config.LegacyJackson2Config;
import org.apache.rocketmq.studio.common.domain.enums.ConsumeType;
import org.apache.rocketmq.studio.common.domain.enums.Protocol;
import org.apache.rocketmq.studio.common.domain.enums.SubscriptionMode;
@@ -152,10 +154,64 @@ class ConsumerGroupReadToolHandlersTest {
assertThat(output.progress().queues()).singleElement()
.extracting(GroupDetailOutput.QueueProgress::queueId)
.isEqualTo(0);
+ assertThat(output.progress().queues()).singleElement()
+ .extracting(GroupDetailOutput.QueueProgress::topic)
+ .isEqualTo("TopicA");
assertThat(output.progress().totalLag()).isEqualTo(12L);
assertThat(output.clients().totalClients()).isEqualTo(1);
}
+ @Test
+ void detailPreservesTopicIdentityForSharedBrokerAndQueueIdTest() {
+ when(metadataService.consumerGroupRuntimeView("instance-a",
"group-a")).thenReturn(group);
+ when(metadataService.consumerGroupConfigurations("instance-a",
"group-a")).thenReturn(List.of(group));
+ when(metadataService.getGroupSubscriptions("instance-a",
"group-a")).thenReturn(List.of());
+ when(metadataService.getGroupProgress("instance-a", "group-a"))
+ .thenReturn(List.of(
+
QueueProgressVO.builder().topic("TopicA").broker("broker-a").queueId(0)
+
.brokerOffset(20).consumerOffset(8).diffTotal(12).build(),
+
QueueProgressVO.builder().topic("TopicB").broker("broker-a").queueId(0)
+
.brokerOffset(20).consumerOffset(8).diffTotal(12).build()));
+
+ GroupDetailOutput output = new GroupDetailToolHandler(metadataService)
+ .execute(new GroupDetailInput("instance-a", "group-a", null),
context());
+ JsonNode progress = new
LegacyJackson2Config().jackson2ObjectMapper().valueToTree(output.progress());
+
+ assertThat(progress.path("queues")).hasSize(2);
+
assertThat(progress.at("/queues/0/topic").asText()).isEqualTo("TopicA");
+
assertThat(progress.at("/queues/1/topic").asText()).isEqualTo("TopicB");
+ for (JsonNode queue : progress.path("queues")) {
+ assertThat(queue.path("broker").asText()).isEqualTo("broker-a");
+ assertThat(queue.path("queueId").asInt()).isZero();
+ assertThat(queue.path("brokerOffset").asLong()).isEqualTo(20);
+ assertThat(queue.path("consumerOffset").asLong()).isEqualTo(8);
+ assertThat(queue.path("lag").asLong()).isEqualTo(12);
+ }
+ assertThat(progress.path("totalLag").asLong()).isEqualTo(24);
+ }
+
+ @Test
+ void detailPreservesUnknownOffsetsAndTopiclessAggregateProgressTest() {
+ when(metadataService.consumerGroupRuntimeView("instance-a",
"group-a")).thenReturn(group);
+ when(metadataService.consumerGroupConfigurations("instance-a",
"group-a")).thenReturn(List.of(group));
+ when(metadataService.getGroupSubscriptions("instance-a",
"group-a")).thenReturn(List.of());
+ when(metadataService.getGroupProgress("instance-a", "group-a"))
+
.thenReturn(List.of(QueueProgressVO.builder().broker("total").queueId(0)
+ .brokerOffset(QueueProgressVO.UNKNOWN_OFFSET)
+
.consumerOffset(QueueProgressVO.UNKNOWN_OFFSET).diffTotal(-1).build()));
+
+ GroupDetailOutput output = new GroupDetailToolHandler(metadataService)
+ .execute(new GroupDetailInput("instance-a", "group-a", null),
context());
+ JsonNode progress = new
LegacyJackson2Config().jackson2ObjectMapper().valueToTree(output.progress());
+
+ assertThat(progress.path("queues")).hasSize(1);
+ assertThat(progress.at("/queues/0").has("topic")).isFalse();
+
assertThat(progress.at("/queues/0/brokerOffset").asLong()).isEqualTo(-1);
+
assertThat(progress.at("/queues/0/consumerOffset").asLong()).isEqualTo(-1);
+ assertThat(progress.at("/queues/0/lag").asLong()).isEqualTo(-1);
+ assertThat(progress.path("totalLag").asLong()).isEqualTo(-1);
+ }
+
@Test
void detailMarksHealthUnknownWhenConsumerConnectionsAreUnavailableTest() {
group.setOnlineInstances(-1);
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/ai/tool/service/ToolOutputSchemaContractTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/tool/service/ToolOutputSchemaContractTest.java
index 2fe9839e0..c99683808 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/ops/ai/tool/service/ToolOutputSchemaContractTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/tool/service/ToolOutputSchemaContractTest.java
@@ -285,7 +285,7 @@ class ToolOutputSchemaContractTest {
new GroupDetailOutput.Health("HEALTHY", List.of()),
List.of(groupItem),
new GroupDetailOutput.Progress(100L, List.of(
- new GroupDetailOutput.QueueProgress("broker-a", 0,
120L, 90L, 30L))),
+ new GroupDetailOutput.QueueProgress("orders",
"broker-a", 0, 120L, 90L, 30L))),
new GroupDetailOutput.Clients(1, List.of(
new GroupDetailOutput.Client(
"client-1", "gRPC", "127.0.0.1:50000", "JAVA",
"5.0.7",
@@ -300,7 +300,8 @@ class ToolOutputSchemaContractTest {
new GroupDetailOutput.Health(
"UNKNOWN", List.of("Consumer connection
information is unavailable.")),
List.of(unknownConnectionsItem),
- null,
+ new GroupDetailOutput.Progress(-1L, List.of(
+ new GroupDetailOutput.QueueProgress(null,
"total", 0, -1L, -1L, -1L))),
null)));
samples.put("rmq.group.update", List.of(planned(),
executed(groupItem)));
samples.put("rmq.group.delete", List.of(planned(), executedVoid()));