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..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 @@ -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,35 @@ 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)); + + // 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()); + } + @Test public void qoS1PubAuthFailed() { // not by pass