Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
16 commits
Select commit Hold shift + click to select a range
942c547
feat: implement MQTT Last Will and Testament (LWT)
wy471x Aug 9, 2026
ab95921
Potential fix for pull request finding
yu199195 Aug 11, 2026
22a7e8f
Merge branch 'master' into feat_Last-Will-&-Testament-entirely-unimpl…
wy471x Aug 14, 2026
74d51d4
Merge branch 'master' into feat_Last-Will-&-Testament-entirely-unimpl…
Aias00 Aug 14, 2026
b37ff8b
Merge branch 'master' into feat_Last-Will-&-Testament-entirely-unimpl…
wy471x Aug 14, 2026
20e61b2
Merge branch 'master' into feat_Last-Will-&-Testament-entirely-unimpl…
wy471x Aug 14, 2026
cfc9cb8
Merge branch 'master' into feat_Last-Will-&-Testament-entirely-unimpl…
wy471x Aug 15, 2026
61384ab
feat: route will delivery through wildcard topic matching
wy471x Aug 15, 2026
21b42ba
Merge branch 'master' into feat_Last-Will-&-Testament-entirely-unimpl…
wy471x Aug 17, 2026
f40dc11
Merge branch 'master' into feat_Last-Will-&-Testament-entirely-unimpl…
wy471x Sep 17, 2026
2a8a5d9
fix(mqtt): make SubscribeRepository updates synchronous and race-free
wy471x Sep 19, 2026
1510139
Merge branch 'master' into feat_Last-Will-&-Testament-entirely-unimpl…
wy471x Sep 20, 2026
8dc98ea
Merge branch 'master' into feat_Last-Will-&-Testament-entirely-unimpl…
Aias00 Sep 23, 2026
ed1f32d
Merge branch 'master' into feat_Last-Will-&-Testament-entirely-unimpl…
wy471x Oct 1, 2026
b29821d
fix(mqtt): resolve merge conflicts in SubscribeRepository and its tests
wy471x Oct 1, 2026
71fb63a
fix(mqtt): cap will qos at the subscriber granted qos
wy471x Oct 1, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 12 additions & 0 deletions shenyu-protocol/shenyu-protocol-mqtt/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,18 @@
<artifactId>junit-jupiter</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.mockito</groupId>
<artifactId>mockito-junit-jupiter</artifactId>
<version>${mockito.version}</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.mockito</groupId>
<artifactId>mockito-core</artifactId>
<version>${mockito.version}</version>
<scope>test</scope>
</dependency>
</dependencies>

</project>
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@
import org.apache.commons.lang3.StringUtils;
import org.apache.shenyu.common.utils.Singleton;
import org.apache.shenyu.protocol.mqtt.repositories.ChannelRepository;
import org.apache.shenyu.protocol.mqtt.repositories.WillRepository;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

