From 70d707f1b9fe7144ec0e7e1df051fadccf7b9346 Mon Sep 17 00:00:00 2001 From: liuhy Date: Mon, 3 Aug 2026 17:19:03 -0700 Subject: [PATCH 1/2] fix: log proxy remoting cleanup failures --- .../proxy/remoting/RemotingProtocolServer.java | 17 ++++++++++++++++- .../remoting/RemotingProtocolServerTest.java | 13 +++++++++++++ 2 files changed, 29 insertions(+), 1 deletion(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/remoting/RemotingProtocolServer.java b/proxy/src/main/java/org/apache/rocketmq/proxy/remoting/RemotingProtocolServer.java index 4f16a530c76..91473e3c4c3 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/remoting/RemotingProtocolServer.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/remoting/RemotingProtocolServer.java @@ -393,11 +393,26 @@ protected void cleanExpiredRequestInQueue(ThreadPoolExecutor threadPoolExecutor, } else { break; } - } catch (Throwable ignored) { + } catch (Throwable t) { + log.warn("clean expired remoting request failed. queueSize:{}, maxWaitTimeMillsInQueue:{}", + safeQueueSize(threadPoolExecutor), maxWaitTimeMillsInQueue, t); + break; } } } + private int safeQueueSize(ThreadPoolExecutor threadPoolExecutor) { + try { + BlockingQueue queue = threadPoolExecutor == null ? null : threadPoolExecutor.getQueue(); + if (queue == null) { + return -1; + } + return queue.size(); + } catch (Throwable ignored) { + return -1; + } + } + private RequestTask castRunnable(final Runnable runnable) { try { if (runnable instanceof FutureTaskExt) { diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/remoting/RemotingProtocolServerTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/remoting/RemotingProtocolServerTest.java index acd9c1c2d44..da1f057708c 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/remoting/RemotingProtocolServerTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/remoting/RemotingProtocolServerTest.java @@ -17,6 +17,8 @@ package org.apache.rocketmq.proxy.remoting; +<<<<<<< HEAD +import java.util.concurrent.ThreadPoolExecutor; import org.apache.rocketmq.remoting.RemotingServer; import org.apache.rocketmq.remoting.protocol.RequestHeaderRegistry; import org.junit.Test; @@ -44,4 +46,15 @@ public void shouldInitializeRequestHeaderRegistryWhenRegisteringProcessors() { verify(requestHeaderRegistry).initialize(); } } + + @Test(timeout = 1000) + public void testCleanExpiredRequestInQueueBreaksWhenQueueAccessFails() { + RemotingProtocolServer server = Mockito.mock(RemotingProtocolServer.class, Mockito.CALLS_REAL_METHODS); + ThreadPoolExecutor executor = Mockito.mock(ThreadPoolExecutor.class); + Mockito.when(executor.getQueue()).thenThrow(new RuntimeException("queue unavailable")); + + server.cleanExpiredRequestInQueue(executor, 1); + + Mockito.verify(executor, Mockito.atMost(2)).getQueue(); + } } From aed0ef59179290c299de258fafe181d57a3a8156 Mon Sep 17 00:00:00 2001 From: liuhy Date: Sat, 15 Aug 2026 00:05:35 -0700 Subject: [PATCH 2/2] test: remove remoting protocol conflict marker --- .../rocketmq/proxy/remoting/RemotingProtocolServerTest.java | 1 - 1 file changed, 1 deletion(-) diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/remoting/RemotingProtocolServerTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/remoting/RemotingProtocolServerTest.java index da1f057708c..ca9a99715ed 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/remoting/RemotingProtocolServerTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/remoting/RemotingProtocolServerTest.java @@ -17,7 +17,6 @@ package org.apache.rocketmq.proxy.remoting; -<<<<<<< HEAD import java.util.concurrent.ThreadPoolExecutor; import org.apache.rocketmq.remoting.RemotingServer; import org.apache.rocketmq.remoting.protocol.RequestHeaderRegistry;