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

Reply via email to