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]