View Javadoc
1   /*
2    * Copyright 2014 The Netty Project
3    *
4    * The Netty Project licenses this file to you under the Apache License,
5    * version 2.0 (the "License"); you may not use this file except in compliance
6    * with the License. You may obtain a copy of the License at:
7    *
8    *   https://www.apache.org/licenses/LICENSE-2.0
9    *
10   * Unless required by applicable law or agreed to in writing, software
11   * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
12   * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
13   * License for the specific language governing permissions and limitations
14   * under the License.
15   */
16  
17  package io.netty.handler.codec.mqtt;
18  
19  import io.netty.buffer.ByteBuf;
20  import io.netty.buffer.ByteBufAllocator;
21  import io.netty.buffer.Unpooled;
22  import io.netty.channel.ChannelHandler;
23  import io.netty.channel.ChannelHandlerContext;
24  import io.netty.handler.codec.EncoderException;
25  import io.netty.handler.codec.MessageToMessageEncoder;
26  import io.netty.util.internal.EmptyArrays;
27  
28  import java.util.List;
29  
30  import static io.netty.buffer.ByteBufUtil.*;
31  import static io.netty.handler.codec.mqtt.MqttCodecUtil.getMqttVersion;
32  import static io.netty.handler.codec.mqtt.MqttCodecUtil.isValidClientId;
33  import static io.netty.handler.codec.mqtt.MqttCodecUtil.isValidPublishTopicName;
34  import static io.netty.handler.codec.mqtt.MqttCodecUtil.isValidUserName;
35  import static io.netty.handler.codec.mqtt.MqttCodecUtil.setMqttVersion;
36  import static io.netty.handler.codec.mqtt.MqttConstant.DEFAULT_MAX_CLIENT_ID_LENGTH;
37  import static io.netty.handler.codec.mqtt.MqttProperties.ASSIGNED_CLIENT_IDENTIFIER;
38  import static io.netty.handler.codec.mqtt.MqttProperties.AUTHENTICATION_DATA;
39  import static io.netty.handler.codec.mqtt.MqttProperties.AUTHENTICATION_METHOD;
40  import static io.netty.handler.codec.mqtt.MqttProperties.CONTENT_TYPE;
41  import static io.netty.handler.codec.mqtt.MqttProperties.CORRELATION_DATA;
42  import static io.netty.handler.codec.mqtt.MqttProperties.MAXIMUM_PACKET_SIZE;
43  import static io.netty.handler.codec.mqtt.MqttProperties.MAXIMUM_QOS;
44  import static io.netty.handler.codec.mqtt.MqttProperties.PAYLOAD_FORMAT_INDICATOR;
45  import static io.netty.handler.codec.mqtt.MqttProperties.PUBLICATION_EXPIRY_INTERVAL;
46  import static io.netty.handler.codec.mqtt.MqttProperties.REASON_STRING;
47  import static io.netty.handler.codec.mqtt.MqttProperties.RECEIVE_MAXIMUM;
48  import static io.netty.handler.codec.mqtt.MqttProperties.REQUEST_PROBLEM_INFORMATION;
49  import static io.netty.handler.codec.mqtt.MqttProperties.REQUEST_RESPONSE_INFORMATION;
50  import static io.netty.handler.codec.mqtt.MqttProperties.RESPONSE_INFORMATION;
51  import static io.netty.handler.codec.mqtt.MqttProperties.RESPONSE_TOPIC;
52  import static io.netty.handler.codec.mqtt.MqttProperties.RETAIN_AVAILABLE;
53  import static io.netty.handler.codec.mqtt.MqttProperties.SERVER_KEEP_ALIVE;
54  import static io.netty.handler.codec.mqtt.MqttProperties.SERVER_REFERENCE;
55  import static io.netty.handler.codec.mqtt.MqttProperties.SESSION_EXPIRY_INTERVAL;
56  import static io.netty.handler.codec.mqtt.MqttProperties.SHARED_SUBSCRIPTION_AVAILABLE;
57  import static io.netty.handler.codec.mqtt.MqttProperties.SUBSCRIPTION_IDENTIFIER;
58  import static io.netty.handler.codec.mqtt.MqttProperties.SUBSCRIPTION_IDENTIFIER_AVAILABLE;
59  import static io.netty.handler.codec.mqtt.MqttProperties.TOPIC_ALIAS;
60  import static io.netty.handler.codec.mqtt.MqttProperties.TOPIC_ALIAS_MAXIMUM;
61  import static io.netty.handler.codec.mqtt.MqttProperties.USER_PROPERTY;
62  import static io.netty.handler.codec.mqtt.MqttProperties.WILDCARD_SUBSCRIPTION_AVAILABLE;
63  import static io.netty.handler.codec.mqtt.MqttProperties.WILL_DELAY_INTERVAL;
64  
65  /**
66   * Encodes Mqtt messages into bytes following the protocol specification v3.1
67   * as described here <a href="https://public.dhe.ibm.com/software/dw/webservices/ws-mqtt/mqtt-v3r1.html">MQTTV3.1</a>
68   * or v5.0 as described here <a href="https://docs.oasis-open.org/mqtt/mqtt/v5.0/mqtt-v5.0.html">MQTTv5.0</a> -
69   * depending on the version specified in the first CONNECT message that goes through the channel.
70   */
71  @ChannelHandler.Sharable
72  public final class MqttEncoder extends MessageToMessageEncoder<MqttMessage> {
73  
74      public static final MqttEncoder INSTANCE = new MqttEncoder();
75  
76      private MqttEncoder() {
77          super(MqttMessage.class);
78      }
79  
80      @Override
81      protected void encode(ChannelHandlerContext ctx, MqttMessage msg, List<Object> out) throws Exception {
82          out.add(doEncode(ctx, msg));
83      }
84  
85      /**
86       * This is the main encoding method.
87       * It's only visible for testing.
88       *
89       * @param message MQTT message to encode
90       * @return ByteBuf with encoded bytes
91       */
92      static ByteBuf doEncode(ChannelHandlerContext ctx,
93                       MqttMessage message) {
94  
95          switch (message.fixedHeader().messageType()) {
96              case CONNECT:
97                  return encodeConnectMessage(ctx, (MqttConnectMessage) message);
98  
99              case CONNACK:
100                 return encodeConnAckMessage(ctx, (MqttConnAckMessage) message);
101 
102             case PUBLISH:
103                 return encodePublishMessage(ctx, (MqttPublishMessage) message);
104 
105             case SUBSCRIBE:
106                 return encodeSubscribeMessage(ctx, (MqttSubscribeMessage) message);
107 
108             case UNSUBSCRIBE:
109                 return encodeUnsubscribeMessage(ctx,  (MqttUnsubscribeMessage) message);
110 
111             case SUBACK:
112                 return encodeSubAckMessage(ctx, (MqttSubAckMessage) message);
113 
114             case UNSUBACK:
115                 if (message instanceof MqttUnsubAckMessage) {
116                     return encodeUnsubAckMessage(ctx, (MqttUnsubAckMessage) message);
117                 }
118                 return encodeMessageWithOnlySingleByteFixedHeaderAndMessageId(ctx.alloc(), message);
119 
120             case PUBACK:
121             case PUBREC:
122             case PUBREL:
123             case PUBCOMP:
124                 return encodePubReplyMessage(ctx, message);
125 
126             case DISCONNECT:
127             case AUTH:
128                 return encodeReasonCodePlusPropertiesMessage(ctx, message);
129 
130             case PINGREQ:
131             case PINGRESP:
132                 return encodeMessageWithOnlySingleByteFixedHeader(ctx.alloc(), message);
133 
134             default:
135                 throw new IllegalArgumentException(
136                         "Unknown message type: " + message.fixedHeader().messageType().value());
137         }
138     }
139 
140     private static ByteBuf encodeConnectMessage(
141             ChannelHandlerContext ctx,
142             MqttConnectMessage message) {
143         int payloadBufferSize = 0;
144 
145         MqttFixedHeader mqttFixedHeader = message.fixedHeader();
146         MqttConnectVariableHeader variableHeader = message.variableHeader();
147         MqttConnectPayload payload = message.payload();
148         MqttVersion mqttVersion = MqttVersion.fromProtocolNameAndLevel(variableHeader.name(),
149                 (byte) variableHeader.version());
150         setMqttVersion(ctx, mqttVersion);
151 
152         // MQTT 3.1 and 3.1.1 require the Password Flag to be 0 when the User Name Flag is 0.
153         if ((mqttVersion == MqttVersion.MQTT_3_1 || mqttVersion == MqttVersion.MQTT_3_1_1) &&
154                 !variableHeader.hasUserName() && variableHeader.hasPassword()) {
155             throw new EncoderException("Without a username, the password MUST be not set");
156         }
157 
158         // Client id
159         String clientIdentifier = payload.clientIdentifier();
160         if (!isValidClientId(mqttVersion, DEFAULT_MAX_CLIENT_ID_LENGTH, clientIdentifier, false)) {
161             throw new MqttIdentifierRejectedException("invalid clientIdentifier");
162         }
163         int clientIdentifierBytes = utf8Bytes(clientIdentifier);
164         payloadBufferSize += 2 + clientIdentifierBytes;
165 
166         // Will topic and message
167         String willTopic = payload.willTopic();
168         int willTopicBytes = nullableUtf8Bytes(willTopic);
169         byte[] willMessage = payload.willMessageInBytes();
170         byte[] willMessageBytes = willMessage != null ? willMessage : EmptyArrays.EMPTY_BYTES;
171         if (variableHeader.isWillFlag()) {
172             if (!isValidPublishTopicName(willTopic)) {
173                 throw new MqttIdentifierRejectedException("invalid willTopic");
174             }
175             payloadBufferSize += 2 + willTopicBytes;
176             payloadBufferSize += 2 + willMessageBytes.length;
177         }
178 
179         String userName = payload.userName();
180         int userNameBytes = nullableUtf8Bytes(userName);
181         if (variableHeader.hasUserName()) {
182             if (!isValidUserName(userName)) {
183                 throw new MqttIdentifierRejectedException("invalid userName");
184             }
185             payloadBufferSize += 2 + userNameBytes;
186         }
187 
188         byte[] password = payload.passwordInBytes();
189         byte[] passwordBytes = password != null ? password : EmptyArrays.EMPTY_BYTES;
190         if (variableHeader.hasPassword()) {
191             payloadBufferSize += 2 + passwordBytes.length;
192         }
193 
194         // Fixed and variable header
195         byte[] protocolNameBytes = mqttVersion.protocolNameBytes();
196         ByteBuf propertiesBuf = encodePropertiesIfNeeded(
197                 mqttVersion,
198                 ctx.alloc(),
199                 message.variableHeader().properties());
200         try {
201             final ByteBuf willPropertiesBuf;
202             if (variableHeader.isWillFlag()) {
203                 willPropertiesBuf = encodePropertiesIfNeeded(mqttVersion, ctx.alloc(), payload.willProperties());
204                 payloadBufferSize += willPropertiesBuf.readableBytes();
205             } else {
206                 willPropertiesBuf = Unpooled.EMPTY_BUFFER;
207             }
208             try {
209                 int variableHeaderBufferSize = 2 + protocolNameBytes.length + 4 + propertiesBuf.readableBytes();
210 
211                 int variablePartSize = variableHeaderBufferSize + payloadBufferSize;
212                 int fixedHeaderBufferSize = 1 + getVariableLengthInt(variablePartSize);
213                 ByteBuf buf = ctx.alloc().buffer(fixedHeaderBufferSize + variablePartSize);
214                 buf.writeByte(getFixedHeaderByte1(mqttFixedHeader));
215                 writeVariableLengthInt(buf, variablePartSize);
216 
217                 buf.writeShort(protocolNameBytes.length);
218                 buf.writeBytes(protocolNameBytes);
219 
220                 buf.writeByte(variableHeader.version());
221                 buf.writeByte(getConnVariableHeaderFlag(variableHeader));
222                 buf.writeShort(variableHeader.keepAliveTimeSeconds());
223                 buf.writeBytes(propertiesBuf);
224 
225                 // Payload
226                 writeExactUTF8String(buf, clientIdentifier, clientIdentifierBytes);
227                 if (variableHeader.isWillFlag()) {
228                     buf.writeBytes(willPropertiesBuf);
229                     writeExactUTF8String(buf, willTopic, willTopicBytes);
230                     buf.writeShort(willMessageBytes.length);
231                     buf.writeBytes(willMessageBytes, 0, willMessageBytes.length);
232                 }
233                 if (variableHeader.hasUserName()) {
234                     writeExactUTF8String(buf, userName, userNameBytes);
235                 }
236                 if (variableHeader.hasPassword()) {
237                     buf.writeShort(passwordBytes.length);
238                     buf.writeBytes(passwordBytes, 0, passwordBytes.length);
239                 }
240                 return buf;
241             } finally {
242                 willPropertiesBuf.release();
243             }
244         } finally {
245             propertiesBuf.release();
246         }
247     }
248 
249     private static int getConnVariableHeaderFlag(MqttConnectVariableHeader variableHeader) {
250         int flagByte = 0;
251         if (variableHeader.hasUserName()) {
252             flagByte |= 0x80;
253         }
254         if (variableHeader.hasPassword()) {
255             flagByte |= 0x40;
256         }
257         if (variableHeader.isWillRetain()) {
258             flagByte |= 0x20;
259         }
260         flagByte |= (variableHeader.willQos() & 0x03) << 3;
261         if (variableHeader.isWillFlag()) {
262             flagByte |= 0x04;
263         }
264         if (variableHeader.isCleanSession()) {
265             flagByte |= 0x02;
266         }
267         return flagByte;
268     }
269 
270     private static ByteBuf encodeConnAckMessage(
271             ChannelHandlerContext ctx,
272             MqttConnAckMessage message) {
273         final MqttVersion mqttVersion = getMqttVersion(ctx);
274         ByteBuf propertiesBuf = encodePropertiesIfNeeded(mqttVersion,
275                 ctx.alloc(),
276                 message.variableHeader().properties());
277 
278         try {
279             ByteBuf buf = ctx.alloc().buffer(4 + propertiesBuf.readableBytes());
280             buf.writeByte(getFixedHeaderByte1(message.fixedHeader()));
281             writeVariableLengthInt(buf, 2 + propertiesBuf.readableBytes());
282             buf.writeByte(message.variableHeader().isSessionPresent() ? 0x01 : 0x00);
283             buf.writeByte(message.variableHeader().connectReturnCode().byteValue());
284             buf.writeBytes(propertiesBuf);
285             return buf;
286         } finally {
287             propertiesBuf.release();
288         }
289     }
290 
291     private static ByteBuf encodeSubscribeMessage(
292             ChannelHandlerContext ctx,
293             MqttSubscribeMessage message) {
294         MqttVersion mqttVersion = getMqttVersion(ctx);
295         ByteBuf propertiesBuf = encodePropertiesIfNeeded(mqttVersion,
296                 ctx.alloc(),
297                 message.idAndPropertiesVariableHeader().properties());
298 
299         try {
300             final int variableHeaderBufferSize = 2 + propertiesBuf.readableBytes();
301             int payloadBufferSize = 0;
302 
303             MqttFixedHeader mqttFixedHeader = message.fixedHeader();
304             MqttMessageIdVariableHeader variableHeader = message.variableHeader();
305             MqttSubscribePayload payload = message.payload();
306 
307             for (MqttTopicSubscription topic : payload.topicSubscriptions()) {
308                 String topicName = topic.topicName();
309                 int topicNameBytes = utf8Bytes(topicName);
310                 payloadBufferSize += 2 + topicNameBytes;
311                 payloadBufferSize += 1;
312             }
313 
314             int variablePartSize = variableHeaderBufferSize + payloadBufferSize;
315             int fixedHeaderBufferSize = 1 + getVariableLengthInt(variablePartSize);
316 
317             ByteBuf buf = ctx.alloc().buffer(fixedHeaderBufferSize + variablePartSize);
318             buf.writeByte(getFixedHeaderByte1(mqttFixedHeader));
319             writeVariableLengthInt(buf, variablePartSize);
320 
321             // Variable Header
322             int messageId = variableHeader.messageId();
323             buf.writeShort(messageId);
324             buf.writeBytes(propertiesBuf);
325 
326             // Payload
327             for (MqttTopicSubscription topic : payload.topicSubscriptions()) {
328                 writeEagerUTF8String(buf, topic.topicName());
329                 if (mqttVersion == MqttVersion.MQTT_3_1_1 || mqttVersion == MqttVersion.MQTT_3_1) {
330                     buf.writeByte(topic.qualityOfService().value());
331                 } else {
332                     final MqttSubscriptionOption option = topic.option();
333 
334                     int optionEncoded = option.retainHandling().value() << 4;
335                     if (option.isRetainAsPublished()) {
336                         optionEncoded |= 0x08;
337                     }
338                     if (option.isNoLocal()) {
339                         optionEncoded |= 0x04;
340                     }
341                     optionEncoded |= option.qos().value();
342 
343                     buf.writeByte(optionEncoded);
344                 }
345             }
346 
347             return buf;
348         } finally {
349             propertiesBuf.release();
350         }
351     }
352 
353     private static ByteBuf encodeUnsubscribeMessage(
354             ChannelHandlerContext ctx,
355             MqttUnsubscribeMessage message) {
356         MqttVersion mqttVersion = getMqttVersion(ctx);
357         ByteBuf propertiesBuf = encodePropertiesIfNeeded(mqttVersion,
358                 ctx.alloc(),
359                 message.idAndPropertiesVariableHeader().properties());
360 
361         try {
362             final int variableHeaderBufferSize = 2 + propertiesBuf.readableBytes();
363             int payloadBufferSize = 0;
364 
365             MqttFixedHeader mqttFixedHeader = message.fixedHeader();
366             MqttMessageIdVariableHeader variableHeader = message.variableHeader();
367             MqttUnsubscribePayload payload = message.payload();
368 
369             for (String topicName : payload.topics()) {
370                 int topicNameBytes = utf8Bytes(topicName);
371                 payloadBufferSize += 2 + topicNameBytes;
372             }
373 
374             int variablePartSize = variableHeaderBufferSize + payloadBufferSize;
375             int fixedHeaderBufferSize = 1 + getVariableLengthInt(variablePartSize);
376 
377             ByteBuf buf = ctx.alloc().buffer(fixedHeaderBufferSize + variablePartSize);
378             buf.writeByte(getFixedHeaderByte1(mqttFixedHeader));
379             writeVariableLengthInt(buf, variablePartSize);
380 
381             // Variable Header
382             int messageId = variableHeader.messageId();
383             buf.writeShort(messageId);
384             buf.writeBytes(propertiesBuf);
385 
386             // Payload
387             for (String topicName : payload.topics()) {
388                 writeEagerUTF8String(buf, topicName);
389             }
390 
391             return buf;
392         } finally {
393             propertiesBuf.release();
394         }
395     }
396 
397     private static ByteBuf encodeSubAckMessage(
398             ChannelHandlerContext ctx,
399             MqttSubAckMessage message) {
400         MqttVersion mqttVersion = getMqttVersion(ctx);
401         ByteBuf propertiesBuf = encodePropertiesIfNeeded(mqttVersion,
402                 ctx.alloc(),
403                 message.idAndPropertiesVariableHeader().properties());
404         try {
405             int variableHeaderBufferSize = 2 + propertiesBuf.readableBytes();
406             int payloadBufferSize = message.payload().grantedQoSLevels().size();
407             int variablePartSize = variableHeaderBufferSize + payloadBufferSize;
408             int fixedHeaderBufferSize = 1 + getVariableLengthInt(variablePartSize);
409             ByteBuf buf = ctx.alloc().buffer(fixedHeaderBufferSize + variablePartSize);
410             buf.writeByte(getFixedHeaderByte1(message.fixedHeader()));
411             writeVariableLengthInt(buf, variablePartSize);
412             buf.writeShort(message.variableHeader().messageId());
413             buf.writeBytes(propertiesBuf);
414             for (int code: message.payload().reasonCodes()) {
415                 buf.writeByte(code);
416             }
417 
418             return buf;
419         } finally {
420             propertiesBuf.release();
421         }
422     }
423 
424     private static ByteBuf encodeUnsubAckMessage(
425             ChannelHandlerContext ctx,
426             MqttUnsubAckMessage message) {
427         if (message.variableHeader() instanceof  MqttMessageIdAndPropertiesVariableHeader) {
428             MqttVersion mqttVersion = getMqttVersion(ctx);
429             // Reason Codes were introduced in MQTT 5.0 only. MQTT 3.1.1 (and 3.1) UNSUBACK packets
430             // have no payload, so reason codes must be suppressed for older protocol versions
431             // even when the caller populated them via MqttMessageBuilders.
432             final boolean writeReasonCodes = mqttVersion == MqttVersion.MQTT_5;
433             ByteBuf propertiesBuf = encodePropertiesIfNeeded(mqttVersion,
434                     ctx.alloc(),
435                     message.idAndPropertiesVariableHeader().properties());
436             try {
437                 int variableHeaderBufferSize = 2 + propertiesBuf.readableBytes();
438                 MqttUnsubAckPayload payload = message.payload();
439                 int payloadBufferSize = writeReasonCodes && payload != null
440                         ? payload.unsubscribeReasonCodes().size() : 0;
441                 int variablePartSize = variableHeaderBufferSize + payloadBufferSize;
442                 int fixedHeaderBufferSize = 1 + getVariableLengthInt(variablePartSize);
443                 ByteBuf buf = ctx.alloc().buffer(fixedHeaderBufferSize + variablePartSize);
444                 buf.writeByte(getFixedHeaderByte1(message.fixedHeader()));
445                 writeVariableLengthInt(buf, variablePartSize);
446                 buf.writeShort(message.variableHeader().messageId());
447                 buf.writeBytes(propertiesBuf);
448 
449                 if (writeReasonCodes && payload != null) {
450                     for (Short reasonCode : payload.unsubscribeReasonCodes()) {
451                         buf.writeByte(reasonCode);
452                     }
453                 }
454 
455                 return buf;
456             } finally {
457                 propertiesBuf.release();
458             }
459         } else {
460             return encodeMessageWithOnlySingleByteFixedHeaderAndMessageId(ctx.alloc(), message);
461         }
462     }
463 
464     private static ByteBuf encodePublishMessage(
465             ChannelHandlerContext ctx,
466             MqttPublishMessage message) {
467         MqttVersion mqttVersion = getMqttVersion(ctx);
468         MqttFixedHeader mqttFixedHeader = message.fixedHeader();
469         MqttPublishVariableHeader variableHeader = message.variableHeader();
470         ByteBuf payload = message.payload().duplicate();
471 
472         String topicName = variableHeader.topicName();
473         if (!isValidPublishTopicName(topicName)) {
474             throw new MqttIdentifierRejectedException("invalid topicName");
475         }
476         int topicNameBytes = utf8Bytes(topicName);
477 
478         ByteBuf propertiesBuf = encodePropertiesIfNeeded(mqttVersion,
479                 ctx.alloc(),
480                 message.variableHeader().properties());
481 
482         try {
483             boolean qosLevelGreaterZero = mqttFixedHeader.qosLevel().value() > 0;
484             int variableHeaderBufferSize = 2 + topicNameBytes +
485                     (qosLevelGreaterZero ? 2 : 0) + propertiesBuf.readableBytes();
486             int payloadBufferSize = payload.readableBytes();
487             int variablePartSize = variableHeaderBufferSize + payloadBufferSize;
488             int fixedHeaderBufferSize = 1 + getVariableLengthInt(variablePartSize);
489 
490             ByteBuf buf = ctx.alloc().buffer(fixedHeaderBufferSize + variablePartSize);
491             buf.writeByte(getFixedHeaderByte1(mqttFixedHeader));
492             writeVariableLengthInt(buf, variablePartSize);
493             writeExactUTF8String(buf, topicName, topicNameBytes);
494             if (qosLevelGreaterZero) {
495                 buf.writeShort(variableHeader.packetId());
496             }
497             buf.writeBytes(propertiesBuf);
498             buf.writeBytes(payload);
499 
500             return buf;
501         } finally {
502             propertiesBuf.release();
503         }
504     }
505 
506     private static ByteBuf encodePubReplyMessage(ChannelHandlerContext ctx,
507                                           MqttMessage message) {
508         if (message.variableHeader() instanceof MqttPubReplyMessageVariableHeader) {
509             MqttFixedHeader mqttFixedHeader = message.fixedHeader();
510             MqttPubReplyMessageVariableHeader variableHeader =
511                     (MqttPubReplyMessageVariableHeader) message.variableHeader();
512             int msgId = variableHeader.messageId();
513 
514             final ByteBuf propertiesBuf;
515             final boolean includeReasonCode;
516             final int variableHeaderBufferSize;
517             final MqttVersion mqttVersion = getMqttVersion(ctx);
518             if (mqttVersion == MqttVersion.MQTT_5 &&
519                     (variableHeader.reasonCode() != MqttPubReplyMessageVariableHeader.REASON_CODE_OK ||
520                             !variableHeader.properties().isEmpty())) {
521                 propertiesBuf = encodeProperties(ctx.alloc(), variableHeader.properties());
522                 includeReasonCode = true;
523                 variableHeaderBufferSize = 3 + propertiesBuf.readableBytes();
524             } else {
525                 propertiesBuf = Unpooled.EMPTY_BUFFER;
526                 includeReasonCode = false;
527                 variableHeaderBufferSize = 2;
528             }
529 
530             try {
531                 final int fixedHeaderBufferSize = 1 + getVariableLengthInt(variableHeaderBufferSize);
532                 ByteBuf buf = ctx.alloc().buffer(fixedHeaderBufferSize + variableHeaderBufferSize);
533                 buf.writeByte(getFixedHeaderByte1(mqttFixedHeader));
534                 writeVariableLengthInt(buf, variableHeaderBufferSize);
535                 buf.writeShort(msgId);
536                 if (includeReasonCode) {
537                     buf.writeByte(variableHeader.reasonCode());
538                 }
539                 buf.writeBytes(propertiesBuf);
540 
541                 return buf;
542             } finally {
543                 propertiesBuf.release();
544             }
545         } else {
546             return encodeMessageWithOnlySingleByteFixedHeaderAndMessageId(ctx.alloc(), message);
547         }
548     }
549 
550     private static ByteBuf encodeMessageWithOnlySingleByteFixedHeaderAndMessageId(
551             ByteBufAllocator byteBufAllocator,
552             MqttMessage message) {
553         MqttFixedHeader mqttFixedHeader = message.fixedHeader();
554         MqttMessageIdVariableHeader variableHeader = (MqttMessageIdVariableHeader) message.variableHeader();
555         int msgId = variableHeader.messageId();
556 
557         int variableHeaderBufferSize = 2; // variable part only has a message id
558         int fixedHeaderBufferSize = 1 + getVariableLengthInt(variableHeaderBufferSize);
559         ByteBuf buf = byteBufAllocator.buffer(fixedHeaderBufferSize + variableHeaderBufferSize);
560         buf.writeByte(getFixedHeaderByte1(mqttFixedHeader));
561         writeVariableLengthInt(buf, variableHeaderBufferSize);
562         buf.writeShort(msgId);
563 
564         return buf;
565     }
566 
567     private static ByteBuf encodeReasonCodePlusPropertiesMessage(
568             ChannelHandlerContext ctx,
569             MqttMessage message) {
570         if (message.variableHeader() instanceof MqttReasonCodeAndPropertiesVariableHeader) {
571             MqttVersion mqttVersion = getMqttVersion(ctx);
572             MqttFixedHeader mqttFixedHeader = message.fixedHeader();
573             MqttReasonCodeAndPropertiesVariableHeader variableHeader =
574                     (MqttReasonCodeAndPropertiesVariableHeader) message.variableHeader();
575 
576             final ByteBuf propertiesBuf;
577             final boolean includeReasonCode;
578             final int variableHeaderBufferSize;
579             if (mqttVersion == MqttVersion.MQTT_5 &&
580                     (variableHeader.reasonCode() != MqttReasonCodeAndPropertiesVariableHeader.REASON_CODE_OK ||
581                             !variableHeader.properties().isEmpty())) {
582                 propertiesBuf = encodeProperties(ctx.alloc(), variableHeader.properties());
583                 includeReasonCode = true;
584                 variableHeaderBufferSize = 1 + propertiesBuf.readableBytes();
585             } else {
586                 propertiesBuf = Unpooled.EMPTY_BUFFER;
587                 includeReasonCode = false;
588                 variableHeaderBufferSize = 0;
589             }
590 
591             try {
592                 final int fixedHeaderBufferSize = 1 + getVariableLengthInt(variableHeaderBufferSize);
593                 ByteBuf buf = ctx.alloc().buffer(fixedHeaderBufferSize + variableHeaderBufferSize);
594                 buf.writeByte(getFixedHeaderByte1(mqttFixedHeader));
595                 writeVariableLengthInt(buf, variableHeaderBufferSize);
596                 if (includeReasonCode) {
597                     buf.writeByte(variableHeader.reasonCode());
598                 }
599                 buf.writeBytes(propertiesBuf);
600 
601                 return buf;
602             } finally {
603                 propertiesBuf.release();
604             }
605         } else {
606             return encodeMessageWithOnlySingleByteFixedHeader(ctx.alloc(), message);
607         }
608     }
609 
610     private static ByteBuf encodeMessageWithOnlySingleByteFixedHeader(
611             ByteBufAllocator byteBufAllocator,
612             MqttMessage message) {
613         MqttFixedHeader mqttFixedHeader = message.fixedHeader();
614         ByteBuf buf = byteBufAllocator.buffer(2);
615         buf.writeByte(getFixedHeaderByte1(mqttFixedHeader));
616         buf.writeByte(0);
617 
618         return buf;
619     }
620 
621     private static ByteBuf encodePropertiesIfNeeded(MqttVersion mqttVersion,
622                                              ByteBufAllocator byteBufAllocator,
623                                              MqttProperties mqttProperties) {
624         if (mqttVersion == MqttVersion.MQTT_5) {
625             return encodeProperties(byteBufAllocator, mqttProperties);
626         }
627         return Unpooled.EMPTY_BUFFER;
628     }
629 
630     private static ByteBuf encodeProperties(ByteBufAllocator byteBufAllocator,
631                                             MqttProperties mqttProperties) {
632         ByteBuf propertiesHeaderBuf = byteBufAllocator.buffer();
633         // encode also the Properties part
634         try {
635             ByteBuf propertiesBuf = byteBufAllocator.buffer();
636             try {
637                 for (MqttProperties.MqttProperty property : mqttProperties.listAll()) {
638                     int propertyId = property.propertyId;
639                     switch (propertyId) {
640                         case PAYLOAD_FORMAT_INDICATOR:
641                         case REQUEST_PROBLEM_INFORMATION:
642                         case REQUEST_RESPONSE_INFORMATION:
643                         case MAXIMUM_QOS:
644                         case RETAIN_AVAILABLE:
645                         case WILDCARD_SUBSCRIPTION_AVAILABLE:
646                         case SUBSCRIPTION_IDENTIFIER_AVAILABLE:
647                         case SHARED_SUBSCRIPTION_AVAILABLE:
648                             writeVariableLengthInt(propertiesBuf, propertyId);
649                             final byte bytePropValue = ((MqttProperties.IntegerProperty) property).value.byteValue();
650                             propertiesBuf.writeByte(bytePropValue);
651                             break;
652                         case SERVER_KEEP_ALIVE:
653                         case RECEIVE_MAXIMUM:
654                         case TOPIC_ALIAS_MAXIMUM:
655                         case TOPIC_ALIAS:
656                             writeVariableLengthInt(propertiesBuf, propertyId);
657                             final short twoBytesInPropValue =
658                                     ((MqttProperties.IntegerProperty) property).value.shortValue();
659                             propertiesBuf.writeShort(twoBytesInPropValue);
660                             break;
661                         case PUBLICATION_EXPIRY_INTERVAL:
662                         case SESSION_EXPIRY_INTERVAL:
663                         case WILL_DELAY_INTERVAL:
664                         case MAXIMUM_PACKET_SIZE:
665                             writeVariableLengthInt(propertiesBuf, propertyId);
666                             final int fourBytesIntPropValue = ((MqttProperties.IntegerProperty) property).value;
667                             propertiesBuf.writeInt(fourBytesIntPropValue);
668                             break;
669                         case SUBSCRIPTION_IDENTIFIER:
670                             writeVariableLengthInt(propertiesBuf, propertyId);
671                             final int vbi = ((MqttProperties.IntegerProperty) property).value;
672                             writeVariableLengthInt(propertiesBuf, vbi);
673                             break;
674                         case CONTENT_TYPE:
675                         case RESPONSE_TOPIC:
676                         case ASSIGNED_CLIENT_IDENTIFIER:
677                         case AUTHENTICATION_METHOD:
678                         case RESPONSE_INFORMATION:
679                         case SERVER_REFERENCE:
680                         case REASON_STRING:
681                             writeVariableLengthInt(propertiesBuf, propertyId);
682                             writeEagerUTF8String(propertiesBuf, ((MqttProperties.StringProperty) property).value);
683                             break;
684                         case USER_PROPERTY:
685                             final List<MqttProperties.StringPair> pairs =
686                                     ((MqttProperties.UserProperties) property).value;
687                             for (MqttProperties.StringPair pair : pairs) {
688                                 writeVariableLengthInt(propertiesBuf, propertyId);
689                                 writeEagerUTF8String(propertiesBuf, pair.key);
690                                 writeEagerUTF8String(propertiesBuf, pair.value);
691                             }
692                             break;
693                         case CORRELATION_DATA:
694                         case AUTHENTICATION_DATA:
695                             writeVariableLengthInt(propertiesBuf, propertyId);
696                             final byte[] binaryPropValue = ((MqttProperties.BinaryProperty) property).value;
697                             propertiesBuf.writeShort(binaryPropValue.length);
698                             propertiesBuf.writeBytes(binaryPropValue, 0, binaryPropValue.length);
699                             break;
700                         default:
701                             //shouldn't reach here
702                             throw new EncoderException("Unknown property type: " + propertyId);
703                     }
704                 }
705                 writeVariableLengthInt(propertiesHeaderBuf, propertiesBuf.readableBytes());
706                 propertiesHeaderBuf.writeBytes(propertiesBuf);
707 
708                 return propertiesHeaderBuf;
709             } finally {
710                 propertiesBuf.release();
711             }
712         } catch (RuntimeException e) {
713             propertiesHeaderBuf.release();
714             throw e;
715         }
716     }
717 
718     private static int getFixedHeaderByte1(MqttFixedHeader header) {
719         int ret = 0;
720         ret |= header.messageType().value() << 4;
721         if (header.isDup()) {
722             ret |= 0x08;
723         }
724         ret |= header.qosLevel().value() << 1;
725         if (header.isRetain()) {
726             ret |= 0x01;
727         }
728         return ret;
729     }
730 
731     private static void writeVariableLengthInt(ByteBuf buf, int num) {
732         do {
733             int digit = num & 0x7F;
734             num >>>= 7;
735             if (num > 0) {
736                 digit |= 0x80;
737             }
738             buf.writeByte(digit);
739         } while (num > 0);
740     }
741 
742     private static int nullableUtf8Bytes(String s) {
743         return s == null? 0 : utf8Bytes(s);
744     }
745 
746     private static int nullableMaxUtf8Bytes(String s) {
747         return s == null? 0 : utf8MaxBytes(s);
748     }
749 
750     private static void writeExactUTF8String(ByteBuf buf, String s, int utf8Length) {
751         buf.ensureWritable(utf8Length + 2);
752         buf.writeShort(utf8Length);
753         if (utf8Length > 0) {
754             final int writtenUtf8Length = reserveAndWriteUtf8(buf, s, utf8Length);
755             assert writtenUtf8Length == utf8Length;
756         }
757     }
758 
759     private static void writeEagerUTF8String(ByteBuf buf, String s) {
760         final int maxUtf8Length = nullableMaxUtf8Bytes(s);
761         buf.ensureWritable(maxUtf8Length + 2);
762         final int writerIndex = buf.writerIndex();
763         final int startUtf8String = writerIndex + 2;
764         buf.writerIndex(startUtf8String);
765         final int utf8Length = s != null? reserveAndWriteUtf8(buf, s, maxUtf8Length) : 0;
766         buf.setShort(writerIndex, utf8Length);
767     }
768 
769     private static int getVariableLengthInt(int num) {
770         if (num < 128) {
771             return 1;
772         }
773         if (num < 16_384) { // 128 * 128
774             return 2;
775         }
776         if (num < 2_097_152) { // 128 * 128 * 128
777             return 3;
778         }
779         return 4;
780     }
781 
782 }