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 c3cbae020 fix(web): ignore stale consumer subscription responses
(#4968)
c3cbae020 is described below
commit c3cbae0206f0a8b72e0522709547634717acdd1b
Author: Yexi Xiang <[email protected]>
AuthorDate: Thu Sep 24 18:22:12 2026 +0800
fix(web): ignore stale consumer subscription responses (#4968)
Merge upstream/rocketmq-studio into fix/4961-consumer-subscription-race
fix(consumer): ignore stale subscription responses
---
.../pages/instance/__tests__/ConsumerPage.test.tsx | 54 ++++++++++++++++++++++
web/src/pages/instance/consumer.tsx | 17 +++++--
2 files changed, 66 insertions(+), 5 deletions(-)
diff --git a/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
b/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
index f5dd217c0..c8b28faa6 100644
--- a/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
+++ b/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
@@ -1327,6 +1327,60 @@ describe('Consumer page', () => {
expect(await screen.findByText('全部 2 个订阅配置一致')).toBeInTheDocument();
});
+ it('ignores stale subscription responses after a newer diagnostic request
completes', async () => {
+ const firstRequest =
+ deferred<Awaited<ReturnType<typeof
consumerService.getConsumerSubscriptions>>>();
+ const latestRequest =
+ deferred<Awaited<ReturnType<typeof
consumerService.getConsumerSubscriptions>>>();
+ vi.mocked(consumerService.getConsumerSubscriptions)
+ .mockReturnValueOnce(firstRequest.promise)
+ .mockReturnValueOnce(latestRequest.promise);
+
+ 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 healthPanel = await screen.findByRole('tabpanel', { name: /健康诊断/ });
+ await user.click(within(healthPanel).getByRole('button', { name: /重新诊断/
}));
+ expect(consumerService.getConsumerSubscriptions).toHaveBeenCalledTimes(2);
+
+ await act(async () => {
+ latestRequest.resolve([
+ {
+ topic: 'remote-topic',
+ expression: '*',
+ type: 'NORMAL',
+ filterMode: '全量',
+ consistency: 'consistent',
+ },
+ {
+ topic: 'new-topic',
+ expression: '*',
+ type: 'NORMAL',
+ filterMode: '全量',
+ consistency: 'consistent',
+ },
+ ]);
+ await latestRequest.promise;
+ });
+ expect(screen.getByText('全部 2 个订阅配置一致')).toBeInTheDocument();
+
+ await act(async () => {
+ firstRequest.resolve([
+ {
+ topic: 'stale-topic',
+ expression: 'important',
+ type: 'NORMAL',
+ filterMode: 'Tag 过滤',
+ consistency: 'inconsistent',
+ },
+ ]);
+ await firstRequest.promise;
+ });
+ expect(screen.getByText('全部 2 个订阅配置一致')).toBeInTheDocument();
+ });
+
it('keeps unknown consistency values separate from mismatches', async () => {
vi.mocked(consumerService.getConsumerSubscriptions).mockResolvedValue([
{
diff --git a/web/src/pages/instance/consumer.tsx
b/web/src/pages/instance/consumer.tsx
index cd0ce23b5..7a57c1e10 100644
--- a/web/src/pages/instance/consumer.tsx
+++ b/web/src/pages/instance/consumer.tsx
@@ -296,6 +296,7 @@ const ConsumerPageContent = ({
const groupRequestIdRef = useRef(0);
const progressRequestIdRef = useRef<Record<string, number>>({});
+ const subscriptionRequestIdRef = useRef<Record<string, number>>({});
const stackRequestIdRef = useRef(0);
const settingsRequestIdRef = useRef(0);
// Consumption switches as loaded from the broker, used to detect high-risk
changes
@@ -388,6 +389,8 @@ const ConsumerPageContent = ({
async (groupName: string, force = false, silent = false) => {
const cacheKey = diagnosticCacheKey(selectedInstanceId, groupName);
if (!force && subscriptionsByGroup[cacheKey]) return;
+ const requestId = (subscriptionRequestIdRef.current[cacheKey] ?? 0) + 1;
+ subscriptionRequestIdRef.current[cacheKey] = requestId;
if (!silent) {
setSubscriptionLoadingByGroup((prev) => ({ ...prev, [cacheKey]: true
}));
}
@@ -397,14 +400,18 @@ const ConsumerPageContent = ({
groupName,
selectedInstanceId || undefined,
);
- setSubscriptionsByGroup((prev) => ({ ...prev, [cacheKey]:
subscriptions }));
+ if (subscriptionRequestIdRef.current[cacheKey] === requestId) {
+ setSubscriptionsByGroup((prev) => ({ ...prev, [cacheKey]:
subscriptions }));
+ }
} catch {
- setSubscriptionErrorByGroup((prev) => ({ ...prev, [cacheKey]: true }));
- if (!silent) {
- message.error(t('consumer.fetchSubscriptionsFailed', { name:
groupName }));
+ if (subscriptionRequestIdRef.current[cacheKey] === requestId) {
+ setSubscriptionErrorByGroup((prev) => ({ ...prev, [cacheKey]: true
}));
+ if (!silent) {
+ message.error(t('consumer.fetchSubscriptionsFailed', { name:
groupName }));
+ }
}
} finally {
- if (!silent) {
+ if (subscriptionRequestIdRef.current[cacheKey] === requestId &&
!silent) {
setSubscriptionLoadingByGroup((prev) => ({ ...prev, [cacheKey]:
false }));
}
}