Expand Down Expand Up @@ -74,6 +75,17 @@ public void connect(final ChannelHandlerContext ctx, final MqttConnectMessage ms

// record connect
Singleton.INST.get(ChannelRepository.class).add(ctx.channel(), clientId);

// store will if present
if (msg.variableHeader().isWillFlag()) {
WillRepository.WillEntry will = new WillRepository.WillEntry(
msg.payload().willTopic(),
msg.payload().willMessageInBytes(),
msg.variableHeader().willQos(),
msg.variableHeader().isWillRetain());
Singleton.INST.get(WillRepository.class).add(ctx.channel(), will);
}

MqttConnAckMessage ackMessage = MqttMessageBuilders.connAck()
.returnCode(MqttConnectReturnCode.CONNECTION_ACCEPTED)
.sessionPresent(true)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
import org.apache.shenyu.common.utils.Singleton;
import org.apache.shenyu.protocol.mqtt.repositories.ChannelRepository;
import org.apache.shenyu.protocol.mqtt.utils.MqttPacketIdGenerator;
import org.apache.shenyu.protocol.mqtt.repositories.WillRepository;

/**
* The DISCONNECT message is sent from the client to the server to indicate
Expand All @@ -37,8 +38,7 @@ public class Disconnect extends MessageType {

@Override
public void disconnect(final ChannelHandlerContext ctx) {
//// todo Last words
//// todo Clean session
Singleton.INST.get(WillRepository.class).remove(ctx.channel());
cleanChannel(ctx.channel());
ctx.close();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -68,8 +68,10 @@ public void connect() {
case PUBREL:
messageType.pubRel(ctx, msg);
break;
case PUBACK:
case DISCONNECT:
messageType.disconnect(ctx);
break;
case PUBACK:
default:
break;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,9 @@
import org.apache.shenyu.protocol.mqtt.repositories.ChannelRepository;
import org.apache.shenyu.protocol.mqtt.repositories.SubscribeRepository;
import org.apache.shenyu.protocol.mqtt.utils.MqttPacketIdGenerator;
import org.apache.shenyu.protocol.mqtt.repositories.WillRepository;

import java.util.Objects;

/**
* mqtt transport handler.
Expand All @@ -51,8 +54,19 @@ public void channelRead(final ChannelHandlerContext ctx, final Object msg) throw

@Override
public void channelInactive(final ChannelHandlerContext ctx) throws Exception {
Singleton.INST.get(ChannelRepository.class).remove(ctx.channel());
ctx.fireChannelInactive();
final Channel channel = ctx.channel();
Singleton.INST.get(ChannelRepository.class).remove(channel);

final WillRepository willRepository = Singleton.INST.get(WillRepository.class);
final WillRepository.WillEntry will = willRepository.get(channel);
if (Objects.nonNull(will)) {
// a will is published at most once, and the repository keeps a strong reference
// to the channel, so it must be removed even if publishing fails.
willRepository.remove(channel);
Publish.publishWill(will);
}
// local state is consistent now, notify the rest of the pipeline exactly once.
super.channelInactive(ctx);
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,9 @@
import org.apache.shenyu.protocol.mqtt.repositories.TopicRepository;
import org.apache.shenyu.protocol.mqtt.utils.MqttPacketIdGenerator;

import org.apache.shenyu.protocol.mqtt.repositories.WillRepository;

import java.util.Objects;
import java.util.Map;
import java.util.concurrent.CompletableFuture;

Expand Down Expand Up @@ -143,4 +146,31 @@ private void send(final String topic, final ByteBuf payload, final MqttQoS publi
private static MqttQoS minQoS(final MqttQoS publishQoS, final MqttQoS grantedQoS) {
return publishQoS.value() <= grantedQoS.value() ? publishQoS : grantedQoS;
}

/**
* Publish a Last Will message to all subscribers of the will topic.
*
* @param will the will entry containing topic, message, qos, and retain flag
*/
static void publishWill(final WillRepository.WillEntry will) {
if (Objects.isNull(will) || Objects.isNull(will.getTopic()) || Objects.isNull(will.getMessage())) {
return;
}
final Map<Channel, MqttQoS> subscribers = Singleton.INST.get(SubscribeRepository.class).getChannelsByTopic(will.getTopic());
final MqttQoS willQos = MqttQoS.valueOf(will.getQos());
subscribers.entrySet().parallelStream().forEach(entry -> {
Channel channel = entry.getKey();
if (channel.isActive()) {
MqttQoS qos = minQoS(willQos, entry.getValue());
int packetId = MqttQoS.AT_MOST_ONCE == qos
? 0
: java.util.concurrent.ThreadLocalRandom.current().nextInt(1, 65536);
MqttFixedHeader mqttFixedHeader = new MqttFixedHeader(MqttMessageType.PUBLISH, false, qos, will.isRetain(), 0);
MqttPublishVariableHeader mqttPublishVariableHeader = new MqttPublishVariableHeader(will.getTopic(), packetId);
MqttPublishMessage mqttPublishMessage = new MqttPublishMessage(mqttFixedHeader, mqttPublishVariableHeader,
Unpooled.wrappedBuffer(will.getMessage()));
channel.writeAndFlush(mqttPublishMessage);
}
});
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,103 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/

package org.apache.shenyu.protocol.mqtt;

import java.util.Objects;

/**
* MQTT topic filter matching per MQTT-4.7.
*
* <p>+ matches exactly one topic level.</p>
*
* <p># matches any number of subsequent levels (must appear at the end of the filter).</p>
*/
public final class TopicMatcher {

private TopicMatcher() {
}

/**
* Check whether a topic filter matches a topic name.
*
* @param filter the subscription topic filter (may contain + and # wildcards)
* @param topic the published topic name (no wildcards)
* @return true if the filter matches the topic
*/
public static boolean matches(final String filter, final String topic) {
if (Objects.isNull(filter) || Objects.isNull(topic)) {
return false;
}

// $ topics must not be matched by wildcards at the first level
if (topic.startsWith("$") && filter.length() > 0 && (filter.charAt(0) == '+' || filter.charAt(0) == '#')) {
return false;
}

String[] filterLevels = filter.split("/", -1);
String[] topicLevels = topic.split("/", -1);

int filterLen = filterLevels.length;
int topicLen = topicLevels.length;

for (int i = 0; i < filterLen; i++) {
String f = filterLevels[i];

if ("#".equals(f)) {
// MQTT-4.7.1-2: # matches any number of levels including the parent level
return i == filterLen - 1;
}

if (i >= topicLen) {
return false;
}

if (!"+".equals(f) && !f.equals(topicLevels[i])) {
return false;
}
}

return filterLen == topicLen;
}

/**
* Validate a topic filter per MQTT-4.7.1: wildcards must occupy an entire
* level, and # must be the last level. Filters must not be empty or
* contain the null character.
*
* @param filter the subscription topic filter
* @return true if the filter is valid
*/
public static boolean isValidFilter(final String filter) {
if (Objects.isNull(filter) || filter.isEmpty() || filter.indexOf((char) 0) >= 0) {
return false;
}
String[] levels = filter.split("/", -1);
for (int i = 0; i < levels.length; i++) {
String level = levels[i];
if (level.indexOf('+') >= 0 || level.indexOf('#') >= 0) {
if (level.length() > 1) {
return false;
}
if ("#".equals(level) && i != levels.length - 1) {
return false;
}
}
}
return true;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
import io.netty.channel.Channel;
import io.netty.handler.codec.mqtt.MqttQoS;
import io.netty.handler.codec.mqtt.MqttTopicSubscription;
import org.apache.shenyu.protocol.mqtt.TopicMatcher;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

Expand Down Expand Up @@ -109,4 +110,35 @@ private static MqttQoS maxQoS(final MqttQoS qos1, final MqttQoS qos2) {
return qos1.value() >= qos2.value() ? qos1 : qos2;
}

/**
* Get the channels whose subscription filter matches the published topic,
* mapped to the maximum qos granted across all their matching filters.
* Supports MQTT wildcards: + (single-level) and # (multi-level).
*
* @param topic the published topic name
* @return matching channels with their maximum granted qos
*/
public Map<Channel, MqttQoS> getChannelsByTopic(final String topic) {
// MQTT requires at most one delivery per publish per client, so merge the
// granted qos when overlapping filters (e.g. sport/# and #) both match.
Map<Channel, MqttQoS> result = new ConcurrentHashMap<>();

// fast path: exact subscription, no wildcard scan needed
Map<Channel, MqttQoS> exactMatch = TOPIC_CHANNEL_FACTORY.get(topic);
if (Objects.nonNull(exactMatch)) {
result.putAll(exactMatch);
}

for (Map.Entry<String, Map<Channel, MqttQoS>> entry : TOPIC_CHANNEL_FACTORY.entrySet()) {
String filter = entry.getKey();
if (filter.equals(topic) || filter.indexOf('+') < 0 && filter.indexOf('#') < 0) {
continue;
}
if (TopicMatcher.matches(filter, topic)) {
entry.getValue().forEach((channel, qos) -> result.merge(channel, qos, SubscribeRepository::maxQoS));
}
}
return result;
}

}
Original file line number Diff line number Diff line change
@@ -0,0 +1,85 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/

package org.apache.shenyu.protocol.mqtt.repositories;

import io.netty.channel.Channel;

import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;

/**
* Stores Last Will and Testament for connected clients.
* Will is set on CONNECT and cleared on graceful DISCONNECT.
* On ungraceful disconnect (channelInactive with will present), the will is published.
*/
public class WillRepository implements BaseRepository<Channel, WillRepository.WillEntry> {

private static final Map<Channel, WillEntry> WILL_FACTORY = new ConcurrentHashMap<>();

@Override
public void add(final Channel channel, final WillEntry willEntry) {
WILL_FACTORY.put(channel, willEntry);
}

@Override
public void remove(final Channel channel) {
WILL_FACTORY.remove(channel);
}

@Override
public WillEntry get(final Channel channel) {
return WILL_FACTORY.get(channel);
}

/**
* Holds the will message fields from a CONNECT payload.
*/
public static class WillEntry {

private final String topic;

private final byte[] message;

private final int qos;

private final boolean retain;

public WillEntry(final String topic, final byte[] message, final int qos, final boolean retain) {
this.topic = topic;
this.message = message;
this.qos = qos;
this.retain = retain;
}

public String getTopic() {
return topic;
}

public byte[] getMessage() {
return message;
}

public int getQos() {
return qos;
}

public boolean isRetain() {
return retain;
}
}
}
Loading
Loading