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 c230fc989 fix(message): preserve replacement queue pull locks (#4501)
c230fc989 is described below

commit c230fc9891967496c53a17d2819547eb5c0b06cd
Author: aias00 <[email protected]>
AuthorDate: Sun Sep 20 21:10:21 2026 -0700

    fix(message): preserve replacement queue pull locks (#4501)
    
    The finally block of handlePull dropped the queue's in-flight lock
    unconditionally, so a request from a previous generation could delete the 
lock
    belonging to the request that replaced it. That allowed a second concurrent 
pull
    on the same queue and cleared the spinner while the newer request was still
    running.
    
    Only release the lock when the finishing request is still the current one, 
which
    matches how the two other request-sequence guards in this hook already 
behave.
---
 web/src/components/QueueBrowser.tsx                |  6 +++--
 web/src/components/__tests__/QueueBrowser.test.tsx | 31 ++++++++++++++++++++++
 2 files changed, 35 insertions(+), 2 deletions(-)

diff --git a/web/src/components/QueueBrowser.tsx 
b/web/src/components/QueueBrowser.tsx
index 4ebcb2ba7..c7252d139 100644
--- a/web/src/components/QueueBrowser.tsx
+++ b/web/src/components/QueueBrowser.tsx
@@ -139,8 +139,10 @@ export const useQueueBrowser = (instanceId?: string) => {
         message.error(err instanceof Error ? err.message : '拉取消息失败');
       }
     } finally {
-      pullingRef.current.delete(key);
-      if (requestId === requestSeqRef.current) setPulling(new 
Set(pullingRef.current));
+      if (requestId === requestSeqRef.current) {
+        pullingRef.current.delete(key);
+        setPulling(new Set(pullingRef.current));
+      }
     }
   };
 
diff --git a/web/src/components/__tests__/QueueBrowser.test.tsx 
b/web/src/components/__tests__/QueueBrowser.test.tsx
index 79958cb21..bcafa44e1 100644
--- a/web/src/components/__tests__/QueueBrowser.test.tsx
+++ b/web/src/components/__tests__/QueueBrowser.test.tsx
@@ -201,6 +201,37 @@ describe('QueueBrowser request ownership', () => {
     await act(async () => pull.resolve(messageRecord('message-a')));
   });
 
+  it('keeps a replacement pull locked when a stale pull settles', async () => {
+    const stalePull = createDeferred<MessageRecord | null>();
+    const currentPull = createDeferred<MessageRecord | null>();
+    const unexpectedPull = createDeferred<MessageRecord | null>();
+    vi.mocked(getQueueOffsets).mockResolvedValue([queue('broker-a')]);
+    vi.mocked(pullMessageAtOffset)
+      .mockReturnValueOnce(stalePull.promise)
+      .mockReturnValueOnce(currentPull.promise)
+      .mockReturnValue(unexpectedPull.promise);
+    const user = userEvent.setup();
+    render(<QueueBrowserProbe />);
+
+    await user.click(screen.getByRole('button', { name: 'topic-a' }));
+    await user.click(screen.getByRole('button', { name: 'load' }));
+    await waitFor(() => 
expect(screen.getByLabelText('queues')).toHaveTextContent('broker-a'));
+
+    await user.click(screen.getByRole('button', { name: 'pull' }));
+    await waitFor(() => expect(pullMessageAtOffset).toHaveBeenCalledTimes(1));
+
+    await user.click(screen.getByRole('button', { name: 'load' }));
+    await waitFor(() => 
expect(screen.getByLabelText('queues')).toHaveTextContent('broker-a'));
+    await user.click(screen.getByRole('button', { name: 'pull' }));
+    await waitFor(() => expect(pullMessageAtOffset).toHaveBeenCalledTimes(2));
+
+    await act(async () => stalePull.resolve(messageRecord('stale-message')));
+    await user.click(screen.getByRole('button', { name: 'pull' }));
+
+    expect(pullMessageAtOffset).toHaveBeenCalledTimes(2);
+    await act(async () => 
currentPull.resolve(messageRecord('current-message')));
+  });
+
   it('deduplicates queue loads before loading state renders', async () => {
     const queues = createDeferred<QueueOffset[]>();
     vi.mocked(getQueueOffsets).mockReturnValue(queues.promise);

Reply via email to