From 805838bbe7b2111264930dbef0a9c766cbdf8afa Mon Sep 17 00:00:00 2001 From: ImDanXie Date: Thu, 24 Sep 2026 21:15:05 +0800 Subject: [PATCH 1/2] fix(mqtt): do not advance the inbox watermark on out-of-order PUBACK (#286) confirm() unconditionally called onConfirm(confirmingMsg.seq) with the PASSED-IN message, even when the drain loop removed nothing because the head of unconfirmedPacketIds was still un-acked. A client acknowledging packetId 2 while packetId 1 is in flight (the spec places no ordering requirement on PUBACK) would then clear stagingBuffer up to the second message's seq and commit sendBufferUpToSeq past it - deleting the never-acknowledged first message from the inbox. On session resume it is not redelivered: permanent loss. Advance the watermark only when the drain loop actually removed a contiguous prefix; when zero entries were removed, do not call onConfirm at all. Regression test: outOfOrderPubAckMustNotAdvanceInboxWatermark. Control experiment: fails on unfixed code, passes with the fix. Full PersistentSessionHandlerTest (24) passes. Fixes #286 --- .../mqtt/handler/MQTTSessionHandler.java | 10 ++++++-- .../v3/MQTT3PersistentSessionHandlerTest.java | 25 +++++++++++++++++++ 2 files changed, 33 insertions(+), 2 deletions(-) diff --git a/bifromq-mqtt/bifromq-mqtt-server/src/main/java/org/apache/bifromq/mqtt/handler/MQTTSessionHandler.java b/bifromq-mqtt/bifromq-mqtt-server/src/main/java/org/apache/bifromq/mqtt/handler/MQTTSessionHandler.java index 97c859075..9811fea04 100644 --- a/bifromq-mqtt/bifromq-mqtt-server/src/main/java/org/apache/bifromq/mqtt/handler/MQTTSessionHandler.java +++ b/bifromq-mqtt/bifromq-mqtt-server/src/main/java/org/apache/bifromq/mqtt/handler/MQTTSessionHandler.java @@ -897,12 +897,14 @@ private void confirm(ConfirmingMessage confirmingMsg, boolean delivered) { long now = sessionCtx.nanoTime(); confirmingMsg.setAcked(); Iterator packetIdItr = unconfirmedPacketIds.keySet().iterator(); + boolean anyConfirmed = false; while (packetIdItr.hasNext()) { int packetId = packetIdItr.next(); ConfirmingMessage head = unconfirmedPacketIds.get(packetId); if (head.acked) { packetIdItr.remove(); confirmingMsg = head; + anyConfirmed = true; long lastSentTimestamp = head.resendTimestamp > 0 ? head.resendTimestamp : head.timestamp; RoutedMessage confirmed = confirmingMsg.message; switch (confirmed.qos()) { @@ -958,8 +960,12 @@ private void confirm(ConfirmingMessage confirmingMsg, boolean delivered) { break; } } - // confirm up to the current seq - onConfirm(confirmingMsg.seq); + // confirm up to the last contiguously acknowledged seq; if the head of the + // window is still unacknowledged nothing is confirmed, so the watermark must + // not advance past it + if (anyConfirmed) { + onConfirm(confirmingMsg.seq); + } } protected abstract void onConfirm(long seq); diff --git a/bifromq-mqtt/bifromq-mqtt-server/src/test/java/org/apache/bifromq/mqtt/handler/v3/MQTT3PersistentSessionHandlerTest.java b/bifromq-mqtt/bifromq-mqtt-server/src/test/java/org/apache/bifromq/mqtt/handler/v3/MQTT3PersistentSessionHandlerTest.java index 806fa2aa4..ca0d969a4 100644 --- a/bifromq-mqtt/bifromq-mqtt-server/src/test/java/org/apache/bifromq/mqtt/handler/v3/MQTT3PersistentSessionHandlerTest.java +++ b/bifromq-mqtt/bifromq-mqtt-server/src/test/java/org/apache/bifromq/mqtt/handler/v3/MQTT3PersistentSessionHandlerTest.java @@ -42,10 +42,12 @@ import static org.apache.bifromq.type.QoS.EXACTLY_ONCE; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.argThat; +import static org.mockito.Mockito.never; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertNotNull; import static org.testng.Assert.assertNull; import static org.testng.Assert.assertTrue; @@ -477,6 +479,29 @@ public void qoS1PubAndNotAllAck() { verify(inboxClient, times(0)).commit(argThat(CommitRequest::hasSendBufferUpToSeq)); } + @Test + public void outOfOrderPubAckMustNotAdvanceInboxWatermark() { + // #286: PUBACK carries no ordering requirement — a client may acknowledge + // packetId 2 while packetId 1 is still in flight. Confirming the inbox + // watermark up to the second message's seq would delete the un-acked first + // message from the inbox (permanent loss on session resume). + mockCheckPermission(true); + mockInboxCommit(CommitReply.Code.OK); + inboxFetchConsumer.accept(fetch(2, 128, QoS.AT_LEAST_ONCE)); + channel.runPendingTasks(); + MqttPublishMessage first = channel.readOutbound(); + MqttPublishMessage second = channel.readOutbound(); + assertNotNull(first); + assertNotNull(second); + // PUBACK the SECOND message only — the head (first) is still un-acked + channel.writeInbound(MQTTMessageUtils.pubAckMessage(second.variableHeader().packetId())); + channel.runPendingTasks(); + // nothing was contiguously confirmed -> the watermark must not advance + verify(inboxClient, never()).commit(argThat(CommitRequest::hasSendBufferUpToSeq)); + // the still-un-acked first message must survive in the session + assertTrue(channel.isOpen()); + } + @Test public void qoS1PubAuthFailed() { // not by pass From 0321772d7b233a1d58c48c5b87d4523cf1a8e880 Mon Sep 17 00:00:00 2001 From: ImDanXie Date: Thu, 24 Sep 2026 21:36:58 +0800 Subject: [PATCH 2/2] test: verify the watermark catches up after the delayed contiguous ack --- .../handler/v3/MQTT3PersistentSessionHandlerTest.java | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/bifromq-mqtt/bifromq-mqtt-server/src/test/java/org/apache/bifromq/mqtt/handler/v3/MQTT3PersistentSessionHandlerTest.java b/bifromq-mqtt/bifromq-mqtt-server/src/test/java/org/apache/bifromq/mqtt/handler/v3/MQTT3PersistentSessionHandlerTest.java index ca0d969a4..b051c544a 100644 --- a/bifromq-mqtt/bifromq-mqtt-server/src/test/java/org/apache/bifromq/mqtt/handler/v3/MQTT3PersistentSessionHandlerTest.java +++ b/bifromq-mqtt/bifromq-mqtt-server/src/test/java/org/apache/bifromq/mqtt/handler/v3/MQTT3PersistentSessionHandlerTest.java @@ -498,7 +498,13 @@ public void outOfOrderPubAckMustNotAdvanceInboxWatermark() { channel.runPendingTasks(); // nothing was contiguously confirmed -> the watermark must not advance verify(inboxClient, never()).commit(argThat(CommitRequest::hasSendBufferUpToSeq)); - // the still-un-acked first message must survive in the session + + // phase 2: PUBACK the first — the delayed contiguous confirmation must now + // drain both acked entries and advance the watermark past the second message + channel.writeInbound(MQTTMessageUtils.pubAckMessage(first.variableHeader().packetId())); + channel.runPendingTasks(); + verify(inboxClient, times(1)).commit(argThat(req -> + req.hasSendBufferUpToSeq() && req.getSendBufferUpToSeq() >= 1)); assertTrue(channel.isOpen()); }