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 836944e2f fix(group): report cloud queue offsets as unknown instead of 
zero (#4902)
836944e2f is described below

commit 836944e2f280cb48d9a8176ef31768bee637efe3
Author: 烤化の初雪 <[email protected]>
AuthorDate: Thu Sep 24 18:21:31 2026 +0800

    fix(group): report cloud queue offsets as unknown instead of zero (#4902)
    
    fix(group): report cloud queue offsets as unknown instead of zero
    
    The Tencent and Aliyun lag APIs report the lag per topic and expose no
    per-queue offsets, but both providers filled brokerOffset and consumerOffset
    with 0 for every row they build (TencentInstanceProvider#getGroupProgress,
    AliyunConverters#toQueueProgressRows). The streamed offsets therefore read 
as
    real measurements in the 消费进度 table: a row showed "Broker Offset 0" and
    "Consumer Offset 0" next to a real 堆积量, and the reset-offset preview took
    the same zero as the current consumer offset.
    
    Report the unknown sentinel instead, the value this codebase already uses 
for
    "cannot be determined" (ConsumerLagResolver.UNKNOWN, the -1 the broker sends
    for a gRPC lag, onlineInstances, the reset preview's own offsets), and 
render a
    negative offset as unavailable in the progress table, which is what
    formatOffsetValue already does for the reset preview.
    
    QueueProgressVO gains UNKNOWN_OFFSET so the providers and the console share 
one
    name for the sentinel; the AI group-detail contract keeps receiving 
numbers, so
    its required schema is unaffected.
    
    Regression tests: two Aliyun converter cases and one Tencent provider case 
fail
    with "expected: -1L but was: 0L" before the change, and a ConsumerPage case
    pins that such an offset renders as unavailable instead of a number.
---
 .../studio/instance/group/QueueProgressVO.java     |  9 ++++++
 .../studio/provider/alibaba/AliyunConverters.java  | 10 ++++---
 .../provider/tencent/TencentInstanceProvider.java  |  7 +++--
 .../provider/alibaba/AliyunConvertersLagTest.java  | 30 ++++++++++++++++++++
 .../tencent/TencentInstanceProviderTest.java       | 19 +++++++++++++
 .../pages/instance/__tests__/ConsumerPage.test.tsx | 33 ++++++++++++++++++++++
 web/src/pages/instance/consumer.tsx                |  7 +++--
 7 files changed, 107 insertions(+), 8 deletions(-)

diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/instance/group/QueueProgressVO.java
 
b/server/src/main/java/org/apache/rocketmq/studio/instance/group/QueueProgressVO.java
index d7149b6df..9fc0599f5 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/instance/group/QueueProgressVO.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/instance/group/QueueProgressVO.java
@@ -26,6 +26,15 @@ import lombok.NoArgsConstructor;
 @NoArgsConstructor
 @AllArgsConstructor
 public class QueueProgressVO {
+
+    /**
+     * Sentinel for an offset the provider cannot report, e.g. the per-topic 
lag rows of the cloud
+     * providers, which carry no per-queue offsets at all. It matches the 
{@code -1} the broker uses
+     * for an undeterminable lag, and the console renders a negative offset as 
unavailable instead
+     * of a number that would read like a measurement.
+     */
+    public static final long UNKNOWN_OFFSET = -1L;
+
     private String topic;
     private String broker;
     private int queueId;
diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConverters.java
 
b/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConverters.java
index 8209b18fc..0ccf5e2c7 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConverters.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConverters.java
@@ -195,12 +195,14 @@ final class AliyunConverters {
             for (Map.Entry<String, DataTopicLagMapValue> entry : 
topicLagMap.entrySet()) {
                 long ready = entry.getValue() == null || 
entry.getValue().getReadyCount() == null
                         ? 0L : entry.getValue().getReadyCount();
+                // The Aliyun API reports the lag per topic, so the row 
carries no queue offsets;
+                // report the unknown sentinel instead of a zero that reads 
like a measurement.
                 rows.add(QueueProgressVO.builder()
                         .topic(entry.getKey())
                         .broker("topic:" + entry.getKey())
                         .queueId(0)
-                        .brokerOffset(0L)
-                        .consumerOffset(0L)
+                        .brokerOffset(QueueProgressVO.UNKNOWN_OFFSET)
+                        .consumerOffset(QueueProgressVO.UNKNOWN_OFFSET)
                         .diffTotal(ready)
                         .build());
             }
@@ -213,8 +215,8 @@ final class AliyunConverters {
             rows.add(QueueProgressVO.builder()
                     .broker("total")
                     .queueId(0)
-                    .brokerOffset(0L)
-                    .consumerOffset(0L)
+                    .brokerOffset(QueueProgressVO.UNKNOWN_OFFSET)
+                    .consumerOffset(QueueProgressVO.UNKNOWN_OFFSET)
                     .diffTotal(totalLag.getReadyCount())
                     .build());
         }
diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProvider.java
 
b/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProvider.java
index 95f2990aa..93a6f0752 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProvider.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProvider.java
@@ -535,8 +535,11 @@ public class TencentInstanceProvider implements 
InstanceProvider {
                     .topic(subscription.getTopic())
                     .broker("topic:" + subscription.getTopic())
                     .queueId(0)
-                    .brokerOffset(0L)
-                    .consumerOffset(0L)
+                    // The Tencent API reports the lag per topic, so this row 
carries no queue
+                    // offsets; report the unknown sentinel instead of a zero 
that the console
+                    // would render as a real measurement next to the real lag.
+                    .brokerOffset(QueueProgressVO.UNKNOWN_OFFSET)
+                    .consumerOffset(QueueProgressVO.UNKNOWN_OFFSET)
                     .diffTotal(subscription.getConsumerLag() == null ? 0L : 
subscription.getConsumerLag())
                     .build());
         }
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConvertersLagTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConvertersLagTest.java
index 9fd14f418..c9d980cfd 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConvertersLagTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConvertersLagTest.java
@@ -57,4 +57,34 @@ class AliyunConvertersLagTest {
                     assertThat(row.getDiffTotal()).isEqualTo(100L);
                 });
     }
+
+    @Test
+    void queueOffsetsShouldBeUnknownWhenOnlyTheTopicLagIsKnownTest() {
+        GetConsumerGroupLagResponseBody.Data data = 
GetConsumerGroupLagResponseBody.Data.builder()
+                .topicLagMap(Map.of("orders", 
DataTopicLagMapValue.builder().readyCount(40L).build()))
+                .build();
+
+        // The Aliyun API reports the lag per topic and no per-queue offsets, 
so the row must not
+        // claim offsets of zero next to the real lag.
+        assertThat(AliyunConverters.toQueueProgressRows(data)).singleElement()
+                .satisfies(row -> {
+                    
assertThat(row.getBrokerOffset()).isEqualTo(QueueProgressVO.UNKNOWN_OFFSET);
+                    
assertThat(row.getConsumerOffset()).isEqualTo(QueueProgressVO.UNKNOWN_OFFSET);
+                    assertThat(row.getDiffTotal()).isEqualTo(40L);
+                });
+    }
+
+    @Test
+    void aggregateRowOffsetsShouldBeUnknownTooTest() {
+        GetConsumerGroupLagResponseBody.Data data = 
GetConsumerGroupLagResponseBody.Data.builder()
+                
.totalLag(GetConsumerGroupLagResponseBody.TotalLag.builder().readyCount(100L).build())
+                .build();
+
+        assertThat(AliyunConverters.toQueueProgressRows(data)).singleElement()
+                .satisfies(row -> {
+                    assertThat(row.getBroker()).isEqualTo("total");
+                    
assertThat(row.getBrokerOffset()).isEqualTo(QueueProgressVO.UNKNOWN_OFFSET);
+                    
assertThat(row.getConsumerOffset()).isEqualTo(QueueProgressVO.UNKNOWN_OFFSET);
+                });
+    }
 }
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProviderTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProviderTest.java
index 67371481a..91393f953 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProviderTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProviderTest.java
@@ -730,6 +730,25 @@ class TencentInstanceProviderTest {
         assertThat(subscriptions.get(1).getConsistency()).isNull();
     }
 
+    @Test
+    void getGroupProgressShouldReportUnknownQueueOffsetsTest() throws 
Exception {
+        SubscriptionData subscription = new SubscriptionData();
+        subscription.setTopic("orders");
+        subscription.setConsumerLag(42L);
+        DescribeTopicListByGroupResponse response = new 
DescribeTopicListByGroupResponse();
+        response.setData(new SubscriptionData[]{subscription});
+        when(client.DescribeTopicListByGroup(any())).thenReturn(response);
+
+        // The Tencent API exposes the lag per topic and no per-queue offsets, 
so the row must not
+        // claim offsets of zero next to the real lag.
+        assertThat(provider.getGroupProgress(STUDIO_INSTANCE_ID, 
"GID_test")).singleElement()
+                .satisfies(row -> {
+                    
assertThat(row.getBrokerOffset()).isEqualTo(QueueProgressVO.UNKNOWN_OFFSET);
+                    
assertThat(row.getConsumerOffset()).isEqualTo(QueueProgressVO.UNKNOWN_OFFSET);
+                    assertThat(row.getDiffTotal()).isEqualTo(42L);
+                });
+    }
+
     @Test
     void previewResetOffsetShouldAllowLimitedCloudPreviewTest() throws 
Exception {
         SubscriptionData subscription = new SubscriptionData();
diff --git a/web/src/pages/instance/__tests__/ConsumerPage.test.tsx 
b/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
index 4e371a1a9..f5dd217c0 100644
--- a/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
+++ b/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
@@ -703,6 +703,39 @@ describe('Consumer page', () => {
     await waitFor(() => 
expect(within(panel).queryByText(/消费进度加载失败/)).toBeInTheDocument());
   });
 
+  it('renders an offset the provider cannot report as unavailable', async () 
=> {
+    vi.mocked(consumerService.getConsumerProgress).mockResolvedValue([
+      {
+        topic: 'remote-topic',
+        // A cloud provider reports the lag per topic and no per-queue 
offsets, so it sends the
+        // negative sentinel the backend uses for a value it cannot determine.
+        broker: 'topic:remote-topic',
+        queueId: 0,
+        brokerOffset: -1,
+        consumerOffset: -1,
+        diffTotal: 42,
+      },
+    ]);
+    const user = userEvent.setup({ pointerEventsCheck: 0 });
+    renderWithProviders(<ConsumerPage />);
+
+    await user.click(await screen.findByRole('button', { name: /详情/ }));
+    await user.click(await screen.findByRole('tab', { name: /消费进度/ }));
+    const progressPanel = await screen.findByRole('tabpanel', { name: /消费进度/ 
});
+    await waitFor(() =>
+      
expect(within(progressPanel).getByText('remote-topic')).toBeInTheDocument(),
+    );
+
+    const row = within(progressPanel)
+      .getAllByRole('row')
+      .find((candidate) => within(candidate).queryByText('remote-topic'));
+    expect(row).toBeDefined();
+    const cells = within(row!)
+      .getAllByRole('cell')
+      .map((cell) => cell.textContent?.trim());
+    expect(cells.filter((cell) => cell === '-')).toHaveLength(2);
+  });
+
   it('shows group health diagnostics from subscriptions, progress and 
clients', async () => {
     const riskyGroup: ConsumerGroup = {
       ...group,
diff --git a/web/src/pages/instance/consumer.tsx 
b/web/src/pages/instance/consumer.tsx
index 4e2036d91..cd0ce23b5 100644
--- a/web/src/pages/instance/consumer.tsx
+++ b/web/src/pages/instance/consumer.tsx
@@ -1268,8 +1268,11 @@ const ConsumerPageContent = ({
       key: 'brokerOffset',
       width: 140,
       align: 'right',
+      // Cloud providers report the lag per topic and cannot supply per-queue 
offsets, so they
+      // send the negative sentinel: show it as unavailable instead of a 
number that would read
+      // like a measurement next to the real lag.
       render: (offset: number) => (
-        <Text style={{ fontFamily: 'monospace' 
}}>{offset.toLocaleString()}</Text>
+        <Text style={{ fontFamily: 'monospace' 
}}>{formatOffsetValue(offset)}</Text>
       ),
     },
     {
@@ -1279,7 +1282,7 @@ const ConsumerPageContent = ({
       width: 150,
       align: 'right',
       render: (offset: number) => (
-        <Text style={{ fontFamily: 'monospace' 
}}>{offset.toLocaleString()}</Text>
+        <Text style={{ fontFamily: 'monospace' 
}}>{formatOffsetValue(offset)}</Text>
       ),
     },
     {

Reply via email to