youngkermit8-coder opened a new issue, #5544:
URL: https://github.com/apache/rocketmq-dashboard/issues/5544

   ### Before Creating the Bug Report
   
   - [x] Searched open issues and PRs for DLQ interruption, cancellation and 
topology lookup. I did not find an equivalent report in those searches.
   - [x] This concerns RocketMQ Studio backend DLQ processing, not a cluster 
usage question.
   - [x] Exact branch and commit are stated below.
   
   ### Studio Version
   
   `rocketmq-studio` at `6a68042fa73a513f40436ca721944b4de48181ec`, built from 
source for local provider tests. This is not a report from a deployed Studio 
HTTP instance.
   
   ### Runtime Environment
   
   Windows, JDK 21.0.12, Maven 3.9.9. Provider tests use mocked clients and 
audit collaborators; MySQL and browser behavior are outside this reproduction.
   
   ### Connected RocketMQ Cluster
   
   RocketMQ Java SDK 5.5.0. The deterministic topology reproduction below uses 
a mocked admin client. A separate earlier send-path experiment used a loopback 
RocketMQ 5.5.0 NameServer and Broker with synthetic messages and direct 
remoting access.
   
   ### Describe the Bug
   
   Direct `InterruptedException` during DLQ scanning, selected-message lookup 
or sending is caught as an ordinary failure, allowing subsequent work and 
losing the interrupt flag.
   
   There is also an indirect path: 
`BrokerTopologyGuards.isWithinKnownBrokerTopology` catches interruption from 
`examineBrokerClusterInfo` and returns `false`. During selected-message resend, 
this can skip the interrupted lookup but still send messages collected before 
it. Handling interruption only around `viewMessage` does not cover this path.
   
   ### Steps to Reproduce
   
   In the existing `RocketMQDLQProviderTest` fixture, resolve one selected DLQ 
message and then interrupt the topology lookup for a valid offset-style ID. The 
following diagnostic uses the fixture's existing mocks; fully qualified names 
avoid extra imports.
   
   ```java
   @Test
   void interruptedTopologyLookupMustAbortSelectedResendTest() throws Exception 
{
       String topic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
       String offsetId = MessageDecoder.createMessageId(
               new java.net.InetSocketAddress("127.0.0.1", 10911), 12345L);
       MessageExt message = new MessageExt();
       message.setMsgId("before-topology-interruption");
       message.setTopic(topic);
       message.setBody(new byte[] {1});
       MessageAccessor.putProperty(message, 
MessageConst.PROPERTY_DLQ_ORIGIN_TOPIC, "orders");
       when(adminExt.viewMessage(topic, 
message.getMsgId())).thenReturn(message);
       when(adminExt.examineBrokerClusterInfo()).thenThrow(new 
InterruptedException("topology interrupted"));
       java.util.concurrent.atomic.AtomicInteger sends = new 
java.util.concurrent.atomic.AtomicInteger();
       
org.mockito.Mockito.lenient().when(dlqProducer.send(any(Message.class))).thenAnswer(invocation
 -> {
           sends.incrementAndGet();
           SendResult sent = new SendResult();
           sent.setSendStatus(SendStatus.SEND_OK);
           return sent;
       });
       try {
           Throwable error = org.assertj.core.api.Assertions.catchThrowable(() 
-> provider.resendMessages(
                   "instance-a", "group-a", List.of(message.getMsgId(), 
offsetId), null));
           org.assertj.core.api.SoftAssertions.assertSoftly(softly -> {
               softly.assertThat(error).as("operation 
aborts").isInstanceOf(BusinessException.class);
               
softly.assertThat(Thread.currentThread().isInterrupted()).as("interrupt 
preserved").isTrue();
               softly.assertThat(sends.get()).as("sends after 
interruption").isZero();
           });
       } finally {
           Thread.interrupted();
       }
   }
   ```
   
   Run from `server` with the development test profile:
   
   ```text
   mvn -B -ntp -Dspring.profiles.active=dev 
-Dtest=RocketMQDLQProviderTest#interruptedTopologyLookupMustAbortSelectedResendTest
 test
   ```
   
   The diagnostic observed no thrown abort, a cleared interrupt flag and one 
producer send. All three assertions failed before the topology fix and passed 
after it.
   
   ### What Did You Expect to See
   
   Stop subsequent work and preserve interruption. Interrupted scanning or 
lookup should not lead to resending a partially collected batch. Interrupted 
sending should retain already confirmed successes and account for unfinished 
entries without additional sends.
   
   ### What Did You See Instead
   
   The topology interruption was treated as an unverified ID, and the 
previously resolved message was sent. In the separate real SDK send-path 
experiment, interrupting a blocked first send also allowed the original 
provider to send the second message.
   
   ### Additional Context
   
   An AI-assisted local prototype catches direct interruption in the provider 
stages and adds a package-private interruptible topology guard for DLQ. 
Existing guard callers retain their boolean, fail-closed interface and 
broker-address validation. The prototype was prepared locally before opening 
this issue; it has not been submitted as a PR. Feedback on the abort and 
result-accounting semantics is requested before proceeding.
   
   The candidate passes 231 related tests across 13 classes. A separate 
expanded check passes 251 test executions, including 10 repetitions of one 
worker-thread interruption test with mocked SDK calls. Two local mutation 
checks detect a swallowed topology interruption and a lost interrupt flag. 
Checkstyle reports zero violations.
   
   Packaging remains blocked at `binary-license-gate`: baseline and candidate 
have the same 165 manual-review findings. This check was not bypassed. Full 
repository CI has not been run. The earlier real Broker experiment was not 
rerun against the topology revision.
   
   An in-flight send may already have reached the Broker; this proposal cannot 
retract it or promise exactly-once delivery. HTTP cancellation propagation, 
durable audit behavior, wrapped SDK interruptions and every pre-existing 
interrupt-flag path are not claimed as covered. This differs from #5386, which 
concerns closing the frontend modal while a request is pending.
   
   ### Are You Willing to Submit a Pull Request
   
   - [x] Yes, after feedback on this proposed fix. The PR would target 
`rocketmq-studio` as specified in `CONTRIBUTING.md`.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to