ai-yang opened a new issue, #10985:
URL: https://github.com/apache/rocketmq/issues/10985
## Before Creating the Bug Report
- [x] I found a bug, not just a question for GitHub Discussions.
- [x] I searched GitHub Issues, PRs, and Discussions and believe this is not
a duplicate.
- [x] I confirmed that this bug belongs to the `apache/rocketmq` repository.
## Runtime platform environment
Reproduced on Ubuntu 22.04.4 LTS, Linux 5.15.0-186, x86_64. The defect is in
the broker's asynchronous POP revive control flow and is not expected to be
OS-specific.
## RocketMQ version
- Branch: `develop`
- Version: current `5.5.1` development sources
- Commit: `e348efa66b08eb645ee123706ea6492fa9a3ad35`
## JDK Version
OpenJDK 8u502 (Maven runtime).
## Describe the Bug
Proposed severity: **High**, because the failure can violate POP's
at-least-once redelivery guarantee for an unacknowledged message.
`PopReviveService.reviveMsgFromCk` starts an asynchronous business-message
read for every unacknowledged offset and records the checkpoint in
`inflightReviveRequestMap`. `mergeAndRevive` then advances and commits the
revive offset after scheduling that work.
If `EscapeBridge.getMessageAsync` (including its decode/transform chain)
completes exceptionally, the corresponding child future remains exceptional.
`CompletableFuture.allOf(...).whenComplete(...)` still invokes its callback,
but the callback calls `future.getNow(...)`; for an exceptionally completed
future, `getNow` throws `CompletionException`. The callback therefore exits
before it can:
- call `rePutCK` to preserve retryability;
- mark the in-flight checkpoint complete;
- remove it from `inflightReviveRequestMap`.
The only timeout fallback is inside `while (inflightReviveRequestMap.size()
> 3)`. With one to three failed checkpoints and no later traffic, that branch
is never entered. Because the revive offset has already moved past the original
checkpoint, the unacknowledged message can remain unredelivered indefinitely.
This is not a security issue.
## Steps to Reproduce
A deterministic broker-free unit reproduction is sufficient:
1. Build one `PopCheckPoint` containing a single unacknowledged message.
2. Make `EscapeBridge.getMessageAsync(...)` return an exceptionally
completed `CompletableFuture`.
3. Invoke `PopReviveService.mergeAndRevive(...)`.
4. Assert that the revive offset and committed offset advanced, the
checkpoint remains in the in-flight map, and no replacement CK was written.
The core setup is:
```java
CompletableFuture<Triple<MessageExt, String, Boolean>> failed = new
CompletableFuture<>();
failed.completeExceptionally(new RuntimeException("store read failed"));
when(escapeBridge.getMessageAsync(anyString(), anyLong(), anyInt(),
anyString(), anyBoolean()))
.thenReturn(failed);
popReviveService.mergeAndRevive(reviveObj);
```
On the unmodified baseline, two independent runs reproduced the same state:
```text
reviveObj.newOffset = 1
committed revive offset = 1
inflightReviveRequestMap.size() = 1
messageStore.putMessage(rewritten CK) = 0 calls
```
The regression test consequently fails because one CK rewrite was expected
but none occurred.
## What Did You Expect to See?
An exceptional business-message read should be treated as a retryable read
failure. The service should rewrite the checkpoint for that message offset,
finish/remove the in-flight record, and preserve the existing offset
progression without losing the message's retry path.
## What Did You See Instead?
The original revive offset is committed, but the replacement CK is never
written. At low traffic the incomplete in-flight record is not old enough to
trigger any independent scan and the size-gated timeout loop is never entered,
so the message can be skipped indefinitely.
## Additional Context
A minimal fix is to use `CompletableFuture.handle` on the `getBizMessage`
stage. If the upstream read completes exceptionally, log it and return
`(msgOffset, false)`, which feeds the existing `rePutCK` path. The handler
should remain scoped to the read stage so exceptions raised later while
processing/writing a retry retain their current semantics.
With that fix, the full `PopReviveServiceTest` class passes 13/13 tests.
Checkstyle, SpotBugs, and `git diff --check` also pass.
Exact GitHub searches for `PopReviveService getNow`, `PopReviveService
rePutCK exception`, `getBizMessage PopReviveService exception`, and
`PopReviveService CompletionException` found no issue or PR match. Related
[#10658](https://github.com/apache/rocketmq/issues/10658)/[#10659](https://github.com/apache/rocketmq/pull/10659)
concern head-of-line blocking in the newer `PopConsumerService` popkv batch
path. [#10667](https://github.com/apache/rocketmq/issues/10667) concerns
discarded futures in `PopConsumerCache.cleanupRecords` when buffer merging is
enabled. Neither covers the legacy/default `PopReviveService` path committing
an offset without preserving retry after an exceptional read.
--
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]