From 793d44e1e92649739ddd6eabb1e40a7caef58694 Mon Sep 17 00:00:00 2001 From: liaodongnian <361485583@qq.com> Date: Thu, 27 Aug 2026 19:21:08 +0800 Subject: [PATCH] Fix duplicate local delivery results by prioritizing NO_RECEIVER --- .../mqtt/service/LocalDistService.java | 1 + .../mqtt/service/LocalDistServiceTest.java | 46 +++++++++++++++++++ 2 files changed, 47 insertions(+) diff --git a/bifromq-mqtt/bifromq-mqtt-server/src/main/java/org/apache/bifromq/mqtt/service/LocalDistService.java b/bifromq-mqtt/bifromq-mqtt-server/src/main/java/org/apache/bifromq/mqtt/service/LocalDistService.java index 6422b802c..a7f849d13 100644 --- a/bifromq-mqtt/bifromq-mqtt-server/src/main/java/org/apache/bifromq/mqtt/service/LocalDistService.java +++ b/bifromq-mqtt/bifromq-mqtt-server/src/main/java/org/apache/bifromq/mqtt/service/LocalDistService.java @@ -194,6 +194,7 @@ public CompletableFuture dist(DeliveryRequest request) { tenantMeter.recordSummary(MqttTransientFanOutBytes, totalFanOutBytes); // don't include duplicated matchInfo in the result // treat skip as ok + noSub.removeAll(noReceiver); Sets.difference(Sets.union(ok, skip), Sets.union(noSub, noReceiver)).forEach(matchInfo -> resultsBuilder.addResult( DeliveryResult.newBuilder().setMatchInfo(matchInfo).setCode(DeliveryResult.Code.OK).build())); diff --git a/bifromq-mqtt/bifromq-mqtt-server/src/test/java/org/apache/bifromq/mqtt/service/LocalDistServiceTest.java b/bifromq-mqtt/bifromq-mqtt-server/src/test/java/org/apache/bifromq/mqtt/service/LocalDistServiceTest.java index 60ea83091..2d611601d 100644 --- a/bifromq-mqtt/bifromq-mqtt-server/src/test/java/org/apache/bifromq/mqtt/service/LocalDistServiceTest.java +++ b/bifromq-mqtt/bifromq-mqtt-server/src/test/java/org/apache/bifromq/mqtt/service/LocalDistServiceTest.java @@ -403,6 +403,52 @@ public void deliverToNoLocalRoute() { assertEquals(result.getCode(), DeliveryResult.Code.NO_RECEIVER); } + @Test + public void noReceiverWinsNoSubAcrossDeliveryPacksRegardlessOfOrder() { + String noSubFirstTenantId = "noSubFirst"; + String noReceiverFirstTenantId = "noReceiverFirst"; + String topicFilter = "testTopic/#"; + String channelId = "channel0"; + MatchInfo matchInfo = MatchInfo.newBuilder() + .setMatcher(TopicUtil.from(topicFilter)) + .setReceiverId("receiverId") + .build(); + DeliveryPack deliveryPack = DeliveryPack.newBuilder() + .setMessagePack(TopicMessagePack.newBuilder().setTopic("testTopic").build()) + .addMatchInfo(matchInfo) + .build(); + DeliveryPackage deliveryPackage = DeliveryPackage.newBuilder() + .addPack(deliveryPack) + .addPack(deliveryPack) + .build(); + DeliveryRequest request = DeliveryRequest.newBuilder() + .putPackage(noSubFirstTenantId, deliveryPackage) + .putPackage(noReceiverFirstTenantId, deliveryPackage) + .build(); + + IMQTTTransientSession session = mock(IMQTTTransientSession.class); + when(session.publish(any(), any())).thenAnswer(invocation -> invocation.getArgument(1)); + when(localSessionRegistry.get(channelId)).thenReturn(session); + ILocalTopicRouter.ILocalRoutes localRoutes = mock(ILocalTopicRouter.ILocalRoutes.class); + when(localRoutes.localReceiverId()).thenReturn("receiverId"); + when(localRoutes.routesInfo()).thenReturn(Map.of(channelId, 1L)); + when(localTopicRouter.getTopicRoutes(noSubFirstTenantId, matchInfo)) + .thenReturn(Optional.of(CompletableFuture.completedFuture(localRoutes))) + .thenReturn(Optional.empty()); + when(localTopicRouter.getTopicRoutes(noReceiverFirstTenantId, matchInfo)) + .thenReturn(Optional.empty()) + .thenReturn(Optional.of(CompletableFuture.completedFuture(localRoutes))); + + DeliveryReply reply = localDistService.dist(request).join(); + + DeliveryResults noSubFirstResults = reply.getResultMap().get(noSubFirstTenantId); + assertEquals(noSubFirstResults.getResultCount(), 1); + assertEquals(noSubFirstResults.getResult(0).getCode(), DeliveryResult.Code.NO_RECEIVER); + DeliveryResults noReceiverFirstResults = reply.getResultMap().get(noReceiverFirstTenantId); + assertEquals(noReceiverFirstResults.getResultCount(), 1); + assertEquals(noReceiverFirstResults.getResult(0).getCode(), DeliveryResult.Code.NO_RECEIVER); + } + @Test public void deliverToEmptyLocalRoutes() { String tenantId = "tenant1";