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 3357bf0e fix: scope consumer group diagnostics by selected instance
(#1248)
3357bf0e is described below
commit 3357bf0e70b29541cc133b2b2bfaea6f212d5e8e
Author: aias00 <[email protected]>
AuthorDate: Mon Aug 10 14:17:13 2026 +0800
fix: scope consumer group diagnostics by selected instance (#1248)
* [ISSUE #1247] Scope Consumer Group list enrichment by instance
* [ISSUE #1249] Surface Topic consumer diagnostics failures
* [ISSUE #1251] Surface Consumer Group subscription failures
* fix: scope consumer diagnostics caches by instance
Signed-off-by: liuhy <[email protected]>
* fix: stop resource loading without an instance
Signed-off-by: liuhy <[email protected]>
* fix: preserve selected instance for group details
---------
Signed-off-by: liuhy <[email protected]>
---
.../provider/apache/ApacheInstanceProvider.java | 4 +-
.../studio/provider/apache/MetadataProvider.java | 5 ++
.../provider/apache/RocketMQMetadataProvider.java | 31 +++++++++---
.../apache/ApacheInstanceProviderTest.java | 10 ++++
.../apache/RocketMQMetadataProviderTest.java | 39 +++++++++++++++
.../pages/instance/__tests__/ConsumerPage.test.tsx | 53 +++++++++++++++++++++
.../pages/instance/__tests__/TopicPage.test.tsx | 2 +
web/src/pages/instance/consumer.tsx | 55 +++++++++++++++-------
web/src/pages/instance/topic.tsx | 4 ++
web/src/services/consumerService.test.ts | 25 +++++++++-
web/src/services/consumerService.ts | 7 ++-
11 files changed, 205 insertions(+), 30 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProvider.java
index eb1faa50..045f1564 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProvider.java
@@ -90,9 +90,7 @@ public class ApacheInstanceProvider implements
InstanceProvider {
@Override
public List<ConsumerGroupVO> listConsumerGroups(String instanceId, String
search) {
- return metadataProvider.listConsumerGroups(null, search).stream()
- .filter(group -> matchesInstance(group.getInstanceId(),
instanceId))
- .toList();
+ return metadataProvider.listConsumerGroups(instanceId, null, search);
}
@Override
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/MetadataProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/MetadataProvider.java
index 9acfe40e..7d243478 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/MetadataProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/MetadataProvider.java
@@ -28,6 +28,11 @@ import java.util.List;
public interface MetadataProvider {
List<TopicVO> listTopics(String clusterId, String type, String search);
List<ConsumerGroupVO> listConsumerGroups(String clusterId, String search);
+
+ default List<ConsumerGroupVO> listConsumerGroups(String instanceId, String
clusterId, String search) {
+ return listConsumerGroups(clusterId, search);
+ }
+
List<BrokerRouteVO> getTopicRoutes(String instanceId, String name);
List<TopicConsumerVO> getTopicConsumers(String instanceId, String name);
List<QueueProgressVO> getGroupProgress(String instanceId, String name);
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
index f63a1128..992644d4 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
@@ -156,7 +156,13 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
@Override
public List<ConsumerGroupVO> listConsumerGroups(String clusterId, String
search) {
+ return listConsumerGroups(null, clusterId, search);
+ }
+
+ @Override
+ public List<ConsumerGroupVO> listConsumerGroups(String instanceId, String
clusterId, String search) {
LambdaQueryWrapper<RmqGroup> query = new LambdaQueryWrapper<RmqGroup>()
+ .eq(StringUtils.hasText(instanceId), RmqGroup::getInstanceId,
instanceId)
.eq(StringUtils.hasText(clusterId), RmqGroup::getClusterId,
clusterId)
.like(StringUtils.hasText(search), RmqGroup::getName, search)
.orderByAsc(RmqGroup::getName);
@@ -175,8 +181,10 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
vo.setCreatedAt(entity.getCreatedAt());
vo.setUpdatedAt(entity.getUpdatedAt());
- if (hasAdmin()) {
- enrichGroupWithConnectionInfo(vo, entity.getName());
+ if (StringUtils.hasText(instanceId)) {
+ enrichGroupWithConnectionInfo(vo, entity.getName(),
instanceId);
+ } else if (hasAdmin()) {
+ enrichGroupWithConnectionInfo(vo, entity.getName(), null);
}
result.add(vo);
}
@@ -321,7 +329,7 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
return consumers;
} catch (Exception e) {
log.warn("Failed to get consumers for topic {}: {}", name,
e.getMessage());
- return Collections.emptyList();
+ throw new BusinessException(502, "Failed to get consumers for
topic " + name + ": " + e.getMessage());
}
}
@@ -412,15 +420,17 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
return subscriptions;
} catch (Exception e) {
log.warn("Failed to get subscriptions for group {}: {}", name,
e.getMessage());
- return Collections.emptyList();
+ throw new BusinessException(502,
+ "Failed to get subscriptions for group " + name + ": " +
e.getMessage());
}
}
// ── Helper methods ──────────────────────────────────────────────────
- private void enrichGroupWithConnectionInfo(ConsumerGroupVO vo, String
groupName) {
+ private void enrichGroupWithConnectionInfo(ConsumerGroupVO vo, String
groupName, String instanceId) {
try {
- ConsumerConnection conn = adminExecute(admin ->
admin.examineConsumerConnectionInfo(groupName));
+ ConsumerConnection conn = executeForInstance(instanceId,
+ admin -> admin.examineConsumerConnectionInfo(groupName));
if (conn != null) {
if (conn.getConnectionSet() != null) {
vo.setOnlineInstances(conn.getConnectionSet().size());
@@ -435,7 +445,7 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
// Try to get lag info
try {
- ConsumeStats stats = adminExecute(admin ->
admin.examineConsumeStats(groupName));
+ ConsumeStats stats = executeForInstance(instanceId, admin ->
admin.examineConsumeStats(groupName));
if (stats != null && stats.getOffsetTable() != null) {
long totalLag = 0;
for (OffsetWrapper ow : stats.getOffsetTable().values()) {
@@ -448,6 +458,13 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
}
}
+ private <T> T executeForInstance(String instanceId,
MqAdminExtFactory.AdminAction<T> action) {
+ if (StringUtils.hasText(instanceId)) {
+ return runtimeAdminClientResolver.execute(instanceId, action);
+ }
+ return adminExecute(action);
+ }
+
private String filterMode(String expressionType) {
if ("SQL92".equals(expressionType)) {
return "SQL";
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProviderTest.java
index 62db28a8..c0fa9b78 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProviderTest.java
@@ -26,6 +26,7 @@ import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import static org.assertj.core.api.Assertions.assertThat;
+import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@ExtendWith(MockitoExtension.class)
@@ -64,4 +65,13 @@ class ApacheInstanceProviderTest {
assertThat(provider.countGroups("inst-1")).isEqualTo(2);
}
+
+ @Test
+ void listConsumerGroupsShouldPassTheSelectedInstanceToMetadataProvider() {
+ when(metadataProvider.listConsumerGroups("inst-1", null,
"orders")).thenReturn(java.util.List.of());
+
+ assertThat(provider.listConsumerGroups("inst-1", "orders")).isEmpty();
+
+ verify(metadataProvider).listConsumerGroups("inst-1", null, "orders");
+ }
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProviderTest.java
index 19459427..41def9bc 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProviderTest.java
@@ -44,6 +44,7 @@ import static org.mockito.Mockito.lenient;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@@ -102,6 +103,21 @@ class RocketMQMetadataProviderTest {
assertThat(groups.get(0).getConsumeType()).isEqualTo(ConsumeType.CLUSTERING);
}
+ @Test
+ void listConsumerGroupsShouldUseSelectedInstanceForRuntimeEnrichment() {
+ RmqGroup entity = new RmqGroup();
+ entity.setName("group-a");
+ entity.setInstanceId("instance-a");
+ entity.setClusterId("cluster-a");
+ when(groupMapper.selectList(any())).thenReturn(List.of(entity));
+ when(runtimeAdminClientResolver.execute(eq("instance-a"),
any())).thenReturn(null);
+
+ List<ConsumerGroupVO> groups =
newProvider().listConsumerGroups("instance-a", null, null);
+
+
assertThat(groups).singleElement().extracting(ConsumerGroupVO::getName).isEqualTo("group-a");
+ verify(runtimeAdminClientResolver, times(2)).execute(eq("instance-a"),
any());
+ }
+
@Test
void getTopicRoutesShouldUseSelectedInstanceRuntimeClient() {
List<BrokerRouteVO> routes =
List.of(BrokerRouteVO.builder().brokerName("broker-a").build());
@@ -147,6 +163,17 @@ class RocketMQMetadataProviderTest {
.satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(502));
}
+ @Test
+ void getTopicConsumersSurfacesAdminFailure() throws Exception {
+ DefaultMQAdminExt admin =
org.mockito.Mockito.mock(DefaultMQAdminExt.class);
+ when(admin.queryTopicConsumeByWho("TopicA")).thenThrow(new
IllegalStateException("broker unavailable"));
+
+ assertThatThrownBy(() ->
newLiveProvider(admin).getTopicConsumers(null, "TopicA"))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Failed to get consumers for topic TopicA: broker
unavailable")
+ .satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(502));
+ }
+
@Test
void getGroupProgressSurfacesAdminFailure() throws Exception {
DefaultMQAdminExt admin =
org.mockito.Mockito.mock(DefaultMQAdminExt.class);
@@ -158,6 +185,18 @@ class RocketMQMetadataProviderTest {
.satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(502));
}
+ @Test
+ void getGroupSubscriptionsSurfacesAdminFailure() throws Exception {
+ DefaultMQAdminExt admin =
org.mockito.Mockito.mock(DefaultMQAdminExt.class);
+ when(admin.examineConsumerConnectionInfo("group-a"))
+ .thenThrow(new IllegalStateException("broker unavailable"));
+
+ assertThatThrownBy(() ->
newLiveProvider(admin).getGroupSubscriptions(null, "group-a"))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Failed to get subscriptions for group group-a:
broker unavailable")
+ .satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(502));
+ }
+
private RocketMQMetadataProvider newLiveProvider(MQAdminExt admin) throws
Exception {
MqAdminExtFactory factory = mock(MqAdminExtFactory.class);
RocketMQProperties liveProperties = new RocketMQProperties();
diff --git a/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
b/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
index 4bd93b8f..66605e8b 100644
--- a/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
+++ b/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
@@ -262,6 +262,57 @@ describe('Consumer page', () => {
);
});
+ it('reloads same-named group diagnostics after changing the selected
instance', async () => {
+ vi.mocked(instanceService.listInstances).mockResolvedValue([
+ {
+ id: 'instance-a',
+ name: 'Instance A',
+ remark: '',
+ type: 'DIRECT',
+ endpoint: '127.0.0.1:9876',
+ topicCount: 0,
+ consumerGroupCount: 0,
+ createdAt: '2026-07-23T00:00:00Z',
+ updatedAt: '2026-07-23T00:00:00Z',
+ },
+ {
+ id: 'instance-b',
+ name: 'Instance B',
+ remark: '',
+ type: 'DIRECT',
+ endpoint: '127.0.0.2:9876',
+ topicCount: 0,
+ consumerGroupCount: 0,
+ createdAt: '2026-07-23T00:00:00Z',
+ updatedAt: '2026-07-23T00:00:00Z',
+ },
+ ]);
+ vi.mocked(consumerService.listConsumerGroups).mockImplementation(async
(params) => [
+ { ...group, instanceId: params?.instanceId ?? 'instance-a' },
+ ]);
+ const user = userEvent.setup();
+ renderWithProviders(<ConsumerPage />, '/instance/instance-a/consumer');
+
+ await user.click(await screen.findByRole('button', { name: /详情/ }));
+ await waitFor(() =>
+
expect(consumerService.getConsumerSubscriptions).toHaveBeenCalledWith('remote-cg',
'instance-a'),
+ );
+
+ await user.click(screen.getByText('Instance A'));
+ await user.click(await screen.findByText('Instance B'));
+ await waitFor(() =>
+ expect(consumerService.listConsumerGroups).toHaveBeenCalledWith({
instanceId: 'instance-b' }),
+ );
+ await user.click(await screen.findByRole('button', { name: /详情/ }));
+
+ await waitFor(() =>
+
expect(consumerService.getConsumerSubscriptions).toHaveBeenCalledWith('remote-cg',
'instance-b'),
+ );
+ await waitFor(() =>
+
expect(consumerService.getConsumerProgress).toHaveBeenCalledWith('remote-cg',
'instance-b'),
+ );
+ });
+
it('highlights inconsistent subscriptions and refreshes the check result',
async () => {
vi.mocked(consumerService.getConsumerSubscriptions)
.mockResolvedValueOnce([
@@ -421,6 +472,8 @@ describe('Consumer page', () => {
renderWithProviders(<ConsumerPage />);
expect(await screen.findByText('选择实例')).toBeInTheDocument();
+ expect(consumerService.listConsumerGroups).not.toHaveBeenCalled();
+ expect(document.querySelector('.ant-spin-spinning')).toBeNull();
expect(screen.getByRole('button', { name: /导入/ })).toBeDisabled();
expect(screen.getByRole('button', { name: '创建 Group' })).toBeDisabled();
});
diff --git a/web/src/pages/instance/__tests__/TopicPage.test.tsx
b/web/src/pages/instance/__tests__/TopicPage.test.tsx
index 8713a88c..fde8eaa9 100644
--- a/web/src/pages/instance/__tests__/TopicPage.test.tsx
+++ b/web/src/pages/instance/__tests__/TopicPage.test.tsx
@@ -391,6 +391,8 @@ describe('TopicPage', () => {
renderWithProviders();
expect(await screen.findByText('共 0 个 Topic')).toBeInTheDocument();
+ expect(topicServiceMocks.listTopics).not.toHaveBeenCalled();
+ expect(document.querySelector('.ant-spin-spinning')).toBeNull();
expect(screen.getByRole('button', { name: /导入/ })).toBeDisabled();
expect(screen.getByRole('button', { name: /创建 Topic/ })).toBeDisabled();
});
diff --git a/web/src/pages/instance/consumer.tsx
b/web/src/pages/instance/consumer.tsx
index b6c2b117..a11b4cba 100644
--- a/web/src/pages/instance/consumer.tsx
+++ b/web/src/pages/instance/consumer.tsx
@@ -176,6 +176,9 @@ const isConsistentSubscription = (subscription:
SubscriptionEntry): boolean =>
const isInconsistentSubscription = (subscription: SubscriptionEntry): boolean
=>
isInconsistentValue(subscription.consistency);
+export const diagnosticCacheKey = (instanceId: string, groupName: string) =>
+ `${instanceId}\u0000${groupName}`;
+
/* ═══════════════════════════════════════════
ConsumerPage
═══════════════════════════════════════════ */
@@ -222,8 +225,19 @@ const ConsumerPage = () => {
const groupRequestIdRef = useRef(0);
+ useEffect(() => {
+ setSelectedGroup(null);
+ setModalOpen(false);
+ setResetGroup(null);
+ setResetModalOpen(false);
+ }, [selectedInstanceId]);
+
useEffect(() => {
if (!selectedInstanceId) {
+ groupRequestIdRef.current += 1;
+ setGroups([]);
+ setSelectedRowKeys([]);
+ setLoading(false);
return;
}
const requestId = ++groupRequestIdRef.current;
@@ -247,20 +261,21 @@ const ConsumerPage = () => {
const loadSubscriptions = useCallback(
async (groupName: string, force = false) => {
- if (!force && subscriptionsByGroup[groupName]) return;
- setSubscriptionLoadingByGroup((prev) => ({ ...prev, [groupName]: true
}));
- setSubscriptionErrorByGroup((prev) => ({ ...prev, [groupName]: false }));
+ const cacheKey = diagnosticCacheKey(selectedInstanceId, groupName);
+ if (!force && subscriptionsByGroup[cacheKey]) return;
+ setSubscriptionLoadingByGroup((prev) => ({ ...prev, [cacheKey]: true }));
+ setSubscriptionErrorByGroup((prev) => ({ ...prev, [cacheKey]: false }));
try {
const subscriptions = await getConsumerSubscriptions(
groupName,
selectedInstanceId || undefined,
);
- setSubscriptionsByGroup((prev) => ({ ...prev, [groupName]:
subscriptions }));
+ setSubscriptionsByGroup((prev) => ({ ...prev, [cacheKey]:
subscriptions }));
} catch {
- setSubscriptionErrorByGroup((prev) => ({ ...prev, [groupName]: true
}));
+ setSubscriptionErrorByGroup((prev) => ({ ...prev, [cacheKey]: true }));
message.error(t('consumer.fetchSubscriptionsFailed', { name: groupName
}));
} finally {
- setSubscriptionLoadingByGroup((prev) => ({ ...prev, [groupName]: false
}));
+ setSubscriptionLoadingByGroup((prev) => ({ ...prev, [cacheKey]: false
}));
}
},
[subscriptionsByGroup, t, selectedInstanceId],
@@ -268,10 +283,11 @@ const ConsumerPage = () => {
const loadProgress = useCallback(
async (groupName: string) => {
- if (progressByGroup[groupName]) return;
+ const cacheKey = diagnosticCacheKey(selectedInstanceId, groupName);
+ if (progressByGroup[cacheKey]) return;
try {
const progress = await getConsumerProgress(groupName,
selectedInstanceId || undefined);
- setProgressByGroup((prev) => ({ ...prev, [groupName]: progress }));
+ setProgressByGroup((prev) => ({ ...prev, [cacheKey]: progress }));
} catch {
message.error(t('consumer.fetchProgressFailed', { name: groupName }));
}
@@ -311,8 +327,11 @@ const ConsumerPage = () => {
void loadProgress(group.name);
};
+ const selectedDiagnosticKey = selectedGroup
+ ? diagnosticCacheKey(selectedInstanceId, selectedGroup.name)
+ : '';
const selectedSubscriptions = selectedGroup
- ? (subscriptionsByGroup[selectedGroup.name] ?? [])
+ ? (subscriptionsByGroup[selectedDiagnosticKey] ?? [])
: [];
const inconsistentSubscriptions =
selectedSubscriptions.filter(isInconsistentSubscription);
const unknownSubscriptions = selectedSubscriptions.filter(
@@ -322,7 +341,7 @@ const ConsumerPage = () => {
const visibleSubscriptions = showOnlyInconsistent
? inconsistentSubscriptions
: selectedSubscriptions;
- const selectedProgress = selectedGroup ?
(progressByGroup[selectedGroup.name] ?? []) : [];
+ const selectedProgress = selectedGroup ?
(progressByGroup[selectedDiagnosticKey] ?? []) : [];
const handleImportFile = async (file: File) => {
if (!selectedInstanceId) {
@@ -895,9 +914,11 @@ const ConsumerPage = () => {
<div style={{ padding: '8px 0' }}>
<Table
columns={subscriptionSubColumns}
- dataSource={subscriptionsByGroup[record.name] ?? []}
+ dataSource={
+
subscriptionsByGroup[diagnosticCacheKey(selectedInstanceId, record.name)] ?? []
+ }
rowKey="topic"
- loading={subscriptionLoadingByGroup[record.name]}
+
loading={subscriptionLoadingByGroup[diagnosticCacheKey(selectedInstanceId,
record.name)]}
pagination={false}
size="small"
/>
@@ -1073,7 +1094,7 @@ const ConsumerPage = () => {
<Button
size="small"
icon={<ArrowsClockwise size={14} />}
-
loading={subscriptionLoadingByGroup[selectedGroup.name]}
+
loading={subscriptionLoadingByGroup[selectedDiagnosticKey]}
onClick={() => {
setShowOnlyInconsistent(false);
void loadSubscriptions(selectedGroup.name, true);
@@ -1085,7 +1106,7 @@ const ConsumerPage = () => {
<Alert
showIcon
type={
- subscriptionErrorByGroup[selectedGroup.name]
+ subscriptionErrorByGroup[selectedDiagnosticKey]
? 'error'
: inconsistentSubscriptions.length > 0 ||
unknownSubscriptions.length > 0
@@ -1095,9 +1116,9 @@ const ConsumerPage = () => {
: 'info'
}
message={
- subscriptionErrorByGroup[selectedGroup.name]
+ subscriptionErrorByGroup[selectedDiagnosticKey]
? '订阅一致性检查失败,当前保留上次检查结果'
- : subscriptionLoadingByGroup[selectedGroup.name] &&
+ :
subscriptionLoadingByGroup[selectedDiagnosticKey] &&
selectedSubscriptions.length === 0
? '正在检查订阅一致性'
: inconsistentSubscriptions.length > 0
@@ -1123,7 +1144,7 @@ const ConsumerPage = () => {
columns={subscriptionSubColumns}
dataSource={visibleSubscriptions}
rowKey="topic"
-
loading={subscriptionLoadingByGroup[selectedGroup.name]}
+
loading={subscriptionLoadingByGroup[selectedDiagnosticKey]}
pagination={false}
size="small"
/>
diff --git a/web/src/pages/instance/topic.tsx b/web/src/pages/instance/topic.tsx
index 19c05dd6..64f74386 100644
--- a/web/src/pages/instance/topic.tsx
+++ b/web/src/pages/instance/topic.tsx
@@ -307,6 +307,10 @@ const TopicPage = () => {
useEffect(() => {
if (!selectedInstanceId) {
+ topicRequestIdRef.current += 1;
+ setTopics([]);
+ setSelectedRowKeys([]);
+ setLoading(false);
return;
}
const requestId = ++topicRequestIdRef.current;
diff --git a/web/src/services/consumerService.test.ts
b/web/src/services/consumerService.test.ts
index 4d61b45d..27ef23b2 100644
--- a/web/src/services/consumerService.test.ts
+++ b/web/src/services/consumerService.test.ts
@@ -24,10 +24,16 @@ import {
listConsumerGroups,
} from './consumerService';
-vi.mock('./dataMode', () => ({ isMockMode: () => true }));
+const { mode, metadataApi } = vi.hoisted(() => ({
+ mode: { mock: true },
+ metadataApi: { getConsumerGroup: vi.fn() },
+}));
+
+vi.mock('./dataMode', () => ({ isMockMode: () => mode.mock }));
vi.mock('../config', () => ({
API_BASE_URL: '/api',
}));
+vi.mock('../api/metadata', () => metadataApi);
describe('consumer service mock data', () => {
it('returns copied consumer group rows', async () => {
@@ -93,4 +99,21 @@ describe('consumer service mock data', () => {
expect(detail.subscribedTopics).toEqual(['created-topic']);
expect(detail).not.toBe(created);
});
+
+ it('forwards the selected instance when loading consumer group details in
API mode', async () => {
+ mode.mock = false;
+ const detail = {
+ name: 'cg-orders',
+ subscribedTopics: [],
+ instances: [],
+ };
+ try {
+ metadataApi.getConsumerGroup.mockResolvedValue(detail);
+
+ await expect(getConsumerGroup('cg-orders',
'instance-a')).resolves.toEqual(detail);
+ expect(metadataApi.getConsumerGroup).toHaveBeenCalledWith('cg-orders',
'instance-a');
+ } finally {
+ mode.mock = true;
+ }
+ });
});
diff --git a/web/src/services/consumerService.ts
b/web/src/services/consumerService.ts
index b48f209a..dfabd684 100644
--- a/web/src/services/consumerService.ts
+++ b/web/src/services/consumerService.ts
@@ -61,13 +61,16 @@ export async function getConsumerProgress(
return metadataApi.getConsumerProgress(name, instanceId);
}
-export async function getConsumerGroup(name: string):
Promise<ConsumerGroupDetail> {
+export async function getConsumerGroup(
+ name: string,
+ instanceId?: string,
+): Promise<ConsumerGroupDetail> {
if (isMockMode()) {
const group = mockConsumerGroups.find((item) => item.name === name);
if (!group) throw new Error(`Consumer group not found: ${name}`);
return copyConsumerGroup(group as unknown as ConsumerGroupDetail) as
ConsumerGroupDetail;
}
- return metadataApi.getConsumerGroup(name);
+ return metadataApi.getConsumerGroup(name, instanceId);
}
export async function getConsumerSubscriptions(