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 d3335d33 fix(consumer): discard stale thread-stack responses by
request id (#1630)
d3335d33 is described below
commit d3335d33af896f95fa5355cd0a716f03b7c95463
Author: 0 <[email protected]>
AuthorDate: Tue Aug 11 21:13:25 2026 +0800
fix(consumer): discard stale thread-stack responses by request id (#1630)
---
.../pages/instance/__tests__/ConsumerPage.test.tsx | 92 +++++++++++++++++++++-
web/src/pages/instance/consumer.tsx | 22 ++++--
2 files changed, 104 insertions(+), 10 deletions(-)
diff --git a/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
b/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
index 019f1b37..0dd8b5b8 100644
--- a/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
+++ b/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
@@ -16,7 +16,7 @@
*/
import { App } from 'antd';
-import { render, screen, waitFor } from '@testing-library/react';
+import { act, render, screen, waitFor } from '@testing-library/react';
import userEvent from '@testing-library/user-event';
import type React from 'react';
import { MemoryRouter } from 'react-router-dom';
@@ -353,7 +353,10 @@ describe('Consumer page', () => {
await user.click(await screen.findByRole('button', { name: /详情/ }));
await waitFor(() =>
-
expect(consumerService.getConsumerSubscriptions).toHaveBeenCalledWith('remote-cg',
'instance-a'),
+ expect(consumerService.getConsumerSubscriptions).toHaveBeenCalledWith(
+ 'remote-cg',
+ 'instance-a',
+ ),
);
await user.click(screen.getByText('Instance A'));
@@ -364,13 +367,96 @@ describe('Consumer page', () => {
await user.click(await screen.findByRole('button', { name: /详情/ }));
await waitFor(() =>
-
expect(consumerService.getConsumerSubscriptions).toHaveBeenCalledWith('remote-cg',
'instance-b'),
+ expect(consumerService.getConsumerSubscriptions).toHaveBeenCalledWith(
+ 'remote-cg',
+ 'instance-b',
+ ),
);
await waitFor(() =>
expect(consumerService.getConsumerProgress).toHaveBeenCalledWith('remote-cg',
'instance-b'),
);
});
+ it('keeps the latest client stack when an older request resolves last',
async () => {
+ const firstStack = {
+ groupName: 'remote-cg',
+ clientId: 'client-1',
+ capturedAt: '2026-07-23T00:00:00Z',
+ threadCount: 1,
+ threads: [
+ {
+ threadName: 'OldClientThread',
+ threadId: 1,
+ state: 'RUNNABLE',
+ blockedTime: 0,
+ waitedTime: 0,
+ stackTrace: ['old.Stack.run(Old.java:1)'],
+ },
+ ],
+ };
+ const secondStack = {
+ ...firstStack,
+ clientId: 'client-2',
+ threads: [
+ {
+ ...firstStack.threads[0],
+ threadName: 'LatestClientThread',
+ stackTrace: ['latest.Stack.run(Latest.java:2)'],
+ },
+ ],
+ };
+ let resolveFirst!: (value: typeof firstStack) => void;
+ let resolveSecond!: (value: typeof secondStack) => void;
+ vi.mocked(consumerService.getConsumerStack)
+ .mockReturnValueOnce(
+ new Promise((resolve) => {
+ resolveFirst = resolve;
+ }),
+ )
+ .mockReturnValueOnce(
+ new Promise((resolve) => {
+ resolveSecond = resolve;
+ }),
+ );
+ vi.mocked(consumerService.listConsumerGroups).mockResolvedValue([
+ {
+ ...group,
+ instances: [
+ {
+ clientId: 'client-1',
+ protocol: 'Remoting',
+ address: '10.0.0.1:1',
+ subscribedTopics: [],
+ lastHeartbeat: '',
+ topicLag: {},
+ },
+ {
+ clientId: 'client-2',
+ protocol: 'Remoting',
+ address: '10.0.0.2:2',
+ subscribedTopics: [],
+ lastHeartbeat: '',
+ topicLag: {},
+ },
+ ],
+ },
+ ]);
+ const user = userEvent.setup();
+ renderWithProviders(<ConsumerPage />);
+
+ await user.click(await screen.findByRole('button', { name: /详情/ }));
+ await user.click(await screen.findByRole('tab', { name: /在线实例/ }));
+ const stackButtons = await screen.findAllByRole('button', { name: /线程栈/ });
+ await user.click(stackButtons[0]);
+ await user.click(stackButtons[1]);
+
+ await act(async () => resolveSecond(secondStack));
+ expect(await screen.findByText('LatestClientThread')).toBeInTheDocument();
+ await act(async () => resolveFirst(firstStack));
+ expect(screen.getByText('LatestClientThread')).toBeInTheDocument();
+ expect(screen.queryByText('OldClientThread')).not.toBeInTheDocument();
+ });
+
it('loads a consumer client stack trace from the selected instance', async
() => {
vi.mocked(consumerService.listConsumerGroups).mockResolvedValue([
{
diff --git a/web/src/pages/instance/consumer.tsx
b/web/src/pages/instance/consumer.tsx
index b1611145..adb92bae 100644
--- a/web/src/pages/instance/consumer.tsx
+++ b/web/src/pages/instance/consumer.tsx
@@ -231,8 +231,10 @@ const ConsumerPage = () => {
const [importing, setImporting] = useState(false);
const groupRequestIdRef = useRef(0);
+ const stackRequestIdRef = useRef(0);
useEffect(() => {
+ stackRequestIdRef.current += 1;
// eslint-disable-next-line react-hooks/set-state-in-effect -- clear state
owned by the previous instance
setSelectedGroup(null);
setModalOpen(false);
@@ -365,21 +367,25 @@ const ConsumerPage = () => {
const openStackModal = async (consumerInstance: ConsumerInstance) => {
if (!selectedGroup) return;
+ const requestId = ++stackRequestIdRef.current;
+ const groupName = selectedGroup.name;
setSelectedStackClient(consumerInstance);
setSelectedStack(null);
setStackModalOpen(true);
setStackLoading(true);
try {
const stack = await getConsumerStack(
- selectedGroup.name,
+ groupName,
consumerInstance.clientId,
selectedInstanceId || undefined,
);
- setSelectedStack(stack);
+ if (requestId === stackRequestIdRef.current) setSelectedStack(stack);
} catch {
- message.error(`客户端 ${consumerInstance.clientId} 线程栈获取失败`);
+ if (requestId === stackRequestIdRef.current) {
+ message.error(`客户端 ${consumerInstance.clientId} 线程栈获取失败`);
+ }
} finally {
- setStackLoading(false);
+ if (requestId === stackRequestIdRef.current) setStackLoading(false);
}
};
@@ -974,7 +980,9 @@ const ConsumerPage = () => {
subscriptionsByGroup[diagnosticCacheKey(selectedInstanceId, record.name)] ?? []
}
rowKey="topic"
-
loading={subscriptionLoadingByGroup[diagnosticCacheKey(selectedInstanceId,
record.name)]}
+ loading={
+
subscriptionLoadingByGroup[diagnosticCacheKey(selectedInstanceId, record.name)]
+ }
pagination={false}
size="small"
/>
@@ -1299,6 +1307,7 @@ const ConsumerPage = () => {
}
open={stackModalOpen}
onCancel={() => {
+ stackRequestIdRef.current += 1;
setStackModalOpen(false);
setSelectedStack(null);
setSelectedStackClient(null);
@@ -1615,8 +1624,7 @@ const ConsumerPage = () => {
}}
confirmLoading={resetSubmitting}
okButtonProps={{
- disabled:
- !resetTopic ||
Boolean(subscriptionLoadingByGroup[resetDiagnosticKey]),
+ disabled: !resetTopic ||
Boolean(subscriptionLoadingByGroup[resetDiagnosticKey]),
}}
okText="确认重置"
cancelText="取消"